Skip to main content

chimera_core/transport/
json_reader.rs

1use tokio::io::{AsyncBufReadExt, BufReader};
2use tokio::process::ChildStdout;
3
4use crate::AgentError;
5
6const MAX_BUFFER_SIZE: usize = 1024 * 1024; // 1 MB
7
8pub struct JsonLineReader {
9    reader: BufReader<ChildStdout>,
10    buffer: String,
11}
12
13impl JsonLineReader {
14    pub fn new(stdout: ChildStdout) -> Self {
15        Self {
16            reader: BufReader::new(stdout),
17            buffer: String::new(),
18        }
19    }
20
21    /// Read the next complete JSON value from stdout.
22    ///
23    /// Handles JSON objects spanning multiple lines by buffering and
24    /// retrying the parse. Returns `None` on EOF.
25    pub async fn next_value<T: serde::de::DeserializeOwned>(
26        &mut self,
27    ) -> Option<Result<T, AgentError>> {
28        loop {
29            let mut line = String::new();
30            match self.reader.read_line(&mut line).await {
31                Ok(0) => {
32                    if self.buffer.is_empty() {
33                        return None;
34                    }
35                    let buf = std::mem::take(&mut self.buffer);
36                    return Some(
37                        serde_json::from_str(&buf)
38                            .map_err(|source| AgentError::JsonParse { line: buf, source }),
39                    );
40                }
41                Ok(_) => {
42                    self.buffer.push_str(&line);
43                    if self.buffer.len() > MAX_BUFFER_SIZE {
44                        self.buffer.clear();
45                        return Some(Err(AgentError::BufferOverflow {
46                            limit: MAX_BUFFER_SIZE,
47                        }));
48                    }
49                    match take_next_value(&mut self.buffer) {
50                        Ok(Some(value)) => return Some(Ok(value)),
51                        Ok(None) => continue,
52                        Err(source) => {
53                            let line = std::mem::take(&mut self.buffer);
54                            return Some(Err(AgentError::JsonParse { line, source }));
55                        }
56                    }
57                }
58                Err(e) => {
59                    return Some(Err(AgentError::Other {
60                        message: "failed to read stdout".into(),
61                        source: Some(Box::new(e)),
62                    }));
63                }
64            }
65        }
66    }
67}
68
69fn take_next_value<T: serde::de::DeserializeOwned>(
70    buffer: &mut String,
71) -> Result<Option<T>, serde_json::Error> {
72    let trimmed = buffer.trim_start_matches(['\r', '\n']);
73    if trimmed.len() != buffer.len() {
74        *buffer = trimmed.to_string();
75    }
76    if buffer.is_empty() {
77        return Ok(None);
78    }
79
80    let mut stream = serde_json::Deserializer::from_str(buffer).into_iter::<T>();
81    match stream.next() {
82        Some(Ok(value)) => {
83            let consumed = stream.byte_offset();
84            let remainder = buffer[consumed..].trim_start_matches(['\r', '\n']);
85            *buffer = remainder.to_string();
86            Ok(Some(value))
87        }
88        Some(Err(err)) if err.is_eof() => Ok(None),
89        Some(Err(err)) => Err(err),
90        None => Ok(None),
91    }
92}
93
94#[cfg(test)]
95mod tests {
96    use super::*;
97
98    #[tokio::test]
99    async fn json_line_reader_single_line() {
100        // JsonLineReader requires ChildStdout, so we test the parsing
101        // logic it relies on directly.
102        let json = r#"{"key": "value"}"#;
103        let parsed: serde_json::Value = serde_json::from_str(json).unwrap();
104        assert_eq!(parsed["key"], "value");
105    }
106
107    #[test]
108    fn buffer_overflow_threshold() {
109        assert_eq!(MAX_BUFFER_SIZE, 1024 * 1024);
110    }
111
112    #[test]
113    fn take_next_value_consumes_one_value_and_keeps_remainder() {
114        let mut buffer = "{\"a\":1}\n{\"b\":2}\n".to_string();
115        let first = take_next_value::<serde_json::Value>(&mut buffer)
116            .expect("first value should parse")
117            .expect("first value should be ready");
118        assert_eq!(first["a"], 1);
119        assert_eq!(buffer, "{\"b\":2}\n");
120
121        let second = take_next_value::<serde_json::Value>(&mut buffer)
122            .expect("second value should parse")
123            .expect("second value should be ready");
124        assert_eq!(second["b"], 2);
125        assert!(buffer.is_empty());
126    }
127
128    #[test]
129    fn take_next_value_waits_for_partial_json() {
130        let mut buffer = "{\"a\":".to_string();
131        let value = take_next_value::<serde_json::Value>(&mut buffer)
132            .expect("partial JSON should not hard-fail");
133        assert!(value.is_none());
134        assert_eq!(buffer, "{\"a\":");
135    }
136}