diff --git a/crates/switchyard-translation/src/codecs/responses/stream.rs b/crates/switchyard-translation/src/codecs/responses/stream.rs index a179d4ee7..7776af9d1 100644 --- a/crates/switchyard-translation/src/codecs/responses/stream.rs +++ b/crates/switchyard-translation/src/codecs/responses/stream.rs @@ -355,7 +355,7 @@ fn encode_responses_stream( .push_str(signature); } } - match encrypted_reasoning_item_id(&details) { + match reasoning_details_item_id(&details) { Some(id) if item.started && item.item_id.as_deref() != Some(id.as_str()) => { // The item already opened under another id; the payload would fail // verification under it, so drop the payload rather than poison the replay. @@ -373,6 +373,13 @@ fn encode_responses_stream( } } } + // A summary that arrives before its item id would open the item under a synthesized + // id, and the encrypted payload that follows binds only to the provider's id. Hold + // the text until that id arrives, or until later output shows none is coming. + if !item.started && item.item_id.is_none() && summary_lacks_item_id(&details) { + item.pending_text.push_str(&text); + return ensure_responses_created(state); + } let mut out = ensure_responses_reasoning_started(state, index); if !text.is_empty() { out.extend(encode_responses_reasoning_delta(state, index, text)); @@ -418,7 +425,8 @@ fn finish_responses_stream(state: &mut StreamTranslationState) -> Vec { ("response.completed", "completed") }; let incomplete_details = is_truncated.then(|| json!({ "reason": "max_output_tokens" })); - let mut out = ensure_responses_created(state); + let mut out = open_held_reasoning(state, None); + out.extend(ensure_responses_created(state)); if state.response_text_started && let Some(output_index) = state.response_text_output_index { @@ -860,7 +868,8 @@ fn add_sequence_numbers(state: &mut StreamTranslationState, mut events: Vec Vec { - let mut out = ensure_responses_created(state); + let mut out = open_held_reasoning(state, None); + out.extend(ensure_responses_created(state)); if !state.response_text_started { state.response_text_started = true; let output_index = state.next_response_output_index; @@ -897,8 +906,65 @@ fn encode_responses_text_delta(state: &mut StreamTranslationState, text: String) out } -// The id of the reasoning item encoded for a source index: the provider's own id when -// encrypted reasoning binds to it, else synthesized from the emitted output index. +// The provider item id named by a delta's reasoning details. The encrypted payload's id wins +// because a replay verifies against it; a summary or text detail names the same item when it +// arrives first. +fn reasoning_details_item_id(details: &[Value]) -> Option { + encrypted_reasoning_item_id(details).or_else(|| { + details + .iter() + .filter_map(Value::as_object) + .filter(|detail| { + matches!( + detail.get("type").and_then(Value::as_str), + Some("reasoning.summary" | "reasoning.text") + ) + }) + .find_map(|detail| { + detail + .get("id") + .and_then(Value::as_str) + .filter(|id| !id.is_empty()) + }) + .map(ToOwned::to_owned) + }) +} + +// True when a `reasoning.summary` detail arrived without the id of the item it belongs to. +fn summary_lacks_item_id(details: &[Value]) -> bool { + details.iter().filter_map(Value::as_object).any(|detail| { + detail.get("type").and_then(Value::as_str) == Some("reasoning.summary") + && detail + .get("id") + .and_then(Value::as_str) + .is_none_or(str::is_empty) + }) +} + +// Opens every reasoning item still holding summary text, under a synthesized id. Output that +// is about to open another item proves no provider id is coming for them, and they must keep +// their place ahead of it. When that output is itself a reasoning item, `before` limits the +// flush to held items at earlier indices; the item's own later details can still name its id. +fn open_held_reasoning(state: &mut StreamTranslationState, before: Option) -> Vec { + let held: Vec = state + .response_reasoning + .iter() + .filter(|(index, item)| { + before.is_none_or(|bound| **index < bound) + && !item.started + && !item.pending_text.is_empty() + }) + .map(|(index, _)| *index) + .collect(); + let mut out = Vec::new(); + for index in held { + out.extend(ensure_responses_reasoning_started(state, index)); + } + out +} + +// The id of the reasoning item encoded for a source index: the provider's own id when a +// reasoning detail named one, else synthesized from the emitted output index. fn responses_reasoning_item_id(state: &StreamTranslationState, index: usize) -> String { let item = state.response_reasoning.get(&index); item.and_then(|item| item.item_id.clone()) @@ -921,11 +987,14 @@ fn ensure_responses_reasoning_started( { return out; } + // Held summaries of earlier items open first, so they keep their place in the output. + out.extend(open_held_reasoning(state, Some(index))); let output_index = state.next_response_output_index; state.next_response_output_index += 1; let item = state.response_reasoning.entry(index).or_default(); item.started = true; item.output_index = Some(output_index); + let pending = std::mem::take(&mut item.pending_text); let item_id = responses_reasoning_item_id(state, index); // Standard Responses shape: reasoning text lives in `summary` as `summary_text` // parts. Clients such as Codex record reasoning items only in this shape. @@ -939,6 +1008,10 @@ fn ensure_responses_reasoning_started( "summary": [], }, })); + // Summary text held for the item's id streams now, under the id it opened with. + if !pending.is_empty() { + out.extend(encode_responses_reasoning_delta(state, index, pending)); + } out } @@ -972,6 +1045,14 @@ fn encode_responses_reasoning_delta( index: usize, text: String, ) -> Vec { + // Text for an item still waiting on its id joins the held summary. + if let Some(item) = state.response_reasoning.get_mut(&index) + && !item.started + && !item.pending_text.is_empty() + { + item.pending_text.push_str(&text); + return Vec::new(); + } // An empty delta carries nothing to show and must not open a part that would never close. if text.is_empty() { return ensure_responses_reasoning_started(state, index); @@ -1016,8 +1097,11 @@ fn encode_responses_tool_delta( let Some(name) = tool.name.clone() else { return out; }; + // The call opens its item next; held reasoning opens first to keep its place. + out.extend(open_held_reasoning(state, None)); let output_index = state.next_response_output_index; state.next_response_output_index += 1; + let tool = state.tool_states.entry(index).or_default(); tool.response_output_index = Some(output_index); tool.response_item_id = Some(responses_item_id_from(&resp_id, "fc", output_index)); tool.started = true; diff --git a/crates/switchyard-translation/src/codecs/stream.rs b/crates/switchyard-translation/src/codecs/stream.rs index 87da2d649..db83015a7 100644 --- a/crates/switchyard-translation/src/codecs/stream.rs +++ b/crates/switchyard-translation/src/codecs/stream.rs @@ -109,6 +109,10 @@ pub(crate) struct ResponseReasoningState { pub(crate) encrypted: Option, #[serde(default)] pub(crate) anthropic_signature: Option, + /// Summary text held back while the item's provider id is still unknown. The item opens + /// once that id arrives, or once later output shows none is coming. + #[serde(default)] + pub(crate) pending_text: String, } // Tracks an in-progress streamed tool call across provider-specific deltas. diff --git a/crates/switchyard-translation/tests/stream_translation.rs b/crates/switchyard-translation/tests/stream_translation.rs index 01e9cc0da..6869e4545 100644 --- a/crates/switchyard-translation/tests/stream_translation.rs +++ b/crates/switchyard-translation/tests/stream_translation.rs @@ -1277,6 +1277,218 @@ fn openai_chat_stream_uses_summary_when_detail_text_is_empty() -> TestResult { Ok(()) } +// A Chat `reasoning_details` stream may send summary text before the id of the item it +// belongs to. The Responses encoder must not open the item under a synthesized id, because the +// encrypted payload that follows binds only to the provider's id. It holds the summary until +// that id arrives, or until later output shows none is coming, and it takes the id from a +// summary detail that carries one. +#[test] +fn openai_chat_stream_summary_before_id_keeps_reasoning_id_and_payload() -> TestResult { + let engine = TranslationEngine::default(); + let chunk = |delta: Value, finish_reason: Option<&str>| { + json!({ + "id": "chatcmpl-reasoning", + "object": "chat.completion.chunk", + "model": REASONING_MODEL, + "choices": [{"index": 0, "delta": delta, "finish_reason": finish_reason}] + }) + }; + let summary = |id: Option<&str>| { + let mut detail = json!({"type": "reasoning.summary", "summary": "summary first"}); + if let Some(id) = id { + detail["id"] = json!(id); + } + json!({"reasoning_details": [detail]}) + }; + let encrypted = json!({"reasoning_details": [{ + "type": "reasoning.encrypted", "id": "rs_provider", "data": "opaque-stream" + }]}); + // Translates each delta as its own chunk and returns the events per chunk, then the + // events of the stop chunk and of finishing the stream. + let translate = + |deltas: Vec| -> Result>, Box> { + let mut state = + StreamTranslationState::new(WireFormat::OpenAiChat, WireFormat::OpenAiResponses); + let mut batches = Vec::new(); + for (delta, finish_reason) in deltas + .into_iter() + .map(|delta| (delta, None)) + .chain([(json!({}), Some("stop"))]) + { + batches.push(engine.translate_event( + &mut state, + WireFormat::OpenAiChat, + WireFormat::OpenAiResponses, + &chunk(delta, finish_reason), + )?); + } + batches.push(engine.finish_stream(&mut state, WireFormat::OpenAiResponses)?); + Ok(batches) + }; + let types = |batch: &[Value]| -> Vec { + batch + .iter() + .filter_map(|event| event["type"].as_str().map(ToOwned::to_owned)) + .collect() + }; + let final_reasoning = |batches: &[Vec]| -> Result { + batches + .iter() + .flatten() + .find(|event| event["type"] == "response.completed") + .and_then(|event| event["response"]["output"].as_array()) + .and_then(|output| output.iter().find(|item| item["type"] == "reasoning")) + .cloned() + .ok_or("expected a reasoning item in the completed output") + }; + let opened = [ + "response.output_item.added", + "response.reasoning_summary_part.added", + "response.reasoning_summary_text.delta", + ]; + let expected = json!({ + "type": "reasoning", + "id": "rs_provider", + "status": "completed", + "summary": [{"type": "summary_text", "text": "summary first"}], + "encrypted_content": "opaque-stream", + }); + + // Summary before the id: the item waits, then opens under the id the payload names. + let batches = translate(vec![ + json!({"role": "assistant"}), + summary(None), + encrypted.clone(), + ])?; + assert_eq!(types(&batches[1]), Vec::::new()); + assert_eq!(types(&batches[2]), opened); + assert_eq!(final_reasoning(&batches)?, expected); + for event in batches.iter().flatten().filter(|event| { + event["type"] + .as_str() + .is_some_and(|kind| kind.starts_with("response.reasoning_summary")) + }) { + assert_eq!(event["item_id"], "rs_provider", "{event}"); + } + + // Summary carrying the id: the item opens at once under it and the payload still attaches. + let batches = translate(vec![ + json!({"role": "assistant"}), + summary(Some("rs_provider")), + encrypted, + ])?; + assert_eq!(types(&batches[1]), opened); + assert_eq!(batches[1][0]["item"]["id"], "rs_provider"); + assert_eq!(final_reasoning(&batches)?, expected); + + // No id ever comes: the answer opens the held item, under a synthesized id, ahead of itself. + let batches = translate(vec![ + json!({"role": "assistant"}), + summary(None), + json!({"content": "answer"}), + ])?; + let opened_items: Vec<&str> = batches + .iter() + .flatten() + .filter(|event| event["type"] == "response.output_item.added") + .filter_map(|event| event["item"]["type"].as_str()) + .collect(); + assert_eq!(opened_items, ["reasoning", "message"]); + let reasoning = final_reasoning(&batches)?; + assert_eq!( + reasoning["summary"], + json!([{"type": "summary_text", "text": "summary first"}]) + ); + assert!( + reasoning["id"] + .as_str() + .is_some_and(|id| id.starts_with("rs_")), + "{reasoning}" + ); + assert!(reasoning.get("encrypted_content").is_none(), "{reasoning}"); + Ok(()) +} + +// A held summary opens ahead of any later item. A tool delta that carries only an id emits no +// item yet, so it leaves the summary held; a reasoning item at a later index opens the held one +// first, under a synthesized id, so the output keeps stream order. +#[test] +fn responses_stream_held_reasoning_opens_ahead_of_later_items() -> TestResult { + let engine = TranslationEngine::default(); + let format = WireFormat::OpenAiResponses; + let mut state = StreamTranslationState::new(format, format); + let mut events = engine.encode_stream_event( + &mut state, + format, + LlmResponseStreamEvent::new(vec![ + LlmResponseChunk::MessageStart { + id: Some("resp_1".into()), + model: Some(REASONING_MODEL.into()), + }, + LlmResponseChunk::ReasoningDetailsDelta { + index: 0, + details: vec![json!({"type": "reasoning.summary", "summary": "held"})], + text: "held".into(), + }, + LlmResponseChunk::ToolCallDelta { + index: 1, + id: Some("call_1".into()), + name: None, + arguments_delta: None, + }, + ]), + )?; + assert!( + events + .iter() + .all(|event| event["type"] != "response.output_item.added"), + "{events:?}" + ); + events.extend(engine.encode_stream_event( + &mut state, + format, + LlmResponseStreamEvent::new(vec![LlmResponseChunk::ReasoningDetailsDelta { + index: 2, + details: vec![json!({"type": "reasoning.encrypted", "id": "rs_2", "data": "opaque"})], + text: String::new(), + }]), + )?); + events.extend(engine.finish_stream(&mut state, format)?); + + let added: Vec<(u64, &str)> = events + .iter() + .filter(|event| event["type"] == "response.output_item.added") + .map(|event| { + ( + event["output_index"].as_u64().unwrap_or(u64::MAX), + event["item"]["id"].as_str().unwrap_or_default(), + ) + }) + .collect(); + assert_eq!(added.len(), 2, "{added:?}"); + assert_eq!(added[0].0, 0); + assert!( + added[0].1.starts_with("rs_") && added[0].1 != "rs_2", + "{added:?}" + ); + assert_eq!(added[1], (1, "rs_2")); + let completed = events + .iter() + .find(|event| event["type"] == "response.completed") + .ok_or("expected response.completed")?; + let output = completed["response"]["output"] + .as_array() + .ok_or("output should be an array")?; + assert_eq!(output.len(), 2, "{output:?}"); + assert_eq!( + output[0]["summary"], + json!([{"type": "summary_text", "text": "held"}]) + ); + assert_eq!(output[1]["id"], "rs_2"); + assert_eq!(output[1]["encrypted_content"], "opaque"); + Ok(()) +} + // Verifies Anthropic thinking deltas become OpenAI reasoning_content, not content. #[test] fn anthropic_thinking_stream_deltas_do_not_become_openai_chat_content() -> TestResult {