fix(ecstore): fail closed on multipart cleanup gaps

This commit is contained in:
cxymds
2026-09-02 21:53:51 +08:00
parent 2861e15d04
commit 4255e0ca9a
5 changed files with 288 additions and 67 deletions
+58 -11
View File
@@ -232,7 +232,13 @@ mod decommission_lock_order_tests {
C: FnOnce() -> F + Send + 'static,
F: Future<Output = ()> + 'static,
{
const STACK_SIZE: usize = 32 * 1024 * 1024;
const STACK_SIZE: usize = if cfg!(debug_assertions) {
8 * rustfs_config::DEFAULT_THREAD_STACK_SIZE
} else if cfg!(target_os = "macos") {
2 * rustfs_config::DEFAULT_THREAD_STACK_SIZE
} else {
rustfs_config::DEFAULT_THREAD_STACK_SIZE
};
std::thread::Builder::new()
.name(name.to_string())
.stack_size(STACK_SIZE)
@@ -255,7 +261,13 @@ mod decommission_lock_order_tests {
C: FnOnce() -> F + Send + 'static,
F: Future<Output = ()> + 'static,
{
const STACK_SIZE: usize = 32 * 1024 * 1024;
const STACK_SIZE: usize = if cfg!(debug_assertions) {
8 * rustfs_config::DEFAULT_THREAD_STACK_SIZE
} else if cfg!(target_os = "macos") {
2 * rustfs_config::DEFAULT_THREAD_STACK_SIZE
} else {
rustfs_config::DEFAULT_THREAD_STACK_SIZE
};
std::thread::Builder::new()
.name(name.to_string())
.stack_size(STACK_SIZE)
@@ -1496,7 +1508,7 @@ mod decommission_lock_order_tests {
#[tokio::test]
#[serial_test::serial]
async fn multipart_cleanup_proves_target_absence_after_acquiring_capacity_gate() {
async fn multipart_cleanup_rediscovers_uploads_and_proves_target_under_capacity_gate() {
let (_temp_dirs, store, _other_store) = test_three_pool_stores_with_isolated_node_contexts(None).await;
let bucket = test_bucket("multipart-cleanup-target-race");
let object = "published-between-cleanup-proof-and-finalize.bin";
@@ -1549,6 +1561,7 @@ mod decommission_lock_order_tests {
upload_identity.clone(),
);
owner.apply_to(&mut upload_opts);
let concurrent_upload_opts = upload_opts.clone();
let mut cleanup_opts = ObjectOptions {
data_movement: true,
src_pool_idx: 0,
@@ -1627,6 +1640,9 @@ mod decommission_lock_order_tests {
result = &mut cleanup => panic!("cleanup finished before its target-gate proof point: {result:?}"),
}
let concurrent_upload = new_multipart_upload(&store, target_pool_index, &bucket, object, concurrent_upload_opts)
.await
.expect("stage a same-identity upload after cleanup reaches its target-gate boundary");
let mut target_data = PutObjReader::from_vec(body);
let published = store.pools[target_pool_index]
.put_object(
@@ -1655,6 +1671,17 @@ mod decommission_lock_order_tests {
.expect("cleanup should preserve the now-published target intent");
drop(barrier);
let remaining_uploads = store.pools[target_pool_index]
.get_disks_by_key(object)
.data_movement_multipart_upload_ids(&bucket, object, Some(incarnation), &upload_identity)
.await
.expect("check same-identity uploads after gated cleanup");
assert!(
remaining_uploads.is_empty(),
"gated cleanup must rediscover and remove the concurrent same-identity upload {}",
concurrent_upload.upload_id
);
let pool_meta = store.pool_meta.read().await;
let reservation = pool_meta.pools[0]
.decommission
@@ -1672,7 +1699,7 @@ mod decommission_lock_order_tests {
#[test]
#[serial_test::serial]
fn data_movement_multipart_restart_reconciles_published_capacity_before_new_upload() {
run_large_stack_async_test(
run_large_stack_current_thread_async_test(
"multipart-published-capacity-restart",
data_movement_multipart_restart_reconciles_published_capacity_before_new_upload_case,
);
@@ -2009,7 +2036,7 @@ mod decommission_lock_order_tests {
#[test]
#[serial_test::serial]
fn data_movement_multipart_restart_cleans_partial_upload_before_exact_fit_retry() {
run_large_stack_async_test(
run_large_stack_current_thread_async_test(
"multipart-partial-upload-restart",
data_movement_multipart_restart_cleans_partial_upload_before_exact_fit_retry_case,
);
@@ -2311,18 +2338,22 @@ mod decommission_lock_order_tests {
.await
.expect("activate the abort fence reservation");
let owner = decommission_capacity_owner(&*store.pool_meta.read().await).with_mutation_id(uuid::Uuid::new_v4());
let version_id = uuid::Uuid::new_v4().to_string();
let mod_time = time::OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(83);
let mut upload_opts = ObjectOptions {
data_movement: true,
src_pool_idx: 0,
versioned: true,
version_id: Some(uuid::Uuid::new_v4().to_string()),
version_id: Some(version_id.clone()),
mod_time: Some(mod_time),
expected_bucket_incarnation_id: Some(incarnation),
..Default::default()
};
let upload_identity = data_movement::data_movement_upload_identity_from_options(&upload_opts);
rustfs_utils::http::insert_str(
&mut upload_opts.user_defined,
rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD,
"v1:abort-fence:1".to_string(),
upload_identity,
);
let upload = new_multipart_upload(&store, 2, &bucket, object, upload_opts)
.await
@@ -2333,6 +2364,9 @@ mod decommission_lock_order_tests {
let mut abort_opts = ObjectOptions {
data_movement: true,
src_pool_idx: 0,
versioned: true,
version_id: Some(version_id),
mod_time: Some(mod_time),
expected_bucket_incarnation_id: Some(incarnation),
..Default::default()
};
@@ -2387,9 +2421,16 @@ mod decommission_lock_order_tests {
assert!(upload_info.parts.is_empty());
}
#[tokio::test]
#[test]
#[serial_test::serial]
async fn data_movement_multipart_abort_restart_reconciles_inflight_after_release_save_loss() {
fn data_movement_multipart_abort_restart_reconciles_inflight_after_release_save_loss() {
run_large_stack_current_thread_async_test(
"multipart-abort-restart",
data_movement_multipart_abort_restart_reconciles_inflight_after_release_save_loss_case,
);
}
async fn data_movement_multipart_abort_restart_reconciles_inflight_after_release_save_loss_case() {
let (_temp_dirs, store, other_store) = test_three_pool_stores_with_isolated_node_contexts(None).await;
let bucket = test_bucket("multipart-abort-restart");
let object = "deleted-upload-before-release-save.bin";
@@ -2418,15 +2459,18 @@ mod decommission_lock_order_tests {
.await
.expect("activate the abort restart reservation");
let owner = decommission_capacity_owner(&*store.pool_meta.read().await).with_mutation_id(uuid::Uuid::new_v4());
let upload_identity = format!("v1:{}:{}", uuid::Uuid::new_v4(), 91);
let version_id = uuid::Uuid::new_v4().to_string();
let mod_time = time::OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(91);
let mut new_opts = ObjectOptions {
data_movement: true,
src_pool_idx: 0,
versioned: true,
version_id: Some(uuid::Uuid::new_v4().to_string()),
version_id: Some(version_id.clone()),
mod_time: Some(mod_time),
expected_bucket_incarnation_id: Some(incarnation),
..Default::default()
};
let upload_identity = data_movement::data_movement_upload_identity_from_options(&new_opts);
rustfs_utils::http::insert_str(
&mut new_opts.user_defined,
rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD,
@@ -2496,6 +2540,9 @@ mod decommission_lock_order_tests {
let mut abort_opts = ObjectOptions {
data_movement: true,
src_pool_idx: 0,
versioned: true,
version_id: Some(version_id),
mod_time: Some(mod_time),
expected_bucket_incarnation_id: Some(incarnation),
..Default::default()
};
+12 -7
View File
@@ -269,16 +269,21 @@ fn insert_data_movement_checksum(user_defined: &mut HashMap<String, String>, obj
}
}
fn data_movement_upload_identity(object_info: &ObjectInfo) -> String {
let version_id = object_info
.version_id
.map_or_else(|| "none".to_string(), |version_id| version_id.to_string());
let mod_time = object_info
.mod_time
.map_or_else(|| "none".to_string(), |mod_time| mod_time.unix_timestamp_nanos().to_string());
fn data_movement_upload_identity_parts(version_id: Option<&str>, mod_time: Option<time::OffsetDateTime>) -> String {
let version_id = version_id.unwrap_or("none");
let mod_time = mod_time.map_or_else(|| "none".to_string(), |mod_time| mod_time.unix_timestamp_nanos().to_string());
format!("v1:{version_id}:{mod_time}")
}
fn data_movement_upload_identity(object_info: &ObjectInfo) -> String {
let version_id = object_info.version_id.map(|version_id| version_id.to_string());
data_movement_upload_identity_parts(version_id.as_deref(), object_info.mod_time)
}
pub(crate) fn data_movement_upload_identity_from_options(opts: &ObjectOptions) -> String {
data_movement_upload_identity_parts(opts.version_id.as_deref(), opts.mod_time)
}
fn data_movement_new_multipart_opts(object_info: &ObjectInfo, src_pool_idx: usize) -> ObjectOptions {
let mut user_defined = data_movement_user_defined(object_info);
let upload_identity = data_movement_upload_identity(object_info);
+121 -7
View File
@@ -837,7 +837,7 @@ impl SetDisks {
orig_bucket: &str,
error_path: &str,
root_prefix: &str,
) -> Result<(Vec<Option<DiskStore>>, Vec<String>, usize)> {
) -> Result<(Vec<Option<DiskStore>>, Vec<String>, usize, bool)> {
let disks = self.disks.read().await.clone();
if disks.is_empty() {
return Err(Error::ErasureReadQuorum);
@@ -882,16 +882,24 @@ impl SetDisks {
return Err(to_object_err(err.into(), vec![orig_bucket, error_path]));
}
let mut has_minority_candidate = false;
let mut candidate_paths = candidate_counts
.into_iter()
.filter_map(|(path, count)| (count >= discovery_quorum).then_some(path))
.filter_map(|(path, count)| {
if count >= discovery_quorum {
Some(path)
} else {
has_minority_candidate = true;
None
}
})
.collect::<Vec<_>>();
candidate_paths.sort_unstable();
Ok((disks, candidate_paths, discovery_quorum))
Ok((disks, candidate_paths, discovery_quorum, has_minority_candidate))
}
pub(crate) async fn first_multipart_upload_path_for_decommission(&self, bucket: &str) -> Result<Option<String>> {
let (_, paths, _) = self
let (_, paths, _, _) = self
.discover_multipart_upload_paths(bucket, RUSTFS_META_MULTIPART_BUCKET, "")
.await?;
Ok(paths.into_iter().next())
@@ -905,7 +913,14 @@ impl SetDisks {
upload_identity: &str,
) -> Result<Vec<String>> {
let expected_parent = format!("{DATA_MOVEMENT_MULTIPART_PREFIX}/{}", Self::get_multipart_sha_dir(bucket, object));
let (_, candidate_paths, _) = self.discover_multipart_upload_paths(bucket, object, &expected_parent).await?;
let (_, candidate_paths, _, has_minority_candidate) =
self.discover_multipart_upload_paths(bucket, object, &expected_parent).await?;
if has_minority_candidate {
return Err(Error::DecommissionCapacityBlocked {
message: "data movement multipart cleanup found an upload path on fewer than the discovery quorum of disks"
.to_string(),
});
}
let mut upload_ids = Vec::new();
for upload_path in candidate_paths {
let Some((parent, raw_upload_id)) = upload_path.rsplit_once('/') else {
@@ -921,7 +936,11 @@ impl SetDisks {
{
Ok((file_info, _)) => file_info,
Err(err) if crate::error::is_err_invalid_upload_id(&err) || crate::error::is_err_object_not_found(&err) => {
continue;
return Err(Error::DecommissionCapacityBlocked {
message: format!(
"data movement multipart cleanup found quorum-visible upload path {upload_path} without verifiable metadata: {err}"
),
});
}
Err(err) => return Err(err),
};
@@ -1190,7 +1209,7 @@ impl SetDisks {
max_uploads: usize,
expected_incarnation_id: Option<Uuid>,
) -> Result<ListMultipartsInfo> {
let (disks, candidate_paths, discovery_quorum) = self.discover_multipart_upload_paths(bucket, prefix, "").await?;
let (disks, candidate_paths, discovery_quorum, _) = self.discover_multipart_upload_paths(bucket, prefix, "").await?;
let listed_uploads = stream::iter(candidate_paths)
.map(|upload_path| {
let disks = &disks;
@@ -7551,6 +7570,101 @@ mod tests {
assert_eq!(bucket_wide.uploads[0].object, "blobs/data/layer.bin");
}
#[tokio::test]
async fn data_movement_cleanup_discovery_rejects_quorum_visible_upload_without_metadata() {
let (temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "data-movement-unverifiable-upload";
let object = "staged/object.bin";
for disk in &disk_stores {
disk.make_volume(bucket).await.expect("bucket volume should be created");
}
let upload_identity = format!("v1:{}:{}", Uuid::new_v4(), OffsetDateTime::now_utc().unix_timestamp_nanos());
let mut metadata = HashMap::new();
rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD, upload_identity.clone());
let opts = ObjectOptions {
data_movement: true,
user_defined: metadata,
..Default::default()
};
let upload = set_disks
.new_multipart_upload(bucket, object, &opts)
.await
.expect("data movement upload should be created");
let upload_path = SetDisks::get_multipart_upload_dir(bucket, object, &upload.upload_id, true);
for temp_dir in &temp_dirs {
tokio::fs::remove_file(
temp_dir
.path()
.join(RUSTFS_META_MULTIPART_BUCKET)
.join(&upload_path)
.join("xl.meta"),
)
.await
.expect("upload metadata should be removable while preserving its quorum-visible directory");
}
let err = set_disks
.data_movement_multipart_upload_ids(bucket, object, None, &upload_identity)
.await
.expect_err("cleanup discovery must fail closed on an unverifiable upload path");
assert!(matches!(err, Error::DecommissionCapacityBlocked { .. }));
}
#[tokio::test]
async fn data_movement_cleanup_discovery_rejects_minority_upload_until_disks_recover() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "data-movement-minority-upload";
let object = "staged/object.bin";
for disk in &disk_stores {
disk.make_volume(bucket).await.expect("bucket volume should be created");
}
{
let mut disks = set_disks.disks.write().await;
disks[3] = None;
}
let upload_identity = format!("v1:{}:{}", Uuid::new_v4(), OffsetDateTime::now_utc().unix_timestamp_nanos());
let mut metadata = HashMap::new();
rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD, upload_identity.clone());
let upload = set_disks
.new_multipart_upload(
bucket,
object,
&ObjectOptions {
data_movement: true,
user_defined: metadata,
..Default::default()
},
)
.await
.expect("write quorum should create the upload while one disk is offline");
{
let mut disks = set_disks.disks.write().await;
disks[1] = None;
disks[2] = None;
disks[3] = Some(disk_stores[3].clone());
}
let err = set_disks
.data_movement_multipart_upload_ids(bucket, object, None, &upload_identity)
.await
.expect_err("a minority-observed upload must block a destructive absence proof");
assert!(matches!(err, Error::DecommissionCapacityBlocked { .. }));
{
let mut disks = set_disks.disks.write().await;
disks[1] = Some(disk_stores[1].clone());
disks[2] = Some(disk_stores[2].clone());
}
let recovered = set_disks
.data_movement_multipart_upload_ids(bucket, object, None, &upload_identity)
.await
.expect("recovered quorum should make the staged upload verifiable");
assert_eq!(recovered.len(), 1);
assert_eq!(upload_uuid_suffix(&recovered[0]), upload_uuid_suffix(&upload.upload_id));
}
/// Regression (issue #5716): a single upload directory whose `xl.meta` was
/// destroyed (crash mid-write, torn disk state) must degrade to that upload
/// alone. Failing the whole ListMultipartUploads turns one piece of stale
+20 -7
View File
@@ -923,7 +923,13 @@ mod tests {
C: FnOnce() -> F + Send + 'static,
F: Future<Output = ()> + 'static,
{
const STACK_SIZE: usize = 32 * 1024 * 1024;
const STACK_SIZE: usize = if cfg!(debug_assertions) {
8 * rustfs_config::DEFAULT_THREAD_STACK_SIZE
} else if cfg!(target_os = "macos") {
2 * rustfs_config::DEFAULT_THREAD_STACK_SIZE
} else {
rustfs_config::DEFAULT_THREAD_STACK_SIZE
};
std::thread::Builder::new()
.name(name.to_string())
.stack_size(STACK_SIZE)
@@ -2418,19 +2424,23 @@ mod tests {
let capacity_owner = test_decommission_capacity_owner(store.as_ref(), 0)
.await
.with_mutation_id(Uuid::new_v4());
let mut upload_metadata = source.user_defined.as_ref().clone();
rustfs_utils::http::insert_str(
&mut upload_metadata,
rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD,
"mpu-staging-test".to_string(),
);
let upload_metadata = source.user_defined.as_ref().clone();
let mut staging_opts = ObjectOptions {
data_movement: true,
src_pool_idx: 0,
versioned: source.version_id.is_some(),
version_id: source.version_id.map(|version_id| version_id.to_string()),
mod_time: source.mod_time,
user_defined: upload_metadata,
expected_bucket_incarnation_id: Some(bucket_incarnation_id),
..Default::default()
};
let upload_identity = crate::data_movement::data_movement_upload_identity_from_options(&staging_opts);
rustfs_utils::http::insert_str(
&mut staging_opts.user_defined,
rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD,
upload_identity,
);
capacity_owner.apply_to(&mut staging_opts);
let (upload, target_pool_idx, staged_incarnation_id) = store
.handle_new_multipart_upload_with_pool_idx(&bucket, object, &staging_opts, None)
@@ -2500,6 +2510,9 @@ mod tests {
let mut abort_opts = ObjectOptions {
data_movement: true,
src_pool_idx: 0,
versioned: source.version_id.is_some(),
version_id: source.version_id.map(|version_id| version_id.to_string()),
mod_time: source.mod_time,
expected_bucket_incarnation_id: Some(bucket_incarnation_id),
..Default::default()
};
+77 -35
View File
@@ -787,8 +787,16 @@ impl ECStore {
upload_id: &str,
opts: &ObjectOptions,
) -> Result<()> {
self.abort_multipart_uploads_for_data_movement(target_pool_idx, bucket, object, &[upload_id.to_owned()], None, opts)
.await
let upload_identity = crate::data_movement::data_movement_upload_identity_from_options(opts);
self.abort_multipart_uploads_for_data_movement(
target_pool_idx,
bucket,
object,
&[upload_id.to_owned()],
&upload_identity,
opts,
)
.await
}
pub(crate) async fn reconcile_multipart_uploads_for_data_movement(
@@ -799,8 +807,7 @@ impl ECStore {
upload_identity: &str,
opts: &ObjectOptions,
) -> Result<()> {
let pool = self
.pools
self.pools
.get(target_pool_idx)
.ok_or_else(|| Error::other(format!("data movement target pool {target_pool_idx} is out of range")))?;
let owner = DecommissionCapacityOwner::from_options(opts)
@@ -830,11 +837,7 @@ impl ECStore {
.entry(self.id)
.or_default() += 1;
}
let set = pool.get_disks_by_key(object);
let upload_ids = set
.data_movement_multipart_upload_ids(bucket, object, opts.expected_bucket_incarnation_id, upload_identity)
.await?;
self.abort_multipart_uploads_for_data_movement(target_pool_idx, bucket, object, &upload_ids, Some(upload_identity), opts)
self.abort_multipart_uploads_for_data_movement(target_pool_idx, bucket, object, &[], upload_identity, opts)
.await
}
@@ -844,7 +847,7 @@ impl ECStore {
bucket: &str,
object: &str,
upload_ids: &[String],
expected_upload_identity: Option<&str>,
expected_upload_identity: &str,
opts: &ObjectOptions,
) -> Result<()> {
check_new_multipart_args(bucket, object)?;
@@ -861,32 +864,78 @@ impl ECStore {
.pools
.get(target_pool_idx)
.ok_or_else(|| Error::other(format!("data movement target pool {target_pool_idx} is out of range")))?;
let set = pool.get_disks_by_key(object);
let mut guards = Vec::with_capacity(upload_ids.len());
for upload_id in upload_ids {
if let Some(guard) = set
.lock_data_movement_multipart_abort(bucket, object, upload_id, expected_upload_identity, &opts)
.await?
{
guard.add_namespace_lock_fence(&mut opts);
guards.push(guard);
}
}
opts.no_lock = true;
// Keep every upload namespace guard alive through the final capacity progress save.
let cleanup_decision_error = self
let set = pool.get_disks_by_key(object);
// Discover and lock uploads only after the target capacity gate is held.
// Return the guards so they remain alive through the final capacity save.
let (cleanup_decision_error, guards) = self
.run_decommission_capacity_temporary_release_with_capacity_lease(target_pool_idx, capacity_owner, |capacity_lease| {
let mut delete_opts = opts.clone();
let guards = &guards;
let set = &set;
let pool = &pool;
async move {
if let Some(capacity_lease) = capacity_lease.as_ref() {
delete_opts.add_namespace_lock_lost_signal(Arc::clone(capacity_lease));
}
// The target capacity gate held by the surrounding
// transaction makes this proof and its pending-ledger
// decision one critical section with cleanup finalize.
let mut candidate_upload_ids = upload_ids.to_vec();
candidate_upload_ids.extend(
set.data_movement_multipart_upload_ids(
bucket,
object,
delete_opts.expected_bucket_incarnation_id,
expected_upload_identity,
)
.await?,
);
candidate_upload_ids.sort_unstable();
candidate_upload_ids.dedup();
let mut guards = Vec::with_capacity(candidate_upload_ids.len());
for upload_id in &candidate_upload_ids {
match set
.lock_data_movement_multipart_abort(
bucket,
object,
upload_id,
Some(expected_upload_identity),
&delete_opts,
)
.await
{
Ok(Some(guard)) => {
guard.add_namespace_lock_fence(&mut delete_opts);
guards.push(guard);
}
Ok(None) => {}
Err(err) => return Err(err),
}
}
for guard in &guards {
match guard.delete(set, bucket, object, &delete_opts).await {
Ok(()) => {}
Err(err) if is_err_invalid_upload_id(&err) => {}
Err(err) => return Err(err),
}
}
if !set
.data_movement_multipart_upload_ids(
bucket,
object,
delete_opts.expected_bucket_incarnation_id,
expected_upload_identity,
)
.await?
.is_empty()
{
return Err(Error::DecommissionCapacityBlocked {
message: "multipart cleanup could not prove the exact staged uploads are absent".to_string(),
});
}
// The target capacity gate makes the exact target proof,
// upload absence proof, and pending-ledger decision one
// critical section with cleanup finalize.
let (clear_pending, cleanup_decision_error) = if capacity_owner.is_some() {
let mut lookup_opts = ObjectOptions {
versioned: delete_opts.versioned,
@@ -911,14 +960,7 @@ impl ECStore {
} else {
(true, None)
};
for guard in guards {
match guard.delete(set, bucket, object, &delete_opts).await {
Ok(()) => {}
Err(err) if is_err_invalid_upload_id(&err) => {}
Err(err) => return Err(err),
}
}
Ok((cleanup_decision_error, clear_pending))
Ok(((cleanup_decision_error, guards), clear_pending))
}
})
.await?;