Skip to content

Commit ee1034a

Browse files
presempathy-awbfarhan-syah
authored andcommitted
fix(rekey): preserve every durable main page
Re-encrypt durable free-list entries and residual copy-on-write pages before retiring the source epoch. Add regression coverage for free-list integrity, pinned reader snapshots, and monotonic epoch validation.
1 parent 9f87a65 commit ee1034a

2 files changed

Lines changed: 216 additions & 18 deletions

File tree

src/txn/db/rekey/main.rs

Lines changed: 94 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
//! Main-database rekey transition and durable intent publication.
22
3-
use std::collections::BTreeMap;
3+
use std::collections::{BTreeMap, BTreeSet};
44
use std::sync::atomic::Ordering;
55

66
use subtle::ConstantTimeEq;
@@ -11,9 +11,10 @@ use crate::catalog::codec::{Catalog, RekeyIntent, RekeyStage, RekeyStateRow};
1111
use crate::crypto::kdf::{derive_hk, derive_mk};
1212
use crate::crypto::{CipherId, DerivedKey, MasterKey, SecretKey};
1313
use crate::errors::PagedbError;
14+
use crate::pager::PageKind;
1415
use crate::pager::header::commit_header;
1516
use crate::pager::structural_header::MainDbHeaderFields;
16-
use crate::vfs::Vfs;
17+
use crate::vfs::{OpenMode, Vfs, VfsFile};
1718

1819
#[cfg(test)]
1920
use super::super::core::RekeyTestFault;
@@ -27,6 +28,30 @@ use super::intent::{intent_proof, migrate_legacy, validate_intent_for_current_ci
2728
/// retention runs.
2829
const HISTORY_ROOT_BATCH: usize = 512;
2930

31+
async fn read_main_page_kind<F: VfsFile>(file: &F, offset: u64) -> Result<Option<PageKind>> {
32+
let mut envelope = [0u8; 2];
33+
let mut filled = 0;
34+
while filled < envelope.len() {
35+
let read = file
36+
.read_at(offset + filled as u64, &mut envelope[filled..])
37+
.await?;
38+
if read == 0 {
39+
return Err(PagedbError::Io(std::io::Error::from(
40+
std::io::ErrorKind::UnexpectedEof,
41+
)));
42+
}
43+
filled += read;
44+
}
45+
if envelope == [0, 0] {
46+
return Ok(None);
47+
}
48+
let kind = PageKind::from_byte(envelope[1])?;
49+
if !kind.is_main_db() {
50+
return Err(PagedbError::IllegalPageKind);
51+
}
52+
Ok(Some(kind))
53+
}
54+
3055
impl<V: Vfs + Clone> Db<V> {
3156
/// Rekey the reachable main database and every catalog-linked immutable
3257
/// segment under `new_mk_epoch`.
@@ -219,23 +244,8 @@ impl<V: Vfs + Clone> Db<V> {
219244
self.rewrite_retained_history_roots(state, &mut rewritten)
220245
.await?;
221246
}
222-
if state.free_list_root_page_id != 0 {
223-
let (_, chain_pages) = crate::pager::freelist::read_chain(
224-
&self.pager,
225-
self.realm_id,
226-
state.free_list_root_page_id,
227-
)
247+
self.rewrite_free_list_pages(state.free_list_root_page_id)
228248
.await?;
229-
for page_id in chain_pages {
230-
self.pager
231-
.rewrite_page_under_current_epoch(
232-
page_id,
233-
self.realm_id,
234-
crate::pager::PageKind::Free,
235-
)
236-
.await?;
237-
}
238-
}
239249
// Keep the catalog path carrying the intent source-readable until
240250
// target-header publication. This is the durable admission anchor
241251
// for an ordinary open that can verify only the stale A/B side.
@@ -248,6 +258,70 @@ impl<V: Vfs + Clone> Db<V> {
248258
Ok(())
249259
}
250260

261+
/// Re-encrypt both the durable free-list chain and every reusable page it
262+
/// names. Rewriting only the chain would leave those pages sealed under
263+
/// the retired epoch, so a later allocation could not authenticate them.
264+
async fn rewrite_free_list_pages(&self, root_page_id: u64) -> Result<()> {
265+
if root_page_id == 0 {
266+
return Ok(());
267+
}
268+
269+
let (entries, chain_pages) =
270+
crate::pager::freelist::read_chain(&self.pager, self.realm_id, root_page_id).await?;
271+
let mut reusable_pages = entries
272+
.into_iter()
273+
.map(|(_, page_id)| page_id)
274+
.collect::<BTreeSet<_>>();
275+
276+
for page_id in chain_pages {
277+
self.pager
278+
.rewrite_page_under_current_epoch(
279+
page_id,
280+
self.realm_id,
281+
crate::pager::PageKind::Free,
282+
)
283+
.await?;
284+
reusable_pages.remove(&page_id);
285+
}
286+
if reusable_pages.is_empty() {
287+
return Ok(());
288+
}
289+
290+
let file = self.vfs.open(&self.main_db_path, OpenMode::Read).await?;
291+
for page_id in reusable_pages {
292+
let offset = page_id
293+
.checked_mul(self.page_size as u64)
294+
.ok_or_else(|| PagedbError::arithmetic_overflow("free-list page offset"))?;
295+
let Some(kind) = read_main_page_kind(&file, offset).await? else {
296+
continue;
297+
};
298+
self.pager
299+
.rewrite_page_under_current_epoch(page_id, self.realm_id, kind)
300+
.await?;
301+
}
302+
Ok(())
303+
}
304+
305+
/// Re-encrypt pages superseded by copy-on-write catalog transitions after
306+
/// the target-authenticated header is durable. No live root discovers
307+
/// these residual pages, but physical integrity scans still authenticate
308+
/// every non-zero page and the source epoch is retired on success.
309+
async fn rewrite_rekey_residual_main_pages(&self, state: &WriterState) -> Result<()> {
310+
let file = self.vfs.open(&self.main_db_path, OpenMode::Read).await?;
311+
for page_id in 4..state.next_page_id {
312+
let offset = page_id
313+
.checked_mul(self.page_size as u64)
314+
.ok_or_else(|| PagedbError::arithmetic_overflow("main-db page offset"))?;
315+
let Some(kind) = read_main_page_kind(&file, offset).await? else {
316+
continue;
317+
};
318+
self.pager
319+
.rewrite_page_under_current_epoch(page_id, self.realm_id, kind)
320+
.await?;
321+
}
322+
Ok(())
323+
}
324+
251325
/// Rewrite the commit-history index and every reader-visible root its
252326
/// retained rows still name.
253327
///
@@ -353,6 +427,8 @@ impl<V: Vfs + Clone> Db<V> {
353427
target_header_key: &DerivedKey,
354428
) -> Result<()> {
355429
if matches!(intent.stage, RekeyStage::HeaderTargetPublished) {
430+
self.rewrite_rekey_residual_main_pages(state).await?;
431+
self.pager.flush_main(self.realm_id).await?;
356432
intent.stage = RekeyStage::MainDone;
357433
self.write_rekey_intent_locked(
358434
state,

tests/rekey_basic.rs

Lines changed: 122 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -379,3 +379,125 @@ async fn mixed_epoch_pages_readable() {
379379
}
380380
drop(rx2);
381381
}
382+
383+
/// An existing reader keeps its pinned logical snapshot even after rekey
384+
/// rewrites durable main-db pages and clean cache entries are evicted.
385+
#[tokio::test(flavor = "current_thread")]
386+
async fn rekey_preserves_preexisting_main_reader_snapshot_after_cache_evict() {
387+
let (_vfs, db) = fresh_db().await;
388+
let old_large_value = vec![0xA5; PAGE];
389+
{
390+
let mut tx = db.begin_write().await.unwrap();
391+
tx.put(b"versioned", b"before").await.unwrap();
392+
tx.put(b"large", &old_large_value).await.unwrap();
393+
tx.commit().await.unwrap();
394+
}
395+
396+
let reader = db.begin_read().await.unwrap();
397+
{
398+
let mut tx = db.begin_write().await.unwrap();
399+
tx.put(b"versioned", b"after").await.unwrap();
400+
tx.put(b"large", b"replacement").await.unwrap();
401+
tx.commit().await.unwrap();
402+
}
403+
404+
db.rekey_db(KEK0, 1).await.unwrap();
405+
db.evict_main_pages(REALM);
406+
407+
assert_eq!(
408+
reader.get(b"versioned").await.unwrap().as_deref(),
409+
Some(b"before".as_slice())
410+
);
411+
assert_eq!(
412+
reader.get(b"large").await.unwrap().as_deref(),
413+
Some(old_large_value.as_slice())
414+
);
415+
}
416+
417+
/// Rekey must re-encrypt the durable free-list chain and every reusable page
418+
/// it names; otherwise physical integrity checks cannot authenticate those
419+
/// pages after the source epoch is retired.
420+
#[tokio::test(flavor = "current_thread")]
421+
async fn rekey_preserves_durable_free_list_across_reopen() {
422+
let vfs = MemVfs::new();
423+
let options = OpenOptions::default().with_commit_history_retain(RetainPolicy::Disabled);
424+
let db = Db::open_internal_with_options(vfs.clone(), KEK0, PAGE, REALM, options.clone())
425+
.await
426+
.unwrap();
427+
let mut insert = db.begin_write().await.unwrap();
428+
for index in 0..300_u32 {
429+
insert
430+
.put(format!("free-{index:05}").as_bytes(), &[0x7A; 128])
431+
.await
432+
.unwrap();
433+
}
434+
insert.commit().await.unwrap();
435+
let mut delete = db.begin_write().await.unwrap();
436+
for index in 0..250_u32 {
437+
delete
438+
.delete(format!("free-{index:05}").as_bytes())
439+
.await
440+
.unwrap();
441+
}
442+
delete.commit().await.unwrap();
443+
assert!(db.stats().await.unwrap().free_list_pending_entries > 0);
444+
445+
db.rekey_db(KEK0, 1).await.unwrap();
446+
drop(db);
447+
448+
let reopened = Db::open_existing_with_options(vfs, KEK0, PAGE, REALM, options)
449+
.await
450+
.unwrap();
451+
assert!(reopened.stats().await.unwrap().free_list_pending_entries > 0);
452+
let report = pagedb::recovery::deep_walk::run_deep_walk(&reopened)
453+
.await
454+
.unwrap();
455+
assert!(report.is_clean(), "free-list rekey report: {report:?}");
456+
}
457+
458+
/// Epochs are monotonic. Rejecting the current or an older epoch must leave
459+
/// the live store usable and must not persist a rekey intent.
460+
#[tokio::test(flavor = "current_thread")]
461+
async fn rekey_rejects_non_advancing_epoch() {
462+
let (vfs, db) = fresh_db().await;
463+
{
464+
let mut tx = db.begin_write().await.unwrap();
465+
tx.put(b"stable", b"value").await.unwrap();
466+
tx.commit().await.unwrap();
467+
}
468+
469+
assert!(matches!(
470+
db.rekey_db(KEK0, 0).await,
471+
Err(PagedbError::RekeyStateInvalid { .. })
472+
));
473+
assert_eq!(
474+
db.begin_read()
475+
.await
476+
.unwrap()
477+
.get(b"stable")
478+
.await
479+
.unwrap()
480+
.as_deref(),
481+
Some(b"value".as_slice())
482+
);
483+
484+
db.rekey_db(KEK0, 2).await.unwrap();
485+
assert!(matches!(
486+
db.rekey_db(KEK0, 1).await,
487+
Err(PagedbError::RekeyStateInvalid { .. })
488+
));
489+
drop(db);
490+
491+
let reopened = Db::open_existing(vfs, KEK0, PAGE, REALM).await.unwrap();
492+
assert_eq!(
493+
reopened
494+
.begin_read()
495+
.await
496+
.unwrap()
497+
.get(b"stable")
498+
.await
499+
.unwrap()
500+
.as_deref(),
501+
Some(b"value".as_slice())
502+
);
503+
}

0 commit comments

Comments
 (0)