chimera-core 0.2.0

Shared traits, types, and transport for the chimera AI agent SDK
Documentation
use tokio::io::{AsyncBufReadExt, BufReader};
use tokio::process::ChildStdout;

use crate::AgentError;

const MAX_BUFFER_SIZE: usize = 1024 * 1024; // 1 MB

pub struct JsonLineReader {
    reader: BufReader<ChildStdout>,
    buffer: String,
}

impl JsonLineReader {
    pub fn new(stdout: ChildStdout) -> Self {
        Self {
            reader: BufReader::new(stdout),
            buffer: String::new(),
        }
    }

    /// Read the next complete JSON value from stdout.
    ///
    /// Handles JSON objects spanning multiple lines by buffering and
    /// retrying the parse. Returns `None` on EOF.
    pub async fn next_value<T: serde::de::DeserializeOwned>(
        &mut self,
    ) -> Option<Result<T, AgentError>> {
        loop {
            let mut line = String::new();
            match self.reader.read_line(&mut line).await {
                Ok(0) => {
                    if self.buffer.is_empty() {
                        return None;
                    }
                    let buf = std::mem::take(&mut self.buffer);
                    return Some(
                        serde_json::from_str(&buf)
                            .map_err(|source| AgentError::JsonParse { line: buf, source }),
                    );
                }
                Ok(_) => {
                    self.buffer.push_str(&line);
                    if self.buffer.len() > MAX_BUFFER_SIZE {
                        self.buffer.clear();
                        return Some(Err(AgentError::BufferOverflow {
                            limit: MAX_BUFFER_SIZE,
                        }));
                    }
                    match take_next_value(&mut self.buffer) {
                        Ok(Some(value)) => return Some(Ok(value)),
                        Ok(None) => continue,
                        Err(source) => {
                            let line = std::mem::take(&mut self.buffer);
                            return Some(Err(AgentError::JsonParse { line, source }));
                        }
                    }
                }
                Err(e) => {
                    return Some(Err(AgentError::Other {
                        message: "failed to read stdout".into(),
                        source: Some(Box::new(e)),
                    }));
                }
            }
        }
    }
}

fn take_next_value<T: serde::de::DeserializeOwned>(
    buffer: &mut String,
) -> Result<Option<T>, serde_json::Error> {
    let trimmed = buffer.trim_start_matches(['\r', '\n']);
    if trimmed.len() != buffer.len() {
        *buffer = trimmed.to_string();
    }
    if buffer.is_empty() {
        return Ok(None);
    }

    let mut stream = serde_json::Deserializer::from_str(buffer).into_iter::<T>();
    match stream.next() {
        Some(Ok(value)) => {
            let consumed = stream.byte_offset();
            let remainder = buffer[consumed..].trim_start_matches(['\r', '\n']);
            *buffer = remainder.to_string();
            Ok(Some(value))
        }
        Some(Err(err)) if err.is_eof() => Ok(None),
        Some(Err(err)) => Err(err),
        None => Ok(None),
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    #[tokio::test]
    async fn json_line_reader_single_line() {
        // JsonLineReader requires ChildStdout, so we test the parsing
        // logic it relies on directly.
        let json = r#"{"key": "value"}"#;
        let parsed: serde_json::Value = serde_json::from_str(json).unwrap();
        assert_eq!(parsed["key"], "value");
    }

    #[test]
    fn buffer_overflow_threshold() {
        assert_eq!(MAX_BUFFER_SIZE, 1024 * 1024);
    }

    #[test]
    fn take_next_value_consumes_one_value_and_keeps_remainder() {
        let mut buffer = "{\"a\":1}\n{\"b\":2}\n".to_string();
        let first = take_next_value::<serde_json::Value>(&mut buffer)
            .expect("first value should parse")
            .expect("first value should be ready");
        assert_eq!(first["a"], 1);
        assert_eq!(buffer, "{\"b\":2}\n");

        let second = take_next_value::<serde_json::Value>(&mut buffer)
            .expect("second value should parse")
            .expect("second value should be ready");
        assert_eq!(second["b"], 2);
        assert!(buffer.is_empty());
    }

    #[test]
    fn take_next_value_waits_for_partial_json() {
        let mut buffer = "{\"a\":".to_string();
        let value = take_next_value::<serde_json::Value>(&mut buffer)
            .expect("partial JSON should not hard-fail");
        assert!(value.is_none());
        assert_eq!(buffer, "{\"a\":");
    }
}