Skip to main content

omgbase_store/
derived.rs

1//! Derived tables (`spec/store/README.md` §3.2, §4.5, §7): sections, the
2//! FTS5 index over block text, rebuild and garbage collection, and the pool
3//! sweep.
4
5use std::collections::HashSet;
6
7use omgbase_format::hash::hex;
8use rusqlite::{Connection, OptionalExtension, params};
9
10use crate::error::Result;
11use crate::tree::{from_hex, parse_tree_entries};
12
13// ---- sections (§4.5) ------------------------------------------------------------
14
15/// Rebuild `sections` for one document over its top-level live blocks.
16pub fn rebuild_sections(conn: &Connection, doc_id: &str) -> Result<()> {
17    conn.execute("DELETE FROM sections WHERE doc_id = ?1", params![doc_id])?;
18    let tops: Vec<(String, i64, String, String)> = {
19        let mut stmt = conn.prepare(
20            "SELECT block_id, ordinal, type, attrs FROM blocks
21             WHERE doc_id = ?1 AND parent_block IS NULL AND deleted_commit IS NULL
22             ORDER BY ordinal",
23        )?;
24        let rows = stmt.query_map(params![doc_id], |r| {
25            Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?))
26        })?;
27        rows.collect::<std::result::Result<Vec<_>, _>>()?
28    };
29    let max_ordinal = tops.last().map_or(-1, |t| t.1);
30    let headings: Vec<(&str, i64, i64)> = tops
31        .iter()
32        .filter(|t| t.2 == "heading")
33        .map(|t| {
34            let level = serde_json::from_str::<serde_json::Value>(&t.3)
35                .ok()
36                .and_then(|v| v.get("level").and_then(serde_json::Value::as_i64))
37                .unwrap_or(1);
38            (t.0.as_str(), t.1, level)
39        })
40        .collect();
41    let mut insert = conn.prepare(
42        "INSERT INTO sections (heading_block, doc_id, level, first_ordinal, last_ordinal)
43         VALUES (?1, ?2, ?3, ?4, ?5)",
44    )?;
45    for (i, &(block_id, ordinal, level)) in headings.iter().enumerate() {
46        let last = headings[i + 1..]
47            .iter()
48            .find(|h| h.2 <= level)
49            .map_or(max_ordinal, |h| h.1 - 1);
50        insert.execute(params![block_id, doc_id, level, ordinal, last])?;
51    }
52    Ok(())
53}
54
55// ---- FTS (external content, spec/search §1.1) -----------------------------------------
56
57/// SQL predicate: the `blocks` row aliased `b` is a live **leaf** — no live
58/// row of the same document names it as `parent_block` (search 1.2: a
59/// container's text is its children's text joined, so only leaves are
60/// indexed). `c.doc_id = b.doc_id` lets the subquery use `idx_blocks_doc`.
61pub const LIVE_LEAF_SQL: &str = "b.deleted_commit IS NULL AND NOT EXISTS (SELECT 1 FROM blocks c WHERE c.doc_id = b.doc_id AND c.parent_block = b.block_id AND c.deleted_commit IS NULL)";
62
63fn live_leaf_rows(conn: &Connection, doc_id: &str) -> Result<Vec<(i64, String)>> {
64    let mut stmt = conn.prepare_cached(&format!(
65        "SELECT b.rowid, b.text FROM blocks b WHERE b.doc_id = ?1 AND {LIVE_LEAF_SQL}"
66    ))?;
67    let it = stmt.query_map(params![doc_id], |r| Ok((r.get(0)?, r.get(1)?)))?;
68    Ok(it.collect::<std::result::Result<Vec<_>, _>>()?)
69}
70
71/// Remove a document's live leaf rows from `blocks_fts` (external-content
72/// index: a `'delete'` command with the original text, before the rows
73/// change). Only live leaves are indexed, so only those are deleted — a
74/// `'delete'` for an unindexed row skews the index statistics.
75pub fn fts_delete_doc(conn: &Connection, doc_id: &str) -> Result<()> {
76    let mut del =
77        conn.prepare("INSERT INTO blocks_fts(blocks_fts, rowid, text) VALUES('delete', ?1, ?2)")?;
78    for (rowid, text) in live_leaf_rows(conn, doc_id)? {
79        del.execute(params![rowid, text])?;
80    }
81    Ok(())
82}
83
84/// Index a document's live leaf rows.
85pub fn fts_index_doc(conn: &Connection, doc_id: &str) -> Result<()> {
86    let mut ins = conn.prepare("INSERT INTO blocks_fts(rowid, text) VALUES(?1, ?2)")?;
87    for (rowid, text) in live_leaf_rows(conn, doc_id)? {
88        ins.execute(params![rowid, text])?;
89    }
90    Ok(())
91}
92
93/// Rebuild the whole index: `'delete-all'`, then every live leaf of every
94/// live document (FTS5's own `'rebuild'` reads the content table wholesale,
95/// containers included, and is not used).
96pub fn fts_rebuild(conn: &Connection) -> Result<()> {
97    conn.execute_batch("INSERT INTO blocks_fts(blocks_fts) VALUES('delete-all')")?;
98    let doc_ids: Vec<String> = {
99        let mut stmt = conn.prepare("SELECT doc_id FROM docs WHERE deleted_commit IS NULL")?;
100        let it = stmt.query_map([], |r| r.get(0))?;
101        it.collect::<std::result::Result<Vec<_>, _>>()?
102    };
103    for id in &doc_ids {
104        fts_index_doc(conn, id)?;
105    }
106    Ok(())
107}
108
109/// A `blocks` row about to be evicted out of another document's row set
110/// (spec/store §5.4 step 8).
111pub(crate) struct EvictRow {
112    pub rowid: i64,
113    pub block_id: String,
114    pub doc_id: String,
115    pub parent_block: Option<String>,
116    pub text: String,
117    pub live: bool,
118}
119
120/// Keep the index equal to the table's live leaves across the deletion of
121/// `rows` (the caller deletes them right after). Leaf-ness is decided over
122/// the table **before** any row goes, so a list moving with its items is
123/// order-independent (the items lose their entries; the list never had one),
124/// and a live parent that keeps its row but loses its last live child becomes
125/// a leaf and gains an entry.
126pub(crate) fn fts_before_evict_rows(conn: &Connection, rows: &[EvictRow]) -> Result<()> {
127    if rows.is_empty() {
128        return Ok(());
129    }
130    let evicted: HashSet<&str> = rows.iter().map(|r| r.block_id.as_str()).collect();
131    let mut has_live_child = conn.prepare_cached(
132        "SELECT 1 FROM blocks c WHERE c.doc_id = ?1 AND c.parent_block = ?2 AND c.deleted_commit IS NULL LIMIT 1",
133    )?;
134    let mut leaves = Vec::new();
135    for r in rows.iter().filter(|r| r.live) {
136        let child: Option<i64> = has_live_child
137            .query_row(params![r.doc_id, r.block_id], |row| row.get(0))
138            .optional()?;
139        if child.is_none() {
140            leaves.push(r);
141        }
142    }
143    let mut del =
144        conn.prepare("INSERT INTO blocks_fts(blocks_fts, rowid, text) VALUES('delete', ?1, ?2)")?;
145    for r in &leaves {
146        del.execute(params![r.rowid, r.text])?;
147    }
148    // Parents that stay live but are left childless are leaves from now on.
149    let mut parent_row = conn.prepare(
150        "SELECT b.rowid, b.text FROM blocks b WHERE b.block_id = ?1 AND b.doc_id = ?2 AND b.deleted_commit IS NULL
151           AND NOT EXISTS (SELECT 1 FROM blocks c WHERE c.doc_id = b.doc_id AND c.parent_block = b.block_id
152                           AND c.deleted_commit IS NULL AND c.block_id NOT IN (SELECT value FROM json_each(?3)))",
153    )?;
154    let evicted_json = serde_json::to_string(&rows.iter().map(|r| &r.block_id).collect::<Vec<_>>())
155        .expect("ids serialize");
156    let mut ins = conn.prepare("INSERT INTO blocks_fts(rowid, text) VALUES(?1, ?2)")?;
157    let mut seen: HashSet<&str> = HashSet::new();
158    for r in rows.iter().filter(|r| r.live) {
159        let Some(parent) = r.parent_block.as_deref() else {
160            continue;
161        };
162        if evicted.contains(parent) || !seen.insert(parent) {
163            continue;
164        }
165        let p: Option<(i64, String)> = parent_row
166            .query_row(params![parent, r.doc_id, evicted_json], |row| {
167                Ok((row.get(0)?, row.get(1)?))
168            })
169            .optional()?;
170        if let Some((rowid, text)) = p {
171            ins.execute(params![rowid, text])?;
172        }
173    }
174    Ok(())
175}
176
177// ---- rebuild (§7) -------------------------------------------------------------------
178
179/// What [`rebuild_index`] recomputes.
180#[derive(Clone, Copy, Debug, PartialEq, Eq)]
181pub enum RebuildTarget {
182    Sections,
183    /// The `doc_edges` rollup of the open `edges` rows (`spec/graph` §3.4).
184    Edges,
185    Fts,
186    BlockChanges,
187    All,
188}
189
190/// Recompute derived tables from the durable ones: `sections` per live doc,
191/// `doc_edges` per live doc (the rollup of its open edges), the FTS index
192/// (`'delete-all'` + every live leaf, `spec/search` §1.1), `block_changes`
193/// from `dispositions`.
194pub fn rebuild_index(conn: &Connection, target: RebuildTarget) -> Result<()> {
195    let tx = conn.unchecked_transaction()?;
196    let live_docs = |tx: &Connection| -> Result<Vec<String>> {
197        let mut stmt = tx.prepare("SELECT doc_id FROM docs WHERE deleted_commit IS NULL")?;
198        let it = stmt.query_map([], |r| r.get(0))?;
199        Ok(it.collect::<std::result::Result<Vec<_>, _>>()?)
200    };
201    if matches!(target, RebuildTarget::Sections | RebuildTarget::All) {
202        for id in &live_docs(&tx)? {
203            rebuild_sections(&tx, id)?;
204        }
205    }
206    if matches!(target, RebuildTarget::Edges | RebuildTarget::All) {
207        for id in &live_docs(&tx)? {
208            crate::graph::rebuild_doc_edges(&tx, id)?;
209        }
210    }
211    if matches!(target, RebuildTarget::Fts | RebuildTarget::All) {
212        fts_rebuild(&tx)?;
213    }
214    if matches!(target, RebuildTarget::BlockChanges | RebuildTarget::All) {
215        tx.execute_batch(
216            "DELETE FROM block_changes;
217             INSERT OR IGNORE INTO block_changes (block_id, commit_id, kind)
218               SELECT block_id, commit_id, kind FROM dispositions;",
219        )?;
220    }
221    tx.commit()?;
222    Ok(())
223}
224
225// ---- garbage (§7) -------------------------------------------------------------------
226
227/// Rows a garbage collection sweeps (or would sweep).
228#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
229pub struct GcResult {
230    pub blobs_swept: usize,
231    pub tree_nodes_swept: usize,
232}
233
234/// The mark phase: tree nodes reachable from every `revisions.root_tree`
235/// (walking `child_tree_hash_hex`) and the blobs they name (raw, trivia)
236/// plus every `frontmatter_blob`. Hex sets.
237pub fn mark_reachable(conn: &Connection) -> Result<(HashSet<String>, HashSet<String>)> {
238    let mut trees = HashSet::new();
239    let mut blobs = HashSet::new();
240    {
241        let mut stmt = conn.prepare(
242            "SELECT DISTINCT frontmatter_blob FROM revisions WHERE frontmatter_blob IS NOT NULL",
243        )?;
244        for row in stmt.query_map([], |r| r.get::<_, Vec<u8>>(0))? {
245            blobs.insert(hex(&row?));
246        }
247    }
248    let roots: Vec<Vec<u8>> = {
249        let mut stmt = conn.prepare("SELECT DISTINCT root_tree FROM revisions")?;
250        let it = stmt.query_map([], |r| r.get(0))?;
251        it.collect::<std::result::Result<Vec<_>, _>>()?
252    };
253    let mut node = conn.prepare("SELECT entries FROM tree_nodes WHERE hash = ?1")?;
254    let mut stack: Vec<String> = roots.iter().map(|h| hex(h)).collect();
255    while let Some(h) = stack.pop() {
256        if !trees.insert(h.clone()) {
257            continue;
258        }
259        let entries: Option<String> = {
260            use rusqlite::OptionalExtension;
261            node.query_row(params![from_hex(&h)?], |r| r.get(0))
262                .optional()?
263        };
264        let Some(text) = entries else {
265            continue;
266        };
267        for e in parse_tree_entries(&text)? {
268            blobs.insert(e.raw_hash_hex);
269            if let Some(t) = e.trivia_hash_hex {
270                blobs.insert(t);
271            }
272            if let Some(c) = e.child_tree_hash_hex {
273                stack.push(c);
274            }
275        }
276    }
277    Ok((trees, blobs))
278}
279
280/// `(tree node hashes, blob hashes)` that no revision reaches.
281type Unreachable = (Vec<Vec<u8>>, Vec<Vec<u8>>);
282
283fn unreachable_rows(conn: &Connection) -> Result<Unreachable> {
284    let (trees, blobs) = mark_reachable(conn)?;
285    let all_trees: Vec<Vec<u8>> = {
286        let mut stmt = conn.prepare("SELECT hash FROM tree_nodes")?;
287        let it = stmt.query_map([], |r| r.get(0))?;
288        it.collect::<std::result::Result<Vec<_>, _>>()?
289    };
290    let all_blobs: Vec<Vec<u8>> = {
291        let mut stmt = conn.prepare("SELECT hash FROM blobs")?;
292        let it = stmt.query_map([], |r| r.get(0))?;
293        it.collect::<std::result::Result<Vec<_>, _>>()?
294    };
295    Ok((
296        all_trees
297            .into_iter()
298            .filter(|h| !trees.contains(&hex(h)))
299            .collect(),
300        all_blobs
301            .into_iter()
302            .filter(|h| !blobs.contains(&hex(h)))
303            .collect(),
304    ))
305}
306
307/// What a collection would sweep, without deleting anything (§8 I5).
308pub fn gc_dry_run(conn: &Connection) -> Result<GcResult> {
309    let (trees, blobs) = unreachable_rows(conn)?;
310    Ok(GcResult {
311        blobs_swept: blobs.len(),
312        tree_nodes_swept: trees.len(),
313    })
314}
315
316/// Mark-and-sweep, flag-gated (`enabled = false` is a no-op). Because
317/// revisions are never pruned, nothing is unreachable in practice.
318pub fn run_gc(conn: &Connection, enabled: bool) -> Result<GcResult> {
319    if !enabled {
320        return Ok(GcResult::default());
321    }
322    let tx = conn.unchecked_transaction()?;
323    let (trees, blobs) = unreachable_rows(&tx)?;
324    for h in &trees {
325        tx.execute("DELETE FROM tree_nodes WHERE hash = ?1", params![h])?;
326    }
327    for h in &blobs {
328        tx.execute("DELETE FROM blobs WHERE hash = ?1", params![h])?;
329    }
330    tx.commit()?;
331    Ok(GcResult {
332        blobs_swept: blobs.len(),
333        tree_nodes_swept: trees.len(),
334    })
335}
336
337/// §5.5: `DELETE FROM resurrection_pool WHERE expires_ts <= ts`; rows deleted.
338pub fn sweep_pool(conn: &Connection, ts: &str) -> Result<usize> {
339    Ok(conn.execute(
340        "DELETE FROM resurrection_pool WHERE expires_ts <= ?1",
341        params![ts],
342    )?)
343}