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
82 changes: 61 additions & 21 deletions crates/tw-gateway/src/glm.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,8 @@
//! 问的纪律:
//! - **界面来要时问**,60 秒内的多次合成一次;**有请求去这一家时顺手问**,5 分钟最多一次。
//! - 失败了退避(30 秒、60 秒、120 秒、300 秒),不跟着界面的每一次刷新重试。
//! - 这把 key 没有开通套餐:记成「没有额度数据」,隔一小时才再问。
//! - 这把 key 没有开通套餐:记成「没有额度数据」,隔一小时才再问。刚才还有额度的 key
//! 要连着说两次才算(一次临时的 500 和「没有套餐」长得一样)。
//!
//! 额度用完时模型请求回 429,业务码在 body 里([`exhausted`])。**那一刻就记成用完**,
//! 不等下一次问额度。
Expand Down Expand Up @@ -324,6 +325,8 @@ struct Slot {
quiet_until: u64,
/// 正在问。**同一时刻只问一次**
running: bool,
/// 上一次问来的是额度
had_quota: bool,
}

impl Slot {
Expand All @@ -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;
Expand All @@ -357,6 +370,7 @@ impl Slot {
self.quiet_until = now_ms + BACKOFF_MS[i];
}
}
answer
}
}

Expand Down Expand Up @@ -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<Answer> {
let mut g = self.slots.lock().ok()?;
let slot = g.get_mut(provider).filter(|s| s.ident == ident)?;
Some(slot.settle(answer, now_ms))
}
}

Expand Down Expand Up @@ -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));
}
Expand All @@ -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));
Expand All @@ -778,40 +796,62 @@ 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));
}

#[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);
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));
}
}
7 changes: 7 additions & 0 deletions crates/tw-gateway/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Expand Down
21 changes: 14 additions & 7 deletions crates/tw-gateway/src/server/pipeline/hop.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down Expand Up @@ -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::<Vec<_>>()
.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| {
Expand Down
17 changes: 15 additions & 2 deletions crates/tw-gateway/src/server/pipeline/relay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 丢掉响应体是另一种 —— 两种
Expand All @@ -106,17 +107,27 @@ pub(super) fn respond(
let mut ending = ending;
let mut chunks = std::pin::pin!(chunks);
let mut broke: Option<GatewayError> = 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, std::io::Error>(Bytes::from_static(PING));
}
continue;
Expand All @@ -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 是真相,而拿不到它就只能估。
//
Expand All @@ -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, std::io::Error>(Bytes::from(out));
}
if let Some(err) = cut {
Expand Down
3 changes: 3 additions & 0 deletions crates/tw-gateway/src/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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();
Expand Down
15 changes: 13 additions & 2 deletions crates/tw-gateway/src/state/glm.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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」
Expand Down Expand Up @@ -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()),
// 没有套餐就没有额度数据:之前记着的也不再作数
Expand All @@ -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());
}))
}

Expand Down Expand Up @@ -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());
}
}
Loading
Loading