Skip to content

Commit 8fdb14e

Browse files
fix(catalog): reject malformed persisted metadata
1 parent 5635bee commit 8fdb14e

4 files changed

Lines changed: 217 additions & 19 deletions

File tree

src/txn/db/catalog.rs

Lines changed: 99 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -44,14 +44,41 @@ impl<V: Vfs + Clone> Db<V> {
4444
let Some(key) = hist.first_key().await? else {
4545
return Ok(None);
4646
};
47-
if key.len() < 8 {
48-
return Ok(None);
47+
if key.len() != 8 {
48+
return Err(PagedbError::catalog_row_invalid("commit_history.key"));
4949
}
5050
let mut b = [0u8; 8];
5151
b.copy_from_slice(&key[..8]);
5252
Ok(Some(u64::from_be_bytes(b)))
5353
}
5454

55+
/// Authenticate and decode every persisted named-counter row during open.
56+
///
57+
/// Named counters are already atomic with catalog-root publication, so
58+
/// recovery validates their encoding but never rewrites their values.
59+
pub(super) async fn validate_counter_rows(
60+
&self,
61+
catalog_root_page_id: u64,
62+
next_page_id: u64,
63+
) -> Result<()> {
64+
if catalog_root_page_id == 0 {
65+
return Ok(());
66+
}
67+
68+
let prefix = [crate::catalog::codec::CatalogRowKind::Counter as u8];
69+
let tree = BTree::open(
70+
self.pager.clone(),
71+
self.realm_id,
72+
catalog_root_page_id,
73+
next_page_id,
74+
self.page_size,
75+
);
76+
for (_key, value) in tree.scan_prefix(&prefix).await? {
77+
Catalog::decode_counter(&value)?;
78+
}
79+
Ok(())
80+
}
81+
5582
/// Write per-realm quota caps into the catalog B+ tree and persist the
5683
/// updated catalog root to the A/B header.
5784
pub async fn set_realm_quotas(&self, realm: RealmId, quotas: RealmQuotas) -> Result<()> {
@@ -307,3 +334,73 @@ impl<V: Vfs + Clone> Db<V> {
307334
Ok(freed)
308335
}
309336
}
337+
338+
#[cfg(test)]
339+
mod tests {
340+
use crate::vfs::memory::MemVfs;
341+
use crate::{Db, PagedbError, RealmId};
342+
343+
use super::*;
344+
345+
const PAGE: usize = 4096;
346+
const REALM: RealmId = RealmId::new([0xA7; 16]);
347+
348+
#[tokio::test(flavor = "current_thread")]
349+
async fn counter_recovery_surfaces_malformed_counter_row() {
350+
let db = Db::open_internal(MemVfs::new(), [9u8; 32], PAGE, REALM)
351+
.await
352+
.unwrap();
353+
{
354+
let mut txn = db.begin_write().await.unwrap();
355+
let mut counter = txn.counter("bad-counter").unwrap();
356+
counter.set(5).await.unwrap();
357+
drop(counter);
358+
txn.commit().await.unwrap();
359+
}
360+
361+
let (catalog_root, next_page_id) = {
362+
let state = db.writer.lock().await;
363+
(state.catalog_root_page_id, state.next_page_id)
364+
};
365+
let mut tree = BTree::open(
366+
db.pager.clone(),
367+
db.realm_id,
368+
catalog_root,
369+
next_page_id,
370+
db.page_size,
371+
);
372+
tree.put(&Catalog::counter_key(&[0xFF]).unwrap(), b"bad")
373+
.await
374+
.unwrap();
375+
tree.flush().await.unwrap();
376+
377+
let err = db
378+
.validate_counter_rows(tree.root_page_id(), tree.next_page_id())
379+
.await
380+
.expect_err("malformed counter row must surface during recovery validation");
381+
assert!(matches!(err, PagedbError::Corruption(_)));
382+
}
383+
384+
#[tokio::test(flavor = "current_thread")]
385+
async fn oldest_retained_history_commit_surfaces_malformed_history_key() {
386+
for malformed_key in [b"x".as_slice(), b"123456789".as_slice()] {
387+
let db = Db::open_internal(MemVfs::new(), [9u8; 32], PAGE, REALM)
388+
.await
389+
.unwrap();
390+
let next_page_id = db.writer.lock().await.next_page_id;
391+
let mut history =
392+
BTree::open(db.pager.clone(), db.realm_id, 0, next_page_id, db.page_size);
393+
history
394+
.put(malformed_key, b"malformed history")
395+
.await
396+
.unwrap();
397+
history.flush().await.unwrap();
398+
399+
let err = db
400+
.oldest_retained_history_commit(history.root_page_id(), history.next_page_id())
401+
.await
402+
.expect_err("malformed history key must surface");
403+
assert!(matches!(err, PagedbError::Corruption(_)));
404+
}
405+
}
406+
}

src/txn/db/misc.rs

Lines changed: 68 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -160,11 +160,11 @@ impl<V: Vfs + Clone> Db<V> {
160160
let seg_start = vec![CatalogRowKind::Segment as u8];
161161
let seg_rows = tree.scan_prefix(&seg_start).await?;
162162
let seg_count = u32::try_from(seg_rows.len()).unwrap_or(u32::MAX);
163-
let seg_bytes: u64 = seg_rows
164-
.iter()
165-
.filter_map(|(_k, v)| Catalog::decode_segment_meta(v).ok())
166-
.map(|m| m.total_bytes)
167-
.sum();
163+
let mut seg_bytes = 0u64;
164+
for (_key, value) in &seg_rows {
165+
seg_bytes =
166+
seg_bytes.saturating_add(Catalog::decode_segment_meta(value)?.total_bytes);
167+
}
168168

169169
(seg_count, seg_bytes)
170170
};
@@ -189,3 +189,66 @@ impl<V: Vfs + Clone> Db<V> {
189189
})
190190
}
191191
}
192+
193+
#[cfg(test)]
194+
mod tests {
195+
use crate::btree::BTree;
196+
use crate::catalog::codec::Catalog;
197+
use crate::vfs::memory::MemVfs;
198+
use crate::{Db, PagedbError, RealmId, SegmentKind, SegmentPageKind};
199+
200+
const PAGE: usize = 4096;
201+
const REALM: RealmId = RealmId::new([0xA5; 16]);
202+
203+
#[tokio::test(flavor = "current_thread")]
204+
async fn stats_surfaces_malformed_segment_catalog_row() {
205+
let db = Db::open_internal(MemVfs::new(), [9u8; 32], PAGE, REALM)
206+
.await
207+
.unwrap();
208+
let mut segment = db
209+
.create_segment(REALM, SegmentKind::Unspecified)
210+
.await
211+
.unwrap();
212+
segment
213+
.append_page(SegmentPageKind::Data, b"stats")
214+
.await
215+
.unwrap();
216+
let meta = segment.seal().await.unwrap();
217+
{
218+
let mut txn = db.begin_write().await.unwrap();
219+
txn.link_segment("good", &meta).await.unwrap();
220+
txn.commit().await.unwrap();
221+
}
222+
223+
let (catalog_root, next_page_id) = {
224+
let state = db.writer.lock().await;
225+
(state.catalog_root_page_id, state.next_page_id)
226+
};
227+
let mut tree = BTree::open(
228+
db.pager.clone(),
229+
db.realm_id,
230+
catalog_root,
231+
next_page_id,
232+
db.page_size,
233+
);
234+
tree.put(
235+
&Catalog::segment_key(REALM, b"bad").unwrap(),
236+
b"not a segment meta",
237+
)
238+
.await
239+
.unwrap();
240+
tree.flush().await.unwrap();
241+
{
242+
let mut state = db.writer.lock().await;
243+
state.catalog_root_page_id = tree.root_page_id();
244+
state.next_page_id = state.next_page_id.max(tree.next_page_id());
245+
db.publish_snapshot(&state);
246+
}
247+
248+
let err = db
249+
.stats()
250+
.await
251+
.expect_err("malformed segment metadata must surface");
252+
assert!(matches!(err, PagedbError::Corruption(_)));
253+
}
254+
}

src/txn/db/open/recovery.rs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -72,5 +72,8 @@ pub(super) async fn recover_open_state<V: Vfs + Clone>(
7272
.await?;
7373
}
7474

75+
db.validate_counter_rows(catalog_root_page_id, next_page_id)
76+
.await?;
77+
7578
Ok(())
7679
}

src/txn/db/reader.rs

Lines changed: 47 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -262,17 +262,18 @@ impl<V: Vfs + Clone> Db<V> {
262262
false,
263263
))
264264
} else {
265-
// Find the oldest available: scan the whole history tree. An
266-
// exclusive `u64::MAX` upper bound would hide the newest commit
267-
// once ids reach the top of the range.
268-
let oldest = tree.collect_all().await?.into_iter().next().map_or(
269-
CommitId(latest_commit_id),
270-
|(k, _)| {
271-
let mut b = [0u8; 8];
272-
b.copy_from_slice(&k[..8]);
273-
CommitId(u64::from_be_bytes(b))
274-
},
275-
);
265+
// Only the leftmost key is needed. This avoids scanning retained
266+
// history while still validating its fixed-width key encoding.
267+
let oldest = if let Some(key) = tree.first_key().await? {
268+
if key.len() != 8 {
269+
return Err(PagedbError::catalog_row_invalid("commit_history.key"));
270+
}
271+
let mut bytes = [0u8; 8];
272+
bytes.copy_from_slice(&key[..8]);
273+
CommitId(u64::from_be_bytes(bytes))
274+
} else {
275+
CommitId(latest_commit_id)
276+
};
276277
Err(PagedbError::CommitGone {
277278
commit,
278279
oldest_available: oldest,
@@ -285,17 +286,51 @@ impl<V: Vfs + Clone> Db<V> {
285286
mod tests {
286287
use std::sync::Arc;
287288

289+
use crate::btree::BTree;
288290
use crate::catalog::codec::SegmentKind;
291+
use crate::errors::PagedbError;
289292
use crate::segment::types::SegmentPageKind;
290293
use crate::vfs::memory::MemVfs;
291-
use crate::{Db, RealmId};
294+
use crate::{CommitId, Db, RealmId};
292295

293296
use super::super::core::VisibilityTestHook;
294297

295298
const PAGE: usize = 4096;
296299
const KEK: [u8; 32] = [0xD1; 32];
297300
const REALM: RealmId = RealmId::new([0xD2; 16]);
298301

302+
#[tokio::test(flavor = "current_thread")]
303+
async fn begin_read_at_surfaces_malformed_oldest_history_key() {
304+
for malformed_key in [b"x".as_slice(), b"123456789".as_slice()] {
305+
let db = Db::open_internal(MemVfs::new(), [9u8; 32], PAGE, REALM)
306+
.await
307+
.unwrap();
308+
let next_page_id = db.writer.lock().await.next_page_id;
309+
let mut history =
310+
BTree::open(db.pager.clone(), db.realm_id, 0, next_page_id, db.page_size);
311+
history
312+
.put(malformed_key, b"malformed history")
313+
.await
314+
.unwrap();
315+
history.flush().await.unwrap();
316+
{
317+
let mut state = db.writer.lock().await;
318+
state.commit_history_root_page_id = history.root_page_id();
319+
state.next_page_id = state.next_page_id.max(history.next_page_id());
320+
db.publish_snapshot(&state);
321+
}
322+
323+
let err = match db.begin_read_at(CommitId::new(42)).await {
324+
Ok(read) => {
325+
drop(read);
326+
panic!("malformed oldest history key must not produce a reader");
327+
}
328+
Err(err) => err,
329+
};
330+
assert!(matches!(err, PagedbError::Corruption(_)));
331+
}
332+
}
333+
299334
async fn linked_segment(db: &Db<MemVfs>, name: &str, bytes: &[u8]) {
300335
let mut writer = db
301336
.create_segment(REALM, SegmentKind::Unspecified)

0 commit comments

Comments
 (0)