diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/user_input.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/user_input.rs index 3cab9e326..7b7cc68d4 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/user_input.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/user_input.rs @@ -100,18 +100,9 @@ async fn enqueue_direct_codex_turn_typed( let detail = redact_agent_runtime_error(root, &error, 1800); DirectTurnError::EnvironmentNotReady { detail } })?; - // 入队:到这里这一条已经过了全部检查,剩下的就是排队等放行。prompt 与 canonical 形状在这里 - // 冻结——放行不重算,所以放行没有失败出口。 - let pending = PendingTurn::prepare( - turn_id, - user_item, - user_prompt, - creation_type, - direct_tool_call_now_ms(), - ) - .map_err(|error| DirectTurnError::InputRejected { - detail: error.to_string(), - })?; + // 入队:到这里这一条已经过了全部检查,剩下的就是排队等放行。条目只带走它自己的事实 + // (用户条目、创建类型、入队时刻),canonical 形状与 prompt 放行时从它重投影——放行没有失败出口。 + let pending = PendingTurn::new(turn_id, user_item, creation_type, direct_tool_call_now_ms()); match enqueue_pending_turn(&thread_id, pending) { Ok(_) => {} Err(EnqueueRejection::QueueFull) => { diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/dispatch.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/dispatch.rs index 7035391b1..22e016807 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/dispatch.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/dispatch.rs @@ -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_tool_call_now_ms, redact_agent_runtime_error, + direct_codex_user_item_to_prompt, direct_tool_call_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, }; @@ -62,7 +62,7 @@ impl TurnReservation { /// 测试用:直接开一轮占用,等价于"入队 + 认领",不经过 Tauri 命令与真实的用户条目。 #[cfg(test)] pub(crate) fn accept_for_test(thread_id: &str, client_turn_id: &str) -> Self { - let pending = PendingTurn::prepare( + let pending = PendingTurn::new( client_turn_id.to_string(), serde_json::from_value(serde_json::json!({ "type": "message", @@ -71,11 +71,9 @@ impl TurnReservation { "id": format!("direct-codex:{client_turn_id}:user"), })) .expect("canonical user item"), - "测试消息".to_string(), None, direct_tool_call_now_ms(), - ) - .expect("prepare pending turn"); + ); super::enqueue_pending_turn(thread_id, pending).expect("enqueue test turn"); let dispatched = claim_pending_turn(thread_id).expect("claim test turn"); Self::resume(thread_id, &dispatched) @@ -128,10 +126,12 @@ async fn run_dispatched_direct_turn( release_identity_generation: u64, ) { let turn_id = dispatched.pending.client_turn_id.clone(); + // 放行只搬运条目自己的事实:canonical 形状与 prompt 都从冻结过的条目重投影,不重算、不写盘。 + let canonical_user_item = serde_json::to_value(&dispatched.pending.user_item) + .expect("冻结过的 canonical 条目一定可序列化"); + let prompt = direct_codex_user_item_to_prompt(&dispatched.pending.user_item); // 落盘即回合成立:放行之后必须留下这条用户消息,哪怕这一轮随后失败。 - if let Err(error) = - append_direct_project_user_message_at(&root, &dispatched.pending.canonical_user_item) - { + if let Err(error) = append_direct_project_user_message_at(&root, &canonical_user_item) { // 不继续起整轮:历史是这条对话的单一事实源,用户消息没落盘时继续跑只会得到一条没有开口用户 // 消息的助手回复,而且失败会被静默掉。 let failure = DirectTurnError::EnvironmentNotReady { @@ -146,18 +146,15 @@ async fn run_dispatched_direct_turn( } // 用户条目落盘成功即下发:这一轮从"放行"到"起 codex"之间的一切失败(连不上 app-server、执行器 // 未通过验收、历史注入失败)都靠它把失败说明挂回自己那一轮。 - crate::agent::codex_app_server::emit_direct_thread_user_item( - &root, - &dispatched.pending.canonical_user_item, - ); + crate::agent::codex_app_server::emit_direct_thread_user_item(&root, &canonical_user_item); let capture = crate::analytics::gui::capture_writer_context(); let emitter = DirectGameCreatorTurnUpdateEmitter::new(&root, turn_id); let outcome = run_direct_game_creator_turn_at_with_creation_type_and_emitter( &root, - &dispatched.pending.prompt, + &prompt, dispatched.pending.creation_type.as_deref(), Some(&emitter), - Some(dispatched.pending.canonical_user_item.clone()), + Some(canonical_user_item.clone()), capture, release_identity_generation, ) @@ -202,9 +199,9 @@ mod tests { .collect() } - /// 一条待发消息。放行路径只读它的身份与产物,所以这里用最小可用形状。 + /// 一条待发消息。放行路径只读它自己的事实(身份、条目、创建类型),所以这里用最小可用形状。 fn pending_turn(client_turn_id: &str) -> PendingTurn { - PendingTurn::prepare( + PendingTurn::new( client_turn_id.to_string(), serde_json::from_value(serde_json::json!({ "type": "message", @@ -213,11 +210,9 @@ mod tests { "id": format!("direct-codex:{client_turn_id}:user"), })) .expect("canonical user item"), - format!("消息 {client_turn_id}"), None, direct_tool_call_now_ms(), ) - .expect("prepare pending turn") } /// 真实存在的项目目录不是这条用例的判据,用唯一的假路径当线程身份即可:认领与占用的语义只认 diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/mod.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/mod.rs index 7e9de8d2f..0055d42e5 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/mod.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/mod.rs @@ -38,12 +38,6 @@ struct StoredEvent { seq: u64, bytes: usize, cleanable: bool, - /// 这条事件在队期间挂着的宿主侧产物。 - /// - /// 只有 `queue.enqueued` 会带它,而且**它本身就是队列成员身份**:`Some` = 这条待发消息还在队, - /// 被取走 = 已经离开队列。队列的先后就是事件的先后,所以没有第二张队列表要对齐,也不会出现 - /// "事件说在队、队列表说不在"的中间态(见 `thread_manager::queue`)。 - pending: Option, } #[derive(Clone, Debug)] @@ -194,16 +188,6 @@ impl ThreadManager { } pub(crate) fn append(&mut self, thread_id: &str, event: ThreadEvent) -> ThreadEvent { - self.append_inner(thread_id, event, None) - } - - /// 追加事件。`pending` 用来把宿主侧产物挂到这条事件上(只有入队会传)。 - fn append_inner( - &mut self, - thread_id: &str, - event: ThreadEvent, - pending: Option, - ) -> ThreadEvent { let thread = self.threads.entry(thread_id.to_string()).or_default(); thread.next_seq = thread.next_seq.saturating_add(1); let seq = thread.next_seq; @@ -217,7 +201,6 @@ impl ThreadManager { seq, bytes, cleanable, - pending, }); if let Some(item_id) = event.item_id() { Self::mark_item_events_cleanable(thread, item_id, seq); @@ -418,13 +401,26 @@ impl ThreadManager { } /// 在队的待发消息,队首在前。 - fn pending_turns(thread: &ThreadState) -> Vec<&PendingTurn> { - thread - .events - .iter() - .skip(thread.head) - .filter_map(|stored| stored.pending.as_ref()) - .collect() + /// + /// 队列就是事件列表:一条待发消息在队 = 事件窗口里"有 `queue.enqueued`、没有配对 + /// `queue.removed`"。一次折叠既判成员也定顺序,不需要第二张表;在队条目的入队事件永不可回收, + /// 所以它不会被前缀回收吃掉(见 `thread_manager::queue` 与 [`Self::observe_event`])。 + fn pending_turns(thread: &ThreadState) -> Vec { + let mut pending: Vec = Vec::new(); + for stored in thread.events.iter().skip(thread.head) { + match &stored.event { + ThreadEvent::QueueEnqueued { .. } => { + if let Some(turn) = PendingTurn::from_event(&stored.event) { + pending.push(turn); + } + } + ThreadEvent::QueueRemoved { client_turn_id, .. } => { + pending.retain(|turn| &turn.client_turn_id != client_turn_id); + } + _ => {} + } + } + pending } /// 入队:把一条已经通过入队检查的待发消息排到队尾。 @@ -436,6 +432,7 @@ impl ThreadManager { thread_id: &str, turn: PendingTurn, ) -> Result { + let event = turn.enqueued_event(); { let thread = self.threads.entry(thread_id.to_string()).or_default(); let already_known = Self::pending_turns(thread) @@ -450,8 +447,7 @@ impl ThreadManager { } queue_has_room(Self::pending_turns(thread).len())?; } - let event = turn.enqueued_event(); - self.append_inner(thread_id, event, Some(turn)); + self.append(thread_id, event); Ok(EnqueueOutcome::Enqueued) } @@ -464,35 +460,24 @@ impl ThreadManager { client_turn_id: &str, at_ms: u64, ) -> QueueRemovalOutcome { - let removed = { - let Some(thread) = self.threads.get_mut(thread_id) else { - return QueueRemovalOutcome::NotFound; - }; - if thread - .active_turn - .as_ref() - .is_some_and(|active| active.turn_id == client_turn_id) - { - return QueueRemovalOutcome::AlreadyDispatched; - } - let head = thread.head; - let mut removed = false; - for stored in thread.events.iter_mut().skip(head) { - if stored - .pending - .as_ref() - .is_some_and(|pending| pending.client_turn_id == client_turn_id) - { - // 取走产物即离开队列:这一幕在临界区里发生,外面看不到"事件还在、条目已经不在"。 - stored.pending = None; - removed = true; - } - } - removed + let Some(thread) = self.threads.get(thread_id) else { + return QueueRemovalOutcome::NotFound; }; - if !removed { + if thread + .active_turn + .as_ref() + .is_some_and(|active| active.turn_id == client_turn_id) + { + return QueueRemovalOutcome::AlreadyDispatched; + } + if !Self::pending_turns(thread) + .iter() + .any(|pending| pending.client_turn_id == client_turn_id) + { return QueueRemovalOutcome::NotFound; } + // 离开队列这件事由 `queue.removed` 自己表达:追加即出队(同一条入队事件就此转成可回收), + // 所以入队成功之后这里一定找得到那一条。 self.append( thread_id, ThreadEvent::queue_removed( @@ -515,17 +500,11 @@ impl ThreadManager { started_at_ms: u64, ) -> Option { let pending = { - let thread = self.threads.get_mut(thread_id)?; + let thread = self.threads.get(thread_id)?; if thread.active_turn.is_some() { return None; } - let head = thread.head; - let stored = thread - .events - .iter_mut() - .skip(head) - .find(|stored| stored.pending.is_some())?; - stored.pending.take().expect("filtered on a pending turn") + Self::pending_turns(thread).into_iter().next()? }; let token = Uuid::new_v4().to_string(); let turn_id = pending.client_turn_id.clone(); @@ -745,10 +724,8 @@ impl ThreadManager { /// 一条待发消息离开队列:它和它对应的入队事件一起变成可回收。 fn mark_queue_events_cleanable(thread: &mut ThreadState, client_turn_id: &str) { for stored in &mut thread.events { - // 仍挂着产物的条目永远不算可回收:在队条目被回收就等于队列悄悄丢了一条消息。 - if stored.pending.is_none() - && stored.event.queue_client_turn_id() == Some(client_turn_id) - { + // 配对的两条一起转可回收——在队条目的入队事件本来就不会走到这里。 + if stored.event.queue_client_turn_id() == Some(client_turn_id) { stored.cleanable = true; } } @@ -1064,7 +1041,10 @@ pub(crate) fn unsubscribe_thread(subscription_id: &str) { #[cfg(test)] mod tests { use super::*; - use crate::agent::{ThreadDeltaKind, ThreadItem, ThreadRequestKind, MAX_PENDING_TURNS}; + use crate::agent::{ + direct_codex_user_item_to_prompt, ThreadDeltaKind, ThreadItem, ThreadRequestKind, + MAX_PENDING_TURNS, + }; fn message(item_id: &str) -> ThreadItem { ThreadItem::Message { @@ -1471,7 +1451,7 @@ mod tests { } fn pending_turn(client_turn_id: &str) -> PendingTurn { - PendingTurn::prepare( + PendingTurn::new( client_turn_id.to_string(), serde_json::from_value(serde_json::json!({ "type": "message", @@ -1480,11 +1460,9 @@ mod tests { "id": format!("direct-codex:{client_turn_id}:user"), })) .expect("canonical user item"), - format!("消息 {client_turn_id}"), None, FIXED_AT_MS, ) - .expect("prepare pending turn") } fn enqueued_ids(events: &[ThreadEvent]) -> Vec<&str> { @@ -1680,7 +1658,10 @@ mod tests { .claim_pending_turn("thread-1", 2_000) .expect("claim head"); assert_eq!(claimed.pending.client_turn_id, "turn-1"); - assert_eq!(claimed.pending.prompt, "消息 turn-1"); + assert_eq!( + direct_codex_user_item_to_prompt(&claimed.pending.user_item), + "消息 turn-1" + ); assert_eq!(manager.pending_turn_ids("thread-1"), vec!["turn-2"]); let events = manager diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/queue.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/queue.rs index 6ea36e168..12fc6879c 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/queue.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/queue.rs @@ -4,11 +4,10 @@ //! `queue.enqueued` 事件不可回收;离开队列(取消或放行)时才转成可回收。新订阅者的 bootstrap 因此 //! 天然看得见当前队列,不需要第二张队列表,也不会有"事件与队列不一致"的窗口。 //! -//! 这个模块只放三件事:一条待发消息带走什么([`PendingTurn`])、容量规则 -//! ([`MAX_PENDING_TURNS`])、以及它在线上长什么样([`PendingTurn::enqueued_event`])。 -//! 它不碰锁、不碰 Tauri、不写盘:入队检查在命令侧,放行顺序在 Thread Manager。 - -use serde_json::Value; +//! 这个模块只放三件事:一条待发消息在线上长什么样([`PendingTurn`] 与 +//! [`PendingTurn::enqueued_event`])、容量规则([`MAX_PENDING_TURNS`])、以及"它还在不在队" +//! 这条判据的成员折叠(`ThreadManager::pending_turns`)。 +//! 它不碰锁、不碰 Tauri、不写盘、不重算 prompt:入队检查在命令侧,放行顺序在 Thread Manager。 use crate::agent::{ direct_codex_user_item_id_for_client_turn_id, DirectCodexUserItem, ThreadEvent, @@ -22,48 +21,57 @@ pub(crate) const MAX_PENDING_TURNS: usize = 5; /// 一条已经通过入队检查、正在等放行的用户消息。 /// -/// 只在内存里,进程重启即消失(与 ADR 记的边界一致)。它同时是**放行时要用的全部输入**:放行 -/// 没有失败出口,所以检查产物在入队时就地冻结,放行只搬运、不重算。 +/// 它**就是 `queue.enqueued` 的投影**([`PendingTurn::from_event`] / [`PendingTurn::enqueued_event`]): +/// 字段与事件载荷一一对应,宿主不为它另存第二份产物——canonical 形状与 prompt 都在放行时从这条条目 +/// 重投影。在队 = 事件窗口里"有 `queue.enqueued`、没有配对 `queue.removed`",成员与顺序因此只从事件 +/// 折出来,没有"事件说在队、队列表说不在"的中间态。 #[derive(Clone, Debug)] pub(crate) struct PendingTurn { /// 这条消息的回合身份;放行后同一轮的 `turn.started` / `turn.completed` 用它。 pub(crate) client_turn_id: String, /// canonical 用户条目:事件与界面 chip 都读它,Rust 不渲染展示形状。 pub(crate) user_item: DirectCodexUserItem, - /// canonical 用户条目的 JSON 形状:入队时算好,放行时直接落盘。 - pub(crate) canonical_user_item: Value, - /// 入队检查产出的 prompt:放行不重算。 - pub(crate) prompt: String, pub(crate) creation_type: Option, /// 入队那一刻的宿主毫秒钟。 pub(crate) at: u64, } impl PendingTurn { - /// 组一条待发消息:入队检查已经全部通过,这里只把放行要用的产物冻结下来。 - /// - /// 冻结是刻意的:放行没有失败出口,所以任何可能在放行时才失败的计算都必须提前到这里 - /// (canonical 形状与 prompt 都是)。 - pub(crate) fn prepare( + pub(crate) fn new( client_turn_id: String, user_item: DirectCodexUserItem, - prompt: String, creation_type: Option, at: u64, - ) -> Result { - let canonical_user_item = serde_json::to_value(&user_item)?; - Ok(Self { + ) -> Self { + Self { client_turn_id, user_item, - canonical_user_item, - prompt, creation_type, at, + } + } + + /// 从 `queue.enqueued` 事件折出这条待发消息;其它事件返回 `None`。 + pub(crate) fn from_event(event: &ThreadEvent) -> Option { + let ThreadEvent::QueueEnqueued { + client_turn_id, + user_item, + creation_type, + at, + } = event + else { + return None; + }; + Some(Self { + client_turn_id: client_turn_id.clone(), + user_item: user_item.clone(), + creation_type: creation_type.clone(), + at: *at, }) } - /// 入队事件的投影:带 canonical 用户条目与 `creationType`,**不带 prompt**(prompt 只留在宿主的 - /// 队列条目里,它不是要下发的展示形状)。 + /// 入队事件的投影:带 canonical 用户条目与 `creationType`,**不带 prompt**(prompt 是入队检查的 + /// 产物,放行时由条目重投影,不是要下发的展示形状),也不带埋点身份(成绩由宿主自己结算)。 pub(crate) fn enqueued_event(&self) -> ThreadEvent { ThreadEvent::queue_enqueued( self.client_turn_id.clone(), @@ -119,14 +127,12 @@ mod tests { } fn pending(client_turn_id: &str, creation_type: Option<&str>) -> PendingTurn { - PendingTurn::prepare( + PendingTurn::new( client_turn_id.to_string(), user_item("生成一个游戏", "direct-codex:turn-1:user"), - "生成一个游戏".to_string(), creation_type.map(str::to_string), 1_700_000_000_123, ) - .expect("prepare pending turn") } /// 入队事件只带 canonical 用户条目与 `creationType`:prompt 是宿主的入队检查产物, @@ -162,18 +168,21 @@ mod tests { assert_eq!(event.queue_removal_reason(), None); } - /// canonical 形状在入队时就冻结:放行时落盘的就是这一份,不再重算。 + /// 条目就是事件投影:由自己的入队事件折回来的条目与原件逐字段相同。 #[test] - fn prepare_freezes_the_canonical_item() { - let turn = pending("turn-1", None); - assert_eq!( - turn.canonical_user_item["id"], - json!("direct-codex:turn-1:user") - ); - assert_eq!( - turn.canonical_user_item, - serde_json::to_value(&turn.user_item).expect("serialize user item") - ); + fn pending_turn_round_trips_through_its_enqueued_event() { + for turn in [pending("turn-1", Some("web-game")), pending("turn-2", None)] { + let restored = + PendingTurn::from_event(&turn.enqueued_event()).expect("queue.enqueued projects"); + assert_eq!(restored.client_turn_id, turn.client_turn_id); + assert_eq!(restored.creation_type, turn.creation_type); + assert_eq!(restored.at, turn.at); + assert_eq!( + serde_json::to_value(&restored.user_item).expect("serialize user item"), + serde_json::to_value(&turn.user_item).expect("serialize user item") + ); + assert!(PendingTurn::from_event(&ThreadEvent::turn_started(1)).is_none()); + } } /// 用户条目身份由 `clientTurnId` 派生,与落盘 / 下发用的是同一个函数。 diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/wire.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/wire.rs index 16ba49cfe..f5d328482 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/wire.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/thread_manager/wire.rs @@ -354,7 +354,8 @@ pub(crate) enum ThreadEvent { /// 这条事件在条目仍在队期间**不可回收**,离开队列(取消或放行)时才转成可回收——新订阅者 /// 靠这一点在 bootstrap 里看到当前队列,`is_bootstrap_event` 不需要为它加特例。 /// - /// 事件不带 prompt:prompt 是入队检查的产物,只留在宿主的队列条目里。 + /// 事件就是这条待发消息的**全部**事实:宿主不为它另存产物,prompt 与 canonical 形状都在放行时 + /// 从这条条目重投影。 #[serde(rename = "queue.enqueued")] QueueEnqueued { /// 这条待发消息的回合身份;放行后同一轮的 `turn.started` / `turn.completed` 用它。