use std::collections::{HashMap, HashSet};
use serde_json::{Value, json};
use crate::completion::{Cost, FinishReason, Usage};
use crate::error::ProviderError;
use crate::json_utils::Lenient;
use crate::message::{Source, SourceLocation};
use crate::operation::{Block, CallFragment, Completion, Finish};
use crate::providers::internal::wire;
use crate::wire::{Decoder, Flow, Out, WireCitation, WireEvent, WireFrame};
const ITEM_EVENTS: &str = "output_item.added output_item.done content_part.added content_part.done \
output_text.delta output_text.done refusal.delta refusal.done function_call_arguments.delta \
function_call_arguments.done custom_tool_call_input.delta custom_tool_call_input.done \
reasoning_summary_part.added reasoning_summary_part.done reasoning_summary_text.delta \
reasoning_summary_text.done reasoning_text.delta reasoning_text.done";
fn is_known_responses_event_type(kind: &str) -> bool {
kind == "error"
|| is_lifecycle_event(kind)
|| kind
.strip_prefix("response.")
.is_some_and(|event| ITEM_EVENTS.split_whitespace().any(|known| known == event))
}
pub(crate) fn is_lifecycle_event(kind: &str) -> bool {
kind.strip_prefix("response.").is_some_and(|event| {
"created queued in_progress completed failed incomplete"
.split(' ')
.any(|known| known == event)
})
}
#[derive(Debug)]
pub enum ResponsesEvent {
Frame {
kind: String,
frame: Value,
raw: String,
},
Whole(Value),
Failure(String),
Sentinel,
}
const WHOLE_BODY_MARKERS: &[&str] = &["object", "output", "status", "id", "error"];
struct Tagged(Value);
impl<'de> serde::Deserialize<'de> for Tagged {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let value = Value::deserialize(deserializer)?;
if value.str("type").is_some() {
Ok(Self(value))
} else {
Err(serde::de::Error::custom("the payload names no `type`"))
}
}
}
pub fn classify_responses_payload(data: &str) -> WireEvent<ResponsesEvent> {
if data.trim() == "[DONE]" {
return WireEvent::Known(ResponsesEvent::Sentinel);
}
wire::classify_or_untagged(
data,
"type",
|data| {
wire::classify_tagged_frame::<Tagged>(data, "type", is_known_responses_event_type).map(
|Tagged(frame)| match frame.str("type") {
Some("error") => ResponsesEvent::Failure(data.to_owned()),
kind => ResponsesEvent::Frame {
kind: kind.unwrap_or_default().to_owned(),
raw: data.to_owned(),
frame,
},
},
)
},
|data| {
wire::classify_marker_keyed_frame::<Value>(data, WHOLE_BODY_MARKERS).map(|body| {
let error = body.get("error").is_some_and(|error| !error.is_null());
if error && body.get("output").is_none() && body.get("status").is_none() {
ResponsesEvent::Failure(data.to_owned())
} else {
ResponsesEvent::Whole(body)
}
})
},
)
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
enum Kind {
Message,
Reasoning,
Call,
Opaque,
}
impl Kind {
fn of(item: &Value) -> Self {
match item.str("type") {
Some("message") => Self::Message,
Some("reasoning") => Self::Reasoning,
Some("function_call" | "custom_tool_call") => Self::Call,
_ => Self::Opaque,
}
}
}
struct Slot {
at: usize,
kind: Kind,
id: Option<String>,
open: bool,
text: String,
field: Option<&'static str>,
part: u64,
arguments: String,
sent: String,
custom: bool,
named: bool,
titled: bool,
}
const CLIENT_EXECUTED: &[&str] = &[
"computer_call",
"local_shell_call",
"shell_call",
"apply_patch_call",
"mcp_approval_request",
];
fn item_id(item: &Value) -> Option<&str> {
item.str("id").filter(|id| !id.is_empty())
}
fn arguments_of(item: &Value) -> Option<String> {
if item.str("type") == Some("custom_tool_call") {
return item
.str("input")
.map(|input| json!({ "input": input }).to_string());
}
match item.get("arguments")? {
Value::String(arguments) if arguments.is_empty() => None,
Value::String(arguments) => Some(arguments.clone()),
Value::Null => None,
arguments => Some(arguments.to_string()),
}
}
fn text_of(item: &Value) -> String {
let texts = |key: &str| -> Vec<&str> {
match item.get(key) {
Some(Value::String(text)) => vec![text.as_str()],
_ => item
.arr(key)
.iter()
.filter_map(|part| {
part.as_str()
.or_else(|| part.str("text"))
.or_else(|| part.str("refusal"))
})
.collect(),
}
};
match Kind::of(item) {
Kind::Message => texts("content").concat(),
Kind::Reasoning => Some(texts("summary"))
.filter(|summary| !summary.is_empty())
.unwrap_or_else(|| texts("content"))
.join("\n\n"),
Kind::Call | Kind::Opaque => String::new(),
}
}
fn stating(mut item: Value, kind: Kind, text: &str) -> Value {
let (key, part) = match kind {
Kind::Message => (
"content",
json!({"type": "output_text", "text": text, "annotations": []}),
),
_ => ("summary", json!({"type": "summary_text", "text": text})),
};
if let Some(fields) = item.as_object_mut() {
fields.insert(key.to_owned(), json!([part]));
}
item
}
fn provider_message(response: &Value, fallback: &str) -> String {
let parts: Vec<&str> = ["/error/code", "/error/message"]
.iter()
.filter_map(|pointer| response.at(pointer).and_then(Value::as_str))
.collect();
if parts.is_empty() {
fallback.to_owned()
} else {
parts.join(": ")
}
}
pub(crate) fn finish_reason_of(response: &Value) -> (FinishReason, Option<String>) {
let reason = response
.at("/incomplete_details/reason")
.and_then(Value::as_str)
.filter(|reason| !reason.is_empty());
let other = |reason: &str, error: String| (FinishReason::Other(reason.to_owned()), Some(error));
match response.str("status") {
None | Some("completed") => (FinishReason::Stop, None),
Some("incomplete") => match reason {
Some("max_output_tokens") => (FinishReason::Length, None),
Some("content_filter") => (FinishReason::ContentFilter, None),
Some(reason) => other(
&format!("incomplete: {reason}"),
format!("Response incomplete: {reason}"),
),
None => other(
"incomplete",
"Response incomplete without a provider reason".to_owned(),
),
},
Some(status @ ("failed" | "cancelled")) => other(
status,
provider_message(response, &format!("Response {status}")),
),
Some(status @ ("queued" | "in_progress")) => {
other(status, format!("Response ended while {status}"))
}
Some(status) => other(
status,
format!("Response ended with the unknown status `{status}`"),
),
}
}
pub(crate) fn usage_of(usage: &Value) -> Usage {
let count = |pointer: &str| usage.at(pointer).and_then(Lenient::as_u64_lenient);
Usage {
input_tokens: count("/input_tokens"),
output_tokens: count("/output_tokens"),
total_tokens: count("/total_tokens"),
cached_input_tokens: count("/input_tokens_details/cached_tokens"),
cache_creation_input_tokens: count("/input_tokens_details/cache_write_tokens"),
reasoning_tokens: count("/output_tokens_details/reasoning_tokens"),
..Usage::default()
}
.cost(reported_cost(usage))
}
const TICKS_PER_USD: f64 = 1e10;
fn reported_cost(usage: &Value) -> Option<Cost> {
let total = match usage.u64("cost_in_usd_ticks") {
Some(ticks) => ticks as f64 / TICKS_PER_USD,
None => usage.f64("cost").filter(|cost| cost.is_finite())?,
};
Some(Cost::from_total(total))
}
fn citations_of(item: &Value) -> Vec<WireCitation> {
item.arr("content")
.iter()
.flat_map(|part| part.arr("annotations"))
.filter_map(source_of)
.map(|source| WireCitation::new(None, vec![source]))
.collect()
}
fn source_of(annotation: &Value) -> Option<Source> {
let text = |key: &str| {
annotation
.str(key)
.filter(|text| !text.is_empty())
.map(str::to_owned)
};
let location = match annotation.str("type")? {
"url_citation" => SourceLocation::Url { url: text("url")? },
"file_citation" | "file_path" | "container_file_citation" => SourceLocation::File {
file_id: text("file_id")?,
filename: text("filename"),
container_id: text("container_id"),
},
_ => return None,
};
let source = Source::new(location);
Some(match text("title") {
Some(title) => source.title(title),
None => source,
})
}
#[derive(Default)]
pub struct ResponsesDecoder {
slots: Vec<Slot>,
indexed: HashMap<usize, usize>,
current: Option<usize>,
held: Vec<usize>,
}
impl ResponsesDecoder {
pub fn new() -> Self {
Self::default()
}
fn addressed(&self, frame: &Value, kind: Kind) -> Result<Option<usize>, ProviderError> {
let index = output_index(frame)?;
if let Some(slot) = index.and_then(|index| self.indexed.get(&index)) {
return Ok(Some(*slot));
}
let id = frame
.at("/item/id")
.and_then(Value::as_str)
.or_else(|| frame.str("item_id"));
if let Some(slot) = id.filter(|id| !id.is_empty()).and_then(|id| {
self.slots
.iter()
.rposition(|slot| slot.id.as_deref() == Some(id))
}) {
return Ok(Some(slot));
}
Ok(self.current.filter(|current| {
(index.is_none() || !self.indexed.values().any(|slot| slot == current))
&& self
.slots
.get(*current)
.is_some_and(|slot| slot.open && slot.kind == kind)
}))
}
fn added(
&mut self,
index: Option<usize>,
item: &Value,
out: &mut Out<'_, Completion>,
) -> Result<usize, ProviderError> {
if let Some(previous) = index.and_then(|index| self.indexed.get(&index).copied()) {
self.vacate(previous, out)?;
}
let at = match index {
Some(index) if !self.slots.iter().any(|slot| slot.at == index) => index,
_ => out.fresh_index(),
};
let kind = Kind::of(item);
let custom = item.str("type") == Some("custom_tool_call");
match kind {
Kind::Message => out.open(at, Block::Text, Value::Null)?,
Kind::Reasoning => out.open(at, Block::Reasoning { redacted: false }, Value::Null)?,
Kind::Call => out.fragment(
Some(at),
CallFragment {
id: item.str("call_id"),
name: item.str("name"),
arguments: None,
},
)?,
Kind::Opaque => {
let replay = item
.str("type")
.is_some_and(|kind| !CLIENT_EXECUTED.contains(&kind));
out.open(at, Block::Opaque { replay }, item.clone())?;
}
}
let arguments = item.str(if custom { "input" } else { "arguments" });
self.slots.push(Slot {
at,
kind,
id: item_id(item).map(str::to_owned),
open: true,
text: String::new(),
field: None,
part: 0,
arguments: arguments.unwrap_or_default().to_owned(),
sent: String::new(),
custom,
named: item.str("call_id").is_some_and(|id| !id.is_empty()),
titled: item.str("name").is_some_and(|name| !name.is_empty()),
});
let slot = self.slots.len() - 1;
if let Some(index) = index {
self.indexed.insert(index, slot);
}
self.current = Some(slot);
if let Some(call) = self
.slots
.get_mut(slot)
.filter(|slot| slot.kind == Kind::Call)
{
send(call, out)?;
}
Ok(slot)
}
fn vacate(&mut self, slot: usize, out: &mut Out<'_, Completion>) -> Result<(), ProviderError> {
let Some(open) = self.slots.get_mut(slot).filter(|slot| slot.open) else {
return Ok(());
};
open.open = false;
let at = open.at;
if open.kind == Kind::Call {
flush(open, out)?;
}
out.close(at)?;
self.release(false, out)
}
fn release(
&mut self,
complete: bool,
out: &mut Out<'_, Completion>,
) -> Result<(), ProviderError> {
std::mem::take(&mut self.held)
.into_iter()
.try_for_each(|at| {
if complete {
out.finish(at)
} else {
out.close(at)
}
})
}
fn text(
&mut self,
frame: &Value,
field: &'static str,
part: &str,
out: &mut Out<'_, Completion>,
) -> Result<(), ProviderError> {
let delta = frame.str("delta").unwrap_or_default();
if delta.is_empty() {
return Ok(());
}
let part = frame.u64(part).unwrap_or(0);
let kind = match field {
"message" => Kind::Message,
_ => Kind::Reasoning,
};
let known = self.addressed(frame, kind)?.and_then(|at| {
let slot = self.slots.get(at)?;
Some((at, slot.open, slot.kind == kind))
});
let slot = match known {
Some((slot, true, true)) => slot,
Some((_, false, true)) => return Ok(()),
_ => self.added(
output_index(frame)?,
&json!({"type": if kind == Kind::Message { "message" } else { "reasoning" }}),
out,
)?,
};
let Some(streamed) = self.slots.get_mut(slot).filter(|slot| slot.open) else {
return Ok(());
};
if streamed.field.is_some_and(|known| known != field) {
return Ok(());
}
if streamed.part != part && kind != Kind::Message && !streamed.text.is_empty() {
streamed.text.push_str("\n\n");
out.push(streamed.at, "\n\n")?;
}
streamed.field = Some(field);
streamed.part = part;
streamed.text.push_str(delta);
out.push(streamed.at, delta)
}
fn arguments(
&mut self,
frame: &Value,
out: &mut Out<'_, Completion>,
) -> Result<(), ProviderError> {
let slot = self.addressed(frame, Kind::Call)?;
if let Some(call) = slot
.and_then(|slot| self.slots.get_mut(slot))
.filter(|slot| slot.open && slot.kind == Kind::Call)
{
call.arguments
.push_str(frame.str("delta").unwrap_or_default());
send(call, out)?;
}
Ok(())
}
fn done(
&mut self,
slot: usize,
item: Value,
out: &mut Out<'_, Completion>,
) -> Result<(), ProviderError> {
let Some(done) = self.slots.get_mut(slot) else {
return Ok(());
};
done.open = false;
if done.id.is_none() {
done.id = item_id(&item).map(str::to_owned);
}
let at = done.at;
let mut item = item;
let mut nameless = false;
match done.kind {
Kind::Call => {
nameless = !done.titled && item.str("name").is_none_or(str::is_empty);
let arguments = arguments_of(&item).unwrap_or_else(|| streamed_arguments(done));
let rest = remainder(done, &arguments, out)?;
out.fragment(
Some(at),
CallFragment {
id: item.str("call_id").filter(|_| !done.named),
name: item.str("name"),
arguments: Some(&rest),
},
)?;
}
Kind::Message | Kind::Reasoning => {
let text = text_of(&item);
if text.is_empty() {
if !done.text.is_empty() && done.field != Some("reasoning") {
item = stating(item, done.kind, &done.text);
}
} else if let Some(rest) = text.strip_prefix(done.text.as_str()) {
out.push(at, rest)?;
} else {
out.restate(at, &text)?;
}
let citations = citations_of(&item);
if done.kind == Kind::Message && !citations.is_empty() {
out.set_citations(at, citations);
}
}
Kind::Opaque => {}
}
let reasoning = done.kind == Kind::Reasoning;
out.edit(at, |native| *native = item)?;
if reasoning {
self.held.push(at);
return Ok(());
}
out.finish(at)?;
self.release(!nameless, out)
}
fn item_event(
&mut self,
kind: &str,
frame: Value,
out: &mut Out<'_, Completion>,
) -> Result<(), ProviderError> {
match kind {
"response.output_item.added" | "response.output_item.done" => {
let Some(item) = frame.get("item").filter(|item| item.is_object()) else {
return Ok(());
};
let index = output_index(&frame)?;
if kind == "response.output_item.added" {
return self.added(index, item, out).map(|_| ());
}
let known = self.addressed(&frame, Kind::of(item))?;
let restates = |slot: &Slot| {
slot.kind == Kind::of(item)
&& (index.is_some()
|| item_id(item)
.zip(slot.id.as_deref())
.is_none_or(|(id, known)| id == known))
};
let known = known.and_then(|at| {
let slot = self.slots.get(at)?;
let same = item_id(item).is_some() && slot.id.as_deref() == item_id(item);
Some((at, slot.open && restates(slot), same))
});
let slot = match known {
Some((slot, true, _)) => slot,
Some((_, false, true)) => return Ok(()),
_ => self.added(index, item, out)?,
};
self.done(slot, item.clone(), out)
}
"response.output_text.delta" | "response.refusal.delta" => {
self.text(&frame, "message", "content_index", out)
}
"response.reasoning_summary_text.delta" => {
self.text(&frame, "summary", "summary_index", out)
}
"response.reasoning_text.delta" => self.text(&frame, "reasoning", "content_index", out),
"response.function_call_arguments.delta" | "response.custom_tool_call_input.delta" => {
self.arguments(&frame, out)
}
_ => Ok(()),
}
}
fn restated(
&self,
index: usize,
item: &Value,
output: &[Value],
taken: &HashSet<usize>,
) -> Option<usize> {
let id = item_id(item);
let free = |at: &usize| {
!taken.contains(at)
&& self
.slots
.get(*at)
.is_some_and(|slot| slot.kind == Kind::of(item))
};
let find = |fits: &dyn Fn(usize, &Slot) -> bool| {
self.slots
.iter()
.enumerate()
.find(|(at, slot)| free(at) && fits(*at, slot))
.map(|(at, _)| at)
};
id.and_then(|id| find(&|_, slot| slot.id.as_deref() == Some(id)))
.or_else(|| self.indexed.get(&index).copied().filter(free))
.or_else(|| {
find(&|at, slot| {
!self.indexed.values().any(|indexed| *indexed == at)
&& (id.is_none() || slot.id.is_none())
})
})
.or_else(|| {
find(&|at, slot| {
(id.is_none() || slot.id.is_none())
&& self.indexed.iter().any(|(streamed, indexed)| {
*indexed == at
&& output
.get(*streamed)
.is_none_or(|other| Kind::of(other) != slot.kind)
})
})
})
}
fn finish(
&mut self,
response: Value,
mut out: Out<'_, Completion>,
) -> Result<Flow, ProviderError> {
let output = response.arr("output");
if self.slots.is_empty()
&& !output.iter().any(|item| Kind::of(item) == Kind::Reasoning)
&& let Some(reasoning) = response.str("reasoning").filter(|text| !text.is_empty())
{
out.whole(
0,
Block::Reasoning { redacted: false },
Value::Null,
reasoning,
)?;
}
let mut taken = HashSet::new();
for (index, item) in output
.iter()
.enumerate()
.filter(|(_, item)| item.is_object())
{
let known = self.restated(index, item, output, &taken).and_then(|at| {
let slot = self.slots.get(at)?;
Some((at, slot.open, slot.kind, slot.at))
});
let slot = match known {
Some((slot, false, kind, at)) => {
taken.insert(slot);
let ciphertext = item
.get("encrypted_content")
.filter(|cipher| cipher.as_str().is_some_and(|cipher| !cipher.is_empty()));
if let (Kind::Reasoning, Some(ciphertext)) = (kind, ciphertext) {
out.edit(at, |native| {
let stated = native
.str("encrypted_content")
.is_some_and(|cipher| !cipher.is_empty());
if let (false, Some(fields)) = (stated, native.as_object_mut()) {
fields.insert("encrypted_content".to_owned(), ciphertext.clone());
}
})?;
}
continue;
}
Some((slot, ..)) => slot,
None => self.added(
Some(index).filter(|index| !self.indexed.contains_key(index)),
item,
&mut out,
)?,
};
taken.insert(slot);
self.done(slot, item.clone(), &mut out)?;
}
let open = self.slots.iter().any(|slot| slot.open);
self.release(!open, &mut out)?;
for slot in self.slots.iter_mut().filter(|slot| slot.open) {
if slot.kind == Kind::Call {
flush(slot, &mut out)?;
}
}
let (reason, error) = finish_reason_of(&response);
let end = Finish {
usage: response.get("usage").map(usage_of).unwrap_or_default(),
reason: Some(reason),
response_id: item_id(&response).map(str::to_owned),
model: response
.str("model")
.filter(|model| !model.is_empty())
.map(str::to_owned),
error,
};
Ok(out.end(end))
}
}
fn flush(slot: &mut Slot, out: &mut Out<'_, Completion>) -> Result<(), ProviderError> {
let arguments = streamed_arguments(slot);
let rest = remainder(slot, &arguments, out)?;
let fragment = CallFragment {
arguments: Some(&rest),
..CallFragment::default()
};
out.fragment(Some(slot.at), fragment)
}
fn send(slot: &mut Slot, out: &mut Out<'_, Completion>) -> Result<(), ProviderError> {
if slot.arguments.is_empty() {
return Ok(());
}
let streamed = if slot.custom {
let quoted = Value::String(slot.arguments.clone()).to_string();
format!("{{\"input\":{}", "ed[..quoted.len() - 1])
} else {
slot.arguments.clone()
};
let Some(rest) = streamed
.strip_prefix(slot.sent.as_str())
.filter(|rest| !rest.is_empty())
else {
return Ok(());
};
let fragment = CallFragment {
arguments: Some(rest),
..CallFragment::default()
};
out.fragment(Some(slot.at), fragment)?;
slot.sent = streamed;
Ok(())
}
fn remainder(
slot: &mut Slot,
arguments: &str,
out: &mut Out<'_, Completion>,
) -> Result<String, ProviderError> {
let rest = match arguments.strip_prefix(slot.sent.as_str()) {
Some(rest) => rest.to_owned(),
None => {
out.restate(slot.at, arguments)?;
String::new()
}
};
arguments.clone_into(&mut slot.sent);
Ok(rest)
}
fn streamed_arguments(slot: &Slot) -> String {
if slot.custom {
json!({ "input": slot.arguments }).to_string()
} else {
slot.arguments.clone()
}
}
fn output_index(frame: &Value) -> Result<Option<usize>, ProviderError> {
frame
.u64("output_index")
.map(|index| {
usize::try_from(index).map_err(|_| {
ProviderError::Response(format!("output_index {index} is out of range"))
})
})
.transpose()
}
impl<'id> Decoder<'id, Completion> for ResponsesDecoder {
type Event = ResponsesEvent;
fn classify(&self, frame: WireFrame) -> WireEvent<ResponsesEvent> {
classify_responses_payload(&frame.as_str())
}
fn decode(
&mut self,
event: ResponsesEvent,
mut out: Out<'id, Completion>,
) -> Result<Flow, ProviderError> {
out.order_by_index();
match event {
ResponsesEvent::Frame { kind, frame, raw } => match kind.as_str() {
"response.completed" | "response.incomplete" => {
let mut response = frame
.get("response")
.filter(|response| response.is_object())
.cloned()
.unwrap_or_else(|| json!({}));
if let (true, Some(fields)) =
(kind == "response.incomplete", response.as_object_mut())
{
fields.insert("status".to_owned(), json!("incomplete"));
}
self.finish(response, out)
}
"response.failed" => Err(ProviderError::from_provider_body(raw)),
kind if is_lifecycle_event(kind) => Ok(Flow::More),
kind => {
self.item_event(kind, frame, &mut out)?;
Ok(Flow::More)
}
},
ResponsesEvent::Whole(body) => self.finish(body, out),
ResponsesEvent::Failure(raw) => Err(ProviderError::from_provider_body(raw)),
ResponsesEvent::Sentinel => Ok(Flow::More),
}
}
}
pub(crate) mod document;
#[cfg(test)]
mod tests;