diff --git a/crates/ecstore/src/data_movement/mod.rs b/crates/ecstore/src/data_movement/mod.rs index 3d7b34fab..147e0846d 100644 --- a/crates/ecstore/src/data_movement/mod.rs +++ b/crates/ecstore/src/data_movement/mod.rs @@ -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 { diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 4bc01d274..7fc9958fa 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -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() { diff --git a/crates/ecstore/src/store/multipart.rs b/crates/ecstore/src/store/multipart.rs index 2c0b18b2a..505938cb7 100644 --- a/crates/ecstore/src/store/multipart.rs +++ b/crates/ecstore/src/store/multipart.rs @@ -400,7 +400,7 @@ impl ECStore { object: &str, opts: &ObjectOptions, ) -> Result { - 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)> { 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, opts: &ObjectOptions, + mutation_fence: Option<&ObjectLockDiagGuard>, ) -> Result { 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 diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 384956bc5..ccde2194b 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -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)> { 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