diff --git a/CHANGELOG.md b/CHANGELOG.md index acf49b6..526feab 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -29,7 +29,9 @@ version, since none exists yet. - **Cross-platform VFS** — Linux (`io_uring`), Windows (IOCP), macOS/iOS (Grand Central Dispatch), Android, WASM/OPFS, and WASI backends, plus a tokio thread-pool fallback and an in-memory backend, with format-bit identity - across targets. + across targets. Backends may complete a positional read or write in several + calls; PageDB finishes the caller's whole buffer before treating authenticated + metadata as present or durable. - **Snapshots** — `snapshot_to`, `restore_from`, and incremental apply, each authenticated against the exact state its manifest describes. Destinations must be empty; malformed or incomplete artifacts fail closed. @@ -38,6 +40,12 @@ version, since none exists yet. - **Online rekey** — rekey the database under a new key with mixed-cipher and mixed-epoch page coexistence (no full-file migration). - **Handle modes** — `Standalone`, `Follower`, `ReadOnly`, and `Observer`. +- **Failures report themselves** — an unreadable free-list chain, main file, or + segment catalog fails `stats()` instead of reporting zero; compaction never + skips a catalog entry whose file it cannot open; segment open distinguishes a + missing file from a permission or backend error; and only genuine contention + is reported as contention. Persisted named-counter rows are validated at open, + and commit-history keys are rejected unless exactly eight bytes. ### Known limitations diff --git a/src/btree/tree/scan.rs b/src/btree/tree/scan.rs index e37fcf9..e3352c6 100644 --- a/src/btree/tree/scan.rs +++ b/src/btree/tree/scan.rs @@ -114,6 +114,29 @@ impl BTree { } } + /// Collect at most `limit` records at or after `start` whose keys still + /// carry `prefix`, in ascending key order. + /// + /// The bounded counterpart to [`Self::scan_prefix`]. Rows sharing a prefix + /// sort contiguously, so the first key outside it ends the range — the scan + /// stops there rather than reading the rest of the tree. + /// + /// Resume as with [`Self::collect_batch_from`]: pass the last returned key + /// with a `0x00` byte appended. A batch shorter than `limit` means the + /// prefix range ended. + pub async fn collect_prefix_batch_from( + &self, + prefix: &[u8], + start: &[u8], + limit: usize, + ) -> Result, Vec)>> { + let mut batch = self.collect_batch_from(start, limit).await?; + if let Some(end) = batch.iter().position(|(key, _)| !key.starts_with(prefix)) { + batch.truncate(end); + } + Ok(batch) + } + /// Return the smallest key in the tree, or `None` if the tree is empty. /// Descends the leftmost spine only — O(tree height), not O(tree size). pub async fn first_key(&self) -> Result>> { diff --git a/src/compaction/full.rs b/src/compaction/full.rs index 48dd4b3..5f3fbf2 100644 --- a/src/compaction/full.rs +++ b/src/compaction/full.rs @@ -84,10 +84,12 @@ async fn compact_now_inner(db: &Db) -> Result { let all_segments = list_all_segments_inner(&db.pager, db.realm_id, &state).await?; for meta in all_segments { let live = crate::segment::writer::live_path(&meta.segment_id); - let file_size = match db.vfs.open(&live, crate::vfs::types::OpenMode::Read).await { - Ok(f) => f.len().await.unwrap_or(meta.total_bytes), - Err(_) => continue, - }; + let file_size = db + .vfs + .open(&live, crate::vfs::types::OpenMode::Read) + .await? + .len() + .await?; // Skip segments with < 5% garbage. let threshold = meta.total_bytes.saturating_add(meta.total_bytes / 20); if file_size <= threshold { diff --git a/src/errors.rs b/src/errors.rs index 1ff9316..d9dd3b4 100644 --- a/src/errors.rs +++ b/src/errors.rs @@ -169,6 +169,19 @@ pub enum PagedbError { #[error("recorded rekey replacement segment {replacement_segment_id:?} is missing or invalid")] RekeyReplacementMissing { replacement_segment_id: [u8; 16] }, + /// A VFS backend reported positional-I/O progress that cannot be true — + /// more bytes transferred than the caller's remaining buffer. + /// + /// Not corruption: nothing on disk is known to be wrong. The backend broke + /// the [`VfsFile`](crate::vfs::VfsFile) contract, so the reported count + /// cannot be reasoned about at all and the transfer stops instead of + /// advancing by a length it cannot trust. + #[error("vfs backend violated the {operation} contract: {detail}")] + VfsContractViolated { + operation: &'static str, + detail: &'static str, + }, + #[error("unsupported by backend")] Unsupported, @@ -473,6 +486,13 @@ impl PagedbError { Self::ArithmeticOverflow { operation } } + /// Canonical constructor for a VFS backend that broke the positional-I/O + /// contract. + #[must_use] + pub const fn vfs_contract_violated(operation: &'static str, detail: &'static str) -> Self { + Self::VfsContractViolated { operation, detail } + } + /// Canonical constructor for a rekey admission that needs both KEKs. #[must_use] pub const fn rekey_resume_key_required(source_epoch: u64, target_epoch: u64) -> Self { diff --git a/src/pager/core.rs b/src/pager/core.rs index a8562c4..64ada90 100644 --- a/src/pager/core.rs +++ b/src/pager/core.rs @@ -23,7 +23,7 @@ use crate::pager::format::data_page::{ use crate::pager::format::page_kind::PageKind; use crate::txn::db::rekey::EpochKeyring; use crate::vfs::types::{OpenMode, WriteReq}; -use crate::vfs::{Vfs, VfsFile}; +use crate::vfs::{Vfs, VfsFile, read_exact_at}; use crate::{RealmId, Result}; use rayon::prelude::*; @@ -741,13 +741,14 @@ impl Pager { } let mut buf = vec![0u8; page_size]; { - let f = file_handle.lock().await; - let n = f.read_at(page_offset, &mut buf).await?; - if n < page_size { - for b in &mut buf[n..] { - *b = 0; - } - } + // Returned rather than retried, unlike the AEAD failures below. + // The transient this loop absorbs is a *torn* read — a + // full-length buffer of mixed old and new bytes — which is why + // retrying it can succeed. A transfer that ends short or errors + // is not that condition, and retrying it would only mask a + // one-shot backend fault. + let mut f = file_handle.lock().await; + read_exact_at(&mut *f, page_offset, &mut buf).await?; } // Extract the cipher_id and mk_epoch recorded in this specific page's diff --git a/src/pager/header.rs b/src/pager/header.rs index c63a18d..5e817fb 100644 --- a/src/pager/header.rs +++ b/src/pager/header.rs @@ -11,7 +11,7 @@ use crate::pager::format::structural_header::{ MainDbHeaderFields, decode_main_db_header, encode_main_db_header, }; use crate::vfs::types::OpenMode; -use crate::vfs::{Vfs, VfsFile}; +use crate::vfs::{Vfs, VfsFile, read_exact_at, write_all_at}; /// Which header slot is the authoritative current header. #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -57,7 +57,7 @@ pub async fn bootstrap_header( let bytes = encode_main_db_header(initial, hk, page_size)?; let mut f = vfs.open(path, OpenMode::CreateNew).await?; // Slot A at offset 0. - f.write_at(0, &bytes).await?; + write_all_at(&mut f, 0, &bytes).await?; // Slot B at offset page_size — write a zero-filled page so the slot is // materialised on disk. Decode of a zero buffer fails the magic check // and returns `Corruption(HeaderUnverifiable)`, which `open_header` @@ -65,7 +65,7 @@ pub async fn bootstrap_header( let zero = vec![0u8; page_size]; let page_size_u64 = u64::try_from(page_size) .map_err(|_| PagedbError::Io(std::io::Error::other("page_size > u64")))?; - f.write_at(page_size_u64, &zero).await?; + write_all_at(&mut f, page_size_u64, &zero).await?; f.sync().await?; // Make the directory entry for the newly created file durable so a // power loss immediately after creation does not lose the file. @@ -82,6 +82,28 @@ pub async fn bootstrap_header( /// Open an existing main.db, return the active header fields and which slot /// they came from. /// +/// Read one header slot into `buf`, treating an incomplete slot as an +/// unverifiable one rather than as a fatal error. +/// +/// The A/B protocol survives one unusable slot by design, and a file truncated +/// mid-slot is exactly that case: the surviving copy must still open the +/// database. `buf` keeps whatever prefix was transferred, so the caller's +/// ordinary decode rejects it on the magic or HK-MAC check like any other +/// damaged slot. Every other backend error propagates unchanged — only +/// end-of-file is absence. +pub(crate) async fn read_header_slot( + file: &mut F, + offset: u64, + buf: &mut [u8], +) -> Result<()> { + match read_exact_at(file, offset, buf).await { + Err(PagedbError::Io(ref error)) if error.kind() == std::io::ErrorKind::UnexpectedEof => { + Ok(()) + } + other => other, + } +} + /// Reads both slots; verifies each via HK-MAC; picks the one with the /// greater `seq`. If only one verifies, it wins. If neither verifies, returns /// `Corruption(HeaderUnverifiable)` — unrecoverable from inside the header @@ -92,13 +114,13 @@ pub async fn open_header( hk: &DerivedKey, page_size: usize, ) -> Result<(MainDbHeaderFields, ActiveSlot)> { - let f = vfs.open(path, OpenMode::ReadWrite).await?; + let mut f = vfs.open(path, OpenMode::ReadWrite).await?; let mut buf_a = vec![0u8; page_size]; let mut buf_b = vec![0u8; page_size]; - let _ = f.read_at(0, &mut buf_a).await?; + read_header_slot(&mut f, 0, &mut buf_a).await?; let page_size_u64 = u64::try_from(page_size) .map_err(|_| PagedbError::Io(std::io::Error::other("page_size > u64")))?; - let _ = f.read_at(page_size_u64, &mut buf_b).await?; + read_header_slot(&mut f, page_size_u64, &mut buf_b).await?; let a = decode_main_db_header(&buf_a, hk, page_size).ok(); let b = decode_main_db_header(&buf_b, hk, page_size).ok(); match (a, b) { @@ -139,7 +161,7 @@ pub async fn commit_header( .ok() .map(|s| next.page_id().saturating_mul(s)) .ok_or_else(|| PagedbError::Io(std::io::Error::other("offset arithmetic overflow")))?; - f.write_at(offset, &bytes).await?; + write_all_at(&mut f, offset, &bytes).await?; f.sync().await?; // No `sync_dir` here: a header rewrite is a data write to an existing, // already-durable inode (main.db). Architecture §883 requires `sync_dir` diff --git a/src/segment/authenticated_metadata.rs b/src/segment/authenticated_metadata.rs index 059ab9c..13fc7c8 100644 --- a/src/segment/authenticated_metadata.rs +++ b/src/segment/authenticated_metadata.rs @@ -11,8 +11,7 @@ use crate::pager::format::data_page::{body, extract_page_header_ids, open_data_p use crate::pager::format::page_kind::PageKind; use crate::pager::format::segment_footer::{SegmentFooterFields, decode_segment_footer}; use crate::pager::format::structural_header::{SegmentHeaderFields, decode_segment_header}; -use crate::vfs::Vfs; -use crate::vfs::VfsFile; +use crate::vfs::{Vfs, VfsFile, read_exact_at_borrowed}; use super::types::{EXTENT_INDEX_ENTRY_LEN, ExtentIndexEntry}; use super::writer::{live_path, staging_path}; @@ -85,18 +84,12 @@ pub(crate) async fn authenticate_segment_metadata( let master_key = pager.mk_for(meta.mk_epoch, cipher_id)?; let hk = derive_hk(&master_key)?; let mut header_bytes = vec![0u8; page_size]; - let header_read = file.read_at(0, &mut header_bytes).await?; - if header_read != page_size { - return Err(PagedbError::segment_geometry_invalid("header_read")); - } + read_exact_at_borrowed!(file, 0, &mut header_bytes[..])?; let header = decode_segment_header(&header_bytes, &hk, page_size)?; validate_header(&header, meta, parent_file_id, page_size)?; let mut footer_bytes = vec![0u8; page_size]; - let footer_read = file.read_at(footer_offset, &mut footer_bytes).await?; - if footer_read != page_size { - return Err(PagedbError::segment_geometry_invalid("footer_read")); - } + read_exact_at_borrowed!(file, footer_offset, &mut footer_bytes[..])?; let (footer, manifest) = { let mut lru = pager.dek_lru().lock(); let cipher = lru.get_or_derive(meta.realm_id, meta.mk_epoch, cipher_id, &master_key)?; @@ -332,9 +325,7 @@ async fn collect_and_decode_index_page( .checked_mul(page_size_u64) .ok_or_else(|| PagedbError::segment_geometry_invalid("index.offset"))?; let mut page = vec![0u8; context.page_size]; - if context.file.read_at(offset, &mut page).await? != context.page_size { - return Err(PagedbError::segment_geometry_invalid("index.read")); - } + read_exact_at_borrowed!(context.file, offset, &mut page[..])?; let (cipher_id, mk_epoch) = extract_page_header_ids(&page)?; if cipher_id != context.cipher_id { return Err(PagedbError::segment_metadata_mismatch( diff --git a/src/segment/reader.rs b/src/segment/reader.rs index db9a6a6..cb92528 100644 --- a/src/segment/reader.rs +++ b/src/segment/reader.rs @@ -14,7 +14,7 @@ use crate::pager::Pager; use crate::pager::format::data_page::{body, extract_page_header_ids, open_data_page}; use crate::pager::format::page_kind::PageKind; use crate::vfs::types::OpenMode; -use crate::vfs::{Vfs, VfsFile}; +use crate::vfs::{Vfs, VfsFile, read_exact_at_borrowed}; use super::authenticated_metadata::{ ExtentIndexDecodeContext, authenticate_segment_metadata, decode_extent_index, @@ -67,11 +67,14 @@ impl SegmentReader { ) -> Result { let page_size = pager.page_size(); let live = live_path(&catalog_meta.segment_id); - let file = pager - .vfs() - .open(&live, OpenMode::Read) - .await - .map_err(|_| PagedbError::NotFound)?; + let file = match pager.vfs().open(&live, OpenMode::Read).await { + Ok(file) => file, + Err(PagedbError::Io(error)) if error.kind() == std::io::ErrorKind::NotFound => { + return Err(PagedbError::NotFound); + } + Err(PagedbError::NotFound) => return Err(PagedbError::NotFound), + Err(error) => return Err(error), + }; Self::finish_open( pager, catalog_meta, @@ -225,10 +228,7 @@ impl SegmentReader { .checked_mul(page_size) .ok_or_else(|| PagedbError::arithmetic_overflow("segment page offset"))?; let mut buf = vec![0u8; self.page_size]; - let n = self.file.read_at(offset, &mut buf).await?; - if n < self.page_size { - return Err(PagedbError::NotFound); - } + read_exact_at_borrowed!(self.file, offset, &mut buf[..])?; // Try each segment page kind; AAD binding rejects wrong ones. let try_kinds = [ diff --git a/src/segment/writer.rs b/src/segment/writer.rs index b18d301..99c914b 100644 --- a/src/segment/writer.rs +++ b/src/segment/writer.rs @@ -16,7 +16,7 @@ use crate::pager::format::segment_footer::{ }; use crate::pager::format::structural_header::{SegmentHeaderFields, encode_segment_header}; use crate::vfs::types::OpenMode; -use crate::vfs::{Vfs, VfsFile}; +use crate::vfs::{Vfs, VfsFile, write_all_at}; use crate::{RealmId, Result}; use tracing; @@ -89,7 +89,7 @@ impl SegmentWriter { flags: 0, }; let header_bytes = encode_segment_header(&header_fields, &hk, page_size)?; - file.write_at(0, &header_bytes).await?; + write_all_at(&mut file, 0, &header_bytes).await?; let total_bytes = u64::try_from(page_size) .map_err(|_| PagedbError::Io(std::io::Error::other("page_size > u64")))?; @@ -185,7 +185,7 @@ impl SegmentWriter { let offset = page_id .checked_mul(self.page_size as u64) .ok_or_else(|| PagedbError::Io(std::io::Error::other("offset overflow")))?; - self.file.write_at(offset, &buf).await?; + write_all_at(&mut self.file, offset, &buf).await?; self.next_page_id += 1; self.total_bytes = self.total_bytes.saturating_add(self.page_size as u64); Ok(page_id) @@ -327,7 +327,7 @@ impl SegmentWriter { let offset = page_id .checked_mul(self.page_size as u64) .ok_or_else(|| PagedbError::Io(std::io::Error::other("offset overflow")))?; - self.file.write_at(offset, &buf).await?; + write_all_at(&mut self.file, offset, &buf).await?; self.next_page_id += 1; self.total_bytes = self.total_bytes.saturating_add(self.page_size as u64); } @@ -362,7 +362,7 @@ impl SegmentWriter { let offset = footer_page_id .checked_mul(self.page_size as u64) .ok_or_else(|| PagedbError::Io(std::io::Error::other("offset overflow")))?; - self.file.write_at(offset, &footer_bytes).await?; + write_all_at(&mut self.file, offset, &footer_bytes).await?; self.file.sync().await?; self.pager.vfs().sync_dir("seg/.staging").await?; diff --git a/src/txn/db/catalog.rs b/src/txn/db/catalog.rs index d3786a4..481d99a 100644 --- a/src/txn/db/catalog.rs +++ b/src/txn/db/catalog.rs @@ -15,6 +15,11 @@ use super::core::{ encode_root_ref, }; +/// Named-counter rows read per batch while validating them at open. Rows are a +/// fixed-width authenticated value, so this is a few KiB resident regardless of +/// how many counters the embedder has named. +const COUNTER_ROW_BATCH: usize = 512; + impl Db { /// The oldest commit id still retained in the commit-history index, or /// `None` when history is disabled or the index is empty. Pages reachable @@ -44,14 +49,62 @@ impl Db { let Some(key) = hist.first_key().await? else { return Ok(None); }; - if key.len() < 8 { - return Ok(None); + if key.len() != 8 { + return Err(PagedbError::catalog_row_invalid("commit_history.key")); } let mut b = [0u8; 8]; b.copy_from_slice(&key[..8]); Ok(Some(u64::from_be_bytes(b))) } + /// Authenticate and decode every persisted named-counter row during open. + /// + /// Named counters are already atomic with catalog-root publication, so + /// recovery validates their encoding but never rewrites their values. + /// + /// Rows are streamed in bounded batches: how many counters an embedder has + /// named is its business, and open must not size an allocation by it. + pub(super) async fn validate_counter_rows( + &self, + catalog_root_page_id: u64, + next_page_id: u64, + ) -> Result<()> { + if catalog_root_page_id == 0 { + return Ok(()); + } + + let prefix = [crate::catalog::codec::CatalogRowKind::Counter as u8]; + let tree = BTree::open( + self.pager.clone(), + self.realm_id, + catalog_root_page_id, + next_page_id, + self.page_size, + ); + let mut cursor: Vec = prefix.to_vec(); + loop { + let batch = tree + .collect_prefix_batch_from(&prefix, &cursor, COUNTER_ROW_BATCH) + .await?; + let Some((last_key, _)) = batch.last() else { + return Ok(()); + }; + cursor.clear(); + cursor.extend_from_slice(last_key); + // The exact successor of `last_key` in the key ordering: resume + // strictly past the row just validated without re-reading it. + cursor.push(0); + let exhausted = batch.len() < COUNTER_ROW_BATCH; + + for (_key, value) in &batch { + Catalog::decode_counter(value)?; + } + if exhausted { + return Ok(()); + } + } + } + /// Write per-realm quota caps into the catalog B+ tree and persist the /// updated catalog root to the A/B header. pub async fn set_realm_quotas(&self, realm: RealmId, quotas: RealmQuotas) -> Result<()> { @@ -307,3 +360,73 @@ impl Db { Ok(freed) } } + +#[cfg(test)] +mod tests { + use crate::vfs::memory::MemVfs; + use crate::{Db, PagedbError, RealmId}; + + use super::*; + + const PAGE: usize = 4096; + const REALM: RealmId = RealmId::new([0xA7; 16]); + + #[tokio::test(flavor = "current_thread")] + async fn counter_recovery_surfaces_malformed_counter_row() { + let db = Db::open_internal(MemVfs::new(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + { + let mut txn = db.begin_write().await.unwrap(); + let mut counter = txn.counter("bad-counter").unwrap(); + counter.set(5).await.unwrap(); + drop(counter); + txn.commit().await.unwrap(); + } + + let (catalog_root, next_page_id) = { + let state = db.writer.lock().await; + (state.catalog_root_page_id, state.next_page_id) + }; + let mut tree = BTree::open( + db.pager.clone(), + db.realm_id, + catalog_root, + next_page_id, + db.page_size, + ); + tree.put(&Catalog::counter_key(&[0xFF]).unwrap(), b"bad") + .await + .unwrap(); + tree.flush().await.unwrap(); + + let err = db + .validate_counter_rows(tree.root_page_id(), tree.next_page_id()) + .await + .expect_err("malformed counter row must surface during recovery validation"); + assert!(matches!(err, PagedbError::Corruption(_))); + } + + #[tokio::test(flavor = "current_thread")] + async fn oldest_retained_history_commit_surfaces_malformed_history_key() { + for malformed_key in [b"x".as_slice(), b"123456789".as_slice()] { + let db = Db::open_internal(MemVfs::new(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let next_page_id = db.writer.lock().await.next_page_id; + let mut history = + BTree::open(db.pager.clone(), db.realm_id, 0, next_page_id, db.page_size); + history + .put(malformed_key, b"malformed history") + .await + .unwrap(); + history.flush().await.unwrap(); + + let err = db + .oldest_retained_history_commit(history.root_page_id(), history.next_page_id()) + .await + .expect_err("malformed history key must surface"); + assert!(matches!(err, PagedbError::Corruption(_))); + } + } +} diff --git a/src/txn/db/misc.rs b/src/txn/db/misc.rs index 7daf362..b33277c 100644 --- a/src/txn/db/misc.rs +++ b/src/txn/db/misc.rs @@ -5,6 +5,7 @@ use crate::Result; use crate::btree::BTree; use crate::catalog::codec::{Catalog, CatalogRowKind}; use crate::crypto::SecretKey; +use crate::errors::PagedbError; use crate::observability::DbStats; use crate::vfs::types::OpenMode; use crate::vfs::{Vfs, VfsFile}; @@ -13,6 +14,11 @@ use std::sync::atomic::Ordering as AtOrd; use super::super::mode::DbMode; use super::core::Db; +/// Segment catalog rows read per batch while aggregating `stats()`. Rows are a +/// fixed-width authenticated value, so this is a few KiB resident regardless of +/// how many segments a realm has linked. +const SEGMENT_ROW_BATCH: usize = 512; + impl Db { /// Return the mode this handle was opened with. pub fn mode(&self) -> DbMode { @@ -112,20 +118,17 @@ impl Db { let latest_commit_id = snapshot.commit_id; // Durable free-list depth (chain rooted at the header's free_list_root). - let free_list_pending_entries = - crate::pager::freelist::read_chain(&self.pager, self.realm_id, free_list_root) - .await - .map_or(0, |(entries, _)| entries.len() as u64); + let (free_list_entries, _) = + crate::pager::freelist::read_chain(&self.pager, self.realm_id, free_list_root).await?; + let free_list_pending_entries = free_list_entries.len() as u64; // Main database file size. - let main_db_bytes = match self + let main_db_bytes = self .vfs .open(&self.main_db_path, crate::vfs::types::OpenMode::Read) - .await - { - Ok(f) => f.len().await.unwrap_or(0), - Err(_) => 0, - }; + .await? + .len() + .await?; // Buffer pool stats from cache. let buffer_pool_pages = { self.pager.inner.buffer_pool.lock().len() as u64 }; @@ -160,14 +163,37 @@ impl Db { self.page_size, ); - let seg_start = vec![CatalogRowKind::Segment as u8]; - let seg_rows = tree.scan_prefix(&seg_start).await.unwrap_or_default(); - let seg_count = u32::try_from(seg_rows.len()).unwrap_or(u32::MAX); - let seg_bytes: u64 = seg_rows - .iter() - .filter_map(|(_k, v)| Catalog::decode_segment_meta(v).ok()) - .map(|m| m.total_bytes) - .sum(); + // Streamed in bounded batches: how many segments a realm has linked + // is the embedder's business, and a metrics call must not size an + // allocation by it. + let seg_prefix = [CatalogRowKind::Segment as u8]; + let mut cursor: Vec = seg_prefix.to_vec(); + let mut seg_count = 0u32; + let mut seg_bytes = 0u64; + loop { + let batch = tree + .collect_prefix_batch_from(&seg_prefix, &cursor, SEGMENT_ROW_BATCH) + .await?; + let Some((last_key, _)) = batch.last() else { + break; + }; + cursor.clear(); + cursor.extend_from_slice(last_key); + cursor.push(0); + let exhausted = batch.len() < SEGMENT_ROW_BATCH; + + for (_key, value) in &batch { + seg_count = seg_count.saturating_add(1); + seg_bytes = seg_bytes + .checked_add(Catalog::decode_segment_meta(value)?.total_bytes) + .ok_or_else(|| { + PagedbError::arithmetic_overflow("stats.segments_total_bytes") + })?; + } + if exhausted { + break; + } + } (seg_count, seg_bytes) }; @@ -192,3 +218,66 @@ impl Db { }) } } + +#[cfg(test)] +mod tests { + use crate::btree::BTree; + use crate::catalog::codec::Catalog; + use crate::vfs::memory::MemVfs; + use crate::{Db, PagedbError, RealmId, SegmentKind, SegmentPageKind}; + + const PAGE: usize = 4096; + const REALM: RealmId = RealmId::new([0xA5; 16]); + + #[tokio::test(flavor = "current_thread")] + async fn stats_surfaces_malformed_segment_catalog_row() { + let db = Db::open_internal(MemVfs::new(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let mut segment = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + segment + .append_page(SegmentPageKind::Data, b"stats") + .await + .unwrap(); + let meta = segment.seal().await.unwrap(); + { + let mut txn = db.begin_write().await.unwrap(); + txn.link_segment("good", &meta).await.unwrap(); + txn.commit().await.unwrap(); + } + + let (catalog_root, next_page_id) = { + let state = db.writer.lock().await; + (state.catalog_root_page_id, state.next_page_id) + }; + let mut tree = BTree::open( + db.pager.clone(), + db.realm_id, + catalog_root, + next_page_id, + db.page_size, + ); + tree.put( + &Catalog::segment_key(REALM, b"bad").unwrap(), + b"not a segment meta", + ) + .await + .unwrap(); + tree.flush().await.unwrap(); + { + let mut state = db.writer.lock().await; + state.catalog_root_page_id = tree.root_page_id(); + state.next_page_id = state.next_page_id.max(tree.next_page_id()); + db.publish_snapshot(&state); + } + + let err = db + .stats() + .await + .expect_err("malformed segment metadata must surface"); + assert!(matches!(err, PagedbError::Corruption(_))); + } +} diff --git a/src/txn/db/open/existing.rs b/src/txn/db/open/existing.rs index 10b957c..7a8078e 100644 --- a/src/txn/db/open/existing.rs +++ b/src/txn/db/open/existing.rs @@ -10,9 +10,10 @@ use crate::crypto::{CipherId, SecretKey}; use crate::errors::PagedbError; use crate::options::OpenOptions; use crate::pager::header::ActiveSlot; +use crate::pager::header::read_header_slot; use crate::pager::structural_header::MainDbHeaderFields; use crate::pager::{Pager, PagerConfig}; -use crate::vfs::{Vfs, VfsFile}; +use crate::vfs::Vfs; use crate::{RealmId, Result}; use super::super::super::mode::DbMode; @@ -104,13 +105,13 @@ impl Db { let capabilities = mode.open_capabilities(); let file_mode = capabilities.main_db_open_mode(); let read_only = capabilities.read_only_file_access(); - let f = vfs.open(&main_db_path, file_mode).await?; + let mut f = vfs.open(&main_db_path, file_mode).await?; let mut buf_a = vec![0u8; page_size]; let mut buf_b = vec![0u8; page_size]; - let _ = f.read_at(0, &mut buf_a).await?; + read_header_slot(&mut f, 0, &mut buf_a).await?; let page_size_u64 = u64::try_from(page_size) .map_err(|_| PagedbError::Io(std::io::Error::other("page_size > u64")))?; - let _ = f.read_at(page_size_u64, &mut buf_b).await?; + read_header_slot(&mut f, page_size_u64, &mut buf_b).await?; drop(f); let try_decode = |buf: &[u8]| -> Option<(MainDbHeaderFields, bool)> { diff --git a/src/txn/db/open/modes.rs b/src/txn/db/open/modes.rs index c1aff67..06d6aee 100644 --- a/src/txn/db/open/modes.rs +++ b/src/txn/db/open/modes.rs @@ -257,18 +257,26 @@ impl Db { .vfs .lock_exclusive(FROZEN_READERS_LOCK_PATH) .await - .map_err(|_| { - crate::diag::lock_rejected("follower", FROZEN_READERS_LOCK_PATH, "readers_present"); - PagedbError::ReadersPresent + .map_err(|error| { + map_lock_contention(error, || { + crate::diag::lock_rejected( + "follower", + FROZEN_READERS_LOCK_PATH, + "readers_present", + ); + PagedbError::ReadersPresent + }) })?; drop(frozen_probe); let writer_lock = self .vfs .lock_exclusive(WRITER_LOCK_PATH) .await - .map_err(|_| { - crate::diag::lock_rejected("follower", WRITER_LOCK_PATH, "already_open"); - PagedbError::AlreadyOpen + .map_err(|error| { + map_lock_contention(error, || { + crate::diag::lock_rejected("follower", WRITER_LOCK_PATH, "already_open"); + PagedbError::AlreadyOpen + }) })?; drop(acquisition); crate::diag::lock_acquired("follower", WRITER_LOCK_PATH); @@ -301,22 +309,33 @@ pub(super) async fn acquire_interlocking_mode_lock( let frozen_probe = vfs.lock_exclusive(FROZEN_READERS_LOCK_PATH) .await - .map_err(|_| { - crate::diag::lock_rejected( - "writer", - FROZEN_READERS_LOCK_PATH, - "readers_present", - ); - PagedbError::ReadersPresent + .map_err(|error| { + map_lock_contention(error, || { + crate::diag::lock_rejected( + "writer", + FROZEN_READERS_LOCK_PATH, + "readers_present", + ); + PagedbError::ReadersPresent + }) })?; drop(frozen_probe); acquire_long_lived_lock(vfs, LongLivedLock::Writer, locks).await } LongLivedLock::FrozenReader => { - let writer_probe = vfs.lock_exclusive(WRITER_LOCK_PATH).await.map_err(|_| { - crate::diag::lock_rejected("frozen_reader", WRITER_LOCK_PATH, "writer_present"); - PagedbError::WriterPresent - })?; + let writer_probe = vfs + .lock_exclusive(WRITER_LOCK_PATH) + .await + .map_err(|error| { + map_lock_contention(error, || { + crate::diag::lock_rejected( + "frozen_reader", + WRITER_LOCK_PATH, + "writer_present", + ); + PagedbError::WriterPresent + }) + })?; drop(writer_probe); acquire_long_lived_lock(vfs, LongLivedLock::FrozenReader, locks).await } @@ -332,27 +351,50 @@ async fn acquire_long_lived_lock( locks: &mut Vec, ) -> Result<()> { let handle = match lock { - LongLivedLock::Writer => vfs.lock_exclusive(WRITER_LOCK_PATH).await.map_err(|_| { - crate::diag::lock_rejected("writer", WRITER_LOCK_PATH, "already_open"); - PagedbError::AlreadyOpen - })?, + LongLivedLock::Writer => vfs + .lock_exclusive(WRITER_LOCK_PATH) + .await + .map_err(|error| { + map_lock_contention(error, || { + crate::diag::lock_rejected("writer", WRITER_LOCK_PATH, "already_open"); + PagedbError::AlreadyOpen + }) + })?, LongLivedLock::FrozenReader => { vfs.lock_shared(FROZEN_READERS_LOCK_PATH) .await - .map_err(|_| { - crate::diag::lock_rejected( - "frozen_reader", - FROZEN_READERS_LOCK_PATH, - "already_locked", - ); - PagedbError::AlreadyLocked + .map_err(|error| { + map_lock_contention(error, || { + crate::diag::lock_rejected( + "frozen_reader", + FROZEN_READERS_LOCK_PATH, + "already_locked", + ); + PagedbError::AlreadyLocked + }) })? } - LongLivedLock::Observer => vfs.lock_shared(OBSERVERS_LOCK_PATH).await.map_err(|_| { - crate::diag::lock_rejected("observer", OBSERVERS_LOCK_PATH, "already_locked"); - PagedbError::AlreadyLocked - })?, + LongLivedLock::Observer => vfs + .lock_shared(OBSERVERS_LOCK_PATH) + .await + .map_err(|error| { + map_lock_contention(error, || { + crate::diag::lock_rejected("observer", OBSERVERS_LOCK_PATH, "already_locked"); + PagedbError::AlreadyLocked + }) + })?, }; locks.push(handle); Ok(()) } + +fn map_lock_contention( + error: PagedbError, + on_contention: impl FnOnce() -> PagedbError, +) -> PagedbError { + if matches!(error, PagedbError::AlreadyLocked) { + on_contention() + } else { + error + } +} diff --git a/src/txn/db/open/recovery.rs b/src/txn/db/open/recovery.rs index 218ac73..338a4ef 100644 --- a/src/txn/db/open/recovery.rs +++ b/src/txn/db/open/recovery.rs @@ -72,5 +72,8 @@ pub(super) async fn recover_open_state( .await?; } + db.validate_counter_rows(catalog_root_page_id, next_page_id) + .await?; + Ok(()) } diff --git a/src/txn/db/reader.rs b/src/txn/db/reader.rs index 3bb6c70..200925a 100644 --- a/src/txn/db/reader.rs +++ b/src/txn/db/reader.rs @@ -262,17 +262,18 @@ impl Db { false, )) } else { - // Find the oldest available: scan the whole history tree. An - // exclusive `u64::MAX` upper bound would hide the newest commit - // once ids reach the top of the range. - let oldest = tree.collect_all().await?.into_iter().next().map_or( - CommitId(latest_commit_id), - |(k, _)| { - let mut b = [0u8; 8]; - b.copy_from_slice(&k[..8]); - CommitId(u64::from_be_bytes(b)) - }, - ); + // Only the leftmost key is needed. This avoids scanning retained + // history while still validating its fixed-width key encoding. + let oldest = if let Some(key) = tree.first_key().await? { + if key.len() != 8 { + return Err(PagedbError::catalog_row_invalid("commit_history.key")); + } + let mut bytes = [0u8; 8]; + bytes.copy_from_slice(&key[..8]); + CommitId(u64::from_be_bytes(bytes)) + } else { + CommitId(latest_commit_id) + }; Err(PagedbError::CommitGone { commit, oldest_available: oldest, @@ -285,10 +286,12 @@ impl Db { mod tests { use std::sync::Arc; + use crate::btree::BTree; use crate::catalog::codec::SegmentKind; + use crate::errors::PagedbError; use crate::segment::types::SegmentPageKind; use crate::vfs::memory::MemVfs; - use crate::{Db, RealmId}; + use crate::{CommitId, Db, RealmId}; use super::super::core::VisibilityTestHook; @@ -296,6 +299,38 @@ mod tests { const KEK: [u8; 32] = [0xD1; 32]; const REALM: RealmId = RealmId::new([0xD2; 16]); + #[tokio::test(flavor = "current_thread")] + async fn begin_read_at_surfaces_malformed_oldest_history_key() { + for malformed_key in [b"x".as_slice(), b"123456789".as_slice()] { + let db = Db::open_internal(MemVfs::new(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let next_page_id = db.writer.lock().await.next_page_id; + let mut history = + BTree::open(db.pager.clone(), db.realm_id, 0, next_page_id, db.page_size); + history + .put(malformed_key, b"malformed history") + .await + .unwrap(); + history.flush().await.unwrap(); + { + let mut state = db.writer.lock().await; + state.commit_history_root_page_id = history.root_page_id(); + state.next_page_id = state.next_page_id.max(history.next_page_id()); + db.publish_snapshot(&state); + } + + let err = match db.begin_read_at(CommitId::new(42)).await { + Ok(read) => { + drop(read); + panic!("malformed oldest history key must not produce a reader"); + } + Err(err) => err, + }; + assert!(matches!(err, PagedbError::Corruption(_))); + } + } + async fn linked_segment(db: &Db, name: &str, bytes: &[u8]) { let mut writer = db .create_segment(REALM, SegmentKind::Unspecified) diff --git a/src/txn/db/util.rs b/src/txn/db/util.rs index 32a1726..dac4322 100644 --- a/src/txn/db/util.rs +++ b/src/txn/db/util.rs @@ -4,6 +4,7 @@ use crate::Result; use crate::crypto::kdf::{derive_hk, derive_mk}; use crate::errors::PagedbError; +use crate::pager::header::read_header_slot; use crate::vfs::Vfs; pub(super) fn page_size_log2(page_size: usize) -> Result { @@ -37,16 +38,15 @@ pub(super) async fn peek_restore_mode( kek: &[u8; 32], page_size: usize, ) -> Result { - use crate::vfs::VfsFile; use crate::vfs::types::OpenMode; - let f = vfs.open("/main.db", OpenMode::Read).await?; + let mut f = vfs.open("/main.db", OpenMode::Read).await?; let mut buf_a = vec![0u8; page_size]; let mut buf_b = vec![0u8; page_size]; - f.read_at(0, &mut buf_a).await?; + read_header_slot(&mut f, 0, &mut buf_a).await?; let page_size_u64 = u64::try_from(page_size) .map_err(|_| PagedbError::Io(std::io::Error::other("page_size > u64")))?; - f.read_at(page_size_u64, &mut buf_b).await?; + read_header_slot(&mut f, page_size_u64, &mut buf_b).await?; drop(f); for buf in [&buf_a, &buf_b] { diff --git a/src/txn/write/spill.rs b/src/txn/write/spill.rs index fbc13a6..9da2ce7 100644 --- a/src/txn/write/spill.rs +++ b/src/txn/write/spill.rs @@ -7,7 +7,7 @@ use crate::crypto::kdf::derive_spill_key; use crate::crypto::nonce::Nonce; use crate::errors::{PagedbError, QuotaKind}; use crate::vfs::types::OpenMode; -use crate::vfs::{Vfs, VfsFile}; +use crate::vfs::{Vfs, VfsFile, read_exact_at, write_all_at}; use super::txn::WriteTxn; @@ -17,6 +17,9 @@ use super::txn::WriteTxn; pub struct ScratchOffset(u64); /// Metadata for one ciphertext chunk in the per-txn spill scratch file. +/// Length of the AEAD tag appended to every spill frame. +pub(crate) const SPILL_TAG_LEN: usize = 16; + /// Stored in memory only; the tmp file is discarded at commit/abort. #[derive(Clone)] pub(crate) struct SpillSegmentMeta { @@ -24,8 +27,8 @@ pub(crate) struct SpillSegmentMeta { pub offset: u64, /// Length of the original plaintext in bytes. pub plaintext_len: u32, - /// Length of ciphertext body (without the 16-byte tag) in bytes. - /// Total on-disk size for this chunk = `ciphertext_len + 16`. + /// Length of ciphertext body (without the AEAD tag) in bytes. Total + /// on-disk size for this chunk = `ciphertext_len + SPILL_TAG_LEN`. pub ciphertext_len: u32, /// The 12-byte nonce used to encrypt this chunk, stored verbatim so we /// can reconstruct a `Nonce` on read without an additional lookup. @@ -106,8 +109,8 @@ impl SpillScope<'_, '_, V> { .expect("derived above") .encrypt(&nonce, &aad, &mut body)?; - // ciphertext body + 16-byte tag. - let pers_len = body.len() as u64 + 16; + // ciphertext body + AEAD tag. + let pers_len = body.len() as u64 + SPILL_TAG_LEN as u64; let new_total = self.txn.spill_bytes_used.saturating_add(pers_len); if new_total > limit { return Err(PagedbError::quota( @@ -123,13 +126,16 @@ impl SpillScope<'_, '_, V> { let mut file = self.txn.db.vfs.open(&path, OpenMode::CreateOrOpen).await?; let body_offset = self.txn.spill_bytes_used; - let tag_offset = body_offset + body.len() as u64; - file.write_at(body_offset, &body).await?; - file.write_at(tag_offset, &tag).await?; + let ciphertext_len = body.len(); + // Body and tag are one contiguous frame: a handle must never be + // returned for a file holding only half of it. + body.extend_from_slice(&tag); + write_all_at(&mut file, body_offset, &body).await?; file.sync().await?; let plaintext_len = u32::try_from(bytes.len()).map_err(|_| PagedbError::PayloadTooLarge)?; - let ciphertext_len = u32::try_from(body.len()).map_err(|_| PagedbError::PayloadTooLarge)?; + let ciphertext_len = + u32::try_from(ciphertext_len).map_err(|_| PagedbError::PayloadTooLarge)?; self.txn.spill_segments.push(SpillSegmentMeta { offset: body_offset, @@ -158,14 +164,20 @@ impl SpillScope<'_, '_, V> { .clone(); let path = self.txn.spill_path.as_ref().ok_or(PagedbError::NotFound)?; - let file = self.txn.db.vfs.open(path, OpenMode::Read).await?; + let mut file = self.txn.db.vfs.open(path, OpenMode::Read).await?; let body_len = meta.ciphertext_len as usize; - let mut body = vec![0u8; body_len]; - let mut tag = [0u8; 16]; - file.read_at(meta.offset, &mut body).await?; - file.read_at(meta.offset + body_len as u64, &mut tag) - .await?; + // `usize` is 32-bit on wasm32, where a maximal `ciphertext_len` plus + // the tag genuinely overflows it. + let frame_len = body_len + .checked_add(SPILL_TAG_LEN) + .ok_or_else(|| PagedbError::arithmetic_overflow("spill frame length"))?; + let mut frame = vec![0u8; frame_len]; + read_exact_at(&mut file, meta.offset, &mut frame).await?; + // The split is exact by construction, so the tail is the whole tag. + let (body, tag_bytes) = frame.split_at_mut(body_len); + let mut tag = [0u8; SPILL_TAG_LEN]; + tag.copy_from_slice(tag_bytes); let cipher = self.txn.spill_cipher_readonly()?; let nonce = Nonce::from_bytes(meta.nonce_bytes); @@ -182,9 +194,9 @@ impl SpillScope<'_, '_, V> { segment_id: self.txn.db.file_id, }); - cipher.decrypt(&nonce, &aad, &mut body, &tag)?; - body.truncate(meta.plaintext_len as usize); - Ok(body) + cipher.decrypt(&nonce, &aad, body, &tag)?; + frame.truncate(meta.plaintext_len as usize); + Ok(frame) } } diff --git a/src/vfs/mod.rs b/src/vfs/mod.rs index 9ec7edd..e905946 100644 --- a/src/vfs/mod.rs +++ b/src/vfs/mod.rs @@ -39,6 +39,9 @@ pub use iouring::{IouringFile, IouringVfs}; #[cfg(feature = "opfs")] pub use opfs::OpfsVfs; pub use traits::{Vfs, VfsFile}; +pub(crate) use traits::{ + checked_read_progress, read_exact_at, read_exact_at_borrowed, write_all_at, +}; pub use types::{OpenMode, ReadReq, WriteReq}; pub use wasi::WasiVfs; diff --git a/src/vfs/traits.rs b/src/vfs/traits.rs index baff963..104f9c9 100644 --- a/src/vfs/traits.rs +++ b/src/vfs/traits.rs @@ -3,6 +3,7 @@ use std::future::Future; use crate::Result; +use crate::errors::PagedbError; use super::types::{OpenMode, ReadReq, WriteReq}; @@ -88,3 +89,172 @@ pub trait VfsFile: Send { fn is_empty(&self) -> impl Future> + Send; fn supports_direct_io(&self) -> bool; } + +/// Read until `buf` is full, or fail if the backend stops making progress. +/// +/// Takes `&mut F` rather than `&F` deliberately. `read_at` only needs `&self`, +/// but a future holding `&F` across an await is `Send` only when `F: Sync`, +/// and `VfsFile` does not require `Sync`. `&mut F` keeps the future `Send` for +/// every backend. Callers that own their handle use this; those that read +/// through a borrowed handle use [`read_exact_at_borrowed!`], which runs the +/// same loop in place. +#[inline] +pub(crate) async fn read_exact_at( + file: &mut F, + mut offset: u64, + mut buf: &mut [u8], +) -> Result<()> { + while !buf.is_empty() { + let read = file.read_at(offset, buf).await?; + checked_read_progress(&mut offset, read, buf.len())?; + buf = buf.split_at_mut(read).1; + } + Ok(()) +} + +/// Validate one positional read result and advance its offset. +/// +/// A backend may legally satisfy a request in several calls, so a short read is +/// not itself failure. Zero progress is end-of-file; a count above the +/// remaining buffer is a backend contract violation, not an on-disk defect. +#[inline] +pub(crate) fn checked_read_progress(offset: &mut u64, read: usize, remaining: usize) -> Result<()> { + if read == 0 { + return Err(PagedbError::Io(std::io::Error::from( + std::io::ErrorKind::UnexpectedEof, + ))); + } + checked_transfer_progress(offset, read, remaining, "read_at", "positional read offset") +} + +/// Write until `buf` is complete, or fail if the backend reports impossible +/// progress. See [`read_exact_at`] for why this takes `&mut F`. +#[inline] +pub(crate) async fn write_all_at( + file: &mut F, + mut offset: u64, + mut buf: &[u8], +) -> Result<()> { + while !buf.is_empty() { + let written = file.write_at(offset, buf).await?; + if written == 0 { + return Err(PagedbError::Io(std::io::Error::from( + std::io::ErrorKind::WriteZero, + ))); + } + checked_transfer_progress( + &mut offset, + written, + buf.len(), + "write_at", + "positional write offset", + )?; + buf = &buf[written..]; + } + Ok(()) +} + +/// Run [`read_exact_at`]'s loop over a handle the caller only borrows. +/// +/// A function taking `&F` would be `Send` only under `F: Sync`, which +/// `VfsFile` does not require and which would spread as a bound through every +/// segment reader. Expanding in place keeps the borrow inside the caller's own +/// future, so the rule stays in one place without costing a trait bound. +/// +/// `$buf` must be a `&mut [u8]`; the result is a `Result<()>` to be `?`-ed. +macro_rules! read_exact_at_borrowed { + ($file:expr, $offset:expr, $buf:expr $(,)?) => {{ + let mut offset: u64 = $offset; + let mut remaining: &mut [u8] = $buf; + loop { + if remaining.is_empty() { + break Ok(()); + } + match $file.read_at(offset, remaining).await { + Ok(read) => { + if let Err(error) = + $crate::vfs::checked_read_progress(&mut offset, read, remaining.len()) + { + break Err(error); + } + remaining = remaining.split_at_mut(read).1; + } + Err(error) => break Err(error), + } + } + }}; +} +pub(crate) use read_exact_at_borrowed; + +/// Shared progress rule for both directions: reject a count the caller never +/// asked for, then advance the offset without wrapping. +#[inline] +fn checked_transfer_progress( + offset: &mut u64, + transferred: usize, + remaining: usize, + operation: &'static str, + offset_label: &'static str, +) -> Result<()> { + if transferred > remaining { + return Err(PagedbError::vfs_contract_violated( + operation, + "reported more bytes than the caller requested", + )); + } + let transferred = + u64::try_from(transferred).map_err(|_| PagedbError::arithmetic_overflow(offset_label))?; + *offset = offset + .checked_add(transferred) + .ok_or_else(|| PagedbError::arithmetic_overflow(offset_label))?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn a_short_read_advances_by_exactly_what_was_transferred() { + let mut offset = 4096; + checked_read_progress(&mut offset, 100, 4096).unwrap(); + assert_eq!(offset, 4196); + } + + #[test] + fn zero_progress_is_end_of_file_not_a_contract_violation() { + let mut offset = 0; + let err = checked_read_progress(&mut offset, 0, 4096).unwrap_err(); + assert!( + matches!(err, PagedbError::Io(ref io) if io.kind() == std::io::ErrorKind::UnexpectedEof), + "expected UnexpectedEof, got {err:?}" + ); + assert_eq!(offset, 0, "a failed transfer must not advance the offset"); + } + + #[test] + fn a_count_above_the_remaining_buffer_is_a_backend_contract_violation() { + let mut offset = 0; + let err = checked_read_progress(&mut offset, 4097, 4096).unwrap_err(); + assert!( + matches!( + err, + PagedbError::VfsContractViolated { + operation: "read_at", + .. + } + ), + "expected VfsContractViolated, got {err:?}" + ); + } + + #[test] + fn an_offset_that_would_wrap_is_reported_as_overflow() { + let mut offset = u64::MAX; + let err = checked_read_progress(&mut offset, 1, 4096).unwrap_err(); + assert!( + matches!(err, PagedbError::ArithmeticOverflow { .. }), + "expected ArithmeticOverflow, got {err:?}" + ); + } +} diff --git a/tests/durability/metadata_errors.rs b/tests/durability/metadata_errors.rs new file mode 100644 index 0000000..b80d440 --- /dev/null +++ b/tests/durability/metadata_errors.rs @@ -0,0 +1,1156 @@ +//! Durable metadata error propagation and positional-I/O completion. +//! +//! Metadata publication depends on `rename`/`sync_dir`, while authenticated +//! pages may be transferred by backends that legally complete positional I/O +//! in several calls. These tests verify that transient failures are surfaced, +//! partial transfers finish before success is reported, and retries preserve +//! usable state. + +use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; +use std::sync::{Arc, Mutex}; + +use pagedb::options::{OpenOptions, RetainPolicy}; +use pagedb::recovery::{JournalAction, execute_journal_actions}; +use pagedb::vfs::memory::MemVfs; +use pagedb::vfs::traits::{Vfs, VfsFile}; +use pagedb::vfs::types::{OpenMode, ReadReq, WriteReq}; +use pagedb::{Db, PagedbError, RealmId, Result, SegmentKind, SegmentPageKind}; + +const PAGE: usize = 4096; +const REALM: RealmId = RealmId::new([1u8; 16]); + +/// MemVfs wrapper that injects one-shot metadata and main.db I/O failures. +#[derive(Clone)] +struct FailOnceVfs { + inner: MemVfs, + fail_live_segment_open_path: Arc>>, + fail_main_db_read_at: Arc, + fail_main_db_read_open: Arc, + fail_main_db_write_at: Arc, + fail_rename: Arc, + fail_sync_dir: Arc, + fail_tombstone_list_dir: Arc, + fail_tombstone_len: Arc, + fail_tombstone_remove: Arc, + fail_writer_lock: Arc, + main_db_write_at_count: Arc, + short_write_at: Arc>>, + short_read_at: Arc>>, + short_main_db_read_at_or_after: Arc>>, +} + +impl FailOnceVfs { + fn new(inner: MemVfs) -> Self { + FailOnceVfs { + inner, + fail_live_segment_open_path: Arc::new(Mutex::new(None)), + fail_main_db_read_at: Arc::new(AtomicBool::new(false)), + fail_main_db_read_open: Arc::new(AtomicBool::new(false)), + fail_main_db_write_at: Arc::new(AtomicBool::new(false)), + fail_rename: Arc::new(AtomicBool::new(false)), + fail_sync_dir: Arc::new(AtomicBool::new(false)), + fail_tombstone_list_dir: Arc::new(AtomicBool::new(false)), + fail_tombstone_len: Arc::new(AtomicBool::new(false)), + fail_tombstone_remove: Arc::new(AtomicBool::new(false)), + fail_writer_lock: Arc::new(AtomicBool::new(false)), + main_db_write_at_count: Arc::new(AtomicUsize::new(0)), + short_write_at: Arc::new(Mutex::new(None)), + short_read_at: Arc::new(Mutex::new(None)), + short_main_db_read_at_or_after: Arc::new(Mutex::new(None)), + } + } + + fn injected() -> PagedbError { + PagedbError::Io(std::io::Error::other("injected metadata fault")) + } +} + +struct FailFile { + inner: F, + path: String, + fail_main_db_read_at: Arc, + fail_main_db_write_at: Arc, + fail_tombstone_len: Arc, + main_db_write_at_count: Arc, + short_write_at: Arc>>, + short_read_at: Arc>>, + short_main_db_read_at_or_after: Arc>>, +} + +impl Vfs for FailOnceVfs { + type File = FailFile<::File>; + type LockHandle = ::LockHandle; + + async fn open(&self, path: &str, mode: OpenMode) -> Result { + let fail_live_segment_open = if matches!(mode, OpenMode::Read) { + let mut guard = self.fail_live_segment_open_path.lock().unwrap(); + if guard.as_deref() == Some(path) { + guard.take(); + true + } else { + false + } + } else { + false + }; + if fail_live_segment_open { + return Err(PagedbError::Io(std::io::Error::from( + std::io::ErrorKind::PermissionDenied, + ))); + } + if self.fail_main_db_read_open.swap(false, Ordering::SeqCst) + && matches!(mode, OpenMode::Read) + && path == "/main.db" + { + return Err(PagedbError::Io(std::io::Error::from( + std::io::ErrorKind::PermissionDenied, + ))); + } + Ok(FailFile { + inner: self.inner.open(path, mode).await?, + path: path.to_string(), + fail_main_db_read_at: self.fail_main_db_read_at.clone(), + fail_main_db_write_at: self.fail_main_db_write_at.clone(), + fail_tombstone_len: self.fail_tombstone_len.clone(), + main_db_write_at_count: self.main_db_write_at_count.clone(), + short_write_at: self.short_write_at.clone(), + short_read_at: self.short_read_at.clone(), + short_main_db_read_at_or_after: self.short_main_db_read_at_or_after.clone(), + }) + } + async fn remove(&self, path: &str) -> Result<()> { + if self.fail_tombstone_remove.swap(false, Ordering::SeqCst) + && path.starts_with("seg/.tombstone/") + { + return Err(PagedbError::Io(std::io::Error::from( + std::io::ErrorKind::PermissionDenied, + ))); + } + self.inner.remove(path).await + } + async fn rename(&self, from: &str, to: &str) -> Result<()> { + if self.fail_rename.swap(false, Ordering::SeqCst) { + return Err(Self::injected()); + } + self.inner.rename(from, to).await + } + async fn list_dir(&self, path: &str) -> Result> { + if self.fail_tombstone_list_dir.swap(false, Ordering::SeqCst) && path == "seg/.tombstone" { + return Err(PagedbError::Io(std::io::Error::from( + std::io::ErrorKind::PermissionDenied, + ))); + } + self.inner.list_dir(path).await + } + async fn mkdir_all(&self, path: &str) -> Result<()> { + self.inner.mkdir_all(path).await + } + async fn sync_dir(&self, path: &str) -> Result<()> { + if self.fail_sync_dir.swap(false, Ordering::SeqCst) { + return Err(Self::injected()); + } + self.inner.sync_dir(path).await + } + async fn lock_exclusive(&self, path: &str) -> Result { + if path == ".writer.lock" && self.fail_writer_lock.swap(false, Ordering::SeqCst) { + return Err(PagedbError::Io(std::io::Error::from( + std::io::ErrorKind::PermissionDenied, + ))); + } + self.inner.lock_exclusive(path).await + } + async fn lock_shared(&self, path: &str) -> Result { + self.inner.lock_shared(path).await + } +} + +impl VfsFile for FailFile { + async fn read_at(&self, offset: u64, buf: &mut [u8]) -> Result { + if self.fail_main_db_read_at.swap(false, Ordering::SeqCst) && self.path == "/main.db" { + return Err(PagedbError::Io(std::io::Error::from( + std::io::ErrorKind::PermissionDenied, + ))); + } + let short_read = { + let mut guard = self.short_read_at.lock().unwrap(); + if guard.as_ref() == Some(&(self.path.clone(), offset)) { + guard.take(); + true + } else { + false + } + }; + let short_main_db_read = { + let mut guard = self.short_main_db_read_at_or_after.lock().unwrap(); + if self.path == "/main.db" && guard.is_some_and(|min_offset| offset >= min_offset) { + guard.take(); + true + } else { + false + } + }; + if short_read || short_main_db_read { + let n = buf.len().saturating_sub(1); + if n == 0 { + return Ok(0); + } + return self.inner.read_at(offset, &mut buf[..n]).await; + } + self.inner.read_at(offset, buf).await + } + async fn read_at_vectored(&self, reqs: &mut [ReadReq<'_>]) -> Result<()> { + self.inner.read_at_vectored(reqs).await + } + async fn write_at(&mut self, offset: u64, buf: &[u8]) -> Result { + if self.path == "/main.db" { + self.main_db_write_at_count.fetch_add(1, Ordering::SeqCst); + } + if self.fail_main_db_write_at.swap(false, Ordering::SeqCst) && self.path == "/main.db" { + return Err(PagedbError::Io(std::io::Error::from( + std::io::ErrorKind::PermissionDenied, + ))); + } + let short_write = { + let mut guard = self.short_write_at.lock().unwrap(); + if guard.as_ref() == Some(&(self.path.clone(), offset)) { + guard.take(); + true + } else { + false + } + }; + if short_write { + let n = buf.len().saturating_sub(1); + if n == 0 { + return Ok(0); + } + return self.inner.write_at(offset, &buf[..n]).await; + } + self.inner.write_at(offset, buf).await + } + async fn write_at_vectored(&mut self, reqs: &[WriteReq<'_>]) -> Result<()> { + self.inner.write_at_vectored(reqs).await + } + async fn sync(&mut self) -> Result<()> { + self.inner.sync().await + } + async fn truncate(&mut self, len: u64) -> Result<()> { + self.inner.truncate(len).await + } + async fn len(&self) -> Result { + if self.fail_tombstone_len.swap(false, Ordering::SeqCst) + && self.path.starts_with("seg/.tombstone/") + { + return Err(PagedbError::Io(std::io::Error::from( + std::io::ErrorKind::PermissionDenied, + ))); + } + self.inner.len().await + } + async fn is_empty(&self) -> Result { + self.inner.is_empty().await + } + fn supports_direct_io(&self) -> bool { + self.inner.supports_direct_io() + } +} + +#[tokio::test(flavor = "current_thread")] +async fn journal_actions_surface_rename_error_then_retry_succeeds() { + let mem = MemVfs::new(); + mem.mkdir_all("seg/.staging").await.unwrap(); + let id = [7u8; 16]; + let staging = format!("seg/.staging/{}", hex(&id)); + let mut f = mem.open(&staging, OpenMode::CreateNew).await.unwrap(); + f.write_at(0, b"payload").await.unwrap(); + drop(f); + + let vfs = FailOnceVfs::new(mem.clone()); + let actions = vec![JournalAction::Promote { segment_id: id }]; + + vfs.fail_rename.store(true, Ordering::SeqCst); + let err = execute_journal_actions(&vfs, &actions) + .await + .expect_err("failed rename must surface, not be swallowed"); + assert!(matches!(err, PagedbError::Io(_))); + + // Retry (fault cleared) completes the action. + execute_journal_actions(&vfs, &actions).await.unwrap(); + let live = format!("seg/{}", hex(&id)); + assert!(mem.open(&live, OpenMode::Read).await.is_ok()); +} + +#[tokio::test(flavor = "current_thread")] +async fn journal_actions_tolerate_already_completed_rename() { + // Source missing = action already ran; replay must be a no-op success. + let mem = MemVfs::new(); + mem.mkdir_all("seg").await.unwrap(); + let id = [8u8; 16]; + let live = format!("seg/{}", hex(&id)); + mem.open(&live, OpenMode::CreateNew).await.unwrap(); + + let actions = vec![JournalAction::Promote { segment_id: id }]; + execute_journal_actions(&mem, &actions).await.unwrap(); + assert!(mem.open(&live, OpenMode::Read).await.is_ok()); +} + +#[tokio::test(flavor = "current_thread")] +async fn gc_sync_dir_error_surfaces_then_retry_succeeds() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let options = OpenOptions::default().with_commit_history_retain(RetainPolicy::Disabled); + let db = Db::open_internal_with_options(vfs.clone(), [9u8; 32], PAGE, REALM, options) + .await + .unwrap(); + let mut w = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + w.append_page(SegmentPageKind::Data, b"x").await.unwrap(); + let m = w.seal().await.unwrap(); + { + let mut t = db.begin_write().await.unwrap(); + t.link_segment("dead", &m).await.unwrap(); + t.commit().await.unwrap(); + } + { + let mut t = db.begin_write().await.unwrap(); + t.unlink_segment("dead").await.unwrap(); + t.commit().await.unwrap(); + } + + // First gc_now hits the injected sync_dir failure and must SURFACE it. + vfs.fail_sync_dir.store(true, Ordering::SeqCst); + let err = db + .gc_now() + .await + .expect_err("gc sync_dir failure must surface"); + assert!(matches!(err, PagedbError::Io(_))); + + // Retry succeeds and the segment is gone. + db.gc_now().await.unwrap(); + let err = db.open_segment(REALM, "dead").await.err().unwrap(); + assert!(matches!(err, PagedbError::NotFound)); +} + +#[tokio::test(flavor = "current_thread")] +async fn journal_actions_surface_sync_dir_error_then_retry_succeeds() { + let mem = MemVfs::new(); + mem.mkdir_all("seg/.staging").await.unwrap(); + let id = [9u8; 16]; + let staging = format!("seg/.staging/{}", hex(&id)); + let mut f = mem.open(&staging, OpenMode::CreateNew).await.unwrap(); + f.write_at(0, b"payload").await.unwrap(); + drop(f); + + let vfs = FailOnceVfs::new(mem.clone()); + let actions = vec![JournalAction::Promote { segment_id: id }]; + + vfs.fail_sync_dir.store(true, Ordering::SeqCst); + let err = execute_journal_actions(&vfs, &actions) + .await + .expect_err("journal action sync_dir failures must surface before the journal is cleared"); + assert!(matches!(err, PagedbError::Io(_))); + + // Retry must complete even if the prior attempt already performed the + // rename but failed before making the directory entry durable. + execute_journal_actions(&vfs, &actions).await.unwrap(); + let live = format!("seg/{}", hex(&id)); + assert!(mem.open(&live, OpenMode::Read).await.is_ok()); +} + +#[tokio::test(flavor = "current_thread")] +async fn reconcile_surfaces_live_segment_open_error_then_retry_succeeds() { + let vfs = FailOnceVfs::new(MemVfs::new()); + { + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let mut w = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + w.append_page(SegmentPageKind::Data, b"recoverable") + .await + .unwrap(); + let m = w.seal().await.unwrap(); + let live_path = format!("seg/{}", hex(&m.segment_id)); + let mut t = db.begin_write().await.unwrap(); + t.link_segment("open-error", &m).await.unwrap(); + t.commit().await.unwrap(); + *vfs.fail_live_segment_open_path.lock().unwrap() = Some(live_path); + } + + let err = match Db::open_existing(vfs.clone(), [9u8; 32], PAGE, REALM).await { + Ok(_) => panic!("reconcile must surface live segment open I/O errors"), + Err(err) => err, + }; + assert!( + matches!(err, PagedbError::Io(ref io) if io.kind() == std::io::ErrorKind::PermissionDenied), + "expected live segment PermissionDenied to surface, got {err:?}" + ); + + let db = Db::open_existing(vfs, [9u8; 32], PAGE, REALM) + .await + .expect("retry after transient open fault must succeed"); + let reader = db.open_segment(REALM, "open-error").await.unwrap(); + let page = reader.read_page(1).await.unwrap(); + assert!(page.starts_with(b"recoverable")); +} + +#[tokio::test(flavor = "current_thread")] +async fn reconcile_retries_short_live_segment_footer_read() { + let vfs = FailOnceVfs::new(MemVfs::new()); + { + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let mut w = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + w.append_page(SegmentPageKind::Data, b"reconcile-short-footer") + .await + .unwrap(); + let m = w.seal().await.unwrap(); + let live_path = format!("seg/{}", hex(&m.segment_id)); + let footer_offset = (m.page_count - 1) * PAGE as u64; + let mut t = db.begin_write().await.unwrap(); + t.link_segment("reconcile-short-footer", &m).await.unwrap(); + t.commit().await.unwrap(); + *vfs.short_read_at.lock().unwrap() = Some((live_path, footer_offset)); + } + + let db = Db::open_existing(vfs, [9u8; 32], PAGE, REALM) + .await + .expect("reconcile must retry transient short live segment footer reads"); + let reader = db + .open_segment(REALM, "reconcile-short-footer") + .await + .unwrap(); + let page = reader.read_page(1).await.unwrap(); + assert!(page.starts_with(b"reconcile-short-footer")); +} + +#[tokio::test(flavor = "current_thread")] +async fn reconcile_surfaces_promote_sync_dir_error_then_retry_succeeds() { + let vfs = FailOnceVfs::new(MemVfs::new()); + { + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let mut w = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + w.append_page(SegmentPageKind::Data, b"staged-promote") + .await + .unwrap(); + let m = w.seal().await.unwrap(); + let live_path = format!("seg/{}", hex(&m.segment_id)); + let staging_path = format!("seg/.staging/{}", hex(&m.segment_id)); + let mut t = db.begin_write().await.unwrap(); + t.link_segment("staged-promote", &m).await.unwrap(); + t.commit().await.unwrap(); + + vfs.inner.rename(&live_path, &staging_path).await.unwrap(); + } + + vfs.fail_sync_dir.store(true, Ordering::SeqCst); + let err = match Db::open_existing(vfs.clone(), [9u8; 32], PAGE, REALM).await { + Ok(db) => { + drop(db); + panic!("reconcile promote sync_dir errors must not be swallowed during open"); + } + Err(err) => err, + }; + assert!( + matches!(err, PagedbError::Io(_)), + "expected staged-promote sync_dir error to surface, got {err:?}" + ); + + let db = Db::open_existing(vfs, [9u8; 32], PAGE, REALM) + .await + .expect("retry after transient reconcile sync fault must succeed"); + let reader = db.open_segment(REALM, "staged-promote").await.unwrap(); + let page = reader.read_page(1).await.unwrap(); + assert!(page.starts_with(b"staged-promote")); +} + +#[tokio::test(flavor = "current_thread")] +async fn open_segment_surfaces_live_segment_open_error_then_retry_succeeds() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let mut w = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + w.append_page(SegmentPageKind::Data, b"reader") + .await + .unwrap(); + let m = w.seal().await.unwrap(); + let live_path = format!("seg/{}", hex(&m.segment_id)); + let mut t = db.begin_write().await.unwrap(); + t.link_segment("reader-open", &m).await.unwrap(); + t.commit().await.unwrap(); + + *vfs.fail_live_segment_open_path.lock().unwrap() = Some(live_path); + let err = match db.open_segment(REALM, "reader-open").await { + Ok(_) => panic!("live segment open I/O errors must not be swallowed"), + Err(err) => err, + }; + assert!( + matches!(err, PagedbError::Io(ref io) if io.kind() == std::io::ErrorKind::PermissionDenied), + "expected live segment PermissionDenied to surface, got {err:?}" + ); + + let reader = db + .open_segment(REALM, "reader-open") + .await + .expect("retry after transient open fault must succeed"); + let page = reader.read_page(1).await.unwrap(); + assert!(page.starts_with(b"reader")); +} + +#[tokio::test(flavor = "current_thread")] +async fn open_segment_retries_short_header_read() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let mut w = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + w.append_page(SegmentPageKind::Data, b"short-header") + .await + .unwrap(); + let m = w.seal().await.unwrap(); + let live_path = format!("seg/{}", hex(&m.segment_id)); + let mut t = db.begin_write().await.unwrap(); + t.link_segment("short-header", &m).await.unwrap(); + t.commit().await.unwrap(); + + *vfs.short_read_at.lock().unwrap() = Some((live_path, 0)); + let reader = db + .open_segment(REALM, "short-header") + .await + .expect("open_segment must retry transient short header reads"); + let page = reader.read_page(1).await.unwrap(); + assert!(page.starts_with(b"short-header")); +} + +#[tokio::test(flavor = "current_thread")] +async fn segment_read_retries_short_data_page_read() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let mut w = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + w.append_page(SegmentPageKind::Data, b"short-read") + .await + .unwrap(); + let m = w.seal().await.unwrap(); + let live_path = format!("seg/{}", hex(&m.segment_id)); + let mut t = db.begin_write().await.unwrap(); + t.link_segment("short-read", &m).await.unwrap(); + t.commit().await.unwrap(); + + let reader = db.open_segment(REALM, "short-read").await.unwrap(); + *vfs.short_read_at.lock().unwrap() = Some((live_path, PAGE as u64)); + let page = reader + .read_page(1) + .await + .expect("segment page reads must retry transient short reads"); + assert!(page.starts_with(b"short-read")); +} + +#[tokio::test(flavor = "current_thread")] +async fn open_segment_retries_short_footer_read() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let mut w = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + w.append_page(SegmentPageKind::Data, b"short-footer") + .await + .unwrap(); + let m = w.seal().await.unwrap(); + let live_path = format!("seg/{}", hex(&m.segment_id)); + let footer_offset = (m.page_count - 1) * PAGE as u64; + let mut t = db.begin_write().await.unwrap(); + t.link_segment("short-footer", &m).await.unwrap(); + t.commit().await.unwrap(); + + *vfs.short_read_at.lock().unwrap() = Some((live_path, footer_offset)); + let reader = db + .open_segment(REALM, "short-footer") + .await + .expect("open_segment must retry transient short footer reads"); + let page = reader.read_page(1).await.unwrap(); + assert!(page.starts_with(b"short-footer")); +} + +#[tokio::test(flavor = "current_thread")] +async fn find_extent_retries_short_index_page_read() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let mut w = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + let extent = w.append_extent(&[b"indexed-page"]).await.unwrap(); + let m = w.seal().await.unwrap(); + let live_path = format!("seg/{}", hex(&m.segment_id)); + let index_offset = 2 * PAGE as u64; + let mut t = db.begin_write().await.unwrap(); + t.link_segment("short-index", &m).await.unwrap(); + t.commit().await.unwrap(); + + let reader = db.open_segment(REALM, "short-index").await.unwrap(); + assert_eq!(extent.start_page_id, 1); + *vfs.short_read_at.lock().unwrap() = Some((live_path, index_offset)); + let pages = reader + .find_extent(extent.start_page_id) + .await + .expect("extent index load must retry transient short index-page reads"); + assert_eq!(pages.len(), 1); + assert!(pages[0].starts_with(b"indexed-page")); +} + +#[tokio::test(flavor = "current_thread")] +async fn segment_writer_retries_short_data_page_write_before_publish() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let mut w = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + let staging_entries = vfs.inner.list_dir("seg/.staging").await.unwrap(); + assert_eq!(staging_entries.len(), 1); + let staging_path = format!("seg/.staging/{}", staging_entries[0]); + *vfs.short_write_at.lock().unwrap() = Some((staging_path, PAGE as u64)); + + w.append_page(SegmentPageKind::Data, b"short-write") + .await + .expect("segment writer must retry short data page writes"); + let m = w.seal().await.unwrap(); + let mut t = db.begin_write().await.unwrap(); + t.link_segment("short-write", &m).await.unwrap(); + t.commit().await.unwrap(); + + let reader = db.open_segment(REALM, "short-write").await.unwrap(); + let page = reader.read_page(1).await.unwrap(); + assert!(page.starts_with(b"short-write")); +} + +#[tokio::test(flavor = "current_thread")] +async fn spill_append_retries_short_body_write_before_returning_handle() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let opts = OpenOptions::default().with_scratch_bytes(1024 * 1024); + let db = Db::open_internal_with_options(vfs.clone(), [9u8; 32], PAGE, REALM, opts) + .await + .unwrap(); + let mut t = db.begin_write().await.unwrap(); + + *vfs.short_write_at.lock().unwrap() = Some(("tmp/scratch-1".to_string(), 0)); + let got = { + let mut spill = t.spill_scope(); + let handle = spill + .append(b"spill-short-write") + .await + .expect("spill append must retry short ciphertext body writes"); + + spill + .read(handle) + .await + .expect("spill handle must decrypt after a transient short write") + }; + assert_eq!(got, b"spill-short-write"); + + t.abort().await; +} + +#[tokio::test(flavor = "current_thread")] +async fn spill_read_retries_short_body_read_before_decrypting() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let opts = OpenOptions::default().with_scratch_bytes(1024 * 1024); + let db = Db::open_internal_with_options(vfs.clone(), [9u8; 32], PAGE, REALM, opts) + .await + .unwrap(); + let mut t = db.begin_write().await.unwrap(); + + let got = { + let mut spill = t.spill_scope(); + let handle = spill.append(b"spill-short-read").await.unwrap(); + + *vfs.short_read_at.lock().unwrap() = Some(("tmp/scratch-1".to_string(), 0)); + spill + .read(handle) + .await + .expect("spill reads must retry partial ciphertext body reads before decrypting") + }; + assert_eq!(got, b"spill-short-read"); + + t.abort().await; +} + +#[tokio::test(flavor = "current_thread")] +async fn commit_retries_short_header_slot_write_before_reporting_success() { + let vfs = FailOnceVfs::new(MemVfs::new()); + { + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let mut t = db.begin_write().await.unwrap(); + t.put(b"header-short-write", b"durable").await.unwrap(); + *vfs.short_write_at.lock().unwrap() = Some(("/main.db".to_string(), PAGE as u64)); + t.commit() + .await + .expect("header commit must retry short inactive-slot writes"); + } + + let reopened = Db::open_existing(vfs, [9u8; 32], PAGE, REALM) + .await + .expect("database should reopen after committed header short write retry"); + let r = reopened.begin_read().await.unwrap(); + assert_eq!( + r.get(b"header-short-write").await.unwrap().as_deref(), + Some(&b"durable"[..]) + ); +} + +#[tokio::test(flavor = "current_thread")] +async fn open_existing_retries_short_latest_header_slot_read() { + let vfs = FailOnceVfs::new(MemVfs::new()); + { + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let mut t = db.begin_write().await.unwrap(); + t.put(b"header-short-read", b"latest").await.unwrap(); + t.commit().await.unwrap(); + } + + *vfs.short_read_at.lock().unwrap() = Some(("/main.db".to_string(), PAGE as u64)); + let reopened = Db::open_existing(vfs, [9u8; 32], PAGE, REALM) + .await + .expect("open_existing must retry transient short header-slot reads"); + let r = reopened.begin_read().await.unwrap(); + assert_eq!( + r.get(b"header-short-read").await.unwrap().as_deref(), + Some(&b"latest"[..]) + ); +} + +#[tokio::test(flavor = "current_thread")] +async fn read_retries_short_main_data_page_read() { + let vfs = FailOnceVfs::new(MemVfs::new()); + { + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let mut t = db.begin_write().await.unwrap(); + t.put(b"main-short-read", b"value").await.unwrap(); + t.commit().await.unwrap(); + } + + let reopened = Db::open_read_only(vfs.clone(), [9u8; 32], PAGE, REALM, OpenOptions::default()) + .await + .unwrap(); + + *vfs.short_main_db_read_at_or_after.lock().unwrap() = Some((PAGE * 2) as u64); + let r = reopened.begin_read().await.unwrap(); + assert_eq!( + r.get(b"main-short-read").await.unwrap().as_deref(), + Some(&b"value"[..]) + ); +} + +#[tokio::test(flavor = "current_thread")] +async fn compact_now_surfaces_live_segment_open_error_then_retry_succeeds() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + let mut w = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + w.append_page(SegmentPageKind::Data, b"compact-reader") + .await + .unwrap(); + let m = w.seal().await.unwrap(); + let live_path = format!("seg/{}", hex(&m.segment_id)); + let mut t = db.begin_write().await.unwrap(); + t.link_segment("compact-open", &m).await.unwrap(); + t.commit().await.unwrap(); + + *vfs.fail_live_segment_open_path.lock().unwrap() = Some(live_path); + let err = match db.compact_now().await { + Ok(_) => panic!("compact_now must not skip live segment open I/O errors"), + Err(err) => err, + }; + assert!( + matches!(err, PagedbError::Io(ref io) if io.kind() == std::io::ErrorKind::PermissionDenied), + "expected compact live segment PermissionDenied to surface, got {err:?}" + ); + + db.compact_now() + .await + .expect("retry after transient compact segment-open fault must succeed"); + let reader = db.open_segment(REALM, "compact-open").await.unwrap(); + let page = reader.read_page(1).await.unwrap(); + assert!(page.starts_with(b"compact-reader")); +} + +#[tokio::test(flavor = "current_thread")] +async fn compact_now_surfaces_dense_repack_sync_dir_error_then_reopen_succeeds() { + let vfs = FailOnceVfs::new(MemVfs::new()); + { + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + { + let mut t = db.begin_write().await.unwrap(); + t.put(b"sentinel", b"kept").await.unwrap(); + for i in 0..80u32 { + let key = format!("dead-{i:04}"); + t.put(key.as_bytes(), &[i as u8; 96]).await.unwrap(); + } + t.commit().await.unwrap(); + } + { + let mut t = db.begin_write().await.unwrap(); + for i in 0..80u32 { + let key = format!("dead-{i:04}"); + t.delete(key.as_bytes()).await.unwrap(); + } + t.commit().await.unwrap(); + } + + vfs.fail_sync_dir.store(true, Ordering::SeqCst); + let err = db + .compact_now() + .await + .expect_err("dense repack root sync_dir failures must surface"); + assert!( + matches!(err, PagedbError::DurablyCommittedButUnpublished { .. }), + "post-rename sync failure must report unknown publication state, got {err:?}" + ); + } + + let reopened = Db::open_existing(vfs, [9u8; 32], PAGE, REALM) + .await + .expect("reopen after transient dense-repack sync fault must succeed"); + let r = reopened.begin_read().await.unwrap(); + assert_eq!( + r.get(b"sentinel").await.unwrap().as_deref(), + Some(&b"kept"[..]) + ); +} + +#[tokio::test(flavor = "current_thread")] +async fn gc_now_surfaces_tombstone_list_dir_error_then_retry_succeeds() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let options = OpenOptions::default().with_commit_history_retain(RetainPolicy::Disabled); + let db = Db::open_internal_with_options(vfs.clone(), [9u8; 32], PAGE, REALM, options) + .await + .unwrap(); + let mut w = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + w.append_page(SegmentPageKind::Data, b"gc-list") + .await + .unwrap(); + let m = w.seal().await.unwrap(); + { + let mut t = db.begin_write().await.unwrap(); + t.link_segment("gc-list", &m).await.unwrap(); + t.commit().await.unwrap(); + } + { + let mut t = db.begin_write().await.unwrap(); + t.unlink_segment("gc-list").await.unwrap(); + t.commit().await.unwrap(); + } + + vfs.fail_tombstone_list_dir.store(true, Ordering::SeqCst); + let err = db + .gc_now() + .await + .expect_err("tombstone list_dir I/O errors must surface"); + assert!( + matches!(err, PagedbError::Io(ref io) if io.kind() == std::io::ErrorKind::PermissionDenied), + "expected tombstone list_dir PermissionDenied to surface, got {err:?}" + ); + + let stats = db.gc_now().await.unwrap(); + assert!(stats.reclaimed_segments >= 1); +} + +#[tokio::test(flavor = "current_thread")] +async fn stats_surfaces_main_db_open_error_then_retry_succeeds() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let db = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + + vfs.fail_main_db_read_open.store(true, Ordering::SeqCst); + let err = db + .stats() + .await + .expect_err("main.db read errors must not be reported as zero bytes"); + assert!( + matches!(err, PagedbError::Io(ref io) if io.kind() == std::io::ErrorKind::PermissionDenied), + "expected stats PermissionDenied to surface, got {err:?}" + ); + + let stats = db.stats().await.unwrap(); + assert!(stats.main_db_bytes > 0); +} + +#[tokio::test(flavor = "current_thread")] +async fn stats_surfaces_free_list_read_error_then_retry_succeeds() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let opts = OpenOptions::default().with_commit_history_retain(RetainPolicy::Disabled); + let value = [0xABu8; 256]; + { + let db = Db::open_internal_with_options(vfs.clone(), [9u8; 32], PAGE, REALM, opts.clone()) + .await + .unwrap(); + { + let mut w = db.begin_write().await.unwrap(); + for i in 0u32..400 { + w.put(format!("k{i:05}").as_bytes(), &value).await.unwrap(); + } + w.commit().await.unwrap(); + } + { + let mut w = db.begin_write().await.unwrap(); + for i in 0u32..400 { + w.delete(format!("k{i:05}").as_bytes()).await.unwrap(); + } + w.commit().await.unwrap(); + } + assert!( + db.stats().await.unwrap().free_list_pending_entries > 20, + "test setup must create durable free-list entries" + ); + } + + let db = Db::open_existing_with_options(vfs.clone(), [9u8; 32], PAGE, REALM, opts) + .await + .unwrap(); + vfs.fail_main_db_read_at.store(true, Ordering::SeqCst); + let err = db + .stats() + .await + .expect_err("free-list read errors must not be reported as zero pending entries"); + assert!( + matches!(err, PagedbError::Io(ref io) if io.kind() == std::io::ErrorKind::PermissionDenied), + "expected stats free-list PermissionDenied to surface, got {err:?}" + ); + + let stats = db.stats().await.unwrap(); + assert!(stats.free_list_pending_entries > 20); +} + +#[tokio::test(flavor = "current_thread")] +async fn open_surfaces_writer_lock_io_error_then_retry_succeeds() { + let vfs = FailOnceVfs::new(MemVfs::new()); + + vfs.fail_writer_lock.store(true, Ordering::SeqCst); + let err = match Db::open(vfs.clone(), [9u8; 32], PAGE, REALM, Default::default()).await { + Ok(_) => panic!("writer lock I/O errors must not be collapsed to AlreadyOpen"), + Err(err) => err, + }; + assert!( + matches!(err, PagedbError::Io(ref io) if io.kind() == std::io::ErrorKind::PermissionDenied), + "expected writer lock PermissionDenied to surface, got {err:?}" + ); + + Db::open(vfs, [9u8; 32], PAGE, REALM, Default::default()) + .await + .expect("retry after transient writer lock fault must succeed"); +} + +#[tokio::test(flavor = "current_thread")] +async fn open_surfaces_main_db_probe_error_then_retry_succeeds() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let existing = Db::open_internal(vfs.clone(), [9u8; 32], PAGE, REALM) + .await + .unwrap(); + drop(existing); + + vfs.fail_main_db_read_open.store(true, Ordering::SeqCst); + let err = match Db::open(vfs.clone(), [9u8; 32], PAGE, REALM, Default::default()).await { + Ok(_) => panic!("main.db probe I/O errors must not be discarded"), + Err(err) => err, + }; + assert!( + matches!(err, PagedbError::Io(ref io) if io.kind() == std::io::ErrorKind::PermissionDenied), + "expected main.db probe PermissionDenied to surface, got {err:?}" + ); + + Db::open(vfs, [9u8; 32], PAGE, REALM, Default::default()) + .await + .expect("retry after transient main.db probe fault must succeed"); +} + +#[tokio::test(flavor = "current_thread")] +async fn open_preserves_named_counter_when_nonce_anchor_is_higher() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let opts = OpenOptions::default().with_buffer_pool_pages(64); + { + let db = Db::open_internal_with_options(vfs.clone(), [9u8; 32], PAGE, REALM, opts.clone()) + .await + .unwrap(); + let mut txn = db.begin_write().await.unwrap(); + let mut counter = txn.counter("recovery-write").unwrap(); + counter.set(1).await.unwrap(); + drop(counter); + txn.commit().await.unwrap(); + } + + vfs.fail_main_db_write_at.store(true, Ordering::SeqCst); + let db = Db::open_existing_with_options(vfs.clone(), [9u8; 32], PAGE, REALM, opts) + .await + .expect("open must not rewrite named counters from the page-nonce anchor"); + assert!( + vfs.fail_main_db_write_at.load(Ordering::SeqCst), + "the injected main.db write fault must remain unused during open" + ); + let mut txn = db.begin_write().await.unwrap(); + let counter = txn.counter("recovery-write").unwrap(); + assert_eq!(counter.get().await.unwrap(), 1); + drop(counter); + txn.abort().await; +} + +#[tokio::test(flavor = "current_thread")] +async fn open_read_only_does_not_rewrite_named_counters() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let opts = OpenOptions::default().with_buffer_pool_pages(64); + { + let db = Db::open_internal_with_options(vfs.clone(), [9u8; 32], PAGE, REALM, opts.clone()) + .await + .unwrap(); + let mut txn = db.begin_write().await.unwrap(); + let mut counter = txn.counter("read-only-recovery").unwrap(); + counter.set(1).await.unwrap(); + drop(counter); + txn.commit().await.unwrap(); + } + + let writes_before = vfs.main_db_write_at_count.load(Ordering::SeqCst); + let db = Db::open_read_only(vfs.clone(), [9u8; 32], PAGE, REALM, opts) + .await + .expect("read-only open should not need recovery writes"); + assert!(!db.is_writer()); + drop(db); + let writes_after = vfs.main_db_write_at_count.load(Ordering::SeqCst); + assert_eq!( + writes_after, writes_before, + "read-only open must not mutate main.db while recovering counters" + ); +} + +#[tokio::test(flavor = "current_thread")] +async fn gc_now_surfaces_tombstone_len_error_then_retry_succeeds() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let options = OpenOptions::default().with_commit_history_retain(RetainPolicy::Disabled); + let db = Db::open_internal_with_options(vfs.clone(), [9u8; 32], PAGE, REALM, options) + .await + .unwrap(); + let mut w = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + w.append_page(SegmentPageKind::Data, b"gc-len") + .await + .unwrap(); + let m = w.seal().await.unwrap(); + { + let mut t = db.begin_write().await.unwrap(); + t.link_segment("gc-len", &m).await.unwrap(); + t.commit().await.unwrap(); + } + { + let mut t = db.begin_write().await.unwrap(); + t.unlink_segment("gc-len").await.unwrap(); + t.commit().await.unwrap(); + } + + vfs.fail_tombstone_len.store(true, Ordering::SeqCst); + let err = db + .gc_now() + .await + .expect_err("tombstone len I/O errors must surface before delete"); + assert!( + matches!(err, PagedbError::Io(ref io) if io.kind() == std::io::ErrorKind::PermissionDenied), + "expected tombstone len PermissionDenied to surface, got {err:?}" + ); + + let stats = db.gc_now().await.unwrap(); + assert!(stats.reclaimed_segments >= 1); + assert!(stats.reclaimed_bytes > 0); +} + +#[tokio::test(flavor = "current_thread")] +async fn gc_now_surfaces_tombstone_remove_error_then_retry_succeeds() { + let vfs = FailOnceVfs::new(MemVfs::new()); + let options = OpenOptions::default().with_commit_history_retain(RetainPolicy::Disabled); + let db = Db::open_internal_with_options(vfs.clone(), [9u8; 32], PAGE, REALM, options) + .await + .unwrap(); + let mut segment = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + segment + .append_page(SegmentPageKind::Data, b"gc-remove") + .await + .unwrap(); + let meta = segment.seal().await.unwrap(); + { + let mut txn = db.begin_write().await.unwrap(); + txn.link_segment("gc-remove", &meta).await.unwrap(); + txn.commit().await.unwrap(); + } + { + let mut txn = db.begin_write().await.unwrap(); + txn.unlink_segment("gc-remove").await.unwrap(); + txn.commit().await.unwrap(); + } + + vfs.fail_tombstone_remove.store(true, Ordering::SeqCst); + let err = db + .gc_now() + .await + .expect_err("tombstone remove I/O errors must surface"); + assert!( + matches!(err, PagedbError::Io(ref io) if io.kind() == std::io::ErrorKind::PermissionDenied), + "expected tombstone remove PermissionDenied to surface, got {err:?}" + ); + + let stats = db.gc_now().await.unwrap(); + assert!(stats.reclaimed_segments >= 1); + assert!(stats.reclaimed_bytes > 0); +} + +fn hex(bytes: &[u8; 16]) -> String { + bytes.iter().map(|b| format!("{b:02x}")).collect() +} diff --git a/tests/durability/mod.rs b/tests/durability/mod.rs index 4765299..d23af1a 100644 --- a/tests/durability/mod.rs +++ b/tests/durability/mod.rs @@ -1,5 +1,6 @@ pub mod apply_journal_crash; pub mod cross_process_lock; +pub mod metadata_errors; pub mod reader_publication; pub mod second_open_rejected; pub mod sync_dir_on_commit; diff --git a/tests/header.rs b/tests/header.rs index dff4fa0..c312090 100644 --- a/tests/header.rs +++ b/tests/header.rs @@ -259,3 +259,55 @@ async fn open_picks_higher_seq_when_both_verify() { assert_eq!(slot, ActiveSlot::B); assert_eq!(got.counter_anchor, 9000); } + +/// A file truncated mid-slot must still open on the surviving copy. +/// +/// The A/B protocol exists so that one unusable slot is survivable. Reading a +/// slot short is the same condition as a slot that fails its HK-MAC: absent, +/// not fatal. Rejecting the open outright would throw away the intact header +/// sitting in the other slot. +#[tokio::test(flavor = "current_thread")] +async fn open_survives_a_file_truncated_inside_the_inactive_slot() { + let vfs = MemVfs::new(); + let hk = hk(); + bootstrap_header(&vfs, "/main.db", &hk, &sample(1, 12, 4242), 4096) + .await + .unwrap(); + + // Slot A is intact; slot B is cut in half. + let mut f = vfs.open("/main.db", OpenMode::ReadWrite).await.unwrap(); + f.truncate(4096 + 2048).await.unwrap(); + f.sync().await.unwrap(); + drop(f); + + let (got, slot) = open_header(&vfs, "/main.db", &hk, 4096) + .await + .expect("the intact slot must still open the database"); + assert_eq!(slot, ActiveSlot::A); + assert_eq!(got.seq, 1); + assert_eq!(got.counter_anchor, 4242); +} + +/// With neither slot readable, the header layer still reports corruption +/// rather than leaking a raw end-of-file error. +#[tokio::test(flavor = "current_thread")] +async fn open_reports_corruption_when_truncation_takes_both_slots() { + let vfs = MemVfs::new(); + let hk = hk(); + bootstrap_header(&vfs, "/main.db", &hk, &sample(1, 12, 0), 4096) + .await + .unwrap(); + + let mut f = vfs.open("/main.db", OpenMode::ReadWrite).await.unwrap(); + f.truncate(0).await.unwrap(); + f.sync().await.unwrap(); + drop(f); + + let err = open_header(&vfs, "/main.db", &hk, 4096) + .await + .expect_err("a file with no readable slot must not open"); + assert!( + matches!(err, PagedbError::Corruption(_)), + "expected Corruption, got {err:?}" + ); +} diff --git a/tests/smoke.rs b/tests/smoke.rs index 02cfed2..28a920f 100644 --- a/tests/smoke.rs +++ b/tests/smoke.rs @@ -42,6 +42,10 @@ fn every_variant_displays() { PagedbError::ChecksumFailure, PagedbError::corruption(CorruptionDetail::HeaderUnverifiable), PagedbError::structural_header_invalid("main.db", "magic"), + PagedbError::vfs_contract_violated( + "read_at", + "reported more bytes than the caller requested", + ), PagedbError::footer_framing_invalid("magic"), PagedbError::node_body_malformed("slot_directory"), PagedbError::node_kind_mismatch(Some(7), "leaf", "internal"),