1use 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#[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 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#[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}