mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-27 08:38:58 +00:00
test(ilm): cover manual transition restart recovery (#5318)
test: cover manual transition restart recovery Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -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<dyn std::error::Error + Send + Sync>> {
|
||||
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?;
|
||||
|
||||
Reference in New Issue
Block a user