chimera_core/transport/
json_reader.rs1use tokio::io::{AsyncBufReadExt, BufReader};
2use tokio::process::ChildStdout;
3
4use crate::AgentError;
5
6const MAX_BUFFER_SIZE: usize = 1024 * 1024; pub 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 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 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}