diff --git a/CONTEXT.md b/CONTEXT.md index c22344ea0..00f7c178e 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -192,7 +192,7 @@ _Avoid_: mock 先行堆积、前后端各自发散、先做排行榜 UI ## 项目开发对话(DirectProject) **DirectProject 专属聊天模块**: -AGC 普通项目聊天的独立容器,拥有 DirectProject 的聊天状态、运行态订阅、历史读取、发送队列、附件和中止交互,并把聊天投影交给专属表现层渲染;它不承接 Supervisor、Design Agent 或 Planning V2 的运行态。 +AGC 普通项目聊天的独立容器,拥有 DirectProject 的聊天状态、运行态订阅、历史读取、待发消息队列的投影、附件和中止交互,并把聊天投影交给专属表现层渲染;它不承接 Supervisor、Design Agent 或 Planning V2 的运行态。 _Avoid_: 把 DirectProject 作为项目总控聊天的一个布尔分支、把四种 Agent 会话抽象成同一事实源 **项目工作台布局**: @@ -208,20 +208,32 @@ Thread Manager 向订阅者推送的当前回合原始事件流,只服务运 _Avoid_: 进度通知、快照轮询、第二套历史 **逻辑回合**: -Thread Manager 拥有的一对回合边界(开始与结束),由接单动作开启、由这一轮的占用对象写出,不镜像 Codex 原生回合;界面忙碌态与回合结果只认它。 +Thread Manager 拥有的一对回合边界(开始与结束),由放行动作开启、由这一轮的占用对象写出,不镜像 Codex 原生回合;界面忙碌态与回合结果只认它。 _Avoid_: Codex 原生回合、原生日志、进程生命周期 -**接单**: -把一条用户消息交给宿主开始执行的动作,成立即表示这一轮已经存在;此后结果只由运行态事件回答。 -_Avoid_: 发送成功、命令调用、接口返回 +**待发消息队列**: +Thread Manager 按项目持有的待发用户消息序列,只支持按入队顺序追加与按身份移除,状态由运行态事件派生,不落盘、不构成第二份事实源。 +_Avoid_: 前端本地队列、队列副本、待发消息的持久化记录 -**拒单**: -接单成立之前拒绝这次请求(并发、权限、目录、参数、工程准备未就绪),只回一条可展示原因,不产生回合事件,也不写用户条目。 +**待发消息**: +已经通过入队检查、等待被放行的用户消息;它在放行之前不是回合,不写用户条目、不产生回合事件。 +_Avoid_: 回合、在途回合、草稿 + +**入队**: +把一条用户消息交给宿主的动作:宿主跑完入队检查后把它放进待发消息队列;入队成立只表示这条消息会按顺序被放行。 +_Avoid_: 发送成功、已经开跑、回合成立 + +**入队失败**: +入队检查未通过(身份、形状、容量、权限、目录、参数、工程准备未就绪)时拒绝这次请求,只回一条可展示原因,不入队、不产生回合事件,也不写用户条目。 _Avoid_: 回合失败、执行失败、失败事件 +**放行**: +Thread Manager 在一个回合收口之后把队首的待发消息送进回合:同一临界区里登记占用、落盘用户条目、发出逻辑回合开始事件并起整轮;放行之后的结果只由运行态事件回答。 +_Avoid_: 前端放行、定时轮询、放行失败 + **在途回合**: -界面本地已经把这条用户消息发出去、宿主还没有对应回合开始事件的那一小段状态。 -_Avoid_: 运行中回合、乐观锁、发送队列 +界面本地已经入队、宿主还没有对应回合开始事件的那一小段状态。 +_Avoid_: 运行中回合、乐观锁、前端发送队列 **聊天投影**: 把项目对话历史条目与运行态事件转换成消息气泡和工具卡片的读取期转换;不持久化,也不构成事实源。 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 347ec3139..781aad21d 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent.rs @@ -29,6 +29,7 @@ mod direct_project_history; mod direct_project_turn_history; mod direct_runtime; mod direct_thread_manager; +mod direct_thread_queue; mod direct_thread_wire; mod direct_tool_bridge; mod direct_tool_calls; @@ -65,6 +66,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_queue::*; pub(crate) use direct_thread_wire::*; pub(crate) use direct_tool_bridge::*; pub(crate) use direct_tool_calls::*; diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_codex_user_item/model.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_codex_user_item/model.rs index b65f25ba4..2b537d5f6 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_codex_user_item/model.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_codex_user_item/model.rs @@ -2,7 +2,7 @@ use serde::{Deserialize, Serialize}; use ts_rs::TS; /// DirectProject 本轮 user input 的唯一结构化入口。 -#[derive(Clone, Debug, Deserialize, Serialize, TS)] +#[derive(Clone, Debug, PartialEq, Deserialize, Serialize, TS)] #[serde(tag = "type", deny_unknown_fields)] #[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))] pub(crate) enum DirectCodexUserItem { @@ -10,7 +10,7 @@ pub(crate) enum DirectCodexUserItem { Message(DirectCodexUserMessageItem), } -#[derive(Clone, Debug, Deserialize, Serialize, TS)] +#[derive(Clone, Debug, PartialEq, Deserialize, Serialize, TS)] #[serde(rename_all = "camelCase", deny_unknown_fields)] #[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))] pub(crate) struct DirectCodexUserMessageItem { @@ -19,14 +19,14 @@ pub(crate) struct DirectCodexUserMessageItem { pub(crate) id: String, } -#[derive(Clone, Debug, Deserialize, Serialize, TS)] +#[derive(Clone, Debug, PartialEq, Deserialize, Serialize, TS)] #[serde(rename_all = "lowercase")] #[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))] pub(crate) enum DirectCodexUserRole { User, } -#[derive(Clone, Debug, Deserialize, Serialize, TS)] +#[derive(Clone, Debug, PartialEq, Deserialize, Serialize, TS)] #[serde(tag = "type", rename_all_fields = "camelCase", deny_unknown_fields)] #[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))] pub(crate) enum DirectCodexUserContentPart { @@ -43,7 +43,7 @@ pub(crate) enum DirectCodexUserContentPart { AgcAttachmentReference(DirectCodexUserAttachmentReferencePart), } -#[derive(Clone, Debug, Deserialize, Serialize, TS)] +#[derive(Clone, Debug, PartialEq, Deserialize, Serialize, TS)] #[serde(rename_all = "camelCase", deny_unknown_fields)] #[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))] pub(crate) struct DirectCodexUserAttachmentReferencePart { @@ -55,7 +55,7 @@ pub(crate) struct DirectCodexUserAttachmentReferencePart { pub(crate) status: String, } -#[derive(Clone, Debug, Deserialize, Serialize, TS)] +#[derive(Clone, Debug, PartialEq, Deserialize, Serialize, TS)] #[serde(rename_all = "camelCase", deny_unknown_fields)] #[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))] pub(crate) struct DirectCodexUserRuntimeRegionPart { 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 2b624df55..a83e68baa 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 @@ -9,7 +9,9 @@ use std::sync::{Mutex, OnceLock}; use uuid::Uuid; use crate::agent::{ + direct_tool_call_now_ms, queue_has_room, DirectQueueRemovalOutcome, DirectQueueRemovalReason, DirectThreadConsumeResult, DirectThreadEvent, DirectThreadSubscriptionBootstrap, + EnqueueOutcome, EnqueueRejection, PendingDirectTurn, }; const DEFAULT_MAX_EVENTS: usize = 8_192; @@ -28,6 +30,12 @@ struct StoredEvent { seq: u64, bytes: usize, cleanable: bool, + /// 这条事件在队期间挂着的宿主侧产物。 + /// + /// 只有 `queue.enqueued` 会带它,而且**它本身就是队列成员身份**:`Some` = 这条待发消息还在队, + /// 被取走 = 已经离开队列。队列的先后就是事件的先后,所以没有第二张队列表要对齐,也不会出现 + /// "事件说在队、队列表说不在"的中间态(见 `direct_thread_queue`)。 + pending: Option, } #[derive(Clone, Debug)] @@ -75,6 +83,16 @@ pub(crate) struct DirectActiveTurnSnapshot { pub(crate) sequence: u64, } +/// 一次放行的结果:队首那条待发消息,以及这次放行的占用身份。 +/// +/// `token` 是这一轮的占用身份(终态出口只认它),`pending` 是放行要用的**全部**输入——放行不再 +/// 重跑任何检查,也就不需要再回到调用方取参数。 +#[derive(Clone, Debug)] +pub(crate) struct DirectDispatchedTurn { + pub(crate) token: String, + pub(crate) pending: PendingDirectTurn, +} + #[derive(Clone, Debug)] struct ThreadState { next_seq: u64, @@ -142,6 +160,16 @@ impl DirectThreadManager { &mut self, thread_id: &str, event: DirectThreadEvent, + ) -> DirectThreadEvent { + self.append_inner(thread_id, event, None) + } + + /// 追加事件。`pending` 用来把宿主侧产物挂到这条事件上(只有入队会传)。 + fn append_inner( + &mut self, + thread_id: &str, + event: DirectThreadEvent, + pending: Option, ) -> DirectThreadEvent { let thread = self.threads.entry(thread_id.to_string()).or_default(); thread.next_seq = thread.next_seq.saturating_add(1); @@ -156,6 +184,7 @@ impl DirectThreadManager { seq, bytes, cleanable, + pending, }); if let Some(item_id) = event.item_id() { Self::mark_item_events_cleanable(thread, item_id, seq); @@ -331,6 +360,153 @@ impl DirectThreadManager { .is_some_and(|thread| thread.active_turn.is_some()) } + /// 在队的待发消息,队首在前。 + fn pending_turns(thread: &ThreadState) -> Vec<&PendingDirectTurn> { + thread + .events + .iter() + .skip(thread.head) + .filter_map(|stored| stored.pending.as_ref()) + .collect() + } + + /// 入队:把一条已经通过入队检查的待发消息排到队尾。 + /// + /// 幂等按 `clientTurnId`,判重范围是「**在队 ∪ 正在跑的那一轮**」:同一身份重复入队返回 + /// [`EnqueueOutcome::AlreadyKnown`],不排第二条、不发第二条事件。容量只数在队条目。 + fn enqueue_pending_turn( + &mut self, + thread_id: &str, + turn: PendingDirectTurn, + ) -> Result { + { + let thread = self.threads.entry(thread_id.to_string()).or_default(); + let already_known = Self::pending_turns(thread) + .iter() + .any(|pending| pending.client_turn_id == turn.client_turn_id) + || thread + .active_turn + .as_ref() + .is_some_and(|active| active.turn_id == turn.client_turn_id); + if already_known { + return Ok(EnqueueOutcome::AlreadyKnown); + } + queue_has_room(Self::pending_turns(thread).len())?; + } + let event = turn.enqueued_event(); + self.append_inner(thread_id, event, Some(turn)); + Ok(EnqueueOutcome::Enqueued) + } + + /// 取消一条待发消息。typed 结果告诉调用方它到底是被取消、已经放行,还是本来就不在队里。 + /// + /// 判「已经被放行」用的是正在跑的那一轮的回合身份:认领之后就再也按待发消息取消不了了。 + fn remove_pending_turn( + &mut self, + thread_id: &str, + client_turn_id: &str, + at_ms: u64, + ) -> DirectQueueRemovalOutcome { + let removed = { + let Some(thread) = self.threads.get_mut(thread_id) else { + return DirectQueueRemovalOutcome::NotFound; + }; + if thread + .active_turn + .as_ref() + .is_some_and(|active| active.turn_id == client_turn_id) + { + return DirectQueueRemovalOutcome::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 + }; + if !removed { + return DirectQueueRemovalOutcome::NotFound; + } + self.append( + thread_id, + DirectThreadEvent::queue_removed( + client_turn_id.to_string(), + DirectQueueRemovalReason::Cancelled, + at_ms, + ), + ); + DirectQueueRemovalOutcome::Removed + } + + /// 放行:同一个临界区里取队首 → 登记占用 → 写 `turn.started` 与 `queue.removed{dispatched}`。 + /// + /// 返回 `None` 表示现在**不该放行**:这个 thread 已经有未收口的回合,或者队列是空的。这两件事 + /// 都由这里判,调用方(kick)不需要先探一遍再决定——探两次就是两个中间态。 + /// + /// 认领认的是队首,不是调用方点名的某一条:待发消息的顺序就是事件顺序,没有第三种来源。 + fn claim_pending_turn( + &mut self, + thread_id: &str, + started_at_ms: u64, + ) -> Option { + let pending = { + let thread = self.threads.get_mut(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") + }; + let token = Uuid::new_v4().to_string(); + let turn_id = pending.client_turn_id.clone(); + let project_name = std::path::Path::new(thread_id) + .file_name() + .and_then(|name| name.to_str()) + .map(str::to_string); + if let Some(thread) = self.threads.get_mut(thread_id) { + thread.active_turn = Some(ActiveDirectTurn { + token: token.clone(), + turn_id: turn_id.clone(), + project_name, + started_at: started_at_ms, + // 与"还没有任何进度事件"的状态一致:运行时给出的第一条进度会覆盖它。 + status: "accepted".to_string(), + activity: Some("request-accepted".to_string()), + updated_at: started_at_ms, + sequence: 0, + }); + } + let user_item_id = pending.user_item_id(); + self.append( + thread_id, + DirectThreadEvent::turn_started(started_at_ms) + .with_user_item_id(user_item_id.as_deref()), + ); + self.append( + thread_id, + DirectThreadEvent::queue_removed( + turn_id, + DirectQueueRemovalReason::Dispatched, + started_at_ms, + ), + ); + Some(DirectDispatchedTurn { token, pending }) + } + fn subscriber_ids(&self, thread_id: &str) -> Vec { self.threads .get(thread_id) @@ -395,6 +571,20 @@ impl DirectThreadManager { }) } + /// 当前在队的待发消息身份,队首在前。 + #[cfg(test)] + fn pending_turn_ids(&self, thread_id: &str) -> Vec { + self.threads + .get(thread_id) + .map(|thread| { + Self::pending_turns(thread) + .iter() + .map(|pending| pending.client_turn_id.clone()) + .collect() + }) + .unwrap_or_default() + } + /// 观察一条事件,返回它本身是否可回收。 fn observe_event(thread: &mut ThreadState, seq: u64, event: &DirectThreadEvent) -> bool { match event { @@ -428,6 +618,14 @@ impl DirectThreadManager { true } DirectThreadEvent::ItemDelta { .. } => true, + // 待发消息在队期间不可回收:队列的成员与顺序**就是**这条事件本身 + // (见 `direct_thread_queue`)。bootstrap 因此天然看得见当前队列。 + DirectThreadEvent::QueueEnqueued { .. } => false, + // 离开队列(取消 / 放行)才把同一条待发消息的入队事件一起转成可回收。 + DirectThreadEvent::QueueRemoved { client_turn_id, .. } => { + Self::mark_queue_events_cleanable(thread, client_turn_id); + true + } } } @@ -471,6 +669,18 @@ impl DirectThreadManager { } } + /// 一条待发消息离开队列:它和它对应的入队事件一起变成可回收。 + 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) + { + stored.cleanable = true; + } + } + } + fn trim_prefix(thread: &mut ThreadState) { let min_cursor = thread .subscribers @@ -584,6 +794,56 @@ pub(crate) fn accept_direct_thread_turn( Ok(()) } +/// 入队一条待发消息(队列的成员与顺序就是事件列表,见 `direct_thread_queue`)。 +pub(crate) fn enqueue_direct_pending_turn( + thread_id: &str, + turn: PendingDirectTurn, +) -> Result { + let outcome = { + global_direct_thread_manager() + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .enqueue_pending_turn(thread_id, turn) + }; + if matches!(outcome, Ok(EnqueueOutcome::Enqueued)) { + notify_direct_thread_subscribers(thread_id); + } + outcome +} + +/// 取消一条待发消息:结果告诉调用方它是被取消、已经放行,还是本来就不在队里。 +pub(crate) fn remove_direct_pending_turn( + thread_id: &str, + client_turn_id: &str, +) -> DirectQueueRemovalOutcome { + let outcome = { + global_direct_thread_manager() + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .remove_pending_turn(thread_id, client_turn_id, direct_tool_call_now_ms()) + }; + if outcome == DirectQueueRemovalOutcome::Removed { + notify_direct_thread_subscribers(thread_id); + } + outcome +} + +/// 认领队首并登记占用(放行的原子步骤,见 [`DirectThreadManager::claim_pending_turn`])。 +/// +/// `None` = 现在不该放行(已有未收口的回合,或队列为空)。 +pub(crate) fn claim_direct_pending_turn(thread_id: &str) -> Option { + let claimed = { + global_direct_thread_manager() + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .claim_pending_turn(thread_id, direct_tool_call_now_ms()) + }; + if claimed.is_some() { + notify_direct_thread_subscribers(thread_id); + } + claimed +} + /// 运行时回填某一轮逻辑回合的进度(状态 / 活动 / 序号)。返回是否真的写进去了。 pub(crate) fn update_direct_thread_active_turn( thread_id: &str, @@ -683,7 +943,9 @@ pub(crate) fn consume_direct_thread( #[cfg(test)] mod tests { use super::*; - use crate::agent::{DirectThreadDeltaKind, DirectThreadItem, DirectThreadRequestKind}; + use crate::agent::{ + DirectThreadDeltaKind, DirectThreadItem, DirectThreadRequestKind, MAX_PENDING_DIRECT_TURNS, + }; fn message(item_id: &str) -> DirectThreadItem { DirectThreadItem::Message { @@ -1093,4 +1355,206 @@ mod tests { assert_eq!(conflict.as_deref(), Some("turn-1")); } + + fn pending_turn(client_turn_id: &str) -> PendingDirectTurn { + PendingDirectTurn::prepare( + client_turn_id.to_string(), + serde_json::from_value(serde_json::json!({ + "type": "message", + "role": "user", + "content": [{"type": "input_text", "text": format!("消息 {client_turn_id}")}], + "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: &[DirectThreadEvent]) -> Vec<&str> { + events + .iter() + .filter_map(|event| match event { + DirectThreadEvent::QueueEnqueued { client_turn_id, .. } => { + Some(client_turn_id.as_str()) + } + _ => None, + }) + .collect() + } + + /// 队列就是事件列表:新订阅者的 bootstrap 看得见当前队列,取消之后那条就从队列里消失。 + #[test] + fn new_subscribers_see_the_current_queue_until_it_is_removed() { + let mut manager = DirectThreadManager::with_limits(100, 100_000); + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-1")) + .expect("enqueue head"); + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-2")) + .expect("enqueue tail"); + + let bootstrap = manager.subscribe("thread-1"); + assert_eq!(enqueued_ids(&bootstrap.events), vec!["turn-1", "turn-2"]); + + let watcher = manager.subscribe("thread-1"); + let _ = manager.consume(&watcher.subscription_id).expect("drain"); + assert_eq!( + manager.remove_pending_turn("thread-1", "turn-1", FIXED_AT_MS), + DirectQueueRemovalOutcome::Removed + ); + // 离开队列的事件带 typed 原因:这条是用户取消。 + let events = manager + .consume(&watcher.subscription_id) + .expect("consume removal") + .events; + assert!( + matches!( + events.as_slice(), + [DirectThreadEvent::QueueRemoved { client_turn_id, reason, .. }] + if client_turn_id == "turn-1" && *reason == DirectQueueRemovalReason::Cancelled + ), + "{events:?}" + ); + // 队尾那条仍在,队首那条不再对新订阅者可见。 + let bootstrap = manager.subscribe("thread-1"); + assert_eq!(enqueued_ids(&bootstrap.events), vec!["turn-2"]); + } + + /// 入队幂等:判重范围是「在队 ∪ 正在跑的那一轮」,重复入队不排第二条、不发第二条事件。 + #[test] + fn enqueue_is_idempotent_per_client_turn_id() { + let mut manager = DirectThreadManager::with_limits(100, 100_000); + assert_eq!( + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-1")) + .expect("enqueue"), + EnqueueOutcome::Enqueued + ); + assert_eq!( + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-1")) + .expect("enqueue again"), + EnqueueOutcome::AlreadyKnown + ); + assert_eq!(manager.pending_turn_ids("thread-1"), vec!["turn-1"]); + + // 正在跑的那一轮同样算已知:它的身份已经在占用里,不该再排一条等放行。 + manager + .claim_pending_turn("thread-1", FIXED_AT_MS) + .expect("claim head"); + assert_eq!( + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-1")) + .expect("enqueue while running"), + EnqueueOutcome::AlreadyKnown + ); + assert!(manager.pending_turn_ids("thread-1").is_empty()); + } + + /// 容量只数在队条目:正在跑的那一轮不占等待位。 + #[test] + fn capacity_counts_only_pending_turns() { + let mut manager = DirectThreadManager::with_limits(100, 100_000); + for index in 0..MAX_PENDING_DIRECT_TURNS { + assert_eq!( + manager + .enqueue_pending_turn("thread-1", pending_turn(&format!("turn-{index}"))) + .expect("enqueue"), + EnqueueOutcome::Enqueued + ); + } + assert_eq!( + manager.enqueue_pending_turn("thread-1", pending_turn("turn-overflow")), + Err(EnqueueRejection::QueueFull) + ); + + manager + .claim_pending_turn("thread-1", FIXED_AT_MS) + .expect("claim head"); + assert_eq!( + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-overflow")) + .expect("enqueue after dispatch"), + EnqueueOutcome::Enqueued + ); + } + + /// 放行认领的是**队首**,而且开始事件与离开队列事件同批写:中间没有第二个中间态。 + #[test] + fn dispatch_claims_the_head_and_writes_both_events_in_one_batch() { + let mut manager = DirectThreadManager::with_limits(100, 100_000); + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-1")) + .expect("enqueue head"); + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-2")) + .expect("enqueue tail"); + let watcher = manager.subscribe("thread-1"); + let _ = manager.consume(&watcher.subscription_id).expect("drain"); + + let claimed = manager + .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!(manager.pending_turn_ids("thread-1"), vec!["turn-2"]); + + let events = manager + .consume(&watcher.subscription_id) + .expect("consume dispatch") + .events; + assert!( + matches!( + events.as_slice(), + [ + DirectThreadEvent::TurnStarted { at, .. }, + DirectThreadEvent::QueueRemoved { client_turn_id, reason, at: removed_at }, + ] if *at == Some(2_000) + && client_turn_id == "turn-1" + && *reason == DirectQueueRemovalReason::Dispatched + && *removed_at == Some(2_000) + ), + "{events:?}" + ); + + // 已经有未收口的回合:认领不再放行第二条。 + assert!(manager.claim_pending_turn("thread-1", 2_100).is_none()); + assert_eq!(manager.pending_turn_ids("thread-1"), vec!["turn-2"]); + } + + /// 取消的三种结果分得开:真的移除了、已经被放行了、本来就不在队里。 + #[test] + fn removal_tells_cancelled_dispatched_and_unknown_apart() { + let mut manager = DirectThreadManager::with_limits(100, 100_000); + assert_eq!( + manager.remove_pending_turn("thread-1", "turn-1", FIXED_AT_MS), + DirectQueueRemovalOutcome::NotFound + ); + + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-1")) + .expect("enqueue"); + assert_eq!( + manager.remove_pending_turn("thread-1", "turn-1", FIXED_AT_MS), + DirectQueueRemovalOutcome::Removed + ); + assert_eq!( + manager.remove_pending_turn("thread-1", "turn-1", FIXED_AT_MS), + DirectQueueRemovalOutcome::NotFound + ); + + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-2")) + .expect("enqueue"); + manager + .claim_pending_turn("thread-1", FIXED_AT_MS) + .expect("claim head"); + assert_eq!( + manager.remove_pending_turn("thread-1", "turn-2", FIXED_AT_MS), + DirectQueueRemovalOutcome::AlreadyDispatched + ); + } } diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_queue.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_queue.rs new file mode 100644 index 000000000..2785c87cb --- /dev/null +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_queue.rs @@ -0,0 +1,200 @@ +//! DirectProject 待发消息队列的条目与规则。 +//! +//! 队列的成员与顺序**就是 Thread Manager 的事件列表本身**:一条待发消息在队期间,它的 +//! `queue.enqueued` 事件不可回收;离开队列(取消或放行)时才转成可回收。新订阅者的 bootstrap 因此 +//! 天然看得见当前队列,不需要第二张队列表,也不会有"事件与队列不一致"的窗口。 +//! +//! 这个模块只放三件事:一条待发消息带走什么([`PendingDirectTurn`])、容量规则 +//! ([`MAX_PENDING_DIRECT_TURNS`])、以及它在线上长什么样([`PendingDirectTurn::enqueued_event`])。 +//! 它不碰锁、不碰 Tauri、不写盘:入队检查在命令侧,放行顺序在 Thread Manager。 + +use serde_json::Value; + +use crate::agent::{ + direct_codex_user_item_id_for_client_turn_id, DirectCodexUserItem, DirectThreadEvent, +}; + +/// 一个项目最多能同时排队的待发消息条数。 +/// +/// 只数**在队**条目,不算正在跑的那一轮。上限只落在宿主这一处:前端不再自己数,满队由命令返回 +/// typed 入队失败。 +pub(crate) const MAX_PENDING_DIRECT_TURNS: usize = 5; + +/// 一条已经通过入队检查、正在等放行的用户消息。 +/// +/// 只在内存里,进程重启即消失(与 ADR 记的边界一致)。它同时是**放行时要用的全部输入**:放行 +/// 没有失败出口,所以检查产物在入队时就地冻结,放行只搬运、不重算。 +#[derive(Clone, Debug)] +pub(crate) struct PendingDirectTurn { + /// 这条消息的回合身份;放行后同一轮的 `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 PendingDirectTurn { + /// 组一条待发消息:入队检查已经全部通过,这里只把放行要用的产物冻结下来。 + /// + /// 冻结是刻意的:放行没有失败出口,所以任何可能在放行时才失败的计算都必须提前到这里 + /// (canonical 形状与 prompt 都是)。 + pub(crate) fn prepare( + 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 { + client_turn_id, + user_item, + canonical_user_item, + prompt, + creation_type, + at, + }) + } + + /// 入队事件的投影:带 canonical 用户条目与 `creationType`,**不带 prompt**(prompt 只留在宿主的 + /// 队列条目里,它不是要下发的展示形状)。 + pub(crate) fn enqueued_event(&self) -> DirectThreadEvent { + DirectThreadEvent::queue_enqueued( + self.client_turn_id.clone(), + self.user_item.clone(), + self.creation_type.clone(), + self.at, + ) + } + + /// 这条待发消息在历史里的用户条目 id:前端用它把 chip 与回合边界对上。 + pub(crate) fn user_item_id(&self) -> Option { + direct_codex_user_item_id_for_client_turn_id(&self.client_turn_id) + } +} + +/// 入队的结果。 +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub(crate) enum EnqueueOutcome { + /// 这次真的排进队尾了。 + Enqueued, + /// 同一 `clientTurnId` 已经在队(或正在跑):按幂等返回成功,不排第二条、不发事件。 + AlreadyKnown, +} + +/// 入队被检查挡下来的原因。满队之外的原因由入队半自己的 typed 错误表达。 +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub(crate) enum EnqueueRejection { + /// 在队条目已达 [`MAX_PENDING_DIRECT_TURNS`]。 + QueueFull, +} + +/// 队列还有没有位置。`pending_count` 只数在队条目。 +pub(crate) fn queue_has_room(pending_count: usize) -> Result<(), EnqueueRejection> { + if pending_count >= MAX_PENDING_DIRECT_TURNS { + return Err(EnqueueRejection::QueueFull); + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + fn user_item(text: &str, id: &str) -> DirectCodexUserItem { + serde_json::from_value(json!({ + "type": "message", + "role": "user", + "content": [{"type": "input_text", "text": text}], + "id": id, + })) + .expect("canonical user item") + } + + fn pending(client_turn_id: &str, creation_type: Option<&str>) -> PendingDirectTurn { + PendingDirectTurn::prepare( + 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 是宿主的入队检查产物, + /// 不许顺着事件下发。 + #[test] + fn enqueued_event_carries_the_canonical_item_and_no_prompt() { + let event = pending("turn-1", Some("web-game")).enqueued_event(); + let value = serde_json::to_value(&event).expect("serialize queue.enqueued"); + assert_eq!( + value, + json!({ + "type": "queue.enqueued", + "clientTurnId": "turn-1", + "userItem": { + "type": "message", + "role": "user", + "content": [{"type": "input_text", "text": "生成一个游戏"}], + "id": "direct-codex:turn-1:user", + }, + "creationType": "web-game", + "at": 1_700_000_000_123u64, + }) + ); + assert!(value.get("prompt").is_none(), "{value}"); + + // 没有创建类型时不写字段,也不补 `null`。 + let bare = serde_json::to_value(pending("turn-2", None).enqueued_event()) + .expect("serialize queue.enqueued without creation type"); + assert!(bare.get("creationType").is_none(), "{bare}"); + + // 事件读得出自己的待发消息身份。 + assert_eq!(event.queue_client_turn_id(), Some("turn-1")); + 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") + ); + } + + /// 用户条目身份由 `clientTurnId` 派生,与落盘 / 下发用的是同一个函数。 + #[test] + fn user_item_id_derives_from_the_client_turn_id() { + assert_eq!( + pending("turn-1", None).user_item_id().as_deref(), + Some("direct-codex:turn-1:user") + ); + assert_eq!(pending(" ", None).user_item_id(), None); + } + + /// 容量只数在队条目,5 条封口;在跑的那一轮不算进去。 + #[test] + fn capacity_closes_at_five_pending_turns() { + for count in 0..MAX_PENDING_DIRECT_TURNS { + assert_eq!(queue_has_room(count), Ok(()), "{count} 条时仍该有位置"); + } + assert_eq!( + queue_has_room(MAX_PENDING_DIRECT_TURNS), + Err(EnqueueRejection::QueueFull) + ); + } +} 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 index 952dd3a24..a4649938e 100644 --- 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 @@ -13,6 +13,7 @@ use crate::agent::redact_secret_tokens; use crate::agent::sanitize_error_context; +use crate::agent::DirectCodexUserItem; use crate::agent::DirectTurnFailureKind; use crate::redact_absolute_path_tokens; use serde::{Deserialize, Serialize}; @@ -212,6 +213,33 @@ impl DirectThreadRequestKind { } } +/// 一条待发消息离开队列的原因。 +/// +/// typed 枚举,取值即语义:取消是用户在输入盒上撤掉这条消息,放行是它已经接单并成为回合 +/// (同一临界区里另有 `turn.started`)。界面按它分流,不解析字符串。 +#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize, TS)] +#[serde(rename_all = "camelCase")] +#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))] +pub(crate) enum DirectQueueRemovalReason { + /// 用户取消了这条待发消息。 + Cancelled, + /// 放行:这条待发消息已经接单并成为回合。 + Dispatched, +} + +/// 取消一条待发消息的结果。 +#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize, TS)] +#[serde(rename_all = "camelCase")] +#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))] +pub(crate) enum DirectQueueRemovalOutcome { + /// 已从队列移除。 + Removed, + /// 这条消息已经被放行(正在跑的那一轮就是它),不能按待发消息取消。 + AlreadyDispatched, + /// 队列里没有这个身份,也没有在跑的一轮是它。 + NotFound, +} + /// 失败终态的可下发载荷(`turn.completed.status == "failed"` 时必有,其余终态没有)。 /// /// `kind` 是稳定分类,只给界面选语气,不参与流程分支;`message` 是**已在宿主侧脱敏并截断**的 @@ -259,7 +287,7 @@ impl DirectTurnFailure { /// 不带回合身份,这个字段只用来把"这一轮的边界属于哪条用户消息"讲清楚:前端在只有生命周期锚点 /// + 历史切片、运行态一直为空时也能按身份认领开口条目,不必靠时间戳猜。缺失表示身份不可证明 /// (旧事件、没有开口用户条目、取消时拿不到 clientTurnId),此时前端不得补造。 -#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, TS)] +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, TS)] #[serde(tag = "type", rename_all_fields = "camelCase", deny_unknown_fields)] #[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))] pub(crate) enum DirectThreadEvent { @@ -321,6 +349,34 @@ pub(crate) enum DirectThreadEvent { #[serde(default)] request_id: Option, }, + /// 待发消息入队:数组顺序就是队首到队尾的顺序。 + /// + /// 这条事件在条目仍在队期间**不可回收**,离开队列(取消或放行)时才转成可回收——新订阅者 + /// 靠这一点在 bootstrap 里看到当前队列,`is_bootstrap_event` 不需要为它加特例。 + /// + /// 事件不带 prompt:prompt 是入队检查的产物,只留在宿主的队列条目里。 + #[serde(rename = "queue.enqueued")] + QueueEnqueued { + /// 这条待发消息的回合身份;放行后同一轮的 `turn.started` / `turn.completed` 用它。 + client_turn_id: String, + /// canonical 用户条目:前端据此派生 chip 文案,Rust 不渲染展示形状。 + user_item: DirectCodexUserItem, + #[serde(default, skip_serializing_if = "Option::is_none")] + #[ts(optional)] + creation_type: Option, + /// 入队那一刻的宿主毫秒钟。 + #[ts(as = "f64")] + at: u64, + }, + /// 待发消息离开队列:`reason` 是取消还是放行。 + #[serde(rename = "queue.removed")] + QueueRemoved { + client_turn_id: String, + reason: DirectQueueRemovalReason, + #[serde(default, skip_serializing_if = "Option::is_none")] + #[ts(optional, as = "Option")] + at: Option, + }, } impl DirectThreadEvent { @@ -417,6 +473,51 @@ impl DirectThreadEvent { Self::Request { kind, request_id } } + /// 待发消息入队事件。 + pub(crate) fn queue_enqueued( + client_turn_id: String, + user_item: DirectCodexUserItem, + creation_type: Option, + at: u64, + ) -> Self { + Self::QueueEnqueued { + client_turn_id, + user_item, + creation_type, + at, + } + } + + /// 待发消息离开队列事件。 + pub(crate) fn queue_removed( + client_turn_id: String, + reason: DirectQueueRemovalReason, + at: u64, + ) -> Self { + Self::QueueRemoved { + client_turn_id, + reason, + at: Some(at), + } + } + + /// 这条事件属于哪条待发消息:只有队列事件有。 + pub(crate) fn queue_client_turn_id(&self) -> Option<&str> { + match self { + Self::QueueEnqueued { client_turn_id, .. } + | Self::QueueRemoved { client_turn_id, .. } => Some(client_turn_id), + _ => None, + } + } + + /// 待发消息离开队列的原因:只有 `queue.removed` 有。 + pub(crate) fn queue_removal_reason(&self) -> Option { + match self { + Self::QueueRemoved { reason, .. } => Some(*reason), + _ => None, + } + } + /// 事件级阶段时间(毫秒):只有四种生命周期事件有,其余事件返回 `None`。 /// /// 只读已存入事件的值,不在读取时取钟——重放要用的就是原事件的时间。 @@ -427,6 +528,8 @@ impl DirectThreadEvent { | Self::TurnCompleted { at, .. } | Self::ItemStarted { at, .. } | Self::ItemCompleted { at, .. } => *at, + Self::QueueEnqueued { at, .. } => Some(*at), + Self::QueueRemoved { at, .. } => *at, Self::ItemDelta { .. } | Self::Request { .. } => None, } } @@ -456,7 +559,7 @@ impl DirectThreadEvent { } } -#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, TS)] +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, TS)] #[serde(rename_all = "camelCase", deny_unknown_fields)] #[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))] pub(crate) struct DirectThreadSubscriptionBootstrap { @@ -468,7 +571,7 @@ pub(crate) struct DirectThreadSubscriptionBootstrap { pub(crate) events: Vec, } -#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, TS)] +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, TS)] #[serde(rename_all = "camelCase", deny_unknown_fields)] #[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))] pub(crate) struct DirectThreadConsumeResult { @@ -1466,4 +1569,61 @@ mod tests { .expect("failed turn without failure payload"); assert_eq!(sparse.failure(), None); } + + /// 待发消息离开队列:`reason` 是 typed 枚举(`cancelled` / `dispatched`),界面按取值分流, + /// 不解析字符串。`at` 缺省时反序列化仍是 `None`。 + #[test] + fn queue_removed_carries_a_typed_reason() { + let dispatched = DirectThreadEvent::queue_removed( + "turn-1".to_string(), + DirectQueueRemovalReason::Dispatched, + 2_000, + ); + assert_eq!( + serde_json::to_value(&dispatched).expect("serialize queue.removed"), + json!({ + "type": "queue.removed", + "clientTurnId": "turn-1", + "reason": "dispatched", + "at": 2_000u64, + }) + ); + assert_eq!( + serde_json::from_value::( + serde_json::to_value(&dispatched).expect("serialize") + ) + .expect("round trip"), + dispatched + ); + assert_eq!(dispatched.queue_client_turn_id(), Some("turn-1")); + assert_eq!( + dispatched.queue_removal_reason(), + Some(DirectQueueRemovalReason::Dispatched) + ); + + // 未识别的取值必须失败关闭:队列归属是宿主事实,不能让界面猜。 + serde_json::from_value::(json!({ + "type": "queue.removed", + "clientTurnId": "turn-1", + "reason": "timeout", + })) + .expect_err("unknown queue removal reason must fail closed"); + + // 没有 `at` 的老形状仍能反序列化,回写不补 `null`。 + let legacy: DirectThreadEvent = serde_json::from_value(json!({ + "type": "queue.removed", + "clientTurnId": "turn-1", + "reason": "cancelled", + })) + .expect("queue.removed without at"); + assert_eq!(legacy.at(), None); + assert_eq!( + serde_json::to_value(legacy).expect("serialize legacy"), + json!({ + "type": "queue.removed", + "clientTurnId": "turn-1", + "reason": "cancelled", + }) + ); + } } diff --git a/apps/ai-game-creator-shell/src/view/project-development/chat/generated/DirectQueueRemovalOutcome.ts b/apps/ai-game-creator-shell/src/view/project-development/chat/generated/DirectQueueRemovalOutcome.ts new file mode 100644 index 000000000..5a70720f9 --- /dev/null +++ b/apps/ai-game-creator-shell/src/view/project-development/chat/generated/DirectQueueRemovalOutcome.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 DirectQueueRemovalOutcome = + | 'removed' + | 'alreadyDispatched' + | 'notFound'; diff --git a/apps/ai-game-creator-shell/src/view/project-development/chat/generated/DirectQueueRemovalReason.ts b/apps/ai-game-creator-shell/src/view/project-development/chat/generated/DirectQueueRemovalReason.ts new file mode 100644 index 000000000..90f8b8247 --- /dev/null +++ b/apps/ai-game-creator-shell/src/view/project-development/chat/generated/DirectQueueRemovalReason.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. + +/** + * 一条待发消息离开队列的原因。 + * + * typed 枚举,取值即语义:取消是用户在输入盒上撤掉这条消息,放行是它已经接单并成为回合 + * (同一临界区里另有 `turn.started`)。界面按它分流,不解析字符串。 + */ +export type DirectQueueRemovalReason = 'cancelled' | 'dispatched'; diff --git a/apps/ai-game-creator-shell/src/view/project-development/chat/generated/DirectThreadEvent.ts b/apps/ai-game-creator-shell/src/view/project-development/chat/generated/DirectThreadEvent.ts index da0136295..630a41e95 100644 --- a/apps/ai-game-creator-shell/src/view/project-development/chat/generated/DirectThreadEvent.ts +++ b/apps/ai-game-creator-shell/src/view/project-development/chat/generated/DirectThreadEvent.ts @@ -1,4 +1,6 @@ // 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'; +import type { DirectQueueRemovalReason } from './DirectQueueRemovalReason'; import type { DirectThreadDeltaKind } from './DirectThreadDeltaKind'; import type { DirectThreadItem } from './DirectThreadItem'; import type { DirectThreadRequestKind } from './DirectThreadRequestKind'; @@ -84,8 +86,26 @@ export type DirectThreadEvent = kind: DirectThreadDeltaKind; delta: string; } + | { type: 'request'; kind: DirectThreadRequestKind; requestId: string | null } | { - type: 'request'; - kind: DirectThreadRequestKind; - requestId: string | null; + type: 'queue.enqueued'; + /** + * 这条待发消息的回合身份;放行后同一轮的 `turn.started` / `turn.completed` 用它。 + */ + clientTurnId: string; + /** + * canonical 用户条目:前端据此派生 chip 文案,Rust 不渲染展示形状。 + */ + userItem: DirectCodexUserItem; + creationType?: string; + /** + * 入队那一刻的宿主毫秒钟。 + */ + at: number; + } + | { + type: 'queue.removed'; + clientTurnId: string; + reason: DirectQueueRemovalReason; + at?: number; }; diff --git a/docs/README.md b/docs/README.md index 20a519f46..e183cb872 100644 --- a/docs/README.md +++ b/docs/README.md @@ -48,6 +48,8 @@ - [引用候选由宿主注入](./adr/【ADR】引用候选由宿主注入-2026-09-22.md):引用输入区只接受宿主注入的引用 provider,素材选择面板独立成组件,附件芯片成为本轮附件唯一事实源。 - [DirectProject 命令接单化](./adr/【ADR】DirectProject命令接单化-2026-09-23.md):命令只负责接单、事件流回答整轮结果;拒单前置、失败后置。 - [DirectProject 命令接单化实施计划](./technical/【实施计划】DirectProject命令接单化-2026-09-23.md):四步落地顺序、每步不变式与验收;四步均已落地。 +- [DirectProject 命令入队化与待发消息队列归宿主](./adr/【ADR】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md):命令只负责入队,放行归 Thread Manager;待发消息队列作为运行态事件归宿主、前端只投影;CLI 直连入口与调用身份守卫一并退役。 +- [DirectProject 命令入队化与待发消息队列归宿主实施计划](./technical/【实施计划】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md):五步落地顺序、每步不变式与验收;待实施。 - [GameAgent 对话工具调用卡片](./technical/【技术方案】GameAgent对话工具调用卡片-2026-09-14.md):把右侧对话里的执行命令 / 写文件投影成 Codex 风格可折叠卡片,含采集、独立历史文件、事件字段与回读契约。 - [DirectProject 客户端 Skill 与 MCP 扩展导入方案](./technical/【技术方案】DirectProject客户端Skill与MCP扩展导入方案-2026-08-31.md):客户端扩展导入、按独立 Skill/MCP 拆分、命名、启用和启动时注入边界。 - [AGC 通用插件宿主与编辑器适配](./technical/【技术方案】AGC通用插件宿主与编辑器适配-2026-09-09.md):通用插件宿主、SDK、权限审计、UI 挂载和 Cocos 编辑器适配边界。 diff --git a/docs/adr/【ADR】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md b/docs/adr/【ADR】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md new file mode 100644 index 000000000..d4e741723 --- /dev/null +++ b/docs/adr/【ADR】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md @@ -0,0 +1,139 @@ +# 【ADR】DirectProject命令入队化与待发消息队列归宿主 + +状态:已接受(2026-09-24 设计定稿;**尚未落地**,实施顺序与验收见 +[`【实施计划】DirectProject命令入队化与待发消息队列归宿主-2026-09-24`](../technical/【实施计划】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md)) + +## 背景 + +待发消息队列今天是纯前端状态:回合运行中用户再发送就进本地 FIFO +(`chat/components/DirectProjectComposer/chatComposerQueue.ts`),回合终态事件到达后由**恰好开着的那个窗口**放行队首 +(`useDirectProjectChatController.ts` 的完成计数 effect 与 `startTurn` 的 `finally`)。三个问题: + +1. 队列是这条对话里唯一没有宿主持有者的事实:另一个窗口、另一个订阅者,或只是离开工作台再回来,都看不到已经排了什么队。 +2. 放行落在"哪个窗口恰好开着"上。队列一旦共享(本 ADR 要做的),两个窗口都会去放行队首,必然双发。 +3. 命令边界今天写的是"接单":校验通过就登记占用、落盘用户条目、起整轮。但用户按下发送时想要的是"这条消息会被依次处理"—— + 命令的成功含义与用户意图之间隔着一次长度未定义的等待。 + +前置口径:`【ADR】DirectProject命令接单化-2026-09-23` 已把逻辑回合收归 Thread Manager,并在 §8 留了 TODO +「以后这条队列挪到 Rust 端,落点就是 Thread Manager 的接单动作」,同时把 Rust 端发送队列列进"明确不做"。 +本 ADR 就是那条 TODO 的收口,并顺带收掉两处已经没有现役价值的实现。 + +## 决策 + +### 1. 词表:入队 / 入队失败 / 放行 + +- 命令边界的成功与失败改叫 **入队 / 入队失败**;「接单」「拒单」两个词退役,不再出现在文档、注释、标识符与测试名里。 +- 旧「接单」在语义上的角色(这一轮真正成立的那一刻)改叫 **放行**。 +- 因此旧句「接单成立 ⇔ 事件流里有开始有结束」要改写成「**放行成立** ⇔ 事件流里有开始有结束」。这不是换词而是角色搬家: + 逐处改写时按角色判,不做字面替换(命令边界→入队/入队失败;回合成立→放行;检查归属→入队时的检查)。 + +### 2. 命令 = 入队 + +`chat_with_game_creator_direct_codex` 改义改名(建议 `enqueue_direct_codex_turn`),主体是今天"接单前"那条链**原样**跑完: +`clientTurnId` 校验 → 工作流恢复 → 用户条目校验 → prompt 投影 → 前置条件 → 工程准备;通过后**只入队**—— +追加 `queue.enqueued`,不登记占用、不落盘用户条目、不发回合事件、不起 codex。 +任何一步失败就是**入队失败**,走命令返回的 typed 载荷(`DirectTurnRejection` → `DirectTurnEnqueueFailure`),用户就在现场。 + +入队按 `clientTurnId` 幂等,判重范围是"在队 ∪ 正在跑的那一轮";重复入队返回同一次成功,不再是并发拒单。 + +### 3. 队列归 Thread Manager + +- 每个 thread(线程身份就是项目规范路径)一条 FIFO;本期只支持**按顺序追加**与**按身份移除**,不做重排、优先级、编辑。 +- 队列状态**就是事件列表本身**:`queue.enqueued` 在条目仍在队期间不可回收,离开队列(取消或放行)时转为可回收并被既有规则回收。 + 不新增第二张队列表;`subscribe` 的 live-set bootstrap 因此天然把当前队列交给新订阅者——这就是"加入者看到的那几条"。 +- 上限 5 只数**在队条目**(不数正在跑的那一轮),由 Rust 持有;满队时入队失败并给出既有提示文案。 + +### 4. 线上形状 + +- `queue.enqueued { clientTurnId, userItem, creationType?, at }`:`userItem` 是 canonical 用户条目,前端据此派生 chip 文案, + Rust 不渲染、不裁成展示形状。事件**不带** prompt——prompt 是入队检查的产物,只留在宿主的队列条目里。 +- `queue.removed { clientTurnId, reason }`:`reason` 是 ts-rs 导出的 **typed 枚举**(`cancelled | dispatched`),永不用字符串, + 形状与 `turn.completed{status, failure?}` 同构。 +- 两条事件与其它运行态事件同一条流、同一个 reducer。 + +### 5. 放行 = 旧「接单」的语义角色 + +Thread Manager 在一个回合收口**之后**原子地做:取队首 → 登记占用 → 发 `turn.started` → 落盘用户条目 → 下发用户条目 → 起整轮, +并在同一临界区写 `queue.removed{ dispatched }`(chip 消失与气泡出现在同一批 consume 里,中间没有空窗)。 + +**放行不重跑任何检查,也不存在"放行失败"这种状态**:放行之后的一切失败都是**回合失败**,走既有 `turn.completed.failure` 通道; +不新增任何失败通道,也不为放行补拒单出口。 + +### 6. 放行的触发与监护顺序 + +- 唯一放行点:回合任务收尾之后的 `kick`(正常 / 失败 / 中止三条路径共用,外加一个 drop 守卫盖 panic),幂等,并且在临界区里原子认领队首。 +- 入队时也踢一脚:「队列非空 + 线程空闲」是合法状态,对应今天的"直接发送"。 +- 入队**不取** `DirectTaonierActiveInvocationGuard`:它必须整轮持有(它是这一轮的调用身份,付费美术、执行会话、MCP、校验、 + 上下文预取都靠它把工作归属到自己那一轮),入队若取它等于"回合运行中不能入队"。入队路径的并发由工程准备自身的项目写锁与队列兜。 + +### 7. 渲染侧 + +前端不再持有队列副本:chip 只由事件投影(入队被拒时只由命令返回值给反馈,保留既有提示文案与"不丢草稿"行为); +忙态只由事件投影加「队列非空」指示。排队消息在放行前**不写** `project.jsonl`,用户气泡仍然只来自宿主条目("落盘即放行"不变)。 + +### 8. 作用域与寿命 + +每项目一条、全进程共享:A 窗口排队 B 窗口可见可取消;切项目、离开工作台、关窗口都不影响队列继续放行;进程结束队列消失。 +**不做跨进程持久化**。 + +### 9. 顺带退役 + +- `--direct-codex-chat` 整个退役(解析、派发,以及只服务它的 `run_direct_game_creator_turn_at` 一对函数)。 + 它的历史用途只有一个:手工生产验证夹具 `scripts/direct-execution-production-fixture.mjs`(PR #439 引入,不在 CI、 + 没有 npm 脚本或 harness 注册、没有任何测试钉它,唯一硬依赖是进程退出码)。没有产品入口价值,也不该在入队化之后 + 成为第二条直接起回合的路径;夹具脚本一并退役。 +- CLI 一退,`DirectTurnError::TurnAlreadyRunning` 的两个生产点(调用身份守卫、占用登记)都没有调用方, + 它连同前端"同一轮消息仍在处理中"文案、专属分支与测试一起删。 +- `DirectTaonierActiveInvocationGuard` 的**身份**与 Thread Manager 的 `active_turn.turn_id` 是同一件事的两份记录 + (GUI 路径下同源字符串),而 09-23 ADR 立的是"同一件事只许有一处真相"。CLI 退役后它的硬阻塞消失:第二步让那五个读者 + 改读 Thread Manager 的活动回合身份,守卫连同它的测试与 60 秒卡死兜底一起删。删之前必须确认取消路径的无条件终态 + 仍然解得开"回合卡死"(那次兜底来自真实事故 `d833ca9d3`)。 + +## 备选方案与取舍 + +1. **状态归宿主、放行留前端**(再用一条 `claim` 命令做 CAS 认领):前端机器原样保留,但队列的继续推进依赖至少一个窗口活着, + 且"谁去认领"要靠竞态解决——正是要消掉的东西。作废。 +2. **入队只做形状校验、把前置检查留到放行**:会造出"放行期拒单"——一个没有调用方在等的失败,得为它发明新通道; + 而且"首个回合还在建工程、第二条已经排队"这类合法流程会被入队误拒。作废(检查跟着入队走)。 +3. **一条 `queue.changed{items:[…]}` 快照事件**代替两条细粒度事件:reducer 更傻,但每次变更搬全量、与既有细粒度事件风格不一致。不选。 +4. **`reason` 用字符串**:前端只能猜、无法穷举、无法在类型层穷尽分支。不选,用 typed 枚举。 +5. **保留 CLI**:省掉夹具改写,但等于为手工验证工具长期保留第二条直接起回合的路径,与"命令只有入队一个入口"冲突。不选。 + +## 影响与代价 + +- **入队即写盘**:引用 UI 设计文档的条目在 prompt 投影时会生成 `ui/generated-.js`(`ui_editor/persistence.rs`); + 从队列里取消不撤该文件。 +- **检查是时间点事实**:放行不重跑,按入队那一刻的结论放行;manifest、权限、目录在入队之后变化也照旧放行,偏差落到回合失败。 +- **队列随进程消失**:待发消息只在内存与事件流里,`kill -9` / 退出后重进看不到(与 09-16 ADR 的 `kill -9` 口径一致)。 +- **退役面**:前端 `completionPendingRef`、`handledCompletedTurnCountRef`、`busyBaselineTurnCountRef`、`dispatchNextQueuedTurn`、 + `turnBusyRef`/`beginTurnBusy`/`endTurnBusy`、`queueSequenceRef`、`queuedTurns*`、排队条目那部分 `pendingRunAnalyticsRef` + 与 `chatComposerQueue.ts` 的队列逻辑整体退役;`directProjectTurnStatus` 的"命令在飞"分支收成 IPC 在飞; + 埋点句柄改由宿主在放行时开。 +- **词表切换是一次性跨文档动作**:仓库里「接单/拒单」共 384 处,并非全属同一个域 + (`features/agent-runtime` 的"拒单文案"属另一个域,改成"请求被拒";`单测` 这类是假阳性)。 + 两份已接受的 ADR 保留正文与文件名,顶部加词表注记。 +- **失去手工验证手段**:CLI 退役同时带走"真实二进制驱动执行层生产验证"这条手工路径(夹具脚本一并退役)。 + 它不在 CI,损失的是排障时的一次性手段,不是门禁。 +- **埋点**:attempt id 只用于宿主内存里的候选配对与前端结算(同 `session_id` 的 `Settle` 才提交), + 前端不再需要为排队条目提前持有它;平台会话代际翻转的 `discard` 判定仍是渲染侧职责,落地时逐条核对。 + +## 明确不做 + +- 不做跨进程持久化(不在入队时写 `project.jsonl`:那会在历史里留下一条永远不会跑的假消息,破坏"用户气泡只来自宿主条目")。 +- 不做重排、优先级、编辑待发条目;也不做"队列挂起 / 继续"状态——入队化之后队列里不存在会被拒的条目。 +- 不给放行新增失败通道,不为"放行不重跑检查"补兜底。 +- 不保留 CLI,也不为夹具保留别名或兼容入口。 + +## 落地时要同步的文档与注释 + +- `CONTEXT.md`:词条(待发消息队列 / 待发消息 / 入队 / 入队失败 / 放行 / 逻辑回合 / 在途回合)。 +- `docs/adr/【ADR】DirectProject命令接单化-2026-09-23.md`:§8 的 TODO 与"明确不做:Rust 端发送队列"、§7 的"命令在飞"措辞、 + CLI 保持 await 的分工,全部改为入队化口径(正文保留,顶部加词表注记)。 +- `docs/technical/【技术方案】DirectProject Codex原始历史与异常恢复-2026-09-04.md`:事件协议一节补两个队列事件与放行时序。 +- `docs/technical/【实施计划】DirectProject命令接单化-2026-09-23.md`:第 4 步的队列 TODO 指向本 ADR。 +- `docs/technical/【技术方案】Direct回合行为审计账本-2026-08-31.md`、`【技术方案】DirectProject本轮附件路径映射-2026-08-31.md`、 + `docs/project-memory/plans/【里程碑】退役AGC项目对话斜杠命令与终端swarm chat入口-2026-09-22.md`:删掉对 CLI 入口的承诺 + (最后那处现在写在"范围外(保留)"里,口径反转)。 +- `docs/project-memory/shared-memory/decision-log.md`:追加本决定,并修正"CLI 保持 await"那条。 +- 代码注释:`agent/direct_runtime/user_input.rs` 模块注释里的 CLI 分工一段删掉;`direct_turn_accept.rs`、 + `direct_thread_manager.rs`、`useDirectProjectChatController.ts` 的旧措辞与 TODO 一并改。 diff --git a/docs/project-memory/shared-memory/decision-log.md b/docs/project-memory/shared-memory/decision-log.md index 53b75ec40..f8304d518 100644 --- a/docs/project-memory/shared-memory/decision-log.md +++ b/docs/project-memory/shared-memory/decision-log.md @@ -9578,3 +9578,36 @@ CI 上 `background_agent_runtime_recovers_stale_running_before_pending_task` 在 - 验证:`cargo test -p module-runtime --lib agc_models::`(4 passed)、`cargo test -p api-server --bin api-server agc` 与 `llm::`、AGC 客户端 `configuration::`、admin-web 页面定向 Vitest 与 typecheck、两套 workspace 的 `cargo fmt -- --check`、`check:encoding`/`check:doc-index`/`check:spacetime-schema`/`git diff --check`。 - 验证(真实上游 smoke,本地 dev DB):清空 `agc_model_catalog` 后启动 api-server → 日志 `已按上游模型列表初始化 AGC 模型目录 revision=1 model_count=6`;登录后 `GET /api/llm/models` 返回同一批模型、`displayName` 即上游原名、默认项为排序后第一项;上游不可达/非 2xx 时启动只记录 error、AGC 接口 `503` 且目录保持未初始化;目录已存在时重启不重写。 - 边界(未验证/残留):上游在售模型超过 32 条时同步会失败(目录项上限未改);`qwen-image-3.0` 这类图像模型会一起进入目录,是否对 AGC 隐藏由 owner 在后台停用;混合版本期间未升级的 api-server 会把自己的 AGC 接口打到 `503`,module 与 api-server 必须同批发布/回滚。 + +## 2026-09-24 命令入队化与待发消息队列归宿主:放行归 Thread Manager,CLI 直连入口退役 + +- 决策(词表):「接单 / 拒单」退役,命令边界的成功与失败改叫「入队 / 入队失败」;旧「接单」的语义角色 + (这一轮真正成立的那一刻)改叫「放行」。所以旧句"接单成立 ⇔ 事件流里有开始有结束"改成"放行成立 ⇔ …", + 逐处按角色改、不做字面替换;`features/agent-runtime` 的"拒单文案"属另一个域,改"请求被拒"。 +- 决策(命令 = 入队):`chat_with_game_creator_direct_codex` 改义改名(暂定 `enqueue_direct_codex_turn`), + 把今天"接单前"那条检查链原样跑完(身份 → 工作流恢复 → 用户条目校验 → prompt 投影 → 前置条件 → 工程准备), + 通过后只入队:不登记占用、不落盘用户条目、不发回合事件、不起 codex;失败走命令返回的 typed 载荷 + (`DirectTurnRejection` → `DirectTurnEnqueueFailure`)。入队按 `clientTurnId` 幂等(判重范围 = 在队 ∪ 在跑)。 +- 决策(队列归 Thread Manager):每项目一条 FIFO,只支持按顺序追加与按身份移除;队列的成员与顺序就是事件列表本身 + (`queue.enqueued` 在队期间不可回收、离开队列后可回收),`subscribe` 的 live-set bootstrap 因此天然把当前队列 + 交给中途加入的订阅者;上限 5 只数在队条目,由 Rust 持有。 +- 决策(线上形状):`queue.enqueued{clientTurnId,userItem,creationType?,at}` + `queue.removed{clientTurnId,reason}`, + `reason` 是 ts-rs 导出的 typed 枚举(`cancelled|dispatched`),不用字符串;事件不带 prompt—— + prompt 是入队检查的产物,只留在宿主的队列条目里。 +- 决策(放行):Thread Manager 在回合收口之后原子地「取队首 → 登记占用 → `turn.started` + `queue.removed{dispatched}`」, + 再落盘用户条目、下发用户条目、起整轮。**放行不重跑检查、不存在放行失败**,放行之后的一切失败都是回合失败, + 走既有 `turn.completed.failure`,不新增通道。kick 点 = 回合任务收尾(含 drop 守卫盖 panic)+ 中止路径 + 入队之后, + 幂等且在临界区里认领队首;入队不取 `DirectTaonierActiveInvocationGuard`(它必须整轮持有,是这一轮的调用身份)。 +- 决策(前端):不再持有队列副本,chip 只由事件投影,入队失败只由命令返回值给反馈(保留提示文案与"不丢草稿"); + 埋点句柄改由宿主在放行时开。 +- 决策(顺带退役):`--direct-codex-chat` 整个退役——它唯一的实际消费者是手工夹具 + `scripts/direct-execution-production-fixture.mjs`(PR #439 引入、不在 CI、无 npm/harness 注册、无测试钉它), + 夹具一并退役;`TurnAlreadyRunning` 与前端"同一轮消息仍在处理中"分支随之删除。 + `DirectTaonierActiveInvocationGuard` 的身份与 Thread Manager 的 `active_turn.turn_id` 是同一件事的两份记录, + CLI 退役后让那五个读者改读 Thread Manager,再在第二步删掉守卫与它的 60 秒卡死兜底 + (删前必须保住"回合卡死可被取消解开"这条由 `d833ca9d3` 事故换来的保证)。 +- 代价:入队即写盘(引用 UI 设计文档的条目会生成 `ui/generated-.js`,取消不撤);检查是时间点事实、 + 放行不重跑(manifest / 权限 / 目录变化后照旧放行,偏差落到回合失败);队列随进程消失,不做跨进程持久化; + 失去"真实二进制驱动执行层生产验证"这条手工路径。 +- 验证方式(设计稿,尚未实施):`docs/adr/【ADR】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md`、 + `docs/technical/【实施计划】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md`。 diff --git a/docs/technical/【实施计划】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md b/docs/technical/【实施计划】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md new file mode 100644 index 000000000..150ef125a --- /dev/null +++ b/docs/technical/【实施计划】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md @@ -0,0 +1,144 @@ +# DirectProject 命令入队化与待发消息队列归宿主实施计划 + +更新时间:`2026-09-24` + +状态:**待实施**(设计已定稿,尚未开工) + +设计口径见 [`【ADR】DirectProject命令入队化与待发消息队列归宿主-2026-09-24`](../adr/【ADR】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md)。 +本文件只排实施顺序、不变式与验收,不重复设计理由。 + +## 第 0 步:词表切换(与代码同批,不单独提交) + +「接单 / 拒单」退役,改成「入队 / 入队失败 / 放行」。逐处按角色改,**不做字面替换**: + +| 旧写法 | 新写法 | +| --- | --- | +| 命令的成功 / 失败(接单 / 拒单) | 入队 / 入队失败 | +| 这一轮真正成立的那一刻(接单) | 放行 | +| 接单前的检查 | 入队时的检查 | +| 接单成立 ⇔ 事件流里有开始有结束 | **放行成立** ⇔ 事件流里有开始有结束 | +| 接单后的失败都是回合失败 | **放行后**的失败都是回合失败 | + +三类出现点必须分开处理(全仓 384 处): + +1. DirectProject 域(文档、注释、标识符、测试名):按上表改。主要落点 `docs/adr/【ADR】DirectProject命令接单化-2026-09-23.md`、 + `docs/technical/【实施计划】DirectProject命令接单化-2026-09-23.md`、 + `docs/technical/【技术方案】DirectProject Codex原始历史与异常恢复-2026-09-04.md`、`CONTEXT.md`、 + `agent/direct_runtime/user_input.rs`、`agent/direct_turn_accept.rs`、`agent/direct_turn_error.rs`、`agent/direct_thread_manager.rs`、 + `chat/controller/useDirectProjectChatController.ts`、`chat/conversation/directTurnPresentation.ts`、`tests/appSurface/chat-composer.suite.ts`。 +2. `features/agent-runtime` 域的「拒单文案」(`model.ts`、`tests/agentRuntimeModel.test.ts`):那里没有队列, + 改成「请求被拒 / 拒绝」,不要写成「入队失败」。 +3. 假阳性:`单测` 这类词不动(如 `server-rs/crates/api-server/src/editor_project.rs`)。 + +标识符同批改(映射表): + +| 旧 | 新 | +| --- | --- | +| `chat_with_game_creator_direct_codex` | `enqueue_direct_codex_turn` | +| `DirectTurnRejection`(ts-rs 导出) | `DirectTurnEnqueueFailure` | +| 前端 `readDirectTurnRejection` / `directTurnRejectionNotice*` | `readDirectTurnEnqueueFailure` / `directTurnEnqueueFailureNotice*` | +| Rust `direct_turn_rejection` | `direct_turn_enqueue_failure` | +| `DirectTurnReservation::accept` | `DirectTurnReservation::start` | +| `accept_direct_thread_turn` / `DirectThreadManager::accept_turn` | `start_direct_thread_turn` / `start_turn` | +| `agent/direct_turn_accept.rs` | `agent/direct_turn_dispatch.rs` | +| `DirectTurnError::TurnAlreadyRunning` | 删除(第 4 步判据确认零调用方后) | + +两份已接受的 ADR 保留正文与文件名,只在顶部加一行词表注记,避免正文里的旧词变成假命题。 + +验收:`rg -n "接单|拒单"` 在 DirectProject 域与 `agent-runtime` 域均为 0;生成绑定重跑(`cargo test export_bindings`)后 +`git diff` 只剩映射表内的改动。 + +## 第 1 步:命令 = 入队(Rust) + +改动点: + +- `agent/direct_runtime/user_input.rs`:把命令主体拆成两半。 + **入队半**:`clientTurnId` 校验 → 工作流恢复 → 用户条目校验 → prompt 投影 → 前置条件 → 工程准备 → 入队; + 任何一步失败返回 typed 入队失败。**放行半**(第 2 步)从占用登记起。 +- 队列条目(宿主侧产物,只在内存):`PendingDirectTurn { client_turn_id, user_item: Value, prompt: String, creation_type: Option, at: u64 }`。 + `prompt` 与 `creation_type` 是入队检查的产物,放行不再重算;事件里**不带** `prompt`。 +- `agent/direct_thread_manager.rs`: + - `MAX_PENDING_DIRECT_TURNS = 5` 落在 Rust,只数在队条目;满队 → typed 入队失败。 + - 队列的成员与顺序**就是事件列表本身**:宿主侧产物挂在对应的 `queue.enqueued` 事件上(`StoredEvent` 增加一个非序列化的 + 可选产物字段),不另建队列表。 + - `observe_event`:`queue.enqueued` 在队期间**不可回收**;`queue.removed` 可回收,并把同 `clientTurnId` 的 enqueued 标记为可回收。 + `is_bootstrap_event` 不改——live-set bootstrap 因此自动把当前队列交给新订阅者。 + - 入队幂等:判重范围「在队 ∪ 正在跑的那一轮」,重复入队返回同一次成功。 + - `remove_direct_project_pending_turn(project_path, client_turn_id)`:typed 结果枚举 `Removed | AlreadyDispatched | NotFound`, + 在临界区里判「是否仍未被认领」。 +- 线上形状:`DirectThreadEvent::QueueEnqueued { client_turn_id, user_item, creation_type?, at }`(`queue.enqueued`)、 + `QueueRemoved { client_turn_id, reason: DirectQueueRemovalReason }`(`queue.removed`), + `DirectQueueRemovalReason` 是 ts-rs 导出的 typed 枚举 `cancelled | dispatched`。 + +不变式:入队不登记占用、不落盘、不发回合事件、不起 codex;入队失败不写用户条目、不产生事件。 + +## 第 2 步:放行与 kick + +- `start_direct_thread_turn`:同一个临界区里「取队首 → 占用登记 → `turn.started` + `queue.removed{dispatched}`」; + 之后落盘用户条目 → 下发用户条目 → spawn 整轮(顺序与今天的接单后半段一致,`accept` 必须早于落盘与 `turn/start`)。 +- `kick_direct_queue_dispatch(thread_id)`:幂等;在临界区里判「无占用 + 队首存在 + 未被认领」,认领后 spawn 放行任务。 + 调用点三个:回合任务收尾(正常 / 失败共用)、中止路径、入队之后。 +- panic 兜底:回合任务里的一个 drop 守卫负责踢一脚,保证任务 panic 或 future 被丢弃时队列不会永久停住。 +- 入队不取 `DirectTaonierActiveInvocationGuard`;该守卫继续由整轮持有。 +- 不变式:放行不重跑检查、没有放行失败;放行之后的一切失败都走 `turn.completed.failure`。 + +验收(Rust 单测):放行原子性(`queue.removed{dispatched}` 与 `turn.started` 同批、无中间窗口); +kick 幂等(并发两次只认领一次);队首在放行后被移除、remove 对已放行条目返回 `AlreadyDispatched`; +回合失败 / 中止后队列继续放行下一条;任务 panic 后队列仍能继续;多订阅者游标各自独立时 bootstrap 仍重建完整队列。 + +## 第 3 步:前端收口 + +退役: + +- `chatComposerQueue.ts` 的 `enqueueChatTurn` / `dequeueChatTurn` / `removeQueuedChatTurn` / `isChatTurnQueueFull` / + `MAX_QUEUED_CHAT_TURNS`(提示文案 `chatQueueFullNotice` 保留,改由 Rust 的 typed 入队失败驱动)。 +- controller 的 `queuedTurns` / `queuedTurnsRef` / `queueSequenceRef` / `completionPendingRef` / `handledCompletedTurnCountRef` / + `busyBaselineTurnCountRef` / `dispatchNextQueuedTurn` / `beginTurnBusy` / `endTurnBusy` / `turnBusyRef` 与排序用的计数 effect。 +- 排队条目那部分 `pendingRunAnalyticsRef` 与 `beginDirectRunAnalytics` 的调用时序(埋点句柄改由宿主在放行时开)。 + +保留与改写: + +- chip 由运行态事件的投影驱动(新增 pending 列表投影与两条队列事件的 reducer 分支); + chip 文案仍用现成的派生(`directCodexContentToPromptText` + `resourceLabelResolver`),只是输入换成事件里的 `userItem`。 +- 取消 chip 改调 `remove_direct_project_pending_turn`;入队失败只由命令返回值驱动提示(保留"不丢草稿"行为)。 +- 忙态 = 事件投影 + 「队列非空」指示;`directProjectTurnStatus` 的"命令在飞"分支收成 IPC 在飞。 +- 写权限门 `ensureConversationWriteAllowed` 留在入队之前(确认框必须在用户在场时弹)。 + +验收:`tests/directThreadChat.test.ts` 补队列事件投影用例;`tests/appSurface/chat-composer.suite.ts` 的排队 / 取消 / 满队 / +放行三组用例改成新语义;`npm --workspace apps/ai-game-creator-shell run typecheck` 通过。 + +## 第 4 步:CLI 与夹具退役 + +- `cli.rs`:删 `CliCommand::DirectCodexChat` 变体、`project_path_mut` 分支(`cli.rs:181`)与派发分支(`cli.rs:899-935`)。 +- `agent/direct_runtime/mod.rs`:删 `run_direct_game_creator_turn_at` 与 `run_direct_game_creator_turn_at_with_creation_type` + (各自只有彼此与 CLI 一个调用方)。 +- 删 `apps/ai-game-creator-shell/scripts/direct-execution-production-fixture.mjs`(CLI 的唯一消费者)。 +- 删 `DirectTurnError::TurnAlreadyRunning`(两个生产点都没了)与前端"同一轮消息仍在处理中"文案、专属分支、 + `tests/appSurface/project-conversation.suite.ts` 的对应断言。 +- 文档同步:09-22 里程碑把 `--direct-codex-chat` 从"范围外(保留)"改成退役项;两份 Direct 技术方案的 CLI 承诺删掉; + 09-23 ADR 与 `decision-log.md` 里"CLI 保持 await"的口径改掉。 + +验收:`rg -n -- "--direct-codex-chat"` 与 `rg -n "run_direct_game_creator_turn_at"` 零命中; +`cargo check --tests` 无新增 `dead_code` 告警;`check-config.mjs` 与 `.gitea/workflows/project-ci.yml` 不受影响(已核实无引用)。 + +## 第 5 步:守卫清理(CLI 退役之后) + +- 五个身份读者(`agent/direct_execution.rs`、`agent/direct_tool_bridge.rs`、`agent/direct_runtime/mod.rs` 的付费美术重生成、 + `agent/direct_validation.rs`、`agent/direct_project_context.rs`)改读 Thread Manager 的活动回合身份 + (新增一个按项目路径取 `active_turn.turn_id` 的只读入口)。 +- 删 `DirectTaonierActiveInvocationGuard` 与 `release_stale_direct_taonier_active_invocation` 及其测试 + (`direct_runtime/mod.rs` 的守卫用例、`direct_project_context.rs`、`direct_tool_bridge.rs`、`direct_tools_mcp.rs`、 + `user_input.rs`、`codex_app_server/mod.rs` 的集成用例),取消路径改为直接走无条件的 `complete_direct_thread_turn("aborted")`。 +- 前置判据:补一条测试证明"回合任务被 park / 泄漏时,取消仍能清空占用并让队列继续放行"——那次 60 秒兜底来自真实事故 + (`d833ca9d3`),删它必须有等价保证。 + +## 验收与证据 + +- Rust:第 1、2、5 步各自的单测;`cargo test` 定向 + `cargo check --tests` 无新增告警。 +- Node:`npx vitest run tests/appSurface.test.ts`、`npm --workspace apps/ai-game-creator-shell run typecheck`。 +- 端到端(`chat-composer.suite.ts`):回合运行中入队两条 → 取消一条 → 终态后只放行剩下那条; + 入队只发生一次 IPC、没有第二次发送命令;满队提示;入队失败时草稿不丢。 +- 手工:两个窗口看同一项目(A 排队 B 可见可取消);离开工作台再回来队列仍在并继续放行; + `kill -9` 后重进队列消失(与 ADR 的已知边界一致)。 +- 全仓:`npm run check:encoding`、`npm run check:doc-index`、`git diff --check`。 +- 不涉及 SpacetimeDB schema,不需要 `npm run check:spacetime-schema`。