mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 04:16:38 +00:00
fix(ecstore): reuse fixed fence for reverse decommission
This commit is contained in:
@@ -1447,11 +1447,8 @@ async fn migrate_object_inner(
|
||||
if should_use_multipart_data_movement(&object_info, has_part_checksums) {
|
||||
let mut new_multipart_opts = data_movement_new_multipart_opts(&object_info, pool_idx);
|
||||
new_multipart_opts.expected_bucket_incarnation_id = source_bucket_incarnation_id;
|
||||
if let Some(fence) = mutation_fence {
|
||||
fence.add_namespace_lock_fence(&mut new_multipart_opts);
|
||||
}
|
||||
let (res, target_pool_idx, expected_bucket_incarnation_id) = match store
|
||||
.handle_new_multipart_upload_with_pool_idx(&bucket, &object_info.name, &new_multipart_opts)
|
||||
.handle_new_multipart_upload_with_pool_idx(&bucket, &object_info.name, &new_multipart_opts, mutation_fence)
|
||||
.await
|
||||
{
|
||||
Ok(res) => res,
|
||||
@@ -1546,9 +1543,6 @@ async fn migrate_object_inner(
|
||||
)
|
||||
})?;
|
||||
complete_multipart_opts.expected_bucket_incarnation_id = expected_bucket_incarnation_id;
|
||||
if let Some(fence) = mutation_fence {
|
||||
fence.add_namespace_lock_fence(&mut complete_multipart_opts);
|
||||
}
|
||||
if let Err(err) = store
|
||||
.clone()
|
||||
.complete_multipart_upload_for_data_movement(
|
||||
@@ -1558,6 +1552,7 @@ async fn migrate_object_inner(
|
||||
&res.upload_id,
|
||||
parts,
|
||||
&complete_multipart_opts,
|
||||
mutation_fence,
|
||||
)
|
||||
.await
|
||||
{
|
||||
@@ -1712,11 +1707,8 @@ async fn migrate_object_inner(
|
||||
|
||||
let mut put_opts = data_movement_put_object_opts(&object_info, pool_idx);
|
||||
put_opts.expected_bucket_incarnation_id = source_bucket_incarnation_id;
|
||||
if let Some(fence) = mutation_fence {
|
||||
fence.add_namespace_lock_fence(&mut put_opts);
|
||||
}
|
||||
let (target_pool_idx, put_result) = store
|
||||
.put_object_for_data_movement(&bucket, &object_info.name, &mut data, &put_opts)
|
||||
.put_object_for_data_movement(&bucket, &object_info.name, &mut data, &put_opts, mutation_fence)
|
||||
.await
|
||||
.map_err(|err| data_movement_stage_error(op_label, "prepare_put_object", &bucket, &object_info.name, err))?;
|
||||
if let Err(err) = put_result {
|
||||
|
||||
@@ -2969,6 +2969,196 @@ mod tests {
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
async fn reverse_decommission_reuses_fixed_target_fence_for_put_and_multipart() {
|
||||
let temp_dir = tempfile::tempdir().expect("create reverse decommission store dir");
|
||||
let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store_with_layout(
|
||||
temp_dir.path(),
|
||||
"reverse-decommission-fixed-target",
|
||||
&[(1, 4), (1, 4)],
|
||||
CancellationToken::new(),
|
||||
))
|
||||
.await;
|
||||
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
|
||||
|
||||
let bucket = format!("reverse-decommission-fixed-target-{}", uuid::Uuid::new_v4());
|
||||
let object = "ordinary.bin";
|
||||
let object_body = b"reverse ordinary generation".to_vec();
|
||||
let multipart_object = "multipart.bin";
|
||||
let first_part = vec![b'm'; 5 * 1024 * 1024];
|
||||
let second_part = b"reverse multipart tail".to_vec();
|
||||
let mut multipart_body = first_part.clone();
|
||||
multipart_body.extend_from_slice(&second_part);
|
||||
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("create reverse decommission bucket");
|
||||
let mut source = PutObjReader::from_vec(object_body.clone());
|
||||
store.pools[1]
|
||||
.put_object(&bucket, object, &mut source, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("write ordinary source object to pool 1");
|
||||
|
||||
let upload = store.pools[1]
|
||||
.new_multipart_upload(&bucket, multipart_object, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("create source multipart upload in pool 1");
|
||||
let mut completed_parts = Vec::with_capacity(2);
|
||||
for (part_number, bytes) in [(1, first_part.as_slice()), (2, second_part.as_slice())] {
|
||||
let mut reader = PutObjReader::from_vec(bytes.to_vec());
|
||||
let part = store.pools[1]
|
||||
.put_object_part(
|
||||
&bucket,
|
||||
multipart_object,
|
||||
&upload.upload_id,
|
||||
part_number,
|
||||
&mut reader,
|
||||
&ObjectOptions::default(),
|
||||
)
|
||||
.await
|
||||
.expect("write source multipart part");
|
||||
completed_parts.push(crate::storage_api_contracts::multipart::CompletePart {
|
||||
part_num: part.part_num,
|
||||
etag: part.etag,
|
||||
..Default::default()
|
||||
});
|
||||
}
|
||||
store.pools[1]
|
||||
.clone()
|
||||
.complete_multipart_upload(&bucket, multipart_object, &upload.upload_id, completed_parts, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("complete source multipart object in pool 1");
|
||||
|
||||
{
|
||||
let mut pool_meta = store.pool_meta.write().await;
|
||||
pool_meta.pools[1].decommission = Some(PoolDecommissionInfo {
|
||||
start_time: Some(OffsetDateTime::now_utc()),
|
||||
..Default::default()
|
||||
});
|
||||
}
|
||||
assert!(store.is_suspended(1).await, "pool 1 must be the reverse decommission source");
|
||||
|
||||
let commit_barrier = crate::set_disk::PutObjectCommitBarrier::install(
|
||||
&bucket,
|
||||
object,
|
||||
crate::set_disk::PutObjectCommitPause::BeforeNamespace,
|
||||
);
|
||||
let source_set = store.pools[1].get_disks_by_key(object);
|
||||
let worker_store = Arc::clone(&store);
|
||||
let worker_bucket = bucket.clone();
|
||||
let worker = tokio::spawn(async move {
|
||||
worker_store
|
||||
.decommission_entry_for_test(
|
||||
1,
|
||||
MetaCacheEntry {
|
||||
name: object.to_string(),
|
||||
..Default::default()
|
||||
},
|
||||
worker_bucket,
|
||||
source_set,
|
||||
)
|
||||
.await
|
||||
});
|
||||
commit_barrier.wait_until_paused().await;
|
||||
|
||||
let delete_barrier = crate::store::object::DeleteAfterObjectLockSnapshotBarrier::install(&bucket);
|
||||
let delete_store = Arc::clone(&store);
|
||||
let delete_bucket = bucket.clone();
|
||||
let delete = tokio::spawn(async move {
|
||||
delete_store
|
||||
.delete_object(&delete_bucket, object, ObjectOptions::default())
|
||||
.await
|
||||
});
|
||||
delete_barrier.wait_until_paused().await;
|
||||
delete_barrier.release_and_wait_until_namespace_pending().await;
|
||||
assert!(
|
||||
!delete_barrier.namespace_acquired() && !delete.is_finished(),
|
||||
"the reverse target commit must keep DELETE behind the fixed read fence"
|
||||
);
|
||||
delete.abort();
|
||||
assert!(
|
||||
delete
|
||||
.await
|
||||
.expect_err("the blocked DELETE should be canceled")
|
||||
.is_cancelled(),
|
||||
"canceling the blocked DELETE must not mutate either pool"
|
||||
);
|
||||
drop(delete_barrier);
|
||||
|
||||
commit_barrier.release();
|
||||
drop(commit_barrier);
|
||||
tokio::time::timeout(Duration::from_secs(60), worker)
|
||||
.await
|
||||
.expect("reverse ordinary decommission must not self-deadlock on the fixed target set")
|
||||
.expect("reverse ordinary decommission worker should join")
|
||||
.expect("reverse ordinary decommission should complete");
|
||||
|
||||
let mut ordinary_reader = store.pools[0]
|
||||
.get_object_reader(&bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
|
||||
.await
|
||||
.expect("read the ordinary object from the fixed target set");
|
||||
let mut ordinary_target_body = Vec::new();
|
||||
ordinary_reader
|
||||
.stream
|
||||
.read_to_end(&mut ordinary_target_body)
|
||||
.await
|
||||
.expect("drain the ordinary target body");
|
||||
assert_eq!(ordinary_target_body, object_body, "ordinary migration must preserve the full body");
|
||||
let ordinary_source_err = store.pools[1]
|
||||
.get_object_info(&bucket, object, &ObjectOptions::default())
|
||||
.await
|
||||
.expect_err("ordinary source generation must be cleaned after migration");
|
||||
assert!(matches!(ordinary_source_err, StorageError::ObjectNotFound(_, _)));
|
||||
|
||||
let multipart_source_set = store.pools[1].get_disks_by_key(multipart_object);
|
||||
let multipart_store = Arc::clone(&store);
|
||||
let multipart_bucket = bucket.clone();
|
||||
let multipart_worker = tokio::spawn(async move {
|
||||
multipart_store
|
||||
.decommission_entry_for_test(
|
||||
1,
|
||||
MetaCacheEntry {
|
||||
name: multipart_object.to_string(),
|
||||
..Default::default()
|
||||
},
|
||||
multipart_bucket,
|
||||
multipart_source_set,
|
||||
)
|
||||
.await
|
||||
});
|
||||
tokio::time::timeout(Duration::from_secs(60), multipart_worker)
|
||||
.await
|
||||
.expect("reverse multipart decommission must not self-deadlock on new or complete")
|
||||
.expect("reverse multipart decommission worker should join")
|
||||
.expect("reverse multipart decommission should complete");
|
||||
|
||||
let target_info = store.pools[0]
|
||||
.get_object_info(&bucket, multipart_object, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("read migrated multipart metadata from the fixed target set");
|
||||
assert!(target_info.is_multipart(), "migration must retain multipart identity");
|
||||
let mut multipart_reader = store.pools[0]
|
||||
.get_object_reader(&bucket, multipart_object, None, HeaderMap::new(), &ObjectOptions::default())
|
||||
.await
|
||||
.expect("read migrated multipart object from the fixed target set");
|
||||
let mut multipart_target_body = Vec::new();
|
||||
multipart_reader
|
||||
.stream
|
||||
.read_to_end(&mut multipart_target_body)
|
||||
.await
|
||||
.expect("drain the multipart target body");
|
||||
assert_eq!(multipart_target_body, multipart_body, "multipart migration must preserve the full body");
|
||||
let multipart_source_err = store.pools[1]
|
||||
.get_object_info(&bucket, multipart_object, &ObjectOptions::default())
|
||||
.await
|
||||
.expect_err("multipart source generation must be cleaned after migration");
|
||||
assert!(matches!(multipart_source_err, StorageError::ObjectNotFound(_, _)));
|
||||
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
async fn batch_delete_real_path_preserves_source_pool_errors_in_any_pool_order() {
|
||||
|
||||
@@ -400,7 +400,7 @@ impl ECStore {
|
||||
object: &str,
|
||||
opts: &ObjectOptions,
|
||||
) -> Result<MultipartUploadResult> {
|
||||
self.handle_new_multipart_upload_with_pool_idx(bucket, object, opts)
|
||||
self.handle_new_multipart_upload_with_pool_idx(bucket, object, opts, None)
|
||||
.await
|
||||
.map(|(res, _, _)| res)
|
||||
}
|
||||
@@ -410,20 +410,21 @@ impl ECStore {
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
opts: &ObjectOptions,
|
||||
mutation_fence: Option<&ObjectLockDiagGuard>,
|
||||
) -> Result<(MultipartUploadResult, usize, Option<Uuid>)> {
|
||||
check_new_multipart_args(bucket, object)?;
|
||||
let (opts, _bucket_lifecycle_guard) = self.guard_multipart_bucket_incarnation(bucket, opts).await?;
|
||||
let opts = &opts;
|
||||
let (mut opts, _bucket_lifecycle_guard) = self.guard_multipart_bucket_incarnation(bucket, opts).await?;
|
||||
|
||||
if self.single_pool() {
|
||||
self.apply_decommission_target_mutation_fence(0, object, &mut opts, mutation_fence);
|
||||
return self.pools[0]
|
||||
.new_multipart_upload(bucket, object, opts)
|
||||
.new_multipart_upload(bucket, object, &opts)
|
||||
.await
|
||||
.map(|res| (res, 0, opts.expected_bucket_incarnation_id));
|
||||
}
|
||||
|
||||
if opts.data_movement && opts.version_id.is_some() {
|
||||
let idx = self.select_data_movement_pool_idx(bucket, object, -1, opts, false).await?;
|
||||
let idx = self.select_data_movement_pool_idx(bucket, object, -1, &opts, false).await?;
|
||||
if idx == opts.src_pool_idx {
|
||||
return Err(StorageError::DataMovementOverwriteErr(
|
||||
bucket.to_owned(),
|
||||
@@ -431,7 +432,8 @@ impl ECStore {
|
||||
opts.version_id.clone().unwrap_or_default(),
|
||||
));
|
||||
}
|
||||
let res = self.pools[idx].new_multipart_upload(bucket, object, opts).await?;
|
||||
self.apply_decommission_target_mutation_fence(idx, object, &mut opts, mutation_fence);
|
||||
let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?;
|
||||
return Ok((res, idx, opts.expected_bucket_incarnation_id));
|
||||
}
|
||||
|
||||
@@ -454,7 +456,8 @@ impl ECStore {
|
||||
.await?;
|
||||
|
||||
if !res.uploads.is_empty() {
|
||||
let res = self.pools[idx].new_multipart_upload(bucket, object, opts).await?;
|
||||
self.apply_decommission_target_mutation_fence(idx, object, &mut opts, mutation_fence);
|
||||
let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?;
|
||||
return Ok((res, idx, opts.expected_bucket_incarnation_id));
|
||||
}
|
||||
}
|
||||
@@ -467,7 +470,8 @@ impl ECStore {
|
||||
));
|
||||
}
|
||||
|
||||
let res = self.pools[idx].new_multipart_upload(bucket, object, opts).await?;
|
||||
self.apply_decommission_target_mutation_fence(idx, object, &mut opts, mutation_fence);
|
||||
let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?;
|
||||
Ok((res, idx, opts.expected_bucket_incarnation_id))
|
||||
}
|
||||
|
||||
@@ -710,6 +714,7 @@ impl ECStore {
|
||||
upload_id: &str,
|
||||
uploaded_parts: Vec<CompletePart>,
|
||||
opts: &ObjectOptions,
|
||||
mutation_fence: Option<&ObjectLockDiagGuard>,
|
||||
) -> Result<ObjectInfo> {
|
||||
check_complete_multipart_args(bucket, object, upload_id)?;
|
||||
if !opts.data_movement {
|
||||
@@ -739,6 +744,7 @@ impl ECStore {
|
||||
snapshot.add_lock_fences(&mut opts);
|
||||
opts.object_lock_config_snapshot = Some(snapshot);
|
||||
}
|
||||
self.apply_decommission_target_mutation_fence(target_pool_idx, object, &mut opts, mutation_fence);
|
||||
#[cfg(test)]
|
||||
pause_data_movement_multipart_before_selected_completion(bucket).await;
|
||||
let pool = self
|
||||
|
||||
@@ -1834,6 +1834,25 @@ impl ECStore {
|
||||
.ok_or_else(|| Error::other("decommission object migration failed to acquire its namespace fence"))
|
||||
}
|
||||
|
||||
pub(super) fn apply_decommission_target_mutation_fence(
|
||||
&self,
|
||||
target_pool_idx: usize,
|
||||
object: &str,
|
||||
opts: &mut ObjectOptions,
|
||||
mutation_fence: Option<&ObjectLockDiagGuard>,
|
||||
) {
|
||||
let Some(mutation_fence) = mutation_fence else {
|
||||
return;
|
||||
};
|
||||
|
||||
mutation_fence.add_namespace_lock_fence(opts);
|
||||
let fixed_set = self.pools.first().and_then(|pool| pool.disk_set.first());
|
||||
let target_set = self.pools.get(target_pool_idx).map(|pool| pool.get_disks_by_key(object));
|
||||
// The fixed read fence can replace the target write acquisition only
|
||||
// when both names resolve to the exact same SetDisks namespace.
|
||||
opts.no_lock = matches!((fixed_set, target_set), (Some(fixed), Some(target)) if Arc::ptr_eq(fixed, &target));
|
||||
}
|
||||
|
||||
pub(crate) async fn acquire_decommission_source_cleanup_fence(
|
||||
&self,
|
||||
bucket: &str,
|
||||
@@ -2316,14 +2335,16 @@ impl ECStore {
|
||||
object: &str,
|
||||
data: &mut PutObjReader,
|
||||
opts: &ObjectOptions,
|
||||
mutation_fence: Option<&ObjectLockDiagGuard>,
|
||||
) -> Result<(usize, Result<ObjectInfo>)> {
|
||||
if !opts.data_movement {
|
||||
return Err(Error::other("data movement PUT requires data_movement options"));
|
||||
}
|
||||
let (object, opts) = self.prepare_put_object(bucket, object, opts).await?;
|
||||
let (object, mut opts) = self.prepare_put_object(bucket, object, opts).await?;
|
||||
let idx = self
|
||||
.select_put_object_pool_idx(bucket, object.as_str(), data.size(), &opts)
|
||||
.await?;
|
||||
self.apply_decommission_target_mutation_fence(idx, object.as_str(), &mut opts, mutation_fence);
|
||||
let result = self.pools[idx]
|
||||
.put_object_with_old_current_size(bucket, &object, data, &opts)
|
||||
.await
|
||||
|
||||
Reference in New Issue
Block a user