diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 6e853da38..69308b3d2 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -129,6 +129,42 @@ fn reduce_quorum_part_numbers(object_parts: Vec>, read_quorum: usize } impl SetDisks { + async fn acquire_multipart_upload_read_lock( + &self, + op: &'static str, + bucket: &str, + object: &str, + upload_id: &str, + opts: &ObjectOptions, + ) -> Result> { + if opts.no_lock { + return Ok(None); + } + + let upload_id_path = Self::get_upload_id_dir(bucket, object, upload_id); + self.acquire_read_lock_diag(op, RUSTFS_META_MULTIPART_BUCKET, &upload_id_path) + .await + .map(Some) + } + + async fn acquire_multipart_upload_write_lock( + &self, + op: &'static str, + bucket: &str, + object: &str, + upload_id: &str, + opts: &ObjectOptions, + ) -> Result> { + if opts.no_lock { + return Ok(None); + } + + let upload_id_path = Self::get_upload_id_dir(bucket, object, upload_id); + self.acquire_write_lock_diag(op, RUSTFS_META_MULTIPART_BUCKET, &upload_id_path) + .await + .map(Some) + } + pub(super) async fn list_parts( disks: &[Option], part_path: &str, @@ -937,9 +973,11 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { bucket: &str, object: &str, upload_id: &str, - _opts: &ObjectOptions, + opts: &ObjectOptions, ) -> Result { - // TODO: nslock + let _upload_guard = self + .acquire_multipart_upload_read_lock("get_multipart_info", bucket, object, upload_id, opts) + .await?; let (mut fi, _) = self .check_upload_id_exists(bucket, object, upload_id, false) .await @@ -958,7 +996,10 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { } #[tracing::instrument(skip(self))] - async fn abort_multipart_upload(&self, bucket: &str, object: &str, upload_id: &str, _opts: &ObjectOptions) -> Result<()> { + async fn abort_multipart_upload(&self, bucket: &str, object: &str, upload_id: &str, opts: &ObjectOptions) -> Result<()> { + let _upload_guard = self + .acquire_multipart_upload_write_lock("abort_multipart_upload", bucket, object, upload_id, opts) + .await?; self.check_upload_id_exists(bucket, object, upload_id, false).await?; let upload_id_path = Self::get_upload_id_dir(bucket, object, upload_id); @@ -1452,8 +1493,13 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { mod tests { use super::*; use crate::disk::DiskAPI as _; + use crate::disk::{endpoint::Endpoint, format::FormatV3}; use crate::set_disk::ops::object::hermetic_set_disks_support::hermetic_set_disks; + use crate::storage_api_contracts::namespace::NamespaceLocking as _; + use rustfs_lock::{LockClient, client::local::LocalClient}; + use serial_test::serial; use tempfile::TempDir; + use tokio::sync::RwLock; async fn non_trash_tmp_entries(temp_dirs: &[TempDir]) -> Vec { let mut leftovers = Vec::new(); @@ -1474,6 +1520,125 @@ mod tests { leftovers } + async fn make_multipart_lock_test_set_disks() -> Arc { + let endpoints = vec![ + Endpoint::try_from("http://127.0.0.1:9000/data").expect("first endpoint should parse"), + Endpoint::try_from("http://127.0.0.1:9001/data").expect("second endpoint should parse"), + ]; + let lockers: Vec> = vec![Arc::new(LocalClient::with_manager(Arc::new( + rustfs_lock::GlobalLockManager::new(), + )))]; + + SetDisks::new( + "test-owner".to_string(), + Arc::new(RwLock::new(vec![None, None])), + 2, + 1, + 0, + 0, + endpoints, + FormatV3::new(1, 2), + lockers, + ) + .await + } + + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn multipart_upload_read_lock_waits_for_upload_writer() { + let set_disks = make_multipart_lock_test_set_disks().await; + let upload_id_path = SetDisks::get_upload_id_dir("bucket", "object", "upload-id"); + let _writer = set_disks + .new_ns_lock(RUSTFS_META_MULTIPART_BUCKET, &upload_id_path) + .await + .expect("namespace lock should be created") + .get_write_lock(Duration::from_secs(1)) + .await + .expect("outer upload write lock should be acquired"); + + let result = temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], async { + set_disks + .acquire_multipart_upload_read_lock( + "test_multipart_upload_read_lock", + "bucket", + "object", + "upload-id", + &ObjectOptions::default(), + ) + .await + }) + .await; + let err = match result { + Ok(_) => panic!("read lock must wait behind an upload write lock"), + Err(err) => err, + }; + + assert!(matches!(err, StorageError::Lock(rustfs_lock::LockError::Timeout { .. }))); + + let no_lock_guard = set_disks + .acquire_multipart_upload_read_lock( + "test_multipart_upload_read_lock_no_lock", + "bucket", + "object", + "upload-id", + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + .expect("no_lock read path should not acquire an upload lock"); + assert!(no_lock_guard.is_none()); + } + + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn multipart_upload_write_lock_waits_for_upload_reader() { + let set_disks = make_multipart_lock_test_set_disks().await; + let upload_id_path = SetDisks::get_upload_id_dir("bucket", "object", "upload-id"); + let _reader = set_disks + .new_ns_lock(RUSTFS_META_MULTIPART_BUCKET, &upload_id_path) + .await + .expect("namespace lock should be created") + .get_read_lock(Duration::from_secs(1)) + .await + .expect("outer upload read lock should be acquired"); + + let result = temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], async { + set_disks + .acquire_multipart_upload_write_lock( + "test_multipart_upload_write_lock", + "bucket", + "object", + "upload-id", + &ObjectOptions::default(), + ) + .await + }) + .await; + let err = match result { + Ok(_) => panic!("write lock must wait behind an upload read lock"), + Err(err) => err, + }; + + assert!(matches!(err, StorageError::Lock(rustfs_lock::LockError::Timeout { .. }))); + + let no_lock_guard = set_disks + .acquire_multipart_upload_write_lock( + "test_multipart_upload_write_lock_no_lock", + "bucket", + "object", + "upload-id", + &ObjectOptions { + no_lock: true, + ..Default::default() + }, + ) + .await + .expect("no_lock write path should not acquire an upload lock"); + assert!(no_lock_guard.is_none()); + } + #[tokio::test] async fn collect_list_parts_results_fails_early_when_quorum_is_impossible() { let started = Instant::now();