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_with_timeout, storage_limit_from_io,
878        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 = load_snapshot(&created.path, &SnapshotReadLimits::default())?;
985        assert_eq!(witness.info, created);
986        assert_eq!(witness.entries.len(), 2);
987        assert_eq!(witness.entries[0].key, b"alpha");
988        assert_eq!(witness.entries[0].value, b"first");
989        assert!(matches!(
990            load_snapshot_with_timeout(
991                &created.path,
992                &SnapshotReadLimits::default(),
993                Duration::ZERO
994            ),
995            Err(source) if source.is_timeout()
996        ));
997        assert!(matches!(
998            load_snapshot(
999                &created.path,
1000                &SnapshotReadLimits {
1001                    entries: 1,
1002                    ..SnapshotReadLimits::default()
1003                }
1004            ),
1005            Err(SnapshotError::EntryLimitExceeded {
1006                actual: 2,
1007                maximum: 1
1008            })
1009        ));
1010
1011        let corrupted_path = temporary.path().join("corrupted.hysnap");
1012        fs::copy(&created.path, &corrupted_path)?;
1013        let mut corrupted = OpenOptions::new()
1014            .read(true)
1015            .write(true)
1016            .open(&corrupted_path)?;
1017        corrupted.seek(SeekFrom::End(-1))?;
1018        corrupted.write_all(&[0xff])?;
1019        corrupted.sync_all()?;
1020        drop(corrupted);
1021
1022        assert!(matches!(
1023            load_snapshot(
1024                &corrupted_path,
1025                &SnapshotReadLimits {
1026                    entries: 1,
1027                    ..SnapshotReadLimits::default()
1028                }
1029            ),
1030            Err(SnapshotError::EntryLimitExceeded {
1031                actual: 2,
1032                maximum: 1
1033            })
1034        ));
1035        assert!(matches!(
1036            load_snapshot(
1037                &corrupted_path,
1038                &SnapshotReadLimits {
1039                    decoded_bytes: 0,
1040                    ..SnapshotReadLimits::default()
1041                }
1042            ),
1043            Err(SnapshotError::DecodedBytesLimitExceeded { maximum: 0 })
1044        ));
1045        assert!(matches!(
1046            verify_snapshot(&corrupted_path),
1047            Err(SnapshotError::Invalid {
1048                reason: "CRC32C mismatch"
1049            })
1050        ));
1051        Ok(())
1052    }
1053
1054    #[test]
1055    fn empty_storage_has_a_canonical_empty_snapshot() -> Result<(), Box<dyn Error>> {
1056        let temporary = TestDirectory::new("storage-empty-snapshot")?;
1057        let opened = StorageEngine::open(temporary.path().join("data"))?;
1058
1059        let snapshot = opened.storage.snapshot()?;
1060        assert_eq!(snapshot.checkpoint_sequence, 0);
1061        assert_eq!(snapshot.checkpoint_digest, None);
1062        assert_eq!(snapshot.entry_count, 0);
1063        assert_eq!(snapshot.vector_space_count, 0);
1064        assert_eq!(snapshot.vector_count, 0);
1065        assert_eq!(snapshot.lexical_index_count, 0);
1066        assert_eq!(snapshot.receipt_count, 0);
1067        assert_eq!(snapshot.file_bytes, 136);
1068        assert_eq!(verify_snapshot(&snapshot.path)?, snapshot);
1069        Ok(())
1070    }
1071
1072    #[test]
1073    fn snapshot_rebuilds_kv_and_idempotency_state() -> Result<(), Box<dyn Error>> {
1074        let temporary = TestDirectory::new("storage-snapshot-restore")?;
1075        let root = temporary.path().join("data");
1076        let transaction_id = Uuid::now_v7();
1077        let mut opened = StorageEngine::open(&root)?;
1078        let outcome = opened.storage.write(
1079            transaction_id,
1080            &[
1081                Mutation::put(b"alpha", b"one"),
1082                Mutation::put(b"beta", b"two"),
1083            ],
1084        )?;
1085        let AppendOutcome::Committed(receipt) = outcome else {
1086            return Err("new transaction was not committed".into());
1087        };
1088        let snapshot = opened.storage.snapshot()?;
1089        drop(opened);
1090
1091        let restored_path = root.join("tmp/restored.redb");
1092        assert_eq!(
1093            MaterializedIndex::restore_from_snapshot(&restored_path, &snapshot.path)?,
1094            snapshot
1095        );
1096        let restored = MaterializedIndex::open(&restored_path)?;
1097        assert_eq!(restored.get(b"alpha")?, Some(b"one".to_vec()));
1098        assert_eq!(restored.get(b"beta")?, Some(b"two".to_vec()));
1099        assert_eq!(restored.receipt(transaction_id)?, Some(receipt));
1100        Ok(())
1101    }
1102
1103    #[test]
1104    fn compaction_retires_history_and_snapshot_rebuilds_the_index() -> Result<(), Box<dyn Error>> {
1105        let temporary = TestDirectory::new("storage-compaction")?;
1106        let root = temporary.path().join("data");
1107        let first_id = Uuid::now_v7();
1108        let mut opened = StorageEngine::open(&root)?;
1109        let first = opened
1110            .storage
1111            .write(first_id, &[Mutation::put(b"before", b"one")])?;
1112        let AppendOutcome::Committed(first_receipt) = first else {
1113            return Err("first transaction was not committed".into());
1114        };
1115
1116        let compacted = opened.storage.compact()?;
1117        let CompactionOutcome::Compacted(report) = compacted else {
1118            return Err("committed history was not compacted".into());
1119        };
1120        assert_eq!(report.generation, 2);
1121        assert!(report.retired_segment_removed);
1122        assert!(!report.retired_segment.exists());
1123        assert!(root.join("log/00000000000000000002.hylog").is_file());
1124        assert!(matches!(
1125            opened.storage.compact()?,
1126            CompactionOutcome::NoChanges { .. }
1127        ));
1128
1129        assert_eq!(
1130            opened
1131                .storage
1132                .write(first_id, &[Mutation::put(b"before", b"one")])?,
1133            AppendOutcome::Existing(first_receipt)
1134        );
1135        let second_id = Uuid::now_v7();
1136        let second = opened
1137            .storage
1138            .write(second_id, &[Mutation::put(b"after", b"two")])?;
1139        let AppendOutcome::Committed(second_receipt) = second else {
1140            return Err("second transaction was not committed".into());
1141        };
1142        assert_eq!(
1143            second_receipt.commit_sequence,
1144            first_receipt.commit_sequence + 3
1145        );
1146        drop(opened);
1147
1148        fs::remove_file(root.join("indexes/primary.redb"))?;
1149        let mut rebuilt = StorageEngine::open(&root)?;
1150        assert_eq!(rebuilt.recovery.replayed_transactions, 1);
1151        assert_eq!(rebuilt.storage.get(b"before")?, Some(b"one".to_vec()));
1152        assert_eq!(rebuilt.storage.get(b"after")?, Some(b"two".to_vec()));
1153        assert_eq!(
1154            rebuilt
1155                .storage
1156                .write(first_id, &[Mutation::put(b"before", b"one")])?,
1157            AppendOutcome::Existing(first_receipt)
1158        );
1159        Ok(())
1160    }
1161
1162    #[test]
1163    fn open_and_compaction_limits_fail_without_advancing_durable_state()
1164    -> Result<(), Box<dyn Error>> {
1165        let temporary = TestDirectory::new("bounded-open-compaction")?;
1166        let root = temporary.path().join("data");
1167        let mut opened = StorageEngine::open(&root)?;
1168        opened.storage.write(
1169            Uuid::now_v7(),
1170            &[Mutation::put(b"a", b"one"), Mutation::put(b"b", b"two")],
1171        )?;
1172        let active_log = opened.storage.directory.active_log_path();
1173        let log_bytes = fs::metadata(&active_log)?.len();
1174        let generation = opened.storage.directory.manifest().generation;
1175
1176        let too_few_records = MaintenanceLimits {
1177            snapshot: SnapshotReadLimits {
1178                entries: 1,
1179                ..SnapshotReadLimits::default()
1180            },
1181            ..MaintenanceLimits::default()
1182        };
1183        assert!(matches!(
1184            opened.storage.compact_with_limits(&too_few_records),
1185            Err(StorageError::Snapshot { source })
1186                if matches!(
1187                    source.as_ref(),
1188                    SnapshotError::EntryLimitExceeded {
1189                        actual: 2,
1190                        maximum: 1
1191                    }
1192                )
1193        ));
1194        assert_eq!(opened.storage.directory.manifest().generation, generation);
1195        assert!(active_log.is_file());
1196        drop(opened);
1197
1198        let limits = StorageLimits {
1199            recovery: crate::RecoveryLimits {
1200                max_log_file_bytes: log_bytes - 1,
1201                ..crate::RecoveryLimits::default()
1202            },
1203            ..StorageLimits::default()
1204        };
1205        assert!(matches!(
1206            StorageEngine::open_with_limits(&root, limits),
1207            Err(StorageError::Log(crate::LogError::Io(source)))
1208                if matches!(
1209                    storage_limit_from_io(&source),
1210                    Some(StorageLimitError::LogFileBytesExceeded { actual, maximum })
1211                        if *actual == log_bytes && *maximum == log_bytes - 1
1212                )
1213        ));
1214        assert_eq!(fs::metadata(&active_log)?.len(), log_bytes);
1215        Ok(())
1216    }
1217
1218    #[test]
1219    fn compaction_snapshot_policy_is_reopenable_at_the_exact_recovery_limit()
1220    -> Result<(), Box<dyn Error>> {
1221        let temporary = TestDirectory::new("compaction-recovery-snapshot-intersection")?;
1222
1223        let exact_root = temporary.path().join("exact");
1224        let exact_limits = StorageLimits {
1225            recovery: crate::RecoveryLimits {
1226                snapshot: SnapshotReadLimits {
1227                    entries: 2,
1228                    ..SnapshotReadLimits::default()
1229                },
1230                ..crate::RecoveryLimits::default()
1231            },
1232            maintenance: MaintenanceLimits {
1233                snapshot: SnapshotReadLimits {
1234                    entries: 10,
1235                    ..SnapshotReadLimits::default()
1236                },
1237                ..MaintenanceLimits::default()
1238            },
1239        };
1240        let mut exact = StorageEngine::open_with_limits(&exact_root, exact_limits.clone())?;
1241        exact.storage.write(
1242            Uuid::now_v7(),
1243            &[Mutation::put(b"a", b"one"), Mutation::put(b"b", b"two")],
1244        )?;
1245        assert!(matches!(
1246            exact.storage.compact()?,
1247            CompactionOutcome::Compacted(_)
1248        ));
1249        drop(exact);
1250        let reopened = StorageEngine::open_with_limits(&exact_root, exact_limits)?;
1251        assert_eq!(reopened.storage.get(b"a")?, Some(b"one".to_vec()));
1252        assert_eq!(reopened.storage.get(b"b")?, Some(b"two".to_vec()));
1253        drop(reopened);
1254
1255        let rejected_root = temporary.path().join("exact-plus-one");
1256        let rejected_limits = StorageLimits {
1257            recovery: crate::RecoveryLimits {
1258                snapshot: SnapshotReadLimits {
1259                    entries: 1,
1260                    ..SnapshotReadLimits::default()
1261                },
1262                ..crate::RecoveryLimits::default()
1263            },
1264            maintenance: MaintenanceLimits {
1265                snapshot: SnapshotReadLimits {
1266                    entries: 2,
1267                    ..SnapshotReadLimits::default()
1268                },
1269                ..MaintenanceLimits::default()
1270            },
1271        };
1272        let mut rejected =
1273            StorageEngine::open_with_limits(&rejected_root, rejected_limits.clone())?;
1274        rejected.storage.write(
1275            Uuid::now_v7(),
1276            &[Mutation::put(b"a", b"one"), Mutation::put(b"b", b"two")],
1277        )?;
1278        assert!(matches!(
1279            rejected.storage.compact(),
1280            Err(StorageError::Snapshot { source })
1281                if matches!(
1282                    source.as_ref(),
1283                    SnapshotError::EntryLimitExceeded {
1284                        actual: 2,
1285                        maximum: 1
1286                    }
1287                )
1288        ));
1289        assert_eq!(rejected.storage.directory.manifest().generation, 1);
1290        drop(rejected);
1291        let reopened = StorageEngine::open_with_limits(&rejected_root, rejected_limits)?;
1292        assert_eq!(reopened.storage.get(b"a")?, Some(b"one".to_vec()));
1293        assert_eq!(reopened.storage.get(b"b")?, Some(b"two".to_vec()));
1294        Ok(())
1295    }
1296
1297    #[test]
1298    fn compaction_reserves_directory_entries_before_creating_any_artifact()
1299    -> Result<(), Box<dyn Error>> {
1300        let temporary = TestDirectory::new("compaction-directory-reservation")?;
1301        let root = temporary.path().join("data");
1302        let limits = StorageLimits {
1303            recovery: crate::RecoveryLimits {
1304                max_directory_entries: 1,
1305                ..crate::RecoveryLimits::default()
1306            },
1307            ..StorageLimits::default()
1308        };
1309        let mut opened = StorageEngine::open_with_limits(&root, limits.clone())?;
1310        opened
1311            .storage
1312            .write(Uuid::now_v7(), &[Mutation::put(b"key", b"value")])?;
1313
1314        assert!(matches!(
1315            opened.storage.compact(),
1316            Err(StorageError::DataDirectory(
1317                crate::DataDirectoryError::Manifest(ManifestError::Io(source))
1318            ))
1319                if matches!(
1320                    storage_limit_from_io(&source),
1321                    Some(StorageLimitError::DirectoryEntriesExceeded { maximum: 1 })
1322                )
1323        ));
1324        assert_eq!(fs::read_dir(root.join("snapshots"))?.count(), 0);
1325        assert_eq!(fs::read_dir(root.join("log"))?.count(), 1);
1326        assert_eq!(fs::read_dir(root.join("manifest"))?.count(), 1);
1327        assert_eq!(fs::read_dir(root.join("tmp"))?.count(), 0);
1328        assert_eq!(opened.storage.directory.manifest().generation, 1);
1329        drop(opened);
1330
1331        let reopened = StorageEngine::open_with_limits(&root, limits)?;
1332        assert_eq!(reopened.storage.get(b"key")?, Some(b"value".to_vec()));
1333        Ok(())
1334    }
1335
1336    #[test]
1337    fn existing_compaction_targets_do_not_consume_another_directory_entry()
1338    -> Result<(), Box<dyn Error>> {
1339        let temporary = TestDirectory::new("compaction-existing-targets")?;
1340        let root = temporary.path().join("data");
1341        let mut opened = StorageEngine::open(&root)?;
1342        opened
1343            .storage
1344            .write(Uuid::now_v7(), &[Mutation::put(b"key", b"value")])?;
1345        let snapshot = opened.storage.snapshot()?;
1346        let base_digest = snapshot
1347            .checkpoint_digest
1348            .ok_or("snapshot checkpoint digest is absent")?;
1349        let (prepared, recovery) = DurableLog::open_file_at(
1350            opened.storage.directory.log_path(2),
1351            snapshot.checkpoint_sequence,
1352            base_digest,
1353        )?;
1354        assert_eq!(recovery.valid_bytes, 0);
1355        drop(prepared);
1356        drop(opened);
1357
1358        let limits = StorageLimits {
1359            recovery: crate::RecoveryLimits {
1360                max_directory_entries: 2,
1361                ..crate::RecoveryLimits::default()
1362            },
1363            ..StorageLimits::default()
1364        };
1365        let mut reopened = StorageEngine::open_with_limits(&root, limits.clone())?;
1366        assert!(matches!(
1367            reopened.storage.compact()?,
1368            CompactionOutcome::Compacted(_)
1369        ));
1370        drop(reopened);
1371
1372        let reopened = StorageEngine::open_with_limits(&root, limits)?;
1373        assert_eq!(reopened.storage.get(b"key")?, Some(b"value".to_vec()));
1374        assert_eq!(reopened.storage.directory.manifest().generation, 2);
1375        Ok(())
1376    }
1377
1378    #[test]
1379    fn orphan_prepared_segment_is_ignored_until_manifest_commit() -> Result<(), Box<dyn Error>> {
1380        let temporary = TestDirectory::new("storage-compaction-orphan")?;
1381        let root = temporary.path().join("data");
1382        let mut opened = StorageEngine::open(&root)?;
1383        opened
1384            .storage
1385            .write(Uuid::now_v7(), &[Mutation::put(b"key", b"value")])?;
1386        let snapshot = opened.storage.snapshot()?;
1387        let base_digest = snapshot
1388            .checkpoint_digest
1389            .ok_or("snapshot checkpoint digest is absent")?;
1390        let orphan_path = opened.storage.directory.log_path(2);
1391        let (orphan, recovery) =
1392            DurableLog::open_file_at(&orphan_path, snapshot.checkpoint_sequence, base_digest)?;
1393        assert_eq!(recovery.valid_bytes, 0);
1394        drop(orphan);
1395        drop(opened);
1396
1397        let mut reopened = StorageEngine::open(&root)?;
1398        assert_eq!(reopened.storage.directory.manifest().generation, 1);
1399        assert_eq!(reopened.storage.get(b"key")?, Some(b"value".to_vec()));
1400        assert!(matches!(
1401            reopened.storage.compact()?,
1402            CompactionOutcome::Compacted(_)
1403        ));
1404        assert_eq!(reopened.storage.directory.manifest().generation, 2);
1405        Ok(())
1406    }
1407
1408    #[test]
1409    fn committed_manifest_wins_before_retired_log_cleanup() -> Result<(), Box<dyn Error>> {
1410        let temporary = TestDirectory::new("storage-compaction-committed")?;
1411        let root = temporary.path().join("data");
1412        let mut opened = StorageEngine::open(&root)?;
1413        opened
1414            .storage
1415            .write(Uuid::now_v7(), &[Mutation::put(b"key", b"value")])?;
1416        let snapshot = opened.storage.snapshot()?;
1417        let base_digest = snapshot
1418            .checkpoint_digest
1419            .ok_or("snapshot checkpoint digest is absent")?;
1420        let next = StorageManifest {
1421            generation: 2,
1422            active_segment: 2,
1423            base_sequence: snapshot.checkpoint_sequence,
1424            base_digest,
1425            snapshot_digest: snapshot.snapshot_digest,
1426        };
1427        let (prepared, recovery) = DurableLog::open_file_at(
1428            opened.storage.directory.log_path(2),
1429            next.base_sequence,
1430            next.base_digest,
1431        )?;
1432        assert_eq!(recovery.valid_bytes, 0);
1433        drop(prepared);
1434        opened.storage.directory.commit_manifest(next)?;
1435        let retired_path = opened.storage.directory.log_path(1);
1436        assert!(retired_path.is_file());
1437        drop(opened);
1438
1439        let reopened = StorageEngine::open(&root)?;
1440        assert_eq!(reopened.storage.directory.manifest().generation, 2);
1441        assert_eq!(reopened.storage.get(b"key")?, Some(b"value".to_vec()));
1442        assert!(!retired_path.exists());
1443        Ok(())
1444    }
1445
1446    #[test]
1447    fn uncertain_log_sync_blocks_the_handle_until_recovery() -> Result<(), Box<dyn Error>> {
1448        let temporary = TestDirectory::new("storage-injected-log-sync")?;
1449        let root = temporary.path().join("data");
1450        let mut opened = StorageEngine::open(&root)?;
1451        opened.storage.log.inject_sync_failure();
1452
1453        let result = opened
1454            .storage
1455            .write(Uuid::now_v7(), &[Mutation::put(b"recovered", b"yes")]);
1456        assert!(matches!(
1457            result,
1458            Err(StorageError::Log(crate::LogError::Io(_)))
1459        ));
1460        assert!(matches!(
1461            opened.storage.get(b"recovered"),
1462            Err(StorageError::StaleIndex)
1463        ));
1464        assert!(matches!(
1465            opened.storage.snapshot(),
1466            Err(StorageError::StaleIndex)
1467        ));
1468        assert!(matches!(
1469            opened.storage.compact(),
1470            Err(StorageError::StaleIndex)
1471        ));
1472        drop(opened);
1473
1474        let reopened = StorageEngine::open(&root)?;
1475        assert_eq!(reopened.recovery.replayed_transactions, 1);
1476        assert_eq!(reopened.storage.get(b"recovered")?, Some(b"yes".to_vec()));
1477        Ok(())
1478    }
1479
1480    #[test]
1481    fn uncertain_manifest_commit_blocks_every_operation_until_reopen() -> Result<(), Box<dyn Error>>
1482    {
1483        let temporary = TestDirectory::new("storage-injected-manifest-commit")?;
1484        let root = temporary.path().join("data");
1485        let mut opened = StorageEngine::open(&root)?;
1486        opened
1487            .storage
1488            .write(Uuid::now_v7(), &[Mutation::put(b"durable", b"yes")])?;
1489        opened
1490            .storage
1491            .directory
1492            .inject_manifest_commit_failure_after_write();
1493
1494        assert!(matches!(
1495            opened.storage.compact(),
1496            Err(StorageError::DataDirectory(crate::DataDirectoryError::Io {
1497                action: "complete injected manifest commit",
1498                ..
1499            }))
1500        ));
1501        assert!(matches!(
1502            opened.storage.get(b"durable"),
1503            Err(StorageError::StaleIndex)
1504        ));
1505        assert!(matches!(
1506            opened
1507                .storage
1508                .write(Uuid::now_v7(), &[Mutation::put(b"blocked", b"yes")]),
1509            Err(StorageError::StaleIndex)
1510        ));
1511        assert!(matches!(
1512            opened.storage.snapshot(),
1513            Err(StorageError::StaleIndex)
1514        ));
1515        assert!(matches!(
1516            opened.storage.compact(),
1517            Err(StorageError::StaleIndex)
1518        ));
1519        drop(opened);
1520
1521        let mut reopened = StorageEngine::open(&root)?;
1522        assert_eq!(reopened.storage.directory.manifest().generation, 2);
1523        assert_eq!(reopened.storage.get(b"durable")?, Some(b"yes".to_vec()));
1524        reopened
1525            .storage
1526            .write(Uuid::now_v7(), &[Mutation::put(b"after", b"reopen")])?;
1527        assert_eq!(reopened.storage.get(b"after")?, Some(b"reopen".to_vec()));
1528        Ok(())
1529    }
1530
1531    #[test]
1532    fn post_commit_index_failure_recovers_from_the_log() -> Result<(), Box<dyn Error>> {
1533        let temporary = TestDirectory::new("storage-injected-index")?;
1534        let root = temporary.path().join("data");
1535        let mut opened = StorageEngine::open(&root)?;
1536        opened.storage.index.inject_apply_failure();
1537
1538        let result = opened
1539            .storage
1540            .write(Uuid::now_v7(), &[Mutation::put(b"durable", b"yes")]);
1541        assert!(matches!(
1542            result,
1543            Err(StorageError::CommittedButNotIndexed { .. })
1544        ));
1545        assert!(matches!(
1546            opened.storage.get(b"durable"),
1547            Err(StorageError::StaleIndex)
1548        ));
1549        drop(opened);
1550
1551        let reopened = StorageEngine::open(&root)?;
1552        assert_eq!(reopened.recovery.replayed_transactions, 1);
1553        assert_eq!(reopened.storage.get(b"durable")?, Some(b"yes".to_vec()));
1554        Ok(())
1555    }
1556
1557    #[test]
1558    fn lexical_document_limit_rejects_exact_plus_one_before_append_and_rebuilds()
1559    -> Result<(), Box<dyn Error>> {
1560        let temporary = TestDirectory::new("storage-lexical-document-limit")?;
1561        let root = temporary.path().join("data");
1562        let limits = StorageLimits {
1563            recovery: crate::RecoveryLimits {
1564                max_lexical_documents: 2,
1565                max_lexical_tokens: 100,
1566                ..crate::RecoveryLimits::default()
1567            },
1568            ..StorageLimits::default()
1569        };
1570        let first = lexical_document("one")?;
1571        let second = lexical_document("two")?;
1572        let rejected = lexical_document("three")?;
1573        let mut opened = StorageEngine::open_with_limits(&root, limits.clone())?;
1574        opened.storage.write(
1575            Uuid::now_v7(),
1576            &[
1577                Mutation::put(b"first", first.clone()),
1578                Mutation::put(b"second", second.clone()),
1579            ],
1580        )?;
1581        opened.storage.write(
1582            Uuid::now_v7(),
1583            &[Mutation::define_lexical_index(lexical_definition(
1584                "documents.limit",
1585            )?)],
1586        )?;
1587
1588        let active_log = opened.storage.directory.active_log_path();
1589        let accepted_log_bytes = fs::metadata(&active_log)?.len();
1590        let result = opened
1591            .storage
1592            .write(Uuid::now_v7(), &[Mutation::put(b"rejected", rejected)]);
1593        assert!(matches!(
1594            result,
1595            Err(StorageError::Index { source })
1596                if matches!(
1597                    source.as_ref(),
1598                    crate::MaterializedIndexError::Lexical(
1599                        LexicalError::DocumentBudgetExceeded { maximum: 2 }
1600                    )
1601                )
1602        ));
1603        assert_eq!(opened.storage.get(b"rejected")?, None);
1604        assert_eq!(opened.storage.get(b"first")?, Some(first));
1605        assert_eq!(opened.storage.get(b"second")?, Some(second));
1606        assert_eq!(fs::metadata(&active_log)?.len(), accepted_log_bytes);
1607        drop(opened);
1608
1609        fs::remove_file(root.join("indexes/primary.redb"))?;
1610        let rebuilt = StorageEngine::open_with_limits(&root, limits)?;
1611        let rebuilt_definition = lexical_definition("documents.limit")?;
1612        assert_eq!(rebuilt.storage.get(b"rejected")?, None);
1613        assert_eq!(
1614            rebuilt
1615                .storage
1616                .lexical_corpus(&rebuilt_definition, &[], 10, Duration::from_secs(1))?
1617                .document_count,
1618            2
1619        );
1620        Ok(())
1621    }
1622
1623    #[test]
1624    fn lexical_token_limit_rejects_exact_plus_one_before_append_and_rebuilds()
1625    -> Result<(), Box<dyn Error>> {
1626        let temporary = TestDirectory::new("storage-lexical-token-limit")?;
1627        let root = temporary.path().join("data");
1628        let limits = StorageLimits {
1629            recovery: crate::RecoveryLimits {
1630                max_lexical_documents: 10,
1631                max_lexical_tokens: 2,
1632                ..crate::RecoveryLimits::default()
1633            },
1634            ..StorageLimits::default()
1635        };
1636        let accepted = lexical_document("one two")?;
1637        let rejected = lexical_document("one two three")?;
1638        let mut opened = StorageEngine::open_with_limits(&root, limits.clone())?;
1639        opened.storage.write(
1640            Uuid::now_v7(),
1641            &[Mutation::put(b"document", accepted.clone())],
1642        )?;
1643        opened.storage.write(
1644            Uuid::now_v7(),
1645            &[Mutation::define_lexical_index(lexical_definition(
1646                "tokens.limit",
1647            )?)],
1648        )?;
1649
1650        let active_log = opened.storage.directory.active_log_path();
1651        let accepted_log_bytes = fs::metadata(&active_log)?.len();
1652        let result = opened
1653            .storage
1654            .write(Uuid::now_v7(), &[Mutation::put(b"document", rejected)]);
1655        assert!(matches!(
1656            result,
1657            Err(StorageError::Index { source })
1658                if matches!(
1659                    source.as_ref(),
1660                    crate::MaterializedIndexError::Lexical(
1661                        LexicalError::TokenBudgetExceeded { maximum: 2 }
1662                    )
1663            )
1664        ));
1665        assert_eq!(opened.storage.get(b"document")?, Some(accepted.clone()));
1666        assert_eq!(fs::metadata(&active_log)?.len(), accepted_log_bytes);
1667        drop(opened);
1668
1669        fs::remove_file(root.join("indexes/primary.redb"))?;
1670        let rebuilt = StorageEngine::open_with_limits(&root, limits)?;
1671        let rebuilt_definition = lexical_definition("tokens.limit")?;
1672        assert_eq!(rebuilt.storage.get(b"document")?, Some(accepted));
1673        assert_eq!(
1674            rebuilt
1675                .storage
1676                .lexical_corpus(&rebuilt_definition, &[], 10, Duration::from_secs(1))?
1677                .token_count,
1678            2
1679        );
1680        Ok(())
1681    }
1682
1683    #[test]
1684    fn lexical_define_over_limit_is_read_only_before_append() -> Result<(), Box<dyn Error>> {
1685        let temporary = TestDirectory::new("storage-lexical-define-limit")?;
1686        let root = temporary.path().join("data");
1687        let limits = StorageLimits {
1688            recovery: crate::RecoveryLimits {
1689                max_lexical_documents: 1,
1690                max_lexical_tokens: 100,
1691                ..crate::RecoveryLimits::default()
1692            },
1693            ..StorageLimits::default()
1694        };
1695        let mut opened = StorageEngine::open_with_limits(&root, limits)?;
1696        opened.storage.write(
1697            Uuid::now_v7(),
1698            &[
1699                Mutation::put(b"first", lexical_document("one")?),
1700                Mutation::put(b"second", lexical_document("two")?),
1701            ],
1702        )?;
1703
1704        let active_log = opened.storage.directory.active_log_path();
1705        let format_path = root.join("FORMAT");
1706        let accepted_log_bytes = fs::metadata(&active_log)?.len();
1707        let accepted_format = fs::read(&format_path)?;
1708
1709        let definition = lexical_definition("define.limit")?;
1710        let result = opened.storage.write(
1711            Uuid::now_v7(),
1712            &[Mutation::define_lexical_index(definition.clone())],
1713        );
1714        assert!(matches!(
1715            result,
1716            Err(StorageError::Index { source })
1717                if matches!(
1718                    source.as_ref(),
1719                    crate::MaterializedIndexError::Lexical(
1720                        LexicalError::DocumentBudgetExceeded { maximum: 1 }
1721                    )
1722            )
1723        ));
1724        assert_eq!(opened.storage.lexical_index(&definition.name)?, None);
1725        assert_eq!(fs::metadata(&active_log)?.len(), accepted_log_bytes);
1726        assert_eq!(fs::read(&format_path)?, accepted_format);
1727        Ok(())
1728    }
1729
1730    #[test]
1731    fn lexical_preflight_simulates_ordered_batches_and_idempotent_retries()
1732    -> Result<(), Box<dyn Error>> {
1733        let temporary = TestDirectory::new("storage-lexical-ordered-preflight")?;
1734        let root = temporary.path().join("put-before-define");
1735        let limits = StorageLimits {
1736            recovery: crate::RecoveryLimits {
1737                max_lexical_documents: 1,
1738                max_lexical_tokens: 2,
1739                ..crate::RecoveryLimits::default()
1740            },
1741            ..StorageLimits::default()
1742        };
1743        let definition = lexical_definition("ordered.limit")?;
1744        let first_transaction = Uuid::now_v7();
1745        let first_batch = [
1746            Mutation::put(b"first", lexical_document("one two")?),
1747            Mutation::define_lexical_index(definition.clone()),
1748        ];
1749        let mut opened = StorageEngine::open_with_limits(&root, limits.clone())?;
1750        let first_outcome = opened.storage.write(first_transaction, &first_batch)?;
1751        assert!(matches!(first_outcome, AppendOutcome::Committed(_)));
1752
1753        opened.storage.write(
1754            Uuid::now_v7(),
1755            &[
1756                Mutation::delete(b"first"),
1757                Mutation::put(b"second", lexical_document("three")?),
1758            ],
1759        )?;
1760        assert_eq!(opened.storage.get(b"first")?, None);
1761        assert!(opened.storage.get(b"second")?.is_some());
1762        assert!(matches!(
1763            opened.storage.write(
1764                Uuid::now_v7(),
1765                &[
1766                    Mutation::put(b"third", lexical_document("four")?),
1767                    Mutation::delete(b"second"),
1768                ],
1769            ),
1770            Err(StorageError::Index { source })
1771                if matches!(
1772                    source.as_ref(),
1773                    crate::MaterializedIndexError::Lexical(
1774                        LexicalError::DocumentBudgetExceeded { maximum: 1 }
1775                    )
1776                )
1777        ));
1778        assert!(matches!(
1779            opened.storage.write(first_transaction, &first_batch)?,
1780            AppendOutcome::Existing(_)
1781        ));
1782        drop(opened);
1783
1784        let define_first_root = temporary.path().join("define-before-put");
1785        let mut define_first = StorageEngine::open_with_limits(&define_first_root, limits)?;
1786        define_first.storage.write(
1787            Uuid::now_v7(),
1788            &[
1789                Mutation::define_lexical_index(lexical_definition("ordered.second")?),
1790                Mutation::put(b"document", lexical_document("one two")?),
1791            ],
1792        )?;
1793        assert!(define_first.storage.get(b"document")?.is_some());
1794        Ok(())
1795    }
1796
1797    #[test]
1798    fn kv_scan_pages_are_strictly_ordered_and_exclusive() -> Result<(), Box<dyn Error>> {
1799        let temporary = TestDirectory::new("storage-scan-page")?;
1800        let mut opened = StorageEngine::open(temporary.path().join("data"))?;
1801        opened.storage.write(
1802            Uuid::now_v7(),
1803            &[
1804                Mutation::put(b"c", b"three"),
1805                Mutation::put(b"a", b"one"),
1806                Mutation::put(b"b", b"two"),
1807            ],
1808        )?;
1809
1810        let first = opened.storage.scan_page(None, 2)?;
1811        assert_eq!(
1812            first
1813                .entries
1814                .iter()
1815                .map(|entry| entry.key.as_slice())
1816                .collect::<Vec<_>>(),
1817            [b"a".as_slice(), b"b".as_slice()]
1818        );
1819        assert_eq!(first.next_after, Some(b"b".to_vec()));
1820
1821        let second = opened.storage.scan_page(first.next_after.as_deref(), 2)?;
1822        assert_eq!(second.entries[0].key, b"c");
1823        assert_eq!(second.next_after, None);
1824
1825        let total_bytes = u64::try_from(
1826            b"a".len() + b"one".len() + b"b".len() + b"two".len() + b"c".len() + b"three".len(),
1827        )?;
1828        assert_eq!(
1829            opened
1830                .storage
1831                .scan_page_with_byte_limit(None, 3, total_bytes)?
1832                .entries
1833                .len(),
1834            3
1835        );
1836        assert!(matches!(
1837            opened
1838                .storage
1839                .scan_page_with_byte_limit(None, 3, total_bytes - 1),
1840            Err(ScanPageError::ByteBudgetExceeded { maximum })
1841                if maximum == total_bytes - 1
1842        ));
1843        let first_bytes = u64::try_from(b"a".len() + b"one".len())?;
1844        let first_only = opened
1845            .storage
1846            .scan_page_with_byte_limit(None, 1, first_bytes)?;
1847        assert_eq!(first_only.entries.len(), 1);
1848        assert_eq!(first_only.next_after, Some(b"a".to_vec()));
1849        Ok(())
1850    }
1851}