Skip to main content

leviathan/
source.rs

1//! Record sources: JSONL / NDJSON, JSON arrays, CSV / TSV (each optionally
2//! gzip), SQLite tables or queries, and stdin. Every source yields JSON
3//! objects; line-oriented formats stream, so inputs far larger than memory
4//! are fine.
5//!
6//! Any other database works through its export tool and stdin, e.g.
7//! `psql --csv -c 'select ...' | leviathan index - --format csv`.
8
9use std::collections::BTreeSet;
10use std::fs::{self, File};
11use std::io::{self, BufRead, BufReader, Read};
12use std::path::{Path, PathBuf};
13use std::time::UNIX_EPOCH;
14
15use anyhow::{Context, Result, bail};
16use flate2::read::MultiGzDecoder;
17use serde_json::{Map, Value};
18
19use crate::config::Format;
20
21const BUF: usize = 1 << 20;
22
23#[derive(Debug, Clone)]
24pub struct Source {
25    pub path: PathBuf,
26    pub format: Format,
27}
28
29impl Source {
30    pub fn is_stdin(&self) -> bool {
31        self.path.as_os_str() == "-"
32    }
33
34    /// Short label for synthetic ids and messages.
35    pub fn label(&self) -> String {
36        if self.is_stdin() {
37            return "stdin".into();
38        }
39        self.path.file_name().map(|n| n.to_string_lossy().into_owned()).unwrap_or_default()
40    }
41}
42
43/// One unit read from a source.
44pub enum Item {
45    Record { line: u64, value: Value, raw: Option<String> },
46    Bad { line: u64, error: String },
47}
48
49fn format_of(path: &Path) -> Option<Format> {
50    let name = path.file_name()?.to_str()?.to_ascii_lowercase();
51    let name = name.strip_suffix(".gz").unwrap_or(&name);
52    let ext = name.rsplit_once('.')?.1;
53    Some(match ext {
54        "jsonl" | "ndjson" => Format::Jsonl,
55        "json" => Format::Json,
56        "csv" => Format::Csv,
57        "tsv" | "tab" => Format::Tsv,
58        "db" | "sqlite" | "sqlite3" => Format::Sqlite,
59        _ => return None,
60    })
61}
62
63/// Expand files and directories into sources. Directories contribute every
64/// JSONL, JSON, CSV and TSV file (optionally `.gz`) beneath them; SQLite
65/// files must be named explicitly.
66pub fn discover(paths: &[PathBuf], format: Format) -> Result<Vec<Source>> {
67    let mut out = Vec::new();
68    for path in paths {
69        if path.as_os_str() == "-" {
70            out.push(Source { path: path.clone(), format });
71        } else if path.is_dir() {
72            let mut found = BTreeSet::new();
73            walk(path, &mut found)?;
74            for file in found {
75                let detected = format_of(&file).unwrap_or(Format::Auto);
76                out.push(Source {
77                    path: file,
78                    format: if format == Format::Auto { detected } else { format },
79                });
80            }
81        } else if path.is_file() {
82            let detected = format_of(path).unwrap_or(Format::Auto);
83            out.push(Source {
84                path: path.clone(),
85                format: if format == Format::Auto { detected } else { format },
86            });
87        } else {
88            bail!("source path not found: {}", path.display());
89        }
90    }
91    if out.is_empty() {
92        bail!("no data files found (looked for .jsonl .ndjson .json .csv .tsv, optionally .gz)");
93    }
94    Ok(out)
95}
96
97fn walk(dir: &Path, out: &mut BTreeSet<PathBuf>) -> Result<()> {
98    for entry in fs::read_dir(dir).with_context(|| format!("read {}", dir.display()))? {
99        let path = entry?.path();
100        let hidden = path.file_name().and_then(|n| n.to_str()).is_some_and(|n| n.starts_with('.'));
101        if hidden {
102            continue;
103        }
104        if path.is_dir() {
105            walk(&path, out)?;
106        } else if matches!(format_of(&path), Some(Format::Jsonl | Format::Json | Format::Csv | Format::Tsv)) {
107            out.insert(path);
108        }
109    }
110    Ok(())
111}
112
113/// Paths, sizes and mtimes; empty when a source cannot be fingerprinted
114/// (stdin), which forces a rebuild.
115pub fn manifest(sources: &[Source], extra: &str) -> Result<String> {
116    let mut entries = Vec::new();
117    for s in sources {
118        if s.is_stdin() {
119            return Ok(String::new());
120        }
121        let meta = fs::metadata(&s.path).with_context(|| format!("stat {}", s.path.display()))?;
122        let mtime = meta.modified()?.duration_since(UNIX_EPOCH)?.as_nanos();
123        entries.push(serde_json::json!({
124            "path": s.path.canonicalize().unwrap_or_else(|_| s.path.clone()),
125            "size": meta.len(),
126            "mtime_ns": mtime.to_string(),
127        }));
128    }
129    Ok(serde_json::to_string(&serde_json::json!({"files": entries, "mapping": extra}))?)
130}
131
132pub fn byte_size(sources: &[Source]) -> u64 {
133    sources.iter().filter(|s| !s.is_stdin()).filter_map(|s| fs::metadata(&s.path).ok()).map(|m| m.len()).sum()
134}
135
136fn open(source: &Source) -> Result<Box<dyn BufRead>> {
137    if source.is_stdin() {
138        return Ok(Box::new(BufReader::with_capacity(BUF, io::stdin())));
139    }
140    let file = File::open(&source.path).with_context(|| format!("open {}", source.path.display()))?;
141    let gz = source.path.extension().is_some_and(|e| e.eq_ignore_ascii_case("gz"));
142    let reader: Box<dyn Read> =
143        if gz { Box::new(MultiGzDecoder::new(BufReader::with_capacity(BUF, file))) } else { Box::new(file) };
144    Ok(Box::new(BufReader::with_capacity(BUF, reader)))
145}
146
147fn first_byte(reader: &mut Box<dyn BufRead>) -> Result<Option<u8>> {
148    loop {
149        let buf = reader.fill_buf()?;
150        if buf.is_empty() {
151            return Ok(None);
152        }
153        if let Some(pos) = buf.iter().position(|b| !b.is_ascii_whitespace()) {
154            return Ok(Some(buf[pos]));
155        }
156        let n = buf.len();
157        reader.consume(n);
158    }
159}
160
161/// Stream every record of `source` into `on`. Return `false` from `on` to
162/// stop early (sampling).
163pub fn read(source: &Source, sql: Option<&str>, on: &mut dyn FnMut(Item) -> Result<bool>) -> Result<()> {
164    if source.format == Format::Sqlite {
165        return read_sqlite(&source.path, sql, on);
166    }
167    let mut reader = open(source)?;
168    let format = match source.format {
169        Format::Auto | Format::Json => match first_byte(&mut reader)? {
170            None => return Ok(()),
171            Some(b'[') => Format::Json,
172            Some(b'{') => Format::Jsonl,
173            Some(_) if source.format == Format::Json => Format::Jsonl,
174            Some(_) => Format::Csv,
175        },
176        f => f,
177    };
178    match format {
179        Format::Json => read_json_array(reader, on),
180        Format::Csv => read_delimited(reader, b',', on),
181        Format::Tsv => read_delimited(reader, b'\t', on),
182        _ => read_jsonl(reader, on),
183    }
184}
185
186fn read_jsonl(mut reader: Box<dyn BufRead>, on: &mut dyn FnMut(Item) -> Result<bool>) -> Result<()> {
187    let mut line = String::new();
188    let mut line_no = 0u64;
189    loop {
190        line.clear();
191        if reader.read_line(&mut line)? == 0 {
192            return Ok(());
193        }
194        line_no += 1;
195        let raw = line.trim();
196        if raw.is_empty() {
197            continue;
198        }
199        let item = match serde_json::from_str::<Value>(raw) {
200            Ok(value @ Value::Object(_)) => Item::Record { line: line_no, value, raw: Some(raw.to_string()) },
201            Ok(_) => Item::Bad { line: line_no, error: "not a JSON object".into() },
202            Err(err) => Item::Bad { line: line_no, error: err.to_string() },
203        };
204        if !on(item)? {
205            return Ok(());
206        }
207    }
208}
209
210fn read_json_array(reader: Box<dyn BufRead>, on: &mut dyn FnMut(Item) -> Result<bool>) -> Result<()> {
211    let items: Vec<Value> = serde_json::from_reader(reader).context("parse JSON array")?;
212    for (i, value) in items.into_iter().enumerate() {
213        let line = i as u64 + 1;
214        let item = match value {
215            Value::Object(_) => Item::Record { line, value, raw: None },
216            _ => Item::Bad { line, error: "array element is not an object".into() },
217        };
218        if !on(item)? {
219            break;
220        }
221    }
222    Ok(())
223}
224
225fn read_delimited(
226    reader: Box<dyn BufRead>,
227    delimiter: u8,
228    on: &mut dyn FnMut(Item) -> Result<bool>,
229) -> Result<()> {
230    let mut csv = csv::ReaderBuilder::new().delimiter(delimiter).flexible(true).from_reader(reader);
231    let headers: Vec<String> = csv.headers()?.iter().map(|h| h.trim().to_string()).collect();
232    for row in csv.records() {
233        let item = match row {
234            Ok(row) => {
235                let line = row.position().map(|p| p.line()).unwrap_or(0);
236                let mut map = Map::new();
237                for (h, v) in headers.iter().zip(row.iter()) {
238                    if !v.trim().is_empty() && !h.is_empty() {
239                        map.insert(h.clone(), Value::String(v.to_string()));
240                    }
241                }
242                Item::Record { line, value: Value::Object(map), raw: None }
243            }
244            Err(err) => {
245                let line = err.position().map(|p| p.line()).unwrap_or(0);
246                Item::Bad { line, error: err.to_string() }
247            }
248        };
249        if !on(item)? {
250            break;
251        }
252    }
253    Ok(())
254}
255
256fn read_sqlite(path: &Path, sql: Option<&str>, on: &mut dyn FnMut(Item) -> Result<bool>) -> Result<()> {
257    use rusqlite::types::ValueRef;
258    let conn =
259        crate::index::open_ro(path).with_context(|| format!("open SQLite source {}", path.display()))?;
260    let query = match sql {
261        Some(q) => q.to_string(),
262        None => {
263            let tables: Vec<String> = conn
264                .prepare("SELECT name FROM sqlite_master WHERE type IN ('table','view') AND name NOT LIKE 'sqlite_%'")?
265                .query_map([], |r| r.get(0))?
266                .collect::<rusqlite::Result<_>>()?;
267            match tables.as_slice() {
268                [one] => format!("SELECT * FROM \"{}\"", one.replace('"', "\"\"")),
269                _ => bail!(
270                    "{} has {} tables ({}); choose one with --sql 'SELECT * FROM <table>'",
271                    path.display(),
272                    tables.len(),
273                    tables.join(", ")
274                ),
275            }
276        }
277    };
278    let mut stmt = conn.prepare(&query).with_context(|| format!("prepare {query:?}"))?;
279    let names: Vec<String> = stmt.column_names().into_iter().map(str::to_string).collect();
280    let mut rows = stmt.query([])?;
281    let mut line = 0u64;
282    while let Some(row) = rows.next()? {
283        line += 1;
284        let mut map = Map::new();
285        for (i, name) in names.iter().enumerate() {
286            let value = match row.get_ref(i)? {
287                ValueRef::Null | ValueRef::Blob(_) => continue,
288                ValueRef::Integer(n) => Value::from(n),
289                ValueRef::Real(f) => {
290                    serde_json::Number::from_f64(f).map(Value::Number).unwrap_or(Value::Null)
291                }
292                ValueRef::Text(t) => {
293                    let s = String::from_utf8_lossy(t).into_owned();
294                    match s.trim_start().as_bytes().first() {
295                        Some(b'{' | b'[') => serde_json::from_str(&s).unwrap_or(Value::String(s)),
296                        _ => Value::String(s),
297                    }
298                }
299            };
300            map.insert(name.clone(), value);
301        }
302        if !on(Item::Record { line, value: Value::Object(map), raw: None })? {
303            break;
304        }
305    }
306    Ok(())
307}