diff --git a/apps/aether-gateway/src/handlers/proxy/websocket/responses/client.rs b/apps/aether-gateway/src/handlers/proxy/websocket/responses/client.rs index e3e447f24..3165965d8 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/client.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/client.rs @@ -10,7 +10,7 @@ use wreq::ws::message::Message as WreqWsMessage; use super::adapter::{resolve_responses_websocket_adapter, ResponsesWebSocketDrainDirective}; use super::lifecycle::{ await_pending_turn_finalization, queue_turn_finalization, - send_responses_websocket_turn_start_error, ActiveResponsesWebSocketTurn, + send_responses_websocket_turn_start_error, ActiveProviderAttempt, }; use super::quota::{mark_active_response_retry_unsafe, send_previous_response_not_found}; use super::redaction::redact_responses_websocket_client_event; @@ -369,7 +369,7 @@ pub(super) async fn forward_client_message( turn.set_provider_response_headers(bound.upstream_response_headers.clone()); bound.turn_state.begin( LogicalTurn::new(client_event.clone(), turn_index, logical_turn_id), - ActiveResponsesWebSocketTurn::new(state, turn), + ActiveProviderAttempt::new(state, turn), ); bound.next_turn_index = bound.next_turn_index.saturating_add(1); @@ -697,7 +697,7 @@ async fn forward_replanned_response_create( bound.body_normalization = normalization; bound.turn_state.begin( LogicalTurn::new(client_event.clone(), turn_index, logical_turn_id.clone()), - ActiveResponsesWebSocketTurn::new(state, turn), + ActiveProviderAttempt::new(state, turn), ); bound.next_turn_index = bound.next_turn_index.saturating_add(1); debug!( @@ -769,7 +769,7 @@ async fn forward_replanned_response_create( bound.binding_identity = replacement.binding_identity; bound.turn_state.begin( LogicalTurn::new(client_event, turn_index, logical_turn_id), - ActiveResponsesWebSocketTurn::new(state, turn), + ActiveProviderAttempt::new(state, turn), ); bound.next_turn_index = bound.next_turn_index.saturating_add(1); bound.upstream_response_headers = replacement.upstream_response_headers; diff --git a/apps/aether-gateway/src/handlers/proxy/websocket/responses/connection.rs b/apps/aether-gateway/src/handlers/proxy/websocket/responses/connection.rs index 8ed84d208..748da4f88 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/connection.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/connection.rs @@ -12,7 +12,7 @@ use super::client::{adapter_drain_ready, forward_client_message, RelayDispositio use super::frame::ParsedResponsesWebSocketFrame; use super::lifecycle::{ await_pending_adapter_observation, finalize_active_turn, queue_turn_finalization, - ActiveResponsesWebSocketTurn, + ActiveProviderAttempt, }; use super::quota::{ active_continuation_can_retry_from_full_input, detach_exhausted_upstream, @@ -25,7 +25,7 @@ use super::relay_policy::{ }; use super::state::BoundResponsesConnection; use super::turn::{ - ResponsesWebSocketTurn, ResponsesWebSocketTurnObservation, ResponsesWebSocketTurnOutcome, + ResponsesProviderAttempt, ResponsesWebSocketTurnObservation, ResponsesWebSocketTurnOutcome, }; use super::upstream::{close_bound_upstream, receive_optional_upstream}; use crate::handlers::proxy::websocket::ingress::WebSocketRequestContext; @@ -383,7 +383,7 @@ pub(super) async fn relay_bound_connection( let mut quota_relay_action = classify_quota_relay(quota_facts); if matches!(quota_relay_action, QuotaRelayAction::AttemptTransparentRetry) { // detach_attempt 保留 logical turn:重试是同一轮请求的下一个 attempt。 - let mut retry_turn = bound.turn_state.detach_attempt().map(ActiveResponsesWebSocketTurn::disarm); + let mut retry_turn = bound.turn_state.detach_attempt().map(ActiveProviderAttempt::disarm); if let Some(turn) = retry_turn.as_mut() { turn.release_admission().await; } @@ -404,7 +404,7 @@ pub(super) async fn relay_bound_connection( if let Some(turn) = retry_turn { let restored = bound .turn_state - .resume(ActiveResponsesWebSocketTurn::new(state, turn)); + .resume(ActiveProviderAttempt::new(state, turn)); debug_assert!( restored.is_ok(), "a failed transparent retry must be able to restore its attempt" @@ -433,7 +433,7 @@ pub(super) async fn relay_bound_connection( error_code = "previous_response_not_found", "gateway will ask the client to retry the continuation with complete input" ); - let mut turn = bound.turn_state.end().map(ActiveResponsesWebSocketTurn::disarm); + let mut turn = bound.turn_state.end().map(ActiveProviderAttempt::disarm); if let Some(active_turn) = turn.as_mut() { active_turn.release_admission().await; } diff --git a/apps/aether-gateway/src/handlers/proxy/websocket/responses/lifecycle.rs b/apps/aether-gateway/src/handlers/proxy/websocket/responses/lifecycle.rs index 15b41f417..ba8814382 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/lifecycle.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/lifecycle.rs @@ -11,7 +11,7 @@ use tokio::time::timeout; use super::state::BoundResponsesConnection; use super::turn::{ - spawn_responses_websocket_turn_finalization, ResponsesWebSocketTurn, + spawn_responses_websocket_turn_finalization, ResponsesProviderAttempt, ResponsesWebSocketTurnOutcome, }; use crate::handlers::proxy::websocket::session::{ @@ -37,13 +37,13 @@ macro_rules! warn { /// would otherwise be discarded with its usage row left `Pending`, its /// candidate row left `Streaming`, and its distributed pool key lease leaked /// until the lease expires. Mirrors the HTTP path's `DirectPassthroughFinalizer`. -pub(super) struct ActiveResponsesWebSocketTurn { - turn: Option, +pub(super) struct ActiveProviderAttempt { + turn: Option, state: AppState, } -impl ActiveResponsesWebSocketTurn { - pub(super) fn new(state: &AppState, turn: ResponsesWebSocketTurn) -> Self { +impl ActiveProviderAttempt { + pub(super) fn new(state: &AppState, turn: ResponsesProviderAttempt) -> Self { Self { turn: Some(turn), state: state.clone(), @@ -51,15 +51,15 @@ impl ActiveResponsesWebSocketTurn { } /// Hands the turn back to a caller that will finalize it explicitly. - pub(super) fn disarm(mut self) -> ResponsesWebSocketTurn { + pub(super) fn disarm(mut self) -> ResponsesProviderAttempt { self.turn .take() .expect("an armed active turn always holds its turn") } } -impl std::ops::Deref for ActiveResponsesWebSocketTurn { - type Target = ResponsesWebSocketTurn; +impl std::ops::Deref for ActiveProviderAttempt { + type Target = ResponsesProviderAttempt; fn deref(&self) -> &Self::Target { self.turn @@ -68,7 +68,7 @@ impl std::ops::Deref for ActiveResponsesWebSocketTurn { } } -impl std::ops::DerefMut for ActiveResponsesWebSocketTurn { +impl std::ops::DerefMut for ActiveProviderAttempt { fn deref_mut(&mut self) -> &mut Self::Target { self.turn .as_mut() @@ -76,7 +76,7 @@ impl std::ops::DerefMut for ActiveResponsesWebSocketTurn { } } -impl Drop for ActiveResponsesWebSocketTurn { +impl Drop for ActiveProviderAttempt { fn drop(&mut self) { let Some(turn) = self.turn.take() else { return; @@ -120,7 +120,7 @@ pub(super) async fn finalize_active_turn( pub(super) async fn queue_turn_finalization( bound: &mut BoundResponsesConnection, state: &AppState, - turn: ResponsesWebSocketTurn, + turn: ResponsesProviderAttempt, outcome: ResponsesWebSocketTurnOutcome, ) { await_pending_adapter_observation(bound).await; @@ -161,7 +161,7 @@ pub(super) async fn await_pending_adapter_observation(bound: &mut BoundResponses pub(super) async fn finalize_unbound_turn( state: AppState, - turn: ResponsesWebSocketTurn, + turn: ResponsesProviderAttempt, outcome: ResponsesWebSocketTurnOutcome, ) -> JoinHandle<()> { spawn_responses_websocket_turn_finalization(state, turn, outcome).await diff --git a/apps/aether-gateway/src/handlers/proxy/websocket/responses/mod.rs b/apps/aether-gateway/src/handlers/proxy/websocket/responses/mod.rs index d16929ed9..88672426f 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/mod.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/mod.rs @@ -19,6 +19,7 @@ mod redaction; mod relay_policy; mod request; mod session; +mod settlement; mod state; mod turn; mod turn_state; diff --git a/apps/aether-gateway/src/handlers/proxy/websocket/responses/quota.rs b/apps/aether-gateway/src/handlers/proxy/websocket/responses/quota.rs index e2501ff6e..c3d1ca449 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/quota.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/quota.rs @@ -10,7 +10,7 @@ use super::adapter::{ resolve_responses_websocket_adapter, ResponsesWebSocketDrainDirective, ResponsesWebSocketRebindSafety, }; -use super::lifecycle::{queue_turn_finalization, ActiveResponsesWebSocketTurn}; +use super::lifecycle::{queue_turn_finalization, ActiveProviderAttempt}; use super::request::{ build_planning_parts, planned_response_create_event, response_create_has_previous_response_id, }; @@ -304,7 +304,7 @@ pub(super) async fn retry_active_turn_after_quota_exhaustion( // pending usage 行、占着 candidate 和 pool key lease 的 attempt。 if let Err(orphan) = bound .turn_state - .resume(ActiveResponsesWebSocketTurn::new(state, turn)) + .resume(ActiveProviderAttempt::new(state, turn)) { drop(orphan); return false; diff --git a/apps/aether-gateway/src/handlers/proxy/websocket/responses/session.rs b/apps/aether-gateway/src/handlers/proxy/websocket/responses/session.rs index 990825c63..64f9cb201 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/session.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/session.rs @@ -20,7 +20,7 @@ use super::connection::relay_bound_connection; use super::lifecycle::{ await_pending_adapter_observation, await_pending_turn_finalization, await_turn_finalization_handle, finalize_unbound_turn, responses_websocket_turn_start_close, - send_responses_websocket_turn_start_error, ActiveResponsesWebSocketTurn, + send_responses_websocket_turn_start_error, ActiveProviderAttempt, }; use super::redaction::redact_responses_websocket_client_event; use super::request::{build_planning_parts, planned_response_create_event}; @@ -415,7 +415,7 @@ pub(super) async fn run_responses_websocket( first_turn.set_provider_response_headers(bound.upstream_response_headers.clone()); bound.turn_state.begin( LogicalTurn::new(first_event, 1, first_logical_turn_id), - ActiveResponsesWebSocketTurn::new(&state, first_turn), + ActiveProviderAttempt::new(&state, first_turn), ); relay_bound_connection( diff --git a/apps/aether-gateway/src/handlers/proxy/websocket/responses/settlement.rs b/apps/aether-gateway/src/handlers/proxy/websocket/responses/settlement.rs new file mode 100644 index 000000000..5c7ec89e8 --- /dev/null +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/settlement.rs @@ -0,0 +1,838 @@ +//! 一次 provider attempt 的结算判定。 +//! +//! 评审第 5 条:现状 `ResponsesWebSocketTurnOutcome` 一个枚举同时表达 +//! 「供应商这一轮怎么结束的」和「内容有没有完整交付给客户端」,`finalize()` +//! 再用 `outcome.cancelled()` 一个布尔驱动 billing、candidate 状态和供应商效果。 +//! 于是 provider 终态已经到达、只是最后一跳写客户端失败时,供应商事实会被 +//! 覆盖掉。这里把两件事拆成正交事实,并把结算动作收进一张可逐行测试的表。 + +use super::turn::ResponsesWebSocketTurnOutcome; + +/// 客户端取消/断开时对外记录的状态码。 +const CLIENT_CANCELLED_STATUS_CODE: u16 = 499; + +/// 流式超时状态码;现状只有它会额外投射 pool stream timeout 效果。 +const STREAM_TIMEOUT_STATUS_CODE: u16 = 504; + +/// provider 侧观察到的终态。 +/// +/// 形状刻意保持 transport 中立:HTTP 流式与 WS turn 的差异只在事实从哪来。 +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(super) enum AttemptProviderOutcome { + /// 观察到了供应商的终态事件。 + Terminal { + status_code: u16, + /// 供应商自己声明这一轮被取消(`response.cancelled`)。 + cancelled_by_provider: bool, + }, + /// 供应商没能给出终态:断链、超时、gateway 侧失败。 + Aborted { + status_code: u16, + reason: &'static str, + stream_timeout: bool, + }, +} + +impl AttemptProviderOutcome { + pub(super) const fn status_code(self) -> u16 { + match self { + Self::Terminal { status_code, .. } | Self::Aborted { status_code, .. } => status_code, + } + } + + pub(super) const fn cancelled_by_provider(self) -> bool { + matches!( + self, + Self::Terminal { + cancelled_by_provider: true, + .. + } + ) + } + + pub(super) const fn stream_timeout(self) -> bool { + matches!(self, Self::Aborted { + stream_timeout: true, + .. + }) + } + + pub(super) const fn is_terminal(self) -> bool { + matches!(self, Self::Terminal { .. }) + } +} + +/// 这一个 attempt 的内容是否完整交付给了客户端。 +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(super) enum AttemptClientDelivery { + Complete, + Aborted { reason: &'static str }, +} + +impl AttemptClientDelivery { + pub(super) const fn aborted_reason(self) -> Option<&'static str> { + match self { + Self::Complete => None, + Self::Aborted { reason } => Some(reason), + } + } + + pub(super) const fn is_aborted(self) -> bool { + matches!(self, Self::Aborted { .. }) + } +} + +/// 一次 attempt 结算时的两个正交事实。 +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(super) struct AttemptTerminalFacts { + pub(super) provider: AttemptProviderOutcome, + pub(super) delivery: AttemptClientDelivery, +} + +impl AttemptTerminalFacts { + /// 记入 usage / candidate / 效果的人类可读原因。 + pub(super) const fn reason(self) -> &'static str { + if let Some(reason) = self.delivery.aborted_reason() { + return reason; + } + match self.provider { + AttemptProviderOutcome::Terminal { + cancelled_by_provider: true, + .. + } => "provider cancelled the response", + AttemptProviderOutcome::Terminal { .. } => { + "provider returned a terminal response event" + } + AttemptProviderOutcome::Aborted { reason, .. } => reason, + } + } + + /// 供应商侧强制错误原因:只有「供应商没给出终态、且内容已完整交付客户端」 + /// 才算,用于给终态摘要补 `parser_error`。 + /// + /// 客户端投递失败不是供应商的错误,所以那一侧返回 `None`——与现状 + /// `ResponsesWebSocketTurnOutcome::forced_error()` 对 `Cancelled` 返回 + /// `None` 一致。 + pub(super) const fn forced_error(self) -> Option<&'static str> { + match (self.provider, self.delivery) { + ( + AttemptProviderOutcome::Aborted { reason, .. }, + AttemptClientDelivery::Complete, + ) => Some(reason), + _ => None, + } + } +} + +/// 这条 usage 记录是否计费。`Void` 等价于现状传给 +/// `record_stream_terminal(.., cancelled = true)` 的那一侧。 +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(super) enum AttemptBilling { + Billed, + Void, +} + +impl AttemptBilling { + pub(super) const fn is_void(self) -> bool { + matches!(self, Self::Void) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(super) enum AttemptCandidateStatus { + Success, + Failed, + Cancelled, +} + +/// candidate 行上记录的错误分类。 +/// +/// 与 [`AttemptCandidateStatus`] 刻意分开:`missing_terminal` 为真而记账层 +/// 判定不算失败(report kind 不要求观察到终态事件)时,现状会写出 +/// 「状态 Success + error_type=stream_missing_terminal_event」的组合, +/// 这里必须原样保留。 +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(super) enum AttemptCandidateError { + None, + Cancelled, + MissingTerminal, + TerminalError, +} + +/// 一轮 turn 结束后要投射给供应商/密钥池的效果。 +/// +/// 每个分支都会释放 pool key lease:`ProviderFailure` 由 `PoolError` 释放, +/// `ProviderSuccess` 由 `PoolSuccessStream` 释放,其余情况直接释放。少一条 +/// 分支就会把 lease 挂到 TTL 过期,等于短时间占死一把 key。 +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(super) 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 释放。 +pub(super) 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 + } +} + +/// 把「结算触发信号」+「已观察到的 provider 终态」映射成两个正交事实。 +/// +/// `ResponsesWebSocketTurnOutcome` 描述的是 relay loop 为什么现在结算这一 +/// attempt,它对 provider 的信息量并不总是完整的: +/// +/// - `ProviderTerminal` / `Failure` 本身就在描述供应商这一轮的结果,是权威的。 +/// - `Cancelled` 只说明「我们为客户端或连接层面的原因停下了」,不携带任何 +/// provider 信息。已经观察到的 provider 终态是独立事实,不能被它覆盖—— +/// 这正是评审第 5 条要求分开记录的那一处。 +pub(super) fn attempt_facts_for_outcome( + observed_provider_terminal: Option, + settling: ResponsesWebSocketTurnOutcome, +) -> AttemptTerminalFacts { + match settling { + ResponsesWebSocketTurnOutcome::ProviderTerminal { + status_code, + cancelled, + } => AttemptTerminalFacts { + provider: AttemptProviderOutcome::Terminal { + status_code, + cancelled_by_provider: cancelled, + }, + delivery: AttemptClientDelivery::Complete, + }, + ResponsesWebSocketTurnOutcome::Failure { + status_code, + reason, + } => AttemptTerminalFacts { + provider: AttemptProviderOutcome::Aborted { + status_code, + reason, + // 现状只有 504 一族(首事件/终态超时)会投射 pool stream timeout。 + stream_timeout: status_code == STREAM_TIMEOUT_STATUS_CODE, + }, + delivery: AttemptClientDelivery::Complete, + }, + ResponsesWebSocketTurnOutcome::Cancelled { reason } => AttemptTerminalFacts { + provider: observed_provider_terminal.unwrap_or(AttemptProviderOutcome::Aborted { + status_code: CLIENT_CANCELLED_STATUS_CODE, + reason, + stream_timeout: false, + }), + delivery: AttemptClientDelivery::Aborted { reason }, + }, + } +} + +/// attempt 对外记录的状态码。 +/// +/// 投递失败一律记 499:与拆分前 `ResponsesWebSocketTurnOutcome::Cancelled` +/// 走的分支一致。payload 需要在结算判定之前就知道状态码,所以单独暴露。 +pub(super) const fn attempt_status_code(facts: AttemptTerminalFacts) -> u16 { + if facts.delivery.is_aborted() { + CLIENT_CANCELLED_STATUS_CODE + } else { + facts.provider.status_code() + } +} + +/// 结算判定的输入:两个正交事实 + 记账层对这条 report 的判定 + 终态摘要事实。 +#[derive(Debug, Clone, Copy)] +pub(super) struct AttemptSettlementInputs { + pub(super) facts: AttemptTerminalFacts, + /// `aether_usage_runtime::stream_report_represents_failure(payload)` 的结果。 + pub(super) report_represents_failure: bool, + /// 终态摘要里是否观察到了 finish。 + pub(super) observed_finish: bool, + /// 终态摘要里是否带解析错误。 + pub(super) has_parser_error: bool, +} + +/// 结算动作。 +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(super) struct AttemptSettlement { + pub(super) status_code: u16, + pub(super) billing: AttemptBilling, + pub(super) candidate_status: AttemptCandidateStatus, + pub(super) candidate_error: AttemptCandidateError, + pub(super) provider_effect: ResponsesWebSocketTurnEffect, + pub(super) submit_execution_report: bool, +} + +/// 由两个正交事实推出结算动作。唯一的判定入口,表驱动测试逐行锁死。 +/// +/// 当前口径与拆分前的 `finalize()` 完全一致:客户端投递失败与「供应商声明取消」 +/// 落在同一侧(作废账单、candidate 记 Cancelled、只释放 lease、不提交 execution +/// report)。即使已经观察到 provider 终态也是如此——这一行是刻意保留的现状, +/// 修正它是独立的一步。 +pub(super) const fn classify_attempt_settlement( + inputs: AttemptSettlementInputs, +) -> AttemptSettlement { + let AttemptSettlementInputs { + facts, + report_represents_failure, + observed_finish, + has_parser_error, + } = inputs; + + let cancelled = facts.provider.cancelled_by_provider() || facts.delivery.is_aborted(); + let status_code = attempt_status_code(facts); + let failed = !cancelled && report_represents_failure; + let missing_terminal = !cancelled && !observed_finish; + let projects_provider_failure = !cancelled + && (status_code >= 400 + || facts.forced_error().is_some() + || has_parser_error + || missing_terminal); + + let candidate_status = if cancelled { + AttemptCandidateStatus::Cancelled + } else if failed { + AttemptCandidateStatus::Failed + } else { + AttemptCandidateStatus::Success + }; + let candidate_error = if cancelled { + AttemptCandidateError::Cancelled + } else if missing_terminal { + AttemptCandidateError::MissingTerminal + } else if failed { + AttemptCandidateError::TerminalError + } else { + AttemptCandidateError::None + }; + + AttemptSettlement { + status_code, + billing: if cancelled { + AttemptBilling::Void + } else { + AttemptBilling::Billed + }, + candidate_status, + candidate_error, + provider_effect: classify_responses_websocket_turn_effect( + cancelled, + projects_provider_failure, + failed, + ), + submit_execution_report: !cancelled, + } +} + +#[cfg(test)] +mod tests { + use super::{ + attempt_facts_for_outcome, classify_attempt_settlement, + classify_responses_websocket_turn_effect, AttemptBilling, AttemptCandidateError, + AttemptCandidateStatus, AttemptClientDelivery, AttemptProviderOutcome, AttemptSettlement, + AttemptSettlementInputs, AttemptTerminalFacts, ResponsesWebSocketTurnEffect, + }; + use super::super::turn::ResponsesWebSocketTurnOutcome; + + fn settle( + provider: AttemptProviderOutcome, + delivery: AttemptClientDelivery, + report_represents_failure: bool, + observed_finish: bool, + has_parser_error: bool, + ) -> AttemptSettlement { + classify_attempt_settlement(AttemptSettlementInputs { + facts: AttemptTerminalFacts { provider, delivery }, + report_represents_failure, + observed_finish, + has_parser_error, + }) + } + + const fn terminal(status_code: u16) -> AttemptProviderOutcome { + AttemptProviderOutcome::Terminal { + status_code, + cancelled_by_provider: false, + } + } + + const fn provider_cancelled() -> AttemptProviderOutcome { + AttemptProviderOutcome::Terminal { + status_code: 499, + cancelled_by_provider: true, + } + } + + const fn aborted(status_code: u16, reason: &'static str) -> AttemptProviderOutcome { + AttemptProviderOutcome::Aborted { + status_code, + reason, + stream_timeout: status_code == 504, + } + } + + /// §1.6 现状 outcome → 双事实映射表,逐行。 + #[test] + fn every_settle_signal_maps_to_a_provider_outcome_and_a_client_delivery() { + assert_eq!( + attempt_facts_for_outcome( + None, + ResponsesWebSocketTurnOutcome::ProviderTerminal { + status_code: 200, + cancelled: false, + }, + ), + AttemptTerminalFacts { + provider: terminal(200), + delivery: AttemptClientDelivery::Complete, + } + ); + assert_eq!( + attempt_facts_for_outcome( + None, + ResponsesWebSocketTurnOutcome::ProviderTerminal { + status_code: 499, + cancelled: true, + }, + ), + AttemptTerminalFacts { + provider: provider_cancelled(), + delivery: AttemptClientDelivery::Complete, + } + ); + assert_eq!( + attempt_facts_for_outcome(None, ResponsesWebSocketTurnOutcome::upstream_closed()), + AttemptTerminalFacts { + provider: aborted( + 502, + "upstream WebSocket closed before provider terminal event" + ), + delivery: AttemptClientDelivery::Complete, + } + ); + assert_eq!( + attempt_facts_for_outcome(None, ResponsesWebSocketTurnOutcome::client_disconnected()), + AttemptTerminalFacts { + provider: aborted(499, "client disconnected before provider terminal event"), + delivery: AttemptClientDelivery::Aborted { + reason: "client disconnected before provider terminal event", + }, + } + ); + + // 超时一族必须保留 stream_timeout 标记,否则 pool stream timeout 效果丢失。 + let first_event_timeout = + attempt_facts_for_outcome(None, ResponsesWebSocketTurnOutcome::first_event_timeout()); + assert!(first_event_timeout.provider.stream_timeout()); + let terminal_timeout = + attempt_facts_for_outcome(None, ResponsesWebSocketTurnOutcome::terminal_timeout()); + assert!(terminal_timeout.provider.stream_timeout()); + // 非 504 的失败不得被当成流式超时。 + assert!(!attempt_facts_for_outcome( + None, + ResponsesWebSocketTurnOutcome::upstream_closed() + ) + .provider + .stream_timeout()); + // provider 终态即使状态码是 504 也不投射 stream timeout:现状 + // `stream_timeout()` 只匹配 Failure 分支。 + assert!(!attempt_facts_for_outcome( + None, + ResponsesWebSocketTurnOutcome::ProviderTerminal { + status_code: 504, + cancelled: false, + }, + ) + .provider + .stream_timeout()); + } + + /// `Cancelled` 不携带 provider 信息,已观察到的终态不能被它覆盖; + /// `ProviderTerminal` / `Failure` 本身就是权威的 provider 事实。 + #[test] + fn an_observed_provider_terminal_survives_a_client_side_cancellation() { + let observed = terminal(200); + + let facts = attempt_facts_for_outcome( + Some(observed), + ResponsesWebSocketTurnOutcome::client_disconnected(), + ); + assert_eq!(facts.provider, observed); + assert_eq!( + facts.delivery, + AttemptClientDelivery::Aborted { + reason: "client disconnected before provider terminal event", + } + ); + + // 权威信号不被已记录事实改写。 + let facts = attempt_facts_for_outcome( + Some(observed), + ResponsesWebSocketTurnOutcome::upstream_closed(), + ); + assert_eq!( + facts.provider, + aborted(502, "upstream WebSocket closed before provider terminal event") + ); + assert_eq!(facts.delivery, AttemptClientDelivery::Complete); + } + + /// 投递失败时 `forced_error` 必须为 `None`:客户端走了不是供应商的错误。 + /// 与现状 `ResponsesWebSocketTurnOutcome::forced_error()` 对 `Cancelled` + /// 返回 `None` 一致。 + #[test] + fn only_a_provider_abort_with_complete_delivery_is_a_forced_error() { + assert_eq!( + AttemptTerminalFacts { + provider: aborted(502, "upstream failed"), + delivery: AttemptClientDelivery::Complete, + } + .forced_error(), + Some("upstream failed") + ); + assert_eq!( + AttemptTerminalFacts { + provider: aborted(499, "client went away"), + delivery: AttemptClientDelivery::Aborted { + reason: "client went away" + }, + } + .forced_error(), + None + ); + assert_eq!( + AttemptTerminalFacts { + provider: terminal(200), + delivery: AttemptClientDelivery::Complete, + } + .forced_error(), + None + ); + } + + #[test] + fn the_recorded_reason_prefers_the_client_delivery_failure() { + assert_eq!( + AttemptTerminalFacts { + provider: terminal(200), + delivery: AttemptClientDelivery::Aborted { + reason: "client went away" + }, + } + .reason(), + "client went away" + ); + assert_eq!( + AttemptTerminalFacts { + provider: provider_cancelled(), + delivery: AttemptClientDelivery::Complete, + } + .reason(), + "provider cancelled the response" + ); + assert_eq!( + AttemptTerminalFacts { + provider: terminal(200), + delivery: AttemptClientDelivery::Complete, + } + .reason(), + "provider returned a terminal response event" + ); + assert_eq!( + AttemptTerminalFacts { + provider: aborted(502, "upstream failed"), + delivery: AttemptClientDelivery::Complete, + } + .reason(), + "upstream failed" + ); + } + + /// §1.6 结算表,逐行。 + #[test] + fn settlement_table_row_provider_cancelled_is_void_regardless_of_delivery() { + for delivery in [ + AttemptClientDelivery::Complete, + AttemptClientDelivery::Aborted { reason: "gone" }, + ] { + for report_represents_failure in [false, true] { + let settlement = + settle(provider_cancelled(), delivery, report_represents_failure, true, false); + assert_eq!( + settlement, + AttemptSettlement { + status_code: 499, + billing: AttemptBilling::Void, + candidate_status: AttemptCandidateStatus::Cancelled, + candidate_error: AttemptCandidateError::Cancelled, + provider_effect: ResponsesWebSocketTurnEffect::ReleasePoolKeyLease, + submit_execution_report: false, + }, + "delivery={delivery:?} report_failure={report_represents_failure}" + ); + } + } + } + + #[test] + fn settlement_table_row_aborted_provider_with_aborted_delivery_is_void() { + let settlement = settle( + aborted(499, "client went away"), + AttemptClientDelivery::Aborted { + reason: "client went away", + }, + true, + false, + false, + ); + assert_eq!( + settlement, + AttemptSettlement { + status_code: 499, + billing: AttemptBilling::Void, + candidate_status: AttemptCandidateStatus::Cancelled, + candidate_error: AttemptCandidateError::Cancelled, + provider_effect: ResponsesWebSocketTurnEffect::ReleasePoolKeyLease, + submit_execution_report: false, + } + ); + } + + #[test] + fn settlement_table_row_clean_provider_terminal_is_a_billed_success() { + let settlement = settle(terminal(200), AttemptClientDelivery::Complete, false, true, false); + assert_eq!( + settlement, + AttemptSettlement { + status_code: 200, + billing: AttemptBilling::Billed, + candidate_status: AttemptCandidateStatus::Success, + candidate_error: AttemptCandidateError::None, + provider_effect: ResponsesWebSocketTurnEffect::ProviderSuccess, + submit_execution_report: true, + } + ); + } + + /// 合法 `response.incomplete`:记账层判失败,但供应商工作正常, + /// 不扣健康分、只释放 lease,并且账单照记。 + #[test] + fn settlement_table_row_legitimate_incomplete_is_billed_without_provider_failure() { + let settlement = settle(terminal(200), AttemptClientDelivery::Complete, true, true, false); + assert_eq!( + settlement, + AttemptSettlement { + status_code: 200, + billing: AttemptBilling::Billed, + candidate_status: AttemptCandidateStatus::Failed, + candidate_error: AttemptCandidateError::TerminalError, + provider_effect: ResponsesWebSocketTurnEffect::ReleasePoolKeyLease, + submit_execution_report: true, + } + ); + } + + #[test] + fn settlement_table_row_provider_abort_projects_a_provider_failure() { + let settlement = settle( + aborted(502, "upstream WebSocket closed before provider terminal event"), + AttemptClientDelivery::Complete, + true, + false, + false, + ); + assert_eq!( + settlement, + AttemptSettlement { + status_code: 502, + billing: AttemptBilling::Billed, + candidate_status: AttemptCandidateStatus::Failed, + candidate_error: AttemptCandidateError::MissingTerminal, + provider_effect: ResponsesWebSocketTurnEffect::ProviderFailure, + submit_execution_report: true, + } + ); + } + + /// ✱ C2 保留现状的那一行:provider 终态已到达,但客户端投递失败 ⇒ 仍作废账单。 + /// 这一行是下一步唯一要改的地方,先在这里锁住现状。 + #[test] + fn settlement_table_row_client_delivery_failure_currently_voids_a_reached_terminal() { + let settlement = settle( + terminal(200), + AttemptClientDelivery::Aborted { + reason: "client disconnected before provider terminal event", + }, + false, + true, + false, + ); + assert_eq!( + settlement, + AttemptSettlement { + status_code: 499, + billing: AttemptBilling::Void, + candidate_status: AttemptCandidateStatus::Cancelled, + candidate_error: AttemptCandidateError::Cancelled, + provider_effect: ResponsesWebSocketTurnEffect::ReleasePoolKeyLease, + submit_execution_report: false, + } + ); + } + + /// 记账层判 Success,但摘要没观察到 finish:现状会写出 + /// 「candidate=Success + error_type=stream_missing_terminal_event」, + /// 所以状态与错误分类必须各自独立。 + #[test] + fn a_missing_terminal_can_coexist_with_a_successful_candidate_status() { + let settlement = settle(terminal(200), AttemptClientDelivery::Complete, false, false, false); + assert_eq!(settlement.candidate_status, AttemptCandidateStatus::Success); + assert_eq!( + settlement.candidate_error, + AttemptCandidateError::MissingTerminal + ); + // missing_terminal 仍然要投射供应商失败。 + assert_eq!( + settlement.provider_effect, + ResponsesWebSocketTurnEffect::ProviderFailure + ); + } + + #[test] + fn a_parser_error_projects_a_provider_failure_even_on_a_clean_status_code() { + let settlement = settle(terminal(200), AttemptClientDelivery::Complete, true, true, true); + assert_eq!( + settlement.provider_effect, + ResponsesWebSocketTurnEffect::ProviderFailure + ); + assert_eq!(settlement.billing, AttemptBilling::Billed); + } + + #[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" + ); + } + } + + /// 每一个结算分支都必须释放 lease:这条不变量跨越整张结算表。 + #[test] + fn every_settlement_branch_releases_the_pool_key_lease() { + let providers = [ + terminal(200), + terminal(429), + provider_cancelled(), + aborted(502, "upstream failed"), + aborted(504, "timed out"), + ]; + let deliveries = [ + AttemptClientDelivery::Complete, + AttemptClientDelivery::Aborted { reason: "gone" }, + ]; + for provider in providers { + for delivery in deliveries { + for report_represents_failure in [false, true] { + for observed_finish in [false, true] { + for has_parser_error in [false, true] { + let settlement = settle( + provider, + delivery, + report_represents_failure, + observed_finish, + has_parser_error, + ); + assert!( + settlement.provider_effect.releases_pool_key_lease(), + "provider={provider:?} delivery={delivery:?}" + ); + // 作废账单的分支一律不提交 execution report。 + assert_eq!( + settlement.submit_execution_report, + !settlement.billing.is_void(), + "provider={provider:?} delivery={delivery:?}" + ); + } + } + } + } + } + } +} 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 177b49cbf..96ae5ea30 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn.rs @@ -33,6 +33,11 @@ use tracing::warn; use super::adapter::ResponsesWebSocketProtocolAdapter; use super::admission::ResponsesWebSocketTurnAdmission; use super::frame::ParsedResponsesWebSocketFrame; +use super::settlement::{ + attempt_facts_for_outcome, attempt_status_code, classify_attempt_settlement, + AttemptCandidateError, AttemptCandidateStatus, AttemptProviderOutcome, AttemptSettlementInputs, + AttemptTerminalFacts, ResponsesWebSocketTurnEffect, +}; use crate::ai_serving::api::StreamingStandardTerminalObserver; use crate::ai_serving::{build_openai_responses_stream_plan_from_decision, AiExecutionDecision}; use crate::clock::current_unix_ms; @@ -203,82 +208,13 @@ impl ResponsesWebSocketTurnOutcome { Self::Cancelled { .. } => 499, } } - - const fn cancelled(self) -> bool { - matches!( - self, - Self::ProviderTerminal { - cancelled: true, - .. - } | Self::Cancelled { .. } - ) - } - - const fn forced_error(self) -> Option<&'static str> { - match self { - Self::Failure { reason, .. } => Some(reason), - Self::ProviderTerminal { .. } | Self::Cancelled { .. } => None, - } - } - - const fn stream_timeout(self) -> bool { - matches!( - self, - Self::Failure { - status_code: 504, - .. - } - ) - } } -/// 一轮 turn 结束后要投射给供应商/密钥池的效果。 +/// 一次 provider attempt:一条上游执行,也是一条独立的 usage/candidate 记录。 /// -/// 每个分支都会释放 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 { +/// 与 [`super::turn_state::LogicalTurn`] 分工明确:logical turn 是客户端看到的 +/// 一轮请求(可能包含多个 attempt),attempt 只负责这一次上游执行的记账事实。 +pub(super) struct ResponsesProviderAttempt { plan: ExecutionPlan, trace_id: String, report_kind: String, @@ -298,6 +234,9 @@ pub(super) struct ResponsesWebSocketTurn { terminal_timeout: Duration, admission: Option, terminal_error_body: Option, + /// 观察到的 provider 终态事实,与「为什么现在结算」这个信号分开保存。 + /// 客户端投递失败不会把它擦掉。 + provider_outcome: Option, } /// 组装一轮 turn 的 decision。 @@ -345,7 +284,7 @@ pub(super) async fn begin_responses_websocket_turn( control_decision: &GatewayControlDecision, decision: AiExecutionDecision, client_event: &Value, -) -> Result { +) -> Result { let planned_report_context = decision.report_context.clone(); let effective_control_decision = match refresh_websocket_turn_auth_context(state, control_decision, parts, client_event) @@ -482,7 +421,7 @@ pub(super) async fn begin_responses_websocket_turn( ) .await; - Ok(ResponsesWebSocketTurn { + Ok(ResponsesProviderAttempt { trace_id: plan.request_id.clone(), plan, report_kind, @@ -502,6 +441,7 @@ pub(super) async fn begin_responses_websocket_turn( terminal_timeout, admission: Some(admission), terminal_error_body: None, + provider_outcome: None, }) } @@ -583,7 +523,7 @@ fn websocket_auth_rejection_error(rejection: GatewayLocalAuthRejection) -> Gatew } } -impl ResponsesWebSocketTurn { +impl ResponsesProviderAttempt { /// Releases all per-turn capacity before terminal persistence starts. /// Provider-pool runtime tokens normally use an awaited removal. The /// bounded wait prevents a broken runtime backend from stalling the relay; @@ -675,6 +615,10 @@ impl ResponsesWebSocketTurn { .and_then(|event| serde_json::to_string(event).ok()); } if let Some(outcome) = provider_terminal_outcome(frame) { + // provider 的终态是独立事实:先记下来,之后即使客户端投递失败、 + // 结算信号变成 Cancelled,这条事实也不会被擦掉。 + self.provider_outcome + .get_or_insert(attempt_facts_for_outcome(None, outcome).provider); return Some(ResponsesWebSocketTurnObservation::Terminal(outcome)); } if frame.is_started() { @@ -751,13 +695,17 @@ impl ResponsesWebSocketTurn { .await; } - async fn finalize(mut self, state: &AppState, outcome: ResponsesWebSocketTurnOutcome) { - let summary = self.finish_summary(outcome); - let cancelled = outcome.cancelled(); - let status_code = outcome.status_code(); - let missing_terminal = !cancelled && !summary.observed_finish; + /// 结算这一个 attempt。 + /// + /// `outcome` 是「为什么现在结算」的信号,不是供应商事实本身: + /// [`attempt_facts_for_outcome`] 把它和已观察到的 provider 终态一起, + /// 拆成 provider outcome 与 client delivery 两个正交事实,再由 + /// [`classify_attempt_settlement`] 一张表推出账单、candidate 状态和效果。 + async fn settle(mut self, state: &AppState, outcome: ResponsesWebSocketTurnOutcome) { + let facts = attempt_facts_for_outcome(self.provider_outcome, outcome); + let summary = self.finish_summary(facts); let terminal_error_body = self.terminal_error_body.take(); - let outcome_reason = outcome_reason(outcome); + let outcome_reason = facts.reason().to_string(); let telemetry = Some(self.telemetry()); let (provider_body_base64, provider_body_state) = encode_stream_capture(&self.provider_capture, self.provider_capture_truncated); @@ -767,7 +715,7 @@ impl ResponsesWebSocketTurn { trace_id: self.trace_id.clone(), report_kind: self.report_kind, report_context: self.report_context, - status_code, + status_code: attempt_status_code(facts), headers: self.provider_headers, provider_body_base64, provider_body_state, @@ -776,7 +724,13 @@ impl ResponsesWebSocketTurn { terminal_summary: Some(summary.clone()), telemetry, }; - let failed = !cancelled && stream_report_represents_failure(&payload); + let settlement = classify_attempt_settlement(AttemptSettlementInputs { + facts, + report_represents_failure: stream_report_represents_failure(&payload), + observed_finish: summary.observed_finish, + has_parser_error: summary.parser_error.is_some(), + }); + let billing_void = settlement.billing.is_void(); // Do not hold gateway/provider capacity while usage and audit writes // run. The turn has a complete terminal payload at this point. @@ -799,34 +753,31 @@ impl ResponsesWebSocketTurn { let usage_data = Arc::clone(state.usage_lifecycle_data_state()); await_detachable_lifecycle_stage(&self.trace_id, "usage_terminal", async move { usage_runtime - .record_stream_terminal(usage_data.as_ref(), context_seed, payload_seed, cancelled) + .record_stream_terminal(usage_data.as_ref(), context_seed, payload_seed, billing_void) .await; }) .await; - let (error_type, error_message) = if cancelled { - ( + let (error_type, error_message) = match settlement.candidate_error { + AttemptCandidateError::Cancelled => ( Some("websocket_cancelled".to_string()), Some(outcome_reason.clone()), - ) - } else if missing_terminal { - ( + ), + AttemptCandidateError::MissingTerminal => ( Some("stream_missing_terminal_event".to_string()), Some(summary.parser_error.clone().unwrap_or_else(|| { "upstream Responses WebSocket ended before a provider terminal event" .to_string() })), - ) - } else if failed { - ( + ), + AttemptCandidateError::TerminalError => ( Some("stream_terminal_error".to_string()), summary .parser_error .clone() .or_else(|| Some(outcome_reason.clone())), - ) - } else { - (None, None) + ), + AttemptCandidateError::None => (None, None), }; let _ = await_websocket_lifecycle_stage( &self.trace_id, @@ -836,14 +787,12 @@ impl ResponsesWebSocketTurn { &self.plan, payload.report_context.as_ref(), SchedulerRequestCandidateStatusUpdate { - status: if cancelled { - RequestCandidateStatus::Cancelled - } else if failed { - RequestCandidateStatus::Failed - } else { - RequestCandidateStatus::Success + status: match settlement.candidate_status { + AttemptCandidateStatus::Cancelled => RequestCandidateStatus::Cancelled, + AttemptCandidateStatus::Failed => RequestCandidateStatus::Failed, + AttemptCandidateStatus::Success => RequestCandidateStatus::Success, }, - status_code: Some(status_code), + status_code: Some(settlement.status_code), error_type, error_message, latency_ms: payload @@ -864,16 +813,9 @@ impl ResponsesWebSocketTurn { plan: &self.plan, report_context: payload.report_context.as_ref(), }; - let projects_provider_failure = !cancelled - && (status_code >= 400 - || 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 { - match provider_effect { + match settlement.provider_effect { ResponsesWebSocketTurnEffect::ReleasePoolKeyLease => { release_local_pool_key_lease(state, effect_context).await; } @@ -883,11 +825,11 @@ impl ResponsesWebSocketTurn { .or(summary.parser_error.as_deref()) .unwrap_or(outcome_reason.as_str()); let mut effect = LocalStreamFailureEffect::new( - status_code, + settlement.status_code, &payload.headers, Some(response_text), ); - if outcome.stream_timeout() { + if facts.provider.stream_timeout() { effect = effect.with_stream_timeout(); } apply_local_stream_failure_effects(state, effect_context, effect).await; @@ -911,7 +853,7 @@ impl ResponsesWebSocketTurn { // The normal execution runtime does not submit a stream report after a // downstream disconnect either. The terminal usage record above still // captures cancellation without applying provider-success side effects. - if !cancelled { + if settlement.submit_execution_report { if let Some(Err(error)) = await_websocket_lifecycle_stage( &self.trace_id, "execution_report", @@ -940,10 +882,13 @@ impl ResponsesWebSocketTurn { ); } - fn finish_summary( - &mut self, - outcome: ResponsesWebSocketTurnOutcome, - ) -> ExecutionStreamTerminalSummary { + /// 终态摘要。 + /// + /// 只消费两个正交事实:`forced_error` 只在「供应商没给出终态且内容已完整 + /// 交付」时补 parser_error,`cancelled` 覆盖供应商声明取消与客户端投递失败 + /// 两种情形——与拆分前 `outcome.forced_error()` / `outcome.cancelled()` 的 + /// 取值逐一对应。 + fn finish_summary(&mut self, facts: AttemptTerminalFacts) -> ExecutionStreamTerminalSummary { let fallback_context = json!({ "provider_api_format": "openai:responses", "client_api_format": "openai:responses", @@ -957,12 +902,12 @@ impl ResponsesWebSocketTurn { self.observer.latest_summary().cloned().unwrap_or_default() } }; - if let Some(reason) = outcome.forced_error() { + if let Some(reason) = facts.forced_error() { if summary.parser_error.is_none() { summary.parser_error = Some(reason.to_string()); } } - if outcome.cancelled() { + if facts.provider.cancelled_by_provider() || facts.delivery.is_aborted() { summary.observed_finish = true; if summary.finish_reason.is_none() { summary.finish_reason = Some("cancelled".to_string()); @@ -984,7 +929,7 @@ impl ResponsesWebSocketTurn { } } -impl ResponsesWebSocketTurn { +impl ResponsesProviderAttempt { /// Finalizes a turn whose owner is already gone, releasing admission first. /// /// The normal path releases admission before spawning the finalizer; a @@ -995,18 +940,18 @@ impl ResponsesWebSocketTurn { outcome: ResponsesWebSocketTurnOutcome, ) { self.release_admission().await; - self.finalize(state, outcome).await; + self.settle(state, outcome).await; } } pub(super) async fn spawn_responses_websocket_turn_finalization( state: AppState, - mut turn: ResponsesWebSocketTurn, + mut turn: ResponsesProviderAttempt, outcome: ResponsesWebSocketTurnOutcome, ) -> tokio::task::JoinHandle<()> { turn.release_admission().await; tokio::spawn(async move { - turn.finalize(&state, outcome).await; + turn.settle(&state, outcome).await; }) } @@ -1130,19 +1075,6 @@ fn resolve_responses_websocket_turn_timeouts( ) } -fn outcome_reason(outcome: ResponsesWebSocketTurnOutcome) -> String { - match outcome { - ResponsesWebSocketTurnOutcome::ProviderTerminal { - cancelled: true, .. - } => "provider cancelled the response".to_string(), - ResponsesWebSocketTurnOutcome::ProviderTerminal { - cancelled: false, .. - } => "provider returned a terminal response event".to_string(), - ResponsesWebSocketTurnOutcome::Cancelled { reason } - | ResponsesWebSocketTurnOutcome::Failure { reason, .. } => reason.to_string(), - } -} - fn websocket_event_as_sse_line(event: &Value) -> Vec { let payload = serde_json::to_string(event).unwrap_or_else(|_| { json!({ @@ -1241,11 +1173,15 @@ mod tests { use crate::ai_serving::api::StreamingStandardTerminalObserver; use super::super::frame::ParsedResponsesWebSocketFrame; + use super::super::settlement::{ + attempt_facts_for_outcome, classify_attempt_settlement, AttemptBilling, + AttemptSettlementInputs, ResponsesWebSocketTurnEffect, + }; use super::{ - 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, + prepare_websocket_report_context, provider_terminal_outcome, + resolve_responses_websocket_turn_timeouts, websocket_event_as_sse_line, + ResponsesWebSocketTurnDeadline, ResponsesWebSocketTurnOutcome, + ResponsesWebSocketTurnTimeoutPhase, }; #[test] @@ -1394,9 +1330,11 @@ mod tests { }) ); let outcome = outcome.expect("incomplete should end the turn"); - assert!(!outcome.cancelled()); - assert!(outcome.forced_error().is_none()); - assert!(!outcome.stream_timeout()); + let facts = attempt_facts_for_outcome(None, outcome); + assert!(!facts.provider.cancelled_by_provider()); + assert!(facts.forced_error().is_none()); + assert!(!facts.provider.stream_timeout()); + assert!(facts.provider.is_terminal()); let report_context = json!({ "provider_api_format": "openai:responses", @@ -1419,16 +1357,21 @@ mod tests { 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); + // 结算表用这些事实决定是否投射供应商失败:合法 incomplete 即使被记账层 + // 判成失败,也必须落在「不扣健康分、只释放 lease」一侧,且账单照记。 + let settlement = classify_attempt_settlement(AttemptSettlementInputs { + facts, + report_represents_failure: true, + observed_finish: summary.observed_finish, + has_parser_error: summary.parser_error.is_some(), + }); + assert_eq!(settlement.status_code, 200); + assert_eq!(settlement.billing, AttemptBilling::Billed); + assert_eq!( + settlement.provider_effect, + ResponsesWebSocketTurnEffect::ReleasePoolKeyLease + ); + assert!(settlement.submit_execution_report); } #[test] @@ -1445,67 +1388,6 @@ mod tests { ); } - #[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!({ @@ -1543,10 +1425,19 @@ mod tests { // cancellation: cancelled turns skip the stream report entirely, which // would defeat the point of reclaiming it. let outcome = ResponsesWebSocketTurnOutcome::relay_task_abandoned(); + let facts = attempt_facts_for_outcome(None, outcome); + let settlement = classify_attempt_settlement(AttemptSettlementInputs { + facts, + report_represents_failure: true, + observed_finish: false, + has_parser_error: false, + }); - assert_eq!(outcome.status_code(), 500); - assert!(!outcome.cancelled()); - assert!(outcome.forced_error().is_some()); + assert_eq!(settlement.status_code, 500); + assert_eq!(settlement.billing, AttemptBilling::Billed); + assert!(!facts.provider.cancelled_by_provider()); + assert!(facts.forced_error().is_some()); + assert!(settlement.submit_execution_report); } #[test] diff --git a/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn_state.rs b/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn_state.rs index d2cdd9812..eb4e42b22 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn_state.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn_state.rs @@ -7,7 +7,7 @@ use serde_json::Value; -use super::lifecycle::ActiveResponsesWebSocketTurn; +use super::lifecycle::ActiveProviderAttempt; use super::request::response_create_has_previous_response_id; /// 客户端一次 `response.create` 对应的 logical turn。 @@ -58,9 +58,9 @@ impl LogicalTurn { /// 连接上「有没有正在进行的 logical turn」这一唯一事实。 /// /// 类型参数只为测试留出注入点:生产代码一律用默认的 -/// [`ActiveResponsesWebSocketTurn`],测试用轻量替身驱动同一套转换逻辑, +/// [`ActiveProviderAttempt`],测试用轻量替身驱动同一套转换逻辑, /// 不必构造 `AppState` 和真实 socket。 -pub(super) enum ResponsesTurnState { +pub(super) enum ResponsesTurnState { /// 没有进行中的 logical turn。上游可能仍绑定,也可能已被 detach。 Idle, /// 一个 logical turn 正在等待 provider 终态:logical 与当前 attempt 同时存在。