mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 20:59:05 +00:00
fix(ecstore): unblock decommission delete fences
This commit is contained in:
@@ -4484,15 +4484,20 @@ impl ECStore {
|
|||||||
) -> Result<()> {
|
) -> Result<()> {
|
||||||
warn!("decommission_object: start {} {}", &bucket, &rd.object_info.name);
|
warn!("decommission_object: start {} {}", &bucket, &rd.object_info.name);
|
||||||
let object_name = rd.object_info.name.clone();
|
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,
|
self,
|
||||||
pool_idx,
|
pool_idx,
|
||||||
bucket.clone(),
|
bucket.clone(),
|
||||||
rd,
|
rd,
|
||||||
expected_bucket_incarnation_id,
|
expected_bucket_incarnation_id,
|
||||||
"decommission_object",
|
"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() {
|
if result.is_ok() {
|
||||||
warn!("decommission_object: migrated {} {}", &bucket, &object_name);
|
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")]
|
#[cfg(target_os = "linux")]
|
||||||
const EVENT_DISK_LOCAL_URING_LATCH_OFF: &str = "disk_local_uring_latch_off";
|
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_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_CHECK_PARTS: &str = "disk_local_check_parts";
|
||||||
const EVENT_DISK_LOCAL_ACCESS_FAILED: &str = "disk_local_access_failed";
|
const EVENT_DISK_LOCAL_ACCESS_FAILED: &str = "disk_local_access_failed";
|
||||||
const EVENT_DISK_LOCAL_VOLUME_SETUP_FAILED: &str = "disk_local_volume_setup_failed";
|
const EVENT_DISK_LOCAL_VOLUME_SETUP_FAILED: &str = "disk_local_volume_setup_failed";
|
||||||
@@ -6106,6 +6107,43 @@ impl LocalDisk {
|
|||||||
Ok((bytes, modtime))
|
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<()> {
|
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 volume_dir = self.io_get_bucket_path(volume)?;
|
||||||
let xlpath = self.io_get_object_path(volume, format!("{path}/{STORAGE_FORMAT_FILE}").as_str())?;
|
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;
|
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() {
|
if data.is_empty() {
|
||||||
return Err(DiskError::FileNotFound);
|
return Err(DiskError::FileNotFound);
|
||||||
@@ -10422,29 +10473,9 @@ impl DiskAPI for LocalDisk {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if fi.deleted && force_del_marker {
|
if fi.deleted && force_del_marker {
|
||||||
if let Some(rollback_dir) = rollback_dir {
|
return self
|
||||||
let rollback_path = file_path.join(rollback_dir.to_string());
|
.write_missing_delete_marker(volume, path, fi, file_path.as_path(), &xl_path, rollback_dir)
|
||||||
fs::create_dir_all(&rollback_path).await.map_err(to_file_error)?;
|
.await;
|
||||||
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 if fi.version_id.is_some() {
|
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 {
|
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) {
|
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;
|
let mut delete_marker = opts.versioned;
|
||||||
|
|
||||||
if opts.version_id.is_some() {
|
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) {
|
if dobj.version_id.is_none() && (version_suspended || versioned) {
|
||||||
vr.mod_time = Some(OffsetDateTime::now_utc());
|
vr.mod_time = Some(OffsetDateTime::now_utc());
|
||||||
vr.deleted = true;
|
vr.deleted = true;
|
||||||
|
vr.mark_deleted = true;
|
||||||
if versioned {
|
if versioned {
|
||||||
vr.version_id = Some(Uuid::new_v4());
|
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?;
|
let (mut opts, _bucket_lifecycle_guard) = self.guard_multipart_bucket_incarnation(bucket, opts).await?;
|
||||||
|
|
||||||
if self.single_pool() {
|
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]
|
return self.pools[0]
|
||||||
.new_multipart_upload(bucket, object, &opts)
|
.new_multipart_upload(bucket, object, &opts)
|
||||||
.await
|
.await
|
||||||
@@ -432,7 +433,8 @@ impl ECStore {
|
|||||||
opts.version_id.clone().unwrap_or_default(),
|
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?;
|
let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?;
|
||||||
return Ok((res, idx, opts.expected_bucket_incarnation_id));
|
return Ok((res, idx, opts.expected_bucket_incarnation_id));
|
||||||
}
|
}
|
||||||
@@ -456,7 +458,8 @@ impl ECStore {
|
|||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
if !res.uploads.is_empty() {
|
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?;
|
let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?;
|
||||||
return Ok((res, idx, opts.expected_bucket_incarnation_id));
|
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?;
|
let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?;
|
||||||
Ok((res, idx, opts.expected_bucket_incarnation_id))
|
Ok((res, idx, opts.expected_bucket_incarnation_id))
|
||||||
}
|
}
|
||||||
@@ -744,7 +748,8 @@ impl ECStore {
|
|||||||
snapshot.add_lock_fences(&mut opts);
|
snapshot.add_lock_fences(&mut opts);
|
||||||
opts.object_lock_config_snapshot = Some(snapshot);
|
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)]
|
#[cfg(test)]
|
||||||
pause_data_movement_multipart_before_selected_completion(bucket).await;
|
pause_data_movement_multipart_before_selected_completion(bucket).await;
|
||||||
let pool = self
|
let pool = self
|
||||||
|
|||||||
@@ -916,7 +916,7 @@ fn resolve_latest_object_access(
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn should_create_delete_marker_for_missing_object(opts: &ObjectOptions) -> bool {
|
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)]
|
#[cfg(test)]
|
||||||
@@ -1932,7 +1932,7 @@ impl ECStore {
|
|||||||
Ok(guard)
|
Ok(guard)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(super) fn apply_decommission_target_mutation_fence(
|
pub(super) async fn apply_decommission_target_mutation_fence(
|
||||||
&self,
|
&self,
|
||||||
target_pool_idx: usize,
|
target_pool_idx: usize,
|
||||||
object: &str,
|
object: &str,
|
||||||
@@ -1944,11 +1944,16 @@ impl ECStore {
|
|||||||
};
|
};
|
||||||
|
|
||||||
mutation_fence.add_namespace_lock_fence(opts);
|
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 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));
|
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
|
// Local locks share one manager without set-qualified resource keys.
|
||||||
// when both names resolve to the exact same SetDisks namespace.
|
// Distributed locks overlap when both sets use the same client domain.
|
||||||
opts.no_lock = matches!((fixed_set, target_set), (Some(fixed), Some(target)) if Arc::ptr_eq(fixed, &target));
|
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(
|
pub(crate) async fn acquire_decommission_source_cleanup_fence(
|
||||||
@@ -1967,10 +1972,9 @@ impl ECStore {
|
|||||||
let test_namespace_lock_fence =
|
let test_namespace_lock_fence =
|
||||||
decommission_mutation_fence_for_test(bucket, object, DecommissionMutationFenceTestPhase::SourceCleanup);
|
decommission_mutation_fence_for_test(bucket, object, DecommissionMutationFenceTestPhase::SourceCleanup);
|
||||||
let object = encode_dir_object(object);
|
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]);
|
let fixed_set = Arc::clone(&self.pools[0].disk_set[0]);
|
||||||
// SetDisks namespaces include pool/set identity, so only the canonical
|
let source_lock_covered = !distributed || same_distributed_lock_domain(&fixed_set.lockers, &source_set.lockers);
|
||||||
// fixed set is covered by the fixed mutation guard.
|
|
||||||
let source_lock_covered = std::ptr::eq(fixed_set.as_ref(), source_set);
|
|
||||||
// Lock order: fixed store mutation domain first; source cleanup takes its
|
// Lock order: fixed store mutation domain first; source cleanup takes its
|
||||||
// hashed source-domain lock second only when this guard does not cover it.
|
// hashed source-domain lock second only when this guard does not cover it.
|
||||||
let guard = self
|
let guard = self
|
||||||
@@ -2451,7 +2455,8 @@ impl ECStore {
|
|||||||
let idx = self
|
let idx = self
|
||||||
.select_put_object_pool_idx(bucket, object.as_str(), data.size(), &opts)
|
.select_put_object_pool_idx(bucket, object.as_str(), data.size(), &opts)
|
||||||
.await?;
|
.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]
|
let result = self.pools[idx]
|
||||||
.put_object_with_old_current_size(bucket, &object, data, &opts)
|
.put_object_with_old_current_size(bucket, &object, data, &opts)
|
||||||
.await
|
.await
|
||||||
|
|||||||
Reference in New Issue
Block a user