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