Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 1 addition & 5 deletions src/pager/core.rs
Original file line number Diff line number Diff line change
Expand Up @@ -949,11 +949,7 @@ impl<V: Vfs> Pager<V> {
// 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
Expand Down
40 changes: 40 additions & 0 deletions src/pager/page_space.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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<u64> {
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());
}
}
159 changes: 132 additions & 27 deletions src/recovery/deep_walk.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down Expand Up @@ -260,7 +261,7 @@ pub async fn run_deep_walk<V: Vfs + Clone>(db: &Db<V>) -> Result<DeepWalkReport>

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 {
Expand All @@ -275,24 +276,25 @@ pub async fn run_deep_walk<V: Vfs + Clone>(db: &Db<V>) -> Result<DeepWalkReport>
// (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;
Expand Down Expand Up @@ -435,7 +437,7 @@ async fn check_segment<V: Vfs + Clone>(
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(),
Expand Down Expand Up @@ -464,7 +466,16 @@ async fn check_segment<V: Vfs + Clone>(
// 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 {
Expand All @@ -483,24 +494,24 @@ async fn check_segment<V: Vfs + Clone>(
// 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;
}
Expand Down Expand Up @@ -553,6 +564,34 @@ async fn check_segment<V: Vfs + Clone>(
}
}

/// 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<F: VfsFile>(
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.
Expand Down Expand Up @@ -877,10 +916,10 @@ async fn diagnose_overflow_chain<V: Vfs + Clone>(
#[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]);
Expand Down Expand Up @@ -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<MemVfs>, 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();
Expand Down
61 changes: 60 additions & 1 deletion src/recovery/reconcile.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Expand Down Expand Up @@ -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());
}
}
Loading