mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-06 01:17:46 +08:00
fix(formats): repair live cross-format conversion gaps
Verified against a live Antigravity + xAI deployment: - Gemini and Claude clients calling a forced-stream Responses upstream (xAI, Codex) without streaming always failed: the aggregated body echoes request metadata (parallel_tool_calls, tools, encrypted reasoning) that the strict cross-format check refuses, and the gateway then wrapped the raw SSE capture in a client error body sent with HTTP 200. Project the validated aggregate to every client format, as the Chat path already does, and return 502 instead of raw provider bytes when a successful cross-format response still cannot be converted. - Gemini stream decoding keyed tool calls by part position, so parallel calls arriving in separate chunks (all at parts[0]) merged into one call with concatenated arguments. Key them by arrival order; ids cannot be used because they are optional and the Antigravity envelope synthesizes per-chunk ids that repeat across chunks. Generated call_auto_N ids now follow arrival order. - Non-stream Responses output reported truncated or filtered cross-format answers as completed; derive incomplete + incomplete_details from the canonical stop reason. - Gemini request parsing ignored parametersJsonSchema and responseJsonSchema and passed OpenAPI upper-case type names (OBJECT, STRING) through to JSON Schema targets, which xAI rejects. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
d30268f80f
commit
75bc32cfe9
@@ -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<Option<Response<Body>>, 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<bool, GatewayError> {
|
||||
@@ -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::<serde_json::Value>(&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::<serde_json::Value>(
|
||||
&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(
|
||||
|
||||
@@ -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\"}"
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user