From 3edde89f1063561921c9e427a2bd279b437e3141 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 25 Sep 2026 16:12:52 +0000 Subject: [PATCH] feat(quota): read the GLM Coding Plan quota for Z.ai and BigModel upstreams GLM Coding Plan responses carry no quota headers, so the quota view stayed empty for Z.ai / BigModel accounts. The only source is the account's quota endpoint (/api/monitor/usage/quota/limit), and a used-up window only shows up as a business code inside a 429 body. - An upstream whose base_url host is api.z.ai or open.bigmodel.cn is a GLM Coding Plan upstream, whether the login wrote it or it was added by hand. - GET /quota asks each due GLM upstream (at most once a minute) and waits up to 5s; traffic to one asks at most every 5 minutes. Failures back off 30s/60s/120s/300s, and a key without a plan is left alone for an hour. - Windows are told apart by unit/number, never by reset order: 5h, weekly, and monthly (the old plans' MCP calls). A reset time past the window's length is dropped rather than shown as a wrong countdown. - Credit-based plans carry total/used/remaining as the endpoint gives them (QuotaWindow.credits, new QuotaCredits type). - A 429 with 1308, 1310 or 1316-1321 marks the window used up and emits QuotaExhausted; the reset time comes from the quota endpoint when known, otherwise from the message (read as UTC+8). 1309, 1313, 1302 and 1305 are not a used-up quota. CONTROL_API_VERSION is now 23: QuotaWindow gained `credits` and the window vocabulary gained `monthly`. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01WjXngih1kA6oBMoqCkXDVx --- crates/tw-api/src/lib.rs | 26 +- crates/tw-control/src/chatgpt.rs | 2 + crates/tw-control/src/lib.rs | 18 +- crates/tw-gateway/src/glm.rs | 817 ++++++++++++++++++ crates/tw-gateway/src/lib.rs | 1 + crates/tw-gateway/src/quota.rs | 24 +- crates/tw-gateway/src/server/pipeline/hop.rs | 5 + .../tw-gateway/src/server/pipeline/relay.rs | 2 + crates/tw-gateway/src/state.rs | 4 + crates/tw-gateway/src/state/glm.rs | 283 ++++++ crates/tw-gateway/src/state/upstream.rs | 21 +- crates/tw-gateway/tests/glm_quota.rs | 222 +++++ 12 files changed, 1400 insertions(+), 25 deletions(-) create mode 100644 crates/tw-gateway/src/glm.rs create mode 100644 crates/tw-gateway/src/state/glm.rs create mode 100644 crates/tw-gateway/tests/glm_quota.rs diff --git a/crates/tw-api/src/lib.rs b/crates/tw-api/src/lib.rs index 1c337f42..d5231b72 100644 --- a/crates/tw-api/src/lib.rs +++ b/crates/tw-api/src/lib.rs @@ -589,7 +589,11 @@ pub const MSG_CODES: &str = include_str!("../msg-codes.txt"); /// ChatGPT 登录的结果([`ChatgptLoginStatus`])不再带 `plan`,换成 `account`:和 /// [`OAuthView::account`] 同一块(邮箱、套餐),从同一个 access token 里读。照 21 写的 /// 界面会把 `/summary/routes` 当成数组去读,在登录结果里找不到套餐。 -pub const CONTROL_API_VERSION: u32 = 22; +/// +/// **23 起额度也来自 GLM Coding Plan**(Z.ai / BigModel 的上游):窗口多了 `monthly`, +/// 积分制套餐的窗口带上 [`QuotaWindow::credits`](总额、已用、剩余)。照 22 写的界面 +/// 不认 `monthly`,也看不到剩余积分。 +pub const CONTROL_API_VERSION: u32 = 23; #[derive(Debug, Clone, Serialize, Deserialize)] #[cfg_attr(feature = "ts", derive(ts_rs::TS))] @@ -1232,14 +1236,15 @@ pub struct RoutingView { pub attempts: Vec, } -/// 一个订阅额度窗口。**每个字段都直接来自上游的响应头。** +/// 一个订阅额度窗口。**每个字段都直接来自上游**:响应头,或者账号的额度接口。 /// /// 我们自己推断的东西不放进这个结构 —— 界面上必须能区分「上游说的」和 /// 「我们猜的」,而混在一个类型里就区分不了了。 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] #[cfg_attr(feature = "ts", derive(ts_rs::TS))] pub struct QuotaWindow { - /// `5h` / `7d`(Anthropic)/ `weekly`(Codex) + /// `5h` / `7d`(Anthropic)/ `weekly`(Codex、GLM)/ `monthly`(GLM 老套餐每月的 + /// MCP 调用次数) pub window: String, pub used_percent: f64, /// 什么时候重置,Unix 毫秒。**是时刻,不是「还有多少秒」**:上游报的秒数 @@ -1249,6 +1254,21 @@ pub struct QuotaWindow { pub resets_at_ms: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub status: Option, + /// 积分制套餐(GLM Coding Plan)这个窗口的积分。别的套餐没有 + #[serde(default, skip_serializing_if = "Option::is_none")] + pub credits: Option, +} + +/// 一个额度窗口的积分,三个数都是上游给的原数。 +/// +/// **剩余不是总额减已用算出来的**:上游给的三个数不一定对得上(实测总额 2000、已用 23、 +/// 剩余 1976),界面要显示剩余就显示它说的剩余。 +#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)] +#[cfg_attr(feature = "ts", derive(ts_rs::TS))] +pub struct QuotaCredits { + pub total: f64, + pub used: f64, + pub remaining: f64, } /// 一次调用的用量。 diff --git a/crates/tw-control/src/chatgpt.rs b/crates/tw-control/src/chatgpt.rs index d0027abb..510d3c76 100644 --- a/crates/tw-control/src/chatgpt.rs +++ b/crates/tw-control/src/chatgpt.rs @@ -998,6 +998,7 @@ async fn usage( used_percent: w.used_percent, resets_at_ms: w.resets_at_ms, status: w.status.clone(), + credits: None, }) .collect(), }, @@ -1022,6 +1023,7 @@ fn parse_usage(v: &Value, now_ms: u64) -> tw_api::ChatgptUsage { .and_then(|x| x.as_u64()) .map(|secs| now_ms.saturating_add(secs.saturating_mul(1000))), status: (used >= 100.0).then(|| "rejected".to_string()), + credits: None, }) }; let limits = &v["rate_limit"]; diff --git a/crates/tw-control/src/lib.rs b/crates/tw-control/src/lib.rs index 846986d1..d1897877 100644 --- a/crates/tw-control/src/lib.rs +++ b/crates/tw-control/src/lib.rs @@ -984,29 +984,27 @@ async fn request_detail( /// /// 按量付费的账号没有这些头,那时这个列表是空的 —— 界面据此决定显示 /// 金额还是百分比,两种人格共用同一块地方。 +/// +/// GLM Coding Plan 的额度不在响应头里:**界面来要时去问**(60 秒内合成一次),问得慢的 +/// 不等,问完由 `QuotaSeen` 补上 async fn quota(State(s): State) -> Json> { + s.gateway.refresh_glm_quotas(GLM_QUOTA_WAIT).await; let mut out: Vec = s .gateway .quotas() .into_iter() .map(|(provider, q)| tw_api::ProviderQuota { provider, - windows: q - .windows - .into_iter() - .map(|w| tw_api::QuotaWindow { - window: w.window, - used_percent: w.used_percent, - resets_at_ms: w.resets_at_ms, - status: w.status, - }) - .collect(), + windows: q.windows.iter().map(Into::into).collect(), }) .collect(); out.sort_by(|a, b| a.provider.cmp(&b.provider)); Json(out) } +/// 界面来要额度时,最多等 GLM 的额度接口多久 +const GLM_QUOTA_WAIT: std::time::Duration = std::time::Duration::from_secs(5); + /// 按上游分的延迟。**「哪家 TTFT 最差」问的是这个。** /// /// 和按模型分是两个问题:前者的下一步是换上游,后者是换模型。 diff --git a/crates/tw-gateway/src/glm.rs b/crates/tw-gateway/src/glm.rs new file mode 100644 index 00000000..d03cf387 --- /dev/null +++ b/crates/tw-gateway/src/glm.rs @@ -0,0 +1,817 @@ +//! GLM Coding Plan(Z.ai / BigModel)的套餐额度。 +//! +//! **响应头里没有额度**,和 Anthropic、Codex 不一样:只能去问账号的额度接口 +//! (`/api/monitor/usage/quota/limit`)。这个接口没有公开文档,形状照 Z.ai 官方的 +//! 开源工具读出来的逻辑写。 +//! +//! 问的纪律: +//! - **界面来要时问**,60 秒内的多次合成一次;**有请求去这一家时顺手问**,5 分钟最多一次。 +//! - 失败了退避(30 秒、60 秒、120 秒、300 秒),不跟着界面的每一次刷新重试。 +//! - 这把 key 没有开通套餐:记成「没有额度数据」,隔一小时才再问。 +//! +//! 额度用完时模型请求回 429,业务码在 body 里([`exhausted`])。**那一刻就记成用完**, +//! 不等下一次问额度。 + +use std::collections::HashMap; +use std::sync::Mutex; + +use serde_json::Value; + +use crate::quota::{Quota, Window}; + +/// 额度接口在两个站上的路径 +const QUOTA_PATH: &str = "/api/monitor/usage/quota/limit"; + +/// 额度接口所在的两个站。**平时是它们的**,测试里换成本机的假服务器 +#[derive(Debug, Clone)] +pub struct Sites { + /// `https://api.z.ai` + pub zai: String, + /// `https://open.bigmodel.cn` + pub bigmodel: String, +} + +impl Default for Sites { + fn default() -> Self { + Self { + zai: "https://api.z.ai".to_string(), + bigmodel: "https://open.bigmodel.cn".to_string(), + } + } +} + +impl Sites { + /// 这个上游是 GLM Coding Plan 的话,它的额度接口地址。 + /// + /// **只按主机认**:配置里没有专门的标记,登录生成的和手动添加的都只能从 `base_url` + /// 看出来。Anthropic 格式(`/api/anthropic`)和 OpenAI 格式(`/api/coding/paas/v4`) + /// 的地址都在同一个主机上,额度跟着 key 走,两种都算 + pub fn quota_url(&self, base_url: &str) -> Option { + let at = host_of(base_url)?; + [&self.zai, &self.bigmodel] + .into_iter() + .find(|site| host_of(site).as_ref() == Some(&at)) + .map(|site| format!("{}{QUOTA_PATH}", site.trim_end_matches('/'))) + } +} + +/// 主机(小写)和写明了的端口。**不看协议**:`http://api.z.ai` 也是那一家 +fn host_of(url: &str) -> Option<(String, Option)> { + let u = reqwest::Url::parse(url.trim()).ok()?; + Some((u.host_str()?.to_ascii_lowercase(), u.port())) +} + +// ---------------------------------------------------------------- 额度接口的回答 + +/// 问一次额度接口的结果。 +#[derive(Debug, Clone, PartialEq)] +pub enum Answer { + /// 套餐额度,至少有一个认得出来的窗口 + Quota(Quota), + /// 这把 key 没有开通套餐,或者回答里没有一个认得出来的窗口:**没有额度数据** + NoPlan, + /// 接口不认这把 key + Rejected, + /// 别的失败:5xx、读不懂、业务码不认识 + Failed, +} + +/// 读额度接口的回答。`now_ms`:收到回答的时刻,用来筛掉说不通的重置时刻。 +/// +/// 成败只看 `success` 和 `code`。**`msg` 随语言变**(「当前用户不存在coding plan」), +/// 不能拿它判断 +pub fn parse(status: u16, body: &[u8], now_ms: u64) -> Answer { + // 鉴权失败多半是 200 里的业务码,但 HTTP 层的 401 / 403 也要认 + if status == 401 || status == 403 { + return Answer::Rejected; + } + if !(200..300).contains(&status) { + return Answer::Failed; + } + let Ok(v) = serde_json::from_slice::(body) else { + return Answer::Failed; + }; + let code = v.get("code").and_then(code_of); + let succeeded = v.get("success").and_then(Value::as_bool) != Some(false) + && matches!(code, None | Some(0) | Some(200)); + if !succeeded { + return match code { + Some(401 | 1000 | 1001) => Answer::Rejected, + // `{"code":500,"msg":"当前用户不存在coding plan","success":false}` + Some(500) => Answer::NoPlan, + _ => Answer::Failed, + }; + } + let windows = v["data"]["limits"] + .as_array() + .map(|xs| windows(xs, now_ms)) + .unwrap_or_default(); + if windows.is_empty() { + Answer::NoPlan + } else { + Answer::Quota(Quota { windows }) + } +} + +/// 业务码:数字,或者写成字符串的数字 +fn code_of(v: &Value) -> Option { + match v { + Value::Number(n) => n.as_i64(), + Value::String(s) => s.trim().parse().ok(), + _ => None, + } +} + +const MINUTE_MS: u64 = 60_000; +const HOUR_MS: u64 = 60 * MINUTE_MS; +const DAY_MS: u64 = 24 * HOUR_MS; + +/// 重置时刻可以比窗口长度晚这么多:两边的钟差一点,不该因此丢掉一个真实的时刻 +const CLOCK_SLACK_MS: u64 = MINUTE_MS; + +/// `limits` 里每一项是哪个窗口、那个窗口最长多久。 +/// +/// **按 `unit` / `number` 分,不按 `nextResetTime` 的先后猜**:周期末尾每周窗口会比 +/// 5 小时窗口先重置,按先后猜就把两个标反了。认不出来的一项不要 —— 没有名字的百分比 +/// 放上界面只会被当成另一个窗口。 +fn window_of(item: &Value) -> Option<(&'static str, u64)> { + let kind = item["type"].as_str().unwrap_or_default(); + // 5 小时和每周是 token(老套餐)或积分(积分制);每月是老套餐的 MCP 调用次数 + let usage = + kind.eq_ignore_ascii_case("TOKENS_LIMIT") || kind.eq_ignore_ascii_case("CREDIT_LIMIT"); + let calls = kind.eq_ignore_ascii_case("TIME_LIMIT"); + match (item["unit"].as_i64(), item["number"].as_i64()) { + (Some(3), Some(5)) if usage => Some(("5h", 5 * HOUR_MS)), + // 每周窗口的 `number` 见过 7 也见过 1,只认 `unit` + (Some(6), _) if usage => Some(("weekly", 7 * DAY_MS)), + (Some(5), Some(1)) if calls => Some(("monthly", 31 * DAY_MS)), + _ => None, + } +} + +fn windows(items: &[Value], now_ms: u64) -> Vec { + let mut out: Vec = Vec::new(); + for item in items { + let Some((name, span)) = window_of(item) else { + continue; + }; + // 同一个窗口出现两次:留第一个 + if out.iter().any(|w| w.window == name) { + continue; + } + let Some(used_percent) = item["percentage"].as_f64() else { + continue; + }; + let used_percent = used_percent.clamp(0.0, 100.0); + let credits = match ( + item["usage"].as_f64(), + item["currentValue"].as_f64(), + item["remaining"].as_f64(), + ) { + (Some(total), Some(used), Some(remaining)) => Some(tw_api::QuotaCredits { + total, + used, + remaining, + }), + _ => None, + }; + let spent = used_percent >= 100.0 || credits.is_some_and(|c| c.remaining <= 0.0); + out.push(Window { + window: name.to_string(), + used_percent, + // 用量为 0 的 5 小时窗口可能没有 `nextResetTime`。**比窗口还长的不要**:有报告 + // 说 5 小时窗口的重置时刻落在 5 小时以后,照着显示就是一个错的倒计时 + resets_at_ms: item["nextResetTime"] + .as_f64() + .filter(|t| t.is_finite() && *t > 0.0) + .map(|t| t as u64) + .filter(|t| *t > now_ms && *t <= now_ms + span + CLOCK_SLACK_MS), + // 用满就是被拒 —— 这是上游数字的直接结论,不是推断 + status: spent.then(|| "rejected".to_string()), + credits, + }); + } + out +} + +// ---------------------------------------------------------------- 429 里的额度用完 + +/// 一个 429 说的「额度用完了」。 +#[derive(Debug, Clone, PartialEq)] +pub struct Exhausted { + /// 上游的业务码 + pub code: i64, + /// 哪个窗口。**业务码和消息都说不清时是 None**,由调用方按已知的额度挑 + pub window: Option<&'static str>, + /// 消息里说的重置时刻。**以额度接口的为准**,那边没有才用它 + pub resets_at_ms: Option, +} + +/// 读 429 的 body:是额度用完的话,是哪个窗口、什么时候重置。 +/// +/// 业务码在 Anthropic 端点上是 `error.type`(字符串),在 OpenAI 端点上是 `error.code`。 +/// **算用完的只有这几个**:1308(5 小时)、1310(每周或每月)、1316–1321(团队超额)。 +/// 套餐到期(1309)、公平使用限制(1313)、临时限流(1302、1305)都不是额度用完, +/// 记成用完会让界面报一个重置之后也好不了的「用完了」。 +pub fn exhausted(body: &[u8]) -> Option { + let v: Value = serde_json::from_slice(body).ok()?; + let err = v.get("error")?; + let code = [err.get("type"), err.get("code")] + .into_iter() + .flatten() + .find_map(code_of)?; + if !matches!(code, 1308 | 1310 | 1316..=1321) { + return None; + } + let message = err["message"].as_str().unwrap_or_default(); + let window = window_in(message).or(match code { + 1308 => Some("5h"), + _ => None, + }); + Some(Exhausted { + code, + window, + resets_at_ms: reset_in(message), + }) +} + +/// 消息里说的是哪个窗口。中英文两种说法 +fn window_in(message: &str) -> Option<&'static str> { + let m = message.to_ascii_lowercase(); + if m.contains("5 hour") || m.contains("5-hour") || m.contains("5小时") || m.contains("5 小时") + { + Some("5h") + } else if m.contains("week") || m.contains("周") { + Some("weekly") + } else if m.contains("month") || m.contains("月") { + Some("monthly") + } else { + None + } +} + +/// 消息里的重置时刻:`2025-10-03 08:23:14`。**不带时区,是北京时间(UTC+8)** +fn reset_in(message: &str) -> Option { + const LEN: usize = "2025-10-03 08:23:14".len(); + let bytes = message.as_bytes(); + let shape = |w: &[u8]| { + w.iter().enumerate().all(|(i, b)| match i { + 4 | 7 => *b == b'-', + 10 => *b == b' ', + 13 | 16 => *b == b':', + _ => b.is_ascii_digit(), + }) + }; + let at = (0..bytes.len().checked_sub(LEN - 1)?).find(|&i| shape(&bytes[i..i + LEN]))?; + // 这一段全是 ASCII,切在哪儿都是字符边界 + let naive = + chrono::NaiveDateTime::parse_from_str(&message[at..at + LEN], "%Y-%m-%d %H:%M:%S").ok()?; + let beijing = chrono::FixedOffset::east_opt(8 * 3600)?; + naive + .and_local_timezone(beijing) + .single()? + .timestamp_millis() + .try_into() + .ok() +} + +/// 业务码和消息都没说是哪个窗口时,按已知的额度挑:候选里用得最多的那个。 +/// 一个都不知道就是每周 —— 1310 说的是「每周或每月」,新套餐只有每周 +pub fn guess_window(code: i64, known: &Quota) -> String { + let candidates: &[&str] = match code { + 1310 => &["weekly", "monthly"], + _ => &["5h", "weekly", "monthly"], + }; + known + .windows + .iter() + .filter(|w| candidates.contains(&w.window.as_str())) + .max_by(|a, b| a.used_percent.total_cmp(&b.used_percent)) + .map(|w| w.window.clone()) + .unwrap_or_else(|| "weekly".to_string()) +} + +// ---------------------------------------------------------------- 什么时候问 + +/// 为什么去问。 +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Why { + /// 界面来要 + Demand, + /// 有请求去了这一家 + Traffic, +} + +/// 界面来要时,这么久之内问过就不再问 +const DEMAND_EVERY_MS: u64 = MINUTE_MS; +/// 有请求时,这么久最多问一次 +const TRAFFIC_EVERY_MS: u64 = 5 * MINUTE_MS; +/// 连续失败时,第 n 次之后等多久 +const BACKOFF_MS: [u64; 4] = [30_000, 60_000, 120_000, 300_000]; +/// 没有开通套餐的 key 隔多久再问一次。**不是再也不问**:用户可能刚买了套餐 +const NO_PLAN_RECHECK_MS: u64 = HOUR_MS; + +/// 一个上游问额度的节奏。 +#[derive(Debug, Default, Clone)] +struct Slot { + /// 凭据的指纹。**换了 key 就从头来**:退避和「没有套餐」说的是旧的那把 + ident: String, + /// 上一次问完的时刻,不论结果 + asked_at: Option, + /// 连续失败了几次 + failures: u32, + /// 这之前不问:退避中,或者没有套餐 + quiet_until: u64, + /// 正在问。**同一时刻只问一次** + running: bool, +} + +impl Slot { + fn due(&self, why: Why, now_ms: u64) -> bool { + if self.running || now_ms < self.quiet_until { + return false; + } + let every = match why { + Why::Demand => DEMAND_EVERY_MS, + Why::Traffic => TRAFFIC_EVERY_MS, + }; + self.asked_at + .is_none_or(|t| now_ms.saturating_sub(t) >= every) + } + + fn settle(&mut self, answer: &Answer, now_ms: u64) { + self.running = false; + self.asked_at = Some(now_ms); + match answer { + Answer::Quota(_) => { + self.failures = 0; + self.quiet_until = 0; + } + Answer::NoPlan => { + self.failures = 0; + self.quiet_until = now_ms + NO_PLAN_RECHECK_MS; + } + Answer::Rejected | Answer::Failed => { + self.failures = self.failures.saturating_add(1); + let i = (self.failures as usize - 1).min(BACKOFF_MS.len() - 1); + self.quiet_until = now_ms + BACKOFF_MS[i]; + } + } + } +} + +/// 每个 GLM 上游问额度的节奏,和额度接口在哪。**跨重载存活**:改一条规则不该让 +/// 所有上游的退避清零 +#[derive(Default)] +pub struct Tracker { + sites: Mutex, + slots: Mutex>, +} + +impl Tracker { + pub fn sites(&self) -> Sites { + self.sites.lock().map(|g| g.clone()).unwrap_or_default() + } + + pub fn set_sites(&self, sites: Sites) { + if let Ok(mut g) = self.sites.lock() { + *g = sites; + } + } + + /// 该问这一家了吗。该问就占上,问完必须 [`Tracker::settle`] + pub fn claim(&self, provider: &str, ident: &str, why: Why, now_ms: u64) -> bool { + let Ok(mut g) = self.slots.lock() else { + return false; + }; + let slot = g.entry(provider.to_string()).or_default(); + if slot.ident != ident { + *slot = Slot { + ident: ident.to_string(), + ..Default::default() + }; + } + let due = slot.due(why, now_ms); + if due { + slot.running = true; + } + 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); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + /// 收到回答的时刻:真实示例里 5 小时窗口重置前约 1 小时 + const NOW: u64 = 1_790_030_000_000; + + fn quota(body: &str) -> Quota { + match parse(200, body.as_bytes(), NOW) { + Answer::Quota(q) => q, + other => panic!("{other:?}"), + } + } + + fn window<'a>(q: &'a Quota, name: &str) -> &'a Window { + q.windows + .iter() + .find(|w| w.window == name) + .unwrap_or_else(|| panic!("no {name} in {q:?}")) + } + + #[test] + fn the_hosts_of_both_sites_are_glm_whatever_the_path() { + let s = Sites::default(); + assert_eq!( + s.quota_url("https://api.z.ai/api/anthropic").as_deref(), + Some("https://api.z.ai/api/monitor/usage/quota/limit") + ); + assert_eq!( + s.quota_url("https://open.bigmodel.cn/api/coding/paas/v4") + .as_deref(), + Some("https://open.bigmodel.cn/api/monitor/usage/quota/limit") + ); + // 手写的地址大小写不一、末尾带斜杠,一样认 + assert!(s.quota_url("https://API.Z.AI/api/anthropic/").is_some()); + assert!(s.quota_url("https://api.anthropic.com").is_none()); + // 只是名字里含着的不算 + assert!( + s.quota_url("https://api.z.ai.example.com/api/anthropic") + .is_none() + ); + assert!(s.quota_url("https://bigmodel.cn.example.com").is_none()); + assert!(s.quota_url("不是地址").is_none()); + } + + #[test] + fn a_test_site_is_told_apart_by_its_port() { + let s = Sites { + zai: "http://127.0.0.1:4000".into(), + bigmodel: "http://127.0.0.1:4001".into(), + }; + assert_eq!( + s.quota_url("http://127.0.0.1:4001/api/anthropic") + .as_deref(), + Some("http://127.0.0.1:4001/api/monitor/usage/quota/limit") + ); + assert!(s.quota_url("http://127.0.0.1:4002/api/anthropic").is_none()); + } + + /// 积分制 Lite 套餐的真实回答 + #[test] + fn the_credit_plan_example_reads_as_two_windows_with_credits() { + let q = quota( + r#"{"code":200,"msg":"Operation successful","data":{"limits":[{"type":"CREDIT_LIMIT","unit":3,"number":5,"usage":2000,"currentValue":23,"remaining":1976,"percentage":1,"nextResetTime":1790033645897},{"type":"CREDIT_LIMIT","unit":6,"number":1,"usage":10000,"currentValue":268,"remaining":9731,"percentage":2,"nextResetTime":1790292019984}],"level":"lite"},"success":true}"#, + ); + assert_eq!(q.windows.len(), 2); + let five = window(&q, "5h"); + assert_eq!(five.used_percent, 1.0); + assert_eq!(five.resets_at_ms, Some(1_790_033_645_897)); + assert_eq!(five.status, None); + // **剩余是上游说的剩余**,不是 2000 − 23 + assert_eq!( + five.credits, + Some(tw_api::QuotaCredits { + total: 2000.0, + used: 23.0, + remaining: 1976.0 + }) + ); + let week = window(&q, "weekly"); + assert_eq!(week.used_percent, 2.0); + assert_eq!(week.resets_at_ms, Some(1_790_292_019_984)); + assert_eq!(week.credits.map(|c| c.remaining), Some(9731.0)); + } + + /// 老套餐 V1:一个 5 小时的 token 窗口,加每月的 MCP 调用次数。 + /// 每月窗口放在前面,它先重置 —— **不按重置先后猜** + #[test] + fn an_old_v1_plan_has_five_hours_and_monthly_calls() { + let q = quota(&format!( + r#"{{"code":200,"success":true,"data":{{"level":"pro","limits":[ + {{"type":"TIME_LIMIT","unit":5,"number":1,"usage":1000,"currentValue":40,"remaining":960,"percentage":4,"nextResetTime":{m}}}, + {{"type":"TOKENS_LIMIT","unit":3,"number":5,"percentage":37,"nextResetTime":{h}}} + ]}}}}"#, + m = NOW + 10 * MINUTE_MS, + h = NOW + 3 * HOUR_MS, + )); + assert_eq!(q.windows.len(), 2); + assert_eq!(window(&q, "5h").used_percent, 37.0); + assert_eq!(window(&q, "5h").credits, None); + let month = window(&q, "monthly"); + assert_eq!(month.used_percent, 4.0); + assert_eq!(month.resets_at_ms, Some(NOW + 10 * MINUTE_MS)); + } + + /// 老套餐 V2:多了每周的 token 窗口,`number` 是 7 + #[test] + fn an_old_v2_plan_adds_a_weekly_window() { + let q = quota(&format!( + r#"{{"code":0,"data":{{"limits":[ + {{"type":"TOKENS_LIMIT","unit":6,"number":7,"percentage":81,"nextResetTime":{w}}}, + {{"type":"TOKENS_LIMIT","unit":3,"number":5,"percentage":100,"nextResetTime":{h}}}, + {{"type":"TIME_LIMIT","unit":5,"number":1,"percentage":0}} + ]}}}}"#, + w = NOW + 2 * HOUR_MS, + h = NOW + 4 * HOUR_MS, + )); + let names: Vec<_> = q.windows.iter().map(|w| w.window.as_str()).collect(); + assert_eq!(names, ["weekly", "5h", "monthly"]); + assert_eq!(window(&q, "weekly").used_percent, 81.0); + // 用满就是被拒 + assert_eq!(window(&q, "5h").status.as_deref(), Some("rejected")); + assert_eq!(window(&q, "monthly").resets_at_ms, None); + } + + #[test] + fn a_key_without_a_plan_has_no_quota_data() { + assert_eq!( + parse( + 200, + r#"{"code":500,"msg":"当前用户不存在coding plan","success":false}"#.as_bytes(), + NOW + ), + Answer::NoPlan + ); + // 英文的同一句也一样:**不看 msg** + assert_eq!( + parse( + 200, + br#"{"code":500,"msg":"No coding plan for the current user","success":false}"#, + NOW + ), + Answer::NoPlan + ); + // 成功了却一个窗口都没有,也是没有额度数据,不是「用了 0%」 + assert_eq!( + parse( + 200, + br#"{"code":200,"success":true,"data":{"limits":[]}}"#, + NOW + ), + Answer::NoPlan + ); + } + + #[test] + fn an_auth_failure_is_told_apart_from_other_failures() { + for body in [ + r#"{"code":401,"msg":"令牌已过期或验证不正确","success":false}"#, + r#"{"code":1000,"msg":"Authentication failed","success":false}"#, + r#"{"code":"1001","msg":"Header中未收到Authorization参数","success":false}"#, + // 业务码对了、`success` 没给也算失败 + r#"{"code":1000,"msg":"x"}"#, + ] { + assert_eq!(parse(200, body.as_bytes(), NOW), Answer::Rejected, "{body}"); + } + assert_eq!(parse(401, b"", NOW), Answer::Rejected); + assert_eq!(parse(502, b"", NOW), Answer::Failed); + assert_eq!(parse(200, b"", NOW), Answer::Failed); + assert_eq!( + parse(200, br#"{"code":1234,"success":false}"#, NOW), + Answer::Failed + ); + // `success: false` 自己就是失败,哪怕业务码是 200 + assert_eq!( + parse(200, br#"{"code":200,"success":false}"#, NOW), + Answer::Failed + ); + } + + #[test] + fn a_missing_reset_time_is_none_not_made_up() { + // 用量为 0 的 5 小时窗口可能没有 `nextResetTime` + let q = quota( + r#"{"success":true,"data":{"limits":[{"type":"TOKENS_LIMIT","unit":3,"number":5,"percentage":0}]}}"#, + ); + assert_eq!(window(&q, "5h").resets_at_ms, None); + assert_eq!(window(&q, "5h").used_percent, 0.0); + } + + #[test] + fn a_reset_time_that_cannot_be_right_is_dropped() { + let q = quota(&format!( + r#"{{"success":true,"data":{{"limits":[ + {{"type":"TOKENS_LIMIT","unit":3,"number":5,"percentage":20,"nextResetTime":{late}}}, + {{"type":"TOKENS_LIMIT","unit":6,"number":1,"percentage":20,"nextResetTime":{past}}}, + {{"type":"TIME_LIMIT","unit":5,"number":1,"percentage":20,"nextResetTime":"明天"}} + ]}}}}"#, + // 5 小时窗口却在 6 小时后重置 + late = NOW + 6 * HOUR_MS, + // 已经过去了 + past = NOW - MINUTE_MS, + )); + assert!(q.windows.iter().all(|w| w.resets_at_ms.is_none()), "{q:?}"); + // 差一点的钟不该让一个真实的时刻被丢掉 + let q = quota(&format!( + r#"{{"success":true,"data":{{"limits":[{{"type":"TOKENS_LIMIT","unit":3,"number":5,"percentage":20,"nextResetTime":{t}}}]}}}}"#, + t = NOW + 5 * HOUR_MS + 30_000, + )); + assert_eq!( + window(&q, "5h").resets_at_ms, + Some(NOW + 5 * HOUR_MS + 30_000) + ); + } + + #[test] + fn a_window_it_cannot_name_is_left_out() { + let q = quota( + r#"{"success":true,"data":{"limits":[ + {"type":"TOKENS_LIMIT","percentage":50}, + {"type":"TOKENS_LIMIT","unit":4,"number":1,"percentage":50}, + {"type":"TIME_LIMIT","unit":3,"number":5,"percentage":50}, + {"type":"TOKENS_LIMIT","unit":3,"number":5,"percentage":150} + ]}}"#, + ); + assert_eq!(q.windows.len(), 1, "{q:?}"); + assert_eq!(window(&q, "5h").used_percent, 100.0, "超过 100 的按 100"); + } + + #[test] + fn credits_used_up_are_rejected_even_below_a_hundred_percent() { + let q = quota( + r#"{"success":true,"data":{"limits":[{"type":"CREDIT_LIMIT","unit":6,"number":1,"usage":10000,"currentValue":10000,"remaining":0,"percentage":99}]}}"#, + ); + assert_eq!(window(&q, "weekly").status.as_deref(), Some("rejected")); + } + + // ------------------------------------------------------------ 429 + + #[test] + fn a_five_hour_limit_on_the_anthropic_endpoint() { + let e = exhausted( + br#"{"type":"error","error":{"type":"1308","message":"Usage limit reached for 5 hour. Your limit will reset at 2025-10-03 08:23:14"}}"#, + ) + .unwrap(); + assert_eq!(e.code, 1308); + assert_eq!(e.window, Some("5h")); + // 北京时间 08:23:14 是 UTC 00:23:14 + let utc = chrono::DateTime::parse_from_rfc3339("2025-10-03T00:23:14Z").unwrap(); + assert_eq!(e.resets_at_ms, Some(utc.timestamp_millis() as u64)); + } + + #[test] + fn a_weekly_limit_on_the_openai_endpoint_in_chinese() { + let e = exhausted( + r#"{"error":{"code":"1310","message":"已达到每周使用上限。您的限额将在 2025-10-06 10:00:00 重置。"}}"# + .as_bytes(), + ) + .unwrap(); + assert_eq!(e.code, 1310); + assert_eq!(e.window, Some("weekly")); + let utc = chrono::DateTime::parse_from_rfc3339("2025-10-06T02:00:00Z").unwrap(); + assert_eq!(e.resets_at_ms, Some(utc.timestamp_millis() as u64)); + } + + #[test] + fn the_chinese_five_hour_message_and_a_numeric_code() { + let e = exhausted( + r#"{"error":{"code":1308,"message":"已达到 5 小时的使用上限。您的限额将在 2025-10-03 08:23:14 重置。"}}"# + .as_bytes(), + ) + .unwrap(); + assert_eq!(e.window, Some("5h")); + assert!(e.resets_at_ms.is_some()); + } + + #[test] + fn monthly_and_team_limits_count_too() { + let e = exhausted( + br#"{"type":"error","error":{"type":"1310","message":"Monthly limit exhausted."}}"#, + ) + .unwrap(); + assert_eq!(e.window, Some("monthly")); + assert_eq!(e.resets_at_ms, None, "消息里没有时刻就没有"); + for code in 1316..=1321 { + let body = + format!(r#"{{"error":{{"code":"{code}","message":"Team quota exceeded"}}}}"#); + let e = exhausted(body.as_bytes()).unwrap_or_else(|| panic!("{code}")); + assert_eq!(e.window, None, "{code}:说不清是哪个窗口"); + } + } + + #[test] + fn other_429_codes_are_not_a_used_up_quota() { + // 套餐到期、公平使用限制、临时限流:**重置之后也好不了,或者一会儿就好** + for code in ["1309", "1313", "1302", "1305", "1311"] { + let body = format!( + r#"{{"type":"error","error":{{"type":"{code}","message":"Usage limit reached for 5 hour."}}}}"# + ); + assert_eq!(exhausted(body.as_bytes()), None, "{code}"); + } + // 业务码之外的 429:普通的限流 + assert_eq!( + exhausted( + br#"{"type":"error","error":{"type":"rate_limit_error","message":"slow down"}}"# + ), + None + ); + assert_eq!(exhausted(b"Too Many Requests"), None); + } + + #[test] + fn an_unnamed_window_is_the_tightest_known_candidate() { + let known = Quota { + windows: vec![w("5h", 100.0), w("weekly", 40.0), w("monthly", 90.0)], + }; + // 1310 说的是每周或每月,不会是 5 小时 + assert_eq!(guess_window(1310, &known), "monthly"); + assert_eq!(guess_window(1318, &known), "5h"); + assert_eq!(guess_window(1310, &Quota::default()), "weekly"); + } + + fn w(name: &str, used: f64) -> Window { + Window { + window: name.into(), + used_percent: used, + resets_at_ms: None, + status: None, + credits: None, + } + } + + // ------------------------------------------------------------ 节奏 + + #[test] + fn demands_within_a_minute_are_asked_once() { + let t = Tracker::default(); + 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); + assert!(!t.claim("glm", "k", Why::Demand, NOW + 30_000)); + assert!(t.claim("glm", "k", Why::Demand, NOW + 100 + MINUTE_MS)); + } + + #[test] + 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); + 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()), + NOW + 2 * MINUTE_MS, + ); + assert!(!t.claim("glm", "k", Why::Traffic, NOW + 6 * MINUTE_MS)); + assert!(t.claim("glm", "k", Why::Traffic, NOW + 7 * MINUTE_MS)); + } + + #[test] + fn failures_back_off_30_60_120_300_seconds() { + let t = Tracker::default(); + 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); + // 界面每分钟来要一次也不提前 + 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); + assert!(t.claim("glm", "k", Why::Demand, 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)); + } + + #[test] + 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); + 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); + assert!(t.claim("glm", "k2", Why::Demand, NOW + 90 * 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); + // 新 key 还在问,旧的结论没把它记成「没有套餐」 + 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 293977b7..f032d52d 100644 --- a/crates/tw-gateway/src/lib.rs +++ b/crates/tw-gateway/src/lib.rs @@ -15,6 +15,7 @@ pub mod ending; pub mod error; pub mod fixture; pub mod forward; +pub mod glm; pub mod guard; pub mod health; pub mod hint; diff --git a/crates/tw-gateway/src/quota.rs b/crates/tw-gateway/src/quota.rs index a9b61802..e90c4501 100644 --- a/crates/tw-gateway/src/quota.rs +++ b/crates/tw-gateway/src/quota.rs @@ -12,6 +12,9 @@ //! 例外是他在界面上手动点「立即刷新」—— 那是他明确的意图,而且他知道 //! 代价。 //! +//! GLM Coding Plan 是另一个例外:它的响应头里根本没有额度,只能去问账号的额度 +//! 接口。那个接口不占用户的额度,问的节奏和退避见 [`crate::glm`]。 +//! //! 一条贯穿这里的纪律:**分清真实信号和本地猜测。**上游给的数字可以 //! 直接显示,我们推断的要标成推断。理由和三态成本完全一样: //! **一个编出来的精确数字,比一个诚实的「不知道」更有害。** @@ -19,10 +22,10 @@ use http::HeaderMap; use serde::{Deserialize, Serialize}; -/// 一个额度窗口的状态。**每个字段都直接来自响应头,没有一个是推算的。** +/// 一个额度窗口的状态。**每个字段都直接来自上游(响应头或额度接口),没有一个是推算的。** #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct Window { - /// 哪个窗口:`5h` / `7d`(Anthropic)/ `weekly`(Codex) + /// 哪个窗口:`5h` / `7d`(Anthropic)/ `weekly`(Codex、GLM)/ `monthly`(GLM) pub window: String, /// 用了百分之多少。0–100 pub used_percent: f64, @@ -36,6 +39,21 @@ pub struct Window { /// `allowed` / `allowed_warning` / `rejected` #[serde(default, skip_serializing_if = "Option::is_none")] pub status: Option, + /// 积分制套餐的积分(GLM Coding Plan) + #[serde(default, skip_serializing_if = "Option::is_none")] + pub credits: Option, +} + +impl From<&Window> for tw_api::QuotaWindow { + fn from(w: &Window) -> Self { + Self { + window: w.window.clone(), + used_percent: w.used_percent, + resets_at_ms: w.resets_at_ms, + status: w.status.clone(), + credits: w.credits, + } + } } impl Window { @@ -105,6 +123,7 @@ pub fn from_headers(h: &HeaderMap, now_ms: u64) -> Quota { .and_then(|v| parse_reset(v, now_ms)), status: get(&format!("anthropic-ratelimit-unified-{key}-status")) .map(|s| s.to_string()), + credits: None, }); } // 7 天窗口的越线标记是个独立的头 @@ -146,6 +165,7 @@ pub fn from_headers(h: &HeaderMap, now_ms: u64) -> Quota { .map(|secs| after(now_ms, secs)), // Codex 不报状态。用满就是被拒 —— 这是上游数字的直接结论,不是推断 status: (used_percent >= 100.0).then(|| "rejected".to_string()), + credits: None, }); } Quota { windows } diff --git a/crates/tw-gateway/src/server/pipeline/hop.rs b/crates/tw-gateway/src/server/pipeline/hop.rs index 392c65c8..761c3f69 100644 --- a/crates/tw-gateway/src/server/pipeline/hop.rs +++ b/crates/tw-gateway/src/server/pipeline/hop.rs @@ -206,6 +206,11 @@ pub(super) async fn try_upstreams<'a>( "Upstream `{upstream}` answered {status}." )) }); + // GLM Coding Plan 的额度用完不在响应头里,在 429 的 body 里 + if r.status() == 429 { + state.note_glm_429(id, provider, r).await; + } + state.glm_traffic(provider); continue; } Ok(r) => { diff --git a/crates/tw-gateway/src/server/pipeline/relay.rs b/crates/tw-gateway/src/server/pipeline/relay.rs index 7169bdd3..a28fa18f 100644 --- a/crates/tw-gateway/src/server/pipeline/relay.rs +++ b/crates/tw-gateway/src/server/pipeline/relay.rs @@ -59,6 +59,8 @@ pub(super) fn respond( // 订阅额度。**零成本** —— 这些头本来就在响应里,读一下 // 就有了。按量付费的账号没有它们,那时什么都不发。 state.note_quota(id, &provider.name, upstream.headers()); + // GLM Coding Plan 的额度不在响应头里:有请求的时候隔一阵去问一次 + state.glm_traffic(provider); // 上游收不收我们的凭据、经过的代理通不通:**都是状态变化,各只报一次** state.note_auth(&provider.name, status.as_u16()); state.note_proxy_ok(&provider.proxy); diff --git a/crates/tw-gateway/src/state.rs b/crates/tw-gateway/src/state.rs index 07ccaa66..8202749a 100644 --- a/crates/tw-gateway/src/state.rs +++ b/crates/tw-gateway/src/state.rs @@ -11,6 +11,7 @@ use crate::outbound::{base_client_builder, client_for_provider, proxy_shape}; use tw_types::msg; mod credentials; +mod glm; pub use credentials::credential_failed; mod upstream; @@ -151,6 +152,8 @@ pub struct AppState { /// 五分钟前的百分比,价值几乎为零,而它会让「重启之后显示的是旧 /// 数字」变成一个要解释的问题。下一个请求回来就有新的了。 quotas: Arc>>, + /// 每个 GLM Coding Plan 上游问额度的节奏(见 [`crate::glm`])。**跨重载存活** + glm: Arc, /// 监听地址变了。**这是「温」那一级**(三级热重载) —— /// 换端口不能只换配置:监听器是启动时建的,不重建的话新端口上什么 /// 都没有,而旧端口还在服务。那种「改了没反应」比报错难查得多。 @@ -232,6 +235,7 @@ impl AppState { models, body_sink: Arc::new(std::sync::Mutex::new(None)), quotas: Arc::new(std::sync::Mutex::new(Default::default())), + glm: Default::default(), relisten: Arc::new(tokio::sync::Notify::new()), listening: Arc::new(std::sync::Mutex::new(Default::default())), oauth: Arc::new(crate::oauth::Cache::new()), diff --git a/crates/tw-gateway/src/state/glm.rs b/crates/tw-gateway/src/state/glm.rs new file mode 100644 index 00000000..dd480fd6 --- /dev/null +++ b/crates/tw-gateway/src/state/glm.rs @@ -0,0 +1,283 @@ +//! GLM Coding Plan 的额度:什么时候去问、问来的和 429 说的怎么记(见 [`crate::glm`])。 + +use std::time::Duration; + +use super::AppState; +use crate::glm::{self, Answer, Why}; +use crate::server::now_ms; + +/// 问一次额度接口最多等多久 +const ASK_TIMEOUT: Duration = Duration::from_secs(15); + +/// 429 的 body 最多读多少、等多久。**这一跳已经失败了**,读它只为认出额度用完, +/// 不能因此拖住下一家 +const BODY_LIMIT: usize = 16 * 1024; +const BODY_TIMEOUT: Duration = Duration::from_secs(2); + +impl AppState { + /// 额度接口换成别的地址。**给测试用**:平时就是 Z.ai 和 BigModel 自己的 + pub fn set_glm_sites(&self, sites: glm::Sites) { + self.glm.set_sites(sites); + } + + /// 这一家是 GLM Coding Plan 的话,额度接口的地址和凭据的指纹。 + /// + /// 没有 key 的不算:额度跟着 key 走,没有 key 就没有可问的 + fn glm_target(&self, p: &tw_config::Provider) -> Option<(String, String)> { + let key = p.key.as_ref()?; + let url = self.glm.sites().quota_url(&p.base_url)?; + // 指纹,不是 key 本身:它只用来认出「换了 key」 + let ident = blake3::hash(key.raw().as_bytes()).to_hex().to_string(); + Some((url, ident)) + } + + /// 界面来要额度:该问的 GLM 上游都问一遍,**最多等 `wait`**。 + /// + /// 等不到的照样在后台问完,结果进 `/quota` 和 `QuotaSeen`。60 秒内问过的、 + /// 退避中的、没有套餐的,这一次都不问 + pub async fn refresh_glm_quotas(&self, wait: Duration) { + let cfg = self.config(); + let asks: Vec<_> = cfg + .providers + .iter() + .filter_map(|p| self.spawn_glm_ask(p, Why::Demand)) + .collect(); + if asks.is_empty() { + return; + } + let _ = tokio::time::timeout(wait, futures::future::join_all(asks)).await; + } + + /// 有请求去了这一家:到时候了就在后台问一次额度。**不挡转发** + pub(crate) fn glm_traffic(&self, p: &tw_config::Provider) { + let _ = self.spawn_glm_ask(p, Why::Traffic); + } + + /// 该问就起一个任务去问。**任务自己跑完**:等的一方放弃了,节奏也照样记上 + fn spawn_glm_ask( + &self, + p: &tw_config::Provider, + why: Why, + ) -> Option> { + let (url, ident) = self.glm_target(p)?; + if !self.glm.claim(&p.name, &ident, why, now_ms()) { + return None; + } + let state = self.clone(); + let p = p.clone(); + Some(tokio::spawn(async move { + let answer = state.ask_glm(&p, &url).await; + match &answer { + Answer::Quota(q) => state.record_quota(state.bus.next_id(), &p.name, q.clone()), + // 没有套餐就没有额度数据:之前记着的也不再作数 + Answer::NoPlan => state.forget_quota(&p.name), + Answer::Rejected => { + tracing::warn!(provider = %p.name, "the GLM quota endpoint rejected the key") + } + Answer::Failed => { + tracing::debug!(provider = %p.name, "the GLM quota could not be read") + } + } + state.glm.settle(&p.name, &ident, &answer, now_ms()); + })) + } + + /// 问一次额度接口。**key 原样放在 `Authorization` 里,不加 `Bearer`**:那个接口就是 + /// 这么认的 + async fn ask_glm(&self, p: &tw_config::Provider, url: &str) -> Answer { + let Some(Ok(key)) = p.key.as_ref().map(|k| k.resolve()) else { + return Answer::Failed; + }; + // 这一家自己的 client:走它该走的代理 + let http = self.client_for(&p.name); + let sent = http + .get(url) + .header(reqwest::header::AUTHORIZATION, key) + .header(reqwest::header::ACCEPT, "application/json") + .timeout(ASK_TIMEOUT) + .send() + .await; + let Ok(r) = sent else { + return Answer::Failed; + }; + let status = r.status().as_u16(); + match r.bytes().await { + Ok(body) => glm::parse(status, &body, now_ms()), + Err(_) => Answer::Failed, + } + } + + /// 一个 GLM 上游回了 429:body 里说额度用完了的话,**那一刻就记成用完**,报一次 + /// `QuotaExhausted`。 + /// + /// 重置时刻以额度接口问来的为准;还没问到过,才用消息里说的 + pub(crate) async fn note_glm_429( + &self, + id: u64, + p: &tw_config::Provider, + r: reqwest::Response, + ) { + if self.glm_target(p).is_none() { + return; + } + let Some(body) = read_capped(r).await else { + return; + }; + self.note_glm_exhausted(id, &p.name, &body); + } + + /// 读出来的 429 body 说额度用完了的话,记下来。 + pub fn note_glm_exhausted(&self, id: u64, provider: &str, body: &[u8]) { + let Some(hit) = glm::exhausted(body) else { + return; + }; + let now = now_ms(); + let mut quota = self.quotas().remove(provider).unwrap_or_default(); + let window = hit + .window + .map(str::to_string) + .unwrap_or_else(|| glm::guess_window(hit.code, "a)); + let known = quota + .windows + .iter() + .find(|w| w.window == window) + .and_then(|w| w.resets_at_ms) + .filter(|t| *t > now); + let resets_at_ms = known.or(hit.resets_at_ms.filter(|t| *t > now)); + tracing::info!(provider, code = hit.code, %window, "GLM says the quota is used up"); + // 用完就是用满 —— 上游说的,不是推断 + match quota.windows.iter_mut().find(|w| w.window == window) { + Some(w) => { + w.used_percent = 100.0; + w.status = Some("rejected".to_string()); + w.resets_at_ms = resets_at_ms; + } + None => quota.windows.push(crate::quota::Window { + window, + used_percent: 100.0, + resets_at_ms, + status: Some("rejected".to_string()), + credits: None, + }), + } + self.record_quota(id, provider, quota); + } +} + +/// 读 body,最多 [`BODY_LIMIT`] 字节、[`BODY_TIMEOUT`]。读不完就算了 +async fn read_capped(mut r: reqwest::Response) -> Option> { + let read = async { + let mut out = Vec::new(); + while let Some(chunk) = r.chunk().await.ok()? { + out.extend_from_slice(&chunk); + if out.len() > BODY_LIMIT { + return None; + } + } + Some(out) + }; + tokio::time::timeout(BODY_TIMEOUT, read).await.ok()? +} + +#[cfg(test)] +mod tests { + use super::*; + use tw_config::{Client, Config, Provider}; + + fn state() -> AppState { + AppState::new(Config { + version: 1, + clients: vec![Client { + name: "default".into(), + key: "tw-good".into(), + ..Default::default() + }], + providers: vec![Provider { + name: "glm".into(), + base_url: "https://api.z.ai/api/anthropic".into(), + key: Some("sk-test".into()), + ..Default::default() + }], + ..Default::default() + }) + .unwrap() + } + + const FIVE_HOURS: &[u8] = br#"{"type":"error","error":{"type":"1308","message":"Usage limit reached for 5 hour. Your limit will reset at 2099-10-03 08:23:14"}}"#; + + #[tokio::test] + async fn a_used_up_window_is_reported_once_with_the_message_reset_time() { + let s = state(); + let mut rx = s.bus.subscribe(); + s.note_glm_exhausted(7, "glm", FIVE_HOURS); + let q = &s.quotas()["glm"]; + assert_eq!(q.windows.len(), 1); + let w = &q.windows[0]; + assert_eq!(w.window, "5h"); + assert!(w.rejected()); + let utc = chrono::DateTime::parse_from_rfc3339("2099-10-03T00:23:14Z").unwrap(); + assert_eq!(w.resets_at_ms, Some(utc.timestamp_millis() as u64)); + + let mut exhausted = 0; + s.note_glm_exhausted(8, "glm", FIVE_HOURS); + while let Ok(e) = rx.try_recv() { + if let tw_api::Event::QuotaExhausted { window, .. } = e { + assert_eq!(window, "5h"); + exhausted += 1; + } + } + assert_eq!(exhausted, 1, "同一个窗口用完只报一次"); + } + + #[tokio::test] + async fn the_quota_endpoint_reset_time_wins_over_the_message() { + let s = state(); + let at = now_ms() + 3_600_000; + s.record_quota( + 1, + "glm", + crate::quota::Quota { + windows: vec![crate::quota::Window { + window: "5h".into(), + used_percent: 97.0, + resets_at_ms: Some(at), + status: None, + credits: Some(tw_api::QuotaCredits { + total: 2000.0, + used: 1940.0, + remaining: 60.0, + }), + }], + }, + ); + s.note_glm_exhausted(2, "glm", FIVE_HOURS); + let w = &s.quotas()["glm"].windows[0]; + assert_eq!(w.resets_at_ms, Some(at)); + assert_eq!(w.used_percent, 100.0); + assert!(w.rejected()); + } + + #[tokio::test] + async fn a_rate_limit_that_is_not_the_quota_changes_nothing() { + let s = state(); + s.note_glm_exhausted( + 1, + "glm", + br#"{"type":"error","error":{"type":"1313","message":"fair use"}}"#, + ); + assert!(s.quotas().is_empty()); + } + + #[test] + fn only_a_glm_upstream_with_a_key_is_asked() { + let s = state(); + let cfg = s.config(); + assert!(s.glm_target(&cfg.providers[0]).is_some()); + let mut other = cfg.providers[0].clone(); + other.base_url = "https://api.anthropic.com".into(); + assert!(s.glm_target(&other).is_none()); + let mut keyless = cfg.providers[0].clone(); + keyless.key = None; + assert!(s.glm_target(&keyless).is_none()); + } +} diff --git a/crates/tw-gateway/src/state/upstream.rs b/crates/tw-gateway/src/state/upstream.rs index c46a9ab4..b706b907 100644 --- a/crates/tw-gateway/src/state/upstream.rs +++ b/crates/tw-gateway/src/state/upstream.rs @@ -236,16 +236,7 @@ impl AppState { self.bus.emit(tw_api::Event::QuotaSeen { id, provider: provider.to_string(), - windows: quota - .windows - .iter() - .map(|w| tw_api::QuotaWindow { - window: w.window.clone(), - used_percent: w.used_percent, - resets_at_ms: w.resets_at_ms, - status: w.status.clone(), - }) - .collect(), + windows: quota.windows.iter().map(Into::into).collect(), at_ms: now_ms(), }); for w in "a.windows { @@ -275,6 +266,16 @@ impl AppState { } } + /// 这一家没有额度数据了(比如 key 没有开通套餐):记着的额度和「用完」都不再作数 + pub(crate) fn forget_quota(&self, provider: &str) { + if let Ok(mut g) = self.quotas.lock() { + g.remove(provider); + } + if let Ok(mut g) = self.exhausted.lock() { + g.retain(|(p, _)| p != provider); + } + } + /// 每个上游最近一次报的订阅额度。 pub fn quotas(&self) -> std::collections::HashMap { self.quotas.lock().map(|g| g.clone()).unwrap_or_default() diff --git a/crates/tw-gateway/tests/glm_quota.rs b/crates/tw-gateway/tests/glm_quota.rs new file mode 100644 index 00000000..94ead7ca --- /dev/null +++ b/crates/tw-gateway/tests/glm_quota.rs @@ -0,0 +1,222 @@ +//! GLM Coding Plan 的额度:去额度接口问,和从 429 里认出额度用完。 +//! +//! **全程只对本机的假 GLM**:模型接口和额度接口都是假的,key 也是编的。 + +use std::net::SocketAddr; +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +use axum::Router; +use axum::extract::State; +use axum::http::HeaderMap; +use axum::response::IntoResponse; +use axum::routing::{any, get}; +use tw_config::{Client, Config, Provider}; + +/// 假 GLM:模型接口回什么、额度接口回什么,以及额度接口收到的 `Authorization` +#[derive(Default)] +struct Glm { + /// 模型接口的回答:(状态码, body) + model: Mutex<(u16, String)>, + /// 额度接口的 body + quota: Mutex, + /// 额度接口每次收到的 `Authorization` + asked: Mutex>, +} + +async fn model(State(g): State>) -> axum::response::Response { + let (status, body) = g.model.lock().unwrap().clone(); + ( + axum::http::StatusCode::from_u16(status).unwrap(), + [("content-type", "application/json")], + body, + ) + .into_response() +} + +async fn quota(State(g): State>, headers: HeaderMap) -> axum::response::Response { + let auth = headers + .get("authorization") + .and_then(|v| v.to_str().ok()) + .unwrap_or_default() + .to_string(); + g.asked.lock().unwrap().push(auth); + let body = g.quota.lock().unwrap().clone(); + ([("content-type", "application/json")], body).into_response() +} + +async fn start_glm(g: Arc) -> SocketAddr { + let app = Router::new() + .route("/api/monitor/usage/quota/limit", get(quota)) + .fallback(any(model)) + .with_state(g); + let l = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let a = l.local_addr().unwrap(); + tokio::spawn(async move { axum::serve(l, app).await.unwrap() }); + a +} + +/// 一个指向假 GLM 的上游,额度接口也换成假 GLM 的 +async fn state_for(glm: SocketAddr) -> tw_gateway::AppState { + let state = tw_gateway::AppState::new(Config { + clients: vec![Client { + name: "c".into(), + key: "tw-k".into(), + ..Default::default() + }], + providers: vec![Provider { + name: "glm".into(), + base_url: format!("http://{glm}/api/anthropic"), + key: Some("fake-glm-key".into()), + protocol: Some(tw_config::Protocol::Anthropic), + ..Default::default() + }], + ..Default::default() + }) + .unwrap(); + state.set_glm_sites(tw_gateway::glm::Sites { + zai: format!("http://{glm}"), + bigmodel: "https://open.bigmodel.cn".into(), + }); + state +} + +fn credit_plan(now_ms: u64) -> String { + format!( + r#"{{"code":200,"msg":"Operation successful","data":{{"limits":[{{"type":"CREDIT_LIMIT","unit":3,"number":5,"usage":2000,"currentValue":23,"remaining":1976,"percentage":1,"nextResetTime":{five}}},{{"type":"CREDIT_LIMIT","unit":6,"number":1,"usage":10000,"currentValue":268,"remaining":9731,"percentage":2,"nextResetTime":{week}}}],"level":"lite"}},"success":true}}"#, + five = now_ms + 3_600_000, + week = now_ms + 3 * 86_400_000, + ) +} + +fn now_ms() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_millis() as u64 +} + +#[tokio::test] +async fn asking_for_the_quota_reads_the_windows_and_credits() { + let g = Arc::new(Glm::default()); + *g.quota.lock().unwrap() = credit_plan(now_ms()); + let state = state_for(start_glm(g.clone()).await).await; + + state.refresh_glm_quotas(Duration::from_secs(5)).await; + let q = &state.quotas()["glm"]; + let names: Vec<_> = q.windows.iter().map(|w| w.window.as_str()).collect(); + assert_eq!(names, ["5h", "weekly"]); + assert_eq!( + q.windows[0].credits, + Some(tw_api::QuotaCredits { + total: 2000.0, + used: 23.0, + remaining: 1976.0 + }) + ); + // **key 原样放在 Authorization 里,不加 Bearer** + assert_eq!(*g.asked.lock().unwrap(), ["fake-glm-key"]); + + // 一分钟之内再来要:不再问 + state.refresh_glm_quotas(Duration::from_secs(5)).await; + assert_eq!(g.asked.lock().unwrap().len(), 1); +} + +#[tokio::test] +async fn a_key_without_a_plan_has_no_quota_and_is_not_asked_again_soon() { + let g = Arc::new(Glm::default()); + *g.quota.lock().unwrap() = + r#"{"code":500,"msg":"当前用户不存在coding plan","success":false}"#.into(); + let state = state_for(start_glm(g.clone()).await).await; + + state.refresh_glm_quotas(Duration::from_secs(5)).await; + assert!(state.quotas().is_empty(), "没有套餐不是「用了 0%」"); + state.refresh_glm_quotas(Duration::from_secs(5)).await; + assert_eq!(g.asked.lock().unwrap().len(), 1); +} + +#[tokio::test] +async fn a_used_up_quota_in_a_429_is_reported_and_the_quota_is_asked() { + let g = Arc::new(Glm::default()); + *g.model.lock().unwrap() = ( + 429, + r#"{"type":"error","error":{"type":"1308","message":"Usage limit reached for 5 hour. Your limit will reset at 2099-10-03 08:23:14"}}"#.into(), + ); + // 额度接口还没跟上:5 小时窗口只报了 99%,也没说什么时候重置 + *g.quota.lock().unwrap() = r#"{"success":true,"data":{"limits":[{"type":"TOKENS_LIMIT","unit":3,"number":5,"percentage":99}]}}"#.into(); + let state = state_for(start_glm(g.clone()).await).await; + let mut rx = state.bus.subscribe(); + let gw = tw_gateway::serve(state.clone(), ([127, 0, 0, 1], 0).into()) + .await + .unwrap(); + tokio::time::sleep(Duration::from_millis(50)).await; + + let r = reqwest::Client::new() + .post(format!("http://{gw}/v1/messages")) + .header("x-api-key", "tw-k") + .body(r#"{"model":"glm-4.6","messages":[]}"#) + .send() + .await + .unwrap(); + assert_eq!(r.status(), 429, "429 要保住 429"); + + let (window, resets_at) = loop { + match tokio::time::timeout(Duration::from_secs(3), rx.recv()).await { + Ok(Ok(tw_api::Event::QuotaExhausted { + window, + resets_at_ms, + .. + })) => break (window, resets_at_ms), + Ok(Ok(_)) => continue, + other => panic!("没等到额度用完的事件:{other:?}"), + } + }; + assert_eq!(window, "5h"); + // 额度接口没说,才用消息里的北京时间 + let utc = chrono::DateTime::parse_from_rfc3339("2099-10-03T00:23:14Z").unwrap(); + assert_eq!(resets_at, Some(utc.timestamp_millis() as u64)); + + // 有请求去了这一家:顺手问了一次额度 + let asked = async { + while g.asked.lock().unwrap().is_empty() { + tokio::time::sleep(Duration::from_millis(20)).await; + } + }; + tokio::time::timeout(Duration::from_secs(3), asked) + .await + .expect("有请求时应当去问额度"); +} + +#[tokio::test] +async fn a_429_that_is_not_the_quota_leaves_the_quota_alone() { + let g = Arc::new(Glm::default()); + *g.model.lock().unwrap() = ( + 429, + r#"{"type":"error","error":{"type":"1302","message":"High concurrency, slow down"}}"# + .into(), + ); + *g.quota.lock().unwrap() = r#"{"code":500,"success":false}"#.into(); + let state = state_for(start_glm(g.clone()).await).await; + let mut rx = state.bus.subscribe(); + let gw = tw_gateway::serve(state.clone(), ([127, 0, 0, 1], 0).into()) + .await + .unwrap(); + tokio::time::sleep(Duration::from_millis(50)).await; + + let r = reqwest::Client::new() + .post(format!("http://{gw}/v1/messages")) + .header("x-api-key", "tw-k") + .body(r#"{"model":"glm-4.6","messages":[]}"#) + .send() + .await + .unwrap(); + assert_eq!(r.status(), 429); + tokio::time::sleep(Duration::from_millis(300)).await; + while let Ok(e) = rx.try_recv() { + assert!( + !matches!(e, tw_api::Event::QuotaExhausted { .. }), + "临时限流不是额度用完:{e:?}" + ); + } + assert!(state.quotas().is_empty()); +}