From 66d6c17d2d80ddebf05c682f06793818f3b8158d Mon Sep 17 00:00:00 2001 From: ZheFox <77232781+zhefox@users.noreply.github.com> Date: Fri, 4 Sep 2026 16:36:40 +0800 Subject: [PATCH] fix antigravity reasoning streaming and terminal errors --- .../execution_runtime/stream/commit_policy.rs | 14 +- .../src/execution_runtime/stream/execution.rs | 33 +++- .../formats/gemini/generate_content/stream.rs | 131 ++++++++++++++ .../shared/stream_core/format_matrix.rs | 162 +++++++++++++++++- 4 files changed, 319 insertions(+), 21 deletions(-) diff --git a/apps/aether-gateway/src/execution_runtime/stream/commit_policy.rs b/apps/aether-gateway/src/execution_runtime/stream/commit_policy.rs index 3cbf9f641..9cfd0cae1 100644 --- a/apps/aether-gateway/src/execution_runtime/stream/commit_policy.rs +++ b/apps/aether-gateway/src/execution_runtime/stream/commit_policy.rs @@ -455,7 +455,10 @@ fn gemini_part_is_client_semantic(part: &Value) -> bool { return true; } if part.get("thought").and_then(Value::as_bool) == Some(true) { - return false; + return part + .get("text") + .and_then(Value::as_str) + .is_some_and(|text| !text.is_empty()); } if part.keys().all(|key| key == "thoughtSignature") { return false; @@ -584,17 +587,12 @@ mod tests { } #[test] - fn gemini_gate_waits_through_thought_and_commits_on_text() { + fn gemini_gate_commits_on_first_nonempty_thought() { let mut gate = StreamCommitGate::new(gemini_policy()); let thought = b"data: {\"response\":{\"candidates\":[{\"content\":{\"role\":\"model\",\"parts\":[{\"thought\":true,\"text\":\"checking\"}]}}]}}\n\n"; - let text = b"data: {\"response\":{\"candidates\":[{\"content\":{\"role\":\"model\",\"parts\":[{\"text\":\"answer\"}]}}]}}\n\n"; assert_eq!( gate.observe_provider_bytes(thought), - StreamPrecommitObservation::Pending - ); - assert_eq!( - gate.observe_provider_bytes(text), StreamPrecommitObservation::Commit ); assert_eq!(gate.state(), StreamCommitState::Committed); @@ -615,7 +613,7 @@ mod tests { #[test] fn gemini_gate_rejects_malformed_function_call_before_commit() { let mut gate = StreamCommitGate::new(gemini_policy()); - let thought = b"data: {\"response\":{\"candidates\":[{\"content\":{\"role\":\"model\",\"parts\":[{\"thought\":true,\"text\":\"calling\"}]}}]}}\n\n"; + let thought = b"data: {\"response\":{\"candidates\":[{\"content\":{\"role\":\"model\",\"parts\":[{\"thoughtSignature\":\"signature\",\"text\":\"\"}]}}]}}\n\n"; let malformed = b"data: {\"response\":{\"candidates\":[{\"content\":{\"role\":\"model\",\"parts\":[{\"thoughtSignature\":\"signature\",\"text\":\"\"}]},\"finishReason\":\"MALFORMED_FUNCTION_CALL\",\"finishMessage\":\"Malformed function call: Function call is empty - no input to parse.\"}]}}\n\n"; assert_eq!( diff --git a/apps/aether-gateway/src/execution_runtime/stream/execution.rs b/apps/aether-gateway/src/execution_runtime/stream/execution.rs index 800c09d12..536326765 100644 --- a/apps/aether-gateway/src/execution_runtime/stream/execution.rs +++ b/apps/aether-gateway/src/execution_runtime/stream/execution.rs @@ -10561,7 +10561,7 @@ mod tests { } #[tokio::test] - async fn malformed_antigravity_function_call_retries_before_stream_commit() { + async fn malformed_antigravity_function_call_streams_thought_then_fails_in_band() { let request_id = "req-antigravity-malformed-function-call"; let plan = antigravity_gemini_stream_plan(request_id); let provider_catalog = provider_catalog_for_plan( @@ -10640,10 +10640,35 @@ mod tests { None, ) .await - .expect("malformed Antigravity stream should resolve through failover"); + .expect("malformed Antigravity stream should return a client stream") + .expect("the first reasoning delta should commit the selected candidate"); - assert!(response.is_none()); - assert_eq!(retry_scope, AiAttemptRetryScope::Candidate); + assert_eq!(response.status(), StatusCode::OK); + let body = to_bytes(response.into_body(), usize::MAX) + .await + .expect("response body should read"); + let body = String::from_utf8(body.to_vec()).expect("response body should be utf8"); + assert!( + body.contains("event: response.reasoning_summary_text.delta\n"), + "{body}" + ); + assert!( + body.contains("\"delta\":\"Validating the document.\""), + "{body}" + ); + assert!(body.contains("event: response.failed\n"), "{body}"); + assert!( + body.contains("\"code\":\"MALFORMED_FUNCTION_CALL\""), + "{body}" + ); + assert!( + body.contains( + "\"message\":\"Malformed function call: Function call is empty - no input to parse.\"" + ), + "{body}" + ); + assert!(!body.contains("unsupported_finish_reason"), "{body}"); + assert_eq!(retry_scope, AiAttemptRetryScope::Provider); } fn tunnel_proxy_snapshot(base_url: String) -> aether_contracts::ProxySnapshot { diff --git a/crates/aether-ai/formats/src/formats/gemini/generate_content/stream.rs b/crates/aether-ai/formats/src/formats/gemini/generate_content/stream.rs index 80d04e4af..5f5875ab2 100644 --- a/crates/aether-ai/formats/src/formats/gemini/generate_content/stream.rs +++ b/crates/aether-ai/formats/src/formats/gemini/generate_content/stream.rs @@ -104,10 +104,25 @@ impl GeminiProviderState { let Some(candidate_object) = candidate.as_object() else { continue; }; + let (response_id, response_model) = self.identity(report_context); + let terminal_error = gemini_stream_terminal_error_payload( + candidate_object, + response_id.as_str(), + response_model.as_str(), + event_object.get("usageMetadata"), + ); let Some(content) = candidate_object.get("content").and_then(Value::as_object) else { + if let Some(payload) = terminal_error { + out.push(self.unknown_frame(report_context, payload)); + self.finished = true; + } continue; }; let Some(parts) = content.get("parts").and_then(Value::as_array) else { + if let Some(payload) = terminal_error { + out.push(self.unknown_frame(report_context, payload)); + self.finished = true; + } continue; }; if !parts.is_empty() { @@ -306,6 +321,11 @@ impl GeminiProviderState { }); } } + if let Some(payload) = terminal_error { + out.push(self.unknown_frame(report_context, payload)); + self.finished = true; + continue; + } if let Some(finish_reason) = candidate_object.get("finishReason").and_then(Value::as_str) { @@ -350,6 +370,60 @@ impl GeminiProviderState { } } +fn gemini_stream_terminal_error_payload( + candidate: &Map, + response_id: &str, + model: &str, + usage_metadata: Option<&Value>, +) -> Option { + let finish_reason = candidate + .get("finishReason") + .or_else(|| candidate.get("finish_reason")) + .and_then(Value::as_str) + .map(str::trim) + .filter(|value| { + matches!( + *value, + "MALFORMED_FUNCTION_CALL" + | "UNEXPECTED_TOOL_CALL" + | "TOO_MANY_TOOL_CALLS" + | "MISSING_THOUGHT_SIGNATURE" + | "MALFORMED_RESPONSE" + ) + })?; + let message = candidate + .get("finishMessage") + .or_else(|| candidate.get("finish_message")) + .and_then(Value::as_str) + .map(str::trim) + .filter(|value| !value.is_empty()) + .map(ToOwned::to_owned) + .unwrap_or_else(|| format!("Gemini stream ended with {finish_reason}")); + + let mut response = json!({ + "id": response_id, + "object": "response", + "model": model, + "status": "failed", + "error": { + "type": "upstream_gemini_finish_error", + "code": finish_reason, + "message": message, + "upstream_status": 200 + } + }); + if let Some(usage) = canonical_usage_from_gemini_usage(usage_metadata) + .map(|usage| openai_responses_usage_from_usage(&usage)) + { + response["usage"] = usage; + } + + Some(json!({ + "type": "response.failed", + "response": response + })) +} + fn map_gemini_stream_finish_reason(value: &str) -> Option<&str> { match value { "STOP" => Some("stop"), @@ -1011,6 +1085,63 @@ mod tests { ))); } + #[test] + fn gemini_provider_state_emits_terminal_error_for_malformed_function_call() { + let mut state = GeminiProviderState::default(); + let report_context = json!({}); + let frames = state + .push_line( + &report_context, + data_line(json!({ + "response": { + "responseId": "resp_malformed_tool_call", + "modelVersion": "gemini-3.7-flash-tiered", + "candidates": [{ + "index": 0, + "content": { + "role": "model", + "parts": [{ + "text": "", + "thoughtSignature": "opaque-thought-signature" + }] + }, + "finishReason": "MALFORMED_FUNCTION_CALL", + "finishMessage": "Malformed function call: Function call is empty - no input to parse." + }], + "usageMetadata": { + "promptTokenCount": 206744, + "cachedContentTokenCount": 203947, + "thoughtsTokenCount": 1130, + "totalTokenCount": 207874 + } + } + })), + ) + .expect("malformed function call terminal should parse"); + + assert!(frames.iter().any(|frame| matches!( + &frame.event, + CanonicalStreamEvent::UnknownEvent(payload) + if payload["type"] == "response.failed" + && payload["response"]["status"] == "failed" + && payload["response"]["id"] == "resp_malformed_tool_call" + && payload["response"]["model"] == "gemini-3.7-flash-tiered" + && payload["response"]["error"]["code"] == "MALFORMED_FUNCTION_CALL" + && payload["response"]["error"]["message"] + == "Malformed function call: Function call is empty - no input to parse." + && payload["response"]["usage"]["input_tokens"] == 206744 + && payload["response"]["usage"]["output_tokens"] == 1130 + && payload["response"]["usage"]["total_tokens"] == 207874 + ))); + assert!(!frames + .iter() + .any(|frame| matches!(frame.event, CanonicalStreamEvent::Finish { .. }))); + assert!(state + .finish(&report_context) + .expect("finished error stream should not synthesize success") + .is_empty()); + } + #[test] fn gemini_provider_state_parses_function_response_as_tool_result() { let mut state = GeminiProviderState::default(); 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 4f5e86d6f..6983f7e94 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 @@ -17,8 +17,9 @@ use crate::formats::shared::error_body::{ }; use crate::formats::shared::sse::encode_json_sse; use crate::formats::shared::stream_core::common::{ - decode_json_data_line, openai_stream_terminal_error_body, openai_stream_terminal_error_message, - unsupported_stream_event_message, CanonicalStreamEvent, CanonicalStreamFrame, CanonicalUsage, + canonical_usage_from_openai_usage, decode_json_data_line, openai_stream_terminal_error_body, + openai_stream_terminal_error_message, unsupported_stream_event_message, CanonicalStreamEvent, + CanonicalStreamFrame, CanonicalUsage, }; use crate::formats::shared::AiSurfaceFinalizeError; @@ -134,7 +135,11 @@ impl StreamingStandardFormatMatrix { } if let CanonicalStreamEvent::UnknownEvent(payload) = &frame.event { self.terminated = true; - out.extend(client.emit_unknown_event(payload)?); + if openai_stream_terminal_error_body(payload).is_some() { + out.extend(client.emit_terminal_error_frame(frame)?); + } else { + out.extend(client.emit_unknown_event(payload)?); + } break; } if let CanonicalStreamEvent::OpenAiResponsesOutputItem { raw_event, .. } = &frame.event @@ -368,6 +373,10 @@ impl StreamingStandardTerminalObserver { summary.observed_finish = true; summary.finish_reason = Some("error".to_string()); summary.parser_error = openai_stream_terminal_error_message(&payload); + summary.standardized_usage = payload + .pointer("/response/usage") + .and_then(|usage| canonical_usage_from_openai_usage(Some(usage))) + .map(standardized_usage_from_canonical); } CanonicalStreamEvent::UnknownEvent(_) => { summary.unknown_event_count = summary.unknown_event_count.saturating_add(1); @@ -613,6 +622,44 @@ impl ClientStreamEmitter { self.emit_error(error_body) } + fn emit_terminal_error_frame( + &mut self, + frame: CanonicalStreamFrame, + ) -> Result, AiSurfaceFinalizeError> { + if matches!( + self, + ClientStreamEmitter::OpenAIChat(_) | ClientStreamEmitter::OpenAIResponses(_) + ) { + return self.emit(frame); + } + let CanonicalStreamEvent::UnknownEvent(payload) = frame.event else { + return self.emit(frame); + }; + let Some(source_error_body) = openai_stream_terminal_error_body(&payload) else { + return self.emit_unknown_event(&payload); + }; + let Some(error) = source_error_body.get("error") else { + return self.emit_unknown_event(&payload); + }; + let message = error + .get("message") + .and_then(Value::as_str) + .unwrap_or("Upstream stream ended with an error"); + let code = error.get("code").and_then(|value| match value { + Value::String(value) => Some(value.as_str()), + _ => None, + }); + let Some(error_body) = build_core_error_body_for_client_format( + self.api_format(), + message, + code, + LocalCoreSyncErrorKind::ServerError, + ) else { + return Ok(Vec::new()); + }; + self.emit_error(error_body) + } + fn emit_unsupported_finish_reason( &mut self, finish_reason: &str, @@ -811,7 +858,13 @@ mod tests { }, "finishReason": "MALFORMED_FUNCTION_CALL", "finishMessage": "Malformed function call: Function call is empty - no input to parse." - }] + }], + "usageMetadata": { + "promptTokenCount": 206744, + "cachedContentTokenCount": 203947, + "thoughtsTokenCount": 1130, + "totalTokenCount": 207874 + } }, "responseId": "resp_malformed_tool_call" })), @@ -824,14 +877,105 @@ mod tests { .expect("Gemini terminal frame should produce a summary"); assert!(summary.observed_finish); - assert_eq!( - summary.finish_reason.as_deref(), - Some("MALFORMED_FUNCTION_CALL") - ); + assert_eq!(summary.finish_reason.as_deref(), Some("error")); assert_eq!( summary.parser_error.as_deref(), - Some("unsupported provider stream finish reason: MALFORMED_FUNCTION_CALL") + Some("Malformed function call: Function call is empty - no input to parse.") ); + let usage = summary + .standardized_usage + .expect("failed Gemini terminal should preserve usage"); + assert_eq!(usage.input_tokens, 206744); + assert_eq!(usage.output_tokens, 1130); + assert_eq!(usage.cache_read_tokens, 203947); + } + + #[test] + fn streams_gemini_thought_text_to_openai_responses_immediately() { + let context = report_context("gemini:generate_content", "openai:responses"); + let mut matrix = StreamingStandardFormatMatrix::default(); + let output = matrix + .transform_line( + &context, + data_line(json!({ + "response": { + "responseId": "resp_reasoning_123", + "modelVersion": "gemini-3.7-flash-tiered", + "candidates": [{ + "index": 0, + "content": { + "role": "model", + "parts": [{"thought": true, "text": "checking"}] + } + }] + } + })), + ) + .expect("first Gemini thought chunk should transform"); + let sse = String::from_utf8(output).expect("reasoning SSE should be utf8"); + + assert!( + sse.contains("event: response.reasoning_summary_text.delta\n"), + "{sse}" + ); + assert!(sse.contains("\"delta\":\"checking\""), "{sse}"); + } + + #[test] + fn transforms_malformed_gemini_function_call_to_responses_failed() { + let context = report_context("gemini:generate_content", "openai:responses"); + let mut matrix = StreamingStandardFormatMatrix::default(); + let output = matrix + .transform_line( + &context, + data_line(json!({ + "response": { + "responseId": "resp_malformed_tool_call", + "modelVersion": "gemini-3.7-flash-tiered", + "candidates": [{ + "index": 0, + "content": { + "role": "model", + "parts": [{ + "text": "", + "thoughtSignature": "opaque-thought-signature" + }] + }, + "finishReason": "MALFORMED_FUNCTION_CALL", + "finishMessage": "Malformed function call: Function call is empty - no input to parse." + }], + "usageMetadata": { + "promptTokenCount": 206744, + "cachedContentTokenCount": 203947, + "thoughtsTokenCount": 1130, + "totalTokenCount": 207874 + } + } + })), + ) + .expect("malformed Gemini terminal should transform to a stream error"); + let sse = String::from_utf8(output).expect("failed response SSE should be utf8"); + + assert!(sse.contains("event: response.failed\n"), "{sse}"); + assert!(sse.contains("\"type\":\"response.failed\""), "{sse}"); + assert!( + sse.contains("\"code\":\"MALFORMED_FUNCTION_CALL\""), + "{sse}" + ); + assert!( + sse.contains( + "\"message\":\"Malformed function call: Function call is empty - no input to parse.\"" + ), + "{sse}" + ); + assert!(sse.contains("\"input_tokens\":206744"), "{sse}"); + assert!(sse.contains("\"output_tokens\":1130"), "{sse}"); + assert!(sse.contains("\"cached_tokens\":203947"), "{sse}"); + assert!(!sse.contains("unsupported_finish_reason"), "{sse}"); + assert!(matrix + .finish(&context) + .expect("failed matrix should stay terminated") + .is_empty()); } #[test]