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 748da4f88..a224d33a3 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/connection.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/connection.rs @@ -23,6 +23,7 @@ use super::quota::{ use super::relay_policy::{ classify_quota_relay, fatal_relay_policy, FatalRelaySignal, QuotaRelayAction, QuotaRelayFacts, }; +use super::settlement::settle_signal_for_client_delivery_failure; use super::state::BoundResponsesConnection; use super::turn::{ ResponsesProviderAttempt, ResponsesWebSocketTurnObservation, ResponsesWebSocketTurnOutcome, @@ -41,6 +42,11 @@ use crate::AppState; const LOG_TARGET: &str = "aether_gateway::handlers::proxy::responses_ws"; +/// 写客户端 socket 失败时记录的投递失败原因。刻意不说「客户端在终态前断开」: +/// 供应商的终态可能已经到达,只是最后一跳没送出去。 +const CLIENT_DELIVERY_FAILED_REASON: &str = + "gateway could not relay the provider event to the client"; + macro_rules! debug { ($($arg:tt)*) => { tracing::debug!(target: LOG_TARGET, $($arg)*) @@ -484,12 +490,18 @@ pub(super) async fn relay_bound_connection( websocket = true, trace_id = %context.trace_id, error_code = error.as_str(), + provider_terminal_reached = terminal_outcome.is_some(), "gateway could not relay a provider event to the client" ); + // 投递失败是独立事实,不能覆盖已经到达的 provider 终态: + // 供应商已经完成推理并消耗 token,账单按它的终态计。 + bound + .turn_state + .record_client_delivery_aborted(CLIENT_DELIVERY_FAILED_REASON); finalize_active_turn( bound, state, - ResponsesWebSocketTurnOutcome::client_disconnected(), + settle_signal_for_client_delivery_failure(terminal_outcome), ).await; close_bound_upstream(bound).await; break; diff --git a/apps/aether-gateway/src/handlers/proxy/websocket/responses/settlement.rs b/apps/aether-gateway/src/handlers/proxy/websocket/responses/settlement.rs index 5c7ec89e8..77e37e48a 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/settlement.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/settlement.rs @@ -155,6 +155,9 @@ pub(super) enum AttemptCandidateStatus { pub(super) enum AttemptCandidateError { None, Cancelled, + /// 供应商已经给出终态,但这一轮内容没能完整交付给客户端。账单照记, + /// candidate 行上留下这条事实。 + ClientDeliveryFailed, MissingTerminal, TerminalError, } @@ -205,7 +208,8 @@ pub(super) const fn classify_responses_websocket_turn_effect( } } -/// 把「结算触发信号」+「已观察到的 provider 终态」映射成两个正交事实。 +/// 把「结算触发信号」+「已观察到的 provider 终态」+「已记录的投递结果」映射成 +/// 两个正交事实。 /// /// `ResponsesWebSocketTurnOutcome` 描述的是 relay loop 为什么现在结算这一 /// attempt,它对 provider 的信息量并不总是完整的: @@ -214,11 +218,16 @@ pub(super) const fn classify_responses_websocket_turn_effect( /// - `Cancelled` 只说明「我们为客户端或连接层面的原因停下了」,不携带任何 /// provider 信息。已经观察到的 provider 终态是独立事实,不能被它覆盖—— /// 这正是评审第 5 条要求分开记录的那一处。 +/// +/// `recorded_delivery` 是 relay loop 明确记下的投递失败(写客户端 socket 失败)。 +/// 它与结算信号推出的投递结果取「只要有一侧失败就是失败」,并优先保留明确记录 +/// 的原因。 pub(super) fn attempt_facts_for_outcome( observed_provider_terminal: Option, + recorded_delivery: AttemptClientDelivery, settling: ResponsesWebSocketTurnOutcome, ) -> AttemptTerminalFacts { - match settling { + let facts = match settling { ResponsesWebSocketTurnOutcome::ProviderTerminal { status_code, cancelled, @@ -249,19 +258,46 @@ pub(super) fn attempt_facts_for_outcome( }), delivery: AttemptClientDelivery::Aborted { reason }, }, + }; + AttemptTerminalFacts { + delivery: match recorded_delivery { + AttemptClientDelivery::Aborted { .. } => recorded_delivery, + AttemptClientDelivery::Complete => facts.delivery, + }, + ..facts } } +/// 客户端投递失败时应该用哪个结算信号。 +/// +/// provider 终态已经到达就用那条终态:它是权威的 provider 事实,绝不能被 +/// `client_disconnected()` 覆盖掉——那正是把已完成响应记成 void billing 的原因。 +/// 供应商还没给出终态时,客户端断开才是这一 attempt 的全部结论。 +pub(super) fn settle_signal_for_client_delivery_failure( + terminal_outcome: Option, +) -> ResponsesWebSocketTurnOutcome { + terminal_outcome.unwrap_or_else(ResponsesWebSocketTurnOutcome::client_disconnected) +} + +/// 这一个 attempt 的账单是否作废。 +/// +/// 只有两种情况作废:供应商自己声明取消,或者供应商根本没给出终态而客户端 +/// 又已经走了。**供应商已经给出终态时,客户端最后一跳投递失败不作废账单**: +/// 供应商已经完成推理并消耗了 token,客户端还能用 `previous_response_id` +/// 续取这条响应,把成本记成 0 等于让上游账单凭空消失。 +pub(super) const fn attempt_billing_is_void(facts: AttemptTerminalFacts) -> bool { + facts.provider.cancelled_by_provider() + || (facts.delivery.is_aborted() && !facts.provider.is_terminal()) +} + /// attempt 对外记录的状态码。 /// -/// 投递失败一律记 499:与拆分前 `ResponsesWebSocketTurnOutcome::Cancelled` -/// 走的分支一致。payload 需要在结算判定之前就知道状态码,所以单独暴露。 +/// 状态码现在纯粹是 provider 事实:客户端投递失败不再把一条已经拿到 200 +/// 终态的记录改写成 499。作废分支的 provider 状态码本身就是 499 +/// (`response.cancelled` 映射 499,`Cancelled` 信号的兜底也是 499), +/// 所以这些行的取值不变。 pub(super) const fn attempt_status_code(facts: AttemptTerminalFacts) -> u16 { - if facts.delivery.is_aborted() { - CLIENT_CANCELLED_STATUS_CODE - } else { - facts.provider.status_code() - } + facts.provider.status_code() } /// 结算判定的输入:两个正交事实 + 记账层对这条 report 的判定 + 终态摘要事实。 @@ -289,10 +325,8 @@ pub(super) struct AttemptSettlement { /// 由两个正交事实推出结算动作。唯一的判定入口,表驱动测试逐行锁死。 /// -/// 当前口径与拆分前的 `finalize()` 完全一致:客户端投递失败与「供应商声明取消」 -/// 落在同一侧(作废账单、candidate 记 Cancelled、只释放 lease、不提交 execution -/// report)。即使已经观察到 provider 终态也是如此——这一行是刻意保留的现状, -/// 修正它是独立的一步。 +/// provider 终态已到达时,客户端投递失败只影响 candidate 的错误分类,不再作废 +/// 账单、不再把状态码改成 499、也不再跳过供应商效果和 execution report。 pub(super) const fn classify_attempt_settlement( inputs: AttemptSettlementInputs, ) -> AttemptSettlement { @@ -303,25 +337,30 @@ pub(super) const fn classify_attempt_settlement( has_parser_error, } = inputs; - let cancelled = facts.provider.cancelled_by_provider() || facts.delivery.is_aborted(); + let void = attempt_billing_is_void(facts); let status_code = attempt_status_code(facts); - let failed = !cancelled && report_represents_failure; - let missing_terminal = !cancelled && !observed_finish; - let projects_provider_failure = !cancelled + let failed = !void && report_represents_failure; + let missing_terminal = !void && !observed_finish; + let projects_provider_failure = !void && (status_code >= 400 || facts.forced_error().is_some() || has_parser_error || missing_terminal); - let candidate_status = if cancelled { + let candidate_status = if void { AttemptCandidateStatus::Cancelled } else if failed { AttemptCandidateStatus::Failed } else { AttemptCandidateStatus::Success }; - let candidate_error = if cancelled { + // 投递失败排在供应商侧分类之前:这条记录之所以特别,正是因为内容没送到 + // 客户端手上。供应商侧的判定仍然通过 candidate_status 和 error_message + // 保留下来。 + let candidate_error = if void { AttemptCandidateError::Cancelled + } else if facts.delivery.is_aborted() { + AttemptCandidateError::ClientDeliveryFailed } else if missing_terminal { AttemptCandidateError::MissingTerminal } else if failed { @@ -332,7 +371,7 @@ pub(super) const fn classify_attempt_settlement( AttemptSettlement { status_code, - billing: if cancelled { + billing: if void { AttemptBilling::Void } else { AttemptBilling::Billed @@ -340,11 +379,11 @@ pub(super) const fn classify_attempt_settlement( candidate_status, candidate_error, provider_effect: classify_responses_websocket_turn_effect( - cancelled, + void, projects_provider_failure, failed, ), - submit_execution_report: !cancelled, + submit_execution_report: !void, } } @@ -352,9 +391,10 @@ pub(super) const fn classify_attempt_settlement( 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, + classify_responses_websocket_turn_effect, settle_signal_for_client_delivery_failure, + AttemptBilling, AttemptCandidateError, AttemptCandidateStatus, AttemptClientDelivery, + AttemptProviderOutcome, AttemptSettlement, AttemptSettlementInputs, AttemptTerminalFacts, + ResponsesWebSocketTurnEffect, }; use super::super::turn::ResponsesWebSocketTurnOutcome; @@ -401,6 +441,7 @@ mod tests { assert_eq!( attempt_facts_for_outcome( None, + AttemptClientDelivery::Complete, ResponsesWebSocketTurnOutcome::ProviderTerminal { status_code: 200, cancelled: false, @@ -414,6 +455,7 @@ mod tests { assert_eq!( attempt_facts_for_outcome( None, + AttemptClientDelivery::Complete, ResponsesWebSocketTurnOutcome::ProviderTerminal { status_code: 499, cancelled: true, @@ -425,7 +467,7 @@ mod tests { } ); assert_eq!( - attempt_facts_for_outcome(None, ResponsesWebSocketTurnOutcome::upstream_closed()), + attempt_facts_for_outcome(None, AttemptClientDelivery::Complete, ResponsesWebSocketTurnOutcome::upstream_closed()), AttemptTerminalFacts { provider: aborted( 502, @@ -435,7 +477,7 @@ mod tests { } ); assert_eq!( - attempt_facts_for_outcome(None, ResponsesWebSocketTurnOutcome::client_disconnected()), + attempt_facts_for_outcome(None, AttemptClientDelivery::Complete, ResponsesWebSocketTurnOutcome::client_disconnected()), AttemptTerminalFacts { provider: aborted(499, "client disconnected before provider terminal event"), delivery: AttemptClientDelivery::Aborted { @@ -446,14 +488,15 @@ mod tests { // 超时一族必须保留 stream_timeout 标记,否则 pool stream timeout 效果丢失。 let first_event_timeout = - attempt_facts_for_outcome(None, ResponsesWebSocketTurnOutcome::first_event_timeout()); + attempt_facts_for_outcome(None, AttemptClientDelivery::Complete, ResponsesWebSocketTurnOutcome::first_event_timeout()); assert!(first_event_timeout.provider.stream_timeout()); let terminal_timeout = - attempt_facts_for_outcome(None, ResponsesWebSocketTurnOutcome::terminal_timeout()); + attempt_facts_for_outcome(None, AttemptClientDelivery::Complete, ResponsesWebSocketTurnOutcome::terminal_timeout()); assert!(terminal_timeout.provider.stream_timeout()); // 非 504 的失败不得被当成流式超时。 assert!(!attempt_facts_for_outcome( None, + AttemptClientDelivery::Complete, ResponsesWebSocketTurnOutcome::upstream_closed() ) .provider @@ -462,6 +505,7 @@ mod tests { // `stream_timeout()` 只匹配 Failure 分支。 assert!(!attempt_facts_for_outcome( None, + AttemptClientDelivery::Complete, ResponsesWebSocketTurnOutcome::ProviderTerminal { status_code: 504, cancelled: false, @@ -479,6 +523,7 @@ mod tests { let facts = attempt_facts_for_outcome( Some(observed), + AttemptClientDelivery::Complete, ResponsesWebSocketTurnOutcome::client_disconnected(), ); assert_eq!(facts.provider, observed); @@ -492,6 +537,7 @@ mod tests { // 权威信号不被已记录事实改写。 let facts = attempt_facts_for_outcome( Some(observed), + AttemptClientDelivery::Complete, ResponsesWebSocketTurnOutcome::upstream_closed(), ); assert_eq!( @@ -678,14 +724,17 @@ mod tests { ); } - /// ✱ C2 保留现状的那一行:provider 终态已到达,但客户端投递失败 ⇒ 仍作废账单。 - /// 这一行是下一步唯一要改的地方,先在这里锁住现状。 + /// ✱ 修正后的那一行:provider 终态已到达,客户端投递失败不再作废账单。 + /// + /// 供应商已经完成推理并消耗 token,客户端还能用 `previous_response_id` + /// 续取这条响应;把成本记成 0 等于让上游账单凭空消失。投递失败作为独立 + /// 事实留在 candidate 的错误分类里。 #[test] - fn settlement_table_row_client_delivery_failure_currently_voids_a_reached_terminal() { + fn settlement_table_row_client_delivery_failure_keeps_a_reached_terminal_billed() { let settlement = settle( terminal(200), AttemptClientDelivery::Aborted { - reason: "client disconnected before provider terminal event", + reason: "gateway could not relay the provider event to the client", }, false, true, @@ -694,14 +743,100 @@ mod tests { 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, + status_code: 200, + billing: AttemptBilling::Billed, + candidate_status: AttemptCandidateStatus::Success, + candidate_error: AttemptCandidateError::ClientDeliveryFailed, + provider_effect: ResponsesWebSocketTurnEffect::ProviderSuccess, + submit_execution_report: true, } ); + + // 除了 candidate 的错误分类,其余判定与「投递成功」完全一致。 + let delivered = settle(terminal(200), AttemptClientDelivery::Complete, false, true, false); + assert_eq!(settlement.status_code, delivered.status_code); + assert_eq!(settlement.billing, delivered.billing); + assert_eq!(settlement.candidate_status, delivered.candidate_status); + assert_eq!(settlement.provider_effect, delivered.provider_effect); + assert_eq!( + settlement.submit_execution_report, + delivered.submit_execution_report + ); + assert_ne!(settlement.candidate_error, delivered.candidate_error); + } + + /// 供应商还没给出终态时,客户端投递失败仍然作废账单:这一轮确实没有产出。 + #[test] + fn a_delivery_failure_without_a_provider_terminal_still_voids_the_bill() { + let settlement = settle( + aborted(499, "client went away"), + AttemptClientDelivery::Aborted { + reason: "client went away", + }, + false, + false, + false, + ); + assert_eq!(settlement.status_code, 499); + assert_eq!(settlement.billing, AttemptBilling::Void); + assert_eq!( + settlement.candidate_status, + AttemptCandidateStatus::Cancelled + ); + assert_eq!(settlement.candidate_error, AttemptCandidateError::Cancelled); + assert!(!settlement.submit_execution_report); + } + + /// 供应商自己声明取消时,即使内容送到了客户端也不计费。 + #[test] + fn a_provider_declared_cancellation_is_void_even_when_delivered() { + let settlement = + settle(provider_cancelled(), AttemptClientDelivery::Complete, false, true, false); + assert_eq!(settlement.billing, AttemptBilling::Void); + assert_eq!(settlement.candidate_error, AttemptCandidateError::Cancelled); + } + + /// 结算信号的选择:provider 终态已到达就用它,否则才是 client 断开。 + /// 这是修正的核心——旧实现无条件用 client_disconnected() 覆盖, + /// 于是已完成的响应被记成 void billing。 + #[test] + fn a_reached_terminal_is_the_settle_signal_for_a_delivery_failure() { + let terminal_outcome = ResponsesWebSocketTurnOutcome::ProviderTerminal { + status_code: 200, + cancelled: false, + }; + assert_eq!( + settle_signal_for_client_delivery_failure(Some(terminal_outcome)), + terminal_outcome + ); + assert_eq!( + settle_signal_for_client_delivery_failure(None), + ResponsesWebSocketTurnOutcome::client_disconnected() + ); + } + + /// 明确记录的投递失败不会被结算信号推出的「投递成功」覆盖。 + #[test] + fn a_recorded_delivery_failure_survives_a_provider_terminal_settle_signal() { + let facts = attempt_facts_for_outcome( + Some(terminal(200)), + AttemptClientDelivery::Aborted { + reason: "write failed" + }, + ResponsesWebSocketTurnOutcome::ProviderTerminal { + status_code: 200, + cancelled: false, + }, + ); + assert_eq!(facts.provider, terminal(200)); + assert_eq!( + facts.delivery, + AttemptClientDelivery::Aborted { + reason: "write failed" + } + ); + // 投递失败不是供应商的错误,摘要不该因此补 parser_error。 + assert_eq!(facts.forced_error(), None); } /// 记账层判 Success,但摘要没观察到 finish:现状会写出 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 96ae5ea30..04cd721e1 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn.rs @@ -34,9 +34,10 @@ 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, + attempt_billing_is_void, attempt_facts_for_outcome, attempt_status_code, + classify_attempt_settlement, AttemptCandidateError, AttemptCandidateStatus, + AttemptClientDelivery, AttemptProviderOutcome, AttemptSettlementInputs, AttemptTerminalFacts, + ResponsesWebSocketTurnEffect, }; use crate::ai_serving::api::StreamingStandardTerminalObserver; use crate::ai_serving::{build_openai_responses_stream_plan_from_decision, AiExecutionDecision}; @@ -61,6 +62,10 @@ const WEBSOCKET_CONNECTION_TRACE_REPORT_CONTEXT_FIELD: &str = "websocket_connect const WEBSOCKET_TURN_INDEX_REPORT_CONTEXT_FIELD: &str = "websocket_turn_index"; const WEBSOCKET_LOGICAL_TURN_ID_REPORT_CONTEXT_FIELD: &str = "websocket_logical_turn_id"; const WEBSOCKET_TURN_ATTEMPT_REPORT_CONTEXT_FIELD: &str = "websocket_turn_attempt"; +const WEBSOCKET_CLIENT_DELIVERY_REPORT_CONTEXT_FIELD: &str = "websocket_client_delivery"; +const WEBSOCKET_CLIENT_DELIVERY_ABORTED: &str = "aborted"; +const WEBSOCKET_CLIENT_DELIVERY_REASON_REPORT_CONTEXT_FIELD: &str = + "websocket_client_delivery_reason"; const DEFAULT_WEBSOCKET_FIRST_EVENT_TIMEOUT_MS: u64 = 30_000; const RESPONSES_WEBSOCKET_LIFECYCLE_STAGE_TIMEOUT: Duration = Duration::from_secs(5); @@ -237,6 +242,8 @@ pub(super) struct ResponsesProviderAttempt { /// 观察到的 provider 终态事实,与「为什么现在结算」这个信号分开保存。 /// 客户端投递失败不会把它擦掉。 provider_outcome: Option, + /// 这一个 attempt 的内容是否完整交付给了客户端。与 provider 终态正交。 + client_delivery: AttemptClientDelivery, } /// 组装一轮 turn 的 decision。 @@ -442,6 +449,7 @@ pub(super) async fn begin_responses_websocket_turn( admission: Some(admission), terminal_error_body: None, provider_outcome: None, + client_delivery: AttemptClientDelivery::Complete, }) } @@ -618,7 +626,10 @@ impl ResponsesProviderAttempt { // provider 的终态是独立事实:先记下来,之后即使客户端投递失败、 // 结算信号变成 Cancelled,这条事实也不会被擦掉。 self.provider_outcome - .get_or_insert(attempt_facts_for_outcome(None, outcome).provider); + .get_or_insert( + attempt_facts_for_outcome(None, AttemptClientDelivery::Complete, outcome) + .provider, + ); return Some(ResponsesWebSocketTurnObservation::Terminal(outcome)); } if frame.is_started() { @@ -652,6 +663,16 @@ impl ResponsesProviderAttempt { )) } + /// 记录「这一个 attempt 的内容没能完整交付给客户端」。 + /// + /// 与 provider 终态分开记录:供应商已经给出终态时,这条事实只影响 + /// candidate 的错误分类和审计里的投递标记,不作废账单。 + pub(super) fn record_client_delivery_aborted(&mut self, reason: &'static str) { + if matches!(self.client_delivery, AttemptClientDelivery::Complete) { + self.client_delivery = AttemptClientDelivery::Aborted { reason }; + } + } + pub(super) fn capture_client_frame(&mut self, event: &Value) { append_capture( &mut self.client_capture, @@ -702,7 +723,12 @@ impl ResponsesProviderAttempt { /// 拆成 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 facts = + attempt_facts_for_outcome(self.provider_outcome, self.client_delivery, outcome); + if let Some(reason) = facts.delivery.aborted_reason() { + self.report_context = + attach_client_delivery_to_report_context(self.report_context.take(), reason); + } let summary = self.finish_summary(facts); let terminal_error_body = self.terminal_error_body.take(); let outcome_reason = facts.reason().to_string(); @@ -763,6 +789,10 @@ impl ResponsesProviderAttempt { Some("websocket_cancelled".to_string()), Some(outcome_reason.clone()), ), + AttemptCandidateError::ClientDeliveryFailed => ( + Some("client_delivery_failed".to_string()), + Some(outcome_reason.clone()), + ), AttemptCandidateError::MissingTerminal => ( Some("stream_missing_terminal_event".to_string()), Some(summary.parser_error.clone().unwrap_or_else(|| { @@ -907,7 +937,9 @@ impl ResponsesProviderAttempt { summary.parser_error = Some(reason.to_string()); } } - if facts.provider.cancelled_by_provider() || facts.delivery.is_aborted() { + // 只有作废账单的那一侧才把摘要改写成 cancelled。provider 终态已到达时 + // 摘要必须保留真实的 finish_reason 和 usage,否则计费记录会被写坏。 + if attempt_billing_is_void(facts) { summary.observed_finish = true; if summary.finish_reason.is_none() { summary.finish_reason = Some("cancelled".to_string()); @@ -1045,6 +1077,30 @@ fn prepare_websocket_report_context( Value::Object(object) } +/// 在审计/用量 report context 上标记这一 attempt 的内容没能交付给客户端。 +/// +/// 只增字段,不改既有字段:账单本身按 provider 终态计,投递失败作为独立事实 +/// 留在记录里,便于事后区分「客户端拿到了」和「客户端没拿到但已计费」。 +fn attach_client_delivery_to_report_context( + report_context: Option, + reason: &str, +) -> Option { + let mut object = match report_context { + Some(Value::Object(object)) => object, + Some(other) => Map::from_iter([("seed".to_string(), other)]), + None => Map::new(), + }; + object.insert( + WEBSOCKET_CLIENT_DELIVERY_REPORT_CONTEXT_FIELD.to_string(), + Value::String(WEBSOCKET_CLIENT_DELIVERY_ABORTED.to_string()), + ); + object.insert( + WEBSOCKET_CLIENT_DELIVERY_REASON_REPORT_CONTEXT_FIELD.to_string(), + Value::String(reason.to_string()), + ); + Some(Value::Object(object)) +} + fn provider_terminal_outcome( frame: &ParsedResponsesWebSocketFrame<'_>, ) -> Option { @@ -1174,11 +1230,14 @@ mod tests { use super::super::frame::ParsedResponsesWebSocketFrame; use super::super::settlement::{ - attempt_facts_for_outcome, classify_attempt_settlement, AttemptBilling, - AttemptSettlementInputs, ResponsesWebSocketTurnEffect, + attempt_facts_for_outcome, classify_attempt_settlement, + settle_signal_for_client_delivery_failure, AttemptBilling, AttemptCandidateError, + AttemptCandidateStatus, AttemptClientDelivery, AttemptSettlementInputs, + ResponsesWebSocketTurnEffect, }; use super::{ - prepare_websocket_report_context, provider_terminal_outcome, + attach_client_delivery_to_report_context, prepare_websocket_report_context, + provider_terminal_outcome, resolve_responses_websocket_turn_timeouts, websocket_event_as_sse_line, ResponsesWebSocketTurnDeadline, ResponsesWebSocketTurnOutcome, ResponsesWebSocketTurnTimeoutPhase, @@ -1330,7 +1389,7 @@ mod tests { }) ); let outcome = outcome.expect("incomplete should end the turn"); - let facts = attempt_facts_for_outcome(None, outcome); + let facts = attempt_facts_for_outcome(None, AttemptClientDelivery::Complete, outcome); assert!(!facts.provider.cancelled_by_provider()); assert!(facts.forced_error().is_none()); assert!(!facts.provider.stream_timeout()); @@ -1388,6 +1447,139 @@ mod tests { ); } + /// relay 级:provider 终态帧已经到达,但客户端 socket 已经关闭。 + /// + /// 走的是 relay loop 写客户端失败时的完整决策链——真实帧解析 → + /// 记录 provider 事实 → 记录投递失败 → 选结算信号 → 结算表。旧实现在这里 + /// 用 client_disconnected() 覆盖 outcome,于是一条已经产出 token 的响应被 + /// 记成 void billing、不投射供应商效果、也不提交 execution report。 + #[test] + fn a_provider_terminal_that_reaches_a_closed_client_socket_is_still_billed() { + let completed = json!({ + "type": "response.completed", + "response": { + "id": "resp_ws_delivery_failed", + "model": "gpt-5.6", + "usage": {"input_tokens": 3, "output_tokens": 5, "total_tokens": 8} + } + }); + let raw = serde_json::to_string(&completed).expect("event should serialize"); + let frame = ParsedResponsesWebSocketFrame::parse(&raw).expect("event should parse"); + + // relay loop 观察到终态帧:attempt 记下 provider 事实。 + let observed = provider_terminal_outcome(&frame).expect("completed ends the turn"); + let recorded_provider = + attempt_facts_for_outcome(None, AttemptClientDelivery::Complete, observed).provider; + + // 随后写客户端失败:记录投递失败,并按 provider 终态选结算信号。 + let delivery = AttemptClientDelivery::Aborted { + reason: "gateway could not relay the provider event to the client", + }; + let signal = settle_signal_for_client_delivery_failure(Some(observed)); + assert_eq!(signal, observed, "a reached terminal must remain the signal"); + + let facts = attempt_facts_for_outcome(Some(recorded_provider), delivery, signal); + + // 终态摘要保留真实 usage 与 finish_reason,不被改写成 cancelled。 + 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(&completed)) + .expect("the terminal event 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!(summary.parser_error.is_none()); + let usage = summary + .standardized_usage + .clone() + .expect("usage must survive a delivery failure"); + assert_eq!(usage.input_tokens, 3); + assert_eq!(usage.output_tokens, 5); + + let settlement = classify_attempt_settlement(AttemptSettlementInputs { + facts, + report_represents_failure: false, + observed_finish: summary.observed_finish, + has_parser_error: summary.parser_error.is_some(), + }); + assert_eq!(settlement.billing, AttemptBilling::Billed); + assert_eq!(settlement.status_code, 200); + assert_eq!( + settlement.candidate_status, + AttemptCandidateStatus::Success + ); + assert_eq!( + settlement.provider_effect, + ResponsesWebSocketTurnEffect::ProviderSuccess + ); + assert!(settlement.submit_execution_report); + // 投递失败仍然留痕。 + assert_eq!( + settlement.candidate_error, + AttemptCandidateError::ClientDeliveryFailed + ); + } + + /// 同一条链路,但供应商还没给出终态:仍然作废账单、不提交 report。 + #[test] + fn a_closed_client_socket_before_any_terminal_still_voids_the_bill() { + let signal = settle_signal_for_client_delivery_failure(None); + let facts = attempt_facts_for_outcome( + None, + AttemptClientDelivery::Aborted { + reason: "gateway could not relay the provider event to the client", + }, + signal, + ); + let settlement = classify_attempt_settlement(AttemptSettlementInputs { + facts, + report_represents_failure: false, + observed_finish: false, + has_parser_error: false, + }); + + assert_eq!(settlement.billing, AttemptBilling::Void); + assert_eq!(settlement.status_code, 499); + assert_eq!( + settlement.candidate_status, + AttemptCandidateStatus::Cancelled + ); + assert_eq!( + settlement.provider_effect, + ResponsesWebSocketTurnEffect::ReleasePoolKeyLease + ); + assert!(!settlement.submit_execution_report); + } + + /// 投递失败只往 report context 里加字段,不改既有字段。 + #[test] + fn the_client_delivery_marker_only_adds_report_context_fields() { + let context = attach_client_delivery_to_report_context( + Some(json!({ + "request_id": "turn-2", + "websocket_mode": true, + "original_request_body": {"model": "public"}, + })), + "gateway could not relay the provider event to the client", + ) + .expect("marker should produce a report context"); + + assert_eq!(context["websocket_client_delivery"], "aborted"); + assert_eq!( + context["websocket_client_delivery_reason"], + "gateway could not relay the provider event to the client" + ); + assert_eq!(context["request_id"], "turn-2"); + assert_eq!(context["websocket_mode"], true); + assert_eq!(context["original_request_body"]["model"], "public"); + } + #[test] fn error_event_uses_the_top_level_status_code() { let event = json!({ @@ -1425,7 +1617,7 @@ 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 facts = attempt_facts_for_outcome(None, AttemptClientDelivery::Complete, outcome); let settlement = classify_attempt_settlement(AttemptSettlementInputs { facts, report_represents_failure: true, 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 eb4e42b22..2b03acf57 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 @@ -168,6 +168,19 @@ impl ResponsesTurnState { } } +impl ResponsesTurnState { + /// 记录「当前 attempt 的内容没能完整交付给客户端」。 + /// + /// 这条事实写在 attempt 上而不是 logical turn 上:结算是按 attempt 进行的, + /// 而每个 attempt 的投递结果各自独立(配额透明重试时旧 attempt 可能已经把 + /// 部分事件交付出去,新 attempt 从零开始)。 + pub(super) fn record_client_delivery_aborted(&mut self, reason: &'static str) { + if let Some(attempt) = self.attempt_mut() { + attempt.record_client_delivery_aborted(reason); + } + } +} + #[cfg(test)] mod tests { use serde_json::json; diff --git a/crates/aether-testing/integration/tests/responses_websocket_e2e.rs b/crates/aether-testing/integration/tests/responses_websocket_e2e.rs index a89bae268..da592caa6 100644 --- a/crates/aether-testing/integration/tests/responses_websocket_e2e.rs +++ b/crates/aether-testing/integration/tests/responses_websocket_e2e.rs @@ -156,12 +156,21 @@ async fn continuation_reuses_one_upstream_connection_and_bills_both_turns() -> R Ok(()) } -/// A client that walks away mid-turn must still be billed for what it started. +/// A client that walks away before the provider produced anything must settle +/// as a void row: nothing was produced, so nothing is billed. /// /// This is the path with no protocol event to announce it: the relay loop owns /// the turn, and losing the client is an exit the upstream never reports. +/// +/// The mirror case — the provider *did* reach a terminal event and only the last +/// hop to the client failed — is billed instead. That one cannot be pinned here: +/// it depends on the relay loop's `select!` observing the upstream terminal frame +/// before it observes the closed client socket, which is a race by construction. +/// It is covered deterministically by the relay-level unit tests +/// `a_provider_terminal_that_reaches_a_closed_client_socket_is_still_billed` and +/// `a_closed_client_socket_before_any_terminal_still_voids_the_bill`. #[tokio::test] -async fn client_disconnect_mid_turn_still_settles_the_usage_row() -> Result<(), BoxError> { +async fn client_disconnect_before_any_provider_output_settles_a_void_row() -> Result<(), BoxError> { let harness = Harness::start(UpstreamBehavior::StallAfterCreated).await?; let mut client = harness.connect().await?; @@ -183,6 +192,17 @@ async fn client_disconnect_mid_turn_still_settles_the_usage_row() -> Result<(), !is_pending(audit), "an abandoned turn must not be left pending: {audit:?}" ); + // The provider never emitted a terminal event, so this row stays void. + // Only a reached provider terminal survives a client delivery failure. + assert!( + !is_billed(audit), + "a turn with no provider output must not be billed: {audit:?}" + ); + assert_eq!( + audit.status, "cancelled", + "a client that left before any provider output settles as cancelled: {audit:?}" + ); + assert_eq!(audit.status_code, Some(499)); Ok(()) }