diff --git a/crates/tw-gateway/src/glm.rs b/crates/tw-gateway/src/glm.rs index d03cf387..6baa79c7 100644 --- a/crates/tw-gateway/src/glm.rs +++ b/crates/tw-gateway/src/glm.rs @@ -7,7 +7,8 @@ //! 问的纪律: //! - **界面来要时问**,60 秒内的多次合成一次;**有请求去这一家时顺手问**,5 分钟最多一次。 //! - 失败了退避(30 秒、60 秒、120 秒、300 秒),不跟着界面的每一次刷新重试。 -//! - 这把 key 没有开通套餐:记成「没有额度数据」,隔一小时才再问。 +//! - 这把 key 没有开通套餐:记成「没有额度数据」,隔一小时才再问。刚才还有额度的 key +//! 要连着说两次才算(一次临时的 500 和「没有套餐」长得一样)。 //! //! 额度用完时模型请求回 429,业务码在 body 里([`exhausted`])。**那一刻就记成用完**, //! 不等下一次问额度。 @@ -324,6 +325,8 @@ struct Slot { quiet_until: u64, /// 正在问。**同一时刻只问一次** running: bool, + /// 上一次问来的是额度 + had_quota: bool, } impl Slot { @@ -339,10 +342,20 @@ impl Slot { .is_none_or(|t| now_ms.saturating_sub(t) >= every) } - fn settle(&mut self, answer: &Answer, now_ms: u64) { + /// 记下这一次的结果,交回该按哪个结论办。 + /// + /// **刚才还有额度的 key 头一次说「没有套餐」,当成一次失败**:「没有套餐」认的是 + /// 业务码 500,而一次临时的 500 长得一模一样。照「没有套餐」办要清掉额度、一小时 + /// 不再问;当成失败只是按退避过一会儿再问。连着第二次还这么说,才是真没有了 + fn settle(&mut self, answer: Answer, now_ms: u64) -> Answer { + let answer = match answer { + Answer::NoPlan if self.had_quota => Answer::Failed, + a => a, + }; + self.had_quota = matches!(answer, Answer::Quota(_)); self.running = false; self.asked_at = Some(now_ms); - match answer { + match &answer { Answer::Quota(_) => { self.failures = 0; self.quiet_until = 0; @@ -357,6 +370,7 @@ impl Slot { self.quiet_until = now_ms + BACKOFF_MS[i]; } } + answer } } @@ -398,14 +412,18 @@ impl Tracker { due } - /// 问完了。凭据在这中间换了的话,这个结果说的是旧的 key,不记 - pub fn settle(&self, provider: &str, ident: &str, answer: &Answer, now_ms: u64) { - if let Ok(mut g) = self.slots.lock() - && let Some(slot) = g.get_mut(provider) - && slot.ident == ident - { - slot.settle(answer, now_ms); - } + /// 问完了:记下节奏,交回该按哪个结论办(见 `Slot::settle`)。凭据在这中间换了 + /// 的话,这个结果说的是旧的 key,不记也不办,是 None + pub fn settle( + &self, + provider: &str, + ident: &str, + answer: Answer, + now_ms: u64, + ) -> Option { + let mut g = self.slots.lock().ok()?; + let slot = g.get_mut(provider).filter(|s| s.ident == ident)?; + Some(slot.settle(answer, now_ms)) } } @@ -749,7 +767,7 @@ mod tests { assert!(t.claim("glm", "k", Why::Demand, NOW)); // 正在问:再来的不另问 assert!(!t.claim("glm", "k", Why::Demand, NOW + 1)); - t.settle("glm", "k", &Answer::Quota(Quota::default()), NOW + 100); + t.settle("glm", "k", Answer::Quota(Quota::default()), NOW + 100); assert!(!t.claim("glm", "k", Why::Demand, NOW + 30_000)); assert!(t.claim("glm", "k", Why::Demand, NOW + 100 + MINUTE_MS)); } @@ -758,14 +776,14 @@ mod tests { fn traffic_asks_at_most_every_five_minutes() { let t = Tracker::default(); assert!(t.claim("glm", "k", Why::Traffic, NOW)); - t.settle("glm", "k", &Answer::Quota(Quota::default()), NOW); + t.settle("glm", "k", Answer::Quota(Quota::default()), NOW); assert!(!t.claim("glm", "k", Why::Traffic, NOW + 4 * MINUTE_MS)); // 界面来要的照样按一分钟算 assert!(t.claim("glm", "k", Why::Demand, NOW + 2 * MINUTE_MS)); t.settle( "glm", "k", - &Answer::Quota(Quota::default()), + Answer::Quota(Quota::default()), NOW + 2 * MINUTE_MS, ); assert!(!t.claim("glm", "k", Why::Traffic, NOW + 6 * MINUTE_MS)); @@ -778,16 +796,16 @@ mod tests { let mut now = NOW; for wait in [30_000, 60_000, 120_000, 300_000, 300_000] { assert!(t.claim("glm", "k", Why::Demand, now), "at {}", now - NOW); - t.settle("glm", "k", &Answer::Failed, now); + t.settle("glm", "k", Answer::Failed, now); // 界面每分钟来要一次也不提前 assert!(!t.claim("glm", "k", Why::Demand, now + wait - 1)); now += wait.max(DEMAND_EVERY_MS); } // 成功一次,退避清零 assert!(t.claim("glm", "k", Why::Demand, now)); - t.settle("glm", "k", &Answer::Quota(Quota::default()), now); + t.settle("glm", "k", Answer::Quota(Quota::default()), now); assert!(t.claim("glm", "k", Why::Demand, now + DEMAND_EVERY_MS)); - t.settle("glm", "k", &Answer::Rejected, now + DEMAND_EVERY_MS); + t.settle("glm", "k", Answer::Rejected, now + DEMAND_EVERY_MS); assert!(t.claim("glm", "k", Why::Demand, now + 2 * DEMAND_EVERY_MS)); } @@ -795,23 +813,45 @@ mod tests { fn a_key_without_a_plan_is_asked_again_only_after_an_hour_or_a_new_key() { let t = Tracker::default(); assert!(t.claim("glm", "k", Why::Demand, NOW)); - t.settle("glm", "k", &Answer::NoPlan, NOW); + t.settle("glm", "k", Answer::NoPlan, NOW); assert!(!t.claim("glm", "k", Why::Demand, NOW + 30 * MINUTE_MS)); assert!(!t.claim("glm", "k", Why::Traffic, NOW + 30 * MINUTE_MS)); // 换了 key:之前的结论说的是旧的那把 assert!(t.claim("glm", "k2", Why::Demand, NOW + 30 * MINUTE_MS)); - t.settle("glm", "k2", &Answer::NoPlan, NOW + 30 * MINUTE_MS); + t.settle("glm", "k2", Answer::NoPlan, NOW + 30 * MINUTE_MS); assert!(t.claim("glm", "k2", Why::Demand, NOW + 90 * MINUTE_MS)); } + /// 刚才还有额度的 key 头一次说「没有套餐」:多半是一次临时的 500,当成失败, + /// 过一会儿再问;连着第二次才信 + #[test] + fn a_key_that_had_a_plan_must_say_no_plan_twice() { + let t = Tracker::default(); + let q = Answer::Quota(Quota::default()); + assert!(t.claim("glm", "k", Why::Demand, NOW)); + assert_eq!(t.settle("glm", "k", q.clone(), NOW), Some(q)); + assert!(t.claim("glm", "k", Why::Demand, NOW + MINUTE_MS)); + assert_eq!( + t.settle("glm", "k", Answer::NoPlan, NOW + MINUTE_MS), + Some(Answer::Failed) + ); + let later = NOW + 2 * MINUTE_MS; + assert!(t.claim("glm", "k", Why::Demand, later)); + assert_eq!( + t.settle("glm", "k", Answer::NoPlan, later), + Some(Answer::NoPlan) + ); + assert!(!t.claim("glm", "k", Why::Demand, later + 30 * MINUTE_MS)); + } + #[test] fn an_answer_for_a_replaced_key_is_not_kept() { let t = Tracker::default(); assert!(t.claim("glm", "old", Why::Demand, NOW)); assert!(t.claim("glm", "new", Why::Demand, NOW)); - t.settle("glm", "old", &Answer::NoPlan, NOW); + t.settle("glm", "old", Answer::NoPlan, NOW); // 新 key 还在问,旧的结论没把它记成「没有套餐」 - t.settle("glm", "new", &Answer::Quota(Quota::default()), NOW); + t.settle("glm", "new", Answer::Quota(Quota::default()), NOW); assert!(t.claim("glm", "new", Why::Demand, NOW + DEMAND_EVERY_MS)); } } diff --git a/crates/tw-gateway/src/lib.rs b/crates/tw-gateway/src/lib.rs index f032d52d..4457cd9a 100644 --- a/crates/tw-gateway/src/lib.rs +++ b/crates/tw-gateway/src/lib.rs @@ -64,6 +64,13 @@ pub use state::{AppState, Runtime}; /// 那个上限,又不至于让一条正常的流塞满心跳 pub const PING_EVERY: std::time::Duration = std::time::Duration::from_secs(15); +/// 上游一个字节都不发(连 SSE 注释都没有)超过这么久,就不再替它补 `ping`。 +/// +/// 心跳让客户端相信连接还活着;上游的连接半开了(NAT、代理把它丢了)时,网关自己 +/// 永远等不到下一个字节,**一直补下去这个请求就永远挂着**。十分钟之后停下,客户端 +/// 自己的静默计时会断开它。正常在排队的上游会发 `: keep-alive` 这样的注释,不受影响 +pub const PING_FOR: std::time::Duration = std::time::Duration::from_secs(600); + /// 请求来自谁。**如实写 ThinkWatch** —— 我们从不把自己报成别的客户端。 pub const ORIGINATOR: &str = "thinkwatch"; diff --git a/crates/tw-gateway/src/server/pipeline/hop.rs b/crates/tw-gateway/src/server/pipeline/hop.rs index aa352cc0..58bc7975 100644 --- a/crates/tw-gateway/src/server/pipeline/hop.rs +++ b/crates/tw-gateway/src/server/pipeline/hop.rs @@ -347,8 +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(); + // DeepSeek Harness 的扩展只有 DeepSeek 官方认。**发给它时一个字节都不改**。 + // 不止生成请求:数 token 这样的请求一样带着它的头,也可能带着会话日志 + let harness = 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(); @@ -525,11 +526,17 @@ async fn send( // 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 beta = (harness && target.is_none()) + .then(|| { + headers + .get_all("anthropic-beta") + .iter() + .filter_map(|v| v.to_str().ok()) + .collect::>() + .join(",") + }) + .and_then(|v| tw_dialect::harness::anthropic_beta(&v)); let build = |upstream_headers: &[(String, String)]| { let mut req = http.request(method.clone(), &url); req = forward::forward_headers_filtered(req, headers, |n| { diff --git a/crates/tw-gateway/src/server/pipeline/relay.rs b/crates/tw-gateway/src/server/pipeline/relay.rs index a28fa18f..6bddbb09 100644 --- a/crates/tw-gateway/src/server/pipeline/relay.rs +++ b/crates/tw-gateway/src/server/pipeline/relay.rs @@ -97,6 +97,7 @@ pub(super) fn respond( && plan.client_sse && status.is_success()) .then_some(state.ping_every); + let ping_for = state.ping_for; let stream = async_stream::stream! { // **通行证跟着响应体走。**这个流被丢掉的时候它才还回去:正常 // 发完是一种,客户端中途断开、hyper 丢掉响应体是另一种 —— 两种 @@ -106,17 +107,27 @@ pub(super) fn respond( let mut ending = ending; let mut chunks = std::pin::pin!(chunks); let mut broke: Option = None; + // 客户端上一次收到字节的时刻。**心跳看的是客户端那一边的静默**,不是上游的: + // 上游在排队时发的 SSE 注释(DeepSeek 的 `: keep-alive`、OpenRouter 的 + // `: OPENROUTER PROCESSING`)转换时被丢掉,上游不算静默,客户端却什么都没收到 + let mut quiet_since = tokio::time::Instant::now(); + // 上游上一次发来任何字节(注释也算)的时刻。**上游整个没了声音太久就不再补**: + // 半开的连接永远等不到下一个字节,一直补心跳的话客户端也永远不会放弃,这个 + // 请求就挂在那儿了。停下之后由客户端自己的静默计时来断 + let mut upstream_since = tokio::time::Instant::now(); loop { let next = match ping_every { None => chunks.next().await, // `next()` 被超时丢掉不丢数据:它只是去问一次流,没拿走任何东西 - Some(every) => match tokio::time::timeout(every, chunks.next()).await { + Some(every) => match tokio::time::timeout_at(quiet_since + every, chunks.next()).await { Ok(next) => next, Err(_) => { + // 停在一帧中间时这一轮不补,也要重新计时,否则会原地空转 + quiet_since = tokio::time::Instant::now(); // **不经过留档、计量和审查**:心跳不是上游说的话,不进请求记录, // 也不算输出。**只在帧的边界上插**,上游停在一帧中间时插进去 // 会把那一帧拆坏 —— 那时宁可不补 - if relay.between_frames() { + if relay.between_frames() && upstream_since.elapsed() < ping_for { yield Ok::(Bytes::from_static(PING)); } continue; @@ -126,6 +137,7 @@ pub(super) fn respond( let Some(item) = next else { break }; match item { Ok(chunk) => { + upstream_since = tokio::time::Instant::now(); // **旁路嗅探和留档,不缓冲**:字节照常流向客户端,同时 // 喂它一份。上游返回的 usage 是真相,而拿不到它就只能估。 // @@ -136,6 +148,7 @@ pub(super) fn respond( let (out, cut) = relay.chunk(&chunk); relay.sent(&out); if !out.is_empty() { + quiet_since = tokio::time::Instant::now(); yield Ok::(Bytes::from(out)); } if let Some(err) = cut { diff --git a/crates/tw-gateway/src/state.rs b/crates/tw-gateway/src/state.rs index 8202749a..84eb9e11 100644 --- a/crates/tw-gateway/src/state.rs +++ b/crates/tw-gateway/src/state.rs @@ -210,6 +210,8 @@ pub struct AppState { /// Anthropic 流里上游静默多久就补一个 `ping`(见 `relay`)。**测试会把它调短**, /// 否则一条心跳的测试要干等十五秒 pub ping_every: std::time::Duration, + /// 上游整个静默(连注释都没有)超过这么久,就不再补 `ping`(见 `relay`)。**测试会把它调短** + pub ping_for: std::time::Duration, } impl AppState { @@ -260,6 +262,7 @@ impl AppState { live: crate::live::Live::default(), sessions: Default::default(), ping_every: crate::PING_EVERY, + ping_for: crate::PING_FOR, }; // 手写的清单马上可用;向上游问是后台的事,不挡启动 state.publish_catalog(); diff --git a/crates/tw-gateway/src/state/glm.rs b/crates/tw-gateway/src/state/glm.rs index dd480fd6..0d670ed1 100644 --- a/crates/tw-gateway/src/state/glm.rs +++ b/crates/tw-gateway/src/state/glm.rs @@ -22,8 +22,12 @@ impl AppState { /// 这一家是 GLM Coding Plan 的话,额度接口的地址和凭据的指纹。 /// - /// 没有 key 的不算:额度跟着 key 走,没有 key 就没有可问的 + /// 没有 key 的不算:额度跟着 key 走,没有 key 就没有可问的。**停用的也不算**: + /// 用户不想用它了,不该再拿它的 key 去问 fn glm_target(&self, p: &tw_config::Provider) -> Option<(String, String)> { + if p.disabled { + return None; + } let key = p.key.as_ref()?; let url = self.glm.sites().quota_url(&p.base_url)?; // 指纹,不是 key 本身:它只用来认出「换了 key」 @@ -67,6 +71,11 @@ impl AppState { let p = p.clone(); Some(tokio::spawn(async move { let answer = state.ask_glm(&p, &url).await; + // 节奏先记上,**按记下的那个结论办**:刚才还有套餐的 key 头一次说没有, + // 当成一次失败(见 `Tracker::settle`) + let Some(answer) = state.glm.settle(&p.name, &ident, answer, now_ms()) else { + return; + }; match &answer { Answer::Quota(q) => state.record_quota(state.bus.next_id(), &p.name, q.clone()), // 没有套餐就没有额度数据:之前记着的也不再作数 @@ -78,7 +87,6 @@ impl AppState { tracing::debug!(provider = %p.name, "the GLM quota could not be read") } } - state.glm.settle(&p.name, &ident, &answer, now_ms()); })) } @@ -279,5 +287,8 @@ mod tests { let mut keyless = cfg.providers[0].clone(); keyless.key = None; assert!(s.glm_target(&keyless).is_none()); + let mut disabled = cfg.providers[0].clone(); + disabled.disabled = true; + assert!(s.glm_target(&disabled).is_none()); } } diff --git a/crates/tw-gateway/tests/claude_desktop.rs b/crates/tw-gateway/tests/claude_desktop.rs index a566ccfd..2d4dfa34 100644 --- a/crates/tw-gateway/tests/claude_desktop.rs +++ b/crates/tw-gateway/tests/claude_desktop.rs @@ -97,6 +97,10 @@ struct Gateway { } async fn gateway(p: Provider) -> Gateway { + gateway_pinging_for(p, tw_gateway::PING_FOR).await +} + +async fn gateway_pinging_for(p: Provider, ping_for: Duration) -> Gateway { let cfg = Config { version: 1, listen: Listen::default(), @@ -111,6 +115,7 @@ async fn gateway(p: Provider) -> Gateway { let mut state = tw_gateway::AppState::new(cfg).unwrap(); // 心跳的间隔调短,一条测试不必干等十五秒 state.ping_every = Duration::from_millis(100); + state.ping_for = ping_for; let events = state.bus.subscribe(); let (tx, bodies) = tokio::sync::mpsc::channel(16); state.set_body_sink(tx); @@ -343,6 +348,68 @@ async fn a_silent_converted_upstream_is_covered_with_pings() { assert!(!recorded.contains("ping"), "心跳进了请求记录:{recorded}"); } +/// 上游在排队时只发 SSE 注释(DeepSeek 的 `: keep-alive`)。转换时注释被丢掉, +/// 客户端什么都收不到 —— **心跳要看客户端那一边的静默**,不能被上游的注释推迟 +#[tokio::test] +async fn upstream_comments_that_never_reach_the_client_do_not_hold_pings_back() { + let mut parts = vec![Some(CHAT_FIRST)]; + for _ in 0..14 { + parts.push(Some(": keep-alive\n\n")); + parts.push(None); + } + parts.push(Some(CHAT_REST)); + let up = stalling_upstream(parts, Duration::from_millis(50)).await; + let gw = gateway(provider(up, Protocol::OpenaiChat, &[])).await; + let text = stream_from( + &gw, + "/v1/messages", + json!({"model": "gpt-x", "max_tokens": 16, "stream": true, + "messages": [{"role": "user", "content": "hi"}]}), + ) + .await; + let ping = "event: ping\ndata: {\"type\": \"ping\"}\n\n"; + assert!(text.matches(ping).count() >= 2, "{text}"); + let rest = text.replace(ping, ""); + assert!(!rest.contains("keep-alive"), "{rest}"); + assert!( + rest.trim_end().ends_with("\"type\":\"message_stop\"}"), + "{rest}" + ); +} + +/// 上游一个字节都不发太久(连接半开了):**不再补心跳**,让客户端自己的静默计时断开它。 +/// 一直补的话这个请求永远挂着 +#[tokio::test] +async fn pings_stop_once_the_upstream_has_been_silent_too_long() { + let up = stalling_upstream( + vec![Some(CHAT_FIRST), None, Some(CHAT_REST)], + Duration::from_millis(1200), + ) + .await; + let gw = gateway_pinging_for( + provider(up, Protocol::OpenaiChat, &[]), + Duration::from_millis(350), + ) + .await; + let text = stream_from( + &gw, + "/v1/messages", + json!({"model": "gpt-x", "max_tokens": 16, "stream": true, + "messages": [{"role": "user", "content": "hi"}]}), + ) + .await; + let ping = "event: ping\ndata: {\"type\": \"ping\"}\n\n"; + // 静默 1.2 秒、每 0.1 秒一次:一直补的话有十一个。只补前 0.35 秒的 + let n = text.matches(ping).count(); + assert!((1..=4).contains(&n), "{n}: {text}"); + assert!( + text.replace(ping, "") + .trim_end() + .ends_with("\"type\":\"message_stop\"}"), + "{text}" + ); +} + #[tokio::test] async fn a_ping_never_splits_a_frame_the_upstream_left_half_sent() { const START: &str = "event: message_start\ndata: {\"type\":\"message_start\",\"message\":{\"id\":\"m1\",\"type\":\"message\",\"role\":\"assistant\",\"model\":\"claude-sonnet-4-5\",\"content\":[],\"usage\":{\"input_tokens\":5,\"output_tokens\":1}}}\n\n"; diff --git a/crates/tw-gateway/tests/harness.rs b/crates/tw-gateway/tests/harness.rs index 8c6f83b6..c29bf815 100644 --- a/crates/tw-gateway/tests/harness.rs +++ b/crates/tw-gateway/tests/harness.rs @@ -345,6 +345,57 @@ async fn messages_to_another_anthropic_upstream_are_cleaned_in_place() { ); } +/// `anthropic-beta` 分两行发:清理的那一行去掉,**另一行照发** +#[tokio::test] +async fn betas_sent_on_separate_lines_all_survive_the_cleaning() { + let (gw, _rx, log) = gateway(elsewhere(Protocol::Anthropic).await).await; + let (status, body) = send( + gw, + "/v1/messages", + &[ + ("anthropic-version", "2023-06-01"), + ("anthropic-beta", TOOL_BETA), + ("anthropic-beta", "files-api-2025-04-14"), + ], + &messages_request().to_string(), + ) + .await; + assert_eq!(status, 200, "{body}"); + let s = generation(&log); + assert_clean(&s); + let betas: Vec<_> = s.headers.get_all("anthropic-beta").iter().collect(); + assert_eq!(betas, ["files-api-2025-04-14"]); +} + +/// 不生成回答的请求(数 token)也是 dsh 发的:发给别家时一样不带它的头和会话日志 +#[tokio::test] +async fn counting_tokens_elsewhere_is_cleaned_too() { + let (gw, _rx, log) = gateway(elsewhere(Protocol::Anthropic).await).await; + let mut req = messages_request(); + req.as_object_mut().unwrap().remove("max_tokens"); + let (status, body) = send( + gw, + "/v1/messages/count_tokens", + &[ + ("anthropic-version", "2023-06-01"), + ("anthropic-beta", TOOL_BETA), + ], + &req.to_string(), + ) + .await; + assert_eq!(status, 200, "{body}"); + let s = log + .lock() + .unwrap() + .iter() + .rev() + .find(|s| s.uri.ends_with("/count_tokens")) + .cloned() + .expect("the upstream got no count_tokens request"); + assert_clean(&s); + assert!(s.headers.get("anthropic-beta").is_none()); +} + #[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;