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