合并 origin/master(33 个提交):DirectProject 入队化与事件投影
Project CI / AI game creator shell Rust crates (pull_request) Successful in 1m30s
Project CI / AI game creator shell Rust lane 1/2 (pull_request) Failing after 1m54s
Project CI / AI game creator shell Rust smoke (pull_request) Successful in 1m57s
Project CI / Frontend tests (pull_request) Successful in 3m39s
Project CI / Backend tests (pull_request) Successful in 6m13s
Project CI / AI game creator shell Rust lane 2/2 (pull_request) Failing after 7m59s
Project CI / Repository checks (pull_request) Successful in 3m29s
Project CI / Native shell tests (pull_request) Successful in 8m0s
Project CI / AI game creator shell web tests (pull_request) Successful in 2m58s

- 合并 origin/master(425fbcf63..9ee34809b),唯一冲突是共享决策日志两侧各插条目,取「两侧并存」
- 本轮 master 未触及随包资源声明、准备步骤与三份 tauri 配置,build.rs 只读校验保持不变
- 决策日志与踩坑记录同步两侧条目;docs/README.md 与项目索引按 master 收敛
- 构建期重新生成的 ts-rs 绑定经 prettier 归一后与 master 逐字节一致(合并未丢类型)
- 校验:cargo test --no-run 全 target 通过、AGC 应用 vitest 191 files / 1889 tests passed、check-package-layout、prepare-bundled-resources 18 passed、check-config、check:encoding、prettier、cargo fmt --check 全绿
This commit is contained in:
2026-09-30 17:38:35 +08:00
90 changed files with 5473 additions and 5291 deletions
File diff suppressed because it is too large Load Diff
@@ -29,12 +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_wire;
mod direct_tool_bridge;
mod direct_tool_calls;
mod direct_tools_mcp;
mod direct_turn_accept;
mod direct_turn_error;
mod direct_turn_failure;
mod direct_turn_stream;
@@ -49,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,
@@ -58,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::{
@@ -71,12 +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_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_error::*;
pub(crate) use direct_turn_failure::*;
pub(crate) use direct_turn_stream::*;
@@ -91,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();
};
File diff suppressed because it is too large Load Diff
@@ -5,8 +5,7 @@ mod validation;
mod wire;
pub(crate) use model::DirectCodexUserItem;
pub(crate) use validation::validate_direct_codex_user_item;
pub(crate) use wire::{
direct_codex_user_item_to_codex_turn_input, direct_codex_user_item_to_prompt,
direct_codex_user_item_to_response_item,
direct_codex_user_item_to_response_item, freeze_direct_codex_user_item,
};
@@ -2,7 +2,7 @@ use serde::{Deserialize, Serialize};
use ts_rs::TS;
/// DirectProject 本轮 user input 的唯一结构化入口。
#[derive(Clone, Debug, Deserialize, Serialize, TS)]
#[derive(Clone, Debug, PartialEq, Deserialize, Serialize, TS)]
#[serde(tag = "type", deny_unknown_fields)]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) enum DirectCodexUserItem {
@@ -10,7 +10,7 @@ pub(crate) enum DirectCodexUserItem {
Message(DirectCodexUserMessageItem),
}
#[derive(Clone, Debug, Deserialize, Serialize, TS)]
#[derive(Clone, Debug, PartialEq, Deserialize, Serialize, TS)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) struct DirectCodexUserMessageItem {
@@ -19,21 +19,28 @@ pub(crate) struct DirectCodexUserMessageItem {
pub(crate) id: String,
}
#[derive(Clone, Debug, Deserialize, Serialize, TS)]
#[derive(Clone, Debug, PartialEq, Deserialize, Serialize, TS)]
#[serde(rename_all = "lowercase")]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) enum DirectCodexUserRole {
User,
}
#[derive(Clone, Debug, Deserialize, Serialize, TS)]
#[derive(Clone, Debug, PartialEq, Deserialize, Serialize, TS)]
#[serde(tag = "type", rename_all_fields = "camelCase", deny_unknown_fields)]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) enum DirectCodexUserContentPart {
#[serde(rename = "input_text")]
InputText { text: String },
#[serde(rename = "agc_resource_reference")]
AgcResourceReference { resource_id: String },
AgcResourceReference {
resource_id: String,
/// 这个引用被宿主解析出来的文本:素材摘要,加上引用的 UI 设计文档展开的代码上下文。
/// 入队冻结时写入一次,此后随条目持久化,prompt / 历史回读 / turn input 共用这一份。
/// 缺省合法(旧历史没有它照样解析),缺省时才退回按当前 manifest 现算。
#[serde(default, skip_serializing_if = "Option::is_none")]
resolved_text: Option<String>,
},
#[serde(rename = "agc_skill_reference")]
AgcSkillReference { name: String },
#[serde(rename = "agc_runtime_region_reference")]
@@ -43,7 +50,7 @@ pub(crate) enum DirectCodexUserContentPart {
AgcAttachmentReference(DirectCodexUserAttachmentReferencePart),
}
#[derive(Clone, Debug, Deserialize, Serialize, TS)]
#[derive(Clone, Debug, PartialEq, Deserialize, Serialize, TS)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) struct DirectCodexUserAttachmentReferencePart {
@@ -55,7 +62,7 @@ pub(crate) struct DirectCodexUserAttachmentReferencePart {
pub(crate) status: String,
}
#[derive(Clone, Debug, Deserialize, Serialize, TS)]
#[derive(Clone, Debug, PartialEq, Deserialize, Serialize, TS)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) struct DirectCodexUserRuntimeRegionPart {
@@ -89,6 +96,7 @@ mod tests {
role: DirectCodexUserRole::User,
content: vec![DirectCodexUserContentPart::AgcResourceReference {
resource_id: "asset-hero".to_string(),
resolved_text: None,
}],
id: "turn-1".to_string(),
});
@@ -122,6 +130,32 @@ mod tests {
assert!(error.to_string().contains("unknown field"), "{error}");
}
/// 旧历史里的资源引用没有 `resolvedText`:缺省必须合法,照样解析成 `None` 的引用。
#[test]
fn resource_reference_without_resolved_text_still_parses() {
let item: DirectCodexUserItem = serde_json::from_value(json!({
"type": "message",
"role": "user",
"content": [{
"type": "agc_resource_reference",
"resourceId": "asset-hero"
}],
"id": "turn-1"
}))
.expect("legacy resource reference must keep parsing");
let DirectCodexUserItem::Message(message) = item;
match &message.content[0] {
DirectCodexUserContentPart::AgcResourceReference {
resource_id,
resolved_text,
} => {
assert_eq!(resource_id, "asset-hero");
assert_eq!(resolved_text, &None);
}
other => panic!("expected a resource reference, got {other:?}"),
}
}
#[test]
fn skill_reference_serializes_with_only_the_stable_name() {
let item: DirectCodexUserItem = serde_json::from_value(json!({
@@ -37,7 +37,7 @@ pub(crate) fn validate_direct_codex_user_item(
for part in &message.content {
match part {
DirectCodexUserContentPart::InputText { .. } => {}
DirectCodexUserContentPart::AgcResourceReference { resource_id } => {
DirectCodexUserContentPart::AgcResourceReference { resource_id, .. } => {
reference_count = reference_count.saturating_add(1);
validate_resource_id_and_manifest(&manifest, resource_id)?;
}
@@ -202,6 +202,7 @@ mod tests {
assert!(content_has_meaningful_input(&[
DirectCodexUserContentPart::AgcResourceReference {
resource_id: "asset-hero".to_string(),
resolved_text: None,
},
]));
}
@@ -4,8 +4,8 @@ use super::model::{
};
use super::validation::validate_direct_codex_user_item;
use crate::agent::{
read_manifest_for_project, sanitize_attachment_local_path, sanitize_attachment_media_type,
sanitize_attachment_name, GameCreationAppManifest,
sanitize_attachment_local_path, sanitize_attachment_media_type, sanitize_attachment_name,
GameCreationAppManifest,
};
use crate::ui_editor::persistence::{
generate_ui_design_code_at, GenerateUiDesignCodeInput, UI_DESIGN_DOC_ASSET_KIND,
@@ -14,6 +14,53 @@ use crate::ui_editor::persistence::{
use serde_json::Value;
use std::path::Path;
/// 入队检查产出的**冻结条目**:校验 → 把每个引用 part 的解析文本写进它自己的
/// [`DirectCodexUserContentPart::AgcResourceReference::resolved_text`] → 判空。
///
/// 这是这条消息**唯一**会算片段、会写盘(UI 设计文档代码导出)的地方,也是唯一的失败出口:
/// 之后的 prompt 折叠([`direct_codex_user_item_to_prompt`])与放行都不再重算、不再读 manifest,
/// 因此也没有失败可言。冻结结果随条目一路走到历史、事件与 turn input——它们读的都是同一份文本。
pub(crate) fn freeze_direct_codex_user_item(
root: &Path,
item: &DirectCodexUserItem,
) -> Result<DirectCodexUserItem, String> {
let manifest = validate_direct_codex_user_item(root, item)?;
let DirectCodexUserItem::Message(message) = item;
let mut content = Vec::with_capacity(message.content.len());
for part in &message.content {
match part {
DirectCodexUserContentPart::AgcResourceReference {
resource_id,
resolved_text,
} => {
// 客户端不产出这个字段;它只可能来自"已冻结的条目"(旧历史缺省)。缺省才现算。
let resolved_text = match resolved_text {
Some(text) => Some(text.clone()),
None => Some(resolved_resource_reference_text(
root,
&manifest,
resource_id,
)?),
};
content.push(DirectCodexUserContentPart::AgcResourceReference {
resource_id: resource_id.clone(),
resolved_text,
});
}
other => content.push(other.clone()),
}
}
let frozen = DirectCodexUserItem::Message(DirectCodexUserMessageItem {
role: message.role.clone(),
content,
id: message.id.clone(),
});
if direct_codex_user_item_to_prompt(&frozen).trim().is_empty() {
return Err("DirectProject user item 不能转换为空 prompt".to_string());
}
Ok(frozen)
}
/// 将历史中的 canonical user item 投影为 Codex `response_item` message。
/// 非 user message 的标准 Response item 原样返回;未知形状直接失败。
pub(crate) fn direct_codex_user_item_to_response_item(
@@ -134,9 +181,15 @@ pub(crate) fn direct_codex_user_item_to_wire_input(
for part in &message.content {
let text = match part {
DirectCodexUserContentPart::InputText { text } => text.clone(),
DirectCodexUserContentPart::AgcResourceReference { resource_id } => {
resource_reference_summary(&manifest, resource_id)?
}
DirectCodexUserContentPart::AgcResourceReference {
resource_id,
resolved_text,
} => match resolved_text {
// 冻结过的条目读自己那一份:历史回放因此就是当初发出去的那条消息。
Some(text) => text.clone(),
// 旧历史没有这份文本,退回按当前 manifest 现算(只有摘要,不产生写副作用)。
None => resource_reference_summary(&manifest, resource_id)?,
},
DirectCodexUserContentPart::AgcSkillReference { name } => {
format!("${}", name.trim())
}
@@ -165,10 +218,16 @@ pub(crate) fn direct_codex_user_item_to_codex_turn_input(
DirectCodexUserContentPart::InputText { text } => {
input.push(serde_json::json!({ "type": "text", "text": text }));
}
DirectCodexUserContentPart::AgcResourceReference { resource_id } => {
DirectCodexUserContentPart::AgcResourceReference {
resource_id,
resolved_text,
} => {
input.push(serde_json::json!({
"type": "text",
"text": resource_reference_summary(&manifest, resource_id)?,
"text": match resolved_text {
Some(text) => text.clone(),
None => resource_reference_summary(&manifest, resource_id)?,
},
}));
}
DirectCodexUserContentPart::AgcSkillReference { name } => {
@@ -201,82 +260,88 @@ pub(crate) fn direct_codex_user_item_to_codex_turn_input(
Ok(Value::Array(input))
}
pub(crate) fn direct_codex_user_item_to_prompt(
root: &Path,
item: &DirectCodexUserItem,
) -> Result<String, String> {
let wire = direct_codex_user_item_to_wire_input(root, item)?;
/// canonical 条目到 prompt 的**纯折叠**:零 IO、零校验、无失败出口。
///
/// 它只读条目自己的事实——文本 part 的正文与引用 part 的 `resolved_text`(入队时
/// [`freeze_direct_codex_user_item`] 冻结,随条目持久化)。条目里没有引用之外的东西要算,
/// 所以这里也不再需要 root。调用方要保证条目是冻结过的;缺省 `resolved_text` 的旧条目
/// 走不了这一条(历史回读用 [`direct_codex_user_item_to_response_item`])。
pub(crate) fn direct_codex_user_item_to_prompt(item: &DirectCodexUserItem) -> String {
let DirectCodexUserItem::Message(message) = item;
let mut prompt = wire
.as_array()
.ok_or_else(|| "DirectProject user item wire input 不是数组".to_string())?
.iter()
.filter_map(|part| part.get("text").and_then(Value::as_str))
.collect::<String>();
if let Some(code_context) = render_ui_design_code_context(root, message)? {
prompt.push('\n');
prompt.push_str(&code_context);
let mut prompt = String::new();
for part in &message.content {
match part {
DirectCodexUserContentPart::InputText { text } => prompt.push_str(text),
DirectCodexUserContentPart::AgcResourceReference { resolved_text, .. } => {
if let Some(text) = resolved_text {
prompt.push_str(text);
}
}
DirectCodexUserContentPart::AgcSkillReference { name } => {
prompt.push_str(&format!("${}", name.trim()));
}
DirectCodexUserContentPart::AgcRuntimeRegionReference(reference) => {
prompt.push_str(&runtime_region_summary(reference));
}
DirectCodexUserContentPart::AgcAttachmentReference(reference) => {
prompt.push_str(&attachment_reference_summary(reference));
}
}
}
if prompt.trim().is_empty() {
return Err("DirectProject user item 不能转换为空 prompt".to_string());
}
Ok(prompt)
prompt
}
/// 本轮 prompt 的 UI 设计文档引用上下文:复用 UI Editor 代码导出,把带文档注释的
/// JS 片段路径交给模型;生成失败只追加原始错误,不阻断本轮引用,其它引用继续处理。
/// 一个资源引用 part 的解析文本:素材摘要,加上引用的 UI 设计文档展开的代码上下文。
///
/// 只在生成本轮 prompt 时展开:历史 item 回读走 `direct_codex_user_item_to_response_item`
/// 的纯投影,不得在这里产生项目写副作用。
fn render_ui_design_code_context(
/// 只在 [`freeze_direct_codex_user_item`] 里现算一次;算出来的就是这一 part 此后唯一的文本。
fn resolved_resource_reference_text(
root: &Path,
message: &DirectCodexUserMessageItem,
) -> Result<Option<String>, String> {
let referenced_ids = message
.content
.iter()
.filter_map(|part| match part {
DirectCodexUserContentPart::AgcResourceReference { resource_id } => {
Some(resource_id.trim())
}
_ => None,
})
.collect::<Vec<_>>();
if referenced_ids.is_empty() {
return Ok(None);
manifest: &GameCreationAppManifest,
resource_id: &str,
) -> Result<String, String> {
let mut text = resource_reference_summary(manifest, resource_id)?;
if let Some(code_context) = ui_design_code_context_for(root, manifest, resource_id) {
text.push('\n');
text.push_str(&code_context);
}
let manifest = read_manifest_for_project(root)?;
let mut lines = Vec::new();
for resource_id in referenced_ids {
let is_ui_design_doc = manifest.assets.iter().any(|asset| {
asset.id == resource_id
&& asset.kind == UI_DESIGN_DOC_ASSET_KIND
&& asset.media_type == UI_DESIGN_DOC_MEDIA_TYPE
});
if !is_ui_design_doc {
continue;
}
lines.push(
match generate_ui_design_code_at(GenerateUiDesignCodeInput {
project_path: root.to_string_lossy().into_owned(),
expected_project_id: manifest.project_id.clone(),
asset_id: resource_id.to_string(),
}) {
Ok(result) => format!(
prompt_text!("projectContext.uiDesign.codeContext"),
result.relative_path
),
Err(error) => format!(
prompt_text!("projectContext.uiDesign.generationErrorContext"),
error = error
),
},
);
Ok(text)
}
/// 引用的 UI 设计文档上下文:复用 UI Editor 代码导出,把带文档注释的 JS 片段路径交给模型;
/// 生成失败只带回原始错误,不阻断本轮引用。非 UI 设计文档返回 `None`。
///
/// 这是整条链上唯一会写盘(生成 `ui/generated-*.js`)的投影处,因此只允许在入队冻结时调用。
/// 历史回读走 `direct_codex_user_item_to_response_item` 的纯投影,不得在这里产生项目写副作用。
fn ui_design_code_context_for(
root: &Path,
manifest: &GameCreationAppManifest,
resource_id: &str,
) -> Option<String> {
let resource_id = resource_id.trim();
let is_ui_design_doc = manifest.assets.iter().any(|asset| {
asset.id == resource_id
&& asset.kind == UI_DESIGN_DOC_ASSET_KIND
&& asset.media_type == UI_DESIGN_DOC_MEDIA_TYPE
});
if !is_ui_design_doc {
return None;
}
if lines.is_empty() {
return Ok(None);
}
Ok(Some(lines.join("\n")))
Some(
match generate_ui_design_code_at(GenerateUiDesignCodeInput {
project_path: root.to_string_lossy().into_owned(),
expected_project_id: manifest.project_id.clone(),
asset_id: resource_id.to_string(),
}) {
Ok(result) => format!(
prompt_text!("projectContext.uiDesign.codeContext"),
result.relative_path
),
Err(error) => format!(
prompt_text!("projectContext.uiDesign.generationErrorContext"),
error = error
),
},
)
}
#[cfg(test)]
@@ -284,7 +349,7 @@ mod tests {
use super::{
direct_codex_user_item_to_codex_turn_input, direct_codex_user_item_to_prompt,
direct_codex_user_item_to_response_item, direct_codex_user_item_to_wire_input,
validate_direct_codex_user_item,
freeze_direct_codex_user_item, validate_direct_codex_user_item,
};
use crate::agent::direct_codex_user_item::model::DirectCodexUserItem;
use crate::ui_editor::persistence::UI_DESIGN_DOC_MEDIA_TYPE;
@@ -375,8 +440,8 @@ mod tests {
let item: super::DirectCodexUserItem =
serde_json::from_value(user_item_with_resource_reference(&asset_id))
.expect("canonical user item");
let prompt =
direct_codex_user_item_to_prompt(project.path(), &item).expect("prompt projection");
let frozen = freeze_direct_codex_user_item(project.path(), &item).expect("freeze");
let prompt = direct_codex_user_item_to_prompt(&frozen);
assert!(
prompt.contains(&format!(
"[素材引用 resourceId={asset_id};项目路径=ui/design.json]"
@@ -387,6 +452,14 @@ mod tests {
prompt.contains("请先阅读生成的带有文档的代码片段: ui/generated-"),
"{prompt}"
);
// 解析文本随条目持久化:历史回放读的就是这一份,不再重算、不再写盘。
let super::DirectCodexUserItem::Message(message) = &frozen;
match &message.content[0] {
super::DirectCodexUserContentPart::AgcResourceReference { resolved_text, .. } => {
assert_eq!(resolved_text.as_deref(), Some(prompt.as_str()));
}
other => panic!("expected a frozen resource reference, got {other:?}"),
}
}
#[test]
@@ -418,8 +491,8 @@ mod tests {
let item: super::DirectCodexUserItem =
serde_json::from_value(user_item_with_resource_reference(&asset_id))
.expect("canonical user item");
let prompt =
direct_codex_user_item_to_prompt(project.path(), &item).expect("prompt projection");
let frozen = freeze_direct_codex_user_item(project.path(), &item).expect("freeze");
let prompt = direct_codex_user_item_to_prompt(&frozen);
assert!(prompt.contains("素材引用 resourceId="), "{prompt}");
assert!(prompt.contains("生成代码遇到错误"), "{prompt}");
}
@@ -444,8 +517,8 @@ mod tests {
],
}))
.expect("canonical user item");
let prompt =
direct_codex_user_item_to_prompt(project.path(), &item).expect("prompt projection");
let frozen = freeze_direct_codex_user_item(project.path(), &item).expect("freeze");
let prompt = direct_codex_user_item_to_prompt(&frozen);
let ui_doc_at = prompt
.find(&failing_ui_doc_id)
@@ -477,8 +550,8 @@ mod tests {
let item: super::DirectCodexUserItem =
serde_json::from_value(user_item_with_resource_reference(&asset_id))
.expect("canonical user item");
let prompt =
direct_codex_user_item_to_prompt(project.path(), &item).expect("prompt projection");
let frozen = freeze_direct_codex_user_item(project.path(), &item).expect("freeze");
let prompt = direct_codex_user_item_to_prompt(&frozen);
assert!(prompt.contains("素材引用 resourceId="), "{prompt}");
assert!(
!prompt.contains("请先阅读生成的带有文档的代码片段"),
@@ -518,6 +591,42 @@ mod tests {
assert_ne!(projected["content"][0]["type"], "text");
}
/// 历史回读读条目自己的 `resolvedText`:引用被冻结过就是当初那条消息,没有就退回当前 manifest 的摘要。
#[test]
fn history_projection_reads_the_persisted_resolved_text() {
let root = prompt_context_project();
let asset_id = register_fixture_asset(
root.path(),
"assets/hero.png",
GameCreationAppAssetKind::Character,
"image/png",
);
for resolved_text in [None, Some("冻结时的解析文本".to_string())] {
let mut reference = json!({
"type": "agc_resource_reference",
"resourceId": asset_id,
});
if let Some(text) = &resolved_text {
reference["resolvedText"] = json!(text);
}
let item = json!({
"type": "message",
"role": "user",
"id": "turn-1:user",
"content": [reference],
});
let projected = direct_codex_user_item_to_response_item(root.path(), &item)
.expect("history projection");
let text = projected["content"][0]["text"].as_str().expect("text");
assert_eq!(
text,
resolved_text.unwrap_or_else(|| format!(
"[素材引用 resourceId={asset_id};项目路径=assets/hero.png]"
))
);
}
}
#[test]
fn text_projection_preserves_empty_parts_line_breaks_and_trailing_whitespace() {
let root = tempfile::tempdir().expect("temp project");
@@ -541,7 +650,7 @@ mod tests {
assert_eq!(projected["content"], item["content"]);
let canonical = serde_json::from_value(item).expect("canonical user item");
assert_eq!(
direct_codex_user_item_to_prompt(root.path(), &canonical).expect("multiline prompt"),
direct_codex_user_item_to_prompt(&canonical),
"你好\n第二段\n "
);
}
@@ -603,11 +712,9 @@ mod tests {
"type": "message", "role": "user", "id": "turn-1:user", "content": content
});
let canonical = serde_json::from_value(item.clone()).expect("canonical user item");
assert_eq!(
direct_codex_user_item_to_prompt(root.path(), &canonical)
.expect("reference prompt"),
expected_prompt
);
let frozen = freeze_direct_codex_user_item(root.path(), &canonical)
.expect("reference freeze");
assert_eq!(direct_codex_user_item_to_prompt(&frozen), expected_prompt);
let projected = direct_codex_user_item_to_response_item(root.path(), &item)
.expect("reference history projection");
assert_eq!(
@@ -432,7 +432,7 @@ pub(super) async fn begin(
let root = root.to_path_buf();
let prompt_hash = hash(prompt.as_bytes());
let session = tokio::task::spawn_blocking(move || {
let turn = super::direct_taonier_active_invocation_id_at(&root)?;
let turn = super::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,11 +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_taonier_active_invocation_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 = [
@@ -623,14 +619,9 @@ 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(
&direct_thread_id_for_project(&root),
"turn-before",
None,
)
.unwrap(),
TurnReservation::accept_for_test(&thread_id_for_project(&root), "turn-before"),
)));
let swap = Arc::clone(&owner);
let result = read_batch_with(
@@ -641,14 +632,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(TurnReservation::accept_for_test(
&thread_id_for_project(r),
"turn-after",
));
result
},
)
@@ -811,14 +798,9 @@ mod tests {
#[tokio::test]
async fn host_prefetch_keeps_data_out_of_system_rules_and_matches_active_turn() {
let (_temp, root) = project();
// 调用身份(预取闸门)与逻辑回合(上下文身份)是两件事,生产入口两步都做。
let _guard = DirectTaonierActiveInvocationGuard::enter(&root, "prefetch-turn").unwrap();
let _turn = DirectTurnReservation::accept(
&direct_thread_id_for_project(&root),
"prefetch-turn",
None,
)
.unwrap();
// 预取闸门与上下文身份读的是同一处:Thread Manager 上这一轮的活动回合登记。
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();
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
@@ -267,7 +267,7 @@ impl DirectToolBridge {
impl DirectToolBridgeState {
fn begin_user_turn(self: &Arc<Self>) -> Result<DirectToolBridgeTurnGuard, String> {
let turn_id = direct_taonier_active_invocation_id_at(&self.root)?;
let turn_id = active_turn_id_at(&self.root)?;
let mut authorization = self
.turn_authorization
.lock()
@@ -4202,6 +4202,14 @@ mod tests {
let temporary = tempfile::tempdir().unwrap();
let root = temporary.path().join("project");
let config = temporary.path().join("config");
// 成绩由宿主在终态直写:这一轮跑在哪个身份下,就看放行代次与终态代次是否同代。
let _platform_session = crate::platform_session::install_test_platform_session(
"A",
"token-a",
"https://dev.genarrative.world",
);
let release_identity_generation =
crate::platform_session::current_platform_session_write_state().identity_generation;
let (metadata, context, writer) = analytics_test_writer(&config);
let original_capture = Some((metadata.context.clone(), writer.clone()));
let mut lifecycle = crate::analytics::gui::LifecycleFixture::start(
@@ -4264,13 +4272,12 @@ mod tests {
);
assert_eq!(session.analytics_output_revision(), Some(revision.clone()));
lease.finish(true, true, None).unwrap();
let attempt = uuid::Uuid::new_v4().to_string();
run::direct_finished(
Some((context.clone(), writer.clone())),
&root,
&project_id,
&metadata,
Some(&attempt),
release_identity_generation,
run::Outcome {
turn_id: Some("analytics-write".into()),
end_reason: RunEndReason::Failed,
@@ -4280,7 +4287,6 @@ mod tests {
revision_id: session.analytics_output_revision(),
},
);
run::settle(Some((context.clone(), writer.clone())), &attempt, false);
// 默认项目使用 npm:预览服务读取真实构建目录,夹具提供构建入口,不调用构建器。
let served_root = crate::project_game_root(&root);
fs::create_dir_all(&served_root).unwrap();
@@ -4840,8 +4846,7 @@ mod tests {
assert!(state.begin_user_turn().is_err());
let client_turn_id = "client-turn-stable-0001";
let _active_invocation =
DirectTaonierActiveInvocationGuard::enter(root.path(), client_turn_id)
.expect("client-owned stable invocation");
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};
@@ -3845,19 +3845,19 @@ mod tests {
(handle, receiver)
}
/// 真实项目 + 真实工具桥 + 真实 Direct 回合身份。
/// 真实项目 + 真实工具桥 + 真实 Direct 回合占用(这一轮的调用身份)。
async fn tool_chain_start(
root: &Path,
) -> (
super::super::direct_tool_bridge::DirectToolBridge,
DirectTaonierActiveInvocationGuard,
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 = DirectTaonierActiveInvocationGuard::enter(root, "tool-chain-turn")
.expect("arm tool chain direct 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;
@@ -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);
}
}
@@ -5,16 +5,16 @@
//! 通道断开要带宿主诊断。用不同变体各带各的字段,分流靠 `match`,不靠 `kind` 字段 + 共用字段的
//! 伪结构化,也不靠对错误文本做子串匹配。
//!
//! 走哪条通道由**发生位置**决定,不由错误种类决定(`接单化` 之后的口径):
//! - **接单前**发生的 = 拒单:只出提示 / 横幅,不做失败载荷、不写失败诊断、不上报成
//! 走哪条通道由**发生位置**决定,不由错误种类决定(`入队化` 之后的口径):
//! - **入队前**发生的 = 入队失败:只出提示 / 横幅,不做失败载荷、不写失败诊断、不上报成
//! "智能创作失败"。命令返回 `Err` 的就是这一类。
//! - **接单后**发生的 = 回合失败:事件载荷、横幅、应用日志、错误上报池四处一致;命令早已返回
//! - **放行后**发生的 = 回合失败:事件载荷、横幅、应用日志、错误上报池四处一致;命令早已返回
//! `Ok`,所以一律由宿主侧的占用对象投影成 `turn.completed.failure`。
//!
//! 所以"同一种错误在接单前后走不同通道"是正常的:[`DirectTurnError::EnvironmentNotReady`] 两边
//! 所以"同一种错误在入队与放行走不同通道"是正常的:[`DirectTurnError::EnvironmentNotReady`] 两边
//! 都可能出现,位置说了算。这里**没有**、也不该有"这个变体是不是回合失败"的判据。
//!
//! 事件载荷(`direct_thread_wire::DirectTurnFailure`)仍然只有 `{kind, message}` 两个字段:
//! 事件载荷(`thread_manager::wire::DirectTurnFailure`)仍然只有 `{kind, message}` 两个字段:
//! 那是**线上协议**,由 [`DirectTurnError::wire_kind`] 与 `Display` 在这一个出口投影出来,不是
//! 另一种状态模型。跨进程边界(`#[tauri::command]`)同样只给前端一个字符串:那是**序列化**,
//! 由 `Display` 一处生成;Rust 侧任何地方都不再解析这个字符串。
@@ -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,18 +385,15 @@ impl DirectModelCallKind {
)]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) enum DirectTurnError {
// ───────── 拒单:接单之前发生,这一轮没有开始 ─────────
// ───────── 入队失败:这一轮没有开始,也没有进队列 ─────────
/// `clientTurnId` 没给:没有稳定回合身份,拒绝创建可计费身份。
ClientTurnIdMissing,
/// `clientTurnId` 形状非法:长度与字符集由宿主定,界面按同一份约束生成。
ClientTurnIdMalformed { min_chars: usize, max_chars: usize },
/// 同一项目已有另一条回合在跑(或同一 `clientTurnId` 并发复用)。
/// 待发消息队列已满:这一条没进队,等前面几条发完再发。
///
/// 两个身份都要带上,因为"撞的是哪一轮"决定界面该不该动当前回合。
TurnAlreadyRunning {
existing_invocation_id: String,
incoming_invocation_id: String,
},
/// 上限只落在宿主这一处(`MAX_PENDING_TURNS`),随载荷带出去,界面不自己数一份。
QueueFull { limit: usize },
/// 项目目录锚不定(符号链接 / 权限 / 目录被删)。
ProjectRootUnanchored { cause: String },
/// 项目目录不存在或不是绝对路径。
@@ -407,13 +404,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 +439,7 @@ pub(crate) enum DirectTurnError {
TurnFailedUnclassified { detail: String },
}
/// 拒单载荷:命令边界交给前端的**结构化拒绝**。
/// 入队失败载荷:命令边界交给前端的**结构化入队失败**。
///
/// 为什么不是只给一句话:界面要按变体分流——认得的"前置条件不满足 / 用户参数无效"给一条与用户
/// 消息同级的提示且不上报,认不得的原样抛出交给既有捕获链路。文案只是给人看的最后一步,仍由
@@ -450,14 +447,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 +464,12 @@ impl DirectTurnRejection {
}
impl DirectTurnError {
/// 命令边界要不要为这条**拒单**补一份运行错误诊断。
/// 命令边界要不要为这条**入队失败**补一份运行错误诊断。
///
/// 只有"宿主 / 环境的事实故障、用户自己改不了"才值得进 `.agent/runtime/errors` 与应用日志;
/// 空内容、`clientTurnId` 形状、另一轮在跑、权限策略、目录锚不定 / 不是绝对路径都是用户自己
/// 就能修的操作结果,留痕只会变成噪声;它们仍按 `Display` 给用户一句可读的话。判据按变体分,
/// 不看文案。
/// 空内容、`clientTurnId` 形状、队列满、另一轮在跑、权限策略、目录锚不定 / 不是绝对路径都是
/// 用户自己就能修的操作结果,留痕只会变成噪声;它们仍按 `Display` 给用户一句可读的话。
/// 判据按变体分,不看文案。
///
/// 回合级失败恒为 `false`:它们在上游(`record_direct_codex_failure`)已经写过诊断,边界再写一次
/// 就是同一件事留两份。
@@ -482,7 +479,7 @@ impl DirectTurnError {
Self::EnvironmentNotReady { .. } | Self::HostStateUnavailable { .. } => true,
Self::ClientTurnIdMissing
| Self::ClientTurnIdMalformed { .. }
| Self::TurnAlreadyRunning { .. }
| Self::QueueFull { .. }
// 项目目录锚不定(符号链接 / 权限 / 目录被删)与目录不存在同类:都是用户能自己修好的
// 文件系统事实,诊断文案不该顶替那句"无法锚定 Direct 调用项目目录:{cause}"。
| Self::ProjectRootUnanchored { .. }
@@ -644,23 +641,10 @@ impl fmt::Display for DirectTurnError {
formatter,
"clientTurnId 必须为 {min_chars} 到 {max_chars} 位 ASCII 字母、数字或连字符,且首位必须为字母或数字"
),
Self::TurnAlreadyRunning {
existing_invocation_id,
incoming_invocation_id,
} => {
// 两条文案按身份是否相同分岔,但**不再有机器前缀**:界面按 `error.type` 与其
// 两个身份字段分流,不解析文案(前缀曾经是协议约定,现在只是噪声)。
if existing_invocation_id == incoming_invocation_id {
formatter.write_str(
"同一轮消息仍在处理中,已拒绝并发复用同一 clientTurnId;请等它结束或点「终止」后再发送",
)
} else {
write!(
formatter,
"当前项目已有另一条 Direct 客户端回合正在运行,已拒绝混用付费生成身份;可在输入盒点「终止」结束它,或等它结束后再发送"
)
}
}
Self::QueueFull { limit } => write!(
formatter,
"待发消息已达上限(最多 {limit} 条),请等前面几条发完再发送"
),
Self::ProjectRootUnanchored { cause } => {
write!(formatter, "无法锚定 Direct 调用项目目录:{cause}")
}
@@ -1013,11 +997,6 @@ mod tests {
policy_detail: "项目权限策略拒绝执行:conversation.write".into(),
}
.is_reportable());
assert!(!DirectTurnError::TurnAlreadyRunning {
existing_invocation_id: "turn-1".into(),
incoming_invocation_id: "turn-2".into(),
}
.is_reportable());
// 回合级失败在上游已经写过诊断。
assert!(
!DirectTurnError::turn_failed(DirectCodexFailureStage::CodeGeneration, "模型失败")
@@ -1025,26 +1004,6 @@ mod tests {
);
}
/// 并发拒绝的两条文案按身份是否相同分岔,身份必须原样带出来。
#[test]
fn concurrent_rejection_keeps_both_invocations_and_splits_the_copy() {
let same = DirectTurnError::TurnAlreadyRunning {
existing_invocation_id: "turn-1".into(),
incoming_invocation_id: "turn-1".into(),
};
assert!(same.to_string().contains("同一轮消息仍在处理中"));
let different = DirectTurnError::TurnAlreadyRunning {
existing_invocation_id: "turn-1".into(),
incoming_invocation_id: "turn-2".into(),
};
assert!(different
.to_string()
.contains("另一条 Direct 客户端回合正在运行"));
assert!(!different
.to_string()
.starts_with("direct-codex-turn-already-running:"));
}
/// 边界序列化:字符串只由 Display 生成,且与改造前的可见文本一致。
#[test]
fn boundary_serialization_uses_display() {
@@ -5,10 +5,10 @@
//! 1. [`direct_turn_terminal`]:拿这一轮的事实判定终态——是不是失败、原因是什么、状态写什么;
//! 2. [`DirectTurnTerminal::event`]:把终态投影成 `turn.completed` 事件。
//!
//! 终态的**出口**(谁写、什么时候兜底)不在这里,在 `direct_turn_accept.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_taonier_active_invocation_id_at(root)?;
let turn_id = active_turn_id_at(root)?;
let root = root.to_path_buf();
tokio::task::spawn_blocking(move || {
if direct_taonier_active_invocation_id_at(&root)? != turn_id {
if 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,
@@ -0,0 +1,391 @@
//! DirectProject 的放行:把队首那条待发消息变成一轮真的回合。
//!
//! 这个模块两件事,别再往里加第三件:
//! 1. [`TurnReservation`]:这一轮的占用对象——持有它代表还没收口,终态出口只认它的 token;
//! 2. [`kick_queue_dispatch`]:认领队首 → 起整轮(落盘用户条目 → 下发 → 跑回合)。
//!
//! 这一轮的**调用身份**(付费美术、执行会话、MCP、校验、上下文预取读的那个"这一轮是谁")由认领
//! 在 Thread Manager 上登记,认领成功即成立;这里不再另取一次、也不再持第二份身份。
//!
//! 放行**不重跑任何入队检查,也没有"放行失败"**:入队时的检查已经判过,这一轮一旦放行,之后的
//! 一切失败都是**回合失败**,由占用对象收口成 `turn.completed` 的失败载荷。设计见
//! `docs/adr/【ADR】DirectProject命令入队化与待发消息队列归宿主-2026-09-24.md`。
use std::path::{Path, PathBuf};
#[cfg(test)]
use crate::agent::PendingTurn;
use crate::agent::{
append_direct_project_user_message_at, claim_pending_turn, complete_turn_if_reserved,
direct_codex_user_item_to_prompt, direct_tool_call_now_ms, redact_agent_runtime_error,
run_direct_game_creator_turn_at_with_creation_type_and_emitter, thread_id_for_project,
DirectGameCreatorTurnUpdateEmitter, DirectTurnError, DirectTurnTerminal, DispatchedTurn,
};
/// 一次放行的占用。持有它就代表这一轮还没收口。
///
/// 生命周期由调用方决定:登记与开始事件已经由 [`kick_queue_dispatch`] 在同一临界区里写好,
/// 这里只接住这一轮的终态出口,并随放行任务一起 drop(正常 / 失败 / panic 都走同一个出口)。
pub(crate) struct TurnReservation {
thread_id: String,
token: String,
user_item_id: Option<String>,
/// Rust 侧会话保活:只在一条 Direct 回合存活期间运行,随这一轮的占用同生共死。
///
/// 渲染层不再按固定间隔触发续期(见 `useDirectProjectChatController` 的说明):会话保活的
/// 生命周期就是"这一轮在不在跑",占用登记时启动、占用释放(回合收口)时随本对象一起 drop。
_session_keepalive: Option<crate::auth_session::ClientSessionKeepalive>,
}
impl TurnReservation {
/// 接管一次放行的占用:占用登记与 `turn.started` 已由 Thread Manager 的认领写好。
pub(crate) fn resume(thread_id: &str, dispatched: &DispatchedTurn) -> Self {
Self {
thread_id: thread_id.to_string(),
token: dispatched.token.clone(),
user_item_id: dispatched.pending.user_item_id(),
_session_keepalive: crate::auth_session::spawn_client_session_keepalive(),
}
}
/// 放行之后还没走到深层终态就失败的收口口:只有这一轮仍被自己占用时才写。
///
/// 深层(真正跑完这一轮的代码)已经写出终态时返回 `false`,兜底不覆盖真实结果。
pub(crate) fn finish_if_unfinished(&self, terminal: DirectTurnTerminal) -> bool {
complete_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 = PendingTurn::new(
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"),
None,
direct_tool_call_now_ms(),
);
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 TurnReservation {
fn drop(&mut self) {
// 兜底:任务 panic、future 被丢弃、或今后在终态之前新增的 `?` 早退。
// 这类失败说不出原因,只给分类;能说清原因的错误必须由调用方在更早的地方显式收口。
let _ = self.finish_if_unfinished(DirectTurnTerminal::host_dropped());
// 占用的释放就是"这个项目腾出了跑回合的位置":踢一脚,队里排着的消息不必等下一次用户动作。
// panic 也走这里,所以队列不会因为一个任务炸掉而永久停住。
kick_queue_dispatch(Path::new(&self.thread_id));
}
}
/// 踢一脚:让队首那条待发消息有机会变成一轮真的回合。
///
/// 幂等,三个调用点:**入队之后**(队列空且没有回合在跑时,放行就是这一脚,不必再等一个调度周期)、
/// **占用释放**(正常 / 失败 / panic 共用,见 [`TurnReservation`] 的 `Drop`)、**中止路径**
/// (兜底释放时那一轮不会再有人替它收尾)。
///
/// 它什么都不返回:没有"放行失败"。踢不动就是现在不该跑——已经有回合在跑,或者队列是空的;
/// 这两种情况都会在下一次释放时再踢。
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;
};
// 放行这一刻的身份代次:这一轮的成绩只在"放行时与终态同代"时才入账(见 `analytics::run`)。
let release_identity_generation =
crate::platform_session::current_platform_session_write_state().identity_generation;
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, release_identity_generation)
.await;
});
}
/// 放行之后这一段:取这一轮的调用身份 → 落盘用户条目 → 下发 → 跑整轮。
///
/// 这一段的每个失败都是**回合失败**:占用与开始事件已经写好,只走占用对象的终态出口,不再回到任何
/// 命令的返回值上。
async fn run_dispatched_direct_turn(
root: PathBuf,
dispatched: DispatchedTurn,
reservation: TurnReservation,
release_identity_generation: u64,
) {
let turn_id = dispatched.pending.client_turn_id.clone();
// 放行只搬运条目自己的事实:canonical 形状与 prompt 都从冻结过的条目重投影,不重算、不写盘。
let canonical_user_item = serde_json::to_value(&dispatched.pending.user_item)
.expect("冻结过的 canonical 条目一定可序列化");
let prompt = direct_codex_user_item_to_prompt(&dispatched.pending.user_item);
// 落盘即回合成立:放行之后必须留下这条用户消息,哪怕这一轮随后失败。
if let Err(error) = append_direct_project_user_message_at(&root, &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, &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,
&prompt,
dispatched.pending.creation_type.as_deref(),
Some(&emitter),
Some(canonical_user_item.clone()),
capture,
release_identity_generation,
)
.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));
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::agent::{
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_thread(thread_id);
let _ = consume_thread(&bootstrap.subscription_id);
bootstrap.subscription_id
}
fn pending(subscription_id: &str) -> Vec<ThreadEvent> {
consume_thread(subscription_id).expect("consume").events
}
fn turn_completed_events(events: &[ThreadEvent]) -> Vec<&ThreadEvent> {
events
.iter()
.filter(|event| matches!(event, ThreadEvent::TurnCompleted { .. }))
.collect()
}
/// 一条待发消息。放行路径只读它自己的事实(身份、条目、创建类型),所以这里用最小可用形状。
fn pending_turn(client_turn_id: &str) -> PendingTurn {
PendingTurn::new(
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"),
None,
direct_tool_call_now_ms(),
)
}
/// 真实存在的项目目录不是这条用例的判据,用唯一的假路径当线程身份即可:认领与占用的语义只认
/// 这个字符串。
fn unique_thread(label: &str) -> String {
format!("dispatch-test-{label}-{}", Uuid::new_v4())
}
fn enqueue(thread_id: &str, client_turn_id: &str) {
enqueue_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_pending_turn(&thread).expect("claim head");
let owner = TurnReservation::resume(&thread, &dispatched);
let foreign = TurnReservation::resume(
&thread,
&DispatchedTurn {
token: "foreign-token".to_string(),
pending: pending_turn("turn-1"),
},
);
assert!(!foreign.finish_if_unfinished(DirectTurnTerminal::host_dropped()));
assert!(thread_turn_is_active(&thread));
assert!(owner.finish_if_unfinished(DirectTurnTerminal::host_dropped()));
assert!(!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_pending_turn(&thread).expect("claim head");
let reservation = TurnReservation::resume(&thread, &dispatched);
// 深层收口:真正跑完这一轮的代码算出来的终态。
let deep = ThreadEvent::turn_completed_failed(
DirectTurnFailure::new(
DirectTurnFailureKind::Timeout,
"等待模型回执超时".to_string(),
),
2_000,
)
.with_user_item_id(Some("direct-codex:turn-1:user"));
crate::agent::complete_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] {
ThreadEvent::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_pending_turn(&thread).expect("claim head");
let reservation = TurnReservation::resume(&thread, &dispatched);
drop(reservation);
// 第一条按宿主丢弃收口,第二条已经接上:放行不等人,也不必等下一次用户动作。
let events = pending(&subscription);
let terminal = events
.iter()
.position(|event| matches!(event, ThreadEvent::TurnCompleted { .. }))
.expect("第一条必须被收口");
assert!(
matches!(
events.get(terminal),
Some(ThreadEvent::TurnCompleted { status, failure: Some(failure), .. })
if status == "failed" && failure.kind == DirectTurnFailureKind::HostDropped
),
"{events:?}"
);
let next_started = events.iter().position(|event| {
matches!(
event,
ThreadEvent::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_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_pending_turn(&busy).expect("claim head");
let reservation = TurnReservation::resume(&busy, &dispatched);
kick_queue_dispatch(Path::new(&busy));
// 仍在跑的那一轮没有被顶掉:占用身份还是第一条。
assert!(thread_turn_is_active(&busy));
drop(reservation);
}
/// 前置判据:删掉调用身份守卫之后,"回合任务泄漏、app-server 侧没有可中断句柄"的兜底释放
/// 仍要能解开占用并把队首放行出去——`d833ca9d3` 的"重进会话被堵死"不得复活。
#[test]
fn stale_cancel_releases_the_occupancy_and_dispatches_the_next_pending_turn() {
let thread = unique_thread("stale-cancel");
let subscription = watch(&thread);
// 这一轮已经被放行(占用登记在 Thread Manager 上),但任务泄漏:拿不到可中断句柄。
let leaked = TurnReservation::accept_for_test(&thread, "turn-stale-1");
enqueue(&thread, "turn-stale-2");
// "从没进执行器"只对过了启动窗口的残留放行;真实用例等不了 60 秒。
crate::agent::backdate_direct_active_turn_for_test(&thread, 61_000);
let view =
crate::agent::cancel_direct_codex_turn_at(Path::new(&thread), Some("turn-stale-1"))
.expect("stale release");
assert_eq!(
view.outcome,
crate::agent::codex_app_server::DIRECT_TURN_CANCEL_OUTCOME_RELEASED
);
assert_eq!(view.client_turn_id, "turn-stale-1");
// 占用解开:第一条被兜底收口,第二条接上——终止这一脚就是"继续放行"。
let events = pending(&subscription);
assert!(
events.iter().any(|event| matches!(
event,
ThreadEvent::TurnCompleted { status, .. } if status == "aborted"
)),
"{events:?}"
);
assert!(
events.iter().any(|event| matches!(
event,
ThreadEvent::TurnStarted { user_item_id, .. }
if user_item_id.as_deref() == Some("direct-codex:turn-stale-2:user")
)),
"队首必须在兜底释放之后被放行:{events:?}"
);
drop(leaked);
}
}
File diff suppressed because it is too large Load Diff

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