fix: stabilize openai responses reasoning streams

This commit is contained in:
fawney19
2026-05-14 00:06:11 +08:00
parent 2d354fa294
commit 697b6e0653
3 changed files with 446 additions and 37 deletions

View File

@@ -46,6 +46,7 @@ pub struct OpenAIResponsesProviderState {
finished: bool,
text: String,
reasoning: String,
reasoning_parts: BTreeMap<usize, String>,
tool_calls: BTreeMap<usize, OpenAIResponsesProviderToolState>,
tool_results: BTreeMap<usize, OpenAIResponsesProviderToolResultState>,
tool_index_by_key: BTreeMap<String, usize>,
@@ -469,6 +470,45 @@ impl OpenAIResponsesProviderState {
});
}
fn emit_missing_reasoning_part_text(
&mut self,
report_context: &Value,
out: &mut Vec<CanonicalStreamFrame>,
summary_index: usize,
text: &str,
) {
if text.is_empty() {
return;
}
let missing = {
let current = self.reasoning_parts.entry(summary_index).or_default();
let missing = if text.starts_with(current.as_str()) {
text[current.len()..].to_string()
} else if current.as_str() == text {
String::new()
} else if current.is_empty() {
text.to_string()
} else {
String::new()
};
if !missing.is_empty() {
current.push_str(&missing);
}
missing
};
if missing.is_empty() {
return;
}
self.ensure_started(report_context, out);
self.reasoning.push_str(&missing);
let (id, model) = self.identity(report_context);
out.push(CanonicalStreamFrame {
id,
model,
event: CanonicalStreamEvent::ReasoningDelta(missing),
});
}
fn emit_tool_call_item(
&mut self,
report_context: &Value,
@@ -741,9 +781,23 @@ impl OpenAIResponsesProviderState {
}
}
"response.reasoning_summary_part.added" | "response.reasoning_summary_part.done" => {
// SKIP — CPA ignores these structural events entirely.
// Only reasoning_summary_text.delta carries incremental text;
// reasoning_summary_text.done signals a paragraph boundary.
if let Some(part) = value.get("part").and_then(Value::as_object) {
if part.get("type").and_then(Value::as_str) == Some("summary_text") {
if let Some(text) = part.get("text").and_then(Value::as_str) {
let summary_index = value
.get("summary_index")
.and_then(Value::as_u64)
.map(|value| value as usize)
.unwrap_or(0);
self.emit_missing_reasoning_part_text(
report_context,
&mut out,
summary_index,
text,
);
}
}
}
}
"response.output_text.done" => {
let text = value
@@ -767,8 +821,17 @@ impl OpenAIResponsesProviderState {
.and_then(Value::as_str)
.unwrap_or_default();
if !piece.is_empty() {
let summary_index = value
.get("summary_index")
.and_then(Value::as_u64)
.map(|value| value as usize)
.unwrap_or(0);
self.ensure_started(report_context, &mut out);
self.reasoning.push_str(piece);
self.reasoning_parts
.entry(summary_index)
.or_default()
.push_str(piece);
let (id, model) = self.identity(report_context);
out.push(CanonicalStreamFrame {
id,
@@ -778,9 +841,30 @@ impl OpenAIResponsesProviderState {
}
}
"response.reasoning_summary_text.done" => {
// CPA strategy: emit a section-end signal so downstream can insert
// paragraph separators (e.g. "\n\n" for Chat). Do NOT re-emit
// the full text — that would duplicate what deltas already sent.
let text = value
.get("text")
.and_then(Value::as_str)
.or_else(|| {
value
.get("part")
.and_then(Value::as_object)
.and_then(|part| part.get("text"))
.and_then(Value::as_str)
})
.unwrap_or_default();
if !text.is_empty() {
let summary_index = value
.get("summary_index")
.and_then(Value::as_u64)
.map(|value| value as usize)
.unwrap_or(0);
self.emit_missing_reasoning_part_text(
report_context,
&mut out,
summary_index,
text,
);
}
self.ensure_started(report_context, &mut out);
let (id, model) = self.identity(report_context);
out.push(CanonicalStreamFrame {
@@ -808,8 +892,6 @@ impl OpenAIResponsesProviderState {
self.emit_message_item(report_context, &mut out, item);
}
"reasoning" => {
// SKIP — CPA ignores reasoning output_item.added to avoid
// re-emitting full summary text that deltas already sent.
self.ensure_started(report_context, &mut out);
}
_ => {
@@ -1023,8 +1105,7 @@ impl OpenAIResponsesProviderState {
self.emit_message_item(report_context, &mut out, item);
}
"reasoning" => {
// SKIP — CPA ignores reasoning output_item.done to avoid
// re-emitting full summary text that deltas already sent.
self.emit_reasoning_item(report_context, &mut out, item);
}
_ => {
out.push(self.unknown_frame(report_context, Value::Object(item.clone())));
@@ -1069,7 +1150,7 @@ impl OpenAIResponsesProviderState {
);
}
"reasoning" => {
// SKIP — deltas already captured via reasoning_summary_text.delta
self.emit_reasoning_item(report_context, &mut out, item);
}
_ => {
out.push(
@@ -1170,6 +1251,8 @@ pub struct OpenAIResponsesClientEmitter {
message_output_index: Option<usize>,
text: String,
reasoning: String,
reasoning_part: String,
reasoning_summary_parts: Vec<String>,
tool_calls: BTreeMap<usize, OpenAIResponsesClientToolState>,
tool_results: BTreeMap<usize, OpenAIResponsesClientToolResultState>,
}
@@ -1333,7 +1416,7 @@ impl OpenAIChatClientEmitter {
Ok(out)
}
CanonicalStreamEvent::ToolResultDelta {
index,
index: _,
tool_use_id,
name,
content,
@@ -1531,6 +1614,10 @@ impl OpenAIResponsesClientEmitter {
output_index
}
fn current_reasoning_summary_index(&self) -> usize {
self.reasoning_summary_parts.len()
}
fn ensure_message_output_index(&mut self) -> usize {
if let Some(output_index) = self.message_output_index {
return output_index;
@@ -1587,6 +1674,7 @@ impl OpenAIResponsesClientEmitter {
self.reasoning_item_started = true;
}
if !self.reasoning_part_started {
let summary_index = self.current_reasoning_summary_index();
out.extend(self.encode_response_event(
"response.reasoning_summary_part.added",
json!({
@@ -1594,7 +1682,7 @@ impl OpenAIResponsesClientEmitter {
"response_id": self.response_id(),
"item_id": item_id,
"output_index": output_index,
"summary_index": 0,
"summary_index": summary_index,
"part": {
"type": "summary_text",
"text": "",
@@ -1714,6 +1802,8 @@ impl OpenAIResponsesClientEmitter {
let item_id = self.reasoning_item_id();
let mut out = Vec::new();
if self.reasoning_part_started {
let summary_index = self.current_reasoning_summary_index();
let part_text = self.reasoning_part.clone();
out.extend(self.encode_response_event(
"response.reasoning_summary_text.done",
json!({
@@ -1721,8 +1811,8 @@ impl OpenAIResponsesClientEmitter {
"response_id": self.response_id(),
"item_id": item_id.clone(),
"output_index": output_index,
"summary_index": 0,
"text": self.reasoning.as_str(),
"summary_index": summary_index,
"text": part_text.as_str(),
}),
)?);
out.extend(self.encode_response_event(
@@ -1732,14 +1822,37 @@ impl OpenAIResponsesClientEmitter {
"response_id": self.response_id(),
"item_id": item_id.clone(),
"output_index": output_index,
"summary_index": 0,
"summary_index": summary_index,
"part": {
"type": "summary_text",
"text": self.reasoning.as_str(),
"text": part_text.as_str(),
}
}),
)?);
self.reasoning_summary_parts.push(part_text);
self.reasoning_part.clear();
self.reasoning_part_started = false;
}
let summary = if self.reasoning_summary_parts.is_empty() {
if self.reasoning.trim().is_empty() {
Vec::new()
} else {
vec![json!({
"type": "summary_text",
"text": self.reasoning.as_str(),
})]
}
} else {
self.reasoning_summary_parts
.iter()
.map(|text| {
json!({
"type": "summary_text",
"text": text,
})
})
.collect::<Vec<_>>()
};
out.extend(self.encode_response_event(
"response.output_item.done",
json!({
@@ -1749,10 +1862,7 @@ impl OpenAIResponsesClientEmitter {
"item": {
"type": "reasoning",
"id": item_id,
"summary": [{
"type": "summary_text",
"text": self.reasoning.as_str(),
}],
"summary": summary,
}
}),
)?);
@@ -1848,17 +1958,34 @@ impl OpenAIResponsesClientEmitter {
fn completed_response(&self, usage: CanonicalUsage) -> Value {
let mut ordered_output = Vec::new();
if !self.reasoning.trim().is_empty() {
let summary = if self.reasoning_summary_parts.is_empty() {
if self.reasoning.trim().is_empty() {
Vec::new()
} else {
vec![json!({
"type": "summary_text",
"text": self.reasoning.as_str(),
})]
}
} else {
self.reasoning_summary_parts
.iter()
.map(|text| {
json!({
"type": "summary_text",
"text": text,
})
})
.collect::<Vec<_>>()
};
if !summary.is_empty() {
ordered_output.push((
self.reasoning_output_index.unwrap_or(0),
json!({
"type": "reasoning",
"id": self.reasoning_item_id(),
"status": "completed",
"summary": [{
"type": "summary_text",
"text": self.reasoning.as_str(),
}]
"summary": summary,
}),
));
}
@@ -1983,6 +2110,7 @@ impl OpenAIResponsesClientEmitter {
CanonicalStreamEvent::ReasoningDelta(text) => {
let mut out = self.ensure_reasoning_item_started()?;
self.reasoning.push_str(&text);
self.reasoning_part.push_str(&text);
out.extend(self.encode_response_event(
"response.reasoning_summary_text.delta",
json!({
@@ -1990,21 +2118,22 @@ impl OpenAIResponsesClientEmitter {
"response_id": self.response_id(),
"item_id": self.reasoning_item_id(),
"output_index": self.reasoning_output_index.unwrap_or(0),
"summary_index": 0,
"summary_index": self.current_reasoning_summary_index(),
"delta": text,
}),
)?);
Ok(out)
}
CanonicalStreamEvent::ReasoningSummaryDone => {
// CPA strategy for Responses downstream: close the current summary
// text/part and reset state so the next ReasoningDelta starts a
// fresh part within the same reasoning item.
// Close the current reasoning part and reset state so the next
// ReasoningDelta starts a fresh part within the same item.
if !self.reasoning_item_started || !self.reasoning_part_started {
return Ok(Vec::new());
}
let output_index = self.reasoning_output_index.unwrap_or(0);
let item_id = self.reasoning_item_id();
let summary_index = self.current_reasoning_summary_index();
let part_text = self.reasoning_part.clone();
let mut out = Vec::new();
out.extend(self.encode_response_event(
"response.reasoning_summary_text.done",
@@ -2013,8 +2142,8 @@ impl OpenAIResponsesClientEmitter {
"response_id": self.response_id(),
"item_id": item_id.clone(),
"output_index": output_index,
"summary_index": 0,
"text": self.reasoning.as_str(),
"summary_index": summary_index,
"text": part_text.as_str(),
}),
)?);
out.extend(self.encode_response_event(
@@ -2024,14 +2153,15 @@ impl OpenAIResponsesClientEmitter {
"response_id": self.response_id(),
"item_id": item_id,
"output_index": output_index,
"summary_index": 0,
"summary_index": summary_index,
"part": {
"type": "summary_text",
"text": self.reasoning.as_str(),
"text": part_text.as_str(),
}
}),
)?);
// Reset part state so next ReasoningDelta opens a new part
self.reasoning_summary_parts.push(part_text);
self.reasoning_part.clear();
self.reasoning_part_started = false;
Ok(out)
}
@@ -2298,6 +2428,38 @@ mod tests {
sequence_numbers
}
fn response_reasoning_text_done_parts(sse: &str) -> Vec<(u64, String)> {
let mut parts = Vec::new();
for block in sse.split("\n\n") {
let mut event_name = None;
let mut data = None;
for line in block.lines() {
if let Some(value) = line.strip_prefix("event: ") {
event_name = Some(value);
} else if let Some(value) = line.strip_prefix("data: ") {
data = Some(value);
}
}
if event_name != Some("response.reasoning_summary_text.done") {
continue;
}
let Some(data) = data else {
continue;
};
let Ok(value) = serde_json::from_str::<Value>(data) else {
continue;
};
let Some(summary_index) = value.get("summary_index").and_then(Value::as_u64) else {
continue;
};
let Some(text) = value.get("text").and_then(Value::as_str) else {
continue;
};
parts.push((summary_index, text.to_string()));
}
parts
}
#[test]
fn openai_chat_provider_state_emits_unknown_events_for_unrecognized_deltas() {
let mut state = OpenAIChatProviderState::default();
@@ -3039,6 +3201,72 @@ mod tests {
assert_eq!(response_sequence_numbers(&sse), (1..=9).collect::<Vec<_>>());
}
#[test]
fn openai_responses_client_emitter_closes_distinct_reasoning_parts() {
let mut emitter = OpenAIResponsesClientEmitter::default();
let mut bytes = emitter
.emit(CanonicalStreamFrame {
id: "resp_789".to_string(),
model: "gpt-5.4".to_string(),
event: CanonicalStreamEvent::Start,
})
.expect("start should encode");
bytes.extend(
emitter
.emit(CanonicalStreamFrame {
id: "resp_789".to_string(),
model: "gpt-5.4".to_string(),
event: CanonicalStreamEvent::ReasoningDelta("alpha".to_string()),
})
.expect("first reasoning should encode"),
);
bytes.extend(
emitter
.emit(CanonicalStreamFrame {
id: "resp_789".to_string(),
model: "gpt-5.4".to_string(),
event: CanonicalStreamEvent::ReasoningSummaryDone,
})
.expect("first boundary should encode"),
);
bytes.extend(
emitter
.emit(CanonicalStreamFrame {
id: "resp_789".to_string(),
model: "gpt-5.4".to_string(),
event: CanonicalStreamEvent::ReasoningDelta("beta".to_string()),
})
.expect("second reasoning should encode"),
);
bytes.extend(
emitter
.emit(CanonicalStreamFrame {
id: "resp_789".to_string(),
model: "gpt-5.4".to_string(),
event: CanonicalStreamEvent::ReasoningSummaryDone,
})
.expect("second boundary should encode"),
);
bytes.extend(
emitter
.emit(CanonicalStreamFrame {
id: "resp_789".to_string(),
model: "gpt-5.4".to_string(),
event: CanonicalStreamEvent::Finish {
finish_reason: Some("stop".to_string()),
usage: None,
},
})
.expect("finish should encode"),
);
let sse = String::from_utf8(bytes).expect("sse should be utf8");
assert_eq!(
response_reasoning_text_done_parts(&sse),
vec![(0, "alpha".to_string()), (1, "beta".to_string())]
);
}
#[test]
fn openai_responses_client_emitter_emits_failed_event_with_sequence_number() {
let mut emitter = OpenAIResponsesClientEmitter::default();
@@ -3156,4 +3384,181 @@ mod tests {
CanonicalStreamEvent::ReasoningDelta(ref text) if text == "step"
));
}
#[test]
fn openai_responses_provider_state_accepts_reasoning_done_without_delta() {
let mut state = OpenAIResponsesProviderState::default();
let report_context = json!({});
let mut frames = Vec::new();
frames.extend(
state
.push_line(
&report_context,
data_line(json!({
"type": "response.created",
"response": {
"id": "resp_done_only",
"model": "gpt-5.4",
}
})),
)
.expect("created should parse"),
);
frames.extend(
state
.push_line(
&report_context,
data_line(json!({
"type": "response.reasoning_summary_text.done",
"response_id": "resp_done_only",
"item_id": "resp_done_only_rs_0",
"output_index": 0,
"summary_index": 0,
"text": "fallback reasoning",
})),
)
.expect("reasoning done should parse"),
);
assert!(frames.iter().any(|frame| matches!(
frame.event,
CanonicalStreamEvent::ReasoningDelta(ref text) if text == "fallback reasoning"
)));
assert!(frames
.iter()
.any(|frame| matches!(frame.event, CanonicalStreamEvent::ReasoningSummaryDone)));
}
#[test]
fn openai_responses_provider_state_does_not_duplicate_part_scoped_reasoning_done() {
let mut state = OpenAIResponsesProviderState::default();
let report_context = json!({});
let mut frames = Vec::new();
for event in [
json!({
"type": "response.reasoning_summary_text.delta",
"response_id": "resp_parts",
"item_id": "resp_parts_rs_0",
"output_index": 0,
"summary_index": 0,
"delta": "alpha",
}),
json!({
"type": "response.reasoning_summary_text.done",
"response_id": "resp_parts",
"item_id": "resp_parts_rs_0",
"output_index": 0,
"summary_index": 0,
"text": "alpha",
}),
json!({
"type": "response.reasoning_summary_text.delta",
"response_id": "resp_parts",
"item_id": "resp_parts_rs_0",
"output_index": 0,
"summary_index": 1,
"delta": "beta",
}),
json!({
"type": "response.reasoning_summary_text.done",
"response_id": "resp_parts",
"item_id": "resp_parts_rs_0",
"output_index": 0,
"summary_index": 1,
"text": "beta",
}),
] {
frames.extend(
state
.push_line(&report_context, data_line(event))
.expect("reasoning event should parse"),
);
}
let reasoning = frames
.iter()
.filter_map(|frame| match &frame.event {
CanonicalStreamEvent::ReasoningDelta(text) => Some(text.as_str()),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(reasoning, vec!["alpha", "beta"]);
}
#[test]
fn openai_responses_provider_state_does_not_duplicate_added_reasoning_item_summary() {
let mut state = OpenAIResponsesProviderState::default();
let report_context = json!({});
let mut frames = Vec::new();
for event in [
json!({
"type": "response.output_item.added",
"response_id": "resp_added_summary",
"output_index": 0,
"item": {
"type": "reasoning",
"id": "resp_added_summary_rs_0",
"summary": [{
"type": "summary_text",
"text": "alpha",
}]
}
}),
json!({
"type": "response.reasoning_summary_text.delta",
"response_id": "resp_added_summary",
"item_id": "resp_added_summary_rs_0",
"output_index": 0,
"summary_index": 0,
"delta": "alpha",
}),
] {
frames.extend(
state
.push_line(&report_context, data_line(event))
.expect("reasoning event should parse"),
);
}
let reasoning = frames
.iter()
.filter_map(|frame| match &frame.event {
CanonicalStreamEvent::ReasoningDelta(text) => Some(text.as_str()),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(reasoning, vec!["alpha"]);
}
#[test]
fn openai_responses_provider_state_uses_reasoning_item_as_fallback() {
let mut state = OpenAIResponsesProviderState::default();
let report_context = json!({});
let frames = state
.push_line(
&report_context,
data_line(json!({
"type": "response.output_item.done",
"response_id": "resp_item_fallback",
"output_index": 0,
"item": {
"type": "reasoning",
"id": "resp_item_fallback_rs_0",
"summary": [{
"type": "summary_text",
"text": "item fallback reasoning",
}]
}
})),
)
.expect("reasoning item should parse");
assert!(frames.iter().any(|frame| matches!(
frame.event,
CanonicalStreamEvent::ReasoningDelta(ref text) if text == "item fallback reasoning"
)));
}
}

View File

@@ -1018,8 +1018,10 @@ fn convert_openai_responses_canonical_responses_response(
.as_ref()
.and_then(|canonical| canonical_to_gemini_response(canonical, report_context))
.or_else(|| {
let openai_chat =
convert_openai_responses_response_to_openai_chat(body_json, report_context)?;
let openai_chat = convert_openai_responses_response_to_openai_chat(
body_json,
report_context,
)?;
convert_openai_chat_response_to_gemini_chat(&openai_chat, report_context)
})
}