From a774bc07daf3975a8705d19d204463d070f9b025 Mon Sep 17 00:00:00 2001 From: cxymds Date: Mon, 20 Jul 2026 12:07:38 +0800 Subject: [PATCH] feat(heal): add authenticated control RPC contract (#4993) * fix(rpc): bind internode auth to exact targets * fix(heal): initialize the runtime atomically * fix(heal): aggregate status across cluster nodes * fix(heal): return canonical tokens for duplicate starts * feat(heal): add authenticated control RPC contract * fix(heal): return canonical tokens for duplicate starts (#4992) --------- Co-authored-by: Zhengchao An --- crates/ecstore/src/cluster/rpc/client.rs | 18 +- crates/ecstore/src/cluster/rpc/http_auth.rs | 9 +- .../src/cluster/rpc/peer_rest_client.rs | 89 +++++- .../src/generated/proto_gen/node_service.rs | 259 ++++++++++++++++++ crates/protos/src/lib.rs | 58 ++++ crates/protos/src/node.proto | 16 ++ rustfs/src/server/http.rs | 100 ++++++- rustfs/src/storage/rpc/mod.rs | 2 +- rustfs/src/storage/rpc/node_service.rs | 161 ++++++++++- rustfs/src/storage/storage_api.rs | 19 +- rustfs/src/storage/tonic_service.rs | 2 +- rustfs/src/storage_api.rs | 4 +- 12 files changed, 705 insertions(+), 32 deletions(-) diff --git a/crates/ecstore/src/cluster/rpc/client.rs b/crates/ecstore/src/cluster/rpc/client.rs index ff3d53731..0bb2946cf 100644 --- a/crates/ecstore/src/cluster/rpc/client.rs +++ b/crates/ecstore/src/cluster/rpc/client.rs @@ -18,7 +18,8 @@ use crate::disk::error::{DiskError, Error as DiskErrorType}; use crate::runtime::sources as runtime_sources; use http::Uri; use rustfs_protos::{ - ChannelClass, create_new_channel, get_channel_for_class, proto_gen::node_service::node_service_client::NodeServiceClient, + ChannelClass, create_new_channel, get_channel_for_class, + proto_gen::node_service::{heal_control_service_client::HealControlServiceClient, node_service_client::NodeServiceClient}, }; use std::{error::Error, io::ErrorKind}; use tonic::{service::interceptor::InterceptedService, transport::Channel}; @@ -37,6 +38,21 @@ pub async fn node_service_time_out_client( node_service_time_out_client_for_class(addr, interceptor, ChannelClass::Control).await } +pub async fn heal_control_time_out_client( + addr: &str, + interceptor: TonicInterceptor, +) -> Result>, Box> { + let interceptor = interceptor.with_rpc_audience(addr)?; + let channel = match runtime_sources::cached_node_channel(addr).await { + Some(channel) => channel, + None => create_new_channel(addr).await?, + }; + let max_message_size = rustfs_protos::HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE; + Ok(HealControlServiceClient::with_interceptor(channel, interceptor) + .max_decoding_message_size(max_message_size) + .max_encoding_message_size(max_message_size)) +} + /// Build a `NodeServiceClient` bound to the [`ChannelClass`]-appropriate channel for `addr`. /// /// Bulk `bytes`-carrying RPCs (ReadAll/WriteAll/ReadMultiple/BatchReadVersion) pass diff --git a/crates/ecstore/src/cluster/rpc/http_auth.rs b/crates/ecstore/src/cluster/rpc/http_auth.rs index 87e4169b5..2dd609443 100644 --- a/crates/ecstore/src/cluster/rpc/http_auth.rs +++ b/crates/ecstore/src/cluster/rpc/http_auth.rs @@ -1104,12 +1104,13 @@ mod tests { .metadata() .get(RPC_CONTENT_SHA256_HEADER) .and_then(|value| value.to_str().ok()); - let headers = gen_tonic_signature_headers("node-a:9000", "node_service.NodeService", "HealControl", content_sha256) - .expect("body-bound auth headers should build"); + let headers = + gen_tonic_signature_headers("node-a:9000", "node_service.HealControlService", "HealControl", content_sha256) + .expect("body-bound auth headers should build"); request.metadata_mut().as_mut().extend(headers.clone()); - assert!(verify_tonic_rpc_signature("node-a:9000", "/node_service.NodeService/HealControl", &headers).is_ok()); - let replay = verify_tonic_rpc_signature("node-a:9000", "/node_service.NodeService/HealControl", &headers) + assert!(verify_tonic_rpc_signature("node-a:9000", "/node_service.HealControlService/HealControl", &headers).is_ok()); + let replay = verify_tonic_rpc_signature("node-a:9000", "/node_service.HealControlService/HealControl", &headers) .expect_err("reusing a body-bound nonce must fail"); assert_eq!(replay.to_string(), "RPC request replay detected"); assert!(verify_tonic_canonical_body_digest(&request, body).is_ok()); diff --git a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs index 619dcb51c..c5b35079f 100644 --- a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs @@ -12,7 +12,10 @@ // See the License for the specific language governing permissions and // limitations under the License. -use crate::cluster::rpc::client::{TonicInterceptor, gen_tonic_signature_interceptor, node_service_time_out_client}; +use crate::cluster::rpc::client::{ + TonicInterceptor, gen_tonic_signature_interceptor, heal_control_time_out_client, node_service_time_out_client, +}; +use crate::cluster::rpc::set_tonic_canonical_body_digest; use crate::error::{Error, Result}; use crate::{ disk::disk_store::{get_drive_active_check_interval, get_drive_active_check_timeout}, @@ -32,11 +35,11 @@ use rustfs_protos::proto_gen::node_service::{ BackgroundHealStatusRequest, CancelDecommissionRequest, ClearDecommissionRequest, DeleteBucketMetadataRequest, DeletePolicyRequest, DeleteServiceAccountRequest, DeleteUserRequest, GetCpusRequest, GetLiveEventsRequest, GetMemInfoRequest, GetMetricsRequest, GetNetInfoRequest, GetOsInfoRequest, GetPartitionsRequest, GetProcInfoRequest, GetSeLinuxInfoRequest, - GetSysConfigRequest, GetSysErrorsRequest, LoadBucketMetadataRequest, LoadGroupRequest, LoadPolicyMappingRequest, - LoadPolicyRequest, LoadRebalanceMetaRequest, LoadServiceAccountRequest, LoadTransitionTierConfigRequest, LoadUserRequest, - LocalStorageInfoRequest, Mss, ReloadPoolMetaRequest, ReloadSiteReplicationConfigRequest, ScannerActivityRequest, - ScannerActivityResponse, ServerInfoRequest, SignalServiceRequest, StartDecommissionRequest, StartProfilingRequest, - StopRebalanceRequest, node_service_client::NodeServiceClient, + GetSysConfigRequest, GetSysErrorsRequest, HealControlRequest, LoadBucketMetadataRequest, LoadGroupRequest, + LoadPolicyMappingRequest, LoadPolicyRequest, LoadRebalanceMetaRequest, LoadServiceAccountRequest, + LoadTransitionTierConfigRequest, LoadUserRequest, LocalStorageInfoRequest, Mss, ReloadPoolMetaRequest, + ReloadSiteReplicationConfigRequest, ScannerActivityRequest, ScannerActivityResponse, ServerInfoRequest, SignalServiceRequest, + StartDecommissionRequest, StartProfilingRequest, StopRebalanceRequest, node_service_client::NodeServiceClient, }; use rustfs_utils::XHost; use serde::{Deserialize, Serialize as _}; @@ -61,6 +64,8 @@ pub const PEER_RESTDRY_RUN: &str = "dry-run"; pub const SERVICE_SIGNAL_REFRESH_CONFIG: u64 = 1; pub const SERVICE_SIGNAL_RELOAD_DYNAMIC: u64 = 2; const BACKGROUND_HEAL_STATUS_MAX_MESSAGE_SIZE: usize = 64 * 1024; +const HEAL_CONTROL_FINGERPRINT_MAX_SIZE: usize = 256; +const HEAL_CONTROL_PAYLOAD_MAX_SIZE: usize = 64 * 1024; const PEER_REST_RECOVERY_MAX_ATTEMPTS: u32 = 60; const PEER_REST_RECOVERY_MAX_BACKOFF: Duration = Duration::from_secs(30); const SCANNER_ACTIVITY_MAX_MESSAGE_SIZE: usize = 1024; @@ -167,6 +172,29 @@ impl PeerRestClient { }) } + async fn get_heal_control_client( + &self, + ) -> Result< + rustfs_protos::proto_gen::node_service::heal_control_service_client::HealControlServiceClient< + InterceptedService, + >, + > { + if self.offline.load(Ordering::Acquire) { + self.mark_offline_and_spawn_recovery(); + return Err(Error::other(format!("peer {} is temporarily offline", self.grid_host))); + } + + heal_control_time_out_client(&self.grid_host, TonicInterceptor::Signature(gen_tonic_signature_interceptor())) + .await + .map_err(|err| { + let storage_err = Error::other(format!("can not get heal control client, err: {err}")); + if Self::is_network_like_error(&storage_err) { + self.mark_offline_and_spawn_recovery(); + } + storage_err + }) + } + /// Evict the connection to this peer from the global cache. /// This should be called when communication with this peer fails. pub async fn evict_connection(&self) { @@ -697,6 +725,43 @@ impl PeerRestClient { .await } + pub async fn heal_control(&self, version: u32, topology_fingerprint: String, command: Vec) -> Result> { + if topology_fingerprint.len() > HEAL_CONTROL_FINGERPRINT_MAX_SIZE { + return Err(Error::other("heal control topology fingerprint exceeds size limit")); + } + if command.len() > HEAL_CONTROL_PAYLOAD_MAX_SIZE { + return Err(Error::other("heal control command exceeds size limit")); + } + self.finalize_result( + async { + let mut client = self + .get_heal_control_client() + .await? + .max_encoding_message_size(rustfs_protos::HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE) + .max_decoding_message_size(rustfs_protos::HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE); + let canonical_body = rustfs_protos::canonical_heal_control_request_body(version, &topology_fingerprint, &command) + .map_err(|_| Error::other("heal control request length cannot be represented"))?; + let mut request = Request::new(HealControlRequest { + version, + topology_fingerprint, + command: command.into(), + }); + set_tonic_canonical_body_digest(&mut request, &canonical_body)?; + let response = client.heal_control(request).await?.into_inner(); + if !response.success { + return Err(Error::other( + response + .error_info + .unwrap_or_else(|| "peer heal control failed without an error".to_string()), + )); + } + Ok(response.result.to_vec()) + } + .await, + ) + .await + } + pub async fn load_bucket_metadata(&self, bucket: &str, scanner_maintenance_change: bool) -> Result<()> { self.finalize_result( async { @@ -1297,6 +1362,18 @@ mod tests { assert!(err.to_string().contains("temporarily offline")); } + #[tokio::test] + async fn peer_rest_client_rejects_oversized_heal_control_before_dialing() { + let client = test_peer_client(); + let err = client + .heal_control(1, "fingerprint".to_string(), vec![0; HEAL_CONTROL_PAYLOAD_MAX_SIZE + 1]) + .await + .expect_err("oversized heal control payload must fail locally"); + + assert!(err.to_string().contains("exceeds size limit")); + assert!(!client.offline.load(Ordering::Acquire)); + } + #[tokio::test] async fn peer_rest_client_prepare_retry_clears_offline_gate() { // finalize_result sets the offline gate on a network error; without diff --git a/crates/protos/src/generated/proto_gen/node_service.rs b/crates/protos/src/generated/proto_gen/node_service.rs index 795ac3e5b..f637b09b1 100644 --- a/crates/protos/src/generated/proto_gen/node_service.rs +++ b/crates/protos/src/generated/proto_gen/node_service.rs @@ -586,6 +586,8 @@ pub struct DeleteVersionRequest { pub force_del_marker: bool, #[prost(string, tag = "6")] pub opts: ::prost::alloc::string::String, + /// msgpack payloads mirroring the JSON fields above (grpc-optimization P2). Senders dual-write + /// both; receivers prefer the \*\_bin form and fall back to the JSON string when it is empty. #[prost(bytes = "bytes", tag = "7")] pub file_info_bin: ::prost::bytes::Bytes, #[prost(bytes = "bytes", tag = "8")] @@ -610,6 +612,8 @@ pub struct DeleteVersionsRequest { pub versions: ::prost::alloc::vec::Vec<::prost::alloc::string::String>, #[prost(string, tag = "4")] pub opts: ::prost::alloc::string::String, + /// msgpack payloads mirroring the JSON fields above (grpc-optimization P2). Senders dual-write + /// both; receivers prefer the \*\_bin form and fall back to the JSON string when it is empty. #[prost(bytes = "bytes", repeated, tag = "5")] pub versions_bin: ::prost::alloc::vec::Vec<::prost::bytes::Bytes>, #[prost(bytes = "bytes", tag = "6")] @@ -650,6 +654,9 @@ pub struct DeleteVolumeRequest { pub disk: ::prost::alloc::string::String, #[prost(string, tag = "2")] pub volume: ::prost::alloc::string::String, + /// When false (default), a non-empty volume is refused (VolumeNotEmpty); when + /// true, the volume is deleted recursively. Old peers omit this field, which + /// decodes to false — the safe, non-recursive behavior (backlog#799 B1). #[prost(bool, tag = "3")] pub force: bool, } @@ -1091,6 +1098,24 @@ pub struct BackgroundHealStatusResponse { pub error_info: ::core::option::Option<::prost::alloc::string::String>, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct HealControlRequest { + #[prost(uint32, tag = "1")] + pub version: u32, + #[prost(string, tag = "2")] + pub topology_fingerprint: ::prost::alloc::string::String, + #[prost(bytes = "bytes", tag = "3")] + pub command: ::prost::bytes::Bytes, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct HealControlResponse { + #[prost(bool, tag = "1")] + pub success: bool, + #[prost(bytes = "bytes", tag = "2")] + pub result: ::prost::bytes::Bytes, + #[prost(string, optional, tag = "3")] + pub error_info: ::core::option::Option<::prost::alloc::string::String>, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct GetMetacacheListingRequest { #[prost(bytes = "bytes", tag = "1")] pub opts: ::prost::bytes::Bytes, @@ -5329,3 +5354,237 @@ pub mod node_service_server { const NAME: &'static str = SERVICE_NAME; } } +/// Generated client implementations. +pub mod heal_control_service_client { + #![allow(unused_variables, dead_code, missing_docs, clippy::wildcard_imports, clippy::let_unit_value)] + use tonic::codegen::http::Uri; + use tonic::codegen::*; + #[derive(Debug, Clone)] + pub struct HealControlServiceClient { + inner: tonic::client::Grpc, + } + impl HealControlServiceClient { + /// Attempt to create a new client by connecting to a given endpoint. + pub async fn connect(dst: D) -> Result + where + D: TryInto, + D::Error: Into, + { + let conn = tonic::transport::Endpoint::new(dst)?.connect().await?; + Ok(Self::new(conn)) + } + } + impl HealControlServiceClient + where + T: tonic::client::GrpcService, + T::Error: Into, + T::ResponseBody: Body + std::marker::Send + 'static, + ::Error: Into + std::marker::Send, + { + pub fn new(inner: T) -> Self { + let inner = tonic::client::Grpc::new(inner); + Self { inner } + } + pub fn with_origin(inner: T, origin: Uri) -> Self { + let inner = tonic::client::Grpc::with_origin(inner, origin); + Self { inner } + } + pub fn with_interceptor(inner: T, interceptor: F) -> HealControlServiceClient> + where + F: tonic::service::Interceptor, + T::ResponseBody: Default, + T: tonic::codegen::Service< + http::Request, + Response = http::Response<>::ResponseBody>, + >, + >>::Error: + Into + std::marker::Send + std::marker::Sync, + { + HealControlServiceClient::new(InterceptedService::new(inner, interceptor)) + } + /// Compress requests with the given encoding. + /// + /// This requires the server to support it otherwise it might respond with an + /// error. + #[must_use] + pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self { + self.inner = self.inner.send_compressed(encoding); + self + } + /// Enable decompressing responses. + #[must_use] + pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self { + self.inner = self.inner.accept_compressed(encoding); + self + } + /// Limits the maximum size of a decoded message. + /// + /// Default: `4MB` + #[must_use] + pub fn max_decoding_message_size(mut self, limit: usize) -> Self { + self.inner = self.inner.max_decoding_message_size(limit); + self + } + /// Limits the maximum size of an encoded message. + /// + /// Default: `usize::MAX` + #[must_use] + pub fn max_encoding_message_size(mut self, limit: usize) -> Self { + self.inner = self.inner.max_encoding_message_size(limit); + self + } + pub async fn heal_control( + &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.HealControlService/HealControl"); + let mut req = request.into_request(); + req.extensions_mut() + .insert(GrpcMethod::new("node_service.HealControlService", "HealControl")); + self.inner.unary(req, path, codec).await + } + } +} +/// Generated server implementations. +pub mod heal_control_service_server { + #![allow(unused_variables, dead_code, missing_docs, clippy::wildcard_imports, clippy::let_unit_value)] + use tonic::codegen::*; + /// Generated trait containing gRPC methods that should be implemented for use with HealControlServiceServer. + #[async_trait] + pub trait HealControlService: std::marker::Send + std::marker::Sync + 'static { + async fn heal_control( + &self, + request: tonic::Request, + ) -> std::result::Result, tonic::Status>; + } + #[derive(Debug)] + pub struct HealControlServiceServer { + inner: Arc, + accept_compression_encodings: EnabledCompressionEncodings, + send_compression_encodings: EnabledCompressionEncodings, + max_decoding_message_size: Option, + max_encoding_message_size: Option, + } + impl HealControlServiceServer { + pub fn new(inner: T) -> Self { + Self::from_arc(Arc::new(inner)) + } + pub fn from_arc(inner: Arc) -> Self { + Self { + inner, + accept_compression_encodings: Default::default(), + send_compression_encodings: Default::default(), + max_decoding_message_size: None, + max_encoding_message_size: None, + } + } + pub fn with_interceptor(inner: T, interceptor: F) -> InterceptedService + where + F: tonic::service::Interceptor, + { + InterceptedService::new(Self::new(inner), interceptor) + } + /// Enable decompressing requests with the given encoding. + #[must_use] + pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self { + self.accept_compression_encodings.enable(encoding); + self + } + /// Compress responses with the given encoding, if the client supports it. + #[must_use] + pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self { + self.send_compression_encodings.enable(encoding); + self + } + /// Limits the maximum size of a decoded message. + /// + /// Default: `4MB` + #[must_use] + pub fn max_decoding_message_size(mut self, limit: usize) -> Self { + self.max_decoding_message_size = Some(limit); + self + } + /// Limits the maximum size of an encoded message. + /// + /// Default: `usize::MAX` + #[must_use] + pub fn max_encoding_message_size(mut self, limit: usize) -> Self { + self.max_encoding_message_size = Some(limit); + self + } + } + impl tonic::codegen::Service> for HealControlServiceServer + where + T: HealControlService, + B: Body + std::marker::Send + 'static, + B::Error: Into + std::marker::Send + 'static, + { + type Response = http::Response; + type Error = std::convert::Infallible; + type Future = BoxFuture; + fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } + fn call(&mut self, req: http::Request) -> Self::Future { + match req.uri().path() { + "/node_service.HealControlService/HealControl" => { + #[allow(non_camel_case_types)] + struct HealControlSvc(pub Arc); + impl tonic::server::UnaryService for HealControlSvc { + type Response = super::HealControlResponse; + type Future = BoxFuture, tonic::Status>; + fn call(&mut self, request: tonic::Request) -> Self::Future { + let inner = Arc::clone(&self.0); + let fut = async move { ::heal_control(&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 = HealControlSvc(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) + } + _ => Box::pin(async move { + let mut response = http::Response::new(tonic::body::Body::default()); + let headers = response.headers_mut(); + headers.insert(tonic::Status::GRPC_STATUS, (tonic::Code::Unimplemented as i32).into()); + headers.insert(http::header::CONTENT_TYPE, tonic::metadata::GRPC_CONTENT_TYPE); + Ok(response) + }), + } + } + } + impl Clone for HealControlServiceServer { + fn clone(&self) -> Self { + let inner = self.inner.clone(); + Self { + inner, + accept_compression_encodings: self.accept_compression_encodings, + send_compression_encodings: self.send_compression_encodings, + max_decoding_message_size: self.max_decoding_message_size, + max_encoding_message_size: self.max_encoding_message_size, + } + } + } + /// Generated gRPC service name + pub const SERVICE_NAME: &str = "node_service.HealControlService"; + impl tonic::server::NamedService for HealControlServiceServer { + const NAME: &'static str = SERVICE_NAME; + } +} diff --git a/crates/protos/src/lib.rs b/crates/protos/src/lib.rs index f51e298c5..2ae7fb03e 100644 --- a/crates/protos/src/lib.rs +++ b/crates/protos/src/lib.rs @@ -140,6 +140,64 @@ pub fn internode_rpc_max_message_size() -> usize { rustfs_utils::get_env_usize(rustfs_config::ENV_INTERNODE_RPC_MAX_MESSAGE_SIZE, DEFAULT_GRPC_SERVER_MESSAGE_LEN) } +pub const HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE: usize = 65 * 1024; + +/// Builds the stable byte representation authenticated for a heal-control request. +/// +/// This deliberately does not reuse protobuf encoding: mixed-version peers may +/// retain unknown protobuf fields, which must not change the signed contract. +pub fn canonical_heal_control_request_body( + version: u32, + topology_fingerprint: &str, + command: &[u8], +) -> Result, std::num::TryFromIntError> { + const DOMAIN: &[u8] = b"rustfs-heal-control-v1\0"; + + let fingerprint = topology_fingerprint.as_bytes(); + let mut body = Vec::with_capacity(DOMAIN.len() + 4 + 8 + fingerprint.len() + 8 + command.len()); + body.extend_from_slice(DOMAIN); + body.extend_from_slice(&version.to_be_bytes()); + body.extend_from_slice(&u64::try_from(fingerprint.len())?.to_be_bytes()); + body.extend_from_slice(fingerprint); + body.extend_from_slice(&u64::try_from(command.len())?.to_be_bytes()); + body.extend_from_slice(command); + Ok(body) +} + +#[cfg(test)] +mod heal_control_tests { + use super::canonical_heal_control_request_body; + + #[test] + fn canonical_heal_control_body_binds_every_field_and_boundary() { + let baseline = canonical_heal_control_request_body(1, "ab", b"c").expect("small request should encode"); + let mut golden = b"rustfs-heal-control-v1\0".to_vec(); + golden.extend_from_slice(&1_u32.to_be_bytes()); + golden.extend_from_slice(&2_u64.to_be_bytes()); + golden.extend_from_slice(b"ab"); + golden.extend_from_slice(&1_u64.to_be_bytes()); + golden.extend_from_slice(b"c"); + assert_eq!(baseline, golden); + + assert_ne!( + baseline, + canonical_heal_control_request_body(2, "ab", b"c").expect("small request should encode") + ); + assert_ne!( + baseline, + canonical_heal_control_request_body(1, "ac", b"c").expect("small request should encode") + ); + assert_ne!( + baseline, + canonical_heal_control_request_body(1, "ab", b"d").expect("small request should encode") + ); + assert_ne!( + canonical_heal_control_request_body(1, "ab", b"c").expect("small request should encode"), + canonical_heal_control_request_body(1, "a", b"bc").expect("small request should encode") + ); + } +} + /// Whether internode metadata RPCs should send only the msgpack `_bin` payloads and leave the JSON /// compatibility strings empty (grpc-optimization P2-1). Shared by the client (`remote_disk`) and /// server (`node_service`) send paths. Defaults to `false` (dual-write); see diff --git a/crates/protos/src/node.proto b/crates/protos/src/node.proto index 3543d589b..694c6027f 100644 --- a/crates/protos/src/node.proto +++ b/crates/protos/src/node.proto @@ -773,6 +773,18 @@ message BackgroundHealStatusResponse { optional string error_info = 3; } +message HealControlRequest { + uint32 version = 1; + string topology_fingerprint = 2; + bytes command = 3; +} + +message HealControlResponse { + bool success = 1; + bytes result = 2; + optional string error_info = 3; +} + message GetMetacacheListingRequest { bytes opts = 1; } @@ -966,3 +978,7 @@ service NodeService { rpc LoadTransitionTierConfig(LoadTransitionTierConfigRequest) returns (LoadTransitionTierConfigResponse) {}; rpc GetLiveEvents(GetLiveEventsRequest) returns (GetLiveEventsResponse) {}; } + +service HealControlService { + rpc HealControl(HealControlRequest) returns (HealControlResponse) {}; +} diff --git a/rustfs/src/server/http.rs b/rustfs/src/server/http.rs index ce5772922..3df55d6e7 100644 --- a/rustfs/src/server/http.rs +++ b/rustfs/src/server/http.rs @@ -35,7 +35,7 @@ use crate::server::{ use crate::storage_api::server::http as storage; use crate::storage_api::server::http::request_context::{RequestContext, extract_request_id_from_headers}; use crate::storage_api::server::http::rpc::InternodeRpcService; -use crate::storage_api::server::http::tonic_service::make_server; +use crate::storage_api::server::http::tonic_service::{make_heal_control_server, make_server}; use crate::storage_api::server::http::{ ServerContextSlot, TONIC_RPC_PREFIX, normalize_tonic_rpc_audience, verify_tonic_rpc_signature, }; @@ -55,7 +55,9 @@ use rustfs_common::GlobalReadiness; use rustfs_keystone::KeystoneAuthLayer; #[cfg(feature = "swift")] use rustfs_protocols::SwiftService; -use rustfs_protos::proto_gen::node_service::node_service_server::NodeServiceServer; +use rustfs_protos::proto_gen::node_service::{ + heal_control_service_server::HealControlServiceServer, node_service_server::NodeServiceServer, +}; use rustfs_trusted_proxies::ClientInfo; use rustfs_utils::net::parse_and_resolve_address; use s3s::{ @@ -74,6 +76,7 @@ use std::task::{Context, Poll}; use std::time::Duration; use tokio::net::{TcpListener, TcpStream}; use tokio::sync::{OwnedSemaphorePermit, Semaphore}; +use tonic::service::Routes; use tonic::service::interceptor::InterceptedService; use tonic::{Request, Status}; use tower::{Service, ServiceBuilder}; @@ -124,6 +127,7 @@ const EVENT_HTTP_CONNECTION_DRAIN: &str = "http_connection_drain"; const EVENT_PEER_ADDR_UNAVAILABLE: &str = "peer_addr_unavailable"; const EVENT_RPC_SIGNATURE_VERIFICATION_FAILED: &str = "rpc_signature_verification_failed"; const EVENT_GRPC_TRACE_CONTEXT_PROPAGATION_FAILED: &str = "grpc_trace_context_propagation_failed"; +const HEAL_CONTROL_TONIC_RPC_PATH: &str = "/node_service.HealControlService/HealControl"; static ACTIVE_HTTP_REQUESTS: AtomicU64 = AtomicU64::new(0); @@ -1297,16 +1301,23 @@ fn process_connection( // Align the server codec limit with the client (both default to // `DEFAULT_GRPC_SERVER_MESSAGE_LEN`, 100 MiB) so `bytes`-carrying unary RPCs are not // capped by tonic's 4 MiB default. Env-overridable via RUSTFS_INTERNODE_RPC_MAX_MESSAGE_SIZE. - // The codec size limits live on `NodeServiceServer`, so set them before wrapping the - // service in the auth interceptor (which returns an `InterceptedService` without those - // methods). + // Codec size limits live on the generated service servers, so set them before wrapping + // each service in the auth interceptor. let rpc_max_message_size = rustfs_protos::internode_rpc_max_message_size(); - let rpc_service = RpcRequestPathService::new(InterceptedService::new( + let node_service = InterceptedService::new( NodeServiceServer::new(make_server()) .max_decoding_message_size(rpc_max_message_size) .max_encoding_message_size(rpc_max_message_size), check_auth, - )); + ); + let heal_control_max_message_size = rustfs_protos::HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE; + let heal_control_service = InterceptedService::new( + HealControlServiceServer::new(make_heal_control_server()) + .max_decoding_message_size(heal_control_max_message_size) + .max_encoding_message_size(heal_control_max_message_size), + check_auth, + ); + let rpc_service = RpcRequestPathService::new(Routes::new(node_service).add_service(heal_control_service).prepare()); #[cfg(feature = "swift")] let http_service = SwiftService::new(true, None, s3_service); @@ -1844,6 +1855,7 @@ fn check_auth(req: Request<()>) -> std::result::Result, Status> { .path() .strip_prefix(TONIC_RPC_PREFIX) .and_then(|suffix| suffix.strip_prefix('/')) + .or_else(|| (target.uri.path() == HEAL_CONTROL_TONIC_RPC_PATH).then_some("HealControl")) .filter(|method| !method.is_empty() && !method.contains('/')) .ok_or_else(|| Status::unauthenticated("Invalid RPC request path"))?; debug_assert!(!rpc_method.is_empty()); @@ -2247,6 +2259,32 @@ mod tests { }); assert!(check_auth(request).is_ok(), "matching trusted node audience should authenticate"); + let heal_headers = + storage::gen_tonic_signature_headers("127.0.0.1:9000", "node_service.HealControlService", "HealControl", None) + .expect("heal control auth headers should build"); + let mut heal_request = Request::new(()); + heal_request.metadata_mut().as_mut().extend(heal_headers); + heal_request.extensions_mut().insert(RpcRequestTarget { + uri: "http://127.0.0.1:9000/node_service.HealControlService/HealControl" + .parse() + .expect("heal control URI should parse"), + method: Method::POST, + }); + assert!(check_auth(heal_request).is_ok(), "heal control service path should authenticate"); + + let replay_headers = storage::gen_tonic_signature_headers("127.0.0.1:9000", "node_service.NodeService", "Ping", None) + .expect("node service auth headers should build"); + let mut cross_service_replay = Request::new(()); + cross_service_replay.metadata_mut().as_mut().extend(replay_headers); + cross_service_replay.extensions_mut().insert(RpcRequestTarget { + uri: HEAL_CONTROL_TONIC_RPC_PATH.parse().expect("heal control path should parse"), + method: Method::POST, + }); + assert!( + check_auth(cross_service_replay).is_err(), + "node service signature must not replay to heal control" + ); + rustfs_common::set_global_local_node_name("127.0.0.1:9001").await; let mut replay_to_other_node = Request::new(()); replay_to_other_node.metadata_mut().as_mut().extend(headers.clone()); @@ -2272,6 +2310,54 @@ mod tests { rustfs_common::set_global_local_node_name(&previous_node_name).await; } + #[tokio::test] + #[serial_test::serial] + async fn peer_rest_heal_control_uses_production_auth_and_stays_online_when_disabled() { + let _ = rustfs_credentials::set_global_rpc_secret("rpc-http-test-secret".to_string()); + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("test listener should bind: {err}"), + }; + let addr = listener.local_addr().expect("listener address should be available"); + let previous_node_name = rustfs_common::get_global_local_node_name().await; + rustfs_common::set_global_local_node_name(&addr.to_string()).await; + + let node_service = InterceptedService::new(NodeServiceServer::new(make_server()), check_auth); + let heal_control_service = InterceptedService::new( + HealControlServiceServer::new(make_heal_control_server()) + .max_decoding_message_size(rustfs_protos::HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE) + .max_encoding_message_size(rustfs_protos::HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE), + check_auth, + ); + let service = RpcRequestPathService::new(Routes::new(node_service).add_service(heal_control_service).prepare()); + let server = tokio::spawn(async move { + let (socket, _) = listener.accept().await.expect("test server should accept a connection"); + let builder = ConnBuilder::new(TokioExecutor::new()); + builder + .serve_connection(TokioIo::new(socket), TowerToHyperService::new(service)) + .await + .expect("test connection should complete"); + }); + + let grid_host = format!("http://{addr}"); + let host = rustfs_utils::XHost::try_from(addr.to_string()).expect("test address should resolve"); + let client = storage::PeerRestClient::new(host, grid_host); + for _ in 0..2 { + let error = client + .heal_control(1, "fingerprint".to_string(), b"query".to_vec()) + .await + .expect_err("dormant heal control must report disabled"); + let message = error.to_string(); + assert!(message.contains("routing is not enabled")); + assert!(!message.contains("temporarily offline")); + } + + client.evict_connection().await; + server.abort(); + rustfs_common::set_global_local_node_name(&previous_node_name).await; + } + #[derive(Clone)] struct MarkerService { name: &'static str, diff --git a/rustfs/src/storage/rpc/mod.rs b/rustfs/src/storage/rpc/mod.rs index 38c249061..4b43fdbb8 100644 --- a/rustfs/src/storage/rpc/mod.rs +++ b/rustfs/src/storage/rpc/mod.rs @@ -16,7 +16,7 @@ pub mod http_service; pub mod node_service; pub use http_service::InternodeRpcService; -pub use node_service::{NodeService, make_server}; +pub use node_service::{HealControlRpcService, NodeService, make_heal_control_server, make_server}; use rmp_serde::Serializer; use serde::Serialize; diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index 9886f1f11..45cefd955 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -26,6 +26,7 @@ use crate::storage::storage_api::rpc_consumer::node_service::{ reload_transition_tier_config, }; use crate::storage::storage_api::runtime_sources_consumer::runtime_sources; +use crate::storage::storage_api::verify_tonic_canonical_body_digest; use bytes::Bytes; use futures::Stream; use futures_util::future::join_all; @@ -49,6 +50,8 @@ use tracing::{debug, error, info, warn}; pub(crate) mod heal; const LOG_COMPONENT_STORAGE: &str = "storage"; +const HEAL_CONTROL_FINGERPRINT_MAX_SIZE: usize = 256; +const HEAL_CONTROL_PAYLOAD_MAX_SIZE: usize = 64 * 1024; const LOG_SUBSYSTEM_RPC: &str = "rpc"; const LOG_SUBSYSTEM_REBALANCE: &str = "rebalance"; const EVENT_RPC_REQUEST_REJECTED: &str = "rpc_request_rejected"; @@ -195,6 +198,33 @@ pub fn make_server_for_context(context: Option> NodeService { local_peer, context } } +#[derive(Clone, Copy, Debug, Default)] +pub struct HealControlRpcService; + +pub fn make_heal_control_server() -> HealControlRpcService { + HealControlRpcService +} + +#[tonic::async_trait] +impl heal_control_service_server::HealControlService for HealControlRpcService { + async fn heal_control(&self, request: Request) -> Result, Status> { + if request.get_ref().topology_fingerprint.len() > HEAL_CONTROL_FINGERPRINT_MAX_SIZE + || request.get_ref().command.len() > HEAL_CONTROL_PAYLOAD_MAX_SIZE + { + return Err(Status::invalid_argument("heal control request exceeds size limit")); + } + let body = rustfs_protos::canonical_heal_control_request_body( + request.get_ref().version, + &request.get_ref().topology_fingerprint, + &request.get_ref().command, + ) + .map_err(|_| Status::invalid_argument("heal control request length cannot be represented"))?; + verify_tonic_canonical_body_digest(&request, &body) + .map_err(|err| Status::permission_denied(format!("heal control authentication failed: {err}")))?; + Err(Status::unimplemented("heal control routing is not enabled")) + } +} + impl NodeService { fn resolve_object_store(&self) -> Option> { let context = self.context.clone().or_else(runtime_sources::current_app_context); @@ -1253,10 +1283,12 @@ impl Node for NodeService { #[allow(unused_imports)] mod tests { use super::{ - CollectMetricsOpts, Error, MetricType, Node as _, NodeService, PEER_RESTSIGNAL, PEER_RESTSUB_SYS, - SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC, STORAGE_CLASS_SUB_SYS, - background_rebalance_start_error_message, make_server, scanner_activity_response, stop_rebalance_response, + CollectMetricsOpts, Error, HEAL_CONTROL_PAYLOAD_MAX_SIZE, MetricType, Node as _, NodeService, PEER_RESTSIGNAL, + PEER_RESTSUB_SYS, SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC, STORAGE_CLASS_SUB_SYS, + background_rebalance_start_error_message, make_heal_control_server, make_server, scanner_activity_response, + stop_rebalance_response, }; + use crate::storage::storage_api::set_tonic_canonical_body_digest; use bytes::Bytes; use rustfs_protos::models::PingBodyBuilder; use rustfs_protos::proto_gen::node_service::{ @@ -1266,14 +1298,17 @@ mod tests { GetAllBucketStatsRequest, GetBucketInfoRequest, GetBucketStatsDataRequest, GetCpusRequest, GetMemInfoRequest, GetMetacacheListingRequest, GetMetricsRequest, GetNetInfoRequest, GetOsInfoRequest, GetPartitionsRequest, GetProcInfoRequest, GetSeLinuxInfoRequest, GetSrMetricsDataRequest, GetSysConfigRequest, GetSysErrorsRequest, - HealBucketRequest, ListBucketRequest, ListDirRequest, ListVolumesRequest, LoadBucketMetadataRequest, LoadGroupRequest, - LoadPolicyMappingRequest, LoadPolicyRequest, LoadRebalanceMetaRequest, LoadServiceAccountRequest, + HealBucketRequest, HealControlRequest, ListBucketRequest, ListDirRequest, ListVolumesRequest, LoadBucketMetadataRequest, + LoadGroupRequest, LoadPolicyMappingRequest, LoadPolicyRequest, LoadRebalanceMetaRequest, LoadServiceAccountRequest, LoadTransitionTierConfigRequest, LoadUserRequest, LocalStorageInfoRequest, MakeBucketRequest, MakeVolumeRequest, MakeVolumesRequest, Mss, PingRequest, ReadAllRequest, ReadAtRequest, ReadMultipleRequest, ReadVersionRequest, ReadXlRequest, ReloadPoolMetaRequest, ReloadSiteReplicationConfigRequest, RenameDataRequest, RenameFileRequest, RenamePartRequest, ScannerActivityRequest, ServerInfoRequest, SignalServiceRequest, StartProfilingRequest, StatVolumeRequest, StopRebalanceRequest, UpdateMetacacheListingRequest, UpdateMetadataRequest, VerifyFileRequest, - WriteAllRequest, WriteMetadataRequest, WriteRequest, node_service_client::NodeServiceClient, + WriteAllRequest, WriteMetadataRequest, WriteRequest, + heal_control_service_client::HealControlServiceClient, + heal_control_service_server::{HealControlService as _, HealControlServiceServer}, + node_service_client::NodeServiceClient, node_service_server::NodeServiceServer, }; use std::collections::HashMap; @@ -1292,6 +1327,56 @@ mod tests { assert!(format!("{:?}", service.local_peer).contains("LocalPeerS3Client")); } + fn heal_control_request(command: &[u8]) -> Request { + Request::new(HealControlRequest { + version: 1, + topology_fingerprint: "fingerprint".to_string(), + command: Bytes::copy_from_slice(command), + }) + } + + fn mark_v2_authenticated(request: &mut Request) { + request + .metadata_mut() + .insert("x-rustfs-rpc-auth-version", "2".parse().expect("valid metadata value")); + } + + #[tokio::test] + async fn heal_control_requires_body_bound_auth_before_reporting_disabled() { + let service = make_heal_control_server(); + let unsigned = service + .heal_control(heal_control_request(b"query")) + .await + .expect_err("unsigned request must fail"); + assert_eq!(unsigned.code(), tonic::Code::PermissionDenied); + + let mut tampered = heal_control_request(b"query"); + let other_body = + rustfs_protos::canonical_heal_control_request_body(1, "fingerprint", b"cancel").expect("small request should encode"); + set_tonic_canonical_body_digest(&mut tampered, &other_body).expect("digest metadata should encode"); + mark_v2_authenticated(&mut tampered); + let tampered = service.heal_control(tampered).await.expect_err("tampered request must fail"); + assert_eq!(tampered.code(), tonic::Code::PermissionDenied); + + let mut signed = heal_control_request(b"query"); + let body = + rustfs_protos::canonical_heal_control_request_body(1, "fingerprint", b"query").expect("small request should encode"); + set_tonic_canonical_body_digest(&mut signed, &body).expect("digest metadata should encode"); + mark_v2_authenticated(&mut signed); + let disabled = service.heal_control(signed).await.expect_err("routing must remain disabled"); + assert_eq!(disabled.code(), tonic::Code::Unimplemented); + } + + #[tokio::test] + async fn heal_control_rejects_oversized_command_before_canonical_copy() { + let service = make_heal_control_server(); + let oversized = service + .heal_control(heal_control_request(&vec![0; HEAL_CONTROL_PAYLOAD_MAX_SIZE + 1])) + .await + .expect_err("oversized request must fail"); + assert_eq!(oversized.code(), tonic::Code::InvalidArgument); + } + #[tokio::test] async fn test_ping_success() { let service = create_test_node_service(); @@ -3054,6 +3139,70 @@ mod tests { ) } + async fn connect_test_heal_control_client() -> Option> { + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return None, + Err(err) => panic!("test listener should bind: {err}"), + }; + let addr = listener.local_addr().expect("listener local address should be available"); + + tokio::spawn(async move { + tonic::transport::Server::builder() + .add_service( + HealControlServiceServer::new(make_heal_control_server()) + .max_decoding_message_size(rustfs_protos::HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE) + .max_encoding_message_size(rustfs_protos::HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE), + ) + .serve_with_incoming(TcpListenerStream::new(listener)) + .await + .expect("heal control test server should run"); + }); + + Some( + HealControlServiceClient::connect(format!("http://{addr}")) + .await + .expect("heal control test client should connect"), + ) + } + + #[tokio::test] + async fn heal_control_transport_enforces_codec_limit_and_reaches_disabled_handler() { + let Some(mut client) = connect_test_heal_control_client().await else { + return; + }; + let mut request = heal_control_request(b"query"); + let body = + rustfs_protos::canonical_heal_control_request_body(1, "fingerprint", b"query").expect("small request should encode"); + set_tonic_canonical_body_digest(&mut request, &body).expect("digest metadata should encode"); + mark_v2_authenticated(&mut request); + let disabled = client.heal_control(request).await.expect_err("routing must remain disabled"); + assert_eq!(disabled.code(), tonic::Code::Unimplemented); + + let max_command = vec![0; HEAL_CONTROL_PAYLOAD_MAX_SIZE]; + let mut max_request = heal_control_request(&max_command); + let max_body = rustfs_protos::canonical_heal_control_request_body(1, "fingerprint", &max_command) + .expect("maximum request should encode"); + set_tonic_canonical_body_digest(&mut max_request, &max_body).expect("digest metadata should encode"); + mark_v2_authenticated(&mut max_request); + let disabled = client + .heal_control(max_request) + .await + .expect_err("maximum valid request must reach disabled handler"); + assert_eq!(disabled.code(), tonic::Code::Unimplemented); + + let oversized = Request::new(HealControlRequest { + version: 1, + topology_fingerprint: "fingerprint".to_string(), + command: Bytes::from(vec![0; rustfs_protos::HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE]), + }); + let rejected = client + .heal_control(oversized) + .await + .expect_err("oversized protobuf message must fail in codec"); + assert_eq!(rejected.code(), tonic::Code::OutOfRange); + } + #[tokio::test] async fn test_write_stream_unimplemented() { let Some(mut client) = connect_test_node_service_client().await else { diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index 2cbd9c4d0..326c5a24b 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -319,7 +319,7 @@ pub(crate) mod timeout_wrapper_consumer { } pub(crate) mod tonic_service_consumer { - pub(crate) use super::super::tonic_service::make_server; + pub(crate) use super::super::tonic_service::{make_heal_control_server, make_server}; } #[cfg(test)] @@ -462,13 +462,13 @@ pub(crate) mod ecstore_rio { } pub(crate) mod ecstore_rpc { - #[cfg(test)] - pub(crate) use rustfs_ecstore::api::rpc::gen_tonic_signature_headers; pub(crate) use rustfs_ecstore::api::rpc::{ LocalPeerS3Client, PEER_RESTSIGNAL, PEER_RESTSUB_SYS, PeerRestClient, PeerS3Client, SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC, TONIC_RPC_PREFIX, normalize_tonic_rpc_audience, verify_rpc_signature, - verify_tonic_rpc_signature, + verify_tonic_canonical_body_digest, verify_tonic_rpc_signature, }; + #[cfg(test)] + pub(crate) use rustfs_ecstore::api::rpc::{gen_tonic_signature_headers, set_tonic_canonical_body_digest}; } pub(crate) mod ecstore_object { @@ -571,6 +571,8 @@ pub(crate) type HashReader = ecstore_rio::HashReader; pub(crate) type InstanceContext = ecstore_runtime::InstanceContext; pub(crate) type ServerContextSlot = crate::storage::runtime_sources::ServerContextSlot; pub(crate) type LocalPeerS3Client = ecstore_rpc::LocalPeerS3Client; +#[cfg(test)] +pub(crate) type PeerRestClient = ecstore_rpc::PeerRestClient; pub(crate) type MetricType = ecstore_metrics::MetricType; pub(crate) type ObjectPartInfo = rustfs_filemeta::ObjectPartInfo; pub(crate) type ObjectLockBlockReason = ecstore_bucket::object_lock::objectlock_sys::ObjectLockBlockReason; @@ -1403,6 +1405,15 @@ pub(crate) fn verify_tonic_rpc_signature(audience: &str, path: &str, headers: &h ecstore_rpc::verify_tonic_rpc_signature(audience, path, headers) } +pub(crate) fn verify_tonic_canonical_body_digest(request: &tonic::Request, canonical_body: &[u8]) -> std::io::Result<()> { + ecstore_rpc::verify_tonic_canonical_body_digest(request, canonical_body) +} + +#[cfg(test)] +pub(crate) fn set_tonic_canonical_body_digest(request: &mut tonic::Request, canonical_body: &[u8]) -> std::io::Result<()> { + ecstore_rpc::set_tonic_canonical_body_digest(request, canonical_body) +} + pub(crate) fn to_s3s_etag(etag: &str) -> s3s::dto::ETag { ecstore_client::object_api_utils::to_s3s_etag(etag) } diff --git a/rustfs/src/storage/tonic_service.rs b/rustfs/src/storage/tonic_service.rs index 9fa6d09b4..66f130178 100644 --- a/rustfs/src/storage/tonic_service.rs +++ b/rustfs/src/storage/tonic_service.rs @@ -12,6 +12,6 @@ // See the License for the specific language governing permissions and // limitations under the License. -pub use crate::storage::rpc::make_server; +pub use crate::storage::rpc::{make_heal_control_server, make_server}; #[allow(dead_code)] pub type NodeService = crate::storage::rpc::NodeService; diff --git a/rustfs/src/storage_api.rs b/rustfs/src/storage_api.rs index a395b5a2e..85dca1b29 100644 --- a/rustfs/src/storage_api.rs +++ b/rustfs/src/storage_api.rs @@ -102,7 +102,7 @@ pub(crate) mod server { } #[cfg(test)] - pub(crate) use crate::storage::storage_api::gen_tonic_signature_headers; + pub(crate) use crate::storage::storage_api::{PeerRestClient, gen_tonic_signature_headers}; pub(crate) mod ecfs { pub(crate) type FS = crate::storage::storage_api::FS; @@ -128,7 +128,7 @@ pub(crate) mod server { } pub(crate) mod tonic_service { - pub(crate) use crate::storage::storage_api::tonic_service_consumer::make_server; + pub(crate) use crate::storage::storage_api::tonic_service_consumer::{make_heal_control_server, make_server}; } }