Skip to main content

rig_core/providers/openai/responses_api/
streaming.rs

1//! Responses frame classification, event decoding, and terminal metadata.
2//!
3//! ```
4//! use rig_core::providers::openai::responses_api::streaming::ResponsesDecoder;
5//! let decoder = ResponsesDecoder::new("openai");
6//! ```
7
8use 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/// Response lifecycle event or output-item event.
22#[derive(Debug, Serialize, Deserialize, Clone)]
23#[serde(untagged)]
24pub enum StreamingCompletionChunk {
25    Response(ResponseChunk),
26    Delta(ItemChunk),
27}
28
29/// What the terminal response event says about the turn, before it maps to
30/// the reply's end.
31#[derive(Debug, Serialize, Deserialize, Clone)]
32pub struct StreamingCompletionResponse {
33    /// Token usage from the terminal response event; `None` when the event
34    /// carried no `usage` object.
35    #[serde(default, skip_serializing_if = "Option::is_none")]
36    pub usage: Option<ResponsesUsage>,
37    /// The complete object-shaped reasoning metadata from the terminal response event.
38    #[serde(default, skip_serializing_if = "Option::is_none")]
39    pub reasoning_metadata: Option<serde_json::Map<String, serde_json::Value>>,
40    /// The effective reasoning context from the terminal response event.
41    #[serde(default, skip_serializing_if = "Option::is_none")]
42    pub reasoning_context: Option<String>,
43    /// The `status` reported by the terminal `response.completed` event.
44    #[serde(default, skip_serializing_if = "Option::is_none")]
45    pub status: Option<ResponseStatus>,
46    /// Why the response stopped short, when the provider said so.
47    #[serde(default, skip_serializing_if = "Option::is_none")]
48    pub incomplete_details: Option<IncompleteDetailsReason>,
49    /// The assistant message ID (`msg_...`) carried by the terminal response's
50    /// output items.
51    ///
52    /// Distinct from [`Self::response_id`] (`resp_...`), which names the whole
53    /// response.
54    #[serde(default, skip_serializing_if = "Option::is_none")]
55    pub message_id: Option<String>,
56    /// The response ID (`resp_...`) reported by the terminal
57    /// `response.completed` event.
58    #[serde(default, skip_serializing_if = "Option::is_none")]
59    pub response_id: Option<String>,
60    /// The model identifier reported by the terminal response event.
61    #[serde(default, skip_serializing_if = "Option::is_none")]
62    pub model: Option<String>,
63}
64
65impl StreamingCompletionResponse {
66    /// Create a terminal record carrying only usage; the remaining metadata is
67    /// filled in from the terminal `response.completed` event as it arrives.
68    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
82/// The provider's end of the reply, from the Responses API's terminal
83/// record, and the issuer of its reasoning for a gateway relaying an
84/// upstream's.
85///
86/// The provider descriptor name is an input for the same reason it is on the
87/// unary conversion: ChatGPT and Copilot stream this exact wire shape, so a
88/// baked-in `"openai"` would mislabel them. The finish reason is left exactly
89/// as the provider reported it; the fold reconciles it with the tool calls
90/// the reply carried.
91fn 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
114/// Combine summaries, content, and encrypted data into one reasoning restatement.
115/// Preserve `provider_id` and wire field order. Return `None` for empty content;
116/// the caller must close the existing block with the returned restatement.
117pub(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    // Same builder as the unary decode, so the restatement and the
125    // non-streaming conversion of one item cannot drift.
126    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/// A response chunk from OpenAI's response API.
145#[derive(Debug, Serialize, Deserialize, Clone)]
146pub struct ResponseChunk {
147    /// The response chunk type
148    #[serde(rename = "type")]
149    pub kind: ResponseChunkKind,
150    /// The response itself
151    pub response: CompletionResponse,
152    /// The item sequence
153    pub sequence_number: u64,
154}
155
156/// Response chunk type.
157/// Renames are used to ensure that this type gets (de)serialized properly.
158#[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
172/// Whether `kind` is a Responses SSE event type this client models.
173///
174/// The union of [`ResponseChunkKind`]'s and [`ItemChunkKind`]'s wire names: a
175/// frame carrying one of these that still fails to deserialize is a data-level
176/// defect in a known event, not an unknown event type, and must surface as an
177/// error rather than be skipped.
178fn 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/// Classify a tagged Responses event with strict decoding for known event types.
206/// Callers must handle `error` and WebSocket `response.done` separately.
207#[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
212/// The assistant message ID (`msg_...`) a terminal response object carries,
213/// which is deliberately not the response's own `resp_...` id.
214fn 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
221/// Fill absent sequence, output, content, and summary indices with zero.
222/// Preserve existing fields and content. Return `None` for invalid or non-object
223/// JSON or serialization failure. Missing output indices can merge distinct items
224/// into slot zero; callers must enable repair only for compatible dialects.
225fn 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
241/// Classified Responses stream frame, whole reply, error, or sentinel.
242pub enum ResponsesEvent {
243    /// Decoded stream frame with raw bytes retained for provider errors.
244    Frame {
245        /// The frame's payload, verbatim.
246        raw: String,
247        /// The frame, decoded.
248        chunk: StreamingCompletionChunk,
249    },
250    /// The unary reply: the response object itself, which carries no `type`
251    /// discriminator because it is not an event.
252    Whole(Box<CompletionResponse>),
253    /// A success whose body is the provider's error envelope instead of a
254    /// response, with the raw body the error preserves. Both the stream's
255    /// own `error` event and a 200 whose whole body is an envelope reach
256    /// this.
257    Failure(String),
258    /// `[DONE]` sentinel. Does not independently establish successful completion.
259    Sentinel,
260}
261
262/// Top-level presence markers for response bodies, including error envelopes.
263const WHOLE_BODY_MARKERS: &[&str] = &["object", "output", "status", "error"];
264
265/// Whether valid JSON carries the `error` event tag.
266fn 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/// The error envelope a success body can carry instead of a response. The
272/// error itself is built from the raw body; this only proves the shape.
273#[derive(Deserialize)]
274struct ErrorEnvelope {
275    // Validate envelope presence while retaining raw bytes for the reported error.
276    #[allow(dead_code)]
277    error: serde_json::Value,
278}
279
280/// The OpenAI Responses wire's decoder: one state machine for the SSE
281/// stream, the unary body and the websocket session.
282pub struct ResponsesDecoder<'id> {
283    /// Stable descriptor name the reply is attributed to: ChatGPT and
284    /// Copilot stream this exact wire shape, so it is an input rather than
285    /// a baked-in `"openai"`.
286    provider: String,
287    /// The terminal event's response object: the provider's own document
288    /// of the turn, and the response's `raw`.
289    document: Option<serde_json::Value>,
290    /// Whether to repair absent envelope indices before retrying classification.
291    /// Selected by dialect, independently of unary or streaming mode.
292    repair_envelopes: bool,
293    /// Whether reasoning belongs to the upstream model's family rather than
294    /// to `provider`, a gateway ([`crate::providers::openai::wire::upstream_reasoning_issuer`]).
295    upstream_reasoning_issuer: bool,
296    /// The terminal record under assembly: what the terminal event says
297    /// about the turn, filled in as the stream reports it.
298    terminal: StreamingCompletionResponse,
299    /// The text part of each message item. A text or refusal delta carrying
300    /// another `item_id` extends that item's part, so two `message` output
301    /// items aggregate as two distinct text parts instead of concatenating.
302    texts: std::collections::HashMap<String, TextPart<'id>>,
303    /// The message item whose text part fragments without an `item_id`
304    /// extend (ChatGPT's envelope-less replays).
305    current_text_item: Option<String>,
306    /// The part fragments without any item id extend.
307    anonymous_text: Option<TextPart<'id>>,
308    /// The message items whose visible text a delta already delivered, and
309    /// whether any fragment arrived that could not be attributed to one.
310    /// The terminal restates the whole turn's output, so its message text
311    /// is published only where no delta delivered it.
312    delta_text_items: std::collections::HashSet<String>,
313    /// Output slots with delivered text. Slot tracking prevents duplicate terminal
314    /// text when a gateway changes item IDs between deltas and restatements.
315    delta_text_slots: std::collections::HashSet<u64>,
316    unattributed_text_delta: bool,
317    /// The message items, by id and output slot, whose content-part extras
318    /// are on their text part. Every snapshot restates them, so only the
319    /// first attaches.
320    extras_items: std::collections::HashSet<String>,
321    extras_slots: std::collections::HashSet<u64>,
322    /// The reasoning part of each output slot, fixed by its first fragment.
323    reasoning: std::collections::HashMap<u64, ReasoningPart<'id>>,
324    /// The output slot of each keyed text part, and of the anonymous one.
325    /// A gateway may name a message's deltas with ids its item events do
326    /// not use, so the slot is what ties a part to its message item.
327    text_slots: std::collections::HashMap<String, u64>,
328    anonymous_text_slot: Option<u64>,
329    /// The message items the reply stated, in the order they opened, and
330    /// the one each output slot holds.
331    message_items: Vec<MessageItem>,
332    message_slots: std::collections::HashMap<u64, usize>,
333}
334
335/// What a `message` output item says about itself rather than its content.
336#[derive(Default)]
337struct MessageItem {
338    id: String,
339    phase: Option<String>,
340    /// Whether `output_item.done` stated the item. What it states wins, as
341    /// its id does for the reply's message id: a gateway may name one item
342    /// differently in each event.
343    done: bool,
344}
345
346/// The event that states a message item.
347#[derive(Clone, Copy, PartialEq, Eq)]
348enum Statement {
349    Added,
350    Done,
351    Terminal,
352}
353
354impl<'id> ResponsesDecoder<'id> {
355    /// A decoder for one reply of `provider`'s Responses endpoint.
356    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    /// Salvage replayed frames that omit their envelope bookkeeping.
380    pub fn with_envelope_repair(mut self) -> Self {
381        self.repair_envelopes = true;
382        self
383    }
384
385    /// Record reasoning as issued by the upstream model's family, for a
386    /// gateway that relays each upstream's own reasoning state.
387    pub fn with_upstream_reasoning_issuer(mut self) -> Self {
388        self.upstream_reasoning_issuer = true;
389        self
390    }
391
392    /// Seed the terminal's usage for a replayed body whose frames may not
393    /// carry one (the unary Responses body's own `usage`).
394    pub fn with_initial_usage(mut self, usage: Option<ResponsesUsage>) -> Self {
395        self.terminal.usage = usage;
396        self
397    }
398
399    /// Recognize sentinels and error events, then try stream, whole-body, and
400    /// error-envelope classifiers in order. Fall through only on corrupt results.
401    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            // Error envelopes must retain the provider's diagnostic on every dialect.
414            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    /// The text part a fragment of the message item `item_id` at
430    /// `output_index` extends.
431    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    /// Record what `statement` says about the message item at
457    /// `output_index`. An item is found by its id first: a terminal may
458    /// shift positions, and a stream whose indices were repaired to zero
459    /// puts every item in one slot. `output_item.added` of an unknown id
460    /// opens a new item, as does `output_item.done` when the slot's item is
461    /// already done; otherwise a statement falls back to the slot. An empty
462    /// id or an absent `phase` erases nothing.
463    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        // The terminal's positions are not the stream's.
489        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    /// Put each message item's `phase` on the text part its content built,
506    /// once, at the end of the reply, when every statement has been read. A
507    /// part is its item's by id, or else by output slot: a gateway may name
508    /// a message's deltas with ids its item events do not use. When the
509    /// reply carries several message items, each part also records its
510    /// item's id, so replay can send each item back as itself.
511    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    /// Record that a delta delivered the visible text of a message item.
561    ///
562    /// The output slot is always recorded; the item id is recorded on top
563    /// of it, because a delta the wire did not attribute extends whichever
564    /// text part is current and is credited to that item, and with no part
565    /// current there is nothing to attribute it to at all.
566    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    /// Whether a delta already delivered the visible text of the message
581    /// item at `output_index` carrying `item_id`.
582    ///
583    /// An unattributable fragment counts for every item: its text is
584    /// already in the choice and nothing on the wire says which item the
585    /// terminal restates, so the merge withholds rather than risk stating
586    /// one turn's text twice.
587    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    /// Publish one message item's visible text as the fragments that built
594    /// it, recording what it delivered so a terminal restating the same
595    /// item merges nothing.
596    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    /// Record that the message item at `output_index` has its extras on its
617    /// text part.
618    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    /// Attach a message item's content-part extras, such as its citation
626    /// annotations, to the text part its deltas built, in content-part order.
627    /// Nothing attaches when a snapshot already did or no delta built a part
628    /// for the item: text stated only by a snapshot publishes its extras
629    /// with it.
630    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    /// Publish nonempty terminal message content when no delta delivered it,
658    /// and otherwise only the extras no `output_item.done` attached.
659    /// Match by output position, item ID, or the unattributed-delta safeguard.
660    fn merge_terminal_body_text(
661        &mut self,
662        response: &CompletionResponse,
663        out: &mut Out<'id, Completion>,
664    ) {
665        // The item's position in `output[]` IS the `output_index` its
666        // stream events carried, which is how a restatement is matched to
667        // the deltas that already delivered it.
668        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    /// Write one output-item event into the reply.
686    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                // A function call interleaving a message item: a later
703                // fragment without an item id cannot belong to that item.
704                self.current_text_item = None;
705                // Without a call_id the item id cannot become a fabricated
706                // tool-result correlator: rig issues the id.
707                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                // Any completed item ends the one it carried; a fragment
723                // arriving afterwards names its own item.
724                self.current_text_item = None;
725                self.push_output_item_done(message.item, output_index, out)?;
726            }
727            // Text and refusal deltas are the same visible-text stream: a
728            // refusal is the assistant's message for that turn.
729            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            // Summary and raw-reasoning deltas differ only in which wire
736            // event carries them; both are fragments of the output item's
737            // reasoning part.
738            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                // The done item restates the call: its name and ids win. An
773                // empty call_id means no provider identity; the fc_* item id
774                // is never substituted, since replay requires a real call_id.
775                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                    // Parsed restatements win.
786                    Ok(arguments) => out.announce_pending(index, arguments),
787                    // A raw restatement is buffered only when no fragment
788                    // preceded it, so its bytes are not stated twice.
789                    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                // The done item completes the call.
799                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                    // The restatement supersedes the part's fragments.
820                    (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                    // A contentless identified item is still replay state.
830                    (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            // An unmodeled output item (e.g. a hosted-tool result such as
848            // `web_search_call`): surfaced raw to the consumer, as the
849            // non-streaming decode preserves it on `CompletionResponse.output`.
850            Output::Unknown(value) => {
851                out.unknown(value.into());
852            }
853            // A compaction item: surfaced raw like an unmodeled item so a
854            // stateless consumer can capture it from the stream.
855            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    /// Record a terminal event's facts: the text no delta delivered, the
868    /// extras no item snapshot attached, how the turn ended, which model
869    /// answered, and which assistant message (`msg_...`, not the response's
870    /// `resp_...`) carried the output.
871    fn record_terminal(&mut self, response: CompletionResponse, out: &mut Out<'id, Completion>) {
872        self.document = serde_json::to_value(&response).ok();
873        // The terminal restates the whole turn, so the message text no delta
874        // delivered is published here: a gateway that states its answer only
875        // in the terminal body still lands it in the choice, and one that
876        // streamed the text first does not state it twice.
877        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    /// The provider ended the turn: close the calls whose done event never
903    /// came (their incomplete arguments drop), then end the reply.
904    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    /// Replay a whole body as the items the stream sends, then end with the
922    /// terminal the body itself is. Structured reasoning suppresses the
923    /// top-level reasoning display string.
924    fn replay_whole_response(
925        &mut self,
926        response: CompletionResponse,
927        mut out: Out<'id, Completion>,
928    ) -> Result<Flow, ProviderError> {
929        // A compatible backend that reports its reasoning as one top-level
930        // string has no stream event for it, so the body is the only place it
931        // is stated. Structured `reasoning` items supersede it: publishing
932        // both would carry one chain of thought twice.
933        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            // Immediate publication keeps the output's order.
952            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        // Replayed bodies omit envelope bookkeeping fields; salvage them
968        // through the same classifier, with the error wording a buffered
969        // body reports.
970        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                    // `response.incomplete` is a genuine terminal (e.g.
1016                    // hitting `max_output_tokens`): the partial output and
1017                    // usage are kept, and the status maps to the finish
1018                    // reason as on the unary path.
1019                    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            // The unary reply is the same turn stated at once.
1033            ResponsesEvent::Whole(response) => self.replay_whole_response(*response, out),
1034            ResponsesEvent::Failure(raw) => {
1035                Err(crate::error::ProviderError::from_provider_body(&raw))
1036            }
1037            // Nothing to write: the provider's end is `response.completed`.
1038            ResponsesEvent::Sentinel => Ok(Flow::More),
1039        }
1040    }
1041}
1042
1043/// Output-item event with its slot index and optional provider item ID.
1044#[derive(Debug, Serialize, Deserialize, Clone)]
1045pub struct ItemChunk {
1046    /// Item ID. Optional.
1047    pub item_id: Option<String>,
1048    /// The output index of the item from a given streamed response.
1049    pub output_index: u64,
1050    /// The item type chunk, as well as the inner data.
1051    #[serde(flatten)]
1052    pub data: ItemChunkKind,
1053}
1054
1055/// The item chunk type from OpenAI's Responses API.
1056#[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    /// Raw-reasoning text restatement. Decoded but not emitted to avoid
1090    /// duplicating accumulated reasoning deltas.
1091    #[serde(rename = "response.reasoning_text.done")]
1092    ReasoningTextDone(OutputTextChunk),
1093    // No `#[serde(other)]` catch-all: unknown event types are triaged by the
1094    // classify layer (`classify_responses_frame` checks the `type` tag against
1095    // `is_known_responses_event_type` BEFORE decoding), so a frame that
1096    // reaches this decoder with an unmodeled tag is a known-set/enum drift
1097    // and must fail loudly (`Corrupt`) rather than be silently absorbed.
1098}
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    /// Unmodeled content part retained verbatim without emitting content.
1123    /// Visible content is delivered by the corresponding delta events.
1124    #[serde(untagged)]
1125    Unknown(serde_json::Value),
1126}
1127
1128/// Decode known tags only with a string text field; preserve unknown or absent tags.
1129/// Nonstring tags return an error. Duplicate keys use the last value retained by
1130/// `serde_json::Value`.
1131impl<'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    // `response.reasoning_summary_text.delta` carries `delta`;
1215    // the `.done` sibling carries the full `text` under the same shape.
1216    #[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;