From b65210b1db008edcebd0c0e539cf28d46496983e Mon Sep 17 00:00:00 2001 From: cxymds Date: Wed, 29 Jul 2026 16:16:40 +0800 Subject: [PATCH] fix(ilm): preserve restore source version mode (#5406) * fix(ilm): preserve restore source version mode * test(ilm): cover suspended null-version restore * fix(ilm): reconcile worker results from journal --- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 20 ++- .../bucket/lifecycle/manual_transition_job.rs | 30 +---- .../src/app/lifecycle_transition_api_test.rs | 126 +++++++++++++++++- rustfs/src/app/object_usecase.rs | 9 +- 4 files changed, 151 insertions(+), 34 deletions(-) diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index e13a8caa3..cb7e8d3e0 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -8974,7 +8974,7 @@ mod tests { .await .expect("first worker result should persist"); assert_eq!(first.state, ManualTransitionJobState::Running); - assert_eq!(first.report.transition_completed, 1); + assert_eq!(first.report.transition_completed, 0); assert_eq!(first.report.transition_failed, 0); let duplicate = record_manual_transition_worker_result( @@ -8987,11 +8987,11 @@ mod tests { .await .expect("duplicate worker result should be idempotent"); assert_eq!(duplicate.state, ManualTransitionJobState::Running); - assert_eq!(duplicate.report.transition_completed, 1); + assert_eq!(duplicate.report.transition_completed, 0); assert_eq!(duplicate.report.transition_failed, 0); let second_key = manual_transition_worker_result_task_key(&bucket, "logs/b", None); - let final_record = record_manual_transition_worker_result( + let pending_record = record_manual_transition_worker_result( ecstore.clone(), job_id, &second_key, @@ -9000,7 +9000,14 @@ mod tests { ) .await .expect("second distinct worker result should persist"); + assert_eq!(pending_record.state, ManualTransitionJobState::Running); + assert_eq!(pending_record.report.transition_completed, 0); + assert_eq!(pending_record.report.transition_failed, 0); + let final_record = + reconcile_manual_transition_worker_results(ecstore.clone(), job_id, ManualTransitionQueueSnapshot::default()) + .await + .expect("worker result journal should reconcile"); assert_eq!(final_record.state, ManualTransitionJobState::Partial); assert_eq!(final_record.report.transition_completed, 1); assert_eq!(final_record.report.transition_failed, 1); @@ -9035,7 +9042,7 @@ mod tests { .expect("worker result job record should save"); let task_key = manual_transition_worker_result_task_key(&bucket, "logs/fail", None); - let final_record = record_manual_transition_worker_result_with_reason( + let pending_record = record_manual_transition_worker_result_with_reason( ecstore.clone(), job_id, &task_key, @@ -9045,7 +9052,12 @@ mod tests { ) .await .expect("worker result with failure reason should persist"); + assert!(pending_record.report.tier_failure_by_reason.is_empty()); + assert_eq!(pending_record.report.transition_failed, 0); + let final_record = reconcile_manual_transition_worker_results(ecstore, job_id, ManualTransitionQueueSnapshot::default()) + .await + .expect("worker failure reason should reconcile"); assert_eq!( final_record .report diff --git a/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs b/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs index 2555e2359..c340da298 100644 --- a/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs +++ b/crates/ecstore/src/bucket/lifecycle/manual_transition_job.rs @@ -1646,40 +1646,14 @@ pub async fn record_manual_transition_worker_result_with_reason( job_id: Uuid, task_key: &str, result: ManualTransitionWorkerResult, - queue_snapshot: ManualTransitionQueueSnapshot, + _queue_snapshot: ManualTransitionQueueSnapshot, failure_reason: Option, ) -> EcstoreResult { let result_record = ManualTransitionWorkerResultRecord::new_with_reason(job_id, task_key, result, failure_reason); if !save_manual_transition_worker_result_if_absent(api.clone(), &result_record).await? { return load_manual_transition_job_record(api, job_id).await; } - - for _ in 0..4 { - let (mut record, etag) = load_manual_transition_job_record_with_etag(api.clone(), job_id).await?; - if record.is_terminal() { - return Ok(record); - } - record.record_worker_result_with_reason(result, queue_snapshot, failure_reason); - 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), - } - } - Err(Error::PreconditionFailed) + load_manual_transition_job_record(api, job_id).await } pub async fn renew_manual_transition_job_lease( diff --git a/rustfs/src/app/lifecycle_transition_api_test.rs b/rustfs/src/app/lifecycle_transition_api_test.rs index 2053ebe5c..fc129194e 100644 --- a/rustfs/src/app/lifecycle_transition_api_test.rs +++ b/rustfs/src/app/lifecycle_transition_api_test.rs @@ -60,6 +60,7 @@ use uuid::Uuid; static GLOBAL_ENV: OnceLock<(Vec, Arc)> = OnceLock::new(); static INIT: Once = Once::new(); const TRANSITION_WAIT_TIMEOUT: Duration = Duration::from_secs(15); +const RESTORE_SUSPENDED_WAIT_TIMEOUT: Duration = Duration::from_secs(15); const ENV_GET_CODEC_STREAMING_ENABLE: &str = "RUSTFS_GET_CODEC_STREAMING_ENABLE"; const ENV_GET_CODEC_STREAMING_ROLLOUT: &str = "RUSTFS_GET_CODEC_STREAMING_ROLLOUT"; const ENV_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED: &str = "RUSTFS_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED"; @@ -161,6 +162,23 @@ async fn create_test_bucket(ecstore: &Arc, bucket_name: &str) { .expect("Failed to create test bucket"); } +async fn suspend_test_bucket(bucket: &str) { + DefaultBucketUsecase::from_global() + .execute_put_bucket_versioning(build_request( + PutBucketVersioningInput::builder() + .bucket(bucket.to_string()) + .versioning_configuration(VersioningConfiguration { + status: Some(BucketVersioningStatus::from_static(BucketVersioningStatus::SUSPENDED)), + ..Default::default() + }) + .build() + .expect("suspended versioning request should build"), + Method::PUT, + )) + .await + .expect("bucket versioning should be suspended"); +} + async fn upload_test_object(ecstore: &Arc, bucket: &str, object: &str, data: &[u8]) -> ObjectInfo { let mut reader = PutObjReader::from_vec(data.to_vec()); (**ecstore) @@ -338,6 +356,45 @@ async fn wait_for_transition(ecstore: &Arc, bucket: &str, object: &str, } } +async fn wait_for_restore_completion( + ecstore: &Arc, + backend: &MockWarmBackend, + bucket: &str, + object: &str, + timeout: Duration, +) -> Result { + let deadline = tokio::time::Instant::now() + timeout; + let mut last_state = None; + + loop { + if tokio::time::Instant::now() >= deadline { + let tier_gets = backend.get_count().await; + let op_log = backend.op_log().await; + return Err(format!( + "restore copy-back should complete within {timeout:?}; tier_gets={tier_gets}, op_log={op_log:?}; last observed state: {}", + last_state.unwrap_or_else(|| "no object info observed".to_string()) + )); + } + + match (**ecstore).get_object_info(bucket, object, &ObjectOptions::default()).await { + Ok(info) => { + if !info.restore_ongoing && info.restore_expires.is_some() { + return Ok(info); + } + last_state = Some(format!( + "restore_ongoing={}, restore_expires={:?}, transitioned_status={}", + info.restore_ongoing, info.restore_expires, info.transitioned_object.status + )); + } + Err(err) => { + last_state = Some(format!("get_object_info failed: {err}")); + } + } + + tokio::time::sleep(Duration::from_millis(500)).await; + } +} + // SAFETY: this helper is used only by `#[serial]` tests and runs under the single-threaded Tokio // runtime (`worker_threads = 1`), so no concurrent test can mutate process environment during the // `env::set_var` / `env::remove_var` window. @@ -2084,9 +2141,9 @@ async fn restore_object_usecase_reports_ongoing_conflict() { create_test_bucket(&ecstore, bucket.as_str()).await; let uploaded = upload_test_object(&ecstore, bucket.as_str(), object, &payload).await; + assert!(uploaded.version_id.is_none(), "fixture must create the null version"); let _ = transition_uploaded_object_directly(&ecstore, bucket.as_str(), object, &tier_name, &uploaded).await; backend.clear_op_log().await; - let get_barrier = backend.arm_get_barrier().await; let restore_request = || RestoreRequest { @@ -2138,6 +2195,73 @@ async fn restore_object_usecase_reports_ongoing_conflict() { get_barrier.release(); } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +#[serial] +#[ignore = "global-state ILM integration test: runs serialized in the CI ILM Integration (serial) lane, see ci.yml test-ilm-integration-serial and rustfs/backlog#4879"] +async fn restore_object_usecase_completes_suspended_null_version_in_place() { + let (_disk_paths, ecstore) = setup_test_env().await; + let usecase = DefaultObjectUsecase::from_global(); + let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + let backend = register_mock_tier(&tier_name).await; + let bucket = format!("test-api-restore-suspended-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object = "test/restore/suspended-null.bin"; + let payload: Vec = (0..128 * 1024).map(|i| (i % 251) as u8).collect(); + + create_test_bucket(&ecstore, bucket.as_str()).await; + let uploaded = upload_test_object(&ecstore, bucket.as_str(), object, &payload).await; + assert!(uploaded.version_id.is_none(), "fixture must create the null version"); + let _ = transition_uploaded_object_directly(&ecstore, bucket.as_str(), object, &tier_name, &uploaded).await; + suspend_test_bucket(bucket.as_str()).await; + backend.clear_op_log().await; + let tier_gets_before_restore = backend.get_count().await; + let get_barrier = backend.arm_get_barrier().await; + + Box::pin( + usecase.execute_restore_object(build_request( + RestoreObjectInput::builder() + .bucket(bucket.clone()) + .key(object.to_string()) + .restore_request(Some(RestoreRequest { + days: Some(1), + description: None, + glacier_job_parameters: None, + output_location: None, + select_parameters: None, + tier: None, + type_: None, + })) + .build() + .expect("restore request should build"), + Method::POST, + )), + ) + .await + .expect("suspended null-version restore should be accepted"); + + get_barrier.wait_until_paused().await; + get_barrier.release(); + let completed = wait_for_restore_completion(&ecstore, &backend, bucket.as_str(), object, RESTORE_SUSPENDED_WAIT_TIMEOUT) + .await + .unwrap_or_else(|err| panic!("{err}")); + + assert!(!completed.restore_ongoing, "the original null version must complete in place"); + assert!(completed.restore_expires.is_some(), "completed restore must carry an expiry"); + assert!( + completed.version_id.is_none() || completed.version_id.is_some_and(|version_id| version_id.is_nil()), + "suspended restore must remain on the null version" + ); + assert_eq!( + live_object_version_count(&ecstore, bucket.as_str(), object).await, + 1, + "suspended restore must not create a UUID version" + ); + assert_eq!( + backend.get_count().await - tier_gets_before_restore, + 1, + "suspended restore copy-back must fetch the tier exactly once" + ); +} + /// backlog#1304: the restore-accept compare-and-set itself, under real /// concurrency. Two POST ?restore for the same transitioned object race each /// other from separate tasks; the accept guard must let exactly one through — diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index e77f7c0e6..e361dd4ee 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -7783,7 +7783,12 @@ impl DefaultObjectUsecase { let bucket_clone = bucket.clone(); let object_clone = object.clone(); let rreq_clone = rreq.clone(); - let version_id_clone = obj_info_.version_id.map(|v| v.to_string()); + let version_id_clone = obj_info_ + .version_id + .map(|v| v.to_string()) + .or_else(|| (opts.versioned || opts.version_suspended).then(|| Uuid::nil().to_string())); + let versioned = opts.versioned; + let version_suspended = opts.version_suspended; let mut restore_operation_metadata = HashMap::new(); if let Some(id) = restore_operation_id { insert_str(&mut restore_operation_metadata, SUFFIX_RESTORE_OPERATION_ID, id.to_string()); @@ -7797,6 +7802,8 @@ impl DefaultObjectUsecase { ..Default::default() }, version_id: version_id_clone, + versioned, + version_suspended, user_defined: restore_operation_metadata, ..Default::default() };