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 <heihutu@gmail.com>
This commit is contained in:
houseme
2026-07-26 18:34:08 +08:00
committed by GitHub
parent caeaa4bf34
commit 26663e0d5d
2 changed files with 126 additions and 43 deletions
+67 -39
View File
@@ -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<ManualTransitionJobStatusResponse, Box<dyn std::error::Error + Send + Sync>> {
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<usize, Box<dyn std::error::Error + Send + Sync>> {
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);
@@ -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()));