Skip to main content

leviathan/
index.rs

1//! Build and update the on-disk index.
2//!
3//! A full build writes `<index>.building` and renames it over the live index
4//! only on success, so a running server never sees a partial database.
5//! Builds are skipped when the sources (paths, sizes, mtimes) and the mapping
6//! are unchanged.
7
8use 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    /// Existing ids replaced (in a full build: duplicate ids in the sources).
85    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    /// Set when no config or flags were given and the mapping was inferred.
97    #[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    /// Fail on the first unusable record instead of counting and skipping it.
109    pub strict: bool,
110    pub quiet: bool,
111}
112
113/// Full rebuild into a fresh database, atomically swapped into place.
114pub 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
198/// Insert or replace records (by id) in an existing index, using the mapping
199/// it was built with.
200pub 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
226/// Remove records by id.
227pub 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
446/// Recompute group rows from their records. The name comes from each group's
447/// most recent record (SQLite's bare-column-with-MAX rule).
448fn 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
498/// The mapping an index was built with.
499pub 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
511/// Open an existing index read-only, refusing other schema versions.
512pub 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}