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 88672426f..f26018a6d 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/mod.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/mod.rs @@ -14,6 +14,7 @@ mod client; mod connection; mod frame; mod lifecycle; +mod observation; mod quota; mod redaction; mod relay_policy; diff --git a/apps/aether-gateway/src/handlers/proxy/websocket/responses/observation.rs b/apps/aether-gateway/src/handlers/proxy/websocket/responses/observation.rs new file mode 100644 index 000000000..9e2387d66 --- /dev/null +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/observation.rs @@ -0,0 +1,156 @@ +//! Responses WebSocket 的终态观测入口。 +//! +//! 这条传输收到的本来就是结构化的 Responses 协议事件。之前为了复用面向 SSE 的 +//! `push_line`,观测路径要先把每个事件序列化成 `data: {json}\n\n`,解析器再把它 +//! 解码回 `Value`——一次纯粹的往返,而且这个「伪 SSE」形状是随手拼的,一旦 +//! 上游事件里出现需要转义的内容,或者以后有人给拼装函数加了换行/分块逻辑, +//! 观测结果就会和真实事件悄悄分叉。 +//! +//! 现在观测走 [`StreamingStandardTerminalObserver::push_event`],直接吃 +//! `frame.protocol_events()` 借出的事件,不再序列化、不再解码。 +//! +//! **body capture 不走这条路,仍然保持 SSE 形状**(`data: {json}\n\n`): +//! `aether_usage_runtime::report` 用 `line.strip_prefix("data:")` 解析被捕获的 +//! body 来判定 `StreamCapturedTerminalState`,而它是 `stream_report_represents_failure` +//! 的一个 OR 项。把捕获内容换成结构化 JSON 会让终态判定恒为 Missing。 +//! 也就是说这一层只换「观测」,不换「捕获」——见 +//! [`super::turn::ResponsesProviderAttempt::capture_client_frame`] 一侧仍在用 +//! SSE 编码。 + +use serde_json::Value; + +use crate::ai_serving::api::StreamingStandardTerminalObserver; +use aether_contracts::ExecutionStreamTerminalSummary; + +/// 包一层 [`StreamingStandardTerminalObserver`],只暴露结构化入口。 +/// +/// 存在的意义是让「WS 不再拼 SSE」成为类型层面的事实:这里没有任何接受字节的 +/// 方法,所以不可能有人不小心把观测路径改回 `push_line`。 +#[derive(Default)] +pub(super) struct ResponsesStructuredTerminalObserver { + inner: StreamingStandardTerminalObserver, +} + +impl ResponsesStructuredTerminalObserver { + /// 观测一帧里的全部协议事件。 + /// + /// 第一个被拒绝的事件就停止推进并把摘要标成 parser_error:解析器的状态机是 + /// 有顺序的,跳过一个事件继续喂后面的只会得到更没意义的摘要。 + pub(super) fn observe_events(&mut self, report_context: &Value, events: &[&Value]) { + for event in events { + if let Err(error) = self.inner.push_event(report_context, event) { + self.inner.disable_with_error(error.to_string()); + break; + } + } + } + + pub(super) fn disable_with_error(&mut self, parser_error: impl Into) { + self.inner.disable_with_error(parser_error); + } + + pub(super) fn finish(&mut self, report_context: &Value) -> ExecutionStreamTerminalSummary { + match self.inner.finish(report_context) { + Ok(Some(summary)) => summary, + Ok(None) => ExecutionStreamTerminalSummary::default(), + Err(error) => { + self.inner.disable_with_error(error.to_string()); + self.inner.latest_summary().cloned().unwrap_or_default() + } + } + } +} + +#[cfg(test)] +mod tests { + use serde_json::json; + + use super::ResponsesStructuredTerminalObserver; + + fn report_context() -> serde_json::Value { + json!({ + "provider_api_format": "openai:responses", + "client_api_format": "openai:responses", + }) + } + + #[test] + fn structured_events_reach_the_terminal_summary_without_sse_text() { + let context = report_context(); + let created = + json!({"type": "response.created", "response": {"id": "resp_ws", "model": "gpt-5-codex"}}); + let completed = json!({ + "type": "response.completed", + "response": { + "id": "resp_ws", + "model": "gpt-5-codex", + "status": "completed", + "usage": {"input_tokens": 9, "output_tokens": 4, "total_tokens": 13}, + }, + }); + + let mut observer = ResponsesStructuredTerminalObserver::default(); + observer.observe_events(&context, &[&created, &completed]); + let summary = observer.finish(&context); + + assert!(summary.observed_finish); + assert_eq!(summary.response_id.as_deref(), Some("resp_ws")); + let usage = summary + .standardized_usage + .as_ref() + .expect("a completed response carries usage"); + assert_eq!(usage.input_tokens, 9); + assert_eq!(usage.output_tokens, 4); + assert!(summary.parser_error.is_none()); + } + + /// 批量帧里的多个事件按顺序喂入,usage 不能因为批量而丢失。 + #[test] + fn a_batched_frame_keeps_the_usage_of_its_last_event() { + let context = report_context(); + let events = [ + json!({"type": "response.created", "response": {"id": "resp_ws", "model": "m"}}), + json!({ + "type": "response.output_text.delta", + "item_id": "msg", + "output_index": 0, + "content_index": 0, + "delta": "hi", + }), + json!({ + "type": "response.completed", + "response": { + "id": "resp_ws", + "model": "m", + "status": "completed", + "usage": {"input_tokens": 3, "output_tokens": 1, "total_tokens": 4}, + }, + }), + ]; + let borrowed: Vec<&serde_json::Value> = events.iter().collect(); + + let mut observer = ResponsesStructuredTerminalObserver::default(); + observer.observe_events(&context, &borrowed); + let summary = observer.finish(&context); + + let usage = summary + .standardized_usage + .as_ref() + .expect("usage survives batching"); + assert_eq!(usage.input_tokens, 3); + assert_eq!(usage.output_tokens, 1); + assert_eq!(usage.dimensions.get("total_tokens"), Some(&json!(4))); + } + + #[test] + fn a_disabled_observer_reports_the_parser_error() { + let context = report_context(); + let mut observer = ResponsesStructuredTerminalObserver::default(); + observer.disable_with_error("upstream event was not valid JSON"); + let summary = observer.finish(&context); + assert_eq!( + summary.parser_error.as_deref(), + Some("upstream event was not valid JSON") + ); + } +} 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 b2e9e7492..046b1fa7b 100644 --- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn.rs +++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn.rs @@ -33,13 +33,13 @@ use tracing::warn; use super::adapter::ResponsesWebSocketProtocolAdapter; use super::admission::ResponsesWebSocketTurnAdmission; use super::frame::ParsedResponsesWebSocketFrame; +use super::observation::ResponsesStructuredTerminalObserver; use super::settlement::attempt_facts_for_outcome; use crate::execution_runtime::attempt_lifecycle::{ attempt_billing_is_void, AttemptBodyCapture, AttemptClientDelivery, AttemptLifecycleSeed, AttemptProviderOutcome, AttemptStageGuard, AttemptTerminalFacts, AttemptTerminalFactsInput, ExecutionAttemptLifecycle, }; -use crate::ai_serving::api::StreamingStandardTerminalObserver; use crate::ai_serving::{build_openai_responses_stream_plan_from_decision, AiExecutionDecision}; use crate::clock::current_unix_ms; use crate::control::{ @@ -227,7 +227,7 @@ pub(super) struct ResponsesProviderAttempt { lifecycle: ExecutionAttemptLifecycle, started_at: Instant, provider_headers: BTreeMap, - observer: StreamingStandardTerminalObserver, + observer: ResponsesStructuredTerminalObserver, provider_capture: AttemptBodyCapture, client_capture: AttemptBodyCapture, upstream_bytes: u64, @@ -414,7 +414,7 @@ pub(super) async fn begin_responses_websocket_turn( lifecycle, started_at: Instant::now(), provider_headers: BTreeMap::new(), - observer: StreamingStandardTerminalObserver::default(), + observer: ResponsesStructuredTerminalObserver::default(), provider_capture: AttemptBodyCapture::default(), client_capture: AttemptBodyCapture::default(), upstream_bytes: 0, @@ -572,12 +572,13 @@ impl ResponsesProviderAttempt { self.first_event_elapsed_ms = Some(elapsed_ms(self.started_at)); } - // A batched frame carries several events; the usage observer parses one - // Responses event per SSE line, so the batch must be unwrapped or its - // token usage is lost. + // 一帧可以带多个协议事件(批量帧),必须拆开:观测器按事件推进状态机, + // 整帧当一个事件喂会丢掉批量里最后那个 completed 的 usage。 let events = frame.protocol_events(); let mut report_context = self.lifecycle.take_report_context(); for event in &events { + // 观测已经走结构化入口,但捕获仍然必须是 SSE 形状:usage runtime 按 + // `data:` 行解析被捕获的 body 判定终态。 self.capture_sse_event(event); adapter.decorate_turn_report_context(&mut report_context, event); } @@ -587,15 +588,7 @@ impl ResponsesProviderAttempt { "client_api_format": "openai:responses", }); let report_context = self.lifecycle.report_context().unwrap_or(&fallback_context); - for event in &events { - if let Err(error) = self - .observer - .push_line(report_context, websocket_event_as_sse_line(event)) - { - self.observer.disable_with_error(error.to_string()); - break; - } - } + self.observer.observe_events(report_context, &events); let event_type = frame.event_type().unwrap_or_default(); if matches!(event_type, "error" | "response.failed") { @@ -735,14 +728,7 @@ impl ResponsesProviderAttempt { "client_api_format": "openai:responses", }); let report_context = self.lifecycle.report_context().unwrap_or(&fallback_context); - let mut summary = match self.observer.finish(report_context) { - Ok(Some(summary)) => summary, - Ok(None) => ExecutionStreamTerminalSummary::default(), - Err(error) => { - self.observer.disable_with_error(error.to_string()); - self.observer.latest_summary().cloned().unwrap_or_default() - } - }; + let mut summary = self.observer.finish(report_context); if let Some(reason) = facts.forced_error() { if summary.parser_error.is_none() { summary.parser_error = Some(reason.to_string()); @@ -967,7 +953,7 @@ mod tests { use aether_contracts::ExecutionTimeouts; use serde_json::json; - use crate::ai_serving::api::StreamingStandardTerminalObserver; + use super::super::observation::ResponsesStructuredTerminalObserver; use super::super::frame::ParsedResponsesWebSocketFrame; use super::super::settlement::{ @@ -1085,14 +1071,9 @@ mod tests { "provider_api_format": "openai:responses", "client_api_format": "openai:responses", }); - let mut observer = StreamingStandardTerminalObserver::default(); - observer - .push_line(&report_context, capture.into_bytes()) - .expect("WebSocket terminal event should be accepted by the usage observer"); - let summary = observer - .finish(&report_context) - .expect("WebSocket terminal observer should finish") - .expect("WebSocket terminal observer should produce a summary"); + let mut observer = ResponsesStructuredTerminalObserver::default(); + observer.observe_events(&report_context, &[&event]); + let summary = observer.finish(&report_context); let usage = summary .standardized_usage .expect("response.completed usage must reach the terminal summary"); @@ -1141,14 +1122,9 @@ mod tests { "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"); + let mut observer = ResponsesStructuredTerminalObserver::default(); + observer.observe_events(&report_context, &[&event]); + let summary = observer.finish(&report_context); assert!(summary.observed_finish); assert_eq!(summary.finish_reason.as_deref(), Some("length")); assert!(summary.parser_error.is_none()); @@ -1227,14 +1203,9 @@ mod tests { "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"); + let mut observer = ResponsesStructuredTerminalObserver::default(); + observer.observe_events(&report_context, &[&completed]); + let summary = observer.finish(&report_context); assert!(summary.observed_finish); assert!(summary.parser_error.is_none()); let usage = summary diff --git a/crates/aether-ai/formats/src/formats/openai/chat/stream.rs b/crates/aether-ai/formats/src/formats/openai/chat/stream.rs index b97dffc33..946efdf32 100644 --- a/crates/aether-ai/formats/src/formats/openai/chat/stream.rs +++ b/crates/aether-ai/formats/src/formats/openai/chat/stream.rs @@ -1264,6 +1264,11 @@ impl OpenAIResponsesProviderState { } } + /// SSE 入口:剥掉 `data:` 包装后交给 [`Self::push_event`]。 + /// + /// 解码是这个函数唯一做的事,协议状态机全在 `push_event` 里。已经持有结构化 + /// 事件的传输(Responses WebSocket)应当直接调用 `push_event`,不要为了复用 + /// 这个入口先把事件拼回 SSE 文本。 pub fn push_line( &mut self, report_context: &Value, @@ -1272,6 +1277,17 @@ impl OpenAIResponsesProviderState { let Some(value) = decode_json_data_line(&line) else { return Ok(Vec::new()); }; + self.push_event(report_context, &value) + } + + /// 结构化入口:消费一个已经解析好的 Responses 协议事件。 + /// + /// 取借用而不是所有权:持有结构化事件的传输不必为了调用它先克隆一份。 + pub fn push_event( + &mut self, + report_context: &Value, + value: &Value, + ) -> Result, AiSurfaceFinalizeError> { let mut out = Vec::new(); if let Some(response) = value.get("response").and_then(Value::as_object) { self.response_id = response @@ -1301,12 +1317,12 @@ impl OpenAIResponsesProviderState { } "response.output_text.delta" | "response.outtext.delta" => match value.get("delta") { Some(Value::String(piece)) if !piece.is_empty() => { - let key = Self::text_part_key_from_event(&value); + let key = Self::text_part_key_from_event(value); self.emit_text_delta(report_context, &mut out, key, piece); } Some(Value::Object(delta)) => { if let Some(text) = delta.get("text").and_then(Value::as_str) { - let key = Self::text_part_key_from_event(&value); + let key = Self::text_part_key_from_event(value); self.emit_missing_text(report_context, &mut out, key, text); } } @@ -1317,7 +1333,7 @@ impl OpenAIResponsesProviderState { if part.get("type").and_then(Value::as_str) == Some("output_text") { if let Some(text) = part.get("text").and_then(Value::as_str) { if !text.is_empty() { - let key = Self::text_part_key_from_event(&value); + let key = Self::text_part_key_from_event(value); self.emit_missing_text(report_context, &mut out, key, text); } } @@ -1356,7 +1372,7 @@ impl OpenAIResponsesProviderState { }) .unwrap_or_default(); if !text.is_empty() { - let key = Self::text_part_key_from_event(&value); + let key = Self::text_part_key_from_event(value); self.emit_missing_text(report_context, &mut out, key, text); } } @@ -1366,7 +1382,7 @@ impl OpenAIResponsesProviderState { .and_then(Value::as_str) .unwrap_or_default(); if !piece.is_empty() { - let key = Self::text_part_key_from_event(&value); + let key = Self::text_part_key_from_event(value); self.emit_text_delta(report_context, &mut out, key, piece); } } @@ -1383,7 +1399,7 @@ impl OpenAIResponsesProviderState { }) .unwrap_or_default(); if !refusal.is_empty() { - let key = Self::text_part_key_from_event(&value); + let key = Self::text_part_key_from_event(value); self.emit_missing_text(report_context, &mut out, key, refusal); } } @@ -1393,7 +1409,7 @@ impl OpenAIResponsesProviderState { .and_then(Value::as_str) .unwrap_or_default(); if !piece.is_empty() { - let key = Self::text_part_key_from_event(&value); + let key = Self::text_part_key_from_event(value); self.emit_text_delta(report_context, &mut out, key, piece); } } @@ -1404,7 +1420,7 @@ impl OpenAIResponsesProviderState { .or_else(|| value.get("text").and_then(Value::as_str)) .unwrap_or_default(); if !transcript.is_empty() { - let key = Self::text_part_key_from_event(&value); + let key = Self::text_part_key_from_event(value); self.emit_missing_text(report_context, &mut out, key, transcript); } } @@ -1477,7 +1493,7 @@ impl OpenAIResponsesProviderState { self.emit_output_item_event( report_context, &mut out, - &value, + value, item, output_index, false, @@ -1724,7 +1740,7 @@ impl OpenAIResponsesProviderState { self.emit_output_item_event( report_context, &mut out, - &value, + value, item, output_index, true, @@ -1742,7 +1758,7 @@ impl OpenAIResponsesProviderState { id, model, event: CanonicalStreamEvent::Finish { - finish_reason: Some(openai_responses_incomplete_finish_reason(&value)), + finish_reason: Some(openai_responses_incomplete_finish_reason(value)), usage: canonical_usage_from_openai_usage(response.get("usage")), }, }); @@ -1752,14 +1768,14 @@ impl OpenAIResponsesProviderState { event_type if openai_responses_stream_event_is_known_noop(event_type) => { self.ensure_started(report_context, &mut out); } - event_type if openai_stream_payload_is_terminal_error(&value) => { + event_type if openai_stream_payload_is_terminal_error(value) => { self.finished = true; let mut payload = value.clone(); if event_type != "response.failed" && event_type != "response.incomplete" && event_type != "error" { - payload = openai_stream_terminal_error_body(&value).unwrap_or(payload); + payload = openai_stream_terminal_error_body(value).unwrap_or(payload); if let Some(object) = payload.as_object_mut() { object.insert( "type".to_string(), diff --git a/crates/aether-ai/formats/src/formats/shared/stream_core/format_matrix.rs b/crates/aether-ai/formats/src/formats/shared/stream_core/format_matrix.rs index 7a485d2eb..091332c3e 100644 --- a/crates/aether-ai/formats/src/formats/shared/stream_core/format_matrix.rs +++ b/crates/aether-ai/formats/src/formats/shared/stream_core/format_matrix.rs @@ -251,6 +251,44 @@ impl StreamingStandardTerminalObserver { Ok(()) } + /// 结构化入口:给已经持有解析好的协议事件的传输用(Responses WebSocket), + /// 避免为了复用 [`Self::push_line`] 把事件重新拼成 `data: {json}` 再解析回来。 + /// + /// 只要 provider 的协议状态机本身接受结构化事件,这条路径与 `push_line` + /// 完全等价——`push_line` 现在就是「解码 + `push_event`」。 + /// + /// `openai:image` 的终态状态机没有结构化入口(它按 SSE 行做增量解析), + /// 这里返回 `Err`,由调用方 `disable_with_error` 把摘要标成 parser_error, + /// 而不是静默丢事件。 + pub fn push_event( + &mut self, + report_context: &Value, + event: &Value, + ) -> Result<(), AiSurfaceFinalizeError> { + self.ensure_initialized(report_context); + let Some(provider) = self.provider.as_mut() else { + return Ok(()); + }; + match provider { + TerminalStreamParser::Standard(provider) => { + let frames = provider.push_event(report_context, event)?; + let actual_service_tier = provider.actual_service_tier().map(ToOwned::to_owned); + self.observe_frames(frames); + if let Some(actual_service_tier) = actual_service_tier { + self.latest_summary + .get_or_insert_with(ExecutionStreamTerminalSummary::default) + .provider_actual_service_tier = Some(actual_service_tier); + } + } + TerminalStreamParser::OpenAIImage(_) => { + return Err(AiSurfaceFinalizeError::new( + "openai:image terminal observation has no structured event entry", + )); + } + } + Ok(()) + } + pub fn finish( &mut self, report_context: &Value, @@ -404,6 +442,24 @@ impl ProviderStreamParser { } } + /// 结构化入口。目前只有 `openai:responses` 有传输会走它(Responses + /// WebSocket);其余格式的协议状态机同样可以按「解码 + push_event」机械拆分, + /// 等到真有非 SSE 传输需要时再拆,不做无调用方的接口。 + fn push_event( + &mut self, + report_context: &Value, + event: &Value, + ) -> Result, AiSurfaceFinalizeError> { + match self { + ProviderStreamParser::OpenAIResponses(state) => state.push_event(report_context, event), + ProviderStreamParser::OpenAIChat(_) + | ProviderStreamParser::Claude(_) + | ProviderStreamParser::Gemini(_) => Err(AiSurfaceFinalizeError::new( + "this provider stream parser has no structured event entry", + )), + } + } + fn finish( &mut self, report_context: &Value, @@ -2490,3 +2546,257 @@ mod tests { ); } } + +#[cfg(test)] +mod structured_entry_tests { + use super::StreamingStandardTerminalObserver; + use aether_contracts::ExecutionStreamTerminalSummary; + use serde_json::{json, Value}; + + fn report_context() -> Value { + json!({ + "provider_api_format": "openai:responses", + "client_api_format": "openai:responses", + "mapped_model": "gpt-5-codex", + }) + } + + /// 用 SSE 入口观测一组事件。这是 C5 之前 WebSocket 走的路径:把结构化事件 + /// 拼成 `data: {json}` 再交给解析器。 + fn summary_via_push_line(events: &[Value]) -> ExecutionStreamTerminalSummary { + let context = report_context(); + let mut observer = StreamingStandardTerminalObserver::default(); + for event in events { + observer + .push_line(&context, format!("data: {event}\n\n").into_bytes()) + .expect("the SSE entry must accept these events"); + } + observer + .finish(&context) + .expect("the observer must finish") + .unwrap_or_default() + } + + /// 用结构化入口观测同一组事件。这是 C5 之后的路径。 + fn summary_via_push_event(events: &[Value]) -> ExecutionStreamTerminalSummary { + let context = report_context(); + let mut observer = StreamingStandardTerminalObserver::default(); + for event in events { + observer + .push_event(&context, event) + .expect("the structured entry must accept these events"); + } + observer + .finish(&context) + .expect("the observer must finish") + .unwrap_or_default() + } + + fn assert_entries_agree(label: &str, events: &[Value]) { + let via_line = summary_via_push_line(events); + let via_event = summary_via_push_event(events); + assert_eq!( + via_line, via_event, + "the SSE entry and the structured entry must produce identical summaries for {label}" + ); + } + + fn created() -> Value { + json!({"type": "response.created", "response": {"id": "resp_diff", "model": "gpt-5-codex"}}) + } + + fn text_delta(piece: &str) -> Value { + json!({ + "type": "response.output_text.delta", + "item_id": "msg_diff", + "output_index": 0, + "content_index": 0, + "delta": piece, + }) + } + + /// 批量事件:WS 一帧可以带多个协议事件,逐个喂入的结果必须和逐行喂入一致。 + #[test] + fn a_batched_delta_sequence_agrees_across_both_entries() { + let events = vec![ + created(), + text_delta("he"), + text_delta("ll"), + text_delta("o"), + json!({ + "type": "response.completed", + "response": { + "id": "resp_diff", + "model": "gpt-5-codex", + "status": "completed", + "usage": {"input_tokens": 11, "output_tokens": 3, "total_tokens": 14}, + }, + }), + ]; + assert_entries_agree("a batched delta sequence", &events); + let summary = summary_via_push_event(&events); + assert!(summary.observed_finish); + assert_eq!(summary.response_id.as_deref(), Some("resp_diff")); + let usage = summary + .standardized_usage + .as_ref() + .expect("completed carries usage"); + assert_eq!(usage.input_tokens, 11); + assert_eq!(usage.output_tokens, 3); + } + + /// 合法 `response.incomplete`:C1 定过的语义(终态、可计费),两条入口必须 + /// 得到同一个摘要,尤其是 finish_reason 与 parser_error 的取值。 + #[test] + fn a_legitimate_incomplete_agrees_across_both_entries() { + let events = vec![ + created(), + text_delta("partial"), + json!({ + "type": "response.incomplete", + "response": { + "id": "resp_diff", + "model": "gpt-5-codex", + "status": "incomplete", + "incomplete_details": {"reason": "max_output_tokens"}, + "usage": {"input_tokens": 7, "output_tokens": 5, "total_tokens": 12}, + }, + }), + ]; + assert_entries_agree("a legitimate incomplete", &events); + let summary = summary_via_push_event(&events); + assert!(summary.observed_finish); + assert!( + summary.parser_error.is_none(), + "a legitimate incomplete is not a parser error: {:?}", + summary.parser_error + ); + } + + #[test] + fn a_terminal_error_agrees_across_both_entries() { + let events = vec![ + created(), + json!({ + "type": "error", + "error": {"type": "server_error", "message": "upstream exploded"}, + }), + ]; + assert_entries_agree("a terminal error", &events); + let summary = summary_via_push_event(&events); + assert!(summary.observed_finish); + assert_eq!(summary.finish_reason.as_deref(), Some("error")); + assert!(summary.parser_error.is_some()); + } + + #[test] + fn a_response_failed_event_agrees_across_both_entries() { + assert_entries_agree( + "a response.failed event", + &[ + created(), + json!({ + "type": "response.failed", + "response": { + "id": "resp_diff", + "model": "gpt-5-codex", + "status": "failed", + "error": {"type": "server_error", "message": "generation failed"}, + }, + }), + ], + ); + } + + /// 未知事件只增计数、不改终态判定,两条入口的计数必须一致。 + #[test] + fn unknown_events_agree_across_both_entries() { + let events = vec![ + created(), + json!({"type": "response.some_future_event", "payload": {"anything": true}}), + json!({"type": "response.another_future_event"}), + json!({ + "type": "response.completed", + "response": { + "id": "resp_diff", + "model": "gpt-5-codex", + "status": "completed", + "usage": {"input_tokens": 1, "output_tokens": 1, "total_tokens": 2}, + }, + }), + ]; + assert_entries_agree("unknown events", &events); + let summary = summary_via_push_event(&events); + assert!(summary.observed_finish); + assert!(summary.unknown_event_count > 0, "unknown events are counted"); + } + + /// 供应商声明的 service tier 通过两条入口都要落到摘要上。 + #[test] + fn a_service_tier_agrees_across_both_entries() { + assert_entries_agree( + "a declared service tier", + &[ + json!({ + "type": "response.created", + "response": { + "id": "resp_diff", + "model": "gpt-5-codex", + "service_tier": "priority", + }, + }), + json!({ + "type": "response.completed", + "response": { + "id": "resp_diff", + "model": "gpt-5-codex", + "status": "completed", + "service_tier": "priority", + "usage": {"input_tokens": 2, "output_tokens": 2, "total_tokens": 4}, + }, + }), + ], + ); + } + + /// 没有任何供应商终态事件。注意 `finish()` 会补一个 `stop`(这是 HTTP 与 + /// WebSocket 共享的既有行为,C5 不改),所以 `observed_finish` 为真而 usage + /// 缺失——真正的「缺终态」判定看的是 usage 与被捕获的 body。这里要钉住的是 + /// 两条入口在这种不完整序列上仍然给出同一个摘要。 + #[test] + fn a_missing_terminal_agrees_across_both_entries() { + let events = vec![created(), text_delta("truncated")]; + assert_entries_agree("a missing terminal", &events); + let summary = summary_via_push_event(&events); + assert_eq!(summary.finish_reason.as_deref(), Some("stop")); + assert!( + summary.standardized_usage.is_none(), + "a synthesized finish carries no usage" + ); + } + + /// `openai:image` 没有结构化入口:必须显式报错,让调用方标记 parser_error, + /// 而不是静默丢掉事件、把摘要留成「未观察到终态」。 + #[test] + fn the_image_format_rejects_the_structured_entry() { + let context = json!({ + "provider_api_format": "openai:image", + "client_api_format": "openai:image", + "mapped_model": "gpt-image-1", + }); + let mut observer = StreamingStandardTerminalObserver::default(); + let error = observer + .push_event(&context, &json!({"type": "image_generation.completed"})) + .expect_err("openai:image has no structured entry"); + assert!( + error.to_string().contains("structured event entry"), + "the error must name the missing entry: {error}" + ); + + observer.disable_with_error(error.to_string()); + let summary = observer + .latest_summary() + .expect("disable_with_error records a summary"); + assert!(summary.parser_error.is_some()); + } +}