use serde_json::Value;
#[derive(Debug, Default)]
pub struct Accumulator {
terminal: Option<Value>,
outcome: Option<Outcome>,
generated: bool,
id: Option<String>,
error: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Outcome {
Completed,
Incomplete,
Failed,
}
impl Accumulator {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub const fn generated(&self) -> bool {
self.generated
}
#[must_use]
pub const fn outcome(&self) -> Option<Outcome> {
self.outcome
}
#[must_use]
pub const fn terminal(&self) -> Option<&Value> {
self.terminal.as_ref()
}
#[must_use]
pub fn id(&self) -> Option<&str> {
self.id.as_deref()
}
#[must_use]
pub fn error(&self) -> Option<&str> {
self.error.as_deref()
}
pub fn event(&mut self, name: &str, data: &str) {
let Ok(value) = serde_json::from_str::<Value>(data) else {
return;
};
let kind = if name.is_empty() {
value
.get("type")
.and_then(Value::as_str)
.unwrap_or_default()
} else {
name
};
match kind {
"response.created" => {
self.id = value
.get("response")
.and_then(|r| r.get("id"))
.and_then(Value::as_str)
.map(ToOwned::to_owned);
}
"response.output_text.delta"
| "response.refusal.delta"
| "response.reasoning_summary_text.delta"
| "response.function_call_arguments.delta" => self.generated = true,
"response.completed" => self.finish(&value, Outcome::Completed),
"response.incomplete" => self.finish(&value, Outcome::Incomplete),
"response.failed" => self.finish(&value, Outcome::Failed),
"error" => {
self.error = Some(
value
.get("message")
.and_then(Value::as_str)
.unwrap_or("the provider sent an error event")
.to_owned(),
);
}
_ => {}
}
}
fn finish(&mut self, value: &Value, outcome: Outcome) {
self.outcome = Some(outcome);
if let Some(r) = value.get("response") {
self.terminal = Some(r.clone());
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn feed(acc: &mut Accumulator, events: &[(&str, &str)]) {
for (name, data) in events {
acc.event(name, data);
}
}
const CREATED: (&str, &str) = (
"response.created",
r#"{"type":"response.created","sequence_number":0,
"response":{"id":"resp_abc","status":"in_progress","usage":null}}"#,
);
const DELTA: (&str, &str) = (
"response.output_text.delta",
r#"{"type":"response.output_text.delta","sequence_number":3,"delta":"Hel"}"#,
);
const COMPLETED: (&str, &str) = (
"response.completed",
r#"{"type":"response.completed","sequence_number":9,
"response":{"id":"resp_abc","status":"completed",
"output":[{"type":"message","content":[{"type":"output_text","text":"Hello"}]}],
"usage":{"input_tokens":25,"output_tokens":15,"total_tokens":40}}}"#,
);
#[test]
fn a_whole_stream_yields_the_terminal_response_object() {
let mut acc = Accumulator::new();
feed(&mut acc, &[CREATED, DELTA, COMPLETED]);
assert_eq!(acc.outcome(), Some(Outcome::Completed));
assert_eq!(acc.id(), Some("resp_abc"));
let terminal = acc.terminal().expect("the terminal event carries it");
assert_eq!(terminal["usage"]["output_tokens"], 15);
assert_eq!(terminal["status"], "completed");
}
#[test]
fn nothing_before_the_terminal_event_says_what_it_cost() {
let mut acc = Accumulator::new();
feed(&mut acc, &[CREATED, DELTA]);
assert!(
acc.terminal().is_none(),
"usage arrives only in the terminal event; a severed stream cannot \
bill what it does not know"
);
assert!(acc.generated(), "but it certainly generated");
}
#[test]
fn generation_is_not_claimed_before_any_output() {
let mut acc = Accumulator::new();
feed(&mut acc, &[CREATED]);
assert!(
!acc.generated(),
"an id and an in_progress status are not evidence that tokens were \
produced — treating them as such makes every failed handshake \
un-retryable"
);
}
#[test]
fn reasoning_refusal_and_tool_deltas_all_count_as_generation() {
for event in [
"response.refusal.delta",
"response.reasoning_summary_text.delta",
"response.function_call_arguments.delta",
] {
let mut acc = Accumulator::new();
acc.event(event, &format!(r#"{{"type":"{event}","delta":"x"}}"#));
assert!(acc.generated(), "{event} is billed output");
}
}
#[test]
fn an_incomplete_response_is_terminal_and_carries_usage() {
let mut acc = Accumulator::new();
feed(
&mut acc,
&[
CREATED,
DELTA,
(
"response.incomplete",
r#"{"type":"response.incomplete","response":{"status":"incomplete",
"incomplete_details":{"reason":"max_output_tokens"},
"usage":{"input_tokens":25,"output_tokens":999}}}"#,
),
],
);
assert_eq!(acc.outcome(), Some(Outcome::Incomplete));
let t = acc.terminal().unwrap();
assert_eq!(t["incomplete_details"]["reason"], "max_output_tokens");
assert_eq!(t["usage"]["output_tokens"], 999);
}
#[test]
fn a_failed_response_is_terminal() {
let mut acc = Accumulator::new();
feed(
&mut acc,
&[
CREATED,
(
"response.failed",
r#"{"type":"response.failed","response":{"status":"failed",
"error":{"code":"server_error","message":"boom"},
"usage":{"input_tokens":25,"output_tokens":3}}}"#,
),
],
);
assert_eq!(acc.outcome(), Some(Outcome::Failed));
assert_eq!(acc.terminal().unwrap()["status"], "failed");
}
#[test]
fn an_error_event_is_captured() {
let mut acc = Accumulator::new();
acc.event("error", r#"{"type":"error","message":"upstream exploded"}"#);
assert_eq!(acc.error(), Some("upstream exploded"));
}
#[test]
fn unknown_event_types_are_ignored() {
let mut acc = Accumulator::new();
feed(
&mut acc,
&[
CREATED,
(
"response.some_new_tool.delta",
r#"{"type":"response.some_new_tool.delta"}"#,
),
COMPLETED,
],
);
assert_eq!(acc.outcome(), Some(Outcome::Completed));
}
#[test]
fn a_malformed_event_does_not_lose_the_rest() {
let mut acc = Accumulator::new();
acc.event("response.created", "{not json");
feed(&mut acc, &[CREATED, COMPLETED]);
assert_eq!(acc.outcome(), Some(Outcome::Completed));
}
}