From 83e7b3ef45c358ebd429631ea6d72f82128199ac Mon Sep 17 00:00:00 2001 From: tison Date: Mon, 14 Sep 2026 19:43:28 +0800 Subject: [PATCH 1/2] refactor: use auto-reset events for drain notifications --- Cargo.lock | 5 +- Cargo.toml | 8 ++- cache2/src/region/runtime/mod.rs | 85 +++++++++++++++++++++++--------- deny.toml | 2 + 4 files changed, 72 insertions(+), 28 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 393454c..bb13bca 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -60,9 +60,8 @@ checksum = "330a5ed07fa54e4702c9d6c4174f74427fc0ef6e214bbd677ae50a5099946470" [[package]] name = "asyncband" -version = "0.7.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2e52766975a4f080528a898235c51e82e65df9db713419067b5368040eeb5659" +version = "0.7.2" +source = "git+https://github.com/tisonkun/asyncband.git?rev=085390437eba05a76efdcfd7e4388b1ea0023dbc#085390437eba05a76efdcfd7e4388b1ea0023dbc" [[package]] name = "benchmarks" diff --git a/Cargo.toml b/Cargo.toml index c5ec144..b32a388 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -27,8 +27,14 @@ rust-version = "1.98" # Workspace crates cache2 = { path = "cache2", version = "0.4.0" } +# Unreleased AutoResetEvent implementation; use its crates.io release before merging. +asyncband = { version = "0.7.2", git = "https://github.com/tisonkun/asyncband.git", rev = "085390437eba05a76efdcfd7e4388b1ea0023dbc", features = [ + "barrier", + "event", + "semaphore", +] } + # Crates.io dependencies -asyncband = { version = "0.7.1", features = ["barrier", "semaphore", "watch"] } cargo_metadata = { version = "0.23.1" } clap = { version = "4.6.5", features = ["derive"] } crc-fast = { version = "1.10.0", default-features = false, features = ["std"] } diff --git a/cache2/src/region/runtime/mod.rs b/cache2/src/region/runtime/mod.rs index 182155c..9206055 100644 --- a/cache2/src/region/runtime/mod.rs +++ b/cache2/src/region/runtime/mod.rs @@ -35,9 +35,9 @@ use std::thread::JoinHandle; use std::time::Duration; use std::time::Instant; +use asyncband::event::AutoResetEvent; use asyncband::semaphore::OwnedSemaphorePermit; use asyncband::semaphore::Semaphore; -use asyncband::watch; use self::metrics::RuntimeMetrics; use crate::IoEngineConfig; @@ -160,17 +160,17 @@ struct MutationGate { state: AtomicUsize, quiescent: Mutex<()>, quiescent_changed: Condvar, - async_changed: watch::Sender<()>, + // MutationDrainGuard admits only one async drain; synchronous close waits use the condvar. + async_changed: AutoResetEvent, } impl MutationGate { fn new() -> Self { - let (async_changed, _) = watch::channel(()); Self { state: AtomicUsize::new(0), quiescent: Mutex::new(()), quiescent_changed: Condvar::new(), - async_changed, + async_changed: AutoResetEvent::new(), } } @@ -231,14 +231,10 @@ impl MutationGate { } async fn wait_quiescent_async(&self) { - // Subscribe before inspecting the predicate so a transition racing the check advances the - // receiver version and cannot be missed. - let mut changed = self.async_changed.subscribe(); + // A completion between the check and wait leaves a stored signal. A signal left by a + // cancelled drain only causes another predicate check in the next drain. while self.active_mutations() != 0 { - changed - .changed() - .await - .expect("the mutation gate retains its watch sender"); + self.async_changed.wait().await; } } @@ -252,7 +248,7 @@ impl MutationGate { .unwrap_or_else(|poisoned| poisoned.into_inner()); self.quiescent_changed.notify_all(); drop(quiescent); - self.async_changed.send_replace(()); + self.async_changed.set(); } } } @@ -615,16 +611,17 @@ struct ShardControlState { struct ShardControl { state: Mutex, changed: Condvar, - async_changed: watch::Sender<()>, + // drain_async holds MutationDrainGuard across all shard waits, so only one async observer + // competes for this signal. Synchronous drain/close observers use changed instead. + async_changed: AutoResetEvent, } impl ShardControl { fn new() -> Self { - let (async_changed, _) = watch::channel(()); Self { state: Mutex::new(ShardControlState::default()), changed: Condvar::new(), - async_changed, + async_changed: AutoResetEvent::new(), } } @@ -697,9 +694,8 @@ impl ShardControl { } async fn wait_for_drain_async(&self, generation: u64) -> io::Result<()> { - // Subscribe before inspecting the predicate so a completion racing the check advances the - // receiver version and cannot be missed. - let mut changed = self.async_changed.subscribe(); + // Completion is recorded in state; the event only prompts a recheck. Signals from older + // generations may remain stored but cannot satisfy this generation's predicate. loop { { let state = self.lock()?; @@ -710,10 +706,7 @@ impl ShardControl { return Ok(()); } } - changed - .changed() - .await - .expect("the shard control retains its watch sender"); + self.async_changed.wait().await; } } @@ -727,7 +720,7 @@ impl ShardControl { .get_or_insert_with(|| ShardFailure::from_error(error)); drop(state); self.changed.notify_all(); - self.async_changed.send_replace(()); + self.async_changed.set(); } fn lock(&self) -> io::Result> { @@ -1954,7 +1947,7 @@ fn complete_shard_drain(control: &ShardControl, generation: u64) -> io::Result<( state.drain_completed = state.drain_completed.max(generation); control.changed.notify_all(); drop(state); - control.async_changed.send_replace(()); + control.async_changed.set(); Ok(()) } @@ -2357,6 +2350,50 @@ mod tests { drop(mutation); } + #[test] + fn async_drain_remains_exclusive_and_can_retry_after_a_selected_wait_is_cancelled() { + let gate = MutationGate::new(); + let mutation = gate.try_enter().unwrap(); + let drain = gate.begin_drain().unwrap(); + let mut wait = Box::pin(drain.wait_async()); + let mut context = Context::from_waker(Waker::noop()); + assert!(wait.as_mut().poll(&mut context).is_pending()); + assert!( + matches!(gate.begin_drain(), Err(error) if error.kind() == io::ErrorKind::WouldBlock) + ); + + drop(mutation); + drop(wait); // Cancel after the completion signal was assigned, before observing it. + drop(drain); + + let mutation = gate.try_enter().unwrap(); + let drain = gate.begin_drain().unwrap(); + let mut retry = Box::pin(drain.wait_async()); + assert!(retry.as_mut().poll(&mut context).is_pending()); + drop(mutation); + assert!(retry.as_mut().poll(&mut context).is_ready()); + } + + #[test] + fn old_shard_signals_do_not_complete_new_drain_generations() { + let control = ShardControl::new(); + let first = control.request_drain(false).unwrap(); + complete_shard_drain(&control, first).unwrap(); + + let second = control.request_drain(false).unwrap(); + let mut cancelled = Box::pin(control.wait_for_drain_async(second)); + let mut context = Context::from_waker(Waker::noop()); + assert!(cancelled.as_mut().poll(&mut context).is_pending()); + complete_shard_drain(&control, second).unwrap(); + drop(cancelled); + + let third = control.request_drain(false).unwrap(); + let mut next = Box::pin(control.wait_for_drain_async(third)); + assert!(next.as_mut().poll(&mut context).is_pending()); + complete_shard_drain(&control, third).unwrap(); + assert!(next.as_mut().poll(&mut context).is_ready()); + } + #[test] fn shard_failure_wakes_async_waiters_after_releasing_state_lock() { let control = Arc::new(ShardControl::new()); diff --git a/deny.toml b/deny.toml index e3fe86f..c551b60 100644 --- a/deny.toml +++ b/deny.toml @@ -25,3 +25,5 @@ wildcards = "deny" [sources] unknown-git = "deny" unknown-registry = "deny" +# Temporary source for the revision pinned in Cargo.toml; remove with the crates.io migration. +allow-git = ["https://github.com/tisonkun/asyncband.git"] From ad1e2f896849fc8053334a24e7199e7a2b22f55d Mon Sep 17 00:00:00 2001 From: tison Date: Mon, 14 Sep 2026 20:40:42 +0800 Subject: [PATCH 2/2] ci: verify packages against the pinned asyncband revision --- .cargo/config.toml | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/.cargo/config.toml b/.cargo/config.toml index 541e3bb..9493b3f 100644 --- a/.cargo/config.toml +++ b/.cargo/config.toml @@ -17,3 +17,8 @@ x = "run --package x --" [env] CARGO_WORKSPACE_DIR = { value = "", relative = true } + +# Cargo normalizes Git dependencies to registry versions when verifying a package. +# Keep this temporary patch in sync with Cargo.toml until AutoResetEvent is released. +[patch.crates-io] +asyncband = { git = "https://github.com/tisonkun/asyncband.git", rev = "085390437eba05a76efdcfd7e4388b1ea0023dbc" }