From fc5609bbb025cb9359b061bef19060572cda2b14 Mon Sep 17 00:00:00 2001 From: RustFS Date: Mon, 28 Sep 2026 19:08:10 +0800 Subject: [PATCH] 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 --- crates/ecstore/src/api/mod.rs | 2 +- crates/ecstore/src/config/com.rs | 1 + crates/ecstore/src/object_api/types.rs | 5 + crates/ecstore/src/set_disk/ops/multipart.rs | 419 ++++++++++++++++++- crates/ecstore/src/store/mod.rs | 5 +- crates/ecstore/src/store/multipart.rs | 30 +- crates/utils/src/http/metadata_compat.rs | 3 + rustfs/src/app/multipart_usecase.rs | 225 ++++++---- rustfs/src/app/storage_api.rs | 2 +- rustfs/src/storage/storage_api.rs | 8 +- 10 files changed, 590 insertions(+), 110 deletions(-) diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index ec54f2a4d..b30addffc 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -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, }; } diff --git a/crates/ecstore/src/config/com.rs b/crates/ecstore/src/config/com.rs index 4081a3e3a..2521729a3 100644 --- a/crates/ecstore/src/config/com.rs +++ b/crates/ecstore/src/config/com.rs @@ -3154,6 +3154,7 @@ mod tests { version_purge_status: Default::default(), replication_decision: String::new(), checksum: None, + multipart_completion_replayed: false, } } } diff --git a/crates/ecstore/src/object_api/types.rs b/crates/ecstore/src/object_api/types.rs index 53fa10f14..0fda80ac5 100644 --- a/crates/ecstore/src/object_api/types.rs +++ b/crates/ecstore/src/object_api/types.rs @@ -1444,6 +1444,10 @@ pub struct ObjectInfo { pub version_purge_status: VersionPurgeStatusType, pub replication_decision: String, pub checksum: Option, + /// 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, } } diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index eb100eaec..1d3a4ee50 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -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> { + 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, + parts: &[ObjectPartInfo], + delete_marker: bool, + upload_id: &str, + uploaded_parts: &[CompletePart], + bucket: &str, + object: &str, +) -> Result { + 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 diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index 2e6338993..7ae4d2571 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -403,7 +403,10 @@ async fn enqueue_transition_after_write( ) -> Result { 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; diff --git a/crates/ecstore/src/store/multipart.rs b/crates/ecstore/src/store/multipart.rs index e0ce0e79c..a2b1f8896 100644 --- a/crates/ecstore/src/store/multipart.rs +++ b/crates/ecstore/src/store/multipart.rs @@ -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"))] diff --git a/crates/utils/src/http/metadata_compat.rs b/crates/utils/src/http/metadata_compat.rs index 2a60be1e5..305933cd8 100644 --- a/crates/utils/src/http/metadata_compat.rs +++ b/crates/utils/src/http/metadata_compat.rs @@ -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"; diff --git a/rustfs/src/app/multipart_usecase.rs b/rustfs/src/app/multipart_usecase.rs index f40679248..a6754bfe5 100644 --- a/rustfs/src/app/multipart_usecase.rs +++ b/rustfs/src/app/multipart_usecase.rs @@ -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 diff --git a/rustfs/src/app/storage_api.rs b/rustfs/src/app/storage_api.rs index 55cc875d0..23e63777b 100644 --- a/rustfs/src/app/storage_api.rs +++ b/rustfs/src/app/storage_api.rs @@ -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; diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 445e6968d..b5f6d8090 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -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) }