mirror of
https://github.com/fawney19/Aether.git
synced 2026-10-04 00:17:45 +08:00
fix(finalize): 严格聚合并投影同步 Responses 流
This commit is contained in:
@@ -10,7 +10,8 @@ use aether_ai_formats::formats::conversion::response::{
|
||||
use aether_ai_formats::formats::openai::responses::response::ensure_modern_openai_responses_response_fields;
|
||||
use aether_ai_formats::formats::registry::{convert_response, FormatContext, FormatError};
|
||||
use aether_ai_formats::{
|
||||
canonical_to_claude_response, canonical_to_embedding_response, canonical_to_gemini_response,
|
||||
canonical_response_unknown_block_count, canonical_to_claude_response,
|
||||
canonical_to_embedding_response, canonical_to_gemini_response,
|
||||
canonical_to_openai_chat_response, canonical_to_openai_responses_compact_response,
|
||||
canonical_to_openai_responses_response, from_claude_to_canonical_response,
|
||||
from_embedding_to_canonical_response, from_gemini_to_canonical_response,
|
||||
@@ -41,6 +42,12 @@ pub struct StandardCrossFormatSyncProduct {
|
||||
pub provider_body_json: Value,
|
||||
}
|
||||
|
||||
struct OpenAiCrossFormatProviderBody {
|
||||
body_json: Value,
|
||||
// Only validated stream aggregation may use a client projection that omits provider-only fields.
|
||||
aggregated_from_stream: bool,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, PartialEq)]
|
||||
pub enum StandardSyncFinalizeNormalizedProduct {
|
||||
SuccessBody(Value),
|
||||
@@ -198,7 +205,7 @@ pub fn maybe_build_openai_chat_cross_format_sync_product_from_normalized_payload
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
let Some(provider_body_json) =
|
||||
let Some(provider_body) =
|
||||
maybe_build_openai_cross_format_provider_body_from_normalized_payload(
|
||||
body_json,
|
||||
body_base64,
|
||||
@@ -208,6 +215,7 @@ pub fn maybe_build_openai_chat_cross_format_sync_product_from_normalized_payload
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
let provider_body_json = provider_body.body_json;
|
||||
let Some(client_body_json) = (match provider_api_format.as_str() {
|
||||
"claude:messages" => {
|
||||
convert_claude_canonical_response_to_openai_chat(&provider_body_json, report_context)
|
||||
@@ -217,6 +225,17 @@ pub fn maybe_build_openai_chat_cross_format_sync_product_from_normalized_payload
|
||||
}
|
||||
"openai:responses" => {
|
||||
convert_openai_responses_response_to_openai_chat(&provider_body_json, report_context)
|
||||
.or_else(|| {
|
||||
provider_body
|
||||
.aggregated_from_stream
|
||||
.then(|| {
|
||||
project_validated_openai_responses_stream_to_openai_chat(
|
||||
&provider_body_json,
|
||||
report_context,
|
||||
)
|
||||
})
|
||||
.flatten()
|
||||
})
|
||||
}
|
||||
_ => None,
|
||||
}) else {
|
||||
@@ -267,7 +286,7 @@ pub fn maybe_build_openai_responses_cross_format_sync_product_from_normalized_pa
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
let Some(provider_body_json) =
|
||||
let Some(provider_body) =
|
||||
maybe_build_openai_cross_format_provider_body_from_normalized_payload(
|
||||
body_json,
|
||||
body_base64,
|
||||
@@ -276,6 +295,7 @@ pub fn maybe_build_openai_responses_cross_format_sync_product_from_normalized_pa
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
let provider_body_json = provider_body.body_json;
|
||||
|
||||
let normalized_provider_api_format =
|
||||
normalize_openai_responses_family_api_format(&provider_api_format);
|
||||
@@ -847,7 +867,16 @@ fn maybe_build_openai_responses_same_family_stream_sync_body(
|
||||
return Ok(None);
|
||||
};
|
||||
let body_bytes = base64::engine::general_purpose::STANDARD.decode(body_base64)?;
|
||||
if let Some(terminal_body) = terminal_openai_responses_stream_response(&body_bytes) {
|
||||
// Same-family clients retain the authoritative terminal body verbatim, including future
|
||||
// output item fields, but unknown intermediate event types still fail closed.
|
||||
ensure_no_unknown_openai_responses_stream_events(&body_bytes, true)?;
|
||||
let terminal_body = validated_terminal_openai_responses_stream_response(&body_bytes)?;
|
||||
if terminal_body.get("status").and_then(Value::as_str) != Some("completed")
|
||||
|| terminal_body
|
||||
.get("output")
|
||||
.and_then(Value::as_array)
|
||||
.is_some_and(|output| !output.is_empty())
|
||||
{
|
||||
return Ok(Some(client_body_with_report_context_model(
|
||||
terminal_body,
|
||||
report_context,
|
||||
@@ -855,43 +884,78 @@ fn maybe_build_openai_responses_same_family_stream_sync_body(
|
||||
)));
|
||||
}
|
||||
Ok(
|
||||
try_aggregate_openai_responses_stream_sync_response(&body_bytes)?.map(|body| {
|
||||
aggregate_openai_responses_stream_sync_response_from_validated_terminal(
|
||||
&body_bytes,
|
||||
terminal_body,
|
||||
)
|
||||
.map(|body| {
|
||||
client_body_with_report_context_model(body, report_context, &client_api_format)
|
||||
}),
|
||||
)
|
||||
}
|
||||
|
||||
fn terminal_openai_responses_stream_response(body: &[u8]) -> Option<Value> {
|
||||
parse_stream_json_events(body)?
|
||||
.into_iter()
|
||||
.rev()
|
||||
.find_map(|event| {
|
||||
let event = event.as_object()?;
|
||||
let event_type = event.get("type").and_then(Value::as_str)?;
|
||||
if !matches!(
|
||||
event_type,
|
||||
"response.completed" | "response.done" | "response.incomplete" | "response.failed"
|
||||
) {
|
||||
return None;
|
||||
}
|
||||
let response = event.get("response").and_then(Value::as_object).cloned()?;
|
||||
if matches!(event_type, "response.completed" | "response.done")
|
||||
&& response
|
||||
.get("output")
|
||||
.and_then(Value::as_array)
|
||||
.is_none_or(|output| output.is_empty())
|
||||
{
|
||||
return None;
|
||||
}
|
||||
Some(Value::Object(response))
|
||||
})
|
||||
fn validated_terminal_openai_responses_stream_response(
|
||||
body: &[u8],
|
||||
) -> Result<Value, AiSurfaceFinalizeError> {
|
||||
let events = parse_stream_json_events(body).ok_or_else(|| {
|
||||
AiSurfaceFinalizeError::new("OpenAI Responses stream contains invalid event framing")
|
||||
})?;
|
||||
let mut terminal_response = None;
|
||||
|
||||
for event in events {
|
||||
let event = event.as_object().ok_or_else(|| {
|
||||
AiSurfaceFinalizeError::new("OpenAI Responses stream contains a non-object event")
|
||||
})?;
|
||||
let Some(event_type) = event.get("type").and_then(Value::as_str) else {
|
||||
continue;
|
||||
};
|
||||
let expected_status = match event_type {
|
||||
"response.completed" | "response.done" => "completed",
|
||||
"response.incomplete" => "incomplete",
|
||||
"response.failed" => "failed",
|
||||
_ => continue,
|
||||
};
|
||||
if terminal_response.is_some() {
|
||||
return Err(AiSurfaceFinalizeError::new(
|
||||
"OpenAI Responses stream contains multiple authoritative terminal events",
|
||||
));
|
||||
}
|
||||
let mut response = event
|
||||
.get("response")
|
||||
.and_then(Value::as_object)
|
||||
.cloned()
|
||||
.ok_or_else(|| {
|
||||
AiSurfaceFinalizeError::new(
|
||||
"OpenAI Responses terminal event is missing its response object",
|
||||
)
|
||||
})?;
|
||||
if response
|
||||
.get("status")
|
||||
.and_then(Value::as_str)
|
||||
.is_some_and(|status| status != expected_status)
|
||||
{
|
||||
return Err(AiSurfaceFinalizeError::new(
|
||||
"OpenAI Responses terminal event conflicts with response.status",
|
||||
));
|
||||
}
|
||||
response
|
||||
.entry("status".to_string())
|
||||
.or_insert_with(|| Value::String(expected_status.to_string()));
|
||||
terminal_response = Some(Value::Object(response));
|
||||
}
|
||||
|
||||
terminal_response.ok_or_else(|| {
|
||||
AiSurfaceFinalizeError::new(
|
||||
"OpenAI Responses stream is missing an authoritative terminal event",
|
||||
)
|
||||
})
|
||||
}
|
||||
|
||||
fn maybe_build_openai_cross_format_provider_body_from_normalized_payload(
|
||||
body_json: Option<&Value>,
|
||||
body_base64: Option<&str>,
|
||||
provider_api_format: &str,
|
||||
) -> Result<Option<Value>, AiSurfaceFinalizeError> {
|
||||
) -> Result<Option<OpenAiCrossFormatProviderBody>, AiSurfaceFinalizeError> {
|
||||
let aggregated_stream_body = match body_base64 {
|
||||
Some(body_base64) => {
|
||||
let body_bytes = base64::engine::general_purpose::STANDARD.decode(body_base64)?;
|
||||
@@ -911,8 +975,14 @@ fn maybe_build_openai_cross_format_provider_body_from_normalized_payload(
|
||||
None => None,
|
||||
};
|
||||
|
||||
let aggregated_from_stream = aggregated_stream_body.is_some();
|
||||
let provider_body_json = aggregated_stream_body.or_else(|| body_json.cloned());
|
||||
Ok(provider_body_json.filter(|value| !is_error_like_sync_body(value)))
|
||||
Ok(provider_body_json
|
||||
.filter(|value| !is_error_like_sync_body(value))
|
||||
.map(|body_json| OpenAiCrossFormatProviderBody {
|
||||
body_json,
|
||||
aggregated_from_stream,
|
||||
}))
|
||||
}
|
||||
|
||||
fn is_error_like_sync_body(value: &Value) -> bool {
|
||||
@@ -921,6 +991,7 @@ fn is_error_like_sync_body(value: &Value) -> bool {
|
||||
};
|
||||
|
||||
object.get("error").is_some_and(|error| !error.is_null())
|
||||
|| object.get("status").and_then(Value::as_str) == Some("failed")
|
||||
|| object
|
||||
.get("type")
|
||||
.and_then(Value::as_str)
|
||||
@@ -1316,6 +1387,28 @@ fn convert_openai_chat_canonical_response_to_openai_chat(
|
||||
Some(canonical_to_openai_chat_response(&canonical))
|
||||
}
|
||||
|
||||
fn project_validated_openai_responses_stream_to_openai_chat(
|
||||
body_json: &Value,
|
||||
report_context: &Value,
|
||||
) -> Option<Value> {
|
||||
// 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.
|
||||
if !matches!(
|
||||
body_json.get("status").and_then(Value::as_str),
|
||||
Some("completed" | "incomplete")
|
||||
) {
|
||||
return None;
|
||||
}
|
||||
|
||||
let mut canonical = from_openai_responses_to_canonical_response(body_json)?;
|
||||
if canonical_response_unknown_block_count(&canonical) > 0 {
|
||||
return None;
|
||||
}
|
||||
|
||||
apply_report_context_model_fallback(&mut canonical.model, report_context);
|
||||
Some(canonical_to_openai_chat_response(&canonical))
|
||||
}
|
||||
|
||||
fn openai_chat_response_can_use_single_response_canonical(body_json: &Value) -> bool {
|
||||
body_json
|
||||
.get("choices")
|
||||
@@ -1841,12 +1934,78 @@ fn try_aggregate_openai_chat_stream_sync_response(
|
||||
fn try_aggregate_openai_responses_stream_sync_response(
|
||||
body: &[u8],
|
||||
) -> Result<Option<Value>, AiSurfaceFinalizeError> {
|
||||
ensure_no_unknown_openai_responses_stream_events(body, false)?;
|
||||
let terminal_response = validated_terminal_openai_responses_stream_response(body)?;
|
||||
Ok(
|
||||
aggregate_openai_responses_stream_sync_response_from_validated_terminal(
|
||||
body,
|
||||
terminal_response,
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
fn ensure_no_unknown_openai_responses_stream_events(
|
||||
body: &[u8],
|
||||
allow_authoritative_terminal_output_extensions: bool,
|
||||
) -> Result<(), AiSurfaceFinalizeError> {
|
||||
let Ok(text) = std::str::from_utf8(body) else {
|
||||
return Ok(());
|
||||
};
|
||||
let report_context = Value::Object(Map::new());
|
||||
let mut provider = OpenAIResponsesProviderState::default();
|
||||
ensure_no_unknown_provider_stream_events(body, |line| {
|
||||
provider.push_line(&report_context, line)
|
||||
})?;
|
||||
Ok(aggregate_openai_responses_stream_sync_response(body))
|
||||
let mut declared_event_type: Option<String> = None;
|
||||
for raw_line in text.lines() {
|
||||
let frames = provider.push_line(&report_context, raw_line.as_bytes().to_vec())?;
|
||||
let raw_payload = raw_line.trim();
|
||||
if let Some(event_type) = raw_payload.strip_prefix("event:") {
|
||||
declared_event_type = Some(event_type.trim().to_string());
|
||||
continue;
|
||||
}
|
||||
let payload = raw_payload
|
||||
.strip_prefix("data:")
|
||||
.map(str::trim)
|
||||
.unwrap_or(raw_payload);
|
||||
let declared_event_type = declared_event_type.take();
|
||||
let payload_event_type = serde_json::from_str::<Value>(payload)
|
||||
.ok()
|
||||
.and_then(|payload| {
|
||||
payload
|
||||
.get("type")
|
||||
.and_then(Value::as_str)
|
||||
.map(ToOwned::to_owned)
|
||||
});
|
||||
if matches!(
|
||||
(declared_event_type.as_deref(), payload_event_type.as_deref()),
|
||||
(Some(declared), Some(payload)) if declared != payload
|
||||
) {
|
||||
return Err(AiSurfaceFinalizeError::new(
|
||||
"OpenAI Responses SSE event name conflicts with payload.type",
|
||||
));
|
||||
}
|
||||
let event_type = payload_event_type.or(declared_event_type);
|
||||
let terminal_output_extensions_are_lossless = allow_authoritative_terminal_output_extensions
|
||||
&& event_type.as_deref().is_some_and(|event_type| {
|
||||
matches!(
|
||||
event_type,
|
||||
"response.completed"
|
||||
| "response.done"
|
||||
| "response.incomplete"
|
||||
| "response.failed"
|
||||
)
|
||||
});
|
||||
if let Some(payload) = frames.iter().find_map(|frame| match &frame.event {
|
||||
CanonicalStreamEvent::UnknownEvent(payload)
|
||||
if payload.get("type").and_then(Value::as_str) != Some("response.failed")
|
||||
&& !terminal_output_extensions_are_lossless =>
|
||||
{
|
||||
Some(payload)
|
||||
}
|
||||
_ => None,
|
||||
}) {
|
||||
return Err(unsupported_stream_event_finalize_error(payload));
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn try_aggregate_claude_stream_sync_response(
|
||||
@@ -2054,12 +2213,21 @@ pub fn aggregate_openai_chat_stream_sync_response(body: &[u8]) -> Option<Value>
|
||||
}
|
||||
|
||||
pub fn aggregate_openai_responses_stream_sync_response(body: &[u8]) -> Option<Value> {
|
||||
try_aggregate_openai_responses_stream_sync_response(body)
|
||||
.ok()
|
||||
.flatten()
|
||||
}
|
||||
|
||||
fn aggregate_openai_responses_stream_sync_response_from_validated_terminal(
|
||||
body: &[u8],
|
||||
terminal_response: Value,
|
||||
) -> Option<Value> {
|
||||
let events = parse_stream_json_events(body)?;
|
||||
if events.is_empty() {
|
||||
return None;
|
||||
}
|
||||
|
||||
let mut response_object: Option<Map<String, Value>> = None;
|
||||
let mut response = terminal_response.as_object()?.clone();
|
||||
let mut response_id: Option<String> = None;
|
||||
let mut model: Option<String> = None;
|
||||
let mut message_states: BTreeMap<usize, OpenAIResponsesSyncMessageState> = BTreeMap::new();
|
||||
@@ -2093,12 +2261,7 @@ pub fn aggregate_openai_responses_stream_sync_response(body: &[u8]) -> Option<Va
|
||||
.and_then(Value::as_str)
|
||||
.unwrap_or_default()
|
||||
{
|
||||
"response.created" | "response.in_progress" if response_object.is_none() => {
|
||||
response_object = event_object
|
||||
.get("response")
|
||||
.and_then(Value::as_object)
|
||||
.cloned();
|
||||
}
|
||||
"response.created" | "response.in_progress" => {}
|
||||
"response.output_text.delta" | "response.outtext.delta" => {
|
||||
let output_index = openai_responses_event_output_index(event_object).unwrap_or(0);
|
||||
let content_index = openai_responses_event_content_index(event_object);
|
||||
@@ -2413,12 +2576,7 @@ pub fn aggregate_openai_responses_stream_sync_response(body: &[u8]) -> Option<Va
|
||||
output_index,
|
||||
);
|
||||
}
|
||||
"response.completed" | "response.done" => {
|
||||
response_object = event_object
|
||||
.get("response")
|
||||
.and_then(Value::as_object)
|
||||
.cloned()
|
||||
.or(response_object);
|
||||
"response.completed" | "response.done" | "response.incomplete" | "response.failed" => {
|
||||
let Some(response) = event_object.get("response").and_then(Value::as_object) else {
|
||||
continue;
|
||||
};
|
||||
@@ -2461,18 +2619,19 @@ pub fn aggregate_openai_responses_stream_sync_response(body: &[u8]) -> Option<Va
|
||||
}
|
||||
}
|
||||
|
||||
let mut response = response_object.unwrap_or_else(|| {
|
||||
let mut response = Map::new();
|
||||
if let Some(response_id) = response_id.as_ref() {
|
||||
response.insert("id".to_string(), Value::String(response_id.clone()));
|
||||
}
|
||||
response.insert("object".to_string(), Value::String("response".to_string()));
|
||||
response.insert("status".to_string(), Value::String("completed".to_string()));
|
||||
if let Some(model) = model.as_ref() {
|
||||
response.insert("model".to_string(), Value::String(model.clone()));
|
||||
}
|
||||
response
|
||||
.entry("object".to_string())
|
||||
.or_insert_with(|| Value::String("response".to_string()));
|
||||
if let Some(response_id) = response_id.as_ref() {
|
||||
response
|
||||
});
|
||||
.entry("id".to_string())
|
||||
.or_insert_with(|| Value::String(response_id.clone()));
|
||||
}
|
||||
if let Some(model) = model.as_ref() {
|
||||
response
|
||||
.entry("model".to_string())
|
||||
.or_insert_with(|| Value::String(model.clone()));
|
||||
}
|
||||
|
||||
let response_id = response
|
||||
.get("id")
|
||||
@@ -3776,7 +3935,7 @@ mod tests {
|
||||
maybe_build_standard_cross_format_sync_product_from_normalized_payload,
|
||||
maybe_build_standard_same_format_sync_body_from_normalized_payload,
|
||||
maybe_build_standard_sync_finalize_product_from_normalized_payload,
|
||||
StandardSyncFinalizeNormalizedProduct,
|
||||
try_aggregate_openai_responses_stream_sync_response, StandardSyncFinalizeNormalizedProduct,
|
||||
};
|
||||
use aether_ai_formats::formats::conversion::response::{
|
||||
convert_claude_chat_response_to_openai_chat, convert_gemini_chat_response_to_openai_chat,
|
||||
@@ -4647,9 +4806,7 @@ mod tests {
|
||||
}
|
||||
});
|
||||
let body = format!(
|
||||
"event: response.program.delta\ndata: {{\"type\":\"response.program.delta\",\"delta\":\"print\"}}\n\n\
|
||||
event: response.multi_agent_call.in_progress\ndata: {{\"type\":\"response.multi_agent_call.in_progress\",\"item_id\":\"agent-call-1\"}}\n\n\
|
||||
event: response.completed\ndata: {}\n\n",
|
||||
"event: response.completed\ndata: {}\n\n",
|
||||
json!({"type": "response.completed", "response": terminal_response})
|
||||
);
|
||||
let report_context = json!({
|
||||
@@ -4665,7 +4822,7 @@ mod tests {
|
||||
None,
|
||||
Some(&base64::engine::general_purpose::STANDARD.encode(body)),
|
||||
)
|
||||
.expect("same-family terminal snapshot should bypass unknown-event aggregation")
|
||||
.expect("same-family terminal snapshot extensions should remain lossless")
|
||||
.expect("terminal response should become the sync response");
|
||||
|
||||
assert_eq!(body_json, terminal_response);
|
||||
@@ -4704,15 +4861,20 @@ mod tests {
|
||||
"model": "gpt-5.6-sol",
|
||||
"status": "incomplete",
|
||||
"output": [{
|
||||
"type": "agent_message",
|
||||
"author": "researcher",
|
||||
"recipient": "assistant",
|
||||
"encrypted_content": "encrypted-partial-message"
|
||||
"type": "message",
|
||||
"id": "msg_gpt56_incomplete",
|
||||
"role": "assistant",
|
||||
"status": "incomplete",
|
||||
"content": [{
|
||||
"type": "output_text",
|
||||
"text": "partial",
|
||||
"annotations": []
|
||||
}]
|
||||
}],
|
||||
"incomplete_details": {"reason": "max_output_tokens"}
|
||||
});
|
||||
let body = format!(
|
||||
"event: response.agent_message.delta\ndata: {{\"type\":\"response.agent_message.delta\",\"delta\":\"partial\"}}\n\n\
|
||||
"event: response.output_text.delta\ndata: {{\"type\":\"response.output_text.delta\",\"output_index\":0,\"content_index\":0,\"delta\":\"partial\"}}\n\n\
|
||||
event: response.incomplete\ndata: {}\n\n",
|
||||
json!({"type": "response.incomplete", "response": terminal_response})
|
||||
);
|
||||
@@ -5948,6 +6110,75 @@ mod tests {
|
||||
.contains("field $.type = \"response.future.delta\""));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn responses_same_family_stream_rejects_unknown_event_before_valid_terminal_snapshot() {
|
||||
let report_context = json!({
|
||||
"provider_api_format": "openai:responses",
|
||||
"client_api_format": "openai:responses",
|
||||
"needs_conversion": false,
|
||||
"model": "gpt-5",
|
||||
"mapped_model": "gpt-5",
|
||||
});
|
||||
let stream_body = concat!(
|
||||
"data: {\"type\":\"response.future.delta\",\"payload\":{\"must_not_drop\":true}}\n\n",
|
||||
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_unknown_before_terminal\",\"object\":\"response\",\"model\":\"gpt-5\",\"status\":\"completed\",\"output\":[{\"type\":\"message\",\"id\":\"msg_unknown_before_terminal\",\"role\":\"assistant\",\"status\":\"completed\",\"content\":[{\"type\":\"output_text\",\"text\":\"must not pass\",\"annotations\":[]}]}]}}\n\n",
|
||||
);
|
||||
|
||||
let error = maybe_build_openai_responses_same_family_sync_body_from_normalized_payload(
|
||||
"openai_responses_sync_finalize",
|
||||
200,
|
||||
Some(&report_context),
|
||||
None,
|
||||
Some(&base64::engine::general_purpose::STANDARD.encode(stream_body)),
|
||||
)
|
||||
.expect_err("unknown event must fail closed even before a valid terminal snapshot");
|
||||
|
||||
assert!(error
|
||||
.to_string()
|
||||
.contains("Unsupported provider stream event cannot be converted losslessly"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn responses_stream_rejects_conflicting_or_multiple_authoritative_terminals() {
|
||||
let cases = [
|
||||
(
|
||||
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_conflicting_terminal\",\"object\":\"response\",\"status\":\"incomplete\",\"output\":[]}}\n\n",
|
||||
"conflicts with response.status",
|
||||
),
|
||||
(
|
||||
concat!(
|
||||
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_multiple_terminal\",\"object\":\"response\",\"status\":\"completed\",\"output\":[]}}\n\n",
|
||||
"data: {\"type\":\"response.failed\",\"response\":{\"id\":\"resp_multiple_terminal\",\"object\":\"response\",\"status\":\"failed\",\"output\":[],\"error\":{\"code\":\"server_error\",\"message\":\"late failure\"}}}\n\n",
|
||||
),
|
||||
"multiple authoritative terminal events",
|
||||
),
|
||||
(
|
||||
concat!(
|
||||
"event: response.future.delta\n",
|
||||
"data: {\"type\":\"response.output_text.delta\",\"output_index\":0,\"content_index\":0,\"delta\":\"disguised\"}\n\n",
|
||||
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_disguised_unknown\",\"object\":\"response\",\"status\":\"completed\",\"output\":[]}}\n\n",
|
||||
),
|
||||
"SSE event name conflicts with payload.type",
|
||||
),
|
||||
(
|
||||
concat!(
|
||||
"event: response.failed\n",
|
||||
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_disguised_terminal\",\"object\":\"response\",\"status\":\"completed\",\"output\":[]}}\n\n",
|
||||
),
|
||||
"SSE event name conflicts with payload.type",
|
||||
),
|
||||
];
|
||||
|
||||
for (stream_body, expected_message) in cases {
|
||||
let error = try_aggregate_openai_responses_stream_sync_response(stream_body.as_bytes())
|
||||
.expect_err("invalid terminal lifecycle must fail closed");
|
||||
assert!(
|
||||
error.to_string().contains(expected_message),
|
||||
"unexpected error: {error}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn standard_sync_finalize_product_prefers_same_format_success_body() {
|
||||
let report_context = json!({
|
||||
@@ -6057,7 +6288,6 @@ mod tests {
|
||||
"mapped_model": "gpt-5",
|
||||
"needs_conversion": true,
|
||||
});
|
||||
|
||||
let product = maybe_build_standard_sync_finalize_product_from_normalized_payload(
|
||||
"claude_chat_sync_finalize",
|
||||
200,
|
||||
@@ -6132,6 +6362,23 @@ mod tests {
|
||||
"content_index": 0,
|
||||
"text": "Hello capture"
|
||||
},
|
||||
{
|
||||
"type": "response.output_item.done",
|
||||
"output_index": 0,
|
||||
"item": {
|
||||
"id": "msg_capture_123",
|
||||
"type": "message",
|
||||
"status": "completed",
|
||||
"role": "assistant",
|
||||
"phase": "final_answer",
|
||||
"content": [{
|
||||
"type": "output_text",
|
||||
"text": "Hello capture",
|
||||
"annotations": [],
|
||||
"logprobs": []
|
||||
}]
|
||||
}
|
||||
},
|
||||
{
|
||||
"type": "response.completed",
|
||||
"response": {
|
||||
@@ -6139,6 +6386,8 @@ mod tests {
|
||||
"object": "response",
|
||||
"status": "completed",
|
||||
"model": "gpt-5.6-luna",
|
||||
"background": false,
|
||||
"max_output_tokens": null,
|
||||
"output": [],
|
||||
"usage": {
|
||||
"input_tokens": 2,
|
||||
@@ -6183,6 +6432,213 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn standard_sync_finalize_projects_validated_responses_stream_metadata_to_chat() {
|
||||
let report_context = json!({
|
||||
"provider_api_format": "openai:responses",
|
||||
"provider_stream_event_api_format": "openai:responses",
|
||||
"client_api_format": "openai:chat",
|
||||
"model": "gpt-5.6-luna",
|
||||
"mapped_model": "gpt-5.6-luna",
|
||||
"needs_conversion": true,
|
||||
});
|
||||
let stream_body = concat!(
|
||||
"data: {\"type\":\"response.output_text.delta\",\"response_id\":\"resp_stream_metadata_123\",\"output_index\":0,\"content_index\":0,\"delta\":\"stream projection\"}\n\n",
|
||||
"data: {\"type\":\"response.output_text.done\",\"response_id\":\"resp_stream_metadata_123\",\"output_index\":0,\"content_index\":0,\"text\":\"stream projection\"}\n\n",
|
||||
"data: {\"type\":\"response.output_item.done\",\"output_index\":0,\"item\":{\"id\":\"msg_stream_metadata_123\",\"type\":\"message\",\"status\":\"completed\",\"role\":\"assistant\",\"phase\":\"final_answer\",\"content\":[{\"type\":\"output_text\",\"text\":\"stream projection\",\"annotations\":[],\"logprobs\":[]}]}}\n\n",
|
||||
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_stream_metadata_123\",\"object\":\"response\",\"status\":\"completed\",\"model\":\"gpt-5.6-luna\",\"background\":false,\"max_output_tokens\":null,\"output\":[],\"usage\":{\"input_tokens\":2,\"output_tokens\":3,\"total_tokens\":5}}}\n\n",
|
||||
);
|
||||
|
||||
let product = maybe_build_standard_sync_finalize_product_from_normalized_payload(
|
||||
"openai_chat_sync_finalize",
|
||||
200,
|
||||
Some(&report_context),
|
||||
None,
|
||||
Some(&base64::engine::general_purpose::STANDARD.encode(stream_body)),
|
||||
)
|
||||
.expect("validated Responses stream should aggregate")
|
||||
.expect("validated Responses stream should produce a Chat projection");
|
||||
|
||||
let StandardSyncFinalizeNormalizedProduct::CrossFormat(product) = product else {
|
||||
panic!("Responses stream should retain a cross-format provider product")
|
||||
};
|
||||
assert_eq!(product.provider_body_json["background"], false);
|
||||
assert_eq!(
|
||||
product.provider_body_json["output"][0]["phase"],
|
||||
"final_answer"
|
||||
);
|
||||
assert_eq!(
|
||||
product.client_body_json["choices"][0]["message"]["content"],
|
||||
"stream projection"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn standard_sync_finalize_projects_authoritative_incomplete_responses_stream() {
|
||||
let report_context = json!({
|
||||
"provider_api_format": "openai:responses",
|
||||
"provider_stream_event_api_format": "openai:responses",
|
||||
"client_api_format": "openai:chat",
|
||||
"model": "gpt-5.6-luna",
|
||||
"mapped_model": "gpt-5.6-luna",
|
||||
"needs_conversion": true,
|
||||
});
|
||||
let stream_body = concat!(
|
||||
"data: {\"type\":\"response.output_text.delta\",\"response_id\":\"resp_incomplete_stream_123\",\"output_index\":0,\"content_index\":0,\"delta\":\"partial output\"}\n\n",
|
||||
"data: {\"type\":\"response.incomplete\",\"response\":{\"id\":\"resp_incomplete_stream_123\",\"object\":\"response\",\"status\":\"incomplete\",\"model\":\"gpt-5.6-luna\",\"background\":false,\"output\":[],\"incomplete_details\":{\"reason\":\"max_output_tokens\"}}}\n\n",
|
||||
);
|
||||
|
||||
let product = maybe_build_standard_sync_finalize_product_from_normalized_payload(
|
||||
"openai_chat_sync_finalize",
|
||||
200,
|
||||
Some(&report_context),
|
||||
None,
|
||||
Some(&base64::engine::general_purpose::STANDARD.encode(stream_body)),
|
||||
)
|
||||
.expect("authoritative incomplete stream should aggregate")
|
||||
.expect("authoritative incomplete stream should produce a Chat projection");
|
||||
|
||||
let StandardSyncFinalizeNormalizedProduct::CrossFormat(product) = product else {
|
||||
panic!("incomplete Responses stream should retain a cross-format provider product")
|
||||
};
|
||||
assert_eq!(product.provider_body_json["status"], "incomplete");
|
||||
assert_eq!(
|
||||
product.client_body_json["choices"][0]["message"]["content"],
|
||||
"partial output"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn standard_sync_finalize_rejects_responses_stream_without_authoritative_terminal() {
|
||||
let report_context = json!({
|
||||
"provider_api_format": "openai:responses",
|
||||
"provider_stream_event_api_format": "openai:responses",
|
||||
"client_api_format": "openai:chat",
|
||||
"model": "gpt-5.6-luna",
|
||||
"mapped_model": "gpt-5.6-luna",
|
||||
"needs_conversion": true,
|
||||
});
|
||||
let stream_body = concat!(
|
||||
"data: {\"type\":\"response.output_text.delta\",\"response_id\":\"resp_truncated_stream_123\",\"output_index\":0,\"content_index\":0,\"delta\":\"truncated\"}\n\n",
|
||||
"data: {\"type\":\"response.output_text.done\",\"response_id\":\"resp_truncated_stream_123\",\"output_index\":0,\"content_index\":0,\"text\":\"truncated\"}\n\n",
|
||||
);
|
||||
|
||||
assert!(aggregate_openai_responses_stream_sync_response(stream_body.as_bytes()).is_none());
|
||||
let error = maybe_build_standard_sync_finalize_product_from_normalized_payload(
|
||||
"openai_chat_sync_finalize",
|
||||
200,
|
||||
Some(&report_context),
|
||||
None,
|
||||
Some(&base64::engine::general_purpose::STANDARD.encode(stream_body)),
|
||||
)
|
||||
.expect_err("a truncated stream must fail closed");
|
||||
assert!(
|
||||
error
|
||||
.to_string()
|
||||
.contains("missing an authoritative terminal event"),
|
||||
"unexpected error: {error}"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn standard_sync_finalize_rejects_unknown_canonical_output_from_known_stream_event() {
|
||||
let report_context = json!({
|
||||
"provider_api_format": "openai:responses",
|
||||
"provider_stream_event_api_format": "openai:responses",
|
||||
"client_api_format": "openai:chat",
|
||||
"model": "gpt-5.6-luna",
|
||||
"mapped_model": "gpt-5.6-luna",
|
||||
"needs_conversion": true,
|
||||
});
|
||||
let stream_body = concat!(
|
||||
"data: {\"type\":\"response.output_item.done\",\"output_index\":0,\"item\":{\"id\":\"future_item_123\",\"type\":\"future_output\",\"payload\":\"must-not-drop\"}}\n\n",
|
||||
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_unknown_output_123\",\"object\":\"response\",\"status\":\"completed\",\"model\":\"gpt-5.6-luna\",\"background\":false,\"output\":[]}}\n\n",
|
||||
);
|
||||
|
||||
let product = maybe_build_standard_sync_finalize_product_from_normalized_payload(
|
||||
"openai_chat_sync_finalize",
|
||||
200,
|
||||
Some(&report_context),
|
||||
None,
|
||||
Some(&base64::engine::general_purpose::STANDARD.encode(stream_body)),
|
||||
)
|
||||
.expect("known event with unknown output should be inspected");
|
||||
assert!(
|
||||
product.is_none(),
|
||||
"an unknown canonical output block must not be dropped by the Chat projection"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn standard_sync_finalize_keeps_direct_responses_metadata_conversion_strict() {
|
||||
let report_context = json!({
|
||||
"provider_api_format": "openai:responses",
|
||||
"client_api_format": "openai:chat",
|
||||
"model": "gpt-5.6-luna",
|
||||
"mapped_model": "gpt-5.6-luna",
|
||||
"needs_conversion": true,
|
||||
});
|
||||
let provider_body_json = json!({
|
||||
"id": "resp_direct_metadata_123",
|
||||
"object": "response",
|
||||
"status": "completed",
|
||||
"model": "gpt-5.6-luna",
|
||||
"background": false,
|
||||
"output": [{
|
||||
"id": "msg_direct_metadata_123",
|
||||
"type": "message",
|
||||
"status": "completed",
|
||||
"role": "assistant",
|
||||
"content": [{
|
||||
"type": "output_text",
|
||||
"text": "direct body",
|
||||
"annotations": []
|
||||
}]
|
||||
}]
|
||||
});
|
||||
|
||||
let product = maybe_build_standard_sync_finalize_product_from_normalized_payload(
|
||||
"openai_chat_sync_finalize",
|
||||
200,
|
||||
Some(&report_context),
|
||||
Some(&provider_body_json),
|
||||
None,
|
||||
)
|
||||
.expect("direct response inspection should not fail");
|
||||
|
||||
assert!(
|
||||
product.is_none(),
|
||||
"direct provider JSON must keep the strict lossless conversion boundary"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn standard_sync_finalize_never_projects_failed_responses_stream_as_chat_success() {
|
||||
let report_context = json!({
|
||||
"provider_api_format": "openai:responses",
|
||||
"provider_stream_event_api_format": "openai:responses",
|
||||
"client_api_format": "openai:chat",
|
||||
"model": "gpt-5.6-luna",
|
||||
"mapped_model": "gpt-5.6-luna",
|
||||
"needs_conversion": true,
|
||||
});
|
||||
let stream_body = "data: {\"type\":\"response.failed\",\"response\":{\"id\":\"resp_failed_stream_123\",\"object\":\"response\",\"status\":\"failed\",\"model\":\"gpt-5.6-luna\",\"error\":null,\"output\":[{\"type\":\"message\",\"id\":\"msg_failed_stream_123\",\"role\":\"assistant\",\"status\":\"incomplete\",\"content\":[{\"type\":\"output_text\",\"text\":\"partial failure\",\"annotations\":[]}]}]}}\n\n";
|
||||
|
||||
let product = maybe_build_standard_sync_finalize_product_from_normalized_payload(
|
||||
"openai_chat_sync_finalize",
|
||||
200,
|
||||
Some(&report_context),
|
||||
None,
|
||||
Some(&base64::engine::general_purpose::STANDARD.encode(stream_body)),
|
||||
)
|
||||
.expect("failed stream should remain available to the error boundary");
|
||||
|
||||
assert!(
|
||||
product.is_none(),
|
||||
"a failed Responses stream must not become a Chat success body"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn standard_sync_finalize_aggregates_openai_responses_capture_envelope_for_same_family() {
|
||||
let report_context = json!({
|
||||
|
||||
Reference in New Issue
Block a user