Skip to main content

tpt_jsonl_stream/
lib.rs

1#![doc = include_str!("../README.md")]
2#![warn(missing_docs)]
3
4use std::fmt;
5use std::io::{self, BufRead, Write};
6
7/// The kind of error that occurred while reading a JSON Lines stream.
8#[derive(Debug)]
9pub enum JsonlErrorKind {
10    /// An I/O error from the underlying reader.
11    Io(io::Error),
12    /// A JSON parse error on a specific line.
13    Json(serde_json::Error),
14}
15
16impl fmt::Display for JsonlErrorKind {
17    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
18        match self {
19            Self::Io(e) => write!(f, "I/O error: {}", e),
20            Self::Json(e) => write!(f, "JSON error: {}", e),
21        }
22    }
23}
24
25/// An error produced while reading or parsing a JSON Lines stream.
26///
27/// Includes the 1-based line number where the error occurred.
28#[derive(Debug)]
29pub struct JsonlError {
30    /// The 1-based line number where the error occurred.
31    pub line: u64,
32    /// The underlying error kind.
33    pub kind: JsonlErrorKind,
34}
35
36impl fmt::Display for JsonlError {
37    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
38        write!(f, "jsonl error on line {}: {}", self.line, self.kind)
39    }
40}
41
42impl std::error::Error for JsonlError {
43    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
44        match &self.kind {
45            JsonlErrorKind::Io(e) => Some(e),
46            JsonlErrorKind::Json(e) => Some(e),
47        }
48    }
49}
50
51/// A streaming JSON Lines reader.
52///
53/// Wraps any [`BufRead`] and yields one [`serde_json::Value`] per non-empty line.
54/// Blank lines are silently skipped. Parse errors carry the line number.
55///
56/// # Example
57///
58/// ```
59/// use tpt_jsonl_stream::JsonlReader;
60/// use std::io::BufReader;
61///
62/// let data = b"{\"a\":1}\n{\"b\":2}\n";
63/// let mut reader = JsonlReader::new(BufReader::new(data.as_slice()));
64/// let first = reader.next().unwrap().unwrap();
65/// assert_eq!(first["a"], 1);
66/// ```
67pub struct JsonlReader<R: BufRead> {
68    reader: R,
69    buf: String,
70    line: u64,
71}
72
73impl<R: BufRead> JsonlReader<R> {
74    /// Create a new `JsonlReader` wrapping the given buffered reader.
75    pub fn new(reader: R) -> Self {
76        Self {
77            reader,
78            buf: String::new(),
79            line: 0,
80        }
81    }
82
83    /// The 1-based line number most recently read (or 0 before any reads).
84    pub fn line_number(&self) -> u64 {
85        self.line
86    }
87}
88
89impl<R: BufRead> Iterator for JsonlReader<R> {
90    type Item = Result<serde_json::Value, JsonlError>;
91
92    fn next(&mut self) -> Option<Self::Item> {
93        loop {
94            self.buf.clear();
95            match self.reader.read_line(&mut self.buf) {
96                Err(e) => {
97                    self.line += 1;
98                    return Some(Err(JsonlError {
99                        line: self.line,
100                        kind: JsonlErrorKind::Io(e),
101                    }));
102                }
103                Ok(0) => return None, // EOF
104                Ok(_) => {
105                    self.line += 1;
106                    let trimmed = self.buf.trim();
107                    if trimmed.is_empty() {
108                        continue; // skip blank lines
109                    }
110                    #[cfg(feature = "simd")]
111                    {
112                        let mut bytes = trimmed.as_bytes().to_vec();
113                        match simd_json::from_slice(&mut bytes) {
114                            Ok(v) => return Some(Ok(v)),
115                            Err(_) => {
116                                // simd-json failed. Re-derive an error and value from
117                                // serde_json so we never panic on parser divergence:
118                                // if serde_json also rejects the line we surface its
119                                // precise error; if it accepts the line (divergence)
120                                // we yield the parsed value rather than crashing.
121                                match serde_json::from_str::<serde_json::Value>(trimmed) {
122                                    Ok(v) => return Some(Ok(v)),
123                                    Err(e) => {
124                                        return Some(Err(JsonlError {
125                                            line: self.line,
126                                            kind: JsonlErrorKind::Json(e),
127                                        }));
128                                    }
129                                }
130                            }
131                        }
132                    }
133                    #[cfg(not(feature = "simd"))]
134                    match serde_json::from_str::<serde_json::Value>(trimmed) {
135                        Ok(v) => return Some(Ok(v)),
136                        Err(e) => {
137                            return Some(Err(JsonlError {
138                                line: self.line,
139                                kind: JsonlErrorKind::Json(e),
140                            }))
141                        }
142                    }
143                }
144            }
145        }
146    }
147}
148
149/// Create a [`JsonlReader`] from any [`BufRead`].
150///
151/// Convenience wrapper around [`JsonlReader::new`].
152///
153/// # Example
154///
155/// ```
156/// use tpt_jsonl_stream::parse_jsonl;
157/// use std::io::BufReader;
158///
159/// let data = b"{\"x\":1}\n\n{\"x\":2}\n";
160/// let values: Vec<_> = parse_jsonl(BufReader::new(data.as_slice()))
161///     .collect::<Result<_, _>>()
162///     .unwrap();
163/// assert_eq!(values.len(), 2);
164/// ```
165pub fn parse_jsonl<R: BufRead>(reader: R) -> JsonlReader<R> {
166    JsonlReader::new(reader)
167}
168
169/// A streaming JSON Lines writer.
170///
171/// Wraps any [`Write`] and emits one JSON value per line. Each call to
172/// [`JsonlWriter::write`] serializes the value with `serde_json` and appends a
173/// trailing newline. Parse errors carry the 1-based line number of the write
174/// that failed.
175///
176/// # Example
177///
178/// ```
179/// use tpt_jsonl_stream::JsonlWriter;
180/// use std::io::Cursor;
181///
182/// let mut buf = Cursor::new(Vec::new());
183/// {
184///     let mut writer = JsonlWriter::new(&mut buf);
185///     writer.write(&serde_json::json!({"a": 1})).unwrap();
186///     writer.write(&serde_json::json!({"b": 2})).unwrap();
187/// }
188/// let out = String::from_utf8(buf.into_inner()).unwrap();
189/// assert_eq!(out, "{\"a\":1}\n{\"b\":2}\n");
190/// ```
191pub struct JsonlWriter<W: Write> {
192    writer: W,
193    line: u64,
194}
195
196impl<W: Write> JsonlWriter<W> {
197    /// Create a new `JsonlWriter` wrapping the given writer.
198    pub fn new(writer: W) -> Self {
199        Self { writer, line: 0 }
200    }
201
202    /// The number of lines (values) written so far.
203    pub fn line_number(&self) -> u64 {
204        self.line
205    }
206
207    /// Serialize `value` as a single JSON Lines record (one line, newline-terminated).
208    pub fn write<T: serde::Serialize>(&mut self, value: &T) -> Result<(), JsonlError> {
209        serde_json::to_writer(&mut self.writer, value).map_err(|e| JsonlError {
210            line: self.line + 1,
211            kind: JsonlErrorKind::Json(e),
212        })?;
213        self.writer.write_all(b"\n").map_err(|e| JsonlError {
214            line: self.line + 1,
215            kind: JsonlErrorKind::Io(e),
216        })?;
217        self.line += 1;
218        Ok(())
219    }
220
221    /// Flush the underlying writer.
222    pub fn flush(&mut self) -> io::Result<()> {
223        self.writer.flush()
224    }
225}
226
227/// Write a sequence of values as JSON Lines to the given [`Write`] sink.
228///
229/// # Example
230///
231/// ```
232/// use tpt_jsonl_stream::write_jsonl;
233/// use std::io::Cursor;
234///
235/// let mut buf = Cursor::new(Vec::new());
236/// let values = vec![serde_json::json!(1), serde_json::json!(2)];
237/// write_jsonl(&mut buf, values.iter()).unwrap();
238/// let out = String::from_utf8(buf.into_inner()).unwrap();
239/// assert_eq!(out, "1\n2\n");
240/// ```
241pub fn write_jsonl<W: Write, I, T>(writer: W, values: I) -> Result<(), JsonlError>
242where
243    I: IntoIterator<Item = T>,
244    T: serde::Serialize,
245{
246    let mut w = JsonlWriter::new(writer);
247    for v in values {
248        w.write(&v)?;
249    }
250    w.flush().map_err(|e| JsonlError {
251        line: w.line_number() + 1,
252        kind: JsonlErrorKind::Io(e),
253    })?;
254    Ok(())
255}
256
257#[cfg(test)]
258mod tests {
259    use super::*;
260    use std::io::BufReader;
261
262    fn read_all(data: &[u8]) -> Vec<serde_json::Value> {
263        parse_jsonl(BufReader::new(data))
264            .collect::<Result<_, _>>()
265            .unwrap()
266    }
267
268    #[test]
269    fn empty_input() {
270        assert!(read_all(b"").is_empty());
271    }
272
273    #[test]
274    fn single_line() {
275        let vals = read_all(b"{\"k\":1}\n");
276        assert_eq!(vals.len(), 1);
277        assert_eq!(vals[0]["k"], 1);
278    }
279
280    #[test]
281    fn multi_line() {
282        let vals = read_all(b"{\"a\":1}\n{\"b\":2}\n{\"c\":3}\n");
283        assert_eq!(vals.len(), 3);
284    }
285
286    #[test]
287    fn blank_lines_skipped() {
288        let vals = read_all(b"{\"a\":1}\n\n\n{\"b\":2}\n");
289        assert_eq!(vals.len(), 2);
290    }
291
292    #[test]
293    fn malformed_json_error_has_correct_line() {
294        let data = b"{\"a\":1}\nNOT_JSON\n{\"c\":3}\n";
295        let mut reader = parse_jsonl(BufReader::new(data.as_slice()));
296        reader.next().unwrap().unwrap(); // line 1 ok
297        let err = reader.next().unwrap().unwrap_err();
298        assert_eq!(err.line, 2);
299    }
300
301    #[test]
302    fn line_counter_exposed() {
303        let data = b"{\"a\":1}\n{\"b\":2}\n";
304        let mut reader = parse_jsonl(BufReader::new(data.as_slice()));
305        assert_eq!(reader.line_number(), 0);
306        reader.next();
307        assert_eq!(reader.line_number(), 1);
308    }
309
310    #[test]
311    fn no_trailing_newline() {
312        let vals = read_all(b"{\"x\":42}");
313        assert_eq!(vals.len(), 1);
314        assert_eq!(vals[0]["x"], 42);
315    }
316
317    #[test]
318    fn writer_round_trips_with_reader() {
319        let mut buf: Vec<u8> = Vec::new();
320        {
321            let mut writer = JsonlWriter::new(&mut buf);
322            writer.write(&serde_json::json!({"a": 1})).unwrap();
323            writer.write(&serde_json::json!({"b": "two"})).unwrap();
324            writer.write(&serde_json::json!([1, 2, 3])).unwrap();
325            writer.flush().unwrap();
326        }
327        let written = String::from_utf8(buf.clone()).unwrap();
328        assert_eq!(written, "{\"a\":1}\n{\"b\":\"two\"}\n[1,2,3]\n");
329
330        // The reader should recover the exact same values.
331        let back: Vec<serde_json::Value> = parse_jsonl(BufReader::new(buf.as_slice()))
332            .collect::<Result<_, _>>()
333            .unwrap();
334        assert_eq!(back.len(), 3);
335        assert_eq!(back[0]["a"], 1);
336        assert_eq!(back[1]["b"], "two");
337        assert_eq!(back[2], serde_json::json!([1, 2, 3]));
338    }
339
340    #[test]
341    fn write_jsonl_helper() {
342        let mut buf: Vec<u8> = Vec::new();
343        let values = [serde_json::json!(1), serde_json::json!(2)];
344        write_jsonl(&mut buf, values.iter()).unwrap();
345        assert_eq!(String::from_utf8(buf).unwrap(), "1\n2\n");
346    }
347
348    #[test]
349    fn writer_tracks_line_numbers() {
350        let mut buf: Vec<u8> = Vec::new();
351        let mut writer = JsonlWriter::new(&mut buf);
352        writer.write(&serde_json::json!({"ok": true})).unwrap();
353        assert_eq!(writer.line_number(), 1);
354    }
355}