From 1cb10dfcb29e30ac51d29febbcf2673007d3cea8 Mon Sep 17 00:00:00 2001 From: Andrew Briscoe Date: Sun, 26 Jul 2026 11:33:43 -0600 Subject: [PATCH 1/7] fix(storage): complete metadata I/O and surface backend errors --- src/compaction/full.rs | 10 +- src/pager/core.rs | 11 +- src/pager/header.rs | 14 +- src/segment/authenticated_metadata.rs | 29 +- src/segment/reader.rs | 24 +- src/segment/writer.rs | 10 +- src/txn/db/misc.rs | 19 +- src/txn/db/open/existing.rs | 8 +- src/txn/db/open/modes.rs | 106 ++- src/txn/db/util.rs | 9 +- src/txn/write/spill.rs | 33 +- src/vfs/mod.rs | 1 + src/vfs/traits.rs | 66 ++ tests/durability/metadata_errors.rs | 1156 +++++++++++++++++++++++++ tests/durability/mod.rs | 1 + 15 files changed, 1388 insertions(+), 109 deletions(-) create mode 100644 tests/durability/metadata_errors.rs 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/pager/core.rs b/src/pager/core.rs index a8562c4..3a549c7 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,8 @@ 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; - } - } + 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..aa25ff3 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. @@ -92,13 +92,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_exact_at(&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_exact_at(&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 +139,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..ee1a89e 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, checked_read_progress}; use super::types::{EXTENT_INDEX_ENTRY_LEN, ExtentIndexEntry}; use super::writer::{live_path, staging_path}; @@ -85,17 +84,23 @@ 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")); + let mut header_offset = 0; + let mut header_remaining = &mut header_bytes[..]; + while !header_remaining.is_empty() { + let read = file.read_at(header_offset, header_remaining).await?; + checked_read_progress(&mut header_offset, read, header_remaining.len())?; + header_remaining = header_remaining.split_at_mut(read).1; } 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")); + let mut next_footer_offset = footer_offset; + let mut footer_remaining = &mut footer_bytes[..]; + while !footer_remaining.is_empty() { + let read = file.read_at(next_footer_offset, footer_remaining).await?; + checked_read_progress(&mut next_footer_offset, read, footer_remaining.len())?; + footer_remaining = footer_remaining.split_at_mut(read).1; } let (footer, manifest) = { let mut lru = pager.dek_lru().lock(); @@ -332,8 +337,12 @@ 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")); + let mut read_offset = offset; + let mut remaining = &mut page[..]; + while !remaining.is_empty() { + let read = context.file.read_at(read_offset, remaining).await?; + checked_read_progress(&mut read_offset, read, remaining.len())?; + remaining = remaining.split_at_mut(read).1; } let (cipher_id, mk_epoch) = extract_page_header_ids(&page)?; if cipher_id != context.cipher_id { diff --git a/src/segment/reader.rs b/src/segment/reader.rs index db9a6a6..84f9dc5 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, checked_read_progress}; 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,9 +228,12 @@ 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); + let mut read_offset = offset; + let mut remaining = &mut buf[..]; + while !remaining.is_empty() { + let read = self.file.read_at(read_offset, remaining).await?; + checked_read_progress(&mut read_offset, read, remaining.len())?; + remaining = remaining.split_at_mut(read).1; } // Try each segment page kind; AAD binding rejects wrong ones. 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/misc.rs b/src/txn/db/misc.rs index 7daf362..99ad2c7 100644 --- a/src/txn/db/misc.rs +++ b/src/txn/db/misc.rs @@ -112,20 +112,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 }; @@ -161,7 +158,7 @@ impl Db { ); let seg_start = vec![CatalogRowKind::Segment as u8]; - let seg_rows = tree.scan_prefix(&seg_start).await.unwrap_or_default(); + let seg_rows = tree.scan_prefix(&seg_start).await?; let seg_count = u32::try_from(seg_rows.len()).unwrap_or(u32::MAX); let seg_bytes: u64 = seg_rows .iter() diff --git a/src/txn/db/open/existing.rs b/src/txn/db/open/existing.rs index 10b957c..bd1fa44 100644 --- a/src/txn/db/open/existing.rs +++ b/src/txn/db/open/existing.rs @@ -12,7 +12,7 @@ use crate::options::OpenOptions; use crate::pager::header::ActiveSlot; use crate::pager::structural_header::MainDbHeaderFields; use crate::pager::{Pager, PagerConfig}; -use crate::vfs::{Vfs, VfsFile}; +use crate::vfs::{Vfs, read_exact_at}; use crate::{RealmId, Result}; use super::super::super::mode::DbMode; @@ -104,13 +104,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_exact_at(&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_exact_at(&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/util.rs b/src/txn/db/util.rs index 32a1726..7803417 100644 --- a/src/txn/db/util.rs +++ b/src/txn/db/util.rs @@ -4,7 +4,7 @@ use crate::Result; use crate::crypto::kdf::{derive_hk, derive_mk}; use crate::errors::PagedbError; -use crate::vfs::Vfs; +use crate::vfs::{Vfs, read_exact_at}; pub(super) fn page_size_log2(page_size: usize) -> Result { match page_size { @@ -37,16 +37,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_exact_at(&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_exact_at(&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..9c687d5 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; @@ -123,13 +123,14 @@ 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.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 +159,18 @@ 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?; + let frame_len = body_len + .checked_add(16) + .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?; + let (body, tag_bytes) = frame.split_at_mut(body_len); + let tag: &[u8; 16] = (&*tag_bytes) + .try_into() + .map_err(|_| PagedbError::ChecksumFailure)?; let cipher = self.txn.spill_cipher_readonly()?; let nonce = Nonce::from_bytes(meta.nonce_bytes); @@ -182,9 +187,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..bb806ea 100644 --- a/src/vfs/mod.rs +++ b/src/vfs/mod.rs @@ -39,6 +39,7 @@ 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, 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..d463b30 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,68 @@ 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. +#[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. +#[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, + ))); + } + if read > remaining { + return Err(PagedbError::Io(std::io::Error::other( + "read_at overreported bytes", + ))); + } + let read_u64 = u64::try_from(read) + .map_err(|_| PagedbError::Io(std::io::Error::other("read count overflow")))?; + *offset = offset + .checked_add(read_u64) + .ok_or_else(|| PagedbError::Io(std::io::Error::other("offset overflow")))?; + Ok(()) +} + +/// Write until `buf` is complete, or fail if the backend reports impossible progress. +#[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, + ))); + } + if written > buf.len() { + return Err(PagedbError::Io(std::io::Error::other( + "write_at overreported bytes", + ))); + } + let written_u64 = u64::try_from(written) + .map_err(|_| PagedbError::Io(std::io::Error::other("write count overflow")))?; + offset = offset + .checked_add(written_u64) + .ok_or_else(|| PagedbError::Io(std::io::Error::other("offset overflow")))?; + buf = &buf[written..]; + } + Ok(()) +} 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; From 249305e51f33ce5f9ec7dbded22ea6d27bb03cd2 Mon Sep 17 00:00:00 2001 From: Andrew Briscoe Date: Sun, 26 Jul 2026 13:04:17 -0600 Subject: [PATCH 2/7] fix(catalog): reject malformed persisted metadata --- src/txn/db/catalog.rs | 101 +++++++++++++++++++++++++++++++++++- src/txn/db/misc.rs | 73 ++++++++++++++++++++++++-- src/txn/db/open/recovery.rs | 3 ++ src/txn/db/reader.rs | 59 ++++++++++++++++----- 4 files changed, 217 insertions(+), 19 deletions(-) diff --git a/src/txn/db/catalog.rs b/src/txn/db/catalog.rs index d3786a4..aa2014f 100644 --- a/src/txn/db/catalog.rs +++ b/src/txn/db/catalog.rs @@ -44,14 +44,41 @@ 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. + 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, + ); + for (_key, value) in tree.scan_prefix(&prefix).await? { + Catalog::decode_counter(&value)?; + } + 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 +334,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 99ad2c7..52bf86d 100644 --- a/src/txn/db/misc.rs +++ b/src/txn/db/misc.rs @@ -160,11 +160,11 @@ impl Db { let seg_start = vec![CatalogRowKind::Segment as u8]; let seg_rows = tree.scan_prefix(&seg_start).await?; 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(); + let mut seg_bytes = 0u64; + for (_key, value) in &seg_rows { + seg_bytes = + seg_bytes.saturating_add(Catalog::decode_segment_meta(value)?.total_bytes); + } (seg_count, seg_bytes) }; @@ -189,3 +189,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/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) From 1f59cf900b488e971b4d6081516751ca72294d38 Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Mon, 27 Jul 2026 05:47:26 +0800 Subject: [PATCH 3/7] fix(vfs): reject impossible positional-I/O progress reports A backend that reports transferring more bytes than the caller's remaining buffer has broken the read_at/write_at contract, not corrupted the on-disk format. Give that case its own error variant instead of a generic io::Error, and share the progress-checking logic between reads and writes. Also add a borrowed-handle variant of the read loop for callers that only have a &F, since a future holding &F across an await is Send only when F: Sync, which VfsFile does not require. --- src/errors.rs | 20 +++++++ src/pager/core.rs | 6 ++ src/vfs/mod.rs | 4 +- src/vfs/traits.rs | 148 +++++++++++++++++++++++++++++++++++++++------- tests/smoke.rs | 4 ++ 5 files changed, 159 insertions(+), 23 deletions(-) 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 3a549c7..64ada90 100644 --- a/src/pager/core.rs +++ b/src/pager/core.rs @@ -741,6 +741,12 @@ impl Pager { } let mut buf = vec![0u8; page_size]; { + // 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?; } diff --git a/src/vfs/mod.rs b/src/vfs/mod.rs index bb806ea..e905946 100644 --- a/src/vfs/mod.rs +++ b/src/vfs/mod.rs @@ -39,7 +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, write_all_at}; +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 d463b30..104f9c9 100644 --- a/src/vfs/traits.rs +++ b/src/vfs/traits.rs @@ -91,6 +91,13 @@ pub trait VfsFile: Send { } /// 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, @@ -106,6 +113,10 @@ pub(crate) async fn read_exact_at( } /// 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 { @@ -113,20 +124,11 @@ pub(crate) fn checked_read_progress(offset: &mut u64, read: usize, remaining: us std::io::ErrorKind::UnexpectedEof, ))); } - if read > remaining { - return Err(PagedbError::Io(std::io::Error::other( - "read_at overreported bytes", - ))); - } - let read_u64 = u64::try_from(read) - .map_err(|_| PagedbError::Io(std::io::Error::other("read count overflow")))?; - *offset = offset - .checked_add(read_u64) - .ok_or_else(|| PagedbError::Io(std::io::Error::other("offset overflow")))?; - Ok(()) + checked_transfer_progress(offset, read, remaining, "read_at", "positional read offset") } -/// Write until `buf` is complete, or fail if the backend reports impossible progress. +/// 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, @@ -140,17 +142,119 @@ pub(crate) async fn write_all_at( std::io::ErrorKind::WriteZero, ))); } - if written > buf.len() { - return Err(PagedbError::Io(std::io::Error::other( - "write_at overreported bytes", - ))); - } - let written_u64 = u64::try_from(written) - .map_err(|_| PagedbError::Io(std::io::Error::other("write count overflow")))?; - offset = offset - .checked_add(written_u64) - .ok_or_else(|| PagedbError::Io(std::io::Error::other("offset overflow")))?; + 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/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"), From 97e3f53fd3ee2f1de2878332d405d4e3a0aee68e Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Mon, 27 Jul 2026 05:47:33 +0800 Subject: [PATCH 4/7] fix(pager): treat a truncated header slot as absent, not fatal A file cut short mid-slot is the same recoverable condition as a slot that fails its HK-MAC check: the A/B protocol is designed to survive exactly one unusable slot. Add read_header_slot, which maps a short/EOF read to an absent slot so the caller's usual decode rejects it like any other damaged copy, and route both header-reading call sites through it so the surviving slot still opens the database. --- src/pager/header.rs | 26 +++++++++++++++++-- src/txn/db/open/existing.rs | 7 ++--- src/txn/db/util.rs | 7 ++--- tests/header.rs | 52 +++++++++++++++++++++++++++++++++++++ 4 files changed, 84 insertions(+), 8 deletions(-) diff --git a/src/pager/header.rs b/src/pager/header.rs index aa25ff3..5e817fb 100644 --- a/src/pager/header.rs +++ b/src/pager/header.rs @@ -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 @@ -95,10 +117,10 @@ pub async fn open_header( 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]; - read_exact_at(&mut f, 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")))?; - read_exact_at(&mut f, 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) { diff --git a/src/txn/db/open/existing.rs b/src/txn/db/open/existing.rs index bd1fa44..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, read_exact_at}; +use crate::vfs::Vfs; use crate::{RealmId, Result}; use super::super::super::mode::DbMode; @@ -107,10 +108,10 @@ impl Db { 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]; - read_exact_at(&mut f, 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")))?; - read_exact_at(&mut f, 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/util.rs b/src/txn/db/util.rs index 7803417..dac4322 100644 --- a/src/txn/db/util.rs +++ b/src/txn/db/util.rs @@ -4,7 +4,8 @@ use crate::Result; use crate::crypto::kdf::{derive_hk, derive_mk}; use crate::errors::PagedbError; -use crate::vfs::{Vfs, read_exact_at}; +use crate::pager::header::read_header_slot; +use crate::vfs::Vfs; pub(super) fn page_size_log2(page_size: usize) -> Result { match page_size { @@ -42,10 +43,10 @@ pub(super) async fn peek_restore_mode( 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]; - read_exact_at(&mut f, 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")))?; - read_exact_at(&mut f, 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/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:?}" + ); +} From 6fb4fd97d8cd963a0b7dea84441c3d636f84246e Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Mon, 27 Jul 2026 05:47:45 +0800 Subject: [PATCH 5/7] refactor(segment): use the borrowed positional-read helper for metadata Segment header/footer/index reads only borrow their file handle, so they were hand-rolling the same read-until-complete loop that read_exact_at already implements. Replace them with the read_exact_at_borrowed macro to remove the duplication. --- CHANGELOG.md | 4 +++- src/segment/authenticated_metadata.rs | 26 ++++---------------------- src/segment/reader.rs | 10 ++-------- 3 files changed, 9 insertions(+), 31 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index acf49b6..71c3a25 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. diff --git a/src/segment/authenticated_metadata.rs b/src/segment/authenticated_metadata.rs index ee1a89e..13fc7c8 100644 --- a/src/segment/authenticated_metadata.rs +++ b/src/segment/authenticated_metadata.rs @@ -11,7 +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, VfsFile, checked_read_progress}; +use crate::vfs::{Vfs, VfsFile, read_exact_at_borrowed}; use super::types::{EXTENT_INDEX_ENTRY_LEN, ExtentIndexEntry}; use super::writer::{live_path, staging_path}; @@ -84,24 +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 mut header_offset = 0; - let mut header_remaining = &mut header_bytes[..]; - while !header_remaining.is_empty() { - let read = file.read_at(header_offset, header_remaining).await?; - checked_read_progress(&mut header_offset, read, header_remaining.len())?; - header_remaining = header_remaining.split_at_mut(read).1; - } + 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 mut next_footer_offset = footer_offset; - let mut footer_remaining = &mut footer_bytes[..]; - while !footer_remaining.is_empty() { - let read = file.read_at(next_footer_offset, footer_remaining).await?; - checked_read_progress(&mut next_footer_offset, read, footer_remaining.len())?; - footer_remaining = footer_remaining.split_at_mut(read).1; - } + 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)?; @@ -337,13 +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]; - let mut read_offset = offset; - let mut remaining = &mut page[..]; - while !remaining.is_empty() { - let read = context.file.read_at(read_offset, remaining).await?; - checked_read_progress(&mut read_offset, read, remaining.len())?; - remaining = remaining.split_at_mut(read).1; - } + 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 84f9dc5..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, checked_read_progress}; +use crate::vfs::{Vfs, VfsFile, read_exact_at_borrowed}; use super::authenticated_metadata::{ ExtentIndexDecodeContext, authenticate_segment_metadata, decode_extent_index, @@ -228,13 +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 mut read_offset = offset; - let mut remaining = &mut buf[..]; - while !remaining.is_empty() { - let read = self.file.read_at(read_offset, remaining).await?; - checked_read_progress(&mut read_offset, read, remaining.len())?; - remaining = remaining.split_at_mut(read).1; - } + read_exact_at_borrowed!(self.file, offset, &mut buf[..])?; // Try each segment page kind; AAD binding rejects wrong ones. let try_kinds = [ From 0f6a1a29c018cc24a736a07201d985a4fe151520 Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Mon, 27 Jul 2026 05:47:55 +0800 Subject: [PATCH 6/7] fix(catalog): validate counter and segment rows in bounded batches Reading an entire prefix scan into memory to validate it at open, or to sum segment bytes for stats(), sizes an allocation by however many counters or segments the embedder happens to have. Add collect_prefix_batch_from to the B+ tree for bounded, resumable prefix scans, and stream both catalog walks through it. Also replace a saturating segment-bytes sum with a checked one that reports overflow instead of silently wrapping. --- CHANGELOG.md | 6 ++++++ src/btree/tree/scan.rs | 23 +++++++++++++++++++++++ src/txn/db/catalog.rs | 32 +++++++++++++++++++++++++++++--- src/txn/db/misc.rs | 41 +++++++++++++++++++++++++++++++++++------ 4 files changed, 93 insertions(+), 9 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 71c3a25..526feab 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -40,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/txn/db/catalog.rs b/src/txn/db/catalog.rs index aa2014f..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 @@ -56,6 +61,9 @@ impl Db { /// /// 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, @@ -73,10 +81,28 @@ impl Db { next_page_id, self.page_size, ); - for (_key, value) in tree.scan_prefix(&prefix).await? { - Catalog::decode_counter(&value)?; + 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(()); + } } - Ok(()) } /// Write per-realm quota caps into the catalog B+ tree and persist the diff --git a/src/txn/db/misc.rs b/src/txn/db/misc.rs index 52bf86d..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 { @@ -157,13 +163,36 @@ impl Db { self.page_size, ); - let seg_start = vec![CatalogRowKind::Segment as u8]; - let seg_rows = tree.scan_prefix(&seg_start).await?; - let seg_count = u32::try_from(seg_rows.len()).unwrap_or(u32::MAX); + // 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; - for (_key, value) in &seg_rows { - seg_bytes = - seg_bytes.saturating_add(Catalog::decode_segment_meta(value)?.total_bytes); + 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) From 57e56e55e160607985ed5ba2fca571dcd5dd0b8b Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Mon, 27 Jul 2026 05:48:00 +0800 Subject: [PATCH 7/7] refactor(spill): name the AEAD tag length constant Replace the repeated literal 16 with SPILL_TAG_LEN, and read the tag into a fixed-size array by copying instead of a fallible slice conversion, since the split is already exact by construction. --- src/txn/write/spill.rs | 25 ++++++++++++++++--------- 1 file changed, 16 insertions(+), 9 deletions(-) diff --git a/src/txn/write/spill.rs b/src/txn/write/spill.rs index 9c687d5..9da2ce7 100644 --- a/src/txn/write/spill.rs +++ b/src/txn/write/spill.rs @@ -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( @@ -124,6 +127,8 @@ impl SpillScope<'_, '_, V> { let mut file = self.txn.db.vfs.open(&path, OpenMode::CreateOrOpen).await?; let body_offset = self.txn.spill_bytes_used; 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?; @@ -162,15 +167,17 @@ impl SpillScope<'_, '_, V> { let mut file = self.txn.db.vfs.open(path, OpenMode::Read).await?; let body_len = meta.ciphertext_len as usize; + // `usize` is 32-bit on wasm32, where a maximal `ciphertext_len` plus + // the tag genuinely overflows it. let frame_len = body_len - .checked_add(16) + .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 tag: &[u8; 16] = (&*tag_bytes) - .try_into() - .map_err(|_| PagedbError::ChecksumFailure)?; + 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); @@ -187,7 +194,7 @@ impl SpillScope<'_, '_, V> { segment_id: self.txn.db.file_id, }); - cipher.decrypt(&nonce, &aad, body, tag)?; + cipher.decrypt(&nonce, &aad, body, &tag)?; frame.truncate(meta.plaintext_len as usize); Ok(frame) }