Skip to main content

lora_io/
format.rs

1use lora_executor::{LoraValue, Row};
2
3/// Wire format for row-level import/export.
4///
5/// JSONL is the lossless default — every [`LoraValue`] variant
6/// round-trips through the tagged JSON shape produced by
7/// [`super::lora_value_to_json`]. CSV is for spreadsheet-friendly
8/// scalar data; non-scalar columns are JSON-encoded into a single
9/// cell with a `name:json` typed header. The JSON-array variant is
10/// the same shape as JSONL but wrapped in a top-level `[ ... ]`.
11#[derive(Debug, Clone, Copy, PartialEq, Eq)]
12pub enum Format {
13    /// One JSON object per line (`\n`-delimited). Streaming-friendly.
14    Jsonl,
15    /// A single top-level JSON array of objects.
16    Json,
17    /// RFC 4180 CSV with optional `name:type` typed headers.
18    Csv,
19}
20
21impl Format {
22    /// Guess the format from a filename suffix. Returns `None` when
23    /// the extension isn't recognized.
24    pub fn from_extension(name: &str) -> Option<Self> {
25        let lower = name.to_ascii_lowercase();
26        let ext = lower.rsplit('.').next()?;
27        match ext {
28            "jsonl" | "ndjson" => Some(Format::Jsonl),
29            "json" => Some(Format::Json),
30            "csv" => Some(Format::Csv),
31            _ => None,
32        }
33    }
34
35    pub fn content_type(&self) -> &'static str {
36        match self {
37            Format::Jsonl => "application/x-ndjson",
38            Format::Json => "application/json",
39            Format::Csv => "text/csv",
40        }
41    }
42
43    pub fn as_str(&self) -> &'static str {
44        match self {
45            Format::Jsonl => "jsonl",
46            Format::Json => "json",
47            Format::Csv => "csv",
48        }
49    }
50
51    /// Parse a lowercase format name. Accepts the same tokens the
52    /// playground exposes in its menu.
53    pub fn parse(name: &str) -> Option<Self> {
54        match name.to_ascii_lowercase().as_str() {
55            "jsonl" | "ndjson" => Some(Format::Jsonl),
56            "json" => Some(Format::Json),
57            "csv" => Some(Format::Csv),
58            _ => None,
59        }
60    }
61}
62
63/// Streaming row encoder. Implementors own the underlying writer
64/// and may flush at chunk boundaries.
65pub trait RowEncoder {
66    /// Emit any header bytes the format needs (CSV header row,
67    /// opening `[` for the JSON array). Must be called exactly once
68    /// before any `write_row`.
69    fn begin(&mut self, columns: &[String]) -> std::io::Result<()>;
70
71    /// Write one row. The encoder pulls column names from the row's
72    /// [`Row::iter_named`] iterator — they must align with the
73    /// columns passed to [`Self::begin`], in the same order, for CSV.
74    fn write_row(&mut self, row: &Row) -> std::io::Result<()>;
75
76    /// Write the same shape as [`Self::write_row`] from a flat
77    /// `(name, value)` slice. Used by encoders that consume rows
78    /// from outside the engine (testing, replay).
79    fn write_named_row(&mut self, columns: &[(String, LoraValue)]) -> std::io::Result<()>;
80
81    /// Emit any trailer (closing `]` for the JSON array) and flush
82    /// the writer.
83    fn finish(&mut self) -> std::io::Result<()>;
84}
85
86/// Pull-based row decoder. Reads from a [`std::io::BufRead`] and
87/// yields completed rows one at a time. Used by native Rust callers
88/// that already have the input as a stream-like reader.
89pub trait RowDecoder {
90    /// Returns the column names declared in the file header. CSV
91    /// uses this; JSONL/JSON return `None` (every row carries its
92    /// own keys).
93    fn header(&mut self) -> std::io::Result<Option<Vec<String>>>;
94
95    /// Pull the next row as a flat `(name, value)` vector. Returns
96    /// `Ok(None)` once the input is exhausted.
97    fn next_row(&mut self) -> std::io::Result<Option<Vec<(String, LoraValue)>>>;
98}
99
100/// Push-based streaming row decoder. The caller feeds bytes one
101/// chunk at a time (e.g. from `File.stream().getReader()` in the
102/// browser); the decoder accumulates partial records internally and
103/// emits completed rows via [`Self::drain`]. Designed for the WASM
104/// streaming-import path where the engine and the source live on
105/// different sides of the worker boundary.
106///
107/// Memory bound: the decoder retains at most one in-progress record
108/// plus the bytes between the most-recently-completed record and
109/// the end of the most-recently-fed chunk.
110pub trait StreamingRowDecoder {
111    /// Append a chunk of bytes to the internal buffer and parse any
112    /// records that became complete. Idempotent: feeding zero bytes
113    /// is a no-op.
114    fn feed(&mut self, chunk: &[u8]) -> std::io::Result<()>;
115
116    /// Pull all records completed since the previous `drain` /
117    /// `finish` call. Returns an empty `Vec` when no full record
118    /// has been parsed yet.
119    fn drain(&mut self) -> std::io::Result<Vec<Vec<(String, LoraValue)>>>;
120
121    /// Signal that no more bytes will be fed. Returns any records
122    /// emitted by handling the residual buffer — for JSONL/CSV this
123    /// covers the case where the file ends without a trailing
124    /// newline.
125    fn finish(&mut self) -> std::io::Result<Vec<Vec<(String, LoraValue)>>>;
126
127    /// Column names declared in the file header. Populated after
128    /// the first record arrives for CSV; always `None` for JSONL.
129    fn header(&self) -> Option<&[String]>;
130
131    /// Total bytes accepted via `feed` since construction. Used
132    /// for progress reporting.
133    fn bytes_fed(&self) -> u64;
134
135    /// Total records emitted via `drain` + `finish` so far.
136    fn rows_emitted(&self) -> u64;
137
138    /// Switch the decoder into permissive mode. In permissive mode,
139    /// per-record parse failures are accumulated for retrieval via
140    /// [`Self::take_errors`] instead of bubbling out of `feed` /
141    /// `finish`. Fatal errors (encoding issues that desync the byte
142    /// stream itself) still bubble. Must be called before the first
143    /// `feed`; later calls take effect at the next record boundary.
144    /// Default impl is a no-op for decoders that don't support it.
145    fn set_permissive(&mut self, _on: bool) {}
146
147    /// Drain the parse errors accumulated since the previous call.
148    /// Empty when permissive mode is off or no failures occurred.
149    fn take_errors(&mut self) -> Vec<RowParseError> {
150        Vec::new()
151    }
152}
153
154/// Structured per-record parse failure. Carries the user-visible row
155/// number, an optional column attribution (CSV cell parses), a
156/// truncated sample of the raw bytes that failed, and the underlying
157/// message. Boxed into [`std::io::Error`] via [`row_parse_io_error`]
158/// so it travels through the standard `feed` / `finish` return type;
159/// recoverable via [`downcast_row_parse_error`].
160#[derive(Debug, Clone)]
161pub struct RowParseError {
162    /// 1-indexed record number. For CSV the header is record 0 and
163    /// the first data row is record 1; for JSONL/JSON-array the
164    /// first object is record 1. Records that fail mid-parse still
165    /// take a slot in this sequence.
166    pub row: u64,
167    /// Column name when the failure is attributable to a single
168    /// cell (CSV typed cell). `None` for whole-record failures and
169    /// JSON-based formats.
170    pub column: Option<String>,
171    /// Truncated raw bytes (or characters) from the offending
172    /// record. Capped to keep error payloads bounded.
173    pub raw_sample: String,
174    /// Human-readable description of what went wrong.
175    pub message: String,
176}
177
178/// Cap on `raw_sample` length to keep error payloads small enough to
179/// transport cheaply across the WASM boundary and render in the UI.
180pub const RAW_SAMPLE_MAX_CHARS: usize = 200;
181
182impl RowParseError {
183    /// Build a `raw_sample` from the offending text, replacing any
184    /// non-printable control characters with `·` (middle dot) and
185    /// truncating at [`RAW_SAMPLE_MAX_CHARS`] with a trailing `…`.
186    pub fn make_sample(raw: &str) -> String {
187        let cleaned: String = raw
188            .chars()
189            .map(|c| if c.is_control() { '·' } else { c })
190            .collect();
191        if cleaned.chars().count() <= RAW_SAMPLE_MAX_CHARS {
192            cleaned
193        } else {
194            let truncated: String = cleaned.chars().take(RAW_SAMPLE_MAX_CHARS).collect();
195            format!("{truncated}…")
196        }
197    }
198
199    /// Build a sample from a byte slice, treating non-UTF-8 bytes
200    /// with the lossy replacement char.
201    pub fn make_sample_from_bytes(raw: &[u8]) -> String {
202        Self::make_sample(&String::from_utf8_lossy(raw))
203    }
204}
205
206impl std::fmt::Display for RowParseError {
207    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
208        write!(f, "row {}", self.row)?;
209        if let Some(col) = &self.column {
210            write!(f, ", column `{col}`")?;
211        }
212        write!(f, ": {}", self.message)?;
213        if !self.raw_sample.is_empty() {
214            write!(f, " (raw: `{}`)", self.raw_sample)?;
215        }
216        Ok(())
217    }
218}
219
220impl std::error::Error for RowParseError {}
221
222/// Wrap a [`RowParseError`] in [`std::io::Error`] so it flows through
223/// the existing `Result<_, std::io::Error>` channels. Recover the
224/// structured form with [`downcast_row_parse_error`].
225pub fn row_parse_io_error(err: RowParseError) -> std::io::Error {
226    std::io::Error::new(std::io::ErrorKind::InvalidData, err)
227}
228
229/// Attempt to recover a [`RowParseError`] from an `io::Error`
230/// previously produced via [`row_parse_io_error`]. Returns `None`
231/// for any other error.
232pub fn downcast_row_parse_error(err: &std::io::Error) -> Option<&RowParseError> {
233    err.get_ref()?.downcast_ref::<RowParseError>()
234}
235
236/// Convenience: drive any [`RowEncoder`] through a row iterator and
237/// return how many rows were written. Used by [`super::import`] and
238/// can also be called directly when a caller already has a
239/// materialized result in hand.
240pub fn write_all_rows<E, I>(encoder: &mut E, columns: &[String], rows: I) -> std::io::Result<u64>
241where
242    E: RowEncoder,
243    I: IntoIterator<Item = Row>,
244{
245    encoder.begin(columns)?;
246    let mut count = 0u64;
247    for row in rows {
248        encoder.write_row(&row)?;
249        count += 1;
250    }
251    encoder.finish()?;
252    Ok(count)
253}
254
255/// Wrap an [`std::io::Error`] with extra context. Used by decoders
256/// when surfacing parser errors as I/O errors.
257pub fn invalid_data<E: Into<Box<dyn std::error::Error + Send + Sync>>>(err: E) -> std::io::Error {
258    std::io::Error::new(std::io::ErrorKind::InvalidData, err)
259}