From 0902538ceb6409264097c0088ca43dd295b45104 Mon Sep 17 00:00:00 2001 From: houseme Date: Mon, 27 Jul 2026 01:33:59 +0800 Subject: [PATCH] test(ilm): cover manual transition restart recovery (#5318) test: cover manual transition restart recovery Co-authored-by: heihutu --- crates/e2e_test/src/reliant/tiering.rs | 138 +++++++++++++++++++++++++ 1 file changed, 138 insertions(+) diff --git a/crates/e2e_test/src/reliant/tiering.rs b/crates/e2e_test/src/reliant/tiering.rs index d2d020eb9..ad447d0c6 100644 --- a/crates/e2e_test/src/reliant/tiering.rs +++ b/crates/e2e_test/src/reliant/tiering.rs @@ -80,6 +80,7 @@ const MANUAL_ASYNC_PARALLEL_BUCKET_B: &str = "ilm7-manual-async-parallel-b"; const MANUAL_TIER_FAILURE_BUCKET: &str = "ilm7-manual-tier-failure"; const MANUAL_WORKER_FAILURE_BUCKET: &str = "ilm7-manual-worker-failure"; const MANUAL_ACTIVE_CANCEL_BUCKET: &str = "ilm7-manual-active-cancel"; +const MANUAL_RESTART_CANCEL_BUCKET: &str = "ilm7-manual-restart-cancel"; const MANUAL_QUEUE_PRESSURE_PREFIX: &str = "manual-queue-pressure/"; const MANUAL_CONTINUATION_PREFIX: &str = "manual-continuation/"; const MANUAL_ASYNC_LIMIT_PREFIX: &str = "manual-async-limit/"; @@ -90,10 +91,13 @@ const MANUAL_ASYNC_PARALLEL_PREFIX: &str = "manual-async-parallel/"; 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_RESTART_CANCEL_PREFIX: &str = "manual-restart-cancel/"; const MANUAL_ASYNC_CONFLICT_OBJECTS: usize = 512; const MANUAL_ASYNC_PARALLEL_OBJECTS: usize = 64; const MANUAL_ACTIVE_CANCEL_OBJECTS: usize = 512; +const MANUAL_RESTART_CANCEL_OBJECTS: usize = 512; const MANUAL_ACTIVE_CANCEL_RUNNING_TIMEOUT: StdDuration = StdDuration::from_secs(15); +const MANUAL_RESTART_RECOVERY_TIMEOUT: StdDuration = StdDuration::from_secs(80); const OBJECT_KEY: &str = "tier/鲁A12345/report.bin"; const MANUAL_DUE_KEY: &str = "manual-due/report.bin"; const MANUAL_DRY_RUN_KEY: &str = "manual-dry-run/report.bin"; @@ -495,6 +499,7 @@ async fn manual_transition_job_status( let (status, body) = signed_admin_request(&hot.url, Method::GET, status_endpoint, None, &hot.access_key, &hot.secret_key).await?; assert_eq!(status, reqwest::StatusCode::OK, "manual transition job status response: {body}"); + assert_manual_transition_job_response_contract(&body, "manual transition job status response")?; Ok(serde_json::from_str(&body)?) } @@ -505,9 +510,48 @@ async fn manual_transition_job_cancel( let (status, body) = signed_admin_request(&hot.url, Method::DELETE, status_endpoint, None, &hot.access_key, &hot.secret_key).await?; assert_eq!(status, reqwest::StatusCode::OK, "manual transition job cancel response: {body}"); + assert_manual_transition_job_response_contract(&body, "manual transition job cancel response")?; Ok(serde_json::from_str(&body)?) } +fn assert_manual_transition_job_response_contract( + body: &str, + context: &str, +) -> Result<(), Box> { + let value: serde_json::Value = serde_json::from_str(body)?; + assert_eq!( + value.get("mode").and_then(serde_json::Value::as_str), + Some("durable_job"), + "{context} must keep the durable job mode: {body}" + ); + assert!( + value.get("job_id").and_then(serde_json::Value::as_str).is_some(), + "{context} must include job_id: {body}" + ); + assert!( + value.get("status_endpoint").and_then(serde_json::Value::as_str).is_some(), + "{context} must include status_endpoint: {body}" + ); + assert!( + value.get("cancel_endpoint").and_then(serde_json::Value::as_str).is_some(), + "{context} must include cancel_endpoint: {body}" + ); + for incompatible in [ + "scope_key", + "next_marker", + "next_version_idmarker", + "marker", + "version_marker", + "versionMarker", + ] { + assert!( + !body.contains(&format!("\"{incompatible}\"")), + "{context} must not expose incompatible field {incompatible}: {body}" + ); + } + Ok(()) +} + async fn wait_for_manual_transition_job_terminal( hot: &RustFSTestEnvironment, status_endpoint: &str, @@ -1732,6 +1776,100 @@ async fn test_manual_transition_async_active_cancel_reports_terminal_cancelled() Ok(()) } +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn test_manual_transition_async_cancel_after_process_restart_recovers_terminal() -> TestResult { + let mut cold = RustFSTestEnvironment::new().await?; + cold.access_key = "manualrestartcancelcoldadmin".to_string(); + cold.secret_key = "manualrestartcancelcoldsecret".to_string(); + cold.start_rustfs_server_without_cleanup(vec![]).await?; + let cold_client = cold.create_s3_client(); + cold_client.create_bucket().bucket(TIER_BUCKET).send().await?; + + let restart_env = [("RUSTFS_SCANNER_ENABLED", "false"), ("RUSTFS_SCANNER_CYCLE", "3600")]; + let mut hot = RustFSTestEnvironment::new().await?; + hot.start_rustfs_server_with_env(vec![], &restart_env).await?; + let hot_client = hot.create_s3_client(); + add_rustfs_tier(&hot, &cold).await?; + + hot_client.create_bucket().bucket(MANUAL_RESTART_CANCEL_BUCKET).send().await?; + put_lifecycle_transition_rule( + &hot_client, + MANUAL_RESTART_CANCEL_BUCKET, + "manual-restart-cancel", + MANUAL_RESTART_CANCEL_PREFIX, + 1, + ) + .await?; + for idx in 0..MANUAL_RESTART_CANCEL_OBJECTS { + let key = format!("{MANUAL_RESTART_CANCEL_PREFIX}obj-{idx:03}"); + put_single_part_object(&hot_client, MANUAL_RESTART_CANCEL_BUCKET, &key, b"manual restart cancel payload").await?; + } + + let before_remote_count = cold_tier_object_count(&cold_client).await?; + let accepted = + manual_transition_async_run(&hot, MANUAL_RESTART_CANCEL_BUCKET, MANUAL_RESTART_CANCEL_PREFIX, true, 10_000).await?; + let job_id = accepted.job_id.as_deref().ok_or("async response must include job_id")?; + let status_endpoint = accepted + .status_endpoint + .as_deref() + .ok_or("async response must include status_endpoint")?; + assert_eq!(accepted.cancel_endpoint.as_deref(), Some(status_endpoint)); + + 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!(!active.cancel_requested); + + hot.restart_server_preserving_data(vec![], &restart_env).await?; + + let restarted = manual_transition_job_status(&hot, status_endpoint).await?; + assert_eq!(restarted.job_id, job_id); + assert_eq!(restarted.status_endpoint, status_endpoint); + assert_eq!(restarted.cancel_endpoint, status_endpoint); + assert!( + matches!(restarted.status.as_str(), "running" | "completed" | "partial" | "cancelled" | "unknown"), + "restart status response must be running or terminal: {restarted:#?}" + ); + + let cancel_after_restart = manual_transition_job_cancel(&hot, status_endpoint).await?; + assert_eq!(cancel_after_restart.job_id, job_id); + assert_eq!(cancel_after_restart.status_endpoint, status_endpoint); + assert_eq!(cancel_after_restart.cancel_endpoint, status_endpoint); + if cancel_after_restart.status == "running" { + assert!( + cancel_after_restart.cancel_requested, + "cancel after restart must durably mark a still-running job: {cancel_after_restart:#?}" + ); + } + + let terminal = wait_for_manual_transition_job_terminal(&hot, status_endpoint, MANUAL_RESTART_RECOVERY_TIMEOUT).await?; + assert_eq!(terminal.job_id, job_id); + assert_eq!(terminal.status_endpoint, status_endpoint); + assert_eq!(terminal.cancel_endpoint, status_endpoint); + assert!( + matches!(terminal.status.as_str(), "unknown" | "cancelled" | "completed" | "partial"), + "post-restart manual transition job must fail closed or recover terminal: {terminal:#?}" + ); + if terminal.status == "unknown" { + assert!( + terminal + .failure_reason + .as_deref() + .is_some_and(|reason| reason.contains("restart") || reason.contains("owner loss")), + "unknown terminal status must explain the restart/owner-loss boundary: {terminal:#?}" + ); + } + assert_eq!(terminal.report.bucket, MANUAL_RESTART_CANCEL_BUCKET); + assert_eq!(terminal.report.prefix, MANUAL_RESTART_CANCEL_PREFIX); + assert!(terminal.report.dry_run); + assert_eq!( + cold_tier_object_count(&cold_client).await?, + before_remote_count, + "dry-run restart recovery must not create remote tier objects" + ); + + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn test_manual_transition_job_status_cancel_reject_unknown_jobs() -> TestResult { let mut hot = RustFSTestEnvironment::new().await?;