fix(s3): make retried CompleteMultipartUpload idempotent (#8153)

Record the upload id on the completed object so a lost-response retry
returns that object's ETag instead of NoSuchUpload, while a different
part list or a replaced object keeps the existing errors.

Signed-off-by: loverustfs <155562731+loverustfs@users.noreply.github.com>
Co-authored-by: Hauser <housemecn@gmail.com>
This commit is contained in:
RustFS
2026-09-28 19:08:10 +08:00
committed by GitHub
parent 64bd224b3b
commit fc5609bbb0
10 changed files with 590 additions and 110 deletions
+1 -1
View File
@@ -448,7 +448,7 @@ pub mod disk {
pub mod error {
pub use crate::error::{
Error, PoolMetadataError, PoolMetadataFailure, Result, StorageError, classify_system_path_failure_reason,
is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found,
is_err_bucket_not_found, is_err_invalid_upload_id, is_err_object_not_found, is_err_version_not_found,
};
}
+1
View File
@@ -3154,6 +3154,7 @@ mod tests {
version_purge_status: Default::default(),
replication_decision: String::new(),
checksum: None,
multipart_completion_replayed: false,
}
}
}
+5
View File
@@ -1444,6 +1444,10 @@ pub struct ObjectInfo {
pub version_purge_status: VersionPurgeStatusType,
pub replication_decision: String,
pub checksum: Option<Bytes>,
/// True when this `CompleteMultipartUpload` result is the object already
/// published by the same upload. Callers must not repeat quota, replication,
/// lifecycle, or object-created side effects for that response.
pub multipart_completion_replayed: bool,
}
impl Clone for ObjectInfo {
@@ -1484,6 +1488,7 @@ impl Clone for ObjectInfo {
version_purge_status: self.version_purge_status.clone(),
replication_decision: self.replication_decision.clone(),
checksum: self.checksum.clone(),
multipart_completion_replayed: self.multipart_completion_replayed,
expires: self.expires,
}
}
+407 -12
View File
@@ -1451,6 +1451,103 @@ impl SetDisks {
delimiter: delimiter.to_owned(),
})
}
/// Answer a CompleteMultipartUpload whose staging upload is already gone.
///
/// AWS keeps Complete idempotent for the same upload id and part list while
/// the completed object is still the result of that upload (the first 200
/// can be lost, and SDK retries depend on a second 200 with the same ETag).
/// Returns `Ok(None)` when this key has no such object, so the caller keeps
/// `InvalidUploadID`. A matching upload id with a different part list is
/// `InvalidPart` and must not fall through to `NoSuchUpload`.
async fn replay_completed_multipart_upload(
&self,
bucket: &str,
object: &str,
upload_id: &str,
uploaded_parts: &[CompletePart],
opts: &ObjectOptions,
) -> Result<Option<ObjectInfo>> {
if upload_id.is_empty() {
return Ok(None);
}
let read_opts = ObjectOptions {
no_lock: true,
metadata_cache_safe: false,
versioned: opts.versioned,
version_suspended: opts.version_suspended,
..Default::default()
};
match self.get_object_info(bucket, object, &read_opts).await {
Ok(info) => {
match completed_multipart_upload_matches(
&info.user_defined,
info.parts.as_ref(),
info.delete_marker,
upload_id,
uploaded_parts,
bucket,
object,
) {
Ok(true) => return Ok(Some(mark_multipart_completion_replayed(info))),
Ok(false) => {}
Err(err) => return Err(err),
}
}
Err(err) if is_err_object_not_found(&err) || is_err_version_not_found(&err) => {}
Err(err) => return Err(err),
}
// Unversioned overwrite replaces the only slot. A versioned bucket can
// keep the completed version under a newer one; that version is still
// the result of this upload until it is deleted.
if !(opts.versioned || opts.version_suspended) {
return Ok(None);
}
let Some(versions) = self.load_file_info_versions_exact(bucket, object).await? else {
return Ok(None);
};
for fi in versions.versions {
if fi.deleted || fi.tier_free_version() || fi.is_canonical_delete_marker() {
continue;
}
match completed_multipart_upload_matches(&fi.metadata, &fi.parts, false, upload_id, uploaded_parts, bucket, object) {
Ok(true) => {
let info = ObjectInfo::from_file_info(&fi, bucket, object, true);
return Ok(Some(mark_multipart_completion_replayed(info)));
}
Ok(false) => {}
Err(err) => return Err(err),
}
}
Ok(None)
}
/// Best-effort removal of staging left behind when commit succeeded and the
/// process stopped before `delete_all`. A cleanup miss must not turn the
/// idempotent 200 into an error; the object is already durable.
async fn reclaim_replayed_multipart_staging(&self, bucket: &str, object: &str, upload_id: &str, upload_id_path: &str) {
if let Err(err) = self
.delete_all_with_quorum(RUSTFS_META_MULTIPART_BUCKET, upload_id_path, self.default_write_quorum())
.await
{
warn!(
target: "rustfs_ecstore::set_disk",
event = EVENT_SET_DISK_MULTIPART,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_SET_DISK,
op = "complete_multipart_upload",
result = "cleanup_quorum_missed",
bucket = %bucket,
object = %object,
upload_id = %upload_id,
error = %err,
"replayed multipart completion left staging behind"
);
}
}
}
#[async_trait::async_trait]
@@ -2337,6 +2434,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
}
result
}
// complete_multipart_upload finished
#[tracing::instrument(skip(self))]
async fn complete_multipart_upload(
@@ -2400,9 +2498,27 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
.await?;
let expected_restore_operation_id = restore_commit_operation_id_from_metadata(&opts.user_defined)?;
let (mut fi, files_metas) = self
let (mut fi, files_metas) = match self
.check_upload_id_exists_with_opts(bucket, object, upload_id, true, opts)
.await?;
.await
{
Ok(found) => found,
Err(err) if crate::error::is_err_invalid_upload_id(&err) => {
match self
.replay_completed_multipart_upload(bucket, object, upload_id, &uploaded_parts, opts)
.await
{
Ok(Some(existing)) => {
self.reclaim_replayed_multipart_staging(bucket, object, upload_id, &upload_id_path)
.await;
return Ok(existing);
}
Ok(None) => return Err(err),
Err(replay_err) => return Err(replay_err),
}
}
Err(err) => return Err(err),
};
ensure_data_movement_upload_access(&fi, bucket, object, upload_id, opts)?;
ensure_multipart_bucket_incarnation(&self.ctx, &fi, bucket, object, upload_id, opts.expected_bucket_incarnation_id)
.await?;
@@ -2988,6 +3104,14 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
rustfs_utils::http::remove_str(&mut fi.metadata, rustfs_filemeta::shard_integrity::SUFFIX_UPLOAD_INTEGRITY);
fi.persist_shard_integrity()?;
// The staging directory is removed after commit. Recording the upload id
// on the object is what lets a retried CompleteMultipartUpload return this
// version instead of NoSuchUpload. insert_str writes both internal prefixes;
// a later conflicting pair fails closed in the replay matcher.
if !upload_id.is_empty() {
insert_str(&mut fi.metadata, rustfs_utils::http::SUFFIX_MULTIPART_UPLOAD_ID, upload_id.to_owned());
}
for meta in parts_metadatas.iter_mut() {
if meta.has_valid_erasure_geometry() {
meta.size = fi.size;
@@ -3522,6 +3646,58 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
}
}
fn mark_multipart_completion_replayed(mut info: ObjectInfo) -> ObjectInfo {
info.multipart_completion_replayed = true;
info
}
/// `Ok(true)` when `upload_id` is the completion that published this version and
/// the requested parts are that version's part list. `Ok(false)` means this
/// version is not that completion. `Err` is a definite part-list or metadata
/// failure and must not be reported as `NoSuchUpload`.
fn completed_multipart_upload_matches(
metadata: &HashMap<String, String>,
parts: &[ObjectPartInfo],
delete_marker: bool,
upload_id: &str,
uploaded_parts: &[CompletePart],
bucket: &str,
object: &str,
) -> Result<bool> {
if delete_marker || upload_id.is_empty() {
return Ok(false);
}
let Some(stored_upload_id) = rustfs_utils::http::get_consistent_str(metadata, rustfs_utils::http::SUFFIX_MULTIPART_UPLOAD_ID)
else {
if rustfs_utils::http::contains_key_str(metadata, rustfs_utils::http::SUFFIX_MULTIPART_UPLOAD_ID) {
return Err(Error::FileCorrupt);
}
return Ok(false);
};
if stored_upload_id.is_empty() || stored_upload_id != upload_id {
return Ok(false);
}
if parts.len() != uploaded_parts.len() {
let part_num = uploaded_parts.first().map(|part| part.part_num).unwrap_or(0);
return Err(Error::InvalidPart(part_num, bucket.to_owned(), object.to_owned()));
}
for (stored, requested) in parts.iter().zip(uploaded_parts) {
if stored.number != requested.part_num {
return Err(Error::InvalidPart(requested.part_num, bucket.to_owned(), object.to_owned()));
}
let stored_etag = rustfs_utils::path::trim_etag(&stored.etag);
let client_etag = requested.etag.as_deref().map(rustfs_utils::path::trim_etag);
if client_etag.as_deref() != Some(stored_etag.as_str()) {
return Err(Error::InvalidPart(
requested.part_num,
stored.etag.clone(),
requested.etag.clone().unwrap_or_default(),
));
}
}
Ok(true)
}
/// Final ETag for a completed multipart object. An authorized replication
/// request preserves the source ETag so the replication HEAD comparison
/// converges even when the source ETag is not derivable from the uploaded
@@ -8630,6 +8806,219 @@ mod tests {
assert_eq!(paged, vec!["u0", "u1", "u2", "u3"]);
}
#[test]
fn completed_multipart_upload_match_requires_same_upload_and_parts() {
let upload_id = "upload-1";
let mut metadata = HashMap::new();
insert_str(&mut metadata, rustfs_utils::http::SUFFIX_MULTIPART_UPLOAD_ID, upload_id.to_string());
let parts = vec![ObjectPartInfo {
number: 1,
etag: "abc".to_string(),
..Default::default()
}];
let requested = vec![CompletePart {
part_num: 1,
etag: Some("\"abc\"".to_string()),
..Default::default()
}];
assert_eq!(
completed_multipart_upload_matches(&metadata, &parts, false, upload_id, &requested, "bucket", "object")
.expect("quoted ETag must match the stored part"),
true
);
let wrong_etag = vec![CompletePart {
part_num: 1,
etag: Some("def".to_string()),
..Default::default()
}];
assert!(matches!(
completed_multipart_upload_matches(&metadata, &parts, false, upload_id, &wrong_etag, "bucket", "object"),
Err(StorageError::InvalidPart(1, _, _))
));
assert_eq!(
completed_multipart_upload_matches(&metadata, &parts, false, "other-upload", &requested, "bucket", "object")
.expect("a different upload id is not this completion"),
false
);
assert_eq!(
completed_multipart_upload_matches(&metadata, &parts, true, upload_id, &requested, "bucket", "object")
.expect("a delete marker is not the completed object"),
false
);
metadata.insert(
format!(
"{}{}",
rustfs_utils::http::MINIO_INTERNAL_PREFIX,
rustfs_utils::http::SUFFIX_MULTIPART_UPLOAD_ID
),
"disagrees".to_string(),
);
assert!(
matches!(
completed_multipart_upload_matches(&metadata, &parts, false, upload_id, &requested, "bucket", "object"),
Err(StorageError::FileCorrupt)
),
"conflicting internal upload ids must fail closed"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn retried_complete_multipart_upload_returns_the_committed_object() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "multipart-complete-retry";
let object = "object";
make_bucket_on_all(&disk_stores, bucket).await;
let payload = vec![0x11; 4096];
let (upload_id, parts) =
stage_upload_with_create_opts(&set_disks, bucket, object, &payload, &ObjectOptions::default()).await;
let first = set_disks
.clone()
.complete_multipart_upload(bucket, object, &upload_id, parts.clone(), &ObjectOptions::default())
.await
.expect("the first completion should publish the object");
assert!(!first.multipart_completion_replayed, "the first completion publishes a new object");
assert_eq!(
rustfs_utils::http::get_consistent_str(&first.user_defined, rustfs_utils::http::SUFFIX_MULTIPART_UPLOAD_ID),
Some(upload_id.as_str()),
"the completed object must record the upload id"
);
let mut quoted_parts = parts.clone();
quoted_parts[0].etag = quoted_parts[0].etag.as_ref().map(|etag| format!("\"{etag}\""));
let retried = set_disks
.clone()
.complete_multipart_upload(bucket, object, &upload_id, quoted_parts, &ObjectOptions::default())
.await
.expect("a retry with the same upload id and parts must succeed");
assert!(retried.multipart_completion_replayed);
assert_eq!(retried.etag, first.etag);
assert_eq!(retried.size, first.size);
assert_eq!(retried.version_id, first.version_id);
let mismatched = parts.clone();
let mut mismatched = mismatched;
mismatched[0].etag = Some("not-the-part".to_string());
let mismatch_err = set_disks
.clone()
.complete_multipart_upload(bucket, object, &upload_id, mismatched, &ObjectOptions::default())
.await
.expect_err("a retry with a different part ETag must not pretend the upload is missing");
assert!(matches!(mismatch_err, StorageError::InvalidPart(..)));
let unknown = set_disks
.clone()
.complete_multipart_upload(bucket, object, "not-a-real-upload", parts.clone(), &ObjectOptions::default())
.await
.expect_err("an unknown upload id must stay InvalidUploadID");
assert!(matches!(unknown, StorageError::InvalidUploadID(..)));
let mut reader = PutObjReader::from_vec(b"replaced".to_vec());
set_disks
.put_object(bucket, object, &mut reader, &ObjectOptions::default())
.await
.expect("an overwrite should replace the completed object");
let overwritten = set_disks
.clone()
.complete_multipart_upload(bucket, object, &upload_id, parts, &ObjectOptions::default())
.await
.expect_err("a retry after the object is no longer that upload must stay InvalidUploadID");
assert!(matches!(overwritten, StorageError::InvalidUploadID(..)));
let current = set_disks
.get_object_info(bucket, object, &ObjectOptions::default())
.await
.expect("the overwrite must remain readable");
assert_eq!(current.size, b"replaced".len() as i64);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn concurrent_complete_multipart_upload_returns_one_object() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "multipart-complete-concurrent";
let object = "object";
make_bucket_on_all(&disk_stores, bucket).await;
let (upload_id, parts) =
stage_upload_with_create_opts(&set_disks, bucket, object, &[0x22; 4096], &ObjectOptions::default()).await;
let first_set = set_disks.clone();
let second_set = set_disks.clone();
let first_parts = parts.clone();
let second_parts = parts.clone();
let first_upload = upload_id.clone();
let second_upload = upload_id.clone();
let (first, second) = tokio::join!(
async move {
first_set
.complete_multipart_upload(bucket, object, &first_upload, first_parts, &ObjectOptions::default())
.await
},
async move {
second_set
.complete_multipart_upload(bucket, object, &second_upload, second_parts, &ObjectOptions::default())
.await
},
);
let first = first.expect("one completion must succeed");
let second = second.expect("the duplicate completion must succeed");
assert_eq!(first.etag, second.etag, "both completions must report the same ETag");
assert_eq!(first.version_id, second.version_id);
assert_eq!(
usize::from(first.multipart_completion_replayed) + usize::from(second.multipart_completion_replayed),
1
);
let info = set_disks
.get_object_info(bucket, object, &ObjectOptions::default())
.await
.expect("the object must be readable after duplicate completion");
assert_eq!(info.etag, first.etag);
assert_eq!(info.size, 4096);
}
#[tokio::test]
async fn versioned_complete_retry_returns_the_completed_version() {
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "multipart-complete-versioned";
let object = "object";
make_bucket_on_all(&disk_stores, bucket).await;
let versioned = ObjectOptions {
versioned: true,
..Default::default()
};
let (upload_id, parts) = stage_upload_with_create_opts(&set_disks, bucket, object, &[0x33; 4096], &versioned).await;
let completed = set_disks
.clone()
.complete_multipart_upload(bucket, object, &upload_id, parts.clone(), &versioned)
.await
.expect("versioned completion should publish a version");
let completed_version = completed.version_id.expect("versioned completion assigns a version id");
let mut reader = PutObjReader::from_vec(b"newer-version".to_vec());
let successor = set_disks
.put_object(bucket, object, &mut reader, &versioned)
.await
.expect("a later versioned put should add a version");
assert_ne!(successor.version_id, Some(completed_version));
let retried = set_disks
.clone()
.complete_multipart_upload(bucket, object, &upload_id, parts, &versioned)
.await
.expect("the retry must still find the version this upload published");
assert!(retried.multipart_completion_replayed);
assert_eq!(retried.etag, completed.etag);
assert_eq!(retried.version_id, Some(completed_version));
let latest = set_disks
.get_object_info(bucket, object, &versioned)
.await
.expect("the latest version must stay the successor put");
assert_eq!(latest.version_id, successor.version_id);
assert_ne!(latest.etag, completed.etag);
}
/// Crash-consistency for the two `complete_multipart_upload` commit windows.
///
/// rustfs/backlog#864: a fault that interrupts a completion must never mutate
@@ -9299,18 +9688,24 @@ mod tests {
"the post-commit crash must leave the upload listable for reclamation"
);
// A retried CompleteMultipartUpload is answered deterministically: the
// commit rename already consumed the upload's metadata, so the retry
// resolves to InvalidUploadID (NoSuchUpload to the S3 client, the
// standard answer for a completed-then-retried upload) — never a torn
// state, and the committed object is untouched by the retry.
let retried = complete(&set_disks, bucket, object, &u_new, parts_retry).await;
// The commit rename already consumed the upload's xl.meta, so the
// retry cannot read staging. It must still return the committed
// object: AWS CompleteMultipartUpload is idempotent for the same
// upload id and part list, and the retry must not tear that object.
let (body_before_retry, etag_before_retry) = read_object(&set_disks, bucket, object).await;
assert_eq!(body_before_retry, new, "a post-commit crash must leave the whole new version readable");
let retried = complete(&set_disks, bucket, object, &u_new, parts_retry)
.await
.expect("a retried complete after the commit landed must return the committed object");
assert!(retried.multipart_completion_replayed, "the retry must not publish a second object");
assert_eq!(retried.etag, etag_before_retry, "the retry must return the committed ETag");
let (body_after_retry, etag_after_retry) = read_object(&set_disks, bucket, object).await;
assert_eq!(body_after_retry, new, "the retry must not disturb the committed object");
assert_eq!(etag_after_retry, etag_before_retry, "the retry must not rewrite the committed ETag");
assert!(
matches!(retried, Err(StorageError::InvalidUploadID(..))),
"a retried complete after the commit landed must resolve to InvalidUploadID, got {retried:?}"
!upload_is_listed(&set_disks, bucket, object, &u_new).await,
"the idempotent retry must reclaim staging left behind by the crash"
);
let (body_after_retry, _) = read_object(&set_disks, bucket, object).await;
assert_eq!(body_after_retry, new, "the failed retry must not disturb the committed object");
// Reclaim the leftover exactly as the production tail does: delete_all
// on the upload path (abort_multipart_upload cannot — the upload's
+4 -1
View File
@@ -403,7 +403,10 @@ async fn enqueue_transition_after_write(
) -> Result<ObjectInfo> {
match result {
Ok(oi) => {
if should_enqueue_transition_immediately(&oi) {
// A retried completion did not publish a new object. Immediate
// transition/expiry already ran for the original commit; the scanner
// still sees the object if that commit never reached this hook.
if !oi.multipart_completion_replayed && should_enqueue_transition_immediately(&oi) {
enqueue_transition_immediate(&oi, src.clone()).await;
if let Ok(api) = metadata_sys::object_store_in(&store.ctx).await {
enqueue_immediate_expiry(api, &oi, src, opts).await;
+26 -4
View File
@@ -993,10 +993,32 @@ impl ECStore {
.await;
}
let pool_idx = self.multipart_upload_pool_idx(bucket, object, upload_id, opts).await?;
let pool = self.pools[pool_idx].clone();
pool.complete_multipart_upload(bucket, object, upload_id, uploaded_parts, opts)
.await
match self.multipart_upload_pool_idx(bucket, object, upload_id, opts).await {
Ok(pool_idx) => {
self.pools[pool_idx]
.clone()
.complete_multipart_upload(bucket, object, upload_id, uploaded_parts, opts)
.await
}
// The staging upload is already gone. The completed object lives on
// the pool that published it; other pools answer InvalidUploadID.
Err(err) if is_err_invalid_upload_id(&err) => {
let mut missing = err;
for pool_idx in self.existing_multipart_pool_order().await {
match self.pools[pool_idx]
.clone()
.complete_multipart_upload(bucket, object, upload_id, uploaded_parts.clone(), opts)
.await
{
Ok(info) => return Ok(info),
Err(err) if is_err_invalid_upload_id(&err) => missing = err,
Err(err) => return Err(err),
}
}
Err(missing)
}
Err(err) => Err(err),
}
}
#[cfg(all(test, feature = "test-util"))]
+3
View File
@@ -70,6 +70,9 @@ pub const SUFFIX_RESTORE_OPERATION_ID: &str = "restore-operation-id";
pub const SUFFIX_RESTORE_WORKER_LOCK: &str = "restore-worker-lock";
pub const RESTORE_WORKER_LOCK_PROTOCOL_V1: &str = "v1";
pub const SUFFIX_BUCKET_INCARNATION_ID: &str = "bucket-incarnation-id";
/// Upload ID of the multipart completion that published this object version.
/// A retried CompleteMultipartUpload matches it after the staging directory is gone.
pub const SUFFIX_MULTIPART_UPLOAD_ID: &str = "multipart-upload-id";
pub const SUFFIX_OBJECT_TRANSACTION_EPOCH: &str = "object-transaction-epoch";
/// Active rebalance run id mirrored onto `rebalance.bin` object metadata.
pub const SUFFIX_REBALANCE_RUN_ID: &str = "rebalance-run-id";
+136 -89
View File
@@ -40,7 +40,9 @@ use super::storage_api::multipart_usecase::contract::range::HTTPRangeSpec;
use super::storage_api::multipart_usecase::data_usage::{
quota_object_size, record_bucket_object_version_write_memory, record_bucket_object_write_memory,
};
use super::storage_api::multipart_usecase::error::{StorageError, is_err_object_not_found, is_err_version_not_found};
use super::storage_api::multipart_usecase::error::{
StorageError, is_err_invalid_upload_id, is_err_object_not_found, is_err_version_not_found,
};
use super::storage_api::multipart_usecase::helper::OperationHelper;
#[cfg(test)]
use super::storage_api::multipart_usecase::io::{DecryptReader, EncryptReader, HardLimitReader, boxed_reader, wrap_reader};
@@ -675,43 +677,39 @@ impl DefaultMultipartUsecase {
}
};
let multipart_info = store
.get_multipart_info(&bucket, &key, &upload_id, &opts)
.await
.map_err(ApiError::from)?;
// A lost Complete response retries the same upload id after staging is
// gone. `get_multipart_info` is NoSuchUpload in that case; completion
// below replays the committed object instead of failing the retry.
let multipart_info = match store.get_multipart_info(&bucket, &key, &upload_id, &opts).await {
Ok(info) => Some(info),
Err(err) if is_err_invalid_upload_id(&err) => None,
Err(err) => return Err(ApiError::from(err).into()),
};
// A ciphertext-passthrough session stores encrypted parts verbatim and
// completes without the customer key (the replication client has none),
// so the SSE-C completion check must be skipped for it.
if !contains_key_str(&multipart_info.user_defined, SUFFIX_REPLICATION_PRESERVE_CIPHERTEXT) {
EncryptionRequest {
bucket: &bucket,
key: &key,
server_side_encryption: None,
ssekms_key_id: None,
ssekms_context: None,
sse_customer_algorithm,
sse_customer_key,
sse_customer_key_md5,
content_size: 0,
principal: None,
if let Some(info) = multipart_info.as_ref() {
if !contains_key_str(&info.user_defined, SUFFIX_REPLICATION_PRESERVE_CIPHERTEXT) {
EncryptionRequest {
bucket: &bucket,
key: &key,
server_side_encryption: None,
ssekms_key_id: None,
ssekms_context: None,
sse_customer_algorithm: sse_customer_algorithm.clone(),
sse_customer_key: sse_customer_key.clone(),
sse_customer_key_md5: sse_customer_key_md5.clone(),
content_size: 0,
principal: None,
}
.validate_complete_multipart_ssec(&info.user_defined)?;
}
.validate_complete_multipart_ssec(&multipart_info.user_defined)?;
validate_complete_multipart_checksum_type(&req.headers, &info.user_defined)?;
}
validate_complete_multipart_checksum_type(&req.headers, &multipart_info.user_defined)?;
let cache_adapter = self.object_data_cache();
let _ = invalidate_object_data_cache_before_mutation(&cache_adapter, &bucket, &key).await;
let server_side_encryption = multipart_info
.user_defined
.get("x-amz-server-side-encryption")
.map(|s| ServerSideEncryption::from(s.clone()));
let ssekms_key_id = match server_side_encryption.as_ref() {
Some(sse) if sse.as_str() == ServerSideEncryption::AWS_KMS => multipart_info
.user_defined
.get("x-amz-server-side-encryption-aws-kms-key-id")
.cloned(),
_ => None,
};
let upload_user_defined = multipart_info.as_ref().map(|info| info.user_defined.clone());
let quota_metadata_sys = self.bucket_metadata_sys();
let quota_tracking = quota_metadata_sys.is_some();
@@ -735,38 +733,44 @@ impl DefaultMultipartUsecase {
// at multipart-session creation. Configuration and matching targets
// can change while parts are uploaded, so persist the exact decision's
// generation and PENDING set atomically with the completed object and
// reuse the same decision for scheduling below.
let completion_replication_decision = must_replicate_object(
&bucket,
&key,
&multipart_info.user_defined,
"".to_string(),
opts.delete_marker_replication_status(),
opts.clone(),
)
.await;
let mut completion_replication_metadata = HashMap::new();
if completion_replication_decision.replicate_any() {
insert_str(
&mut completion_replication_metadata,
SUFFIX_REPLICATION_GENERATION,
Uuid::new_v4().to_string(),
);
insert_str(
&mut completion_replication_metadata,
SUFFIX_REPLICATION_TIMESTAMP,
jiff::Zoned::now().to_string(),
);
insert_str(
&mut completion_replication_metadata,
SUFFIX_REPLICATION_STATUS,
completion_replication_decision.pending_status().unwrap_or_default(),
);
}
// `Some(empty)` deliberately means that completion re-evaluated the
// session as not admitted; storage removes stale Create-MPU admission
// metadata in the same final-object commit.
opts.eval_metadata = Some(completion_replication_metadata);
// reuse the same decision for scheduling below. A retry whose upload
// is already gone must not re-evaluate or rewrite that decision.
let completion_replication_decision = if let Some(user_defined) = upload_user_defined.as_ref() {
let completion_replication_decision = must_replicate_object(
&bucket,
&key,
user_defined,
"".to_string(),
opts.delete_marker_replication_status(),
opts.clone(),
)
.await;
let mut completion_replication_metadata = HashMap::new();
if completion_replication_decision.replicate_any() {
insert_str(
&mut completion_replication_metadata,
SUFFIX_REPLICATION_GENERATION,
Uuid::new_v4().to_string(),
);
insert_str(
&mut completion_replication_metadata,
SUFFIX_REPLICATION_TIMESTAMP,
jiff::Zoned::now().to_string(),
);
insert_str(
&mut completion_replication_metadata,
SUFFIX_REPLICATION_STATUS,
completion_replication_decision.pending_status().unwrap_or_default(),
);
}
// `Some(empty)` deliberately means that completion re-evaluated the
// session as not admitted; storage removes stale Create-MPU admission
// metadata in the same final-object commit.
opts.eval_metadata = Some(completion_replication_metadata);
Some(completion_replication_decision)
} else {
None
};
let complete_commit = spawn_traced_join({
let store = Arc::clone(&store);
@@ -781,30 +785,38 @@ impl DefaultMultipartUsecase {
.await
.map_err(ApiError::from)?;
let _ = invalidate_object_data_cache_after_complete_multipart_success(&cache_adapter, &bucket, &key).await;
// Staging reclaim on a replay can free bytes; the dirty mark is
// not a usage counter. Quota, replication, lifecycle, and the
// object-created scanner signal belong to the commit that
// published the object.
record_capacity_write(Some(capacity_scope_token)).await;
if quota_tracking {
let committed_size = quota_accounting_object_size(&obj_info, quota_enabled)?;
if !obj_info.multipart_completion_replayed {
if quota_tracking {
let committed_size = quota_accounting_object_size(&obj_info, quota_enabled)?;
if versioned {
record_bucket_object_version_write_memory(&bucket, previous_current_size, committed_size).await;
} else {
record_bucket_object_write_memory(&bucket, previous_current_size, committed_size).await;
if versioned {
record_bucket_object_version_write_memory(&bucket, previous_current_size, committed_size).await;
} else {
record_bucket_object_write_memory(&bucket, previous_current_size, committed_size).await;
}
}
enqueue_transition_immediate(&obj_info, LcEventSrc::S3CompleteMultipartUpload).await;
if let Some(decision) = completion_replication_decision
&& decision.replicate_any()
{
warn!("need multipart replication");
schedule_object_replication(obj_info.clone(), store, decision).await;
}
rustfs_scanner::record_dirty_usage_object_from_producer(
&bucket,
&key,
rustfs_scanner::SegmentInvalidationProducerIdentity::CompleteMultipartUpload,
);
}
enqueue_transition_immediate(&obj_info, LcEventSrc::S3CompleteMultipartUpload).await;
if completion_replication_decision.replicate_any() {
warn!("need multipart replication");
schedule_object_replication(obj_info.clone(), store, completion_replication_decision).await;
}
rustfs_scanner::record_dirty_usage_object_from_producer(
&bucket,
&key,
rustfs_scanner::SegmentInvalidationProducerIdentity::CompleteMultipartUpload,
);
Ok::<_, S3Error>(obj_info)
}
});
@@ -815,6 +827,42 @@ impl DefaultMultipartUsecase {
)
})??;
// The upload record is gone on a lost-response retry, so SSE-C is
// checked against the committed object instead of the staging metadata.
if upload_user_defined.is_none() && !contains_key_str(&obj_info.user_defined, SUFFIX_REPLICATION_PRESERVE_CIPHERTEXT) {
EncryptionRequest {
bucket: &bucket,
key: &key,
server_side_encryption: None,
ssekms_key_id: None,
ssekms_context: None,
sse_customer_algorithm,
sse_customer_key,
sse_customer_key_md5,
content_size: 0,
principal: None,
}
.validate_complete_multipart_ssec(obj_info.user_defined.as_ref())?;
}
let encryption_metadata = upload_user_defined.as_ref().unwrap_or(obj_info.user_defined.as_ref());
let server_side_encryption = encryption_metadata
.get("x-amz-server-side-encryption")
.map(|s| ServerSideEncryption::from(s.clone()));
let ssekms_key_id = match server_side_encryption.as_ref() {
Some(sse) if sse.as_str() == ServerSideEncryption::AWS_KMS => encryption_metadata
.get("x-amz-server-side-encryption-aws-kms-key-id")
.cloned(),
_ => None,
};
let ssec_algorithm = encryption_metadata
.get("x-amz-server-side-encryption-customer-algorithm")
.cloned();
let ssec_key_md5 = encryption_metadata
.get("x-amz-server-side-encryption-customer-key-md5")
.cloned();
let replayed = obj_info.multipart_completion_replayed;
let mpu_version = if versioned {
obj_info.version_id.map(|v| v.to_string())
} else {
@@ -862,26 +910,25 @@ impl DefaultMultipartUsecase {
let mut response = S3Response::new(output);
crate::app::object_usecase::inject_additional_checksum_headers(&mut response.headers, &complete_extra_checksum_headers);
if let Some(algorithm) = multipart_info
.user_defined
.get("x-amz-server-side-encryption-customer-algorithm")
{
if let Some(algorithm) = ssec_algorithm.as_deref() {
let value = HeaderValue::from_str(algorithm)
.map_err(|_| s3_error!(InternalError, "Invalid stored SSE-C algorithm metadata"))?;
response
.headers
.insert("x-amz-server-side-encryption-customer-algorithm", value);
}
if let Some(key_md5) = multipart_info
.user_defined
.get("x-amz-server-side-encryption-customer-key-md5")
{
if let Some(key_md5) = ssec_key_md5.as_deref() {
let value =
HeaderValue::from_str(key_md5).map_err(|_| s3_error!(InternalError, "Invalid stored SSE-C key metadata"))?;
response
.headers
.insert("x-amz-server-side-encryption-customer-key-md5", value);
}
if replayed {
// The object-created event was emitted by the commit that published
// the version. The retry is still a successful API call.
helper = helper.suppress_event();
}
let result = Ok(response);
let _ = helper.complete(&result);
result
+1 -1
View File
@@ -1041,7 +1041,7 @@ pub(crate) mod ecfs {
pub(crate) mod error {
pub(crate) use crate::storage::storage_api::{
StorageError, is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found,
StorageError, is_err_bucket_not_found, is_err_invalid_upload_id, is_err_object_not_found, is_err_version_not_found,
};
pub(crate) type Error = StorageError;
+6 -2
View File
@@ -492,8 +492,8 @@ pub(crate) mod ecstore_error {
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::error::PoolMetadataFailure;
pub(crate) use rustfs_ecstore::api::error::{
Error, PoolMetadataError, Result, StorageError, is_err_bucket_not_found, is_err_object_not_found,
is_err_version_not_found,
Error, PoolMetadataError, Result, StorageError, is_err_bucket_not_found, is_err_invalid_upload_id,
is_err_object_not_found, is_err_version_not_found,
};
}
@@ -2018,6 +2018,10 @@ pub(crate) fn is_err_bucket_not_found(err: &Error) -> bool {
ecstore_error::is_err_bucket_not_found(err)
}
pub(crate) fn is_err_invalid_upload_id(err: &Error) -> bool {
ecstore_error::is_err_invalid_upload_id(err)
}
pub(crate) fn is_err_object_not_found(err: &Error) -> bool {
ecstore_error::is_err_object_not_found(err)
}