From 0fbb5ba87be24e7f1a04019ebd93bf0052235e47 Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 21 Jul 2026 23:52:46 +0800 Subject: [PATCH] fix(tier): add peer mutation control client (#5096) * fix(tier): lock tier config mutations Co-Authored-By: heihutu * fix(tier): add mutation RPC auth contract (#5082) Co-authored-by: heihutu * fix(tier): add peer mutation handler core (#5084) Co-authored-by: heihutu * fix(tier): add mutation control rpc service (#5087) Co-authored-by: heihutu * fix(tier): recover prepared mutation drains (#5093) Recover prepared tier mutation intent records into the local tier runtime so a restarted peer fails closed before issuing new remote-tier operation leases or conflicting admin publishes. Reconcile the recovered block map on each scan so committed, aborted, or removed intents clear stale local blocks instead of wedging the peer until process restart. Co-authored-by: heihutu * fix(tier): prove zero references before tier removal (#5092) Signed-off-by: houseme Co-authored-by: heihutu * fix(tier): clear peer mutation runtime blocks (#5094) Install prepared mutation runtime blocks when peer prepare requests are applied or replayed so followers fail closed immediately before restart recovery. Clear the in-memory block once peer commit or abort reaches a durable terminal state, including delayed duplicate prepare requests that observe a committed or aborted record. Co-authored-by: heihutu * fix(tier): add peer mutation control client Add signed tier mutation prepare, commit, and abort client calls for peer fanout while preserving the existing protobuf and RPC contract. Verify response proofs before interpreting peer outcomes, fail closed on invalid states, and reject oversized payloads before dialing peers. Co-Authored-By: heihutu --------- Signed-off-by: houseme Co-authored-by: heihutu --- crates/ecstore/src/cluster/rpc/client.rs | 20 +- .../src/cluster/rpc/peer_rest_client.rs | 332 +++++++++++++++++- 2 files changed, 349 insertions(+), 3 deletions(-) diff --git a/crates/ecstore/src/cluster/rpc/client.rs b/crates/ecstore/src/cluster/rpc/client.rs index 0bb2946cf..83a0d50a2 100644 --- a/crates/ecstore/src/cluster/rpc/client.rs +++ b/crates/ecstore/src/cluster/rpc/client.rs @@ -19,7 +19,10 @@ 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::{heal_control_service_client::HealControlServiceClient, node_service_client::NodeServiceClient}, + proto_gen::node_service::{ + heal_control_service_client::HealControlServiceClient, node_service_client::NodeServiceClient, + tier_mutation_control_service_client::TierMutationControlServiceClient, + }, }; use std::{error::Error, io::ErrorKind}; use tonic::{service::interceptor::InterceptedService, transport::Channel}; @@ -53,6 +56,21 @@ pub async fn heal_control_time_out_client( .max_encoding_message_size(max_message_size)) } +pub async fn tier_mutation_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::TIER_MUTATION_RPC_MAX_MESSAGE_SIZE; + Ok(TierMutationControlServiceClient::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/peer_rest_client.rs b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs index dd9ee7a5e..c85b0f206 100644 --- a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs @@ -14,6 +14,7 @@ use crate::cluster::rpc::client::{ TonicInterceptor, gen_tonic_signature_interceptor, heal_control_time_out_client, node_service_time_out_client, + tier_mutation_control_time_out_client, }; use crate::cluster::rpc::{set_tonic_canonical_body_digest, verify_tonic_rpc_response_proof}; use crate::error::{Error, Result}; @@ -23,6 +24,7 @@ use crate::{ runtime::sources as runtime_sources, services::metrics_realtime::{CollectMetricsOpts, MetricType}, }; +use bytes::Bytes; use rmp_serde::{Deserializer, Serializer}; use rustfs_madmin::{ ServerProperties, @@ -30,7 +32,6 @@ use rustfs_madmin::{ metrics::RealtimeMetrics, net::NetInfo, }; -use rustfs_protos::evict_failed_connection; use rustfs_protos::proto_gen::node_service::{ BackgroundHealStatusRequest, CancelDecommissionRequest, ClearDecommissionRequest, DeleteBucketMetadataRequest, DeletePolicyRequest, DeleteServiceAccountRequest, DeleteUserRequest, GetCpusRequest, GetLiveEventsRequest, GetMemInfoRequest, @@ -39,8 +40,11 @@ use rustfs_protos::proto_gen::node_service::{ LoadPolicyMappingRequest, LoadPolicyRequest, LoadRebalanceMetaRequest, LoadServiceAccountRequest, LoadTransitionTierConfigRequest, LoadUserRequest, LocalStorageInfoRequest, Mss, ReloadPoolMetaRequest, ReloadSiteReplicationConfigRequest, ScannerActivityRequest, ScannerActivityResponse, ServerInfoRequest, SignalServiceRequest, - StartDecommissionRequest, StartProfilingRequest, StopRebalanceRequest, node_service_client::NodeServiceClient, + StartDecommissionRequest, StartProfilingRequest, StopRebalanceRequest, TierMutationAbortRequest, TierMutationCommitRequest, + TierMutationControlResponse, TierMutationPeerState, TierMutationPrepareRequest, node_service_client::NodeServiceClient, + tier_mutation_control_service_client::TierMutationControlServiceClient, }; +use rustfs_protos::{TierMutationRpcPhase, evict_failed_connection}; use rustfs_utils::XHost; use serde::{Deserialize, Serialize as _}; use std::{ @@ -57,6 +61,7 @@ use tonic::Request; use tonic::service::interceptor::InterceptedService; use tonic::transport::Channel; use tracing::{debug, info, warn}; +use uuid::Uuid; pub const PEER_RESTSIGNAL: &str = "signal"; pub const PEER_RESTSUB_SYS: &str = "sub-sys"; @@ -119,6 +124,69 @@ pub struct PeerRestClient { recovery_running: Arc, } +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum PeerTierMutationState { + Prepared, + Committed, + Aborted, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct PeerTierMutationOutcome { + pub state: PeerTierMutationState, + pub applied: bool, +} + +fn validate_tier_mutation_response_proof( + version: u32, + phase: TierMutationRpcPhase, + mutation_id: Uuid, + canonical_payload: &[u8], + response: &TierMutationControlResponse, +) -> Result<()> { + let canonical_response = + rustfs_protos::canonical_tier_mutation_rpc_response_body(rustfs_protos::TierMutationRpcResponseProofInput { + version, + phase, + mutation_id, + canonical_payload, + success: response.success, + state: response.state, + applied: response.applied, + error_info: response.error_info.as_deref(), + }) + .map_err(|_| Error::other("tier mutation response length cannot be represented"))?; + verify_tonic_rpc_response_proof(&canonical_response, &response.response_proof) + .map_err(|_| Error::other("peer returned an invalid tier mutation response proof")) +} + +fn decode_tier_mutation_peer_state(state: i32) -> Result { + match TierMutationPeerState::try_from(state).map_err(|_| Error::other("peer returned an invalid tier mutation state"))? { + TierMutationPeerState::Prepared => Ok(PeerTierMutationState::Prepared), + TierMutationPeerState::Committed => Ok(PeerTierMutationState::Committed), + TierMutationPeerState::Aborted => Ok(PeerTierMutationState::Aborted), + TierMutationPeerState::Unspecified => Err(Error::other("peer returned an unspecified tier mutation state")), + } +} + +fn validate_tier_mutation_payload_len(phase: TierMutationRpcPhase, payload_len: usize) -> Result<()> { + let limit = match phase { + TierMutationRpcPhase::Prepare => rustfs_protos::TIER_MUTATION_RPC_MAX_PREPARE_PAYLOAD_SIZE, + TierMutationRpcPhase::Commit => rustfs_protos::TIER_MUTATION_RPC_MAX_COMMIT_PAYLOAD_SIZE, + TierMutationRpcPhase::Abort => { + if payload_len == 0 { + return Ok(()); + } + return Err(Error::other("tier mutation abort payload must be empty")); + } + _ => return Err(Error::other("tier mutation rpc phase is unsupported")), + }; + if payload_len > limit { + return Err(Error::other("tier mutation payload exceeds size limit")); + } + Ok(()) +} + impl PeerRestClient { fn recovery_monitor_span(grid_host: &str) -> tracing::Span { tracing::info_span!( @@ -205,6 +273,25 @@ impl PeerRestClient { }) } + async fn get_tier_mutation_control_client( + &self, + ) -> Result>> { + if self.offline.load(Ordering::Acquire) { + self.mark_offline_and_spawn_recovery(); + return Err(Error::other(format!("peer {} is temporarily offline", self.grid_host))); + } + + tier_mutation_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 tier mutation 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) { @@ -735,6 +822,89 @@ impl PeerRestClient { .await } + pub async fn prepare_tier_mutation(&self, mutation_id: Uuid, canonical_payload: Bytes) -> Result { + self.tier_mutation_control(TierMutationRpcPhase::Prepare, mutation_id, canonical_payload) + .await + } + + pub async fn commit_tier_mutation(&self, mutation_id: Uuid, canonical_payload: Bytes) -> Result { + self.tier_mutation_control(TierMutationRpcPhase::Commit, mutation_id, canonical_payload) + .await + } + + pub async fn abort_tier_mutation(&self, mutation_id: Uuid) -> Result { + self.tier_mutation_control(TierMutationRpcPhase::Abort, mutation_id, Bytes::new()) + .await + } + + async fn tier_mutation_control( + &self, + phase: TierMutationRpcPhase, + mutation_id: Uuid, + canonical_payload: Bytes, + ) -> Result { + validate_tier_mutation_payload_len(phase, canonical_payload.len())?; + let version = rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION; + self.finalize_result( + async { + let mut client = self + .get_tier_mutation_control_client() + .await? + .max_encoding_message_size(rustfs_protos::TIER_MUTATION_RPC_MAX_MESSAGE_SIZE) + .max_decoding_message_size(rustfs_protos::TIER_MUTATION_RPC_MAX_MESSAGE_SIZE); + let canonical_body = + rustfs_protos::canonical_tier_mutation_rpc_body(version, phase, mutation_id, canonical_payload.as_ref()) + .map_err(|_| Error::other("tier mutation request length cannot be represented"))?; + let mutation_id_text = mutation_id.to_string(); + let response = match phase { + TierMutationRpcPhase::Prepare => { + let mut request = Request::new(TierMutationPrepareRequest { + version, + mutation_id: mutation_id_text.clone(), + canonical_payload: canonical_payload.clone(), + }); + set_tonic_canonical_body_digest(&mut request, &canonical_body)?; + client.prepare_tier_mutation(request).await?.into_inner() + } + TierMutationRpcPhase::Commit => { + let mut request = Request::new(TierMutationCommitRequest { + version, + mutation_id: mutation_id_text.clone(), + canonical_payload: canonical_payload.clone(), + }); + set_tonic_canonical_body_digest(&mut request, &canonical_body)?; + client.commit_tier_mutation(request).await?.into_inner() + } + TierMutationRpcPhase::Abort => { + let mut request = Request::new(TierMutationAbortRequest { + version, + mutation_id: mutation_id_text, + canonical_payload: canonical_payload.clone(), + }); + set_tonic_canonical_body_digest(&mut request, &canonical_body)?; + client.abort_tier_mutation(request).await?.into_inner() + } + _ => return Err(Error::other("tier mutation rpc phase is unsupported")), + }; + validate_tier_mutation_response_proof(version, phase, mutation_id, &canonical_payload, &response)?; + if !response.success { + return Err(Error::other( + response + .error_info + .unwrap_or_else(|| "peer tier mutation failed without an error".to_string()), + )); + } + let state = decode_tier_mutation_peer_state(response.state)?; + Ok(PeerTierMutationOutcome { + state, + applied: response.applied, + }) + } + .await, + ) + .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")); @@ -1443,6 +1613,164 @@ mod tests { } } + struct TierMutationResponseFixture<'a> { + version: u32, + phase: TierMutationRpcPhase, + mutation_id: Uuid, + canonical_payload: &'a [u8], + success: bool, + state: i32, + applied: bool, + error_info: Option<&'a str>, + } + + fn signed_tier_mutation_response(input: TierMutationResponseFixture<'_>) -> TierMutationControlResponse { + let canonical = + rustfs_protos::canonical_tier_mutation_rpc_response_body(rustfs_protos::TierMutationRpcResponseProofInput { + version: input.version, + phase: input.phase, + mutation_id: input.mutation_id, + canonical_payload: input.canonical_payload, + success: input.success, + state: input.state, + applied: input.applied, + error_info: input.error_info, + }) + .expect("small tier mutation response should encode"); + let response_proof = + crate::cluster::rpc::sign_tonic_rpc_response_proof(&canonical).expect("tier mutation response should sign"); + TierMutationControlResponse { + success: input.success, + state: input.state, + applied: input.applied, + error_info: input.error_info.map(str::to_string), + response_proof: response_proof.into(), + } + } + + #[test] + fn tier_mutation_response_proof_binds_phase_payload_state_and_error() { + runtime_sources::ensure_test_rpc_secret(); + let mutation_id = Uuid::new_v4(); + let payload = b"tier-mutation-prepare"; + let response = signed_tier_mutation_response(TierMutationResponseFixture { + version: rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, + phase: TierMutationRpcPhase::Prepare, + mutation_id, + canonical_payload: payload, + success: true, + state: TierMutationPeerState::Prepared as i32, + applied: true, + error_info: None, + }); + validate_tier_mutation_response_proof( + rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Prepare, + mutation_id, + payload, + &response, + ) + .expect("matching tier mutation response proof should verify"); + + for tampered in [ + TierMutationControlResponse { + success: false, + error_info: Some("peer failed".to_string()), + ..response.clone() + }, + TierMutationControlResponse { + state: TierMutationPeerState::Committed as i32, + ..response.clone() + }, + TierMutationControlResponse { + applied: false, + ..response.clone() + }, + ] { + let err = validate_tier_mutation_response_proof( + rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Prepare, + mutation_id, + payload, + &tampered, + ) + .expect_err("tampered tier mutation response proof must fail"); + assert!(err.to_string().contains("invalid tier mutation response proof")); + } + + let err = validate_tier_mutation_response_proof( + rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Commit, + mutation_id, + payload, + &response, + ) + .expect_err("response proof must bind request phase"); + assert!(err.to_string().contains("invalid tier mutation response proof")); + } + + #[test] + fn tier_mutation_peer_state_decode_fails_closed() { + assert_eq!( + decode_tier_mutation_peer_state(TierMutationPeerState::Prepared as i32).expect("prepared state should decode"), + PeerTierMutationState::Prepared + ); + assert_eq!( + decode_tier_mutation_peer_state(TierMutationPeerState::Committed as i32).expect("committed state should decode"), + PeerTierMutationState::Committed + ); + assert_eq!( + decode_tier_mutation_peer_state(TierMutationPeerState::Aborted as i32).expect("aborted state should decode"), + PeerTierMutationState::Aborted + ); + assert!(decode_tier_mutation_peer_state(TierMutationPeerState::Unspecified as i32).is_err()); + assert!(decode_tier_mutation_peer_state(99).is_err()); + } + + #[test] + fn tier_mutation_payload_guard_rejects_invalid_lengths() { + validate_tier_mutation_payload_len( + TierMutationRpcPhase::Prepare, + rustfs_protos::TIER_MUTATION_RPC_MAX_PREPARE_PAYLOAD_SIZE, + ) + .expect("max prepare payload should fit"); + assert!( + validate_tier_mutation_payload_len( + TierMutationRpcPhase::Prepare, + rustfs_protos::TIER_MUTATION_RPC_MAX_PREPARE_PAYLOAD_SIZE + 1, + ) + .is_err() + ); + validate_tier_mutation_payload_len( + TierMutationRpcPhase::Commit, + rustfs_protos::TIER_MUTATION_RPC_MAX_COMMIT_PAYLOAD_SIZE, + ) + .expect("max commit payload should fit"); + assert!( + validate_tier_mutation_payload_len( + TierMutationRpcPhase::Commit, + rustfs_protos::TIER_MUTATION_RPC_MAX_COMMIT_PAYLOAD_SIZE + 1, + ) + .is_err() + ); + validate_tier_mutation_payload_len(TierMutationRpcPhase::Abort, 0).expect("empty abort payload should fit"); + assert!(validate_tier_mutation_payload_len(TierMutationRpcPhase::Abort, 1).is_err()); + } + + #[tokio::test] + async fn peer_rest_client_rejects_oversized_tier_prepare_before_dialing() { + let client = test_peer_client(); + let err = client + .prepare_tier_mutation( + Uuid::new_v4(), + Bytes::from(vec![0; rustfs_protos::TIER_MUTATION_RPC_MAX_PREPARE_PAYLOAD_SIZE + 1]), + ) + .await + .expect_err("oversized tier prepare should fail before dialing"); + assert!(err.to_string().contains("tier mutation payload 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