From ecac9d8e15e5dd94766a8e569048baeb52bcbe76 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=8E=8B=E5=BE=B7=E5=AE=87?= Date: Wed, 30 Sep 2026 10:39:32 +0800 Subject: [PATCH] =?UTF-8?q?DirectProject=20=E8=B0=83=E7=94=A8=E8=BA=AB?= =?UTF-8?q?=E4=BB=BD=E5=BD=92=20Thread=20Manager=EF=BC=9A=E5=88=A0?= =?UTF-8?q?=E6=8E=89=E7=AC=AC=E4=BA=8C=E4=BB=BD=E8=BF=9B=E7=A8=8B=E5=86=85?= =?UTF-8?q?=E5=AE=88=E5=8D=AB=E8=A1=A8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 新增 Thread Manager 的只读身份入口:DirectTurnIdentity + read_direct_turn_identity,占用仍只由终态出口解除 - 五个身份读者(执行会话、工具桥、校验预约、上下文预取、付费美术重生成)改读 direct_active_turn_id_at - 删 DirectTaonierActiveInvocationGuard / DIRECT_TAONIER_ACTIVE_INVOCATIONS / DirectActiveTurnView 与只读探测 - 删 release_stale_direct_taonier_active_invocation,改为 direct_stale_turn_for_release 只做前置校验(身份一致 + 启动窗口),释放仍由 complete_direct_thread_turn + kick_direct_queue_dispatch 完成 - DirectStaleTurnReleaseReason 取代 DirectTaonierStaleGuardReason;终止兜底的注释与文案改成占用口径 - 放行不再另取调用身份:认领时已在 Thread Manager 登记,删掉认领后的守卫进入与失败分支 - 新增前置判据用例 stale_cancel_releases_the_occupancy_and_dispatches_the_next_pending_turn,并补背景拨时测试钩子 - 集成用例改用 DirectTurnReservation::accept_for_test 建占用,不再进守卫 --- .../src/agent/codex_app_server/mod.rs | 45 +-- .../src-tauri/src/agent/direct_execution.rs | 2 +- .../src/agent/direct_project_context.rs | 9 +- .../src-tauri/src/agent/direct_runtime/mod.rs | 381 +++++------------- .../src/agent/direct_thread_manager.rs | 54 +++ .../src-tauri/src/agent/direct_tool_bridge.rs | 9 +- .../src-tauri/src/agent/direct_tools_mcp.rs | 10 +- .../src/agent/direct_turn_dispatch.rs | 64 ++- .../src-tauri/src/agent/direct_validation.rs | 4 +- 9 files changed, 229 insertions(+), 349 deletions(-) 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 4538e1c57..4af07737d 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 @@ -4428,7 +4428,7 @@ fn register_active_direct_codex_turn( /// 已向正在跑的回合发出中断:界面等这一轮自己的收尾复位。 pub(crate) const DIRECT_TURN_CANCEL_OUTCOME_INTERRUPTED: &str = "interrupted"; -/// 这一轮已经没有人替它收尾,本地守卫已被兜底释放:界面必须自己复位。 +/// 这一轮已经没有人替它收尾,占用已被兜底解除:界面必须自己复位。 pub(crate) const DIRECT_TURN_CANCEL_OUTCOME_RELEASED: &str = "released"; /// `cancel_direct_codex_turn` 的返回值:界面据此决定是自己复位,还是等回合自己收尾。 @@ -4444,12 +4444,12 @@ pub(crate) struct DirectTurnCancelView { pub(crate) client_turn_id: String, } -/// "终止"这一步要作用在哪:发中断,还是走残留守卫兜底释放。 +/// "终止"这一步要作用在哪:发中断,还是走残留回合的兜底释放。 enum DirectCodexTurnCancelTarget { /// app-server 侧还有活句柄:正常发 `turn/interrupt`。 Interrupt(Arc), /// app-server 侧已经拿不到可中断的活句柄;带上是哪种情况。 - Stale(DirectTaonierStaleGuardReason), + Stale(DirectStaleTurnReleaseReason), } /// 终止当前项目正在运行的 Direct 回合。 @@ -4459,13 +4459,13 @@ enum DirectCodexTurnCancelTarget { /// 不动任何既有事件或命令语义。 /// /// 兜底路径:app-server 侧已经拿不到可中断的活句柄时,说明这一轮不会再有人替它收尾。 -/// 只发中断会让本地守卫(`DirectTaonierActiveInvocationGuard`)永远留在进程内,用户此后 -/// 每条消息都会被"已有另一条回合正在运行"拒绝——这正是"重进会话被堵死"的死锁形态。 -/// 这时显式释放这条守卫并把可读原因返回给界面。释放条件见 -/// [`release_stale_direct_taonier_active_invocation`] 的注释;"正在跑的是另一轮"仍然 -/// 保持原拒绝语义,什么都不释放。 +/// 只发中断会让 Thread Manager 上这一轮的占用永远解不开,这个项目此后每条消息都会被 +/// "已有另一条回合正在跑"挡在放行之外——这正是"重进会话被堵死"的死锁形态。这时显式写一条 +/// `turn.completed{aborted}` 解除占用并把可读原因返回给界面。释放条件见 +/// [`direct_stale_turn_for_release`] 的注释;"正在跑的是另一轮"仍然保持原拒绝语义, +/// 什么都不释放。 /// -/// 兜底终态带 `userItemId`:身份取 `release_stale_direct_taonier_active_invocation` 返回的 +/// 兜底终态带 `userItemId`:身份取 [`direct_stale_turn_for_release`] 返回的 /// clientTurnId(客户端回合身份的唯一来源),与正常路径的开口条目 id 同一份 canonical 口径。 /// 拿不到 clientTurnId 就留空——这一轮不会再有原生终态,猜一个身份会让前端把边界盖到别人身上。 fn direct_stale_cancel_turn_completed_event(client_turn_id: &str) -> DirectThreadEvent { @@ -4502,10 +4502,10 @@ pub(crate) fn cancel_direct_codex_turn_at( DirectCodexTurnCancelTarget::Interrupt(Arc::clone(cancellation)) } Ok(_) => { - DirectCodexTurnCancelTarget::Stale(DirectTaonierStaleGuardReason::ExecutorExited) + DirectCodexTurnCancelTarget::Stale(DirectStaleTurnReleaseReason::ExecutorExited) } Err(_) => DirectCodexTurnCancelTarget::Stale( - DirectTaonierStaleGuardReason::NeverReachedExecutor, + DirectStaleTurnReleaseReason::NeverReachedExecutor, ), } }; @@ -4523,17 +4523,17 @@ pub(crate) fn cancel_direct_codex_turn_at( }) } DirectCodexTurnCancelTarget::Stale(reason) => { - let released = - release_stale_direct_taonier_active_invocation(root, client_turn_id, reason)?; + // 先做前置校验:确实登记着一轮、身份对得上、过了启动窗口才允许判成残留。 + let released = direct_stale_turn_for_release(root, client_turn_id, reason)?; // 这一轮不会再有人替它发终态事件(执行进程已退出 / 从没进执行器), - // 兜底补一条,否则前端的"最新回合是否在跑"会永远停在运行中。 - // 走 Thread Manager 的深层出口而不是裸 append:这是**为这一轮写的终态**,占用必须 - // 同时解除,否则这个 thread 会一直被认为是"还有没收口的回合",挡住后面的放行。 + // 兜底补一条,否则前端的"最新回合是否在跑"会永远停在运行中。走 Thread Manager 的 + // 深层出口而不是裸 append:这也是**解除占用的唯一出口**,解除不了这个 thread 会一直 + // 被认为是"还有没收口的回合",挡住后面的放行。 complete_direct_thread_turn( &direct_thread_id_for_project(root), direct_stale_cancel_turn_completed_event(&released), ); - // 这一轮不会再有人替它收尾(守卫刚被兜底释放),占用也在这里解除了:踢一脚让队列继续。 + // 这一轮不会再有人替它收尾(占用刚被兜底解除):踢一脚让队列继续。 // 少了这一脚,排在这条后面的待发消息要等到下一次用户动作才会被放行。 crate::agent::kick_direct_queue_dispatch(root); Ok(DirectTurnCancelView { @@ -8048,9 +8048,6 @@ done }); let thread_id = direct_thread_id_for_project(&project); let bootstrap = crate::agent::subscribe_direct_thread(&thread_id); - let _active_invocation = - crate::agent::DirectTaonierActiveInvocationGuard::enter(&project, "turn-0001") - .expect("enter direct invocation"); let _reservation = crate::agent::DirectTurnReservation::accept_for_test(&thread_id, "turn-0001"); crate::agent::append_direct_project_user_message_at(&project, &user_item) @@ -8207,12 +8204,8 @@ done bootstrap.events.is_empty(), "订阅发生在回合之前,bootstrap 必须为空" ); - // 生产入口(`chat_with_game_creator_direct_codex`)在发起回合前登记本客户端的 - // 付费生成身份;工具桥在这一轮里按它绑定付费调用,这里补上同一步。 - let _active_invocation = - crate::agent::DirectTaonierActiveInvocationGuard::enter(&project, "turn-0001") - .expect("enter direct invocation"); - // 生产入口(`enqueue_direct_codex_turn` + 放行半)在起 codex 之前先放行,再把用户条目落盘: + // 生产入口(`enqueue_direct_codex_turn` + 放行半)在起 codex 之前先放行(认领时登记调用身份), + // 再把用户条目落盘: // 这里补上同一步,于是这一轮的边界仍在同一个订阅里成对出现,历史里也有那条用户消息。 let _reservation = crate::agent::DirectTurnReservation::accept_for_test(&thread_id, "turn-0001"); diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_execution.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_execution.rs index 73a283b85..400202c54 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_execution.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_execution.rs @@ -432,7 +432,7 @@ pub(super) async fn begin( let root = root.to_path_buf(); let prompt_hash = hash(prompt.as_bytes()); let session = tokio::task::spawn_blocking(move || { - let turn = super::direct_taonier_active_invocation_id_at(&root)?; + let turn = super::direct_active_turn_id_at(&root)?; let host = crate::game_creator_runtime_config_dir() .ok_or("direct-execution-host: 需要客户端私有配置目录,CLI 请提供 --config-dir")?; open_with_analytics_at( diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_project_context.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_project_context.rs index 971d41931..772166b30 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_project_context.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_project_context.rs @@ -349,11 +349,7 @@ pub(super) async fn prefetch_turn_input( let scan_root = root.to_path_buf(); let expected_turn = client_turn_id.to_string(); let files = tokio::task::spawn_blocking(move || { - if direct_taonier_active_invocation_id_at(&scan_root) - .ok() - .as_deref() - != Some(expected_turn.as_str()) - { + if direct_active_turn_id_at(&scan_root).ok().as_deref() != Some(expected_turn.as_str()) { return Vec::new(); } let candidates = [ @@ -590,8 +586,7 @@ mod tests { #[tokio::test] async fn host_prefetch_keeps_data_out_of_system_rules_and_matches_active_turn() { let (_temp, root) = project(); - // 调用身份(预取闸门)与逻辑回合(上下文身份)是两件事,生产入口两步都做。 - let _guard = DirectTaonierActiveInvocationGuard::enter(&root, "prefetch-turn").unwrap(); + // 预取闸门与上下文身份读的是同一处:Thread Manager 上这一轮的活动回合登记。 let _turn = DirectTurnReservation::accept_for_test( &direct_thread_id_for_project(&root), "prefetch-turn", diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/mod.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/mod.rs index 3fa5066b0..83040048a 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/mod.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_runtime/mod.rs @@ -1,11 +1,9 @@ use super::*; use base64::Engine as _; use std::collections::BTreeMap; -use std::collections::HashMap; use std::io::Write; use std::path::{Path, PathBuf}; -use std::sync::{Mutex, OnceLock}; -use std::time::{SystemTime, UNIX_EPOCH}; +use std::sync::Mutex; mod user_input; pub(crate) use user_input::enqueue_direct_codex_turn; @@ -453,100 +451,16 @@ fn direct_taonier_regeneration_invocation_sha256(invocation_id: &str) -> String format!("{:x}", hasher.finalize()) } -#[derive(Debug)] -struct DirectTaonierActiveInvocation { - invocation_id: String, - started_at: u64, -} - -static DIRECT_TAONIER_ACTIVE_INVOCATIONS: OnceLock< - Mutex>, -> = OnceLock::new(); - -/// 一条"app-server 侧完全没有登记"的守卫,只有在存在时间超过这个量级后才允许被 -/// "终止"兜底释放。一轮 Direct 回合在进入 app-server 之前只做本地准备(读配置、 -/// 读 manifest、拼系统提示、开审计),是秒级的;超过这个窗口还没登记,说明这一轮 -/// 不可能再进入执行器,守卫是残留。 -const DIRECT_TAONIER_STALE_GUARD_MIN_AGE_MS: u64 = 60_000; - -fn direct_taonier_active_now_millis() -> u64 { - unix_millis().min(u128::from(u64::MAX)) as u64 -} - -#[derive(Debug)] -pub(crate) struct DirectTaonierActiveInvocationGuard { - root: PathBuf, - invocation_id: String, -} - -impl DirectTaonierActiveInvocationGuard { - pub(crate) fn enter(root: &Path, invocation_id: &str) -> Result { - let root = root - .canonicalize() - .map_err(|error| DirectTurnError::ProjectRootUnanchored { - cause: error.to_string(), - })?; - let active = DIRECT_TAONIER_ACTIVE_INVOCATIONS.get_or_init(|| Mutex::new(HashMap::new())); - let mut active = active - .lock() - .map_err(|_| DirectTurnError::HostStateUnavailable { - detail: "Direct 调用身份锁已损坏".to_string(), - })?; - match active.get(&root) { - Some(existing) => { - return Err(DirectTurnError::TurnAlreadyRunning { - existing_invocation_id: existing.invocation_id.clone(), - incoming_invocation_id: invocation_id.to_string(), - }); - } - None => { - let started_at = SystemTime::now() - .duration_since(UNIX_EPOCH) - .map(|duration| duration.as_millis() as u64) - .unwrap_or_default(); - active.insert( - root.clone(), - DirectTaonierActiveInvocation { - invocation_id: invocation_id.to_string(), - started_at, - }, - ); - } - } - Ok(Self { - root, - invocation_id: invocation_id.to_string(), - }) - } -} - -impl Drop for DirectTaonierActiveInvocationGuard { - fn drop(&mut self) { - let Some(active) = DIRECT_TAONIER_ACTIVE_INVOCATIONS.get() else { - return; - }; - let Ok(mut active) = active.lock() else { - return; - }; - let remove = active - .get(&self.root) - .is_some_and(|existing| existing.invocation_id == self.invocation_id); - if remove { - active.remove(&self.root); - } - } -} - -pub(crate) fn direct_taonier_active_invocation_id_at(root: &Path) -> Result { - let root = root - .canonicalize() - .map_err(|error| format!("无法锚定 Direct 调用项目目录:{error}"))?; - DIRECT_TAONIER_ACTIVE_INVOCATIONS - .get_or_init(|| Mutex::new(HashMap::new())) - .lock() - .map_err(|_| "Direct 调用身份锁已损坏".to_string())? - .get(&root) - .map(|active| active.invocation_id.clone()) +/// 这个项目上正在跑的那一轮身份(`clientTurnId`),取自 **Thread Manager 的活动回合**。 +/// +/// 它是"这一轮是谁 / 在不在跑"的唯一登记:首页「运行中的项目」快照、放行闸门(有未收口回合就不认领 +/// 队首)读的都是这一处,不再另立第二份守卫表。 +/// +/// 拿不到身份时返回带前缀的错误:调用方(付费美术执行会话、工具桥、校验预约、上下文预取)都靠 +/// "拿不到身份"来拒绝"在没有回合的情况下创建 / 清理生成账本"。 +pub(crate) fn direct_active_turn_id_at(root: &Path) -> Result { + read_direct_turn_identity(&direct_thread_id_for_project(root)) + .map(|identity| identity.client_turn_id) .ok_or_else(|| { format!( "{DIRECT_TAONIER_RESULT_UNKNOWN_PREFIX} 当前付费美术调用缺少稳定 clientTurnId,已拒绝创建或清理生成账本" @@ -554,49 +468,26 @@ pub(crate) fn direct_taonier_active_invocation_id_at(root: &Path) -> Result Result, String> { - let root = root - .canonicalize() - .map_err(|error| format!("无法锚定 Direct 调用项目目录:{error}"))?; - Ok(DIRECT_TAONIER_ACTIVE_INVOCATIONS - .get_or_init(|| Mutex::new(HashMap::new())) - .lock() - .map_err(|_| "Direct 调用身份锁已损坏".to_string())? - .get(&root) - .map(|active| DirectActiveTurnView { - client_turn_id: active.invocation_id.clone(), - started_at: active.started_at, - })) +/// 上限用来挡"刚放行、还在本地准备(读配置 / manifest / 拼系统提示 / 开审计)的回合被终止误伤": +/// 这期间执行器侧还没有这一轮的登记,但这一轮马上就会进去,提前释放等于放开并发。 +const DIRECT_STALE_TURN_RELEASE_MIN_AGE_MS: u64 = 60_000; + +fn direct_taonier_active_now_millis() -> u64 { + unix_millis().min(u128::from(u64::MAX)) as u64 } -/// "终止"拿不到可中断句柄时的分类,决定是否允许强制释放本地守卫。 +/// "终止"拿不到可中断句柄时的分类,决定这一轮能不能被兜底释放。 #[derive(Clone, Copy, Debug, Eq, PartialEq)] -pub(crate) enum DirectTaonierStaleGuardReason { +pub(crate) enum DirectStaleTurnReleaseReason { /// app-server 侧登记着这一轮,但执行进程已经退出:这一轮不可能再有收尾。 ExecutorExited, /// app-server 侧完全没有这一轮的登记:只有过了正常启动窗口才允许释放。 NeverReachedExecutor, } -impl DirectTaonierStaleGuardReason { +impl DirectStaleTurnReleaseReason { pub(crate) fn message(self) -> &'static str { match self { Self::ExecutorExited => "陶泥儿执行进程已退出", @@ -605,53 +496,42 @@ impl DirectTaonierStaleGuardReason { } } -/// 强制释放某项目登记的 Direct 活跃回合占用("终止"的兜底出口)。 +/// 兜底释放前的前置校验:这一轮是不是**可以**被判成残留,返回它的 `clientTurnId`。 /// -/// 释放条件(四条必须同时成立,这段注释就是契约): -/// 1. 项目路径能 canonicalize,且守卫表里确实登记了这一轮; +/// 条件(三条必须同时成立,这段注释就是契约): +/// 1. 这个 thread 上确实登记着一轮(Thread Manager 的活动回合); /// 2. 传了 `expected_client_turn_id` 时必须与登记一致——绝不误伤另一条回合; -/// 3. 调用方已确认 app-server 侧没有可中断的活句柄,即 `reason` 成立; -/// 4. `reason == NeverReachedExecutor` 时,这条登记的年龄必须超过 -/// [`DIRECT_TAONIER_STALE_GUARD_MIN_AGE_MS`],排除"刚进入、还在本地准备阶段" -/// 的正常启动窗口——那种情况下这一轮马上就会去执行器,释放等于放开并发。 +/// 3. `reason == NeverReachedExecutor` 时,登记的年龄必须超过 +/// [`DIRECT_STALE_TURN_RELEASE_MIN_AGE_MS`],排除"刚放行、还在本地准备"的正常启动窗口。 /// -/// 移除后原守卫的 `Drop` 变成空操作(`invocation_id` 已不在表里),所以释放是幂等的; -/// 释放只影响"能否开始新回合",不动任何正在跑的回合事件。 -pub(crate) fn release_stale_direct_taonier_active_invocation( +/// **释放本身不在这里**:调用方拿到身份后往 Thread Manager 写一条 `turn.completed{aborted}`, +/// 那才是解除占用的唯一出口(也因此天然幂等——占用已经没了就返回"没有正在运行的回合")。 +pub(crate) fn direct_stale_turn_for_release( root: &Path, expected_client_turn_id: Option<&str>, - reason: DirectTaonierStaleGuardReason, + reason: DirectStaleTurnReleaseReason, ) -> Result { - let root = root - .canonicalize() - .map_err(|error| format!("无法锚定 Direct 调用项目目录:{error}"))?; - let mut active = DIRECT_TAONIER_ACTIVE_INVOCATIONS - .get_or_init(|| Mutex::new(HashMap::new())) - .lock() - .map_err(|_| "Direct 调用身份锁已损坏,无法释放".to_string())?; - let Some(existing) = active.get(&root) else { - return Err("当前项目没有正在运行的陶泥儿回合,无法终止".to_string()); - }; + let thread_id = direct_thread_id_for_project(root); + let active = read_direct_turn_identity(&thread_id) + .ok_or_else(|| "当前项目没有正在运行的陶泥儿回合,无法终止".to_string())?; if let Some(expected) = expected_client_turn_id .map(str::trim) .filter(|value| !value.is_empty()) { - if existing.invocation_id != expected { + if active.client_turn_id != expected { return Err("正在运行的是另一条 Direct 客户端回合,已拒绝终止".to_string()); } } - if reason == DirectTaonierStaleGuardReason::NeverReachedExecutor { - let age_ms = direct_taonier_active_now_millis().saturating_sub(existing.started_at); - if age_ms < DIRECT_TAONIER_STALE_GUARD_MIN_AGE_MS { + if reason == DirectStaleTurnReleaseReason::NeverReachedExecutor { + let age_ms = direct_taonier_active_now_millis().saturating_sub(active.started_at); + if age_ms < DIRECT_STALE_TURN_RELEASE_MIN_AGE_MS { return Err(format!( "这一轮 Direct 客户端回合刚开始 {} 秒、还在准备中,暂不能强制释放;请稍后再试", age_ms / 1000 )); } } - let released = existing.invocation_id.clone(); - active.remove(&root); - Ok(released) + Ok(active.client_turn_id) } fn direct_taonier_regeneration_project_id(root: &Path) -> Result { @@ -3262,7 +3142,7 @@ pub(crate) async fn ensure_direct_taonier_art_package_at( )); } let mut regeneration_workflow = if mode.regenerates_existing() { - let invocation_id = direct_taonier_active_invocation_id_at(root)?; + let invocation_id = direct_active_turn_id_at(root)?; Some(prepare_direct_taonier_regeneration_workflow_at( root, prompt, @@ -5639,170 +5519,93 @@ mod tests { assert!(normalize_direct_client_turn_id(None).is_err()); } + /// "这一轮是谁"只有 Thread Manager 一处真相;读取是只读的,不改占用。 #[test] - fn active_client_turn_id_is_exclusive_until_the_outer_turn_finishes() { - let root = tempfile::tempdir().expect("active invocation root"); - let first = DirectTaonierActiveInvocationGuard::enter(root.path(), "client-turn-0001") - .expect("first client turn"); - let duplicate = DirectTaonierActiveInvocationGuard::enter(root.path(), "client-turn-0001") - .expect_err("same stable turn is already running"); - assert!( - matches!( - &duplicate, - DirectTurnError::TurnAlreadyRunning { - existing_invocation_id, - incoming_invocation_id, - } if existing_invocation_id == "client-turn-0001" - && incoming_invocation_id == "client-turn-0001" - ), - "{duplicate:?}" - ); - // 界面按 typed 变体的两个身份字段分流,不解析文案:同一条身份 != 另一条身份。 - assert_eq!( - duplicate.to_string(), - "同一轮消息仍在处理中,已拒绝并发复用同一 clientTurnId;请等它结束或点「终止」后再发送" - ); - let different = DirectTaonierActiveInvocationGuard::enter(root.path(), "client-turn-0002") - .expect_err("different turn cannot take over the project"); - assert!( - different - .to_string() - .contains("已有另一条 Direct 客户端回合正在运行"), - "{different}" - ); - drop(first); - DirectTaonierActiveInvocationGuard::enter(root.path(), "client-turn-0001") - .expect("lost-response replay after the original turn finishes"); - } + fn reading_the_active_turn_identity_does_not_take_over_the_occupancy() { + let root = tempfile::tempdir().expect("active turn root"); + let thread_id = direct_thread_id_for_project(root.path()); + assert_eq!(read_direct_turn_identity(&thread_id), None); - #[test] - fn read_direct_active_turn_reports_the_registered_turn_and_disappears_after_drop() { - let root = tempfile::tempdir().expect("active invocation root"); - assert_eq!( - read_direct_taonier_active_invocation_at(root.path()).expect("read idle project"), - None - ); - - let first = DirectTaonierActiveInvocationGuard::enter(root.path(), "client-turn-read-1") - .expect("first client turn"); - let running = read_direct_taonier_active_invocation_at(root.path()) - .expect("read running project") - .expect("running turn is visible to the read-only probe"); + let reservation = DirectTurnReservation::accept_for_test(&thread_id, "client-turn-read-1"); + let running = read_direct_turn_identity(&thread_id).expect("running turn is readable"); assert_eq!(running.client_turn_id, "client-turn-read-1"); assert!(running.started_at > 0, "{running:?}"); - // camelCase 契约:前端按 `clientTurnId` / `startedAt` 取值。 assert_eq!( - serde_json::to_value(&running).expect("serialize view"), - serde_json::json!({ - "clientTurnId": "client-turn-read-1", - "startedAt": running.started_at, - }) + direct_active_turn_id_at(root.path()).expect("identity for the running turn"), + "client-turn-read-1" ); - // 只读探测不占有、不释放:探测之后同项目第二次进入仍然被拒。 - let duplicate = - DirectTaonierActiveInvocationGuard::enter(root.path(), "client-turn-read-2") - .expect_err("read-only probe must not take over the project"); - assert!(duplicate - .to_string() - .contains("已有另一条 Direct 客户端回合正在运行")); + // 只读探测不释放占用:认领仍然被这一轮挡着。 + assert!(claim_direct_pending_turn(&thread_id).is_none()); - drop(first); - assert_eq!( - read_direct_taonier_active_invocation_at(root.path()).expect("read idle project"), - None - ); + drop(reservation); + assert_eq!(read_direct_turn_identity(&thread_id), None); + assert!(direct_active_turn_id_at(root.path()).is_err()); } + /// 兜底释放的前置判据:身份对得上,且"从没进执行器"必须过启动窗口。 #[test] - fn stale_guard_release_requires_a_matching_identity_and_only_after_the_start_window() { - let root = tempfile::tempdir().expect("active invocation root"); - let first = DirectTaonierActiveInvocationGuard::enter(root.path(), "client-turn-stale-1") - .expect("first client turn"); + fn stale_turn_release_requires_a_matching_identity_and_only_after_the_start_window() { + let root = tempfile::tempdir().expect("stale turn root"); + let thread_id = direct_thread_id_for_project(root.path()); + let reservation = DirectTurnReservation::accept_for_test(&thread_id, "client-turn-stale-1"); // ① clientTurnId 不匹配:拒绝,且不误伤正在跑的那一轮。 - let mismatch = release_stale_direct_taonier_active_invocation( + let mismatch = direct_stale_turn_for_release( root.path(), Some("client-turn-stale-2"), - DirectTaonierStaleGuardReason::ExecutorExited, + DirectStaleTurnReleaseReason::ExecutorExited, ) .expect_err("another turn must not be released"); assert!(mismatch.contains("另一条"), "{mismatch}"); - assert!( - DirectTaonierActiveInvocationGuard::enter(root.path(), "client-turn-stale-2").is_err() + assert_eq!( + read_direct_turn_identity(&thread_id).map(|turn| turn.client_turn_id), + Some("client-turn-stale-1".to_string()) ); - // ② 刚登记、还没进执行器:正常启动窗口内不许释放(释放等于放开并发)。 - let young = release_stale_direct_taonier_active_invocation( + // ② 刚放行、还没进执行器:正常启动窗口内不许释放(释放等于放开并发)。 + let young = direct_stale_turn_for_release( root.path(), Some("client-turn-stale-1"), - DirectTaonierStaleGuardReason::NeverReachedExecutor, + DirectStaleTurnReleaseReason::NeverReachedExecutor, ) - .expect_err("a freshly registered turn is still starting"); + .expect_err("a freshly dispatched turn is still starting"); assert!(young.contains("暂不能强制释放"), "{young}"); - // ③ 同一条登记老过窗口:判定为残留守卫,释放后同项目可以再次进入。 - backdate_active_direct_invocation(root.path(), DIRECT_TAONIER_STALE_GUARD_MIN_AGE_MS + 1); - let released = release_stale_direct_taonier_active_invocation( + // ③ 老过窗口:判成残留,调用方拿到身份(解除占用由调用方写终态完成)。 + backdate_direct_active_turn_for_test(&thread_id, DIRECT_STALE_TURN_RELEASE_MIN_AGE_MS + 1); + let released = direct_stale_turn_for_release( root.path(), Some("client-turn-stale-1"), - DirectTaonierStaleGuardReason::NeverReachedExecutor, + DirectStaleTurnReleaseReason::NeverReachedExecutor, ) - .expect("stale guard is released"); + .expect("stale turn is releasable"); assert_eq!(released, "client-turn-stale-1"); - let second = DirectTaonierActiveInvocationGuard::enter(root.path(), "client-turn-stale-2") - .expect("a new turn can start once the stale guard is released"); - // 释放是幂等的:原 guard 的 Drop 不会影响后来登记的那一轮。 - drop(first); - let still_running = read_direct_taonier_active_invocation_at(root.path()) - .expect("read running project") - .expect("the newer turn survives the stale guard drop"); - assert_eq!(still_running.client_turn_id, "client-turn-stale-2"); - drop(second); + // ④ "执行进程已退出"不看启动窗口:刚放行也允许释放。 + assert_eq!( + direct_stale_turn_for_release( + root.path(), + None, + DirectStaleTurnReleaseReason::ExecutorExited, + ) + .expect("executor exited ignores the start window"), + "client-turn-stale-1" + ); + drop(reservation); } + /// 没有任何活动回合时,兜底释放给出可读原因,而不是静默成功。 #[test] - fn stale_guard_release_after_the_executor_exited_frees_the_project() { - let root = tempfile::tempdir().expect("active invocation root"); - let first = DirectTaonierActiveInvocationGuard::enter(root.path(), "client-turn-exited-1") - .expect("client turn"); - let released = release_stale_direct_taonier_active_invocation( + fn stale_turn_release_without_an_active_turn_says_so() { + let root = tempfile::tempdir().expect("empty turn root"); + let nothing = direct_stale_turn_for_release( root.path(), None, - DirectTaonierStaleGuardReason::ExecutorExited, - ) - .expect("executor exited: this guard is residue"); - assert_eq!(released, "client-turn-exited-1"); - assert_eq!( - read_direct_taonier_active_invocation_at(root.path()).expect("read idle project"), - None - ); - let _second = - DirectTaonierActiveInvocationGuard::enter(root.path(), "client-turn-exited-2") - .expect("a new turn can start after the residue is released"); - drop(first); - - // 没有任何登记时给出可读原因,而不是静默成功。 - let empty = tempfile::tempdir().expect("empty invocation root"); - let nothing = release_stale_direct_taonier_active_invocation( - empty.path(), - None, - DirectTaonierStaleGuardReason::ExecutorExited, + DirectStaleTurnReleaseReason::ExecutorExited, ) .expect_err("nothing to release"); assert!(nothing.contains("没有正在运行"), "{nothing}"); } - /// 把某项目当前登记的活跃回合往前拨 `age_ms`,用于覆盖"守卫年龄"分支。 - fn backdate_active_direct_invocation(root: &Path, age_ms: u64) { - let root = root.canonicalize().expect("canonical root"); - let mut active = DIRECT_TAONIER_ACTIVE_INVOCATIONS - .get_or_init(|| Mutex::new(HashMap::new())) - .lock() - .expect("active invocation lock"); - let entry = active.get_mut(&root).expect("registered invocation"); - entry.started_at = entry.started_at.saturating_sub(age_ms); - } - #[test] fn direct_success_reply_is_persisted_once_with_the_stable_client_turn_identity() { let root = tempfile::tempdir().expect("temp dir"); @@ -6780,9 +6583,10 @@ mod tests { .expect("completed replay authority is the stable outer turn, not a resampled brief"); assert_eq!(completed_replay, persisted); - let active = - DirectTaonierActiveInvocationGuard::enter(root.path(), "client-turn-completed-1") - .expect("restore the same client invocation after response loss"); + let _active_turn = DirectTurnReservation::accept_for_test( + &direct_thread_id_for_project(root.path()), + "client-turn-completed-1", + ); let replay = ensure_direct_taonier_art_package_at( root.path(), "第一套美术但偷偷换 brief", @@ -6799,7 +6603,6 @@ mod tests { .state, DirectTaonierRegenerationWorkflowState::Completed ); - drop(active); } #[test] 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 f3b844034..ea7dfb10d 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 @@ -83,6 +83,19 @@ pub(crate) struct DirectActiveTurnSnapshot { pub(crate) sequence: u64, } +/// 一个 thread 上正在跑的那一轮的**只读身份**:给"这一轮是谁"的读者用(付费美术执行会话、 +/// 工具桥、校验预约、上下文预取)。 +/// +/// 身份取自 Thread Manager 的活动回合登记——"这一轮在不在跑"只有这一处真相,界面快照 +/// ([`list_direct_active_turns`])与放行闸门读的也是它,不再另立第二份记录。 +#[derive(Clone, Debug, Eq, PartialEq)] +pub(crate) struct DirectTurnIdentity { + /// 与事件 `turnId` 同一个身份(`clientTurnId`)。 + pub(crate) client_turn_id: String, + /// 这一轮被放行(登记占用)的时刻(Unix 毫秒)。 + pub(crate) started_at: u64, +} + /// 一次放行的结果:队首那条待发消息,以及这次放行的占用身份。 /// /// `token` 是这一轮的占用身份(终态出口只认它),`pending` 是放行要用的**全部**输入——放行不再 @@ -314,6 +327,19 @@ impl DirectThreadManager { true } + /// 这个 thread 上正在跑的那一轮身份;`None` = 没有未收口的回合。 + fn active_turn_identity(&self, thread_id: &str) -> Option { + self.threads.get(thread_id).and_then(|thread| { + thread + .active_turn + .as_ref() + .map(|active| DirectTurnIdentity { + client_turn_id: active.turn_id.clone(), + started_at: active.started_at, + }) + }) + } + fn turn_is_active(&self, thread_id: &str) -> bool { self.threads .get(thread_id) @@ -836,6 +862,16 @@ pub(crate) fn complete_direct_thread_turn_if_reserved( written } +/// 这个 thread 上正在跑的那一轮身份。 +/// +/// 只读:占用仍然只由终态出口(深层终态 / 占用对象兜底 / 兜底释放)解除,这里不改任何状态。 +pub(crate) fn read_direct_turn_identity(thread_id: &str) -> Option { + global_direct_thread_manager() + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .active_turn_identity(thread_id) +} + /// 这个 thread 是否还有没收口的逻辑回合。 pub(crate) fn direct_thread_turn_is_active(thread_id: &str) -> bool { global_direct_thread_manager() @@ -844,6 +880,24 @@ pub(crate) fn direct_thread_turn_is_active(thread_id: &str) -> bool { .turn_is_active(thread_id) } +/// 测试用:把某 thread 上正在跑那一轮的登记时刻往前拨 `age_ms`。 +/// +/// 只服务"启动窗口闸门"这一条判据:终止的兜底释放只对"登记已存在很久"的残留放行, +/// 真实用例等不了 60 秒。 +#[cfg(test)] +pub(crate) fn backdate_direct_active_turn_for_test(thread_id: &str, age_ms: u64) { + let mut manager = global_direct_thread_manager() + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + if let Some(active) = manager + .threads + .get_mut(thread_id) + .and_then(|thread| thread.active_turn.as_mut()) + { + active.started_at = active.started_at.saturating_sub(age_ms); + } +} + fn notify_direct_thread_subscribers(thread_id: &str) { let subscriber_ids = { let manager = global_direct_thread_manager() diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_tool_bridge.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_tool_bridge.rs index 93b0d5741..18ebdac42 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_tool_bridge.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_tool_bridge.rs @@ -267,7 +267,7 @@ impl DirectToolBridge { impl DirectToolBridgeState { fn begin_user_turn(self: &Arc) -> Result { - let turn_id = direct_taonier_active_invocation_id_at(&self.root)?; + let turn_id = direct_active_turn_id_at(&self.root)?; let mut authorization = self .turn_authorization .lock() @@ -4665,9 +4665,10 @@ mod tests { let state = direct_tool_bridge_state(root.path().to_path_buf()); assert!(state.begin_user_turn().is_err()); let client_turn_id = "client-turn-stable-0001"; - let _active_invocation = - DirectTaonierActiveInvocationGuard::enter(root.path(), client_turn_id) - .expect("client-owned stable invocation"); + let _active_invocation = DirectTurnReservation::accept_for_test( + &direct_thread_id_for_project(root.path()), + client_turn_id, + ); let active_turn = state .begin_user_turn() .expect("client turn authorization state"); diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_tools_mcp.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_tools_mcp.rs index 28c506608..e89980666 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_tools_mcp.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_tools_mcp.rs @@ -3838,19 +3838,21 @@ mod tests { (handle, receiver) } - /// 真实项目 + 真实工具桥 + 真实 Direct 回合身份。 + /// 真实项目 + 真实工具桥 + 真实 Direct 回合占用(这一轮的调用身份)。 async fn tool_chain_start( root: &Path, ) -> ( super::super::direct_tool_bridge::DirectToolBridge, - DirectTaonierActiveInvocationGuard, + DirectTurnReservation, super::super::direct_tool_bridge::DirectExecutionTestFixture, ) { let bridge = super::super::direct_tool_bridge::start_direct_tool_bridge(root, false) .await .expect("start tool chain bridge"); - let turn = DirectTaonierActiveInvocationGuard::enter(root, "tool-chain-turn") - .expect("arm tool chain direct turn"); + let turn = DirectTurnReservation::accept_for_test( + &direct_thread_id_for_project(root), + "tool-chain-turn", + ); let execution = super::super::direct_tool_bridge::direct_execution_fixture(root, "tool-chain-turn") .await; diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_turn_dispatch.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_turn_dispatch.rs index 689298c53..734c5b852 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_turn_dispatch.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_turn_dispatch.rs @@ -2,7 +2,10 @@ //! //! 这个模块两件事,别再往里加第三件: //! 1. [`DirectTurnReservation`]:这一轮的占用对象——持有它代表还没收口,终态出口只认它的 token; -//! 2. [`kick_direct_queue_dispatch`]:认领队首 → 起整轮(取调用身份 → 落盘用户条目 → 下发 → 跑回合)。 +//! 2. [`kick_direct_queue_dispatch`]:认领队首 → 起整轮(落盘用户条目 → 下发 → 跑回合)。 +//! +//! 这一轮的**调用身份**(付费美术、执行会话、MCP、校验、上下文预取读的那个"这一轮是谁")由认领 +//! 在 Thread Manager 上登记,认领成功即成立;这里不再另取一次、也不再持第二份身份。 //! //! 放行**不重跑任何入队检查,也没有"放行失败"**:入队时的检查已经判过,这一轮一旦放行,之后的 //! 一切失败都是**回合失败**,由占用对象收口成 `turn.completed` 的失败载荷。设计见 @@ -14,8 +17,8 @@ use super::{ append_direct_project_user_message_at, claim_direct_pending_turn, complete_direct_thread_turn_if_reserved, direct_thread_id_for_project, direct_tool_call_now_ms, redact_agent_runtime_error, run_direct_game_creator_turn_at_with_creation_type_and_emitter, - DirectDispatchedTurn, DirectGameCreatorTurnUpdateEmitter, DirectTaonierActiveInvocationGuard, - DirectTurnError, DirectTurnTerminal, PendingDirectTurn, + DirectDispatchedTurn, DirectGameCreatorTurnUpdateEmitter, DirectTurnError, DirectTurnTerminal, + PendingDirectTurn, }; /// 一次放行的占用。持有它就代表这一轮还没收口。 @@ -88,7 +91,7 @@ impl Drop for DirectTurnReservation { /// /// 幂等,三个调用点:**入队之后**(队列空且没有回合在跑时,放行就是这一脚,不必再等一个调度周期)、 /// **占用释放**(正常 / 失败 / panic 共用,见 [`DirectTurnReservation`] 的 `Drop`)、**中止路径** -/// (守卫被兜底释放时那一轮不会再有人替它收尾)。 +/// (兜底释放时那一轮不会再有人替它收尾)。 /// /// 它什么都不返回:没有"放行失败"。踢不动就是现在不该跑——已经有回合在跑,或者队列是空的; /// 这两种情况都会在下一次释放时再踢。 @@ -114,15 +117,6 @@ async fn run_dispatched_direct_turn( reservation: DirectTurnReservation, ) { let turn_id = dispatched.pending.client_turn_id.clone(); - // 调用身份整轮持有:付费美术、执行会话、MCP、校验、上下文预取都读它。这里拿不到身份说明这个 - // 项目上还有另一条调用没放(残留守卫 / 别处的入口),按回合失败说清原因;队列照常继续。 - let active_invocation = match DirectTaonierActiveInvocationGuard::enter(&root, &turn_id) { - Ok(guard) => guard, - Err(failure) => { - reservation.finish_if_unfinished(DirectTurnTerminal::failed(&root, &failure)); - return; - } - }; // 落盘即回合成立:放行之后必须留下这条用户消息,哪怕这一轮随后失败。 if let Err(error) = append_direct_project_user_message_at(&root, &dispatched.pending.canonical_user_item) @@ -168,9 +162,6 @@ async fn run_dispatched_direct_turn( reservation.finish_if_unfinished(DirectTurnTerminal::failed(&root, &error)); } } - // 调用身份先放,再让占用对象的 `Drop` 收尾并踢下一脚:顺序反过来就会出现"守卫还握着、队首却已经 - // 可以放行"的窗口——那一脚会被自己的守卫挡成一条假的回合失败。 - drop(active_invocation); } #[cfg(test)] @@ -354,4 +345,45 @@ mod tests { assert!(direct_thread_turn_is_active(&busy)); drop(reservation); } + + /// 前置判据:删掉调用身份守卫之后,"回合任务泄漏、app-server 侧没有可中断句柄"的兜底释放 + /// 仍要能解开占用并把队首放行出去——`d833ca9d3` 的"重进会话被堵死"不得复活。 + #[test] + fn stale_cancel_releases_the_occupancy_and_dispatches_the_next_pending_turn() { + let thread = unique_thread("stale-cancel"); + let subscription = watch(&thread); + // 这一轮已经被放行(占用登记在 Thread Manager 上),但任务泄漏:拿不到可中断句柄。 + let leaked = DirectTurnReservation::accept_for_test(&thread, "turn-stale-1"); + enqueue(&thread, "turn-stale-2"); + + // "从没进执行器"只对过了启动窗口的残留放行;真实用例等不了 60 秒。 + crate::agent::backdate_direct_active_turn_for_test(&thread, 61_000); + let view = + crate::agent::cancel_direct_codex_turn_at(Path::new(&thread), Some("turn-stale-1")) + .expect("stale release"); + assert_eq!( + view.outcome, + super::super::codex_app_server::DIRECT_TURN_CANCEL_OUTCOME_RELEASED + ); + assert_eq!(view.client_turn_id, "turn-stale-1"); + + // 占用解开:第一条被兜底收口,第二条接上——终止这一脚就是"继续放行"。 + let events = pending(&subscription); + assert!( + events.iter().any(|event| matches!( + event, + DirectThreadEvent::TurnCompleted { status, .. } if status == "aborted" + )), + "{events:?}" + ); + assert!( + events.iter().any(|event| matches!( + event, + DirectThreadEvent::TurnStarted { user_item_id, .. } + if user_item_id.as_deref() == Some("direct-codex:turn-stale-2:user") + )), + "队首必须在兜底释放之后被放行:{events:?}" + ); + drop(leaked); + } } diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_validation.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_validation.rs index 33b849435..2cffc419b 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_validation.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_validation.rs @@ -368,10 +368,10 @@ async fn start_for_current_turn( key: String, build: bool, ) -> Result { - let turn_id = direct_taonier_active_invocation_id_at(root)?; + let turn_id = direct_active_turn_id_at(root)?; let root = root.to_path_buf(); tokio::task::spawn_blocking(move || { - if direct_taonier_active_invocation_id_at(&root)? != turn_id { + if direct_active_turn_id_at(&root)? != turn_id { return Err("validation-turn-changed: 当前验证所属回合已结束".into()); } let session = super::direct_execution::current(&root)?;