1#![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
124pub const SPEC_VERSION: &str = "13.5";
127
128pub const VERSION: &str = env!("CARGO_PKG_VERSION");
130
131pub const FORMAT_MARKDOWN: &str = "markdown";
133
134pub 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#[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 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
177 Self::open_with_minter(path, Box::new(RandomMinter))
178 }
179
180 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 pub fn open_in_memory() -> Result<Self> {
194 Self::from_connection(Connection::open_in_memory()?, Box::new(RandomMinter))
195 }
196
197 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 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 #[must_use]
222 pub fn conn(&self) -> &Connection {
223 &self.conn
224 }
225
226 pub fn minter(&mut self) -> Mint<'_> {
228 self.ids.at(&self.conn)
229 }
230
231 pub fn mint(&mut self, prefix: &str) -> Result<String> {
233 self.minter().mint(prefix)
234 }
235
236 pub fn transaction(&self) -> Result<Transaction<'_>> {
238 Ok(self.conn.unchecked_transaction()?)
239 }
240
241 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 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 pub fn user_version(&self) -> Result<i64> {
266 schema::user_version(&self.conn)
267 }
268
269 pub fn put_blob(&self, text: &str) -> Result<String> {
273 writers::put_blob(&self.conn, text)
274 }
275
276 pub fn put_tree_node(&self, entries: &[TreeEntry]) -> Result<String> {
278 writers::put_tree_node(&self.conn, entries)
279 }
280
281 pub fn write_block_tree(&self, blocks: &[TreeInputBlock]) -> Result<String> {
283 writers::write_block_tree(&self.conn, blocks)
284 }
285
286 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 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 pub fn reconstruct(&self, doc_id: &str) -> Result<Option<String>> {
300 read::reconstruct(&self.conn, doc_id)
301 }
302
303 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 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 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 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 pub fn properties(&self, doc_id: &str) -> Result<Vec<PropertyRow>> {
334 properties::read_doc_properties(&self.conn, doc_id)
335 }
336
337 pub fn properties_grouped(&self, doc_id: &str) -> Result<serde_json::Value> {
340 Ok(omgbase_properties::grouped(&self.properties(doc_id)?))
341 }
342
343 pub fn properties_merged(&self, doc_id: &str) -> Result<serde_json::Value> {
345 Ok(omgbase_properties::merged(&self.properties(doc_id)?))
346 }
347
348 pub fn rebuild_doc_edges(&self, doc_id: &str) -> Result<()> {
353 graph::rebuild_doc_edges(&self.conn, doc_id)
354 }
355
356 pub fn rebuild_sections(&self, doc_id: &str) -> Result<()> {
360 derived::rebuild_sections(&self.conn, doc_id)
361 }
362
363 pub fn rebuild_index(&self, target: RebuildTarget) -> Result<()> {
365 derived::rebuild_index(&self.conn, target)
366 }
367
368 pub fn gc(&self, enabled: bool) -> Result<GcResult> {
370 derived::run_gc(&self.conn, enabled)
371 }
372
373 pub fn gc_dry_run(&self) -> Result<GcResult> {
375 derived::gc_dry_run(&self.conn)
376 }
377
378 pub fn sweep_pool(&self, ts: &str) -> Result<usize> {
380 derived::sweep_pool(&self.conn, ts)
381 }
382
383 pub fn close(self) -> Result<()> {
385 self.conn.close().map_err(|(_, e)| Error::Sqlite(e))
386 }
387}
388
389#[cfg(test)]
390mod tests;