mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-24 13:16:28 +00:00
fix(ilm): harden manual transition partial reports
Preserve the null-version cursor contract, stop manual scans on enqueue pressure without skipping the failed object, and keep raw resume markers out of admin JSON responses. Co-Authored-By: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -2472,7 +2472,9 @@ pub struct ManualTransitionRunReport {
|
|||||||
pub skipped_queue_closed: u64,
|
pub skipped_queue_closed: u64,
|
||||||
pub skipped_queue_timeout: u64,
|
pub skipped_queue_timeout: u64,
|
||||||
pub truncated_by_limit: bool,
|
pub truncated_by_limit: bool,
|
||||||
|
#[serde(skip_serializing)]
|
||||||
pub next_marker: Option<String>,
|
pub next_marker: Option<String>,
|
||||||
|
#[serde(skip_serializing)]
|
||||||
pub next_version_idmarker: Option<String>,
|
pub next_version_idmarker: Option<String>,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2511,6 +2513,8 @@ pub async fn enqueue_transition_for_existing_objects_scoped(
|
|||||||
report.lifecycle_config_found = true;
|
report.lifecycle_config_found = true;
|
||||||
let mut marker = options.marker.clone();
|
let mut marker = options.marker.clone();
|
||||||
let mut version_marker = options.version_marker.clone();
|
let mut version_marker = options.version_marker.clone();
|
||||||
|
let mut previous_marker = marker.clone();
|
||||||
|
let mut previous_version_marker = version_marker.clone();
|
||||||
let src = LcEventSrc::Scanner;
|
let src = LcEventSrc::Scanner;
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
@@ -2522,14 +2526,21 @@ pub async fn enqueue_transition_for_existing_objects_scoped(
|
|||||||
for (index, object) in page.objects.iter().enumerate() {
|
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 report.has_partial_enqueue() {
|
||||||
|
report.next_marker.clone_from(&previous_marker);
|
||||||
|
report.next_version_idmarker.clone_from(&previous_version_marker);
|
||||||
|
return Ok(report);
|
||||||
|
}
|
||||||
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) {
|
||||||
if manual_transition_has_more_after_limit(index, page.objects.len(), page.is_truncated) {
|
if manual_transition_has_more_after_limit(index, page.objects.len(), page.is_truncated) {
|
||||||
report.truncated_by_limit = true;
|
report.truncated_by_limit = true;
|
||||||
report.next_marker = Some(object.name.clone());
|
report.next_marker = Some(object.name.clone());
|
||||||
report.next_version_idmarker = object.version_id.map(|version| version.to_string());
|
report.next_version_idmarker = Some(manual_transition_version_marker(object));
|
||||||
}
|
}
|
||||||
return Ok(report);
|
return Ok(report);
|
||||||
}
|
}
|
||||||
|
previous_marker = Some(object.name.clone());
|
||||||
|
previous_version_marker = Some(manual_transition_version_marker(object));
|
||||||
}
|
}
|
||||||
|
|
||||||
if !page.is_truncated {
|
if !page.is_truncated {
|
||||||
@@ -2538,6 +2549,8 @@ pub async fn enqueue_transition_for_existing_objects_scoped(
|
|||||||
|
|
||||||
marker = page.next_marker;
|
marker = page.next_marker;
|
||||||
version_marker = page.next_version_idmarker;
|
version_marker = page.next_version_idmarker;
|
||||||
|
previous_marker = marker.clone();
|
||||||
|
previous_version_marker = version_marker.clone();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2731,6 +2744,12 @@ fn manual_transition_has_more_after_limit(page_index: usize, page_len: usize, pa
|
|||||||
page_index.saturating_add(1) < page_len || page_is_truncated
|
page_index.saturating_add(1) < page_len || page_is_truncated
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn manual_transition_version_marker(oi: &ObjectInfo) -> String {
|
||||||
|
oi.version_id
|
||||||
|
.map(|version| version.to_string())
|
||||||
|
.unwrap_or_else(|| "null".to_string())
|
||||||
|
}
|
||||||
|
|
||||||
async fn enqueue_transition_with_lifecycle_report(
|
async fn enqueue_transition_with_lifecycle_report(
|
||||||
oi: &ObjectInfo,
|
oi: &ObjectInfo,
|
||||||
lc: &BucketLifecycleConfiguration,
|
lc: &BucketLifecycleConfiguration,
|
||||||
@@ -2793,7 +2812,7 @@ async fn enqueue_transition_with_lifecycle_report(
|
|||||||
|
|
||||||
async fn enqueue_transition_with_lifecycle(oi: &ObjectInfo, lc: &BucketLifecycleConfiguration, src: &LcEventSrc) -> bool {
|
async fn enqueue_transition_with_lifecycle(oi: &ObjectInfo, lc: &BucketLifecycleConfiguration, src: &LcEventSrc) -> bool {
|
||||||
let options = ManualTransitionRunOptions::default();
|
let options = ManualTransitionRunOptions::default();
|
||||||
let mut report = ManualTransitionRunReport::new(&oi.bucket, &options);
|
let mut report = ManualTransitionRunReport::default();
|
||||||
enqueue_transition_with_lifecycle_report(oi, lc, src, &options, &mut report).await
|
enqueue_transition_with_lifecycle_report(oi, lc, src, &options, &mut report).await
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -3630,12 +3649,13 @@ 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,
|
||||||
manual_transition_has_more_after_limit, mark_delete_opts_skip_decommissioned_on_remote_success,
|
manual_transition_has_more_after_limit, manual_transition_version_marker,
|
||||||
merge_stale_multipart_candidate, replication_state_for_delete, resolve_transition_queue_capacity,
|
mark_delete_opts_skip_decommissioned_on_remote_success, merge_stale_multipart_candidate, replication_state_for_delete,
|
||||||
resolve_transition_queue_send_timeout, resolve_transition_worker_count, resolve_transition_workers_absolute_max,
|
resolve_transition_queue_capacity, resolve_transition_queue_send_timeout, resolve_transition_worker_count,
|
||||||
run_tier_free_version_recovery_loop, select_restore_s3_location, set_recovered_free_version_enqueue_observer,
|
resolve_transition_workers_absolute_max, run_tier_free_version_recovery_loop, select_restore_s3_location,
|
||||||
should_defer_date_expiry_for_recent_config_update, should_reuse_lifecycle_delete_replication_state,
|
set_recovered_free_version_enqueue_observer, should_defer_date_expiry_for_recent_config_update,
|
||||||
transitioned_cleanup_tuple, transitioned_object_delete_opts, wait_for_tier_free_version_recovery,
|
should_reuse_lifecycle_delete_replication_state, 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};
|
||||||
@@ -5784,6 +5804,37 @@ mod tests {
|
|||||||
assert_eq!(state.transition_rx.len(), 1);
|
assert_eq!(state.transition_rx.len(), 1);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial]
|
||||||
|
async fn queue_transition_task_outcome_reports_queue_full_separately() {
|
||||||
|
let state = TransitionState::new_with_capacity(1);
|
||||||
|
let first_object = ObjectInfo {
|
||||||
|
bucket: "bucket".to_string(),
|
||||||
|
name: "first".to_string(),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
let second_object = ObjectInfo {
|
||||||
|
bucket: "bucket".to_string(),
|
||||||
|
name: "second".to_string(),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
let event = crate::bucket::lifecycle::lifecycle::Event {
|
||||||
|
action: IlmAction::TransitionAction,
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let first = state
|
||||||
|
.queue_transition_task_outcome(&first_object, &event, &LcEventSrc::Scanner)
|
||||||
|
.await;
|
||||||
|
let second = state
|
||||||
|
.queue_transition_task_outcome(&second_object, &event, &LcEventSrc::Scanner)
|
||||||
|
.await;
|
||||||
|
|
||||||
|
assert_eq!(first, TransitionEnqueueOutcome::Queued);
|
||||||
|
assert_eq!(second, TransitionEnqueueOutcome::QueueFull);
|
||||||
|
assert_eq!(state.transition_rx.len(), 1);
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial]
|
#[serial]
|
||||||
async fn queue_transition_task_dedupes_immediate_and_scanner_sources_for_same_version() {
|
async fn queue_transition_task_dedupes_immediate_and_scanner_sources_for_same_version() {
|
||||||
@@ -6277,6 +6328,22 @@ mod tests {
|
|||||||
assert_eq!(report.skipped_tier, 1);
|
assert_eq!(report.skipped_tier, 1);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn manual_transition_version_marker_preserves_null_version_cursor() {
|
||||||
|
let null_version = ObjectInfo {
|
||||||
|
version_id: None,
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
let version_id = Uuid::new_v4();
|
||||||
|
let versioned = ObjectInfo {
|
||||||
|
version_id: Some(version_id),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
assert_eq!(manual_transition_version_marker(&null_version), "null");
|
||||||
|
assert_eq!(manual_transition_version_marker(&versioned), version_id.to_string());
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn manual_transition_limit_boundary_only_truncates_when_more_objects_exist() {
|
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, 10, false));
|
||||||
|
|||||||
@@ -212,4 +212,45 @@ mod tests {
|
|||||||
|
|
||||||
assert_eq!(response_state(&report), "partial");
|
assert_eq!(response_state(&report), "partial");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn manual_transition_response_omits_raw_resume_markers() {
|
||||||
|
let report = ManualTransitionRunReport {
|
||||||
|
truncated_by_limit: true,
|
||||||
|
next_marker: Some("private/object".to_string()),
|
||||||
|
next_version_idmarker: Some("null".to_string()),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
let response = ManualTransitionRunResponse {
|
||||||
|
state: response_state(&report),
|
||||||
|
mode: "enqueue_only",
|
||||||
|
job_id: None,
|
||||||
|
status_endpoint: None,
|
||||||
|
report,
|
||||||
|
};
|
||||||
|
|
||||||
|
let value = serde_json::to_value(response).expect("response should serialize");
|
||||||
|
assert!(value.pointer("/report/next_marker").is_none());
|
||||||
|
assert!(value.pointer("/report/next_version_idmarker").is_none());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn manual_transition_handler_requires_set_tier_action() {
|
||||||
|
let src = include_str!("ilm_transition.rs");
|
||||||
|
let auth_block = extract_block_between_markers(src, "async fn authorize_manual_transition_request", "fn response_state");
|
||||||
|
|
||||||
|
assert!(auth_block.contains("AdminAction::SetTierAction"));
|
||||||
|
assert!(!auth_block.contains("AdminAction::ServerInfoAdminAction"));
|
||||||
|
}
|
||||||
|
|
||||||
|
fn extract_block_between_markers<'a>(src: &'a str, start_marker: &str, end_marker: &str) -> &'a str {
|
||||||
|
let start = src
|
||||||
|
.find(start_marker)
|
||||||
|
.unwrap_or_else(|| panic!("expected start marker `{start_marker}`"));
|
||||||
|
let after_start = &src[start..];
|
||||||
|
let end = after_start
|
||||||
|
.find(end_marker)
|
||||||
|
.unwrap_or_else(|| panic!("expected end marker `{end_marker}` after `{start_marker}`"));
|
||||||
|
&after_start[..end]
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user