fix(ai): 按 Responses 文本分片去重快照

upstream 已有 8abedecb 处理单个 OpenAI Responses 文本流中 delta 与 done/completed 快照重复输出的问题。

本提交保留该方向,并把去重状态从全局文本扩展为按 output_index/item_id 与 content_index 分片记录,避免多个 message item 或多个 text content part 共用同一段快照状态。
This commit is contained in:
stabey
2026-05-31 02:54:11 +08:00
parent 8d1e54eba6
commit 3dfafbc379
@@ -45,7 +45,7 @@ pub struct OpenAIResponsesProviderState {
model: Option<String>, model: Option<String>,
started: bool, started: bool,
finished: bool, finished: bool,
text: String, text_parts: BTreeMap<String, String>,
reasoning: String, reasoning: String,
reasoning_parts: BTreeMap<usize, String>, reasoning_parts: BTreeMap<usize, String>,
tool_calls: BTreeMap<usize, OpenAIResponsesProviderToolState>, tool_calls: BTreeMap<usize, OpenAIResponsesProviderToolState>,
@@ -420,24 +420,87 @@ impl OpenAIResponsesProviderState {
index index
} }
fn text_part_key_from_event(value: &Value) -> String {
let item_key = value
.get("output_index")
.and_then(Value::as_u64)
.map(|value| format!("output:{value}"))
.or_else(|| {
value
.get("item_id")
.or_else(|| value.get("id"))
.and_then(Value::as_str)
.map(|value| format!("item:{value}"))
})
.unwrap_or_else(|| "output:default".to_string());
let content_index = value
.get("content_index")
.and_then(Value::as_u64)
.unwrap_or(0);
format!("{item_key}:content:{content_index}")
}
fn text_part_key_from_message_item(
output_index: Option<usize>,
item: &Map<String, Value>,
content_index: usize,
) -> String {
let item_key = output_index
.map(|value| format!("output:{value}"))
.or_else(|| {
item.get("id")
.and_then(Value::as_str)
.map(|value| format!("item:{value}"))
})
.unwrap_or_else(|| "output:default".to_string());
format!("{item_key}:content:{content_index}")
}
fn emit_text_delta(
&mut self,
report_context: &Value,
out: &mut Vec<CanonicalStreamFrame>,
key: String,
text: &str,
) {
if text.is_empty() {
return;
}
self.text_parts.entry(key).or_default().push_str(text);
self.ensure_started(report_context, out);
let (id, model) = self.identity(report_context);
out.push(CanonicalStreamFrame {
id,
model,
event: CanonicalStreamEvent::TextDelta(text.to_string()),
});
}
fn emit_missing_text( fn emit_missing_text(
&mut self, &mut self,
report_context: &Value, report_context: &Value,
out: &mut Vec<CanonicalStreamFrame>, out: &mut Vec<CanonicalStreamFrame>,
key: String,
text: &str, text: &str,
) { ) {
let missing = if text.starts_with(&self.text) { let missing = {
text[self.text.len()..].to_string() let current = self.text_parts.entry(key).or_default();
} else if self.text == text || self.text.starts_with(text) { let missing = if text.starts_with(current.as_str()) {
text[current.len()..].to_string()
} else if current.as_str() == text || current.starts_with(text) {
String::new() String::new()
} else { } else {
text.to_string() text.to_string()
}; };
if !missing.is_empty() {
current.push_str(&missing);
}
missing
};
if missing.is_empty() { if missing.is_empty() {
return; return;
} }
self.ensure_started(report_context, out); self.ensure_started(report_context, out);
self.text.push_str(&missing);
let (id, model) = self.identity(report_context); let (id, model) = self.identity(report_context);
out.push(CanonicalStreamFrame { out.push(CanonicalStreamFrame {
id, id,
@@ -695,28 +758,33 @@ impl OpenAIResponsesProviderState {
report_context: &Value, report_context: &Value,
out: &mut Vec<CanonicalStreamFrame>, out: &mut Vec<CanonicalStreamFrame>,
item: &Map<String, Value>, item: &Map<String, Value>,
output_index: Option<usize>,
) { ) {
if item.get("type").and_then(Value::as_str) != Some("message") { if item.get("type").and_then(Value::as_str) != Some("message") {
return; return;
} }
let mut completed_text = String::new(); for (content_index, raw_content) in item
for raw_content in item
.get("content") .get("content")
.and_then(Value::as_array) .and_then(Value::as_array)
.into_iter() .into_iter()
.flatten() .flatten()
.enumerate()
{ {
let Some(content) = raw_content.as_object() else { let Some(content) = raw_content.as_object() else {
continue; continue;
}; };
if content.get("type").and_then(Value::as_str) == Some("output_text") { if content.get("type").and_then(Value::as_str) == Some("output_text") {
if let Some(text) = content.get("text").and_then(Value::as_str) { if let Some(text) = content.get("text").and_then(Value::as_str) {
completed_text.push_str(text); if !text.is_empty() {
let key = Self::text_part_key_from_message_item(
output_index,
item,
content_index,
);
self.emit_missing_text(report_context, out, key, text);
} }
} }
} }
if !completed_text.is_empty() {
self.emit_missing_text(report_context, out, &completed_text);
} }
} }
@@ -831,18 +899,13 @@ impl OpenAIResponsesProviderState {
} }
"response.output_text.delta" | "response.outtext.delta" => match value.get("delta") { "response.output_text.delta" | "response.outtext.delta" => match value.get("delta") {
Some(Value::String(piece)) if !piece.is_empty() => { Some(Value::String(piece)) if !piece.is_empty() => {
self.ensure_started(report_context, &mut out); let key = Self::text_part_key_from_event(&value);
self.text.push_str(piece); self.emit_text_delta(report_context, &mut out, key, piece);
let (id, model) = self.identity(report_context);
out.push(CanonicalStreamFrame {
id,
model,
event: CanonicalStreamEvent::TextDelta(piece.clone()),
});
} }
Some(Value::Object(delta)) => { Some(Value::Object(delta)) => {
if let Some(text) = delta.get("text").and_then(Value::as_str) { if let Some(text) = delta.get("text").and_then(Value::as_str) {
self.emit_missing_text(report_context, &mut out, text); let key = Self::text_part_key_from_event(&value);
self.emit_missing_text(report_context, &mut out, key, text);
} }
} }
_ => {} _ => {}
@@ -852,7 +915,8 @@ impl OpenAIResponsesProviderState {
if part.get("type").and_then(Value::as_str) == Some("output_text") { if part.get("type").and_then(Value::as_str) == Some("output_text") {
if let Some(text) = part.get("text").and_then(Value::as_str) { if let Some(text) = part.get("text").and_then(Value::as_str) {
if !text.is_empty() { if !text.is_empty() {
self.emit_missing_text(report_context, &mut out, text); let key = Self::text_part_key_from_event(&value);
self.emit_missing_text(report_context, &mut out, key, text);
} }
} }
} }
@@ -890,7 +954,8 @@ impl OpenAIResponsesProviderState {
}) })
.unwrap_or_default(); .unwrap_or_default();
if !text.is_empty() { if !text.is_empty() {
self.emit_missing_text(report_context, &mut out, text); let key = Self::text_part_key_from_event(&value);
self.emit_missing_text(report_context, &mut out, key, text);
} }
} }
"response.reasoning_summary_text.delta" => { "response.reasoning_summary_text.delta" => {
@@ -967,7 +1032,7 @@ impl OpenAIResponsesProviderState {
self.emit_tool_result_item(report_context, &mut out, item, output_index); self.emit_tool_result_item(report_context, &mut out, item, output_index);
} }
"message" => { "message" => {
self.emit_message_item(report_context, &mut out, item); self.emit_message_item(report_context, &mut out, item, output_index);
} }
"reasoning" => { "reasoning" => {
self.ensure_started(report_context, &mut out); self.ensure_started(report_context, &mut out);
@@ -1139,7 +1204,7 @@ impl OpenAIResponsesProviderState {
self.emit_tool_result_item(report_context, &mut out, item, output_index); self.emit_tool_result_item(report_context, &mut out, item, output_index);
} }
"message" => { "message" => {
self.emit_message_item(report_context, &mut out, item); self.emit_message_item(report_context, &mut out, item, output_index);
} }
"reasoning" => { "reasoning" => {
self.emit_reasoning_item(report_context, &mut out, item); self.emit_reasoning_item(report_context, &mut out, item);
@@ -1194,7 +1259,12 @@ impl OpenAIResponsesProviderState {
}; };
match item.get("type").and_then(Value::as_str).unwrap_or_default() { match item.get("type").and_then(Value::as_str).unwrap_or_default() {
"message" => { "message" => {
self.emit_message_item(report_context, &mut out, item); self.emit_message_item(
report_context,
&mut out,
item,
Some(output_index),
);
} }
"function_call" => { "function_call" => {
self.emit_tool_call_item( self.emit_tool_call_item(
@@ -3235,6 +3305,102 @@ mod tests {
assert_eq!(text, "Hello world"); assert_eq!(text, "Hello world");
} }
#[test]
fn openai_responses_provider_state_dedupes_text_snapshots_per_output_item() {
let mut state = OpenAIResponsesProviderState::default();
let report_context = json!({});
let mut frames = Vec::new();
for event in [
json!({
"type": "response.output_text.delta",
"response_id": "resp_multi_message",
"output_index": 0,
"item_id": "msg_1",
"content_index": 0,
"delta": "First message.",
}),
json!({
"type": "response.output_text.done",
"response_id": "resp_multi_message",
"output_index": 0,
"item_id": "msg_1",
"content_index": 0,
"text": "First message.",
}),
json!({
"type": "response.output_item.done",
"response_id": "resp_multi_message",
"output_index": 0,
"item": {
"type": "message",
"id": "msg_1",
"status": "completed",
"content": [{
"type": "output_text",
"text": "First message.",
}],
},
}),
json!({
"type": "response.output_text.delta",
"response_id": "resp_multi_message",
"output_index": 1,
"item_id": "msg_2",
"content_index": 0,
"delta": "Second message.",
}),
json!({
"type": "response.output_text.done",
"response_id": "resp_multi_message",
"output_index": 1,
"item_id": "msg_2",
"content_index": 0,
"text": "Second message.",
}),
json!({
"type": "response.content_part.done",
"response_id": "resp_multi_message",
"output_index": 1,
"item_id": "msg_2",
"content_index": 0,
"part": {
"type": "output_text",
"text": "Second message.",
},
}),
json!({
"type": "response.output_item.done",
"response_id": "resp_multi_message",
"output_index": 1,
"item": {
"type": "message",
"id": "msg_2",
"status": "completed",
"content": [{
"type": "output_text",
"text": "Second message.",
}],
},
}),
] {
frames.extend(
state
.push_line(&report_context, data_line(event))
.expect("responses text event should parse"),
);
}
let text = frames
.iter()
.filter_map(|frame| match &frame.event {
CanonicalStreamEvent::TextDelta(text) => Some(text.as_str()),
_ => None,
})
.collect::<String>();
assert_eq!(text, "First message.Second message.");
}
#[test] #[test]
fn openai_responses_provider_state_delays_arguments_until_tool_name_is_known() { fn openai_responses_provider_state_delays_arguments_until_tool_name_is_known() {
let mut state = OpenAIResponsesProviderState::default(); let mut state = OpenAIResponsesProviderState::default();