From 26663e0d5dc0b356f6729d55b1354802cc96e1ef Mon Sep 17 00:00:00 2001 From: houseme Date: Sun, 26 Jul 2026 18:34:08 +0800 Subject: [PATCH] test(ilm): cover durable progress checkpoint renewal (#5287) Add regression coverage for durable manual transition progress checkpoints so page progress persists the report, queue snapshot, and renewed scope admission lease together. Co-authored-by: heihutu --- crates/e2e_test/src/reliant/tiering.rs | 106 +++++++++++------- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 63 ++++++++++- 2 files changed, 126 insertions(+), 43 deletions(-) diff --git a/crates/e2e_test/src/reliant/tiering.rs b/crates/e2e_test/src/reliant/tiering.rs index d29a9c425..fe2ab0bf4 100644 --- a/crates/e2e_test/src/reliant/tiering.rs +++ b/crates/e2e_test/src/reliant/tiering.rs @@ -87,6 +87,7 @@ const MANUAL_ASYNC_SAME_SCOPE_CONFLICT_PREFIX: &str = "manual-async-same-scope-c const MANUAL_TIER_FAILURE_PREFIX: &str = "manual-tier-failure/"; const MANUAL_WORKER_FAILURE_PREFIX: &str = "manual-worker-failure/"; const MANUAL_ACTIVE_CANCEL_PREFIX: &str = "manual-active-cancel/"; +const MANUAL_ASYNC_CONFLICT_OBJECTS: usize = 512; const MANUAL_ACTIVE_CANCEL_OBJECTS: usize = 512; const MANUAL_ACTIVE_CANCEL_RUNNING_TIMEOUT: StdDuration = StdDuration::from_secs(15); const OBJECT_KEY: &str = "tier/鲁A12345/report.bin"; @@ -523,6 +524,34 @@ async fn wait_for_manual_transition_job_terminal( } } +async fn wait_for_manual_transition_job_running( + hot: &RustFSTestEnvironment, + status_endpoint: &str, + deadline: StdDuration, +) -> Result> { + let start = Instant::now(); + loop { + let status = manual_transition_job_status(hot, status_endpoint).await?; + if status.status == "running" { + return Ok(status); + } + if matches!(status.status.as_str(), "completed" | "partial" | "cancelled" | "failed" | "unknown") { + return Err(format!( + "manual transition job reached terminal state before it became observable as running: {status:#?}" + ) + .into()); + } + if start.elapsed() >= deadline { + return Err(format!( + "manual transition job at {status_endpoint} did not become running within {}s; last={status:#?}", + deadline.as_secs() + ) + .into()); + } + tokio::time::sleep(StdDuration::from_millis(50)).await; + } +} + /// Number of objects currently stored in the cold-tier bucket. async fn cold_tier_object_count(cold_client: &Client) -> Result> { let resp = cold_client.list_objects_v2().bucket(TIER_BUCKET).send().await?; @@ -1002,7 +1031,7 @@ async fn test_manual_transition_async_overlapping_scope_conflict_reports_active_ add_rustfs_tier(&hot, &cold).await?; hot_client.create_bucket().bucket(MANUAL_ASYNC_CONFLICT_BUCKET).send().await?; - for idx in 0..50 { + for idx in 0..MANUAL_ASYNC_CONFLICT_OBJECTS { let key = format!("{MANUAL_ASYNC_CONFLICT_NESTED_PREFIX}obj-{idx:02}"); put_single_part_object(&hot_client, MANUAL_ASYNC_CONFLICT_BUCKET, &key, b"async conflict payload").await?; } @@ -1015,22 +1044,15 @@ async fn test_manual_transition_async_overlapping_scope_conflict_reports_active_ ) .await?; - let (first, second) = tokio::join!( - manual_transition_async_run_raw(&hot, MANUAL_ASYNC_CONFLICT_BUCKET, MANUAL_ASYNC_CONFLICT_PREFIX, false, 50), - manual_transition_async_run_raw(&hot, MANUAL_ASYNC_CONFLICT_BUCKET, MANUAL_ASYNC_CONFLICT_NESTED_PREFIX, false, 50) - ); - let responses = [first?, second?]; - let accepted = responses - .iter() - .find(|(status, _)| *status == reqwest::StatusCode::ACCEPTED) - .ok_or("one concurrent async run must be accepted")?; - let conflict = responses - .iter() - .find(|(status, _)| *status == reqwest::StatusCode::CONFLICT) - .ok_or("one concurrent async run must report conflict")?; - - let accepted: ManualTransitionRunResponse = serde_json::from_str(&accepted.1)?; - let conflict: ManualTransitionJobConflictResponse = serde_json::from_str(&conflict.1)?; + let before_remote_count = cold_tier_object_count(&cold_client).await?; + let accepted = manual_transition_async_run( + &hot, + MANUAL_ASYNC_CONFLICT_BUCKET, + MANUAL_ASYNC_CONFLICT_PREFIX, + false, + MANUAL_ASYNC_CONFLICT_OBJECTS as u64, + ) + .await?; let job_id = accepted .job_id .as_deref() @@ -1042,6 +1064,22 @@ async fn test_manual_transition_async_overlapping_scope_conflict_reports_active_ assert_eq!(accepted.state, "accepted"); assert_eq!(accepted.mode, "durable_job"); + let active = wait_for_manual_transition_job_running(&hot, status_endpoint, MANUAL_ACTIVE_CANCEL_RUNNING_TIMEOUT).await?; + assert_eq!(active.job_id, job_id); + let (conflict_status, conflict_body) = manual_transition_async_run_raw( + &hot, + MANUAL_ASYNC_CONFLICT_BUCKET, + MANUAL_ASYNC_CONFLICT_NESTED_PREFIX, + false, + MANUAL_ASYNC_CONFLICT_OBJECTS as u64, + ) + .await?; + assert_eq!( + conflict_status, + reqwest::StatusCode::CONFLICT, + "overlapping async run response: {conflict_body}" + ); + let conflict: ManualTransitionJobConflictResponse = serde_json::from_str(&conflict_body)?; assert_eq!(conflict.state, "conflict"); assert_eq!(conflict.mode, "durable_job"); assert_eq!(conflict.active_job_id, job_id); @@ -1054,16 +1092,19 @@ async fn test_manual_transition_async_overlapping_scope_conflict_reports_active_ assert_eq!(terminal.status, "completed", "terminal conflict winner response: {terminal:#?}"); assert!(!terminal.report.dry_run); assert_eq!(terminal.report.bucket, MANUAL_ASYNC_CONFLICT_BUCKET); - assert!( - terminal.report.prefix == MANUAL_ASYNC_CONFLICT_PREFIX || terminal.report.prefix == MANUAL_ASYNC_CONFLICT_NESTED_PREFIX, + assert_eq!(terminal.report.prefix, MANUAL_ASYNC_CONFLICT_PREFIX); + assert_eq!( + terminal.report.scanned, MANUAL_ASYNC_CONFLICT_OBJECTS as u64, + "terminal conflict winner response: {terminal:#?}" + ); + assert_eq!( + terminal.report.eligible, MANUAL_ASYNC_CONFLICT_OBJECTS as u64, "terminal conflict winner response: {terminal:#?}" ); - assert_eq!(terminal.report.scanned, 50, "terminal conflict winner response: {terminal:#?}"); - assert_eq!(terminal.report.eligible, 50, "terminal conflict winner response: {terminal:#?}"); assert_eq!(terminal.report.dry_run_eligible, 0, "terminal conflict winner response: {terminal:#?}"); assert_eq!( terminal.report.enqueued + terminal.report.skipped_already_in_flight, - 50, + MANUAL_ASYNC_CONFLICT_OBJECTS as u64, "terminal conflict winner response: {terminal:#?}" ); assert_eq!( @@ -1072,7 +1113,9 @@ async fn test_manual_transition_async_overlapping_scope_conflict_reports_active_ ); assert_eq!(terminal.report.transition_failed, 0, "terminal conflict winner response: {terminal:#?}"); assert_eq!(terminal.report.tier_failure, 0, "terminal conflict winner response: {terminal:#?}"); - assert!(cold_tier_object_count(&cold_client).await? <= 50); + let after_remote_count = cold_tier_object_count(&cold_client).await?; + assert!(after_remote_count >= before_remote_count); + assert!(after_remote_count <= before_remote_count + MANUAL_ASYNC_CONFLICT_OBJECTS); Ok(()) } @@ -1420,22 +1463,7 @@ async fn test_manual_transition_async_active_cancel_reports_terminal_cancelled() .as_deref() .ok_or("async response must include status_endpoint")?; - let start = Instant::now(); - let active = loop { - let status = manual_transition_job_status(&hot, status_endpoint).await?; - if status.status == "running" { - break status; - } - assert!( - !matches!(status.status.as_str(), "completed" | "partial" | "cancelled" | "failed" | "unknown"), - "manual transition job reached terminal state before active cancel could be requested: {status:#?}" - ); - assert!( - start.elapsed() < MANUAL_ACTIVE_CANCEL_RUNNING_TIMEOUT, - "manual transition job did not become running before active cancel timeout: {status:#?}" - ); - tokio::time::sleep(StdDuration::from_millis(50)).await; - }; + let active = wait_for_manual_transition_job_running(&hot, status_endpoint, MANUAL_ACTIVE_CANCEL_RUNNING_TIMEOUT).await?; assert_eq!(active.job_id, job_id); assert_eq!(active.status_endpoint, status_endpoint); assert!(!active.cancel_requested); diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index d5e7aa4d2..e25326860 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -4627,10 +4627,10 @@ mod tests { load_manual_transition_job_record, load_manual_transition_scope_admission, load_manual_transition_scope_admission_with_etag, manual_transition_scope_record_object_name, manual_transition_worker_result_object_name, manual_transition_worker_result_task_key, - reconcile_manual_transition_worker_results, record_manual_transition_worker_result, renew_manual_transition_job_lease, - request_manual_transition_job_cancel, save_manual_transition_job_record, - save_manual_transition_scope_admission_if_absent, save_manual_transition_scope_admission_if_current, - save_manual_transition_worker_result_if_absent, + persist_manual_transition_job_progress, reconcile_manual_transition_worker_results, + record_manual_transition_worker_result, renew_manual_transition_job_lease, request_manual_transition_job_cancel, + save_manual_transition_job_record, save_manual_transition_scope_admission_if_absent, + save_manual_transition_scope_admission_if_current, save_manual_transition_worker_result_if_absent, }; use crate::bucket::lifecycle::replication_sink::{ ReplicateDecision, ReplicateTargetDecision, ReplicationStatusType, VersionPurgeStatusType, @@ -7568,6 +7568,61 @@ mod tests { assert!(value.get("next_version_idmarker").is_none()); } + #[tokio::test] + #[serial] + async fn manual_transition_progress_persists_checkpoint_and_renews_admission() { + let (_paths, ecstore) = setup_test_env().await; + let job_id = Uuid::new_v4(); + let options = ManualTransitionRunOptions { + prefix: "logs/".to_string(), + ..Default::default() + }; + let record = ManualTransitionJobRecord::new(job_id, "manual-progress-journal-bucket", &options, "owner-a"); + let scope_key = record.scope_key.clone(); + save_manual_transition_job_record(ecstore.clone(), &record) + .await + .expect("running job record should save"); + save_manual_transition_scope_admission_if_absent(ecstore.clone(), &ManualTransitionScopeAdmission::from_job(&record)) + .await + .expect("running scope admission should save"); + let continuation_token = + encode_manual_transition_continuation_token(Some("logs/page-end".to_string()), Some("null".to_string())) + .expect("resume cursor should encode"); + let queue_snapshot = ManualTransitionQueueSnapshot { + queued: 1, + active: 2, + workers: 3, + ..Default::default() + }; + let report = ManualTransitionRunReport { + bucket: "manual-progress-journal-bucket".to_string(), + prefix: "logs/".to_string(), + scanned: 1000, + continuation_token: Some(continuation_token.clone()), + ..Default::default() + }; + + let persisted = persist_manual_transition_job_progress(ecstore.clone(), job_id, &report, queue_snapshot) + .await + .expect("page checkpoint should persist to the job record"); + + assert_eq!(persisted.state, ManualTransitionJobState::Running); + assert_eq!(persisted.report.scanned, 1000); + assert_eq!(persisted.report.continuation_token.as_deref(), Some(continuation_token.as_str())); + assert_eq!(persisted.queue_snapshot, queue_snapshot); + let loaded = load_manual_transition_job_record(ecstore.clone(), job_id) + .await + .expect("checkpointed job should reload"); + assert_eq!(loaded.report.continuation_token.as_deref(), Some(continuation_token.as_str())); + assert_eq!(loaded.queue_snapshot, queue_snapshot); + let admission = load_manual_transition_scope_admission(ecstore, &scope_key) + .await + .expect("running checkpoint must keep the scope admission alive"); + assert_eq!(admission.job_id, job_id); + assert_eq!(admission.lease_id, loaded.lease_id); + assert_eq!(admission.updated_at_unix_nanos, loaded.updated_at_unix_nanos); + } + #[tokio::test] async fn manual_transition_page_checkpoint_persists_resume_cursor() { let observed = Arc::new(StdMutex::new(Vec::new()));