fix(ilm): enqueue committed tier free versions (#7041)

* fix(ilm): enqueue committed tier free versions

* fix(ilm): stabilize causal cleanup CI coverage

* test(ilm): make expire GET race deterministic

* test(ilm): synchronize expiry with active GET
This commit is contained in:
cxymds
2026-09-02 19:08:28 +08:00
committed by GitHub
parent 922552083f
commit afc66b7182
11 changed files with 1729 additions and 161 deletions
+1 -1
View File
@@ -1 +1 @@
sha256=dbebfbab9b9efd4eff31211e69dd32235dc00e207f2ab0dd919a1b2ac9e724c2
sha256=db9bd8cdcb0abe43461aa6b36499b17cabd4098e5b34e300b1a0f0d0f34d9884
+5
View File
@@ -440,6 +440,11 @@ pub mod object {
PrepareSelectObjectSnapshotError, PreparedGetObjectReader, SelectObjectSnapshot, SelectObjectSnapshotReadError,
SnapshotConsistencyError,
};
#[cfg(feature = "test-util")]
pub mod test_util {
pub use crate::store::DeleteAfterObjectLockSnapshotBarrier;
}
}
pub mod rebalance {
@@ -921,7 +921,7 @@ impl ExpiryState {
Ok(())
}
pub fn enqueue_free_version(&mut self, oi: ObjectInfo) -> bool {
pub fn enqueue_free_version(&self, oi: ObjectInfo) -> bool {
let task = FreeVersionTask(oi);
let wrkr = self.get_worker_ch(task.op_hash());
if wrkr.is_none() {
@@ -1215,6 +1215,22 @@ impl ExpiryState {
}
}
pub(crate) async fn enqueue_committed_free_versions(api: &ECStore, free_versions: Vec<ObjectInfo>) -> usize {
if free_versions.is_empty() {
return 0;
}
let expiry_state = api.ctx.expiry_state();
let state = expiry_state.read().await;
let mut queued = 0;
for free_version in free_versions {
if state.enqueue_free_version(free_version) {
queued += 1;
}
}
queued
}
async fn enqueue_recovered_free_version_with_state(state: &Arc<RwLock<ExpiryState>>, oi: ObjectInfo) -> bool {
let task = FreeVersionTask(oi);
let hash = task.op_hash();
@@ -6627,7 +6643,7 @@ mod tests {
async fn enqueue_free_version_reports_false_without_worker_channel() {
let state = ExpiryState::new();
let recovery_notify = Arc::clone(&state.read().await.recovery_notify);
let mut state = state.write().await;
let state = state.write().await;
let oi = ObjectInfo {
bucket: "bucket".to_string(),
name: "object".to_string(),
@@ -6819,7 +6835,7 @@ mod tests {
},
..Default::default()
};
let mut state = state.write().await;
let state = state.write().await;
assert!(state.enqueue_free_version(oi.clone()));
assert!(recovery_notify.notified().now_or_never().is_none());
+7
View File
@@ -6444,6 +6444,9 @@ impl LocalDisk {
.abort_reserved_version_delete(object_dir, rollback_dir, volume, path, "delete_versions_commit_intent", err)
.await);
}
if should_fail_after_delete_commit(self.root.as_path(), path) {
return Err(DiskError::Unexpected);
}
return Ok(());
}
@@ -6514,6 +6517,10 @@ impl LocalDisk {
.await);
}
if should_fail_after_delete_commit(self.root.as_path(), path) {
return Err(DiskError::Unexpected);
}
Ok(())
}
+568
View File
@@ -19,6 +19,7 @@ use crate::storage_api_contracts::{
HTTPPreconditions, ObjectLockRetentionOptions, ObjectPreconditionError, ObjectPreconditionPart, ObjectPreconditionState,
},
};
use std::io;
use std::sync::atomic::{AtomicBool, AtomicU8, Ordering};
use tokio::sync::{Mutex, Notify, OwnedRwLockReadGuard};
use tokio_util::sync::CancellationToken;
@@ -678,6 +679,196 @@ pub struct DecommissionCapacityOptions {
pub(crate) mutation_id: Option<Uuid>,
}
/// Opaque storage-owned collection point for post-commit tier free-version
/// cleanup receipts. This type is public only because workspace crates build
/// [`ObjectOptions`] with struct literals; callers outside `ecstore` must leave
/// the corresponding option unset.
#[doc(hidden)]
#[derive(Clone)]
pub struct TierFreeVersionReceiptSink {
inner: Arc<parking_lot::Mutex<TierFreeVersionReceiptSinkState>>,
}
struct TierFreeVersionReceiptSinkState {
receipts: Option<HashMap<TierFreeVersionReceiptIdentity, TierFreeVersionReceiptPayload>>,
}
#[derive(PartialEq, Eq, Hash)]
struct TierFreeVersionReceiptIdentity {
bucket: String,
logical_name: String,
tier: String,
remote_name: String,
remote_version_state: TierFreeVersionReceiptVersionState,
remote_version: String,
backend_identity: crate::services::tier::tier::TierDestinationId,
}
struct TierFreeVersionReceiptPayload {
local_free_version_id: Uuid,
mod_time: Option<OffsetDateTime>,
}
#[derive(Clone, Copy, PartialEq, Eq, Hash)]
enum TierFreeVersionReceiptVersionState {
KnownDisabled,
SuspendedNull,
Exact,
}
impl TierFreeVersionReceiptSink {
/// Only the delete wrapper may originate a sink. The public type exists so
/// workspace struct literals can carry it, but external crates cannot
/// create an undrainable collector accidentally.
pub(crate) fn new() -> Self {
Self {
inner: Arc::new(parking_lot::Mutex::new(TierFreeVersionReceiptSinkState {
receipts: Some(HashMap::new()),
})),
}
}
}
impl std::fmt::Debug for TierFreeVersionReceiptSink {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let state = self.inner.lock();
f.debug_struct("TierFreeVersionReceiptSink")
.field("drained", &state.receipts.is_none())
.field("receipt_count", &state.receipts.as_ref().map(HashMap::len).unwrap_or_default())
.finish()
}
}
impl TierFreeVersionReceiptVersionState {
fn persisted(self) -> rustfs_filemeta::TransitionVersionState {
match self {
Self::KnownDisabled => rustfs_filemeta::TransitionVersionState::KnownDisabled,
Self::SuspendedNull => rustfs_filemeta::TransitionVersionState::SuspendedNull,
Self::Exact => rustfs_filemeta::TransitionVersionState::Exact,
}
}
}
impl TierFreeVersionReceiptIdentity {
fn into_object_info(self, payload: TierFreeVersionReceiptPayload) -> ObjectInfo {
let mut metadata = HashMap::with_capacity(2);
rustfs_utils::http::metadata_compat::insert_str(
&mut metadata,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
rustfs_utils::crypto::hex(self.backend_identity),
);
ObjectInfo {
bucket: self.bucket,
name: self.logical_name,
mod_time: payload.mod_time,
user_defined: Arc::new(metadata),
version_id: Some(payload.local_free_version_id),
delete_marker: true,
transitioned_object: TransitionedObject {
name: self.remote_name,
version_id: self.remote_version,
tier: self.tier,
free_version: true,
status: String::new(),
},
transition_version_state: self.remote_version_state.persisted(),
..Default::default()
}
}
}
fn tier_free_version_scheduling_receipt_from_source(
source: &ObjectInfo,
local_free_version_id: Uuid,
) -> io::Result<Option<(TierFreeVersionReceiptIdentity, TierFreeVersionReceiptPayload)>> {
if source.transitioned_object.status != rustfs_filemeta::TRANSITION_COMPLETE
|| source.transitioned_object.free_version
|| source.delete_marker
|| source.bucket.is_empty()
|| source.name.is_empty()
|| source.transitioned_object.tier.is_empty()
|| source.transitioned_object.name.is_empty()
{
return Ok(None);
}
if local_free_version_id.is_nil() {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"tier free-version receipt has a nil local version identity",
));
}
let remote_version = source.transitioned_object.version_id.as_str();
let remote_version_state = match source.transition_version_state {
rustfs_filemeta::TransitionVersionState::Unknown => return Ok(None),
rustfs_filemeta::TransitionVersionState::KnownDisabled if remote_version.is_empty() => {
TierFreeVersionReceiptVersionState::KnownDisabled
}
rustfs_filemeta::TransitionVersionState::SuspendedNull if remote_version == "null" => {
TierFreeVersionReceiptVersionState::SuspendedNull
}
rustfs_filemeta::TransitionVersionState::Exact if !remote_version.is_empty() && remote_version != "null" => {
TierFreeVersionReceiptVersionState::Exact
}
_ => return Ok(None),
};
let Some(backend_identity) = crate::services::tier::tier::tier_destination_id_from_metadata(&source.user_defined)
.map_err(|err| io::Error::new(io::ErrorKind::InvalidData, err))?
else {
return Ok(None);
};
Ok(Some((
TierFreeVersionReceiptIdentity {
bucket: source.bucket.clone(),
logical_name: decode_dir_object(&source.name),
tier: source.transitioned_object.tier.clone(),
remote_name: source.transitioned_object.name.clone(),
remote_version_state,
remote_version: source.transitioned_object.version_id.clone(),
backend_identity,
},
TierFreeVersionReceiptPayload {
local_free_version_id,
mod_time: source.mod_time,
},
)))
}
impl TierFreeVersionReceiptSink {
/// Record one committed free-version cleanup target. Cloned options share
/// this sink; tuple-equivalent physical copies collapse to one worker task.
/// `false` means the source cannot safely identify a destructive cleanup.
pub(crate) fn record(&self, source: &ObjectInfo, local_free_version_id: Uuid) -> io::Result<bool> {
let Some((identity, payload)) = tier_free_version_scheduling_receipt_from_source(source, local_free_version_id)? else {
return Ok(false);
};
let mut state = self.inner.lock();
let receipts = state
.receipts
.as_mut()
.ok_or_else(|| io::Error::new(io::ErrorKind::BrokenPipe, "tier free-version receipt sink was already drained"))?;
receipts.entry(identity).or_insert(payload);
Ok(true)
}
/// Consume every receipt exactly once. A second drain is a caller bug: it
/// could otherwise make two outer wrappers believe they own the same tasks.
pub(crate) fn drain(&self) -> io::Result<Vec<ObjectInfo>> {
let mut state = self.inner.lock();
let receipts = state
.receipts
.take()
.ok_or_else(|| io::Error::new(io::ErrorKind::BrokenPipe, "tier free-version receipt sink was already drained"))?;
drop(state);
Ok(receipts
.into_iter()
.map(|(identity, payload)| identity.into_object_info(payload))
.collect())
}
}
#[derive(Default, Clone)]
pub struct ObjectOptions {
// Use the maximum parity (N/2), used when saving server configuration files
@@ -716,6 +907,12 @@ pub struct ObjectOptions {
pub skip_rebalancing: bool,
pub skip_free_version: bool,
/// Storage-owned, per-request hand-off for committed tier free-version
/// cleanup work. The outer delete wrapper installs and drains it; clones
/// below that boundary share the same opaque sink.
#[doc(hidden)]
pub tier_free_version_receipt_sink: Option<TierFreeVersionReceiptSink>,
/// Cooperative cancellation for an owned PutObject before authoritative
/// rename begins. Storage ignores it after entering the durable commit.
#[doc(hidden)]
@@ -851,6 +1048,7 @@ impl std::fmt::Debug for ObjectOptions {
.field("skip_decommissioned", &self.skip_decommissioned)
.field("skip_rebalancing", &self.skip_rebalancing)
.field("skip_free_version", &self.skip_free_version)
.field("tier_free_version_receipt_sink", &self.tier_free_version_receipt_sink)
.field("put_object_cancellation", &self.put_object_cancellation.is_some())
.field("scanner_publication_commit_scope", &self.scanner_publication_commit_scope)
.field("data_movement", &self.data_movement)
@@ -2622,11 +2820,381 @@ mod tests {
assert!(default_cloned.parts.is_empty());
}
fn transitioned_receipt_source(
bucket: &str,
object: &str,
remote_version: &str,
version_state: rustfs_filemeta::TransitionVersionState,
identity_hex: Option<&str>,
) -> ObjectInfo {
let mut metadata = HashMap::new();
if let Some(identity_hex) = identity_hex {
rustfs_utils::http::metadata_compat::insert_str(
&mut metadata,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
identity_hex.to_string(),
);
}
ObjectInfo {
bucket: bucket.to_string(),
name: object.to_string(),
mod_time: Some(OffsetDateTime::UNIX_EPOCH),
user_defined: Arc::new(metadata),
transitioned_object: TransitionedObject {
name: format!("remote/{object}"),
version_id: remote_version.to_string(),
tier: "WARM".to_string(),
status: TRANSITION_COMPLETE.to_string(),
..Default::default()
},
transition_version_state: version_state,
..Default::default()
}
}
#[test]
fn tier_free_version_receipt_matches_persisted_free_version_worker_fields() {
let bucket = "receipt-bucket";
let object = "archive/object.bin";
let source_version_id = Uuid::from_u128(1);
let local_free_version_id = Uuid::from_u128(2);
let remote_version_id = Uuid::from_u128(3);
let source_mod_time = OffsetDateTime::UNIX_EPOCH + time::Duration::hours(4);
let identity_hex = "ab".repeat(32);
let mut source_metadata = HashMap::from([
("etag".to_string(), "source-etag".to_string()),
("x-amz-meta-private".to_string(), "must-not-enter-receipt".to_string()),
]);
rustfs_utils::http::metadata_compat::insert_str(
&mut source_metadata,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
identity_hex.clone(),
);
let source_file_info = FileInfo {
volume: bucket.to_string(),
name: object.to_string(),
version_id: Some(source_version_id),
transition_status: TRANSITION_COMPLETE.to_string(),
transitioned_objname: "remote/receipt-object".to_string(),
transition_tier: "WARM".to_string(),
transition_version_id: Some(remote_version_id),
transition_version: Some(remote_version_id.to_string()),
transition_version_state: rustfs_filemeta::TransitionVersionState::Exact,
mod_time: Some(source_mod_time),
size: 8192,
data_dir: Some(Uuid::from_u128(4)),
metadata: source_metadata,
..Default::default()
};
let source = ObjectInfo::from_file_info(&source_file_info, bucket, object, true);
let mut persisted = FileMeta::new();
persisted
.add_version(source_file_info)
.expect("transitioned receipt source should be persisted");
let mut delete_file_info = FileInfo {
volume: bucket.to_string(),
name: object.to_string(),
version_id: Some(source_version_id),
mod_time: Some(source_mod_time + time::Duration::minutes(1)),
..Default::default()
};
delete_file_info.set_tier_free_version_id(&local_free_version_id.to_string());
persisted
.delete_version(&delete_file_info)
.expect("transitioned source delete should create a free-version");
let encoded = persisted.marshal_msg().expect("free-version metadata should encode");
let decoded = FileMeta::load(&encoded).expect("free-version metadata should decode");
let persisted_free_version = decoded
.get_all_file_info_versions(bucket, object, true)
.expect("decoded free-version should produce FileInfo")
.versions
.into_iter()
.find(|version| version.tier_free_version())
.expect("decoded metadata should contain the persisted free-version");
let persisted_object_info = ObjectInfo::from_file_info(&persisted_free_version, bucket, object, true);
let sink = TierFreeVersionReceiptSink::new();
assert!(
sink.record(&source, local_free_version_id)
.expect("valid transitioned source should produce a receipt")
);
let mut receipts = sink.drain().expect("receipt owner should drain exactly once");
assert_eq!(receipts.len(), 1);
let receipt = receipts.pop().expect("one receipt should be present");
assert_eq!(receipt.bucket, persisted_object_info.bucket);
assert_eq!(receipt.name, persisted_object_info.name);
assert_eq!(receipt.version_id, persisted_object_info.version_id);
assert_eq!(receipt.mod_time, persisted_object_info.mod_time);
assert_eq!(receipt.delete_marker, persisted_object_info.delete_marker);
assert_eq!(receipt.transitioned_object.name, persisted_object_info.transitioned_object.name);
assert_eq!(
receipt.transitioned_object.version_id,
persisted_object_info.transitioned_object.version_id
);
assert_eq!(receipt.transitioned_object.tier, persisted_object_info.transitioned_object.tier);
assert_eq!(
receipt.transitioned_object.free_version,
persisted_object_info.transitioned_object.free_version
);
assert_eq!(receipt.transitioned_object.status, persisted_object_info.transitioned_object.status);
assert_eq!(receipt.transition_version_state, persisted_object_info.transition_version_state);
assert_eq!(
crate::services::tier::tier::tier_destination_id_from_metadata(&receipt.user_defined)
.expect("receipt identity should decode"),
crate::services::tier::tier::tier_destination_id_from_metadata(&persisted_object_info.user_defined)
.expect("persisted identity should decode")
);
assert_eq!(
receipt.user_defined.len(),
2,
"receipt should carry only the two compatibility identity keys"
);
assert_eq!(
receipt.user_defined.get("x-rustfs-internal-transition-tier-destination-id"),
Some(&identity_hex)
);
assert_eq!(
receipt.user_defined.get("x-minio-internal-transition-tier-destination-id"),
Some(&identity_hex)
);
assert!(!receipt.user_defined.contains_key("x-amz-meta-private"));
assert_eq!(receipt.size, 0);
assert_eq!(receipt.actual_size, 0);
assert!(receipt.parts.is_empty());
assert!(receipt.etag.is_none());
assert!(receipt.checksum.is_none());
assert!(receipt.data_dir.is_none());
}
#[test]
fn tier_free_version_receipt_sink_deduplicates_remote_target_and_drains_once() {
let identity_hex = "11".repeat(32);
let source = transitioned_receipt_source(
"bucket",
"object",
"remote-version",
rustfs_filemeta::TransitionVersionState::Exact,
Some(&identity_hex),
);
let other_object = transitioned_receipt_source(
"bucket",
"other-object",
"remote-version",
rustfs_filemeta::TransitionVersionState::Exact,
Some(&identity_hex),
);
let sink = TierFreeVersionReceiptSink::new();
let clone = sink.clone();
assert!(
sink.record(&source, Uuid::from_u128(10))
.expect("first physical receipt should record")
);
assert!(
clone
.record(&source, Uuid::from_u128(11))
.expect("tuple-equivalent physical receipt should be represented")
);
assert!(
clone
.record(&other_object, Uuid::from_u128(12))
.expect("a different logical key should retain its own task")
);
let mut receipts = sink.drain().expect("owner should drain shared receipts");
receipts.sort_by(|left, right| left.name.cmp(&right.name));
assert_eq!(receipts.len(), 2);
assert_eq!(receipts[0].name, "object");
assert_eq!(receipts[0].version_id, Some(Uuid::from_u128(10)));
assert_eq!(receipts[1].name, "other-object");
assert_eq!(receipts[1].version_id, Some(Uuid::from_u128(12)));
assert_eq!(
clone.drain().expect_err("a shared sink must drain only once").kind(),
io::ErrorKind::BrokenPipe
);
assert_eq!(
clone
.record(&source, Uuid::from_u128(13))
.expect_err("recording after drain must fail")
.kind(),
io::ErrorKind::BrokenPipe
);
}
#[test]
fn tier_free_version_receipt_identity_covers_every_destructive_dimension() {
let identity_hex = "44".repeat(32);
let baseline = transitioned_receipt_source(
"bucket",
"directory/",
"remote-version",
rustfs_filemeta::TransitionVersionState::Exact,
Some(&identity_hex),
);
let mut encoded_duplicate = baseline.clone();
encoded_duplicate.name = "directory__XLDIR__".to_string();
let mut variants = Vec::new();
let mut changed = baseline.clone();
changed.bucket = "other-bucket".to_string();
variants.push(changed);
let mut changed = baseline.clone();
changed.name = "other-directory/".to_string();
variants.push(changed);
let mut changed = baseline.clone();
changed.transitioned_object.tier = "COLD".to_string();
variants.push(changed);
let mut changed = baseline.clone();
changed.transitioned_object.name = "remote/other-directory/".to_string();
variants.push(changed);
let mut changed = baseline.clone();
changed.transitioned_object.version_id = "other-remote-version".to_string();
variants.push(changed);
let mut changed = baseline.clone();
changed.transition_version_state = rustfs_filemeta::TransitionVersionState::SuspendedNull;
changed.transitioned_object.version_id = "null".to_string();
variants.push(changed);
let mut changed = baseline.clone();
rustfs_utils::http::metadata_compat::insert_str(
Arc::make_mut(&mut changed.user_defined),
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
"55".repeat(32),
);
variants.push(changed);
let sink = TierFreeVersionReceiptSink::new();
assert!(
sink.record(&baseline, Uuid::from_u128(20))
.expect("baseline receipt should record")
);
assert!(
sink.record(&encoded_duplicate, Uuid::from_u128(21))
.expect("the encoded spelling of one logical key should deduplicate")
);
for (offset, variant) in variants.iter().enumerate() {
assert!(
sink.record(variant, Uuid::from_u128(30 + offset as u128))
.expect("each distinct cleanup identity should record")
);
}
let receipts = sink.drain().expect("identity matrix should drain once");
assert_eq!(receipts.len(), 8, "every destructive identity dimension must prevent deduplication");
let baseline_receipt = receipts
.iter()
.find(|receipt| {
receipt.bucket == "bucket"
&& receipt.name == "directory/"
&& receipt.transitioned_object.tier == "WARM"
&& receipt.transitioned_object.name == "remote/directory/"
&& receipt.transitioned_object.version_id == "remote-version"
&& receipt.transition_version_state == rustfs_filemeta::TransitionVersionState::Exact
&& crate::services::tier::tier::tier_destination_id_from_metadata(&receipt.user_defined)
.is_ok_and(|identity| identity == Some([0x44; 32]))
})
.expect("baseline cleanup identity should remain present");
assert_eq!(
baseline_receipt.version_id,
Some(Uuid::from_u128(20)),
"deduplication must retain the first UUID"
);
}
#[test]
fn tier_free_version_receipt_source_validation_fails_closed() {
let identity_hex = "22".repeat(32);
for (state, remote_version) in [
(rustfs_filemeta::TransitionVersionState::KnownDisabled, ""),
(rustfs_filemeta::TransitionVersionState::SuspendedNull, "null"),
(rustfs_filemeta::TransitionVersionState::Exact, "opaque-version"),
] {
let source = transitioned_receipt_source("bucket", "object", remote_version, state, Some(&identity_hex));
assert!(
TierFreeVersionReceiptSink::new()
.record(&source, Uuid::new_v4())
.expect("canonical remote-version state should be eligible"),
"state={state:?} remote_version={remote_version:?}"
);
}
let unknown = transitioned_receipt_source(
"bucket",
"object",
"opaque-version",
rustfs_filemeta::TransitionVersionState::Unknown,
Some(&identity_hex),
);
assert!(
!TierFreeVersionReceiptSink::new()
.record(&unknown, Uuid::new_v4())
.expect("unknown remote version state should defer to recovery")
);
let missing_identity = transitioned_receipt_source(
"bucket",
"object",
"opaque-version",
rustfs_filemeta::TransitionVersionState::Exact,
None,
);
assert!(
!TierFreeVersionReceiptSink::new()
.record(&missing_identity, Uuid::new_v4())
.expect("missing durable identity should defer to recovery")
);
let invalid_exact = transitioned_receipt_source(
"bucket",
"object",
"",
rustfs_filemeta::TransitionVersionState::Exact,
Some(&identity_hex),
);
assert!(
!TierFreeVersionReceiptSink::new()
.record(&invalid_exact, Uuid::new_v4())
.expect("conflicting remote state should defer to recovery")
);
let mut conflicting = transitioned_receipt_source(
"bucket",
"object",
"opaque-version",
rustfs_filemeta::TransitionVersionState::Exact,
Some(&identity_hex),
);
Arc::make_mut(&mut conflicting.user_defined)
.insert("x-minio-internal-transition-tier-destination-id".to_string(), "33".repeat(32));
assert_eq!(
TierFreeVersionReceiptSink::new()
.record(&conflicting, Uuid::new_v4())
.expect_err("conflicting identity aliases must fail closed")
.kind(),
io::ErrorKind::InvalidData
);
let valid = transitioned_receipt_source(
"bucket",
"object",
"opaque-version",
rustfs_filemeta::TransitionVersionState::Exact,
Some(&identity_hex),
);
assert_eq!(
TierFreeVersionReceiptSink::new()
.record(&valid, Uuid::nil())
.expect_err("nil local free-version identity must be rejected")
.kind(),
io::ErrorKind::InvalidInput
);
}
#[test]
fn object_options_default_does_not_allocate_lifecycle_delete_all_journal() {
let mut opts = ObjectOptions::default();
assert!(opts.lifecycle_delete_all_journal().is_none());
assert!(opts.tier_free_version_receipt_sink.is_none());
opts.ensure_lifecycle_delete_all_journal();
assert!(opts.lifecycle_delete_all_journal().is_some());
}
+296 -16
View File
@@ -217,7 +217,8 @@ use crate::bucket::lifecycle::{
use crate::bucket::quota::reservation;
use crate::bucket::replication::{
DeleteReplicationConfigSnapshot, ReplicationLifecycleBridge, ReplicationStatusType, VersionPurgeStatusType,
replication_state_to_filemeta, replication_status_from_filemeta, version_purge_status_to_filemeta,
replication_state_to_filemeta, replication_status_from_filemeta, version_purge_status_from_filemeta,
version_purge_status_to_filemeta,
};
use crate::data_usage::quota_object_size;
use crate::diagnostics::get::GetObjectFailureReason;
@@ -283,12 +284,95 @@ fn record_transitioned_delete_cleanup_owner(bucket: &str, object: &str, batch: b
);
}
async fn acquire_single_tier_delete_lease(
/// A causal free-version receipt is only valid when the delete request really
/// removes the locked transitioned source. In particular, replication may turn
/// an otherwise successful delete into a metadata-only purge-state update, and
/// a versioned delete without a version ID writes a new delete marker instead
/// of removing the source selected by `goi`.
fn transitioned_delete_publishes_free_version(source: &ObjectInfo, delete_request: &FileInfo, skip_free_version: bool) -> bool {
if source.delete_marker
|| source.transitioned_object.status != TRANSITION_COMPLETE
|| skip_free_version
|| delete_request.skip_tier_free_version()
|| delete_request.expire_restored
|| delete_request.transition_status == TRANSITION_COMPLETE
|| delete_file_info_version_id(source.version_id) != delete_request.version_id
{
return false;
}
// Keep this predicate aligned with FileMeta::delete_version's Object
// branch: a non-delete-marker request with a nonterminal purge status (or
// mark_deleted with no purge status) updates replication metadata in place
// and never calls MetaObject::init_free_version.
let purge_status = version_purge_status_from_filemeta(delete_request.version_purge_status());
let metadata_only = !delete_request.deleted
&& ((purge_status.is_empty() && delete_request.mark_deleted)
|| (!purge_status.is_empty() && purge_status != VersionPurgeStatusType::Complete));
!metadata_only
}
fn record_committed_tier_free_version_receipt(
opts: &ObjectOptions,
bucket: &str,
object: &str,
opts: &ObjectOptions,
source: &ObjectInfo,
) -> Result<Option<TierOperationLease>> {
free_version_id: Uuid,
batch: bool,
) {
if let Some(sink) = opts.tier_free_version_receipt_sink.as_ref()
&& let Err(err) = sink.record(source, free_version_id)
{
warn!(
event = EVENT_LIFECYCLE_TRANSITIONED_DELETE_CLEANUP_OWNER,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
bucket,
object,
batch,
error = ?err,
"Failed to retain the in-memory tier free-version scheduling receipt"
);
}
record_transitioned_delete_cleanup_owner(bucket, object, batch);
}
struct TierFreeVersionReceiptCandidate {
source: ObjectInfo,
free_version_id: Uuid,
}
fn committed_tier_free_version_receipt_indices(
versions: &[FileInfoVersions],
delete_errors: &[Option<Error>],
candidates: &HashMap<usize, TierFreeVersionReceiptCandidate>,
) -> Vec<usize> {
if candidates.is_empty() {
return Vec::new();
}
let mut committed = Vec::with_capacity(candidates.len());
for group in versions {
let should_rollback = group
.versions
.iter()
.any(|version| delete_errors.get(version.idx).is_none_or(|error| error.is_some()));
if should_rollback {
continue;
}
committed.extend(
group
.versions
.iter()
.map(|version| version.idx)
.filter(|idx| candidates.contains_key(idx)),
);
}
committed
}
async fn acquire_single_tier_delete_lease(opts: &ObjectOptions, source: &ObjectInfo) -> Result<Option<TierOperationLease>> {
let Some(api) = opts.tier_delete_journal_api.as_ref() else {
return Ok(None);
};
@@ -308,7 +392,6 @@ async fn acquire_single_tier_delete_lease(
None => TierConfigMgr::acquire_operation_lease(&api.tier_config_mgr(), &source.transitioned_object.tier).await,
}
.map_err(Error::other)?;
record_transitioned_delete_cleanup_owner(bucket, object, false);
Ok(Some(lease))
}
@@ -384,6 +467,168 @@ mod scanner_publication_lease_fence_tests {
}
}
#[cfg(test)]
mod tier_free_version_receipt_eligibility_tests {
use super::*;
use crate::bucket::replication::ReplicationState;
fn transitioned_source(version_id: Option<Uuid>) -> ObjectInfo {
let mut source = ObjectInfo {
version_id,
..Default::default()
};
source.transitioned_object.status = TRANSITION_COMPLETE.to_string();
source.transitioned_object.tier = "WARM".to_string();
source.transitioned_object.name = "remote/object".to_string();
source
}
fn delete_request(version_id: Option<Uuid>) -> FileInfo {
let mut request = FileInfo {
version_id,
..Default::default()
};
request.set_tier_free_version_id(&Uuid::new_v4().to_string());
request
}
#[test]
fn accepts_exact_transitioned_source_removal_and_suspended_null_replacement() {
let version_id = Uuid::new_v4();
assert!(transitioned_delete_publishes_free_version(
&transitioned_source(Some(version_id)),
&delete_request(Some(version_id)),
false,
));
let mut suspended_null_delete = delete_request(None);
suspended_null_delete.deleted = true;
suspended_null_delete.mark_deleted = true;
assert!(transitioned_delete_publishes_free_version(
&transitioned_source(Some(Uuid::nil())),
&suspended_null_delete,
false,
));
}
#[test]
fn rejects_new_marker_version_and_non_transitioned_or_delete_marker_sources() {
let source_id = Uuid::new_v4();
assert!(!transitioned_delete_publishes_free_version(
&transitioned_source(Some(source_id)),
&delete_request(Some(Uuid::new_v4())),
false,
));
let mut ordinary = transitioned_source(Some(source_id));
ordinary.transitioned_object.status.clear();
assert!(!transitioned_delete_publishes_free_version(
&ordinary,
&delete_request(Some(source_id)),
false,
));
let mut delete_marker = transitioned_source(Some(source_id));
delete_marker.delete_marker = true;
assert!(!transitioned_delete_publishes_free_version(
&delete_marker,
&delete_request(Some(source_id)),
false,
));
}
#[test]
fn rejects_skip_restore_and_transition_metadata_updates() {
let version_id = Uuid::new_v4();
let source = transitioned_source(Some(version_id));
assert!(!transitioned_delete_publishes_free_version(
&source,
&delete_request(Some(version_id)),
true,
));
let mut skip_request = delete_request(Some(version_id));
skip_request.set_skip_tier_free_version();
assert!(!transitioned_delete_publishes_free_version(&source, &skip_request, false));
let mut restore_request = delete_request(Some(version_id));
restore_request.expire_restored = true;
assert!(!transitioned_delete_publishes_free_version(&source, &restore_request, false));
let mut transition_update = delete_request(Some(version_id));
transition_update.transition_status = TRANSITION_COMPLETE.to_string();
assert!(!transitioned_delete_publishes_free_version(&source, &transition_update, false));
}
#[test]
fn rejects_nonterminal_replication_metadata_only_update_but_accepts_complete_purge() {
let version_id = Uuid::new_v4();
let source = transitioned_source(Some(version_id));
let mut pending = delete_request(Some(version_id));
pending.replication_state_internal = Some(replication_state_to_filemeta(&ReplicationState {
version_purge_status_internal: Some("PENDING".to_string()),
..Default::default()
}));
assert!(!transitioned_delete_publishes_free_version(&source, &pending, false));
let mut mark_deleted = delete_request(Some(version_id));
mark_deleted.mark_deleted = true;
assert!(!transitioned_delete_publishes_free_version(&source, &mark_deleted, false));
let mut complete = delete_request(Some(version_id));
complete.replication_state_internal = Some(replication_state_to_filemeta(&ReplicationState {
version_purge_status_internal: Some("COMPLETE".to_string()),
..Default::default()
}));
assert!(transitioned_delete_publishes_free_version(&source, &complete, false));
}
#[test]
fn whole_physical_object_group_must_commit_before_any_receipt_is_retained() {
let candidate = || TierFreeVersionReceiptCandidate {
source: transitioned_source(Some(Uuid::new_v4())),
free_version_id: Uuid::new_v4(),
};
let candidates = HashMap::from([(0, candidate()), (2, candidate())]);
let versions = vec![
FileInfoVersions {
versions: vec![
FileInfo {
idx: 0,
..Default::default()
},
FileInfo {
idx: 1,
..Default::default()
},
],
..Default::default()
},
FileInfoVersions {
versions: vec![FileInfo {
idx: 2,
..Default::default()
}],
..Default::default()
},
];
let one_sibling_failed = vec![None, Some(Error::other("injected quorum failure")), None];
assert_eq!(
committed_tier_free_version_receipt_indices(&versions, &one_sibling_failed, &candidates),
vec![2],
"a sibling failure must suppress every receipt from the rolled-back xl.meta group"
);
assert_eq!(
committed_tier_free_version_receipt_indices(&versions, &[None, None, None], &candidates),
vec![0, 2],
"independent fully committed groups should retain their sparse receipts"
);
}
}
struct PutObjectCommitCancellation {
token: CancellationToken,
armed: bool,
@@ -6977,7 +7222,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
};
let mut vers_map: HashMap<&String, FileInfoVersions> = HashMap::new();
let mut tier_reference_leases: Vec<(usize, String, Option<TierDestinationId>)> = Vec::new();
let mut transitioned_cleanup_items = vec![false; objects.len()];
let mut tier_free_version_receipt_candidates: HashMap<usize, TierFreeVersionReceiptCandidate> = HashMap::new();
for (i, dobj) in objects.iter().enumerate() {
if del_errs[i].is_some() {
@@ -7069,7 +7314,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
{
match tier_destination_id_from_metadata(&goi.user_defined) {
Ok(identity) => {
transitioned_cleanup_items[i] = true;
tier_reference_leases.push((i, goi.transitioned_object.tier.clone(), identity));
}
Err(err) => {
@@ -7110,7 +7354,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
..Default::default()
};
vr.set_tier_free_version_id(&Uuid::new_v4().to_string());
let tier_free_version_id = Uuid::new_v4();
vr.set_tier_free_version_id(&tier_free_version_id.to_string());
// Delete
// del_objects[i].object_name.clone_from(&vr.name);
@@ -7195,6 +7440,19 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
});
}
if opts.tier_free_version_receipt_sink.is_some()
&& !dobj.synthetic_version_id
&& transitioned_delete_publishes_free_version(&goi, &vr, opts.skip_free_version)
{
tier_free_version_receipt_candidates.insert(
i,
TierFreeVersionReceiptCandidate {
source: goi,
free_version_id: tier_free_version_id,
},
);
}
// Only add to vers_map if we hold the lock
if locked_objects.contains(&dobj.object_name) {
vers_map.insert(&dobj.object_name, v);
@@ -7271,12 +7529,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
}
}
for (idx, transitioned) in transitioned_cleanup_items.into_iter().enumerate() {
if transitioned && del_errs[idx].is_none() {
record_transitioned_delete_cleanup_owner(bucket, &decode_dir_object(&objects[idx].object_name), true);
}
}
// Keep backend generations pinned through the source mutation, its
// free-version write quorum, and any local rollback. Ordinary
// single/batch deletes never transfer cleanup ownership to a journal.
@@ -7388,6 +7640,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
self.record_capacity_scope_if_needed(opts.capacity_scope_token, &disks);
let mut rollback_futures = Vec::new();
let committed_receipt_indices =
committed_tier_free_version_receipt_indices(&vers, &del_errs, &tier_free_version_receipt_candidates);
for fi_vers in &vers {
// delete_versions commits one xl.meta per object group, so rollback must use the same boundary.
let should_rollback = fi_vers.versions.iter().any(|fi| del_errs[fi.idx].is_some());
@@ -7462,6 +7716,20 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
join_all(rollback_futures).await;
for idx in committed_receipt_indices {
let Some(candidate) = tier_free_version_receipt_candidates.remove(&idx) else {
continue;
};
record_committed_tier_free_version_receipt(
&opts,
bucket,
&decode_dir_object(&objects[idx].object_name),
&candidate.source,
candidate.free_version_id,
true,
);
}
// TODO(backlog): support partial object deletion for multi-part objects
if dist_erasure {
@@ -7813,7 +8081,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
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(bucket, object, &opts, &goi).await?;
let _tier_delete_lease = acquire_single_tier_delete_lease(&opts, &goi).await?;
if opts.skip_free_version {
fi.set_skip_tier_free_version();
}
@@ -7823,6 +8091,12 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
#[cfg(test)]
pause_delete_object_commit_after_publish(bucket, object).await;
if opts.tier_free_version_receipt_sink.is_some()
&& transitioned_delete_publishes_free_version(&goi, &fi, opts.skip_free_version)
{
record_committed_tier_free_version_receipt(&opts, bucket, object, &goi, find_vid, false);
}
let disks = self.disk_inventory().await;
self.record_capacity_scope_if_needed(opts.capacity_scope_token, &disks);
@@ -7855,7 +8129,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
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(bucket, object, &opts, &goi).await?;
let _tier_delete_lease = acquire_single_tier_delete_lease(&opts, &goi).await?;
if opts.skip_free_version {
dfi.set_skip_tier_free_version();
}
@@ -7865,6 +8139,12 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
#[cfg(test)]
pause_delete_object_commit_after_publish(bucket, object).await;
if opts.tier_free_version_receipt_sink.is_some()
&& transitioned_delete_publishes_free_version(&goi, &dfi, opts.skip_free_version)
{
record_committed_tier_free_version_receipt(&opts, bucket, object, &goi, find_vid, false);
}
let disks = self.disk_inventory().await;
self.record_capacity_scope_if_needed(opts.capacity_scope_token, &disks);
+559 -5
View File
@@ -10496,6 +10496,8 @@ mod tests {
bucket: &str,
object: &str,
restore_before_delete: bool,
causal_enqueue: bool,
delete_with_journal: bool,
) {
let temp_dir = tempfile::tempdir().expect("create transitioned delete store dir");
let (ctx, store, _shutdown) =
@@ -10559,11 +10561,51 @@ mod tests {
);
}
backend.set_remove_failure(true);
store
.delete_object_with_tier_delete_journal(bucket, object, ObjectOptions::default())
if causal_enqueue {
ExpiryState::resize_workers(1, store.clone()).await;
}
backend.set_remove_failure(!causal_enqueue);
if delete_with_journal {
store
.delete_object_with_tier_delete_journal(bucket, object, ObjectOptions::default())
.await
.expect("transitioned source journal-wrapper delete should commit");
} else {
store
.delete_object(bucket, object, ObjectOptions::default())
.await
.expect("transitioned source plain object-layer delete should commit");
}
if causal_enqueue {
tokio::time::timeout(Duration::from_secs(30), async {
loop {
let metadata_absent = store.pools[0]
.get_disks_by_key(object)
.load_file_info_versions_exact(bucket, object)
.await
.expect("causal free-version cleanup metadata should remain readable")
.is_none();
if metadata_absent && backend.remove_versions().await.len() == 1 {
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("transitioned source delete should commit");
.expect("committed free-version should be cleaned without a recovery scan");
assert_eq!(backend.object_count().await, 0, "causal cleanup should remove the remote object");
assert_eq!(
tier_delete_journal_count(store.clone()).await,
0,
"ordinary causal cleanup must not create a journal"
);
store
.delete_bucket(bucket, &DeleteBucketOptions::default())
.await
.expect("bucket delete should succeed after causal free-version cleanup");
return;
}
let local_versions = store.pools[0]
.get_disks_by_key(object)
@@ -10642,6 +10684,8 @@ mod tests {
"transitioned-delete-journal-owner-bucket",
"transition/archive.bin",
false,
false,
true,
)
.await;
}
@@ -10656,6 +10700,40 @@ mod tests {
"restored-transitioned-delete-journal-owner-bucket",
"transition/archive.bin",
true,
false,
true,
)
.await;
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn transitioned_delete_causally_enqueues_free_version() {
run_transitioned_delete_free_version_owner_case(
"transitioned-delete-causal-enqueue",
"DELETE-CAUSAL",
"transitioned-delete-causal-enqueue-bucket",
"transition/archive.bin",
false,
true,
false,
)
.await;
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn restored_transitioned_delete_causally_enqueues_free_version() {
run_transitioned_delete_free_version_owner_case(
"restored-transitioned-delete-causal-enqueue",
"RESTORE-DELETE-CAUSAL",
"restored-transitioned-delete-causal-enqueue-bucket",
"transition/archive.bin",
true,
true,
true,
)
.await;
}
@@ -12335,12 +12413,262 @@ mod tests {
.is_none(),
"free-version recovery must remove the exact cleanup owner"
);
let causal = "causal.bin";
let mut causal_reader = PutObjReader::from_vec(vec![b'c'; 1024 * 1024]);
let causal_source = store
.put_object(bucket, causal, &mut causal_reader, &ObjectOptions::default())
.await
.expect("causal batch source should be written");
store
.transition_object(
bucket,
causal,
&ObjectOptions {
transition: TransitionOptions {
status: TRANSITION_PENDING.to_string(),
tier: tier_name.to_string(),
etag: causal_source.etag.clone().expect("causal batch source should have an etag"),
..Default::default()
},
mod_time: causal_source.mod_time,
..Default::default()
},
)
.await
.expect("causal batch source transition should commit");
let (_deleted, errors) = store
.delete_objects(
bucket,
vec![
ObjectToDelete {
object_name: causal.to_string(),
..Default::default()
},
ObjectToDelete {
object_name: causal.to_string(),
..Default::default()
},
],
ObjectOptions::default(),
)
.await;
assert!(
errors.iter().all(Option::is_none),
"duplicate causal batch deletes should remain idempotent: {errors:?}"
);
tokio::time::timeout(Duration::from_secs(30), async {
loop {
let metadata_absent = store.pools[0]
.get_disks_by_key(causal)
.load_file_info_versions_exact(bucket, causal)
.await
.expect("causal batch cleanup metadata should remain readable")
.is_none();
if metadata_absent && backend.remove_versions().await.len() >= 2 {
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("batch free-version receipt should converge without another recovery scan");
assert_eq!(
backend.remove_versions().await.len(),
2,
"duplicate batch requests must cause only one remote delete for the causal object"
);
crate::bucket::metadata_sys::update_in(
&ctx,
bucket,
BUCKET_VERSIONING_CONFIG,
b"<VersioningConfiguration><Status>Enabled</Status></VersioningConfiguration>".to_vec(),
)
.await
.expect("causal batch bucket versioning should be enabled");
let versioned_causal = "versioned-causal.bin";
let mut versioned_reader = PutObjReader::from_vec(vec![b'v'; 1024 * 1024]);
let versioned_source = store
.put_object(
bucket,
versioned_causal,
&mut versioned_reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("versioned causal batch source should be written");
let versioned_source_id = versioned_source
.version_id
.expect("versioned causal batch source should have an identity");
store
.transition_object(
bucket,
versioned_causal,
&ObjectOptions {
version_id: Some(versioned_source_id.to_string()),
versioned: true,
transition: TransitionOptions {
status: TRANSITION_PENDING.to_string(),
tier: tier_name.to_string(),
etag: versioned_source
.etag
.clone()
.expect("versioned causal batch source should have an etag"),
..Default::default()
},
mod_time: versioned_source.mod_time,
..Default::default()
},
)
.await
.expect("versioned causal batch source transition should commit");
let (_deleted, errors) = store
.delete_objects(
bucket,
vec![ObjectToDelete {
object_name: versioned_causal.to_string(),
version_id: Some(versioned_source_id),
..Default::default()
}],
ObjectOptions::default(),
)
.await;
assert!(
errors.iter().all(Option::is_none),
"explicit-version causal batch delete should commit: {errors:?}"
);
tokio::time::timeout(Duration::from_secs(30), async {
loop {
let metadata_absent = store.pools[0]
.get_disks_by_key(versioned_causal)
.load_file_info_versions_exact(bucket, versioned_causal)
.await
.expect("versioned causal batch cleanup metadata should remain readable")
.is_none();
if metadata_absent && backend.remove_versions().await.len() == 3 {
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("explicit-version batch receipt should converge without a recovery scan");
store
.delete_bucket(bucket, &DeleteBucketOptions::default())
.await
.expect("batch source bucket should be physically empty");
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn batch_transitioned_delete_aggregate_error_still_enqueues_committed_free_version() {
let temp_dir = tempfile::tempdir().expect("create aggregate-error batch delete store dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "batch-transitioned-aggregate-error", &[4, 4]))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let tier_name = "BATCH-AGGREGATE-ERROR";
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
let bucket = "batch-transitioned-aggregate-error-bucket";
let object = "archive.bin";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("aggregate-error source bucket should be created");
let mut reader = PutObjReader::from_vec(vec![b'a'; 1024 * 1024]);
let source = store.pools[0]
.put_object(bucket, object, &mut reader, &ObjectOptions::default())
.await
.expect("aggregate-error source should be written");
store.pools[0]
.transition_object(
bucket,
object,
&ObjectOptions {
transition: TransitionOptions {
status: TRANSITION_PENDING.to_string(),
tier: tier_name.to_string(),
etag: source.etag.clone().expect("aggregate-error source should have an etag"),
..Default::default()
},
mod_time: source.mod_time,
..Default::default()
},
)
.await
.expect("aggregate-error source should transition");
// Model a data-movement copy: both pools own the same logical source
// and exact remote tuple, but each batch delete creates its own local
// free-version UUID in the shared request sink.
for disk_index in 0..4 {
let source_meta = temp_dir
.path()
.join(format!("pool0/set0/disk{disk_index}/{bucket}/{object}/{STORAGE_FORMAT_FILE}"));
let target_meta = temp_dir
.path()
.join(format!("pool1/set0/disk{disk_index}/{bucket}/{object}/{STORAGE_FORMAT_FILE}"));
tokio::fs::create_dir_all(target_meta.parent().expect("target xl.meta should have a parent"))
.await
.expect("second-pool object directory should be created");
tokio::fs::copy(&source_meta, &target_meta)
.await
.expect("transitioned xl.meta should copy exactly to the second pool");
}
ExpiryState::resize_workers(1, store.clone()).await;
let injection = crate::store::object::BatchDeletePoolErrorInjection::install(
bucket,
1,
vec![(object.to_string(), StorageError::ErasureWriteQuorum)],
);
let (deleted, errors) = store
.delete_objects(
bucket,
vec![ObjectToDelete {
object_name: object.to_string(),
..Default::default()
}],
ObjectOptions::default(),
)
.await;
assert_eq!(injection.observed(), 1, "the second pool should inject one post-commit aggregate error");
assert_eq!(errors, vec![Some(StorageError::ErasureWriteQuorum)]);
assert!(deleted[0].found, "the aggregate error must retain the committed pool result");
drop(injection);
tokio::time::timeout(Duration::from_secs(30), async {
loop {
let mut metadata_absent = true;
for pool in &store.pools {
metadata_absent &= pool
.get_disks_by_key(object)
.load_file_info_versions_exact(bucket, object)
.await
.expect("aggregate-error cleanup metadata should remain readable")
.is_none();
}
if metadata_absent && backend.remove_count().await == 1 {
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("aggregate failure must not suppress committed receipt dispatch");
assert_eq!(backend.object_count().await, 0, "the shared remote object should be removed exactly once");
store
.delete_bucket(bucket, &DeleteBucketOptions::default())
.await
.expect("aggregate-error bucket should be physically empty");
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
@@ -12572,6 +12900,217 @@ mod tests {
.expect("retry should leave the source bucket empty");
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn batch_transitioned_delete_post_commit_failures_roll_back_without_free_version_receipt() {
let temp_dir = tempfile::tempdir().expect("create failed batch delete store dir");
let (ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "batch-delete-local-failure", &[4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let tier_name = "BATCH-DELETE-LOCAL-FAIL";
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await;
let bucket = "batch-delete-local-failure-bucket";
let object = "archive.bin";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("failed batch source bucket should be created");
crate::bucket::metadata_sys::update_in(
&ctx,
bucket,
BUCKET_VERSIONING_CONFIG,
b"<VersioningConfiguration><Status>Enabled</Status></VersioningConfiguration>".to_vec(),
)
.await
.expect("failed batch bucket versioning should be enabled");
let mut transitioned_reader = PutObjReader::from_vec(vec![b't'; 1024 * 1024]);
let transitioned_source = store
.put_object(
bucket,
object,
&mut transitioned_reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("failed batch transitioned version should be written");
let transitioned_version_id = transitioned_source
.version_id
.expect("failed batch transitioned source should have a version identity");
store
.transition_object(
bucket,
object,
&ObjectOptions {
version_id: Some(transitioned_version_id.to_string()),
versioned: true,
transition: TransitionOptions {
status: TRANSITION_PENDING.to_string(),
tier: tier_name.to_string(),
etag: transitioned_source
.etag
.clone()
.expect("failed batch transitioned source should have an etag"),
..Default::default()
},
mod_time: transitioned_source.mod_time,
..Default::default()
},
)
.await
.expect("failed batch source version should transition");
let mut ordinary_reader = PutObjReader::from_vec(vec![b'o'; 1024 * 1024]);
let ordinary_source = store
.put_object(
bucket,
object,
&mut ordinary_reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("failed batch ordinary sibling should be written");
let ordinary_version_id = ordinary_source
.version_id
.expect("failed batch ordinary sibling should have a version identity");
let delete_requests = || {
vec![
ObjectToDelete {
object_name: object.to_string(),
version_id: Some(transitioned_version_id),
..Default::default()
},
ObjectToDelete {
object_name: object.to_string(),
version_id: Some(ordinary_version_id),
..Default::default()
},
]
};
let set = store.pools[0].get_disks_by_key(object);
let disks = set.disks.read().await;
assert_eq!(disks.len(), 4, "the rollback fixture must use four disks");
// Keep discovery fully online, then make a quorum of disks report an
// error only after their batch metadata commit has completed.
for disk in disks.iter().take(3) {
let disk = disk.as_ref().expect("injected rollback disks should be online");
crate::disk::local::set_delete_version_fail_after_commit(disk.path().as_path(), object);
}
drop(disks);
let receipt_sink = crate::object_api::TierFreeVersionReceiptSink::new();
let (_deleted, errors) = store
.delete_objects(
bucket,
delete_requests(),
ObjectOptions {
tier_free_version_receipt_sink: Some(receipt_sink.clone()),
..Default::default()
},
)
.await;
assert_eq!(
errors,
vec![Some(StorageError::Unexpected), Some(StorageError::Unexpected),],
"three post-commit disk errors must fail batch delete before receipts publish"
);
assert!(
receipt_sink
.drain()
.expect("the test-owned failed-batch sink should drain exactly once")
.is_empty(),
"a rolled-back physical group must publish no cleanup receipt"
);
let retained_transitioned = store
.get_object_info(
bucket,
object,
&ObjectOptions {
version_id: Some(transitioned_version_id.to_string()),
versioned: true,
..Default::default()
},
)
.await
.expect("failed batch delete must restore the transitioned sibling");
assert_eq!(retained_transitioned.transitioned_object.status, rustfs_filemeta::TRANSITION_COMPLETE);
let retained_ordinary = store
.get_object_info(
bucket,
object,
&ObjectOptions {
version_id: Some(ordinary_version_id.to_string()),
versioned: true,
..Default::default()
},
)
.await
.expect("failed batch delete must restore the ordinary sibling");
assert_ne!(retained_ordinary.transitioned_object.status, rustfs_filemeta::TRANSITION_COMPLETE);
let retained_versions = set
.load_file_info_versions_exact(bucket, object)
.await
.expect("rolled-back batch metadata should decode")
.expect("rolled-back batch source should remain on disk");
assert_eq!(
retained_versions
.versions
.iter()
.chain(retained_versions.free_versions.iter())
.filter(|version| version.tier_free_version())
.count(),
0,
"failed batch quorum must not retain a free-version owner"
);
let retained_version_ids = retained_versions
.versions
.iter()
.filter_map(|version| version.version_id)
.collect::<std::collections::HashSet<_>>();
assert_eq!(
retained_version_ids,
std::collections::HashSet::from([transitioned_version_id, ordinary_version_id]),
"the physical-group rollback must restore both explicit siblings"
);
assert_eq!(backend.object_count().await, 1, "failed batch commit must retain the remote object");
assert_eq!(backend.remove_count().await, 0, "failed batch commit must not dispatch remote cleanup");
ExpiryState::resize_workers(1, store.clone()).await;
let (_deleted, retry_errors) = store
.delete_objects(bucket, delete_requests(), ObjectOptions::default())
.await;
assert!(
retry_errors.iter().all(Option::is_none),
"retry after disk recovery should commit: {retry_errors:?}"
);
tokio::time::timeout(Duration::from_secs(30), async {
loop {
let metadata_absent = set
.load_file_info_versions_exact(bucket, object)
.await
.expect("retry cleanup metadata should remain readable")
.is_none();
if metadata_absent && backend.remove_count().await == 1 {
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("successful batch retry should converge without a recovery scan");
store
.delete_bucket(bucket, &DeleteBucketOptions::default())
.await
.expect("successful batch retry should leave the bucket empty");
shutdown.cancel();
}
#[cfg(feature = "test-util")]
async fn run_multi_pool_same_remote_tuple_delete_case(batch: bool) {
let temp_dir = tempfile::tempdir().expect("create shared-tuple multi-pool store dir");
@@ -12645,8 +13184,10 @@ mod tests {
);
assert_eq!(backend.object_count().await, 1);
let receipt_sink = crate::object_api::TierFreeVersionReceiptSink::new();
let mut delete_opts = ObjectOptions {
tier_delete_journal_api: Some(store.clone()),
tier_free_version_receipt_sink: Some(receipt_sink.clone()),
..Default::default()
};
if batch {
@@ -12774,7 +13315,20 @@ mod tests {
);
backend.set_remove_failure(false);
wait_for_tier_free_version_recovery(store.clone(), &backend, 1).await;
let receipts = receipt_sink
.drain()
.expect("the simulated outer multi-pool wrapper should drain exactly once");
assert_eq!(
receipts.len(),
1,
"the same physical key and remote tuple must collapse to one causal task"
);
assert_eq!(
crate::bucket::lifecycle::bucket_lifecycle_ops::enqueue_committed_free_versions(&store, receipts).await,
1,
"the committed shared-tuple task should enter the running worker"
);
wait_for_expiry_workers_idle(&store).await;
assert_eq!(backend.remove_count().await, 1, "shared remote tuple should be deleted exactly once");
for pool_idx in 0..2 {
assert!(
+2
View File
@@ -425,6 +425,8 @@ pub(crate) mod init_format;
pub(crate) mod list_objects;
mod multipart;
mod object;
#[cfg(any(test, feature = "test-util"))]
pub use object::DeleteAfterObjectLockSnapshotBarrier;
pub(crate) use object::{
DecommissionFixedReadAnchor, ObjectLockDiagGuard, RemoteTuplePublicationCommitGuard, RemoteTuplePublicationFence,
SourceCleanupMutationFence, tiered_data_movement_source_matches,
+94 -19
View File
@@ -14,7 +14,7 @@
use super::*;
use crate::bucket::lifecycle::{
bucket_lifecycle_ops::eval_action_from_lifecycle,
bucket_lifecycle_ops::{enqueue_committed_free_versions, eval_action_from_lifecycle},
get_expiry_configs,
tier_delete_journal::{
ActiveTierDeleteDispatch, EVENT_LIFECYCLE_TIER_DELETE_JOURNAL, LOG_COMPONENT_ECSTORE, LOG_SUBSYSTEM_LIFECYCLE,
@@ -39,6 +39,7 @@ use crate::core::pools::{DecommissionCapacityOwner, ensure_decommission_capacity
use crate::disk::OldCurrentSize;
use crate::object_api::{
NamespaceLockFence, ObjectLockConfigSnapshot, ScannerPublicationCommitScopeGuard, ScannerPublicationCommitState,
TierFreeVersionReceiptSink,
};
use crate::services::notification_sys::acquire_tier_delete_journal_fleet_proof;
use crate::services::tier::tier::{TierConfigMgr, TierDestinationId, TierOperationLease, tier_destination_id_from_metadata};
@@ -71,6 +72,26 @@ const RECURSIVE_DELETE_VERSION_SCAN_PAGE_SIZE: i32 = 1000;
#[cfg(test)]
const RECURSIVE_DELETE_VERSION_SCAN_PAGE_SIZE: i32 = 2;
fn install_tier_free_version_receipt_sink(opts: &mut ObjectOptions) -> Option<TierFreeVersionReceiptSink> {
if opts.tier_free_version_receipt_sink.is_some() || opts.skip_free_version || opts.delete_prefix {
return None;
}
let sink = TierFreeVersionReceiptSink::new();
opts.tier_free_version_receipt_sink = Some(sink.clone());
Some(sink)
}
async fn enqueue_recorded_tier_free_versions(store: &ECStore, sink: Option<TierFreeVersionReceiptSink>) -> usize {
let Some(sink) = sink else {
return 0;
};
let Ok(receipts) = sink.drain() else {
return 0;
};
enqueue_committed_free_versions(store, receipts).await
}
fn build_tier_delete_journal_entry(
bucket: &str,
object: &str,
@@ -1293,7 +1314,7 @@ fn should_create_delete_marker_for_missing_object(opts: &ObjectOptions) -> bool
(opts.versioned || opts.version_suspended) && opts.version_id.is_none() && !opts.delete_marker && !opts.data_movement
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
struct DeleteAfterObjectLockSnapshotBarrierState {
bucket: String,
arrived: tokio::sync::Notify,
@@ -1302,19 +1323,19 @@ struct DeleteAfterObjectLockSnapshotBarrierState {
namespace_acquired: AtomicBool,
}
#[cfg(test)]
pub(crate) struct DeleteAfterObjectLockSnapshotBarrier {
#[cfg(any(test, feature = "test-util"))]
pub struct DeleteAfterObjectLockSnapshotBarrier {
state: Arc<DeleteAfterObjectLockSnapshotBarrierState>,
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
static DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER: std::sync::OnceLock<
std::sync::Mutex<Option<Arc<DeleteAfterObjectLockSnapshotBarrierState>>>,
> = std::sync::OnceLock::new();
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
impl DeleteAfterObjectLockSnapshotBarrier {
pub(crate) fn install(bucket: &str) -> Self {
pub fn install(bucket: &str) -> Self {
let state = Arc::new(DeleteAfterObjectLockSnapshotBarrierState {
bucket: bucket.to_string(),
arrived: tokio::sync::Notify::new(),
@@ -1331,15 +1352,15 @@ impl DeleteAfterObjectLockSnapshotBarrier {
Self { state }
}
pub(crate) async fn wait_until_paused(&self) {
pub async fn wait_until_paused(&self) {
self.state.arrived.notified().await;
}
pub(crate) fn release(&self) {
pub fn release(&self) {
self.state.release.notify_one();
}
pub(crate) async fn release_and_wait_until_namespace_pending(&self) {
pub 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)
@@ -1347,12 +1368,12 @@ impl DeleteAfterObjectLockSnapshotBarrier {
.expect("delete should proceed to its namespace lock after leaving the snapshot barrier");
}
pub(crate) fn namespace_acquired(&self) -> bool {
pub fn namespace_acquired(&self) -> bool {
self.state.namespace_acquired.load(Ordering::Acquire)
}
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
impl Drop for DeleteAfterObjectLockSnapshotBarrier {
fn drop(&mut self) {
self.state.release.notify_one();
@@ -1365,7 +1386,7 @@ impl Drop for DeleteAfterObjectLockSnapshotBarrier {
}
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
async fn pause_delete_after_object_lock_snapshot(bucket: &str) {
let state = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
@@ -1377,11 +1398,24 @@ async fn pause_delete_after_object_lock_snapshot(bucket: &str) {
if let Some(state) = state {
state.arrived.notify_one();
state.release.notified().await;
}
}
#[cfg(any(test, feature = "test-util"))]
fn notify_delete_namespace_pending(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_pending.notify_one();
}
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
fn notify_delete_namespace_acquired(bucket: &str) {
let state = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
@@ -2564,6 +2598,10 @@ impl ECStore {
let diag_enabled = is_object_lock_diag_enabled();
let ns_lock = self.handle_new_ns_lock(bucket, object).await?;
let acquire_start = Instant::now();
#[cfg(any(test, feature = "test-util"))]
if matches!(op, "delete_object" | "delete_objects") {
notify_delete_namespace_pending(bucket);
}
let guard = ns_lock
.get_write_lock(get_lock_acquire_timeout())
.await
@@ -4155,7 +4193,16 @@ impl ECStore {
opts: ObjectOptions,
tier_journal_api: Option<Arc<ECStore>>,
) -> Result<ObjectInfo> {
Box::pin(self.handle_delete_object_with_journal_inner(bucket, object, opts, tier_journal_api)).await
Box::pin(async move {
let mut opts = opts;
let receipt_sink = install_tier_free_version_receipt_sink(&mut opts);
let result = self
.handle_delete_object_with_journal_inner(bucket, object, opts, tier_journal_api)
.await;
enqueue_recorded_tier_free_versions(self, receipt_sink).await;
result
})
.await
}
async fn handle_delete_object_with_journal_inner(
@@ -4229,7 +4276,7 @@ impl ECStore {
if opts.delete_prefix && opts.expected_bucket_incarnation_id.is_none() {
opts.expected_bucket_incarnation_id = current_bucket_incarnation_id;
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
pause_delete_after_object_lock_snapshot(bucket).await;
if opts.delete_prefix && !opts.delete_prefix_object {
@@ -4243,7 +4290,7 @@ impl ECStore {
} else {
None
};
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
if _object_lock_guard.is_some() {
notify_delete_namespace_acquired(bucket);
}
@@ -4533,6 +4580,25 @@ impl ECStore {
objects: Vec<ObjectToDelete>,
opts: ObjectOptions,
tier_journal_api: Option<Arc<ECStore>>,
) -> (Vec<DeletedObject>, Vec<Option<Error>>, Vec<Option<DeleteAccounting>>) {
Box::pin(async move {
let mut opts = opts;
let receipt_sink = install_tier_free_version_receipt_sink(&mut opts);
let result = self
.handle_delete_objects_with_journal_and_accounting_inner(bucket, objects, opts, tier_journal_api)
.await;
enqueue_recorded_tier_free_versions(self, receipt_sink).await;
result
})
.await
}
async fn handle_delete_objects_with_journal_and_accounting_inner(
&self,
bucket: &str,
objects: Vec<ObjectToDelete>,
opts: ObjectOptions,
tier_journal_api: Option<Arc<ECStore>>,
) -> (Vec<DeletedObject>, Vec<Option<Error>>, Vec<Option<DeleteAccounting>>) {
// encode object name
let objects: Vec<ObjectToDelete> = objects
@@ -4617,7 +4683,7 @@ impl ECStore {
StorageError::BucketNotFound(bucket.to_string()),
);
}
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
if current_bucket_incarnation_id.is_some() {
pause_delete_after_object_lock_snapshot(bucket).await;
}
@@ -4625,7 +4691,7 @@ impl ECStore {
Ok(guards) => guards,
Err(err) => return return_batch_delete_lock_error_with_accounting(objects.as_slice(), err),
};
#[cfg(test)]
#[cfg(any(test, feature = "test-util"))]
if !_object_lock_guards.is_empty() {
notify_delete_namespace_acquired(bucket);
}
@@ -7528,6 +7594,15 @@ mod tests {
);
drop(unified_future);
let batch_future =
store.handle_delete_objects_with_journal_and_accounting("bucket", Vec::new(), ObjectOptions::default(), None);
let batch_future_size = std::mem::size_of_val(&batch_future);
assert!(
batch_future_size <= 4 * 1024,
"batch delete handler future must remain stack-bounded; measured {batch_future_size} bytes"
);
drop(batch_future);
let outer_future = store.handle_delete_object("bucket", "object", ObjectOptions::default());
let outer_future_size = std::mem::size_of_val(&outer_future);
assert!(
+167 -109
View File
@@ -40,14 +40,15 @@ use uuid::Uuid;
mod storage_api;
use storage_api::lifecycle::{
BUCKET_LIFECYCLE_CONFIG, BucketOperations, BucketOptions, BucketVersioningSys, CompletePart, DiskOption, ECStore,
EcstoreError, Endpoint, EndpointServerPools, Endpoints, IlmAction, LcEvent, LcEventSrc, ListOperations as _,
MakeBucketOptions, MockWarmBackend, MultipartOperations as _, ObjectIO as _, ObjectOperations as _, PoolEndpoints,
STORAGE_FORMAT_FILE, TRANSITION_PENDING, TransitionCleanupStoreBarrier, TransitionOptions, assert_transition_meta_consistent,
enqueue_transition_for_existing_objects, expire_transitioned_object, free_version_count, get_bucket_metadata,
get_global_tier_config_mgr, init_background_expiry, init_bucket_metadata_sys, init_local_disks, is_err_object_not_found,
is_err_version_not_found, new_disk, path2_bucket_object_with_base_path, recover_transition_transaction_records,
register_mock_tier_util, update_bucket_metadata, wait_for_free_version_absence,
BUCKET_LIFECYCLE_CONFIG, BucketOperations, BucketOptions, BucketVersioningSys, CompletePart,
DeleteAfterObjectLockSnapshotBarrier, DiskOption, ECStore, EcstoreError, Endpoint, EndpointServerPools, Endpoints,
ExpiryState, IlmAction, LcEvent, LcEventSrc, ListOperations as _, MakeBucketOptions, MockWarmBackend,
MultipartOperations as _, ObjectIO as _, ObjectOperations as _, PoolEndpoints, STORAGE_FORMAT_FILE, TRANSITION_PENDING,
TransitionCleanupStoreBarrier, TransitionOptions, assert_transition_meta_consistent, enqueue_transition_for_existing_objects,
expire_transitioned_object, free_version_count, get_bucket_metadata, get_global_tier_config_mgr, init_background_expiry,
init_bucket_metadata_sys, init_local_disks, is_err_object_not_found, is_err_version_not_found, new_disk,
path2_bucket_object_with_base_path, recover_transition_transaction_records, register_mock_tier_util, update_bucket_metadata,
wait_for_free_version_absence,
};
static GLOBAL_ENV: OnceLock<(Vec<PathBuf>, Arc<ECStore>)> = OnceLock::new();
@@ -601,24 +602,21 @@ mod serial_tests {
/// persisted free-version recovery -- so no live local metadata ever points
/// at an already-removed remote version.
///
/// This test pins the FIXED contract two complementary ways, both
/// revert-proof (reverting to remote-first ordering turns them red):
/// This test pins the fixed contract with deterministic GET and DELETE
/// barriers (reverting to remote-first ordering turns it red):
///
/// 1. Ordering (deterministic): immediately after
/// `expire_transitioned_object` returns, the remote tier object is still
/// present and the mock recorded **zero** remote `remove` calls --
/// proving the local delete happened with no synchronous remote removal
/// (local-first). Remote-first ordering loses the object and records a
/// `remove`.
/// 2. Concurrent GET (user-visible): a tight GET loop runs concurrently
/// with the expiry; every observation must be either a full, correct
/// body (GET won) or a clean object/version-not-found (expiry won). A
/// tier-fetch failure -- the #3491 symptom -- is never tolerated.
/// 1. A GET that already resolved the transitioned metadata keeps its read
/// lock and returns the complete remote body while expiry waits.
/// 2. Expiry returns after committing the local free-version without
/// waiting for the post-commit worker's remote DELETE. While that DELETE
/// is paused, the durable marker and remote body must both still exist.
/// 3. A later GET observes a clean object/version-not-found, never a tier
/// fetch or read-quorum failure.
#[tokio::test(flavor = "multi_thread", worker_threads = 1)]
#[serial]
#[ignore = "global-state ILM integration test: runs serialized in the CI ILM Integration (serial) lane, see ci.yml test-ilm-integration-serial and rustfs/backlog#1148 (ilm-2)"]
async fn test_expire_transitioned_object_never_races_concurrent_get() {
let (_disk_paths, ecstore) = setup_isolated_test_env(false).await;
let (disk_paths, ecstore) = setup_isolated_test_env(false).await;
let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
let backend = register_mock_tier(&tier_name).await;
@@ -664,59 +662,7 @@ mod serial_tests {
"the regression must exercise an unversioned remote tier"
);
// Concurrent GET loop: hammer GET while the expiry runs. Every outcome
// must be a full correct body or a clean not-found -- never a tier-fetch
// failure.
let get_store = ecstore.clone();
let get_bucket = bucket_name.clone();
let get_object = object_name.to_string();
let expected = payload.clone();
let get_loop = tokio::spawn(async move {
let mut saw_full_body = 0usize;
let mut saw_not_found = 0usize;
for _ in 0..400 {
match get_store
.get_object_reader(
get_bucket.as_str(),
get_object.as_str(),
None,
http::HeaderMap::new(),
&ObjectOptions::default(),
)
.await
{
Ok(mut reader) => {
let mut data = Vec::new();
match reader.stream.read_to_end(&mut data).await {
Ok(_) => {
assert_eq!(
data, expected,
"a successful GET during expiry must return the complete, correct body"
);
saw_full_body += 1;
}
Err(err) => {
panic!("GET during expiry streamed a truncated/failed body (expire/GET race regression): {err:?}")
}
}
}
Err(err) => {
let ec: &EcstoreError = &err;
assert!(
is_err_object_not_found(ec) || is_err_version_not_found(ec),
"GET during expiry may only fail with a clean object/version-not-found (expiry won \
the race); a tier-fetch failure is the #3491 regression: {err:?}"
);
saw_not_found += 1;
}
}
tokio::task::yield_now().await;
}
(saw_full_body, saw_not_found)
});
// Run the exact expiry action the scanner drives for a transitioned
// current version.
ExpiryState::resize_workers(1, ecstore.clone()).await;
let lc_event = LcEvent {
action: IlmAction::DeleteAction,
..Default::default()
@@ -725,39 +671,127 @@ mod serial_tests {
.bucket_incarnation_id(bucket_name.as_str())
.await
.expect("read bucket incarnation");
expire_transitioned_object(ecstore.clone(), &oi, &lc_event, &LcEventSrc::Scanner, bucket_incarnation_id)
// Pause one real tier GET after it has resolved local transition
// metadata. The reader still owns the object read lock, so local expiry
// cannot commit until this GET finishes.
let get_barrier = backend.arm_get_barrier().await;
let get_store = ecstore.clone();
let get_bucket = bucket_name.clone();
let get_object = object_name.to_string();
let in_flight_get = tokio::spawn(async move {
let mut reader = get_store
.get_object_reader(
get_bucket.as_str(),
get_object.as_str(),
None,
http::HeaderMap::new(),
&ObjectOptions::default(),
)
.await
.map_err(|err| format!("in-flight GET failed before streaming: {err:?}"))?;
let mut data = Vec::new();
reader
.stream
.read_to_end(&mut data)
.await
.map_err(|err| format!("in-flight GET returned a failed or truncated stream: {err:?}"))?;
Ok::<_, String>(data)
});
tokio::time::timeout(TRANSITION_WAIT_TIMEOUT, get_barrier.wait_until_paused())
.await
.expect("the in-flight GET should reach the remote read barrier");
// The next remote DELETE pauses and then fails. A correct local-first
// expiry returns while this barrier is still held; synchronous cleanup
// (remote-first or local-first) instead times out here.
let delete_start_barrier = DeleteAfterObjectLockSnapshotBarrier::install(bucket_name.as_str());
let remove_barrier = backend.arm_failing_remove_barrier().await;
let expiry_store = ecstore.clone();
let expiry_oi = oi.clone();
let expiry_event = lc_event.clone();
let mut expiry = tokio::spawn(async move {
expire_transitioned_object(expiry_store, &expiry_oi, &expiry_event, &LcEventSrc::Scanner, bucket_incarnation_id).await
});
tokio::time::timeout(TRANSITION_WAIT_TIMEOUT, delete_start_barrier.wait_until_paused())
.await
.expect("expiry should reach the store delete path while the GET remains paused");
delete_start_barrier.release_and_wait_until_namespace_pending().await;
assert!(
!delete_start_barrier.namespace_acquired() && !expiry.is_finished(),
"expiry must wait for the in-flight GET's object read lock before committing the local delete"
);
get_barrier.release();
let expiry_outcome = tokio::time::timeout(TRANSITION_WAIT_TIMEOUT, &mut expiry).await;
let delete_lock_acquired_after_get = delete_start_barrier.namespace_acquired();
drop(delete_start_barrier);
let remove_arrival = tokio::time::timeout(TRANSITION_WAIT_TIMEOUT, remove_barrier.wait_until_paused()).await;
// Snapshot only lock-free observables while the cleanup worker holds
// the object write lock. Store API reads wait until the barrier is
// released below.
let free_version_persisted = free_version_count(&disk_paths[0], bucket_name.as_str(), object_name).await > 0;
let remote_present_at_cleanup = backend.contains(&remote_object).await;
remove_barrier.release();
let remove_operation_dropped = if remove_arrival.is_ok() {
Some(tokio::time::timeout(TRANSITION_WAIT_TIMEOUT, remove_barrier.wait_until_operation_dropped()).await)
} else {
None
};
if expiry_outcome.is_err() && tokio::time::timeout(TRANSITION_WAIT_TIMEOUT, &mut expiry).await.is_err() {
expiry.abort();
let _ = expiry.await;
}
let in_flight_get_outcome = tokio::time::timeout(TRANSITION_WAIT_TIMEOUT, in_flight_get).await;
let post_expiry_get = tokio::time::timeout(
TRANSITION_WAIT_TIMEOUT,
ecstore.get_object_reader(bucket_name.as_str(), object_name, None, http::HeaderMap::new(), &ObjectOptions::default()),
)
.await;
expiry_outcome
.expect("expire_transitioned_object must not wait for asynchronous remote-tier cleanup")
.expect("the expiry task should not panic")
.expect("expire_transitioned_object should succeed");
assert!(
delete_lock_acquired_after_get,
"expiry must acquire the object write lock only after the in-flight GET releases its read lock"
);
remove_arrival.expect("the post-commit free-version worker should reach the remote DELETE barrier");
remove_operation_dropped
.expect("the remote DELETE should have reached the barrier")
.expect("the injected remote DELETE should finish after release");
assert!(
free_version_persisted,
"the durable free-version marker must exist before asynchronous remote cleanup"
);
assert!(
remote_present_at_cleanup,
"the remote object must remain readable until the paused cleanup DELETE is released"
);
// --- Ordering contract (deterministic revert-proof) ----------------
// #3491 defers remote cleanup to free-version recovery, so immediately
// after expiry the remote object is still present and NO synchronous
// remote `remove` was issued. Reverting to remote-first ordering makes
// both assertions fail.
let in_flight_body = in_flight_get_outcome
.expect("the in-flight GET should finish within the test deadline")
.expect("the in-flight GET task should not panic")
.expect("a GET that wins the expiry race must return a complete body");
assert_eq!(
backend.remove_count().await,
0,
"expire_transitioned_object must NOT issue a synchronous remote-tier removal (local-first \
ordering, #3491); remote cleanup is deferred to free-version recovery"
);
assert!(
backend.contains(&remote_object).await,
"remote tier object must still exist immediately after expiry (deferred cleanup, #3491)"
in_flight_body, payload,
"a GET that resolved transitioned metadata before expiry must return the complete, correct body"
);
// Local metadata is gone: the object is atomically unreachable.
assert!(
wait_for_object_absence(&ecstore, bucket_name.as_str(), object_name, Duration::from_secs(5)).await,
"local metadata for the expired transitioned object should be gone"
);
// Drain the concurrent GET loop; its internal asserts already guarantee
// no #3491-style tier-fetch failure was ever observed.
let (saw_full_body, saw_not_found) = get_loop.await.expect("concurrent GET loop task panicked");
assert!(
saw_full_body + saw_not_found > 0,
"the concurrent GET loop should have observed at least one GET outcome"
);
match post_expiry_get.expect("the post-expiry GET should finish within the test deadline") {
Ok(_) => panic!("the locally expired transitioned object must no longer be readable"),
Err(err) => {
let ec: &EcstoreError = &err;
assert!(
is_err_object_not_found(ec) || is_err_version_not_found(ec),
"a GET after expiry may only fail with a clean object/version-not-found; \
a tier-fetch or read-quorum failure is the #3491 regression: {err:?}"
);
}
}
}
#[test]
@@ -1469,10 +1503,18 @@ mod serial_tests {
let stale_remote_object = transitioned.transitioned_object.name.clone();
assert!(backend.contains(&stale_remote_object).await);
ecstore
.delete_object(bucket_name.as_str(), object_name, ObjectOptions::default())
ExpiryState::resize_workers(1, ecstore.clone()).await;
let remove_barrier = backend.arm_failing_remove_barrier().await;
tokio::time::timeout(
Duration::from_secs(5),
ecstore.delete_object(bucket_name.as_str(), object_name, ObjectOptions::default()),
)
.await
.expect("DeleteObject must not wait for asynchronous remote-tier cleanup")
.expect("Failed to delete transitioned object before scanner fallback");
tokio::time::timeout(Duration::from_secs(5), remove_barrier.wait_until_paused())
.await
.expect("Failed to delete transitioned object without expiry workers");
.expect("the immediate free-version worker should reach the injected remote DELETE barrier");
assert!(
free_version_count(&disk_paths[0], bucket_name.as_str(), object_name).await > 0,
@@ -1483,8 +1525,12 @@ mod serial_tests {
"stale transitioned remote object should still exist before scanner fallback runs"
);
init_background_expiry(ecstore.clone()).await;
// Queue the scanner fallback while the causal task is still blocked.
// Releasing the barrier fails only that first task, so the queued
// scanner task can prove durable-marker recovery on a healthy backend.
scan_object_metadata(&disk_paths[0], bucket_name.as_str(), object_name).await;
remove_barrier.release();
remove_barrier.wait_until_operation_dropped().await;
assert!(
backend
@@ -1531,10 +1577,18 @@ mod serial_tests {
let stale_remote_object = transitioned.transitioned_object.name.clone();
assert!(backend.contains(&stale_remote_object).await);
ecstore
.delete_object(bucket_name.as_str(), object_name, ObjectOptions::default())
ExpiryState::resize_workers(1, ecstore.clone()).await;
let remove_barrier = backend.arm_failing_remove_barrier().await;
tokio::time::timeout(
Duration::from_secs(5),
ecstore.delete_object(bucket_name.as_str(), object_name, ObjectOptions::default()),
)
.await
.expect("DeleteObject must not wait for asynchronous remote-tier cleanup")
.expect("Failed to delete transitioned object after compensation-driven transition");
tokio::time::timeout(Duration::from_secs(5), remove_barrier.wait_until_paused())
.await
.expect("Failed to delete transitioned object after compensation-driven transition");
.expect("the immediate free-version worker should reach the injected remote DELETE barrier");
assert!(
free_version_count(&disk_paths[0], bucket_name.as_str(), object_name).await > 0,
@@ -1545,8 +1599,12 @@ mod serial_tests {
"stale transitioned remote object should still exist before scanner cleanup runs"
);
init_background_expiry(ecstore.clone()).await;
// Enqueue the scanner fallback before the first, causal cleanup task is
// released into its injected failure. This keeps attribution
// deterministic and proves the durable marker drives convergence.
scan_object_metadata(&disk_paths[0], bucket_name.as_str(), object_name).await;
remove_barrier.release();
remove_barrier.wait_until_operation_dropped().await;
assert!(
backend
+11 -8
View File
@@ -15,7 +15,9 @@
pub(crate) use rustfs_ecstore::api::bucket::lifecycle::transition_transaction::recover_transition_transaction_records;
pub(crate) use rustfs_ecstore::api::bucket::lifecycle::{
bucket_lifecycle_audit::LcEventSrc,
bucket_lifecycle_ops::{enqueue_transition_for_existing_objects, expire_transitioned_object, init_background_expiry},
bucket_lifecycle_ops::{
ExpiryState, enqueue_transition_for_existing_objects, expire_transitioned_object, init_background_expiry,
},
lifecycle::{Event as LcEvent, IlmAction, TRANSITION_PENDING, TransitionOptions},
};
pub(crate) use rustfs_ecstore::api::bucket::metadata::BUCKET_LIFECYCLE_CONFIG;
@@ -27,6 +29,7 @@ pub(crate) use rustfs_ecstore::api::capacity::path2_bucket_object_with_base_path
pub(crate) use rustfs_ecstore::api::disk::{DiskOption, STORAGE_FORMAT_FILE, endpoint::Endpoint, new_disk};
pub(crate) use rustfs_ecstore::api::error::{Error as EcstoreError, is_err_object_not_found, is_err_version_not_found};
pub(crate) use rustfs_ecstore::api::layout::{EndpointServerPools, Endpoints, PoolEndpoints};
pub(crate) use rustfs_ecstore::api::object::test_util::DeleteAfterObjectLockSnapshotBarrier;
pub(crate) use rustfs_ecstore::api::runtime::global_tier_config_mgr as get_global_tier_config_mgr;
pub(crate) use rustfs_ecstore::api::storage::{ECStore, init_local_disks};
// Shared lifecycle/tier test utilities (rustfs/backlog#1148 ilm-6). The mock
@@ -45,12 +48,12 @@ pub(crate) mod lifecycle {
};
pub(crate) use super::{
BUCKET_LIFECYCLE_CONFIG, BucketVersioningSys, DiskOption, ECStore, EcstoreError, Endpoint, EndpointServerPools,
Endpoints, IlmAction, LcEvent, LcEventSrc, MockWarmBackend, PoolEndpoints, STORAGE_FORMAT_FILE, TRANSITION_PENDING,
TransitionCleanupStoreBarrier, TransitionOptions, assert_transition_meta_consistent,
enqueue_transition_for_existing_objects, expire_transitioned_object, free_version_count, get_bucket_metadata,
get_global_tier_config_mgr, init_background_expiry, init_bucket_metadata_sys, init_local_disks, is_err_object_not_found,
is_err_version_not_found, new_disk, path2_bucket_object_with_base_path, recover_transition_transaction_records,
register_mock_tier_util, update_bucket_metadata, wait_for_free_version_absence,
BUCKET_LIFECYCLE_CONFIG, BucketVersioningSys, DeleteAfterObjectLockSnapshotBarrier, DiskOption, ECStore, EcstoreError,
Endpoint, EndpointServerPools, Endpoints, ExpiryState, IlmAction, LcEvent, LcEventSrc, MockWarmBackend, PoolEndpoints,
STORAGE_FORMAT_FILE, TRANSITION_PENDING, TransitionCleanupStoreBarrier, TransitionOptions,
assert_transition_meta_consistent, enqueue_transition_for_existing_objects, expire_transitioned_object,
free_version_count, get_bucket_metadata, get_global_tier_config_mgr, init_background_expiry, init_bucket_metadata_sys,
init_local_disks, is_err_object_not_found, is_err_version_not_found, new_disk, path2_bucket_object_with_base_path,
recover_transition_transaction_records, register_mock_tier_util, update_bucket_metadata, wait_for_free_version_absence,
};
}