From db24bc1a2196ada0cd4979b9065dc8dedf4881b2 Mon Sep 17 00:00:00 2001 From: fylorn <249551762+fylorn@users.noreply.github.com> Date: Fri, 2 Oct 2026 22:54:24 +0800 Subject: [PATCH] Record an error the upstream answered as a failed request; TurnView carries the status When an upstream answered 4xx (or 3xx) and the gateway handed the answer on to the client, the ending was RequestFinished with that status and an empty error. Everything that counts failures reads `error`, so the request looked successful everywhere: the traffic list, the overview's failed count, the session's errors, the upstream check-up, the history search's failed filter. Session detail (TurnView) treated the turn as done while the transcript, which reads the status, gave it no answer: the desktop app's conversation showed only the user's message, with no answer and no reason. The gateway now reports these as RequestFailed (source `upstream`): - relay marks a non-2xx answer on the Ending (`Ending::refused`). The Ending keeps the head of the error body whether or not bodies are stored, and on `finished` reports the failure with what the upstream said: gw.upstream.status_message {upstream, status, message} ("Upstream `x` answered 400: prompt is too long..."), read in the upstream's format with tw_dialect::convert::error_message (now shared with Session::error). Plain text is used as it is; an empty body or an HTML page falls back to gw.upstream.status. The words are masked with the request's redaction, like the stored response body, and capped at 500 characters. - A Bedrock credential refusal already replaces AWS's words with ours; that message is now the recorded reason. The `[ThinkWatch]` prefix moves out of those four messages into the body sent to the client, as with the gateway's own errors. - A client that leaves while the error is being passed on is still a cancellation. WebSocket upgrades (101) are unchanged. The row's status still comes from RequestHeaders, so one rule holds in every place that counts: a request failed if and only if `error` is set. No SQL changes, and cancellations are counted as before. TurnView gains `status` (the upstream's status code, null when the request never got one), so a failed turn can say what the upstream answered. CONTROL_API_VERSION stays at 32 (unreleased); its note mentions both changes. Tests: the Ending (the upstream's words, nothing readable, masking, our words for Bedrock, a cancellation during an error answer); a passed-on 400 reaches the client unchanged and ends once as a failure; the Bedrock refusal's recorded reason; rows, turns, session and summary in the recorder; and an end-to-end control-plane test through a real gateway and store with a 200, a passed-on 400 and a cancelled turn in one session, read back from /sessions/{id}, /history, /summary, /sessions and the transcript. Co-Authored-By: Claude Opus 5.5 --- crates/tw-api/msg-codes.txt | 1 + crates/tw-api/src/lib.rs | 35 +- crates/tw-api/src/ts.rs | 14 + crates/tw-control/src/lib.rs | 1 + crates/tw-control/tests/error_answers.rs | 354 ++++++++++++++++++ crates/tw-dialect/src/convert.rs | 79 +++- crates/tw-gateway/src/ending.rs | 269 ++++++++++++- crates/tw-gateway/src/server/pipeline/hop.rs | 41 +- .../tw-gateway/src/server/pipeline/relay.rs | 6 + crates/tw-gateway/tests/bedrock.rs | 16 +- crates/tw-gateway/tests/endings.rs | 46 +++ crates/tw-store/src/db.rs | 8 +- crates/tw-store/src/recorder.rs | 101 +++++ 13 files changed, 928 insertions(+), 43 deletions(-) create mode 100644 crates/tw-control/tests/error_answers.rs diff --git a/crates/tw-api/msg-codes.txt b/crates/tw-api/msg-codes.txt index d4146dd6..872eaada 100644 --- a/crates/tw-api/msg-codes.txt +++ b/crates/tw-api/msg-codes.txt @@ -286,6 +286,7 @@ gw.upstream.forward_failed gw.upstream.rate_limited gw.upstream.sign_failed gw.upstream.status +gw.upstream.status_message gw.upstream.stream_error gw.upstream.stream_exception gw.upstream.stream_opening_error diff --git a/crates/tw-api/src/lib.rs b/crates/tw-api/src/lib.rs index 494a8edd..48186b3f 100644 --- a/crates/tw-api/src/lib.rs +++ b/crates/tw-api/src/lib.rs @@ -675,6 +675,12 @@ pub const MSG_CODES: &str = include_str!("../msg-codes.txt"); /// **32 起会话能读成一段对话**:新端点 `GET /sessions/{id}/transcript`([`Transcript`]) /// 从存下来的正文里读出每一轮新说的话、回答、推理、工具调用和结果,读不到的地方逐轮说出来 /// ([`TranscriptGap`])。照 31 写的界面只有每一轮的用量和金额。 +/// +/// 同一版起**上游回了错误、原样交给客户端的请求是失败的**:结局是 [`Event::RequestFailed`] +/// (`upstream`,上游在错误正文里说的话是 `gw.upstream.status_message`),记录的 `error` +/// 有值,概览、会话、上游体检都数它。以前它是 `RequestFinished`、`error` 为空,对话里那 +/// 一轮只剩用户的话,没有回答,也说不出为什么。[`TurnView`] 多了 `status`,说得出上游回了 +/// 什么。照 31 写的界面看不到那个状态码。 pub const CONTROL_API_VERSION: u32 = 32; #[derive(Debug, Clone, Serialize, Deserialize)] @@ -838,7 +844,9 @@ pub enum Event { /// 只有成功的流式响应有。非流式的整段一起到,没有「第一个」;不带 `alt=sse` 的 /// Gemini 流和 WebSocket 那条路不在这里解析,也没有。 RequestFirstToken { id: u64, ttft_ms: u64 }, - /// 结束了 + /// 结束了:上游回的是成功的状态码(2xx;WebSocket 那条路是升级成功的 101),回答 + /// 交完了。**上游回了别的、原样交给了客户端的不是这一条**,是 `RequestFailed` —— + /// 客户端拿到的是上游的错误,不是回答 RequestFinished { id: u64, /// 客户端要的模型名,和 `RequestStarted` 里的是同一个。 @@ -877,6 +885,10 @@ pub enum Event { /// (`auth` / `config` / `upstream` / `request` / `rate_limited` / /// `denied`),另外多一个 `internal`:网关自己的代码 /// 崩掉了。它只出现在这里 —— 那时往往已经没有一个 HTTP 响应能带上它。 + /// + /// **上游回了错误(不是 2xx)、原样交给客户端的也是失败**(`upstream`):先有那个 + /// 状态码的 `RequestHeaders`,`message` 是上游在错误正文里说的话 + /// (`gw.upstream.status_message`,读不出来的是 `gw.upstream.status`)。 RequestFailed { id: u64, /// 模型名。理由见 `RequestFinished::model` @@ -1323,7 +1335,8 @@ slug_enum! { /// 尝试链里一跳的结果。 pub enum AttemptOutcome { /// 这一跳接下了请求,尝试链到此为止。上游回的是 4xx 也算 —— 请求本身有 - /// 问题,换一个上游也一样被拒 + /// 问题,换一个上游也一样被拒;那个错误原样交给客户端,**请求本身记成失败** + /// (见 [`Event::RequestFailed`]) Served = "served", /// 上游返回 5xx 或 429,换下一个上游 Status = "status", @@ -3274,6 +3287,8 @@ pub enum CostDim { #[cfg_attr(feature = "ts", derive(ts_rs::TS))] pub struct Summary { pub requests: i64, + /// 失败的请求([`HistoryRow::error`] 有值的):网关没转发成的,和上游回了错误、原样 + /// 交给客户端的。客户端先走了的不算(见 [`HistoryRow::cancelled`]) pub failed: i64, /// 本地应答的次数。**是个正向数字**,单独显示 pub locally_answered: i64, @@ -3294,7 +3309,7 @@ pub struct Summary { /// /// 和 `unpriced_requests` 一样让金额合计偏低,但配价格解决不了它 —— /// 界面上是两句不同的话。上游确实接下了的才算:成功的响应和客户端 - /// 取消的,失败的和上游回了 4xx 的不算。 + /// 取消的,失败的不算(上游回了错误的也是失败,那种响应不计费)。 pub no_usage_requests: i64, /// 用了缓存之后净省下多少微分。 /// @@ -3531,7 +3546,11 @@ pub struct HistoryRow { pub cost_micros: Option, /// 这个成本是估的吗。**界面上要标出来** pub cost_estimated: bool, - /// 失败的原因。**带着码** —— 翻历史时界面照样能说自己那句话; + /// 失败的原因。**带着码** —— 翻历史时界面照样能说自己那句话。 + /// + /// **有它就是失败**,数失败的地方都按它数(概览、会话、上游体检、搜索的筛选):网关 + /// 没转发成的(连不上、被拒、断在半路),和上游回了错误(不是 2xx)、原样交给客户端 + /// 的 —— 那时 `status` 是上游回的那个状态码,这一句是它在错误正文里说的话 pub error: Option, /// 本地应答的 pub local: bool, @@ -4021,6 +4040,14 @@ pub struct TurnView { /// **没有价格就是 None,不是 0** pub cost_micros: Option, pub duration_ms: Option, + /// 上游回的状态码,和 [`HistoryRow::status`] 同一个。没走到上游的没有:连不上、 + /// 被规则拒绝、客户端在响应头到之前就走了 + pub status: Option, + /// 这一轮为什么失败(见 [`HistoryRow::error`])。没失败是 None。 + /// + /// **上游回了错误、原样交给客户端的也在这里**:`status` 是那个状态码,这一句是 + /// 上游在错误正文里说的话(`gw.upstream.status_message`,读不出来的是 + /// `gw.upstream.status`)。网关自己没转发成的没有 `status`,原因只在这一句里 pub error: Option, /// 客户端没等到这一轮结束就走了(见 `HistoryRow::cancelled`) pub cancelled: bool, diff --git a/crates/tw-api/src/ts.rs b/crates/tw-api/src/ts.rs index f2f28142..0752ba92 100644 --- a/crates/tw-api/src/ts.rs +++ b/crates/tw-api/src/ts.rs @@ -212,6 +212,20 @@ mod tests { assert!(!login.contains("plan"), "{login}"); } + /// 会话里的一轮说得出上游回了什么:状态码和失败的原因都是必有的字段,没有时是 null + #[test] + fn a_turn_says_what_the_upstream_answered() { + let ts = typescript(); + let turn = decl_of(&ts, "TurnView"); + for field in [ + "status: number | null", + "error: Msg | null", + "cancelled: boolean", + ] { + assert!(turn.contains(field), "{field}: {turn}"); + } + } + /// 命中数带着记录从哪一刻起是全的:必有的字段,没有记录时是 null,不是省掉。 #[test] fn route_stats_say_where_their_history_starts() { diff --git a/crates/tw-control/src/lib.rs b/crates/tw-control/src/lib.rs index 93ac21a6..24f871ee 100644 --- a/crates/tw-control/src/lib.rs +++ b/crates/tw-control/src/lib.rs @@ -1455,6 +1455,7 @@ fn turn_view(t: &tw_store::db::TurnRow) -> tw_api::TurnView { cache_read_tokens: t.cache_read_tokens, cost_micros: t.cost_micros, duration_ms: t.duration_ms, + status: t.status, error: t.error.clone(), cancelled: t.cancelled, cost_estimated: t.cost_estimated, diff --git a/crates/tw-control/tests/error_answers.rs b/crates/tw-control/tests/error_answers.rs new file mode 100644 index 00000000..57ee4e54 --- /dev/null +++ b/crates/tw-control/tests/error_answers.rs @@ -0,0 +1,354 @@ +//! 上游回了错误、原样交给客户端的那一轮**是失败的**,界面读的每一处都这么说。 +//! +//! 整条路走一遍:真的网关、一个按最后一句话回答 / 回 400 / 吐一帧就不再说话的假上游、 +//! 存储层写进临时目录、控制面。同一次会话里三轮:成功的、上游回了 400 的、客户端中途 +//! 走掉的。然后从界面读的几个端点看它们:会话详情(`TurnView`)、流量(`/history`)、 +//! 概览(`/summary`)、会话列表、对话记录。 +//! +//! 以前上游的 4xx 落库时 `error` 是空的:会话详情把那一轮当成成功,对话里只剩用户的那句 +//! 话,没有回答,也没有原因。 + +use std::net::SocketAddr; +use std::sync::Arc; +use std::time::Duration; + +use axum::body::Body; +use axum::http::{Request, StatusCode}; +use futures::StreamExt; +use serde_json::{Value, json}; +use tower::ServiceExt; +use tw_control::{ConfigManager, ControlState}; + +const TOO_LONG: &str = "prompt is too long: 212000 tokens > 200000 maximum"; + +const MESSAGE_START: &str = "event: message_start\ndata: {\"type\":\"message_start\",\"message\":{\"model\":\"claude-sonnet-4-5\",\"usage\":{\"input_tokens\":10,\"output_tokens\":1}}}\n\n"; + +/// 假上游,看最后一句话:`too long` 回 400(Anthropic 的错误体),`slow` 吐一帧之后不再 +/// 说话(客户端会走掉),别的好好回答 +async fn upstream() -> SocketAddr { + let app = axum::Router::new().route( + "/v1/messages", + axum::routing::post(|body: bytes::Bytes| async move { + let v: Value = serde_json::from_slice(&body).unwrap_or_default(); + let last = v["messages"] + .as_array() + .and_then(|m| m.last()) + .and_then(|m| m["content"].as_str()) + .unwrap_or_default() + .to_string(); + let answer = axum::response::Response::builder(); + match last.as_str() { + "too long" => answer + .status(400) + .header("content-type", "application/json") + .body(Body::from( + json!({"type": "error", "error": {"type": "invalid_request_error", "message": TOO_LONG}}) + .to_string(), + )) + .unwrap(), + "slow" => { + let first = futures::stream::once(async { + Ok::<_, std::io::Error>(bytes::Bytes::from_static(MESSAGE_START.as_bytes())) + }); + let stall = futures::stream::once(async { + tokio::time::sleep(Duration::from_secs(60)).await; + Ok(bytes::Bytes::new()) + }); + answer + .header("content-type", "text/event-stream") + .body(Body::from_stream(first.chain(stall))) + .unwrap() + } + _ => answer + .header("content-type", "application/json") + .body(Body::from( + json!({"type": "message", "model": "claude-sonnet-4-5", "role": "assistant", + "content": [{"type": "text", "text": "好"}], "stop_reason": "end_turn", + "usage": {"input_tokens": 10, "output_tokens": 5}}) + .to_string(), + )) + .unwrap(), + } + }), + ); + let l = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = l.local_addr().unwrap(); + tokio::spawn(async move { axum::serve(l, app).await.unwrap() }); + addr +} + +struct World { + _dir: tempfile::TempDir, + blobs: std::path::PathBuf, + gw: SocketAddr, + store: Arc>, + app: axum::Router, +} + +/// 网关、存储层、控制面,**和 `twcore` 一样接**(正文也落盘:对话记录要读它)。 +async fn world() -> World { + let up = upstream().await; + let d = tempfile::tempdir().unwrap(); + let yaml = format!( + "version: 1\nlisten:\n control:\n key: c0ffee00c0ffee00c0ffee00c0ffee00c0ffee00c0ffee00c0ffee00c0ffee00\n\ + clients:\n - name: 我\n key: tw-k\n\ + providers:\n - name: 中转\n base_url: http://{up}\n key: sk-upstream\n protocol: anthropic\n" + ); + let p = d.path().join("config.yaml"); + std::fs::write(&p, &yaml).unwrap(); + let cfg: tw_config::Config = serde_yaml_ng::from_str(&yaml).unwrap(); + let gw = tw_gateway::AppState::new(cfg).unwrap(); + + let (sink, mut bodies) = tw_gateway::bodies::channel(); + gw.set_body_sink(sink); + let (tx, rx) = tokio::sync::mpsc::channel(1); + tokio::spawn(async move { + while let Some(b) = bodies.recv().await { + let disk = tokio::task::spawn_blocking(move || b.for_disk()) + .await + .unwrap(); + let which = match disk.kind { + tw_gateway::bodies::BodyKind::Request => tw_store::Which::Request, + tw_gateway::bodies::BodyKind::Response => tw_store::Which::Response, + }; + let stored = tw_store::StoredBody { + id: disk.id, + at_ms: disk.at_ms, + which, + body: disk.body, + original_len: disk.original_len, + }; + if tx.send(stored).await.is_err() { + return; + } + } + }); + let blobs = d.path().join("blobs"); + let store = tw_store::task::spawn( + tw_store::Recorder::new( + tw_store::Db::open(&d.path().join("data.db")).unwrap(), + tw_store::Blobs::new(blobs.clone()), + tw_pricing::shared(tw_pricing::PriceBook::builtin().unwrap()), + ), + gw.bus.subscribe(), + rx, + ); + let state = ControlState { + shutdown: Default::default(), + remote: Default::default(), + cfg: Arc::new(ConfigManager::new(p, gw.clone(), gw.bus.clone())), + gateway: gw.clone(), + store: Some(store.clone()), + started: std::time::Instant::now(), + price_updater: Default::default(), + chatgpt: Default::default(), + zai: Default::default(), + }; + let app = tw_control::router(state); + let addr = tw_gateway::serve(gw, ([127, 0, 0, 1], 0).into()) + .await + .unwrap(); + tokio::time::sleep(Duration::from_millis(50)).await; + World { + _dir: d, + blobs, + gw: addr, + store, + app, + } +} + +/// 同一次会话的一轮:开头那句话都一样(会话按它认)。`later` 是第一轮之后带着的:上一轮 +/// 的回答,和这一轮新说的 +fn turn(later: &[&str], stream: bool) -> String { + let mut messages = vec![json!({"role": "user", "content": "把这个函数改短一点"})]; + for (i, text) in later.iter().enumerate() { + let role = if i % 2 == 0 { "assistant" } else { "user" }; + messages.push(json!({"role": role, "content": text})); + } + json!({ + "model": "claude-sonnet-4-5", + "max_tokens": 64, + "stream": stream, + "messages": messages, + }) + .to_string() +} + +fn send(w: &World, body: String) -> reqwest::RequestBuilder { + reqwest::Client::builder() + .no_proxy() + .build() + .unwrap() + .post(format!("http://{}/v1/messages", w.gw)) + .header("x-api-key", "tw-k") + .header("content-type", "application/json") + .body(body) +} + +/// 等第 `n` 条请求落库、它的请求体和响应体都落了盘,交回它的号 +async fn landed(w: &World, n: i64) -> i64 { + let deadline = std::time::Instant::now() + Duration::from_secs(10); + loop { + let id = { + let g = w.store.lock().await; + (g.db().count().unwrap() >= n).then(|| g.db().recent(None, 1).unwrap()[0].id) + }; + if let Some(id) = id { + let names = stored_names(&w.blobs); + if [format!("{id}.req"), format!("{id}.res")] + .iter() + .all(|f| names.contains(f)) + { + return id; + } + } + assert!( + std::time::Instant::now() < deadline, + "第 {n} 条请求没有落库,或者正文没有落盘" + ); + tokio::time::sleep(Duration::from_millis(20)).await; + } +} + +/// 正文目录里每一个文件的名字 +fn stored_names(root: &std::path::Path) -> Vec { + std::fs::read_dir(root) + .into_iter() + .flatten() + .flatten() + .flat_map(|day| { + std::fs::read_dir(day.path()) + .into_iter() + .flatten() + .flatten() + }) + .map(|f| f.file_name().to_string_lossy().to_string()) + .collect() +} + +async fn get(app: &axum::Router, path: &str) -> Value { + let r = app + .clone() + .oneshot(Request::get(path).body(Body::empty()).unwrap()) + .await + .unwrap(); + assert_eq!(r.status(), StatusCode::OK, "{path}"); + let b = axum::body::to_bytes(r.into_body(), usize::MAX) + .await + .unwrap(); + serde_json::from_slice(&b).unwrap() +} + +#[tokio::test] +async fn an_error_the_upstream_answered_is_a_failed_turn_everywhere_the_ui_reads() { + let w = world().await; + + // 第一轮:好好回答了 + let r = send(&w, turn(&[], false)).send().await.unwrap(); + assert_eq!(r.status(), 200); + let done = landed(&w, 1).await; + + // 第二轮:上游回 400,原话原样交给客户端 + let r = send(&w, turn(&["好", "too long"], false)) + .send() + .await + .unwrap(); + assert_eq!(r.status(), 400); + assert!(r.text().await.unwrap().contains(TOO_LONG)); + let refused = landed(&w, 2).await; + + // 第三轮:客户端拿到第一帧就走了(Claude Code 里按一下 Esc) + let r = send(&w, turn(&["好", "slow"], true)).send().await.unwrap(); + assert_eq!(r.status(), 200); + let mut body = r.bytes_stream(); + body.next().await.unwrap().unwrap(); + drop(body); + let left = landed(&w, 3).await; + + // ── 会话详情:每一轮带着状态码;上游回了错误的那一轮是失败,原因是上游的原话 + let sessions = get(&w.app, "/sessions").await; + let sessions = sessions.as_array().unwrap(); + assert_eq!(sessions.len(), 1, "三轮该是同一次会话:{sessions:?}"); + let s = &sessions[0]; + assert_eq!( + (s["turns"].as_u64(), s["errors"].as_u64()), + (Some(3), Some(1)), + "{s}" + ); + let id = s["id"].as_str().unwrap(); + let detail = get(&w.app, &format!("/sessions/{id}")).await; + let turns = detail["turns"].as_array().unwrap(); + let ids: Vec = turns.iter().map(|t| t["id"].as_i64().unwrap()).collect(); + assert_eq!(ids, [done, refused, left]); + + let (ok, failed, cancelled) = (&turns[0], &turns[1], &turns[2]); + assert_eq!(ok["status"], 200, "{ok}"); + assert_eq!(ok["error"], Value::Null, "{ok}"); + assert_eq!(ok["cancelled"], false); + + assert_eq!(failed["status"], 400, "{failed}"); + assert_eq!(failed["cancelled"], false); + let why = &failed["error"]; + assert_eq!(why["code"], "gw.upstream.status_message", "{why}"); + assert_eq!(why["args"]["upstream"], "中转"); + assert_eq!(why["args"]["status"], "400"); + assert_eq!(why["args"]["message"], TOO_LONG); + assert_eq!( + why["text"].as_str().unwrap(), + format!("Upstream `中转` answered 400: {TOO_LONG}") + ); + + // 取消照旧:不是失败,带着响应头那一刻的状态码 + assert_eq!(cancelled["cancelled"], true, "{cancelled}"); + assert_eq!( + cancelled["error"], + Value::Null, + "取消被算成了失败:{cancelled}" + ); + assert_eq!(cancelled["status"], 200); + + // ── 流量:同一批行,同一个说法 + let history = get(&w.app, "/history").await; + let row = |rid: i64| { + history + .as_array() + .unwrap() + .iter() + .find(|r| r["id"].as_i64() == Some(rid)) + .unwrap_or_else(|| panic!("流量里没有 {rid}:{history}")) + .clone() + }; + let (ok, failed, cancelled) = (row(done), row(refused), row(left)); + assert_eq!((&ok["status"], &ok["error"]), (&json!(200), &Value::Null)); + assert_eq!(failed["status"], 400); + assert_eq!(failed["error"], *why); + assert_eq!( + (&cancelled["cancelled"], &cancelled["error"]), + (&json!(true), &Value::Null) + ); + + // ── 概览:失败一条,取消的不算 + let summary = get(&w.app, "/summary?from_ms=0&to_ms=9999999999999").await; + assert_eq!(summary["requests"], 3, "{summary}"); + assert_eq!(summary["failed"], 1, "{summary}"); + + // ── 对话:失败的那一轮没有回答,也不算缺 —— 原因在会话详情里 + let transcript = get(&w.app, &format!("/sessions/{id}/transcript")).await; + // 对话里的请求号是字符串 + let key = refused.to_string(); + let t = transcript["turns"] + .as_array() + .unwrap() + .iter() + .find(|t| t["id"] == key.as_str()) + .unwrap_or_else(|| panic!("对话里没有 {refused}:{transcript}")) + .clone(); + assert_eq!(t["output"], json!([]), "{t}"); + assert_eq!(t["gaps"], json!([]), "{t}"); + assert_eq!( + t["input"][0]["parts"][0], + json!({"kind": "text", "text": "too long"}), + "{t}" + ); +} diff --git a/crates/tw-dialect/src/convert.rs b/crates/tw-dialect/src/convert.rs index b2e097a5..d0775178 100644 --- a/crates/tw-dialect/src/convert.rs +++ b/crates/tw-dialect/src/convert.rs @@ -307,23 +307,15 @@ impl Session { /// 上游的错误响应 → 客户端格式的错误体。状态码不变,说明取上游的原话 pub fn error(&self, status: u16, body: &[u8]) -> Vec { - let message = serde_json::from_slice::(body) - .ok() - .and_then(|v| match self.upstream { - Dialect::Anthropic => anthropic::response::error_message(&v), - Dialect::Chat | Dialect::Responses => chat::response::error_message(&v), - Dialect::Gemini => gemini::response::error_message(&v), - Dialect::Bedrock => bedrock::response::error_message(&v), - }) - .unwrap_or_else(|| { - let text = String::from_utf8_lossy(body); - let text = text.trim(); - if text.is_empty() { - format!("The upstream answered HTTP {status} with nothing else.") - } else { - text.chars().take(2000).collect() - } - }); + let message = error_message(self.upstream, body).unwrap_or_else(|| { + let text = String::from_utf8_lossy(body); + let text = text.trim(); + if text.is_empty() { + format!("The upstream answered HTTP {status} with nothing else.") + } else { + text.chars().take(2000).collect() + } + }); error_body(self.client, status, &message) } @@ -380,6 +372,19 @@ impl Session { } } +/// 上游的错误体里那句说明,按它的格式读:Anthropic、OpenAI、Gemini 的 `error.message` +/// (OpenAI 有时只给一个字符串),Bedrock 的 `message`。不是 JSON、或者没有这一项的是 None +/// —— 正文里别的东西怎么办,由调用方定。 +pub fn error_message(upstream: Dialect, body: &[u8]) -> Option { + let v = serde_json::from_slice::(body).ok()?; + match upstream { + Dialect::Anthropic => anthropic::response::error_message(&v), + Dialect::Chat | Dialect::Responses => chat::response::error_message(&v), + Dialect::Gemini => gemini::response::error_message(&v), + Dialect::Bedrock => bedrock::response::error_message(&v), + } +} + /// 给某种格式的客户端的错误体。 pub fn error_body(client: Dialect, status: u16, message: &str) -> Vec { let v = match client { @@ -1193,6 +1198,46 @@ mod tests { assert!(v["error"]["message"].as_str().unwrap().contains("502")); } + /// 每种格式的错误体里那句说明。读不出来的交回 None,不拿正文凑一句 + #[test] + fn the_message_in_an_error_body_is_read_in_the_upstreams_shape() { + let said = |d, body: &str| error_message(d, body.as_bytes()); + assert_eq!( + said( + Dialect::Anthropic, + r#"{"type":"error","error":{"type":"invalid_request_error","message":"prompt is too long"}}"# + ) + .as_deref(), + Some("prompt is too long") + ); + assert_eq!( + said( + Dialect::Responses, + r#"{"error":{"message":"max_output_tokens is too large","code":"invalid_value"}}"# + ) + .as_deref(), + Some("max_output_tokens is too large") + ); + assert_eq!( + said(Dialect::Chat, r#"{"error":"model not found"}"#).as_deref(), + Some("model not found") + ); + assert_eq!( + said( + Dialect::Gemini, + r#"[{"error":{"code":400,"message":"API key not valid","status":"INVALID_ARGUMENT"}}]"# + ) + .as_deref(), + Some("API key not valid") + ); + assert_eq!( + said(Dialect::Bedrock, r#"{"Message":"Malformed input request"}"#).as_deref(), + Some("Malformed input request") + ); + assert_eq!(said(Dialect::Anthropic, r#"{"detail":"Not Found"}"#), None); + assert_eq!(said(Dialect::Anthropic, "Bad Gateway"), None); + } + #[test] fn a_freeform_call_from_a_json_only_upstream_is_unwrapped_in_the_stream() { let mut s = Session::for_test(Dialect::Responses, Dialect::Anthropic); diff --git a/crates/tw-gateway/src/ending.rs b/crates/tw-gateway/src/ending.rs index efd4fa85..e0ef4eb2 100644 --- a/crates/tw-gateway/src/ending.rs +++ b/crates/tw-gateway/src/ending.rs @@ -1,8 +1,8 @@ //! 一个请求怎么收场。 //! //! 发出 `RequestStarted` 的那一刻起,这个请求就欠总线一个结局:跑完了是 -//! `RequestFinished`,出错了是 `RequestFailed`,客户端先走了是 -//! `RequestCancelled`。**恰好一个** —— 少一个,存储层永远等不到它:那一行 +//! `RequestFinished`,出错了是 `RequestFailed`(上游回的不是 2xx、原样交给客户端的 +//! 也是),客户端先走了是 `RequestCancelled`。**恰好一个** —— 少一个,存储层永远等不到它:那一行 //! 不落库,上游已经计的费从账上消失,界面上那一行也永远停在「进行中」; //! 多一个,同一行会被写两遍。 //! @@ -69,6 +69,9 @@ pub struct Ending { /// 上游在流里报的错:哪一家、原话。**有它就不是成功** —— 响应头是 200,回答却断在了 /// 半路 upstream_error: Option<(String, String)>, + /// 上游回的不是 2xx、原样交给了客户端(见 [`Ending::refused`])。**有它就不是成功**: + /// 客户端拿到的是上游的错误,不是回答 + refusal: Option, /// 这段对话这一次由谁回答(见 [`crate::affinity`])。**成功走完了才记**:失败的、 /// 半路断了的不算回答过,下一次照常排序 answer: Option, @@ -76,6 +79,70 @@ pub struct Ending { told: bool, } +/// 上游在错误正文里说的那句话,最多留多少个字。错误说明都很短,太长的多半是把请求 +/// 整段回显了出来 +const SAID_MAX: usize = 500; + +/// 上游回的不是 2xx(见 [`Ending::refused`]):失败的原因从这里读。 +struct Refusal { + /// 哪一家。报出去的那句话要点名 + provider: String, + /// 它说的格式:错误正文按它读 + dialect: ir::Dialect, + /// 网关替它说的那句话。有它就不读正文 + ours: Option, + /// 错误正文的开头,最多 [`crate::failure::BODY_PEEK`]。**不管留不留档都攒**:没有它 + /// 就说不出上游为什么拒绝 + head: Vec, +} + +impl Refusal { + fn feed(&mut self, chunk: &[u8]) { + if self.ours.is_some() { + return; + } + let room = crate::failure::BODY_PEEK.saturating_sub(self.head.len()); + self.head.extend_from_slice(&chunk[..chunk.len().min(room)]); + } + + /// 这个请求为什么失败。上游的原话和存下来的正文同一套打码(`redaction`):存下来的 + /// 那份打了码,这里照原样留着就白打了 + fn why(self, status: u16, redaction: Option<&Redaction>) -> Msg { + if let Some(ours) = self.ours { + return ours; + } + let said = said(self.dialect, &self.head).map(|s| match redaction { + Some(r) => r.apply(&s), + None => Redaction::default().apply(&s), + }); + let upstream = self.provider; + match said { + Some(message) => msg!( + "gw.upstream.status_message", + upstream = upstream, status = status, message = message => + "Upstream `{upstream}` answered {status}: {message}" + ), + None => msg!( + "gw.upstream.status", upstream = upstream, status = status => + "Upstream `{upstream}` answered {status}." + ), + } + } +} + +/// 上游在错误正文里说的那句话:按它的格式读出来的说明,读不出来的就是正文本身。**网页 +/// 不算**(代理、防火墙回的那种错误页):一页 HTML 的开头说明不了什么。最多 [`SAID_MAX`] +/// 个字。 +fn said(dialect: ir::Dialect, head: &[u8]) -> Option { + let text = tw_dialect::convert::error_message(dialect, head) + .unwrap_or_else(|| String::from_utf8_lossy(head).into_owned()); + let text = text.trim(); + if text.is_empty() || text.starts_with('<') { + return None; + } + Some(text.chars().take(SAID_MAX).collect()) +} + /// 盯着流里的错误帧。 struct Watch { /// 哪一家上游在说话。报出去的那句话要点名 @@ -137,6 +204,7 @@ impl Ending { opened: None, watch: None, upstream_error: None, + refusal: None, answer: None, told: false, } @@ -163,6 +231,25 @@ impl Ending { }); } + /// 上游 `provider` 回的不是 2xx,原样交给了客户端(4xx 是请求本身的问题,或者没有 + /// 下一家可换了;3xx 交还客户端,由它决定跟不跟)。 + /// + /// **这个请求是失败的。**客户端拿到的是上游的错误,不是回答:收尾时报失败(见 + /// [`Ending::finished`]),原因是上游在错误正文里说的那句话,按它的格式 `upstream` 读。 + /// 记成结束的话,流量、概览、会话里这一轮都像是成功的,对话里只剩用户的那句话, + /// 回答没有,原因也没有。`ours` 是网关替它说的那句(Bedrock 拒绝凭证时,AWS 的原话 + /// 点名账号,交出去的是它),有它就不读正文。 + /// + /// 客户端没等错误交完就走了的,照旧报取消(见 `Drop`)。 + pub fn refused(&mut self, upstream: ir::Dialect, provider: &str, ours: Option) { + self.refusal = Some(Refusal { + provider: provider.to_string(), + dialect: upstream, + ours, + head: Vec::new(), + }); + } + /// 响应体落盘之前按什么换、打码。**和请求体用同一套**:开始时生效的规则,拦截档下 /// 这个请求的账本 —— 上游回答里的占位符和存下来的请求对得上号。没交代的按出厂规则 /// 打码(见 [`Redaction::default`]) @@ -190,6 +277,9 @@ impl Ending { if self.sink.is_some() { self.tap.feed(chunk); } + if let Some(r) = self.refusal.as_mut() { + r.feed(chunk); + } self.spot(chunk); self.watch_for_errors(chunk); } @@ -248,9 +338,15 @@ impl Ending { self.bytes += bytes as u64; } - /// 走完了。上游在流里报过错的,报的是失败(见 [`Ending::streaming`])。 + /// 走完了。上游在流里报过错的、回的不是 2xx 的,报的是失败(见 [`Ending::streaming`]、 + /// [`Ending::refused`])。 pub fn finished(mut self, status: u16) { self.status = Some(status); + if let Some(r) = self.refusal.take() { + let why = r.why(status, self.redaction.as_ref()); + self.failed(tw_api::FailureSource::Upstream, why); + return; + } // 最后一帧后面不带空行的上游:收尾时再看一眼 if let Some(mut w) = self.watch.take() { let frames = w.frames.flush(); @@ -580,6 +676,173 @@ mod tests { ); } + /// 一个响应头已经到了、回的不是 2xx 的请求,原样交给客户端(见 `Ending::refused`) + fn refused(bus: &tw_observe::EventBus, status: u16, provider: &str) -> Ending { + let mut e = Ending::new(bus.clone(), 7, MODEL.into(), Instant::now(), 1_000, None); + e.responded(status); + e.refused(ir::Dialect::Anthropic, provider, None); + e + } + + /// 那一条失败的原因 + fn failure(rx: &mut tokio::sync::broadcast::Receiver) -> Msg { + match drain(rx).as_slice() { + [ + Event::RequestFailed { + source: tw_api::FailureSource::Upstream, + message, + .. + }, + ] => message.clone(), + other => panic!("该是一条来自上游的失败,实际 {other:?}"), + } + } + + /// 上游回了 4xx、原样交给了客户端:**不是成功**。结局是一条失败,原因是上游在错误 + /// 正文里说的那句话,点名是哪一家、回了什么。记成结束的话,这一轮在流量、概览、 + /// 会话里都像是成功的,对话里只剩用户的那句话 + #[test] + fn an_error_answer_passed_on_is_a_failure_in_the_upstreams_words() { + const TOO_LONG: &[u8] = br#"{"type":"error","error":{"type":"invalid_request_error","message":"prompt is too long: 212000 tokens > 200000 maximum"}}"#; + let bus = tw_observe::EventBus::new(); + let mut rx = bus.subscribe(); + let mut e = refused(&bus, 400, "官方"); + // 错误正文分几块到 + for part in TOO_LONG.chunks(7) { + e.feed(part); + } + e.finished(400); + + match drain(&mut rx).as_slice() { + [ + Event::RequestFailed { + source, + message, + bytes: Some(bytes), + duration_ms: Some(_), + usage: None, + .. + }, + ] => { + assert_eq!(*source, tw_api::FailureSource::Upstream); + assert_eq!(message.code, "gw.upstream.status_message"); + assert_eq!(message.arg("upstream"), "官方"); + assert_eq!(message.arg("status"), "400"); + assert_eq!( + message.arg("message"), + "prompt is too long: 212000 tokens > 200000 maximum" + ); + assert_eq!( + message.text, + "Upstream `官方` answered 400: prompt is too long: 212000 tokens > 200000 maximum" + ); + assert_eq!(*bytes, TOO_LONG.len() as u64); + } + other => panic!("该是一条失败,实际 {other:?}"), + } + } + + /// 正文里读不出一句话的(空的、一页网页)只说状态码;不是 JSON 的纯文本、不认得的 + /// JSON 照原样说 + #[test] + fn an_error_answer_without_readable_words_says_its_status() { + let bus = tw_observe::EventBus::new(); + let mut rx = bus.subscribe(); + let said = |rx: &mut tokio::sync::broadcast::Receiver, body: &[u8]| { + let mut e = refused(&bus, 403, "中转"); + e.feed(body); + e.finished(403); + failure(rx) + }; + + let nothing: [&[u8]; 3] = [ + b"", + b" \n", + b"

403 Forbidden

", + ]; + for body in nothing { + let m = said(&mut rx, body); + assert_eq!(m.code, "gw.upstream.status", "{m:?}"); + assert_eq!(m.text, "Upstream `中转` answered 403."); + assert_eq!(m.arg("status"), "403"); + } + let m = said(&mut rx, b"Forbidden\n"); + assert_eq!( + (m.code.as_str(), m.arg("message")), + ("gw.upstream.status_message", "Forbidden") + ); + let m = said(&mut rx, br#"{"detail":"Not Found"}"#); + assert_eq!(m.arg("message"), r#"{"detail":"Not Found"}"#); + // 太长的只留开头 + let long = format!(r#"{{"error":{{"message":"{}"}}}}"#, "很".repeat(2_000)); + assert_eq!( + said(&mut rx, long.as_bytes()) + .arg("message") + .chars() + .count(), + SAID_MAX + ); + } + + /// 上游的原话和存下来的正文同一套打码:它回显了请求里的密钥,记录里那一句也不能有 + #[test] + fn what_the_upstream_said_is_masked_like_the_stored_answer() { + const KEY: &str = "sk-ant-api03-USERSOWNKEYAAAAAAAAAAAAAA"; + let bus = tw_observe::EventBus::new(); + let mut rx = bus.subscribe(); + let mut e = refused(&bus, 401, "中转"); + e.feed( + format!(r#"{{"type":"error","error":{{"type":"authentication_error","message":"invalid x-api-key {KEY}"}}}}"#) + .as_bytes(), + ); + e.finished(401); + let m = failure(&mut rx); + assert_eq!(m.code, "gw.upstream.status_message"); + assert!(m.arg("message").starts_with("invalid x-api-key "), "{m:?}"); + assert!( + !m.text.contains(KEY) && !m.arg("message").contains(KEY), + "{m:?}" + ); + } + + /// 网关替上游说了话的(Bedrock 拒绝凭证),原因就是那一句,不读正文 + #[test] + fn when_the_gateway_spoke_for_the_upstream_its_words_are_the_reason() { + let bus = tw_observe::EventBus::new(); + let mut rx = bus.subscribe(); + let ours = msg!( + "gw.upstream.bedrock_refused_unnamed", upstream = "bedrock", status = 403u16 => + "AWS refused the credential of upstream `{upstream}` (HTTP {status})." + ); + let mut e = Ending::new(bus.clone(), 7, MODEL.into(), Instant::now(), 1_000, None); + e.responded(403); + e.refused(ir::Dialect::Bedrock, "bedrock", Some(ours.clone())); + e.feed(br#"{"message":"[ThinkWatch] AWS refused the credential"}"#); + e.finished(403); + assert_eq!(failure(&mut rx), ours); + } + + /// 错误还没交完客户端就走了:**照旧是取消**,带着那个状态码 —— 取消怎么数不因为 + /// 回的是错误而变 + #[test] + fn a_client_that_leaves_during_an_error_answer_cancelled_it() { + let bus = tw_observe::EventBus::new(); + let mut rx = bus.subscribe(); + let mut e = refused(&bus, 400, "官方"); + e.feed(br#"{"type":"error","#); + drop(e); + assert!( + matches!( + drain(&mut rx).as_slice(), + [Event::RequestCancelled { + status: Some(400), + .. + }] + ), + "取消被报成了别的" + ); + } + /// 不是流的(非流式、错误响应)不看:那些不经过这里认错误 #[test] fn only_a_stream_is_watched() { diff --git a/crates/tw-gateway/src/server/pipeline/hop.rs b/crates/tw-gateway/src/server/pipeline/hop.rs index 7d8456ba..80ce9978 100644 --- a/crates/tw-gateway/src/server/pipeline/hop.rs +++ b/crates/tw-gateway/src/server/pipeline/hop.rs @@ -26,6 +26,9 @@ pub(super) struct Served<'a> { /// 成功那一跳的转换。**必须是成功那一次的** —— 故障转移从 Anthropic 上游 /// 切到 OpenAI 上游时,两跳转成的格式不一样;直通时是 None pub(super) session: Option, + /// 交出去的是网关替上游说的一句话,不是它的原话(Bedrock 拒绝凭证,见 + /// [`bedrock_refusal`]):这个请求失败的原因就是这一句。别的都是 None + pub(super) refusal: Option, } /// 这个请求的着落。 @@ -381,10 +384,11 @@ pub(super) async fn try_upstreams<'a>( }; note_health(&state.bus, &state.health, &provider.name, change); // Bedrock 拒绝凭证时的原话会点名账号和 IAM 身份:换成我们自己的话再交出去 - let r = if provider.is_bedrock() && matches!(status, 401 | 403) { - bedrock_refusal(provider, r).await + let (r, refusal) = if provider.is_bedrock() && matches!(status, 401 | 403) { + let (r, ours) = bedrock_refusal(provider, r).await; + (r, Some(ours)) } else { - r + (r, None) }; chain.push(hop( &provider.name, @@ -398,6 +402,7 @@ pub(super) async fn try_upstreams<'a>( provider, ledger, session: out.session, + refusal, }); break; } @@ -494,6 +499,7 @@ pub(super) async fn try_upstreams<'a>( provider, ledger, session: out.session, + refusal: None, }); break; } @@ -1149,6 +1155,7 @@ async fn send( } /// Bedrock 拒绝了凭证(401/403):状态码和响应头照原样,正文换成我们自己的一句话。 +/// 交回换好的响应和那句话 —— 请求记录里失败的原因也是它。 /// /// **AWS 的原话会点名账号 ID 和 IAM 身份**(`User: arn:aws:iam::…` is not authorized…), /// 那不能交给客户端,也不能进请求记录。能说的是异常名:凭证被拒、过期、没有权限,是 @@ -1156,7 +1163,7 @@ async fn send( async fn bedrock_refusal( provider: &tw_config::Provider, r: reqwest::Response, -) -> reqwest::Response { +) -> (reqwest::Response, tw_types::Msg) { let status = r.status(); let mut headers = r.headers().clone(); let named = headers @@ -1171,35 +1178,37 @@ async fn bedrock_refusal( (Some("ExpiredTokenException"), Some(profile)) => msg!( "gw.upstream.aws_profile_expired", upstream = provider.name.clone(), profile = profile => - "[ThinkWatch] The temporary AWS credential of upstream `{upstream}` has expired. \ - Refresh AWS profile `{profile}`; the next request reads it again." + "The temporary AWS credential of upstream `{upstream}` has expired. Refresh AWS \ + profile `{profile}`; the next request reads it again." ), (Some("ExpiredTokenException"), None) => msg!( "gw.upstream.aws_token_expired", upstream = provider.name.clone() => - "[ThinkWatch] The temporary AWS credential of upstream `{upstream}` has expired. \ - Replace its session token and the access keys that came with it." + "The temporary AWS credential of upstream `{upstream}` has expired. Replace its \ + session token and the access keys that came with it." ), (Some(kind), _) => msg!( "gw.upstream.bedrock_refused", upstream = provider.name.clone(), status = status.as_u16(), kind = kind => - "[ThinkWatch] AWS refused the credential of upstream `{upstream}` (HTTP {status}, \ - {kind}). Check that the credential is valid and may use this model. AWS's own \ - message names the account, so it is not passed on." + "AWS refused the credential of upstream `{upstream}` (HTTP {status}, {kind}). Check \ + that the credential is valid and may use this model. AWS's own message names the \ + account, so it is not passed on." ), // AWS 没说是哪种异常:另一句话,不拿一个英文词组去填 `{kind}` —— 那一格在译文里 // 就是一段没翻译的英文 (None, _) => msg!( "gw.upstream.bedrock_refused_unnamed", upstream = provider.name.clone(), status = status.as_u16() => - "[ThinkWatch] AWS refused the credential of upstream `{upstream}` (HTTP {status}). \ - Check that the credential is valid and may use this model. AWS's own message names \ - the account, so it is not passed on." + "AWS refused the credential of upstream `{upstream}` (HTTP {status}). Check that the \ + credential is valid and may use this model. AWS's own message names the account, \ + so it is not passed on." ), }; for h in ["content-length", "content-type", "transfer-encoding"] { headers.remove(h); } - let body = serde_json::json!({ "message": text.text }).to_string(); + // 前缀只在交给客户端的这一份上(和网关自己的错误一样,见 `crate::error`):记录里的 + // 那一句是这次请求失败的原因,界面按码翻译 + let body = serde_json::json!({ "message": format!("[ThinkWatch] {}", text.text) }).to_string(); let mut resp = http::Response::new(body); *resp.status_mut() = status; *resp.headers_mut() = headers; @@ -1207,7 +1216,7 @@ async fn bedrock_refusal( http::header::CONTENT_TYPE, http::HeaderValue::from_static("application/json"), ); - reqwest::Response::from(resp) + (reqwest::Response::from(resp), text) } /// 读错误响应的开头(最多 [`crate::failure::BODY_PEEK`]),交回读到的和一个照旧能从头 diff --git a/crates/tw-gateway/src/server/pipeline/relay.rs b/crates/tw-gateway/src/server/pipeline/relay.rs index c66b05fb..50e6cca5 100644 --- a/crates/tw-gateway/src/server/pipeline/relay.rs +++ b/crates/tw-gateway/src/server/pipeline/relay.rs @@ -37,6 +37,7 @@ pub(super) fn respond( provider, ledger, session, + refusal, } = served; let status = StatusCode::from_u16(upstream.status().as_u16()).unwrap_or(StatusCode::BAD_GATEWAY); @@ -114,6 +115,11 @@ pub(super) fn respond( if generates && plan.is_sse && status.is_success() { ending.streaming(upstream_dialect, &provider.name); } + // 回的不是 2xx:原样交给客户端,**这个请求照样是失败的** —— 客户端拿到的是上游的 + // 错误,不是回答。原因在错误正文里,交完时读(见 `Ending::refused`) + if !status.is_success() { + ending.refused(upstream_dialect, &provider.name, refusal); + } let mut relay = Relay::new( state, rt, diff --git a/crates/tw-gateway/tests/bedrock.rs b/crates/tw-gateway/tests/bedrock.rs index 0e6e051f..7fdde388 100644 --- a/crates/tw-gateway/tests/bedrock.rs +++ b/crates/tw-gateway/tests/bedrock.rs @@ -424,7 +424,7 @@ async fn a_refused_credential_does_not_pass_on_what_aws_said_about_the_account() .unwrap() })) .await; - let (gw, _, _) = gateway(with_keys(up)).await; + let (gw, _, mut rx) = gateway(with_keys(up)).await; let (status, body) = post( gw, "/v1/messages", @@ -437,6 +437,20 @@ async fn a_refused_credential_does_not_pass_on_what_aws_said_about_the_account() assert!(!body.contains("alice"), "{body}"); assert!(body.contains("AccessDeniedException"), "{body}"); assert!(body.contains("[ThinkWatch]"), "{body}"); + // 请求记录里失败的原因也是这一句(不带给客户端看的前缀),一样不点名账号 + match ending(&mut rx).await { + tw_api::Event::RequestFailed { + message, source, .. + } => { + assert_eq!(source, tw_api::FailureSource::Upstream); + assert_eq!(message.code, "gw.upstream.bedrock_refused", "{message:?}"); + assert_eq!(message.arg("kind"), "AccessDeniedException"); + assert_eq!(message.arg("status"), "403"); + assert!(!message.text.contains("123456789012"), "{message:?}"); + assert!(!message.text.starts_with("[ThinkWatch]"), "{message:?}"); + } + other => panic!("{other:?}"), + } } /// AWS 没说是哪种异常:换一句话说,不在句子里填一个英文词组(译文里那一格会是英文) diff --git a/crates/tw-gateway/tests/endings.rs b/crates/tw-gateway/tests/endings.rs index 33325494..c24aa1a3 100644 --- a/crates/tw-gateway/tests/endings.rs +++ b/crates/tw-gateway/tests/endings.rs @@ -606,6 +606,52 @@ async fn a_request_every_upstream_refused_is_failed_once_for_the_reason_it_was_r assert_eq!(model_of(&got[0]), MODEL); } +/// 上游回了 4xx(请求本身的问题,换一家也一样被拒):**原话原样交给客户端,结局是一次 +/// 失败**,来自 `upstream`,原因是上游在错误正文里说的那句话。以前报的是结束 —— 那一行 +/// 在流量、概览、会话里都像是成功的,对话里只剩用户的那句话 +#[tokio::test] +async fn an_error_answer_passed_on_to_the_client_is_failed_once_in_the_upstreams_words() { + const SAID: &str = r#"{"type":"error","error":{"type":"invalid_request_error","message":"prompt is too long: 212000 tokens > 200000 maximum"}}"#; + let up = listen(Router::new().fallback(any(|| async { + axum::response::Response::builder() + .status(400) + .header("content-type", "application/json") + .body(axum::body::Body::from(SAID)) + .unwrap() + }))) + .await; + let (gw, mut events) = serve(cfg(provider(up))).await; + let r = post(gw).send().await.unwrap(); + assert_eq!(r.status(), 400); + // 那是上游说的话,不是网关的错误 + assert!(r.headers().get("x-thinkwatch-error").is_none()); + assert_eq!(r.text().await.unwrap(), SAID); + + let got = endings(&mut events).await; + assert_eq!(got.len(), 1, "该恰好有一个结局:{got:?}"); + match &got[0] { + Event::RequestFailed { + source, + message, + bytes, + .. + } => { + assert_eq!(source, "upstream"); + assert_eq!(message.code, "gw.upstream.status_message"); + assert_eq!(message.arg("upstream"), "up"); + assert_eq!(message.arg("status"), "400"); + assert_eq!( + message.arg("message"), + "prompt is too long: 212000 tokens > 200000 maximum" + ); + // 响应头到了:收到的字节是有的 + assert_eq!(*bytes, Some(SAID.len() as u64)); + } + other => panic!("该是一次失败,实际 {other:?}"), + } + assert_eq!(model_of(&got[0]), MODEL); +} + // ---------------------------------------------------------------- WebSocket /// 一次 Codex 会话结束了。**以前 WS 这条路只有开始、没有结局**,每一条连接 diff --git a/crates/tw-store/src/db.rs b/crates/tw-store/src/db.rs index 9f247eb3..979196a3 100644 --- a/crates/tw-store/src/db.rs +++ b/crates/tw-store/src/db.rs @@ -472,7 +472,10 @@ pub struct TurnRow { pub cache_read_tokens: Option, pub cost_micros: Option, pub duration_ms: Option, - /// 失败的原因,带着码 + /// 上游回的状态码。没走到上游的是 None + pub status: Option, + /// 失败的原因,带着码。上游回了错误、原样交给客户端的也有(那时 `status` 是那个 + /// 状态码),见 `tw_api::TurnView::error` pub error: Option, pub cancelled: bool, /// 这一轮的金额是估算。**瀑布图上要带记号** @@ -552,7 +555,7 @@ impl Db { let mut st = self.conn.prepare( "SELECT id, at_ms, model, provider, input_tokens, output_tokens, cache_read_tokens, cost_micros, duration_ms, error, cancelled, - cost_estimated, billing, error_code, error_args + cost_estimated, billing, error_code, error_args, status FROM requests WHERE session = ?1 AND local = 0 ORDER BY at_ms, id", )?; let rows = st.query_map([session], |r| { @@ -566,6 +569,7 @@ impl Db { cache_read_tokens: r.get(6)?, cost_micros: r.get(7)?, duration_ms: r.get(8)?, + status: r.get(15)?, error: error_from(r)?, cancelled: r.get::<_, i64>(10)? != 0, cost_estimated: r.get::<_, i64>(11)? != 0, diff --git a/crates/tw-store/src/recorder.rs b/crates/tw-store/src/recorder.rs index 25f3514e..6fc1a336 100644 --- a/crates/tw-store/src/recorder.rs +++ b/crates/tw-store/src/recorder.rs @@ -2453,6 +2453,107 @@ mod failure_tests { assert_eq!(row.duration_ms, Some(20_000)); } + /// 上游回了 4xx、原样交给了客户端(网关报的是失败,见 `tw_gateway::ending`):这一行 + /// 是失败,状态码是上游回的那个。**会话里那一轮也一样**:状态码和原因都在,数进会话和 + /// 概览的失败,不数进「没有用量」(那种响应不计费)。同一次会话里成功的、取消的那两轮 + /// 照旧 + #[test] + fn an_error_answer_is_a_failed_turn_with_the_status_the_upstream_answered() { + let (_d, mut r) = rec(); + let start = |id| { + let mut ev = started(id, "claude-sonnet-4-5"); + if let Event::RequestStarted { session, .. } = &mut ev { + *session = Some("s".into()); + } + ev + }; + let headers = |id, status| Event::RequestHeaders { + id, + status, + ttfb_ms: 300, + }; + r.on_event(&start(1)); + r.on_event(&headers(1, 200)); + r.on_event(&finished( + 1, + Some(UsageView { + input: 1_000, + output: 10, + ..Default::default() + }), + )); + let said = tw_api::Msg { + code: "gw.upstream.status_message".into(), + args: [ + ("upstream", "官方"), + ("status", "400"), + ("message", "prompt is too long"), + ] + .into_iter() + .map(|(k, v)| (k.to_string(), v.to_string())) + .collect(), + text: "Upstream `官方` answered 400: prompt is too long".into(), + }; + r.on_event(&start(2)); + r.on_event(&headers(2, 400)); + r.on_event(&Event::RequestFailed { + id: 2, + model: String::new(), + source: tw_api::FailureSource::Upstream, + message: said.clone(), + bytes: Some(120), + duration_ms: Some(320), + usage: None, + answered_model: None, + }); + r.on_event(&start(3)); + r.on_event(&headers(3, 200)); + r.on_event(&Event::RequestCancelled { + id: 3, + model: String::new(), + status: Some(200), + bytes: 40, + duration_ms: 900, + usage: None, + answered_model: None, + }); + + let row = r.db().get(2).unwrap().unwrap(); + assert_eq!(row.status, Some(400), "状态码来自响应头那个事件"); + assert_eq!(row.error.as_ref(), Some(&said)); + assert!(!row.cancelled); + assert_eq!(row.cost_micros, None); + + let turns = r.db().turns("s").unwrap(); + let seen: Vec<_> = turns + .iter() + .map(|t| { + ( + t.id, + t.status, + t.error.as_ref().map(|e| e.code.as_str()), + t.cancelled, + ) + }) + .collect(); + assert_eq!( + seen, + [ + (1, Some(200), None, false), + (2, Some(400), Some("gw.upstream.status_message"), false), + (3, Some(200), None, true), + ] + ); + let s = &r.db().sessions(None, 10).unwrap()[0]; + assert_eq!((s.turns, s.errors, s.no_usage_turns), (3, 1, 1), "{s:?}"); + let sum = r.db().summary(0, i64::MAX).unwrap(); + assert_eq!( + (sum.requests, sum.failed, sum.no_usage_requests), + (3, 1, 1), + "{sum:?}" + ); + } + /// 算出来的价钱照样报回总线,**标着估算**。 #[test] fn the_price_of_a_failure_goes_back_onto_the_bus_as_an_estimate() {