diff --git a/src/pager/cache.rs b/src/pager/cache.rs index 1f618b7..6cfa73c 100644 --- a/src/pager/cache.rs +++ b/src/pager/cache.rs @@ -269,6 +269,16 @@ impl PageCache { self.dirty.contains(&key) } + /// How many entries are currently dirty, across every file. + /// + /// Dirty entries are never evicted, so this is the part of the cache that + /// can push it past its configured capacity. Callers that seal pages in a + /// long loop watch this to decide when they must flush. + #[must_use] + pub fn dirty_len(&self) -> usize { + self.dirty.len() + } + /// Sorted iterator over the dirty page ids for one file, ascending by /// `page_id`. Used by the Pager to flush in physical-id order. #[must_use] diff --git a/src/pager/core.rs b/src/pager/core.rs index 64ada90..c707c7b 100644 --- a/src/pager/core.rs +++ b/src/pager/core.rs @@ -608,6 +608,15 @@ impl Pager { /// /// `read_main_page` now always routes to the on-disk epoch/cipher, so this /// simply reads then marks dirty. + /// + /// The dirty set is kept inside the configured buffer-pool budget. A rekey + /// re-seals every page a root can reach, and dirty entries are never + /// evicted (the cache's eviction scan skips them), so without an intermediate + /// flush the walk would pin the whole decrypted database in memory no + /// matter what budget the caller asked for. Flushing mid-walk is safe and + /// idempotent: the active epoch is already the target, so a page re-read + /// after being re-sealed opens under the same key, and a crash resumes + /// from the durable rekey intent. pub async fn rewrite_page_under_current_epoch( &self, page_id: u64, @@ -618,10 +627,15 @@ impl Pager { .read_main_page(page_id, realm_id, expected_kind) .await?; let file = FileKey::Main; - let mut cache = self.inner.cache_for_key(file).lock(); - cache.mark_dirty((file, page_id)); - drop(cache); + let over_budget = { + let mut cache = self.inner.cache_for_key(file).lock(); + cache.mark_dirty((file, page_id)); + cache.dirty_len() >= self.cfg.buffer_pool_pages.max(1) + }; drop(guard); + if over_budget { + self.flush_main(realm_id).await?; + } Ok(()) } diff --git a/src/txn/db/rekey/main.rs b/src/txn/db/rekey/main.rs index 53dfc24..5fb7b73 100644 --- a/src/txn/db/rekey/main.rs +++ b/src/txn/db/rekey/main.rs @@ -1,6 +1,6 @@ //! Main-database rekey transition and durable intent publication. -use std::collections::BTreeMap; +use std::collections::{BTreeMap, BTreeSet}; use std::sync::atomic::Ordering; use subtle::ConstantTimeEq; @@ -11,9 +11,10 @@ use crate::catalog::codec::{Catalog, RekeyIntent, RekeyStage, RekeyStateRow}; use crate::crypto::kdf::{derive_hk, derive_mk}; use crate::crypto::{CipherId, DerivedKey, MasterKey, SecretKey}; use crate::errors::PagedbError; +use crate::pager::PageKind; use crate::pager::header::commit_header; use crate::pager::structural_header::MainDbHeaderFields; -use crate::vfs::Vfs; +use crate::vfs::{OpenMode, Vfs, VfsFile, read_exact_at}; #[cfg(test)] use super::super::core::RekeyTestFault; @@ -27,6 +28,40 @@ use super::intent::{intent_proof, migrate_legacy, validate_intent_for_current_ci /// retention runs. const HISTORY_ROOT_BATCH: usize = 512; +/// One catalog rewrite a rekey stage is about to make durable. +/// +/// The new root and allocation cursor, the pages the rewrite superseded, and +/// the segment side effects it publishes all describe the same transition, so +/// they travel together rather than as loose positional arguments. +pub(super) struct RekeyCatalogCommit<'a> { + pub catalog_root_page_id: u64, + pub next_page_id: u64, + pub freed_pages: &'a [u64], + pub effects: &'a [crate::txn::write::SegmentSideEffect], +} + +/// Read the cleartext kind byte of a `main.db` page, or `None` when the page +/// has been cleared. +/// +/// Only selects which AAD binding to authenticate under; the pager still +/// verifies the full envelope against that same kind, so a tampered byte +/// cannot launder a page into another role — it fails the tag instead. +async fn read_main_page_kind( + file: &mut F, + offset: u64, +) -> Result> { + let mut envelope = [0u8; 2]; + read_exact_at(file, offset, &mut envelope).await?; + if envelope == [0, 0] { + return Ok(None); + } + let kind = PageKind::from_byte(envelope[1])?; + if !kind.is_main_db() { + return Err(PagedbError::IllegalPageKind); + } + Ok(Some(kind)) +} + impl Db { /// Rekey the reachable main database and every catalog-linked immutable /// segment under `new_mk_epoch`. @@ -219,23 +254,8 @@ impl Db { self.rewrite_retained_history_roots(state, &mut rewritten) .await?; } - if state.free_list_root_page_id != 0 { - let (_, chain_pages) = crate::pager::freelist::read_chain( - &self.pager, - self.realm_id, - state.free_list_root_page_id, - ) + self.rewrite_free_list_pages(state.free_list_root_page_id) .await?; - for page_id in chain_pages { - self.pager - .rewrite_page_under_current_epoch( - page_id, - self.realm_id, - crate::pager::PageKind::Free, - ) - .await?; - } - } // Keep the catalog path carrying the intent source-readable until // target-header publication. This is the durable admission anchor // for an ordinary open that can verify only the stale A/B side. @@ -248,6 +268,77 @@ impl Db { Ok(()) } + /// Re-encrypt both the durable free-list chain and every reusable page it + /// names. Rewriting only the chain would leave those pages sealed under + /// the retired epoch, so a later allocation could not authenticate them. + async fn rewrite_free_list_pages(&self, root_page_id: u64) -> Result<()> { + if root_page_id == 0 { + return Ok(()); + } + + let (entries, chain_pages) = + crate::pager::freelist::read_chain(&self.pager, self.realm_id, root_page_id).await?; + let mut reusable_pages = entries + .into_iter() + .map(|(_, page_id)| page_id) + .collect::>(); + + for page_id in chain_pages { + self.pager + .rewrite_page_under_current_epoch( + page_id, + self.realm_id, + crate::pager::PageKind::Free, + ) + .await?; + reusable_pages.remove(&page_id); + } + if reusable_pages.is_empty() { + return Ok(()); + } + + let mut file = self.vfs.open(&self.main_db_path, OpenMode::Read).await?; + for page_id in reusable_pages { + let offset = page_id + .checked_mul(self.page_size as u64) + .ok_or_else(|| PagedbError::arithmetic_overflow("free-list page offset"))?; + let Some(kind) = read_main_page_kind(&mut file, offset).await? else { + continue; + }; + self.pager + .rewrite_page_under_current_epoch(page_id, self.realm_id, kind) + .await?; + } + Ok(()) + } + + /// Re-seal the live catalog tree under the target epoch. + /// + /// Every other reader-visible root is re-sealed before the target header is + /// published. The catalog is the one that cannot be: until the target + /// header is durable it must stay source-readable, because it carries the + /// rekey intent that an open able to verify only the stale A/B side uses to + /// admit recovery. So it is re-sealed here instead, once that anchor is in + /// place and while both keys are still installed. + /// + /// The walk is the same bounded, deduplicating traversal the data tree + /// uses — proportional to unique reachable catalog pages, not to the size + /// of the file. + async fn rewrite_rekey_catalog_pages(&self, state: &WriterState) -> Result<()> { + if state.catalog_root_page_id == 0 { + return Ok(()); + } + let catalog = BTree::open( + self.pager.clone(), + self.realm_id, + state.catalog_root_page_id, + state.next_page_id, + self.page_size, + ); + catalog.rekey_walk_unique(&mut BTreeMap::new()).await?; + Ok(()) + } + /// Rewrite the commit-history index and every reader-visible root its /// retained rows still name. /// @@ -353,6 +444,14 @@ impl Db { target_header_key: &DerivedKey, ) -> Result<()> { if matches!(intent.stage, RekeyStage::HeaderTargetPublished) { + self.rewrite_rekey_catalog_pages(state).await?; + // The pages the intent write above superseded are already on the + // durable free list, so this pass also re-seals them; from here on + // every catalog page is target-sealed and later transitions can + // only supersede target-epoch pages. + self.rewrite_free_list_pages(state.free_list_root_page_id) + .await?; + self.pager.flush_main(self.realm_id).await?; intent.stage = RekeyStage::MainDone; self.write_rekey_intent_locked( state, @@ -444,13 +543,17 @@ impl Db { ) .await?; tree.flush().await?; + let freed_pages = tree.drain_freed(); self.commit_rekey_catalog_root( state, - tree.root_page_id(), - tree.next_page_id(), + RekeyCatalogCommit { + catalog_root_page_id: tree.root_page_id(), + next_page_id: tree.next_page_id(), + freed_pages: &freed_pages, + effects: &[], + }, header_epoch, header_hk, - &[], ) .await .map(|_| ()) @@ -474,27 +577,85 @@ impl Db { ); let _ = tree.delete(&Catalog::rekey_state_key()).await?; tree.flush().await?; + let freed_pages = tree.drain_freed(); self.commit_rekey_catalog_root( state, - tree.root_page_id(), - tree.next_page_id(), + RekeyCatalogCommit { + catalog_root_page_id: tree.root_page_id(), + next_page_id: tree.next_page_id(), + freed_pages: &freed_pages, + effects: &[], + }, header_epoch, header_hk, - &[], ) .await .map(|_| ()) } - pub(super) async fn commit_rekey_catalog_root( + /// Fold pages a rekey-time catalog rewrite superseded into the durable free + /// list, returning the allocation cursor the header must record. + /// + /// Every ordinary commit routes its copy-on-write leftovers here. Without + /// it a rekey abandons pages at each intent transition: unreachable from + /// any root, absent from the free list, and — once the source epoch is + /// retired — sealed under a key that no longer exists. Entries carry the + /// current commit id, so the reclamation floor still withholds any page a + /// live reader or retained history root can still name. + async fn record_rekey_freed_pages( &self, state: &mut WriterState, - catalog_root_page_id: u64, + freed_pages: &[u64], next_page_id: u64, + ) -> Result { + if freed_pages.is_empty() { + return Ok(next_page_id); + } + let (mut entries, old_chain) = crate::pager::freelist::read_chain( + &self.pager, + self.realm_id, + state.free_list_root_page_id, + ) + .await?; + entries.extend( + freed_pages + .iter() + .map(|&page_id| (state.latest_commit_id, page_id)), + ); + // Chain pages are writer-only metadata that no reader snapshot ever + // traverses, so they carry a commit id below every real floor and are + // immediately recyclable — the same rule the ordinary commit path uses. + entries.extend(old_chain.into_iter().map(|page_id| (0, page_id))); + let (new_free_list_root, new_next_page_id) = crate::pager::freelist::rewrite_chain( + &self.pager, + self.realm_id, + self.page_size, + entries, + Vec::new(), + next_page_id, + ) + .await?; + state.free_list_root_page_id = new_free_list_root; + self.pager.flush_main(self.realm_id).await?; + Ok(new_next_page_id) + } + + pub(super) async fn commit_rekey_catalog_root( + &self, + state: &mut WriterState, + catalog: RekeyCatalogCommit<'_>, header_epoch: u64, hk: &DerivedKey, - effects: &[crate::txn::write::SegmentSideEffect], ) -> Result { + let RekeyCatalogCommit { + catalog_root_page_id, + next_page_id, + freed_pages, + effects, + } = catalog; + let next_page_id = self + .record_rekey_freed_pages(state, freed_pages, next_page_id) + .await?; let next_page_id = next_page_id.max(state.next_page_id); let new_seq = state .seq @@ -877,13 +1038,17 @@ mod tests { .unwrap(); catalog.flush().await.unwrap(); let source_header_key = db.hk.read().clone(); + let freed_pages = catalog.drain_freed(); db.commit_rekey_catalog_root( &mut state, - catalog.root_page_id(), - catalog.next_page_id(), + RekeyCatalogCommit { + catalog_root_page_id: catalog.root_page_id(), + next_page_id: catalog.next_page_id(), + freed_pages: &freed_pages, + effects: &[], + }, 0, &source_header_key, - &[], ) .await .unwrap(); diff --git a/src/txn/db/rekey/segments.rs b/src/txn/db/rekey/segments.rs index 6e15b7f..e225724 100644 --- a/src/txn/db/rekey/segments.rs +++ b/src/txn/db/rekey/segments.rs @@ -18,6 +18,7 @@ use crate::vfs::Vfs; #[cfg(test)] use super::super::core::RekeyTestFault; use super::super::core::{Db, WriterState}; +use super::main::RekeyCatalogCommit; struct SegmentEntry { key: Vec, @@ -249,6 +250,7 @@ impl Db { tree.put(&source.key, &Catalog::encode_segment_meta(replacement)) .await?; tree.flush().await?; + let freed_pages = tree.drain_freed(); let effects = [ SegmentSideEffect::Promote { segment_id: replacement.segment_id, @@ -260,11 +262,14 @@ impl Db { ]; self.commit_rekey_catalog_root( state, - tree.root_page_id(), - tree.next_page_id(), + RekeyCatalogCommit { + catalog_root_page_id: tree.root_page_id(), + next_page_id: tree.next_page_id(), + freed_pages: &freed_pages, + effects: &effects, + }, intent.target_mk_epoch, target_hk, - &effects, ) .await .map(|_| ()) @@ -285,13 +290,17 @@ impl Db { ) .await?; tree.flush().await?; + let freed_pages = tree.drain_freed(); self.commit_rekey_catalog_root( state, - tree.root_page_id(), - tree.next_page_id(), + RekeyCatalogCommit { + catalog_root_page_id: tree.root_page_id(), + next_page_id: tree.next_page_id(), + freed_pages: &freed_pages, + effects: &[], + }, header_epoch, header_hk, - &[], ) .await?; #[cfg(test)] @@ -311,13 +320,17 @@ impl Db { .delete(&Catalog::rekey_segment_progress_key(source_id)) .await?; tree.flush().await?; + let freed_pages = tree.drain_freed(); self.commit_rekey_catalog_root( state, - tree.root_page_id(), - tree.next_page_id(), + RekeyCatalogCommit { + catalog_root_page_id: tree.root_page_id(), + next_page_id: tree.next_page_id(), + freed_pages: &freed_pages, + effects: &[], + }, header_epoch, header_hk, - &[], ) .await?; #[cfg(test)] diff --git a/tests/rekey_basic.rs b/tests/rekey_basic.rs index 2e95327..4922d53 100644 --- a/tests/rekey_basic.rs +++ b/tests/rekey_basic.rs @@ -379,3 +379,230 @@ async fn mixed_epoch_pages_readable() { } drop(rx2); } + +/// An existing reader keeps its pinned logical snapshot even after rekey +/// rewrites durable main-db pages and clean cache entries are evicted. +#[tokio::test(flavor = "current_thread")] +async fn rekey_preserves_preexisting_main_reader_snapshot_after_cache_evict() { + let (_vfs, db) = fresh_db().await; + let old_large_value = vec![0xA5; PAGE]; + { + let mut tx = db.begin_write().await.unwrap(); + tx.put(b"versioned", b"before").await.unwrap(); + tx.put(b"large", &old_large_value).await.unwrap(); + tx.commit().await.unwrap(); + } + + let reader = db.begin_read().await.unwrap(); + { + let mut tx = db.begin_write().await.unwrap(); + tx.put(b"versioned", b"after").await.unwrap(); + tx.put(b"large", b"replacement").await.unwrap(); + tx.commit().await.unwrap(); + } + + db.rekey_db(KEK0, 1).await.unwrap(); + db.evict_main_pages(REALM); + + assert_eq!( + reader.get(b"versioned").await.unwrap().as_deref(), + Some(b"before".as_slice()) + ); + assert_eq!( + reader.get(b"large").await.unwrap().as_deref(), + Some(old_large_value.as_slice()) + ); +} + +/// Rekey must re-encrypt the durable free-list chain and every reusable page +/// it names; otherwise physical integrity checks cannot authenticate those +/// pages after the source epoch is retired. +#[tokio::test(flavor = "current_thread")] +async fn rekey_preserves_durable_free_list_across_reopen() { + let vfs = MemVfs::new(); + let options = OpenOptions::default().with_commit_history_retain(RetainPolicy::Disabled); + let db = Db::open_internal_with_options(vfs.clone(), KEK0, PAGE, REALM, options.clone()) + .await + .unwrap(); + let mut insert = db.begin_write().await.unwrap(); + for index in 0..300_u32 { + insert + .put(format!("free-{index:05}").as_bytes(), &[0x7A; 128]) + .await + .unwrap(); + } + insert.commit().await.unwrap(); + let mut delete = db.begin_write().await.unwrap(); + for index in 0..250_u32 { + delete + .delete(format!("free-{index:05}").as_bytes()) + .await + .unwrap(); + } + delete.commit().await.unwrap(); + assert!(db.stats().await.unwrap().free_list_pending_entries > 0); + + db.rekey_db(KEK0, 1).await.unwrap(); + drop(db); + + let reopened = Db::open_existing_with_options(vfs, KEK0, PAGE, REALM, options) + .await + .unwrap(); + assert!(reopened.stats().await.unwrap().free_list_pending_entries > 0); + let report = pagedb::recovery::deep_walk::run_deep_walk(&reopened) + .await + .unwrap(); + assert!(report.is_clean(), "free-list rekey report: {report:?}"); +} + +/// Epochs are monotonic. Rejecting the current or an older epoch must leave +/// the live store usable and must not persist a rekey intent. +#[tokio::test(flavor = "current_thread")] +async fn rekey_rejects_non_advancing_epoch() { + let (vfs, db) = fresh_db().await; + { + let mut tx = db.begin_write().await.unwrap(); + tx.put(b"stable", b"value").await.unwrap(); + tx.commit().await.unwrap(); + } + + assert!(matches!( + db.rekey_db(KEK0, 0).await, + Err(PagedbError::RekeyStateInvalid { .. }) + )); + assert_eq!( + db.begin_read() + .await + .unwrap() + .get(b"stable") + .await + .unwrap() + .as_deref(), + Some(b"value".as_slice()) + ); + + db.rekey_db(KEK0, 2).await.unwrap(); + assert!(matches!( + db.rekey_db(KEK0, 1).await, + Err(PagedbError::RekeyStateInvalid { .. }) + )); + drop(db); + + let reopened = Db::open_existing(vfs, KEK0, PAGE, REALM).await.unwrap(); + assert_eq!( + reopened + .begin_read() + .await + .unwrap() + .get(b"stable") + .await + .unwrap() + .as_deref(), + Some(b"value".as_slice()) + ); +} + +/// A catalog too large to fit one page must be re-sealed in full. +/// +/// The catalog is the one reader-visible root that cannot be re-sealed before +/// the target header is published — it carries the intent that admits recovery +/// from a stale header. Only the pages a rekey happens to write through get +/// re-sealed by that write, so a catalog spanning several pages is where a +/// missing catalog traversal shows up: everything off the written path stays +/// sealed under an epoch that retirement then destroys. +fn counter_name(index: u32) -> String { + // Long names on purpose: the catalog must span more pages than any single + // copy-on-write path through it touches. + format!("catalog-counter-{index:05}-{}", "n".repeat(200)) +} + +#[tokio::test(flavor = "current_thread")] +async fn rekey_reseals_a_catalog_larger_than_one_page() { + let vfs = MemVfs::new(); + // Retention disabled on purpose: with history retained, the newest retained + // row names the live catalog root, so the retained-root walk re-seals the + // catalog incidentally and hides whether the rekey covers it on its own. + let options = OpenOptions::default().with_commit_history_retain(RetainPolicy::Disabled); + let db = Db::open_internal_with_options(vfs.clone(), KEK0, PAGE, REALM, options.clone()) + .await + .unwrap(); + + // Enough counter rows to push the catalog tree past a single page. + { + let mut tx = db.begin_write().await.unwrap(); + for index in 0..400_u32 { + let mut counter = tx.counter(&counter_name(index)).unwrap(); + counter.set(u64::from(index) + 1).await.unwrap(); + drop(counter); + } + tx.commit().await.unwrap(); + } + let catalog_pages_before = db.stats().await.unwrap().main_db_next_page_id; + assert!( + catalog_pages_before > 8, + "test setup must build a multi-page catalog, got {catalog_pages_before} pages" + ); + + db.rekey_db(KEK0, 1).await.unwrap(); + drop(db); + + let reopened = Db::open_existing_with_options(vfs, KEK0, PAGE, REALM, options) + .await + .unwrap(); + let report = pagedb::recovery::deep_walk::run_deep_walk(&reopened) + .await + .unwrap(); + assert!( + report.is_clean(), + "multi-page catalog rekey report: {report:?}" + ); + + // Every counter must still decode under the target epoch. + let mut tx = reopened.begin_write().await.unwrap(); + for index in [0_u32, 199, 399] { + let counter = tx.counter(&counter_name(index)).unwrap(); + assert_eq!(counter.get().await.unwrap(), u64::from(index) + 1); + drop(counter); + } + tx.commit().await.unwrap(); +} + +/// Rekey must not abandon the pages its own catalog transitions supersede. +/// +/// Each durable stage rewrites the catalog copy-on-write. Without routing the +/// superseded pages onto the free list they are unreachable from every root and +/// absent from the free list — a permanent leak that grows with every rekey, +/// and one the source epoch's retirement makes unauthenticatable. +#[tokio::test(flavor = "current_thread")] +async fn rekey_leaves_no_leaked_pages() { + let vfs = MemVfs::new(); + let db = Db::open_internal(vfs.clone(), KEK0, PAGE, REALM) + .await + .unwrap(); + { + let mut tx = db.begin_write().await.unwrap(); + for index in 0..200_u32 { + tx.put(format!("leak-{index:05}").as_bytes(), &[0x5C; 96]) + .await + .unwrap(); + } + tx.commit().await.unwrap(); + } + + for epoch in 1..=3_u64 { + db.rekey_db(KEK0, epoch).await.unwrap(); + } + drop(db); + + let reopened = Db::open_existing(vfs, KEK0, PAGE, REALM).await.unwrap(); + let report = pagedb::recovery::deep_walk::run_deep_walk(&reopened) + .await + .unwrap(); + assert!(report.is_clean(), "repeated rekey report: {report:?}"); + assert!( + report.orphan_page_ids.is_empty(), + "rekey leaked {} pages: {:?}", + report.orphan_page_ids.len(), + report.orphan_page_ids + ); +}