1use 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 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
43pub 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
63pub 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
113pub 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
161pub 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}