diff --git a/crates/tw-api/msg-codes.txt b/crates/tw-api/msg-codes.txt index d4146dd..872eaad 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 494a8ed..48186b3 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 f2f2814..0752ba9 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 93ac21a..24f871e 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 0000000..57ee4e5 --- /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 b2e097a..d077517 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 efd4fa8..e0ef4eb 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 7d8456b..80ce997 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 c66b05f..50e6cca 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 0e6e051..7fdde38 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 3332549..c24aa1a 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 9f247eb..979196a 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 25f3514..6fc1a33 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() {