mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-05 19:55:37 +00:00
fix(ecstore): wait for multipart copy readiness (#7065)
This commit is contained in:
@@ -66,7 +66,6 @@ use crate::disk::new_disk;
|
|||||||
use crate::multipart_listing::paginate_multipart_listing;
|
use crate::multipart_listing::paginate_multipart_listing;
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
use crate::object_api::ObjectLockConfigSnapshot;
|
use crate::object_api::ObjectLockConfigSnapshot;
|
||||||
use crate::set_disk::core::io_primitives::finish_rename_tail_heal;
|
|
||||||
use crate::set_disk::mem;
|
use crate::set_disk::mem;
|
||||||
use crate::set_disk::metadata_sys;
|
use crate::set_disk::metadata_sys;
|
||||||
use crate::set_disk::runtime_sources;
|
use crate::set_disk::runtime_sources;
|
||||||
@@ -3124,8 +3123,11 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
|||||||
let commit_object_lock_guard = object_lock_guard.take();
|
let commit_object_lock_guard = object_lock_guard.take();
|
||||||
let commit_decommission_object_lock_guard = decommission_object_lock_guard.take();
|
let commit_decommission_object_lock_guard = decommission_object_lock_guard.take();
|
||||||
let commit_decommission_capacity_guard = decommission_capacity_guard.take();
|
let commit_decommission_capacity_guard = decommission_capacity_guard.take();
|
||||||
let commit_allows_early_ack = !(opts.data_movement && opts.has_decommission_capacity_reservation())
|
// CompleteMultipartUpload is an S3 publication boundary: after a
|
||||||
&& (commit_object_lock_guard.is_some() || commit_decommission_object_lock_guard.is_some());
|
// successful response, the object must be immediately readable and
|
||||||
|
// usable as a CopyObject source. Do not return on rename quorum while
|
||||||
|
// a tail owner may still hold the object guard and finish shard moves.
|
||||||
|
let commit_allows_early_ack = false;
|
||||||
let detach_commit_owner = commit_allows_early_ack || upload_guard.is_some() || quota_mutation_fence;
|
let detach_commit_owner = commit_allows_early_ack || upload_guard.is_some() || quota_mutation_fence;
|
||||||
let commit = async move {
|
let commit = async move {
|
||||||
let mut _object_lock_guard = commit_object_lock_guard;
|
let mut _object_lock_guard = commit_object_lock_guard;
|
||||||
@@ -3256,105 +3258,15 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
|||||||
commit_allows_early_ack,
|
commit_allows_early_ack,
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
let mut rename_guard_release = None;
|
|
||||||
let mut needs_immediate_heal = false;
|
|
||||||
let mut tail_owns_staging_cleanup = false;
|
|
||||||
if let Ok(rename_commit) = rename_result.as_mut() {
|
if let Ok(rename_commit) = rename_result.as_mut() {
|
||||||
commit_set.record_capacity_scope_if_needed(commit_capacity_scope_token, &rename_commit.capacity_disks);
|
commit_set.record_capacity_scope_if_needed(commit_capacity_scope_token, &rename_commit.capacity_disks);
|
||||||
// Install the tail watcher before any post-commit await. The
|
debug_assert!(
|
||||||
// latch keeps namespace guards through their prior handoff point.
|
rename_commit.tail_drain.is_none(),
|
||||||
needs_immediate_heal = rename_commit.needs_immediate_heal();
|
"multipart completion disables early ACK and must not detach a rename tail"
|
||||||
if let Some(rename_tail_drain) = rename_commit.tail_drain.take() {
|
);
|
||||||
tail_owns_staging_cleanup = true;
|
|
||||||
let mut request = rustfs_heal_contracts::heal_channel::create_heal_request_with_options(
|
|
||||||
commit_bucket.clone(),
|
|
||||||
Some(commit_object.clone()),
|
|
||||||
false,
|
|
||||||
Some(HealChannelPriority::Normal),
|
|
||||||
Some(commit_set.pool_index),
|
|
||||||
Some(commit_set.set_index),
|
|
||||||
);
|
|
||||||
request.object_version_id = fi
|
|
||||||
.version_id
|
|
||||||
.or_else(|| commit_version_suspended.then(Uuid::nil))
|
|
||||||
.map(|version_id| version_id.to_string());
|
|
||||||
let object_lock_guard = _object_lock_guard.take();
|
|
||||||
let upload_guard = _upload_guard.take();
|
|
||||||
let decommission_object_lock_guard = _decommission_object_lock_guard.take();
|
|
||||||
let decommission_capacity_guard = _decommission_capacity_guard.take();
|
|
||||||
let cleanup_bucket = commit_bucket.clone();
|
|
||||||
let cleanup_object = commit_object.clone();
|
|
||||||
let heal_set = commit_set.clone();
|
|
||||||
let cleanup_set = commit_set.clone();
|
|
||||||
let committed_data_dir = fi.data_dir;
|
|
||||||
let cleanup_parts = parts.clone();
|
|
||||||
let cleanup_upload_path = commit_upload_id_path.clone();
|
|
||||||
let cleanup_upload_id = commit_upload_id.clone();
|
|
||||||
let fence_disks = commit_disks.clone();
|
|
||||||
let fence_tokens = quota_fence_tokens.clone();
|
|
||||||
let fence_bucket = commit_bucket.clone();
|
|
||||||
let fence_object = commit_object.clone();
|
|
||||||
let (guard_release_tx, guard_release_rx) = tokio::sync::oneshot::channel();
|
|
||||||
rename_guard_release = Some(guard_release_tx);
|
|
||||||
tokio::spawn(finish_rename_tail_heal(
|
|
||||||
rename_tail_drain,
|
|
||||||
guard_release_rx,
|
|
||||||
(
|
|
||||||
object_lock_guard,
|
|
||||||
upload_guard,
|
|
||||||
decommission_object_lock_guard,
|
|
||||||
decommission_capacity_guard,
|
|
||||||
),
|
|
||||||
request,
|
|
||||||
move || async move {
|
|
||||||
if quota_mutation_fence {
|
|
||||||
let _ = SetDisks::release_quota_mutation_fences(
|
|
||||||
&fence_disks,
|
|
||||||
&fence_tokens,
|
|
||||||
&fence_bucket,
|
|
||||||
&fence_object,
|
|
||||||
write_quorum,
|
|
||||||
)
|
|
||||||
.await;
|
|
||||||
}
|
|
||||||
},
|
|
||||||
move |(object_lock_guard, upload_guard, decommission_object_lock_guard, decommission_capacity_guard),
|
|
||||||
targets| async move {
|
|
||||||
drop(object_lock_guard);
|
|
||||||
cleanup_set.cleanup_multipart_path(&cleanup_parts).await;
|
|
||||||
cleanup_set
|
|
||||||
.cleanup_rename_tail(
|
|
||||||
targets,
|
|
||||||
&cleanup_bucket,
|
|
||||||
&cleanup_object,
|
|
||||||
committed_data_dir,
|
|
||||||
transaction_epoch,
|
|
||||||
)
|
|
||||||
.await;
|
|
||||||
if let Err(err) = cleanup_set
|
|
||||||
.delete_all_with_quorum(RUSTFS_META_MULTIPART_BUCKET, &cleanup_upload_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"
|
|
||||||
);
|
|
||||||
}
|
|
||||||
drop(upload_guard);
|
|
||||||
drop(decommission_object_lock_guard);
|
|
||||||
drop(decommission_capacity_guard);
|
|
||||||
},
|
|
||||||
|request| async move { heal_set.submit_rename_tail_heal(request).await },
|
|
||||||
));
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
if !tail_owns_staging_cleanup {
|
drop(_decommission_capacity_guard.take());
|
||||||
drop(_decommission_capacity_guard.take());
|
if quota_mutation_fence {
|
||||||
}
|
|
||||||
if quota_mutation_fence && !tail_owns_staging_cleanup {
|
|
||||||
let _ = SetDisks::release_quota_mutation_fences(
|
let _ = SetDisks::release_quota_mutation_fences(
|
||||||
&commit_disks,
|
&commit_disks,
|
||||||
"a_fence_tokens,
|
"a_fence_tokens,
|
||||||
@@ -3371,6 +3283,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
|||||||
Ok(result) => result,
|
Ok(result) => result,
|
||||||
Err(err) => return Err(err.into()),
|
Err(err) => return Err(err.into()),
|
||||||
};
|
};
|
||||||
|
let needs_immediate_heal = rename_commit.needs_immediate_heal();
|
||||||
let op_old_dir = rename_commit.data_dir;
|
let op_old_dir = rename_commit.data_dir;
|
||||||
let cleanup_disks = rename_commit.cleanup_disks;
|
let cleanup_disks = rename_commit.cleanup_disks;
|
||||||
let committed_file_info = rename_commit.committed_file_info;
|
let committed_file_info = rename_commit.committed_file_info;
|
||||||
@@ -3413,9 +3326,6 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
|||||||
// parts are swept by a retried completion or upload GC (rustfs/backlog#946).
|
// parts are swept by a retried completion or upload GC (rustfs/backlog#946).
|
||||||
// Compiles to a no-op outside `#[cfg(test)]`.
|
// Compiles to a no-op outside `#[cfg(test)]`.
|
||||||
if crash_inject::should_crash_at(CrashPoint::MultipartAfterCommitBeforePartsCleanup, &commit_object) {
|
if crash_inject::should_crash_at(CrashPoint::MultipartAfterCommitBeforePartsCleanup, &commit_object) {
|
||||||
if let Some(release) = rename_guard_release.take() {
|
|
||||||
let _ = release.send(false);
|
|
||||||
}
|
|
||||||
return Err(StorageError::Unexpected);
|
return Err(StorageError::Unexpected);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -3431,10 +3341,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
|||||||
.invalidate_get_object_metadata_cache(&commit_bucket, &commit_object)
|
.invalidate_get_object_metadata_cache(&commit_bucket, &commit_object)
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
if let Some(release) = rename_guard_release.take() {
|
drop(_object_lock_guard.take()); // release the object lock before multipart cleanup IO.
|
||||||
let _ = release.send(true);
|
|
||||||
}
|
|
||||||
drop(_object_lock_guard.take()); // release the object lock before multipart cleanup tail IO.
|
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
pause_multipart_commit(&commit_bucket, &commit_object, MultipartCommitPause::AfterObjectPublication).await;
|
pause_multipart_commit(&commit_bucket, &commit_object, MultipartCommitPause::AfterObjectPublication).await;
|
||||||
@@ -3447,9 +3354,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
|||||||
// parts; deleting them before the commit would strand the upload
|
// parts; deleting them before the commit would strand the upload
|
||||||
// permanently. This mirrors the "clean up only after commit" pattern
|
// permanently. This mirrors the "clean up only after commit" pattern
|
||||||
// already used for the old data-dir GC and the upload-dir delete_all below.
|
// already used for the old data-dir GC and the upload-dir delete_all below.
|
||||||
if !tail_owns_staging_cleanup {
|
commit_set.cleanup_multipart_path(&parts).await;
|
||||||
commit_set.cleanup_multipart_path(&parts).await;
|
|
||||||
}
|
|
||||||
|
|
||||||
if let Some(old_dir) = op_old_dir {
|
if let Some(old_dir) = op_old_dir {
|
||||||
// backlog#898: best-effort reclaim of the dereferenced old data dir.
|
// backlog#898: best-effort reclaim of the dereferenced old data dir.
|
||||||
@@ -3480,10 +3385,9 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
pause_multipart_commit(&commit_bucket, &commit_object, MultipartCommitPause::AfterRename).await;
|
pause_multipart_commit(&commit_bucket, &commit_object, MultipartCommitPause::AfterRename).await;
|
||||||
|
|
||||||
if !tail_owns_staging_cleanup
|
if let Err(err) = commit_set
|
||||||
&& let Err(err) = commit_set
|
.delete_all_with_quorum(RUSTFS_META_MULTIPART_BUCKET, &commit_upload_id_path, write_quorum)
|
||||||
.delete_all_with_quorum(RUSTFS_META_MULTIPART_BUCKET, &commit_upload_id_path, write_quorum)
|
.await
|
||||||
.await
|
|
||||||
{
|
{
|
||||||
warn!(
|
warn!(
|
||||||
bucket = %commit_bucket,
|
bucket = %commit_bucket,
|
||||||
@@ -4010,7 +3914,7 @@ mod tests {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial(capacity_dirty_scope)]
|
#[serial(capacity_dirty_scope)]
|
||||||
async fn early_ack_multipart_holds_quota_fences_and_re_marks_capacity_after_tail_drain() {
|
async fn complete_multipart_waits_for_tail_before_releasing_guards_and_marking_capacity() {
|
||||||
use rustfs_object_capacity::capacity_scope::drain_global_dirty_scopes;
|
use rustfs_object_capacity::capacity_scope::drain_global_dirty_scopes;
|
||||||
|
|
||||||
temp_env::async_with_vars([(ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, Some("true"))], async {
|
temp_env::async_with_vars([(ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, Some("true"))], async {
|
||||||
@@ -4056,10 +3960,9 @@ mod tests {
|
|||||||
.collect::<HashSet<_>>();
|
.collect::<HashSet<_>>();
|
||||||
let _ = drain_global_dirty_scopes();
|
let _ = drain_global_dirty_scopes();
|
||||||
|
|
||||||
let rename_tasks = rename_fanout_barrier::observe_tasks(object);
|
|
||||||
let rename_barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME);
|
let rename_barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME);
|
||||||
let complete_store = Arc::clone(&set_disks);
|
let complete_store = Arc::clone(&set_disks);
|
||||||
let complete = tokio::spawn(async move {
|
let mut complete = tokio::spawn(async move {
|
||||||
let mut opts = ObjectOptions::default();
|
let mut opts = ObjectOptions::default();
|
||||||
assert!(opts.set_quota_admission(0, u64::MAX));
|
assert!(opts.set_quota_admission(0, u64::MAX));
|
||||||
complete_store
|
complete_store
|
||||||
@@ -4069,20 +3972,15 @@ mod tests {
|
|||||||
tokio::time::timeout(Duration::from_secs(30), rename_barrier.wait_until_paused())
|
tokio::time::timeout(Duration::from_secs(30), rename_barrier.wait_until_paused())
|
||||||
.await
|
.await
|
||||||
.expect("multipart completion should pause one tail disk during rename");
|
.expect("multipart completion should pause one tail disk during rename");
|
||||||
let cleanup_barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_CLEANUP);
|
|
||||||
complete
|
|
||||||
.await
|
|
||||||
.expect("early-ACK multipart task should join before tail release")
|
|
||||||
.expect("multipart completion should return after write quorum");
|
|
||||||
assert!(
|
assert!(
|
||||||
rename_tasks.running() >= 1,
|
tokio::time::timeout(Duration::from_millis(100), &mut complete).await.is_err(),
|
||||||
"the paused multipart tail disk must remain in flight after quorum ACK"
|
"multipart completion must not publish success while a tail rename is still paused"
|
||||||
);
|
);
|
||||||
|
|
||||||
let initial = drain_global_dirty_scopes().into_iter().collect::<HashSet<_>>();
|
let initial = drain_global_dirty_scopes().into_iter().collect::<HashSet<_>>();
|
||||||
assert!(
|
assert!(
|
||||||
expected.is_subset(&initial),
|
initial.is_empty(),
|
||||||
"the multipart quorum ACK must mark every candidate disk dirty"
|
"capacity must not be marked as committed before the full multipart rename finishes"
|
||||||
);
|
);
|
||||||
|
|
||||||
let abort_store = Arc::clone(&set_disks);
|
let abort_store = Arc::clone(&set_disks);
|
||||||
@@ -4092,7 +3990,7 @@ mod tests {
|
|||||||
.await
|
.await
|
||||||
});
|
});
|
||||||
signaling.wait_for_attempts(2).await;
|
signaling.wait_for_attempts(2).await;
|
||||||
assert!(!abort.is_finished(), "the detached tail owner must retain the multipart upload guard");
|
assert!(!abort.is_finished(), "the in-flight completion must retain the multipart upload guard");
|
||||||
|
|
||||||
let retained_staging = futures::future::join_all(
|
let retained_staging = futures::future::join_all(
|
||||||
disk_stores
|
disk_stores
|
||||||
@@ -4117,25 +4015,20 @@ mod tests {
|
|||||||
.await
|
.await
|
||||||
});
|
});
|
||||||
signaling.wait_for_attempts(object_attempt).await;
|
signaling.wait_for_attempts(object_attempt).await;
|
||||||
assert!(!object_probe.is_finished(), "the detached tail owner must retain the object guard");
|
assert!(!object_probe.is_finished(), "the in-flight completion must retain the object guard");
|
||||||
|
|
||||||
rename_barrier.release();
|
rename_barrier.release();
|
||||||
|
complete
|
||||||
|
.await
|
||||||
|
.expect("multipart task should join after tail release")
|
||||||
|
.expect("multipart completion should return after every tail rename finishes");
|
||||||
object_probe
|
object_probe
|
||||||
.await
|
.await
|
||||||
.expect("object guard probe should join after the tail releases")
|
.expect("object guard probe should join after completion releases")
|
||||||
.expect("object guard probe should acquire after the tail releases");
|
.expect("object guard probe should acquire after completion releases");
|
||||||
tokio::time::timeout(Duration::from_secs(30), cleanup_barrier.wait_until_paused())
|
|
||||||
.await
|
|
||||||
.expect("the multipart tail should pause before reclaiming its old body");
|
|
||||||
let after_tail = drain_global_dirty_scopes().into_iter().collect::<HashSet<_>>();
|
|
||||||
assert!(
|
|
||||||
expected.is_subset(&after_tail),
|
|
||||||
"the multipart rename tail must re-mark capacity after the first scope was drained"
|
|
||||||
);
|
|
||||||
cleanup_barrier.release();
|
|
||||||
let abort_err = abort
|
let abort_err = abort
|
||||||
.await
|
.await
|
||||||
.expect("abort task should join after the tail releases")
|
.expect("abort task should join after completion releases")
|
||||||
.expect_err("the committed upload should no longer exist");
|
.expect_err("the committed upload should no longer exist");
|
||||||
assert!(matches!(abort_err, StorageError::InvalidUploadID(..)));
|
assert!(matches!(abort_err, StorageError::InvalidUploadID(..)));
|
||||||
|
|
||||||
@@ -4148,7 +4041,7 @@ mod tests {
|
|||||||
let after_cleanup = drain_global_dirty_scopes().into_iter().collect::<HashSet<_>>();
|
let after_cleanup = drain_global_dirty_scopes().into_iter().collect::<HashSet<_>>();
|
||||||
assert!(
|
assert!(
|
||||||
expected.is_subset(&after_cleanup),
|
expected.is_subset(&after_cleanup),
|
||||||
"the multipart tail cleanup must re-mark capacity after its preceding scope was drained"
|
"the completed multipart commit must mark every candidate disk dirty"
|
||||||
);
|
);
|
||||||
})
|
})
|
||||||
.await;
|
.await;
|
||||||
@@ -4283,23 +4176,26 @@ mod tests {
|
|||||||
],
|
],
|
||||||
async {
|
async {
|
||||||
let rename_barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME);
|
let rename_barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME);
|
||||||
set_disks
|
let complete_store = Arc::clone(&set_disks);
|
||||||
.clone()
|
let mut complete = tokio::spawn(async move {
|
||||||
.complete_multipart_upload(bucket, object, &upload_id, parts.clone(), &ObjectOptions::default())
|
complete_store
|
||||||
.await
|
.complete_multipart_upload(bucket, object, &upload_id, parts.clone(), &ObjectOptions::default())
|
||||||
.expect("fenced multipart completion should commit with a live proof");
|
.await
|
||||||
|
});
|
||||||
|
|
||||||
tokio::time::timeout(Duration::from_secs(30), rename_barrier.wait_until_paused())
|
tokio::time::timeout(Duration::from_secs(30), rename_barrier.wait_until_paused())
|
||||||
.await
|
.await
|
||||||
.expect("multipart completion should leave one rename tail in flight after quorum ACK");
|
.expect("multipart completion should pause one tail disk during rename");
|
||||||
let disks = disk_stores.clone();
|
|
||||||
let mut epochs = tokio::spawn(async move { object_transaction_epochs(&disks, bucket, object).await });
|
|
||||||
assert!(
|
assert!(
|
||||||
tokio::time::timeout(Duration::from_millis(100), &mut epochs).await.is_err(),
|
tokio::time::timeout(Duration::from_millis(100), &mut complete).await.is_err(),
|
||||||
"epoch read-back should wait for the lagging rename tail"
|
"fenced multipart completion must wait for every rename tail before returning"
|
||||||
);
|
);
|
||||||
rename_barrier.release();
|
rename_barrier.release();
|
||||||
epochs.await.expect("epoch read-back should finish after the rename tail")
|
complete
|
||||||
|
.await
|
||||||
|
.expect("fenced multipart task should join after tail release")
|
||||||
|
.expect("fenced multipart completion should commit with a live proof");
|
||||||
|
object_transaction_epochs(&disk_stores, bucket, object).await
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
@@ -8392,29 +8288,18 @@ mod tests {
|
|||||||
let new = payload(0xC3);
|
let new = payload(0xC3);
|
||||||
let (u_new, parts_new) = stage_upload(&set_disks, bucket, object, &new).await;
|
let (u_new, parts_new) = stage_upload(&set_disks, bucket, object, &new).await;
|
||||||
let parts_retry = parts_new.clone();
|
let parts_retry = parts_new.clone();
|
||||||
let rename_tasks = rename_fanout_barrier::observe_tasks(object);
|
|
||||||
let rename_barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME);
|
|
||||||
crash_inject::arm(CrashPoint::MultipartAfterCommitBeforePartsCleanup, object);
|
crash_inject::arm(CrashPoint::MultipartAfterCommitBeforePartsCleanup, object);
|
||||||
let crashed = complete(&set_disks, bucket, object, &u_new, parts_new).await;
|
let crashed = complete(&set_disks, bucket, object, &u_new, parts_new).await;
|
||||||
assert!(
|
assert!(
|
||||||
matches!(crashed, Err(StorageError::Unexpected)),
|
matches!(crashed, Err(StorageError::Unexpected)),
|
||||||
"the armed post-commit crash point must be the failure that surfaced, got {crashed:?}"
|
"the armed post-commit crash point must be the failure that surfaced, got {crashed:?}"
|
||||||
);
|
);
|
||||||
assert!(rename_tasks.running() >= 1, "the crash must interrupt an actual early-ACK tail handoff");
|
|
||||||
crash_inject::disarm(CrashPoint::MultipartAfterCommitBeforePartsCleanup, object);
|
crash_inject::disarm(CrashPoint::MultipartAfterCommitBeforePartsCleanup, object);
|
||||||
rename_barrier.release();
|
|
||||||
tokio::time::timeout(Duration::from_secs(30), async {
|
|
||||||
while rename_tasks.running() != 0 {
|
|
||||||
tokio::task::yield_now().await;
|
|
||||||
}
|
|
||||||
})
|
|
||||||
.await
|
|
||||||
.expect("the crash-interrupted rename tail should drain after release");
|
|
||||||
drop(
|
drop(
|
||||||
set_disks
|
set_disks
|
||||||
.acquire_write_lock_diag("post_commit_crash_tail_probe", bucket, object)
|
.acquire_write_lock_diag("post_commit_crash_tail_probe", bucket, object)
|
||||||
.await
|
.await
|
||||||
.expect("the crash-interrupted tail should release its object guard"),
|
.expect("the failed post-commit completion should release its object guard"),
|
||||||
);
|
);
|
||||||
|
|
||||||
// The commit landed: the new version reads back whole and correct.
|
// The commit landed: the new version reads back whole and correct.
|
||||||
@@ -8487,29 +8372,18 @@ mod tests {
|
|||||||
|
|
||||||
let new = payload(0x52);
|
let new = payload(0x52);
|
||||||
let (u_new, parts_new) = stage_upload(&set_disks, bucket, object, &new).await;
|
let (u_new, parts_new) = stage_upload(&set_disks, bucket, object, &new).await;
|
||||||
let rename_tasks = rename_fanout_barrier::observe_tasks(object);
|
|
||||||
let rename_barrier = rename_fanout_barrier::arm(object, 0, rename_fanout_barrier::PHASE_RENAME);
|
|
||||||
crash_inject::arm(CrashPoint::MultipartAfterCommitBeforePartsCleanup, object);
|
crash_inject::arm(CrashPoint::MultipartAfterCommitBeforePartsCleanup, object);
|
||||||
let crashed = complete(&set_disks, bucket, object, &u_new, parts_new).await;
|
let crashed = complete(&set_disks, bucket, object, &u_new, parts_new).await;
|
||||||
assert!(
|
assert!(
|
||||||
matches!(crashed, Err(StorageError::Unexpected)),
|
matches!(crashed, Err(StorageError::Unexpected)),
|
||||||
"the post-commit crash point must surface as unexpected, got {crashed:?}"
|
"the post-commit crash point must surface as unexpected, got {crashed:?}"
|
||||||
);
|
);
|
||||||
assert!(rename_tasks.running() >= 1, "the crash must interrupt an actual early-ACK tail handoff");
|
|
||||||
crash_inject::disarm(CrashPoint::MultipartAfterCommitBeforePartsCleanup, object);
|
crash_inject::disarm(CrashPoint::MultipartAfterCommitBeforePartsCleanup, object);
|
||||||
rename_barrier.release();
|
|
||||||
tokio::time::timeout(Duration::from_secs(30), async {
|
|
||||||
while rename_tasks.running() != 0 {
|
|
||||||
tokio::task::yield_now().await;
|
|
||||||
}
|
|
||||||
})
|
|
||||||
.await
|
|
||||||
.expect("the crash-interrupted rename tail should drain after release");
|
|
||||||
drop(
|
drop(
|
||||||
set_disks
|
set_disks
|
||||||
.acquire_write_lock_diag("post_commit_receipt_tail_probe", bucket, object)
|
.acquire_write_lock_diag("post_commit_receipt_tail_probe", bucket, object)
|
||||||
.await
|
.await
|
||||||
.expect("the crash-interrupted tail should release its object guard"),
|
.expect("the failed post-commit completion should release its object guard"),
|
||||||
);
|
);
|
||||||
|
|
||||||
let (body, _) = read_object(&set_disks, bucket, object).await;
|
let (body, _) = read_object(&set_disks, bucket, object).await;
|
||||||
@@ -8523,8 +8397,8 @@ mod tests {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
receipts, 3,
|
receipts, 4,
|
||||||
"the committed quorum must persist receipts while the crash-interrupted tail preserves staging"
|
"the completed rename must persist old-data cleanup receipts on every disk before surfacing the post-commit crash"
|
||||||
);
|
);
|
||||||
|
|
||||||
let restarted_endpoints = temp_dirs
|
let restarted_endpoints = temp_dirs
|
||||||
@@ -8570,12 +8444,12 @@ mod tests {
|
|||||||
.reconcile_old_data_cleanup_receipts(bucket, object)
|
.reconcile_old_data_cleanup_receipts(bucket, object)
|
||||||
.await
|
.await
|
||||||
.expect("restart receipt reconciliation should succeed");
|
.expect("restart receipt reconciliation should succeed");
|
||||||
assert_eq!(removed, 3, "restart receipt reconciliation should delete the committed quorum's targets");
|
assert_eq!(removed, 4, "restart receipt reconciliation should delete every committed target");
|
||||||
let reclaimed = restarted_set
|
let reclaimed = restarted_set
|
||||||
.reclaim_orphan_data_dirs(bucket, object)
|
.reclaim_orphan_data_dirs(bucket, object)
|
||||||
.await
|
.await
|
||||||
.expect("restart orphan reconciliation should succeed");
|
.expect("restart orphan reconciliation should succeed");
|
||||||
assert_eq!(reclaimed, 1, "the late commit without a receipt must remain reclaimable as an orphan");
|
assert_eq!(reclaimed, 0, "the post-commit crash should leave no receipt-less late commit orphan");
|
||||||
for disk in &reloaded {
|
for disk in &reloaded {
|
||||||
assert!(
|
assert!(
|
||||||
!data_dir_exists(disk, bucket, object, old_dir).await,
|
!data_dir_exists(disk, bucket, object, old_dir).await,
|
||||||
|
|||||||
@@ -2305,6 +2305,200 @@ mod tests {
|
|||||||
shutdown.cancel();
|
shutdown.cancel();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||||
|
#[serial_test::serial(storage_class_env)]
|
||||||
|
async fn copy_object_immediately_reads_small_completed_multipart_source() {
|
||||||
|
let temp_dir = tempfile::tempdir().expect("create small multipart copy store dir");
|
||||||
|
let (_ctx, store, shutdown) =
|
||||||
|
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "small-multipart-copy", &[1])).await;
|
||||||
|
crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await;
|
||||||
|
|
||||||
|
let bucket = format!("small-multipart-copy-{}", Uuid::new_v4());
|
||||||
|
let source_object = "docker/registry/v2/repositories/example/_uploads/upload-id/data";
|
||||||
|
let target_object = "docker/registry/v2/blobs/sha256/c0/digest/data";
|
||||||
|
let payload = vec![0xAB; 273];
|
||||||
|
|
||||||
|
store
|
||||||
|
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("create bucket for small multipart copy");
|
||||||
|
let upload = store
|
||||||
|
.new_multipart_upload(&bucket, source_object, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("create source multipart upload");
|
||||||
|
let mut part_reader = PutObjReader::from_vec(payload.clone());
|
||||||
|
let part = store
|
||||||
|
.put_object_part(&bucket, source_object, &upload.upload_id, 1, &mut part_reader, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("stage small multipart source part");
|
||||||
|
let completed = store
|
||||||
|
.clone()
|
||||||
|
.complete_multipart_upload(
|
||||||
|
&bucket,
|
||||||
|
source_object,
|
||||||
|
&upload.upload_id,
|
||||||
|
vec![crate::storage_api_contracts::multipart::CompletePart {
|
||||||
|
part_num: part.part_num,
|
||||||
|
etag: part.etag,
|
||||||
|
..Default::default()
|
||||||
|
}],
|
||||||
|
&ObjectOptions::default(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("complete the small multipart source");
|
||||||
|
assert_eq!(completed.get_actual_size().expect("completed object logical size"), payload.len() as i64);
|
||||||
|
|
||||||
|
let source_reader = store
|
||||||
|
.get_object_reader(&bucket, source_object, None, HeaderMap::new(), &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("completed multipart source should be immediately readable");
|
||||||
|
let mut copy_info = source_reader.object_info.clone();
|
||||||
|
let actual_size = copy_info.get_actual_size().expect("copy source logical size should resolve");
|
||||||
|
assert_eq!(actual_size, payload.len() as i64);
|
||||||
|
let copy_reader = rustfs_rio::HashReader::from_stream(source_reader.stream, actual_size, actual_size, None, None, false)
|
||||||
|
.expect("copy source hash reader should build");
|
||||||
|
copy_info.put_object_reader = Some(PutObjReader::new(copy_reader));
|
||||||
|
|
||||||
|
store
|
||||||
|
.copy_object(
|
||||||
|
&bucket,
|
||||||
|
source_object,
|
||||||
|
&bucket,
|
||||||
|
target_object,
|
||||||
|
&mut copy_info,
|
||||||
|
&ObjectOptions::default(),
|
||||||
|
&ObjectOptions::default(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("CopyObject should accept a freshly completed multipart source");
|
||||||
|
|
||||||
|
let mut target_reader = store
|
||||||
|
.get_object_reader(&bucket, target_object, None, HeaderMap::new(), &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("copied target should be readable");
|
||||||
|
let mut target_body = Vec::new();
|
||||||
|
target_reader
|
||||||
|
.stream
|
||||||
|
.read_to_end(&mut target_body)
|
||||||
|
.await
|
||||||
|
.expect("target body should stream");
|
||||||
|
assert_eq!(target_body, payload);
|
||||||
|
shutdown.cancel();
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(feature = "test-util")]
|
||||||
|
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||||
|
#[serial_test::serial(storage_class_env)]
|
||||||
|
async fn complete_multipart_waits_for_tail_rename_before_copy_source_visibility() {
|
||||||
|
let temp_dir = tempfile::tempdir().expect("create early-ack multipart copy store dir");
|
||||||
|
let (_ctx, store, shutdown) =
|
||||||
|
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "early-ack-multipart-copy", &[4])).await;
|
||||||
|
crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await;
|
||||||
|
|
||||||
|
let bucket = format!("early-ack-multipart-copy-{}", Uuid::new_v4());
|
||||||
|
let source_object = "docker/registry/v2/repositories/example/_uploads/upload-id/data";
|
||||||
|
let target_object = "docker/registry/v2/blobs/sha256/c0/digest/data";
|
||||||
|
let payload = vec![0xCD; 273];
|
||||||
|
|
||||||
|
store
|
||||||
|
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("create bucket for early-ack multipart copy");
|
||||||
|
let upload = store
|
||||||
|
.new_multipart_upload(&bucket, source_object, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("create source multipart upload");
|
||||||
|
let mut part_reader = PutObjReader::from_vec(payload.clone());
|
||||||
|
let part = store
|
||||||
|
.put_object_part(&bucket, source_object, &upload.upload_id, 1, &mut part_reader, &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("stage small multipart source part");
|
||||||
|
let completed_parts = vec![crate::storage_api_contracts::multipart::CompletePart {
|
||||||
|
part_num: part.part_num,
|
||||||
|
etag: part.etag,
|
||||||
|
..Default::default()
|
||||||
|
}];
|
||||||
|
|
||||||
|
temp_env::async_with_vars([(crate::set_disk::ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, Some("true"))], async {
|
||||||
|
let rename_tasks = crate::set_disk::rename_fanout_barrier::observe_tasks(source_object);
|
||||||
|
let rename_barrier = crate::set_disk::rename_fanout_barrier::arm(
|
||||||
|
source_object,
|
||||||
|
0,
|
||||||
|
crate::set_disk::rename_fanout_barrier::PHASE_RENAME,
|
||||||
|
);
|
||||||
|
let complete_store = Arc::clone(&store);
|
||||||
|
let complete_bucket = bucket.clone();
|
||||||
|
let complete_upload_id = upload.upload_id.clone();
|
||||||
|
let mut complete = tokio::spawn(async move {
|
||||||
|
complete_store
|
||||||
|
.complete_multipart_upload(
|
||||||
|
&complete_bucket,
|
||||||
|
source_object,
|
||||||
|
&complete_upload_id,
|
||||||
|
completed_parts,
|
||||||
|
&ObjectOptions::default(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
});
|
||||||
|
tokio::time::timeout(Duration::from_secs(30), rename_barrier.wait_until_paused())
|
||||||
|
.await
|
||||||
|
.expect("multipart completion should pause one tail disk during rename");
|
||||||
|
assert!(
|
||||||
|
rename_tasks.running() >= 1,
|
||||||
|
"the paused multipart tail disk must remain in flight before completion returns"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
tokio::time::timeout(Duration::from_millis(100), &mut complete).await.is_err(),
|
||||||
|
"CompleteMultipartUpload must not return while a copy source rename tail is still pending"
|
||||||
|
);
|
||||||
|
rename_barrier.release();
|
||||||
|
complete
|
||||||
|
.await
|
||||||
|
.expect("multipart completion task should join")
|
||||||
|
.expect("multipart completion should return after every rename tail finishes");
|
||||||
|
|
||||||
|
let source_reader = store
|
||||||
|
.get_object_reader(&bucket, source_object, None, HeaderMap::new(), &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("completed multipart source should be immediately readable after success");
|
||||||
|
let mut copy_info = source_reader.object_info.clone();
|
||||||
|
let actual_size = copy_info.get_actual_size().expect("copy source logical size should resolve");
|
||||||
|
assert_eq!(actual_size, payload.len() as i64);
|
||||||
|
let copy_reader =
|
||||||
|
rustfs_rio::HashReader::from_stream(source_reader.stream, actual_size, actual_size, None, None, false)
|
||||||
|
.expect("copy source hash reader should build");
|
||||||
|
copy_info.put_object_reader = Some(PutObjReader::new(copy_reader));
|
||||||
|
|
||||||
|
store
|
||||||
|
.copy_object(
|
||||||
|
&bucket,
|
||||||
|
source_object,
|
||||||
|
&bucket,
|
||||||
|
target_object,
|
||||||
|
&mut copy_info,
|
||||||
|
&ObjectOptions::default(),
|
||||||
|
&ObjectOptions::default(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("CopyObject should accept a freshly completed multipart source");
|
||||||
|
})
|
||||||
|
.await;
|
||||||
|
|
||||||
|
let mut target_reader = store
|
||||||
|
.get_object_reader(&bucket, target_object, None, HeaderMap::new(), &ObjectOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("copied target should be readable after tail release");
|
||||||
|
let mut target_body = Vec::new();
|
||||||
|
target_reader
|
||||||
|
.stream
|
||||||
|
.read_to_end(&mut target_body)
|
||||||
|
.await
|
||||||
|
.expect("target body should stream");
|
||||||
|
assert_eq!(target_body, payload);
|
||||||
|
shutdown.cancel();
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(feature = "test-util")]
|
#[cfg(feature = "test-util")]
|
||||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||||
#[serial_test::serial(storage_class_env)]
|
#[serial_test::serial(storage_class_env)]
|
||||||
|
|||||||
Reference in New Issue
Block a user