diff --git a/src/pager/core.rs b/src/pager/core.rs index 4662462..4719200 100644 --- a/src/pager/core.rs +++ b/src/pager/core.rs @@ -949,11 +949,7 @@ impl Pager { // mk_epoch before constructing AAD and selecting the DEK. self.inner.record_miss(file); let page_size = self.cfg.page_size; - let page_size_u64 = - u64::try_from(page_size).map_err(|_| PagedbError::arithmetic_overflow("page size"))?; - let page_offset = page_id - .checked_mul(page_size_u64) - .ok_or_else(|| PagedbError::arithmetic_overflow("page read offset"))?; + let page_offset = crate::pager::page_space::page_offset(page_id, page_size, "page read")?; let file_handle = self.open_file_handle(file).await?; // Observer-mode retry loop: on AEAD failure retry up to diff --git a/src/pager/page_space.rs b/src/pager/page_space.rs index bf78966..7a176b5 100644 --- a/src/pager/page_space.rs +++ b/src/pager/page_space.rs @@ -15,6 +15,9 @@ //! pages, so a live tree pointer that reaches one is a wild pointer or a //! use-after-free that recycled a reserved id — never a benign condition. +use crate::Result; +use crate::errors::PagedbError; + /// First page id the allocator may hand out. Ids below this are reserved. pub const FIRST_ALLOCATABLE_PAGE_ID: u64 = 4; @@ -24,3 +27,40 @@ pub const FIRST_ALLOCATABLE_PAGE_ID: u64 = 4; pub const fn is_reserved(page_id: u64) -> bool { page_id < FIRST_ALLOCATABLE_PAGE_ID } + +/// Byte offset of `page_id` in a paged file, or an arithmetic error. +/// +/// Page ids reach this from disk — a catalog record, a free-list entry, a +/// header field — so the product is not trusted to fit. `operation` names the +/// caller in the resulting error, since a wrapped offset and a rejected one +/// are indistinguishable by the time a diagnostic reports them. +pub fn page_offset(page_id: u64, page_size: usize, operation: &'static str) -> Result { + let page_size = + u64::try_from(page_size).map_err(|_| PagedbError::arithmetic_overflow(operation))?; + page_id + .checked_mul(page_size) + .ok_or_else(|| PagedbError::arithmetic_overflow(operation)) +} + +#[cfg(test)] +mod tests { + use super::*; + + const PAGE: usize = 4096; + + #[test] + fn offset_scales_by_page_size() { + assert_eq!(page_offset(0, PAGE, "test").unwrap(), 0); + assert_eq!(page_offset(4, PAGE, "test").unwrap(), 16_384); + } + + #[test] + fn offset_rejects_a_page_id_that_overflows_the_address_space() { + let first_unrepresentable = (u64::MAX / PAGE as u64) + 1; + assert!(matches!( + page_offset(first_unrepresentable, PAGE, "test"), + Err(PagedbError::ArithmeticOverflow { .. }) + )); + assert!(page_offset(first_unrepresentable - 1, PAGE, "test").is_ok()); + } +} diff --git a/src/recovery/deep_walk.rs b/src/recovery/deep_walk.rs index d5bcfb0..6d64de0 100644 --- a/src/recovery/deep_walk.rs +++ b/src/recovery/deep_walk.rs @@ -14,14 +14,15 @@ use crate::btree::leaf::{Leaf, LeafValue}; use crate::btree::overflow; use crate::catalog::codec::{Catalog, SegmentMeta}; use crate::crypto::aad::{Aad, AadFields, MAIN_DB_SEGMENT_ID}; +use crate::errors::PagedbError; use crate::pager::format::data_page::extract_page_header_ids; use crate::pager::format::page_kind::PageKind; -use crate::pager::page_space::{FIRST_ALLOCATABLE_PAGE_ID, is_reserved}; +use crate::pager::page_space::{FIRST_ALLOCATABLE_PAGE_ID, is_reserved, page_offset}; use crate::pager::{PageGuard, Pager}; use crate::segment::authenticated_metadata::authenticate_segment_metadata; use crate::txn::db::Db; use crate::vfs::types::OpenMode; -use crate::vfs::{Vfs, VfsFile}; +use crate::vfs::{Vfs, VfsFile, read_exact_at}; /// A single page-level issue found during deep walk. #[non_exhaustive] @@ -260,7 +261,7 @@ pub async fn run_deep_walk(db: &Db) -> Result let vfs: &V = &db.vfs; let main_file_res = vfs.open(main_db_path, OpenMode::Read).await; - let main_file = match main_file_res { + let mut main_file = match main_file_res { Ok(f) => f, Err(e) => { report.page_issues.push(PageIssue { @@ -275,24 +276,25 @@ pub async fn run_deep_walk(db: &Db) -> Result // (HK-MAC, cleartext) and are already verified by `Db::open`. Skip them. // Page 2 and 3 are reserved (apply-journal). Walk from page 4. for page_id in 4..next_page_id { - let offset = page_id * page_size as u64; - let mut buf = vec![0u8; page_size]; - match main_file.read_at(offset, &mut buf).await { - Ok(n) if n < page_size => { - // Short read at the tail — the file may be smaller than expected. - // Report but continue. + let offset = match page_offset(page_id, page_size, "deep-walk main page offset") { + Ok(offset) => offset, + Err(error) => { report.page_issues.push(PageIssue { page_id, - description: format!("short read: expected {page_size} bytes, got {n}"), + description: format!("{error}"), }); report.pages_examined += 1; continue; } - Ok(_) => {} - Err(e) => { + }; + let mut buf = vec![0u8; page_size]; + match read_exact_at(&mut main_file, offset, &mut buf).await { + Ok(()) => {} + Err(error) => { + let description = describe_page_read_failure(&mut main_file, offset, error).await; report.page_issues.push(PageIssue { page_id, - description: format!("read error: {e}"), + description, }); report.pages_examined += 1; continue; @@ -435,7 +437,7 @@ async fn check_segment( let page_size = pager.page_size(); // Check file exists. - let Ok(file) = vfs.open(&live, OpenMode::Read).await else { + let Ok(mut file) = vfs.open(&live, OpenMode::Read).await else { report.drift_issues.push(DriftIssue { segment_id: meta.segment_id, description: "segment file missing from seg/".to_string(), @@ -464,7 +466,16 @@ async fn check_segment( // We don't have a metadata API, but we can check via read: try reading one // byte past the expected end. If it succeeds (on some VFS) we skip the // check; if we read exactly `page_count * page_size` bytes we're consistent. - let expected_size = meta.page_count * page_size as u64; + let expected_size = match page_offset(meta.page_count, page_size, "segment expected size") { + Ok(offset) => offset, + Err(error) => { + report.segment_issues.push(SegmentIssue { + segment_id: meta.segment_id, + description: format!("{error}"), + }); + return; + } + }; let mut probe = vec![0u8; 1]; let over_read = file.read_at(expected_size, &mut probe).await; match over_read { @@ -483,24 +494,24 @@ async fn check_segment( // Walk data pages (1 .. page_count - 1, skipping header=0 and footer=last). let last_data = footer_page_id; for page_id in 1..last_data { - let offset = page_id * page_size as u64; - let mut buf = vec![0u8; page_size]; - let read_res = file.read_at(offset, &mut buf).await; - match read_res { - Ok(n) if n < page_size => { + let offset = match page_offset(page_id, page_size, "segment data page offset") { + Ok(offset) => offset, + Err(error) => { report.segment_issues.push(SegmentIssue { segment_id: meta.segment_id, - description: format!( - "short read at page {page_id}: expected {page_size} bytes, got {n}" - ), + description: format!("page {page_id}: {error}"), }); continue; } - Ok(_) => {} - Err(e) => { + }; + let mut buf = vec![0u8; page_size]; + match read_exact_at(&mut file, offset, &mut buf).await { + Ok(()) => {} + Err(error) => { + let description = describe_page_read_failure(&mut file, offset, error).await; report.segment_issues.push(SegmentIssue { segment_id: meta.segment_id, - description: format!("read error at page {page_id}: {e}"), + description: format!("page {page_id}: {description}"), }); continue; } @@ -553,6 +564,34 @@ async fn check_segment( } } +/// Turn a failed full-page read into a description that keeps the byte counts. +/// +/// `read_exact_at` completes a legal short read and only gives up once the +/// backend stops making progress, so the partial count it consumed never +/// reaches the caller — a truncated file arrives here as a bare +/// `UnexpectedEof`. On a diagnostic surface those numbers are the product: an +/// operator needs to know the file is short and by how much, not merely that a +/// read ended. Any other error is already self-describing and passes through. +async fn describe_page_read_failure( + file: &mut F, + offset: u64, + error: PagedbError, +) -> String { + let is_eof = matches!( + &error, + PagedbError::Io(io) if io.kind() == std::io::ErrorKind::UnexpectedEof + ); + if !is_eof { + return format!("read error: {error}"); + } + match file.len().await { + Ok(len) => format!("truncated: page starts at offset {offset}, file is {len} bytes"), + Err(len_error) => { + format!("read error: {error} (file length unavailable: {len_error})") + } + } +} + /// Collect the set of all page IDs reachable from the main B+ tree root, /// the catalog root, the commit-history root, and the free-list root. /// Pages 0..=3 (reserved) are always considered reachable. @@ -877,10 +916,10 @@ async fn diagnose_overflow_chain( #[cfg(test)] mod tests { use super::*; - use crate::OpenOptions; use crate::btree::node::body_capacity; use crate::pager::format::data_page::ENVELOPE_OVERHEAD; use crate::vfs::memory::MemVfs; + use crate::{OpenOptions, SegmentKind, SegmentPageKind}; const PAGE: usize = 4096; const REALM: crate::RealmId = crate::RealmId::new([0xD3; 16]); @@ -1074,6 +1113,72 @@ mod tests { ); } + /// A catalog `page_count` large enough to overflow a byte offset must be + /// rejected as a structured issue, never reach the arithmetic that would + /// wrap it, and never abort the walk. Authenticated metadata validation is + /// what stops it, ahead of any footer index or loop bound; the checked + /// `page_offset` behind it is the second line, covered directly in + /// `pager::page_space`. + #[tokio::test(flavor = "current_thread")] + async fn deep_walk_rejects_impossible_catalog_page_count() { + let db = open_db().await; + let mut segment = db + .create_segment(REALM, SegmentKind::Unspecified) + .await + .unwrap(); + segment + .append_page(SegmentPageKind::Data, b"deep-walk") + .await + .unwrap(); + let mut meta = segment.seal().await.unwrap(); + { + let mut txn = db.begin_write().await.unwrap(); + txn.link_segment("overflow", &meta).await.unwrap(); + txn.commit().await.unwrap(); + } + + meta.page_count = (u64::MAX / PAGE as u64) + 2; + 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, + ); + let key = Catalog::segment_key(REALM, b"overflow").unwrap(); + tree.put(&key, &Catalog::encode_segment_meta(&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()); + } + + let report = run_deep_walk(&db).await.unwrap(); + assert!( + report.segment_issues.iter().any(|issue| { + issue.segment_id == meta.segment_id + && issue + .description + .contains("authenticated segment metadata invalid") + }), + "impossible segment geometry must become a structured issue, got {report:?}" + ); + assert!( + report.segment_issues.iter().all(|issue| { + !issue.description.contains("segment expected size") + && !issue.description.contains("segment data page offset") + }), + "validation must reject the record before any offset arithmetic runs: {report:?}" + ); + } + async fn sole_overflow_root(db: &Db, leaf_page_id: u64) -> u64 { let (guard, _) = db.pager.read_main_node(leaf_page_id, REALM).await.unwrap(); let leaf = Leaf::decode(guard.body_ref()).unwrap(); diff --git a/src/recovery/reconcile.rs b/src/recovery/reconcile.rs index ebac59b..7ab2c03 100644 --- a/src/recovery/reconcile.rs +++ b/src/recovery/reconcile.rs @@ -277,7 +277,13 @@ mod tests { use crate::vfs::{Vfs, VfsFile}; use crate::{RealmId, btree::BTree}; - use super::repair_catalog; + use super::{repair_catalog, sweep_orphans}; + + fn segment_id(value: u64) -> [u8; 16] { + let mut id = [0; 16]; + id[..8].copy_from_slice(&value.to_le_bytes()); + id + } #[tokio::test(flavor = "current_thread")] async fn malformed_catalog_key_prevents_reconciliation_mutation() { @@ -347,4 +353,57 @@ mod tests { )); assert!(vfs.open(marker, OpenMode::Read).await.is_ok()); } + + /// The sweep has three outcomes and every one of them is destructive if it + /// fires on the wrong file: an expected live segment must survive, a live + /// segment the catalog does not name must become a tombstone rather than a + /// deletion, and an unnamed staging file must be removed outright. The + /// expected set is large enough that a membership test which silently + /// matched on a prefix, a truncated id, or the first entry alone would + /// misclassify one of the three. + #[tokio::test(flavor = "current_thread")] + async fn sweep_orphans_tombstones_live_orphans_and_removes_staged_orphans() { + let vfs = MemVfs::new(); + vfs.mkdir_all("seg/.staging").await.unwrap(); + let expected: Vec<[u8; 16]> = (0..1024).map(segment_id).collect(); + + for id in expected.iter().take(8) { + let path = crate::segment::writer::live_path(id); + let mut file = vfs.open(&path, OpenMode::CreateOrOpen).await.unwrap(); + file.write_at(0, b"live").await.unwrap(); + } + + let live_orphan = segment_id(10_000); + let live_orphan_path = crate::segment::writer::live_path(&live_orphan); + let mut live_file = vfs + .open(&live_orphan_path, OpenMode::CreateOrOpen) + .await + .unwrap(); + live_file.write_at(0, b"orphan").await.unwrap(); + + let staged_orphan = segment_id(10_001); + let staged_orphan_path = crate::segment::writer::staging_path(&staged_orphan); + let mut staged_file = vfs + .open(&staged_orphan_path, OpenMode::CreateOrOpen) + .await + .unwrap(); + staged_file.write_at(0, b"orphan").await.unwrap(); + + sweep_orphans(&vfs, &expected, 77).await.unwrap(); + + for id in expected.iter().take(8) { + let path = crate::segment::writer::live_path(id); + assert!( + vfs.open(&path, OpenMode::Read).await.is_ok(), + "a segment the catalog names must survive the sweep: {path}" + ); + } + assert!(vfs.open(&live_orphan_path, OpenMode::Read).await.is_err()); + let tombstone = format!( + "seg/.tombstone/{}.77", + crate::hex::to_hex_lower(&live_orphan) + ); + assert!(vfs.open(&tombstone, OpenMode::Read).await.is_ok()); + assert!(vfs.open(&staged_orphan_path, OpenMode::Read).await.is_err()); + } } diff --git a/src/recovery/tests/apply_journal_crash.rs b/src/recovery/tests/apply_journal_crash.rs index 7440d7d..f186398 100644 --- a/src/recovery/tests/apply_journal_crash.rs +++ b/src/recovery/tests/apply_journal_crash.rs @@ -133,14 +133,111 @@ async fn promote_action_renames_staging_to_live() { assert_eq!(&buf, b"segment_content"); } +/// A tombstone's goal state is that the live file is *gone*, so its absence is +/// itself the proof that the action completed — unlike a promote, whose goal +/// state is a file that exists and whose absence therefore proves nothing. +/// +/// The two are not symmetric and must not be made so. `seg/.tombstone/` is +/// reclaimed wholesale by `recovery::gc::delete_tombstone_files`, so "live gone +/// and tombstone gone" is the ordinary state of a delete that finished and was +/// then collected. Treating it as unproven would make replay of an already +/// completed journal fail permanently, and would diverge from the pin-aware +/// `Db::tombstone_segment` this fixture models, which reports the same state as +/// complete. #[tokio::test(flavor = "current_thread")] async fn tombstone_action_is_idempotent_when_live_absent() { - // If the live file is absent, the tombstone rename is a no-op. use crate::recovery::journal::execute_journal_actions; + use crate::vfs::types::OpenMode; + + let vfs = MemVfs::new(); + let segment_id = [0xEE; 16]; + let actions = vec![JournalAction::Tombstone { + segment_id, + tombstone_commit_id: 5, + }]; + execute_journal_actions(&vfs, &actions).await.unwrap(); + + // A no-op, not a rename of nothing into something: replay must not leave a + // tombstone behind for a segment whose bytes are already gone. + let tombstone = tombstone_path(&segment_id, 5); + assert!( + vfs.open(&tombstone, OpenMode::Read).await.is_err(), + "replay must not fabricate a tombstone for an already-deleted segment" + ); +} + +/// Replay after a crash between the rename and the journal clear finds the +/// destination already in place. It must leave those bytes untouched — a +/// second rename would either fail or overwrite the retained copy that a +/// reader may still be pinning. +#[tokio::test(flavor = "current_thread")] +async fn tombstone_action_preserves_an_existing_destination() { + use crate::recovery::journal::execute_journal_actions; + use crate::vfs::VfsFile; + use crate::vfs::types::OpenMode; + let vfs = MemVfs::new(); + let segment_id = [0xEF; 16]; + let tombstone = tombstone_path(&segment_id, 5); + vfs.mkdir_all("seg/.tombstone").await.unwrap(); + { + let mut file = vfs.open(&tombstone, OpenMode::CreateNew).await.unwrap(); + file.write_at(0, b"already-tombstoned").await.unwrap(); + file.sync().await.unwrap(); + } + let actions = vec![JournalAction::Tombstone { - segment_id: [0xEE; 16], + segment_id, tombstone_commit_id: 5, }]; execute_journal_actions(&vfs, &actions).await.unwrap(); + + let file = vfs.open(&tombstone, OpenMode::Read).await.unwrap(); + let mut bytes = vec![0u8; b"already-tombstoned".len()]; + let read = file.read_at(0, &mut bytes).await.unwrap(); + assert_eq!(read, bytes.len()); + assert_eq!(&bytes, b"already-tombstoned"); +} + +#[tokio::test(flavor = "current_thread")] +async fn tombstone_action_renames_live_to_tombstone() { + use crate::recovery::journal::execute_journal_actions; + use crate::vfs::VfsFile; + use crate::vfs::types::OpenMode; + + let vfs = MemVfs::new(); + let segment_id = [0xF0; 16]; + let live = crate::segment::writer::live_path(&segment_id); + vfs.mkdir_all("seg").await.unwrap(); + { + let mut file = vfs.open(&live, OpenMode::CreateNew).await.unwrap(); + file.write_at(0, b"segment_content").await.unwrap(); + file.sync().await.unwrap(); + } + + let actions = vec![JournalAction::Tombstone { + segment_id, + tombstone_commit_id: 9, + }]; + execute_journal_actions(&vfs, &actions).await.unwrap(); + + assert!( + vfs.open(&live, OpenMode::Read).await.is_err(), + "the live path must be vacated by the rename" + ); + let file = vfs + .open(&tombstone_path(&segment_id, 9), OpenMode::Read) + .await + .unwrap(); + let mut bytes = vec![0u8; b"segment_content".len()]; + let read = file.read_at(0, &mut bytes).await.unwrap(); + assert_eq!(read, bytes.len()); + assert_eq!(&bytes, b"segment_content"); +} + +fn tombstone_path(segment_id: &[u8; 16], tombstone_commit_id: u64) -> String { + format!( + "seg/.tombstone/{}.{tombstone_commit_id}", + crate::hex::to_hex_lower(segment_id) + ) } diff --git a/tests/fsck_deep.rs b/tests/fsck_deep.rs index 237a47b..f154241 100644 --- a/tests/fsck_deep.rs +++ b/tests/fsck_deep.rs @@ -1,11 +1,15 @@ //! Tests for the deep-walk integrity checker. +use std::sync::{Arc, Mutex}; + use pagedb::options::RetainPolicy; use pagedb::vfs::memory::MemVfs; +use pagedb::vfs::{OpenMode, ReadReq, Vfs, VfsFile, WriteReq}; use pagedb::{Db, OpenOptions, PagedbError, RealmId, run_deep_walk}; const KEK: [u8; 32] = [3u8; 32]; const REALM: RealmId = RealmId::new([1u8; 16]); +const PAGE: usize = 4096; async fn open_db() -> Db { let opts = OpenOptions::default().with_buffer_pool_pages(64); @@ -228,3 +232,235 @@ async fn retained_history_pages_are_not_reported_as_orphans() { report.orphan_page_ids ); } + +/// A `MemVfs` that shortens exactly one positional read by one byte, then +/// behaves normally. +/// +/// Every PageDB VFS is allowed to satisfy a request across several `read_at` +/// calls, so a positive short read is a legal step in a transfer, not evidence +/// about the file. The wrapper leaves the missing byte available to the next +/// call, which is what makes the resulting report a statement about deep walk +/// rather than about the mock: a clean report is only reachable by completing +/// the same legal byte stream. +#[derive(Clone)] +struct ShortReadVfs { + inner: MemVfs, + short_read_at: Arc>>, +} + +impl ShortReadVfs { + fn new() -> Self { + Self { + inner: MemVfs::new(), + short_read_at: Arc::new(Mutex::new(None)), + } + } + + fn short_once_at(&self, path: impl Into, offset: u64) { + *self.short_read_at.lock().unwrap() = Some((path.into(), offset)); + } + + /// Whether the armed short read has been consumed. A test that never fires + /// it proves nothing, so the assertion belongs next to the report. + fn fired(&self) -> bool { + self.short_read_at.lock().unwrap().is_none() + } +} + +struct ShortReadFile { + inner: F, + path: String, + short_read_at: Arc>>, +} + +impl Vfs for ShortReadVfs { + type File = ShortReadFile<::File>; + type LockHandle = ::LockHandle; + + async fn open(&self, path: &str, mode: OpenMode) -> pagedb::Result { + Ok(ShortReadFile { + inner: self.inner.open(path, mode).await?, + path: path.to_string(), + short_read_at: self.short_read_at.clone(), + }) + } + + async fn remove(&self, path: &str) -> pagedb::Result<()> { + self.inner.remove(path).await + } + + async fn rename(&self, from: &str, to: &str) -> pagedb::Result<()> { + self.inner.rename(from, to).await + } + + async fn list_dir(&self, path: &str) -> pagedb::Result> { + self.inner.list_dir(path).await + } + + async fn mkdir_all(&self, path: &str) -> pagedb::Result<()> { + self.inner.mkdir_all(path).await + } + + async fn sync_dir(&self, path: &str) -> pagedb::Result<()> { + self.inner.sync_dir(path).await + } + + async fn lock_exclusive(&self, path: &str) -> pagedb::Result { + self.inner.lock_exclusive(path).await + } + + async fn lock_shared(&self, path: &str) -> pagedb::Result { + self.inner.lock_shared(path).await + } +} + +impl VfsFile for ShortReadFile { + async fn read_at(&self, offset: u64, buf: &mut [u8]) -> pagedb::Result { + let should_shorten = { + let mut armed = self.short_read_at.lock().unwrap(); + let matches_here = armed + .as_ref() + .is_some_and(|(path, at)| path == &self.path && *at == offset); + if matches_here { + armed.take(); + } + matches_here + }; + if should_shorten && buf.len() > 1 { + let short_len = buf.len() - 1; + return self.inner.read_at(offset, &mut buf[..short_len]).await; + } + self.inner.read_at(offset, buf).await + } + + async fn read_at_vectored(&self, reqs: &mut [ReadReq<'_>]) -> pagedb::Result<()> { + self.inner.read_at_vectored(reqs).await + } + + async fn write_at(&mut self, offset: u64, buf: &[u8]) -> pagedb::Result { + self.inner.write_at(offset, buf).await + } + + async fn write_at_vectored(&mut self, reqs: &[WriteReq<'_>]) -> pagedb::Result<()> { + self.inner.write_at_vectored(reqs).await + } + + async fn sync(&mut self) -> pagedb::Result<()> { + self.inner.sync().await + } + + async fn truncate(&mut self, len: u64) -> pagedb::Result<()> { + self.inner.truncate(len).await + } + + async fn len(&self) -> pagedb::Result { + self.inner.len().await + } + + async fn is_empty(&self) -> pagedb::Result { + self.inner.is_empty().await + } + + fn supports_direct_io(&self) -> bool { + self.inner.supports_direct_io() + } +} + +#[tokio::test(flavor = "current_thread")] +async fn deep_walk_completes_a_short_main_data_page_read() { + let vfs = ShortReadVfs::new(); + let opts = OpenOptions::default().with_buffer_pool_pages(64); + let db = Db::open(vfs.clone(), KEK, PAGE, REALM, opts).await.unwrap(); + + let mut txn = db.begin_write().await.unwrap(); + for i in 0u64..10 { + let key = format!("short-main-{i:04}"); + txn.put(key.as_bytes(), &[0xCC; 128]).await.unwrap(); + } + txn.commit().await.unwrap(); + + vfs.short_once_at("/main.db", (PAGE * 4) as u64); + let report = run_deep_walk(&db).await.unwrap(); + + assert!(vfs.fired(), "the armed short read never reached deep walk"); + assert!( + report.page_issues.is_empty(), + "a legal short main.db read must be completed, not reported: {:?}", + report.page_issues + ); + assert!(report.is_clean(), "report should be clean: {report:?}"); +} + +#[tokio::test(flavor = "current_thread")] +async fn deep_walk_completes_a_short_segment_data_page_read() { + let vfs = ShortReadVfs::new(); + let opts = OpenOptions::default().with_buffer_pool_pages(64); + let db = Db::open(vfs.clone(), KEK, PAGE, REALM, opts).await.unwrap(); + + let mut segment = db + .create_segment(REALM, pagedb::SegmentKind::Unspecified) + .await + .unwrap(); + segment + .append_page(pagedb::SegmentPageKind::Data, b"deep-walk-short-segment") + .await + .unwrap(); + let meta = segment.seal().await.unwrap(); + let mut txn = db.begin_write().await.unwrap(); + txn.link_segment("short-segment", &meta).await.unwrap(); + txn.commit().await.unwrap(); + + let segment_path = format!("seg/{}", hex_lower(&meta.segment_id)); + vfs.short_once_at(segment_path, PAGE as u64); + + let report = run_deep_walk(&db).await.unwrap(); + + assert!(vfs.fired(), "the armed short read never reached deep walk"); + assert!( + report.segment_issues.is_empty(), + "a legal short segment read must be completed, not reported: {:?}", + report.segment_issues + ); + assert!(report.is_clean(), "report should be clean: {report:?}"); +} + +fn hex_lower(bytes: &[u8; 16]) -> String { + bytes.iter().map(|byte| format!("{byte:02x}")).collect() +} + +/// A genuinely short file must still be reported, and the report must carry the +/// numbers. `read_exact_at` completes legal short reads and surfaces real +/// truncation as a bare end-of-file, so the walk restores the offset and the +/// file length an operator needs to size the damage. +#[tokio::test(flavor = "current_thread")] +async fn deep_walk_reports_a_truncated_main_db_with_its_extent() { + let vfs = MemVfs::new(); + let opts = OpenOptions::default().with_buffer_pool_pages(64); + let db = Db::open(vfs.clone(), KEK, PAGE, REALM, opts).await.unwrap(); + + let mut txn = db.begin_write().await.unwrap(); + for i in 0u64..400 { + let key = format!("truncated-{i:04}"); + txn.put(key.as_bytes(), &[0xAB; 128]).await.unwrap(); + } + txn.commit().await.unwrap(); + + let kept_bytes = (PAGE * 5) as u64; + { + let mut main = vfs.open("/main.db", OpenMode::CreateOrOpen).await.unwrap(); + main.truncate(kept_bytes).await.unwrap(); + main.sync().await.unwrap(); + } + + let report = run_deep_walk(&db).await.unwrap(); + assert!( + report.page_issues.iter().any(|issue| { + issue.description.contains("truncated:") + && issue + .description + .contains(&format!("file is {kept_bytes} bytes")) + }), + "a truncated main.db must be reported with its extent: {:?}", + report.page_issues + ); +}