xagent-pi 0.2.10

Self-contained local brain (chat UI + API + SSE) for the Pi agent, tunneled into xagent-service.
//! Accumulates Pi `message_update` deltas; flushes on `message_end` (and `turn_end`
//! as a safety net). Flush order is reasoning-first, text-second.

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(),
        }
    }

    /// Returns chunks to emit (non-empty only on flush boundaries).
    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() {
        // deltas before message_start accumulate but never flush.
        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())]);
        // a second message_end without a new message_start is a noop
        assert!(acc.handle(&msg_end()).is_empty());
        // and a new message_start starts fresh
        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())]
        );
    }
}