Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions crates/tw-gateway/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,14 @@ pub use server::{router, serve};
pub use state::credential_failed;
pub use state::{AppState, Runtime};

/// Anthropic 流里上游静默多久补一个 `ping`。
///
/// Claude Code(和内嵌它的 Claude Desktop)数的是网关发来的每一个字节:一条流静默
/// 五分钟就被放弃,只有 ping 在来的话还肯多等一段。Anthropic 上游思考时自己会发
/// ping,**转换别的上游时没有人发**,长时间的推理就会被客户端当成断线。十五秒远低于
/// 那个上限,又不至于让一条正常的流塞满心跳
pub const PING_EVERY: std::time::Duration = std::time::Duration::from_secs(15);

/// 请求来自谁。**如实写 ThinkWatch** —— 我们从不把自己报成别的客户端。
pub const ORIGINATOR: &str = "thinkwatch";

Expand Down
4 changes: 4 additions & 0 deletions crates/tw-gateway/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,10 @@ mod upgrade;
pub fn router(state: AppState) -> Router {
Router::new()
.route("/healthz", get(|| async { "ok" }))
// Claude Code(和内嵌它的 Claude Desktop)启动时发一个 `HEAD /api/hello` 预热
// 连接。**不鉴权、不转发**:它不带密钥也不需要上游,掉进透传的话会被当成一次
// 没有密钥的请求拒成 401。`get` 也接 HEAD
.route("/api/hello", get(|| async {}))
// **和准入共用同一个函数** —— 列表和准入不可能不一致。
.route("/v1/models", get(list_models))
// **单点查询要走同一道准入**。不接这条的话它掉进
Expand Down
143 changes: 140 additions & 3 deletions crates/tw-gateway/src/server/listing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -90,8 +90,14 @@ pub(super) async fn list_models(
"name": format!("models/{m}"),
})).collect::<Vec<_>>()
}),
// Anthropic 和 OpenAI 的 /v1/models 形状一样
_ => serde_json::json!({
ListingShape::Anthropic => serde_json::json!({
"object": "list",
"data": models.iter().map(|m| anthropic_model(m, now)).collect::<Vec<_>>(),
"has_more": false,
"first_id": models.first(),
"last_id": models.last(),
}),
ListingShape::Openai => serde_json::json!({
"object": "list",
"data": models.iter().map(|m| serde_json::json!({
"id": m, "object": "model", "created": now,
Expand All @@ -101,6 +107,88 @@ pub(super) async fn list_models(
Ok(axum::Json(body).into_response())
}

/// Anthropic 格式的一个模型对象。
///
/// **是 OpenAI 那个对象的超集**:Anthropic 的字段(`type`、`display_name`、
/// `created_at`)之外,`object` 和 `created` 照样在。只放 `x-api-key` 的客户端也被
/// 认成 Anthropic,其中有按 OpenAI 的形状读列表的,不能让它们读不出来。
///
/// Claude 的模型再带上 `anthropic_family_tier`:Claude Desktop 按它把模型归到
/// opus / sonnet / haiku,配置里写的 `sonnet` 这样的简称靠它解析。**只看模型名**,
/// 名字里看不出是 Claude 的一律不标 —— 把别家的模型标成 Claude 是在替客户端撒谎。
fn anthropic_model(id: &str, now: u64) -> serde_json::Value {
let created_at = chrono::DateTime::from_timestamp(now as i64, 0)
.unwrap_or_default()
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let mut m = serde_json::json!({
"type": "model",
"id": id,
"display_name": display_name(id).unwrap_or_else(|| id.to_string()),
"created_at": created_at,
"object": "model",
"created": now,
});
if let Some(tier) = family_tier(id) {
m["anthropic_family_tier"] = tier.into();
}
m
}

/// Claude 模型的名字:`claude-sonnet-4-5-20250929` → `Claude Sonnet 4.5`。
///
/// 只认最后一段(`/` 之后)以 `claude-` 开头、每一节都是字母或数字的;日期那一节
/// 去掉,相邻的数字用点连起来。**认不出就是 `None`**,调用方用模型 ID 本身 ——
/// 客户端看到和 ID 一样的名字时会自己想办法,一个猜错的名字它却会照着显示。
fn display_name(id: &str) -> Option<String> {
let last = id.rsplit('/').next()?;
let rest = last.strip_prefix("claude-")?;
let mut words: Vec<String> = vec!["Claude".into()];
let mut number = false;
for part in rest.split('-') {
if part.is_empty() {
return None;
}
if part.len() == 8 && part.bytes().all(|b| b.is_ascii_digit()) {
// 发布日期,不是名字的一部分
continue;
}
if part.bytes().all(|b| b.is_ascii_digit() || b == b'.') {
match words.last_mut() {
Some(w) if number => {
w.push('.');
w.push_str(part);
}
_ => words.push(part.to_string()),
}
number = true;
} else if part.bytes().all(|b| b.is_ascii_alphanumeric()) {
let mut c = part.chars();
let first = c.next()?.to_ascii_uppercase();
words.push(std::iter::once(first).chain(c).collect());
number = false;
} else {
return None;
}
}
(words.len() > 1).then(|| words.join(" "))
}

/// 名字里看得出是哪一档的 Claude 模型:`opus`、`sonnet` 或 `haiku`。
fn family_tier(id: &str) -> Option<&'static str> {
let lower = id.to_ascii_lowercase();
if !lower.contains("claude") && !lower.contains("anthropic") {
return None;
}
let mut tiers = ["opus", "sonnet", "haiku"]
.into_iter()
.filter(|t| lower.contains(t));
// 名字里同时出现两档的(一个路由别名)不猜
match (tiers.next(), tiers.next()) {
(Some(t), None) => Some(t),
_ => None,
}
}

/// `GET /v1/models/:model`。
///
/// **不许可就当它不存在(404),不是 403。**回 403 等于告诉对方
Expand Down Expand Up @@ -146,7 +234,56 @@ pub(super) async fn get_model(
ListingShape::Gemini => {
serde_json::json!({ "name": format!("models/{model}") })
}
_ => serde_json::json!({ "id": model, "object": "model", "created": now }),
ListingShape::Anthropic => anthropic_model(&model, now),
ListingShape::Openai => {
serde_json::json!({ "id": model, "object": "model", "created": now })
}
};
Ok(axum::Json(body).into_response())
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn claude_model_ids_get_their_names() {
for (id, name) in [
("claude-sonnet-4-5", "Claude Sonnet 4.5"),
("claude-sonnet-4-5-20250929", "Claude Sonnet 4.5"),
("claude-opus-4-1-20250805", "Claude Opus 4.1"),
("claude-3-5-haiku-20241022", "Claude 3.5 Haiku"),
("claude-fable-5-1", "Claude Fable 5.1"),
("anthropic/claude-sonnet-4.5", "Claude Sonnet 4.5"),
] {
assert_eq!(display_name(id).as_deref(), Some(name), "{id}");
}
}

#[test]
fn a_name_that_cannot_be_read_is_not_guessed() {
for id in [
"deepseek-chat",
"gpt-5",
"claude-",
"claude--x",
"us.anthropic.claude-sonnet-4-5-20250929-v1:0",
] {
assert_eq!(display_name(id), None, "{id}");
}
}

#[test]
fn only_claude_models_are_given_a_family_tier() {
assert_eq!(family_tier("claude-opus-4-1"), Some("opus"));
assert_eq!(family_tier("claude-3-5-haiku-20241022"), Some("haiku"));
assert_eq!(
family_tier("us.anthropic.claude-sonnet-4-5-20250929-v1:0"),
Some("sonnet")
);
assert_eq!(family_tier("claude-fable-5-1"), None);
// 名字里有 sonnet,但看不出是 Claude
assert_eq!(family_tier("my-sonnet-alias"), None);
assert_eq!(family_tier("claude-opus-or-sonnet"), None);
}
}
44 changes: 41 additions & 3 deletions crates/tw-gateway/src/server/pipeline/relay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,12 @@ pub(super) fn respond(
let mut relay = Relay::new(state, rt, req, plan, &ledger, session, provider, id);
let chunks = upstream.bytes_stream();
let dialect = req.dialect;
// 上游静默时补心跳:**只给 Anthropic Messages 的流**。那是客户端按字节计时、
// 认得 `ping` 事件的格式;别的格式没有这个事件,塞进去就是一帧解析不了的东西
let ping_every = (req.api == Some(crate::client_api::ClientApi::AnthropicMessages)
&& plan.client_sse
&& status.is_success())
.then_some(state.ping_every);
let stream = async_stream::stream! {
// **通行证跟着响应体走。**这个流被丢掉的时候它才还回去:正常
// 发完是一种,客户端中途断开、hyper 丢掉响应体是另一种 —— 两种
Expand All @@ -98,7 +104,24 @@ pub(super) fn respond(
let mut ending = ending;
let mut chunks = std::pin::pin!(chunks);
let mut broke: Option<GatewayError> = None;
while let Some(item) = chunks.next().await {
loop {
let next = match ping_every {
None => chunks.next().await,
// `next()` 被超时丢掉不丢数据:它只是去问一次流,没拿走任何东西
Some(every) => match tokio::time::timeout(every, chunks.next()).await {
Ok(next) => next,
Err(_) => {
// **不经过留档、计量和审查**:心跳不是上游说的话,不进请求记录,
// 也不算输出。**只在帧的边界上插**,上游停在一帧中间时插进去
// 会把那一帧拆坏 —— 那时宁可不补
if relay.between_frames() {
yield Ok::<Bytes, std::io::Error>(Bytes::from_static(PING));
}
continue;
}
},
};
let Some(item) = next else { break };
match item {
Ok(chunk) => {
// **旁路嗅探和留档,不缓冲**:字节照常流向客户端,同时
Expand Down Expand Up @@ -161,6 +184,9 @@ pub(super) fn respond(
resp
}

/// Anthropic 的心跳帧,和它自己的 API 发的一样。
const PING: &[u8] = b"event: ping\ndata: {\"type\": \"ping\"}\n\n";

/// 响应头到手时就定下的处理方式。
#[derive(Clone, Copy)]
struct Plan {
Expand Down Expand Up @@ -317,6 +343,8 @@ struct Relay {
/// **切断时要把数组收好**(见 [`Relay::error_tail`])
array_opened: bool,
array_element: bool,
/// 发给客户端的最后一段停在帧的边界上(或者还什么都没发)。心跳只能插在这里
at_boundary: bool,
bus: tw_observe::EventBus,
id: u64,
provider: String,
Expand Down Expand Up @@ -396,6 +424,7 @@ impl Relay {
client_dialect,
array_opened: false,
array_element: false,
at_boundary: true,
bus: state.bus.clone(),
id,
provider: provider.name.clone(),
Expand Down Expand Up @@ -601,9 +630,13 @@ impl Relay {
(tail, None)
}

/// 记下发给客户端的这一段。**只有直通的 JSON 数组流要记**:切断的位置总在元素
/// 边界上(分隔符算在后面那个元素上),所以只要知道 `[` 之后有没有过 `{`
/// 记下发给客户端的这一段:停没停在帧的边界上(心跳要看)。直通的 JSON 数组流
/// 还要记数组发到哪儿了:切断的位置总在元素边界上(分隔符算在后面那个元素上),
/// 所以只要知道 `[` 之后有没有过 `{`
fn sent(&mut self, out: &[u8]) {
if !out.is_empty() {
self.at_boundary = out.ends_with(b"\n\n") || out.ends_with(b"\r\n\r\n");
}
if !self.plan.client_json_stream || self.session.is_some() {
return;
}
Expand All @@ -619,6 +652,11 @@ impl Relay {
}
}

/// 现在插一帧会不会拆坏客户端正在收的那一帧。
fn between_frames(&self) -> bool {
self.at_boundary
}

/// 流断了之后还能对客户端说的最后一句:按它收到的格式收尾。
fn error_tail(&mut self, err: &GatewayError) -> Option<Vec<u8>> {
if let Some(c) = self.back.as_mut() {
Expand Down
4 changes: 4 additions & 0 deletions crates/tw-gateway/src/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -204,6 +204,9 @@ pub struct AppState {
pub live: crate::live::Live,
/// 每段对话此刻归到哪一次会话(见 [`crate::session::Sessions`])。**跨重载存活**
pub sessions: Arc<crate::session::Sessions>,
/// Anthropic 流里上游静默多久就补一个 `ping`(见 `relay`)。**测试会把它调短**,
/// 否则一条心跳的测试要干等十五秒
pub ping_every: std::time::Duration,
}

impl AppState {
Expand Down Expand Up @@ -252,6 +255,7 @@ impl AppState {
proxies: Arc::new(std::sync::Mutex::new(Default::default())),
live: crate::live::Live::default(),
sessions: Default::default(),
ping_every: crate::PING_EVERY,
};
// 手写的清单马上可用;向上游问是后台的事,不挡启动
state.publish_catalog();
Expand Down
Loading
Loading