diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 828bfc423..349bfc863 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -3379,85 +3379,98 @@ mod tests { #[tokio::test(flavor = "multi_thread")] #[serial] async fn complete_holds_object_then_upload_lock_through_commit() { - temp_env::async_with_vars([(crate::object_api::ENV_RUSTFS_ENCRYPTED_RANGE_SEEK, Some("true"))], async { - let manager = Arc::new(rustfs_lock::GlobalLockManager::new()); - let signaling = Arc::new(SignalingLockClient::new(Arc::new(LocalClient::with_manager(manager)))); - let lockers: Vec> = vec![signaling.clone()]; - let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_with_lockers(4, 0, 2, lockers).await; - let bucket = "multipart-layout-lock-order-bucket"; - let object = "object"; - make_bucket_on_all(&disk_stores, bucket).await; - let create_opts = ObjectOptions { - user_defined: HashMap::from([(SSEC_ALGORITHM_HEADER.to_string(), "AES256".to_string())]), - ..Default::default() - }; - let (upload_id, parts) = stage_upload_with_create_opts(&set_disks, bucket, object, &[0x43; 4096], &create_opts).await; - let upload_id_path = SetDisks::get_upload_id_dir(bucket, object, &upload_id); - let upload_resource = rustfs_lock::ObjectKey::new(RUSTFS_META_MULTIPART_BUCKET, upload_id_path); - let object_resource = rustfs_lock::ObjectKey::new(bucket, object); - signaling.set_target(upload_resource.clone()); - signaling.clear_observed(); - let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await; - let barrier = MultipartCommitBarrier::install(bucket, object, MultipartCommitPause::AfterRename); + temp_env::async_with_vars( + [ + (crate::object_api::ENV_RUSTFS_ENCRYPTED_RANGE_SEEK, Some("true")), + (rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("60")), + ], + async { + let manager = Arc::new(rustfs_lock::GlobalLockManager::new()); + let signaling = Arc::new(SignalingLockClient::new(Arc::new(LocalClient::with_manager(manager)))); + let lockers: Vec> = vec![signaling.clone()]; + let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_with_lockers(4, 0, 2, lockers).await; + let bucket = "multipart-layout-lock-order-bucket"; + let object = "object"; + make_bucket_on_all(&disk_stores, bucket).await; + let create_opts = ObjectOptions { + user_defined: HashMap::from([(SSEC_ALGORITHM_HEADER.to_string(), "AES256".to_string())]), + ..Default::default() + }; + let (upload_id, parts) = + stage_upload_with_create_opts(&set_disks, bucket, object, &[0x43; 4096], &create_opts).await; + let upload_id_path = SetDisks::get_upload_id_dir(bucket, object, &upload_id); + let upload_resource = rustfs_lock::ObjectKey::new(RUSTFS_META_MULTIPART_BUCKET, upload_id_path); + let object_resource = rustfs_lock::ObjectKey::new(bucket, object); + signaling.set_target(upload_resource.clone()); + signaling.clear_observed(); + let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await; + let barrier = MultipartCommitBarrier::install(bucket, object, MultipartCommitPause::AfterRename); - let complete_store = set_disks.clone(); - let complete_upload_id = upload_id.clone(); - let complete = tokio::spawn(async move { - complete_store - .complete_multipart_upload(bucket, object, &complete_upload_id, parts, &ObjectOptions::default()) + let complete_store = set_disks.clone(); + let complete_upload_id = upload_id.clone(); + let complete = tokio::spawn(async move { + complete_store + .complete_multipart_upload(bucket, object, &complete_upload_id, parts, &ObjectOptions::default()) + .await + }); + barrier.wait_until_paused().await; + + let observed = signaling.observed(); + let object_position = observed + .iter() + .position(|resource| resource == &object_resource) + .expect("completion should acquire the object lock"); + let upload_position = observed + .iter() + .position(|resource| resource == &upload_resource) + .expect("completion should acquire the upload lock"); + assert!(object_position < upload_position, "completion must acquire object before upload"); + + let abort_store = set_disks.clone(); + let abort_upload_id = upload_id.clone(); + let abort = tokio::spawn(async move { + abort_store + .abort_multipart_upload(bucket, object, &abort_upload_id, &ObjectOptions::default()) + .await + }); + signaling.wait_for_attempts(2).await; + tokio::task::yield_now().await; + assert!(!abort.is_finished(), "abort must wait until completion releases the upload lock"); + + let list_store = set_disks.clone(); + let list_upload_id = upload_id.clone(); + let list = tokio::spawn(async move { + list_store + .list_object_parts(bucket, object, &list_upload_id, None, MAX_PARTS_COUNT, &ObjectOptions::default()) + .await + }); + signaling.wait_for_attempts(3).await; + tokio::task::yield_now().await; + assert!(!list.is_finished(), "ListParts must wait until completion releases the upload lock"); + + barrier.release(); + complete .await - }); - barrier.wait_until_paused().await; - - let observed = signaling.observed(); - let object_position = observed - .iter() - .position(|resource| resource == &object_resource) - .expect("completion should acquire the object lock"); - let upload_position = observed - .iter() - .position(|resource| resource == &upload_resource) - .expect("completion should acquire the upload lock"); - assert!(object_position < upload_position, "completion must acquire object before upload"); - - let abort_store = set_disks.clone(); - let abort_upload_id = upload_id.clone(); - let abort = tokio::spawn(async move { - abort_store - .abort_multipart_upload(bucket, object, &abort_upload_id, &ObjectOptions::default()) + .expect("completion task should not panic") + .expect("completion should commit after the barrier is released"); + let abort_err = abort .await - }); - signaling.wait_for_attempts(2).await; - tokio::task::yield_now().await; - assert!(!abort.is_finished(), "abort must wait until completion releases the upload lock"); - - let list_store = set_disks.clone(); - let list_upload_id = upload_id.clone(); - let list = tokio::spawn(async move { - list_store - .list_object_parts(bucket, object, &list_upload_id, None, MAX_PARTS_COUNT, &ObjectOptions::default()) + .expect("abort task should not panic") + .expect_err("the committed upload should no longer exist when abort acquires the lock"); + assert!( + matches!(abort_err, StorageError::InvalidUploadID(..)), + "abort should return InvalidUploadID after completion, got {abort_err:?}" + ); + let list_err = list .await - }); - signaling.wait_for_attempts(3).await; - tokio::task::yield_now().await; - assert!(!list.is_finished(), "ListParts must wait until completion releases the upload lock"); - - barrier.release(); - complete - .await - .expect("completion task should not panic") - .expect("completion should commit after the barrier is released"); - let abort_err = abort - .await - .expect("abort task should not panic") - .expect_err("the committed upload should no longer exist when abort acquires the lock"); - assert!(matches!(abort_err, StorageError::InvalidUploadID(..))); - let list_err = list - .await - .expect("ListParts task should not panic") - .expect_err("the committed upload should no longer exist when ListParts acquires the lock"); - assert!(matches!(list_err, StorageError::InvalidUploadID(..))); - }) + .expect("ListParts task should not panic") + .expect_err("the committed upload should no longer exist when ListParts acquires the lock"); + assert!( + matches!(list_err, StorageError::InvalidUploadID(..)), + "ListParts should return InvalidUploadID after completion, got {list_err:?}" + ); + }, + ) .await; }