1use std::path::{Path, PathBuf};
4
5use hyphae_core::{VectorSpaceDefinition, VectorSpaceName};
6use hyphae_retrieval::{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, MaterializedIndexError, Mutation, MutationError, RecoveredTransaction,
17 RecoveryReport, SnapshotError, SnapshotInfo,
18 index::{MaterializedIndex, VectorEntry},
19 manifest::StorageManifest,
20 mutation::validate_key,
21 snapshot::{create_snapshot, verify_snapshot},
22};
23
24#[derive(Debug, Error)]
26pub enum StorageError {
27 #[error(transparent)]
29 DataDirectory(#[from] DataDirectoryError),
30
31 #[error(transparent)]
33 Log(#[from] LogError),
34
35 #[error("materialized index failure: {source}")]
37 Index {
38 #[source]
40 source: Box<MaterializedIndexError>,
41 },
42
43 #[error(transparent)]
45 Mutation(#[from] MutationError),
46
47 #[error("transaction {receipt:?} is durable but not materialized; reopen to recover")]
49 CommittedButNotIndexed {
50 receipt: CommitReceipt,
52 #[source]
54 source: Box<MaterializedIndexError>,
55 },
56
57 #[error("materialized index is stale; reopen storage to replay the durable log")]
59 StaleIndex,
60
61 #[error("snapshot failure: {source}")]
63 Snapshot {
64 #[source]
66 source: Box<SnapshotError>,
67 },
68
69 #[error("storage manifest generation space is exhausted")]
71 ManifestGenerationExhausted,
72
73 #[error("prepared compaction segment is not empty: {path}")]
75 PreparedSegmentNotEmpty {
76 path: PathBuf,
78 },
79
80 #[error("scan page size {requested} is outside 1..={maximum}")]
82 InvalidScanLimit {
83 requested: usize,
85 maximum: usize,
87 },
88}
89
90#[derive(Clone, Debug, Eq, PartialEq)]
92pub struct StorageRecoveryReport {
93 pub log: RecoveryReport,
95 pub replayed_transactions: u64,
97}
98
99#[derive(Debug)]
101pub struct OpenedStorage {
102 pub storage: StorageEngine,
104 pub recovery: StorageRecoveryReport,
106}
107
108#[derive(Clone, Debug, Eq, PartialEq)]
110pub struct CompactionReport {
111 pub generation: u64,
113 pub snapshot: SnapshotInfo,
115 pub retired_segment: PathBuf,
117 pub retired_segment_removed: bool,
119}
120
121#[derive(Clone, Debug, Eq, PartialEq)]
123pub struct KvEntry {
124 pub key: Vec<u8>,
126 pub value: Vec<u8>,
128}
129
130#[derive(Clone, Debug, Eq, PartialEq)]
132pub struct KvPage {
133 pub entries: Vec<KvEntry>,
135 pub next_after: Option<Vec<u8>>,
137}
138
139#[derive(Clone, Debug, Eq, PartialEq)]
141pub enum CompactionOutcome {
142 NoChanges {
144 snapshot: SnapshotInfo,
146 },
147 Compacted(CompactionReport),
149}
150
151#[derive(Debug)]
153pub struct StorageEngine {
154 log: DurableLog,
155 index: MaterializedIndex,
156 index_stale: bool,
157 directory: DataDirectory,
158}
159
160impl StorageEngine {
161 pub fn open(path: impl AsRef<Path>) -> Result<OpenedStorage, StorageError> {
168 let directory = DataDirectory::open(path)?;
169 let index_path = directory.path().join("indexes").join("primary.redb");
170 ensure_snapshot_base(&directory, &index_path)?;
171 let (base_sequence, base_digest) = directory.log_anchor();
172 let (log, log_recovery) = DurableLog::open_file_at_version(
173 directory.active_log_path(),
174 base_sequence,
175 base_digest,
176 directory.disk_format_version(),
177 )?;
178 let index = MaterializedIndex::open(index_path)?;
179 let replayed_transactions = index.replay(&log_recovery)?;
180 let _cleanup_complete = directory.cleanup_retired_logs();
181 let storage = Self {
182 log,
183 index,
184 index_stale: false,
185 directory,
186 };
187 Ok(OpenedStorage {
188 storage,
189 recovery: StorageRecoveryReport {
190 log: log_recovery,
191 replayed_transactions,
192 },
193 })
194 }
195
196 pub fn data_path(&self) -> &Path {
198 self.directory.path()
199 }
200
201 pub fn backup(&self, destination: impl AsRef<Path>) -> Result<BackupInfo, BackupError> {
209 crate::backup::create_backup(self, destination.as_ref())
210 }
211
212 pub fn write(
223 &mut self,
224 transaction_id: Uuid,
225 mutations: &[Mutation],
226 ) -> Result<AppendOutcome, StorageError> {
227 if self.index_stale {
228 return Err(StorageError::StaleIndex);
229 }
230 self.index.validate_mutations(mutations)?;
231 if mutations.iter().any(|mutation| {
232 matches!(
233 mutation,
234 Mutation::DefineVectorSpace { .. }
235 | Mutation::UpsertVector { .. }
236 | Mutation::DeleteVector { .. }
237 | Mutation::DefineLexicalIndex { .. }
238 )
239 }) {
240 self.directory.promote_format()?;
241 self.log
242 .set_disk_format_version(self.directory.disk_format_version())?;
243 }
244 let operations = mutations
245 .iter()
246 .map(Mutation::encode)
247 .collect::<Result<Vec<_>, _>>()?;
248 let operation_count =
249 u32::try_from(operations.len()).map_err(|_| LogError::TooManyOperations)?;
250 let requested_digest = transaction_digest(&operations, operation_count)?;
251 if let Some(receipt) = self.index.receipt(transaction_id)? {
252 return if receipt.transaction_digest == requested_digest {
253 Ok(AppendOutcome::Existing(receipt))
254 } else {
255 Err(LogError::IdempotencyConflict { transaction_id }.into())
256 };
257 }
258 let outcome = match self.log.append_transaction(transaction_id, &operations) {
259 Ok(outcome) => outcome,
260 Err(source) => {
261 if self.log.is_poisoned() {
262 self.index_stale = true;
263 }
264 return Err(source.into());
265 }
266 };
267 let AppendOutcome::Committed(receipt) = outcome else {
268 return Ok(outcome);
269 };
270
271 let transaction = RecoveredTransaction {
272 receipt,
273 operations,
274 };
275 if let Err(source) = self.index.apply(&transaction) {
276 self.index_stale = true;
277 return Err(StorageError::CommittedButNotIndexed {
278 receipt,
279 source: Box::new(source),
280 });
281 }
282 Ok(outcome)
283 }
284
285 pub fn get(&self, key: &[u8]) -> Result<Option<Vec<u8>>, StorageError> {
292 if self.index_stale {
293 return Err(StorageError::StaleIndex);
294 }
295 validate_key(key)?;
296 Ok(self.index.get(key)?)
297 }
298
299 pub fn vector_space(
305 &self,
306 name: &VectorSpaceName,
307 ) -> Result<Option<VectorSpaceDefinition>, StorageError> {
308 if self.index_stale {
309 return Err(StorageError::StaleIndex);
310 }
311 Ok(self.index.vector_space(name)?)
312 }
313
314 pub fn lexical_index(
320 &self,
321 name: &VectorSpaceName,
322 ) -> Result<Option<LexicalIndexDefinition>, StorageError> {
323 if self.index_stale {
324 return Err(StorageError::StaleIndex);
325 }
326 Ok(self.index.lexical_index(name)?)
327 }
328
329 pub fn lexical_corpus(
337 &self,
338 definition: &LexicalIndexDefinition,
339 query_tokens: &[String],
340 max_candidates: u64,
341 timeout: std::time::Duration,
342 ) -> Result<LexicalMaterializedCorpus, StorageError> {
343 if self.index_stale {
344 return Err(StorageError::StaleIndex);
345 }
346 Ok(self
347 .index
348 .lexical_corpus(definition, query_tokens, max_candidates, timeout)?)
349 }
350
351 pub fn vector_entries(
359 &self,
360 name: &VectorSpaceName,
361 max_candidates: u64,
362 max_bytes: u64,
363 ) -> Result<Vec<VectorEntry>, StorageError> {
364 if self.index_stale {
365 return Err(StorageError::StaleIndex);
366 }
367 Ok(self.index.scan_vectors(name, max_candidates, max_bytes)?)
368 }
369
370 pub fn scan_page(&self, after: Option<&[u8]>, limit: usize) -> Result<KvPage, StorageError> {
380 if self.index_stale {
381 return Err(StorageError::StaleIndex);
382 }
383 if let Some(key) = after {
384 validate_key(key)?;
385 }
386 if limit == 0 || limit > MAX_SCAN_PAGE_ENTRIES {
387 return Err(StorageError::InvalidScanLimit {
388 requested: limit,
389 maximum: MAX_SCAN_PAGE_ENTRIES,
390 });
391 }
392 let mut raw = self.index.scan_after(after, limit.saturating_add(1))?;
393 let has_more = raw.len() > limit;
394 raw.truncate(limit);
395 let next_after = has_more
396 .then(|| raw.last().map(|(key, _)| key.clone()))
397 .flatten();
398 Ok(KvPage {
399 entries: raw
400 .into_iter()
401 .map(|(key, value)| KvEntry { key, value })
402 .collect(),
403 next_after,
404 })
405 }
406
407 pub fn index_path(&self) -> PathBuf {
409 self.directory.path().join("indexes").join("primary.redb")
410 }
411
412 pub fn snapshot(&self) -> Result<SnapshotInfo, StorageError> {
419 if self.index_stale {
420 return Err(StorageError::StaleIndex);
421 }
422 let snapshots = self.directory.path().join("snapshots");
423 let temporary = self.directory.path().join("tmp");
424 Ok(create_snapshot(
425 &self.index,
426 &snapshots,
427 &temporary,
428 self.directory.disk_format_version(),
429 )?)
430 }
431
432 pub fn compact(&mut self) -> Result<CompactionOutcome, StorageError> {
443 if self.index_stale {
444 return Err(StorageError::StaleIndex);
445 }
446 let snapshot = self.snapshot()?;
447 let current = self.directory.manifest();
448 if snapshot.checkpoint_sequence == 0
449 || snapshot.checkpoint_sequence == current.base_sequence
450 {
451 return Ok(CompactionOutcome::NoChanges { snapshot });
452 }
453 let generation = current
454 .generation
455 .checked_add(1)
456 .ok_or(StorageError::ManifestGenerationExhausted)?;
457 let Some(base_digest) = snapshot.checkpoint_digest else {
458 return Err(SnapshotError::Invalid {
459 reason: "compaction snapshot lacks a checkpoint digest",
460 }
461 .into());
462 };
463 let next = StorageManifest {
464 generation,
465 active_segment: generation,
466 base_sequence: snapshot.checkpoint_sequence,
467 base_digest,
468 snapshot_digest: snapshot.snapshot_digest,
469 };
470 let next_segment = self.directory.log_path(generation);
471 let (next_log, prepared) = DurableLog::open_file_at_version(
472 &next_segment,
473 next.base_sequence,
474 next.base_digest,
475 self.directory.disk_format_version(),
476 )?;
477 if prepared.valid_bytes != 0 {
478 return Err(StorageError::PreparedSegmentNotEmpty { path: next_segment });
479 }
480
481 let retired_segment = self.directory.active_log_path();
482 self.directory.commit_manifest(next)?;
483 let retired_log = std::mem::replace(&mut self.log, next_log);
484 drop(retired_log);
485 let retired_segment_removed = remove_retired_segment(&retired_segment);
486 Ok(CompactionOutcome::Compacted(CompactionReport {
487 generation,
488 snapshot,
489 retired_segment,
490 retired_segment_removed,
491 }))
492 }
493}
494
495fn remove_retired_segment(path: &Path) -> bool {
496 match std::fs::remove_file(path) {
497 Ok(()) => {
498 #[cfg(unix)]
499 if let Some(parent) = path.parent()
500 && sync_directory(parent).is_err()
501 {
502 return false;
503 }
504 true
505 }
506 Err(source) if source.kind() == std::io::ErrorKind::NotFound => true,
507 Err(_) => false,
508 }
509}
510
511fn ensure_snapshot_base(directory: &DataDirectory, index_path: &Path) -> Result<(), StorageError> {
512 let manifest = directory.manifest();
513 if manifest.base_sequence == 0 {
514 return Ok(());
515 }
516 let snapshot_path = directory.snapshot_path(manifest.base_sequence);
517 let verified = verify_snapshot(&snapshot_path)?;
518 if verified.checkpoint_sequence != manifest.base_sequence
519 || verified.checkpoint_digest != Some(manifest.base_digest)
520 || verified.snapshot_digest != manifest.snapshot_digest
521 {
522 return Err(SnapshotError::Invalid {
523 reason: "snapshot does not match active storage manifest",
524 }
525 .into());
526 }
527 if index_path.exists() {
528 return Ok(());
529 }
530
531 let temporary_path = directory
532 .path()
533 .join("tmp")
534 .join(format!("index-restore-{}.redb.tmp", Uuid::now_v7()));
535 let restored = MaterializedIndex::restore_from_snapshot(&temporary_path, &snapshot_path)?;
536 if restored != verified {
537 return Err(SnapshotError::Invalid {
538 reason: "snapshot changed while rebuilding the materialized index",
539 }
540 .into());
541 }
542 std::fs::rename(&temporary_path, index_path)
543 .map_err(SnapshotError::from)
544 .map_err(StorageError::from)?;
545 #[cfg(unix)]
546 sync_directory(index_path.parent().ok_or(SnapshotError::Invalid {
547 reason: "materialized index path has no parent",
548 })?)?;
549 Ok(())
550}
551
552#[cfg(unix)]
553fn sync_directory(path: &Path) -> Result<(), StorageError> {
554 std::fs::File::open(path)
555 .and_then(|directory| directory.sync_all())
556 .map_err(SnapshotError::from)
557 .map_err(StorageError::from)
558}
559
560impl From<MaterializedIndexError> for StorageError {
561 fn from(source: MaterializedIndexError) -> Self {
562 Self::Index {
563 source: Box::new(source),
564 }
565 }
566}
567
568impl From<SnapshotError> for StorageError {
569 fn from(source: SnapshotError) -> Self {
570 Self::Snapshot {
571 source: Box::new(source),
572 }
573 }
574}
575
576#[cfg(test)]
577mod tests {
578 use std::{
579 error::Error,
580 fs::{self, OpenOptions},
581 io::{Seek, SeekFrom, Write},
582 };
583
584 use uuid::Uuid;
585
586 use super::{CompactionOutcome, DurableLog, StorageEngine, StorageError, StorageManifest};
587 use crate::{
588 AppendOutcome, DataDirectory, Mutation, SnapshotError, SnapshotReadLimits,
589 index::MaterializedIndex, load_snapshot, test_support::TestDirectory, verify_snapshot,
590 };
591
592 #[test]
593 fn atomic_batches_persist_and_delete() -> Result<(), Box<dyn Error>> {
594 let temporary = TestDirectory::new("storage-kv")?;
595 let root = temporary.path().join("data");
596 let mut opened = StorageEngine::open(&root)?;
597 opened.storage.write(
598 Uuid::now_v7(),
599 &[Mutation::put(b"a", b"one"), Mutation::put(b"b", b"two")],
600 )?;
601 assert_eq!(opened.storage.get(b"a")?, Some(b"one".to_vec()));
602 opened
603 .storage
604 .write(Uuid::now_v7(), &[Mutation::delete(b"a")])?;
605 drop(opened);
606
607 let reopened = StorageEngine::open(&root)?;
608 assert_eq!(reopened.storage.get(b"a")?, None);
609 assert_eq!(reopened.storage.get(b"b")?, Some(b"two".to_vec()));
610 assert_eq!(reopened.recovery.replayed_transactions, 0);
611 Ok(())
612 }
613
614 #[test]
615 fn exact_retry_does_not_reapply_a_batch() -> Result<(), Box<dyn Error>> {
616 let temporary = TestDirectory::new("storage-idempotency")?;
617 let transaction_id = Uuid::now_v7();
618 let mutations = [Mutation::put(b"key", b"value")];
619 let mut opened = StorageEngine::open(temporary.path())?;
620
621 let first = opened.storage.write(transaction_id, &mutations)?;
622 let second = opened.storage.write(transaction_id, &mutations)?;
623 assert!(matches!(first, AppendOutcome::Committed(_)));
624 assert!(matches!(second, AppendOutcome::Existing(_)));
625 assert_eq!(opened.storage.get(b"key")?, Some(b"value".to_vec()));
626
627 let conflict = opened
628 .storage
629 .write(transaction_id, &[Mutation::put(b"key", b"different")]);
630 assert!(matches!(
631 conflict,
632 Err(super::StorageError::Log(
633 crate::LogError::IdempotencyConflict { .. }
634 ))
635 ));
636 Ok(())
637 }
638
639 #[test]
640 fn reopen_replays_a_commit_missing_from_the_index() -> Result<(), Box<dyn Error>> {
641 let temporary = TestDirectory::new("storage-index-replay")?;
642 let root = temporary.path().join("data");
643 let directory = DataDirectory::open(&root)?;
644 let mutation = Mutation::put(b"recovered", b"yes");
645 let mut log = directory.open_log()?;
646 log.log
647 .append_transaction(Uuid::now_v7(), &[mutation.encode()?])?;
648 drop(log);
649 drop(directory);
650
651 let reopened = StorageEngine::open(&root)?;
652 assert_eq!(reopened.recovery.replayed_transactions, 1);
653 assert_eq!(reopened.storage.get(b"recovered")?, Some(b"yes".to_vec()));
654 Ok(())
655 }
656
657 #[test]
658 fn logical_snapshot_is_stable_and_detects_payload_corruption() -> Result<(), Box<dyn Error>> {
659 let temporary = TestDirectory::new("storage-snapshot")?;
660 let mut opened = StorageEngine::open(temporary.path().join("data"))?;
661 opened.storage.write(
662 Uuid::now_v7(),
663 &[
664 Mutation::put(b"beta", b"second"),
665 Mutation::put(b"alpha", b"first"),
666 ],
667 )?;
668
669 let created = opened.storage.snapshot()?;
670 assert_eq!(created.checkpoint_sequence, 4);
671 assert!(created.checkpoint_digest.is_some());
672 assert_eq!(created.entry_count, 2);
673 assert_eq!(created.receipt_count, 1);
674 assert_eq!(verify_snapshot(&created.path)?, created);
675 assert_eq!(opened.storage.snapshot()?, created);
676 let witness = load_snapshot(&created.path, &SnapshotReadLimits::default())?;
677 assert_eq!(witness.info, created);
678 assert_eq!(witness.entries.len(), 2);
679 assert_eq!(witness.entries[0].key, b"alpha");
680 assert_eq!(witness.entries[0].value, b"first");
681 assert!(matches!(
682 load_snapshot(
683 &created.path,
684 &SnapshotReadLimits {
685 entries: 1,
686 ..SnapshotReadLimits::default()
687 }
688 ),
689 Err(SnapshotError::EntryLimitExceeded {
690 actual: 2,
691 maximum: 1
692 })
693 ));
694
695 let corrupted_path = temporary.path().join("corrupted.hysnap");
696 fs::copy(&created.path, &corrupted_path)?;
697 let mut corrupted = OpenOptions::new()
698 .read(true)
699 .write(true)
700 .open(&corrupted_path)?;
701 corrupted.seek(SeekFrom::End(-1))?;
702 corrupted.write_all(&[0xff])?;
703 corrupted.sync_all()?;
704 drop(corrupted);
705
706 assert!(matches!(
707 verify_snapshot(&corrupted_path),
708 Err(SnapshotError::Invalid {
709 reason: "CRC32C mismatch"
710 })
711 ));
712 Ok(())
713 }
714
715 #[test]
716 fn empty_storage_has_a_canonical_empty_snapshot() -> Result<(), Box<dyn Error>> {
717 let temporary = TestDirectory::new("storage-empty-snapshot")?;
718 let opened = StorageEngine::open(temporary.path().join("data"))?;
719
720 let snapshot = opened.storage.snapshot()?;
721 assert_eq!(snapshot.checkpoint_sequence, 0);
722 assert_eq!(snapshot.checkpoint_digest, None);
723 assert_eq!(snapshot.entry_count, 0);
724 assert_eq!(snapshot.vector_space_count, 0);
725 assert_eq!(snapshot.vector_count, 0);
726 assert_eq!(snapshot.lexical_index_count, 0);
727 assert_eq!(snapshot.receipt_count, 0);
728 assert_eq!(snapshot.file_bytes, 136);
729 assert_eq!(verify_snapshot(&snapshot.path)?, snapshot);
730 Ok(())
731 }
732
733 #[test]
734 fn snapshot_rebuilds_kv_and_idempotency_state() -> Result<(), Box<dyn Error>> {
735 let temporary = TestDirectory::new("storage-snapshot-restore")?;
736 let root = temporary.path().join("data");
737 let transaction_id = Uuid::now_v7();
738 let mut opened = StorageEngine::open(&root)?;
739 let outcome = opened.storage.write(
740 transaction_id,
741 &[
742 Mutation::put(b"alpha", b"one"),
743 Mutation::put(b"beta", b"two"),
744 ],
745 )?;
746 let AppendOutcome::Committed(receipt) = outcome else {
747 return Err("new transaction was not committed".into());
748 };
749 let snapshot = opened.storage.snapshot()?;
750 drop(opened);
751
752 let restored_path = root.join("tmp/restored.redb");
753 assert_eq!(
754 MaterializedIndex::restore_from_snapshot(&restored_path, &snapshot.path)?,
755 snapshot
756 );
757 let restored = MaterializedIndex::open(&restored_path)?;
758 assert_eq!(restored.get(b"alpha")?, Some(b"one".to_vec()));
759 assert_eq!(restored.get(b"beta")?, Some(b"two".to_vec()));
760 assert_eq!(restored.receipt(transaction_id)?, Some(receipt));
761 Ok(())
762 }
763
764 #[test]
765 fn compaction_retires_history_and_snapshot_rebuilds_the_index() -> Result<(), Box<dyn Error>> {
766 let temporary = TestDirectory::new("storage-compaction")?;
767 let root = temporary.path().join("data");
768 let first_id = Uuid::now_v7();
769 let mut opened = StorageEngine::open(&root)?;
770 let first = opened
771 .storage
772 .write(first_id, &[Mutation::put(b"before", b"one")])?;
773 let AppendOutcome::Committed(first_receipt) = first else {
774 return Err("first transaction was not committed".into());
775 };
776
777 let compacted = opened.storage.compact()?;
778 let CompactionOutcome::Compacted(report) = compacted else {
779 return Err("committed history was not compacted".into());
780 };
781 assert_eq!(report.generation, 2);
782 assert!(report.retired_segment_removed);
783 assert!(!report.retired_segment.exists());
784 assert!(root.join("log/00000000000000000002.hylog").is_file());
785 assert!(matches!(
786 opened.storage.compact()?,
787 CompactionOutcome::NoChanges { .. }
788 ));
789
790 assert_eq!(
791 opened
792 .storage
793 .write(first_id, &[Mutation::put(b"before", b"one")])?,
794 AppendOutcome::Existing(first_receipt)
795 );
796 let second_id = Uuid::now_v7();
797 let second = opened
798 .storage
799 .write(second_id, &[Mutation::put(b"after", b"two")])?;
800 let AppendOutcome::Committed(second_receipt) = second else {
801 return Err("second transaction was not committed".into());
802 };
803 assert_eq!(
804 second_receipt.commit_sequence,
805 first_receipt.commit_sequence + 3
806 );
807 drop(opened);
808
809 fs::remove_file(root.join("indexes/primary.redb"))?;
810 let mut rebuilt = StorageEngine::open(&root)?;
811 assert_eq!(rebuilt.recovery.replayed_transactions, 1);
812 assert_eq!(rebuilt.storage.get(b"before")?, Some(b"one".to_vec()));
813 assert_eq!(rebuilt.storage.get(b"after")?, Some(b"two".to_vec()));
814 assert_eq!(
815 rebuilt
816 .storage
817 .write(first_id, &[Mutation::put(b"before", b"one")])?,
818 AppendOutcome::Existing(first_receipt)
819 );
820 Ok(())
821 }
822
823 #[test]
824 fn orphan_prepared_segment_is_ignored_until_manifest_commit() -> Result<(), Box<dyn Error>> {
825 let temporary = TestDirectory::new("storage-compaction-orphan")?;
826 let root = temporary.path().join("data");
827 let mut opened = StorageEngine::open(&root)?;
828 opened
829 .storage
830 .write(Uuid::now_v7(), &[Mutation::put(b"key", b"value")])?;
831 let snapshot = opened.storage.snapshot()?;
832 let base_digest = snapshot
833 .checkpoint_digest
834 .ok_or("snapshot checkpoint digest is absent")?;
835 let orphan_path = opened.storage.directory.log_path(2);
836 let (orphan, recovery) =
837 DurableLog::open_file_at(&orphan_path, snapshot.checkpoint_sequence, base_digest)?;
838 assert_eq!(recovery.valid_bytes, 0);
839 drop(orphan);
840 drop(opened);
841
842 let mut reopened = StorageEngine::open(&root)?;
843 assert_eq!(reopened.storage.directory.manifest().generation, 1);
844 assert_eq!(reopened.storage.get(b"key")?, Some(b"value".to_vec()));
845 assert!(matches!(
846 reopened.storage.compact()?,
847 CompactionOutcome::Compacted(_)
848 ));
849 assert_eq!(reopened.storage.directory.manifest().generation, 2);
850 Ok(())
851 }
852
853 #[test]
854 fn committed_manifest_wins_before_retired_log_cleanup() -> Result<(), Box<dyn Error>> {
855 let temporary = TestDirectory::new("storage-compaction-committed")?;
856 let root = temporary.path().join("data");
857 let mut opened = StorageEngine::open(&root)?;
858 opened
859 .storage
860 .write(Uuid::now_v7(), &[Mutation::put(b"key", b"value")])?;
861 let snapshot = opened.storage.snapshot()?;
862 let base_digest = snapshot
863 .checkpoint_digest
864 .ok_or("snapshot checkpoint digest is absent")?;
865 let next = StorageManifest {
866 generation: 2,
867 active_segment: 2,
868 base_sequence: snapshot.checkpoint_sequence,
869 base_digest,
870 snapshot_digest: snapshot.snapshot_digest,
871 };
872 let (prepared, recovery) = DurableLog::open_file_at(
873 opened.storage.directory.log_path(2),
874 next.base_sequence,
875 next.base_digest,
876 )?;
877 assert_eq!(recovery.valid_bytes, 0);
878 drop(prepared);
879 opened.storage.directory.commit_manifest(next)?;
880 let retired_path = opened.storage.directory.log_path(1);
881 assert!(retired_path.is_file());
882 drop(opened);
883
884 let reopened = StorageEngine::open(&root)?;
885 assert_eq!(reopened.storage.directory.manifest().generation, 2);
886 assert_eq!(reopened.storage.get(b"key")?, Some(b"value".to_vec()));
887 assert!(!retired_path.exists());
888 Ok(())
889 }
890
891 #[test]
892 fn uncertain_log_sync_blocks_the_handle_until_recovery() -> Result<(), Box<dyn Error>> {
893 let temporary = TestDirectory::new("storage-injected-log-sync")?;
894 let root = temporary.path().join("data");
895 let mut opened = StorageEngine::open(&root)?;
896 opened.storage.log.inject_sync_failure();
897
898 let result = opened
899 .storage
900 .write(Uuid::now_v7(), &[Mutation::put(b"recovered", b"yes")]);
901 assert!(matches!(
902 result,
903 Err(StorageError::Log(crate::LogError::Io(_)))
904 ));
905 assert!(matches!(
906 opened.storage.get(b"recovered"),
907 Err(StorageError::StaleIndex)
908 ));
909 assert!(matches!(
910 opened.storage.snapshot(),
911 Err(StorageError::StaleIndex)
912 ));
913 assert!(matches!(
914 opened.storage.compact(),
915 Err(StorageError::StaleIndex)
916 ));
917 drop(opened);
918
919 let reopened = StorageEngine::open(&root)?;
920 assert_eq!(reopened.recovery.replayed_transactions, 1);
921 assert_eq!(reopened.storage.get(b"recovered")?, Some(b"yes".to_vec()));
922 Ok(())
923 }
924
925 #[test]
926 fn post_commit_index_failure_recovers_from_the_log() -> Result<(), Box<dyn Error>> {
927 let temporary = TestDirectory::new("storage-injected-index")?;
928 let root = temporary.path().join("data");
929 let mut opened = StorageEngine::open(&root)?;
930 opened.storage.index.inject_apply_failure();
931
932 let result = opened
933 .storage
934 .write(Uuid::now_v7(), &[Mutation::put(b"durable", b"yes")]);
935 assert!(matches!(
936 result,
937 Err(StorageError::CommittedButNotIndexed { .. })
938 ));
939 assert!(matches!(
940 opened.storage.get(b"durable"),
941 Err(StorageError::StaleIndex)
942 ));
943 drop(opened);
944
945 let reopened = StorageEngine::open(&root)?;
946 assert_eq!(reopened.recovery.replayed_transactions, 1);
947 assert_eq!(reopened.storage.get(b"durable")?, Some(b"yes".to_vec()));
948 Ok(())
949 }
950
951 #[test]
952 fn kv_scan_pages_are_strictly_ordered_and_exclusive() -> Result<(), Box<dyn Error>> {
953 let temporary = TestDirectory::new("storage-scan-page")?;
954 let mut opened = StorageEngine::open(temporary.path().join("data"))?;
955 opened.storage.write(
956 Uuid::now_v7(),
957 &[
958 Mutation::put(b"c", b"three"),
959 Mutation::put(b"a", b"one"),
960 Mutation::put(b"b", b"two"),
961 ],
962 )?;
963
964 let first = opened.storage.scan_page(None, 2)?;
965 assert_eq!(
966 first
967 .entries
968 .iter()
969 .map(|entry| entry.key.as_slice())
970 .collect::<Vec<_>>(),
971 [b"a".as_slice(), b"b".as_slice()]
972 );
973 assert_eq!(first.next_after, Some(b"b".to_vec()));
974
975 let second = opened.storage.scan_page(first.next_after.as_deref(), 2)?;
976 assert_eq!(second.entries[0].key, b"c");
977 assert_eq!(second.next_after, None);
978 Ok(())
979 }
980}