1use std::collections::HashMap;
9use std::fs;
10use std::path::{Path, PathBuf};
11use std::time::Instant;
12
13use anyhow::{Context, Result, bail};
14use rusqlite::{Connection, OpenFlags, OptionalExtension, params};
15use serde::Serialize;
16
17use crate::config::{Config, Fields};
18use crate::fields::Mapping;
19use crate::source::{self, Item, Source};
20use crate::text::normalize_name;
21
22pub const SCHEMA_VERSION: &str = "2";
23
24fn schema_sql(title_weight: f64) -> String {
25 format!(
26 r#"
27CREATE TABLE meta (
28 key TEXT PRIMARY KEY,
29 value TEXT NOT NULL
30);
31
32CREATE TABLE groups (
33 key TEXT PRIMARY KEY,
34 key_lower TEXT,
35 name TEXT,
36 name_normalized TEXT,
37 record_count INTEGER NOT NULL DEFAULT 0
38);
39CREATE INDEX idx_groups_key_lower ON groups (key_lower);
40CREATE INDEX idx_groups_name_norm ON groups (name_normalized);
41
42CREATE TABLE records (
43 rowid INTEGER PRIMARY KEY,
44 id TEXT NOT NULL UNIQUE,
45 grp TEXT,
46 grp_name TEXT,
47 date TEXT,
48 boost REAL NOT NULL DEFAULT 1.0,
49 facets TEXT,
50 doc TEXT NOT NULL
51);
52CREATE INDEX idx_records_grp_date ON records (grp, date);
53CREATE INDEX idx_records_date ON records (date);
54
55CREATE TABLE facets (
56 field TEXT NOT NULL,
57 key TEXT NOT NULL,
58 value TEXT NOT NULL,
59 count INTEGER NOT NULL,
60 PRIMARY KEY (field, key)
61) WITHOUT ROWID;
62
63-- Contentless: text is tokenized, never stored twice. `names` holds the
64-- group key and name, searched only outside a group scope (inside one they
65-- match every record). `tags` holds one synthetic token per group and per
66-- filter value, so scoping and filtering are posting-list intersections,
67-- not post-filters over every match.
68CREATE VIRTUAL TABLE record_fts USING fts5(
69 title, body, names, tags,
70 content = '', contentless_delete = 1,
71 tokenize = 'porter unicode61'
72);
73INSERT INTO record_fts (record_fts, rank) VALUES ('rank', 'bm25({title_weight}, 1.0, 0.5, 0.0)');
74"#
75 )
76}
77
78#[derive(Debug, Default, Serialize)]
79pub struct IngestReport {
80 pub built: bool,
81 pub index_path: PathBuf,
82 pub records_read: u64,
83 pub inserted: u64,
84 pub updated: u64,
86 #[serde(skip_serializing_if = "is_zero")]
87 pub deleted: u64,
88 pub skipped_lines: u64,
89 #[serde(skip_serializing_if = "Vec::is_empty")]
90 pub skipped_examples: Vec<String>,
91 pub record_count: i64,
92 pub group_count: i64,
93 pub source_bytes: u64,
94 pub index_bytes: u64,
95 pub elapsed_seconds: f64,
96 #[serde(skip_serializing_if = "Option::is_none")]
98 pub inferred_mapping: Option<Fields>,
99}
100
101fn is_zero(n: &u64) -> bool {
102 *n == 0
103}
104
105#[derive(Debug, Clone, Copy, Default)]
106pub struct BuildOptions {
107 pub force: bool,
108 pub strict: bool,
110 pub quiet: bool,
111}
112
113pub fn build(
115 index_path: &Path,
116 sources: &[Source],
117 config: &Config,
118 inferred: bool,
119 opts: BuildOptions,
120) -> Result<IngestReport> {
121 let mapping = Mapping::new(config)?;
122 let config_json = serde_json::to_string(config)?;
123 let manifest = source::manifest(sources, &config_json)?;
124 if !opts.force && !manifest.is_empty() && index_is_current(index_path, &manifest) {
125 let conn = open_ro(index_path)?;
126 return Ok(IngestReport {
127 index_path: index_path.to_path_buf(),
128 index_bytes: fs::metadata(index_path)?.len(),
129 record_count: count(&conn, "records")?,
130 group_count: count(&conn, "groups")?,
131 ..Default::default()
132 });
133 }
134
135 let started = Instant::now();
136 if let Some(parent) = index_path.parent().filter(|p| !p.as_os_str().is_empty()) {
137 fs::create_dir_all(parent)?;
138 }
139 let building = sibling(index_path, ".building");
140 let _ = fs::remove_file(&building);
141
142 let result = (|| -> Result<IngestReport> {
143 let mut conn = Connection::open(&building)?;
144 conn.execute_batch(
145 "PRAGMA journal_mode = OFF; PRAGMA synchronous = OFF; \
146 PRAGMA temp_store = MEMORY; PRAGMA cache_size = -262144; PRAGMA page_size = 8192;",
147 )?;
148 conn.execute_batch(&schema_sql(config.rank.title_weight))?;
149 let tx = conn.transaction()?;
150 let mut ingest = Ingest::new(&tx, &mapping, opts, false)?;
151 for s in sources {
152 ingest.source(s, config.source.sql.as_deref())?;
153 }
154 let mut report = ingest.finish()?;
155 rebuild_groups(&tx, false)?;
156 log(opts, "optimizing full-text index");
157 tx.execute("INSERT INTO record_fts (record_fts) VALUES ('optimize')", [])?;
158 report.built = true;
159 report.source_bytes = source::byte_size(sources);
160 report.record_count = count(&tx, "records")?;
161 report.group_count = count(&tx, "groups")?;
162 for (key, value) in [
163 ("schema_version", SCHEMA_VERSION.to_string()),
164 ("leviathan_version", env!("CARGO_PKG_VERSION").to_string()),
165 ("built_at", crate::now_rfc3339()),
166 ("config", config_json.clone()),
167 ("mapping_inferred", inferred.to_string()),
168 ("skipped_lines", report.skipped_lines.to_string()),
169 ("source_bytes", report.source_bytes.to_string()),
170 ("source_manifest", manifest.clone()),
171 ] {
172 set_meta(&tx, key, &value)?;
173 }
174 refresh_counts(&tx)?;
175 tx.commit()?;
176 conn.close().map_err(|(_, e)| e)?;
177 Ok(report)
178 })();
179
180 let mut report = match result {
181 Ok(report) => report,
182 Err(err) => {
183 let _ = fs::remove_file(&building);
184 return Err(err);
185 }
186 };
187 fs::rename(&building, index_path)
188 .with_context(|| format!("install index at {}", index_path.display()))?;
189 report.index_path = index_path.to_path_buf();
190 report.index_bytes = fs::metadata(index_path)?.len();
191 report.elapsed_seconds = round1(started.elapsed().as_secs_f64());
192 if inferred {
193 report.inferred_mapping = Some(config.fields.clone());
194 }
195 Ok(report)
196}
197
198pub fn upsert(index_path: &Path, sources: &[Source], opts: BuildOptions) -> Result<IngestReport> {
201 let started = Instant::now();
202 let mut conn = open_rw(index_path)?;
203 let config = stored_config(&conn)?;
204 let mapping = Mapping::new(&config)?;
205 if !mapping.has_id() {
206 bail!(
207 "this index has no `id` field in its mapping, so records cannot be matched for replacement; rebuild instead"
208 );
209 }
210 let tx = conn.transaction()?;
211 let mut ingest = Ingest::new(&tx, &mapping, opts, true)?;
212 for s in sources {
213 ingest.source(s, config.source.sql.as_deref())?;
214 }
215 let mut report = ingest.finish()?;
216 finish_incremental(&tx, &mut report)?;
217 tx.commit()?;
218 report.built = true;
219 report.index_path = index_path.to_path_buf();
220 report.source_bytes = source::byte_size(sources);
221 report.index_bytes = fs::metadata(index_path)?.len();
222 report.elapsed_seconds = round1(started.elapsed().as_secs_f64());
223 Ok(report)
224}
225
226pub fn delete(index_path: &Path, ids: &[String]) -> Result<IngestReport> {
228 let started = Instant::now();
229 let mut conn = open_rw(index_path)?;
230 let tx = conn.transaction()?;
231 tx.execute_batch("CREATE TEMP TABLE touched_groups (key TEXT PRIMARY KEY)")?;
232 let mut deltas = FacetDeltas::default();
233 let mut report = IngestReport::default();
234 for id in ids.iter().map(|s| s.trim()).filter(|s| !s.is_empty()) {
235 if let Some((rowid, grp, facets)) = existing(&tx, id)? {
236 tx.prepare_cached("DELETE FROM records WHERE rowid = ?1")?.execute([rowid])?;
237 tx.prepare_cached("DELETE FROM record_fts WHERE rowid = ?1")?.execute([rowid])?;
238 deltas.remove_json(facets.as_deref());
239 if let Some(g) = grp {
240 touch(&tx, &g)?;
241 }
242 report.deleted += 1;
243 }
244 }
245 deltas.apply(&tx)?;
246 finish_incremental(&tx, &mut report)?;
247 tx.commit()?;
248 report.built = true;
249 report.index_path = index_path.to_path_buf();
250 report.index_bytes = fs::metadata(index_path)?.len();
251 report.elapsed_seconds = round1(started.elapsed().as_secs_f64());
252 Ok(report)
253}
254
255fn finish_incremental(tx: &Connection, report: &mut IngestReport) -> Result<()> {
256 rebuild_groups(tx, true)?;
257 report.record_count = count(tx, "records")?;
258 report.group_count = count(tx, "groups")?;
259 set_meta(tx, "updated_at", &crate::now_rfc3339())?;
260 refresh_counts(tx)
261}
262
263fn refresh_counts(tx: &Connection) -> Result<()> {
264 let (min, max): (Option<String>, Option<String>) =
265 tx.query_row("SELECT MIN(date), MAX(date) FROM records", [], |r| Ok((r.get(0)?, r.get(1)?)))?;
266 set_meta(tx, "record_count", &count(tx, "records")?.to_string())?;
267 set_meta(tx, "group_count", &count(tx, "groups")?.to_string())?;
268 set_meta(tx, "date_min", &min.unwrap_or_default())?;
269 set_meta(tx, "date_max", &max.unwrap_or_default())
270}
271
272#[derive(Default)]
273struct FacetDeltas(HashMap<(String, String), (String, i64)>);
274
275impl FacetDeltas {
276 fn add(&mut self, field: &str, value: &str, delta: i64) {
277 let entry =
278 self.0.entry((field.to_string(), value.to_lowercase())).or_insert_with(|| (value.to_string(), 0));
279 entry.1 += delta;
280 }
281
282 fn remove_json(&mut self, facets: Option<&str>) {
283 let pairs: Vec<(String, String)> =
284 facets.and_then(|f| serde_json::from_str(f).ok()).unwrap_or_default();
285 for (f, v) in pairs {
286 self.add(&f, &v, -1);
287 }
288 }
289
290 fn apply(&mut self, conn: &Connection) -> Result<()> {
291 let mut upsert = conn.prepare_cached(
292 "INSERT INTO facets (field, key, value, count) VALUES (?1, ?2, ?3, ?4) \
293 ON CONFLICT (field, key) DO UPDATE SET count = count + excluded.count",
294 )?;
295 for ((field, key), (value, delta)) in self.0.drain() {
296 if delta != 0 {
297 upsert.execute(params![field, key, value, delta])?;
298 }
299 }
300 conn.execute("DELETE FROM facets WHERE count <= 0", [])?;
301 Ok(())
302 }
303}
304
305struct Ingest<'a> {
306 conn: &'a Connection,
307 mapping: &'a Mapping,
308 opts: BuildOptions,
309 report: IngestReport,
310 deltas: FacetDeltas,
311 incremental: bool,
312}
313
314impl<'a> Ingest<'a> {
315 fn new(
316 conn: &'a Connection,
317 mapping: &'a Mapping,
318 opts: BuildOptions,
319 incremental: bool,
320 ) -> Result<Self> {
321 conn.execute_batch("CREATE TEMP TABLE IF NOT EXISTS touched_groups (key TEXT PRIMARY KEY)")?;
322 Ok(Self {
323 conn,
324 mapping,
325 opts,
326 report: IngestReport::default(),
327 deltas: FacetDeltas::default(),
328 incremental,
329 })
330 }
331
332 fn source(&mut self, source: &Source, sql: Option<&str>) -> Result<()> {
333 log(self.opts, &format!("reading {}", source.path.display()));
334 let label = source.label();
335 source::read(source, sql, &mut |item| {
336 match item {
337 Item::Bad { line, error } => self.skip(&label, line, &error)?,
338 Item::Record { line, value, raw } => self.record(&label, line, &value, raw)?,
339 }
340 Ok(true)
341 })
342 }
343
344 fn skip(&mut self, label: &str, line: u64, why: &str) -> Result<()> {
345 let message = format!("{label}:{line}: {why}");
346 if self.opts.strict {
347 bail!("unusable record (strict mode) {message}");
348 }
349 self.report.skipped_lines += 1;
350 if self.report.skipped_examples.len() < 5 {
351 self.report.skipped_examples.push(message);
352 }
353 Ok(())
354 }
355
356 fn record(
357 &mut self,
358 label: &str,
359 line: u64,
360 value: &serde_json::Value,
361 raw: Option<String>,
362 ) -> Result<()> {
363 let p = self.mapping.prepare(value);
364 let id = match (&p.id, self.mapping.has_id()) {
365 (Some(id), _) => id.clone(),
366 (None, false) => format!("{label}:{line}"),
367 (None, true) => {
368 let field = self.mapping.config.fields.id.as_deref().unwrap_or_default();
369 return self.skip(label, line, &format!("missing id field `{field}`"));
370 }
371 };
372 self.report.records_read += 1;
373 let doc = match raw {
374 Some(raw) => raw,
375 None => serde_json::to_string(value)?,
376 };
377 let facets_json = (!p.facets.is_empty()).then(|| serde_json::to_string(&p.facets)).transpose()?;
378 let row = params![id, p.group, p.group_name, p.date, p.boost, facets_json, doc];
379 let inserted = self
380 .conn
381 .prepare_cached(
382 "INSERT OR IGNORE INTO records (id, grp, grp_name, date, boost, facets, doc) \
383 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
384 )?
385 .execute(row)?;
386 let rowid = if inserted == 1 {
387 self.report.inserted += 1;
388 self.conn.last_insert_rowid()
389 } else {
390 let (rowid, old_group, old_facets) = existing(self.conn, &id)?.expect("conflicting row exists");
391 self.conn
392 .prepare_cached(
393 "UPDATE records SET grp = ?2, grp_name = ?3, date = ?4, boost = ?5, facets = ?6, doc = ?7 \
394 WHERE id = ?1",
395 )?
396 .execute(row)?;
397 self.conn.prepare_cached("DELETE FROM record_fts WHERE rowid = ?1")?.execute([rowid])?;
398 self.deltas.remove_json(old_facets.as_deref());
399 if self.incremental
400 && let Some(g) = old_group
401 {
402 touch(self.conn, &g)?;
403 }
404 self.report.updated += 1;
405 rowid
406 };
407 self.conn
408 .prepare_cached(
409 "INSERT INTO record_fts (rowid, title, body, names, tags) VALUES (?1, ?2, ?3, ?4, ?5)",
410 )?
411 .execute(params![rowid, p.title, p.body, p.names, p.tags])?;
412 for (f, v) in &p.facets {
413 self.deltas.add(f, v, 1);
414 }
415 if self.incremental
416 && let Some(g) = &p.group
417 {
418 touch(self.conn, g)?;
419 }
420 if self.report.records_read.is_multiple_of(100_000) {
421 log(self.opts, &format!(" ...{} records", self.report.records_read));
422 }
423 Ok(())
424 }
425
426 fn finish(mut self) -> Result<IngestReport> {
427 self.deltas.apply(self.conn)?;
428 Ok(self.report)
429 }
430}
431
432type Existing = (i64, Option<String>, Option<String>);
433
434fn existing(conn: &Connection, id: &str) -> Result<Option<Existing>> {
435 Ok(conn
436 .prepare_cached("SELECT rowid, grp, facets FROM records WHERE id = ?1")?
437 .query_row([id], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))
438 .optional()?)
439}
440
441fn touch(conn: &Connection, key: &str) -> Result<()> {
442 conn.prepare_cached("INSERT OR IGNORE INTO temp.touched_groups (key) VALUES (?1)")?.execute([key])?;
443 Ok(())
444}
445
446fn rebuild_groups(conn: &Connection, only_touched: bool) -> Result<()> {
449 let (filter, scope) = if only_touched {
450 (
451 "AND grp IN (SELECT key FROM temp.touched_groups)",
452 "WHERE key IN (SELECT key FROM temp.touched_groups)",
453 )
454 } else {
455 ("", "")
456 };
457 if only_touched {
458 conn.execute(&format!("UPDATE groups SET record_count = 0 {scope}"), [])?;
459 }
460 conn.execute(
461 &format!(
462 "INSERT INTO groups (key, name, record_count) \
463 SELECT grp, grp_name, n FROM ( \
464 SELECT grp, grp_name, MAX(COALESCE(date, '')), COUNT(*) AS n \
465 FROM records WHERE grp IS NOT NULL {filter} GROUP BY grp) WHERE true \
466 ON CONFLICT (key) DO UPDATE SET \
467 name = COALESCE(excluded.name, groups.name), \
468 record_count = excluded.record_count"
469 ),
470 [],
471 )?;
472 conn.execute("DELETE FROM groups WHERE record_count = 0", [])?;
473 let rows: Vec<(String, Option<String>)> = conn
474 .prepare(&format!("SELECT key, name FROM groups {scope}"))?
475 .query_map([], |r| Ok((r.get(0)?, r.get(1)?)))?
476 .collect::<rusqlite::Result<_>>()?;
477 let mut update = conn.prepare("UPDATE groups SET key_lower = ?2, name_normalized = ?3 WHERE key = ?1")?;
478 for (key, name) in rows {
479 update.execute(params![key, key.to_lowercase(), name.map(|n| normalize_name(&n))])?;
480 }
481 conn.execute("DELETE FROM temp.touched_groups", [])?;
482 Ok(())
483}
484
485fn set_meta(conn: &Connection, key: &str, value: &str) -> Result<()> {
486 conn.prepare_cached("INSERT OR REPLACE INTO meta (key, value) VALUES (?1, ?2)")?
487 .execute(params![key, value])?;
488 Ok(())
489}
490
491pub(crate) fn meta(conn: &Connection, key: &str) -> Result<Option<String>> {
492 Ok(conn
493 .prepare_cached("SELECT value FROM meta WHERE key = ?1")?
494 .query_row([key], |r| r.get(0))
495 .optional()?)
496}
497
498pub fn stored_config(conn: &Connection) -> Result<Config> {
500 let raw = meta(conn, "config")?.context("index has no stored mapping")?;
501 serde_json::from_str(&raw).context("stored mapping is unreadable")
502}
503
504fn index_is_current(index_path: &Path, manifest: &str) -> bool {
505 let Ok(conn) = open_ro(index_path) else { return false };
506 let get = |key: &str| meta(&conn, key).ok().flatten();
507 get("schema_version").as_deref() == Some(SCHEMA_VERSION)
508 && get("source_manifest").as_deref() == Some(manifest)
509}
510
511pub fn check_schema(conn: &Connection, path: &Path) -> Result<()> {
513 let version = meta(conn, "schema_version")
514 .with_context(|| format!("{} is not a Leviathan index", path.display()))?;
515 if version.as_deref() != Some(SCHEMA_VERSION) {
516 bail!(
517 "index {} has schema {:?}; this build reads schema {SCHEMA_VERSION}. Rebuild it with `leviathan index --force`",
518 path.display(),
519 version.unwrap_or_default()
520 );
521 }
522 Ok(())
523}
524
525fn open_rw(path: &Path) -> Result<Connection> {
526 if !path.exists() {
527 bail!("no index at {}; run `leviathan index` first", path.display());
528 }
529 let conn = Connection::open(path)?;
530 conn.busy_timeout(std::time::Duration::from_secs(30))?;
531 check_schema(&conn, path)?;
532 Ok(conn)
533}
534
535pub(crate) fn open_ro(path: &Path) -> Result<Connection> {
536 Ok(Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX)?)
537}
538
539fn count(conn: &Connection, table: &str) -> Result<i64> {
540 Ok(conn.query_row(&format!("SELECT COUNT(*) FROM {table}"), [], |r| r.get(0))?)
541}
542
543fn sibling(path: &Path, suffix: &str) -> PathBuf {
544 let mut name = path.file_name().unwrap_or_default().to_os_string();
545 name.push(suffix);
546 path.with_file_name(name)
547}
548
549fn round1(v: f64) -> f64 {
550 (v * 10.0).round() / 10.0
551}
552
553fn log(opts: BuildOptions, message: &str) {
554 if !opts.quiet {
555 eprintln!("[leviathan] {message}");
556 }
557}