1use crate::error::ProviderError;
9use crate::operation::{
10 CallFragment, Completion, Finish, IfMalformed, ReasoningPart, Seal, TextPart,
11};
12use crate::providers::internal::wire;
13use crate::providers::openai::responses_api::{
14 IncompleteDetailsReason, ReasoningSummary, ResponseStatus, ResponsesUsage,
15};
16use crate::wire::{Decoder, Flow, Out, WireEvent, WireFrame};
17use serde::{Deserialize, Serialize};
18
19use super::{CompletionResponse, Output};
20
21#[derive(Debug, Serialize, Deserialize, Clone)]
23#[serde(untagged)]
24pub enum StreamingCompletionChunk {
25 Response(ResponseChunk),
26 Delta(ItemChunk),
27}
28
29#[derive(Debug, Serialize, Deserialize, Clone)]
32pub struct StreamingCompletionResponse {
33 #[serde(default, skip_serializing_if = "Option::is_none")]
36 pub usage: Option<ResponsesUsage>,
37 #[serde(default, skip_serializing_if = "Option::is_none")]
39 pub reasoning_metadata: Option<serde_json::Map<String, serde_json::Value>>,
40 #[serde(default, skip_serializing_if = "Option::is_none")]
42 pub reasoning_context: Option<String>,
43 #[serde(default, skip_serializing_if = "Option::is_none")]
45 pub status: Option<ResponseStatus>,
46 #[serde(default, skip_serializing_if = "Option::is_none")]
48 pub incomplete_details: Option<IncompleteDetailsReason>,
49 #[serde(default, skip_serializing_if = "Option::is_none")]
55 pub message_id: Option<String>,
56 #[serde(default, skip_serializing_if = "Option::is_none")]
59 pub response_id: Option<String>,
60 #[serde(default, skip_serializing_if = "Option::is_none")]
62 pub model: Option<String>,
63}
64
65impl StreamingCompletionResponse {
66 pub fn new(usage: Option<ResponsesUsage>) -> Self {
69 Self {
70 usage,
71 reasoning_metadata: None,
72 reasoning_context: None,
73 status: None,
74 incomplete_details: None,
75 message_id: None,
76 response_id: None,
77 model: None,
78 }
79 }
80}
81
82fn finish_of(
92 provider: &str,
93 upstream_reasoning_issuer: bool,
94 response: StreamingCompletionResponse,
95) -> (Finish, Option<String>) {
96 let issuer = upstream_reasoning_issuer
97 .then_some(response.model.as_deref())
98 .flatten()
99 .map(|model| crate::providers::openai::wire::upstream_reasoning_issuer(provider, model));
100 let finish_reason = response
101 .status
102 .as_ref()
103 .and_then(|status| super::map_finish_reason(status, response.incomplete_details.as_ref()));
104 let finish = Finish {
105 usage: crate::completion::Usage::from(&response),
106 reason: finish_reason,
107 message_id: response.message_id,
108 response_id: response.response_id,
109 model: response.model,
110 };
111 (finish, issuer)
112}
113
114pub(crate) fn reasoning_from_done_item(
118 provider_id: Option<&str>,
119 summary: Vec<ReasoningSummary>,
120 content: Vec<String>,
121 encrypted_content: Option<String>,
122 signature: Option<String>,
123) -> Option<crate::message::Reasoning> {
124 let blocks = super::reasoning_content_blocks(summary, content, encrypted_content, signature);
127
128 if blocks.is_empty() {
129 return None;
130 }
131
132 Some(crate::message::Reasoning {
133 id: provider_id.map(str::to_owned),
134 content: blocks,
135 })
136}
137
138impl From<&StreamingCompletionResponse> for crate::completion::Usage {
139 fn from(response: &StreamingCompletionResponse) -> Self {
140 response.usage.as_ref().map(Self::from).unwrap_or_default()
141 }
142}
143
144#[derive(Debug, Serialize, Deserialize, Clone)]
146pub struct ResponseChunk {
147 #[serde(rename = "type")]
149 pub kind: ResponseChunkKind,
150 pub response: CompletionResponse,
152 pub sequence_number: u64,
154}
155
156#[derive(Debug, Serialize, Deserialize, Clone, Copy)]
159pub enum ResponseChunkKind {
160 #[serde(rename = "response.created")]
161 ResponseCreated,
162 #[serde(rename = "response.in_progress")]
163 ResponseInProgress,
164 #[serde(rename = "response.completed")]
165 ResponseCompleted,
166 #[serde(rename = "response.failed")]
167 ResponseFailed,
168 #[serde(rename = "response.incomplete")]
169 ResponseIncomplete,
170}
171
172fn is_known_responses_event_type(kind: &str) -> bool {
179 matches!(
180 kind,
181 "response.created"
182 | "response.in_progress"
183 | "response.completed"
184 | "response.failed"
185 | "response.incomplete"
186 | "response.output_item.added"
187 | "response.output_item.done"
188 | "response.content_part.added"
189 | "response.content_part.done"
190 | "response.output_text.delta"
191 | "response.output_text.done"
192 | "response.refusal.delta"
193 | "response.refusal.done"
194 | "response.function_call_arguments.delta"
195 | "response.function_call_arguments.done"
196 | "response.reasoning_summary_part.added"
197 | "response.reasoning_summary_part.done"
198 | "response.reasoning_summary_text.delta"
199 | "response.reasoning_summary_text.done"
200 | "response.reasoning_text.delta"
201 | "response.reasoning_text.done"
202 )
203}
204
205#[doc(hidden)]
208pub fn classify_responses_frame(data: &str) -> WireEvent<StreamingCompletionChunk> {
209 wire::classify_tagged_frame(data, "type", is_known_responses_event_type)
210}
211
212fn message_id_from_response(response: &CompletionResponse) -> Option<String> {
215 response.output.iter().find_map(|item| match item {
216 Output::Message(message) => Some(message.id.clone()),
217 _ => None,
218 })
219}
220
221fn repair_envelope_less_frame(data: &str) -> Option<String> {
226 let mut value = serde_json::from_str::<serde_json::Value>(data).ok()?;
227 let object = value.as_object_mut()?;
228 for field in [
229 "sequence_number",
230 "output_index",
231 "content_index",
232 "summary_index",
233 ] {
234 object
235 .entry(field)
236 .or_insert_with(|| serde_json::Value::from(0));
237 }
238 serde_json::to_string(&value).ok()
239}
240
241pub enum ResponsesEvent {
243 Frame {
245 raw: String,
247 chunk: StreamingCompletionChunk,
249 },
250 Whole(Box<CompletionResponse>),
253 Failure(String),
258 Sentinel,
260}
261
262const WHOLE_BODY_MARKERS: &[&str] = &["object", "output", "status", "error"];
264
265fn is_error_event(data: &str) -> bool {
267 serde_json::from_str::<serde_json::Value>(data)
268 .is_ok_and(|value| value.get("type").and_then(serde_json::Value::as_str) == Some("error"))
269}
270
271#[derive(Deserialize)]
274struct ErrorEnvelope {
275 #[allow(dead_code)]
277 error: serde_json::Value,
278}
279
280pub struct ResponsesDecoder<'id> {
283 provider: String,
287 document: Option<serde_json::Value>,
290 repair_envelopes: bool,
293 upstream_reasoning_issuer: bool,
296 terminal: StreamingCompletionResponse,
299 texts: std::collections::HashMap<String, TextPart<'id>>,
303 current_text_item: Option<String>,
306 anonymous_text: Option<TextPart<'id>>,
308 delta_text_items: std::collections::HashSet<String>,
313 delta_text_slots: std::collections::HashSet<u64>,
316 unattributed_text_delta: bool,
317 extras_items: std::collections::HashSet<String>,
321 extras_slots: std::collections::HashSet<u64>,
322 reasoning: std::collections::HashMap<u64, ReasoningPart<'id>>,
324 text_slots: std::collections::HashMap<String, u64>,
328 anonymous_text_slot: Option<u64>,
329 message_items: Vec<MessageItem>,
332 message_slots: std::collections::HashMap<u64, usize>,
333}
334
335#[derive(Default)]
337struct MessageItem {
338 id: String,
339 phase: Option<String>,
340 done: bool,
344}
345
346#[derive(Clone, Copy, PartialEq, Eq)]
348enum Statement {
349 Added,
350 Done,
351 Terminal,
352}
353
354impl<'id> ResponsesDecoder<'id> {
355 pub fn new(provider: &str) -> Self {
357 Self {
358 provider: provider.to_owned(),
359 document: None,
360 repair_envelopes: false,
361 upstream_reasoning_issuer: false,
362 terminal: StreamingCompletionResponse::new(None),
363 texts: std::collections::HashMap::new(),
364 current_text_item: None,
365 anonymous_text: None,
366 delta_text_items: std::collections::HashSet::new(),
367 delta_text_slots: std::collections::HashSet::new(),
368 unattributed_text_delta: false,
369 extras_items: std::collections::HashSet::new(),
370 extras_slots: std::collections::HashSet::new(),
371 reasoning: std::collections::HashMap::new(),
372 text_slots: std::collections::HashMap::new(),
373 anonymous_text_slot: None,
374 message_items: Vec::new(),
375 message_slots: std::collections::HashMap::new(),
376 }
377 }
378
379 pub fn with_envelope_repair(mut self) -> Self {
381 self.repair_envelopes = true;
382 self
383 }
384
385 pub fn with_upstream_reasoning_issuer(mut self) -> Self {
388 self.upstream_reasoning_issuer = true;
389 self
390 }
391
392 pub fn with_initial_usage(mut self, usage: Option<ResponsesUsage>) -> Self {
395 self.terminal.usage = usage;
396 self
397 }
398
399 fn classify_payload(&self, data: &str) -> WireEvent<ResponsesEvent> {
402 if data.trim() == "[DONE]" {
403 return WireEvent::Known(ResponsesEvent::Sentinel);
404 }
405 if is_error_event(data) {
406 return WireEvent::Known(ResponsesEvent::Failure(data.to_owned()));
407 }
408 let body = |data: &str| {
409 wire::classify_marker_keyed_frame::<CompletionResponse>(data, WHOLE_BODY_MARKERS)
410 .map(|response| ResponsesEvent::Whole(Box::new(response)))
411 };
412 let envelope = |data: &str| {
413 wire::classify_marker_keyed_frame::<ErrorEnvelope>(data, &["error"])
415 .map(|_| ResponsesEvent::Failure(data.to_owned()))
416 };
417 wire::classify_or(
418 data,
419 |data| {
420 classify_responses_frame(data).map(|chunk| ResponsesEvent::Frame {
421 raw: data.to_owned(),
422 chunk,
423 })
424 },
425 |data| wire::classify_or(data, body, envelope),
426 )
427 }
428
429 fn text_part(
432 &mut self,
433 output_index: u64,
434 item_id: Option<&str>,
435 out: &mut Out<'id, Completion>,
436 ) -> &TextPart<'id> {
437 let item_id = item_id
438 .filter(|id| !id.is_empty())
439 .map(str::to_owned)
440 .or_else(|| self.current_text_item.clone());
441 match item_id {
442 Some(item_id) => {
443 self.current_text_item = Some(item_id.clone());
444 self.text_slots
445 .entry(item_id.clone())
446 .or_insert(output_index);
447 self.texts.entry(item_id).or_insert_with(|| out.text())
448 }
449 None => {
450 self.anonymous_text_slot.get_or_insert(output_index);
451 self.anonymous_text.get_or_insert_with(|| out.text())
452 }
453 }
454 }
455
456 fn note_message_item(
464 &mut self,
465 output_index: u64,
466 message: &super::OutputMessage,
467 statement: Statement,
468 ) {
469 let known = self
470 .message_items
471 .iter()
472 .position(|item| !message.id.is_empty() && item.id == message.id);
473 let slot = self.message_slots.get(&output_index).copied();
474 let slot = match statement {
475 Statement::Added => None,
476 Statement::Done => {
477 slot.filter(|at| self.message_items.get(*at).is_some_and(|item| !item.done))
478 }
479 Statement::Terminal => slot,
480 };
481 let at = match known.or(slot) {
482 Some(at) => at,
483 None => {
484 self.message_items.push(MessageItem::default());
485 self.message_items.len() - 1
486 }
487 };
488 if statement != Statement::Terminal || !self.message_slots.contains_key(&output_index) {
490 self.message_slots.insert(output_index, at);
491 }
492 let Some(item) = self.message_items.get_mut(at) else {
493 return;
494 };
495 let wins = statement == Statement::Done || !item.done;
496 if !message.id.is_empty() && (wins || item.id.is_empty()) {
497 item.id.clone_from(&message.id);
498 }
499 if message.phase.is_some() && (wins || item.phase.is_none()) {
500 item.phase.clone_from(&message.phase);
501 }
502 item.done |= statement == Statement::Done;
503 }
504
505 fn attach_message_items(&mut self, out: &mut Out<'id, Completion>) {
512 let item_of = |key: Option<&str>, slot: Option<u64>| {
513 key.and_then(|key| {
514 self.message_items
515 .iter()
516 .position(|item| !item.id.is_empty() && item.id == key)
517 })
518 .or_else(|| self.message_slots.get(&slot?).copied())
519 .and_then(|at| self.message_items.get(at))
520 };
521 let parts: Vec<(&TextPart<'id>, &MessageItem)> = self
522 .texts
523 .iter()
524 .filter_map(|(key, part)| {
525 Some((part, item_of(Some(key), self.text_slots.get(key).copied())?))
526 })
527 .chain(
528 self.anonymous_text
529 .as_ref()
530 .and_then(|part| Some((part, item_of(None, self.anonymous_text_slot)?))),
531 )
532 .collect();
533 let several = self.message_items.len() > 1;
534 for (part, item) in parts {
535 let mut extras = serde_json::Map::new();
536 if let Some(phase) = &item.phase {
537 extras.insert(
538 super::OPENAI_RESPONSES_PHASE_KEY.to_owned(),
539 serde_json::Value::String(phase.clone()),
540 );
541 }
542 if several && !item.id.is_empty() {
543 extras.insert(
544 super::OPENAI_RESPONSES_MESSAGE_ID_KEY.to_owned(),
545 serde_json::Value::String(item.id.clone()),
546 );
547 }
548 if extras.is_empty() {
549 continue;
550 }
551 if let Some(params) = crate::message::AdditionalParams::from_entries(Some((
552 super::OPENAI_RESPONSES_EXTRAS_KEY,
553 serde_json::Value::Object(extras),
554 ))) {
555 out.text_params(part, params);
556 }
557 }
558 }
559
560 fn note_text_delta(&mut self, output_index: u64, item_id: Option<&str>) {
567 self.delta_text_slots.insert(output_index);
568 match item_id
569 .filter(|id| !id.is_empty())
570 .map(str::to_owned)
571 .or_else(|| self.current_text_item.clone())
572 {
573 Some(id) => {
574 self.delta_text_items.insert(id);
575 }
576 None => self.unattributed_text_delta = true,
577 }
578 }
579
580 fn delta_delivered_text(&self, output_index: u64, item_id: &str) -> bool {
588 self.unattributed_text_delta
589 || self.delta_text_slots.contains(&output_index)
590 || self.delta_text_items.contains(item_id)
591 }
592
593 fn publish_message_text(
597 &mut self,
598 output_index: u64,
599 message: &super::OutputMessage,
600 out: &mut Out<'id, Completion>,
601 ) {
602 if !message.content.is_empty() {
603 self.note_text_delta(output_index, Some(&message.id));
604 }
605 self.note_extras(output_index, &message.id);
606 for content in message.content.iter().cloned() {
607 let text = super::text_block(content);
608 let part = self.text_part(output_index, Some(&message.id), out);
609 out.push_text(part, &text.text);
610 if let Some(additional_params) = text.additional_params {
611 out.text_params(part, additional_params);
612 }
613 }
614 }
615
616 fn note_extras(&mut self, output_index: u64, item_id: &str) {
619 self.extras_slots.insert(output_index);
620 if !item_id.is_empty() {
621 self.extras_items.insert(item_id.to_owned());
622 }
623 }
624
625 fn attach_message_extras(
631 &mut self,
632 output_index: u64,
633 message: &super::OutputMessage,
634 out: &mut Out<'id, Completion>,
635 ) {
636 if self.extras_slots.contains(&output_index) || self.extras_items.contains(&message.id) {
637 return;
638 }
639 let Some(part) = self.texts.get(&message.id) else {
640 return;
641 };
642 let extras = message
643 .content
644 .iter()
645 .cloned()
646 .filter_map(|content| super::text_block(content).additional_params)
647 .reduce(|mut extras, next| {
648 extras.merge(next);
649 extras
650 });
651 if let Some(extras) = extras {
652 out.text_params(part, extras);
653 }
654 self.note_extras(output_index, &message.id);
655 }
656
657 fn merge_terminal_body_text(
661 &mut self,
662 response: &CompletionResponse,
663 out: &mut Out<'id, Completion>,
664 ) {
665 for (output_index, item) in response.output.iter().enumerate() {
669 let output_index = output_index as u64;
670 let Output::Message(message) = item else {
671 continue;
672 };
673 self.note_message_item(output_index, message, Statement::Terminal);
674 if message.content.is_empty() {
675 continue;
676 }
677 if self.delta_delivered_text(output_index, &message.id) {
678 self.attach_message_extras(output_index, message, out);
679 } else {
680 self.publish_message_text(output_index, message, out);
681 }
682 }
683 }
684
685 fn decode_item_chunk(
687 &mut self,
688 chunk: ItemChunk,
689 out: &mut Out<'id, Completion>,
690 ) -> Result<(), ProviderError> {
691 let ItemChunk {
692 item_id: outer_item_id,
693 output_index,
694 data: item,
695 } = chunk;
696
697 match item {
698 ItemChunkKind::OutputItemAdded(StreamingItemDoneOutput {
699 item: Output::FunctionCall(func),
700 ..
701 }) => {
702 self.current_text_item = None;
705 out.call_fragment(
708 output_index as usize,
709 CallFragment {
710 id: Some(func.call_id.as_str()),
711 item_id: (!func.call_id.is_empty()).then_some(func.id.as_str()),
712 name: Some(func.name.as_str()),
713 ..CallFragment::default()
714 },
715 )?;
716 }
717 ItemChunkKind::OutputItemAdded(StreamingItemDoneOutput {
718 item: Output::Message(message),
719 ..
720 }) => self.note_message_item(output_index, &message, Statement::Added),
721 ItemChunkKind::OutputItemDone(message) => {
722 self.current_text_item = None;
725 self.push_output_item_done(message.item, output_index, out)?;
726 }
727 ItemChunkKind::OutputTextDelta(DeltaTextChunk { delta, .. })
730 | ItemChunkKind::RefusalDelta(DeltaTextChunk { delta, .. }) => {
731 self.note_text_delta(output_index, outer_item_id.as_deref());
732 let part = self.text_part(output_index, outer_item_id.as_deref(), out);
733 out.push_text(part, &delta);
734 }
735 ItemChunkKind::ReasoningSummaryTextDelta(SummaryTextChunk { delta, .. })
739 | ItemChunkKind::ReasoningTextDelta(DeltaTextChunkWithItemId { delta, .. }) => {
740 self.current_text_item = None;
741 let part = self
742 .reasoning
743 .entry(output_index)
744 .or_insert_with(|| out.reasoning());
745 out.push_reasoning(part, &delta);
746 }
747 ItemChunkKind::FunctionCallArgsDelta(delta) => {
748 self.current_text_item = None;
749 out.call_fragment(
750 output_index as usize,
751 CallFragment {
752 arguments: Some(delta.delta.as_str()),
753 ..CallFragment::default()
754 },
755 )?;
756 }
757 _ => {}
758 }
759 Ok(())
760 }
761
762 fn push_output_item_done(
763 &mut self,
764 item: Output,
765 output_index: u64,
766 out: &mut Out<'id, Completion>,
767 ) -> Result<(), ProviderError> {
768 match item {
769 Output::FunctionCall(func) => {
770 let index = output_index as usize;
771 let streamed = out.pending_has_arguments(index);
772 out.call_fragment(
776 index,
777 CallFragment {
778 id: Some(func.call_id.as_str()),
779 item_id: (!func.call_id.is_empty()).then_some(func.id.as_str()),
780 name: Some(func.name.as_str()),
781 ..CallFragment::default()
782 },
783 )?;
784 match func.arguments.parse() {
785 Ok(arguments) => out.announce_pending(index, arguments),
787 Err(_) if !streamed => out.call_fragment(
790 index,
791 CallFragment {
792 arguments: Some(func.arguments.as_str()),
793 ..CallFragment::default()
794 },
795 )?,
796 Err(_) => {}
797 }
798 out.close_pending(index, IfMalformed::Drop)?;
800 }
801 Output::Reasoning {
802 id,
803 summary,
804 content,
805 encrypted_content,
806 signature,
807 ..
808 } => {
809 let provider_id = (!id.is_empty()).then_some(id);
810 let part = self.reasoning.remove(&output_index);
811 let restated = reasoning_from_done_item(
812 provider_id.as_deref(),
813 summary,
814 content,
815 encrypted_content,
816 signature,
817 );
818 match (part, restated) {
819 (Some(part), restated) => out.close_reasoning(
821 part,
822 Seal {
823 id: provider_id,
824 restated,
825 ..Seal::default()
826 },
827 ),
828 (None, Some(restated)) => out.reasoning_block(restated),
829 (None, None) => {
831 if let Some(id) = provider_id {
832 out.reasoning_block(crate::message::Reasoning {
833 id: Some(id),
834 content: Vec::new(),
835 });
836 }
837 }
838 }
839 }
840 Output::Message(message) => {
841 self.note_message_item(output_index, &message, Statement::Done);
842 self.attach_message_extras(output_index, &message, out);
843 if !message.id.is_empty() {
844 out.message_id(message.id);
845 }
846 }
847 Output::Unknown(value) => {
851 out.unknown(value.into());
852 }
853 Output::Compaction(fields) => {
856 let mut map = fields;
857 map.insert(
858 "type".to_string(),
859 serde_json::Value::String("compaction".to_string()),
860 );
861 out.unknown(serde_json::Value::Object(map).into());
862 }
863 }
864 Ok(())
865 }
866
867 fn record_terminal(&mut self, response: CompletionResponse, out: &mut Out<'id, Completion>) {
872 self.document = serde_json::to_value(&response).ok();
873 self.merge_terminal_body_text(&response, out);
878 if let Some(message_id) = message_id_from_response(&response) {
879 self.terminal.message_id = Some(message_id);
880 }
881 if !response.id.is_empty() {
882 self.terminal.response_id = Some(response.id.clone());
883 }
884 if !response.model.is_empty() {
885 self.terminal.model = Some(response.model.clone());
886 }
887 self.terminal.status = Some(response.status);
888 if response.incomplete_details.is_some() {
889 self.terminal.incomplete_details = response.incomplete_details;
890 }
891 if response.usage.is_some() {
892 self.terminal.usage = response.usage;
893 }
894 if response.reasoning_metadata.is_some() {
895 self.terminal.reasoning_metadata = response.reasoning_metadata;
896 }
897 if response.reasoning_context.is_some() {
898 self.terminal.reasoning_context = response.reasoning_context;
899 }
900 }
901
902 fn end(&mut self, mut out: Out<'id, Completion>) -> Result<Flow, ProviderError> {
905 for index in out.pending_calls() {
906 out.close_pending(index, IfMalformed::Drop)?;
907 }
908 self.attach_message_items(&mut out);
909 if let Some(document) = self.document.take() {
910 out.raw(document);
911 }
912 let terminal =
913 std::mem::replace(&mut self.terminal, StreamingCompletionResponse::new(None));
914 let (finish, issuer) = finish_of(&self.provider, self.upstream_reasoning_issuer, terminal);
915 if let Some(issuer) = issuer {
916 out.issued_by(issuer);
917 }
918 Ok(out.end(finish))
919 }
920
921 fn replay_whole_response(
925 &mut self,
926 response: CompletionResponse,
927 mut out: Out<'id, Completion>,
928 ) -> Result<Flow, ProviderError> {
929 let structured_reasoning = response
934 .output
935 .iter()
936 .any(|item| matches!(item, Output::Reasoning { .. }));
937 if !structured_reasoning
938 && let Some(reasoning) = response
939 .provider_reasoning
940 .as_deref()
941 .filter(|reasoning| !reasoning.is_empty())
942 {
943 out.reasoning_block(crate::message::Reasoning::new(reasoning));
944 }
945
946 for (output_index, item) in response.output.iter().cloned().enumerate() {
947 let output_index = output_index as u64;
948 if let Output::Message(message) = &item {
949 self.publish_message_text(output_index, message, &mut out);
950 }
951 self.push_output_item_done(item, output_index, &mut out)?;
953 }
954 self.record_terminal(response, &mut out);
955 self.end(out)
956 }
957}
958
959impl<'id> Decoder<'id, Completion> for ResponsesDecoder<'id> {
960 type Event = ResponsesEvent;
961
962 fn classify(&self, frame: WireFrame) -> WireEvent<ResponsesEvent> {
963 let data = frame.as_str().into_owned();
964 if !self.repair_envelopes {
965 return self.classify_payload(&data);
966 }
967 wire::classify_with_repair(
971 &data,
972 |data| self.classify_payload(data),
973 repair_envelope_less_frame,
974 |corrupt| {
975 <serde_json::Error as serde::de::Error>::custom(format!(
976 "invalid JSON frame in buffered Responses SSE body: {corrupt}"
977 ))
978 },
979 || {
980 let kind = serde_json::from_str::<serde_json::Value>(&data)
981 .ok()
982 .and_then(|value| {
983 value
984 .get("type")
985 .and_then(serde_json::Value::as_str)
986 .map(ToOwned::to_owned)
987 })
988 .unwrap_or_default();
989 <serde_json::Error as serde::de::Error>::custom(format!(
990 "malformed `{kind}` event in buffered Responses SSE body"
991 ))
992 },
993 )
994 }
995
996 fn decode(
997 &mut self,
998 event: ResponsesEvent,
999 mut out: Out<'id, Completion>,
1000 ) -> Result<Flow, ProviderError> {
1001 match event {
1002 ResponsesEvent::Frame {
1003 chunk: StreamingCompletionChunk::Delta(chunk),
1004 ..
1005 } => {
1006 self.decode_item_chunk(chunk, &mut out)?;
1007 Ok(Flow::More)
1008 }
1009 ResponsesEvent::Frame {
1010 raw,
1011 chunk: StreamingCompletionChunk::Response(chunk),
1012 } => {
1013 let ResponseChunk { kind, response, .. } = chunk;
1014 match kind {
1015 ResponseChunkKind::ResponseCompleted
1020 | ResponseChunkKind::ResponseIncomplete => {
1021 self.record_terminal(response, &mut out);
1022 self.end(out)
1023 }
1024 ResponseChunkKind::ResponseFailed => {
1025 Err(crate::error::ProviderError::from_provider_body(&raw))
1026 }
1027 ResponseChunkKind::ResponseCreated | ResponseChunkKind::ResponseInProgress => {
1028 Ok(Flow::More)
1029 }
1030 }
1031 }
1032 ResponsesEvent::Whole(response) => self.replay_whole_response(*response, out),
1034 ResponsesEvent::Failure(raw) => {
1035 Err(crate::error::ProviderError::from_provider_body(&raw))
1036 }
1037 ResponsesEvent::Sentinel => Ok(Flow::More),
1039 }
1040 }
1041}
1042
1043#[derive(Debug, Serialize, Deserialize, Clone)]
1045pub struct ItemChunk {
1046 pub item_id: Option<String>,
1048 pub output_index: u64,
1050 #[serde(flatten)]
1052 pub data: ItemChunkKind,
1053}
1054
1055#[derive(Debug, Serialize, Deserialize, Clone)]
1057#[serde(tag = "type")]
1058pub enum ItemChunkKind {
1059 #[serde(rename = "response.output_item.added")]
1060 OutputItemAdded(StreamingItemDoneOutput),
1061 #[serde(rename = "response.output_item.done")]
1062 OutputItemDone(StreamingItemDoneOutput),
1063 #[serde(rename = "response.content_part.added")]
1064 ContentPartAdded(ContentPartChunk),
1065 #[serde(rename = "response.content_part.done")]
1066 ContentPartDone(ContentPartChunk),
1067 #[serde(rename = "response.output_text.delta")]
1068 OutputTextDelta(DeltaTextChunk),
1069 #[serde(rename = "response.output_text.done")]
1070 OutputTextDone(OutputTextChunk),
1071 #[serde(rename = "response.refusal.delta")]
1072 RefusalDelta(DeltaTextChunk),
1073 #[serde(rename = "response.refusal.done")]
1074 RefusalDone(RefusalTextChunk),
1075 #[serde(rename = "response.function_call_arguments.delta")]
1076 FunctionCallArgsDelta(DeltaTextChunkWithItemId),
1077 #[serde(rename = "response.function_call_arguments.done")]
1078 FunctionCallArgsDone(ArgsTextChunk),
1079 #[serde(rename = "response.reasoning_summary_part.added")]
1080 ReasoningSummaryPartAdded(SummaryPartChunk),
1081 #[serde(rename = "response.reasoning_summary_part.done")]
1082 ReasoningSummaryPartDone(SummaryPartChunk),
1083 #[serde(rename = "response.reasoning_summary_text.delta")]
1084 ReasoningSummaryTextDelta(SummaryTextChunk),
1085 #[serde(rename = "response.reasoning_summary_text.done")]
1086 ReasoningSummaryTextDone(SummaryTextChunk),
1087 #[serde(rename = "response.reasoning_text.delta")]
1088 ReasoningTextDelta(DeltaTextChunkWithItemId),
1089 #[serde(rename = "response.reasoning_text.done")]
1092 ReasoningTextDone(OutputTextChunk),
1093 }
1099
1100#[derive(Debug, Serialize, Deserialize, Clone)]
1101pub struct StreamingItemDoneOutput {
1102 pub sequence_number: u64,
1103 pub item: Output,
1104}
1105
1106#[derive(Debug, Serialize, Deserialize, Clone)]
1107pub struct ContentPartChunk {
1108 pub content_index: u64,
1109 pub sequence_number: u64,
1110 pub part: ContentPartChunkPart,
1111}
1112
1113#[derive(Debug, Serialize, Clone)]
1114#[serde(tag = "type", rename_all = "snake_case")]
1115pub enum ContentPartChunkPart {
1116 OutputText {
1117 text: String,
1118 },
1119 SummaryText {
1120 text: String,
1121 },
1122 #[serde(untagged)]
1125 Unknown(serde_json::Value),
1126}
1127
1128impl<'de> Deserialize<'de> for ContentPartChunkPart {
1132 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1133 where
1134 D: serde::Deserializer<'de>,
1135 {
1136 let value = serde_json::Value::deserialize(deserializer)?;
1137 let text_field = |part: &str| -> Result<String, D::Error> {
1138 value
1139 .get("text")
1140 .and_then(serde_json::Value::as_str)
1141 .map(ToOwned::to_owned)
1142 .ok_or_else(|| {
1143 serde::de::Error::custom(format!(
1144 "`{part}` content part is missing a string `text` field"
1145 ))
1146 })
1147 };
1148 match value.get("type").cloned() {
1149 Some(serde_json::Value::String(tag)) => match tag.as_str() {
1150 "output_text" => Ok(Self::OutputText {
1151 text: text_field("output_text")?,
1152 }),
1153 "summary_text" => Ok(Self::SummaryText {
1154 text: text_field("summary_text")?,
1155 }),
1156 _ => Ok(Self::Unknown(value)),
1157 },
1158 Some(_) => Err(serde::de::Error::custom(
1159 "content part `type` must be a string",
1160 )),
1161 None => Ok(Self::Unknown(value)),
1162 }
1163 }
1164}
1165
1166#[derive(Debug, Serialize, Deserialize, Clone)]
1167pub struct DeltaTextChunk {
1168 pub content_index: u64,
1169 pub sequence_number: u64,
1170 pub delta: String,
1171}
1172
1173#[derive(Debug, Serialize, Deserialize, Clone)]
1174pub struct DeltaTextChunkWithItemId {
1175 #[serde(default, skip_serializing_if = "Option::is_none")]
1176 pub content_index: Option<u64>,
1177 pub sequence_number: u64,
1178 pub delta: String,
1179}
1180
1181#[derive(Debug, Serialize, Deserialize, Clone)]
1182pub struct OutputTextChunk {
1183 pub content_index: u64,
1184 pub sequence_number: u64,
1185 pub text: String,
1186}
1187
1188#[derive(Debug, Serialize, Deserialize, Clone)]
1189pub struct RefusalTextChunk {
1190 pub content_index: u64,
1191 pub sequence_number: u64,
1192 pub refusal: String,
1193}
1194
1195#[derive(Debug, Serialize, Deserialize, Clone)]
1196pub struct ArgsTextChunk {
1197 #[serde(default, skip_serializing_if = "Option::is_none")]
1198 pub content_index: Option<u64>,
1199 pub sequence_number: u64,
1200 pub arguments: serde_json::Value,
1201}
1202
1203#[derive(Debug, Serialize, Deserialize, Clone)]
1204pub struct SummaryPartChunk {
1205 pub summary_index: u64,
1206 pub sequence_number: u64,
1207 pub part: SummaryPartChunkPart,
1208}
1209
1210#[derive(Debug, Serialize, Deserialize, Clone)]
1211pub struct SummaryTextChunk {
1212 pub summary_index: u64,
1213 pub sequence_number: u64,
1214 #[serde(alias = "text")]
1217 pub delta: String,
1218}
1219
1220#[derive(Debug, Serialize, Deserialize, Clone)]
1221#[serde(tag = "type", rename_all = "snake_case")]
1222pub enum SummaryPartChunkPart {
1223 SummaryText { text: String },
1224}
1225
1226#[cfg(test)]
1227mod tests;