Refactor/整理项目对话错误相关代码 #593
@@ -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?;
|
||||
@@ -996,7 +996,7 @@ pub(crate) async fn direct_game_creator_claude_code_chat_at(
|
||||
// 聊天区是按 `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(
|
||||
@@ -1054,7 +1054,7 @@ fn persist_direct_claude_assistant_reply_at(
|
||||
|
||||
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())?;
|
||||
@@ -1082,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));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1171,7 +1171,7 @@ mod tests {
|
||||
assert_eq!(text, "最终回复");
|
||||
assert!(matches!(
|
||||
observed.as_slice(),
|
||||
[DirectCodexTurnObservation::AgentMessageSegment(_)]
|
||||
[TurnObservation::AgentMessageSegment(_)]
|
||||
));
|
||||
}
|
||||
|
||||
|
||||
@@ -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));
|
||||
}
|
||||
|
||||
@@ -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"
|
||||
));
|
||||
}
|
||||
}
|
||||
@@ -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
@@ -0,0 +1,8 @@
|
||||
/// 宿主观测时刻:Unix 毫秒。
|
||||
pub(crate) fn now_ms() -> u64 {
|
||||
std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_millis()
|
||||
.min(u64::MAX as u128) as u64
|
||||
}
|
||||
@@ -0,0 +1,50 @@
|
||||
//! 失败终态的**线上载荷**。
|
||||
//!
|
||||
//! [`TurnFailure`] 是 `turn.completed.status == "failed"` 时必有的下发形状,形状属于线上协议,
|
||||
//! 所以留在这个 wire 模块里;`turn.completed` 事件本身在 [`super::turn`]。
|
||||
//!
|
||||
//! "什么算失败、原因怎么写"是宿主内部策略,不在 wire 里:那是
|
||||
//! [`crate::agent::TurnCompletion`](`thread_manager::turn_completion`)。
|
||||
//!
|
||||
//! 载荷与脱敏都由 [`TurnError::classify`] 一处投影:Rust 侧没有第二个地方再拼它,也没有任何
|
||||
//! 地方再解析它。这里不碰事件队列的搬运规则,也不自己认 `LlmError`。
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
use ts_rs::TS;
|
||||
|
||||
use crate::agent::{
|
||||
EnvironmentNotReady, HostStateUnavailable, ModelCallFailed, ProjectRootUnanchored,
|
||||
SuperErrorFromStringPlusStage, TimedOut, TransportClosed, TurnInterrupted, Unclassified,
|
||||
};
|
||||
|
||||
/// 失败终态的可下发载荷(`turn.completed.status == "failed"` 时必有,其余终态没有)。
|
||||
///
|
||||
/// 载荷直接携带**typed 变体**:没有 `kind` 粗分类、也没有预拼的 `message`。前端按变体选语气、
|
||||
/// 按变体拼文案;宿主原始事实(`detail` / `cause` / `diagnostic`)留在字段里,只用于分流与诊断、
|
||||
/// 不直接上屏。
|
||||
///
|
||||
/// 唯一投影点是 [`crate::agent::TurnError::classify`]:控制流(返修要求)落进
|
||||
/// [`crate::agent::TurnErrorClassified::ShouldContinue`],所以控制流既不会出现在这里,前端也
|
||||
/// 不需要为它写分支。
|
||||
// TODO badnaming: 这个名字没有表达出它只是"回合终态的失败载荷",先留着待改名。
|
||||
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, TS)]
|
||||
#[serde(
|
||||
tag = "type",
|
||||
rename_all = "camelCase",
|
||||
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 TurnFailure {
|
||||
ProjectRootUnanchored(ProjectRootUnanchored),
|
||||
EnvironmentNotReady(EnvironmentNotReady),
|
||||
HostStateUnavailable(HostStateUnavailable),
|
||||
ModelCallFailed(ModelCallFailed),
|
||||
TransportClosed(TransportClosed),
|
||||
TimedOut(TimedOut),
|
||||
TurnInterrupted(TurnInterrupted),
|
||||
SuperErrorFromStringPlusStage(SuperErrorFromStringPlusStage),
|
||||
Unclassified(Unclassified),
|
||||
/// 宿主任务提前结束(panic / 被取消):说不出原因的那一种兜底。
|
||||
HostDropped,
|
||||
}
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user