Skip to main content

scirs2_io/
ndjson_streaming.rs

1//! NDJSON / CSV streaming types and TSV utilities
2//!
3//! Provides memory-efficient streaming interfaces for NDJSON (Newline-Delimited JSON),
4//! CSV, and TSV formats. Each type reads or writes one record at a time, making them
5//! suitable for datasets that cannot fit in memory.
6//!
7//! # Examples
8//!
9//! ```rust,no_run
10//! use scirs2_io::ndjson_streaming::{NdjsonReader, NdjsonWriter};
11//! use std::io::BufReader;
12//! use std::fs::File;
13//!
14//! // Write records
15//! let file = File::create("/tmp/records.ndjson").unwrap();
16//! let mut writer = NdjsonWriter::new(file);
17//! let record = serde_json::json!({"id": 1, "value": 3.14});
18//! writer.write_record(&record).unwrap();
19//! writer.flush().unwrap();
20//!
21//! // Read records back
22//! let file = File::open("/tmp/records.ndjson").unwrap();
23//! let mut reader = NdjsonReader::new(BufReader::new(file));
24//! while let Some(rec) = reader.next_record().unwrap() {
25//!     println!("{:?}", rec);
26//! }
27//! ```
28
29use std::fs::{File, OpenOptions};
30use std::io::{BufRead, BufReader, BufWriter, Write};
31use std::path::Path;
32
33use crate::error::{IoError, Result};
34
35// ─────────────────────────────── NdjsonReader ────────────────────────────────
36
37/// NDJSON (Newline-Delimited JSON) streaming reader.
38///
39/// Reads one JSON value per line, skipping blank lines and `#`-prefixed comments.
40/// Memory usage is proportional to the largest single record, not the whole file.
41pub struct NdjsonReader<R: BufRead> {
42    reader: R,
43    line_number: usize,
44}
45
46impl<R: BufRead> NdjsonReader<R> {
47    /// Create a new reader wrapping any `BufRead` source.
48    pub fn new(reader: R) -> Self {
49        Self {
50            reader,
51            line_number: 0,
52        }
53    }
54
55    /// Read and parse the next JSON record.
56    ///
57    /// Returns `Ok(None)` at end-of-file.  Blank lines and `#` comments are skipped.
58    pub fn next_record(&mut self) -> Result<Option<serde_json::Value>> {
59        let mut line = String::new();
60        loop {
61            line.clear();
62            let n = self
63                .reader
64                .read_line(&mut line)
65                .map_err(|e| IoError::FileError(format!("line {}: {e}", self.line_number + 1)))?;
66            if n == 0 {
67                return Ok(None);
68            }
69            self.line_number += 1;
70            let trimmed = line.trim();
71            if trimmed.is_empty() || trimmed.starts_with('#') {
72                continue;
73            }
74            let val = serde_json::from_str(trimmed)
75                .map_err(|e| IoError::ParseError(format!("line {}: {e}", self.line_number)))?;
76            return Ok(Some(val));
77        }
78    }
79
80    /// Collect every remaining record into a `Vec`.
81    pub fn collect_all(&mut self) -> Result<Vec<serde_json::Value>> {
82        let mut out = Vec::new();
83        while let Some(v) = self.next_record()? {
84            out.push(v);
85        }
86        Ok(out)
87    }
88
89    /// Count remaining records without storing them.
90    pub fn count_records(&mut self) -> Result<usize> {
91        let mut count = 0usize;
92        while self.next_record()?.is_some() {
93            count += 1;
94        }
95        Ok(count)
96    }
97
98    /// Current 1-based line number (lines consumed so far).
99    pub fn line_number(&self) -> usize {
100        self.line_number
101    }
102}
103
104// ─────────────────────────────── NdjsonWriter ────────────────────────────────
105
106/// NDJSON (Newline-Delimited JSON) streaming writer.
107///
108/// Serialises each [`serde_json::Value`] as a single line and appends `\n`.
109pub struct NdjsonWriter<W: Write> {
110    writer: W,
111}
112
113impl<W: Write> NdjsonWriter<W> {
114    /// Create a new writer wrapping any `Write` sink.
115    pub fn new(writer: W) -> Self {
116        Self { writer }
117    }
118
119    /// Serialise `record` and write it as one line.
120    pub fn write_record(&mut self, record: &serde_json::Value) -> Result<()> {
121        let json = serde_json::to_string(record)
122            .map_err(|e| IoError::SerializationError(format!("JSON serialization failed: {e}")))?;
123        self.writer
124            .write_all(json.as_bytes())
125            .map_err(|e| IoError::FileError(format!("write failed: {e}")))?;
126        self.writer
127            .write_all(b"\n")
128            .map_err(|e| IoError::FileError(format!("write newline failed: {e}")))?;
129        Ok(())
130    }
131
132    /// Flush the underlying writer.
133    pub fn flush(&mut self) -> Result<()> {
134        self.writer
135            .flush()
136            .map_err(|e| IoError::FileError(format!("flush failed: {e}")))
137    }
138}
139
140// ─────────────────────────────── CsvValue ────────────────────────────────────
141
142/// A single type-inferred cell from a CSV / TSV row.
143#[derive(Debug, Clone, PartialEq)]
144pub enum CsvValue {
145    /// Parsed as a 64-bit integer.
146    Integer(i64),
147    /// Parsed as a 64-bit float.
148    Float(f64),
149    /// Parsed as a boolean (`true`/`false`/`yes`/`no`/`1`/`0`).
150    Boolean(bool),
151    /// Raw text that could not be parsed as any structured type.
152    Text(String),
153    /// Empty field or explicit null sentinel.
154    Null,
155}
156
157impl CsvValue {
158    fn infer(s: &str) -> Self {
159        let trimmed = s.trim();
160        if trimmed.is_empty() || trimmed.eq_ignore_ascii_case("null") || trimmed == "NA" {
161            return CsvValue::Null;
162        }
163        // Integer check before Boolean to avoid treating "1"/"0" as booleans.
164        if let Ok(i) = trimmed.parse::<i64>() {
165            return CsvValue::Integer(i);
166        }
167        // Boolean (word forms only; numeric "1"/"0" handled above as integers)
168        match trimmed.to_lowercase().as_str() {
169            "true" | "yes" => return CsvValue::Boolean(true),
170            "false" | "no" => return CsvValue::Boolean(false),
171            _ => {}
172        }
173        // Float
174        if let Ok(f) = trimmed.parse::<f64>() {
175            return CsvValue::Float(f);
176        }
177        CsvValue::Text(trimmed.to_string())
178    }
179}
180
181// ─────────────────────────────── CsvStreamReader ─────────────────────────────
182
183/// Streaming CSV reader with per-row type inference.
184///
185/// Reads one row at a time; does not load the entire file into memory.
186/// Handles quoted fields and configurable delimiters.
187pub struct CsvStreamReader<R: BufRead> {
188    reader: R,
189    delimiter: u8,
190    headers: Option<Vec<String>>,
191    has_header: bool,
192    line_number: usize,
193    finished: bool,
194}
195
196impl<R: BufRead> CsvStreamReader<R> {
197    /// Create a new streaming CSV reader.
198    ///
199    /// If `has_header` is `true` the first non-blank line is consumed immediately
200    /// and stored; subsequent calls to `next_row` return data rows only.
201    pub fn new(mut reader: R, has_header: bool, delimiter: u8) -> Result<Self> {
202        let headers = if has_header {
203            let mut line = String::new();
204            loop {
205                line.clear();
206                let n = reader
207                    .read_line(&mut line)
208                    .map_err(|e| IoError::FileError(format!("header read error: {e}")))?;
209                if n == 0 {
210                    break None;
211                }
212                let trimmed = line.trim();
213                if !trimmed.is_empty() {
214                    let hdrs = parse_csv_row(trimmed, delimiter);
215                    break Some(hdrs);
216                }
217            }
218        } else {
219            None
220        };
221
222        Ok(Self {
223            reader,
224            delimiter,
225            headers,
226            has_header,
227            line_number: if has_header { 1 } else { 0 },
228            finished: false,
229        })
230    }
231
232    /// Return the parsed header row, or `None` if `has_header` was `false`.
233    pub fn headers(&self) -> Option<&[String]> {
234        self.headers.as_deref()
235    }
236
237    /// Read the next raw (string) row.  Returns `Ok(None)` at end-of-file.
238    pub fn next_row(&mut self) -> Result<Option<Vec<String>>> {
239        if self.finished {
240            return Ok(None);
241        }
242        let mut line = String::new();
243        loop {
244            line.clear();
245            let n = self
246                .reader
247                .read_line(&mut line)
248                .map_err(|e| IoError::FileError(format!("line {}: {e}", self.line_number + 1)))?;
249            if n == 0 {
250                self.finished = true;
251                return Ok(None);
252            }
253            self.line_number += 1;
254            let trimmed = line.trim();
255            if trimmed.is_empty() {
256                continue;
257            }
258            return Ok(Some(parse_csv_row(trimmed, self.delimiter)));
259        }
260    }
261
262    /// Read the next row with automatic type inference applied to each field.
263    pub fn next_typed_row(&mut self) -> Result<Option<Vec<CsvValue>>> {
264        match self.next_row()? {
265            None => Ok(None),
266            Some(fields) => Ok(Some(fields.iter().map(|s| CsvValue::infer(s)).collect())),
267        }
268    }
269}
270
271/// Parse a single CSV/TSV line respecting double-quoted fields.
272fn parse_csv_row(line: &str, delimiter: u8) -> Vec<String> {
273    let sep = delimiter as char;
274    let mut fields = Vec::new();
275    let mut current = String::new();
276    let mut in_quotes = false;
277    let mut chars = line.chars().peekable();
278
279    while let Some(ch) = chars.next() {
280        if ch == '"' {
281            if in_quotes {
282                // Check for escaped quote (`""`)
283                if chars.peek() == Some(&'"') {
284                    chars.next();
285                    current.push('"');
286                } else {
287                    in_quotes = false;
288                }
289            } else {
290                in_quotes = true;
291            }
292        } else if ch == sep && !in_quotes {
293            fields.push(current.trim().to_string());
294            current.clear();
295        } else {
296            current.push(ch);
297        }
298    }
299    fields.push(current.trim().to_string());
300    fields
301}
302
303// ─────────────────────────────── TSV helpers ─────────────────────────────────
304
305/// Read an entire TSV file into `(headers, rows)`.
306///
307/// The first row is treated as the header.  All subsequent rows are data.
308pub fn read_tsv(path: &Path) -> Result<(Vec<String>, Vec<Vec<String>>)> {
309    let file =
310        File::open(path).map_err(|e| IoError::FileError(format!("cannot open {:?}: {e}", path)))?;
311    let mut reader = CsvStreamReader::new(BufReader::new(file), true, b'\t')?;
312
313    let headers = reader
314        .headers()
315        .ok_or_else(|| IoError::FormatError("TSV file appears empty".to_string()))?
316        .to_vec();
317
318    let mut rows = Vec::new();
319    while let Some(row) = reader.next_row()? {
320        rows.push(row);
321    }
322    Ok((headers, rows))
323}
324
325/// Write headers and rows to a TSV file.
326pub fn write_tsv(path: &Path, headers: &[String], data: &[Vec<String>]) -> Result<()> {
327    let file = File::create(path)
328        .map_err(|e| IoError::FileError(format!("cannot create {:?}: {e}", path)))?;
329    let mut writer = BufWriter::new(file);
330
331    writer
332        .write_all(headers.join("\t").as_bytes())
333        .map_err(|e| IoError::FileError(format!("write header failed: {e}")))?;
334    writer
335        .write_all(b"\n")
336        .map_err(|e| IoError::FileError(format!("write newline failed: {e}")))?;
337
338    for row in data {
339        writer
340            .write_all(row.join("\t").as_bytes())
341            .map_err(|e| IoError::FileError(format!("write row failed: {e}")))?;
342        writer
343            .write_all(b"\n")
344            .map_err(|e| IoError::FileError(format!("write newline failed: {e}")))?;
345    }
346    writer
347        .flush()
348        .map_err(|e| IoError::FileError(format!("flush failed: {e}")))
349}
350
351// ─────────────────────────────── File-level helpers ──────────────────────────
352
353/// Open an NDJSON file and return a reader over it.
354pub fn open_ndjson_file(path: &Path) -> Result<NdjsonReader<BufReader<File>>> {
355    let file =
356        File::open(path).map_err(|e| IoError::FileError(format!("cannot open {:?}: {e}", path)))?;
357    Ok(NdjsonReader::new(BufReader::new(file)))
358}
359
360/// Create / overwrite an NDJSON file and return a buffered writer over it.
361pub fn create_ndjson_file(path: &Path) -> Result<NdjsonWriter<BufWriter<File>>> {
362    let file = File::create(path)
363        .map_err(|e| IoError::FileError(format!("cannot create {:?}: {e}", path)))?;
364    Ok(NdjsonWriter::new(BufWriter::new(file)))
365}
366
367/// Append to an existing NDJSON file (creates it if absent).
368pub fn append_ndjson_file(path: &Path) -> Result<NdjsonWriter<BufWriter<File>>> {
369    let file = OpenOptions::new()
370        .create(true)
371        .append(true)
372        .open(path)
373        .map_err(|e| IoError::FileError(format!("cannot open {:?} for append: {e}", path)))?;
374    Ok(NdjsonWriter::new(BufWriter::new(file)))
375}
376
377// ─────────────────────────────── Tests ───────────────────────────────────────
378
379#[cfg(test)]
380mod tests {
381    use super::*;
382    use std::io::BufReader;
383
384    fn ndjson_bytes(lines: &[&str]) -> Vec<u8> {
385        lines.join("\n").into_bytes()
386    }
387
388    // ── NdjsonReader ──────────────────────────────────────────────────────────
389
390    #[test]
391    fn test_ndjson_reader_single_record() {
392        let src = ndjson_bytes(&[r#"{"id":1,"v":2.5}"#]);
393        let mut r = NdjsonReader::new(BufReader::new(src.as_slice()));
394        let rec = r
395            .next_record()
396            .expect("should parse")
397            .expect("should have record");
398        assert_eq!(rec["id"], 1);
399        assert!((rec["v"].as_f64().expect("float") - 2.5).abs() < 1e-10);
400        assert!(r.next_record().expect("no error").is_none());
401    }
402
403    #[test]
404    fn test_ndjson_reader_multi_record() {
405        let src = ndjson_bytes(&[r#"{"a":1}"#, r#"{"a":2}"#, r#"{"a":3}"#]);
406        let mut r = NdjsonReader::new(BufReader::new(src.as_slice()));
407        let all = r.collect_all().expect("collect");
408        assert_eq!(all.len(), 3);
409        assert_eq!(all[2]["a"], 3);
410    }
411
412    #[test]
413    fn test_ndjson_reader_skips_blank_and_comment_lines() {
414        let src = ndjson_bytes(&["", "# comment", r#"{"x":42}"#, "", r#"{"x":99}"#]);
415        let mut r = NdjsonReader::new(BufReader::new(src.as_slice()));
416        assert_eq!(r.count_records().expect("count"), 2);
417    }
418
419    #[test]
420    fn test_ndjson_reader_empty_source() {
421        let src: &[u8] = b"";
422        let mut r = NdjsonReader::new(BufReader::new(src));
423        assert!(r.next_record().expect("no error").is_none());
424    }
425
426    // ── NdjsonWriter ─────────────────────────────────────────────────────────
427
428    #[test]
429    fn test_ndjson_writer_produces_newline_delimited_json() {
430        let mut buf: Vec<u8> = Vec::new();
431        let mut w = NdjsonWriter::new(&mut buf);
432        w.write_record(&serde_json::json!({"k": "v1"}))
433            .expect("write");
434        w.write_record(&serde_json::json!({"k": "v2"}))
435            .expect("write");
436        w.flush().expect("flush");
437
438        let text = String::from_utf8(buf).expect("utf8");
439        let lines: Vec<_> = text.lines().collect();
440        assert_eq!(lines.len(), 2);
441        let v: serde_json::Value = serde_json::from_str(lines[1]).expect("parse");
442        assert_eq!(v["k"], "v2");
443    }
444
445    // ── NdjsonWriter + NdjsonReader round-trip ────────────────────────────────
446
447    #[test]
448    fn test_ndjson_roundtrip_via_temp_file() {
449        let dir = std::env::temp_dir().join("scirs2_io_ndjson_rt_test");
450        std::fs::create_dir_all(&dir).expect("mkdir");
451        let path = dir.join("roundtrip.ndjson");
452
453        let records = vec![
454            serde_json::json!({"id": 1, "name": "alpha", "score": 9.5}),
455            serde_json::json!({"id": 2, "name": "beta",  "score": 7.2}),
456            serde_json::json!({"id": 3, "name": "gamma", "score": 8.8}),
457        ];
458
459        {
460            let mut w = create_ndjson_file(&path).expect("create");
461            for rec in &records {
462                w.write_record(rec).expect("write");
463            }
464            w.flush().expect("flush");
465        }
466
467        let mut r = open_ndjson_file(&path).expect("open");
468        let loaded = r.collect_all().expect("collect");
469
470        assert_eq!(loaded.len(), 3);
471        for (orig, loaded_rec) in records.iter().zip(loaded.iter()) {
472            assert_eq!(orig["id"], loaded_rec["id"]);
473            assert_eq!(orig["name"], loaded_rec["name"]);
474        }
475
476        let _ = std::fs::remove_dir_all(&dir);
477    }
478
479    // ── CsvStreamReader ───────────────────────────────────────────────────────
480
481    #[test]
482    fn test_csv_stream_reader_headers_and_rows() {
483        let csv = b"name,age,city\nAlice,30,London\nBob,25,Paris\n";
484        let mut r =
485            CsvStreamReader::new(BufReader::new(csv.as_ref()), true, b',').expect("new reader");
486
487        let hdrs = r.headers().expect("headers").to_vec();
488        assert_eq!(hdrs, vec!["name", "age", "city"]);
489
490        let row1 = r.next_row().expect("row1 err").expect("row1 some");
491        assert_eq!(row1, vec!["Alice", "30", "London"]);
492
493        let row2 = r.next_row().expect("row2 err").expect("row2 some");
494        assert_eq!(row2, vec!["Bob", "25", "Paris"]);
495
496        assert!(r.next_row().expect("eof err").is_none());
497    }
498
499    #[test]
500    fn test_csv_stream_reader_no_header() {
501        let csv = b"1,2,3\n4,5,6\n";
502        let mut r =
503            CsvStreamReader::new(BufReader::new(csv.as_ref()), false, b',').expect("new reader");
504        assert!(r.headers().is_none());
505        let row = r.next_row().expect("row").expect("some");
506        assert_eq!(row, vec!["1", "2", "3"]);
507    }
508
509    #[test]
510    fn test_csv_stream_reader_typed_row() {
511        let csv = b"id,active,value,label\n1,true,3.14,hello\n2,false,,NA\n";
512        let mut r =
513            CsvStreamReader::new(BufReader::new(csv.as_ref()), true, b',').expect("new reader");
514
515        let row = r.next_typed_row().expect("row").expect("some");
516        assert!(matches!(row[0], CsvValue::Integer(1)));
517        assert!(matches!(row[1], CsvValue::Boolean(true)));
518        assert!(matches!(row[2], CsvValue::Float(_)));
519        assert!(matches!(row[3], CsvValue::Text(_)));
520
521        let row2 = r.next_typed_row().expect("row2").expect("some2");
522        assert!(matches!(row2[2], CsvValue::Null));
523        assert!(matches!(row2[3], CsvValue::Null));
524    }
525
526    #[test]
527    fn test_csv_stream_reader_tsv_delimiter() {
528        let tsv = b"a\tb\tc\n10\t20\t30\n";
529        let mut r =
530            CsvStreamReader::new(BufReader::new(tsv.as_ref()), true, b'\t').expect("new reader");
531        let hdrs = r.headers().expect("hdrs").to_vec();
532        assert_eq!(hdrs, vec!["a", "b", "c"]);
533        let row = r.next_row().expect("row").expect("some");
534        assert_eq!(row, vec!["10", "20", "30"]);
535    }
536
537    // ── TSV read/write round-trip ─────────────────────────────────────────────
538
539    #[test]
540    fn test_tsv_roundtrip() {
541        let dir = std::env::temp_dir().join("scirs2_io_tsv_rt_test");
542        std::fs::create_dir_all(&dir).expect("mkdir");
543        let path = dir.join("data.tsv");
544
545        let headers = vec!["x".to_string(), "y".to_string(), "z".to_string()];
546        let data = vec![
547            vec!["1".to_string(), "2".to_string(), "3".to_string()],
548            vec!["4".to_string(), "5".to_string(), "6".to_string()],
549        ];
550
551        write_tsv(&path, &headers, &data).expect("write tsv");
552        let (read_hdrs, read_data) = read_tsv(&path).expect("read tsv");
553
554        assert_eq!(read_hdrs, headers);
555        assert_eq!(read_data, data);
556
557        let _ = std::fs::remove_dir_all(&dir);
558    }
559
560    // ── CsvValue::infer edge-cases ────────────────────────────────────────────
561
562    #[test]
563    fn test_csv_value_infer() {
564        assert!(matches!(CsvValue::infer(""), CsvValue::Null));
565        assert!(matches!(CsvValue::infer("null"), CsvValue::Null));
566        assert!(matches!(CsvValue::infer("NA"), CsvValue::Null));
567        assert!(matches!(CsvValue::infer("true"), CsvValue::Boolean(true)));
568        assert!(matches!(CsvValue::infer("False"), CsvValue::Boolean(false)));
569        assert!(matches!(CsvValue::infer("42"), CsvValue::Integer(42)));
570        assert!(matches!(CsvValue::infer("3.14"), CsvValue::Float(_)));
571        assert!(matches!(CsvValue::infer("hello"), CsvValue::Text(_)));
572    }
573}