Skip to main content

hyphae_storage/
engine.rs

1// SPDX-License-Identifier: Apache-2.0
2
3use std::path::{Path, PathBuf};
4
5use hyphae_core::{VectorSpaceDefinition, VectorSpaceName};
6use hyphae_retrieval::{LexicalIndexDefinition, LexicalMaterializedCorpus};
7use thiserror::Error;
8use uuid::Uuid;
9
10/// Maximum KV entries returned by one ordered storage scan page.
11pub const MAX_SCAN_PAGE_ENTRIES: usize = 4_096;
12
13use crate::log::transaction_digest;
14use crate::{
15    AppendOutcome, BackupError, BackupInfo, CommitReceipt, DataDirectory, DataDirectoryError,
16    DurableLog, LogError, MaterializedIndexError, Mutation, MutationError, RecoveredTransaction,
17    RecoveryReport, SnapshotError, SnapshotInfo,
18    index::{MaterializedIndex, VectorEntry},
19    manifest::StorageManifest,
20    mutation::validate_key,
21    snapshot::{create_snapshot, verify_snapshot},
22};
23
24/// Failure while opening or operating the durable embedded storage engine.
25#[derive(Debug, Error)]
26pub enum StorageError {
27    /// The data directory could not be initialized or exclusively locked.
28    #[error(transparent)]
29    DataDirectory(#[from] DataDirectoryError),
30
31    /// The authoritative log rejected or failed an operation.
32    #[error(transparent)]
33    Log(#[from] LogError),
34
35    /// The rebuildable materialized index failed before a new log commit.
36    #[error("materialized index failure: {source}")]
37    Index {
38        /// Rebuildable-index failure.
39        #[source]
40        source: Box<MaterializedIndexError>,
41    },
42
43    /// A mutation violates the stable binary codec.
44    #[error(transparent)]
45    Mutation(#[from] MutationError),
46
47    /// The log commit is durable but its index update failed.
48    #[error("transaction {receipt:?} is durable but not materialized; reopen to recover")]
49    CommittedButNotIndexed {
50        /// Receipt proving that the log commit succeeded.
51        receipt: CommitReceipt,
52        /// Rebuildable-index failure.
53        #[source]
54        source: Box<MaterializedIndexError>,
55    },
56
57    /// Reads and further writes are blocked after an index update failure.
58    #[error("materialized index is stale; reopen storage to replay the durable log")]
59    StaleIndex,
60
61    /// Snapshot creation or verification failed.
62    #[error("snapshot failure: {source}")]
63    Snapshot {
64        /// Underlying snapshot failure.
65        #[source]
66        source: Box<SnapshotError>,
67    },
68
69    /// The immutable manifest generation space is exhausted.
70    #[error("storage manifest generation space is exhausted")]
71    ManifestGenerationExhausted,
72
73    /// A prepared compaction segment unexpectedly contains complete frames.
74    #[error("prepared compaction segment is not empty: {path}")]
75    PreparedSegmentNotEmpty {
76        /// Unexpected nonempty segment path.
77        path: PathBuf,
78    },
79
80    /// A KV scan page size is zero or exceeds the hard storage bound.
81    #[error("scan page size {requested} is outside 1..={maximum}")]
82    InvalidScanLimit {
83        /// Requested page size.
84        requested: usize,
85        /// Hard maximum page size.
86        maximum: usize,
87    },
88}
89
90/// Recovery evidence returned when the complete embedded storage layer opens.
91#[derive(Clone, Debug, Eq, PartialEq)]
92pub struct StorageRecoveryReport {
93    /// Authoritative log verification and tail-repair evidence.
94    pub log: RecoveryReport,
95    /// Durable transactions newly applied to the materialized index.
96    pub replayed_transactions: u64,
97}
98
99/// A newly opened storage engine and its recovery evidence.
100#[derive(Debug)]
101pub struct OpenedStorage {
102    /// Ready-to-use storage engine.
103    pub storage: StorageEngine,
104    /// Evidence from log validation and index replay.
105    pub recovery: StorageRecoveryReport,
106}
107
108/// Evidence for one successfully committed compaction generation.
109#[derive(Clone, Debug, Eq, PartialEq)]
110pub struct CompactionReport {
111    /// Newly active immutable manifest generation.
112    pub generation: u64,
113    /// Snapshot anchoring the retired log prefix.
114    pub snapshot: SnapshotInfo,
115    /// Segment that became inactive after the manifest commit.
116    pub retired_segment: PathBuf,
117    /// Whether best-effort physical cleanup removed the retired segment.
118    pub retired_segment_removed: bool,
119}
120
121/// One binary KV entry in canonical key order.
122#[derive(Clone, Debug, Eq, PartialEq)]
123pub struct KvEntry {
124    /// Binary key.
125    pub key: Vec<u8>,
126    /// Opaque binary value.
127    pub value: Vec<u8>,
128}
129
130/// One bounded ordered KV scan page.
131#[derive(Clone, Debug, Eq, PartialEq)]
132pub struct KvPage {
133    /// Entries strictly after the requested cursor.
134    pub entries: Vec<KvEntry>,
135    /// Last emitted key when more entries remain.
136    pub next_after: Option<Vec<u8>>,
137}
138
139/// Result of an online compaction request.
140#[derive(Clone, Debug, Eq, PartialEq)]
141pub enum CompactionOutcome {
142    /// No committed frames exist beyond the already active snapshot anchor.
143    NoChanges {
144        /// Current verified snapshot.
145        snapshot: SnapshotInfo,
146    },
147    /// A new manifest and anchored segment were committed.
148    Compacted(CompactionReport),
149}
150
151/// Single-writer durable KV storage composed from the log and rebuildable redb index.
152#[derive(Debug)]
153pub struct StorageEngine {
154    log: DurableLog,
155    index: MaterializedIndex,
156    index_stale: bool,
157    directory: DataDirectory,
158}
159
160impl StorageEngine {
161    /// Opens a data directory, verifies its log, and catches the index up before use.
162    ///
163    /// # Errors
164    ///
165    /// Returns an error for directory contention, log corruption, invalid committed
166    /// mutations, a divergent index checkpoint, or filesystem failures.
167    pub fn open(path: impl AsRef<Path>) -> Result<OpenedStorage, StorageError> {
168        let directory = DataDirectory::open(path)?;
169        let index_path = directory.path().join("indexes").join("primary.redb");
170        ensure_snapshot_base(&directory, &index_path)?;
171        let (base_sequence, base_digest) = directory.log_anchor();
172        let (log, log_recovery) = DurableLog::open_file_at_version(
173            directory.active_log_path(),
174            base_sequence,
175            base_digest,
176            directory.disk_format_version(),
177        )?;
178        let index = MaterializedIndex::open(index_path)?;
179        let replayed_transactions = index.replay(&log_recovery)?;
180        let _cleanup_complete = directory.cleanup_retired_logs();
181        let storage = Self {
182            log,
183            index,
184            index_stale: false,
185            directory,
186        };
187        Ok(OpenedStorage {
188            storage,
189            recovery: StorageRecoveryReport {
190                log: log_recovery,
191                replayed_transactions,
192            },
193        })
194    }
195
196    /// Returns the owned data-directory path.
197    pub fn data_path(&self) -> &Path {
198        self.directory.path()
199    }
200
201    /// Creates an atomic, independently verifiable backup at the current checkpoint.
202    ///
203    /// # Errors
204    ///
205    /// Returns an error when the snapshot cannot be created, the destination
206    /// exists or is inside the live data directory, or synchronized promotion
207    /// of the complete backup fails.
208    pub fn backup(&self, destination: impl AsRef<Path>) -> Result<BackupInfo, BackupError> {
209        crate::backup::create_backup(self, destination.as_ref())
210    }
211
212    /// Durably commits an atomic batch and then materializes it.
213    ///
214    /// The log is synchronized before redb is updated. If redb fails, the error
215    /// includes the durable commit receipt and this handle blocks reads and writes
216    /// until reopen replays the log.
217    ///
218    /// # Errors
219    ///
220    /// Returns an error for invalid mutations, idempotency conflicts, log I/O,
221    /// or materialized-index failures.
222    pub fn write(
223        &mut self,
224        transaction_id: Uuid,
225        mutations: &[Mutation],
226    ) -> Result<AppendOutcome, StorageError> {
227        if self.index_stale {
228            return Err(StorageError::StaleIndex);
229        }
230        self.index.validate_mutations(mutations)?;
231        if mutations.iter().any(|mutation| {
232            matches!(
233                mutation,
234                Mutation::DefineVectorSpace { .. }
235                    | Mutation::UpsertVector { .. }
236                    | Mutation::DeleteVector { .. }
237                    | Mutation::DefineLexicalIndex { .. }
238            )
239        }) {
240            self.directory.promote_format()?;
241            self.log
242                .set_disk_format_version(self.directory.disk_format_version())?;
243        }
244        let operations = mutations
245            .iter()
246            .map(Mutation::encode)
247            .collect::<Result<Vec<_>, _>>()?;
248        let operation_count =
249            u32::try_from(operations.len()).map_err(|_| LogError::TooManyOperations)?;
250        let requested_digest = transaction_digest(&operations, operation_count)?;
251        if let Some(receipt) = self.index.receipt(transaction_id)? {
252            return if receipt.transaction_digest == requested_digest {
253                Ok(AppendOutcome::Existing(receipt))
254            } else {
255                Err(LogError::IdempotencyConflict { transaction_id }.into())
256            };
257        }
258        let outcome = match self.log.append_transaction(transaction_id, &operations) {
259            Ok(outcome) => outcome,
260            Err(source) => {
261                if self.log.is_poisoned() {
262                    self.index_stale = true;
263                }
264                return Err(source.into());
265            }
266        };
267        let AppendOutcome::Committed(receipt) = outcome else {
268            return Ok(outcome);
269        };
270
271        let transaction = RecoveredTransaction {
272            receipt,
273            operations,
274        };
275        if let Err(source) = self.index.apply(&transaction) {
276            self.index_stale = true;
277            return Err(StorageError::CommittedButNotIndexed {
278                receipt,
279                source: Box::new(source),
280            });
281        }
282        Ok(outcome)
283    }
284
285    /// Reads a binary value from the caught-up materialized index.
286    ///
287    /// # Errors
288    ///
289    /// Returns an error for an empty or oversized key, an index read failure,
290    /// or a handle made stale by a prior post-commit index failure.
291    pub fn get(&self, key: &[u8]) -> Result<Option<Vec<u8>>, StorageError> {
292        if self.index_stale {
293            return Err(StorageError::StaleIndex);
294        }
295        validate_key(key)?;
296        Ok(self.index.get(key)?)
297    }
298
299    /// Returns an immutable vector-space definition when present.
300    ///
301    /// # Errors
302    ///
303    /// Returns an error for a stale handle or malformed materialized state.
304    pub fn vector_space(
305        &self,
306        name: &VectorSpaceName,
307    ) -> Result<Option<VectorSpaceDefinition>, StorageError> {
308        if self.index_stale {
309            return Err(StorageError::StaleIndex);
310        }
311        Ok(self.index.vector_space(name)?)
312    }
313
314    /// Returns an immutable lexical-index definition when present.
315    ///
316    /// # Errors
317    ///
318    /// Returns an error for a stale handle or malformed materialized state.
319    pub fn lexical_index(
320        &self,
321        name: &VectorSpaceName,
322    ) -> Result<Option<LexicalIndexDefinition>, StorageError> {
323        if self.index_stale {
324            return Err(StorageError::StaleIndex);
325        }
326        Ok(self.index.lexical_index(name)?)
327    }
328
329    /// Reads bounded query-relevant statistics from the rebuildable lexical
330    /// projection.
331    ///
332    /// # Errors
333    ///
334    /// Returns an error for a stale handle, malformed projection, exhausted
335    /// candidate budget, or elapsed deadline.
336    pub fn lexical_corpus(
337        &self,
338        definition: &LexicalIndexDefinition,
339        query_tokens: &[String],
340        max_candidates: u64,
341        timeout: std::time::Duration,
342    ) -> Result<LexicalMaterializedCorpus, StorageError> {
343        if self.index_stale {
344            return Err(StorageError::StaleIndex);
345        }
346        Ok(self
347            .index
348            .lexical_corpus(definition, query_tokens, max_candidates, timeout)?)
349    }
350
351    /// Reads one vector space in strict binary-key order under explicit
352    /// candidate and decoded-byte budgets.
353    ///
354    /// # Errors
355    ///
356    /// Returns an error for a stale handle, malformed state, or an exhausted
357    /// count/byte budget. No partial list is returned.
358    pub fn vector_entries(
359        &self,
360        name: &VectorSpaceName,
361        max_candidates: u64,
362        max_bytes: u64,
363    ) -> Result<Vec<VectorEntry>, StorageError> {
364        if self.index_stale {
365            return Err(StorageError::StaleIndex);
366        }
367        Ok(self.index.scan_vectors(name, max_candidates, max_bytes)?)
368    }
369
370    /// Scans one bounded page in strict binary-key order.
371    ///
372    /// `after` is exclusive. A returned `next_after` is present only when at
373    /// least one additional entry exists.
374    ///
375    /// # Errors
376    ///
377    /// Returns an error for a stale handle, invalid cursor key, invalid page
378    /// size, or materialized-index failure.
379    pub fn scan_page(&self, after: Option<&[u8]>, limit: usize) -> Result<KvPage, StorageError> {
380        if self.index_stale {
381            return Err(StorageError::StaleIndex);
382        }
383        if let Some(key) = after {
384            validate_key(key)?;
385        }
386        if limit == 0 || limit > MAX_SCAN_PAGE_ENTRIES {
387            return Err(StorageError::InvalidScanLimit {
388                requested: limit,
389                maximum: MAX_SCAN_PAGE_ENTRIES,
390            });
391        }
392        let mut raw = self.index.scan_after(after, limit.saturating_add(1))?;
393        let has_more = raw.len() > limit;
394        raw.truncate(limit);
395        let next_after = has_more
396            .then(|| raw.last().map(|(key, _)| key.clone()))
397            .flatten();
398        Ok(KvPage {
399            entries: raw
400                .into_iter()
401                .map(|(key, value)| KvEntry { key, value })
402                .collect(),
403            next_after,
404        })
405    }
406
407    /// Returns the internal materialized-index path for diagnostics.
408    pub fn index_path(&self) -> PathBuf {
409        self.directory.path().join("indexes").join("primary.redb")
410    }
411
412    /// Creates or reuses a verified logical snapshot at the current index checkpoint.
413    ///
414    /// # Errors
415    ///
416    /// Returns an error when the live index is stale, cannot be streamed, or the
417    /// snapshot cannot be synchronized, verified, and atomically promoted.
418    pub fn snapshot(&self) -> Result<SnapshotInfo, StorageError> {
419        if self.index_stale {
420            return Err(StorageError::StaleIndex);
421        }
422        let snapshots = self.directory.path().join("snapshots");
423        let temporary = self.directory.path().join("tmp");
424        Ok(create_snapshot(
425            &self.index,
426            &snapshots,
427            &temporary,
428            self.directory.disk_format_version(),
429        )?)
430    }
431
432    /// Retires the active log prefix behind a verified logical snapshot.
433    ///
434    /// The new empty segment is synchronized before an immutable manifest
435    /// generation selects it. Physical deletion of the retired segment happens
436    /// only after that commit and is reported independently.
437    ///
438    /// # Errors
439    ///
440    /// Returns an error while preparing the snapshot, segment, or manifest. A
441    /// poisoned or stale handle must be reopened before compaction.
442    pub fn compact(&mut self) -> Result<CompactionOutcome, StorageError> {
443        if self.index_stale {
444            return Err(StorageError::StaleIndex);
445        }
446        let snapshot = self.snapshot()?;
447        let current = self.directory.manifest();
448        if snapshot.checkpoint_sequence == 0
449            || snapshot.checkpoint_sequence == current.base_sequence
450        {
451            return Ok(CompactionOutcome::NoChanges { snapshot });
452        }
453        let generation = current
454            .generation
455            .checked_add(1)
456            .ok_or(StorageError::ManifestGenerationExhausted)?;
457        let Some(base_digest) = snapshot.checkpoint_digest else {
458            return Err(SnapshotError::Invalid {
459                reason: "compaction snapshot lacks a checkpoint digest",
460            }
461            .into());
462        };
463        let next = StorageManifest {
464            generation,
465            active_segment: generation,
466            base_sequence: snapshot.checkpoint_sequence,
467            base_digest,
468            snapshot_digest: snapshot.snapshot_digest,
469        };
470        let next_segment = self.directory.log_path(generation);
471        let (next_log, prepared) = DurableLog::open_file_at_version(
472            &next_segment,
473            next.base_sequence,
474            next.base_digest,
475            self.directory.disk_format_version(),
476        )?;
477        if prepared.valid_bytes != 0 {
478            return Err(StorageError::PreparedSegmentNotEmpty { path: next_segment });
479        }
480
481        let retired_segment = self.directory.active_log_path();
482        self.directory.commit_manifest(next)?;
483        let retired_log = std::mem::replace(&mut self.log, next_log);
484        drop(retired_log);
485        let retired_segment_removed = remove_retired_segment(&retired_segment);
486        Ok(CompactionOutcome::Compacted(CompactionReport {
487            generation,
488            snapshot,
489            retired_segment,
490            retired_segment_removed,
491        }))
492    }
493}
494
495fn remove_retired_segment(path: &Path) -> bool {
496    match std::fs::remove_file(path) {
497        Ok(()) => {
498            #[cfg(unix)]
499            if let Some(parent) = path.parent()
500                && sync_directory(parent).is_err()
501            {
502                return false;
503            }
504            true
505        }
506        Err(source) if source.kind() == std::io::ErrorKind::NotFound => true,
507        Err(_) => false,
508    }
509}
510
511fn ensure_snapshot_base(directory: &DataDirectory, index_path: &Path) -> Result<(), StorageError> {
512    let manifest = directory.manifest();
513    if manifest.base_sequence == 0 {
514        return Ok(());
515    }
516    let snapshot_path = directory.snapshot_path(manifest.base_sequence);
517    let verified = verify_snapshot(&snapshot_path)?;
518    if verified.checkpoint_sequence != manifest.base_sequence
519        || verified.checkpoint_digest != Some(manifest.base_digest)
520        || verified.snapshot_digest != manifest.snapshot_digest
521    {
522        return Err(SnapshotError::Invalid {
523            reason: "snapshot does not match active storage manifest",
524        }
525        .into());
526    }
527    if index_path.exists() {
528        return Ok(());
529    }
530
531    let temporary_path = directory
532        .path()
533        .join("tmp")
534        .join(format!("index-restore-{}.redb.tmp", Uuid::now_v7()));
535    let restored = MaterializedIndex::restore_from_snapshot(&temporary_path, &snapshot_path)?;
536    if restored != verified {
537        return Err(SnapshotError::Invalid {
538            reason: "snapshot changed while rebuilding the materialized index",
539        }
540        .into());
541    }
542    std::fs::rename(&temporary_path, index_path)
543        .map_err(SnapshotError::from)
544        .map_err(StorageError::from)?;
545    #[cfg(unix)]
546    sync_directory(index_path.parent().ok_or(SnapshotError::Invalid {
547        reason: "materialized index path has no parent",
548    })?)?;
549    Ok(())
550}
551
552#[cfg(unix)]
553fn sync_directory(path: &Path) -> Result<(), StorageError> {
554    std::fs::File::open(path)
555        .and_then(|directory| directory.sync_all())
556        .map_err(SnapshotError::from)
557        .map_err(StorageError::from)
558}
559
560impl From<MaterializedIndexError> for StorageError {
561    fn from(source: MaterializedIndexError) -> Self {
562        Self::Index {
563            source: Box::new(source),
564        }
565    }
566}
567
568impl From<SnapshotError> for StorageError {
569    fn from(source: SnapshotError) -> Self {
570        Self::Snapshot {
571            source: Box::new(source),
572        }
573    }
574}
575
576#[cfg(test)]
577mod tests {
578    use std::{
579        error::Error,
580        fs::{self, OpenOptions},
581        io::{Seek, SeekFrom, Write},
582    };
583
584    use uuid::Uuid;
585
586    use super::{CompactionOutcome, DurableLog, StorageEngine, StorageError, StorageManifest};
587    use crate::{
588        AppendOutcome, DataDirectory, Mutation, SnapshotError, SnapshotReadLimits,
589        index::MaterializedIndex, load_snapshot, test_support::TestDirectory, verify_snapshot,
590    };
591
592    #[test]
593    fn atomic_batches_persist_and_delete() -> Result<(), Box<dyn Error>> {
594        let temporary = TestDirectory::new("storage-kv")?;
595        let root = temporary.path().join("data");
596        let mut opened = StorageEngine::open(&root)?;
597        opened.storage.write(
598            Uuid::now_v7(),
599            &[Mutation::put(b"a", b"one"), Mutation::put(b"b", b"two")],
600        )?;
601        assert_eq!(opened.storage.get(b"a")?, Some(b"one".to_vec()));
602        opened
603            .storage
604            .write(Uuid::now_v7(), &[Mutation::delete(b"a")])?;
605        drop(opened);
606
607        let reopened = StorageEngine::open(&root)?;
608        assert_eq!(reopened.storage.get(b"a")?, None);
609        assert_eq!(reopened.storage.get(b"b")?, Some(b"two".to_vec()));
610        assert_eq!(reopened.recovery.replayed_transactions, 0);
611        Ok(())
612    }
613
614    #[test]
615    fn exact_retry_does_not_reapply_a_batch() -> Result<(), Box<dyn Error>> {
616        let temporary = TestDirectory::new("storage-idempotency")?;
617        let transaction_id = Uuid::now_v7();
618        let mutations = [Mutation::put(b"key", b"value")];
619        let mut opened = StorageEngine::open(temporary.path())?;
620
621        let first = opened.storage.write(transaction_id, &mutations)?;
622        let second = opened.storage.write(transaction_id, &mutations)?;
623        assert!(matches!(first, AppendOutcome::Committed(_)));
624        assert!(matches!(second, AppendOutcome::Existing(_)));
625        assert_eq!(opened.storage.get(b"key")?, Some(b"value".to_vec()));
626
627        let conflict = opened
628            .storage
629            .write(transaction_id, &[Mutation::put(b"key", b"different")]);
630        assert!(matches!(
631            conflict,
632            Err(super::StorageError::Log(
633                crate::LogError::IdempotencyConflict { .. }
634            ))
635        ));
636        Ok(())
637    }
638
639    #[test]
640    fn reopen_replays_a_commit_missing_from_the_index() -> Result<(), Box<dyn Error>> {
641        let temporary = TestDirectory::new("storage-index-replay")?;
642        let root = temporary.path().join("data");
643        let directory = DataDirectory::open(&root)?;
644        let mutation = Mutation::put(b"recovered", b"yes");
645        let mut log = directory.open_log()?;
646        log.log
647            .append_transaction(Uuid::now_v7(), &[mutation.encode()?])?;
648        drop(log);
649        drop(directory);
650
651        let reopened = StorageEngine::open(&root)?;
652        assert_eq!(reopened.recovery.replayed_transactions, 1);
653        assert_eq!(reopened.storage.get(b"recovered")?, Some(b"yes".to_vec()));
654        Ok(())
655    }
656
657    #[test]
658    fn logical_snapshot_is_stable_and_detects_payload_corruption() -> Result<(), Box<dyn Error>> {
659        let temporary = TestDirectory::new("storage-snapshot")?;
660        let mut opened = StorageEngine::open(temporary.path().join("data"))?;
661        opened.storage.write(
662            Uuid::now_v7(),
663            &[
664                Mutation::put(b"beta", b"second"),
665                Mutation::put(b"alpha", b"first"),
666            ],
667        )?;
668
669        let created = opened.storage.snapshot()?;
670        assert_eq!(created.checkpoint_sequence, 4);
671        assert!(created.checkpoint_digest.is_some());
672        assert_eq!(created.entry_count, 2);
673        assert_eq!(created.receipt_count, 1);
674        assert_eq!(verify_snapshot(&created.path)?, created);
675        assert_eq!(opened.storage.snapshot()?, created);
676        let witness = load_snapshot(&created.path, &SnapshotReadLimits::default())?;
677        assert_eq!(witness.info, created);
678        assert_eq!(witness.entries.len(), 2);
679        assert_eq!(witness.entries[0].key, b"alpha");
680        assert_eq!(witness.entries[0].value, b"first");
681        assert!(matches!(
682            load_snapshot(
683                &created.path,
684                &SnapshotReadLimits {
685                    entries: 1,
686                    ..SnapshotReadLimits::default()
687                }
688            ),
689            Err(SnapshotError::EntryLimitExceeded {
690                actual: 2,
691                maximum: 1
692            })
693        ));
694
695        let corrupted_path = temporary.path().join("corrupted.hysnap");
696        fs::copy(&created.path, &corrupted_path)?;
697        let mut corrupted = OpenOptions::new()
698            .read(true)
699            .write(true)
700            .open(&corrupted_path)?;
701        corrupted.seek(SeekFrom::End(-1))?;
702        corrupted.write_all(&[0xff])?;
703        corrupted.sync_all()?;
704        drop(corrupted);
705
706        assert!(matches!(
707            verify_snapshot(&corrupted_path),
708            Err(SnapshotError::Invalid {
709                reason: "CRC32C mismatch"
710            })
711        ));
712        Ok(())
713    }
714
715    #[test]
716    fn empty_storage_has_a_canonical_empty_snapshot() -> Result<(), Box<dyn Error>> {
717        let temporary = TestDirectory::new("storage-empty-snapshot")?;
718        let opened = StorageEngine::open(temporary.path().join("data"))?;
719
720        let snapshot = opened.storage.snapshot()?;
721        assert_eq!(snapshot.checkpoint_sequence, 0);
722        assert_eq!(snapshot.checkpoint_digest, None);
723        assert_eq!(snapshot.entry_count, 0);
724        assert_eq!(snapshot.vector_space_count, 0);
725        assert_eq!(snapshot.vector_count, 0);
726        assert_eq!(snapshot.lexical_index_count, 0);
727        assert_eq!(snapshot.receipt_count, 0);
728        assert_eq!(snapshot.file_bytes, 136);
729        assert_eq!(verify_snapshot(&snapshot.path)?, snapshot);
730        Ok(())
731    }
732
733    #[test]
734    fn snapshot_rebuilds_kv_and_idempotency_state() -> Result<(), Box<dyn Error>> {
735        let temporary = TestDirectory::new("storage-snapshot-restore")?;
736        let root = temporary.path().join("data");
737        let transaction_id = Uuid::now_v7();
738        let mut opened = StorageEngine::open(&root)?;
739        let outcome = opened.storage.write(
740            transaction_id,
741            &[
742                Mutation::put(b"alpha", b"one"),
743                Mutation::put(b"beta", b"two"),
744            ],
745        )?;
746        let AppendOutcome::Committed(receipt) = outcome else {
747            return Err("new transaction was not committed".into());
748        };
749        let snapshot = opened.storage.snapshot()?;
750        drop(opened);
751
752        let restored_path = root.join("tmp/restored.redb");
753        assert_eq!(
754            MaterializedIndex::restore_from_snapshot(&restored_path, &snapshot.path)?,
755            snapshot
756        );
757        let restored = MaterializedIndex::open(&restored_path)?;
758        assert_eq!(restored.get(b"alpha")?, Some(b"one".to_vec()));
759        assert_eq!(restored.get(b"beta")?, Some(b"two".to_vec()));
760        assert_eq!(restored.receipt(transaction_id)?, Some(receipt));
761        Ok(())
762    }
763
764    #[test]
765    fn compaction_retires_history_and_snapshot_rebuilds_the_index() -> Result<(), Box<dyn Error>> {
766        let temporary = TestDirectory::new("storage-compaction")?;
767        let root = temporary.path().join("data");
768        let first_id = Uuid::now_v7();
769        let mut opened = StorageEngine::open(&root)?;
770        let first = opened
771            .storage
772            .write(first_id, &[Mutation::put(b"before", b"one")])?;
773        let AppendOutcome::Committed(first_receipt) = first else {
774            return Err("first transaction was not committed".into());
775        };
776
777        let compacted = opened.storage.compact()?;
778        let CompactionOutcome::Compacted(report) = compacted else {
779            return Err("committed history was not compacted".into());
780        };
781        assert_eq!(report.generation, 2);
782        assert!(report.retired_segment_removed);
783        assert!(!report.retired_segment.exists());
784        assert!(root.join("log/00000000000000000002.hylog").is_file());
785        assert!(matches!(
786            opened.storage.compact()?,
787            CompactionOutcome::NoChanges { .. }
788        ));
789
790        assert_eq!(
791            opened
792                .storage
793                .write(first_id, &[Mutation::put(b"before", b"one")])?,
794            AppendOutcome::Existing(first_receipt)
795        );
796        let second_id = Uuid::now_v7();
797        let second = opened
798            .storage
799            .write(second_id, &[Mutation::put(b"after", b"two")])?;
800        let AppendOutcome::Committed(second_receipt) = second else {
801            return Err("second transaction was not committed".into());
802        };
803        assert_eq!(
804            second_receipt.commit_sequence,
805            first_receipt.commit_sequence + 3
806        );
807        drop(opened);
808
809        fs::remove_file(root.join("indexes/primary.redb"))?;
810        let mut rebuilt = StorageEngine::open(&root)?;
811        assert_eq!(rebuilt.recovery.replayed_transactions, 1);
812        assert_eq!(rebuilt.storage.get(b"before")?, Some(b"one".to_vec()));
813        assert_eq!(rebuilt.storage.get(b"after")?, Some(b"two".to_vec()));
814        assert_eq!(
815            rebuilt
816                .storage
817                .write(first_id, &[Mutation::put(b"before", b"one")])?,
818            AppendOutcome::Existing(first_receipt)
819        );
820        Ok(())
821    }
822
823    #[test]
824    fn orphan_prepared_segment_is_ignored_until_manifest_commit() -> Result<(), Box<dyn Error>> {
825        let temporary = TestDirectory::new("storage-compaction-orphan")?;
826        let root = temporary.path().join("data");
827        let mut opened = StorageEngine::open(&root)?;
828        opened
829            .storage
830            .write(Uuid::now_v7(), &[Mutation::put(b"key", b"value")])?;
831        let snapshot = opened.storage.snapshot()?;
832        let base_digest = snapshot
833            .checkpoint_digest
834            .ok_or("snapshot checkpoint digest is absent")?;
835        let orphan_path = opened.storage.directory.log_path(2);
836        let (orphan, recovery) =
837            DurableLog::open_file_at(&orphan_path, snapshot.checkpoint_sequence, base_digest)?;
838        assert_eq!(recovery.valid_bytes, 0);
839        drop(orphan);
840        drop(opened);
841
842        let mut reopened = StorageEngine::open(&root)?;
843        assert_eq!(reopened.storage.directory.manifest().generation, 1);
844        assert_eq!(reopened.storage.get(b"key")?, Some(b"value".to_vec()));
845        assert!(matches!(
846            reopened.storage.compact()?,
847            CompactionOutcome::Compacted(_)
848        ));
849        assert_eq!(reopened.storage.directory.manifest().generation, 2);
850        Ok(())
851    }
852
853    #[test]
854    fn committed_manifest_wins_before_retired_log_cleanup() -> Result<(), Box<dyn Error>> {
855        let temporary = TestDirectory::new("storage-compaction-committed")?;
856        let root = temporary.path().join("data");
857        let mut opened = StorageEngine::open(&root)?;
858        opened
859            .storage
860            .write(Uuid::now_v7(), &[Mutation::put(b"key", b"value")])?;
861        let snapshot = opened.storage.snapshot()?;
862        let base_digest = snapshot
863            .checkpoint_digest
864            .ok_or("snapshot checkpoint digest is absent")?;
865        let next = StorageManifest {
866            generation: 2,
867            active_segment: 2,
868            base_sequence: snapshot.checkpoint_sequence,
869            base_digest,
870            snapshot_digest: snapshot.snapshot_digest,
871        };
872        let (prepared, recovery) = DurableLog::open_file_at(
873            opened.storage.directory.log_path(2),
874            next.base_sequence,
875            next.base_digest,
876        )?;
877        assert_eq!(recovery.valid_bytes, 0);
878        drop(prepared);
879        opened.storage.directory.commit_manifest(next)?;
880        let retired_path = opened.storage.directory.log_path(1);
881        assert!(retired_path.is_file());
882        drop(opened);
883
884        let reopened = StorageEngine::open(&root)?;
885        assert_eq!(reopened.storage.directory.manifest().generation, 2);
886        assert_eq!(reopened.storage.get(b"key")?, Some(b"value".to_vec()));
887        assert!(!retired_path.exists());
888        Ok(())
889    }
890
891    #[test]
892    fn uncertain_log_sync_blocks_the_handle_until_recovery() -> Result<(), Box<dyn Error>> {
893        let temporary = TestDirectory::new("storage-injected-log-sync")?;
894        let root = temporary.path().join("data");
895        let mut opened = StorageEngine::open(&root)?;
896        opened.storage.log.inject_sync_failure();
897
898        let result = opened
899            .storage
900            .write(Uuid::now_v7(), &[Mutation::put(b"recovered", b"yes")]);
901        assert!(matches!(
902            result,
903            Err(StorageError::Log(crate::LogError::Io(_)))
904        ));
905        assert!(matches!(
906            opened.storage.get(b"recovered"),
907            Err(StorageError::StaleIndex)
908        ));
909        assert!(matches!(
910            opened.storage.snapshot(),
911            Err(StorageError::StaleIndex)
912        ));
913        assert!(matches!(
914            opened.storage.compact(),
915            Err(StorageError::StaleIndex)
916        ));
917        drop(opened);
918
919        let reopened = StorageEngine::open(&root)?;
920        assert_eq!(reopened.recovery.replayed_transactions, 1);
921        assert_eq!(reopened.storage.get(b"recovered")?, Some(b"yes".to_vec()));
922        Ok(())
923    }
924
925    #[test]
926    fn post_commit_index_failure_recovers_from_the_log() -> Result<(), Box<dyn Error>> {
927        let temporary = TestDirectory::new("storage-injected-index")?;
928        let root = temporary.path().join("data");
929        let mut opened = StorageEngine::open(&root)?;
930        opened.storage.index.inject_apply_failure();
931
932        let result = opened
933            .storage
934            .write(Uuid::now_v7(), &[Mutation::put(b"durable", b"yes")]);
935        assert!(matches!(
936            result,
937            Err(StorageError::CommittedButNotIndexed { .. })
938        ));
939        assert!(matches!(
940            opened.storage.get(b"durable"),
941            Err(StorageError::StaleIndex)
942        ));
943        drop(opened);
944
945        let reopened = StorageEngine::open(&root)?;
946        assert_eq!(reopened.recovery.replayed_transactions, 1);
947        assert_eq!(reopened.storage.get(b"durable")?, Some(b"yes".to_vec()));
948        Ok(())
949    }
950
951    #[test]
952    fn kv_scan_pages_are_strictly_ordered_and_exclusive() -> Result<(), Box<dyn Error>> {
953        let temporary = TestDirectory::new("storage-scan-page")?;
954        let mut opened = StorageEngine::open(temporary.path().join("data"))?;
955        opened.storage.write(
956            Uuid::now_v7(),
957            &[
958                Mutation::put(b"c", b"three"),
959                Mutation::put(b"a", b"one"),
960                Mutation::put(b"b", b"two"),
961            ],
962        )?;
963
964        let first = opened.storage.scan_page(None, 2)?;
965        assert_eq!(
966            first
967                .entries
968                .iter()
969                .map(|entry| entry.key.as_slice())
970                .collect::<Vec<_>>(),
971            [b"a".as_slice(), b"b".as_slice()]
972        );
973        assert_eq!(first.next_after, Some(b"b".to_vec()));
974
975        let second = opened.storage.scan_page(first.next_after.as_deref(), 2)?;
976        assert_eq!(second.entries[0].key, b"c");
977        assert_eq!(second.next_after, None);
978        Ok(())
979    }
980}