use serde_json::json;
use super::chunks;
use crate::canonical::{Content, ContentKind, Delta, Event, FinishReason};
use crate::ingress::state::{IngressState, Slot, ThinkAcc, ToolAcc};
pub(crate) fn encode_response(event: &Event, state: &mut IngressState) -> Vec<u8> {
match event {
Event::MessageStart { id, model, .. } => {
state.id.clone_from(id);
state.model.clone_from(model);
chunks::emit(state, json!({"content": "", "role": "assistant"}), None)
}
Event::ContentStart { index, kind } => start(*index, kind, state),
Event::ContentDelta { index, delta } => fragment(*index, delta, state),
Event::ContentStop { index } => {
state.close(*index); Vec::new()
}
Event::Usage(u) => {
state.usage = Some(chunks::usage_json(u));
if state.include_usage {
chunks::emit_usage(state) } else {
Vec::new()
}
}
Event::Finish { reason } => finish(reason, state),
Event::Error(e) => chunks::error(e, state),
Event::End => {
state.finish_stash(); if state.stream {
chunks::sentinel(state)
} else {
chunks::body(state)
}
}
Event::Raw(_) | Event::Other => Vec::new(),
}
}
fn start(index: u32, kind: &ContentKind, state: &mut IngressState) -> Vec<u8> {
match kind {
ContentKind::Text {} => {
state.slots.insert(index, Slot::Text);
Vec::new()
}
ContentKind::ToolUse { id, name } => {
let t = state.tools.len(); state.tools.push(ToolAcc {
id: id.clone(),
name: name.clone(),
args: String::new(),
signature: None,
});
state.slots.insert(index, Slot::Tool(t));
chunks::emit(
state,
json!({"tool_calls": [{
"function": {"arguments": "", "name": name},
"id": id, "index": t, "type": "function"}]}),
None,
)
}
ContentKind::Thinking { id } => {
state.slots.insert(
index,
Slot::Thinking(ThinkAcc {
id: id.clone(),
..ThinkAcc::default()
}),
);
Vec::new()
}
ContentKind::RedactedThinking { data } => {
state
.blocks
.push(Content::RedactedThinking { data: data.clone() });
state.slots.insert(index, Slot::Skip);
Vec::new()
}
_ => {
state.slots.insert(index, Slot::Skip);
Vec::new()
}
}
}
fn fragment(index: u32, delta: &Delta, state: &mut IngressState) -> Vec<u8> {
let payload = match (state.slots.get_mut(&index), delta) {
(Some(Slot::Text), Delta::TextDelta(t)) => {
state.text.push_str(t);
Some(json!({"content": t}))
}
(Some(Slot::Tool(i)), Delta::JsonDelta(a)) => {
let i = *i;
state.tools[i].args.push_str(a);
Some(json!({"tool_calls": [{"function": {"arguments": a}, "index": i}]}))
}
(Some(Slot::Tool(i)), Delta::SignatureDelta(s)) => {
let i = *i; state.tools[i]
.signature
.get_or_insert_with(String::new)
.push_str(s);
None
}
(Some(Slot::Thinking(t)), Delta::ThinkingDelta(d)) => {
t.text.push_str(d);
None
}
(Some(Slot::Thinking(t)), Delta::SignatureDelta(s)) => {
t.signature.get_or_insert_with(String::new).push_str(s);
None
}
(Some(Slot::Thinking(t)), Delta::EncryptedReasoningDelta(e)) => {
t.encrypted.get_or_insert_with(String::new).push_str(e);
None
}
_ => None,
};
payload.map_or_else(Vec::new, |p| chunks::emit(state, p, None))
}
fn finish(reason: &FinishReason, state: &mut IngressState) -> Vec<u8> {
let (refusal, fr) = match reason {
FinishReason::Stop | FinishReason::StopSequence | FinishReason::Pause => (None, "stop"),
FinishReason::Length => (None, "length"),
FinishReason::ToolUse => (None, "tool_calls"),
FinishReason::Refusal {
explanation: Some(t),
..
} => (Some(t.clone()), "stop"),
FinishReason::Refusal { .. } => (None, "content_filter"),
FinishReason::Other(s) => (None, s.as_str()),
};
let fr = fr.to_owned();
let mut out = Vec::new();
if let Some(t) = refusal {
state.refusal.push_str(&t);
out.extend(chunks::emit(state, json!({"refusal": t}), None));
}
out.extend(chunks::emit(state, json!({}), Some(&fr)));
state.finish = Some(fr);
out
}