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::{ExactRetrievalError, 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, MaintenanceLimits, MaterializedIndexError, Mutation, MutationError,
17    RecoveredTransaction, RecoveryLimits, RecoveryReport, SnapshotError, SnapshotInfo,
18    SnapshotReadLimits, StorageLimitError, StorageLimits,
19    index::{KvScanError, MaterializedIndex, VectorEntry, VectorScanError},
20    limits::OperationDeadline,
21    manifest::StorageManifest,
22    mutation::validate_key,
23    snapshot::{create_snapshot, verify_snapshot_with_policy},
24};
25
26/// Failure while opening or operating the durable embedded storage engine.
27#[derive(Debug, Error)]
28pub enum StorageError {
29    /// The data directory could not be initialized or exclusively locked.
30    #[error(transparent)]
31    DataDirectory(#[from] DataDirectoryError),
32
33    /// The authoritative log rejected or failed an operation.
34    #[error(transparent)]
35    Log(#[from] LogError),
36
37    /// The rebuildable materialized index failed before a new log commit.
38    #[error("materialized index failure: {source}")]
39    Index {
40        /// Rebuildable-index failure.
41        #[source]
42        source: Box<MaterializedIndexError>,
43    },
44
45    /// A mutation violates the stable binary codec.
46    #[error(transparent)]
47    Mutation(#[from] MutationError),
48
49    /// The log commit is durable but its index update failed.
50    #[error("transaction {receipt:?} is durable but not materialized; reopen to recover")]
51    CommittedButNotIndexed {
52        /// Receipt proving that the log commit succeeded.
53        receipt: CommitReceipt,
54        /// Rebuildable-index failure.
55        #[source]
56        source: Box<MaterializedIndexError>,
57    },
58
59    /// Reads and further writes are blocked after an index update failure.
60    #[error("materialized index is stale; reopen storage to replay the durable log")]
61    StaleIndex,
62
63    /// Snapshot creation or verification failed.
64    #[error("snapshot failure: {source}")]
65    Snapshot {
66        /// Underlying snapshot failure.
67        #[source]
68        source: Box<SnapshotError>,
69    },
70
71    /// The immutable manifest generation space is exhausted.
72    #[error("storage manifest generation space is exhausted")]
73    ManifestGenerationExhausted,
74
75    /// A prepared compaction segment unexpectedly contains complete frames.
76    #[error("prepared compaction segment is not empty: {path}")]
77    PreparedSegmentNotEmpty {
78        /// Unexpected nonempty segment path.
79        path: PathBuf,
80    },
81
82    /// A KV scan page size is zero or exceeds the hard storage bound.
83    #[error("scan page size {requested} is outside 1..={maximum}")]
84    InvalidScanLimit {
85        /// Requested page size.
86        requested: usize,
87        /// Hard maximum page size.
88        maximum: usize,
89    },
90}
91
92impl From<StorageLimitError> for StorageError {
93    fn from(source: StorageLimitError) -> Self {
94        Self::Snapshot {
95            source: Box::new(SnapshotError::from(source)),
96        }
97    }
98}
99
100/// Failure while reading a vector space under an explicit end-to-end
101/// retrieval deadline.
102#[derive(Debug, Error)]
103pub enum VectorEntriesError {
104    /// Ordinary durable storage failure.
105    #[error(transparent)]
106    Storage(#[from] StorageError),
107    /// The exact-retrieval deadline expired.
108    #[error(transparent)]
109    ExactRetrieval(#[from] ExactRetrievalError),
110}
111
112impl From<MaterializedIndexError> for VectorEntriesError {
113    fn from(source: MaterializedIndexError) -> Self {
114        Self::Storage(StorageError::from(source))
115    }
116}
117
118/// Failure while scanning one page under an explicit aggregate byte budget.
119#[derive(Debug, Error)]
120pub enum ScanPageError {
121    /// Ordinary durable storage failure.
122    #[error(transparent)]
123    Storage(#[from] StorageError),
124    /// Aggregate key and value bytes exceeded caller policy.
125    #[error("scan page byte budget exceeded: {maximum}")]
126    ByteBudgetExceeded {
127        /// Maximum aggregate key and value bytes permitted.
128        maximum: u64,
129    },
130}
131
132/// Recovery evidence returned when the complete embedded storage layer opens.
133#[derive(Clone, Debug, Eq, PartialEq)]
134pub struct StorageRecoveryReport {
135    /// Authoritative log verification and tail-repair evidence.
136    pub log: RecoveryReport,
137    /// Durable transactions newly applied to the materialized index.
138    pub replayed_transactions: u64,
139}
140
141/// A newly opened storage engine and its recovery evidence.
142#[derive(Debug)]
143pub struct OpenedStorage {
144    /// Ready-to-use storage engine.
145    pub storage: StorageEngine,
146    /// Evidence from log validation and index replay.
147    pub recovery: StorageRecoveryReport,
148}
149
150/// Evidence for one successfully committed compaction generation.
151#[derive(Clone, Debug, Eq, PartialEq)]
152pub struct CompactionReport {
153    /// Newly active immutable manifest generation.
154    pub generation: u64,
155    /// Snapshot anchoring the retired log prefix.
156    pub snapshot: SnapshotInfo,
157    /// Segment that became inactive after the manifest commit.
158    pub retired_segment: PathBuf,
159    /// Whether best-effort physical cleanup removed the retired segment.
160    pub retired_segment_removed: bool,
161}
162
163/// One binary KV entry in canonical key order.
164#[derive(Clone, Debug, Eq, PartialEq)]
165pub struct KvEntry {
166    /// Binary key.
167    pub key: Vec<u8>,
168    /// Opaque binary value.
169    pub value: Vec<u8>,
170}
171
172/// One bounded ordered KV scan page.
173#[derive(Clone, Debug, Eq, PartialEq)]
174pub struct KvPage {
175    /// Entries strictly after the requested cursor.
176    pub entries: Vec<KvEntry>,
177    /// Last emitted key when more entries remain.
178    pub next_after: Option<Vec<u8>>,
179}
180
181/// Result of an online compaction request.
182#[derive(Clone, Debug, Eq, PartialEq)]
183pub enum CompactionOutcome {
184    /// No committed frames exist beyond the already active snapshot anchor.
185    NoChanges {
186        /// Current verified snapshot.
187        snapshot: SnapshotInfo,
188    },
189    /// A new manifest and anchored segment were committed.
190    Compacted(CompactionReport),
191}
192
193/// Single-writer durable KV storage composed from the log and rebuildable redb index.
194#[derive(Debug)]
195pub struct StorageEngine {
196    log: DurableLog,
197    index: MaterializedIndex,
198    index_stale: bool,
199    directory: DataDirectory,
200    limits: StorageLimits,
201}
202
203impl StorageEngine {
204    /// Opens a data directory, verifies its log, and catches the index up before use.
205    ///
206    /// # Errors
207    ///
208    /// Returns an error for directory contention, log corruption, invalid committed
209    /// mutations, a divergent index checkpoint, or filesystem failures.
210    pub fn open(path: impl AsRef<Path>) -> Result<OpenedStorage, StorageError> {
211        Self::open_with_limits(path, StorageLimits::compatibility())
212    }
213
214    /// Opens and fully recovers a data directory under finite shared limits.
215    ///
216    /// # Errors
217    ///
218    /// Returns an error for invalid limits, directory contention, corruption,
219    /// exhausted recovery work/bytes, timeout, divergence, or I/O failure.
220    pub fn open_with_limits(
221        path: impl AsRef<Path>,
222        limits: StorageLimits,
223    ) -> Result<OpenedStorage, StorageError> {
224        limits.validate()?;
225        let deadline = OperationDeadline::new(limits.recovery.timeout);
226        let directory =
227            DataDirectory::open_with_limits_and_deadline(path, &limits.recovery, &deadline)?;
228        let index_path = directory.path().join("indexes").join("primary.redb");
229        ensure_snapshot_base(&directory, &index_path, &limits.recovery, &deadline)?;
230        let (base_sequence, base_digest) = directory.log_anchor();
231        let (log, log_recovery) = DurableLog::open_file_at_version_with_limits(
232            directory.active_log_path(),
233            base_sequence,
234            base_digest,
235            directory.disk_format_version(),
236            &limits.recovery,
237            &deadline,
238        )?;
239        let index = MaterializedIndex::open(index_path)?;
240        let replayed_transactions =
241            index.replay_with_limits(&log_recovery, &limits.recovery, &deadline)?;
242        let _cleanup_complete =
243            directory.cleanup_retired_logs_with_limits(&limits.recovery, &deadline);
244        deadline.check()?;
245        let storage = Self {
246            log,
247            index,
248            index_stale: false,
249            directory,
250            limits,
251        };
252        Ok(OpenedStorage {
253            storage,
254            recovery: StorageRecoveryReport {
255                log: log_recovery,
256                replayed_transactions,
257            },
258        })
259    }
260
261    /// Returns the owned data-directory path.
262    pub fn data_path(&self) -> &Path {
263        self.directory.path()
264    }
265
266    /// Creates an atomic, independently verifiable backup at the current checkpoint.
267    ///
268    /// # Errors
269    ///
270    /// Returns an error when the snapshot cannot be created, the destination
271    /// exists or is inside the live data directory, or synchronized promotion
272    /// of the complete backup fails.
273    pub fn backup(&self, destination: impl AsRef<Path>) -> Result<BackupInfo, BackupError> {
274        crate::backup::create_backup(self, destination.as_ref())
275    }
276
277    /// Durably commits an atomic batch and then materializes it.
278    ///
279    /// The log is synchronized before redb is updated. If redb fails, the error
280    /// includes the durable commit receipt and this handle blocks reads and writes
281    /// until reopen replays the log.
282    ///
283    /// # Errors
284    ///
285    /// Returns an error for invalid mutations, idempotency conflicts, log I/O,
286    /// or materialized-index failures.
287    pub fn write(
288        &mut self,
289        transaction_id: Uuid,
290        mutations: &[Mutation],
291    ) -> Result<AppendOutcome, StorageError> {
292        if self.index_stale {
293            return Err(StorageError::StaleIndex);
294        }
295        self.index.validate_mutations(mutations)?;
296        let operations = mutations
297            .iter()
298            .map(Mutation::encode)
299            .collect::<Result<Vec<_>, _>>()?;
300        let operation_count =
301            u32::try_from(operations.len()).map_err(|_| LogError::TooManyOperations)?;
302        let requested_digest = transaction_digest(&operations, operation_count)?;
303        if let Some(receipt) = self.index.receipt(transaction_id)? {
304            return if receipt.transaction_digest == requested_digest {
305                Ok(AppendOutcome::Existing(receipt))
306            } else {
307                Err(LogError::IdempotencyConflict { transaction_id }.into())
308            };
309        }
310        let lexical_deadline = OperationDeadline::new(self.limits.recovery.timeout);
311        self.index.preflight_lexical_mutations(
312            mutations,
313            &self.limits.recovery,
314            &lexical_deadline,
315        )?;
316        if mutations.iter().any(|mutation| {
317            matches!(
318                mutation,
319                Mutation::DefineVectorSpace { .. }
320                    | Mutation::UpsertVector { .. }
321                    | Mutation::DeleteVector { .. }
322                    | Mutation::DefineLexicalIndex { .. }
323            )
324        }) {
325            self.directory.promote_format()?;
326            self.log
327                .set_disk_format_version(self.directory.disk_format_version())?;
328        }
329        let outcome = match self.log.append_transaction(transaction_id, &operations) {
330            Ok(outcome) => outcome,
331            Err(source) => {
332                if self.log.is_poisoned() {
333                    self.index_stale = true;
334                }
335                return Err(source.into());
336            }
337        };
338        let AppendOutcome::Committed(receipt) = outcome else {
339            return Ok(outcome);
340        };
341
342        let transaction = RecoveredTransaction {
343            receipt,
344            operations,
345        };
346        if let Err(source) =
347            self.index
348                .apply_with_limits(&transaction, &self.limits.recovery, &lexical_deadline)
349        {
350            self.index_stale = true;
351            return Err(StorageError::CommittedButNotIndexed {
352                receipt,
353                source: Box::new(source),
354            });
355        }
356        Ok(outcome)
357    }
358
359    /// Reads a binary value from the caught-up materialized index.
360    ///
361    /// # Errors
362    ///
363    /// Returns an error for an empty or oversized key, an index read failure,
364    /// or a handle made stale by a prior post-commit index failure.
365    pub fn get(&self, key: &[u8]) -> Result<Option<Vec<u8>>, StorageError> {
366        if self.index_stale {
367            return Err(StorageError::StaleIndex);
368        }
369        validate_key(key)?;
370        Ok(self.index.get(key)?)
371    }
372
373    /// Returns an immutable vector-space definition when present.
374    ///
375    /// # Errors
376    ///
377    /// Returns an error for a stale handle or malformed materialized state.
378    pub fn vector_space(
379        &self,
380        name: &VectorSpaceName,
381    ) -> Result<Option<VectorSpaceDefinition>, StorageError> {
382        if self.index_stale {
383            return Err(StorageError::StaleIndex);
384        }
385        Ok(self.index.vector_space(name)?)
386    }
387
388    /// Returns an immutable lexical-index definition when present.
389    ///
390    /// # Errors
391    ///
392    /// Returns an error for a stale handle or malformed materialized state.
393    pub fn lexical_index(
394        &self,
395        name: &VectorSpaceName,
396    ) -> Result<Option<LexicalIndexDefinition>, StorageError> {
397        if self.index_stale {
398            return Err(StorageError::StaleIndex);
399        }
400        Ok(self.index.lexical_index(name)?)
401    }
402
403    /// Reads bounded query-relevant statistics from the rebuildable lexical
404    /// projection.
405    ///
406    /// # Errors
407    ///
408    /// Returns an error for a stale handle, malformed projection, exhausted
409    /// candidate budget, or elapsed deadline.
410    pub fn lexical_corpus(
411        &self,
412        definition: &LexicalIndexDefinition,
413        query_tokens: &[String],
414        max_candidates: u64,
415        timeout: std::time::Duration,
416    ) -> Result<LexicalMaterializedCorpus, StorageError> {
417        if self.index_stale {
418            return Err(StorageError::StaleIndex);
419        }
420        Ok(self
421            .index
422            .lexical_corpus(definition, query_tokens, max_candidates, timeout)?)
423    }
424
425    /// Reads one vector space in strict binary-key order under explicit
426    /// candidate and decoded-byte budgets.
427    ///
428    /// # Errors
429    ///
430    /// Returns an error for a stale handle, malformed state, or an exhausted
431    /// count/byte budget. No partial list is returned.
432    pub fn vector_entries(
433        &self,
434        name: &VectorSpaceName,
435        max_candidates: u64,
436        max_bytes: u64,
437    ) -> Result<Vec<VectorEntry>, StorageError> {
438        if self.index_stale {
439            return Err(StorageError::StaleIndex);
440        }
441        Ok(self.index.scan_vectors(name, max_candidates, max_bytes)?)
442    }
443
444    /// Reads one vector space under count, decoded-byte, and wall-clock bounds.
445    ///
446    /// # Errors
447    ///
448    /// Returns an error for a stale handle, malformed state, exhausted
449    /// count/byte budget, or elapsed exact-retrieval deadline. No partial list
450    /// is returned.
451    pub fn vector_entries_with_timeout(
452        &self,
453        name: &VectorSpaceName,
454        max_candidates: u64,
455        max_bytes: u64,
456        timeout: std::time::Duration,
457    ) -> Result<Vec<VectorEntry>, VectorEntriesError> {
458        if self.index_stale {
459            return Err(StorageError::StaleIndex.into());
460        }
461        match self
462            .index
463            .scan_vectors_with_timeout(name, max_candidates, max_bytes, timeout)
464        {
465            Ok(entries) => Ok(entries),
466            Err(VectorScanError::Index(source)) => Err(source.into()),
467            Err(VectorScanError::ExactRetrieval(source)) => Err(source.into()),
468        }
469    }
470
471    /// Scans one bounded page in strict binary-key order.
472    ///
473    /// `after` is exclusive. A returned `next_after` is present only when at
474    /// least one additional entry exists.
475    ///
476    /// # Errors
477    ///
478    /// Returns an error for a stale handle, invalid cursor key, invalid page
479    /// size, or materialized-index failure.
480    pub fn scan_page(&self, after: Option<&[u8]>, limit: usize) -> Result<KvPage, StorageError> {
481        match self.scan_page_with_byte_limit(after, limit, u64::MAX) {
482            Ok(page) => Ok(page),
483            Err(ScanPageError::Storage(source)) => Err(source),
484            Err(ScanPageError::ByteBudgetExceeded { .. }) => {
485                unreachable!("an in-memory page cannot exceed the u64 byte limit")
486            }
487        }
488    }
489
490    /// Scans one ordered page under an aggregate key/value byte budget.
491    ///
492    /// The byte limit is checked inside the materialized-index iterator before
493    /// a key or value is cloned into the returned page.
494    ///
495    /// # Errors
496    ///
497    /// Returns an error for a stale handle, invalid cursor or page size,
498    /// exhausted byte budget, or materialized-index failure.
499    pub fn scan_page_with_byte_limit(
500        &self,
501        after: Option<&[u8]>,
502        limit: usize,
503        max_bytes: u64,
504    ) -> Result<KvPage, ScanPageError> {
505        if self.index_stale {
506            return Err(StorageError::StaleIndex.into());
507        }
508        if let Some(key) = after {
509            validate_key(key).map_err(StorageError::from)?;
510        }
511        if limit == 0 || limit > MAX_SCAN_PAGE_ENTRIES {
512            return Err(StorageError::InvalidScanLimit {
513                requested: limit,
514                maximum: MAX_SCAN_PAGE_ENTRIES,
515            }
516            .into());
517        }
518        let (raw, has_more) = match self
519            .index
520            .scan_after_with_byte_limit(after, limit, max_bytes)
521        {
522            Err(KvScanError::ByteBudgetExceeded { maximum }) => {
523                return Err(ScanPageError::ByteBudgetExceeded { maximum });
524            }
525            Err(KvScanError::Index(source)) => return Err(StorageError::from(source).into()),
526            Ok(raw) => raw,
527        };
528        let next_after = has_more
529            .then(|| raw.last().map(|(key, _)| key.clone()))
530            .flatten();
531        Ok(KvPage {
532            entries: raw
533                .into_iter()
534                .map(|(key, value)| KvEntry { key, value })
535                .collect(),
536            next_after,
537        })
538    }
539
540    /// Returns the internal materialized-index path for diagnostics.
541    pub fn index_path(&self) -> PathBuf {
542        self.directory.path().join("indexes").join("primary.redb")
543    }
544
545    /// Creates or reuses a verified logical snapshot at the current index checkpoint.
546    ///
547    /// # Errors
548    ///
549    /// Returns an error when the live index is stale, cannot be streamed, or the
550    /// snapshot cannot be synchronized, verified, and atomically promoted.
551    pub fn snapshot(&self) -> Result<SnapshotInfo, StorageError> {
552        let limits = self.limits.maintenance.clone();
553        self.snapshot_with_limits(&limits)
554    }
555
556    /// Creates or reuses a verified snapshot under explicit finite limits.
557    ///
558    /// # Errors
559    ///
560    /// Returns an error for invalid limits, stale state, exhausted records or
561    /// bytes, timeout, verification failure, or I/O failure.
562    pub fn snapshot_with_limits(
563        &self,
564        limits: &MaintenanceLimits,
565    ) -> Result<SnapshotInfo, StorageError> {
566        limits.validate()?;
567        let deadline = OperationDeadline::new(limits.timeout);
568        self.snapshot_with_deadline(limits, &deadline)
569    }
570
571    fn snapshot_with_deadline(
572        &self,
573        limits: &MaintenanceLimits,
574        deadline: &OperationDeadline,
575    ) -> Result<SnapshotInfo, StorageError> {
576        if self.index_stale {
577            return Err(StorageError::StaleIndex);
578        }
579        let snapshots = self.directory.path().join("snapshots");
580        let temporary = self.directory.path().join("tmp");
581        let checkpoint = self.index.checkpoint()?;
582        let snapshot_target = self.directory.snapshot_path(checkpoint.sequence);
583        let snapshot_is_new = self.directory.reserve_target_entry(
584            "snapshots",
585            &snapshot_target,
586            &self.limits.recovery,
587            deadline,
588        )?;
589        if snapshot_is_new {
590            self.directory
591                .reserve_directory_entries("tmp", 1, &self.limits.recovery, deadline)?;
592        }
593        Ok(create_snapshot(
594            &self.index,
595            &snapshots,
596            &temporary,
597            self.directory.disk_format_version(),
598            &limits.snapshot,
599            deadline,
600        )?)
601    }
602
603    /// Retires the active log prefix behind a verified logical snapshot.
604    ///
605    /// The new empty segment is synchronized before an immutable manifest
606    /// generation selects it. Physical deletion of the retired segment happens
607    /// only after that commit and is reported independently.
608    ///
609    /// # Errors
610    ///
611    /// Returns an error while preparing the snapshot, segment, or manifest. A
612    /// poisoned or stale handle must be reopened before compaction.
613    pub fn compact(&mut self) -> Result<CompactionOutcome, StorageError> {
614        let limits = self.limits.maintenance.clone();
615        self.compact_with_limits(&limits)
616    }
617
618    /// Compacts the active generation under one shared finite deadline.
619    ///
620    /// # Errors
621    ///
622    /// Returns an error before the manifest commit for invalid limits, stale
623    /// state, exhausted snapshot policy, timeout, corruption, or I/O failure.
624    pub fn compact_with_limits(
625        &mut self,
626        limits: &MaintenanceLimits,
627    ) -> Result<CompactionOutcome, StorageError> {
628        limits.validate()?;
629        let deadline = OperationDeadline::new(limits.timeout);
630        if self.index_stale {
631            return Err(StorageError::StaleIndex);
632        }
633        let effective_limits = MaintenanceLimits {
634            timeout: limits.timeout,
635            snapshot: intersect_snapshot_limits(&limits.snapshot, &self.limits.recovery.snapshot),
636        };
637        let current = self.directory.manifest();
638        let checkpoint = self.index.checkpoint()?;
639        if checkpoint.sequence == 0 || checkpoint.sequence == current.base_sequence {
640            let snapshot = self.snapshot_with_deadline(&effective_limits, &deadline)?;
641            return Ok(CompactionOutcome::NoChanges { snapshot });
642        }
643        let generation = current
644            .generation
645            .checked_add(1)
646            .ok_or(StorageError::ManifestGenerationExhausted)?;
647        let prospective_manifest = StorageManifest {
648            generation,
649            active_segment: generation,
650            base_sequence: checkpoint.sequence,
651            base_digest: checkpoint.digest.unwrap_or([0; 32]),
652            snapshot_digest: [0; 32],
653        };
654        let snapshot_target = self.directory.snapshot_path(checkpoint.sequence);
655        let snapshot_is_new = self.directory.reserve_target_entry(
656            "snapshots",
657            &snapshot_target,
658            &self.limits.recovery,
659            &deadline,
660        )?;
661        let next_segment = self.directory.log_path(generation);
662        self.directory.reserve_target_entry(
663            "log",
664            &next_segment,
665            &self.limits.recovery,
666            &deadline,
667        )?;
668        let manifest_target = prospective_manifest.path(self.directory.path());
669        let manifest_is_new = self.directory.reserve_target_entry(
670            "manifest",
671            &manifest_target,
672            &self.limits.recovery,
673            &deadline,
674        )?;
675        if snapshot_is_new || manifest_is_new {
676            self.directory
677                .reserve_directory_entries("tmp", 1, &self.limits.recovery, &deadline)?;
678        }
679        let snapshot = self.snapshot_with_deadline(&effective_limits, &deadline)?;
680        let Some(base_digest) = snapshot.checkpoint_digest else {
681            return Err(SnapshotError::Invalid {
682                reason: "compaction snapshot lacks a checkpoint digest",
683            }
684            .into());
685        };
686        let next = StorageManifest {
687            generation,
688            active_segment: generation,
689            base_sequence: snapshot.checkpoint_sequence,
690            base_digest,
691            snapshot_digest: snapshot.snapshot_digest,
692        };
693        let (next_log, prepared) = DurableLog::open_file_at_version_with_limits(
694            &next_segment,
695            next.base_sequence,
696            next.base_digest,
697            self.directory.disk_format_version(),
698            &self.limits.recovery,
699            &deadline,
700        )?;
701        if prepared.valid_bytes != 0 {
702            return Err(StorageError::PreparedSegmentNotEmpty { path: next_segment });
703        }
704
705        let retired_segment = self.directory.active_log_path();
706        deadline.check()?;
707        if let Err(source) =
708            self.directory
709                .commit_manifest_with_limits(next, &self.limits.recovery, &deadline)
710        {
711            self.index_stale = true;
712            return Err(source.into());
713        }
714        let retired_log = std::mem::replace(&mut self.log, next_log);
715        drop(retired_log);
716        let retired_segment_removed = remove_retired_segment(&retired_segment);
717        Ok(CompactionOutcome::Compacted(CompactionReport {
718            generation,
719            snapshot,
720            retired_segment,
721            retired_segment_removed,
722        }))
723    }
724}
725
726fn intersect_snapshot_limits(
727    maintenance: &SnapshotReadLimits,
728    recovery: &SnapshotReadLimits,
729) -> SnapshotReadLimits {
730    SnapshotReadLimits {
731        file_bytes: maintenance.file_bytes.min(recovery.file_bytes),
732        entries: maintenance.entries.min(recovery.entries),
733        decoded_bytes: maintenance.decoded_bytes.min(recovery.decoded_bytes),
734    }
735}
736
737fn remove_retired_segment(path: &Path) -> bool {
738    match std::fs::remove_file(path) {
739        Ok(()) => {
740            #[cfg(unix)]
741            if let Some(parent) = path.parent()
742                && sync_directory(parent).is_err()
743            {
744                return false;
745            }
746            true
747        }
748        Err(source) if source.kind() == std::io::ErrorKind::NotFound => true,
749        Err(_) => false,
750    }
751}
752
753fn ensure_snapshot_base(
754    directory: &DataDirectory,
755    index_path: &Path,
756    limits: &RecoveryLimits,
757    deadline: &OperationDeadline,
758) -> Result<(), StorageError> {
759    deadline.check()?;
760    let manifest = directory.manifest();
761    if manifest.base_sequence == 0 {
762        return Ok(());
763    }
764    let snapshot_path = directory.snapshot_path(manifest.base_sequence);
765    let verified = verify_snapshot_with_policy(&snapshot_path, &limits.snapshot, deadline)?;
766    if verified.checkpoint_sequence != manifest.base_sequence
767        || verified.checkpoint_digest != Some(manifest.base_digest)
768        || verified.snapshot_digest != manifest.snapshot_digest
769    {
770        return Err(SnapshotError::Invalid {
771            reason: "snapshot does not match active storage manifest",
772        }
773        .into());
774    }
775    if index_path.exists() {
776        return Ok(());
777    }
778
779    let temporary_path = directory
780        .path()
781        .join("tmp")
782        .join(format!("index-restore-{}.redb.tmp", Uuid::now_v7()));
783    let mut temporary_guard = TemporaryFileGuard::new(temporary_path.clone());
784    let restored = MaterializedIndex::restore_from_snapshot_with_limits(
785        &temporary_path,
786        &snapshot_path,
787        &limits.snapshot,
788        limits,
789        deadline,
790    )?;
791    if restored != verified {
792        return Err(SnapshotError::Invalid {
793            reason: "snapshot changed while rebuilding the materialized index",
794        }
795        .into());
796    }
797    deadline.check()?;
798    std::fs::rename(&temporary_path, index_path)
799        .map_err(SnapshotError::from)
800        .map_err(StorageError::from)?;
801    temporary_guard.disarm();
802    #[cfg(unix)]
803    sync_directory(index_path.parent().ok_or(SnapshotError::Invalid {
804        reason: "materialized index path has no parent",
805    })?)?;
806    Ok(())
807}
808
809struct TemporaryFileGuard {
810    path: PathBuf,
811    armed: bool,
812}
813
814impl TemporaryFileGuard {
815    fn new(path: PathBuf) -> Self {
816        Self { path, armed: true }
817    }
818
819    fn disarm(&mut self) {
820        self.armed = false;
821    }
822}
823
824impl Drop for TemporaryFileGuard {
825    fn drop(&mut self) {
826        if self.armed {
827            let _ignored = std::fs::remove_file(&self.path);
828        }
829    }
830}
831
832#[cfg(unix)]
833fn sync_directory(path: &Path) -> Result<(), StorageError> {
834    std::fs::File::open(path)
835        .and_then(|directory| directory.sync_all())
836        .map_err(SnapshotError::from)
837        .map_err(StorageError::from)
838}
839
840impl From<MaterializedIndexError> for StorageError {
841    fn from(source: MaterializedIndexError) -> Self {
842        Self::Index {
843            source: Box::new(source),
844        }
845    }
846}
847
848impl From<SnapshotError> for StorageError {
849    fn from(source: SnapshotError) -> Self {
850        Self::Snapshot {
851            source: Box::new(source),
852        }
853    }
854}
855
856#[cfg(test)]
857mod tests {
858    use std::{
859        collections::BTreeMap,
860        error::Error,
861        fs::{self, OpenOptions},
862        io::{Seek, SeekFrom, Write},
863        time::Duration,
864    };
865
866    use hyphae_core::VectorSpaceName;
867    use hyphae_query::{FieldPath, Value, encode_document};
868    use hyphae_retrieval::{LexicalError, LexicalField, LexicalIndexDefinition};
869    use uuid::Uuid;
870
871    use super::{
872        CompactionOutcome, DurableLog, ScanPageError, StorageEngine, StorageError, StorageManifest,
873    };
874    use crate::{
875        AppendOutcome, DataDirectory, MaintenanceLimits, ManifestError, Mutation, SnapshotError,
876        SnapshotReadLimits, StorageLimitError, StorageLimits, index::MaterializedIndex,
877        load_snapshot, load_snapshot_for_migration, load_snapshot_with_timeout,
878        storage_limit_from_io, test_support::TestDirectory, verify_snapshot,
879    };
880
881    fn lexical_document(text: &str) -> Result<Vec<u8>, hyphae_query::DocumentError> {
882        encode_document(&Value::Object(BTreeMap::from([(
883            "body".to_owned(),
884            Value::String(text.to_owned()),
885        )])))
886    }
887
888    fn lexical_definition(
889        name: &str,
890    ) -> Result<LexicalIndexDefinition, Box<dyn std::error::Error>> {
891        Ok(LexicalIndexDefinition::new(
892            VectorSpaceName::new(name)?,
893            vec![LexicalField {
894                path: FieldPath::field("body"),
895                weight_micros: 1_000_000,
896            }],
897        )?)
898    }
899
900    #[test]
901    fn atomic_batches_persist_and_delete() -> Result<(), Box<dyn Error>> {
902        let temporary = TestDirectory::new("storage-kv")?;
903        let root = temporary.path().join("data");
904        let mut opened = StorageEngine::open(&root)?;
905        opened.storage.write(
906            Uuid::now_v7(),
907            &[Mutation::put(b"a", b"one"), Mutation::put(b"b", b"two")],
908        )?;
909        assert_eq!(opened.storage.get(b"a")?, Some(b"one".to_vec()));
910        opened
911            .storage
912            .write(Uuid::now_v7(), &[Mutation::delete(b"a")])?;
913        drop(opened);
914
915        let reopened = StorageEngine::open(&root)?;
916        assert_eq!(reopened.storage.get(b"a")?, None);
917        assert_eq!(reopened.storage.get(b"b")?, Some(b"two".to_vec()));
918        assert_eq!(reopened.recovery.replayed_transactions, 0);
919        Ok(())
920    }
921
922    #[test]
923    fn exact_retry_does_not_reapply_a_batch() -> Result<(), Box<dyn Error>> {
924        let temporary = TestDirectory::new("storage-idempotency")?;
925        let transaction_id = Uuid::now_v7();
926        let mutations = [Mutation::put(b"key", b"value")];
927        let mut opened = StorageEngine::open(temporary.path())?;
928
929        let first = opened.storage.write(transaction_id, &mutations)?;
930        let second = opened.storage.write(transaction_id, &mutations)?;
931        assert!(matches!(first, AppendOutcome::Committed(_)));
932        assert!(matches!(second, AppendOutcome::Existing(_)));
933        assert_eq!(opened.storage.get(b"key")?, Some(b"value".to_vec()));
934
935        let conflict = opened
936            .storage
937            .write(transaction_id, &[Mutation::put(b"key", b"different")]);
938        assert!(matches!(
939            conflict,
940            Err(super::StorageError::Log(
941                crate::LogError::IdempotencyConflict { .. }
942            ))
943        ));
944        Ok(())
945    }
946
947    #[test]
948    fn reopen_replays_a_commit_missing_from_the_index() -> Result<(), Box<dyn Error>> {
949        let temporary = TestDirectory::new("storage-index-replay")?;
950        let root = temporary.path().join("data");
951        let directory = DataDirectory::open(&root)?;
952        let mutation = Mutation::put(b"recovered", b"yes");
953        let mut log = directory.open_log()?;
954        log.log
955            .append_transaction(Uuid::now_v7(), &[mutation.encode()?])?;
956        drop(log);
957        drop(directory);
958
959        let reopened = StorageEngine::open(&root)?;
960        assert_eq!(reopened.recovery.replayed_transactions, 1);
961        assert_eq!(reopened.storage.get(b"recovered")?, Some(b"yes".to_vec()));
962        Ok(())
963    }
964
965    #[test]
966    fn logical_snapshot_is_stable_and_detects_payload_corruption() -> Result<(), Box<dyn Error>> {
967        let temporary = TestDirectory::new("storage-snapshot")?;
968        let mut opened = StorageEngine::open(temporary.path().join("data"))?;
969        opened.storage.write(
970            Uuid::now_v7(),
971            &[
972                Mutation::put(b"beta", b"second"),
973                Mutation::put(b"alpha", b"first"),
974            ],
975        )?;
976
977        let created = opened.storage.snapshot()?;
978        assert_eq!(created.checkpoint_sequence, 4);
979        assert!(created.checkpoint_digest.is_some());
980        assert_eq!(created.entry_count, 2);
981        assert_eq!(created.receipt_count, 1);
982        assert_eq!(verify_snapshot(&created.path)?, created);
983        assert_eq!(opened.storage.snapshot()?, created);
984        let (witness, receipts) =
985            load_snapshot_for_migration(&created.path, &SnapshotReadLimits::default())?;
986        assert_eq!(witness.info, created);
987        assert_eq!(witness.entries.len(), 2);
988        assert_eq!(receipts.0.len(), 1);
989        assert_eq!(witness.entries[0].key, b"alpha");
990        assert_eq!(witness.entries[0].value, b"first");
991        assert!(matches!(
992            load_snapshot_with_timeout(
993                &created.path,
994                &SnapshotReadLimits::default(),
995                Duration::ZERO
996            ),
997            Err(source) if source.is_timeout()
998        ));
999        assert!(matches!(
1000            load_snapshot(
1001                &created.path,
1002                &SnapshotReadLimits {
1003                    entries: 1,
1004                    ..SnapshotReadLimits::default()
1005                }
1006            ),
1007            Err(SnapshotError::EntryLimitExceeded {
1008                actual: 2,
1009                maximum: 1
1010            })
1011        ));
1012
1013        let corrupted_path = temporary.path().join("corrupted.hysnap");
1014        fs::copy(&created.path, &corrupted_path)?;
1015        let mut corrupted = OpenOptions::new()
1016            .read(true)
1017            .write(true)
1018            .open(&corrupted_path)?;
1019        corrupted.seek(SeekFrom::End(-1))?;
1020        corrupted.write_all(&[0xff])?;
1021        corrupted.sync_all()?;
1022        drop(corrupted);
1023
1024        assert!(matches!(
1025            load_snapshot(
1026                &corrupted_path,
1027                &SnapshotReadLimits {
1028                    entries: 1,
1029                    ..SnapshotReadLimits::default()
1030                }
1031            ),
1032            Err(SnapshotError::EntryLimitExceeded {
1033                actual: 2,
1034                maximum: 1
1035            })
1036        ));
1037        assert!(matches!(
1038            load_snapshot(
1039                &corrupted_path,
1040                &SnapshotReadLimits {
1041                    decoded_bytes: 0,
1042                    ..SnapshotReadLimits::default()
1043                }
1044            ),
1045            Err(SnapshotError::DecodedBytesLimitExceeded { maximum: 0 })
1046        ));
1047        assert!(matches!(
1048            verify_snapshot(&corrupted_path),
1049            Err(SnapshotError::Invalid {
1050                reason: "CRC32C mismatch"
1051            })
1052        ));
1053        Ok(())
1054    }
1055
1056    #[test]
1057    fn empty_storage_has_a_canonical_empty_snapshot() -> Result<(), Box<dyn Error>> {
1058        let temporary = TestDirectory::new("storage-empty-snapshot")?;
1059        let opened = StorageEngine::open(temporary.path().join("data"))?;
1060
1061        let snapshot = opened.storage.snapshot()?;
1062        assert_eq!(snapshot.checkpoint_sequence, 0);
1063        assert_eq!(snapshot.checkpoint_digest, None);
1064        assert_eq!(snapshot.entry_count, 0);
1065        assert_eq!(snapshot.vector_space_count, 0);
1066        assert_eq!(snapshot.vector_count, 0);
1067        assert_eq!(snapshot.lexical_index_count, 0);
1068        assert_eq!(snapshot.receipt_count, 0);
1069        assert_eq!(snapshot.file_bytes, 136);
1070        assert_eq!(verify_snapshot(&snapshot.path)?, snapshot);
1071        Ok(())
1072    }
1073
1074    #[test]
1075    fn snapshot_rebuilds_kv_and_idempotency_state() -> Result<(), Box<dyn Error>> {
1076        let temporary = TestDirectory::new("storage-snapshot-restore")?;
1077        let root = temporary.path().join("data");
1078        let transaction_id = Uuid::now_v7();
1079        let mut opened = StorageEngine::open(&root)?;
1080        let outcome = opened.storage.write(
1081            transaction_id,
1082            &[
1083                Mutation::put(b"alpha", b"one"),
1084                Mutation::put(b"beta", b"two"),
1085            ],
1086        )?;
1087        let AppendOutcome::Committed(receipt) = outcome else {
1088            return Err("new transaction was not committed".into());
1089        };
1090        let snapshot = opened.storage.snapshot()?;
1091        drop(opened);
1092
1093        let restored_path = root.join("tmp/restored.redb");
1094        assert_eq!(
1095            MaterializedIndex::restore_from_snapshot(&restored_path, &snapshot.path)?,
1096            snapshot
1097        );
1098        let restored = MaterializedIndex::open(&restored_path)?;
1099        assert_eq!(restored.get(b"alpha")?, Some(b"one".to_vec()));
1100        assert_eq!(restored.get(b"beta")?, Some(b"two".to_vec()));
1101        assert_eq!(restored.receipt(transaction_id)?, Some(receipt));
1102        Ok(())
1103    }
1104
1105    #[test]
1106    fn compaction_retires_history_and_snapshot_rebuilds_the_index() -> Result<(), Box<dyn Error>> {
1107        let temporary = TestDirectory::new("storage-compaction")?;
1108        let root = temporary.path().join("data");
1109        let first_id = Uuid::now_v7();
1110        let mut opened = StorageEngine::open(&root)?;
1111        let first = opened
1112            .storage
1113            .write(first_id, &[Mutation::put(b"before", b"one")])?;
1114        let AppendOutcome::Committed(first_receipt) = first else {
1115            return Err("first transaction was not committed".into());
1116        };
1117
1118        let compacted = opened.storage.compact()?;
1119        let CompactionOutcome::Compacted(report) = compacted else {
1120            return Err("committed history was not compacted".into());
1121        };
1122        assert_eq!(report.generation, 2);
1123        assert!(report.retired_segment_removed);
1124        assert!(!report.retired_segment.exists());
1125        assert!(root.join("log/00000000000000000002.hylog").is_file());
1126        assert!(matches!(
1127            opened.storage.compact()?,
1128            CompactionOutcome::NoChanges { .. }
1129        ));
1130
1131        assert_eq!(
1132            opened
1133                .storage
1134                .write(first_id, &[Mutation::put(b"before", b"one")])?,
1135            AppendOutcome::Existing(first_receipt)
1136        );
1137        let second_id = Uuid::now_v7();
1138        let second = opened
1139            .storage
1140            .write(second_id, &[Mutation::put(b"after", b"two")])?;
1141        let AppendOutcome::Committed(second_receipt) = second else {
1142            return Err("second transaction was not committed".into());
1143        };
1144        assert_eq!(
1145            second_receipt.commit_sequence,
1146            first_receipt.commit_sequence + 3
1147        );
1148        drop(opened);
1149
1150        fs::remove_file(root.join("indexes/primary.redb"))?;
1151        let mut rebuilt = StorageEngine::open(&root)?;
1152        assert_eq!(rebuilt.recovery.replayed_transactions, 1);
1153        assert_eq!(rebuilt.storage.get(b"before")?, Some(b"one".to_vec()));
1154        assert_eq!(rebuilt.storage.get(b"after")?, Some(b"two".to_vec()));
1155        assert_eq!(
1156            rebuilt
1157                .storage
1158                .write(first_id, &[Mutation::put(b"before", b"one")])?,
1159            AppendOutcome::Existing(first_receipt)
1160        );
1161        Ok(())
1162    }
1163
1164    #[test]
1165    fn open_and_compaction_limits_fail_without_advancing_durable_state()
1166    -> Result<(), Box<dyn Error>> {
1167        let temporary = TestDirectory::new("bounded-open-compaction")?;
1168        let root = temporary.path().join("data");
1169        let mut opened = StorageEngine::open(&root)?;
1170        opened.storage.write(
1171            Uuid::now_v7(),
1172            &[Mutation::put(b"a", b"one"), Mutation::put(b"b", b"two")],
1173        )?;
1174        let active_log = opened.storage.directory.active_log_path();
1175        let log_bytes = fs::metadata(&active_log)?.len();
1176        let generation = opened.storage.directory.manifest().generation;
1177
1178        let too_few_records = MaintenanceLimits {
1179            snapshot: SnapshotReadLimits {
1180                entries: 1,
1181                ..SnapshotReadLimits::default()
1182            },
1183            ..MaintenanceLimits::default()
1184        };
1185        assert!(matches!(
1186            opened.storage.compact_with_limits(&too_few_records),
1187            Err(StorageError::Snapshot { source })
1188                if matches!(
1189                    source.as_ref(),
1190                    SnapshotError::EntryLimitExceeded {
1191                        actual: 2,
1192                        maximum: 1
1193                    }
1194                )
1195        ));
1196        assert_eq!(opened.storage.directory.manifest().generation, generation);
1197        assert!(active_log.is_file());
1198        drop(opened);
1199
1200        let limits = StorageLimits {
1201            recovery: crate::RecoveryLimits {
1202                max_log_file_bytes: log_bytes - 1,
1203                ..crate::RecoveryLimits::default()
1204            },
1205            ..StorageLimits::default()
1206        };
1207        assert!(matches!(
1208            StorageEngine::open_with_limits(&root, limits),
1209            Err(StorageError::Log(crate::LogError::Io(source)))
1210                if matches!(
1211                    storage_limit_from_io(&source),
1212                    Some(StorageLimitError::LogFileBytesExceeded { actual, maximum })
1213                        if *actual == log_bytes && *maximum == log_bytes - 1
1214                )
1215        ));
1216        assert_eq!(fs::metadata(&active_log)?.len(), log_bytes);
1217        Ok(())
1218    }
1219
1220    #[test]
1221    fn compaction_snapshot_policy_is_reopenable_at_the_exact_recovery_limit()
1222    -> Result<(), Box<dyn Error>> {
1223        let temporary = TestDirectory::new("compaction-recovery-snapshot-intersection")?;
1224
1225        let exact_root = temporary.path().join("exact");
1226        let exact_limits = StorageLimits {
1227            recovery: crate::RecoveryLimits {
1228                snapshot: SnapshotReadLimits {
1229                    entries: 2,
1230                    ..SnapshotReadLimits::default()
1231                },
1232                ..crate::RecoveryLimits::default()
1233            },
1234            maintenance: MaintenanceLimits {
1235                snapshot: SnapshotReadLimits {
1236                    entries: 10,
1237                    ..SnapshotReadLimits::default()
1238                },
1239                ..MaintenanceLimits::default()
1240            },
1241        };
1242        let mut exact = StorageEngine::open_with_limits(&exact_root, exact_limits.clone())?;
1243        exact.storage.write(
1244            Uuid::now_v7(),
1245            &[Mutation::put(b"a", b"one"), Mutation::put(b"b", b"two")],
1246        )?;
1247        assert!(matches!(
1248            exact.storage.compact()?,
1249            CompactionOutcome::Compacted(_)
1250        ));
1251        drop(exact);
1252        let reopened = StorageEngine::open_with_limits(&exact_root, exact_limits)?;
1253        assert_eq!(reopened.storage.get(b"a")?, Some(b"one".to_vec()));
1254        assert_eq!(reopened.storage.get(b"b")?, Some(b"two".to_vec()));
1255        drop(reopened);
1256
1257        let rejected_root = temporary.path().join("exact-plus-one");
1258        let rejected_limits = StorageLimits {
1259            recovery: crate::RecoveryLimits {
1260                snapshot: SnapshotReadLimits {
1261                    entries: 1,
1262                    ..SnapshotReadLimits::default()
1263                },
1264                ..crate::RecoveryLimits::default()
1265            },
1266            maintenance: MaintenanceLimits {
1267                snapshot: SnapshotReadLimits {
1268                    entries: 2,
1269                    ..SnapshotReadLimits::default()
1270                },
1271                ..MaintenanceLimits::default()
1272            },
1273        };
1274        let mut rejected =
1275            StorageEngine::open_with_limits(&rejected_root, rejected_limits.clone())?;
1276        rejected.storage.write(
1277            Uuid::now_v7(),
1278            &[Mutation::put(b"a", b"one"), Mutation::put(b"b", b"two")],
1279        )?;
1280        assert!(matches!(
1281            rejected.storage.compact(),
1282            Err(StorageError::Snapshot { source })
1283                if matches!(
1284                    source.as_ref(),
1285                    SnapshotError::EntryLimitExceeded {
1286                        actual: 2,
1287                        maximum: 1
1288                    }
1289                )
1290        ));
1291        assert_eq!(rejected.storage.directory.manifest().generation, 1);
1292        drop(rejected);
1293        let reopened = StorageEngine::open_with_limits(&rejected_root, rejected_limits)?;
1294        assert_eq!(reopened.storage.get(b"a")?, Some(b"one".to_vec()));
1295        assert_eq!(reopened.storage.get(b"b")?, Some(b"two".to_vec()));
1296        Ok(())
1297    }
1298
1299    #[test]
1300    fn compaction_reserves_directory_entries_before_creating_any_artifact()
1301    -> Result<(), Box<dyn Error>> {
1302        let temporary = TestDirectory::new("compaction-directory-reservation")?;
1303        let root = temporary.path().join("data");
1304        let limits = StorageLimits {
1305            recovery: crate::RecoveryLimits {
1306                max_directory_entries: 1,
1307                ..crate::RecoveryLimits::default()
1308            },
1309            ..StorageLimits::default()
1310        };
1311        let mut opened = StorageEngine::open_with_limits(&root, limits.clone())?;
1312        opened
1313            .storage
1314            .write(Uuid::now_v7(), &[Mutation::put(b"key", b"value")])?;
1315
1316        assert!(matches!(
1317            opened.storage.compact(),
1318            Err(StorageError::DataDirectory(
1319                crate::DataDirectoryError::Manifest(ManifestError::Io(source))
1320            ))
1321                if matches!(
1322                    storage_limit_from_io(&source),
1323                    Some(StorageLimitError::DirectoryEntriesExceeded { maximum: 1 })
1324                )
1325        ));
1326        assert_eq!(fs::read_dir(root.join("snapshots"))?.count(), 0);
1327        assert_eq!(fs::read_dir(root.join("log"))?.count(), 1);
1328        assert_eq!(fs::read_dir(root.join("manifest"))?.count(), 1);
1329        assert_eq!(fs::read_dir(root.join("tmp"))?.count(), 0);
1330        assert_eq!(opened.storage.directory.manifest().generation, 1);
1331        drop(opened);
1332
1333        let reopened = StorageEngine::open_with_limits(&root, limits)?;
1334        assert_eq!(reopened.storage.get(b"key")?, Some(b"value".to_vec()));
1335        Ok(())
1336    }
1337
1338    #[test]
1339    fn existing_compaction_targets_do_not_consume_another_directory_entry()
1340    -> Result<(), Box<dyn Error>> {
1341        let temporary = TestDirectory::new("compaction-existing-targets")?;
1342        let root = temporary.path().join("data");
1343        let mut opened = StorageEngine::open(&root)?;
1344        opened
1345            .storage
1346            .write(Uuid::now_v7(), &[Mutation::put(b"key", b"value")])?;
1347        let snapshot = opened.storage.snapshot()?;
1348        let base_digest = snapshot
1349            .checkpoint_digest
1350            .ok_or("snapshot checkpoint digest is absent")?;
1351        let (prepared, recovery) = DurableLog::open_file_at(
1352            opened.storage.directory.log_path(2),
1353            snapshot.checkpoint_sequence,
1354            base_digest,
1355        )?;
1356        assert_eq!(recovery.valid_bytes, 0);
1357        drop(prepared);
1358        drop(opened);
1359
1360        let limits = StorageLimits {
1361            recovery: crate::RecoveryLimits {
1362                max_directory_entries: 2,
1363                ..crate::RecoveryLimits::default()
1364            },
1365            ..StorageLimits::default()
1366        };
1367        let mut reopened = StorageEngine::open_with_limits(&root, limits.clone())?;
1368        assert!(matches!(
1369            reopened.storage.compact()?,
1370            CompactionOutcome::Compacted(_)
1371        ));
1372        drop(reopened);
1373
1374        let reopened = StorageEngine::open_with_limits(&root, limits)?;
1375        assert_eq!(reopened.storage.get(b"key")?, Some(b"value".to_vec()));
1376        assert_eq!(reopened.storage.directory.manifest().generation, 2);
1377        Ok(())
1378    }
1379
1380    #[test]
1381    fn orphan_prepared_segment_is_ignored_until_manifest_commit() -> Result<(), Box<dyn Error>> {
1382        let temporary = TestDirectory::new("storage-compaction-orphan")?;
1383        let root = temporary.path().join("data");
1384        let mut opened = StorageEngine::open(&root)?;
1385        opened
1386            .storage
1387            .write(Uuid::now_v7(), &[Mutation::put(b"key", b"value")])?;
1388        let snapshot = opened.storage.snapshot()?;
1389        let base_digest = snapshot
1390            .checkpoint_digest
1391            .ok_or("snapshot checkpoint digest is absent")?;
1392        let orphan_path = opened.storage.directory.log_path(2);
1393        let (orphan, recovery) =
1394            DurableLog::open_file_at(&orphan_path, snapshot.checkpoint_sequence, base_digest)?;
1395        assert_eq!(recovery.valid_bytes, 0);
1396        drop(orphan);
1397        drop(opened);
1398
1399        let mut reopened = StorageEngine::open(&root)?;
1400        assert_eq!(reopened.storage.directory.manifest().generation, 1);
1401        assert_eq!(reopened.storage.get(b"key")?, Some(b"value".to_vec()));
1402        assert!(matches!(
1403            reopened.storage.compact()?,
1404            CompactionOutcome::Compacted(_)
1405        ));
1406        assert_eq!(reopened.storage.directory.manifest().generation, 2);
1407        Ok(())
1408    }
1409
1410    #[test]
1411    fn committed_manifest_wins_before_retired_log_cleanup() -> Result<(), Box<dyn Error>> {
1412        let temporary = TestDirectory::new("storage-compaction-committed")?;
1413        let root = temporary.path().join("data");
1414        let mut opened = StorageEngine::open(&root)?;
1415        opened
1416            .storage
1417            .write(Uuid::now_v7(), &[Mutation::put(b"key", b"value")])?;
1418        let snapshot = opened.storage.snapshot()?;
1419        let base_digest = snapshot
1420            .checkpoint_digest
1421            .ok_or("snapshot checkpoint digest is absent")?;
1422        let next = StorageManifest {
1423            generation: 2,
1424            active_segment: 2,
1425            base_sequence: snapshot.checkpoint_sequence,
1426            base_digest,
1427            snapshot_digest: snapshot.snapshot_digest,
1428        };
1429        let (prepared, recovery) = DurableLog::open_file_at(
1430            opened.storage.directory.log_path(2),
1431            next.base_sequence,
1432            next.base_digest,
1433        )?;
1434        assert_eq!(recovery.valid_bytes, 0);
1435        drop(prepared);
1436        opened.storage.directory.commit_manifest(next)?;
1437        let retired_path = opened.storage.directory.log_path(1);
1438        assert!(retired_path.is_file());
1439        drop(opened);
1440
1441        let reopened = StorageEngine::open(&root)?;
1442        assert_eq!(reopened.storage.directory.manifest().generation, 2);
1443        assert_eq!(reopened.storage.get(b"key")?, Some(b"value".to_vec()));
1444        assert!(!retired_path.exists());
1445        Ok(())
1446    }
1447
1448    #[test]
1449    fn uncertain_log_sync_blocks_the_handle_until_recovery() -> Result<(), Box<dyn Error>> {
1450        let temporary = TestDirectory::new("storage-injected-log-sync")?;
1451        let root = temporary.path().join("data");
1452        let mut opened = StorageEngine::open(&root)?;
1453        opened.storage.log.inject_sync_failure();
1454
1455        let result = opened
1456            .storage
1457            .write(Uuid::now_v7(), &[Mutation::put(b"recovered", b"yes")]);
1458        assert!(matches!(
1459            result,
1460            Err(StorageError::Log(crate::LogError::Io(_)))
1461        ));
1462        assert!(matches!(
1463            opened.storage.get(b"recovered"),
1464            Err(StorageError::StaleIndex)
1465        ));
1466        assert!(matches!(
1467            opened.storage.snapshot(),
1468            Err(StorageError::StaleIndex)
1469        ));
1470        assert!(matches!(
1471            opened.storage.compact(),
1472            Err(StorageError::StaleIndex)
1473        ));
1474        drop(opened);
1475
1476        let reopened = StorageEngine::open(&root)?;
1477        assert_eq!(reopened.recovery.replayed_transactions, 1);
1478        assert_eq!(reopened.storage.get(b"recovered")?, Some(b"yes".to_vec()));
1479        Ok(())
1480    }
1481
1482    #[test]
1483    fn uncertain_manifest_commit_blocks_every_operation_until_reopen() -> Result<(), Box<dyn Error>>
1484    {
1485        let temporary = TestDirectory::new("storage-injected-manifest-commit")?;
1486        let root = temporary.path().join("data");
1487        let mut opened = StorageEngine::open(&root)?;
1488        opened
1489            .storage
1490            .write(Uuid::now_v7(), &[Mutation::put(b"durable", b"yes")])?;
1491        opened
1492            .storage
1493            .directory
1494            .inject_manifest_commit_failure_after_write();
1495
1496        assert!(matches!(
1497            opened.storage.compact(),
1498            Err(StorageError::DataDirectory(crate::DataDirectoryError::Io {
1499                action: "complete injected manifest commit",
1500                ..
1501            }))
1502        ));
1503        assert!(matches!(
1504            opened.storage.get(b"durable"),
1505            Err(StorageError::StaleIndex)
1506        ));
1507        assert!(matches!(
1508            opened
1509                .storage
1510                .write(Uuid::now_v7(), &[Mutation::put(b"blocked", b"yes")]),
1511            Err(StorageError::StaleIndex)
1512        ));
1513        assert!(matches!(
1514            opened.storage.snapshot(),
1515            Err(StorageError::StaleIndex)
1516        ));
1517        assert!(matches!(
1518            opened.storage.compact(),
1519            Err(StorageError::StaleIndex)
1520        ));
1521        drop(opened);
1522
1523        let mut reopened = StorageEngine::open(&root)?;
1524        assert_eq!(reopened.storage.directory.manifest().generation, 2);
1525        assert_eq!(reopened.storage.get(b"durable")?, Some(b"yes".to_vec()));
1526        reopened
1527            .storage
1528            .write(Uuid::now_v7(), &[Mutation::put(b"after", b"reopen")])?;
1529        assert_eq!(reopened.storage.get(b"after")?, Some(b"reopen".to_vec()));
1530        Ok(())
1531    }
1532
1533    #[test]
1534    fn post_commit_index_failure_recovers_from_the_log() -> Result<(), Box<dyn Error>> {
1535        let temporary = TestDirectory::new("storage-injected-index")?;
1536        let root = temporary.path().join("data");
1537        let mut opened = StorageEngine::open(&root)?;
1538        opened.storage.index.inject_apply_failure();
1539
1540        let result = opened
1541            .storage
1542            .write(Uuid::now_v7(), &[Mutation::put(b"durable", b"yes")]);
1543        assert!(matches!(
1544            result,
1545            Err(StorageError::CommittedButNotIndexed { .. })
1546        ));
1547        assert!(matches!(
1548            opened.storage.get(b"durable"),
1549            Err(StorageError::StaleIndex)
1550        ));
1551        drop(opened);
1552
1553        let reopened = StorageEngine::open(&root)?;
1554        assert_eq!(reopened.recovery.replayed_transactions, 1);
1555        assert_eq!(reopened.storage.get(b"durable")?, Some(b"yes".to_vec()));
1556        Ok(())
1557    }
1558
1559    #[test]
1560    fn lexical_document_limit_rejects_exact_plus_one_before_append_and_rebuilds()
1561    -> Result<(), Box<dyn Error>> {
1562        let temporary = TestDirectory::new("storage-lexical-document-limit")?;
1563        let root = temporary.path().join("data");
1564        let limits = StorageLimits {
1565            recovery: crate::RecoveryLimits {
1566                max_lexical_documents: 2,
1567                max_lexical_tokens: 100,
1568                ..crate::RecoveryLimits::default()
1569            },
1570            ..StorageLimits::default()
1571        };
1572        let first = lexical_document("one")?;
1573        let second = lexical_document("two")?;
1574        let rejected = lexical_document("three")?;
1575        let mut opened = StorageEngine::open_with_limits(&root, limits.clone())?;
1576        opened.storage.write(
1577            Uuid::now_v7(),
1578            &[
1579                Mutation::put(b"first", first.clone()),
1580                Mutation::put(b"second", second.clone()),
1581            ],
1582        )?;
1583        opened.storage.write(
1584            Uuid::now_v7(),
1585            &[Mutation::define_lexical_index(lexical_definition(
1586                "documents.limit",
1587            )?)],
1588        )?;
1589
1590        let active_log = opened.storage.directory.active_log_path();
1591        let accepted_log_bytes = fs::metadata(&active_log)?.len();
1592        let result = opened
1593            .storage
1594            .write(Uuid::now_v7(), &[Mutation::put(b"rejected", rejected)]);
1595        assert!(matches!(
1596            result,
1597            Err(StorageError::Index { source })
1598                if matches!(
1599                    source.as_ref(),
1600                    crate::MaterializedIndexError::Lexical(
1601                        LexicalError::DocumentBudgetExceeded { maximum: 2 }
1602                    )
1603                )
1604        ));
1605        assert_eq!(opened.storage.get(b"rejected")?, None);
1606        assert_eq!(opened.storage.get(b"first")?, Some(first));
1607        assert_eq!(opened.storage.get(b"second")?, Some(second));
1608        assert_eq!(fs::metadata(&active_log)?.len(), accepted_log_bytes);
1609        drop(opened);
1610
1611        fs::remove_file(root.join("indexes/primary.redb"))?;
1612        let rebuilt = StorageEngine::open_with_limits(&root, limits)?;
1613        let rebuilt_definition = lexical_definition("documents.limit")?;
1614        assert_eq!(rebuilt.storage.get(b"rejected")?, None);
1615        assert_eq!(
1616            rebuilt
1617                .storage
1618                .lexical_corpus(&rebuilt_definition, &[], 10, Duration::from_secs(1))?
1619                .document_count,
1620            2
1621        );
1622        Ok(())
1623    }
1624
1625    #[test]
1626    fn lexical_token_limit_rejects_exact_plus_one_before_append_and_rebuilds()
1627    -> Result<(), Box<dyn Error>> {
1628        let temporary = TestDirectory::new("storage-lexical-token-limit")?;
1629        let root = temporary.path().join("data");
1630        let limits = StorageLimits {
1631            recovery: crate::RecoveryLimits {
1632                max_lexical_documents: 10,
1633                max_lexical_tokens: 2,
1634                ..crate::RecoveryLimits::default()
1635            },
1636            ..StorageLimits::default()
1637        };
1638        let accepted = lexical_document("one two")?;
1639        let rejected = lexical_document("one two three")?;
1640        let mut opened = StorageEngine::open_with_limits(&root, limits.clone())?;
1641        opened.storage.write(
1642            Uuid::now_v7(),
1643            &[Mutation::put(b"document", accepted.clone())],
1644        )?;
1645        opened.storage.write(
1646            Uuid::now_v7(),
1647            &[Mutation::define_lexical_index(lexical_definition(
1648                "tokens.limit",
1649            )?)],
1650        )?;
1651
1652        let active_log = opened.storage.directory.active_log_path();
1653        let accepted_log_bytes = fs::metadata(&active_log)?.len();
1654        let result = opened
1655            .storage
1656            .write(Uuid::now_v7(), &[Mutation::put(b"document", rejected)]);
1657        assert!(matches!(
1658            result,
1659            Err(StorageError::Index { source })
1660                if matches!(
1661                    source.as_ref(),
1662                    crate::MaterializedIndexError::Lexical(
1663                        LexicalError::TokenBudgetExceeded { maximum: 2 }
1664                    )
1665            )
1666        ));
1667        assert_eq!(opened.storage.get(b"document")?, Some(accepted.clone()));
1668        assert_eq!(fs::metadata(&active_log)?.len(), accepted_log_bytes);
1669        drop(opened);
1670
1671        fs::remove_file(root.join("indexes/primary.redb"))?;
1672        let rebuilt = StorageEngine::open_with_limits(&root, limits)?;
1673        let rebuilt_definition = lexical_definition("tokens.limit")?;
1674        assert_eq!(rebuilt.storage.get(b"document")?, Some(accepted));
1675        assert_eq!(
1676            rebuilt
1677                .storage
1678                .lexical_corpus(&rebuilt_definition, &[], 10, Duration::from_secs(1))?
1679                .token_count,
1680            2
1681        );
1682        Ok(())
1683    }
1684
1685    #[test]
1686    fn lexical_define_over_limit_is_read_only_before_append() -> Result<(), Box<dyn Error>> {
1687        let temporary = TestDirectory::new("storage-lexical-define-limit")?;
1688        let root = temporary.path().join("data");
1689        let limits = StorageLimits {
1690            recovery: crate::RecoveryLimits {
1691                max_lexical_documents: 1,
1692                max_lexical_tokens: 100,
1693                ..crate::RecoveryLimits::default()
1694            },
1695            ..StorageLimits::default()
1696        };
1697        let mut opened = StorageEngine::open_with_limits(&root, limits)?;
1698        opened.storage.write(
1699            Uuid::now_v7(),
1700            &[
1701                Mutation::put(b"first", lexical_document("one")?),
1702                Mutation::put(b"second", lexical_document("two")?),
1703            ],
1704        )?;
1705
1706        let active_log = opened.storage.directory.active_log_path();
1707        let format_path = root.join("FORMAT");
1708        let accepted_log_bytes = fs::metadata(&active_log)?.len();
1709        let accepted_format = fs::read(&format_path)?;
1710
1711        let definition = lexical_definition("define.limit")?;
1712        let result = opened.storage.write(
1713            Uuid::now_v7(),
1714            &[Mutation::define_lexical_index(definition.clone())],
1715        );
1716        assert!(matches!(
1717            result,
1718            Err(StorageError::Index { source })
1719                if matches!(
1720                    source.as_ref(),
1721                    crate::MaterializedIndexError::Lexical(
1722                        LexicalError::DocumentBudgetExceeded { maximum: 1 }
1723                    )
1724            )
1725        ));
1726        assert_eq!(opened.storage.lexical_index(&definition.name)?, None);
1727        assert_eq!(fs::metadata(&active_log)?.len(), accepted_log_bytes);
1728        assert_eq!(fs::read(&format_path)?, accepted_format);
1729        Ok(())
1730    }
1731
1732    #[test]
1733    fn lexical_preflight_simulates_ordered_batches_and_idempotent_retries()
1734    -> Result<(), Box<dyn Error>> {
1735        let temporary = TestDirectory::new("storage-lexical-ordered-preflight")?;
1736        let root = temporary.path().join("put-before-define");
1737        let limits = StorageLimits {
1738            recovery: crate::RecoveryLimits {
1739                max_lexical_documents: 1,
1740                max_lexical_tokens: 2,
1741                ..crate::RecoveryLimits::default()
1742            },
1743            ..StorageLimits::default()
1744        };
1745        let definition = lexical_definition("ordered.limit")?;
1746        let first_transaction = Uuid::now_v7();
1747        let first_batch = [
1748            Mutation::put(b"first", lexical_document("one two")?),
1749            Mutation::define_lexical_index(definition.clone()),
1750        ];
1751        let mut opened = StorageEngine::open_with_limits(&root, limits.clone())?;
1752        let first_outcome = opened.storage.write(first_transaction, &first_batch)?;
1753        assert!(matches!(first_outcome, AppendOutcome::Committed(_)));
1754
1755        opened.storage.write(
1756            Uuid::now_v7(),
1757            &[
1758                Mutation::delete(b"first"),
1759                Mutation::put(b"second", lexical_document("three")?),
1760            ],
1761        )?;
1762        assert_eq!(opened.storage.get(b"first")?, None);
1763        assert!(opened.storage.get(b"second")?.is_some());
1764        assert!(matches!(
1765            opened.storage.write(
1766                Uuid::now_v7(),
1767                &[
1768                    Mutation::put(b"third", lexical_document("four")?),
1769                    Mutation::delete(b"second"),
1770                ],
1771            ),
1772            Err(StorageError::Index { source })
1773                if matches!(
1774                    source.as_ref(),
1775                    crate::MaterializedIndexError::Lexical(
1776                        LexicalError::DocumentBudgetExceeded { maximum: 1 }
1777                    )
1778                )
1779        ));
1780        assert!(matches!(
1781            opened.storage.write(first_transaction, &first_batch)?,
1782            AppendOutcome::Existing(_)
1783        ));
1784        drop(opened);
1785
1786        let define_first_root = temporary.path().join("define-before-put");
1787        let mut define_first = StorageEngine::open_with_limits(&define_first_root, limits)?;
1788        define_first.storage.write(
1789            Uuid::now_v7(),
1790            &[
1791                Mutation::define_lexical_index(lexical_definition("ordered.second")?),
1792                Mutation::put(b"document", lexical_document("one two")?),
1793            ],
1794        )?;
1795        assert!(define_first.storage.get(b"document")?.is_some());
1796        Ok(())
1797    }
1798
1799    #[test]
1800    fn kv_scan_pages_are_strictly_ordered_and_exclusive() -> Result<(), Box<dyn Error>> {
1801        let temporary = TestDirectory::new("storage-scan-page")?;
1802        let mut opened = StorageEngine::open(temporary.path().join("data"))?;
1803        opened.storage.write(
1804            Uuid::now_v7(),
1805            &[
1806                Mutation::put(b"c", b"three"),
1807                Mutation::put(b"a", b"one"),
1808                Mutation::put(b"b", b"two"),
1809            ],
1810        )?;
1811
1812        let first = opened.storage.scan_page(None, 2)?;
1813        assert_eq!(
1814            first
1815                .entries
1816                .iter()
1817                .map(|entry| entry.key.as_slice())
1818                .collect::<Vec<_>>(),
1819            [b"a".as_slice(), b"b".as_slice()]
1820        );
1821        assert_eq!(first.next_after, Some(b"b".to_vec()));
1822
1823        let second = opened.storage.scan_page(first.next_after.as_deref(), 2)?;
1824        assert_eq!(second.entries[0].key, b"c");
1825        assert_eq!(second.next_after, None);
1826
1827        let total_bytes = u64::try_from(
1828            b"a".len() + b"one".len() + b"b".len() + b"two".len() + b"c".len() + b"three".len(),
1829        )?;
1830        assert_eq!(
1831            opened
1832                .storage
1833                .scan_page_with_byte_limit(None, 3, total_bytes)?
1834                .entries
1835                .len(),
1836            3
1837        );
1838        assert!(matches!(
1839            opened
1840                .storage
1841                .scan_page_with_byte_limit(None, 3, total_bytes - 1),
1842            Err(ScanPageError::ByteBudgetExceeded { maximum })
1843                if maximum == total_bytes - 1
1844        ));
1845        let first_bytes = u64::try_from(b"a".len() + b"one".len())?;
1846        let first_only = opened
1847            .storage
1848            .scan_page_with_byte_limit(None, 1, first_bytes)?;
1849        assert_eq!(first_only.entries.len(), 1);
1850        assert_eq!(first_only.next_after, Some(b"a".to_vec()));
1851        Ok(())
1852    }
1853}