use tokio::io::{AsyncBufReadExt, BufReader};
use tokio::process::ChildStdout;
use crate::AgentError;
const MAX_BUFFER_SIZE: usize = 1024 * 1024;
pub struct JsonLineReader {
reader: BufReader<ChildStdout>,
buffer: String,
}
impl JsonLineReader {
pub fn new(stdout: ChildStdout) -> Self {
Self {
reader: BufReader::new(stdout),
buffer: String::new(),
}
}
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() {
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\":");
}
}