Merge remote-tracking branch 'origin/master' into fix/web-preflight-browser-recovery
Project CI / AI game creator shell Rust crates (pull_request) Successful in 3m3s
Project CI / AI game creator shell Rust lane 1/2 (pull_request) Successful in 4m56s
Project CI / AI game creator shell Rust lane 2/2 (pull_request) Successful in 5m4s
Project CI / Backend tests (pull_request) Successful in 5m59s
Project CI / Frontend tests (pull_request) Successful in 2m52s
Project CI / AI game creator shell web tests (pull_request) Successful in 3m5s
Project CI / Native shell tests (pull_request) Successful in 6m37s
Project CI / Repository checks (pull_request) Successful in 6m24s
Project CI / AI game creator shell Rust crates (pull_request) Successful in 3m3s
Project CI / AI game creator shell Rust lane 1/2 (pull_request) Successful in 4m56s
Project CI / AI game creator shell Rust lane 2/2 (pull_request) Successful in 5m4s
Project CI / Backend tests (pull_request) Successful in 5m59s
Project CI / Frontend tests (pull_request) Successful in 2m52s
Project CI / AI game creator shell web tests (pull_request) Successful in 3m5s
Project CI / Native shell tests (pull_request) Successful in 6m37s
Project CI / Repository checks (pull_request) Successful in 6m24s
This commit is contained in:
@@ -164,6 +164,10 @@ module.exports = {
|
||||
'server-rs/target-*',
|
||||
'apps/desktop-shell/src-tauri/target',
|
||||
'apps/ai-game-creator-shell/src/features/ui-editor/types/**',
|
||||
// ts-rs 生成绑定:不接受 eslint --fix 二次改写,必须与原始输出逐字节一致
|
||||
'packages/shared/src/contracts/generated/**',
|
||||
'apps/ai-game-creator-shell/src/view/project-development/chat/generated/**',
|
||||
'apps/ai-game-creator-shell/src/services/generated/**',
|
||||
'apps/ai-game-creator-shell/src/features/project-workspace/generated/**',
|
||||
'target',
|
||||
'src/main.tsx',
|
||||
|
||||
+7
-2
@@ -24,5 +24,10 @@
|
||||
*.anim text
|
||||
*.controller text
|
||||
|
||||
# Rust ts-rs 生成的共享契约:保留在仓库中供 TS 消费,但不作为手写源文件统计。
|
||||
packages/shared/src/contracts/generated/** linguist-generated=true
|
||||
# ts-rs 生成绑定:保留在仓库中供 TS 消费,但不作为手写源文件统计。
|
||||
# ts-rs 原始输出在多行对象 / 枚举变体后带行尾空格,生成目录统一豁免行尾空白检查;
|
||||
# 一致性由 `npm run check:generated-bindings` 逐字节比对原始输出来保证。
|
||||
packages/shared/src/contracts/generated/** linguist-generated=true whitespace=-trailing-space
|
||||
apps/ai-game-creator-shell/src/features/ui-editor/types/** linguist-generated=true whitespace=-trailing-space
|
||||
apps/ai-game-creator-shell/src/view/project-development/chat/generated/** linguist-generated=true whitespace=-trailing-space
|
||||
apps/ai-game-creator-shell/src/services/generated/** linguist-generated=true whitespace=-trailing-space
|
||||
|
||||
@@ -3,7 +3,11 @@ node_modules
|
||||
.git
|
||||
.codex-logs
|
||||
public/Icons
|
||||
# ts-rs 生成绑定:一律以原始输出提交,禁止 prettier 二次改写
|
||||
packages/shared/src/contracts/generated/
|
||||
apps/ai-game-creator-shell/src/features/ui-editor/types/
|
||||
apps/ai-game-creator-shell/src/view/project-development/chat/generated/
|
||||
apps/ai-game-creator-shell/src/services/generated/
|
||||
apps/ai-game-creator-shell/src-tauri/resources/agc-skills/
|
||||
media
|
||||
*.log
|
||||
|
||||
+1
-10
@@ -1,14 +1,5 @@
|
||||
{
|
||||
"singleQuote": true,
|
||||
"semi": true,
|
||||
"trailingComma": "all",
|
||||
"overrides": [
|
||||
{
|
||||
"files": "packages/shared/src/contracts/generated/**/*.ts",
|
||||
"options": {
|
||||
"printWidth": 1000,
|
||||
"singleQuote": false
|
||||
}
|
||||
}
|
||||
]
|
||||
"trailingComma": "all"
|
||||
}
|
||||
|
||||
@@ -29,8 +29,6 @@ mod direct_project_turn_history;
|
||||
mod direct_runtime;
|
||||
mod direct_tool_bridge;
|
||||
mod direct_tools_mcp;
|
||||
mod direct_turn_error;
|
||||
mod direct_turn_failure;
|
||||
mod direct_validation;
|
||||
mod generation;
|
||||
pub(crate) mod json_sidecar;
|
||||
@@ -50,11 +48,12 @@ pub(crate) use claude_code_cli::{
|
||||
game_creator_claude_code_cli_version_identity,
|
||||
};
|
||||
pub(crate) use codex_app_server::direct_game_creator_codex_chat_at;
|
||||
pub(crate) use codex_app_server::turn_error::*;
|
||||
use codex_app_server::*;
|
||||
#[cfg(not(test))]
|
||||
pub(crate) use codex_app_server::{
|
||||
cancel_direct_codex_turn_at, direct_game_creator_home_codex_chat, thread_id_for_project,
|
||||
DirectTurnCancelView,
|
||||
TurnCancelView,
|
||||
};
|
||||
use codex_cli::*;
|
||||
pub(crate) use codex_cli::{
|
||||
@@ -72,8 +71,6 @@ pub(crate) use direct_project_turn_history::*;
|
||||
pub(crate) use direct_runtime::*;
|
||||
pub(crate) use direct_tool_bridge::*;
|
||||
pub(crate) use direct_tools_mcp::*;
|
||||
pub(crate) use direct_turn_error::*;
|
||||
pub(crate) use direct_turn_failure::*;
|
||||
pub(crate) use direct_validation::DirectValidationConfig;
|
||||
pub(crate) use generation::*;
|
||||
pub(crate) use json_sidecar::*;
|
||||
|
||||
@@ -375,7 +375,7 @@ async fn run_sidecar_turn(
|
||||
pub(crate) fn cancel_direct_claude_code_turn_at(
|
||||
root: &Path,
|
||||
client_turn_id: Option<&str>,
|
||||
) -> Result<Option<super::codex_app_server::DirectTurnCancelView>, String> {
|
||||
) -> Result<Option<super::codex_app_server::TurnCancelView>, String> {
|
||||
let key = claude_project_key(root);
|
||||
let turns = active_claude_code_turns()
|
||||
.lock()
|
||||
@@ -394,7 +394,7 @@ pub(crate) fn cancel_direct_claude_code_turn_at(
|
||||
active.alive.store(false, Ordering::Release);
|
||||
kill_claude_code_process_tree(active.pid);
|
||||
let client_turn_id = active.client_turn_id.clone();
|
||||
Ok(Some(super::codex_app_server::DirectTurnCancelView {
|
||||
Ok(Some(super::codex_app_server::TurnCancelView {
|
||||
outcome: super::codex_app_server::DIRECT_TURN_CANCEL_OUTCOME_INTERRUPTED.to_string(),
|
||||
message: "已向正在运行的 cc 回合发出终止".to_string(),
|
||||
client_turn_id,
|
||||
@@ -917,7 +917,7 @@ pub(crate) async fn direct_game_creator_claude_code_chat_at(
|
||||
system_prompt: String,
|
||||
user_prompt: String,
|
||||
client_turn_id: Option<&str>,
|
||||
observer: Option<&mut (dyn FnMut(DirectCodexTurnObservation) + Send)>,
|
||||
observer: Option<&mut (dyn FnMut(TurnObservation) + Send)>,
|
||||
) -> Result<String, String> {
|
||||
direct_turn_trace("claude-executor-enter");
|
||||
let (mcp_url, mcp_token) = start_external_mcp_loopback(root, llm.web_search_enabled).await?;
|
||||
@@ -984,25 +984,19 @@ pub(crate) async fn direct_game_creator_claude_code_chat_at(
|
||||
// 项目对话历史是这条对话的单一事实源:回复不落盘,UI 就看不到本轮结果。codex 路径由
|
||||
// app-server 的 collect-history 负责写 assistant 条目,cc 没有那一步——只补终态会让
|
||||
// 用户看到"回合结束但没有回复"。落盘失败按回合失败收口,不吞。
|
||||
let item_id = match client_turn_id {
|
||||
Some(client_turn_id) => format!("direct-codex:{client_turn_id}:assistant"),
|
||||
None => format!("direct-codex:{}:assistant", uuid::Uuid::new_v4()),
|
||||
let item_id = match persist_direct_claude_assistant_reply_at(root, client_turn_id, text)
|
||||
{
|
||||
Ok(item_id) => item_id,
|
||||
Err(error) => {
|
||||
eprintln!("[agc-cc-direct] persist assistant failed: {error}");
|
||||
direct_turn_trace("claude-parse-error");
|
||||
return Err(format!("写入本项目对话历史失败:{error}"));
|
||||
}
|
||||
};
|
||||
let item = serde_json::json!({
|
||||
"type": "message",
|
||||
"role": "assistant",
|
||||
"id": item_id.clone(),
|
||||
"content": [{ "type": "output_text", "text": text }],
|
||||
});
|
||||
if let Err(error) = append_direct_project_history_item_at(root, &item) {
|
||||
eprintln!("[agc-cc-direct] persist assistant failed: {error}");
|
||||
direct_turn_trace("claude-parse-error");
|
||||
return Err(format!("写入本项目对话历史失败:{error}"));
|
||||
}
|
||||
// 聊天区是按 `item.completed` 事件流投影的(codex 路径在 rawResponseItem/completed
|
||||
// 时下发同款事件),只落盘历史不会让本轮回复出现在界面上——重进项目才看得到。
|
||||
// 条目身份与落盘的历史条目保持同一个,重进项目按 id 去重。
|
||||
let at = crate::agent::direct_now_ms();
|
||||
let at = crate::agent::now_ms();
|
||||
crate::agent::append_thread_event(
|
||||
&crate::agent::thread_id_for_project(root),
|
||||
ThreadEvent::item_completed(
|
||||
@@ -1026,9 +1020,41 @@ pub(crate) async fn direct_game_creator_claude_code_chat_at(
|
||||
parsed
|
||||
}
|
||||
|
||||
fn persist_direct_claude_assistant_reply_at(
|
||||
root: &Path,
|
||||
client_turn_id: Option<&str>,
|
||||
text: &str,
|
||||
) -> Result<String, String> {
|
||||
let stable_item_id = match client_turn_id {
|
||||
Some(client_turn_id) => format!("direct-codex:{client_turn_id}:assistant"),
|
||||
None => format!("direct-codex:{}:assistant", uuid::Uuid::new_v4()),
|
||||
};
|
||||
let mut item = serde_json::json!({
|
||||
"type": "message",
|
||||
"role": "assistant",
|
||||
"id": stable_item_id.clone(),
|
||||
"content": [{ "type": "output_text", "text": text }],
|
||||
});
|
||||
match append_direct_project_history_item_at(root, &item) {
|
||||
Ok(()) => Ok(stable_item_id),
|
||||
Err(error) if error.starts_with("DirectProject 历史 item id 冲突:") => {
|
||||
// 同一 client turn 的验收反馈会再次进入 cc 执行器。首个回复保留稳定身份,
|
||||
// 后续不同回复必须另用 item id;否则历史层会把反馈收尾误判成失败。
|
||||
let feedback_item_id = format!("{stable_item_id}:{}", uuid::Uuid::new_v4());
|
||||
item["id"] = serde_json::Value::String(feedback_item_id.clone());
|
||||
append_direct_project_history_item_at(root, &item)
|
||||
.map(|()| feedback_item_id)
|
||||
.map_err(|retry_error| {
|
||||
format!("反馈回复历史追加失败:{retry_error}(首个条目冲突:{error})")
|
||||
})
|
||||
}
|
||||
Err(error) => Err(error),
|
||||
}
|
||||
}
|
||||
|
||||
fn parse_direct_stream_result(
|
||||
stdout: &[u8],
|
||||
mut observer: Option<&mut (dyn FnMut(DirectCodexTurnObservation) + Send)>,
|
||||
mut observer: Option<&mut (dyn FnMut(TurnObservation) + Send)>,
|
||||
) -> Result<String, String> {
|
||||
let text = std::str::from_utf8(stdout)
|
||||
.map_err(|_| "Claude Agent SDK stream-json 不是 UTF-8".to_string())?;
|
||||
@@ -1056,7 +1082,7 @@ fn parse_direct_stream_result(
|
||||
.collect::<String>();
|
||||
if !visible.is_empty() {
|
||||
if let Some(observer) = observer.as_deref_mut() {
|
||||
observer(DirectCodexTurnObservation::AgentMessageSegment(visible));
|
||||
observer(TurnObservation::AgentMessageSegment(visible));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1145,7 +1171,7 @@ mod tests {
|
||||
assert_eq!(text, "最终回复");
|
||||
assert!(matches!(
|
||||
observed.as_slice(),
|
||||
[DirectCodexTurnObservation::AgentMessageSegment(_)]
|
||||
[TurnObservation::AgentMessageSegment(_)]
|
||||
));
|
||||
}
|
||||
|
||||
@@ -1160,6 +1186,28 @@ mod tests {
|
||||
assert!(error.contains("Authentication failed"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn direct_claude_feedback_reply_does_not_fail_on_a_reused_client_turn_id() {
|
||||
let root = tempfile::tempdir().expect("temp dir");
|
||||
crate::init_local_game_project_at(root.path(), "cc-feedback", "cc feedback")
|
||||
.expect("init project");
|
||||
let turn_id = Some("client-turn-feedback-0001");
|
||||
|
||||
let first = persist_direct_claude_assistant_reply_at(root.path(), turn_id, "首轮回复")
|
||||
.expect("persist first reply");
|
||||
let identical = persist_direct_claude_assistant_reply_at(root.path(), turn_id, "首轮回复")
|
||||
.expect("an identical retry remains idempotent");
|
||||
let second = persist_direct_claude_assistant_reply_at(root.path(), turn_id, "反馈后回复")
|
||||
.expect("a host feedback retry must persist a second reply");
|
||||
|
||||
assert_eq!(first, identical);
|
||||
assert_ne!(first, second);
|
||||
assert!(second.starts_with("direct-codex:client-turn-feedback-0001:assistant:"));
|
||||
let items =
|
||||
crate::agent::read_direct_project_history_items_at(root.path()).expect("read history");
|
||||
assert_eq!(items.len(), 2);
|
||||
}
|
||||
|
||||
fn test_platform_session() -> crate::platform_session::PlatformSessionSnapshot {
|
||||
crate::platform_session::PlatformSessionSnapshot {
|
||||
user_id: "user-1".to_string(),
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
//! Native / third-party approval adapter. The host execution session owns policy
|
||||
//! and persistence; this module only binds the app-server protocol to its leases.
|
||||
|
||||
use super::super::{direct_delivery, direct_execution, direct_validation, DirectTurnError};
|
||||
use super::super::{direct_delivery, direct_execution, direct_validation, TurnError};
|
||||
use super::{shutdown_game_creator_codex_app_server_inner, CodexAppServerInner};
|
||||
use direct_execution::{EffectKind, ExecutionLease, ExecutionPhase, ExecutionSession};
|
||||
use serde_json::{json, Value};
|
||||
@@ -152,12 +152,10 @@ pub(super) const HOST_OUTCOME_REPAIR_REQUIRED_DETAIL: &str =
|
||||
|
||||
impl HostOutcomeText {
|
||||
/// 投影成这一轮的收尾结果:正常报告是文本,返修要求是控制流(走 `Err` 侧自己的变体)。
|
||||
pub(super) fn into_run_result(self) -> Result<String, super::DirectTurnRunFailure> {
|
||||
pub(super) fn into_run_result(self) -> Result<String, super::RunFailure> {
|
||||
match self {
|
||||
Self::Report(text) => Ok(text),
|
||||
Self::RepairRequired { detail } => {
|
||||
Err(super::DirectTurnRunFailure::RepairRequired { detail })
|
||||
}
|
||||
Self::RepairRequired { detail } => Err(super::RunFailure::RepairRequired { detail }),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -246,7 +244,7 @@ pub(super) struct ExecutionAdapter {
|
||||
outcome: watch::Sender<Option<HostOutcome>>,
|
||||
/// 宿主自己判定的"本轮以失败收口":`(分类, 原因)`。有值就代表本轮终态必须是失败,
|
||||
/// 原因与交付报告同一份文本。
|
||||
turn_failure: Mutex<Option<DirectTurnError>>,
|
||||
turn_failure: Mutex<Option<TurnError>>,
|
||||
/// 用户/宿主是否主动要求终止这一轮(界面的「终止」按钮)。用户主动终止不是失败。
|
||||
host_stop_requested: AtomicBool,
|
||||
/// 已留痕的拒绝原因(`method\u{1}reason`)。同一回合内同因只记一次,避免模型重试刷屏。
|
||||
@@ -921,8 +919,8 @@ impl ExecutionAdapter {
|
||||
///
|
||||
/// 只记第一份:第一份最接近现场(连接终止时带 exitStatus / stderr 摘要),后面更粗的收束理由
|
||||
/// 不得覆盖它。
|
||||
pub(super) async fn fail_turn(&self, failure: DirectTurnError) {
|
||||
let reason = failure.to_string();
|
||||
pub(super) async fn fail_turn(&self, failure: TurnError) {
|
||||
let reason = failure.diagnostic_detail();
|
||||
if self.is_closed() {
|
||||
// 宿主自己收尾:连接是我们先关的,紧随其后的 `TransportClosed` 只是收尾的副产物。
|
||||
// 只把原因留给报告,不改阶段——否则正常的宿主收尾会被改写成 `interrupted`
|
||||
@@ -944,7 +942,7 @@ impl ExecutionAdapter {
|
||||
}
|
||||
|
||||
/// 本轮以什么理由失败;有值就是宿主记下的 typed 事实。终态判定只读这一次。
|
||||
pub(super) fn turn_failure(&self) -> Option<DirectTurnError> {
|
||||
pub(super) fn turn_failure(&self) -> Option<TurnError> {
|
||||
self.turn_failure.lock().ok().and_then(|slot| slot.clone())
|
||||
}
|
||||
|
||||
@@ -1041,7 +1039,7 @@ impl ExecutionAdapter {
|
||||
|
||||
pub(super) fn lifecycle_status(&self, fallback: &str) -> String {
|
||||
// 只按收尾阶段归类。失败事实(`fail_turn` 记下的)不在这里翻案:终态由
|
||||
// `direct_turn_terminal` 拿事实判定——否则"模型已经判失败"的一轮会被这里的
|
||||
// `turn_terminal` 拿事实判定——否则"模型已经判失败"的一轮会被这里的
|
||||
// `Interrupted` 抹成一次没有原因的"已结束"。
|
||||
match self.session.snapshot().map(|state| state.phase) {
|
||||
Ok(ExecutionPhase::Completed) => "completed",
|
||||
@@ -1335,7 +1333,7 @@ pub(super) async fn wait_outcome(
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::super::DirectTurnDeadline;
|
||||
use super::super::{Deadline, TimedOut, TransportClosed};
|
||||
|
||||
use super::*;
|
||||
|
||||
@@ -1391,34 +1389,28 @@ mod tests {
|
||||
assert!(!adapter.host_stop_requested());
|
||||
|
||||
adapter
|
||||
.fail_turn(DirectTurnError::TransportClosed {
|
||||
.fail_turn(TurnError::TransportClosed(TransportClosed {
|
||||
diagnostic: "Codex app-server 已退出;exitStatus=signal: 9 (SIGKILL)".into(),
|
||||
})
|
||||
}))
|
||||
.await;
|
||||
|
||||
// 终态判定读这份事实,界面才有理由把它当失败讲,而不是"本轮已结束"。
|
||||
let failure = adapter.turn_failure().expect("host fact must be recorded");
|
||||
assert_eq!(
|
||||
failure.wire_kind(),
|
||||
Some(super::super::DirectTurnFailureKind::TransportFailed)
|
||||
);
|
||||
assert!(failure.to_string().contains("SIGKILL"));
|
||||
assert!(matches!(failure, TurnError::TransportClosed(_)));
|
||||
assert!(failure.diagnostic_detail().contains("SIGKILL"));
|
||||
// 报告与事件载荷同一份原因:用户看到的现象和交付状态对得上。
|
||||
assert!(adapter.report().contains("SIGKILL"));
|
||||
|
||||
// 只认第一份原因:后续更粗的收束理由不得覆盖真实诊断。
|
||||
adapter
|
||||
.fail_turn(DirectTurnError::TimedOut {
|
||||
deadline: DirectTurnDeadline::ResponseIdle,
|
||||
})
|
||||
.fail_turn(TurnError::TimedOut(TimedOut {
|
||||
deadline: Deadline::ResponseIdle,
|
||||
}))
|
||||
.await;
|
||||
let failure = adapter.turn_failure().expect("first reason is kept");
|
||||
assert_eq!(
|
||||
failure.wire_kind(),
|
||||
Some(super::super::DirectTurnFailureKind::TransportFailed)
|
||||
);
|
||||
assert!(failure.to_string().contains("SIGKILL"));
|
||||
assert!(!failure.to_string().contains("超时"));
|
||||
assert!(matches!(failure, TurnError::TransportClosed(_)));
|
||||
assert!(failure.diagnostic_detail().contains("SIGKILL"));
|
||||
assert!(!failure.diagnostic_detail().contains("超时"));
|
||||
}
|
||||
|
||||
/// 宿主自己关的连接不算失败:正常终态、用户主动停止、预算与交付收尾都会关掉连接,回合事件通道
|
||||
@@ -1430,9 +1422,9 @@ mod tests {
|
||||
adapter.closed.store(true, Ordering::Release);
|
||||
|
||||
adapter
|
||||
.fail_turn(DirectTurnError::TransportClosed {
|
||||
.fail_turn(TurnError::TransportClosed(TransportClosed {
|
||||
diagnostic: "模型本次执行结束,回收原生后台子树".into(),
|
||||
})
|
||||
}))
|
||||
.await;
|
||||
|
||||
assert!(adapter.turn_failure().is_none());
|
||||
@@ -1457,9 +1449,9 @@ mod tests {
|
||||
adapter.request_host_stop();
|
||||
|
||||
adapter
|
||||
.fail_turn(DirectTurnError::TransportClosed {
|
||||
.fail_turn(TurnError::TransportClosed(TransportClosed {
|
||||
diagnostic: "Codex app-server 已退出;exitStatus=signal: 9 (SIGKILL)".into(),
|
||||
})
|
||||
}))
|
||||
.await;
|
||||
|
||||
assert!(adapter.turn_failure().is_none());
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
+589
-434
File diff suppressed because it is too large
Load Diff
@@ -4,7 +4,7 @@ mod model;
|
||||
mod validation;
|
||||
mod wire;
|
||||
|
||||
pub(crate) use model::DirectCodexUserItem;
|
||||
pub(crate) use model::UserItem;
|
||||
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, freeze_direct_codex_user_item,
|
||||
|
||||
@@ -5,31 +5,31 @@ use ts_rs::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 {
|
||||
pub(crate) enum UserItem {
|
||||
#[serde(rename = "message")]
|
||||
Message(DirectCodexUserMessageItem),
|
||||
Message(UserMessageItem),
|
||||
}
|
||||
|
||||
#[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 {
|
||||
pub(crate) role: DirectCodexUserRole,
|
||||
pub(crate) content: Vec<DirectCodexUserContentPart>,
|
||||
pub(crate) struct UserMessageItem {
|
||||
pub(crate) role: UserRole,
|
||||
pub(crate) content: Vec<UserContentPart>,
|
||||
pub(crate) id: String,
|
||||
}
|
||||
|
||||
#[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 {
|
||||
pub(crate) enum UserRole {
|
||||
User,
|
||||
}
|
||||
|
||||
#[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 {
|
||||
pub(crate) enum UserContentPart {
|
||||
#[serde(rename = "input_text")]
|
||||
InputText { text: String },
|
||||
#[serde(rename = "agc_resource_reference")]
|
||||
@@ -44,16 +44,16 @@ pub(crate) enum DirectCodexUserContentPart {
|
||||
#[serde(rename = "agc_skill_reference")]
|
||||
AgcSkillReference { name: String },
|
||||
#[serde(rename = "agc_runtime_region_reference")]
|
||||
AgcRuntimeRegionReference(DirectCodexUserRuntimeRegionPart),
|
||||
AgcRuntimeRegionReference(UserRuntimeRegionPart),
|
||||
/// Uploaded project attachment kept inline in canonical content.
|
||||
#[serde(rename = "agc_attachment_reference")]
|
||||
AgcAttachmentReference(DirectCodexUserAttachmentReferencePart),
|
||||
AgcAttachmentReference(UserAttachmentReferencePart),
|
||||
}
|
||||
|
||||
#[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 {
|
||||
pub(crate) struct UserAttachmentReferencePart {
|
||||
pub(crate) name: String,
|
||||
pub(crate) media_type: String,
|
||||
#[ts(type = "number")]
|
||||
@@ -65,7 +65,7 @@ pub(crate) struct DirectCodexUserAttachmentReferencePart {
|
||||
#[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 {
|
||||
pub(crate) struct UserRuntimeRegionPart {
|
||||
pub(crate) label: String,
|
||||
#[serde(default)]
|
||||
pub(crate) run_id: Option<String>,
|
||||
@@ -92,9 +92,9 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn resource_reference_serializes_with_only_camel_case_resource_id() {
|
||||
let item = DirectCodexUserItem::Message(DirectCodexUserMessageItem {
|
||||
role: DirectCodexUserRole::User,
|
||||
content: vec![DirectCodexUserContentPart::AgcResourceReference {
|
||||
let item = UserItem::Message(UserMessageItem {
|
||||
role: UserRole::User,
|
||||
content: vec![UserContentPart::AgcResourceReference {
|
||||
resource_id: "asset-hero".to_string(),
|
||||
resolved_text: None,
|
||||
}],
|
||||
@@ -116,7 +116,7 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn resource_reference_rejects_extra_identity_fields() {
|
||||
let error = serde_json::from_value::<DirectCodexUserItem>(json!({
|
||||
let error = serde_json::from_value::<UserItem>(json!({
|
||||
"type": "message",
|
||||
"role": "user",
|
||||
"content": [{
|
||||
@@ -133,7 +133,7 @@ mod tests {
|
||||
/// 旧历史里的资源引用没有 `resolvedText`:缺省必须合法,照样解析成 `None` 的引用。
|
||||
#[test]
|
||||
fn resource_reference_without_resolved_text_still_parses() {
|
||||
let item: DirectCodexUserItem = serde_json::from_value(json!({
|
||||
let item: UserItem = serde_json::from_value(json!({
|
||||
"type": "message",
|
||||
"role": "user",
|
||||
"content": [{
|
||||
@@ -143,9 +143,9 @@ mod tests {
|
||||
"id": "turn-1"
|
||||
}))
|
||||
.expect("legacy resource reference must keep parsing");
|
||||
let DirectCodexUserItem::Message(message) = item;
|
||||
let UserItem::Message(message) = item;
|
||||
match &message.content[0] {
|
||||
DirectCodexUserContentPart::AgcResourceReference {
|
||||
UserContentPart::AgcResourceReference {
|
||||
resource_id,
|
||||
resolved_text,
|
||||
} => {
|
||||
@@ -158,7 +158,7 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn skill_reference_serializes_with_only_the_stable_name() {
|
||||
let item: DirectCodexUserItem = serde_json::from_value(json!({
|
||||
let item: UserItem = serde_json::from_value(json!({
|
||||
"type": "message",
|
||||
"role": "user",
|
||||
"content": [{"type": "agc_skill_reference", "name": "agc-web-game-development"}],
|
||||
@@ -172,7 +172,7 @@ mod tests {
|
||||
|
||||
// canonical part 不接受正文、路径或凭据类附加字段:它们只可能来自宿主私密状态。
|
||||
for forbidden in ["path", "body", "content", "token", "apiKey"] {
|
||||
let error = serde_json::from_value::<DirectCodexUserItem>(json!({
|
||||
let error = serde_json::from_value::<UserItem>(json!({
|
||||
"type": "message",
|
||||
"role": "user",
|
||||
"content": [{
|
||||
@@ -192,7 +192,7 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn unknown_content_part_fails_closed() {
|
||||
serde_json::from_value::<DirectCodexUserItem>(json!({
|
||||
serde_json::from_value::<UserItem>(json!({
|
||||
"type": "message",
|
||||
"role": "user",
|
||||
"content": [{"type": "future_part", "value": "x"}],
|
||||
|
||||
+16
-19
@@ -1,7 +1,4 @@
|
||||
use super::model::{
|
||||
DirectCodexUserContentPart, DirectCodexUserItem, DirectCodexUserRole,
|
||||
DirectCodexUserRuntimeRegionPart,
|
||||
};
|
||||
use super::model::{UserContentPart, UserItem, UserRole, UserRuntimeRegionPart};
|
||||
use crate::agent::{
|
||||
read_manifest_for_project, sanitize_attachment_local_path, GameCreationAppManifest,
|
||||
MAX_DIRECT_CODEX_ATTACHMENTS, MAX_DIRECT_CODEX_ATTACHMENT_MEDIA_TYPE_CHARS,
|
||||
@@ -16,10 +13,10 @@ pub(crate) const MAX_DIRECT_CODEX_SKILL_REFERENCES: usize = 32;
|
||||
|
||||
pub(crate) fn validate_direct_codex_user_item(
|
||||
root: &Path,
|
||||
item: &DirectCodexUserItem,
|
||||
item: &UserItem,
|
||||
) -> Result<GameCreationAppManifest, String> {
|
||||
let DirectCodexUserItem::Message(message) = item;
|
||||
if !matches!(message.role, DirectCodexUserRole::User) {
|
||||
let UserItem::Message(message) = item;
|
||||
if !matches!(message.role, UserRole::User) {
|
||||
return Err("DirectProject 只接受 user message item".to_string());
|
||||
}
|
||||
if message.id.trim().is_empty() {
|
||||
@@ -36,12 +33,12 @@ pub(crate) fn validate_direct_codex_user_item(
|
||||
let mut skill_count = 0usize;
|
||||
for part in &message.content {
|
||||
match part {
|
||||
DirectCodexUserContentPart::InputText { .. } => {}
|
||||
DirectCodexUserContentPart::AgcResourceReference { resource_id, .. } => {
|
||||
UserContentPart::InputText { .. } => {}
|
||||
UserContentPart::AgcResourceReference { resource_id, .. } => {
|
||||
reference_count = reference_count.saturating_add(1);
|
||||
validate_resource_id_and_manifest(&manifest, resource_id)?;
|
||||
}
|
||||
DirectCodexUserContentPart::AgcSkillReference { name } => {
|
||||
UserContentPart::AgcSkillReference { name } => {
|
||||
skill_count = skill_count.saturating_add(1);
|
||||
if skill_count > MAX_DIRECT_CODEX_SKILL_REFERENCES {
|
||||
return Err(format!(
|
||||
@@ -61,11 +58,11 @@ pub(crate) fn validate_direct_codex_user_item(
|
||||
return Err("引用的 Skill 名称无效,请移除后重新选择".to_string());
|
||||
}
|
||||
}
|
||||
DirectCodexUserContentPart::AgcRuntimeRegionReference(reference) => {
|
||||
UserContentPart::AgcRuntimeRegionReference(reference) => {
|
||||
reference_count = reference_count.saturating_add(1);
|
||||
validate_runtime_region_reference(&manifest, reference)?;
|
||||
}
|
||||
DirectCodexUserContentPart::AgcAttachmentReference(reference) => {
|
||||
UserContentPart::AgcAttachmentReference(reference) => {
|
||||
attachment_count = attachment_count.saturating_add(1);
|
||||
if attachment_count > MAX_DIRECT_CODEX_ATTACHMENTS {
|
||||
return Err(format!(
|
||||
@@ -112,9 +109,9 @@ pub(crate) fn validate_direct_codex_user_item(
|
||||
}
|
||||
|
||||
/// 整条 content 是否还有有效输入:任何一段非空白文本、或任何一个非文本 part 都算。
|
||||
pub(crate) fn content_has_meaningful_input(content: &[DirectCodexUserContentPart]) -> bool {
|
||||
pub(crate) fn content_has_meaningful_input(content: &[UserContentPart]) -> bool {
|
||||
content.iter().any(|part| match part {
|
||||
DirectCodexUserContentPart::InputText { text } => !text.trim().is_empty(),
|
||||
UserContentPart::InputText { text } => !text.trim().is_empty(),
|
||||
_ => true,
|
||||
})
|
||||
}
|
||||
@@ -142,7 +139,7 @@ pub(crate) fn validate_resource_id_and_manifest(
|
||||
|
||||
fn validate_runtime_region_reference(
|
||||
manifest: &GameCreationAppManifest,
|
||||
reference: &DirectCodexUserRuntimeRegionPart,
|
||||
reference: &UserRuntimeRegionPart,
|
||||
) -> Result<(), String> {
|
||||
if reference.label.trim().is_empty() {
|
||||
return Err("运行画面区域缺少名称".to_string());
|
||||
@@ -164,11 +161,11 @@ mod tests {
|
||||
content_has_meaningful_input, validate_direct_codex_user_item,
|
||||
MAX_DIRECT_CODEX_SKILL_REFERENCES,
|
||||
};
|
||||
use crate::agent::direct_codex_user_item::model::DirectCodexUserContentPart;
|
||||
use crate::agent::direct_codex_user_item::model::UserContentPart;
|
||||
use serde_json::json;
|
||||
|
||||
fn input_text(text: &str) -> DirectCodexUserContentPart {
|
||||
DirectCodexUserContentPart::InputText {
|
||||
fn input_text(text: &str) -> UserContentPart {
|
||||
UserContentPart::InputText {
|
||||
text: text.to_string(),
|
||||
}
|
||||
}
|
||||
@@ -200,7 +197,7 @@ mod tests {
|
||||
#[test]
|
||||
fn non_text_parts_always_count_as_input() {
|
||||
assert!(content_has_meaningful_input(&[
|
||||
DirectCodexUserContentPart::AgcResourceReference {
|
||||
UserContentPart::AgcResourceReference {
|
||||
resource_id: "asset-hero".to_string(),
|
||||
resolved_text: None,
|
||||
},
|
||||
|
||||
@@ -1,6 +1,5 @@
|
||||
use super::model::{
|
||||
DirectCodexUserAttachmentReferencePart, DirectCodexUserContentPart, DirectCodexUserItem,
|
||||
DirectCodexUserMessageItem, DirectCodexUserRuntimeRegionPart,
|
||||
UserAttachmentReferencePart, UserContentPart, UserItem, UserMessageItem, UserRuntimeRegionPart,
|
||||
};
|
||||
use super::validation::validate_direct_codex_user_item;
|
||||
use crate::agent::{
|
||||
@@ -15,21 +14,21 @@ use serde_json::Value;
|
||||
use std::path::Path;
|
||||
|
||||
/// 入队检查产出的**冻结条目**:校验 → 把每个引用 part 的解析文本写进它自己的
|
||||
/// [`DirectCodexUserContentPart::AgcResourceReference::resolved_text`] → 判空。
|
||||
/// [`UserContentPart::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> {
|
||||
item: &UserItem,
|
||||
) -> Result<UserItem, String> {
|
||||
let manifest = validate_direct_codex_user_item(root, item)?;
|
||||
let DirectCodexUserItem::Message(message) = item;
|
||||
let UserItem::Message(message) = item;
|
||||
let mut content = Vec::with_capacity(message.content.len());
|
||||
for part in &message.content {
|
||||
match part {
|
||||
DirectCodexUserContentPart::AgcResourceReference {
|
||||
UserContentPart::AgcResourceReference {
|
||||
resource_id,
|
||||
resolved_text,
|
||||
} => {
|
||||
@@ -42,7 +41,7 @@ pub(crate) fn freeze_direct_codex_user_item(
|
||||
resource_id,
|
||||
)?),
|
||||
};
|
||||
content.push(DirectCodexUserContentPart::AgcResourceReference {
|
||||
content.push(UserContentPart::AgcResourceReference {
|
||||
resource_id: resource_id.clone(),
|
||||
resolved_text,
|
||||
});
|
||||
@@ -50,7 +49,7 @@ pub(crate) fn freeze_direct_codex_user_item(
|
||||
other => content.push(other.clone()),
|
||||
}
|
||||
}
|
||||
let frozen = DirectCodexUserItem::Message(DirectCodexUserMessageItem {
|
||||
let frozen = UserItem::Message(UserMessageItem {
|
||||
role: message.role.clone(),
|
||||
content,
|
||||
id: message.id.clone(),
|
||||
@@ -75,7 +74,7 @@ pub(crate) fn direct_codex_user_item_to_response_item(
|
||||
}
|
||||
return Err("DirectProject 历史 item 缺少 type,无法投影为 Codex item".to_string());
|
||||
}
|
||||
let canonical: DirectCodexUserItem = serde_json::from_value(item.clone())
|
||||
let canonical: UserItem = serde_json::from_value(item.clone())
|
||||
.map_err(|error| format!("DirectProject user item 无法转换为 Codex item:{error}"))?;
|
||||
let content = direct_codex_user_item_to_response_content(root, &canonical)?;
|
||||
let mut projected = serde_json::json!({
|
||||
@@ -91,7 +90,7 @@ pub(crate) fn direct_codex_user_item_to_response_item(
|
||||
|
||||
fn direct_codex_user_item_to_response_content(
|
||||
root: &Path,
|
||||
item: &DirectCodexUserItem,
|
||||
item: &UserItem,
|
||||
) -> Result<Vec<Value>, String> {
|
||||
let Value::Array(input) = direct_codex_user_item_to_wire_input(root, item)? else {
|
||||
return Err("DirectProject user item wire content 不是数组".to_string());
|
||||
@@ -125,7 +124,7 @@ fn resource_reference_summary(
|
||||
))
|
||||
}
|
||||
|
||||
fn runtime_region_summary(reference: &DirectCodexUserRuntimeRegionPart) -> String {
|
||||
fn runtime_region_summary(reference: &UserRuntimeRegionPart) -> String {
|
||||
let resources = reference
|
||||
.resource_ids
|
||||
.iter()
|
||||
@@ -153,7 +152,7 @@ fn runtime_region_summary(reference: &DirectCodexUserRuntimeRegionPart) -> Strin
|
||||
///
|
||||
/// turn 输入与 history/prompt 投影共用这一份清洗:文件名取 basename 并去控制字符、
|
||||
/// media type 与项目路径同样过白名单,避免两条路径对同一个引用给出不同摘要。
|
||||
fn attachment_reference_summary(reference: &DirectCodexUserAttachmentReferencePart) -> String {
|
||||
fn attachment_reference_summary(reference: &UserAttachmentReferencePart) -> String {
|
||||
let name = sanitize_attachment_name(&reference.name);
|
||||
let media_type = sanitize_attachment_media_type(&reference.media_type);
|
||||
let mut summary = format!(
|
||||
@@ -172,16 +171,16 @@ fn attachment_reference_summary(reference: &DirectCodexUserAttachmentReferencePa
|
||||
/// AGC 私有 part 只在这里投影为安全摘要,canonical item 本身不被修改。
|
||||
pub(crate) fn direct_codex_user_item_to_wire_input(
|
||||
root: &Path,
|
||||
item: &DirectCodexUserItem,
|
||||
item: &UserItem,
|
||||
) -> Result<Value, String> {
|
||||
// validate 已经读过清单并返回它,不要再读一次(seed task 变更也会被重复触发)。
|
||||
let manifest = validate_direct_codex_user_item(root, item)?;
|
||||
let DirectCodexUserItem::Message(message) = item;
|
||||
let UserItem::Message(message) = item;
|
||||
let mut input = Vec::with_capacity(message.content.len());
|
||||
for part in &message.content {
|
||||
let text = match part {
|
||||
DirectCodexUserContentPart::InputText { text } => text.clone(),
|
||||
DirectCodexUserContentPart::AgcResourceReference {
|
||||
UserContentPart::InputText { text } => text.clone(),
|
||||
UserContentPart::AgcResourceReference {
|
||||
resource_id,
|
||||
resolved_text,
|
||||
} => match resolved_text {
|
||||
@@ -190,13 +189,13 @@ pub(crate) fn direct_codex_user_item_to_wire_input(
|
||||
// 旧历史没有这份文本,退回按当前 manifest 现算(只有摘要,不产生写副作用)。
|
||||
None => resource_reference_summary(&manifest, resource_id)?,
|
||||
},
|
||||
DirectCodexUserContentPart::AgcSkillReference { name } => {
|
||||
UserContentPart::AgcSkillReference { name } => {
|
||||
format!("${}", name.trim())
|
||||
}
|
||||
DirectCodexUserContentPart::AgcRuntimeRegionReference(reference) => {
|
||||
UserContentPart::AgcRuntimeRegionReference(reference) => {
|
||||
runtime_region_summary(reference)
|
||||
}
|
||||
DirectCodexUserContentPart::AgcAttachmentReference(reference) => {
|
||||
UserContentPart::AgcAttachmentReference(reference) => {
|
||||
attachment_reference_summary(reference)
|
||||
}
|
||||
};
|
||||
@@ -207,18 +206,18 @@ pub(crate) fn direct_codex_user_item_to_wire_input(
|
||||
|
||||
pub(crate) fn direct_codex_user_item_to_codex_turn_input(
|
||||
root: &Path,
|
||||
item: &DirectCodexUserItem,
|
||||
item: &UserItem,
|
||||
skill_roots: &[std::path::PathBuf],
|
||||
) -> Result<Value, String> {
|
||||
let manifest = validate_direct_codex_user_item(root, item)?;
|
||||
let DirectCodexUserItem::Message(message) = item;
|
||||
let UserItem::Message(message) = item;
|
||||
let mut input = Vec::with_capacity(message.content.len());
|
||||
for part in &message.content {
|
||||
match part {
|
||||
DirectCodexUserContentPart::InputText { text } => {
|
||||
UserContentPart::InputText { text } => {
|
||||
input.push(serde_json::json!({ "type": "text", "text": text }));
|
||||
}
|
||||
DirectCodexUserContentPart::AgcResourceReference {
|
||||
UserContentPart::AgcResourceReference {
|
||||
resource_id,
|
||||
resolved_text,
|
||||
} => {
|
||||
@@ -230,7 +229,7 @@ pub(crate) fn direct_codex_user_item_to_codex_turn_input(
|
||||
},
|
||||
}));
|
||||
}
|
||||
DirectCodexUserContentPart::AgcSkillReference { name } => {
|
||||
UserContentPart::AgcSkillReference { name } => {
|
||||
let name = name.trim();
|
||||
let path = skill_roots
|
||||
.iter()
|
||||
@@ -243,13 +242,13 @@ pub(crate) fn direct_codex_user_item_to_codex_turn_input(
|
||||
"path": path,
|
||||
}));
|
||||
}
|
||||
DirectCodexUserContentPart::AgcRuntimeRegionReference(reference) => {
|
||||
UserContentPart::AgcRuntimeRegionReference(reference) => {
|
||||
input.push(serde_json::json!({
|
||||
"type": "text",
|
||||
"text": runtime_region_summary(reference),
|
||||
}));
|
||||
}
|
||||
DirectCodexUserContentPart::AgcAttachmentReference(reference) => {
|
||||
UserContentPart::AgcAttachmentReference(reference) => {
|
||||
input.push(serde_json::json!({
|
||||
"type": "text",
|
||||
"text": attachment_reference_summary(reference),
|
||||
@@ -266,24 +265,24 @@ pub(crate) fn direct_codex_user_item_to_codex_turn_input(
|
||||
/// [`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;
|
||||
pub(crate) fn direct_codex_user_item_to_prompt(item: &UserItem) -> String {
|
||||
let UserItem::Message(message) = item;
|
||||
let mut prompt = String::new();
|
||||
for part in &message.content {
|
||||
match part {
|
||||
DirectCodexUserContentPart::InputText { text } => prompt.push_str(text),
|
||||
DirectCodexUserContentPart::AgcResourceReference { resolved_text, .. } => {
|
||||
UserContentPart::InputText { text } => prompt.push_str(text),
|
||||
UserContentPart::AgcResourceReference { resolved_text, .. } => {
|
||||
if let Some(text) = resolved_text {
|
||||
prompt.push_str(text);
|
||||
}
|
||||
}
|
||||
DirectCodexUserContentPart::AgcSkillReference { name } => {
|
||||
UserContentPart::AgcSkillReference { name } => {
|
||||
prompt.push_str(&format!("${}", name.trim()));
|
||||
}
|
||||
DirectCodexUserContentPart::AgcRuntimeRegionReference(reference) => {
|
||||
UserContentPart::AgcRuntimeRegionReference(reference) => {
|
||||
prompt.push_str(&runtime_region_summary(reference));
|
||||
}
|
||||
DirectCodexUserContentPart::AgcAttachmentReference(reference) => {
|
||||
UserContentPart::AgcAttachmentReference(reference) => {
|
||||
prompt.push_str(&attachment_reference_summary(reference));
|
||||
}
|
||||
}
|
||||
@@ -351,7 +350,7 @@ mod tests {
|
||||
direct_codex_user_item_to_response_item, direct_codex_user_item_to_wire_input,
|
||||
freeze_direct_codex_user_item, validate_direct_codex_user_item,
|
||||
};
|
||||
use crate::agent::direct_codex_user_item::model::DirectCodexUserItem;
|
||||
use crate::agent::direct_codex_user_item::model::UserItem;
|
||||
use crate::ui_editor::persistence::UI_DESIGN_DOC_MEDIA_TYPE;
|
||||
use serde_json::json;
|
||||
use shared_contracts::game_creation_app::{
|
||||
@@ -437,7 +436,7 @@ mod tests {
|
||||
#[test]
|
||||
fn ui_design_doc_reference_appends_generated_code_context() {
|
||||
let (project, asset_id) = ui_design_doc_fixture(true);
|
||||
let item: super::DirectCodexUserItem =
|
||||
let item: super::UserItem =
|
||||
serde_json::from_value(user_item_with_resource_reference(&asset_id))
|
||||
.expect("canonical user item");
|
||||
let frozen = freeze_direct_codex_user_item(project.path(), &item).expect("freeze");
|
||||
@@ -453,9 +452,9 @@ mod tests {
|
||||
"{prompt}"
|
||||
);
|
||||
// 解析文本随条目持久化:历史回放读的就是这一份,不再重算、不再写盘。
|
||||
let super::DirectCodexUserItem::Message(message) = &frozen;
|
||||
let super::UserItem::Message(message) = &frozen;
|
||||
match &message.content[0] {
|
||||
super::DirectCodexUserContentPart::AgcResourceReference { resolved_text, .. } => {
|
||||
super::UserContentPart::AgcResourceReference { resolved_text, .. } => {
|
||||
assert_eq!(resolved_text.as_deref(), Some(prompt.as_str()));
|
||||
}
|
||||
other => panic!("expected a frozen resource reference, got {other:?}"),
|
||||
@@ -488,7 +487,7 @@ mod tests {
|
||||
#[test]
|
||||
fn ui_design_generation_failure_keeps_reference_and_reports_error() {
|
||||
let (project, asset_id) = ui_design_doc_fixture(false);
|
||||
let item: super::DirectCodexUserItem =
|
||||
let item: super::UserItem =
|
||||
serde_json::from_value(user_item_with_resource_reference(&asset_id))
|
||||
.expect("canonical user item");
|
||||
let frozen = freeze_direct_codex_user_item(project.path(), &item).expect("freeze");
|
||||
@@ -506,7 +505,7 @@ mod tests {
|
||||
GameCreationAppAssetKind::Character,
|
||||
"image/png",
|
||||
);
|
||||
let item: super::DirectCodexUserItem = serde_json::from_value(json!({
|
||||
let item: super::UserItem = serde_json::from_value(json!({
|
||||
"type": "message",
|
||||
"role": "user",
|
||||
"id": "turn-multi-1:user",
|
||||
@@ -547,7 +546,7 @@ mod tests {
|
||||
GameCreationAppAssetKind::Character,
|
||||
"image/png",
|
||||
);
|
||||
let item: super::DirectCodexUserItem =
|
||||
let item: super::UserItem =
|
||||
serde_json::from_value(user_item_with_resource_reference(&asset_id))
|
||||
.expect("canonical user item");
|
||||
let frozen = freeze_direct_codex_user_item(project.path(), &item).expect("freeze");
|
||||
@@ -774,8 +773,7 @@ mod tests {
|
||||
{"type": "agc_skill_reference", "name": "missing-skill"}
|
||||
]
|
||||
});
|
||||
let user_item: DirectCodexUserItem =
|
||||
serde_json::from_value(item).expect("parse canonical item");
|
||||
let user_item: UserItem = serde_json::from_value(item).expect("parse canonical item");
|
||||
// 已启用目录里没有这个 Skill:转换必须在启动回合前失败关闭,
|
||||
// 不能把不可用的引用降级成正文放行。
|
||||
let error = direct_codex_user_item_to_codex_turn_input(
|
||||
@@ -804,8 +802,7 @@ mod tests {
|
||||
{"type": "input_text", "text": "然后创建菜单"}
|
||||
]
|
||||
});
|
||||
let user_item: DirectCodexUserItem =
|
||||
serde_json::from_value(item).expect("parse canonical item");
|
||||
let user_item: UserItem = serde_json::from_value(item).expect("parse canonical item");
|
||||
let input =
|
||||
direct_codex_user_item_to_codex_turn_input(root.path(), &user_item, &[skill_root])
|
||||
.expect("available skill should convert");
|
||||
@@ -884,7 +881,7 @@ mod tests {
|
||||
let root = tempfile::tempdir().expect("temp project");
|
||||
crate::init_local_game_project_at(root.path(), "wire-test", "wire 投影测试")
|
||||
.expect("init project");
|
||||
let item: DirectCodexUserItem = serde_json::from_value(json!({
|
||||
let item: UserItem = serde_json::from_value(json!({
|
||||
"type": "message",
|
||||
"role": "user",
|
||||
"id": "turn-1:user",
|
||||
@@ -908,7 +905,7 @@ mod tests {
|
||||
let root = tempfile::tempdir().expect("temp project");
|
||||
crate::init_local_game_project_at(root.path(), "wire-test", "wire 投影测试")
|
||||
.expect("init project");
|
||||
let item: DirectCodexUserItem = serde_json::from_value(json!({
|
||||
let item: UserItem = serde_json::from_value(json!({
|
||||
"type": "message",
|
||||
"role": "user",
|
||||
"id": "turn-1:user",
|
||||
|
||||
@@ -698,12 +698,12 @@ pub(super) async fn finish_sealing(
|
||||
|
||||
/// 回合末的宿主复核:返回要交付的答复,或者一个"还没完,按这份证据继续修"的要求。
|
||||
///
|
||||
/// 返修要求是**控制流**([`DirectTurnError::ReviewRequired`]),不是失败:调用方据此把要求写回
|
||||
/// 返修要求是**控制流**([`TurnError::ReviewRequired`]),不是失败:调用方据此把要求写回
|
||||
/// prompt 再跑一轮,界面不该看到失败文案。其余错误都是真的回合失败,按 typed 错误交给上层。
|
||||
pub(super) async fn review_reply(
|
||||
root: &Path,
|
||||
session: &Arc<ExecutionSession>,
|
||||
) -> Result<Option<String>, DirectTurnError> {
|
||||
) -> Result<Option<String>, TurnError> {
|
||||
if let Some(report) = terminal_report(session) {
|
||||
return Ok(Some(report));
|
||||
}
|
||||
@@ -754,7 +754,7 @@ pub(super) async fn review_reply(
|
||||
.map_err(|_| "delivery-review-worker-exited")??;
|
||||
return Ok(Some(report));
|
||||
}
|
||||
Err(DirectTurnError::ReviewRequired {
|
||||
Err(TurnError::ReviewRequired {
|
||||
detail: format!("delivery-review-required: {detail}"),
|
||||
})
|
||||
}
|
||||
@@ -919,7 +919,7 @@ mod tests {
|
||||
for _ in 0..2 {
|
||||
assert!(matches!(
|
||||
review_reply(new_game.path(), &required).await.unwrap_err(),
|
||||
DirectTurnError::ReviewRequired { .. }
|
||||
TurnError::ReviewRequired { .. }
|
||||
));
|
||||
}
|
||||
assert!(review_reply(new_game.path(), &required)
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -7,9 +7,9 @@ use super::*;
|
||||
|
||||
pub(crate) fn normalize_direct_client_turn_id(
|
||||
client_turn_id: Option<&str>,
|
||||
) -> Result<String, DirectTurnError> {
|
||||
) -> Result<String, EnqueueError> {
|
||||
let Some(client_turn_id) = client_turn_id else {
|
||||
return Err(DirectTurnError::ClientTurnIdMissing);
|
||||
return Err(EnqueueError::ClientTurnIdMissing);
|
||||
};
|
||||
let client_turn_id = client_turn_id.trim();
|
||||
let valid_length = (MIN_DIRECT_CLIENT_TURN_ID_CHARS..=MAX_DIRECT_CLIENT_TURN_ID_CHARS)
|
||||
@@ -20,10 +20,10 @@ pub(crate) fn normalize_direct_client_turn_id(
|
||||
.is_some_and(|byte| byte.is_ascii_alphanumeric());
|
||||
let valid_rest = bytes.all(|byte| byte.is_ascii_alphanumeric() || byte == b'-');
|
||||
if !valid_length || !valid_first || !valid_rest {
|
||||
return Err(DirectTurnError::ClientTurnIdMalformed {
|
||||
return Err(EnqueueError::ClientTurnIdMalformed(ClientTurnIdMalformed {
|
||||
min_chars: MIN_DIRECT_CLIENT_TURN_ID_CHARS,
|
||||
max_chars: MAX_DIRECT_CLIENT_TURN_ID_CHARS,
|
||||
});
|
||||
}));
|
||||
}
|
||||
Ok(client_turn_id.to_string())
|
||||
}
|
||||
@@ -37,7 +37,7 @@ pub(crate) fn normalize_direct_client_turn_id(
|
||||
/// 于是"这一轮跑成什么"仍然只有订阅事件一个来源:命令返回 `Ok` 只说明**入队成立**。真正的回合边界
|
||||
/// (`turn.started` / `turn.completed`)由 Thread Manager 在**放行**时写出(见 `thread_manager::dispatch`),
|
||||
/// 入队失败不写用户条目、不产生任何事件。可留痕的调用级失败(宿主 / 环境事实)仍在边界补一份运行
|
||||
/// 错误诊断,返回串不带诊断引用。
|
||||
/// 错误诊断;命令返回的就是 typed 变体本身。
|
||||
///
|
||||
/// 它是 DirectProject 唯一的命令入口:终端入口 `--direct-codex-chat`(它要保持 await 才能把回复打到
|
||||
/// 终端上)已经退役,不要再为"手工跑一轮"新增第二条直接起回合的路径。
|
||||
@@ -47,10 +47,10 @@ pub(crate) fn normalize_direct_client_turn_id(
|
||||
#[tauri::command]
|
||||
pub(crate) async fn enqueue_direct_codex_turn(
|
||||
project_path: String,
|
||||
user_item: DirectCodexUserItem,
|
||||
user_item: UserItem,
|
||||
creation_type: Option<String>,
|
||||
client_turn_id: Option<String>,
|
||||
) -> Result<(), DirectTurnEnqueueFailure> {
|
||||
) -> Result<(), EnqueueError> {
|
||||
let root = Path::new(project_path.trim());
|
||||
let boundary_turn_id = client_turn_id.clone();
|
||||
enqueue_direct_codex_turn_typed(root, user_item, creation_type, client_turn_id)
|
||||
@@ -67,23 +67,23 @@ pub(crate) async fn enqueue_direct_codex_turn(
|
||||
/// 回合事件与整轮都在放行那一侧。
|
||||
async fn enqueue_direct_codex_turn_typed(
|
||||
root: &Path,
|
||||
user_item: DirectCodexUserItem,
|
||||
user_item: UserItem,
|
||||
creation_type: Option<String>,
|
||||
client_turn_id: Option<String>,
|
||||
) -> Result<(), DirectTurnError> {
|
||||
) -> Result<(), EnqueueError> {
|
||||
let turn_id = normalize_direct_client_turn_id(client_turn_id.as_deref())?;
|
||||
recover_direct_taonier_regeneration_workflow_at(root).map_err(|error| {
|
||||
DirectTurnError::HostStateUnavailable {
|
||||
EnqueueError::HostStateUnavailable(HostStateUnavailable {
|
||||
detail: redact_agent_runtime_error(
|
||||
root,
|
||||
&format!("恢复上一轮陶泥儿整包事务失败:{error}"),
|
||||
500,
|
||||
),
|
||||
}
|
||||
})
|
||||
})?;
|
||||
// 冻结是这条消息唯一的算片段 / 写盘时机:引用解析文本就此写进条目自身,放行只重投影。
|
||||
let user_item = freeze_direct_codex_user_item(root, &user_item)
|
||||
.map_err(|detail| DirectTurnError::InputRejected { detail })?;
|
||||
.map_err(|detail| EnqueueError::InputRejected(InputRejected { detail }))?;
|
||||
let user_prompt = direct_codex_user_item_to_prompt(&user_item);
|
||||
check_direct_turn_preconditions(root, &user_prompt, creation_type.as_deref())?;
|
||||
let thread_id = thread_id_for_project(root);
|
||||
@@ -91,25 +91,27 @@ async fn enqueue_direct_codex_turn_typed(
|
||||
// 准备——那是分钟级的活(`npm ci` + Vite 构建),还会在磁盘上留下产物。权威判据仍然是入队那一刻
|
||||
// 临界区里的容量检查(下面 `enqueue_pending_turn`):这里只是快速失败,中间被别人的消息
|
||||
// 挤满时那一条照样拦得住。
|
||||
queue_has_room(pending_turn_count(&thread_id)).map_err(|_| DirectTurnError::QueueFull {
|
||||
limit: MAX_PENDING_TURNS,
|
||||
queue_has_room(pending_turn_count(&thread_id)).map_err(|_| {
|
||||
EnqueueError::QueueFull(QueueFull {
|
||||
limit: MAX_PENDING_TURNS,
|
||||
})
|
||||
})?;
|
||||
// 创建类型来自结构化用户入口;实际工程和可信脚手架由宿主复核。
|
||||
crate::environment_check::prepare_new_web_project_at(root, creation_type.as_deref())
|
||||
.await
|
||||
.map_err(|error| {
|
||||
let detail = redact_agent_runtime_error(root, &error, 1800);
|
||||
DirectTurnError::EnvironmentNotReady { detail }
|
||||
EnqueueError::EnvironmentNotReady(EnvironmentNotReady { detail })
|
||||
})?;
|
||||
// 入队:到这里这一条已经过了全部检查,剩下的就是排队等放行。条目只带走它自己的事实
|
||||
// (用户条目、创建类型、入队时刻),canonical 形状与 prompt 放行时从它重投影——放行没有失败出口。
|
||||
let pending = PendingTurn::new(turn_id, user_item, creation_type, direct_now_ms());
|
||||
let pending = PendingTurn::new(turn_id, user_item, creation_type, now_ms());
|
||||
match enqueue_pending_turn(&thread_id, pending) {
|
||||
Ok(_) => {}
|
||||
Err(EnqueueRejection::QueueFull) => {
|
||||
return Err(DirectTurnError::QueueFull {
|
||||
return Err(EnqueueError::QueueFull(QueueFull {
|
||||
limit: MAX_PENDING_TURNS,
|
||||
})
|
||||
}))
|
||||
}
|
||||
}
|
||||
// 入队之后立刻踢一脚:队列空且没有回合在跑时,放行就是这一脚,用户点发送不必再等一个调度周期。
|
||||
@@ -146,7 +148,23 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
fn user_item(text: &str) -> DirectCodexUserItem {
|
||||
/// 载荷里那份宿主原文:测试只关心"有没有原因 / 是哪一类",不关心前端怎么拼文案。
|
||||
fn failure_detail(failure: &TurnFailure) -> String {
|
||||
match failure {
|
||||
TurnFailure::ProjectRootUnanchored(payload) => payload.cause.clone(),
|
||||
TurnFailure::EnvironmentNotReady(payload) => payload.detail.clone(),
|
||||
TurnFailure::HostStateUnavailable(payload) => payload.detail.clone(),
|
||||
TurnFailure::ModelCallFailed(payload) => payload.detail.clone(),
|
||||
TurnFailure::TransportClosed(payload) => payload.diagnostic.clone(),
|
||||
TurnFailure::TimedOut(payload) => format!("{:?}", payload.deadline),
|
||||
TurnFailure::TurnInterrupted(payload) => payload.detail.clone(),
|
||||
TurnFailure::SuperErrorFromStringPlusStage(payload) => payload.detail.clone(),
|
||||
TurnFailure::Unclassified(payload) => payload.detail.clone(),
|
||||
TurnFailure::HostDropped => String::new(),
|
||||
}
|
||||
}
|
||||
|
||||
fn user_item(text: &str) -> UserItem {
|
||||
serde_json::from_value(serde_json::json!({
|
||||
"type": "message",
|
||||
"role": "user",
|
||||
@@ -188,10 +206,10 @@ mod tests {
|
||||
.collect::<Vec<_>>();
|
||||
assert_eq!(terminal.len(), 1, "一轮只许有一条终态:{events:?}");
|
||||
let (status, failure, user_item_id) = terminal[0];
|
||||
assert_eq!(status, "failed");
|
||||
assert_eq!(*status, TurnCompletedStatus::Failed);
|
||||
let failure = failure.as_ref().expect("失败终态必须带载荷");
|
||||
assert!(
|
||||
!failure.message.trim().is_empty(),
|
||||
!failure_detail(failure).trim().is_empty(),
|
||||
"放行之后的失败必须带上原因"
|
||||
);
|
||||
assert_eq!(user_item_id.as_deref(), Some("direct-codex:turn-1:user"));
|
||||
@@ -315,20 +333,13 @@ mod tests {
|
||||
.collect::<Vec<_>>();
|
||||
assert_eq!(terminals.len(), 1, "一轮只许有一条终态:{events:?}");
|
||||
let (status, failure) = terminals[0];
|
||||
assert_eq!(status, "failed");
|
||||
assert_eq!(*status, TurnCompletedStatus::Failed);
|
||||
let failure = failure.as_ref().expect("失败终态必须带载荷");
|
||||
assert!(
|
||||
failure.message.contains("写入本项目对话历史失败"),
|
||||
"{}",
|
||||
failure.message
|
||||
);
|
||||
let detail = failure_detail(failure);
|
||||
assert!(detail.contains("写入本项目对话历史失败"), "{detail}");
|
||||
// 这一轮已经放行,所以走的是**回合失败**:入队失败那套 `direct-codex-failure:v2` 收口文案
|
||||
// 不许出现在这里(它只属于可留痕的入队失败)。
|
||||
assert!(
|
||||
!failure.message.contains("direct-codex-failure"),
|
||||
"{}",
|
||||
failure.message
|
||||
);
|
||||
assert!(!detail.contains("direct-codex-failure"), "{detail}");
|
||||
// 占用已释放:下一轮还能继续。
|
||||
assert!(!crate::agent::thread_turn_is_active(&thread_id));
|
||||
}
|
||||
|
||||
@@ -3464,14 +3464,16 @@ mod tests {
|
||||
"title": "工具链远端画布"
|
||||
}]}}),
|
||||
);
|
||||
} else if route.starts_with("/api/editor/assets/library") {
|
||||
} else if request_line.starts_with(&format!(
|
||||
"GET /api/editor/assets/folders/{TOOL_CHAIN_ASSET_FOLDER_ID} "
|
||||
)) {
|
||||
tool_chain_write_json(
|
||||
&mut stream,
|
||||
"200 OK",
|
||||
json!({"data": {"library": {"folders": [{
|
||||
json!({"data": {"folder": {
|
||||
"folderId": TOOL_CHAIN_ASSET_FOLDER_ID,
|
||||
"label": "工具链远端目录"
|
||||
}]}}}),
|
||||
}}}),
|
||||
);
|
||||
} else if request_line.starts_with(expected_submission) {
|
||||
tool_chain_write_json(
|
||||
@@ -4017,6 +4019,9 @@ mod tests {
|
||||
);
|
||||
|
||||
assert_eq!(requests.len(), 6, "{requests:?}");
|
||||
assert!(requests.iter().any(|request| request.starts_with(&format!(
|
||||
"GET /api/editor/assets/folders/{TOOL_CHAIN_ASSET_FOLDER_ID} "
|
||||
))));
|
||||
let submission = requests
|
||||
.iter()
|
||||
.find(|request| request.starts_with("POST /api/editor/images/edits "))
|
||||
|
||||
@@ -1,283 +0,0 @@
|
||||
//! 失败终态的宿主侧策略:把"这一轮为什么失败"翻译成可下发的 `failure` 载荷,并在宿主自己
|
||||
//! 提前收场时补一条失败终态。
|
||||
//!
|
||||
//! 这个模块只有三件事,别再往里加第四件:
|
||||
//! 1. [`direct_turn_terminal`]:拿这一轮的事实判定终态——是不是失败、原因是什么、状态写什么;
|
||||
//! 2. [`DirectTurnTerminal::event`]:把终态投影成 `turn.completed` 事件。
|
||||
//!
|
||||
//! 终态的**出口**(谁写、什么时候兜底)不在这里,在 `thread_manager::dispatch` 的放行占用对象里:
|
||||
//! 这个模块只负责"什么算失败、原因怎么写"。
|
||||
//!
|
||||
//! 失败载荷的**形状**属于线上协议,定义在 `thread_manager::wire`(`DirectTurnFailure`);
|
||||
//! 载荷的 `kind` 与 `message` 由 [`DirectTurnError`] 投影而来(`kind` 的取值表见
|
||||
//! [`DirectTurnError::wire_kind`]);这里只负责"什么算失败、原因怎么写、什么时候兜底",
|
||||
//! 不碰事件队列的搬运规则,也不自己认 `LlmError`。
|
||||
|
||||
use std::path::Path;
|
||||
|
||||
use super::{
|
||||
redact_agent_runtime_error, DirectTurnError, DirectTurnFailure, DirectTurnFailureKind,
|
||||
ThreadEvent,
|
||||
};
|
||||
|
||||
/// `turn.completed.failure.message` 的字符上限:与本地错误文案同一档——够说清原因,又不至于
|
||||
/// 把整段上游报文塞进事件队列。
|
||||
const DIRECT_TURN_FAILURE_MESSAGE_MAX_CHARS: usize = 600;
|
||||
|
||||
/// 宿主任务提前结束(panic / future 被丢弃 / 终态之前的早退)时的分类与文案。
|
||||
const DIRECT_TURN_FAILURE_HOST_DROPPED_MESSAGE: &str =
|
||||
"陶泥儿回合的宿主任务提前结束(崩溃或任务被取消),本轮已按失败收口,请重试。";
|
||||
|
||||
/// 一轮的终态:写进事件的 `status` 与(失败时的)载荷。**状态由载荷反推**,不由收尾阶段推。
|
||||
pub(crate) struct DirectTurnTerminal {
|
||||
pub(crate) status: String,
|
||||
pub(crate) failure: Option<DirectTurnFailure>,
|
||||
}
|
||||
|
||||
impl DirectTurnTerminal {
|
||||
/// 一次成功终态:只带 `status="completed"`,不带失败载荷。
|
||||
///
|
||||
/// 供给没有"深层终态出口"的执行器(cc / Claude Code sidecar)用:它们整轮成功返回后,
|
||||
/// 线程仍被放行占用,必须由放行侧补写这条终态。
|
||||
pub(crate) fn completed() -> Self {
|
||||
Self {
|
||||
status: "completed".to_string(),
|
||||
failure: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// 终态事件:失败时同一个 `turn.completed` 带载荷,其余只带 `status`。
|
||||
pub(crate) fn event(self, completed_at: u64, user_item_id: Option<&str>) -> ThreadEvent {
|
||||
let event = match self.failure {
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
/// 拿这一轮的**事实**判定终态。判据按优先级:
|
||||
/// 1. `host_failure`:宿主自己观察 / 判定的失败(执行通道断开、等待超时、app-server 单方面中断…),
|
||||
/// 原因就用宿主当场写下的那句——它比交付报告更接近现场,报告只说明"收束到哪一步";
|
||||
/// 2. `collect_outcome` 是错误:真失败(模型 / 传输 / 历史落盘)。模型自报失败也走这一档:
|
||||
/// 原生 `turn/completed.status="failed"` 的 `error` 由调用点投影成 [`DirectTurnError`] 再进来;
|
||||
/// 3. `session_status` 已经判成 `failed`、而拿到的只是一份交付报告:原因用那份报告兜底——收尾
|
||||
/// 阶段的账本读不出来时只有它可用。
|
||||
///
|
||||
/// **有载荷就一定是 `failed`,没载荷就用收尾阶段的 `session_status`。** 这条反推关系是这个模块存在
|
||||
/// 的理由:`session_status` 是宿主收尾时按 ledger 阶段推的,收尾本身会把阶段推成 `Interrupted`,
|
||||
/// 于是"模型已经判失败"的一轮会被写成 `status="interrupted"` 且不带载荷——界面只剩"本轮已结束",
|
||||
/// 用户看不到任何原因(连接/上游断开时就是这个现象)。事实判失败就必须报失败。
|
||||
///
|
||||
/// 载荷的 `kind` 与 `message` 在这一个出口从 typed 错误投影:`kind` 决定界面语气,`message` 是脱敏
|
||||
/// 截断后的原因文本;Rust 侧没有第二个地方再解析它。
|
||||
pub(crate) fn direct_turn_terminal(
|
||||
session_status: &str,
|
||||
collect_outcome: Result<&str, DirectTurnError>,
|
||||
host_failure: Option<&DirectTurnError>,
|
||||
history_root: &Path,
|
||||
) -> DirectTurnTerminal {
|
||||
let failure = match (host_failure, collect_outcome) {
|
||||
(Some(failure), _) => Some(failure.clone()),
|
||||
(None, Err(error)) => Some(error.clone()),
|
||||
// 账本读不出来时没有 typed 原因可用:报告文本就是这一轮唯一的收口依据,按未分类失败发出去,
|
||||
// 不能让界面停在"已结束、没原因"。
|
||||
(None, Ok(report)) if session_status == "failed" => {
|
||||
Some(DirectTurnError::TurnFailedUnclassified {
|
||||
detail: report.to_string(),
|
||||
})
|
||||
}
|
||||
(None, Ok(_)) => None,
|
||||
};
|
||||
match failure {
|
||||
Some(failure) => DirectTurnTerminal::failed(history_root, &failure),
|
||||
None => DirectTurnTerminal {
|
||||
status: session_status.to_string(),
|
||||
failure: None,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
impl DirectTurnTerminal {
|
||||
/// 一次失败终态:`kind` 与 `message` 只在这一个出口从 typed 错误投影。
|
||||
pub(crate) fn failed(history_root: &Path, failure: &DirectTurnError) -> Self {
|
||||
Self {
|
||||
status: "failed".to_string(),
|
||||
failure: Some(DirectTurnFailure::new(
|
||||
failure
|
||||
.wire_kind()
|
||||
.unwrap_or(DirectTurnFailureKind::ModelFailed),
|
||||
redact_agent_runtime_error(
|
||||
history_root,
|
||||
&failure.to_string(),
|
||||
DIRECT_TURN_FAILURE_MESSAGE_MAX_CHARS,
|
||||
),
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
/// 宿主任务提前结束(panic / future 被丢弃 / 取消)的兜底终态。
|
||||
///
|
||||
/// 这类收场说不出原因,只给分类;能说清原因的一律走 [`Self::failed`]。
|
||||
pub(crate) fn host_dropped() -> Self {
|
||||
Self {
|
||||
status: "failed".to_string(),
|
||||
failure: Some(DirectTurnFailure::new(
|
||||
DirectTurnFailureKind::HostDropped,
|
||||
DIRECT_TURN_FAILURE_HOST_DROPPED_MESSAGE.to_string(),
|
||||
)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::agent::{consume_thread, subscribe_thread};
|
||||
use platform_llm::LlmError;
|
||||
|
||||
fn history_root() -> std::path::PathBuf {
|
||||
std::path::PathBuf::from("/tmp/direct-turn-failure-test")
|
||||
}
|
||||
|
||||
/// 正常收场:不带载荷,`status` 就用收尾阶段推出来的那个。
|
||||
#[test]
|
||||
fn non_failure_terminals_keep_the_session_status() {
|
||||
for status in ["completed", "interrupted", "aborted"] {
|
||||
let terminal = direct_turn_terminal(status, Ok("报告不重要"), None, &history_root());
|
||||
assert!(terminal.failure.is_none(), "{status} 不该带失败载荷");
|
||||
assert_eq!(terminal.status, status);
|
||||
}
|
||||
}
|
||||
|
||||
/// 拿得到错误:分类与原因都取自错误。
|
||||
#[test]
|
||||
fn collect_error_becomes_a_failure_terminal() {
|
||||
let error = DirectTurnError::from_model_call(&LlmError::Transport(
|
||||
"DirectProject 收尾历史失败:写入 project.jsonl 失败".into(),
|
||||
));
|
||||
let terminal = direct_turn_terminal("completed", Err(error), None, &history_root());
|
||||
let failure = terminal
|
||||
.failure
|
||||
.expect("transport error must fail the turn");
|
||||
assert_eq!(terminal.status, "failed");
|
||||
assert_eq!(failure.kind, DirectTurnFailureKind::TransportFailed);
|
||||
assert!(failure.message.contains("收尾历史失败"));
|
||||
}
|
||||
|
||||
/// **收尾阶段的中断不能把已经失败的一轮讲成"已结束"。** 模型自报失败在调用点被投影成 typed
|
||||
/// 错误(原因带 `codex-app-server-error:<kind>` 前缀),宿主收尾自己又把 ledger 阶段推成
|
||||
/// `Interrupted`(`session_status` 因此是 `interrupted`):事实就是失败、原因就是那份投影,
|
||||
/// 必须原样发出去——否则界面只剩"本轮已结束",用户看不到任何东西。
|
||||
#[test]
|
||||
fn projected_native_failure_outranks_the_interrupted_session_status() {
|
||||
let error = DirectTurnError::from_model_call(&LlmError::InvalidRequest(
|
||||
"codex-app-server-error:context-window-exceeded".into(),
|
||||
));
|
||||
let terminal = direct_turn_terminal("interrupted", Err(error), None, &history_root());
|
||||
let failure = terminal.failure.expect("native failure must fail the turn");
|
||||
assert_eq!(terminal.status, "failed");
|
||||
assert_eq!(failure.kind, DirectTurnFailureKind::RequestRejected);
|
||||
assert_eq!(
|
||||
failure.message,
|
||||
"codex-app-server-error:context-window-exceeded"
|
||||
);
|
||||
}
|
||||
|
||||
/// 收尾阶段的账本读不出来(`session_status` 只能是 `failed`)时没有错误可用:用交付报告兜底,
|
||||
/// 但照样要带载荷发出去,不能让界面停在"已结束、没原因"。
|
||||
#[test]
|
||||
fn unreadable_session_ledger_still_reports_a_payload() {
|
||||
let terminal = direct_turn_terminal("failed", Ok("报告"), None, &history_root());
|
||||
assert_eq!(terminal.status, "failed");
|
||||
let failure = terminal
|
||||
.failure
|
||||
.expect("unreadable ledger must fail the turn");
|
||||
assert_eq!(failure.kind, DirectTurnFailureKind::ModelFailed);
|
||||
assert_eq!(failure.message, "报告");
|
||||
}
|
||||
|
||||
/// 宿主自己记下的失败排在最前面:它比交付报告更接近现场。
|
||||
#[test]
|
||||
fn host_recorded_failure_outranks_every_other_source() {
|
||||
let diagnostic = "Codex app-server 已退出;exitStatus=signal: 9 (SIGKILL);\
|
||||
stderrClass=nonempty;stderrBytes=1000";
|
||||
let host_failure = DirectTurnError::TransportClosed {
|
||||
diagnostic: diagnostic.to_string(),
|
||||
};
|
||||
let terminal = direct_turn_terminal(
|
||||
"interrupted",
|
||||
Ok("执行连接已结束,正在核对自有子进程与在途操作。"),
|
||||
Some(&host_failure),
|
||||
&history_root(),
|
||||
);
|
||||
let failure = terminal.failure.expect("host fact must fail the turn");
|
||||
assert_eq!(terminal.status, "failed");
|
||||
assert_eq!(failure.kind, DirectTurnFailureKind::TransportFailed);
|
||||
assert!(failure.message.contains("SIGKILL"));
|
||||
assert!(!failure.message.contains("正在核对自有子进程"));
|
||||
|
||||
// 即使同时拿到了错误,宿主亲眼看到的事实仍然是第一顺位。
|
||||
let error = DirectTurnError::from_model_call(&LlmError::Transport(
|
||||
"DirectProject 收尾历史失败".into(),
|
||||
));
|
||||
let host_failure = DirectTurnError::TurnInterrupted {
|
||||
detail: "本轮模型执行被中断".into(),
|
||||
};
|
||||
let terminal = direct_turn_terminal(
|
||||
"interrupted",
|
||||
Err(error),
|
||||
Some(&host_failure),
|
||||
&history_root(),
|
||||
);
|
||||
let failure = terminal.failure.expect("host fact must fail the turn");
|
||||
assert_eq!(failure.kind, DirectTurnFailureKind::TurnInterrupted);
|
||||
assert!(failure.message.contains("本轮模型执行被中断"));
|
||||
}
|
||||
|
||||
/// 终态事件的形状:失败时同一个 `turn.completed` 带载荷,其余只带 `status`。
|
||||
#[test]
|
||||
fn terminal_event_carries_the_payload_and_the_opening_identity() {
|
||||
let error = DirectTurnError::from_model_call(&LlmError::Upstream {
|
||||
status_code: 502,
|
||||
message: "上游 502".into(),
|
||||
});
|
||||
let failing = direct_turn_terminal("interrupted", Err(error), None, &history_root());
|
||||
let event = failing.event(2_000, Some("direct-codex:turn-1:user"));
|
||||
assert_eq!(
|
||||
event.failure().map(|failure| failure.kind),
|
||||
Some(DirectTurnFailureKind::ModelFailed)
|
||||
);
|
||||
assert_eq!(event.user_item_id(), Some("direct-codex:turn-1:user"));
|
||||
assert_eq!(event.at(), Some(2_000));
|
||||
|
||||
let quiet = direct_turn_terminal("completed", Ok("本轮交付已完成"), None, &history_root());
|
||||
let event = quiet.event(3_000, None);
|
||||
assert!(event.failure().is_none());
|
||||
assert!(matches!(
|
||||
event,
|
||||
ThreadEvent::TurnCompleted { ref status, .. } if status == "completed"
|
||||
));
|
||||
}
|
||||
|
||||
/// 兜底终态:说不出原因的那一种只给分类,不冒充真实原因。
|
||||
#[test]
|
||||
fn host_dropped_terminal_only_carries_the_classification() {
|
||||
let terminal = DirectTurnTerminal::host_dropped();
|
||||
assert_eq!(terminal.status, "failed");
|
||||
let failure = terminal.failure.expect("host-dropped must fail the turn");
|
||||
assert_eq!(failure.kind, DirectTurnFailureKind::HostDropped);
|
||||
assert!(!failure.message.trim().is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn completed_terminal_carries_no_failure_payload() {
|
||||
let terminal = DirectTurnTerminal::completed();
|
||||
assert_eq!(terminal.status, "completed");
|
||||
assert!(terminal.failure.is_none());
|
||||
assert!(matches!(
|
||||
terminal.event(1_700_000_000_000, Some("item-1")),
|
||||
ThreadEvent::TurnCompleted { ref status, .. } if status == "completed"
|
||||
));
|
||||
}
|
||||
}
|
||||
File diff suppressed because it is too large
Load Diff
@@ -629,7 +629,7 @@ const ERROR_REDACTED_KEY: &str = "[redacted-sensitive-field]";
|
||||
/// information whenever a safe error line mentions a credential field.
|
||||
///
|
||||
/// 逐行脱敏,但**保留每个 chunk 末尾的换行**:本函数会被流式增量逐段调用
|
||||
/// (`thread_delta_text` → 前端把各段拼成一条消息再交给 Markdown 渲染)。用
|
||||
/// (脱敏后的各段由前端拼成一条消息再交给 Markdown 渲染)。用
|
||||
/// `lines()` + `join("\n")` 会把「以换行结尾的段」的末尾换行吃掉,拼接后段落、列表项和表格行
|
||||
/// 会并进同一行,整条消息的 Markdown 结构(尤其表格)就作废了。
|
||||
pub(crate) fn sanitize_error_context(value: &str) -> String {
|
||||
|
||||
@@ -19,10 +19,10 @@ use futures::FutureExt;
|
||||
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_now_ms, record_direct_codex_failure,
|
||||
redact_agent_runtime_error, run_direct_game_creator_turn_at_with_creation_type_and_emitter,
|
||||
thread_id_for_project, DirectCodexFailureStage, DirectGameCreatorTurnUpdateEmitter,
|
||||
DirectTurnError, DirectTurnTerminal, DispatchedTurn,
|
||||
direct_codex_user_item_to_prompt, now_ms, record_direct_codex_failure,
|
||||
run_direct_game_creator_turn_at_with_creation_type_and_emitter, thread_id_for_project,
|
||||
DirectGameCreatorTurnUpdateEmitter, DispatchedTurn, FailureStage, TurnCompletion, TurnError,
|
||||
TurnErrorClassified,
|
||||
};
|
||||
|
||||
/// 一次放行的占用。持有它就代表这一轮还没收口。
|
||||
@@ -54,11 +54,11 @@ impl TurnReservation {
|
||||
/// 放行之后还没走到深层终态就失败的收口口:只有这一轮仍被自己占用时才写。
|
||||
///
|
||||
/// 深层(真正跑完这一轮的代码)已经写出终态时返回 `false`,兜底不覆盖真实结果。
|
||||
pub(crate) fn finish_if_unfinished(&self, terminal: DirectTurnTerminal) -> bool {
|
||||
pub(crate) fn finish_if_unfinished(&self, terminal: TurnCompletion) -> bool {
|
||||
complete_turn_if_reserved(
|
||||
&self.thread_id,
|
||||
&self.token,
|
||||
terminal.event(direct_now_ms(), self.user_item_id.as_deref()),
|
||||
terminal.event(now_ms(), self.user_item_id.as_deref()),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -75,7 +75,7 @@ impl TurnReservation {
|
||||
}))
|
||||
.expect("canonical user item"),
|
||||
None,
|
||||
direct_now_ms(),
|
||||
now_ms(),
|
||||
);
|
||||
super::enqueue_pending_turn(thread_id, pending).expect("enqueue test turn");
|
||||
let dispatched = claim_pending_turn(thread_id).expect("claim test turn");
|
||||
@@ -87,7 +87,7 @@ impl Drop for TurnReservation {
|
||||
fn drop(&mut self) {
|
||||
// 兜底:任务 panic、future 被丢弃、或今后在终态之前新增的 `?` 早退。
|
||||
// 这类失败说不出原因,只给分类;能说清原因的错误必须由调用方在更早的地方显式收口。
|
||||
let finalized_by_guard = self.finish_if_unfinished(DirectTurnTerminal::host_dropped());
|
||||
let finalized_by_guard = self.finish_if_unfinished(TurnCompletion::host_dropped());
|
||||
// 兜底一旦真的收口,就说明这一轮**从未写下终态**:深层既没成功也没失败地退出了。
|
||||
// 这条必须留应用日志,否则离线只剩一个 `phase=working` 的账本,无从判断是哪一层
|
||||
// 提前退出(真实案例:2026-10-02 连续两轮只留下 working 账本,errors/ 与
|
||||
@@ -96,16 +96,16 @@ impl Drop for TurnReservation {
|
||||
if finalized_by_guard {
|
||||
// 项目侧也留一份脱敏诊断:用户在项目目录里就能看到这一轮的收场,
|
||||
// 不必只依赖 AppData 的应用日志。
|
||||
let failure = DirectTurnError::TurnFailed {
|
||||
stage: DirectCodexFailureStage::CodeGeneration,
|
||||
// 字段名用 `tt`(不是 `turnToken`):`turnToken=` 会命中应用日志的凭据标记,
|
||||
// 整行被替换成 `<sensitive diagnostic details redacted>`,离线就只剩一个
|
||||
// 说不出原因的 HostDropped。
|
||||
detail: format!(
|
||||
// 字段名用 `tt`(不是 `turnToken`):`turnToken=` 会命中应用日志的凭据标记,
|
||||
// 整行被替换成 `<sensitive diagnostic details redacted>`,离线就只剩一个
|
||||
// 说不出原因的 HostDropped。
|
||||
let failure = TurnError::turn_failed(
|
||||
FailureStage::CodeGeneration,
|
||||
format!(
|
||||
"DirectProject 宿主任务提前结束(panic、future 被丢弃或被取消),本轮未写下终态;threadId={} tt={}",
|
||||
self.thread_id, self.token
|
||||
),
|
||||
};
|
||||
);
|
||||
let _ = record_direct_codex_failure(Path::new(&self.thread_id), &failure, None);
|
||||
app_log!(
|
||||
"agent.direct_turn.host_dropped threadId={} tt={} userItemId={} panicking={}",
|
||||
@@ -222,14 +222,17 @@ async fn run_dispatched_direct_turn(
|
||||
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));
|
||||
// 原文不在这里脱敏:`classify` 在投影成失败载荷时统一脱敏 + 截断(错误文案的脱敏保留)。
|
||||
let failure = TurnError::environment_not_ready(format!("写入本项目对话历史失败:{error}"));
|
||||
// 这条错误是本地构造的**真失败**:classify 只可能给 ShouldStop。
|
||||
match failure.classify(&root) {
|
||||
TurnErrorClassified::ShouldStop(payload) => {
|
||||
reservation.finish_if_unfinished(TurnCompletion::failed(payload));
|
||||
}
|
||||
TurnErrorClassified::ShouldContinue { detail } => {
|
||||
unreachable!("environment_not_ready 只可能是回合失败,不该分类成控制流:{detail}")
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
// 用户条目落盘成功即下发:这一轮从"放行"到"起 codex"之间的一切失败(连不上 app-server、执行器
|
||||
@@ -257,16 +260,17 @@ async fn run_dispatched_direct_turn(
|
||||
{
|
||||
Ok(result) => result,
|
||||
Err(payload) => {
|
||||
let failure = DirectTurnError::turn_failed(
|
||||
DirectCodexFailureStage::CodeGeneration,
|
||||
let failure = TurnError::turn_failed(
|
||||
FailureStage::CodeGeneration,
|
||||
format!(
|
||||
"DirectProject 宿主任务 panic:{}",
|
||||
direct_turn_panic_detail(payload.as_ref())
|
||||
),
|
||||
);
|
||||
let stage = failure.turn_failure_stage();
|
||||
let detail = record_direct_codex_failure(&root, &failure, Some(turn_id.as_str()));
|
||||
Err(DirectTurnError::TurnFailed { stage, detail })
|
||||
// 审计照写,载荷原样透传:panic 本身已经是一条 `SuperErrorFromStringPlusStage`,
|
||||
// 不需要再包一层。
|
||||
let _ = record_direct_codex_failure(&root, &failure, Some(turn_id.as_str()));
|
||||
Err(failure)
|
||||
}
|
||||
};
|
||||
match outcome {
|
||||
@@ -277,12 +281,24 @@ async fn run_dispatched_direct_turn(
|
||||
// 线程仍被这条放行占用,不补写终态的话,占用对象 Drop 时的兜底会把一轮已经拿到回复的
|
||||
// 回合收成 `HostDropped`(真实案例:2026-10-03 `claude-parse-done chars=176` 之后立刻
|
||||
// `host_dropped panicking=false`)。codex 路径已写过终态,这里是空操作。
|
||||
reservation.finish_if_unfinished(DirectTurnTerminal::completed());
|
||||
reservation.finish_if_unfinished(TurnCompletion::Completed);
|
||||
}
|
||||
Err(error) => {
|
||||
// 放行之后的失败一律是回合失败:失败诊断与失败说明已由上层写过,这里补终态事件。
|
||||
// 深层已经写出终态时它不覆盖(同一轮只允许一条终态)。
|
||||
reservation.finish_if_unfinished(DirectTurnTerminal::failed(&root, &error));
|
||||
// 按 `classify` 显式分流,不留静默分支:控制流(返修 / 复核要求继续)本应被
|
||||
// `direct_runtime` 的返修循环消化成下一轮;两条来源(交付复核、app-server 封口复核)
|
||||
// 都在循环里接住了。漏到这里就是上游缺陷——直接炸出来,不再静默吞掉。
|
||||
match error.classify(&root) {
|
||||
TurnErrorClassified::ShouldStop(payload) => {
|
||||
reservation.finish_if_unfinished(TurnCompletion::failed(payload));
|
||||
}
|
||||
TurnErrorClassified::ShouldContinue { detail } => {
|
||||
unreachable!(
|
||||
"控制流错误到达回合失败收口(应在 direct_runtime 的返修循环内消化):{detail}"
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -291,8 +307,8 @@ async fn run_dispatched_direct_turn(
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::agent::{
|
||||
consume_thread, enqueue_pending_turn, subscribe_thread, thread_turn_is_active,
|
||||
DirectTurnFailure, DirectTurnFailureKind, PendingTurn, ThreadEvent,
|
||||
consume_thread, enqueue_pending_turn, subscribe_thread, thread_turn_is_active, Deadline,
|
||||
PendingTurn, ThreadEvent, TimedOut, TurnCompletedStatus, TurnFailure,
|
||||
};
|
||||
use uuid::Uuid;
|
||||
|
||||
@@ -326,7 +342,7 @@ mod tests {
|
||||
}))
|
||||
.expect("canonical user item"),
|
||||
None,
|
||||
direct_now_ms(),
|
||||
now_ms(),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -355,9 +371,9 @@ mod tests {
|
||||
},
|
||||
);
|
||||
|
||||
assert!(!foreign.finish_if_unfinished(DirectTurnTerminal::host_dropped()));
|
||||
assert!(!foreign.finish_if_unfinished(TurnCompletion::host_dropped()));
|
||||
assert!(thread_turn_is_active(&thread));
|
||||
assert!(owner.finish_if_unfinished(DirectTurnTerminal::host_dropped()));
|
||||
assert!(owner.finish_if_unfinished(TurnCompletion::host_dropped()));
|
||||
assert!(!thread_turn_is_active(&thread));
|
||||
// 显式收口之后 Drop 不再补第二条:兜底只负责"没人写过"的那一种。
|
||||
drop(foreign);
|
||||
@@ -375,17 +391,16 @@ mod tests {
|
||||
|
||||
// 深层收口:真正跑完这一轮的代码算出来的终态。
|
||||
let deep = ThreadEvent::turn_completed_failed(
|
||||
DirectTurnFailure::new(
|
||||
DirectTurnFailureKind::Timeout,
|
||||
"等待模型回执超时".to_string(),
|
||||
),
|
||||
TurnFailure::TimedOut(TimedOut {
|
||||
deadline: Deadline::ResponseIdle,
|
||||
}),
|
||||
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()),
|
||||
!reservation.finish_if_unfinished(TurnCompletion::host_dropped()),
|
||||
"深层已收口时兜底不许再写"
|
||||
);
|
||||
drop(reservation);
|
||||
@@ -396,8 +411,10 @@ mod tests {
|
||||
match completed[0] {
|
||||
ThreadEvent::TurnCompleted { failure, .. } => {
|
||||
assert_eq!(
|
||||
failure.as_ref().map(|f| f.kind),
|
||||
Some(DirectTurnFailureKind::Timeout)
|
||||
failure,
|
||||
&Some(TurnFailure::TimedOut(TimedOut {
|
||||
deadline: Deadline::ResponseIdle,
|
||||
}))
|
||||
);
|
||||
}
|
||||
other => panic!("expected a terminal, got {other:?}"),
|
||||
@@ -426,7 +443,7 @@ mod tests {
|
||||
matches!(
|
||||
events.get(terminal),
|
||||
Some(ThreadEvent::TurnCompleted { status, failure: Some(failure), .. })
|
||||
if status == "failed" && failure.kind == DirectTurnFailureKind::HostDropped
|
||||
if *status == TurnCompletedStatus::Failed && matches!(failure, TurnFailure::HostDropped)
|
||||
),
|
||||
"{events:?}"
|
||||
);
|
||||
@@ -489,7 +506,7 @@ mod tests {
|
||||
assert!(
|
||||
events.iter().any(|event| matches!(
|
||||
event,
|
||||
ThreadEvent::TurnCompleted { status, .. } if status == "aborted"
|
||||
ThreadEvent::TurnCompleted { status, .. } if *status == TurnCompletedStatus::Aborted
|
||||
)),
|
||||
"{events:?}"
|
||||
);
|
||||
|
||||
@@ -6,10 +6,12 @@
|
||||
|
||||
pub(crate) mod dispatch;
|
||||
pub(crate) mod queue;
|
||||
pub(crate) mod turn_completion;
|
||||
pub(crate) mod wire;
|
||||
|
||||
pub(crate) use dispatch::*;
|
||||
pub(crate) use queue::*;
|
||||
pub(crate) use turn_completion::*;
|
||||
pub(crate) use wire::*;
|
||||
|
||||
use std::collections::{HashMap, HashSet};
|
||||
@@ -17,7 +19,7 @@ use std::sync::{Mutex, OnceLock};
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::agent::{
|
||||
direct_codex_user_item_id_for_client_turn_id, direct_now_ms, queue_has_room, ConsumeResult,
|
||||
direct_codex_user_item_id_for_client_turn_id, now_ms, queue_has_room, ConsumeResult,
|
||||
EnqueueOutcome, EnqueueRejection, PendingTurn, QueueRemovalOutcome, QueueRemovalReason,
|
||||
SubscriptionBootstrap, ThreadEvent,
|
||||
};
|
||||
@@ -377,7 +379,8 @@ impl ThreadManager {
|
||||
// `turn.started` 的口径同一份(`PendingTurn::user_item_id`)。
|
||||
self.append(
|
||||
thread_id,
|
||||
ThreadEvent::turn_completed("aborted".to_string(), now_ms).with_user_item_id(
|
||||
TurnCompletion::Aborted.event(
|
||||
now_ms,
|
||||
direct_codex_user_item_id_for_client_turn_id(&released).as_deref(),
|
||||
),
|
||||
);
|
||||
@@ -854,7 +857,7 @@ pub(crate) fn remove_pending_turn(thread_id: &str, client_turn_id: &str) -> Queu
|
||||
global_thread_manager()
|
||||
.lock()
|
||||
.unwrap_or_else(|poisoned| poisoned.into_inner())
|
||||
.remove_pending_turn(thread_id, client_turn_id, direct_now_ms())
|
||||
.remove_pending_turn(thread_id, client_turn_id, now_ms())
|
||||
};
|
||||
if outcome == QueueRemovalOutcome::Removed {
|
||||
notify_subscribers(thread_id);
|
||||
@@ -873,7 +876,7 @@ pub(crate) fn claim_pending_turn(thread_id: &str) -> Option<DispatchedTurn> {
|
||||
global_thread_manager()
|
||||
.lock()
|
||||
.unwrap_or_else(|poisoned| poisoned.into_inner())
|
||||
.claim_pending_turn(thread_id, direct_now_ms())
|
||||
.claim_pending_turn(thread_id, now_ms())
|
||||
};
|
||||
if claimed.is_some() {
|
||||
notify_subscribers(thread_id);
|
||||
@@ -950,12 +953,7 @@ pub(crate) fn release_stale_direct_turn(
|
||||
global_thread_manager()
|
||||
.lock()
|
||||
.unwrap_or_else(|poisoned| poisoned.into_inner())
|
||||
.release_stale_turn(
|
||||
thread_id,
|
||||
expected_client_turn_id,
|
||||
min_age_ms,
|
||||
direct_now_ms(),
|
||||
)
|
||||
.release_stale_turn(thread_id, expected_client_turn_id, min_age_ms, now_ms())
|
||||
};
|
||||
if matches!(outcome, StaleTurnRelease::Released(_)) {
|
||||
notify_subscribers(thread_id);
|
||||
@@ -1243,13 +1241,13 @@ mod tests {
|
||||
let mut manager = ThreadManager::with_limits(100, 100_000);
|
||||
manager.append(
|
||||
"thread-1",
|
||||
ThreadEvent::turn_completed("completed".to_string(), FIXED_AT_MS),
|
||||
ThreadEvent::turn_completed(TurnOutcome::Completed, FIXED_AT_MS),
|
||||
);
|
||||
let bootstrap = manager.subscribe("thread-1");
|
||||
assert!(matches!(
|
||||
bootstrap.events.as_slice(),
|
||||
[ThreadEvent::TurnCompleted { status, at, .. }]
|
||||
if status == "completed" && *at == Some(FIXED_AT_MS)
|
||||
if *status == TurnCompletedStatus::Completed && *at == Some(FIXED_AT_MS)
|
||||
));
|
||||
}
|
||||
|
||||
@@ -1261,22 +1259,18 @@ mod tests {
|
||||
manager.append("thread-1", ThreadEvent::turn_started(1_000));
|
||||
manager.append(
|
||||
"thread-1",
|
||||
ThreadEvent::turn_completed_failed(
|
||||
crate::agent::DirectTurnFailure::new(
|
||||
crate::agent::DirectTurnFailureKind::HostDropped,
|
||||
"回合宿主任务提前结束",
|
||||
),
|
||||
FIXED_AT_MS,
|
||||
),
|
||||
ThreadEvent::turn_completed_failed(crate::agent::TurnFailure::HostDropped, FIXED_AT_MS),
|
||||
);
|
||||
|
||||
let bootstrap = manager.subscribe("thread-1");
|
||||
assert!(matches!(
|
||||
bootstrap.events.as_slice(),
|
||||
[ThreadEvent::TurnCompleted { status, failure, at, .. }]
|
||||
if status == "failed"
|
||||
&& failure.as_ref().is_some_and(|failure| failure.kind
|
||||
== crate::agent::DirectTurnFailureKind::HostDropped)
|
||||
if *status == TurnCompletedStatus::Failed
|
||||
&& failure.as_ref().is_some_and(|failure| matches!(
|
||||
failure,
|
||||
crate::agent::TurnFailure::HostDropped
|
||||
))
|
||||
&& *at == Some(FIXED_AT_MS)
|
||||
));
|
||||
}
|
||||
@@ -1356,7 +1350,7 @@ mod tests {
|
||||
);
|
||||
manager.append(
|
||||
"thread-1",
|
||||
ThreadEvent::turn_completed("completed".to_string(), 3_000)
|
||||
ThreadEvent::turn_completed(TurnOutcome::Completed, 3_000)
|
||||
.with_user_item_id(Some("direct-codex:turn-1:user")),
|
||||
);
|
||||
let second = manager.subscribe("thread-1");
|
||||
@@ -1445,7 +1439,7 @@ mod tests {
|
||||
|
||||
manager.complete_turn(
|
||||
thread_id,
|
||||
ThreadEvent::turn_completed("completed".to_string(), 5_000),
|
||||
ThreadEvent::turn_completed(TurnOutcome::Completed, 5_000),
|
||||
);
|
||||
assert!(snapshot_of(&manager, thread_id).is_none());
|
||||
}
|
||||
@@ -1739,7 +1733,7 @@ mod tests {
|
||||
assert!(
|
||||
matches!(
|
||||
events[0],
|
||||
ThreadEvent::TurnCompleted { ref status, .. } if status == "aborted"
|
||||
ThreadEvent::TurnCompleted { ref status, .. } if *status == TurnCompletedStatus::Aborted
|
||||
),
|
||||
"{events:?}"
|
||||
);
|
||||
@@ -1838,7 +1832,7 @@ mod tests {
|
||||
|
||||
complete_turn(
|
||||
&thread_id,
|
||||
ThreadEvent::turn_completed("completed".to_string(), 5_000),
|
||||
ThreadEvent::turn_completed(TurnOutcome::Completed, 5_000),
|
||||
);
|
||||
assert_eq!(
|
||||
crate::agent::direct_active_turns_event_test_count(),
|
||||
|
||||
@@ -9,9 +9,7 @@
|
||||
//! 这条判据的成员折叠(`ThreadManager::pending_turns`)。
|
||||
//! 它不碰锁、不碰 Tauri、不写盘、不重算 prompt:入队检查在命令侧,放行顺序在 Thread Manager。
|
||||
|
||||
use crate::agent::{
|
||||
direct_codex_user_item_id_for_client_turn_id, DirectCodexUserItem, ThreadEvent,
|
||||
};
|
||||
use crate::agent::{direct_codex_user_item_id_for_client_turn_id, ThreadEvent, UserItem};
|
||||
|
||||
/// 一个项目最多能同时排队的待发消息条数。
|
||||
///
|
||||
@@ -30,7 +28,7 @@ pub(crate) struct PendingTurn {
|
||||
/// 这条消息的回合身份;放行后同一轮的 `turn.started` / `turn.completed` 用它。
|
||||
pub(crate) client_turn_id: String,
|
||||
/// canonical 用户条目:事件与界面 chip 都读它,Rust 不渲染展示形状。
|
||||
pub(crate) user_item: DirectCodexUserItem,
|
||||
pub(crate) user_item: UserItem,
|
||||
pub(crate) creation_type: Option<String>,
|
||||
/// 入队那一刻的宿主毫秒钟。
|
||||
pub(crate) at: u64,
|
||||
@@ -39,7 +37,7 @@ pub(crate) struct PendingTurn {
|
||||
impl PendingTurn {
|
||||
pub(crate) fn new(
|
||||
client_turn_id: String,
|
||||
user_item: DirectCodexUserItem,
|
||||
user_item: UserItem,
|
||||
creation_type: Option<String>,
|
||||
at: u64,
|
||||
) -> Self {
|
||||
@@ -116,7 +114,7 @@ mod tests {
|
||||
use super::*;
|
||||
use serde_json::json;
|
||||
|
||||
fn user_item(text: &str, id: &str) -> DirectCodexUserItem {
|
||||
fn user_item(text: &str, id: &str) -> UserItem {
|
||||
serde_json::from_value(json!({
|
||||
"type": "message",
|
||||
"role": "user",
|
||||
|
||||
@@ -0,0 +1,317 @@
|
||||
//! 一轮的**终态**:宿主侧的内部判别联合,不下发、不导出。
|
||||
//!
|
||||
//! 属于线上协议的只有失败**载荷** [`TurnFailure`](`turn.completed.failure` 的形状,留在
|
||||
//! [`super::wire`]);"这一轮是怎么收场的"是宿主自己的概念,所以放在这里、不进 wire,也绝不
|
||||
//! 加 `Serialize` / `TS`——一旦能序列化就会被误当成线上形状。
|
||||
//!
|
||||
//! 这里只有两件事,别再往里加第三件:
|
||||
//! 1. [`turn_terminal`]:拿这一轮的事实判定终态——是不是失败、原因是什么、状态写什么;
|
||||
//! 2. [`TurnCompletion::event`]:把终态投影成 `turn.completed` 事件(`status` 由变体反推)。
|
||||
//!
|
||||
//! 终态的**出口**(谁写、什么时候兜底)不在这里,在 [`super::dispatch`] 的放行占用对象里:
|
||||
//! 这里只负责"什么算失败、原因怎么写"。
|
||||
//!
|
||||
//! 载荷与脱敏都由 [`TurnError::classify`] 一处投影:Rust 侧没有第二个地方再拼它,也没有任何
|
||||
//! 地方再解析它。
|
||||
|
||||
use std::path::Path;
|
||||
|
||||
use crate::agent::{ThreadEvent, TurnError, TurnErrorClassified, TurnOutcome, Unclassified};
|
||||
|
||||
use super::wire::TurnFailure;
|
||||
|
||||
/// 一轮的收场:正常收场只有前三档,失败必须带载荷。
|
||||
///
|
||||
/// **有载荷就一定是失败,没载荷才看收尾阶段推出来的 `status`**——这条反推关系是这个类型存在的
|
||||
/// 理由:`session_status` 是宿主收尾时按 ledger 阶段推的,收尾本身会把阶段推成 `Interrupted`,
|
||||
/// 于是"模型已经判失败"的一轮会被写成 `status="interrupted"` 且不带载荷——界面只剩"本轮已结束",
|
||||
/// 用户看不到任何原因(连接/上游断开时就是这个现象)。事实判失败就必须报失败。
|
||||
#[derive(Clone, Debug, Eq, PartialEq)]
|
||||
pub(crate) enum TurnCompletion {
|
||||
Completed,
|
||||
Interrupted,
|
||||
/// 执行进程已退出 / 这一轮从没进执行器:不会再有人替它发终态的兜底收场。
|
||||
Aborted,
|
||||
Failed(TurnFailure),
|
||||
}
|
||||
|
||||
impl TurnCompletion {
|
||||
/// 终态事件:失败时同一个 `turn.completed` 带载荷,其余只带 `status`。
|
||||
pub(crate) fn event(self, completed_at: u64, user_item_id: Option<&str>) -> ThreadEvent {
|
||||
let event = match self {
|
||||
Self::Failed(failure) => ThreadEvent::turn_completed_failed(failure, completed_at),
|
||||
Self::Completed => ThreadEvent::turn_completed(TurnOutcome::Completed, completed_at),
|
||||
Self::Interrupted => {
|
||||
ThreadEvent::turn_completed(TurnOutcome::Interrupted, completed_at)
|
||||
}
|
||||
Self::Aborted => ThreadEvent::turn_completed(TurnOutcome::Aborted, completed_at),
|
||||
};
|
||||
event.with_user_item_id(user_item_id)
|
||||
}
|
||||
|
||||
/// 一次失败终态:载荷已由 [`TurnError::classify`] 投影并脱敏。
|
||||
pub(crate) fn failed(failure: TurnFailure) -> Self {
|
||||
Self::Failed(failure)
|
||||
}
|
||||
|
||||
/// 宿主任务提前结束(panic / future 被丢弃 / 取消)的兜底终态。
|
||||
///
|
||||
/// 这类收场说不出原因;能说清原因的一律走 [`Self::failed`]。
|
||||
pub(crate) fn host_dropped() -> Self {
|
||||
Self::Failed(TurnFailure::HostDropped)
|
||||
}
|
||||
}
|
||||
|
||||
/// 拿这一轮的**事实**判定终态。判据按优先级:
|
||||
/// 1. `host_failure`:宿主自己观察 / 判定的失败(执行通道断开、等待超时、app-server 单方面中断…),
|
||||
/// 原因就用宿主当场写下的那句——它比交付报告更接近现场,报告只说明"收束到哪一步";
|
||||
/// 2. `collect_outcome` 是错误:真失败(模型 / 传输 / 历史落盘)。模型自报失败也走这一档:
|
||||
/// 原生 `turn/completed.status="failed"` 的 `error` 由调用点投影成 [`TurnError`] 再进来;
|
||||
/// 3. `session_status` 已经判成 `failed`、而拿到的只是一份交付报告:原因用那份报告兜底——收尾
|
||||
/// 阶段的账本读不出来时只有它可用。
|
||||
///
|
||||
/// 错误由 [`TurnError::classify`] 分成两层:`ShouldStop` 的载荷直接成失败终态;`ShouldContinue`
|
||||
/// (返修 / 复核要求继续)说明这一轮还没结束,在产生 / 消费它的那一层就被消化,
|
||||
/// `codex_app_server` 的终态投影也把它过滤成了 `None`。它不该、也不能进入终态判定:
|
||||
/// 漏进来就是上游缺陷,直接 `unreachable!`,不落任何终态。
|
||||
///
|
||||
/// 载荷在这一个出口从 typed 错误投影([`TurnError::classify`]):脱敏与截断也在那一处完成,
|
||||
/// Rust 侧没有第二个地方再拼它、也没有任何地方再解析它。
|
||||
pub(crate) fn turn_terminal(
|
||||
session_status: &str,
|
||||
collect_outcome: Result<&str, TurnError>,
|
||||
host_failure: Option<&TurnError>,
|
||||
history_root: &Path,
|
||||
) -> TurnCompletion {
|
||||
let error = match (host_failure, collect_outcome) {
|
||||
(Some(failure), _) => Some(failure.clone()),
|
||||
// `collect_outcome` 按值匹配,`Err(error)` 已经把 `TurnError` 移出来,不必再 clone。
|
||||
(None, Err(error)) => Some(error),
|
||||
// 账本读不出来时(`session_status == "failed"`)没有 typed 原因可用:报告文本就是这一轮
|
||||
// 唯一的收口依据,按未分类失败发出去,不能让界面停在"已结束、没原因"。
|
||||
(None, Ok(report)) if session_status == "failed" => {
|
||||
Some(TurnError::Unclassified(Unclassified {
|
||||
detail: report.to_string(),
|
||||
}))
|
||||
}
|
||||
(None, Ok(_)) => None,
|
||||
};
|
||||
if let Some(error) = error {
|
||||
match error.classify(history_root) {
|
||||
TurnErrorClassified::ShouldStop(payload) => return TurnCompletion::Failed(payload),
|
||||
// 控制流:这一轮还没结束,不该走到终态判定。上游的返修循环与 `direct_turn_terminal_write`
|
||||
// 都已把它挡在外面;漏进来就是缺陷,直接炸出来,不落任何终态。
|
||||
TurnErrorClassified::ShouldContinue { detail } => {
|
||||
unreachable!("控制流错误不应进入终态判定:{detail}")
|
||||
}
|
||||
}
|
||||
}
|
||||
session_completion(session_status).unwrap_or_else(|| {
|
||||
// 收尾阶段的 `status` 认不出来(当前不可能发生):宁可报一条说不出原因的失败,也不冒充
|
||||
// 正常收场;载荷照样从 typed 错误投影,保持"只在一处拼载荷"。
|
||||
let error = TurnError::Unclassified(Unclassified {
|
||||
detail: format!("收尾阶段给出的回合终态无法识别:{session_status}"),
|
||||
});
|
||||
match error.classify(history_root) {
|
||||
TurnErrorClassified::ShouldStop(payload) => TurnCompletion::Failed(payload),
|
||||
// `Unclassified` 恒为真失败;这一臂写全只是把"分类层不允许静默"补齐。
|
||||
TurnErrorClassified::ShouldContinue { .. } => {
|
||||
unreachable!("Unclassified 只可能是回合失败,不该分类成控制流")
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
/// 收尾阶段账本给出的**非失败** `status` → 终态;认不出的值返回 `None` 由 [`turn_terminal`] 按
|
||||
/// 失败兜底。
|
||||
///
|
||||
/// 这就是原 `SessionOutcome` 的全部内容——它只是"没有 `Failed` 的 [`TurnCompletion`]",并进来
|
||||
/// 少一个同义类型。
|
||||
fn session_completion(status: &str) -> Option<TurnCompletion> {
|
||||
match status {
|
||||
"completed" => Some(TurnCompletion::Completed),
|
||||
"interrupted" => Some(TurnCompletion::Interrupted),
|
||||
"aborted" => Some(TurnCompletion::Aborted),
|
||||
// `failed` 不走这里:它要么带 typed 错误、要么按报告文本兜底。
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::agent::{
|
||||
ModelCallKind, ThreadEvent, TransportClosed, TurnCompletedStatus, TurnError,
|
||||
TurnInterrupted, Unclassified,
|
||||
};
|
||||
use platform_llm::LlmError;
|
||||
|
||||
fn history_root() -> std::path::PathBuf {
|
||||
std::path::PathBuf::from("/tmp/direct-turn-completion-test")
|
||||
}
|
||||
|
||||
fn model_call(payload: &TurnFailure) -> (&ModelCallKind, &str) {
|
||||
match payload {
|
||||
TurnFailure::ModelCallFailed(payload) => (&payload.kind, &payload.detail),
|
||||
other => panic!("expected a model call failure, got {other:?}"),
|
||||
}
|
||||
}
|
||||
|
||||
/// 正常收场:不带载荷,变体就是收尾阶段推出来的那个。
|
||||
#[test]
|
||||
fn non_failure_terminals_keep_the_session_status() {
|
||||
assert_eq!(
|
||||
turn_terminal("completed", Ok("报告不重要"), None, &history_root()),
|
||||
TurnCompletion::Completed
|
||||
);
|
||||
assert_eq!(
|
||||
turn_terminal("interrupted", Ok("报告不重要"), None, &history_root()),
|
||||
TurnCompletion::Interrupted
|
||||
);
|
||||
assert_eq!(
|
||||
turn_terminal("aborted", Ok("报告不重要"), None, &history_root()),
|
||||
TurnCompletion::Aborted
|
||||
);
|
||||
}
|
||||
|
||||
/// 认不出的 `status` 不冒充正常收场:按未分类失败发出去,并把原文留在原因里。
|
||||
#[test]
|
||||
fn unknown_session_status_fails_closed() {
|
||||
let completion = turn_terminal("something-new", Ok("报告"), None, &history_root());
|
||||
match completion {
|
||||
TurnCompletion::Failed(TurnFailure::Unclassified(payload)) => {
|
||||
assert!(payload.detail.contains("something-new"));
|
||||
}
|
||||
other => panic!("unknown status must fail closed, got {other:?}"),
|
||||
}
|
||||
}
|
||||
|
||||
/// 拿得到错误:分类与原因都取自错误。
|
||||
#[test]
|
||||
fn collect_error_becomes_a_failure_terminal() {
|
||||
let error = TurnError::from_model_call(&LlmError::Transport(
|
||||
"DirectProject 收尾历史失败:写入 project.jsonl 失败".into(),
|
||||
));
|
||||
let completion = turn_terminal("completed", Err(error), None, &history_root());
|
||||
let TurnCompletion::Failed(failure) = completion else {
|
||||
panic!("transport error must fail the turn");
|
||||
};
|
||||
let (kind, detail) = model_call(&failure);
|
||||
assert_eq!(*kind, ModelCallKind::TransportBroken);
|
||||
assert!(detail.contains("project.jsonl"));
|
||||
}
|
||||
|
||||
/// 收尾阶段的账本读不出来(`session_status` 只能是 `failed`)时没有错误可用:用交付报告兜底,
|
||||
/// 但照样要带载荷发出去,不能让界面停在"已结束、没原因"。
|
||||
#[test]
|
||||
fn unreadable_session_ledger_still_reports_a_payload() {
|
||||
assert_eq!(
|
||||
turn_terminal("failed", Ok("报告"), None, &history_root()),
|
||||
TurnCompletion::Failed(TurnFailure::Unclassified(Unclassified {
|
||||
detail: "报告".into()
|
||||
}))
|
||||
);
|
||||
}
|
||||
|
||||
/// 宿主自己记下的失败排在最前面:它比交付报告更接近现场。
|
||||
#[test]
|
||||
fn host_recorded_failure_outranks_every_other_source() {
|
||||
let diagnostic = "Codex app-server 已退出;exitStatus=signal: 9 (SIGKILL);\
|
||||
stderrClass=nonempty;stderrBytes=1000";
|
||||
let host_failure = TurnError::TransportClosed(TransportClosed {
|
||||
diagnostic: diagnostic.to_string(),
|
||||
});
|
||||
let completion = turn_terminal(
|
||||
"interrupted",
|
||||
Ok("执行连接已结束,正在核对自有子进程与在途操作。"),
|
||||
Some(&host_failure),
|
||||
&history_root(),
|
||||
);
|
||||
let TurnCompletion::Failed(failure) = completion else {
|
||||
panic!("host fact must fail the turn");
|
||||
};
|
||||
match &failure {
|
||||
TurnFailure::TransportClosed(payload) => {
|
||||
assert!(payload.diagnostic.contains("SIGKILL"));
|
||||
assert!(!payload.diagnostic.contains("正在核对自有子进程"));
|
||||
}
|
||||
other => panic!("expected a transport-closed payload, got {other:?}"),
|
||||
}
|
||||
|
||||
// 即使同时拿到了错误,宿主亲眼看到的事实仍然是第一顺位。
|
||||
let error =
|
||||
TurnError::from_model_call(&LlmError::Transport("DirectProject 收尾历史失败".into()));
|
||||
let host_failure = TurnError::TurnInterrupted(TurnInterrupted {
|
||||
detail: "本轮模型执行被中断".into(),
|
||||
});
|
||||
let completion = turn_terminal(
|
||||
"interrupted",
|
||||
Err(error),
|
||||
Some(&host_failure),
|
||||
&history_root(),
|
||||
);
|
||||
let TurnCompletion::Failed(failure) = completion else {
|
||||
panic!("host fact must fail the turn");
|
||||
};
|
||||
match &failure {
|
||||
TurnFailure::TurnInterrupted(payload) => {
|
||||
assert!(payload.detail.contains("本轮模型执行被中断"))
|
||||
}
|
||||
other => panic!("expected a turn-interrupted payload, got {other:?}"),
|
||||
}
|
||||
}
|
||||
|
||||
/// 终态事件的形状:失败时同一个 `turn.completed` 带载荷,其余只带 `status`。
|
||||
#[test]
|
||||
fn terminal_event_carries_the_payload_and_the_opening_identity() {
|
||||
let error = TurnError::from_model_call(&LlmError::Upstream {
|
||||
status_code: 502,
|
||||
message: "上游 502".into(),
|
||||
});
|
||||
let failing = turn_terminal("interrupted", Err(error), None, &history_root());
|
||||
let event = failing.event(2_000, Some("direct-codex:turn-1:user"));
|
||||
match event.failure() {
|
||||
Some(TurnFailure::ModelCallFailed(payload)) => assert_eq!(
|
||||
payload.kind,
|
||||
ModelCallKind::UpstreamFailed {
|
||||
status_code: 502,
|
||||
native: None
|
||||
}
|
||||
),
|
||||
other => panic!("expected a model call failure payload, got {other:?}"),
|
||||
}
|
||||
assert_eq!(event.user_item_id(), Some("direct-codex:turn-1:user"));
|
||||
assert_eq!(event.at(), Some(2_000));
|
||||
|
||||
let quiet = turn_terminal("completed", Ok("本轮交付已完成"), None, &history_root());
|
||||
let event = quiet.event(3_000, None);
|
||||
assert!(event.failure().is_none());
|
||||
assert!(matches!(
|
||||
event,
|
||||
ThreadEvent::TurnCompleted { ref status, .. } if *status == TurnCompletedStatus::Completed
|
||||
));
|
||||
}
|
||||
|
||||
/// 兜底终态:说不出原因的那一种只给分类,不冒充真实原因。
|
||||
#[test]
|
||||
fn host_dropped_terminal_only_carries_the_classification() {
|
||||
assert_eq!(
|
||||
TurnCompletion::host_dropped(),
|
||||
TurnCompletion::Failed(TurnFailure::HostDropped)
|
||||
);
|
||||
}
|
||||
|
||||
/// 成功终态:没有失败载荷,`status` 就是 `completed`。
|
||||
///
|
||||
/// 供给没有"深层终态出口"的执行器(cc / Claude Code sidecar)用:它们整轮成功返回后,
|
||||
/// 线程仍被放行占用,必须由放行侧补写这条终态。
|
||||
#[test]
|
||||
fn completed_terminal_carries_no_failure_payload() {
|
||||
let event = TurnCompletion::Completed.event(1_700_000_000_000, Some("item-1"));
|
||||
assert!(event.failure().is_none());
|
||||
assert!(matches!(
|
||||
event,
|
||||
ThreadEvent::TurnCompleted { ref status, .. } if *status == TurnCompletedStatus::Completed
|
||||
));
|
||||
}
|
||||
}
|
||||
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
Reference in New Issue
Block a user