mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-28 16:07:05 +00:00
fix(multipart): serialize complete and abort (#5356)
* fix(multipart): serialize complete and abort * test(multipart): order abort-first finalization * fix(multipart): enforce quorum staging cleanup * fix(multipart): remove stale mutable binding
This commit is contained in:
@@ -29,6 +29,12 @@ impl SetDisks {
|
|||||||
pub async fn delete_all(&self, bucket: &str, prefix: &str) -> Result<()> {
|
pub async fn delete_all(&self, bucket: &str, prefix: &str) -> Result<()> {
|
||||||
ListOperations::new(self.ctx()).delete_all(bucket, prefix).await
|
ListOperations::new(self.ctx()).delete_all(bucket, prefix).await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn delete_all_with_quorum(&self, bucket: &str, prefix: &str, write_quorum: usize) -> Result<()> {
|
||||||
|
ListOperations::new(self.ctx())
|
||||||
|
.delete_all_with_quorum(bucket, prefix, write_quorum)
|
||||||
|
.await
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// List/prefix maintenance operations, borrowing the `SetDisks` core state
|
/// List/prefix maintenance operations, borrowing the `SetDisks` core state
|
||||||
@@ -48,6 +54,14 @@ impl<'a> ListOperations<'a> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn delete_all(&self, bucket: &str, prefix: &str) -> Result<()> {
|
pub(crate) async fn delete_all(&self, bucket: &str, prefix: &str) -> Result<()> {
|
||||||
|
self.delete_all_inner(bucket, prefix, None).await
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn delete_all_with_quorum(&self, bucket: &str, prefix: &str, write_quorum: usize) -> Result<()> {
|
||||||
|
self.delete_all_inner(bucket, prefix, Some(write_quorum)).await
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn delete_all_inner(&self, bucket: &str, prefix: &str, write_quorum: Option<usize>) -> Result<()> {
|
||||||
let disks = self.ctx.disks().read().await;
|
let disks = self.ctx.disks().read().await;
|
||||||
|
|
||||||
let disks = disks.clone();
|
let disks = disks.clone();
|
||||||
@@ -79,6 +93,9 @@ impl<'a> ListOperations<'a> {
|
|||||||
Ok(_) => {
|
Ok(_) => {
|
||||||
errors.push(None);
|
errors.push(None);
|
||||||
}
|
}
|
||||||
|
Err(DiskError::FileNotFound | DiskError::PathNotFound | DiskError::VolumeNotFound) => {
|
||||||
|
errors.push(None);
|
||||||
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
errors.push(Some(e));
|
errors.push(Some(e));
|
||||||
}
|
}
|
||||||
@@ -97,6 +114,12 @@ impl<'a> ListOperations<'a> {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if let Some(write_quorum) = write_quorum
|
||||||
|
&& let Some(err) = reduce_write_quorum_errs(&errors, OBJECT_OP_IGNORED_ERRS, write_quorum)
|
||||||
|
{
|
||||||
|
return Err(err.into());
|
||||||
|
}
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -783,6 +783,9 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
|||||||
mut max_parts: usize,
|
mut max_parts: usize,
|
||||||
opts: &ObjectOptions,
|
opts: &ObjectOptions,
|
||||||
) -> Result<ListPartsInfo> {
|
) -> Result<ListPartsInfo> {
|
||||||
|
let _upload_guard = self
|
||||||
|
.acquire_multipart_upload_read_lock("list_object_parts", bucket, object, upload_id, opts)
|
||||||
|
.await?;
|
||||||
let (fi, _) = self.check_upload_id_exists(bucket, object, upload_id, false).await?;
|
let (fi, _) = self.check_upload_id_exists(bucket, object, upload_id, false).await?;
|
||||||
|
|
||||||
let upload_id_path = Self::get_upload_id_dir(bucket, object, upload_id);
|
let upload_id_path = Self::get_upload_id_dir(bucket, object, upload_id);
|
||||||
@@ -1260,10 +1263,15 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
|||||||
let _upload_guard = self
|
let _upload_guard = self
|
||||||
.acquire_multipart_upload_write_lock("abort_multipart_upload", bucket, object, upload_id, opts)
|
.acquire_multipart_upload_write_lock("abort_multipart_upload", bucket, object, upload_id, opts)
|
||||||
.await?;
|
.await?;
|
||||||
self.check_upload_id_exists(bucket, object, upload_id, false).await?;
|
let (fi, _) = self.check_upload_id_exists(bucket, object, upload_id, true).await?;
|
||||||
let upload_id_path = Self::get_upload_id_dir(bucket, object, upload_id);
|
let upload_id_path = Self::get_upload_id_dir(bucket, object, upload_id);
|
||||||
|
|
||||||
self.delete_all(RUSTFS_META_MULTIPART_BUCKET, &upload_id_path).await
|
self.delete_all_with_quorum(
|
||||||
|
RUSTFS_META_MULTIPART_BUCKET,
|
||||||
|
&upload_id_path,
|
||||||
|
fi.write_quorum(self.default_write_quorum()),
|
||||||
|
)
|
||||||
|
.await
|
||||||
}
|
}
|
||||||
// complete_multipart_upload finished
|
// complete_multipart_upload finished
|
||||||
#[tracing::instrument(skip(self))]
|
#[tracing::instrument(skip(self))]
|
||||||
@@ -1821,7 +1829,36 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
|||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
pause_multipart_commit(bucket, object, MultipartCommitPause::AfterRename).await;
|
pause_multipart_commit(bucket, object, MultipartCommitPause::AfterRename).await;
|
||||||
drop(upload_guard);
|
|
||||||
|
let cleanup_store = self.clone();
|
||||||
|
let cleanup_upload_id_path = upload_id_path.clone();
|
||||||
|
let cleanup_bucket = bucket.to_owned();
|
||||||
|
let cleanup_object = object.to_owned();
|
||||||
|
let cleanup_upload_id = upload_id.to_owned();
|
||||||
|
let cleanup_handle = tokio::spawn(async move {
|
||||||
|
let _upload_guard = upload_guard;
|
||||||
|
if let Err(err) = cleanup_store
|
||||||
|
.delete_all_with_quorum(RUSTFS_META_MULTIPART_BUCKET, &cleanup_upload_id_path, write_quorum)
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
warn!(
|
||||||
|
bucket = %cleanup_bucket,
|
||||||
|
object = %cleanup_object,
|
||||||
|
upload_id = %cleanup_upload_id,
|
||||||
|
error = ?err,
|
||||||
|
"completed multipart upload staging cleanup did not reach write quorum"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
if let Err(err) = cleanup_handle.await {
|
||||||
|
warn!(
|
||||||
|
bucket = %bucket,
|
||||||
|
object = %object,
|
||||||
|
upload_id = %upload_id,
|
||||||
|
error = ?err,
|
||||||
|
"completed multipart upload staging cleanup task failed"
|
||||||
|
);
|
||||||
|
}
|
||||||
drop(object_lock_guard); // drop object lock guard to release the lock
|
drop(object_lock_guard); // drop object lock guard to release the lock
|
||||||
|
|
||||||
// backlog#1321: enqueue heal only when the committed replicas actually
|
// backlog#1321: enqueue heal only when the committed replicas actually
|
||||||
@@ -1866,12 +1903,6 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
let upload_id_path = upload_id_path.clone();
|
|
||||||
let store = self.clone();
|
|
||||||
let _cleanup_handle = tokio::spawn(async move {
|
|
||||||
let _ = store.delete_all(RUSTFS_META_MULTIPART_BUCKET, &upload_id_path).await;
|
|
||||||
});
|
|
||||||
|
|
||||||
for (i, op_disk) in online_disks.iter().enumerate() {
|
for (i, op_disk) in online_disks.iter().enumerate() {
|
||||||
if let Some(disk) = op_disk
|
if let Some(disk) = op_disk
|
||||||
&& disk.is_online().await
|
&& disk.is_online().await
|
||||||
@@ -2183,6 +2214,122 @@ mod tests {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn assert_complete_first_linearizes(bucket: &'static str, object: &'static str, create_opts: ObjectOptions) {
|
||||||
|
let manager = Arc::new(rustfs_lock::GlobalLockManager::new());
|
||||||
|
let signaling = Arc::new(SignalingLockClient::new(Arc::new(LocalClient::with_manager(manager))));
|
||||||
|
let lockers: Vec<Arc<dyn LockClient>> = vec![signaling.clone()];
|
||||||
|
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_with_lockers(4, 0, 2, lockers).await;
|
||||||
|
make_bucket_on_all(&disk_stores, bucket).await;
|
||||||
|
let (upload_id, parts) = stage_upload_with_create_opts(&set_disks, bucket, object, &[0x47; 4096], &create_opts).await;
|
||||||
|
let upload_id_path = SetDisks::get_upload_id_dir(bucket, object, &upload_id);
|
||||||
|
signaling.set_target(rustfs_lock::ObjectKey::new(RUSTFS_META_MULTIPART_BUCKET, upload_id_path));
|
||||||
|
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())
|
||||||
|
.await
|
||||||
|
});
|
||||||
|
barrier.wait_until_paused().await;
|
||||||
|
|
||||||
|
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;
|
||||||
|
assert!(!abort.is_finished(), "abort must wait for the completion upload lock");
|
||||||
|
|
||||||
|
barrier.release();
|
||||||
|
complete
|
||||||
|
.await
|
||||||
|
.expect("completion task should not panic")
|
||||||
|
.expect("completion should win the upload finalization");
|
||||||
|
let abort_err = abort
|
||||||
|
.await
|
||||||
|
.expect("abort task should not panic")
|
||||||
|
.expect_err("abort must observe the upload as finalized");
|
||||||
|
assert!(matches!(abort_err, StorageError::InvalidUploadID(..)));
|
||||||
|
set_disks
|
||||||
|
.get_object_info(bucket, object, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("complete-first must leave the committed object readable");
|
||||||
|
assert!(matches!(
|
||||||
|
set_disks.check_upload_id_exists(bucket, object, &upload_id, false).await,
|
||||||
|
Err(StorageError::InvalidUploadID(..))
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn assert_abort_first_linearizes(bucket: &'static str, object: &'static str, create_opts: ObjectOptions) {
|
||||||
|
let manager = Arc::new(rustfs_lock::GlobalLockManager::new());
|
||||||
|
let signaling = Arc::new(SignalingLockClient::new(Arc::new(LocalClient::with_manager(manager))));
|
||||||
|
let lockers: Vec<Arc<dyn LockClient>> = vec![signaling.clone()];
|
||||||
|
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_with_lockers(4, 0, 2, lockers).await;
|
||||||
|
make_bucket_on_all(&disk_stores, bucket).await;
|
||||||
|
let (upload_id, parts) = stage_upload_with_create_opts(&set_disks, bucket, object, &[0x48; 4096], &create_opts).await;
|
||||||
|
let upload_id_path = SetDisks::get_upload_id_dir(bucket, object, &upload_id);
|
||||||
|
signaling.set_target(rustfs_lock::ObjectKey::new(RUSTFS_META_MULTIPART_BUCKET, upload_id_path.clone()));
|
||||||
|
let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await;
|
||||||
|
let object_holder = set_disks
|
||||||
|
.new_ns_lock(bucket, object)
|
||||||
|
.await
|
||||||
|
.expect("object namespace lock should be created")
|
||||||
|
.get_write_lock(Duration::from_secs(5))
|
||||||
|
.await
|
||||||
|
.expect("test should hold the object lock");
|
||||||
|
let holder = set_disks
|
||||||
|
.new_ns_lock(RUSTFS_META_MULTIPART_BUCKET, &upload_id_path)
|
||||||
|
.await
|
||||||
|
.expect("upload namespace lock should be created")
|
||||||
|
.get_write_lock(Duration::from_secs(5))
|
||||||
|
.await
|
||||||
|
.expect("test should hold the upload lock");
|
||||||
|
signaling.wait_for_attempts(1).await;
|
||||||
|
|
||||||
|
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;
|
||||||
|
|
||||||
|
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
|
||||||
|
});
|
||||||
|
drop(holder);
|
||||||
|
|
||||||
|
abort
|
||||||
|
.await
|
||||||
|
.expect("abort task should not panic")
|
||||||
|
.expect("abort should win the upload finalization");
|
||||||
|
drop(object_holder);
|
||||||
|
let complete_err = complete
|
||||||
|
.await
|
||||||
|
.expect("completion task should not panic")
|
||||||
|
.expect_err("completion must observe the aborted upload");
|
||||||
|
assert!(matches!(complete_err, StorageError::InvalidUploadID(..)));
|
||||||
|
let object_err = set_disks
|
||||||
|
.get_object_info(bucket, object, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect_err("abort-first must not publish an object");
|
||||||
|
assert!(matches!(object_err, StorageError::ObjectNotFound(..)));
|
||||||
|
assert!(matches!(
|
||||||
|
set_disks.check_upload_id_exists(bucket, object, &upload_id, false).await,
|
||||||
|
Err(StorageError::InvalidUploadID(..))
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
async fn assert_quorum_minus_one_retry_preserves_completable_part(
|
async fn assert_quorum_minus_one_retry_preserves_completable_part(
|
||||||
disk_count: usize,
|
disk_count: usize,
|
||||||
parity: usize,
|
parity: usize,
|
||||||
@@ -3045,6 +3192,81 @@ mod tests {
|
|||||||
.await;
|
.await;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test(flavor = "multi_thread")]
|
||||||
|
#[serial]
|
||||||
|
async fn abort_and_complete_linearize_for_plain_sse_and_legacy_layouts() {
|
||||||
|
assert_complete_first_linearizes("multipart-complete-first-plain", "object", ObjectOptions::default()).await;
|
||||||
|
assert_abort_first_linearizes("multipart-abort-first-plain", "object", ObjectOptions::default()).await;
|
||||||
|
|
||||||
|
let encrypted_opts = ObjectOptions {
|
||||||
|
user_defined: HashMap::from([(SSEC_ALGORITHM_HEADER.to_string(), "AES256".to_string())]),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
temp_env::async_with_vars([(crate::object_api::ENV_RUSTFS_ENCRYPTED_RANGE_SEEK, Some("true"))], async {
|
||||||
|
assert_complete_first_linearizes("multipart-complete-first-sse", "object", encrypted_opts.clone()).await;
|
||||||
|
assert_abort_first_linearizes("multipart-abort-first-sse", "object", encrypted_opts.clone()).await;
|
||||||
|
})
|
||||||
|
.await;
|
||||||
|
temp_env::async_with_vars([(crate::object_api::ENV_RUSTFS_ENCRYPTED_RANGE_SEEK, Some("false"))], async {
|
||||||
|
assert_complete_first_linearizes("multipart-complete-first-legacy", "object", encrypted_opts.clone()).await;
|
||||||
|
assert_abort_first_linearizes("multipart-abort-first-legacy", "object", encrypted_opts).await;
|
||||||
|
})
|
||||||
|
.await;
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn abort_enforces_delete_write_quorum_boundary() {
|
||||||
|
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
|
||||||
|
let bucket = "multipart-abort-delete-quorum";
|
||||||
|
let object = "object";
|
||||||
|
make_bucket_on_all(&disk_stores, bucket).await;
|
||||||
|
let quorum_upload = set_disks
|
||||||
|
.new_multipart_upload(bucket, object, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("multipart upload should be created");
|
||||||
|
|
||||||
|
let saved_disks = {
|
||||||
|
let mut disks = set_disks.disks.write().await;
|
||||||
|
let saved = disks.clone();
|
||||||
|
disks[3] = None;
|
||||||
|
saved
|
||||||
|
};
|
||||||
|
set_disks
|
||||||
|
.abort_multipart_upload(bucket, object, &quorum_upload.upload_id, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("abort should succeed at the exact delete write quorum");
|
||||||
|
*set_disks.disks.write().await = saved_disks;
|
||||||
|
assert!(matches!(
|
||||||
|
set_disks
|
||||||
|
.check_upload_id_exists(bucket, object, &quorum_upload.upload_id, false)
|
||||||
|
.await,
|
||||||
|
Err(StorageError::InvalidUploadID(..))
|
||||||
|
));
|
||||||
|
|
||||||
|
let below_quorum_upload = set_disks
|
||||||
|
.new_multipart_upload(bucket, object, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("second multipart upload should be created");
|
||||||
|
let saved_disks = {
|
||||||
|
let mut disks = set_disks.disks.write().await;
|
||||||
|
let saved = disks.clone();
|
||||||
|
disks[2] = None;
|
||||||
|
disks[3] = None;
|
||||||
|
saved
|
||||||
|
};
|
||||||
|
let err = set_disks
|
||||||
|
.abort_multipart_upload(bucket, object, &below_quorum_upload.upload_id, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect_err("abort must report a delete below write quorum");
|
||||||
|
assert!(matches!(err, StorageError::ErasureWriteQuorum));
|
||||||
|
|
||||||
|
*set_disks.disks.write().await = saved_disks;
|
||||||
|
set_disks
|
||||||
|
.check_upload_id_exists(bucket, object, &below_quorum_upload.upload_id, false)
|
||||||
|
.await
|
||||||
|
.expect("failed abort must leave quorum-visible staging on the restored disks");
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test(flavor = "multi_thread")]
|
#[tokio::test(flavor = "multi_thread")]
|
||||||
#[serial]
|
#[serial]
|
||||||
async fn complete_revalidates_layout_candidate_after_upload_lock() {
|
async fn complete_revalidates_layout_candidate_after_upload_lock() {
|
||||||
@@ -3176,6 +3398,17 @@ mod tests {
|
|||||||
tokio::task::yield_now().await;
|
tokio::task::yield_now().await;
|
||||||
assert!(!abort.is_finished(), "abort must wait until completion releases the upload lock");
|
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();
|
barrier.release();
|
||||||
complete
|
complete
|
||||||
.await
|
.await
|
||||||
@@ -3186,10 +3419,68 @@ mod tests {
|
|||||||
.expect("abort task should not panic")
|
.expect("abort task should not panic")
|
||||||
.expect_err("the committed upload should no longer exist when abort acquires the lock");
|
.expect_err("the committed upload should no longer exist when abort acquires the lock");
|
||||||
assert!(matches!(abort_err, StorageError::InvalidUploadID(..)));
|
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(..)));
|
||||||
})
|
})
|
||||||
.await;
|
.await;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test(flavor = "multi_thread")]
|
||||||
|
#[serial]
|
||||||
|
async fn complete_validates_parts_after_an_inflight_upload_part_commit() {
|
||||||
|
let manager = Arc::new(rustfs_lock::GlobalLockManager::new());
|
||||||
|
let signaling = Arc::new(SignalingLockClient::new(Arc::new(LocalClient::with_manager(manager))));
|
||||||
|
let lockers: Vec<Arc<dyn LockClient>> = vec![signaling.clone()];
|
||||||
|
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks_with_lockers(4, 0, 2, lockers).await;
|
||||||
|
let bucket = "multipart-complete-put-part-race-bucket";
|
||||||
|
let object = "object";
|
||||||
|
make_bucket_on_all(&disk_stores, bucket).await;
|
||||||
|
let (upload_id, original_parts) =
|
||||||
|
stage_upload_with_create_opts(&set_disks, bucket, object, &[0x49; 4096], &ObjectOptions::default()).await;
|
||||||
|
let upload_id_path = SetDisks::get_upload_id_dir(bucket, object, &upload_id);
|
||||||
|
signaling.set_target(rustfs_lock::ObjectKey::new(RUSTFS_META_MULTIPART_BUCKET, upload_id_path));
|
||||||
|
let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await;
|
||||||
|
let barrier = MultipartCommitBarrier::install(bucket, object, MultipartCommitPause::PutPartBeforeLockLost);
|
||||||
|
|
||||||
|
let put_store = set_disks.clone();
|
||||||
|
let put_upload_id = upload_id.clone();
|
||||||
|
let put = tokio::spawn(async move {
|
||||||
|
let mut reader = PutObjReader::from_vec(vec![0x4a; 4096]);
|
||||||
|
put_store
|
||||||
|
.put_object_part(bucket, object, &put_upload_id, 1, &mut reader, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
});
|
||||||
|
barrier.wait_until_paused().await;
|
||||||
|
|
||||||
|
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, original_parts, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
});
|
||||||
|
signaling.wait_for_attempts(2).await;
|
||||||
|
tokio::task::yield_now().await;
|
||||||
|
assert!(!complete.is_finished(), "completion must wait for the UploadPart commit lock");
|
||||||
|
|
||||||
|
barrier.release();
|
||||||
|
put.await
|
||||||
|
.expect("UploadPart task should not panic")
|
||||||
|
.expect("UploadPart replacement should commit");
|
||||||
|
let err = complete
|
||||||
|
.await
|
||||||
|
.expect("completion task should not panic")
|
||||||
|
.expect_err("completion must reject the stale ETag after UploadPart wins");
|
||||||
|
assert!(matches!(err, StorageError::InvalidPart(..)));
|
||||||
|
set_disks
|
||||||
|
.list_object_parts(bucket, object, &upload_id, None, MAX_PARTS_COUNT, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("failed completion must leave the upload retryable");
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test(start_paused = true)]
|
#[tokio::test(start_paused = true)]
|
||||||
#[serial]
|
#[serial]
|
||||||
async fn complete_fences_upload_lock_loss_before_commit() {
|
async fn complete_fences_upload_lock_loss_before_commit() {
|
||||||
|
|||||||
Reference in New Issue
Block a user