fix(ecstore): retry manual ILM job CAS updates (#6012)

Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
houseme
2026-08-13 02:40:10 +08:00
committed by GitHub
parent 59d8d93832
commit 398d2d87c8
7 changed files with 735 additions and 221 deletions
+5 -3
View File
@@ -61,9 +61,11 @@ pub mod bucket {
delete_manual_transition_scope_admission_if_current, load_manual_transition_job_record,
load_manual_transition_job_record_with_etag, load_manual_transition_scope_admission,
manual_transition_job_lease_expired, manual_transition_scope_admission_lease_expired,
manual_transition_scope_key, persist_manual_transition_job_progress, renew_manual_transition_job_lease,
request_manual_transition_job_cancel, save_manual_transition_job_record,
save_manual_transition_job_record_if_current, save_manual_transition_scope_admission_if_absent,
manual_transition_scope_key, persist_manual_transition_job_progress,
persist_manual_transition_job_progress_if_owned, renew_manual_transition_job_lease,
renew_manual_transition_job_lease_if_owned, request_manual_transition_job_cancel,
save_manual_transition_job_record, save_manual_transition_job_record_if_current,
save_manual_transition_scope_admission_if_absent, update_manual_transition_job_record,
};
}
@@ -27,9 +27,10 @@ use crate::bucket::lifecycle::manual_transition_job::{
ManualTransitionWorkerResult, claim_manual_transition_scope_admission, delete_manual_transition_scope_admission_if_current,
load_manual_transition_job_record, load_manual_transition_job_record_with_etag, load_manual_transition_pending_task_records,
manual_transition_job_id_from_record_object_name, manual_transition_job_lease_expired,
manual_transition_worker_result_task_key, persist_manual_transition_job_progress, reconcile_manual_transition_worker_results,
record_manual_transition_worker_result, record_manual_transition_worker_result_with_reason,
renew_manual_transition_job_lease, save_manual_transition_job_record_if_current, save_manual_transition_task_if_absent,
manual_transition_worker_result_task_key, persist_manual_transition_job_progress_if_owned,
reconcile_manual_transition_worker_results_if_owned, record_manual_transition_worker_result,
record_manual_transition_worker_result_with_reason, renew_manual_transition_job_lease_if_owned,
save_manual_transition_job_record_if_current, save_manual_transition_task_if_absent, update_manual_transition_job_record,
};
use crate::bucket::lifecycle::replication_sink;
use crate::bucket::lifecycle::replication_sink::{
@@ -2212,7 +2213,18 @@ async fn recover_manual_transition_job(
let recovery_unknown_snapshot = ManualTransitionQueueSnapshot::default();
if record.scan_completed {
let reconciled = reconcile_manual_transition_worker_results(api.clone(), job_id, recovery_unknown_snapshot).await?;
let reconciled = match reconcile_manual_transition_worker_results_if_owned(
api.clone(),
job_id,
record.lease_id,
recovery_unknown_snapshot,
)
.await
{
Ok(record) => record,
Err(Error::PreconditionFailed) => return Ok(ManualTransitionJobRecoveryOutcome::Skipped),
Err(err) => return Err(err),
};
if reconciled.is_terminal() {
release_manual_transition_recovery_admission(api, &reconciled).await;
return match reconciled.state {
@@ -2265,34 +2277,41 @@ async fn recover_manual_transition_job(
replay,
ManualTransitionPendingTaskReplay::Queued | ManualTransitionPendingTaskReplay::Deferred
) {
spawn_manual_transition_recovery_heartbeat(api, job_id);
spawn_manual_transition_recovery_heartbeat(api, job_id, recovery_lease_id);
return Ok(ManualTransitionJobRecoveryOutcome::Resumed);
}
let (mut record, etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?;
if record.mark_unknown_if_worker_results_lost(recovery_unknown_snapshot)
|| record.mark_unknown_if_recovery_would_skip_pending_page(recovery_unknown_snapshot)
let mut marked_unknown = false;
let record = match update_manual_transition_job_record(api.clone(), job_id, Some(recovery_lease_id), |record| {
marked_unknown = record.mark_unknown_if_worker_results_lost(recovery_unknown_snapshot)
|| record.mark_unknown_if_recovery_would_skip_pending_page(recovery_unknown_snapshot);
marked_unknown
})
.await
{
return match save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await {
Ok(()) => {
release_manual_transition_recovery_admission(api, &record).await;
Ok(ManualTransitionJobRecoveryOutcome::Unknown)
}
Err(Error::PreconditionFailed) => Ok(ManualTransitionJobRecoveryOutcome::Skipped),
Err(err) => Err(err),
};
Ok(record) => record,
Err(Error::PreconditionFailed) => return Ok(ManualTransitionJobRecoveryOutcome::Skipped),
Err(err) => return Err(err),
};
if marked_unknown {
release_manual_transition_recovery_admission(api, &record).await;
return Ok(ManualTransitionJobRecoveryOutcome::Unknown);
}
let mut options = record.resume_options();
options.job_id = Some(job_id);
options.cancel_check = Some(manual_transition_recovery_cancel_check(api.clone(), job_id));
options.progress_sink = Some(manual_transition_recovery_progress_sink(api.clone(), job_id));
options.progress_sink = Some(manual_transition_recovery_progress_sink(api.clone(), job_id, recovery_lease_id));
let result = enqueue_transition_for_existing_objects_scoped(api.clone(), &record.bucket, options).await;
let final_record = finalize_recovered_manual_transition_job(api.clone(), job_id, result).await?;
let final_record = match finalize_recovered_manual_transition_job(api.clone(), job_id, recovery_lease_id, result).await {
Ok(record) => record,
Err(Error::PreconditionFailed) => return Ok(ManualTransitionJobRecoveryOutcome::Skipped),
Err(err) => return Err(err),
};
if final_record.is_terminal() {
release_manual_transition_recovery_admission(api, &final_record).await;
} else {
spawn_manual_transition_recovery_heartbeat(api, job_id);
spawn_manual_transition_recovery_heartbeat(api, job_id, recovery_lease_id);
}
Ok(ManualTransitionJobRecoveryOutcome::Resumed)
}
@@ -2376,11 +2395,11 @@ fn manual_transition_recovery_cancel_check(api: Arc<ECStore>, job_id: Uuid) -> M
})
}
fn manual_transition_recovery_progress_sink(api: Arc<ECStore>, job_id: Uuid) -> ManualTransitionProgressSink {
fn manual_transition_recovery_progress_sink(api: Arc<ECStore>, job_id: Uuid, lease_id: Uuid) -> ManualTransitionProgressSink {
Arc::new(move |report| {
let api = api.clone();
Box::pin(async move {
persist_manual_transition_job_progress(api, job_id, &report, manual_transition_queue_snapshot())
persist_manual_transition_job_progress_if_owned(api, job_id, lease_id, &report, manual_transition_queue_snapshot())
.await
.map(|_| ())
})
@@ -2390,24 +2409,20 @@ fn manual_transition_recovery_progress_sink(api: Arc<ECStore>, job_id: Uuid) ->
async fn finalize_recovered_manual_transition_job(
api: Arc<ECStore>,
job_id: Uuid,
expected_lease_id: Uuid,
result: Result<ManualTransitionRunReport, Error>,
) -> Result<ManualTransitionJobRecord, Error> {
for _ in 0..4 {
let (mut record, etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?;
update_manual_transition_job_record(api, job_id, Some(expected_lease_id), |record| {
if record.is_terminal() {
return Ok(record);
return false;
}
match &result {
Ok(report) => record.complete(report.clone(), manual_transition_queue_snapshot()),
Err(err) => record.fail(format!("manual transition recovery failed: {err}")),
}
match save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await {
Ok(()) => return Ok(record),
Err(Error::PreconditionFailed) => continue,
Err(err) => return Err(err),
}
}
Err(Error::PreconditionFailed)
true
})
.await
}
async fn release_manual_transition_recovery_admission(api: Arc<ECStore>, record: &ManualTransitionJobRecord) {
@@ -2426,18 +2441,20 @@ async fn release_manual_transition_recovery_admission(api: Arc<ECStore>, record:
}
}
fn spawn_manual_transition_recovery_heartbeat(api: Arc<ECStore>, job_id: Uuid) {
fn spawn_manual_transition_recovery_heartbeat(api: Arc<ECStore>, job_id: Uuid, lease_id: Uuid) {
tokio::spawn(async move {
let mut interval = tokio::time::interval(std::time::Duration::from_secs(5));
loop {
interval.tick().await;
match renew_manual_transition_job_lease(api.clone(), job_id, manual_transition_queue_snapshot()).await {
match renew_manual_transition_job_lease_if_owned(api.clone(), job_id, lease_id, manual_transition_queue_snapshot())
.await
{
Ok(record) if record.is_terminal() => {
release_manual_transition_recovery_admission(api, &record).await;
return;
}
Ok(_) => {}
Err(Error::ConfigNotFound) => return,
Err(Error::ConfigNotFound | Error::PreconditionFailed) => return,
Err(err) => {
warn!(
event = EVENT_LIFECYCLE_WORKER_STATE,
@@ -2455,23 +2472,18 @@ fn spawn_manual_transition_recovery_heartbeat(api: Arc<ECStore>, job_id: Uuid) {
}
async fn abandon_manual_transition_recovery_lease(api: Arc<ECStore>, job_id: Uuid, lease_id: Uuid) -> Result<(), Error> {
for _ in 0..4 {
let (mut record, etag) = match load_manual_transition_job_record_with_etag(api.clone(), job_id).await {
Ok(record) => record,
Err(Error::ConfigNotFound) => return Ok(()),
Err(err) => return Err(err),
};
if record.lease_id != lease_id || record.is_terminal() {
return Ok(());
match update_manual_transition_job_record(api, job_id, Some(lease_id), |record| {
if record.is_terminal() {
return false;
}
record.abandon_recovery_lease(lease_id);
match save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await {
Ok(()) => return Ok(()),
Err(Error::PreconditionFailed) => continue,
Err(err) => return Err(err),
}
true
})
.await
{
Ok(_) | Err(Error::ConfigNotFound | Error::PreconditionFailed) => Ok(()),
Err(err) => Err(err),
}
Ok(())
}
fn tier_free_version_recovery_enabled() -> bool {
@@ -5089,12 +5101,13 @@ mod tests {
lifecycle_rule_has_date_expiration, manual_transition_duration_elapsed, manual_transition_has_more_after_limit,
manual_transition_recovery_progress_sink, manual_transition_version_marker, manual_transition_worker_failure_reason,
mark_delete_opts_skip_decommissioned_on_remote_success, merge_stale_multipart_candidate,
persist_manual_transition_job_progress, persist_manual_transition_page_checkpoint, recover_manual_transition_job,
recover_manual_transition_jobs, resolve_tier_free_version_recovery_enabled, resolve_transition_queue_capacity,
resolve_transition_queue_send_timeout, resolve_transition_worker_count, resolve_transition_workers_absolute_max,
run_tier_free_version_recovery_loop, select_restore_s3_location, set_lifecycle_observability_observer,
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,
persist_manual_transition_job_progress_if_owned, persist_manual_transition_page_checkpoint,
recover_manual_transition_job, recover_manual_transition_jobs, resolve_tier_free_version_recovery_enabled,
resolve_transition_queue_capacity, resolve_transition_queue_send_timeout, resolve_transition_worker_count,
resolve_transition_workers_absolute_max, run_tier_free_version_recovery_loop, select_restore_s3_location,
set_lifecycle_observability_observer, 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,
};
#[cfg(feature = "test-util")]
use super::{delete_free_version_remote_object_then, encode_dir_object, get_transitioned_object_reader_with_tier_manager};
@@ -5104,18 +5117,19 @@ mod tests {
};
use crate::bucket::lifecycle::config_boundary;
use crate::bucket::lifecycle::manual_transition_job::{
ManualTransitionJobRecord, ManualTransitionJobState, ManualTransitionScopeAdmission, ManualTransitionScopeAdmissionClaim,
ManualTransitionTaskRecord, ManualTransitionWorkerFailureReason, ManualTransitionWorkerResult,
ManualTransitionWorkerResultRecord, claim_manual_transition_scope_admission,
ManualTransitionJobCasBarrier, ManualTransitionJobRecord, ManualTransitionJobState, ManualTransitionScopeAdmission,
ManualTransitionScopeAdmissionClaim, ManualTransitionTaskRecord, ManualTransitionWorkerFailureReason,
ManualTransitionWorkerResult, ManualTransitionWorkerResultRecord, claim_manual_transition_scope_admission,
delete_manual_transition_scope_admission_if_current, legacy_manual_transition_scope_key,
load_manual_transition_job_record, load_manual_transition_scope_admission,
load_manual_transition_job_record, load_manual_transition_job_record_with_etag, load_manual_transition_scope_admission,
load_manual_transition_scope_admission_with_etag, load_manual_transition_task_record,
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, record_manual_transition_worker_result_with_reason,
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_task_if_absent, save_manual_transition_worker_result_if_absent,
renew_manual_transition_job_lease_if_owned, request_manual_transition_job_cancel, save_manual_transition_job_record,
save_manual_transition_job_record_if_current, save_manual_transition_scope_admission_if_absent,
save_manual_transition_scope_admission_if_current, save_manual_transition_task_if_absent,
save_manual_transition_worker_result_if_absent,
};
use crate::bucket::lifecycle::replication_sink::{ReplicationStatusType, VersionPurgeStatusType};
use crate::bucket::lifecycle::runtime_boundary as runtime_sources;
@@ -8553,9 +8567,10 @@ mod tests {
..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");
let persisted =
persist_manual_transition_job_progress_if_owned(ecstore.clone(), job_id, record.lease_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);
@@ -8574,6 +8589,232 @@ mod tests {
assert_eq!(admission.updated_at_unix_nanos, loaded.updated_at_unix_nanos);
}
#[tokio::test]
#[serial]
async fn manual_transition_progress_retries_heartbeat_cas_without_losing_checkpoint() {
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-cas-bucket", &options, "owner-a");
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 lease_id = record.lease_id;
let barrier = ManualTransitionJobCasBarrier::install(job_id);
let progress_store = ecstore.clone();
let progress = tokio::spawn(async move {
persist_manual_transition_job_progress_if_owned(
progress_store,
job_id,
lease_id,
&ManualTransitionRunReport {
bucket: "manual-progress-cas-bucket".to_string(),
prefix: "logs/".to_string(),
scanned: 1000,
eligible: 900,
enqueued: 800,
continuation_token: Some("opaque-page-cursor".to_string()),
..Default::default()
},
ManualTransitionQueueSnapshot {
queued: 7,
active: 3,
..Default::default()
},
)
.await
});
barrier.wait_until_paused().await;
let heartbeat = renew_manual_transition_job_lease_if_owned(
ecstore.clone(),
job_id,
lease_id,
ManualTransitionQueueSnapshot {
queued: 2,
active: 1,
..Default::default()
},
)
.await
.expect("heartbeat should win the first CAS write");
barrier.release();
let checkpointed = progress
.await
.expect("progress task should join")
.expect("progress should retry its stale ETag");
assert_eq!(checkpointed.lease_id, heartbeat.lease_id);
assert_eq!(checkpointed.report.scanned, 1000);
assert_eq!(checkpointed.report.eligible, 900);
assert_eq!(checkpointed.report.enqueued, 800);
assert_eq!(checkpointed.report.continuation_token.as_deref(), Some("opaque-page-cursor"));
assert_eq!(checkpointed.queue_snapshot.queued, 7);
assert_eq!(checkpointed.queue_snapshot.active, 3);
}
#[tokio::test]
#[serial]
async fn manual_transition_progress_rejects_stale_recovery_lease() {
let (_paths, ecstore) = setup_test_env().await;
let job_id = Uuid::new_v4();
let record = ManualTransitionJobRecord::new(
job_id,
"manual-progress-stale-lease-bucket",
&ManualTransitionRunOptions::default(),
"owner-a",
);
let stale_lease_id = record.lease_id;
save_manual_transition_job_record(ecstore.clone(), &record)
.await
.expect("running job record should save");
let (mut recovered, etag) = load_manual_transition_job_record_with_etag(ecstore.clone(), job_id)
.await
.expect("running job record should load");
recovered.lease_id = Uuid::new_v4();
recovered.owner_id = "owner-b".to_string();
save_manual_transition_job_record_if_current(ecstore.clone(), &recovered, &etag)
.await
.expect("recovery owner should replace the lease");
let error = persist_manual_transition_job_progress_if_owned(
ecstore.clone(),
job_id,
stale_lease_id,
&ManualTransitionRunReport {
scanned: 1000,
continuation_token: Some("stale-owner-cursor".to_string()),
..Default::default()
},
ManualTransitionQueueSnapshot::default(),
)
.await
.expect_err("the stale owner must not update the recovered job");
let heartbeat_error = renew_manual_transition_job_lease_if_owned(
ecstore.clone(),
job_id,
stale_lease_id,
ManualTransitionQueueSnapshot::default(),
)
.await
.expect_err("the stale owner must not renew the recovered job");
assert_eq!(error, Error::PreconditionFailed);
assert_eq!(heartbeat_error, Error::PreconditionFailed);
let loaded = load_manual_transition_job_record(ecstore, job_id)
.await
.expect("recovered job record should load");
assert_eq!(loaded.lease_id, recovered.lease_id);
assert_eq!(loaded.owner_id, "owner-b");
assert_eq!(loaded.report.scanned, 0);
assert!(loaded.report.continuation_token.is_none());
}
#[tokio::test]
#[serial]
async fn manual_transition_reconcile_rejects_lease_takeover_during_cas() {
let (_paths, ecstore) = setup_test_env().await;
let job_id = Uuid::new_v4();
let bucket = format!("manual-reconcile-lease-race-{}", job_id.simple());
let mut record = ManualTransitionJobRecord::new(job_id, &bucket, &ManualTransitionRunOptions::default(), "owner-a");
record.scan_completed = true;
let stale_lease_id = record.lease_id;
save_manual_transition_job_record(ecstore.clone(), &record)
.await
.expect("running job record should save");
let task_key = manual_transition_worker_result_task_key(&bucket, "logs/a", None);
let task = ManualTransitionTaskRecord::new(job_id, &task_key, &bucket, "logs/a", None, "WARM");
assert!(
save_manual_transition_task_if_absent(ecstore.clone(), &task)
.await
.expect("task journal marker should save")
);
let barrier = ManualTransitionJobCasBarrier::install(job_id);
let heartbeat_store = ecstore.clone();
let heartbeat = tokio::spawn(async move {
renew_manual_transition_job_lease_if_owned(
heartbeat_store,
job_id,
stale_lease_id,
ManualTransitionQueueSnapshot::default(),
)
.await
});
barrier.wait_until_paused().await;
let (mut recovered, etag) = load_manual_transition_job_record_with_etag(ecstore.clone(), job_id)
.await
.expect("running job record should load during reconciliation");
recovered.lease_id = Uuid::new_v4();
recovered.owner_id = "owner-b".to_string();
save_manual_transition_job_record_if_current(ecstore.clone(), &recovered, &etag)
.await
.expect("recovery owner should replace the lease");
barrier.release();
let error = heartbeat
.await
.expect("heartbeat task should join")
.expect_err("stale reconciliation must reject the recovery lease");
assert_eq!(error, Error::PreconditionFailed);
let loaded = load_manual_transition_job_record(ecstore, job_id)
.await
.expect("recovered job record should load");
assert_eq!(loaded.lease_id, recovered.lease_id);
assert_eq!(loaded.owner_id, "owner-b");
assert_eq!(loaded.state, ManualTransitionJobState::Running);
assert_eq!(loaded.report.enqueued, 0);
}
#[tokio::test]
#[serial]
async fn manual_transition_progress_does_not_regress_newer_admission_lease() {
let (_paths, ecstore) = setup_test_env().await;
let job_id = Uuid::new_v4();
let record = ManualTransitionJobRecord::new(
job_id,
"manual-progress-admission-order-bucket",
&ManualTransitionRunOptions::default(),
"owner-a",
);
save_manual_transition_job_record(ecstore.clone(), &record)
.await
.expect("running job record should save");
let mut newer_admission = ManualTransitionScopeAdmission::from_job(&record);
newer_admission.lease_expires_at_unix_nanos = newer_admission.lease_expires_at_unix_nanos.saturating_add(60_000_000_000);
newer_admission.updated_at_unix_nanos = newer_admission.updated_at_unix_nanos.saturating_add(60_000_000_000);
save_manual_transition_scope_admission_if_absent(ecstore.clone(), &newer_admission)
.await
.expect("newer scope admission should save");
persist_manual_transition_job_progress_if_owned(
ecstore.clone(),
job_id,
record.lease_id,
&ManualTransitionRunReport {
scanned: 1000,
..Default::default()
},
ManualTransitionQueueSnapshot::default(),
)
.await
.expect("progress should preserve the newer admission lease");
let admission = load_manual_transition_scope_admission(ecstore, &record.scope_key)
.await
.expect("scope admission should load");
assert_eq!(admission.lease_expires_at_unix_nanos, newer_admission.lease_expires_at_unix_nanos);
assert_eq!(admission.updated_at_unix_nanos, newer_admission.updated_at_unix_nanos);
}
#[tokio::test]
async fn manual_transition_page_checkpoint_persists_resume_cursor() {
let observed = Arc::new(StdMutex::new(Vec::new()));
@@ -8638,7 +8879,7 @@ mod tests {
.await
.expect("expired scope admission should save");
let checkpoint_options = ManualTransitionRunOptions {
progress_sink: Some(manual_transition_recovery_progress_sink(ecstore.clone(), job_id)),
progress_sink: Some(manual_transition_recovery_progress_sink(ecstore.clone(), job_id, record.lease_id)),
..options
};
let report = ManualTransitionRunReport {
@@ -8727,7 +8968,7 @@ mod tests {
prefix: prefix.to_string(),
tier: Some("WARM".to_string()),
dry_run: true,
progress_sink: Some(manual_transition_recovery_progress_sink(ecstore.clone(), job_id)),
progress_sink: Some(manual_transition_recovery_progress_sink(ecstore.clone(), job_id, record.lease_id)),
..Default::default()
};
let final_report = enqueue_transition_for_existing_objects_scoped(ecstore.clone(), &bucket, production_path_options)
@@ -9427,9 +9668,14 @@ mod tests {
"new worker result marker must be created"
);
let renewed = renew_manual_transition_job_lease(ecstore.clone(), job_id, ManualTransitionQueueSnapshot::default())
.await
.expect("heartbeat should reconcile marker before unknown fallback");
let renewed = renew_manual_transition_job_lease_if_owned(
ecstore.clone(),
job_id,
record.lease_id,
ManualTransitionQueueSnapshot::default(),
)
.await
.expect("heartbeat should reconcile marker before unknown fallback");
assert_eq!(renewed.state, ManualTransitionJobState::Completed);
assert_eq!(renewed.report.transition_completed, 1);
@@ -9472,9 +9718,14 @@ mod tests {
"new worker result marker must be created"
);
let renewed = renew_manual_transition_job_lease(ecstore.clone(), job_id, ManualTransitionQueueSnapshot::default())
.await
.expect("heartbeat should reconcile task and result journals");
let renewed = renew_manual_transition_job_lease_if_owned(
ecstore.clone(),
job_id,
record.lease_id,
ManualTransitionQueueSnapshot::default(),
)
.await
.expect("heartbeat should reconcile task and result journals");
assert_eq!(renewed.state, ManualTransitionJobState::Completed);
assert_eq!(renewed.report.enqueued, 1);
@@ -9789,9 +10040,10 @@ mod tests {
.await
.expect("running scope admission should save");
let checkpointed = persist_manual_transition_job_progress(
let checkpointed = persist_manual_transition_job_progress_if_owned(
ecstore.clone(),
job_id,
record.lease_id,
&ManualTransitionRunReport {
bucket: bucket.to_string(),
prefix: "logs/".to_string(),
@@ -9880,7 +10132,7 @@ mod tests {
compensation_running: 1,
};
let renewed = renew_manual_transition_job_lease(ecstore.clone(), job_id, queue_snapshot)
let renewed = renew_manual_transition_job_lease_if_owned(ecstore.clone(), job_id, record.lease_id, queue_snapshot)
.await
.expect("running job heartbeat should persist queue pressure status");
@@ -9931,9 +10183,14 @@ mod tests {
.await
.expect("running job admission should save");
let renewed = renew_manual_transition_job_lease(ecstore.clone(), job_id, ManualTransitionQueueSnapshot::default())
.await
.expect("lost worker result should persist unknown state");
let renewed = renew_manual_transition_job_lease_if_owned(
ecstore.clone(),
job_id,
record.lease_id,
ManualTransitionQueueSnapshot::default(),
)
.await
.expect("lost worker result should persist unknown state");
assert_eq!(renewed.state, ManualTransitionJobState::Unknown);
assert!(renewed.completed_at_unix_nanos.is_some());
@@ -86,6 +86,21 @@ where
com::save_config_with_opts(api, file, data, opts).await
}
pub(crate) async fn save_config_with_opts_quiet<S>(api: Arc<S>, file: &str, data: Vec<u8>, opts: &ObjectOptions) -> Result<()>
where
S: ObjectIO<
Error = Error,
RangeSpec = HTTPRangeSpec,
HeaderMap = HeaderMap,
ObjectOptions = ObjectOptions,
ObjectInfo = ObjectInfo,
GetObjectReader = GetObjectReader,
PutObjectReader = PutObjReader,
>,
{
com::save_config_with_opts_quiet(api, file, data, opts).await
}
pub(crate) async fn delete_config<S>(api: Arc<S>, file: &str) -> Result<()>
where
S: ObjectOperations<
@@ -45,6 +45,104 @@ const MANUAL_TRANSITION_JOB_LEASE_SECONDS: i128 = 60;
const MANUAL_TRANSITION_LEGACY_SCOPE_SCAN_LIMIT: i32 = 1000;
const MANUAL_TRANSITION_TASK_SCAN_LIMIT: i32 = 1000;
const MANUAL_TRANSITION_WORKER_RESULT_SCAN_LIMIT: i32 = 1000;
const MANUAL_TRANSITION_JOB_CAS_RETRIES: usize = 4;
#[cfg(test)]
struct ManualTransitionJobCasBarrierState {
job_id: Uuid,
paused: std::sync::atomic::AtomicBool,
arrived: tokio::sync::Notify,
release: tokio::sync::Semaphore,
}
#[cfg(test)]
pub(crate) struct ManualTransitionJobCasBarrier {
state: Arc<ManualTransitionJobCasBarrierState>,
}
#[cfg(test)]
static MANUAL_TRANSITION_JOB_CAS_BARRIER: std::sync::OnceLock<std::sync::Mutex<Option<Arc<ManualTransitionJobCasBarrierState>>>> =
std::sync::OnceLock::new();
#[cfg(test)]
impl ManualTransitionJobCasBarrier {
pub(crate) fn install(job_id: Uuid) -> Self {
let state = Arc::new(ManualTransitionJobCasBarrierState {
job_id,
paused: std::sync::atomic::AtomicBool::new(false),
arrived: tokio::sync::Notify::new(),
release: tokio::sync::Semaphore::new(0),
});
let mut slot = MANUAL_TRANSITION_JOB_CAS_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("manual transition progress CAS barrier mutex should not poison");
assert!(
slot.is_none(),
"manual transition job CAS barrier must be installed by one test at a time"
);
*slot = Some(Arc::clone(&state));
drop(slot);
Self { state }
}
pub(crate) async fn wait_until_paused(&self) {
tokio::time::timeout(std::time::Duration::from_secs(30), async {
loop {
let arrived = self.state.arrived.notified();
if self.state.paused.load(std::sync::atomic::Ordering::Acquire) {
return;
}
arrived.await;
}
})
.await
.expect("manual transition job update should reach the deterministic CAS barrier");
}
pub(crate) fn release(&self) {
self.state.release.add_permits(1);
}
}
#[cfg(test)]
impl Drop for ManualTransitionJobCasBarrier {
fn drop(&mut self) {
self.release();
let mut slot = MANUAL_TRANSITION_JOB_CAS_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("manual transition progress CAS barrier mutex should not poison");
if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) {
*slot = None;
}
}
}
#[cfg(test)]
async fn pause_manual_transition_job_before_first_cas(job_id: Uuid) {
let barrier = MANUAL_TRANSITION_JOB_CAS_BARRIER
.get_or_init(|| std::sync::Mutex::new(None))
.lock()
.expect("manual transition progress CAS barrier mutex should not poison")
.as_ref()
.filter(|barrier| barrier.job_id == job_id)
.cloned();
if let Some(barrier) = barrier
&& barrier
.paused
.compare_exchange(false, true, std::sync::atomic::Ordering::AcqRel, std::sync::atomic::Ordering::Acquire)
.is_ok()
{
barrier.arrived.notify_one();
barrier
.release
.acquire()
.await
.expect("manual transition job CAS barrier should remain open")
.forget();
}
}
fn is_false(value: &bool) -> bool {
!*value
@@ -148,7 +246,6 @@ impl ManualTransitionJobRecord {
pub fn fail(&mut self, error: impl Into<String>) {
self.state = ManualTransitionJobState::Failed;
self.report.tier_failure = self.report.tier_failure.saturating_add(1);
self.error = Some(error.into());
self.mark_updated_terminal();
}
@@ -1040,7 +1137,7 @@ pub async fn save_manual_transition_job_record_if_current(
}
let object = manual_transition_job_record_object_name(job.job_id).map_err(manual_transition_job_store_error)?;
let data = job.encode().map_err(manual_transition_job_store_error)?;
config_boundary::save_config_with_opts(
config_boundary::save_config_with_opts_quiet(
api,
&object,
data,
@@ -1056,6 +1153,54 @@ pub async fn save_manual_transition_job_record_if_current(
.await
}
/// Applies a job-record mutation with optimistic concurrency control.
///
/// The mutation returns whether the record needs to be persisted. When a lease
/// is supplied, ownership is checked again after every conflicting write.
pub async fn update_manual_transition_job_record<F>(
api: Arc<ECStore>,
job_id: Uuid,
expected_lease_id: Option<Uuid>,
update: F,
) -> EcstoreResult<ManualTransitionJobRecord>
where
F: FnMut(&mut ManualTransitionJobRecord) -> bool,
{
update_manual_transition_job_record_from(api, job_id, expected_lease_id, None, update).await
}
async fn update_manual_transition_job_record_from<F>(
api: Arc<ECStore>,
job_id: Uuid,
expected_lease_id: Option<Uuid>,
mut current: Option<(ManualTransitionJobRecord, String)>,
mut update: F,
) -> EcstoreResult<ManualTransitionJobRecord>
where
F: FnMut(&mut ManualTransitionJobRecord) -> bool,
{
for _ in 0..MANUAL_TRANSITION_JOB_CAS_RETRIES {
let (mut record, etag) = match current.take() {
Some(current) => current,
None => load_manual_transition_job_record_with_etag(api.clone(), job_id).await?,
};
if expected_lease_id.is_some_and(|lease_id| record.lease_id != lease_id) {
return Err(Error::PreconditionFailed);
}
if !update(&mut record) {
return Ok(record);
}
#[cfg(test)]
pause_manual_transition_job_before_first_cas(job_id).await;
match save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await {
Ok(()) => return Ok(record),
Err(Error::PreconditionFailed) => continue,
Err(err) => return Err(err),
}
}
Err(Error::PreconditionFailed)
}
pub(crate) async fn save_manual_transition_worker_result_if_absent(
api: Arc<ECStore>,
record: &ManualTransitionWorkerResultRecord,
@@ -1314,99 +1459,113 @@ pub async fn reconcile_manual_transition_worker_results(
api: Arc<ECStore>,
job_id: Uuid,
queue_snapshot: ManualTransitionQueueSnapshot,
) -> EcstoreResult<ManualTransitionJobRecord> {
reconcile_manual_transition_worker_results_inner(api, job_id, None, queue_snapshot, false).await
}
pub(crate) async fn reconcile_manual_transition_worker_results_if_owned(
api: Arc<ECStore>,
job_id: Uuid,
expected_lease_id: Uuid,
queue_snapshot: ManualTransitionQueueSnapshot,
) -> EcstoreResult<ManualTransitionJobRecord> {
reconcile_manual_transition_worker_results_inner(api, job_id, Some(expected_lease_id), queue_snapshot, false).await
}
async fn reconcile_manual_transition_worker_results_inner(
api: Arc<ECStore>,
job_id: Uuid,
expected_lease_id: Option<Uuid>,
queue_snapshot: ManualTransitionQueueSnapshot,
mark_missing_results_unknown: bool,
) -> EcstoreResult<ManualTransitionJobRecord> {
let task_stats = match scan_manual_transition_task_journal(api.clone(), job_id).await? {
ManualTransitionTaskJournal::Stats(stats) => stats,
ManualTransitionTaskJournal::Corrupt(error) => {
return mark_manual_transition_job_unknown_for_task_journal_error(api, job_id, error, queue_snapshot).await;
return mark_manual_transition_job_unknown_for_task_journal_error(
api,
job_id,
expected_lease_id,
error,
queue_snapshot,
)
.await;
}
};
let stats = match scan_manual_transition_worker_result_journal(api.clone(), job_id).await? {
ManualTransitionWorkerResultJournal::Stats(stats) => stats,
ManualTransitionWorkerResultJournal::Corrupt(error) => {
return mark_manual_transition_job_unknown_for_worker_result_journal_error(api, job_id, error, queue_snapshot).await;
return mark_manual_transition_job_unknown_for_worker_result_journal_error(
api,
job_id,
expected_lease_id,
error,
queue_snapshot,
)
.await;
}
};
for _ in 0..4 {
let (mut record, etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?;
let changed = record.apply_worker_result_counts(
let mut changed = false;
let record = update_manual_transition_job_record(api.clone(), job_id, expected_lease_id, |record| {
let counts_changed = record.apply_worker_result_counts(
stats.stats.completed,
stats.stats.failed,
&stats.stats.tier_failure_by_reason,
task_stats.queued,
queue_snapshot,
);
if !changed {
return Ok(record);
}
match save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await {
Ok(()) => {
if record.is_terminal() {
delete_manual_transition_scope_admission_if_current(
api.clone(),
&record.scope_key,
record.job_id,
record.lease_id,
)
.await?;
} else {
renew_manual_transition_scope_admission_from_job(api, &record).await?;
}
return Ok(record);
}
Err(Error::PreconditionFailed) => continue,
Err(err) => return Err(err),
}
let became_unknown = mark_missing_results_unknown && record.mark_unknown_if_worker_results_lost(queue_snapshot);
changed = counts_changed || became_unknown;
changed
})
.await?;
if !changed {
return Ok(record);
}
Err(Error::PreconditionFailed)
if record.is_terminal() {
delete_manual_transition_scope_admission_if_current(api, &record.scope_key, record.job_id, record.lease_id).await?;
} else {
renew_manual_transition_scope_admission_from_job(api, &record).await?;
}
Ok(record)
}
async fn mark_manual_transition_job_unknown_for_task_journal_error(
api: Arc<ECStore>,
job_id: Uuid,
expected_lease_id: Option<Uuid>,
error: String,
queue_snapshot: ManualTransitionQueueSnapshot,
) -> EcstoreResult<ManualTransitionJobRecord> {
for _ in 0..4 {
let (mut record, etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?;
if !record.mark_unknown_for_task_journal_error(error.clone(), queue_snapshot) {
return Ok(record);
}
match save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await {
Ok(()) => {
delete_manual_transition_scope_admission_if_current(api, &record.scope_key, record.job_id, record.lease_id)
.await?;
return Ok(record);
}
Err(Error::PreconditionFailed) => continue,
Err(err) => return Err(err),
}
let mut changed = false;
let record = update_manual_transition_job_record(api.clone(), job_id, expected_lease_id, |record| {
changed = record.mark_unknown_for_task_journal_error(error.clone(), queue_snapshot);
changed
})
.await?;
if changed && record.is_terminal() {
delete_manual_transition_scope_admission_if_current(api, &record.scope_key, record.job_id, record.lease_id).await?;
}
Err(Error::PreconditionFailed)
Ok(record)
}
async fn mark_manual_transition_job_unknown_for_worker_result_journal_error(
api: Arc<ECStore>,
job_id: Uuid,
expected_lease_id: Option<Uuid>,
error: String,
queue_snapshot: ManualTransitionQueueSnapshot,
) -> EcstoreResult<ManualTransitionJobRecord> {
for _ in 0..4 {
let (mut record, etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?;
if !record.mark_unknown_for_worker_result_journal_error(error.clone(), queue_snapshot) {
return Ok(record);
}
match save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await {
Ok(()) => {
delete_manual_transition_scope_admission_if_current(api, &record.scope_key, record.job_id, record.lease_id)
.await?;
return Ok(record);
}
Err(Error::PreconditionFailed) => continue,
Err(err) => return Err(err),
}
let mut changed = false;
let record = update_manual_transition_job_record(api.clone(), job_id, expected_lease_id, |record| {
changed = record.mark_unknown_for_worker_result_journal_error(error.clone(), queue_snapshot);
changed
})
.await?;
if changed && record.is_terminal() {
delete_manual_transition_scope_admission_if_current(api, &record.scope_key, record.job_id, record.lease_id).await?;
}
Err(Error::PreconditionFailed)
Ok(record)
}
pub async fn save_manual_transition_scope_admission_if_absent(
@@ -1603,19 +1762,14 @@ async fn find_active_legacy_manual_transition_scope_conflict(
}
pub async fn request_manual_transition_job_cancel(api: Arc<ECStore>, job_id: Uuid) -> EcstoreResult<ManualTransitionJobRecord> {
for _ in 0..4 {
let (mut record, etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?;
update_manual_transition_job_record(api, job_id, None, |record| {
if record.is_terminal() || record.cancel_requested {
return Ok(record);
return false;
}
record.mark_cancel_requested();
match save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await {
Ok(()) => return Ok(record),
Err(Error::PreconditionFailed) => continue,
Err(err) => return Err(err),
}
}
Err(Error::PreconditionFailed)
true
})
.await
}
pub async fn persist_manual_transition_job_progress(
@@ -1624,10 +1778,39 @@ pub async fn persist_manual_transition_job_progress(
report: &ManualTransitionRunReport,
queue_snapshot: ManualTransitionQueueSnapshot,
) -> EcstoreResult<ManualTransitionJobRecord> {
let (mut record, etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?;
record.update_running_progress(report.clone(), queue_snapshot);
save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await?;
renew_manual_transition_scope_admission_from_job(api, &record).await?;
let current = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?;
persist_manual_transition_job_progress_inner(api, job_id, current.0.lease_id, Some(current), report, queue_snapshot).await
}
pub async fn persist_manual_transition_job_progress_if_owned(
api: Arc<ECStore>,
job_id: Uuid,
expected_lease_id: Uuid,
report: &ManualTransitionRunReport,
queue_snapshot: ManualTransitionQueueSnapshot,
) -> EcstoreResult<ManualTransitionJobRecord> {
persist_manual_transition_job_progress_inner(api, job_id, expected_lease_id, None, report, queue_snapshot).await
}
async fn persist_manual_transition_job_progress_inner(
api: Arc<ECStore>,
job_id: Uuid,
expected_lease_id: Uuid,
current: Option<(ManualTransitionJobRecord, String)>,
report: &ManualTransitionRunReport,
queue_snapshot: ManualTransitionQueueSnapshot,
) -> EcstoreResult<ManualTransitionJobRecord> {
let record = update_manual_transition_job_record_from(api.clone(), job_id, Some(expected_lease_id), current, |record| {
if record.state != ManualTransitionJobState::Running {
return false;
}
record.update_running_progress(report.clone(), queue_snapshot);
true
})
.await?;
if record.state == ManualTransitionJobState::Running {
renew_manual_transition_scope_admission_from_job(api, &record).await?;
}
Ok(record)
}
@@ -1661,25 +1844,58 @@ pub async fn renew_manual_transition_job_lease(
job_id: Uuid,
queue_snapshot: ManualTransitionQueueSnapshot,
) -> EcstoreResult<ManualTransitionJobRecord> {
let (mut record, mut etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?;
if record.state == ManualTransitionJobState::Running {
if record.scan_completed && queue_snapshot.queued == 0 && queue_snapshot.active == 0 {
record = reconcile_manual_transition_worker_results(api.clone(), job_id, queue_snapshot).await?;
if record.is_terminal() || !record.report.worker_transition_pending() {
return Ok(record);
let current = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?;
renew_manual_transition_job_lease_inner(api, job_id, current.0.lease_id, Some(current), queue_snapshot).await
}
pub async fn renew_manual_transition_job_lease_if_owned(
api: Arc<ECStore>,
job_id: Uuid,
expected_lease_id: Uuid,
queue_snapshot: ManualTransitionQueueSnapshot,
) -> EcstoreResult<ManualTransitionJobRecord> {
renew_manual_transition_job_lease_inner(api, job_id, expected_lease_id, None, queue_snapshot).await
}
async fn renew_manual_transition_job_lease_inner(
api: Arc<ECStore>,
job_id: Uuid,
expected_lease_id: Uuid,
current: Option<(ManualTransitionJobRecord, String)>,
queue_snapshot: ManualTransitionQueueSnapshot,
) -> EcstoreResult<ManualTransitionJobRecord> {
let (current, current_etag) = match current {
Some(current) => current,
None => load_manual_transition_job_record_with_etag(api.clone(), job_id).await?,
};
if current.lease_id != expected_lease_id {
return Err(Error::PreconditionFailed);
}
if current.state != ManualTransitionJobState::Running {
return Ok(current);
}
if current.scan_completed && queue_snapshot.queued == 0 && queue_snapshot.active == 0 {
return reconcile_manual_transition_worker_results_inner(api, job_id, Some(expected_lease_id), queue_snapshot, true)
.await;
}
let record = update_manual_transition_job_record_from(
api.clone(),
job_id,
Some(expected_lease_id),
Some((current, current_etag)),
|record| {
if record.state != ManualTransitionJobState::Running {
return false;
}
(record, etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?;
}
let became_terminal = record.mark_unknown_if_worker_results_lost(queue_snapshot);
if !became_terminal {
record.renew_lease(queue_snapshot);
}
save_manual_transition_job_record_if_current(api.clone(), &record, &etag).await?;
if became_terminal {
delete_manual_transition_scope_admission_if_current(api, &record.scope_key, record.job_id, record.lease_id).await?;
} else {
renew_manual_transition_scope_admission_from_job(api, &record).await?;
}
true
},
)
.await?;
if record.is_terminal() {
delete_manual_transition_scope_admission_if_current(api, &record.scope_key, record.job_id, record.lease_id).await?;
} else if record.state == ManualTransitionJobState::Running {
renew_manual_transition_scope_admission_from_job(api, &record).await?;
}
Ok(record)
}
@@ -1688,15 +1904,31 @@ async fn renew_manual_transition_scope_admission_from_job(
api: Arc<ECStore>,
record: &ManualTransitionJobRecord,
) -> EcstoreResult<()> {
if let Ok((admission, admission_etag)) =
load_manual_transition_scope_admission_with_etag(api.clone(), &record.scope_key).await
&& admission.job_id == record.job_id
&& admission.lease_id == record.lease_id
{
let renewed_admission = ManualTransitionScopeAdmission::from_job(record);
save_manual_transition_scope_admission_if_current(api, &renewed_admission, &admission_etag).await?;
for _ in 0..MANUAL_TRANSITION_JOB_CAS_RETRIES {
let (admission, admission_etag) =
match load_manual_transition_scope_admission_with_etag(api.clone(), &record.scope_key).await {
Ok(admission) => admission,
Err(Error::ConfigNotFound) => return Ok(()),
Err(err) => return Err(err),
};
if admission.job_id != record.job_id || admission.lease_id != record.lease_id {
return Err(Error::PreconditionFailed);
}
let mut renewed_admission = ManualTransitionScopeAdmission::from_job(record);
renewed_admission.lease_expires_at_unix_nanos = renewed_admission
.lease_expires_at_unix_nanos
.max(admission.lease_expires_at_unix_nanos);
renewed_admission.updated_at_unix_nanos = renewed_admission.updated_at_unix_nanos.max(admission.updated_at_unix_nanos);
if renewed_admission == admission {
return Ok(());
}
match save_manual_transition_scope_admission_if_current(api.clone(), &renewed_admission, &admission_etag).await {
Ok(()) => return Ok(()),
Err(Error::PreconditionFailed) => continue,
Err(err) => return Err(err),
}
}
Ok(())
Err(Error::PreconditionFailed)
}
pub async fn delete_manual_transition_scope_admission_if_current(
@@ -2386,14 +2618,14 @@ mod tests {
}
#[test]
fn manual_transition_job_record_failure_counts_tier_failure() {
fn manual_transition_job_record_control_plane_failure_does_not_count_tier_failure() {
let options = ManualTransitionRunOptions::default();
let mut record = ManualTransitionJobRecord::new(Uuid::new_v4(), "bucket", &options, TEST_OWNER);
record.fail("missing tier");
assert_eq!(record.state, ManualTransitionJobState::Failed);
assert_eq!(record.report.tier_failure, 1);
assert_eq!(record.report.tier_failure, 0);
assert_eq!(record.error.as_deref(), Some("missing tier"));
}