fix(ws): treat max_output_tokens incomplete as legitimate terminal

This commit is contained in:
AAEE86
2026-08-17 14:52:33 +08:00
committed by ZheFox
parent 1353d76e07
commit f70ae68273
2 changed files with 360 additions and 22 deletions
@@ -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<ResponsesWebSocketFrameTerminal> {
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(
@@ -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!({