mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-06 13:27:43 +00:00
fix(ecstore): retain namespace read locks for streaming GETs to prevent quota-bypass (#5699)
* fix(ecstore): retain locks for streaming GETs * fix(rpc): keep disabled snapshot leases lint-clean
This commit is contained in:
+135
-680
@@ -370,254 +370,6 @@ struct SetDiskLockGuardedReader {
|
||||
guard: Option<ObjectLockDiagGuard>,
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
struct SnapshotLease {
|
||||
disk: DiskStore,
|
||||
token: SnapshotLeaseToken,
|
||||
}
|
||||
|
||||
struct SnapshotLeaseState {
|
||||
volume: Arc<str>,
|
||||
path: Arc<str>,
|
||||
leases: parking_lot::Mutex<Vec<SnapshotLease>>,
|
||||
renewal_failed: AtomicBool,
|
||||
renewal_waker: AtomicWaker,
|
||||
cancel: CancellationToken,
|
||||
runtime: tokio::runtime::Handle,
|
||||
}
|
||||
|
||||
impl SnapshotLeaseState {
|
||||
fn apply_renewal_results(&self, results: Vec<std::result::Result<SnapshotLeaseToken, DiskError>>) -> 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<SnapshotLeaseState>);
|
||||
|
||||
impl SnapshotLeaseHandle {
|
||||
fn new(volume: Arc<str>, path: Arc<str>, leases: Vec<SnapshotLease>) -> 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<F, Fut>(
|
||||
leases: &[SnapshotLease],
|
||||
volume: Arc<str>,
|
||||
path: Arc<str>,
|
||||
renewal_timeout: Duration,
|
||||
renew: F,
|
||||
) -> Vec<std::result::Result<SnapshotLeaseToken, DiskError>>
|
||||
where
|
||||
F: Fn(DiskStore, Arc<str>, Arc<str>, SnapshotLeaseToken) -> Fut,
|
||||
Fut: Future<Output = std::result::Result<SnapshotLeaseToken, DiskError>>,
|
||||
{
|
||||
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<dyn AsyncRead + Unpin + Send + Sync>,
|
||||
lease: Option<SnapshotLeaseHandle>,
|
||||
terminal_error: bool,
|
||||
}
|
||||
|
||||
impl AsyncRead for SnapshotLeaseReader {
|
||||
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<std::io::Result<()>> {
|
||||
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<DiskStore>],
|
||||
volume: &str,
|
||||
path: &str,
|
||||
read_quorum: usize,
|
||||
) -> Option<SnapshotLeaseHandle> {
|
||||
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<F, Fut>(
|
||||
disks: &[Option<DiskStore>],
|
||||
volume: &str,
|
||||
path: &str,
|
||||
read_quorum: usize,
|
||||
candidate_timeout: Duration,
|
||||
acquire: F,
|
||||
) -> Option<SnapshotLeaseHandle>
|
||||
where
|
||||
F: Fn(DiskStore, Arc<str>, Arc<str>) -> Fut,
|
||||
Fut: Future<Output = std::result::Result<SnapshotLeaseToken, DiskError>> + Send + 'static,
|
||||
{
|
||||
let candidates = disks.iter().cloned().collect::<Option<Vec<_>>>()?;
|
||||
if candidates.len() < read_quorum {
|
||||
return None;
|
||||
}
|
||||
let volume: Arc<str> = Arc::from(volume);
|
||||
let path: Arc<str> = 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<std::io::Result<()>> {
|
||||
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<ObjectLockDiagGuard>,
|
||||
snapshot_lease: Option<SnapshotLeaseHandle>,
|
||||
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<AtomicBool>);
|
||||
|
||||
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<Option<DiskStore>>) -> Arc<SetDisks> {
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -380,7 +380,6 @@ mod metrics;
|
||||
pub struct NodeService {
|
||||
local_peer: LocalPeerS3Client,
|
||||
context: Option<Arc<runtime_sources::AppContext>>,
|
||||
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<Arc<runtime_sources::AppContext>>) -> 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();
|
||||
|
||||
@@ -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<SnapshotLeaseExpiryCommand>,
|
||||
}
|
||||
|
||||
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<SnapshotLeaseExpiry> = 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<Duration, Status> {
|
||||
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<T: DeserializeOwned>(
|
||||
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]
|
||||
|
||||
@@ -1099,8 +1099,6 @@ pub(crate) trait StorageDiskRpcExt {
|
||||
) -> DiskResult<()>;
|
||||
async fn read_metadata(&self, volume: &str, path: &str) -> DiskResult<bytes::Bytes>;
|
||||
async fn delete_paths(&self, volume: &str, paths: &[String]) -> DiskResult<()>;
|
||||
async fn acquire_snapshot_lease(&self, volume: &str, path: &str) -> DiskResult<SnapshotLeaseToken>;
|
||||
async fn renew_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> DiskResult<SnapshotLeaseToken>;
|
||||
async fn release_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> DiskResult<()>;
|
||||
async fn stat_volume(&self, volume: &str) -> DiskResult<VolumeInfo>;
|
||||
async fn list_volumes(&self) -> DiskResult<Vec<VolumeInfo>>;
|
||||
@@ -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<SnapshotLeaseToken> {
|
||||
ecstore_disk::DiskAPI::acquire_snapshot_lease(self, volume, path).await
|
||||
}
|
||||
|
||||
async fn renew_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> DiskResult<SnapshotLeaseToken> {
|
||||
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
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user