use serde_json::{json, Value};
use crate::canonical::{CanonicalError, ContentKind, Delta, Event, Role};
use crate::protocol::json::{http_error, parse, text_of, u32_at};
use crate::protocol::{DecodeState, Frame, OpenBlock};
mod terminal;
pub(super) fn decode_full(
body: &[u8],
state: &mut DecodeState,
) -> Result<Vec<Event>, CanonicalError> {
let response = parse(body)?;
let mut out = event(&created(&response), state);
for (oi, item) in response["output"]
.as_array()
.into_iter()
.flatten()
.enumerate()
{
explode_item(oi as u32, item, state, &mut out);
}
out.extend(event(&completed(&response), state));
Ok(out)
}
fn explode_item(oi: u32, item: &Value, state: &mut DecodeState, out: &mut Vec<Event>) {
let added = json!({ "type": "response.output_item.added", "output_index": oi, "item": item });
out.extend(event(&added, state));
match item["type"].as_str().unwrap_or_default() {
"message" => {
for (ci, part) in item["content"].as_array().into_iter().flatten().enumerate() {
explode_part(oi, ci as u32, part, state, out);
}
}
"function_call" => out.extend(event(
&arg_delta(
oi,
"response.function_call_arguments.delta",
&text_of(item, "arguments"),
),
state,
)),
"reasoning" => {
for s in item["summary"].as_array().into_iter().flatten() {
out.extend(event(
&arg_delta(
oi,
"response.reasoning_summary_text.delta",
&text_of(s, "text"),
),
state,
));
}
}
_ => {}
}
let done = json!({ "type": "response.output_item.done", "output_index": oi, "item": item });
out.extend(event(&done, state));
}
fn explode_part(oi: u32, ci: u32, part: &Value, state: &mut DecodeState, out: &mut Vec<Event>) {
let added = json!({ "type": "response.content_part.added", "output_index": oi, "content_index": ci, "part": part });
out.extend(event(&added, state));
let delta = json!({ "type": "response.output_text.delta", "output_index": oi, "content_index": ci, "delta": part["text"] });
out.extend(event(&delta, state));
}
fn arg_delta(oi: u32, ty: &str, delta: &str) -> Value {
json!({ "type": ty, "output_index": oi, "delta": delta })
}
fn created(response: &Value) -> Value {
json!({ "type": "response.created", "response": response })
}
fn completed(response: &Value) -> Value {
json!({ "type": "response.completed", "response": response })
}
pub(super) fn decode(frame: Frame, state: &mut DecodeState) -> Result<Vec<Event>, CanonicalError> {
if let Some(status) = frame.status {
return Ok(vec![Event::Error(http_error(&frame.data, status))]); }
Ok(event(&parse(&frame.data)?, state))
}
fn event(v: &Value, state: &mut DecodeState) -> Vec<Event> {
match v["type"].as_str().unwrap_or_default() {
"response.created" | "response.in_progress" => message_start(v, state),
"response.output_item.added" => item_added(v, state),
"response.content_part.added" => part_added(v, state),
"response.output_text.delta" => delta(v, state, Delta::TextDelta),
"response.function_call_arguments.delta" => delta(v, state, Delta::JsonDelta),
"response.reasoning_summary_text.delta" => delta(v, state, Delta::ThinkingDelta),
"response.reasoning_text.delta" => delta(v, state, Delta::ThinkingDelta), "response.refusal.delta" => {
state.refusal.push_str(&text_of(v, "delta")); vec![]
}
"response.output_item.done" => item_done(v, state),
"response.completed" => terminal::completed(v, state),
"response.incomplete" => terminal::incomplete(v, state),
"response.failed" | "response.error" => vec![Event::Error(terminal::stream_error(v))], _ => vec![],
}
}
fn message_start(v: &Value, state: &mut DecodeState) -> Vec<Event> {
if state.started {
return vec![];
}
state.started = true;
let r = &v["response"];
vec![Event::message_start(
r["id"].as_str().map(str::to_owned),
r["model"].as_str().map(str::to_owned),
Role::Assistant,
)]
}
fn item_added(v: &Value, state: &mut DecodeState) -> Vec<Event> {
let item = &v["item"];
let kind = match item["type"].as_str() {
Some("function_call") => ContentKind::ToolUse {
id: text_of(item, "call_id"),
name: text_of(item, "name"),
},
Some("reasoning") => ContentKind::Thinking {},
_ => return vec![],
};
let index = canonical(state, part_key(v)); open(state, index, kind.clone());
vec![Event::ContentStart { index, kind }]
}
fn part_added(v: &Value, state: &mut DecodeState) -> Vec<Event> {
if v["part"]["type"].as_str() != Some("output_text") {
return vec![];
}
let index = canonical(state, part_key(v));
open(state, index, ContentKind::Text {});
vec![Event::ContentStart {
index,
kind: ContentKind::Text {},
}]
}
fn delta(v: &Value, state: &mut DecodeState, wrap: fn(String) -> Delta) -> Vec<Event> {
let Some(&index) = state.part_index.get(&part_key(v)) else {
return vec![]; };
if !state.open.contains_key(&index) {
return vec![]; }
let frag = text_of(v, "delta");
vec![Event::ContentDelta {
index,
delta: wrap(frag),
}]
}
fn item_done(v: &Value, state: &mut DecodeState) -> Vec<Event> {
let oi = u32_at(v, "output_index");
let mut indices: Vec<u32> = state
.part_index
.iter()
.filter(|((o, _), c)| *o == oi && state.open.contains_key(c))
.map(|(_, &c)| c)
.collect();
indices.sort_unstable();
indices
.into_iter()
.map(|index| {
state.open.remove(&index);
Event::ContentStop { index }
})
.collect()
}
fn open(state: &mut DecodeState, index: u32, kind: ContentKind) {
state.open.insert(index, OpenBlock { kind });
}
fn canonical(state: &mut DecodeState, key: (u32, u32)) -> u32 {
let next = state.part_index.len() as u32;
*state.part_index.entry(key).or_insert(next)
}
fn part_key(v: &Value) -> (u32, u32) {
(u32_at(v, "output_index"), u32_at(v, "content_index"))
}