Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
94 changes: 89 additions & 5 deletions crates/switchyard-translation/src/codecs/responses/stream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
return ensure_responses_created(state);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
let mut out = ensure_responses_reasoning_started(state, index);
if !text.is_empty() {
out.extend(encode_responses_reasoning_delta(state, index, text));
Expand Down Expand Up @@ -418,7 +425,8 @@ fn finish_responses_stream(state: &mut StreamTranslationState) -> Vec<Value> {
("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
{
Expand Down Expand Up @@ -860,7 +868,8 @@ fn add_sequence_numbers(state: &mut StreamTranslationState, mut events: Vec<Valu

// Accumulates assistant text and emits Responses text delta events.
fn encode_responses_text_delta(state: &mut StreamTranslationState, text: String) -> Vec<Value> {
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;
Expand Down Expand Up @@ -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<String> {
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<usize>) -> Vec<Value> {
let held: Vec<usize> = 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())
Expand All @@ -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.
Expand All @@ -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
}

Expand Down Expand Up @@ -972,6 +1045,14 @@ fn encode_responses_reasoning_delta(
index: usize,
text: String,
) -> Vec<Value> {
// 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);
Expand Down Expand Up @@ -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;
Expand Down
4 changes: 4 additions & 0 deletions crates/switchyard-translation/src/codecs/stream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,10 @@ pub(crate) struct ResponseReasoningState {
pub(crate) encrypted: Option<String>,
#[serde(default)]
pub(crate) anthropic_signature: Option<String>,
/// 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.
Expand Down
212 changes: 212 additions & 0 deletions crates/switchyard-translation/tests/stream_translation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Value>| -> Result<Vec<Vec<Value>>, Box<dyn std::error::Error + Send + Sync>> {
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<String> {
batch
.iter()
.filter_map(|event| event["type"].as_str().map(ToOwned::to_owned))
.collect()
};
let final_reasoning = |batches: &[Vec<Value>]| -> Result<Value, &'static str> {
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::<String>::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 {
Expand Down
Loading