DirectProject 命令入队化:放行接线、队列移除命令与词表切换

- 命令改名 enqueue_direct_codex_turn:跑完入队检查后只入队,返回结构化的入队失败载荷
- direct_turn_accept.rs 改为 direct_turn_dispatch.rs:放行占用对象改为接管 Thread Manager 认领好的那一轮
- 新增 kick_direct_queue_dispatch:入队、占用释放与中止路径共用同一个幂等放行点
- 放行之后取调用身份、落盘用户条目、下发条目、跑整轮,失败一律收口成回合失败
- 新增 remove_direct_project_pending_turn 命令,取消一条还没放行的待发消息
- 队列条目补 analytics_attempt_id:埋点身份跟着入队走,放行不重算
- DirectTurnRejection 改名 DirectTurnEnqueueFailure,新增 QueueFull 变体
- 退役 accept_turn / accept_direct_thread_turn,测试改用 accept_for_test 建占用
- 前端与生成绑定同步改名,DirectTurnRejection.ts 换成 DirectTurnEnqueueFailure.ts
This commit is contained in:
2026-09-24 21:36:22 +08:00
parent 053e0bd932
commit 132f96581a
34 changed files with 800 additions and 841 deletions
@@ -34,7 +34,7 @@ mod direct_thread_wire;
mod direct_tool_bridge;
mod direct_tool_calls;
mod direct_tools_mcp;
mod direct_turn_accept;
mod direct_turn_dispatch;
mod direct_turn_error;
mod direct_turn_failure;
mod direct_turn_stream;
@@ -71,7 +71,7 @@ pub(crate) use direct_thread_wire::*;
pub(crate) use direct_tool_bridge::*;
pub(crate) use direct_tool_calls::*;
pub(crate) use direct_tools_mcp::*;
pub(crate) use direct_turn_accept::*;
pub(crate) use direct_turn_dispatch::*;
pub(crate) use direct_turn_error::*;
pub(crate) use direct_turn_failure::*;
pub(crate) use direct_turn_stream::*;
@@ -831,10 +831,10 @@ fn direct_thread_visible_item(
direct_thread_event_item(root, item)
}
/// 下发本轮的开口用户条目:接单之后、起 codex 之前的第一条运行态条目。
/// 下发本轮的开口用户条目:放行之后、起 codex 之前的第一条运行态条目。
///
/// **顺序是这条通道的全部意义**:整轮里任何失败说明都靠"属于哪一轮"归位,而归属只认这一轮的
/// 开口用户条目。发点在 `turn/start` 之后时,"接单到 `turn/start` 之间"的失败(连不上
/// 开口用户条目。发点在 `turn/start` 之后时,"放行到 `turn/start` 之间"的失败(连不上
/// app-server、执行器未通过验收、历史注入失败)没有用户条目可以挂,说明会按位置落进**上一轮**
/// 的分区里:界面显示成"错误在用户消息上面",上一轮还顶替本轮显示耗时,本轮的用户气泡再自成
/// 一个 0.0 秒的假回合。
@@ -846,7 +846,7 @@ pub(crate) fn emit_direct_thread_user_item(
item: &serde_json::Value,
) -> Option<DirectThreadItem> {
let entry_item = direct_thread_event_item(root, item)?;
// 条目时间是落盘 / 观测时间。前端乐观用户气泡已删(ADR「DirectProject命令接单化」后续更新),
// 条目时间是落盘 / 观测时间。前端乐观用户气泡已删(ADR「DirectProject命令入队化」),
// 所以这就是界面显示这条用户消息的唯一时间口径:它晚于用户按下发送,但不再有第二份更早的时间。
let at = entry_item.at();
append_direct_thread_event(
@@ -3300,11 +3300,11 @@ impl CodexAppServerConnection {
.map(str::trim)
.filter(|turn_id| !turn_id.is_empty())
.map(str::to_string);
// 用户条目由 GUI 命令在**接单之后、起 codex 之前**落盘("落盘即接单"),这里不再重复写;
// 用户条目由放行半在**放行之后、起 codex 之前**落盘("落盘即回合成立"),这里不再重复写;
// 本函数只取它的身份,把它作为本回合的第一条运行态条目下发,前端就能用同一个 itemId 把
// "本地乐观气泡"和"历史里的同一条"合成一条。
// 没有条目但给了 `clientTurnId` 的调用方(不是 GUI 命令那条路)只拿到一份本地投影:
// 不落盘,因为落盘的时机属于接单动作,不属于这里。
// 不落盘,因为落盘的时机属于放行动作,不属于这里。
let mut direct_turn_user_item: Option<serde_json::Value> = None;
if self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject {
let current_prompt = direct_codex_current_user_prompt(&request).trim();
@@ -3569,8 +3569,8 @@ impl CodexAppServerConnection {
let direct_turn_user_item_id = direct_turn_user_item
.as_ref()
.and_then(direct_thread_item_identity);
// 逻辑回合的**边界**不在这里:开始事件由接单动作发出、兜底由接单占用对象持有
// (`direct_turn_accept.rs`)。本轮的用户条目也不在这里下发——发点在接单之后、
// 逻辑回合的**边界**不在这里:开始事件由认领动作发出、兜底由放行占用对象持有
// (`direct_turn_dispatch.rs`)。本轮的用户条目也不在这里下发——发点在放行之后、
// 起 codex 之前(`emit_direct_thread_user_item`),见那条注释。
//
// 下面这个毫秒钟与逻辑回合无关,只服务模型终态的**完成时刻**:上游 Turn 的
@@ -4528,11 +4528,14 @@ pub(crate) fn cancel_direct_codex_turn_at(
// 这一轮不会再有人替它发终态事件(执行进程已退出 / 从没进执行器),
// 兜底补一条,否则前端的"最新回合是否在跑"会永远停在运行中。
// 走 Thread Manager 的深层出口而不是裸 append:这是**为这一轮写的终态**,占用必须
// 同时解除,否则这个 thread 会一直被认为是"还有没收口的回合",挡住后面的接单。
// 同时解除,否则这个 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 {
outcome: DIRECT_TURN_CANCEL_OUTCOME_RELEASED.to_string(),
message: format!(
@@ -8048,12 +8051,8 @@ done
let _active_invocation =
crate::agent::DirectTaonierActiveInvocationGuard::enter(&project, "turn-0001")
.expect("enter direct invocation");
let _reservation = crate::agent::direct_turn_accept::DirectTurnReservation::accept(
&thread_id,
"turn-0001",
Some("direct-codex:turn-0001:user"),
)
.expect("accept logical turn");
let _reservation =
crate::agent::DirectTurnReservation::accept_for_test(&thread_id, "turn-0001");
crate::agent::append_direct_project_user_message_at(&project, &user_item)
.expect("persist opener user item");
// 生产入口在落盘成功、起 codex 之前就把本轮的用户条目下发(`emit_direct_thread_user_item`):
@@ -8213,14 +8212,10 @@ done
let _active_invocation =
crate::agent::DirectTaonierActiveInvocationGuard::enter(&project, "turn-0001")
.expect("enter direct invocation");
// 生产入口(`chat_with_game_creator_direct_codex`)在起 codex 之前先接单,再把用户条目落盘:
// 生产入口(`enqueue_direct_codex_turn` + 放行半)在起 codex 之前先放行,再把用户条目落盘:
// 这里补上同一步,于是这一轮的边界仍在同一个订阅里成对出现,历史里也有那条用户消息。
let _reservation = crate::agent::direct_turn_accept::DirectTurnReservation::accept(
&thread_id,
"turn-0001",
Some("direct-codex:turn-0001:user"),
)
.expect("accept logical turn");
let _reservation =
crate::agent::DirectTurnReservation::accept_for_test(&thread_id, "turn-0001");
crate::agent::append_direct_project_user_message_at(&project, &user_item)
.expect("persist opener user item");
// 生产入口在落盘成功、起 codex 之前就把本轮的用户条目下发(`emit_direct_thread_user_item`):
@@ -437,14 +437,12 @@ mod tests {
async fn a_replaced_active_turn_marks_the_batch_stale() {
let (_temp, root) = project();
std::fs::write(root.join("code.js"), "unchanged").unwrap();
// 身份来自逻辑回合(Thread Manager):接单才是"这一轮在跑"的唯一登记。
// 身份来自逻辑回合(Thread Manager):放行才是"这一轮在跑"的唯一登记。
let owner = Arc::new(std::sync::Mutex::new(Some(
DirectTurnReservation::accept(
DirectTurnReservation::accept_for_test(
&direct_thread_id_for_project(&root),
"turn-before",
None,
)
.unwrap(),
),
)));
let swap = Arc::clone(&owner);
let result = read_batch_with(
@@ -455,14 +453,10 @@ mod tests {
let result = read_file(r, f, b);
let mut guard = swap.lock().unwrap();
drop(guard.take());
*guard = Some(
DirectTurnReservation::accept(
&direct_thread_id_for_project(r),
"turn-after",
None,
)
.unwrap(),
);
*guard = Some(DirectTurnReservation::accept_for_test(
&direct_thread_id_for_project(r),
"turn-after",
));
result
},
)
@@ -598,12 +592,10 @@ mod tests {
let (_temp, root) = project();
// 调用身份(预取闸门)与逻辑回合(上下文身份)是两件事,生产入口两步都做。
let _guard = DirectTaonierActiveInvocationGuard::enter(&root, "prefetch-turn").unwrap();
let _turn = DirectTurnReservation::accept(
let _turn = DirectTurnReservation::accept_for_test(
&direct_thread_id_for_project(&root),
"prefetch-turn",
None,
)
.unwrap();
);
let data = prefetch_turn_input(&root, "prefetch-turn")
.await
.unwrap()
@@ -8,7 +8,7 @@ use std::sync::{Mutex, OnceLock};
use std::time::{SystemTime, UNIX_EPOCH};
mod user_input;
pub(crate) use user_input::chat_with_game_creator_direct_codex;
pub(crate) use user_input::enqueue_direct_codex_turn;
#[cfg(test)]
pub(crate) use user_input::normalize_direct_client_turn_id;
@@ -2142,13 +2142,13 @@ pub(crate) fn record_direct_codex_failure(
recovery_hint,
diagnostics_suffix,
);
// 失败也进错误上报池:命令接单化之后前端 catch 只剩"接单被拒",池不能只靠前端填。
// 失败也进错误上报池:命令入队化之后前端 catch 只剩"入队失败",池不能只靠前端填。
// 两条通道同时上报也不会变成两条——池按 fingerprint 合并同一份文案。
let _ = crate::error_report::report_agent_runtime_error("direct-codex", &public_text);
public_text
}
/// 命令边界的错误文本:可留痕的调用级拒绝在这里补一份运行错误诊断(文案里不带诊断引用),
/// 命令边界的错误文本:可留痕的调用级入队失败在这里补一份运行错误诊断(文案里不带诊断引用),
/// 其余只输出 [`DirectTurnError`] 的 `Display`。
///
/// 分层改成 typed 之前,这几条"宿主 / 环境事实"是在回合失败通道里被写进诊断的;分层之后它们不再
@@ -2159,23 +2159,23 @@ pub(crate) fn direct_turn_error_boundary_text(
client_turn_id: Option<&str>,
failure: DirectTurnError,
) -> String {
direct_turn_rejection(root, client_turn_id, failure).message
direct_turn_enqueue_failure(root, client_turn_id, failure).message
}
/// 拒单边界:可留痕的调用级拒绝(宿主 / 环境事实)在这里补一份运行错误诊断,然后连同**结构化
/// 入队失败边界:可留痕的调用级失败(宿主 / 环境事实)在这里补一份运行错误诊断,然后连同**结构化
/// 变体**一起交给前端;其余只输出 [`DirectTurnError`] 的 `Display`。
///
/// GUI 命令与 CLI 边界共用这一份判据(CLI 只要文本,走上面的 `..._text`),禁止在各自边界再写一套。
pub(crate) fn direct_turn_rejection(
pub(crate) fn direct_turn_enqueue_failure(
root: &Path,
client_turn_id: Option<&str>,
failure: DirectTurnError,
) -> DirectTurnRejection {
) -> DirectTurnEnqueueFailure {
if !failure.is_reportable() {
return DirectTurnRejection::new(failure);
return DirectTurnEnqueueFailure::new(failure);
}
let message = record_direct_codex_failure(root, &failure, client_turn_id);
DirectTurnRejection {
DirectTurnEnqueueFailure {
error: failure,
message,
}
@@ -4409,9 +4409,9 @@ pub(crate) fn build_direct_codex_system_prompt_with_creation_type(
.collect())
}
/// 接单前必须成立的前置条件:目录可用、读写权限、正文非空、创建类型合法。
/// 入队前必须成立的前置条件:目录可用、读写权限、正文非空、创建类型合法。
///
/// GUI 命令在接单前调用(不成立就是**拒单**),CLI 入口在起回合前调用。两处共用这一份判据,
/// GUI 命令在入队半调用(不成立就是**入队失败**),CLI 入口在起回合前调用。两处共用这一份判据,
/// 不要再各自复制一遍条件。
pub(crate) fn check_direct_turn_preconditions(
root: &Path,
@@ -4450,8 +4450,8 @@ pub(crate) async fn run_direct_game_creator_turn_at_with_creation_type(
prompt: &str,
creation_type: Option<&str>,
) -> Result<String, DirectTurnError> {
// CLI 入口没有"接单"这一步(它 await 整轮,要那段回复文本),前置条件在这里自己过一遍;
// GUI 命令在同一步骤之后才接单,两边共用这一份判据。
// CLI 入口没有"入队 / 放行"这两步(它 await 整轮,要那段回复文本),前置条件在这里自己过一遍;
// GUI 命令在入队半里过同一份判据。
check_direct_turn_preconditions(root, prompt, creation_type)?;
run_direct_game_creator_turn_at_with_creation_type_and_emitter(
root,
@@ -4465,7 +4465,7 @@ pub(crate) async fn run_direct_game_creator_turn_at_with_creation_type(
.await
}
async fn run_direct_game_creator_turn_at_with_creation_type_and_emitter(
pub(crate) async fn run_direct_game_creator_turn_at_with_creation_type_and_emitter(
root: &Path,
prompt: &str,
creation_type: Option<&str>,
@@ -4495,7 +4495,7 @@ async fn run_direct_game_creator_turn_at_with_creation_type_and_emitter(
{
Ok(reply) => Ok(reply),
// 走到这里的一切失败都是**回合失败**:判据已经从"错误种类"改成"发生位置"——前置条件
// 在接单前就查过,能到这条通道的只有接单之后的事(连接、配置、历史注入、`turn/start`
// 在入队前就查过,能到这条通道的只有放行之后的事(连接、配置、历史注入、`turn/start`
// 被拒、模型与交付)。所以不再有"直通调用方"的分支。
Err(failure) => {
let stage = failure.turn_failure_stage();
@@ -1,7 +1,7 @@
//! DirectProject 用户输入命令适配器。
//!
//! Tauri 只在这里接收前端 item,校验与 canonical→prompt 投影交给 user-item
//! 深模块,回合编排仍由父模块负责。
//! Tauri 只在这里接收前端 item,校验与 canonical→prompt 投影交给 user-item 深模块,队列交给
//! Thread Manager;放行在 `direct_turn_dispatch`,不在这里。
use super::*;
@@ -28,31 +28,33 @@ pub(crate) fn normalize_direct_client_turn_id(
Ok(client_turn_id.to_string())
}
/// DirectProject 聊天命令:**只接单**,不再 await 整轮。
/// DirectProject 聊天命令:**只入队**,不 await、也不起回合。
///
/// 边界文案仍只在这里生成一次(`Display`);但 `Err` 的含义收窄成**拒单**——接单成立之后的
/// 一切失败(连不上 app-server、配置 / 凭据未就绪、历史注入失败、`turn/start` 被拒、模型与
/// 交付失败)都由这一轮的占用对象收口成 `turn.completed` 带失败载荷,不再回到这条返回值上。
/// 入队半就是"这条消息能不能收下":`clientTurnId` 校验 → 工作流恢复 → 用户条目校验 → prompt
/// 投影 → 前置条件 → 工程准备,任何一步失败都是**入队失败**(typed 载荷)。全部通过之后只做一件事
/// ——把待发消息排进队列:不登记占用、不落盘、不发回合事件、不起 codex。
///
/// 于是"这一轮跑成什么"只有订阅事件一个来源:命令返回 `Ok` 只说明**接单成立**。可留痕的调用级
/// 拒绝(宿主 / 环境事实)仍在边界补一份运行错误诊断,返回串不带诊断引用。
/// 于是"这一轮跑成什么"仍然只有订阅事件一个来源:命令返回 `Ok` 只说明**入队成立**。真正的回合边界
/// (`turn.started` / `turn.completed`)由 Thread Manager 在**放行**时写出(见 `direct_turn_dispatch`),
/// 入队失败不写用户条目、不产生任何事件。可留痕的调用级失败(宿主 / 环境事实)仍在边界补一份运行
/// 错误诊断,返回串不带诊断引用。
///
/// 与 CLI 的分工:CLI 入口(`cli.rs` 的 `direct-codex.chat`)**保持 await**——它要把那段回复文本
/// 打到终端上,没有事件订阅可用;它复用同一份接单前检查与同一个命令主体,只是自己等整轮的返回值。
/// 两个入口共用 [`direct_turn_error_boundary_text`] / [`direct_turn_rejection`],不要再各写一套判据。
/// 打到终端上,没有事件订阅可用;两个入口共用 [`direct_turn_error_boundary_text`] /
/// [`direct_turn_enqueue_failure`],不要再各写一套判据。
///
/// 设计见 `docs/adr/【ADR】DirectProject命令接单化-2026-09-23.md`。
/// 设计见 `docs/adr/【ADR】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md`。
#[tauri::command]
pub(crate) async fn chat_with_game_creator_direct_codex(
pub(crate) async fn enqueue_direct_codex_turn(
project_path: String,
user_item: DirectCodexUserItem,
creation_type: Option<String>,
client_turn_id: Option<String>,
analytics_attempt_id: Option<String>,
) -> Result<(), DirectTurnRejection> {
) -> Result<(), DirectTurnEnqueueFailure> {
let root = Path::new(project_path.trim());
let boundary_turn_id = client_turn_id.clone();
chat_with_game_creator_direct_codex_typed(
enqueue_direct_codex_turn_typed(
root,
user_item,
creation_type,
@@ -60,17 +62,16 @@ pub(crate) async fn chat_with_game_creator_direct_codex(
analytics_attempt_id,
)
.await
.map_err(|failure| direct_turn_rejection(root, boundary_turn_id.as_deref(), failure))
.map_err(|failure| direct_turn_enqueue_failure(root, boundary_turn_id.as_deref(), failure))
}
/// 命令主体:全程 typed。顺序固定,**每一步失败都还是拒单**:
/// `clientTurnId` 校验 → 占用调用身份 → 工作流恢复 → 用户条目校验 → 前置条件 → 工程准备
/// → 接单 → 落盘用户条目 → 后台起整轮。
/// 入队半:全程 typed,顺序固定,**每一步失败都还是入队失败**:
/// `clientTurnId` 校验 → 工作流恢复 → 用户条目校验 → prompt 投影 → 前置条件 → 工程准备 → 入队。
///
/// 这个顺序不是风格问题:接单(`DirectTurnReservation::accept`)必须在所有"接单前就能判定"的
/// 检查之后,也必须早于用户条目落盘与 `turn/start`,否则并发拒单会晚于副作用、逻辑回合的开始
/// 事件会排在用户消息之后。
async fn chat_with_game_creator_direct_codex_typed(
/// 这个顺序不是风格问题:工程准备(可能建工程、铺脚手架)必须在"这条消息真的被收下"之前判定;
/// 没通过检查的东西不进队列。入队本身只做判重与容量,不再有第二种副作用——落盘、回合事件与整轮
/// 都在放行那一侧。
async fn enqueue_direct_codex_turn_typed(
root: &Path,
user_item: DirectCodexUserItem,
creation_type: Option<String>,
@@ -78,10 +79,6 @@ async fn chat_with_game_creator_direct_codex_typed(
analytics_attempt_id: Option<String>,
) -> Result<(), DirectTurnError> {
let turn_id = normalize_direct_client_turn_id(client_turn_id.as_deref())?;
// 占用调用身份:并发拒单要早于工程准备,避免两个请求同时改同一个项目。它只挡并发,**不是**
// 首页"运行中的项目"的来源(那张表由 Thread Manager 的逻辑回合导出),但仍必须与整轮同生
// 共死——随任务一起搬进后台。
let active_invocation = DirectTaonierActiveInvocationGuard::enter(root, &turn_id)?;
recover_direct_taonier_regeneration_workflow_at(root).map_err(|error| {
DirectTurnError::HostStateUnavailable {
detail: redact_agent_runtime_error(
@@ -96,10 +93,6 @@ async fn chat_with_game_creator_direct_codex_typed(
let user_prompt = direct_codex_user_item_to_prompt(root, &user_item)
.map_err(|detail| DirectTurnError::InputRejected { detail })?;
check_direct_turn_preconditions(root, &user_prompt, creation_type.as_deref())?;
let canonical_user_item =
serde_json::to_value(&user_item).map_err(|error| DirectTurnError::InputRejected {
detail: error.to_string(),
})?;
// 创建类型来自结构化用户入口;实际工程和可信脚手架由宿主复核。
crate::environment_check::prepare_new_web_project_at(root, creation_type.as_deref())
.await
@@ -107,96 +100,31 @@ async fn chat_with_game_creator_direct_codex_typed(
let detail = redact_agent_runtime_error(root, &error, 1800);
DirectTurnError::EnvironmentNotReady { detail }
})?;
// 接单:从这里开始这一轮就成立了。开始事件的身份由 `clientTurnId` 推导,**不读盘回填**
// ——开始事件发生在用户条目落盘之前,而落盘本身也可能失败。
let thread_id = direct_thread_id_for_project(root);
let user_item_id = direct_codex_user_item_id_for_client_turn_id(&turn_id);
let reservation = DirectTurnReservation::accept(&thread_id, &turn_id, user_item_id.as_deref())?;
// 落盘即接单:接单成功就必须在历史里留下这条用户消息,哪怕这一轮随后失败。
if let Err(error) = append_direct_project_user_message_at(root, &canonical_user_item) {
// 这一轮**已经接单**,所以收口只能走占用对象:写出失败终态(事件流里的那条失败说明就是
// 界面唯一一份解释),然后返回 `Ok`——命令的 `Err` 只表示**拒单**,回到那里会让同一个失败
// 同时从事件与横幅两条通道下发,也会让前端把"已经开始的回合"读成"没开始"。
// 不继续起整轮:历史是这条对话的单一事实源,用户消息没落盘时继续跑只会得到一条没有开口
// 用户消息的助手回复,而且失败会被静默掉。
let failure = DirectTurnError::EnvironmentNotReady {
detail: redact_agent_runtime_error(
root,
&format!("写入本项目对话历史失败:{error}"),
600,
),
};
reservation.finish_if_unfinished(DirectTurnTerminal::failed(root, &failure));
return Ok(());
}
// 用户条目落盘成功即下发:这一轮从"接单"到"起 codex"之间的一切失败(连不上
// app-server、执行器未通过验收、历史注入失败)都靠它把失败说明挂回自己那一轮;晚到
// `turn/start` 之后才发,这些失败就没有用户条目可挂,界面会把说明显示在用户消息上面。
crate::agent::codex_app_server::emit_direct_thread_user_item(root, &canonical_user_item);
let capture = crate::analytics::gui::capture_writer_context();
let root = root.to_path_buf();
tauri::async_runtime::spawn(async move {
run_accepted_direct_turn(
root,
turn_id,
user_prompt,
creation_type,
canonical_user_item,
capture,
analytics_attempt_id,
active_invocation,
reservation,
)
.await;
});
Ok(())
}
/// 接单之后的整轮:命令不再 await 它,它的收场只走事件流。
///
/// 三条收场路径都在这里收口:正常(深层的终态出口写 `turn.completed`)、失败(没有深层终态的
/// 早退由这里的占用对象补)、任务被丢弃 / panic(占用对象的 `Drop` 补 `host-dropped`)。
///
/// 两个守卫都**必须活到整轮结束**,所以随任务搬进来,不留在命令里:
/// `_active_invocation` 是这一轮的调用身份(并发拒单与首页在途回合都读它),`reservation`
/// 是逻辑回合的占用。
#[allow(clippy::too_many_arguments)]
async fn run_accepted_direct_turn(
root: std::path::PathBuf,
turn_id: String,
user_prompt: String,
creation_type: Option<String>,
canonical_user_item: serde_json::Value,
capture: Option<(
crate::analytics::contract::Context,
crate::analytics::store::AnalyticsWriter,
)>,
analytics_attempt_id: Option<String>,
_active_invocation: DirectTaonierActiveInvocationGuard,
reservation: DirectTurnReservation,
) {
let emitter = DirectGameCreatorTurnUpdateEmitter::new(&root, turn_id);
let outcome = run_direct_game_creator_turn_at_with_creation_type_and_emitter(
&root,
&user_prompt,
creation_type.as_deref(),
Some(&emitter),
Some(canonical_user_item),
capture,
analytics_attempt_id.as_deref(),
// 入队:到这里这一条已经过了全部检查,剩下的就是排队等放行。prompt 与 canonical 形状在这里
// 冻结——放行不重算,所以放行没有失败出口。
let pending = PendingDirectTurn::prepare(
turn_id,
user_item,
user_prompt,
creation_type,
analytics_attempt_id,
direct_tool_call_now_ms(),
)
.await;
match outcome {
Ok(reply) => {
// 深层的终态出口已经在 `run_turn` 里写出 `turn.completed`;这里只补最后一条回合更新。
emitter.emit("completed", Some("none"), Some(reply), None);
}
Err(error) => {
// 接单之后的失败一律是回合失败:失败诊断与失败说明已由上层写过,这里补终态事件。
// 深层已经写出终态时它不覆盖(同一轮只允许一条终态)。
reservation.finish_if_unfinished(DirectTurnTerminal::failed(&root, &error));
.map_err(|error| DirectTurnError::InputRejected {
detail: error.to_string(),
})?;
let thread_id = direct_thread_id_for_project(root);
match enqueue_direct_pending_turn(&thread_id, pending) {
Ok(_) => {}
Err(EnqueueRejection::QueueFull) => {
return Err(DirectTurnError::QueueFull {
limit: MAX_PENDING_DIRECT_TURNS,
})
}
}
// 入队之后立刻踢一脚:队列空且没有回合在跑时,放行就是这一脚,用户点发送不必再等一个调度周期。
kick_direct_queue_dispatch(root);
Ok(())
}
#[cfg(test)]
@@ -204,46 +132,64 @@ mod tests {
use super::*;
use crate::agent::{consume_direct_thread, subscribe_direct_thread, DirectThreadEvent};
/// 接单之后的早退也必须有终态。
/// 放行之后整轮是异步的:测试只能等事件。返回从订阅之后看到的全部事件。
async fn wait_for_completion(subscription_id: &str) -> Vec<DirectThreadEvent> {
let mut events = Vec::new();
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(30);
loop {
events.extend(
consume_direct_thread(subscription_id)
.expect("consume logical turn")
.events,
);
if events
.iter()
.any(|event| matches!(event, DirectThreadEvent::TurnCompleted { .. }))
{
return events;
}
assert!(
std::time::Instant::now() < deadline,
"回合没有收口:{events:?}"
);
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
}
fn user_item(text: &str) -> DirectCodexUserItem {
serde_json::from_value(serde_json::json!({
"type": "message",
"role": "user",
"id": "direct-codex:turn-1:user",
"content": [{ "type": "input_text", "text": text }],
}))
.expect("canonical user item")
}
/// 放行之后发现的失败也必须收口。
///
/// 这里用一个"目录存在但不是项目"的根制造一条**接单之后**才发现的失败(连 `run_turn` 的
/// 收尾都走不到)。命令此时早已返回 `Ok`,前端唯一的收口依据就是事件流,所以占用对象必须
/// 补出 `turn.completed`——这正是接单化要买的那条不变式。
/// 命令此刻早就返回了 `Ok`(它只回报入队成立),前端唯一的收口依据就是事件流:占用对象必须补出
/// `turn.completed`。这里用一个"目录存在但不是项目"的根制造一条放行之后才发现的失败。
#[tokio::test]
async fn a_failure_after_accept_still_closes_the_logical_turn() {
async fn a_failure_after_dispatch_still_closes_the_logical_turn() {
let temp = tempfile::tempdir().expect("temp dir");
let root = temp.path().join("not-a-project");
std::fs::create_dir_all(&root).expect("create project dir");
let thread_id = direct_thread_id_for_project(&root);
let subscription = subscribe_direct_thread(&thread_id);
let _ = consume_direct_thread(&subscription.subscription_id);
let reservation =
DirectTurnReservation::accept(&thread_id, "turn-1", Some("direct-codex:turn-1:user"))
.expect("accept logical turn");
let invocation =
DirectTaonierActiveInvocationGuard::enter(&root, "turn-1").expect("enter invocation");
run_accepted_direct_turn(
root.clone(),
"turn-1".to_string(),
"你好".to_string(),
enqueue_direct_codex_turn_typed(
&root,
user_item("你好"),
None,
serde_json::json!({
"type": "message",
"role": "user",
"id": "direct-codex:turn-1:user",
"content": [{ "type": "input_text", "text": "你好" }],
}),
Some("turn-1".to_string()),
None,
None,
invocation,
reservation,
)
.await;
.await
.expect("入队成立:命令只回报入队");
let events = consume_direct_thread(&subscription.subscription_id)
.expect("consume logical turn")
.events;
let events = wait_for_completion(&subscription.subscription_id).await;
let terminal = events
.iter()
.filter_map(|event| match event {
@@ -262,20 +208,22 @@ mod tests {
let failure = failure.as_ref().expect("失败终态必须带载荷");
assert!(
!failure.message.trim().is_empty(),
"接单之后的失败必须带上原因"
"放行之后的失败必须带上原因"
);
assert_eq!(user_item_id.as_deref(), Some("direct-codex:turn-1:user"));
// 占用已释放:队列继续走得动。
assert!(!crate::agent::direct_thread_turn_is_active(&thread_id));
}
/// 本轮的开口用户条目必须先于整轮里任何可能失败的东西下发。
///
/// 现场(用户可见的坏体验):命令接单、用户条目落盘之后,整轮在 `turn/start` 之前就失败
/// (连不上 app-server 一类)。这时如果用户条目还没下发,界面就只剩一条失败说明——它按位置
/// 落进**上一轮**的分区里,于是"错误显示在用户消息上面"、上一轮顶替本轮显示耗时,本轮的用户
/// 气泡再自成一个 0.0 秒的假回合。
/// 现场(用户可见的坏体验):放行、用户条目落盘之后,整轮在 `turn/start` 之前就失败(连不上
/// app-server 一类)。这时如果用户条目还没下发,界面就只剩一条失败说明——它按位置落进**上一轮**
/// 的分区里,于是"错误显示在用户消息上面"、上一轮顶替本轮显示耗时,本轮的用户气泡再自成一个
/// 0.0 秒的假回合。
///
/// 判据取事件流的前两条:命令体是顺序执行的,后台整轮是它之后才起的,所以"开始 → 用户条目"
/// 一定在最前面,之后才可能有失败终态。
/// 判据取事件流的前两条:放行任务的顺序是"开始 → 下发用户条目 → 跑整轮",所以终态一定在它们
/// 之后。
#[tokio::test]
async fn the_opening_user_item_is_emitted_before_anything_that_can_fail_in_the_turn() {
let temp = tempfile::tempdir().expect("temp dir");
@@ -285,48 +233,50 @@ mod tests {
let thread_id = direct_thread_id_for_project(&root);
let subscription = subscribe_direct_thread(&thread_id);
let _ = consume_direct_thread(&subscription.subscription_id);
let user_item: DirectCodexUserItem = serde_json::from_value(serde_json::json!({
"type": "message",
"role": "user",
"id": "direct-codex:turn-1:user",
"content": [{ "type": "input_text", "text": "hello" }],
}))
.expect("canonical user item");
chat_with_game_creator_direct_codex_typed(
enqueue_direct_codex_turn_typed(
&root,
user_item,
user_item("hello"),
None,
Some("turn-1".to_string()),
None,
)
.await
.expect("接单成立:命令只回报接单");
.expect("入队成立:命令只回报入队");
let events = consume_direct_thread(&subscription.subscription_id)
.expect("consume logical turn")
.events;
let events = wait_for_completion(&subscription.subscription_id).await;
// 队列事件先于放行:`queue.enqueued` → 放行(开始 + 离开队列)→ 开口用户条目 → 终态。
let started = events
.iter()
.position(|event| matches!(event, DirectThreadEvent::TurnStarted { .. }))
.expect("这一轮必须被放行");
assert!(
matches!(
events.first(),
events.get(started),
Some(DirectThreadEvent::TurnStarted { user_item_id, .. })
if user_item_id.as_deref() == Some("direct-codex:turn-1:user")
),
"第一条必须是带身份的回合开始:{events:?}"
"放行的第一条事件必须是带身份的回合开始:{events:?}"
);
// 开口用户条目必须是这一轮里的**第一条条目事件**:它之前只许有回合边界与队列事件。
let opening = events
.iter()
.position(|event| event.item_id().is_some())
.expect("本轮的开口用户条目必须下发");
assert!(
matches!(
events.get(1),
events.get(opening),
Some(DirectThreadEvent::ItemCompleted { item, .. })
if item.item_id() == "direct-codex:turn-1:user"
),
"第二条必须是本轮的开口用户条目:{events:?}"
"第一条条目事件必须是本轮的开口用户条目:{events:?}"
);
assert!(opening > started, "用户条目只能在放行之后:{events:?}");
let terminal = events
.iter()
.position(|event| matches!(event, DirectThreadEvent::TurnCompleted { .. }));
assert!(
terminal.is_none_or(|index| index > 1),
terminal.is_none_or(|index| index > opening),
"终态只能在用户条目之后:{events:?}"
);
// 落盘与下发同一份身份:历史里的条目 id 就是事件里的 itemId。
@@ -338,14 +288,15 @@ mod tests {
);
}
/// 接单之后的落盘失败:**只走占用对象的失败终态**,命令返回 `Ok`。
/// 放行之后的落盘失败:**只走占用对象的失败终态**,命令返回 `Ok`。
///
/// 这条路径的 `turn.started` 已经发过,命令再回一个 `Err` 就等于同一个失败下发两次(事件一条
/// 说明、横幅又一份),而且 `Err` 的含义是**拒单**——前端会把它读成"这一轮没开始"。历史追加写
/// 有一条测试注入(`.agent/runtime/test-fail-next-direct-project-history-append`),用它把这条
/// 路径钉成确定性:恰好一条失败终态、命令 `Ok`、占用释放(下一轮还能接单)。
/// 说明、横幅又一份),而且 `Err` 的含义是**入队失败**——前端会把它读成"这一轮没开始"。历史
/// 追加写有一条测试注入(`.agent/runtime/test-fail-next-direct-project-history-append`),用它把
/// 这条路径钉成确定性:恰好一条失败终态、命令 `Ok`、占用释放(下一轮还能继续)。
#[tokio::test]
async fn a_history_write_failure_after_accept_closes_the_turn_instead_of_rejecting() {
async fn a_history_write_failure_after_dispatch_closes_the_turn_instead_of_failing_the_command()
{
let temp = tempfile::tempdir().expect("temp dir");
let root = temp.path().join("direct-history-write-failure");
crate::init_local_game_project_at(&root, "direct-history-write", "落盘失败")
@@ -359,27 +310,18 @@ mod tests {
"9",
)
.expect("write history contention injection");
let user_item: DirectCodexUserItem = serde_json::from_value(serde_json::json!({
"type": "message",
"role": "user",
"id": "direct-codex:turn-1:user",
"content": [{ "type": "input_text", "text": "生成一个游戏" }],
}))
.expect("canonical user item");
chat_with_game_creator_direct_codex_typed(
enqueue_direct_codex_turn_typed(
&root,
user_item,
user_item("生成一个游戏"),
None,
Some("turn-1".to_string()),
None,
)
.await
.expect("接单之后的失败不再回到命令返回值:命令只回报接单成立");
.expect("入队成立:命令只回报入队");
let events = consume_direct_thread(&subscription.subscription_id)
.expect("consume logical turn")
.events;
let events = wait_for_completion(&subscription.subscription_id).await;
let terminals = events
.iter()
.filter_map(|event| match event {
@@ -398,20 +340,14 @@ mod tests {
"{}",
failure.message
);
// 这一轮已经接单,所以走的是**回合失败**:拒单那套 `direct-codex-failure:v2` 收口文案
// 不许出现在这里(它只属于可留痕的拒单)。
// 这一轮已经放行,所以走的是**回合失败**:入队失败那套 `direct-codex-failure:v2` 收口文案
// 不许出现在这里(它只属于可留痕的入队失败)。
assert!(
!failure.message.contains("direct-codex-failure"),
"{}",
failure.message
);
// 占用已释放:下一轮还能接单。
// 占用已释放:下一轮还能继续。
assert!(!crate::agent::direct_thread_turn_is_active(&thread_id));
assert!(DirectTurnReservation::accept(
&thread_id,
"turn-2",
Some("direct-codex:turn-2:user")
)
.is_ok());
}
}
@@ -43,14 +43,14 @@ struct SubscriberState {
cursor: u64,
}
/// 一条正在跑的逻辑回合的占用:接单时登记,终态写出时解除。
/// 一条正在跑的逻辑回合的占用:放行时登记,终态写出时解除。
///
/// 它同时是首页「运行中的项目」快照的**唯一事实源**([`list_direct_active_turns`]):这一格的
/// 生命周期就是"这一轮在不在跑",进度字段由运行时那一侧经 [`update_direct_thread_active_turn`]
/// 回填。任务侧不再另建一张活动回合表——同一件事只许有一处真相。
///
/// 两个身份别混:
/// - `token` 是这一次接单的占用身份:终态出口只有拿着同一个 token 的占用对象才能写兜底终态,
/// - `token` 是这一次放行的占用身份:终态出口只有拿着同一个 token 的占用对象才能写兜底终态,
/// 避免迟到的旧占用把新回合的边界顶掉。它不对外。
/// - `turn_id` 是给界面看的回合身份(`clientTurnId` 派生),只服务快照与进度回填的匹配。
#[derive(Clone, Debug)]
@@ -227,46 +227,6 @@ impl DirectThreadManager {
}
}
/// 接单:同一个临界区里拒绝并发、登记占用、追加逻辑回合开始事件。
///
/// 返回 `Err(existing_turn_id)` 表示这个 thread 已经有一条没收口的回合——此时不动队列,
/// 由调用方把它投影成接单拒绝。回的是**回合身份**(调用方接单时给的 `turn_id`)而不是占用
/// `token`:占用 token 只活在这个进程里,界面拿它匹配不了自己发出的那一轮,也没法判断
/// "撞的是同一轮还是另一轮"。
fn accept_turn(
&mut self,
thread_id: &str,
token: &str,
turn_id: &str,
user_item_id: Option<&str>,
started_at_ms: u64,
) -> Result<DirectThreadEvent, String> {
{
let thread = self.threads.entry(thread_id.to_string()).or_default();
if let Some(active) = thread.active_turn.as_ref() {
return Err(active.turn_id.clone());
}
thread.active_turn = Some(ActiveDirectTurn {
token: token.to_string(),
turn_id: turn_id.to_string(),
project_name: std::path::Path::new(thread_id)
.file_name()
.and_then(|name| name.to_str())
.map(str::to_string),
started_at: started_at_ms,
// 与"还没有任何进度事件"的状态一致:运行时给出的第一条进度会覆盖它。
status: "accepted".to_string(),
activity: Some("request-accepted".to_string()),
updated_at: started_at_ms,
sequence: 0,
});
}
Ok(self.append(
thread_id,
DirectThreadEvent::turn_started(started_at_ms).with_user_item_id(user_item_id),
))
}
/// 运行时回填这一轮的进度。只认"仍在跑 + 回合身份一致 + 序号不倒退"的那一次。
///
/// 返回是否真的写进去了:没有未收口的回合、身份对不上(上一轮迟到的进度)、序号倒退
@@ -449,10 +409,9 @@ impl DirectThreadManager {
/// 放行:同一个临界区里取队首 → 登记占用 → 写 `turn.started` 与 `queue.removed{dispatched}`。
///
/// 返回 `None` 表示现在**不该放行**:这个 thread 已经有未收口的回合,或者队列是空的。这两件事
/// 都由这里判,调用方(kick)不需要先探一遍再决定——探两次就是两个中间态。
///
/// 认领认的是队首,不是调用方点名的某一条:待发消息的顺序就是事件顺序,没有第三种来源。
/// 返回 `None` 表示现在**不该放行**:这个 thread 已经有未收口的回合,或者队列是空的——两件事
/// 都由这里判,调用方(kick)不必先探一遍再决定。认领认的永远是队首:待发消息的顺序就是事件
/// 顺序,没有第二种来源。
fn claim_pending_turn(
&mut self,
thread_id: &str,
@@ -773,27 +732,6 @@ pub(crate) fn append_direct_thread_event(
event
}
/// 接单:拒绝并发 + 登记占用 + 发逻辑回合开始事件(见 [`DirectThreadManager::accept_turn`])。
/// `Err` 是这一轮**已有的回合身份**(`turn_id`,也就是调用方的 `clientTurnId`),不是占用 token:
/// 调用方拿它投影成 `TurnAlreadyRunning` 的两个身份字段,界面按"撞的是同一轮还是另一轮"决定要
/// 不要动当前回合。
pub(crate) fn accept_direct_thread_turn(
thread_id: &str,
token: &str,
turn_id: &str,
user_item_id: Option<&str>,
started_at_ms: u64,
) -> Result<(), String> {
{
global_direct_thread_manager()
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.accept_turn(thread_id, token, turn_id, user_item_id, started_at_ms)?;
}
notify_direct_thread_subscribers(thread_id);
Ok(())
}
/// 入队一条待发消息(队列的成员与顺序就是事件列表,见 `direct_thread_queue`)。
pub(crate) fn enqueue_direct_pending_turn(
thread_id: &str,
@@ -830,7 +768,7 @@ pub(crate) fn remove_direct_pending_turn(
/// 认领队首并登记占用(放行的原子步骤,见 [`DirectThreadManager::claim_pending_turn`])。
///
/// `None` = 现在不该放行(已有未收口的回合,或队列为空)。
/// `None` = 现在不该放行:已有未收口的回合,或队列为空。
pub(crate) fn claim_direct_pending_turn(thread_id: &str) -> Option<DirectDispatchedTurn> {
let claimed = {
global_direct_thread_manager()
@@ -1281,7 +1219,7 @@ mod tests {
);
}
/// 首页快照就是逻辑回合的导出:接单即出现、进度按序号回填、收口即消失。
/// 首页快照就是逻辑回合的导出:放行即出现、进度按序号回填、收口即消失。
fn snapshot_of(
manager: &DirectThreadManager,
thread_id: &str,
@@ -1299,8 +1237,11 @@ mod tests {
assert!(snapshot_of(&manager, thread_id).is_none());
manager
.accept_turn(thread_id, "token-1", "turn-1", Some("u-1"), FIXED_AT_MS)
.expect("accept");
.enqueue_pending_turn(thread_id, pending_turn("turn-1"))
.expect("enqueue");
manager
.claim_pending_turn(thread_id, FIXED_AT_MS)
.expect("claim head");
let accepted = snapshot_of(&manager, thread_id).expect("accepted turn is visible");
assert_eq!(accepted.turn_id, "turn-1");
assert_eq!(accepted.project_name.as_deref(), Some("快照项目"));
@@ -1340,22 +1281,6 @@ mod tests {
assert!(snapshot_of(&manager, thread_id).is_none());
}
/// 并发接单回给调用方的是**回合身份**(`turn_id`),不是占用 `token`:token 只活在这个进程
/// 里,界面拿它匹配不了自己发出的那一轮。
#[test]
fn accept_conflict_returns_the_existing_turn_id() {
let mut manager = DirectThreadManager::with_limits(100, 100_000);
manager
.accept_turn("thread-1", "token-1", "turn-1", None, FIXED_AT_MS)
.expect("accept");
let conflict = manager
.accept_turn("thread-1", "token-2", "turn-2", None, FIXED_AT_MS)
.err();
assert_eq!(conflict.as_deref(), Some("turn-1"));
}
fn pending_turn(client_turn_id: &str) -> PendingDirectTurn {
PendingDirectTurn::prepare(
client_turn_id.to_string(),
@@ -1368,6 +1293,7 @@ mod tests {
.expect("canonical user item"),
format!("消息 {client_turn_id}"),
None,
None,
FIXED_AT_MS,
)
.expect("prepare pending turn")
@@ -35,6 +35,8 @@ pub(crate) struct PendingDirectTurn {
/// 入队检查产出的 prompt:放行不重算。
pub(crate) prompt: String,
pub(crate) creation_type: Option<String>,
/// 这一轮的运行埋点句柄:由前端在入队时生成,放行开跑时用同一份身份。
pub(crate) analytics_attempt_id: Option<String>,
/// 入队那一刻的宿主毫秒钟。
pub(crate) at: u64,
}
@@ -49,6 +51,7 @@ impl PendingDirectTurn {
user_item: DirectCodexUserItem,
prompt: String,
creation_type: Option<String>,
analytics_attempt_id: Option<String>,
at: u64,
) -> Result<Self, serde_json::Error> {
let canonical_user_item = serde_json::to_value(&user_item)?;
@@ -58,6 +61,7 @@ impl PendingDirectTurn {
canonical_user_item,
prompt,
creation_type,
analytics_attempt_id,
at,
})
}
@@ -124,6 +128,7 @@ mod tests {
user_item("生成一个游戏", "direct-codex:turn-1:user"),
"生成一个游戏".to_string(),
creation_type.map(str::to_string),
None,
1_700_000_000_123,
)
.expect("prepare pending turn")
@@ -215,7 +215,7 @@ impl DirectThreadRequestKind {
/// 一条待发消息离开队列的原因。
///
/// typed 枚举,取值即语义:取消是用户在输入盒上撤掉这条消息,放行是它已经接单并成为回合
/// typed 枚举,取值即语义:取消是用户在输入盒上撤掉这条消息,放行是它已经离开队列并成为回合
/// (同一临界区里另有 `turn.started`)。界面按它分流,不解析字符串。
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize, TS)]
#[serde(rename_all = "camelCase")]
@@ -223,7 +223,7 @@ impl DirectThreadRequestKind {
pub(crate) enum DirectQueueRemovalReason {
/// 用户取消了这条待发消息。
Cancelled,
/// 放行:这条待发消息已经接单并成为回合。
/// 放行:这条待发消息已经离开队列,成为正在跑的那一轮。
Dispatched,
}
@@ -293,7 +293,7 @@ impl DirectTurnFailure {
pub(crate) enum DirectThreadEvent {
#[serde(rename = "turn.started")]
TurnStarted {
/// 本轮开始的阶段时间(毫秒):**接单**那一刻的宿主毫秒钟(逻辑回合的起点,不是
/// 本轮开始的阶段时间(毫秒):**放行**那一刻的宿主毫秒钟(逻辑回合的起点,不是
/// `turn/start` 的时刻)。
#[serde(default, skip_serializing_if = "Option::is_none")]
#[ts(optional, as = "Option<f64>")]
@@ -1,282 +0,0 @@
//! DirectProject 的接单:把"一条用户消息被接单"变成 Thread Manager 里一对必然成对的逻辑回合事件。
//!
//! 这个模块只有一件事,别再往里加第二件:**接单成立的那一刻**在同一个临界区里拒绝并发、登记占用、
//! 发出逻辑回合开始事件;占用对象持有这一轮的终态出口——正常 / 失败 / 接单后的前置失败谁先写谁算,
//! 都没写时由 `Drop` 补一条 `host-dropped`。
//!
//! 为什么回合边界不能继续镜像 Codex 原生回合:`turn/start` 之前的失败(连不上 app-server、配置未
//! 就绪失败、历史注入失败)根本没有原生回合可以镜像,而它们同样是"这一轮已经成立"。设计见
//! `docs/adr/【ADR】DirectProject命令接单化-2026-09-23.md`。
use uuid::Uuid;
use super::{
accept_direct_thread_turn, complete_direct_thread_turn_if_reserved, direct_tool_call_now_ms,
DirectTurnError, DirectTurnTerminal,
};
/// 一次接单的占用。持有它就代表这一轮还没收口。
///
/// 生命周期由调用方决定:命令把整轮任务 spawn 出去时把它一起搬进任务,任务结束(正常或失败)
/// 时它随任务一起 drop。**持有顺序要与单飞锁一致**:单飞锁先声明、占用后声明,drop 时占用先收尾,
/// 新回合不可能插到中间。
pub(crate) struct DirectTurnReservation {
thread_id: String,
token: String,
user_item_id: Option<String>,
}
impl DirectTurnReservation {
/// 接单:登记占用并发出逻辑回合开始事件。
///
/// 失败表示这个 thread 已经有一条没收口的回合(并发接单),此时不改队列、不发事件。
/// `client_turn_id` 是给界面看的回合身份(首页快照与进度回填按它匹配),与占用身份 `token`
/// 是两件事:前者来自调用方,后者只活在这个进程里。
///
/// 拒单载荷里的两个身份都取**回合身份**:`existing` 是已在跑的那一轮的 `clientTurnId`
/// (Thread Manager 回的就是它),`incoming` 是本次请求的 `clientTurnId`。同一轮重发时两者
/// 相等,界面才走得到"同一轮消息仍在处理中"那条文案。
pub(crate) fn accept(
thread_id: &str,
client_turn_id: &str,
user_item_id: Option<&str>,
) -> Result<Self, DirectTurnError> {
let token = Uuid::new_v4().to_string();
accept_direct_thread_turn(
thread_id,
&token,
client_turn_id,
user_item_id,
direct_tool_call_now_ms(),
)
.map_err(|existing| DirectTurnError::TurnAlreadyRunning {
existing_invocation_id: existing,
incoming_invocation_id: client_turn_id.to_string(),
})?;
Ok(Self {
thread_id: thread_id.to_string(),
token,
user_item_id: user_item_id.map(str::to_string),
})
}
pub(crate) fn thread_id(&self) -> &str {
&self.thread_id
}
/// 接单之后还没走到深层终态就失败的收口口:只有这一轮仍被自己占用时才写。
///
/// 深层(真正跑完这一轮的代码)已经写出终态时返回 `false`,兜底不覆盖真实结果。
pub(crate) fn finish_if_unfinished(&self, terminal: DirectTurnTerminal) -> bool {
complete_direct_thread_turn_if_reserved(
&self.thread_id,
&self.token,
terminal.event(direct_tool_call_now_ms(), self.user_item_id.as_deref()),
)
}
}
impl Drop for DirectTurnReservation {
fn drop(&mut self) {
// 兜底:任务 panic、future 被丢弃、或今后在终态之前新增的 `?` 早退。
// 这类失败说不出原因,只给分类;能说清原因的错误必须由调用方在更早的地方显式收口。
let _ = self.finish_if_unfinished(DirectTurnTerminal::host_dropped());
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::agent::{
consume_direct_thread, direct_thread_turn_is_active, subscribe_direct_thread,
DirectThreadEvent, DirectTurnFailure, DirectTurnFailureKind,
};
/// 订阅并把 bootstrap 拿掉:之后的 `consume` 只返回这次订阅之后产生的事件。
fn watch(thread_id: &str) -> String {
let bootstrap = subscribe_direct_thread(thread_id);
let _ = consume_direct_thread(&bootstrap.subscription_id);
bootstrap.subscription_id
}
fn pending(subscription_id: &str) -> Vec<DirectThreadEvent> {
consume_direct_thread(subscription_id)
.expect("consume")
.events
}
fn turn_completed_events(events: &[DirectThreadEvent]) -> Vec<&DirectThreadEvent> {
events
.iter()
.filter(|event| matches!(event, DirectThreadEvent::TurnCompleted { .. }))
.collect()
}
fn unique_thread(label: &str) -> String {
format!("accept-test-{label}-{}", Uuid::new_v4())
}
#[test]
fn accept_emits_a_logical_turn_started_and_holds_the_turn() {
let thread = unique_thread("started");
let subscription = watch(&thread);
let reservation =
DirectTurnReservation::accept(&thread, "turn-1", Some("u-1")).expect("accept");
let events = pending(&subscription);
assert_eq!(events.len(), 1, "{events:?}");
match &events[0] {
DirectThreadEvent::TurnStarted { user_item_id, .. } => {
assert_eq!(user_item_id.as_deref(), Some("u-1"));
}
other => panic!("expected turn.started, got {other:?}"),
}
assert!(direct_thread_turn_is_active(&thread));
drop(reservation);
}
#[test]
fn a_second_accept_is_rejected_while_the_turn_is_open() {
let thread = unique_thread("busy");
let subscription = watch(&thread);
let reservation =
DirectTurnReservation::accept(&thread, "turn-1", Some("u-1")).expect("accept");
// 先取走第一条接单自己的开始事件,之后的"空"才只说明被拒的这一次没写东西。
assert_eq!(pending(&subscription).len(), 1);
let rejected = DirectTurnReservation::accept(&thread, "turn-2", Some("u-2"));
assert!(matches!(
rejected,
Err(DirectTurnError::TurnAlreadyRunning { .. })
));
let events = pending(&subscription);
assert_eq!(events.len(), 0, "被拒的接单不许产生事件:{events:?}");
drop(reservation);
}
/// 拒单载荷里的两个身份都是**回合身份**:撞的是同一轮时两者相等,界面才走得到"同一轮消息仍在
/// 处理中";撞的是另一轮时两者不等,界面才敢提示"另一条回合在运行"。占用 token 只活在本进程,
/// 一旦漏进载荷,这两个分支就都判不出来(UUID 永远不等于界面的 `clientTurnId`)。
#[test]
fn accept_conflict_reports_client_turn_ids_not_reservation_tokens() {
let thread = unique_thread("same-turn-conflict");
let reservation =
DirectTurnReservation::accept(&thread, "turn-1", Some("u-1")).expect("accept");
let same_turn = match DirectTurnReservation::accept(&thread, "turn-1", Some("u-1")) {
Ok(_) => panic!("同一 thread 的第二条回合必须被拒"),
Err(error) => error,
};
let DirectTurnError::TurnAlreadyRunning {
existing_invocation_id,
incoming_invocation_id,
} = &same_turn
else {
panic!("expected a concurrency rejection, got {same_turn:?}");
};
assert_eq!(existing_invocation_id, "turn-1");
assert_eq!(incoming_invocation_id, "turn-1");
assert!(
same_turn.to_string().contains("同一轮消息仍在处理中"),
"{same_turn}"
);
let other_turn = match DirectTurnReservation::accept(&thread, "turn-2", Some("u-2")) {
Ok(_) => panic!("同一 thread 的第二条回合必须被拒"),
Err(error) => error,
};
assert!(matches!(
&other_turn,
DirectTurnError::TurnAlreadyRunning {
existing_invocation_id,
incoming_invocation_id,
} if existing_invocation_id == "turn-1" && incoming_invocation_id == "turn-2"
));
assert!(
other_turn.to_string().contains("另一条 Direct 客户端回合"),
"{other_turn}"
);
drop(reservation);
}
#[test]
fn drop_without_a_terminal_writes_a_host_dropped_terminal() {
let thread = unique_thread("drop");
let subscription = watch(&thread);
let reservation =
DirectTurnReservation::accept(&thread, "turn-1", Some("u-1")).expect("accept");
assert!(reservation.finish_if_unfinished(DirectTurnTerminal::host_dropped()));
assert!(!direct_thread_turn_is_active(&thread));
// 显式收口之后 Drop 不再补第二条:兜底只负责"没人写过"的那一种。
drop(reservation);
let events = pending(&subscription);
let completed = turn_completed_events(&events);
assert_eq!(completed.len(), 1, "{events:?}");
match completed[0] {
DirectThreadEvent::TurnCompleted {
user_item_id,
failure: Some(failure),
..
} => {
assert_eq!(failure.kind, DirectTurnFailureKind::HostDropped);
assert_eq!(user_item_id.as_deref(), Some("u-1"));
}
other => panic!("expected a failed terminal, got {other:?}"),
}
}
#[test]
fn the_deep_terminal_wins_and_the_fallback_stays_silent() {
let thread = unique_thread("deep");
let subscription = watch(&thread);
let reservation =
DirectTurnReservation::accept(&thread, "turn-1", Some("u-1")).expect("accept");
// 深层收口:真正跑完这一轮的代码算出来的终态。
let deep = DirectThreadEvent::turn_completed_failed(
DirectTurnFailure::new(
DirectTurnFailureKind::Timeout,
"等待模型回执超时".to_string(),
),
2_000,
)
.with_user_item_id(Some("u-1"));
crate::agent::complete_direct_thread_turn(&thread, deep);
assert!(
!reservation.finish_if_unfinished(DirectTurnTerminal::host_dropped()),
"深层已收口时兜底不许再写"
);
drop(reservation);
let events = pending(&subscription);
let completed = turn_completed_events(&events);
assert_eq!(completed.len(), 1, "一轮只许有一条终态:{events:?}");
match completed[0] {
DirectThreadEvent::TurnCompleted { failure, .. } => {
assert_eq!(
failure.as_ref().map(|f| f.kind),
Some(DirectTurnFailureKind::Timeout)
);
}
other => panic!("expected a terminal, got {other:?}"),
}
}
#[test]
fn the_thread_can_be_accepted_again_after_the_turn_is_settled() {
let thread = unique_thread("again");
let first = DirectTurnReservation::accept(&thread, "turn-1", Some("u-1")).expect("accept");
drop(first);
let second =
DirectTurnReservation::accept(&thread, "turn-2", Some("u-2")).expect("second accept");
assert!(direct_thread_turn_is_active(&thread));
drop(second);
}
}
@@ -0,0 +1,357 @@
//! DirectProject 的放行:把队首那条待发消息变成一轮真的回合。
//!
//! 这个模块两件事,别再往里加第三件:
//! 1. [`DirectTurnReservation`]:这一轮的占用对象——持有它代表还没收口,终态出口只认它的 token;
//! 2. [`kick_direct_queue_dispatch`]:认领队首 → 起整轮(取调用身份 → 落盘用户条目 → 下发 → 跑回合)。
//!
//! 放行**不重跑任何入队检查,也没有"放行失败"**:入队时的检查已经判过,这一轮一旦放行,之后的
//! 一切失败都是**回合失败**,由占用对象收口成 `turn.completed` 的失败载荷。设计见
//! `docs/adr/【ADR】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md`。
use std::path::{Path, PathBuf};
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,
};
/// 一次放行的占用。持有它就代表这一轮还没收口。
///
/// 生命周期由调用方决定:登记与开始事件已经由 [`kick_direct_queue_dispatch`] 在同一临界区里写好,
/// 这里只接住这一轮的终态出口,并随放行任务一起 drop(正常 / 失败 / panic 都走同一个出口)。
pub(crate) struct DirectTurnReservation {
thread_id: String,
token: String,
user_item_id: Option<String>,
}
impl DirectTurnReservation {
/// 接管一次放行的占用:占用登记与 `turn.started` 已由 Thread Manager 的认领写好。
pub(crate) fn resume(thread_id: &str, dispatched: &DirectDispatchedTurn) -> Self {
Self {
thread_id: thread_id.to_string(),
token: dispatched.token.clone(),
user_item_id: dispatched.pending.user_item_id(),
}
}
/// 放行之后还没走到深层终态就失败的收口口:只有这一轮仍被自己占用时才写。
///
/// 深层(真正跑完这一轮的代码)已经写出终态时返回 `false`,兜底不覆盖真实结果。
pub(crate) fn finish_if_unfinished(&self, terminal: DirectTurnTerminal) -> bool {
complete_direct_thread_turn_if_reserved(
&self.thread_id,
&self.token,
terminal.event(direct_tool_call_now_ms(), self.user_item_id.as_deref()),
)
}
/// 测试用:直接开一轮占用,等价于"入队 + 认领",不经过 Tauri 命令与真实的用户条目。
#[cfg(test)]
pub(crate) fn accept_for_test(thread_id: &str, client_turn_id: &str) -> Self {
let pending = PendingDirectTurn::prepare(
client_turn_id.to_string(),
serde_json::from_value(serde_json::json!({
"type": "message",
"role": "user",
"content": [{"type": "input_text", "text": "测试消息"}],
"id": format!("direct-codex:{client_turn_id}:user"),
}))
.expect("canonical user item"),
"测试消息".to_string(),
None,
None,
direct_tool_call_now_ms(),
)
.expect("prepare pending turn");
super::enqueue_direct_pending_turn(thread_id, pending).expect("enqueue test turn");
let dispatched = claim_direct_pending_turn(thread_id).expect("claim test turn");
Self::resume(thread_id, &dispatched)
}
}
impl Drop for DirectTurnReservation {
fn drop(&mut self) {
// 兜底:任务 panic、future 被丢弃、或今后在终态之前新增的 `?` 早退。
// 这类失败说不出原因,只给分类;能说清原因的错误必须由调用方在更早的地方显式收口。
let _ = self.finish_if_unfinished(DirectTurnTerminal::host_dropped());
// 占用的释放就是"这个项目腾出了跑回合的位置":踢一脚,队里排着的消息不必等下一次用户动作。
// panic 也走这里,所以队列不会因为一个任务炸掉而永久停住。
kick_direct_queue_dispatch(Path::new(&self.thread_id));
}
}
/// 踢一脚:让队首那条待发消息有机会变成一轮真的回合。
///
/// 幂等,三个调用点:**入队之后**(队列空且没有回合在跑时,放行就是这一脚,不必再等一个调度周期)、
/// **占用释放**(正常 / 失败 / panic 共用,见 [`DirectTurnReservation`] 的 `Drop`)、**中止路径**
/// (守卫被兜底释放时那一轮不会再有人替它收尾)。
///
/// 它什么都不返回:没有"放行失败"。踢不动就是现在不该跑——已经有回合在跑,或者队列是空的;
/// 这两种情况都会在下一次释放时再踢。
pub(crate) fn kick_direct_queue_dispatch(root: &Path) {
let thread_id = direct_thread_id_for_project(root);
let Some(dispatched) = claim_direct_pending_turn(&thread_id) else {
return;
};
let reservation = DirectTurnReservation::resume(&thread_id, &dispatched);
let root = root.to_path_buf();
tauri::async_runtime::spawn(async move {
run_dispatched_direct_turn(root, dispatched, reservation).await;
});
}
/// 放行之后这一段:取这一轮的调用身份 → 落盘用户条目 → 下发 → 跑整轮。
///
/// 这一段的每个失败都是**回合失败**:占用与开始事件已经写好,只走占用对象的终态出口,不再回到任何
/// 命令的返回值上。
async fn run_dispatched_direct_turn(
root: PathBuf,
dispatched: DirectDispatchedTurn,
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)
{
// 不继续起整轮:历史是这条对话的单一事实源,用户消息没落盘时继续跑只会得到一条没有开口用户
// 消息的助手回复,而且失败会被静默掉。
let failure = DirectTurnError::EnvironmentNotReady {
detail: redact_agent_runtime_error(
&root,
&format!("写入本项目对话历史失败:{error}"),
600,
),
};
reservation.finish_if_unfinished(DirectTurnTerminal::failed(&root, &failure));
return;
}
// 用户条目落盘成功即下发:这一轮从"放行"到"起 codex"之间的一切失败(连不上 app-server、执行器
// 未通过验收、历史注入失败)都靠它把失败说明挂回自己那一轮。
crate::agent::codex_app_server::emit_direct_thread_user_item(
&root,
&dispatched.pending.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,
dispatched.pending.creation_type.as_deref(),
Some(&emitter),
Some(dispatched.pending.canonical_user_item.clone()),
capture,
dispatched.pending.analytics_attempt_id.as_deref(),
)
.await;
match outcome {
Ok(reply) => {
// 深层的终态出口已经在 `run_turn` 里写出 `turn.completed`;这里只补最后一条回合更新。
emitter.emit("completed", Some("none"), Some(reply), None);
}
Err(error) => {
// 放行之后的失败一律是回合失败:失败诊断与失败说明已由上层写过,这里补终态事件。
// 深层已经写出终态时它不覆盖(同一轮只允许一条终态)。
reservation.finish_if_unfinished(DirectTurnTerminal::failed(&root, &error));
}
}
// 调用身份先放,再让占用对象的 `Drop` 收尾并踢下一脚:顺序反过来就会出现"守卫还握着、队首却已经
// 可以放行"的窗口——那一脚会被自己的守卫挡成一条假的回合失败。
drop(active_invocation);
}
#[cfg(test)]
mod tests {
use super::*;
use crate::agent::{
consume_direct_thread, direct_thread_turn_is_active, enqueue_direct_pending_turn,
subscribe_direct_thread, DirectThreadEvent, DirectTurnFailure, DirectTurnFailureKind,
PendingDirectTurn,
};
use uuid::Uuid;
/// 订阅并把 bootstrap 拿掉:之后的 `consume` 只返回这次订阅之后产生的事件。
fn watch(thread_id: &str) -> String {
let bootstrap = subscribe_direct_thread(thread_id);
let _ = consume_direct_thread(&bootstrap.subscription_id);
bootstrap.subscription_id
}
fn pending(subscription_id: &str) -> Vec<DirectThreadEvent> {
consume_direct_thread(subscription_id)
.expect("consume")
.events
}
fn turn_completed_events(events: &[DirectThreadEvent]) -> Vec<&DirectThreadEvent> {
events
.iter()
.filter(|event| matches!(event, DirectThreadEvent::TurnCompleted { .. }))
.collect()
}
/// 一条待发消息。放行路径只读它的身份与产物,所以这里用最小可用形状。
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,
None,
direct_tool_call_now_ms(),
)
.expect("prepare pending turn")
}
/// 真实存在的项目目录不是这条用例的判据,用唯一的假路径当线程身份即可:认领与占用的语义只认
/// 这个字符串。
fn unique_thread(label: &str) -> String {
format!("dispatch-test-{label}-{}", Uuid::new_v4())
}
fn enqueue(thread_id: &str, client_turn_id: &str) {
enqueue_direct_pending_turn(thread_id, pending_turn(client_turn_id)).expect("enqueue");
}
/// 放行的占用只认**自己的 token**:认领写下的占用身份与这一轮绑死,别人的终态盖不上它。
#[test]
fn a_foreign_token_cannot_close_the_turn() {
let thread = unique_thread("token");
enqueue(&thread, "turn-1");
let dispatched = claim_direct_pending_turn(&thread).expect("claim head");
let owner = DirectTurnReservation::resume(&thread, &dispatched);
let foreign = DirectTurnReservation::resume(
&thread,
&DirectDispatchedTurn {
token: "foreign-token".to_string(),
pending: pending_turn("turn-1"),
},
);
assert!(!foreign.finish_if_unfinished(DirectTurnTerminal::host_dropped()));
assert!(direct_thread_turn_is_active(&thread));
assert!(owner.finish_if_unfinished(DirectTurnTerminal::host_dropped()));
assert!(!direct_thread_turn_is_active(&thread));
// 显式收口之后 Drop 不再补第二条:兜底只负责"没人写过"的那一种。
drop(foreign);
drop(owner);
}
/// 深层终态先写,占用对象的兜底就闭嘴:一轮只许有一条终态。
#[test]
fn the_deep_terminal_wins_and_the_fallback_stays_silent() {
let thread = unique_thread("deep");
let subscription = watch(&thread);
enqueue(&thread, "turn-1");
let dispatched = claim_direct_pending_turn(&thread).expect("claim head");
let reservation = DirectTurnReservation::resume(&thread, &dispatched);
// 深层收口:真正跑完这一轮的代码算出来的终态。
let deep = DirectThreadEvent::turn_completed_failed(
DirectTurnFailure::new(
DirectTurnFailureKind::Timeout,
"等待模型回执超时".to_string(),
),
2_000,
)
.with_user_item_id(Some("direct-codex:turn-1:user"));
crate::agent::complete_direct_thread_turn(&thread, deep);
assert!(
!reservation.finish_if_unfinished(DirectTurnTerminal::host_dropped()),
"深层已收口时兜底不许再写"
);
drop(reservation);
let events = pending(&subscription);
let completed = turn_completed_events(&events);
assert_eq!(completed.len(), 1, "一轮只许有一条终态:{events:?}");
match completed[0] {
DirectThreadEvent::TurnCompleted { failure, .. } => {
assert_eq!(
failure.as_ref().map(|f| f.kind),
Some(DirectTurnFailureKind::Timeout)
);
}
other => panic!("expected a terminal, got {other:?}"),
}
}
/// 没有终态的 Drop 会补一条 `host-dropped`,并且**踢一脚**让队里排着的下一条接上。
#[test]
fn dropping_a_reservation_closes_it_and_lets_the_queue_move_on() {
let thread = unique_thread("drop");
let subscription = watch(&thread);
enqueue(&thread, "turn-1");
enqueue(&thread, "turn-2");
let dispatched = claim_direct_pending_turn(&thread).expect("claim head");
let reservation = DirectTurnReservation::resume(&thread, &dispatched);
drop(reservation);
// 第一条按宿主丢弃收口,第二条已经接上:放行不等人,也不必等下一次用户动作。
let events = pending(&subscription);
let terminal = events
.iter()
.position(|event| matches!(event, DirectThreadEvent::TurnCompleted { .. }))
.expect("第一条必须被收口");
assert!(
matches!(
events.get(terminal),
Some(DirectThreadEvent::TurnCompleted { status, failure: Some(failure), .. })
if status == "failed" && failure.kind == DirectTurnFailureKind::HostDropped
),
"{events:?}"
);
let next_started = events.iter().position(|event| {
matches!(
event,
DirectThreadEvent::TurnStarted { user_item_id, .. }
if user_item_id.as_deref() == Some("direct-codex:turn-2:user")
)
});
assert!(
next_started.is_some_and(|index| index > terminal),
"队列里排着的下一条必须在收口之后立刻放行:{events:?}"
);
}
/// 已经有回合在跑、或队列为空时,踢一脚什么都不做。
#[test]
fn kicking_does_nothing_while_a_turn_is_open_or_the_queue_is_empty() {
let empty = unique_thread("empty");
kick_direct_queue_dispatch(Path::new(&empty));
assert!(!direct_thread_turn_is_active(&empty));
let busy = unique_thread("busy");
enqueue(&busy, "turn-1");
enqueue(&busy, "turn-2");
let dispatched = claim_direct_pending_turn(&busy).expect("claim head");
let reservation = DirectTurnReservation::resume(&busy, &dispatched);
kick_direct_queue_dispatch(Path::new(&busy));
// 仍在跑的那一轮没有被顶掉:占用身份还是第一条。
assert!(direct_thread_turn_is_active(&busy));
drop(reservation);
}
}
@@ -5,13 +5,13 @@
//! 通道断开要带宿主诊断。用不同变体各带各的字段,分流靠 `match`,不靠 `kind` 字段 + 共用字段的
//! 伪结构化,也不靠对错误文本做子串匹配。
//!
//! 走哪条通道由**发生位置**决定,不由错误种类决定(`接单化` 之后的口径):
//! - **接单前**发生的 = 拒单:只出提示 / 横幅,不做失败载荷、不写失败诊断、不上报成
//! 走哪条通道由**发生位置**决定,不由错误种类决定(`入队化` 之后的口径):
//! - **入队前**发生的 = 入队失败:只出提示 / 横幅,不做失败载荷、不写失败诊断、不上报成
//! "智能创作失败"。命令返回 `Err` 的就是这一类。
//! - **接单后**发生的 = 回合失败:事件载荷、横幅、应用日志、错误上报池四处一致;命令早已返回
//! - **放行后**发生的 = 回合失败:事件载荷、横幅、应用日志、错误上报池四处一致;命令早已返回
//! `Ok`,所以一律由宿主侧的占用对象投影成 `turn.completed.failure`。
//!
//! 所以"同一种错误在接单前后走不同通道"是正常的:[`DirectTurnError::EnvironmentNotReady`] 两边
//! 所以"同一种错误在入队与放行走不同通道"是正常的:[`DirectTurnError::EnvironmentNotReady`] 两边
//! 都可能出现,位置说了算。这里**没有**、也不该有"这个变体是不是回合失败"的判据。
//!
//! 事件载荷(`direct_thread_wire::DirectTurnFailure`)仍然只有 `{kind, message}` 两个字段:
@@ -38,7 +38,7 @@ const DIRECT_CODEX_NATIVE_KIND_PREFIX: &str = "codex-app-server-error:";
/// 失败发生在交付的哪一段。与错误分类正交:分类说明"怎么回事",阶段说明"走到哪一步"。
///
/// 线上取值跟着拒单 / 失败载荷一起给前端(`art-preparation` 这类),所以也要导出。
/// 线上取值跟着入队失败 / 回合失败载荷一起给前端(`art-preparation` 这类),所以也要导出。
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, TS)]
#[serde(rename_all = "kebab-case")]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
@@ -62,7 +62,7 @@ impl DirectCodexFailureStage {
/// 宿主等不到模型回执时,撞的是哪一条上限。
///
/// 跟着拒单 / 失败载荷一起给前端,界面不靠文案区分这两条。
/// 跟着入队失败 / 回合失败载荷一起给前端,界面不靠文案区分这两条。
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, TS)]
#[serde(rename_all = "kebab-case")]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
@@ -227,7 +227,7 @@ pub(crate) enum DirectTurnFailureKind {
TransportFailed,
/// app-server 或上游明确拒绝了这次请求。
RequestRejected,
/// 接单之后的连接 / 配置 / 凭据 / 脚手架未就绪(不是模型的错,界面语气也不同)。
/// 放行之后的连接 / 配置 / 凭据 / 脚手架未就绪(不是模型的错,界面语气也不同)。
EnvironmentNotReady,
/// app-server 单方面把这一轮判成中断(用户没要求停止、宿主也没在收尾)。
TurnInterrupted,
@@ -235,7 +235,7 @@ pub(crate) enum DirectTurnFailureKind {
HostDropped,
}
/// 模型调用失败(app-server 一次 `turn` 的结果)的分类,跟着拒单 / 失败载荷一起给前端。
/// 模型调用失败(app-server 一次 `turn` 的结果)的分类,跟着入队失败 / 回合失败载荷一起给前端。
///
/// 每个变体对应平台层 `LlmError` 的一个分支,于是 [`DirectTurnError::wire_kind`] 的取值与改造前
/// 完全一致:事件的 `failure.kind` 就是这一份取值,界面按它选语气,不拿它做流程分支。
@@ -385,7 +385,7 @@ impl DirectModelCallKind {
)]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) enum DirectTurnError {
// ───────── 拒单:接单之前发生,这一轮没有开始 ─────────
// ───────── 入队失败:这一轮没有开始,也没有进队列 ─────────
/// `clientTurnId` 没给:没有稳定回合身份,拒绝创建可计费身份。
ClientTurnIdMissing,
/// `clientTurnId` 形状非法:长度与字符集由宿主定,界面按同一份约束生成。
@@ -397,6 +397,10 @@ pub(crate) enum DirectTurnError {
existing_invocation_id: String,
incoming_invocation_id: String,
},
/// 待发消息队列已满:这一条没进队,等前面几条发完再发。
///
/// 上限只落在宿主这一处(`MAX_PENDING_DIRECT_TURNS`),随载荷带出去,界面不自己数一份。
QueueFull { limit: usize },
/// 项目目录锚不定(符号链接 / 权限 / 目录被删)。
ProjectRootUnanchored { cause: String },
/// 项目目录不存在或不是绝对路径。
@@ -407,13 +411,13 @@ pub(crate) enum DirectTurnError {
InputRejected { detail: String },
/// 结构化消息既没有正文也没有任何引用。
ContentEmpty,
/// 环境 / 凭据 / 脚手架未就绪。**接单前后都可能出现**:接单前是拒单(工程 / 凭据还没准备好),
/// 接单后是回合失败(分类 `environment-not-ready`,例如 `turn/start` 之前连不上 app-server)。
/// 环境 / 凭据 / 脚手架未就绪。**入队与放行都可能出现**:入队时是入队失败(工程 / 凭据还没准备
/// 好),放行后是回合失败(分类 `environment-not-ready`,例如 `turn/start` 之前连不上 app-server)。
EnvironmentNotReady { detail: String },
/// 宿主执行账本取不到(初始化失败、归属锁被占、状态损坏、时钟回退)。
HostStateUnavailable { detail: String },
// ───────── 回合失败:接单之后发生,这一轮已经开始 ─────────
// ───────── 回合失败:放行之后发生,这一轮已经开始 ─────────
/// 模型调用失败:`kind` 是分类,`detail` 是平台层原文(就是给用户看的那句话)。
ModelCallFailed {
kind: DirectModelCallKind,
@@ -442,7 +446,7 @@ pub(crate) enum DirectTurnError {
TurnFailedUnclassified { detail: String },
}
/// 拒单载荷:命令边界交给前端的**结构化拒绝**。
/// 入队失败载荷:命令边界交给前端的**结构化入队失败**。
///
/// 为什么不是只给一句话:界面要按变体分流——认得的"前置条件不满足 / 用户参数无效"给一条与用户
/// 消息同级的提示且不上报,认不得的原样抛出交给既有捕获链路。文案只是给人看的最后一步,仍由
@@ -450,14 +454,14 @@ pub(crate) enum DirectTurnError {
#[derive(Clone, Debug, PartialEq, Eq, Serialize, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) struct DirectTurnRejection {
pub(crate) struct DirectTurnEnqueueFailure {
/// 结构化变体:界面按 `error.type` 分流,不解析文案。
pub(crate) error: DirectTurnError,
/// 可展示文案(`Display` 的唯一出口)。
pub(crate) message: String,
}
impl DirectTurnRejection {
impl DirectTurnEnqueueFailure {
pub(crate) fn new(error: DirectTurnError) -> Self {
Self {
message: error.to_string(),
@@ -467,12 +471,12 @@ impl DirectTurnRejection {
}
impl DirectTurnError {
/// 命令边界要不要为这条**拒单**补一份运行错误诊断。
/// 命令边界要不要为这条**入队失败**补一份运行错误诊断。
///
/// 只有"宿主 / 环境的事实故障、用户自己改不了"才值得进 `.agent/runtime/errors` 与应用日志;
/// 空内容、`clientTurnId` 形状、另一轮在跑、权限策略、目录锚不定 / 不是绝对路径都是用户自己
/// 就能修的操作结果,留痕只会变成噪声;它们仍按 `Display` 给用户一句可读的话。判据按变体分,
/// 不看文案。
/// 空内容、`clientTurnId` 形状、队列满、另一轮在跑、权限策略、目录锚不定 / 不是绝对路径都是
/// 用户自己就能修的操作结果,留痕只会变成噪声;它们仍按 `Display` 给用户一句可读的话。
/// 判据按变体分,不看文案。
///
/// 回合级失败恒为 `false`:它们在上游(`record_direct_codex_failure`)已经写过诊断,边界再写一次
/// 就是同一件事留两份。
@@ -482,6 +486,7 @@ impl DirectTurnError {
Self::EnvironmentNotReady { .. } | Self::HostStateUnavailable { .. } => true,
Self::ClientTurnIdMissing
| Self::ClientTurnIdMalformed { .. }
| Self::QueueFull { .. }
| Self::TurnAlreadyRunning { .. }
// 项目目录锚不定(符号链接 / 权限 / 目录被删)与目录不存在同类:都是用户能自己修好的
// 文件系统事实,诊断文案不该顶替那句"无法锚定 Direct 调用项目目录:{cause}"。
@@ -644,6 +649,10 @@ impl fmt::Display for DirectTurnError {
formatter,
"clientTurnId 必须为 {min_chars} 到 {max_chars} 位 ASCII 字母、数字或连字符,且首位必须为字母或数字"
),
Self::QueueFull { limit } => write!(
formatter,
"待发消息已达上限(最多 {limit} 条),请等前面几条发完再发送"
),
Self::TurnAlreadyRunning {
existing_invocation_id,
incoming_invocation_id,
@@ -5,7 +5,7 @@
//! 1. [`direct_turn_terminal`]:拿这一轮的事实判定终态——是不是失败、原因是什么、状态写什么;
//! 2. [`DirectTurnTerminal::event`]:把终态投影成 `turn.completed` 事件。
//!
//! 终态的**出口**(谁写、什么时候兜底)不在这里,在 `direct_turn_accept.rs` 的接单占用对象里:
//! 终态的**出口**(谁写、什么时候兜底)不在这里,在 `direct_turn_dispatch.rs` 的放行占用对象里:
//! 这个模块只负责"什么算失败、原因怎么写"。
//!
//! 失败载荷的**形状**属于线上协议,定义在 `direct_thread_wire.rs`(`DirectTurnFailure`);
@@ -5567,6 +5567,29 @@ pub(crate) fn list_game_creator_direct_active_turns(
list_direct_active_turns()
}
/// 取消一条**还没放行**的待发消息。
///
/// 结果分三种:真的从队列移除了(并以 `queue.removed{cancelled}` 告知所有订阅者)、这一条已经被
/// 放行(正在跑的那一轮就是它,不能按待发消息取消)、队列里本来就没有这个身份。判据只看队列与占用,
/// 不看界面状态——同一份事实源在宿主。
#[tauri::command]
pub(crate) async fn remove_direct_project_pending_turn(
project_path: String,
client_turn_id: String,
) -> Result<DirectQueueRemovalOutcome, String> {
tauri::async_runtime::spawn_blocking(move || {
let root = Path::new(project_path.trim());
enforce_project_permission_policy(root, "conversation.write")?;
let thread_id = direct_thread_id_for_project(root);
Ok(remove_direct_pending_turn(
&thread_id,
client_turn_id.trim(),
))
})
.await
.map_err(|error| format!("取消 DirectProject 待发消息后台任务失败:{error}"))?
}
#[tauri::command]
pub(crate) async fn subscribe_direct_project_thread(
project_path: String,
@@ -2619,7 +2619,8 @@ fn main() {
chat_with_game_creator_agent,
chat_with_game_creator_role_agent,
chat_with_game_creator_role_agent_stream,
chat_with_game_creator_direct_codex,
enqueue_direct_codex_turn,
remove_direct_project_pending_turn,
preflight_web_game_creation,
cancel_direct_codex_turn,
select_game_creator_reasoning_effort,
@@ -946,7 +946,7 @@ function directCodexDiagnosticFailureParts(
return { stageLabel, summary, hint, retryable: retryable === 'true' };
}
/** 收口文案的一句话:阶段标签只有"回合失败"才成立,拒单那边传 `null`(见下面的导出函数)。 */
/** 收口文案的一句话:阶段标签只有"回合失败"才成立,请求被拒那边传 `null`(见下面的导出函数)。 */
function directDiagnosticSentence(
parts: DirectDiagnosticParts,
stageLabel: string | null,
@@ -962,10 +962,10 @@ function directCodexDiagnosticFailureDetail(message: string) {
}
/**
* **拒单**的可见文案:宿主收口文案里那段已脱敏的摘要与建议。
* **请求被拒**的可见文案:宿主收口文案里那段已脱敏的摘要与建议。
*
* 与 [`projectRuntimeVisibleError`] 的差别只有一处——**不带阶段标签**:阶段说的是"失败发生在交付的
* 哪一步",而拒单是"这一轮没有开始",阶段只会是默认值(`code-generation`),套上去会把没发生的
* 哪一步",而请求被拒是"这一轮没有开始",阶段只会是默认值(`code-generation`),套上去会把没发生的
* 事讲成发生了。宿主的文案不是收口形状时退回同一份运行错误映射。
*/
export function projectRuntimeVisibleRejectionError(
@@ -163,7 +163,7 @@ export function DirectProjectChatView({
turns: directTurns,
});
// 状态条的起点只认**未结束**的最新一轮:只有它才拿得到本轮的 `turn.started.at`(运行中读实时值)。
// 只可能是最后一轮:原生 `turnRunning` 只赋给最新一轮。接单窗口里还没有这一轮的条目,
// 只可能是最后一轮:原生 `turnRunning` 只赋给最新一轮。入队窗口里还没有这一轮的条目,
// 于是这里是 0,卡片只报"正在处理"、不读秒(宿主开始事件一到就开始读秒)。
const latestTurn = directTurns.at(-1) ?? null;
const activeTurnStartedAt =
@@ -37,10 +37,10 @@ import {
directCodexConversationMessageId,
directCodexPolicyRetryInput,
type DirectProjectTurnInput,
directTurnRejectionNotice,
directTurnRejectionNoticeMessageId,
directTurnUnrecognizedRejectionNoticeText,
readDirectTurnRejection,
directTurnEnqueueFailureNotice,
directTurnEnqueueFailureNoticeMessageId,
directTurnUnrecognizedEnqueueFailureNoticeText,
readDirectTurnEnqueueFailure,
} from '../conversation/directCodexConversation';
import { DIRECT_CODEX_SESSION_KEEPALIVE_MS } from '../conversation/directCodexSessionKeepalive';
import {
@@ -103,7 +103,7 @@ export type DirectProjectChatControllerProps = {
* 传输层(三条通道,前端各拉各的)
* A 运行态:notify → invoke consume_direct_project_thread → events[](实时)
* B 历史: invoke read_direct_project_history_slice → items[](分页,文件尾反向扫描)
* C 本地: 前端自己造(忙态、拒单提示、终止说明)——**不造用户消息**
* C 本地: 前端自己造(忙态、入队失败提示、终止说明)——**不造用户消息**
* ▼
* 前端
* useDirectThreadChatSubscription reducer:A + B 进同一份 state(turnRunning / history / live)
@@ -119,29 +119,29 @@ export type DirectProjectChatControllerProps = {
* - 项目对话历史(`.agent/conversations/project.jsonl`):持久,只有条目、**没有回合边界**,
* 经历史切片读取(首屏按 `lastCompletedItemId` 锚定)。
* - 运行态事件(subscribe / consume / notify):进程内;`turn.started` / `turn.completed` 是**逻辑回合**
* 活跃与否的**唯一**判据(接单时成对发出,不再镜像 Codex 原生回合);可回收事件被回收后靠
* 活跃与否的**唯一**判据(放行时成对发出,不再镜像 Codex 原生回合);可回收事件被回收后靠
* `lifecycle_anchor` 保住最新一条生命周期事件。
* 失败也走这条流:`turn.completed.failure` 自己带脱敏后的原因,reducer 把它落成本轮说明条目;
* 命令返回那条通道只提供横幅与诊断,不再写聊天文案。
* - 本地说明:只存在于本次会话,`projectPath` 变化即清空;只有壳层 `announce` 与拒单提示两种。
* 用户消息一律来自宿主条目——本地乐观气泡已删(见 ADR「DirectProject命令接单化」的后续更新),
* - 本地说明:只存在于本次会话,`projectPath` 变化即清空;只有壳层 `announce` 与入队失败提示两种。
* 用户消息一律来自宿主条目——本地乐观气泡已删(见 ADR「DirectProject命令入队化」),
* 所以"用户那句话说没说出去"只有宿主条目一个来源。
*
* 一次发送的时序(第 2 → 3 步之间就是「本地已发出、宿主还没确认」的空窗:聊天区里没有这一轮的
* 任何条目,只有 composer 忙态与状态行):
* 1. 按下发送:`turnBusy=true`(同帧)。聊天区不动——这一轮在宿主认领之前不存在。
* 2. `invoke('chat_with_game_creator_direct_codex')`:Rust 走完接单前的检查 → 接单(登记占用 +
* append `turn.started`)→ 落盘用户条目 → spawn 整轮 → **立刻返回**。命令返回只说明接单成立,
* 整轮的结果不再从这条通道回来;拒单则返回结构化的 typed 错误。
* 1. 按下发送:`turnBusy=true`(同帧)。聊天区不动——这一轮在宿主放行之前不存在。
* 2. `invoke('enqueue_direct_codex_turn')`:Rust 走完入队检查 → **入队**(排进待发消息队列)→
* **立刻返回**。命令返回只说明入队成立;放行(登记占用 + `turn.started` + 离开队列)由宿主
* 在队首就绪时自己完成,整轮的结果不再从这条通道回来,入队失败则返回结构化的 typed 错误。
* 3. notify → consume → `turn.started`:reducer 的 `turnRunning=true`、`turnStartedAt`、`turnUserItemId`。
* 4. `item.completed`(本轮用户条目下发):同身份条目已在历史里就合并进去,否则进 `live`。这是这一轮
* 的用户气泡**第一次**出现在聊天区(发点在接单之后、起 codex 之前)。
* 的用户气泡**第一次**出现在聊天区(发点在放行之后、起 codex 之前)。
* 5. `item.delta` / `item.started` / `item.completed`:正文追加、工具卡片 upsert(先到定形、后到只补空)。
* 6. `turn.completed`:`live` 并入 `history` 后清空,`turnEndedAt` 冻结,边界按身份盖到本轮开口条目上,
* 收口计数 +1;带 `failure` 载荷时,说明条目已经在上一步由 reducer 落进 `live`,随本轮一起并入历史。
* 7. **回合终态**(第 6 步的收口计数变化):结算本轮埋点 → 放行发送队列,顺序固定在这一处。
* 8. 命令收尾(`finally`):刷新清单;只在**没接单**时放掉忙态并出队(权限被拒那种
* 「本轮从未发出但要继续出队」的路径也在这里收口),接单成立的那一轮交给第 3 / 7 步。
* 8. 命令收尾(`finally`):刷新清单;只在**入队失败**时放掉忙态并出队(权限被拒那种
* 「本轮从未发出但要继续出队」的路径也在这里收口),入队成立的那一轮交给第 3 / 7 步。
*
* 状态变量归属:reducer 的回合字段、收口计数与 `history` / `live` 只由 `directThreadChat.ts` 写;
* 本文件的 `turnBusy`(唯一入口 `beginTurnBusy` / `endTurnBusy`)、`localMessages`、发送队列、
@@ -213,7 +213,7 @@ export function useDirectProjectChatController({
const completionPendingRef = useRef(false);
/**
* 本轮的埋点句柄。它在命令返回之后仍然要活着:成绩是**回合末**才在宿主侧入账的,
* 接单返回时结算只会静默丢掉这一次埋点(见 `settlePendingRunAnalytics`)。
* 入队返回时结算只会静默丢掉这一次埋点(见 `settlePendingRunAnalytics`)。
*/
const pendingRunAnalyticsRef = useRef<{
clientTurnId: string;
@@ -266,12 +266,12 @@ export function useDirectProjectChatController({
}, [turnBusy, currentTurnRunning, completedTurnCount]);
/**
* 回合终态是队列放行与埋点结算的唯一出口(接单被拒走命令那条路,见 `startTurn` 的收尾)。
* 回合终态是队列放行与埋点结算的唯一出口(入队失败走命令那条路,见 `startTurn` 的收尾)。
*
* 判据用 reducer 的**单调计数**而不是 `turnRunning` 的下降沿:一轮可能在同一次 consume 里
* 开始并结束,那时下降沿永远不会出现,队列就永久卡住了。
*
* TODO(发送队列):这条队列整体挪到 Rust 端,放行点就是 Thread Manager 的接单动作。
* TODO(发送队列):这条队列整体挪到宿主的待发消息队列(放行由 Thread Manager 自己驱动)。
*/
useEffect(() => {
if (!enabled) {
@@ -318,10 +318,10 @@ export function useDirectProjectChatController({
}, [enabled, projectPath]);
/**
* 「本地命令在飞」的唯一起止点:按下发送时置上,宿主认领这一轮(`turn.started` 落进 reducer)
* 「本地命令在飞」的唯一起止点:按下发送时置上,宿主放行这一轮(`turn.started` 落进 reducer)
* 或这一轮明确没成立时放掉。
*
* 它与原生忙态是两件事,所以**不在这里**按回合身份收口:接单化之后命令只等到接单就返回,
* 它与原生忙态是两件事,所以**不在这里**按回合身份收口:入队化之后命令只等到入队就返回,
* 「宿主认领了吗」由 reducer 的 `turnRunning` 回答(`displayBusy` 是两者的并集)。
*/
function beginTurnBusy() {
@@ -336,8 +336,8 @@ export function useDirectProjectChatController({
}
/**
* 结算本轮埋点。只在**回合终态**调用:成绩是回合末才在宿主侧入账的,接单返回时就结算
* 会变成一次空操作(宿主找不到候选,静默丢弃)。拒单那一轮没有候选,句柄由 `runTurn` 自己清掉。
* 结算本轮埋点。只在**回合终态**调用:成绩是回合末才在宿主侧入账的,入队返回时就结算
* 会变成一次空操作(宿主找不到候选,静默丢弃)。入队失败那一轮没有候选,句柄由 `runTurn` 自己清掉。
*/
function settlePendingRunAnalytics() {
const pending = pendingRunAnalyticsRef.current;
@@ -516,7 +516,7 @@ export function useDirectProjectChatController({
/**
* 发起一轮 DirectProject 回合:写权限门 + invoke + 本地忙态与错误收尾。
*
* 这里**不往聊天里写用户消息**:这一轮的用户气泡只来自宿主条目,所以"接单窗口期聊天区没有这一
* 这里**不往聊天里写用户消息**:这一轮的用户气泡只来自宿主条目,所以"入队窗口期聊天区没有这一
* 轮"是正常现象(忙态与状态行负责告知)。权限确认后重跑的是同一份输入,只跳过权限检查。
*/
function startTurn(input: DirectProjectTurnInput) {
@@ -582,14 +582,14 @@ export function useDirectProjectChatController({
if (invoked) {
await refreshDirectManifest(nextProjectPath);
}
// 收尾分两种:接单成立的整轮交给宿主的事件(`turn.started` 时退场、`turn.completed`
// 时出队与结算,见上面两个 effect);没成立的那些路径没有任何事件会来,只能在这里收口。
// 收尾分两种:入队成立的整轮交给宿主的事件(`turn.started` 时退场、`turn.completed`
// 时出队与结算,见上面两个 effect);入队失败的那些路径没有任何事件会来,只能在这里收口。
if (!turnAccepted) {
// 忙态必须早于出队放掉:出队会同步开始下一轮并设上它自己的忙态,清在它后面就等于
// 把下一轮的忙态抹掉(composer 会以为可以并发发送)。
endTurnBusy();
// 出队条件:权限被拒的一轮从未发出,不受项目切换影响,照旧出队;命令真发出过又返回
// 拒单时要求项目没被换掉;没有 invoke(非 Tauri 环境)时不出队,避免空转。
// 入队失败时要求项目没被换掉;没有 invoke(非 Tauri 环境)时不出队,避免空转。
if (
queueAdvance ||
(invoked && projectPathRef.current === nextProjectPath)
@@ -606,7 +606,7 @@ export function useDirectProjectChatController({
nextProjectPath: string,
input: DirectProjectTurnInput,
): Promise<boolean> {
// 命令返回 `Ok` 只说明**接单成立**:整轮怎么收场只由 `turn.completed` 事件回答。
// 命令返回 `Ok` 只说明**入队成立**:整轮怎么收场只由 `turn.completed` 事件回答。
// 所以这里的返回值只服务队列放行——"这一轮有没有真的开始"。
let turnAccepted = false;
try {
@@ -619,7 +619,7 @@ export function useDirectProjectChatController({
clientTurnId: input.clientTurnId,
runAnalytics,
};
await invoke<string>('chat_with_game_creator_direct_codex', {
await invoke<string>('enqueue_direct_codex_turn', {
projectPath: nextProjectPath,
clientTurnId: input.clientTurnId,
userItem: input.userItem,
@@ -629,11 +629,11 @@ export function useDirectProjectChatController({
turnAccepted = true;
// 清单刷新统一交给 startTurn 的 finally:成功与报错路径都覆盖,且只读一次。
} catch (error) {
// 命令的拒单是**结构化的**:命令返回 `Ok` 只说明接单成立,所以这条 catch 从接单化之后
// 只剩"拒单"一种输入(整轮结果由 `turn.completed` 事件回答,不再回到这里)。
const rejection = readDirectTurnRejection(error);
// 只有结构化拒单能证明"这一轮没接单"(拒单不产生回合事件、也就不会有埋点候选),句柄才
// 只清不发;非结构化错误(IPC 失败、命令 panic)可能发生在接单之后,那时必须留着句柄等
// 命令的入队失败是**结构化的**:命令返回 `Ok` 只说明入队成立,所以这条 catch 从入队化之后
// 只剩"入队失败"一种输入(整轮结果由 `turn.completed` 事件回答,不再回到这里)。
const rejection = readDirectTurnEnqueueFailure(error);
// 只有结构化入队失败能证明"这一轮没成立"(入队失败不产生回合事件、也就不会有埋点候选),
// 句柄才只清不发;非结构化错误(IPC 失败、命令 panic)可能发生在放行之后,那时必须留着句柄等
// `turn.completed` 来结算——提前清掉会让宿主侧这一轮的候选永远没有人结算。
if (
rejection &&
@@ -642,9 +642,9 @@ export function useDirectProjectChatController({
pendingRunAnalyticsRef.current = null;
}
if (rejection) {
const notice = directTurnRejectionNotice(rejection);
const notice = directTurnEnqueueFailureNotice(rejection);
if (notice) {
// 认得的前置条件 / 参数类拒单:写成与用户消息同级的提示,不占状态行、不写运行错误、
// 认得的前置条件 / 参数 / 队列已满类入队失败:写成与用户消息同级的提示,不占状态行、不写运行错误、
// 也不上报(用户自己就能改,上报只会变成噪声)。
if (projectPathRef.current === nextProjectPath) {
onRuntimeError('');
@@ -652,7 +652,7 @@ export function useDirectProjectChatController({
role: 'assistant',
text: notice,
runtimeOwned: true,
messageId: directTurnRejectionNoticeMessageId(
messageId: directTurnEnqueueFailureNoticeMessageId(
directCodexConversationMessageId(input.clientTurnId, 'user'),
),
updatedAt: Date.now(),
@@ -662,12 +662,12 @@ export function useDirectProjectChatController({
}
}
if (projectPathRef.current !== nextProjectPath) return false;
// 认不出的拒单(宿主 / 环境事实)与其它非结构化错误走同一条通道:上报 + 横幅。
// 认不出的入队失败(宿主 / 环境事实)与其它非结构化错误走同一条通道:上报 + 横幅。
void captureAgentRuntimeError(error, DIRECT_CODEX_AGENT_ID);
// 拒单文案优先:它是宿主生成的唯一一份(`Display` 或脱敏收口文案),比 `Error` 的形状更可信;
// 拒单不带阶段标签(这一轮没有开始),所以走拒单那一份映射。
// 入队失败文案优先:它是宿主生成的唯一一份(`Display` 或脱敏收口文案),比 `Error` 的形状更可信;
// 入队失败不带阶段标签(这一轮没有开始),所以走入队失败那一份映射。
const visibleMessage = rejection
? directTurnUnrecognizedRejectionNoticeText(rejection)
? directTurnUnrecognizedEnqueueFailureNoticeText(rejection)
: projectRuntimeVisibleError(
error instanceof Error ? error.message : String(error),
'陶泥儿智能创作',
@@ -675,15 +675,15 @@ export function useDirectProjectChatController({
);
if (projectPathRef.current !== nextProjectPath) return false;
// 回合失败的说明不由这里写:宿主已经把它放进了 `turn.completed.failure`,reducer 会把它落成
// 本轮最后一条条目(唯一来源)。**拒单没有这条出口**——拒单不产生回合事件,聊天里那条乐观
// 用户气泡后面永远不会再有任何说明,所以这里必须补一条同级提示;上报与横幅照旧保留。
// 非结构化错误同样不写聊天:它可能发生在接单之后,说明由事件流负责。
// 本轮最后一条条目(唯一来源)。**入队失败没有这条出口**——入队失败不产生回合事件,聊天里
// 那条乐观用户气泡后面永远不会再有任何说明,所以这里必须补一条同级提示;上报与横幅照旧保留。
// 非结构化错误同样不写聊天:它可能发生在放行之后,说明由事件流负责。
if (rejection) {
appendLocalMessage({
role: 'assistant',
text: visibleMessage,
runtimeOwned: true,
messageId: directTurnRejectionNoticeMessageId(
messageId: directTurnEnqueueFailureNoticeMessageId(
directCodexConversationMessageId(input.clientTurnId, 'user'),
),
updatedAt: Date.now(),
@@ -15,8 +15,8 @@ import type {
* - `nativeRunning`:**原生真相**。只由订阅 reducer 的 `turnRunning` 给出(`turn.started`
* 已到、`turn.completed` 未到)。
* - `commandInFlight`:**本地真相**。本次会话是否有一条本地在飞的回合(写权限门 → invoke →
* 宿主认领这一轮);它从按下发送那一刻就为真,与原生是否已经开始无关。接单化之后命令只
* 等到接单就返回,所以它不等于"命令还没返回"——它活到宿主那一轮的开始事件被观察到为止。
* 宿主放行这一轮);它从按下发送那一刻就为真,与原生是否已经开始无关。入队化之后命令只
* 等到入队就返回,所以它不等于"命令还没返回"——它活到宿主那一轮的开始事件被观察到为止。
* - `displayBusy`:header / composer / 「陶泥儿正在处理」卡片该读的忙态,就是两者的并集:
* 只要有一条成立就不能再接受新的发送。卡片读它而不是 `nativeRunning`:`turn.started`
* 要等宿主应答返回才发出,只认原生真相会让模型首 token 之前那十来秒没有任何「正在处理」
@@ -2,7 +2,7 @@ import type { ChatMessage } from '../../../../app/types';
import { projectRuntimeVisibleRejectionError } from '../../../../features/agent-runtime';
import type { HomeCreationType } from '../../../home';
import type { DirectCodexUserItem } from '../generated/DirectCodexUserItem';
import type { DirectTurnRejection } from '../generated/DirectTurnRejection';
import type { DirectTurnEnqueueFailure } from '../generated/DirectTurnEnqueueFailure';
export const DIRECT_CODEX_AGENT_ID = 'direct-codex';
export const DIRECT_CODEX_CONVERSATION_MESSAGE_ID_PREFIX = 'direct-codex:';
@@ -63,14 +63,14 @@ export function directCodexConversationMessageId(
}
/**
* 命令的**拒单**载荷(Rust 侧 `DirectTurnRejection`):结构化变体 + 宿主生成的文案。
* 命令的**入队失败**载荷(Rust 侧 `DirectTurnEnqueueFailure`):结构化变体 + 宿主生成的文案。
*
* `invoke` 拒绝时拿到的就是这份值(不是 `Error`)。这里只做一次形状读取,分流一律看
* `error.type`——文案是给人看的,不参与任何判断。
*/
export function readDirectTurnRejection(
export function readDirectTurnEnqueueFailure(
error: unknown,
): DirectTurnRejection | null {
): DirectTurnEnqueueFailure | null {
if (!error || typeof error !== 'object') return null;
const candidate = error as { error?: unknown; message?: unknown };
const variant = candidate.error;
@@ -78,12 +78,12 @@ export function readDirectTurnRejection(
const type = (variant as { type?: unknown }).type;
if (typeof type !== 'string' || !type) return null;
if (typeof candidate.message !== 'string') return null;
return candidate as DirectTurnRejection;
return candidate as DirectTurnEnqueueFailure;
}
/**
* 认得的拒单(前置条件不满足 / 用户参数无效)→ 与用户消息同级的提示文案;认不得的返回 `null`,
* 由调用方原样抛出交给既有捕获链路(上报 + 横幅)。
* 认得的入队失败(前置条件不满足 / 用户参数无效 / 队列已满)→ 与用户消息同级的提示文案;
* 认不得的返回 `null`,由调用方原样抛出交给既有捕获链路(上报 + 横幅)。
*
* 文案是宿主 `Display` 生成的**唯一一份**,界面原样显示:不套运行错误映射,也不在界面另写一份
* ——那一套会把"聊天内容不能为空"这类前置条件压成"执行失败,请稍后重试"。
@@ -95,49 +95,47 @@ export function readDirectTurnRejection(
* 时返回 `null`,而不是让调用方拿到一条空串(`''` 显示不出任何东西,却会被按 `!== null` 判据的
* 调用方当成"有提示")。
*/
export function directTurnRejectionNotice(
rejection: DirectTurnRejection,
export function directTurnEnqueueFailureNotice(
failure: DirectTurnEnqueueFailure,
): string | null {
switch (rejection.error.type) {
switch (failure.error.type) {
case 'clientTurnIdMissing':
case 'clientTurnIdMalformed':
case 'queueFull':
case 'turnAlreadyRunning':
case 'projectRootUnanchored':
case 'projectRootUnusable':
case 'permissionRejected':
case 'inputRejected':
case 'contentEmpty':
return rejection.message.trim() || null;
return failure.message.trim() || null;
default:
return null;
}
}
/**
* 拒绝提示在同一条用户消息里的展示身份:与失败说明(`:failure`)同一套派生规则但不同后缀,
* 入队失败提示在同一条用户消息里的展示身份:与失败说明(`:failure`)同一套派生规则但不同后缀,
* 两条通道永远不会合并成一条。
*/
export function directTurnRejectionNoticeMessageId(userItemId: string) {
export function directTurnEnqueueFailureNoticeMessageId(userItemId: string) {
return `${userItemId}:rejected`;
}
/**
* **认不出的**拒单(宿主 / 环境事实)在聊天区末尾自成一组提示的文案。
* **认不出的**入队失败(宿主 / 环境事实)在聊天区末尾自成一组提示的文案。
*
* 这两类拒单不产生 `turn.completed`(拒单没有接单),而本地也不再造用户气泡,所以这一轮在聊天区里
* 本来什么都不剩——说明只能由命令边界补一条,否则用户只看得到一条会消失的横幅。它带自己的身份
* (`…:rejected`),投影据此自成一组,不挂进上一轮。上报与横幅照旧保留:两件事不是同一份
* 这两类失败不产生 `turn.completed`(入队失败的这一轮从未成立),而本地也不再造用户气泡,所以
* 这一轮在聊天区里本来什么都不剩——说明只能由命令边界补一条,否则用户只看得到一条会消失的横幅。
* 它带自己的身份(`…:rejected`),投影据此自成一组,不挂进上一轮。上报与横幅照旧保留:两件事不是同一份
* (一个是给用户看的话,一个是把现场送进上报池与 `.agent/runtime/errors`)。
*
* 文案不能原样用宿主给的 `message`:这类拒单的 `message` 是宿主的收口文案(带 `stage=` / `code=`
* 文案不能原样用宿主给的 `message`:这类失败的 `message` 是宿主的收口文案(带 `stage=` / `code=`
* 这类机器字段),先过与失败说明同一份可见文案映射再进聊天;映射认不出形状时给一句通用兜底,
* 绝不把内部字段塞进聊天。
*/
export function directTurnUnrecognizedRejectionNoticeText(
rejection: DirectTurnRejection,
export function directTurnUnrecognizedEnqueueFailureNoticeText(
failure: DirectTurnEnqueueFailure,
): string {
return projectRuntimeVisibleRejectionError(
rejection.message,
'陶泥儿智能创作',
);
return projectRuntimeVisibleRejectionError(failure.message, '陶泥儿智能创作');
}
@@ -89,7 +89,7 @@ export type DirectThreadChatState = {
* 已经收口的回合数(单调递增,项目切换时随整份状态重置)。
*
* 它是"回合完成"这个事实**唯一的计数**,给上层放行发送队列与结算埋点用。为什么要计数
* 而不是看 `turnRunning` 的下降沿:一轮可能在**同一次 consume** 里开始并结束(接单后
* 而不是看 `turnRunning` 的下降沿:一轮可能在**同一次 consume** 里开始并结束(放行后
* 立刻失败),那时 `turnRunning` 从头到尾没有被观察到真,下降沿永远不会来。计数是状态,
* 批量到达也一样看得见。
*/
@@ -324,7 +324,7 @@ export function reduceDirectThreadEvent(
const eventAt = readDirectThreadEventAt(event);
// 回合身份只认**事件流顺序**,不拿时间戳大小当身份:原生回合时间是秒级精度、
// 宿主收口时间可能带毫秒,"上一轮结束之后又来一条 turn.started"就是新回合,
// 哪怕它落在同一秒。开始事件由 Thread Manager 在接单时发一次,线上不再有同一轮的
// 哪怕它落在同一秒。开始事件由 Thread Manager 在放行时发一次,线上不再有同一轮的
// 重复起点,所以这里不做"保留第一次"的兼容。
const turnStartedAt = eventAt;
// 本轮的 canonical user identity 跟着事件走:新回合就换成新的;旧原生不带身份时
@@ -6,7 +6,7 @@
*
* **回合归属只认身份**:开口用户条目的 canonical `itemId`(`direct-codex:{clientTurnId}:user`)
* 就是这一轮的回合身份,本轮用户条目与失败说明按它归进同一轮。本地只保留"说明"类消息
* (拒单提示、壳层 `announce`),它们不是回合条目,也不参与回合身份。
* (入队失败提示、壳层 `announce`),它们不是回合条目,也不参与回合身份。
*/
import type { ChatMessage } from '../../../../app/types';
@@ -64,11 +64,11 @@ export type DirectChatBlock =
* 回显条目的 `at` 还是宿主落盘 / 观测时间(晚于用户真实发送)。
*
* 这里曾经有过第三个态 `awaiting-start`(本地已发出、宿主还没认领):它只服务本地乐观用户气泡。
* 气泡已按"回合只认宿主条目"删除,所以"接单窗口"不再是一种展示态——窗口期聊天区里没有这一轮的
* 任何条目,只有 composer 的忙态与状态行(见 ADR「DirectProject命令接单化」的后续更新)。
* 气泡已按"回合只认宿主条目"删除,所以"入队窗口"不再是一种展示态——窗口期聊天区里没有这一轮的
* 任何条目,只有 composer 的忙态与状态行(见 ADR「DirectProject命令入队化」)。
*
* 两个容易读错的地方:
* - 命令返回 **≠** 这一轮结束了:`invoke` 只等到接单,宿主确认这一轮靠的是开始事件;
* - 命令返回 **≠** 这一轮结束了:`invoke` 只等到入队,宿主确认这一轮靠的是开始事件;
* 所以"命令还没回来"不是任何展示态依据,只有显式事件是。
* - `endedAt === 0` **≠** 还在跑:历史回合没有边界元数据(`turnEndedAt` 只是会话内展示
* 缓存),它们必须落 `finished`。
@@ -105,7 +105,7 @@ export type DirectChatTurn = {
type DirectChatTurnEntries = {
key: string;
entries: DirectChatEntry[];
/** 本地说明(拒单提示、壳层 `announce`):不是回合条目,按挂点归组。 */
/** 本地说明(入队失败提示、壳层 `announce`):不是回合条目,按挂点归组。 */
notices: ChatMessage[];
/** 原生回合是否在跑;只有一个来源——reducer 的 `turnRunning`。 */
nativeRunning: boolean;
@@ -157,7 +157,7 @@ function blockFromLocalMessage(
): DirectChatBlock | null {
const text = message.text.trim();
if (!text) return null;
// 本地消息只剩"说明"一种(拒单提示、壳层 `announce`):用户消息一律来自宿主条目。
// 本地消息只剩"说明"一种(入队失败提示、壳层 `announce`):用户消息一律来自宿主条目。
return {
kind: 'assistant',
key: message.messageId ?? `local:${index}`,
@@ -198,7 +198,7 @@ function directEntryTurnKey(entry: DirectChatEntry): string {
* 其余条目(工具、思考、正文)不带身份,跟着当前回合走。
*
* 本地消息只剩说明一种,用户消息一律来自宿主条目(本地乐观气泡已删,见 ADR「DirectProject命令
* 接单化」的后续更新):带身份的说明(拒单提示,id 形如 `…:user:rejected`)说的是"这一轮从没
* 入队化」的后续更新):带身份的说明(入队失败提示,id 形如 `…:user:rejected`)说的是"这一轮从没
* 成立过",不属于任何回合,自己在会话末尾开一组;不带头身份的是壳层 `announce`,挂到当前回合末
* 尾。同身份的本地消息不重复渲染:条目赢。
*
@@ -269,7 +269,7 @@ export function buildDirectChatTurns({
// 本地不再造用户消息(乐观气泡已删):万一还来了一条,既不进聊天区、也不开回合。
if (message.role === 'user') return;
if (message.messageId && entryIds.has(message.messageId)) return;
// 带身份的本地说明(拒单提示)说的是"这一轮从没成立过":它不属于任何回合,也不能按位置挂进
// 带身份的本地说明(入队失败提示)说的是"这一轮从没成立过":它不属于任何回合,也不能按位置挂进
// 上一轮(那会被读成上一轮的问题)。它自己在会话末尾开一组提示:没有用户条目的回合不会凭空
// 多出一条耗时文案(`startedAt` 拿不到时终态文案整条隐藏)。
if (message.messageId) {
@@ -3,7 +3,7 @@
/**
* 失败发生在交付的哪一段。与错误分类正交:分类说明"怎么回事",阶段说明"走到哪一步"。
*
* 线上取值跟着拒单 / 失败载荷一起给前端(`art-preparation` 这类),所以也要导出。
* 线上取值跟着入队失败 / 回合失败载荷一起给前端(`art-preparation` 这类),所以也要导出。
*/
export type DirectCodexFailureStage =
| 'art-preparation'
@@ -2,7 +2,7 @@
import type { DirectCodexNativeKind } from './DirectCodexNativeKind';
/**
* 模型调用失败(app-server 一次 `turn` 的结果)的分类,跟着拒单 / 失败载荷一起给前端。
* 模型调用失败(app-server 一次 `turn` 的结果)的分类,跟着入队失败 / 回合失败载荷一起给前端。
*
* 每个变体对应平台层 `LlmError` 的一个分支,于是 [`DirectTurnError::wire_kind`] 的取值与改造前
* 完全一致:事件的 `failure.kind` 就是这一份取值,界面按它选语气,不拿它做流程分支。
@@ -3,7 +3,7 @@
/**
* 一条待发消息离开队列的原因。
*
* typed 枚举,取值即语义:取消是用户在输入盒上撤掉这条消息,放行是它已经接单并成为回合
* typed 枚举,取值即语义:取消是用户在输入盒上撤掉这条消息,放行是它已经离开队列并成为回合
* (同一临界区里另有 `turn.started`)。界面按它分流,不解析字符串。
*/
export type DirectQueueRemovalReason = 'cancelled' | 'dispatched';
@@ -35,7 +35,7 @@ export type DirectThreadEvent =
| {
type: 'turn.started';
/**
* 本轮开始的阶段时间(毫秒):**接单**那一刻的宿主毫秒钟(逻辑回合的起点,不是
* 本轮开始的阶段时间(毫秒):**放行**那一刻的宿主毫秒钟(逻辑回合的起点,不是
* `turn/start` 的时刻)。
*/
at?: number;

Some files were not shown because too many files have changed in this diff Show More