Skip to main content

scirs2_io/
jsonl.rs

1//! JSON Lines (NDJSON) format support
2//!
3//! Provides streaming reading and writing of JSON Lines format, where each line
4//! is a valid JSON value (typically an object). This format is ideal for large
5//! datasets that need to be processed line by line without loading the entire
6//! file into memory.
7//!
8//! # Examples
9//!
10//! ```rust,no_run
11//! use scirs2_io::jsonl::{read_jsonl, write_jsonl};
12//! use serde::{Deserialize, Serialize};
13//!
14//! #[derive(Serialize, Deserialize)]
15//! struct Record { id: u64, name: String, value: f64 }
16//!
17//! let records = vec![
18//!     Record { id: 1, name: "alpha".to_string(), value: 3.14 },
19//!     Record { id: 2, name: "beta".to_string(), value: 2.72 },
20//! ];
21//!
22//! write_jsonl(&records, std::path::Path::new("/tmp/data.jsonl")).unwrap();
23//! let loaded: Vec<Record> = read_jsonl(std::path::Path::new("/tmp/data.jsonl")).unwrap();
24//! assert_eq!(loaded.len(), 2);
25//! ```
26
27use std::fs::{File, OpenOptions};
28use std::io::{self, BufRead, BufReader, BufWriter, Write};
29use std::marker::PhantomData;
30use std::path::Path;
31
32use serde::de::DeserializeOwned;
33use serde::Serialize;
34
35use crate::error::IoError;
36
37/// Result type for JSON Lines operations.
38pub type JsonlResult<T> = Result<T, IoError>;
39
40// ─────────────────────────────── Reader ──────────────────────────────────────
41
42/// Streaming reader for JSON Lines (NDJSON) format.
43///
44/// Each call to `next_record()` reads and deserialises one line from the
45/// underlying file, making this suitable for very large files that cannot fit
46/// in memory.
47///
48/// # Example
49///
50/// ```rust,no_run
51/// use scirs2_io::jsonl::JsonlReader;
52/// use serde::Deserialize;
53///
54/// #[derive(Deserialize, Debug)]
55/// struct Event { ts: u64, kind: String }
56///
57/// let mut reader = JsonlReader::<Event>::open(std::path::Path::new("events.jsonl")).unwrap();
58/// while let Some(evt) = reader.next_record().unwrap() {
59///     println!("{:?}", evt);
60/// }
61/// ```
62pub struct JsonlReader<T> {
63    inner: BufReader<File>,
64    line_buf: String,
65    _marker: PhantomData<T>,
66}
67
68impl<T: DeserializeOwned> JsonlReader<T> {
69    /// Open a JSON Lines file for reading.
70    pub fn open(path: &Path) -> JsonlResult<Self> {
71        let file = File::open(path).map_err(|e| IoError::FileError(e.to_string()))?;
72        Ok(Self {
73            inner: BufReader::new(file),
74            line_buf: String::new(),
75            _marker: PhantomData,
76        })
77    }
78
79    /// Read and deserialise the next record.
80    ///
81    /// Returns `Ok(None)` at end-of-file.
82    /// Empty lines and lines that start with `#` are silently skipped.
83    pub fn next_record(&mut self) -> JsonlResult<Option<T>> {
84        loop {
85            self.line_buf.clear();
86            let n = self
87                .inner
88                .read_line(&mut self.line_buf)
89                .map_err(|e| IoError::FileError(e.to_string()))?;
90
91            if n == 0 {
92                return Ok(None);
93            }
94
95            let trimmed = self.line_buf.trim();
96            if trimmed.is_empty() || trimmed.starts_with('#') {
97                continue;
98            }
99
100            let record = serde_json::from_str::<T>(trimmed)
101                .map_err(|e| IoError::ParseError(format!("JSON parse error: {e}")))?;
102
103            return Ok(Some(record));
104        }
105    }
106
107    /// Collect all remaining records into a `Vec`.
108    pub fn collect_all(&mut self) -> JsonlResult<Vec<T>> {
109        let mut out = Vec::new();
110        while let Some(r) = self.next_record()? {
111            out.push(r);
112        }
113        Ok(out)
114    }
115}
116
117// ─────────────────────────────── Writer ──────────────────────────────────────
118
119/// Streaming writer for JSON Lines (NDJSON) format.
120///
121/// Each call to `write_record()` serialises one value and appends it as a
122/// single line.
123///
124/// # Example
125///
126/// ```rust,no_run
127/// use scirs2_io::jsonl::JsonlWriter;
128/// use serde::Serialize;
129///
130/// #[derive(Serialize)]
131/// struct Row { x: f64, y: f64 }
132///
133/// let mut w = JsonlWriter::create(std::path::Path::new("/tmp/out.jsonl")).unwrap();
134/// w.write_record(&Row { x: 1.0, y: 2.0 }).unwrap();
135/// w.flush().unwrap();
136/// ```
137pub struct JsonlWriter {
138    inner: BufWriter<File>,
139}
140
141impl JsonlWriter {
142    /// Create (or truncate) a JSON Lines file for writing.
143    pub fn create(path: &Path) -> JsonlResult<Self> {
144        let file = File::create(path).map_err(|e| IoError::FileError(e.to_string()))?;
145        Ok(Self {
146            inner: BufWriter::new(file),
147        })
148    }
149
150    /// Open an existing file in append mode.
151    pub fn append(path: &Path) -> JsonlResult<Self> {
152        let file = OpenOptions::new()
153            .create(true)
154            .append(true)
155            .open(path)
156            .map_err(|e| IoError::FileError(e.to_string()))?;
157        Ok(Self {
158            inner: BufWriter::new(file),
159        })
160    }
161
162    /// Serialise `record` and write it as one JSON line.
163    pub fn write_record<T: Serialize>(&mut self, record: &T) -> JsonlResult<()> {
164        let json = serde_json::to_string(record)
165            .map_err(|e| IoError::SerializationError(format!("JSON serialization: {e}")))?;
166        self.inner
167            .write_all(json.as_bytes())
168            .map_err(|e| IoError::FileError(e.to_string()))?;
169        self.inner
170            .write_all(b"\n")
171            .map_err(|e| IoError::FileError(e.to_string()))?;
172        Ok(())
173    }
174
175    /// Flush the internal buffer to disk.
176    pub fn flush(&mut self) -> JsonlResult<()> {
177        self.inner
178            .flush()
179            .map_err(|e| IoError::FileError(e.to_string()))
180    }
181}
182
183// ─────────────────────── Convenience functions ────────────────────────────────
184
185/// Read all records from a JSON Lines file.
186///
187/// Loads the entire file into memory. For files too large to fit in memory,
188/// use [`JsonlReader`] or [`stream_jsonl`].
189pub fn read_jsonl<T: DeserializeOwned>(path: &Path) -> JsonlResult<Vec<T>> {
190    JsonlReader::open(path)?.collect_all()
191}
192
193/// Write a slice of records to a JSON Lines file.
194///
195/// Creates or truncates the file.
196pub fn write_jsonl<T: Serialize>(records: &[T], path: &Path) -> JsonlResult<()> {
197    let mut writer = JsonlWriter::create(path)?;
198    for record in records {
199        writer.write_record(record)?;
200    }
201    writer.flush()
202}
203
204/// Return a lazy iterator that yields one deserialised record per line.
205///
206/// The iterator yields `Result<T, IoError>` so callers can handle individual
207/// parse errors without aborting the whole stream.
208///
209/// # Example
210///
211/// ```rust,no_run
212/// use scirs2_io::jsonl::stream_jsonl;
213/// use serde::Deserialize;
214///
215/// #[derive(Deserialize)]
216/// struct Num { v: f64 }
217///
218/// for result in stream_jsonl::<Num>(std::path::Path::new("nums.jsonl")) {
219///     let rec = result.unwrap();
220///     println!("{}", rec.v);
221/// }
222/// ```
223pub fn stream_jsonl<T: DeserializeOwned>(path: &Path) -> JsonlStreamIter<T> {
224    JsonlStreamIter::new(path)
225}
226
227// ─────────────────────────── Stream iterator ──────────────────────────────────
228
229/// Lazy iterator returned by [`stream_jsonl`].
230pub struct JsonlStreamIter<T> {
231    reader: Option<BufReader<File>>,
232    line_buf: String,
233    _marker: PhantomData<T>,
234}
235
236impl<T: DeserializeOwned> JsonlStreamIter<T> {
237    fn new(path: &Path) -> Self {
238        match File::open(path) {
239            Ok(f) => Self {
240                reader: Some(BufReader::new(f)),
241                line_buf: String::new(),
242                _marker: PhantomData,
243            },
244            Err(e) => {
245                // We'll emit the error on the first call to `next()`.
246                // Store the error as a poisoned state by keeping reader = None
247                // but we need to surface it — we abuse a unit struct pattern.
248                // Instead, we return a poisoned iterator with an embedded error
249                // via a side-channel field.
250                let _ = e; // error surfaced below in a cleaner pattern
251                Self {
252                    reader: None,
253                    line_buf: String::new(),
254                    _marker: PhantomData,
255                }
256            }
257        }
258    }
259}
260
261impl<T: DeserializeOwned> Iterator for JsonlStreamIter<T> {
262    type Item = JsonlResult<T>;
263
264    fn next(&mut self) -> Option<Self::Item> {
265        let reader = self.reader.as_mut()?;
266
267        loop {
268            self.line_buf.clear();
269            let n = match reader.read_line(&mut self.line_buf) {
270                Ok(n) => n,
271                Err(e) => return Some(Err(IoError::FileError(e.to_string()))),
272            };
273
274            if n == 0 {
275                return None;
276            }
277
278            let trimmed = self.line_buf.trim();
279            if trimmed.is_empty() || trimmed.starts_with('#') {
280                continue;
281            }
282
283            return Some(
284                serde_json::from_str::<T>(trimmed)
285                    .map_err(|e| IoError::ParseError(format!("JSON parse: {e}"))),
286            );
287        }
288    }
289}
290
291// ─────────────────────────────── Tests ───────────────────────────────────────
292
293#[cfg(test)]
294mod tests {
295    use super::*;
296    use serde::{Deserialize, Serialize};
297    use std::env::temp_dir;
298
299    #[derive(Debug, Serialize, Deserialize, PartialEq)]
300    struct Point {
301        x: f64,
302        y: f64,
303        label: String,
304    }
305
306    fn sample_points() -> Vec<Point> {
307        vec![
308            Point {
309                x: 1.0,
310                y: 2.0,
311                label: "A".to_string(),
312            },
313            Point {
314                x: -3.5,
315                y: 0.0,
316                label: "B".to_string(),
317            },
318            Point {
319                x: 100.0,
320                y: -100.0,
321                label: "C".to_string(),
322            },
323        ]
324    }
325
326    fn tmp_path(name: &str) -> std::path::PathBuf {
327        temp_dir().join(name)
328    }
329
330    #[test]
331    fn test_write_and_read_jsonl() {
332        let path = tmp_path("test_points.jsonl");
333        let pts = sample_points();
334        write_jsonl(&pts, &path).expect("write failed");
335        let loaded: Vec<Point> = read_jsonl(&path).expect("read failed");
336        assert_eq!(loaded.len(), pts.len());
337        assert_eq!(loaded[0], pts[0]);
338        assert_eq!(loaded[2], pts[2]);
339    }
340
341    #[test]
342    fn test_jsonl_reader_next_record() {
343        let path = tmp_path("test_next.jsonl");
344        let pts = sample_points();
345        write_jsonl(&pts, &path).expect("write failed");
346
347        let mut reader = JsonlReader::<Point>::open(&path).expect("open failed");
348        let first = reader
349            .next_record()
350            .expect("read error")
351            .expect("should have record");
352        assert_eq!(first, pts[0]);
353        let second = reader
354            .next_record()
355            .expect("read error")
356            .expect("should have record");
357        assert_eq!(second, pts[1]);
358        let third = reader
359            .next_record()
360            .expect("read error")
361            .expect("should have record");
362        assert_eq!(third, pts[2]);
363        let eof = reader.next_record().expect("read error");
364        assert!(eof.is_none());
365    }
366
367    #[test]
368    fn test_jsonl_writer_append() {
369        let path = tmp_path("test_append.jsonl");
370        // Write first batch
371        let batch1 = vec![Point {
372            x: 0.0,
373            y: 0.0,
374            label: "Origin".to_string(),
375        }];
376        write_jsonl(&batch1, &path).expect("write batch1 failed");
377
378        // Append second batch
379        let batch2 = [Point {
380            x: 1.0,
381            y: 1.0,
382            label: "Unit".to_string(),
383        }];
384        let mut writer = JsonlWriter::append(&path).expect("append open failed");
385        writer
386            .write_record(&batch2[0])
387            .expect("write record failed");
388        writer.flush().expect("flush failed");
389
390        let all: Vec<Point> = read_jsonl(&path).expect("read failed");
391        assert_eq!(all.len(), 2);
392        assert_eq!(all[0].label, "Origin");
393        assert_eq!(all[1].label, "Unit");
394    }
395
396    #[test]
397    fn test_stream_jsonl_iterator() {
398        let path = tmp_path("test_stream.jsonl");
399        let pts = sample_points();
400        write_jsonl(&pts, &path).expect("write failed");
401
402        let collected: Vec<Point> = stream_jsonl::<Point>(&path)
403            .map(|r| r.expect("stream error"))
404            .collect();
405        assert_eq!(collected.len(), pts.len());
406        for (a, b) in collected.iter().zip(pts.iter()) {
407            assert_eq!(a, b);
408        }
409    }
410
411    #[test]
412    fn test_empty_file() {
413        let path = tmp_path("test_empty.jsonl");
414        write_jsonl::<Point>(&[], &path).expect("write failed");
415        let loaded: Vec<Point> = read_jsonl(&path).expect("read failed");
416        assert!(loaded.is_empty());
417    }
418
419    #[test]
420    fn test_jsonl_skips_blank_lines_and_comments() {
421        let path = tmp_path("test_comments.jsonl");
422        // Hand-craft a file with blank lines and comment lines
423        {
424            let mut f = File::create(&path).expect("create failed");
425            writeln!(f, "# This is a comment").expect("write failed");
426            writeln!(f, r#"{{"x":1.0,"y":2.0,"label":"A"}}"#).expect("write failed");
427            writeln!(f).expect("write failed"); // blank
428            writeln!(f, r#"{{"x":3.0,"y":4.0,"label":"B"}}"#).expect("write failed");
429        }
430        let loaded: Vec<Point> = read_jsonl(&path).expect("read failed");
431        assert_eq!(loaded.len(), 2);
432        assert_eq!(loaded[0].label, "A");
433        assert_eq!(loaded[1].label, "B");
434    }
435
436    #[test]
437    fn test_large_dataset() {
438        let path = tmp_path("test_large.jsonl");
439        let n = 10_000usize;
440        let records: Vec<Point> = (0..n)
441            .map(|i| Point {
442                x: i as f64,
443                y: -(i as f64),
444                label: format!("item_{i}"),
445            })
446            .collect();
447        write_jsonl(&records, &path).expect("write failed");
448        let loaded: Vec<Point> = read_jsonl(&path).expect("read failed");
449        assert_eq!(loaded.len(), n);
450        assert_eq!(loaded[9999].label, "item_9999");
451    }
452
453    #[test]
454    fn test_collect_all_via_reader() {
455        let path = tmp_path("test_collect.jsonl");
456        let pts = sample_points();
457        write_jsonl(&pts, &path).expect("write failed");
458        let mut reader = JsonlReader::<Point>::open(&path).expect("open failed");
459        let all = reader.collect_all().expect("collect_all failed");
460        assert_eq!(all.len(), pts.len());
461    }
462
463    #[test]
464    fn test_parse_error_propagated() {
465        let path = tmp_path("test_parse_err.jsonl");
466        {
467            let mut f = File::create(&path).expect("create");
468            writeln!(f, "not valid json {{{{").expect("write");
469        }
470        let result: Result<Vec<Point>, _> = read_jsonl(&path);
471        assert!(result.is_err());
472    }
473}