1use std::path::{Path, PathBuf};
4
5use hyphae_core::{VectorSpaceDefinition, VectorSpaceName};
6use hyphae_retrieval::{ExactRetrievalError, LexicalIndexDefinition, LexicalMaterializedCorpus};
7use thiserror::Error;
8use uuid::Uuid;
9
10pub 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#[derive(Debug, Error)]
28pub enum StorageError {
29 #[error(transparent)]
31 DataDirectory(#[from] DataDirectoryError),
32
33 #[error(transparent)]
35 Log(#[from] LogError),
36
37 #[error("materialized index failure: {source}")]
39 Index {
40 #[source]
42 source: Box<MaterializedIndexError>,
43 },
44
45 #[error(transparent)]
47 Mutation(#[from] MutationError),
48
49 #[error("transaction {receipt:?} is durable but not materialized; reopen to recover")]
51 CommittedButNotIndexed {
52 receipt: CommitReceipt,
54 #[source]
56 source: Box<MaterializedIndexError>,
57 },
58
59 #[error("materialized index is stale; reopen storage to replay the durable log")]
61 StaleIndex,
62
63 #[error("snapshot failure: {source}")]
65 Snapshot {
66 #[source]
68 source: Box<SnapshotError>,
69 },
70
71 #[error("storage manifest generation space is exhausted")]
73 ManifestGenerationExhausted,
74
75 #[error("prepared compaction segment is not empty: {path}")]
77 PreparedSegmentNotEmpty {
78 path: PathBuf,
80 },
81
82 #[error("scan page size {requested} is outside 1..={maximum}")]
84 InvalidScanLimit {
85 requested: usize,
87 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#[derive(Debug, Error)]
103pub enum VectorEntriesError {
104 #[error(transparent)]
106 Storage(#[from] StorageError),
107 #[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#[derive(Debug, Error)]
120pub enum ScanPageError {
121 #[error(transparent)]
123 Storage(#[from] StorageError),
124 #[error("scan page byte budget exceeded: {maximum}")]
126 ByteBudgetExceeded {
127 maximum: u64,
129 },
130}
131
132#[derive(Clone, Debug, Eq, PartialEq)]
134pub struct StorageRecoveryReport {
135 pub log: RecoveryReport,
137 pub replayed_transactions: u64,
139}
140
141#[derive(Debug)]
143pub struct OpenedStorage {
144 pub storage: StorageEngine,
146 pub recovery: StorageRecoveryReport,
148}
149
150#[derive(Clone, Debug, Eq, PartialEq)]
152pub struct CompactionReport {
153 pub generation: u64,
155 pub snapshot: SnapshotInfo,
157 pub retired_segment: PathBuf,
159 pub retired_segment_removed: bool,
161}
162
163#[derive(Clone, Debug, Eq, PartialEq)]
165pub struct KvEntry {
166 pub key: Vec<u8>,
168 pub value: Vec<u8>,
170}
171
172#[derive(Clone, Debug, Eq, PartialEq)]
174pub struct KvPage {
175 pub entries: Vec<KvEntry>,
177 pub next_after: Option<Vec<u8>>,
179}
180
181#[derive(Clone, Debug, Eq, PartialEq)]
183pub enum CompactionOutcome {
184 NoChanges {
186 snapshot: SnapshotInfo,
188 },
189 Compacted(CompactionReport),
191}
192
193#[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 pub fn open(path: impl AsRef<Path>) -> Result<OpenedStorage, StorageError> {
211 Self::open_with_limits(path, StorageLimits::compatibility())
212 }
213
214 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 pub fn data_path(&self) -> &Path {
263 self.directory.path()
264 }
265
266 pub fn backup(&self, destination: impl AsRef<Path>) -> Result<BackupInfo, BackupError> {
274 crate::backup::create_backup(self, destination.as_ref())
275 }
276
277 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 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 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 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 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 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 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 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 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 pub fn index_path(&self) -> PathBuf {
542 self.directory.path().join("indexes").join("primary.redb")
543 }
544
545 pub fn snapshot(&self) -> Result<SnapshotInfo, StorageError> {
552 let limits = self.limits.maintenance.clone();
553 self.snapshot_with_limits(&limits)
554 }
555
556 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 pub fn compact(&mut self) -> Result<CompactionOutcome, StorageError> {
614 let limits = self.limits.maintenance.clone();
615 self.compact_with_limits(&limits)
616 }
617
618 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}