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}