Skip to main content

uqa_storage_sqlite/
lib.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! `SQLite`-backed physical `KeyValue` storage.
8//!
9//! The logical catalog, document store, inverted index, and vector index live
10//! in `uqa-storage::key_value`. This crate provides the `SQLite` implementation
11//! of the ordered byte-key store they require.
12
13use std::path::Path;
14use std::sync::atomic::{AtomicBool, Ordering};
15use std::sync::Arc;
16
17use rusqlite::{params, OptionalExtension};
18use uqa_storage::key_value::{
19    prefix_upper_bound, KeyValueBatch, KeyValueCatalog, KeyValueStorageBackend, KeyValueStore,
20};
21use uqa_storage::sqlite::{ManagedConnection, Result as SQLiteResult, SQLiteError};
22use uqa_storage::{
23    CatalogFacade, PersistentStorageBackend, PersistentStorageIdentity, PersistentStorageProvider,
24    PersistentStorageSession, StorageBackendError, StorageBackendResult,
25};
26
27const KEY_VALUE_TABLE: &str = "_key_value";
28
29#[derive(Debug, Clone)]
30enum SQLiteKeyValueBatchOperation {
31    Put(Vec<u8>, Vec<u8>),
32    Delete(Vec<u8>),
33    DeletePrefix(Vec<u8>),
34}
35
36/// `SQLite` physical implementation of [`KeyValueStore`].
37#[derive(Clone)]
38pub struct SQLiteKeyValueStore {
39    conn: ManagedConnection,
40    table_ready: Arc<AtomicBool>,
41}
42
43impl SQLiteKeyValueStore {
44    pub fn open(path: &Path) -> SQLiteResult<Self> {
45        Self::new(ManagedConnection::open(path)?)
46    }
47
48    pub fn open_in_memory() -> SQLiteResult<Self> {
49        Self::new(ManagedConnection::open_in_memory()?)
50    }
51
52    pub fn new(conn: ManagedConnection) -> SQLiteResult<Self> {
53        let store = Self {
54            conn,
55            table_ready: Arc::new(AtomicBool::new(false)),
56        };
57        store.ensure_table()?;
58        Ok(store)
59    }
60
61    pub fn connection(&self) -> ManagedConnection {
62        self.conn.clone()
63    }
64
65    /// Create an isolated transaction session over the same `SQLite` database.
66    pub fn new_session(&self) -> Self {
67        Self {
68            conn: self.conn.new_session(),
69            table_ready: Arc::clone(&self.table_ready),
70        }
71    }
72
73    fn ensure_table(&self) -> SQLiteResult<()> {
74        if self.table_ready.load(Ordering::Acquire) {
75            return Ok(());
76        }
77        self.conn.with(|conn| {
78            conn.execute(
79                &format!(
80                    "CREATE TABLE IF NOT EXISTS {KEY_VALUE_TABLE} (
81                        key   BLOB PRIMARY KEY,
82                        value BLOB NOT NULL
83                    ) WITHOUT ROWID"
84                ),
85                [],
86            )?;
87            Ok(())
88        })?;
89        self.table_ready.store(true, Ordering::Release);
90        Ok(())
91    }
92}
93
94impl KeyValueStore for SQLiteKeyValueStore {
95    fn storage_identity(&self) -> StorageBackendResult<Option<PersistentStorageIdentity>> {
96        let Some(path) = self.conn.database_path() else {
97            return Ok(None);
98        };
99        PersistentStorageIdentity::for_database_path(path).map(Some)
100    }
101
102    fn open_session(&self) -> StorageBackendResult<Arc<dyn KeyValueStore>> {
103        Ok(Arc::new(self.new_session()))
104    }
105
106    fn get(&self, key: &[u8]) -> StorageBackendResult<Option<Vec<u8>>> {
107        self.ensure_table()?;
108        Ok(self.conn.with(|conn| {
109            conn.query_row(
110                &format!("SELECT value FROM {KEY_VALUE_TABLE} WHERE key = ?1"),
111                params![key],
112                |row| row.get::<_, Vec<u8>>(0),
113            )
114            .optional()
115            .map_err(SQLiteError::from)
116        })?)
117    }
118
119    fn contains_key(&self, key: &[u8]) -> StorageBackendResult<bool> {
120        self.ensure_table()?;
121        Ok(self.conn.with(|conn| {
122            conn.query_row(
123                &format!("SELECT 1 FROM {KEY_VALUE_TABLE} WHERE key = ?1 LIMIT 1"),
124                params![key],
125                |_| Ok(()),
126            )
127            .optional()
128            .map(|value| value.is_some())
129            .map_err(SQLiteError::from)
130        })?)
131    }
132
133    fn put(&self, key: &[u8], value: &[u8]) -> StorageBackendResult<()> {
134        self.ensure_table()?;
135        self.conn.with(|conn| {
136            conn.execute(
137                &format!("INSERT OR REPLACE INTO {KEY_VALUE_TABLE} (key, value) VALUES (?1, ?2)"),
138                params![key, value],
139            )?;
140            Ok(())
141        })?;
142        Ok(())
143    }
144
145    fn delete(&self, key: &[u8]) -> StorageBackendResult<()> {
146        self.ensure_table()?;
147        self.conn.with(|conn| {
148            conn.execute(
149                &format!("DELETE FROM {KEY_VALUE_TABLE} WHERE key = ?1"),
150                params![key],
151            )?;
152            Ok(())
153        })?;
154        Ok(())
155    }
156
157    fn scan_prefix(&self, prefix: &[u8]) -> StorageBackendResult<Vec<(Vec<u8>, Vec<u8>)>> {
158        self.ensure_table()?;
159        let upper = prefix_upper_bound(prefix);
160        self.conn
161            .with(|conn| {
162                let mut rows = Vec::new();
163                if let Some(upper) = upper {
164                    let mut stmt = conn.prepare(&format!(
165                        "SELECT key, value FROM {KEY_VALUE_TABLE}
166                         WHERE key >= ?1 AND key < ?2
167                         ORDER BY key"
168                    ))?;
169                    let iter = stmt.query_map(params![prefix, upper], |row| {
170                        Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, Vec<u8>>(1)?))
171                    })?;
172                    for row in iter {
173                        rows.push(row?);
174                    }
175                } else {
176                    let mut stmt = conn.prepare(&format!(
177                        "SELECT key, value FROM {KEY_VALUE_TABLE}
178                         WHERE key >= ?1
179                         ORDER BY key"
180                    ))?;
181                    let iter = stmt.query_map(params![prefix], |row| {
182                        Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, Vec<u8>>(1)?))
183                    })?;
184                    for row in iter {
185                        rows.push(row?);
186                    }
187                }
188                Ok(rows)
189            })
190            .map_err(StorageBackendError::from)
191    }
192
193    fn scan_prefix_after(
194        &self,
195        prefix: &[u8],
196        after: Option<&[u8]>,
197        limit: usize,
198    ) -> StorageBackendResult<Vec<(Vec<u8>, Vec<u8>)>> {
199        if limit == 0 {
200            return Ok(Vec::new());
201        }
202        self.ensure_table()?;
203        let upper = prefix_upper_bound(prefix);
204        let after = after.filter(|after| *after >= prefix);
205        let limit = i64::try_from(limit).map_err(|_| {
206            StorageBackendError::Other(format!(
207                "key/value cursor limit {limit} is outside SQLite's integer range"
208            ))
209        })?;
210        self.conn
211            .with(|connection| {
212                let mut output = Vec::new();
213                match (after, upper) {
214                    (Some(after), Some(upper)) => {
215                        let mut statement = connection.prepare_cached(&format!(
216                            "SELECT key, value FROM {KEY_VALUE_TABLE}
217                             WHERE key > ?1 AND key < ?2
218                             ORDER BY key LIMIT ?3"
219                        ))?;
220                        let rows = statement.query_map(params![after, upper, limit], |row| {
221                            Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, Vec<u8>>(1)?))
222                        })?;
223                        for row in rows {
224                            output.push(row?);
225                        }
226                    }
227                    (Some(after), None) => {
228                        let mut statement = connection.prepare_cached(&format!(
229                            "SELECT key, value FROM {KEY_VALUE_TABLE}
230                             WHERE key > ?1
231                             ORDER BY key LIMIT ?2"
232                        ))?;
233                        let rows = statement.query_map(params![after, limit], |row| {
234                            Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, Vec<u8>>(1)?))
235                        })?;
236                        for row in rows {
237                            output.push(row?);
238                        }
239                    }
240                    (None, Some(upper)) => {
241                        let mut statement = connection.prepare_cached(&format!(
242                            "SELECT key, value FROM {KEY_VALUE_TABLE}
243                             WHERE key >= ?1 AND key < ?2
244                             ORDER BY key LIMIT ?3"
245                        ))?;
246                        let rows = statement.query_map(params![prefix, upper, limit], |row| {
247                            Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, Vec<u8>>(1)?))
248                        })?;
249                        for row in rows {
250                            output.push(row?);
251                        }
252                    }
253                    (None, None) => {
254                        let mut statement = connection.prepare_cached(&format!(
255                            "SELECT key, value FROM {KEY_VALUE_TABLE}
256                             WHERE key >= ?1
257                             ORDER BY key LIMIT ?2"
258                        ))?;
259                        let rows = statement.query_map(params![prefix, limit], |row| {
260                            Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, Vec<u8>>(1)?))
261                        })?;
262                        for row in rows {
263                            output.push(row?);
264                        }
265                    }
266                }
267                Ok(output)
268            })
269            .map_err(StorageBackendError::from)
270    }
271
272    fn scan_prefix_keys_after(
273        &self,
274        prefix: &[u8],
275        after: Option<&[u8]>,
276        limit: usize,
277    ) -> StorageBackendResult<Vec<Vec<u8>>> {
278        if limit == 0 {
279            return Ok(Vec::new());
280        }
281        self.ensure_table()?;
282        let upper = prefix_upper_bound(prefix);
283        let after = after.filter(|after| *after >= prefix);
284        let limit = i64::try_from(limit).map_err(|_| {
285            StorageBackendError::Other(format!(
286                "key/value cursor limit {limit} is outside SQLite's integer range"
287            ))
288        })?;
289        self.conn
290            .with(|connection| {
291                let mut keys = Vec::new();
292                match (after, upper) {
293                    (Some(after), Some(upper)) => {
294                        let mut stmt = connection.prepare_cached(&format!(
295                            "SELECT key FROM {KEY_VALUE_TABLE}
296                             WHERE key > ?1 AND key < ?2
297                             ORDER BY key LIMIT ?3"
298                        ))?;
299                        let rows =
300                            stmt.query_map(params![after, upper, limit], |row| row.get(0))?;
301                        for key in rows {
302                            keys.push(key?);
303                        }
304                    }
305                    (Some(after), None) => {
306                        let mut stmt = connection.prepare_cached(&format!(
307                            "SELECT key FROM {KEY_VALUE_TABLE}
308                             WHERE key > ?1
309                             ORDER BY key LIMIT ?2"
310                        ))?;
311                        let rows = stmt.query_map(params![after, limit], |row| row.get(0))?;
312                        for key in rows {
313                            keys.push(key?);
314                        }
315                    }
316                    (None, Some(upper)) => {
317                        let mut stmt = connection.prepare_cached(&format!(
318                            "SELECT key FROM {KEY_VALUE_TABLE}
319                             WHERE key >= ?1 AND key < ?2
320                             ORDER BY key LIMIT ?3"
321                        ))?;
322                        let rows =
323                            stmt.query_map(params![prefix, upper, limit], |row| row.get(0))?;
324                        for key in rows {
325                            keys.push(key?);
326                        }
327                    }
328                    (None, None) => {
329                        let mut stmt = connection.prepare_cached(&format!(
330                            "SELECT key FROM {KEY_VALUE_TABLE}
331                             WHERE key >= ?1
332                             ORDER BY key LIMIT ?2"
333                        ))?;
334                        let rows = stmt.query_map(params![prefix, limit], |row| row.get(0))?;
335                        for key in rows {
336                            keys.push(key?);
337                        }
338                    }
339                }
340                Ok(keys)
341            })
342            .map_err(StorageBackendError::from)
343    }
344
345    fn first_prefix_after(
346        &self,
347        prefix: &[u8],
348        after: Option<&[u8]>,
349    ) -> StorageBackendResult<Option<(Vec<u8>, Vec<u8>)>> {
350        self.ensure_table()?;
351        let upper = prefix_upper_bound(prefix);
352        let after = after.filter(|after| *after >= prefix);
353        self.conn
354            .with(|connection| {
355                let row = match (after, upper) {
356                    (Some(after), Some(upper)) => connection
357                        .prepare_cached(&format!(
358                            "SELECT key, value FROM {KEY_VALUE_TABLE}
359                             WHERE key > ?1 AND key < ?2
360                             ORDER BY key LIMIT 1"
361                        ))?
362                        .query_row(params![after, upper], |row| {
363                            Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, Vec<u8>>(1)?))
364                        })
365                        .optional()?,
366                    (Some(after), None) => connection
367                        .prepare_cached(&format!(
368                            "SELECT key, value FROM {KEY_VALUE_TABLE}
369                             WHERE key > ?1
370                             ORDER BY key LIMIT 1"
371                        ))?
372                        .query_row(params![after], |row| {
373                            Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, Vec<u8>>(1)?))
374                        })
375                        .optional()?,
376                    (None, Some(upper)) => connection
377                        .prepare_cached(&format!(
378                            "SELECT key, value FROM {KEY_VALUE_TABLE}
379                             WHERE key >= ?1 AND key < ?2
380                             ORDER BY key LIMIT 1"
381                        ))?
382                        .query_row(params![prefix, upper], |row| {
383                            Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, Vec<u8>>(1)?))
384                        })
385                        .optional()?,
386                    (None, None) => connection
387                        .prepare_cached(&format!(
388                            "SELECT key, value FROM {KEY_VALUE_TABLE}
389                             WHERE key >= ?1
390                             ORDER BY key LIMIT 1"
391                        ))?
392                        .query_row(params![prefix], |row| {
393                            Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, Vec<u8>>(1)?))
394                        })
395                        .optional()?,
396                };
397                Ok(row.filter(|(key, _)| key.starts_with(prefix)))
398            })
399            .map_err(StorageBackendError::from)
400    }
401
402    fn delete_prefix(&self, prefix: &[u8]) -> StorageBackendResult<usize> {
403        self.ensure_table()?;
404        let upper = prefix_upper_bound(prefix);
405        let deleted = self.conn.with(|conn| {
406            let deleted = if let Some(upper) = upper {
407                conn.execute(
408                    &format!(
409                        "DELETE FROM {KEY_VALUE_TABLE}
410                         WHERE key >= ?1 AND key < ?2"
411                    ),
412                    params![prefix, upper],
413                )?
414            } else {
415                conn.execute(
416                    &format!("DELETE FROM {KEY_VALUE_TABLE} WHERE key >= ?1"),
417                    params![prefix],
418                )?
419            };
420            Ok(deleted)
421        })?;
422        Ok(deleted)
423    }
424
425    fn batch(&self) -> Box<dyn KeyValueBatch + '_> {
426        Box::new(SQLiteKeyValueBatch {
427            store: self,
428            operations: Vec::new(),
429        })
430    }
431
432    fn begin_transaction(&self) -> StorageBackendResult<()> {
433        self.conn.begin_transaction()?;
434        Ok(())
435    }
436
437    fn begin_read_transaction(&self) -> StorageBackendResult<()> {
438        self.conn.begin_deferred_transaction()?;
439        Ok(())
440    }
441
442    fn in_transaction(&self) -> bool {
443        self.conn.in_transaction()
444    }
445
446    fn transaction_has_written(&self) -> StorageBackendResult<bool> {
447        Ok(self.conn.transaction_has_written()?)
448    }
449
450    fn change_version(&self) -> StorageBackendResult<Option<u64>> {
451        Ok(self.conn.data_version()?)
452    }
453
454    fn change_version_monitor_is_nonblocking(&self) -> StorageBackendResult<bool> {
455        Ok(self.conn.data_version_monitor_is_nonblocking()?)
456    }
457
458    fn pin_transaction_snapshot(&self) -> StorageBackendResult<()> {
459        self.conn.pin_transaction_snapshot()?;
460        Ok(())
461    }
462
463    fn commit_transaction(&self) -> StorageBackendResult<()> {
464        self.conn.commit_transaction()?;
465        Ok(())
466    }
467
468    fn rollback_transaction(&self) -> StorageBackendResult<()> {
469        self.conn.rollback_transaction()?;
470        Ok(())
471    }
472
473    fn savepoint(&self, name: &str) -> StorageBackendResult<()> {
474        self.conn.savepoint(name)?;
475        Ok(())
476    }
477
478    fn release_savepoint(&self, name: &str) -> StorageBackendResult<()> {
479        self.conn.release_savepoint(name)?;
480        Ok(())
481    }
482
483    fn rollback_to_savepoint(&self, name: &str) -> StorageBackendResult<()> {
484        self.conn.rollback_to_savepoint(name)?;
485        Ok(())
486    }
487}
488
489struct SQLiteKeyValueBatch<'a> {
490    store: &'a SQLiteKeyValueStore,
491    operations: Vec<SQLiteKeyValueBatchOperation>,
492}
493
494impl KeyValueBatch for SQLiteKeyValueBatch<'_> {
495    fn put(&mut self, key: &[u8], value: &[u8]) -> StorageBackendResult<()> {
496        self.operations.push(SQLiteKeyValueBatchOperation::Put(
497            key.to_vec(),
498            value.to_vec(),
499        ));
500        Ok(())
501    }
502
503    fn delete(&mut self, key: &[u8]) -> StorageBackendResult<()> {
504        self.operations
505            .push(SQLiteKeyValueBatchOperation::Delete(key.to_vec()));
506        Ok(())
507    }
508
509    fn delete_prefix(&mut self, prefix: &[u8]) -> StorageBackendResult<()> {
510        self.operations
511            .push(SQLiteKeyValueBatchOperation::DeletePrefix(prefix.to_vec()));
512        Ok(())
513    }
514
515    fn commit(self: Box<Self>) -> StorageBackendResult<()> {
516        self.store.ensure_table()?;
517        self.store.conn.with_mut(|conn| {
518            let tx = conn.savepoint()?;
519            for operation in self.operations {
520                match operation {
521                    SQLiteKeyValueBatchOperation::Put(key, value) => {
522                        tx.execute(
523                            &format!(
524                                "INSERT OR REPLACE INTO {KEY_VALUE_TABLE} (key, value)
525                                 VALUES (?1, ?2)"
526                            ),
527                            params![key, value],
528                        )?;
529                    }
530                    SQLiteKeyValueBatchOperation::Delete(key) => {
531                        tx.execute(
532                            &format!("DELETE FROM {KEY_VALUE_TABLE} WHERE key = ?1"),
533                            params![key],
534                        )?;
535                    }
536                    SQLiteKeyValueBatchOperation::DeletePrefix(prefix) => {
537                        if let Some(upper) = prefix_upper_bound(&prefix) {
538                            tx.execute(
539                                &format!(
540                                    "DELETE FROM {KEY_VALUE_TABLE}
541                                     WHERE key >= ?1 AND key < ?2"
542                                ),
543                                params![prefix, upper],
544                            )?;
545                        } else {
546                            tx.execute(
547                                &format!("DELETE FROM {KEY_VALUE_TABLE} WHERE key >= ?1"),
548                                params![prefix],
549                            )?;
550                        }
551                    }
552                }
553            }
554            tx.commit()?;
555            Ok(())
556        })?;
557        Ok(())
558    }
559}
560
561/// Shared `SQLite` `KeyValue` storage handle with catalog and backend factories.
562#[derive(Clone)]
563pub struct SQLiteKeyValueStorage {
564    store: Arc<SQLiteKeyValueStore>,
565}
566
567pub type SQLiteKeyValueCatalog = KeyValueCatalog;
568pub type SQLiteKeyValueStorageBackend = KeyValueStorageBackend;
569
570impl SQLiteKeyValueStorage {
571    pub fn open(path: &Path) -> SQLiteResult<Self> {
572        Ok(Self {
573            store: Arc::new(SQLiteKeyValueStore::open(path)?),
574        })
575    }
576
577    pub fn open_in_memory() -> SQLiteResult<Self> {
578        Ok(Self {
579            store: Arc::new(SQLiteKeyValueStore::open_in_memory()?),
580        })
581    }
582
583    pub fn from_connection(conn: ManagedConnection) -> SQLiteResult<Self> {
584        Ok(Self {
585            store: Arc::new(SQLiteKeyValueStore::new(conn)?),
586        })
587    }
588
589    pub fn store(&self) -> Arc<SQLiteKeyValueStore> {
590        Arc::clone(&self.store)
591    }
592
593    pub fn catalog(&self) -> SQLiteKeyValueCatalog {
594        let store: Arc<dyn KeyValueStore> = self.store.clone();
595        KeyValueCatalog::new(store)
596    }
597
598    pub fn backend(&self) -> SQLiteKeyValueStorageBackend {
599        let store: Arc<dyn KeyValueStore> = self.store.clone();
600        KeyValueStorageBackend::new(store)
601    }
602}
603
604impl PersistentStorageProvider for SQLiteKeyValueStorage {
605    fn open_session(&self) -> StorageBackendResult<PersistentStorageSession> {
606        let store: Arc<dyn KeyValueStore> = Arc::new(self.store.new_session());
607        let catalog: Arc<dyn CatalogFacade> = Arc::new(KeyValueCatalog::new(Arc::clone(&store)));
608        let backend: Arc<dyn PersistentStorageBackend> =
609            Arc::new(KeyValueStorageBackend::new(store));
610        Ok(PersistentStorageSession::new(catalog, backend))
611    }
612
613    fn storage_identity(
614        &self,
615    ) -> StorageBackendResult<Option<uqa_storage::PersistentStorageIdentity>> {
616        let connection = self.store.connection();
617        let Some(path) = connection.database_path() else {
618            return Ok(None);
619        };
620        uqa_storage::PersistentStorageIdentity::for_database_path(path)
621            .map(Some)
622            .map_err(|error| {
623                uqa_storage::StorageBackendError::Other(format!(
624                    "resolve SQLite key/value database identity `{}`: {error}",
625                    path.display()
626                ))
627            })
628    }
629}
630
631#[cfg(test)]
632mod tests {
633    use std::collections::BTreeMap;
634
635    use uqa_analysis::standard_analyzer;
636    use uqa_core::Value;
637    use uqa_storage::catalog::{ColumnStatsInput, TableSchema};
638    use uqa_storage::{
639        CatalogFacade, PersistentStorageBackend, VectorIndexOpenMode, VectorIndexSpec,
640    };
641
642    use super::*;
643
644    #[test]
645    fn sqlite_key_value_store_round_trips_and_reopens() {
646        let dir = tempfile::tempdir().unwrap();
647        let path = dir.path().join("keyvalue.sqlite3");
648        {
649            let store = SQLiteKeyValueStore::open(&path).unwrap();
650            store.put(b"apple/1", b"red").unwrap();
651            store.put(b"apple/2", b"green").unwrap();
652            store.put(b"banana/1", b"yellow").unwrap();
653            store.put(&[0x10, 0xff, 0x01], b"binary-prefix").unwrap();
654            store.put(&[0x11, 0x00], b"binary-neighbour").unwrap();
655            assert_eq!(store.get(b"apple/1").unwrap().as_deref(), Some(&b"red"[..]));
656            assert_eq!(store.scan_prefix(b"apple/").unwrap().len(), 2);
657            assert_eq!(
658                store
659                    .scan_prefix_keys_after(b"apple/", Some(b"apple/1"), 1)
660                    .unwrap(),
661                vec![b"apple/2".to_vec()]
662            );
663            assert_eq!(
664                store
665                    .scan_prefix_keys_after(b"apple/", Some(b"a"), 2)
666                    .unwrap(),
667                vec![b"apple/1".to_vec(), b"apple/2".to_vec()]
668            );
669            assert!(store
670                .scan_prefix_keys_after(b"apple/", Some(b"z"), 2)
671                .unwrap()
672                .is_empty());
673            assert_eq!(
674                store.first_prefix_after(b"apple/", Some(b"a")).unwrap(),
675                Some((b"apple/1".to_vec(), b"red".to_vec()))
676            );
677            assert!(store
678                .scan_prefix_keys_after(b"apple/", None, 0)
679                .unwrap()
680                .is_empty());
681            assert_eq!(store.scan_prefix(&[0x10, 0xff]).unwrap().len(), 1);
682            store.delete_prefix(&[0x10, 0xff]).unwrap();
683            assert_eq!(
684                store.get(&[0x11, 0x00]).unwrap().as_deref(),
685                Some(&b"binary-neighbour"[..])
686            );
687        }
688        {
689            let store = SQLiteKeyValueStore::open(&path).unwrap();
690            assert_eq!(
691                store.get(b"apple/2").unwrap().as_deref(),
692                Some(&b"green"[..])
693            );
694            store.delete_prefix(b"apple/").unwrap();
695            assert!(store.get(b"apple/1").unwrap().is_none());
696            assert_eq!(
697                store.get(b"banana/1").unwrap().as_deref(),
698                Some(&b"yellow"[..])
699            );
700        }
701    }
702
703    #[test]
704    fn sqlite_store_passes_the_reusable_backend_contract() {
705        let directory = tempfile::tempdir().unwrap();
706        let reader =
707            SQLiteKeyValueStore::open(&directory.path().join("conformance.sqlite3")).unwrap();
708        let writer = reader.new_session();
709        uqa_storage::key_value::conformance::verify_store(&reader).unwrap();
710        uqa_storage::key_value::conformance::verify_session_isolation(&reader, &writer).unwrap();
711    }
712
713    #[test]
714    fn sqlite_key_value_batch_is_atomic() {
715        let store = SQLiteKeyValueStore::open_in_memory().unwrap();
716        let mut batch = store.batch();
717        batch.put(b"k1", b"v1").unwrap();
718        batch.put(b"k2", b"v2").unwrap();
719        batch.commit().unwrap();
720        assert_eq!(store.get(b"k1").unwrap().as_deref(), Some(&b"v1"[..]));
721
722        store.begin_transaction().unwrap();
723        store.put(b"k3", b"v3").unwrap();
724        store.rollback_transaction().unwrap();
725        assert!(store.get(b"k3").unwrap().is_none());
726    }
727
728    #[test]
729    fn sqlite_key_value_storage_supports_existing_store_contracts() {
730        let storage = SQLiteKeyValueStorage::open_in_memory().unwrap();
731        let backend = storage.backend();
732
733        let mut docs = backend.document_store("articles");
734        docs.put(
735            1,
736            BTreeMap::from([("title".to_string(), Value::Str("rust search".into()))]),
737        )
738        .unwrap();
739        assert_eq!(
740            docs.get_field(1, "title").unwrap(),
741            Some(Value::Str("rust search".into()))
742        );
743
744        let mut index = backend.inverted_index("articles", standard_analyzer("english"));
745        index
746            .add_document(1, BTreeMap::from([("title".into(), "rust search".into())]))
747            .unwrap();
748        assert_eq!(index.doc_freq("title", "rust").unwrap(), 1);
749
750        let mut vectors = backend
751            .vector_index(
752                "articles",
753                "embedding",
754                2,
755                VectorIndexSpec::BruteForce,
756                VectorIndexOpenMode::Create,
757            )
758            .unwrap();
759        vectors.add(1, vec![1.0, 0.0]).unwrap();
760        vectors.add(2, vec![0.0, 1.0]).unwrap();
761        let hits = vectors.search_knn(&[1.0, 0.0], 1).unwrap();
762        assert_eq!(hits.entries()[0].doc_id, 1);
763    }
764
765    #[test]
766    fn sqlite_key_value_catalog_supports_existing_registry_contracts() {
767        let storage = SQLiteKeyValueStorage::open_in_memory().unwrap();
768        let catalog = storage.catalog();
769        catalog.set_metadata("schema_version", "keyvalue").unwrap();
770        catalog.save_schema("public").unwrap();
771        catalog
772            .save_table(&TableSchema {
773                relation: uqa_storage::RelationIdentity::new("public", "docs"),
774                analyzer_json: "{}".into(),
775                fts_fields: vec!["title".into()],
776                vector_fields: Vec::new(),
777                columns_json: "[]".into(),
778                constraints_json: String::new(),
779            })
780            .unwrap();
781        catalog
782            .save_analyzer("ko", "{\"name\":\"standard\"}")
783            .unwrap();
784        catalog
785            .save_table_field_analyzer("docs", "title", "index", "ko")
786            .unwrap();
787        catalog
788            .save_foreign_server("fs", "memory", "{\"root\":\"/tmp\"}")
789            .unwrap();
790        catalog
791            .save_catalog_index("idx_docs_title", "gin", "docs", "[\"title\"]", "{}")
792            .unwrap();
793        catalog
794            .save_column_stats(ColumnStatsInput::basic(
795                "docs",
796                "title",
797                4,
798                0,
799                Some("a"),
800                Some("z"),
801                10,
802            ))
803            .unwrap();
804
805        assert_eq!(
806            catalog.get_metadata("schema_version").unwrap().as_deref(),
807            Some("keyvalue")
808        );
809        assert_eq!(
810            catalog.load_tables().unwrap()[0].relation.qualified_name(),
811            "public.docs"
812        );
813        assert_eq!(catalog.load_analyzers().unwrap()[0].0, "ko");
814        assert_eq!(catalog.load_foreign_servers().unwrap()[0].0, "fs");
815        assert_eq!(
816            catalog.load_catalog_indexes().unwrap()[0].name,
817            "idx_docs_title"
818        );
819        assert_eq!(
820            catalog.load_column_stats("docs").unwrap()[0].distinct_count,
821            4
822        );
823    }
824}