From 9ce9ec22d13484fd66cb3ce3b91ea9adcff41b6f Mon Sep 17 00:00:00 2001 From: GatewayJ <835269233@qq.com> Date: Mon, 1 Jun 2026 07:19:21 +0800 Subject: [PATCH] fix(ecstore): tighten object copy rename handling (#3131) * fix(ecstore): tighten object copy rename handling * fix(ecstore): narrow copy lock lifetime * test(ecstore): cover reverse copy concurrency * fix(multipart): ignore preconditions for internal lookup * fix(ecstore): clean precondition lock bindings --------- Co-authored-by: houseme --- crates/ecstore/src/disk/disk_store.rs | 3 +- crates/ecstore/src/disk/os.rs | 24 +- crates/ecstore/src/rpc/remote_disk.rs | 3 +- crates/ecstore/src/set_disk.rs | 274 ++++++++++++++++-- crates/ecstore/src/sets.rs | 11 +- crates/ecstore/src/store/object.rs | 257 +++++++++++++--- crates/ecstore/src/store/rebalance.rs | 14 +- .../src/app/lifecycle_transition_api_test.rs | 203 ++++++++++++- rustfs/src/app/multipart_usecase.rs | 33 ++- rustfs/src/app/object_usecase.rs | 93 +++++- 10 files changed, 814 insertions(+), 101 deletions(-) diff --git a/crates/ecstore/src/disk/disk_store.rs b/crates/ecstore/src/disk/disk_store.rs index 9c864ac8d..e02276d55 100644 --- a/crates/ecstore/src/disk/disk_store.rs +++ b/crates/ecstore/src/disk/disk_store.rs @@ -1241,7 +1241,8 @@ impl DiskAPI for LocalDiskWrapper { dst_volume: &str, dst_path: &str, ) -> Result { - self.track_disk_health( + self.track_disk_health_with_op( + "rename_data", || async { self.disk.rename_data(src_volume, src_path, fi, dst_volume, dst_path).await }, get_max_timeout_duration(), ) diff --git a/crates/ecstore/src/disk/os.rs b/crates/ecstore/src/disk/os.rs index c4172d8c9..b70cdb449 100644 --- a/crates/ecstore/src/disk/os.rs +++ b/crates/ecstore/src/disk/os.rs @@ -154,10 +154,6 @@ async fn reliable_rename( let mut i = 0; loop { if let Err(e) = super::fs::rename_std(src_file_path.as_ref(), dst_file_path.as_ref()) { - if e.kind() == io::ErrorKind::NotFound { - break; - } - if i == 0 { i += 1; continue; @@ -248,3 +244,23 @@ pub async fn os_mkdir_all(dir_path: impl AsRef, base_dir: impl AsRef pub fn file_exists(path: impl AsRef) -> bool { std::fs::metadata(path.as_ref()).map(|_| true).unwrap_or(false) } + +#[cfg(test)] +mod tests { + use super::*; + use tempfile::tempdir; + + #[tokio::test] + async fn rename_all_missing_source_returns_file_not_found() { + let temp_dir = tempdir().expect("create temp dir"); + let src = temp_dir.path().join("missing"); + let dst = temp_dir.path().join("dst"); + + let err = rename_all(&src, &dst, temp_dir.path()) + .await + .expect_err("missing source must fail"); + + assert!(matches!(err, DiskError::FileNotFound)); + assert!(!dst.exists()); + } +} diff --git a/crates/ecstore/src/rpc/remote_disk.rs b/crates/ecstore/src/rpc/remote_disk.rs index 29f7cad86..9ce5ec1fd 100644 --- a/crates/ecstore/src/rpc/remote_disk.rs +++ b/crates/ecstore/src/rpc/remote_disk.rs @@ -1148,7 +1148,8 @@ impl DiskAPI for RemoteDisk { ) -> Result { info!("rename_data {}/{}/{}/{}", self.addr, self.endpoint.to_string(), dst_volume, dst_path); - self.execute_with_timeout( + self.execute_with_timeout_for_op( + "rename_data", || async { let file_info = serde_json::to_string(&fi)?; let mut client = self diff --git a/crates/ecstore/src/set_disk.rs b/crates/ecstore/src/set_disk.rs index 678b5598f..898aa8a92 100644 --- a/crates/ecstore/src/set_disk.rs +++ b/crates/ecstore/src/set_disk.rs @@ -759,7 +759,7 @@ impl ObjectIO for SetDisks { let mut object_lock_guard = None; - if let Some(http_preconditions) = opts.http_preconditions.clone() { + if opts.http_preconditions.is_some() { if !opts.no_lock { let ns_lock = self.new_ns_lock(bucket, object).await?; object_lock_guard = Some( @@ -1408,8 +1408,8 @@ impl ObjectOperations for SetDisks { &self, src_bucket: &str, src_object: &str, - _dst_bucket: &str, - _dst_object: &str, + dst_bucket: &str, + dst_object: &str, src_info: &mut ObjectInfo, src_opts: &ObjectOptions, dst_opts: &ObjectOptions, @@ -1420,13 +1420,27 @@ impl ObjectOperations for SetDisks { return Err(StorageError::NotImplemented); } - // Guard lock for source object metadata update - let _lock_guard = self - .new_ns_lock(src_bucket, src_object) - .await? - .get_write_lock(get_lock_acquire_timeout()) - .await - .map_err(|e| self.map_namespace_lock_error(src_bucket, src_object, "write", e))?; + if path_join_buf(&[src_bucket, src_object]) != path_join_buf(&[dst_bucket, dst_object]) { + return Err(StorageError::NotImplemented); + } + + let _lock_guard = if dst_opts.no_lock { + None + } else { + Some( + self.new_ns_lock(dst_bucket, dst_object) + .await? + .get_write_lock(get_lock_acquire_timeout()) + .await + .map_err(|e| self.map_namespace_lock_error(dst_bucket, dst_object, "write", e))?, + ) + }; + + if dst_opts.http_preconditions.is_some() + && let Some(err) = self.check_write_precondition(dst_bucket, dst_object, dst_opts).await + { + return Err(err); + } let disks = self.get_disks_internal().await; @@ -1807,7 +1821,7 @@ impl ObjectOperations for SetDisks { #[tracing::instrument(skip(self))] async fn delete_object(&self, bucket: &str, object: &str, mut opts: ObjectOptions) -> Result { // Guard lock for single object delete - let _lock_guard = if !opts.delete_prefix { + let _lock_guard = if (!opts.delete_prefix || opts.delete_prefix_object) && !opts.no_lock { Some( self.new_ns_lock(bucket, object) .await? @@ -3072,18 +3086,18 @@ impl MultipartOperations for SetDisks { #[tracing::instrument(skip(self))] async fn new_multipart_upload(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result { - if let Some(http_preconditions) = opts.http_preconditions.clone() { - let object_lock_guard = if !opts.no_lock { + let mut _object_lock_guard = None; + + if opts.http_preconditions.is_some() { + if !opts.no_lock { let ns_lock = self.new_ns_lock(bucket, object).await?; - Some( + _object_lock_guard = Some( ns_lock .get_write_lock(get_lock_acquire_timeout()) .await .map_err(|e| self.map_namespace_lock_error(bucket, object, "write", e))?, - ) - } else { - None - }; + ); + } if let Some(err) = self.check_write_precondition(bucket, object, opts).await { return Err(err); @@ -3253,8 +3267,7 @@ impl MultipartOperations for SetDisks { ) -> Result { let mut object_lock_guard = None; - // Acquire per-object exclusive lock via RAII guard. It auto-releases asynchronously on drop. - if let Some(http_preconditions) = opts.http_preconditions.clone() { + if opts.http_preconditions.is_some() { if !opts.no_lock { let ns_lock = self.new_ns_lock(bucket, object).await?; object_lock_guard = Some( @@ -5003,6 +5016,227 @@ mod tests { ); } + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn copy_object_honors_no_lock_when_outer_write_lock_is_held() { + let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::Erasure).await; + let set_disks = make_test_set_disks(vec![Arc::new(LocalClient::with_manager(Arc::new( + rustfs_lock::GlobalLockManager::new(), + )))]) + .await; + + let _outer_guard = set_disks + .new_ns_lock("bucket", "object") + .await + .expect("namespace lock should be created") + .get_write_lock(Duration::from_secs(1)) + .await + .expect("outer write lock should be acquired"); + + let mut src_info = ObjectInfo { + metadata_only: true, + ..Default::default() + }; + let dst_opts = ObjectOptions { + no_lock: true, + ..Default::default() + }; + + let result = tokio::time::timeout( + Duration::from_secs(1), + set_disks.copy_object( + "bucket", + "object", + "bucket", + "object", + &mut src_info, + &ObjectOptions::default(), + &dst_opts, + ), + ) + .await + .expect("no_lock copy path must not wait for the outer lock"); + + let err = result.expect_err("empty test disks should fail after bypassing the inner lock"); + assert!( + !err.to_string().to_ascii_lowercase().contains("lock"), + "copy_object returned a lock error despite no_lock=true: {err}" + ); + } + + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn copy_object_rejects_metadata_only_cross_key() { + let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::Erasure).await; + let set_disks = make_test_set_disks(vec![Arc::new(LocalClient::with_manager(Arc::new( + rustfs_lock::GlobalLockManager::new(), + )))]) + .await; + + let mut src_info = ObjectInfo { + metadata_only: true, + ..Default::default() + }; + + let err = set_disks + .copy_object( + "bucket", + "source", + "bucket", + "dest", + &mut src_info, + &ObjectOptions::default(), + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + .expect_err("metadata-only lower copy is only valid for self-copy updates"); + + assert!(matches!(err, StorageError::NotImplemented)); + } + + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn delete_object_honors_no_lock_when_outer_write_lock_is_held() { + let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::Erasure).await; + let set_disks = make_test_set_disks(vec![Arc::new(LocalClient::with_manager(Arc::new( + rustfs_lock::GlobalLockManager::new(), + )))]) + .await; + + let _outer_guard = set_disks + .new_ns_lock("bucket", "object") + .await + .expect("namespace lock should be created") + .get_write_lock(Duration::from_secs(1)) + .await + .expect("outer write lock should be acquired"); + + let result = tokio::time::timeout( + Duration::from_secs(1), + set_disks.delete_object( + "bucket", + "object", + ObjectOptions { + no_lock: true, + ..Default::default() + }, + ), + ) + .await + .expect("no_lock delete path must not wait for the outer lock"); + + let err = result.expect_err("empty test disks should fail after bypassing the inner lock"); + assert!( + !err.to_string().to_ascii_lowercase().contains("lock"), + "delete_object returned a lock error despite no_lock=true: {err}" + ); + } + + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn delete_prefix_does_not_lock_literal_prefix_key() { + let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::Erasure).await; + let set_disks = make_test_set_disks(vec![Arc::new(LocalClient::with_manager(Arc::new( + rustfs_lock::GlobalLockManager::new(), + )))]) + .await; + + let _outer_guard = set_disks + .new_ns_lock("bucket", "prefix") + .await + .expect("namespace lock should be created") + .get_write_lock(Duration::from_secs(1)) + .await + .expect("outer write lock should be acquired"); + + tokio::time::timeout( + Duration::from_secs(1), + set_disks.delete_object( + "bucket", + "prefix", + ObjectOptions { + delete_prefix: true, + ..Default::default() + }, + ), + ) + .await + .expect("broad prefix delete must not wait on a literal prefix namespace lock") + .expect("empty test disks should allow broad prefix cleanup"); + } + + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn delete_prefix_object_honors_no_lock_when_outer_write_lock_is_held() { + let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::Erasure).await; + let set_disks = make_test_set_disks(vec![Arc::new(LocalClient::with_manager(Arc::new( + rustfs_lock::GlobalLockManager::new(), + )))]) + .await; + + let _outer_guard = set_disks + .new_ns_lock("bucket", "object") + .await + .expect("namespace lock should be created") + .get_write_lock(Duration::from_secs(1)) + .await + .expect("outer write lock should be acquired"); + + tokio::time::timeout( + Duration::from_secs(1), + set_disks.delete_object( + "bucket", + "object", + ObjectOptions { + delete_prefix: true, + delete_prefix_object: true, + no_lock: true, + ..Default::default() + }, + ), + ) + .await + .expect("no_lock exact prefix delete path must not wait for the outer lock") + .expect("empty test disks should allow exact prefix cleanup"); + } + + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn delete_prefix_object_locks_real_object_key() { + let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::Erasure).await; + let set_disks = make_test_set_disks(vec![Arc::new(LocalClient::with_manager(Arc::new( + rustfs_lock::GlobalLockManager::new(), + )))]) + .await; + + let _outer_guard = set_disks + .new_ns_lock("bucket", "object") + .await + .expect("namespace lock should be created") + .get_write_lock(Duration::from_secs(1)) + .await + .expect("outer write lock should be acquired"); + + let result = tokio::time::timeout( + Duration::from_millis(50), + set_disks.delete_object( + "bucket", + "object", + ObjectOptions { + delete_prefix: true, + delete_prefix_object: true, + ..Default::default() + }, + ), + ) + .await; + + assert!(result.is_err(), "exact prefix delete should wait on the real object namespace lock"); + } + #[tokio::test(flavor = "multi_thread")] #[serial] async fn test_acquire_dist_delete_object_locks_batch_succeeds_with_two_healthy_lockers() { diff --git a/crates/ecstore/src/sets.rs b/crates/ecstore/src/sets.rs index 48502a4b6..e06d29aab 100644 --- a/crates/ecstore/src/sets.rs +++ b/crates/ecstore/src/sets.rs @@ -306,12 +306,10 @@ impl Sets { // unimplemented!() // } - async fn delete_prefix(&self, bucket: &str, object: &str) -> Result<()> { + async fn delete_prefix(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<()> { let mut futures = Vec::new(); - let opt = ObjectOptions { - delete_prefix: true, - ..Default::default() - }; + let mut opt = opts.clone(); + opt.delete_prefix = true; for set in self.disk_set.iter() { futures.push(set.delete_object(bucket, object, opt.clone())); @@ -463,6 +461,7 @@ impl ObjectOperations for Sets { versioned: dst_opts.versioned, version_id: dst_opts.version_id.clone(), mod_time: dst_opts.mod_time, + http_preconditions: dst_opts.http_preconditions.clone(), ..Default::default() }; @@ -485,7 +484,7 @@ impl ObjectOperations for Sets { #[tracing::instrument(skip(self))] async fn delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result { if opts.delete_prefix && !opts.delete_prefix_object { - self.delete_prefix(bucket, object).await?; + self.delete_prefix(bucket, object, &opts).await?; return Ok(ObjectInfo::default()); } diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 1d6ba282f..ec5bdf0aa 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -13,6 +13,30 @@ // limitations under the License. use super::*; +use crate::set_disk::{get_lock_acquire_timeout, is_lock_optimization_enabled}; +use std::{ + pin::Pin, + task::{Context, Poll}, +}; +use tokio::io::{AsyncRead, ReadBuf}; + +struct LockGuardedReader { + inner: Box, + guard: Option, +} + +impl AsyncRead for LockGuardedReader { + fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { + let had_capacity = buf.remaining() > 0; + let filled_before = buf.filled().len(); + let poll = Pin::new(&mut self.inner).poll_read(cx, buf); + if had_capacity && matches!(poll, Poll::Ready(Ok(()))) && buf.filled().len() == filled_before { + self.guard.take(); + } + poll + } +} + fn select_data_movement_target_pool( existing_pool_idx: Result, src_pool_idx: usize, @@ -89,6 +113,74 @@ fn data_movement_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> Object } impl ECStore { + fn map_namespace_lock_error(bucket: &str, object: &str, mode: &'static str, err: rustfs_lock::LockError) -> StorageError { + match err { + rustfs_lock::LockError::QuorumNotReached { required, achieved } => StorageError::NamespaceLockQuorumUnavailable { + mode, + bucket: bucket.to_string(), + object: object.to_string(), + required, + achieved, + }, + other => StorageError::other(format!("Failed to acquire {mode} lock on {bucket}/{object}: {other}")), + } + } + + async fn acquire_object_write_lock_if_needed( + &self, + bucket: &str, + object: &str, + opts: &mut ObjectOptions, + ) -> Result> { + if opts.no_lock { + return Ok(None); + } + + let ns_lock = self.handle_new_ns_lock(bucket, object).await?; + let guard = ns_lock + .get_write_lock(get_lock_acquire_timeout()) + .await + .map_err(|err| Self::map_namespace_lock_error(bucket, object, "write", err))?; + opts.no_lock = true; + + Ok(Some(guard)) + } + + async fn acquire_object_read_lock_if_needed( + &self, + bucket: &str, + object: &str, + opts: &mut ObjectOptions, + ) -> Result> { + if opts.no_lock { + return Ok(None); + } + + let ns_lock = self.handle_new_ns_lock(bucket, object).await?; + let guard = ns_lock + .get_read_lock(get_lock_acquire_timeout()) + .await + .map_err(|err| Self::map_namespace_lock_error(bucket, object, "read", err))?; + opts.no_lock = true; + + Ok(Some(guard)) + } + + fn attach_read_lock_guard(mut reader: GetObjectReader, guard: Option) -> GetObjectReader { + if is_lock_optimization_enabled() { + return reader; + } + + if let Some(guard) = guard { + reader.stream = Box::new(LockGuardedReader { + inner: reader.stream, + guard: Some(guard), + }); + } + + reader + } + async fn get_latest_accessible_object_info_with_idx( &self, bucket: &str, @@ -193,23 +285,23 @@ impl ECStore { check_get_obj_args(bucket, object)?; let object = encode_dir_object(object); - - if self.single_pool() { - return self.pools[0].get_object_reader(bucket, object.as_str(), range, h, opts).await; - } - - // TODO: nslock - let mut opts = opts.clone(); + let read_lock_guard = self.acquire_object_read_lock_if_needed(bucket, &object, &mut opts).await?; - opts.no_lock = true; + let reader = if self.single_pool() { + self.pools[0] + .get_object_reader(bucket, object.as_str(), range, h, &opts) + .await? + } else { + let (_, idx) = self + .get_latest_accessible_object_info_with_idx(bucket, &object, &opts) + .await?; + self.pools[idx] + .get_object_reader(bucket, object.as_str(), range, h, &opts) + .await? + }; - let (_, idx) = self - .get_latest_accessible_object_info_with_idx(bucket, &object, &opts) - .await?; - self.pools[idx] - .get_object_reader(bucket, object.as_str(), range, h, &opts) - .await + Ok(Self::attach_read_lock_guard(reader, read_lock_guard)) } #[instrument(level = "debug", skip(self, data))] @@ -224,6 +316,8 @@ impl ECStore { let object = encode_dir_object(object); + // Keep PUT atomic-read friendly: SetDisks takes the object write lock only + // around precondition checks and the final rename/commit. if self.single_pool() { return self.pools[0].put_object(bucket, object.as_str(), data, opts).await; } @@ -251,16 +345,16 @@ impl ECStore { check_object_args(bucket, object)?; let object = encode_dir_object(object); + let mut opts = opts.clone(); + let _object_lock_guard = self.acquire_object_read_lock_if_needed(bucket, &object, &mut opts).await?; - if self.single_pool() { - return self.pools[0].get_object_info(bucket, object.as_str(), opts).await; - } - - // TODO: nslock - - let (info, _) = self - .get_latest_accessible_object_info_with_idx(bucket, object.as_str(), opts) - .await?; + let info = if self.single_pool() { + self.pools[0].get_object_info(bucket, object.as_str(), &opts).await? + } else { + self.get_latest_accessible_object_info_with_idx(bucket, object.as_str(), &opts) + .await? + .0 + }; opts.precondition_check(&info)?; Ok(info) } @@ -285,43 +379,56 @@ impl ECStore { let cp_src_dst_same = path_join_buf(&[src_bucket, &src_object]) == path_join_buf(&[dst_bucket, &dst_object]); - // TODO: nslock - - let pool_idx = self - .get_pool_info_existing_with_opts(src_bucket, &src_object, &version_aware_lookup_opts(src_opts, true)) - .await? - .0 - .index; + let mut dst_opts = dst_opts.clone(); + let _dst_lock_guard = if cp_src_dst_same { + self.acquire_object_write_lock_if_needed(dst_bucket, &dst_object, &mut dst_opts) + .await? + } else { + None + }; if cp_src_dst_same { + let pool_idx = self + .get_pool_info_existing_with_opts(src_bucket, &src_object, &version_aware_lookup_opts(src_opts, true)) + .await? + .0 + .index; + if let (Some(src_vid), Some(dst_vid)) = (&src_opts.version_id, &dst_opts.version_id) && src_vid == dst_vid { return self.pools[pool_idx] - .copy_object(src_bucket, &src_object, dst_bucket, &dst_object, src_info, src_opts, dst_opts) + .copy_object(src_bucket, &src_object, dst_bucket, &dst_object, src_info, src_opts, &dst_opts) .await; } if !dst_opts.versioned && src_opts.version_id.is_none() { return self.pools[pool_idx] - .copy_object(src_bucket, &src_object, dst_bucket, &dst_object, src_info, src_opts, dst_opts) + .copy_object(src_bucket, &src_object, dst_bucket, &dst_object, src_info, src_opts, &dst_opts) .await; } if dst_opts.versioned && src_opts.version_id != dst_opts.version_id { src_info.version_only = true; return self.pools[pool_idx] - .copy_object(src_bucket, &src_object, dst_bucket, &dst_object, src_info, src_opts, dst_opts) + .copy_object(src_bucket, &src_object, dst_bucket, &dst_object, src_info, src_opts, &dst_opts) .await; } } + let pool_idx = if dst_opts.no_lock { + self.get_pool_idx_no_lock(dst_bucket, &dst_object, src_info.size).await? + } else { + self.get_pool_idx(dst_bucket, &dst_object, src_info.size).await? + }; + let put_opts = ObjectOptions { user_defined: src_info.user_defined.clone(), versioned: dst_opts.versioned, version_id: dst_opts.version_id.clone(), - no_lock: true, + no_lock: dst_opts.no_lock, mod_time: dst_opts.mod_time, + http_preconditions: dst_opts.http_preconditions.clone(), ..Default::default() }; @@ -342,15 +449,27 @@ impl ECStore { pub(super) async fn handle_delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result { check_del_obj_args(bucket, object)?; - if opts.delete_prefix { - self.delete_prefix(bucket, object).await?; + let object = if opts.delete_prefix && !opts.delete_prefix_object { + object.to_owned() + } else { + encode_dir_object(object) + }; + let object = object.as_str(); + let mut opts = opts; + + if opts.delete_prefix && !opts.delete_prefix_object { + // Prefix deletes cover multiple object keys; an exact lock on the prefix string + // would not protect child objects. + self.delete_prefix(bucket, object, &opts).await?; return Ok(ObjectInfo::default()); } - // TODO: nslock + let _object_lock_guard = self.acquire_object_write_lock_if_needed(bucket, object, &mut opts).await?; - let object = encode_dir_object(object); - let object = object.as_str(); + if opts.delete_prefix { + self.delete_prefix(bucket, object, &opts).await?; + return Ok(ObjectInfo::default()); + } let gopts = version_aware_lookup_opts(&opts, true); @@ -775,6 +894,8 @@ impl ECStore { #[cfg(test)] mod tests { use super::*; + use std::io::Cursor; + use tokio::io::AsyncReadExt; #[test] fn delete_marker_data_movement_falls_back_when_only_source_pool_has_object() { @@ -983,4 +1104,62 @@ mod tests { assert!(lookup_opts.skip_decommissioned); assert!(lookup_opts.skip_rebalancing); } + + #[tokio::test] + #[serial_test::serial] + async fn reader_lock_is_held_when_optimization_is_disabled() { + temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some("false"))], async { + let manager = Arc::new(rustfs_lock::GlobalLockManager::new()); + let lock = rustfs_lock::NamespaceLock::with_local_manager("test".to_string(), manager); + let key = rustfs_lock::ObjectKey::new("bucket", "object"); + let read_guard = lock + .get_read_lock(key.clone(), "reader", Duration::from_secs(1)) + .await + .expect("read lock should be acquired"); + let reader = GetObjectReader { + stream: Box::new(Cursor::new(Vec::::new())), + object_info: ObjectInfo::default(), + }; + + let reader = ECStore::attach_read_lock_guard(reader, Some(read_guard)); + + lock.get_write_lock(key.clone(), "writer", Duration::from_millis(20)) + .await + .expect_err("reader should hold the read lock"); + drop(reader); + lock.get_write_lock(key, "writer", Duration::from_secs(1)) + .await + .expect("dropping the reader should release the read lock"); + }) + .await; + } + + #[tokio::test] + #[serial_test::serial] + async fn reader_lock_is_released_after_stream_eof() { + temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some("false"))], async { + let manager = Arc::new(rustfs_lock::GlobalLockManager::new()); + let lock = rustfs_lock::NamespaceLock::with_local_manager("test".to_string(), manager); + let key = rustfs_lock::ObjectKey::new("bucket", "object"); + let read_guard = lock + .get_read_lock(key.clone(), "reader", Duration::from_secs(1)) + .await + .expect("read lock should be acquired"); + let reader = GetObjectReader { + stream: Box::new(Cursor::new(vec![1, 2, 3])), + object_info: ObjectInfo::default(), + }; + + let mut reader = ECStore::attach_read_lock_guard(reader, Some(read_guard)); + let mut output = Vec::new(); + reader.stream.read_to_end(&mut output).await.expect("reader should reach EOF"); + assert_eq!(output, vec![1, 2, 3]); + + lock.get_write_lock(key, "writer", Duration::from_secs(1)) + .await + .expect("EOF should release the read lock before the reader is dropped"); + drop(reader); + }) + .await; + } } diff --git a/crates/ecstore/src/store/rebalance.rs b/crates/ecstore/src/store/rebalance.rs index de788cfd5..64ff51a05 100644 --- a/crates/ecstore/src/store/rebalance.rs +++ b/crates/ecstore/src/store/rebalance.rs @@ -183,17 +183,11 @@ impl ECStore { Ok(()) } - pub(super) async fn delete_prefix(&self, bucket: &str, object: &str) -> Result<()> { + pub(super) async fn delete_prefix(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<()> { for pool in self.pools.iter() { - pool.delete_object( - bucket, - object, - ObjectOptions { - delete_prefix: true, - ..Default::default() - }, - ) - .await?; + let mut opts = opts.clone(); + opts.delete_prefix = true; + pool.delete_object(bucket, object, opts).await?; } Ok(()) diff --git a/rustfs/src/app/lifecycle_transition_api_test.rs b/rustfs/src/app/lifecycle_transition_api_test.rs index a02a22644..830804f2b 100644 --- a/rustfs/src/app/lifecycle_transition_api_test.rs +++ b/rustfs/src/app/lifecycle_transition_api_test.rs @@ -18,8 +18,8 @@ use crate::storage::ecfs::FS; use bytes::Bytes; use futures::FutureExt; use futures::stream; -use http::{Extensions, HeaderMap, Method, Uri}; -use rustfs_config::ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT; +use http::{Extensions, HeaderMap, HeaderValue, Method, Uri, header::IF_NONE_MATCH}; +use rustfs_config::{ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, ENV_TEST_FORCE_IMMEDIATE_TRANSITION_ENQUEUE_TIMEOUT}; use rustfs_ecstore::{ bucket::metadata::{BUCKET_LIFECYCLE_CONFIG, OBJECT_LOCK_CONFIG}, bucket::metadata_sys, @@ -53,7 +53,7 @@ use std::{ }; use tokio::fs; use tokio::io::AsyncReadExt; -use tokio::sync::Mutex; +use tokio::sync::{Barrier, Mutex}; use tokio_util::sync::CancellationToken; use uuid::Uuid; @@ -484,6 +484,203 @@ fn streaming_blob_from_bytes(data: &[u8]) -> StreamingBlob { StreamingBlob::wrap::<_, Infallible>(stream::once(async move { Ok(body) })) } +async fn read_object_bytes(ecstore: &Arc, bucket: &str, object: &str) -> Vec { + let mut reader = (**ecstore) + .get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("Failed to read object"); + let mut buf = Vec::new(); + reader + .stream + .read_to_end(&mut buf) + .await + .expect("Failed to drain object reader"); + buf +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 1)] +#[serial] +#[ignore = "requires isolated global object layer state"] +async fn put_object_if_none_match_existing_object_returns_precondition_failed() { + let (_disk_paths, ecstore) = setup_test_env().await; + let fs = FS::new(); + let usecase = DefaultObjectUsecase::without_context(); + + let bucket = format!("test-put-if-none-match-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object = "test/object.txt"; + let initial_payload = b"initial conditional put payload"; + let replacement_payload = b"replacement conditional put payload"; + + create_test_bucket(&ecstore, bucket.as_str()).await; + + let initial_input = PutObjectInput::builder() + .bucket(bucket.clone()) + .key(object.to_string()) + .body(Some(streaming_blob_from_bytes(initial_payload))) + .content_length(Some(initial_payload.len() as i64)) + .build() + .unwrap(); + Box::pin(usecase.execute_put_object(&fs, build_request(initial_input, Method::PUT))) + .await + .expect("Failed to upload initial object through usecase"); + + let existing_info = ecstore + .get_object_info(bucket.as_str(), object, &ObjectOptions::default()) + .await + .expect("Failed to fetch existing object info"); + let existing_etag = existing_info.etag.expect("existing object should have an ETag"); + + let replacement_input = PutObjectInput::builder() + .bucket(bucket.clone()) + .key(object.to_string()) + .body(Some(streaming_blob_from_bytes(replacement_payload))) + .content_length(Some(replacement_payload.len() as i64)) + .build() + .unwrap(); + let mut req = build_request(replacement_input, Method::PUT); + req.headers + .insert(IF_NONE_MATCH, HeaderValue::from_str(existing_etag.as_str()).unwrap()); + + let err = Box::pin(usecase.execute_put_object(&fs, req)).await.unwrap_err(); + + assert_eq!(err.code(), &s3s::S3ErrorCode::PreconditionFailed); + assert_eq!( + read_object_bytes(&ecstore, bucket.as_str(), object).await, + initial_payload, + "failed conditional PutObject must not overwrite the current object" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 1)] +#[serial] +#[ignore = "requires isolated global object layer state"] +async fn copy_object_if_none_match_existing_destination_returns_precondition_failed() { + let (_disk_paths, ecstore) = setup_test_env().await; + let usecase = DefaultObjectUsecase::without_context(); + + let src_bucket = format!("test-copy-if-none-match-src-{}", &Uuid::new_v4().simple().to_string()[..8]); + let dst_bucket = format!("test-copy-if-none-match-dst-{}", &Uuid::new_v4().simple().to_string()[..8]); + let src_object = "test/source.txt"; + let dst_object = "test/destination.txt"; + let src_payload = b"conditional copy source payload"; + let dst_payload = b"conditional copy destination payload"; + + create_test_bucket(&ecstore, src_bucket.as_str()).await; + create_test_bucket(&ecstore, dst_bucket.as_str()).await; + let _ = upload_test_object(&ecstore, src_bucket.as_str(), src_object, src_payload).await; + let _ = upload_test_object(&ecstore, dst_bucket.as_str(), dst_object, dst_payload).await; + + let dst_info = ecstore + .get_object_info(dst_bucket.as_str(), dst_object, &ObjectOptions::default()) + .await + .expect("Failed to fetch destination object info"); + let dst_etag = dst_info.etag.expect("destination object should have an ETag"); + + let copy_input = CopyObjectInput::builder() + .copy_source(CopySource::Bucket { + bucket: src_bucket.clone().into(), + key: src_object.to_string().into(), + version_id: None, + }) + .bucket(dst_bucket.clone()) + .key(dst_object.to_string()) + .build() + .unwrap(); + let mut req = build_request(copy_input, Method::PUT); + req.headers + .insert(IF_NONE_MATCH, HeaderValue::from_str(dst_etag.as_str()).unwrap()); + + let err = Box::pin(usecase.execute_copy_object(req)).await.unwrap_err(); + + assert_eq!(err.code(), &s3s::S3ErrorCode::PreconditionFailed); + assert_eq!( + read_object_bytes(&ecstore, dst_bucket.as_str(), dst_object).await, + dst_payload, + "failed conditional CopyObject must not overwrite the current destination object" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +#[serial] +#[ignore = "requires isolated global object layer state"] +async fn concurrent_reverse_copy_object_does_not_deadlock_with_reader_locks() { + temp_env::async_with_vars([(ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some("false"))], async { + let (_disk_paths, ecstore) = setup_test_env().await; + + let bucket = format!("test-reverse-copy-{}", &Uuid::new_v4().simple().to_string()[..8]); + let object_a = "test/a.bin"; + let object_b = "test/b.bin"; + let payload_a = vec![b'a'; 4 * 1024 * 1024]; + let payload_b = vec![b'b'; 4 * 1024 * 1024]; + + create_test_bucket(&ecstore, bucket.as_str()).await; + let _ = upload_test_object(&ecstore, bucket.as_str(), object_a, &payload_a).await; + let _ = upload_test_object(&ecstore, bucket.as_str(), object_b, &payload_b).await; + + let a_to_b_input = CopyObjectInput::builder() + .copy_source(CopySource::Bucket { + bucket: bucket.clone().into(), + key: object_a.to_string().into(), + version_id: None, + }) + .bucket(bucket.clone()) + .key(object_b.to_string()) + .build() + .unwrap(); + let b_to_a_input = CopyObjectInput::builder() + .copy_source(CopySource::Bucket { + bucket: bucket.clone().into(), + key: object_b.to_string().into(), + version_id: None, + }) + .bucket(bucket.clone()) + .key(object_a.to_string()) + .build() + .unwrap(); + + let start = Arc::new(Barrier::new(3)); + let start_a = start.clone(); + let a_to_b = tokio::spawn(async move { + start_a.wait().await; + let usecase = DefaultObjectUsecase::without_context(); + Box::pin(usecase.execute_copy_object(build_request(a_to_b_input, Method::PUT))).await + }); + let start_b = start.clone(); + let b_to_a = tokio::spawn(async move { + start_b.wait().await; + let usecase = DefaultObjectUsecase::without_context(); + Box::pin(usecase.execute_copy_object(build_request(b_to_a_input, Method::PUT))).await + }); + + start.wait().await; + + let (a_to_b_result, b_to_a_result) = tokio::time::timeout(Duration::from_secs(10), async { + let (a_to_b_result, b_to_a_result) = tokio::join!(a_to_b, b_to_a); + ( + a_to_b_result.expect("A-to-B CopyObject task should not panic"), + b_to_a_result.expect("B-to-A CopyObject task should not panic"), + ) + }) + .await + .expect("reverse CopyObject operations should not deadlock"); + + a_to_b_result.expect("A-to-B CopyObject should succeed"); + b_to_a_result.expect("B-to-A CopyObject should succeed"); + + let final_a = read_object_bytes(&ecstore, bucket.as_str(), object_a).await; + let final_b = read_object_bytes(&ecstore, bucket.as_str(), object_b).await; + assert!( + final_a == payload_a || final_a == payload_b, + "object A must contain a complete copied payload" + ); + assert!( + final_b == payload_a || final_b == payload_b, + "object B must contain a complete copied payload" + ); + }) + .await; +} + #[tokio::test(flavor = "multi_thread", worker_threads = 1)] #[serial] #[ignore = "requires isolated global object layer state"] diff --git a/rustfs/src/app/multipart_usecase.rs b/rustfs/src/app/multipart_usecase.rs index b68b93717..d11eebaac 100644 --- a/rustfs/src/app/multipart_usecase.rs +++ b/rustfs/src/app/multipart_usecase.rs @@ -148,6 +148,11 @@ fn has_complete_multipart_object_lock_headers(headers: &HeaderMap) -> bool { || has_bypass_governance_header(headers) } +fn internal_object_info_lookup_opts(mut opts: ObjectOptions) -> ObjectOptions { + opts.http_preconditions = None; + opts +} + fn encode_s3_path(path: &str) -> String { path.split('/') .map(|part| encode(part).to_string()) @@ -336,9 +341,11 @@ impl DefaultMultipartUsecase { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; - let current_opts = get_opts(&bucket, &key, None, None, &req.headers) - .await - .map_err(ApiError::from)?; + let current_opts = internal_object_info_lookup_opts( + get_opts(&bucket, &key, None, None, &req.headers) + .await + .map_err(ApiError::from)?, + ); let previous_current_size = match store.get_object_info(&bucket, &key, ¤t_opts).await { Ok(existing_obj_info) => { validate_existing_object_lock_for_write(&existing_obj_info)?; @@ -1317,6 +1324,26 @@ mod tests { assert_eq!(location, "/bucket/nested/object"); } + #[test] + fn internal_object_info_lookup_opts_drops_http_preconditions() { + let opts = ObjectOptions { + version_id: Some(Uuid::new_v4().to_string()), + no_lock: true, + http_preconditions: Some(rustfs_ecstore::store_api::HTTPPreconditions { + if_none_match: Some("*".to_string()), + if_match: Some("\"etag\"".to_string()), + ..Default::default() + }), + ..Default::default() + }; + + let lookup_opts = internal_object_info_lookup_opts(opts); + + assert!(lookup_opts.http_preconditions.is_none()); + assert!(lookup_opts.no_lock); + assert!(lookup_opts.version_id.is_some()); + } + #[test] fn merge_part_encryption_metadata_keeps_source_metadata_unchanged() { let multipart_metadata = HashMap::from([ diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index fa7eee0a9..d30160590 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -70,9 +70,9 @@ use rustfs_ecstore::config::storageclass; use rustfs_ecstore::disk::{error::DiskError, error_reduce::is_all_buckets_not_found}; use rustfs_ecstore::error::{StorageError, is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found}; use rustfs_ecstore::new_object_layer_fn; -use rustfs_ecstore::set_disk::is_valid_storage_class; +use rustfs_ecstore::set_disk::{get_lock_acquire_timeout, is_valid_storage_class}; use rustfs_ecstore::store_api::{ - HTTPRangeSpec, ObjectIO, ObjectInfo, ObjectOperations, ObjectOptions, ObjectToDelete, PutObjReader, + HTTPRangeSpec, ObjectIO, ObjectInfo, ObjectOperations, ObjectOptions, ObjectToDelete, PutObjReader, StorageAPI, }; use rustfs_filemeta::{ REPLICATE_INCOMING_DELETE, ReplicateDecision, ReplicateTargetDecision, ReplicationState, ReplicationStatusType, @@ -80,6 +80,7 @@ use rustfs_filemeta::{ version_purge_statuses_map, }; use rustfs_io_metrics; +use rustfs_lock::NamespaceLockGuard; use rustfs_notify::EventArgsBuilder; use rustfs_policy::policy::action::{Action, S3Action}; use rustfs_rio::{CompressReader, DynReader, EncryptReader, HashReader, wrap_reader}; @@ -103,7 +104,7 @@ use rustfs_utils::http::{ }, insert_str, remove_str, }; -use rustfs_utils::path::{is_dir_object, path_join_buf}; +use rustfs_utils::path::{encode_dir_object, is_dir_object, path_join_buf}; use rustfs_zip::CompressionFormat; use s3s::dto::*; use s3s::header::{X_AMZ_RESTORE, X_AMZ_RESTORE_OUTPUT_PATH}; @@ -679,6 +680,36 @@ fn should_use_existing_delete_replication_info(opts: &ObjectOptions) -> bool { opts.version_id.is_some() && !opts.delete_marker } +fn internal_object_info_lookup_opts(mut opts: ObjectOptions) -> ObjectOptions { + opts.http_preconditions = None; + opts +} + +fn copy_namespace_lock_error(bucket: &str, object: &str, mode: &'static str, err: rustfs_lock::LockError) -> StorageError { + match err { + rustfs_lock::LockError::QuorumNotReached { required, achieved } => StorageError::NamespaceLockQuorumUnavailable { + mode, + bucket: bucket.to_owned(), + object: object.to_owned(), + required, + achieved, + }, + other => StorageError::other(format!("Failed to acquire {mode} lock on {bucket}/{object}: {other}")), + } +} + +async fn acquire_self_copy_namespace_lock( + store: &S, + bucket: &str, + object: &str, +) -> S3Result { + let object = encode_dir_object(object); + let lock = store.new_ns_lock(bucket, &object).await.map_err(ApiError::from)?; + lock.get_write_lock(get_lock_acquire_timeout()) + .await + .map_err(|err| ApiError::from(copy_namespace_lock_error(bucket, &object, "write", err)).into()) +} + fn delete_replication_state_source<'a>( opts: &ObjectOptions, existing_object_info: Option<&'a ObjectInfo>, @@ -1812,9 +1843,11 @@ impl DefaultObjectUsecase { ) .await?; - let current_opts: ObjectOptions = get_opts(&bucket, &key, version_id.clone(), None, &req.headers) - .await - .map_err(ApiError::from)?; + let current_opts: ObjectOptions = internal_object_info_lookup_opts( + get_opts(&bucket, &key, version_id.clone(), None, &req.headers) + .await + .map_err(ApiError::from)?, + ); let previous_current_size = match store.get_object_info(&bucket, &key, ¤t_opts).await { Ok(existing_obj_info) => { validate_existing_object_lock_for_write(&existing_obj_info)?; @@ -2671,23 +2704,34 @@ impl DefaultObjectUsecase { ..Default::default() }; - let dst_opts = copy_dst_opts(&bucket, &key, version_id, &req.headers, HashMap::new()) + let mut dst_opts = copy_dst_opts(&bucket, &key, version_id, &req.headers, HashMap::new()) .await .map_err(ApiError::from)?; let cp_src_dst_same = path_join_buf(&[&src_bucket, &src_key]) == path_join_buf(&[&bucket, &key]); - if cp_src_dst_same { - src_get_opts.no_lock = true; - } - let Some(store) = new_object_layer_fn() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; - let current_opts: ObjectOptions = get_opts(&bucket, &key, dest_version_id.clone(), None, &req.headers) - .await - .map_err(ApiError::from)?; + let _self_copy_lock_guard = if cp_src_dst_same { + let guard = acquire_self_copy_namespace_lock(store.as_ref(), &bucket, &key).await?; + src_opts.no_lock = true; + src_get_opts.no_lock = true; + dst_opts.no_lock = true; + Some(guard) + } else { + None + }; + + let mut current_opts: ObjectOptions = internal_object_info_lookup_opts( + get_opts(&bucket, &key, dest_version_id.clone(), None, &req.headers) + .await + .map_err(ApiError::from)?, + ); + if cp_src_dst_same { + current_opts.no_lock = true; + } let previous_current_size = match store.get_object_info(&bucket, &key, ¤t_opts).await { Ok(existing_obj_info) => { validate_existing_object_lock_for_write(&existing_obj_info)?; @@ -4446,6 +4490,27 @@ mod tests { } } + #[test] + fn internal_object_info_lookup_opts_drops_http_preconditions() { + let version_id = Uuid::new_v4().to_string(); + let opts = ObjectOptions { + version_id: Some(version_id.clone()), + no_lock: true, + http_preconditions: Some(rustfs_ecstore::store_api::HTTPPreconditions { + if_none_match: Some("\"etag\"".to_string()), + if_match: Some("\"other\"".to_string()), + ..Default::default() + }), + ..Default::default() + }; + + let lookup_opts = internal_object_info_lookup_opts(opts); + + assert!(lookup_opts.http_preconditions.is_none()); + assert_eq!(lookup_opts.version_id.as_deref(), Some(version_id.as_str())); + assert!(lookup_opts.no_lock); + } + #[tokio::test] async fn build_put_like_object_lock_metadata_rejects_mode_without_retain_until_date() { let err = build_put_like_object_lock_metadata(