Skip to main content

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