diff --git a/apps/aether-gateway/src/execution_runtime/submission.rs b/apps/aether-gateway/src/execution_runtime/submission.rs index 3ee844002..ff1d72fcf 100644 --- a/apps/aether-gateway/src/execution_runtime/submission.rs +++ b/apps/aether-gateway/src/execution_runtime/submission.rs @@ -196,6 +196,78 @@ fn maybe_build_invalid_provider_success_finalize_response( )?)) } +fn local_sync_needs_conversion(payload: &GatewaySyncReportRequest) -> bool { + payload + .report_context + .as_ref() + .and_then(|value| value.get("needs_conversion")) + .and_then(|value| value.as_bool()) + .unwrap_or(false) +} + +/// A successful upstream response that needed conversion but could not be +/// converted must not reach the client in the provider's own format. +fn maybe_build_unconverted_cross_format_success_response( + trace_id: &str, + decision: &GatewayControlDecision, + payload: &GatewaySyncReportRequest, +) -> Result>, GatewayError> { + if payload.status_code >= 400 + || !local_sync_needs_conversion(payload) + || !is_core_error_finalize_kind(payload.report_kind.as_str()) + { + return Ok(None); + } + + let client_api_format = resolve_local_sync_client_api_format(payload); + let provider_api_format = resolve_local_sync_provider_api_format(payload); + warn!( + event_name = "local_core_finalize_cross_format_success_unconverted", + log_type = "event", + trace_id = %trace_id, + report_kind = %payload.report_kind, + status_code = payload.status_code, + client_api_format = %client_api_format, + provider_api_format = %provider_api_format, + "gateway could not convert a successful provider response to the client format" + ); + let message = format!( + "Provider returned HTTP {} but its {provider_api_format} response could not be converted to {client_api_format}.", + payload.status_code + ); + let body_json = build_core_error_body_for_client_format( + &client_api_format, + &message, + Some("response_conversion_failed"), + LocalCoreSyncErrorKind::ServerError, + ) + .unwrap_or_else(|| { + serde_json::json!({ + "error": { + "message": message, + "type": "server_error", + "code": "response_conversion_failed" + } + }) + }); + + let mut response_headers = payload.headers.clone(); + response_headers.remove("content-encoding"); + response_headers.remove("content-length"); + response_headers.insert("content-type".to_string(), "application/json".to_string()); + let body_bytes = + serde_json::to_vec(&body_json).map_err(|err| GatewayError::Internal(err.to_string()))?; + response_headers.insert("content-length".to_string(), body_bytes.len().to_string()); + + Ok(Some(build_client_response_from_parts( + StatusCode::BAD_GATEWAY.as_u16(), + &response_headers, + Body::from(body_bytes), + trace_id, + Some(decision), + )?)) +} + fn local_core_sync_finalize_has_invalid_provider_success( payload: &GatewaySyncReportRequest, ) -> Result { @@ -274,6 +346,12 @@ pub(crate) fn resolve_local_core_error_response_body_json( return Ok(Some(body_json)); } + // A 2xx cross-format body that is not JSON (e.g. an aggregated SSE capture) + // carries no upstream error; wrapping it as one would ship raw provider + // bytes to the client under the success status. + if payload.status_code < 400 && local_sync_needs_conversion(payload) { + return Ok(None); + } let Some(body_text) = decode_local_sync_body_text(payload)? else { return Ok(None); }; @@ -626,6 +704,10 @@ pub(crate) async fn submit_local_core_error_or_sync_finalize( maybe_build_local_core_error_response(trace_id, decision, &payload)? { response + } else if let Some(response) = + maybe_build_unconverted_cross_format_success_response(trace_id, decision, &payload)? + { + response } else { warn!( event_name = "local_core_finalize_fallback_raw_response_body", @@ -937,6 +1019,128 @@ mod tests { ); } + #[tokio::test] + async fn local_core_sync_finalize_converts_forced_responses_stream_for_gemini_client() { + use base64::Engine as _; + + // Forced-stream xAI shape: the terminal response echoes request + // metadata and encrypted reasoning next to the real answer. + let raw_sse = concat!( + "event: response.created\n", + "data: {\"type\":\"response.created\",\"sequence_number\":0,\"response\":{\"id\":\"resp_xai_123\",\"object\":\"response\",\"status\":\"in_progress\",\"model\":\"grok-4.7-build\",\"output\":[],\"parallel_tool_calls\":true,\"tools\":[]}}\n\n", + "event: response.output_item.done\n", + "data: {\"type\":\"response.output_item.done\",\"sequence_number\":1,\"output_index\":0,\"item\":{\"id\":\"rs_xai_123\",\"type\":\"reasoning\",\"status\":\"completed\",\"summary\":[],\"encrypted_content\":\"opaque-xai-reasoning\"}}\n\n", + "event: response.output_text.delta\n", + "data: {\"type\":\"response.output_text.delta\",\"sequence_number\":2,\"item_id\":\"msg_xai_123\",\"output_index\":1,\"content_index\":0,\"delta\":\"Hi there, friend\"}\n\n", + "event: response.output_item.done\n", + "data: {\"type\":\"response.output_item.done\",\"sequence_number\":3,\"output_index\":1,\"item\":{\"id\":\"msg_xai_123\",\"type\":\"message\",\"status\":\"completed\",\"role\":\"assistant\",\"content\":[{\"type\":\"output_text\",\"text\":\"Hi there, friend\",\"annotations\":[]}]}}\n\n", + "event: response.completed\n", + "data: {\"type\":\"response.completed\",\"sequence_number\":4,\"response\":{\"id\":\"resp_xai_123\",\"object\":\"response\",\"status\":\"completed\",\"model\":\"grok-4.7-build\",\"output\":[],\"parallel_tool_calls\":true,\"tool_choice\":\"auto\",\"tools\":[],\"text\":{\"format\":{\"type\":\"text\"}},\"temperature\":0.7,\"store\":false,\"usage\":{\"input_tokens\":1249,\"output_tokens\":12,\"total_tokens\":1261}}}\n\n", + ); + let mut payload = core_finalize_payload( + "gemini_chat_sync_finalize", + "gemini:generate_content", + "openai:responses", + 200, + json!(null), + ); + payload.body_json = None; + payload.body_base64 = Some(base64::engine::general_purpose::STANDARD.encode(raw_sse)); + payload.report_context = Some(json!({ + "client_api_format": "gemini:generate_content", + "provider_api_format": "openai:responses", + "provider_stream_event_api_format": "openai:responses", + "model": "grok-4.7", + "mapped_model": "grok-4.7", + "needs_conversion": true, + })); + + let state = AppState::new().expect("state should build"); + let response = submit_local_core_error_or_sync_finalize( + &state, + "trace-forced-responses-gemini", + &test_decision(), + payload, + ) + .await + .expect("finalize should build a response"); + + assert_eq!(response.status(), http::StatusCode::OK); + let body_bytes = to_bytes(response.into_body(), usize::MAX) + .await + .expect("body should read"); + let body = + serde_json::from_slice::(&body_bytes).expect("body should decode"); + assert!(body.get("error").is_none(), "unexpected error body: {body}"); + let parts = body["candidates"][0]["content"]["parts"] + .as_array() + .expect("gemini parts"); + assert!(parts.iter().any(|part| part["text"] == "Hi there, friend")); + let text = String::from_utf8_lossy(&body_bytes); + assert!(!text.contains("opaque-xai-reasoning") && !text.contains("response.created")); + } + + #[tokio::test] + async fn local_core_sync_finalize_never_wraps_unconvertible_success_sse_as_client_error() { + use base64::Engine as _; + + // A complete stream whose output the Gemini client cannot represent. + let raw_sse = concat!( + "event: response.output_item.done\n", + "data: {\"type\":\"response.output_item.done\",\"output_index\":0,\"item\":{\"id\":\"future_item_123\",\"type\":\"future_output\",\"payload\":\"must-not-drop\"}}\n\n", + "event: response.completed\n", + "data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_raw_123\",\"object\":\"response\",\"status\":\"completed\",\"model\":\"grok-4.7\",\"output\":[]}}\n\n", + ); + let mut payload = core_finalize_payload( + "gemini_chat_sync_finalize", + "gemini:generate_content", + "openai:responses", + 200, + json!(null), + ); + payload.body_json = None; + payload.body_base64 = Some(base64::engine::general_purpose::STANDARD.encode(raw_sse)); + payload.report_context = Some(json!({ + "client_api_format": "gemini:generate_content", + "provider_api_format": "openai:responses", + "provider_stream_event_api_format": "openai:responses", + "needs_conversion": true, + })); + + assert!(maybe_build_local_core_error_response( + "trace-raw-success-sse", + &test_decision(), + &payload, + ) + .expect("response build should not error") + .is_none()); + + let state = AppState::new().expect("state should build"); + let response = submit_local_core_error_or_sync_finalize( + &state, + "trace-raw-success-sse", + &test_decision(), + payload, + ) + .await + .expect("finalize should build a response"); + + assert_eq!(response.status(), http::StatusCode::BAD_GATEWAY); + let body = serde_json::from_slice::( + &to_bytes(response.into_body(), usize::MAX) + .await + .expect("body should read"), + ) + .expect("body should decode"); + let message = body["error"]["message"] + .as_str() + .expect("error message should exist"); + assert!( + message.contains("could not be converted") && !message.contains("must-not-drop"), + "unexpected message: {message}" + ); + } + #[tokio::test] async fn submit_local_core_finalize_keeps_http_200_for_success_image_body() { let payload = core_finalize_payload( diff --git a/apps/aether-gateway/src/tests/ai_execute/finalize_local_cli/cross_format.rs b/apps/aether-gateway/src/tests/ai_execute/finalize_local_cli/cross_format.rs index f199ee5f3..795829d64 100644 --- a/apps/aether-gateway/src/tests/ai_execute/finalize_local_cli/cross_format.rs +++ b/apps/aether-gateway/src/tests/ai_execute/finalize_local_cli/cross_format.rs @@ -893,8 +893,8 @@ async fn gateway_executes_openai_responses_cross_format_function_call_upstream_s }, { "type": "function_call", - "id": "call_auto_1", - "call_id": "call_auto_1", + "id": "call_auto_0", + "call_id": "call_auto_0", "name": "get_weather", "arguments": "{\"location\":\"Tokyo\"}" } 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 8300258fc..b751794d4 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 @@ -323,6 +323,12 @@ impl GeminiProviderState { self.observed_tool_calls = true; continue; } + // Gemini streams are incremental and every functionCall part is a + // complete call, so parallel calls arriving in separate chunks all + // sit at parts[0]. Key calls by arrival order, not part position. + // Ids cannot disambiguate: they are optional, and the Antigravity + // envelope synthesizes per-chunk ids that repeat across chunks. + let index = self.tool_calls.len(); let tool_state = self.tool_calls.entry(index).or_default(); tool_state.call_id = function_call .get("id") @@ -1526,6 +1532,80 @@ mod tests { assert!(signature_index < call_index); } + #[test] + fn gemini_provider_state_keeps_parallel_function_calls_from_separate_chunks() { + let mut state = GeminiProviderState::default(); + let report_context = json!({}); + let chunk = |call: Value| { + data_line(json!({ + "responseId": "resp_parallel_123", + "modelVersion": "gemini-3.8-flash", + "candidates": [{ + "index": 0, + "content": {"role": "model", "parts": [{"functionCall": call}]} + }] + })) + }; + let mut frames = Vec::new(); + for call in [ + json!({"id": "call_a", "name": "get_weather", "args": {"city": "Paris"}}), + json!({"id": "call_b", "name": "get_weather", "args": {"city": "Tokyo"}}), + json!({"name": "get_time", "args": {"city": "Paris"}}), + json!({"name": "get_time", "args": {"city": "Tokyo"}}), + ] { + frames.extend( + state + .push_line(&report_context, chunk(call)) + .expect("function call chunk should parse"), + ); + } + + let starts = frames + .iter() + .filter_map(|frame| match &frame.event { + CanonicalStreamEvent::ToolCallStart { + index, + call_id, + name, + } => Some((*index, call_id.clone(), name.clone())), + _ => None, + }) + .collect::>(); + assert_eq!(starts.len(), 4); + assert_eq!( + starts + .iter() + .map(|(index, _, _)| *index) + .collect::>(), + vec![0, 1, 2, 3] + ); + assert_eq!(starts[0].1, "call_a"); + assert_eq!(starts[1].1, "call_b"); + assert_eq!( + starts + .iter() + .map(|(_, _, name)| name.as_str()) + .collect::>(), + vec!["get_weather", "get_weather", "get_time", "get_time"] + ); + assert_ne!(starts[2].1, starts[3].1); + + let mut arguments = BTreeMap::::new(); + for frame in &frames { + if let CanonicalStreamEvent::ToolCallArgumentsDelta { + index, + arguments: delta, + } = &frame.event + { + arguments.entry(*index).or_default().push_str(delta); + } + } + assert_eq!(arguments[&0], "{\"city\":\"Paris\"}"); + assert_eq!(arguments[&1], "{\"city\":\"Tokyo\"}"); + assert_eq!(arguments[&2], "{\"city\":\"Paris\"}"); + assert_eq!(arguments[&3], "{\"city\":\"Tokyo\"}"); + } + #[test] fn gemini_client_emitter_marks_reasoning_parts_as_thoughts() { let mut emitter = GeminiClientEmitter::default(); diff --git a/crates/aether-ai/formats/src/formats/openai/responses/response.rs b/crates/aether-ai/formats/src/formats/openai/responses/response.rs index 2b0f85163..86d44dc4b 100644 --- a/crates/aether-ai/formats/src/formats/openai/responses/response.rs +++ b/crates/aether-ai/formats/src/formats/openai/responses/response.rs @@ -114,6 +114,14 @@ fn openai_responses_incomplete_stop_reason(body: &Map) -> Canonic } } +fn canonical_incomplete_reason(canonical: &CanonicalResponse) -> Option<&'static str> { + match canonical.stop_reason.as_ref()? { + CanonicalStopReason::MaxTokens => Some("max_output_tokens"), + CanonicalStopReason::ContentFiltered => Some("content_filter"), + _ => None, + } +} + pub fn to_raw(canonical: &CanonicalResponse, report_context: &Value, compact: bool) -> Value { let namespace_tool_aliases = NamespaceToolAliases::from_report_context(report_context); let mut response = Map::new(); @@ -143,6 +151,17 @@ pub fn to_raw(canonical: &CanonicalResponse, report_context: &Value, compact: bo .cloned() { response.insert("status".to_string(), raw_status); + } else if let Some(reason) = canonical_incomplete_reason(canonical) { + // Cross-format sources carry no Responses status of their own; a + // truncated or filtered answer must not be reported as completed. + response.insert( + "status".to_string(), + Value::String("incomplete".to_string()), + ); + response.insert( + "incomplete_details".to_string(), + json!({ "reason": reason }), + ); } } @@ -931,4 +950,73 @@ mod tests { }) if id == "call_ws_1" && name == "web_search" && input["query"] == "today tech") ); } + + #[test] + fn responses_response_builder_reports_cross_format_truncation_as_incomplete() { + let response = |stop_reason| CanonicalResponse { + id: "gemini-resp".to_string(), + model: "gemini-3.8-flash".to_string(), + content: vec![CanonicalContentBlock::Text { + text: "partial".to_string(), + extensions: BTreeMap::new(), + }], + outputs: Vec::new(), + stop_reason: Some(stop_reason), + usage: None, + extensions: BTreeMap::new(), + }; + + let truncated = to_raw(&response(CanonicalStopReason::MaxTokens), &json!({}), false); + assert_eq!(truncated["status"], "incomplete"); + assert_eq!( + truncated["incomplete_details"], + json!({"reason": "max_output_tokens"}) + ); + + let filtered = to_raw( + &response(CanonicalStopReason::ContentFiltered), + &json!({}), + false, + ); + assert_eq!(filtered["status"], "incomplete"); + assert_eq!( + filtered["incomplete_details"], + json!({"reason": "content_filter"}) + ); + + let finished = to_raw(&response(CanonicalStopReason::EndTurn), &json!({}), false); + assert_eq!(finished["status"], "completed"); + assert!(finished.get("incomplete_details").is_none()); + } + + #[test] + fn gemini_max_tokens_response_converts_to_incomplete_responses_body() { + let gemini = json!({ + "responseId": "gemini-trunc-123", + "modelVersion": "gemini-3.8-flash", + "candidates": [{ + "content": {"role": "model", "parts": [{"text": "The printing press"}]}, + "finishReason": "MAX_TOKENS" + }], + "usageMetadata": { + "promptTokenCount": 19, + "candidatesTokenCount": 256, + "totalTokenCount": 275 + } + }); + + let body = crate::formats::registry::convert_response( + "gemini:generate_content", + "openai:responses", + &gemini, + &FormatContext::default(), + ) + .expect("gemini response should convert"); + + assert_eq!(body["status"], "incomplete"); + assert_eq!( + body["incomplete_details"], + json!({"reason": "max_output_tokens"}) + ); + } } diff --git a/crates/aether-ai/formats/src/formats/shared/sync_products.rs b/crates/aether-ai/formats/src/formats/shared/sync_products.rs index 0bc9b58a3..93fe07764 100644 --- a/crates/aether-ai/formats/src/formats/shared/sync_products.rs +++ b/crates/aether-ai/formats/src/formats/shared/sync_products.rs @@ -115,11 +115,26 @@ pub fn maybe_build_standard_cross_format_sync_product_from_normalized_payload( .as_deref() .unwrap_or(provider_api_format); + let aggregated_from_stream = aggregated_stream_body.is_some(); let Some(provider_body_json) = aggregated_stream_body.or_else(|| body_json.cloned()) else { return Ok(None); }; + let projection_fallback_body = aggregated_from_stream.then(|| provider_body_json.clone()); - Ok(maybe_build_standard_cross_format_sync_product( + let product = maybe_build_standard_cross_format_sync_product( + report_kind, + provider_body_api_format, + client_api_format, + report_context, + provider_body_json, + ); + if product.is_some() { + return Ok(product); + } + let Some(provider_body_json) = projection_fallback_body else { + return Ok(None); + }; + Ok(project_validated_openai_responses_stream_sync_product( report_kind, provider_body_api_format, client_api_format, @@ -128,6 +143,46 @@ pub fn maybe_build_standard_cross_format_sync_product_from_normalized_payload( )) } +/// Forced-stream Responses upstreams (Codex, xAI) echo request metadata such as +/// `parallel_tool_calls`, `tools` and encrypted reasoning back in the aggregated +/// body, which the strict cross-format response check refuses. The aggregated +/// body stays the provider body, so the client projection may drop those +/// provider-only fields — mirroring the OpenAI Chat client path. +fn project_validated_openai_responses_stream_sync_product( + report_kind: &str, + provider_api_format: &str, + client_api_format: &str, + report_context: &Value, + provider_body_json: Value, +) -> Option { + let provider_api_format = normalize_openai_responses_family_api_format(provider_api_format); + if !matches!( + provider_api_format.as_str(), + "openai:responses" | "openai:responses:compact" + ) { + return None; + } + let client_api_format = client_api_format.trim().to_ascii_lowercase(); + if is_standard_chat_finalize_kind(report_kind) { + sync_chat_response_conversion_kind(&provider_api_format, &client_api_format)?; + } else if is_standard_cli_finalize_kind(report_kind) { + sync_cli_response_conversion_kind(&provider_api_format, &client_api_format)?; + } else { + return None; + } + let client_body_json = project_validated_openai_responses_stream_to_client( + &provider_body_json, + &client_api_format, + report_context, + )?; + let client_body_json = + client_body_with_report_context_model(client_body_json, report_context, &client_api_format); + Some(StandardCrossFormatSyncProduct { + client_body_json, + provider_body_json, + }) +} + pub fn maybe_build_standard_same_format_sync_body_from_normalized_payload( report_kind: &str, status_code: u16, @@ -1474,6 +1529,14 @@ fn convert_openai_chat_canonical_response_to_openai_chat( fn project_validated_openai_responses_stream_to_openai_chat( body_json: &Value, report_context: &Value, +) -> Option { + project_validated_openai_responses_stream_to_client(body_json, "openai:chat", report_context) +} + +fn project_validated_openai_responses_stream_to_client( + body_json: &Value, + client_api_format: &str, + report_context: &Value, ) -> Option { // The caller retains body_json as provider_body_json. This projection is therefore allowed // to omit provider-only response metadata, but never unknown canonical output blocks. @@ -1490,7 +1553,12 @@ fn project_validated_openai_responses_stream_to_openai_chat( } apply_report_context_model_fallback(&mut canonical.model, report_context); - Some(canonical_to_openai_chat_response(&canonical)) + match client_api_format { + "openai:chat" => Some(canonical_to_openai_chat_response(&canonical)), + "claude:messages" => Some(canonical_to_claude_response(&canonical)), + "gemini:generate_content" => canonical_to_gemini_response(&canonical, report_context), + _ => None, + } } fn openai_chat_response_can_use_single_response_canonical(body_json: &Value) -> bool { @@ -7143,6 +7211,84 @@ mod tests { ); } + #[test] + fn standard_sync_finalize_projects_forced_responses_stream_to_gemini_and_claude_clients() { + // Shape of a forced-stream xAI / Codex upstream: the terminal response + // echoes request metadata and carries encrypted reasoning, which the + // strict cross-format check refuses. + let stream_body = concat!( + "data: {\"type\":\"response.created\",\"sequence_number\":0,\"response\":{\"id\":\"resp_forced_123\",\"object\":\"response\",\"status\":\"in_progress\",\"model\":\"grok-4.7-build\",\"output\":[],\"parallel_tool_calls\":true,\"tool_choice\":\"auto\",\"tools\":[],\"temperature\":0.7}}\n\n", + "data: {\"type\":\"response.output_item.done\",\"sequence_number\":1,\"output_index\":0,\"item\":{\"id\":\"rs_forced_123\",\"type\":\"reasoning\",\"status\":\"completed\",\"summary\":[{\"type\":\"summary_text\",\"text\":\"greet briefly\"}],\"encrypted_content\":\"opaque-xai-reasoning\"}}\n\n", + "data: {\"type\":\"response.output_text.delta\",\"sequence_number\":2,\"item_id\":\"msg_forced_123\",\"output_index\":1,\"content_index\":0,\"delta\":\"Hello there friend\"}\n\n", + "data: {\"type\":\"response.output_item.done\",\"sequence_number\":3,\"output_index\":1,\"item\":{\"id\":\"msg_forced_123\",\"type\":\"message\",\"status\":\"completed\",\"role\":\"assistant\",\"content\":[{\"type\":\"output_text\",\"text\":\"Hello there friend\",\"annotations\":[]}]}}\n\n", + "data: {\"type\":\"response.output_item.done\",\"sequence_number\":4,\"output_index\":2,\"item\":{\"id\":\"fc_forced_123\",\"type\":\"function_call\",\"status\":\"completed\",\"call_id\":\"call_forced_123\",\"name\":\"search\",\"arguments\":\"{\\\"q\\\":\\\"aether\\\"}\"}}\n\n", + "data: {\"type\":\"response.completed\",\"sequence_number\":5,\"response\":{\"id\":\"resp_forced_123\",\"object\":\"response\",\"status\":\"completed\",\"model\":\"grok-4.7-build\",\"output\":[],\"parallel_tool_calls\":true,\"tool_choice\":\"auto\",\"tools\":[{\"type\":\"function\",\"name\":\"search\",\"parameters\":{\"type\":\"object\"}}],\"text\":{\"format\":{\"type\":\"text\"}},\"reasoning\":{\"effort\":null,\"summary\":null},\"temperature\":0.7,\"top_p\":0.95,\"store\":false,\"usage\":{\"input_tokens\":1249,\"input_tokens_details\":{\"cached_tokens\":1152},\"output_tokens\":40,\"output_tokens_details\":{\"reasoning_tokens\":31},\"total_tokens\":1289}}}\n\n", + ); + let encoded = base64::engine::general_purpose::STANDARD.encode(stream_body); + + for (report_kind, client_api_format) in [ + ("gemini_chat_sync_finalize", "gemini:generate_content"), + ("gemini_cli_sync_finalize", "gemini:generate_content"), + ("claude_chat_sync_finalize", "claude:messages"), + ("claude_cli_sync_finalize", "claude:messages"), + ] { + let report_context = json!({ + "provider_api_format": "openai:responses", + "provider_stream_event_api_format": "openai:responses", + "client_api_format": client_api_format, + "model": "grok-4.7", + "mapped_model": "grok-4.7", + "needs_conversion": true, + }); + let product = maybe_build_standard_sync_finalize_product_from_normalized_payload( + report_kind, + 200, + Some(&report_context), + None, + Some(&encoded), + ) + .expect("forced Responses stream should aggregate") + .unwrap_or_else(|| panic!("{report_kind} should receive a projection")); + let StandardSyncFinalizeNormalizedProduct::CrossFormat(product) = product else { + panic!("{report_kind}: Responses stream should stay a cross-format product") + }; + assert_eq!(product.provider_body_json["parallel_tool_calls"], true); + let client = product.client_body_json.to_string(); + assert!( + !client.contains("opaque-xai-reasoning") && !client.contains("response.created"), + "{report_kind}: provider-only data leaked into the client body: {client}" + ); + if client_api_format == "gemini:generate_content" { + let parts = product.client_body_json["candidates"][0]["content"]["parts"] + .as_array() + .expect("gemini parts"); + assert!(parts + .iter() + .any(|part| part["text"] == "Hello there friend" + && part.get("thought").is_none())); + assert!(parts + .iter() + .any(|part| part["functionCall"]["name"] == "search" + && part["functionCall"]["args"]["q"] == "aether")); + assert_eq!( + product.client_body_json["usageMetadata"]["promptTokenCount"], + 1249 + ); + } else { + let content = product.client_body_json["content"] + .as_array() + .expect("claude content"); + assert!(content + .iter() + .any(|block| block["type"] == "text" && block["text"] == "Hello there friend")); + assert!(content.iter().any(|block| block["type"] == "tool_use" + && block["name"] == "search" + && block["input"]["q"] == "aether")); + assert_eq!(product.client_body_json["stop_reason"], "tool_use"); + } + } + } + #[test] fn standard_sync_finalize_projects_authoritative_incomplete_responses_stream() { let report_context = json!({ diff --git a/crates/aether-ai/formats/src/protocol/canonical.rs b/crates/aether-ai/formats/src/protocol/canonical.rs index 8a1abc49d..52db51749 100644 --- a/crates/aether-ai/formats/src/protocol/canonical.rs +++ b/crates/aether-ai/formats/src/protocol/canonical.rs @@ -4935,13 +4935,22 @@ pub(crate) fn gemini_response_format_to_canonical( if response_mime_type != "application/json" { return None; } - let json_schema = gemini_value_by_case(generation_config, "responseSchema", "response_schema") - .map(|schema| { - json!({ - "name": "response_schema", - "schema": schema, - }) - }); + let json_schema = gemini_value_by_case( + generation_config, + "responseJsonSchema", + "response_json_schema", + ) + .cloned() + .or_else(|| { + gemini_value_by_case(generation_config, "responseSchema", "response_schema") + .map(gemini_openapi_schema_to_json_schema) + }) + .map(|schema| { + json!({ + "name": "response_schema", + "schema": schema, + }) + }); Some(CanonicalResponseFormat { format_type: if json_schema.is_some() { "json_schema".to_string() @@ -4953,6 +4962,68 @@ pub(crate) fn gemini_response_format_to_canonical( }) } +/// `parametersJsonSchema` is already standard JSON Schema; the legacy +/// `parameters` field is Gemini's OpenAPI subset with upper-case type names. +fn gemini_declaration_parameters_to_json_schema(declaration: &Map) -> Option { + gemini_value_by_case( + declaration, + "parametersJsonSchema", + "parameters_json_schema", + ) + .cloned() + .or_else(|| { + declaration + .get("parameters") + .map(gemini_openapi_schema_to_json_schema) + }) +} + +/// Gemini's OpenAPI-style `Schema` spells types in upper case (`OBJECT`, +/// `STRING`, ...); other protocols expect JSON Schema's lower-case names. +pub(crate) fn gemini_openapi_schema_to_json_schema(schema: &Value) -> Value { + fn normalize(value: &mut Value) { + match value { + Value::Object(object) => { + for (key, child) in object.iter_mut() { + if key == "type" { + match child { + Value::String(type_name) => lowercase_schema_type(type_name), + Value::Array(type_names) => { + for type_name in type_names.iter_mut() { + if let Value::String(type_name) = type_name { + lowercase_schema_type(type_name); + } + } + } + other => normalize(other), + } + } else if key != "enum" + && key != "const" + && key != "default" + && key != "example" + { + normalize(child); + } + } + } + Value::Array(items) => items.iter_mut().for_each(normalize), + _ => {} + } + } + fn lowercase_schema_type(type_name: &mut String) { + if matches!( + type_name.as_str(), + "OBJECT" | "STRING" | "INTEGER" | "NUMBER" | "BOOLEAN" | "ARRAY" | "NULL" + ) { + *type_name = type_name.to_ascii_lowercase(); + } + } + + let mut schema = schema.clone(); + normalize(&mut schema); + schema +} + pub(crate) type GeminiCanonicalTools = ( Vec, Vec, @@ -5126,12 +5197,18 @@ pub(crate) fn gemini_tools_to_canonical(value: Option<&Value>) -> Option