diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 08bd28a29..b0fe789ad 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -4334,15 +4334,20 @@ impl ECStore { ) -> Result<()> { warn!("decommission_object: start {} {}", &bucket, &rd.object_info.name); let object_name = rd.object_info.name.clone(); - let result = data_movement::migrate_decommission_object( + let mut migration = tokio::task::JoinSet::new(); + migration.spawn(data_movement::migrate_decommission_object( self, pool_idx, bucket.clone(), rd, expected_bucket_incarnation_id, "decommission_object", - ) - .await; + )); + let result = migration + .join_next() + .await + .ok_or_else(|| Error::other("decommission migration task was not started"))? + .map_err(|err| Error::other(format!("decommission migration task join error: {err}")))?; if result.is_ok() { warn!("decommission_object: migrated {} {}", &bucket, &object_name); } diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 2d7270cff..9994f8612 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -858,6 +858,7 @@ const EVENT_DISK_LOCAL_DIRECT_IO_FALLBACK: &str = "disk_local_direct_io_fallback #[cfg(target_os = "linux")] const EVENT_DISK_LOCAL_URING_LATCH_OFF: &str = "disk_local_uring_latch_off"; const EVENT_DISK_LOCAL_DELETE_FAILED: &str = "disk_local_delete_failed"; +const EVENT_DISK_LOCAL_DELETE_ROLLBACK_FAILED: &str = "disk_local_delete_rollback_failed"; const EVENT_DISK_LOCAL_CHECK_PARTS: &str = "disk_local_check_parts"; const EVENT_DISK_LOCAL_ACCESS_FAILED: &str = "disk_local_access_failed"; const EVENT_DISK_LOCAL_VOLUME_SETUP_FAILED: &str = "disk_local_volume_setup_failed"; @@ -6106,6 +6107,43 @@ impl LocalDisk { Ok((bytes, modtime)) } + async fn write_missing_delete_marker( + &self, + volume: &str, + path: &str, + fi: FileInfo, + object_dir: &Path, + xl_path: &Path, + rollback_dir: Option, + ) -> Result<()> { + if let Some(rollback_dir) = rollback_dir { + let rollback_path = object_dir.join(rollback_dir.to_string()); + fs::create_dir_all(&rollback_path).await.map_err(to_file_error)?; + fs::write(rollback_path.join(DELETE_MARKER_ROLLBACK_FILE), []) + .await + .map_err(to_file_error)?; + } + if let Err(err) = self.write_metadata("", volume, path, fi).await { + if let Some(rollback_dir) = rollback_dir + && let Err(restore_err) = restore_delete_rollback(object_dir, xl_path, rollback_dir, &self.publication_root).await + { + warn!( + event = EVENT_DISK_LOCAL_DELETE_ROLLBACK_FAILED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_DISK_LOCAL, + result = "failed", + volume, + path, + rollback_dir = %rollback_dir, + error = ?restore_err, + "Disk local delete rollback failed" + ); + } + return Err(err); + } + Ok(()) + } + async fn delete_versions_internal(&self, volume: &str, path: &str, fis: &[FileInfo], opts: &DeleteOptions) -> Result<()> { let volume_dir = self.io_get_bucket_path(volume)?; let xlpath = self.io_get_object_path(volume, format!("{path}/{STORAGE_FORMAT_FILE}").as_str())?; @@ -6123,7 +6161,20 @@ impl LocalDisk { return restore_metadata_backup(object_dir, &xlpath, rollback_dir, &self.publication_root).await; } - let (data, _) = self.read_all_data_with_dmtime(volume, volume_dir.as_path(), &xlpath).await?; + let (data, _) = match self.read_all_data_with_dmtime(volume, volume_dir.as_path(), &xlpath).await { + Ok(data) => data, + Err(DiskError::FileNotFound) => { + // `deleted` alone can be an explicit marker purge; only + // `mark_deleted` may create metadata that was not present. + let Some(delete_marker) = fis.iter().find(|fi| fi.deleted && fi.mark_deleted).cloned() else { + return Err(DiskError::FileNotFound); + }; + return self + .write_missing_delete_marker(volume, path, delete_marker, object_dir, &xlpath, opts.old_data_dir) + .await; + } + Err(err) => return Err(err), + }; if data.is_empty() { return Err(DiskError::FileNotFound); @@ -10422,29 +10473,9 @@ impl DiskAPI for LocalDisk { } if fi.deleted && force_del_marker { - if let Some(rollback_dir) = rollback_dir { - let rollback_path = file_path.join(rollback_dir.to_string()); - fs::create_dir_all(&rollback_path).await.map_err(to_file_error)?; - fs::write(rollback_path.join(DELETE_MARKER_ROLLBACK_FILE), []) - .await - .map_err(to_file_error)?; - } - if let Err(err) = self.write_metadata("", volume, path, fi).await { - if let Some(rollback_dir) = rollback_dir - && let Err(restore_err) = - restore_delete_rollback(file_path.as_path(), &xl_path, rollback_dir, &self.publication_root).await - { - warn!( - volume, - path, - rollback_dir = %rollback_dir, - error = ?restore_err, - "failed to restore metadata after delete marker commit error" - ); - } - return Err(err); - } - return Ok(()); + return self + .write_missing_delete_marker(volume, path, fi, file_path.as_path(), &xl_path, rollback_dir) + .await; } return if fi.version_id.is_some() { diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index b39240996..f0e454880 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -4588,11 +4588,11 @@ fn should_preserve_delete_replication_state(opts: &ObjectOptions) -> bool { } fn should_force_delete_marker_for_missing_version(opts: &ObjectOptions) -> bool { - opts.delete_marker || (opts.versioned && opts.version_id.is_none() && !opts.data_movement) + opts.delete_marker || ((opts.versioned || opts.version_suspended) && opts.version_id.is_none() && !opts.data_movement) } fn resolve_delete_version_state(opts: &ObjectOptions, goi: &ObjectInfo, version_found: bool) -> (bool, bool) { - let mut mark_delete = goi.version_id.is_some() || (opts.versioned && opts.version_id.is_none()); + let mut mark_delete = goi.version_id.is_some() || ((opts.versioned || opts.version_suspended) && opts.version_id.is_none()); let mut delete_marker = opts.versioned; if opts.version_id.is_some() { diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index b57dbb7a8..eedd0b4b5 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -5900,6 +5900,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { if dobj.version_id.is_none() && (version_suspended || versioned) { vr.mod_time = Some(OffsetDateTime::now_utc()); vr.deleted = true; + vr.mark_deleted = true; if versioned { vr.version_id = Some(Uuid::new_v4()); } diff --git a/crates/ecstore/src/store/multipart.rs b/crates/ecstore/src/store/multipart.rs index 505938cb7..9a855c220 100644 --- a/crates/ecstore/src/store/multipart.rs +++ b/crates/ecstore/src/store/multipart.rs @@ -416,7 +416,8 @@ impl ECStore { 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); + self.apply_decommission_target_mutation_fence(0, object, &mut opts, mutation_fence) + .await; return self.pools[0] .new_multipart_upload(bucket, object, &opts) .await @@ -432,7 +433,8 @@ impl ECStore { opts.version_id.clone().unwrap_or_default(), )); } - self.apply_decommission_target_mutation_fence(idx, object, &mut opts, mutation_fence); + self.apply_decommission_target_mutation_fence(idx, object, &mut opts, mutation_fence) + .await; let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?; return Ok((res, idx, opts.expected_bucket_incarnation_id)); } @@ -456,7 +458,8 @@ impl ECStore { .await?; if !res.uploads.is_empty() { - self.apply_decommission_target_mutation_fence(idx, object, &mut opts, mutation_fence); + self.apply_decommission_target_mutation_fence(idx, object, &mut opts, mutation_fence) + .await; let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?; return Ok((res, idx, opts.expected_bucket_incarnation_id)); } @@ -470,7 +473,8 @@ impl ECStore { )); } - self.apply_decommission_target_mutation_fence(idx, object, &mut opts, mutation_fence); + self.apply_decommission_target_mutation_fence(idx, object, &mut opts, mutation_fence) + .await; let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?; Ok((res, idx, opts.expected_bucket_incarnation_id)) } @@ -744,7 +748,8 @@ 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); + self.apply_decommission_target_mutation_fence(target_pool_idx, object, &mut opts, mutation_fence) + .await; #[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 e0faa82bd..3f0c6b712 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -916,7 +916,7 @@ fn resolve_latest_object_access( } fn should_create_delete_marker_for_missing_object(opts: &ObjectOptions) -> bool { - opts.versioned && opts.version_id.is_none() && !opts.delete_marker && !opts.data_movement + (opts.versioned || opts.version_suspended) && opts.version_id.is_none() && !opts.delete_marker && !opts.data_movement } #[cfg(test)] @@ -1932,7 +1932,7 @@ impl ECStore { Ok(guard) } - pub(super) fn apply_decommission_target_mutation_fence( + pub(super) async fn apply_decommission_target_mutation_fence( &self, target_pool_idx: usize, object: &str, @@ -1944,11 +1944,16 @@ impl ECStore { }; mutation_fence.add_namespace_lock_fence(opts); + let distributed = self.ctx.is_dist_erasure().await; 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)); + // Local locks share one manager without set-qualified resource keys. + // Distributed locks overlap when both sets use the same client domain. + opts.no_lock = matches!( + (fixed_set, target_set), + (Some(fixed), Some(target)) + if !distributed || same_distributed_lock_domain(&fixed.lockers, &target.lockers) + ); } pub(crate) async fn acquire_decommission_source_cleanup_fence( @@ -1967,10 +1972,9 @@ impl ECStore { let test_namespace_lock_fence = decommission_mutation_fence_for_test(bucket, object, DecommissionMutationFenceTestPhase::SourceCleanup); let object = encode_dir_object(object); + let distributed = self.ctx.is_dist_erasure().await; let fixed_set = Arc::clone(&self.pools[0].disk_set[0]); - // SetDisks namespaces include pool/set identity, so only the canonical - // fixed set is covered by the fixed mutation guard. - let source_lock_covered = std::ptr::eq(fixed_set.as_ref(), source_set); + let source_lock_covered = !distributed || same_distributed_lock_domain(&fixed_set.lockers, &source_set.lockers); // Lock order: fixed store mutation domain first; source cleanup takes its // hashed source-domain lock second only when this guard does not cover it. let guard = self @@ -2451,7 +2455,8 @@ impl ECStore { 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); + self.apply_decommission_target_mutation_fence(idx, object.as_str(), &mut opts, mutation_fence) + .await; let result = self.pools[idx] .put_object_with_old_current_size(bucket, &object, data, &opts) .await