Skip to main content

omgbase_store/
lib.rs

1//! # omgbase-store
2//!
3//! The omgbase store, Rust implementation: the embedded SQLite database that
4//! owns block **identity**, **history** and the **current state** of a
5//! repository of authored files. Files stay the source of truth for content;
6//! the store adds stable block ids, an append-only history of revisions and
7//! commits, the dispositions the matcher recorded, and the derived indexes
8//! the query surfaces read. The contract is `spec/store/README.md` in the
9//! omgbase repository — `schema.sql` (embedded verbatim as [`SCHEMA_SQL`]) is
10//! the DDL, `spec/store/cases/*.json` the executable fixtures — and this
11//! crate opens the same databases the TypeScript reference writes.
12//!
13//! ```
14//! use omgbase_reconcile::Config;
15//! use omgbase_store::{BatchItem, Store};
16//!
17//! let mut store = Store::open_in_memory()?;
18//! let repo = store.create_repo("notes")?;
19//! let ts = "2026-09-26T10:00:00.000Z";
20//! let items = [BatchItem::observed("a.md", "# Title\n\nFirst.\n")];
21//! let outcomes = store.observe_batch(&repo, &items, ts, &Config::default())?;
22//! let o = outcomes[0].as_observed().unwrap();
23//! assert!(o.converged && !o.echo);
24//! assert_eq!(o.dispositions["inserted"], 2);
25//! assert_eq!(store.reconstruct(&o.doc_id)?.as_deref(), Some("# Title\n\nFirst.\n"));
26//!
27//! // Re-observing what the store holds is an echo: no commit, no mint.
28//! let again = store.observe_batch(&repo, &items, ts, &Config::default())?;
29//! assert!(again[0].as_observed().unwrap().echo);
30//! # Ok::<(), omgbase_store::Error>(())
31//! ```
32//!
33//! ## Layering
34//!
35//! [`schema`] (§1, §3: the opener and migrations), [`ids`] (§2.1–2.2),
36//! [`time`] (§2.4), [`tree`] (§4.1 encodings), [`order_key`] (§4.3),
37//! [`writers`] (blobs, tree nodes, commits, revisions), [`observe`] (§5),
38//! [`properties`] (the `properties` rows of `spec/properties`, written in
39//! §5.4), [`history`] (the `changes_since` feed of `spec/sync` §6), [`graph`] (the `nodes`, `external_nodes`, `edges` and `doc_edges`
40//! rows of `spec/graph`, written in §5.4), [`search`] (`spec/search`:
41//! `text_search`, the embedding drain over the `embeddings`/`doc_embeddings`
42//! caches, vector, hybrid and `resolve`), [`read`] (§5.2, §6), [`derived`]
43//! (§4.5 sections, FTS, §7 rebuild and GC), and the `spec/mutate` host side
44//! (store 13.4): [`mutate`] (loading the working tree, `apply` with the
45//! file-CAS + known-id `api` commit), [`doc_store`] (the write-target seam),
46//! [`macros`], [`links`] (link-destination rewriting), [`docs_ops`] (create /
47//! move / delete / set_meta, [`yaml_emit`]) and [`plan`] (the whole-document
48//! update planner).
49//!
50//! Minted ids are opaque; the store asks its [`IdMinter`] for each one. The
51//! default is the CSPRNG-backed [`RandomMinter`]; a fixture runner installs a
52//! [`SequentialMinter`] through [`Store::open_in_memory_with_minter`]. Every
53//! draw passes through the checking [`Mint`] (§2.1, 13.5): an id already
54//! issued in this process or naming a row of its prefix's table is redrawn.
55
56#![forbid(unsafe_code)]
57
58pub mod derived;
59pub mod doc_store;
60pub mod docs_ops;
61pub mod error;
62pub mod graph;
63pub mod history;
64pub mod ids;
65pub mod links;
66pub mod macros;
67pub mod mint;
68pub mod mutate;
69pub mod observe;
70pub mod order_key;
71pub mod plan;
72pub mod properties;
73pub mod read;
74pub mod schema;
75pub mod search;
76pub mod time;
77pub mod tree;
78pub mod writers;
79pub mod yaml_emit;
80
81use std::path::Path;
82
83use rusqlite::functions::{Context, FunctionFlags};
84use rusqlite::types::ValueRef;
85use rusqlite::{Connection, Transaction, params};
86
87pub use derived::{GcResult, LIVE_LEAF_SQL, RebuildTarget};
88pub use doc_store::{DocStore, FsDocStore, MemDocStore, NullDocStore};
89pub use docs_ops::{DocMoveResult, DocOpContext, DocOpResult, Retargeted};
90pub use error::{Error, Result};
91pub use graph::ResolvedEdge;
92pub use history::{ChangesPage, CommitDigest, DigestRevision};
93pub use ids::{IdMinter, RandomMinter, RepeatingMinter, SequentialMinter, is_valid_id, prefix_of};
94pub use links::InboundLink;
95pub use macros::{LinkRepair, LinkRepairCount, LinkRepairPlan, RetargetHit};
96pub use mint::{Deferred, MINT_GIVE_UP_AFTER, Mint, id_in_use};
97pub use mutate::{
98    ApplyOrigin, ApplyRequest, ApplyResult, Diff, DocInfo, Revision, SetFrontmatter,
99    find_doc_by_ref, is_id_ref, load_mut_doc,
100};
101pub use observe::{
102    BatchItem, BatchOutcome, Committed, DeleteOutcome, ObserveOutcome, has_conflict_markers,
103};
104pub use omgbase_graph::{EdgeDescriptor, ProjectedNode};
105pub use omgbase_mutate::{
106    self as mutate_kernel, At, Expect, MutBlock, MutDoc, MutationError, Op, OpResult, Opset,
107    Parent, PlanOp, To,
108};
109pub use omgbase_properties::PropertyRow;
110pub use omgbase_reconcile::{Config, MatchBlock, PoolEntry};
111pub use omgbase_search::{
112    Boosts, DocEmbedBlockRef, DocEmbedMethod, DocEmbedTask, EmbedTask, EmbeddingProvider, Evidence,
113};
114pub use read::RevisionRead;
115pub use schema::{SCHEMA_SQL, SCHEMA_VERSION};
116pub use search::{
117    BlockContexts, ContextScope, DocEmbedStats, DocVectorHit, DocVectorRow, DrainStats, EmbedStats,
118    ForeignVectors, HybridHit, HybridQuery, QueryVector, ResolveHit, TextHit, TextSearchResult,
119    VectorHit, block_vector,
120};
121pub use tree::{TreeEntry, canonical_attrs, canonical_json, serialize_tree_entries, tree_hash};
122pub use writers::{NewCommit, NewRevision, Origin, TreeInputBlock};
123
124/// The `spec/store/VERSION` this crate implements (`major.minor`); the major
125/// is [`SCHEMA_VERSION`].
126pub const SPEC_VERSION: &str = "13.5";
127
128/// This crate's own version (`Cargo.toml`), for the surface's `version` tool.
129pub const VERSION: &str = env!("CARGO_PKG_VERSION");
130
131/// The one format this crate ingests (`docs.format`).
132pub const FORMAT_MARKDOWN: &str = "markdown";
133
134/// An open store: one connection, one id minter (behind the in-use check of
135/// [`mint`]). All writes go through it.
136pub struct Store {
137    conn: Connection,
138    ids: mint::IdSource,
139}
140
141impl std::fmt::Debug for Store {
142    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
143        f.debug_struct("Store").finish_non_exhaustive()
144    }
145}
146
147/// Cosine similarity of two little-endian float32 vector blobs
148/// (`spec/search` §3, the reference's `cosineFloat32`): products and sums in
149/// f64 over the shorter length (each blob floored to a multiple of 4 bytes);
150/// `0.0` when either norm is zero.
151#[must_use]
152pub fn cosine_bytes(a: &[u8], b: &[u8]) -> f64 {
153    omgbase_search::cosine_bytes(a, b)
154}
155
156fn cosine_udf(ctx: &Context<'_>) -> rusqlite::Result<Option<f64>> {
157    let blob = |i: usize| -> rusqlite::Result<Option<&[u8]>> {
158        match ctx.get_raw(i) {
159            ValueRef::Null => Ok(None),
160            ValueRef::Blob(b) => Ok(Some(b)),
161            other => Err(rusqlite::Error::InvalidFunctionParameterType(
162                i,
163                other.data_type(),
164            )),
165        }
166    };
167    match (blob(0)?, blob(1)?) {
168        (Some(a), Some(b)) => Ok(Some(cosine_bytes(a, b))),
169        _ => Ok(None),
170    }
171}
172
173impl Store {
174    /// Open (creating parent directories and the file as needed) with the
175    /// production minter. `:memory:` is a valid path.
176    pub fn open(path: impl AsRef<Path>) -> Result<Self> {
177        Self::open_with_minter(path, Box::new(RandomMinter))
178    }
179
180    /// [`Store::open`] with a replaceable minter (§2.2).
181    pub fn open_with_minter(path: impl AsRef<Path>, minter: Box<dyn IdMinter>) -> Result<Self> {
182        let path = path.as_ref();
183        if path.as_os_str() != ":memory:" {
184            if let Some(dir) = path.parent().filter(|d| !d.as_os_str().is_empty()) {
185                std::fs::create_dir_all(dir)
186                    .map_err(|e| Error::Other(format!("cannot create {}: {e}", dir.display())))?;
187            }
188        }
189        Self::from_connection(Connection::open(path)?, minter)
190    }
191
192    /// A fresh in-memory store (tests, fixtures).
193    pub fn open_in_memory() -> Result<Self> {
194        Self::from_connection(Connection::open_in_memory()?, Box::new(RandomMinter))
195    }
196
197    /// A fresh in-memory store with a replaceable minter.
198    pub fn open_in_memory_with_minter(minter: Box<dyn IdMinter>) -> Result<Self> {
199        Self::from_connection(Connection::open_in_memory()?, minter)
200    }
201
202    /// Adopt an existing connection: apply the §1 pragmas, register
203    /// `cosine`, run the opener (fresh → `schema.sql`; older → migrations;
204    /// newer → [`Error::SchemaTooNew`]).
205    pub fn from_connection(conn: Connection, minter: Box<dyn IdMinter>) -> Result<Self> {
206        conn.pragma_update(None, "journal_mode", "WAL")?;
207        conn.pragma_update(None, "synchronous", "NORMAL")?;
208        conn.pragma_update(None, "foreign_keys", "ON")?;
209        conn.create_scalar_function(
210            "cosine",
211            2,
212            FunctionFlags::SQLITE_UTF8 | FunctionFlags::SQLITE_DETERMINISTIC,
213            cosine_udf,
214        )?;
215        let mut ids = mint::IdSource::new(minter);
216        schema::migrate(&conn, &mut ids)?;
217        Ok(Self { conn, ids })
218    }
219
220    /// The underlying connection.
221    #[must_use]
222    pub fn conn(&self) -> &Connection {
223        &self.conn
224    }
225
226    /// The store's checking minter (§2.1) over its connection.
227    pub fn minter(&mut self) -> Mint<'_> {
228        self.ids.at(&self.conn)
229    }
230
231    /// Mint an id with `prefix` (§2.1): never one in use.
232    pub fn mint(&mut self, prefix: &str) -> Result<String> {
233        self.minter().mint(prefix)
234    }
235
236    /// Begin a write transaction (one per commit, §1); rolls back on drop.
237    pub fn transaction(&self) -> Result<Transaction<'_>> {
238        Ok(self.conn.unchecked_transaction()?)
239    }
240
241    /// Create a repo with `slug` (**mints `rp`**); returns its id.
242    pub fn create_repo(&mut self, slug: &str) -> Result<String> {
243        let repo_id = self.mint("rp")?;
244        self.conn.execute(
245            "INSERT INTO repos (repo_id, slug) VALUES (?1, ?2)",
246            params![repo_id, slug],
247        )?;
248        Ok(repo_id)
249    }
250
251    /// The repo id for `slug`, if any.
252    pub fn repo_by_slug(&self, slug: &str) -> Result<Option<String>> {
253        use rusqlite::OptionalExtension;
254        Ok(self
255            .conn
256            .query_row(
257                "SELECT repo_id FROM repos WHERE slug = ?1",
258                params![slug],
259                |r| r.get(0),
260            )
261            .optional()?)
262    }
263
264    /// `PRAGMA user_version`.
265    pub fn user_version(&self) -> Result<i64> {
266        schema::user_version(&self.conn)
267    }
268
269    // ---- writers ------------------------------------------------------------------
270
271    /// [`writers::put_blob`].
272    pub fn put_blob(&self, text: &str) -> Result<String> {
273        writers::put_blob(&self.conn, text)
274    }
275
276    /// [`writers::put_tree_node`].
277    pub fn put_tree_node(&self, entries: &[TreeEntry]) -> Result<String> {
278        writers::put_tree_node(&self.conn, entries)
279    }
280
281    /// [`writers::write_block_tree`].
282    pub fn write_block_tree(&self, blocks: &[TreeInputBlock]) -> Result<String> {
283        writers::write_block_tree(&self.conn, blocks)
284    }
285
286    /// [`writers::new_commit`]; returns `(commit_id, seq)`.
287    pub fn new_commit(&mut self, input: &NewCommit<'_>) -> Result<(String, i64)> {
288        writers::new_commit(&self.conn, &mut self.ids.at(&self.conn), input)
289    }
290
291    /// [`writers::write_revision`]; returns `(rev_id, seq)`.
292    pub fn write_revision(&mut self, input: &NewRevision<'_>) -> Result<(String, i64)> {
293        writers::write_revision(&self.conn, &mut self.ids.at(&self.conn), input)
294    }
295
296    // ---- reads --------------------------------------------------------------------
297
298    /// §6.1: the current bytes of a live doc; `None` when tombstoned or unknown.
299    pub fn reconstruct(&self, doc_id: &str) -> Result<Option<String>> {
300        read::reconstruct(&self.conn, doc_id)
301    }
302
303    /// §6.2: the bytes at a revision, from its Merkle root.
304    pub fn read_at_revision(&self, doc_id: &str, rev_id: &str) -> Result<Option<RevisionRead>> {
305        read::read_at_revision(&self.conn, doc_id, rev_id)
306    }
307
308    /// The change feed (`spec/sync` §6): the repo's commits with `seq >
309    /// cursor`, optionally of one `origin`, `limit + 1` fetched to set
310    /// `truncated`.
311    pub fn changes_since(
312        &self,
313        repo_id: &str,
314        cursor: i64,
315        limit: usize,
316        origin: Option<&str>,
317    ) -> Result<ChangesPage> {
318        history::changes_since(&self.conn, repo_id, cursor, limit, origin)
319    }
320
321    /// §5.2: the matcher's old side for a doc, from the live `blocks` rows.
322    pub fn load_old_match_blocks(&self, doc_id: &str) -> Result<Vec<MatchBlock>> {
323        read::load_old_match_blocks(&self.conn, doc_id)
324    }
325
326    /// §5.1: the repo's unexpired pool at `ts`.
327    pub fn load_pool(&self, repo_id: &str, ts: &str) -> Result<Vec<PoolEntry>> {
328        read::load_pool(&self.conn, repo_id, ts)
329    }
330
331    /// The document's live `properties` rows (`spec/properties` §1), in
332    /// write order.
333    pub fn properties(&self, doc_id: &str) -> Result<Vec<PropertyRow>> {
334        properties::read_doc_properties(&self.conn, doc_id)
335    }
336
337    /// `spec/properties` §5 **grouped**: `{frontmatter, inline, computed}`
338    /// (the `docs_read.properties` shape).
339    pub fn properties_grouped(&self, doc_id: &str) -> Result<serde_json::Value> {
340        Ok(omgbase_properties::grouped(&self.properties(doc_id)?))
341    }
342
343    /// `spec/properties` §5 **merged**: `{key: shape}` over every row.
344    pub fn properties_merged(&self, doc_id: &str) -> Result<serde_json::Value> {
345        Ok(omgbase_properties::merged(&self.properties(doc_id)?))
346    }
347
348    // ---- graph (spec/graph) -----------------------------------------------------------
349
350    /// `spec/graph` §3.4: recompute one document's `doc_edges` from its open
351    /// edges.
352    pub fn rebuild_doc_edges(&self, doc_id: &str) -> Result<()> {
353        graph::rebuild_doc_edges(&self.conn, doc_id)
354    }
355
356    // ---- derived ------------------------------------------------------------------
357
358    /// §4.5: rebuild `sections` for one document.
359    pub fn rebuild_sections(&self, doc_id: &str) -> Result<()> {
360        derived::rebuild_sections(&self.conn, doc_id)
361    }
362
363    /// §7: recompute derived tables.
364    pub fn rebuild_index(&self, target: RebuildTarget) -> Result<()> {
365        derived::rebuild_index(&self.conn, target)
366    }
367
368    /// §7: mark-and-sweep garbage collection, flag-gated (off → no-op).
369    pub fn gc(&self, enabled: bool) -> Result<GcResult> {
370        derived::run_gc(&self.conn, enabled)
371    }
372
373    /// What [`Store::gc`] would sweep, without deleting (§8 I5).
374    pub fn gc_dry_run(&self) -> Result<GcResult> {
375        derived::gc_dry_run(&self.conn)
376    }
377
378    /// §5.5: delete expired pool rows; returns how many.
379    pub fn sweep_pool(&self, ts: &str) -> Result<usize> {
380        derived::sweep_pool(&self.conn, ts)
381    }
382
383    /// Close the connection.
384    pub fn close(self) -> Result<()> {
385        self.conn.close().map_err(|(_, e)| Error::Sqlite(e))
386    }
387}
388
389#[cfg(test)]
390mod tests;