diff --git a/apps/aether-gateway/src/handlers/proxy/websocket/responses/frame.rs b/apps/aether-gateway/src/handlers/proxy/websocket/responses/frame.rs index 65e568608..d756bd652 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/frame.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/frame.rs @@ -148,14 +148,69 @@ fn event_is_started(event: &Value) -> bool { ) } +/// `response.incomplete` 的合法终态 reason 白名单。 +/// +/// 这些 reason 表示上游按规则正常结束了本轮响应(写满 `max_output_tokens`、 +/// 命中内容过滤、按工具调用截断),标准流解析里它们会变成 `length` / +/// `content_filter` / `tool_calls` 这类正常 finish,和 +/// `openai_responses_incomplete_finish_reason` 的既有映射保持一致,因此不能 +/// 当成 provider failure 记账。 +const LEGITIMATE_RESPONSES_INCOMPLETE_REASONS: [&str; 5] = [ + "max_output_tokens", + "max_tokens", + "content_filter", + "tool_calls", + "function_call", +]; + +/// 读取 `response.incomplete` 携带的 `incomplete_details.reason`。 +/// +/// 标准位置是 `response.incomplete_details.reason`;批量封装偶尔把 +/// `incomplete_details` 直接放在事件顶层,两处都要看,否则合法终态会被漏判。 +fn responses_incomplete_reason(event: &Value) -> Option<&str> { + [ + event.pointer("/response/incomplete_details/reason"), + event.pointer("/incomplete_details/reason"), + ] + .into_iter() + .flatten() + .filter_map(Value::as_str) + .map(str::trim) + .find(|reason| !reason.is_empty()) +} + +/// 判断一个 `response.incomplete` 是否是合法终态。 +/// +/// reason 缺失或不在白名单内(例如 `error`、`server_error`)时继续按 +/// provider failure 处理:这类 incomplete 说明上游确实没能正常收尾,仍应扣 +/// 供应商健康分。 +fn responses_incomplete_is_legitimate_terminal(event: &Value) -> bool { + responses_incomplete_reason(event).is_some_and(|reason| { + LEGITIMATE_RESPONSES_INCOMPLETE_REASONS + .iter() + .any(|candidate| reason.eq_ignore_ascii_case(candidate)) + }) +} + fn terminal_for_event(event: &Value) -> Option { match event_type_of(event).unwrap_or_default() { "response.completed" => Some(ResponsesWebSocketFrameTerminal { status_code: websocket_event_status_code(event, 200), cancelled: false, }), + // 合法 incomplete(例如写满 max_output_tokens)是正常终态,默认按 200 + // 记账,不再一律当 502 provider failure;reason 缺失或未知时保留原来的 + // 502 默认值。显式 `status_code` 和 error code 映射仍然优先于默认值, + // 所以带 `rate_limit_exceeded` 的 incomplete 依旧是 429。 "response.incomplete" => Some(ResponsesWebSocketFrameTerminal { - status_code: websocket_event_status_code(event, 502), + status_code: websocket_event_status_code( + event, + if responses_incomplete_is_legitimate_terminal(event) { + 200 + } else { + 502 + }, + ), cancelled: false, }), "response.cancelled" => Some(ResponsesWebSocketFrameTerminal { @@ -281,6 +336,95 @@ mod tests { assert_eq!(failed.status(), Some(429)); } + #[test] + fn a_legitimate_incomplete_is_a_terminal_but_not_a_provider_failure() { + for reason in [ + "max_output_tokens", + "max_tokens", + "content_filter", + "tool_calls", + "function_call", + "MAX_OUTPUT_TOKENS", + ] { + let raw = format!( + r#"{{"type":"response.incomplete","response":{{"status":"incomplete","incomplete_details":{{"reason":"{reason}"}}}}}}"# + ); + let frame = ParsedResponsesWebSocketFrame::parse(&raw).expect("valid frame"); + + assert!(frame.is_terminal(), "{reason} should end the turn"); + assert_eq!( + frame + .terminal() + .map(|terminal| (terminal.status_code, terminal.cancelled)), + Some((200, false)), + "{reason} is a legitimate terminal result, not a 502 provider failure" + ); + } + } + + #[test] + fn a_top_level_incomplete_details_reason_is_also_honored() { + let frame = ParsedResponsesWebSocketFrame::parse( + r#"{"type":"response.incomplete","incomplete_details":{"reason":"max_output_tokens"}}"#, + ) + .expect("valid frame"); + + assert_eq!(frame.status(), Some(200)); + } + + #[test] + fn an_incomplete_without_a_legitimate_reason_stays_a_provider_failure() { + for raw in [ + r#"{"type":"response.incomplete"}"#, + r#"{"type":"response.incomplete","response":{"incomplete_details":null}}"#, + r#"{"type":"response.incomplete","response":{"incomplete_details":{"reason":""}}}"#, + r#"{"type":"response.incomplete","response":{"incomplete_details":{"reason":"error"}}}"#, + r#"{"type":"response.incomplete","response":{"incomplete_details":{"reason":"server_error"}}}"#, + ] { + let frame = ParsedResponsesWebSocketFrame::parse(raw).expect("valid frame"); + + assert_eq!( + frame.status(), + Some(502), + "an incomplete without a known-good reason must stay a provider failure: {raw}" + ); + } + } + + #[test] + fn a_legitimate_incomplete_still_respects_an_explicit_provider_status() { + let explicit = ParsedResponsesWebSocketFrame::parse( + r#"{"type":"response.incomplete","status_code":503,"response":{"incomplete_details":{"reason":"max_output_tokens"}}}"#, + ) + .expect("valid frame"); + assert_eq!(explicit.status(), Some(503)); + + let quota = ParsedResponsesWebSocketFrame::parse( + r#"{"type":"response.incomplete","response":{"error":{"code":"rate_limit_exceeded"},"incomplete_details":{"reason":"max_output_tokens"}}}"#, + ) + .expect("valid frame"); + assert_eq!(quota.status(), Some(429)); + } + + #[test] + fn a_legitimate_incomplete_batched_inside_a_chunks_envelope_is_not_a_failure() { + let frame = ParsedResponsesWebSocketFrame::parse( + r#"{"chunks":[{"type":"response.output_text.delta","delta":"hi"},{"type":"response.incomplete","response":{"incomplete_details":{"reason":"max_output_tokens"},"usage":{"total_tokens":9}}}]}"#, + ) + .expect("valid frame"); + + assert!(frame.is_chunked()); + assert!(frame.is_terminal()); + assert_eq!(frame.status(), Some(200)); + assert_eq!(frame.event_type(), Some("response.incomplete")); + assert_eq!( + frame.terminal_event().and_then(|event| event + .pointer("/response/usage/total_tokens") + .and_then(serde_json::Value::as_u64)), + Some(9) + ); + } + #[test] fn detects_a_terminal_batched_inside_a_chunks_envelope() { let frame = ParsedResponsesWebSocketFrame::parse( diff --git a/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn.rs b/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn.rs index 499fe3946..dbd31644d 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn.rs @@ -232,6 +232,52 @@ impl ResponsesWebSocketTurnOutcome { } } +/// 一轮 turn 结束后要投射给供应商/密钥池的效果。 +/// +/// 每个分支都会释放 pool key lease:`ProviderFailure` 由 `PoolError` 释放, +/// `ProviderSuccess` 由 `PoolSuccessStream` 释放,其余情况直接释放。少一条 +/// 分支就会把 lease 挂到 TTL 过期,等于短时间占死一把 key。 +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum ResponsesWebSocketTurnEffect { + /// 既不投射成功也不投射失败,只把 lease 还回去。 + ReleasePoolKeyLease, + ProviderFailure, + ProviderSuccess, +} + +#[cfg(test)] +impl ResponsesWebSocketTurnEffect { + /// 把「每个分支都必须释放 lease」这条不变量显式化,便于测试锁住 + /// 「没进任何分支导致 lease 泄漏」这类回归。 + const fn releases_pool_key_lease(self) -> bool { + match self { + Self::ReleasePoolKeyLease | Self::ProviderFailure | Self::ProviderSuccess => true, + } + } +} + +/// 判定一轮 turn 结束后要投射的效果。 +/// +/// 关键分支是「记账层判成 failed,但这一轮没有投射供应商失败」:例如合法的 +/// `response.incomplete`(写满 max_output_tokens)。共享 usage 判定目前仍会 +/// 把这类终态记成失败,但供应商本身工作正常,既不该扣健康分,也不能因为落 +/// 不到任何分支而漏掉 lease 释放。 +const fn classify_responses_websocket_turn_effect( + cancelled: bool, + projects_provider_failure: bool, + failed: bool, +) -> ResponsesWebSocketTurnEffect { + if cancelled { + ResponsesWebSocketTurnEffect::ReleasePoolKeyLease + } else if projects_provider_failure { + ResponsesWebSocketTurnEffect::ProviderFailure + } else if failed { + ResponsesWebSocketTurnEffect::ReleasePoolKeyLease + } else { + ResponsesWebSocketTurnEffect::ProviderSuccess + } +} + pub(super) struct ResponsesWebSocketTurn { plan: ExecutionPlan, trace_id: String, @@ -816,26 +862,32 @@ impl ResponsesWebSocketTurn { || outcome.forced_error().is_some() || summary.parser_error.is_some() || missing_terminal); + let provider_effect = + classify_responses_websocket_turn_effect(cancelled, projects_provider_failure, failed); let effects_completed = await_websocket_lifecycle_stage(&self.trace_id, "provider_effects", async { - if cancelled { - release_local_pool_key_lease(state, effect_context).await; - } else if projects_provider_failure { - let response_text = terminal_error_body - .as_deref() - .or(summary.parser_error.as_deref()) - .unwrap_or(outcome_reason.as_str()); - let mut effect = LocalStreamFailureEffect::new( - status_code, - &payload.headers, - Some(response_text), - ); - if outcome.stream_timeout() { - effect = effect.with_stream_timeout(); + match provider_effect { + ResponsesWebSocketTurnEffect::ReleasePoolKeyLease => { + release_local_pool_key_lease(state, effect_context).await; + } + ResponsesWebSocketTurnEffect::ProviderFailure => { + let response_text = terminal_error_body + .as_deref() + .or(summary.parser_error.as_deref()) + .unwrap_or(outcome_reason.as_str()); + let mut effect = LocalStreamFailureEffect::new( + status_code, + &payload.headers, + Some(response_text), + ); + if outcome.stream_timeout() { + effect = effect.with_stream_timeout(); + } + apply_local_stream_failure_effects(state, effect_context, effect).await; + } + ResponsesWebSocketTurnEffect::ProviderSuccess => { + apply_local_stream_success_effects(state, effect_context, &payload).await; } - apply_local_stream_failure_effects(state, effect_context, effect).await; - } else if !failed { - apply_local_stream_success_effects(state, effect_context, &payload).await; } }) .await @@ -1174,10 +1226,10 @@ mod tests { use super::super::frame::ParsedResponsesWebSocketFrame; use super::{ - prepare_websocket_report_context, provider_terminal_outcome, - resolve_responses_websocket_turn_timeouts, websocket_event_as_sse_line, - ResponsesWebSocketTurnDeadline, ResponsesWebSocketTurnOutcome, - ResponsesWebSocketTurnTimeoutPhase, + classify_responses_websocket_turn_effect, prepare_websocket_report_context, + provider_terminal_outcome, resolve_responses_websocket_turn_timeouts, + websocket_event_as_sse_line, ResponsesWebSocketTurnDeadline, ResponsesWebSocketTurnEffect, + ResponsesWebSocketTurnOutcome, ResponsesWebSocketTurnTimeoutPhase, }; #[test] @@ -1296,6 +1348,148 @@ mod tests { assert_eq!(usage.dimensions.get("total_tokens"), Some(&json!(8))); } + #[test] + fn a_legitimate_incomplete_is_a_successful_provider_terminal_that_keeps_its_usage() { + // 写满 max_output_tokens 的 incomplete 是合法终态:状态码不再是 502, + // usage 观测器照样能看到 finish 和 token,记账层不该把它当解析失败。 + let event = json!({ + "type": "response.incomplete", + "response": { + "id": "resp_ws_incomplete_123", + "model": "gpt-5.6", + "status": "incomplete", + "incomplete_details": {"reason": "max_output_tokens"}, + "output": [], + "usage": { + "input_tokens": 4, + "output_tokens": 7, + "total_tokens": 11 + } + } + }); + let raw = serde_json::to_string(&event).expect("event should serialize"); + let frame = ParsedResponsesWebSocketFrame::parse(&raw).expect("event should parse"); + let outcome = provider_terminal_outcome(&frame); + assert_eq!( + outcome, + Some(ResponsesWebSocketTurnOutcome::ProviderTerminal { + status_code: 200, + cancelled: false + }) + ); + let outcome = outcome.expect("incomplete should end the turn"); + assert!(!outcome.cancelled()); + assert!(outcome.forced_error().is_none()); + assert!(!outcome.stream_timeout()); + + let report_context = json!({ + "provider_api_format": "openai:responses", + "client_api_format": "openai:responses", + }); + let mut observer = StreamingStandardTerminalObserver::default(); + observer + .push_line(&report_context, websocket_event_as_sse_line(&event)) + .expect("a legitimate incomplete must be accepted by the usage observer"); + let summary = observer + .finish(&report_context) + .expect("terminal observer should finish") + .expect("terminal observer should produce a summary"); + assert!(summary.observed_finish); + assert_eq!(summary.finish_reason.as_deref(), Some("length")); + assert!(summary.parser_error.is_none()); + let usage = summary + .standardized_usage + .expect("incomplete usage must reach the terminal summary"); + assert_eq!(usage.input_tokens, 4); + assert_eq!(usage.output_tokens, 7); + + // finalize() 用这些事实决定是否投射供应商失败:合法 incomplete 必须 + // 全部落在“非失败”一侧。 + let missing_terminal = !outcome.cancelled() && !summary.observed_finish; + let projects_provider_failure = !outcome.cancelled() + && (outcome.status_code() >= 400 + || outcome.forced_error().is_some() + || summary.parser_error.is_some() + || missing_terminal); + assert!(!missing_terminal); + assert!(!projects_provider_failure); + } + + #[test] + fn an_incomplete_without_a_legitimate_reason_still_projects_a_provider_failure() { + let raw = r#"{"type":"response.incomplete","response":{"incomplete_details":{"reason":"error"}}}"#; + let frame = ParsedResponsesWebSocketFrame::parse(raw).expect("event should parse"); + + assert_eq!( + provider_terminal_outcome(&frame), + Some(ResponsesWebSocketTurnOutcome::ProviderTerminal { + status_code: 502, + cancelled: false + }) + ); + } + + #[test] + fn a_legitimate_incomplete_still_releases_the_pool_key_lease() { + // 共享 usage 判定目前仍把 response.incomplete 记成终态失败,于是会出现 + // failed=true 而 projects_provider_failure=false 的组合。这种组合必须 + // 明确落到“只释放 lease”的分支,否则 lease 会挂到 TTL 过期。 + let effect = classify_responses_websocket_turn_effect(false, false, true); + + assert_eq!(effect, ResponsesWebSocketTurnEffect::ReleasePoolKeyLease); + assert!(effect.releases_pool_key_lease()); + } + + #[test] + fn every_turn_effect_releases_the_pool_key_lease() { + for (cancelled, projects_provider_failure, failed, expected) in [ + ( + true, + false, + false, + ResponsesWebSocketTurnEffect::ReleasePoolKeyLease, + ), + ( + true, + true, + true, + ResponsesWebSocketTurnEffect::ReleasePoolKeyLease, + ), + ( + false, + true, + true, + ResponsesWebSocketTurnEffect::ProviderFailure, + ), + ( + false, + false, + true, + ResponsesWebSocketTurnEffect::ReleasePoolKeyLease, + ), + ( + false, + false, + false, + ResponsesWebSocketTurnEffect::ProviderSuccess, + ), + ] { + let effect = classify_responses_websocket_turn_effect( + cancelled, + projects_provider_failure, + failed, + ); + assert_eq!( + effect, expected, + "cancelled={cancelled} projects_provider_failure={projects_provider_failure} failed={failed}" + ); + assert!( + effect.releases_pool_key_lease(), + "every effect branch must release the pool key lease" + ); + } + } + #[test] fn error_event_uses_the_top_level_status_code() { let event = json!({