feat(ai-formats): add reasoning summary boundaries for streams

This commit is contained in:
zhefox
2026-05-13 17:41:34 +08:00
parent 1726df1169
commit 1adbf23be4
5 changed files with 125 additions and 36 deletions

View File

@@ -25,6 +25,11 @@ pub struct ClaudeProviderState {
finished: bool, finished: bool,
usage: Option<CanonicalUsage>, usage: Option<CanonicalUsage>,
tool_calls: BTreeMap<usize, ClaudeProviderToolState>, tool_calls: BTreeMap<usize, ClaudeProviderToolState>,
/// True while we are inside a thinking content block. Used to emit
/// `ReasoningSummaryDone` when the block closes, so that downstream
/// emitters can insert paragraph separators between distinct thinking
/// blocks (CPA strategy).
in_thinking_block: bool,
} }
impl ClaudeProviderState { impl ClaudeProviderState {
@@ -215,6 +220,7 @@ impl ClaudeProviderState {
.and_then(Value::as_str) .and_then(Value::as_str)
.unwrap_or_default(); .unwrap_or_default();
if block_type == "thinking" { if block_type == "thinking" {
self.in_thinking_block = true;
let Some(piece) = block let Some(piece) = block
.get("thinking") .get("thinking")
.and_then(Value::as_str) .and_then(Value::as_str)
@@ -331,7 +337,22 @@ impl ClaudeProviderState {
}); });
self.finished = true; self.finished = true;
} }
"content_block_stop" | "message_stop" | "ping" => {} "content_block_stop" => {
// CPA strategy: when a thinking block closes, emit
// ReasoningSummaryDone so downstream emitters can insert
// paragraph separators between distinct thinking blocks.
if self.in_thinking_block {
self.in_thinking_block = false;
self.ensure_started(report_context, &mut out);
let (id, model) = self.identity(report_context);
out.push(CanonicalStreamFrame {
id,
model,
event: CanonicalStreamEvent::ReasoningSummaryDone,
});
}
}
"message_stop" | "ping" => {}
_ => { _ => {
out.push(self.unknown_frame(report_context, value.clone())); out.push(self.unknown_frame(report_context, value.clone()));
} }
@@ -746,6 +767,12 @@ impl ClaudeClientEmitter {
self.finished = true; self.finished = true;
Ok(out) Ok(out)
} }
CanonicalStreamEvent::ReasoningSummaryDone => {
// CPA strategy: close the current thinking block so the next
// ReasoningDelta opens a fresh one. Each reasoning paragraph
// becomes its own thinking block in Claude's wire format.
Ok(self.close_open_block().unwrap_or_default())
}
} }
} }

View File

@@ -551,6 +551,7 @@ impl GeminiClientEmitter {
self.finished = true; self.finished = true;
Ok(out) Ok(out)
} }
CanonicalStreamEvent::ReasoningSummaryDone => Ok(Vec::new()),
} }
} }

View File

@@ -741,15 +741,9 @@ impl OpenAIResponsesProviderState {
} }
} }
"response.reasoning_summary_part.added" | "response.reasoning_summary_part.done" => { "response.reasoning_summary_part.added" | "response.reasoning_summary_part.done" => {
if let Some(part) = value.get("part").and_then(Value::as_object) { // SKIP — CPA ignores these structural events entirely.
if part.get("type").and_then(Value::as_str) == Some("summary_text") { // Only reasoning_summary_text.delta carries incremental text;
if let Some(text) = part.get("text").and_then(Value::as_str) { // reasoning_summary_text.done signals a paragraph boundary.
if !text.is_empty() {
self.emit_missing_reasoning(report_context, &mut out, text);
}
}
}
}
} }
"response.output_text.done" => { "response.output_text.done" => {
let text = value let text = value
@@ -784,20 +778,16 @@ impl OpenAIResponsesProviderState {
} }
} }
"response.reasoning_summary_text.done" => { "response.reasoning_summary_text.done" => {
let text = value // CPA strategy: emit a section-end signal so downstream can insert
.get("text") // paragraph separators (e.g. "\n\n" for Chat). Do NOT re-emit
.and_then(Value::as_str) // the full text — that would duplicate what deltas already sent.
.or_else(|| { self.ensure_started(report_context, &mut out);
value let (id, model) = self.identity(report_context);
.get("part") out.push(CanonicalStreamFrame {
.and_then(Value::as_object) id,
.and_then(|part| part.get("text")) model,
.and_then(Value::as_str) event: CanonicalStreamEvent::ReasoningSummaryDone,
}) });
.unwrap_or_default();
if !text.is_empty() {
self.emit_missing_reasoning(report_context, &mut out, text);
}
} }
"response.output_item.added" => { "response.output_item.added" => {
let Some(item) = value.get("item").and_then(Value::as_object) else { let Some(item) = value.get("item").and_then(Value::as_object) else {
@@ -818,7 +808,9 @@ impl OpenAIResponsesProviderState {
self.emit_message_item(report_context, &mut out, item); self.emit_message_item(report_context, &mut out, item);
} }
"reasoning" => { "reasoning" => {
self.emit_reasoning_item(report_context, &mut out, item); // 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);
} }
_ => { _ => {
out.push(self.unknown_frame(report_context, Value::Object(item.clone()))); out.push(self.unknown_frame(report_context, Value::Object(item.clone())));
@@ -1031,7 +1023,8 @@ impl OpenAIResponsesProviderState {
self.emit_message_item(report_context, &mut out, item); self.emit_message_item(report_context, &mut out, item);
} }
"reasoning" => { "reasoning" => {
self.emit_reasoning_item(report_context, &mut out, item); // SKIP — CPA ignores reasoning output_item.done to avoid
// re-emitting full summary text that deltas already sent.
} }
_ => { _ => {
out.push(self.unknown_frame(report_context, Value::Object(item.clone()))); out.push(self.unknown_frame(report_context, Value::Object(item.clone())));
@@ -1076,7 +1069,7 @@ impl OpenAIResponsesProviderState {
); );
} }
"reasoning" => { "reasoning" => {
self.emit_reasoning_item(report_context, &mut out, item); // SKIP — deltas already captured via reasoning_summary_text.delta
} }
_ => { _ => {
out.push( out.push(
@@ -1244,6 +1237,29 @@ impl OpenAIChatClientEmitter {
)?); )?);
Ok(out) Ok(out)
} }
CanonicalStreamEvent::ReasoningSummaryDone => {
// CPA strategy: emit "\n\n" as paragraph separator between
// reasoning sections, matching CPA's Chat downstream behavior.
let mut out = self.ensure_started()?;
out.extend(encode_json_sse(
None,
&json!({
"id": self.response_id
.as_deref()
.unwrap_or("chatcmpl-local-stream"),
"object": "chat.completion.chunk",
"model": self.model.as_deref().unwrap_or("unknown"),
"choices": [{
"index": 0,
"delta": {
"reasoning_content": "\n\n",
},
"finish_reason": Value::Null
}]
}),
)?);
Ok(out)
}
CanonicalStreamEvent::ReasoningSignature(_) => Ok(Vec::new()), CanonicalStreamEvent::ReasoningSignature(_) => Ok(Vec::new()),
CanonicalStreamEvent::ContentPart(part) => { CanonicalStreamEvent::ContentPart(part) => {
let placeholder = openai_stream_placeholder_for_content_part(&part); let placeholder = openai_stream_placeholder_for_content_part(&part);
@@ -1317,10 +1333,10 @@ impl OpenAIChatClientEmitter {
Ok(out) Ok(out)
} }
CanonicalStreamEvent::ToolResultDelta { CanonicalStreamEvent::ToolResultDelta {
index,
tool_use_id, tool_use_id,
name, name,
content, content,
..
} => { } => {
let mut out = self.ensure_started()?; let mut out = self.ensure_started()?;
let mut delta = Map::new(); let mut delta = Map::new();
@@ -1808,10 +1824,15 @@ impl OpenAIResponsesClientEmitter {
); );
item.insert("id".to_string(), Value::String(format!("{item_id}_output"))); item.insert("id".to_string(), Value::String(format!("{item_id}_output")));
item.insert("call_id".to_string(), Value::String(item_id)); item.insert("call_id".to_string(), Value::String(item_id));
if let Some(name) = state.name.filter(|value| !value.trim().is_empty()) { if let Some(name) = state
.name
.as_ref()
.filter(|value| !value.trim().is_empty())
.cloned()
{
item.insert("name".to_string(), Value::String(name)); item.insert("name".to_string(), Value::String(name));
} }
item.insert("output".to_string(), Value::String(state.content)); item.insert("output".to_string(), Value::String(state.content.clone()));
out.extend(self.encode_response_event( out.extend(self.encode_response_event(
"response.output_item.done", "response.output_item.done",
json!({ json!({
@@ -1975,6 +1996,45 @@ impl OpenAIResponsesClientEmitter {
)?); )?);
Ok(out) 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.
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 mut out = Vec::new();
out.extend(self.encode_response_event(
"response.reasoning_summary_text.done",
json!({
"type": "response.reasoning_summary_text.done",
"response_id": self.response_id(),
"item_id": item_id.clone(),
"output_index": output_index,
"summary_index": 0,
"text": self.reasoning.as_str(),
}),
)?);
out.extend(self.encode_response_event(
"response.reasoning_summary_part.done",
json!({
"type": "response.reasoning_summary_part.done",
"response_id": self.response_id(),
"item_id": item_id,
"output_index": output_index,
"summary_index": 0,
"part": {
"type": "summary_text",
"text": self.reasoning.as_str(),
}
}),
)?);
// Reset part state so next ReasoningDelta opens a new part
self.reasoning_part_started = false;
Ok(out)
}
CanonicalStreamEvent::ReasoningSignature(_) => Ok(Vec::new()), CanonicalStreamEvent::ReasoningSignature(_) => Ok(Vec::new()),
CanonicalStreamEvent::ContentPart(part) => { CanonicalStreamEvent::ContentPart(part) => {
let placeholder = openai_stream_placeholder_for_content_part(&part); let placeholder = openai_stream_placeholder_for_content_part(&part);

View File

@@ -1018,10 +1018,8 @@ fn convert_openai_responses_canonical_responses_response(
.as_ref() .as_ref()
.and_then(|canonical| canonical_to_gemini_response(canonical, report_context)) .and_then(|canonical| canonical_to_gemini_response(canonical, report_context))
.or_else(|| { .or_else(|| {
let openai_chat = convert_openai_responses_response_to_openai_chat( let openai_chat =
body_json, convert_openai_responses_response_to_openai_chat(body_json, report_context)?;
report_context,
)?;
convert_openai_chat_response_to_gemini_chat(&openai_chat, report_context) convert_openai_chat_response_to_gemini_chat(&openai_chat, report_context)
}) })
} }
@@ -2593,6 +2591,7 @@ pub fn aggregate_gemini_stream_sync_response(body: &[u8]) -> Option<Value> {
)); ));
} }
CanonicalStreamEvent::UnknownEvent(_) => {} CanonicalStreamEvent::UnknownEvent(_) => {}
CanonicalStreamEvent::ReasoningSummaryDone => {}
CanonicalStreamEvent::Finish { CanonicalStreamEvent::Finish {
finish_reason: frame_finish_reason, finish_reason: frame_finish_reason,
usage, usage,
@@ -3630,8 +3629,9 @@ mod tests {
fn rejects_openai_responses_same_family_error_body_json() { fn rejects_openai_responses_same_family_error_body_json() {
let report_context = json!({ let report_context = json!({
"provider_api_format": "openai:responses", "provider_api_format": "openai:responses",
"client_api_format": "openai:responses", "client_api_format": "openai:responses:compact",
"needs_conversion": false, "model": "gpt-5",
"mapped_model": "gpt-5",
}); });
let provider_body_json = json!({ let provider_body_json = json!({
"error": { "error": {

View File

@@ -35,6 +35,7 @@ pub enum CanonicalStreamEvent {
Start, Start,
TextDelta(String), TextDelta(String),
ReasoningDelta(String), ReasoningDelta(String),
ReasoningSummaryDone,
ReasoningSignature(String), ReasoningSignature(String),
ContentPart(CanonicalContentPart), ContentPart(CanonicalContentPart),
ToolCallStart { ToolCallStart {