mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-02 03:19:19 +00:00
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
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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<ManualTransitionWorkerFailureReason>,
|
||||
) -> EcstoreResult<ManualTransitionJobRecord> {
|
||||
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(
|
||||
|
||||
@@ -60,6 +60,7 @@ use uuid::Uuid;
|
||||
static GLOBAL_ENV: OnceLock<(Vec<PathBuf>, Arc<ECStore>)> = 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<ECStore>, 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<ECStore>, 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<ECStore>, bucket: &str, object: &str,
|
||||
}
|
||||
}
|
||||
|
||||
async fn wait_for_restore_completion(
|
||||
ecstore: &Arc<ECStore>,
|
||||
backend: &MockWarmBackend,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
timeout: Duration,
|
||||
) -> Result<ObjectInfo, String> {
|
||||
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<u8> = (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 —
|
||||
|
||||
@@ -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()
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user