From 2da2fc29d9c5b06237527d610f358abdc14a2860 Mon Sep 17 00:00:00 2001 From: Andrew Date: Tue, 11 Aug 2026 06:03:16 +0000 Subject: [PATCH 1/4] fix(rekey): a reader leaving mid-deferral must not strand the retirement --- Cargo.lock | 2 ++ src/txn/db/core.rs | 45 +++++++++++++++++++++++++-- src/txn/db/open/create.rs | 2 ++ src/txn/db/open/existing.rs | 2 ++ src/txn/db/rekey/main.rs | 61 +++++++++++++++++++++++++++++++++++++ 5 files changed, 109 insertions(+), 3 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 66c1f34..2ef864d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -521,6 +521,8 @@ checksum = "da7c62ceae207dd37ea5b845da6a0696c799f85e97da1ab5b7910be3c1c80223" [[package]] name = "faultbox" version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ecb75815becab6e976d1af62d50f11d88f2600ab80d36a7191cdcba057f7d8e3" dependencies = [ "serde", "serde_json", diff --git a/src/txn/db/core.rs b/src/txn/db/core.rs index 6aad0b3..adc6b36 100644 --- a/src/txn/db/core.rs +++ b/src/txn/db/core.rs @@ -244,8 +244,17 @@ pub struct Db { pub(crate) visibility_test_hook: parking_lot::Mutex>>, #[cfg(test)] pub(crate) rekey_test_fault: parking_lot::Mutex>, + #[cfg(test)] + pub(crate) retirement_defer_hook: parking_lot::Mutex>>, } +/// Runs inside `retire_rekey_source_when_safe`'s critical section, once it has +/// decided to defer and before the obligation is recorded. That instant is the +/// only place a test can stand to prove the two are indivisible, because the +/// failure it guards against needs the last reader to leave precisely there. +#[cfg(test)] +pub(crate) type RetirementDeferHook = Arc) + Send + Sync>; + /// Reader-visible state, refreshed by the writer at commit time. #[derive(Debug, Clone, Copy)] #[allow(clippy::struct_field_names)] @@ -329,16 +338,28 @@ impl Db { /// Retire an obsolete source epoch immediately when no reader can still /// resolve a pre-cutover snapshot; otherwise defer retirement until the /// tracked reader set drains. + /// + /// Deciding to defer and recording what was deferred happen under one + /// lock, held across both. They are one step, not two: the reader set is + /// only a reason to defer for as long as the obligation is not yet visible + /// to the drain, and the last reader can leave at any instant. Split them + /// and that reader's drain reads an obligation list that is still empty, + /// then this pushes into a list nothing will visit again — the rotation + /// reports success while the superseded key stays leasable for the life of + /// the handle. `drain_pending_key_retirements` takes the same two locks in + /// the same order. pub(crate) fn retire_rekey_source_when_safe( &self, epoch: u64, cipher_id: CipherId, ) -> Result<()> { + let mut retirements = self.pending_key_retirements.lock(); if self.tracked_readers.lock().is_empty() { return self.pager.retire_mk_epoch(epoch, cipher_id); } + #[cfg(test)] + self.await_retirement_defer_hook(); let pending = PendingKeyRetirement { epoch, cipher_id }; - let mut retirements = self.pending_key_retirements.lock(); if !retirements.contains(&pending) { retirements.push(pending); } @@ -346,14 +367,19 @@ impl Db { } /// Drain deferred source-epoch retirements once no tracked reader remains. + /// + /// An entry leaves the list only once its retirement has succeeded, so a + /// failure keeps every obligation — including the one that failed — queued + /// for the next drain instead of discarding the untried remainder. pub(crate) fn drain_pending_key_retirements(&self) -> Result<()> { + let mut retirements = self.pending_key_retirements.lock(); if !self.tracked_readers.lock().is_empty() { return Ok(()); } - let pending = std::mem::take(&mut *self.pending_key_retirements.lock()); - for retirement in pending { + while let Some(retirement) = retirements.first().copied() { self.pager .retire_mk_epoch(retirement.epoch, retirement.cipher_id)?; + retirements.remove(0); } Ok(()) } @@ -425,6 +451,19 @@ impl Db { *self.visibility_test_hook.lock() = Some(hook); } + #[cfg(test)] + pub(crate) fn install_retirement_defer_hook(&self, hook: RetirementDeferHook) { + *self.retirement_defer_hook.lock() = Some(hook); + } + + #[cfg(test)] + fn await_retirement_defer_hook(&self) { + let hook = self.retirement_defer_hook.lock().clone(); + if let Some(hook) = hook { + hook(self); + } + } + /// Detach the hook so the rest of a test runs on unrehearsed paths. Every /// pause point is a rendezvous that only completes when its counterpart /// notifies, so a hook left installed past the interleaving it was staging diff --git a/src/txn/db/open/create.rs b/src/txn/db/open/create.rs index 50cc421..eb93727 100644 --- a/src/txn/db/open/create.rs +++ b/src/txn/db/open/create.rs @@ -134,6 +134,8 @@ impl Db { #[cfg(test)] visibility_test_hook: parking_lot::Mutex::new(None), #[cfg(test)] + retirement_defer_hook: parking_lot::Mutex::new(None), + #[cfg(test)] rekey_test_fault: parking_lot::Mutex::new(None), }) } diff --git a/src/txn/db/open/existing.rs b/src/txn/db/open/existing.rs index 6ca0cb6..ad37d2b 100644 --- a/src/txn/db/open/existing.rs +++ b/src/txn/db/open/existing.rs @@ -325,6 +325,8 @@ impl Db { #[cfg(test)] visibility_test_hook: parking_lot::Mutex::new(None), #[cfg(test)] + retirement_defer_hook: parking_lot::Mutex::new(None), + #[cfg(test)] rekey_test_fault: parking_lot::Mutex::new(None), }; diff --git a/src/txn/db/rekey/main.rs b/src/txn/db/rekey/main.rs index 785fdc3..5068cfd 100644 --- a/src/txn/db/rekey/main.rs +++ b/src/txn/db/rekey/main.rs @@ -944,6 +944,67 @@ mod tests { ); } + /// The last reader may leave *during* the deferral, and the retirement must + /// still complete. + /// + /// Deferring is two facts that have to agree: a reader exists, and an + /// obligation is recorded where the drain will find it. Nothing pins the + /// first fact in place — the reader's `drop` is a plain unregister that can + /// land at any instant, including between the two. When it did, the drain + /// read an empty obligation list and the record arrived after the only + /// thing that ever reads it had gone: `rekey_db` returned `Ok(())` while the + /// superseded master key stayed leasable for the life of the handle, with no + /// error, no log line, and nothing to query. This drives the reader out at + /// exactly that instant and requires the epoch to retire anyway. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn a_reader_leaving_mid_deferral_still_retires_the_source_epoch() { + let db = Db::open_internal(MemVfs::new(), SOURCE_KEK, PAGE, REALM) + .await + .unwrap(); + // Move the active epoch off 0 so retiring it is a real retirement + // rather than the active-epoch refusal. + let target_mk = derive_mk(&TARGET_KEK, &db.kek_salt, 1).unwrap(); + db.pager.set_active_mk_epoch(target_mk, 1); + + db.install_retirement_defer_hook(std::sync::Arc::new(|db: &Db| { + for _ in 0..5_000 { + if db.tracked_readers.lock().is_empty() { + return; + } + std::thread::sleep(std::time::Duration::from_millis(1)); + } + panic!("the reader never left; the window under test was not reached"); + })); + + let read = db.begin_read().await.unwrap(); + assert!( + db.pager.mk_for(0, db.cipher_id).is_ok(), + "the source epoch must still be leasable while a reader holds it" + ); + + std::thread::scope(|scope| { + scope.spawn(|| drop(read)); + db.retire_rekey_source_when_safe(0, db.cipher_id) + }) + .unwrap(); + + assert!( + db.tracked_readers.lock().is_empty(), + "the reader must be gone before the retirement is judged" + ); + assert!( + db.pending_key_retirements.lock().is_empty(), + "an obligation no reader is waiting on must not survive the deferral" + ); + assert!( + matches!( + db.pager.mk_for(0, db.cipher_id), + Err(PagedbError::MissingPersistedKey { mk_epoch: 0, .. }) + ), + "a reader leaving mid-deferral must not leave the source epoch leasable" + ); + } + #[tokio::test(flavor = "current_thread")] async fn direct_segment_reader_keeps_source_lease_after_retirement() { let db = Db::open_internal(MemVfs::new(), SOURCE_KEK, PAGE, REALM) From d2e930207e340e7e58cbefc660b22b612c2bbddb Mon Sep 17 00:00:00 2001 From: Andrew Date: Tue, 11 Aug 2026 06:57:22 +0000 Subject: [PATCH 2/4] docs(changelog): record the rekey retirement fix under Unreleased --- CHANGELOG.md | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1b82abe..ce98be5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,12 @@ All notable changes to this project are documented here. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [Unreleased] + +### Fixed + +- **A rekey's source-epoch retirement is no longer lost when the last reader leaves mid-deferral.** `retire_rekey_source_when_safe` decided to defer under one lock and recorded the obligation under another. A reader dropping between the two ran its drain against a still-empty obligation list, so the record landed where nothing would visit it again and the superseded master key stayed leasable for the life of the handle — while `rekey_db` returned `Ok(())` with no error and no completion state to query. The deferral is now one step under one lock, and a drain that fails keeps every obligation queued instead of discarding the untried remainder. + ## [0.1.0] - 2026-07-28 The first release. Pre-releases were published as `0.1.0-beta.N`; the entries below describe `0.1.0` as a whole rather than deltas against a shipped version, since none exists yet. From 460e7dd75c48b91cc0d4429b1ef54aec035a19a9 Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Wed, 12 Aug 2026 09:42:28 +0800 Subject: [PATCH 3/4] fix(rekey): retire independent obligations even if one refuses retire_recorded previously stopped draining at the first failed retirement, leaving every obligation behind it stuck queued as well. Retirements are independent epochs, so it now retains only the entries that actually fail and keeps retiring the rest, reporting the first error once the pass is done. --- src/txn/db/core.rs | 61 +++++++++++++++++------ src/txn/db/rekey/main.rs | 104 ++++++++++++++++++++++++++++++++++++++- 2 files changed, 149 insertions(+), 16 deletions(-) diff --git a/src/txn/db/core.rs b/src/txn/db/core.rs index adc6b36..87a10ac 100644 --- a/src/txn/db/core.rs +++ b/src/txn/db/core.rs @@ -348,40 +348,65 @@ impl Db { /// reports success while the superseded key stays leasable for the life of /// the handle. `drain_pending_key_retirements` takes the same two locks in /// the same order. + /// + /// The immediate path records the obligation too, then retires the whole + /// list: it is the only place that sees an empty reader set together with + /// the backlog, and the drain that would otherwise retry a failed entry + /// only runs when the last reader leaves — which by then has happened. pub(crate) fn retire_rekey_source_when_safe( &self, epoch: u64, cipher_id: CipherId, ) -> Result<()> { let mut retirements = self.pending_key_retirements.lock(); - if self.tracked_readers.lock().is_empty() { - return self.pager.retire_mk_epoch(epoch, cipher_id); + let readers_present = !self.tracked_readers.lock().is_empty(); + if readers_present { + #[cfg(test)] + self.await_retirement_defer_hook(); } - #[cfg(test)] - self.await_retirement_defer_hook(); let pending = PendingKeyRetirement { epoch, cipher_id }; if !retirements.contains(&pending) { retirements.push(pending); } - Ok(()) + if readers_present { + return Ok(()); + } + self.retire_recorded(&mut retirements) } /// Drain deferred source-epoch retirements once no tracked reader remains. - /// - /// An entry leaves the list only once its retirement has succeeded, so a - /// failure keeps every obligation — including the one that failed — queued - /// for the next drain instead of discarding the untried remainder. pub(crate) fn drain_pending_key_retirements(&self) -> Result<()> { let mut retirements = self.pending_key_retirements.lock(); if !self.tracked_readers.lock().is_empty() { return Ok(()); } - while let Some(retirement) = retirements.first().copied() { - self.pager - .retire_mk_epoch(retirement.epoch, retirement.cipher_id)?; - retirements.remove(0); + self.retire_recorded(&mut retirements) + } + + /// Retire every recorded obligation, keeping the ones that fail. + /// + /// An entry leaves the list only once its retirement has succeeded, so a + /// failure keeps that obligation queued for the next attempt — and keeps + /// only that one: the rest are independent epochs, and stopping at the + /// first error would let one unretirable epoch stand in front of them + /// forever. The first error is reported once the whole pass is done. + fn retire_recorded(&self, retirements: &mut Vec) -> Result<()> { + let mut first_error = None; + retirements.retain(|pending| { + match self.pager.retire_mk_epoch(pending.epoch, pending.cipher_id) { + Ok(()) => false, + Err(error) => { + if first_error.is_none() { + first_error = Some(error); + } + true + } + } + }); + match first_error { + Some(error) => Err(error), + None => Ok(()), } - Ok(()) } pub(crate) fn ensure_usable(&self) -> Result<()> { @@ -456,6 +481,14 @@ impl Db { *self.retirement_defer_hook.lock() = Some(hook); } + /// Detach the deferral hook, for the same reason as + /// [`Self::clear_visibility_test_hook`]: it is a rendezvous, and a later + /// deferral would park on a counterpart that is no longer coming. + #[cfg(test)] + pub(crate) fn clear_retirement_defer_hook(&self) { + *self.retirement_defer_hook.lock() = None; + } + #[cfg(test)] fn await_retirement_defer_hook(&self) { let hook = self.retirement_defer_hook.lock().clone(); diff --git a/src/txn/db/rekey/main.rs b/src/txn/db/rekey/main.rs index 5068cfd..f8c9fb5 100644 --- a/src/txn/db/rekey/main.rs +++ b/src/txn/db/rekey/main.rs @@ -966,7 +966,19 @@ mod tests { let target_mk = derive_mk(&TARGET_KEK, &db.kek_salt, 1).unwrap(); db.pager.set_active_mk_epoch(target_mk, 1); - db.install_retirement_defer_hook(std::sync::Arc::new(|db: &Db| { + // Left to race, the reader can be gone before the critical section is + // even entered, the deferral is never taken, and every assertion below + // passes against the unfixed code. So the hook drives the reader out + // itself, from the one instant the bug needed. + let departure = std::sync::Arc::new(std::sync::Barrier::new(2)); + let deferrals = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let hook_departure = departure.clone(); + let hook_deferrals = deferrals.clone(); + db.install_retirement_defer_hook(std::sync::Arc::new(move |db: &Db| { + hook_deferrals.fetch_add(1, Ordering::Relaxed); + hook_departure.wait(); + // `unregister_read` releases `tracked_readers` before it touches + // the retirement list held here, so this cannot deadlock. for _ in 0..5_000 { if db.tracked_readers.lock().is_empty() { return; @@ -983,11 +995,21 @@ mod tests { ); std::thread::scope(|scope| { - scope.spawn(|| drop(read)); + scope.spawn(|| { + departure.wait(); + drop(read); + }); db.retire_rekey_source_when_safe(0, db.cipher_id) }) .unwrap(); + db.clear_retirement_defer_hook(); + assert_eq!( + deferrals.load(Ordering::Relaxed), + 1, + "the retirement must have been deferred; otherwise the reader left \ + before the window and the race was never staged" + ); assert!( db.tracked_readers.lock().is_empty(), "the reader must be gone before the retirement is judged" @@ -1005,6 +1027,84 @@ mod tests { ); } + /// Queue two obligations, then make the first one refuse to retire. + /// Retirements are independent epochs, so the refusal must be kept and + /// stepped over, not treated as a wall. + #[tokio::test(flavor = "current_thread")] + async fn a_refused_retirement_does_not_hold_up_the_obligations_behind_it() { + let db = Db::open_internal(MemVfs::new(), SOURCE_KEK, PAGE, REALM) + .await + .unwrap(); + let target_mk = derive_mk(&TARGET_KEK, &db.kek_salt, 1).unwrap(); + db.pager.set_active_mk_epoch(target_mk, 1); + let spare_mk = derive_mk(&TARGET_KEK, &db.kek_salt, 2).unwrap(); + db.pager.install_mk_epoch(spare_mk, 2, db.cipher_id); + + let read = db.begin_read().await.unwrap(); + db.retire_rekey_source_when_safe(0, db.cipher_id).unwrap(); + db.retire_rekey_source_when_safe(2, db.cipher_id).unwrap(); + + // Epoch 0 is the active one by the time the drain runs, which is the + // one condition `retire_mk_epoch` refuses. + let source_mk = derive_mk(&SOURCE_KEK, &db.kek_salt, 0).unwrap(); + db.pager.set_active_mk_epoch(source_mk, 0); + drop(read); + + assert!( + db.pager.mk_for(2, db.cipher_id).is_err(), + "an epoch queued behind a refusal must still be retired" + ); + let queued: Vec = db + .pending_key_retirements + .lock() + .iter() + .map(|pending| pending.epoch) + .collect(); + assert_eq!( + queued, + vec![0], + "the refused epoch — and only it — must stay queued for the next attempt" + ); + } + + /// A drain runs only when the last reader leaves, so an obligation it + /// failed to clear has no second drain coming while the reader set stays + /// empty. The next rotation is the only thing left that can retry it. + #[tokio::test(flavor = "current_thread")] + async fn a_rotation_with_no_readers_retires_the_backlog_it_inherits() { + let db = Db::open_internal(MemVfs::new(), SOURCE_KEK, PAGE, REALM) + .await + .unwrap(); + let target_mk = derive_mk(&TARGET_KEK, &db.kek_salt, 1).unwrap(); + db.pager.set_active_mk_epoch(target_mk, 1); + + let read = db.begin_read().await.unwrap(); + db.retire_rekey_source_when_safe(0, db.cipher_id).unwrap(); + let source_mk = derive_mk(&SOURCE_KEK, &db.kek_salt, 0).unwrap(); + db.pager.set_active_mk_epoch(source_mk, 0); + drop(read); + assert_eq!( + db.pending_key_retirements.lock().len(), + 1, + "the refused epoch must have been kept" + ); + + let target_mk = derive_mk(&TARGET_KEK, &db.kek_salt, 1).unwrap(); + db.pager.set_active_mk_epoch(target_mk, 1); + let spare_mk = derive_mk(&TARGET_KEK, &db.kek_salt, 2).unwrap(); + db.pager.install_mk_epoch(spare_mk, 2, db.cipher_id); + db.retire_rekey_source_when_safe(2, db.cipher_id).unwrap(); + + assert!( + db.pending_key_retirements.lock().is_empty(), + "a rotation that finds no reader must clear the backlog, not just its own epoch" + ); + assert!( + db.pager.mk_for(0, db.cipher_id).is_err(), + "the inherited obligation must actually have been retired" + ); + } + #[tokio::test(flavor = "current_thread")] async fn direct_segment_reader_keeps_source_lease_after_retirement() { let db = Db::open_internal(MemVfs::new(), SOURCE_KEK, PAGE, REALM) From 88efa8f1695a345595983f092d355f404601d68f Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Wed, 12 Aug 2026 09:42:31 +0800 Subject: [PATCH 4/4] docs(changelog): fold the retirement fix into the rekey release bullet --- CHANGELOG.md | 8 +------- 1 file changed, 1 insertion(+), 7 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index ce98be5..4b5d00c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,12 +2,6 @@ All notable changes to this project are documented here. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). -## [Unreleased] - -### Fixed - -- **A rekey's source-epoch retirement is no longer lost when the last reader leaves mid-deferral.** `retire_rekey_source_when_safe` decided to defer under one lock and recorded the obligation under another. A reader dropping between the two ran its drain against a still-empty obligation list, so the record landed where nothing would visit it again and the superseded master key stayed leasable for the life of the handle — while `rekey_db` returned `Ok(())` with no error and no completion state to query. The deferral is now one step under one lock, and a drain that fails keeps every obligation queued instead of discarding the untried remainder. - ## [0.1.0] - 2026-07-28 The first release. Pre-releases were published as `0.1.0-beta.N`; the entries below describe `0.1.0` as a whole rather than deltas against a shipped version, since none exists yet. @@ -22,7 +16,7 @@ The first release. Pre-releases were published as `0.1.0-beta.N`; the entries be - **Cross-platform VFS** — Linux (`io_uring`), Windows (IOCP), macOS/iOS (Grand Central Dispatch), Android, WASM/OPFS and WASI backends, plus a tokio thread-pool fallback and an in-memory backend, with format-bit identity across targets. On Linux the backend is chosen at run time: a kernel that refuses an `io_uring` ring falls back to the thread pool with a warning instead of failing the open. All native backends share one advisory-lock implementation, so processes on different backends still exclude each other on one store. - **Snapshots** — `snapshot_to`, `restore_from`, and incremental apply, each authenticated against the state its manifest describes. Destinations must be empty; malformed or incomplete artifacts fail closed. - **Recovery** — open-flow GC, apply-journal replay, deep-walk `fsck`, and the `pagedb-fsck` binary. -- **Online rekey** — rekey under a new key with mixed-cipher and mixed-epoch page coexistence; no full-file migration. +- **Online rekey** — rekey under a new key with mixed-cipher and mixed-epoch page coexistence; no full-file migration. A rotation whose source epoch is still pinned by a reader defers retiring it and completes the retirement once the reader set drains — including when the last reader leaves during the deferral itself, so the superseded master key never stays leasable behind an `Ok(())`. A retirement that cannot be taken yet stays queued for the next attempt without holding up the others. - **Handle modes** — `Standalone`, `Follower`, `ReadOnly`, and `Observer`. - **Open refusals name the parameter, not the store** — `KeyMismatch`, `PageSizeMismatch`, and `RealmMismatch`, each decided before anything is read or written, and none reported as corruption. - **Failures report themselves** — an unreadable free-list chain, main file, or segment catalog fails `stats()` instead of reporting zero; compaction never skips a catalog entry whose file it cannot open; segment open distinguishes a missing file from a permission or backend error; and only genuine contention is reported as contention. Persisted named-counter rows are validated at open, and commit-history keys are rejected unless exactly eight bytes.