fix(heal): retire stale delete markers after bucket recreation (#7743)

* fix(heal): prove completed historical version cleanup

* test(heal): use debug runtime stack for C06 regression

* fix(heal): retire stale delete markers after bucket recreation

* test(ecstore): fix Clippy in retired marker regressions
This commit is contained in:
cxymds
2026-09-13 21:26:11 +08:00
committed by GitHub
parent 17ddecb075
commit 244e7dfb99
31 changed files with 1508 additions and 24 deletions
@@ -583,6 +583,13 @@ impl NodeService for MinimalLockNodeService {
Err(Status::unimplemented("lock-only test server"))
}
async fn delete_retired_marker(
&self,
_request: Request<rustfs_protos::proto_gen::node_service::DeleteVersionRequest>,
) -> Result<Response<rustfs_protos::proto_gen::node_service::DeleteVersionResponse>, Status> {
Err(Status::unimplemented("lock-only test server"))
}
async fn delete_versions(
&self,
_request: Request<rustfs_protos::proto_gen::node_service::DeleteVersionsRequest>,
+1
View File
@@ -30,6 +30,7 @@ pub mod policy_sys;
pub mod quota;
pub mod remote_s3_client;
pub mod replication;
pub(crate) mod retirement;
pub mod sealed_credentials;
pub mod tagging;
pub mod target;
+89
View File
@@ -0,0 +1,89 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! Immutable evidence of a completed bucket-generation deletion. These objects
//! deliberately live outside the bucket metadata cleanup prefix. Offline disks
//! have no bounded return time, so retirement evidence must not expire by age.
use crate::config::com::{read_config_limited_preserve_empty, save_config_with_opts};
use crate::error::{Error, Result};
use crate::object_api::ObjectOptions;
use crate::store::ECStore;
use serde::{Deserialize, Serialize};
use std::sync::Arc;
use uuid::Uuid;
pub(crate) struct MarkerRetirementContext<'a> {
pub store: Option<Arc<ECStore>>,
pub current_incarnation: Option<Uuid>,
pub lifecycle_guard: &'a rustfs_lock::NamespaceLockGuard,
}
#[derive(Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct RetirementRecord {
version: u8,
deployment_id: Uuid,
bucket: String,
incarnation_id: Uuid,
}
fn expected_record(store: &ECStore, bucket: &str, incarnation_id: Uuid) -> Result<RetirementRecord> {
let deployment_id = store
.pools
.first()
.ok_or_else(|| Error::other("missing retirement deployment"))?
.format
.id;
if deployment_id.is_nil() || incarnation_id.is_nil() {
return Err(Error::other("retirement requires non-nil deployment and incarnation IDs"));
}
Ok(RetirementRecord {
version: 1,
deployment_id,
bucket: bucket.to_owned(),
incarnation_id,
})
}
fn record_path(bucket: &str, incarnation_id: Uuid) -> String {
format!("bucket-retirements/{bucket}/{incarnation_id}.json")
}
/// The caller owns the bucket lifecycle WRITE guard and has completed physical
/// deletion on every selected set, including the existing rollback decision.
/// A failed or uncertain publication is an error, never a successful DELETE ACK.
pub(crate) async fn commit_retirement(
store: Arc<ECStore>,
bucket: &str,
incarnation_id: Uuid,
opts: &ObjectOptions,
) -> Result<()> {
let record = expected_record(&store, bucket, incarnation_id)?;
save_config_with_opts(store, &record_path(bucket, incarnation_id), serde_json::to_vec(&record)?, opts).await
}
pub(crate) async fn is_retired(store: Arc<ECStore>, bucket: &str, incarnation_id: Uuid) -> Result<bool> {
let expected = expected_record(&store, bucket, incarnation_id)?;
let bytes = match read_config_limited_preserve_empty(store, &record_path(bucket, incarnation_id), 4096).await {
Ok(bytes) => bytes,
Err(Error::ConfigNotFound) => return Ok(false),
Err(error) => return Err(error),
};
let record: RetirementRecord = serde_json::from_slice(&bytes)?;
if record != expected {
return Err(Error::other("bucket retirement identity mismatch"));
}
Ok(true)
}
@@ -222,6 +222,8 @@ pub struct ListPathRawOptions {
pub filter_prefix: Option<String>,
pub forward_to: Option<String>,
pub min_disks: usize,
/// Deliver every replica to the partial callback, including matching headers.
pub preserve_replica_metadata: bool,
pub report_not_found: bool,
pub per_disk_limit: i32,
pub skip_walkdir_total_timeout: bool,
@@ -254,6 +256,7 @@ impl Clone for ListPathRawOptions {
filter_prefix: self.filter_prefix.clone(),
forward_to: self.forward_to.clone(),
min_disks: self.min_disks,
preserve_replica_metadata: self.preserve_replica_metadata,
report_not_found: self.report_not_found,
per_disk_limit: self.per_disk_limit,
skip_walkdir_total_timeout: self.skip_walkdir_total_timeout,
@@ -918,7 +921,7 @@ async fn list_path_raw_inner(
break;
}
if agree == readers.len() {
if agree == readers.len() && !opts.preserve_replica_metadata {
for r in readers.iter_mut() {
let _ = r.skip(1).await;
}
+52 -1
View File
@@ -2504,6 +2504,7 @@ impl DiskAPI for RemoteDisk {
// JSON + msgpack until its fallback counter has read zero across a release window.
let file_info_bin = encode_file_info_msgpack(&fi)?;
let opts_bin = encode_msgpack(&opts)?;
let conditional_marker = opts.expected_delete_marker.is_some();
let file_info = serde_json::to_string(&fi)?;
let opts = serde_json::to_string(&opts)?;
@@ -2521,7 +2522,14 @@ impl DiskAPI for RemoteDisk {
let canonical_body = rustfs_protos::canonical_delete_version_request_body(request.get_ref());
attach_mutation_body_digest(&mut request, canonical_body, "delete_version")?;
let response = client.delete_version(request).await?.into_inner();
// Unknown RPC methods fail closed on older peers. Never retry a
// conditional delete through the legacy unconditional method.
let response = if conditional_marker {
client.delete_retired_marker(request).await?
} else {
client.delete_version(request).await?
}
.into_inner();
if !response.success {
return Err(response.error.unwrap_or_default().into());
@@ -3908,6 +3916,49 @@ mod tests {
}
}
#[test]
fn retired_marker_options_preserve_legacy_positional_wire_shape() {
#[derive(Debug, serde::Serialize, serde::Deserialize, PartialEq)]
struct LegacyDeleteOptions {
recursive: bool,
immediate: bool,
undo_write: bool,
undo_delete: bool,
old_data_dir: Option<Uuid>,
}
let legacy = LegacyDeleteOptions {
recursive: false,
immediate: false,
undo_write: false,
undo_delete: false,
old_data_dir: None,
};
let original = encode_msgpack(&legacy).unwrap();
let current = encode_msgpack(&DeleteOptions::default()).unwrap();
assert_eq!(current, original, "ordinary deletes must retain the older peer's positional payload");
assert_eq!(rmp_serde::from_slice::<LegacyDeleteOptions>(&current).unwrap(), legacy);
assert!(
rmp_serde::from_slice::<DeleteOptions>(&original)
.unwrap()
.expected_delete_marker
.is_none()
);
let mut marker = FileInfo {
deleted: true,
version_id: Some(Uuid::new_v4()),
mod_time: Some(::time::OffsetDateTime::now_utc()),
..Default::default()
};
marker.set_delete_marker_incarnation(Uuid::new_v4());
let marker = rustfs_filemeta::MetaDeleteMarker::from(marker);
let options = DeleteOptions {
expected_delete_marker: Some(marker.clone()),
..Default::default()
};
let decoded: DeleteOptions = rmp_serde::from_slice(&encode_msgpack(&options).unwrap()).unwrap();
assert_eq!(decoded.expected_delete_marker, Some(marker));
}
#[test]
fn delete_versions_response_preserves_typed_item_errors() {
let errors = decode_delete_versions_errors(
+2 -1
View File
@@ -1483,11 +1483,12 @@ impl Sets {
object: &str,
version_id: &str,
opts: &HealOpts,
retirement: Option<&crate::bucket::retirement::MarkerRetirementContext<'_>>,
) -> Result<(HealResultItem, Option<Error>, Option<crate::set_disk::HealedObjectAbsence>)> {
let mut absence = None;
let (item, error) = self
.get_disks_for_heal_object(object, opts)?
.heal_object_with_absence(bucket, object, version_id, opts, &mut absence)
.heal_object_with_retirement(bucket, object, version_id, opts, &mut absence, retirement)
.await?;
// A caller-owned lock does not expose its lease to this boundary.
// Keep cleanup unverified when that lease cannot be checked here.
+1
View File
@@ -1597,6 +1597,7 @@ impl LocalDiskWrapper {
undo_write: false,
undo_delete: false,
old_data_dir: None,
expected_delete_marker: None,
},
)
.await?;
+31
View File
@@ -41,6 +41,10 @@ struct DanglingDeleteGraceError {
grace_secs: i64,
}
#[derive(Debug, Clone, thiserror::Error)]
#[error("retired delete marker cleanup deferred: {0}")]
struct RetiredMarkerDeferred(String);
/// Marks a conditional-file write that failed before its publication rename.
/// Callers may choose another owner only while this marker is preserved; every
/// unmarked error remains commit-ambiguous and must fail closed.
@@ -334,6 +338,19 @@ impl DiskError {
Some(io::Error::new(error.kind(), grace.clone()))
}
pub(crate) fn retired_marker_deferred(reason: impl Into<String>) -> Self {
Self::other(RetiredMarkerDeferred(reason.into()))
}
pub fn io_error_is_retired_marker_deferred(error: &io::Error) -> bool {
error.get_ref().is_some_and(|source| source.is::<RetiredMarkerDeferred>())
}
pub(crate) fn clone_retired_marker_deferred(error: &io::Error) -> Option<io::Error> {
let deferred = error.get_ref()?.downcast_ref::<RetiredMarkerDeferred>()?;
Some(io::Error::other(deferred.clone()))
}
pub fn dangling_delete_retry_after(&self) -> Option<std::time::Duration> {
match self {
Self::Io(error) => Self::io_error_dangling_delete_retry_after(error),
@@ -685,6 +702,7 @@ impl Clone for DiskError {
),
DiskError::Io(io_error) => DiskError::Io(
Self::clone_dangling_delete_grace(io_error)
.or_else(|| Self::clone_retired_marker_deferred(io_error))
.or_else(|| rustfs_rio::clone_internode_http_io_error(io_error))
.and_then(std::io::Error::into_inner)
// The helper derives a kind from the source; Clone must retain the original outer kind.
@@ -877,6 +895,19 @@ mod tests {
use super::*;
use std::collections::HashMap;
#[test]
fn retired_marker_deferral_survives_disk_and_storage_clones() {
let original = DiskError::retired_marker_deferred("missing retirement record");
let disk = original.clone();
let storage = crate::error::StorageError::from(disk);
let cloned = storage.clone();
assert!(cloned.is_retired_marker_deferred());
assert!(crate::error::StorageError::from(original).is_retired_marker_deferred());
assert!(storage.is_retired_marker_deferred());
assert!(!crate::error::StorageError::other(storage.to_string()).is_retired_marker_deferred());
assert!(!crate::error::StorageError::FileVersionNotFound.is_retired_marker_deferred());
}
#[test]
fn dangling_grace_retry_timing_survives_disk_and_storage_clones() {
let original = super::DiskError::dangling_delete_grace(21, 3600);
+145
View File
@@ -5828,6 +5828,25 @@ impl LocalDisk {
opts,
namespace_owner,
} = mutation;
if let Some(expected) = opts.expected_delete_marker.as_ref() {
let expected_info = expected.into_fileinfo(volume, path, false)?;
if force_del_marker
|| opts.recursive
|| opts.immediate
|| opts.undo_write
|| opts.undo_delete
|| opts.old_data_dir.is_some()
|| path.starts_with(SLASH_SEPARATOR)
|| fi.deleted
|| fi.mark_deleted
|| fi.version_id.is_none_or(|id| id.is_nil())
|| fi.version_id != expected.version_id
|| expected_info.delete_marker_incarnation().is_none()
|| !expected_info.is_canonical_delete_marker()
{
return Err(DiskError::FileCorrupt);
}
}
if path.starts_with(SLASH_SEPARATOR) {
return self
.delete_with_namespace_owner(
@@ -5879,6 +5898,23 @@ impl LocalDisk {
};
let mut meta = FileMeta::load(&buf)?;
if let Some(expected) = opts.expected_delete_marker.as_ref() {
let Some(version) = meta
.versions
.iter()
.find(|version| version.header.version_id == fi.version_id)
else {
return Err(DiskError::FileVersionNotFound);
};
let actual = version.parse_version_meta()?;
if actual.version_type != rustfs_filemeta::VersionType::Delete
|| actual.object.is_some()
|| actual.legacy_object.is_some()
|| actual.delete_marker.as_ref() != Some(expected)
{
return Err(DiskError::FileCorrupt);
}
}
let old_dir = meta.delete_version(&fi)?;
let mut reserved_version_delete = false;
if let Some(rollback_dir) = rollback_dir {
@@ -16384,6 +16420,114 @@ mod test {
}
}
#[tokio::test]
async fn retired_marker_condition_preserves_replacements_and_other_versions() {
let dir = tempfile::tempdir().unwrap();
let endpoint = Endpoint::try_from(dir.path().to_str().unwrap()).unwrap();
let disk = Arc::new(LocalDisk::new(&endpoint, false).await.unwrap());
let bucket = "retired-marker-condition";
let object = "reused.bin";
ensure_test_volume(&disk, bucket).await;
let version = Uuid::new_v4();
let old_incarnation = Uuid::new_v4();
let mut old = FileInfo {
name: object.into(),
version_id: Some(version),
deleted: true,
mod_time: Some(time::OffsetDateTime::now_utc()),
..Default::default()
};
old.set_delete_marker_incarnation(old_incarnation);
let condition = rustfs_filemeta::MetaDeleteMarker::from(old.clone());
let request = FileInfo {
name: object.into(),
version_id: Some(version),
..Default::default()
};
let options = DeleteOptions {
expected_delete_marker: Some(condition),
..Default::default()
};
let mut unrelated =
test_file_info(object, Uuid::new_v4(), Some(Uuid::new_v4()), Some(Bytes::from_static(b"current bytes")));
unrelated.set_inline_data();
disk.write_metadata("", bucket, object, unrelated.clone()).await.unwrap();
disk.write_metadata("", bucket, object, old.clone()).await.unwrap();
for replacement in [
{
let mut current = old.clone();
current.set_delete_marker_incarnation(Uuid::new_v4());
current
},
{
let mut changed = old.clone();
changed.mod_time = Some(changed.mod_time.unwrap() + time::Duration::seconds(1));
changed
},
test_file_info(object, version, Some(Uuid::new_v4()), Some(Bytes::from_static(b"replacement data"))),
] {
// The request was formed from the old observation. Publication must
// still compare against the metadata present when it obtains the lease.
disk.write_metadata("", bucket, object, replacement).await.unwrap();
let path = dir.path().join(bucket).join(object).join(STORAGE_FORMAT_FILE);
let before = fs::read(&path).await.unwrap();
assert!(
disk.delete_version(bucket, object, request.clone(), false, options.clone())
.await
.is_err()
);
assert_eq!(
fs::read(&path).await.unwrap(),
before,
"a stale condition must not publish any metadata change"
);
}
disk.write_metadata("", bucket, object, old.clone()).await.unwrap();
let lease = os::acquire_metadata_mutation_lease(&disk.get_object_path(bucket, object).unwrap(), None).await;
let pending = tokio::spawn({
let disk = disk.clone();
let request = request.clone();
let options = options.clone();
async move { disk.delete_version(bucket, object, request, false, options).await }
});
// Publish a competing generation while owning the actual metadata
// lease. The delayed delete must read this replacement after release.
let mut replacement = old.clone();
replacement.set_delete_marker_incarnation(Uuid::new_v4());
disk.write_metadata_with_namespace_owner(bucket, object, replacement, Some(lease.clone()))
.await
.unwrap();
let path = dir.path().join(bucket).join(object).join(STORAGE_FORMAT_FILE);
let after_replacement = fs::read(&path).await.unwrap();
assert!(!pending.is_finished(), "conditional deletion must wait for the metadata lease");
drop(lease);
assert!(pending.await.unwrap().is_err());
assert_eq!(fs::read(&path).await.unwrap(), after_replacement);
disk.write_metadata("", bucket, object, old).await.unwrap();
disk.delete_version(bucket, object, request, false, options)
.await
.expect("matching retired marker condition");
assert!(matches!(
disk.read_version("", bucket, object, &version.to_string(), &ReadOptions::default())
.await,
Err(DiskError::FileVersionNotFound)
));
let retained = disk
.read_version(
"",
bucket,
object,
&unrelated.version_id.unwrap().to_string(),
&ReadOptions {
read_data: true,
..Default::default()
},
)
.await
.expect("unrelated new version remains");
assert_eq!(retained.data, unrelated.data);
}
#[tokio::test]
async fn test_delete_version_undo_restores_backup_to_object_root() {
use tempfile::tempdir;
@@ -19659,6 +19803,7 @@ mod test {
undo_write: false,
undo_delete: false,
old_data_dir: None,
expected_delete_marker: None,
};
disk.delete("test-volume", "test-file.txt", delete_opts)
.await
+5
View File
@@ -1518,6 +1518,10 @@ pub struct DeleteOptions {
#[serde(default)]
pub undo_delete: bool,
pub old_data_dir: Option<Uuid>,
/// Full marker precondition checked under the actual metadata mutation lease.
/// Remote calls carrying it must use DeleteRetiredMarker, never DeleteVersion.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub expected_delete_marker: Option<rustfs_filemeta::MetaDeleteMarker>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -1806,6 +1810,7 @@ mod tests {
undo_write: true,
undo_delete: false,
old_data_dir: Some(Uuid::new_v4()),
expected_delete_marker: None,
};
assert!(opts.recursive);
+11 -1
View File
@@ -406,6 +406,14 @@ impl StorageError {
pub fn is_dangling_delete_grace(&self) -> bool {
matches!(self, StorageError::Io(io_error) if DiskError::io_error_is_dangling_delete_grace(io_error))
}
pub fn is_retired_marker_deferred(&self) -> bool {
matches!(self, StorageError::Io(error) if DiskError::io_error_is_retired_marker_deferred(error))
}
pub fn retired_marker_deferred(reason: impl Into<String>) -> Self {
DiskError::retired_marker_deferred(reason).into()
}
}
impl From<HTTPRangeError> for StorageError {
@@ -615,7 +623,9 @@ impl Clone for StorageError {
fn clone(&self) -> Self {
match self {
StorageError::Io(e) => {
if let Some(error) = DiskError::clone_dangling_delete_grace(e) {
if let Some(error) =
DiskError::clone_dangling_delete_grace(e).or_else(|| DiskError::clone_retired_marker_deferred(e))
{
return StorageError::Io(error);
}
if let Some(context) = self.pool_metadata_failure() {
+108
View File
@@ -34,6 +34,114 @@ struct FileInfoIdentityGroup {
}
impl SetDisks {
/// Resolve each listed version using the same metadata authority as an exact read.
pub(crate) fn resolve_listed_versions(
bucket: &str,
entries: rustfs_filemeta::MetaCacheEntries,
read_errors: &[Option<DiskError>],
disk_count: usize,
default_parity: usize,
) -> disk::error::Result<Option<rustfs_filemeta::MetaCacheEntry>> {
use rustfs_filemeta::{FileMeta, MetaCacheEntry};
if disk_count == 0 || entries.0.len() > disk_count || entries.0.len() != read_errors.len() {
return Err(DiskError::ErasureReadQuorum);
}
let Some(first) = entries.0.iter().flatten().next() else {
return Ok(None);
};
let name = first.name.clone();
if first.is_dir() {
let quorum = if default_parity == 0 {
disk_count
} else {
disk_count.div_ceil(2)
};
let matches = entries
.0
.iter()
.flatten()
.filter(|entry| entry.is_dir() && entry.name == name)
.count();
return Ok((matches >= quorum).then(|| first.clone()));
}
let mut versions = HashMap::<Option<Uuid>, Vec<FileInfo>>::new();
let mut missing = vec![Some(DiskError::DiskNotFound); disk_count];
for (slot, (entry, error)) in entries.0.into_iter().zip(read_errors).enumerate() {
missing[slot] = error.clone().or(Some(DiskError::FileVersionNotFound));
let Some(entry) = entry else { continue };
if error.is_some() {
continue;
}
if entry.is_dir() || entry.name != name {
missing[slot] = Some(DiskError::FileCorrupt);
continue;
}
let decoded = match entry.file_info_versions_with_free_versions(bucket) {
Ok(decoded) => decoded,
Err(error) => {
missing[slot] = Some(error.into());
continue;
}
};
let decoded = decoded.versions.into_iter().chain(decoded.free_versions).collect::<Vec<_>>();
let mut seen = HashSet::new();
if decoded
.iter()
.any(|info| !file_info_is_valid_for_metadata(info) || !seen.insert(info.version_id))
{
missing[slot] = Some(DiskError::FileCorrupt);
continue;
}
for info in decoded {
let version_id = info.version_id;
versions
.entry(version_id)
.or_insert_with(|| vec![FileInfo::default(); disk_count])[slot] = info;
}
}
let complete = missing.iter().all(|error| {
matches!(
error,
Some(DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound)
)
});
let mut metadata = FileMeta::new();
for parts in versions.into_values() {
let errors = parts
.iter()
.zip(&missing)
.map(|(info, error)| {
if file_info_is_valid_for_metadata(info) {
None
} else {
error.clone()
}
})
.collect::<Vec<_>>();
let read_quorum = match Self::object_quorum_from_meta(&parts, &errors, default_parity) {
Ok((quorum, _)) => usize::try_from(quorum).map_err(|_| DiskError::FileCorrupt)?,
Err(DiskError::FileNotFound | DiskError::FileVersionNotFound) if complete => continue,
Err(DiskError::FileNotFound | DiskError::FileVersionNotFound) => return Err(DiskError::ErasureReadQuorum),
Err(DiskError::ErasureReadQuorum) if complete => continue,
Err(error) => return Err(error),
};
let mod_time = Self::common_time(&Self::list_object_modtimes(&parts, &errors), read_quorum);
let selected = Self::pick_valid_fileinfo(&parts, mod_time, None, read_quorum)?;
metadata.add_version(selected)?;
}
if metadata.versions.is_empty() {
return Ok(None);
}
Ok(Some(MetaCacheEntry {
name,
metadata: metadata.marshal_msg()?,
cached: Some(metadata),
reusable: false,
}))
}
pub(super) fn all_not_found_metadata(errs: &[Option<DiskError>]) -> bool {
!errs.is_empty()
&& errs.iter().all(|err| match err {
+2
View File
@@ -866,6 +866,8 @@ mod ctx;
mod metadata;
mod ops;
pub(crate) use ops::bucket::BucketInfoQuorum;
#[cfg(test)]
pub(crate) use ops::heal::DanglingDeleteFailure;
pub(crate) use ops::heal::HealedObjectAbsence;
#[cfg(test)]
+213 -10
View File
@@ -22,6 +22,15 @@ use super::super::{
heal_bucket_local_on_disks, is_object_dir_dangling, join_all, load_format_erasure_all, path_join_buf, save_format_file,
should_heal_object_on_disk, stat_all_dirs, to_object_err, warn,
};
use crate::bucket::retirement::MarkerRetirementContext;
#[derive(Clone, Copy)]
struct ExplicitVersionHeal<'a> {
opts: &'a HealOpts,
allow_regeneration: bool,
retirement: Option<&'a MarkerRetirementContext<'a>>,
}
use crate::disk::DataDirDeleteStatus;
use crate::disk::DiskAPI;
use crate::disk::local::{DELETE_DATA_DIR_MARKER_PREFIX, metadata_less_part_file};
@@ -359,7 +368,7 @@ fn injected_dangling_check_parts_error(bucket: &str, object: &str, disk_index: u
}
#[cfg(test)]
struct DanglingDeleteFailure {
pub(crate) struct DanglingDeleteFailure {
key: DanglingDeleteFailureKey,
}
@@ -377,7 +386,7 @@ fn dangling_delete_failures() -> &'static std::sync::Mutex<DanglingDeleteFailure
#[cfg(test)]
impl DanglingDeleteFailure {
fn install(bucket: &str, object: &str, disk_index: usize, error: DiskError) -> Self {
pub(crate) fn install(bucket: &str, object: &str, disk_index: usize, error: DiskError) -> Self {
let key = (bucket.to_string(), object.to_string(), disk_index);
let previous = dangling_delete_failures()
.lock()
@@ -558,7 +567,18 @@ impl SetDisks {
version_id: &str,
opts: &HealOpts,
) -> disk::error::Result<(HealResultItem, Option<DiskError>)> {
Box::pin(self.heal_object_with_explicit_version_regen(bucket, object, version_id, opts, true, &mut None)).await
Box::pin(self.heal_object_with_explicit_version_regen(
bucket,
object,
version_id,
ExplicitVersionHeal {
opts,
allow_regeneration: true,
retirement: None,
},
&mut None,
))
.await
}
async fn read_repair_commit_fingerprint(
@@ -668,10 +688,14 @@ impl SetDisks {
bucket: &str,
object: &str,
version_id: &str,
opts: &HealOpts,
allow_explicit_version_regen: bool,
run: ExplicitVersionHeal<'_>,
absence: &mut Option<HealedObjectAbsence>,
) -> disk::error::Result<(HealResultItem, Option<DiskError>)> {
let ExplicitVersionHeal {
opts,
allow_regeneration: allow_explicit_version_regen,
retirement,
} = run;
trace!(
event = EVENT_SET_DISK_HEAL,
component = LOG_COMPONENT_ECSTORE,
@@ -1663,15 +1687,48 @@ impl SetDisks {
}
}
Err(err) => {
if let Some(retirement) = retirement
&& parts_metadata
.iter()
.zip(&errs)
.any(|(info, error)| error.is_none() && info.is_canonical_delete_marker())
{
let cleanup = match Self::retired_marker_candidate(&parts_metadata, &errs) {
Ok(marker) => {
self.remove_retired_marker(retirement, bucket, object, marker, opts, &disks)
.await
}
Err(error) => Err(error),
};
let absent = cleanup.is_ok();
let item = self
.dangling_heal_result(FileInfo::default(), &errs, bucket, object, version_id, absent)
.await;
if absent {
*absence = Some(HealedObjectAbsence {
pool_index: self.pool_index,
set_index: self.set_index,
removed: true,
});
}
return Ok((item, cleanup.err()));
}
if allow_explicit_version_regen
&& !version_id.is_empty()
&& self
.try_regenerate_explicit_version_meta(bucket, object, version_id, &parts_metadata, &errs, &disks)
.await?
{
return Box::pin(
self.heal_object_with_explicit_version_regen(bucket, object, version_id, opts, false, absence),
)
return Box::pin(self.heal_object_with_explicit_version_regen(
bucket,
object,
version_id,
ExplicitVersionHeal {
allow_regeneration: false,
..run
},
absence,
))
.await;
}
@@ -1733,6 +1790,119 @@ impl SetDisks {
}
}
fn retired_marker_candidate(
metadata: &[FileInfo],
errors: &[Option<DiskError>],
) -> disk::error::Result<rustfs_filemeta::MetaDeleteMarker> {
let mut candidate = None;
for (info, error) in metadata.iter().zip(errors) {
match error {
Some(DiskError::FileNotFound | DiskError::FileVersionNotFound) => continue,
Some(_) => return Err(DiskError::retired_marker_deferred("not every replica is readable")),
None => {}
}
if !info.is_canonical_delete_marker() || info.delete_marker_incarnation().is_none() {
return Err(DiskError::retired_marker_deferred("marker has no trustworthy bucket incarnation"));
}
let marker = rustfs_filemeta::MetaDeleteMarker::from(info.clone());
if candidate.as_ref().is_some_and(|previous| previous != &marker) {
return Err(DiskError::retired_marker_deferred("surviving marker identities conflict"));
}
candidate = Some(marker);
}
candidate.ok_or_else(|| DiskError::retired_marker_deferred("no exact marker candidate"))
}
async fn remove_retired_marker(
&self,
context: &MarkerRetirementContext<'_>,
bucket: &str,
object: &str,
marker: rustfs_filemeta::MetaDeleteMarker,
opts: &HealOpts,
disks: &[Option<DiskStore>],
) -> disk::error::Result<()> {
let defer = DiskError::retired_marker_deferred;
if opts.dry_run || !opts.remove {
return Err(defer("cleanup requires remove=true and dry_run=false"));
}
if disks.len() != self.set_drive_count || disks.iter().any(Option::is_none) {
return Err(defer("every target disk must be online"));
}
let incarnation = marker.into_fileinfo(bucket, object, false)?.delete_marker_incarnation();
let Some(old) = incarnation else { return Err(defer("marker incarnation is unavailable")) };
let Some(current) = context.current_incarnation else {
return Err(defer("current bucket incarnation is unavailable"));
};
if old == current || context.lifecycle_guard.is_lock_lost() {
return Err(defer("current generation or lost bucket lifecycle fence"));
}
let Some(store) = context.store.as_ref() else {
return Err(defer("retirement owner is unavailable"));
};
match crate::bucket::retirement::is_retired(store.clone(), bucket, old).await {
Ok(true) => {}
Ok(false) => return Err(defer("no committed bucket retirement record")),
Err(error) => return Err(DiskError::retired_marker_deferred(format!("retirement read failed: {error}"))),
}
let Some(version) = marker.version_id.filter(|id| !id.is_nil()) else {
return Err(defer("cleanup requires an explicit non-null version"));
};
if context.lifecycle_guard.is_lock_lost() {
return Err(defer("bucket lifecycle fence was lost"));
}
let request = FileInfo {
volume: bucket.to_owned(),
name: object.to_owned(),
version_id: Some(version),
..Default::default()
};
let results = join_all(disks.iter().enumerate().map(|(disk_index, disk)| {
let request = request.clone();
let marker = marker.clone();
async move {
if let Some(error) = injected_dangling_delete_error(bucket, object, disk_index) {
return Err(error);
}
let Some(disk) = disk else { return Err(DiskError::DiskNotFound) };
disk.delete_version(
bucket,
object,
request,
false,
DeleteOptions {
expected_delete_marker: Some(marker),
..Default::default()
},
)
.await
}
}))
.await;
for result in results {
if let Err(error) = result
&& !matches!(error, DiskError::FileNotFound | DiskError::FileVersionNotFound)
{
return Err(DiskError::retired_marker_deferred(format!("conditional marker deletion failed: {error}")));
}
}
let (_, errors) = Self::read_all_fileinfo(disks, "", bucket, object, &version.to_string(), false, false, false).await?;
let selected = self.get_disks_internal().await;
let same_targets = selected.len() == disks.len() && selected.iter().zip(disks).all(|(current, original)| {
matches!((current, original), (Some(current), Some(original)) if std::sync::Arc::ptr_eq(current, original))
});
if context.lifecycle_guard.is_lock_lost()
|| !same_targets
|| !errors
.iter()
.all(|error| matches!(error, Some(DiskError::FileNotFound | DiskError::FileVersionNotFound)))
{
return Err(defer("complete marker absence could not be verified under the original fence"));
}
self.invalidate_get_object_metadata_cache(bucket, object).await;
Ok(())
}
async fn try_regenerate_explicit_version_meta(
&self,
bucket: &str,
@@ -2709,6 +2879,19 @@ impl SetDisks {
version_id: &str,
opts: &HealOpts,
absence: &mut Option<HealedObjectAbsence>,
) -> Result<(HealResultItem, Option<Error>)> {
self.heal_object_with_retirement(bucket, object, version_id, opts, absence, None)
.await
}
pub(crate) async fn heal_object_with_retirement(
&self,
bucket: &str,
object: &str,
version_id: &str,
opts: &HealOpts,
absence: &mut Option<HealedObjectAbsence>,
retirement: Option<&MarkerRetirementContext<'_>>,
) -> Result<(HealResultItem, Option<Error>)> {
*absence = None;
let _write_lock_guard = if !opts.no_lock {
@@ -2796,7 +2979,17 @@ impl SetDisks {
let mut inner_opts = *opts;
inner_opts.no_lock = true;
let (result, err) = self
.heal_object_with_explicit_version_regen(bucket, object, version_id, &inner_opts, true, absence)
.heal_object_with_explicit_version_regen(
bucket,
object,
version_id,
ExplicitVersionHeal {
opts: &inner_opts,
allow_regeneration: true,
retirement,
},
absence,
)
.await
.map_err(|e| to_object_err(e.into(), vec![bucket, object]))?;
if let Some(err) = err.as_ref() {
@@ -2806,7 +2999,17 @@ impl SetDisks {
// during a normal heal scan, heal again with bitrot flag enabled.
inner_opts.scan_mode = HealScanMode::Deep;
let (result, err) = self
.heal_object_with_explicit_version_regen(bucket, object, version_id, &inner_opts, true, absence)
.heal_object_with_explicit_version_regen(
bucket,
object,
version_id,
ExplicitVersionHeal {
opts: &inner_opts,
allow_regeneration: true,
retirement,
},
absence,
)
.await
.map_err(|e| to_object_err(e.into(), vec![bucket, object]))?;
if _write_lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost()) {
+12
View File
@@ -8154,6 +8154,10 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
vr.mod_time = goi.mod_time;
}
if let Some(incarnation) = opts.expected_bucket_incarnation_id {
vr.set_delete_marker_incarnation(incarnation);
}
let v = {
if vers_map.contains_key(&dobj.object_name) {
let val = vers_map.get_mut(&dobj.object_name).unwrap();
@@ -8835,6 +8839,10 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
fi.set_tier_free_version_id(&find_vid.to_string());
if let Some(incarnation) = opts.expected_bucket_incarnation_id {
fi.set_delete_marker_incarnation(incarnation);
}
fi.version_id = if let Some(vid) = opts.version_id.as_ref() {
let vid = Uuid::parse_str(vid.as_str())?;
(!opts.version_suspended || !vid.is_nil()).then_some(vid)
@@ -8894,6 +8902,10 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
dfi.set_tier_free_version_id(&find_vid.to_string());
if let Some(incarnation) = opts.expected_bucket_incarnation_id {
dfi.set_delete_marker_incarnation(incarnation);
}
ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?;
begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?;
let _tier_delete_lease = acquire_single_tier_delete_lease(&opts, &goi).await?;
+35
View File
@@ -1077,6 +1077,20 @@ impl ECStore {
.await?;
}
// Capture the authoritative old identity before its namespace disappears.
// Legacy buckets without a stamp cannot produce retirement authority.
let retirement = if bucket_exists && bucket_lifecycle_guard.is_some() && !is_meta_bucketname(bucket) {
if let Some(store) = metadata_sys::object_store_if_initialized_in(&self.ctx).await {
crate::bucket::metadata::load_bucket_incarnation(store.clone(), bucket)
.await?
.map(|incarnation| (store, incarnation))
} else {
None
}
} else {
None
};
let delete_result = await_bucket_namespace_operation(
bucket_lifecycle_guard.as_ref(),
bucket,
@@ -1100,6 +1114,27 @@ impl ECStore {
return Err(err);
}
if let Some((store, incarnation)) = retirement {
let mut record_opts = ObjectOptions {
max_parity: true,
..Default::default()
};
if let Some(guard) = bucket_lifecycle_guard.as_ref() {
record_opts.add_bucket_lifecycle_lock_guard(guard);
}
if let Some(guard) = ns_guard.as_ref() {
record_opts.add_namespace_lock_guard(guard);
}
await_bucket_lifecycle_operation(
bucket_lifecycle_guard.as_ref(),
ns_guard.as_ref(),
bucket,
"bucket retirement publication",
crate::bucket::retirement::commit_retirement(store, bucket, incarnation, &record_opts),
)
.await?;
}
self.cleanup_bucket_usage_best_effort(bucket, ns_guard.as_ref()).await;
if let Err(err) = self
+16 -4
View File
@@ -484,7 +484,7 @@ impl ECStore {
version_id: &str,
opts: &HealOpts,
) -> Result<(HealResultItem, Option<Error>)> {
self.handle_heal_object_with_absence(bucket, object, version_id, opts, &mut None)
self.handle_heal_object_with_absence(bucket, object, version_id, opts, &mut None, None)
.await
}
@@ -507,9 +507,18 @@ impl ECStore {
// Match object publication: bucket lifecycle before capacity and object
// namespace locks. Keep the incarnation pinned through proof delivery.
let guard = self.acquire_bucket_lifecycle_read_lock(bucket).await?;
let retirement = crate::bucket::retirement::MarkerRetirementContext {
store: crate::bucket::metadata_sys::object_store_if_initialized_in(&self.ctx).await,
current_incarnation: self
.bucket_incarnation_id_from_disk(bucket)
.await
.ok()
.filter(|id| !id.is_nil()),
lifecycle_guard: &guard,
};
let mut proofs = None;
let (item, mut error) = self
.handle_heal_object_with_absence(bucket, object, version_id, opts, &mut proofs)
.handle_heal_object_with_absence(bucket, object, version_id, opts, &mut proofs, Some(&retirement))
.await?;
// Read the authoritative incarnation only for an absence candidate.
// The lifecycle guard has pinned it throughout the storage operation.
@@ -547,6 +556,7 @@ impl ECStore {
version_id: &str,
opts: &HealOpts,
absence: &mut Option<Vec<crate::set_disk::HealedObjectAbsence>>,
retirement: Option<&crate::bucket::retirement::MarkerRetirementContext<'_>>,
) -> Result<(HealResultItem, Option<Error>)> {
trace!(
event = EVENT_HEAL_OBJECT_STARTED,
@@ -637,7 +647,8 @@ impl ECStore {
}
#[cfg(test)]
crate::core::pools::notify_decommission_external_heal_operation_started(store_id);
pool.heal_object_with_absence(bucket, &pool_object, version_id, &opts).await
pool.heal_object_with_absence(bucket, &pool_object, version_id, &opts, retirement)
.await
}
});
let results = join_all(futures).await;
@@ -660,7 +671,8 @@ impl ECStore {
move |opts| async move {
#[cfg(test)]
crate::core::pools::notify_decommission_external_heal_operation_started(store_id);
pool.heal_object_with_absence(bucket, &pool_object, version_id, &opts).await
pool.heal_object_with_absence(bucket, &pool_object, version_id, &opts, retirement)
.await
},
));
}
+232
View File
@@ -6786,7 +6786,74 @@ impl SetDisks {
Ok(result)
}
async fn list_versions_authoritatively(
&self,
rx: CancellationToken,
opts: ListPathOptions,
sender: Sender<MetaCacheEntry>,
) -> Result<()> {
let (disks, _, _) = self.get_online_disks_with_healing_and_info(true).await;
let disk_count = self.set_drive_count;
let parity = self.default_parity_count;
let read_quorum = if parity == 0 { disk_count } else { disk_count.div_ceil(2) };
if disk_count == 0 || disks.len() < read_quorum {
return Err(DiskError::ErasureReadQuorum.into());
}
let first_error = Arc::new(tokio::sync::Mutex::new(None));
let error_sink = Arc::clone(&first_error);
let bucket = opts.bucket.clone();
let cancel = rx.clone();
let result = list_path_raw_with_claim_tracker(
rx,
ListPathRawOptions {
disks: disks.into_iter().map(Some).collect(),
bucket: opts.bucket,
path: opts.base_dir,
recursive: opts.recursive,
incl_deleted: true,
skip_hidden_prefix_check: opts.skip_hidden_prefix_check,
filter_prefix: opts.filter_prefix,
forward_to: opts.marker,
min_disks: read_quorum,
preserve_replica_metadata: true,
// A reader's local page may be consumed by stale entries. Only
// the resolved, merged page may stop the authoritative walk.
per_disk_limit: 0,
skip_walkdir_total_timeout: true,
walkdir_timeout: opts.walkdir_timeout,
walkdir_stall_timeout: opts.walkdir_stall_timeout,
partial: Some(Box::new(move |entries, errors| {
let resolved = Self::resolve_listed_versions(&bucket, entries, errors, disk_count, parity);
let sender = sender.clone();
let cancel = cancel.clone();
let first_error = Arc::clone(&error_sink);
Box::pin(async move {
match resolved {
Ok(Some(entry)) => {
let _ = send_or_cancel(&cancel, &sender, entry).await;
}
Ok(None) => {}
Err(error) => {
first_error.lock().await.get_or_insert(error);
}
}
})
})),
..Default::default()
},
FallbackClaimTracker::default(),
)
.await;
if let Some(error) = first_error.lock().await.take() {
return Err(error.into());
}
result.map_err(Into::into)
}
pub async fn list_path(&self, rx: CancellationToken, opts: ListPathOptions, sender: Sender<MetaCacheEntry>) -> Result<()> {
if opts.versioned {
return self.list_versions_authoritatively(rx, opts, sender).await;
}
let list_path_started = std::time::Instant::now();
let (mut disks, infos, _) = self.get_online_disks_with_healing_and_info(true).await;
@@ -7140,6 +7207,7 @@ mod test {
use crate::disk::{DiskAPI, DiskOption, STORAGE_FORMAT_FILE, endpoint::Endpoint, error::DiskError, new_disk};
use crate::error::{Result, StorageError};
use crate::object_api::ObjectInfo;
use crate::set_disk::SetDisks;
use rustfs_filemeta::{
FileInfo, FileMeta, FileMetaVersion, MetaCacheEntries, MetaCacheEntriesSorted, MetaCacheEntry, MetaDeleteMarker,
ObjectPartInfo, VersionType,
@@ -7404,6 +7472,8 @@ mod test {
metadata.insert("etag".to_string(), (*etag).to_string());
let mut fi = FileInfo::new(name, *data_blocks, *parity_blocks);
fi.erasure.index = 1;
fi.data_dir = Some(Uuid::from_u128(0x1234));
fi.volume = "bucket".to_owned();
fi.name = name.to_owned();
let version_idx = u128::try_from(idx + 1).expect("test version index should fit u128");
@@ -9590,6 +9660,168 @@ mod test {
assert_eq!(resolver.requested_versions, 1);
}
#[test]
fn version_listing_rejects_stale_markers_at_full_set_read_quorum() {
let marker = test_delete_marker_meta_entry(
"marker.bin",
time::OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("valid fixture time"),
);
let sampled = MetaCacheEntries((0..8).map(|slot| (slot < 4).then(|| marker.clone())).collect());
assert!(
resolve_listing_entries(sampled, list_metadata_resolution_params("bucket".into(), 4, 12, true, 0), false).is_some(),
"the former sampled resolver accepts the four stale replicas"
);
for copies in [1, 4, 7, 8, 9, 16] {
let entries = MetaCacheEntries((0..16).map(|slot| (slot < copies).then(|| marker.clone())).collect());
let resolved = SetDisks::resolve_listed_versions("bucket", entries, &vec![None; 16], 16, 4)
.expect("known marker replicas and exact absence must be decidable");
assert_eq!(resolved.is_some(), copies >= 8, "marker copies: {copies}");
}
}
#[test]
fn version_listing_checks_history_after_a_quorum_marker() {
let time = time::OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("valid fixture time");
let mut marker = test_delete_marker_meta_entry("history.bin", time + time::Duration::seconds(1));
let history = test_object_meta_entry_with_erasure_versions("history.bin", &[(time, "old", 12, 4)]);
marker
.cached
.as_mut()
.expect("marker fixture")
.versions
.extend(history.cached.expect("history fixture").versions);
marker.metadata = marker
.cached
.as_ref()
.expect("combined fixture")
.marshal_msg()
.expect("encode history");
for copies in [8, 11, 12, 16] {
let entries = MetaCacheEntries((0..16).map(|slot| (slot < copies).then(|| marker.clone())).collect());
let resolved = SetDisks::resolve_listed_versions("bucket", entries, &vec![None; 16], 16, 4)
.expect("known history must resolve")
.expect("marker has read quorum");
let versions = resolved.file_info_versions("bucket").expect("read selected history").versions;
assert!(versions[0].deleted, "marker must remain latest");
assert_eq!(versions.len(), if copies >= 12 { 2 } else { 1 });
}
}
#[test]
fn version_listing_preserves_metadata_uncertainty() {
let marker = test_delete_marker_meta_entry(
"marker.bin",
time::OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("valid fixture time"),
);
let entries = MetaCacheEntries(vec![Some(marker); 4]);
assert!(matches!(
SetDisks::resolve_listed_versions("bucket", entries, &vec![None; 4], 16, 4),
Err(DiskError::ErasureReadQuorum)
));
}
#[test]
fn version_listing_tolerates_corrupt_replicas_only_with_version_quorum() {
let entry = test_object_meta_entry_with_erasure_versions(
"readable.bin",
&[(time::OffsetDateTime::from_unix_timestamp(1_700_000_000).unwrap(), "current", 12, 4)],
);
let corrupt = MetaCacheEntry {
name: entry.name.clone(),
metadata: vec![0xff],
..Default::default()
};
for copies in [11, 12, 15] {
let entries = MetaCacheEntries(
(0..16)
.map(|slot| Some(if slot < copies { entry.clone() } else { corrupt.clone() }))
.collect(),
);
let resolved = SetDisks::resolve_listed_versions("bucket", entries, &vec![None; 16], 16, 4);
if copies >= 12 {
assert!(resolved.unwrap().is_some(), "a readable version must tolerate corrupt replicas");
} else {
assert!(resolved.is_err(), "corruption cannot establish the missing version's absence");
}
}
}
#[tokio::test]
async fn version_listing_does_not_expose_four_stale_marker_replicas() {
let bucket = "version-list-stale-marker";
let (dirs, set) = crate::ecstore_validation_blackbox::make_local_set_disks(16, 4).await;
let marker = test_delete_marker_meta_entry(
"marker.bin",
time::OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("valid fixture time"),
);
let current = test_object_meta_entry_with_erasure_versions(
"new.bin",
&[(
time::OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("valid fixture time")
+ time::Duration::seconds(1),
"new",
12,
4,
)],
);
for (slot, dir) in dirs.iter().enumerate() {
let current_dir = dir.path().join(bucket).join("new.bin");
tokio::fs::create_dir_all(&current_dir)
.await
.expect("create current object directory");
tokio::fs::write(current_dir.join("xl.meta"), &current.metadata)
.await
.expect("persist current version");
if slot < 4 {
let stale_dir = dir.path().join(bucket).join("marker.bin");
tokio::fs::create_dir_all(&stale_dir)
.await
.expect("create stale marker directory");
tokio::fs::write(stale_dir.join("xl.meta"), &marker.metadata)
.await
.expect("persist stale marker");
}
}
for quorum_mode in ["disk", "reduced", "optimal", "strict"] {
let (sender, mut receiver) = mpsc::channel(1);
let walk = set.list_path(
CancellationToken::new(),
ListPathOptions {
bucket: bucket.to_string(),
versioned: true,
incl_deleted: true,
recursive: true,
ask_disks: quorum_mode.to_string(),
limit: 100,
..Default::default()
},
sender,
);
let drain = async {
let mut names = Vec::new();
while let Some(entry) = receiver.recv().await {
if !entry.is_dir() {
names.push(entry.name);
}
}
names
};
let (result, names) = timeout(Duration::from_secs(10), async { tokio::join!(walk, drain) })
.await
.expect("native listing must terminate");
result.expect("all drives are readable");
assert_eq!(names, ["new.bin"], "version authority cannot depend on {quorum_mode} sampling");
}
for dir in dirs.iter().take(4) {
assert_eq!(
tokio::fs::read(dir.path().join(bucket).join("marker.bin/xl.meta"))
.await
.expect("listing must preserve unresolved physical marker"),
marker.metadata
);
}
}
#[test]
fn list_metadata_resolution_params_keeps_all_versions_for_version_listing() {
let resolver = list_metadata_resolution_params("bucket".to_string(), 3, 5, true, 0);
+290 -1
View File
@@ -4741,7 +4741,7 @@ impl ECStore {
{
return Err(StorageError::BucketNotFound(bucket.to_string()));
}
if opts.delete_prefix && opts.expected_bucket_incarnation_id.is_none() {
if opts.expected_bucket_incarnation_id.is_none() {
opts.expected_bucket_incarnation_id = current_bucket_incarnation_id;
}
#[cfg(any(test, feature = "test-util"))]
@@ -5179,6 +5179,9 @@ impl ECStore {
StorageError::BucketNotFound(bucket.to_string()),
);
}
if opts.expected_bucket_incarnation_id.is_none() {
opts.expected_bucket_incarnation_id = current_bucket_incarnation_id;
}
#[cfg(any(test, feature = "test-util"))]
if current_bucket_incarnation_id.is_some() {
pause_delete_after_object_lock_snapshot(bucket).await;
@@ -8343,6 +8346,292 @@ mod tests {
(dirs, store)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn retired_marker_converges_after_bucket_recreation_on_sixteen_disks() {
let bucket = "retired-marker-c11";
let object = "marker.bin";
let ctx = Arc::new(crate::runtime::instance::InstanceContext::new());
let (dirs, original) = make_local_set_disks_with_ctx(16, 4, ctx.clone()).await;
let store = Arc::new(new_prepared_reader_test_store_with_ctx(&[original], ctx).await);
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let create = MakeBucketOptions {
versioning_enabled: true,
..Default::default()
};
store.handle_make_bucket(bucket, &create).await.expect("create old bucket");
let old_incarnation = store.bucket_incarnation_id_from_disk(bucket).await.expect("old identity");
let mut historical = Vec::new();
for value in 1..=3 {
historical.push(
store
.put_object(
bucket,
"history.bin",
&mut PutObjReader::from_vec(vec![value; 4097]),
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("write old version")
.version_id
.unwrap(),
);
}
let marker = store
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("create old marker")
.version_id
.expect("explicit marker version");
assert!(
store
.handle_delete_bucket(bucket, &crate::storage_api_contracts::bucket::DeleteBucketOptions::default())
.await
.is_err()
);
assert!(
!crate::bucket::retirement::is_retired(store.clone(), bucket, old_incarnation)
.await
.unwrap(),
"failed bucket deletion cannot publish retirement"
);
let set = &store.pools[0].disk_set[0];
let before = set.disks.read().await.clone();
for disk in before.iter().flatten() {
let info = disk
.read_version("", bucket, object, &marker.to_string(), &crate::disk::ReadOptions::default())
.await
.expect("read original marker");
assert_eq!(info.delete_marker_incarnation(), Some(old_incarnation));
}
// Taking four actual members offline preserves their existing xl.meta.
// All deletion/recreation work below goes through the normal store API.
for slot in 12..16 {
set.disks.write().await[slot] = None;
}
for version in historical {
store
.delete_object(
bucket,
"history.bin",
ObjectOptions {
versioned: true,
version_id: Some(version.to_string()),
..Default::default()
},
)
.await
.expect("delete exact old data version with four drives offline");
}
store
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
version_id: Some(marker.to_string()),
..Default::default()
},
)
.await
.expect("delete exact old marker with four drives offline");
store
.handle_delete_bucket(bucket, &crate::storage_api_contracts::bucket::DeleteBucketOptions::default())
.await
.expect("delete empty old bucket");
assert!(
crate::bucket::retirement::is_retired(store.clone(), bucket, old_incarnation)
.await
.expect("durable retirement")
);
store
.handle_make_bucket(bucket, &create)
.await
.expect("recreate same bucket name");
let current = store.bucket_incarnation_id_from_disk(bucket).await.expect("new identity");
assert_ne!(old_incarnation, current);
let mut new_versions = Vec::new();
for (key, payload) in [
(object, b"new generation reused key".as_slice()),
("new.bin", b"new generation control".as_slice()),
] {
let info = store
.put_object(
bucket,
key,
&mut PutObjReader::from_vec(payload.to_vec()),
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("write new generation");
new_versions.push((key, info.version_id.unwrap(), payload));
}
*set.disks.write().await = before.clone();
let heal = rustfs_heal_contracts::heal_channel::HealOpts {
remove: true,
scan_mode: rustfs_heal_contracts::heal_channel::HealScanMode::Deep,
..Default::default()
};
let observed = store
.clone()
.inner_list_object_versions(bucket, "", None, None, None, 100)
.await
.expect("list after rejoin");
assert_eq!(observed.objects.len(), 2, "stale marker must not contaminate listing");
assert!(observed.objects.iter().all(|info| !info.delete_marker));
for (key, version, _) in &new_versions {
let repaired = store
.heal_object_with_proof(bucket, key, &version.to_string(), &heal)
.await
.expect("heal new generation");
assert!(repaired.error.is_none(), "new version heal: {:?}", repaired.error);
}
let no_remove = store
.heal_object_with_proof(
bucket,
object,
&marker.to_string(),
&rustfs_heal_contracts::heal_channel::HealOpts { remove: false, ..heal },
)
.await
.expect("inspect marker without removal");
assert!(no_remove.error.as_ref().is_some_and(Error::is_retired_marker_deferred));
assert!(no_remove.absence.is_none());
let dry = store
.heal_object_with_proof(
bucket,
object,
&marker.to_string(),
&rustfs_heal_contracts::heal_channel::HealOpts { dry_run: true, ..heal },
)
.await
.expect("dry run");
assert!(dry.absence.is_none());
for disk in before.iter().skip(12).flatten() {
assert!(
disk.read_version("", bucket, object, &marker.to_string(), &crate::disk::ReadOptions::default())
.await
.is_ok()
);
}
let original_marker = before[12]
.as_ref()
.unwrap()
.read_version("", bucket, object, &marker.to_string(), &crate::disk::ReadOptions::default())
.await
.unwrap();
for stamp in [None, Some(current), Some(Uuid::new_v4())] {
for disk in before.iter().skip(12).flatten() {
let mut changed = original_marker.clone();
changed.metadata.clear();
if let Some(stamp) = stamp {
changed.set_delete_marker_incarnation(stamp);
}
disk.write_metadata("", bucket, object, changed).await.unwrap();
}
let unresolved = store
.heal_object_with_proof(bucket, object, &marker.to_string(), &heal)
.await
.unwrap();
assert!(
unresolved.error.as_ref().is_some_and(Error::is_retired_marker_deferred),
"legacy/current/unproven markers must be deferred: {:?}",
unresolved.error
);
assert!(unresolved.absence.is_none());
for disk in before.iter().skip(12).flatten() {
assert!(
disk.read_version("", bucket, object, &marker.to_string(), &crate::disk::ReadOptions::default())
.await
.is_ok()
);
disk.write_metadata("", bucket, object, original_marker.clone())
.await
.unwrap();
}
}
set.disks.write().await[15] = None;
let incomplete = store
.heal_object_with_proof(bucket, object, &marker.to_string(), &heal)
.await
.unwrap();
assert!(incomplete.error.as_ref().is_some_and(Error::is_retired_marker_deferred));
assert!(incomplete.absence.is_none(), "offline targets cannot be omitted from a cleanup receipt");
set.disks.write().await[15] = before[15].clone();
let failure =
crate::set_disk::DanglingDeleteFailure::install(bucket, object, 15, crate::disk::error::DiskError::FaultyDisk);
let partial = store
.heal_object_with_proof(bucket, object, &marker.to_string(), &heal)
.await
.unwrap();
assert!(partial.error.as_ref().is_some_and(Error::is_retired_marker_deferred));
assert!(
partial.absence.is_none(),
"fifteen successful/absent disks cannot hide one remaining marker"
);
assert!(
before[15]
.as_ref()
.unwrap()
.read_version("", bucket, object, &marker.to_string(), &crate::disk::ReadOptions::default())
.await
.is_ok()
);
drop(failure);
let cleaned = store
.heal_object_with_proof(bucket, object, &marker.to_string(), &heal)
.await
.expect("cleanup retired marker");
assert!(cleaned.error.is_none(), "cleanup error: {:?}", cleaned.error);
let proof = cleaned.absence.expect("complete cleanup receipt");
assert!(proof.removed);
assert_eq!(proof.bucket_incarnation_id, current);
for disk in before.iter().flatten() {
assert!(matches!(
disk.read_version("", bucket, object, &marker.to_string(), &crate::disk::ReadOptions::default())
.await,
Err(crate::disk::error::DiskError::FileVersionNotFound | crate::disk::error::DiskError::FileNotFound)
));
}
let replay = store
.heal_object_with_proof(bucket, object, &marker.to_string(), &heal)
.await
.expect("idempotent replay");
assert!(replay.error.is_none());
assert!(!replay.absence.expect("authoritative absence receipt").removed);
for (key, version, payload) in new_versions {
let mut reader = store
.handle_get_object_reader(
bucket,
key,
None,
HeaderMap::new(),
&ObjectOptions {
version_id: Some(version.to_string()),
..Default::default()
},
)
.await
.expect("new generation still readable");
let mut actual = Vec::new();
reader.stream.read_to_end(&mut actual).await.expect("read new bytes");
assert_eq!(actual, payload);
}
assert_eq!(dirs.len(), 16);
}
#[tokio::test]
async fn multipool_delete_marker_stays_with_existing_versions() {
let bucket = "multipool-marker-routing";
+21
View File
@@ -1207,6 +1207,27 @@ impl FileInfo {
insert_str(&mut self.metadata, SUFFIX_OBJECT_TRANSACTION_EPOCH, epoch.to_string());
}
/// Bind a newly created delete marker to the destination bucket generation.
pub fn set_delete_marker_incarnation(&mut self, incarnation: Uuid) {
if self.deleted && !incarnation.is_nil() {
insert_str(
&mut self.metadata,
rustfs_utils::http::metadata_compat::SUFFIX_BUCKET_INCARNATION_ID,
incarnation.to_string(),
);
}
}
/// Legacy, malformed, or conflicting stamps cannot authorize marker cleanup.
pub fn delete_marker_incarnation(&self) -> Option<Uuid> {
if !self.deleted || self.tier_free_version() {
return None;
}
get_consistent_str(&self.metadata, rustfs_utils::http::metadata_compat::SUFFIX_BUCKET_INCARNATION_ID)
.and_then(|value| Uuid::parse_str(value).ok())
.filter(|id| !id.is_nil())
}
pub fn object_transaction_epoch(&self) -> Result<Option<Uuid>> {
if !contains_key_str(&self.metadata, SUFFIX_OBJECT_TRANSACTION_EPOCH) {
return Ok(None);
+9
View File
@@ -722,6 +722,15 @@ impl FileMeta {
mod_time: fi.mod_time,
..Default::default()
});
if let Some(incarnation) = fi.delete_marker_incarnation()
&& let Some(marker) = ventry.delete_marker.as_mut()
{
insert_bytes(
&mut marker.meta_sys,
rustfs_utils::http::metadata_compat::SUFFIX_BUCKET_INCARNATION_ID,
incarnation.to_string().into_bytes(),
);
}
}
let mut update_version = false;
+1
View File
@@ -48,6 +48,7 @@ pub struct HealObjectIdentity {
#[serde(rename_all = "snake_case")]
pub enum HealDeferredReason {
DanglingDeleteGrace,
RetiredMarkerProof,
TransientUsageCache,
TransientExistenceCheck,
Deadline,
+20
View File
@@ -903,6 +903,26 @@ impl HealTask {
true
}
async fn skip_retired_marker_error(&self, err: &Error) -> bool {
if !matches!(err, Error::Storage(source) if source.is_retired_marker_deferred()) {
return false;
}
if let Some(identity) = self.single_object_identity() {
let mut outcome = self.outcome.write().await;
outcome.attempt_failed();
outcome.record(HealObjectOutcome {
identity,
disposition: HealObjectDisposition::Deferred {
reason: HealDeferredReason::RetiredMarkerProof,
retry_not_before: None,
},
detail: Some(err.to_string()),
});
}
self.progress.write().await.update_stage(3, 3);
true
}
async fn skip_dangling_delete_grace_error(&self, bucket: &str, object: &str, err: &Error) -> bool {
if !Self::is_dangling_delete_grace_error(err) {
return false;
+7 -1
View File
@@ -676,7 +676,13 @@ impl HealTask {
_ => {}
}
detail = Some(err.to_string());
if Self::is_dangling_delete_grace_error(&err) {
if matches!(&err, Error::Storage(source) if source.is_retired_marker_deferred()) {
disposition = HealObjectDisposition::Deferred {
reason: HealDeferredReason::RetiredMarkerProof,
retry_not_before: None,
};
telemetry_unknown |= !increment_counter(&mut skipped);
} else if Self::is_dangling_delete_grace_error(&err) {
disposition = HealObjectDisposition::Deferred {
reason: HealDeferredReason::DanglingDeleteGrace,
retry_not_before: err.dangling_delete_retry_not_before(),
+3 -2
View File
@@ -183,7 +183,8 @@ impl HealTask {
let result = storage_result.item;
let error = storage_result.error;
if let Some(e) = error {
if self.skip_dangling_delete_grace_error(bucket, object, &e).await {
if self.skip_retired_marker_error(&e).await || self.skip_dangling_delete_grace_error(bucket, object, &e).await
{
return Ok(());
}
@@ -277,7 +278,7 @@ impl HealTask {
Err(Error::TaskCancelled) => Err(Error::TaskCancelled),
Err(Error::TaskTimeout) => Err(Error::TaskTimeout),
Err(e) => {
if self.skip_dangling_delete_grace_error(bucket, object, &e).await {
if self.skip_retired_marker_error(&e).await || self.skip_dangling_delete_grace_error(bucket, object, &e).await {
return Ok(());
}
+52
View File
@@ -528,6 +528,49 @@ mod canonical_outcome {
assert_eq!(legacy.disposition, HealObjectDisposition::Unknown);
}
#[tokio::test]
async fn retired_marker_is_deferred_for_bucket_and_single_object_tasks() {
let storage = Arc::new(MockStorage::default());
storage
.heal_object_outcomes
.lock()
.unwrap()
.insert("object-a".into(), VecDeque::from([MockHealObjectOutcome::RetiredMarkerDeferred]));
let bucket = bucket_task(storage);
let object = HealTask::from_request(
HealRequest::object("bucket-a".into(), "marker.bin".into(), Some(Uuid::new_v4().to_string())),
Arc::new(MockStorage {
heal_object_outcome: Mutex::new(Some(MockHealObjectOutcome::RetiredMarkerDeferred)),
..Default::default()
}),
);
for task in [&bucket, &object] {
task.execute().await.expect("unproven marker permits traversal completion");
let outcome = task.get_outcome().await;
assert_eq!(outcome.counters.failed, 0);
assert_eq!(outcome.counters.healed, 0);
let deferred = outcome
.objects
.iter()
.find(|item| {
matches!(
item.disposition,
HealObjectDisposition::Deferred {
reason: HealDeferredReason::RetiredMarkerProof,
..
}
)
})
.expect("typed deferral");
assert!(
deferred
.detail
.as_ref()
.is_some_and(|detail| detail.contains("no committed retirement record"))
);
}
}
#[tokio::test]
async fn grace_single_object_is_completed_but_deferred() {
let storage = Arc::new(MockStorage {
@@ -1755,6 +1798,7 @@ enum MockHealObjectOutcome {
OkWithReadQuorum,
ErrOther(&'static str),
DanglingGraceDeferred,
RetiredMarkerDeferred,
UnavailableDrive(DriveState),
RetryableReadQuorum,
RetryableSlowDown,
@@ -1884,6 +1928,10 @@ impl HealStorageAPI for MockStorage {
.and_then(VecDeque::pop_front)
{
return match outcome {
MockHealObjectOutcome::RetiredMarkerDeferred => Ok((
HealResultItem::default(),
Some(Error::Storage(EcstoreError::retired_marker_deferred("no committed retirement record"))),
)),
MockHealObjectOutcome::DanglingGraceDeferred => Ok((
HealResultItem::default(),
Some(Error::Disk(DiskError::other(
@@ -1926,6 +1974,10 @@ impl HealStorageAPI for MockStorage {
}
if let Some(outcome) = self.heal_object_outcome.lock().unwrap().take() {
return match outcome {
MockHealObjectOutcome::RetiredMarkerDeferred => Ok((
HealResultItem::default(),
Some(Error::Storage(EcstoreError::retired_marker_deferred("no committed retirement record"))),
)),
MockHealObjectOutcome::DanglingGraceDeferred => Ok((
HealResultItem::default(),
Some(Error::Disk(DiskError::other(
@@ -2386,6 +2386,21 @@ pub mod node_service_client {
.insert(GrpcMethod::new("node_service.NodeService", "DeleteVersion"));
self.inner.unary(req, path, codec).await
}
pub async fn delete_retired_marker(
&mut self,
request: impl tonic::IntoRequest<super::DeleteVersionRequest>,
) -> std::result::Result<tonic::Response<super::DeleteVersionResponse>, tonic::Status> {
self.inner
.ready()
.await
.map_err(|e| tonic::Status::unknown(format!("Service was not ready: {}", e.into())))?;
let codec = tonic_prost::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static("/node_service.NodeService/DeleteRetiredMarker");
let mut req = request.into_request();
req.extensions_mut()
.insert(GrpcMethod::new("node_service.NodeService", "DeleteRetiredMarker"));
self.inner.unary(req, path, codec).await
}
pub async fn delete_versions(
&mut self,
request: impl tonic::IntoRequest<super::DeleteVersionsRequest>,
@@ -3399,6 +3414,10 @@ pub mod node_service_server {
&self,
request: tonic::Request<super::DeleteVersionRequest>,
) -> std::result::Result<tonic::Response<super::DeleteVersionResponse>, tonic::Status>;
async fn delete_retired_marker(
&self,
request: tonic::Request<super::DeleteVersionRequest>,
) -> std::result::Result<tonic::Response<super::DeleteVersionResponse>, tonic::Status>;
async fn delete_versions(
&self,
request: tonic::Request<super::DeleteVersionsRequest>,
@@ -4735,6 +4754,34 @@ pub mod node_service_server {
};
Box::pin(fut)
}
"/node_service.NodeService/DeleteRetiredMarker" => {
#[allow(non_camel_case_types)]
struct DeleteRetiredMarkerSvc<T: NodeService>(pub Arc<T>);
impl<T: NodeService> tonic::server::UnaryService<super::DeleteVersionRequest> for DeleteRetiredMarkerSvc<T> {
type Response = super::DeleteVersionResponse;
type Future = BoxFuture<tonic::Response<Self::Response>, tonic::Status>;
fn call(&mut self, request: tonic::Request<super::DeleteVersionRequest>) -> Self::Future {
let inner = Arc::clone(&self.0);
let fut = async move { <T as NodeService>::delete_retired_marker(&inner, request).await };
Box::pin(fut)
}
}
let accept_compression_encodings = self.accept_compression_encodings;
let send_compression_encodings = self.send_compression_encodings;
let max_decoding_message_size = self.max_decoding_message_size;
let max_encoding_message_size = self.max_encoding_message_size;
let inner = self.inner.clone();
let fut = async move {
let method = DeleteRetiredMarkerSvc(inner);
let codec = tonic_prost::ProstCodec::default();
let mut grpc = tonic::server::Grpc::new(codec)
.apply_compression_config(accept_compression_encodings, send_compression_encodings)
.apply_max_message_size_config(max_decoding_message_size, max_encoding_message_size);
let res = grpc.unary(method, req).await;
Ok(res)
};
Box::pin(fut)
}
"/node_service.NodeService/DeleteVersions" => {
#[allow(non_camel_case_types)]
struct DeleteVersionsSvc<T: NodeService>(pub Arc<T>);
+1
View File
@@ -1213,6 +1213,7 @@ service NodeService {
rpc BatchReadVersion(BatchReadVersionRequest) returns (BatchReadVersionResponse) {}; // auth-policy: read-only
rpc ReadXL(ReadXLRequest) returns (ReadXLResponse) {}; // auth-policy: read-only
rpc DeleteVersion(DeleteVersionRequest) returns (DeleteVersionResponse) {}; // auth-policy: body-bound
rpc DeleteRetiredMarker(DeleteVersionRequest) returns (DeleteVersionResponse) {}; // auth-policy: body-bound
rpc DeleteVersions(DeleteVersionsRequest) returns (DeleteVersionsResponse) {}; // auth-policy: body-bound
rpc ReadMultiple(ReadMultipleRequest) returns (ReadMultipleResponse) {}; // auth-policy: read-only
rpc DeleteVolume(DeleteVolumeRequest) returns (DeleteVolumeResponse) {}; // auth-policy: body-bound
@@ -17,6 +17,18 @@ Heal and every foreground or background write path serialize on the same object-
## Heal lock scope
### Retired bucket delete markers
An explicit version heal may remove a subquorum delete marker from a deleted bucket generation only through `heal_object_with_proof` in `crates/ecstore/src/store/heal.rs`. The coordinator holds the bucket lifecycle read fence before entering the existing capacity and object lock scopes. The current, non-nil bucket incarnation must differ from the marker's persisted incarnation.
Normal single and batch marker creation stamps both `x-rustfs-internal-bucket-incarnation-id` and `x-minio-internal-bucket-incarnation-id`. A successful bucket deletion publishes an immutable, deployment-bound retirement record under the system bucket's `bucket-retirements` namespace, outside ordinary bucket metadata cleanup (`crates/ecstore/src/bucket/retirement.rs`). Publication follows successful physical deletion and the existing all-set rollback decision, precedes metadata cleanup and the successful DELETE response, and retains the lifecycle guards through the metadata write. Failed physical deletion does not publish retirement authority. Failed or uncertain record publication returns an error; it cannot authorize cleanup without an authoritative read of the record. Records have no age-based expiration because offline members can return arbitrarily late.
Cleanup requires `remove=true`, `dry_run=false`, an explicit non-null version, a committed retirement record, and agreement on the full surviving marker metadata. The `DeleteRetiredMarker` RPC carries the full marker precondition in the existing authenticated request body. `LocalDisk::delete_version_inner` checks it after acquiring the metadata mutation lease and reading the actual `xl.meta`; a replaced marker or data version is preserved. Older nodes that do not implement this RPC reject it; coordinators do not fall back to `DeleteVersion`.
Every selected disk must be online and confirm exact version absence after cleanup. Changed topology, lost fences, failed mutations, unreadable evidence, legacy markers without an incarnation, and conflicting identities produce no completion receipt. Unproven markers remain deferred with `retired_marker_proof`; a normal not-found error is not globally reclassified as success. An interrupted operation can leave partial cleanup, and replay checks the remaining actual replicas before producing an absence receipt. Deployments upgraded after an old deletion without stamps or retirement evidence retain those residual markers for separate recovery; new successful DELETE operations provide the evidence needed for automatic convergence.
`ListObjectVersions` independently reads all reachable replicas and resolves each version with the full erasure-set metadata read quorum (`SetDisks::resolve_listed_versions` in `crates/ecstore/src/set_disk/metadata.rs`). Sample settings in `RUSTFS_LIST_OBJECTS_QUORUM` continue to govern ordinary object listing; they do not lower version-listing authority. Local reader page limits cannot truncate this discovery before the merged page is resolved. This increases metadata reads relative to sampled listing, but does not issue a separate GET for every version. Heal's union walk still discovers subquorum repair candidates.
`heal_object` delegates to `heal_object_with_explicit_version_regen`, which takes the namespace write lock at entry unless `opts.no_lock` is set and binds the guard to the function scope. The guard covers the quorum metadata read, EC reconstruction, per-disk rename commit, tmp cleanup, the `HEAL_RENAME_INCOMPLETE` partial-commit return, and orphan `data_dir` reclamation (`reclaim_orphan_data_dirs`).
Read-repair heals (`opts.read_repair`) hold a shared lock (`HealObjectLockKind::Read`) during reconstruction so readers keep flowing, then `acquire_revalidated_read_repair_commit_lock` takes the write lock and re-reads a commit fingerprint; a changed fingerprint aborts the commit (`read_repair_commit_stale`).
+63 -2
View File
@@ -1668,7 +1668,14 @@ impl Node for NodeService {
}
async fn delete_version(&self, request: Request<DeleteVersionRequest>) -> Result<Response<DeleteVersionResponse>, Status> {
self.handle_delete_version(request).await
self.handle_delete_version(request, false).await
}
async fn delete_retired_marker(
&self,
request: Request<DeleteVersionRequest>,
) -> Result<Response<DeleteVersionResponse>, Status> {
self.handle_delete_version(request, true).await
}
async fn delete_versions(&self, request: Request<DeleteVersionsRequest>) -> Result<Response<DeleteVersionsResponse>, Status> {
@@ -2915,9 +2922,10 @@ mod tests {
use tonic::{Request, Response, Status};
use uuid::Uuid;
const DISK_MUTATION_RPC_METHODS: [&str; 18] = [
const DISK_MUTATION_RPC_METHODS: [&str; 19] = [
"renamedata",
"deleteversion",
"deleteretiredmarker",
"deleteversions",
"writemetadata",
"updatemetadata",
@@ -4275,6 +4283,20 @@ mod tests {
},
rustfs_protos::canonical_delete_version_request_body
);
assert_gated!(
delete_retired_marker,
DeleteVersionRequest {
disk: disk.clone(),
volume: "v".into(),
path: "p".into(),
file_info: "{}".into(),
force_del_marker: false,
opts: "{}".into(),
file_info_bin: vec![0x80].into(),
opts_bin: vec![0x80].into(),
},
rustfs_protos::canonical_delete_version_request_body
);
assert_gated!(
delete_versions,
DeleteVersionsRequest {
@@ -6567,6 +6589,45 @@ mod tests {
assert!(read_response.raw_file_info.is_empty());
}
#[tokio::test]
async fn retired_marker_rpc_rejects_missing_or_combined_conditions() {
use crate::storage::storage_api::ecstore_disk::DeleteOptions;
use rustfs_filemeta::{FileInfo, MetaDeleteMarker};
let service = create_test_node_service();
let mut marker = FileInfo {
deleted: true,
version_id: Some(Uuid::new_v4()),
mod_time: Some(time::OffsetDateTime::now_utc()),
..Default::default()
};
marker.set_delete_marker_incarnation(Uuid::new_v4());
for opts in [
DeleteOptions::default(),
DeleteOptions {
undo_write: true,
expected_delete_marker: Some(MetaDeleteMarker::from(marker)),
..Default::default()
},
] {
let mut request = Request::new(DeleteVersionRequest {
disk: "invalid-disk-path".into(),
volume: "bucket".into(),
path: "marker.bin".into(),
file_info: serde_json::to_string(&FileInfo::default()).unwrap(),
opts: serde_json::to_string(&opts).unwrap(),
..Default::default()
});
let body = rustfs_protos::canonical_delete_version_request_body(request.get_ref()).unwrap();
set_tonic_canonical_body_digest(&mut request, &body).unwrap();
mark_v2_authenticated(&mut request);
let error = service
.delete_retired_marker(request)
.await
.expect_err("invalid conditional request must be rejected before disk lookup");
assert_eq!(error.code(), tonic::Code::InvalidArgument);
}
}
#[tokio::test]
async fn test_delete_version_invalid_disk() {
let service = create_test_node_service();
@@ -726,6 +726,7 @@ impl NodeService {
pub(super) async fn handle_delete_version(
&self,
request: Request<DeleteVersionRequest>,
require_marker_condition: bool,
) -> Result<Response<DeleteVersionResponse>, Status> {
verify_disk_mutation_digest(
&request,
@@ -753,6 +754,21 @@ impl NodeService {
}));
}
};
if require_marker_condition && opts.expected_delete_marker.is_none() {
return Err(Status::invalid_argument("retired marker deletion requires a marker precondition"));
}
if opts.expected_delete_marker.is_some()
&& (request.force_del_marker
|| opts.undo_write
|| opts.undo_delete
|| opts.recursive
|| opts.immediate
|| opts.old_data_dir.is_some())
{
return Err(Status::invalid_argument(
"retired marker preconditions cannot be combined with other mutations",
));
}
let result = if opts.undo_write {
if request.force_del_marker {
Err(DiskError::other("undo_write cannot force a delete marker"))