diff --git a/crates/tw-api/msg-codes.txt b/crates/tw-api/msg-codes.txt index 1ab08c40..65e6ba0a 100644 --- a/crates/tw-api/msg-codes.txt +++ b/crates/tw-api/msg-codes.txt @@ -206,6 +206,7 @@ gw.config.proxy_unusable gw.config.security_rules passthrough gw.content.refused gw.convert.failed +gw.files.unsupported gw.hidden_text.refused_message gw.hidden_text.refused_tool_result gw.internal diff --git a/crates/tw-api/src/lib.rs b/crates/tw-api/src/lib.rs index d5231b72..5392a914 100644 --- a/crates/tw-api/src/lib.rs +++ b/crates/tw-api/src/lib.rs @@ -593,6 +593,10 @@ pub const MSG_CODES: &str = include_str!("../msg-codes.txt"); /// **23 起额度也来自 GLM Coding Plan**(Z.ai / BigModel 的上游):窗口多了 `monthly`, /// 积分制套餐的窗口带上 [`QuotaWindow::credits`](总额、已用、剩余)。照 22 写的界面 /// 不认 `monthly`,也看不到剩余积分。 +/// +/// 同一版起**记录说得出请求带没带 DeepSeek Harness 的会话日志**:`RequestStarted` 和 +/// [`HistoryRow`] 多了 `session_log_bytes`(`dsh_session_log` 的字节数,没带的没有)。 +/// 照 22 写的界面看不到它。 pub const CONTROL_API_VERSION: u32 = 23; #[derive(Debug, Clone, Serialize, Deserialize)] @@ -723,6 +727,15 @@ pub enum Event { model: String, method: String, path: String, + /// 请求带着 DeepSeek Harness 的会话日志(`dsh_session_log`):这是它序列化之后 + /// 的字节数。没带是 None。 + /// + /// **会话日志是整段对话**(工作目录、系统提示、每一轮的输入输出、工具的参数和 + /// 结果),默认开着,每个请求补上一段,单次最多 8 MiB。它不进模型输入,只有 + /// DeepSeek 官方收:发给别的上游之前网关会去掉它,界面要能说出「这个请求带着 + /// 会话日志、有多大」 + #[serde(default, skip_serializing_if = "Option::is_none")] + session_log_bytes: Option, at_ms: u64, }, /// 收到上游响应头。**这个事件单独存在是有意的**:流式请求从这里 @@ -3105,6 +3118,10 @@ pub struct HistoryRow { /// 请求带的那把网关密钥打码后的样子(`tw-re…wb4e`),请求那一刻的 #[serde(default, skip_serializing_if = "Option::is_none")] pub key_masked: Option, + /// 请求带着 DeepSeek Harness 的会话日志:它的字节数。没带的没有,理由见 + /// `Event::RequestStarted::session_log_bytes` + #[serde(default, skip_serializing_if = "Option::is_none")] + pub session_log_bytes: Option, /// 这次请求在两项防护上留下的记录(和安全日志同一份)。没有就是空的。 /// /// **流量页的徽标靠它。**以前徽标只来自实时事件,关窗再开就没了 —— @@ -4270,6 +4287,7 @@ mod tests { model: "m".into(), method: "POST".into(), path: "/v1/messages".into(), + session_log_bytes: None, at_ms: 0, }, Event::RequestHeaders { @@ -4327,6 +4345,7 @@ mod tests { model: "m".into(), method: "POST".into(), path: "/v1/messages".into(), + session_log_bytes: None, at_ms: 0, }; let v = serde_json::to_value(&e).unwrap(); diff --git a/crates/tw-control/src/lib.rs b/crates/tw-control/src/lib.rs index d1897877..46db48ed 100644 --- a/crates/tw-control/src/lib.rs +++ b/crates/tw-control/src/lib.rs @@ -1106,6 +1106,7 @@ fn history_row( client_hint: r.client_hint, peer: r.peer, key_masked: r.key_masked, + session_log_bytes: r.session_log_bytes, security, } } diff --git a/crates/tw-control/tests/bundle.rs b/crates/tw-control/tests/bundle.rs index aa95239c..bccd51d5 100644 --- a/crates/tw-control/tests/bundle.rs +++ b/crates/tw-control/tests/bundle.rs @@ -82,6 +82,7 @@ async fn replaying_a_truncated_body_is_refused_rather_than_misleading() { let db = tw_store::Db::open(&d.path().join("data.db")).unwrap(); let blobs = tw_store::Blobs::new(d.path().join("blobs")); let mut row = tw_store::db::RequestRow { + session_log_bytes: None, key_masked: None, peer: None, id: 1, diff --git a/crates/tw-control/tests/in_flight.rs b/crates/tw-control/tests/in_flight.rs index fb50273a..5e5bec80 100644 --- a/crates/tw-control/tests/in_flight.rs +++ b/crates/tw-control/tests/in_flight.rs @@ -51,6 +51,7 @@ fn started(id: u64, model: &str) -> tw_api::Event { model: model.into(), method: "POST".into(), path: "/v1/messages".into(), + session_log_bytes: None, at_ms: 1_000 + id, } } diff --git a/crates/tw-control/tests/live_state.rs b/crates/tw-control/tests/live_state.rs index 18b60050..0cc68d12 100644 --- a/crates/tw-control/tests/live_state.rs +++ b/crates/tw-control/tests/live_state.rs @@ -111,6 +111,7 @@ fn started(id: u64) -> tw_api::Event { model: "claude-sonnet-4-5".into(), method: "POST".into(), path: "/v1/messages".into(), + session_log_bytes: None, at_ms: 1_758_000_000_000, } } diff --git a/crates/tw-control/tests/replay.rs b/crates/tw-control/tests/replay.rs index 41c56732..c6236d02 100644 --- a/crates/tw-control/tests/replay.rs +++ b/crates/tw-control/tests/replay.rs @@ -87,6 +87,7 @@ fn system_proxy() { fn row(id: i64, provider: &str) -> tw_store::db::RequestRow { tw_store::db::RequestRow { + session_log_bytes: None, key_masked: None, peer: None, id, diff --git a/crates/tw-control/tests/security.rs b/crates/tw-control/tests/security.rs index eb709844..5008f2d8 100644 --- a/crates/tw-control/tests/security.rs +++ b/crates/tw-control/tests/security.rs @@ -44,6 +44,7 @@ impl Bed { fn request(id: i64, at_ms: i64, provider: &str) -> tw_store::db::RequestRow { tw_store::db::RequestRow { + session_log_bytes: None, key_masked: None, peer: None, id, diff --git a/crates/tw-control/tests/sessions.rs b/crates/tw-control/tests/sessions.rs index 7f6b6b44..2a686933 100644 --- a/crates/tw-control/tests/sessions.rs +++ b/crates/tw-control/tests/sessions.rs @@ -16,6 +16,7 @@ const BASE: &str = "version: 1\nlisten:\n control:\n key: c0ffee00c0ffee00c0 fn turn(id: i64, billing: &str) -> tw_store::db::RequestRow { let per_token = billing == "per-token"; tw_store::db::RequestRow { + session_log_bytes: None, key_masked: None, peer: None, id, diff --git a/crates/tw-control/tests/summary.rs b/crates/tw-control/tests/summary.rs index fc530ebe..e081c4b2 100644 --- a/crates/tw-control/tests/summary.rs +++ b/crates/tw-control/tests/summary.rs @@ -15,6 +15,7 @@ const BASE: &str = "version: 1\nlisten:\n control:\n key: c0ffee00c0ffee00c0 /// 一条有用量、有价格的请求。 fn req(id: i64, model: &str) -> tw_store::db::RequestRow { tw_store::db::RequestRow { + session_log_bytes: None, key_masked: None, peer: None, id, diff --git a/crates/tw-dialect/src/anthropic/request.rs b/crates/tw-dialect/src/anthropic/request.rs index 3421d2a0..21cf2cab 100644 --- a/crates/tw-dialect/src/anthropic/request.rs +++ b/crates/tw-dialect/src/anthropic/request.rs @@ -40,12 +40,23 @@ pub fn decode_request(v: &Value, dropped: &mut Dropped) -> Result Role::Assistant, - // 消息里的 system 角色:并进系统提示,位置信息丢失但内容保留 + // 消息里的 system 角色:并进系统提示,位置信息丢失但内容保留。DeepSeek + // Harness 在这里放改过的系统提示,还有对话中途增删工具的 `tool_addition` / + // `tool_removal`:别家没有这种写法 Some("system") => { - let t = text_of(m.get("content").unwrap_or(&Value::Null)); + let content = m.get("content").unwrap_or(&Value::Null); + let t = text_of(content); if !t.is_empty() { r.system.push(t); } + for b in content.as_array().into_iter().flatten() { + match str_of(b, "type") { + Some("text") => {} + other => { + dropped.path(format!("messages.content.{}", other.unwrap_or("unknown"))) + } + } + } continue; } _ => Role::User, @@ -59,6 +70,11 @@ pub fn decode_request(v: &Value, dropped: &mut Dropped) -> Result r.tools.push(Tool { name: str_of(t, "name").unwrap_or_default().to_string(), diff --git a/crates/tw-dialect/src/chat/request.rs b/crates/tw-dialect/src/chat/request.rs index 7743d882..2868e284 100644 --- a/crates/tw-dialect/src/chat/request.rs +++ b/crates/tw-dialect/src/chat/request.rs @@ -139,6 +139,27 @@ pub fn decode_request( summary: e.is_some(), }); } + // DeepSeek、GLM、Kimi 的写法:`thinking.type` 开关推理,强度另写在 `reasoning_effort`。 + // 关掉时以它为准 + match v.get("thinking").and_then(|t| str_of(t, "type")) { + Some("disabled") => { + r.reasoning = Some(Reasoning { + enabled: false, + effort: None, + budget: None, + summary: false, + }) + } + Some("enabled") if r.reasoning.is_none() => { + r.reasoning = Some(Reasoning { + enabled: true, + effort: None, + budget: None, + summary: true, + }) + } + _ => {} + } r.format = match v.get("response_format").and_then(|f| str_of(f, "type")) { Some("json_object") => Some(Format::JsonObject), diff --git a/crates/tw-dialect/src/convert.rs b/crates/tw-dialect/src/convert.rs index 95950163..7eff6776 100644 --- a/crates/tw-dialect/src/convert.rs +++ b/crates/tw-dialect/src/convert.rs @@ -8,7 +8,7 @@ //! ``` //! //! **同格式不经过这里**:客户端和上游说同一种格式时,网关原样转发,一个字节都不改。 -//! 唯一的例外见 [`strip_carried`]。 +//! 例外见 [`strip_carried`] 和 [`crate::harness::clean`]。 use std::collections::{HashMap, HashSet}; diff --git a/crates/tw-dialect/src/harness.rs b/crates/tw-dialect/src/harness.rs new file mode 100644 index 00000000..65fb2210 --- /dev/null +++ b/crates/tw-dialect/src/harness.rs @@ -0,0 +1,459 @@ +//! DeepSeek Harness(dsh)发的请求:认出它,发给别家之前去掉只有 DeepSeek 认的东西。 +//! +//! dsh 是 DeepSeek 官方开源的编码代理(0.1.7 发 Anthropic Messages,0.1.5 发 Chat +//! Completions)。**不管配置的地址是哪里,它都照 DeepSeek 官方接口的样子发**: +//! +//! - 请求头:User-Agent 是 `deepseek-harness/<版本> (+<地址>)`,另有 +//! `x-deepseek-harness-user-id` / `-session-id` / `-compact` +//! - 顶层的 `dsh_*` 字段:`dsh_session_log`(整段会话日志,默认开着,单次最多 8 MiB)、 +//! `dsh_plugin_packages` +//! - messages 里 `role: system` 的条目(改过的系统提示),其中可以有 `tool_addition` / +//! `tool_removal` 块(对话中途增删工具,配着 `anthropic-beta: +//! mid-conversation-tool-changes-…` 和工具上的 `defer_loading`) +//! - 思考只写 `thinking.type: enabled` 加 `output_config.effort`,不带预算 +//! +//! **上游是 DeepSeek 官方时一个字节都不改**:这些本来就是发给它的。别的上游不认,有的 +//! 兼容实现还会因为不认识的字段整个拒绝;会话日志带着整段对话,更不该送给一个本来不收 +//! 它的上游。那时由 [`clean`] 去掉(同格式直通),或者由解码器在转换时丢掉。请求头那 +//! 两件([`own_header`]、[`anthropic_beta`])在网关里做。 + +use serde_json::{Map, Value, json}; + +use crate::ir::{Dialect, Effort, str_of}; +use crate::think; + +/// User-Agent 的产品名。实测:`deepseek-harness/0.1.7 (+https://github.com/…)` +const PRODUCT: &str = "deepseek-harness/"; +/// dsh 自己的请求头的前缀 +const HEADER_PREFIX: &str = "x-deepseek-harness-"; +/// 请求体里 dsh 扩展字段的前缀。官方文档把这个前缀整个留给了 dsh +const FIELD_PREFIX: &str = "dsh_"; +const SESSION_LOG: &str = "dsh_session_log"; +/// 对话中途增删工具要开的 beta,日期部分会变 +const TOOL_CHANGES_BETA: &str = "mid-conversation-tool-changes-"; + +/// 认出来的 dsh 请求。 +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +pub struct Harness { + /// 请求带着的会话日志(`dsh_session_log`)序列化之后有多少字节。没带是 None + pub session_log_bytes: Option, +} + +/// 这个请求是不是 dsh 发的:User-Agent 是它的,或者请求体里有 `dsh_*` 字段。 +/// +/// **两条任一即可**:User-Agent 可以被配置改掉,扩展字段也可以关掉。 +pub fn detect(user_agent: Option<&str>, body: Option<&Value>) -> Option { + let by_agent = user_agent + .and_then(|ua| ua.get(..PRODUCT.len())) + .is_some_and(|p| p.eq_ignore_ascii_case(PRODUCT)); + let fields = body.and_then(Value::as_object); + let by_fields = fields.is_some_and(|o| o.keys().any(|k| k.starts_with(FIELD_PREFIX))); + if !by_agent && !by_fields { + return None; + } + Some(Harness { + session_log_bytes: fields + .and_then(|o| o.get(SESSION_LOG)) + .filter(|v| !v.is_null()) + .map(json_len), + }) +} + +/// 一个 JSON 值写出来有多少字节。**只数不存**:会话日志可以有 8 MiB +fn json_len(v: &Value) -> u64 { + struct Count(u64); + impl std::io::Write for Count { + fn write(&mut self, b: &[u8]) -> std::io::Result { + self.0 += b.len() as u64; + Ok(b.len()) + } + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + let mut c = Count(0); + // 写进一个只计数的 writer,不会失败 + let _ = serde_json::to_writer(&mut c, v); + c.0 +} + +/// 这是 dsh 自己的请求头(`x-deepseek-harness-*`)。发给 DeepSeek 官方以外的上游时不转发: +/// 里面是一个匿名用户 id 和会话 id,别家用不上。 +pub fn own_header(name: &str) -> bool { + name.get(..HEADER_PREFIX.len()) + .is_some_and(|p| p.eq_ignore_ascii_case(HEADER_PREFIX)) +} + +/// `anthropic-beta` 去掉对话中途增删工具那一项(那些块已经被 [`clean`] 去掉了)。 +/// 一项都不剩时是 None,这个头就不发了。 +pub fn anthropic_beta(value: &str) -> Option { + let kept: Vec<&str> = value + .split(',') + .map(str::trim) + .filter(|b| !b.is_empty() && !b.starts_with(TOOL_CHANGES_BETA)) + .collect(); + (!kept.is_empty()).then(|| kept.join(",")) +} + +/// 清理过的请求体。 +#[derive(Debug, Clone, PartialEq)] +pub struct Cleaned { + pub body: Vec, + /// 转不过去、被丢掉的字段,写法和格式转换列的一样 + pub dropped: Vec, +} + +/// 同格式直通、上游又不是 DeepSeek 官方时,去掉 dsh 的扩展。没有可去的时候是 None。 +/// +/// - 所有格式:顶层的 `dsh_*` 字段去掉 +/// - Anthropic:messages 里 `role: system` 的条目按顺序并进 `system`(Anthropic 的 +/// messages 里没有这个角色),里面的 `tool_addition` / `tool_removal` 去掉并列为丢弃; +/// 工具上的 `defer_loading` 去掉(没有了增加工具的块,被推迟的工具永远不会出现),也 +/// 列为丢弃;只写了 `thinking.type: enabled` 的,按模型补成 Anthropic 认的写法 +/// +/// Chat 的 messages 本来就有 `system` 角色,dsh 0.1.5 也不发增删工具的块,只去字段。 +pub fn clean(client: Dialect, body: &[u8]) -> Option { + let mut v: Value = serde_json::from_slice(body).ok()?; + let o = v.as_object_mut()?; + let before = o.len(); + o.retain(|k, _| !k.starts_with(FIELD_PREFIX)); + let mut changed = o.len() != before; + let mut dropped = Vec::new(); + if client == Dialect::Anthropic { + changed |= system_updates(o, &mut dropped); + changed |= deferred_tools(o, &mut dropped); + changed |= thinking(o); + } + changed.then(|| Cleaned { + body: v.to_string().into_bytes(), + dropped, + }) +} + +/// 请求转换成另一种格式、发给 DeepSeek 官方时,把客户端请求里的 `dsh_*` 字段照原样 +/// 带过去:直连 DeepSeek 时它们本来就在,两种格式它都收。客户端没带时是 None。 +pub fn carry(client_body: &[u8], into: &[u8]) -> Option> { + let client: Value = serde_json::from_slice(client_body).ok()?; + let fields: Vec<(&String, &Value)> = client + .as_object()? + .iter() + .filter(|(k, _)| k.starts_with(FIELD_PREFIX)) + .collect(); + if fields.is_empty() { + return None; + } + let mut v: Value = serde_json::from_slice(into).ok()?; + let o = v.as_object_mut()?; + for (k, x) in fields { + o.insert(k.clone(), x.clone()); + } + Some(v.to_string().into_bytes()) +} + +/// messages 里 `role: system` 的条目并进 `system`,按出现的顺序接在后面。 +fn system_updates(o: &mut Map, dropped: &mut Vec) -> bool { + let Some(messages) = o.get_mut("messages").and_then(Value::as_array_mut) else { + return false; + }; + let mut texts: Vec = Vec::new(); + let before = messages.len(); + messages.retain(|m| { + if str_of(m, "role") != Some("system") { + return true; + } + match m.get("content") { + Some(Value::String(s)) if !s.is_empty() => texts.push(s.clone()), + Some(Value::Array(blocks)) => { + for b in blocks { + match str_of(b, "type") { + Some("text") => { + if let Some(t) = str_of(b, "text").filter(|t| !t.is_empty()) { + texts.push(t.to_string()); + } + } + // tool_addition、tool_removal:别家没有对话中途增删工具的写法 + other => note( + dropped, + format!("messages.content.{}", other.unwrap_or("unknown")), + ), + } + } + } + _ => {} + } + false + }); + if messages.len() == before { + return false; + } + if !texts.is_empty() { + match o.get_mut("system") { + Some(Value::Array(blocks)) => { + blocks.extend( + texts + .into_iter() + .map(|t| json!({ "type": "text", "text": t })), + ); + } + // dsh 自己拼系统提示也是用空行隔开 + Some(Value::String(s)) if !s.is_empty() => { + for t in texts { + s.push_str("\n\n"); + s.push_str(&t); + } + } + _ => { + o.insert("system".into(), Value::String(texts.join("\n\n"))); + } + } + } + true +} + +/// 工具上的 `defer_loading` 去掉:它等着一个 `tool_addition` 来加载,而那些块发不过去。 +/// 留着的话,Anthropic 会要一个工具搜索工具,别家可能根本不认这个字段。 +fn deferred_tools(o: &mut Map, dropped: &mut Vec) -> bool { + let mut any = false; + if let Some(tools) = o.get_mut("tools").and_then(Value::as_array_mut) { + for t in tools.iter_mut().filter_map(Value::as_object_mut) { + any |= t.remove("defer_loading").is_some(); + } + } + if any { + note(dropped, "tools.defer_loading".into()); + } + any +} + +/// dsh 的思考只写 `thinking.type: enabled` 加 `output_config.effort`,不带预算。 +/// DeepSeek 认这个写法,Anthropic 不认(`enabled` 必须带 `budget_tokens`),按模型改成 +/// 它认的:自适应的模型写 `adaptive`,强度留在 `output_config` 里;别的写成预算,强度 +/// 折进预算里(折法和格式转换用的同一张表,见 [`think`])。 +fn thinking(o: &mut Map) -> bool { + let Some(t) = o.get("thinking") else { + return false; + }; + if str_of(t, "type") != Some("enabled") || t.get("budget_tokens").is_some() { + return false; + } + let model = o.get("model").and_then(Value::as_str).unwrap_or_default(); + if think::claude_adaptive(model) { + o.insert("thinking".into(), json!({ "type": "adaptive" })); + return true; + } + // dsh 不写强度时用的是 high + let effort = o + .get("output_config") + .and_then(|c| str_of(c, "effort")) + .and_then(think::parse_anthropic) + .unwrap_or(Effort::High); + // 预算不少于 1024,且必须小于 max_tokens;放不下就不开 + let budget = think::budget_of_effort(effort); + let budget = match o.get("max_tokens").and_then(Value::as_u64) { + Some(m) => budget.min(m.saturating_sub(1)), + None => budget, + }; + let thinking = if budget < 1024 { + json!({ "type": "disabled" }) + } else { + json!({ "type": "enabled", "budget_tokens": budget }) + }; + o.insert("thinking".into(), thinking); + // 手动预算的模型不认强度,它说的已经折进预算了 + if let Some(c) = o.get_mut("output_config").and_then(Value::as_object_mut) { + c.remove("effort"); + if c.is_empty() { + o.remove("output_config"); + } + } + true +} + +fn note(dropped: &mut Vec, path: String) { + if !dropped.contains(&path) { + dropped.push(path); + } +} + +#[cfg(test)] +pub(crate) mod tests { + use super::*; + + /// dsh 0.1.7 发的 Anthropic 请求:会话日志、插件清单、中途的系统提示和增删工具 + pub(crate) fn anthropic_request() -> Value { + json!({ + "model": "deepseek-v4-pro", + "stream": true, + "max_tokens": 32000, + "system": "You are DeepSeek Harness.", + "thinking": { "type": "enabled" }, + "output_config": { "effort": "max" }, + "tools": [ + { "name": "read", "description": "Read a file", "input_schema": { "type": "object" } }, + { "name": "web", "description": "Search", "input_schema": { "type": "object" }, "defer_loading": true } + ], + "messages": [ + { "role": "user", "content": [{ "type": "text", "text": "hi" }] }, + { "role": "system", "content": [ + { "type": "text", "text": "The project root is /work." }, + { "type": "tool_addition", "tool": { "type": "tool_reference", "name": "web" } } + ] }, + { "role": "assistant", "content": [{ "type": "text", "text": "hello" }] }, + { "role": "user", "content": [{ "type": "text", "text": "search" }] } + ], + "dsh_session_log": { + "version": 1, "sessionFormatVersion": 2, + "session": { "version": 2, "id": "s", "createdAt": 1780000000000u64 }, + "afterSeq": -1, "throughSeq": 0, + "events": [{ "type": "turn/start", "seq": 0, "time": 1780000000001u64, "data": { "turn": 1 } }] + }, + "dsh_plugin_packages": { "version": 1, "packages": [] } + }) + } + + #[test] + fn a_harness_request_is_known_by_its_agent_or_its_fields() { + let ua = "deepseek-harness/0.1.7 (+https://github.com/deepseek-ai/deepseek-harness)"; + assert!(detect(Some(ua), None).is_some()); + assert!(detect(Some("DeepSeek-Harness/0.1.5"), None).is_some()); + let body = anthropic_request(); + let h = detect(Some("node"), Some(&body)).unwrap(); + let log = serde_json::to_vec(&body["dsh_session_log"]).unwrap(); + assert_eq!(h.session_log_bytes, Some(log.len() as u64)); + // 关掉了会话日志的 dsh:认得出,但没有会话日志 + let h = detect(Some(ua), Some(&json!({"model": "m", "messages": []}))).unwrap(); + assert_eq!(h.session_log_bytes, None); + assert!(detect(Some("claude-cli/2.1.0"), Some(&json!({"model": "m"}))).is_none()); + assert!(detect(None, None).is_none()); + // 一个多字节字符不会让前缀比较越界 + assert!(detect(Some("深度求索"), None).is_none()); + } + + #[test] + fn only_its_own_headers_and_its_beta_are_singled_out() { + assert!(own_header("x-deepseek-harness-user-id")); + assert!(own_header("X-DeepSeek-Harness-Session-Id")); + assert!(!own_header("x-deepseek-request-id")); + assert!(!own_header("user-agent")); + assert_eq!( + anthropic_beta("mid-conversation-tool-changes-2026-07-01"), + None + ); + assert_eq!( + anthropic_beta("files-api-2025-04-14, mid-conversation-tool-changes-2026-07-01") + .as_deref(), + Some("files-api-2025-04-14") + ); + } + + #[test] + fn an_anthropic_request_is_cleaned_for_another_upstream() { + let body = anthropic_request().to_string(); + let c = clean(Dialect::Anthropic, body.as_bytes()).unwrap(); + let v: Value = serde_json::from_slice(&c.body).unwrap(); + let o = v.as_object().unwrap(); + assert!(!o.keys().any(|k| k.starts_with("dsh_")), "{v}"); + // 系统提示的更新按顺序接在 system 后面,消息里不再有 system 角色 + assert_eq!( + v["system"], + "You are DeepSeek Harness.\n\nThe project root is /work." + ); + let roles: Vec<&str> = v["messages"] + .as_array() + .unwrap() + .iter() + .map(|m| m["role"].as_str().unwrap()) + .collect(); + assert_eq!(roles, ["user", "assistant", "user"]); + assert!(v["tools"][1].get("defer_loading").is_none()); + assert_eq!( + c.dropped, + ["messages.content.tool_addition", "tools.defer_loading"] + ); + // 不是 Claude 的模型:强度折成预算,output_config 里只有它时整个去掉 + assert_eq!( + v["thinking"], + json!({"type": "enabled", "budget_tokens": 31999}) + ); + assert!(v.get("output_config").is_none()); + } + + #[test] + fn thinking_is_written_the_way_the_model_takes_it() { + let mut r = anthropic_request(); + r["model"] = json!("claude-opus-4-7"); + let c = clean(Dialect::Anthropic, r.to_string().as_bytes()).unwrap(); + let v: Value = serde_json::from_slice(&c.body).unwrap(); + assert_eq!(v["thinking"], json!({"type": "adaptive"})); + assert_eq!(v["output_config"]["effort"], "max"); + + // 预算放不下:不开 + let mut r = anthropic_request(); + r["max_tokens"] = json!(512); + r["output_config"] = + json!({"effort": "low", "format": {"type": "json_schema", "schema": {}}}); + let c = clean(Dialect::Anthropic, r.to_string().as_bytes()).unwrap(); + let v: Value = serde_json::from_slice(&c.body).unwrap(); + assert_eq!(v["thinking"], json!({"type": "disabled"})); + assert!(v["output_config"].get("effort").is_none()); + assert_eq!(v["output_config"]["format"]["type"], "json_schema"); + } + + #[test] + fn system_updates_join_whatever_shape_the_system_prompt_has() { + let mut r = anthropic_request(); + r["system"] = json!([{"type": "text", "text": "A"}]); + let c = clean(Dialect::Anthropic, r.to_string().as_bytes()).unwrap(); + let v: Value = serde_json::from_slice(&c.body).unwrap(); + assert_eq!(v["system"][1]["text"], "The project root is /work."); + + let mut r = anthropic_request(); + r.as_object_mut().unwrap().remove("system"); + let c = clean(Dialect::Anthropic, r.to_string().as_bytes()).unwrap(); + let v: Value = serde_json::from_slice(&c.body).unwrap(); + assert_eq!(v["system"], "The project root is /work."); + } + + #[test] + fn a_chat_request_only_loses_the_extension_fields() { + let r = json!({ + "model": "deepseek-v4-flash", + "messages": [ + {"role": "system", "content": "You are DeepSeek Harness."}, + {"role": "user", "content": "hi"} + ], + "thinking": {"type": "enabled"}, + "reasoning_effort": "high", + "dsh_session_log": {"version": 1, "events": []} + }); + let c = clean(Dialect::Chat, r.to_string().as_bytes()).unwrap(); + let v: Value = serde_json::from_slice(&c.body).unwrap(); + assert!(v.get("dsh_session_log").is_none()); + assert_eq!(v["messages"], r["messages"]); + assert_eq!(v["thinking"], r["thinking"]); + assert!(c.dropped.is_empty()); + } + + #[test] + fn a_request_with_nothing_of_the_harness_is_left_alone() { + let r = json!({"model": "claude-sonnet-4-5", "max_tokens": 10, "messages": [ + {"role": "user", "content": "hi"} + ]}); + assert_eq!(clean(Dialect::Anthropic, r.to_string().as_bytes()), None); + assert_eq!(clean(Dialect::Anthropic, b"not json"), None); + } + + #[test] + fn the_extension_fields_ride_along_to_deepseek_after_a_conversion() { + let client = anthropic_request().to_string(); + let out = carry( + client.as_bytes(), + br#"{"model":"deepseek-v4-pro","messages":[]}"#, + ) + .unwrap(); + let v: Value = serde_json::from_slice(&out).unwrap(); + assert_eq!(v["dsh_session_log"], anthropic_request()["dsh_session_log"]); + assert_eq!(v["dsh_plugin_packages"]["version"], 1); + assert_eq!(carry(br#"{"model":"m"}"#, br#"{"model":"m"}"#), None); + } +} diff --git a/crates/tw-dialect/src/lib.rs b/crates/tw-dialect/src/lib.rs index 1ccf727b..6f7eea1c 100644 --- a/crates/tw-dialect/src/lib.rs +++ b/crates/tw-dialect/src/lib.rs @@ -6,7 +6,8 @@ //! [`ir`] 里的中间表示。这个 crate 只依赖 serde,不碰网络,企业版网关也可以直接用。 //! //! 同样两边都用的还有:从响应里旁路嗅出用量([`usage`],换算和转换共用各家的 -//! `usage()`),以及拼上游地址([`url`])。 +//! `usage()`),拼上游地址([`url`]),以及去掉 DeepSeek Harness 只发给 DeepSeek 的 +//! 扩展([`harness`])。 pub mod anthropic; pub mod bedrock; @@ -14,6 +15,7 @@ pub mod chat; pub mod convert; pub mod frame; pub mod gemini; +pub mod harness; pub mod ir; pub mod official; pub mod responses; diff --git a/crates/tw-dialect/src/official.rs b/crates/tw-dialect/src/official.rs index 84319db2..d2b4624a 100644 --- a/crates/tw-dialect/src/official.rs +++ b/crates/tw-dialect/src/official.rs @@ -35,6 +35,13 @@ pub fn is_official_host(url: &str) -> bool { OFFICIAL.contains(&host.as_str()) || host.ends_with(".amazonaws.com") } +/// 这个地址是不是 DeepSeek 官方的端点。DeepSeek Harness 的扩展(见 +/// [`crate::harness`])只有它认,发给别家之前要去掉。和 [`is_official_host`] 一样 +/// 在 host 上比。 +pub fn is_deepseek_host(url: &str) -> bool { + host_of(url).eq_ignore_ascii_case("api.deepseek.com") +} + /// 地址里的 host:去掉协议头、用户信息和端口,到路径、查询串、片段为止。 fn host_of(url: &str) -> &str { let rest = url.split_once("://").map_or(url, |(_, r)| r); @@ -89,6 +96,15 @@ mod tests { assert!(is_official_host("https://user:pw@api.openai.com/v1")); } + #[test] + fn deepseek_is_recognised_by_its_host_only() { + assert!(is_deepseek_host("https://api.deepseek.com")); + assert!(is_deepseek_host("https://API.DeepSeek.com/anthropic")); + assert!(!is_deepseek_host("https://relay.example.com/deepseek")); + assert!(!is_deepseek_host("https://api.deepseek.com.evil.com")); + assert!(!is_deepseek_host("https://api.deepseek.com@evil.com")); + } + #[test] fn the_fallback_output_limit_goes_by_the_model_name() { assert_eq!(fallback_max_output_tokens("Claude-Opus-4-5"), 32000); diff --git a/crates/tw-gateway/src/client_api.rs b/crates/tw-gateway/src/client_api.rs index 93453d39..f5d1a768 100644 --- a/crates/tw-gateway/src/client_api.rs +++ b/crates/tw-gateway/src/client_api.rs @@ -145,6 +145,8 @@ pub struct Reading { /// 才用得上**:同格式直通时,我们解不开的东西上游可能完全认识 pub decoded: Option>, pub facts: tw_engine::RequestFacts, + /// DeepSeek Harness 发的请求(见 [`tw_dialect::harness`])。要看请求头,由管线填 + pub harness: Option, } pub fn read(path: &str, query: Option<&str>, body: Option<&serde_json::Value>) -> Reading { @@ -179,6 +181,7 @@ pub fn read(path: &str, query: Option<&str>, body: Option<&serde_json::Value>) - generates, decoded, facts, + harness: None, } } diff --git a/crates/tw-gateway/src/hint.rs b/crates/tw-gateway/src/hint.rs index 5da2792b..ae93fa7d 100644 --- a/crates/tw-gateway/src/hint.rs +++ b/crates/tw-gateway/src/hint.rs @@ -39,6 +39,8 @@ pub fn client_hint(headers: &HeaderMap) -> Option { let ua = get("user-agent"); for (needle, id) in [ ("claude-cli", "claude-code"), + // `deepseek-harness/0.1.7 (+https://github.com/deepseek-ai/deepseek-harness)` + ("deepseek-harness", "deepseek-harness"), ("codex", "codex"), ("opencode", "opencode"), ("aider", "aider"), @@ -107,6 +109,17 @@ mod tests { ); } + #[test] + fn deepseek_harness_is_recognised_by_its_user_agent() { + assert_eq!( + client_hint(&h(&[( + "user-agent", + "deepseek-harness/0.1.7 (+https://github.com/deepseek-ai/deepseek-harness)" + )])), + Some("deepseek-harness".into()) + ); + } + #[test] fn an_unknown_client_is_none_rather_than_a_guess() { // 猜错比不猜更糟:观察窗口会因此宣布「已生效」。 diff --git a/crates/tw-gateway/src/server.rs b/crates/tw-gateway/src/server.rs index 016e16cf..2727a916 100644 --- a/crates/tw-gateway/src/server.rs +++ b/crates/tw-gateway/src/server.rs @@ -5,7 +5,7 @@ use std::sync::Arc; use axum::Router; use axum::extract::{OriginalUri, RawQuery, State}; use axum::http::HeaderMap; -use axum::response::Response; +use axum::response::{IntoResponse, Response}; use axum::routing::{any, get}; use bytes::Bytes; @@ -44,12 +44,43 @@ pub fn router(state: AppState) -> Router { "/v1beta/models/{model}", get(get_model).fallback(passthrough), ) + // Files API:网关不代管文件,一律 404(见 `no_files`)。带 `/v1` 和不带的都认, + // 和别的接口一样 + .route("/v1/files", any(no_files)) + .route("/v1/files/{*rest}", any(no_files)) + .route("/files", any(no_files)) + .route("/files/{*rest}", any(no_files)) // M0 只有透传:任何方法、任何路径都往上游送。M1 加路由时, // 这里会先过规则引擎再决定送给谁。 .fallback(any(passthrough)) .with_state(state) } +/// Files API 回 404。 +/// +/// **不能透传**:传上去的文件落在那一刻被选中的那家上游,之后引用它的请求可能被路由 +/// 到另一家(规则、故障转移),那家不认这个 file id。DeepSeek Harness 先把图片传到 +/// `/v1/files`,失败了就改成内联 base64 —— 404 正是让它退回内联的那个回答,内联的 +/// 图片哪家上游都能收(要转换时也转得过去)。 +async fn no_files(headers: HeaderMap) -> Response { + let dialect = if headers.contains_key("anthropic-version") { + tw_dialect::ir::Dialect::Anthropic + } else { + tw_dialect::ir::Dialect::Chat + }; + let m = msg!( + "gw.files.unsupported" => + "The gateway does not host files. Send images and documents inline in the request." + ); + let body = tw_dialect::convert::error_body(dialect, 404, &format!("[ThinkWatch] {}", m.text)); + let mut resp = (axum::http::StatusCode::NOT_FOUND, body).into_response(); + resp.headers_mut().insert( + axum::http::header::CONTENT_TYPE, + axum::http::HeaderValue::from_static("application/json"), + ); + resp +} + async fn passthrough( State(state): State, axum::extract::ConnectInfo(peer): axum::extract::ConnectInfo, diff --git a/crates/tw-gateway/src/server/pipeline.rs b/crates/tw-gateway/src/server/pipeline.rs index a78b4c09..d2e9e6d3 100644 --- a/crates/tw-gateway/src/server/pipeline.rs +++ b/crates/tw-gateway/src/server/pipeline.rs @@ -192,6 +192,12 @@ fn read(req: &Inbound, intent: String) -> (crate::client_api::Reading, Option, /// ChatGPT 账号(Codex 后端):只收流式、不认输出上限、身份头由网关填 chatgpt: bool, + /// DeepSeek Harness 的请求发给 DeepSeek 官方以外的上游:它自己的请求头不转发 + /// (见 [`tw_dialect::harness`]) + harness_elsewhere: bool, /// 回程要用的转换会话:转换过的,或者直通到 Codex 后端、客户端却要整包时 /// 收齐流要用的 session: Option, @@ -344,6 +347,9 @@ fn prepare( let client_dialect = req.api.map(|a| a.dialect()); let target = crate::translate::plan(req.api, generates, provider.effective_protocol()); let chatgpt = generates && provider.effective_protocol() == Some(tw_config::Protocol::Chatgpt); + // DeepSeek Harness 的扩展只有 DeepSeek 官方认。**发给它时一个字节都不改** + let harness = generates && reading.harness.is_some(); + let to_deepseek = tw_dialect::official::is_deepseek_host(&provider.base_url); let mut path = req.uri.path().to_string(); let mut query = req.query.clone(); let mut session: Option = None; @@ -367,20 +373,39 @@ fn prepare( Some(b) => Bytes::from(b), None => out, }; - if chatgpt { - // **只动 Codex 后端不认的那几个字段**,其余原样发(见 `chatgpt` 模块) - let (shaped, dropped) = crate::chatgpt::shape_passthrough(&out); - if !dropped.is_empty() { - let responses = tw_api::Dialect::OpenaiResponses; - state.bus.emit(tw_api::Event::Translated { - id, - provider: provider.name.clone(), - from: responses, - to: responses, - dropped, - at_ms: crate::server::now_ms(), - }); + // DeepSeek Harness 发给别家:去掉只有 DeepSeek 认的扩展。**去掉了什么要说**, + // 和转换丢了字段一样记在这一跳上 + let mut dropped = Vec::new(); + let out = match client_dialect + .filter(|_| harness && !to_deepseek) + .and_then(|d| tw_dialect::harness::clean(d, &out)) + { + Some(c) => { + dropped = c.dropped; + Bytes::from(c.body) } + None => out, + }; + let out = if chatgpt { + // **只动 Codex 后端不认的那几个字段**,其余原样发(见 `chatgpt` 模块) + let (shaped, more) = crate::chatgpt::shape_passthrough(&out); + dropped.extend(more); + shaped + } else { + out + }; + if let Some(d) = client_dialect.filter(|_| !dropped.is_empty()) { + let same = crate::wire::dialect(d); + state.bus.emit(tw_api::Event::Translated { + id, + provider: provider.name.clone(), + from: same, + to: same, + dropped, + at_ms: crate::server::now_ms(), + }); + } + if chatgpt { // 客户端要整包,后端只给流:由网关收齐。收齐要知道客户端的格式,所以要一个会话 if let Some(Ok(d)) = &reading.decoded && !d.request.stream @@ -397,10 +422,8 @@ fn prepare( .session, ); } - shaped - } else { - out } + out } Some(dialect) => { let d = match &reading.decoded { @@ -453,6 +476,11 @@ fn prepare( // **客户端要不要流由会话记着**,发给 Codex 后端的这一份一律是流式 let body = if chatgpt { Bytes::from(crate::chatgpt::force_stream(p.body.clone())) + } else if harness && to_deepseek { + // 转换成另一种格式发给 DeepSeek 官方:直连时它收得到的扩展照样带上 + Bytes::from( + tw_dialect::harness::carry(&req.body, &p.body).unwrap_or(p.body.clone()), + ) } else { Bytes::from(p.body.clone()) }; @@ -472,6 +500,7 @@ fn prepare( query, target, chatgpt, + harness_elsewhere: harness && !to_deepseek, session, }) } @@ -493,6 +522,14 @@ async fn send( let required = target .map(crate::translate::required_headers) .unwrap_or_default(); + // DeepSeek Harness 发给别家:它自己的头不转发;直通时 `anthropic-beta` 去掉对话中途 + // 增删工具那一项(那些块已经去掉了),剩下的照发 + let harness = out.harness_elsewhere; + let beta = headers + .get("anthropic-beta") + .filter(|_| harness && target.is_none()) + .and_then(|v| v.to_str().ok()) + .and_then(tw_dialect::harness::anthropic_beta); let build = |upstream_headers: &[(String, String)]| { let mut req = http.request(method.clone(), &url); req = forward::forward_headers_filtered(req, headers, |n| { @@ -503,10 +540,18 @@ async fn send( // 请求来自谁由网关如实填写:客户端报的来源(比如 Codex CLI 的 originator) // 不转发,请求经过的是 ThinkWatch let identity = chatgpt && !crate::chatgpt::keeps_client_header(n); + let dsh = harness + && (tw_dialect::harness::own_header(n) || n.eq_ignore_ascii_case("anthropic-beta")); !own && !identity + && !dsh && !required.iter().any(|(k, _)| k.eq_ignore_ascii_case(n)) && !forward::overridden(upstream_headers, n) }); + if let Some(b) = &beta + && !forward::overridden(upstream_headers, "anthropic-beta") + { + req = req.header("anthropic-beta", b.as_str()); + } // 目标格式必需的头(Anthropic 的 anthropic-version)。上游配置里写了同名头时以配置为准 for (k, v) in required { if !forward::overridden(upstream_headers, k) { diff --git a/crates/tw-gateway/src/server/upgrade.rs b/crates/tw-gateway/src/server/upgrade.rs index e7af610d..28d8cd92 100644 --- a/crates/tw-gateway/src/server/upgrade.rs +++ b/crates/tw-gateway/src/server/upgrade.rs @@ -58,6 +58,7 @@ pub(super) async fn ws_upgrade( model: String::new(), method: "WS".to_string(), path: uri.path().to_string(), + session_log_bytes: None, at_ms, }); // 这条连接怎么断的,就是这个请求的结局。**跟着连接走**:升级没完成 diff --git a/crates/tw-gateway/src/translate.rs b/crates/tw-gateway/src/translate.rs index 14e1c987..af4f7af8 100644 --- a/crates/tw-gateway/src/translate.rs +++ b/crates/tw-gateway/src/translate.rs @@ -1,8 +1,9 @@ //! 方言互转在管线上的接线。 //! //! **只在客户端格式和上游协议不同、而且调的是生成回答时才醒过来。**同格式那条路 -//! 一个字节都不碰(唯一的例外是去掉转换写出去的推理签名,见 -//! [`tw_dialect::convert::strip_carried`])。转换本身在 `tw-dialect` 里,这里只管 +//! 一个字节都不碰(例外有两个:去掉转换写出去的推理签名,见 +//! [`tw_dialect::convert::strip_carried`];DeepSeek Harness 发给 DeepSeek 以外的上游时 +//! 去掉它的扩展,见 [`tw_dialect::harness`])。转换本身在 `tw-dialect` 里,这里只管 //! 三个接缝: //! //! 1. **要不要转**:上游协议认不出来时不转。猜一个方向的代价是「请求被改成另一个 diff --git a/crates/tw-gateway/tests/harness.rs b/crates/tw-gateway/tests/harness.rs new file mode 100644 index 00000000..8c6f83b6 --- /dev/null +++ b/crates/tw-gateway/tests/harness.rs @@ -0,0 +1,604 @@ +//! DeepSeek Harness(dsh)的请求在网关上:发给 DeepSeek 官方原样透传,发给别家先清理。 +//! +//! 清理本身在 `tw_dialect::harness` 里有单元测试;这里验的是网关把它接对了: +//! +//! - 同格式直通和格式转换两条路,发给别家的请求里都没有 dsh 的扩展字段、中途的 +//! system 条目和增删工具的块,也不带它自己的请求头;丢掉了什么记在这一跳上 +//! - 发给 DeepSeek 官方时一个字节都不改,请求头照发;要转换格式时扩展字段照样带上 +//! - 开始事件说得出请求带没带会话日志、有多大 +//! - `/v1/files` 回 404,dsh 才会退回内联图片 +//! +//! 「DeepSeek 官方」认的是上游地址的 host,测试里的假上游在 127.0.0.1 上。所以那几条 +//! 把上游写成 `http://api.deepseek.com`,再让这家走一个 HTTP 代理,代理就是假上游: +//! 它收到的正是发给 DeepSeek 的那一份。 + +use std::net::SocketAddr; +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +use axum::Router; +use axum::extract::{OriginalUri, State}; +use axum::http::HeaderMap; +use serde_json::{Value, json}; +use tw_config::{Client, Config, Listen, Protocol, Provider}; + +const UA: &str = "deepseek-harness/0.1.7 (+https://github.com/deepseek-ai/deepseek-harness)"; + +#[derive(Debug, Clone)] +struct Seen { + uri: String, + headers: HeaderMap, + body: Vec, +} + +type Log = Arc>>; + +/// 一个假上游:记下收到的每个请求,按上游的格式回一个整包。也可以当 HTTP 代理用 +async fn upstream(protocol: Protocol) -> (SocketAddr, Log) { + let reply = match protocol { + Protocol::Anthropic => json!({ + "id": "msg_1", "type": "message", "role": "assistant", "model": "m", + "content": [{"type": "text", "text": "ok"}], + "stop_reason": "end_turn", + "usage": {"input_tokens": 10, "output_tokens": 2} + }), + _ => json!({ + "id": "chatcmpl-1", "object": "chat.completion", "model": "m", + "choices": [{"index": 0, "message": {"role": "assistant", "content": "ok"}, "finish_reason": "stop"}], + "usage": {"prompt_tokens": 10, "completion_tokens": 2} + }), + }; + let log: Log = Arc::new(Mutex::new(Vec::new())); + let app = Router::new() + .fallback( + move |State(log): State, + OriginalUri(uri): OriginalUri, + headers: HeaderMap, + body: bytes::Bytes| { + let reply = reply.clone(); + async move { + log.lock().unwrap().push(Seen { + uri: uri.to_string(), + headers, + body: body.to_vec(), + }); + axum::response::Response::builder() + .status(200) + .header("content-type", "application/json") + .body(axum::body::Body::from(reply.to_string())) + .unwrap() + } + }, + ) + .with_state(log.clone()); + let l = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = l.local_addr().unwrap(); + tokio::spawn(async move { axum::serve(l, app).await.unwrap() }); + (addr, log) +} + +/// 上游收到的那个生成请求(启动时也许还有别的请求,比如拉模型列表) +fn generation(log: &Log) -> Seen { + log.lock() + .unwrap() + .iter() + .rev() + .find(|s| s.uri.ends_with("/messages") || s.uri.ends_with("/chat/completions")) + .cloned() + .expect("the upstream got no generation request") +} + +/// 一家别的上游(不是 DeepSeek 官方) +async fn elsewhere(protocol: Protocol) -> (Provider, Vec, Log) { + let (addr, log) = upstream(protocol).await; + let p = Provider { + name: "relay".into(), + base_url: format!("http://{addr}"), + key: Some("sk-upstream".into()), + protocol: Some(protocol), + ..Default::default() + }; + (p, Vec::new(), log) +} + +/// DeepSeek 官方:地址是它的,请求经代理落到假上游上 +async fn deepseek(protocol: Protocol) -> (Provider, Vec, Log) { + let (addr, log) = upstream(protocol).await; + let base_url = match protocol { + Protocol::Anthropic => "http://api.deepseek.com/anthropic", + _ => "http://api.deepseek.com", + }; + let p = Provider { + name: "deepseek".into(), + base_url: base_url.into(), + key: Some("sk-upstream".into()), + protocol: Some(protocol), + proxy: "fake".into(), + ..Default::default() + }; + let proxies = vec![tw_config::Proxy { + name: "fake".into(), + kind: tw_config::ProxyKind::Http, + addr: addr.to_string(), + auth: None, + }]; + (p, proxies, log) +} + +async fn gateway( + (p, proxies, log): (Provider, Vec, Log), +) -> ( + SocketAddr, + tokio::sync::broadcast::Receiver, + Log, +) { + let cfg = Config { + version: 1, + listen: Listen::default(), + clients: vec![Client { + name: "c".into(), + key: "tw-k".into(), + ..Default::default() + }], + providers: vec![p], + proxies, + ..Default::default() + }; + let state = tw_gateway::AppState::new(cfg).unwrap(); + let rx = state.bus.subscribe(); + let addr = tw_gateway::serve(state, ([127, 0, 0, 1], 0).into()) + .await + .unwrap(); + tokio::time::sleep(Duration::from_millis(50)).await; + (addr, rx, log) +} + +/// dsh 发请求的样子:它的 User-Agent 和自己的三个头 +async fn send(gw: SocketAddr, path: &str, extra: &[(&str, &str)], body: &str) -> (u16, String) { + let mut r = reqwest::Client::new() + .post(format!("http://{gw}{path}")) + .header("content-type", "application/json") + .header("authorization", "Bearer tw-k") + .header("user-agent", UA) + .header( + "x-deepseek-harness-user-id", + "0b0c6a52-2f0e-4c1a-9d5e-3a1f5e2b7c11", + ) + .header("x-deepseek-harness-session-id", "s-1") + .header("x-deepseek-harness-compact", "1"); + for (k, v) in extra { + r = r.header(*k, *v); + } + let resp = r.body(body.to_string()).send().await.unwrap(); + let status = resp.status().as_u16(); + (status, resp.text().await.unwrap()) +} + +const TOOL_BETA: &str = "mid-conversation-tool-changes-2026-07-01"; + +fn session_log() -> Value { + json!({ + "version": 1, "sessionFormatVersion": 2, + "session": {"version": 2, "id": "s-1", "createdAt": 1780000000000u64, "cwd": "/work"}, + "afterSeq": -1, "throughSeq": 1, + "events": [ + {"type": "turn/start", "seq": 0, "time": 1780000000001u64, "data": {"turn": 1}}, + {"type": "user", "seq": 1, "time": 1780000000002u64, "data": {"text": "hi"}} + ] + }) +} + +/// dsh 0.1.7:Anthropic Messages,带会话日志、中途的系统提示和增删工具 +fn messages_request() -> Value { + json!({ + "model": "deepseek-v4-pro", + "max_tokens": 32000, + "system": "You are DeepSeek Harness.", + "thinking": {"type": "enabled"}, + "output_config": {"effort": "high"}, + "tools": [ + {"name": "read", "description": "Read a file", "input_schema": {"type": "object"}}, + {"name": "web", "description": "Search the web", "input_schema": {"type": "object"}, "defer_loading": true} + ], + "messages": [ + {"role": "user", "content": [{"type": "text", "text": "hi"}]}, + {"role": "system", "content": [ + {"type": "text", "text": "The project root is /work."}, + {"type": "tool_addition", "tool": {"type": "tool_reference", "name": "web"}} + ]}, + {"role": "assistant", "content": [{"type": "text", "text": "hello"}]}, + {"role": "user", "content": [{"type": "text", "text": "search it"}]}, + {"role": "system", "content": [ + {"type": "tool_removal", "tool": {"type": "tool_reference", "name": "web"}} + ]} + ], + "dsh_session_log": session_log(), + "dsh_plugin_packages": {"version": 1, "packages": [{"name": "@deepseek-ai/dsh-example", "version": "0.1.1"}]} + }) +} + +/// dsh 0.1.5:Chat Completions,带会话日志 +fn chat_request() -> Value { + json!({ + "model": "deepseek-v4-flash", + "messages": [ + {"role": "system", "content": "You are DeepSeek Harness."}, + {"role": "user", "content": "hi"}, + {"role": "assistant", "content": "hello", "reasoning_content": "greet"}, + {"role": "system", "content": "The project root is /work."}, + {"role": "user", "content": "search it"} + ], + "thinking": {"type": "enabled"}, + "reasoning_effort": "high", + "dsh_session_log": session_log(), + "dsh_plugin_packages": {"version": 1, "packages": []} + }) +} + +/// 发给别家的那一份里不能再有的东西 +fn assert_clean(s: &Seen) { + let v: Value = serde_json::from_slice(&s.body).unwrap(); + let o = v.as_object().unwrap(); + assert!(!o.keys().any(|k| k.starts_with("dsh_")), "{v}"); + let text = String::from_utf8_lossy(&s.body); + assert!( + !text.contains("tool_addition") && !text.contains("tool_removal"), + "{v}" + ); + assert!(!text.contains("defer_loading"), "{v}"); + for (name, _) in s.headers.iter() { + assert!( + !name.as_str().starts_with("x-deepseek-harness-"), + "{name} was forwarded" + ); + } + if let Some(b) = s.headers.get("anthropic-beta") { + assert!(!b.to_str().unwrap().contains("mid-conversation"), "{b:?}"); + } +} + +/// 这个请求的开始事件里的会话日志大小,和这一跳丢掉的字段 +async fn events( + rx: &mut tokio::sync::broadcast::Receiver, +) -> ( + Option, + Option<(tw_api::Dialect, tw_api::Dialect, Vec)>, +) { + let mut log = None; + let mut translated = None; + while let Ok(Ok(ev)) = tokio::time::timeout(Duration::from_secs(3), rx.recv()).await { + match ev { + tw_api::Event::RequestStarted { + session_log_bytes, .. + } => log = session_log_bytes, + tw_api::Event::Translated { + from, to, dropped, .. + } => translated = Some((from, to, dropped)), + tw_api::Event::RequestFinished { .. } | tw_api::Event::RequestFailed { .. } => break, + _ => {} + } + } + (log, translated) +} + +fn log_len() -> u64 { + serde_json::to_vec(&session_log()).unwrap().len() as u64 +} + +// ───────────────────────────────────────────────────────── dsh 0.1.7(Anthropic) + +#[tokio::test] +async fn messages_to_another_anthropic_upstream_are_cleaned_in_place() { + let (gw, mut rx, log) = gateway(elsewhere(Protocol::Anthropic).await).await; + let beta = format!("files-api-2025-04-14,{TOOL_BETA}"); + let (status, body) = send( + gw, + "/v1/messages", + &[ + ("anthropic-version", "2023-06-01"), + ("anthropic-beta", &beta), + ], + &messages_request().to_string(), + ) + .await; + assert_eq!(status, 200, "{body}"); + + let s = generation(&log); + assert_eq!(s.uri, "/v1/messages"); + assert_clean(&s); + // 别的 beta 照发,User-Agent 照发 + assert_eq!( + s.headers.get("anthropic-beta").unwrap(), + "files-api-2025-04-14" + ); + assert_eq!(s.headers.get("user-agent").unwrap(), UA); + let v: Value = serde_json::from_slice(&s.body).unwrap(); + assert_eq!( + v["system"], + "You are DeepSeek Harness.\n\nThe project root is /work." + ); + let roles: Vec<&str> = v["messages"] + .as_array() + .unwrap() + .iter() + .map(|m| m["role"].as_str().unwrap()) + .collect(); + assert_eq!(roles, ["user", "assistant", "user"]); + // Anthropic 的 `enabled` 要带预算 + assert_eq!(v["thinking"]["type"], "enabled"); + assert!(v["thinking"]["budget_tokens"].as_u64().unwrap() >= 1024); + + let (bytes, translated) = events(&mut rx).await; + assert_eq!(bytes, Some(log_len())); + let (from, to, dropped) = translated.expect("what was dropped is said"); + assert_eq!( + (from, to), + (tw_api::Dialect::Anthropic, tw_api::Dialect::Anthropic) + ); + assert_eq!( + dropped, + [ + "messages.content.tool_addition", + "messages.content.tool_removal", + "tools.defer_loading" + ] + ); +} + +#[tokio::test] +async fn messages_to_an_openai_upstream_are_cleaned_by_the_conversion() { + let (gw, mut rx, log) = gateway(elsewhere(Protocol::OpenaiChat).await).await; + let (status, body) = send( + gw, + "/v1/messages", + &[ + ("anthropic-version", "2023-06-01"), + ("anthropic-beta", TOOL_BETA), + ], + &messages_request().to_string(), + ) + .await; + assert_eq!(status, 200, "{body}"); + let answer: Value = serde_json::from_str(&body).unwrap(); + assert_eq!(answer["content"][0]["text"], "ok"); + + let s = generation(&log); + assert_eq!(s.uri, "/v1/chat/completions"); + 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() + .iter() + .filter(|m| m["role"] == "system") + .map(|m| m["content"].as_str().unwrap()) + .collect(); + assert_eq!( + system + .concat() + .matches("The project root is /work.") + .count(), + 1 + ); + 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; + assert_eq!(bytes, Some(log_len())); + let (from, to, dropped) = translated.unwrap(); + assert_eq!( + (from, to), + (tw_api::Dialect::Anthropic, tw_api::Dialect::OpenaiChat) + ); + for d in [ + "messages.content.tool_addition", + "messages.content.tool_removal", + "tools.defer_loading", + ] { + assert!(dropped.iter().any(|x| x == d), "{d} not in {dropped:?}"); + } +} + +#[tokio::test] +async fn messages_to_deepseek_go_through_untouched() { + let (gw, mut rx, log) = gateway(deepseek(Protocol::Anthropic).await).await; + let body = messages_request().to_string(); + let (status, answer) = send( + gw, + "/v1/messages", + &[ + ("anthropic-version", "2023-06-01"), + ("anthropic-beta", TOOL_BETA), + ], + &body, + ) + .await; + assert_eq!(status, 200, "{answer}"); + + let s = generation(&log); + assert_eq!(s.uri, "http://api.deepseek.com/anthropic/v1/messages"); + assert_eq!(s.body, body.as_bytes(), "a single byte changed"); + assert_eq!( + s.headers.get("x-deepseek-harness-session-id").unwrap(), + "s-1" + ); + assert_eq!(s.headers.get("x-deepseek-harness-compact").unwrap(), "1"); + assert_eq!(s.headers.get("anthropic-beta").unwrap(), TOOL_BETA); + + let (bytes, translated) = events(&mut rx).await; + assert_eq!(bytes, Some(log_len()), "the log is flagged whoever gets it"); + assert!(translated.is_none(), "{translated:?}"); +} + +// ───────────────────────────────────────────────────────── dsh 0.1.5(Chat) + +#[tokio::test] +async fn chat_to_another_openai_upstream_loses_only_the_extensions() { + let (gw, mut rx, log) = gateway(elsewhere(Protocol::OpenaiChat).await).await; + let (status, body) = send(gw, "/chat/completions", &[], &chat_request().to_string()).await; + assert_eq!(status, 200, "{body}"); + + let s = generation(&log); + assert_eq!(s.uri, "/chat/completions"); + assert_clean(&s); + let v: Value = serde_json::from_slice(&s.body).unwrap(); + // Chat 本来就有 system 角色:messages 原样 + assert_eq!(v["messages"], chat_request()["messages"]); + assert_eq!(v["reasoning_effort"], "high"); + + let (bytes, translated) = events(&mut rx).await; + assert_eq!(bytes, Some(log_len())); + // 只去掉了扩展字段:没有转不过去的东西要报 + assert!(translated.is_none(), "{translated:?}"); +} + +#[tokio::test] +async fn chat_to_an_anthropic_upstream_is_cleaned_by_the_conversion() { + let (gw, mut rx, log) = gateway(elsewhere(Protocol::Anthropic).await).await; + let (status, body) = send(gw, "/v1/chat/completions", &[], &chat_request().to_string()).await; + assert_eq!(status, 200, "{body}"); + let answer: Value = serde_json::from_str(&body).unwrap(); + assert_eq!(answer["choices"][0]["message"]["content"], "ok"); + + let s = generation(&log); + assert_eq!(s.uri, "/v1/messages"); + assert_clean(&s); + let v: Value = serde_json::from_slice(&s.body).unwrap(); + let system: Vec<&str> = v["system"] + .as_array() + .unwrap() + .iter() + .map(|b| b["text"].as_str().unwrap()) + .collect(); + assert_eq!( + system, + ["You are DeepSeek Harness.", "The project root is /work."] + ); + assert!( + v["messages"] + .as_array() + .unwrap() + .iter() + .all(|m| m["role"] != "system") + ); + assert_eq!(v["thinking"]["type"], "enabled"); + + let (bytes, translated) = events(&mut rx).await; + assert_eq!(bytes, Some(log_len())); + let (from, to, _) = translated.unwrap(); + assert_eq!( + (from, to), + (tw_api::Dialect::OpenaiChat, tw_api::Dialect::Anthropic) + ); +} + +#[tokio::test] +async fn chat_to_deepseek_goes_through_untouched() { + let (gw, _rx, log) = gateway(deepseek(Protocol::OpenaiChat).await).await; + let body = chat_request().to_string(); + let (status, answer) = send(gw, "/v1/chat/completions", &[], &body).await; + assert_eq!(status, 200, "{answer}"); + + let s = generation(&log); + assert_eq!(s.uri, "http://api.deepseek.com/v1/chat/completions"); + assert_eq!(s.body, body.as_bytes(), "a single byte changed"); + assert_eq!( + s.headers.get("x-deepseek-harness-user-id").unwrap(), + "0b0c6a52-2f0e-4c1a-9d5e-3a1f5e2b7c11" + ); +} + +#[tokio::test] +async fn converted_for_deepseek_the_extensions_ride_along() { + // 客户端说 Anthropic、DeepSeek 这家配成 Chat:格式要转,扩展照样带给 DeepSeek + let (gw, _rx, log) = gateway(deepseek(Protocol::OpenaiChat).await).await; + let (status, answer) = send( + gw, + "/v1/messages", + &[("anthropic-version", "2023-06-01")], + &messages_request().to_string(), + ) + .await; + assert_eq!(status, 200, "{answer}"); + + let s = generation(&log); + assert_eq!(s.uri, "http://api.deepseek.com/v1/chat/completions"); + let v: Value = serde_json::from_slice(&s.body).unwrap(); + assert_eq!(v["dsh_session_log"], session_log()); + assert_eq!(v["dsh_plugin_packages"]["version"], 1); + assert_eq!( + s.headers.get("x-deepseek-harness-session-id").unwrap(), + "s-1" + ); +} + +// ───────────────────────────────────────────────────────── 别的请求 + +#[tokio::test] +async fn a_request_without_the_harness_is_not_touched() { + let (gw, mut rx, log) = gateway(elsewhere(Protocol::Anthropic).await).await; + let body = json!({ + "model": "claude-sonnet-4-5", + "max_tokens": 100, + "thinking": {"type": "enabled", "budget_tokens": 2048}, + "messages": [{"role": "user", "content": "hi"}] + }) + .to_string(); + let resp = reqwest::Client::new() + .post(format!("http://{gw}/v1/messages")) + .header("x-api-key", "tw-k") + .header("anthropic-version", "2023-06-01") + .header("user-agent", "claude-cli/2.1.195 (external, cli)") + .body(body.clone()) + .send() + .await + .unwrap(); + assert_eq!(resp.status(), 200); + assert_eq!(generation(&log).body, body.as_bytes()); + let (bytes, translated) = events(&mut rx).await; + assert_eq!(bytes, None); + assert!(translated.is_none()); +} + +#[tokio::test] +async fn the_files_api_is_not_found_so_images_go_inline() { + let (gw, _rx, log) = gateway(elsewhere(Protocol::Anthropic).await).await; + let c = reqwest::Client::new(); + for (path, anthropic) in [ + ("/v1/files", true), + ("/v1/files", false), + ("/files", false), + ("/v1/files/file-abc", true), + ("/files?limit=10", false), + ] { + let mut r = c + .post(format!("http://{gw}{path}")) + .header("authorization", "Bearer tw-k") + .header("user-agent", UA) + .body("--x\r\n"); + if anthropic { + r = r.header("anthropic-version", "2023-06-01"); + } + let resp = r.send().await.unwrap(); + assert_eq!(resp.status(), 404, "{path}"); + let v: Value = resp.json().await.unwrap(); + // 两种形状都把原因放在 `error.type` / `error.message`,Anthropic 的外面还有一层 `type` + assert_eq!(v["error"]["type"], "not_found_error", "{path}: {v}"); + let message = v["error"]["message"].as_str().unwrap(); + assert!(message.starts_with("[ThinkWatch]"), "{v}"); + assert_eq!(v.get("type").is_some(), anthropic, "{path}: {v}"); + } + let got = c + .get(format!("http://{gw}/v1/files")) + .header("authorization", "Bearer tw-k") + .send() + .await + .unwrap(); + assert_eq!(got.status(), 404); + assert!( + log.lock().unwrap().iter().all(|s| !s.uri.contains("files")), + "the upload reached an upstream" + ); +} diff --git a/crates/tw-observe/src/bus.rs b/crates/tw-observe/src/bus.rs index 106958c6..34865766 100644 --- a/crates/tw-observe/src/bus.rs +++ b/crates/tw-observe/src/bus.rs @@ -352,6 +352,7 @@ mod tests { model: "m".into(), method: "POST".into(), path: "/v1/messages".into(), + session_log_bytes: None, at_ms: 1_000 + id, } } diff --git a/crates/tw-store/src/db.rs b/crates/tw-store/src/db.rs index 55861f3d..55427fcc 100644 --- a/crates/tw-store/src/db.rs +++ b/crates/tw-store/src/db.rs @@ -22,7 +22,7 @@ use tw_api::Msg; /// /// **一列 JSON 的样子变了也算**(比如 `routing` 多了必有的字段):旧的那些行 /// 读出来是坏的,而读的一方会把「解不开」当成「没有」。 -const SCHEMA: i64 = 20; +const SCHEMA: i64 = 21; /// 这一行算不出钱,**因为价目表里没有这个模型**:用量是有的,缺的是单价。 /// @@ -119,6 +119,8 @@ pub struct RequestRow { pub price_source: Option, /// 服务它的那一跳做过的格式转换,JSON。直通的是 None pub translated: Option, + /// 请求带着的 DeepSeek Harness 会话日志有多少字节。没带是 None + pub session_log_bytes: Option, } #[derive(Debug)] @@ -257,7 +259,10 @@ impl Db { error_args TEXT, local INTEGER NOT NULL, -- 客户端没等到响应结束就走了。**不能写进 `error`**:它不算失败 - cancelled INTEGER NOT NULL + cancelled INTEGER NOT NULL, + -- 请求带着 DeepSeek Harness 的会话日志:它的字节数。整段对话都在里面, + -- 界面要说得出哪些请求带着它 + session_log_bytes INTEGER ); -- 几乎所有查询都是「最近的 N 条」或者「某段时间内的」 CREATE INDEX requests_at ON requests (at_ms DESC); @@ -296,8 +301,8 @@ impl Db { input_tokens, output_tokens, cache_read_tokens, cache_write_tokens, cost_micros, cost_estimated, error, local, routing, billing, cache_saved_micros, client_hint, session, cancelled, price_source, translated, - error_code, error_args, peer, key_masked) - VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13,?14,?15,?16,?17,?18,?19,?20,?21,?22,?23,?24,?25,?26,?27,?28,?29,?30)", + error_code, error_args, peer, key_masked, session_log_bytes) + VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13,?14,?15,?16,?17,?18,?19,?20,?21,?22,?23,?24,?25,?26,?27,?28,?29,?30,?31)", params![ r.id, r.at_ms, @@ -332,6 +337,7 @@ impl Db { .map(|e| serde_json::to_string(&e.args).unwrap_or_default()), r.peer, r.key_masked, + r.session_log_bytes, ], )?; Ok(()) @@ -1177,6 +1183,7 @@ fn row_from(r: &rusqlite::Row) -> rusqlite::Result { client_hint: r.get("client_hint")?, peer: r.get("peer")?, key_masked: r.get("key_masked")?, + session_log_bytes: r.get("session_log_bytes")?, session: r.get("session")?, provider: r.get("provider")?, model: r.get("model")?, @@ -1354,6 +1361,7 @@ pub(crate) mod tests { billing: tw_api::Billing::PerToken, cache_saved_micros: None, price_source: None, + session_log_bytes: None, translated: None, } } diff --git a/crates/tw-store/src/recorder.rs b/crates/tw-store/src/recorder.rs index 7003b0e6..916462d8 100644 --- a/crates/tw-store/src/recorder.rs +++ b/crates/tw-store/src/recorder.rs @@ -38,6 +38,8 @@ struct Partial { /// 做过的格式转换,JSON,带着做转换的那一家。**只留服务它的那一跳的**: /// 故障转移前一跳转换过、后一跳直通时,这一行不该说它转换过 translated: Option, + /// 请求带着的 DeepSeek Harness 会话日志的字节数。开始事件带着 + session_log_bytes: Option, } /// 在飞的请求最多攒多少条。 @@ -156,6 +158,7 @@ impl Recorder { cache_saved_micros: None, price_source: None, translated: p.translated.clone(), + session_log_bytes: p.session_log_bytes, }) } @@ -190,6 +193,7 @@ impl Recorder { billing, model, path, + session_log_bytes, at_ms, .. } => { @@ -222,6 +226,7 @@ impl Recorder { }, billing: *billing, translated: None, + session_log_bytes: session_log_bytes.map(|b| b as i64), }, ); } @@ -548,6 +553,7 @@ impl Recorder { cache_saved_micros: None, price_source: None, translated: None, + session_log_bytes: None, }); } // 配置事件、额度事件、扫描告警都不是请求,不落这张表。 @@ -687,6 +693,7 @@ impl Recorder { cache_saved_micros, price_source, translated: p.translated, + session_log_bytes: p.session_log_bytes, }); } @@ -748,6 +755,7 @@ mod tests { model: model.into(), method: "POST".into(), path: "/v1/messages".into(), + session_log_bytes: None, at_ms: 1_000_000, } } @@ -1243,6 +1251,7 @@ mod billing_tests { model: String::new(), method: "WS".into(), path: "/backend-api/codex/responses".into(), + session_log_bytes: None, at_ms: 1_000_000, } } @@ -1338,6 +1347,7 @@ mod billing_tests { model: "claude-sonnet-4-5".into(), method: "POST".into(), path: "/v1/messages".into(), + session_log_bytes: None, at_ms: 1_000_000, }); r.on_event(&Event::RequestCancelled { @@ -1960,6 +1970,7 @@ mod cancellation_tests { model: "claude-sonnet-4-5".into(), method: "POST".into(), path: "/v1/messages".into(), + session_log_bytes: None, at_ms: 1_000_000, }); r.on_event(&cancelled(1, partial())); diff --git a/crates/tw-store/tests/m3_acceptance.rs b/crates/tw-store/tests/m3_acceptance.rs index 2903d3ba..5a803774 100644 --- a/crates/tw-store/tests/m3_acceptance.rs +++ b/crates/tw-store/tests/m3_acceptance.rs @@ -25,6 +25,7 @@ const NOW: i64 = 1_800_000_000_000; fn req(id: i64, at_ms: i64) -> RequestRow { RequestRow { + session_log_bytes: None, key_masked: None, peer: None, id,