Merge remote-tracking branch 'origin/master' into feat/game-works-management
Project CI / AI game creator shell Rust crates (pull_request) Successful in 3m5s
Project CI / AI game creator shell Rust lane 2/2 (pull_request) Failing after 3m28s
Project CI / AI game creator shell Rust lane 1/2 (pull_request) Successful in 5m5s
Project CI / Frontend tests (pull_request) Successful in 2m43s
Project CI / Backend tests (pull_request) Successful in 6m12s
Project CI / AI game creator shell web tests (pull_request) Successful in 2m37s
Project CI / Repository checks (pull_request) Failing after 3m54s
Project CI / Native shell tests (pull_request) Successful in 5m59s

This commit is contained in:
2026-10-03 17:06:04 +08:00
130 changed files with 5151 additions and 4113 deletions
@@ -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
@@ -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"}],
@@ -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,
}
@@ -0,0 +1,479 @@
use serde::{Deserialize, Serialize};
use serde_json::Value;
use ts_rs::TS;
/// 一条文件变更。
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, 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 ThreadFileChange {
pub(crate) path: String,
/// `add` | `update` | `delete`
pub(crate) kind: String,
}
/// 聊天视图的输入条目:一条 Codex 原始条目的原样投影(脱敏与限长在事件进前端状态时统一做)。
///
/// `itemType` 就是 Codex 的原始类型,逐字透传;前端按它决定投影成消息、思考还是工具卡片。
/// 未识别的类型走 [`ThreadItem::Other`],Rust 不替前端决定它是否可见。
///
/// 条目上的 `at` 是只用于显示的毫秒时间戳:ts-rs 默认把 `u64` 映射成 `bigint`,
/// 而 Tauri 的 JSON 通道传过来的是 `number`,因此统一标 `#[ts(as = "f64")]` 对齐。
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, TS)]
#[serde(tag = "itemType", 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 ThreadItem {
#[serde(rename = "message")]
Message {
/// 归一身份:全链路只有这一个 id。
item_id: String,
/// 原始 role(`user` / `assistant` / `system` / …);显示与否由前端判断。
role: String,
text: String,
#[ts(as = "f64")]
at: u64,
},
#[serde(rename = "reasoning")]
Reasoning {
item_id: String,
text: String,
#[ts(as = "f64")]
at: u64,
},
/// 原始 response item 的工具调用:参数在 `arguments`,输出在后续的
/// [`ThreadItem::FunctionCallOutput`](两者共用归一身份)。
#[serde(rename = "function_call")]
FunctionCall {
item_id: String,
name: String,
arguments: String,
#[ts(as = "f64")]
at: u64,
},
#[serde(rename = "function_call_output")]
FunctionCallOutput {
item_id: String,
output: String,
#[ts(as = "f64")]
at: u64,
},
#[serde(rename = "commandExecution")]
CommandExecution {
item_id: String,
command: String,
#[serde(default)]
output: Option<String>,
/// app-server 原始状态:`inProgress` / `completed` / `failed` / `declined` / …
#[serde(default)]
status: Option<String>,
#[serde(default)]
#[ts(as = "Option<i32>")]
exit_code: Option<i64>,
#[ts(as = "f64")]
at: u64,
},
#[serde(rename = "fileChange")]
FileChange {
item_id: String,
changes: Vec<ThreadFileChange>,
#[ts(as = "f64")]
at: u64,
},
#[serde(rename = "mcpToolCall")]
McpToolCall {
item_id: String,
tool: String,
arguments: String,
#[serde(default)]
output: Option<String>,
#[serde(default)]
status: Option<String>,
#[ts(as = "f64")]
at: u64,
},
#[serde(rename = "webSearch")]
WebSearch {
item_id: String,
#[serde(default)]
query: Option<String>,
#[serde(default)]
output: Option<String>,
#[ts(as = "f64")]
at: u64,
},
#[serde(rename = "contextCompaction")]
ContextCompaction {
item_id: String,
#[ts(as = "f64")]
at: u64,
},
/// 未识别的 Codex item 类型:原样透传身份与类型,不投影正文。
#[serde(rename = "other")]
Other {
item_id: String,
raw_type: String,
#[ts(as = "f64")]
at: u64,
},
}
impl ThreadItem {
/// 归一身份:Thread Manager 用它登记与释放未完成条目,前端用它合并同一张卡片。
pub(crate) fn item_id(&self) -> &str {
match self {
Self::Message { item_id, .. }
| Self::Reasoning { item_id, .. }
| Self::FunctionCall { item_id, .. }
| Self::FunctionCallOutput { item_id, .. }
| Self::CommandExecution { item_id, .. }
| Self::FileChange { item_id, .. }
| Self::McpToolCall { item_id, .. }
| Self::WebSearch { item_id, .. }
| Self::ContextCompaction { item_id, .. }
| Self::Other { item_id, .. } => item_id,
}
}
/// 条目展示时间(毫秒)。只用于条目自身的展示,不能当工具的开始 / 完成边界;
/// 那两类边界用事件级 `at`。
pub(crate) fn at(&self) -> u64 {
match self {
Self::Message { at, .. }
| Self::Reasoning { at, .. }
| Self::FunctionCall { at, .. }
| Self::FunctionCallOutput { at, .. }
| Self::CommandExecution { at, .. }
| Self::FileChange { at, .. }
| Self::McpToolCall { at, .. }
| Self::WebSearch { at, .. }
| Self::ContextCompaction { at, .. }
| Self::Other { at, .. } => *at,
}
}
}
/// 增量正文属于哪类条目。
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) enum ThreadDeltaKind {
/// assistant 正文。
Message,
/// 思考正文。
Reasoning,
}
/// 审批 / 提问请求与解决:本轮只透传,不并入聊天状态。
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) enum ThreadRequestKind {
#[serde(rename = "approval.requested")]
ApprovalRequested,
#[serde(rename = "ask.requested")]
AskRequested,
#[serde(rename = "request.resolved")]
RequestResolved,
}
impl ThreadRequestKind {
/// 未解决的请求要留在 bootstrap 里,直到出现对应的解决事件。
pub(crate) fn is_request(&self) -> bool {
matches!(self, Self::ApprovalRequested | Self::AskRequested)
}
pub(crate) fn is_resolution(&self) -> bool {
!self.is_request()
}
}
/// 原始值转文本:字符串原样,其它 JSON 值序列化。脱敏不在这一层做。
fn value_text(value: &Value) -> Option<String> {
match value {
Value::Null => None,
Value::String(text) => (!text.trim().is_empty()).then(|| text.trim().to_string()),
other => serde_json::to_string_pretty(other).ok(),
}
}
fn item_text(item: &Value) -> Option<String> {
item.get("text")
.and_then(Value::as_str)
// 空 / 全空白的 `text` 要在**回落之前**判掉:否则 `Some("")` 会短路 `content` / `summary`
// 兜底,末尾那条非空过滤再把整条条目丢掉(message / reasoning 的正文就此消失)。
.filter(|text| !text.trim().is_empty())
.map(str::to_string)
.or_else(|| {
for key in ["content", "summary"] {
let Some(parts) = item.get(key).and_then(Value::as_array) else {
continue;
};
let joined = parts
.iter()
.filter_map(|part| part.get("text").and_then(Value::as_str))
.collect::<Vec<_>>()
.join("");
if !joined.trim().is_empty() {
return Some(joined);
}
}
None
})
}
fn item_at_ms(item: &Value, observed_at_ms: u64) -> u64 {
let from_metadata = item
.get("internal_chat_message_metadata_passthrough")
.and_then(|meta| meta.get("create_time"))
.and_then(Value::as_f64)
.map(|seconds| (seconds * 1000.0).clamp(0.0, u64::MAX as f64) as u64)
.unwrap_or_default();
if from_metadata > 0 {
return from_metadata;
}
for key in ["startedAtMs", "completedAtMs"] {
let value = item.get(key).and_then(Value::as_u64).unwrap_or_default();
if value > 0 {
return value;
}
}
observed_at_ms
}
/// 原生毫秒时间戳:0(协议里的"缺省")与非法值一样按缺失处理。
fn json_ms(container: &Value, key: &str) -> Option<u64> {
container
.get(key)
.and_then(Value::as_u64)
.filter(|value| *value > 0)
}
/// `item/started` / `item/completed` 的事件级阶段时间(毫秒)。
///
/// 字段位置按当前 app-server 协议:通知层带 `params.startedAtMs` / `params.completedAtMs`,
/// 条目自带时用条目里的同名毫秒字段。完成事件即使同时
/// 带着开始字段也只取**完成**时间;两者都没有、但 `durationMs` 有可靠起点时按
/// 起点 + 时长派生结束。都没有就用宿主处理该事件的钟——原生缺阶段时间时这是唯一诚实的值。
pub(crate) fn thread_item_event_at_ms(
params: &Value,
item: &Value,
completed: bool,
observed_at_ms: u64,
) -> u64 {
let started_ms = || json_ms(params, "startedAtMs").or_else(|| json_ms(item, "startedAtMs"));
if !completed {
return started_ms().unwrap_or(observed_at_ms);
}
if let Some(at) = json_ms(params, "completedAtMs").or_else(|| json_ms(item, "completedAtMs")) {
return at;
}
let duration_ms = json_ms(params, "durationMs").or_else(|| json_ms(item, "durationMs"));
match (started_ms(), duration_ms) {
(Some(started), Some(duration)) => started.saturating_add(duration),
_ => observed_at_ms,
}
}
/// `turn.completed` 的事件级阶段时间(毫秒)。
///
/// Turn 里的 `startedAt` / `completedAt` 是 Unix **秒**(协议 `format: int64`,字段名不带
/// `Ms` 的都是秒),而 `durationMs` 才是毫秒。秒级截断在这里是不能用的:它既撑不起前端
/// 0.1 秒粒度的展示(显示出来的小数位是假精度),也可能让"完成时刻"落进该轮用户消息所在的
/// 同一秒、落在用户真实发送时间之前,前端按"结束早于开始"判成无效边界,于是一轮新回合被
/// 整轮吞掉。因此这里不采用任何秒字段:
/// - 只有 `durationMs` 与**高精度起点**都可靠时才按 `起点 + 时长` 派生结束;
/// - 否则取宿主处理终态的钟,语义与条目侧"没有原生阶段时间就用宿主钟"完全一致。
///
/// `high_precision_started_at_ms` 是本轮开始时宿主记下的那个毫秒起点(即 `turn.started`
/// 事件写入的同一个值),不是从上游秒字段换算出来的,`None` 表示起点也不可证明。
pub(crate) fn thread_turn_completed_at_ms(
turn: &Value,
high_precision_started_at_ms: Option<u64>,
observed_at_ms: u64,
) -> u64 {
match (high_precision_started_at_ms, json_ms(turn, "durationMs")) {
(Some(started), Some(duration)) => started.saturating_add(duration),
_ => observed_at_ms,
}
}
/// 归一身份:工具条目用工具调用 id,其它条目用自己的 `id`;只产出这一个值。
pub(crate) fn thread_item_identity(item: &Value) -> Option<String> {
let call_id = item
.get("call_id")
.or_else(|| item.get("callId"))
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.map(str::to_string);
let id = item
.get("id")
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.map(str::to_string);
call_id.or(id)
}
fn item_changes(item: &Value) -> Vec<ThreadFileChange> {
item.get("changes")
.and_then(Value::as_array)
.map(|changes| {
changes
.iter()
.filter_map(|change| {
let path = change
.get("path")
.and_then(Value::as_str)
.map(str::trim)
.filter(|path| !path.is_empty())?;
Some(ThreadFileChange {
path: path.to_string(),
kind: change
.get("kind")
.and_then(Value::as_str)
.map(str::to_string)
.unwrap_or_else(|| "update".to_string()),
})
})
.collect::<Vec<_>>()
})
.unwrap_or_default()
}
fn field_text(item: &Value, key: &str) -> Option<String> {
item.get(key).and_then(value_text)
}
/// 把一条 Codex 原始条目投影成线上条目;拿不到身份或类型时返回 `None`。
///
/// `observed_at_ms` 只在条目自带时间缺失时兜底(运行态用当前时间,历史用文件记录时间)。
///
/// 这一层**原样透传**字段,不做脱敏、也不限长:脱敏在事件进前端状态时统一做(前端
/// `directThreadSanitize`)。Thread Manager 只归一身形与身份,不替前端决定什么能看。
pub(crate) fn thread_item_from_value(item: &Value, observed_at_ms: u64) -> Option<ThreadItem> {
if !item.is_object() {
return None;
}
let item_id = thread_item_identity(item)?;
let item_type = item
.get("type")
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())?;
let at = item_at_ms(item, observed_at_ms);
let text = item_text(item);
let role = item
.get("role")
.and_then(Value::as_str)
.map(str::trim)
.filter(|role| !role.is_empty())
.map(str::to_string);
Some(match item_type {
"message" | "agentMessage" | "userMessage" => ThreadItem::Message {
item_id,
role: role.unwrap_or_else(|| {
if item_type == "userMessage" {
"user".to_string()
} else {
"assistant".to_string()
}
}),
text: text?,
at,
},
"reasoning" => ThreadItem::Reasoning {
item_id,
text: text?,
at,
},
"function_call" => ThreadItem::FunctionCall {
item_id,
name: item
.get("name")
.and_then(Value::as_str)
.map(str::to_string)
.unwrap_or_default(),
arguments: field_text(item, "arguments").unwrap_or_default(),
at,
},
"function_call_output" => ThreadItem::FunctionCallOutput {
item_id,
output: field_text(item, "output").unwrap_or_default(),
at,
},
"commandExecution" => ThreadItem::CommandExecution {
item_id,
command: item
.get("command")
.and_then(Value::as_str)
.map(str::trim)
.filter(|command| !command.is_empty())
.map(str::to_string)
.unwrap_or_default(),
output: ["aggregatedOutput", "output", "error"]
.iter()
.find_map(|key| field_text(item, key)),
status: item
.get("status")
.and_then(Value::as_str)
.map(str::to_string),
exit_code: item.get("exitCode").and_then(Value::as_i64),
at,
},
"fileChange" => ThreadItem::FileChange {
item_id,
changes: item_changes(item),
at,
},
"mcpToolCall" => ThreadItem::McpToolCall {
item_id,
tool: item
.get("tool")
.and_then(Value::as_str)
.map(str::trim)
.filter(|tool| !tool.is_empty())
.map(str::to_string)
.unwrap_or_default(),
arguments: field_text(item, "arguments").unwrap_or_default(),
output: ["result", "error"]
.iter()
.find_map(|key| field_text(item, key)),
status: item
.get("status")
.and_then(Value::as_str)
.map(str::to_string),
at,
},
"webSearch" => ThreadItem::WebSearch {
item_id,
query: field_text(item, "query").or_else(|| {
item.get("action")
.and_then(|action| action.get("query"))
.and_then(value_text)
}),
output: field_text(item, "output"),
at,
},
"contextCompaction" => ThreadItem::ContextCompaction { item_id, at },
other => ThreadItem::Other {
item_id,
raw_type: other.to_string(),
at,
},
})
}
/// 历史切片投影:保持文件顺序,不做任何合并(同一调用的调用与输出是两条条目)。
///
/// `timestamp_of` 是文件记录时间,仅在条目自带时间缺失时兜底。
pub(crate) fn thread_items_from_history(
items: &[Value],
timestamp_of: impl Fn(&Value) -> u64,
) -> Vec<ThreadItem> {
items
.iter()
.filter_map(|item| thread_item_from_value(item, timestamp_of(item)))
.collect()
}
@@ -0,0 +1,31 @@
//! DirectProject 聊天事件的线上模型与投影。
//!
//! 前端消费的类型由 ts-rs 导出到 `src/view/project-development/chat/generated/`,
//! 与 Rust 定义同源:加一个字段不会只改一边。
//!
//! 本模块只做三件事:挑字段、脱敏、截断。工具卡片的 `kind`、标题、折叠摘要、可见性与
//! 合并规则全部属于前端投影,这里一概不出现。
//!
//! 条目身份在进队列前就归一成**一个** `itemId`:原始文件里工具条目带两个 id(app-server
//! 的调用 id 与 response item id,调用与输出共用前者),归一只在 Rust 边界做一次,
//! Thread Manager 与前端都只认这一个,不暴露第二个 id 概念。
//! 历史分页锚点是另一回事,那是文件里的原始 item id,单独取。
//!
//! 文件按职责拆开:
//! - [`items`]:条目与历史切片的投影、脱敏、截断;
//! - [`turn`]:事件、订阅 / 消费载荷与队列出队语义;
//! - [`failure`]:失败终态的线上载荷;
//! - [`clock`]:宿主观测钟。
mod clock;
mod failure;
mod items;
mod turn;
#[cfg(test)]
mod tests;
pub(crate) use clock::*;
pub(crate) use failure::*;
pub(crate) use items::*;
pub(crate) use turn::*;
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,395 @@
use super::failure::TurnFailure;
use super::items::{ThreadDeltaKind, ThreadItem, ThreadRequestKind};
use crate::agent::UserItem;
use serde::{Deserialize, Serialize};
use ts_rs::TS;
/// 一条待发消息离开队列的原因。
///
/// typed 枚举,取值即语义:取消是用户在输入盒上撤掉这条消息,放行是它已经离开队列并成为回合
/// (同一临界区里另有 `turn.started`)。界面按它分流,不解析字符串。
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) enum QueueRemovalReason {
/// 用户取消了这条待发消息。
Cancelled,
/// 放行:这条待发消息已经离开队列,成为正在跑的那一轮。
Dispatched,
}
/// 取消一条待发消息的结果。
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) enum QueueRemovalOutcome {
/// 已从队列移除。
Removed,
/// 这条消息已经被放行(正在跑的那一轮就是它),不能按待发消息取消。
AlreadyDispatched,
/// 队列里没有这个身份,也没有在跑的一轮是它。
NotFound,
}
/// `turn.completed` 的终态取值。
///
/// 正常收场是前三档、不带载荷;`Failed` **一定**带 [`TurnFailure`] 载荷,只由
/// [`ThreadEvent::turn_completed_failed`] 写出。
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize, TS)]
#[serde(rename_all = "lowercase")]
#[ts(export, export_to = concat!(env!("CARGO_MANIFEST_DIR"), "/../src/view/project-development/chat/generated/"))]
pub(crate) enum TurnCompletedStatus {
Completed,
Interrupted,
Aborted,
Failed,
}
/// 正常收场的三档(不含失败):只作 [`ThreadEvent::turn_completed`] 的构造参数,不进线上形状。
///
/// 失败必须带载荷,所以这里没有 `Failed`——这就是 [`ThreadEvent::turn_completed`] 写不出
/// "failed 但没有载荷"的原因。
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum TurnOutcome {
Completed,
Interrupted,
Aborted,
}
impl From<TurnOutcome> for TurnCompletedStatus {
fn from(outcome: TurnOutcome) -> Self {
match outcome {
TurnOutcome::Completed => Self::Completed,
TurnOutcome::Interrupted => Self::Interrupted,
TurnOutcome::Aborted => Self::Aborted,
}
}
}
/// Thread Manager 下发的运行态事件。
///
/// 顺序由数组顺序给出(同一个 subscriber 的 `consume` 按队列顺序返回),因此不需要 `seq`:
/// 游标是 Thread Manager 的内部事实,不下发。
///
/// 事件不带回合身份:DirectProject 同一时刻只有一个回合在跑,"当前回合是否还在跑"由
/// 生命周期事件在序列中的位置给出,`turn_id` 对前端没有任何额外信息。
///
/// 四种生命周期事件(`turn.started` / `turn.completed` / `item.started` / `item.completed`)
/// 额外带事件级 `at`:它是**该阶段本身**的发生时间(毫秒),不是条目展示时间。条目上的
/// `item.at` 只说明"这条条目什么时候被看到",工具计时不得拿它当开始或完成边界。
/// 条目阶段优先用通知层的毫秒字段(`startedAtMs` / `completedAtMs`),缺失才用宿主钟;
/// 回合阶段没有可用的毫秒上游字段(Turn 只有秒级 `startedAt` / `completedAt`),一律用宿主
/// 在该阶段取的毫秒钟——见 `thread_turn_completed_at_ms` 的说明。
/// `at` 在事件进入 Thread Manager 时就固定:重放(bootstrap / consume)必须沿用原值,
/// 不能在前端收到或重放时重新取当前时间。
///
/// `turn.started` / `turn.completed` 额外带可选的 `userItemId`:本轮开口用户条目的 **canonical
/// itemId**(与同轮那条用户条目事件同源,由宿主按 `clientTurnId` 现算,`direct-codex:{clientTurnId}:user`;
/// **不读盘回填**——开始事件发生在用户条目落盘之前,落盘本身也可能失败)。回合事件本身
/// 不带回合身份,这个字段只用来把"这一轮的边界属于哪条用户消息"讲清楚:前端在只有生命周期锚点
/// + 历史切片、运行态一直为空时也能按身份认领开口条目,不必靠时间戳猜。缺失表示身份不可证明
/// (旧事件、没有开口用户条目、取消时拿不到 clientTurnId),此时前端不得补造。
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, 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 ThreadEvent {
#[serde(rename = "turn.started")]
TurnStarted {
/// 本轮开始的阶段时间(毫秒):**放行**那一刻的宿主毫秒钟(逻辑回合的起点,不是
/// `turn/start` 的时刻)。
#[serde(default, skip_serializing_if = "Option::is_none")]
#[ts(optional, as = "Option<f64>")]
at: Option<u64>,
/// 本轮开口用户条目的 canonical itemId;缺失表示身份不可证明。
#[serde(default, skip_serializing_if = "Option::is_none")]
#[ts(optional, as = "Option<String>")]
user_item_id: Option<String>,
},
#[serde(rename = "turn.completed")]
TurnCompleted {
/// 终态语义:`completed` / `interrupted` / `aborted` 是正常收场,不带载荷;`failed` 是
/// **失败**,一定带 `failure` 载荷。正常收场的构造参数是 [`TurnOutcome`](不含 `Failed`),
/// [`ThreadEvent::turn_completed`] 因此写不出"failed 但没有载荷"。
status: TurnCompletedStatus,
/// 失败载荷:只有 `status == "failed"` 才有;失败原因只从这里下发一次。
#[serde(default, skip_serializing_if = "Option::is_none")]
#[ts(optional)]
failure: Option<TurnFailure>,
/// 本轮终态的阶段时间(毫秒):宿主写下终态的毫秒钟,或 `durationMs` + 高精度起点的派生值。
#[serde(default, skip_serializing_if = "Option::is_none")]
#[ts(optional, as = "Option<f64>")]
at: Option<u64>,
/// 本轮开口用户条目的 canonical itemId:与同一轮的 `turn.started` 同源;缺失表示不可证明。
#[serde(default, skip_serializing_if = "Option::is_none")]
#[ts(optional, as = "Option<String>")]
user_item_id: Option<String>,
},
#[serde(rename = "item.started")]
ItemStarted {
item: ThreadItem,
/// 条目开始执行的原生阶段时间(毫秒);缺失时是宿主观测到该阶段的时间。
#[serde(default, skip_serializing_if = "Option::is_none")]
#[ts(optional, as = "Option<f64>")]
at: Option<u64>,
},
#[serde(rename = "item.completed")]
ItemCompleted {
item: ThreadItem,
/// 条目结束的原生阶段时间(毫秒);缺失时是宿主观测到该阶段的时间。
#[serde(default, skip_serializing_if = "Option::is_none")]
#[ts(optional, as = "Option<f64>")]
at: Option<u64>,
},
#[serde(rename = "item.delta")]
ItemDelta {
item_id: String,
kind: ThreadDeltaKind,
delta: String,
},
#[serde(rename = "request")]
Request {
kind: ThreadRequestKind,
#[serde(default)]
request_id: Option<String>,
},
/// 待发消息入队:数组顺序就是队首到队尾的顺序。
///
/// 这条事件在条目仍在队期间**不可回收**,离开队列(取消或放行)时才转成可回收——新订阅者
/// 靠这一点在 bootstrap 里看到当前队列,`is_bootstrap_event` 不需要为它加特例。
///
/// 事件就是这条待发消息的**全部**事实:宿主不为它另存产物,prompt 与 canonical 形状都在放行时
/// 从这条条目重投影。
#[serde(rename = "queue.enqueued")]
QueueEnqueued {
/// 这条待发消息的回合身份;放行后同一轮的 `turn.started` / `turn.completed` 用它。
client_turn_id: String,
/// canonical 用户条目:前端据此派生 chip 文案,Rust 不渲染展示形状。
user_item: UserItem,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[ts(optional)]
creation_type: Option<String>,
/// 入队那一刻的宿主毫秒钟。
#[ts(as = "f64")]
at: u64,
},
/// 待发消息离开队列:`reason` 是取消还是放行。
#[serde(rename = "queue.removed")]
QueueRemoved {
client_turn_id: String,
reason: QueueRemovalReason,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[ts(optional, as = "Option<f64>")]
at: Option<u64>,
},
}
impl ThreadEvent {
pub(crate) fn turn_started(at: u64) -> Self {
Self::TurnStarted {
at: Some(at),
user_item_id: None,
}
}
pub(crate) fn turn_completed(status: TurnOutcome, at: u64) -> Self {
Self::TurnCompleted {
status: status.into(),
failure: None,
at: Some(at),
user_item_id: None,
}
}
/// 失败终态:`status` 固定 `failed`,原因必须随事件一起带出去。
pub(crate) fn turn_completed_failed(failure: TurnFailure, at: u64) -> Self {
Self::TurnCompleted {
status: TurnCompletedStatus::Failed,
failure: Some(failure),
at: Some(at),
user_item_id: None,
}
}
/// 失败载荷:只有失败终态有。
pub(crate) fn failure(&self) -> Option<&TurnFailure> {
match self {
Self::TurnCompleted { failure, .. } => failure.as_ref(),
_ => None,
}
}
/// 附上本轮开口用户条目的 canonical itemId。
///
/// 只在构造之后补一次身份,避免 `turn.started` / `turn.completed` 的既有调用点(含各处兜底
/// 终态)全部改签名。空串按缺失处理:宁可让前端隐藏未知用时,也不写一个假身份。
pub(crate) fn with_user_item_id(self, user_item_id: Option<&str>) -> Self {
let user_item_id = user_item_id
.map(str::trim)
.filter(|value| !value.is_empty())
.map(str::to_string);
match self {
Self::TurnStarted { at, .. } => Self::TurnStarted { at, user_item_id },
Self::TurnCompleted {
status,
failure,
at,
..
} => Self::TurnCompleted {
status,
failure,
at,
user_item_id,
},
other => other,
}
}
/// 本轮开口用户条目的 canonical itemId:只有生命周期事件有,其余返回 `None`。
///
/// 只读已存入事件的值,不在读取时重算——重放要用的就是原事件的身份。
#[cfg(test)]
pub(crate) fn user_item_id(&self) -> Option<&str> {
match self {
Self::TurnStarted { user_item_id, .. } | Self::TurnCompleted { user_item_id, .. } => {
user_item_id.as_deref()
}
_ => None,
}
}
pub(crate) fn item_started(item: ThreadItem, at: u64) -> Self {
Self::ItemStarted { item, at: Some(at) }
}
pub(crate) fn item_completed(item: ThreadItem, at: u64) -> Self {
Self::ItemCompleted { item, at: Some(at) }
}
pub(crate) fn item_delta(item_id: String, kind: ThreadDeltaKind, delta: String) -> Self {
Self::ItemDelta {
item_id,
kind,
delta,
}
}
pub(crate) fn request(kind: ThreadRequestKind, request_id: Option<String>) -> Self {
Self::Request { kind, request_id }
}
/// 待发消息入队事件。
pub(crate) fn queue_enqueued(
client_turn_id: String,
user_item: UserItem,
creation_type: Option<String>,
at: u64,
) -> Self {
Self::QueueEnqueued {
client_turn_id,
user_item,
creation_type,
at,
}
}
/// 待发消息离开队列事件。
pub(crate) fn queue_removed(
client_turn_id: String,
reason: QueueRemovalReason,
at: u64,
) -> Self {
Self::QueueRemoved {
client_turn_id,
reason,
at: Some(at),
}
}
/// 这条事件属于哪条待发消息:只有队列事件有。
pub(crate) fn queue_client_turn_id(&self) -> Option<&str> {
match self {
Self::QueueEnqueued { client_turn_id, .. }
| Self::QueueRemoved { client_turn_id, .. } => Some(client_turn_id),
_ => None,
}
}
/// 待发消息离开队列的原因:只有 `queue.removed` 有。
pub(crate) fn queue_removal_reason(&self) -> Option<QueueRemovalReason> {
match self {
Self::QueueRemoved { reason, .. } => Some(*reason),
_ => None,
}
}
/// 事件级阶段时间(毫秒):只有四种生命周期事件有,其余事件返回 `None`。
///
/// 只读已存入事件的值,不在读取时取钟——重放要用的就是原事件的时间。
#[cfg(test)]
pub(crate) fn at(&self) -> Option<u64> {
match self {
Self::TurnStarted { at, .. }
| Self::TurnCompleted { at, .. }
| Self::ItemStarted { at, .. }
| Self::ItemCompleted { at, .. } => *at,
Self::QueueEnqueued { at, .. } => Some(*at),
Self::QueueRemoved { at, .. } => *at,
Self::ItemDelta { .. } | Self::Request { .. } => None,
}
}
/// 事件关联的条目身份:只有 item 事件有。
pub(crate) fn item_id(&self) -> Option<&str> {
match self {
Self::ItemStarted { item, .. } | Self::ItemCompleted { item, .. } => {
Some(item.item_id())
}
_ => None,
}
}
pub(crate) fn request_id(&self) -> Option<&str> {
match self {
Self::Request { request_id, .. } => request_id.as_deref(),
_ => None,
}
}
pub(crate) fn request_kind(&self) -> Option<ThreadRequestKind> {
match self {
Self::Request { kind, .. } => Some(*kind),
_ => None,
}
}
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, 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 SubscriptionBootstrap {
pub(crate) subscription_id: String,
/// 首屏历史锚点:`project.jsonl` 里最后一条原始 item id。
#[serde(default)]
pub(crate) last_completed_item_id: Option<String>,
/// 该 subscriber 此刻应当处理的运行态事件(游标已经在队尾)。
pub(crate) events: Vec<ThreadEvent>,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, 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 ConsumeResult {
pub(crate) events: Vec<ThreadEvent>,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, 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 HistorySlice {
/// 脱敏条目,顺序即文件顺序;与运行态事件里的条目同形。
pub(crate) items: Vec<ThreadItem>,
pub(crate) has_more: bool,
/// 本次切片的原始 item id 锚点:无论切片里有没有可显示条目,分页都靠它向前。
#[serde(default)]
pub(crate) first_item_id: Option<String>,
}

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