mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 12:26:37 +00:00
fix(ecstore): unblock decommission delete fences
This commit is contained in:
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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<Uuid>,
|
||||
) -> 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() {
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user