Skip to main content

thingd/
persistent.rs

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
39/// Persistent storage engine implementing all 6 storage traits.
40///
41/// Data directory layout:
42/// - `objects`: `{collection}\0{id}` → serialized `MemoryObject`
43/// - `events`: `{stream}\0{seq:8BE}` → serialized `MemoryEvent`
44/// - `queue_jobs`: `{queue}\0{id}` → serialized `QueueJob`
45/// - `links_by_id`: `{link_id}` → serialized `Link`
46/// - `links_from`: `{from_ref}\0{type}\0{link_id}` → `()`
47/// - `links_to`: `{to_ref}\0{type}\0{link_id}` → `()`
48pub 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";
83// Tantivy requires at least 15 MB per writer thread; this is lower than the
84// previous 50 MB budget while remaining valid for the current Tantivy release.
85const 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/// Search index behavior for a persistent engine.
120#[derive(Debug, Clone, Copy, Default, Eq, PartialEq)]
121pub enum PersistentSearchMode {
122    /// Open the persistent Tantivy index and rebuild it when necessary.
123    #[default]
124    Persistent,
125    /// Open a compatible Tantivy index but do not create or rebuild one.
126    PersistentNoRebuild,
127    /// Do not open or rebuild Tantivy. Search uses the bounded fallback scan.
128    Disabled,
129}
130
131/// Result of validating a native storage directory without opening Fjall.
132#[derive(Debug, Clone, Serialize, Deserialize)]
133pub struct StorageValidationReport {
134    /// Current Thingd storage format version.
135    pub format_version: u32,
136    /// Whether the directory predates the Thingd manifest and can be upgraded on open.
137    pub legacy_manifest: bool,
138    /// Whether the expected lock file is present.
139    pub lock_present: bool,
140    /// Whether the expected keyspace directory is present.
141    pub keyspaces_present: bool,
142    /// Whether the existing Tantivy schema is compatible, when an index exists.
143    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    // Stores created before the manifest was introduced are accepted after the
259    // structural checks above and receive a manifest during the normal open.
260    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/// Options used when opening a persistent database.
295#[derive(Clone)]
296pub struct PersistentOpenOptions {
297    /// Optional authenticated-encryption configuration.
298    pub encryption: Option<EncryptionConfig>,
299    /// Permit an explicit encrypted-to-plaintext migration when used by
300    /// `PersistentEngine::reencrypt_to`. Ignored during normal open.
301    pub allow_plaintext_output: bool,
302    /// Controls whether the persistent Tantivy index is opened.
303    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    /// Open or create a Persistent database at the given path.
330    /// Creates all required keyspaces (partitions) on first open.
331    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    /// Open a persistent database with explicit options.
337    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        // Legacy stores do not have event metadata yet. Scan them once and persist
397        // the derived counters so later opens do not repeat the full reconstruction.
398        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    /// Validate a native storage directory without opening Fjall or mutating files.
504    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    /// Re-encrypt a database into a new destination without modifying the source.
511    ///
512    /// The destination must not already exist. This operation is intended for
513    /// offline migration and key rotation; it never changes a key implicitly.
514    #[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        // A Tantivy index is derived state. If an older SDK wrote an
706        // incompatible schema (for example without `doc_key`), discard only
707        // that derived directory and rebuild it from Fjall records. Never let
708        // a schema mismatch panic or make the primary database unreadable.
709        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    /// Ready jobs index key: {`queue}\0{priority_rev:8BE}\0{created_at}\0{id`}
827    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
1159// ── ObjectStore ──────────────────────────────────────────────────────────────
1160
1161impl 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        // Atomic batch: object data + vector state
1184        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        // Avoid materializing the whole collection when the caller requests
1371        // the natural key order. Sorting still requires the full candidate set.
1372        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        // Atomic batch: object + vector removal
1507        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
1743// ── EventLog ─────────────────────────────────────────────────────────────────
1744
1745impl 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
1950// ── QueueStore ───────────────────────────────────────────────────────────────
1951
1952impl 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        // Index in ready_jobs for O(1) claiming
1961        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        // Release expired leases and re-index into ready_jobs
1985        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        // Scan ready_jobs index in priority order, skipping delayed and stale entries
2004        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            // Read full job from queue_jobs
2020            let qkey = self.make_queue_key(queue, &job_id);
2021            let Some(job_data) = value_to_vec(self.queue_jobs.get(&qkey)?) else {
2022                // Job record missing — remove stale index entry
2023                let _ = self.ready_jobs.remove(&rkey);
2024                continue;
2025            };
2026            let mut job: QueueJob = self.deserialize(&job_data)?;
2027            // Release expired lease if this job was previously leased
2028            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            // Skip (and remove) stale entries for completed or dead jobs
2037            if job.status != QueueJobStatus::Ready {
2038                let _ = self.ready_jobs.remove(&rkey);
2039                continue;
2040            }
2041
2042            // Job is delayed — skip it but keep the index entry so it can be claimed later
2043            if job.available_at_ms > now {
2044                continue;
2045            }
2046
2047            // Claim this job
2048            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                // Re-index if retrying
2109                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
2182// ── Searcher ─────────────────────────────────────────────────────────────────
2183
2184impl Searcher for PersistentEngine {
2185    fn search(&self, query: &str, options: SearchOptions) -> ThingdResult<Vec<SearchHit>> {
2186        // Try Tantivy search first
2187        #[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        // Fallback: naive substring search (same as MemoryEngine)
2193        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        // Remove existing document with the same doc_key to prevent duplicates
2224        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
2604// ── LinkStore ────────────────────────────────────────────────────────────────
2605
2606impl 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
2728// ── AggregateStore ───────────────────────────────────────────────────────────
2729
2730impl 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
2875// ── VectorStore ──────────────────────────────────────────────────────────────
2876
2877impl 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
2996// ── Helper functions ─────────────────────────────────────────────────────────
2997
2998fn 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    /// Create a test engine with a temp directory that stays alive for the caller.
3183    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    // ── ObjectStore ───────────────────────────────────────────────────────
3611
3612    #[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    // ── EventLog ──────────────────────────────────────────────────────────
3830
3831    #[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    // ── QueueStore ────────────────────────────────────────────────────────
3933
3934    #[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    // ── LinkStore ─────────────────────────────────────────────────────────
4055
4056    #[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    // ── Searcher (naive — Tantivy is feature-gated) ──────────────────────
4120
4121    #[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    // ── Tantivy search (feature-gated) ──────────────────────────────────
4163
4164    #[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        // Should be findable after put
4229        assert_eq!(
4230            engine
4231                .search("deletable_content", SearchOptions::default())
4232                .unwrap()
4233                .len(),
4234            1
4235        );
4236        // Delete and verify it's gone from search
4237        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    // ── AggregateStore ────────────────────────────────────────────────────
4289
4290    #[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    // ── ready_jobs index behavior (Persistent-specific) ───────────────────────
4406
4407    #[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        // ready_jobs should iterate in priority order (highest first)
4443        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        // First key should contain "high" — it has higher priority
4454        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        // Slight delay so created_at differs
4468        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        // Nack with no delay — should re-index into ready_jobs
4499        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        // Claim with zero lease so it immediately expires
4516        engine
4517            .claim_job_with_options("q", QueueClaimOptions::new(0))
4518            .unwrap();
4519        // The claim method should have reaped the expired lease and re-indexed
4520        let prefix = b"q\0";
4521        let _count = engine.ready_jobs.prefix(prefix).count();
4522        // claim_job called next will reap expired lease and return the job
4523        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    // ── VectorStore ───────────────────────────────────────────────────────
4538
4539    #[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    // ── Reopen tests ─────────────────────────────────────────────────────────
4709
4710    #[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            // Next event should continue at sequence 3
4728            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            // Sequence 1 should not be overwritten
4736            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            // Same idempotency key — must return existing event, not duplicate
4757            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            // New event should continue
4765            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            // Update without vector — old vector must be removed
4810            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    // ── Shared contract tests ───────────────────────────────────────────────
4824
4825    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}