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 <anzhengchao@gmail.com>
This commit is contained in:
cxymds
2026-07-20 12:07:38 +08:00
committed by GitHub
parent 69f543568b
commit a774bc07da
12 changed files with 705 additions and 32 deletions
+17 -1
View File
@@ -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<HealControlServiceClient<InterceptedService<Channel, TonicInterceptor>>, Box<dyn Error>> {
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
+5 -4
View File
@@ -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());
@@ -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<Channel, TonicInterceptor>,
>,
> {
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<u8>) -> Result<Vec<u8>> {
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
@@ -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<T> {
inner: tonic::client::Grpc<T>,
}
impl HealControlServiceClient<tonic::transport::Channel> {
/// Attempt to create a new client by connecting to a given endpoint.
pub async fn connect<D>(dst: D) -> Result<Self, tonic::transport::Error>
where
D: TryInto<tonic::transport::Endpoint>,
D::Error: Into<StdError>,
{
let conn = tonic::transport::Endpoint::new(dst)?.connect().await?;
Ok(Self::new(conn))
}
}
impl<T> HealControlServiceClient<T>
where
T: tonic::client::GrpcService<tonic::body::Body>,
T::Error: Into<StdError>,
T::ResponseBody: Body<Data = Bytes> + std::marker::Send + 'static,
<T::ResponseBody as Body>::Error: Into<StdError> + 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<F>(inner: T, interceptor: F) -> HealControlServiceClient<InterceptedService<T, F>>
where
F: tonic::service::Interceptor,
T::ResponseBody: Default,
T: tonic::codegen::Service<
http::Request<tonic::body::Body>,
Response = http::Response<<T as tonic::client::GrpcService<tonic::body::Body>>::ResponseBody>,
>,
<T as tonic::codegen::Service<http::Request<tonic::body::Body>>>::Error:
Into<StdError> + 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<super::HealControlRequest>,
) -> std::result::Result<tonic::Response<super::HealControlResponse>, 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<super::HealControlRequest>,
) -> std::result::Result<tonic::Response<super::HealControlResponse>, tonic::Status>;
}
#[derive(Debug)]
pub struct HealControlServiceServer<T> {
inner: Arc<T>,
accept_compression_encodings: EnabledCompressionEncodings,
send_compression_encodings: EnabledCompressionEncodings,
max_decoding_message_size: Option<usize>,
max_encoding_message_size: Option<usize>,
}
impl<T> HealControlServiceServer<T> {
pub fn new(inner: T) -> Self {
Self::from_arc(Arc::new(inner))
}
pub fn from_arc(inner: Arc<T>) -> 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<F>(inner: T, interceptor: F) -> InterceptedService<Self, F>
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<T, B> tonic::codegen::Service<http::Request<B>> for HealControlServiceServer<T>
where
T: HealControlService,
B: Body + std::marker::Send + 'static,
B::Error: Into<StdError> + std::marker::Send + 'static,
{
type Response = http::Response<tonic::body::Body>;
type Error = std::convert::Infallible;
type Future = BoxFuture<Self::Response, Self::Error>;
fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<std::result::Result<(), Self::Error>> {
Poll::Ready(Ok(()))
}
fn call(&mut self, req: http::Request<B>) -> Self::Future {
match req.uri().path() {
"/node_service.HealControlService/HealControl" => {
#[allow(non_camel_case_types)]
struct HealControlSvc<T: HealControlService>(pub Arc<T>);
impl<T: HealControlService> tonic::server::UnaryService<super::HealControlRequest> for HealControlSvc<T> {
type Response = super::HealControlResponse;
type Future = BoxFuture<tonic::Response<Self::Response>, tonic::Status>;
fn call(&mut self, request: tonic::Request<super::HealControlRequest>) -> Self::Future {
let inner = Arc::clone(&self.0);
let fut = async move { <T as HealControlService>::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<T> Clone for HealControlServiceServer<T> {
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<T> tonic::server::NamedService for HealControlServiceServer<T> {
const NAME: &'static str = SERVICE_NAME;
}
}
+58
View File
@@ -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<Vec<u8>, 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
+16
View File
@@ -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) {};
}
+93 -7
View File
@@ -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<Request<()>, 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,
+1 -1
View File
@@ -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;
+155 -6
View File
@@ -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<Arc<runtime_sources::AppContext>>
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<HealControlRequest>) -> Result<Response<HealControlResponse>, 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<Arc<ECStore>> {
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<HealControlRequest> {
Request::new(HealControlRequest {
version: 1,
topology_fingerprint: "fingerprint".to_string(),
command: Bytes::copy_from_slice(command),
})
}
fn mark_v2_authenticated<T>(request: &mut Request<T>) {
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<HealControlServiceClient<tonic::transport::Channel>> {
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 {
+15 -4
View File
@@ -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<T>(request: &tonic::Request<T>, 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<T>(request: &mut tonic::Request<T>, 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)
}
+1 -1
View File
@@ -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;
+2 -2
View File
@@ -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};
}
}