Skip to content

Commit d759bee

Browse files
fix(recovery): harden replay and deep-walk reads
1 parent f9e6cfa commit d759bee

5 files changed

Lines changed: 377 additions & 33 deletions

File tree

src/recovery/deep_walk.rs

Lines changed: 95 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@ use crate::pager::{PageGuard, Pager};
2121
use crate::segment::authenticated_metadata::authenticate_segment_metadata;
2222
use crate::txn::db::Db;
2323
use crate::vfs::types::OpenMode;
24-
use crate::vfs::{Vfs, VfsFile};
24+
use crate::vfs::{Vfs, VfsFile, read_exact_at};
2525

2626
/// A single page-level issue found during deep walk.
2727
#[derive(Debug, Clone)]
@@ -217,7 +217,7 @@ pub async fn run_deep_walk<V: Vfs + Clone>(db: &Db<V>) -> Result<DeepWalkReport>
217217

218218
let vfs: &V = &db.vfs;
219219
let main_file_res = vfs.open(main_db_path, OpenMode::Read).await;
220-
let main_file = match main_file_res {
220+
let mut main_file = match main_file_res {
221221
Ok(f) => f,
222222
Err(e) => {
223223
report.page_issues.push(PageIssue {
@@ -232,24 +232,24 @@ pub async fn run_deep_walk<V: Vfs + Clone>(db: &Db<V>) -> Result<DeepWalkReport>
232232
// (HK-MAC, cleartext) and are already verified by `Db::open`. Skip them.
233233
// Page 2 and 3 are reserved (apply-journal). Walk from page 4.
234234
for page_id in 4..next_page_id {
235-
let offset = page_id * page_size as u64;
236-
let mut buf = vec![0u8; page_size];
237-
match main_file.read_at(offset, &mut buf).await {
238-
Ok(n) if n < page_size => {
239-
// Short read at the tail — the file may be smaller than expected.
240-
// Report but continue.
235+
let offset = match page_offset(page_id, page_size) {
236+
Ok(offset) => offset,
237+
Err(error) => {
241238
report.page_issues.push(PageIssue {
242239
page_id,
243-
description: format!("short read: expected {page_size} bytes, got {n}"),
240+
description: format!("page offset overflow: {error}"),
244241
});
245242
report.pages_examined += 1;
246243
continue;
247244
}
248-
Ok(_) => {}
249-
Err(e) => {
245+
};
246+
let mut buf = vec![0u8; page_size];
247+
match read_exact_at(&mut main_file, offset, &mut buf).await {
248+
Ok(()) => {}
249+
Err(error) => {
250250
report.page_issues.push(PageIssue {
251251
page_id,
252-
description: format!("read error: {e}"),
252+
description: format!("read error: {error}"),
253253
});
254254
report.pages_examined += 1;
255255
continue;
@@ -380,7 +380,7 @@ async fn check_segment<V: Vfs + Clone>(
380380
let page_size = pager.page_size();
381381

382382
// Check file exists.
383-
let Ok(file) = vfs.open(&live, OpenMode::Read).await else {
383+
let Ok(mut file) = vfs.open(&live, OpenMode::Read).await else {
384384
report.drift_issues.push(DriftIssue {
385385
segment_id: meta.segment_id,
386386
description: "segment file missing from seg/".to_string(),
@@ -409,7 +409,16 @@ async fn check_segment<V: Vfs + Clone>(
409409
// We don't have a metadata API, but we can check via read: try reading one
410410
// byte past the expected end. If it succeeds (on some VFS) we skip the
411411
// check; if we read exactly `page_count * page_size` bytes we're consistent.
412-
let expected_size = meta.page_count * page_size as u64;
412+
let expected_size = match page_offset(meta.page_count, page_size) {
413+
Ok(offset) => offset,
414+
Err(error) => {
415+
report.segment_issues.push(SegmentIssue {
416+
segment_id: meta.segment_id,
417+
description: format!("catalog page_count offset overflow: {error}"),
418+
});
419+
return;
420+
}
421+
};
413422
let mut probe = vec![0u8; 1];
414423
let over_read = file.read_at(expected_size, &mut probe).await;
415424
match over_read {
@@ -428,24 +437,23 @@ async fn check_segment<V: Vfs + Clone>(
428437
// Walk data pages (1 .. page_count - 1, skipping header=0 and footer=last).
429438
let last_data = footer_page_id;
430439
for page_id in 1..last_data {
431-
let offset = page_id * page_size as u64;
432-
let mut buf = vec![0u8; page_size];
433-
let read_res = file.read_at(offset, &mut buf).await;
434-
match read_res {
435-
Ok(n) if n < page_size => {
440+
let offset = match page_offset(page_id, page_size) {
441+
Ok(offset) => offset,
442+
Err(error) => {
436443
report.segment_issues.push(SegmentIssue {
437444
segment_id: meta.segment_id,
438-
description: format!(
439-
"short read at page {page_id}: expected {page_size} bytes, got {n}"
440-
),
445+
description: format!("data page {page_id} offset overflow: {error}"),
441446
});
442447
continue;
443448
}
444-
Ok(_) => {}
445-
Err(e) => {
449+
};
450+
let mut buf = vec![0u8; page_size];
451+
match read_exact_at(&mut file, offset, &mut buf).await {
452+
Ok(()) => {}
453+
Err(error) => {
446454
report.segment_issues.push(SegmentIssue {
447455
segment_id: meta.segment_id,
448-
description: format!("read error at page {page_id}: {e}"),
456+
description: format!("read error at page {page_id}: {error}"),
449457
});
450458
continue;
451459
}
@@ -498,6 +506,14 @@ async fn check_segment<V: Vfs + Clone>(
498506
}
499507
}
500508

509+
fn page_offset(page_id: u64, page_size: usize) -> Result<u64> {
510+
let page_size = u64::try_from(page_size)
511+
.map_err(|_| crate::PagedbError::arithmetic_overflow("deep-walk page size"))?;
512+
page_id
513+
.checked_mul(page_size)
514+
.ok_or_else(|| crate::PagedbError::arithmetic_overflow("deep-walk page offset"))
515+
}
516+
501517
/// Collect the set of all page IDs reachable from the main B+ tree root,
502518
/// the catalog root, the commit-history root, and the free-list root.
503519
/// Pages 0..=3 (reserved) are always considered reachable.
@@ -822,10 +838,10 @@ async fn diagnose_overflow_chain<V: Vfs + Clone>(
822838
#[cfg(test)]
823839
mod tests {
824840
use super::*;
825-
use crate::OpenOptions;
826841
use crate::btree::node::body_capacity;
827842
use crate::pager::format::data_page::ENVELOPE_OVERHEAD;
828843
use crate::vfs::memory::MemVfs;
844+
use crate::{OpenOptions, SegmentKind, SegmentPageKind};
829845

830846
const PAGE: usize = 4096;
831847
const REALM: crate::RealmId = crate::RealmId::new([0xD3; 16]);
@@ -1019,6 +1035,59 @@ mod tests {
10191035
);
10201036
}
10211037

1038+
#[tokio::test(flavor = "current_thread")]
1039+
async fn deep_walk_reports_segment_page_count_offset_overflow() {
1040+
let db = open_db().await;
1041+
let mut segment = db
1042+
.create_segment(REALM, SegmentKind::Unspecified)
1043+
.await
1044+
.unwrap();
1045+
segment
1046+
.append_page(SegmentPageKind::Data, b"deep-walk")
1047+
.await
1048+
.unwrap();
1049+
let mut meta = segment.seal().await.unwrap();
1050+
{
1051+
let mut txn = db.begin_write().await.unwrap();
1052+
txn.link_segment("overflow", &meta).await.unwrap();
1053+
txn.commit().await.unwrap();
1054+
}
1055+
1056+
meta.page_count = (u64::MAX / PAGE as u64) + 2;
1057+
let (catalog_root, next_page_id) = {
1058+
let state = db.writer.lock().await;
1059+
(state.catalog_root_page_id, state.next_page_id)
1060+
};
1061+
let mut tree = BTree::open(
1062+
db.pager.clone(),
1063+
db.realm_id,
1064+
catalog_root,
1065+
next_page_id,
1066+
db.page_size,
1067+
);
1068+
let key = Catalog::segment_key(REALM, b"overflow").unwrap();
1069+
tree.put(&key, &Catalog::encode_segment_meta(&meta))
1070+
.await
1071+
.unwrap();
1072+
tree.flush().await.unwrap();
1073+
{
1074+
let mut state = db.writer.lock().await;
1075+
state.catalog_root_page_id = tree.root_page_id();
1076+
state.next_page_id = state.next_page_id.max(tree.next_page_id());
1077+
}
1078+
1079+
let report = run_deep_walk(&db).await.unwrap();
1080+
assert!(
1081+
report.segment_issues.iter().any(|issue| {
1082+
issue.segment_id == meta.segment_id
1083+
&& issue
1084+
.description
1085+
.contains("authenticated segment metadata invalid")
1086+
}),
1087+
"overflowing segment geometry must become a structured issue, got {report:?}"
1088+
);
1089+
}
1090+
10221091
async fn sole_overflow_root(db: &Db<MemVfs>, leaf_page_id: u64) -> u64 {
10231092
let (guard, _) = db.pager.read_main_node(leaf_page_id, REALM).await.unwrap();
10241093
let leaf = Leaf::decode(guard.body_ref()).unwrap();

src/recovery/journal.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -311,7 +311,7 @@ pub async fn execute_journal_actions<V: Vfs>(vfs: &V, actions: &[JournalAction])
311311
crate::hex::to_hex_lower(segment_id),
312312
tombstone_commit_id,
313313
);
314-
if !path_exists(vfs, &dst).await? && path_exists(vfs, &src).await? {
314+
if !path_exists(vfs, &dst).await? {
315315
vfs.rename(&src, &dst).await?;
316316
}
317317
}

src/recovery/reconcile.rs

Lines changed: 48 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -226,7 +226,13 @@ mod tests {
226226
use crate::vfs::{Vfs, VfsFile};
227227
use crate::{RealmId, btree::BTree};
228228

229-
use super::repair_catalog;
229+
use super::{repair_catalog, sweep_orphans};
230+
231+
fn segment_id(value: u64) -> [u8; 16] {
232+
let mut id = [0; 16];
233+
id[..8].copy_from_slice(&value.to_le_bytes());
234+
id
235+
}
230236

231237
#[tokio::test(flavor = "current_thread")]
232238
async fn malformed_catalog_key_prevents_reconciliation_mutation() {
@@ -296,4 +302,45 @@ mod tests {
296302
));
297303
assert!(vfs.open(marker, OpenMode::Read).await.is_ok());
298304
}
305+
306+
#[tokio::test(flavor = "current_thread")]
307+
async fn sweep_orphans_handles_many_expected_segments() {
308+
let vfs = MemVfs::new();
309+
vfs.mkdir_all("seg/.staging").await.unwrap();
310+
let expected: Vec<[u8; 16]> = (0..1024).map(segment_id).collect();
311+
312+
for id in expected.iter().take(8) {
313+
let path = crate::segment::writer::live_path(id);
314+
let mut file = vfs.open(&path, OpenMode::CreateOrOpen).await.unwrap();
315+
file.write_at(0, b"live").await.unwrap();
316+
}
317+
318+
let live_orphan = segment_id(10_000);
319+
let live_orphan_path = crate::segment::writer::live_path(&live_orphan);
320+
let mut live_file = vfs
321+
.open(&live_orphan_path, OpenMode::CreateOrOpen)
322+
.await
323+
.unwrap();
324+
live_file.write_at(0, b"orphan").await.unwrap();
325+
326+
let staged_orphan = segment_id(10_001);
327+
let staged_orphan_path = crate::segment::writer::staging_path(&staged_orphan);
328+
let mut staged_file = vfs
329+
.open(&staged_orphan_path, OpenMode::CreateOrOpen)
330+
.await
331+
.unwrap();
332+
staged_file.write_at(0, b"orphan").await.unwrap();
333+
334+
sweep_orphans(&vfs, &expected, 77).await.unwrap();
335+
336+
let expected_live = crate::segment::writer::live_path(&expected[0]);
337+
assert!(vfs.open(&expected_live, OpenMode::Read).await.is_ok());
338+
assert!(vfs.open(&live_orphan_path, OpenMode::Read).await.is_err());
339+
let tombstone = format!(
340+
"seg/.tombstone/{}.77",
341+
crate::hex::to_hex_lower(&live_orphan)
342+
);
343+
assert!(vfs.open(&tombstone, OpenMode::Read).await.is_ok());
344+
assert!(vfs.open(&staged_orphan_path, OpenMode::Read).await.is_err());
345+
}
299346
}

tests/durability/apply_journal_crash.rs

Lines changed: 47 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
1-
/// Verify apply-journal encode/decode round-trips and idempotent action replay.
2-
/// These tests exercise the journal machinery at the level of the public
3-
/// encode/decode API and the VFS rename layer without requiring a full crash.
1+
//! Verify apply-journal encoding and idempotent action replay at the VFS
2+
//! boundary without requiring a full process crash.
3+
4+
use pagedb::PagedbError;
45
use pagedb::recovery::journal::{ApplyJournalRecord, JournalAction, decode_record, encode_record};
56
use pagedb::vfs::Vfs;
67
use pagedb::vfs::memory::MemVfs;
@@ -134,13 +135,54 @@ async fn promote_action_renames_staging_to_live() {
134135
}
135136

136137
#[tokio::test(flavor = "current_thread")]
137-
async fn tombstone_action_is_idempotent_when_live_absent() {
138-
// If the live file is absent, the tombstone rename is a no-op.
138+
async fn tombstone_action_errors_when_live_and_tombstone_are_absent() {
139139
use pagedb::recovery::journal::execute_journal_actions;
140+
140141
let vfs = MemVfs::new();
141142
let actions = vec![JournalAction::Tombstone {
142143
segment_id: [0xEE; 16],
143144
tombstone_commit_id: 5,
144145
}];
146+
let err = execute_journal_actions(&vfs, &actions)
147+
.await
148+
.expect_err("missing tombstone source is only idempotent when the destination exists");
149+
assert!(
150+
matches!(err, PagedbError::Io(ref io) if io.kind() == std::io::ErrorKind::NotFound),
151+
"expected missing tombstone source and destination to surface NotFound, got {err:?}"
152+
);
153+
}
154+
155+
#[tokio::test(flavor = "current_thread")]
156+
async fn tombstone_action_is_idempotent_when_tombstone_exists() {
157+
use pagedb::recovery::journal::execute_journal_actions;
158+
use pagedb::vfs::VfsFile;
159+
use pagedb::vfs::types::OpenMode;
160+
161+
let vfs = MemVfs::new();
162+
let segment_id = [0xEF; 16];
163+
let tombstone = format!(
164+
"seg/.tombstone/{}.5",
165+
segment_id
166+
.iter()
167+
.map(|byte| format!("{byte:02x}"))
168+
.collect::<String>()
169+
);
170+
vfs.mkdir_all("seg/.tombstone").await.unwrap();
171+
{
172+
let mut file = vfs.open(&tombstone, OpenMode::CreateNew).await.unwrap();
173+
file.write_at(0, b"already-tombstoned").await.unwrap();
174+
file.sync().await.unwrap();
175+
}
176+
177+
let actions = vec![JournalAction::Tombstone {
178+
segment_id,
179+
tombstone_commit_id: 5,
180+
}];
145181
execute_journal_actions(&vfs, &actions).await.unwrap();
182+
183+
let file = vfs.open(&tombstone, OpenMode::Read).await.unwrap();
184+
let mut bytes = vec![0; 18];
185+
let read = file.read_at(0, &mut bytes).await.unwrap();
186+
assert_eq!(read, bytes.len());
187+
assert_eq!(&bytes, b"already-tombstoned");
146188
}

0 commit comments

Comments
 (0)