mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-29 16:37:07 +00:00
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 <heihutu@gmail.com>
This commit is contained in:
@@ -2444,6 +2444,8 @@ pub async fn enqueue_immediate_expiry(oi: &ObjectInfo, src: LcEventSrc) {
|
|||||||
#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize)]
|
#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize)]
|
||||||
pub struct ManualTransitionRunOptions {
|
pub struct ManualTransitionRunOptions {
|
||||||
pub prefix: String,
|
pub prefix: String,
|
||||||
|
pub marker: Option<String>,
|
||||||
|
pub version_marker: Option<String>,
|
||||||
pub tier: Option<String>,
|
pub tier: Option<String>,
|
||||||
pub dry_run: bool,
|
pub dry_run: bool,
|
||||||
pub max_objects: Option<u64>,
|
pub max_objects: Option<u64>,
|
||||||
@@ -2507,8 +2509,8 @@ pub async fn enqueue_transition_for_existing_objects_scoped(
|
|||||||
return Ok(report);
|
return Ok(report);
|
||||||
};
|
};
|
||||||
report.lifecycle_config_found = true;
|
report.lifecycle_config_found = true;
|
||||||
let mut marker = None;
|
let mut marker = options.marker.clone();
|
||||||
let mut version_marker = None;
|
let mut version_marker = options.version_marker.clone();
|
||||||
let src = LcEventSrc::Scanner;
|
let src = LcEventSrc::Scanner;
|
||||||
|
|
||||||
loop {
|
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)
|
.list_object_versions(bucket, &options.prefix, marker.clone(), version_marker.clone(), None, LIST_PAGE_SIZE)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
for object in &page.objects {
|
for (index, object) in page.objects.iter().enumerate() {
|
||||||
report.scanned = report.scanned.saturating_add(1);
|
report.scanned = report.scanned.saturating_add(1);
|
||||||
enqueue_transition_with_lifecycle_report(object, &lc, &src, &options, &mut report).await;
|
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) {
|
if options.max_objects.is_some_and(|max_objects| report.scanned >= max_objects) {
|
||||||
report.truncated_by_limit = true;
|
if manual_transition_has_more_after_limit(index, page.objects.len(), page.is_truncated) {
|
||||||
report.next_marker = Some(object.name.clone());
|
report.truncated_by_limit = true;
|
||||||
report.next_version_idmarker = object.version_id.map(|version| version.to_string());
|
report.next_marker = Some(object.name.clone());
|
||||||
|
report.next_version_idmarker = object.version_id.map(|version| version.to_string());
|
||||||
|
}
|
||||||
return Ok(report);
|
return Ok(report);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -2723,6 +2727,10 @@ pub async fn enqueue_expiry_for_existing_objects(api: Arc<ECStore>, 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(
|
async fn enqueue_transition_with_lifecycle_report(
|
||||||
oi: &ObjectInfo,
|
oi: &ObjectInfo,
|
||||||
lc: &BucketLifecycleConfiguration,
|
lc: &BucketLifecycleConfiguration,
|
||||||
@@ -3622,12 +3630,12 @@ mod tests {
|
|||||||
eval_action_from_lifecycle, jitter_tier_free_version_recovery_delay, lifecycle_action_blocked_by_replication,
|
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_delete_all_versions_replication_scan, lifecycle_deleted_object, lifecycle_replication_blocks_action,
|
||||||
lifecycle_rule_has_date_expiration, lifecycle_version_purge_state_from_completed_targets,
|
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,
|
manual_transition_has_more_after_limit, mark_delete_opts_skip_decommissioned_on_remote_success,
|
||||||
resolve_transition_queue_capacity, resolve_transition_queue_send_timeout, resolve_transition_worker_count,
|
merge_stale_multipart_candidate, replication_state_for_delete, resolve_transition_queue_capacity,
|
||||||
resolve_transition_workers_absolute_max, run_tier_free_version_recovery_loop, select_restore_s3_location,
|
resolve_transition_queue_send_timeout, resolve_transition_worker_count, resolve_transition_workers_absolute_max,
|
||||||
set_recovered_free_version_enqueue_observer, should_defer_date_expiry_for_recent_config_update,
|
run_tier_free_version_recovery_loop, select_restore_s3_location, set_recovered_free_version_enqueue_observer,
|
||||||
should_reuse_lifecycle_delete_replication_state, transitioned_cleanup_tuple, transitioned_object_delete_opts,
|
should_defer_date_expiry_for_recent_config_update, should_reuse_lifecycle_delete_replication_state,
|
||||||
wait_for_tier_free_version_recovery,
|
transitioned_cleanup_tuple, transitioned_object_delete_opts, wait_for_tier_free_version_recovery,
|
||||||
};
|
};
|
||||||
#[cfg(feature = "test-util")]
|
#[cfg(feature = "test-util")]
|
||||||
use super::{delete_free_version_remote_object_then, encode_dir_object, get_transitioned_object_reader_with_tier_manager};
|
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);
|
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]
|
#[tokio::test]
|
||||||
async fn existing_object_lifecycle_allows_expired_marker_after_replication_completed() {
|
async fn existing_object_lifecycle_allows_expired_marker_after_replication_completed() {
|
||||||
let lc = expired_delete_marker_lifecycle();
|
let lc = expired_delete_marker_lifecycle();
|
||||||
|
|||||||
@@ -38,6 +38,9 @@ const MAX_MANUAL_TRANSITION_OBJECTS: u64 = 100_000;
|
|||||||
struct ManualTransitionRunQuery {
|
struct ManualTransitionRunQuery {
|
||||||
bucket: Option<String>,
|
bucket: Option<String>,
|
||||||
prefix: Option<String>,
|
prefix: Option<String>,
|
||||||
|
marker: Option<String>,
|
||||||
|
#[serde(rename = "versionMarker")]
|
||||||
|
version_marker: Option<String>,
|
||||||
tier: Option<String>,
|
tier: Option<String>,
|
||||||
#[serde(rename = "dryRun")]
|
#[serde(rename = "dryRun")]
|
||||||
dry_run: Option<bool>,
|
dry_run: Option<bool>,
|
||||||
@@ -90,6 +93,8 @@ fn parse_manual_transition_query(query: Option<&str>) -> S3Result<(String, Manua
|
|||||||
bucket.to_string(),
|
bucket.to_string(),
|
||||||
ManualTransitionRunOptions {
|
ManualTransitionRunOptions {
|
||||||
prefix: query.prefix.unwrap_or_default(),
|
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()),
|
tier: query.tier.map(|tier| tier.trim().to_string()).filter(|tier| !tier.is_empty()),
|
||||||
dry_run: query.dry_run.unwrap_or(false),
|
dry_run: query.dry_run.unwrap_or(false),
|
||||||
max_objects: Some(max_objects),
|
max_objects: Some(max_objects),
|
||||||
@@ -172,10 +177,13 @@ mod tests {
|
|||||||
#[test]
|
#[test]
|
||||||
fn manual_transition_query_defaults_to_bounded_run() {
|
fn manual_transition_query_defaults_to_bounded_run() {
|
||||||
let (bucket, options) =
|
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!(bucket, "data");
|
||||||
assert_eq!(options.prefix, "logs/");
|
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_eq!(options.tier.as_deref(), Some("warm"));
|
||||||
assert!(!options.dry_run);
|
assert!(!options.dry_run);
|
||||||
assert_eq!(options.max_objects, Some(DEFAULT_MANUAL_TRANSITION_MAX_OBJECTS));
|
assert_eq!(options.max_objects, Some(DEFAULT_MANUAL_TRANSITION_MAX_OBJECTS));
|
||||||
|
|||||||
Reference in New Issue
Block a user