From f7757e6437f117fb4f0329bda7bffda8b6a40fe6 Mon Sep 17 00:00:00 2001 From: cxymds Date: Wed, 29 Jul 2026 15:05:53 +0800 Subject: [PATCH] test(ecstore): rendezvous concurrent resend commits (#5411) * test(ecstore): rendezvous concurrent resend commits * test(ecstore): satisfy multipart barrier lint --- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 19 +++--- crates/ecstore/src/set_disk/ops/multipart.rs | 59 +++++++++++++++---- 2 files changed, 56 insertions(+), 22 deletions(-) diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 5f270668f..e13a8caa3 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -11306,23 +11306,25 @@ mod tests { // Distinct payloads with distinct sizes: a mixed-generation reassembly // would produce bytes matching none of them (or fail the read outright). - let candidates: Vec> = (0..3) + let candidates: Vec> = (0..2) .map(|g| { let len = 4096 + g * 512; vec![b'a' + g as u8; len] }) .collect(); - let commit_barrier = MultipartCommitBarrier::install(&bucket, object, MultipartCommitPause::PutPartBeforeLockLost); - let start = Arc::new(tokio::sync::Barrier::new(candidates.len() + 1)); + let commit_barrier = MultipartCommitBarrier::install_for_arrivals( + &bucket, + object, + MultipartCommitPause::PutPartBeforeLockAcquire, + candidates.len(), + ); let mut tasks = tokio::task::JoinSet::new(); for payload in candidates.iter().cloned() { let store = ecstore.clone(); let bucket = bucket.clone(); let upload_id = upload.upload_id.clone(); - let start = Arc::clone(&start); tasks.spawn(async move { - start.wait().await; let mut data = PutObjReader::from_vec(payload.clone()); store .put_object_part(&bucket, object, &upload_id, 1, &mut data, &ObjectOptions::default()) @@ -11330,11 +11332,10 @@ mod tests { .map(|info| (info, payload)) }); } - start.wait().await; - // The first writer holds the uploadId commit lock while the other - // resends reach the same critical section. Releasing it proves the - // handoff without depending on saturated CI disk latency. + // Both writers finish streaming before racing for the uploadId commit + // lock. Two generations are sufficient to exercise the mixed-shard + // hazard, while each waiter sits behind at most one cross-disk rename. commit_barrier.wait_until_paused().await; commit_barrier.release(); diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 6eaaa3399..828bfc423 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -27,7 +27,7 @@ use crate::multipart_listing::paginate_multipart_listing; use futures::{StreamExt, stream}; use std::future::Future; #[cfg(test)] -use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::atomic::{AtomicUsize, Ordering}; use std::time::Duration; use tokio::task::JoinSet; @@ -36,6 +36,7 @@ const MULTIPART_LIST_IO_CONCURRENCY: usize = 16; #[cfg(test)] #[derive(Clone, Copy, PartialEq, Eq)] pub(crate) enum MultipartCommitPause { + PutPartBeforeLockAcquire, PutPartBeforeLockLost, PutPartAfterRename, BeforeLockLost, @@ -47,9 +48,10 @@ struct MultipartCommitBarrierState { bucket: String, object: String, pause: MultipartCommitPause, - armed: AtomicBool, + expected_arrivals: usize, + arrivals: AtomicUsize, arrived: tokio::sync::Notify, - release: tokio::sync::Notify, + release: tokio::sync::Semaphore, } #[cfg(test)] @@ -64,13 +66,24 @@ static MULTIPART_COMMIT_BARRIER: std::sync::OnceLock Self { + Self::install_for_arrivals(bucket, object, pause, 1) + } + + pub(crate) fn install_for_arrivals( + bucket: &str, + object: &str, + pause: MultipartCommitPause, + expected_arrivals: usize, + ) -> Self { + assert!(expected_arrivals > 0, "multipart commit barrier must wait for at least one arrival"); let state = Arc::new(MultipartCommitBarrierState { bucket: bucket.to_string(), object: object.to_string(), pause, - armed: AtomicBool::new(true), + expected_arrivals, + arrivals: AtomicUsize::new(0), arrived: tokio::sync::Notify::new(), - release: tokio::sync::Notify::new(), + release: tokio::sync::Semaphore::new(0), }); let mut slot = MULTIPART_COMMIT_BARRIER .get_or_init(|| std::sync::Mutex::new(None)) @@ -83,20 +96,28 @@ impl MultipartCommitBarrier { } pub(crate) async fn wait_until_paused(&self) { - tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified()) - .await - .expect("multipart completion should reach the deterministic commit barrier"); + tokio::time::timeout(Duration::from_secs(30), async { + loop { + let arrived = self.state.arrived.notified(); + if self.state.arrivals.load(Ordering::Acquire) >= self.state.expected_arrivals { + return; + } + arrived.await; + } + }) + .await + .expect("multipart completion should reach the deterministic commit barrier"); } pub(crate) fn release(&self) { - self.state.release.notify_one(); + self.state.release.add_permits(self.state.expected_arrivals); } } #[cfg(test)] impl Drop for MultipartCommitBarrier { fn drop(&mut self) { - self.state.release.notify_one(); + self.release(); let mut slot = MULTIPART_COMMIT_BARRIER .get_or_init(|| std::sync::Mutex::new(None)) .lock() @@ -117,10 +138,20 @@ async fn pause_multipart_commit(bucket: &str, object: &str, pause: MultipartComm .filter(|barrier| barrier.bucket == bucket && barrier.object == object && barrier.pause == pause) .cloned(); if let Some(barrier) = barrier - && barrier.armed.swap(false, Ordering::AcqRel) + && let Ok(previous) = barrier.arrivals.fetch_update(Ordering::AcqRel, Ordering::Acquire, |current| { + (current < barrier.expected_arrivals).then_some(current + 1) + }) { - barrier.arrived.notify_one(); - barrier.release.notified().await; + let arrival = previous + 1; + if arrival == barrier.expected_arrivals { + barrier.arrived.notify_one(); + } + barrier + .release + .acquire() + .await + .expect("multipart commit barrier should remain open") + .forget(); } } @@ -693,6 +724,8 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { let part_path = format!("{}/{}/{}", upload_id_path, fi.data_dir.unwrap_or_default(), part_suffix); + #[cfg(test)] + pause_multipart_commit(bucket, object, MultipartCommitPause::PutPartBeforeLockAcquire).await; // Serialize only the commit (rename_part), not the whole upload. Each // concurrent stream writes to its own unique temp dir (see `tmp_part` // above), so the encode/stream phase never conflicts and must stay