From 2ad4627329b0728924e9e1d89a046deaf6480b22 Mon Sep 17 00:00:00 2001 From: houseme Date: Fri, 24 Jul 2026 12:19:31 +0800 Subject: [PATCH] feat(ilm): allow manual transition resume markers Accept additive marker and versionMarker parameters on the bounded manual transition run API so clients can continue from a partial report without changing the existing default scan behavior. Co-Authored-By: heihutu --- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 39 +++++++++++++------ rustfs/src/admin/handlers/ilm_transition.rs | 10 ++++- 2 files changed, 36 insertions(+), 13 deletions(-) diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 276ceda3d..c9ff3f686 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -2444,6 +2444,8 @@ pub async fn enqueue_immediate_expiry(oi: &ObjectInfo, src: LcEventSrc) { #[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize)] pub struct ManualTransitionRunOptions { pub prefix: String, + pub marker: Option, + pub version_marker: Option, pub tier: Option, pub dry_run: bool, pub max_objects: Option, @@ -2507,8 +2509,8 @@ pub async fn enqueue_transition_for_existing_objects_scoped( return Ok(report); }; report.lifecycle_config_found = true; - let mut marker = None; - let mut version_marker = None; + let mut marker = options.marker.clone(); + let mut version_marker = options.version_marker.clone(); let src = LcEventSrc::Scanner; loop { @@ -2517,13 +2519,15 @@ pub async fn enqueue_transition_for_existing_objects_scoped( .list_object_versions(bucket, &options.prefix, marker.clone(), version_marker.clone(), None, LIST_PAGE_SIZE) .await?; - for object in &page.objects { + for (index, object) in page.objects.iter().enumerate() { report.scanned = report.scanned.saturating_add(1); enqueue_transition_with_lifecycle_report(object, &lc, &src, &options, &mut report).await; if options.max_objects.is_some_and(|max_objects| report.scanned >= max_objects) { - report.truncated_by_limit = true; - report.next_marker = Some(object.name.clone()); - report.next_version_idmarker = object.version_id.map(|version| version.to_string()); + if manual_transition_has_more_after_limit(index, page.objects.len(), page.is_truncated) { + report.truncated_by_limit = true; + report.next_marker = Some(object.name.clone()); + report.next_version_idmarker = object.version_id.map(|version| version.to_string()); + } return Ok(report); } } @@ -2723,6 +2727,10 @@ pub async fn enqueue_expiry_for_existing_objects(api: Arc, bucket: &str } } +fn manual_transition_has_more_after_limit(page_index: usize, page_len: usize, page_is_truncated: bool) -> bool { + page_index.saturating_add(1) < page_len || page_is_truncated +} + async fn enqueue_transition_with_lifecycle_report( oi: &ObjectInfo, lc: &BucketLifecycleConfiguration, @@ -3622,12 +3630,12 @@ mod tests { eval_action_from_lifecycle, jitter_tier_free_version_recovery_delay, lifecycle_action_blocked_by_replication, lifecycle_delete_all_versions_replication_scan, lifecycle_deleted_object, lifecycle_replication_blocks_action, lifecycle_rule_has_date_expiration, lifecycle_version_purge_state_from_completed_targets, - mark_delete_opts_skip_decommissioned_on_remote_success, merge_stale_multipart_candidate, replication_state_for_delete, - resolve_transition_queue_capacity, resolve_transition_queue_send_timeout, resolve_transition_worker_count, - resolve_transition_workers_absolute_max, run_tier_free_version_recovery_loop, select_restore_s3_location, - set_recovered_free_version_enqueue_observer, should_defer_date_expiry_for_recent_config_update, - should_reuse_lifecycle_delete_replication_state, transitioned_cleanup_tuple, transitioned_object_delete_opts, - wait_for_tier_free_version_recovery, + manual_transition_has_more_after_limit, mark_delete_opts_skip_decommissioned_on_remote_success, + merge_stale_multipart_candidate, replication_state_for_delete, resolve_transition_queue_capacity, + resolve_transition_queue_send_timeout, resolve_transition_worker_count, resolve_transition_workers_absolute_max, + run_tier_free_version_recovery_loop, select_restore_s3_location, set_recovered_free_version_enqueue_observer, + should_defer_date_expiry_for_recent_config_update, should_reuse_lifecycle_delete_replication_state, + transitioned_cleanup_tuple, transitioned_object_delete_opts, wait_for_tier_free_version_recovery, }; #[cfg(feature = "test-util")] use super::{delete_free_version_remote_object_then, encode_dir_object, get_transitioned_object_reader_with_tier_manager}; @@ -6269,6 +6277,13 @@ mod tests { assert_eq!(report.skipped_tier, 1); } + #[test] + fn manual_transition_limit_boundary_only_truncates_when_more_objects_exist() { + assert!(!manual_transition_has_more_after_limit(9, 10, false)); + assert!(manual_transition_has_more_after_limit(9, 11, false)); + assert!(manual_transition_has_more_after_limit(9, 10, true)); + } + #[tokio::test] async fn existing_object_lifecycle_allows_expired_marker_after_replication_completed() { let lc = expired_delete_marker_lifecycle(); diff --git a/rustfs/src/admin/handlers/ilm_transition.rs b/rustfs/src/admin/handlers/ilm_transition.rs index 9603d2248..beb7d3ff6 100644 --- a/rustfs/src/admin/handlers/ilm_transition.rs +++ b/rustfs/src/admin/handlers/ilm_transition.rs @@ -38,6 +38,9 @@ const MAX_MANUAL_TRANSITION_OBJECTS: u64 = 100_000; struct ManualTransitionRunQuery { bucket: Option, prefix: Option, + marker: Option, + #[serde(rename = "versionMarker")] + version_marker: Option, tier: Option, #[serde(rename = "dryRun")] dry_run: Option, @@ -90,6 +93,8 @@ fn parse_manual_transition_query(query: Option<&str>) -> S3Result<(String, Manua bucket.to_string(), ManualTransitionRunOptions { prefix: query.prefix.unwrap_or_default(), + marker: query.marker.filter(|marker| !marker.is_empty()), + version_marker: query.version_marker.filter(|version_marker| !version_marker.is_empty()), tier: query.tier.map(|tier| tier.trim().to_string()).filter(|tier| !tier.is_empty()), dry_run: query.dry_run.unwrap_or(false), max_objects: Some(max_objects), @@ -172,10 +177,13 @@ mod tests { #[test] fn manual_transition_query_defaults_to_bounded_run() { let (bucket, options) = - parse_manual_transition_query(Some("bucket=data&prefix=logs/&tier=warm")).expect("valid query should parse"); + parse_manual_transition_query(Some("bucket=data&prefix=logs/&marker=logs/a&versionMarker=v1&tier=warm")) + .expect("valid query should parse"); assert_eq!(bucket, "data"); assert_eq!(options.prefix, "logs/"); + assert_eq!(options.marker.as_deref(), Some("logs/a")); + assert_eq!(options.version_marker.as_deref(), Some("v1")); assert_eq!(options.tier.as_deref(), Some("warm")); assert!(!options.dry_run); assert_eq!(options.max_objects, Some(DEFAULT_MANUAL_TRANSITION_MAX_OBJECTS));