fix(tier): preserve cleanup scheduling for overwrites (#7766)

Queue committed tier free-version cleanup receipts for PUT and materialized CopyObject overwrites of transitioned null versions, while keeping remote deletion behind the existing persisted free-version cleanup path.

Tighten data-movement delete-marker retry equivalence by ignoring local bucket-incarnation fencing metadata, avoid retry fallback to the source pool, and keep version-list pagination from manufacturing an empty final page.

Refresh the e2e-distributed selector hash for the current release test set.

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
Hauser
2026-09-14 01:31:12 +08:00
committed by GitHub
parent 568e8cbe0c
commit 8982b4a3d2
4 changed files with 127 additions and 26 deletions
+2 -2
View File
@@ -1,2 +1,2 @@
sha256-linux=563bff8f1171d6dbe166ff8440310dbe98430e466aa3ecd8dc39e3c872b320f7
sha256-darwin=563bff8f1171d6dbe166ff8440310dbe98430e466aa3ecd8dc39e3c872b320f7
sha256-linux=3b17d97e97dacff32414b8094d3086ed8a38d509ba88b0734b884166b30e3157
sha256-darwin=3b17d97e97dacff32414b8094d3086ed8a38d509ba88b0734b884166b30e3157
+82 -8
View File
@@ -383,7 +383,25 @@ fn record_committed_tier_free_version_receipt(
free_version_id: Uuid,
batch: bool,
) {
if let Some(sink) = opts.tier_free_version_receipt_sink.as_ref()
record_committed_tier_free_version_receipt_to_sink(
opts.tier_free_version_receipt_sink.as_ref(),
bucket,
object,
source,
free_version_id,
batch,
);
}
fn record_committed_tier_free_version_receipt_to_sink(
sink: Option<&crate::object_api::TierFreeVersionReceiptSink>,
bucket: &str,
object: &str,
source: &ObjectInfo,
free_version_id: Uuid,
batch: bool,
) {
if let Some(sink) = sink
&& let Err(err) = sink.record(source, free_version_id)
{
warn!(
@@ -3822,13 +3840,18 @@ impl SetDisks {
);
}
fi.metadata = user_defined;
if fi.version_id.is_none_or(|id| id.is_nil()) && !opts.data_movement && expected_restore_operation_id.is_none() {
// Every disk must publish the same cleanup owner alongside a
// replaced null version. This transient key is not persisted
// on the new object; recovery discovers the free-version in
// the committed xl.meta even if this request is cancelled.
fi.set_tier_free_version_id(&Uuid::new_v4().to_string());
}
let put_tier_free_version_id =
if fi.version_id.is_none_or(|id| id.is_nil()) && !opts.data_movement && expected_restore_operation_id.is_none() {
// Every disk must publish the same cleanup owner alongside a
// replaced null version. This transient key is not persisted
// on the new object; recovery discovers the free-version in
// the committed xl.meta even if this request is cancelled.
let free_version_id = Uuid::new_v4();
fi.set_tier_free_version_id(&free_version_id.to_string());
Some(free_version_id)
} else {
None
};
fi.mod_time = mod_time;
fi.size = w_size as i64;
fi.versioned = opts.versioned || opts.version_suspended;
@@ -4304,6 +4327,9 @@ impl SetDisks {
let commit_capacity_scope_token = opts.capacity_scope_token;
let commit_replication_state = replication_state_to_filemeta(&opts.put_replication_state());
let commit_scanner_publication_lease_tokens = scanner_publication_lease_tokens;
let commit_put_tier_free_version_id = put_tier_free_version_id;
let commit_tier_free_version_receipt_sink = opts.tier_free_version_receipt_sink.clone();
let commit_skip_free_version = opts.skip_free_version;
let request_cancellation = operation_cancellation.clone();
tmp_cleanup_owned = true;
@@ -4460,6 +4486,41 @@ impl SetDisks {
return Err(err);
}
let put_tier_free_version_source =
if commit_put_tier_free_version_id.is_some() && commit_tier_free_version_receipt_sink.is_some() {
match commit_set
.get_object_info(
&commit_bucket,
&commit_object,
&ObjectOptions {
no_lock: true,
metadata_cache_safe: false,
versioned: commit_versioned,
version_suspended: commit_version_suspended,
..Default::default()
},
)
.await
{
Ok(source) => Some(source),
Err(err) if is_err_object_not_found(&err) || is_err_version_not_found(&err) => None,
Err(err) => {
debug!(
event = EVENT_LIFECYCLE_TRANSITIONED_DELETE_CLEANUP_OWNER,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
bucket = %commit_bucket,
object = %commit_object,
error = ?err,
"Skipped opportunistic tier free-version receipt source capture"
);
None
}
}
} else {
None
};
Self::assign_rename_data_indexes(&mut parts_metadatas);
let mut rename_result = SetDisks::rename_data_owned_with_fence(
&commit_disks,
@@ -4632,6 +4693,19 @@ impl SetDisks {
let cleanup_disks = rename_commit.cleanup_disks;
let old_current_size = rename_commit.old_current_size;
let mut fi = rename_commit.committed_file_info;
if let (Some(source), Some(free_version_id)) =
(put_tier_free_version_source.as_ref(), commit_put_tier_free_version_id)
&& transitioned_delete_publishes_free_version(source, &fi, commit_skip_free_version)
{
record_committed_tier_free_version_receipt_to_sink(
commit_tier_free_version_receipt_sink.as_ref(),
&commit_bucket,
&commit_object,
source,
free_version_id,
false,
);
}
if needs_immediate_heal {
let mut request = rustfs_heal_contracts::heal_channel::create_heal_request_with_options(
+21 -8
View File
@@ -2216,7 +2216,7 @@ fn list_objects_paginate(
}
}
if !is_truncated && disk_has_more {
if !is_truncated && disk_has_more && !include_version_id {
let visible_count = objects.len() + prefixes.len();
let should_truncate = if delimiter.is_none() {
visible_count > 0
@@ -7312,14 +7312,23 @@ mod test {
fn test_object_meta_entry(name: &str) -> MetaCacheEntry {
let mut meta = FileMeta::new();
meta.add_version(FileInfo {
volume: "bucket".to_owned(),
name: name.to_owned(),
let mut metadata = HashMap::new();
metadata.insert("etag".to_string(), "etag".to_string());
let mut fi = FileInfo::new(name, 2, 2);
fi.erasure.index = 1;
fi.data_dir = Some(Uuid::from_u128(0x1234));
fi.volume = "bucket".to_owned();
fi.name = name.to_owned();
fi.size = 1;
fi.parts = vec![ObjectPartInfo {
number: 1,
size: 1,
mod_time: Some(time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp")),
actual_size: 1,
..Default::default()
})
.expect("test metadata should accept object version");
}];
fi.mod_time = Some(time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp"));
fi.metadata = metadata;
meta.add_version(fi).expect("test metadata should accept object version");
let metadata = meta.marshal_msg().expect("test metadata should marshal");
MetaCacheEntry {
@@ -7844,7 +7853,11 @@ mod test {
}
.expect("version page should list successfully");
let page_size = usize::try_from(max_keys).expect("nonnegative page size");
assert_eq!(result.objects.len() + result.prefixes.len(), (10 - page * page_size).min(page_size));
assert_eq!(
result.objects.len() + result.prefixes.len(),
(10 - page * page_size).min(page_size),
"{kind}, layer {layer}, max_keys {max_keys}, page {page}"
);
let has_more = page + 1 < expected_pages;
assert_eq!(result.is_truncated, has_more, "{kind}, layer {layer}, max_keys {max_keys}, page {page}");
assert_eq!(
+22 -8
View File
@@ -2056,7 +2056,7 @@ fn data_movement_delete_marker_metadata_identity(metadata: &HashMap<String, Stri
local_tier_free_version_id = Some(version_id);
continue;
}
if suffix.eq_ignore_ascii_case(rustfs_utils::http::SUFFIX_BUCKET_INCARNATION_ID) {
if suffix.eq_ignore_ascii_case(rustfs_utils::http::metadata_compat::SUFFIX_BUCKET_INCARNATION_ID) {
continue;
}
@@ -4270,9 +4270,12 @@ impl ECStore {
if !self.single_pool() {
opts.decommission_capacity_admission = crate::bucket::metadata_sys::object_store_if_initialized_in(&self.ctx).await;
}
self.pools[idx]
let receipt_sink = install_tier_free_version_receipt_sink(&mut opts);
let result = self.pools[idx]
.put_object_with_old_current_size(bucket, object.as_str(), data, &opts)
.await
.await;
enqueue_recorded_tier_free_versions(self, receipt_sink).await;
result
}
#[instrument(level = "trace", skip(self))]
@@ -4476,9 +4479,12 @@ impl ECStore {
crate::bucket::metadata_sys::object_store_if_initialized_in(&self.ctx).await;
}
return if let Some(reader) = src_info.put_object_reader.as_mut() {
self.pools[pool_idx]
let receipt_sink = install_tier_free_version_receipt_sink(&mut put_opts);
let result = self.pools[pool_idx]
.put_object(dst_bucket, &dst_object, reader, &put_opts)
.await
.await;
enqueue_recorded_tier_free_versions(self, receipt_sink).await;
result
} else {
Err(StorageError::InvalidArgument(
src_bucket.to_owned(),
@@ -4514,9 +4520,12 @@ impl ECStore {
put_opts.decommission_capacity_admission =
crate::bucket::metadata_sys::object_store_if_initialized_in(&self.ctx).await;
}
return self.pools[pool_idx]
let receipt_sink = install_tier_free_version_receipt_sink(&mut put_opts);
let result = self.pools[pool_idx]
.put_object(dst_bucket, &dst_object, reader, &put_opts)
.await;
enqueue_recorded_tier_free_versions(self, receipt_sink).await;
return result;
}
src_info.version_only = true;
let capacity_object = dst_object.clone();
@@ -4565,9 +4574,12 @@ impl ECStore {
}
if let Some(put_object_reader) = src_info.put_object_reader.as_mut() {
return self.pools[pool_idx]
let receipt_sink = install_tier_free_version_receipt_sink(&mut put_opts);
let result = self.pools[pool_idx]
.put_object(dst_bucket, dst_object_name, put_object_reader, &put_opts)
.await;
enqueue_recorded_tier_free_versions(self, receipt_sink).await;
return result;
}
Err(StorageError::InvalidArgument(
@@ -4819,7 +4831,9 @@ impl ECStore {
if let Some(owner) = DecommissionCapacityOwner::from_options(&opts) {
self.select_decommission_capacity_target_pool(owner, 0).await?
} else {
self.get_pool_idx_no_lock(bucket, object, 0).await?
self.get_available_pool_idx_excluding(bucket, object, 0, opts.src_pool_idx)
.await
.ok_or(Error::DiskFull)?
}
}
};