From f704d015d6d9167254d97d9949b3cd73df614b14 Mon Sep 17 00:00:00 2001 From: houseme Date: Thu, 13 Aug 2026 20:19:33 +0800 Subject: [PATCH] fix(copy): keep copy commit owner alive (#6070) Keep S3 CopyObject's real outer owner task alive across caller cancellation so the source/destination bucket guards, same-key copy guard, storage commit, and post-commit publication hooks complete as one request-owned transaction boundary. Co-authored-by: heihutu --- rustfs/src/app/object_usecase.rs | 66 +++++++++++++++++---------- rustfs/src/app/storage_api.rs | 2 +- rustfs/src/storage/request_context.rs | 9 ++++ rustfs/src/storage/storage_api.rs | 4 +- 4 files changed, 55 insertions(+), 26 deletions(-) diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index 5c4e1e4c6..0e3bbb031 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -86,7 +86,7 @@ use super::storage_api::object_usecase::options::{ namespace_reserved_user_metadata, normalize_content_encoding_for_storage, preserve_unclassified_user_metadata, put_opts_with_replication_authorization, validate_archive_content_encoding, }; -use super::storage_api::object_usecase::request_context::{self, spawn_traced}; +use super::storage_api::object_usecase::request_context::{self, spawn_traced, spawn_traced_join}; use super::storage_api::object_usecase::s3_api::multipart::parse_list_parts_params; use super::storage_api::object_usecase::set_disk::{ get_lock_acquire_timeout, get_object_disk_read_timeout, is_valid_storage_class, @@ -7506,31 +7506,50 @@ impl DefaultObjectUsecase { let cache_adapter = self.object_data_cache(); let _ = invalidate_object_data_cache_before_mutation(&cache_adapter, &bucket, &key).await; - let oi = store - .copy_object(&src_bucket, &src_key, &bucket, &key, &mut src_info, &src_opts, &dst_opts) - .await - .map_err(ApiError::from)?; - drop(_self_copy_lock_guard); + let copy_commit = spawn_traced_join({ + let store = Arc::clone(&store); + let src_bucket = src_bucket.clone(); + let src_key = src_key.clone(); + let bucket = bucket.clone(); + let key = key.clone(); + let src_opts = src_opts.clone(); + let dst_opts = dst_opts.clone(); + async move { + let _source_bucket_lifecycle_guard = source_bucket_lifecycle_guard; + let _destination_bucket_lifecycle_guard_storage = destination_bucket_lifecycle_guard_storage; + let _self_copy_lock_guard = _self_copy_lock_guard; - // Reuse the single pre-commit replication decision (see `dsc` above) so - // the persisted pending marker and the schedule always agree, mirroring - // the PUT path. - if dsc.replicate_any() { - schedule_object_replication(oi.clone(), store.clone(), dsc).await; - } + let oi = store + .copy_object(&src_bucket, &src_key, &bucket, &key, &mut src_info, &src_opts, &dst_opts) + .await + .map_err(ApiError::from)?; - maybe_enqueue_transition_immediate(&oi, LcEventSrc::S3CopyObject).await; - let _ = invalidate_object_data_cache_after_copy_success(&cache_adapter, &bucket, &key).await; + // Reuse the single pre-commit replication decision (see `dsc` above) so + // the persisted pending marker and the schedule always agree, mirroring + // the PUT path. + if dsc.replicate_any() { + schedule_object_replication(oi.clone(), Arc::clone(&store), dsc).await; + } - let dest_versioned = BucketVersioningSys::prefix_enabled(&bucket, &key).await; - // Update quota tracking after successful copy - if has_bucket_metadata { - if dest_versioned { - record_bucket_object_version_write_memory(&bucket, previous_current_size, oi.size.max(0) as u64).await; - } else { - record_bucket_object_write_memory(&bucket, previous_current_size, oi.size.max(0) as u64).await; + maybe_enqueue_transition_immediate(&oi, LcEventSrc::S3CopyObject).await; + let _ = invalidate_object_data_cache_after_copy_success(&cache_adapter, &bucket, &key).await; + + let dest_versioned = BucketVersioningSys::prefix_enabled(&bucket, &key).await; + if has_bucket_metadata { + if dest_versioned { + record_bucket_object_version_write_memory(&bucket, previous_current_size, oi.size.max(0) as u64).await; + } else { + record_bucket_object_write_memory(&bucket, previous_current_size, oi.size.max(0) as u64).await; + } + } + + rustfs_scanner::record_dirty_usage_bucket(&bucket); + Ok::<_, S3Error>((oi, dest_versioned)) } - } + }); + let (oi, dest_versioned) = copy_commit.await.map_err(|err| { + S3Error::with_message(S3ErrorCode::InternalError, format!("copy object commit owner task failed: {err}")) + })??; let raw_dest_version = oi.version_id.map(|v| v.to_string()); let dest_version = if dest_versioned { raw_dest_version } else { None }; @@ -7578,7 +7597,7 @@ impl DefaultObjectUsecase { } } let copy_object_result = CopyObjectResult { - e_tag: oi.etag.map(|etag| to_s3s_etag(&etag)), + e_tag: oi.etag.as_ref().map(|etag| to_s3s_etag(etag)), last_modified: oi.mod_time.map(Timestamp::from), checksum_crc32: response_checksums.crc32, checksum_crc32c: response_checksums.crc32c, @@ -7609,7 +7628,6 @@ impl DefaultObjectUsecase { let result = Ok(S3Response::new(output)); let _ = helper.complete(&result); - rustfs_scanner::record_dirty_usage_bucket(&bucket); result } diff --git a/rustfs/src/app/storage_api.rs b/rustfs/src/app/storage_api.rs index cfb5c5493..c18e8c3b3 100644 --- a/rustfs/src/app/storage_api.rs +++ b/rustfs/src/app/storage_api.rs @@ -998,7 +998,7 @@ pub(crate) mod options { } pub(crate) mod request_context { - pub(crate) use crate::storage::storage_api::request_context_consumer::{RequestContext, spawn_traced}; + pub(crate) use crate::storage::storage_api::request_context_consumer::{RequestContext, spawn_traced, spawn_traced_join}; } pub(crate) mod sse { diff --git a/rustfs/src/storage/request_context.rs b/rustfs/src/storage/request_context.rs index d8bf66074..dddf5de92 100644 --- a/rustfs/src/storage/request_context.rs +++ b/rustfs/src/storage/request_context.rs @@ -257,6 +257,15 @@ where tokio::spawn(tracing::Instrument::instrument(fut, tracing::Span::current())); } +/// Spawn a request-internal task and return its join handle to the caller. +pub fn spawn_traced_join(fut: F) -> tokio::task::JoinHandle +where + F: std::future::Future + Send + 'static, + F::Output: Send + 'static, +{ + tokio::spawn(tracing::Instrument::instrument(fut, tracing::Span::current())) +} + #[cfg(test)] #[allow(unused_imports)] mod tests { diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index d5dfc3af5..b12169c71 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -203,7 +203,9 @@ pub(crate) mod options_consumer { } pub(crate) mod request_context_consumer { - pub(crate) use super::super::request_context::{RequestContext, extract_request_id_from_headers, spawn_traced}; + pub(crate) use super::super::request_context::{ + RequestContext, extract_request_id_from_headers, spawn_traced, spawn_traced_join, + }; } pub(crate) mod rpc_consumer {