From 4e00a3b91fb9c8382bfc7f111209718437009c51 Mon Sep 17 00:00:00 2001 From: fylorn <249551762+fylorn@users.noreply.github.com> Date: Fri, 9 Oct 2026 17:30:59 +0800 Subject: [PATCH] Codex compaction on converted routes; mid-conversation developer messages stay in place Compaction. Codex compacts through any provider named "OpenAI" (which is what ThinkWatch's takeover configures) by sending its history plus a `compaction_trigger` item and requiring exactly one `compaction` output item. Converted routes refused that request, so Codex sessions on Claude, Gemini, Bedrock or Chat upstreams could never compact. - A `compaction_trigger` now becomes a handoff-summary request appended to the unchanged history (same tools, reasoning and tool_choice, so the previous turn's prompt cache still applies). - The upstream's text comes back as the single `compaction` item Codex expects, streamed or whole. Its `encrypted_content` is our carried format, `tw1.c.`; an answer without text fails instead of handing Codex an empty summary. While the summary is collected the stream reports `response.in_progress` so Codex's idle timeout doesn't fire. - A carried compaction item in a later request decodes to the summary, in place, as a system turn. OpenAI's own encrypted compactions are still refused. On passthrough to an OpenAI upstream, `strip_carried` turns a carried item into a developer message with the same text. Mid-conversation system messages. Responses `developer`/`system` items and Chat/Anthropic `system` messages after the conversation has started stay in place as `Role::System` instead of being appended to the system prompt, so the system prefix stays stable from one request to the next. Responses targets get a native developer message; Anthropic, Chat, Gemini and Bedrock get a user turn wrapped in ``, moved past a tool result if it sat between a call and its result. Leading ones still form the system prompt. The transcript and search readers show carried summaries and skip the gateway's summarisation instruction. Co-Authored-By: Claude Opus 5.5 --- crates/tw-dialect/src/anthropic/request.rs | 7 +- crates/tw-dialect/src/bedrock/request.rs | 3 + crates/tw-dialect/src/chat/request.rs | 7 +- crates/tw-dialect/src/compaction.rs | 160 ++++++ crates/tw-dialect/src/convert.rs | 31 +- crates/tw-dialect/src/gemini/request.rs | 3 + crates/tw-dialect/src/ir.rs | 153 ++++++ crates/tw-dialect/src/lib.rs | 4 +- crates/tw-dialect/src/responses/request.rs | 145 ++++- crates/tw-dialect/src/responses/response.rs | 43 +- crates/tw-dialect/src/responses/stream.rs | 91 ++- crates/tw-dialect/tests/caller.rs | 10 +- crates/tw-dialect/tests/codex_compaction.rs | 581 ++++++++++++++++++++ crates/tw-gateway/tests/conversion.rs | 109 ++++ crates/tw-gateway/tests/harness.rs | 40 +- crates/tw-store/src/search/text.rs | 14 +- crates/tw-store/src/transcript/mod.rs | 18 + crates/tw-store/src/transcript/read.rs | 13 +- 18 files changed, 1365 insertions(+), 67 deletions(-) create mode 100644 crates/tw-dialect/src/compaction.rs create mode 100644 crates/tw-dialect/tests/codex_compaction.rs diff --git a/crates/tw-dialect/src/anthropic/request.rs b/crates/tw-dialect/src/anthropic/request.rs index be1af89a..0651656b 100644 --- a/crates/tw-dialect/src/anthropic/request.rs +++ b/crates/tw-dialect/src/anthropic/request.rs @@ -48,14 +48,14 @@ pub fn decode_request(v: &Value, dropped: &mut Dropped) -> Result Role::Assistant, - // 消息里的 system 角色:并进系统提示,位置信息丢失但内容保留。DeepSeek + // 消息里的 system 角色:开头的并进系统提示,对话中途的留在原位。DeepSeek // Harness 在这里放改过的系统提示,还有对话中途增删工具的 `tool_addition` / // `tool_removal`:别家没有这种写法 Some("system") => { let content = m.get("content").unwrap_or(&Value::Null); let t = text_of(content); if !t.is_empty() { - r.system.push(t); + system_turn(&mut r, t); } for b in content.as_array().into_iter().flatten() { match str_of(b, "type") { @@ -336,6 +336,9 @@ fn result_content(c: Option<&Value>, dropped: &mut Dropped) -> Vec { /// 中间表示 → 发给 Anthropic 上游的请求。 pub fn encode_request(r: &Request, t: &Target, dropped: &mut Dropped) -> Value { + // 对话中途的系统消息写成带标记的用户消息(见 `fold_system_turns`) + let folded = fold_system_turns(r); + let r = folded.as_ref(); let mut out = Map::new(); out.insert("model".into(), json!(r.model)); let max_tokens = r.max_tokens.unwrap_or(t.default_max_tokens); diff --git a/crates/tw-dialect/src/bedrock/request.rs b/crates/tw-dialect/src/bedrock/request.rs index 4ed91fd1..3193ca6b 100644 --- a/crates/tw-dialect/src/bedrock/request.rs +++ b/crates/tw-dialect/src/bedrock/request.rs @@ -89,6 +89,9 @@ fn cache_point(ttl: CacheTtl) -> Value { } pub fn encode_request(r: &Request, t: &Target, dropped: &mut Dropped) -> Value { + // 对话中途的系统消息写成带标记的用户消息(见 `fold_system_turns`) + let folded = fold_system_turns(r); + let r = folded.as_ref(); let mut out = Map::new(); // 断点放不放:模型不认就一个都不放,报出来 let cache: &[CachePoint] = if caches(&r.model) { &r.cache } else { &[] }; diff --git a/crates/tw-dialect/src/chat/request.rs b/crates/tw-dialect/src/chat/request.rs index ad52a725..bcd8e5b3 100644 --- a/crates/tw-dialect/src/chat/request.rs +++ b/crates/tw-dialect/src/chat/request.rs @@ -40,7 +40,7 @@ pub fn decode_request( "system" | "developer" => { let t = text_of(content); if !t.is_empty() { - r.system.push(t); + system_turn(&mut r, t); } } "user" => r.messages.push(Message { @@ -298,6 +298,9 @@ fn assistant_parts(m: &Value, dropped: &mut Dropped) -> Vec { /// 中间表示 → 发给 Chat 上游的请求。 pub fn encode_request(r: &Request, t: &Target, dropped: &mut Dropped) -> Value { + // 对话中途的系统消息写成带标记的用户消息(见 `fold_system_turns`) + let folded = fold_system_turns(r); + let r = folded.as_ref(); let mut out = Map::new(); out.insert("model".into(), json!(r.model)); @@ -307,7 +310,7 @@ pub fn encode_request(r: &Request, t: &Target, dropped: &mut Dropped) -> Value { } for m in &r.messages { match m.role { - Role::User => user_messages(m, dropped, &mut messages), + Role::User | Role::System => user_messages(m, dropped, &mut messages), Role::Assistant => { if let Some(a) = assistant_message(m, dropped) { messages.push(a); diff --git a/crates/tw-dialect/src/compaction.rs b/crates/tw-dialect/src/compaction.rs new file mode 100644 index 00000000..ac2cf8c1 --- /dev/null +++ b/crates/tw-dialect/src/compaction.rs @@ -0,0 +1,160 @@ +//! Codex 的远程压缩,在转给别家格式的上游时怎么做。 +//! +//! # Codex 那边的约定 +//! +//! 上游是 OpenAI(Codex 里叫 `OpenAI` 的 provider,ThinkWatch 接管时就是这个名字)时,Codex +//! 压缩前文走「远程压缩」(`codex-rs/core/src/compact_remote_v2.rs`): +//! +//! - 发一个普通的流式 Responses 请求,`input` 是整段历史,末尾加一项 +//! `{"type": "compaction_trigger"}`,工具和平时一样 +//! - 收流:`output_item.done` 里**必须恰好有一个** `{"type": "compaction", "encrypted_content": …}`, +//! 还要有 `response.completed`;别的项不看。少了、多了都算失败 +//! - 之后的历史换成:保留下来的用户消息 + 这个 `compaction` 项 + 重新写的上下文。下一轮 +//! 请求把这个项原样带回来,OpenAI 读它里面的密文接上前文 +//! +//! # 转给别家时 +//! +//! 别家写不出 OpenAI 的密文,也读不懂。所以: +//! +//! - 请求:历史原样,末尾加一条用户消息,请上游写一份交接摘要([`INSTRUCTION`])。工具、 +//! 推理设置、工具选择都不动 —— 和上一轮请求同一个开头,提示缓存照样命中(Anthropic 改了 +//! `tool_choice` 整段对话的缓存就作废)。这条是网关写的,不在原文里,内容过滤不看它 +//! - 回答:上游写的文字收齐,作为**唯一一个** `compaction` 项交回去,`encrypted_content` 是 +//! [`carry`] 写的 `tw1.c.<摘要的 base64url>`。Codex 只把它当一串不透明的字符存着、带回来 +//! - 下一轮:认出这个前缀,解回摘要,放在原位:对话中途的一条系统消息 +//! ([`crate::ir::Role::System`],开头是 [`SUMMARY_PREFIX`])。摘要是上游写的,不是调用方 +//! 说的话,内容过滤和脱敏不看它。OpenAI 自己的密文照旧拒绝:那段前文只有 OpenAI 读得出来 +//! - 这段对话之后换到 OpenAI 的上游直通时,`convert::strip_carried` 把它原位换成同样内容的 +//! 一条 developer 消息 —— OpenAI 不认我们写的「密文」 +//! +//! 写摘要的这一次和别的请求一样计用量和费用:上游的回答按它自己的格式嗅用量。 + +/// 转换写出去的压缩项的前缀:`tw1.` 是转换写出去的东西(见 [`crate::ir::CARRIED`]), +/// `c.` 是压缩 +pub const CARRIED_PREFIX: &str = "tw1.c."; + +/// 请上游写摘要的那条用户消息。 +/// +/// 摘要之后只剩它、用户最近说的话和系统提示,接手的模型要能只凭这些接着干:所以要的是 +/// 交接,不是概括 —— 文件路径、做过的决定、试过不行的路、跑着的东西、没做完的事都要在。 +pub const INSTRUCTION: &str = "\ +The conversation above is about to be compacted: everything before this message will be replaced \ +by a summary, and only the summary, the user's recent messages and the system instructions will \ +remain. Write that summary now, as a handoff that lets the work continue from it alone. + +Do not call any tools. Reply with the summary only, in the language the user has been writing in. + +Include: +- The user's goal, and every explicit request, preference and constraint they stated; quote them \ +where the exact wording matters. +- What has been done: decisions made and why, approaches tried and how they turned out, including \ +failures and their causes. +- The current state: files created, changed or examined (with their paths), commands run and \ +their important results, tests and whether they pass, and anything left half-done. +- State that still matters: running processes or sessions, plans or to-do items and their status, \ +facts learned about the environment (paths, versions, configuration). +- What remains: the next concrete steps, and any open questions or blockers. +- If an earlier summary appears above, carry forward everything in it that is still relevant. + +Be specific: keep exact names, paths, identifiers, error messages and numbers. Leave out anything \ +that no longer matters."; + +/// 下一轮把摘要放回对话时,放在它前面的那句话 +pub const SUMMARY_PREFIX: &str = "\ +The earlier part of this conversation was compacted. This summary of it replaces those messages, \ +which are no longer available:"; + +/// 摘要 → `compaction` 项的 `encrypted_content` +pub fn carry(summary: &str) -> String { + format!("{CARRIED_PREFIX}{}", base64url(summary.as_bytes())) +} + +/// `encrypted_content` 是不是转换写出去的 +pub fn is_carried(encrypted: &str) -> bool { + encrypted.starts_with(CARRIED_PREFIX) +} + +/// 读回转换写出去的摘要。不是这个前缀、或者解不开的是 `None` +pub fn read(encrypted: &str) -> Option { + let bytes = unbase64url(encrypted.strip_prefix(CARRIED_PREFIX)?)?; + String::from_utf8(bytes).ok() +} + +/// 摘要放回对话时的正文 +pub fn restored(summary: &str) -> String { + format!("{SUMMARY_PREFIX}\n\n{summary}") +} + +const ALPHABET: &[u8; 64] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-_"; + +/// base64url,不补 `=`。这个 crate 只依赖 serde,为两个函数不值得多拉一个依赖 +fn base64url(bytes: &[u8]) -> String { + let mut out = String::with_capacity(bytes.len().div_ceil(3) * 4); + for chunk in bytes.chunks(3) { + let n = chunk + .iter() + .enumerate() + .fold(0u32, |n, (i, b)| n | u32::from(*b) << (16 - 8 * i)); + for i in 0..=chunk.len() { + out.push(ALPHABET[((n >> (18 - 6 * i)) & 63) as usize] as char); + } + } + out +} + +fn unbase64url(s: &str) -> Option> { + let digit = |c: u8| ALPHABET.iter().position(|a| *a == c).map(|i| i as u32); + let mut out = Vec::with_capacity(s.len() / 4 * 3 + 2); + for chunk in s.as_bytes().chunks(4) { + // 一个字符凑不出一个字节 + if chunk.len() == 1 { + return None; + } + let mut n = 0u32; + for (i, c) in chunk.iter().enumerate() { + n |= digit(*c)? << (18 - 6 * i); + } + for i in 0..chunk.len() - 1 { + out.push((n >> (16 - 8 * i)) as u8); + } + } + Some(out) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn a_summary_goes_out_opaque_and_comes_back_unchanged() { + for s in [ + "", + "a", + "ab", + "abc", + "Fix the test in src/lib.rs — 还差一步 🚀", + ] { + let c = carry(s); + assert!(is_carried(&c), "{c}"); + assert!( + c[CARRIED_PREFIX.len()..] + .bytes() + .all(|b| ALPHABET.contains(&b)), + "{c}" + ); + assert_eq!(read(&c).as_deref(), Some(s)); + } + assert_eq!(carry("hi?"), "tw1.c.aGk_"); + } + + #[test] + fn something_else_is_not_read_as_a_summary() { + // OpenAI 自己的密文 + assert_eq!(read("gAAAAABo"), None); + // 前缀对、内容坏了 + assert_eq!(read("tw1.c.a"), None); + assert_eq!(read("tw1.c.a+b="), None); + // 不是 UTF-8 + assert_eq!(read(&format!("{CARRIED_PREFIX}_w")), None); + } +} diff --git a/crates/tw-dialect/src/convert.rs b/crates/tw-dialect/src/convert.rs index c9acce47..1d4295d0 100644 --- a/crates/tw-dialect/src/convert.rs +++ b/crates/tw-dialect/src/convert.rs @@ -251,6 +251,11 @@ impl Session { self.shape.tool_search.as_deref() == Some(name) } + /// 这是一次压缩:上游写的文字作为一个 `compaction` 项交回(见 [`crate::compaction`]) + pub fn is_compaction(&self) -> bool { + self.shape.compaction + } + /// 上游的整包响应 → 客户端的整包响应。上游返回的不是 JSON 时是 `None` pub fn response(&self, body: &[u8]) -> Option> { let v: Value = serde_json::from_slice(body).ok()?; @@ -373,6 +378,7 @@ impl Session { include_usage: false, gemini_sse: true, tool_search: None, + compaction: false, }, } } @@ -851,7 +857,8 @@ impl Collector { // ───────────────────────────────────────────────────────── 直通时的清理 -/// 直通请求里去掉转换写出去的推理签名。**没有需要去掉的就返回 `None`,请求一个字节都不改。** +/// 直通请求里去掉转换写出去的推理签名,转换写出去的压缩项换回摘要。**没有需要改的就返回 +/// `None`,请求一个字节都不改。** /// /// 客户端把转换写出去的推理内容原样带回来(这正是签名存在的意义)。之后如果这段对话 /// 换到了和客户端同格式的上游(故障转移、改了路由),请求会直通过去,而那些带 @@ -908,6 +915,28 @@ pub fn strip_carried(client: Dialect, body: &[u8]) -> Option> { && carried(i.get("encrypted_content"))) }); changed = input.len() != before; + // 转换写出去的压缩项:OpenAI 读不了我们写的「密文」,可那是前文仅剩的东西, + // 不能扔 —— 原位换成一条写着摘要的 developer 消息,和转给别家时一样 + for i in input.iter_mut() { + let summary = matches!( + i.get("type").and_then(Value::as_str), + Some("compaction" | "context_compaction") + ) + .then(|| i.get("encrypted_content").and_then(Value::as_str)) + .flatten() + .and_then(crate::compaction::read); + if let Some(summary) = summary { + *i = serde_json::json!({ + "type": "message", + "role": "developer", + "content": [{ + "type": "input_text", + "text": crate::compaction::restored(&summary), + }], + }); + changed = true; + } + } } Dialect::Gemini => { let contents = v.get_mut("contents")?.as_array_mut()?; diff --git a/crates/tw-dialect/src/gemini/request.rs b/crates/tw-dialect/src/gemini/request.rs index d1ac9cb3..06c77e8e 100644 --- a/crates/tw-dialect/src/gemini/request.rs +++ b/crates/tw-dialect/src/gemini/request.rs @@ -303,6 +303,9 @@ fn openapi_to_json_schema(v: &Value) -> Value { /// 中间表示 → 发给 Gemini 上游的请求体。路径由调用方按模型和是否流式拼。 pub fn encode_request(r: &Request, _t: &Target, dropped: &mut Dropped) -> Value { + // 对话中途的系统消息写成带标记的用户消息(见 `fold_system_turns`) + let folded = fold_system_turns(r); + let r = folded.as_ref(); let mut out = Map::new(); if !r.system.is_empty() { out.insert( diff --git a/crates/tw-dialect/src/ir.rs b/crates/tw-dialect/src/ir.rs index 9499b634..be29046b 100644 --- a/crates/tw-dialect/src/ir.rs +++ b/crates/tw-dialect/src/ir.rs @@ -14,6 +14,7 @@ //! `stop`),编码时记。两处记的都是**客户端请求里的字段路径**:用户对照的是自己 //! 发出去的请求,不是我们转成的那一份。 +use std::borrow::Cow; use std::sync::atomic::{AtomicU64, Ordering}; use serde_json::Value; @@ -152,12 +153,21 @@ pub struct ClientShape { /// Responses 客户端自己执行的工具搜索(Codex 的 `tool_search`,`execution: client`) /// 转成的函数工具叫什么。上游调用它时,写回去的是 `tool_search_call`,不是函数调用 pub tool_search: Option, + /// Responses 客户端要的是一次压缩(Codex 的 `compaction_trigger`):上游写的摘要要作为 + /// 一个 `compaction` 项交回去(见 [`crate::compaction`]) + pub compaction: bool, } #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum Role { User, Assistant, + /// 对话中途的系统消息:Responses 的 `developer`、Chat 和 Anthropic 消息里的 `system`。 + /// + /// **开头连着的那几条进 [`Request::system`],之后的留在原位。**都并进系统提示的话, + /// 对话里每多一条这样的消息,系统提示就变一次 —— 从系统提示算起的提示缓存跟着全部 + /// 作废。没有这种写法的格式写成一条带标记的用户消息([`fold_system_turns`]) + System, } #[derive(Debug, Clone, PartialEq)] @@ -811,6 +821,75 @@ pub(crate) fn text_of(v: &Value) -> String { } } +/// 解码时读到一条系统消息:对话还没开始就是系统提示的一段,开始了就是对话里的一条 +/// ([`Role::System`]) +pub fn system_turn(r: &mut Request, text: String) { + if r.messages.is_empty() { + r.system.push(text); + } else { + r.messages.push(Message { + role: Role::System, + parts: vec![Part::Text(text)], + }); + } +} + +/// 对话中途的系统消息在没有这种写法的格式里怎么写:包在 `` 里的一段 +/// 用户的话。 +/// +/// Claude Code 就是这么在对话里插系统消息的,Claude 认得这是系统说的、不是用户说的; +/// 别家的模型看到成对的标签也分得清 +pub fn system_reminder(text: &str) -> String { + format!("\n{text}\n") +} + +/// 对话中途的系统消息换成用户消息(正文见 [`system_reminder`])。**消息的条数和位置不变**, +/// 缓存断点按位置记的,照样对得上。 +/// +/// Anthropic、Gemini、Bedrock 的消息里没有系统角色;Chat 有,可很多别家模型的对话模板只 +/// 认开头那一条系统消息,后面的直接报错,所以 Chat 也这么写。没有这种消息的请求原样借出去。 +/// +/// **夹在工具调用和它的结果中间的,挪到结果后面**(并进带着结果的那条用户消息的末尾,原位 +/// 留一条空消息):Chat 要求调用之后紧跟着结果,Anthropic 和 Bedrock 要求结果排在用户消息 +/// 最前面 —— 留在原位就是一个 400 +pub fn fold_system_turns(r: &Request) -> Cow<'_, Request> { + if !r.messages.iter().any(|m| m.role == Role::System) { + return Cow::Borrowed(r); + } + let has = |m: &Message, f: fn(&Part) -> bool| m.parts.iter().any(f); + let mut out = r.clone(); + for i in 0..r.messages.len() { + if r.messages[i].role != Role::System { + continue; + } + let parts: Vec = std::mem::take(&mut out.messages[i].parts) + .into_iter() + .map(|p| match p { + Part::Text(t) => Part::Text(system_reminder(&t)), + other => other, + }) + .collect(); + out.messages[i].role = Role::User; + let prev = r.messages[..i] + .iter() + .rev() + .find(|m| m.role != Role::System); + let next = (i + 1..r.messages.len()).find(|&j| r.messages[j].role != Role::System); + match next { + Some(j) + if prev.is_some_and(|p| { + p.role == Role::Assistant && has(p, |x| matches!(x, Part::ToolCall(_))) + }) && r.messages[j].role == Role::User + && has(&r.messages[j], |x| matches!(x, Part::ToolResult(_))) => + { + out.messages[j].parts.extend(parts); + } + _ => out.messages[i].parts = parts, + } + } + Cow::Owned(out) +} + /// 同一角色的相邻消息并成一条:Anthropic 和 Gemini 都要求角色交替。 pub fn merge_roles(messages: Vec) -> Vec { let mut out: Vec = Vec::with_capacity(messages.len()); @@ -934,6 +1013,80 @@ mod tests { assert_eq!(out[0].parts.len(), 2); } + #[test] + fn a_system_turn_becomes_a_marked_user_turn_in_its_place() { + let text = |role, t: &str| Message { + role, + parts: vec![Part::Text(t.into())], + }; + let r = Request { + messages: vec![ + text(Role::User, "hi"), + text(Role::Assistant, "hello"), + text(Role::System, "be brief"), + text(Role::User, "and?"), + ], + ..Default::default() + }; + let f = fold_system_turns(&r); + assert_eq!(f.messages.len(), 4); + assert_eq!( + f.messages[2], + text( + Role::User, + "\nbe brief\n" + ) + ); + // 没有系统消息的不复制 + let plain = Request { + messages: vec![text(Role::User, "hi")], + ..Default::default() + }; + assert!(matches!(fold_system_turns(&plain), Cow::Borrowed(_))); + } + + #[test] + fn a_system_turn_between_a_call_and_its_result_moves_after_the_result() { + let call = Message { + role: Role::Assistant, + parts: vec![Part::ToolCall(ToolCall { + id: "c".into(), + name: "f".into(), + input: ToolInput::Json(serde_json::json!({})), + })], + }; + let result = Part::ToolResult(ToolResult { + id: "c".into(), + content: vec![Part::Text("ok".into())], + is_error: false, + }); + let r = Request { + messages: vec![ + call, + Message { + role: Role::System, + parts: vec![Part::Text("note".into())], + }, + Message { + role: Role::User, + parts: vec![result.clone()], + }, + ], + ..Default::default() + }; + let f = fold_system_turns(&r); + // 位置还在,空了;提示跟在结果后面 + assert_eq!(f.messages.len(), 3); + assert!(f.messages[1].parts.is_empty()); + assert_eq!( + f.messages[2].parts, + [ + result, + Part::Text("\nnote\n".into()) + ] + ); + } + #[test] fn usage_reported_in_pieces_merges_to_the_largest_of_each() { let mut u = Usage { diff --git a/crates/tw-dialect/src/lib.rs b/crates/tw-dialect/src/lib.rs index 5c1593e9..14c8b3a5 100644 --- a/crates/tw-dialect/src/lib.rs +++ b/crates/tw-dialect/src/lib.rs @@ -8,12 +8,14 @@ //! 同样两边都用的还有:从响应里旁路嗅出用量([`usage`],换算和转换共用各家的 //! `usage()`),拼上游地址([`url`]),去掉 DeepSeek Harness 只发给 DeepSeek 的 //! 扩展([`harness`]),在原文上找调用方的正文([`caller`],内容过滤读它、删它), -//! 以及读写各格式里名字不同的请求参数([`params`])。 +//! 读写各格式里名字不同的请求参数([`params`]),以及 Codex 的远程压缩转给别家时怎么做 +//! ([`compaction`])。 pub mod anthropic; pub mod bedrock; pub mod caller; pub mod chat; +pub mod compaction; pub mod convert; pub mod frame; pub mod gemini; diff --git a/crates/tw-dialect/src/responses/request.rs b/crates/tw-dialect/src/responses/request.rs index 227150d0..03b6e522 100644 --- a/crates/tw-dialect/src/responses/request.rs +++ b/crates/tw-dialect/src/responses/request.rs @@ -9,7 +9,7 @@ //! | 输入项 | 转成 | //! |---|---| //! | `additional_tools` | 里面的工具和顶层 `tools` 一样解码。Responses Lite 只在这里声明工具,顶层没有 `tools`;对话里可以有好几个(增量声明) | -//! | `message` | `system`、`developer` 并进系统提示,别的是一轮对话。`phase`(commentary、final_answer)别家没有,文字照留 | +//! | `message` | 开头连着的 `system`、`developer` 并进系统提示,对话中途的留在原位([`Role::System`]),别的是一轮对话。`phase`(commentary、final_answer)别家没有,文字照留 | //! | `agent_message` | 别的代理发来的话(发信人和任务名写在正文里),当用户的一轮。只有 OpenAI 读得懂的 `encrypted_content` 记丢弃 | //! | `reasoning` | 推理,签名照规矩带着 | //! | `function_call`、`custom_tool_call` 和各自的 `_output` | 工具调用和结果 | @@ -17,7 +17,9 @@ //! | `tool_search_call`、`tool_search_output` | 一次 `tool_search` 调用和结果。搜到的工具从此可以调用,加进工具列表 | //! | `configuration_update` | 对话中途改的推理强度。最后一个说了算,盖过顶层的 `reasoning.effort` —— Codex 为了保住提示缓存,顶层一直写开头那一档 | //! | `web_search_call`、`image_generation_call` | 丢弃并记下:OpenAI 服务端工具的执行记录,搜到的、画出的写在后面的回答里 | -//! | `compaction`、`context_compaction`、`compaction_trigger`、`item_reference` | 拒绝:内容在 OpenAI 服务端,或者是只有 OpenAI 读得懂的密文 | +//! | `compaction_trigger` | 要压缩前文:历史照转,末尾请上游写一份交接摘要,回答作为一个 `compaction` 项交回(见 [`crate::compaction`]) | +//! | `compaction`、`context_compaction` | 转换写出去的(`tw1.c.`)解回摘要,留在原位;OpenAI 自己的密文拒绝,只有 OpenAI 读得懂 | +//! | `item_reference` | 拒绝:内容在 OpenAI 服务端 | //! //! 顶层字段:`text.verbosity` 写给认它的模型([`Verbosity::understood_by`]),别处记丢弃; //! `service_tier` 别家没有同样的档位,记丢弃。`store`、`include`(转换写出的推理项总是带着 @@ -101,6 +103,7 @@ pub fn decode_request( shape, tools: ToolSet::default(), effort: None, + compaction: false, }; for t in arr_of(v, "tools") { cx.tool(t, None, "tools"); @@ -119,11 +122,21 @@ pub fn decode_request( } let Ctx { dropped, + shape, tools, effort, - .. + compaction, } = cx; r.tools = tools.tools; + // 要压缩前文:历史、工具、推理设置都不动(和上一轮同一个开头,提示缓存照样命中), + // 末尾请上游写摘要 + if compaction { + shape.compaction = true; + r.messages.push(Message { + role: Role::User, + parts: vec![Part::Text(crate::compaction::INSTRUCTION.to_string())], + }); + } // namespace 的说明(MCP 服务器的使用说明就写在这里):别家的工具没有 namespace, // 写进系统提示,模型照样看得到 for (ns, note) in tools.notes { @@ -341,6 +354,8 @@ struct Ctx<'a> { tools: ToolSet, /// 最后一个 `configuration_update` 里的推理强度 effort: Option, + /// 有 `compaction_trigger`:这一次要的是压缩 + compaction: bool, } impl Ctx<'_> { @@ -413,7 +428,7 @@ impl Ctx<'_> { "system" | "developer" => { let t = text_of(content); if !t.is_empty() { - r.system.push(t); + system_turn(r, t); } } role => { @@ -592,29 +607,33 @@ impl Ctx<'_> { .into(), )); } - "compaction" => { - return Err(Rejection( - "A compaction in input is an encrypted, compacted conversation only OpenAI can read, so the request cannot be converted for an upstream of another format." - .into(), - )); - } - "context_compaction" => { - if str_of(item, "encrypted_content").is_some_and(|e| !e.is_empty()) { - return Err(Rejection( - "A context_compaction in input is an encrypted, compacted conversation only OpenAI can read, so the request cannot be converted for an upstream of another format." - .into(), - )); + "compaction" | "context_compaction" => { + match str_of(item, "encrypted_content").filter(|e| !e.is_empty()) { + // 转换写出去的:解回摘要,留在原位。摘要是上游写的、不是调用方说的话, + // 和对话中途的系统消息一样放(内容过滤、脱敏只看调用方的话) + Some(enc) if crate::compaction::is_carried(enc) => { + let summary = crate::compaction::read(enc).ok_or_else(|| { + Rejection(format!( + "A {kind} in input carries a summary written during an earlier conversion, but it is damaged, so the request cannot be converted for an upstream of another format." + )) + })?; + push( + r, + Role::System, + Part::Text(crate::compaction::restored(&summary)), + ); + } + Some(_) => { + return Err(Rejection(format!( + "A {kind} in input is an encrypted, compacted conversation only OpenAI can read, so the request cannot be converted for an upstream of another format." + ))); + } + // 没有内容的只是一个记号 + None => self.dropped.path(format!("input.{kind}")), } - self.dropped.path("input.context_compaction"); - } - // 别家上游答不出 Codex 要的那个加密的 compaction 项:转过去只会白答一轮, - // Codex 再报「没有收到 compaction」 - "compaction_trigger" => { - return Err(Rejection( - "A compaction_trigger in input asks OpenAI's servers to compact the conversation into an encrypted item only OpenAI can read, so the request cannot be converted for an upstream of another format." - .into(), - )); } + // 要压缩前文:历史照转,最后由 `decode_request` 加上写摘要的请求 + "compaction_trigger" => self.compaction = true, // web_search_call、image_generation_call……:服务端工具的执行记录 other => self.dropped.path(format!("input.{other}")), } @@ -702,6 +721,24 @@ pub fn encode_request(r: &Request, _t: &Target, dropped: &mut Dropped) -> Value let mut input = Vec::new(); for m in &r.messages { match m.role { + // Responses 有对话中途的 developer 消息,原样放在原位 + Role::System => { + let content: Vec = m + .parts + .iter() + .filter_map(|p| match p { + Part::Text(t) if !t.is_empty() => { + Some(json!({ "type": "input_text", "text": t })) + } + _ => None, + }) + .collect(); + if !content.is_empty() { + input.push( + json!({ "type": "message", "role": "developer", "content": content }), + ); + } + } Role::User => { let mut content = Vec::new(); for p in &m.parts { @@ -1047,8 +1084,8 @@ mod tests { r#"{"model":"m","input":[{"type":"item_reference","id":"msg_1"}]}"#, r#"{"model":"m","input":[{"type":"compaction","encrypted_content":"x"}]}"#, r#"{"model":"m","input":[{"type":"context_compaction","encrypted_content":"x"}]}"#, - // Codex 的 Responses Lite 要压缩前文时发的:要的是一个只有 OpenAI 写得出的加密项 - r#"{"model":"m","input":[{"type":"message","role":"user","content":"hi"},{"type":"compaction_trigger"}]}"#, + // 前缀是转换写出去的,内容坏了 + r#"{"model":"m","input":[{"type":"compaction","encrypted_content":"tw1.c.a"}]}"#, r#"{"model":"m","background":true,"input":"hi"}"#, ] { let e = decode(body).unwrap_err(); @@ -1063,6 +1100,60 @@ mod tests { assert_eq!(dropped, ["input.context_compaction"]); } + #[test] + fn a_compaction_trigger_asks_the_upstream_for_a_summary_at_the_end() { + let (r, dropped, shape) = decode( + r#"{"model": "m", "input": [ + {"type": "message", "role": "developer", "content": "You are Codex."}, + {"type": "message", "role": "user", "content": "fix it"}, + {"type": "message", "role": "assistant", "content": "done"}, + {"type": "compaction_trigger"} + ]}"#, + ) + .unwrap(); + assert!(shape.compaction); + assert!(dropped.is_empty(), "{dropped:?}"); + assert_eq!(r.system, ["You are Codex."]); + let last = r.messages.last().unwrap(); + assert_eq!(last.role, Role::User); + assert_eq!( + last.parts, + [Part::Text(crate::compaction::INSTRUCTION.into())] + ); + // 没有它就是普通的一轮 + let (r, _, shape) = + decode(r#"{"model": "m", "input": [{"role": "user", "content": "hi"}]}"#).unwrap(); + assert!(!shape.compaction); + assert_eq!(r.messages.len(), 1); + } + + #[test] + fn a_summary_written_during_conversion_comes_back_in_its_place() { + let body = json!({"model": "m", "input": [ + {"type": "message", "role": "developer", "content": "You are Codex."}, + {"type": "message", "role": "user", "content": "fix it"}, + {"type": "compaction", "encrypted_content": crate::compaction::carry("Edited src/a.rs; tests pass.")}, + {"type": "message", "role": "developer", "content": ""}, + {"type": "message", "role": "user", "content": "now the docs"} + ]}) + .to_string(); + let (r, dropped, _) = decode(&body).unwrap(); + assert!(dropped.is_empty(), "{dropped:?}"); + assert_eq!(r.system, ["You are Codex."]); + let roles: Vec = r.messages.iter().map(|m| m.role).collect(); + assert_eq!(roles, [Role::User, Role::System, Role::System, Role::User]); + assert_eq!( + r.messages[1].parts, + [Part::Text(crate::compaction::restored( + "Edited src/a.rs; tests pass." + ))] + ); + assert_eq!( + r.messages[2].parts, + [Part::Text("".into())] + ); + } + /// Codex 的 Responses Lite(`use_responses_lite`):顶层没有 `tools`,所有工具装在 /// `input` 开头的 `additional_tools` 里,函数和自由格式工具在 `functions` 这个 /// namespace 里(`codex-rs/core/src/client.rs`、`codex-rs/tools/src/tool_spec.rs`) diff --git a/crates/tw-dialect/src/responses/response.rs b/crates/tw-dialect/src/responses/response.rs index c1f81f2b..fe500a2c 100644 --- a/crates/tw-dialect/src/responses/response.rs +++ b/crates/tw-dialect/src/responses/response.rs @@ -240,6 +240,20 @@ pub(crate) fn response_id(upstream: Option<&str>) -> String { } } +/// 压缩请求的回答(见 [`crate::compaction`]):一个 `compaction` 项,`encrypted_content` 是 +/// 转换写出去的摘要。`summary` 为 `None` 是刚开始的空项 +pub(crate) fn compaction_item(id: &str, summary: Option<&str>) -> Value { + json!({ + "id": id, + "type": "compaction", + "encrypted_content": summary.map(crate::compaction::carry).unwrap_or_default(), + }) +} + +/// 压缩请求的回答里没有文字时怎么说。Codex 收到失败会重试 +pub(crate) const NO_SUMMARY: &str = + "The upstream answered the compaction request without writing a summary."; + /// 中间表示 → 给 Responses 客户端的整包响应。 pub fn encode_response(r: &Response, s: &Session) -> Value { let (state, details) = status(r.stop.as_ref()); @@ -249,12 +263,31 @@ pub fn encode_response(r: &Response, s: &Session) -> Value { r.model.as_deref().unwrap_or(&s.model), state, ); - out["output"] = Value::Array( - r.blocks + out["output"] = if s.is_compaction() { + // 上游写的文字就是摘要;思考和(不该有的)工具调用都不算 + let texts: Vec<&str> = r + .blocks .iter() - .map(|b| item(b, &item_id(b, s), true, s)) - .collect(), - ); + .filter_map(|b| match b { + Block::Text(t) if !t.trim().is_empty() => Some(t.trim()), + _ => None, + }) + .collect(); + if texts.is_empty() { + out["status"] = json!("failed"); + out["error"] = json!({ "code": "server_error", "message": NO_SUMMARY }); + json!([]) + } else { + json!([compaction_item(&new_id("cmp_"), Some(&texts.join("\n\n")))]) + } + } else { + Value::Array( + r.blocks + .iter() + .map(|b| item(b, &item_id(b, s), true, s)) + .collect(), + ) + }; out["incomplete_details"] = details; if let Some(u) = &r.usage { out["usage"] = usage_json(u); diff --git a/crates/tw-dialect/src/responses/stream.rs b/crates/tw-dialect/src/responses/stream.rs index 46a4deac..6290a86c 100644 --- a/crates/tw-dialect/src/responses/stream.rs +++ b/crates/tw-dialect/src/responses/stream.rs @@ -10,8 +10,8 @@ use std::collections::HashMap; use serde_json::{Value, json}; use super::response::{ - envelope, incomplete, item, item_id, reasoning_signature, reasoning_text, response_id, status, - usage, usage_json, + NO_SUMMARY, compaction_item, envelope, incomplete, item, item_id, reasoning_signature, + reasoning_text, response_id, status, usage, usage_json, }; use crate::convert::Session; use crate::frame::{self, Frame}; @@ -342,8 +342,17 @@ pub struct Writer { order: Vec, usage: Option, stop: Option, + /// 压缩请求(见 [`crate::compaction`]):上游写的文字收在这里,结束时作为一个 + /// `compaction` 项交回。Codex 只认那一个项,中间的文字不往外写 + summary: Option, + /// 压缩请求:上一次报「还在写」之后又收到了几个增量 + quiet: u32, } +/// 压缩请求收着摘要不往外写,每收到这么多个增量报一次 `response.in_progress`:Codex 的流 +/// 五分钟没有事件就当断了,而一份长摘要加上前面的思考写得比这久 +const KEEPALIVE_DELTAS: u32 = 32; + impl Writer { pub fn new(s: &Session) -> Writer { Writer { @@ -358,7 +367,41 @@ impl Writer { order: Vec::new(), usage: None, stop: None, + summary: s.is_compaction().then(String::new), + quiet: 0, + } + } + + /// 压缩请求里的块和增量:文字收进摘要,别的只当作「还在写」。不是压缩请求、或者不是 + /// 块和增量的返回 false,照常写 + fn collect_summary(&mut self, e: &Event, out: &mut String) -> bool { + let Some(summary) = self.summary.as_mut() else { + return false; + }; + match e { + // 几段文字之间空一行 + Event::BlockStart { + kind: BlockKind::Text, + .. + } if !summary.trim_end().is_empty() => summary.push_str("\n\n"), + Event::BlockStart { .. } | Event::BlockStop { .. } => {} + Event::Delta { delta, .. } => { + if let Delta::Text(t) = delta { + summary.push_str(t); + } + self.quiet += 1; + if self.quiet >= KEEPALIVE_DELTAS { + self.quiet = 0; + self.start(out); + let r = envelope(&self.id, self.created, &self.model, "in_progress"); + self.emit("response.in_progress", json!({ "response": r }), out); + } + return true; + } + _ => return false, } + self.start(out); + true } fn emit(&mut self, kind: &str, mut body: Value, out: &mut String) { @@ -381,7 +424,7 @@ impl Writer { pub fn event(&mut self, e: &Event) -> String { let mut out = String::new(); - if self.finished { + if self.finished || self.collect_summary(e, &mut out) { return out; } match e { @@ -597,6 +640,10 @@ impl Writer { } self.start(&mut out); self.finished = true; + if let Some(summary) = self.summary.take() { + self.finish_compaction(summary.trim(), &mut out); + return out; + } for index in self.order.clone() { self.stop_item(index, &mut out); } @@ -623,6 +670,44 @@ impl Writer { } } +impl Writer { + /// 压缩请求的结尾:恰好一个 `compaction` 项,然后 `response.completed`。没写出摘要就是 + /// `response.failed` —— 交回一个空摘要,Codex 会当它压缩成功、把前文全扔掉 + fn finish_compaction(&mut self, summary: &str, out: &mut String) { + if summary.is_empty() { + let mut r = envelope(&self.id, self.created, &self.model, "failed"); + r["error"] = failure(500, NO_SUMMARY); + self.emit("response.failed", json!({ "response": r }), out); + return; + } + let id = new_id("cmp_"); + self.emit( + "response.output_item.added", + json!({ "output_index": 0, "item": compaction_item(&id, None) }), + out, + ); + let done = compaction_item(&id, Some(summary)); + self.emit( + "response.output_item.done", + json!({ "output_index": 0, "item": done }), + out, + ); + let (state, details) = status(self.stop.as_ref()); + let mut r = envelope(&self.id, self.created, &self.model, state); + r["incomplete_details"] = details; + r["output"] = json!([done]); + if let Some(u) = &self.usage { + r["usage"] = usage_json(u); + } + let kind = if state == "incomplete" { + "response.incomplete" + } else { + "response.completed" + }; + self.emit(kind, json!({ "response": r }), out); + } +} + /// 失败时 `response.error` 那一项。**Codex 按 `code` 决定退不退避**:429 说成 /// `server_error` 的话它会立刻重试,而该做的是等一等 fn failure(status: u16, message: &str) -> Value { diff --git a/crates/tw-dialect/tests/caller.rs b/crates/tw-dialect/tests/caller.rs index 1cff6408..3e8fa64c 100644 --- a/crates/tw-dialect/tests/caller.rs +++ b/crates/tw-dialect/tests/caller.rs @@ -58,14 +58,14 @@ fn decoded(d: Dialect, v: &Value) -> Marks { out } -/// 中间表示里系统提示和模型的话带着的记号 +/// 中间表示里系统提示(连同对话中途的系统消息)和模型的话带着的记号 fn not_callers(d: Dialect, v: &Value) -> BTreeSet { let r = decode(d, v, path(d), None).expect("decodes").request; let mut out = Marks::new(); for s in &r.system { marks_in(s, false, &mut out); } - for m in r.messages.iter().filter(|m| m.role == Role::Assistant) { + for m in r.messages.iter().filter(|m| m.role != Role::User) { for p in &m.parts { match p { Part::Text(t) => marks_in(t, false, &mut out), @@ -178,6 +178,8 @@ fn chat_completions() { {"type": "text", "text": "«t3»"}, ]}, {"role": "function", "content": "«x4»"}, + // 对话中途的系统消息留在原位,但不是调用方的话 + {"role": "system", "content": "«s3» mid-conversation"}, ], "tools": [{"type": "function", "function": {"name": "f", "description": "«x5»", "parameters": {}}}], }); @@ -201,6 +203,8 @@ fn responses() { ]}, {"content": "«u5» neither a type nor a role"}, {"type": "message", "role": "assistant", "content": [{"type": "output_text", "text": "«a1»"}]}, + // 对话中途的 developer 消息留在原位,但不是调用方的话 + {"type": "message", "role": "developer", "content": "«s5» mid-conversation"}, {"type": "function_call", "call_id": "c1", "name": "f", "arguments": "{\"q\":\"«x2»\"}"}, {"type": "function_call_output", "call_id": "c1", "output": "«t1»"}, {"type": "custom_tool_call", "call_id": "c2", "name": "g", "input": "«x3»"}, @@ -229,6 +233,8 @@ fn responses() { {"type": "encrypted_content", "encrypted_content": "«x11»"}, ]}, {"type": "configuration_update", "reasoning": {"effort": "high"}}, + // 转换时压缩出的摘要:上游写的,不是调用方的话 + {"type": "compaction", "encrypted_content": tw_dialect::compaction::carry("«s6» summary")}, ], "tools": [ {"type": "function", "name": "f", "description": "«x5»", "parameters": {}}, diff --git a/crates/tw-dialect/tests/codex_compaction.rs b/crates/tw-dialect/tests/codex_compaction.rs new file mode 100644 index 00000000..b31f4aef --- /dev/null +++ b/crates/tw-dialect/tests/codex_compaction.rs @@ -0,0 +1,581 @@ +//! Codex 压缩前文、对话中途的 developer 消息,转给别家格式的上游时。 +//! +//! **压缩**:上游是 OpenAI 时 Codex 走远程压缩(`codex-rs/core/src/compact_remote_v2.rs`): +//! 历史末尾加一项 `compaction_trigger` 发一个流式请求,收流时要求 `output_item.done` 里恰好 +//! 一个 `compaction` 项、还要有 `response.completed`;之后的请求把这个项原样带回来。这里照这个 +//! 约定走一个来回:上游写摘要,Codex 收到一个 `compaction` 项,下一轮带回来时上游读到的是摘要。 +//! +//! **对话中途的 developer 消息**:留在原位,系统提示不跟着变 —— 前后两个请求的开头一样, +//! 从系统提示算起的提示缓存才接得上。 +//! +//! 请求照 Codex 的 serde 类型写(`codex-rs/protocol/src/models.rs` 的 `ResponseItem`)。 + +// 整个请求写成一个 `json!`,嵌套得深 +#![recursion_limit = "512"] + +use serde_json::{Value, json}; +use tw_dialect::compaction; +use tw_dialect::convert::{Prepared, Session, decode, prepare, strip_carried}; +use tw_dialect::frame::Decoder; +use tw_dialect::ir::*; + +const UPSTREAMS: [Dialect; 4] = [ + Dialect::Anthropic, + Dialect::Chat, + Dialect::Gemini, + Dialect::Bedrock, +]; + +const SUMMARY: &str = "The user wants the failing parser test fixed.\n\n\ + - Ran `cargo test`: `parser::tests::nested` failed (unclosed bracket).\n\ + - Patched src/parser.rs to close nested brackets; tests now pass.\n\ + - Next: update docs/syntax.md to describe nesting."; + +fn tools() -> Value { + json!({"id": "at_1", "type": "additional_tools", "role": "developer", "tools": [ + {"type": "namespace", "name": "functions", "description": "", "tools": [ + {"type": "function", "name": "exec_command", "description": "Runs a command.", "strict": false, + "parameters": {"type": "object", "properties": {"cmd": {"type": "string"}}, "required": ["cmd"]}}, + {"type": "custom", "name": "apply_patch", "description": "Edit files.", + "format": {"type": "grammar", "syntax": "lark", "definition": "start: begin_patch hunk+ end_patch"}} + ]} + ]}) +} + +/// 一轮做完的对话:开头的工具声明和 developer 消息,中途 Codex 插了一条 developer 消息 +fn history() -> Vec { + vec![ + tools(), + json!({"id": "msg_b", "type": "message", "role": "developer", + "content": [{"type": "input_text", "text": "You are Codex, a coding agent."}]}), + json!({"type": "message", "role": "developer", + "content": [{"type": "input_text", "text": "sandbox: workspace-write"}]}), + json!({"type": "message", "role": "user", + "content": [{"type": "input_text", "text": "\n /repo\n"}]}), + json!({"type": "message", "role": "user", "content": [{"type": "input_text", "text": "Fix the failing parser test."}]}), + json!({"type": "message", "role": "assistant", "content": [{"type": "output_text", "text": "Running the tests."}], "phase": "commentary"}), + json!({"type": "function_call", "name": "exec_command", "namespace": "functions", "arguments": "{\"cmd\":\"cargo test\"}", "call_id": "call_1"}), + json!({"type": "function_call_output", "call_id": "call_1", "output": "parser::tests::nested FAILED"}), + json!({"type": "message", "role": "developer", + "content": [{"type": "input_text", "text": "Default"}]}), + json!({"type": "custom_tool_call", "call_id": "call_2", "name": "apply_patch", "namespace": "functions", + "input": "*** Begin Patch\n*** Update File: src/parser.rs\n*** End Patch\n"}), + json!({"type": "custom_tool_call_output", "call_id": "call_2", "output": "Success."}), + json!({"type": "message", "role": "assistant", "content": [{"type": "output_text", "text": "Fixed; the tests pass."}], "phase": "final_answer"}), + ] +} + +fn request(input: Vec) -> Value { + json!({ + "model": "gpt-5.4", + "stream": true, + "input": input, + "tool_choice": "auto", + "parallel_tool_calls": false, + "reasoning": {"effort": "medium", "summary": "auto", "context": "all_turns"}, + "store": false, + "include": ["reasoning.encrypted_content"], + "prompt_cache_key": "019a" + }) +} + +/// Codex 要压缩时发的:历史末尾加 `compaction_trigger` +fn compaction_request() -> Value { + let mut input = history(); + input.push(json!({"type": "compaction_trigger"})); + request(input) +} + +fn target(d: Dialect) -> Target { + Target { + dialect: d, + official: false, + default_max_tokens: 8192, + } +} + +fn model(d: Dialect) -> &'static str { + match d { + Dialect::Anthropic => "claude-opus-4-7", + Dialect::Chat => "deepseek-chat", + Dialect::Gemini => "gemini-2.5-pro", + _ => "us.anthropic.claude-sonnet-4-5-20250929-v1:0", + } +} + +/// 网关的做法:解码一次,改成上游的模型名,再按上游的格式编码 +fn prepared(body: &Value, upstream: Dialect) -> Prepared { + let mut d = decode(Dialect::Responses, body, "/v1/responses", None).unwrap(); + d.request.model = model(upstream).to_string(); + d.encode(&target(upstream)) +} + +fn body_of(p: &Prepared) -> Value { + serde_json::from_slice(&p.body).unwrap() +} + +/// 发给上游的请求里:系统提示、工具、对话 +fn parts_of(d: Dialect, v: &Value) -> (Value, Value, Vec) { + match d { + Dialect::Anthropic | Dialect::Bedrock => ( + v["system"].clone(), + v.get("tools") + .or_else(|| v.get("toolConfig")) + .cloned() + .unwrap(), + v["messages"].as_array().unwrap().clone(), + ), + Dialect::Chat => { + let messages = v["messages"].as_array().unwrap(); + assert_eq!(messages[0]["role"], "system"); + ( + messages[0].clone(), + v["tools"].clone(), + messages[1..].to_vec(), + ) + } + Dialect::Gemini => ( + v["systemInstruction"].clone(), + v["tools"].clone(), + v["contents"].as_array().unwrap().clone(), + ), + Dialect::Responses => unreachable!(), + } +} + +// ───────────────────────────────────────────────────────── 对话中途的 developer 消息 + +#[test] +fn a_developer_message_mid_conversation_leaves_the_prefix_alone() { + // 这一轮:历史 + 新的一句 + let mut first = history(); + first.push(json!({"type": "message", "role": "user", "content": [{"type": "input_text", "text": "Now the docs."}]})); + // 下一轮:上游答了一句,Codex 又插了一条 developer 消息,用户接着说 + let mut second = first.clone(); + second.extend([ + json!({"type": "message", "role": "assistant", "content": [{"type": "output_text", "text": "Which file?"}]}), + json!({"type": "message", "role": "developer", "content": [{"type": "input_text", "text": "approval: never"}]}), + json!({"type": "message", "role": "user", "content": [{"type": "input_text", "text": "docs/syntax.md"}]}), + ]); + for upstream in UPSTREAMS { + let a = body_of(&prepared(&request(first.clone()), upstream)); + let b = body_of(&prepared(&request(second.clone()), upstream)); + let (sys_a, tools_a, msgs_a) = parts_of(upstream, &a); + let (sys_b, tools_b, msgs_b) = parts_of(upstream, &b); + // 系统提示和工具一个字都没变:开头那几条 developer 消息在里面,中途的不在 + assert_eq!(sys_a, sys_b, "{upstream:?}"); + assert_eq!(tools_a, tools_b, "{upstream:?}"); + let sys = sys_b.to_string(); + assert!(sys.contains("You are Codex"), "{upstream:?}: {sys}"); + assert!( + sys.contains("sandbox: workspace-write"), + "{upstream:?}: {sys}" + ); + assert!(!sys.contains("collaboration_mode"), "{upstream:?}: {sys}"); + assert!(!sys.contains("approval: never"), "{upstream:?}: {sys}"); + // 前一个请求的对话是后一个的开头 + assert_eq!(msgs_a[..], msgs_b[..msgs_a.len()], "{upstream:?}"); + // 中途的 developer 消息在它原来的位置,标明是系统说的 + let later = Value::Array(msgs_b[msgs_a.len()..].to_vec()).to_string(); + assert!( + later.contains( + "\\napproval: never\\n" + ), + "{upstream:?}: {later}" + ); + let earlier = Value::Array(msgs_a.clone()).to_string(); + assert!( + earlier.contains("\\n"), + "{upstream:?}: {earlier}" + ); + } +} + +#[test] +fn a_mid_conversation_system_message_stays_a_developer_message_for_responses() { + // Chat 客户端对话中途的 system 消息,转给 Responses 的上游:那边有原生的写法 + let body = json!({"model": "gpt-5.4", "messages": [ + {"role": "system", "content": "Be brief."}, + {"role": "user", "content": "hi"}, + {"role": "assistant", "content": "hello"}, + {"role": "system", "content": "The user is on mobile now."}, + {"role": "user", "content": "and?"} + ]}) + .to_string(); + let p = prepare( + Dialect::Chat, + body.as_bytes(), + "/v1/chat/completions", + None, + &target(Dialect::Responses), + ) + .unwrap(); + let v = body_of(&p); + assert_eq!(v["instructions"], "Be brief."); + let roles: Vec<&str> = v["input"] + .as_array() + .unwrap() + .iter() + .map(|i| i["role"].as_str().unwrap()) + .collect(); + assert_eq!(roles, ["user", "assistant", "developer", "user"]); + assert_eq!( + v["input"][2]["content"][0]["text"], + "The user is on mobile now." + ); +} + +// ───────────────────────────────────────────────────────── 压缩 + +#[test] +fn a_compaction_request_asks_for_a_summary_on_the_same_prefix() { + // 压缩请求和同一段历史上的普通一轮:系统提示、工具、对话都一样,只有最后一条不同 + let mut normal = history(); + normal.push(json!({"type": "message", "role": "user", "content": [{"type": "input_text", "text": "Next?"}]})); + for upstream in UPSTREAMS { + let p = prepared(&compaction_request(), upstream); + assert!(p.session.is_compaction(), "{upstream:?}"); + assert!( + !p.dropped.iter().any(|d| d.contains("compaction")), + "{upstream:?}: {:?}", + p.dropped + ); + let c = body_of(&p); + let n = body_of(&prepared(&request(normal.clone()), upstream)); + let (sys_c, tools_c, msgs_c) = parts_of(upstream, &c); + let (sys_n, tools_n, msgs_n) = parts_of(upstream, &n); + assert_eq!(sys_c, sys_n, "{upstream:?}"); + assert_eq!(tools_c, tools_n, "{upstream:?}"); + assert_eq!( + msgs_c[..msgs_c.len() - 1], + msgs_n[..msgs_n.len() - 1], + "{upstream:?}" + ); + // 最后一条是请上游写摘要;工具选择没动(动了 Anthropic 整段对话的缓存就作废) + let last = msgs_c.last().unwrap().to_string(); + assert!( + last.contains("Write that summary now"), + "{upstream:?}: {last}" + ); + assert_eq!(c.get("tool_choice"), n.get("tool_choice"), "{upstream:?}"); + } +} + +fn named(event: &str, data: Value) -> String { + format!("event: {event}\ndata: {data}\n\n") +} + +fn data(v: Value) -> String { + format!("data: {v}\n\n") +} + +/// 上游写的摘要,按它的格式分很多片流回来。Anthropic 先想一想 +fn summary_stream(d: Dialect) -> String { + let pieces: Vec = SUMMARY.split_inclusive(' ').map(str::to_string).collect(); + match d { + Dialect::Anthropic => { + let mut s = vec![ + named( + "message_start", + json!({"type": "message_start", "message": {"id": "msg_up", "model": "claude-opus-4-7", + "usage": {"input_tokens": 1200, "cache_read_input_tokens": 9000, "output_tokens": 1}}}), + ), + named( + "content_block_start", + json!({"type": "content_block_start", "index": 0, "content_block": {"type": "thinking", "thinking": ""}}), + ), + named( + "content_block_delta", + json!({"type": "content_block_delta", "index": 0, "delta": {"type": "thinking_delta", "thinking": "Let me recall the work."}}), + ), + named( + "content_block_delta", + json!({"type": "content_block_delta", "index": 0, "delta": {"type": "signature_delta", "signature": "sig"}}), + ), + named( + "content_block_stop", + json!({"type": "content_block_stop", "index": 0}), + ), + named( + "content_block_start", + json!({"type": "content_block_start", "index": 1, "content_block": {"type": "text", "text": ""}}), + ), + ]; + for p in &pieces { + s.push(named("content_block_delta", json!({"type": "content_block_delta", "index": 1, "delta": {"type": "text_delta", "text": p}}))); + } + s.push(named( + "content_block_stop", + json!({"type": "content_block_stop", "index": 1}), + )); + s.push(named("message_delta", json!({"type": "message_delta", "delta": {"stop_reason": "end_turn"}, "usage": {"output_tokens": 300}}))); + s.push(named("message_stop", json!({"type": "message_stop"}))); + s.concat() + } + Dialect::Chat => { + let chunk = |delta: Value, finish: Value| { + data( + json!({"id": "c", "object": "chat.completion.chunk", "model": "deepseek-chat", + "choices": [{"index": 0, "delta": delta, "finish_reason": finish}]}), + ) + }; + let mut s = vec![chunk( + json!({"role": "assistant", "content": ""}), + Value::Null, + )]; + for p in &pieces { + s.push(chunk(json!({"content": p}), Value::Null)); + } + s.push(chunk(json!({}), json!("stop"))); + s.push(data(json!({"id": "c", "object": "chat.completion.chunk", "model": "deepseek-chat", "choices": [], + "usage": {"prompt_tokens": 10200, "completion_tokens": 300, "prompt_tokens_details": {"cached_tokens": 9000}}}))); + s.push("data: [DONE]\n\n".to_string()); + s.concat() + } + Dialect::Gemini => { + let mut s: Vec = pieces + .iter() + .map(|p| data(json!({"candidates": [{"content": {"role": "model", "parts": [{"text": p}]}}], "modelVersion": "gemini-2.5-pro"}))) + .collect(); + s.push(data(json!({"candidates": [{"content": {"role": "model", "parts": []}, "finishReason": "STOP"}], + "usageMetadata": {"promptTokenCount": 10200, "cachedContentTokenCount": 9000, "candidatesTokenCount": 300}}))); + s.concat() + } + Dialect::Bedrock => { + let mut s = vec![named("messageStart", json!({"role": "assistant"}))]; + for p in &pieces { + s.push(named( + "contentBlockDelta", + json!({"contentBlockIndex": 0, "delta": {"text": p}}), + )); + } + s.push(named("contentBlockStop", json!({"contentBlockIndex": 0}))); + s.push(named("messageStop", json!({"stopReason": "end_turn"}))); + s.push(named("metadata", json!({"usage": {"inputTokens": 1200, "cacheReadInputTokens": 9000, "outputTokens": 300, "totalTokens": 10500}}))); + s.concat() + } + Dialect::Responses => unreachable!(), + } +} + +fn session(upstream: Dialect) -> Session { + prepared(&compaction_request(), upstream).session +} + +fn stream_events(s: &Session, upstream_bytes: &str) -> Vec<(String, Value)> { + let mut c = s.stream(); + let mut out = Vec::new(); + for piece in upstream_bytes.as_bytes().chunks(7) { + out.extend(c.process(piece)); + } + out.extend(c.finish()); + let mut dec = Decoder::default(); + dec.feed(&out) + .into_iter() + .map(|f| (f.event.unwrap(), serde_json::from_str(&f.data).unwrap())) + .collect() +} + +/// Codex 收压缩的流时认的(`collect_compaction_output`):`output_item.done` 里恰好一个 +/// `compaction` 项,有 `response.completed`。返回那个项 +fn as_codex_reads_it(events: &[(String, Value)], at: &str) -> Value { + let done: Vec<&Value> = events + .iter() + .filter(|(k, _)| k == "response.output_item.done") + .map(|(_, v)| &v["item"]) + .collect(); + let compactions: Vec<&&Value> = done.iter().filter(|i| i["type"] == "compaction").collect(); + assert_eq!(compactions.len(), 1, "{at}: {done:?}"); + // 只有这一个项:摘要的文字、思考都不另外往外写 + assert_eq!(done.len(), 1, "{at}: {done:?}"); + let completed = events + .iter() + .find(|(k, _)| k == "response.completed") + .unwrap_or_else(|| panic!("{at}: no response.completed")); + assert!( + completed.1["response"]["id"].is_string(), + "{at}: {completed:?}" + ); + let item = (*compactions[0]).clone(); + // Codex 的 `ResponseItem::Compaction` 要的字段 + assert!(item["encrypted_content"].is_string(), "{at}: {item}"); + item +} + +#[test] +fn the_summary_comes_back_to_codex_as_exactly_one_compaction_item() { + for upstream in UPSTREAMS { + let at = format!("{upstream:?}"); + let events = stream_events(&session(upstream), &summary_stream(upstream)); + let item = as_codex_reads_it(&events, &at); + let enc = item["encrypted_content"].as_str().unwrap(); + assert!(enc.starts_with("tw1.c."), "{at}: {enc}"); + assert_eq!(compaction::read(enc).as_deref(), Some(SUMMARY), "{at}"); + // 写得久也不会被 Codex 当成断了:中间报着「还在写」 + assert!( + events.iter().any(|(k, _)| k == "response.in_progress"), + "{at}" + ); + // 用量照常交回,Codex 记账用 + let (_, done) = events.last().unwrap(); + assert!( + done["response"]["usage"]["output_tokens"].as_u64() > Some(0), + "{at}: {done}" + ); + assert_eq!(done["response"]["output"][0], item, "{at}"); + } +} + +#[test] +fn a_whole_or_collected_answer_is_one_compaction_item_too() { + for upstream in UPSTREAMS { + // 上游给流、客户端要整包 + let mut c = session(upstream).collector(); + c.process(summary_stream(upstream).as_bytes()); + let v: Value = serde_json::from_slice(&c.finish().unwrap()).unwrap(); + let output = v["output"].as_array().unwrap(); + assert_eq!(output.len(), 1, "{upstream:?}: {v}"); + assert_eq!(output[0]["type"], "compaction"); + assert_eq!( + compaction::read(output[0]["encrypted_content"].as_str().unwrap()).as_deref(), + Some(SUMMARY), + "{upstream:?}" + ); + assert_eq!(v["status"], "completed"); + } +} + +#[test] +fn an_answer_without_a_summary_fails_instead_of_wiping_the_conversation() { + // 上游没听话,调了个工具、一个字没写:交回空摘要的话 Codex 会把前文全扔掉 + let s = session(Dialect::Anthropic); + let up = [ + named("message_start", json!({"type": "message_start", "message": {"id": "m", "model": "claude-opus-4-7", "usage": {"input_tokens": 10, "output_tokens": 1}}})), + named("content_block_start", json!({"type": "content_block_start", "index": 0, "content_block": {"type": "tool_use", "id": "t", "name": "exec_command", "input": {}}})), + named("content_block_delta", json!({"type": "content_block_delta", "index": 0, "delta": {"type": "input_json_delta", "partial_json": "{\"cmd\":\"ls\"}"}})), + named("content_block_stop", json!({"type": "content_block_stop", "index": 0})), + named("message_delta", json!({"type": "message_delta", "delta": {"stop_reason": "tool_use"}, "usage": {"output_tokens": 5}})), + named("message_stop", json!({"type": "message_stop"})), + ] + .concat(); + let events = stream_events(&s, &up); + assert!( + !events + .iter() + .any(|(k, _)| k.starts_with("response.output_item")), + "{events:?}" + ); + let (last, v) = events.last().unwrap(); + assert_eq!(last, "response.failed"); + assert!( + v["response"]["error"]["message"] + .as_str() + .unwrap() + .contains("without writing a summary") + ); + let whole = s + .response( + json!({"id": "m", "type": "message", "role": "assistant", "content": [], + "stop_reason": "end_turn", "usage": {"input_tokens": 1, "output_tokens": 1}}) + .to_string() + .as_bytes(), + ) + .unwrap(); + let v: Value = serde_json::from_slice(&whole).unwrap(); + assert_eq!(v["status"], "failed"); + assert_eq!(v["output"], json!([])); +} + +/// 压缩之后 Codex 的历史:保留下来的用户消息、压缩项、重新写的上下文,接着是新的一轮 +/// (`build_v2_compacted_history`、`build_compaction_replacement_history`) +fn after_compaction(item: Value) -> Value { + request(vec![ + tools(), + json!({"id": "msg_b", "type": "message", "role": "developer", + "content": [{"type": "input_text", "text": "You are Codex, a coding agent."}]}), + json!({"type": "message", "role": "user", "content": [{"type": "input_text", "text": "Fix the failing parser test."}]}), + item, + json!({"type": "message", "role": "developer", + "content": [{"type": "input_text", "text": "sandbox: workspace-write"}]}), + json!({"type": "message", "role": "user", + "content": [{"type": "input_text", "text": "\n /repo\n"}]}), + json!({"type": "message", "role": "user", "content": [{"type": "input_text", "text": "Now update the docs."}]}), + ]) +} + +#[test] +fn the_next_request_reads_the_summary_back_in_its_place() { + let item = as_codex_reads_it( + &stream_events( + &session(Dialect::Anthropic), + &summary_stream(Dialect::Anthropic), + ), + "Anthropic", + ); + for upstream in UPSTREAMS { + let p = prepared(&after_compaction(item.clone()), upstream); + assert!(!p.session.is_compaction()); + let sent = String::from_utf8(p.body.clone()).unwrap(); + assert!(!sent.contains("tw1.c."), "{upstream:?}: {sent}"); + let (sys, _, msgs) = parts_of(upstream, &body_of(&p)); + // 摘要不进系统提示,在对话里它原来的位置:保留的那句话之后、新的一轮之前 + assert!( + !sys.to_string().contains("compacted"), + "{upstream:?}: {sys}" + ); + let convo = Value::Array(msgs).to_string(); + let (summary_at, kept_at, new_at) = ( + convo.find("parser test fixed").unwrap(), + convo.find("Fix the failing parser test.").unwrap(), + convo.find("Now update the docs.").unwrap(), + ); + assert!(kept_at < summary_at && summary_at < new_at, "{upstream:?}"); + assert!( + convo.contains( + "\\nThe earlier part of this conversation was compacted" + ), + "{upstream:?}: {convo}" + ); + } + + // OpenAI 自己的压缩项照旧读不了 + let theirs = after_compaction(json!({"type": "compaction", "encrypted_content": "gAAAAABpZx"})); + let e = decode(Dialect::Responses, &theirs, "/v1/responses", None).unwrap_err(); + assert!(e.0.contains("only OpenAI can read"), "{e}"); +} + +#[test] +fn passing_straight_through_to_openai_turns_our_summary_back_into_a_message() { + let item = json!({"id": "cmp_1", "type": "compaction", "encrypted_content": compaction::carry(SUMMARY)}); + let body = after_compaction(item).to_string(); + let out: Value = + serde_json::from_slice(&strip_carried(Dialect::Responses, body.as_bytes()).unwrap()) + .unwrap(); + let before: Value = serde_json::from_str(&body).unwrap(); + let (a, b) = ( + before["input"].as_array().unwrap(), + out["input"].as_array().unwrap(), + ); + assert_eq!(a.len(), b.len()); + // 只换了那一项,原位换成写着摘要的 developer 消息 + for (i, (x, y)) in a.iter().zip(b).enumerate() { + if i == 3 { + assert_eq!(y["type"], "message"); + assert_eq!(y["role"], "developer"); + assert_eq!( + y["content"][0]["text"], + compaction::restored(SUMMARY).as_str() + ); + } else { + assert_eq!(x, y); + } + } + // 没有转换写出去的东西:一个字节都不改 + assert_eq!( + strip_carried( + Dialect::Responses, + compaction_request().to_string().as_bytes() + ), + None + ); +} diff --git a/crates/tw-gateway/tests/conversion.rs b/crates/tw-gateway/tests/conversion.rs index ee11731b..7643dbdc 100644 --- a/crates/tw-gateway/tests/conversion.rs +++ b/crates/tw-gateway/tests/conversion.rs @@ -700,3 +700,112 @@ async fn a_codex_responses_lite_request_to_openai_goes_byte_for_byte() { "同格式直通改了请求体" ); } + +/// Codex 压缩前文:历史末尾加 `compaction_trigger`(`codex-rs/core/src/compact_remote_v2.rs`) +fn codex_compaction_request() -> Value { + let mut v = codex_lite_request(); + let input = v["input"].as_array_mut().unwrap(); + input.extend([ + json!({"type": "message", "role": "assistant", "content": [{"type": "output_text", "text": "src/ main.rs"}]}), + json!({"type": "compaction_trigger"}), + ]); + v +} + +#[tokio::test] +async fn a_codex_compaction_on_a_claude_route_comes_back_as_one_compaction_item() { + let stream = [ + "event: message_start\ndata: {\"type\":\"message_start\",\"message\":{\"id\":\"msg_1\",\"model\":\"claude-opus-4-7\",\"usage\":{\"input_tokens\":40,\"output_tokens\":1}}}\n\n", + "event: content_block_start\ndata: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n\n", + "event: content_block_delta\ndata: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"The user asked to list files; \"}}\n\n", + "event: content_block_delta\ndata: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"they are src/ and main.rs.\"}}\n\n", + "event: content_block_stop\ndata: {\"type\":\"content_block_stop\",\"index\":0}\n\n", + "event: message_delta\ndata: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"},\"usage\":{\"output_tokens\":17}}\n\n", + "event: message_stop\ndata: {\"type\":\"message_stop\"}\n\n", + ] + .concat(); + let (up, seen) = upstream(200, "text/event-stream", stream).await; + let (gw, mut rx) = gateway(provider(up, Protocol::Anthropic), SecurityMode::Observe).await; + let (status, ct, body) = post( + gw, + "/v1/responses", + &[("authorization", "Bearer tw-k")], + codex_compaction_request(), + ) + .await; + assert_eq!(status, 200, "{body}"); + assert_eq!(ct, "text/event-stream"); + + // 发给 Claude 的:同样的历史和工具,末尾请它写摘要 + let sent: Value = serde_json::from_slice(&seen.lock().unwrap().body).unwrap(); + assert_eq!(sent["tools"].as_array().unwrap().len(), 3); + let last = sent["messages"].as_array().unwrap().last().unwrap().clone(); + assert_eq!(last["role"], "user"); + assert!( + last.to_string().contains("Write that summary now"), + "{last}" + ); + + // Codex 收到的:恰好一个 compaction 项,摘要在里面 + let frames = data_frames(&body); + let done: Vec<&Value> = frames + .iter() + .filter(|f| f["type"] == "response.output_item.done") + .map(|f| &f["item"]) + .collect(); + assert_eq!(done.len(), 1, "{body}"); + assert_eq!(done[0]["type"], "compaction"); + assert_eq!( + tw_dialect::compaction::read(done[0]["encrypted_content"].as_str().unwrap()).as_deref(), + Some("The user asked to list files; they are src/ and main.rs.") + ); + assert!( + frames.iter().any(|f| f["type"] == "response.completed"), + "{body}" + ); + + // 写摘要这一次和别的请求一样记用量 + let mut finished = None; + while let Ok(Ok(ev)) = tokio::time::timeout(Duration::from_secs(3), rx.recv()).await { + if let tw_api::Event::RequestFinished { usage, .. } = ev { + finished = usage; + break; + } + } + assert_eq!(finished.map(|u| (u.input, u.output)), Some((40, 17))); +} + +#[tokio::test] +async fn a_summary_written_on_a_claude_route_reaches_openai_as_a_message() { + // 这段对话之前在别家的上游上压缩过,现在这一跳直通 OpenAI:它读不了我们写的「密文」 + let reply = concat!( + "event: response.completed\ndata: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_1\",\"status\":\"completed\",", + "\"output\":[],\"usage\":{\"input_tokens\":5,\"output_tokens\":1}}}\n\n", + ); + let (up, seen) = upstream(200, "text/event-stream", reply.into()).await; + let (gw, _) = gateway( + provider(up, Protocol::OpenaiResponses), + SecurityMode::Observe, + ) + .await; + let mut sent = codex_lite_request(); + sent["input"].as_array_mut().unwrap().insert( + 3, + json!({"type": "compaction", "encrypted_content": tw_dialect::compaction::carry("Listed the files.")}), + ); + let (status, _, body) = post( + gw, + "/v1/responses", + &[("authorization", "Bearer tw-k")], + sent, + ) + .await; + assert_eq!(status, 200, "{body}"); + let got: Value = serde_json::from_slice(&seen.lock().unwrap().body).unwrap(); + assert_eq!(got["input"][3]["role"], "developer"); + assert_eq!( + got["input"][3]["content"][0]["text"], + tw_dialect::compaction::restored("Listed the files.").as_str() + ); + assert!(!got.to_string().contains("tw1.c.")); +} diff --git a/crates/tw-gateway/tests/harness.rs b/crates/tw-gateway/tests/harness.rs index e7117305..cca5da11 100644 --- a/crates/tw-gateway/tests/harness.rs +++ b/crates/tw-gateway/tests/harness.rs @@ -425,22 +425,19 @@ async fn messages_to_an_openai_upstream_are_cleaned_by_the_conversion() { assert_clean(&s); assert!(s.headers.get("anthropic-beta").is_none()); let v: Value = serde_json::from_slice(&s.body).unwrap(); - // 中途的系统提示并进了开头的那一条 - let system: Vec<&str> = v["messages"] - .as_array() - .unwrap() + // 开头的系统提示一条;中途的留在原位,写成标明是系统说的一段(系统提示不跟着变, + // 提示缓存接得上) + let messages = v["messages"].as_array().unwrap(); + let roles: Vec<&str> = messages .iter() - .filter(|m| m["role"] == "system") - .map(|m| m["content"].as_str().unwrap()) + .map(|m| m["role"].as_str().unwrap()) .collect(); + assert_eq!(roles, ["system", "user", "user", "assistant", "user"]); + assert_eq!(messages[0]["content"], "You are DeepSeek Harness."); assert_eq!( - system - .concat() - .matches("The project root is /work.") - .count(), - 1 + messages[2]["content"], + "\nThe project root is /work.\n" ); - assert!(system.concat().starts_with("You are DeepSeek Harness.")); assert_eq!(v["tools"].as_array().unwrap().len(), 2); let (bytes, translated) = events(&mut rx).await; @@ -530,16 +527,17 @@ async fn chat_to_an_anthropic_upstream_is_cleaned_by_the_conversion() { .iter() .map(|b| b["text"].as_str().unwrap()) .collect(); + // 开头的进系统提示;中途的留在原位,写成标明是系统说的一段(Anthropic 的消息里没有 + // 系统角色) + assert_eq!(system, ["You are DeepSeek Harness."]); + let messages = v["messages"].as_array().unwrap(); + assert!(messages.iter().all(|m| m["role"] != "system")); assert_eq!( - system, - ["You are DeepSeek Harness.", "The project root is /work."] - ); - assert!( - v["messages"] - .as_array() - .unwrap() - .iter() - .all(|m| m["role"] != "system") + messages[2]["content"], + json!([ + {"type": "text", "text": "\nThe project root is /work.\n"}, + {"type": "text", "text": "search it"} + ]) ); assert_eq!(v["thinking"]["type"], "enabled"); diff --git a/crates/tw-store/src/search/text.rs b/crates/tw-store/src/search/text.rs index 77579146..520dbe5f 100644 --- a/crates/tw-store/src/search/text.rs +++ b/crates/tw-store/src/search/text.rs @@ -124,7 +124,8 @@ fn trim_for_reading(v: &mut Value, dialect: Dialect) { items.retain(|it| { !matches!( it.get("type").and_then(Value::as_str), - Some("item_reference" | "compaction") + // `compaction_trigger`:解码会在末尾加一条请上游写摘要的话,那不是用户说的 + Some("item_reference" | "compaction" | "compaction_trigger") ) }); } @@ -428,6 +429,17 @@ mod tests { ); } + /// Codex 要压缩前文的请求:末尾的 `compaction_trigger` 在转换时变成一句请上游写摘要的话, + /// 那不是用户说的。这一轮读到的是用户最后说的那句 + #[test] + fn a_compaction_request_reads_the_users_last_words_not_the_gateways() { + let body = json!({"input": [ + {"type": "message", "role": "user", "content": [{"type": "input_text", "text": "改完了吗"}]}, + {"type": "compaction_trigger"} + ]}); + assert_eq!(turn(body, "/v1/responses").unwrap(), "改完了吗"); + } + /// Python 客户端把 ASCII 以外的字全写成 `\uXXXX`:按 JSON 解开之后就是原字 #[test] fn escaped_text_is_read_as_the_characters_it_stands_for() { diff --git a/crates/tw-store/src/transcript/mod.rs b/crates/tw-store/src/transcript/mod.rs index 5b207879..e0f56b9d 100644 --- a/crates/tw-store/src/transcript/mod.rs +++ b/crates/tw-store/src/transcript/mod.rs @@ -2063,6 +2063,24 @@ mod tests { ] ); assert!(t.turns[0].output.is_empty() && t.turns[0].gaps.is_empty()); + + // 转换时压缩出的前文:摘要是上游写的,读得出来 + let mut d = Disk::new(); + d.turn( + "/v1/responses", + &json!({"model": "m", "input": [ + {"type": "compaction", "encrypted_content": tw_dialect::compaction::carry("改过 src/a.rs")}, + {"role": "user", "content": "接着来"} + ]}), + br#"{"output":[]}"#, + ); + assert_eq!( + d.transcript().turns[0].input[0], + msg( + R::System, + vec![text(&tw_dialect::compaction::restored("改过 src/a.rs"))] + ) + ); } // ───────────────────────────────────────────────── 读过的不再读 diff --git a/crates/tw-store/src/transcript/read.rs b/crates/tw-store/src/transcript/read.rs index a34607f4..e2228cac 100644 --- a/crates/tw-store/src/transcript/read.rs +++ b/crates/tw-store/src/transcript/read.rs @@ -611,8 +611,17 @@ pub(super) fn responses_item(it: &Value) -> Item<'_> { pieces, } } - // 压缩过的前文:只有 OpenAI 读得懂的一段密文,在对话里的位置像一条系统消息 - "compaction" => one(Role::System, Piece::Other(kind)), + // 压缩过的前文:在对话里的位置像一条系统消息。转换时写的是上游写的摘要,读得出来; + // OpenAI 的是只有它读得懂的一段密文 + "compaction" => { + match str_of(it, "encrypted_content").and_then(tw_dialect::compaction::read) { + Some(summary) => one( + Role::System, + Piece::Text(Cow::Owned(tw_dialect::compaction::restored(&summary))), + ), + None => one(Role::System, Piece::Other(kind)), + } + } // 别的工具结果(computer_call_output……) k if k.ends_with("_output") => one(Role::Tool, Piece::Other(k)), // 托管工具的调用(web_search_call、image_generation_call、mcp_call、local_shell_call……)