From d4bf17356ba27462ce7deaa4983ec578ea177cba Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=A9=AC=E7=99=BB=E5=B1=B1?= Date: Tue, 28 Jul 2026 20:07:59 +0800 Subject: [PATCH] feat(ecstore): add remote snapshot lease RPCs --- Cargo.lock | 1 + crates/ecstore/src/api/mod.rs | 3 +- crates/ecstore/src/cluster/rpc/remote_disk.rs | 122 ++++++++- crates/ecstore/src/disk/disk_store.rs | 8 + crates/ecstore/src/disk/local.rs | 28 +- crates/ecstore/src/disk/mod.rs | 22 ++ .../src/generated/proto_gen/node_service.rs | 194 ++++++++++++++ crates/protos/src/lib.rs | 93 ++++++- crates/protos/src/node.proto | 37 +++ rustfs/Cargo.toml | 2 +- rustfs/src/storage/rpc/node_service.rs | 25 +- rustfs/src/storage/rpc/node_service/disk.rs | 245 +++++++++++++++++- rustfs/src/storage/storage_api.rs | 18 +- 13 files changed, 784 insertions(+), 14 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index e769854d8..27850b191 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -11829,6 +11829,7 @@ dependencies = [ "futures-util", "libc", "pin-project-lite", + "slab", "tokio", ] diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index fc042f97f..eab2dedfa 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -305,7 +305,8 @@ pub mod disk { CheckPartsResp, DeleteOptions, Disk, DiskAPI, DiskInfo, DiskInfoOptions, DiskLocation, DiskOption, DiskStore, FileInfoVersions, FileReader, FileWriter, HEALING_MARKER_PATH, NsScannerOpenRequest, OldCurrentSize, PartTransactionAction, RUSTFS_META_BUCKET, ReadMultipleReq, ReadMultipleResp, ReadOptions, RenameDataResp, - STORAGE_FORMAT_FILE, UpdateMetadataOpts, VolumeInfo, WalkDirOptions, new_disk, validate_batch_read_version_item_count, + STORAGE_FORMAT_FILE, SnapshotLeaseToken, UpdateMetadataOpts, VolumeInfo, WalkDirOptions, new_disk, + validate_batch_read_version_item_count, }; pub use bytes::Bytes; pub use endpoint::Endpoint; diff --git a/crates/ecstore/src/cluster/rpc/remote_disk.rs b/crates/ecstore/src/cluster/rpc/remote_disk.rs index 300f443d2..ffeaa2791 100644 --- a/crates/ecstore/src/cluster/rpc/remote_disk.rs +++ b/crates/ecstore/src/cluster/rpc/remote_disk.rs @@ -25,7 +25,7 @@ use crate::disk::error::{Error, Result}; use crate::disk::{ BatchReadVersionReq, BatchReadVersionResp, CheckPartsResp, DeleteOptions, DiskAPI, DiskInfo, DiskInfoOptions, DiskLocation, DiskOption, FileInfoVersions, FileReader, FileWriter, PartTransactionAction, ReadMultipleReq, ReadMultipleResp, ReadOptions, - RenameDataResp, UpdateMetadataOpts, VolumeInfo, WalkDirOptions, batch_read_version_one_by_one, + RenameDataResp, SnapshotLeaseToken, UpdateMetadataOpts, VolumeInfo, WalkDirOptions, batch_read_version_one_by_one, disk_store::{ DEFAULT_RUSTFS_DRIVE_ACTIVE_MONITORING, ENV_RUSTFS_DRIVE_ACTIVE_MONITORING, SKIP_IF_SUCCESS_BEFORE, get_drive_active_check_interval, get_drive_active_check_timeout, get_drive_disk_info_timeout, get_drive_list_dir_timeout, @@ -50,8 +50,9 @@ use rustfs_protos::proto_gen::node_service::{ DeleteVersionRequest, DeleteVersionsRequest, DeleteVolumeRequest, DiskInfoRequest, ListDirRequest, ListVolumesRequest, MakeVolumeRequest, MakeVolumesRequest, PreparePartTransactionRequest, ReadAllRequest, ReadMetadataRequest, ReadMultipleRequest, ReadMultipleResponse, ReadPartsRequest, ReadVersionRequest, ReadXlRequest, RenameDataRequest, - RenameFileRequest, SettlePartTransactionRequest, StatVolumeRequest, UpdateMetadataRequest, VerifyFileRequest, - WriteAllRequest, WriteMetadataRequest, node_service_client::NodeServiceClient, + RenameFileRequest, SettlePartTransactionRequest, SnapshotLeaseReleaseRequest, SnapshotLeaseRenewRequest, + SnapshotLeaseRequest, SnapshotLeaseResponse, StatVolumeRequest, UpdateMetadataRequest, VerifyFileRequest, WriteAllRequest, + WriteMetadataRequest, node_service_client::NodeServiceClient, }; use serde::{Serialize, de::DeserializeOwned}; use std::{ @@ -100,6 +101,18 @@ const LOG_COMPONENT_ECSTORE: &str = "ecstore"; const LOG_SUBSYSTEM_REMOTE_DISK: &str = "remote_disk"; const EVENT_REMOTE_DISK_HEALTH: &str = "remote_disk_health"; const EVENT_REMOTE_DISK_RPC: &str = "remote_disk_rpc"; +const SNAPSHOT_LEASE_PROTOCOL_VERSION: u32 = 1; +pub const REMOTE_SNAPSHOT_LEASE_TTL: Duration = Duration::from_secs(60); + +fn snapshot_lease_token_from_response(response: SnapshotLeaseResponse) -> Result { + if !response.success { + return Err(response.error.unwrap_or_default().into()); + } + if response.protocol_version != SNAPSHOT_LEASE_PROTOCOL_VERSION { + return Err(Error::other("remote snapshot lease protocol is incompatible")); + } + SnapshotLeaseToken::from_slice(&response.token) +} /// Bind a mutating disk RPC to its canonical body: the digest lands in the request metadata, and /// the signing interceptor folds it (plus a replay-protected nonce) into the v2 signature scope @@ -1784,6 +1797,81 @@ impl DiskAPI for RemoteDisk { .await } + async fn acquire_snapshot_lease(&self, volume: &str, path: &str) -> Result { + self.execute_with_timeout( + || async { + let mut client = self + .get_client() + .await + .map_err(|err| Error::other(format!("can not get client, err: {err}")))?; + let mut request = Request::new(SnapshotLeaseRequest { + disk: self.endpoint.to_string(), + volume: volume.to_string(), + path: path.to_string(), + ttl_ms: u64::try_from(REMOTE_SNAPSHOT_LEASE_TTL.as_millis()) + .map_err(|_| Error::other("snapshot lease TTL cannot be represented"))?, + }); + let canonical_body = rustfs_protos::canonical_snapshot_lease_request_body(request.get_ref()); + attach_mutation_body_digest(&mut request, canonical_body, "acquire_snapshot_lease")?; + let response = client.acquire_snapshot_lease(request).await?.into_inner(); + snapshot_lease_token_from_response(response) + }, + get_max_timeout_duration(), + ) + .await + } + + async fn renew_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> Result { + self.execute_with_timeout( + || async { + let mut client = self + .get_client() + .await + .map_err(|err| Error::other(format!("can not get client, err: {err}")))?; + let mut request = Request::new(SnapshotLeaseRenewRequest { + disk: self.endpoint.to_string(), + volume: volume.to_string(), + path: path.to_string(), + token: token.as_bytes().to_vec().into(), + ttl_ms: u64::try_from(REMOTE_SNAPSHOT_LEASE_TTL.as_millis()) + .map_err(|_| Error::other("snapshot lease TTL cannot be represented"))?, + }); + let canonical_body = rustfs_protos::canonical_snapshot_lease_renew_request_body(request.get_ref()); + attach_mutation_body_digest(&mut request, canonical_body, "renew_snapshot_lease")?; + let response = client.renew_snapshot_lease(request).await?.into_inner(); + snapshot_lease_token_from_response(response) + }, + get_max_timeout_duration(), + ) + .await + } + + async fn release_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> Result<()> { + self.execute_with_timeout( + || async { + let mut client = self + .get_client() + .await + .map_err(|err| Error::other(format!("can not get client, err: {err}")))?; + let mut request = Request::new(SnapshotLeaseReleaseRequest { + disk: self.endpoint.to_string(), + volume: volume.to_string(), + path: path.to_string(), + token: token.as_bytes().to_vec().into(), + }); + let canonical_body = rustfs_protos::canonical_snapshot_lease_release_request_body(request.get_ref()); + attach_mutation_body_digest(&mut request, canonical_body, "release_snapshot_lease")?; + let response = client.release_snapshot_lease(request).await?.into_inner(); + if !response.success { + return Err(response.error.unwrap_or_default().into()); + } + Ok(()) + }, + get_max_timeout_duration(), + ) + .await + } + #[tracing::instrument(level = "trace", skip_all)] async fn write_metadata(&self, _org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()> { trace!( @@ -2930,6 +3018,34 @@ mod tests { static INIT: Once = Once::new(); + #[test] + fn snapshot_lease_response_requires_current_protocol_and_valid_token() { + let token = SnapshotLeaseToken::new(); + let response = SnapshotLeaseResponse { + success: true, + token: token.as_bytes().to_vec().into(), + protocol_version: SNAPSHOT_LEASE_PROTOCOL_VERSION, + error: None, + }; + assert_eq!(snapshot_lease_token_from_response(response).unwrap(), token); + + let incompatible = SnapshotLeaseResponse { + success: true, + token: token.as_bytes().to_vec().into(), + protocol_version: SNAPSHOT_LEASE_PROTOCOL_VERSION + 1, + error: None, + }; + assert!(snapshot_lease_token_from_response(incompatible).is_err()); + + let malformed = SnapshotLeaseResponse { + success: true, + token: Bytes::from_static(b"not-a-uuid"), + protocol_version: SNAPSHOT_LEASE_PROTOCOL_VERSION, + error: None, + }; + assert!(snapshot_lease_token_from_response(malformed).is_err()); + } + #[test] fn list_volumes_decode_rejects_a_malformed_entry() { let valid = serde_json::to_string(&VolumeInfo { diff --git a/crates/ecstore/src/disk/disk_store.rs b/crates/ecstore/src/disk/disk_store.rs index 940638853..1183d7c18 100644 --- a/crates/ecstore/src/disk/disk_store.rs +++ b/crates/ecstore/src/disk/disk_store.rs @@ -1365,6 +1365,14 @@ impl DiskAPI for LocalDiskWrapper { .await } + async fn renew_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> Result { + self.track_disk_health( + || async { self.disk.renew_snapshot_lease(volume, path, token).await }, + get_max_timeout_duration(), + ) + .await + } + async fn delete_data_dir(&self, volume: &str, path: &str, opts: DeleteOptions) -> Result { self.track_disk_health( || async { self.disk.delete_data_dir(volume, path, opts).await }, diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 991a2248b..c45a9ec32 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -7908,6 +7908,23 @@ impl DiskAPI for LocalDisk { } } + async fn renew_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> Result { + let key = SnapshotLeaseKey { + volume: volume.to_string(), + path: path.to_string(), + }; + let mut registry = self.snapshot_leases.lock().await; + let Some(entry) = registry.entries.get_mut(&key) else { + return Err(DiskError::FileNotFound); + }; + if entry.deleting || !entry.tokens.remove(&token) { + return Err(DiskError::FileNotFound); + } + let renewed = SnapshotLeaseToken::new(); + entry.tokens.insert(renewed); + Ok(renewed) + } + async fn delete_data_dir(&self, volume: &str, path: &str, opts: DeleteOptions) -> Result { let key = SnapshotLeaseKey { volume: volume.to_string(), @@ -14957,6 +14974,13 @@ mod test { .acquire_snapshot_lease(volume, &data_dir) .await .expect("second lease should be acquired"); + let renewed = disk + .renew_snapshot_lease(volume, &data_dir, first) + .await + .expect("first lease should renew atomically"); + disk.release_snapshot_lease(volume, &data_dir, first) + .await + .expect("the superseded token should be idempotent"); let status = disk .delete_data_dir( volume, @@ -14976,9 +15000,9 @@ mod test { Bytes::from_static(b"later") ); - disk.release_snapshot_lease(volume, &data_dir, first) + disk.release_snapshot_lease(volume, &data_dir, renewed) .await - .expect("first lease release should succeed"); + .expect("renewed lease release should succeed"); assert!( disk.read_all(volume, &first_part).await.is_ok(), "one remaining lease must keep the data directory" diff --git a/crates/ecstore/src/disk/mod.rs b/crates/ecstore/src/disk/mod.rs index ee947d291..4ac894373 100644 --- a/crates/ecstore/src/disk/mod.rs +++ b/crates/ecstore/src/disk/mod.rs @@ -79,6 +79,18 @@ impl SnapshotLeaseToken { pub fn new() -> Self { Self(Uuid::new_v4()) } + + pub fn from_slice(bytes: &[u8]) -> Result { + let uuid = Uuid::from_slice(bytes).map_err(|_| Error::other("invalid snapshot lease token"))?; + if uuid.is_nil() { + return Err(Error::other("invalid snapshot lease token")); + } + Ok(Self(uuid)) + } + + pub fn as_bytes(&self) -> &[u8; 16] { + self.0.as_bytes() + } } impl Default for SnapshotLeaseToken { @@ -284,6 +296,13 @@ impl DiskAPI for Disk { } } + async fn renew_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> Result { + match self { + Disk::Local(local_disk) => local_disk.renew_snapshot_lease(volume, path, token).await, + Disk::Remote(remote_disk) => remote_disk.renew_snapshot_lease(volume, path, token).await, + } + } + async fn delete_data_dir(&self, volume: &str, path: &str, opts: DeleteOptions) -> Result { match self { Disk::Local(local_disk) => local_disk.delete_data_dir(volume, path, opts).await, @@ -694,6 +713,9 @@ pub trait DiskAPI: Debug + Send + Sync + 'static { async fn release_snapshot_lease(&self, _volume: &str, _path: &str, _token: SnapshotLeaseToken) -> Result<()> { Err(Error::other("snapshot leases are not supported by this disk")) } + async fn renew_snapshot_lease(&self, _volume: &str, _path: &str, _token: SnapshotLeaseToken) -> Result { + Err(Error::other("snapshot leases are not supported by this disk")) + } async fn delete_data_dir(&self, volume: &str, path: &str, opts: DeleteOptions) -> Result { self.delete(volume, path, opts).await?; Ok(DataDirDeleteStatus::Deleted) diff --git a/crates/protos/src/generated/proto_gen/node_service.rs b/crates/protos/src/generated/proto_gen/node_service.rs index 8bc8a386f..35c320d7a 100644 --- a/crates/protos/src/generated/proto_gen/node_service.rs +++ b/crates/protos/src/generated/proto_gen/node_service.rs @@ -484,6 +484,59 @@ pub struct DeletePathsResponse { pub error: ::core::option::Option, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct SnapshotLeaseRequest { + #[prost(string, tag = "1")] + pub disk: ::prost::alloc::string::String, + #[prost(string, tag = "2")] + pub volume: ::prost::alloc::string::String, + #[prost(string, tag = "3")] + pub path: ::prost::alloc::string::String, + #[prost(uint64, tag = "4")] + pub ttl_ms: u64, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct SnapshotLeaseRenewRequest { + #[prost(string, tag = "1")] + pub disk: ::prost::alloc::string::String, + #[prost(string, tag = "2")] + pub volume: ::prost::alloc::string::String, + #[prost(string, tag = "3")] + pub path: ::prost::alloc::string::String, + #[prost(bytes = "bytes", tag = "4")] + pub token: ::prost::bytes::Bytes, + #[prost(uint64, tag = "5")] + pub ttl_ms: u64, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct SnapshotLeaseReleaseRequest { + #[prost(string, tag = "1")] + pub disk: ::prost::alloc::string::String, + #[prost(string, tag = "2")] + pub volume: ::prost::alloc::string::String, + #[prost(string, tag = "3")] + pub path: ::prost::alloc::string::String, + #[prost(bytes = "bytes", tag = "4")] + pub token: ::prost::bytes::Bytes, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct SnapshotLeaseResponse { + #[prost(bool, tag = "1")] + pub success: bool, + #[prost(bytes = "bytes", tag = "2")] + pub token: ::prost::bytes::Bytes, + #[prost(uint32, tag = "3")] + pub protocol_version: u32, + #[prost(message, optional, tag = "4")] + pub error: ::core::option::Option, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct SnapshotLeaseMutationResponse { + #[prost(bool, tag = "1")] + pub success: bool, + #[prost(message, optional, tag = "2")] + pub error: ::core::option::Option, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct ReadMetadataRequest { #[prost(string, tag = "1")] pub disk: ::prost::alloc::string::String, @@ -1595,6 +1648,51 @@ pub mod node_service_client { .insert(GrpcMethod::new("node_service.NodeService", "Delete")); self.inner.unary(req, path, codec).await } + pub async fn acquire_snapshot_lease( + &mut self, + request: impl tonic::IntoRequest, + ) -> std::result::Result, tonic::Status> { + self.inner + .ready() + .await + .map_err(|e| tonic::Status::unknown(format!("Service was not ready: {}", e.into())))?; + let codec = tonic_prost::ProstCodec::default(); + let path = http::uri::PathAndQuery::from_static("/node_service.NodeService/AcquireSnapshotLease"); + let mut req = request.into_request(); + req.extensions_mut() + .insert(GrpcMethod::new("node_service.NodeService", "AcquireSnapshotLease")); + self.inner.unary(req, path, codec).await + } + pub async fn renew_snapshot_lease( + &mut self, + request: impl tonic::IntoRequest, + ) -> std::result::Result, tonic::Status> { + self.inner + .ready() + .await + .map_err(|e| tonic::Status::unknown(format!("Service was not ready: {}", e.into())))?; + let codec = tonic_prost::ProstCodec::default(); + let path = http::uri::PathAndQuery::from_static("/node_service.NodeService/RenewSnapshotLease"); + let mut req = request.into_request(); + req.extensions_mut() + .insert(GrpcMethod::new("node_service.NodeService", "RenewSnapshotLease")); + self.inner.unary(req, path, codec).await + } + pub async fn release_snapshot_lease( + &mut self, + request: impl tonic::IntoRequest, + ) -> std::result::Result, tonic::Status> { + self.inner + .ready() + .await + .map_err(|e| tonic::Status::unknown(format!("Service was not ready: {}", e.into())))?; + let codec = tonic_prost::ProstCodec::default(); + let path = http::uri::PathAndQuery::from_static("/node_service.NodeService/ReleaseSnapshotLease"); + let mut req = request.into_request(); + req.extensions_mut() + .insert(GrpcMethod::new("node_service.NodeService", "ReleaseSnapshotLease")); + self.inner.unary(req, path, codec).await + } pub async fn verify_file( &mut self, request: impl tonic::IntoRequest, @@ -2784,6 +2882,18 @@ pub mod node_service_server { &self, request: tonic::Request, ) -> std::result::Result, tonic::Status>; + async fn acquire_snapshot_lease( + &self, + request: tonic::Request, + ) -> std::result::Result, tonic::Status>; + async fn renew_snapshot_lease( + &self, + request: tonic::Request, + ) -> std::result::Result, tonic::Status>; + async fn release_snapshot_lease( + &self, + request: tonic::Request, + ) -> std::result::Result, tonic::Status>; async fn verify_file( &self, request: tonic::Request, @@ -3426,6 +3536,90 @@ pub mod node_service_server { }; Box::pin(fut) } + "/node_service.NodeService/AcquireSnapshotLease" => { + #[allow(non_camel_case_types)] + struct AcquireSnapshotLeaseSvc(pub Arc); + impl tonic::server::UnaryService for AcquireSnapshotLeaseSvc { + type Response = super::SnapshotLeaseResponse; + type Future = BoxFuture, tonic::Status>; + fn call(&mut self, request: tonic::Request) -> Self::Future { + let inner = Arc::clone(&self.0); + let fut = async move { ::acquire_snapshot_lease(&inner, request).await }; + Box::pin(fut) + } + } + let accept_compression_encodings = self.accept_compression_encodings; + let send_compression_encodings = self.send_compression_encodings; + let max_decoding_message_size = self.max_decoding_message_size; + let max_encoding_message_size = self.max_encoding_message_size; + let inner = self.inner.clone(); + let fut = async move { + let method = AcquireSnapshotLeaseSvc(inner); + let codec = tonic_prost::ProstCodec::default(); + let mut grpc = tonic::server::Grpc::new(codec) + .apply_compression_config(accept_compression_encodings, send_compression_encodings) + .apply_max_message_size_config(max_decoding_message_size, max_encoding_message_size); + let res = grpc.unary(method, req).await; + Ok(res) + }; + Box::pin(fut) + } + "/node_service.NodeService/RenewSnapshotLease" => { + #[allow(non_camel_case_types)] + struct RenewSnapshotLeaseSvc(pub Arc); + impl tonic::server::UnaryService for RenewSnapshotLeaseSvc { + type Response = super::SnapshotLeaseResponse; + type Future = BoxFuture, tonic::Status>; + fn call(&mut self, request: tonic::Request) -> Self::Future { + let inner = Arc::clone(&self.0); + let fut = async move { ::renew_snapshot_lease(&inner, request).await }; + Box::pin(fut) + } + } + let accept_compression_encodings = self.accept_compression_encodings; + let send_compression_encodings = self.send_compression_encodings; + let max_decoding_message_size = self.max_decoding_message_size; + let max_encoding_message_size = self.max_encoding_message_size; + let inner = self.inner.clone(); + let fut = async move { + let method = RenewSnapshotLeaseSvc(inner); + let codec = tonic_prost::ProstCodec::default(); + let mut grpc = tonic::server::Grpc::new(codec) + .apply_compression_config(accept_compression_encodings, send_compression_encodings) + .apply_max_message_size_config(max_decoding_message_size, max_encoding_message_size); + let res = grpc.unary(method, req).await; + Ok(res) + }; + Box::pin(fut) + } + "/node_service.NodeService/ReleaseSnapshotLease" => { + #[allow(non_camel_case_types)] + struct ReleaseSnapshotLeaseSvc(pub Arc); + impl tonic::server::UnaryService for ReleaseSnapshotLeaseSvc { + type Response = super::SnapshotLeaseMutationResponse; + type Future = BoxFuture, tonic::Status>; + fn call(&mut self, request: tonic::Request) -> Self::Future { + let inner = Arc::clone(&self.0); + let fut = async move { ::release_snapshot_lease(&inner, request).await }; + Box::pin(fut) + } + } + let accept_compression_encodings = self.accept_compression_encodings; + let send_compression_encodings = self.send_compression_encodings; + let max_decoding_message_size = self.max_decoding_message_size; + let max_encoding_message_size = self.max_encoding_message_size; + let inner = self.inner.clone(); + let fut = async move { + let method = ReleaseSnapshotLeaseSvc(inner); + let codec = tonic_prost::ProstCodec::default(); + let mut grpc = tonic::server::Grpc::new(codec) + .apply_compression_config(accept_compression_encodings, send_compression_encodings) + .apply_max_message_size_config(max_decoding_message_size, max_encoding_message_size); + let res = grpc.unary(method, req).await; + Ok(res) + }; + Box::pin(fut) + } "/node_service.NodeService/VerifyFile" => { #[allow(non_camel_case_types)] struct VerifyFileSvc(pub Arc); diff --git a/crates/protos/src/lib.rs b/crates/protos/src/lib.rs index 165ac3c52..5bbdcd2f9 100644 --- a/crates/protos/src/lib.rs +++ b/crates/protos/src/lib.rs @@ -451,6 +451,10 @@ impl CanonicalBodyBuilder { self.body.push(u8::from(field)); } + fn push_u64(&mut self, field: u64) { + self.body.extend_from_slice(&field.to_be_bytes()); + } + fn push_count(&mut self, count: usize) -> Result<(), std::num::TryFromIntError> { self.body.extend_from_slice(&u64::try_from(count)?.to_be_bytes()); Ok(()) @@ -575,6 +579,40 @@ pub fn canonical_delete_paths_request_body( Ok(body.finish()) } +pub fn canonical_snapshot_lease_request_body( + request: &proto_gen::node_service::SnapshotLeaseRequest, +) -> Result, std::num::TryFromIntError> { + let mut body = CanonicalBodyBuilder::new(b"rustfs-snapshot-lease-request-v1\0"); + body.push_str(&request.disk)?; + body.push_str(&request.volume)?; + body.push_str(&request.path)?; + body.push_u64(request.ttl_ms); + Ok(body.finish()) +} + +pub fn canonical_snapshot_lease_renew_request_body( + request: &proto_gen::node_service::SnapshotLeaseRenewRequest, +) -> Result, std::num::TryFromIntError> { + let mut body = CanonicalBodyBuilder::new(b"rustfs-snapshot-lease-renew-request-v1\0"); + body.push_str(&request.disk)?; + body.push_str(&request.volume)?; + body.push_str(&request.path)?; + body.push_bytes(&request.token)?; + body.push_u64(request.ttl_ms); + Ok(body.finish()) +} + +pub fn canonical_snapshot_lease_release_request_body( + request: &proto_gen::node_service::SnapshotLeaseReleaseRequest, +) -> Result, std::num::TryFromIntError> { + let mut body = CanonicalBodyBuilder::new(b"rustfs-snapshot-lease-release-request-v1\0"); + body.push_str(&request.disk)?; + body.push_str(&request.volume)?; + body.push_str(&request.path)?; + body.push_bytes(&request.token)?; + Ok(body.finish()) +} + pub fn canonical_rename_file_request_body( request: &proto_gen::node_service::RenameFileRequest, ) -> Result, std::num::TryFromIntError> { @@ -662,7 +700,8 @@ mod disk_mutation_canonical_tests { use super::proto_gen::node_service::{ DeletePathsRequest, DeleteRequest, DeleteVersionRequest, DeleteVersionsRequest, DeleteVolumeRequest, MakeVolumeRequest, MakeVolumesRequest, PreparePartTransactionRequest, RenameDataRequest, RenameFileRequest, RenamePartRequest, - SettlePartTransactionRequest, UpdateMetadataRequest, WriteAllRequest, WriteMetadataRequest, + SettlePartTransactionRequest, SnapshotLeaseReleaseRequest, SnapshotLeaseRenewRequest, SnapshotLeaseRequest, + UpdateMetadataRequest, WriteAllRequest, WriteMetadataRequest, }; use super::*; @@ -1020,6 +1059,58 @@ mod disk_mutation_canonical_tests { ); } + #[test] + fn snapshot_lease_canonical_bodies_bind_every_field() { + let acquire = SnapshotLeaseRequest { + disk: "d".into(), + volume: "v".into(), + path: "p".into(), + ttl_ms: 60_000, + }; + let mut acquire_bodies = vec![canonical_snapshot_lease_request_body(&acquire).unwrap()]; + for mutate in [ + |r: &mut SnapshotLeaseRequest| r.disk = "d2".into(), + |r: &mut SnapshotLeaseRequest| r.volume = "v2".into(), + |r: &mut SnapshotLeaseRequest| r.path = "p2".into(), + |r: &mut SnapshotLeaseRequest| r.ttl_ms = 60_001, + ] { + let mut request = acquire.clone(); + mutate(&mut request); + acquire_bodies.push(canonical_snapshot_lease_request_body(&request).unwrap()); + } + assert_all_distinct(&acquire_bodies); + + let renew = SnapshotLeaseRenewRequest { + disk: "d".into(), + volume: "v".into(), + path: "p".into(), + token: vec![1; 16].into(), + ttl_ms: 60_000, + }; + let mut changed_token = renew.clone(); + changed_token.token = vec![2; 16].into(); + let mut changed_ttl = renew.clone(); + changed_ttl.ttl_ms += 1; + assert_all_distinct(&[ + canonical_snapshot_lease_renew_request_body(&renew).unwrap(), + canonical_snapshot_lease_renew_request_body(&changed_token).unwrap(), + canonical_snapshot_lease_renew_request_body(&changed_ttl).unwrap(), + ]); + + let release = SnapshotLeaseReleaseRequest { + disk: "d".into(), + volume: "v".into(), + path: "p".into(), + token: vec![1; 16].into(), + }; + let mut changed_release = release.clone(); + changed_release.token = vec![2; 16].into(); + assert_ne!( + canonical_snapshot_lease_release_request_body(&release).unwrap(), + canonical_snapshot_lease_release_request_body(&changed_release).unwrap() + ); + } + #[test] fn disk_mutation_canonical_domains_are_distinct_per_message() { // The same field values must never authenticate one RPC's request as another's. diff --git a/crates/protos/src/node.proto b/crates/protos/src/node.proto index 36df6b5cb..112e88763 100644 --- a/crates/protos/src/node.proto +++ b/crates/protos/src/node.proto @@ -342,6 +342,40 @@ message DeletePathsResponse { optional Error error = 2; } +message SnapshotLeaseRequest { + string disk = 1; + string volume = 2; + string path = 3; + uint64 ttl_ms = 4; +} + +message SnapshotLeaseRenewRequest { + string disk = 1; + string volume = 2; + string path = 3; + bytes token = 4; + uint64 ttl_ms = 5; +} + +message SnapshotLeaseReleaseRequest { + string disk = 1; + string volume = 2; + string path = 3; + bytes token = 4; +} + +message SnapshotLeaseResponse { + bool success = 1; + bytes token = 2; + uint32 protocol_version = 3; + optional Error error = 4; +} + +message SnapshotLeaseMutationResponse { + bool success = 1; + optional Error error = 2; +} + message ReadMetadataRequest { string disk = 1; string volume = 2; @@ -966,6 +1000,9 @@ service NodeService { rpc ReadAll(ReadAllRequest) returns (ReadAllResponse) {}; rpc WriteAll(WriteAllRequest) returns (WriteAllResponse) {}; rpc Delete(DeleteRequest) returns (DeleteResponse) {}; + rpc AcquireSnapshotLease(SnapshotLeaseRequest) returns (SnapshotLeaseResponse) {}; + rpc RenewSnapshotLease(SnapshotLeaseRenewRequest) returns (SnapshotLeaseResponse) {}; + rpc ReleaseSnapshotLease(SnapshotLeaseReleaseRequest) returns (SnapshotLeaseMutationResponse) {}; rpc VerifyFile(VerifyFileRequest) returns (VerifyFileResponse) {}; rpc ReadParts(ReadPartsRequest) returns (ReadPartsResponse) {}; rpc CheckParts(CheckPartsRequest) returns (CheckPartsResponse) {}; diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index 4915a1755..2a1db6b0d 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -124,7 +124,7 @@ tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "sig tokio-rustls = { workspace = true, default-features = false, features = ["logging", "tls12", "aws-lc-rs"] } aws-sdk-s3 = { workspace = true, default-features = false, features = ["sigv4a", "default-https-client", "rt-tokio"] } tokio-stream.workspace = true -tokio-util = { workspace = true, features = ["io", "compat"] } +tokio-util = { workspace = true, features = ["io", "compat", "time"] } tonic = { workspace = true, features = ["gzip", "deflate"] } tower = { workspace = true, features = ["timeout"] } tower-http = { workspace = true, features = ["trace", "compression-full", "cors", "catch-panic", "timeout", "limit", "request-id", "add-extension"] } diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index a7d2b272c..bf2f81f71 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -343,6 +343,7 @@ mod metrics; pub struct NodeService { local_peer: LocalPeerS3Client, context: Option>, + snapshot_lease_expiry: disk::SnapshotLeaseExpiryScheduler, } impl std::fmt::Debug for NodeService { @@ -361,7 +362,11 @@ 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 } + NodeService { + local_peer, + context, + snapshot_lease_expiry: disk::SnapshotLeaseExpiryScheduler::new(), + } } #[derive(Clone, Debug, Default)] @@ -1124,6 +1129,24 @@ impl Node for NodeService { async fn delete_paths(&self, request: Request) -> Result, Status> { self.handle_delete_paths(request).await } + async fn acquire_snapshot_lease( + &self, + request: Request, + ) -> Result, Status> { + self.handle_acquire_snapshot_lease(request).await + } + async fn renew_snapshot_lease( + &self, + request: Request, + ) -> Result, Status> { + self.handle_renew_snapshot_lease(request).await + } + async fn release_snapshot_lease( + &self, + request: Request, + ) -> Result, Status> { + self.handle_release_snapshot_lease(request).await + } async fn read_metadata(&self, request: Request) -> Result, Status> { self.handle_read_metadata(request).await } diff --git a/rustfs/src/storage/rpc/node_service/disk.rs b/rustfs/src/storage/rpc/node_service/disk.rs index e6c3fb60d..ca863d0e5 100644 --- a/rustfs/src/storage/rpc/node_service/disk.rs +++ b/rustfs/src/storage/rpc/node_service/disk.rs @@ -14,11 +14,12 @@ use super::NodeService; use crate::storage::storage_api::rpc_consumer::node_service::{ - BatchReadVersionReq, BatchReadVersionResp, DeleteOptions, DiskError, DiskInfoOptions, FileInfoVersions, ReadMultipleReq, - ReadMultipleResp, ReadOptions, StorageDiskRpcExt as _, UpdateMetadataOpts, validate_batch_read_version_item_count, + BatchReadVersionReq, BatchReadVersionResp, DeleteOptions, DiskError, DiskInfoOptions, DiskStore, 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, verify_tonic_mutation_body_digest}; +use crate::storage::storage_api::{PartTransactionAction, SnapshotLeaseToken, verify_tonic_mutation_body_digest}; use bytes::Bytes; use rustfs_filemeta::FileInfo; use rustfs_io_metrics::internode_metrics::{ @@ -28,13 +29,109 @@ use rustfs_io_metrics::internode_metrics::{ }; use rustfs_protos::proto_gen::node_service::*; use serde::de::DeserializeOwned; -use std::io::Cursor; +use std::{collections::HashMap, io::Cursor, time::Duration}; +use tokio::sync::mpsc; +use tokio_util::time::DelayQueue; use tonic::{Request, Response, Status}; use tracing::debug; /// Initial capacity hint (bytes) for msgpack encode buffers, sized to cover a typical single- /// version `FileInfo` without repeated growth reallocations. Larger payloads still grow as needed. const MSGPACK_ENCODE_CAPACITY_HINT: usize = 512; +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_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) +} + +#[cfg(test)] +mod snapshot_lease_tests { + use super::{SNAPSHOT_LEASE_MAX_TTL, SNAPSHOT_LEASE_MIN_TTL, snapshot_lease_ttl}; + + #[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 decode_msgpack_or_json( binary: &[u8], @@ -165,6 +262,146 @@ fn encode_batch_read_version_response_payloads( } impl NodeService { + pub(super) async fn handle_acquire_snapshot_lease( + &self, + request: Request, + ) -> Result, Status> { + verify_disk_mutation_digest( + &request, + 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()), + })), + } + } + + pub(super) async fn handle_renew_snapshot_lease( + &self, + request: Request, + ) -> Result, Status> { + verify_disk_mutation_digest( + &request, + 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()), + })), + } + } + + pub(super) async fn handle_release_snapshot_lease( + &self, + request: Request, + ) -> Result, Status> { + verify_disk_mutation_digest( + &request, + rustfs_protos::canonical_snapshot_lease_release_request_body(request.get_ref()), + "release_snapshot_lease", + )?; + let request = request.into_inner(); + 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(SnapshotLeaseMutationResponse { + success: false, + error: Some(DiskError::other("cannot find disk").into()), + })); + }; + 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, + })) + } + Err(err) => Ok(Response::new(SnapshotLeaseMutationResponse { + success: false, + error: Some(err.into()), + })), + } + } + pub(super) async fn handle_disk_info(&self, request: Request) -> Result, Status> { let request = request.into_inner(); if let Some(disk) = self.find_disk(&request.disk).await { diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 1735ce2b4..859c9e384 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -431,7 +431,7 @@ pub(crate) mod ecstore_disk { pub(crate) use rustfs_ecstore::api::disk::{ BatchReadVersionReq, BatchReadVersionResp, CheckPartsResp, DeleteOptions, DiskAPI, DiskInfo, DiskInfoOptions, DiskStore, FileInfoVersions, FileReader, FileWriter, OldCurrentSize, PartTransactionAction, RUSTFS_META_BUCKET, ReadMultipleReq, - ReadMultipleResp, ReadOptions, RenameDataResp, UpdateMetadataOpts, VolumeInfo, WalkDirOptions, + ReadMultipleResp, ReadOptions, RenameDataResp, SnapshotLeaseToken, UpdateMetadataOpts, VolumeInfo, WalkDirOptions, get_object_disk_read_timeout, validate_batch_read_version_item_count, }; pub(crate) use rustfs_ecstore::api::disk::{endpoint, error, error_reduce}; @@ -606,6 +606,7 @@ pub(crate) type ExpiryState = ecstore_bucket::lifecycle::bucket_lifecycle_ops::E pub(crate) type FileInfoVersions = ecstore_disk::FileInfoVersions; pub(crate) type FileReader = ecstore_disk::FileReader; pub(crate) type FileWriter = ecstore_disk::FileWriter; +pub(crate) type SnapshotLeaseToken = ecstore_disk::SnapshotLeaseToken; pub(crate) type FS = super::ecfs::FS; pub(crate) type HashReader = ecstore_rio::HashReader; pub(crate) type InstanceContext = ecstore_runtime::InstanceContext; @@ -1075,6 +1076,9 @@ 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>; async fn make_volume(&self, volume: &str) -> DiskResult<()>; @@ -1201,6 +1205,18 @@ 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 + } + async fn stat_volume(&self, volume: &str) -> DiskResult { ecstore_disk::DiskAPI::stat_volume(self, volume).await }