From 2748468d127a2c4c6455a535fe0358b26c603be6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=8E=8B=E5=BE=B7=E5=AE=87?= Date: Wed, 16 Sep 2026 19:38:31 +0800 Subject: [PATCH] =?UTF-8?q?DirectProject=20=E8=81=8A=E5=A4=A9=E4=BA=8B?= =?UTF-8?q?=E4=BB=B6=E6=94=B9=E7=94=A8=20ts-rs=20=E5=AF=BC=E5=87=BA?= =?UTF-8?q?=E7=9A=84=20tagged=20enum=20=E5=B9=B6=E5=88=A0=E6=8E=89=20turn?= =?UTF-8?q?=20id?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 新增 direct_thread_wire.rs:DirectThreadItem / DirectThreadEvent / 订阅与历史切片全部改成 ts-rs 导出的 tagged enum,取代原大而全的可空结构体 - 删除 direct_thread_raw_item.rs,模块注册与直通引用改到 direct_thread_wire - 条目身份只看一个 itemId:工具条目的第二个 id 在 Rust 边界归一,不再对外暴露 - 删除 DirectProject 聊天事件里的 turn id:生命周期用无载荷的 turn.started / turn.completed{status} 表示 - append 直接接收 DirectThreadEvent 并返回同一事件,队列内部自算 seq - 请求事件改为携带 DirectThreadRequestKind,去掉字符串中转 - 思考增量走 ReasoningDelta 通道,与正文增量共用 item.delta - at 用 #[ts(as = "f64")] 对齐 Tauri JSON 通道的 number - 用 cargo test export_bindings 重新生成 project-workspace/generated 绑定 --- .../src-tauri/src/agent.rs | 4 +- .../src/agent/codex_app_server/mod.rs | 224 ++--- .../src/agent/direct_thread_manager.rs | 448 ++++----- .../src/agent/direct_thread_raw_item.rs | 581 ------------ .../src-tauri/src/agent/direct_thread_wire.rs | 864 ++++++++++++++++++ .../src-tauri/src/agent/direct_tool_calls.rs | 2 +- .../src-tauri/src/commands.rs | 2 +- .../generated/DirectCodexUserContentPart.ts | 3 +- .../generated/DirectCodexUserItem.ts | 10 +- .../DirectCodexUserMessageEnvelope.ts | 4 + .../generated/DirectCodexUserMessageItem.ts | 6 +- .../generated/DirectCodexUserRole.ts | 2 +- .../DirectCodexUserRuntimeRegionPart.ts | 18 +- .../generated/DirectThreadConsumeResult.ts | 4 + .../generated/DirectThreadDeltaKind.ts | 6 + .../generated/DirectThreadEvent.ts | 30 + .../generated/DirectThreadFileChange.ts | 12 + .../generated/DirectThreadHistorySlice.ts | 14 + .../generated/DirectThreadItem.ts | 76 ++ .../generated/DirectThreadRequestKind.ts | 9 + .../DirectThreadSubscriptionBootstrap.ts | 14 + 21 files changed, 1315 insertions(+), 1018 deletions(-) delete mode 100644 apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_raw_item.rs create mode 100644 apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_wire.rs create mode 100644 apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserMessageEnvelope.ts create mode 100644 apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadConsumeResult.ts create mode 100644 apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadDeltaKind.ts create mode 100644 apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadEvent.ts create mode 100644 apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadFileChange.ts create mode 100644 apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadHistorySlice.ts create mode 100644 apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadItem.ts create mode 100644 apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadRequestKind.ts create mode 100644 apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadSubscriptionBootstrap.ts diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent.rs b/apps/ai-game-creator-shell/src-tauri/src/agent.rs index f3ca892e7..5c8f7d58a 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent.rs @@ -21,7 +21,7 @@ mod direct_project_history; mod direct_project_turn_history; mod direct_runtime; mod direct_thread_manager; -mod direct_thread_raw_item; +mod direct_thread_wire; mod direct_tool_bridge; mod direct_tool_calls; mod direct_tools_mcp; @@ -56,7 +56,7 @@ pub(crate) use direct_project_history::*; pub(crate) use direct_project_turn_history::*; pub(crate) use direct_runtime::*; pub(crate) use direct_thread_manager::*; -pub(crate) use direct_thread_raw_item::*; +pub(crate) use direct_thread_wire::*; pub(crate) use direct_tool_bridge::*; pub(crate) use direct_tool_calls::*; pub(crate) use direct_tools_mcp::*; diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/codex_app_server/mod.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/codex_app_server/mod.rs index 8b789f2e8..f47ec3b29 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/codex_app_server/mod.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/codex_app_server/mod.rs @@ -564,6 +564,17 @@ enum CodexTurnEvent { item_id: String, delta: String, }, + /// 思考正文增量:app-server `item/reasoning/summaryTextDelta` 的明文思考文本。 + /// + /// `item/reasoning/summaryTextDelta`(core `ReasoningContentDelta`)与 + /// `item/reasoning/textDelta`(core `ReasoningRawContentDelta`)都进这条通道:前者是 + /// reasoning item 的 `summary`,后者是它的 `content`,两段文本都随 `item/completed` + /// 落进 `project.jsonl`、此前也已经在完成时展示给用户。plan 文本与命令输出仍然只降级为 + /// 活动状态,不下发正文。 + ReasoningDelta { + item_id: String, + delta: String, + }, IntermediateText(String), Activity(&'static str), Item { @@ -571,7 +582,7 @@ enum CodexTurnEvent { params: serde_json::Value, }, Request { - event_type: &'static str, + kind: DirectThreadRequestKind, params: serde_json::Value, }, RawItem(serde_json::Value), @@ -741,40 +752,14 @@ fn direct_codex_safe_activity_for_item_value(item: &serde_json::Value) -> &'stat direct_codex_safe_activity_for_item(item_type) } -/// Project an app-server item into the small public payload carried by the -/// DirectProject event queue. Full item contents are persisted in JSONL and -/// must not be forwarded through the runtime event stream. -/// 运行态事件载荷:与历史切片同形的脱敏原始条目;拿不到条目时给空对象。 +/// 运行态事件载荷:与历史切片同形的脱敏原始条目;拿不到身份或类型就整条跳过。 /// /// 这里不生成工具卡片形状:标题、折叠摘要和可见性都是前端投影的职责。 -fn direct_thread_raw_item_payload( +fn direct_thread_event_item( root: &std::path::Path, item: &serde_json::Value, - turn_id: Option<&str>, - completed: bool, - now_ms: u64, -) -> serde_json::Value { - direct_thread_raw_item_from_value(root, item, turn_id, completed, now_ms) - .and_then(|raw| serde_json::to_value(raw).ok()) - .unwrap_or_else(|| serde_json::json!({})) -} - -fn direct_thread_item_id(item: &serde_json::Value) -> Option { - item.get("id") - .and_then(serde_json::Value::as_str) - .filter(|value| !value.is_empty()) - .map(str::to_string) -} - -/// 原始 response item 的调用 id:工具条目的 app-server `itemId` 就是这个值, -/// 所以它是两个 id 空间唯一的对齐点。 -fn direct_thread_item_call_id(item: &serde_json::Value) -> Option { - item.get("call_id") - .or_else(|| item.get("callId")) - .and_then(serde_json::Value::as_str) - .map(str::trim) - .filter(|value| !value.is_empty()) - .map(str::to_string) +) -> Option { + direct_thread_item_from_value(root, item, direct_tool_call_now_ms()) } fn direct_codex_command_is_game_verification(command: &str) -> bool { @@ -991,19 +976,21 @@ fn direct_codex_safe_activity_for_notification(method: &str) -> Option<&'static } } -fn direct_codex_request_event_type(method: &str) -> Option<&'static str> { +fn direct_codex_request_event_type(method: &str) -> Option { match method { "item/fileChange/requestApproval" | "item/commandExecution/requestApproval" - | "item/permissions/requestApproval" => Some("approval.requested"), - "item/tool/requestUserInput" | "item/mcpToolCall/requestUserInput" => Some("ask.requested"), + | "item/permissions/requestApproval" => Some(DirectThreadRequestKind::ApprovalRequested), + "item/tool/requestUserInput" | "item/mcpToolCall/requestUserInput" => { + Some(DirectThreadRequestKind::AskRequested) + } _ => None, } } -fn direct_codex_resolution_event_type(method: &str) -> Option<&'static str> { +fn direct_codex_resolution_event_type(method: &str) -> Option { match method { - "serverRequest/resolved" => Some("request.resolved"), + "serverRequest/resolved" => Some(DirectThreadRequestKind::RequestResolved), _ => None, } } @@ -1076,6 +1063,33 @@ fn direct_codex_notification_event( intermediate_text: Option, safe_activity: Option<&'static str>, ) -> Option { + // 思考正文走独立通道,交给 DirectProject 的运行态事件;它不因为 + // "preparing 活动" 的降级规则被丢掉,否则界面只能等 item/completed 才看到思考。 + // + // 两条通知都下发正文,不下发活动文本: + // - `item/reasoning/summaryTextDelta`(core `ReasoningContentDelta`)→ reasoning item 的 `summary`; + // - `item/reasoning/textDelta`(core `ReasoningRawContentDelta`)→ reasoning item 的 `content`, + // 正是 `project.jsonl` 里保存、并在此前 `item/completed` 已经展示给用户的同一段文本。 + // 因此这里只是把"完成时才看到"提前为"边生成边看到",没有放宽可见文本的范围; + // 未识别的 plan 文本与命令输出仍然只降级为活动状态,不下发正文。 + if matches!( + method, + "item/reasoning/summaryTextDelta" | "item/reasoning/textDelta" + ) { + return params + .get("delta") + .and_then(serde_json::Value::as_str) + .filter(|value| !value.is_empty()) + .map(|delta| CodexTurnEvent::ReasoningDelta { + item_id: params + .get("itemId") + .and_then(serde_json::Value::as_str) + .filter(|value| !value.is_empty()) + .map(str::to_string) + .unwrap_or_else(|| "direct-missing-item".to_string()), + delta: delta.to_string(), + }); + } let (activity, intermediate_text) = match (&intermediate_text, safe_activity) { (Some(_), Some(activity)) if activity == "preparing" => (Some(activity), None), _ => (safe_activity, intermediate_text), @@ -2943,19 +2957,7 @@ impl CodexAppServerConnection { turn_start_guard.armed = false; let direct_thread_id = history_root.to_string_lossy().into_owned(); if self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject { - append_direct_thread_event( - &direct_thread_id, - DirectThreadRawEventDraft { - event_type: "turn.started".to_string(), - turn_id: turn_id.clone(), - item_id: None, - call_id: None, - payload: serde_json::json!({ - "threadId": thread_id, - "turnId": turn_id, - }), - }, - ); + append_direct_thread_event(&direct_thread_id, DirectThreadEvent::turn_started()); } let mut receiver = self.register_turn(&turn_id).await; let mut direct_project_history = DirectProjectHistoryAccumulator::default(); @@ -3010,18 +3012,13 @@ impl CodexAppServerConnection { direct_project_history.observe_delta(&item_id, &delta); append_direct_thread_event( &direct_thread_id, - DirectThreadRawEventDraft { - event_type: "item.delta".to_string(), - turn_id: turn_id.clone(), - item_id: Some(item_id.clone()), - call_id: None, - // 事件自足:增量也要说明它是哪类条目的正文, - // 前端 reducer 不允许靠猜 itemId 的来源决定 kind。 - payload: serde_json::json!({ - "delta": delta.clone(), - "kind": "message", - }), - }, + // 事件自足:增量自带 item 身份与正文类别(正文 / 思考), + // 前端 reducer 不允许靠猜 itemId 的来源决定 kind。 + DirectThreadEvent::item_delta( + item_id.clone(), + DirectThreadDeltaKind::Message, + delta.clone(), + ), ); } streamed_text.push_str(&delta); @@ -3052,6 +3049,18 @@ impl CodexAppServerConnection { }); } } + Some(CodexTurnEvent::ReasoningDelta { item_id, delta }) => { + if self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject { + append_direct_thread_event( + &direct_thread_id, + DirectThreadEvent::item_delta( + item_id, + DirectThreadDeltaKind::Reasoning, + delta, + ), + ); + } + } Some(CodexTurnEvent::IntermediateText(text)) => { if let Some(observer) = direct_observer.as_deref_mut() { observer(DirectCodexTurnObservation::IntermediateText(text)); @@ -3064,13 +3073,7 @@ impl CodexAppServerConnection { "rawResponseItem/completed 缺少 item".to_string(), )); } - let entry_payload = direct_thread_raw_item_payload( - history_root, - &item, - Some(turn_id.as_str()), - true, - direct_tool_call_now_ms(), - ); + let entry_item = direct_thread_event_item(history_root, &item); let history_root = history_root.to_path_buf(); let history_item = item.clone(); tokio::task::spawn_blocking(move || { @@ -3084,20 +3087,15 @@ impl CodexAppServerConnection { })? .map_err(platform_llm::LlmError::InvalidRequest)?; direct_project_history.complete_item(&item); - let item_id = direct_thread_item_id(&item); - append_direct_thread_event( - &direct_thread_id, - DirectThreadRawEventDraft { - event_type: "item.completed".to_string(), - turn_id: turn_id.clone(), - item_id, - call_id: direct_thread_item_call_id(&item), - payload: entry_payload, - }, - ); + if let Some(entry_item) = entry_item { + append_direct_thread_event( + &direct_thread_id, + DirectThreadEvent::item_completed(entry_item), + ); + } } } - Some(CodexTurnEvent::Request { event_type, params }) => { + Some(CodexTurnEvent::Request { kind, params }) => { if self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject { let request_id = params .get("requestId") @@ -3107,15 +3105,7 @@ impl CodexAppServerConnection { .map(str::to_string); append_direct_thread_event( &direct_thread_id, - DirectThreadRawEventDraft { - event_type: event_type.to_string(), - turn_id: turn_id.clone(), - item_id: None, - call_id: None, - payload: request_id - .map(|id| serde_json::json!({ "requestId": id })) - .unwrap_or_else(|| serde_json::json!({})), - }, + DirectThreadEvent::request(kind, request_id), ); } } @@ -3226,23 +3216,14 @@ impl CodexAppServerConnection { && self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject { - let item_id = direct_thread_item_id(item); - append_direct_thread_event( - &direct_thread_id, - DirectThreadRawEventDraft { - event_type: "item.started".to_string(), - turn_id: turn_id.clone(), - item_id, - call_id: None, - payload: direct_thread_raw_item_payload( - history_root, - item, - Some(turn_id.as_str()), - completed, - direct_tool_call_now_ms(), - ), - }, - ); + if let Some(entry_item) = + direct_thread_event_item(history_root, item) + { + append_direct_thread_event( + &direct_thread_id, + DirectThreadEvent::item_started(entry_item), + ); + } } } } @@ -3284,13 +3265,7 @@ impl CodexAppServerConnection { { append_direct_thread_event( &direct_thread_id, - DirectThreadRawEventDraft { - event_type: "turn.completed".to_string(), - turn_id: turn_id.clone(), - item_id: None, - call_id: None, - payload: serde_json::json!({ "status": status }), - }, + DirectThreadEvent::turn_completed(status.to_string()), ); } match status { @@ -3985,8 +3960,8 @@ async fn read_game_creator_codex_app_server_stdout( continue; } } - let event = if let Some(event_type) = direct_codex_resolution_event_type(method) { - CodexTurnEvent::Request { event_type, params } + let event = if let Some(kind) = direct_codex_resolution_event_type(method) { + CodexTurnEvent::Request { kind, params } } else if let Some(activity) = safe_activity { // Preparing notifications may carry private plan/reasoning text; // expose only the safe activity category. Other categories may @@ -4038,8 +4013,8 @@ async fn read_game_creator_codex_app_server_stdout( ), method if direct_codex_request_event_type(method).is_some() => { CodexTurnEvent::Request { - event_type: direct_codex_request_event_type(method) - .expect("request event type checked above"), + kind: direct_codex_request_event_type(method) + .expect("request kind checked above"), params, } } @@ -4550,15 +4525,10 @@ mod tests { "arguments": { "path": "game/index.html", "token": "secret" }, "result": { "content": "large output" } }); - assert_eq!(direct_thread_item_id(&item).as_deref(), Some("item-1")); // 运行态事件必须自足:载荷是脱敏原始条目,前端不需要再按 itemId 取快照。 - let payload = direct_thread_raw_item_payload( - std::path::Path::new("."), - &item, - Some("turn-1"), - false, - 1000, - ); + let projected = direct_thread_event_item(std::path::Path::new("."), &item).expect("item"); + assert_eq!(projected.item_id(), "item-1"); + let payload = serde_json::to_value(&projected).expect("payload"); assert_eq!( payload.get("itemType").and_then(serde_json::Value::as_str), Some("mcpToolCall") @@ -4567,10 +4537,6 @@ mod tests { payload.get("itemId").and_then(serde_json::Value::as_str), Some("item-1") ); - assert_eq!( - payload.get("turnId").and_then(serde_json::Value::as_str), - Some("turn-1") - ); // 卡片标题 / 折叠摘要 / kind 属于前端投影:载荷里不得出现这些 UI 语义。 assert!(payload.get("toolCall").is_none(), "{payload}"); assert!(payload.get("title").is_none(), "{payload}"); diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_manager.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_manager.rs index 20c668a11..a30816d82 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_manager.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_manager.rs @@ -4,13 +4,13 @@ //! 每个 subscriber 一个受保护的消费游标。它不理解前端 reducer,也不负责 JSONL //! 持久化;调用方必须在完成 item 持久化成功后再追加对应完成事件。 -use serde::{Deserialize, Serialize}; -use serde_json::Value; use std::collections::{HashMap, HashSet}; use std::sync::{Mutex, OnceLock}; use uuid::Uuid; -use crate::agent::DirectThreadRawItem; +use crate::agent::{ + DirectThreadConsumeResult, DirectThreadEvent, DirectThreadSubscriptionBootstrap, +}; const DEFAULT_MAX_EVENTS: usize = 8_192; const DEFAULT_MAX_BYTES: usize = 8 * 1024 * 1024; @@ -21,65 +21,11 @@ pub(crate) const DIRECT_THREAD_NOTIFY_EVENT: &str = "game-creator-direct-thread- static DIRECT_THREAD_MANAGER: OnceLock> = OnceLock::new(); static DIRECT_THREAD_MANAGER_APP_HANDLE: OnceLock = OnceLock::new(); -#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] -#[serde(rename_all = "camelCase")] -pub(crate) struct DirectThreadRawEvent { - pub(crate) seq: u64, - #[serde(rename = "type")] - pub(crate) event_type: String, - pub(crate) turn_id: String, - #[serde(skip_serializing_if = "Option::is_none")] - pub(crate) item_id: Option, - /// 合并身份的另一半:工具类条目的 `itemId` 是调用 id,原始 item 的 `id` 是 response item id, - /// 两者只在 `callId` 上对齐。前端按 `callId ?? itemId` 归并同一张卡片。 - #[serde(skip_serializing_if = "Option::is_none")] - pub(crate) call_id: Option, - pub(crate) payload: Value, -} - -#[derive(Clone, Debug, Eq, PartialEq)] -pub(crate) struct DirectThreadRawEventDraft { - pub(crate) event_type: String, - pub(crate) turn_id: String, - pub(crate) item_id: Option, - pub(crate) call_id: Option, - pub(crate) payload: Value, -} - -impl DirectThreadRawEvent { - /// 条目合并身份:有 `callId` 就用它,否则用 `itemId`。 - fn item_identity(&self) -> Option<&str> { - self.call_id.as_deref().or(self.item_id.as_deref()) - } -} - -#[derive(Clone, Debug, Eq, PartialEq, Serialize)] -#[serde(rename_all = "camelCase")] -pub(crate) struct DirectThreadSubscriptionBootstrap { - pub(crate) subscription_id: String, - pub(crate) last_completed_item_id: Option, - pub(crate) events: Vec, -} - -#[derive(Clone, Debug, Eq, PartialEq, Serialize)] -#[serde(rename_all = "camelCase")] -pub(crate) struct DirectThreadConsumeResult { - pub(crate) events: Vec, -} - -#[derive(Clone, Debug, Eq, PartialEq, Serialize)] -#[serde(rename_all = "camelCase")] -pub(crate) struct DirectThreadHistorySlice { - /// 脱敏原始条目,顺序即文件顺序。与运行态事件载荷同形,前端只有一套投影。 - pub(crate) items: Vec, - pub(crate) has_more: bool, - /// 本次切片的原始条目锚点:无论切片里有没有可显示条目,分页都要靠它继续向前。 - pub(crate) first_item_id: Option, -} - +/// 队列里的一个事件。`seq` 只服务内部游标,不下发:前端按 `consume` 返回的数组顺序处理。 #[derive(Clone, Debug)] struct StoredEvent { - event: DirectThreadRawEvent, + event: DirectThreadEvent, + seq: u64, bytes: usize, cleanable: bool, } @@ -97,15 +43,14 @@ struct ThreadState { total_bytes: usize, active_items: HashSet, unresolved_requests: HashSet, - /// 最新一条 `turn.started` / `turn.completed` 的独立拷贝。 + /// 最近一条 `turn.started` / `turn.completed` 的独立拷贝。 /// /// TODO(thread-manager): 这里有意只保留"锚点",因为 replay 队列会回收可回收事件, /// 队列本身不是完美事件日志——被回收的 `turn.started` / `turn.completed` 不会回放, /// 只有这份拷贝保证新订阅仍能判定"最新回合是否还在跑"。若将来需要回放多个回合的 /// 生命周期(回合账本、跨进程恢复、按回合统计),必须另建持久 ledger, /// 不能靠扩大这份拷贝或放宽回收规则来模拟。 - lifecycle_anchor: Option, - last_completed_item_id: Option, + lifecycle_anchor: Option<(u64, DirectThreadEvent)>, subscribers: HashMap, } @@ -119,7 +64,6 @@ impl Default for ThreadState { active_items: HashSet::new(), unresolved_requests: HashSet::new(), lifecycle_anchor: None, - last_completed_item_id: None, subscribers: HashMap::new(), } } @@ -154,34 +98,29 @@ impl DirectThreadManager { pub(crate) fn append( &mut self, thread_id: &str, - draft: DirectThreadRawEventDraft, - ) -> DirectThreadRawEvent { + event: DirectThreadEvent, + ) -> DirectThreadEvent { let thread = self.threads.entry(thread_id.to_string()).or_default(); thread.next_seq = thread.next_seq.saturating_add(1); - let event = DirectThreadRawEvent { - seq: thread.next_seq, - event_type: draft.event_type, - turn_id: draft.turn_id, - item_id: draft.item_id, - call_id: draft.call_id, - payload: draft.payload, - }; - let cleanable = Self::observe_event(thread, &event); + let seq = thread.next_seq; + let cleanable = Self::observe_event(thread, seq, &event); let bytes = serde_json::to_vec(&event) .map(|value| value.len()) .unwrap_or_default(); thread.total_bytes = thread.total_bytes.saturating_add(bytes); thread.events.push(StoredEvent { event: event.clone(), + seq, bytes, cleanable, }); - Self::mark_item_events_cleanable(thread, event.item_identity()); - if matches!( - event.event_type.as_str(), - "approval.resolved" | "request.resolved" | "ask.resolved" - ) { - Self::mark_request_events_cleanable(thread, request_id(&event).as_deref()); + if let Some(item_id) = event.item_id() { + Self::mark_item_events_cleanable(thread, item_id, seq); + } + if let Some(kind) = event.request_kind() { + if kind.is_resolution() { + Self::mark_request_events_cleanable(thread, event.request_id(), seq); + } } self.evict(thread_id); event @@ -199,19 +138,20 @@ impl DirectThreadManager { .events .iter() .skip(thread.head) - .filter(|stored| Self::is_bootstrap_event(thread, &stored.event, stored.cleanable)) - .map(|stored| stored.event.clone()) + .filter(|stored| Self::is_bootstrap_event(thread, stored)) + .map(|stored| (stored.seq, stored.event.clone())) .collect::>(); - if let Some(anchor) = thread.lifecycle_anchor.as_ref() { - if !events.iter().any(|event| event.seq == anchor.seq) { - events.push(anchor.clone()); + // 生命周期锚点独立保存:即使队列里那条事件已被回收,也要作为 bootstrap 事件返回。 + if let Some((anchor_seq, anchor)) = thread.lifecycle_anchor.as_ref() { + if !events.iter().any(|(seq, _)| seq == anchor_seq) { + events.push((*anchor_seq, anchor.clone())); } } - events.sort_by_key(|event| event.seq); + events.sort_by_key(|(seq, _)| *seq); DirectThreadSubscriptionBootstrap { subscription_id, - last_completed_item_id: thread.last_completed_item_id.clone(), - events, + last_completed_item_id: None, + events: events.into_iter().map(|(_, event)| event).collect(), } } @@ -241,26 +181,28 @@ impl DirectThreadManager { let oldest_seq = thread .events .get(thread.head) - .map(|stored| stored.event.seq) + .map(|stored| stored.seq) .unwrap_or(thread.next_seq.saturating_add(1)); if cursor.saturating_add(1) < oldest_seq { thread.subscribers.remove(subscription_id); return Err(SUBSCRIPTION_EXPIRED.to_string()); } + let mut last_seq = cursor; let events = thread .events .iter() .skip(thread.head) - .filter(|stored| stored.event.seq > cursor) - .map(|stored| stored.event.clone()) + .filter(|stored| stored.seq > cursor) + .map(|stored| { + last_seq = stored.seq; + stored.event.clone() + }) .collect::>(); - if let Some(last) = events.last() { - thread - .subscribers - .get_mut(subscription_id) - .expect("subscriber remains registered") - .cursor = last.seq; - } + thread + .subscribers + .get_mut(subscription_id) + .expect("subscriber remains registered") + .cursor = last_seq; let result = DirectThreadConsumeResult { events }; Self::trim_prefix(thread); Ok(result) @@ -277,91 +219,76 @@ impl DirectThreadManager { }) } - fn observe_event(thread: &mut ThreadState, event: &DirectThreadRawEvent) -> bool { - match event.event_type.as_str() { - "item.started" => { - if let Some(identity) = event.item_identity() { - thread.active_items.insert(identity.to_string()); - } + /// 观察一条事件,返回它本身是否可回收。 + fn observe_event(thread: &mut ThreadState, seq: u64, event: &DirectThreadEvent) -> bool { + match event { + DirectThreadEvent::ItemStarted { item, .. } => { + thread.active_items.insert(item.item_id().to_string()); false } - "item.completed" => { - if let Some(identity) = event.item_identity() { - thread.active_items.remove(identity); - } - // 历史锚点必须是 project.jsonl 里的 response item id,不能用 call id。 - if let Some(item_id) = event.item_id.as_deref() { - thread.last_completed_item_id = Some(item_id.to_string()); - } + DirectThreadEvent::ItemCompleted { item, .. } => { + thread.active_items.remove(item.item_id()); true } - "turn.started" | "turn.completed" => { + DirectThreadEvent::TurnStarted { .. } | DirectThreadEvent::TurnCompleted { .. } => { // 队列只保留最新一条生命周期事件,更早的可能已被回收;见 // `ThreadState::lifecycle_anchor` 的 TODO:这不是完整事件日志。 - thread.lifecycle_anchor = Some(event.clone()); + thread.lifecycle_anchor = Some((seq, event.clone())); true } - "approval.requested" | "request.requested" | "ask.requested" => { - let request_id = request_id(event); + DirectThreadEvent::Request { + kind, request_id, .. + } => { + if kind.is_request() { + if let Some(request_id) = request_id.as_deref() { + thread.unresolved_requests.insert(request_id.to_string()); + } + // 没有 request id 的请求事件无法配对,直接视为可回收。 + return request_id.is_none(); + } if let Some(request_id) = request_id.as_deref() { - thread.unresolved_requests.insert(request_id.to_string()); - } - request_id.is_none() - } - "approval.resolved" | "request.resolved" | "ask.resolved" => { - if let Some(request_id) = request_id(event) { - thread.unresolved_requests.remove(&request_id); + thread.unresolved_requests.remove(request_id); } true } - _ => true, + DirectThreadEvent::ItemDelta { .. } => true, } } - fn is_bootstrap_event( - thread: &ThreadState, - event: &DirectThreadRawEvent, - cleanable: bool, - ) -> bool { - if thread - .lifecycle_anchor - .as_ref() - .is_some_and(|anchor| anchor.seq == event.seq) - { - return true; + /// 新订阅此刻需要补的事件:未完成 item 的完整事件、未解决请求、不可回收的事件。 + fn is_bootstrap_event(thread: &ThreadState, stored: &StoredEvent) -> bool { + if let Some(item_id) = stored.event.item_id() { + return thread.active_items.contains(item_id); } - if let Some(identity) = event.item_identity() { - return thread.active_items.contains(identity); + if let Some(request_id) = stored.event.request_id() { + return thread.unresolved_requests.contains(request_id); } - if let Some(request_id) = request_id(event) { - return thread.unresolved_requests.contains(&request_id); - } - !cleanable + // 生命周期锚点单独补,增量正文这类瞬时事件不回放。 + !stored.cleanable } - fn mark_item_events_cleanable(thread: &mut ThreadState, item_id: Option<&str>) { - let Some(item_id) = item_id else { - return; - }; + fn mark_item_events_cleanable(thread: &mut ThreadState, item_id: &str, seq: u64) { if thread.active_items.contains(item_id) { return; } for stored in &mut thread.events { - if stored.event.item_id.as_deref() == Some(item_id) { + if stored.seq <= seq && stored.event.item_id() == Some(item_id) { stored.cleanable = true; } } } - fn mark_request_events_cleanable(thread: &mut ThreadState, resolved_request_id: Option<&str>) { - let Some(resolved_request_id) = resolved_request_id else { + fn mark_request_events_cleanable(thread: &mut ThreadState, request_id: Option<&str>, seq: u64) { + let Some(resolved_request_id) = request_id else { return; }; for stored in &mut thread.events { - if matches!( - stored.event.event_type.as_str(), - "approval.requested" | "request.requested" | "ask.requested" - ) && request_id(&stored.event).as_deref() == Some(resolved_request_id) + if stored.seq <= seq + && stored + .event + .request_kind() + .is_some_and(|kind| kind.is_request()) + && stored.event.request_id() == Some(resolved_request_id) { stored.cleanable = true; } @@ -379,7 +306,7 @@ impl DirectThreadManager { let can_pop = thread .events .get(thread.head) - .is_some_and(|stored| stored.event.seq <= min_cursor && stored.cleanable); + .is_some_and(|stored| stored.seq <= min_cursor && stored.cleanable); if !can_pop { break; } @@ -419,7 +346,7 @@ impl DirectThreadManager { let oldest_seq = thread .events .get(thread.head) - .map(|stored| stored.event.seq) + .map(|stored| stored.seq) .unwrap_or(thread.next_seq.saturating_add(1)); let slowest = thread .subscribers @@ -448,13 +375,13 @@ pub(crate) fn set_direct_thread_manager_app_handle(app: tauri::AppHandle) { pub(crate) fn append_direct_thread_event( thread_id: &str, - draft: DirectThreadRawEventDraft, -) -> DirectThreadRawEvent { + event: DirectThreadEvent, +) -> DirectThreadEvent { let (event, subscriber_ids) = { let mut manager = global_direct_thread_manager() .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()); - let event = manager.append(thread_id, draft); + let event = manager.append(thread_id, event); let subscriber_ids = manager.subscriber_ids(thread_id); (event, subscriber_ids) }; @@ -486,84 +413,48 @@ pub(crate) fn consume_direct_thread( .consume(subscription_id) } -fn request_id(event: &DirectThreadRawEvent) -> Option { - event - .payload - .get("requestId") - .and_then(Value::as_str) - .or_else(|| event.payload.get("id").and_then(Value::as_str)) - .filter(|value| !value.is_empty()) - .map(str::to_string) -} - #[cfg(test)] mod tests { use super::*; + use crate::agent::{DirectThreadDeltaKind, DirectThreadItem, DirectThreadRequestKind}; - fn draft(event_type: &str, turn_id: &str, item_id: Option<&str>) -> DirectThreadRawEventDraft { - DirectThreadRawEventDraft { - event_type: event_type.to_string(), - turn_id: turn_id.to_string(), - item_id: item_id.map(str::to_string), - call_id: None, - payload: serde_json::json!({}), + fn message(item_id: &str) -> DirectThreadItem { + DirectThreadItem::Message { + item_id: item_id.to_string(), + role: "assistant".to_string(), + text: "内容".to_string(), + at: 0, } } - /// 工具条目的两个 id 空间必须靠 `callId` 对齐:app-server `item.started` 的 itemId 是调用 id, - /// 原始 item 的 `item.completed` 的 itemId 是 response item id。 - fn draft_with_call_id( - event_type: &str, - turn_id: &str, - item_id: Option<&str>, - call_id: Option<&str>, - ) -> DirectThreadRawEventDraft { - DirectThreadRawEventDraft { - event_type: event_type.to_string(), - turn_id: turn_id.to_string(), - item_id: item_id.map(str::to_string), - call_id: call_id.map(str::to_string), - payload: serde_json::json!({}), - } + fn item_started(item_id: &str) -> DirectThreadEvent { + DirectThreadEvent::item_started(message(item_id)) } - #[test] - fn call_id_joins_started_and_completed_across_id_spaces() { - let mut manager = DirectThreadManager::with_limits(100, 100_000); - manager.append( - "thread-1", - draft_with_call_id("item.started", "turn-1", Some("call-1"), None), - ); - manager.append( - "thread-1", - draft_with_call_id("item.completed", "turn-1", Some("item-9"), Some("call-1")), - ); - // 活跃条目按合并身份清理,不留下永远收不到完成事件的幽灵条目。 - let bootstrap = manager.subscribe("thread-1"); - assert!( - bootstrap - .events - .iter() - .all(|event| event.event_type != "item.started"), - "已完成的条目不得再作为运行态事件回到 bootstrap" - ); - // 没有订阅者时,可回收事件全部被回收;历史锚点靠独立字段保留,不占队列。 - assert_eq!( - manager.thread_debug("thread-1").map(|debug| debug.0), - Some(0) - ); - // 历史锚点仍然是 response item id,前端才能拿它当分页锚点。 - assert_eq!(bootstrap.last_completed_item_id.as_deref(), Some("item-9")); + fn item_completed(item_id: &str) -> DirectThreadEvent { + DirectThreadEvent::item_completed(message(item_id)) + } + + fn item_delta(item_id: &str) -> DirectThreadEvent { + DirectThreadEvent::item_delta( + item_id.to_string(), + DirectThreadDeltaKind::Message, + "增量".to_string(), + ) + } + + fn request(kind: DirectThreadRequestKind, request_id: Option<&str>) -> DirectThreadEvent { + DirectThreadEvent::request(kind, request_id.map(str::to_string)) } #[test] fn subscribers_have_independent_cursors_on_one_global_queue() { let mut manager = DirectThreadManager::with_limits(100, 100_000); - manager.append("thread-1", draft("turn.started", "turn-1", None)); + manager.append("thread-1", DirectThreadEvent::turn_started()); let first = manager.subscribe("thread-1"); let second = manager.subscribe("thread-1"); - manager.append("thread-1", draft("item.started", "turn-1", Some("item-1"))); - manager.append("thread-1", draft("item.delta", "turn-1", Some("item-1"))); + manager.append("thread-1", item_started("item-1")); + manager.append("thread-1", item_delta("item-1")); let first_batch = manager .consume(&first.subscription_id) @@ -583,52 +474,49 @@ mod tests { #[test] fn bootstrap_contains_lifecycle_anchor_and_unfinished_events_only() { let mut manager = DirectThreadManager::with_limits(100, 100_000); - manager.append("thread-1", draft("turn.started", "turn-1", None)); - manager.append("thread-1", draft("item.started", "turn-1", Some("item-1"))); - manager.append("thread-1", draft("item.delta", "turn-1", Some("item-1"))); - manager.append( - "thread-1", - draft("item.completed", "turn-1", Some("item-1")), - ); - manager.append("thread-1", draft("item.started", "turn-1", Some("item-2"))); + manager.append("thread-1", DirectThreadEvent::turn_started()); + manager.append("thread-1", item_started("item-1")); + manager.append("thread-1", item_delta("item-1")); + manager.append("thread-1", item_completed("item-1")); + manager.append("thread-1", item_started("item-2")); let bootstrap = manager.subscribe("thread-1"); - assert_eq!(bootstrap.last_completed_item_id.as_deref(), Some("item-1")); - assert_eq!( - bootstrap - .events - .iter() - .map(|event| event.event_type.as_str()) - .collect::>(), - vec!["turn.started", "item.started"] - ); + assert!(matches!( + bootstrap.events.as_slice(), + [ + DirectThreadEvent::TurnStarted {}, + DirectThreadEvent::ItemStarted { item, .. }, + ] if item.item_id() == "item-2" + )); } #[test] fn completion_releases_item_events_only_after_the_completion_event_is_appended() { let mut manager = DirectThreadManager::with_limits(100, 100_000); - manager.append("thread-1", draft("item.started", "turn-1", Some("item-1"))); - manager.append("thread-1", draft("item.delta", "turn-1", Some("item-1"))); + manager.append("thread-1", item_started("item-1")); + manager.append("thread-1", item_delta("item-1")); let bootstrap = manager.subscribe("thread-1"); - manager.append( - "thread-1", - draft("item.completed", "turn-1", Some("item-1")), - ); + manager.append("thread-1", item_completed("item-1")); let events = manager .consume(&bootstrap.subscription_id) .expect("consume completion") .events; - assert_eq!(events.len(), 1); - assert_eq!(events[0].event_type, "item.completed"); + assert!(matches!( + events.as_slice(), + [DirectThreadEvent::ItemCompleted { item, .. }] if item.item_id() == "item-1" + )); } #[test] fn slow_subscriber_is_expired_when_queue_limit_is_reached() { let mut manager = DirectThreadManager::with_limits(2, 100_000); let subscription = manager.subscribe("thread-1"); - manager.append("thread-1", draft("approval.resolved", "turn-1", None)); - manager.append("thread-1", draft("approval.resolved", "turn-1", None)); - manager.append("thread-1", draft("approval.resolved", "turn-1", None)); + for _ in 0..3 { + manager.append( + "thread-1", + request(DirectThreadRequestKind::RequestResolved, None), + ); + } assert_eq!( manager.consume(&subscription.subscription_id), Err(SUBSCRIPTION_EXPIRED.to_string()) @@ -638,11 +526,11 @@ mod tests { #[test] fn current_subscriber_is_not_expired_by_pinned_queue_head() { let mut manager = DirectThreadManager::with_limits(2, 100_000); - manager.append("thread-1", draft("item.started", "turn-1", Some("item-1"))); + manager.append("thread-1", item_started("item-1")); let subscription = manager.subscribe("thread-1"); - manager.append("thread-1", draft("item.delta", "turn-1", Some("item-1"))); - manager.append("thread-1", draft("item.delta", "turn-1", Some("item-1"))); - manager.append("thread-1", draft("item.delta", "turn-1", Some("item-1"))); + for _ in 0..3 { + manager.append("thread-1", item_delta("item-1")); + } assert_ne!( manager.consume(&subscription.subscription_id), Err(SUBSCRIPTION_EXPIRED.to_string()) @@ -654,25 +542,16 @@ mod tests { let mut manager = DirectThreadManager::with_limits(100, 100_000); manager.append( "thread-1", - DirectThreadRawEventDraft { - event_type: "approval.requested".to_string(), - turn_id: "turn-1".to_string(), - item_id: None, - call_id: None, - payload: serde_json::json!({"requestId": "request-1"}), - }, + request( + DirectThreadRequestKind::ApprovalRequested, + Some("request-1"), + ), ); let bootstrap = manager.subscribe("thread-1"); assert_eq!(bootstrap.events.len(), 1); manager.append( "thread-1", - DirectThreadRawEventDraft { - event_type: "approval.resolved".to_string(), - turn_id: "turn-1".to_string(), - item_id: None, - call_id: None, - payload: serde_json::json!({"requestId": "request-1"}), - }, + request(DirectThreadRequestKind::RequestResolved, Some("request-1")), ); assert_eq!( manager @@ -689,24 +568,15 @@ mod tests { let mut manager = DirectThreadManager::with_limits(100, 100_000); manager.append( "thread-1", - DirectThreadRawEventDraft { - event_type: "approval.requested".to_string(), - turn_id: "turn-1".to_string(), - item_id: None, - call_id: None, - payload: serde_json::json!({"requestId": "request-1"}), - }, + request( + DirectThreadRequestKind::ApprovalRequested, + Some("request-1"), + ), ); let subscription = manager.subscribe("thread-1"); manager.append( "thread-1", - DirectThreadRawEventDraft { - event_type: "approval.resolved".to_string(), - turn_id: "turn-1".to_string(), - item_id: None, - call_id: None, - payload: serde_json::json!({"requestId": "request-1"}), - }, + request(DirectThreadRequestKind::RequestResolved, Some("request-1")), ); manager .consume(&subscription.subscription_id) @@ -717,17 +587,25 @@ mod tests { #[test] fn turn_completed_anchor_survives_empty_queue_for_new_subscriber() { let mut manager = DirectThreadManager::with_limits(100, 100_000); - manager.append("thread-1", draft("turn.completed", "turn-1", None)); + manager.append( + "thread-1", + DirectThreadEvent::turn_completed("completed".to_string()), + ); let bootstrap = manager.subscribe("thread-1"); - assert_eq!(bootstrap.events.len(), 1); - assert_eq!(bootstrap.events[0].event_type, "turn.completed"); + assert!(matches!( + bootstrap.events.as_slice(), + [DirectThreadEvent::TurnCompleted { status }] if status == "completed" + )); } #[test] fn queue_cleanup_only_removes_a_cleanable_prefix() { let mut manager = DirectThreadManager::with_limits(100, 100_000); - manager.append("thread-1", draft("item.started", "turn-1", Some("item-1"))); - manager.append("thread-1", draft("approval.resolved", "turn-1", None)); + manager.append("thread-1", item_started("item-1")); + manager.append( + "thread-1", + request(DirectThreadRequestKind::RequestResolved, None), + ); let subscription = manager.subscribe("thread-1"); manager .consume(&subscription.subscription_id) diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_raw_item.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_raw_item.rs deleted file mode 100644 index 2dc96f3fe..000000000 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_raw_item.rs +++ /dev/null @@ -1,581 +0,0 @@ -//! DirectProject 「原始条目」搬运层。 -//! -//! Thread Manager 只做搬运:这里把 Codex 原始 response item 与 app-server thread item -//! 收敛成同一形状的**脱敏原始条目**,只做三件事——挑字段、脱敏、截断。 -//! -//! 这里不生成工具卡片,不推导标题 / 折叠摘要 / 可见性 / 排序:`title`、`summary`、 -//! `kind` 和"哪些条目要显示"全部是前端投影的职责。运行态事件与历史切片共用本形状, -//! 前端因此只需要一套投影与一套合并规则。 - -use crate::agent::redact_secret_tokens; -use crate::agent::sanitize_error_context; -use crate::redact_absolute_path_tokens; -use serde::Serialize; -use serde_json::Value; -use std::path::Path; - -/// 正文(消息 / 思考)上限。 -const DIRECT_THREAD_RAW_TEXT_MAX_CHARS: usize = 8000; -/// 工具明细(命令 / 参数 / 输出)上限。 -const DIRECT_THREAD_RAW_DETAIL_MAX_CHARS: usize = 4000; -/// 单条变更路径上限。 -const DIRECT_THREAD_RAW_PATH_MAX_CHARS: usize = 300; - -/// 原始条目里的一条文件变更。 -#[derive(Clone, Debug, Eq, PartialEq, Serialize)] -#[serde(rename_all = "camelCase")] -pub(crate) struct DirectThreadRawChange { - pub(crate) path: String, - /// 原始 kind:`add` | `update` | `delete`(缺失时为 `update`)。 - pub(crate) kind: String, -} - -/// 前端唯一消费的条目形状:脱敏、截断后的 Codex 原始条目。 -#[derive(Clone, Debug, Eq, PartialEq, Serialize)] -#[serde(rename_all = "camelCase")] -pub(crate) struct DirectThreadRawItem { - /// `project.jsonl` 里的 response item id;app-server `item/started` 时是 app-server item id。 - pub(crate) item_id: String, - /// 工具条目的调用 id:两个 id 空间唯一对齐点(原始 item 的 `call_id`)。 - #[serde(skip_serializing_if = "Option::is_none")] - pub(crate) call_id: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub(crate) turn_id: Option, - /// Codex 原始 `type`:`message` / `reasoning` / `function_call` / `function_call_output` / - /// `commandExecution` / `fileChange` / `mcpToolCall` / `webSearch` / …(原样透传,不做映射)。 - pub(crate) item_type: String, - #[serde(skip_serializing_if = "Option::is_none")] - pub(crate) role: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub(crate) text: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub(crate) name: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub(crate) arguments: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub(crate) output: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub(crate) command: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub(crate) tool: Option, - #[serde(skip_serializing_if = "Vec::is_empty")] - pub(crate) changes: Vec, - #[serde(skip_serializing_if = "Option::is_none")] - pub(crate) item_status: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub(crate) exit_code: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub(crate) success: Option, - /// 协议事实:来自 `item/completed` / `rawResponseItem/completed` 为 true。 - pub(crate) completed: bool, - /// 毫秒时间;条目自带时间与记录时间都拿不到时为 0。 - pub(crate) at: u64, -} - -fn bounded(value: &str, max_chars: usize) -> String { - if value.chars().count() <= max_chars { - return value.to_string(); - } - let mut truncated = value.chars().take(max_chars).collect::(); - truncated.push('…'); - truncated -} - -/// 原始值转文本:字符串原样,其它 JSON 值序列化。 -fn value_text(value: &Value) -> Option { - match value { - Value::Null => None, - Value::String(text) => (!text.trim().is_empty()).then(|| text.trim().to_string()), - other => serde_json::to_string_pretty(other).ok(), - } -} - -fn detail_text(root: &Path, value: &str) -> String { - bounded( - &sanitize_detail_text(root, value), - DIRECT_THREAD_RAW_DETAIL_MAX_CHARS, - ) -} - -fn item_text(root: &Path, item: &Value) -> Option { - let raw = item - .get("text") - .and_then(Value::as_str) - .map(str::to_string) - .or_else(|| { - for key in ["content", "summary"] { - let Some(parts) = item.get(key).and_then(Value::as_array) else { - continue; - }; - let joined = parts - .iter() - .filter_map(|part| part.get("text").and_then(Value::as_str)) - .collect::>() - .join(""); - if !joined.trim().is_empty() { - return Some(joined); - } - } - None - }) - .filter(|text| !text.trim().is_empty())?; - Some(bounded( - &sanitize_detail_text(root, &raw), - DIRECT_THREAD_RAW_TEXT_MAX_CHARS, - )) -} - -fn item_turn_id(item: &Value) -> Option { - let from_metadata = item - .get("internal_chat_message_metadata_passthrough") - .and_then(|meta| meta.get("turn_id").or_else(|| meta.get("turnId"))) - .and_then(Value::as_str) - .map(str::trim) - .filter(|value| !value.is_empty()) - .map(str::to_string); - if from_metadata.is_some() { - return from_metadata; - } - // 用户条目没有元数据,但 id 里带回合:`direct-codex::user`。 - let rest = item - .get("id") - .and_then(Value::as_str)? - .trim() - .strip_prefix("direct-codex:")?; - let turn = rest.rsplit_once(':')?.0.trim(); - (!turn.is_empty()).then(|| turn.to_string()) -} - -fn item_at_ms(item: &Value, observed_at_ms: u64) -> u64 { - let from_metadata = item - .get("internal_chat_message_metadata_passthrough") - .and_then(|meta| meta.get("create_time")) - .and_then(Value::as_f64) - .map(|seconds| (seconds * 1000.0).clamp(0.0, u64::MAX as f64) as u64) - .unwrap_or_default(); - if from_metadata > 0 { - return from_metadata; - } - for key in ["startedAtMs", "completedAtMs"] { - let value = item.get(key).and_then(Value::as_u64).unwrap_or_default(); - if value > 0 { - return value; - } - } - observed_at_ms -} - -fn item_identity(item: &Value) -> Option<(String, Option)> { - let call_id = item - .get("call_id") - .or_else(|| item.get("callId")) - .and_then(Value::as_str) - .map(str::trim) - .filter(|value| !value.is_empty()) - .map(str::to_string); - let id = item - .get("id") - .and_then(Value::as_str) - .map(str::trim) - .filter(|value| !value.is_empty()) - .map(str::to_string); - match (id, call_id) { - (Some(id), call_id) => Some((id, call_id)), - (None, Some(call_id)) => Some((call_id.clone(), Some(call_id))), - (None, None) => None, - } -} - -fn item_changes(root: &Path, item: &Value) -> Vec { - item.get("changes") - .and_then(Value::as_array) - .map(|changes| { - changes - .iter() - .filter_map(|change| { - let path = change - .get("path") - .and_then(Value::as_str) - .map(str::trim) - .filter(|path| !path.is_empty())?; - Some(DirectThreadRawChange { - path: bounded( - &sanitize_detail_text(root, path), - DIRECT_THREAD_RAW_PATH_MAX_CHARS, - ), - kind: change - .get("kind") - .and_then(Value::as_str) - .unwrap_or("update") - .to_string(), - }) - }) - .collect::>() - }) - .unwrap_or_default() -} - -/// 把一条 Codex 条目搬运成脱敏原始条目;拿不到身份或类型时返回 `None`。 -/// -/// `observed_at_ms` 只作为条目自带时间缺失时的兜底(运行态用当前时间,历史用文件记录时间)。 -pub(crate) fn direct_thread_raw_item_from_value( - root: &Path, - item: &Value, - turn_id: Option<&str>, - completed: bool, - observed_at_ms: u64, -) -> Option { - if !item.is_object() { - return None; - } - let (item_id, call_id) = item_identity(item)?; - let item_type = item - .get("type") - .and_then(Value::as_str) - .map(str::trim) - .filter(|value| !value.is_empty())? - .to_string(); - let turn_id = turn_id - .map(str::trim) - .filter(|value| !value.is_empty()) - .map(str::to_string) - .or_else(|| item_turn_id(item)); - let at = item_at_ms(item, observed_at_ms); - - let text = match item_type.as_str() { - "message" | "agentMessage" | "userMessage" | "reasoning" => item_text(root, item), - _ => None, - }; - let arguments = item - .get("arguments") - .and_then(value_text) - .map(|value| detail_text(root, &value)); - let output = ["output", "aggregatedOutput", "result", "error"] - .iter() - .find_map(|key| item.get(key).and_then(value_text)) - .map(|value| detail_text(root, &value)); - let command = item - .get("command") - .and_then(Value::as_str) - .map(str::trim) - .filter(|command| !command.is_empty()) - .map(|command| detail_text(root, command)); - - Some(DirectThreadRawItem { - item_id, - call_id, - turn_id, - item_type, - role: item - .get("role") - .and_then(Value::as_str) - .map(str::trim) - .filter(|role| !role.is_empty()) - .map(str::to_string), - text, - name: item - .get("name") - .and_then(Value::as_str) - .map(str::trim) - .filter(|name| !name.is_empty()) - .map(str::to_string), - arguments, - output, - command, - tool: item - .get("tool") - .and_then(Value::as_str) - .map(str::trim) - .filter(|tool| !tool.is_empty()) - .map(|tool| sanitize_detail_text(root, tool)), - changes: item_changes(root, item), - item_status: item - .get("status") - .and_then(Value::as_str) - .map(str::trim) - .filter(|status| !status.is_empty()) - .map(str::to_string), - exit_code: item.get("exitCode").and_then(Value::as_i64), - success: item.get("success").and_then(Value::as_bool), - completed, - at, - }) -} - -/// 历史切片搬运:保持文件顺序,不做任何合并。 -/// -/// 同一调用的 `function_call` 与 `function_call_output` 是两条原始条目,合并属于前端。 -/// `timestamp_of` 是文件记录时间,仅在条目自带时间缺失时兜底。 -pub(crate) fn direct_thread_raw_items_from_history( - root: &Path, - items: &[Value], - timestamp_of: impl Fn(&Value) -> u64, -) -> Vec { - items - .iter() - .filter_map(|item| { - direct_thread_raw_item_from_value(root, item, None, true, timestamp_of(item)) - }) - .collect() -} - -/// 项目根目录之后的路径 token:分隔符统一成 `/`,返回 `(消费到的下标, 项目相对路径)`。 -fn project_relative_path_segment(value: &str, start: usize) -> (usize, String) { - let mut index = start; - let mut relative = String::new(); - while index < value.len() { - let character = value[index..].chars().next().unwrap_or_default(); - if matches!(character, '/' | '\\') { - if !relative.is_empty() { - relative.push('/'); - } - index += character.len_utf8(); - continue; - } - if character.is_whitespace() - || matches!( - character, - '\'' | '"' - | '`' - | ',' - | ';' - | '|' - | '&' - | '(' - | ')' - | '[' - | ']' - | '{' - | '}' - | '<' - | '>' - | ':' - ) - { - break; - } - relative.push(character); - index += character.len_utf8(); - } - while relative.ends_with('/') { - relative.pop(); - } - (index, relative) -} - -/// 把项目根目录前缀换成**项目相对路径**(`/game/src/x.ts` → `game/src/x.ts`)。 -/// -/// 必须排在 `redact_absolute_path_tokens` 之前:后者会把整个绝对路径抹成 -/// ``,之后就再也认不出哪些路径在项目内了。 -/// Windows 上同时匹配 `\` 与 `/` 两种分隔符写法,并按大小写不敏感比较(盘符大小写会变)。 -fn relativize_project_root_paths(root: &Path, value: &str) -> String { - let root_text = root.to_string_lossy(); - let root_text = root_text.trim_end_matches(['/', '\\']); - if root_text.is_empty() { - return value.to_string(); - } - let mut needles = [ - root_text.to_string(), - root_text.replace('\\', "/"), - root_text.replace('/', "\\"), - ] - .into_iter() - .map(|needle| needle.to_ascii_lowercase()) - .filter(|needle| !needle.is_empty()) - .collect::>(); - needles.sort(); - needles.dedup(); - let lower = value.to_ascii_lowercase(); - - let mut output = String::with_capacity(value.len()); - let mut cursor = 0usize; - while cursor < value.len() { - let mut hit: Option<(usize, usize)> = None; - for needle in &needles { - let mut search = cursor; - while let Some(relative) = lower[search..].find(needle.as_str()) { - let start = search + relative; - let end = start + needle.len(); - let left_is_boundary = start == 0 - || lower[..start].chars().next_back().is_some_and(|character| { - !character.is_alphanumeric() && character != '_' && character != '-' - }); - if left_is_boundary && value[end..].starts_with(['/', '\\']) { - if hit.is_none_or(|(best_start, _)| start < best_start) { - hit = Some((start, end)); - } - break; - } - search = end; - } - } - let Some((start, end)) = hit else { - break; - }; - output.push_str(&value[cursor..start]); - let (consumed, relative) = project_relative_path_segment(value, end); - if relative.is_empty() { - // 只写了项目根目录本身(没有后续路径段):按占位形状处理。 - output.push_str(""); - } else { - output.push_str(&relative); - } - cursor = consumed; - } - output.push_str(&value[cursor..]); - output -} - -/// 脱敏:项目内绝对路径先归一化成项目相对路径,再依次做绝对路径、密钥前缀与 -/// 错误上下文脱敏。 -/// -/// 顺序不能反:先抹密钥会把 `sk-…` 之类的 token 换成占位符,但绝对路径里的用户名目录 -/// 仍然会留下;这里先归一化路径 token,再处理密钥。 -pub(crate) fn sanitize_detail_text(root: &Path, value: &str) -> String { - let without_project_root = relativize_project_root_paths(root, value); - let without_absolute = redact_absolute_path_tokens(&without_project_root); - let without_secret = redact_secret_tokens(&without_absolute); - sanitize_error_context(&without_secret) -} - -#[cfg(test)] -mod tests { - use super::*; - use serde_json::json; - use std::path::Path; - - fn root() -> &'static Path { - Path::new(".") - } - - #[test] - fn message_item_carries_role_text_and_turn() { - let item = direct_thread_raw_item_from_value( - root(), - &json!({ - "id": "direct-codex:turn-1:user", - "type": "message", - "role": "user", - "content": [{"type": "input_text", "text": "做一个拼图游戏"}], - }), - None, - true, - 0, - ) - .expect("user item"); - assert_eq!(item.item_type, "message"); - assert_eq!(item.role.as_deref(), Some("user")); - assert_eq!(item.text.as_deref(), Some("做一个拼图游戏")); - assert_eq!(item.turn_id.as_deref(), Some("turn-1")); - } - - #[test] - fn tool_item_keeps_call_id_and_raw_command() { - let item = direct_thread_raw_item_from_value( - root(), - &json!({ - "id": "05dc0af1-8023-47fd-ad22-d54df2837b1b", - "call_id": "call_00_Gpd0s0Ytm9YgIbwbEXva1473", - "type": "function_call", - "name": "exec_command", - "arguments": "{\"cmd\": \"ls\"}", - "internal_chat_message_metadata_passthrough": {"turn_id": "turn-1"}, - }), - None, - true, - 1000, - ) - .expect("tool item"); - assert_eq!(item.item_type, "function_call"); - assert_eq!( - item.call_id.as_deref(), - Some("call_00_Gpd0s0Ytm9YgIbwbEXva1473") - ); - assert_eq!(item.name.as_deref(), Some("exec_command")); - assert_eq!(item.arguments.as_deref(), Some("{\"cmd\": \"ls\"}")); - assert_eq!(item.turn_id.as_deref(), Some("turn-1")); - } - - #[test] - fn reasoning_item_keeps_text_without_role() { - let item = direct_thread_raw_item_from_value( - root(), - &json!({ - "id": "reason-1", - "type": "reasoning", - "content": [{"type": "output_text", "text": "先看目录"}], - }), - Some("turn-1"), - true, - 0, - ) - .expect("reasoning item"); - assert_eq!(item.item_type, "reasoning"); - assert_eq!(item.role, None); - assert_eq!(item.text.as_deref(), Some("先看目录")); - } - - #[test] - fn secrets_and_absolute_paths_are_not_leaked() { - let item = direct_thread_raw_item_from_value( - root(), - &json!({ - "id": "msg-1", - "type": "message", - "role": "assistant", - "content": [{"type": "output_text", "text": "key=sk-abcdefghijklmnop at /root/secret/x"}], - }), - Some("turn-1"), - true, - 0, - ) - .expect("assistant item"); - let text = item.text.unwrap_or_default(); - assert!( - !text.contains("sk-abcdefghijklmnop"), - "不得泄漏明文密钥:{text}" - ); - assert!(!text.contains("/root/secret"), "不得泄漏绝对路径:{text}"); - } - - #[test] - fn history_keeps_call_and_output_as_two_raw_items() { - let items = vec![ - json!({ - "id": "05dc0af1-8023-47fd-ad22-d54df2837b1b", - "call_id": "call_00_Gpd0s0Ytm9YgIbwbEXva1473", - "type": "function_call", - "name": "exec_command", - "arguments": "{\"cmd\": \"ls\"}", - }), - json!({ - "id": "fco_01a06fa5-d636-7452-b337-a641c2e6bc76", - "call_id": "call_00_Gpd0s0Ytm9YgIbwbEXva1473", - "type": "function_call_output", - "output": "assets\ngame\n", - }), - ]; - let projected = direct_thread_raw_items_from_history(root(), &items, |_| 0); - assert_eq!(projected.len(), 2, "搬运层不得替前端做合并"); - assert_eq!(projected[0].item_type, "function_call"); - assert_eq!(projected[1].item_type, "function_call_output"); - assert_eq!( - projected[0].call_id.as_deref(), - projected[1].call_id.as_deref() - ); - } - - #[test] - fn system_role_is_passed_through_for_frontend_to_filter() { - let item = direct_thread_raw_item_from_value( - root(), - &json!({"id": "sys-1", "type": "message", "role": "system", "content": [{"type": "input_text", "text": "x"}]}), - Some("turn-1"), - true, - 0, - ) - .expect("system item"); - assert_eq!(item.role.as_deref(), Some("system")); - assert_eq!(item.item_type, "message"); - } -} diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_wire.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_wire.rs new file mode 100644 index 000000000..d9374a3a0 --- /dev/null +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_wire.rs @@ -0,0 +1,864 @@ +//! DirectProject 聊天事件的线上模型与投影。 +//! +//! 前端消费的类型由 ts-rs 导出到 `src/features/project-workspace/generated/`, +//! 与 Rust 定义同源:加一个字段不会只改一边。 +//! +//! 本模块只做三件事:挑字段、脱敏、截断。工具卡片的 `kind`、标题、折叠摘要、可见性与 +//! 合并规则全部属于前端投影,这里一概不出现。 +//! +//! 条目身份在进队列前就归一成**一个** `itemId`:原始文件里工具条目带两个 id(app-server +//! 的调用 id 与 response item id,调用与输出共用前者),归一只在 Rust 边界做一次, +//! Thread Manager 与前端都只认这一个,不暴露第二个 id 概念。 +//! 历史分页锚点是另一回事,那是文件里的原始 item id,单独取。 + +use crate::agent::redact_secret_tokens; +use crate::agent::sanitize_error_context; +use crate::redact_absolute_path_tokens; +use serde::{Deserialize, Serialize}; +use serde_json::Value; +use std::path::Path; +use ts_rs::TS; + +/// 正文(消息 / 思考)上限。 +const DIRECT_THREAD_TEXT_MAX_CHARS: usize = 8_000; +/// 工具明细(命令 / 参数 / 输出)上限。 +const DIRECT_THREAD_DETAIL_MAX_CHARS: usize = 4_000; +/// 单条变更路径上限。 +const DIRECT_THREAD_PATH_MAX_CHARS: usize = 300; + +/// 一条文件变更。 +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, TS)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] +#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/features/project-workspace/generated/"))] +pub(crate) struct DirectThreadFileChange { + pub(crate) path: String, + /// `add` | `update` | `delete` + pub(crate) kind: String, +} + +/// 聊天视图的输入条目:一条 Codex 原始条目的脱敏投影。 +/// +/// `itemType` 就是 Codex 的原始类型,逐字透传;前端按它决定投影成消息、思考还是工具卡片。 +/// 未识别的类型走 [`DirectThreadItem::Other`],Rust 不替前端决定它是否可见。 +/// +/// 条目上的 `at` 是只用于显示的毫秒时间戳:ts-rs 默认把 `u64` 映射成 `bigint`, +/// 而 Tauri 的 JSON 通道传过来的是 `number`,因此统一标 `#[ts(as = "f64")]` 对齐。 +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, TS)] +#[serde(tag = "itemType", rename_all_fields = "camelCase", deny_unknown_fields)] +#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/features/project-workspace/generated/"))] +pub(crate) enum DirectThreadItem { + #[serde(rename = "message")] + Message { + /// 归一身份:全链路只有这一个 id。 + item_id: String, + /// 原始 role(`user` / `assistant` / `system` / …);显示与否由前端判断。 + role: String, + text: String, + #[ts(as = "f64")] + at: u64, + }, + #[serde(rename = "reasoning")] + Reasoning { + item_id: String, + text: String, + #[ts(as = "f64")] + at: u64, + }, + /// 原始 response item 的工具调用:参数在 `arguments`,输出在后续的 + /// [`DirectThreadItem::FunctionCallOutput`](两者共用归一身份)。 + #[serde(rename = "function_call")] + FunctionCall { + item_id: String, + name: String, + arguments: String, + #[ts(as = "f64")] + at: u64, + }, + #[serde(rename = "function_call_output")] + FunctionCallOutput { + item_id: String, + output: String, + #[ts(as = "f64")] + at: u64, + }, + #[serde(rename = "commandExecution")] + CommandExecution { + item_id: String, + command: String, + #[serde(default)] + output: Option, + /// app-server 原始状态:`inProgress` / `completed` / `failed` / `declined` / … + #[serde(default)] + status: Option, + #[serde(default)] + #[ts(as = "Option")] + exit_code: Option, + #[ts(as = "f64")] + at: u64, + }, + #[serde(rename = "fileChange")] + FileChange { + item_id: String, + changes: Vec, + #[ts(as = "f64")] + at: u64, + }, + #[serde(rename = "mcpToolCall")] + McpToolCall { + item_id: String, + tool: String, + arguments: String, + #[serde(default)] + output: Option, + #[serde(default)] + status: Option, + #[ts(as = "f64")] + at: u64, + }, + #[serde(rename = "webSearch")] + WebSearch { + item_id: String, + #[serde(default)] + query: Option, + #[serde(default)] + output: Option, + #[ts(as = "f64")] + at: u64, + }, + #[serde(rename = "contextCompaction")] + ContextCompaction { + item_id: String, + #[ts(as = "f64")] + at: u64, + }, + /// 未识别的 Codex item 类型:原样透传身份与类型,不投影正文。 + #[serde(rename = "other")] + Other { + item_id: String, + raw_type: String, + #[ts(as = "f64")] + at: u64, + }, +} + +impl DirectThreadItem { + /// 归一身份:Thread Manager 用它登记与释放未完成条目,前端用它合并同一张卡片。 + pub(crate) fn item_id(&self) -> &str { + match self { + Self::Message { item_id, .. } + | Self::Reasoning { item_id, .. } + | Self::FunctionCall { item_id, .. } + | Self::FunctionCallOutput { item_id, .. } + | Self::CommandExecution { item_id, .. } + | Self::FileChange { item_id, .. } + | Self::McpToolCall { item_id, .. } + | Self::WebSearch { item_id, .. } + | Self::ContextCompaction { item_id, .. } + | Self::Other { item_id, .. } => item_id, + } + } +} + +/// 增量正文属于哪类条目。 +#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize, TS)] +#[serde(rename_all = "camelCase")] +#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/features/project-workspace/generated/"))] +pub(crate) enum DirectThreadDeltaKind { + /// assistant 正文。 + Message, + /// 思考正文。 + Reasoning, +} + +/// 审批 / 提问请求与解决:本轮只透传,不并入聊天状态。 +#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize, TS)] +#[serde(rename_all = "camelCase")] +#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/features/project-workspace/generated/"))] +pub(crate) enum DirectThreadRequestKind { + #[serde(rename = "approval.requested")] + ApprovalRequested, + #[serde(rename = "ask.requested")] + AskRequested, + #[serde(rename = "request.resolved")] + RequestResolved, +} + +impl DirectThreadRequestKind { + /// 未解决的请求要留在 bootstrap 里,直到出现对应的解决事件。 + pub(crate) fn is_request(&self) -> bool { + matches!(self, Self::ApprovalRequested | Self::AskRequested) + } + + pub(crate) fn is_resolution(&self) -> bool { + !self.is_request() + } +} + +/// Thread Manager 下发的运行态事件。 +/// +/// 顺序由数组顺序给出(同一个 subscriber 的 `consume` 按队列顺序返回),因此不需要 `seq`: +/// 游标是 Thread Manager 的内部事实,不下发。 +/// +/// 事件不带回合身份:DirectProject 同一时刻只有一个回合在跑,"当前回合是否还在跑"由 +/// 生命周期事件在序列中的位置给出,`turn_id` 对前端没有任何额外信息。 +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, TS)] +#[serde(tag = "type", rename_all_fields = "camelCase", deny_unknown_fields)] +#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/features/project-workspace/generated/"))] +pub(crate) enum DirectThreadEvent { + #[serde(rename = "turn.started")] + TurnStarted, + #[serde(rename = "turn.completed")] + TurnCompleted { status: String }, + #[serde(rename = "item.started")] + ItemStarted { item: DirectThreadItem }, + #[serde(rename = "item.completed")] + ItemCompleted { item: DirectThreadItem }, + #[serde(rename = "item.delta")] + ItemDelta { + item_id: String, + kind: DirectThreadDeltaKind, + delta: String, + }, + #[serde(rename = "request")] + Request { + kind: DirectThreadRequestKind, + #[serde(default)] + request_id: Option, + }, +} + +impl DirectThreadEvent { + pub(crate) fn turn_started() -> Self { + Self::TurnStarted + } + + pub(crate) fn turn_completed(status: String) -> Self { + Self::TurnCompleted { status } + } + + pub(crate) fn item_started(item: DirectThreadItem) -> Self { + Self::ItemStarted { item } + } + + pub(crate) fn item_completed(item: DirectThreadItem) -> Self { + Self::ItemCompleted { item } + } + + pub(crate) fn item_delta(item_id: String, kind: DirectThreadDeltaKind, delta: String) -> Self { + Self::ItemDelta { + item_id, + kind, + delta, + } + } + + pub(crate) fn request(kind: DirectThreadRequestKind, request_id: Option) -> Self { + Self::Request { kind, request_id } + } + + /// 事件关联的条目身份:只有 item 事件有。 + pub(crate) fn item_id(&self) -> Option<&str> { + match self { + Self::ItemStarted { item, .. } | Self::ItemCompleted { item, .. } => { + Some(item.item_id()) + } + _ => None, + } + } + + pub(crate) fn request_id(&self) -> Option<&str> { + match self { + Self::Request { request_id, .. } => request_id.as_deref(), + _ => None, + } + } + + pub(crate) fn request_kind(&self) -> Option { + match self { + Self::Request { kind, .. } => Some(*kind), + _ => None, + } + } +} + +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, TS)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] +#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/features/project-workspace/generated/"))] +pub(crate) struct DirectThreadSubscriptionBootstrap { + pub(crate) subscription_id: String, + /// 首屏历史锚点:`project.jsonl` 里最后一条原始 item id。 + #[serde(default)] + pub(crate) last_completed_item_id: Option, + /// 该 subscriber 此刻应当处理的运行态事件(游标已经在队尾)。 + pub(crate) events: Vec, +} + +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, TS)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] +#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/features/project-workspace/generated/"))] +pub(crate) struct DirectThreadConsumeResult { + pub(crate) events: Vec, +} + +#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, TS)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] +#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/features/project-workspace/generated/"))] +pub(crate) struct DirectThreadHistorySlice { + /// 脱敏条目,顺序即文件顺序;与运行态事件里的条目同形。 + pub(crate) items: Vec, + pub(crate) has_more: bool, + /// 本次切片的原始 item id 锚点:无论切片里有没有可显示条目,分页都靠它向前。 + #[serde(default)] + pub(crate) first_item_id: Option, +} + +fn bounded(value: &str, max_chars: usize) -> String { + if value.chars().count() <= max_chars { + return value.to_string(); + } + let mut truncated = value.chars().take(max_chars).collect::(); + truncated.push('…'); + truncated +} + +fn detail_text(root: &Path, value: &str) -> String { + bounded( + &sanitize_detail_text(root, value), + DIRECT_THREAD_DETAIL_MAX_CHARS, + ) +} + +/// 原始值转文本:字符串原样,其它 JSON 值序列化(调用方随后脱敏)。 +fn value_text(value: &Value) -> Option { + match value { + Value::Null => None, + Value::String(text) => (!text.trim().is_empty()).then(|| text.trim().to_string()), + other => serde_json::to_string_pretty(other).ok(), + } +} + +fn item_text(root: &Path, item: &Value) -> Option { + let raw = item + .get("text") + .and_then(Value::as_str) + .map(str::to_string) + .or_else(|| { + for key in ["content", "summary"] { + let Some(parts) = item.get(key).and_then(Value::as_array) else { + continue; + }; + let joined = parts + .iter() + .filter_map(|part| part.get("text").and_then(Value::as_str)) + .collect::>() + .join(""); + if !joined.trim().is_empty() { + return Some(joined); + } + } + None + }) + .filter(|text| !text.trim().is_empty())?; + Some(bounded( + &sanitize_detail_text(root, &raw), + DIRECT_THREAD_TEXT_MAX_CHARS, + )) +} + +fn item_at_ms(item: &Value, observed_at_ms: u64) -> u64 { + let from_metadata = item + .get("internal_chat_message_metadata_passthrough") + .and_then(|meta| meta.get("create_time")) + .and_then(Value::as_f64) + .map(|seconds| (seconds * 1000.0).clamp(0.0, u64::MAX as f64) as u64) + .unwrap_or_default(); + if from_metadata > 0 { + return from_metadata; + } + for key in ["startedAtMs", "completedAtMs"] { + let value = item.get(key).and_then(Value::as_u64).unwrap_or_default(); + if value > 0 { + return value; + } + } + observed_at_ms +} + +/// 归一身份:工具条目用工具调用 id,其它条目用自己的 `id`;只产出这一个值。 +pub(crate) fn direct_thread_item_identity(item: &Value) -> Option { + let call_id = item + .get("call_id") + .or_else(|| item.get("callId")) + .and_then(Value::as_str) + .map(str::trim) + .filter(|value| !value.is_empty()) + .map(str::to_string); + let id = item + .get("id") + .and_then(Value::as_str) + .map(str::trim) + .filter(|value| !value.is_empty()) + .map(str::to_string); + call_id.or(id) +} + +fn item_changes(root: &Path, item: &Value) -> Vec { + item.get("changes") + .and_then(Value::as_array) + .map(|changes| { + changes + .iter() + .filter_map(|change| { + let path = change + .get("path") + .and_then(Value::as_str) + .map(str::trim) + .filter(|path| !path.is_empty())?; + Some(DirectThreadFileChange { + path: bounded( + &sanitize_detail_text(root, path), + DIRECT_THREAD_PATH_MAX_CHARS, + ), + kind: change + .get("kind") + .and_then(Value::as_str) + .unwrap_or("update") + .to_string(), + }) + }) + .collect::>() + }) + .unwrap_or_default() +} + +fn field_text(root: &Path, item: &Value, key: &str) -> Option { + item.get(key) + .and_then(value_text) + .map(|value| detail_text(root, &value)) +} + +/// 把一条 Codex 原始条目投影成线上条目;拿不到身份或类型时返回 `None`。 +/// +/// `observed_at_ms` 只在条目自带时间缺失时兜底(运行态用当前时间,历史用文件记录时间)。 +pub(crate) fn direct_thread_item_from_value( + root: &Path, + item: &Value, + observed_at_ms: u64, +) -> Option { + if !item.is_object() { + return None; + } + let item_id = direct_thread_item_identity(item)?; + let item_type = item + .get("type") + .and_then(Value::as_str) + .map(str::trim) + .filter(|value| !value.is_empty())?; + let at = item_at_ms(item, observed_at_ms); + let text = item_text(root, item); + let role = item + .get("role") + .and_then(Value::as_str) + .map(str::trim) + .filter(|role| !role.is_empty()) + .map(str::to_string); + + Some(match item_type { + "message" | "agentMessage" | "userMessage" => DirectThreadItem::Message { + item_id, + role: role.unwrap_or_else(|| { + if item_type == "userMessage" { + "user".to_string() + } else { + "assistant".to_string() + } + }), + text: text?, + at, + }, + "reasoning" => DirectThreadItem::Reasoning { + item_id, + text: text?, + at, + }, + "function_call" => DirectThreadItem::FunctionCall { + item_id, + name: item + .get("name") + .and_then(Value::as_str) + .unwrap_or_default() + .to_string(), + arguments: field_text(root, item, "arguments").unwrap_or_default(), + at, + }, + "function_call_output" => DirectThreadItem::FunctionCallOutput { + item_id, + output: field_text(root, item, "output").unwrap_or_default(), + at, + }, + "commandExecution" => DirectThreadItem::CommandExecution { + item_id, + command: item + .get("command") + .and_then(Value::as_str) + .map(str::trim) + .filter(|command| !command.is_empty()) + .map(|command| detail_text(root, command)) + .unwrap_or_default(), + output: ["aggregatedOutput", "output", "error"] + .iter() + .find_map(|key| field_text(root, item, key)), + status: item + .get("status") + .and_then(Value::as_str) + .map(str::to_string), + exit_code: item.get("exitCode").and_then(Value::as_i64), + at, + }, + "fileChange" => DirectThreadItem::FileChange { + item_id, + changes: item_changes(root, item), + at, + }, + "mcpToolCall" => DirectThreadItem::McpToolCall { + item_id, + tool: item + .get("tool") + .and_then(Value::as_str) + .map(str::trim) + .filter(|tool| !tool.is_empty()) + .map(|tool| sanitize_detail_text(root, tool)) + .unwrap_or_default(), + arguments: field_text(root, item, "arguments").unwrap_or_default(), + output: ["result", "error"] + .iter() + .find_map(|key| field_text(root, item, key)), + status: item + .get("status") + .and_then(Value::as_str) + .map(str::to_string), + at, + }, + "webSearch" => DirectThreadItem::WebSearch { + item_id, + query: field_text(root, item, "query").or_else(|| { + item.get("action") + .and_then(|action| action.get("query")) + .and_then(value_text) + .map(|query| detail_text(root, &query)) + }), + output: field_text(root, item, "output"), + at, + }, + "contextCompaction" => DirectThreadItem::ContextCompaction { item_id, at }, + other => DirectThreadItem::Other { + item_id, + raw_type: other.to_string(), + at, + }, + }) +} + +/// 历史切片投影:保持文件顺序,不做任何合并(同一调用的调用与输出是两条条目)。 +/// +/// `timestamp_of` 是文件记录时间,仅在条目自带时间缺失时兜底。 +pub(crate) fn direct_thread_items_from_history( + root: &Path, + items: &[Value], + timestamp_of: impl Fn(&Value) -> u64, +) -> Vec { + items + .iter() + .filter_map(|item| direct_thread_item_from_value(root, item, timestamp_of(item))) + .collect() +} + +/// 项目根目录之后的路径 token:分隔符统一成 `/`,返回 `(消费到的下标, 项目相对路径)`。 +fn project_relative_path_segment(value: &str, start: usize) -> (usize, String) { + let mut index = start; + let mut relative = String::new(); + while index < value.len() { + let character = value[index..].chars().next().unwrap_or_default(); + if matches!(character, '/' | '\\') { + if !relative.is_empty() { + relative.push('/'); + } + index += character.len_utf8(); + continue; + } + if character.is_whitespace() + || matches!( + character, + '\'' | '"' + | '`' + | ',' + | ';' + | '|' + | '&' + | '(' + | ')' + | '[' + | ']' + | '{' + | '}' + | '<' + | '>' + | ':' + ) + { + break; + } + relative.push(character); + index += character.len_utf8(); + } + while relative.ends_with('/') { + relative.pop(); + } + (index, relative) +} + +/// 把项目根目录前缀换成**项目相对路径**(`/game/src/x.ts` → `game/src/x.ts`)。 +/// +/// 必须排在 `redact_absolute_path_tokens` 之前:后者会把整个绝对路径抹成 +/// ``,之后就再也认不出哪些路径在项目内了。 +/// Windows 上同时匹配 `\` 与 `/` 两种分隔符写法,并按大小写不敏感比较(盘符大小写会变)。 +fn relativize_project_root_paths(root: &Path, value: &str) -> String { + let root_text = root.to_string_lossy(); + let root_text = root_text.trim_end_matches(['/', '\\']); + if root_text.is_empty() { + return value.to_string(); + } + let mut needles = [ + root_text.to_string(), + root_text.replace('\\', "/"), + root_text.replace('/', "\\"), + ] + .into_iter() + .map(|needle| needle.to_ascii_lowercase()) + .filter(|needle| !needle.is_empty()) + .collect::>(); + needles.sort(); + needles.dedup(); + let lower = value.to_ascii_lowercase(); + + let mut output = String::with_capacity(value.len()); + let mut cursor = 0usize; + while cursor < value.len() { + let mut hit: Option<(usize, usize)> = None; + for needle in &needles { + let mut search = cursor; + while let Some(relative) = lower[search..].find(needle.as_str()) { + let start = search + relative; + let end = start + needle.len(); + let left_is_boundary = start == 0 + || lower[..start].chars().next_back().is_some_and(|character| { + !character.is_alphanumeric() && character != '_' && character != '-' + }); + if left_is_boundary && value[end..].starts_with(['/', '\\']) { + if hit.is_none_or(|(best_start, _)| start < best_start) { + hit = Some((start, end)); + } + break; + } + search = end; + } + } + let Some((start, end)) = hit else { + break; + }; + output.push_str(&value[cursor..start]); + let (consumed, relative) = project_relative_path_segment(value, end); + if relative.is_empty() { + // 只写了项目根目录本身(没有后续路径段):按占位形状处理。 + output.push_str(""); + } else { + output.push_str(&relative); + } + cursor = consumed; + } + output.push_str(&value[cursor..]); + output +} + +/// 脱敏:项目内绝对路径先归一化成项目相对路径,再依次做绝对路径、密钥前缀与 +/// 错误上下文脱敏。 +/// +/// 顺序不能反:先抹密钥会把 `sk-…` 之类的 token 换成占位符,但绝对路径里的用户名目录 +/// 仍然会留下;这里先归一化路径 token,再处理密钥。 +pub(crate) fn sanitize_detail_text(root: &Path, value: &str) -> String { + let without_project_root = relativize_project_root_paths(root, value); + let without_absolute = redact_absolute_path_tokens(&without_project_root); + let without_secret = redact_secret_tokens(&without_absolute); + sanitize_error_context(&without_secret) +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + use std::path::Path; + + fn root() -> &'static Path { + Path::new(".") + } + + #[test] + fn message_item_carries_role_text_and_turn() { + let item = direct_thread_item_from_value( + root(), + &json!({ + "id": "direct-codex:turn-1:user", + "type": "message", + "role": "user", + "content": [{"type": "input_text", "text": "做一个拼图游戏"}], + }), + 0, + ) + .expect("user item"); + assert_eq!( + item, + DirectThreadItem::Message { + item_id: "direct-codex:turn-1:user".to_string(), + role: "user".to_string(), + text: "做一个拼图游戏".to_string(), + at: 0, + } + ); + } + + #[test] + fn app_server_agent_message_defaults_to_assistant_role() { + let item = direct_thread_item_from_value( + root(), + &json!({"id": "msg-1", "type": "agentMessage", "text": "已执行"}), + 1000, + ) + .expect("agent message"); + assert!(matches!( + item, + DirectThreadItem::Message { role, .. } if role == "assistant" + )); + } + + #[test] + fn tool_item_identity_is_normalized_to_one_id() { + let item = direct_thread_item_from_value( + root(), + &json!({ + "id": "05dc0af1-8023-47fd-ad22-d54df2837b1b", + "call_id": "call_00_Gpd0s0Ytm9YgIbwbEXva1473", + "type": "function_call", + "name": "exec_command", + "arguments": "{\"cmd\": \"ls\"}", + "internal_chat_message_metadata_passthrough": {"turn_id": "turn-1"}, + }), + 1000, + ) + .expect("tool item"); + // 只有唯一身份:工具条目在文件里的另一个 id 不再对外暴露。 + assert_eq!(item.item_id(), "call_00_Gpd0s0Ytm9YgIbwbEXva1473"); + assert!(matches!( + item, + DirectThreadItem::FunctionCall { name, arguments, .. } + if name == "exec_command" && arguments == "{\"cmd\": \"ls\"}" + )); + } + + #[test] + fn command_execution_keeps_raw_status_and_exit_code() { + let item = direct_thread_item_from_value( + root(), + &json!({ + "id": "call-1", + "type": "commandExecution", + "command": "ls", + "status": "failed", + "exitCode": 2, + "aggregatedOutput": "boom", + }), + 0, + ) + .expect("command item"); + assert!(matches!( + item, + DirectThreadItem::CommandExecution { + command, + output: Some(output), + status: Some(status), + exit_code: Some(2), + .. + } if command == "ls" && output == "boom" && status == "failed" + )); + } + + #[test] + fn secrets_and_absolute_paths_are_not_leaked() { + let item = direct_thread_item_from_value( + root(), + &json!({ + "id": "msg-1", + "type": "message", + "role": "assistant", + "content": [{"type": "output_text", "text": "key=sk-abcdefghijklmnop at /root/secret/x"}], + }), + 0, + ) + .expect("assistant item"); + let DirectThreadItem::Message { text, .. } = item else { + panic!("message item"); + }; + assert!( + !text.contains("sk-abcdefghijklmnop"), + "不得泄漏明文密钥:{text}" + ); + assert!(!text.contains("/root/secret"), "不得泄漏绝对路径:{text}"); + } + + #[test] + fn history_keeps_call_and_output_as_two_items_with_one_identity() { + let items = vec![ + json!({ + "id": "05dc0af1-8023-47fd-ad22-d54df2837b1b", + "call_id": "call_00_Gpd0s0Ytm9YgIbwbEXva1473", + "type": "function_call", + "name": "exec_command", + "arguments": "{\"cmd\": \"ls\"}", + }), + json!({ + "id": "fco_01a06fa5-d636-7452-b337-a641c2e6bc76", + "call_id": "call_00_Gpd0s0Ytm9YgIbwbEXva1473", + "type": "function_call_output", + "output": "assets\ngame\n", + }), + ]; + let projected = direct_thread_items_from_history(root(), &items, |_| 0); + assert_eq!(projected.len(), 2, "搬运层不得替前端做合并"); + assert!(matches!( + projected[0], + DirectThreadItem::FunctionCall { .. } + )); + assert!(matches!( + projected[1], + DirectThreadItem::FunctionCallOutput { .. } + )); + // 调用与输出共享同一个归一身份,前端才能把它们并成一张卡片。 + assert_eq!(projected[0].item_id(), projected[1].item_id()); + assert_eq!(projected[0].item_id(), "call_00_Gpd0s0Ytm9YgIbwbEXva1473"); + } + + #[test] + fn unknown_item_types_are_passed_through_without_body() { + let item = direct_thread_item_from_value( + root(), + &json!({"id": "plan-1", "type": "plan", "text": "内部计划"}), + 0, + ) + .expect("unknown item"); + assert!(matches!( + item, + // TODO(direct-thread): 未知类型目前只带类型与身份,前端投影会丢弃它。 + // 哪些类型要显示属于前端可见性决策,需要时改前端,不要在这里加白名单。 + DirectThreadItem::Other { ref raw_type, .. } if raw_type == "plan" + )); + } +} diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_tool_calls.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_tool_calls.rs index 545b1fa11..f73471a56 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_tool_calls.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_tool_calls.rs @@ -8,7 +8,7 @@ //! 为什么不复用 `project.jsonl`:那条链路的回读只投影 `role ∈ {user, assistant}` 的 //! 文本条目,而且会被注入 Codex 上下文。往里面塞新形状既装不下,又有污染模型上下文的风险。 -use super::direct_thread_raw_item::sanitize_detail_text; +use super::direct_thread_wire::sanitize_detail_text; use crate::config::{prepare_game_creator_private_path_for_read, write_game_creator_private_file}; use crate::project::{enforce_project_permission_policy, project_append_lock_for}; use serde::{Deserialize, Serialize}; diff --git a/apps/ai-game-creator-shell/src-tauri/src/commands.rs b/apps/ai-game-creator-shell/src-tauri/src/commands.rs index 8dbdc185b..696aff810 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/commands.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/commands.rs @@ -5388,7 +5388,7 @@ pub(crate) async fn read_direct_project_history_slice( .and_then(|item| item.get("id")) .and_then(serde_json::Value::as_str) .map(str::to_string); - let items = direct_thread_raw_items_from_history(root, &items, |item| { + let items = direct_thread_items_from_history(root, &items, |item| { item.get("id") .and_then(serde_json::Value::as_str) .and_then(|id| item_timestamps.get(id).copied()) diff --git a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserContentPart.ts b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserContentPart.ts index 10c4e0f5a..eec5662e2 100644 --- a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserContentPart.ts +++ b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserContentPart.ts @@ -1,5 +1,4 @@ -// This file is generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. - +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. import type { DirectCodexUserRuntimeRegionPart } from './DirectCodexUserRuntimeRegionPart'; export type DirectCodexUserContentPart = diff --git a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserItem.ts b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserItem.ts index 4c9ea8ca3..aa1ee88dd 100644 --- a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserItem.ts +++ b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserItem.ts @@ -1,5 +1,9 @@ -// This file is generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. - +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. import type { DirectCodexUserMessageItem } from './DirectCodexUserMessageItem'; -export type DirectCodexUserItem = DirectCodexUserMessageItem; +/** + * DirectProject 本轮 user input 的唯一结构化入口。 + */ +export type DirectCodexUserItem = { + type: 'message'; +} & DirectCodexUserMessageItem; diff --git a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserMessageEnvelope.ts b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserMessageEnvelope.ts new file mode 100644 index 000000000..97e0b0dae --- /dev/null +++ b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserMessageEnvelope.ts @@ -0,0 +1,4 @@ +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. +import type { DirectCodexUserItem } from './DirectCodexUserItem'; + +export type DirectCodexUserMessageEnvelope = { item: DirectCodexUserItem }; diff --git a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserMessageItem.ts b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserMessageItem.ts index f7447dc5d..a1101c8cb 100644 --- a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserMessageItem.ts +++ b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserMessageItem.ts @@ -1,11 +1,9 @@ -// This file is generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. - +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. import type { DirectCodexUserContentPart } from './DirectCodexUserContentPart'; import type { DirectCodexUserRole } from './DirectCodexUserRole'; export type DirectCodexUserMessageItem = { - type: 'message'; role: DirectCodexUserRole; - content: DirectCodexUserContentPart[]; + content: Array; id: string; }; diff --git a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserRole.ts b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserRole.ts index e1dc59c59..0c2540122 100644 --- a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserRole.ts +++ b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserRole.ts @@ -1,3 +1,3 @@ -// This file is generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. export type DirectCodexUserRole = 'user'; diff --git a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserRuntimeRegionPart.ts b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserRuntimeRegionPart.ts index 296c0d19a..b134fb7bc 100644 --- a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserRuntimeRegionPart.ts +++ b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectCodexUserRuntimeRegionPart.ts @@ -1,13 +1,13 @@ -// This file is generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. export type DirectCodexUserRuntimeRegionPart = { label: string; - runId?: string; - versionId?: string; - elementTag?: string; - elementRole?: string; - text?: string; - width?: number; - height?: number; - resourceIds: string[]; + runId: string | null; + versionId: string | null; + elementTag: string | null; + elementRole: string | null; + text: string | null; + width: number | null; + height: number | null; + resourceIds: Array; }; diff --git a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadConsumeResult.ts b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadConsumeResult.ts new file mode 100644 index 000000000..ea7e6e112 --- /dev/null +++ b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadConsumeResult.ts @@ -0,0 +1,4 @@ +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. +import type { DirectThreadEvent } from './DirectThreadEvent'; + +export type DirectThreadConsumeResult = { events: Array }; diff --git a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadDeltaKind.ts b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadDeltaKind.ts new file mode 100644 index 000000000..215819e1a --- /dev/null +++ b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadDeltaKind.ts @@ -0,0 +1,6 @@ +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. + +/** + * 增量正文属于哪类条目。 + */ +export type DirectThreadDeltaKind = 'message' | 'reasoning'; diff --git a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadEvent.ts b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadEvent.ts new file mode 100644 index 000000000..a4b0d7688 --- /dev/null +++ b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadEvent.ts @@ -0,0 +1,30 @@ +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. +import type { DirectThreadDeltaKind } from './DirectThreadDeltaKind'; +import type { DirectThreadItem } from './DirectThreadItem'; +import type { DirectThreadRequestKind } from './DirectThreadRequestKind'; + +/** + * Thread Manager 下发的运行态事件。 + * + * 顺序由数组顺序给出(同一个 subscriber 的 `consume` 按队列顺序返回),因此不需要 `seq`: + * 游标是 Thread Manager 的内部事实,不下发。 + * + * 事件不带回合身份:DirectProject 同一时刻只有一个回合在跑,"当前回合是否还在跑"由 + * 生命周期事件在序列中的位置给出,`turn_id` 对前端没有任何额外信息。 + */ +export type DirectThreadEvent = + | { type: 'turn.started' } + | { type: 'turn.completed'; status: string } + | { type: 'item.started'; item: DirectThreadItem } + | { type: 'item.completed'; item: DirectThreadItem } + | { + type: 'item.delta'; + itemId: string; + kind: DirectThreadDeltaKind; + delta: string; + } + | { + type: 'request'; + kind: DirectThreadRequestKind; + requestId: string | null; + }; diff --git a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadFileChange.ts b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadFileChange.ts new file mode 100644 index 000000000..a510ccb68 --- /dev/null +++ b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadFileChange.ts @@ -0,0 +1,12 @@ +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. + +/** + * 一条文件变更。 + */ +export type DirectThreadFileChange = { + path: string; + /** + * `add` | `update` | `delete` + */ + kind: string; +}; diff --git a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadHistorySlice.ts b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadHistorySlice.ts new file mode 100644 index 000000000..9ac30f139 --- /dev/null +++ b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadHistorySlice.ts @@ -0,0 +1,14 @@ +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. +import type { DirectThreadItem } from './DirectThreadItem'; + +export type DirectThreadHistorySlice = { + /** + * 脱敏条目,顺序即文件顺序;与运行态事件里的条目同形。 + */ + items: Array; + hasMore: boolean; + /** + * 本次切片的原始 item id 锚点:无论切片里有没有可显示条目,分页都靠它向前。 + */ + firstItemId: string | null; +}; diff --git a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadItem.ts b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadItem.ts new file mode 100644 index 000000000..f1257ddcc --- /dev/null +++ b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadItem.ts @@ -0,0 +1,76 @@ +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. +import type { DirectThreadFileChange } from './DirectThreadFileChange'; + +/** + * 聊天视图的输入条目:一条 Codex 原始条目的脱敏投影。 + * + * `itemType` 就是 Codex 的原始类型,逐字透传;前端按它决定投影成消息、思考还是工具卡片。 + * 未识别的类型走 [`DirectThreadItem::Other`],Rust 不替前端决定它是否可见。 + * + * 条目上的 `at` 是只用于显示的毫秒时间戳:ts-rs 默认把 `u64` 映射成 `bigint`, + * 而 Tauri 的 JSON 通道传过来的是 `number`,因此统一标 `#[ts(as = "f64")]` 对齐。 + */ +export type DirectThreadItem = + | { + itemType: 'message'; + /** + * 归一身份:全链路只有这一个 id。 + */ + itemId: string; + /** + * 原始 role(`user` / `assistant` / `system` / …);显示与否由前端判断。 + */ + role: string; + text: string; + at: number; + } + | { itemType: 'reasoning'; itemId: string; text: string; at: number } + | { + itemType: 'function_call'; + itemId: string; + name: string; + arguments: string; + at: number; + } + | { + itemType: 'function_call_output'; + itemId: string; + output: string; + at: number; + } + | { + itemType: 'commandExecution'; + itemId: string; + command: string; + output: string | null; + /** + * app-server 原始状态:`inProgress` / `completed` / `failed` / `declined` / … + */ + status: string | null; + exitCode: number | null; + at: number; + } + | { + itemType: 'fileChange'; + itemId: string; + changes: Array; + at: number; + } + | { + itemType: 'mcpToolCall'; + itemId: string; + tool: string; + arguments: string; + output: string | null; + status: string | null; + at: number; + } + | { + itemType: 'webSearch'; + itemId: string; + query: string | null; + output: string | null; + at: number; + } + | { itemType: 'contextCompaction'; itemId: string; at: number } + | { itemType: 'other'; itemId: string; rawType: string; at: number }; diff --git a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadRequestKind.ts b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadRequestKind.ts new file mode 100644 index 000000000..13ee401f3 --- /dev/null +++ b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadRequestKind.ts @@ -0,0 +1,9 @@ +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. + +/** + * 审批 / 提问请求与解决:本轮只透传,不并入聊天状态。 + */ +export type DirectThreadRequestKind = + | 'approval.requested' + | 'ask.requested' + | 'request.resolved'; diff --git a/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadSubscriptionBootstrap.ts b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadSubscriptionBootstrap.ts new file mode 100644 index 000000000..45efbf4b4 --- /dev/null +++ b/apps/ai-game-creator-shell/src/features/project-workspace/generated/DirectThreadSubscriptionBootstrap.ts @@ -0,0 +1,14 @@ +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. +import type { DirectThreadEvent } from './DirectThreadEvent'; + +export type DirectThreadSubscriptionBootstrap = { + subscriptionId: string; + /** + * 首屏历史锚点:`project.jsonl` 里最后一条原始 item id。 + */ + lastCompletedItemId: string | null; + /** + * 该 subscriber 此刻应当处理的运行态事件(游标已经在队尾)。 + */ + events: Array; +};