diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 60d8aa41c..4261073cf 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -370,254 +370,6 @@ struct SetDiskLockGuardedReader { guard: Option, } -#[derive(Clone)] -struct SnapshotLease { - disk: DiskStore, - token: SnapshotLeaseToken, -} - -struct SnapshotLeaseState { - volume: Arc, - path: Arc, - leases: parking_lot::Mutex>, - renewal_failed: AtomicBool, - renewal_waker: AtomicWaker, - cancel: CancellationToken, - runtime: tokio::runtime::Handle, -} - -impl SnapshotLeaseState { - fn apply_renewal_results(&self, results: Vec>) -> bool { - let mut leases = self.leases.lock(); - let mut failed = false; - for (lease, result) in leases.iter_mut().zip(results) { - match result { - Ok(token) => lease.token = token, - Err(_) => failed = true, - } - } - drop(leases); - if failed { - self.renewal_failed.store(true, Ordering::Release); - self.renewal_waker.wake(); - } - failed - } -} - -impl Drop for SnapshotLeaseState { - fn drop(&mut self) { - self.cancel.cancel(); - let leases = mem::take(&mut *self.leases.lock()); - if leases.is_empty() { - return; - } - let volume = Arc::clone(&self.volume); - let path = Arc::clone(&self.path); - self.runtime.spawn(async move { - join_all(leases.into_iter().map(|lease| { - let volume = Arc::clone(&volume); - let path = Arc::clone(&path); - async move { lease.disk.release_snapshot_lease(&volume, &path, lease.token).await } - })) - .await; - }); - } -} - -#[derive(Clone)] -struct SnapshotLeaseHandle(Arc); - -impl SnapshotLeaseHandle { - fn new(volume: Arc, path: Arc, leases: Vec) -> Self { - let state = Arc::new(SnapshotLeaseState { - volume, - path, - leases: parking_lot::Mutex::new(leases), - renewal_failed: AtomicBool::new(false), - renewal_waker: AtomicWaker::new(), - cancel: CancellationToken::new(), - runtime: tokio::runtime::Handle::current(), - }); - let weak = Arc::downgrade(&state); - let cancel = state.cancel.clone(); - tokio::spawn(async move { - let renew_interval = crate::cluster::rpc::remote_disk::REMOTE_SNAPSHOT_LEASE_TTL / 3; - loop { - tokio::select! { - _ = cancel.cancelled() => return, - _ = tokio::time::sleep(renew_interval) => {} - } - let Some(state) = weak.upgrade() else { - return; - }; - let leases = state.leases.lock().clone(); - let results = renew_snapshot_leases_with_timeout( - &leases, - Arc::clone(&state.volume), - Arc::clone(&state.path), - renew_interval, - |disk, volume, path, token| async move { disk.renew_snapshot_lease(&volume, &path, token).await }, - ) - .await; - if state.apply_renewal_results(results) { - return; - } - } - }); - Self(state) - } - - fn renewal_failed(&self) -> bool { - self.0.renewal_failed.load(Ordering::Acquire) - } -} - -async fn renew_snapshot_leases_with_timeout( - leases: &[SnapshotLease], - volume: Arc, - path: Arc, - renewal_timeout: Duration, - renew: F, -) -> Vec> -where - F: Fn(DiskStore, Arc, Arc, SnapshotLeaseToken) -> Fut, - Fut: Future>, -{ - join_all(leases.iter().map(|lease| { - let renew = renew(Arc::clone(&lease.disk), Arc::clone(&volume), Arc::clone(&path), lease.token); - async move { - match timeout(renewal_timeout, renew).await { - Ok(result) => result, - Err(_) => Err(DiskError::Timeout), - } - } - })) - .await -} - -struct SnapshotLeaseReader { - inner: Box, - lease: Option, - terminal_error: bool, -} - -impl AsyncRead for SnapshotLeaseReader { - fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { - if self.terminal_error { - return Poll::Ready(Err(std::io::Error::other("snapshot lease renewal failed"))); - } - let Some(lease) = self.lease.as_ref() else { - return Pin::new(&mut self.inner).poll_read(cx, buf); - }; - if lease.renewal_failed() { - self.lease.take(); - self.terminal_error = true; - return Poll::Ready(Err(std::io::Error::other("snapshot lease renewal failed"))); - } - lease.0.renewal_waker.register(cx.waker()); - if lease.renewal_failed() { - self.lease.take(); - self.terminal_error = true; - return Poll::Ready(Err(std::io::Error::other("snapshot lease renewal failed"))); - } - let filled_before = buf.filled().len(); - let poll = Pin::new(&mut self.inner).poll_read(cx, buf); - if matches!(poll, Poll::Ready(Err(_))) - || matches!(poll, Poll::Ready(Ok(())) if buf.filled().len() == filled_before && buf.remaining() > 0) - { - self.lease.take(); - } - poll - } -} - -async fn acquire_snapshot_leases( - disks: &[Option], - volume: &str, - path: &str, - read_quorum: usize, -) -> Option { - acquire_snapshot_leases_with_timeout( - disks, - volume, - path, - read_quorum, - crate::disk::disk_store::get_drive_metadata_timeout(), - |disk, volume, path| async move { disk.acquire_snapshot_lease(&volume, &path).await }, - ) - .await -} - -async fn acquire_snapshot_leases_with_timeout( - disks: &[Option], - volume: &str, - path: &str, - read_quorum: usize, - candidate_timeout: Duration, - acquire: F, -) -> Option -where - F: Fn(DiskStore, Arc, Arc) -> Fut, - Fut: Future> + Send + 'static, -{ - let candidates = disks.iter().cloned().collect::>>()?; - if candidates.len() < read_quorum { - return None; - } - let volume: Arc = Arc::from(volume); - let path: Arc = Arc::from(path); - let results = join_all(candidates.iter().cloned().map(|disk| { - let volume = Arc::clone(&volume); - let path = Arc::clone(&path); - let acquire_disk = Arc::clone(&disk); - let acquire = acquire(acquire_disk, Arc::clone(&volume), Arc::clone(&path)); - async move { - let mut task = tokio::spawn(acquire); - match timeout(candidate_timeout, &mut task).await { - Ok(Ok(result)) => result, - Ok(Err(err)) => Err(DiskError::Io(std::io::Error::other(format!( - "snapshot lease acquisition task failed: {err}" - )))), - Err(_) => { - tokio::spawn(async move { - match timeout(crate::cluster::rpc::remote_disk::REMOTE_SNAPSHOT_LEASE_TTL, &mut task).await { - Ok(Ok(Ok(token))) => { - let _ = disk.release_snapshot_lease(&volume, &path, token).await; - } - Ok(_) => {} - Err(_) => { - task.abort(); - let _ = task.await; - } - } - }); - Err(DiskError::Timeout) - } - } - } - })) - .await; - let mut leases = Vec::with_capacity(candidates.len()); - let mut failed = false; - for (disk, result) in candidates.into_iter().zip(results) { - match result { - Ok(token) => leases.push(SnapshotLease { disk, token }), - Err(_) => failed = true, - } - } - if failed || leases.len() < read_quorum { - join_all(leases.into_iter().map(|lease| { - let volume = Arc::clone(&volume); - let path = Arc::clone(&path); - async move { lease.disk.release_snapshot_lease(&volume, &path, lease.token).await } - })) - .await; - return None; - } - Some(SnapshotLeaseHandle::new(volume, path, leases)) -} - impl AsyncRead for SetDiskLockGuardedReader { fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { let had_capacity = buf.remaining() > 0; @@ -633,7 +385,6 @@ impl AsyncRead for SetDiskLockGuardedReader { fn finish_set_disk_read_lock( mut reader: GetObjectReader, read_lock_guard: Option, - snapshot_lease: Option, bucket: &str, object: &str, ) -> GetObjectReader { @@ -642,16 +393,6 @@ fn finish_set_disk_read_lock( return reader; } - if let Some(lease) = snapshot_lease { - release_materialized_read_lock(bucket, object, read_lock_guard); - reader.stream = Box::new(SnapshotLeaseReader { - inner: reader.stream, - lease: Some(lease), - terminal_error: false, - }); - return reader; - } - if let Some(guard) = read_lock_guard { reader.stream = Box::new(SetDiskLockGuardedReader { inner: reader.stream, @@ -1294,8 +1035,7 @@ pub fn get_object_lock_diag_slow_hold_threshold() -> Duration { /// Check if lock optimization is enabled. /// Fully materialized reads release the read lock before returning. Streaming -/// reads may replace it with data-directory snapshot leases when every -/// candidate disk supports the lease protocol. +/// reads retain it until the response body completes or is dropped. /// /// **Note**: Cached via `OnceLock` in production — env var changes require /// process restart. In test builds the env var is read directly so that @@ -5617,6 +5357,101 @@ mod tests { ); } + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn streaming_reader_holds_read_lock_until_eof() { + 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 read_guard = set_disks + .new_ns_lock("bucket", "object") + .await + .expect("namespace lock should be created") + .get_read_lock(Duration::from_secs(1)) + .await + .expect("read lock should be acquired"); + let read_guard = ObjectLockDiagGuard::new( + read_guard, + false, + "GetObject", + Some("bucket".to_owned()), + Some("object".to_owned()), + None, + "read", + ); + let mut reader = finish_set_disk_read_lock( + GetObjectReader { + stream: Box::new(Cursor::new(b"body")), + object_info: ObjectInfo::default(), + buffered_body: None, + body_source: Default::default(), + }, + Some(read_guard), + "bucket", + "object", + ); + + let blocked_write = set_disks + .new_ns_lock("bucket", "object") + .await + .expect("namespace lock should be created") + .get_write_lock(Duration::from_millis(100)) + .await; + assert!(blocked_write.is_err(), "a stalled response must block an overwrite"); + + let mut body = Vec::new(); + reader.stream.read_to_end(&mut body).await.expect("stream should reach EOF"); + assert_eq!(body, b"body"); + + let write_guard = set_disks + .new_ns_lock("bucket", "object") + .await + .expect("namespace lock should be created") + .get_write_lock(Duration::from_secs(1)) + .await; + assert!(write_guard.is_ok(), "the overwrite should proceed after EOF releases the read lock"); + drop(write_guard); + + let read_guard = set_disks + .new_ns_lock("bucket", "object") + .await + .expect("namespace lock should be created") + .get_read_lock(Duration::from_secs(1)) + .await + .expect("second read lock should be acquired"); + let reader = finish_set_disk_read_lock( + GetObjectReader { + stream: Box::new(tokio::io::empty()), + object_info: ObjectInfo::default(), + buffered_body: None, + body_source: Default::default(), + }, + Some(ObjectLockDiagGuard::new( + read_guard, + false, + "GetObject", + Some("bucket".to_owned()), + Some("object".to_owned()), + None, + "read", + )), + "bucket", + "object", + ); + drop(reader); + + let write_after_drop = set_disks + .new_ns_lock("bucket", "object") + .await + .expect("namespace lock should be created") + .get_write_lock(Duration::from_secs(1)) + .await; + assert!(write_after_drop.is_ok(), "dropping the response must release the read lock"); + } + #[tokio::test(flavor = "multi_thread")] #[serial] async fn copy_object_honors_no_lock_when_outer_write_lock_is_held() { @@ -5866,400 +5701,6 @@ mod tests { (dir, disk) } - #[tokio::test] - async fn snapshot_lease_acquisition_is_all_or_nothing() { - let (_dir1, disk1) = make_single_local_disk().await; - let (_dir2, disk2) = make_single_local_disk().await; - let bucket = "snapshot-lease-acquire"; - let data_dir = "object/11111111-1111-1111-1111-111111111111"; - let part = format!("{data_dir}/part.1"); - disk1.make_volume(bucket).await.expect("first volume should be created"); - disk2.make_volume(bucket).await.expect("second volume should be created"); - disk1 - .write_all(bucket, &part, Bytes::from_static(b"shard")) - .await - .expect("first shard should be written"); - - let lease = acquire_snapshot_leases(&[Some(disk1.clone()), Some(disk2)], bucket, data_dir, 1).await; - assert!(lease.is_none(), "one candidate missing snapshot data must retain the namespace lock"); - assert_eq!( - disk1 - .delete_data_dir( - bucket, - data_dir, - DeleteOptions { - recursive: true, - ..Default::default() - }, - ) - .await - .expect("released partial lease should not defer cleanup"), - crate::disk::DataDirDeleteStatus::Deleted - ); - } - - #[tokio::test] - async fn snapshot_lease_acquisition_rejects_unavailable_candidate() { - let (_dir, disk) = make_single_local_disk().await; - let bucket = "snapshot-lease-unavailable"; - let data_dir = "object/11111111-1111-1111-1111-111111111111"; - let part = format!("{data_dir}/part.1"); - disk.make_volume(bucket).await.expect("volume should be created"); - disk.write_all(bucket, &part, Bytes::from_static(b"shard")) - .await - .expect("shard should be written"); - - let lease = acquire_snapshot_leases(&[Some(disk.clone()), None], bucket, data_dir, 1).await; - assert!(lease.is_none(), "one unavailable candidate must retain the namespace lock"); - assert_eq!( - disk.delete_data_dir( - bucket, - data_dir, - DeleteOptions { - recursive: true, - ..Default::default() - }, - ) - .await - .expect("unavailable candidate fallback must not leave a lease"), - crate::disk::DataDirDeleteStatus::Deleted - ); - } - - #[tokio::test(start_paused = true)] - async fn snapshot_lease_acquisition_times_out_one_candidate_and_releases_partial_success() { - let (_dir1, disk1) = make_single_local_disk().await; - let (_dir2, disk2) = make_single_local_disk().await; - let bucket = "snapshot-lease-timeout"; - let data_dir = "object/11111111-1111-1111-1111-111111111111"; - let part = format!("{data_dir}/part.1"); - for disk in [&disk1, &disk2] { - disk.make_volume(bucket).await.expect("volume should be created"); - disk.write_all(bucket, &part, Bytes::from_static(b"shard")) - .await - .expect("shard should be written"); - } - let slow_endpoint = disk2.endpoint(); - let release_slow_candidate = Arc::new(tokio::sync::Notify::new()); - let slow_candidate_gate = Arc::clone(&release_slow_candidate); - - let lease = acquire_snapshot_leases_with_timeout( - &[Some(disk1.clone()), Some(disk2.clone())], - bucket, - data_dir, - 1, - Duration::from_millis(10), - move |disk, volume, path| { - let slow = disk.endpoint() == slow_endpoint; - let release_slow_candidate = Arc::clone(&slow_candidate_gate); - async move { - if slow { - release_slow_candidate.notified().await; - } - disk.acquire_snapshot_lease(&volume, &path).await - } - }, - ) - .await; - - assert!(lease.is_none(), "a timed-out candidate must retain the namespace lock"); - assert_eq!( - disk1 - .delete_data_dir( - bucket, - data_dir, - DeleteOptions { - recursive: true, - ..Default::default() - }, - ) - .await - .expect("the successful candidate lease must be released"), - crate::disk::DataDirDeleteStatus::Deleted - ); - release_slow_candidate.notify_one(); - timeout(Duration::from_secs(1), async { - loop { - match disk2 - .delete_data_dir( - bucket, - data_dir, - DeleteOptions { - recursive: true, - ..Default::default() - }, - ) - .await - .expect("a token acquired after the deadline must be released") - { - crate::disk::DataDirDeleteStatus::Deleted => return, - crate::disk::DataDirDeleteStatus::Deferred => tokio::time::sleep(Duration::from_millis(1)).await, - } - } - }) - .await - .expect("late snapshot lease cleanup must finish within the bounded wait"); - } - - #[tokio::test(start_paused = true)] - async fn snapshot_lease_acquisition_aborts_permanently_pending_cleanup_at_ttl() { - struct DropProbe(Arc); - - impl Drop for DropProbe { - fn drop(&mut self) { - self.0.store(true, Ordering::Release); - } - } - - let (_dir, disk) = make_single_local_disk().await; - let dropped = Arc::new(AtomicBool::new(false)); - let acquire_probe = Arc::clone(&dropped); - let lease = acquire_snapshot_leases_with_timeout( - &[Some(disk)], - "snapshot-lease-pending-cleanup", - "object/11111111-1111-1111-1111-111111111111", - 1, - Duration::from_millis(10), - move |_disk, _volume, _path| { - let probe = DropProbe(Arc::clone(&acquire_probe)); - async move { - let _probe = probe; - futures::future::pending().await - } - }, - ) - .await; - assert!(lease.is_none(), "a pending candidate must retain the namespace lock"); - assert!( - !dropped.load(Ordering::Acquire), - "the late cleanup must initially retain the acquire task" - ); - - tokio::task::yield_now().await; - tokio::time::advance(crate::cluster::rpc::remote_disk::REMOTE_SNAPSHOT_LEASE_TTL).await; - for _ in 0..10 { - if dropped.load(Ordering::Acquire) { - break; - } - tokio::task::yield_now().await; - } - assert!( - dropped.load(Ordering::Acquire), - "the TTL fallback must abort and drop a permanently pending acquire task" - ); - } - - #[tokio::test] - async fn snapshot_lease_acquisition_fails_closed_for_an_old_peer() { - let (_dir1, disk1) = make_single_local_disk().await; - let (_dir2, disk2) = make_single_local_disk().await; - let bucket = "snapshot-lease-old-peer"; - let data_dir = "object/11111111-1111-1111-1111-111111111111"; - let part = format!("{data_dir}/part.1"); - for disk in [&disk1, &disk2] { - disk.make_volume(bucket).await.expect("volume should be created"); - disk.write_all(bucket, &part, Bytes::from_static(b"shard")) - .await - .expect("shard should be written"); - } - let old_peer_endpoint = disk2.endpoint(); - - let lease = acquire_snapshot_leases_with_timeout( - &[Some(disk1.clone()), Some(disk2)], - bucket, - data_dir, - 1, - Duration::from_secs(1), - move |disk, volume, path| { - let old_peer = disk.endpoint() == old_peer_endpoint; - async move { - if old_peer { - return Err(DiskError::MethodNotAllowed); - } - disk.acquire_snapshot_lease(&volume, &path).await - } - }, - ) - .await; - - assert!(lease.is_none(), "an old peer without the lease RPC must retain the namespace lock"); - assert_eq!( - disk1 - .delete_data_dir( - bucket, - data_dir, - DeleteOptions { - recursive: true, - ..Default::default() - }, - ) - .await - .expect("the compatible peer lease must be released"), - crate::disk::DataDirDeleteStatus::Deleted - ); - } - - #[tokio::test] - async fn snapshot_lease_reader_fails_when_renewal_fails() { - let state = Arc::new(SnapshotLeaseState { - volume: Arc::from("bucket"), - path: Arc::from("object/data-dir"), - leases: parking_lot::Mutex::new(Vec::new()), - renewal_failed: AtomicBool::new(true), - renewal_waker: AtomicWaker::new(), - cancel: CancellationToken::new(), - runtime: tokio::runtime::Handle::current(), - }); - let mut reader = SnapshotLeaseReader { - inner: Box::new(Cursor::new(Bytes::from_static(b"payload"))), - lease: Some(SnapshotLeaseHandle(state)), - terminal_error: false, - }; - let mut body = Vec::new(); - let err = reader - .read_to_end(&mut body) - .await - .expect_err("renewal failure must terminate the body"); - assert_eq!(err.kind(), std::io::ErrorKind::Other); - assert!(body.is_empty()); - let mut second = [0; 1]; - let second_err = reader - .read(&mut second) - .await - .expect_err("renewal failure must remain terminal on later polls"); - assert_eq!(second_err.kind(), std::io::ErrorKind::Other); - } - - #[tokio::test(start_paused = true)] - async fn snapshot_lease_pending_renewal_hits_deadline_and_terminates_reader() { - let (_dir, disk) = make_single_local_disk().await; - let bucket = "snapshot-lease-renewal-timeout"; - let data_dir = "object/11111111-1111-1111-1111-111111111111"; - let part = format!("{data_dir}/part.1"); - disk.make_volume(bucket).await.expect("volume should be created"); - disk.write_all(bucket, &part, Bytes::from_static(b"shard")) - .await - .expect("shard should be written"); - let token = disk - .acquire_snapshot_lease(bucket, data_dir) - .await - .expect("candidate lease should be acquired"); - let state = Arc::new(SnapshotLeaseState { - volume: Arc::from(bucket), - path: Arc::from(data_dir), - leases: parking_lot::Mutex::new(vec![SnapshotLease { disk, token }]), - renewal_failed: AtomicBool::new(false), - renewal_waker: AtomicWaker::new(), - cancel: CancellationToken::new(), - runtime: tokio::runtime::Handle::current(), - }); - let (pending_reader, _pending_writer) = tokio::io::duplex(1); - let reader_state = Arc::clone(&state); - let read = tokio::spawn(async move { - let mut reader = SnapshotLeaseReader { - inner: Box::new(pending_reader), - lease: Some(SnapshotLeaseHandle(reader_state)), - terminal_error: false, - }; - let mut byte = [0; 1]; - let first = reader.read(&mut byte).await; - let second = reader.read(&mut byte).await; - (first, second) - }); - tokio::task::yield_now().await; - assert!(!read.is_finished(), "the body must be pending before the renewal deadline"); - - let leases = state.leases.lock().clone(); - let results = renew_snapshot_leases_with_timeout( - &leases, - Arc::clone(&state.volume), - Arc::clone(&state.path), - crate::cluster::rpc::remote_disk::REMOTE_SNAPSHOT_LEASE_TTL / 3, - |_disk, _volume, _path, _token| futures::future::pending(), - ) - .await; - assert!( - state.apply_renewal_results(results), - "a pending renewal must fail at the bounded deadline" - ); - - let (first, second) = read.await.expect("the renewal wake must resume the body task"); - let first = first.expect_err("the deadline must terminate the reader before it emits data"); - let second = second.expect_err("the deadline failure must remain terminal"); - assert_eq!(first.kind(), std::io::ErrorKind::Other); - assert_eq!(second.kind(), std::io::ErrorKind::Other); - } - - #[tokio::test(start_paused = true)] - async fn snapshot_lease_renewal_partial_failure_releases_renewed_tokens() { - let (_dir1, disk1) = make_single_local_disk().await; - let (_dir2, disk2) = make_single_local_disk().await; - let bucket = "snapshot-lease-renewal"; - let data_dir = "object/11111111-1111-1111-1111-111111111111"; - let part = format!("{data_dir}/part.1"); - for disk in [&disk1, &disk2] { - disk.make_volume(bucket).await.expect("volume should be created"); - disk.write_all(bucket, &part, Bytes::from_static(b"shard")) - .await - .expect("shard should be written"); - } - let first = disk1 - .acquire_snapshot_lease(bucket, data_dir) - .await - .expect("first candidate lease should be acquired"); - let second = disk2 - .acquire_snapshot_lease(bucket, data_dir) - .await - .expect("second candidate lease should be acquired"); - let lease = SnapshotLeaseHandle::new( - Arc::from(bucket), - Arc::from(data_dir), - vec![ - SnapshotLease { - disk: disk1.clone(), - token: first, - }, - SnapshotLease { - disk: disk2.clone(), - token: second, - }, - ], - ); - disk2 - .release_snapshot_lease(bucket, data_dir, second) - .await - .expect("removing one token should force a real renewal failure"); - - tokio::task::yield_now().await; - tokio::time::advance(crate::cluster::rpc::remote_disk::REMOTE_SNAPSHOT_LEASE_TTL / 3).await; - for _ in 0..10 { - if lease.renewal_failed() { - break; - } - tokio::task::yield_now().await; - } - assert!(lease.renewal_failed(), "one failed renewal must terminate the lease set"); - - drop(lease); - for _ in 0..10 { - tokio::task::yield_now().await; - } - assert_eq!( - disk1 - .delete_data_dir( - bucket, - data_dir, - DeleteOptions { - recursive: true, - ..Default::default() - }, - ) - .await - .expect("the successfully renewed token must be released"), - crate::disk::DataDirDeleteStatus::Deleted - ); - } - async fn make_set_disks_with(disks: Vec>) -> Arc { let drive_count = disks.len(); let endpoints = (0..drive_count) @@ -10111,7 +9552,7 @@ mod tests { #[tokio::test(flavor = "multi_thread")] #[serial] - async fn streaming_get_snapshot_survives_concurrent_overwrite() { + async fn streaming_get_blocks_concurrent_overwrite_until_eof() { temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some("true"))], async { let set_disks = make_local_bucket_test_set_disks().await; let bucket = "snapshot-streaming-overwrite"; @@ -10137,23 +9578,29 @@ mod tests { let overwrite_set = Arc::clone(&set_disks); let overwrite_body = new_body.clone(); let overwrite_opts = opts.clone(); - let overwrite = tokio::spawn(async move { + let mut overwrite = tokio::spawn(async move { let mut reader = PutObjReader::from_vec(overwrite_body); overwrite_set.put_object(bucket, object, &mut reader, &overwrite_opts).await }); - tokio::time::timeout(Duration::from_secs(5), overwrite) - .await - .expect("overwrite should not wait for the response body") - .expect("overwrite task should join") - .expect("overwrite should succeed"); + assert!( + tokio::time::timeout(Duration::from_millis(100), &mut overwrite) + .await + .is_err(), + "overwrite must wait while the response body retains the read lock" + ); let mut restored = Vec::new(); snapshot .stream .read_to_end(&mut restored) .await - .expect("leased snapshot should remain readable"); + .expect("stream should remain readable"); assert_eq!(restored, old_body); + tokio::time::timeout(Duration::from_secs(5), overwrite) + .await + .expect("overwrite should proceed after EOF") + .expect("overwrite task should join") + .expect("overwrite should succeed"); let mut latest = set_disks .get_object_reader(bucket, object, None, HeaderMap::new(), &opts) @@ -10175,7 +9622,7 @@ mod tests { let mut replacement = PutObjReader::from_vec(vec![0x43; 2 * 1024 * 1024]); tokio::time::timeout(Duration::from_secs(5), set_disks.put_object(bucket, object, &mut replacement, &opts)) .await - .expect("reader drop must release its snapshot") + .expect("reader drop must release its read lock") .expect("replacement after cancellation should succeed"); }) .await; @@ -10183,7 +9630,7 @@ mod tests { #[tokio::test(flavor = "multi_thread")] #[serial] - async fn streaming_get_snapshot_survives_concurrent_delete() { + async fn streaming_get_blocks_concurrent_delete_until_eof() { temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some("true"))], async { let set_disks = make_local_bucket_test_set_disks().await; let bucket = "snapshot-streaming-delete"; @@ -10212,20 +9659,24 @@ mod tests { .expect("snapshot reader should open"); let delete_set = Arc::clone(&set_disks); let delete_opts = opts.clone(); - let delete = tokio::spawn(async move { delete_set.delete_object(bucket, object, delete_opts).await }); - tokio::time::timeout(Duration::from_secs(30), delete) - .await - .expect("delete should not wait for the response body") - .expect("delete task should join") - .expect("delete should succeed"); + let mut delete = tokio::spawn(async move { delete_set.delete_object(bucket, object, delete_opts).await }); + assert!( + tokio::time::timeout(Duration::from_millis(100), &mut delete).await.is_err(), + "delete must wait while the response body retains the read lock" + ); let mut restored = Vec::new(); snapshot .stream .read_to_end(&mut restored) .await - .expect("leased snapshot should remain readable after delete"); + .expect("stream should remain readable before delete"); assert_eq!(restored, body); + tokio::time::timeout(Duration::from_secs(30), delete) + .await + .expect("delete should proceed after EOF") + .expect("delete task should join") + .expect("delete should succeed"); let err = match set_disks .get_object_reader(bucket, object, None, HeaderMap::new(), &opts) .await @@ -10240,7 +9691,7 @@ mod tests { #[tokio::test(flavor = "multi_thread")] #[serial] - async fn streaming_get_snapshot_survives_concurrent_delete_objects() { + async fn streaming_get_blocks_concurrent_delete_objects_until_eof() { temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some("true"))], async { let set_disks = make_local_bucket_test_set_disks().await; let bucket = "snapshot-streaming-delete-objects"; @@ -10269,7 +9720,7 @@ mod tests { .expect("snapshot reader should open"); let delete_set = Arc::clone(&set_disks); let delete_opts = opts.clone(); - let delete = tokio::spawn(async move { + let mut delete = tokio::spawn(async move { delete_set .delete_objects( bucket, @@ -10281,19 +9732,23 @@ mod tests { ) .await }); - let (_, errors) = tokio::time::timeout(Duration::from_secs(30), delete) - .await - .expect("batch delete should not wait for the response body") - .expect("batch delete task should join"); - assert!(errors.iter().all(Option::is_none)); + assert!( + tokio::time::timeout(Duration::from_millis(100), &mut delete).await.is_err(), + "batch delete must wait while the response body retains the read lock" + ); let mut restored = Vec::new(); snapshot .stream .read_to_end(&mut restored) .await - .expect("leased snapshot should remain readable after batch delete"); + .expect("stream should remain readable before batch delete"); assert_eq!(restored, body); + let (_, errors) = tokio::time::timeout(Duration::from_secs(30), delete) + .await + .expect("batch delete should proceed after EOF") + .expect("batch delete task should join"); + assert!(errors.iter().all(Option::is_none)); }) .await; } diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 2063e5f25..de44edcd0 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -642,7 +642,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { &self.ctx.tier_config_mgr(), ) .await?; - return Ok(finish_set_disk_read_lock(gr, read_lock_guard.take(), None, bucket, object)); + return Ok(finish_set_disk_read_lock(gr, read_lock_guard.take(), bucket, object)); } // App-layer object data cache probe: metadata (etag/size) is resolved @@ -769,18 +769,6 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { return Ok(reader); } - let snapshot_lease = if lock_optimization_enabled { - match fi.data_dir.filter(|data_dir| !data_dir.is_nil()) { - Some(data_dir) => { - let data_dir_path = format!("{object}/{data_dir}"); - acquire_snapshot_leases(&disks, bucket, &data_dir_path, fi.erasure.data_blocks).await - } - None => None, - } - } else { - None - }; - match codec_streaming_gate.decision { GetCodecStreamingDecision::Use => { match Self::get_object_decode_reader_with_fileinfo( @@ -810,7 +798,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { // Carry the hook probe result so the app layer skips its // now-redundant lookup on the streaming miss path (ODC-16). reader.body_source = body_source; - return Ok(finish_set_disk_read_lock(reader, read_lock_guard.take(), snapshot_lease, bucket, object)); + return Ok(finish_set_disk_read_lock(reader, read_lock_guard.take(), bucket, object)); } core::io_primitives::GetCodecStreamingReaderBuildOutcome::Fallback(reason) => { record_get_codec_streaming_gate_decision( @@ -850,22 +838,8 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { let set_index = self.set_index; let pool_index = self.pool_index; let skip_verify = opts.skip_verify_bitrot; - let producer_snapshot_lease = snapshot_lease.clone(); - if let Some(lease) = snapshot_lease { - release_materialized_read_lock(&bucket, &object, read_lock_guard.take()); - reader.stream = Box::new(SnapshotLeaseReader { - inner: reader.stream, - lease: Some(lease), - terminal_error: false, - }); - debug!(bucket, object, "Lock optimization: replaced read lock with snapshot leases"); - } - - // The producer shares the lease lifetime with the body so cancellation - // cannot release the snapshot while the duplex task is still unwinding. tokio::spawn(async move { let _guard = read_lock_guard; - let _snapshot_lease = producer_snapshot_lease; let mut writer = GetObjectDownstreamWriter::new(wd); // Do not wrap the entire read+write pipeline in `disk_read_timeout`. // `get_object_with_fileinfo` also waits on `writer`, so an outer timeout diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index 26f600191..9a572a707 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -380,7 +380,6 @@ mod metrics; pub struct NodeService { local_peer: LocalPeerS3Client, context: Option>, - snapshot_lease_expiry: disk::SnapshotLeaseExpiryScheduler, } impl std::fmt::Debug for NodeService { @@ -399,11 +398,7 @@ pub fn make_server() -> NodeService { pub fn make_server_for_context(context: Option>) -> NodeService { let local_peer = LocalPeerS3Client::new(None, None); - NodeService { - local_peer, - context, - snapshot_lease_expiry: disk::SnapshotLeaseExpiryScheduler::new(), - } + NodeService { local_peer, context } } #[derive(Clone, Debug, Default)] @@ -2170,7 +2165,7 @@ mod tests { previous_scanner_activity_response, remove_heal_control_replay, scanner_activity_response, stop_rebalance_response, }; use crate::storage::rpc::node_service::heal::heal_topology_fingerprint; - use crate::storage::storage_api::rpc_consumer::node_service::{HealBucketInfo, HealEndpoint}; + use crate::storage::storage_api::rpc_consumer::node_service::{DiskError, HealBucketInfo, HealEndpoint}; use crate::storage::storage_api::set_tonic_canonical_body_digest; use crate::storage::storage_api::{ Endpoint, @@ -2910,6 +2905,52 @@ mod tests { ); } + #[tokio::test] + async fn snapshot_lease_acquire_and_renew_handlers_fail_closed() { + let service = make_server(); + let disk = "http://node-a:9000/data/rustfs0".to_string(); + + let mut acquire = Request::new(SnapshotLeaseRequest { + disk: disk.clone(), + volume: "v".into(), + path: "p".into(), + ttl_ms: 60_000, + }); + let acquire_body = + rustfs_protos::canonical_snapshot_lease_request_body(acquire.get_ref()).expect("acquire request body should encode"); + set_tonic_canonical_body_digest(&mut acquire, &acquire_body).expect("acquire digest metadata should encode"); + mark_v2_authenticated(&mut acquire); + let acquire = service + .acquire_snapshot_lease(acquire) + .await + .expect("disabled acquire should return a protocol response") + .into_inner(); + + let mut renew = Request::new(SnapshotLeaseRenewRequest { + disk, + volume: "v".into(), + path: "p".into(), + token: vec![1; 16].into(), + ttl_ms: 60_000, + }); + let renew_body = rustfs_protos::canonical_snapshot_lease_renew_request_body(renew.get_ref()) + .expect("renew request body should encode"); + set_tonic_canonical_body_digest(&mut renew, &renew_body).expect("renew digest metadata should encode"); + mark_v2_authenticated(&mut renew); + let renew = service + .renew_snapshot_lease(renew) + .await + .expect("disabled renew should return a protocol response") + .into_inner(); + + for response in [acquire, renew] { + assert!(!response.success); + assert!(response.token.is_empty()); + assert_eq!(response.protocol_version, 1); + assert_eq!(response.error, Some(DiskError::UnsupportedDisk.into())); + } + } + #[tokio::test] async fn heal_control_requires_body_bound_auth_before_topology_validation() { let service = make_heal_control_server(); diff --git a/rustfs/src/storage/rpc/node_service/disk.rs b/rustfs/src/storage/rpc/node_service/disk.rs index 31d63bd97..75a16a45f 100644 --- a/rustfs/src/storage/rpc/node_service/disk.rs +++ b/rustfs/src/storage/rpc/node_service/disk.rs @@ -14,9 +14,8 @@ use super::NodeService; use crate::storage::storage_api::rpc_consumer::node_service::{ - BatchReadVersionReq, BatchReadVersionResp, DeleteOptions, DiskError, DiskInfoOptions, DiskStore, FileInfoVersions, - ReadMultipleReq, ReadMultipleResp, ReadOptions, StorageDiskRpcExt as _, UpdateMetadataOpts, - validate_batch_read_version_item_count, + BatchReadVersionReq, BatchReadVersionResp, DeleteOptions, DiskError, DiskInfoOptions, FileInfoVersions, ReadMultipleReq, + ReadMultipleResp, ReadOptions, StorageDiskRpcExt as _, UpdateMetadataOpts, validate_batch_read_version_item_count, }; use crate::storage::storage_api::runtime_sources_consumer::runtime_sources; use crate::storage::storage_api::{PartTransactionAction, SnapshotLeaseToken, verify_tonic_mutation_body_digest}; @@ -29,9 +28,7 @@ use rustfs_io_metrics::internode_metrics::{ }; use rustfs_protos::proto_gen::node_service::*; use serde::de::DeserializeOwned; -use std::{collections::HashMap, io::Cursor, time::Duration}; -use tokio::sync::mpsc; -use tokio_util::time::DelayQueue; +use std::io::Cursor; use tonic::{Request, Response, Status}; use tracing::debug; @@ -39,87 +36,15 @@ use tracing::debug; const MSGPACK_ENCODE_CAPACITY_HINT: usize = 512; const FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT: usize = 1024; const SNAPSHOT_LEASE_PROTOCOL_VERSION: u32 = 1; -const SNAPSHOT_LEASE_MIN_TTL: Duration = Duration::from_secs(5); -const SNAPSHOT_LEASE_MAX_TTL: Duration = Duration::from_secs(5 * 60); -struct SnapshotLeaseExpiry { - disk: DiskStore, - volume: String, - path: String, - token: SnapshotLeaseToken, -} - -pub(super) struct SnapshotLeaseExpiryScheduler { - tx: mpsc::UnboundedSender, -} - -enum SnapshotLeaseExpiryCommand { - Schedule(SnapshotLeaseExpiry, Duration), - Cancel(SnapshotLeaseToken), -} - -impl SnapshotLeaseExpiryScheduler { - pub(super) fn new() -> Self { - let (tx, mut rx) = mpsc::unbounded_channel(); - tokio::spawn(async move { - let mut expirations: DelayQueue = DelayQueue::new(); - let mut keys = HashMap::new(); - loop { - tokio::select! { - Some(command) = rx.recv() => { - match command { - SnapshotLeaseExpiryCommand::Schedule(expiry, ttl) => { - if let Some(key) = keys.remove(&expiry.token) { - expirations.remove(&key); - } - let token = expiry.token; - let key = expirations.insert(expiry, ttl); - keys.insert(token, key); - } - SnapshotLeaseExpiryCommand::Cancel(token) => { - if let Some(key) = keys.remove(&token) { - expirations.remove(&key); - } - } - } - } - Some(expired) = futures_util::StreamExt::next(&mut expirations), if !expirations.is_empty() => { - let expiry = expired.into_inner(); - keys.remove(&expiry.token); - let _ = expiry - .disk - .release_snapshot_lease(&expiry.volume, &expiry.path, expiry.token) - .await; - } - else => break, - } - } - }); - Self { tx } - } - - fn schedule(&self, expiry: SnapshotLeaseExpiry, ttl: Duration) -> Result<(), SnapshotLeaseExpiry> { - self.tx - .send(SnapshotLeaseExpiryCommand::Schedule(expiry, ttl)) - .map_err(|err| match err.0 { - SnapshotLeaseExpiryCommand::Schedule(expiry, _) => expiry, - SnapshotLeaseExpiryCommand::Cancel(_) => unreachable!(), - }) - } - - fn cancel(&self, token: SnapshotLeaseToken) { - let _ = self.tx.send(SnapshotLeaseExpiryCommand::Cancel(token)); +fn snapshot_lease_disabled_response() -> SnapshotLeaseResponse { + SnapshotLeaseResponse { + success: false, + token: Bytes::new(), + protocol_version: SNAPSHOT_LEASE_PROTOCOL_VERSION, + error: Some(DiskError::UnsupportedDisk.into()), } } - -fn snapshot_lease_ttl(ttl_ms: u64) -> Result { - let ttl = Duration::from_millis(ttl_ms); - if !(SNAPSHOT_LEASE_MIN_TTL..=SNAPSHOT_LEASE_MAX_TTL).contains(&ttl) { - return Err(Status::invalid_argument("snapshot lease TTL is outside the supported range")); - } - Ok(ttl) -} - fn decode_msgpack_or_json( binary: &[u8], json: &str, @@ -274,47 +199,7 @@ impl NodeService { rustfs_protos::canonical_snapshot_lease_request_body(request.get_ref()), "acquire_snapshot_lease", )?; - let request = request.into_inner(); - let ttl = snapshot_lease_ttl(request.ttl_ms)?; - let Some(disk) = self.find_disk(&request.disk).await else { - return Ok(Response::new(SnapshotLeaseResponse { - success: false, - token: Bytes::new(), - protocol_version: SNAPSHOT_LEASE_PROTOCOL_VERSION, - error: Some(DiskError::other("cannot find disk").into()), - })); - }; - match disk.acquire_snapshot_lease(&request.volume, &request.path).await { - Ok(token) => { - if let Err(expiry) = self.snapshot_lease_expiry.schedule( - SnapshotLeaseExpiry { - disk, - volume: request.volume, - path: request.path, - token, - }, - ttl, - ) { - let _ = expiry - .disk - .release_snapshot_lease(&expiry.volume, &expiry.path, expiry.token) - .await; - return Err(Status::internal("snapshot lease expiry scheduler is unavailable")); - } - Ok(Response::new(SnapshotLeaseResponse { - success: true, - token: Bytes::copy_from_slice(token.as_bytes()), - protocol_version: SNAPSHOT_LEASE_PROTOCOL_VERSION, - error: None, - })) - } - Err(err) => Ok(Response::new(SnapshotLeaseResponse { - success: false, - token: Bytes::new(), - protocol_version: SNAPSHOT_LEASE_PROTOCOL_VERSION, - error: Some(err.into()), - })), - } + Ok(Response::new(snapshot_lease_disabled_response())) } pub(super) async fn handle_renew_snapshot_lease( @@ -326,50 +211,7 @@ impl NodeService { rustfs_protos::canonical_snapshot_lease_renew_request_body(request.get_ref()), "renew_snapshot_lease", )?; - let request = request.into_inner(); - let ttl = snapshot_lease_ttl(request.ttl_ms)?; - let token = - SnapshotLeaseToken::from_slice(&request.token).map_err(|_| Status::invalid_argument("invalid lease token"))?; - let Some(disk) = self.find_disk(&request.disk).await else { - return Ok(Response::new(SnapshotLeaseResponse { - success: false, - token: Bytes::new(), - protocol_version: SNAPSHOT_LEASE_PROTOCOL_VERSION, - error: Some(DiskError::other("cannot find disk").into()), - })); - }; - match disk.renew_snapshot_lease(&request.volume, &request.path, token).await { - Ok(renewed) => { - if let Err(expiry) = self.snapshot_lease_expiry.schedule( - SnapshotLeaseExpiry { - disk, - volume: request.volume, - path: request.path, - token: renewed, - }, - ttl, - ) { - let _ = expiry - .disk - .release_snapshot_lease(&expiry.volume, &expiry.path, expiry.token) - .await; - return Err(Status::internal("snapshot lease expiry scheduler is unavailable")); - } - self.snapshot_lease_expiry.cancel(token); - Ok(Response::new(SnapshotLeaseResponse { - success: true, - token: Bytes::copy_from_slice(renewed.as_bytes()), - protocol_version: SNAPSHOT_LEASE_PROTOCOL_VERSION, - error: None, - })) - } - Err(err) => Ok(Response::new(SnapshotLeaseResponse { - success: false, - token: Bytes::new(), - protocol_version: SNAPSHOT_LEASE_PROTOCOL_VERSION, - error: Some(err.into()), - })), - } + Ok(Response::new(snapshot_lease_disabled_response())) } pub(super) async fn handle_release_snapshot_lease( @@ -391,13 +233,10 @@ impl NodeService { })); }; match disk.release_snapshot_lease(&request.volume, &request.path, token).await { - Ok(()) => { - self.snapshot_lease_expiry.cancel(token); - Ok(Response::new(SnapshotLeaseMutationResponse { - success: true, - error: None, - })) - } + Ok(()) => Ok(Response::new(SnapshotLeaseMutationResponse { + success: true, + error: None, + })), Err(err) => Ok(Response::new(SnapshotLeaseMutationResponse { success: false, error: Some(err.into()), @@ -1608,10 +1447,10 @@ impl NodeService { #[cfg(test)] mod tests { use super::{ - SNAPSHOT_LEASE_MAX_TTL, SNAPSHOT_LEASE_MIN_TTL, compat_response_json, decode_msgpack_or_json, - encode_batch_read_version_response_payloads, encode_msgpack, encode_msgpack_named, - encode_read_multiple_response_payloads, snapshot_lease_ttl, + compat_response_json, decode_msgpack_or_json, encode_batch_read_version_response_payloads, encode_msgpack, + encode_msgpack_named, encode_read_multiple_response_payloads, snapshot_lease_disabled_response, }; + use crate::storage::storage_api::DiskError; use crate::storage::storage_api::ReadMultipleResp; use crate::storage::storage_api::rpc_consumer::node_service::BatchReadVersionResp; use rustfs_io_metrics::internode_metrics::global_internode_metrics; @@ -1624,11 +1463,14 @@ mod tests { } #[test] - fn snapshot_lease_ttl_rejects_values_outside_server_bounds() { - assert!(snapshot_lease_ttl(4_999).is_err()); - assert_eq!(snapshot_lease_ttl(5_000).unwrap(), SNAPSHOT_LEASE_MIN_TTL); - assert_eq!(snapshot_lease_ttl(300_000).unwrap(), SNAPSHOT_LEASE_MAX_TTL); - assert!(snapshot_lease_ttl(300_001).is_err()); + fn snapshot_lease_acquire_and_renew_fail_closed() { + let response = snapshot_lease_disabled_response(); + let expected_error = DiskError::UnsupportedDisk.into(); + + assert!(!response.success); + assert!(response.token.is_empty()); + assert_eq!(response.protocol_version, 1); + assert_eq!(response.error, Some(expected_error)); } #[test] diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index d90e106df..4b40fe6ca 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -1099,8 +1099,6 @@ pub(crate) trait StorageDiskRpcExt { ) -> DiskResult<()>; async fn read_metadata(&self, volume: &str, path: &str) -> DiskResult; async fn delete_paths(&self, volume: &str, paths: &[String]) -> DiskResult<()>; - async fn acquire_snapshot_lease(&self, volume: &str, path: &str) -> DiskResult; - async fn renew_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> DiskResult; async fn release_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> DiskResult<()>; async fn stat_volume(&self, volume: &str) -> DiskResult; async fn list_volumes(&self) -> DiskResult>; @@ -1228,14 +1226,6 @@ where ecstore_disk::DiskAPI::delete_paths(self, volume, paths).await } - async fn acquire_snapshot_lease(&self, volume: &str, path: &str) -> DiskResult { - ecstore_disk::DiskAPI::acquire_snapshot_lease(self, volume, path).await - } - - async fn renew_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> DiskResult { - ecstore_disk::DiskAPI::renew_snapshot_lease(self, volume, path, token).await - } - async fn release_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> DiskResult<()> { ecstore_disk::DiskAPI::release_snapshot_lease(self, volume, path, token).await }