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
7 changes: 5 additions & 2 deletions crates/tw-api/msg-codes.txt
Original file line number Diff line number Diff line change
Expand Up @@ -107,7 +107,6 @@ config.schema_too_new
config.secret.empty_name
config.secret.env_missing
config.secret.unterminated
config.slow_start_too_short
config.store.conflict
config.store.missing
config.store.read_failed
Expand Down Expand Up @@ -206,6 +205,7 @@ control.replay_bedrock
control.request_body_gone
control.request_body_truncated
control.request_not_found
control.request_not_running
control.request_rejected
control.reset_card_result_unreadable
control.reset_cards_unreadable
Expand All @@ -224,6 +224,7 @@ control.rule.no_such_probe_class
control.rule.no_such_upstream
control.rule.unknown_dialect
control.session_not_found
control.session_not_running
control.sheet_in_use
control.sheet_not_found
control.shutdown
Expand Down Expand Up @@ -372,6 +373,7 @@ gw.probe.connect
gw.probe.key_rejected
gw.probe.request_failed
gw.probe.timeout
gw.request.aborted
gw.request.body_declared_too_large
gw.request.body_over_limit
gw.request.body_unreadable
Expand All @@ -383,7 +385,6 @@ gw.route.protocol_mismatch
gw.route.rule_failed
gw.route.selected_upstream_missing
gw.route.upstream_missing
gw.slow_start
gw.toolcall.connection_cut
gw.toolcall.response_cut
gw.toolcall.response_withheld
Expand All @@ -393,6 +394,8 @@ gw.upstream.bedrock_refused
gw.upstream.bedrock_refused_unnamed
gw.upstream.eventstream_broken
gw.upstream.forward_failed
gw.upstream.idle_timeout
gw.upstream.idle_timeout_mid_stream
gw.upstream.rate_limited
gw.upstream.sign_failed
gw.upstream.status
Expand Down
5 changes: 5 additions & 0 deletions crates/tw-api/src/ep.rs
Original file line number Diff line number Diff line change
Expand Up @@ -60,13 +60,18 @@ endpoints! {
/// 延迟与速度,各带样本数和别家的参照。不给时间窗是最近 7 天
UpstreamHealth: GET "/upstreams/health", api::Window => api::UpstreamHealth;
RequestDetail: GET "/request/{id}" [id], () => api::RequestDetail;
/// 中止一个在跑的请求:和上游的连接断开,客户端收到错误,请求记成手动中止。已经
/// 结束了的是 404
AbortRequest: POST "/request/{id}/abort" [id], () => api::Aborted;
/// 把一条记录变成回放用例(YAML)
Fixture: GET "/request/{id}/fixture" [id], () => String, text;
Sessions: GET "/sessions", api::ListQuery => Vec<api::SessionView>;
SessionDetail: GET "/sessions/{id}" [id], () => api::SessionDetail;
/// 一次会话读成一段对话:每一轮新说的话、回答、工具调用和结果(已脱敏)。可以只要
/// 从某一轮起的那些
SessionTranscript: GET "/sessions/{id}/transcript" [id], api::TranscriptQuery => api::Transcript;
/// 中止一次会话里所有在跑的请求。一个都没有是 404
AbortSession: POST "/sessions/{id}/abort" [id], () => api::Aborted;

// ─────────────────────────────────────────────── 测速、回放、试路由
SpeedQuote: POST "/speed/quote", api::SpeedRunRequest => api::SpeedQuote;
Expand Down
82 changes: 65 additions & 17 deletions crates/tw-api/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -224,6 +224,11 @@ slug_enum! {
Request = "request",
RateLimited = "rate_limited",
Denied = "denied",
/// 在界面上手动中止的(`POST /request/{id}/abort`、`POST /sessions/{id}/abort`)。
/// **不是上游的错**:那一家不停用、不算失败;和客户端自己走掉(`RequestCancelled`)
/// 也不是一回事。`message` 是 [`ABORTED`] 那个码。还没开始回答的,客户端收到 499;
/// 回答到一半的,按它的格式以一条错误收尾
Aborted = "aborted",
/// 网关自己的代码崩掉了
Internal = "internal",
}
Expand Down Expand Up @@ -814,7 +819,24 @@ pub const MSG_CODES: &str = include_str!("../msg-codes.txt");
/// `config.manual_model_blank`、`_padded`、`_wildcard`、`_duplicate`、`_too_long`(400),
/// 整份配置的校验也查这几条。`source` 是 `discovered` 时 `models` 里也有手动添加的那些。
/// 照 41 写的界面分不出哪些模型是手动添加的。
pub const CONTROL_API_VERSION: u32 = 42;
///
/// **43 起上游不出声有了上限,请求能手动中止**:[`FailoverView`] 的 `stream_start_wait_secs`
/// 和 `next_on_slow_start` 删了,换成 `idle_timeout_secs`(无响应超时,默认 300 秒,30 到
/// 3600):从请求发出去起算,每来一段内容重新计时,心跳不算。还没有内容交给客户端时,这一家
/// 记一次失败、换下一家,没有下一家了回 504;已经收到一部分的,回答按客户端的格式以错误收尾。
/// 流式的回答压着第一段内容最多 15 秒,过了就先交出 `200` 和流的响应头、隔一阵发一行 SSE
/// 注释(`: keep-alive`,Gemini 的客户端不发),故障转移照旧,候选用完了在流里报错,不再是 504。
/// 配置里写 `stream_start_wait_secs`、`next_on_slow_start` 加载不了(不认识的字段,一键修复删掉
/// 它们),消息码 `config.slow_start_too_short`、`gw.slow_start` 跟着删。尝试链的结果
/// ([`AttemptOutcome`])删了 `slow_start`,多了 `idle_timeout`(说等了多少秒的
/// `gw.upstream.idle_timeout`)和 `aborted`。新端点 `POST /request/{id}/abort` 和
/// `POST /sessions/{id}/abort`(→ [`Aborted`])叫停一个在跑的请求、一次会话里所有在跑的请求:
/// 和上游的连接立刻断开,客户端收到它自己格式的错误(还没开始回答的是 499),请求照常报
/// [`Event::RequestFailed`],`source` 是新的 [`FailureSource::Aborted`]、`message` 是
/// [`ABORTED`],上游不停用。已经结束了的回 404(`control.request_not_running`、
/// `control.session_not_running`)。流中途停了的另有 `gw.upstream.idle_timeout_mid_stream`。
/// 照 42 写的界面读不到 `idle_timeout_secs`,不认 `aborted` 这个来源。
pub const CONTROL_API_VERSION: u32 = 43;

#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "ts", derive(ts_rs::TS))]
Expand Down Expand Up @@ -1515,11 +1537,17 @@ slug_enum! {
/// (`status` 是它回的那个)。尝试链到此为止,这一行记成网关自己答的
/// ([`HistoryRow::local`]),费用 0
Estimated = "estimated",
/// 流式回答等了 `failover.stream_start_wait_secs` 还没有内容,开着
/// `failover.next_on_slow_start`,放弃这一家、换下一家(连接断开,上游不再生成)。
/// 响应头到了的有 `status`,没到的没有。**这一家不停用、不算失败**。上游可能已经按
/// 输入收了钱:知道多少的在 `usage` 里
SlowStart = "slow_start",
/// 从请求发出去起 `failover.idle_timeout_secs` 秒没有内容(无响应超时),客户端也还
/// 什么都没收到:放弃这一家(连接断开,上游不再生成),换下一家,没有下一家了就回
/// 超时错误。响应头到了的有 `status`,没到的没有。`error` 是 `gw.upstream.idle_timeout`,
/// 说等了多少秒。**这一家记一次失败**(和 5xx 一样算进停用的账)。上游可能已经按输入
/// 收了钱:知道多少的在 `usage` 里
IdleTimeout = "idle_timeout",
/// 这一跳在等上游时被手动中止(`POST /request/{id}/abort`、`POST /sessions/{id}/abort`):
/// 连接断开,上游不再生成,尝试链到此为止。**这一家不停用、不算失败**。已经接下、
/// 在交回答的那一跳不改成它(还是 `served`),请求本身记成手动中止(见
/// [`FailureSource::Aborted`])
Aborted = "aborted",
}
}

Expand All @@ -1545,11 +1573,12 @@ pub struct AttemptView {
/// 上游返回的状态码。`error` 时没有
#[serde(default, skip_serializing_if = "Option::is_none")]
pub status: Option<u16>,
/// `error` 时的说明。和这一跳报给客户端的那条错误是同一句。`slow_start` 时说等了多久
/// `error` 时的说明。和这一跳报给客户端的那条错误是同一句。`idle_timeout` 时说等了多久,
/// `aborted` 时是 `gw.request.aborted`
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<Msg>,
pub ms: u64,
/// 放弃了的这一跳(`slow_start`)上游可能已经收了钱的输入(见 [`AttemptUsage`])。估不
/// 放弃了的这一跳(`idle_timeout`)上游可能已经收了钱的输入(见 [`AttemptUsage`])。估不
/// 出来的(请求解不开)没有。别的结果都没有:接下请求的那一跳的用量在结局里
#[serde(default, skip_serializing_if = "Option::is_none")]
pub usage: Option<AttemptUsage>,
Expand All @@ -1575,12 +1604,19 @@ pub const FAILED_AFTER_SENDING: &[&str] = &[
"gw.upstream.stream_opening_error",
];

/// 手动中止的请求报的码([`FailureSource::Aborted`] 的 `message`,尝试链上 `aborted` 那一跳的
/// `error`)。记录里只存码,上游体检按它把手动中止的和客户端取消的一样数,不算上游失败
pub const ABORTED: &str = "gw.request.aborted";

impl AttemptView {
/// 这一跳发到了上游:上游回了话(`served`、`status`、`slow_start`,`estimated` 里带着
/// 这一跳发到了上游:上游回了话(`served`、`status`、`idle_timeout`、`aborted`,`estimated` 里带着
/// 状态码的),或者发出去之后才失败([`FAILED_AFTER_SENDING`])。
pub fn sent(&self) -> bool {
match self.outcome {
AttemptOutcome::Served | AttemptOutcome::Status | AttemptOutcome::SlowStart => true,
AttemptOutcome::Served
| AttemptOutcome::Status
| AttemptOutcome::IdleTimeout
| AttemptOutcome::Aborted => true,
AttemptOutcome::Estimated => self.status.is_some(),
AttemptOutcome::Error => {
self.skipped.is_none()
Expand Down Expand Up @@ -1609,7 +1645,7 @@ impl RoutingView {
}
}

/// 放弃了的一跳([`AttemptOutcome::SlowStart`])上游可能已经收了钱的输入。
/// 放弃了的一跳([`AttemptOutcome::IdleTimeout`])上游可能已经收了钱的输入。
///
/// 上游在流开头报了的(Anthropic 的 `message_start`)是它报的数;没报的只有 `input`,是网关
/// 估的(`estimated`,和 [`Event::RequestStarted`] 的 `input_estimate` 同一个数)。**输出不知道**:
Expand Down Expand Up @@ -1841,6 +1877,17 @@ pub struct InFlightRequest {
pub events: Vec<Event>,
}

/// 手动中止(`POST /request/{id}/abort`、`POST /sessions/{id}/abort`)叫停了哪几个请求。
///
/// **叫停是立刻的,结局随后到**:每个请求照常报一条 [`Event::RequestFailed`](`source` 是
/// [`FailureSource::Aborted`]),通常在这条响应之前或紧跟着。叫停之后才跑完的不改。
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "ts", derive(ts_rs::TS))]
pub struct Aborted {
/// 叫停的请求号,开始得早的在前
pub requests: Vec<u64>,
}

/// 配置文件没通过校验的那一次。字段和 [`Event::ConfigRejected`] 一样。
///
/// **只记外部改动**(在编辑器里改的、命令行写的):界面自己写坏的根本没落盘,
Expand Down Expand Up @@ -1973,7 +2020,7 @@ pub struct RetentionView {
pub body_bytes_now: u64,
}

/// 上游失败之后停用多久、流开头最多等多久。和配置的 `failover` 一一对应,
/// 上游失败之后停用多久、多久没有内容就不再等。和配置的 `failover` 一一对应,
/// 没写的是默认值 —— **界面显示的就是真在用的数**,不是「空 = 默认」。
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[cfg_attr(feature = "ts", derive(ts_rs::TS))]
Expand All @@ -1990,10 +2037,9 @@ pub struct FailoverView {
pub quota_pause_secs: u64,
/// 限流时按 `Retry-After` 停用,最多多少秒
pub rate_limit_max_pause_secs: u64,
/// 流式回答的开头最多等多少秒
pub stream_start_wait_secs: u64,
/// 等过 `stream_start_wait_secs` 还没有内容就换下一家(最后一家照常等)
pub next_on_slow_start: bool,
/// 上游多少秒没有内容就不再等它(无响应超时):从请求发出去起算,每来一段内容重新
/// 计时。客户端还什么都没收到时换下一家,已经收到一部分的报错收尾
pub idle_timeout_secs: u64,
/// 一个请求合计最多等多少秒:等密钥的分钟、小时上限空出名额,和等满着(`max_concurrent`)
/// 的上游空出位置,共用这一段。0 是不等
pub slot_wait_secs: u64,
Expand Down Expand Up @@ -5833,7 +5879,9 @@ mod tests {
for a in [
attempt(Served, Some(200), None),
attempt(Status, Some(503), None),
attempt(SlowStart, None, Some("gw.slow_start")),
attempt(IdleTimeout, None, Some("gw.upstream.idle_timeout")),
attempt(IdleTimeout, Some(200), Some("gw.upstream.idle_timeout")),
attempt(Aborted, None, Some("gw.request.aborted")),
attempt(Estimated, Some(404), None),
attempt(Error, None, Some("gw.upstream.timeout")),
attempt(Error, None, Some("gw.upstream.stream_opening_error")),
Expand Down
51 changes: 21 additions & 30 deletions crates/tw-config/src/failover.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
//! 故障转移:一家上游失败之后停用多久、流开头最多等多久、等不到内容换不换下一家、上游
//! 满着时最多等多久。
//! 故障转移:一家上游失败之后停用多久、多久没有内容就不再等它、上游满着时最多等多久。

use serde::{Deserialize, Serialize};

Expand Down Expand Up @@ -34,15 +33,15 @@ pub struct Failover {
/// 原因的失败
#[serde(default = "d_rate_limit_max_pause_secs")]
pub rate_limit_max_pause_secs: u64,
/// 流式回答开头最多等多少秒。在第一段内容到达之前,上游报的错误照样换下一家;
/// 等过这么久还没有内容,就不再等,把已经收到的交给客户端
#[serde(default = "d_stream_start_wait_secs")]
pub stream_start_wait_secs: u64,
/// 流式回答等过 [`Self::stream_start_wait_secs`] 还没有内容时,放弃这一家、换下一家。
/// **最后一家不换**,照常等下去;这一家不停用,也不算一次失败。默认关:先想好再
/// 输出的模型开头本来就慢,开着时要把等待调长
#[serde(default)]
pub next_on_slow_start: bool,
/// 上游多少秒没有内容就不再等它(无响应超时)。从请求发给它的那一刻算起,**每来一段
/// 内容重新计时**:正文、推理、工具调用都算,心跳不算(SSE 注释、Anthropic 的 `ping`、
/// 只有角色的空块……)—— 否则一家只发心跳的上游会一直挂着。整包的回答没有「一段段」,
/// 从发出去到整份回来算一段。
///
/// 到点时还没有内容交给客户端的,这一家记一次失败,换下一家;没有下一家了就回超时
/// 错误。已经交出去一部分的换不了(客户端会收到两遍开头),按客户端的格式报错收尾
#[serde(default = "d_idle_timeout_secs")]
pub idle_timeout_secs: u64,
/// 一个请求最多等多少秒,**整个请求合起来算**:准入时等密钥的分钟、小时上限空出名额,
/// 之后上游的并发数满了(`providers[].max_concurrent`)时等空位,共用这一段。留在那一家
/// 的对话等它空出来,候选都满了时等先空出来的那一家;等不到的换下一家,或者回 429。
Expand All @@ -69,8 +68,8 @@ fn d_quota_pause_secs() -> u64 {
fn d_rate_limit_max_pause_secs() -> u64 {
3600
}
fn d_stream_start_wait_secs() -> u64 {
15
fn d_idle_timeout_secs() -> u64 {
300
}
fn d_slot_wait_secs() -> u64 {
30
Expand All @@ -85,8 +84,7 @@ impl Default for Failover {
no_balance_pause_secs: d_no_balance_pause_secs(),
quota_pause_secs: d_quota_pause_secs(),
rate_limit_max_pause_secs: d_rate_limit_max_pause_secs(),
stream_start_wait_secs: d_stream_start_wait_secs(),
next_on_slow_start: false,
idle_timeout_secs: d_idle_timeout_secs(),
slot_wait_secs: d_slot_wait_secs(),
}
}
Expand All @@ -96,12 +94,11 @@ impl Default for Failover {
/// 说的时刻,周额度本来就可能在七天之后
pub const MAX_PAUSE_SECS: u64 = 7 * 24 * 3600;

/// 流开头最多等多少秒。再长的话,一家卡在半路的上游会让客户端先超时
pub const MAX_STREAM_START_WAIT_SECS: u64 = 120;
/// 无响应超时最少写多少秒。再短的话,先想好再输出的模型还没开口就被放弃了
pub const MIN_IDLE_TIMEOUT_SECS: u64 = 30;

/// 开着「开头慢就换下一家」时,流开头至少等多少秒。再短的话,平常的请求还没开口就被
/// 切掉了
pub const MIN_SLOW_START_WAIT_SECS: u64 = 5;
/// 无响应超时最多写多少秒:一小时。再长就等于不设
pub const MAX_IDLE_TIMEOUT_SECS: u64 = 3600;

/// 等空位最多写多少秒。等的时候客户端一个字节都收不到,再长的话它先超时了
pub const MAX_SLOT_WAIT_SECS: u64 = 300;
Expand Down Expand Up @@ -137,21 +134,15 @@ impl Failover {
MAX_PAUSE_SECS,
),
(
"stream_start_wait_secs",
self.stream_start_wait_secs,
1,
MAX_STREAM_START_WAIT_SECS,
"idle_timeout_secs",
self.idle_timeout_secs,
MIN_IDLE_TIMEOUT_SECS,
MAX_IDLE_TIMEOUT_SECS,
),
("slot_wait_secs", self.slot_wait_secs, 0, MAX_SLOT_WAIT_SECS),
];
fields
.into_iter()
.find(|(_, v, min, max)| v < min || v > max)
}

/// 开着「开头慢就换下一家」、等待却短于 [`MIN_SLOW_START_WAIT_SECS`]:写的等待秒数。
pub(crate) fn slow_start_too_short(&self) -> Option<u64> {
(self.next_on_slow_start && self.stream_start_wait_secs < MIN_SLOW_START_WAIT_SECS)
.then_some(self.stream_start_wait_secs)
}
}
3 changes: 1 addition & 2 deletions crates/tw-config/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1201,8 +1201,7 @@ pub fn write(path: &Path, cfg: &Config) -> Result<(), WriteError> {
}

pub use failover::{
Failover, MAX_PAUSE_SECS, MAX_SLOT_WAIT_SECS, MAX_STREAM_START_WAIT_SECS,
MIN_SLOW_START_WAIT_SECS,
Failover, MAX_IDLE_TIMEOUT_SECS, MAX_PAUSE_SECS, MAX_SLOT_WAIT_SECS, MIN_IDLE_TIMEOUT_SECS,
};
pub use probes::{ClientProbes, ProbeAction};
pub use reload::{Rejected, Stage, stand_in, try_parse};
Expand Down
Loading
Loading