diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index e101e139c..1482839ba 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -6142,9 +6142,8 @@ pub(in crate::set_disk) mod rename_fanout_barrier_phase { /// disk_index, phase)`. The matching spawned fan-out task blocks at its /// [`checkpoint`] until the test releases it. The test awaits the pause via /// [`BarrierHandle::wait_until_paused`] (a deterministic `Notify` handshake — -/// no sleeps) and resumes it via [`BarrierHandle::release`]. At most one -/// barrier is armed per object at a time, matching the single-scope style of -/// `disk_call_counters`. +/// no sleeps) and resumes it via [`BarrierHandle::release`]. Different disks +/// and phases may be armed together for the same object. /// /// 2. **Background-task introspection.** A test [`observe_tasks`] for `object`; /// each instrumented fan-out task then holds a [`TaskGuard`] for its whole @@ -6170,17 +6169,38 @@ pub(crate) mod rename_fanout_barrier { use std::sync::{Arc, Mutex, OnceLock}; use tokio::sync::Notify; - pub use super::rename_fanout_barrier_phase::{ - CLEANUP as PHASE_CLEANUP, READ_VERSION as PHASE_READ_VERSION, RENAME as PHASE_RENAME, ROLLBACK as PHASE_ROLLBACK, - }; + #[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] + pub enum Phase { + Rename, + Cleanup, + Rollback, + ReadVersion, + NonInlineHedgeTimer, + } + + impl From<&'static str> for Phase { + fn from(label: &'static str) -> Self { + match label { + super::rename_fanout_barrier_phase::RENAME => Self::Rename, + super::rename_fanout_barrier_phase::CLEANUP => Self::Cleanup, + super::rename_fanout_barrier_phase::ROLLBACK => Self::Rollback, + super::rename_fanout_barrier_phase::READ_VERSION => Self::ReadVersion, + "non_inline_hedge_timer" => Self::NonInlineHedgeTimer, + _ => panic!("unknown commit fault phase: {label}"), + } + } + } + + pub const PHASE_RENAME: Phase = Phase::Rename; + pub const PHASE_CLEANUP: Phase = Phase::Cleanup; + pub const PHASE_ROLLBACK: Phase = Phase::Rollback; + pub const PHASE_READ_VERSION: Phase = Phase::ReadVersion; /// Object-scoped hedge timer checkpoint; slot zero identifies the timer, not a disk. - pub const PHASE_NON_INLINE_HEDGE_TIMER: &str = "non_inline_hedge_timer"; + pub const PHASE_NON_INLINE_HEDGE_TIMER: Phase = Phase::NonInlineHedgeTimer; /// One armed barrier: the fan-out task matching `(disk_index, phase)` pauses. struct Armed { - disk_index: usize, - phase: &'static str, /// Signalled (task -> test) when the target task reaches the checkpoint. arrived: Arc, /// Signalled (test -> task) to release the paused task. @@ -6191,8 +6211,8 @@ pub(crate) mod rename_fanout_barrier { #[derive(Default)] struct Registry { - /// object -> armed barrier (at most one per object). - armed: HashMap, + /// Independently armed checkpoints for an object's fault schedule. + armed: HashMap<(String, usize, Phase), Armed>, /// object -> live in-flight fan-out task count, only for observed objects. observed: HashMap>, } @@ -6211,7 +6231,7 @@ pub(crate) mod rename_fanout_barrier { /// leave a spawned fan-out task wedged. #[must_use] pub struct BarrierHandle { - object: String, + key: (String, usize, Phase), arrived: Arc, release: Arc, paused: Arc, @@ -6219,22 +6239,26 @@ pub(crate) mod rename_fanout_barrier { /// Arm a barrier: the fan-out task for `object` at `(disk_index, phase)` will /// pause at its checkpoint until the returned handle is released or dropped. - pub fn arm(object: &str, disk_index: usize, phase: &'static str) -> BarrierHandle { + pub fn arm(object: &str, disk_index: usize, phase: Phase) -> BarrierHandle { let arrived = Arc::new(Notify::new()); let release = Arc::new(Notify::new()); let paused = Arc::new(AtomicBool::new(false)); - lock().armed.insert( - object.to_string(), + let key = (object.to_string(), disk_index, phase); + let mut reg = lock(); + if reg.armed.contains_key(&key) { + drop(reg); + panic!("commit fault checkpoint already armed: {key:?}"); + } + reg.armed.insert( + key.clone(), Armed { - disk_index, - phase, arrived: arrived.clone(), release: release.clone(), paused: paused.clone(), }, ); BarrierHandle { - object: object.to_string(), + key, arrived, release, paused, @@ -6252,6 +6276,13 @@ pub(crate) mod rename_fanout_barrier { self.arrived.notified().await; } + /// Bound a fixture handshake and identify the unreached checkpoint. + pub async fn wait_until_paused_before(&self, deadline: tokio::time::Instant) { + tokio::time::timeout_at(deadline, self.wait_until_paused()) + .await + .unwrap_or_else(|_| panic!("commit fault checkpoint not reached before deadline: {:?}", self.key)); + } + /// Whether the target task is currently parked at the checkpoint. pub fn is_paused(&self) -> bool { self.paused.load(Ordering::SeqCst) @@ -6268,7 +6299,14 @@ pub(crate) mod rename_fanout_barrier { // Unblock any task still parked at the checkpoint before disarming, so // a dropped handle can never wedge a spawned fan-out task. self.release.notify_one(); - lock().armed.remove(&self.object); + let mut reg = lock(); + if reg + .armed + .get(&self.key) + .is_some_and(|armed| Arc::ptr_eq(&armed.release, &self.release)) + { + reg.armed.remove(&self.key); + } } } @@ -6276,15 +6314,13 @@ pub(crate) mod rename_fanout_barrier { /// unless a barrier is armed for exactly this `(object, disk_index, phase)`. /// The registry mutex is released before awaiting, so it is never held across /// the pause. - pub(in crate::set_disk) async fn checkpoint(object: &str, disk_index: usize, phase: &'static str) { + pub(in crate::set_disk) async fn checkpoint(object: &str, disk_index: usize, phase: impl Into) { + let key = (object.to_string(), disk_index, phase.into()); let hooks = { let reg = lock(); - match reg.armed.get(object) { - Some(a) if a.disk_index == disk_index && a.phase == phase => { - Some((a.arrived.clone(), a.release.clone(), a.paused.clone())) - } - _ => None, - } + reg.armed + .get(&key) + .map(|a| (a.arrived.clone(), a.release.clone(), a.paused.clone())) }; if let Some((arrived, release, paused)) = hooks { paused.store(true, Ordering::SeqCst); @@ -6367,6 +6403,45 @@ mod tests { use tempfile::TempDir; use tokio::io::AsyncReadExt; + #[tokio::test] + async fn commit_fault_schedule_keeps_multiple_disk_phases_independent() { + use rename_fanout_barrier::{PHASE_CLEANUP, PHASE_RENAME, arm, checkpoint}; + + let object = "commit-fault-schedule-multiple-phases"; + let rename = arm(object, 0, PHASE_RENAME); + let cleanup = arm(object, 1, PHASE_CLEANUP); + let rename_task = tokio::spawn(checkpoint(object, 0, PHASE_RENAME)); + let cleanup_task = tokio::spawn(checkpoint(object, 1, PHASE_CLEANUP)); + let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(2); + rename.wait_until_paused_before(deadline).await; + cleanup.wait_until_paused_before(deadline).await; + + rename.release(); + tokio::time::timeout_at(deadline, rename_task) + .await + .expect("rename checkpoint should release independently") + .expect("rename checkpoint task should finish"); + assert!(!cleanup_task.is_finished(), "releasing rename must not release cleanup"); + drop(rename); + assert!(cleanup.is_paused(), "dropping a different handle must preserve cleanup"); + drop(cleanup); + tokio::time::timeout_at(deadline, cleanup_task) + .await + .expect("dropping cleanup should release its task") + .expect("cleanup checkpoint task should finish"); + } + + #[test] + fn commit_fault_schedule_rejects_duplicate_checkpoint_without_disarming_owner() { + let object = "commit-fault-schedule-duplicate-checkpoint"; + let owner = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME); + let duplicate = std::panic::catch_unwind(|| rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME)); + assert!(duplicate.is_err(), "a second handle cannot replace a live checkpoint"); + drop(owner); + let next = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME); + drop(next); + } + #[test] fn orphan_dir_entries_must_be_single_relative_components() { for entry in ["part.1", "child/", "delete-data.00000000-0000-0000-0000-000000000001"] { @@ -9835,7 +9910,7 @@ mod tests { prepare_rename_source_dirs(&dirs, &disks, "source").await; let receipt = RenameRollbackReceipt::default(); let _fault = rollback_fault_injection::arm(&object, 0, fault); - let barrier = rename_fanout_barrier::arm(&object, 0, rename_fanout_barrier_phase::RENAME); + let barrier = rename_fanout_barrier::arm(&object, 0, rename_fanout_barrier::PHASE_RENAME); let mut rename = Box::pin(SetDisks::rename_data_owned_with_fence( &disks, (RUSTFS_META_TMP_BUCKET, "source"), @@ -10413,7 +10488,7 @@ mod tests { } let _rename_fault = rename_fault_injection::fail_rename_on(object, &[2, 3]); let _undo_fault = rollback_fault_injection::arm(object, 0, rollback_fault_injection::Fault::Io); - let barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier_phase::ROLLBACK); + let barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_ROLLBACK); let receipt = RenameRollbackReceipt::default(); let mut rename = Box::pin(SetDisks::rename_data_owned_with_fence( &disks, @@ -10505,7 +10580,7 @@ mod tests { let _rename_fault = rename_fault_injection::fail_rename_on(object, &[2, 3]); let _rollback_fault = rollback_fault_injection::arm(object, 0, rollback_fault_injection::Fault::RollbackCoordinatorPanic); - let barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier_phase::ROLLBACK); + let barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_ROLLBACK); let receipt = RenameRollbackReceipt::default(); let result = tokio::time::timeout( BARRIER_PAUSE_GUARD, diff --git a/docs/testing/ecstore-validation-suite-design.md b/docs/testing/ecstore-validation-suite-design.md index de717d1b9..ba8536db3 100644 --- a/docs/testing/ecstore-validation-suite-design.md +++ b/docs/testing/ecstore-validation-suite-design.md @@ -58,7 +58,7 @@ Fault injection is explicit and deterministic: local disk mocks for unit tests, `ObjectOptions::default()` uses `WriteCompletion::Quorum`: a namespace-lock-owning PUT may acknowledge write quorum while its rename tail retains the lock. A fixture that immediately inspects every disk or primes a metadata generation must set `write_completion: WriteCompletion::TailDrained` and keep normal locking. TailDrained waits for the existing rename fan-out; it does not require every disk to succeed or change fsync policy. Codec-only `no_lock` fixtures do not cover namespace locking. -The object tests reuse `rename_fanout_barrier::arm(object, disk_slot, phase)` and `observe_tasks(object)`. Wait for the barrier with a deadline, observe actual metadata quorum with `wait_for_paused_tail_metadata_quorum`, then release or cancel. The metadata check distinguishes a real quorum from disk tasks that have not started. Assert zero remaining rename tasks after the owned coordinator releases its lock; cancellation tests also wait for staging cleanup. +The object tests reuse `rename_fanout_barrier::arm(object, disk_slot, phase)` and `observe_tasks(object)`. The fixture uses typed rename, rollback, cleanup, read-version and hedge-timer phases; multiple disk/phase checkpoints can be armed for one object. A duplicate live checkpoint is rejected instead of replacing its owner. `wait_until_paused_before(deadline)` reports the exact unreached checkpoint on timeout. Wait for the barrier with a deadline, observe actual metadata quorum with `wait_for_paused_tail_metadata_quorum`, then release or cancel. The metadata check distinguishes a real quorum from disk tasks that have not started. Assert zero remaining rename tasks after the owned coordinator releases its lock; cancellation tests also wait for staging cleanup. | Fixture | Completion boundary | |---|---|