Skip to main content

rig_core/providers/openai/responses_api/
streaming.rs

1//! Responses frame classification and event decoding: one block per output
2//! item, in the order the items arrive, each holding the item as the
3//! provider stated it complete.
4//!
5//! Frames are read as JSON and classified by their `type` alone, and every
6//! field is read leniently, so a gateway that omits or retypes a field no
7//! block or finish needs never fails a reply.
8//!
9//! ```
10//! use rig_core::providers::openai::responses_api::streaming::ResponsesDecoder;
11//! let decoder = ResponsesDecoder::new();
12//! # let _ = decoder;
13//! ```
14
15use std::collections::{HashMap, HashSet};
16
17use serde_json::{Value, json};
18
19use crate::completion::{Cost, FinishReason, Usage};
20use crate::error::ProviderError;
21use crate::json_utils::Lenient;
22use crate::message::{Source, SourceLocation};
23use crate::operation::{Block, CallFragment, Completion, Finish};
24use crate::providers::internal::wire;
25use crate::wire::{Decoder, Flow, Out, WireCitation, WireEvent, WireFrame};
26
27/// The item events this decoder reads, after their `response.` prefix.
28const ITEM_EVENTS: &str = "output_item.added output_item.done content_part.added content_part.done \
29    output_text.delta output_text.done refusal.delta refusal.done function_call_arguments.delta \
30    function_call_arguments.done custom_tool_call_input.delta custom_tool_call_input.done \
31    reasoning_summary_part.added reasoning_summary_part.done reasoning_summary_text.delta \
32    reasoning_summary_text.done reasoning_text.delta reasoning_text.done";
33
34/// Whether `kind` is a Responses event type this decoder reads. A frame of
35/// any other type passes through as unknown.
36fn is_known_responses_event_type(kind: &str) -> bool {
37    kind == "error"
38        || is_lifecycle_event(kind)
39        || kind
40            .strip_prefix("response.")
41            .is_some_and(|event| ITEM_EVENTS.split_whitespace().any(|known| known == event))
42}
43
44/// Whether `kind` is a response lifecycle event, which carries the response
45/// object under `response`.
46pub(crate) fn is_lifecycle_event(kind: &str) -> bool {
47    kind.strip_prefix("response.").is_some_and(|event| {
48        "created queued in_progress completed failed incomplete"
49            .split(' ')
50            .any(|known| known == event)
51    })
52}
53
54/// A classified Responses payload.
55#[derive(Debug)]
56pub enum ResponsesEvent {
57    /// A stream event of a known `type`, as the provider sent it.
58    Frame {
59        /// The event's `type`.
60        kind: String,
61        /// The event.
62        frame: Value,
63        /// The frame's payload, verbatim, for the error a failed reply
64        /// reports.
65        raw: String,
66    },
67    /// The unary reply: the response object itself, which carries no `type`
68    /// because it is not an event.
69    Whole(Value),
70    /// The provider's error, as the stream's `error` event or a success
71    /// body holding an error envelope instead of a response.
72    Failure(String),
73    /// The `[DONE]` sentinel. The provider's end is `response.completed`.
74    Sentinel,
75}
76
77/// The keys a response object, or the error envelope a success body can
78/// hold instead, always has one of.
79const WHOLE_BODY_MARKERS: &[&str] = &["object", "output", "status", "id", "error"];
80
81/// A payload that names its `type`, read as the JSON it is.
82struct Tagged(Value);
83
84impl<'de> serde::Deserialize<'de> for Tagged {
85    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
86    where
87        D: serde::Deserializer<'de>,
88    {
89        let value = Value::deserialize(deserializer)?;
90        if value.str("type").is_some() {
91            Ok(Self(value))
92        } else {
93            Err(serde::de::Error::custom("the payload names no `type`"))
94        }
95    }
96}
97
98/// Classify one payload by its `type`; a payload without one is the unary
99/// response object or an error envelope.
100pub fn classify_responses_payload(data: &str) -> WireEvent<ResponsesEvent> {
101    if data.trim() == "[DONE]" {
102        return WireEvent::Known(ResponsesEvent::Sentinel);
103    }
104    wire::classify_or_untagged(
105        data,
106        "type",
107        |data| {
108            wire::classify_tagged_frame::<Tagged>(data, "type", is_known_responses_event_type).map(
109                |Tagged(frame)| match frame.str("type") {
110                    Some("error") => ResponsesEvent::Failure(data.to_owned()),
111                    kind => ResponsesEvent::Frame {
112                        kind: kind.unwrap_or_default().to_owned(),
113                        raw: data.to_owned(),
114                        frame,
115                    },
116                },
117            )
118        },
119        |data| {
120            wire::classify_marker_keyed_frame::<Value>(data, WHOLE_BODY_MARKERS).map(|body| {
121                let error = body.get("error").is_some_and(|error| !error.is_null());
122                if error && body.get("output").is_none() && body.get("status").is_none() {
123                    ResponsesEvent::Failure(data.to_owned())
124                } else {
125                    ResponsesEvent::Whole(body)
126                }
127            })
128        },
129    )
130}
131
132/// What an output item becomes.
133#[derive(Clone, Copy, PartialEq, Eq, Debug)]
134enum Kind {
135    Message,
136    Reasoning,
137    Call,
138    Opaque,
139}
140
141impl Kind {
142    fn of(item: &Value) -> Self {
143        match item.str("type") {
144            Some("message") => Self::Message,
145            Some("reasoning") => Self::Reasoning,
146            Some("function_call" | "custom_tool_call") => Self::Call,
147            _ => Self::Opaque,
148        }
149    }
150}
151
152/// One output item of the reply, as the decoder saw it.
153struct Slot {
154    /// The writer index its block opened at.
155    at: usize,
156    kind: Kind,
157    /// The item id it states.
158    id: Option<String>,
159    open: bool,
160    /// The text its deltas wrote, and the field and part that last
161    /// extended it: a new part of reasoning starts a paragraph, and a second
162    /// field restating the same reasoning is not appended.
163    text: String,
164    field: Option<&'static str>,
165    part: u64,
166    /// A call's streamed argument text, or a custom call's input.
167    arguments: String,
168    /// The argument JSON already written to the call.
169    sent: String,
170    custom: bool,
171    /// Whether the call's id was stated.
172    named: bool,
173    /// Whether the call's tool name was stated.
174    titled: bool,
175}
176
177/// The item types the client executes, which nothing in rig answers: they
178/// are kept in history and never sent back.
179const CLIENT_EXECUTED: &[&str] = &[
180    "computer_call",
181    "local_shell_call",
182    "shell_call",
183    "apply_patch_call",
184    "mcp_approval_request",
185];
186
187fn item_id(item: &Value) -> Option<&str> {
188    item.str("id").filter(|id| !id.is_empty())
189}
190
191/// The argument JSON a done call states: a function call's `arguments`, a
192/// custom call's `input` as `{"input": ...}`; `None` when it states none.
193fn arguments_of(item: &Value) -> Option<String> {
194    if item.str("type") == Some("custom_tool_call") {
195        return item
196            .str("input")
197            .map(|input| json!({ "input": input }).to_string());
198    }
199    match item.get("arguments")? {
200        Value::String(arguments) if arguments.is_empty() => None,
201        Value::String(arguments) => Some(arguments.clone()),
202        Value::Null => None,
203        arguments => Some(arguments.to_string()),
204    }
205}
206
207/// The text a done item states: a message's parts joined, refusals
208/// included; reasoning's summary, or its raw content when it has none.
209fn text_of(item: &Value) -> String {
210    let texts = |key: &str| -> Vec<&str> {
211        match item.get(key) {
212            Some(Value::String(text)) => vec![text.as_str()],
213            _ => item
214                .arr(key)
215                .iter()
216                .filter_map(|part| {
217                    part.as_str()
218                        .or_else(|| part.str("text"))
219                        .or_else(|| part.str("refusal"))
220                })
221                .collect(),
222        }
223    };
224    match Kind::of(item) {
225        Kind::Message => texts("content").concat(),
226        Kind::Reasoning => Some(texts("summary"))
227            .filter(|summary| !summary.is_empty())
228            .unwrap_or_else(|| texts("content"))
229            .join("\n\n"),
230        Kind::Call | Kind::Opaque => String::new(),
231    }
232}
233
234/// `item`, a done item that states no text, stating `text`, so the item
235/// replays as the block it is.
236fn stating(mut item: Value, kind: Kind, text: &str) -> Value {
237    let (key, part) = match kind {
238        Kind::Message => (
239            "content",
240            json!({"type": "output_text", "text": text, "annotations": []}),
241        ),
242        _ => ("summary", json!({"type": "summary_text", "text": text})),
243    };
244    if let Some(fields) = item.as_object_mut() {
245        fields.insert(key.to_owned(), json!([part]));
246    }
247    item
248}
249
250/// The provider's own words for a failed or cancelled response.
251fn provider_message(response: &Value, fallback: &str) -> String {
252    let parts: Vec<&str> = ["/error/code", "/error/message"]
253        .iter()
254        .filter_map(|pointer| response.at(pointer).and_then(Value::as_str))
255        .collect();
256    if parts.is_empty() {
257        fallback.to_owned()
258    } else {
259        parts.join(": ")
260    }
261}
262
263/// How the turn ended, for every status the API documents. A response
264/// without a status ended as pi reads it, with a stop. A status that is
265/// not a documented end, including one still `queued` or `in_progress`, is
266/// a failed turn.
267pub(crate) fn finish_reason_of(response: &Value) -> (FinishReason, Option<String>) {
268    let reason = response
269        .at("/incomplete_details/reason")
270        .and_then(Value::as_str)
271        .filter(|reason| !reason.is_empty());
272    let other = |reason: &str, error: String| (FinishReason::Other(reason.to_owned()), Some(error));
273    match response.str("status") {
274        None | Some("completed") => (FinishReason::Stop, None),
275        Some("incomplete") => match reason {
276            Some("max_output_tokens") => (FinishReason::Length, None),
277            Some("content_filter") => (FinishReason::ContentFilter, None),
278            Some(reason) => other(
279                &format!("incomplete: {reason}"),
280                format!("Response incomplete: {reason}"),
281            ),
282            None => other(
283                "incomplete",
284                "Response incomplete without a provider reason".to_owned(),
285            ),
286        },
287        Some(status @ ("failed" | "cancelled")) => other(
288            status,
289            provider_message(response, &format!("Response {status}")),
290        ),
291        Some(status @ ("queued" | "in_progress")) => {
292            other(status, format!("Response ended while {status}"))
293        }
294        Some(status) => other(
295            status,
296            format!("Response ended with the unknown status `{status}`"),
297        ),
298    }
299}
300
301/// The usage a Responses `usage` value reports; a counter that is absent or
302/// not a count is unreported. Its cost is the one the provider reports.
303pub(crate) fn usage_of(usage: &Value) -> Usage {
304    let count = |pointer: &str| usage.at(pointer).and_then(Lenient::as_u64_lenient);
305    Usage {
306        input_tokens: count("/input_tokens"),
307        output_tokens: count("/output_tokens"),
308        total_tokens: count("/total_tokens"),
309        cached_input_tokens: count("/input_tokens_details/cached_tokens"),
310        cache_creation_input_tokens: count("/input_tokens_details/cache_write_tokens"),
311        reasoning_tokens: count("/output_tokens_details/reasoning_tokens"),
312        ..Usage::default()
313    }
314    .cost(reported_cost(usage))
315}
316
317/// xAI cost ticks per USD.
318const TICKS_PER_USD: f64 = 1e10;
319
320/// The total a `usage` value reports in USD: xAI's integer
321/// `cost_in_usd_ticks`, or OpenRouter's `cost` in credits of one USD.
322fn reported_cost(usage: &Value) -> Option<Cost> {
323    let total = match usage.u64("cost_in_usd_ticks") {
324        Some(ticks) => ticks as f64 / TICKS_PER_USD,
325        None => usage.f64("cost").filter(|cost| cost.is_finite())?,
326    };
327    Some(Cost::from_total(total))
328}
329
330/// The citations a message item's `output_text` parts state, in part order.
331/// The API does not document the unit of `start_index` and `end_index` and
332/// quotes no text, so each cites the whole block and its offsets stay in the
333/// item.
334fn citations_of(item: &Value) -> Vec<WireCitation> {
335    item.arr("content")
336        .iter()
337        .flat_map(|part| part.arr("annotations"))
338        .filter_map(source_of)
339        .map(|source| WireCitation::new(None, vec![source]))
340        .collect()
341}
342
343/// The source an annotation names; `None` for an annotation that names
344/// none or is of an unknown type.
345fn source_of(annotation: &Value) -> Option<Source> {
346    let text = |key: &str| {
347        annotation
348            .str(key)
349            .filter(|text| !text.is_empty())
350            .map(str::to_owned)
351    };
352    let location = match annotation.str("type")? {
353        "url_citation" => SourceLocation::Url { url: text("url")? },
354        "file_citation" | "file_path" | "container_file_citation" => SourceLocation::File {
355            file_id: text("file_id")?,
356            filename: text("filename"),
357            container_id: text("container_id"),
358        },
359        _ => return None,
360    };
361    let source = Source::new(location);
362    Some(match text("title") {
363        Some(title) => source.title(title),
364        None => source,
365    })
366}
367
368/// The OpenAI Responses wire's decoder: one state machine for the SSE
369/// stream, the unary body and the websocket session.
370///
371/// Each output item becomes one block. It opens when it is announced, or at
372/// its first delta, with no provider item, and its done item closes it,
373/// becoming the block's native: an item the provider never stated complete
374/// replays from its canonical fields. A done item merges into what streamed:
375/// a call keeps the id and name it was announced with and the arguments it
376/// streamed when its done item leaves them out, and text the done item
377/// leaves out stays, written into the item. The terminal response is the
378/// whole output: it finishes the items no done event carried, writes the
379/// ones the stream never stated, which is how a unary body decodes, and
380/// backfills the ciphertext of reasoning done without it.
381#[derive(Default)]
382pub struct ResponsesDecoder {
383    /// Every item, in the order it arrived.
384    slots: Vec<Slot>,
385    /// The latest slot at each output index the provider named.
386    indexed: HashMap<usize, usize>,
387    /// The slot the last item event addressed.
388    current: Option<usize>,
389    /// The writer indices of reasoning done and waiting for the item after
390    /// it to complete.
391    held: Vec<usize>,
392}
393
394impl ResponsesDecoder {
395    /// A decoder for one reply of a Responses endpoint.
396    pub fn new() -> Self {
397        Self::default()
398    }
399
400    /// The slot `frame` addresses, an event for an item of `kind`: the one
401    /// at its output index, else the one its item id names, else the open
402    /// one the stream is on, when the frame names no index or that item
403    /// was opened by a frame that named none (an envelope-less delta, then
404    /// its envelope-full done item). `None` for an item no slot holds yet.
405    fn addressed(&self, frame: &Value, kind: Kind) -> Result<Option<usize>, ProviderError> {
406        let index = output_index(frame)?;
407        if let Some(slot) = index.and_then(|index| self.indexed.get(&index)) {
408            return Ok(Some(*slot));
409        }
410        let id = frame
411            .at("/item/id")
412            .and_then(Value::as_str)
413            .or_else(|| frame.str("item_id"));
414        if let Some(slot) = id.filter(|id| !id.is_empty()).and_then(|id| {
415            self.slots
416                .iter()
417                .rposition(|slot| slot.id.as_deref() == Some(id))
418        }) {
419            return Ok(Some(slot));
420        }
421        Ok(self.current.filter(|current| {
422            (index.is_none() || !self.indexed.values().any(|slot| slot == current))
423                && self
424                    .slots
425                    .get(*current)
426                    .is_some_and(|slot| slot.open && slot.kind == kind)
427        }))
428    }
429
430    /// Open the block of `item` as a new slot at output `index`, closing what
431    /// is still open there.
432    fn added(
433        &mut self,
434        index: Option<usize>,
435        item: &Value,
436        out: &mut Out<'_, Completion>,
437    ) -> Result<usize, ProviderError> {
438        if let Some(previous) = index.and_then(|index| self.indexed.get(&index).copied()) {
439            self.vacate(previous, out)?;
440        }
441        let at = match index {
442            Some(index) if !self.slots.iter().any(|slot| slot.at == index) => index,
443            _ => out.fresh_index(),
444        };
445        let kind = Kind::of(item);
446        let custom = item.str("type") == Some("custom_tool_call");
447        match kind {
448            Kind::Message => out.open(at, Block::Text, Value::Null)?,
449            Kind::Reasoning => out.open(at, Block::Reasoning { redacted: false }, Value::Null)?,
450            Kind::Call => out.fragment(
451                Some(at),
452                CallFragment {
453                    id: item.str("call_id"),
454                    name: item.str("name"),
455                    arguments: None,
456                },
457            )?,
458            // An item that names no type is kept and never sent back.
459            Kind::Opaque => {
460                let replay = item
461                    .str("type")
462                    .is_some_and(|kind| !CLIENT_EXECUTED.contains(&kind));
463                out.open(at, Block::Opaque { replay }, item.clone())?;
464            }
465        }
466        let arguments = item.str(if custom { "input" } else { "arguments" });
467        self.slots.push(Slot {
468            at,
469            kind,
470            id: item_id(item).map(str::to_owned),
471            open: true,
472            text: String::new(),
473            field: None,
474            part: 0,
475            arguments: arguments.unwrap_or_default().to_owned(),
476            sent: String::new(),
477            custom,
478            named: item.str("call_id").is_some_and(|id| !id.is_empty()),
479            titled: item.str("name").is_some_and(|name| !name.is_empty()),
480        });
481        let slot = self.slots.len() - 1;
482        if let Some(index) = index {
483            self.indexed.insert(index, slot);
484        }
485        self.current = Some(slot);
486        if let Some(call) = self
487            .slots
488            .get_mut(slot)
489            .filter(|slot| slot.kind == Kind::Call)
490        {
491            send(call, out)?;
492        }
493        Ok(slot)
494    }
495
496    /// Close the slot's block before another item takes its place: the
497    /// provider never stated it complete. A call keeps what streamed.
498    fn vacate(&mut self, slot: usize, out: &mut Out<'_, Completion>) -> Result<(), ProviderError> {
499        let Some(open) = self.slots.get_mut(slot).filter(|slot| slot.open) else {
500            return Ok(());
501        };
502        open.open = false;
503        let at = open.at;
504        if open.kind == Kind::Call {
505            flush(open, out)?;
506        }
507        out.close(at)?;
508        self.release(false, out)
509    }
510
511    /// End the reasoning held for the item after it: complete when that
512    /// item completed, else as never stated complete.
513    fn release(
514        &mut self,
515        complete: bool,
516        out: &mut Out<'_, Completion>,
517    ) -> Result<(), ProviderError> {
518        std::mem::take(&mut self.held)
519            .into_iter()
520            .try_for_each(|at| {
521                if complete {
522                    out.finish(at)
523                } else {
524                    out.close(at)
525                }
526            })
527    }
528
529    /// Append a text delta of `field` to the item `frame` addresses,
530    /// opening a block for an item the stream never announced.
531    fn text(
532        &mut self,
533        frame: &Value,
534        field: &'static str,
535        part: &str,
536        out: &mut Out<'_, Completion>,
537    ) -> Result<(), ProviderError> {
538        let delta = frame.str("delta").unwrap_or_default();
539        if delta.is_empty() {
540            return Ok(());
541        }
542        let part = frame.u64(part).unwrap_or(0);
543        let kind = match field {
544            "message" => Kind::Message,
545            _ => Kind::Reasoning,
546        };
547        let known = self.addressed(frame, kind)?.and_then(|at| {
548            let slot = self.slots.get(at)?;
549            Some((at, slot.open, slot.kind == kind))
550        });
551        let slot = match known {
552            Some((slot, true, true)) => slot,
553            // A delta for an item that already closed is left to the item,
554            // which states its text.
555            Some((_, false, true)) => return Ok(()),
556            _ => self.added(
557                output_index(frame)?,
558                &json!({"type": if kind == Kind::Message { "message" } else { "reasoning" }}),
559                out,
560            )?,
561        };
562        let Some(streamed) = self.slots.get_mut(slot).filter(|slot| slot.open) else {
563            return Ok(());
564        };
565        if streamed.field.is_some_and(|known| known != field) {
566            return Ok(());
567        }
568        if streamed.part != part && kind != Kind::Message && !streamed.text.is_empty() {
569            streamed.text.push_str("\n\n");
570            out.push(streamed.at, "\n\n")?;
571        }
572        streamed.field = Some(field);
573        streamed.part = part;
574        streamed.text.push_str(delta);
575        out.push(streamed.at, delta)
576    }
577
578    /// Append an argument delta to the call `frame` addresses and write
579    /// it to the call, which streams it.
580    fn arguments(
581        &mut self,
582        frame: &Value,
583        out: &mut Out<'_, Completion>,
584    ) -> Result<(), ProviderError> {
585        let slot = self.addressed(frame, Kind::Call)?;
586        if let Some(call) = slot
587            .and_then(|slot| self.slots.get_mut(slot))
588            .filter(|slot| slot.open && slot.kind == Kind::Call)
589        {
590            call.arguments
591                .push_str(frame.str("delta").unwrap_or_default());
592            send(call, out)?;
593        }
594        Ok(())
595    }
596
597    /// The item at `slot` is done: it merges into what streamed and becomes
598    /// the block's native.
599    fn done(
600        &mut self,
601        slot: usize,
602        item: Value,
603        out: &mut Out<'_, Completion>,
604    ) -> Result<(), ProviderError> {
605        let Some(done) = self.slots.get_mut(slot) else {
606            return Ok(());
607        };
608        done.open = false;
609        if done.id.is_none() {
610            done.id = item_id(&item).map(str::to_owned);
611        }
612        let at = done.at;
613        let mut item = item;
614        let mut nameless = false;
615        match done.kind {
616            Kind::Call => {
617                nameless = !done.titled && item.str("name").is_none_or(str::is_empty);
618                let arguments = arguments_of(&item).unwrap_or_else(|| streamed_arguments(done));
619                let rest = remainder(done, &arguments, out)?;
620                out.fragment(
621                    Some(at),
622                    CallFragment {
623                        id: item.str("call_id").filter(|_| !done.named),
624                        name: item.str("name"),
625                        arguments: Some(&rest),
626                    },
627                )?;
628            }
629            Kind::Message | Kind::Reasoning => {
630                let text = text_of(&item);
631                if text.is_empty() {
632                    // Raw reasoning text the done item leaves out stays out:
633                    // written in as a summary, it would replay one the
634                    // provider never produced.
635                    if !done.text.is_empty() && done.field != Some("reasoning") {
636                        item = stating(item, done.kind, &done.text);
637                    }
638                } else if let Some(rest) = text.strip_prefix(done.text.as_str()) {
639                    out.push(at, rest)?;
640                } else {
641                    // The done item states the whole text, as pi takes it.
642                    out.restate(at, &text)?;
643                }
644                let citations = citations_of(&item);
645                if done.kind == Kind::Message && !citations.is_empty() {
646                    // The item restates every annotation, so its list
647                    // replaces any an earlier snapshot gave.
648                    out.set_citations(at, citations);
649                }
650            }
651            Kind::Opaque => {}
652        }
653        let reasoning = done.kind == Kind::Reasoning;
654        out.edit(at, |native| *native = item)?;
655        if reasoning {
656            // Reasoning is sent only with the item it precedes, so it is
657            // complete only once that item is.
658            self.held.push(at);
659            return Ok(());
660        }
661        out.finish(at)?;
662        // The writer drops a call with no name, so the reasoning it held
663        // has no item to go with.
664        self.release(!nameless, out)
665    }
666
667    /// Write one output-item event into the reply.
668    fn item_event(
669        &mut self,
670        kind: &str,
671        frame: Value,
672        out: &mut Out<'_, Completion>,
673    ) -> Result<(), ProviderError> {
674        match kind {
675            "response.output_item.added" | "response.output_item.done" => {
676                let Some(item) = frame.get("item").filter(|item| item.is_object()) else {
677                    return Ok(());
678                };
679                let index = output_index(&frame)?;
680                if kind == "response.output_item.added" {
681                    return self.added(index, item, out).map(|_| ());
682                }
683                let known = self.addressed(&frame, Kind::of(item))?;
684                // At a named output index the item is the one there, as pi
685                // reads it, whatever id it states.
686                let restates = |slot: &Slot| {
687                    slot.kind == Kind::of(item)
688                        && (index.is_some()
689                            || item_id(item)
690                                .zip(slot.id.as_deref())
691                                .is_none_or(|(id, known)| id == known))
692                };
693                let known = known.and_then(|at| {
694                    let slot = self.slots.get(at)?;
695                    let same = item_id(item).is_some() && slot.id.as_deref() == item_id(item);
696                    Some((at, slot.open && restates(slot), same))
697                });
698                let slot = match known {
699                    Some((slot, true, _)) => slot,
700                    // The item this one restates is already done.
701                    Some((_, false, true)) => return Ok(()),
702                    _ => self.added(index, item, out)?,
703                };
704                self.done(slot, item.clone(), out)
705            }
706            // A refusal is the assistant's message for that turn.
707            "response.output_text.delta" | "response.refusal.delta" => {
708                self.text(&frame, "message", "content_index", out)
709            }
710            "response.reasoning_summary_text.delta" => {
711                self.text(&frame, "summary", "summary_index", out)
712            }
713            "response.reasoning_text.delta" => self.text(&frame, "reasoning", "content_index", out),
714            "response.function_call_arguments.delta" | "response.custom_tool_call_input.delta" => {
715                self.arguments(&frame, out)
716            }
717            _ => Ok(()),
718        }
719    }
720
721    /// The slot the terminal's item at output `index` restates, among those
722    /// no earlier terminal item took: the one with its id; else the one the
723    /// stream gave that index when it is of the item's kind, as pi reads an
724    /// index whatever id it states; else the first of its kind the stream
725    /// gave no index and no other id; else the first of its kind with no
726    /// other id that the stream put at an index the terminal holds no item
727    /// of that kind at, a stream whose indices are shifted against it.
728    fn restated(
729        &self,
730        index: usize,
731        item: &Value,
732        output: &[Value],
733        taken: &HashSet<usize>,
734    ) -> Option<usize> {
735        let id = item_id(item);
736        let free = |at: &usize| {
737            !taken.contains(at)
738                && self
739                    .slots
740                    .get(*at)
741                    .is_some_and(|slot| slot.kind == Kind::of(item))
742        };
743        let find = |fits: &dyn Fn(usize, &Slot) -> bool| {
744            self.slots
745                .iter()
746                .enumerate()
747                .find(|(at, slot)| free(at) && fits(*at, slot))
748                .map(|(at, _)| at)
749        };
750        id.and_then(|id| find(&|_, slot| slot.id.as_deref() == Some(id)))
751            .or_else(|| self.indexed.get(&index).copied().filter(free))
752            .or_else(|| {
753                find(&|at, slot| {
754                    !self.indexed.values().any(|indexed| *indexed == at)
755                        && (id.is_none() || slot.id.is_none())
756                })
757            })
758            .or_else(|| {
759                find(&|at, slot| {
760                    (id.is_none() || slot.id.is_none())
761                        && self.indexed.iter().any(|(streamed, indexed)| {
762                            *indexed == at
763                                && output
764                                    .get(*streamed)
765                                    .is_none_or(|other| Kind::of(other) != slot.kind)
766                        })
767                })
768            })
769    }
770
771    /// The terminal response: finish or write each item of its output,
772    /// backfill the ciphertext of reasoning done without it, then end with
773    /// how the turn ended. A body that reports its reasoning as one
774    /// top-level string and no reasoning item is written that reasoning
775    /// first.
776    fn finish(
777        &mut self,
778        response: Value,
779        mut out: Out<'_, Completion>,
780    ) -> Result<Flow, ProviderError> {
781        let output = response.arr("output");
782        if self.slots.is_empty()
783            && !output.iter().any(|item| Kind::of(item) == Kind::Reasoning)
784            && let Some(reasoning) = response.str("reasoning").filter(|text| !text.is_empty())
785        {
786            // Index 0 orders it before the output; it closes at once, so the
787            // output's first item may open there after it.
788            out.whole(
789                0,
790                Block::Reasoning { redacted: false },
791                Value::Null,
792                reasoning,
793            )?;
794        }
795        let mut taken = HashSet::new();
796        for (index, item) in output
797            .iter()
798            .enumerate()
799            .filter(|(_, item)| item.is_object())
800        {
801            let known = self.restated(index, item, output, &taken).and_then(|at| {
802                let slot = self.slots.get(at)?;
803                Some((at, slot.open, slot.kind, slot.at))
804            });
805            let slot = match known {
806                Some((slot, false, kind, at)) => {
807                    taken.insert(slot);
808                    let ciphertext = item
809                        .get("encrypted_content")
810                        .filter(|cipher| cipher.as_str().is_some_and(|cipher| !cipher.is_empty()));
811                    if let (Kind::Reasoning, Some(ciphertext)) = (kind, ciphertext) {
812                        out.edit(at, |native| {
813                            let stated = native
814                                .str("encrypted_content")
815                                .is_some_and(|cipher| !cipher.is_empty());
816                            if let (false, Some(fields)) = (stated, native.as_object_mut()) {
817                                fields.insert("encrypted_content".to_owned(), ciphertext.clone());
818                            }
819                        })?;
820                    }
821                    continue;
822                }
823                Some((slot, ..)) => slot,
824                None => self.added(
825                    Some(index).filter(|index| !self.indexed.contains_key(index)),
826                    item,
827                    &mut out,
828                )?,
829            };
830            taken.insert(slot);
831            self.done(slot, item.clone(), &mut out)?;
832        }
833        // Reasoning still waiting is complete when nothing after it is left
834        // open.
835        let open = self.slots.iter().any(|slot| slot.open);
836        self.release(!open, &mut out)?;
837        // A call never done shows what streamed, in a turn that fails.
838        for slot in self.slots.iter_mut().filter(|slot| slot.open) {
839            if slot.kind == Kind::Call {
840                flush(slot, &mut out)?;
841            }
842        }
843        let (reason, error) = finish_reason_of(&response);
844        let end = Finish {
845            usage: response.get("usage").map(usage_of).unwrap_or_default(),
846            reason: Some(reason),
847            response_id: item_id(&response).map(str::to_owned),
848            model: response
849                .str("model")
850                .filter(|model| !model.is_empty())
851                .map(str::to_owned),
852            error,
853        };
854        Ok(out.end(end))
855    }
856}
857
858/// Give `slot`'s call what it streamed, before it closes undone.
859fn flush(slot: &mut Slot, out: &mut Out<'_, Completion>) -> Result<(), ProviderError> {
860    let arguments = streamed_arguments(slot);
861    let rest = remainder(slot, &arguments, out)?;
862    let fragment = CallFragment {
863        arguments: Some(&rest),
864        ..CallFragment::default()
865    };
866    out.fragment(Some(slot.at), fragment)
867}
868
869/// Write the argument JSON `slot`'s call streamed so far that the call has
870/// not been given yet. A custom call's input streams inside its
871/// `{"input": ...}` object, escaped as the object's string, so its
872/// fragments join to the object its end states.
873fn send(slot: &mut Slot, out: &mut Out<'_, Completion>) -> Result<(), ProviderError> {
874    if slot.arguments.is_empty() {
875        return Ok(());
876    }
877    let streamed = if slot.custom {
878        let quoted = Value::String(slot.arguments.clone()).to_string();
879        // The string without its closing quote: later input extends it.
880        format!("{{\"input\":{}", &quoted[..quoted.len() - 1])
881    } else {
882        slot.arguments.clone()
883    };
884    let Some(rest) = streamed
885        .strip_prefix(slot.sent.as_str())
886        .filter(|rest| !rest.is_empty())
887    else {
888        return Ok(());
889    };
890    let fragment = CallFragment {
891        arguments: Some(rest),
892        ..CallFragment::default()
893    };
894    out.fragment(Some(slot.at), fragment)?;
895    slot.sent = streamed;
896    Ok(())
897}
898
899/// What of `arguments`, the whole JSON `slot`'s call states, the call has
900/// not been given yet. When `arguments` does not extend what streamed, it
901/// replaces the call's text and nothing more streams: the call's end states
902/// it.
903fn remainder(
904    slot: &mut Slot,
905    arguments: &str,
906    out: &mut Out<'_, Completion>,
907) -> Result<String, ProviderError> {
908    let rest = match arguments.strip_prefix(slot.sent.as_str()) {
909        Some(rest) => rest.to_owned(),
910        None => {
911            out.restate(slot.at, arguments)?;
912            String::new()
913        }
914    };
915    arguments.clone_into(&mut slot.sent);
916    Ok(rest)
917}
918
919/// The argument JSON of what `slot`'s call streamed.
920fn streamed_arguments(slot: &Slot) -> String {
921    if slot.custom {
922        json!({ "input": slot.arguments }).to_string()
923    } else {
924        slot.arguments.clone()
925    }
926}
927
928/// The output index `frame` names.
929fn output_index(frame: &Value) -> Result<Option<usize>, ProviderError> {
930    frame
931        .u64("output_index")
932        .map(|index| {
933            usize::try_from(index).map_err(|_| {
934                ProviderError::Response(format!("output_index {index} is out of range"))
935            })
936        })
937        .transpose()
938}
939
940impl<'id> Decoder<'id, Completion> for ResponsesDecoder {
941    type Event = ResponsesEvent;
942
943    fn classify(&self, frame: WireFrame) -> WireEvent<ResponsesEvent> {
944        classify_responses_payload(&frame.as_str())
945    }
946
947    fn decode(
948        &mut self,
949        event: ResponsesEvent,
950        mut out: Out<'id, Completion>,
951    ) -> Result<Flow, ProviderError> {
952        // A terminal may state items the stream never announced; they take
953        // their place by output index.
954        out.order_by_index();
955        match event {
956            ResponsesEvent::Frame { kind, frame, raw } => match kind.as_str() {
957                // `response.incomplete` is a genuine terminal that keeps the
958                // partial output and usage. It is incomplete whatever its
959                // `status` says, where pi reads a missing status as a stop.
960                "response.completed" | "response.incomplete" => {
961                    let mut response = frame
962                        .get("response")
963                        .filter(|response| response.is_object())
964                        .cloned()
965                        .unwrap_or_else(|| json!({}));
966                    if let (true, Some(fields)) =
967                        (kind == "response.incomplete", response.as_object_mut())
968                    {
969                        fields.insert("status".to_owned(), json!("incomplete"));
970                    }
971                    self.finish(response, out)
972                }
973                "response.failed" => Err(ProviderError::from_provider_body(raw)),
974                kind if is_lifecycle_event(kind) => Ok(Flow::More),
975                kind => {
976                    self.item_event(kind, frame, &mut out)?;
977                    Ok(Flow::More)
978                }
979            },
980            // The unary reply is the terminal response with no stream before it.
981            ResponsesEvent::Whole(body) => self.finish(body, out),
982            ResponsesEvent::Failure(raw) => Err(ProviderError::from_provider_body(raw)),
983            // Nothing to write: the provider's end is `response.completed`.
984            ResponsesEvent::Sentinel => Ok(Flow::More),
985        }
986    }
987}
988
989pub(crate) mod document;
990
991#[cfg(test)]
992mod tests;