From 6623785d08efbe0fe2265657202037d8205f6678 Mon Sep 17 00:00:00 2001 From: fylorn <249551762+fylorn@users.noreply.github.com> Date: Fri, 9 Oct 2026 22:22:23 +0800 Subject: [PATCH] Mask secrets inside carried compaction summaries; fall back when automatic cache breakpoints are refused Compaction summaries at rest. The summary we hand Codex as a compaction item is `tw1.c.` plus base64url text, opaque to rule-based redaction and to shape-based masking. Codex sends it back every turn, so a secret mentioned in the summary could reach the stored request body. Before a body is written, each decodable carried summary is now decoded, masked with the same rules, ledger placeholders and shape masking, re-encoded and put back in place. The rest of the body is masked as before, still reusing the hits found on the request path. Damaged or truncated values are masked as ordinary text. Automatic cache breakpoints. - Bedrock: automatic `cachePoint`s now go only to Claude models AWS lists for prompt caching: every Claude from 4.5 on (unlisted newer ids included), plus 3.7 Sonnet and 3.5 Sonnet v2. Matched by family and version in the id, so cross-region inference profile ids and ARNs pointing at them are recognized. Sonnet 4, Opus 4/4.1, 3.5 Haiku and Claude 3 get none. - Safety net for Anthropic-format and Bedrock upstreams: if the request carries only automatic breakpoints and the upstream answers 400 mentioning `cache_control` / `cachePoint` / prompt caching, the hop is sent once more without them. That upstream and model is remembered (in memory, bounded) so later requests leave them out. Breakpoints the client set itself are never removed. - Speed probes no longer carry automatic breakpoints. New `tw_dialect::cache::{may_have_marks, strip_marks}`; docs updated. Co-Authored-By: Claude Opus 5.5 --- crates/tw-dialect/src/anthropic/request.rs | 4 +- crates/tw-dialect/src/bedrock/request.rs | 122 +++++++++++- crates/tw-dialect/src/cache.rs | 132 +++++++++++++ crates/tw-dialect/src/lib.rs | 5 +- crates/tw-gateway/src/bodies.rs | 191 ++++++++++++++++++- crates/tw-gateway/src/cache_marks.rs | 113 +++++++++++ crates/tw-gateway/src/l3.rs | 5 +- crates/tw-gateway/src/lib.rs | 1 + crates/tw-gateway/src/server/pipeline/hop.rs | 97 +++++++++- crates/tw-gateway/src/state.rs | 4 + crates/tw-gateway/tests/bedrock.rs | 81 ++++++++ crates/tw-gateway/tests/conversion.rs | 120 ++++++++++++ docs/config.md | 7 + docs/config.zh-CN.md | 2 + 14 files changed, 867 insertions(+), 17 deletions(-) create mode 100644 crates/tw-dialect/src/cache.rs create mode 100644 crates/tw-gateway/src/cache_marks.rs diff --git a/crates/tw-dialect/src/anthropic/request.rs b/crates/tw-dialect/src/anthropic/request.rs index b25bd6b7..12d9855f 100644 --- a/crates/tw-dialect/src/anthropic/request.rs +++ b/crates/tw-dialect/src/anthropic/request.rs @@ -222,7 +222,9 @@ fn cache_ttl(b: &Value) -> Option { /// /// 断点不算内容:上一轮标在更早位置的那个这一轮挪走了,缓存照样命中(Anthropic 文档里多轮 /// 对话就是这么标的)。太短(不到模型的最小缓存长度)的断点上游直接忽略,不报错,所以不估 -/// 长度。思考块不能标,标在它前面的那块上 +/// 长度。思考块不能标,标在它前面的那块上。 +/// +/// 不收 `cache_control` 的兼容接口回 400 时,网关去掉这些断点再发一次([`crate::cache`]) fn auto_cache(out: &mut Map) { fn mark(blocks: Option<&mut Value>) { let Some(Value::Array(blocks)) = blocks else { diff --git a/crates/tw-dialect/src/bedrock/request.rs b/crates/tw-dialect/src/bedrock/request.rs index 5cab8b27..a93ac8d7 100644 --- a/crates/tw-dialect/src/bedrock/request.rs +++ b/crates/tw-dialect/src/bedrock/request.rs @@ -80,6 +80,67 @@ fn caches(model: &str) -> bool { m.contains("anthropic.claude") || m.contains("amazon.nova") || m.starts_with("arn:") } +/// 客户端没标断点时,替它标不标:**只给 AWS 列出支持显式提示缓存的 Claude**。 +/// +/// 依据是 AWS 的「Supported models, Regions, and explicit caching limits」表 +/// (,2026-10-09 +/// 读的):Claude 4.5 起的每一个(Haiku 4.5、Sonnet 4.5 / 4.6 / 5 / 5.5、Opus 4.5 ~ 4.8 / 5 / +/// 5.5、Fable、Mythos),加上更早的 Claude 3.7 Sonnet 和 Claude 3.5 Sonnet v2(`20241022`)。 +/// 表里没有的 —— Claude 3 Haiku / Sonnet / Opus、3.5 Sonnet v1、3.5 Haiku、Sonnet 4、 +/// Opus 4 / 4.1 —— 不自动标:不认的模型收到 `cachePoint` 会拒掉整个请求。 +/// +/// **按 id 里的家族和版本认**,不按完整的 id:跨区域推理配置的前缀(`us.`、`eu.`、 +/// `apac.`、`global.`……)、日期和 `-v1:0` 这些尾巴都不影响。表里还没有的新 id 按版本 +/// 判:4.5 起的家族都支持,以后的也算。看不出背后是谁的应用推理配置 ARN 不标。 +/// +/// 客户端自己标的断点不受这条管(见 [`caches`]):那是客户端的决定 +fn caches_automatically(model: &str) -> bool { + let m = model.to_ascii_lowercase(); + // ARN 的最后一段才是模型(推理配置 ARN 里写着 `us.anthropic.claude-…`) + let tail = m.rsplit('/').next().unwrap_or(&m); + let Some(at) = tail.find("anthropic.claude-") else { + return false; + }; + let parts: Vec<&str> = tail[at + "anthropic.claude-".len()..] + .split(['-', ':']) + .collect(); + // 版本号是一两位的数;八位的是日期 + let num = |s: Option<&&str>| { + s.filter(|s| (1..=2).contains(&s.len())) + .and_then(|s| s.parse::().ok()) + }; + let (family, major, minor, rest) = match parts.first() { + // 新的写法:`claude-sonnet-4-5-20250929-v1:0`、`claude-opus-5` + Some(f) if f.chars().all(|c| c.is_ascii_alphabetic()) => { + let Some(major) = num(parts.get(1)) else { + return false; + }; + match num(parts.get(2)) { + Some(minor) => (*f, major, minor, &parts[3..]), + None => (*f, major, 0, &parts[2..]), + } + } + // 老的写法:`claude-3-7-sonnet-20250219-v1:0`、`claude-3-haiku-20240307-v1:0` + Some(_) => { + let Some(major) = num(parts.first()) else { + return false; + }; + let (minor, at) = match num(parts.get(1)) { + Some(minor) => (minor, 2), + None => (0, 1), + }; + let Some(family) = parts.get(at) else { + return false; + }; + (*family, major, minor, &parts[at + 1..]) + } + None => return false, + }; + (major, minor) >= (4, 5) + || (family == "sonnet" && (major, minor) == (3, 7)) + || (family == "sonnet" && (major, minor) == (3, 5) && rest.contains(&"20241022")) +} + /// 客户端没标断点时自动标的,和转给 Anthropic 时同样的四处:工具的末尾、系统提示的末尾、 /// 最后两条用户消息的末尾,各跟一个 5 分钟的 `cachePoint`(Converse 也最多四个) fn auto_cache(out: &mut Map) { @@ -231,10 +292,10 @@ pub fn encode_request(r: &Request, t: &Target, dropped: &mut Dropped) -> Value { } else if r.tool_choice.is_some() { dropped.path("tool_choice"); } - // 客户端自己一个断点都没标,模型又是 Claude:替它标(见 `anthropic::request` 的 - // `auto_cache`)。Nova 和看不出是谁的推理配置 ARN 不自动标:不认的模型收到 - // `cachePoint` 会拒掉整个请求 - if r.cache.is_empty() && is_claude(&r.model) { + // 客户端自己一个断点都没标,模型又是 AWS 列出支持缓存的 Claude:替它标(见 + // `anthropic::request` 的 `auto_cache`、[`caches_automatically`])。Nova 和看不出是谁的 + // 推理配置 ARN 不自动标:不认的模型收到 `cachePoint` 会拒掉整个请求 + if r.cache.is_empty() && caches_automatically(&r.model) { auto_cache(&mut out); } @@ -1024,6 +1085,59 @@ mod tests { (v, d.into_vec()) } + #[test] + fn only_the_claude_models_aws_lists_get_automatic_cache_points() { + for (model, yes) in [ + // AWS 的表里有的 + ("anthropic.claude-haiku-5-5", true), + ("anthropic.claude-sonnet-5-5", true), + ("anthropic.claude-opus-5-5", true), + ("anthropic.claude-fable-5-1", true), + ("anthropic.claude-mythos-5", true), + ("anthropic.claude-opus-5", true), + ("anthropic.claude-opus-4-8", true), + ("anthropic.claude-opus-4-7", true), + ("anthropic.claude-opus-4-6-v1", true), + ("anthropic.claude-opus-4-5-20251101-v1:0", true), + ("anthropic.claude-sonnet-5", true), + ("anthropic.claude-sonnet-4-6", true), + ("anthropic.claude-sonnet-4-5-20250929-v1:0", true), + ("anthropic.claude-haiku-4-5-20251001-v1:0", true), + ("anthropic.claude-3-7-sonnet-20250219-v1:0", true), + ("anthropic.claude-3-5-sonnet-20241022-v2:0", true), + // 跨区域推理配置,和指着它的 ARN + ("us.anthropic.claude-sonnet-4-5-20250929-v1:0", true), + ("global.anthropic.claude-opus-4-7", true), + ("apac.anthropic.claude-3-7-sonnet-20250219-v1:0", true), + ( + "arn:aws:bedrock:us-east-1:123456789012:inference-profile/eu.anthropic.claude-haiku-4-5-20251001-v1:0", + true, + ), + // 表里还没有的新 id:按版本 + ("anthropic.claude-opus-6", true), + ("us.anthropic.claude-sonnet-5-7-20270101-v1:0", true), + // 表里没有的 + ("anthropic.claude-sonnet-4-20250514-v1:0", false), + ("us.anthropic.claude-opus-4-1-20250805-v1:0", false), + ("anthropic.claude-opus-4-20250514-v1:0", false), + ("anthropic.claude-3-5-sonnet-20240620-v1:0", false), + ("anthropic.claude-3-5-haiku-20241022-v1:0", false), + ("anthropic.claude-3-haiku-20240307-v1:0", false), + ("anthropic.claude-3-opus-20240229-v1:0", false), + ("anthropic.claude-v2:1", false), + ("anthropic.claude-instant-v1", false), + // 别家的模型、看不出是谁的 ARN + ("amazon.nova-pro-v1:0", false), + ("meta.llama3-70b-instruct-v1:0", false), + ( + "arn:aws:bedrock:us-east-1:123456789012:application-inference-profile/a1b2c3", + false, + ), + ] { + assert_eq!(caches_automatically(model), yes, "{model}"); + } + } + #[test] fn cache_points_follow_the_blocks_the_client_marked() { let r = claude_code_like("us.anthropic.claude-sonnet-4-5-20250929-v1:0"); diff --git a/crates/tw-dialect/src/cache.rs b/crates/tw-dialect/src/cache.rs new file mode 100644 index 00000000..a7439f43 --- /dev/null +++ b/crates/tw-dialect/src/cache.rs @@ -0,0 +1,132 @@ +//! 提示缓存的断点,在编码好的请求体上。 +//! +//! 转给 Claude 时客户端没标断点的,编码器替它标(`anthropic::request`、`bedrock::request` +//! 里的 `auto_cache`)。有的上游不认 —— 自称 Anthropic 格式的兼容接口、AWS 没列进支持表 +//! 的老模型 —— 拒掉整个请求。网关这时去掉断点、同一家再发一次(见 `tw_gateway` 的 +//! `cache_marks`),靠的就是这里。 +//! +//! **只用在客户端自己一个断点都没标的请求上**:那时请求体里的断点全是编码器加的。客户端 +//! 自己标的是它的决定,去掉了它的缓存就白建了。 + +use serde_json::Value; + +use crate::ir::Dialect; + +/// 请求体里可能有提示缓存断点吗:按字节找那个键名,不解析。有的不一定真是断点(正文里 +/// 提到了这个词),没有的一定不是 +pub fn may_have_marks(dialect: Dialect, body: &[u8]) -> bool { + let key: &[u8] = match dialect { + Dialect::Anthropic => b"cache_control", + Dialect::Bedrock => b"cachePoint", + _ => return false, + }; + memchr::memmem::find(body, key).is_some() +} + +/// 去掉请求体里的提示缓存断点:Anthropic 的 `cache_control`,Converse 的 `cachePoint` 块。 +/// 只看工具、系统提示、消息这三处(编码器只标在这里)。 +/// +/// 没有可去的、不是这两种格式、不是 JSON 的返回 `None`,请求一个字节都不改。 +pub fn strip_marks(dialect: Dialect, body: &[u8]) -> Option> { + // 绝大多数请求在这里就回去了:一个字都没有,不必解析 + if !may_have_marks(dialect, body) { + return None; + } + let mut v: Value = serde_json::from_slice(body).ok()?; + let o = v.as_object_mut()?; + // 一处块列表:Anthropic 的块上去掉 `cache_control`,Converse 去掉 `cachePoint` 块 + let strip = |list: Option<&mut Value>| -> bool { + let Some(Value::Array(list)) = list else { + return false; + }; + if dialect == Dialect::Anthropic { + let mut changed = false; + for b in list.iter_mut().filter_map(Value::as_object_mut) { + changed |= b.remove("cache_control").is_some(); + } + changed + } else { + let before = list.len(); + list.retain(|b| b.get("cachePoint").is_none()); + list.len() != before + } + }; + let mut changed = strip(o.get_mut("system")); + changed |= strip(if dialect == Dialect::Anthropic { + o.get_mut("tools") + } else { + o.get_mut("toolConfig").and_then(|c| c.get_mut("tools")) + }); + if let Some(Value::Array(messages)) = o.get_mut("messages") { + for m in messages.iter_mut() { + changed |= strip(m.get_mut("content")); + } + } + changed.then(|| v.to_string().into_bytes()) +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + #[test] + fn the_breakpoints_and_only_them_are_taken_out() { + let mark = json!({"type": "ephemeral"}); + let anthropic = json!({ + "model": "m", + "tools": [{"name": "f", "input_schema": {}, "cache_control": mark}], + "system": [{"type": "text", "text": "s", "cache_control": mark}], + "messages": [ + {"role": "user", "content": [{"type": "text", "text": "mentions cache_control", "cache_control": mark}]}, + {"role": "assistant", "content": "plain"} + ] + }); + let out: Value = serde_json::from_slice( + &strip_marks(Dialect::Anthropic, anthropic.to_string().as_bytes()).unwrap(), + ) + .unwrap(); + assert_eq!( + out, + json!({ + "model": "m", + "tools": [{"name": "f", "input_schema": {}}], + "system": [{"type": "text", "text": "s"}], + "messages": [ + {"role": "user", "content": [{"type": "text", "text": "mentions cache_control"}]}, + {"role": "assistant", "content": "plain"} + ] + }) + ); + + let point = json!({"cachePoint": {"type": "default"}}); + let converse = json!({ + "system": [{"text": "s"}, point], + "toolConfig": {"tools": [{"toolSpec": {"name": "f"}}, point]}, + "messages": [{"role": "user", "content": [{"text": "hi"}, point]}] + }); + let out: Value = serde_json::from_slice( + &strip_marks(Dialect::Bedrock, converse.to_string().as_bytes()).unwrap(), + ) + .unwrap(); + assert_eq!( + out, + json!({ + "system": [{"text": "s"}], + "toolConfig": {"tools": [{"toolSpec": {"name": "f"}}]}, + "messages": [{"role": "user", "content": [{"text": "hi"}]}] + }) + ); + } + + #[test] + fn nothing_to_take_out_leaves_the_body_alone() { + // 只是在正文里提到了这个词 + let body = json!({"messages": [{"role": "user", "content": "what is cache_control?"}]}) + .to_string(); + assert_eq!(strip_marks(Dialect::Anthropic, body.as_bytes()), None); + assert_eq!(strip_marks(Dialect::Anthropic, b"{}"), None); + assert_eq!(strip_marks(Dialect::Chat, b"{\"cache_control\":1}"), None); + assert_eq!(strip_marks(Dialect::Bedrock, b"not json cachePoint"), None); + } +} diff --git a/crates/tw-dialect/src/lib.rs b/crates/tw-dialect/src/lib.rs index 14c8b3a5..563ead38 100644 --- a/crates/tw-dialect/src/lib.rs +++ b/crates/tw-dialect/src/lib.rs @@ -8,11 +8,12 @@ //! 同样两边都用的还有:从响应里旁路嗅出用量([`usage`],换算和转换共用各家的 //! `usage()`),拼上游地址([`url`]),去掉 DeepSeek Harness 只发给 DeepSeek 的 //! 扩展([`harness`]),在原文上找调用方的正文([`caller`],内容过滤读它、删它), -//! 读写各格式里名字不同的请求参数([`params`]),以及 Codex 的远程压缩转给别家时怎么做 -//! ([`compaction`])。 +//! 读写各格式里名字不同的请求参数([`params`]),Codex 的远程压缩转给别家时怎么做 +//! ([`compaction`]),以及请求体上的提示缓存断点([`cache`])。 pub mod anthropic; pub mod bedrock; +pub mod cache; pub mod caller; pub mod chat; pub mod compaction; diff --git a/crates/tw-gateway/src/bodies.rs b/crates/tw-gateway/src/bodies.rs index df5c11aa..16eac89d 100644 --- a/crates/tw-gateway/src/bodies.rs +++ b/crates/tw-gateway/src/bodies.rs @@ -16,10 +16,15 @@ //! - 别的一律打码,和安全日志里报的是同一种写法 //! - 最后整段再按形状打一遍码,和读的时候是同一个函数([`tw_secret::mask_body`]), //! 兜住规则没认出来的 +//! - **转换时压缩出的前文摘要**(`compaction` 项里 `tw1.c.` 开头的那串,见 +//! [`tw_dialect::compaction`])是 base64 包着的一段文字,上面两道都看不见里面:解开, +//! 照同样的规矩换、打码,再包回原位([`carried_summaries`])。Codex 每一轮都把它带回来, +//! 摘要里提到过的密钥否则每一轮都原样落一次盘 //! //! 以前存的是客户端发来的原文,密钥只在读出来的时候才打码:磁盘上躺着的一直是真值。 use std::collections::HashMap; +use std::ops::Range; use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; @@ -117,9 +122,33 @@ impl Redaction { } }; let sent: HashMap<&str, &str> = self.ledger.replacements().collect(); + // 摘要以外的各段照常(找过的命中照用);摘要解开了重新找、换、打码,再包回去 let mut out = String::with_capacity(text.len()); let mut at = 0; - for h in hits { + for (span, summary) in carried_summaries(text) { + out.push_str(&self.span(text, at..span.start, hits, &sent)); + out.push_str(&tw_dialect::compaction::carry(&self.apply(&summary))); + at = span.end; + } + out.push_str(&self.span(text, at..text.len(), hits, &sent)); + out + } + + /// `text` 里的一段:落在这一段里的命中换掉或打码,再按形状打一遍。各段在 token 的边界上 + /// 分开(见 [`carried_summaries`]),一段一段打和整段一起打是一样的 + fn span( + &self, + text: &str, + range: Range, + hits: &[Hit], + sent: &HashMap<&str, &str>, + ) -> String { + let mut out = String::with_capacity(range.len()); + let mut at = range.start; + for h in hits + .iter() + .filter(|h| h.bytes.start >= range.start && h.bytes.end <= range.end) + { out.push_str(&text[at..h.bytes.start]); let value = &text[h.bytes.clone()]; match sent.get(value) { @@ -128,11 +157,43 @@ impl Redaction { } at = h.bytes.end; } - out.push_str(&text[at..]); + out.push_str(&text[at..range.end]); tw_secret::mask_body(&out) } } +/// 正文里转换时压缩出的前文摘要:`tw1.c.` 开头、解得开的那串([`tw_dialect::compaction`]), +/// 在正文里的位置和解开的摘要,按出现的顺序。 +/// +/// 一串按 [`tw_secret::mask_body`] 认 token 的办法划边界(字母、数字、`-`、`_`、`.`), +/// 所以把它挖出去之后,剩下的各段打码和整段打码一样。**解不开的(截断了的、坏了的) +/// 不算**:它照别的字一样打码 —— 一长串不透明的字符,按形状整串打掉 +fn carried_summaries(text: &str) -> Vec<(Range, String)> { + fn is_tok(b: u8) -> bool { + b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.') + } + let bytes = text.as_bytes(); + let prefix = tw_dialect::compaction::CARRIED_PREFIX.as_bytes(); + let mut out: Vec<(Range, String)> = Vec::new(); + for start in memchr::memmem::find_iter(bytes, prefix) { + // 一串的开头,不是别的字中间 + if start > 0 && is_tok(bytes[start - 1]) { + continue; + } + if out.last().is_some_and(|(r, _)| start < r.end) { + continue; + } + let end = bytes[start..] + .iter() + .position(|&b| !is_tok(b)) + .map_or(bytes.len(), |n| start + n); + if let Some(summary) = tw_dialect::compaction::read(&text[start..end]) { + out.push((start..end, summary)); + } + } + out +} + /// 一个认出来、没换成占位符的值存成什么样。 /// /// 和安全日志里报的是同一种写法([`tw_guard::redact::rules::masked`]),只有内网地址和 @@ -533,6 +594,132 @@ mod tests { assert!(!cut.contains("hunter2hunter2"), "{cut}"); } + const AWS: &str = "AKIAIOSFODNN7EXAMPLE"; + + /// 摘要里提到过的密钥:一个规则认得出(Anthropic 的钥匙),一个只认得出形状(AWS 的 + /// 访问密钥 ID) + fn summary() -> String { + format!("Set ANTHROPIC_API_KEY={KEY} and aws key {AWS} in .env; tests pass.") + } + + /// 存下来的那份里每一个转换时写的摘要,解开 + fn summaries_in(stored: &str) -> Vec { + carried_summaries(stored) + .into_iter() + .map(|(_, s)| s) + .collect() + } + + /// 转换时压缩出的摘要(`tw1.c.` + base64)里的密钥:解开、照样打码、包回原位。Codex + /// 每一轮的请求都把它带回来;直通 OpenAI 时上游的回答里没有它,转换时客户端收到的那份 + /// 不落盘 —— 但正文里凡是出现了,一样处理 + #[test] + fn a_secret_inside_a_carried_summary_is_masked_before_it_is_written() { + let item = tw_dialect::compaction::carry(&summary()); + let request = format!( + r#"{{"model":"gpt-5.4","input":[{{"type":"message","role":"user","content":"fix it"}},{{"type":"compaction","encrypted_content":"{item}"}},{{"type":"message","role":"user","content":"next {KEY}"}}]}}"# + ); + let streamed = format!( + "event: response.output_item.done\ndata: {{\"type\":\"response.output_item.done\",\"item\":{{\"type\":\"compaction\",\"encrypted_content\":\"{item}\"}}}}\n\n" + ); + let whole = + format!(r#"{{"output":[{{"type":"compaction","encrypted_content":"{item}"}}]}}"#); + for (kind, body) in [ + (BodyKind::Request, &request), + (BodyKind::Response, &streamed), + (BodyKind::Response, &whole), + ] { + let stored = written(record(kind, body, Redaction::default())); + let inside = summaries_in(&stored); + assert_eq!(inside.len(), 1, "{stored}"); + let s = &inside[0]; + for secret in [KEY, AWS] { + assert!(!s.contains(secret), "{secret} 原样进了磁盘:{s}"); + assert!(!stored.contains(secret), "{stored}"); + } + assert!(s.contains("ANTHROPIC_API_KEY=sk-an…AAAA"), "{s}"); + assert!(s.contains("tests pass."), "摘要的其余部分要留着:{s}"); + // 摘要外面的照常打码 + if kind == BodyKind::Request { + assert!(stored.contains("next sk-an…AAAA"), "{stored}"); + serde_json::from_str::(&stored).expect("存下来的还是 JSON"); + } + } + } + + /// 请求路上找过的命中照用(找过的是摘要外面的那些),摘要里面另外找;拦截档下摘要里 + /// 的值换成发给上游的那个占位符 + #[test] + fn the_hits_found_on_the_way_in_still_cover_the_rest_and_the_ledger_reaches_inside() { + let item = tw_dialect::compaction::carry(&summary()); + let body = format!( + r#"{{"input":[{{"type":"compaction","encrypted_content":"{item}"}},{{"role":"user","content":"key {KEY} 库 postgres://app:hunter2hunter2@db/x"}}]}}"# + ); + for mode in [ + tw_config::SecurityMode::Observe, + tw_config::SecurityMode::Enforce, + ] { + let rules = RuleSet::defaults(); + let seen = crate::guard::look_hits(mode, &rules, body.as_bytes()); + let r = Redaction { + rules: Arc::new(rules), + ledger: seen.ledger, + }; + let again = written(record(BodyKind::Request, &body, r.clone())); + let reused = + written(record(BodyKind::Request, &body, r).found(seen.hits.map(Arc::from))); + assert_eq!(reused, again, "{mode:?}"); + assert!(!reused.contains("hunter2hunter2"), "{reused}"); + let inside = summaries_in(&reused); + assert!(!inside[0].contains(KEY), "{mode:?}: {inside:?}"); + if mode == tw_config::SecurityMode::Enforce { + // 上游收到的是占位符,摘要里存的也是它 + assert!(inside[0].contains("<>"), "{inside:?}"); + } + } + } + + #[test] + fn a_damaged_carried_summary_is_masked_as_text_and_never_panics() { + for bad in [ + // 一个字符凑不出一个字节 + "tw1.c.a".to_string(), + // 不是 base64url + "tw1.c.!!!!".to_string(), + // 截断在半个 UTF-8 字符上 + format!( + "tw1.c.{}", + &tw_dialect::compaction::carry("密钥")["tw1.c.".len()..][..3] + ), + // 一长串解不开的(base64 的标准字母表) + format!("tw1.c.{}+/=", "QUtJQUlPU0ZPRE5ON0VYQU1QTEU".repeat(3)), + // 前缀之后什么都没有 + "tw1.c.".to_string(), + ] { + let body = + format!(r#"{{"input":[{{"type":"compaction","encrypted_content":"{bad}"}}]}}"#); + let stored = written(record(BodyKind::Request, &body, Redaction::default())); + // 和没有这回事时一样:按形状打码(只剩前缀的那个解出来是空摘要,原样) + assert_eq!(stored, tw_secret::mask_body(&body), "{bad}"); + assert!(summaries_in(&stored).iter().all(String::is_empty), "{bad}"); + } + // 正文被截断在一个摘要中间(存下来的只有开头 4 MB):解得开的那半截照样在里面打码, + // 解不开的整串打掉。哪一种都不会留下原值 + let item = tw_dialect::compaction::carry(&summary()); + for drop in 1..=8 { + let cut = &item[..item.len() - drop]; + let body = format!(r#"{{"encrypted_content":"{cut}"#); + let stored = written(record(BodyKind::Request, &body, Redaction::default())); + let inside = summaries_in(&stored); + for s in &inside { + assert!(!s.contains(KEY) && !s.contains(AWS), "{drop}: {s}"); + } + if inside.is_empty() { + assert!(!stored.contains(&cut[10..40]), "{drop}: {stored}"); + } + } + } + #[test] fn what_the_rules_miss_is_masked_by_its_shape() { // 网关自己的钥匙(`tw-`)没有一条脱敏规则认:读的时候那一道兜住它,写的时候也一样 diff --git a/crates/tw-gateway/src/cache_marks.rs b/crates/tw-gateway/src/cache_marks.rs new file mode 100644 index 00000000..daccc7b4 --- /dev/null +++ b/crates/tw-gateway/src/cache_marks.rs @@ -0,0 +1,113 @@ +//! 转给 Claude 时自动标的提示缓存断点,被上游拒了怎么办。 +//! +//! 客户端没标断点的请求转成 Anthropic 格式、或者转给 Bedrock 上的 Claude 时,编码器替它 +//! 标(`tw_dialect::anthropic::request` 的 `auto_cache`)。有的上游不认:自称 Anthropic +//! 格式的兼容接口不收 `cache_control`,AWS 改了支持哪些模型。不认的回 400,整个请求被拒。 +//! +//! 被拒了就**去掉这些断点、同一家再发一次**(只一次),再记下这一家的这个模型不认,往后 +//! 发给它之前先去掉,不用每次都先被拒一次。和别家封存的推理被拒时一样(见 [`crate::seal`]): +//! 不换一家,这一家好好的,只是不认那几个记号。 +//! +//! **只动编码器加的断点**:客户端自己标了的,请求体里的断点是它的决定,原样发,被拒了也 +//! 原样交回去。 + +use std::collections::HashMap; +use std::sync::Mutex; + +/// 最多记多少个(上游,模型)。表满了丢最早记下的那个:丢了也只是下次再被拒一次、再重发 +const MAX: usize = 1024; + +/// 被拒的响应体最多看这么多。拒绝的原话很短 +const REFUSAL_MAX: usize = 64 * 1024; + +/// 哪些(上游,模型)不认自动标的断点。**只在内存里,跨重载存活**:重启之后丢了也只是 +/// 下一次再被拒一次、再重发一次。 +#[derive(Default)] +pub struct Refused { + by: Mutex>, +} + +impl Refused { + /// `upstream` 上的 `model` 拒了自动标的断点 + pub fn note(&self, upstream: &str, model: &str, at_ms: u64) { + let mut by = self.by.lock().unwrap_or_else(|p| p.into_inner()); + if by.len() >= MAX + && !by.contains_key(&(upstream.to_string(), model.to_string())) + && let Some(oldest) = by.iter().min_by_key(|(_, at)| **at).map(|(k, _)| k.clone()) + { + by.remove(&oldest); + } + by.insert((upstream.to_string(), model.to_string()), at_ms); + } + + /// `upstream` 上的 `model` 拒过自动标的断点吗 + pub fn refused(&self, upstream: &str, model: &str) -> bool { + let by = self.by.lock().unwrap_or_else(|p| p.into_inner()); + by.contains_key(&(upstream.to_string(), model.to_string())) + } + + #[cfg(test)] + fn len(&self) -> usize { + self.by.lock().unwrap().len() + } +} + +/// 上游这个 400 是在拒绝提示缓存的断点吗:原话里说到了 `cache_control`、`cachePoint` +/// 或者提示缓存。 +/// +/// 看原话,不看错误外壳的结构:兼容接口说「Extra inputs are not permitted: +/// messages.0.content.0.cache_control」「unknown field `cache_control`」,Bedrock 说 +/// `ValidationException` 加一句提到 `cachePoint` 或 prompt caching 的话,写在哪个字段里 +/// 各家不一样 +pub fn refusal(said: &[u8]) -> bool { + let said = String::from_utf8_lossy(&said[..said.len().min(REFUSAL_MAX)]).to_ascii_lowercase(); + [ + "cache_control", + "cachepoint", + "cache_point", + "cache point", + "prompt caching", + "prompt_caching", + ] + .iter() + .any(|w| said.contains(w)) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn what_each_upstream_says_when_it_does_not_take_breakpoints() { + for said in [ + r#"{"type":"error","error":{"type":"invalid_request_error","message":"messages.0.content.0.cache_control: Extra inputs are not permitted"}}"#, + r#"{"error":{"message":"unknown field `cache_control`, expected one of `type`, `text`","code":"invalid_request"}}"#, + r#"{"message":"The model returned the following errors: cachePoint is not supported for this model."}"#, + r#"{"message":"Prompt caching is not supported for this model"}"#, + ] { + assert!(refusal(said.as_bytes()), "{said}"); + } + for said in [ + r#"{"type":"error","error":{"type":"invalid_request_error","message":"max_tokens: Field required"}}"#, + r#"{"message":"Malformed input request: #/messages/0/content: expected type: JSONArray"}"#, + ] { + assert!(!refusal(said.as_bytes()), "{said}"); + } + } + + #[test] + fn the_memory_is_per_upstream_and_model_and_stays_bounded() { + let r = Refused::default(); + r.note("relay", "claude-opus-4-7", 1); + assert!(r.refused("relay", "claude-opus-4-7")); + assert!(!r.refused("relay", "claude-sonnet-4-5")); + assert!(!r.refused("official", "claude-opus-4-7")); + for i in 0..MAX + 10 { + r.note("relay", &format!("m{i}"), 10 + i as u64); + } + assert_eq!(r.len(), MAX); + // 最早记下的先丢 + assert!(!r.refused("relay", "claude-opus-4-7")); + assert!(r.refused("relay", &format!("m{}", MAX + 9))); + } +} diff --git a/crates/tw-gateway/src/l3.rs b/crates/tw-gateway/src/l3.rs index c87458b4..da8570c6 100644 --- a/crates/tw-gateway/src/l3.rs +++ b/crates/tw-gateway/src/l3.rs @@ -162,10 +162,13 @@ pub fn probe_request( upstream_headers, )); } + // 编码器给转成 Claude 的请求自动标缓存断点;测速的请求短得缓存不了,标了反而可能被不认 + // 的上游整个拒掉、被当成测速失败(见 `crate::cache_marks`) + let body = tw_dialect::cache::strip_marks(dialect, &prepared.body).unwrap_or(prepared.body); ProbeRequest { path, query, - body: prepared.body, + body, headers, } } diff --git a/crates/tw-gateway/src/lib.rs b/crates/tw-gateway/src/lib.rs index 0539e735..96c899b7 100644 --- a/crates/tw-gateway/src/lib.rs +++ b/crates/tw-gateway/src/lib.rs @@ -13,6 +13,7 @@ pub mod auth; pub mod balance; pub mod bedrock; pub mod bodies; +pub mod cache_marks; pub mod chatgpt; pub mod client_api; pub mod clientprobe; diff --git a/crates/tw-gateway/src/server/pipeline/hop.rs b/crates/tw-gateway/src/server/pipeline/hop.rs index 5407b05e..2db480f6 100644 --- a/crates/tw-gateway/src/server/pipeline/hop.rs +++ b/crates/tw-gateway/src/server/pipeline/hop.rs @@ -110,6 +110,9 @@ struct Outbound { /// 回程要用的转换会话:转换过的,或者直通到 Codex 后端、客户端却要整包时 /// 收齐流要用的 session: Option, + /// 请求体里的提示缓存断点是转换时自动标的(客户端一个都没标):发给的是这个模型。 + /// 上游拒了这些断点就去掉再发一次(见 [`crate::cache_marks`]) + auto_cache: Option, } /// 这一跳没发出去的原因。尝试链里记的和报给客户端的是同一句。 @@ -556,7 +559,7 @@ pub(super) async fn try_upstreams<'a>( aws.as_ref(), ) .await; - // 上游拒绝了别家封存的推理:去掉它们,同一家再发一次 + // 上游拒绝了自动标的缓存断点、别家封存的推理:去掉它们,同一家再发一次 match sent { Ok(r) => { let resend = Resend { @@ -566,6 +569,15 @@ pub(super) async fn try_upstreams<'a>( aws: aws.as_ref(), conversation: started.conversation.as_deref(), }; + let (r, uncached) = + resend_uncached(state, req, provider, http, resend, r).await?; + let resend = Resend { + out: &out, + body: uncached.as_ref().unwrap_or(&body), + headers: &upstream_headers, + aws: aws.as_ref(), + conversation: started.conversation.as_deref(), + }; resend_unsealed(state, req, provider, http, resend, r).await } Err(e) => Err(e), @@ -1168,6 +1180,7 @@ fn prepare( let mut path = asked.path.to_string(); let mut query = req.query.clone(); let mut session: Option = None; + let mut auto_cache: Option = None; let body = match target { None => { // 参数改写。**只在这里动 body,而且只动被点名的那几个字段** —— @@ -1301,6 +1314,22 @@ fn prepare( }); path = p.path.clone(); query = p.query.clone(); + // 客户端一个断点都没标:请求体里的断点是编码器自动标的。这一家的这个模型拒过的话 + // 先去掉,不用再被拒一次(见 `crate::cache_marks`) + let mut encoded = p.body.clone(); + if d.request.cache.is_empty() && tw_dialect::cache::may_have_marks(dialect, &encoded) { + let model = &d.request.model; + if !state.cache_marks.refused(&provider.name, model) { + auto_cache = Some(model.clone()); + } else if let Some(stripped) = tw_dialect::cache::strip_marks(dialect, &encoded) { + tracing::debug!( + provider = %provider.name, + model = %model, + "left out the cache breakpoints this upstream refused before" + ); + encoded = stripped; + } + } // Bedrock 上的 Claude:客户端 `anthropic-beta` 里 Bedrock 认的那几个放进请求体 let claude_on_bedrock = dialect == tw_dialect::ir::Dialect::Bedrock && d.client == tw_dialect::ir::Dialect::Anthropic @@ -1320,16 +1349,14 @@ fn prepare( }; // **客户端要不要流由会话记着**,发给 Codex 后端的这一份一律是流式 let body = if chatgpt { - Bytes::from(crate::chatgpt::force_stream(p.body.clone())) - } else if let Some(b) = tw_bedrock::beta::with_betas(&p.body, &betas) { + Bytes::from(crate::chatgpt::force_stream(encoded)) + } else if let Some(b) = tw_bedrock::beta::with_betas(&encoded, &betas) { Bytes::from(b) } else if harness && to_deepseek { // 转换成另一种格式发给 DeepSeek 官方:直连时它收得到的扩展照样带上 - Bytes::from( - tw_dialect::harness::carry(asked.body, &p.body).unwrap_or(p.body.clone()), - ) + Bytes::from(tw_dialect::harness::carry(asked.body, &encoded).unwrap_or(encoded)) } else { - Bytes::from(p.body.clone()) + Bytes::from(encoded) }; session = Some(p.session); body @@ -1356,6 +1383,7 @@ fn prepare( chatgpt, hop, session, + auto_cache, }) } @@ -1401,6 +1429,61 @@ struct Resend<'a> { conversation: Option<&'a str>, } +/// 上游回 400、拒绝了转换时自动标的提示缓存断点(见 [`crate::cache_marks`]):去掉,同一家 +/// 再发一次,只一次。记下这一家的这个模型不认,往后发给它之前先去掉。 +/// +/// 客户端自己标了断点的不管([`Outbound::auto_cache`] 是 None):那是它的决定,被拒了原样 +/// 交回去。不是这种 400 的也原样交回去。返回再发时用的请求体,没再发是 None +async fn resend_uncached( + state: &AppState, + req: &Inbound, + provider: &tw_config::Provider, + http: &reqwest::Client, + resend: Resend<'_>, + r: reqwest::Response, +) -> Result<(reqwest::Response, Option), SendError> { + let (Some(model), Some(wire)) = (resend.out.auto_cache.as_deref(), wire(req, resend.out)) + else { + return Ok((r, None)); + }; + if r.status() != 400 { + return Ok((r, None)); + } + let Some(stripped) = tw_dialect::cache::strip_marks(wire, resend.body) else { + return Ok((r, None)); + }; + let status = r.status(); + let headers = r.headers().clone(); + let said = r.bytes().await.map_err(SendError::Http)?; + if !crate::cache_marks::refusal(&said) { + let mut back = http::Response::new(said); + *back.status_mut() = status; + *back.headers_mut() = headers; + return Ok((reqwest::Response::from(back), None)); + } + state + .cache_marks + .note(&provider.name, model, crate::server::now_ms()); + tracing::info!( + provider = %provider.name, + model = %model, + "the upstream refused the cache breakpoints added on conversion; sending again without them" + ); + let stripped = Bytes::from(stripped); + let r = send( + state, + req, + provider, + http, + resend.out, + stripped.clone(), + resend.headers.to_vec(), + resend.aws, + ) + .await?; + Ok((r, Some(stripped))) +} + /// 上游回 400、拒绝了请求里别家封存的推理(见 [`crate::seal`]):去掉**全部**封存的 /// 推理,同一家再发一次,只一次。拒过哪些记下来,这段对话往后发给它之前先去掉。 /// diff --git a/crates/tw-gateway/src/state.rs b/crates/tw-gateway/src/state.rs index d2921cc7..4319fb13 100644 --- a/crates/tw-gateway/src/state.rs +++ b/crates/tw-gateway/src/state.rs @@ -234,6 +234,9 @@ pub struct AppState { pub balance: Arc, /// 每段对话里、每一家上游拒过的别家封存的推理(见 [`crate::seal`])。**跨重载存活** pub seals: Arc, + /// 哪些(上游,模型)拒过转换时自动标的提示缓存断点(见 [`crate::cache_marks`])。 + /// **跨重载存活** + pub cache_marks: Arc, /// 脚本插件里跨重载存活的那一半:运行时、插件文件在哪儿、计数和日志、编译缓存 /// (见 [`crate::plugin::Plugins`])。跟着配置换的那一半在 `Runtime::plugins` pub plugins: Arc, @@ -319,6 +322,7 @@ impl AppState { affinity: Default::default(), balance: Default::default(), seals: Default::default(), + cache_marks: Default::default(), plugins, swap: Default::default(), ping_every: crate::PING_EVERY, diff --git a/crates/tw-gateway/tests/bedrock.rs b/crates/tw-gateway/tests/bedrock.rs index 7fdde388..a1377cca 100644 --- a/crates/tw-gateway/tests/bedrock.rs +++ b/crates/tw-gateway/tests/bedrock.rs @@ -692,3 +692,84 @@ async fn counting_tokens_is_not_supported_the_way_claude_code_expects() { "counting tokens reached AWS" ); } + +/// 不认 `cachePoint` 的 Bedrock:请求里带着它就回 400 ValidationException +fn refuses_cache_points() -> Answer { + Arc::new(|s: &Seen| { + if String::from_utf8_lossy(&s.body).contains("cachePoint") { + axum::response::Response::builder() + .status(400) + .header("content-type", "application/json") + .header( + "x-amzn-errortype", + "ValidationException:http://internal.amazon.com/coral/com.amazon.bedrock/", + ) + .body(axum::body::Body::from( + json!({"message": "The model returned the following errors: cachePoint is not supported for this model."}).to_string(), + )) + .unwrap() + } else { + json_answer(200, converse_reply("ok")) + } + }) +} + +#[tokio::test] +async fn automatic_cache_points_a_model_refuses_are_dropped_and_not_sent_again() { + let (up, seen) = bedrock(refuses_cache_points()).await; + let (gw, _, mut rx) = gateway(with_keys(up)).await; + let ask = || { + post( + gw, + "/v1/chat/completions", + &[("authorization", "Bearer tw-k")], + json!({"model": MODEL, "messages": [ + {"role": "system", "content": "Be brief."}, + {"role": "user", "content": "hi"} + ]}), + ) + }; + // Chat 客户端不标断点:转给 Claude 时自动标,被拒了就去掉、同一家再发一次 + let (status, body) = ask().await; + assert_eq!(status, 200, "{body}"); + { + let seen = seen.lock().unwrap(); + assert_eq!(seen.len(), 2, "{seen:?}"); + assert!(String::from_utf8_lossy(&seen[0].body).contains("cachePoint")); + assert!(!String::from_utf8_lossy(&seen[1].body).contains("cachePoint")); + // 再发的那一份按它自己的内容重新签名 + assert_signed_as_sent(&seen[1]); + } + assert!(matches!( + ending(&mut rx).await, + tw_api::Event::RequestFinished { .. } + )); + // 记住了:这一家的这个模型往后直接不标 + let (status, body) = ask().await; + assert_eq!(status, 200, "{body}"); + let seen = seen.lock().unwrap(); + assert_eq!(seen.len(), 3, "{seen:?}"); + assert!(!String::from_utf8_lossy(&seen[2].body).contains("cachePoint")); +} + +#[tokio::test] +async fn cache_points_the_client_set_itself_are_never_dropped() { + // Claude Code 自己标的断点:被拒了原样交回去,不替它去掉 + let (up, seen) = bedrock(refuses_cache_points()).await; + let (gw, _, _) = gateway(with_keys(up)).await; + let (status, body) = post( + gw, + "/v1/messages", + &[("x-api-key", "tw-k")], + json!({ + "model": MODEL, + "max_tokens": 16, + "system": [{"type": "text", "text": "rules", "cache_control": {"type": "ephemeral"}}], + "messages": [{"role": "user", "content": "hi"}] + }), + ) + .await; + assert_eq!(status, 400, "{body}"); + assert!(body.contains("cachePoint"), "{body}"); + assert_eq!(seen.lock().unwrap().len(), 1); +} diff --git a/crates/tw-gateway/tests/conversion.rs b/crates/tw-gateway/tests/conversion.rs index 6a01c92f..5cfe35aa 100644 --- a/crates/tw-gateway/tests/conversion.rs +++ b/crates/tw-gateway/tests/conversion.rs @@ -878,3 +878,123 @@ async fn a_claude_request_without_breakpoints_passes_straight_through_unchanged( assert_eq!(status, 200, "{body}"); assert_eq!(seen.lock().unwrap().body, sent.to_string().into_bytes()); } + +/// 自称 Anthropic 格式、却不收 `cache_control` 的兼容接口:带着它就回 400 +async fn refuses_cache_control() -> (SocketAddr, Arc>>>) { + let seen: Arc>>> = Arc::default(); + let app = Router::new() + .fallback( + move |State(s): State>>>>, body: bytes::Bytes| async move { + let refused = String::from_utf8_lossy(&body).contains("cache_control"); + s.lock().unwrap().push(body.to_vec()); + let (status, ct, reply) = if refused { + ( + 400, + "application/json", + json!({"type": "error", "error": {"type": "invalid_request_error", + "message": "messages.0.content.0.cache_control: Extra inputs are not permitted"}}) + .to_string(), + ) + } else { + (200, "text/event-stream", ANTHROPIC_STREAM.to_string()) + }; + axum::response::Response::builder() + .status(status) + .header("content-type", ct) + .body(axum::body::Body::from(reply)) + .unwrap() + }, + ) + .with_state(seen.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, seen) +} + +#[tokio::test] +async fn automatic_breakpoints_an_upstream_refuses_are_dropped_and_not_sent_again() { + let (up, seen) = refuses_cache_control().await; + let (gw, _) = gateway(provider(up, Protocol::Anthropic), SecurityMode::Observe).await; + let ask = || { + post( + gw, + "/v1/responses", + &[("authorization", "Bearer tw-k")], + codex_lite_request(), + ) + }; + // 自动标了断点、被拒:去掉,同一家再发一次,客户端照常拿到回答 + let (status, _, body) = ask().await; + assert_eq!(status, 200, "{body}"); + assert!(body.contains("你好"), "{body}"); + { + let seen = seen.lock().unwrap(); + assert_eq!(seen.len(), 2); + assert!(String::from_utf8_lossy(&seen[0]).contains("cache_control")); + let again: Value = serde_json::from_slice(&seen[1]).unwrap(); + assert!(!again.to_string().contains("cache_control"), "{again}"); + // 别的一个字都没动 + assert_eq!(again["tools"].as_array().unwrap().len(), 3); + assert_eq!(again["messages"][0]["content"][0]["text"], "list the files"); + } + // 记住了:这一家的这个模型往后直接不标,不再先被拒一次 + let (status, _, body) = ask().await; + assert_eq!(status, 200, "{body}"); + let seen = seen.lock().unwrap(); + assert_eq!(seen.len(), 3); + assert!(!String::from_utf8_lossy(&seen[2]).contains("cache_control")); +} + +#[tokio::test] +async fn a_secret_in_a_carried_summary_does_not_reach_the_disk() { + // 之前转换时压缩出的摘要里提到了一把钥匙,Codex 这一轮把它带回来了 + const KEY: &str = "sk-ant-api03-AAAAAAAAAAAAAAAAAAAAAAAAAA"; + let (up, _) = upstream(200, "text/event-stream", ANTHROPIC_STREAM.into()).await; + let cfg = Config { + version: 1, + listen: Listen::default(), + clients: vec![Client { + name: "c".into(), + key: "tw-k".into(), + ..Default::default() + }], + providers: vec![provider(up, Protocol::Anthropic)], + ..Default::default() + }; + let state = tw_gateway::AppState::new(cfg).unwrap(); + let (tx, mut bodies) = tw_gateway::bodies::channel(); + state.set_body_sink(tx); + let gw = tw_gateway::serve(state, ([127, 0, 0, 1], 0).into()) + .await + .unwrap(); + let mut req = codex_lite_request(); + req["input"].as_array_mut().unwrap().insert( + 3, + json!({"type": "compaction", "encrypted_content": + tw_dialect::compaction::carry(&format!("Exported ANTHROPIC_API_KEY={KEY}; tests pass."))}), + ); + let (status, _, body) = post( + gw, + "/v1/responses", + &[("authorization", "Bearer tw-k")], + req, + ) + .await; + assert_eq!(status, 200, "{body}"); + + let rec = tokio::time::timeout(Duration::from_secs(3), bodies.recv()) + .await + .unwrap() + .unwrap(); + assert_eq!(rec.kind, tw_gateway::bodies::BodyKind::Request); + let stored = String::from_utf8(rec.for_disk().body.to_vec()).unwrap(); + let v: Value = serde_json::from_str(&stored).unwrap(); + let carried = v["input"][3]["encrypted_content"].as_str().unwrap(); + let summary = tw_dialect::compaction::read(carried).expect("摘要还解得开"); + assert!(!summary.contains(KEY), "{summary}"); + assert!( + summary.contains("ANTHROPIC_API_KEY=sk-an…AAAA; tests pass."), + "{summary}" + ); +} diff --git a/docs/config.md b/docs/config.md index 20119977..5775a5ae 100644 --- a/docs/config.md +++ b/docs/config.md @@ -467,6 +467,13 @@ writes and reads are charged at the price table's cache prices. A request that carries its own marks, as Claude Code's do, keeps exactly those, and a request sent on in the upstream's own format is not changed. +On Bedrock, marks are added only for the Claude models AWS lists as +supporting prompt caching: Claude 3.7 Sonnet, Claude 3.5 Sonnet v2, and every +Claude from version 4.5 on, including newer ones not yet listed. Older models, +such as Claude 3 Haiku, Sonnet 4 and Opus 4.1, get none. An upstream that +refuses the marks is sent the request once more without them, and that +upstream is not sent marks for that model again until the core restarts. + A ChatGPT account upstream (`protocol: chatgpt`) takes only the credential the desktop app obtains by signing in; it cannot be written by hand. Claude and Google subscription sign-ins are not supported; use an API key. diff --git a/docs/config.zh-CN.md b/docs/config.zh-CN.md index 57cd8312..bb985241 100644 --- a/docs/config.zh-CN.md +++ b/docs/config.zh-CN.md @@ -335,6 +335,8 @@ providers: 请求转换为 Anthropic 格式、或转给 Bedrock 上的 Claude 时,客户端自己没有标出缓存位置的,由网关标出可以缓存的位置。Codex 等使用 OpenAI 或 Gemini 格式的客户端无从标注:这两家自动缓存重复的提示,Anthropic 只缓存标出的部分。标注的位置是工具列表末尾、系统提示末尾和最后两条用户消息的末尾,最多四处,缓存时长为默认的五分钟。两条用户消息中靠前的那一处正是上一个请求结束的位置,因此每一轮都能读回上一轮写入的缓存,这部分按缓存价而不是全额输入价计费。缓存的写入和读取按价目表中的缓存单价计费。自带标注的请求(如 Claude Code 的)保持原样,按上游自身格式直接转发的请求不做改动。 +Bedrock 上只为 AWS 列出支持提示缓存的 Claude 模型标注:Claude 3.7 Sonnet、Claude 3.5 Sonnet v2,以及 4.5 起的所有 Claude,包括尚未列出的更新版本。更早的模型(如 Claude 3 Haiku、Sonnet 4、Opus 4.1)不标注。上游拒绝这些标注时,网关去掉标注重新发送一次,此后在 core 重启之前不再为这家上游的该模型标注。 + ChatGPT 账号上游(`protocol: chatgpt`)只接受桌面应用登录得到的凭据,不能手写。不支持 Claude 和 Google 的订阅登录,请使用 API 密钥。 有的中转站和账号同时只接受几个请求,多出来的直接拒绝。`max_concurrent` 让网关守住这个数:请求发出时占用这家的一个位置,回答完整交给客户端、或者客户端断开时归还。这家满了的时候,为复用提示缓存而留在这家的对话等空位,别的请求直接换下一家。最多等多久由 `failover.slot_wait_secs` 决定。等待不算失败,这家不会因此停用。只计算 token 数的请求不占位置。Responses 的 WebSocket 连接上,每个 `response.create` 从发出起占一个位置,直到它的回答结束,密钥的 `max_concurrent` 也一样;空闲的连接不占位置。