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 begin_upgradeable_transaction(&self) -> StorageBackendResult<()> {
443        self.conn.begin_deferred_transaction()?;
444        Ok(())
445    }
446
447    fn in_transaction(&self) -> bool {
448        self.conn.in_transaction()
449    }
450
451    fn transaction_has_written(&self) -> StorageBackendResult<bool> {
452        Ok(self.conn.transaction_has_written()?)
453    }
454
455    fn change_version(&self) -> StorageBackendResult<Option<u64>> {
456        Ok(self.conn.data_version()?)
457    }
458
459    fn change_version_monitor_is_nonblocking(&self) -> StorageBackendResult<bool> {
460        Ok(self.conn.data_version_monitor_is_nonblocking()?)
461    }
462
463    fn pin_transaction_snapshot(&self) -> StorageBackendResult<()> {
464        self.conn.pin_transaction_snapshot()?;
465        Ok(())
466    }
467
468    fn commit_transaction(&self) -> StorageBackendResult<()> {
469        self.conn.commit_transaction()?;
470        Ok(())
471    }
472
473    fn rollback_transaction(&self) -> StorageBackendResult<()> {
474        self.conn.rollback_transaction()?;
475        Ok(())
476    }
477
478    fn savepoint(&self, name: &str) -> StorageBackendResult<()> {
479        self.conn.savepoint(name)?;
480        Ok(())
481    }
482
483    fn release_savepoint(&self, name: &str) -> StorageBackendResult<()> {
484        self.conn.release_savepoint(name)?;
485        Ok(())
486    }
487
488    fn rollback_to_savepoint(&self, name: &str) -> StorageBackendResult<()> {
489        self.conn.rollback_to_savepoint(name)?;
490        Ok(())
491    }
492}
493
494struct SQLiteKeyValueBatch<'a> {
495    store: &'a SQLiteKeyValueStore,
496    operations: Vec<SQLiteKeyValueBatchOperation>,
497}
498
499impl KeyValueBatch for SQLiteKeyValueBatch<'_> {
500    fn put(&mut self, key: &[u8], value: &[u8]) -> StorageBackendResult<()> {
501        self.operations.push(SQLiteKeyValueBatchOperation::Put(
502            key.to_vec(),
503            value.to_vec(),
504        ));
505        Ok(())
506    }
507
508    fn delete(&mut self, key: &[u8]) -> StorageBackendResult<()> {
509        self.operations
510            .push(SQLiteKeyValueBatchOperation::Delete(key.to_vec()));
511        Ok(())
512    }
513
514    fn delete_prefix(&mut self, prefix: &[u8]) -> StorageBackendResult<()> {
515        self.operations
516            .push(SQLiteKeyValueBatchOperation::DeletePrefix(prefix.to_vec()));
517        Ok(())
518    }
519
520    fn commit(self: Box<Self>) -> StorageBackendResult<()> {
521        self.store.ensure_table()?;
522        self.store.conn.with_mut(|conn| {
523            let tx = conn.savepoint()?;
524            for operation in self.operations {
525                match operation {
526                    SQLiteKeyValueBatchOperation::Put(key, value) => {
527                        tx.execute(
528                            &format!(
529                                "INSERT OR REPLACE INTO {KEY_VALUE_TABLE} (key, value)
530                                 VALUES (?1, ?2)"
531                            ),
532                            params![key, value],
533                        )?;
534                    }
535                    SQLiteKeyValueBatchOperation::Delete(key) => {
536                        tx.execute(
537                            &format!("DELETE FROM {KEY_VALUE_TABLE} WHERE key = ?1"),
538                            params![key],
539                        )?;
540                    }
541                    SQLiteKeyValueBatchOperation::DeletePrefix(prefix) => {
542                        if let Some(upper) = prefix_upper_bound(&prefix) {
543                            tx.execute(
544                                &format!(
545                                    "DELETE FROM {KEY_VALUE_TABLE}
546                                     WHERE key >= ?1 AND key < ?2"
547                                ),
548                                params![prefix, upper],
549                            )?;
550                        } else {
551                            tx.execute(
552                                &format!("DELETE FROM {KEY_VALUE_TABLE} WHERE key >= ?1"),
553                                params![prefix],
554                            )?;
555                        }
556                    }
557                }
558            }
559            tx.commit()?;
560            Ok(())
561        })?;
562        Ok(())
563    }
564}
565
566/// Shared `SQLite` `KeyValue` storage handle with catalog and backend factories.
567#[derive(Clone)]
568pub struct SQLiteKeyValueStorage {
569    store: Arc<SQLiteKeyValueStore>,
570}
571
572pub type SQLiteKeyValueCatalog = KeyValueCatalog;
573pub type SQLiteKeyValueStorageBackend = KeyValueStorageBackend;
574
575impl SQLiteKeyValueStorage {
576    pub fn open(path: &Path) -> SQLiteResult<Self> {
577        Ok(Self {
578            store: Arc::new(SQLiteKeyValueStore::open(path)?),
579        })
580    }
581
582    pub fn open_in_memory() -> SQLiteResult<Self> {
583        Ok(Self {
584            store: Arc::new(SQLiteKeyValueStore::open_in_memory()?),
585        })
586    }
587
588    pub fn from_connection(conn: ManagedConnection) -> SQLiteResult<Self> {
589        Ok(Self {
590            store: Arc::new(SQLiteKeyValueStore::new(conn)?),
591        })
592    }
593
594    pub fn store(&self) -> Arc<SQLiteKeyValueStore> {
595        Arc::clone(&self.store)
596    }
597
598    pub fn catalog(&self) -> SQLiteKeyValueCatalog {
599        let store: Arc<dyn KeyValueStore> = self.store.clone();
600        KeyValueCatalog::new(store)
601    }
602
603    pub fn backend(&self) -> SQLiteKeyValueStorageBackend {
604        let store: Arc<dyn KeyValueStore> = self.store.clone();
605        KeyValueStorageBackend::new(store)
606    }
607}
608
609impl PersistentStorageProvider for SQLiteKeyValueStorage {
610    fn open_session(&self) -> StorageBackendResult<PersistentStorageSession> {
611        let store: Arc<dyn KeyValueStore> = Arc::new(self.store.new_session());
612        let catalog: Arc<dyn CatalogFacade> = Arc::new(KeyValueCatalog::new(Arc::clone(&store)));
613        let backend: Arc<dyn PersistentStorageBackend> =
614            Arc::new(KeyValueStorageBackend::new(store));
615        Ok(PersistentStorageSession::new(catalog, backend))
616    }
617
618    fn storage_identity(
619        &self,
620    ) -> StorageBackendResult<Option<uqa_storage::PersistentStorageIdentity>> {
621        let connection = self.store.connection();
622        let Some(path) = connection.database_path() else {
623            return Ok(None);
624        };
625        uqa_storage::PersistentStorageIdentity::for_database_path(path)
626            .map(Some)
627            .map_err(|error| {
628                uqa_storage::StorageBackendError::Other(format!(
629                    "resolve SQLite key/value database identity `{}`: {error}",
630                    path.display()
631                ))
632            })
633    }
634}
635
636#[cfg(test)]
637mod tests {
638    use std::collections::BTreeMap;
639
640    use uqa_analysis::standard_analyzer;
641    use uqa_core::Value;
642    use uqa_storage::catalog::{ColumnStatsInput, TableSchema};
643    use uqa_storage::{
644        CatalogFacade, PersistentStorageBackend, VectorIndexOpenMode, VectorIndexSpec,
645    };
646
647    use super::*;
648
649    #[test]
650    fn sqlite_key_value_store_round_trips_and_reopens() {
651        let dir = tempfile::tempdir().unwrap();
652        let path = dir.path().join("keyvalue.sqlite3");
653        {
654            let store = SQLiteKeyValueStore::open(&path).unwrap();
655            store.put(b"apple/1", b"red").unwrap();
656            store.put(b"apple/2", b"green").unwrap();
657            store.put(b"banana/1", b"yellow").unwrap();
658            store.put(&[0x10, 0xff, 0x01], b"binary-prefix").unwrap();
659            store.put(&[0x11, 0x00], b"binary-neighbour").unwrap();
660            assert_eq!(store.get(b"apple/1").unwrap().as_deref(), Some(&b"red"[..]));
661            assert_eq!(store.scan_prefix(b"apple/").unwrap().len(), 2);
662            assert_eq!(
663                store
664                    .scan_prefix_keys_after(b"apple/", Some(b"apple/1"), 1)
665                    .unwrap(),
666                vec![b"apple/2".to_vec()]
667            );
668            assert_eq!(
669                store
670                    .scan_prefix_keys_after(b"apple/", Some(b"a"), 2)
671                    .unwrap(),
672                vec![b"apple/1".to_vec(), b"apple/2".to_vec()]
673            );
674            assert!(store
675                .scan_prefix_keys_after(b"apple/", Some(b"z"), 2)
676                .unwrap()
677                .is_empty());
678            assert_eq!(
679                store.first_prefix_after(b"apple/", Some(b"a")).unwrap(),
680                Some((b"apple/1".to_vec(), b"red".to_vec()))
681            );
682            assert!(store
683                .scan_prefix_keys_after(b"apple/", None, 0)
684                .unwrap()
685                .is_empty());
686            assert_eq!(store.scan_prefix(&[0x10, 0xff]).unwrap().len(), 1);
687            store.delete_prefix(&[0x10, 0xff]).unwrap();
688            assert_eq!(
689                store.get(&[0x11, 0x00]).unwrap().as_deref(),
690                Some(&b"binary-neighbour"[..])
691            );
692        }
693        {
694            let store = SQLiteKeyValueStore::open(&path).unwrap();
695            assert_eq!(
696                store.get(b"apple/2").unwrap().as_deref(),
697                Some(&b"green"[..])
698            );
699            store.delete_prefix(b"apple/").unwrap();
700            assert!(store.get(b"apple/1").unwrap().is_none());
701            assert_eq!(
702                store.get(b"banana/1").unwrap().as_deref(),
703                Some(&b"yellow"[..])
704            );
705        }
706    }
707
708    #[test]
709    fn sqlite_store_passes_the_reusable_backend_contract() {
710        let directory = tempfile::tempdir().unwrap();
711        let reader =
712            SQLiteKeyValueStore::open(&directory.path().join("conformance.sqlite3")).unwrap();
713        let writer = reader.new_session();
714        uqa_storage::key_value::conformance::verify_store(&reader).unwrap();
715        uqa_storage::key_value::conformance::verify_session_isolation(&reader, &writer).unwrap();
716    }
717
718    #[test]
719    fn sqlite_key_value_batch_is_atomic() {
720        let store = SQLiteKeyValueStore::open_in_memory().unwrap();
721        let mut batch = store.batch();
722        batch.put(b"k1", b"v1").unwrap();
723        batch.put(b"k2", b"v2").unwrap();
724        batch.commit().unwrap();
725        assert_eq!(store.get(b"k1").unwrap().as_deref(), Some(&b"v1"[..]));
726
727        store.begin_transaction().unwrap();
728        store.put(b"k3", b"v3").unwrap();
729        store.rollback_transaction().unwrap();
730        assert!(store.get(b"k3").unwrap().is_none());
731    }
732
733    #[test]
734    fn sqlite_key_value_storage_supports_existing_store_contracts() {
735        let storage = SQLiteKeyValueStorage::open_in_memory().unwrap();
736        let backend = storage.backend();
737
738        let mut docs = backend.document_store("articles");
739        docs.put(
740            1,
741            BTreeMap::from([("title".to_string(), Value::Str("rust search".into()))]),
742        )
743        .unwrap();
744        assert_eq!(
745            docs.get_field(1, "title").unwrap(),
746            Some(Value::Str("rust search".into()))
747        );
748
749        let mut index = backend.inverted_index("articles", standard_analyzer("english"));
750        index
751            .add_document(1, BTreeMap::from([("title".into(), "rust search".into())]))
752            .unwrap();
753        assert_eq!(index.doc_freq("title", "rust").unwrap(), 1);
754
755        let mut vectors = backend
756            .vector_index(
757                "articles",
758                "embedding",
759                2,
760                VectorIndexSpec::BruteForce,
761                VectorIndexOpenMode::Create,
762            )
763            .unwrap();
764        vectors.add(1, vec![1.0, 0.0]).unwrap();
765        vectors.add(2, vec![0.0, 1.0]).unwrap();
766        let hits = vectors.search_knn(&[1.0, 0.0], 1).unwrap();
767        assert_eq!(hits.entries()[0].doc_id, 1);
768    }
769
770    #[test]
771    fn sqlite_key_value_catalog_supports_existing_registry_contracts() {
772        let storage = SQLiteKeyValueStorage::open_in_memory().unwrap();
773        let catalog = storage.catalog();
774        catalog.set_metadata("schema_version", "keyvalue").unwrap();
775        catalog.save_schema("public").unwrap();
776        catalog
777            .save_table(&TableSchema {
778                relation: uqa_storage::RelationIdentity::new("public", "docs"),
779                role_owner: "uqa".into(),
780                acl: None,
781                column_acls: std::collections::BTreeMap::default(),
782                object_id: [1; 16],
783                storage_generation: [1; 16],
784                analyzer_json: "{}".into(),
785                fts_fields: vec!["title".into()],
786                vector_fields: Vec::new(),
787                columns_json: "[]".into(),
788                constraints_json: String::new(),
789            })
790            .unwrap();
791        catalog
792            .save_analyzer("ko", "{\"name\":\"standard\"}")
793            .unwrap();
794        catalog
795            .save_table_field_analyzer("docs", "title", "index", "ko")
796            .unwrap();
797        catalog
798            .save_foreign_server("fs", "memory", "{\"root\":\"/tmp\"}")
799            .unwrap();
800        catalog
801            .save_catalog_index(
802                &uqa_storage::RelationIdentity::new("public", "idx_docs_title"),
803                "gin",
804                "public.docs",
805                "[\"title\"]",
806                "{}",
807            )
808            .unwrap();
809        catalog
810            .save_column_stats(ColumnStatsInput::basic(
811                "docs",
812                "title",
813                4,
814                0,
815                Some("a"),
816                Some("z"),
817                10,
818            ))
819            .unwrap();
820
821        assert_eq!(
822            catalog.get_metadata("schema_version").unwrap().as_deref(),
823            Some("keyvalue")
824        );
825        assert_eq!(
826            catalog.load_tables().unwrap()[0].relation.qualified_name(),
827            "public.docs"
828        );
829        assert_eq!(catalog.load_analyzers().unwrap()[0].0, "ko");
830        assert_eq!(catalog.load_foreign_servers().unwrap()[0].0, "fs");
831        assert_eq!(
832            catalog.load_catalog_indexes().unwrap()[0]
833                .relation
834                .qualified_name(),
835            "public.idx_docs_title"
836        );
837        assert_eq!(
838            catalog.load_column_stats("docs").unwrap()[0].distinct_count,
839            4
840        );
841    }
842}