Skip to main content

taquba_workflow/bulk/
io.rs

1//! Input and output adapters for bulk runs.
2//!
3//! Line-delimited JSON (JSONL): one input item per line in, one result
4//! record per line out. [`read_jsonl`] decodes a reader into typed
5//! items; [`OutputSink`] is the write side, with [`JsonlSink`] as the
6//! built-in implementation. Other sources (CSV, S3 prefixes) and sinks can
7//! be added by implementing the traits without touching the runner.
8
9use std::io::{BufRead, Write};
10use std::sync::Mutex;
11
12use serde::de::DeserializeOwned;
13use serde_json::{Map, Value};
14
15use crate::{Error, Result};
16
17/// Decode a JSONL reader into an iterator of typed items. Each non-empty
18/// line is parsed as one `T`; blank lines are skipped. Decode errors are
19/// yielded as `Err` so the caller decides whether to stop or continue.
20pub 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/// One result record handed to an [`OutputSink`] as an item completes.
33#[derive(Debug)]
34pub struct OutputRecord<'a> {
35    /// The completed item's key: the value of
36    /// [`BulkBuilder::key_fn`](crate::bulk::BulkBuilder::key_fn) for its
37    /// input, or the positional `item-{i}` default.
38    pub key: &'a str,
39    /// Terminal status, as the canonical lowercase string
40    /// (`"succeeded"`, `"failed"`, `"cancelled"`).
41    pub status: &'a str,
42    /// The pipeline output, present only for succeeded items.
43    pub output: Option<Value>,
44    /// A failure reason, present for failed or cancelled items when one was
45    /// recorded.
46    pub error: Option<&'a str>,
47}
48
49/// The write side of a bulk run. Implementations receive one
50/// [`OutputRecord`] per item as it reaches a terminal state, possibly from
51/// many tasks concurrently, so `write` takes `&self` and must handle its
52/// own synchronization.
53///
54/// A batch run writes each of its items once. A later run of the same
55/// batch writes its items again, the succeeded ones from their outcome
56/// records, so consumers that must not double-apply a record across runs
57/// deduplicate on `key`.
58pub trait OutputSink: Send + Sync {
59    /// Persist one completed item's record.
60    fn write(&self, record: &OutputRecord<'_>) -> Result<()>;
61
62    /// Flush any buffered output. Called once when the run finishes. The
63    /// default does nothing.
64    fn flush(&self) -> Result<()> {
65        Ok(())
66    }
67}
68
69/// An [`OutputSink`] that writes one JSON object per line to an underlying
70/// writer. Each line holds `key`, `status` and either `output` (for
71/// succeeded items) or `error` (when one is present). Writes are serialized
72/// through a mutex so the sink can be shared across worker tasks.
73pub struct JsonlSink<W: Write> {
74    writer: Mutex<W>,
75}
76
77impl<W: Write> JsonlSink<W> {
78    /// Wrap a writer as a JSONL sink. Pass a buffered writer (e.g.
79    /// [`std::io::BufWriter`]) for file or socket targets.
80    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/// An [`OutputSink`] that discards every record. The default sink, for runs
112/// whose pipeline produces its results as side effects (writing to a
113/// database, calling an API) rather than through the output stream.
114#[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        // Recover the buffer by writing into a fresh sink and reading lines.
185        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}