test(storage): make commit fault checkpoints independent (#8326)

* test(storage): make commit fault checkpoints independent

* test(storage): simplify commit checkpoint lookup
This commit is contained in:
Chris
2026-10-04 01:19:28 +08:00
committed by GitHub
parent 47d208325e
commit a760d567b6
2 changed files with 105 additions and 30 deletions
+104 -29
View File
@@ -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<Notify>,
/// 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<String, Armed>,
/// 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<String, Arc<AtomicUsize>>,
}
@@ -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<Notify>,
release: Arc<Notify>,
paused: Arc<AtomicBool>,
@@ -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<Phase>) {
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,
@@ -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 |
|---|---|