Compare commits

..

18 Commits

Author SHA1 Message Date
overtrue 7b7dcfaaf8 fix(rebalance): preserve access-denied delete errors 2026-08-22 19:20:11 +08:00
overtrue 657835f12c fix(ecstore): satisfy delete fence lint checks 2026-08-22 17:05:47 +08:00
overtrue ca8d4c2ea7 test(ecstore): align decommission fence barriers 2026-08-22 17:05:47 +08:00
overtrue ca999e42f5 fix(ecstore): match decommission lock backend domain 2026-08-22 17:05:46 +08:00
overtrue 65a0eabfbf fix(ecstore): preserve distributed decommission set locks 2026-08-22 17:05:46 +08:00
overtrue 84f6b10186 fix(ecstore): unblock decommission delete fences 2026-08-22 17:05:46 +08:00
overtrue 7202c9937e test(ecstore): fix decommission fence fixtures 2026-08-22 17:05:46 +08:00
overtrue 5677201fc6 fix(ecstore): annotate batch delete fallback 2026-08-22 17:05:46 +08:00
overtrue 77604a5908 fix(ecstore): fence decommission commit loss 2026-08-22 17:05:46 +08:00
overtrue f1bf990588 fix(ecstore): reuse fixed fence for reverse decommission 2026-08-22 17:05:46 +08:00
overtrue 4a07a0a917 test(ecstore): finish decommission delete fence scenario 2026-08-22 17:05:46 +08:00
overtrue 906d982ced test(ecstore): exercise decommission delete fences 2026-08-22 17:05:46 +08:00
overtrue 971b84c4b5 fix(ecstore): retain source-set lock during cleanup 2026-08-22 17:05:46 +08:00
overtrue 631a30d6e2 fix(ecstore): preserve batch delete pool errors 2026-08-22 17:05:46 +08:00
overtrue cd06111dcb fix(ecstore): route batch delete markers to active pools 2026-08-22 17:05:46 +08:00
overtrue d139c8a8c3 fix(ecstore): preserve delete markers during source cleanup 2026-08-22 17:05:46 +08:00
overtrue 3793552885 fix(ecstore): preserve decommission target write locks 2026-08-22 17:05:46 +08:00
overtrue fdd193ff36 fix(ecstore): fence deletes against decommission commits 2026-08-22 17:05:46 +08:00
22 changed files with 2711 additions and 899 deletions
+17 -84
View File
@@ -50,10 +50,10 @@ use rustfs_protos::evict_failed_connection;
use rustfs_protos::proto_gen::node_service::RenamePartRequest;
use rustfs_protos::proto_gen::node_service::{
BatchReadVersionRequest, BatchReadVersionResponse, CheckPartsRequest, DeletePathsRequest, DeleteRequest,
DeleteVersionRequest, DeleteVersionsRequest, DeleteVersionsResponse, DeleteVolumeRequest, DiskInfoRequest, ListDirRequest,
ListVolumesRequest, MakeVolumeRequest, MakeVolumesRequest, PreparePartTransactionRequest, ReadAllRequest,
ReadMetadataRequest, ReadMultipleRequest, ReadMultipleResponse, ReadPartsRequest, ReadVersionRequest, ReadXlRequest,
RenameDataRequest, RenameFileRequest, SettlePartTransactionRequest, SnapshotLeaseReleaseRequest, SnapshotLeaseRenewRequest,
DeleteVersionRequest, DeleteVersionsRequest, DeleteVolumeRequest, DiskInfoRequest, ListDirRequest, ListVolumesRequest,
MakeVolumeRequest, MakeVolumesRequest, PreparePartTransactionRequest, ReadAllRequest, ReadMetadataRequest,
ReadMultipleRequest, ReadMultipleResponse, ReadPartsRequest, ReadVersionRequest, ReadXlRequest, RenameDataRequest,
RenameFileRequest, SettlePartTransactionRequest, SnapshotLeaseReleaseRequest, SnapshotLeaseRenewRequest,
SnapshotLeaseRequest, SnapshotLeaseResponse, StatVolumeRequest, UpdateMetadataRequest, VerifyFileRequest, WriteAllRequest,
WriteMetadataRequest, node_service_client::NodeServiceClient,
};
@@ -112,28 +112,6 @@ const EVENT_REMOTE_DISK_RPC: &str = "remote_disk_rpc";
const SNAPSHOT_LEASE_PROTOCOL_VERSION: u32 = 1;
pub const REMOTE_SNAPSHOT_LEASE_TTL: Duration = Duration::from_secs(60);
fn decode_delete_versions_errors(response: DeleteVersionsResponse, expected_len: usize) -> Vec<Option<Error>> {
if !response.item_errors.is_empty() {
if response.item_errors.len() != expected_len {
return vec![Some(Error::other("malformed delete_versions item errors")); expected_len];
}
return response
.item_errors
.into_iter()
.map(|error| (error.code != 0).then(|| error.into()))
.collect();
}
if response.errors.len() != expected_len {
return vec![Some(Error::other("malformed delete_versions errors")); expected_len];
}
response
.errors
.into_iter()
.map(|error| (!error.is_empty()).then(|| Error::other(error)))
.collect()
}
fn snapshot_lease_token_from_response(response: SnapshotLeaseResponse) -> Result<SnapshotLeaseToken> {
if !response.success {
return Err(response.error.unwrap_or_default().into());
@@ -2428,6 +2406,8 @@ impl DiskAPI for RemoteDisk {
return errors;
}
// TODO(backlog): replace string errors with typed `StorageError` variants
let result = self
.execute_with_timeout(
|| async {
@@ -2459,7 +2439,17 @@ impl DiskAPI for RemoteDisk {
}
return errors;
}
decode_delete_versions_errors(response, versions.len())
response
.errors
.iter()
.map(|error| {
if error.is_empty() {
None
} else {
Some(Error::other(error.to_string()))
}
})
.collect()
}
#[tracing::instrument(level = "trace", skip_all)]
@@ -3770,63 +3760,6 @@ mod tests {
static INIT: Once = Once::new();
#[test]
fn delete_versions_response_preserves_typed_item_errors() {
let errors = decode_delete_versions_errors(
DeleteVersionsResponse {
success: true,
errors: vec!["file not found".to_string(), String::new()],
error: None,
item_errors: vec![
rustfs_protos::proto_gen::node_service::Error {
code: DiskError::FileNotFound.to_u32(),
error_info: "file not found".to_string(),
},
rustfs_protos::proto_gen::node_service::Error::default(),
],
},
2,
);
assert!(matches!(errors.as_slice(), [Some(DiskError::FileNotFound), None]));
}
#[test]
fn delete_versions_response_accepts_legacy_string_errors() {
let errors = decode_delete_versions_errors(
DeleteVersionsResponse {
success: true,
errors: vec!["legacy error".to_string(), String::new()],
error: None,
item_errors: Vec::new(),
},
2,
);
assert_eq!(errors.len(), 2);
assert_eq!(errors[0].as_ref().map(ToString::to_string).as_deref(), Some("io error legacy error"));
assert!(errors[1].is_none());
}
#[test]
fn delete_versions_response_rejects_misaligned_item_errors() {
let errors = decode_delete_versions_errors(
DeleteVersionsResponse {
success: true,
errors: vec!["file not found".to_string()],
error: None,
item_errors: vec![rustfs_protos::proto_gen::node_service::Error {
code: DiskError::FileNotFound.to_u32(),
error_info: "file not found".to_string(),
}],
},
2,
);
assert_eq!(errors.len(), 2);
assert!(errors.iter().all(Option::is_some));
}
#[test]
fn disk_mutation_digest_marks_rolling_compatibility() {
let mut request = Request::new(());
+28 -3
View File
@@ -3416,6 +3416,9 @@ impl ECStore {
)
.await?;
let source_cleanup_mutation_fence = self
.acquire_decommission_source_cleanup_fence(bucket.as_str(), entry.name.as_str(), set.as_ref())
.await?;
let cleanup_result = data_movement::cleanup_source_entry_if_unchanged(
set.clone(),
bucket.as_str(),
@@ -3427,6 +3430,7 @@ impl ECStore {
lifecycle_guard: bucket_incarnation_fence
.as_ref()
.and_then(|guard| guard.namespace_lock_guard()),
object_mutation_fence: Some(&source_cleanup_mutation_fence),
},
"decommission",
)
@@ -3528,6 +3532,22 @@ impl ECStore {
Ok(())
}
#[cfg(test)]
pub(crate) async fn decommission_entry_for_test(
self: &Arc<Self>,
idx: usize,
entry: MetaCacheEntry,
bucket: String,
set: Arc<SetDisks>,
) -> Result<()> {
let worker_permit = Arc::new(Semaphore::new(1))
.acquire_owned()
.await
.map_err(|err| Error::other(format!("decommission test worker permit acquire failed: {err}")))?;
self.decommission_entry(CancellationToken::new(), idx, entry, bucket, set, worker_permit, None, None, None, None)
.await
}
#[tracing::instrument(skip(self, rx))]
async fn decommission_pool(
self: &Arc<Self>,
@@ -4464,15 +4484,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_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);
}
+161 -42
View File
@@ -26,7 +26,7 @@ use crate::storage_api_contracts::{
namespace::NamespaceLocking as _,
object::{HTTPPreconditions, ObjectOperations as _},
};
use crate::store::ECStore;
use crate::store::{ECStore, ObjectLockDiagGuard, SourceCleanupMutationFence};
use bytes::Bytes;
use rustfs_filemeta::{FileInfo, FileInfoVersions, ObjectPartInfo};
use rustfs_rio::{EtagResolvable, HashReader, HashReaderDetector, Index, TryGetIndex};
@@ -856,7 +856,6 @@ fn is_equivalent_data_movement_object(source: &ObjectInfo, target: &ObjectInfo)
fn is_superseding_unversioned_data_movement_object(source: &ObjectInfo, target: &ObjectInfo) -> bool {
is_unversioned_data_movement_object(source)
&& is_unversioned_data_movement_object(target)
&& !target.delete_marker
&& source
.mod_time
.zip(target.mod_time)
@@ -1028,6 +1027,7 @@ pub(crate) enum SourceCleanupError {
pub(crate) struct SourceCleanupBucketFence<'a> {
pub(crate) expected_incarnation_id: Option<uuid::Uuid>,
pub(crate) lifecycle_guard: Option<&'a rustfs_lock::NamespaceLockGuard>,
pub(crate) object_mutation_fence: Option<&'a SourceCleanupMutationFence>,
}
fn ensure_source_cleanup_versions_match(
@@ -1065,7 +1065,9 @@ pub(crate) async fn ensure_source_cleanup_versions_unchanged(
struct SourceCleanupDeleteBarrierState {
bucket: String,
object: String,
fence_pending: tokio::sync::Notify,
arrived: tokio::sync::Notify,
is_paused: AtomicBool,
release: tokio::sync::Notify,
}
@@ -1079,7 +1081,7 @@ pub(crate) struct SourceCleanupDeleteBarrier {
}
#[cfg(test)]
static SOURCE_CLEANUP_DELETE_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option<Arc<SourceCleanupDeleteBarrierState>>>> =
static SOURCE_CLEANUP_DELETE_BARRIERS: std::sync::OnceLock<std::sync::Mutex<Vec<Arc<SourceCleanupDeleteBarrierState>>>> =
std::sync::OnceLock::new();
#[cfg(test)]
@@ -1092,15 +1094,22 @@ impl SourceCleanupDeleteBarrier {
let state = Arc::new(SourceCleanupDeleteBarrierState {
bucket: bucket.to_string(),
object: object.to_string(),
fence_pending: tokio::sync::Notify::new(),
arrived: tokio::sync::Notify::new(),
is_paused: AtomicBool::new(false),
release: tokio::sync::Notify::new(),
});
let mut slot = SOURCE_CLEANUP_DELETE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
let mut barriers = SOURCE_CLEANUP_DELETE_BARRIERS
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("source cleanup delete barrier mutex should not poison");
assert!(slot.is_none(), "source cleanup delete barrier must be unique");
*slot = Some(Arc::clone(&state));
assert!(
!barriers
.iter()
.any(|barrier| barrier.bucket == bucket && barrier.object == object),
"source cleanup delete barrier must be unique per object"
);
barriers.push(Arc::clone(&state));
Self { state }
}
@@ -1110,35 +1119,58 @@ impl SourceCleanupDeleteBarrier {
.expect("source cleanup should reach the pre-delete barrier");
}
pub(crate) async fn wait_until_fence_pending(&self) {
tokio::time::timeout(StdDuration::from_secs(30), self.state.fence_pending.notified())
.await
.expect("source cleanup should attempt the fixed mutation fence");
}
pub(crate) fn is_paused(&self) -> bool {
self.state.is_paused.load(Ordering::Acquire)
}
pub(crate) fn release(&self) {
self.state.release.notify_one();
}
}
#[cfg(test)]
pub(crate) fn notify_source_cleanup_mutation_fence_pending(bucket: &str, object: &str) {
let barrier = SOURCE_CLEANUP_DELETE_BARRIERS
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("source cleanup delete barrier mutex should not poison")
.iter()
.find(|barrier| barrier.bucket == bucket && barrier.object == object)
.cloned();
if let Some(barrier) = barrier {
barrier.fence_pending.notify_one();
}
}
#[cfg(test)]
impl Drop for SourceCleanupDeleteBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
let mut slot = SOURCE_CLEANUP_DELETE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
let mut barriers = SOURCE_CLEANUP_DELETE_BARRIERS
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("source cleanup delete barrier mutex should not poison");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
barriers.retain(|state| !Arc::ptr_eq(state, &self.state));
}
}
#[cfg(test)]
async fn pause_source_cleanup_before_delete(bucket: &str, object: &str) {
let barrier = SOURCE_CLEANUP_DELETE_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
let barrier = SOURCE_CLEANUP_DELETE_BARRIERS
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("source cleanup delete barrier mutex should not poison")
.as_ref()
.filter(|barrier| barrier.bucket == bucket && barrier.object == object)
.iter()
.find(|barrier| barrier.bucket == bucket && barrier.object == object)
.cloned();
if let Some(barrier) = barrier {
barrier.is_paused.store(true, Ordering::Release);
barrier.arrived.notify_one();
barrier.release.notified().await;
}
@@ -1154,11 +1186,20 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
op_label: &str,
) -> std::result::Result<ObjectInfo, SourceCleanupError> {
let cleanup_key = encode_dir_object(object);
let ns_lock = set.new_ns_lock(bucket, cleanup_key.as_str()).await?;
let _guard = ns_lock
.get_write_lock(get_lock_acquire_timeout())
.await
.map_err(Error::from)?;
let source_guard = if bucket_fence
.object_mutation_fence
.is_some_and(SourceCleanupMutationFence::source_lock_covered)
{
None
} else {
let ns_lock = set.new_ns_lock(bucket, cleanup_key.as_str()).await?;
Some(
ns_lock
.get_write_lock(get_lock_acquire_timeout())
.await
.map_err(Error::from)?,
)
};
if bucket_fence
.lifecycle_guard
@@ -1168,6 +1209,14 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
"{op_label}: bucket incarnation fence was lost before source cleanup"
))));
}
if bucket_fence
.object_mutation_fence
.is_some_and(SourceCleanupMutationFence::is_lock_lost)
{
return Err(SourceCleanupError::Storage(Error::other(format!(
"{op_label}: object mutation fence was lost before source cleanup"
))));
}
ensure_source_cleanup_versions_unchanged(set.clone(), bucket, object, expected, allowed_missing, op_label).await?;
@@ -1182,7 +1231,12 @@ pub(crate) async fn cleanup_source_entry_if_unchanged(
expected_bucket_incarnation_id: bucket_fence.expected_incarnation_id,
..Default::default()
};
opts.add_namespace_lock_guard(&_guard);
if let Some(source_guard) = source_guard.as_ref() {
opts.add_namespace_lock_guard(source_guard);
}
if let Some(object_mutation_fence) = bucket_fence.object_mutation_fence {
object_mutation_fence.add_namespace_lock_fence(&mut opts);
}
if let Some(bucket_lifecycle_guard) = bucket_fence.lifecycle_guard {
opts.add_bucket_lifecycle_lock_guard(bucket_lifecycle_guard);
}
@@ -1330,6 +1384,37 @@ fn data_movement_part_upload_failure_stage(err: &Error) -> &'static str {
}
}
pub(crate) async fn migrate_decommission_object(
store: Arc<ECStore>,
pool_idx: usize,
bucket: String,
rd: GetObjectReader,
source_bucket_incarnation_id: Option<uuid::Uuid>,
op_label: &str,
) -> Result<()> {
let source = rd.object_info.clone();
let _mutation_fence = store
.acquire_decommission_object_mutation_fence(&bucket, &source.name)
.await?;
let current = find_data_movement_target_info(store.as_ref(), pool_idx, &bucket, &source)
.await?
.ok_or(Error::FileNotFound)?;
if !is_equivalent_data_movement_object_identity(&source, &current, true, false) {
return Err(Error::FileNotFound);
}
migrate_object_inner(
store,
pool_idx,
bucket,
rd,
source_bucket_incarnation_id,
op_label,
Some(&_mutation_fence),
)
.await
}
pub(crate) async fn migrate_object(
store: Arc<ECStore>,
pool_idx: usize,
@@ -1337,6 +1422,18 @@ pub(crate) async fn migrate_object(
rd: GetObjectReader,
source_bucket_incarnation_id: Option<uuid::Uuid>,
op_label: &str,
) -> Result<()> {
migrate_object_inner(store, pool_idx, bucket, rd, source_bucket_incarnation_id, op_label, None).await
}
async fn migrate_object_inner(
store: Arc<ECStore>,
pool_idx: usize,
bucket: String,
rd: GetObjectReader,
source_bucket_incarnation_id: Option<uuid::Uuid>,
op_label: &str,
mutation_fence: Option<&ObjectLockDiagGuard>,
) -> Result<()> {
let object_info = rd.object_info.clone();
let has_part_checksums = object_info
@@ -1350,7 +1447,7 @@ pub(crate) async fn migrate_object(
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;
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,
@@ -1448,7 +1545,7 @@ pub(crate) async fn migrate_object(
if let Err(err) = store
.clone()
.complete_multipart_upload_for_data_movement(
target_pool_idx,
(target_pool_idx, mutation_fence),
&bucket,
&object_info.name,
&res.upload_id,
@@ -1609,7 +1706,7 @@ pub(crate) async fn migrate_object(
let mut put_opts = data_movement_put_object_opts(&object_info, pool_idx);
put_opts.expected_bucket_incarnation_id = source_bucket_incarnation_id;
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 {
@@ -3541,25 +3638,47 @@ mod tests {
}
#[test]
fn test_precondition_conflict_rejects_newer_delete_marker() {
let source = ObjectInfo {
size: 128,
etag: Some("etag-source".to_string()),
mod_time: Some(OffsetDateTime::UNIX_EPOCH),
..Default::default()
};
let target = ObjectInfo {
delete_marker: true,
etag: None,
mod_time: OffsetDateTime::UNIX_EPOCH.checked_add(time::Duration::SECOND),
..source.clone()
};
fn test_precondition_conflict_accepts_only_newer_null_delete_marker() {
for version_id in [None, Some(Uuid::nil())] {
let source = ObjectInfo {
version_id,
size: 128,
etag: Some("etag-source".to_string()),
mod_time: Some(OffsetDateTime::UNIX_EPOCH),
..Default::default()
};
let target = ObjectInfo {
delete_marker: true,
etag: None,
mod_time: OffsetDateTime::UNIX_EPOCH.checked_add(time::Duration::SECOND),
..source.clone()
};
let should_resume =
resolve_data_movement_overwrite_resume_result(&Error::PreconditionFailed, Ok(Some(target)), &source, 0, 1)
.expect("delete marker conflict should be evaluated");
assert!(
resolve_data_movement_overwrite_resume_result(
&Error::PreconditionFailed,
Ok(Some(target.clone())),
&source,
0,
1,
)
.expect("newer null delete marker should be evaluated")
);
assert!(!should_resume);
let mut same_time = target.clone();
same_time.mod_time = source.mod_time;
assert!(
!resolve_data_movement_overwrite_resume_result(&Error::PreconditionFailed, Ok(Some(same_time)), &source, 0, 1,)
.expect("same-generation null delete marker should be rejected")
);
let mut versioned = target;
versioned.version_id = Some(Uuid::new_v4());
assert!(
!resolve_data_movement_overwrite_resume_result(&Error::PreconditionFailed, Ok(Some(versioned)), &source, 0, 1,)
.expect("a UUID delete marker must not erase a null source version")
);
}
}
#[test]
+55 -24
View File
@@ -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() {
+20 -10
View File
@@ -24,7 +24,7 @@ use crate::storage_api_contracts::{
pub struct NamespaceLockFence {
signals: Arc<Vec<Arc<rustfs_lock::distributed_lock::LockLostSignal>>>,
#[cfg(test)]
forced_lost: Arc<std::sync::atomic::AtomicBool>,
forced_lost: Arc<Vec<Arc<std::sync::atomic::AtomicBool>>>,
}
impl Debug for NamespaceLockFence {
@@ -40,13 +40,17 @@ impl NamespaceLockFence {
Self {
signals: Arc::default(),
#[cfg(test)]
forced_lost: Arc::new(std::sync::atomic::AtomicBool::new(false)),
forced_lost: Arc::new(vec![Arc::new(std::sync::atomic::AtomicBool::new(false))]),
}
}
pub(crate) fn is_lock_lost(&self) -> bool {
#[cfg(test)]
if self.forced_lost.load(std::sync::atomic::Ordering::Acquire) {
if self
.forced_lost
.iter()
.any(|lost| lost.load(std::sync::atomic::Ordering::Acquire))
{
return true;
}
self.signals.iter().any(|signal| signal.is_lost())
@@ -57,27 +61,26 @@ impl NamespaceLockFence {
}
fn extend(&mut self, other: &Self) {
if Arc::ptr_eq(&self.signals, &other.signals) {
return;
if !Arc::ptr_eq(&self.signals, &other.signals) {
Arc::make_mut(&mut self.signals).extend(other.signals.iter().cloned());
}
Arc::make_mut(&mut self.signals).extend(other.signals.iter().cloned());
#[cfg(test)]
if other.forced_lost.load(std::sync::atomic::Ordering::Acquire) {
self.forced_lost.store(true, std::sync::atomic::Ordering::Release);
if !Arc::ptr_eq(&self.forced_lost, &other.forced_lost) {
Arc::make_mut(&mut self.forced_lost).extend(other.forced_lost.iter().cloned());
}
}
#[cfg(test)]
pub(crate) fn lost_for_test() -> Self {
let fence = Self::new();
fence.forced_lost.store(true, std::sync::atomic::Ordering::Release);
fence.forced_lost[0].store(true, std::sync::atomic::Ordering::Release);
fence
}
#[cfg(test)]
pub(crate) fn loss_handle_for_test() -> (Self, Arc<std::sync::atomic::AtomicBool>) {
let fence = Self::new();
(fence.clone(), Arc::clone(&fence.forced_lost))
(fence.clone(), Arc::clone(&fence.forced_lost[0]))
}
}
@@ -411,6 +414,13 @@ impl ObjectOptions {
self.namespace_lock_fence.get_or_insert_with(NamespaceLockFence::new);
}
#[cfg(test)]
pub(crate) fn add_namespace_lock_fence_for_test(&mut self, fence: &NamespaceLockFence) {
self.namespace_lock_fence
.get_or_insert_with(NamespaceLockFence::new)
.extend(fence);
}
pub(crate) fn ensure_lifecycle_delete_all_journal(&mut self) {
self.lifecycle_delete_all_journal
.get_or_insert_with(|| Arc::new(parking_lot::Mutex::new(LifecycleDeleteAllJournalState::default())));
@@ -334,6 +334,7 @@ impl ECStore {
lifecycle_guard: bucket_incarnation_fence
.as_ref()
.and_then(|guard| guard.namespace_lock_guard()),
..Default::default()
},
"rebalance",
),
+25 -2
View File
@@ -735,8 +735,12 @@ pub(crate) use core::io_primitives::disk_call_counters;
mod ctx;
mod metadata;
mod ops;
#[cfg(test)]
pub(crate) use ops::multipart::NewMultipartUploadCommitObservation;
#[cfg(any(test, feature = "test-util"))]
pub use ops::multipart::{MultipartCommitBarrier, MultipartCommitPause};
#[cfg(test)]
pub(crate) use ops::object::DeleteObjectCommitBarrier;
#[cfg(feature = "test-util")]
pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier;
pub(crate) use ops::object::body_cache_plaintext_len;
@@ -3025,6 +3029,16 @@ pub struct SetDisks {
storage_class_config_override: Arc<std::sync::RwLock<Option<Arc<storageclass::Config>>>>,
}
// DistributedLock sends the raw ObjectKey to its clients; LockRegistry clones
// each endpoint's canonical Arc, so an exact Arc set identifies the lock domain.
pub(crate) fn same_distributed_lock_domain(left: &[Arc<dyn LockClient>], right: &[Arc<dyn LockClient>]) -> bool {
left.iter()
.all(|left_client| right.iter().any(|right_client| Arc::ptr_eq(left_client, right_client)))
&& right
.iter()
.all(|right_client| left.iter().any(|left_client| Arc::ptr_eq(left_client, right_client)))
}
const ERASURE_CACHE_MAX_ENTRIES: usize = 32;
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
@@ -3600,6 +3614,15 @@ impl SetDisks {
&self.ctx
}
/// Whether both sets' namespace-lock implementations cover the same object key.
pub(crate) async fn shares_namespace_lock_domain(&self, other: &Self) -> bool {
match (self.ctx.is_dist_erasure().await, other.ctx.is_dist_erasure().await) {
(false, false) => Arc::ptr_eq(&self.local_lock_manager, &other.local_lock_manager),
(true, true) => same_distributed_lock_domain(&self.lockers, &other.lockers),
_ => false,
}
}
/// The lock manager this set actually uses (test-only; Phase 5 Slice 3).
#[cfg(test)]
pub(crate) fn local_lock_manager_for_test(&self) -> &Arc<rustfs_lock::GlobalLockManager> {
@@ -4584,11 +4607,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() {
@@ -32,6 +32,8 @@ use crate::crash_inject::{self, CrashPoint};
use crate::multipart_listing::paginate_multipart_listing;
use futures::{StreamExt, stream};
use std::future::Future;
#[cfg(test)]
use std::sync::atomic::AtomicBool;
#[cfg(any(test, feature = "test-util"))]
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
@@ -65,6 +67,7 @@ impl StaleMultipartCleanupGuard {
#[cfg(any(test, feature = "test-util"))]
#[derive(Clone, Copy, PartialEq, Eq)]
pub enum MultipartCommitPause {
NewUploadBeforeLockLost,
PutPartBeforeLockAcquire,
PutPartBeforeLockLost,
PutPartAfterRename,
@@ -156,6 +159,72 @@ impl Drop for MultipartCommitBarrier {
}
}
#[cfg(test)]
struct NewMultipartUploadCommitObservationState {
bucket: String,
object: String,
committed: AtomicBool,
}
#[cfg(test)]
pub(crate) struct NewMultipartUploadCommitObservation {
state: Arc<NewMultipartUploadCommitObservationState>,
}
#[cfg(test)]
static NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION: std::sync::OnceLock<
std::sync::Mutex<Option<Arc<NewMultipartUploadCommitObservationState>>>,
> = std::sync::OnceLock::new();
#[cfg(test)]
impl NewMultipartUploadCommitObservation {
pub(crate) fn install(bucket: &str, object: &str) -> Self {
let state = Arc::new(NewMultipartUploadCommitObservationState {
bucket: bucket.to_string(),
object: object.to_string(),
committed: AtomicBool::new(false),
});
let mut slot = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("new multipart upload commit observation mutex should not poison");
assert!(slot.is_none(), "new multipart upload commit observation must be unique");
*slot = Some(Arc::clone(&state));
Self { state }
}
pub(crate) fn committed(&self) -> bool {
self.state.committed.load(Ordering::Acquire)
}
}
#[cfg(test)]
impl Drop for NewMultipartUploadCommitObservation {
fn drop(&mut self) {
let mut slot = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("new multipart upload commit observation mutex should not poison");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[cfg(test)]
fn observe_new_multipart_upload_commit(bucket: &str, object: &str) {
let state = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("new multipart upload commit observation mutex should not poison")
.as_ref()
.filter(|state| state.bucket == bucket && state.object == object)
.cloned();
if let Some(state) = state {
state.committed.store(true, Ordering::Release);
}
}
#[cfg(any(test, feature = "test-util"))]
async fn pause_multipart_commit(bucket: &str, object: &str, pause: MultipartCommitPause) {
let barrier = {
@@ -1615,6 +1684,30 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
let upload_path = Self::get_multipart_upload_dir(bucket, object, upload_uuid.as_str(), opts.data_movement);
#[cfg(any(test, feature = "test-util"))]
pause_multipart_commit(bucket, object, MultipartCommitPause::NewUploadBeforeLockLost).await;
if _object_lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost()) {
return Err(StorageError::NamespaceLockQuorumUnavailable {
mode: "new_multipart_upload_commit",
bucket: bucket.to_string(),
object: object.to_string(),
required: 1,
achieved: 0,
});
}
if opts
.namespace_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
{
return Err(StorageError::NamespaceLockQuorumUnavailable {
mode: "new_multipart_upload_outer_lock",
bucket: bucket.to_string(),
object: object.to_string(),
required: 1,
achieved: 0,
});
}
ensure_multipart_bucket_lifecycle_lock_held(bucket, object, opts)?;
Self::write_unique_file_info(
&shuffle_disks,
@@ -1626,6 +1719,8 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
)
.await
.map_err(|e| to_object_err(e.into(), vec![bucket, object]))?;
#[cfg(test)]
observe_new_multipart_upload_commit(bucket, object);
// evalDisks
+29 -4
View File
@@ -2483,6 +2483,7 @@ impl SetDisks {
})
.await?,
);
notify_put_object_commit_namespace_acquired(bucket, object);
}
#[cfg(not(any(test, feature = "test-util")))]
{
@@ -4630,6 +4631,7 @@ struct PutObjectCommitBarrierState {
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
namespace_pending: tokio::sync::Notify,
namespace_acquired: std::sync::atomic::AtomicBool,
}
#[cfg(any(test, feature = "test-util"))]
@@ -4651,6 +4653,7 @@ impl PutObjectCommitBarrier {
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
namespace_pending: tokio::sync::Notify::new(),
namespace_acquired: std::sync::atomic::AtomicBool::new(false),
});
let mut slot = PUT_OBJECT_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
@@ -4685,6 +4688,10 @@ impl PutObjectCommitBarrier {
.await
.expect("put object should wait for the namespace lock after leaving the commit barrier");
}
pub fn namespace_acquired(&self) -> bool {
self.state.namespace_acquired.load(std::sync::atomic::Ordering::Acquire)
}
}
#[cfg(any(test, feature = "test-util"))]
@@ -4741,6 +4748,22 @@ fn notify_put_object_commit_namespace_pending(bucket: &str, object: &str) {
}
}
#[cfg(any(test, feature = "test-util"))]
fn notify_put_object_commit_namespace_acquired(bucket: &str, object: &str) {
let barrier = PUT_OBJECT_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(Vec::new()))
.lock()
.expect("put object commit barrier mutex should not poison")
.iter()
.find(|barrier| {
barrier.bucket == bucket && barrier.object == object && barrier.pause == PutObjectCommitPause::BeforeNamespace
})
.cloned();
if let Some(barrier) = barrier {
barrier.namespace_acquired.store(true, std::sync::atomic::Ordering::Release);
}
}
#[cfg(test)]
struct DeleteObjectCommitBarrierState {
bucket: String,
@@ -4750,7 +4773,7 @@ struct DeleteObjectCommitBarrierState {
}
#[cfg(test)]
struct DeleteObjectCommitBarrier {
pub(crate) struct DeleteObjectCommitBarrier {
state: Arc<DeleteObjectCommitBarrierState>,
}
@@ -4760,7 +4783,7 @@ static DELETE_OBJECT_COMMIT_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option
#[cfg(test)]
impl DeleteObjectCommitBarrier {
fn install(bucket: &str, object: &str) -> Self {
pub(crate) fn install(bucket: &str, object: &str) -> Self {
let state = Arc::new(DeleteObjectCommitBarrierState {
bucket: bucket.to_string(),
object: object.to_string(),
@@ -4776,13 +4799,13 @@ impl DeleteObjectCommitBarrier {
Self { state }
}
async fn wait_until_paused(&self) {
pub(crate) async fn wait_until_paused(&self) {
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
.await
.expect("delete object should reach the deterministic commit barrier");
}
fn release(&self) {
pub(crate) fn release(&self) {
self.state.release.notify_one();
}
}
@@ -5877,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());
}
@@ -11636,6 +11660,7 @@ mod transition_upload_integrity_tests {
crate::data_movement::SourceCleanupBucketFence {
expected_incarnation_id: None,
lifecycle_guard: Some(&bucket_guard),
..Default::default()
},
"test_data_movement",
)
File diff suppressed because it is too large Load Diff
+1 -1
View File
@@ -151,7 +151,7 @@ pub(crate) mod init_format;
pub(crate) mod list_objects;
mod multipart;
mod object;
pub(crate) use object::ObjectLockDiagGuard;
pub(crate) use object::{ObjectLockDiagGuard, SourceCleanupMutationFence};
pub use object::{
PrepareSelectObjectSnapshotError, PreparedGetObjectReader, SelectObjectSnapshot, SelectObjectSnapshotReadError,
SnapshotConsistencyError,
+20 -9
View File
@@ -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,22 @@ 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)
.await;
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 +433,9 @@ 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)
.await;
let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?;
return Ok((res, idx, opts.expected_bucket_incarnation_id));
}
@@ -454,7 +458,9 @@ 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)
.await;
let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?;
return Ok((res, idx, opts.expected_bucket_incarnation_id));
}
}
@@ -467,7 +473,9 @@ 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)
.await;
let res = self.pools[idx].new_multipart_upload(bucket, object, &opts).await?;
Ok((res, idx, opts.expected_bucket_incarnation_id))
}
@@ -704,13 +712,14 @@ impl ECStore {
pub(crate) async fn complete_multipart_upload_for_data_movement(
self: Arc<Self>,
target_pool_idx: usize,
target: (usize, Option<&ObjectLockDiagGuard>),
bucket: &str,
object: &str,
upload_id: &str,
uploaded_parts: Vec<CompletePart>,
opts: &ObjectOptions,
) -> Result<ObjectInfo> {
let (target_pool_idx, mutation_fence) = target;
check_complete_multipart_args(bucket, object, upload_id)?;
if !opts.data_movement {
return Err(Error::other("targeted multipart completion requires data_movement options"));
@@ -739,6 +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)
.await;
#[cfg(test)]
pause_data_movement_multipart_before_selected_completion(bucket).await;
let pool = self
+790 -33
View File
@@ -32,12 +32,13 @@ use crate::bucket::metadata_sys::{
use crate::bucket::object_lock::objectlock_sys::{
check_object_lock_for_deletion_with_state, ensure_recursive_force_delete_allowed_for_state,
};
use crate::bucket::replication::ReplicationObjectBridge;
use crate::bucket::replication::{DeleteReplicationConfigSnapshot, ReplicationObjectBridge};
use crate::bucket::versioning::VersioningApi;
use crate::disk::OldCurrentSize;
use crate::object_api::{NamespaceLockFence, ObjectLockConfigSnapshot};
use crate::set_disk::{
get_lock_acquire_timeout, get_object_lock_diag_slow_acquire_threshold, get_object_lock_diag_slow_hold_threshold,
is_lock_optimization_enabled, is_object_lock_diag_enabled,
SetDisks, get_lock_acquire_timeout, get_object_lock_diag_slow_acquire_threshold, get_object_lock_diag_slow_hold_threshold,
is_lock_optimization_enabled, is_object_lock_diag_enabled, same_distributed_lock_domain,
};
use crate::storage_api_contracts::{
namespace::NamespaceLocking as _,
@@ -352,6 +353,8 @@ impl fmt::Display for ObjectLockDiagMode {
pub(crate) struct ObjectLockDiagGuard {
guard: rustfs_lock::NamespaceLockGuard,
#[cfg(test)]
test_namespace_lock_fence: Option<NamespaceLockFence>,
enabled: bool,
op: &'static str,
bucket: Option<String>,
@@ -373,6 +376,8 @@ impl ObjectLockDiagGuard {
) -> Self {
Self {
guard,
#[cfg(test)]
test_namespace_lock_fence: None,
enabled,
op,
bucket,
@@ -393,6 +398,115 @@ impl ObjectLockDiagGuard {
pub(crate) fn is_lock_lost(&self) -> bool {
self.guard.is_lock_lost()
}
pub(crate) fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) {
opts.ensure_namespace_lock_fence();
if let Some(signal) = self.lock_lost_signal() {
opts.add_namespace_lock_lost_signal(signal);
}
#[cfg(test)]
if let Some(fence) = self.test_namespace_lock_fence.as_ref() {
opts.add_namespace_lock_fence_for_test(fence);
}
}
}
#[cfg(test)]
#[derive(Clone, Copy, PartialEq, Eq)]
pub(crate) enum DecommissionMutationFenceTestPhase {
Migration,
SourceCleanup,
}
#[cfg(test)]
struct DecommissionMutationFenceLossState {
bucket: String,
object: String,
phase: DecommissionMutationFenceTestPhase,
fence: NamespaceLockFence,
loss_handle: Arc<std::sync::atomic::AtomicBool>,
}
#[cfg(test)]
pub(crate) struct DecommissionMutationFenceLossHook {
state: Arc<DecommissionMutationFenceLossState>,
}
#[cfg(test)]
static DECOMMISSION_MUTATION_FENCE_LOSS_HOOK: std::sync::OnceLock<
std::sync::Mutex<Option<Arc<DecommissionMutationFenceLossState>>>,
> = std::sync::OnceLock::new();
#[cfg(test)]
impl DecommissionMutationFenceLossHook {
pub(crate) fn install(bucket: &str, object: &str, phase: DecommissionMutationFenceTestPhase) -> Self {
let (fence, loss_handle) = NamespaceLockFence::loss_handle_for_test();
let state = Arc::new(DecommissionMutationFenceLossState {
bucket: bucket.to_string(),
object: object.to_string(),
phase,
fence,
loss_handle,
});
let mut slot = DECOMMISSION_MUTATION_FENCE_LOSS_HOOK
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission mutation fence loss hooks should not poison");
assert!(slot.is_none(), "decommission mutation fence loss hook must be unique");
*slot = Some(Arc::clone(&state));
Self { state }
}
pub(crate) fn mark_lost(&self) {
self.state.loss_handle.store(true, Ordering::Release);
}
}
#[cfg(test)]
impl Drop for DecommissionMutationFenceLossHook {
fn drop(&mut self) {
let mut slot = DECOMMISSION_MUTATION_FENCE_LOSS_HOOK
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission mutation fence loss hooks should not poison");
if slot.as_ref().is_some_and(|hook| Arc::ptr_eq(hook, &self.state)) {
*slot = None;
}
}
}
#[cfg(test)]
fn decommission_mutation_fence_for_test(
bucket: &str,
object: &str,
phase: DecommissionMutationFenceTestPhase,
) -> Option<NamespaceLockFence> {
DECOMMISSION_MUTATION_FENCE_LOSS_HOOK
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("decommission mutation fence loss hooks should not poison")
.as_ref()
.filter(|hook| hook.bucket == bucket && hook.object == object && hook.phase == phase)
.map(|hook| hook.fence.clone())
}
pub(crate) struct SourceCleanupMutationFence {
guard: ObjectLockDiagGuard,
source_lock_covered: bool,
}
impl SourceCleanupMutationFence {
pub(crate) fn source_lock_covered(&self) -> bool {
self.source_lock_covered
}
pub(crate) fn is_lock_lost(&self) -> bool {
self.guard.is_lock_lost()
}
pub(crate) fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) {
self.guard.add_namespace_lock_fence(opts);
}
}
/// Opaque write-lock guard for the RestoreObject accept path; see
@@ -410,10 +524,7 @@ impl RestoreAcceptGuard {
}
pub fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) {
opts.ensure_namespace_lock_fence();
if let Some(signal) = self.0.lock_lost_signal() {
opts.add_namespace_lock_lost_signal(signal);
}
self.0.add_namespace_lock_fence(opts);
}
}
@@ -690,16 +801,6 @@ impl SelectObjectSnapshotLockLossWake {
}
}
// LockRegistry clones its canonical client Arc for each endpoint host, so an
// exact Arc set identifies one distributed namespace-lock quorum domain.
fn same_distributed_lock_domain(left: &[Arc<dyn rustfs_lock::LockClient>], right: &[Arc<dyn rustfs_lock::LockClient>]) -> bool {
left.iter()
.all(|left_client| right.iter().any(|right_client| Arc::ptr_eq(left_client, right_client)))
&& right
.iter()
.all(|right_client| left.iter().any(|left_client| Arc::ptr_eq(left_client, right_client)))
}
impl AsyncRead for SelectObjectSnapshotReader {
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<std::io::Result<()>> {
if self.lock_loss_wake.poll_lost(cx) || self.lease.is_lost() {
@@ -805,7 +906,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)]
@@ -813,6 +914,8 @@ struct DeleteAfterObjectLockSnapshotBarrierState {
bucket: String,
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
namespace_pending: tokio::sync::Notify,
namespace_acquired: AtomicBool,
}
#[cfg(test)]
@@ -832,6 +935,8 @@ impl DeleteAfterObjectLockSnapshotBarrier {
bucket: bucket.to_string(),
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
namespace_pending: tokio::sync::Notify::new(),
namespace_acquired: AtomicBool::new(false),
});
let mut slot = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
@@ -849,6 +954,18 @@ impl DeleteAfterObjectLockSnapshotBarrier {
pub(crate) fn release(&self) {
self.state.release.notify_one();
}
pub(crate) async fn release_and_wait_until_namespace_pending(&self) {
let namespace_pending = self.state.namespace_pending.notified();
self.release();
tokio::time::timeout(Duration::from_secs(5), namespace_pending)
.await
.expect("delete should proceed to its namespace lock after leaving the snapshot barrier");
}
pub(crate) fn namespace_acquired(&self) -> bool {
self.state.namespace_acquired.load(Ordering::Acquire)
}
}
#[cfg(test)]
@@ -873,6 +990,97 @@ async fn pause_delete_after_object_lock_snapshot(bucket: &str) {
.as_ref()
.filter(|state| state.bucket == bucket)
.cloned();
if let Some(state) = state {
state.arrived.notify_one();
state.release.notified().await;
state.namespace_pending.notify_one();
}
}
#[cfg(test)]
fn notify_delete_namespace_acquired(bucket: &str) {
let state = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("delete snapshot barrier mutex should not poison")
.as_ref()
.filter(|state| state.bucket == bucket)
.cloned();
if let Some(state) = state {
state.namespace_acquired.store(true, Ordering::Release);
}
}
#[cfg(test)]
struct VersionedDeleteMarkerCommitBarrierState {
bucket: String,
object: String,
arrived: tokio::sync::Notify,
release: tokio::sync::Notify,
}
#[cfg(test)]
pub(crate) struct VersionedDeleteMarkerCommitBarrier {
state: Arc<VersionedDeleteMarkerCommitBarrierState>,
}
#[cfg(test)]
static VERSIONED_DELETE_MARKER_COMMIT_BARRIER: std::sync::OnceLock<
std::sync::Mutex<Option<Arc<VersionedDeleteMarkerCommitBarrierState>>>,
> = std::sync::OnceLock::new();
#[cfg(test)]
impl VersionedDeleteMarkerCommitBarrier {
pub(crate) fn install(bucket: &str, object: &str) -> Self {
let state = Arc::new(VersionedDeleteMarkerCommitBarrierState {
bucket: bucket.to_string(),
object: object.to_string(),
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
});
let mut slot = VERSIONED_DELETE_MARKER_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("versioned delete-marker commit barrier mutex should not poison");
assert!(slot.is_none(), "versioned delete-marker commit barrier must be unique");
*slot = Some(Arc::clone(&state));
Self { state }
}
pub(crate) async fn wait_until_paused(&self) {
tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified())
.await
.expect("versioned DELETE should reach the post-marker-commit barrier");
}
pub(crate) fn release(&self) {
self.state.release.notify_one();
}
}
#[cfg(test)]
impl Drop for VersionedDeleteMarkerCommitBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
let mut slot = VERSIONED_DELETE_MARKER_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("versioned delete-marker commit barrier mutex should not poison");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[cfg(test)]
async fn pause_versioned_delete_marker_after_commit(bucket: &str, object: &str) {
let state = VERSIONED_DELETE_MARKER_COMMIT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("versioned delete-marker commit barrier mutex should not poison")
.as_ref()
.filter(|state| state.bucket == bucket && state.object == object)
.cloned();
if let Some(state) = state {
state.arrived.notify_one();
state.release.notified().await;
@@ -913,6 +1121,160 @@ fn writer_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptions
lookup_opts
}
fn delete_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptions {
let mut lookup_opts = writer_pool_lookup_opts(opts, no_lock);
lookup_opts.skip_decommissioned = opts.data_movement;
lookup_opts
}
fn should_delete_from_all_pools(opts: &ObjectOptions, pool_count: usize) -> bool {
pool_count > 0 && (!opts.versioned && !opts.version_suspended || opts.version_id.is_some())
}
fn batch_delete_creates_latest_marker(object: &ObjectToDelete, delete_config_snapshot: &DeleteReplicationConfigSnapshot) -> bool {
if object.version_id.is_some() {
return false;
}
let object_name = decode_dir_object(&object.object_name);
let (versioned, version_suspended) = delete_config_snapshot.versioning_config().delete_state(&object_name);
versioned || version_suspended
}
fn batch_delete_targets_pool(creates_latest_marker: bool, marker_target_pool_idx: Option<usize>, pool_idx: usize) -> bool {
!creates_latest_marker || marker_target_pool_idx == Some(pool_idx)
}
#[cfg(test)]
struct BatchDeletePoolErrorInjectionState {
bucket: String,
pool_idx: usize,
errors: std::collections::HashMap<String, Error>,
observed: std::sync::atomic::AtomicUsize,
}
#[cfg(test)]
pub(crate) struct BatchDeletePoolErrorInjection {
state: Arc<BatchDeletePoolErrorInjectionState>,
}
#[cfg(test)]
static BATCH_DELETE_POOL_ERROR_INJECTION: std::sync::OnceLock<std::sync::Mutex<Option<Arc<BatchDeletePoolErrorInjectionState>>>> =
std::sync::OnceLock::new();
#[cfg(test)]
impl BatchDeletePoolErrorInjection {
pub(crate) fn install(bucket: &str, pool_idx: usize, errors: Vec<(String, Error)>) -> Self {
let state = Arc::new(BatchDeletePoolErrorInjectionState {
bucket: bucket.to_string(),
pool_idx,
errors: errors.into_iter().collect(),
observed: std::sync::atomic::AtomicUsize::new(0),
});
let mut slot = BATCH_DELETE_POOL_ERROR_INJECTION
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("batch delete pool error injection mutex should not poison");
assert!(slot.is_none(), "batch delete pool error injection must be unique");
*slot = Some(Arc::clone(&state));
Self { state }
}
pub(crate) fn observed(&self) -> usize {
self.state.observed.load(Ordering::Acquire)
}
}
#[cfg(test)]
impl Drop for BatchDeletePoolErrorInjection {
fn drop(&mut self) {
let mut slot = BATCH_DELETE_POOL_ERROR_INJECTION
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("batch delete pool error injection mutex should not poison");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[cfg(test)]
fn inject_batch_delete_pool_errors(
bucket: &str,
pool_idx: usize,
object_names: &[String],
result: &mut (Vec<DeletedObject>, Vec<Option<Error>>),
) {
let state = BATCH_DELETE_POOL_ERROR_INJECTION
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("batch delete pool error injection mutex should not poison")
.as_ref()
.filter(|state| state.bucket == bucket && state.pool_idx == pool_idx)
.cloned();
let Some(state) = state else {
return;
};
for (idx, object_name) in object_names.iter().enumerate() {
let Some(error) = state.errors.get(object_name) else {
continue;
};
if result.1[idx].is_none() && result.0[idx].found {
result.1[idx] = Some(error.clone());
state.observed.fetch_add(1, Ordering::AcqRel);
}
}
}
fn resolve_batch_delete_pool_results<'a>(
initial_error: Option<Error>,
pool_results: impl IntoIterator<Item = (&'a DeletedObject, &'a Option<Error>)>,
) -> (Option<DeletedObject>, Option<Error>, bool) {
let mut failure = initial_error.map(|err| (None, err));
let mut deleted = None;
let mut fallback: Option<(DeletedObject, Option<Error>)> = None;
let mut attempted = false;
for (pool_delete, pool_error) in pool_results {
attempted = true;
match pool_error {
Some(err) if is_err_object_not_found(err) || is_err_version_not_found(err) => {
if fallback.as_ref().is_none_or(|(_, error)| error.is_none()) {
fallback = Some(((*pool_delete).clone(), Some(err.clone())));
}
}
Some(err) => {
if failure.is_none() {
failure = Some((Some((*pool_delete).clone()), err.clone()));
}
}
None if pool_delete.found => {
if deleted.is_none() {
deleted = Some((*pool_delete).clone());
}
}
None => {
if fallback.is_none() {
fallback = Some(((*pool_delete).clone(), None));
}
}
}
}
if let Some((failed_delete, err)) = failure {
return (failed_delete, Some(err), attempted);
}
if let Some(deleted) = deleted {
return (Some(deleted), None, attempted);
}
if let Some((deleted, err)) = fallback {
return (Some(deleted), err, attempted);
}
(None, None, attempted)
}
fn transition_restore_pool_opts(opts: &ObjectOptions) -> ObjectOptions {
let mut lookup_opts = opts.clone();
lookup_opts.skip_decommissioned = true;
@@ -1533,6 +1895,89 @@ impl ECStore {
)))
}
pub(crate) async fn acquire_decommission_object_mutation_fence(
&self,
bucket: &str,
object: &str,
) -> Result<ObjectLockDiagGuard> {
if self.ctx.lock_manager().is_disabled() {
return Err(Error::other("decommission object migration requires namespace locking"));
}
#[cfg(test)]
let test_namespace_lock_fence =
decommission_mutation_fence_for_test(bucket, object, DecommissionMutationFenceTestPhase::Migration);
let object = encode_dir_object(object);
let mut opts = ObjectOptions::default();
let guard = self
.acquire_object_read_lock_if_needed("decommission_object", bucket, &object, &mut opts)
.await?
.ok_or_else(|| Error::other("decommission object migration failed to acquire its namespace fence"))?;
#[cfg(test)]
let guard = {
let mut guard = guard;
guard.test_namespace_lock_fence = test_namespace_lock_fence;
guard
};
Ok(guard)
}
pub(super) async 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));
opts.no_lock = match (fixed_set, target_set) {
(Some(fixed), Some(target)) => fixed.shares_namespace_lock_domain(&target).await,
_ => false,
};
}
pub(crate) async fn acquire_decommission_source_cleanup_fence(
&self,
bucket: &str,
object: &str,
source_set: &SetDisks,
) -> Result<SourceCleanupMutationFence> {
if self.ctx.lock_manager().is_disabled() {
return Err(Error::other("decommission source cleanup requires namespace locking"));
}
#[cfg(test)]
crate::data_movement::notify_source_cleanup_mutation_fence_pending(bucket, object);
#[cfg(test)]
let test_namespace_lock_fence =
decommission_mutation_fence_for_test(bucket, object, DecommissionMutationFenceTestPhase::SourceCleanup);
let object = encode_dir_object(object);
let fixed_set = Arc::clone(&self.pools[0].disk_set[0]);
let source_lock_covered = fixed_set.shares_namespace_lock_domain(source_set).await;
// 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
.acquire_object_write_lock("decommission_source_cleanup", bucket, &object)
.await?;
#[cfg(test)]
let guard = {
let mut guard = guard;
guard.test_namespace_lock_fence = test_namespace_lock_fence;
guard
};
Ok(SourceCleanupMutationFence {
guard,
source_lock_covered,
})
}
pub(crate) async fn acquire_all_object_read_locks(
&self,
op: &'static str,
@@ -1986,14 +2431,17 @@ 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)
.await;
let result = self.pools[idx]
.put_object_with_old_current_size(bucket, &object, data, &opts)
.await
@@ -2446,6 +2894,10 @@ impl ECStore {
} else {
None
};
#[cfg(test)]
if _object_lock_guard.is_some() {
notify_delete_namespace_acquired(bucket);
}
if let Some(trigger) = opts.lifecycle_delete_all.as_ref() {
let configs = delete_all_configs.as_ref().ok_or(StorageError::PreconditionFailed)?;
let expected_bucket_incarnation_id = opts.expected_bucket_incarnation_id.ok_or(StorageError::PreconditionFailed)?;
@@ -2479,7 +2931,7 @@ impl ECStore {
return Ok(ObjectInfo::default());
}
let gopts = writer_pool_lookup_opts(&opts, true);
let gopts = delete_pool_lookup_opts(&opts, true);
if opts.data_movement {
let existing_pool_info = self.get_pool_info_existing_with_opts(bucket, object, &gopts).await;
@@ -2584,6 +3036,8 @@ impl ECStore {
Err(err) if is_err_object_not_found(&err) && should_create_delete_marker_for_missing_object(&opts) => {
let target_pool_idx = self.get_pool_idx_no_lock(bucket, object, 0).await?;
let mut obj = self.pools[target_pool_idx].delete_object(bucket, object, opts).await?;
#[cfg(test)]
pause_versioned_delete_marker_after_commit(bucket, object).await;
obj.name = decode_dir_object(object);
return Ok(obj);
}
@@ -2622,7 +3076,7 @@ impl ECStore {
None
};
if !errs.is_empty() && !opts.versioned && !opts.version_suspended {
if should_delete_from_all_pools(&opts, errs.len()) {
let mut obj = match self.delete_object_from_all_pools(bucket, object, &opts, errs).await {
Ok(obj) => obj,
Err(err) => {
@@ -2646,6 +3100,8 @@ impl ECStore {
match pool.delete_object(bucket, object, opts.clone()).await {
Ok(res) => {
#[cfg(test)]
pause_versioned_delete_marker_after_commit(bucket, object).await;
if let (Some(api), Some(je)) = (tier_journal_api.as_ref(), journal_entry.as_ref()) {
commit_prepared_tier_delete_journal_entry(api, je).await;
}
@@ -2776,30 +3232,104 @@ impl ECStore {
Ok(guards) => guards,
Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err),
};
#[cfg(test)]
if !_object_lock_guards.is_empty() {
notify_delete_namespace_acquired(bucket);
}
let delete_config_snapshot = opts
.delete_replication_config_snapshot
.as_deref()
.expect("batch delete replication config snapshot should be loaded");
let latest_marker_objects = objects
.iter()
.map(|object| batch_delete_creates_latest_marker(object, delete_config_snapshot))
.collect::<Vec<_>>();
let marker_target_results = join_all(objects.iter().zip(&latest_marker_objects).map(
|(object, creates_marker)| async move {
if *creates_marker {
Some(self.get_pool_idx_no_lock(bucket, &object.object_name, 0).await)
} else {
None
}
},
))
.await;
let mut marker_target_pool_indices = Vec::with_capacity(objects.len());
for (idx, target_result) in marker_target_results.into_iter().enumerate() {
match target_result {
Some(Ok(pool_idx)) => marker_target_pool_indices.push(Some(pool_idx)),
Some(Err(err)) => {
del_errs[idx] = Some(err);
marker_target_pool_indices.push(None);
}
None => marker_target_pool_indices.push(None),
}
}
let mut futures = Vec::with_capacity(self.pools.len());
for pool in self.pools.iter() {
if self.is_pool_rebalancing(pool.pool_idx).await {
continue;
}
futures.push(pool.delete_objects(bucket, objects.clone(), opts.clone()));
let (object_indices, pool_objects): (Vec<_>, Vec<_>) = objects
.iter()
.enumerate()
.filter(|(idx, _)| {
batch_delete_targets_pool(latest_marker_objects[*idx], marker_target_pool_indices[*idx], pool.pool_idx)
})
.map(|(idx, object)| (idx, object.clone()))
.unzip();
if pool_objects.is_empty() {
continue;
}
let pool_opts = opts.clone();
futures.push(async move {
#[cfg(test)]
let pool_object_names = pool_objects
.iter()
.map(|object| object.object_name.clone())
.collect::<Vec<_>>();
let result = pool.delete_objects(bucket, pool_objects, pool_opts).await;
#[cfg(test)]
let result = {
let mut result = result;
inject_batch_delete_pool_errors(bucket, pool.pool_idx, &pool_object_names, &mut result);
result
};
(object_indices, result)
});
}
let results = join_all(futures).await;
for idx in 0..del_objects.len() {
for (dels, errs) in results.iter() {
if errs[idx].is_none() && dels[idx].found {
del_errs[idx] = None;
del_objects[idx] = dels[idx].clone();
break;
}
let pool_results = results.iter().filter_map(|(object_indices, (dels, errs))| {
let pool_object_idx = object_indices.binary_search(&idx).ok()?;
Some((&dels[pool_object_idx], &errs[pool_object_idx]))
});
let (deleted, error, attempted) = resolve_batch_delete_pool_results(del_errs[idx].take(), pool_results);
if let Some(deleted) = deleted {
del_objects[idx] = deleted;
}
del_errs[idx] = error;
if del_errs[idx].is_none() {
del_errs[idx] = errs[idx].clone();
del_objects[idx] = dels[idx].clone();
}
if !attempted && del_errs[idx].is_none() && latest_marker_objects[idx] {
del_objects[idx] = DeletedObject {
object_name: objects[idx].object_name.clone(),
version_id: objects[idx].version_id,
..Default::default()
};
del_errs[idx] = Some(StorageError::ObjectNotFound(bucket.to_owned(), objects[idx].object_name.clone()));
}
}
#[cfg(test)]
for (idx, object) in objects.iter().enumerate() {
if del_errs[idx].is_none() && del_objects[idx].delete_marker {
pause_versioned_delete_marker_after_commit(bucket, &object.object_name).await;
}
}
@@ -3374,6 +3904,80 @@ mod tests {
assert!(!same_distributed_lock_domain(&[first, second], &[other]));
}
#[tokio::test]
async fn decommission_fence_covers_dist_sets_with_same_clients_despite_different_namespaces() {
let ctx = Arc::new(crate::runtime::instance::InstanceContext::new());
let (_dirs, original_sets) = make_local_two_set_sets_with_ctx(Arc::clone(&ctx)).await;
let mut second_set = (*original_sets.disk_set[1]).clone();
second_set.lockers = original_sets.disk_set[0].lockers.clone();
let mut sets = (*original_sets).clone();
sets.disk_set[1] = Arc::new(second_set);
let sets = Arc::new(sets);
ctx.update_erasure_type(SetupType::DistErasure).await;
assert!(
sets.disk_set[0]
.lockers
.iter()
.zip(&sets.disk_set[1].lockers)
.all(|(fixed, hashed)| Arc::ptr_eq(fixed, hashed)),
"the regression requires identical distributed lock clients"
);
assert_ne!(sets.disk_set[0].set_index, sets.disk_set[1].set_index);
let pool_config = sets.endpoints.clone();
let store = new_prepared_reader_test_store_from_pools(vec![Arc::clone(&sets)], vec![pool_config], ctx);
let object = (0..1_000)
.map(|index| format!("decommission-dist-domain-{index}.bin"))
.find(|candidate| Arc::ptr_eq(&sets.get_disks_by_key(candidate), &sets.disk_set[1]))
.expect("a key should hash to the second set namespace");
let mutation_fence = store
.acquire_decommission_object_mutation_fence("bucket", &object)
.await
.expect("the fixed distributed mutation fence should be acquired");
let target_lock = sets.disk_set[1]
.new_ns_lock("bucket", &object)
.await
.expect("the hashed-set namespace lock should be created");
let target_err = target_lock
.get_write_lock(Duration::from_millis(50))
.await
.expect_err("the fixed read fence must conflict through the shared clients");
assert!(matches!(target_err, rustfs_lock::LockError::Timeout { .. }));
let mut put_opts = ObjectOptions::default();
store
.apply_decommission_target_mutation_fence(0, &object, &mut put_opts, Some(&mutation_fence))
.await;
assert!(put_opts.no_lock, "migration target PUT must reuse the covering fixed fence");
let mut multipart_opts = ObjectOptions::default();
store
.apply_decommission_target_mutation_fence(0, &object, &mut multipart_opts, Some(&mutation_fence))
.await;
assert!(multipart_opts.no_lock, "migration target multipart must reuse the covering fixed fence");
drop(mutation_fence);
let cleanup_object = (0..1_000)
.map(|index| format!("decommission-dist-cleanup-{index}.bin"))
.find(|candidate| Arc::ptr_eq(&sets.get_disks_by_key(candidate), &sets.disk_set[1]))
.expect("a cleanup key should hash to the second set namespace");
let source_fence = store
.acquire_decommission_source_cleanup_fence("bucket", &cleanup_object, sets.disk_set[1].as_ref())
.await
.expect("the fixed distributed cleanup fence should be acquired");
assert!(source_fence.source_lock_covered(), "source cleanup must reuse the covering fixed fence");
let source_lock = sets.disk_set[1]
.new_ns_lock("bucket", &cleanup_object)
.await
.expect("the source-set namespace lock should be created");
let source_err = source_lock
.get_read_lock(Duration::from_millis(50))
.await
.expect_err("the fixed write fence must conflict through the shared clients");
assert!(matches!(source_err, rustfs_lock::LockError::Timeout { .. }));
}
#[test]
fn select_snapshot_version_matching_normalizes_null_and_uuid_forms() {
let nil = Uuid::nil();
@@ -4433,6 +5037,159 @@ mod tests {
assert_eq!(lookup_opts.version_id.as_deref(), Some("vid-1"));
}
#[test]
fn ordinary_delete_lookup_includes_decommission_source_and_skips_rebalance_source() {
let lookup_opts = delete_pool_lookup_opts(&ObjectOptions::default(), true);
assert!(lookup_opts.no_lock);
assert!(!lookup_opts.skip_decommissioned);
assert!(lookup_opts.skip_rebalancing);
let explicit_version = delete_pool_lookup_opts(
&ObjectOptions {
versioned: true,
version_id: Some(uuid::Uuid::new_v4().to_string()),
..Default::default()
},
true,
);
assert!(!explicit_version.skip_decommissioned);
}
#[test]
fn delete_fans_out_for_unversioned_and_explicit_version_mutations() {
assert!(should_delete_from_all_pools(&ObjectOptions::default(), 1));
assert!(should_delete_from_all_pools(
&ObjectOptions {
versioned: true,
version_id: Some(uuid::Uuid::new_v4().to_string()),
..Default::default()
},
2,
));
assert!(!should_delete_from_all_pools(
&ObjectOptions {
versioned: true,
..Default::default()
},
1,
));
assert!(!should_delete_from_all_pools(&ObjectOptions::default(), 0));
}
#[test]
fn batch_delete_identifies_only_latest_versioned_markers() {
let versioned = DeleteReplicationConfigSnapshot::from_configs_for_test(
s3s::dto::VersioningConfiguration {
status: Some(s3s::dto::BucketVersioningStatus::from_static(s3s::dto::BucketVersioningStatus::ENABLED)),
..Default::default()
},
None,
);
let latest = ObjectToDelete {
object_name: "latest".to_string(),
..Default::default()
};
assert!(batch_delete_creates_latest_marker(&latest, &versioned));
assert!(!batch_delete_targets_pool(true, Some(1), 0));
assert!(batch_delete_targets_pool(true, Some(1), 1));
assert!(!batch_delete_targets_pool(true, Some(1), 2));
let explicit = ObjectToDelete {
object_name: "explicit".to_string(),
version_id: Some(uuid::Uuid::new_v4()),
..Default::default()
};
assert!(!batch_delete_creates_latest_marker(&explicit, &versioned));
assert!(batch_delete_targets_pool(false, Some(1), 0));
let unversioned = DeleteReplicationConfigSnapshot::default();
assert!(!batch_delete_creates_latest_marker(&latest, &unversioned));
assert!(batch_delete_targets_pool(false, None, 0));
}
#[test]
fn batch_delete_pool_failures_override_success_in_any_pool_order() {
let success = DeletedObject {
object_name: "object".to_string(),
found: true,
..Default::default()
};
let source_errors = [
StorageError::ErasureWriteQuorum,
StorageError::NamespaceLockQuorumUnavailable {
mode: "delete_objects_commit",
bucket: "bucket".to_string(),
object: "object".to_string(),
required: 1,
achieved: 0,
},
];
for source_error in source_errors {
for source_first in [true, false] {
let failed = (DeletedObject::default(), Some(source_error.clone()));
let succeeded = (success.clone(), None);
let pool_results = if source_first {
vec![failed, succeeded]
} else {
vec![succeeded, failed]
};
let (_, error, attempted) =
resolve_batch_delete_pool_results(None, pool_results.iter().map(|(deleted, error)| (deleted, error)));
assert!(attempted);
assert_eq!(error, Some(source_error.clone()));
}
}
}
#[test]
fn batch_delete_ignores_missing_pool_only_after_another_pool_succeeds() {
let success = DeletedObject {
object_name: "object".to_string(),
found: true,
..Default::default()
};
let missing_errors = [
StorageError::ObjectNotFound("bucket".to_string(), "object".to_string()),
StorageError::VersionNotFound("bucket".to_string(), "object".to_string(), "version".to_string()),
];
for missing_error in missing_errors {
let missing = (DeletedObject::default(), Some(missing_error.clone()));
for missing_first in [true, false] {
let succeeded = (success.clone(), None);
let pool_results = if missing_first {
vec![missing.clone(), succeeded]
} else {
vec![succeeded, missing.clone()]
};
let (deleted, error, attempted) =
resolve_batch_delete_pool_results(None, pool_results.iter().map(|(deleted, error)| (deleted, error)));
assert!(attempted);
let deleted = deleted.expect("successful pool result should be retained");
assert!(deleted.found);
assert_eq!(deleted.object_name, success.object_name.as_str());
assert!(error.is_none());
}
let missing_only = [missing];
let (_, error, attempted) =
resolve_batch_delete_pool_results(None, missing_only.iter().map(|(deleted, error)| (deleted, error)));
assert!(attempted);
assert_eq!(error, Some(missing_error));
}
let silent_missing = [(DeletedObject::default(), None)];
let (_, error, attempted) =
resolve_batch_delete_pool_results(None, silent_missing.iter().map(|(deleted, error)| (deleted, error)));
assert!(attempted);
assert!(error.is_none());
}
#[test]
fn data_movement_pool_lookup_opts_keeps_no_lock_for_tiered_moves() {
let lookup_opts = data_movement_pool_lookup_opts(
+29 -2
View File
@@ -73,7 +73,7 @@ pub(super) fn resolve_rebalance_delete_from_all_pools_result(
object: &str,
) -> Result<ObjectInfo> {
result.map_err(|err| {
if err == Error::PreconditionFailed {
if matches!(&err, Error::PreconditionFailed | Error::PrefixAccessDenied(_, _)) {
err
} else {
Error::other(format!("failed to delete rebalance source object {bucket}/{object}: {err}"))
@@ -86,7 +86,7 @@ fn is_ignorable_rebalance_delete_error(err: &Error) -> bool {
}
fn rebalance_delete_pool_error(pool_idx: usize, bucket: &str, object: &str, err: Error) -> Error {
if err == Error::PreconditionFailed {
if matches!(&err, Error::PreconditionFailed | Error::PrefixAccessDenied(_, _)) {
err
} else {
Error::other(format!("pool {pool_idx} delete failed for {bucket}/{object}: {err}"))
@@ -191,6 +191,18 @@ mod tests {
assert_eq!(err, Error::PreconditionFailed);
}
#[test]
fn rebalance_delete_result_preserves_prefix_access_denied() {
let err = resolve_rebalance_delete_from_all_pools_result(
Err(Error::PrefixAccessDenied("bucket".to_owned(), "object".to_owned())),
"bucket",
"object",
)
.expect_err("prefix access denial should remain structured");
assert_eq!(err, Error::PrefixAccessDenied("bucket".to_owned(), "object".to_owned()));
}
#[test]
fn rebalance_delete_pool_result_preserves_precondition_failed() {
let err = resolve_rebalance_delete_from_all_pools_results(
@@ -205,4 +217,19 @@ mod tests {
assert_eq!(err, Error::PreconditionFailed);
}
#[test]
fn rebalance_delete_pool_result_preserves_prefix_access_denied() {
let err = resolve_rebalance_delete_from_all_pools_results(
vec![RebalanceDeletePoolResult {
pool_idx: 0,
result: Err(Error::PrefixAccessDenied("bucket".to_owned(), "object".to_owned())),
}],
"bucket",
"object",
)
.expect_err("prefix access denial should remain structured");
assert_eq!(err, Error::PrefixAccessDenied("bucket".to_owned(), "object".to_owned()));
}
}
+1 -502
View File
@@ -21,7 +21,6 @@ use arc_swap::ArcSwapOption;
use rmp::Marker;
use serde::{Deserialize, Serialize};
use std::cmp::Ordering;
use std::collections::{HashMap, HashSet};
use std::str::from_utf8;
use std::{
fmt::Debug,
@@ -38,10 +37,8 @@ use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};
use tokio::spawn;
use tokio::sync::Mutex;
use tracing::{debug, warn};
use uuid::Uuid;
const SLASH_SEPARATOR: &str = "/";
pub const MAX_META_CACHE_HEAL_CANDIDATES: usize = 1024;
#[derive(Clone, Debug, Default)]
pub struct MetadataResolutionParams {
@@ -69,41 +66,6 @@ pub struct MetaCacheEntry {
pub reusable: bool,
}
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
pub enum MetaCacheHealCandidateKind {
Object,
DeleteMarker,
UnversionedObject,
}
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
pub struct MetaCacheHealCandidate {
pub object: String,
pub version_id: Option<Uuid>,
pub kind: MetaCacheHealCandidateKind,
/// Number of raw disk entries that carried this validated version.
pub replica_count: usize,
}
impl MetaCacheHealCandidate {
pub fn validated_version(&self) -> Option<Uuid> {
match self.kind {
MetaCacheHealCandidateKind::Object | MetaCacheHealCandidateKind::DeleteMarker => self.version_id,
MetaCacheHealCandidateKind::UnversionedObject => None,
}
}
pub fn is_unversioned(&self) -> bool {
self.kind == MetaCacheHealCandidateKind::UnversionedObject
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct MetaCacheHealDiscovery {
pub candidates: Vec<MetaCacheHealCandidate>,
pub unverified_count: usize,
}
impl MetaCacheEntry {
pub fn marshal_msg(&self) -> Result<Vec<u8>> {
let mut wr = Vec::new();
@@ -408,176 +370,6 @@ impl MetaCacheEntries {
})
}
/// Discover validated object/delete-marker versions and safe unversioned
/// inspection candidates in the raw entries without applying read quorum.
/// This is intentionally separate from [`Self::resolve`]: a sub-quorum
/// version is a valid heal target even though it must not participate in
/// normal reads or writes.
///
/// The validated list is bounded and deduplicated by object, version id,
/// and metadata kind; each candidate retains the number of raw disk
/// entries that carried it so callers can classify sub-quorum versions.
/// Entries whose xl.meta cannot be decoded are counted separately for
/// discovery accounting; they never become versionless destructive heal
/// requests and do not consume the validated quota. An
/// [`MetaCacheHealCandidateKind::UnversionedObject`] is always consumed by
/// a non-destructive scanner request.
pub fn discover_heal_candidates(&self, bucket: &str, max_candidates: usize) -> MetaCacheHealDiscovery {
let limit = max_candidates.min(MAX_META_CACHE_HEAL_CANDIDATES);
if limit == 0 || bucket.is_empty() {
return MetaCacheHealDiscovery::default();
}
let mut discovery = MetaCacheHealDiscovery {
candidates: Vec::<MetaCacheHealCandidate>::with_capacity(limit.min(self.0.len())),
unverified_count: 0,
};
let mut seen: HashMap<(String, Option<Uuid>, MetaCacheHealCandidateKind), usize> =
HashMap::with_capacity(limit.min(self.0.len()));
for entry in self.0.iter().flatten() {
if !valid_heal_candidate_name(bucket, entry) {
continue;
}
let meta = match FileMeta::load(&entry.metadata) {
Ok(meta) => meta,
Err(_) => {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
};
let mut entry_seen = HashSet::new();
for shallow in meta.versions {
let version = match shallow.parse_version_meta() {
Ok(version) if version.valid() => version,
Ok(_) | Err(_) => {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
};
if version.free_version() {
continue;
}
let payload_header = version.header();
if normalize_version_id(shallow.header.version_id) != normalize_version_id(payload_header.version_id)
|| shallow.header.version_type != payload_header.version_type
{
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
let (kind, version_id) = match version.version_type {
VersionType::Object
if version.object.is_some() && version.delete_marker.is_none() && version.legacy_object.is_none() =>
{
match version.object.as_ref().and_then(|object| object.version_id) {
Some(id) if !id.is_nil() => (MetaCacheHealCandidateKind::Object, Some(id)),
Some(_) | None => (MetaCacheHealCandidateKind::UnversionedObject, None),
}
}
VersionType::Delete
if version.delete_marker.is_some() && version.object.is_none() && version.legacy_object.is_none() =>
{
let Some(id) = version.delete_marker.as_ref().and_then(|marker| marker.version_id) else {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
};
if id.is_nil() {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
(MetaCacheHealCandidateKind::DeleteMarker, Some(id))
}
VersionType::Legacy
if version.legacy_object.is_some() && version.object.is_none() && version.delete_marker.is_none() =>
{
let Some(legacy) = version.legacy_object.as_ref() else {
continue;
};
if legacy.version_id.is_empty() {
(MetaCacheHealCandidateKind::UnversionedObject, None)
} else {
let Ok(id) = Uuid::parse_str(&legacy.version_id) else {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
};
if id.is_nil() {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
(MetaCacheHealCandidateKind::Object, Some(id))
}
}
_ => {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
};
if normalize_version_id(payload_header.version_id) != version_id {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
// `all_parts=true` is the trust-boundary check for versioned
// candidates. A null/legacy object may still need the old
// non-destructive inspection fallback when its part arrays
// are parseable but incomplete; never use that fallback for
// a candidate carrying a real version id.
let file_info = match version.clone().into_fileinfo(bucket, &entry.name, true) {
Ok(file_info) => file_info,
Err(_) if version_id.is_none() && matches!(kind, MetaCacheHealCandidateKind::UnversionedObject) => {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
match version.into_fileinfo(bucket, &entry.name, false) {
Ok(file_info) => file_info,
Err(_) => continue,
}
}
Err(_) => {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
};
if file_info.volume != bucket || file_info.name != entry.name {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
let candidate = MetaCacheHealCandidate {
object: entry.name.clone(),
version_id,
kind,
replica_count: 1,
};
let key = (candidate.object.clone(), candidate.version_id, candidate.kind.clone());
if !entry_seen.contains(&key) {
// Keep per-entry dedupe bounded as well as the global
// candidate union. Once the cap is reached, only keys
// already present in the global map may update replica
// counts; novel versions are accounting-only.
if entry_seen.len() >= limit && !seen.contains_key(&key) {
continue;
}
entry_seen.insert(key.clone());
}
if let Some(index) = seen.get(&key).copied() {
discovery.candidates[index].replica_count = discovery.candidates[index].replica_count.saturating_add(1);
} else {
seen.insert(key, discovery.candidates.len());
discovery.candidates.push(candidate);
if discovery.candidates.len() >= limit {
return discovery;
}
}
}
}
discovery
}
fn resolve_inner(&self, mut params: MetadataResolutionParams, enforce_write_quorum: bool) -> Option<MetaCacheEntry> {
if self.0.is_empty() {
debug!(
@@ -754,28 +546,6 @@ impl MetaCacheEntries {
}
}
fn valid_heal_candidate_name(bucket: &str, entry: &MetaCacheEntry) -> bool {
if bucket.is_empty() || entry.name.is_empty() || entry.is_dir() || entry.name.contains('\0') {
return false;
}
// Validate raw key components without normalizing them. A dot component
// could otherwise escape the bucket when the key is later mapped back to
// a disk path; a final empty component is retained for valid keys ending
// in '/'.
let mut components = entry.name.split('/').peekable();
while let Some(component) = components.next() {
if component == "." || component == ".." || (component.is_empty() && components.peek().is_some()) {
return false;
}
}
true
}
fn normalize_version_id(version_id: Option<Uuid>) -> Option<Uuid> {
version_id.filter(|id| !id.is_nil())
}
#[derive(Debug, Default)]
pub struct MetaCacheEntriesSortedResult {
pub entries: Option<MetaCacheEntriesSorted>,
@@ -1221,7 +991,7 @@ impl<T: Clone + Debug + Send + Sync + 'static> Cache<T> {
mod tests {
use super::*;
use crate::test_data::create_real_xlmeta;
use crate::{FileMetaVersion, MetaDeleteMarker, MetaObjectV1, MetaObjectV1Erasure, MetaObjectV1Stat, TRANSITION_COMPLETE};
use crate::{FileMetaVersion, MetaDeleteMarker, TRANSITION_COMPLETE};
use std::collections::HashMap;
use std::io::Cursor;
use std::sync::{
@@ -1822,277 +1592,6 @@ mod tests {
);
}
#[test]
fn discover_heal_candidates_keeps_sub_quorum_versions_and_deduplicates() {
let now = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let entries = MetaCacheEntries(vec![
Some(metacache_entry_single_version(1, now, "one")),
Some(metacache_entry_single_version(2, now, "two")),
Some(metacache_entry_single_version(2, now, "two")),
Some(metacache_entry_single_version(3, now, "three")),
]);
let discovery = entries.discover_heal_candidates("bucket", 16);
let ids: std::collections::HashSet<Uuid> = discovery
.candidates
.iter()
.filter_map(|candidate| candidate.version_id)
.collect();
assert_eq!(
ids,
[Uuid::from_u128(1), Uuid::from_u128(2), Uuid::from_u128(3)]
.into_iter()
.collect()
);
assert_eq!(discovery.candidates.len(), 3, "duplicate tied versions must be emitted once");
assert_eq!(
discovery
.candidates
.iter()
.find(|candidate| candidate.version_id == Some(Uuid::from_u128(2)))
.expect("duplicate version should be discovered")
.replica_count,
2
);
}
#[test]
fn discover_heal_candidates_covers_divergent_quorum_boundaries_n2_n4_n6() {
let now = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
for (disk_count, quorum) in [(2usize, 1usize), (4, 2), (6, 3)] {
let target_id = Uuid::from_u128(0x1000 + disk_count as u128);
for target_replicas in [quorum.saturating_sub(1), quorum, quorum + 1] {
let entries = (0..disk_count)
.map(|disk| {
let version_id = if disk < target_replicas {
target_id
} else {
Uuid::from_u128(0x2000 + disk as u128)
};
Some(metacache_entry_single_version(version_id.as_u128(), now, "divergent"))
})
.collect();
let discovery = MetaCacheEntries(entries).discover_heal_candidates("bucket", 32);
let target = discovery
.candidates
.iter()
.find(|candidate| candidate.version_id == Some(target_id));
assert_eq!(target.is_some(), target_replicas > 0, "N={disk_count}, replicas={target_replicas}");
if let Some(target) = target {
assert_eq!(target.replica_count, target_replicas);
}
}
}
}
#[test]
fn discover_heal_candidates_separates_delete_markers_and_preserves_unversioned_objects() {
let mut marker_meta = FileMeta::new();
marker_meta
.add_version(FileInfo {
volume: "bucket".to_string(),
name: "object".to_string(),
version_id: Some(Uuid::from_u128(99)),
deleted: true,
mod_time: Some(OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp")),
..Default::default()
})
.expect("delete marker should be added");
let marker = MetaCacheEntry {
name: "object".to_string(),
metadata: marker_meta.marshal_msg().expect("delete marker metadata should marshal"),
cached: Some(marker_meta),
reusable: false,
};
let unversioned_entry = metacache_entry_with_mod_time(
OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp"),
"unversioned",
);
let discovery = MetaCacheEntries(vec![Some(marker), Some(unversioned_entry)]).discover_heal_candidates("bucket", 16);
assert!(discovery.candidates.iter().any(|candidate| {
candidate.kind == MetaCacheHealCandidateKind::DeleteMarker && candidate.version_id == Some(Uuid::from_u128(99))
}));
assert!(discovery.candidates.iter().any(|candidate| {
candidate.kind == MetaCacheHealCandidateKind::UnversionedObject && candidate.version_id.is_none()
}));
}
#[test]
fn discover_heal_candidates_skips_free_versions() {
let object_id = Uuid::from_u128(100);
let free_id = Uuid::from_u128(101);
let mut meta = FileMeta::new();
meta.add_version(FileInfo {
volume: "bucket".to_string(),
name: "object".to_string(),
version_id: Some(object_id),
transition_status: TRANSITION_COMPLETE.to_string(),
transitioned_objname: "remote/object".to_string(),
transition_version_id: Some(Uuid::from_u128(102)),
transition_tier: "WARM".to_string(),
mod_time: Some(OffsetDateTime::now_utc()),
..Default::default()
})
.expect("transitioned object should be added");
let mut delete = FileInfo {
volume: "bucket".to_string(),
name: "object".to_string(),
version_id: Some(object_id),
mod_time: Some(OffsetDateTime::now_utc()),
..Default::default()
};
delete.set_tier_free_version_id(&free_id.to_string());
meta.delete_version(&delete).expect("free version should be persisted");
let discovery = MetaCacheEntries(vec![Some(MetaCacheEntry {
name: "object".to_string(),
metadata: meta.marshal_msg().expect("free version metadata should marshal"),
cached: Some(meta),
reusable: false,
})])
.discover_heal_candidates("bucket", 16);
assert!(discovery.candidates.is_empty());
}
#[test]
fn discover_heal_candidates_preserves_unversioned_legacy_object() {
let legacy = MetaObjectV1 {
version: "1.0.1".to_string(),
format: "xl".to_string(),
stat: MetaObjectV1Stat {
size: 1,
mod_time: Some(OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp")),
name: "object".to_string(),
..Default::default()
},
erasure: MetaObjectV1Erasure {
data_blocks: 4,
parity_blocks: 2,
index: 1,
distribution: vec![1, 2, 3, 4, 5, 6],
..Default::default()
},
..Default::default()
};
let version = FileMetaVersion {
version_type: VersionType::Legacy,
legacy_object: Some(legacy),
..Default::default()
};
let mut meta = FileMeta::new();
meta.versions
.push(FileMetaShallowVersion::try_from(version).expect("legacy metadata should marshal"));
let discovery = MetaCacheEntries(vec![Some(MetaCacheEntry {
name: "object".to_string(),
metadata: meta.marshal_msg().expect("legacy metadata should marshal"),
cached: Some(meta),
reusable: false,
})])
.discover_heal_candidates("bucket", 16);
assert!(discovery.candidates.iter().any(|candidate| {
candidate.kind == MetaCacheHealCandidateKind::UnversionedObject && candidate.version_id.is_none()
}));
}
#[test]
fn discover_heal_candidates_rejects_nil_and_malformed_metadata_and_is_bounded() {
let now = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let mut nil = metacache_entry_single_version(1, now, "nil");
let mut nil_meta = FileMeta::load(&nil.metadata).expect("nil fixture should decode");
let mut nil_version = nil_meta.versions[0]
.parse_version_meta()
.expect("nil fixture version should decode");
nil_version.object.as_mut().expect("object fixture").version_id = Some(Uuid::nil());
nil_meta.versions[0] = FileMetaShallowVersion::try_from(nil_version).expect("nil fixture should marshal");
nil.metadata = nil_meta.marshal_msg().expect("nil fixture metadata should marshal");
let mut mismatched = metacache_entry_single_version(2, now, "mismatched");
let mut mismatched_meta = FileMeta::load(&mismatched.metadata).expect("mismatched fixture should decode");
mismatched_meta.versions[0].header.version_id = Some(Uuid::from_u128(200));
mismatched.metadata = mismatched_meta.marshal_msg().expect("mismatched metadata should marshal");
let mut short_parts = metacache_entry_single_version(3, now, "short-parts");
let mut short_parts_meta = FileMeta::load(&short_parts.metadata).expect("short-parts fixture should decode");
let mut short_parts_version = short_parts_meta.versions[0]
.parse_version_meta()
.expect("short-parts fixture version should decode");
let object = short_parts_version.object.as_mut().expect("object fixture");
object.part_numbers = vec![1];
object.part_actual_sizes = vec![1];
object.part_sizes.clear();
short_parts_meta.versions[0] = FileMetaShallowVersion::try_from(short_parts_version).expect("short-parts should marshal");
short_parts.metadata = short_parts_meta.marshal_msg().expect("short-parts metadata should marshal");
let mut short_unversioned = metacache_entry_with_mod_time(now, "short-unversioned");
let mut short_unversioned_meta =
FileMeta::load(&short_unversioned.metadata).expect("short-unversioned fixture should decode");
let mut short_unversioned_version = short_unversioned_meta.versions[0]
.parse_version_meta()
.expect("short-unversioned version should decode");
let unversioned_object = short_unversioned_version.object.as_mut().expect("unversioned object fixture");
unversioned_object.part_numbers = vec![1];
unversioned_object.part_actual_sizes = vec![1];
unversioned_object.part_sizes.clear();
short_unversioned_meta.versions[0] =
FileMetaShallowVersion::try_from(short_unversioned_version).expect("short-unversioned should marshal");
short_unversioned.metadata = short_unversioned_meta
.marshal_msg()
.expect("short-unversioned metadata should marshal");
let mut malformed = nil.clone();
malformed.name = "malformed".to_string();
malformed.metadata = vec![1, 2, 3];
let entries = MetaCacheEntries(
std::iter::once(Some(nil))
.chain(std::iter::once(Some(mismatched)))
.chain(std::iter::once(Some(short_parts)))
.chain(std::iter::once(Some(short_unversioned)))
.chain(std::iter::once(Some(malformed)))
.chain((0..32).map(|id| Some(metacache_entry_single_version(id + 10, now, "bounded"))))
.collect(),
);
let discovery = entries.discover_heal_candidates("bucket", 5);
assert!(discovery.candidates.len() <= 5);
assert!(
!discovery
.candidates
.iter()
.any(|candidate| candidate.version_id == Some(Uuid::nil()))
);
assert!(
!discovery
.candidates
.iter()
.any(|candidate| candidate.version_id == Some(Uuid::from_u128(2)))
);
assert!(
!discovery
.candidates
.iter()
.any(|candidate| candidate.version_id == Some(Uuid::from_u128(3)))
);
assert!(discovery.candidates.iter().any(|candidate| {
candidate.kind == MetaCacheHealCandidateKind::UnversionedObject && candidate.version_id.is_none()
}));
assert!(
discovery.unverified_count >= 1,
"malformed and rejected metadata must remain observable during discovery"
);
for invalid_name in ["../object", "object//", "object\0name"] {
let mut entry = metacache_entry_single_version(400, now, invalid_name);
entry.name = invalid_name.to_string();
let discovery = MetaCacheEntries(vec![Some(entry)]).discover_heal_candidates("bucket", 5);
assert!(
discovery.candidates.is_empty(),
"invalid key should not become a heal candidate: {invalid_name:?}"
);
}
}
#[test]
fn resolve_rejects_partial_latest_and_returns_committed_previous_metadata() {
let old_mod_time = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
@@ -722,10 +722,6 @@ pub struct DeleteVersionsResponse {
pub errors: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
#[prost(message, optional, tag = "3")]
pub error: ::core::option::Option<Error>,
/// Senders dual-write the legacy strings and typed entries. Receivers prefer typed entries
/// when present and fall back to strings for peers that predate this field. Code zero means success.
#[prost(message, repeated, tag = "4")]
pub item_errors: ::prost::alloc::vec::Vec<Error>,
}
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
pub struct ReadMultipleRequest {
-3
View File
@@ -493,9 +493,6 @@ message DeleteVersionsResponse {
bool success = 1;
repeated string errors = 2;
optional Error error = 3;
// Senders dual-write the legacy strings and typed entries. Receivers prefer typed entries
// when present and fall back to strings for peers that predate this field. Code zero means success.
repeated Error item_errors = 4;
}
message ReadMultipleRequest {
+68 -56
View File
@@ -43,7 +43,7 @@ use rustfs_common::metrics::{
UpdateCurrentPathFn, current_path_updater, global_metrics,
};
use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, trace_emit, trace_subscriber_count};
use rustfs_filemeta::{MAX_META_CACHE_HEAL_CANDIDATES, MetaCacheEntries, MetaCacheEntry, MetaCacheHealCandidateKind};
use rustfs_filemeta::{MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams};
use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf};
use s3s::dto::{BucketLifecycleConfiguration, ObjectLockConfiguration, VersioningConfiguration};
use time::OffsetDateTime;
@@ -96,10 +96,6 @@ const METRIC_SCANNER_EXCESS_OBJECT_VERSION_SIZE_TOTAL: &str = "rustfs_scanner_ex
const METRIC_SCANNER_EXCESS_FOLDERS_TOTAL: &str = "rustfs_scanner_excess_folders_total";
const METRIC_SCANNER_PENDING_HEAL_PRUNE_TOTAL: &str = "rustfs_scanner_pending_heal_prune_total";
const METRIC_SCANNER_PENDING_HEAL_MALFORMED_TOTAL: &str = "rustfs_scanner_pending_heal_malformed_total";
const METRIC_SCANNER_HEAL_DISCOVERY_CANDIDATES_TOTAL: &str = "rustfs_scanner_heal_discovery_candidates_total";
const METRIC_SCANNER_HEAL_DISCOVERY_SUB_QUORUM_TOTAL: &str = "rustfs_scanner_heal_discovery_sub_quorum_total";
const METRIC_SCANNER_HEAL_DISCOVERY_UNVERIFIED_TOTAL: &str = "rustfs_scanner_heal_discovery_unverified_total";
const METRIC_SCANNER_HEAL_DISCOVERY_QUEUED_TOTAL: &str = "rustfs_scanner_heal_discovery_queued_total";
const MAX_PENDING_SCANNER_HEAL_RETRIES_PER_BUCKET: usize = 128;
// --- scanner excess alerts as S3 notification events (rustfs/backlog#1868) --
@@ -887,7 +883,7 @@ impl FolderScanner {
object: Option<String>,
version_id: Option<String>,
request: HealChannelRequest,
) -> Result<HealAdmissionResult, ScannerError> {
) -> Result<(), ScannerError> {
let candidate_type = pending_scanner_heal_candidate_type(kind);
let priority = request.priority;
let scan_mode = request.scan_mode.unwrap_or(self.scan_mode);
@@ -915,7 +911,7 @@ impl FolderScanner {
error = %err,
"Scanner deferred heal request after channel error"
);
return Ok(HealAdmissionResult::Full);
return Ok(());
}
};
self.update_pending_scanner_heal_after_admission(
@@ -927,7 +923,7 @@ impl FolderScanner {
result,
);
if result.is_admitted() {
return Ok(result);
return Ok(());
}
record_high_priority_heal_escalation(candidate_type, priority, result);
@@ -948,7 +944,7 @@ impl FolderScanner {
state = "high_priority_not_admitted",
"Scanner high-priority heal admission failed"
);
Ok(result)
Ok(())
}
pub fn set_heal_object_select(&mut self, prob: u32) {
@@ -1740,7 +1736,14 @@ impl FolderScanner {
break;
}
let mut previous_bucket = String::new();
let mut resolver = MetadataResolutionParams {
dir_quorum: self.disks_quorum,
obj_quorum: self.disks_quorum,
bucket: "".to_string(),
strict: false,
..Default::default()
};
for name in abandoned_children {
if !self.should_heal().await {
break;
@@ -1748,7 +1751,7 @@ impl FolderScanner {
let (bucket, prefix) = path2_bucket_object(name.as_str());
if bucket != previous_bucket {
if bucket != resolver.bucket {
self.send_required_scanner_heal_request(
PendingScannerHealKind::Bucket,
bucket.clone(),
@@ -1757,9 +1760,10 @@ impl FolderScanner {
build_bucket_heal_request(bucket.clone(), HealChannelPriority::High),
)
.await?;
previous_bucket = bucket.clone();
}
resolver.bucket = bucket.clone();
let child_ctx = ctx.child_token();
let (agreed_tx, mut agreed_rx) = mpsc::channel::<String>(1);
@@ -1876,7 +1880,6 @@ impl FolderScanner {
let mut agreed_closed = false;
let mut partial_closed = false;
let mut finished_closed = false;
let mut seen_heal_candidates: HashSet<(String, Option<String>, MetaCacheHealCandidateKind)> = HashSet::new();
loop {
if agreed_closed && partial_closed && finished_closed {
@@ -1901,56 +1904,65 @@ impl FolderScanner {
break;
}
let discovery = entries.discover_heal_candidates(&bucket, MAX_META_CACHE_HEAL_CANDIDATES);
counter!(METRIC_SCANNER_HEAL_DISCOVERY_CANDIDATES_TOTAL)
.increment(u64::try_from(discovery.candidates.len()).unwrap_or(u64::MAX));
counter!(METRIC_SCANNER_HEAL_DISCOVERY_SUB_QUORUM_TOTAL).increment(
u64::try_from(
discovery
.candidates
.iter()
.filter(|candidate| candidate.replica_count < disks_quorum)
.count(),
)
.unwrap_or(u64::MAX),
);
counter!(METRIC_SCANNER_HEAL_DISCOVERY_UNVERIFIED_TOTAL).increment(
u64::try_from(discovery.unverified_count).unwrap_or(u64::MAX),
);
let Some(entry) = resolve_object_heal_entry(&entries, resolver.clone()) else {
continue;
};
for candidate in discovery.candidates {
let version_id = candidate.validated_version().map(|id| id.to_string());
let identity = (candidate.object.clone(), version_id.clone(), candidate.kind.clone());
if seen_heal_candidates.len() >= MAX_META_CACHE_HEAL_CANDIDATES
&& !seen_heal_candidates.contains(&identity)
{
(self.update_current_path)(&entry.name).await;
if entry.is_dir() {
continue;
}
let fivs = match entry.file_info_versions(&bucket) {
Ok(fivs) => fivs,
Err(e) => {
error!(
target: "rustfs::scanner::folder",
event = EVENT_SCANNER_FOLDER_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_FOLDER,
bucket = %bucket,
entry = %entry.name,
state = "file_info_versions_failed",
error = %e,
"Scanner list_path_raw failed to resolve file versions"
);
self.send_required_scanner_heal_request(
PendingScannerHealKind::Object,
bucket.clone(),
Some(entry.name.clone()),
None,
build_object_heal_request(
bucket.clone(),
entry.name.clone(),
None,
self.scan_mode,
HealChannelPriority::High,
),
)
.await?;
found_objects = true;
continue;
}
if !seen_heal_candidates.insert(identity) {
continue;
}
let mut request = build_object_heal_request(
bucket.clone(),
candidate.object.clone(),
version_id.clone(),
self.scan_mode,
HealChannelPriority::High,
);
if candidate.is_unversioned() {
request.remove_corrupted = Some(false);
}
(self.update_current_path)(&candidate.object).await;
let admission = self.send_required_scanner_heal_request(
};
for fiv in fivs.versions {
let version_id = fiv.version_id.and_then(|v| if v.is_nil() { None } else { Some(v.to_string()) });
self.send_required_scanner_heal_request(
PendingScannerHealKind::Object,
bucket.clone(),
Some(candidate.object.clone()),
version_id,
request,
Some(entry.name.clone()),
version_id.clone(),
build_object_heal_request(
bucket.clone(),
entry.name.clone(),
version_id,
self.scan_mode,
HealChannelPriority::High,
),
)
.await?;
if admission.is_admitted() {
counter!(METRIC_SCANNER_HEAL_DISCOVERY_QUEUED_TOTAL).increment(1);
}
found_objects = true;
}
@@ -13,8 +13,6 @@
// limitations under the License.
/// Per-object scan actions: ScannerItem, the get-size failure policy, and the heal/ILM admission helpers.
use super::*;
#[cfg(test)]
use rustfs_filemeta::MetadataResolutionParams;
/// Cached folder information for scanning
#[derive(Clone, Debug)]
@@ -90,7 +88,6 @@ pub(super) fn build_object_heal_request(
}
}
#[cfg(test)]
pub(super) fn resolve_object_heal_entry(
entries: &MetaCacheEntries,
resolver: MetadataResolutionParams,
+2 -6
View File
@@ -305,17 +305,13 @@ pub(super) fn build_pending_scanner_heal_request(entry: &PendingScannerHeal) ->
match entry.kind {
PendingScannerHealKind::Bucket => Some(build_bucket_heal_request(entry.bucket.clone(), HealChannelPriority::High)),
PendingScannerHealKind::Object => entry.object.as_ref().map(|object| {
let mut request = build_object_heal_request(
build_object_heal_request(
entry.bucket.clone(),
object.clone(),
entry.version_id.clone(),
entry.scan_mode,
HealChannelPriority::High,
);
if entry.version_id.is_none() {
request.remove_corrupted = Some(false);
}
request
)
}),
}
}
+6 -66
View File
@@ -17,7 +17,7 @@ use crate::SCANNER_SLEEPER;
use super::*;
use crate::storage_api::VersionPurgeStatusType;
use crate::{DiskOption, Endpoint, STORAGE_FORMAT_FILE, TierStats, new_disk, storageclass};
use rustfs_filemeta::{FileInfo, FileMeta, MetadataResolutionParams};
use rustfs_filemeta::{FileInfo, FileMeta};
use std::io::Write;
#[cfg(unix)]
use std::os::unix::fs::{PermissionsExt, symlink};
@@ -1121,17 +1121,6 @@ fn test_pending_heal_reconstructs_object_request_with_version() {
assert_eq!(request.source, HealRequestSource::Scanner);
}
#[test]
fn test_pending_heal_reconstructs_unversioned_request_without_removal() {
let pending = pending_heal(PendingScannerHealKind::Object, "bucket", Some("object"), None, 1, 1);
let request = build_pending_scanner_heal_request(&pending).expect("unversioned object request should rebuild");
assert!(request.object_version_id.is_none());
assert_eq!(request.remove_corrupted, Some(false));
assert_eq!(request.recreate_missing, Some(false));
}
#[test]
fn test_pending_heal_retry_candidates_respect_cap_and_order() {
let pending: Vec<PendingScannerHeal> = (0..(MAX_PENDING_SCANNER_HEAL_RETRIES_PER_BUCKET + 2))
@@ -1334,20 +1323,6 @@ fn metadata_for_object(bucket: &str, object: &str) -> Vec<u8> {
meta.marshal_msg().expect("test metadata should marshal")
}
fn metadata_for_object_version(bucket: &str, object: &str, version_id: Option<Uuid>) -> Vec<u8> {
let mut file_info = FileInfo::new(object, 4, 2);
file_info.volume = bucket.to_string();
file_info.name = object.to_string();
file_info.version_id = version_id;
file_info.versioned = version_id.is_some();
file_info.mod_time = Some(OffsetDateTime::now_utc());
file_info.size = 1;
let mut meta = FileMeta::new();
meta.add_version(file_info).expect("test metadata version should be accepted");
meta.marshal_msg().expect("test metadata should marshal")
}
async fn write_test_object_metadata(root: &std::path::Path, bucket: &str, object: &str) {
write_test_object_metadata_bytes(root, bucket, object, &metadata_for_object(bucket, object)).await;
}
@@ -1751,21 +1726,12 @@ async fn test_scan_folder_exits_when_abandoned_child_listing_finishes() {
let _guard = TestGuard::new(60, 100, &mut scanner, temp_dir.clone());
let heal_starts = Arc::new(AtomicUsize::new(0));
let heal_starts_clone = heal_starts.clone();
let healed_versions = Arc::new(Mutex::new(Vec::<Option<String>>::new()));
let healed_versions_clone = healed_versions.clone();
let mut heal_rx =
rustfs_common::heal_channel::init_heal_channel().expect("heal channel should initialize once for scanner tests");
let _heal_responder = tokio::spawn(async move {
while let Some(command) = heal_rx.recv().await {
if let rustfs_common::heal_channel::HealChannelCommand::Start {
request, response_tx, ..
} = command
{
if let rustfs_common::heal_channel::HealChannelCommand::Start { response_tx, .. } = command {
heal_starts_clone.fetch_add(1, Ordering::Relaxed);
healed_versions_clone
.lock()
.expect("heal version capture lock should not be poisoned")
.push(request.object_version_id);
let _ = response_tx.send(Ok(HealAdmissionResult::Accepted));
}
}
@@ -1773,18 +1739,13 @@ async fn test_scan_folder_exits_when_abandoned_child_listing_finishes() {
let bucket = "src-archive";
let object = "snapshots/37b3f20d941e2f5e6d99114d9bb2f3e67a8a2e5c9c4c5a1b0d6e7f8091a2b3c4";
let orphan_version = Uuid::from_u128(0x1934);
let shared_version = Uuid::from_u128(0x1935);
let orphan_metadata = metadata_for_object_version(bucket, object, Some(orphan_version));
let shared_metadata = metadata_for_object_version(bucket, object, Some(shared_version));
write_test_object_metadata_bytes(&temp_dir, bucket, object, &orphan_metadata).await;
let mut expected_metadata = vec![(temp_dir.join(bucket).join(object).join("xl.meta"), orphan_metadata.clone())];
let metadata = metadata_for_object(bucket, object);
write_test_object_metadata_bytes(&temp_dir, bucket, object, &metadata).await;
let mut disks = vec![scanner.local_disk.clone()];
for disk_name in ["disk2", "disk3", "disk4"] {
let disk_root = temp_dir.join(disk_name);
write_test_object_metadata_bytes(&disk_root, bucket, object, &shared_metadata).await;
expected_metadata.push((disk_root.join(bucket).join(object).join("xl.meta"), shared_metadata.clone()));
write_test_object_metadata_bytes(&disk_root, bucket, object, &metadata).await;
let endpoint = Endpoint::try_from(disk_root.to_string_lossy().as_ref()).expect("failed to create extra disk endpoint");
let disk = new_disk(
&endpoint,
@@ -1833,29 +1794,8 @@ async fn test_scan_folder_exits_when_abandoned_child_listing_finishes() {
.new_cache
.checked_flatten(bucket)
.expect("healed cache must contain canonical child links");
// The fixture intentionally exposes two divergent version histories, so
// the scanner keeps both logical versions visible while discovering heals.
assert_eq!(root.objects, 2);
assert_eq!(root.objects, 1);
assert!(heal_starts.load(Ordering::Relaxed) > 0, "test must execute the heal child-link path");
let orphan_version_text = orphan_version.to_string();
assert!(
healed_versions
.lock()
.expect("heal version capture lock should not be poisoned")
.iter()
.any(|version| version.as_deref() == Some(orphan_version_text.as_str())),
"sub-quorum orphan version must be submitted as an exact heal candidate"
);
for (path, expected) in expected_metadata {
assert_eq!(
tokio::fs::read(&path)
.await
.expect("scanner discovery must not delete metadata"),
expected,
"scanner discovery must not modify candidate metadata: {}",
path.display()
);
}
}
#[tokio::test]
+11 -43
View File
@@ -146,29 +146,6 @@ fn encode_file_info_msgpack(value: &FileInfo) -> std::result::Result<Vec<u8>, Di
encode_msgpack_with_capacity(value, "FileInfo", FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT)
}
fn encode_delete_versions_errors(disk_errors: Vec<Option<DiskError>>) -> (Vec<String>, Vec<Error>) {
let mut errors = Vec::with_capacity(disk_errors.len());
let mut item_errors = Vec::with_capacity(disk_errors.len());
for error in disk_errors {
match error {
Some(error) => {
let code = match &error {
DiskError::Io(source) if source.kind() == std::io::ErrorKind::NotFound => DiskError::FileNotFound.to_u32(),
_ => error.to_u32(),
};
let error_info = error.to_string();
errors.push(error_info.clone());
item_errors.push(Error { code, error_info });
}
None => {
errors.push(String::new());
item_errors.push(Error::default());
}
}
}
(errors, item_errors)
}
fn encode_msgpack_named<T: serde::Serialize>(value: &T, value_name: &str) -> std::result::Result<Vec<u8>, DiskError> {
let mut serializer = rmp_serde::Serializer::new(Vec::with_capacity(MSGPACK_ENCODE_CAPACITY_HINT)).with_struct_map();
value
@@ -575,7 +552,6 @@ impl NodeService {
success: false,
errors: Vec::new(),
error: Some(DiskError::other(format!("decode FileInfoVersions failed: {err}")).into()),
item_errors: Vec::new(),
}));
}
};
@@ -587,26 +563,30 @@ impl NodeService {
success: false,
errors: Vec::new(),
error: Some(DiskError::other(format!("decode DeleteOptions failed: {err}")).into()),
item_errors: Vec::new(),
}));
}
};
let (errors, item_errors) =
encode_delete_versions_errors(disk.delete_versions(&request.volume, versions, opts).await);
let errors = disk
.delete_versions(&request.volume, versions, opts)
.await
.into_iter()
.map(|error| match error {
Some(e) => e.to_string(),
None => "".to_string(),
})
.collect();
Ok(Response::new(DeleteVersionsResponse {
success: true,
errors,
error: None,
item_errors,
}))
} else {
Ok(Response::new(DeleteVersionsResponse {
success: false,
errors: Vec::new(),
error: Some(DiskError::other("cannot find disk".to_string()).into()),
item_errors: Vec::new(),
}))
}
}
@@ -1632,8 +1612,8 @@ impl NodeService {
mod tests {
use super::{
compat_response_json, decode_msgpack_or_json, decode_rename_data_request_file_info,
encode_batch_read_version_response_payloads, encode_delete_versions_errors, encode_file_info_msgpack, encode_msgpack,
encode_msgpack_named, encode_read_multiple_response_payloads, encode_rename_data_response_payloads,
encode_batch_read_version_response_payloads, encode_file_info_msgpack, encode_msgpack, encode_msgpack_named,
encode_read_multiple_response_payloads, encode_rename_data_response_payloads,
};
use crate::storage::rpc::node_service::make_server;
use crate::storage::storage_api::ReadMultipleResp;
@@ -1652,18 +1632,6 @@ mod tests {
count: u32,
}
#[test]
fn delete_versions_response_dual_writes_typed_item_errors() {
let raw_not_found = super::DiskError::Io(std::io::Error::from(std::io::ErrorKind::NotFound));
let (errors, item_errors) = encode_delete_versions_errors(vec![Some(raw_not_found), None]);
assert!(errors[0].starts_with("io error "));
assert!(errors[1].is_empty());
assert_eq!(item_errors[0].code, super::DiskError::FileNotFound.to_u32());
assert_eq!(item_errors[0].error_info, errors[0]);
assert_eq!(item_errors[1].code, 0);
}
#[tokio::test]
#[serial]
async fn handle_read_version_records_attribution_for_missing_disk() {