use crate::pi::converter::Chunk;
use serde_json::Value;
pub struct Accumulator {
active: bool,
text: String,
reasoning: String,
stream_id: String,
}
impl Accumulator {
pub fn new() -> Self {
Self {
active: false,
text: String::new(),
reasoning: String::new(),
stream_id: "pi-stream".to_string(),
}
}
pub fn handle(&mut self, event: &Value) -> Vec<Chunk> {
let t = event.get("type").and_then(|v| v.as_str()).unwrap_or("");
match t {
"message_start" => {
self.active = true;
self.text.clear();
self.reasoning.clear();
self.stream_id = "pi-stream".to_string();
vec![]
}
"message_update" => {
if let Some(ame) = event.get("assistantMessageEvent") {
let sub = ame.get("type").and_then(|v| v.as_str()).unwrap_or("");
match sub {
"text_delta" => {
if let Some(d) = ame.get("delta").and_then(|v| v.as_str()) {
if !d.is_empty() {
self.text.push_str(d);
self.update_stream_id(ame);
}
}
}
"thinking_delta" => {
if let Some(d) = ame.get("delta").and_then(|v| v.as_str()) {
if !d.is_empty() {
self.reasoning.push_str(d);
self.update_stream_id(ame);
}
}
}
_ => {}
}
}
vec![]
}
"message_end" => self.flush(),
"turn_end" => {
if self.active {
self.flush()
} else {
vec![]
}
}
_ => vec![],
}
}
fn update_stream_id(&mut self, ame: &Value) {
if let Some(ci) = ame
.get("contentIndex")
.and_then(|v| v.as_u64().or_else(|| v.as_i64().map(|i| i as u64)))
{
self.stream_id = ci.to_string();
}
}
fn flush(&mut self) -> Vec<Chunk> {
if !self.active {
return vec![];
}
self.active = false;
let mut out = Vec::new();
if !self.reasoning.is_empty() {
out.push(Chunk::Reasoning(std::mem::take(&mut self.reasoning)));
}
if !self.text.is_empty() {
out.push(Chunk::Text(std::mem::take(&mut self.text)));
}
self.stream_id = "pi-stream".to_string();
out
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn msg_start() -> Value {
json!({ "type": "message_start" })
}
fn msg_end() -> Value {
json!({ "type": "message_end" })
}
fn turn_end() -> Value {
json!({ "type": "turn_end" })
}
fn text_delta(d: &str, ci: u64) -> Value {
json!({
"type": "message_update",
"assistantMessageEvent": { "type": "text_delta", "delta": d, "contentIndex": ci }
})
}
fn think_delta(d: &str, ci: u64) -> Value {
json!({
"type": "message_update",
"assistantMessageEvent": { "type": "thinking_delta", "delta": d, "contentIndex": ci }
})
}
#[test]
fn flush_orders_reasoning_before_text() {
let mut acc = Accumulator::new();
assert!(acc.handle(&msg_start()).is_empty());
assert!(acc.handle(&think_delta("hmm", 0)).is_empty());
assert!(acc.handle(&text_delta("hello", 1)).is_empty());
let chunks = acc.handle(&msg_end());
assert_eq!(
chunks,
vec![Chunk::Reasoning("hmm".into()), Chunk::Text("hello".into())]
);
}
#[test]
fn deltas_concatenate_across_updates() {
let mut acc = Accumulator::new();
acc.handle(&msg_start());
acc.handle(&text_delta("hel", 0));
acc.handle(&text_delta("lo", 0));
assert_eq!(acc.handle(&msg_end()), vec![Chunk::Text("hello".into())]);
}
#[test]
fn turn_end_flushes_as_safety_net() {
let mut acc = Accumulator::new();
acc.handle(&msg_start());
acc.handle(&text_delta("hi", 0));
let chunks = acc.handle(&turn_end());
assert_eq!(chunks, vec![Chunk::Text("hi".into())]);
}
#[test]
fn turn_end_without_active_message_is_noop() {
let mut acc = Accumulator::new();
assert!(acc.handle(&turn_end()).is_empty());
}
#[test]
fn updates_without_message_start_are_discarded() {
let mut acc = Accumulator::new();
acc.handle(&text_delta("orphan", 0));
assert!(acc.handle(&msg_end()).is_empty());
}
#[test]
fn empty_deltas_are_ignored() {
let mut acc = Accumulator::new();
acc.handle(&msg_start());
acc.handle(&text_delta("", 0));
acc.handle(&think_delta("", 0));
assert!(acc.handle(&msg_end()).is_empty());
}
#[test]
fn flush_resets_state() {
let mut acc = Accumulator::new();
acc.handle(&msg_start());
acc.handle(&text_delta("x", 0));
assert_eq!(acc.handle(&msg_end()), vec![Chunk::Text("x".into())]);
assert!(acc.handle(&msg_end()).is_empty());
acc.handle(&msg_start());
acc.handle(&text_delta("y", 0));
assert_eq!(acc.handle(&msg_end()), vec![Chunk::Text("y".into())]);
}
#[test]
fn unknown_events_are_noop() {
let mut acc = Accumulator::new();
assert!(acc.handle(&json!({ "type": "agent_start" })).is_empty());
assert!(acc.handle(&json!({})).is_empty());
}
#[test]
fn thinking_without_text_flushes_reasoning_only() {
let mut acc = Accumulator::new();
acc.handle(&msg_start());
acc.handle(&think_delta("just thinking", 0));
assert_eq!(
acc.handle(&msg_end()),
vec![Chunk::Reasoning("just thinking".into())]
);
}
}