taquba_workflow/bulk/
io.rs1use std::io::{BufRead, Write};
10use std::sync::Mutex;
11
12use serde::de::DeserializeOwned;
13use serde_json::{Map, Value};
14
15use crate::{Error, Result};
16
17pub fn read_jsonl<T, R>(reader: R) -> impl Iterator<Item = Result<T>>
21where
22 T: DeserializeOwned,
23 R: BufRead,
24{
25 reader.lines().filter_map(|line| match line {
26 Ok(line) if line.trim().is_empty() => None,
27 Ok(line) => Some(serde_json::from_str::<T>(&line).map_err(Error::from)),
28 Err(err) => Some(Err(Error::from(err))),
29 })
30}
31
32#[derive(Debug)]
34pub struct OutputRecord<'a> {
35 pub key: &'a str,
39 pub status: &'a str,
42 pub output: Option<Value>,
44 pub error: Option<&'a str>,
47}
48
49pub trait OutputSink: Send + Sync {
59 fn write(&self, record: &OutputRecord<'_>) -> Result<()>;
61
62 fn flush(&self) -> Result<()> {
65 Ok(())
66 }
67}
68
69pub struct JsonlSink<W: Write> {
74 writer: Mutex<W>,
75}
76
77impl<W: Write> JsonlSink<W> {
78 pub fn new(writer: W) -> Self {
81 Self {
82 writer: Mutex::new(writer),
83 }
84 }
85}
86
87impl<W: Write + Send> OutputSink for JsonlSink<W> {
88 fn write(&self, record: &OutputRecord<'_>) -> Result<()> {
89 let mut obj = Map::new();
90 obj.insert("key".into(), Value::String(record.key.to_string()));
91 obj.insert("status".into(), Value::String(record.status.to_string()));
92 if let Some(output) = &record.output {
93 obj.insert("output".into(), output.clone());
94 }
95 if let Some(error) = record.error {
96 obj.insert("error".into(), Value::String(error.to_string()));
97 }
98 let line = serde_json::to_string(&Value::Object(obj))?;
99 let mut writer = self.writer.lock().unwrap();
100 writer.write_all(line.as_bytes())?;
101 writer.write_all(b"\n")?;
102 Ok(())
103 }
104
105 fn flush(&self) -> Result<()> {
106 self.writer.lock().unwrap().flush()?;
107 Ok(())
108 }
109}
110
111#[derive(Debug, Default, Clone, Copy)]
115pub struct NullSink;
116
117impl OutputSink for NullSink {
118 fn write(&self, _record: &OutputRecord<'_>) -> Result<()> {
119 Ok(())
120 }
121}
122
123#[cfg(test)]
124mod tests {
125 use super::*;
126 use serde::{Deserialize, Serialize};
127
128 #[derive(Debug, PartialEq, Serialize, Deserialize)]
129 struct Item {
130 id: u32,
131 name: String,
132 }
133
134 #[test]
135 fn read_jsonl_decodes_lines_and_skips_blanks() {
136 let input = "{\"id\":1,\"name\":\"a\"}\n\n{\"id\":2,\"name\":\"b\"}\n";
137 let items: Vec<Item> = read_jsonl(input.as_bytes())
138 .collect::<Result<Vec<_>>>()
139 .unwrap();
140 assert_eq!(
141 items,
142 vec![
143 Item {
144 id: 1,
145 name: "a".into()
146 },
147 Item {
148 id: 2,
149 name: "b".into()
150 },
151 ],
152 );
153 }
154
155 #[test]
156 fn read_jsonl_yields_error_for_bad_line() {
157 let input = "{\"id\":1,\"name\":\"a\"}\nnot json\n";
158 let results: Vec<Result<Item>> = read_jsonl(input.as_bytes()).collect();
159 assert_eq!(results.len(), 2);
160 assert!(results[0].is_ok());
161 assert!(results[1].is_err());
162 }
163
164 #[test]
165 fn jsonl_sink_writes_one_object_per_line() {
166 let buf: Vec<u8> = Vec::new();
167 let sink = JsonlSink::new(buf);
168 sink.write(&OutputRecord {
169 key: "item-0",
170 status: "succeeded",
171 output: Some(serde_json::json!({"n": 42})),
172 error: None,
173 })
174 .unwrap();
175 sink.write(&OutputRecord {
176 key: "item-1",
177 status: "failed",
178 output: None,
179 error: Some("boom"),
180 })
181 .unwrap();
182 sink.flush().unwrap();
183
184 let bytes = sink.writer.into_inner().unwrap();
186 let text = String::from_utf8(bytes).unwrap();
187 let lines: Vec<&str> = text.lines().collect();
188 assert_eq!(lines.len(), 2);
189
190 let first: Value = serde_json::from_str(lines[0]).unwrap();
191 assert_eq!(first["key"], "item-0");
192 assert_eq!(first["status"], "succeeded");
193 assert_eq!(first["output"]["n"], 42);
194 assert!(first.get("error").is_none());
195
196 let second: Value = serde_json::from_str(lines[1]).unwrap();
197 assert_eq!(second["status"], "failed");
198 assert_eq!(second["error"], "boom");
199 assert!(second.get("output").is_none());
200 }
201}