退役 DirectProject 工具调用与回合流账本
Project CI / AI game creator shell Rust lane 2/2 (pull_request) Failing after 1m11s
Project CI / AI game creator shell Rust lane 1/2 (pull_request) Failing after 1m11s
Project CI / AI game creator shell Rust smoke (pull_request) Failing after 1m11s
Project CI / Backend tests (pull_request) Failing after 26s
Project CI / AI game creator shell Rust crates (pull_request) Successful in 1m42s
Project CI / Repository checks (pull_request) Failing after 27s
Project CI / Frontend tests (pull_request) Successful in 2m56s
Project CI / AI game creator shell web tests (pull_request) Successful in 2m42s
Project CI / Native shell tests (pull_request) Successful in 7m0s

- agent/direct_tool_calls.rs、agent/direct_turn_stream.rs:整体删除,agent.rs 去掉模块声明与导出
- direct_runtime/mod.rs:删 DirectToolCallCollector、回合流节流器与全部落盘/发射调用点,observer 只回填 update_active_turn
- codex_app_server/mod.rs:删 DirectCodexTurnObservation::ToolCall,AgentMessageSegment 结构体变体简化为元组变体,direct_tool_call_turn_id / direct_tool_call_now_ms 改名
- claude_code_cli.rs:删 tool_use / tool_result 采集分支与 DirectToolCall 构造
- main.rs、runtime_driver/entrypoints.rs:删 GameCreatorDirectTurnUpdateEvent 与 game-creator-direct-turn-update 事件,emit 收缩为 (status, activity)
- thread_manager/wire.rs:新增 direct_now_ms,替代随 direct_tool_calls 删除的时钟函数
- docs:ADR、工具调用卡片技术方案、DirectProject 聊天真相源收敛里程碑与实施计划、pitfalls 旧路径同步到当前状态,decision-log 新增 2026-09-30 退役条目
- 未修(另开):redact_absolute_path_tokens 的无语境绝对路径扫描仍会误伤 HTML 结束标签与嵌套 JSON 转义
- 验证:cargo test --bin genarrative-ai-game-creator-shell agent::(929 passed / 0 failed / 5 ignored)、npm run check:encoding、npm run check:doc-index、git diff --check
This commit is contained in:
2026-09-30 19:32:05 +08:00
parent 9ee34809be
commit 719db293d2
18 changed files with 328 additions and 2545 deletions
@@ -30,11 +30,9 @@ mod direct_project_history;
mod direct_project_turn_history;
mod direct_runtime;
mod direct_tool_bridge;
mod direct_tool_calls;
mod direct_tools_mcp;
mod direct_turn_error;
mod direct_turn_failure;
mod direct_turn_stream;
mod direct_validation;
mod generation;
mod prompt;
@@ -71,11 +69,9 @@ pub(crate) use direct_project_history::*;
pub(crate) use direct_project_turn_history::*;
pub(crate) use direct_runtime::*;
pub(crate) use direct_tool_bridge::*;
pub(crate) use direct_tool_calls::*;
pub(crate) use direct_tools_mcp::*;
pub(crate) use direct_turn_error::*;
pub(crate) use direct_turn_failure::*;
pub(crate) use direct_turn_stream::*;
pub(crate) use direct_validation::DirectValidationConfig;
pub(crate) use generation::*;
pub(crate) use prompt::*;
@@ -709,32 +709,21 @@ pub(crate) async fn direct_game_creator_claude_code_chat_at(
.filter_map(|event| serde_json::to_string(&event).ok())
.collect::<Vec<_>>();
stream.push(String::new());
parse_direct_stream_result(
stream.join("\n").as_bytes(),
client_turn_id.unwrap_or("claude-code-turn"),
observer,
)
parse_direct_stream_result(stream.join("\n").as_bytes(), observer)
}
fn parse_direct_stream_result(
stdout: &[u8],
turn_id: &str,
mut observer: Option<&mut (dyn FnMut(DirectCodexTurnObservation) + Send)>,
) -> Result<String, String> {
let text = std::str::from_utf8(stdout)
.map_err(|_| "Claude Agent SDK stream-json 不是 UTF-8".to_string())?;
let mut result = None;
let mut tool_calls = std::collections::HashMap::<String, DirectToolCall>::new();
for line in text.lines().filter(|line| !line.trim().is_empty()) {
let event: serde_json::Value = serde_json::from_str(line)
.map_err(|_| "Claude Agent SDK stream-json 包含无效 JSON".to_string())?;
let event_type = event.get("type").and_then(serde_json::Value::as_str);
if event_type == Some("assistant") {
let id = event
.get("uuid")
.and_then(serde_json::Value::as_str)
.unwrap_or("claude-code-assistant")
.to_string();
let content = event
.pointer("/message/content")
.and_then(serde_json::Value::as_array);
@@ -753,77 +742,7 @@ fn parse_direct_stream_result(
.collect::<String>();
if !visible.is_empty() {
if let Some(observer) = observer.as_deref_mut() {
observer(DirectCodexTurnObservation::AgentMessageSegment {
item_id: id,
accumulated_text: visible,
completed: true,
});
}
}
for part in content.iter().filter(|part| {
part.get("type").and_then(serde_json::Value::as_str) == Some("tool_use")
}) {
let Some(call_id) = part.get("id").and_then(serde_json::Value::as_str) else {
continue;
};
let name = part
.get("name")
.and_then(serde_json::Value::as_str)
.unwrap_or("mcp_tool")
.to_string();
let arguments = part
.get("input")
.cloned()
.unwrap_or(serde_json::Value::Null);
let now = direct_tool_call_now_ms();
let call = DirectToolCall {
schema_version: DIRECT_TOOL_CALL_SCHEMA_VERSION.to_string(),
id: call_id.to_string(),
turn_id: turn_id.to_string(),
kind: "mcp_tool".to_string(),
title: "调用工具".to_string(),
summary: name,
status: "running".to_string(),
detail: DirectToolCallDetail {
command: serde_json::to_string(&arguments).ok(),
output: None,
changes: Vec::new(),
},
started_at: now,
updated_at: now,
};
if let Some(observer) = observer.as_deref_mut() {
observer(DirectCodexTurnObservation::ToolCall(call.clone()));
}
tool_calls.insert(call_id.to_string(), call);
}
}
} else if event_type == Some("user") {
if let Some(content) = event
.pointer("/message/content")
.and_then(serde_json::Value::as_array)
{
for part in content.iter().filter(|part| {
part.get("type").and_then(serde_json::Value::as_str) == Some("tool_result")
}) {
let Some(call_id) = part.get("tool_use_id").and_then(serde_json::Value::as_str)
else {
continue;
};
let Some(mut call) = tool_calls.remove(call_id) else {
continue;
};
call.status = "completed".to_string();
call.updated_at = direct_tool_call_now_ms();
call.detail.output = part
.get("content")
.and_then(serde_json::Value::as_str)
.map(str::to_string);
if part.get("is_error").and_then(serde_json::Value::as_bool) == Some(true) {
call.status = "failed".to_string();
}
if let Some(observer) = observer.as_deref_mut() {
observer(DirectCodexTurnObservation::ToolCall(call));
observer(DirectCodexTurnObservation::AgentMessageSegment(visible));
}
}
}
@@ -893,14 +812,13 @@ mod tests {
{"type":"result","is_error":false,"result":"最终回复"}
"#
.as_bytes(),
"turn-1",
Some(&mut |event| observed.push(event)),
)
.expect("stream result");
assert_eq!(text, "最终回复");
assert!(matches!(
observed.as_slice(),
[DirectCodexTurnObservation::AgentMessageSegment { .. }]
[DirectCodexTurnObservation::AgentMessageSegment(_)]
));
}
}
@@ -642,21 +642,12 @@ enum CodexTurnEvent {
#[derive(Clone, Debug, Eq, PartialEq)]
pub(crate) enum DirectCodexTurnObservation {
AccumulatedText(String),
/// 一个 assistant 文本段的当前累计全文。
///
/// `item_id` 是一次 assistant 消息的稳定身份:同一个 id 的后续 delta 属于**同一段**,
/// id 变了就是新的一段。回合流的"文本段 + 工具"顺序用它来分段,而不是按 delta 分。
AgentMessageSegment {
item_id: String,
accumulated_text: String,
completed: bool,
},
/// 一段可见的 assistant 正文(同一 assistant item 的当前累计全文)。
AgentMessageSegment(String),
IntermediateText(String),
/// 模型的思考过程(reasoning item 的明文摘要):流式阶段整段替换下发。
Reasoning(String),
Activity(&'static str),
/// 一条结构化工具调用(`item/started` 与 `item/completed` 各采一次,按 id 幂等)。
ToolCall(crate::DirectToolCall),
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
@@ -809,7 +800,7 @@ fn direct_thread_event_item(
root: &std::path::Path,
item: &serde_json::Value,
) -> Option<ThreadItem> {
thread_item_from_value(root, item, direct_tool_call_now_ms())
thread_item_from_value(root, item, direct_now_ms())
}
/// 运行态条目投影:Codex 回显的用户消息整条跳过。
@@ -3344,9 +3335,9 @@ impl CodexAppServerConnection {
) -> Result<platform_llm::LlmRunResponse, DirectTurnRunFailure> {
let _turn_guard = self.inner.turn_gate.lock().await;
let history_root = direct_history_root.unwrap_or(&self.inner.workspace_path);
// 工具调用卡片的 turnId 用 AGC 客户端回合 id(与实时事件、落盘条目同一口径),
// 不用 Codex app-server 自己的 turnId——前端要按它把卡片挂回对应的那一轮。
let direct_tool_call_turn_id: Option<String> = direct_client_turn_id
// 回合身份用 AGC 客户端回合 id(与实时事件、落盘条目同一口径),
// 不用 Codex app-server 自己的 turnId。
let direct_turn_id: Option<String> = direct_client_turn_id
.map(str::trim)
.filter(|turn_id| !turn_id.is_empty())
.map(str::to_string);
@@ -3525,7 +3516,7 @@ impl CodexAppServerConnection {
}
let context = ProjectModelUsageContext {
root: history_root.to_path_buf(),
client_turn_id: direct_tool_call_turn_id.clone(),
client_turn_id: direct_turn_id.clone(),
thread_id: Some(thread_id.clone()),
requested_model: model.to_string(),
};
@@ -3546,7 +3537,7 @@ impl CodexAppServerConnection {
let turn_start_cancellation =
Arc::new(CodexTurnStartCancellation::new(&self.inner, &thread_id));
// Direct 回合登记为"可终止":终止命令只作用在这一轮上,回合结束时自动注销。
let _active_turn_guard = direct_tool_call_turn_id
let _active_turn_guard = direct_turn_id
.as_deref()
.filter(|_| self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject)
.map(|turn_id| {
@@ -3630,7 +3621,7 @@ impl CodexAppServerConnection {
// `startedAt` / `completedAt` 只有秒级,秒级截断撑不起前端 0.1 秒粒度的展示,也可能
// 让完成时刻落进该轮用户消息的同一秒。因此这里在进入模型往返前取一次宿主毫秒钟,与
// `durationMs` 相加得到终态时刻;拿不到 `durationMs` 时退回观察时刻。
let direct_turn_started_at_ms = direct_tool_call_now_ms();
let direct_turn_started_at_ms = direct_now_ms();
let mut receiver = self.register_turn(&turn_id).await;
let mut direct_project_history = DirectProjectHistoryAccumulator::default();
let mut guard = CodexTurnGuard {
@@ -3730,17 +3721,14 @@ impl CodexAppServerConnection {
observer(DirectCodexTurnObservation::AccumulatedText(
streamed_text.clone(),
));
// 同一个 assistant item 的当前累计全文:回合流按 item 分段,
// 段内只追加、段间才换行,不能拿"整轮累计"当一段。
// 同一 assistant item 的当前累计全文,不能拿"整轮累计"当一段。
let segment_text = direct_project_history
.accumulated_text_for(&item_id)
.unwrap_or_else(|| delta.clone());
if !segment_text.trim().is_empty() {
observer(DirectCodexTurnObservation::AgentMessageSegment {
item_id: item_id.clone(),
accumulated_text: segment_text,
completed: false,
});
observer(DirectCodexTurnObservation::AgentMessageSegment(
segment_text,
));
}
}
if let Some(callback) = on_agent_message_delta.as_deref_mut() {
@@ -3799,10 +3787,7 @@ impl CodexAppServerConnection {
// 通知的钟就是该阶段唯一可证明的时间。
append_thread_event(
&direct_thread_id,
ThreadEvent::item_completed(
entry_item,
direct_tool_call_now_ms(),
),
ThreadEvent::item_completed(entry_item, direct_now_ms()),
);
}
}
@@ -3860,43 +3845,22 @@ impl CodexAppServerConnection {
completed,
&params,
);
// 工具调用卡片:item/started 与 item/completed 各采一次,
// 由下游按 id 幂等 upsert 成同一条。采集失败(拿不到 id /
// 非工具类 item)就静默跳过,不影响这一轮的其它投影。
if let Some(turn_id) = direct_tool_call_turn_id.as_deref() {
if let Some(tool_call) = direct_tool_call_from_item(
history_root,
item,
turn_id,
completed,
direct_tool_call_now_ms(),
) {
if let Some(observer) = direct_observer.as_deref_mut() {
observer(DirectCodexTurnObservation::ToolCall(
tool_call,
));
}
}
}
}
if item_type == "agentMessage" {
// 某些 app-server 实现会在工具开始后停止发送 agentMessage delta,
// 但会在 item/completed 携带完整文本。把这份最终快照补进回合流,
// 让流中的文本段不会停在工具前的短前缀。
// 但会在 item/completed 携带完整文本;把这份最终快照也交给观察者,
// 运行态条目就不会停在工具前的短前缀。
if completed {
if let (Some(item_id), Some(text)) = (
item.get("id").and_then(serde_json::Value::as_str),
item.get("text")
.and_then(serde_json::Value::as_str)
.filter(|value| !value.trim().is_empty()),
) {
if let Some(text) = item
.get("text")
.and_then(serde_json::Value::as_str)
.filter(|value| !value.trim().is_empty())
{
if let Some(observer) = direct_observer.as_deref_mut() {
observer(
DirectCodexTurnObservation::AgentMessageSegment {
item_id: item_id.to_string(),
accumulated_text: text.to_string(),
completed: true,
},
DirectCodexTurnObservation::AgentMessageSegment(
text.to_string(),
),
);
}
}
@@ -3938,7 +3902,7 @@ impl CodexAppServerConnection {
&params,
item,
false,
direct_tool_call_now_ms(),
direct_now_ms(),
),
),
);
@@ -3962,15 +3926,10 @@ impl CodexAppServerConnection {
.filter(|text| !text.trim().is_empty())
{
final_text = Some(text.to_string());
if let (Some(item_id), Some(observer)) = (
item.get("id").and_then(serde_json::Value::as_str),
direct_observer.as_deref_mut(),
) {
observer(DirectCodexTurnObservation::AgentMessageSegment {
item_id: item_id.to_string(),
accumulated_text: text.to_string(),
completed: true,
});
if let Some(observer) = direct_observer.as_deref_mut() {
observer(DirectCodexTurnObservation::AgentMessageSegment(
text.to_string(),
));
}
}
}
@@ -3987,7 +3946,7 @@ impl CodexAppServerConnection {
thread_turn_completed_at_ms(
turn,
Some(direct_turn_started_at_ms),
direct_tool_call_now_ms(),
direct_now_ms(),
),
));
}
@@ -4157,12 +4116,12 @@ impl CodexAppServerConnection {
.unwrap_or_else(|| fallback_status.to_string());
// 有执行许可时,起止时间包含实际宿主收尾;上游模型完成不能提前结束 UI。
let completed_at = if approval_adapter.is_some() {
direct_tool_call_now_ms()
direct_now_ms()
} else {
model_terminal
.as_ref()
.map(|(_, at)| *at)
.unwrap_or_else(direct_tool_call_now_ms)
.unwrap_or_else(direct_now_ms)
};
// 终态判定的**事实**在这里固定,写点留到整轮真正结束之后(见下面的
// `turn_result`):终态只有 `turn.completed` 一种事件,失败时同一个事件带 `failure`
File diff suppressed because it is too large Load Diff
@@ -102,7 +102,7 @@ async fn enqueue_direct_codex_turn_typed(
})?;
// 入队:到这里这一条已经过了全部检查,剩下的就是排队等放行。条目只带走它自己的事实
// (用户条目、创建类型、入队时刻),canonical 形状与 prompt 放行时从它重投影——放行没有失败出口。
let pending = PendingTurn::new(turn_id, user_item, creation_type, direct_tool_call_now_ms());
let pending = PendingTurn::new(turn_id, user_item, creation_type, direct_now_ms());
match enqueue_pending_turn(&thread_id, pending) {
Ok(_) => {}
Err(EnqueueRejection::QueueFull) => {
File diff suppressed because it is too large Load Diff
@@ -1,430 +0,0 @@
//! GameAgent 对话「回合流」的采集与持久化。
//!
//! 顺序真相放在一处:`<projectRoot>/.agent/conversations/turn-stream.jsonl` 按**出现顺序**
//! 记录一个回合里的文本段与工具调用。工具条目只记位置标记(`callId`),工具本身的正文
//! 仍然来自 `tool-calls.jsonl`(同一 id 幂等合并只有一处实现)。
//!
//! 位置稳定:每条条目的 `seq` 在**首次出现**时由观察方分配并落盘,后续更新(同一 id 的
//! 文本追加 / 工具状态变化)只改内容不改 `seq`。因此并发落盘的先后顺序不会让"新工具插到
//! 旧文本前面"——渲染顺序只由 `seq` 决定。
//!
//! `project.jsonl` 保留原始消息;本流补充文本与工具交替的 item 顺序,不能重复展示两份正文。
use crate::agent::sanitize_detail_text;
use crate::config::write_game_creator_private_file;
use crate::project::{enforce_project_permission_policy, project_append_lock_for};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::BTreeMap;
use std::fs::File;
use std::io::{BufRead, BufReader};
use std::path::{Path, PathBuf};
/// 行信封类型,与既有历史文件同构(`{"type": …, "payload": {…}}`)。
pub(crate) const DIRECT_TURN_STREAM_RECORD_TYPE: &str = "turn_stream_item";
/// 条目 schema 版本。
pub(crate) const DIRECT_TURN_STREAM_SCHEMA_VERSION: &str = "agc-turn-stream.v1";
/// 回读上限:只保留最后这么多条(按 `seq` 取最新)。
pub(crate) const DIRECT_TURN_STREAM_LIMIT: usize = 400;
/// 单条文本段的字符上限(与工具明细同口径的截断,避免单段失控)。
const DIRECT_TURN_STREAM_TEXT_MAX_CHARS: usize = 8000;
/// 没有流式分段时,最终回复那一段的固定 item id。
const DIRECT_TURN_STREAM_FINAL_ITEM_ID: &str = "final";
/// 回合失败说明那一段的固定 item id:失败说明也是这一回合的内容,排在流末尾。
pub(crate) const DIRECT_TURN_STREAM_FAILURE_ITEM_ID: &str = "failure";
/// 文本段。
pub(crate) const DIRECT_TURN_STREAM_KIND_TEXT: &str = "text";
/// 工具调用的位置标记。
pub(crate) const DIRECT_TURN_STREAM_KIND_TOOL: &str = "tool";
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct DirectTurnStreamItem {
pub(crate) schema_version: String,
/// 幂等身份:文本段 `text:<turnId>:<itemId>`、工具 `tool:<turnId>:<callId>`。
pub(crate) id: String,
pub(crate) turn_id: String,
/// `text` | `tool`
pub(crate) kind: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub(crate) text: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub(crate) call_id: Option<String>,
/// 首次出现的写入序号:**顺序真相**,同刻按它排序。
pub(crate) seq: u64,
/// 条目首次出现的本机毫秒时刻。
pub(crate) at: u64,
pub(crate) updated_at: u64,
}
#[cfg(test)]
mod snapshot_tests {
use super::*;
fn text(
turn: &str,
id: &str,
seq: u64,
at: u64,
updated: u64,
text: &str,
) -> DirectTurnStreamItem {
direct_turn_stream_text_item(Path::new("."), turn, id, text, seq, at, updated)
}
#[test]
fn late_older_snapshot_cannot_undo_completed_text_or_position() {
let complete = text("turn", "item", 1, 1000, 1002, "正文");
let late = text("turn", "item", 9, 1001, 1001, "更长但已经过期的草稿");
let merged = normalize_stream_items(vec![complete, late]);
assert_eq!(merged.len(), 1);
assert_eq!(merged[0].text.as_deref(), Some("正文"));
assert_eq!(merged[0].seq, 1);
assert_eq!(merged[0].at, 1000);
}
#[test]
fn retention_does_not_treat_new_turn_seq_one_as_oldest() {
let mut snapshots = (1..=DIRECT_TURN_STREAM_LIMIT)
.map(|seq| text("old", &seq.to_string(), seq as u64, 1000, 1000, "旧"))
.collect::<Vec<_>>();
snapshots.push(text("new", "one", 1, 2000, 2000, "新"));
let merged = normalize_stream_items(snapshots);
assert_eq!(merged.len(), DIRECT_TURN_STREAM_LIMIT);
assert_eq!(merged.last().unwrap().turn_id, "new");
}
}
impl DirectTurnStreamItem {
fn order_key(&self) -> (u64, u64, &str) {
(self.seq, self.at, self.id.as_str())
}
}
fn turn_stream_path(root: &Path) -> PathBuf {
root.join(".agent/conversations/turn-stream.jsonl")
}
/// 文本段条目的幂等 id:同一个 Codex assistant item 只占一行。
pub(crate) fn direct_turn_stream_text_item_id(turn_id: &str, item_id: &str) -> String {
format!("text:{}:{}", turn_id.trim(), item_id.trim())
}
/// 工具条目(位置标记)的幂等 id:同一个 callId 只占一行。
pub(crate) fn direct_turn_stream_tool_item_id(turn_id: &str, call_id: &str) -> String {
format!("tool:{}:{}", turn_id.trim(), call_id.trim())
}
/// 构造一条文本段条目:脱敏 + 截断与 `tool-calls.jsonl` 同口径。
pub(crate) fn direct_turn_stream_text_item(
root: &Path,
turn_id: &str,
item_id: &str,
text: &str,
seq: u64,
at: u64,
updated_at: u64,
) -> DirectTurnStreamItem {
DirectTurnStreamItem {
schema_version: DIRECT_TURN_STREAM_SCHEMA_VERSION.to_string(),
id: direct_turn_stream_text_item_id(turn_id, item_id),
turn_id: turn_id.trim().to_string(),
kind: DIRECT_TURN_STREAM_KIND_TEXT.to_string(),
text: Some(sanitize_stream_text(root, text)),
call_id: None,
seq,
at,
updated_at,
}
}
/// 构造一条工具条目:只记位置,正文仍来自 `DirectToolCall`。
pub(crate) fn direct_turn_stream_tool_item(
turn_id: &str,
call: &crate::DirectToolCall,
seq: u64,
at: u64,
) -> DirectTurnStreamItem {
DirectTurnStreamItem {
schema_version: DIRECT_TURN_STREAM_SCHEMA_VERSION.to_string(),
id: direct_turn_stream_tool_item_id(turn_id, &call.id),
turn_id: turn_id.trim().to_string(),
kind: DIRECT_TURN_STREAM_KIND_TOOL.to_string(),
text: None,
call_id: Some(call.id.trim().to_string()),
seq,
at,
updated_at: call.updated_at,
}
}
/// 文本脱敏 + 截断:与 `tool-calls.jsonl` 同一套 `sanitize_detail_text`。
pub(crate) fn sanitize_stream_text(root: &Path, text: &str) -> String {
let sanitized = sanitize_detail_text(root, text);
if sanitized.chars().count() <= DIRECT_TURN_STREAM_TEXT_MAX_CHARS {
return sanitized;
}
let mut truncated = sanitized
.chars()
.take(DIRECT_TURN_STREAM_TEXT_MAX_CHARS)
.collect::<String>();
truncated.push('…');
truncated
}
fn record_line(item: &DirectTurnStreamItem) -> Result<String, String> {
serde_json::to_string(&serde_json::json!({
"type": DIRECT_TURN_STREAM_RECORD_TYPE,
"payload": item,
}))
.map_err(|error| format!("序列化回合流条目失败:{error}"))
}
/// 解析一行信封;坏行 / 非本文件条目都返回 `None`(尽力而为的展示数据,不整体失败)。
fn stream_item_from_line(line: &str) -> Option<DirectTurnStreamItem> {
let trimmed = line.trim();
if trimmed.is_empty() {
return None;
}
let parsed: Value = serde_json::from_str(trimmed).ok()?;
if parsed.get("type").and_then(Value::as_str) != Some(DIRECT_TURN_STREAM_RECORD_TYPE) {
return None;
}
let payload = parsed.get("payload")?;
let mut item: DirectTurnStreamItem = serde_json::from_value(payload.clone()).ok()?;
if item.id.trim().is_empty() || item.turn_id.trim().is_empty() {
return None;
}
if !matches!(
item.kind.as_str(),
DIRECT_TURN_STREAM_KIND_TEXT | DIRECT_TURN_STREAM_KIND_TOOL
) {
return None;
}
if item.schema_version.trim().is_empty() {
item.schema_version = DIRECT_TURN_STREAM_SCHEMA_VERSION.to_string();
}
Some(item)
}
fn read_stream_lines(path: &Path) -> Vec<DirectTurnStreamItem> {
let Ok(file) = File::open(path) else {
return Vec::new();
};
let mut reader = BufReader::new(file);
let mut buffer = Vec::new();
let mut items = Vec::new();
loop {
buffer.clear();
match reader.read_until(b'\n', &mut buffer) {
Ok(0) => break,
// 单行解码失败(非法 UTF-8)只跳过这一行,继续读后面的行。
Ok(_) => match std::str::from_utf8(&buffer) {
Ok(line) => {
if let Some(item) = stream_item_from_line(line) {
items.push(item);
}
}
Err(_) => continue,
},
// 读 I/O 错误:无法再定位下一行边界,停止读取(已读到的照常返回)。
Err(_) => break,
}
}
items
}
/// 同一 id 的重复行合并:`seq` 取最早(位置钉死,后到的不得回退),`at` 取最早非零,
/// `updated_at` 取最大;文本只在更新(或同刻更长)的快照上替换。
fn merge_stream_snapshot(
existing: &DirectTurnStreamItem,
incoming: &DirectTurnStreamItem,
) -> DirectTurnStreamItem {
let text_len = |item: &DirectTurnStreamItem| {
item.text
.as_deref()
.map(str::chars)
.map(Iterator::count)
.unwrap_or_default()
};
// writer 保证更新时间单调;完成快照可以纠正正文,旧快照不能靠更长抢回所有权。
let take_incoming = incoming.updated_at > existing.updated_at
|| (incoming.updated_at == existing.updated_at && text_len(incoming) > text_len(existing));
let mut merged = existing.clone();
if take_incoming {
merged.text = incoming.text.clone();
}
merged.updated_at = merged.updated_at.max(incoming.updated_at);
if merged.call_id.is_none() {
merged.call_id = incoming.call_id.clone();
}
merged.seq = merged.seq.min(incoming.seq);
merged.at = [merged.at, incoming.at]
.into_iter()
.filter(|at| *at > 0)
.min()
.unwrap_or_default();
merged
}
/// 按身份归并;跨回合按起点,回合内按 seq,不能用局部 seq 判断全局新旧。
fn normalize_stream_items(items: Vec<DirectTurnStreamItem>) -> Vec<DirectTurnStreamItem> {
let mut by_id: BTreeMap<String, DirectTurnStreamItem> = BTreeMap::new();
for item in items {
let merged = match by_id.remove(&item.id) {
Some(previous) => merge_stream_snapshot(&previous, &item),
None => item,
};
by_id.insert(merged.id.clone(), merged);
}
let mut normalized = by_id.into_values().collect::<Vec<_>>();
let mut turn_starts = BTreeMap::<String, u64>::new();
for item in &normalized {
turn_starts
.entry(item.turn_id.clone())
.and_modify(|at| *at = (*at).min(item.at))
.or_insert(item.at);
}
normalized.sort_by(|left, right| {
(turn_starts[&left.turn_id], &left.turn_id, left.order_key()).cmp(&(
turn_starts[&right.turn_id],
&right.turn_id,
right.order_key(),
))
});
if normalized.len() > DIRECT_TURN_STREAM_LIMIT {
normalized.drain(..normalized.len() - DIRECT_TURN_STREAM_LIMIT);
}
normalized
}
/// 锁内读改写:整文件重写(追加与就地更新混用,没有纯追加的 JSONL 语义)。
/// 文件规模由 400 条上限与 8000 字符截断兜住。
fn with_locked_stream_items<T>(
root: &Path,
mutate: impl FnOnce(&mut Vec<DirectTurnStreamItem>) -> T,
) -> Result<T, String> {
let path = turn_stream_path(root);
let _project_lock = crate::acquire_game_creator_agent_runtime_project_write_lock_with_wait(
root,
"conversation.write",
)?;
let lock = project_append_lock_for(&path)?;
let _append_guard = lock.lock("回合流写入")?;
let mut items = read_stream_lines(&path);
let outcome = mutate(&mut items);
let normalized = normalize_stream_items(items);
let mut body = String::new();
for item in &normalized {
body.push_str(&record_line(item)?);
body.push('\n');
}
write_game_creator_private_file(&path, body.as_bytes(), "回合流历史")?;
Ok(outcome)
}
/// 幂等 upsert 一条回合流条目。
///
/// 位置(`seq` / `at`)只在第一次出现时确定:同一 id 的后续快照不得回退位置,
/// 也不得把已经写下的文本改短(并发落盘下"后到的旧快照"不会覆盖新快照)。
pub(crate) fn upsert_direct_turn_stream_item_at(
root: &Path,
item: &DirectTurnStreamItem,
) -> Result<(), String> {
enforce_project_permission_policy(root, "conversation.write")?;
with_locked_stream_items(root, |items| {
// normalize_stream_items 在锁内归并全部版本;不得提前删除比较基准。
items.push(item.clone());
})
}
/// 追加一段固定身份的文本段(失败说明等):位置排在当前流末尾。
///
/// 幂等:同一 `(turnId, itemId)` 已经存在时只更新文本与 `updatedAt`(回合重放 / 重复收尾
/// 不会多出一段)。返回写下的那一条,调用方用它下发同一份快照。
pub(crate) fn append_direct_turn_stream_text_at(
root: &Path,
turn_id: &str,
item_id: &str,
text: &str,
) -> Result<Option<DirectTurnStreamItem>, String> {
let turn_id = turn_id.trim();
let text = text.trim();
if turn_id.is_empty() || text.is_empty() {
return Ok(None);
}
let sanitized = sanitize_stream_text(root, text);
let item_id = item_id.trim();
enforce_project_permission_policy(root, "conversation.write")?;
let now = crate::agent::direct_tool_call_now_ms();
with_locked_stream_items(root, |items| {
let existing_id = direct_turn_stream_text_item_id(turn_id, item_id);
if let Some(existing) = items.iter_mut().find(|item| item.id == existing_id) {
// 位置不动:只替换文本与 updatedAt。
existing.text = Some(sanitized.clone());
existing.updated_at = now.max(existing.updated_at);
return Some(existing.clone());
}
// 首次出现:位置钉在末尾(当前最大 seq + 1)。
let next_seq = items.iter().map(|item| item.seq).max().unwrap_or(0) + 1;
let item = DirectTurnStreamItem {
schema_version: DIRECT_TURN_STREAM_SCHEMA_VERSION.to_string(),
id: existing_id,
turn_id: turn_id.to_string(),
kind: DIRECT_TURN_STREAM_KIND_TEXT.to_string(),
text: Some(sanitized),
call_id: None,
seq: next_seq,
at: now,
updated_at: now,
};
items.push(item.clone());
Some(item)
})
}
/// 没有任何 item 文本时补最终回复;已有 item 由完成事件负责,不能猜测覆盖某一段。
pub(crate) fn finalize_direct_turn_stream_reply_at(
root: &Path,
turn_id: &str,
visible_reply: &str,
) -> Result<Option<DirectTurnStreamItem>, String> {
let turn_id = turn_id.trim();
if turn_id.is_empty() || visible_reply.trim().is_empty() {
return Ok(None);
}
// 入口再做一次可见性投影:调用方给的是原始回复时,思考块不能落进对话流。
let visible_reply = crate::agent::project_direct_codex_visible_text(visible_reply)
.unwrap_or_else(|| visible_reply.trim().to_string());
let visible_reply = visible_reply.as_str();
enforce_project_permission_policy(root, "conversation.write")?;
let now = crate::agent::direct_tool_call_now_ms();
with_locked_stream_items(root, |items| {
if items
.iter()
.any(|item| item.turn_id == turn_id && item.kind == DIRECT_TURN_STREAM_KIND_TEXT)
{
None
} else {
let next_seq = items
.iter()
.filter(|item| item.turn_id == turn_id)
.map(|item| item.seq)
.max()
.unwrap_or(0)
+ 1;
let item = direct_turn_stream_text_item(
root,
turn_id,
DIRECT_TURN_STREAM_FINAL_ITEM_ID,
visible_reply,
next_seq,
now,
now,
);
items.push(item.clone());
Some(item)
}
})
}
@@ -60,7 +60,6 @@ pub(crate) fn emit_direct_game_creator_progress(root: &Path, stage: &str, messag
#[derive(Clone)]
pub(crate) struct DirectGameCreatorTurnUpdateEmitter {
project_path: String,
/// Thread Manager 的线程身份:进度只回填到"这一轮仍被占用"的那一格上。
thread_id: String,
turn_id: String,
@@ -70,74 +69,14 @@ pub(crate) struct DirectGameCreatorTurnUpdateEmitter {
impl DirectGameCreatorTurnUpdateEmitter {
pub(crate) fn new(root: &Path, turn_id: String) -> Self {
Self {
project_path: root.to_string_lossy().into_owned(),
thread_id: crate::agent::thread_id_for_project(root),
turn_id,
sequence: Arc::new(AtomicU64::new(0)),
}
}
pub(crate) fn emit(
&self,
status: &'static str,
activity: Option<&'static str>,
accumulated_text: Option<String>,
tool_calls: Option<Vec<crate::DirectToolCall>>,
) {
self.emit_with_reasoning(status, activity, accumulated_text, tool_calls, None);
}
/// 带思考过程的回合更新:`reasoning_text` 为"当前累计的思考全文"(前端整段替换)。
pub(crate) fn emit_with_reasoning(
&self,
status: &'static str,
activity: Option<&'static str>,
accumulated_text: Option<String>,
tool_calls: Option<Vec<crate::DirectToolCall>>,
reasoning_text: Option<String>,
) {
self.emit_full(
status,
activity,
accumulated_text,
tool_calls,
reasoning_text,
Vec::new(),
);
}
/// 带回合流的回合更新:`stream_items` 是"顺序真相"里本次变化的那几条。
///
/// 前端按这些条目的 `seq` 顺序渲染,所以它们必须来自与落盘同一份数据,
/// 不能在前端各算一套顺序。
pub(crate) fn emit_with_stream_items(
&self,
status: &'static str,
activity: Option<&'static str>,
accumulated_text: Option<String>,
tool_calls: Option<Vec<crate::DirectToolCall>>,
stream_items: Vec<crate::DirectTurnStreamItem>,
) {
self.emit_full(
status,
activity,
accumulated_text,
tool_calls,
None,
stream_items,
);
}
#[allow(clippy::too_many_arguments)]
fn emit_full(
&self,
status: &'static str,
activity: Option<&'static str>,
accumulated_text: Option<String>,
tool_calls: Option<Vec<crate::DirectToolCall>>,
reasoning_text: Option<String>,
stream_items: Vec<crate::DirectTurnStreamItem>,
) {
/// 回填这一轮逻辑回合的进度(状态 / 活动 / 序号);首页「运行中的项目」快照读它。
pub(crate) fn emit(&self, status: &'static str, activity: Option<&'static str>) {
let status_is_allowed = matches!(
status,
"accepted" | "running" | "streaming" | "finalizing" | "completed" | "failed"
@@ -176,24 +115,6 @@ impl DirectGameCreatorTurnUpdateEmitter {
sequence,
updated_at,
);
let Some(app) = GAME_CREATOR_AGENT_RUNTIME_UPDATE_APP_HANDLE.get() else {
return;
};
let _ = app.emit(
"game-creator-direct-turn-update",
GameCreatorDirectTurnUpdateEvent {
project_path: self.project_path.clone(),
turn_id: self.turn_id.clone(),
sequence,
status: status.to_string(),
activity: activity.map(str::to_string),
accumulated_text,
tool_calls,
reasoning_text,
stream_items: (!stream_items.is_empty()).then_some(stream_items),
updated_at,
},
);
}
pub(crate) fn turn_id(&self) -> &str {
@@ -17,7 +17,7 @@ use std::path::{Path, PathBuf};
use crate::agent::PendingTurn;
use crate::agent::{
append_direct_project_user_message_at, claim_pending_turn, complete_turn_if_reserved,
direct_codex_user_item_to_prompt, direct_tool_call_now_ms, redact_agent_runtime_error,
direct_codex_user_item_to_prompt, direct_now_ms, redact_agent_runtime_error,
run_direct_game_creator_turn_at_with_creation_type_and_emitter, thread_id_for_project,
DirectGameCreatorTurnUpdateEmitter, DirectTurnError, DirectTurnTerminal, DispatchedTurn,
};
@@ -55,7 +55,7 @@ impl TurnReservation {
complete_turn_if_reserved(
&self.thread_id,
&self.token,
terminal.event(direct_tool_call_now_ms(), self.user_item_id.as_deref()),
terminal.event(direct_now_ms(), self.user_item_id.as_deref()),
)
}
@@ -72,7 +72,7 @@ impl TurnReservation {
}))
.expect("canonical user item"),
None,
direct_tool_call_now_ms(),
direct_now_ms(),
);
super::enqueue_pending_turn(thread_id, pending).expect("enqueue test turn");
let dispatched = claim_pending_turn(thread_id).expect("claim test turn");
@@ -162,7 +162,7 @@ async fn run_dispatched_direct_turn(
match outcome {
Ok(reply) => {
// 深层的终态出口已经在 `run_turn` 里写出 `turn.completed`;这里只补最后一条回合更新。
emitter.emit("completed", Some("none"), Some(reply), None);
emitter.emit("completed", Some("none"));
}
Err(error) => {
// 放行之后的失败一律是回合失败:失败诊断与失败说明已由上层写过,这里补终态事件。
@@ -211,7 +211,7 @@ mod tests {
}))
.expect("canonical user item"),
None,
direct_tool_call_now_ms(),
direct_now_ms(),
)
}
@@ -17,9 +17,9 @@ use std::sync::{Mutex, OnceLock};
use uuid::Uuid;
use crate::agent::{
direct_codex_user_item_id_for_client_turn_id, direct_tool_call_now_ms, queue_has_room,
ConsumeResult, EnqueueOutcome, EnqueueRejection, PendingTurn, QueueRemovalOutcome,
QueueRemovalReason, SubscriptionBootstrap, ThreadEvent,
direct_codex_user_item_id_for_client_turn_id, direct_now_ms, queue_has_room, ConsumeResult,
EnqueueOutcome, EnqueueRejection, PendingTurn, QueueRemovalOutcome, QueueRemovalReason,
SubscriptionBootstrap, ThreadEvent,
};
const DEFAULT_MAX_EVENTS: usize = 8_192;
@@ -854,7 +854,7 @@ pub(crate) fn remove_pending_turn(thread_id: &str, client_turn_id: &str) -> Queu
global_thread_manager()
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.remove_pending_turn(thread_id, client_turn_id, direct_tool_call_now_ms())
.remove_pending_turn(thread_id, client_turn_id, direct_now_ms())
};
if outcome == QueueRemovalOutcome::Removed {
notify_subscribers(thread_id);
@@ -873,7 +873,7 @@ pub(crate) fn claim_pending_turn(thread_id: &str) -> Option<DispatchedTurn> {
global_thread_manager()
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.claim_pending_turn(thread_id, direct_tool_call_now_ms())
.claim_pending_turn(thread_id, direct_now_ms())
};
if claimed.is_some() {
notify_subscribers(thread_id);
@@ -954,7 +954,7 @@ pub(crate) fn release_stale_direct_turn(
thread_id,
expected_client_turn_id,
min_age_ms,
direct_tool_call_now_ms(),
direct_now_ms(),
)
};
if matches!(outcome, StaleTurnRelease::Released(_)) {
@@ -21,6 +21,15 @@ use serde_json::Value;
use std::path::Path;
use ts_rs::TS;
/// 宿主观测时刻:Unix 毫秒。
pub(crate) fn direct_now_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis()
.min(u64::MAX as u128) as u64
}
/// 正文(消息 / 思考)上限。
const THREAD_TEXT_MAX_CHARS: usize = 8_000;
/// 工具明细(命令 / 参数 / 输出)上限。
@@ -676,7 +685,7 @@ fn json_ms(container: &Value, key: &str) -> Option<u64> {
/// `item/started` / `item/completed` 的事件级阶段时间(毫秒)。
///
/// 字段位置按当前 app-server 协议:通知层带 `params.startedAtMs` / `params.completedAtMs`,
/// 条目自带时用条目里的同名毫秒字段(`direct_tool_calls` 读的是同一处)。完成事件即使同时
/// 条目自带时用条目里的同名毫秒字段。完成事件即使同时
/// 带着开始字段也只取**完成**时间;两者都没有、但 `durationMs` 有可靠起点时按
/// 起点 + 时长派生结束。都没有就用宿主处理该事件的钟——原生缺阶段时间时这是唯一诚实的值。
pub(crate) fn thread_item_event_at_ms(
@@ -888,29 +888,6 @@ struct GameCreatorAgentProgressEvent {
message: String,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
struct GameCreatorDirectTurnUpdateEvent {
project_path: String,
turn_id: String,
sequence: u64,
status: String,
activity: Option<String>,
accumulated_text: Option<String>,
/// 本回合内发生变化的结构化工具调用集合(只有变化时才带,老事件没有这个字段)。
/// `skip_serializing_if`:字段缺席时前端拿到 `undefined`,行为与改造前一致。
#[serde(skip_serializing_if = "Option::is_none")]
tool_calls: Option<Vec<crate::DirectToolCall>>,
/// 本回合当前累计的思考过程(流式整段替换);拿不到时字段缺席。
#[serde(skip_serializing_if = "Option::is_none")]
reasoning_text: Option<String>,
/// 本回合**顺序真相**里本次发生变化的那几条(文本段 / 工具位置标记)。
/// `skip_serializing_if`:字段缺席时前端拿到 `undefined`,行为与改造前一致。
#[serde(skip_serializing_if = "Option::is_none")]
stream_items: Option<Vec<crate::DirectTurnStreamItem>>,
updated_at: u64,
}
#[derive(Debug, Eq, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
struct GameCreatorLlmConfigStatus {