fix(ilm): preserve lifecycle version groups (#5239)

* fix(ilm): evaluate complete version groups

Co-Authored-By: heihutu <heihutu@gmail.com>

* feat(ilm): log replication expiry blocks

Co-Authored-By: heihutu <heihutu@gmail.com>

* fix(ilm): diagnose incomplete noncurrent chains

Co-Authored-By: heihutu <heihutu@gmail.com>

* test(ilm): cover purge-pending version groups

Add regression coverage for lifecycle-only version listing so purge-pending versions stay present for ILM evaluation while the public ListObjectVersions projection remains filtered.

Refs rustfs/backlog#1500

Co-Authored-By: heihutu <heihutu@gmail.com>

* feat(ilm): log lifecycle evaluation failures

Emit structured lifecycle_evaluation_failed events when version group loading or evaluation fails during immediate and existing-object expiry scans.

Refs rustfs/backlog#1503

Co-Authored-By: heihutu <heihutu@gmail.com>

* test(ilm): cover incomplete noncurrent chains

Pin the fail-closed behavior for noncurrent expiration when successor_mod_time is missing and reuse a stable structured event name for that skip path.

Refs rustfs/backlog#1502

Co-Authored-By: heihutu <heihutu@gmail.com>

* test(ilm): cover one-day noncurrent expiry boundary

Add a runtime regression test proving NoncurrentVersionExpiration Days=1 stays inactive before expected_expiry_time and deletes exactly at the computed due boundary.

Refs rustfs/backlog#1504

Co-Authored-By: heihutu <heihutu@gmail.com>

* fix(ilm): keep transition checks after incomplete expiry

Co-Authored-By: heihutu <heihutu@gmail.com>

---------

Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
houseme
2026-07-26 01:42:02 +08:00
committed by GitHub
parent 516f7fecc1
commit 23a5db4012
5 changed files with 374 additions and 29 deletions
@@ -112,6 +112,7 @@ const EVENT_LIFECYCLE_WORKER_STATE: &str = "lifecycle_worker_state";
const EVENT_LIFECYCLE_TRANSITION_COMPENSATION: &str = "lifecycle_transition_compensation";
const EVENT_LIFECYCLE_STALE_MULTIPART_CLEANUP: &str = "lifecycle_stale_multipart_cleanup";
const EVENT_LIFECYCLE_SCAN_SKIPPED: &str = "lifecycle_scan_skipped";
const EVENT_LIFECYCLE_EVALUATION_FAILED: &str = "lifecycle_evaluation_failed";
const EVENT_LIFECYCLE_TIER_AUDIT: &str = "lifecycle_tier_audit";
const EVENT_LIFECYCLE_TIER_OPERATION_FAILED: &str = "lifecycle_tier_operation_failed";
const EVENT_LIFECYCLE_DELETE_FAILED: &str = "lifecycle_delete_failed";
@@ -2647,12 +2648,25 @@ pub async fn enqueue_immediate_expiry(oi: &ObjectInfo, src: LcEventSrc) {
let mut object_infos = Vec::new();
loop {
let Ok(page) = api
let page = match api
.clone()
.list_object_versions(&oi.bucket, &oi.name, marker.clone(), version_marker.clone(), None, 1000)
.list_object_versions_for_lifecycle(&oi.bucket, &oi.name, marker.clone(), version_marker.clone(), None, 1000)
.await
else {
return;
{
Ok(page) => page,
Err(err) => {
warn!(
event = EVENT_LIFECYCLE_EVALUATION_FAILED,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
bucket = %oi.bucket,
object = %oi.name,
error = %err,
reason = "list_versions_failed",
"Failed to load lifecycle version group"
);
return;
}
};
object_infos.extend(page.objects.into_iter().filter(|object| object.name == oi.name));
@@ -2677,12 +2691,27 @@ pub async fn enqueue_immediate_expiry(oi: &ObjectInfo, src: LcEventSrc) {
.iter()
.map(lifecycle::object_opts_from_object_info)
.collect::<Vec<ObjectOpts>>();
let Ok(events) = Evaluator::new(Arc::new(lifecycle))
let events = match Evaluator::new(Arc::new(lifecycle))
.with_lock_retention(lock_config)
.eval(&object_opts)
.await
else {
return;
{
Ok(events) => events,
Err(err) => {
warn!(
event = EVENT_LIFECYCLE_EVALUATION_FAILED,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
bucket = %oi.bucket,
object = %oi.name,
expected_version_count = oi.num_versions,
observed_version_count = object_infos.len(),
error = %err,
reason = "version_group_evaluation_failed",
"Failed to evaluate lifecycle version group"
);
return;
}
};
let mut to_delete_objs = Vec::new();
@@ -3099,10 +3128,16 @@ async fn enqueue_expiry_for_existing_object_group(
Ok(events) => events,
Err(err) => {
warn!(
event = EVENT_LIFECYCLE_EVALUATION_FAILED,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
bucket = context.bucket,
object = %object_infos[0].name,
expected_version_count = object_infos[0].num_versions,
observed_version_count = object_infos.len(),
error = %err,
"failed to evaluate lifecycle events for existing object versions"
reason = "version_group_evaluation_failed",
"Failed to evaluate lifecycle version group"
);
return;
}
@@ -3210,7 +3245,7 @@ pub async fn enqueue_expiry_for_existing_objects(api: Arc<ECStore>, bucket: &str
loop {
let page = api
.clone()
.list_object_versions(bucket, "", marker.clone(), version_marker.clone(), None, 1000)
.list_object_versions_for_lifecycle(bucket, "", marker.clone(), version_marker.clone(), None, 1000)
.await?;
for object in page.objects {
@@ -3835,6 +3870,25 @@ pub async fn eval_action_from_lifecycle(
}
if lifecycle_action_blocked_by_replication(event.action, oi) {
let reason = if oi.version_purge_status.is_pending() {
"version_purge_pending"
} else if oi.replication_status == ReplicationStatusType::Failed {
"replication_failed"
} else {
"replication_pending"
};
debug!(
event = EVENT_LIFECYCLE_SCAN_SKIPPED,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
object = %oi.name,
version_id = ?oi.version_id,
action = ?event.action,
replication_status = ?oi.replication_status,
version_purge_status = ?oi.version_purge_status,
reason,
"Skipped lifecycle action because replication is not terminal"
);
return lifecycle::Event::default();
}
+85 -1
View File
@@ -556,6 +556,29 @@ impl ObjectInfo {
prefix: &str,
delimiter: Option<String>,
after_version_marker: Option<VersionMarker>,
) -> Vec<ObjectInfo> {
Self::from_meta_cache_entries_sorted_versions_with_purge(entries, bucket, prefix, delimiter, after_version_marker, false)
.await
}
pub(crate) async fn from_meta_cache_entries_sorted_versions_for_lifecycle(
entries: &MetaCacheEntriesSorted,
bucket: &str,
prefix: &str,
delimiter: Option<String>,
after_version_marker: Option<VersionMarker>,
) -> Vec<ObjectInfo> {
Self::from_meta_cache_entries_sorted_versions_with_purge(entries, bucket, prefix, delimiter, after_version_marker, true)
.await
}
async fn from_meta_cache_entries_sorted_versions_with_purge(
entries: &MetaCacheEntriesSorted,
bucket: &str,
prefix: &str,
delimiter: Option<String>,
after_version_marker: Option<VersionMarker>,
include_version_purge: bool,
) -> Vec<ObjectInfo> {
let vcfg = get_versioning_config(bucket).await.ok();
let mut objects = Vec::with_capacity(entries.entries().len());
@@ -604,7 +627,7 @@ impl ObjectInfo {
};
for fi in versions.iter() {
if !fi.version_purge_status().is_empty() {
if !include_version_purge && !fi.version_purge_status().is_empty() {
continue;
}
@@ -1064,6 +1087,67 @@ mod tests {
assert_eq!(objects[0].num_versions, 1);
}
#[tokio::test]
async fn lifecycle_versions_listing_preserves_purge_pending_versions() {
let visible_version_id = Uuid::new_v4();
let purge_version_id = Uuid::new_v4();
let base_time = OffsetDateTime::now_utc();
let mut fm = FileMeta::new();
fm.add_version(FileInfo {
volume: "bucket".to_string(),
name: "object".to_string(),
version_id: Some(purge_version_id),
mod_time: Some(base_time),
..Default::default()
})
.expect("version pending purge should be added");
fm.add_version(FileInfo {
volume: "bucket".to_string(),
name: "object".to_string(),
version_id: Some(visible_version_id),
mod_time: Some(base_time + time::Duration::seconds(1)),
..Default::default()
})
.expect("visible version should be added");
fm.delete_version(&FileInfo {
volume: "bucket".to_string(),
name: "object".to_string(),
version_id: Some(purge_version_id),
replication_state_internal: Some(crate::bucket::replication::replication_state_to_filemeta(&ReplicationState {
version_purge_status_internal: Some("arn:target-a=PENDING;".to_string()),
purge_targets: version_purge_statuses_map("arn:target-a=PENDING;"),
..Default::default()
})),
..Default::default()
})
.expect("version purge status should be persisted");
let entries = MetaCacheEntriesSorted {
o: rustfs_filemeta::MetaCacheEntries(vec![Some(MetaCacheEntry {
name: "object".to_string(),
metadata: fm.marshal_msg().expect("metadata should marshal"),
..Default::default()
})]),
..Default::default()
};
let public_objects = ObjectInfo::from_meta_cache_entries_sorted_versions(&entries, "bucket", "", None, None).await;
let lifecycle_objects =
ObjectInfo::from_meta_cache_entries_sorted_versions_for_lifecycle(&entries, "bucket", "", None, None).await;
assert_eq!(public_objects.len(), 1);
assert_eq!(public_objects[0].version_id, Some(visible_version_id));
assert_eq!(public_objects[0].num_versions, 2);
assert_eq!(lifecycle_objects.len(), 2);
assert!(
lifecycle_objects
.iter()
.any(|object| object.version_purge_status == VersionPurgeStatusType::Pending)
);
assert!(lifecycle_objects.iter().all(|object| object.num_versions == 2));
}
#[test]
fn get_actual_size_prefers_actual_size_field() {
let info = ObjectInfo {
+13
View File
@@ -55,6 +55,19 @@ impl ECStore {
.await
}
pub(crate) async fn list_object_versions_for_lifecycle(
self: Arc<Self>,
bucket: &str,
prefix: &str,
marker: Option<String>,
version_marker: Option<String>,
delimiter: Option<String>,
max_keys: i32,
) -> Result<ListObjectVersionsInfo> {
self.inner_list_object_versions_for_lifecycle(bucket, prefix, marker, version_marker, delimiter, max_keys)
.await
}
pub(super) async fn handle_walk(
self: Arc<Self>,
rx: CancellationToken,
+70 -8
View File
@@ -81,6 +81,17 @@ type ListObjectsV2Info = StorageListObjectsV2Info<ObjectInfo>;
type ListObjectVersionsInfo = StorageListObjectVersionsInfo<ObjectInfo>;
type ObjectInfoOrErr = StorageObjectInfoOrErr<ObjectInfo, Error>;
type WalkOptions = StorageWalkOptions<fn(&rustfs_filemeta::FileInfo) -> bool>;
struct ListObjectVersionsInput<'a> {
bucket: &'a str,
prefix: &'a str,
marker: Option<String>,
version_marker: Option<String>,
delimiter: Option<String>,
max_keys: i32,
include_version_purge: bool,
}
const LIST_MERGED_INPUT_BUFFER: usize = 1;
fn list_merged_entry_channel() -> (Sender<MetaCacheEntry>, Receiver<MetaCacheEntry>) {
@@ -3888,6 +3899,52 @@ impl ECStore {
delimiter: Option<String>,
max_keys: i32,
) -> Result<ListObjectVersionsInfo> {
self.inner_list_object_versions_with_projection(ListObjectVersionsInput {
bucket,
prefix,
marker,
version_marker,
delimiter,
max_keys,
include_version_purge: false,
})
.await
}
pub(crate) async fn inner_list_object_versions_for_lifecycle(
self: Arc<Self>,
bucket: &str,
prefix: &str,
marker: Option<String>,
version_marker: Option<String>,
delimiter: Option<String>,
max_keys: i32,
) -> Result<ListObjectVersionsInfo> {
self.inner_list_object_versions_with_projection(ListObjectVersionsInput {
bucket,
prefix,
marker,
version_marker,
delimiter,
max_keys,
include_version_purge: true,
})
.await
}
async fn inner_list_object_versions_with_projection(
self: Arc<Self>,
input: ListObjectVersionsInput<'_>,
) -> Result<ListObjectVersionsInfo> {
let ListObjectVersionsInput {
bucket,
prefix,
marker,
version_marker,
delimiter,
max_keys,
include_version_purge,
} = input;
let max_keys = normalize_max_keys(max_keys);
if marker.is_none() && version_marker.is_some() {
return Err(StorageError::NotImplemented);
@@ -3947,14 +4004,19 @@ impl ECStore {
// Last RAW scanned key, captured before folding (ECA-03 / #944).
let last_scanned_key = last_scanned_entry_name(list_result.entries.as_ref());
let get_objects = ObjectInfo::from_meta_cache_entries_sorted_versions(
&list_result.entries.unwrap_or_default(),
bucket,
prefix,
delimiter.clone(),
version_marker,
)
.await;
let entries = list_result.entries.unwrap_or_default();
let get_objects = if include_version_purge {
ObjectInfo::from_meta_cache_entries_sorted_versions_for_lifecycle(
&entries,
bucket,
prefix,
delimiter.clone(),
version_marker,
)
.await
} else {
ObjectInfo::from_meta_cache_entries_sorted_versions(&entries, bucket, prefix, delimiter.clone(), version_marker).await
};
let (objects, prefixes, is_truncated, next_marker, next_version_idmarker) = list_objects_paginate(
get_objects,