1#![allow(
2 clippy::branches_sharing_code,
3 clippy::missing_errors_doc,
4 clippy::match_wildcard_for_single_variants,
5 clippy::assigning_clones,
6 clippy::needless_pass_by_value
7)]
8
9use std::collections::HashMap;
10use std::path::Path;
11use std::sync::{
12 Mutex,
13 atomic::{AtomicU64, Ordering},
14};
15
16use fjall::{Database, Keyspace, KeyspaceCreateOptions, PersistMode};
17use serde::{Deserialize, Serialize};
18use serde_json::Value;
19
20use crate::encryption::{EncryptionConfig, StorageCodec, StorageCrypto, make_codec};
21use crate::error::ThingdResult;
22use crate::model::{
23 AggregateFunction, AggregateGroupResult, AggregateOptions, AggregateResult, IndexDefinition,
24 LinkDirection, LinkQueryOptions, ListEventsOptions, ListObjectsOptions, MigrationRecord,
25 ObjectKey, PutObjectOptions, SchemaOptions, SearchOptions, SortDirection, StoredSchema,
26 TimeBucket, TimeSeriesBucket, TimeSeriesOptions, TimeSeriesResult,
27};
28use crate::replication::{REPLICATION_STATE_COLLECTION, REPLICATION_STREAM};
29use crate::store::{
30 AggregateStore, EventLog, LinkStore, ObjectStore, QueueStore, RetentionOptions,
31 RetentionReport, Searcher,
32};
33use crate::{
34 CollectionSchema, FieldSchema, Link, MemoryEvent, MemoryObject, QueueClaimOptions, QueueJob,
35 QueueJobStatus, QueueNackOptions, SearchHit, ThingdError, VectorSearchHit, VectorSearchOptions,
36};
37use crate::{now_iso_string, unix_timestamp_millis};
38
39pub struct PersistentEngine {
49 #[allow(dead_code)]
50 db: Database,
51 objects: Keyspace,
52 events: Keyspace,
53 event_meta: Keyspace,
54 queue_jobs: Keyspace,
55 ready_jobs: Keyspace,
56 links_by_id: Keyspace,
57 links_from: Keyspace,
58 links_to: Keyspace,
59 schemas: Keyspace,
60 migrations: Keyspace,
61 indexes: Keyspace,
62 next_link_id: AtomicU64,
63 event_seq_counters: HashMap<String, u64>,
64 event_idempotency_keys: HashMap<(String, String), u64>,
65 unique_index_values: HashMap<(String, String, String), ObjectKey>,
66 unique_index_cache_complete: bool,
67 #[cfg(feature = "search")]
68 search_index: Option<tantivy::Index>,
69 #[cfg(feature = "search")]
70 search_writer: Option<Mutex<tantivy::IndexWriter<tantivy::TantivyDocument>>>,
71 #[cfg(feature = "search")]
72 search_reader: Option<Mutex<tantivy::IndexReader>>,
73 #[cfg(feature = "vectors")]
74 vectors: Keyspace,
75 codec: Box<dyn StorageCodec>,
76}
77
78const STORAGE_FORMAT_VERSION: u32 = 1;
79const STORAGE_CONTRACT: &str = "fjall-tantivy-v1";
80const STORAGE_MANIFEST_FILE: &str = ".thingd-storage.json";
81const STORAGE_LOCK_FILE: &str = "lock";
82const STORAGE_KEYSPACES_DIR: &str = "keyspaces";
83const SEARCH_WRITER_MEMORY_BYTES: usize = 15_000_000;
86type UniqueIndexCache = HashMap<(String, String, String), ObjectKey>;
87
88const REQUIRED_KEYSPACES: &[&str] = &[
89 "objects",
90 "events",
91 "queue_jobs",
92 "ready_jobs",
93 "links_by_id",
94 "links_from",
95 "links_to",
96 "schemas",
97 "migrations",
98 "indexes",
99];
100
101#[cfg(feature = "vectors")]
102const REQUIRED_VECTOR_KEYSPACE: &str = "vectors";
103
104#[derive(Debug, Clone, Serialize, Deserialize)]
105struct StorageManifest {
106 format_version: u32,
107 contract: String,
108 keyspaces: Vec<String>,
109 search_schema_version: u32,
110}
111
112#[derive(Debug, Default, Serialize, Deserialize)]
113struct EventMetadata {
114 stream: String,
115 max_sequence: u64,
116 idempotency_keys: HashMap<String, u64>,
117}
118
119#[derive(Debug, Clone, Copy, Default, Eq, PartialEq)]
121pub enum PersistentSearchMode {
122 #[default]
124 Persistent,
125 PersistentNoRebuild,
127 Disabled,
129}
130
131#[derive(Debug, Clone, Serialize, Deserialize)]
133pub struct StorageValidationReport {
134 pub format_version: u32,
136 pub legacy_manifest: bool,
138 pub lock_present: bool,
140 pub keyspaces_present: bool,
142 pub search_index_compatible: Option<bool>,
144}
145
146fn manifest_keyspaces() -> Vec<String> {
147 #[cfg(feature = "vectors")]
148 let extra = std::iter::once(REQUIRED_VECTOR_KEYSPACE);
149 #[cfg(not(feature = "vectors"))]
150 let extra = std::iter::empty();
151
152 REQUIRED_KEYSPACES
153 .iter()
154 .copied()
155 .chain(extra)
156 .map(str::to_string)
157 .collect()
158}
159
160fn search_index_compatible(path: &Path) -> Option<bool> {
161 let search_dir = path.join("search");
162 if !search_dir.exists() {
163 return None;
164 }
165 #[cfg(feature = "search")]
166 {
167 let Ok(index) = tantivy::Index::open_in_dir(search_dir) else {
168 return Some(false);
169 };
170 let schema = index.schema();
171 Some(
172 ["doc_key", "collection", "id", "body", "kind"]
173 .iter()
174 .all(|field| schema.get_field(field).is_ok()),
175 )
176 }
177 #[cfg(not(feature = "search"))]
178 {
179 Some(false)
180 }
181}
182
183fn validate_existing_directory(path: &Path) -> ThingdResult<Option<StorageValidationReport>> {
184 if !path.exists() {
185 return Ok(None);
186 }
187 if !path.is_dir() {
188 return Err(ThingdError::StorageValidation(format!(
189 "database path is not a directory: {}",
190 path.display()
191 )));
192 }
193
194 let has_entries = std::fs::read_dir(path)
195 .map_err(|error| ThingdError::StorageValidation(error.to_string()))?
196 .next()
197 .transpose()
198 .map_err(|error| ThingdError::StorageValidation(error.to_string()))?
199 .is_some();
200 if !has_entries {
201 return Ok(None);
202 }
203
204 let lock_present = path.join(STORAGE_LOCK_FILE).is_file();
205 if !lock_present {
206 return Err(ThingdError::StorageValidation(format!(
207 "missing required lock file: {}",
208 path.join(STORAGE_LOCK_FILE).display()
209 )));
210 }
211
212 let keyspaces_present = path.join(STORAGE_KEYSPACES_DIR).is_dir();
213 if !keyspaces_present {
214 return Err(ThingdError::UnsupportedStorageFormat(format!(
215 "missing keyspaces directory: {}",
216 path.join(STORAGE_KEYSPACES_DIR).display()
217 )));
218 }
219
220 let manifest_path = path.join(STORAGE_MANIFEST_FILE);
221 if manifest_path.exists() {
222 let bytes = std::fs::read(&manifest_path)
223 .map_err(|error| ThingdError::StorageValidation(error.to_string()))?;
224 let manifest: StorageManifest = serde_json::from_slice(&bytes).map_err(|error| {
225 ThingdError::UnsupportedStorageFormat(format!(
226 "invalid {}: {error}",
227 manifest_path.display()
228 ))
229 })?;
230 if manifest.format_version != STORAGE_FORMAT_VERSION {
231 return Err(ThingdError::UnsupportedStorageFormat(format!(
232 "expected format version {}, found {}",
233 STORAGE_FORMAT_VERSION, manifest.format_version
234 )));
235 }
236 if manifest.contract != STORAGE_CONTRACT {
237 return Err(ThingdError::UnsupportedStorageFormat(format!(
238 "expected contract {STORAGE_CONTRACT}, found {}",
239 manifest.contract
240 )));
241 }
242 for required in manifest_keyspaces() {
243 if !manifest.keyspaces.iter().any(|found| found == &required) {
244 return Err(ThingdError::UnsupportedStorageFormat(format!(
245 "manifest does not declare keyspace {required}"
246 )));
247 }
248 }
249 return Ok(Some(StorageValidationReport {
250 format_version: manifest.format_version,
251 legacy_manifest: false,
252 lock_present,
253 keyspaces_present,
254 search_index_compatible: search_index_compatible(path),
255 }));
256 }
257
258 Ok(Some(StorageValidationReport {
261 format_version: STORAGE_FORMAT_VERSION,
262 legacy_manifest: true,
263 lock_present,
264 keyspaces_present,
265 search_index_compatible: search_index_compatible(path),
266 }))
267}
268
269fn write_or_validate_manifest(
270 path: &Path,
271 existing: Option<&StorageValidationReport>,
272) -> ThingdResult<()> {
273 let manifest_path = path.join(STORAGE_MANIFEST_FILE);
274 if manifest_path.exists() {
275 return Ok(());
276 }
277 if let Some(report) = existing
278 && !report.legacy_manifest
279 {
280 return Ok(());
281 }
282 let manifest = StorageManifest {
283 format_version: STORAGE_FORMAT_VERSION,
284 contract: STORAGE_CONTRACT.to_string(),
285 keyspaces: manifest_keyspaces(),
286 search_schema_version: 1,
287 };
288 let bytes = serde_json::to_vec_pretty(&manifest)
289 .map_err(|error| ThingdError::StorageValidation(error.to_string()))?;
290 std::fs::write(manifest_path, bytes)
291 .map_err(|error| ThingdError::StorageValidation(error.to_string()))
292}
293
294#[derive(Clone)]
296pub struct PersistentOpenOptions {
297 pub encryption: Option<EncryptionConfig>,
299 pub allow_plaintext_output: bool,
302 pub search_mode: PersistentSearchMode,
304}
305
306impl Default for PersistentOpenOptions {
307 fn default() -> Self {
308 Self {
309 encryption: None,
310 allow_plaintext_output: false,
311 search_mode: PersistentSearchMode::Persistent,
312 }
313 }
314}
315
316#[cfg(feature = "vectors")]
317#[derive(serde::Serialize, serde::Deserialize)]
318struct StoredVector {
319 collection: String,
320 id: String,
321 vector: Vec<f32>,
322}
323
324fn value_to_vec(v: Option<fjall::Slice>) -> Option<Vec<u8>> {
325 v.map(|c| c.to_vec())
326}
327
328impl PersistentEngine {
329 pub fn open(path: impl AsRef<Path>) -> Result<Self, fjall::Error> {
332 Self::open_with_options(path, PersistentOpenOptions::default())
333 .map_err(|error| fjall::Error::Io(std::io::Error::other(error.to_string())))
334 }
335
336 pub fn open_with_options(
338 path: impl AsRef<Path>,
339 options: PersistentOpenOptions,
340 ) -> ThingdResult<Self> {
341 let path = path.as_ref();
342 let existing = validate_existing_directory(path)?;
343 let crypto = StorageCrypto::open(path, options.encryption.as_ref())?;
344 let codec = make_codec(crypto);
345 let encrypted = codec.encrypted();
346 let db = Database::builder(path).open()?;
347
348 let objects = db.keyspace("objects", KeyspaceCreateOptions::default)?;
349 let events = db.keyspace("events", KeyspaceCreateOptions::default)?;
350 let event_meta = db.keyspace("event_meta", KeyspaceCreateOptions::default)?;
351 let queue_jobs = db.keyspace("queue_jobs", KeyspaceCreateOptions::default)?;
352 let ready_jobs = db.keyspace("ready_jobs", KeyspaceCreateOptions::default)?;
353 let links_by_id = db.keyspace("links_by_id", KeyspaceCreateOptions::default)?;
354 let links_from = db.keyspace("links_from", KeyspaceCreateOptions::default)?;
355 let links_to = db.keyspace("links_to", KeyspaceCreateOptions::default)?;
356 let schemas = db.keyspace("schemas", KeyspaceCreateOptions::default)?;
357 let migrations = db.keyspace("migrations", KeyspaceCreateOptions::default)?;
358 let indexes = db.keyspace("indexes", KeyspaceCreateOptions::default)?;
359
360 #[cfg(feature = "vectors")]
361 let vectors = db.keyspace("vectors", KeyspaceCreateOptions::default)?;
362
363 write_or_validate_manifest(path, existing.as_ref())?;
364
365 let mut next_link_id = 0u64;
366 for kv in links_by_id.iter() {
367 let (_, value) = guard_data(kv)?;
368 let link: Link = {
369 let decoded = codec.decode_value("record", &value)?;
370 serde_json::from_slice(&decoded).map_err(|e| ThingdError::Storage(e.to_string()))?
371 };
372 if let Some(id_str) = link.id.strip_prefix("link-")
373 && let Ok(id) = id_str.parse::<u64>()
374 && id > next_link_id
375 {
376 next_link_id = id;
377 }
378 }
379
380 let mut event_seq_counters: HashMap<String, u64> = HashMap::new();
381 let mut event_idempotency_keys: HashMap<(String, String), u64> = HashMap::new();
382 let mut metadata_loaded = false;
383 for kv in event_meta.iter() {
384 let (_, value) = guard_data(kv)?;
385 let decoded = codec.decode_value("event_metadata", &value)?;
386 let metadata: EventMetadata = serde_json::from_slice(&decoded)
387 .map_err(|e| ThingdError::Storage(e.to_string()))?;
388 metadata_loaded = true;
389 let stream = metadata.stream.clone();
390 event_seq_counters.insert(stream.clone(), metadata.max_sequence);
391 for (idempotency_key, sequence) in metadata.idempotency_keys {
392 event_idempotency_keys.insert((stream.clone(), idempotency_key), sequence);
393 }
394 }
395
396 if !metadata_loaded {
399 for kv in events.iter() {
400 let (_, value) = guard_data(kv)?;
401 let decoded = codec.decode_value("record", &value)?;
402 let event: MemoryEvent = serde_json::from_slice(&decoded)
403 .map_err(|e| ThingdError::Storage(e.to_string()))?;
404 let stream = event.stream.clone();
405 let seq = event.sequence;
406 event_seq_counters
407 .entry(stream.clone())
408 .and_modify(|max| *max = (*max).max(seq))
409 .or_insert(seq);
410 if !event.idempotency_key.is_empty() {
411 event_idempotency_keys.insert((stream, event.idempotency_key), seq);
412 }
413 }
414 }
415
416 let (unique_index_values, unique_index_cache_complete) =
417 Self::build_unique_index_cache(&objects, &indexes, codec.as_ref())?;
418
419 #[cfg(feature = "search")]
420 let (search_index, rebuild_search_index) = match options.search_mode {
421 PersistentSearchMode::Disabled => (None, false),
422 PersistentSearchMode::Persistent if encrypted => {
423 (Self::create_search_index(true), true)
424 },
425 PersistentSearchMode::PersistentNoRebuild if encrypted => (None, false),
426 PersistentSearchMode::Persistent => Self::init_search_index(path, true)?,
427 PersistentSearchMode::PersistentNoRebuild => Self::init_search_index(path, false)?,
428 };
429
430 #[cfg(feature = "search")]
431 let search_writer = search_index
432 .as_ref()
433 .map(|index| {
434 index
435 .writer(SEARCH_WRITER_MEMORY_BYTES)
436 .map(Mutex::new)
437 .map_err(|error| ThingdError::Storage(format!("create search writer: {error}")))
438 })
439 .transpose()?;
440
441 #[cfg(feature = "search")]
442 let search_reader = search_index
443 .as_ref()
444 .map(|index| {
445 index
446 .reader()
447 .map(Mutex::new)
448 .map_err(|error| ThingdError::Storage(format!("create search reader: {error}")))
449 })
450 .transpose()?;
451
452 let engine = Self {
453 db,
454 objects,
455 events,
456 event_meta,
457 queue_jobs,
458 ready_jobs,
459 links_by_id,
460 links_from,
461 links_to,
462 schemas,
463 migrations,
464 indexes,
465 next_link_id: AtomicU64::new(next_link_id + 1),
466 event_seq_counters,
467 event_idempotency_keys,
468 unique_index_values,
469 unique_index_cache_complete,
470 #[cfg(feature = "search")]
471 search_index,
472 #[cfg(feature = "search")]
473 search_writer,
474 #[cfg(feature = "search")]
475 search_reader,
476 #[cfg(feature = "vectors")]
477 vectors,
478 codec,
479 };
480
481 if !metadata_loaded {
482 engine.persist_all_event_metadata()?;
483 }
484
485 #[cfg(feature = "search")]
486 if rebuild_search_index {
487 for entry in engine.objects.iter() {
488 let (_, value) = guard_data(entry)?;
489 let object: MemoryObject = engine.deserialize(&value)?;
490 engine.index_object_for_search_with_commit(&object, false);
491 }
492 for entry in engine.events.iter() {
493 let (_, value) = guard_data(entry)?;
494 let event: MemoryEvent = engine.deserialize(&value)?;
495 engine.index_event_for_search_with_commit(&event, false);
496 }
497 engine.commit_search_index();
498 }
499
500 Ok(engine)
501 }
502
503 pub fn validate_path(path: impl AsRef<Path>) -> ThingdResult<StorageValidationReport> {
505 validate_existing_directory(path.as_ref())?.ok_or_else(|| {
506 ThingdError::StorageValidation("database directory does not exist yet".to_string())
507 })
508 }
509
510 #[allow(clippy::items_after_statements)]
515 pub fn reencrypt_to(
516 source_path: impl AsRef<Path>,
517 destination_path: impl AsRef<Path>,
518 source_options: PersistentOpenOptions,
519 destination_options: PersistentOpenOptions,
520 ) -> ThingdResult<()> {
521 let source_path = source_path.as_ref();
522 let destination_path = destination_path.as_ref();
523 if source_path == destination_path {
524 return Err(ThingdError::EncryptionMigration(
525 "source and destination paths must differ".to_string(),
526 ));
527 }
528 if destination_path.exists() {
529 return Err(ThingdError::EncryptionMigration(
530 "destination already exists".to_string(),
531 ));
532 }
533
534 if source_options.encryption.is_some()
535 && destination_options.encryption.is_none()
536 && !destination_options.allow_plaintext_output
537 {
538 return Err(ThingdError::EncryptionMigration(
539 "encrypted-to-plaintext migration requires explicit opt-in".to_string(),
540 ));
541 }
542
543 let result = (|| {
544 let source = Self::open_with_options(source_path, source_options)?;
545 let destination = Self::open_with_options(destination_path, destination_options)?;
546
547 for entry in source.objects.iter() {
548 let (_, value) = guard_data(entry)?;
549 let object: MemoryObject = source.deserialize(&value)?;
550 let key = destination.make_object_key(&object.key.collection, &object.key.id);
551 let data = destination.serialize(&object)?;
552 destination.objects.insert(&key, &data)?;
553 #[cfg(feature = "vectors")]
554 if let Some(vector) = object.vector {
555 let vkey = destination.make_vector_key(&object.key.collection, &object.key.id);
556 let vdata = destination.serialize(&StoredVector {
557 collection: object.key.collection.clone(),
558 id: object.key.id.clone(),
559 vector,
560 })?;
561 destination.vectors.insert(&vkey, &vdata)?;
562 }
563 }
564
565 for entry in source.events.iter() {
566 let (_, value) = guard_data(entry)?;
567 let event: MemoryEvent = source.deserialize(&value)?;
568 let key = destination.make_event_key(&event.stream, event.sequence);
569 let data = destination.serialize(&event)?;
570 destination.events.insert(&key, &data)?;
571 }
572
573 for entry in source.schemas.iter() {
574 let (_, value) = guard_data(entry)?;
575 let schema: StoredSchema = source.deserialize(&value)?;
576 let key = destination.make_schema_key();
577 destination
578 .schemas
579 .insert(key, destination.serialize(&schema)?)?;
580 }
581
582 for entry in source.migrations.iter() {
583 let (_, value) = guard_data(entry)?;
584 let migration: MigrationRecord = source.deserialize(&value)?;
585 destination.migrations.insert(
586 destination.make_migration_key(&migration.id),
587 destination.serialize(&migration)?,
588 )?;
589 }
590
591 for entry in source.indexes.iter() {
592 let (_, value) = guard_data(entry)?;
593 let index: IndexDefinition = source.deserialize(&value)?;
594 destination.indexes.insert(
595 destination.make_index_key(&index.collection, &index.field),
596 destination.serialize(&index)?,
597 )?;
598 }
599
600 for entry in source.queue_jobs.iter() {
601 let (_, value) = guard_data(entry)?;
602 let job: QueueJob = source.deserialize(&value)?;
603 let key = destination.make_queue_key(&job.queue, &job.id);
604 let data = destination.serialize(&job)?;
605 destination.queue_jobs.insert(&key, &data)?;
606 if job.status == QueueJobStatus::Ready {
607 let ready_key = destination.make_ready_key(
608 &job.queue,
609 job.priority,
610 &job.created_at,
611 &job.id,
612 );
613 let ready_data = destination.serialize(&job.id)?;
614 destination.ready_jobs.insert(&ready_key, &ready_data)?;
615 }
616 }
617
618 #[cfg(feature = "vectors")]
619 for entry in source.vectors.iter() {
620 let (physical_key, value) = guard_data(entry)?;
621 let vector = if let Ok(vector) = source.deserialize::<StoredVector>(&value) {
622 vector
623 } else {
624 let Some(separator) = physical_key.iter().rposition(|byte| *byte == 0) else {
625 return Err(ThingdError::Storage(
626 "legacy vector record has no collection/id separator".to_string(),
627 ));
628 };
629 let collection = String::from_utf8_lossy(&physical_key[..separator]);
630 let id = String::from_utf8_lossy(&physical_key[separator + 1..]);
631 StoredVector {
632 collection: collection.into_owned(),
633 id: id.into_owned(),
634 vector: source.deserialize(&value)?,
635 }
636 };
637 let key = destination.make_vector_key(&vector.collection, &vector.id);
638 let data = destination.serialize(&vector)?;
639 destination.vectors.insert(&key, &data)?;
640 }
641
642 for entry in source.links_by_id.iter() {
643 let (_, value) = guard_data(entry)?;
644 let link: Link = source.deserialize(&value)?;
645 let data = destination.serialize(&link)?;
646 let index_data = destination.serialize(&link.id)?;
647 destination
648 .links_by_id
649 .insert(destination.make_link_id_key(&link.id), &data)?;
650 destination.links_from.insert(
651 destination.make_link_from_key(&link.from_ref, &link.link_type, &link.id),
652 &index_data,
653 )?;
654 destination.links_to.insert(
655 destination.make_link_to_key(&link.to_ref, &link.link_type, &link.id),
656 &index_data,
657 )?;
658 }
659
660 destination
661 .db
662 .persist(PersistMode::SyncAll)
663 .map_err(ThingdError::from)
664 })();
665
666 if result.is_err() {
667 let _ = std::fs::remove_dir_all(destination_path);
668 }
669 result.map_err(|error: ThingdError| {
670 ThingdError::EncryptionMigration(format!("logical database copy failed: {error}"))
671 })
672 }
673
674 #[cfg(feature = "search")]
675 fn search_schema() -> tantivy::schema::Schema {
676 let mut schema_builder = tantivy::schema::Schema::builder();
677 schema_builder.add_text_field("doc_key", tantivy::schema::STRING | tantivy::schema::STORED);
678 schema_builder.add_text_field(
679 "collection",
680 tantivy::schema::STRING | tantivy::schema::STORED,
681 );
682 schema_builder.add_text_field("id", tantivy::schema::STRING | tantivy::schema::STORED);
683 schema_builder.add_text_field("body", tantivy::schema::TEXT | tantivy::schema::STORED);
684 schema_builder.add_text_field("kind", tantivy::schema::STRING | tantivy::schema::STORED);
685 schema_builder.build()
686 }
687
688 #[cfg(feature = "search")]
689 fn create_search_index(in_memory: bool) -> Option<tantivy::Index> {
690 let schema = Self::search_schema();
691 if in_memory {
692 Some(tantivy::Index::create_in_ram(schema))
693 } else {
694 None
695 }
696 }
697
698 #[cfg(feature = "search")]
699 fn init_search_index(
700 path: &Path,
701 rebuild_incompatible: bool,
702 ) -> ThingdResult<(Option<tantivy::Index>, bool)> {
703 let search_dir = path.join("search");
704
705 if let Ok(index) = tantivy::Index::open_in_dir(&search_dir)
710 && Self::search_schema_is_compatible(&index.schema())
711 {
712 return Ok((Some(index), false));
713 }
714
715 if !rebuild_incompatible {
716 return Ok((None, false));
717 }
718
719 if search_dir.exists() {
720 std::fs::remove_dir_all(&search_dir)
721 .map_err(|e| ThingdError::Storage(format!("replace legacy search index: {e}")))?;
722 }
723 std::fs::create_dir_all(&search_dir)
724 .map_err(|e| ThingdError::Storage(format!("recreate search directory: {e}")))?;
725
726 let schema = Self::search_schema();
727
728 let index = tantivy::Index::create_in_dir(&search_dir, schema)
729 .map_err(|e| ThingdError::Storage(format!("create search index: {e}")))?;
730 Ok((Some(index), true))
731 }
732
733 #[cfg(feature = "search")]
734 fn search_schema_is_compatible(schema: &tantivy::schema::Schema) -> bool {
735 ["doc_key", "collection", "id", "body", "kind"]
736 .iter()
737 .all(|field| schema.get_field(field).is_ok())
738 }
739
740 fn serialize<T: serde::Serialize>(&self, value: &T) -> ThingdResult<Vec<u8>> {
741 let data = serde_json::to_vec(value).map_err(|e| ThingdError::Storage(e.to_string()))?;
742 self.codec.encode_value("record", &data)
743 }
744
745 fn deserialize<T: for<'a> serde::Deserialize<'a>>(&self, bytes: &[u8]) -> ThingdResult<T> {
746 let data = self.codec.decode_value("record", bytes)?;
747 serde_json::from_slice(&data).map_err(|e| ThingdError::Storage(e.to_string()))
748 }
749
750 fn persist_event_metadata(&self, stream: &str) -> ThingdResult<()> {
751 let max_sequence = self.event_seq_counters.get(stream).copied().unwrap_or(0);
752 let idempotency_keys = self
753 .event_idempotency_keys
754 .iter()
755 .filter_map(|((stored_stream, key), sequence)| {
756 (stored_stream == stream).then_some((key.clone(), *sequence))
757 })
758 .collect();
759 let metadata = EventMetadata {
760 stream: stream.to_string(),
761 max_sequence,
762 idempotency_keys,
763 };
764 let raw = serde_json::to_vec(&metadata)
765 .map_err(|error| ThingdError::Storage(error.to_string()))?;
766 let encoded = self.codec.encode_value("event_metadata", &raw)?;
767 let key = self
768 .codec
769 .encode_key("event_metadata.stream", stream.as_bytes());
770 self.event_meta.insert(&key, &encoded)?;
771 Ok(())
772 }
773
774 fn persist_all_event_metadata(&self) -> ThingdResult<()> {
775 let streams: Vec<String> = self.event_seq_counters.keys().cloned().collect();
776 for stream in streams {
777 self.persist_event_metadata(&stream)?;
778 }
779 Ok(())
780 }
781
782 fn make_object_key(&self, collection: &str, id: &str) -> Vec<u8> {
783 self.codec
784 .encode_scoped_key("objects", collection.as_bytes(), id.as_bytes())
785 }
786
787 fn make_object_prefix(&self, collection: &str) -> Vec<u8> {
788 self.codec
789 .encode_scoped_prefix("objects", collection.as_bytes())
790 }
791
792 fn make_event_key(&self, stream: &str, sequence: u64) -> Vec<u8> {
793 let seq_be = sequence.to_be_bytes();
794 let mut key = if self.codec.encrypted() {
795 self.codec.encode_key("events.stream", stream.as_bytes())
796 } else {
797 let mut raw = Vec::with_capacity(stream.len() + 1);
798 raw.extend_from_slice(stream.as_bytes());
799 raw.push(0);
800 raw
801 };
802 key.extend_from_slice(&seq_be);
803 key
804 }
805
806 fn make_event_prefix(&self, stream: &str) -> Vec<u8> {
807 if self.codec.encrypted() {
808 self.codec.encode_key("events.stream", stream.as_bytes())
809 } else {
810 let mut prefix = stream.as_bytes().to_vec();
811 prefix.push(0);
812 prefix
813 }
814 }
815
816 fn make_queue_key(&self, queue: &str, id: &str) -> Vec<u8> {
817 self.codec
818 .encode_scoped_key("queue_jobs", queue.as_bytes(), id.as_bytes())
819 }
820
821 fn make_queue_prefix(&self, queue: &str) -> Vec<u8> {
822 self.codec
823 .encode_scoped_prefix("queue_jobs", queue.as_bytes())
824 }
825
826 fn make_ready_key(&self, queue: &str, priority: i32, created_at: &str, id: &str) -> Vec<u8> {
828 let priority_rev = (i32::MAX - priority).to_be_bytes();
829 if self.codec.encrypted() {
830 let mut key = self.codec.encode_key("ready_jobs.queue", queue.as_bytes());
831 key.extend_from_slice(&priority_rev);
832 key.extend_from_slice(created_at.as_bytes());
833 key.extend_from_slice(&self.codec.encode_key("ready_jobs.id", id.as_bytes()));
834 return key;
835 }
836 let mut key = Vec::new();
837 key.extend_from_slice(queue.as_bytes());
838 key.push(b'\0');
839 key.extend_from_slice(&priority_rev);
840 key.push(b'\0');
841 key.extend_from_slice(created_at.as_bytes());
842 key.push(b'\0');
843 key.extend_from_slice(id.as_bytes());
844 key
845 }
846
847 fn make_ready_prefix(&self, queue: &str) -> Vec<u8> {
848 if self.codec.encrypted() {
849 self.codec.encode_key("ready_jobs.queue", queue.as_bytes())
850 } else {
851 let mut prefix = queue.as_bytes().to_vec();
852 prefix.push(0);
853 prefix
854 }
855 }
856
857 #[cfg(feature = "vectors")]
858 fn make_vector_key(&self, collection: &str, id: &str) -> Vec<u8> {
859 self.codec
860 .encode_scoped_key("vectors", collection.as_bytes(), id.as_bytes())
861 }
862
863 #[cfg(feature = "vectors")]
864 fn make_vector_prefix(&self, collection: &str) -> Vec<u8> {
865 self.codec
866 .encode_scoped_prefix("vectors", collection.as_bytes())
867 }
868
869 fn make_link_id_key(&self, link_id: &str) -> Vec<u8> {
870 self.codec.encode_key("links_by_id", link_id.as_bytes())
871 }
872
873 fn make_link_from_key(&self, from_ref: &str, link_type: &str, link_id: &str) -> Vec<u8> {
874 let suffix = format!("{link_type}\0{link_id}");
875 self.codec
876 .encode_scoped_key("links_from", from_ref.as_bytes(), suffix.as_bytes())
877 }
878
879 fn make_link_from_prefix(&self, from_ref: &str) -> Vec<u8> {
880 self.codec
881 .encode_scoped_prefix("links_from", from_ref.as_bytes())
882 }
883
884 fn make_link_to_key(&self, to_ref: &str, link_type: &str, link_id: &str) -> Vec<u8> {
885 let suffix = format!("{link_type}\0{link_id}");
886 self.codec
887 .encode_scoped_key("links_to", to_ref.as_bytes(), suffix.as_bytes())
888 }
889
890 fn make_link_to_prefix(&self, to_ref: &str) -> Vec<u8> {
891 self.codec
892 .encode_scoped_prefix("links_to", to_ref.as_bytes())
893 }
894
895 fn make_schema_key(&self) -> Vec<u8> {
896 self.codec.encode_key("schemas.current", b"current")
897 }
898
899 fn make_migration_key(&self, id: &str) -> Vec<u8> {
900 self.codec.encode_key("migrations.id", id.as_bytes())
901 }
902
903 fn make_index_key(&self, collection: &str, field: &str) -> Vec<u8> {
904 let mut value = collection.as_bytes().to_vec();
905 value.push(0);
906 value.extend_from_slice(field.as_bytes());
907 self.codec.encode_key("indexes.definition", &value)
908 }
909
910 fn unique_cache_key(
911 collection: &str,
912 field: &str,
913 value: &serde_json::Value,
914 ) -> ThingdResult<(String, String, String)> {
915 let value = serde_json::to_string(value)
916 .map_err(|error| ThingdError::Storage(error.to_string()))?;
917 Ok((collection.to_string(), field.to_string(), value))
918 }
919
920 fn build_unique_index_cache(
921 objects: &Keyspace,
922 indexes: &Keyspace,
923 codec: &dyn StorageCodec,
924 ) -> ThingdResult<(UniqueIndexCache, bool)> {
925 let mut unique_indexes = Vec::new();
926 for entry in indexes.iter() {
927 let (_, value) = guard_data(entry)?;
928 let decoded = codec.decode_value("record", &value)?;
929 let index: IndexDefinition = serde_json::from_slice(&decoded)
930 .map_err(|error| ThingdError::Storage(error.to_string()))?;
931 if index.unique {
932 unique_indexes.push(index);
933 }
934 }
935 if unique_indexes.is_empty() {
936 return Ok((HashMap::new(), true));
937 }
938
939 let mut cache = HashMap::new();
940 let mut complete = true;
941 for entry in objects.iter() {
942 let (_, value) = guard_data(entry)?;
943 let decoded = codec.decode_value("record", &value)?;
944 let object: MemoryObject = serde_json::from_slice(&decoded)
945 .map_err(|error| ThingdError::Storage(error.to_string()))?;
946 let Ok(body) = serde_json::from_str::<serde_json::Value>(&object.body) else {
947 complete = false;
948 continue;
949 };
950 for index in unique_indexes
951 .iter()
952 .filter(|index| index.collection == object.key.collection)
953 {
954 let Some(value) = body.get(&index.field).filter(|value| !value.is_null()) else {
955 continue;
956 };
957 let key = Self::unique_cache_key(&index.collection, &index.field, value)?;
958 if cache.insert(key, object.key.clone()).is_some() {
959 complete = false;
960 }
961 }
962 }
963 Ok((cache, complete))
964 }
965
966 fn refresh_unique_index_cache(&mut self) -> ThingdResult<()> {
967 let (cache, complete) =
968 Self::build_unique_index_cache(&self.objects, &self.indexes, self.codec.as_ref())?;
969 self.unique_index_values = cache;
970 self.unique_index_cache_complete = complete;
971 Ok(())
972 }
973
974 fn update_unique_index_cache_for_object(
975 &mut self,
976 object: &MemoryObject,
977 previous: Option<&MemoryObject>,
978 ) -> ThingdResult<()> {
979 if !self.unique_index_cache_complete {
980 return Ok(());
981 }
982 let indexes: Vec<IndexDefinition> = self
983 .list_index_definitions()?
984 .into_iter()
985 .filter(|index| index.unique && index.collection == object.key.collection)
986 .collect();
987 for index in indexes {
988 if let Some(previous) = previous
989 && let Ok(body) = serde_json::from_str::<serde_json::Value>(&previous.body)
990 && let Some(value) = body.get(&index.field).filter(|value| !value.is_null())
991 {
992 let key = Self::unique_cache_key(&index.collection, &index.field, value)?;
993 self.unique_index_values.remove(&key);
994 }
995 if let Ok(body) = serde_json::from_str::<serde_json::Value>(&object.body)
996 && let Some(value) = body.get(&index.field).filter(|value| !value.is_null())
997 {
998 let key = Self::unique_cache_key(&index.collection, &index.field, value)?;
999 self.unique_index_values.insert(key, object.key.clone());
1000 }
1001 }
1002 Ok(())
1003 }
1004
1005 fn remove_unique_index_cache_for_object(&mut self, object: &MemoryObject) -> ThingdResult<()> {
1006 if !self.unique_index_cache_complete {
1007 return Ok(());
1008 }
1009 let indexes: Vec<IndexDefinition> = self
1010 .list_index_definitions()?
1011 .into_iter()
1012 .filter(|index| index.unique && index.collection == object.key.collection)
1013 .collect();
1014 if let Ok(body) = serde_json::from_str::<serde_json::Value>(&object.body) {
1015 for index in indexes {
1016 if let Some(value) = body.get(&index.field).filter(|value| !value.is_null()) {
1017 let key = Self::unique_cache_key(&index.collection, &index.field, value)?;
1018 self.unique_index_values.remove(&key);
1019 }
1020 }
1021 }
1022 Ok(())
1023 }
1024
1025 fn validate_unique_indexes(&self, object: &MemoryObject) -> ThingdResult<()> {
1026 let indexes: Vec<IndexDefinition> = self
1027 .list_index_definitions()?
1028 .into_iter()
1029 .filter(|index| index.unique && index.collection == object.key.collection)
1030 .collect();
1031 if indexes.is_empty() {
1032 return Ok(());
1033 }
1034
1035 let body = serde_json::from_str::<serde_json::Value>(&object.body)
1036 .map_err(|error| ThingdError::InvalidInput(format!("invalid object JSON: {error}")))?;
1037 let existing_bodies = if self.unique_index_cache_complete {
1038 None
1039 } else {
1040 let objects = self.list_objects(
1041 Some(std::slice::from_ref(&object.key.collection)),
1042 &ListObjectsOptions::default(),
1043 )?;
1044 Some(
1045 objects
1046 .into_iter()
1047 .filter_map(|existing| {
1048 serde_json::from_str::<serde_json::Value>(&existing.body)
1049 .ok()
1050 .map(|existing_body| (existing.key, existing_body))
1051 })
1052 .collect::<Vec<_>>(),
1053 )
1054 };
1055
1056 for index in indexes {
1057 let Some(value) = body.get(&index.field) else {
1058 continue;
1059 };
1060 if value.is_null() {
1061 continue;
1062 }
1063 if self.unique_index_cache_complete {
1064 let key = Self::unique_cache_key(&index.collection, &index.field, value)?;
1065 if let Some(existing_key) = self.unique_index_values.get(&key)
1066 && existing_key != &object.key
1067 {
1068 return Err(ThingdError::Conflict(format!(
1069 "unique index {}.{} rejects duplicate value",
1070 index.collection, index.field
1071 )));
1072 }
1073 continue;
1074 }
1075 if existing_bodies.as_ref().is_some_and(|existing_bodies| {
1076 existing_bodies.iter().any(|(existing_key, existing_body)| {
1077 existing_key != &object.key
1078 && existing_body
1079 .get(&index.field)
1080 .is_some_and(|existing_value| existing_value == value)
1081 })
1082 }) {
1083 return Err(ThingdError::Conflict(format!(
1084 "unique index {}.{} rejects duplicate value",
1085 index.collection, index.field
1086 )));
1087 }
1088 }
1089 Ok(())
1090 }
1091}
1092
1093fn guard_data(kv: fjall::Guard) -> ThingdResult<(Vec<u8>, Vec<u8>)> {
1094 let kv = kv.into_inner()?;
1095 let key = kv.0.to_vec();
1096 let val = kv.1.to_vec();
1097 Ok((key, val))
1098}
1099
1100fn timestamp_before(value: &str, cutoff_ms: i64) -> bool {
1101 chrono::DateTime::parse_from_rfc3339(value)
1102 .is_ok_and(|timestamp| timestamp.timestamp_millis() < cutoff_ms)
1103}
1104
1105impl crate::SchemaStore for PersistentEngine {
1106 fn get_schema_document(&self) -> ThingdResult<Option<StoredSchema>> {
1107 let key = self.make_schema_key();
1108 value_to_vec(self.schemas.get(&key)?)
1109 .map(|value| self.deserialize(&value))
1110 .transpose()
1111 }
1112
1113 fn put_schema_document(&mut self, schema: StoredSchema) -> ThingdResult<()> {
1114 let key = self.make_schema_key();
1115 let value = self.serialize(&schema)?;
1116 self.schemas.insert(key, value)?;
1117 Ok(())
1118 }
1119
1120 fn list_migrations(&self) -> ThingdResult<Vec<MigrationRecord>> {
1121 let mut migrations = Vec::new();
1122 for entry in self.migrations.iter() {
1123 let (_, value) = guard_data(entry)?;
1124 migrations.push(self.deserialize(&value)?);
1125 }
1126 migrations.sort_by(|left: &MigrationRecord, right| left.id.cmp(&right.id));
1127 Ok(migrations)
1128 }
1129
1130 fn record_migration(&mut self, migration: MigrationRecord) -> ThingdResult<()> {
1131 let key = self.make_migration_key(&migration.id);
1132 let value = self.serialize(&migration)?;
1133 self.migrations.insert(key, value)?;
1134 Ok(())
1135 }
1136}
1137
1138impl PersistentEngine {
1139 fn safe_replication_cursor(&self) -> ThingdResult<Option<u64>> {
1140 let mut minimum: Option<u64> = None;
1141 for entry in self.objects.iter() {
1142 let (_, value) = guard_data(entry)?;
1143 let object: MemoryObject = self.deserialize(&value)?;
1144 if object.key.collection != REPLICATION_STATE_COLLECTION {
1145 continue;
1146 }
1147 let Some(cursor) = serde_json::from_str::<serde_json::Value>(&object.body)
1148 .ok()
1149 .and_then(|body| body.get("lastAppliedCursor").and_then(Value::as_u64))
1150 else {
1151 continue;
1152 };
1153 minimum = Some(minimum.map_or(cursor, |current| current.min(cursor)));
1154 }
1155 Ok(minimum)
1156 }
1157}
1158
1159impl ObjectStore for PersistentEngine {
1162 fn put_object(&mut self, mut object: MemoryObject) -> ThingdResult<MemoryObject> {
1163 self.validate_unique_indexes(&object)?;
1164 let key = self.make_object_key(&object.key.collection, &object.key.id);
1165 let previous = value_to_vec(self.objects.get(&key)?)
1166 .map(|data| self.deserialize::<MemoryObject>(&data))
1167 .transpose()?;
1168
1169 if object.created_at.is_empty() {
1170 object.created_at = now_iso_string();
1171 }
1172 object.updated_at = now_iso_string();
1173
1174 if let Some(existing_obj) = previous.as_ref() {
1175 object.version = existing_obj.version + 1;
1176 object.created_at.clone_from(&existing_obj.created_at);
1177 } else {
1178 object.version = 1;
1179 }
1180
1181 let data = self.serialize(&object)?;
1182
1183 let mut batch = self.db.batch();
1185 batch.insert(&self.objects, &key, &data);
1186 #[cfg(feature = "vectors")]
1187 {
1188 let vkey = self.make_vector_key(&object.key.collection, &object.key.id);
1189 if let Some(ref vector) = object.vector {
1190 let vdata = self.serialize(&StoredVector {
1191 collection: object.key.collection.clone(),
1192 id: object.key.id.clone(),
1193 vector: vector.clone(),
1194 })?;
1195 batch.insert(&self.vectors, &vkey, vdata);
1196 } else {
1197 batch.remove(&self.vectors, vkey);
1198 }
1199 }
1200 batch
1201 .commit()
1202 .map_err(|e| ThingdError::Storage(e.to_string()))?;
1203
1204 #[cfg(feature = "search")]
1205 self.index_object_for_search(&object);
1206
1207 self.update_unique_index_cache_for_object(&object, previous.as_ref())?;
1208
1209 Ok(object)
1210 }
1211
1212 fn put_object_with_options(
1213 &mut self,
1214 object: MemoryObject,
1215 options: PutObjectOptions,
1216 ) -> ThingdResult<MemoryObject> {
1217 let key = self.make_object_key(&object.key.collection, &object.key.id);
1218
1219 if let Some(expected_version) = options.expected_version {
1220 match value_to_vec(self.objects.get(&key)?) {
1221 Some(existing) => {
1222 let existing_obj: MemoryObject = self.deserialize(&existing)?;
1223 if existing_obj.version != expected_version {
1224 return Err(ThingdError::Conflict(format!(
1225 "expected version {} but current version is {}",
1226 expected_version, existing_obj.version
1227 )));
1228 }
1229 },
1230 None => {
1231 return Err(ThingdError::Conflict(format!(
1232 "object '{}/{}' does not exist",
1233 object.key.collection, object.key.id
1234 )));
1235 },
1236 }
1237 }
1238
1239 self.put_object(object)
1240 }
1241
1242 fn retain(&mut self, options: RetentionOptions) -> ThingdResult<RetentionReport> {
1243 let safe_replication_cursor = if options.include_replication {
1244 self.safe_replication_cursor()?
1245 } else {
1246 None
1247 };
1248 let mut report = RetentionReport {
1249 dry_run: options.dry_run,
1250 safe_replication_cursor,
1251 ..RetentionReport::default()
1252 };
1253 let mut event_deletions = Vec::new();
1254 for entry in self.events.iter() {
1255 let (key, value) = guard_data(entry)?;
1256 let event: MemoryEvent = self.deserialize(&value)?;
1257 let old = timestamp_before(&event.created_at, options.before_unix_ms);
1258 let eligible = if event.stream == REPLICATION_STREAM {
1259 if old
1260 && options.include_replication
1261 && safe_replication_cursor.is_some_and(|cursor| event.sequence <= cursor)
1262 {
1263 true
1264 } else if old {
1265 report.skipped_replication_events += 1;
1266 false
1267 } else {
1268 false
1269 }
1270 } else {
1271 !self.is_protected_stream(&event.stream) && old
1272 };
1273 if eligible {
1274 report.events += 1;
1275 event_deletions.push((key, event));
1276 }
1277 }
1278
1279 let mut job_deletions = Vec::new();
1280 for entry in self.queue_jobs.iter() {
1281 let (key, value) = guard_data(entry)?;
1282 let job: QueueJob = self.deserialize(&value)?;
1283 let eligible = match job.status {
1284 QueueJobStatus::Completed => job
1285 .completed_at_ms
1286 .is_some_and(|timestamp| timestamp < options.before_unix_ms),
1287 QueueJobStatus::Dead => job
1288 .dead_at_ms
1289 .is_some_and(|timestamp| timestamp < options.before_unix_ms),
1290 _ => false,
1291 };
1292 if eligible {
1293 match job.status {
1294 QueueJobStatus::Completed => report.completed_jobs += 1,
1295 QueueJobStatus::Dead => report.dead_jobs += 1,
1296 _ => unreachable!(),
1297 }
1298 job_deletions.push((key, job));
1299 }
1300 }
1301
1302 if options.dry_run {
1303 return Ok(report);
1304 }
1305
1306 let mut batch = self.db.batch();
1307 for (key, _) in &event_deletions {
1308 batch.remove(&self.events, key.clone());
1309 }
1310 for (key, _) in &job_deletions {
1311 batch.remove(&self.queue_jobs, key.clone());
1312 }
1313 if !event_deletions.is_empty() || !job_deletions.is_empty() {
1314 batch
1315 .commit()
1316 .map_err(|error| ThingdError::Storage(error.to_string()))?;
1317 }
1318
1319 let mut affected_streams = std::collections::HashSet::new();
1320 #[cfg(feature = "search")]
1321 let mut search_deletions = Vec::new();
1322 for (_, event) in &event_deletions {
1323 affected_streams.insert(event.stream.clone());
1324 self.event_idempotency_keys.retain(|(stream, _), sequence| {
1325 stream != &event.stream || *sequence != event.sequence
1326 });
1327 #[cfg(feature = "search")]
1328 search_deletions.push((event.stream.clone(), event.sequence));
1329 }
1330 #[cfg(feature = "search")]
1331 self.delete_events_from_search_index(&search_deletions);
1332 for stream in affected_streams {
1333 self.persist_event_metadata(&stream)?;
1334 }
1335 if options.compact && (!event_deletions.is_empty() || !job_deletions.is_empty()) {
1336 self.events
1337 .major_compact()
1338 .map_err(|error| ThingdError::Storage(error.to_string()))?;
1339 self.queue_jobs
1340 .major_compact()
1341 .map_err(|error| ThingdError::Storage(error.to_string()))?;
1342 report.compacted = true;
1343 }
1344 Ok(report)
1345 }
1346
1347 fn get_object(&self, collection: &str, id: &str) -> ThingdResult<Option<MemoryObject>> {
1348 let key = self.make_object_key(collection, id);
1349 match value_to_vec(self.objects.get(&key)?) {
1350 Some(data) => Ok(Some(self.deserialize(&data)?)),
1351 None => Ok(None),
1352 }
1353 }
1354
1355 fn get_objects_batch(
1356 &self,
1357 collection: &str,
1358 ids: &[String],
1359 ) -> ThingdResult<Vec<Option<MemoryObject>>> {
1360 ids.iter()
1361 .map(|id| self.get_object(collection, id))
1362 .collect()
1363 }
1364
1365 fn list_objects(
1366 &self,
1367 collections: Option<&[String]>,
1368 options: &ListObjectsOptions,
1369 ) -> ThingdResult<Vec<MemoryObject>> {
1370 if options.sort_by.is_none() {
1373 let offset = options.offset.unwrap_or(0);
1374 let limit = options.limit.unwrap_or(u64::MAX);
1375 let mut skipped = 0u64;
1376 let mut results = Vec::new();
1377
1378 if let Some(collections) = collections
1379 && collections.len() == 1
1380 {
1381 let prefix = self.make_object_prefix(&collections[0]);
1382 for kv in self.objects.prefix(&prefix) {
1383 let (_, value) = guard_data(kv)?;
1384 let object: MemoryObject = self.deserialize(&value)?;
1385 if !matches_object_filters(&object, &options.filter) {
1386 continue;
1387 }
1388 if skipped < offset {
1389 skipped += 1;
1390 continue;
1391 }
1392 results.push(object);
1393 if results.len() as u64 >= limit {
1394 break;
1395 }
1396 }
1397 } else {
1398 for kv in self.objects.iter() {
1399 let (_, value) = guard_data(kv)?;
1400 let object: MemoryObject = self.deserialize(&value)?;
1401 if let Some(collections) = collections
1402 && !collections.contains(&object.key.collection)
1403 {
1404 continue;
1405 }
1406 if !matches_object_filters(&object, &options.filter) {
1407 continue;
1408 }
1409 if skipped < offset {
1410 skipped += 1;
1411 continue;
1412 }
1413 results.push(object);
1414 if results.len() as u64 >= limit {
1415 break;
1416 }
1417 }
1418 }
1419 return Ok(results);
1420 }
1421
1422 let prefix = if let Some(collections) = collections
1423 && collections.len() == 1
1424 {
1425 Some(self.make_object_prefix(&collections[0]))
1426 } else {
1427 None
1428 };
1429
1430 let mut objects: Vec<MemoryObject> = if let Some(ref prefix) = prefix {
1431 let mut objs = Vec::new();
1432 for kv in self.objects.prefix(prefix) {
1433 let (_, value) = guard_data(kv)?;
1434 objs.push(self.deserialize(&value)?);
1435 }
1436 objs
1437 } else {
1438 let mut objs = Vec::new();
1439 for kv in self.objects.iter() {
1440 let (_, value) = guard_data(kv)?;
1441 objs.push(self.deserialize(&value)?);
1442 }
1443 objs
1444 };
1445
1446 if let Some(cols) = collections
1447 && cols.len() != 1
1448 {
1449 objects.retain(|o| cols.contains(&o.key.collection));
1450 }
1451
1452 if !options.filter.is_empty() {
1453 objects.retain(|object| matches_object_filters(object, &options.filter));
1454 }
1455
1456 if let Some(ref sort_by) = options.sort_by {
1457 let asc = sort_by.direction == SortDirection::Asc;
1458 objects.sort_by(|a, b| {
1459 let cmp = if sort_by.field.starts_with("$.") {
1460 let path = sort_by.field.trim_start_matches('$');
1461 let a_val = serde_json::from_str::<serde_json::Value>(&a.body)
1462 .ok()
1463 .and_then(|v| v.get(path).cloned());
1464 let b_val = serde_json::from_str::<serde_json::Value>(&b.body)
1465 .ok()
1466 .and_then(|v| v.get(path).cloned());
1467 match (&a_val, &b_val) {
1468 (Some(a), Some(b)) => value_compare(a, b),
1469 (Some(_), None) => std::cmp::Ordering::Greater,
1470 (None, Some(_)) => std::cmp::Ordering::Less,
1471 (None, None) => std::cmp::Ordering::Equal,
1472 }
1473 } else {
1474 match sort_by.field.as_str() {
1475 "id" => a.key.id.cmp(&b.key.id),
1476 "collection" => a.key.collection.cmp(&b.key.collection),
1477 "created_at" => a.created_at.cmp(&b.created_at),
1478 "updated_at" => a.updated_at.cmp(&b.updated_at),
1479 "version" => a.version.cmp(&b.version),
1480 _ => std::cmp::Ordering::Equal,
1481 }
1482 };
1483 if asc { cmp } else { cmp.reverse() }
1484 });
1485 }
1486
1487 if let Some(offset) = options.offset {
1488 let skip = usize::try_from(offset).unwrap_or(usize::MAX);
1489 objects = objects.into_iter().skip(skip).collect();
1490 }
1491 if let Some(limit) = options.limit {
1492 let take = usize::try_from(limit).unwrap_or(usize::MAX);
1493 objects.truncate(take);
1494 }
1495
1496 Ok(objects)
1497 }
1498
1499 fn delete_object(&mut self, collection: &str, id: &str) -> ThingdResult<bool> {
1500 let key = self.make_object_key(collection, id);
1501 let existing = value_to_vec(self.objects.get(&key)?)
1502 .map(|data| self.deserialize::<MemoryObject>(&data))
1503 .transpose()?;
1504 let existed = existing.is_some();
1505
1506 let mut batch = self.db.batch();
1508 batch.remove(&self.objects, key);
1509 #[cfg(feature = "vectors")]
1510 {
1511 let vkey = self.make_vector_key(collection, id);
1512 batch.remove(&self.vectors, vkey);
1513 }
1514 batch
1515 .commit()
1516 .map_err(|e| ThingdError::Storage(e.to_string()))?;
1517
1518 #[cfg(feature = "search")]
1519 self.delete_object_from_search_index(collection, id);
1520 if let Some(existing) = existing.as_ref() {
1521 self.remove_unique_index_cache_for_object(existing)?;
1522 }
1523 Ok(existed)
1524 }
1525
1526 fn delete_objects_batch(&mut self, keys: &[(String, String)]) -> ThingdResult<u64> {
1527 let mut count = 0u64;
1528 let mut batch = self.db.batch();
1529 let mut deleted = Vec::new();
1530
1531 for (collection, id) in keys {
1532 let key = self.make_object_key(collection, id);
1533 if self.objects.get(&key)?.is_some() {
1534 if let Some(data) = value_to_vec(self.objects.get(&key)?) {
1535 deleted.push(self.deserialize::<MemoryObject>(&data)?);
1536 }
1537 batch.remove(&self.objects, key);
1538 count += 1;
1539 }
1540 #[cfg(feature = "vectors")]
1541 {
1542 let vkey = self.make_vector_key(collection, id);
1543 batch.remove(&self.vectors, vkey);
1544 }
1545 }
1546
1547 if count > 0 {
1548 batch
1549 .commit()
1550 .map_err(|e| ThingdError::Storage(e.to_string()))?;
1551 }
1552
1553 #[cfg(feature = "search")]
1554 for (collection, id) in keys {
1555 self.delete_object_from_search_index(collection, id);
1556 }
1557 for object in &deleted {
1558 self.remove_unique_index_cache_for_object(object)?;
1559 }
1560
1561 Ok(count)
1562 }
1563
1564 fn count_objects(&self) -> ThingdResult<u64> {
1565 let mut count = 0u64;
1566 for kv in self.objects.iter() {
1567 let _ = kv;
1568 count += 1;
1569 }
1570 Ok(count)
1571 }
1572
1573 fn count_objects_in_collection(&self, collection: &str) -> ThingdResult<u64> {
1574 let prefix = self.make_object_prefix(collection);
1575 let mut count = 0u64;
1576 for kv in self.objects.prefix(&prefix) {
1577 let _ = kv;
1578 count += 1;
1579 }
1580 Ok(count)
1581 }
1582
1583 fn list_collections(&self) -> ThingdResult<Vec<String>> {
1584 let mut collections: Vec<String> = Vec::new();
1585 for kv in self.objects.iter() {
1586 let (_, value) = guard_data(kv)?;
1587 let object: MemoryObject = self.deserialize(&value)?;
1588 if !collections.contains(&object.key.collection) {
1589 collections.push(object.key.collection);
1590 }
1591 }
1592 Ok(collections)
1593 }
1594
1595 fn create_index(&mut self, collection: &str, field: &str) -> ThingdResult<()> {
1596 self.create_index_definition(IndexDefinition {
1597 collection: collection.to_string(),
1598 field: field.to_string(),
1599 unique: false,
1600 })
1601 }
1602
1603 fn create_index_definition(&mut self, index: IndexDefinition) -> ThingdResult<()> {
1604 if index.collection.is_empty() || index.field.is_empty() {
1605 return Err(ThingdError::InvalidInput(
1606 "index collection and field are required".to_string(),
1607 ));
1608 }
1609 if index.unique {
1610 let objects = self.list_objects(
1611 Some(std::slice::from_ref(&index.collection)),
1612 &ListObjectsOptions::default(),
1613 )?;
1614 let mut values = Vec::new();
1615 for object in objects {
1616 let body =
1617 serde_json::from_str::<serde_json::Value>(&object.body).map_err(|error| {
1618 ThingdError::InvalidInput(format!("invalid object JSON: {error}"))
1619 })?;
1620 if let Some(value) = body.get(&index.field).filter(|value| !value.is_null()) {
1621 if values.iter().any(|existing| existing == value) {
1622 return Err(ThingdError::Conflict(format!(
1623 "cannot create unique index {}.{}: existing values are duplicated",
1624 index.collection, index.field
1625 )));
1626 }
1627 values.push(value.clone());
1628 }
1629 }
1630 }
1631 self.indexes.insert(
1632 self.make_index_key(&index.collection, &index.field),
1633 self.serialize(&index)?,
1634 )?;
1635 self.refresh_unique_index_cache()?;
1636 Ok(())
1637 }
1638
1639 fn list_indexes(&self) -> ThingdResult<Vec<(String, String)>> {
1640 Ok(self
1641 .list_index_definitions()?
1642 .into_iter()
1643 .map(|index| (index.collection, index.field))
1644 .collect())
1645 }
1646
1647 fn delete_index(&mut self, collection: &str, field: &str) -> ThingdResult<bool> {
1648 let key = self.make_index_key(collection, field);
1649 let existed = self.indexes.get(&key)?.is_some();
1650 if existed {
1651 self.indexes.remove(key)?;
1652 self.refresh_unique_index_cache()?;
1653 }
1654 Ok(existed)
1655 }
1656
1657 fn list_index_definitions(&self) -> ThingdResult<Vec<IndexDefinition>> {
1658 let mut definitions = Vec::new();
1659 for entry in self.indexes.iter() {
1660 let (_, value) = guard_data(entry)?;
1661 definitions.push(self.deserialize(&value)?);
1662 }
1663 definitions.sort_by(|left: &IndexDefinition, right: &IndexDefinition| {
1664 (&left.collection, &left.field).cmp(&(&right.collection, &right.field))
1665 });
1666 Ok(definitions)
1667 }
1668
1669 fn schema(
1670 &self,
1671 collection: Option<&str>,
1672 options: &SchemaOptions,
1673 ) -> ThingdResult<Vec<CollectionSchema>> {
1674 let sample_size = options.sample_size.unwrap_or(50);
1675 let mut schemas: Vec<CollectionSchema> = Vec::new();
1676
1677 let collections: Vec<String> = if let Some(c) = collection {
1678 vec![c.to_string()]
1679 } else {
1680 self.list_collections()?
1681 };
1682
1683 for col in collections {
1684 let prefix = self.make_object_prefix(&col);
1685 let mut objects: Vec<MemoryObject> = Vec::new();
1686 let mut count = 0u64;
1687 for kv in self.objects.prefix(&prefix) {
1688 count += 1;
1689 if objects.len() < sample_size {
1690 let (_, value) = guard_data(kv)?;
1691 objects.push(self.deserialize(&value)?);
1692 }
1693 }
1694
1695 let object_count = count;
1696 let mut fields: Vec<FieldSchema> = Vec::new();
1697 let mut field_types: HashMap<String, (String, bool, Vec<serde_json::Value>)> =
1698 HashMap::new();
1699
1700 for obj in &objects {
1701 if let Ok(body) = serde_json::from_str::<serde_json::Value>(&obj.body)
1702 && let serde_json::Value::Object(map) = body
1703 {
1704 for (field_name, field_val) in map {
1705 let entry = field_types
1706 .entry(field_name.clone())
1707 .or_insert_with(|| ("unknown".into(), false, Vec::new()));
1708 let t = infer_json_type(&field_val);
1709 if entry.0 == "unknown" {
1710 entry.0 = t.clone();
1711 } else if entry.0 != t {
1712 entry.0 = "string".into();
1713 }
1714 if !field_val.is_null() && entry.2.len() < 3 {
1715 entry.2.push(field_val.clone());
1716 }
1717 }
1718 }
1719 }
1720
1721 for (name, (field_type, _, sample_values)) in field_types {
1722 fields.push(FieldSchema {
1723 name,
1724 field_type,
1725 nullable: false,
1726 sample_values,
1727 });
1728 }
1729
1730 fields.sort_by(|a, b| a.name.cmp(&b.name));
1731
1732 schemas.push(CollectionSchema {
1733 name: col,
1734 object_count,
1735 fields,
1736 });
1737 }
1738
1739 Ok(schemas)
1740 }
1741}
1742
1743impl EventLog for PersistentEngine {
1746 fn is_protected_stream(&self, stream: &str) -> bool {
1747 stream.starts_with("__thingd:")
1748 }
1749
1750 fn append_event(&mut self, mut event: MemoryEvent) -> ThingdResult<MemoryEvent> {
1751 if !event.idempotency_key.is_empty() {
1752 let idem_key = (event.stream.clone(), event.idempotency_key.clone());
1753 if let Some(&existing_seq) = self.event_idempotency_keys.get(&idem_key) {
1754 let ekey = self.make_event_key(&event.stream, existing_seq);
1755 if let Some(data) = value_to_vec(self.events.get(&ekey)?) {
1756 let existing: MemoryEvent = self.deserialize(&data)?;
1757 return Ok(existing);
1758 }
1759 }
1760 }
1761
1762 let seq = self
1763 .event_seq_counters
1764 .entry(event.stream.clone())
1765 .and_modify(|s| *s += 1)
1766 .or_insert(1);
1767 event.sequence = *seq;
1768
1769 if event.created_at.is_empty() {
1770 event.created_at = now_iso_string();
1771 }
1772
1773 let ekey = self.make_event_key(&event.stream, event.sequence);
1774 let data = self.serialize(&event)?;
1775 self.events.insert(&ekey, &data)?;
1776
1777 #[cfg(feature = "search")]
1778 self.index_event_for_search(&event);
1779
1780 if !event.idempotency_key.is_empty() {
1781 self.event_idempotency_keys.insert(
1782 (event.stream.clone(), event.idempotency_key.clone()),
1783 event.sequence,
1784 );
1785 }
1786
1787 self.persist_event_metadata(&event.stream)?;
1788
1789 Ok(event)
1790 }
1791
1792 fn append_events_batch(&mut self, events: Vec<MemoryEvent>) -> ThingdResult<Vec<MemoryEvent>> {
1793 let mut results = Vec::with_capacity(events.len());
1794 for event in events {
1795 results.push(self.append_event(event)?);
1796 }
1797 Ok(results)
1798 }
1799
1800 fn list_events(
1801 &self,
1802 stream: Option<&str>,
1803 options: ListEventsOptions,
1804 ) -> ThingdResult<Vec<MemoryEvent>> {
1805 let mut results: Vec<MemoryEvent> = Vec::new();
1806
1807 if let Some(stream_name) = stream {
1808 for kv in self.events.prefix(self.make_event_prefix(stream_name)) {
1809 let (_, value) = guard_data(kv)?;
1810 let event: MemoryEvent = self.deserialize(&value)?;
1811 if event.stream != stream_name {
1812 continue;
1813 }
1814 if let Some(from_seq) = options.from_sequence
1815 && event.sequence <= from_seq
1816 {
1817 continue;
1818 }
1819 if let Some(ref since) = options.since
1820 && event.created_at.as_str() < since.as_str()
1821 {
1822 continue;
1823 }
1824 results.push(event);
1825 if let Some(limit) = options.limit
1826 && results.len() as u64 >= limit
1827 {
1828 break;
1829 }
1830 }
1831 } else {
1832 for kv in self.events.iter() {
1833 let (_, value) = guard_data(kv)?;
1834 let event: MemoryEvent = self.deserialize(&value)?;
1835 if let Some(ref since) = options.since
1836 && event.created_at.as_str() < since.as_str()
1837 {
1838 continue;
1839 }
1840 results.push(event);
1841 if let Some(limit) = options.limit
1842 && results.len() as u64 >= limit
1843 {
1844 break;
1845 }
1846 }
1847 }
1848
1849 Ok(results)
1850 }
1851
1852 fn delete_last_event(&mut self, stream: &str) -> ThingdResult<Option<MemoryEvent>> {
1853 if self.is_protected_stream(stream) {
1854 return Err(ThingdError::Protected(format!(
1855 "stream '{stream}' is protected and cannot be modified"
1856 )));
1857 }
1858
1859 let mut last_key: Option<Vec<u8>> = None;
1860 let mut last_event: Option<MemoryEvent> = None;
1861
1862 for kv in self.events.iter() {
1863 let (key, value) = guard_data(kv)?;
1864 let event: MemoryEvent = self.deserialize(&value)?;
1865 if event.stream == stream
1866 && last_event
1867 .as_ref()
1868 .is_none_or(|last| event.sequence > last.sequence)
1869 {
1870 last_key = Some(key);
1871 last_event = Some(event);
1872 }
1873 }
1874
1875 if let Some(key) = last_key {
1876 self.events.remove(&key)?;
1877 if let Some(ref ev) = last_event {
1878 self.event_idempotency_keys
1879 .retain(|(stored_stream, _), sequence| {
1880 stored_stream != stream || *sequence != ev.sequence
1881 });
1882 }
1883 #[cfg(feature = "search")]
1884 if let Some(ref ev) = last_event {
1885 self.delete_event_from_search_index(&ev.stream, ev.sequence);
1886 }
1887 self.persist_event_metadata(stream)?;
1888 Ok(last_event)
1889 } else {
1890 Ok(None)
1891 }
1892 }
1893
1894 fn delete_stream(&mut self, stream: &str) -> ThingdResult<u64> {
1895 if self.is_protected_stream(stream) {
1896 return Err(ThingdError::Protected(format!(
1897 "stream '{stream}' is protected and cannot be modified"
1898 )));
1899 }
1900
1901 let mut entries: Vec<(Vec<u8>, u64)> = Vec::new();
1902 for kv in self.events.iter() {
1903 let (key, value) = guard_data(kv)?;
1904 let event: MemoryEvent = self.deserialize(&value)?;
1905 if event.stream == stream {
1906 entries.push((key, event.sequence));
1907 }
1908 }
1909
1910 let count = entries.len() as u64;
1911 for (key, seq) in &entries {
1912 self.events.remove(key)?;
1913 #[cfg(feature = "search")]
1914 self.delete_event_from_search_index(stream, *seq);
1915 }
1916
1917 self.event_seq_counters.remove(stream);
1918 self.event_idempotency_keys.retain(|(s, _), _| s != stream);
1919 let key = self
1920 .codec
1921 .encode_key("event_metadata.stream", stream.as_bytes());
1922 self.event_meta.remove(&key)?;
1923
1924 Ok(count)
1925 }
1926
1927 fn count_events(&self) -> ThingdResult<u64> {
1928 let mut count = 0u64;
1929 for kv in self.events.iter() {
1930 let _ = kv;
1931 count += 1;
1932 }
1933 Ok(count)
1934 }
1935
1936 fn list_streams(&self) -> ThingdResult<Vec<String>> {
1937 let mut streams: Vec<String> = Vec::new();
1938 for kv in self.events.iter() {
1939 let (_, value) = guard_data(kv)?;
1940 let event: MemoryEvent = self.deserialize(&value)?;
1941 let stream = event.stream;
1942 if !streams.contains(&stream) {
1943 streams.push(stream);
1944 }
1945 }
1946 Ok(streams)
1947 }
1948}
1949
1950impl QueueStore for PersistentEngine {
1953 fn push_job(&mut self, mut job: QueueJob) -> ThingdResult<QueueJob> {
1954 if job.created_at.is_empty() {
1955 job.created_at = now_iso_string();
1956 }
1957 let key = self.make_queue_key(&job.queue, &job.id);
1958 let data = self.serialize(&job)?;
1959 self.queue_jobs.insert(&key, &data)?;
1960 if job.status == QueueJobStatus::Ready {
1962 let rkey = self.make_ready_key(&job.queue, job.priority, &job.created_at, &job.id);
1963 let rdata = self.serialize(&job.id)?;
1964 self.ready_jobs.insert(&rkey, &rdata)?;
1965 }
1966 Ok(job)
1967 }
1968
1969 fn push_jobs_batch(&mut self, jobs: Vec<QueueJob>) -> ThingdResult<Vec<QueueJob>> {
1970 let mut results = Vec::with_capacity(jobs.len());
1971 for job in jobs {
1972 results.push(self.push_job(job)?);
1973 }
1974 Ok(results)
1975 }
1976
1977 fn claim_job_with_options(
1978 &mut self,
1979 queue: &str,
1980 options: QueueClaimOptions,
1981 ) -> ThingdResult<Option<QueueJob>> {
1982 let now = unix_timestamp_millis();
1983
1984 let qprefix = self.make_queue_prefix(queue);
1986 for kv in self.queue_jobs.prefix(&qprefix) {
1987 let (key, value) = guard_data(kv)?;
1988 let mut job: QueueJob = self.deserialize(&value)?;
1989 if job.status == QueueJobStatus::Leased
1990 && job.lease_expires_at_ms.is_some_and(|exp| exp <= now)
1991 {
1992 job.status = QueueJobStatus::Ready;
1993 job.leased_at_ms = None;
1994 job.lease_expires_at_ms = None;
1995 let data = self.serialize(&job)?;
1996 self.queue_jobs.insert(&key, &data)?;
1997 let rkey = self.make_ready_key(&job.queue, job.priority, &job.created_at, &job.id);
1998 let rdata = self.serialize(&job.id)?;
1999 self.ready_jobs.insert(&rkey, &rdata)?;
2000 }
2001 }
2002
2003 let prefix = self.make_ready_prefix(queue);
2005 for kv in self.ready_jobs.prefix(&prefix) {
2006 let (key, value) = guard_data(kv)?;
2007 let job_id = if value.is_empty() {
2008 let key_str = String::from_utf8_lossy(&key);
2009 let parts: Vec<&str> = key_str.splitn(4, '\0').collect();
2010 if parts.len() < 4 {
2011 continue;
2012 }
2013 parts[3].to_string()
2014 } else {
2015 self.deserialize::<String>(&value)?
2016 };
2017 let rkey = key;
2018
2019 let qkey = self.make_queue_key(queue, &job_id);
2021 let Some(job_data) = value_to_vec(self.queue_jobs.get(&qkey)?) else {
2022 let _ = self.ready_jobs.remove(&rkey);
2024 continue;
2025 };
2026 let mut job: QueueJob = self.deserialize(&job_data)?;
2027 if job.status == QueueJobStatus::Leased
2029 && job.lease_expires_at_ms.is_some_and(|exp| exp <= now)
2030 {
2031 job.status = QueueJobStatus::Ready;
2032 job.leased_at_ms = None;
2033 job.lease_expires_at_ms = None;
2034 }
2035
2036 if job.status != QueueJobStatus::Ready {
2038 let _ = self.ready_jobs.remove(&rkey);
2039 continue;
2040 }
2041
2042 if job.available_at_ms > now {
2044 continue;
2045 }
2046
2047 self.ready_jobs.remove(&rkey)?;
2049 job.status = QueueJobStatus::Leased;
2050 job.attempts = job.attempts.saturating_add(1);
2051 job.leased_at_ms = Some(now);
2052 job.lease_expires_at_ms = Some(now + options.lease_ms as i64);
2053 let data = self.serialize(&job)?;
2054 self.queue_jobs.insert(&qkey, &data)?;
2055 return Ok(Some(job));
2056 }
2057
2058 Ok(None)
2059 }
2060
2061 fn ack_job(&mut self, queue: &str, id: &str) -> ThingdResult<Option<QueueJob>> {
2062 let key = self.make_queue_key(queue, id);
2063 match value_to_vec(self.queue_jobs.get(&key)?) {
2064 Some(data) => {
2065 let mut job: QueueJob = self.deserialize(&data)?;
2066 if job.status != QueueJobStatus::Leased {
2067 return Ok(None);
2068 }
2069 job.status = QueueJobStatus::Completed;
2070 job.completed_at_ms = Some(unix_timestamp_millis());
2071 let new_data = self.serialize(&job)?;
2072 self.queue_jobs.insert(&key, &new_data)?;
2073 Ok(Some(job))
2074 },
2075 None => Ok(None),
2076 }
2077 }
2078
2079 fn nack_job_with_options(
2080 &mut self,
2081 queue: &str,
2082 id: &str,
2083 options: QueueNackOptions,
2084 ) -> ThingdResult<Option<QueueJob>> {
2085 let key = self.make_queue_key(queue, id);
2086 match value_to_vec(self.queue_jobs.get(&key)?) {
2087 Some(data) => {
2088 let mut job: QueueJob = self.deserialize(&data)?;
2089 if job.status != QueueJobStatus::Leased {
2090 return Ok(None);
2091 }
2092 job.last_error = options.error;
2093 job.leased_at_ms = None;
2094 job.lease_expires_at_ms = None;
2095
2096 let is_dead = job.attempts >= job.max_attempts;
2097 if is_dead {
2098 job.status = QueueJobStatus::Dead;
2099 job.dead_at_ms = Some(unix_timestamp_millis());
2100 } else {
2101 job.status = QueueJobStatus::Ready;
2102 job.available_at_ms = unix_timestamp_millis() + options.delay_ms as i64;
2103 }
2104
2105 let new_data = self.serialize(&job)?;
2106 self.queue_jobs.insert(&key, &new_data)?;
2107
2108 if !is_dead {
2110 let rkey =
2111 self.make_ready_key(&job.queue, job.priority, &job.created_at, &job.id);
2112 let rdata = self.serialize(&job.id)?;
2113 self.ready_jobs.insert(&rkey, &rdata)?;
2114 }
2115
2116 Ok(Some(job))
2117 },
2118 None => Ok(None),
2119 }
2120 }
2121
2122 fn list_jobs(&self, queue: &str) -> ThingdResult<Vec<QueueJob>> {
2123 let prefix = self.make_queue_prefix(queue);
2124 let mut jobs = Vec::new();
2125 for kv in self.queue_jobs.prefix(&prefix) {
2126 let (_, value) = guard_data(kv)?;
2127 jobs.push(self.deserialize(&value)?);
2128 }
2129 Ok(jobs)
2130 }
2131
2132 fn list_dead_jobs(&self, queue: &str) -> ThingdResult<Vec<QueueJob>> {
2133 let prefix = self.make_queue_prefix(queue);
2134 let mut jobs = Vec::new();
2135 for kv in self.queue_jobs.prefix(&prefix) {
2136 let (_, value) = guard_data(kv)?;
2137 let job: QueueJob = self.deserialize(&value)?;
2138 if job.status == QueueJobStatus::Dead {
2139 jobs.push(job);
2140 }
2141 }
2142 Ok(jobs)
2143 }
2144
2145 fn list_queues(&self) -> ThingdResult<Vec<String>> {
2146 let mut queues: Vec<String> = Vec::new();
2147 for kv in self.queue_jobs.iter() {
2148 let (_, value) = guard_data(kv)?;
2149 let job: QueueJob = self.deserialize(&value)?;
2150 if !queues.contains(&job.queue) {
2151 queues.push(job.queue);
2152 }
2153 }
2154 Ok(queues)
2155 }
2156
2157 fn count_active_jobs(&self) -> ThingdResult<u64> {
2158 let mut count = 0u64;
2159 for kv in self.queue_jobs.iter() {
2160 let (_, value) = guard_data(kv)?;
2161 let job: QueueJob = self.deserialize(&value)?;
2162 if job.status == QueueJobStatus::Ready || job.status == QueueJobStatus::Leased {
2163 count += 1;
2164 }
2165 }
2166 Ok(count)
2167 }
2168
2169 fn count_dead_jobs(&self) -> ThingdResult<u64> {
2170 let mut count = 0u64;
2171 for kv in self.queue_jobs.iter() {
2172 let (_, value) = guard_data(kv)?;
2173 let job: QueueJob = self.deserialize(&value)?;
2174 if job.status == QueueJobStatus::Dead {
2175 count += 1;
2176 }
2177 }
2178 Ok(count)
2179 }
2180}
2181
2182impl Searcher for PersistentEngine {
2185 fn search(&self, query: &str, options: SearchOptions) -> ThingdResult<Vec<SearchHit>> {
2186 #[cfg(feature = "search")]
2188 if let (Some(index), Some(reader)) = (&self.search_index, &self.search_reader) {
2189 return self.search_tantivy(index, reader, query, options);
2190 }
2191
2192 self.search_naive(query, options)
2194 }
2195}
2196
2197impl PersistentEngine {
2198 #[cfg(feature = "search")]
2199 fn index_object_for_search(&self, object: &MemoryObject) {
2200 self.index_object_for_search_with_commit(object, true);
2201 }
2202
2203 #[cfg(feature = "search")]
2204 fn index_object_for_search_with_commit(&self, object: &MemoryObject, commit: bool) {
2205 let Some(ref index) = self.search_index else {
2206 return;
2207 };
2208 let schema = index.schema();
2209 let doc_key_field = schema.get_field("doc_key").unwrap();
2210 let body_field = schema.get_field("body").unwrap();
2211 let collection_field = schema.get_field("collection").unwrap();
2212 let id_field = schema.get_field("id").unwrap();
2213 let kind_field = schema.get_field("kind").unwrap();
2214
2215 let Some(ref writer) = self.search_writer else {
2216 return;
2217 };
2218 let Ok(mut writer) = writer.lock() else {
2219 return;
2220 };
2221
2222 let doc_key = format!("{}/{}", object.key.collection, object.key.id);
2223 let term = tantivy::Term::from_field_text(doc_key_field, &doc_key);
2225 let _ = writer.delete_term(term);
2226
2227 let mut doc = tantivy::TantivyDocument::new();
2228 doc.add_text(doc_key_field, &doc_key);
2229 doc.add_text(collection_field, &object.key.collection);
2230 doc.add_text(id_field, &object.key.id);
2231 doc.add_text(body_field, &object.body);
2232 doc.add_text(kind_field, "object");
2233
2234 let _ = writer.add_document(doc);
2235 if commit {
2236 let _ = writer.commit();
2237 drop(writer);
2238 self.reload_search_reader();
2239 }
2240 }
2241
2242 #[cfg(feature = "search")]
2243 fn index_event_for_search(&self, event: &MemoryEvent) {
2244 self.index_event_for_search_with_commit(event, true);
2245 }
2246
2247 #[cfg(feature = "search")]
2248 fn index_event_for_search_with_commit(&self, event: &MemoryEvent, commit: bool) {
2249 let Some(ref index) = self.search_index else {
2250 return;
2251 };
2252 let schema = index.schema();
2253 let doc_key_field = schema.get_field("doc_key").unwrap();
2254 let body_field = schema.get_field("body").unwrap();
2255 let collection_field = schema.get_field("collection").unwrap();
2256 let id_field = schema.get_field("id").unwrap();
2257 let kind_field = schema.get_field("kind").unwrap();
2258
2259 let Some(ref writer) = self.search_writer else {
2260 return;
2261 };
2262 let Ok(mut writer) = writer.lock() else {
2263 return;
2264 };
2265
2266 let mut doc = tantivy::TantivyDocument::new();
2267 let doc_key = format!("event:{}/{}", event.stream, event.sequence);
2268 doc.add_text(doc_key_field, &doc_key);
2269 doc.add_text(collection_field, &event.stream);
2270 doc.add_text(id_field, event.sequence.to_string());
2271 doc.add_text(body_field, &event.body);
2272 doc.add_text(kind_field, "event");
2273
2274 let _ = writer.add_document(doc);
2275 if commit {
2276 let _ = writer.commit();
2277 drop(writer);
2278 self.reload_search_reader();
2279 }
2280 }
2281
2282 #[cfg(feature = "search")]
2283 fn commit_search_index(&self) {
2284 let Some(ref writer) = self.search_writer else {
2285 return;
2286 };
2287 let Ok(mut writer) = writer.lock() else {
2288 return;
2289 };
2290 let _ = writer.commit();
2291 drop(writer);
2292 self.reload_search_reader();
2293 }
2294
2295 #[cfg(feature = "search")]
2296 fn reload_search_reader(&self) {
2297 if let Some(ref reader) = self.search_reader
2298 && let Ok(reader) = reader.lock()
2299 {
2300 let _ = reader.reload();
2301 }
2302 }
2303
2304 #[cfg(feature = "search")]
2305 fn delete_event_from_search_index(&self, stream: &str, sequence: u64) {
2306 let Some(ref index) = self.search_index else {
2307 return;
2308 };
2309 let schema = index.schema();
2310 let doc_key_field = schema.get_field("doc_key").unwrap();
2311
2312 let Some(ref writer) = self.search_writer else {
2313 return;
2314 };
2315 let Ok(mut writer) = writer.lock() else {
2316 return;
2317 };
2318
2319 let doc_key = format!("event:{stream}/{sequence}");
2320 let term = tantivy::Term::from_field_text(doc_key_field, &doc_key);
2321 let _ = writer.delete_term(term);
2322 let _ = writer.commit();
2323 drop(writer);
2324 self.reload_search_reader();
2325 }
2326
2327 #[cfg(feature = "search")]
2328 fn delete_events_from_search_index(&self, events: &[(String, u64)]) {
2329 if events.is_empty() {
2330 return;
2331 }
2332 let Some(ref index) = self.search_index else {
2333 return;
2334 };
2335 let schema = index.schema();
2336 let doc_key_field = schema.get_field("doc_key").unwrap();
2337 let Some(ref writer) = self.search_writer else {
2338 return;
2339 };
2340 let Ok(mut writer) = writer.lock() else {
2341 return;
2342 };
2343 for (stream, sequence) in events {
2344 let doc_key = format!("event:{stream}/{sequence}");
2345 let term = tantivy::Term::from_field_text(doc_key_field, &doc_key);
2346 let _ = writer.delete_term(term);
2347 }
2348 let _ = writer.commit();
2349 drop(writer);
2350 self.reload_search_reader();
2351 }
2352
2353 #[cfg(feature = "search")]
2354 fn delete_object_from_search_index(&self, collection: &str, id: &str) {
2355 let Some(ref index) = self.search_index else {
2356 return;
2357 };
2358 let schema = index.schema();
2359 let doc_key_field = schema.get_field("doc_key").unwrap();
2360
2361 let Some(ref writer) = self.search_writer else {
2362 return;
2363 };
2364 let Ok(mut writer) = writer.lock() else {
2365 return;
2366 };
2367
2368 let doc_key = format!("{collection}/{id}");
2369 let term = tantivy::Term::from_field_text(doc_key_field, &doc_key);
2370 let _ = writer.delete_term(term);
2371 let _ = writer.commit();
2372 drop(writer);
2373 self.reload_search_reader();
2374 }
2375
2376 #[cfg(feature = "search")]
2377 fn search_tantivy(
2378 &self,
2379 index: &tantivy::Index,
2380 reader: &Mutex<tantivy::IndexReader>,
2381 query: &str,
2382 options: SearchOptions,
2383 ) -> ThingdResult<Vec<SearchHit>> {
2384 use tantivy::collector::TopDocs;
2385 use tantivy::query::QueryParser;
2386 use tantivy::schema::Value;
2387
2388 let reader = reader
2389 .lock()
2390 .map_err(|_| ThingdError::Storage("search reader lock poisoned".to_string()))?;
2391 let searcher = reader.searcher();
2392 let schema = index.schema();
2393
2394 let body_field = schema.get_field("body").unwrap();
2395 let collection_field = schema.get_field("collection").unwrap();
2396 let id_field = schema.get_field("id").unwrap();
2397 let kind_field = schema.get_field("kind").unwrap();
2398
2399 let mut parser = QueryParser::for_index(index, vec![body_field]);
2400 parser.set_conjunction_by_default();
2401
2402 let tantivy_query = parser
2403 .parse_query(query)
2404 .map_err(|e| ThingdError::InvalidInput(e.to_string()))?;
2405
2406 let limit = options.limit.unwrap_or(10).min(1000);
2407 let doc_ids = searcher
2408 .search(&tantivy_query, &TopDocs::with_limit(limit).order_by_score())
2409 .map_err(|e| ThingdError::Storage(e.to_string()))?;
2410
2411 let mut hits = Vec::new();
2412
2413 for (score, doc_address) in doc_ids {
2414 let doc = searcher
2415 .doc::<tantivy::TantivyDocument>(doc_address)
2416 .map_err(|e| ThingdError::Storage(e.to_string()))?;
2417
2418 let collection = doc
2419 .get_first(collection_field)
2420 .and_then(|v| v.as_str())
2421 .unwrap_or("")
2422 .to_string();
2423 let id = doc
2424 .get_first(id_field)
2425 .and_then(|v| v.as_str())
2426 .unwrap_or("")
2427 .to_string();
2428 let body = doc
2429 .get_first(body_field)
2430 .and_then(|v| v.as_str())
2431 .unwrap_or("")
2432 .to_string();
2433 let kind = doc
2434 .get_first(kind_field)
2435 .and_then(|v| v.as_str())
2436 .unwrap_or("object")
2437 .to_string();
2438
2439 if let Some(ref collections) = options.collections
2440 && !collections.contains(&collection)
2441 {
2442 continue;
2443 }
2444
2445 if let Some(ref filter) = options.filter
2446 && !matches_filter_memory(&body, filter)
2447 {
2448 continue;
2449 }
2450
2451 let col = collection.clone();
2452 let doc_id = id.clone();
2453 let doc_body = body.clone();
2454
2455 if kind == "object"
2456 && let Some(obj_data) = self
2457 .objects
2458 .get(self.make_object_key(&col, &doc_id))
2459 .ok()
2460 .and_then(value_to_vec)
2461 && let Ok(obj) = self.deserialize::<MemoryObject>(&obj_data)
2462 {
2463 hits.push(SearchHit {
2464 kind: "object".to_string(),
2465 collection: col,
2466 id: doc_id,
2467 text: doc_body.clone(),
2468 score: score.into(),
2469 body: doc_body,
2470 version: Some(obj.version),
2471 created_at: obj.created_at,
2472 updated_at: Some(obj.updated_at),
2473 event_type: None,
2474 });
2475 } else if kind == "event"
2476 && let Ok(seq) = id.parse::<u64>()
2477 && let Some(ev_data) = self
2478 .events
2479 .get(self.make_event_key(&collection, seq))
2480 .ok()
2481 .and_then(value_to_vec)
2482 && let Ok(event) = self.deserialize::<MemoryEvent>(&ev_data)
2483 {
2484 hits.push(SearchHit {
2485 kind: "event".to_string(),
2486 collection: collection.clone(),
2487 id: seq.to_string(),
2488 text: body.clone(),
2489 score: score.into(),
2490 body: body.clone(),
2491 version: None,
2492 created_at: event.created_at,
2493 updated_at: None,
2494 event_type: Some(event.event_type),
2495 });
2496 }
2497 }
2498
2499 drop(reader);
2500 Ok(hits)
2501 }
2502
2503 fn search_naive(&self, query: &str, options: SearchOptions) -> ThingdResult<Vec<SearchHit>> {
2504 let query_words: Vec<String> = query
2505 .split_whitespace()
2506 .map(|w| {
2507 w.to_lowercase()
2508 .chars()
2509 .filter(|c| c.is_alphanumeric())
2510 .collect()
2511 })
2512 .filter(|w: &String| !w.is_empty())
2513 .collect();
2514
2515 if query_words.is_empty() {
2516 return Ok(Vec::new());
2517 }
2518
2519 let mut hits = Vec::new();
2520
2521 for kv in self.objects.iter() {
2522 let (_, value) = guard_data(kv)?;
2523 let object: MemoryObject = self.deserialize(&value)?;
2524
2525 if let Some(ref collections) = options.collections
2526 && !collections.contains(&object.key.collection)
2527 {
2528 continue;
2529 }
2530
2531 if let Some(ref filter) = options.filter
2532 && !matches_filter_memory(&object.body, filter)
2533 {
2534 continue;
2535 }
2536
2537 let text_to_search = format!(
2538 "{} {} {}",
2539 object.key.collection, object.key.id, object.body
2540 )
2541 .to_lowercase();
2542 let matches_all = query_words.iter().all(|word| text_to_search.contains(word));
2543
2544 if matches_all {
2545 hits.push(SearchHit {
2546 kind: "object".to_string(),
2547 collection: object.key.collection.clone(),
2548 id: object.key.id.clone(),
2549 text: object.body.clone(),
2550 score: 1.0,
2551 body: object.body.clone(),
2552 version: Some(object.version),
2553 created_at: object.created_at.clone(),
2554 updated_at: Some(object.updated_at.clone()),
2555 event_type: None,
2556 });
2557 }
2558 }
2559
2560 for kv in self.events.iter() {
2561 let (_, value) = guard_data(kv)?;
2562 let event: MemoryEvent = self.deserialize(&value)?;
2563
2564 if let Some(ref collections) = options.collections
2565 && !collections.contains(&event.stream)
2566 {
2567 continue;
2568 }
2569
2570 if let Some(ref filter) = options.filter
2571 && !matches_filter_memory(&event.body, filter)
2572 {
2573 continue;
2574 }
2575
2576 let text_to_search =
2577 format!("{} {} {}", event.stream, event.event_type, event.body).to_lowercase();
2578 let matches_all = query_words.iter().all(|word| text_to_search.contains(word));
2579
2580 if matches_all {
2581 hits.push(SearchHit {
2582 kind: "event".to_string(),
2583 collection: event.stream.clone(),
2584 id: event.sequence.to_string(),
2585 text: event.body.clone(),
2586 score: 1.0,
2587 body: event.body.clone(),
2588 version: None,
2589 created_at: event.created_at.clone(),
2590 updated_at: None,
2591 event_type: Some(event.event_type.clone()),
2592 });
2593 }
2594 }
2595
2596 if let Some(limit) = options.limit {
2597 hits.truncate(limit);
2598 }
2599
2600 Ok(hits)
2601 }
2602}
2603
2604impl LinkStore for PersistentEngine {
2607 fn create_link(&mut self, mut link: Link) -> ThingdResult<Link> {
2608 let id = self.next_link_id.fetch_add(1, Ordering::Relaxed);
2609 link.id = format!("link-{id}");
2610 if link.created_at.is_empty() {
2611 link.created_at = now_iso_string();
2612 }
2613
2614 let data = self.serialize(&link)?;
2615
2616 let link_id_key = self.make_link_id_key(&link.id);
2617 let index_data = self.serialize(&link.id)?;
2618 self.links_by_id.insert(&link_id_key, &data)?;
2619 self.links_from.insert(
2620 self.make_link_from_key(&link.from_ref, &link.link_type, &link.id),
2621 &index_data,
2622 )?;
2623 self.links_to.insert(
2624 self.make_link_to_key(&link.to_ref, &link.link_type, &link.id),
2625 &index_data,
2626 )?;
2627
2628 Ok(link)
2629 }
2630
2631 fn delete_link(&mut self, id: &str) -> ThingdResult<bool> {
2632 let link_id_key = self.make_link_id_key(id);
2633 if let Some(data) = value_to_vec(self.links_by_id.get(&link_id_key)?) {
2634 let link: Link = self.deserialize(&data)?;
2635 self.links_by_id.remove(&link_id_key)?;
2636 self.links_from
2637 .remove(self.make_link_from_key(&link.from_ref, &link.link_type, id))?;
2638 self.links_to
2639 .remove(self.make_link_to_key(&link.to_ref, &link.link_type, id))?;
2640 Ok(true)
2641 } else {
2642 Ok(false)
2643 }
2644 }
2645
2646 fn get_link(&self, id: &str) -> ThingdResult<Option<Link>> {
2647 match value_to_vec(self.links_by_id.get(self.make_link_id_key(id))?) {
2648 Some(data) => Ok(Some(self.deserialize(&data)?)),
2649 None => Ok(None),
2650 }
2651 }
2652
2653 fn get_neighbors(
2654 &self,
2655 reference: &str,
2656 direction: LinkDirection,
2657 options: LinkQueryOptions,
2658 ) -> ThingdResult<Vec<Link>> {
2659 let mut results: Vec<Link> = Vec::new();
2660 let outgoing_prefix = self.make_link_from_prefix(reference);
2661 let incoming_prefix = self.make_link_to_prefix(reference);
2662
2663 if direction == LinkDirection::Outgoing || direction == LinkDirection::Both {
2664 for kv in self.links_from.prefix(&outgoing_prefix) {
2665 let (key, value) = guard_data(kv)?;
2666 let link_id = if value.is_empty() {
2667 let key_str = String::from_utf8_lossy(&key);
2668 key_str.splitn(3, '\0').nth(2).unwrap_or("").to_string()
2669 } else {
2670 self.deserialize::<String>(&value)?
2671 };
2672 if let Some(data) =
2673 value_to_vec(self.links_by_id.get(self.make_link_id_key(&link_id))?)
2674 {
2675 let link: Link = self.deserialize(&data)?;
2676 if options
2677 .link_type
2678 .as_ref()
2679 .is_none_or(|link_type| link.link_type == *link_type)
2680 {
2681 results.push(link);
2682 }
2683 }
2684 }
2685 }
2686
2687 if direction == LinkDirection::Incoming || direction == LinkDirection::Both {
2688 for kv in self.links_to.prefix(&incoming_prefix) {
2689 let (key, value) = guard_data(kv)?;
2690 let link_id = if value.is_empty() {
2691 let key_str = String::from_utf8_lossy(&key);
2692 key_str.splitn(3, '\0').nth(2).unwrap_or("").to_string()
2693 } else {
2694 self.deserialize::<String>(&value)?
2695 };
2696 if let Some(data) =
2697 value_to_vec(self.links_by_id.get(self.make_link_id_key(&link_id))?)
2698 {
2699 let link: Link = self.deserialize(&data)?;
2700 if options
2701 .link_type
2702 .as_ref()
2703 .is_none_or(|link_type| link.link_type == *link_type)
2704 {
2705 results.push(link);
2706 }
2707 }
2708 }
2709 }
2710
2711 if let Some(limit) = options.limit {
2712 results.truncate(limit);
2713 }
2714
2715 Ok(results)
2716 }
2717
2718 fn count_links(&self) -> ThingdResult<u64> {
2719 let mut count = 0u64;
2720 for kv in self.links_by_id.iter() {
2721 let _ = kv;
2722 count += 1;
2723 }
2724 Ok(count)
2725 }
2726}
2727
2728impl AggregateStore for PersistentEngine {
2731 fn aggregate(
2732 &self,
2733 collection: &str,
2734 options: &AggregateOptions,
2735 ) -> ThingdResult<AggregateResult> {
2736 let prefix = self.make_object_prefix(collection);
2737 let mut objects: Vec<MemoryObject> = Vec::new();
2738
2739 for kv in self.objects.prefix(&prefix) {
2740 let (_, value) = guard_data(kv)?;
2741 let obj: MemoryObject = self.deserialize(&value)?;
2742
2743 if options.filter.is_empty() {
2744 objects.push(obj);
2745 } else {
2746 let Ok(body) = serde_json::from_str::<serde_json::Value>(&obj.body) else {
2747 continue;
2748 };
2749 let matches = options
2750 .filter
2751 .iter()
2752 .all(|(key, expected)| body.get(key.as_str()).is_some_and(|v| v == expected));
2753 if matches {
2754 objects.push(obj);
2755 }
2756 }
2757 }
2758
2759 if let Some(group_field) = &options.group_by {
2760 let mut groups: HashMap<String, Vec<MemoryObject>> = HashMap::new();
2761 for obj in &objects {
2762 let key = extract_field_str(&obj.body, group_field);
2763 groups.entry(key).or_default().push(obj.clone());
2764 }
2765
2766 let mut group_results: Vec<AggregateGroupResult> = groups
2767 .iter()
2768 .map(|(key, objs)| AggregateGroupResult {
2769 key: key.clone(),
2770 value: compute_aggregate(objs, options.function, options.field.as_deref()),
2771 })
2772 .collect();
2773 group_results.sort_by(|a, b| {
2774 b.value
2775 .partial_cmp(&a.value)
2776 .unwrap_or(std::cmp::Ordering::Equal)
2777 });
2778
2779 let total = group_results.iter().map(|g| g.value).sum();
2780
2781 Ok(AggregateResult {
2782 total,
2783 groups: group_results,
2784 })
2785 } else {
2786 let total = compute_aggregate(&objects, options.function, options.field.as_deref());
2787 Ok(AggregateResult {
2788 total,
2789 groups: vec![],
2790 })
2791 }
2792 }
2793
2794 fn timeseries(
2795 &self,
2796 collection: &str,
2797 options: &TimeSeriesOptions,
2798 ) -> ThingdResult<TimeSeriesResult> {
2799 let prefix = self.make_object_prefix(collection);
2800 let mut bucket_map: HashMap<String, Vec<f64>> = HashMap::new();
2801
2802 for kv in self.objects.prefix(&prefix) {
2803 let (_, value) = guard_data(kv)?;
2804 let obj: MemoryObject = self.deserialize(&value)?;
2805
2806 if !options.filter.is_empty() {
2807 let Ok(body) = serde_json::from_str::<serde_json::Value>(&obj.body) else {
2808 continue;
2809 };
2810 let matches = options
2811 .filter
2812 .iter()
2813 .all(|(key, expected)| body.get(key.as_str()).is_some_and(|v| v == expected));
2814 if !matches {
2815 continue;
2816 }
2817 }
2818
2819 let bucket_label = bucket_label_for_date(&obj.created_at, options.bucket);
2820
2821 if let Some(ref from) = options.from
2822 && bucket_label.as_str() < from.as_str()
2823 {
2824 continue;
2825 }
2826 if let Some(ref to) = options.to
2827 && bucket_label.as_str() > to.as_str()
2828 {
2829 continue;
2830 }
2831
2832 if options.function == AggregateFunction::Count {
2833 bucket_map.entry(bucket_label).or_default().push(1.0);
2834 } else if let Some(ref field) = options.field
2835 && let Ok(body) = serde_json::from_str::<serde_json::Value>(&obj.body)
2836 && let Some(val) = body.get(field.as_str()).and_then(serde_json::Value::as_f64)
2837 {
2838 bucket_map.entry(bucket_label).or_default().push(val);
2839 }
2840 }
2841
2842 let mut bucket_list: Vec<(String, f64)> = bucket_map
2843 .into_iter()
2844 .map(|(label, values)| {
2845 let value = match options.function {
2846 AggregateFunction::Count => values.len() as f64,
2847 AggregateFunction::Sum => values.iter().sum(),
2848 AggregateFunction::Avg => {
2849 if values.is_empty() {
2850 0.0
2851 } else {
2852 values.iter().sum::<f64>() / values.len() as f64
2853 }
2854 },
2855 AggregateFunction::Min => values.iter().copied().fold(f64::MAX, f64::min),
2856 AggregateFunction::Max => {
2857 values.iter().copied().fold(f64::NEG_INFINITY, f64::max)
2858 },
2859 };
2860 (label, value)
2861 })
2862 .collect();
2863
2864 bucket_list.sort_by(|a, b| a.0.cmp(&b.0));
2865
2866 let buckets: Vec<TimeSeriesBucket> = bucket_list
2867 .into_iter()
2868 .map(|(label, value)| TimeSeriesBucket { label, value })
2869 .collect();
2870
2871 Ok(TimeSeriesResult { buckets })
2872 }
2873}
2874
2875impl crate::store::VectorStore for PersistentEngine {
2878 fn vector_search(
2879 &self,
2880 collection: &str,
2881 query_vector: &[f32],
2882 options: VectorSearchOptions,
2883 ) -> ThingdResult<Vec<VectorSearchHit>> {
2884 if query_vector.is_empty() {
2885 return Err(ThingdError::InvalidInput(
2886 "query vector must not be empty".to_string(),
2887 ));
2888 }
2889
2890 #[cfg(not(feature = "vectors"))]
2891 {
2892 let _ = (collection, query_vector, options);
2893 Ok(vec![])
2894 }
2895
2896 #[cfg(feature = "vectors")]
2897 {
2898 let prefix = self.make_vector_prefix(collection);
2899 let mut hits: Vec<VectorSearchHit> = Vec::new();
2900
2901 for kv in self.vectors.prefix(&prefix) {
2902 let (physical_key, value) = guard_data(kv)?;
2903 let stored = if let Ok(stored) = self.deserialize::<StoredVector>(&value) {
2904 stored
2905 } else {
2906 let Some(separator) = physical_key.iter().rposition(|byte| *byte == 0) else {
2907 return Err(ThingdError::Storage(
2908 "legacy vector record has no collection/id separator".to_string(),
2909 ));
2910 };
2911 StoredVector {
2912 collection: collection.to_string(),
2913 id: String::from_utf8_lossy(&physical_key[separator + 1..]).into_owned(),
2914 vector: self.deserialize(&value)?,
2915 }
2916 };
2917 let id = stored.id;
2918 let vector = stored.vector;
2919
2920 if vector.len() != query_vector.len() {
2921 return Err(ThingdError::InvalidInput(format!(
2922 "query vector dimension {} does not match stored vector dimension {}",
2923 query_vector.len(),
2924 vector.len()
2925 )));
2926 }
2927
2928 let Some(object) = self.get_object(collection, &id)? else {
2929 continue;
2930 };
2931
2932 if let Some(ref filter) = options.filter
2933 && !matches_filter_memory(&object.body, filter)
2934 {
2935 continue;
2936 }
2937
2938 let score = crate::cosine_similarity(query_vector, &vector);
2939 hits.push(VectorSearchHit {
2940 id: id.clone(),
2941 score,
2942 value: object,
2943 });
2944 }
2945
2946 hits.sort_by(|a, b| {
2947 b.score
2948 .partial_cmp(&a.score)
2949 .unwrap_or(std::cmp::Ordering::Equal)
2950 });
2951
2952 if let Some(top_k) = options.top_k {
2953 hits.truncate(top_k);
2954 }
2955
2956 Ok(hits)
2957 }
2958 }
2959
2960 fn add_vector(&mut self, collection: &str, id: &str, vector: &[f32]) -> ThingdResult<()> {
2961 #[cfg(not(feature = "vectors"))]
2962 {
2963 let _ = (collection, id, vector);
2964 }
2965
2966 #[cfg(feature = "vectors")]
2967 {
2968 let vkey = self.make_vector_key(collection, id);
2969 let vdata = self.serialize(&StoredVector {
2970 collection: collection.to_string(),
2971 id: id.to_string(),
2972 vector: vector.to_vec(),
2973 })?;
2974 self.vectors.insert(&vkey, &vdata)?;
2975 }
2976
2977 Ok(())
2978 }
2979
2980 fn remove_vector(&mut self, collection: &str, id: &str) -> ThingdResult<()> {
2981 #[cfg(not(feature = "vectors"))]
2982 {
2983 let _ = (collection, id);
2984 }
2985
2986 #[cfg(feature = "vectors")]
2987 {
2988 let vkey = self.make_vector_key(collection, id);
2989 let _ = self.vectors.remove(&vkey);
2990 }
2991
2992 Ok(())
2993 }
2994}
2995
2996fn value_compare(a: &serde_json::Value, b: &serde_json::Value) -> std::cmp::Ordering {
2999 match (a, b) {
3000 (serde_json::Value::Number(a), serde_json::Value::Number(b)) => {
3001 let a_f = a.as_f64().unwrap_or(0.0);
3002 let b_f = b.as_f64().unwrap_or(0.0);
3003 a_f.partial_cmp(&b_f).unwrap_or(std::cmp::Ordering::Equal)
3004 },
3005 (serde_json::Value::String(a), serde_json::Value::String(b)) => a.cmp(b),
3006 (serde_json::Value::Bool(a), serde_json::Value::Bool(b)) => a.cmp(b),
3007 _ => format!("{a}").cmp(&format!("{b}")),
3008 }
3009}
3010
3011fn like_match(s: &str, pattern: &str) -> bool {
3012 let parts = pattern.split('%');
3013 let mut pos = 0;
3014 for part in parts {
3015 if part.is_empty() {
3016 continue;
3017 }
3018 if let Some(idx) = s[pos..].find(part) {
3019 pos += idx + part.len();
3020 } else {
3021 return false;
3022 }
3023 }
3024 true
3025}
3026
3027fn matches_filter_memory(body_str: &str, filter: &serde_json::Value) -> bool {
3028 let Ok(body) = serde_json::from_str::<serde_json::Value>(body_str) else {
3029 return false;
3030 };
3031 if let serde_json::Value::Object(map) = filter {
3032 map.iter()
3033 .all(|(key, expected)| body.get(key.as_str()).is_some_and(|v| v == expected))
3034 } else {
3035 false
3036 }
3037}
3038
3039fn matches_object_filters(object: &MemoryObject, filters: &[(String, serde_json::Value)]) -> bool {
3040 if filters.is_empty() {
3041 return true;
3042 }
3043 let Ok(body) = serde_json::from_str::<serde_json::Value>(&object.body) else {
3044 return false;
3045 };
3046 filters.iter().all(|(key, expected)| {
3047 let field_val = body.get(key.as_str());
3048 match expected {
3049 serde_json::Value::Object(ops)
3050 if ops.keys().any(|k| {
3051 matches!(
3052 k.as_str(),
3053 "$gt" | "$gte" | "$lt" | "$lte" | "$ne" | "$in" | "$like"
3054 )
3055 }) =>
3056 {
3057 let Some(fv) = field_val else {
3058 return false;
3059 };
3060 ops.iter().all(|(op, operand)| match op.as_str() {
3061 "$gt" => value_compare(fv, operand) == std::cmp::Ordering::Greater,
3062 "$gte" => matches!(
3063 value_compare(fv, operand),
3064 std::cmp::Ordering::Greater | std::cmp::Ordering::Equal
3065 ),
3066 "$lt" => value_compare(fv, operand) == std::cmp::Ordering::Less,
3067 "$lte" => matches!(
3068 value_compare(fv, operand),
3069 std::cmp::Ordering::Less | std::cmp::Ordering::Equal
3070 ),
3071 "$ne" => value_compare(fv, operand) != std::cmp::Ordering::Equal,
3072 "$in" => operand
3073 .as_array()
3074 .is_some_and(|items| items.iter().any(|item| fv == item)),
3075 "$like" => {
3076 if let (Some(s), Some(pattern)) = (fv.as_str(), operand.as_str()) {
3077 like_match(s, pattern)
3078 } else {
3079 false
3080 }
3081 },
3082 _ => true,
3083 })
3084 },
3085 _ => field_val.is_some_and(|value| value == expected),
3086 }
3087 })
3088}
3089
3090fn extract_field_str(body_str: &str, field: &str) -> String {
3091 if let Ok(body) = serde_json::from_str::<serde_json::Value>(body_str)
3092 && let Some(val) = body.get(field)
3093 {
3094 if let Some(s) = val.as_str() {
3095 return s.to_string();
3096 }
3097 return format!("{val}");
3098 }
3099 String::new()
3100}
3101
3102fn compute_aggregate(
3103 objects: &[MemoryObject],
3104 function: AggregateFunction,
3105 field: Option<&str>,
3106) -> f64 {
3107 if function == AggregateFunction::Count {
3108 objects.len() as f64
3109 } else {
3110 let values: Vec<f64> = objects
3111 .iter()
3112 .filter_map(|obj| {
3113 field.and_then(|f| {
3114 let Ok(body) = serde_json::from_str::<serde_json::Value>(&obj.body) else {
3115 return None;
3116 };
3117 body.get(f).and_then(serde_json::Value::as_f64)
3118 })
3119 })
3120 .collect();
3121
3122 match function {
3123 AggregateFunction::Sum => values.iter().sum(),
3124 AggregateFunction::Avg => {
3125 if values.is_empty() {
3126 0.0
3127 } else {
3128 values.iter().sum::<f64>() / values.len() as f64
3129 }
3130 },
3131 AggregateFunction::Min => values.iter().copied().fold(f64::MAX, f64::min),
3132 AggregateFunction::Max => values.iter().copied().fold(f64::NEG_INFINITY, f64::max),
3133 _ => values.len() as f64,
3134 }
3135 }
3136}
3137
3138fn infer_json_type(value: &serde_json::Value) -> String {
3139 match value {
3140 serde_json::Value::Null => "null".to_string(),
3141 serde_json::Value::Bool(_) => "boolean".to_string(),
3142 serde_json::Value::Number(_) => "number".to_string(),
3143 serde_json::Value::String(s) => {
3144 if s.parse::<chrono::DateTime<chrono::Utc>>().is_ok() {
3145 "date".to_string()
3146 } else {
3147 "string".to_string()
3148 }
3149 },
3150 serde_json::Value::Array(_) => "array".to_string(),
3151 serde_json::Value::Object(_) => "object".to_string(),
3152 }
3153}
3154
3155fn bucket_label_for_date(iso_date: &str, bucket: TimeBucket) -> String {
3156 if let Ok(dt) = chrono::DateTime::parse_from_rfc3339(iso_date) {
3157 match bucket {
3158 TimeBucket::Hour => dt.format("%Y-%m-%dT%H:00:00Z").to_string(),
3159 TimeBucket::Day => dt.format("%Y-%m-%d").to_string(),
3160 TimeBucket::Week => dt.format("%Y-W%V").to_string(),
3161 TimeBucket::Month => dt.format("%Y-%m").to_string(),
3162 }
3163 } else {
3164 iso_date.to_string()
3165 }
3166}
3167
3168#[cfg(test)]
3169#[allow(clippy::float_cmp, clippy::cast_precision_loss)]
3170mod tests {
3171 use super::*;
3172 #[cfg(feature = "vectors")]
3173 use crate::VectorSearchOptions;
3174 #[cfg(feature = "vectors")]
3175 use crate::store::VectorStore;
3176 use crate::store::{AggregateStore, EventLog, LinkStore, ObjectStore, QueueStore, Searcher};
3177 use crate::{
3178 Link, ListObjectsOptions, MemoryEvent, MemoryObject, QueueClaimOptions, QueueJob,
3179 QueueJobStatus, QueueNackOptions, SearchOptions, TimeBucket,
3180 };
3181
3182 fn setup() -> (PersistentEngine, tempfile::TempDir) {
3184 let dir = tempfile::tempdir().unwrap();
3185 let engine = PersistentEngine::open(dir.path()).unwrap();
3186 (engine, dir)
3187 }
3188
3189 #[test]
3190 fn persistent_open_writes_and_validates_storage_manifest() {
3191 let dir = tempfile::tempdir().unwrap();
3192 let _engine = PersistentEngine::open(dir.path()).unwrap();
3193 let report = PersistentEngine::validate_path(dir.path()).unwrap();
3194 assert_eq!(report.format_version, STORAGE_FORMAT_VERSION);
3195 assert!(!report.legacy_manifest);
3196 assert!(report.lock_present);
3197 assert!(report.keyspaces_present);
3198 assert!(dir.path().join(STORAGE_MANIFEST_FILE).is_file());
3199 }
3200
3201 #[test]
3202 fn persistent_retention_is_dry_run_by_default_and_preserves_protected_events() {
3203 let (mut engine, _dir) = setup();
3204 let mut old = MemoryEvent::new("old", "test", "{}");
3205 old.created_at = "2020-01-01T00:00:00Z".to_string();
3206 engine.append_event(old).unwrap();
3207 let mut protected = MemoryEvent::new("__thingd:system", "test", "{}");
3208 protected.created_at = "2020-01-01T00:00:00Z".to_string();
3209 engine.append_event(protected).unwrap();
3210 let mut replication = MemoryEvent::new(REPLICATION_STREAM, "object.upsert", "{}");
3211 replication.created_at = "2020-01-01T00:00:00Z".to_string();
3212 engine.append_event(replication).unwrap();
3213
3214 let preview = engine
3215 .retain(RetentionOptions {
3216 before_unix_ms: 1_700_000_000_000,
3217 dry_run: true,
3218 compact: false,
3219 include_replication: false,
3220 })
3221 .unwrap();
3222 assert_eq!(preview.events, 1);
3223 assert_eq!(preview.skipped_replication_events, 1);
3224 assert_eq!(preview.safe_replication_cursor, None);
3225 assert_eq!(engine.count_events().unwrap(), 3);
3226
3227 engine
3228 .put_object(MemoryObject::new(
3229 REPLICATION_STATE_COLLECTION,
3230 "source:replica-a",
3231 r#"{"sourceId":"replica-a","lastAppliedCursor":1}"#,
3232 ))
3233 .unwrap();
3234 let checkpointed = engine
3235 .retain(RetentionOptions {
3236 before_unix_ms: 1_700_000_000_000,
3237 dry_run: true,
3238 compact: false,
3239 include_replication: true,
3240 })
3241 .unwrap();
3242 assert_eq!(checkpointed.events, 2);
3243 assert_eq!(checkpointed.safe_replication_cursor, Some(1));
3244
3245 let deleted = engine
3246 .retain(RetentionOptions {
3247 before_unix_ms: 1_700_000_000_000,
3248 dry_run: false,
3249 compact: true,
3250 include_replication: false,
3251 })
3252 .unwrap();
3253 assert_eq!(deleted.events, 1);
3254 assert!(deleted.compacted);
3255 assert_eq!(engine.count_events().unwrap(), 2);
3256 assert_eq!(
3257 engine
3258 .list_events(Some("__thingd:system"), ListEventsOptions::default())
3259 .unwrap()
3260 .len(),
3261 1
3262 );
3263 }
3264
3265 #[test]
3266 fn validation_rejects_existing_directory_without_lock() {
3267 let dir = tempfile::tempdir().unwrap();
3268 std::fs::create_dir(dir.path().join(STORAGE_KEYSPACES_DIR)).unwrap();
3269 std::fs::write(dir.path().join("0.jnl"), []).unwrap();
3270 let result = PersistentEngine::validate_path(dir.path());
3271 assert!(matches!(
3272 result,
3273 Err(ThingdError::StorageValidation(message)) if message.contains("lock")
3274 ));
3275 }
3276
3277 #[test]
3278 fn validation_rejects_unsupported_manifest() {
3279 let dir = tempfile::tempdir().unwrap();
3280 std::fs::create_dir(dir.path().join(STORAGE_KEYSPACES_DIR)).unwrap();
3281 std::fs::write(dir.path().join(STORAGE_LOCK_FILE), []).unwrap();
3282 std::fs::write(
3283 dir.path().join(STORAGE_MANIFEST_FILE),
3284 r#"{"format_version":99,"contract":"future","keyspaces":[],"search_schema_version":1}"#,
3285 )
3286 .unwrap();
3287 let result = PersistentEngine::validate_path(dir.path());
3288 assert!(matches!(
3289 result,
3290 Err(ThingdError::UnsupportedStorageFormat(message)) if message.contains("format version")
3291 ));
3292 }
3293
3294 #[cfg(feature = "search")]
3295 #[test]
3296 fn disabled_search_mode_does_not_create_search_directory() {
3297 let dir = tempfile::tempdir().unwrap();
3298 let options = PersistentOpenOptions {
3299 search_mode: PersistentSearchMode::Disabled,
3300 ..PersistentOpenOptions::default()
3301 };
3302 let mut engine = PersistentEngine::open_with_options(dir.path(), options).unwrap();
3303 engine
3304 .put_object(MemoryObject::new("notes", "one", r#"{"text":"hello"}"#))
3305 .unwrap();
3306 assert!(!dir.path().join("search").exists());
3307 assert_eq!(
3308 engine
3309 .search("hello", SearchOptions::default())
3310 .unwrap()
3311 .len(),
3312 1
3313 );
3314 }
3315
3316 #[cfg(feature = "search")]
3317 #[test]
3318 fn no_rebuild_search_mode_uses_fallback_without_an_index() {
3319 let dir = tempfile::tempdir().unwrap();
3320 let options = PersistentOpenOptions {
3321 search_mode: PersistentSearchMode::PersistentNoRebuild,
3322 ..PersistentOpenOptions::default()
3323 };
3324 let mut engine = PersistentEngine::open_with_options(dir.path(), options.clone()).unwrap();
3325 engine
3326 .put_object(MemoryObject::new("notes", "one", r#"{"text":"hello"}"#))
3327 .unwrap();
3328 drop(engine);
3329 let reopened = PersistentEngine::open_with_options(dir.path(), options).unwrap();
3330 assert!(!dir.path().join("search").exists());
3331 assert_eq!(
3332 reopened
3333 .search("hello", SearchOptions::default())
3334 .unwrap()
3335 .len(),
3336 1
3337 );
3338 }
3339
3340 #[test]
3341 fn encrypted_database_reopens_and_rejects_missing_or_wrong_keys() {
3342 let dir = tempfile::tempdir().unwrap();
3343 let key = [9_u8; 32];
3344 let options = PersistentOpenOptions {
3345 encryption: Some(EncryptionConfig::from_key(&key).unwrap()),
3346 ..PersistentOpenOptions::default()
3347 };
3348 {
3349 let mut engine =
3350 PersistentEngine::open_with_options(dir.path(), options.clone()).unwrap();
3351 engine
3352 .put_object(MemoryObject::new("private", "id", r#"{"secret":"value"}"#))
3353 .unwrap();
3354 }
3355 let missing =
3356 PersistentEngine::open_with_options(dir.path(), PersistentOpenOptions::default());
3357 assert!(matches!(missing, Err(ThingdError::EncryptionRequired(_))));
3358 let wrong = PersistentEngine::open_with_options(
3359 dir.path(),
3360 PersistentOpenOptions {
3361 encryption: Some(EncryptionConfig::from_key(&[8_u8; 32]).unwrap()),
3362 ..PersistentOpenOptions::default()
3363 },
3364 );
3365 assert!(matches!(
3366 wrong,
3367 Err(ThingdError::EncryptionAuthentication(_))
3368 ));
3369 let engine = PersistentEngine::open_with_options(dir.path(), options).unwrap();
3370 assert_eq!(
3371 engine.get_object("private", "id").unwrap().unwrap().body,
3372 r#"{"secret":"value"}"#
3373 );
3374 }
3375
3376 #[cfg(feature = "search")]
3377 #[test]
3378 fn encrypted_search_rebuilds_in_memory_without_search_directory() {
3379 let dir = tempfile::tempdir().unwrap();
3380 let key = [7_u8; 32];
3381 let options = PersistentOpenOptions {
3382 encryption: Some(EncryptionConfig::from_key(&key).unwrap()),
3383 ..PersistentOpenOptions::default()
3384 };
3385 {
3386 let mut engine =
3387 PersistentEngine::open_with_options(dir.path(), options.clone()).unwrap();
3388 engine
3389 .put_object(MemoryObject::new(
3390 "private_collection",
3391 "private_id",
3392 r#"{"body":"encrypted_search_term"}"#,
3393 ))
3394 .unwrap();
3395 engine
3396 .append_event(MemoryEvent::new(
3397 "private_stream",
3398 "event-id",
3399 "encrypted event",
3400 ))
3401 .unwrap();
3402 engine
3403 .push_job(QueueJob::new(
3404 "private_queue",
3405 "private_job",
3406 "private payload",
3407 2,
3408 ))
3409 .unwrap();
3410 engine
3411 .put_object(MemoryObject::new("private_nodes", "a", "{}"))
3412 .unwrap();
3413 engine
3414 .put_object(MemoryObject::new("private_nodes", "b", "{}"))
3415 .unwrap();
3416 engine
3417 .create_link(Link::new(
3418 "private_nodes/a",
3419 "private_link",
3420 "private_nodes/b",
3421 ))
3422 .unwrap();
3423 assert_eq!(
3424 engine
3425 .search("encrypted_search_term", SearchOptions::default())
3426 .unwrap()
3427 .len(),
3428 1
3429 );
3430 }
3431 assert!(!dir.path().join("search").exists());
3432 for bytes in walk_files(dir.path()) {
3433 assert!(
3434 !bytes
3435 .windows("private_collection".len())
3436 .any(|w| w == b"private_collection")
3437 );
3438 assert!(
3439 !bytes
3440 .windows("private_id".len())
3441 .any(|w| w == b"private_id")
3442 );
3443 assert!(
3444 !bytes
3445 .windows("encrypted_search_term".len())
3446 .any(|w| w == b"encrypted_search_term")
3447 );
3448 for identifier in [
3449 "private_stream",
3450 "private_queue",
3451 "private_job",
3452 "private_nodes",
3453 "private_link",
3454 "private payload",
3455 ] {
3456 assert!(
3457 !bytes
3458 .windows(identifier.len())
3459 .any(|w| w == identifier.as_bytes())
3460 );
3461 }
3462 }
3463
3464 let mut reopened = PersistentEngine::open_with_options(dir.path(), options).unwrap();
3465 assert_eq!(
3466 reopened
3467 .search("encrypted_search_term", SearchOptions::default())
3468 .unwrap()
3469 .len(),
3470 1
3471 );
3472 reopened
3473 .delete_object("private_collection", "private_id")
3474 .unwrap();
3475 assert!(
3476 reopened
3477 .search("encrypted_search_term", SearchOptions::default())
3478 .unwrap()
3479 .is_empty()
3480 );
3481 }
3482
3483 fn walk_files(path: &std::path::Path) -> Vec<Vec<u8>> {
3484 let mut files = Vec::new();
3485 let Ok(entries) = std::fs::read_dir(path) else {
3486 return files;
3487 };
3488 for entry in entries.flatten() {
3489 let path = entry.path();
3490 if path.is_dir() {
3491 files.extend(walk_files(&path));
3492 } else if let Ok(bytes) = std::fs::read(path) {
3493 files.push(bytes);
3494 }
3495 }
3496 files
3497 }
3498
3499 #[test]
3500 fn reencrypts_without_modifying_source() {
3501 let root = tempfile::tempdir().unwrap();
3502 let source_path = root.path().join("source");
3503 let destination_path = root.path().join("encrypted");
3504 {
3505 let mut source = PersistentEngine::open(&source_path).unwrap();
3506 source
3507 .put_object(MemoryObject::new("users", "alice", r#"{"name":"Alice"}"#))
3508 .unwrap();
3509 }
3510 let destination_options = PersistentOpenOptions {
3511 encryption: Some(EncryptionConfig::from_key(&[3_u8; 32]).unwrap()),
3512 ..PersistentOpenOptions::default()
3513 };
3514 PersistentEngine::reencrypt_to(
3515 &source_path,
3516 &destination_path,
3517 PersistentOpenOptions::default(),
3518 destination_options.clone(),
3519 )
3520 .unwrap();
3521 let source = PersistentEngine::open(&source_path).unwrap();
3522 assert!(source.get_object("users", "alice").unwrap().is_some());
3523 let destination =
3524 PersistentEngine::open_with_options(&destination_path, destination_options).unwrap();
3525 assert_eq!(
3526 destination
3527 .get_object("users", "alice")
3528 .unwrap()
3529 .unwrap()
3530 .body,
3531 r#"{"name":"Alice"}"#
3532 );
3533 }
3534
3535 #[test]
3536 fn rotates_encrypted_database_without_overwriting_source() {
3537 let root = tempfile::tempdir().unwrap();
3538 let source_path = root.path().join("source-encrypted");
3539 let destination_path = root.path().join("rotated-encrypted");
3540 let source_options = PersistentOpenOptions {
3541 encryption: Some(EncryptionConfig::from_key(&[0x11_u8; 32]).unwrap()),
3542 ..PersistentOpenOptions::default()
3543 };
3544 let destination_options = PersistentOpenOptions {
3545 encryption: Some(EncryptionConfig::from_key(&[0x22_u8; 32]).unwrap()),
3546 ..PersistentOpenOptions::default()
3547 };
3548 {
3549 let mut source =
3550 PersistentEngine::open_with_options(&source_path, source_options.clone()).unwrap();
3551 source
3552 .put_object(MemoryObject::new("objects", "object", r#"{"value":true}"#))
3553 .unwrap();
3554 source
3555 .append_event(MemoryEvent::new("events", "created", "event body"))
3556 .unwrap();
3557 source
3558 .push_job(QueueJob::new("queue", "job", "payload", 2))
3559 .unwrap();
3560 }
3561 PersistentEngine::reencrypt_to(
3562 &source_path,
3563 &destination_path,
3564 source_options,
3565 destination_options.clone(),
3566 )
3567 .unwrap();
3568
3569 let destination =
3570 PersistentEngine::open_with_options(&destination_path, destination_options).unwrap();
3571 assert!(
3572 destination
3573 .get_object("objects", "object")
3574 .unwrap()
3575 .is_some()
3576 );
3577 assert_eq!(
3578 destination
3579 .list_events(Some("events"), ListEventsOptions::default())
3580 .unwrap()
3581 .len(),
3582 1
3583 );
3584 assert_eq!(destination.list_jobs("queue").unwrap().len(), 1);
3585 let old_key = PersistentEngine::open_with_options(
3586 &destination_path,
3587 PersistentOpenOptions {
3588 encryption: Some(EncryptionConfig::from_key(&[0x11_u8; 32]).unwrap()),
3589 ..PersistentOpenOptions::default()
3590 },
3591 );
3592 assert!(matches!(
3593 old_key,
3594 Err(ThingdError::EncryptionAuthentication(_))
3595 ));
3596 assert!(
3597 PersistentEngine::reencrypt_to(
3598 &source_path,
3599 root.path().join("plaintext-refused"),
3600 PersistentOpenOptions {
3601 encryption: Some(EncryptionConfig::from_key(&[0x11_u8; 32]).unwrap()),
3602 ..PersistentOpenOptions::default()
3603 },
3604 PersistentOpenOptions::default(),
3605 )
3606 .is_err()
3607 );
3608 }
3609
3610 #[test]
3613 fn persistent_stores_and_reads_objects() {
3614 let (mut engine, _dir) = setup();
3615 let object = engine
3616 .put_object(MemoryObject::new(
3617 "decisions",
3618 "rust-core",
3619 r#"{"text":"Use Rust"}"#,
3620 ))
3621 .unwrap();
3622 let stored = engine
3623 .get_object("decisions", "rust-core")
3624 .unwrap()
3625 .unwrap();
3626 assert_eq!(object.version, 1);
3627 assert_eq!(stored.key.collection, "decisions");
3628 assert_eq!(stored.key.id, "rust-core");
3629 }
3630
3631 #[test]
3632 fn persistent_object_created_at_preserved_on_update() {
3633 let (mut engine, _dir) = setup();
3634 let first = engine
3635 .put_object(MemoryObject::new("col", "id", r#"{"v":1}"#))
3636 .unwrap();
3637 assert!(!first.created_at.is_empty());
3638 let second = engine
3639 .put_object(MemoryObject::new("col", "id", r#"{"v":2}"#))
3640 .unwrap();
3641 assert_eq!(second.created_at, first.created_at);
3642 assert!(second.updated_at >= first.created_at);
3643 }
3644
3645 #[test]
3646 fn persistent_object_version_increments_on_update() {
3647 let (mut engine, _dir) = setup();
3648 let v1 = engine
3649 .put_object(MemoryObject::new("col", "x", "{}"))
3650 .unwrap();
3651 assert_eq!(v1.version, 1);
3652 let v2 = engine
3653 .put_object(MemoryObject::new("col", "x", r#"{"v":2}"#))
3654 .unwrap();
3655 assert_eq!(v2.version, 2);
3656 }
3657
3658 #[test]
3659 fn persistent_lists_objects_with_filter() {
3660 let (mut engine, _dir) = setup();
3661 engine
3662 .put_object(MemoryObject::new("w", "a", r#"{"color":"red","size":1}"#))
3663 .unwrap();
3664 engine
3665 .put_object(MemoryObject::new("w", "b", r#"{"color":"blue","size":2}"#))
3666 .unwrap();
3667 engine
3668 .put_object(MemoryObject::new("w", "c", r#"{"color":"red","size":3}"#))
3669 .unwrap();
3670 let opts = ListObjectsOptions {
3671 filter: vec![("color".into(), serde_json::json!("red"))],
3672 ..Default::default()
3673 };
3674 let results = engine
3675 .list_objects(Some(&["w".to_string()]), &opts)
3676 .unwrap();
3677 assert_eq!(results.len(), 2);
3678 assert!(results.iter().all(|o| o.body.contains("\"red\"")));
3679 }
3680
3681 #[test]
3682 fn persistent_list_objects_pagination() {
3683 let (mut engine, _dir) = setup();
3684 for i in 0..5u32 {
3685 engine
3686 .put_object(MemoryObject::new("col", format!("id-{i}"), "{}"))
3687 .unwrap();
3688 }
3689 let limit_opts = ListObjectsOptions {
3690 limit: Some(3),
3691 ..Default::default()
3692 };
3693 assert_eq!(
3694 engine
3695 .list_objects(Some(&["col".to_string()]), &limit_opts)
3696 .unwrap()
3697 .len(),
3698 3
3699 );
3700 let offset_opts = ListObjectsOptions {
3701 offset: Some(3),
3702 ..Default::default()
3703 };
3704 assert_eq!(
3705 engine
3706 .list_objects(Some(&["col".to_string()]), &offset_opts)
3707 .unwrap()
3708 .len(),
3709 2
3710 );
3711 }
3712
3713 #[test]
3714 fn persistent_list_objects_sort_by_created_at_desc() {
3715 let (mut engine, _dir) = setup();
3716 engine
3717 .put_object(MemoryObject::new("w", "a", r#"{"x":1}"#))
3718 .unwrap();
3719 engine
3720 .put_object(MemoryObject::new("w", "b", r#"{"x":2}"#))
3721 .unwrap();
3722 engine
3723 .put_object(MemoryObject::new("w", "c", r#"{"x":3}"#))
3724 .unwrap();
3725 let opts = ListObjectsOptions {
3726 sort_by: Some(crate::SortBy::desc("created_at")),
3727 ..Default::default()
3728 };
3729 let results = engine
3730 .list_objects(Some(&["w".to_string()]), &opts)
3731 .unwrap();
3732 assert_eq!(results.len(), 3);
3733 }
3734
3735 #[test]
3736 fn persistent_list_objects_sort_by_id_asc() {
3737 let (mut engine, _dir) = setup();
3738 engine
3739 .put_object(MemoryObject::new("w", "c", r#"{"x":3}"#))
3740 .unwrap();
3741 engine
3742 .put_object(MemoryObject::new("w", "a", r#"{"x":1}"#))
3743 .unwrap();
3744 engine
3745 .put_object(MemoryObject::new("w", "b", r#"{"x":2}"#))
3746 .unwrap();
3747 let opts = ListObjectsOptions {
3748 sort_by: Some(crate::SortBy::asc("id")),
3749 ..Default::default()
3750 };
3751 let results = engine
3752 .list_objects(Some(&["w".to_string()]), &opts)
3753 .unwrap();
3754 assert_eq!(results.len(), 3);
3755 assert_eq!(results[0].key.id, "a");
3756 assert_eq!(results[1].key.id, "b");
3757 assert_eq!(results[2].key.id, "c");
3758 }
3759
3760 #[test]
3761 fn persistent_cas_succeeds_on_matching_version() {
3762 let (mut engine, _dir) = setup();
3763 let stored = engine
3764 .put_object(MemoryObject::new("col", "id", r#"{"v":1}"#))
3765 .unwrap();
3766 assert_eq!(stored.version, 1);
3767 let opts = crate::PutObjectOptions {
3768 expected_version: Some(1),
3769 ..Default::default()
3770 };
3771 let updated = engine
3772 .put_object_with_options(MemoryObject::new("col", "id", r#"{"v":2}"#), opts)
3773 .unwrap();
3774 assert_eq!(updated.version, 2);
3775 }
3776
3777 #[test]
3778 fn persistent_cas_fails_on_version_mismatch() {
3779 let (mut engine, _dir) = setup();
3780 engine
3781 .put_object(MemoryObject::new("col", "id", r#"{"v":1}"#))
3782 .unwrap();
3783 let opts = crate::PutObjectOptions {
3784 expected_version: Some(42),
3785 ..Default::default()
3786 };
3787 let err = engine
3788 .put_object_with_options(MemoryObject::new("col", "id", r#"{"v":2}"#), opts)
3789 .unwrap_err();
3790 assert!(matches!(err, crate::ThingdError::Conflict(_)));
3791 }
3792
3793 #[test]
3794 fn persistent_cas_fails_on_nonexistent_object() {
3795 let (mut engine, _dir) = setup();
3796 let opts = crate::PutObjectOptions {
3797 expected_version: Some(1),
3798 ..Default::default()
3799 };
3800 let err = engine
3801 .put_object_with_options(MemoryObject::new("col", "id", r#"{"v":1}"#), opts)
3802 .unwrap_err();
3803 assert!(matches!(err, crate::ThingdError::Conflict(_)));
3804 }
3805
3806 #[test]
3807 fn persistent_delete_objects_batch() {
3808 let (mut engine, _dir) = setup();
3809 engine
3810 .put_object(MemoryObject::new("w", "a", "{}"))
3811 .unwrap();
3812 engine
3813 .put_object(MemoryObject::new("w", "b", "{}"))
3814 .unwrap();
3815 engine
3816 .put_object(MemoryObject::new("w", "c", "{}"))
3817 .unwrap();
3818 let keys = vec![
3819 ("w".to_string(), "a".to_string()),
3820 ("w".to_string(), "b".to_string()),
3821 ];
3822 let deleted = engine.delete_objects_batch(&keys).unwrap();
3823 assert_eq!(deleted, 2);
3824 assert_eq!(engine.count_objects().unwrap(), 1);
3825 assert!(engine.get_object("w", "a").unwrap().is_none());
3826 assert!(engine.get_object("w", "c").unwrap().is_some());
3827 }
3828
3829 #[test]
3832 fn persistent_appends_events_with_sequence_numbers() {
3833 let (mut engine, _dir) = setup();
3834 let event = engine
3835 .append_event(MemoryEvent::new(
3836 "project:thingd",
3837 "decision.made",
3838 "MCP-native object storage",
3839 ))
3840 .unwrap();
3841 assert_eq!(event.sequence, 1);
3842 assert_eq!(
3843 engine
3844 .list_events(Some("project:thingd"), ListEventsOptions::default())
3845 .unwrap()
3846 .len(),
3847 1
3848 );
3849 }
3850
3851 #[test]
3852 fn persistent_event_idempotency() {
3853 let (mut engine, _dir) = setup();
3854 let mut event = MemoryEvent::new("stream", "test", r#"{"key":"val"}"#);
3855 event.idempotency_key = "idem-1".to_string();
3856 let first = engine.append_event(event.clone()).unwrap();
3857 assert_eq!(first.sequence, 1);
3858 let second = engine.append_event(event).unwrap();
3859 assert_eq!(second.sequence, first.sequence);
3860 assert_eq!(second.body, first.body);
3861 }
3862
3863 #[test]
3864 fn persistent_deletes_last_event_from_stream() {
3865 let (mut engine, _dir) = setup();
3866 engine
3867 .append_event(MemoryEvent::new("match:1", "turn.recorded", "{}"))
3868 .unwrap();
3869 engine
3870 .append_event(MemoryEvent::new("match:1", "turn.recorded", "{}"))
3871 .unwrap();
3872 engine
3873 .append_event(MemoryEvent::new("match:2", "turn.recorded", "{}"))
3874 .unwrap();
3875 let deleted = engine.delete_last_event("match:1").unwrap().unwrap();
3876 assert_eq!(deleted.sequence, 2);
3877 let remaining = engine
3878 .list_events(Some("match:1"), ListEventsOptions::default())
3879 .unwrap();
3880 assert_eq!(remaining.len(), 1);
3881 assert_eq!(remaining[0].sequence, 1);
3882 let match2 = engine
3883 .list_events(Some("match:2"), ListEventsOptions::default())
3884 .unwrap();
3885 assert_eq!(match2.len(), 1);
3886 }
3887
3888 #[test]
3889 fn persistent_deletes_stream_and_returns_count() {
3890 let (mut engine, _dir) = setup();
3891 engine
3892 .append_event(MemoryEvent::new("match:1", "turn.recorded", "{}"))
3893 .unwrap();
3894 engine
3895 .append_event(MemoryEvent::new("match:1", "turn.recorded", "{}"))
3896 .unwrap();
3897 engine
3898 .append_event(MemoryEvent::new("match:2", "turn.recorded", "{}"))
3899 .unwrap();
3900 assert_eq!(engine.delete_stream("match:1").unwrap(), 2);
3901 assert_eq!(
3902 engine
3903 .list_events(Some("match:1"), ListEventsOptions::default())
3904 .unwrap()
3905 .len(),
3906 0
3907 );
3908 assert_eq!(
3909 engine
3910 .list_events(Some("match:2"), ListEventsOptions::default())
3911 .unwrap()
3912 .len(),
3913 1
3914 );
3915 }
3916
3917 #[test]
3918 fn persistent_lists_streams() {
3919 let (mut engine, _dir) = setup();
3920 assert!(engine.list_streams().unwrap().is_empty());
3921 engine
3922 .append_event(MemoryEvent::new("s1", "t", "e1"))
3923 .unwrap();
3924 engine
3925 .append_event(MemoryEvent::new("s2", "t", "e2"))
3926 .unwrap();
3927 let mut streams = engine.list_streams().unwrap();
3928 streams.sort();
3929 assert_eq!(streams, vec!["s1", "s2"]);
3930 }
3931
3932 #[test]
3935 fn persistent_claims_and_acks_queue_jobs() {
3936 let (mut engine, _dir) = setup();
3937 engine
3938 .push_job(QueueJob::new("embed", "job-1", "doc-1", 3))
3939 .unwrap();
3940 let claimed = engine.claim_job("embed").unwrap().unwrap();
3941 let acked = engine.ack_job("embed", "job-1").unwrap().unwrap();
3942 assert_eq!(claimed.status, QueueJobStatus::Leased);
3943 assert_eq!(claimed.attempts, 1);
3944 assert_eq!(acked.status, QueueJobStatus::Completed);
3945 }
3946
3947 #[test]
3948 fn persistent_nacks_to_dead_letter_after_max_attempts() {
3949 let (mut engine, _dir) = setup();
3950 engine
3951 .push_job(QueueJob::new("embed", "job-1", "doc-1", 1))
3952 .unwrap();
3953 engine.claim_job("embed").unwrap().unwrap();
3954 let nacked = engine.nack_job("embed", "job-1").unwrap().unwrap();
3955 assert_eq!(nacked.status, QueueJobStatus::Dead);
3956 assert_eq!(engine.list_dead_jobs("embed").unwrap().len(), 1);
3957 }
3958
3959 #[test]
3960 fn persistent_does_not_claim_delayed_jobs() {
3961 let (mut engine, _dir) = setup();
3962 engine
3963 .push_job(QueueJob::new("embed", "job-1", "doc-1", 3).delay_by_ms(60_000))
3964 .unwrap();
3965 assert!(engine.claim_job("embed").unwrap().is_none());
3966 }
3967
3968 #[test]
3969 fn persistent_nacks_with_retry_delay() {
3970 let (mut engine, _dir) = setup();
3971 engine
3972 .push_job(QueueJob::new("embed", "job-1", "doc-1", 3))
3973 .unwrap();
3974 engine.claim_job("embed").unwrap().unwrap();
3975 let retried = engine
3976 .nack_job_with_options("embed", "job-1", QueueNackOptions::new(60_000))
3977 .unwrap()
3978 .unwrap();
3979 assert_eq!(retried.status, QueueJobStatus::Ready);
3980 assert!(engine.claim_job("embed").unwrap().is_none());
3981 }
3982
3983 #[test]
3984 fn persistent_queue_counts() {
3985 let (mut engine, _dir) = setup();
3986 assert_eq!(engine.count_active_jobs().unwrap(), 0);
3987 assert_eq!(engine.count_dead_jobs().unwrap(), 0);
3988 engine
3989 .push_job(QueueJob::new("work", "j1", "p1", 3))
3990 .unwrap();
3991 engine
3992 .push_job(QueueJob::new("work", "j2", "p2", 3))
3993 .unwrap();
3994 engine
3995 .push_job(QueueJob::new("other", "j3", "p3", 1))
3996 .unwrap();
3997 assert_eq!(engine.count_active_jobs().unwrap(), 3);
3998 engine.claim_job("other").unwrap();
3999 engine.nack_job("other", "j3").unwrap();
4000 assert_eq!(engine.count_dead_jobs().unwrap(), 1);
4001 assert_eq!(engine.count_active_jobs().unwrap(), 2);
4002 }
4003
4004 #[test]
4005 fn persistent_lists_queues() {
4006 let (mut engine, _dir) = setup();
4007 engine
4008 .push_job(QueueJob::new("work", "j1", "p1", 3))
4009 .unwrap();
4010 engine
4011 .push_job(QueueJob::new("jobs", "j2", "p2", 3))
4012 .unwrap();
4013 let mut queues = engine.list_queues().unwrap();
4014 queues.sort();
4015 assert_eq!(queues, vec!["jobs", "work"]);
4016 }
4017
4018 #[test]
4019 fn persistent_claim_reclaims_expired_lease() {
4020 let (mut engine, _dir) = setup();
4021 engine
4022 .push_job(QueueJob::new("embed", "job-1", "doc-1", 3))
4023 .unwrap();
4024 let first = engine
4025 .claim_job_with_options("embed", QueueClaimOptions::new(0))
4026 .unwrap()
4027 .unwrap();
4028 let second = engine.claim_job("embed").unwrap().unwrap();
4029 assert_eq!(first.status, QueueJobStatus::Leased);
4030 assert_eq!(second.status, QueueJobStatus::Leased);
4031 assert_eq!(second.attempts, 2);
4032 }
4033
4034 #[test]
4035 fn persistent_priority_ordering() {
4036 let (mut engine, _dir) = setup();
4037 engine
4038 .push_job(QueueJob::new("q", "low", "body", 3).with_priority(0))
4039 .unwrap();
4040 engine
4041 .push_job(QueueJob::new("q", "high", "body", 3).with_priority(10))
4042 .unwrap();
4043 engine
4044 .push_job(QueueJob::new("q", "mid", "body", 3).with_priority(5))
4045 .unwrap();
4046 let first = engine.claim_job("q").unwrap().unwrap();
4047 assert_eq!(first.id, "high", "highest priority claimed first");
4048 let second = engine.claim_job("q").unwrap().unwrap();
4049 assert_eq!(second.id, "mid", "medium priority claimed second");
4050 let third = engine.claim_job("q").unwrap().unwrap();
4051 assert_eq!(third.id, "low", "lowest priority claimed last");
4052 }
4053
4054 #[test]
4057 fn persistent_create_get_delete_link() {
4058 let (mut engine, _dir) = setup();
4059 engine
4060 .put_object(MemoryObject::new("n", "a", "{}"))
4061 .unwrap();
4062 engine
4063 .put_object(MemoryObject::new("n", "b", "{}"))
4064 .unwrap();
4065 let link = engine
4066 .create_link(Link::new("n/a", "connects", "n/b"))
4067 .unwrap();
4068 assert!(!link.id.is_empty());
4069 let fetched = engine.get_link(&link.id).unwrap().unwrap();
4070 assert_eq!(fetched.id, link.id);
4071 assert!(engine.delete_link(&link.id).unwrap());
4072 assert!(engine.get_link(&link.id).unwrap().is_none());
4073 }
4074
4075 #[test]
4076 fn persistent_neighbor_query() {
4077 let (mut engine, _dir) = setup();
4078 engine
4079 .put_object(MemoryObject::new("n", "a", "{}"))
4080 .unwrap();
4081 engine
4082 .put_object(MemoryObject::new("n", "b", "{}"))
4083 .unwrap();
4084 engine
4085 .put_object(MemoryObject::new("n", "c", "{}"))
4086 .unwrap();
4087 engine
4088 .create_link(Link::new("n/a", "knows", "n/b"))
4089 .unwrap();
4090 engine
4091 .create_link(Link::new("n/a", "knows", "n/c"))
4092 .unwrap();
4093 let outgoing = engine
4094 .get_neighbors("n/a", LinkDirection::Outgoing, LinkQueryOptions::default())
4095 .unwrap();
4096 assert_eq!(outgoing.len(), 2);
4097 let incoming = engine
4098 .get_neighbors("n/b", LinkDirection::Incoming, LinkQueryOptions::default())
4099 .unwrap();
4100 assert_eq!(incoming.len(), 1);
4101 }
4102
4103 #[test]
4104 fn persistent_link_count() {
4105 let (mut engine, _dir) = setup();
4106 engine
4107 .put_object(MemoryObject::new("n", "a", "{}"))
4108 .unwrap();
4109 engine
4110 .put_object(MemoryObject::new("n", "b", "{}"))
4111 .unwrap();
4112 assert_eq!(engine.count_links().unwrap(), 0);
4113 engine
4114 .create_link(Link::new("n/a", "knows", "n/b"))
4115 .unwrap();
4116 assert_eq!(engine.count_links().unwrap(), 1);
4117 }
4118
4119 #[test]
4122 fn persistent_search_objects_and_events() {
4123 let (mut engine, _dir) = setup();
4124 engine
4125 .put_object(MemoryObject::new("docs", "a", r#"{"text":"hello world"}"#))
4126 .unwrap();
4127 engine
4128 .put_object(MemoryObject::new(
4129 "docs",
4130 "b",
4131 r#"{"text":"goodbye world"}"#,
4132 ))
4133 .unwrap();
4134 engine
4135 .append_event(MemoryEvent::new("audit", "test", "hello event"))
4136 .unwrap();
4137 let results = engine.search("hello", SearchOptions::default()).unwrap();
4138 assert_eq!(results.len(), 2);
4139 let kinds: Vec<&str> = results.iter().map(|h| h.kind.as_str()).collect();
4140 assert!(kinds.contains(&"object"));
4141 assert!(kinds.contains(&"event"));
4142 }
4143
4144 #[test]
4145 fn persistent_search_with_collections() {
4146 let (mut engine, _dir) = setup();
4147 engine
4148 .put_object(MemoryObject::new("docs", "a", r#"{"text":"hello world"}"#))
4149 .unwrap();
4150 engine
4151 .put_object(MemoryObject::new("notes", "b", r#"{"text":"hello there"}"#))
4152 .unwrap();
4153 let opts = SearchOptions {
4154 collections: Some(vec!["docs".into()]),
4155 ..Default::default()
4156 };
4157 let results = engine.search("hello", opts).unwrap();
4158 assert_eq!(results.len(), 1);
4159 assert_eq!(results[0].collection, "docs");
4160 }
4161
4162 #[cfg(feature = "search")]
4165 #[test]
4166 fn persistent_search_indexes_on_put() {
4167 let (mut engine, _dir) = setup();
4168 engine
4169 .put_object(MemoryObject::new(
4170 "docs",
4171 "a",
4172 r#"{"text":"unique_search_term_xyz"}"#,
4173 ))
4174 .unwrap();
4175 let results = engine
4176 .search("unique_search_term_xyz", SearchOptions::default())
4177 .unwrap();
4178 assert_eq!(
4179 results.len(),
4180 1,
4181 "search must find indexed content immediately after put"
4182 );
4183 assert_eq!(results[0].id, "a");
4184 }
4185
4186 #[cfg(feature = "search")]
4187 #[test]
4188 fn legacy_search_schema_is_rebuilt_without_panicking() {
4189 let dir = tempfile::tempdir().unwrap();
4190 {
4191 let mut engine = PersistentEngine::open(dir.path()).unwrap();
4192 engine
4193 .put_object(MemoryObject::new(
4194 "legacy",
4195 "object",
4196 r#"{"text":"legacy_search_term"}"#,
4197 ))
4198 .unwrap();
4199 }
4200
4201 let search_dir = dir.path().join("search");
4202 std::fs::remove_dir_all(&search_dir).unwrap();
4203 std::fs::create_dir_all(&search_dir).unwrap();
4204 let mut schema_builder = tantivy::schema::Schema::builder();
4205 schema_builder.add_text_field("body", tantivy::schema::TEXT | tantivy::schema::STORED);
4206 let legacy_schema = schema_builder.build();
4207 tantivy::Index::create_in_dir(&search_dir, legacy_schema).unwrap();
4208
4209 let engine = PersistentEngine::open(dir.path()).unwrap();
4210 let results = engine
4211 .search("legacy_search_term", SearchOptions::default())
4212 .unwrap();
4213 assert_eq!(results.len(), 1);
4214 assert_eq!(results[0].id, "object");
4215 }
4216
4217 #[cfg(feature = "search")]
4218 #[test]
4219 fn persistent_search_removes_on_delete() {
4220 let (mut engine, _dir) = setup();
4221 engine
4222 .put_object(MemoryObject::new(
4223 "docs",
4224 "to-delete",
4225 r#"{"text":"deletable_content"}"#,
4226 ))
4227 .unwrap();
4228 assert_eq!(
4230 engine
4231 .search("deletable_content", SearchOptions::default())
4232 .unwrap()
4233 .len(),
4234 1
4235 );
4236 engine.delete_object("docs", "to-delete").unwrap();
4238 let after = engine
4239 .search("deletable_content", SearchOptions::default())
4240 .unwrap();
4241 assert_eq!(
4242 after.len(),
4243 0,
4244 "deleted object must not appear in search results"
4245 );
4246 }
4247
4248 #[cfg(feature = "search")]
4249 #[test]
4250 fn persistent_search_deleted_batch_removes_from_index() {
4251 let (mut engine, _dir) = setup();
4252 engine
4253 .put_object(MemoryObject::new(
4254 "docs",
4255 "a",
4256 r#"{"text":"batch_deleted_a"}"#,
4257 ))
4258 .unwrap();
4259 engine
4260 .put_object(MemoryObject::new(
4261 "docs",
4262 "b",
4263 r#"{"text":"batch_deleted_b"}"#,
4264 ))
4265 .unwrap();
4266 assert_eq!(
4267 engine
4268 .search("batch_deleted", SearchOptions::default())
4269 .unwrap()
4270 .len(),
4271 2
4272 );
4273 let keys = vec![
4274 ("docs".to_string(), "a".to_string()),
4275 ("docs".to_string(), "b".to_string()),
4276 ];
4277 engine.delete_objects_batch(&keys).unwrap();
4278 let after = engine
4279 .search("batch_deleted", SearchOptions::default())
4280 .unwrap();
4281 assert_eq!(
4282 after.len(),
4283 0,
4284 "batch-deleted objects must be removed from search index"
4285 );
4286 }
4287
4288 #[test]
4291 fn persistent_aggregate_count_sum_avg() {
4292 let (mut engine, _dir) = setup();
4293 engine
4294 .put_object(MemoryObject::new("stats", "a", r#"{"val":10}"#))
4295 .unwrap();
4296 engine
4297 .put_object(MemoryObject::new("stats", "b", r#"{"val":20}"#))
4298 .unwrap();
4299 engine
4300 .put_object(MemoryObject::new("stats", "c", r#"{"val":30}"#))
4301 .unwrap();
4302 let count = engine
4303 .aggregate(
4304 "stats",
4305 &AggregateOptions {
4306 function: AggregateFunction::Count,
4307 ..Default::default()
4308 },
4309 )
4310 .unwrap();
4311 assert_eq!(count.total, 3.0);
4312 let sum = engine
4313 .aggregate(
4314 "stats",
4315 &AggregateOptions {
4316 function: AggregateFunction::Sum,
4317 field: Some("val".into()),
4318 ..Default::default()
4319 },
4320 )
4321 .unwrap();
4322 assert_eq!(sum.total, 60.0);
4323 let avg = engine
4324 .aggregate(
4325 "stats",
4326 &AggregateOptions {
4327 function: AggregateFunction::Avg,
4328 field: Some("val".into()),
4329 ..Default::default()
4330 },
4331 )
4332 .unwrap();
4333 assert_eq!(avg.total, 20.0);
4334 }
4335
4336 #[test]
4337 fn persistent_aggregate_group_by() {
4338 let (mut engine, _dir) = setup();
4339 engine
4340 .put_object(MemoryObject::new(
4341 "sales",
4342 "a",
4343 r#"{"region":"EU","val":100}"#,
4344 ))
4345 .unwrap();
4346 engine
4347 .put_object(MemoryObject::new(
4348 "sales",
4349 "b",
4350 r#"{"region":"US","val":200}"#,
4351 ))
4352 .unwrap();
4353 engine
4354 .put_object(MemoryObject::new(
4355 "sales",
4356 "c",
4357 r#"{"region":"EU","val":50}"#,
4358 ))
4359 .unwrap();
4360 let result = engine
4361 .aggregate(
4362 "sales",
4363 &AggregateOptions {
4364 function: AggregateFunction::Sum,
4365 field: Some("val".into()),
4366 group_by: Some("region".into()),
4367 ..Default::default()
4368 },
4369 )
4370 .unwrap();
4371 assert_eq!(result.total, 350.0);
4372 assert_eq!(result.groups.len(), 2);
4373 for group in &result.groups {
4374 match group.key.as_str() {
4375 "EU" => assert_eq!(group.value, 150.0),
4376 "US" => assert_eq!(group.value, 200.0),
4377 _ => panic!("unexpected group key"),
4378 }
4379 }
4380 }
4381
4382 #[test]
4383 fn persistent_timeseries_bucketing() {
4384 let (mut engine, _dir) = setup();
4385 engine
4386 .put_object(MemoryObject::new("events", "a", r#"{"val":1}"#))
4387 .unwrap();
4388 engine
4389 .put_object(MemoryObject::new("events", "b", r#"{"val":2}"#))
4390 .unwrap();
4391 let result = engine
4392 .timeseries(
4393 "events",
4394 &TimeSeriesOptions {
4395 function: AggregateFunction::Count,
4396 bucket: TimeBucket::Day,
4397 ..Default::default()
4398 },
4399 )
4400 .unwrap();
4401 assert_eq!(result.buckets.len(), 1);
4402 assert_eq!(result.buckets[0].value, 2.0);
4403 }
4404
4405 #[test]
4408 fn persistent_ready_jobs_indexes_on_push() {
4409 let (mut engine, _dir) = setup();
4410 engine
4411 .push_job(QueueJob::new("q", "j1", "body", 3))
4412 .unwrap();
4413 let prefix = b"q\0";
4414 let count = engine.ready_jobs.prefix(prefix).count();
4415 assert_eq!(count, 1, "ready_jobs must have one entry after push");
4416 }
4417
4418 #[test]
4419 fn persistent_ready_jobs_removed_on_claim() {
4420 let (mut engine, _dir) = setup();
4421 engine
4422 .push_job(QueueJob::new("q", "j1", "body", 3))
4423 .unwrap();
4424 engine.claim_job("q").unwrap();
4425 let prefix = b"q\0";
4426 let count = engine.ready_jobs.prefix(prefix).count();
4427 assert_eq!(
4428 count, 0,
4429 "ready_jobs must be empty after claiming the only job"
4430 );
4431 }
4432
4433 #[test]
4434 fn persistent_ready_jobs_priority_order() {
4435 let (mut engine, _dir) = setup();
4436 engine
4437 .push_job(QueueJob::new("q", "low", "body", 3).with_priority(0))
4438 .unwrap();
4439 engine
4440 .push_job(QueueJob::new("q", "high", "body", 3).with_priority(10))
4441 .unwrap();
4442 let prefix = b"q\0";
4444 let keys: Vec<Vec<u8>> = engine
4445 .ready_jobs
4446 .prefix(prefix)
4447 .map(|kv| {
4448 let (k, _) = guard_data(kv).unwrap();
4449 k
4450 })
4451 .collect();
4452 assert_eq!(keys.len(), 2);
4453 let first_key_str = String::from_utf8_lossy(&keys[0]);
4455 assert!(
4456 first_key_str.contains("high"),
4457 "first ready entry should be high-priority job; got {first_key_str}"
4458 );
4459 }
4460
4461 #[test]
4462 fn persistent_ready_jobs_fifo_order() {
4463 let (mut engine, _dir) = setup();
4464 engine
4465 .push_job(QueueJob::new("q", "first", "body", 3))
4466 .unwrap();
4467 std::thread::sleep(std::time::Duration::from_millis(5));
4469 engine
4470 .push_job(QueueJob::new("q", "second", "body", 3))
4471 .unwrap();
4472 let prefix = b"q\0";
4473 let keys: Vec<Vec<u8>> = engine
4474 .ready_jobs
4475 .prefix(prefix)
4476 .map(|kv| {
4477 let (k, _) = guard_data(kv).unwrap();
4478 k
4479 })
4480 .collect();
4481 assert_eq!(keys.len(), 2);
4482 let first_key_str = String::from_utf8_lossy(&keys[0]);
4483 assert!(
4484 first_key_str.contains("first"),
4485 "first ready entry should be FIFO; got {first_key_str}"
4486 );
4487 }
4488
4489 #[test]
4490 fn persistent_ready_jobs_reindex_on_nack() {
4491 let (mut engine, _dir) = setup();
4492 engine
4493 .push_job(QueueJob::new("q", "j1", "body", 3))
4494 .unwrap();
4495 engine.claim_job("q").unwrap();
4496 let prefix = b"q\0";
4497 assert_eq!(engine.ready_jobs.prefix(prefix).count(), 0);
4498 engine
4500 .nack_job_with_options("q", "j1", QueueNackOptions::new(0))
4501 .unwrap();
4502 assert_eq!(
4503 engine.ready_jobs.prefix(prefix).count(),
4504 1,
4505 "ready_jobs must have entry after nack with retry"
4506 );
4507 }
4508
4509 #[test]
4510 fn persistent_ready_jobs_reindex_on_lease_expire() {
4511 let (mut engine, _dir) = setup();
4512 engine
4513 .push_job(QueueJob::new("q", "j1", "body", 3))
4514 .unwrap();
4515 engine
4517 .claim_job_with_options("q", QueueClaimOptions::new(0))
4518 .unwrap();
4519 let prefix = b"q\0";
4521 let _count = engine.ready_jobs.prefix(prefix).count();
4522 let claimed = engine.claim_job("q").unwrap();
4524 assert!(
4525 claimed.is_some(),
4526 "job should be claimable after lease expires"
4527 );
4528 let job = claimed.unwrap();
4529 assert_eq!(job.attempts, 2, "second attempt after lease expiry");
4530 assert_eq!(
4531 engine.ready_jobs.prefix(prefix).count(),
4532 0,
4533 "ready_jobs must be empty after re-claiming"
4534 );
4535 }
4536
4537 #[cfg(feature = "vectors")]
4540 #[test]
4541 fn persistent_vector_search_returns_by_cosine_similarity() {
4542 let (mut engine, _dir) = setup();
4543 engine
4544 .put_object(
4545 MemoryObject::new("docs", "a", r#"{"text":"alpha"}"#)
4546 .with_vector(vec![1.0, 0.0, 0.0]),
4547 )
4548 .unwrap();
4549 engine
4550 .put_object(
4551 MemoryObject::new("docs", "b", r#"{"text":"beta"}"#)
4552 .with_vector(vec![0.0, 1.0, 0.0]),
4553 )
4554 .unwrap();
4555
4556 let results = engine
4557 .vector_search("docs", &[0.9, 0.1, 0.0], VectorSearchOptions::default())
4558 .unwrap();
4559 assert_eq!(results.len(), 2);
4560 assert_eq!(results[0].id, "a");
4561 assert!(results[0].score > results[1].score);
4562 }
4563
4564 #[cfg(feature = "vectors")]
4565 #[test]
4566 fn persistent_vector_search_respects_filter() {
4567 let (mut engine, _dir) = setup();
4568 engine
4569 .put_object(
4570 MemoryObject::new("docs", "a", r#"{"tag":"x"}"#).with_vector(vec![1.0, 0.0]),
4571 )
4572 .unwrap();
4573 engine
4574 .put_object(
4575 MemoryObject::new("docs", "b", r#"{"tag":"y"}"#).with_vector(vec![0.0, 1.0]),
4576 )
4577 .unwrap();
4578
4579 let results = engine
4580 .vector_search(
4581 "docs",
4582 &[1.0, 0.0],
4583 VectorSearchOptions {
4584 filter: Some(serde_json::json!({"tag": "x"})),
4585 ..Default::default()
4586 },
4587 )
4588 .unwrap();
4589 assert_eq!(results.len(), 1);
4590 assert_eq!(results[0].id, "a");
4591 }
4592
4593 #[cfg(feature = "vectors")]
4594 #[test]
4595 fn persistent_vector_search_excludes_deleted_objects() {
4596 let (mut engine, _dir) = setup();
4597 engine
4598 .put_object(MemoryObject::new("docs", "a", "{}").with_vector(vec![1.0, 0.0]))
4599 .unwrap();
4600 engine.delete_object("docs", "a").unwrap();
4601 let results = engine
4602 .vector_search("docs", &[1.0, 0.0], VectorSearchOptions::default())
4603 .unwrap();
4604 assert!(results.is_empty());
4605 }
4606
4607 #[cfg(feature = "vectors")]
4608 #[test]
4609 fn persistent_vector_search_respects_top_k() {
4610 let (mut engine, _dir) = setup();
4611 engine
4612 .put_object(
4613 MemoryObject::new("docs", "a", r#"{"text":"alpha"}"#).with_vector(vec![1.0, 0.0]),
4614 )
4615 .unwrap();
4616 engine
4617 .put_object(
4618 MemoryObject::new("docs", "b", r#"{"text":"beta"}"#).with_vector(vec![0.0, 1.0]),
4619 )
4620 .unwrap();
4621
4622 let results = engine
4623 .vector_search(
4624 "docs",
4625 &[1.0, 0.0],
4626 VectorSearchOptions {
4627 top_k: Some(1),
4628 ..Default::default()
4629 },
4630 )
4631 .unwrap();
4632 assert_eq!(results.len(), 1);
4633 assert_eq!(results[0].id, "a");
4634 }
4635
4636 #[cfg(feature = "vectors")]
4637 #[test]
4638 fn persistent_vector_search_empty_collection_returns_empty() {
4639 let (engine, _dir) = setup();
4640 let results = engine
4641 .vector_search("docs", &[1.0, 0.0, 0.0], VectorSearchOptions::default())
4642 .unwrap();
4643 assert!(results.is_empty());
4644 }
4645
4646 #[cfg(feature = "vectors")]
4647 #[test]
4648 fn persistent_vector_search_rejects_dimension_mismatch() {
4649 let (mut engine, _dir) = setup();
4650 engine
4651 .put_object(MemoryObject::new("docs", "a", "{}").with_vector(vec![1.0, 0.0]))
4652 .unwrap();
4653
4654 let error = engine
4655 .vector_search("docs", &[1.0, 0.0, 0.0], VectorSearchOptions::default())
4656 .unwrap_err();
4657 assert!(
4658 matches!(error, ThingdError::InvalidInput(message) if message.contains("dimension"))
4659 );
4660 }
4661
4662 #[cfg(feature = "vectors")]
4663 #[test]
4664 fn persistent_vector_search_rejects_empty_query() {
4665 let (engine, _dir) = setup();
4666 let error = engine
4667 .vector_search("docs", &[], VectorSearchOptions::default())
4668 .unwrap_err();
4669 assert!(matches!(error, ThingdError::InvalidInput(message) if message.contains("empty")));
4670 }
4671
4672 #[cfg(feature = "vectors")]
4673 #[test]
4674 fn persistent_put_object_without_vector_does_not_store_vector() {
4675 let (mut engine, _dir) = setup();
4676 engine
4677 .put_object(MemoryObject::new("docs", "a", "{}"))
4678 .unwrap();
4679 let results = engine
4680 .vector_search("docs", &[1.0, 0.0, 0.0], VectorSearchOptions::default())
4681 .unwrap();
4682 assert!(results.is_empty());
4683 }
4684
4685 #[cfg(feature = "vectors")]
4686 #[test]
4687 fn persistent_vector_search_persists_across_engine_reopen() {
4688 let dir = tempfile::tempdir().unwrap();
4689 {
4690 let mut engine = PersistentEngine::open(dir.path()).unwrap();
4691 engine
4692 .put_object(
4693 MemoryObject::new("docs", "a", r#"{"text":"persist"}"#)
4694 .with_vector(vec![1.0, 0.0, 0.0]),
4695 )
4696 .unwrap();
4697 }
4698 {
4699 let engine = PersistentEngine::open(dir.path()).unwrap();
4700 let results = engine
4701 .vector_search("docs", &[1.0, 0.0, 0.0], VectorSearchOptions::default())
4702 .unwrap();
4703 assert_eq!(results.len(), 1);
4704 assert_eq!(results[0].id, "a");
4705 }
4706 }
4707
4708 #[test]
4711 fn persistent_event_sequence_survives_reopen() {
4712 let dir = tempfile::tempdir().unwrap();
4713 let stream = "test-stream";
4714 {
4715 let mut engine = PersistentEngine::open(dir.path()).unwrap();
4716 let e1 = engine
4717 .append_event(MemoryEvent::new(stream, "t1", "{}"))
4718 .unwrap();
4719 assert_eq!(e1.sequence, 1);
4720 let e2 = engine
4721 .append_event(MemoryEvent::new(stream, "t2", "{}"))
4722 .unwrap();
4723 assert_eq!(e2.sequence, 2);
4724 }
4725 {
4726 let mut engine = PersistentEngine::open(dir.path()).unwrap();
4727 let e3 = engine
4729 .append_event(MemoryEvent::new(stream, "t3", "{}"))
4730 .unwrap();
4731 assert_eq!(
4732 e3.sequence, 3,
4733 "sequence must continue from durable max after reopen"
4734 );
4735 let events = engine
4737 .list_events(Some(stream), ListEventsOptions::default())
4738 .unwrap();
4739 assert_eq!(events.len(), 3);
4740 }
4741 }
4742
4743 #[test]
4744 fn persistent_event_idempotency_survives_reopen() {
4745 let dir = tempfile::tempdir().unwrap();
4746 let stream = "test-stream";
4747 {
4748 let mut engine = PersistentEngine::open(dir.path()).unwrap();
4749 let mut e = MemoryEvent::new(stream, "t1", r#"{"x":1}"#);
4750 e.idempotency_key = "key-1".to_string();
4751 let e1 = engine.append_event(e).unwrap();
4752 assert_eq!(e1.sequence, 1);
4753 }
4754 {
4755 let mut engine = PersistentEngine::open(dir.path()).unwrap();
4756 let mut e = MemoryEvent::new(stream, "t1", r#"{"x":1}"#);
4758 e.idempotency_key = "key-1".to_string();
4759 let e2 = engine.append_event(e).unwrap();
4760 assert_eq!(
4761 e2.sequence, 1,
4762 "idempotency must be preserved across reopen"
4763 );
4764 let e3 = engine
4766 .append_event(MemoryEvent::new(stream, "t2", "{}"))
4767 .unwrap();
4768 assert_eq!(e3.sequence, 2);
4769 }
4770 }
4771
4772 #[cfg(feature = "vectors")]
4773 #[test]
4774 fn persistent_vector_survives_reopen() {
4775 let dir = tempfile::tempdir().unwrap();
4776 {
4777 let mut engine = PersistentEngine::open(dir.path()).unwrap();
4778 engine
4779 .put_object(
4780 MemoryObject::new("docs", "a", r#"{"text":"persist"}"#)
4781 .with_vector(vec![1.0, 0.0, 0.0]),
4782 )
4783 .unwrap();
4784 }
4785 {
4786 let engine = PersistentEngine::open(dir.path()).unwrap();
4787 let results = engine
4788 .vector_search("docs", &[1.0, 0.0, 0.0], VectorSearchOptions::default())
4789 .unwrap();
4790 assert_eq!(results.len(), 1);
4791 assert_eq!(results[0].id, "a");
4792 }
4793 }
4794
4795 #[cfg(feature = "vectors")]
4796 #[test]
4797 fn persistent_vector_removed_on_update_without_vector_reopen() {
4798 let dir = tempfile::tempdir().unwrap();
4799 {
4800 let mut engine = PersistentEngine::open(dir.path()).unwrap();
4801 engine
4802 .put_object(
4803 MemoryObject::new("docs", "a", r#"{"v":1}"#).with_vector(vec![1.0, 0.0, 0.0]),
4804 )
4805 .unwrap();
4806 }
4807 {
4808 let mut engine = PersistentEngine::open(dir.path()).unwrap();
4809 engine
4811 .put_object(MemoryObject::new("docs", "a", r#"{"v":2}"#))
4812 .unwrap();
4813 }
4814 {
4815 let engine = PersistentEngine::open(dir.path()).unwrap();
4816 let results = engine
4817 .vector_search("docs", &[1.0, 0.0, 0.0], VectorSearchOptions::default())
4818 .unwrap();
4819 assert_eq!(results.len(), 0, "vector must survive reopen and removal");
4820 }
4821 }
4822
4823 fn setup_persistent() -> (PersistentEngine, tempfile::TempDir) {
4826 let dir = tempfile::tempdir().unwrap();
4827 let engine = PersistentEngine::open(dir.path()).unwrap();
4828 (engine, dir)
4829 }
4830
4831 #[test]
4832 fn contract_object_lifecycle() {
4833 let (mut engine, _dir) = setup_persistent();
4834 crate::contract_tests::test_contract_object_lifecycle(&mut engine);
4835 }
4836
4837 #[test]
4838 fn contract_vector_lifecycle() {
4839 let (mut engine, _dir) = setup_persistent();
4840 crate::contract_tests::test_contract_vector_lifecycle(&mut engine);
4841 }
4842
4843 #[test]
4844 fn contract_schema_store() {
4845 let (mut engine, _dir) = setup_persistent();
4846 crate::contract_tests::test_contract_schema_store(&mut engine);
4847 crate::contract_tests::test_contract_indexes(&mut engine);
4848 }
4849
4850 #[test]
4851 fn contract_event_idempotency() {
4852 let (mut engine, _dir) = setup_persistent();
4853 crate::contract_tests::test_contract_event_idempotency(&mut engine);
4854 }
4855
4856 #[test]
4857 fn contract_queue_lifecycle() {
4858 let (mut engine, _dir) = setup_persistent();
4859 crate::contract_tests::test_contract_queue_lifecycle(&mut engine);
4860 }
4861
4862 #[test]
4863 fn contract_delayed_job() {
4864 let (mut engine, _dir) = setup_persistent();
4865 crate::contract_tests::test_contract_delayed_job(&mut engine);
4866 }
4867
4868 #[test]
4869 fn contract_lease_expiration() {
4870 let (mut engine, _dir) = setup_persistent();
4871 crate::contract_tests::test_contract_lease_expiration(&mut engine);
4872 }
4873
4874 #[test]
4875 fn contract_nack_dead_letter() {
4876 let (mut engine, _dir) = setup_persistent();
4877 crate::contract_tests::test_contract_nack_dead_letter(&mut engine);
4878 }
4879
4880 #[test]
4881 fn contract_search() {
4882 let (mut engine, _dir) = setup_persistent();
4883 crate::contract_tests::test_contract_search(&mut engine);
4884 }
4885}