线程管理器归位 agent/thread_manager:去掉 Direct 前缀,按职责拆文件

- 四个 direct_* 文件移入 agent/thread_manager:wire.rs(运行态事件契约)、queue.rs(待发消息队列规则)、dispatch.rs(放行与占用)、mod.rs(管理器主体与对外入口)
- 类型去掉 Direct 前缀:ThreadEvent / ThreadManager / ThreadItem / ThreadFileChange / ThreadDeltaKind / ThreadRequestKind / PendingTurn / QueueRemovalReason / QueueRemovalOutcome / DispatchedTurn / TurnReservation / TurnIdentity / StaleTurnRelease / ConsumeResult / SubscriptionBootstrap / ActiveTurn / ActiveTurnSnapshot
- 对外函数改名:subscribe_thread / consume_thread / unsubscribe_thread / append_thread_event / enqueue_pending_turn / remove_pending_turn / claim_pending_turn / update_active_turn / list_active_turns / complete_turn / complete_turn_if_reserved / active_turn_id_at / read_turn_identity / stale_turn_for_release / thread_id_for_project
- 生成绑定跟着改名(DirectThreadEvent→ThreadEvent 等十个文件为 rename),前端八处 import 同步;Tauri 命令名与前端 IPC 契约不变
- 定向验证:cargo test agent::(945 passed,另两条并发用例为既有 flaky、重跑通过)、npx vitest run appSurface/directThreadChat/directHistoryPaging(247 passed)、typecheck 通过
This commit is contained in:
2026-09-30 15:47:22 +08:00
parent e70b338900
commit f121257cca
40 changed files with 747 additions and 867 deletions
@@ -29,13 +29,9 @@ mod direct_project_context;
mod direct_project_history;
mod direct_project_turn_history;
mod direct_runtime;
mod direct_thread_manager;
mod direct_thread_queue;
mod direct_thread_wire;
mod direct_tool_bridge;
mod direct_tool_calls;
mod direct_tools_mcp;
mod direct_turn_dispatch;
mod direct_turn_error;
mod direct_turn_failure;
mod direct_turn_stream;
@@ -50,6 +46,8 @@ mod runtime_protocol;
mod runtime_state;
mod runtime_tools;
mod skill_pack;
mod thread_manager;
use claude_code_cli::*;
pub(crate) use claude_code_cli::{
cancel_direct_claude_code_turn_at, direct_game_creator_claude_code_chat_at,
@@ -59,7 +57,7 @@ pub(crate) use claude_code_cli::{
use codex_app_server::*;
pub(crate) use codex_app_server::{
cancel_direct_codex_turn_at, direct_game_creator_codex_chat_at,
direct_game_creator_home_codex_chat, direct_thread_id_for_project, DirectTurnCancelView,
direct_game_creator_home_codex_chat, thread_id_for_project, DirectTurnCancelView,
};
use codex_cli::*;
pub(crate) use codex_cli::{
@@ -72,13 +70,9 @@ pub(crate) use direct_codex_user_item::*;
pub(crate) use direct_project_history::*;
pub(crate) use direct_project_turn_history::*;
pub(crate) use direct_runtime::*;
pub(crate) use direct_thread_manager::*;
pub(crate) use direct_thread_queue::*;
pub(crate) use direct_thread_wire::*;
pub(crate) use direct_tool_bridge::*;
pub(crate) use direct_tool_calls::*;
pub(crate) use direct_tools_mcp::*;
pub(crate) use direct_turn_dispatch::*;
pub(crate) use direct_turn_error::*;
pub(crate) use direct_turn_failure::*;
pub(crate) use direct_turn_stream::*;
@@ -93,6 +87,10 @@ pub(crate) use runtime_protocol::*;
pub(crate) use runtime_state::*;
pub(crate) use runtime_tools::*;
pub(crate) use skill_pack::*;
pub(crate) use thread_manager::dispatch::*;
pub(crate) use thread_manager::queue::*;
pub(crate) use thread_manager::wire::*;
pub(crate) use thread_manager::*;
pub(crate) fn shutdown_game_creator_codex_app_servers() -> Result<(), String> {
shutdown_game_creator_codex_app_servers_impl()
@@ -30,7 +30,7 @@ pub(crate) fn direct_codex_canonical_project_identity(
/// `direct_codex_canonical_project_identity`)的要求,而线程 id 只是一个路径 key。把
/// manifest 的瞬时抖动混进线程 id,会让同一项目在"订阅那一刻"与"跑回合那一刻"算出两个
/// 字符串(例如调用方给的是符号链接路径),订阅就绑到一条永远不会有事件的空线程上。
pub(crate) fn direct_thread_id_for_project(root: &std::path::Path) -> String {
pub(crate) fn thread_id_for_project(root: &std::path::Path) -> String {
let Ok((canonical_root, _)) = resolve_direct_codex_project_authority(root) else {
return root.to_string_lossy().into_owned();
};
@@ -16,7 +16,7 @@ use process_tree::{OwnedProcessTree, ProcessTreeExitProof};
mod direct_project_identity;
mod execution;
mod model_catalog;
pub(crate) use direct_project_identity::direct_thread_id_for_project;
pub(crate) use direct_project_identity::thread_id_for_project;
use direct_project_identity::*;
use execution::ExecutionAdapter;
@@ -631,7 +631,7 @@ enum CodexTurnEvent {
params: serde_json::Value,
},
Request {
kind: DirectThreadRequestKind,
kind: ThreadRequestKind,
params: serde_json::Value,
},
RawItem(serde_json::Value),
@@ -808,8 +808,8 @@ fn direct_codex_safe_activity_for_item_value(item: &serde_json::Value) -> &'stat
fn direct_thread_event_item(
root: &std::path::Path,
item: &serde_json::Value,
) -> Option<DirectThreadItem> {
direct_thread_item_from_value(root, item, direct_tool_call_now_ms())
) -> Option<ThreadItem> {
thread_item_from_value(root, item, direct_tool_call_now_ms())
}
/// 运行态条目投影:Codex 回显的用户消息整条跳过。
@@ -824,7 +824,7 @@ fn direct_thread_event_item(
fn direct_thread_visible_item(
root: &std::path::Path,
item: &serde_json::Value,
) -> Option<DirectThreadItem> {
) -> Option<ThreadItem> {
if is_direct_project_codex_user_item(item) {
return None;
}
@@ -844,14 +844,14 @@ fn direct_thread_visible_item(
pub(crate) fn emit_direct_thread_user_item(
root: &std::path::Path,
item: &serde_json::Value,
) -> Option<DirectThreadItem> {
) -> Option<ThreadItem> {
let entry_item = direct_thread_event_item(root, item)?;
// 条目时间是落盘 / 观测时间。前端乐观用户气泡已删(ADR「DirectProject命令入队化」),
// 所以这就是界面显示这条用户消息的唯一时间口径:它晚于用户按下发送,但不再有第二份更早的时间。
let at = entry_item.at();
append_direct_thread_event(
&direct_thread_id_for_project(root),
DirectThreadEvent::item_completed(entry_item.clone(), at),
append_thread_event(
&thread_id_for_project(root),
ThreadEvent::item_completed(entry_item.clone(), at),
);
Some(entry_item)
}
@@ -1080,21 +1080,21 @@ fn direct_codex_safe_activity_for_notification(method: &str) -> Option<&'static
}
}
fn direct_codex_request_event_type(method: &str) -> Option<DirectThreadRequestKind> {
fn direct_codex_request_event_type(method: &str) -> Option<ThreadRequestKind> {
match method {
"item/fileChange/requestApproval"
| "item/commandExecution/requestApproval"
| "item/permissions/requestApproval" => Some(DirectThreadRequestKind::ApprovalRequested),
| "item/permissions/requestApproval" => Some(ThreadRequestKind::ApprovalRequested),
"item/tool/requestUserInput" | "item/mcpToolCall/requestUserInput" => {
Some(DirectThreadRequestKind::AskRequested)
Some(ThreadRequestKind::AskRequested)
}
_ => None,
}
}
fn direct_codex_resolution_event_type(method: &str) -> Option<DirectThreadRequestKind> {
fn direct_codex_resolution_event_type(method: &str) -> Option<ThreadRequestKind> {
match method {
"serverRequest/resolved" => Some(DirectThreadRequestKind::RequestResolved),
"serverRequest/resolved" => Some(ThreadRequestKind::RequestResolved),
_ => None,
}
}
@@ -1197,10 +1197,10 @@ fn direct_codex_reasoning_delta_event(
fn direct_codex_thread_delta_event(
root: &std::path::Path,
item_id: String,
kind: DirectThreadDeltaKind,
kind: ThreadDeltaKind,
delta: &str,
) -> DirectThreadEvent {
DirectThreadEvent::item_delta(item_id, kind, direct_thread_delta_text(root, delta))
) -> ThreadEvent {
ThreadEvent::item_delta(item_id, kind, thread_delta_text(root, delta))
}
/// 通知 → 回合事件的唯一分类函数:运行态读取器与单测共用这一份。
///
@@ -3616,14 +3616,14 @@ impl CodexAppServerConnection {
}
}
turn_start_guard.armed = false;
let direct_thread_id = direct_thread_id_for_project(history_root);
let direct_thread_id = thread_id_for_project(history_root);
// 本轮开口用户条目的 canonical id:只从已落盘的那条条目上读身份(`id`,工具条目才用
// `call_id`),不在事件侧重造一份。拿不到就留空,让前端按"归属不可证明"处理。
let direct_turn_user_item_id = direct_turn_user_item
.as_ref()
.and_then(direct_thread_item_identity);
.and_then(thread_item_identity);
// 逻辑回合的**边界**不在这里:开始事件由认领动作发出、兜底由放行占用对象持有
// (`direct_turn_dispatch.rs`)。本轮的用户条目也不在这里下发——发点在放行之后、
// (`thread_manager::dispatch`)。本轮的用户条目也不在这里下发——发点在放行之后、
// 起 codex 之前(`emit_direct_thread_user_item`),见那条注释。
//
// 下面这个毫秒钟与逻辑回合无关,只服务模型终态的**完成时刻**:上游 Turn 的
@@ -3713,14 +3713,14 @@ impl CodexAppServerConnection {
Some(CodexTurnEvent::AgentMessageDelta { item_id, delta }) => {
if self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject {
direct_project_history.observe_delta(&item_id, &delta);
append_direct_thread_event(
append_thread_event(
&direct_thread_id,
// 事件自足:增量自带 item 身份与正文类别(正文 / 思考),
// 前端 reducer 不允许靠猜 itemId 的来源决定 kind。
direct_codex_thread_delta_event(
history_root,
item_id.clone(),
DirectThreadDeltaKind::Message,
ThreadDeltaKind::Message,
&delta,
),
);
@@ -3755,12 +3755,12 @@ impl CodexAppServerConnection {
}
Some(CodexTurnEvent::ReasoningDelta { item_id, delta }) => {
if self.inner.workspace_mode == CodexAppServerWorkspaceMode::DirectProject {
append_direct_thread_event(
append_thread_event(
&direct_thread_id,
direct_codex_thread_delta_event(
history_root,
item_id,
DirectThreadDeltaKind::Reasoning,
ThreadDeltaKind::Reasoning,
&delta,
),
);
@@ -3797,9 +3797,9 @@ impl CodexAppServerConnection {
if let Some(entry_item) = entry_item {
// `rawResponseItem/completed` 不带阶段时间,宿主处理到这条
// 通知的钟就是该阶段唯一可证明的时间。
append_direct_thread_event(
append_thread_event(
&direct_thread_id,
DirectThreadEvent::item_completed(
ThreadEvent::item_completed(
entry_item,
direct_tool_call_now_ms(),
),
@@ -3815,9 +3815,9 @@ impl CodexAppServerConnection {
.or_else(|| params.get("id").and_then(serde_json::Value::as_str))
.filter(|value| !value.is_empty())
.map(str::to_string);
append_direct_thread_event(
append_thread_event(
&direct_thread_id,
DirectThreadEvent::request(kind, request_id),
ThreadEvent::request(kind, request_id),
);
}
}
@@ -3930,11 +3930,11 @@ impl CodexAppServerConnection {
{
// `item/started` 的通知层带 `startedAtMs`:这是工具真正
// 开始的阶段时间,优先于条目展示时间与宿主钟。
append_direct_thread_event(
append_thread_event(
&direct_thread_id,
DirectThreadEvent::item_started(
ThreadEvent::item_started(
entry_item,
direct_thread_item_event_at_ms(
thread_item_event_at_ms(
&params,
item,
false,
@@ -3984,7 +3984,7 @@ impl CodexAppServerConnection {
{
model_terminal = Some((
status.to_string(),
direct_thread_turn_completed_at_ms(
thread_turn_completed_at_ms(
turn,
Some(direct_turn_started_at_ms),
direct_tool_call_now_ms(),
@@ -4311,7 +4311,7 @@ impl DirectTurnTerminalContext {
history_root,
);
// 终态走 Thread Manager 的深出口:解除这一轮的占用并写下 `turn.completed`。
complete_direct_thread_turn(
complete_turn(
&self.thread_id,
terminal.event(self.completed_at, self.user_item_id.as_deref()),
);
@@ -4502,7 +4502,7 @@ enum DirectCodexTurnCancelTarget {
/// app-server 侧还有活句柄:正常发 `turn/interrupt`。
Interrupt(Arc<CodexTurnStartCancellation>),
/// app-server 侧已经拿不到可中断的活句柄;带上是哪种情况。
Stale(DirectStaleTurnReleaseReason),
Stale(StaleTurnReleaseReason),
}
/// 终止当前项目正在运行的 Direct 回合。
@@ -4515,7 +4515,7 @@ enum DirectCodexTurnCancelTarget {
/// 只发中断会让 Thread Manager 上这一轮的占用永远解不开,这个项目此后每条消息都会被
/// "已有另一条回合正在跑"挡在放行之外——这正是"重进会话被堵死"的死锁形态。这时显式写一条
/// `turn.completed{aborted}` 解除占用并把可读原因返回给界面。释放条件见
/// [`direct_stale_turn_for_release`] 的注释;"正在跑的是另一轮"仍然保持原拒绝语义,
/// [`stale_turn_for_release`] 的注释;"正在跑的是另一轮"仍然保持原拒绝语义,
/// 什么都不释放。
///
/// 兜底终态由 Thread Manager 在解除占用的同一个临界区里写下(带 `userItemId`,身份取被释放那一轮
@@ -4552,12 +4552,10 @@ pub(crate) fn cancel_direct_codex_turn_at(
Ok((_, cancellation)) if cancellation.app_server_alive() => {
DirectCodexTurnCancelTarget::Interrupt(Arc::clone(cancellation))
}
Ok(_) => {
DirectCodexTurnCancelTarget::Stale(DirectStaleTurnReleaseReason::ExecutorExited)
Ok(_) => DirectCodexTurnCancelTarget::Stale(StaleTurnReleaseReason::ExecutorExited),
Err(_) => {
DirectCodexTurnCancelTarget::Stale(StaleTurnReleaseReason::NeverReachedExecutor)
}
Err(_) => DirectCodexTurnCancelTarget::Stale(
DirectStaleTurnReleaseReason::NeverReachedExecutor,
),
}
};
match target {
@@ -4577,10 +4575,10 @@ pub(crate) fn cancel_direct_codex_turn_at(
// 校验与释放**同一个临界区**(确实登记着一轮、身份对得上、过了启动窗口才允许判成残留;
// 占用与兜底终态要么一起落地,要么一个字都不写)。校验完了再单独释放不行:两步之间并发
// 认领可能已经登记了新那一轮的占用,无条件解除就会清掉它。
let released = direct_stale_turn_for_release(root, client_turn_id, reason)?;
let released = stale_turn_for_release(root, client_turn_id, reason)?;
// 这一轮不会再有人替它收尾(占用刚被兜底解除):踢一脚让队列继续。
// 少了这一脚,排在这条后面的待发消息要等到下一次用户动作才会被放行。
crate::agent::kick_direct_queue_dispatch(root);
crate::agent::kick_queue_dispatch(root);
Ok(DirectTurnCancelView {
outcome: DIRECT_TURN_CANCEL_OUTCOME_RELEASED.to_string(),
message: format!(
@@ -6004,10 +6002,7 @@ mod tests {
fn direct_thread_message_and_reasoning_deltas_are_sanitized_before_enqueue() {
let root = std::path::Path::new("/workspace/direct-project");
let raw = "key=sk-abcdefghijklmnop project=/workspace/direct-project/assets/a.png private=/root/secret.txt";
for kind in [
DirectThreadDeltaKind::Message,
DirectThreadDeltaKind::Reasoning,
] {
for kind in [ThreadDeltaKind::Message, ThreadDeltaKind::Reasoning] {
let event = direct_codex_thread_delta_event(root, "item-1".to_string(), kind, raw);
let wire = serde_json::to_string(&event).expect("serialize direct thread delta");
assert!(
@@ -6043,12 +6038,12 @@ mod tests {
];
let streamed: String = chunks
.iter()
.map(|chunk| direct_thread_delta_text(root, chunk))
.map(|chunk| thread_delta_text(root, chunk))
.collect();
let whole = chunks.concat();
assert_eq!(streamed, whole, "逐段脱敏不得吃掉段尾换行");
assert_eq!(
direct_thread_delta_text(root, &whole),
thread_delta_text(root, &whole),
streamed,
"逐段脱敏与整段脱敏必须同形"
);
@@ -6070,7 +6065,7 @@ mod tests {
"turn-1"
),
Some(CodexTurnEvent::Request {
kind: DirectThreadRequestKind::RequestResolved,
kind: ThreadRequestKind::RequestResolved,
..
})
));
@@ -6084,7 +6079,7 @@ mod tests {
"turn-1",
),
Some(CodexTurnEvent::Request {
kind: DirectThreadRequestKind::ApprovalRequested,
kind: ThreadRequestKind::ApprovalRequested,
..
})
));
@@ -6162,7 +6157,7 @@ mod tests {
};
let item = event_params.get("item").expect("item payload");
assert_eq!(
direct_thread_item_event_at_ms(&event_params, item, completed, 9_999),
thread_item_event_at_ms(&event_params, item, completed, 9_999),
expected_at_ms,
"{method} 必须用通知层的阶段时间,而不是宿主钟"
);
@@ -6185,7 +6180,7 @@ mod tests {
};
let turn = params.get("turn").unwrap_or(&params);
assert_eq!(
direct_thread_turn_completed_at_ms(turn, Some(1_700_000_000_500), 9_999),
thread_turn_completed_at_ms(turn, Some(1_700_000_000_500), 9_999),
9_999,
"上游只有秒级 completedAt:不采用,取宿主处理终态的毫秒钟"
);
@@ -6198,7 +6193,7 @@ mod tests {
"durationMs": 42_500u64,
});
assert_eq!(
direct_thread_turn_completed_at_ms(&with_duration, Some(1_700_000_000_500), 9_999),
thread_turn_completed_at_ms(&with_duration, Some(1_700_000_000_500), 9_999),
1_700_000_043_000,
"durationMs + 宿主高精度起点才派生结束"
);
@@ -6830,7 +6825,7 @@ mod tests {
let stable_link = temp.path().join("current-project");
symlink(&project, &stable_link).expect("project symlink");
let canonical = direct_thread_id_for_project(&project);
let canonical = thread_id_for_project(&project);
assert_eq!(
canonical,
project
@@ -6840,7 +6835,7 @@ mod tests {
"线程 id 就是 canonical 路径"
);
assert_eq!(
direct_thread_id_for_project(&stable_link),
thread_id_for_project(&stable_link),
canonical,
"符号链接必须归一到同一个线程 id"
);
@@ -6852,7 +6847,7 @@ mod tests {
"前提:manifest 读不到时权威身份确实会失败"
);
assert_eq!(
direct_thread_id_for_project(&stable_link),
thread_id_for_project(&stable_link),
canonical,
"manifest 读失败不能把线程 id 退回调用方原始字符串"
);
@@ -8181,10 +8176,9 @@ done
"id": "direct-codex:turn-0001:user",
"content": [{ "type": "input_text", "text": "请创建菜单" }]
});
let thread_id = direct_thread_id_for_project(&project);
let bootstrap = crate::agent::subscribe_direct_thread(&thread_id);
let _reservation =
crate::agent::DirectTurnReservation::accept_for_test(&thread_id, "turn-0001");
let thread_id = thread_id_for_project(&project);
let bootstrap = crate::agent::subscribe_thread(&thread_id);
let _reservation = crate::agent::TurnReservation::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`):
@@ -8242,13 +8236,13 @@ done
.await
.expect("the turn must finish after the connection is reclaimed");
let consumed = crate::agent::consume_direct_thread(&bootstrap.subscription_id)
.expect("consume events");
let consumed =
crate::agent::consume_thread(&bootstrap.subscription_id).expect("consume events");
let terminal = consumed
.events
.iter()
.find_map(|event| match event {
DirectThreadEvent::TurnCompleted {
ThreadEvent::TurnCompleted {
status, failure, ..
} => Some((status.clone(), failure.clone())),
_ => None,
@@ -8333,8 +8327,8 @@ done
"id": "direct-codex:turn-0001:user",
"content": [{"type": "input_text", "text": "请创建菜单"}]
});
let thread_id = direct_thread_id_for_project(&project);
let bootstrap = crate::agent::subscribe_direct_thread(&thread_id);
let thread_id = thread_id_for_project(&project);
let bootstrap = crate::agent::subscribe_thread(&thread_id);
assert!(
bootstrap.events.is_empty(),
"订阅发生在回合之前,bootstrap 必须为空"
@@ -8342,8 +8336,7 @@ done
// 生产入口(`enqueue_direct_codex_turn` + 放行侧)在起 codex 之前先放行(认领时登记调用身份),
// 再把用户条目落盘:
// 这里补上同一步,于是这一轮的边界仍在同一个订阅里成对出现,历史里也有那条用户消息。
let _reservation =
crate::agent::DirectTurnReservation::accept_for_test(&thread_id, "turn-0001");
let _reservation = crate::agent::TurnReservation::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`):
@@ -8382,8 +8375,8 @@ done
.expect("run direct-project turn");
drop(observer);
let consumed = crate::agent::consume_direct_thread(&bootstrap.subscription_id)
.expect("consume events");
let consumed =
crate::agent::consume_thread(&bootstrap.subscription_id).expect("consume events");
// 回合起止必须与开口用户条目同源:前端在「只有锚点 + 历史、运行态为空」的回合里靠这个
// 身份把边界认领给同一条用户条目,缺了它就只能隐藏未知用时。
let lifecycle_user_item_ids = consumed
@@ -8392,10 +8385,10 @@ done
.filter(|event| {
matches!(
event,
DirectThreadEvent::TurnStarted { .. } | DirectThreadEvent::TurnCompleted { .. }
ThreadEvent::TurnStarted { .. } | ThreadEvent::TurnCompleted { .. }
)
})
.map(DirectThreadEvent::user_item_id)
.map(ThreadEvent::user_item_id)
.collect::<Vec<_>>();
assert_eq!(
lifecycle_user_item_ids,
@@ -8409,12 +8402,13 @@ done
let mut assistant_items = Vec::new();
for event in &consumed.events {
let item = match event {
DirectThreadEvent::ItemStarted { item, .. }
| DirectThreadEvent::ItemCompleted { item, .. } => item,
ThreadEvent::ItemStarted { item, .. } | ThreadEvent::ItemCompleted { item, .. } => {
item
}
_ => continue,
};
match item {
DirectThreadItem::Message {
ThreadItem::Message {
item_id,
role,
text,
@@ -8422,7 +8416,7 @@ done
} if role.as_str() == "user" => {
user_items.push((item_id.clone(), text.clone()));
}
DirectThreadItem::Message { text, .. } if !text.trim().is_empty() => {
ThreadItem::Message { text, .. } if !text.trim().is_empty() => {
assistant_items.push(text.clone());
}
_ => {}
@@ -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_active_turn_id_at(&root)?;
let turn = super::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(
@@ -53,7 +53,7 @@ fn context_identity(root: &Path) -> Result<ContextIdentity, String> {
let manifest = read_existing_manifest_for_project(root)?;
let revision = read_game_creator_agent_runtime_project_revision(root)?.revision;
let canonical = root.canonicalize().map_err(|_| "项目上下文目录无法解析")?;
let active = list_direct_active_turns()?.into_iter().find(|turn| {
let active = list_active_turns()?.into_iter().find(|turn| {
Path::new(&turn.project_path).canonicalize().ok().as_ref() == Some(&canonical)
});
Ok(ContextIdentity {
@@ -506,7 +506,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_active_turn_id_at(&scan_root).ok().as_deref() != Some(expected_turn.as_str()) {
if active_turn_id_at(&scan_root).ok().as_deref() != Some(expected_turn.as_str()) {
return Vec::new();
}
let candidates = [
@@ -621,10 +621,7 @@ mod tests {
std::fs::write(root.join("code.js"), "unchanged").unwrap();
// 身份来自逻辑回合(Thread Manager):放行才是"这一轮在跑"的唯一登记。
let owner = Arc::new(std::sync::Mutex::new(Some(
DirectTurnReservation::accept_for_test(
&direct_thread_id_for_project(&root),
"turn-before",
),
TurnReservation::accept_for_test(&thread_id_for_project(&root), "turn-before"),
)));
let swap = Arc::clone(&owner);
let result = read_batch_with(
@@ -635,8 +632,8 @@ mod tests {
let result = read_file(r, f, b);
let mut guard = swap.lock().unwrap();
drop(guard.take());
*guard = Some(DirectTurnReservation::accept_for_test(
&direct_thread_id_for_project(r),
*guard = Some(TurnReservation::accept_for_test(
&thread_id_for_project(r),
"turn-after",
));
result
@@ -802,10 +799,8 @@ mod tests {
async fn host_prefetch_keeps_data_out_of_system_rules_and_matches_active_turn() {
let (_temp, root) = project();
// 预取闸门与上下文身份读的是同一处:Thread Manager 上这一轮的活动回合登记。
let _turn = DirectTurnReservation::accept_for_test(
&direct_thread_id_for_project(&root),
"prefetch-turn",
);
let _turn =
TurnReservation::accept_for_test(&thread_id_for_project(&root), "prefetch-turn");
let data = prefetch_turn_input(&root, "prefetch-turn")
.await
.unwrap()
@@ -1,5 +1,5 @@
use super::direct_thread_item_identity;
use super::runtime_actions::acquire_game_creator_agent_runtime_project_write_lock_with_wait;
use super::thread_item_identity;
use crate::config::prepare_game_creator_private_path_for_read;
use crate::project::{
append_jsonl_line_unlocked, enforce_project_permission_policy, project_append_lock_for,
@@ -610,7 +610,7 @@ fn direct_project_history_anchor_id(item: &Value) -> Option<String> {
.map(str::trim)
.filter(|value| !value.is_empty())
.map(str::to_string)
.or_else(|| direct_thread_item_identity(item))
.or_else(|| thread_item_identity(item))
}
/// 一屏历史窗口在新端(较新一侧)的锚点。
@@ -719,7 +719,7 @@ pub(crate) fn read_direct_project_history_items_slice_at(
let recorded_at_ms = newest_first
.iter()
.filter_map(|(item, at)| {
let identity = direct_thread_item_identity(item)?;
let identity = thread_item_identity(item)?;
(*at > 0).then_some((identity, *at))
})
.collect();
@@ -458,8 +458,8 @@ fn direct_taonier_regeneration_invocation_sha256(invocation_id: &str) -> String
///
/// 拿不到身份时返回带前缀的错误:调用方(付费美术执行会话、工具桥、校验预约、上下文预取)都靠
/// "拿不到身份"来拒绝"在没有回合的情况下创建 / 清理生成账本"。
pub(crate) fn direct_active_turn_id_at(root: &Path) -> Result<String, String> {
read_direct_turn_identity(&direct_thread_id_for_project(root))
pub(crate) fn active_turn_id_at(root: &Path) -> Result<String, String> {
read_turn_identity(&thread_id_for_project(root))
.map(|identity| identity.client_turn_id)
.ok_or_else(|| {
format!(
@@ -472,7 +472,7 @@ pub(crate) fn direct_active_turn_id_at(root: &Path) -> Result<String, String> {
///
/// 上限用来挡"刚放行、还在本地准备(读配置 / manifest / 拼系统提示 / 开审计)的回合被终止误伤":
/// 这期间执行器侧还没有这一轮的登记,但这一轮马上就会进去,提前释放等于放开并发。
const DIRECT_STALE_TURN_RELEASE_MIN_AGE_MS: u64 = 60_000;
const 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
@@ -480,14 +480,14 @@ fn direct_taonier_active_now_millis() -> u64 {
/// "终止"拿不到可中断句柄时的分类,决定这一轮能不能被兜底释放。
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum DirectStaleTurnReleaseReason {
pub(crate) enum StaleTurnReleaseReason {
/// app-server 侧登记着这一轮,但执行进程已经退出:这一轮不可能再有收尾。
ExecutorExited,
/// app-server 侧完全没有这一轮的登记:只有过了正常启动窗口才允许释放。
NeverReachedExecutor,
}
impl DirectStaleTurnReleaseReason {
impl StaleTurnReleaseReason {
pub(crate) fn message(self) -> &'static str {
match self {
Self::ExecutorExited => "陶泥儿执行进程已退出",
@@ -503,31 +503,31 @@ impl DirectStaleTurnReleaseReason {
/// 1. 这个 thread 上确实登记着一轮(Thread Manager 的活动回合);
/// 2. 传了 `expected_client_turn_id` 时必须与登记一致——绝不误伤另一条回合;
/// 3. `reason == NeverReachedExecutor` 时,登记的年龄必须超过
/// [`DIRECT_STALE_TURN_RELEASE_MIN_AGE_MS`],排除"刚放行、还在本地准备"的正常启动窗口。
/// [`STALE_TURN_RELEASE_MIN_AGE_MS`],排除"刚放行、还在本地准备"的正常启动窗口。
///
/// 不满足任何一条都**什么都不写**(占用、事件都不变),返回可读原因;因此它天然幂等——占用已经
/// 没了就返回"没有正在运行的回合"。这一层只负责把 [`DirectStaleTurnRelease`] 翻成人话:判据与释放
/// 没了就返回"没有正在运行的回合"。这一层只负责把 [`StaleTurnRelease`] 翻成人话:判据与释放
/// 动作都在 Thread Manager 里,因为它们必须在同一个临界区里(拆成两步会让并发认领插进中间,见
/// [`release_stale_direct_turn`])。
pub(crate) fn direct_stale_turn_for_release(
pub(crate) fn stale_turn_for_release(
root: &Path,
expected_client_turn_id: Option<&str>,
reason: DirectStaleTurnReleaseReason,
reason: StaleTurnReleaseReason,
) -> Result<String, String> {
let expected = expected_client_turn_id
.map(str::trim)
.filter(|value| !value.is_empty());
let min_age_ms = (reason == DirectStaleTurnReleaseReason::NeverReachedExecutor)
.then_some(DIRECT_STALE_TURN_RELEASE_MIN_AGE_MS);
match release_stale_direct_turn(&direct_thread_id_for_project(root), expected, min_age_ms) {
DirectStaleTurnRelease::Released(client_turn_id) => Ok(client_turn_id),
DirectStaleTurnRelease::NoActiveTurn => {
let min_age_ms = (reason == StaleTurnReleaseReason::NeverReachedExecutor)
.then_some(STALE_TURN_RELEASE_MIN_AGE_MS);
match release_stale_direct_turn(&thread_id_for_project(root), expected, min_age_ms) {
StaleTurnRelease::Released(client_turn_id) => Ok(client_turn_id),
StaleTurnRelease::NoActiveTurn => {
Err("当前项目没有正在运行的陶泥儿回合,无法终止".to_string())
}
DirectStaleTurnRelease::AnotherTurnRunning => {
StaleTurnRelease::AnotherTurnRunning => {
Err("正在运行的是另一条 Direct 客户端回合,已拒绝终止".to_string())
}
DirectStaleTurnRelease::WithinStartupWindow { age_ms } => Err(format!(
StaleTurnRelease::WithinStartupWindow { age_ms } => Err(format!(
"这一轮 Direct 客户端回合刚开始 {} 秒、还在准备中,暂不能强制释放;请稍后再试",
age_ms / 1000
)),
@@ -3159,7 +3159,7 @@ pub(crate) async fn ensure_direct_taonier_art_package_at(
));
}
let mut regeneration_workflow = if mode.regenerates_existing() {
let invocation_id = direct_active_turn_id_at(root)?;
let invocation_id = active_turn_id_at(root)?;
Some(prepare_direct_taonier_regeneration_workflow_at(
root,
prompt,
@@ -5524,23 +5524,23 @@ mod tests {
#[test]
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);
let thread_id = thread_id_for_project(root.path());
assert_eq!(read_turn_identity(&thread_id), None);
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");
let reservation = TurnReservation::accept_for_test(&thread_id, "client-turn-read-1");
let running = read_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:?}");
assert_eq!(
direct_active_turn_id_at(root.path()).expect("identity for the running turn"),
active_turn_id_at(root.path()).expect("identity for the running turn"),
"client-turn-read-1"
);
// 只读探测不释放占用:认领仍然被这一轮挡着。
assert!(claim_direct_pending_turn(&thread_id).is_none());
assert!(claim_pending_turn(&thread_id).is_none());
drop(reservation);
assert_eq!(read_direct_turn_identity(&thread_id), None);
assert!(direct_active_turn_id_at(root.path()).is_err());
assert_eq!(read_turn_identity(&thread_id), None);
assert!(active_turn_id_at(root.path()).is_err());
}
/// 兜底释放的前置判据:身份对得上,且"从没进执行器"必须过启动窗口;判据不过**一个字都不写**,
@@ -5548,36 +5548,36 @@ mod tests {
#[test]
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");
let watcher = subscribe_direct_thread(&thread_id);
let _ = consume_direct_thread(&watcher.subscription_id).expect("drain bootstrap");
let thread_id = thread_id_for_project(root.path());
let reservation = TurnReservation::accept_for_test(&thread_id, "client-turn-stale-1");
let watcher = subscribe_thread(&thread_id);
let _ = consume_thread(&watcher.subscription_id).expect("drain bootstrap");
// ① clientTurnId 不匹配:拒绝,且不误伤正在跑的那一轮。
let mismatch = direct_stale_turn_for_release(
let mismatch = stale_turn_for_release(
root.path(),
Some("client-turn-stale-2"),
DirectStaleTurnReleaseReason::ExecutorExited,
StaleTurnReleaseReason::ExecutorExited,
)
.expect_err("another turn must not be released");
assert!(mismatch.contains("另一条"), "{mismatch}");
// ② 刚放行、还没进执行器:正常启动窗口内不许释放(释放等于放开并发)。
let young = direct_stale_turn_for_release(
let young = stale_turn_for_release(
root.path(),
Some("client-turn-stale-1"),
DirectStaleTurnReleaseReason::NeverReachedExecutor,
StaleTurnReleaseReason::NeverReachedExecutor,
)
.expect_err("a freshly dispatched turn is still starting");
assert!(young.contains("暂不能强制释放"), "{young}");
// ① ② 两条拒绝路径都**什么都没写**:占用还在,事件流上也没有任何东西。
assert_eq!(
read_direct_turn_identity(&thread_id).map(|turn| turn.client_turn_id),
read_turn_identity(&thread_id).map(|turn| turn.client_turn_id),
Some("client-turn-stale-1".to_string())
);
assert!(
consume_direct_thread(&watcher.subscription_id)
consume_thread(&watcher.subscription_id)
.expect("consume after rejected releases")
.events
.is_empty(),
@@ -5585,20 +5585,16 @@ mod tests {
);
// ③ 老过窗口:判成残留——占用与兜底终态在同一个临界区里一起落地。
backdate_direct_active_turn_for_test(&thread_id, DIRECT_STALE_TURN_RELEASE_MIN_AGE_MS + 1);
let released = direct_stale_turn_for_release(
backdate_direct_active_turn_for_test(&thread_id, STALE_TURN_RELEASE_MIN_AGE_MS + 1);
let released = stale_turn_for_release(
root.path(),
Some("client-turn-stale-1"),
DirectStaleTurnReleaseReason::NeverReachedExecutor,
StaleTurnReleaseReason::NeverReachedExecutor,
)
.expect("stale turn is releasable");
assert_eq!(released, "client-turn-stale-1");
assert_eq!(
read_direct_turn_identity(&thread_id),
None,
"占用必须一起解除"
);
let events = consume_direct_thread(&watcher.subscription_id)
assert_eq!(read_turn_identity(&thread_id), None, "占用必须一起解除");
let events = consume_thread(&watcher.subscription_id)
.expect("consume released terminal")
.events;
assert_eq!(events.len(), 1, "兜底释放只写一条终态:{events:?}");
@@ -5606,7 +5602,7 @@ mod tests {
assert!(
matches!(
terminal,
DirectThreadEvent::TurnCompleted { status, .. } if status == "aborted"
ThreadEvent::TurnCompleted { status, .. } if status == "aborted"
),
"{terminal:?}"
);
@@ -5619,14 +5615,10 @@ mod tests {
// ④ 释放之后这个 thread 立刻腾出来了(占用解除与终态同批落地):下一条能重新认领,
// 而且"执行进程已退出"不看启动窗口——刚放行也允许释放。
let _next = DirectTurnReservation::accept_for_test(&thread_id, "client-turn-stale-2");
let _next = TurnReservation::accept_for_test(&thread_id, "client-turn-stale-2");
assert_eq!(
direct_stale_turn_for_release(
root.path(),
None,
DirectStaleTurnReleaseReason::ExecutorExited,
)
.expect("executor exited ignores the start window"),
stale_turn_for_release(root.path(), None, StaleTurnReleaseReason::ExecutorExited,)
.expect("executor exited ignores the start window"),
"client-turn-stale-2"
);
drop(reservation);
@@ -5715,12 +5707,9 @@ mod tests {
#[test]
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,
DirectStaleTurnReleaseReason::ExecutorExited,
)
.expect_err("nothing to release");
let nothing =
stale_turn_for_release(root.path(), None, StaleTurnReleaseReason::ExecutorExited)
.expect_err("nothing to release");
assert!(nothing.contains("没有正在运行"), "{nothing}");
}
@@ -6701,8 +6690,8 @@ mod tests {
.expect("completed replay authority is the stable outer turn, not a resampled brief");
assert_eq!(completed_replay, persisted);
let _active_turn = DirectTurnReservation::accept_for_test(
&direct_thread_id_for_project(root.path()),
let _active_turn = TurnReservation::accept_for_test(
&thread_id_for_project(root.path()),
"client-turn-completed-1",
);
let replay = ensure_direct_taonier_art_package_at(
@@ -1,7 +1,7 @@
//! DirectProject 用户输入命令适配器。
//!
//! Tauri 只在这里接收前端 item,校验与 canonical→prompt 投影交给 user-item 深模块,队列交给
//! Thread Manager;放行在 `direct_turn_dispatch`,不在这里。
//! Thread Manager;放行在 `thread_manager::dispatch`,不在这里。
use super::*;
@@ -35,7 +35,7 @@ pub(crate) fn normalize_direct_client_turn_id(
/// 之后只做一件事——把待发消息排进队列:不登记占用、不落盘、不发回合事件、不起 codex。
///
/// 于是"这一轮跑成什么"仍然只有订阅事件一个来源:命令返回 `Ok` 只说明**入队成立**。真正的回合边界
/// (`turn.started` / `turn.completed`)由 Thread Manager 在**放行**时写出(见 `direct_turn_dispatch`),
/// (`turn.started` / `turn.completed`)由 Thread Manager 在**放行**时写出(见 `thread_manager::dispatch`),
/// 入队失败不写用户条目、不产生任何事件。可留痕的调用级失败(宿主 / 环境事实)仍在边界补一份运行
/// 错误诊断,返回串不带诊断引用。
///
@@ -93,15 +93,13 @@ async fn enqueue_direct_codex_turn_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 thread_id = direct_thread_id_for_project(root);
let thread_id = thread_id_for_project(root);
// 容量是**最便宜**的一条入队前置条件,先判:满了就直接说满,别让一条注定收不下的消息先去做工程
// 准备——那是分钟级的活(`npm ci` + Vite 构建),还会在磁盘上留下产物。权威判据仍然是入队那一刻
// 临界区里的容量检查(下面 `enqueue_direct_pending_turn`):这里只是快速失败,中间被别人的消息
// 临界区里的容量检查(下面 `enqueue_pending_turn`):这里只是快速失败,中间被别人的消息
// 挤满时那一条照样拦得住。
queue_has_room(direct_pending_turn_count(&thread_id)).map_err(|_| {
DirectTurnError::QueueFull {
limit: MAX_PENDING_DIRECT_TURNS,
}
queue_has_room(pending_turn_count(&thread_id)).map_err(|_| DirectTurnError::QueueFull {
limit: MAX_PENDING_TURNS,
})?;
// 创建类型来自结构化用户入口;实际工程和可信脚手架由宿主复核。
crate::environment_check::prepare_new_web_project_at(root, creation_type.as_deref())
@@ -112,7 +110,7 @@ async fn enqueue_direct_codex_turn_typed(
})?;
// 入队:到这里这一条已经过了全部检查,剩下的就是排队等放行。prompt 与 canonical 形状在这里
// 冻结——放行不重算,所以放行没有失败出口。
let pending = PendingDirectTurn::prepare(
let pending = PendingTurn::prepare(
turn_id,
user_item,
user_prompt,
@@ -123,37 +121,37 @@ async fn enqueue_direct_codex_turn_typed(
.map_err(|error| DirectTurnError::InputRejected {
detail: error.to_string(),
})?;
match enqueue_direct_pending_turn(&thread_id, pending) {
match enqueue_pending_turn(&thread_id, pending) {
Ok(_) => {}
Err(EnqueueRejection::QueueFull) => {
return Err(DirectTurnError::QueueFull {
limit: MAX_PENDING_DIRECT_TURNS,
limit: MAX_PENDING_TURNS,
})
}
}
// 入队之后立刻踢一脚:队列空且没有回合在跑时,放行就是这一脚,用户点发送不必再等一个调度周期。
kick_direct_queue_dispatch(root);
kick_queue_dispatch(root);
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::agent::{consume_direct_thread, subscribe_direct_thread, DirectThreadEvent};
use crate::agent::{consume_thread, subscribe_thread, ThreadEvent};
/// 放行之后整轮是异步的:测试只能等事件。返回从订阅之后看到的全部事件。
async fn wait_for_completion(subscription_id: &str) -> Vec<DirectThreadEvent> {
async fn wait_for_completion(subscription_id: &str) -> Vec<ThreadEvent> {
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)
consume_thread(subscription_id)
.expect("consume logical turn")
.events,
);
if events
.iter()
.any(|event| matches!(event, DirectThreadEvent::TurnCompleted { .. }))
.any(|event| matches!(event, ThreadEvent::TurnCompleted { .. }))
{
return events;
}
@@ -184,9 +182,9 @@ mod tests {
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 thread_id = thread_id_for_project(&root);
let subscription = subscribe_thread(&thread_id);
let _ = consume_thread(&subscription.subscription_id);
enqueue_direct_codex_turn_typed(
&root,
@@ -202,7 +200,7 @@ mod tests {
let terminal = events
.iter()
.filter_map(|event| match event {
DirectThreadEvent::TurnCompleted {
ThreadEvent::TurnCompleted {
status,
failure,
user_item_id,
@@ -221,7 +219,7 @@ mod tests {
);
assert_eq!(user_item_id.as_deref(), Some("direct-codex:turn-1:user"));
// 占用已释放:队列继续走得动。
assert!(!crate::agent::direct_thread_turn_is_active(&thread_id));
assert!(!crate::agent::thread_turn_is_active(&thread_id));
}
/// 本轮的开口用户条目必须先于整轮里任何可能失败的东西下发。
@@ -239,9 +237,9 @@ mod tests {
let root = temp.path().join("direct-user-item-first");
crate::init_local_game_project_at(&root, "direct-user-item-first", "用户条目先下发")
.expect("init project");
let thread_id = direct_thread_id_for_project(&root);
let subscription = subscribe_direct_thread(&thread_id);
let _ = consume_direct_thread(&subscription.subscription_id);
let thread_id = thread_id_for_project(&root);
let subscription = subscribe_thread(&thread_id);
let _ = consume_thread(&subscription.subscription_id);
enqueue_direct_codex_turn_typed(
&root,
@@ -257,12 +255,12 @@ mod tests {
// 队列事件先于放行:`queue.enqueued` → 放行(开始 + 离开队列)→ 开口用户条目 → 终态。
let started = events
.iter()
.position(|event| matches!(event, DirectThreadEvent::TurnStarted { .. }))
.position(|event| matches!(event, ThreadEvent::TurnStarted { .. }))
.expect("这一轮必须被放行");
assert!(
matches!(
events.get(started),
Some(DirectThreadEvent::TurnStarted { user_item_id, .. })
Some(ThreadEvent::TurnStarted { user_item_id, .. })
if user_item_id.as_deref() == Some("direct-codex:turn-1:user")
),
"放行的第一条事件必须是带身份的回合开始:{events:?}"
@@ -275,7 +273,7 @@ mod tests {
assert!(
matches!(
events.get(opening),
Some(DirectThreadEvent::ItemCompleted { item, .. })
Some(ThreadEvent::ItemCompleted { item, .. })
if item.item_id() == "direct-codex:turn-1:user"
),
"第一条条目事件必须是本轮的开口用户条目:{events:?}"
@@ -283,7 +281,7 @@ mod tests {
assert!(opening > started, "用户条目只能在放行之后:{events:?}");
let terminal = events
.iter()
.position(|event| matches!(event, DirectThreadEvent::TurnCompleted { .. }));
.position(|event| matches!(event, ThreadEvent::TurnCompleted { .. }));
assert!(
terminal.is_none_or(|index| index > opening),
"终态只能在用户条目之后:{events:?}"
@@ -310,9 +308,9 @@ mod tests {
let root = temp.path().join("direct-history-write-failure");
crate::init_local_game_project_at(&root, "direct-history-write", "落盘失败")
.expect("init project");
let thread_id = direct_thread_id_for_project(&root);
let subscription = subscribe_direct_thread(&thread_id);
let _ = consume_direct_thread(&subscription.subscription_id);
let thread_id = thread_id_for_project(&root);
let subscription = subscribe_thread(&thread_id);
let _ = consume_thread(&subscription.subscription_id);
// 接下来这次追加写的两次尝试都按"争用失败"返回:确定性地走到落盘失败分支。
std::fs::write(
root.join(".agent/runtime/test-fail-next-direct-project-history-append"),
@@ -334,7 +332,7 @@ mod tests {
let terminals = events
.iter()
.filter_map(|event| match event {
DirectThreadEvent::TurnCompleted {
ThreadEvent::TurnCompleted {
status, failure, ..
} => Some((status, failure)),
_ => None,
@@ -357,6 +355,6 @@ mod tests {
failure.message
);
// 占用已释放:下一轮还能继续。
assert!(!crate::agent::direct_thread_turn_is_active(&thread_id));
assert!(!crate::agent::thread_turn_is_active(&thread_id));
}
}
@@ -267,7 +267,7 @@ impl DirectToolBridge {
impl DirectToolBridgeState {
fn begin_user_turn(self: &Arc<Self>) -> Result<DirectToolBridgeTurnGuard, String> {
let turn_id = direct_active_turn_id_at(&self.root)?;
let turn_id = active_turn_id_at(&self.root)?;
let mut authorization = self
.turn_authorization
.lock()
@@ -4839,10 +4839,8 @@ 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 = DirectTurnReservation::accept_for_test(
&direct_thread_id_for_project(root.path()),
client_turn_id,
);
let _active_invocation =
TurnReservation::accept_for_test(&thread_id_for_project(root.path()), client_turn_id);
let active_turn = state
.begin_user_turn()
.expect("client turn authorization state");
@@ -8,7 +8,7 @@
//! 为什么不复用 `project.jsonl`:那条链路的回读只投影 `role ∈ {user, assistant}` 的
//! 文本条目,而且会被注入 Codex 上下文。往里面塞新形状既装不下,又有污染模型上下文的风险。
use super::direct_thread_wire::sanitize_detail_text;
use super::thread_manager::wire::sanitize_detail_text;
use crate::config::write_game_creator_private_file;
use crate::project::{enforce_project_permission_policy, project_append_lock_for};
use serde::{Deserialize, Serialize};
@@ -3850,16 +3850,14 @@ mod tests {
root: &Path,
) -> (
super::super::direct_tool_bridge::DirectToolBridge,
DirectTurnReservation,
TurnReservation,
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 = DirectTurnReservation::accept_for_test(
&direct_thread_id_for_project(root),
"tool-chain-turn",
);
let turn =
TurnReservation::accept_for_test(&thread_id_for_project(root), "tool-chain-turn");
let execution =
super::super::direct_tool_bridge::direct_execution_fixture(root, "tool-chain-turn")
.await;
@@ -14,7 +14,7 @@
//! 所以"同一种错误在入队与放行走不同通道"是正常的:[`DirectTurnError::EnvironmentNotReady`] 两边
//! 都可能出现,位置说了算。这里**没有**、也不该有"这个变体是不是回合失败"的判据。
//!
//! 事件载荷(`direct_thread_wire::DirectTurnFailure`)仍然只有 `{kind, message}` 两个字段:
//! 事件载荷(`thread_manager::wire::DirectTurnFailure`)仍然只有 `{kind, message}` 两个字段:
//! 那是**线上协议**,由 [`DirectTurnError::wire_kind`] 与 `Display` 在这一个出口投影出来,不是
//! 另一种状态模型。跨进程边界(`#[tauri::command]`)同样只给前端一个字符串:那是**序列化**,
//! 由 `Display` 一处生成;Rust 侧任何地方都不再解析这个字符串。
@@ -392,7 +392,7 @@ pub(crate) enum DirectTurnError {
ClientTurnIdMalformed { min_chars: usize, max_chars: usize },
/// 待发消息队列已满:这一条没进队,等前面几条发完再发。
///
/// 上限只落在宿主这一处(`MAX_PENDING_DIRECT_TURNS`),随载荷带出去,界面不自己数一份。
/// 上限只落在宿主这一处(`MAX_PENDING_TURNS`),随载荷带出去,界面不自己数一份。
QueueFull { limit: usize },
/// 项目目录锚不定(符号链接 / 权限 / 目录被删)。
ProjectRootUnanchored { cause: String },
@@ -5,10 +5,10 @@
//! 1. [`direct_turn_terminal`]:拿这一轮的事实判定终态——是不是失败、原因是什么、状态写什么;
//! 2. [`DirectTurnTerminal::event`]:把终态投影成 `turn.completed` 事件。
//!
//! 终态的**出口**(谁写、什么时候兜底)不在这里,在 `direct_turn_dispatch.rs` 的放行占用对象里:
//! 终态的**出口**(谁写、什么时候兜底)不在这里,在 `thread_manager::dispatch` 的放行占用对象里:
//! 这个模块只负责"什么算失败、原因怎么写"。
//!
//! 失败载荷的**形状**属于线上协议,定义在 `direct_thread_wire.rs`(`DirectTurnFailure`);
//! 失败载荷的**形状**属于线上协议,定义在 `thread_manager::wire`(`DirectTurnFailure`);
//! 载荷的 `kind` 与 `message` 由 [`DirectTurnError`] 投影而来(`kind` 的取值表见
//! [`DirectTurnError::wire_kind`]);这里只负责"什么算失败、原因怎么写、什么时候兜底",
//! 不碰事件队列的搬运规则,也不自己认 `LlmError`。
@@ -16,8 +16,8 @@
use std::path::Path;
use super::{
redact_agent_runtime_error, DirectThreadEvent, DirectTurnError, DirectTurnFailure,
DirectTurnFailureKind,
redact_agent_runtime_error, DirectTurnError, DirectTurnFailure, DirectTurnFailureKind,
ThreadEvent,
};
/// `turn.completed.failure.message` 的字符上限:与本地错误文案同一档——够说清原因,又不至于
@@ -36,10 +36,10 @@ pub(crate) struct DirectTurnTerminal {
impl DirectTurnTerminal {
/// 终态事件:失败时同一个 `turn.completed` 带载荷,其余只带 `status`。
pub(crate) fn event(self, completed_at: u64, user_item_id: Option<&str>) -> DirectThreadEvent {
pub(crate) fn event(self, completed_at: u64, user_item_id: Option<&str>) -> ThreadEvent {
let event = match self.failure {
Some(failure) => DirectThreadEvent::turn_completed_failed(failure, completed_at),
None => DirectThreadEvent::turn_completed(self.status, completed_at),
Some(failure) => ThreadEvent::turn_completed_failed(failure, completed_at),
None => ThreadEvent::turn_completed(self.status, completed_at),
};
event.with_user_item_id(user_item_id)
}
@@ -122,7 +122,7 @@ impl DirectTurnTerminal {
#[cfg(test)]
mod tests {
use super::*;
use crate::agent::{consume_direct_thread, subscribe_direct_thread};
use crate::agent::{consume_thread, subscribe_thread};
use platform_llm::LlmError;
fn history_root() -> std::path::PathBuf {
@@ -245,7 +245,7 @@ stderrClass=nonempty;stderrBytes=1000";
assert!(event.failure().is_none());
assert!(matches!(
event,
DirectThreadEvent::TurnCompleted { ref status, .. } if status == "completed"
ThreadEvent::TurnCompleted { ref status, .. } if status == "completed"
));
}
@@ -368,10 +368,10 @@ async fn start_for_current_turn(
key: String,
build: bool,
) -> Result<ValidationStart, String> {
let turn_id = direct_active_turn_id_at(root)?;
let turn_id = active_turn_id_at(root)?;
let root = root.to_path_buf();
tokio::task::spawn_blocking(move || {
if direct_active_turn_id_at(&root)? != turn_id {
if active_turn_id_at(&root)? != turn_id {
return Err("validation-turn-changed: 当前验证所属回合已结束".into());
}
let session = super::direct_execution::current(&root)?;
@@ -625,7 +625,7 @@ const ERROR_REDACTED_KEY: &str = "[redacted-sensitive-field]";
/// information whenever a safe error line mentions a credential field.
///
/// 逐行脱敏,但**保留每个 chunk 末尾的换行**:本函数会被流式增量逐段调用
/// (`direct_thread_delta_text` → 前端把各段拼成一条消息再交给 Markdown 渲染)。用
/// (`thread_delta_text` → 前端把各段拼成一条消息再交给 Markdown 渲染)。用
/// `lines()` + `join("\n")` 会把「以换行结尾的段」的末尾换行吃掉,拼接后段落、列表项和表格行
/// 会并进同一行,整条消息的 Markdown 结构(尤其表格)就作废了。
pub(crate) fn sanitize_error_context(value: &str) -> String {
@@ -71,7 +71,7 @@ impl DirectGameCreatorTurnUpdateEmitter {
pub(crate) fn new(root: &Path, turn_id: String) -> Self {
Self {
project_path: root.to_string_lossy().into_owned(),
thread_id: crate::agent::direct_thread_id_for_project(root),
thread_id: crate::agent::thread_id_for_project(root),
turn_id,
sequence: Arc::new(AtomicU64::new(0)),
}
@@ -168,7 +168,7 @@ impl DirectGameCreatorTurnUpdateEmitter {
.unwrap_or_default()
.as_millis()
.min(u64::MAX as u128) as u64;
update_direct_thread_active_turn(
update_active_turn(
&self.thread_id,
&self.turn_id,
status,
@@ -1,8 +1,8 @@
//! DirectProject 的放行:把队首那条待发消息变成一轮真的回合。
//!
//! 这个模块两件事,别再往里加第三件:
//! 1. [`DirectTurnReservation`]:这一轮的占用对象——持有它代表还没收口,终态出口只认它的 token;
//! 2. [`kick_direct_queue_dispatch`]:认领队首 → 起整轮(落盘用户条目 → 下发 → 跑回合)。
//! 1. [`TurnReservation`]:这一轮的占用对象——持有它代表还没收口,终态出口只认它的 token;
//! 2. [`kick_queue_dispatch`]:认领队首 → 起整轮(落盘用户条目 → 下发 → 跑回合)。
//!
//! 这一轮的**调用身份**(付费美术、执行会话、MCP、校验、上下文预取读的那个"这一轮是谁")由认领
//! 在 Thread Manager 上登记,认领成功即成立;这里不再另取一次、也不再持第二份身份。
@@ -13,19 +13,19 @@
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, DirectTurnError, DirectTurnTerminal,
PendingDirectTurn,
use crate::agent::{
append_direct_project_user_message_at, claim_pending_turn, complete_turn_if_reserved,
direct_tool_call_now_ms, redact_agent_runtime_error,
run_direct_game_creator_turn_at_with_creation_type_and_emitter, thread_id_for_project,
DirectGameCreatorTurnUpdateEmitter, DirectTurnError, DirectTurnTerminal, DispatchedTurn,
PendingTurn,
};
/// 一次放行的占用。持有它就代表这一轮还没收口。
///
/// 生命周期由调用方决定:登记与开始事件已经由 [`kick_direct_queue_dispatch`] 在同一临界区里写好,
/// 生命周期由调用方决定:登记与开始事件已经由 [`kick_queue_dispatch`] 在同一临界区里写好,
/// 这里只接住这一轮的终态出口,并随放行任务一起 drop(正常 / 失败 / panic 都走同一个出口)。
pub(crate) struct DirectTurnReservation {
pub(crate) struct TurnReservation {
thread_id: String,
token: String,
user_item_id: Option<String>,
@@ -36,9 +36,9 @@ pub(crate) struct DirectTurnReservation {
_session_keepalive: Option<crate::auth_session::ClientSessionKeepalive>,
}
impl DirectTurnReservation {
impl TurnReservation {
/// 接管一次放行的占用:占用登记与 `turn.started` 已由 Thread Manager 的认领写好。
pub(crate) fn resume(thread_id: &str, dispatched: &DirectDispatchedTurn) -> Self {
pub(crate) fn resume(thread_id: &str, dispatched: &DispatchedTurn) -> Self {
Self {
thread_id: thread_id.to_string(),
token: dispatched.token.clone(),
@@ -51,7 +51,7 @@ impl DirectTurnReservation {
///
/// 深层(真正跑完这一轮的代码)已经写出终态时返回 `false`,兜底不覆盖真实结果。
pub(crate) fn finish_if_unfinished(&self, terminal: DirectTurnTerminal) -> bool {
complete_direct_thread_turn_if_reserved(
complete_turn_if_reserved(
&self.thread_id,
&self.token,
terminal.event(direct_tool_call_now_ms(), self.user_item_id.as_deref()),
@@ -61,7 +61,7 @@ impl DirectTurnReservation {
/// 测试用:直接开一轮占用,等价于"入队 + 认领",不经过 Tauri 命令与真实的用户条目。
#[cfg(test)]
pub(crate) fn accept_for_test(thread_id: &str, client_turn_id: &str) -> Self {
let pending = PendingDirectTurn::prepare(
let pending = PendingTurn::prepare(
client_turn_id.to_string(),
serde_json::from_value(serde_json::json!({
"type": "message",
@@ -76,37 +76,37 @@ impl DirectTurnReservation {
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");
super::enqueue_pending_turn(thread_id, pending).expect("enqueue test turn");
let dispatched = claim_pending_turn(thread_id).expect("claim test turn");
Self::resume(thread_id, &dispatched)
}
}
impl Drop for DirectTurnReservation {
impl Drop for TurnReservation {
fn drop(&mut self) {
// 兜底:任务 panic、future 被丢弃、或今后在终态之前新增的 `?` 早退。
// 这类失败说不出原因,只给分类;能说清原因的错误必须由调用方在更早的地方显式收口。
let _ = self.finish_if_unfinished(DirectTurnTerminal::host_dropped());
// 占用的释放就是"这个项目腾出了跑回合的位置":踢一脚,队里排着的消息不必等下一次用户动作。
// panic 也走这里,所以队列不会因为一个任务炸掉而永久停住。
kick_direct_queue_dispatch(Path::new(&self.thread_id));
kick_queue_dispatch(Path::new(&self.thread_id));
}
}
/// 踢一脚:让队首那条待发消息有机会变成一轮真的回合。
///
/// 幂等,三个调用点:**入队之后**(队列空且没有回合在跑时,放行就是这一脚,不必再等一个调度周期)、
/// **占用释放**(正常 / 失败 / panic 共用,见 [`DirectTurnReservation`] 的 `Drop`)、**中止路径**
/// **占用释放**(正常 / 失败 / panic 共用,见 [`TurnReservation`] 的 `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 {
pub(crate) fn kick_queue_dispatch(root: &Path) {
let thread_id = thread_id_for_project(root);
let Some(dispatched) = claim_pending_turn(&thread_id) else {
return;
};
let reservation = DirectTurnReservation::resume(&thread_id, &dispatched);
let reservation = TurnReservation::resume(&thread_id, &dispatched);
let root = root.to_path_buf();
tauri::async_runtime::spawn(async move {
run_dispatched_direct_turn(root, dispatched, reservation).await;
@@ -119,8 +119,8 @@ pub(crate) fn kick_direct_queue_dispatch(root: &Path) {
/// 命令的返回值上。
async fn run_dispatched_direct_turn(
root: PathBuf,
dispatched: DirectDispatchedTurn,
reservation: DirectTurnReservation,
dispatched: DispatchedTurn,
reservation: TurnReservation,
) {
let turn_id = dispatched.pending.client_turn_id.clone();
// 落盘即回合成立:放行之后必须留下这条用户消息,哪怕这一轮随后失败。
@@ -174,35 +174,32 @@ async fn run_dispatched_direct_turn(
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,
consume_thread, enqueue_pending_turn, subscribe_thread, thread_turn_is_active,
DirectTurnFailure, DirectTurnFailureKind, PendingTurn, ThreadEvent,
};
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);
let bootstrap = subscribe_thread(thread_id);
let _ = consume_thread(&bootstrap.subscription_id);
bootstrap.subscription_id
}
fn pending(subscription_id: &str) -> Vec<DirectThreadEvent> {
consume_direct_thread(subscription_id)
.expect("consume")
.events
fn pending(subscription_id: &str) -> Vec<ThreadEvent> {
consume_thread(subscription_id).expect("consume").events
}
fn turn_completed_events(events: &[DirectThreadEvent]) -> Vec<&DirectThreadEvent> {
fn turn_completed_events(events: &[ThreadEvent]) -> Vec<&ThreadEvent> {
events
.iter()
.filter(|event| matches!(event, DirectThreadEvent::TurnCompleted { .. }))
.filter(|event| matches!(event, ThreadEvent::TurnCompleted { .. }))
.collect()
}
/// 一条待发消息。放行路径只读它的身份与产物,所以这里用最小可用形状。
fn pending_turn(client_turn_id: &str) -> PendingDirectTurn {
PendingDirectTurn::prepare(
fn pending_turn(client_turn_id: &str) -> PendingTurn {
PendingTurn::prepare(
client_turn_id.to_string(),
serde_json::from_value(serde_json::json!({
"type": "message",
@@ -226,7 +223,7 @@ mod tests {
}
fn enqueue(thread_id: &str, client_turn_id: &str) {
enqueue_direct_pending_turn(thread_id, pending_turn(client_turn_id)).expect("enqueue");
enqueue_pending_turn(thread_id, pending_turn(client_turn_id)).expect("enqueue");
}
/// 放行的占用只认**自己的 token**:认领写下的占用身份与这一轮绑死,别人的终态盖不上它。
@@ -234,20 +231,20 @@ mod tests {
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(
let dispatched = claim_pending_turn(&thread).expect("claim head");
let owner = TurnReservation::resume(&thread, &dispatched);
let foreign = TurnReservation::resume(
&thread,
&DirectDispatchedTurn {
&DispatchedTurn {
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!(thread_turn_is_active(&thread));
assert!(owner.finish_if_unfinished(DirectTurnTerminal::host_dropped()));
assert!(!direct_thread_turn_is_active(&thread));
assert!(!thread_turn_is_active(&thread));
// 显式收口之后 Drop 不再补第二条:兜底只负责"没人写过"的那一种。
drop(foreign);
drop(owner);
@@ -259,11 +256,11 @@ mod tests {
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 dispatched = claim_pending_turn(&thread).expect("claim head");
let reservation = TurnReservation::resume(&thread, &dispatched);
// 深层收口:真正跑完这一轮的代码算出来的终态。
let deep = DirectThreadEvent::turn_completed_failed(
let deep = ThreadEvent::turn_completed_failed(
DirectTurnFailure::new(
DirectTurnFailureKind::Timeout,
"等待模型回执超时".to_string(),
@@ -271,7 +268,7 @@ mod tests {
2_000,
)
.with_user_item_id(Some("direct-codex:turn-1:user"));
crate::agent::complete_direct_thread_turn(&thread, deep);
crate::agent::complete_turn(&thread, deep);
assert!(
!reservation.finish_if_unfinished(DirectTurnTerminal::host_dropped()),
@@ -283,7 +280,7 @@ mod tests {
let completed = turn_completed_events(&events);
assert_eq!(completed.len(), 1, "一轮只许有一条终态:{events:?}");
match completed[0] {
DirectThreadEvent::TurnCompleted { failure, .. } => {
ThreadEvent::TurnCompleted { failure, .. } => {
assert_eq!(
failure.as_ref().map(|f| f.kind),
Some(DirectTurnFailureKind::Timeout)
@@ -300,8 +297,8 @@ mod tests {
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);
let dispatched = claim_pending_turn(&thread).expect("claim head");
let reservation = TurnReservation::resume(&thread, &dispatched);
drop(reservation);
@@ -309,12 +306,12 @@ mod tests {
let events = pending(&subscription);
let terminal = events
.iter()
.position(|event| matches!(event, DirectThreadEvent::TurnCompleted { .. }))
.position(|event| matches!(event, ThreadEvent::TurnCompleted { .. }))
.expect("第一条必须被收口");
assert!(
matches!(
events.get(terminal),
Some(DirectThreadEvent::TurnCompleted { status, failure: Some(failure), .. })
Some(ThreadEvent::TurnCompleted { status, failure: Some(failure), .. })
if status == "failed" && failure.kind == DirectTurnFailureKind::HostDropped
),
"{events:?}"
@@ -322,7 +319,7 @@ mod tests {
let next_started = events.iter().position(|event| {
matches!(
event,
DirectThreadEvent::TurnStarted { user_item_id, .. }
ThreadEvent::TurnStarted { user_item_id, .. }
if user_item_id.as_deref() == Some("direct-codex:turn-2:user")
)
});
@@ -336,19 +333,19 @@ mod tests {
#[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));
kick_queue_dispatch(Path::new(&empty));
assert!(!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);
let dispatched = claim_pending_turn(&busy).expect("claim head");
let reservation = TurnReservation::resume(&busy, &dispatched);
kick_direct_queue_dispatch(Path::new(&busy));
kick_queue_dispatch(Path::new(&busy));
// 仍在跑的那一轮没有被顶掉:占用身份还是第一条。
assert!(direct_thread_turn_is_active(&busy));
assert!(thread_turn_is_active(&busy));
drop(reservation);
}
@@ -359,7 +356,7 @@ mod tests {
let thread = unique_thread("stale-cancel");
let subscription = watch(&thread);
// 这一轮已经被放行(占用登记在 Thread Manager 上),但任务泄漏:拿不到可中断句柄。
let leaked = DirectTurnReservation::accept_for_test(&thread, "turn-stale-1");
let leaked = TurnReservation::accept_for_test(&thread, "turn-stale-1");
enqueue(&thread, "turn-stale-2");
// "从没进执行器"只对过了启动窗口的残留放行;真实用例等不了 60 秒。
@@ -369,7 +366,7 @@ mod tests {
.expect("stale release");
assert_eq!(
view.outcome,
super::super::codex_app_server::DIRECT_TURN_CANCEL_OUTCOME_RELEASED
crate::agent::codex_app_server::DIRECT_TURN_CANCEL_OUTCOME_RELEASED
);
assert_eq!(view.client_turn_id, "turn-stale-1");
@@ -378,14 +375,14 @@ mod tests {
assert!(
events.iter().any(|event| matches!(
event,
DirectThreadEvent::TurnCompleted { status, .. } if status == "aborted"
ThreadEvent::TurnCompleted { status, .. } if status == "aborted"
)),
"{events:?}"
);
assert!(
events.iter().any(|event| matches!(
event,
DirectThreadEvent::TurnStarted { user_item_id, .. }
ThreadEvent::TurnStarted { user_item_id, .. }
if user_item_id.as_deref() == Some("direct-codex:turn-stale-2:user")
)),
"队首必须在兜底释放之后被放行:{events:?}"
@@ -4,28 +4,28 @@
//! `queue.enqueued` 事件不可回收;离开队列(取消或放行)时才转成可回收。新订阅者的 bootstrap 因此
//! 天然看得见当前队列,不需要第二张队列表,也不会有"事件与队列不一致"的窗口。
//!
//! 这个模块只放三件事:一条待发消息带走什么([`PendingDirectTurn`])、容量规则
//! ([`MAX_PENDING_DIRECT_TURNS`])、以及它在线上长什么样([`PendingDirectTurn::enqueued_event`])。
//! 这个模块只放三件事:一条待发消息带走什么([`PendingTurn`])、容量规则
//! ([`MAX_PENDING_TURNS`])、以及它在线上长什么样([`PendingTurn::enqueued_event`])。
//! 它不碰锁、不碰 Tauri、不写盘:入队检查在命令侧,放行顺序在 Thread Manager。
use serde_json::Value;
use crate::agent::{
direct_codex_user_item_id_for_client_turn_id, DirectCodexUserItem, DirectThreadEvent,
direct_codex_user_item_id_for_client_turn_id, DirectCodexUserItem, ThreadEvent,
};
/// 一个项目最多能同时排队的待发消息条数。
///
/// 只数**在队**条目,不算正在跑的那一轮。上限只落在宿主这一处:前端不再自己数,满队由命令返回
/// typed 入队失败。
pub(crate) const MAX_PENDING_DIRECT_TURNS: usize = 5;
pub(crate) const MAX_PENDING_TURNS: usize = 5;
/// 一条已经通过入队检查、正在等放行的用户消息。
///
/// 只在内存里,进程重启即消失(与 ADR 记的边界一致)。它同时是**放行时要用的全部输入**:放行
/// 没有失败出口,所以检查产物在入队时就地冻结,放行只搬运、不重算。
#[derive(Clone, Debug)]
pub(crate) struct PendingDirectTurn {
pub(crate) struct PendingTurn {
/// 这条消息的回合身份;放行后同一轮的 `turn.started` / `turn.completed` 用它。
pub(crate) client_turn_id: String,
/// canonical 用户条目:事件与界面 chip 都读它,Rust 不渲染展示形状。
@@ -41,7 +41,7 @@ pub(crate) struct PendingDirectTurn {
pub(crate) at: u64,
}
impl PendingDirectTurn {
impl PendingTurn {
/// 组一条待发消息:入队检查已经全部通过,这里只把放行要用的产物冻结下来。
///
/// 冻结是刻意的:放行没有失败出口,所以任何可能在放行时才失败的计算都必须提前到这里
@@ -68,8 +68,8 @@ impl PendingDirectTurn {
/// 入队事件的投影:带 canonical 用户条目与 `creationType`,**不带 prompt**(prompt 只留在宿主的
/// 队列条目里,它不是要下发的展示形状)。
pub(crate) fn enqueued_event(&self) -> DirectThreadEvent {
DirectThreadEvent::queue_enqueued(
pub(crate) fn enqueued_event(&self) -> ThreadEvent {
ThreadEvent::queue_enqueued(
self.client_turn_id.clone(),
self.user_item.clone(),
self.creation_type.clone(),
@@ -95,13 +95,13 @@ pub(crate) enum EnqueueOutcome {
/// 入队被检查挡下来的原因。满队之外的原因由入队侧自己的 typed 错误表达。
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum EnqueueRejection {
/// 在队条目已达 [`MAX_PENDING_DIRECT_TURNS`]。
/// 在队条目已达 [`MAX_PENDING_TURNS`]。
QueueFull,
}
/// 队列还有没有位置。`pending_count` 只数在队条目。
pub(crate) fn queue_has_room(pending_count: usize) -> Result<(), EnqueueRejection> {
if pending_count >= MAX_PENDING_DIRECT_TURNS {
if pending_count >= MAX_PENDING_TURNS {
return Err(EnqueueRejection::QueueFull);
}
Ok(())
@@ -122,8 +122,8 @@ mod tests {
.expect("canonical user item")
}
fn pending(client_turn_id: &str, creation_type: Option<&str>) -> PendingDirectTurn {
PendingDirectTurn::prepare(
fn pending(client_turn_id: &str, creation_type: Option<&str>) -> PendingTurn {
PendingTurn::prepare(
client_turn_id.to_string(),
user_item("生成一个游戏", "direct-codex:turn-1:user"),
"生成一个游戏".to_string(),
@@ -194,11 +194,11 @@ mod tests {
/// 容量只数在队条目,5 条封口;在跑的那一轮不算进去。
#[test]
fn capacity_closes_at_five_pending_turns() {
for count in 0..MAX_PENDING_DIRECT_TURNS {
for count in 0..MAX_PENDING_TURNS {
assert_eq!(queue_has_room(count), Ok(()), "{count} 条时仍该有位置");
}
assert_eq!(
queue_has_room(MAX_PENDING_DIRECT_TURNS),
queue_has_room(MAX_PENDING_TURNS),
Err(EnqueueRejection::QueueFull)
);
}
@@ -5893,9 +5893,8 @@ pub(crate) async fn read_direct_project_conversation(
}
#[tauri::command]
pub(crate) fn list_game_creator_direct_active_turns(
) -> Result<Vec<DirectActiveTurnSnapshot>, String> {
list_direct_active_turns()
pub(crate) fn list_game_creator_direct_active_turns() -> Result<Vec<ActiveTurnSnapshot>, String> {
list_active_turns()
}
/// 取消一条**还没放行**的待发消息。
@@ -5907,15 +5906,12 @@ pub(crate) fn list_game_creator_direct_active_turns(
pub(crate) async fn remove_direct_project_pending_turn(
project_path: String,
client_turn_id: String,
) -> Result<DirectQueueRemovalOutcome, String> {
) -> Result<QueueRemovalOutcome, 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(),
))
let thread_id = thread_id_for_project(root);
Ok(remove_pending_turn(&thread_id, client_turn_id.trim()))
})
.await
.map_err(|error| format!("取消 DirectProject 待发消息后台任务失败:{error}"))?
@@ -5924,12 +5920,12 @@ pub(crate) async fn remove_direct_project_pending_turn(
#[tauri::command]
pub(crate) async fn subscribe_direct_project_thread(
project_path: String,
) -> Result<DirectThreadSubscriptionBootstrap, String> {
) -> Result<SubscriptionBootstrap, String> {
tauri::async_runtime::spawn_blocking(move || {
let root = Path::new(project_path.trim());
enforce_project_permission_policy(root, "conversation.read")?;
let thread_id = direct_thread_id_for_project(root);
let mut bootstrap = subscribe_direct_thread(&thread_id);
let thread_id = thread_id_for_project(root);
let mut bootstrap = subscribe_thread(&thread_id);
if bootstrap.last_completed_item_id.is_none() {
match read_direct_project_last_item_id_at(root) {
Ok(last_completed_item_id) => {
@@ -5938,7 +5934,7 @@ pub(crate) async fn subscribe_direct_project_thread(
Err(error) => {
// 订阅已经登记,后续锚点读取失败时也必须回收 native
// subscriber;否则前端拿不到 subscriptionId,无法自行清理。
unsubscribe_direct_thread(&bootstrap.subscription_id);
unsubscribe_thread(&bootstrap.subscription_id);
return Err(error);
}
}
@@ -5952,13 +5948,13 @@ pub(crate) async fn subscribe_direct_project_thread(
#[tauri::command]
pub(crate) fn consume_direct_project_thread(
subscription_id: String,
) -> Result<DirectThreadConsumeResult, String> {
consume_direct_thread(subscription_id.trim())
) -> Result<ConsumeResult, String> {
consume_thread(subscription_id.trim())
}
#[tauri::command]
pub(crate) fn unsubscribe_direct_project_thread(subscription_id: String) -> Result<(), String> {
unsubscribe_direct_thread(subscription_id.trim());
unsubscribe_thread(subscription_id.trim());
Ok(())
}
@@ -5973,7 +5969,7 @@ pub(crate) async fn read_direct_project_history_slice(
before_item_id: Option<String>,
through_item_id: Option<String>,
limit: Option<usize>,
) -> Result<DirectThreadHistorySlice, String> {
) -> Result<HistorySlice, String> {
tauri::async_runtime::spawn_blocking(move || {
let root = Path::new(project_path.trim());
enforce_project_permission_policy(root, "conversation.read")?;
@@ -5990,12 +5986,12 @@ pub(crate) async fn read_direct_project_history_slice(
};
let (raw_items, has_more, recorded_at_ms, first_item_id) =
read_direct_project_history_items_slice_at(root, anchor, limit.unwrap_or(20))?;
let items = direct_thread_items_from_history(root, &raw_items, |item| {
direct_thread_item_identity(item)
let items = thread_items_from_history(root, &raw_items, |item| {
thread_item_identity(item)
.and_then(|identity| recorded_at_ms.get(&identity).copied())
.unwrap_or_default()
});
Ok(DirectThreadHistorySlice {
Ok(HistorySlice {
items,
has_more,
first_item_id,
@@ -2565,7 +2565,7 @@ fn main() {
})?;
setup_log.append("startup.runner.start.begin");
set_game_creator_agent_runtime_update_app_handle(app.handle().clone());
set_direct_thread_manager_app_handle(app.handle().clone());
set_thread_manager_app_handle(app.handle().clone());
platform_maintenance::set_platform_maintenance_app_handle(app.handle().clone());
auth_session::initialize_auth_session(app.handle());
set_asset_generation_task_app_handle(app.handle().clone());
@@ -35,9 +35,9 @@ import {
toDirectCodexTurnAttachments,
} from '../conversation/directCodexTurnAttachments';
import type { DirectCodexUserContentPart } from '../generated/DirectCodexUserContentPart';
import type { DirectQueueRemovalOutcome } from '../generated/DirectQueueRemovalOutcome';
import type { DirectThreadHistorySlice } from '../generated/DirectThreadHistorySlice';
import type { DirectThreadItem } from '../generated/DirectThreadItem';
import type { HistorySlice } from '../generated/HistorySlice';
import type { QueueRemovalOutcome } from '../generated/QueueRemovalOutcome';
import type { ThreadItem } from '../generated/ThreadItem';
import { directHistoryAnchorGateToWaitFor } from '../history/directHistoryAnchorGate';
import { readDirectHistoryPages } from '../history/directHistoryPaging';
import { canStopDirectProjectTurn } from './useDirectProjectTurnStatus';
@@ -324,7 +324,7 @@ export function useDirectProjectChatController({
return;
}
try {
const outcome = await invoke<DirectQueueRemovalOutcome>(
const outcome = await invoke<QueueRemovalOutcome>(
'remove_direct_project_pending_turn',
{ projectPath: nextProjectPath, clientTurnId },
);
@@ -679,7 +679,7 @@ export function useDirectProjectChatController({
}
function mergeHistoryPages(input: {
items: DirectThreadItem[];
items: ThreadItem[];
hasMore: boolean;
firstItemId: string | null;
}) {
@@ -709,18 +709,15 @@ export function useDirectProjectChatController({
existingEntries: [],
beforeItemId: null,
readSlice: (beforeItemId) =>
invoke<DirectThreadHistorySlice>(
'read_direct_project_history_slice',
{
projectPath: nextProjectPath,
...(beforeItemId
? { beforeItemId, limit: CONVERSATION_INITIAL_VISIBLE_COUNT }
: {
...(throughItemId ? { throughItemId } : {}),
limit: CONVERSATION_INITIAL_VISIBLE_COUNT,
}),
},
),
invoke<HistorySlice>('read_direct_project_history_slice', {
projectPath: nextProjectPath,
...(beforeItemId
? { beforeItemId, limit: CONVERSATION_INITIAL_VISIBLE_COUNT }
: {
...(throughItemId ? { throughItemId } : {}),
limit: CONVERSATION_INITIAL_VISIBLE_COUNT,
}),
}),
});
if (projectPathRef.current !== nextProjectPath) return;
mergeHistoryPages(pages);
@@ -751,15 +748,12 @@ export function useDirectProjectChatController({
existingEntries: directEntries,
beforeItemId: historyOldestItemIdRef.current,
readSlice: (beforeItemId) =>
invoke<DirectThreadHistorySlice>(
'read_direct_project_history_slice',
{
projectPath: nextProjectPath,
...(beforeItemId
? { beforeItemId, limit: DIRECT_HISTORY_PAGE_SIZE }
: { limit: DIRECT_HISTORY_PAGE_SIZE }),
},
),
invoke<HistorySlice>('read_direct_project_history_slice', {
projectPath: nextProjectPath,
...(beforeItemId
? { beforeItemId, limit: DIRECT_HISTORY_PAGE_SIZE }
: { limit: DIRECT_HISTORY_PAGE_SIZE }),
}),
});
if (projectPathRef.current !== nextProjectPath) return;
directThread.mergeHistoryItems(pages.items);
@@ -20,9 +20,9 @@ import {
resolveDirectThreadBootstrap,
selectDirectChatEntries,
} from '../conversation/directThreadChat';
import type { DirectThreadConsumeResult } from '../generated/DirectThreadConsumeResult';
import type { DirectThreadItem } from '../generated/DirectThreadItem';
import type { DirectThreadSubscriptionBootstrap } from '../generated/DirectThreadSubscriptionBootstrap';
import type { ConsumeResult } from '../generated/ConsumeResult';
import type { SubscriptionBootstrap } from '../generated/SubscriptionBootstrap';
import type { ThreadItem } from '../generated/ThreadItem';
import {
type DirectHistoryAnchorGate,
reuseOrOpenDirectHistoryAnchorGate,
@@ -43,7 +43,7 @@ export type DirectThreadChatSubscription = {
/** 订阅回执锚点闸门:首屏历史读取靠它拿到 `lastCompletedItemId`。 */
anchorGateRef: MutableRefObject<DirectHistoryAnchorGate | null>;
/** 历史切片并入同一个 reducer:条目只有这一份事实源。 */
mergeHistoryItems: (items: readonly DirectThreadItem[]) => void;
mergeHistoryItems: (items: readonly ThreadItem[]) => void;
};
/**
@@ -111,7 +111,7 @@ export function useDirectThreadChatSubscription({
try {
do {
consumeAgain = false;
const result = await invoke<DirectThreadConsumeResult>(
const result = await invoke<ConsumeResult>(
'consume_direct_project_thread',
{ subscriptionId },
);
@@ -140,7 +140,7 @@ export function useDirectThreadChatSubscription({
}
}
async function bootstrap() {
const result = await invoke<DirectThreadSubscriptionBootstrap>(
const result = await invoke<SubscriptionBootstrap>(
'subscribe_direct_project_thread',
{ projectPath },
);
@@ -189,7 +189,7 @@ export function useDirectThreadChatSubscription({
}, [enabled, projectPath]);
const mergeHistoryItems = useMemo(
() => (items: readonly DirectThreadItem[]) => {
() => (items: readonly ThreadItem[]) => {
setState((current) => mergeDirectHistoryItems(current, items));
},
[],
@@ -13,11 +13,11 @@
*/
import type { GameCreatorDirectToolCall } from '../../../../app/types';
import type { DirectThreadConsumeResult } from '../generated/DirectThreadConsumeResult';
import type { DirectThreadEvent } from '../generated/DirectThreadEvent';
import type { DirectThreadHistorySlice } from '../generated/DirectThreadHistorySlice';
import type { DirectThreadItem } from '../generated/DirectThreadItem';
import type { DirectThreadSubscriptionBootstrap } from '../generated/DirectThreadSubscriptionBootstrap';
import type { ConsumeResult } from '../generated/ConsumeResult';
import type { HistorySlice } from '../generated/HistorySlice';
import type { SubscriptionBootstrap } from '../generated/SubscriptionBootstrap';
import type { ThreadEvent } from '../generated/ThreadEvent';
import type { ThreadItem } from '../generated/ThreadItem';
import {
type DirectPendingTurn,
enqueuePendingTurn,
@@ -143,7 +143,7 @@ function validBoundaryAt(value: number | null | undefined): number {
* 它与条目里的 `item.at` 是两件事:后者是条目展示时间,不作工具计时的起止边界。
* 字段由生成绑定声明(ts-rs),重放沿用原值,因此这里只读不取当前时间。
*/
export function readDirectThreadEventAt(event: DirectThreadEvent): number {
export function readDirectThreadEventAt(event: ThreadEvent): number {
return validBoundaryAt('at' in event ? event.at : 0);
}
@@ -153,9 +153,7 @@ export function readDirectThreadEventAt(event: DirectThreadEvent): number {
* 字段可选:缺失表示身份不可证明(旧事件保持原有顺序语义),此时不猜历史归属,也不拿
* 时间戳近似。
*/
export function readDirectThreadEventUserItemId(
event: DirectThreadEvent,
): string {
export function readDirectThreadEventUserItemId(event: ThreadEvent): string {
const value = 'userItemId' in event ? event.userItemId : '';
return typeof value === 'string' ? value.trim() : '';
}
@@ -318,7 +316,7 @@ function stampTurnUserEntry(
function appendLiveText(
state: DirectThreadChatState,
event: Extract<DirectThreadEvent, { type: 'item.delta' }>,
event: Extract<ThreadEvent, { type: 'item.delta' }>,
): DirectThreadChatState {
const itemId = event.itemId.trim();
if (!itemId || !event.delta) return state;
@@ -334,7 +332,7 @@ function appendLiveText(
export function reduceDirectThreadEvent(
state: DirectThreadChatState,
event: DirectThreadEvent,
event: ThreadEvent,
): DirectThreadChatState {
switch (event.type) {
case 'turn.started': {
@@ -532,7 +530,7 @@ export function directThreadTurnMatchesUser(
export function reduceDirectThreadEvents(
state: DirectThreadChatState,
events: readonly DirectThreadEvent[],
events: readonly ThreadEvent[],
): DirectThreadChatState {
return events.reduce(reduceDirectThreadEvent, state);
}
@@ -545,7 +543,7 @@ export function reduceDirectThreadEvents(
*/
export function resolveDirectThreadBootstrap(
state: DirectThreadChatState,
bootstrap: DirectThreadSubscriptionBootstrap,
bootstrap: SubscriptionBootstrap,
): DirectThreadChatState {
return reduceDirectThreadEvents(state, bootstrap.events);
}
@@ -553,7 +551,7 @@ export function resolveDirectThreadBootstrap(
/** 事件顺序 = 游标顺序;调用方只需要把 `consume` 的结果喂进来。 */
export function applyDirectThreadConsumeResult(
state: DirectThreadChatState,
result: DirectThreadConsumeResult,
result: ConsumeResult,
): DirectThreadChatState {
return reduceDirectThreadEvents(state, result.events);
}
@@ -580,7 +578,7 @@ export function mergeHistoryEntries(
/** 历史切片条目 → 聊天条目:可见性判定的唯一入口,分页判据也读这一份。 */
export function projectDirectHistoryItems(
items: readonly DirectThreadItem[],
items: readonly ThreadItem[],
): DirectChatEntry[] {
return items
.map((item) => projectDirectThreadItem(item))
@@ -594,7 +592,7 @@ export function projectDirectHistoryItems(
*/
export function mergeDirectHistoryItems(
state: DirectThreadChatState,
items: readonly DirectThreadItem[],
items: readonly ThreadItem[],
): DirectThreadChatState {
const history = mergeHistoryEntries(
projectDirectHistoryItems(items),
@@ -619,7 +617,7 @@ export function mergeDirectHistoryItems(
export function mergeDirectThreadHistorySlice(
state: DirectThreadChatState,
slice: DirectThreadHistorySlice,
slice: HistorySlice,
): DirectThreadChatState {
return mergeDirectHistoryItems(state, slice.items);
}

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