use std::io::{BufRead, Write};
use std::sync::Mutex;
use serde::de::DeserializeOwned;
use serde_json::{Map, Value};
use crate::{Error, Result};
pub fn read_jsonl<T, R>(reader: R) -> impl Iterator<Item = Result<T>>
where
T: DeserializeOwned,
R: BufRead,
{
reader.lines().filter_map(|line| match line {
Ok(line) if line.trim().is_empty() => None,
Ok(line) => Some(serde_json::from_str::<T>(&line).map_err(Error::from)),
Err(err) => Some(Err(Error::from(err))),
})
}
#[derive(Debug)]
pub struct OutputRecord<'a> {
pub key: &'a str,
pub status: &'a str,
pub output: Option<Value>,
pub error: Option<&'a str>,
}
pub trait OutputSink: Send + Sync {
fn write(&self, record: &OutputRecord<'_>) -> Result<()>;
fn flush(&self) -> Result<()> {
Ok(())
}
}
pub struct JsonlSink<W: Write> {
writer: Mutex<W>,
}
impl<W: Write> JsonlSink<W> {
pub fn new(writer: W) -> Self {
Self {
writer: Mutex::new(writer),
}
}
}
impl<W: Write + Send> OutputSink for JsonlSink<W> {
fn write(&self, record: &OutputRecord<'_>) -> Result<()> {
let mut obj = Map::new();
obj.insert("key".into(), Value::String(record.key.to_string()));
obj.insert("status".into(), Value::String(record.status.to_string()));
if let Some(output) = &record.output {
obj.insert("output".into(), output.clone());
}
if let Some(error) = record.error {
obj.insert("error".into(), Value::String(error.to_string()));
}
let line = serde_json::to_string(&Value::Object(obj))?;
let mut writer = self.writer.lock().unwrap();
writer.write_all(line.as_bytes())?;
writer.write_all(b"\n")?;
Ok(())
}
fn flush(&self) -> Result<()> {
self.writer.lock().unwrap().flush()?;
Ok(())
}
}
#[derive(Debug, Default, Clone, Copy)]
pub struct NullSink;
impl OutputSink for NullSink {
fn write(&self, _record: &OutputRecord<'_>) -> Result<()> {
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde::{Deserialize, Serialize};
#[derive(Debug, PartialEq, Serialize, Deserialize)]
struct Item {
id: u32,
name: String,
}
#[test]
fn read_jsonl_decodes_lines_and_skips_blanks() {
let input = "{\"id\":1,\"name\":\"a\"}\n\n{\"id\":2,\"name\":\"b\"}\n";
let items: Vec<Item> = read_jsonl(input.as_bytes())
.collect::<Result<Vec<_>>>()
.unwrap();
assert_eq!(
items,
vec![
Item {
id: 1,
name: "a".into()
},
Item {
id: 2,
name: "b".into()
},
],
);
}
#[test]
fn read_jsonl_yields_error_for_bad_line() {
let input = "{\"id\":1,\"name\":\"a\"}\nnot json\n";
let results: Vec<Result<Item>> = read_jsonl(input.as_bytes()).collect();
assert_eq!(results.len(), 2);
assert!(results[0].is_ok());
assert!(results[1].is_err());
}
#[test]
fn jsonl_sink_writes_one_object_per_line() {
let buf: Vec<u8> = Vec::new();
let sink = JsonlSink::new(buf);
sink.write(&OutputRecord {
key: "item-0",
status: "succeeded",
output: Some(serde_json::json!({"n": 42})),
error: None,
})
.unwrap();
sink.write(&OutputRecord {
key: "item-1",
status: "failed",
output: None,
error: Some("boom"),
})
.unwrap();
sink.flush().unwrap();
let bytes = sink.writer.into_inner().unwrap();
let text = String::from_utf8(bytes).unwrap();
let lines: Vec<&str> = text.lines().collect();
assert_eq!(lines.len(), 2);
let first: Value = serde_json::from_str(lines[0]).unwrap();
assert_eq!(first["key"], "item-0");
assert_eq!(first["status"], "succeeded");
assert_eq!(first["output"]["n"], 42);
assert!(first.get("error").is_none());
let second: Value = serde_json::from_str(lines[1]).unwrap();
assert_eq!(second["status"], "failed");
assert_eq!(second["error"], "boom");
assert!(second.get("output").is_none());
}
}