Files
rustfs/crates/ecstore/src/cluster/rpc/peer_rest_client.rs
T
houseme acfeef55ab feat(scanner): add bounded incarnation-scoped ACK receiver (#7182)
* chore(deps): refresh SDKs and pin clock skew regression coverage

Refresh compatible dependencies for Scanner/Heal V2 batch 1 and verify
the production S3 retry/signing path with a deterministic clock.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

* feat(scanner): add bounded incarnation-scoped ACK receiver

Refs rustfs/backlog#2265 and rustfs/backlog#2240.

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

---------

Co-authored-by: heihutu <heihutu@gmail.com>
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-05 08:54:55 +00:00

4274 lines
185 KiB
Rust

// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use crate::cluster::rpc::client::{
AuthenticatedChannel, TonicInterceptor, embedded_tonic_status, gen_tonic_signature_interceptor, heal_control_time_out_client,
is_network_like_status, message_has_network_needle, node_service_time_out_client, tier_mutation_control_time_out_client,
};
use crate::cluster::rpc::{set_tonic_canonical_body_digest, set_tonic_mutation_body_digest, verify_tonic_rpc_response_proof};
use crate::error::{Error, Result};
use crate::storage_api_contracts::internode::{
SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION, SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION, SCANNER_ACTIVITY_PROTOCOL_VERSION,
SCANNER_ACTIVITY_V6_PROTOCOL_VERSION, SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES,
SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION, SCANNER_DIRTY_USAGE_SNAPSHOT_RPC_MAX_MESSAGE_SIZE,
};
use crate::{
bucket::lifecycle::tier_last_day_stats::{DailyAllTierStats, LastDayTierStats, TierDailyStatsWire},
bucket::replication::BucketStats,
disk::disk_store::{get_drive_active_check_interval, get_drive_active_check_timeout},
layout::endpoints::EndpointServerPools,
runtime::sources as runtime_sources,
services::metrics_realtime::{CollectMetricsOpts, MetricType},
};
use bytes::Bytes;
use rmp_serde::{Deserializer, Serializer};
use rustfs_config::{HEAL_SUB_SYS, SCANNER_SUB_SYS};
use rustfs_madmin::{
ServerProperties,
health::{Cpus, MemInfo, OsInfo, Partitions, ProcInfo, SysConfig, SysErrors, SysServices},
metrics::RealtimeMetrics,
net::NetInfo,
};
use rustfs_protos::proto_gen::node_service::{
BackgroundHealStatusRequest, CancelDecommissionRequest, ClearDecommissionRequest, DeleteBucketMetadataRequest,
DeletePolicyRequest, DeleteServiceAccountRequest, DeleteUserRequest, GetBucketStatsDataRequest, GetBucketStatsDataResponse,
GetCpusRequest, GetLiveEventsRequest, GetMemInfoRequest, GetMetricsRequest, GetNetInfoRequest, GetOsInfoRequest,
GetPartitionsRequest, GetProcInfoRequest, GetSeLinuxInfoRequest, GetSysConfigRequest, GetSysErrorsRequest,
HealControlRequest, LoadBucketMetadataRequest, LoadGroupRequest, LoadPolicyMappingRequest, LoadPolicyRequest,
LoadRebalanceMetaRequest, LoadServiceAccountRequest, LoadTransitionTierConfigRequest, LoadUserRequest,
LocalStorageInfoRequest, Mss, ReloadPoolMetaRequest, ReloadSiteReplicationConfigRequest, ReplacementRecoveryStatusRequest,
ScannerActivityRequest, ScannerActivityResponse, ScannerDirtyUsageSnapshotRequest, ScannerDirtyUsageSnapshotResponse,
ScannerPublicationLeaseReleaseRequest, ScannerPublicationLeaseRequest, ScannerPublicationLeaseResponse, ServerInfoRequest,
SignalServiceRequest, SignalServiceResponse, StartDecommissionRequest, StartProfilingRequest, StopRebalanceRequest,
TierDailyStatsRequest, TierMutationAbortRequest, TierMutationCommitRequest, TierMutationControlResponse,
TierMutationFailureClass, TierMutationPeerState, TierMutationPrepareRequest, node_service_client::NodeServiceClient,
tier_mutation_control_service_client::TierMutationControlServiceClient,
};
pub use rustfs_protos::{PEER_RESTDRY_RUN, PEER_RESTSIGNAL, PEER_RESTSUB_SYS};
use rustfs_protos::{TierMutationRpcPhase, evict_failed_connection};
use rustfs_utils::XHost;
use serde::{Deserialize, Serialize as _};
use std::{
collections::{BTreeMap, HashMap},
io::Cursor,
sync::{
Arc, Weak,
atomic::{AtomicBool, Ordering},
},
time::SystemTime,
};
use tokio::{net::TcpStream, time::Duration};
use tonic::Request;
use tonic::service::interceptor::InterceptedService;
use tracing::{debug, info, warn};
use uuid::Uuid;
pub const SERVICE_SIGNAL_REFRESH_CONFIG: u64 = 1;
pub const SERVICE_SIGNAL_RELOAD_DYNAMIC: u64 = 2;
/// Dynamic config subsystem for the cluster-persisted KMS configuration.
///
/// KMS configuration lives in its own cluster object rather than in the server
/// config document, so it is not a `ServerConfig` subsystem; it only shares the
/// reload signal transport.
pub const KMS_SIGNAL_SUBSYSTEM: &str = "kms";
const BACKGROUND_HEAL_STATUS_MAX_MESSAGE_SIZE: usize = 64 * 1024;
const REPLACEMENT_RECOVERY_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;
/// Reserve time for the acquire response's network/clock uncertainty. The
/// server owns the real expiry; this local deadline is intentionally earlier
/// so a coordinator never starts a bounded persistence operation at the edge
/// of a remote lease.
const SCANNER_PUBLICATION_LEASE_SAFETY_MARGIN: Duration = Duration::from_secs(5);
const REPLICATION_STATS_MAX_MESSAGE_SIZE: usize = 8 * 1024 * 1024;
const BUCKET_METADATA_RELOAD_TIMEOUT: Duration = Duration::from_secs(5);
/// Error for a peer that reported `success = false` without an `error_info` payload.
///
/// Same shape as `peer_s3_client::peer_failure_without_details`, over `StorageError`
/// instead of `DiskError`. The message names the operation (and the bucket, where the
/// operation has one) and nothing else, for two reasons:
///
/// - `finalize_result` classifies failures by message substring, so any text matching
/// `message_has_network_needle` would take an answering peer offline and evict its
/// connection over a plain application-level rejection.
/// - Quorum aggregation (`reduce_errs`) buckets `Io` errors by kind plus rendered
/// message, so a per-peer detail such as the peer address would split one shared
/// failure into single-count buckets and downgrade the dominant error.
fn peer_failure_without_details(op: &str, bucket: Option<&str>) -> Error {
match bucket {
Some(bucket) => Error::other(format!("{op}({bucket}): peer returned failure without error details")),
None => Error::other(format!("{op}: peer returned failure without error details")),
}
}
/// Decode a control-plane response failure. Peers at or above the typed
/// `ControlPlaneErrorCode` change (backlog#1845) carry a machine-readable
/// discriminant beside the legacy `error_info` string; prefer it, then fall
/// back to the string, then to the detail-free per-op failure.
/// RUSTFS_COMPAT_TODO(not-initialized-error-code-v1): string fallback for peers that predate the typed wire code. Remove after the minimum supported RustFS peer version always sends error_code.
fn control_plane_failure(op: &str, bucket: Option<&str>, error_code: Option<i32>, error_info: Option<String>) -> Error {
if error_code == Some(rustfs_protos::proto_gen::node_service::ControlPlaneErrorCode::ControlPlaneErrorNotInitialized as i32) {
return Error::RemoteNotInitialized;
}
if error_code == Some(rustfs_protos::proto_gen::node_service::ControlPlaneErrorCode::ControlPlaneErrorInvalidArgument as i32)
{
return Error::InvalidArgument(
"control-plane".to_string(),
op.to_string(),
error_info.unwrap_or_else(|| format!("{op}: peer rejected invalid argument without details")),
);
}
match error_info {
Some(msg) => Error::other(msg),
None => peer_failure_without_details(op, bucket),
}
}
fn decode_bucket_stats_response(response: GetBucketStatsDataResponse) -> Result<BucketStats> {
if !response.success {
return Err(Error::other(
response
.error_info
.unwrap_or_else(|| "peer replication statistics provider is unavailable".to_string()),
));
}
if response.bucket_stats.len() > REPLICATION_STATS_MAX_MESSAGE_SIZE {
return Err(Error::other("peer replication statistics response exceeds size limit"));
}
let mut buf = Deserializer::new(Cursor::new(response.bucket_stats));
let stats = BucketStats::deserialize(&mut buf).map_err(Error::from)?;
if !stats.replication_stats.provider_available {
return Err(Error::other("peer replication statistics provider is unavailable"));
}
Ok(stats)
}
fn validate_signal_service_protocol(sig: u64, sub_sys: &str, protocol_version: u32) -> Result<()> {
// The version stays pinned to DYNAMIC_CONFIG_PROTOCOL_VERSION rather than
// being bumped per subsystem: the comparison is shared, so raising it would
// retire peers that already converge scanner and heal config correctly.
// Subsystems added after a peer was built are rejected by that peer's own
// subsystem allow-list, which surfaces as an explicit failed signal.
if sig == SERVICE_SIGNAL_RELOAD_DYNAMIC
&& matches!(sub_sys, SCANNER_SUB_SYS | HEAL_SUB_SYS | KMS_SIGNAL_SUBSYSTEM)
&& protocol_version < rustfs_protos::DYNAMIC_CONFIG_PROTOCOL_VERSION
{
return Err(Error::other(format!("peer does not support dynamic {sub_sys} config convergence")));
}
Ok(())
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ScannerPeerActivity {
pub instance_id: String,
pub namespace_generation: u64,
pub maintenance_generation: u64,
pub protocol_version: u32,
pub topology_digest: Option<[u8; 32]>,
pub data_movement_active: Option<bool>,
pub dirty_usage_generation: Option<u64>,
pub dirty_usage_pending: Option<bool>,
pub movement_generation: Option<u64>,
pub publication_blocked: Option<bool>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ScannerPeerDirtyUsageSnapshot {
pub instance_id: String,
pub generation: u64,
pub pending_bucket_count: u64,
pub protocol_version: u32,
pub complete: bool,
pub buckets: BTreeMap<String, u64>,
}
fn scanner_instance_id_is_valid(instance_id: &str) -> bool {
instance_id.len() == 32
&& instance_id
.as_bytes()
.iter()
.all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(byte))
}
fn decode_scanner_activity_with_verifier(
response: ScannerActivityResponse,
challenge: &[u8; 16],
verify_proof: impl FnOnce(&[u8], &[u8]) -> Result<()>,
) -> Result<ScannerPeerActivity> {
let instance_id = &response.instance_id;
if !scanner_instance_id_is_valid(instance_id) {
return Err(Error::other("peer returned an invalid scanner activity instance ID"));
}
let (
topology_digest,
data_movement_active,
dirty_usage_generation,
dirty_usage_pending,
movement_generation,
publication_blocked,
) = match response.protocol_version {
// RUSTFS_COMPAT_TODO(ns-scanner-rpc-v3): legacy response fields are unauthenticated. Remove after protocol v0 peers are unsupported.
SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION
if response.topology_digest.is_empty()
&& response.response_proof.is_empty()
&& !response.data_movement_active
&& response.dirty_usage_generation == 0
&& !response.dirty_usage_pending =>
{
(None, None, None, None, None, None)
}
SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION => {
return Err(Error::other("legacy scanner activity peer returned unexpected extended fields"));
}
SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION => {
if response.dirty_usage_generation != 0 || response.dirty_usage_pending {
return Err(Error::other("scanner activity protocol v4 peer returned unauthenticated v5 fields"));
}
let canonical = rustfs_protos::canonical_scanner_activity_v4_response_body(challenge, &response)
.map_err(|_| Error::other("peer scanner activity response is too large to authenticate"))?;
verify_proof(&canonical, &response.response_proof)?;
(
Some(
response
.topology_digest
.as_ref()
.try_into()
.map_err(|_| Error::other("peer returned an invalid scanner topology digest"))?,
),
Some(response.data_movement_active),
None,
None,
None,
None,
)
}
SCANNER_ACTIVITY_V6_PROTOCOL_VERSION => {
if response.dirty_usage_pending && response.dirty_usage_generation == 0 {
return Err(Error::other("scanner activity peer returned pending dirty usage without a generation"));
}
let canonical = rustfs_protos::canonical_scanner_activity_response_body(challenge, &response)
.map_err(|_| Error::other("peer scanner activity response is too large to authenticate"))?;
verify_proof(&canonical, &response.response_proof)?;
(
Some(
response
.topology_digest
.as_ref()
.try_into()
.map_err(|_| Error::other("peer returned an invalid scanner topology digest"))?,
),
Some(response.data_movement_active),
Some(response.dirty_usage_generation),
Some(response.dirty_usage_pending),
None,
None,
)
}
SCANNER_ACTIVITY_PROTOCOL_VERSION => {
if response.dirty_usage_pending && response.dirty_usage_generation == 0 {
return Err(Error::other("scanner activity peer returned pending dirty usage without a generation"));
}
let movement_generation = response
.movement_generation
.ok_or_else(|| Error::other("scanner activity peer omitted its movement generation"))?;
let publication_blocked = response
.publication_blocked
.ok_or_else(|| Error::other("scanner activity peer omitted its publication blocked state"))?;
if movement_generation == u64::MAX {
return Err(Error::other("scanner activity peer exhausted its movement generation"));
}
let canonical = rustfs_protos::canonical_scanner_activity_v7_response_body(challenge, &response)
.map_err(|_| Error::other("scanner activity peer response is too large to authenticate"))?;
verify_proof(&canonical, &response.response_proof)?;
(
Some(
response
.topology_digest
.as_ref()
.try_into()
.map_err(|_| Error::other("peer returned an invalid scanner topology digest"))?,
),
Some(response.data_movement_active),
Some(response.dirty_usage_generation),
Some(response.dirty_usage_pending),
Some(movement_generation),
Some(publication_blocked),
)
}
version => return Err(Error::other(format!("peer returned unsupported scanner activity protocol {version}"))),
};
Ok(ScannerPeerActivity {
instance_id: response.instance_id,
namespace_generation: response.namespace_generation,
maintenance_generation: response.maintenance_generation,
protocol_version: response.protocol_version,
topology_digest,
data_movement_active,
dirty_usage_generation,
dirty_usage_pending,
movement_generation,
publication_blocked,
})
}
fn decode_scanner_activity(response: ScannerActivityResponse, challenge: &[u8; 16]) -> Result<ScannerPeerActivity> {
decode_scanner_activity_with_verifier(response, challenge, |canonical, proof| {
verify_tonic_rpc_response_proof(canonical, proof)
.map_err(|_| Error::other("peer returned an invalid scanner activity response proof"))
})
}
fn decode_scanner_dirty_usage_snapshot_with_verifier(
response: ScannerDirtyUsageSnapshotResponse,
challenge: &[u8; 16],
verify_proof: impl FnOnce(&[u8], &[u8]) -> Result<()>,
) -> Result<ScannerPeerDirtyUsageSnapshot> {
let canonical = rustfs_protos::canonical_scanner_dirty_usage_snapshot_response_body(challenge, &response)
.map_err(|_| Error::other("peer scanner dirty usage snapshot is too large to authenticate"))?;
verify_proof(&canonical, &response.response_proof)?;
if response.protocol_version != SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION {
return Err(Error::other("peer returned unsupported scanner dirty usage snapshot protocol"));
}
if !scanner_instance_id_is_valid(&response.instance_id) {
return Err(Error::other("peer returned an invalid scanner dirty usage snapshot instance ID"));
}
if response.generation == u64::MAX {
return Err(Error::other("peer scanner dirty usage snapshot exhausted its generation"));
}
if response.pending_bucket_count > 0 && response.generation == 0 {
return Err(Error::other("peer scanner dirty usage snapshot has pending buckets without a generation"));
}
if response.buckets.len() > SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES {
return Err(Error::other("peer scanner dirty usage snapshot exceeds the entry limit"));
}
let bucket_count = u64::try_from(response.buckets.len())
.map_err(|_| Error::other("peer scanner dirty usage snapshot entry count cannot be represented"))?;
let max_entries = u64::try_from(SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES).unwrap_or(u64::MAX);
if response.complete {
if response.pending_bucket_count != bucket_count {
return Err(Error::other(
"complete peer scanner dirty usage snapshot has an inconsistent bucket count",
));
}
} else if !response.buckets.is_empty() || response.pending_bucket_count <= max_entries {
return Err(Error::other(
"incomplete peer scanner dirty usage snapshot must represent an entry-limit overflow",
));
}
for pair in response.buckets.windows(2) {
if pair[0].bucket >= pair[1].bucket {
return Err(Error::other("peer scanner dirty usage snapshot buckets are not strictly ordered"));
}
}
for bucket in &response.buckets {
if bucket.bucket.is_empty() {
return Err(Error::other("peer scanner dirty usage snapshot contains an empty bucket name"));
}
if bucket.generation == 0 || bucket.generation > response.generation {
return Err(Error::other("peer scanner dirty usage snapshot contains an invalid bucket generation"));
}
}
Ok(ScannerPeerDirtyUsageSnapshot {
instance_id: response.instance_id,
generation: response.generation,
pending_bucket_count: response.pending_bucket_count,
protocol_version: response.protocol_version,
complete: response.complete,
buckets: response
.buckets
.into_iter()
.map(|bucket| (bucket.bucket, bucket.generation))
.collect(),
})
}
fn decode_scanner_dirty_usage_snapshot(
response: ScannerDirtyUsageSnapshotResponse,
challenge: &[u8; 16],
) -> Result<ScannerPeerDirtyUsageSnapshot> {
decode_scanner_dirty_usage_snapshot_with_verifier(response, challenge, |canonical, proof| {
verify_tonic_rpc_response_proof(canonical, proof)
.map_err(|_| Error::other("peer returned an invalid scanner dirty usage snapshot response proof"))
})
}
fn scanner_activity_protocol_unsupported(err: &Error) -> bool {
matches!(
err,
Error::Io(io_err)
if embedded_tonic_status(io_err).is_some_and(|status| {
status.code() == tonic::Code::FailedPrecondition
&& status.message().starts_with("unsupported scanner activity request protocol")
})
)
}
fn validate_heal_control_capability_proof(canonical_ack: &[u8], proof: &[u8]) -> Result<()> {
verify_tonic_rpc_response_proof(canonical_ack, proof)
.map_err(|_| Error::other("peer returned an invalid heal control capability proof"))
}
fn validate_heal_control_response_proof(canonical_response: &[u8], proof: &[u8]) -> Result<()> {
verify_tonic_rpc_response_proof(canonical_response, proof)
.map_err(|_| Error::other("peer returned an invalid heal control response proof"))
}
fn decode_remote_version_state_capability(expected_member: &str, result: &[u8]) -> Result<Uuid> {
let (topology_member, process_epoch) = rustfs_protos::decode_remote_version_state_capability(result).map_err(Error::other)?;
if topology_member != expected_member {
return Err(Error::other(
"peer returned a remote version state capability for a different topology member",
));
}
let server_epoch =
Uuid::from_slice(process_epoch).map_err(|_| Error::other("peer returned an invalid remote version state epoch"))?;
if server_epoch.is_nil() {
return Err(Error::other("peer returned a nil remote version state epoch"));
}
Ok(server_epoch)
}
fn decode_cross_pool_fence_capability(expected_member: &str, result: &[u8]) -> Result<(u32, Uuid)> {
let version = result
.get(..4)
.and_then(|value| value.try_into().ok())
.map(u32::from_be_bytes)
.ok_or_else(|| Error::other("peer returned an invalid cross-pool fence capability version"))?;
let epoch = decode_remote_version_state_capability(expected_member, &result[4..])?;
Ok((version, epoch))
}
#[derive(Clone, Debug)]
pub struct PeerLiveEventsBatch {
pub events: Vec<u8>,
pub next_sequence: u64,
pub truncated: bool,
}
#[derive(Clone, Debug)]
pub struct ScannerPublicationLease {
pub token: Uuid,
pub movement_generation: u64,
/// Stable storage owner identity. This is distinct from the activity
/// session and is bound into both acquire and release proofs.
pub owner_id: String,
/// Process/session nonce observed by the final activity probe.
pub session_id: String,
pub expires_at: std::time::Instant,
}
impl ScannerPublicationLease {
pub fn is_valid(&self) -> bool {
std::time::Instant::now() < self.expires_at
}
}
fn validate_scanner_publication_lease_response_fields(
response: &ScannerPublicationLeaseResponse,
expected_session_id: &str,
expected_generation: u64,
) -> Result<(Uuid, String)> {
if !response.success {
return Err(Error::other(
response
.error
.as_ref()
.map(|error| error.error_info.clone())
.unwrap_or_else(|| "peer rejected scanner publication lease".to_string()),
));
}
if response.movement_generation != expected_generation {
return Err(Error::other("peer returned a different scanner publication lease generation"));
}
if response.session_id != expected_session_id {
return Err(Error::other("peer returned a different scanner publication lease session"));
}
let owner_id = Uuid::parse_str(&response.owner_id)
.ok()
.filter(|owner_id| !owner_id.is_nil())
.map(|owner_id| owner_id.to_string())
.ok_or_else(|| Error::other("peer returned an invalid scanner publication lease owner"))?;
if response.lease_ttl_ms != crate::store::SCANNER_PUBLICATION_LEASE_TTL_MS {
return Err(Error::other("peer returned an unsupported scanner publication lease TTL"));
}
let token = Uuid::from_slice(response.token.as_ref())
.map_err(|_| Error::other("peer returned an invalid scanner publication lease token"))?;
Ok((token, owner_id))
}
fn scanner_publication_lease_deadline(
request_started: std::time::Instant,
response_received: std::time::Instant,
lease_ttl_ms: u64,
) -> Result<std::time::Instant> {
let lease_window = Duration::from_millis(lease_ttl_ms)
.checked_sub(SCANNER_PUBLICATION_LEASE_SAFETY_MARGIN)
.ok_or_else(|| Error::other("scanner publication lease TTL is shorter than its safety margin"))?;
let elapsed = response_received
.checked_duration_since(request_started)
.ok_or_else(|| Error::other("scanner publication lease response clock moved backwards"))?;
if elapsed >= lease_window {
return Err(Error::other("scanner publication lease response arrived after its safety window"));
}
request_started
.checked_add(lease_window)
.ok_or_else(|| Error::other("scanner publication lease deadline overflowed"))
}
#[derive(Clone, Debug)]
pub struct PeerRestClient {
pub host: XHost,
pub grid_host: String,
topology_member: String,
offline: Arc<AtomicBool>,
recovery_running: Arc<AtomicBool>,
}
#[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,
}
#[derive(Debug, thiserror::Error)]
#[error("{message}")]
struct TierMutationDefinitelyRejected {
message: String,
}
fn tier_mutation_definitely_rejected_error(message: String) -> Error {
Error::other(TierMutationDefinitelyRejected { message })
}
#[cfg(test)]
pub(crate) fn test_tier_mutation_definitely_rejected_error(message: &str) -> Error {
tier_mutation_definitely_rejected_error(message.to_string())
}
pub(crate) fn tier_mutation_error_is_definitely_rejected(error: &Error) -> bool {
matches!(
error,
Error::Io(io_error)
if io_error
.get_ref()
.is_some_and(|source| source.downcast_ref::<TierMutationDefinitelyRejected>().is_some())
)
}
fn validate_tier_mutation_response_proof(
version: u32,
phase: TierMutationRpcPhase,
mutation_id: Uuid,
canonical_payload: &[u8],
response: &TierMutationControlResponse,
) -> Result<()> {
if response.response_proof.len() > rustfs_protos::TIER_MUTATION_RPC_MAX_RESPONSE_PROOF_SIZE {
return Err(Error::other("peer tier mutation response proof exceeds size limit"));
}
if response
.error_info
.as_ref()
.is_some_and(|error| error.len() > rustfs_protos::TIER_MUTATION_RPC_MAX_ERROR_INFO_SIZE)
{
return Err(Error::other("peer tier mutation error response exceeds size limit"));
}
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(),
failure_class: response.failure_class,
})
.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<PeerTierMutationState> {
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 Err(Error::other("tier mutation abort payload is empty"));
}
rustfs_protos::TIER_MUTATION_RPC_MAX_ABORT_PAYLOAD_SIZE
}
_ => 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(())
}
fn tier_mutation_phase_label(phase: TierMutationRpcPhase) -> &'static str {
match phase {
TierMutationRpcPhase::Prepare => "prepare",
TierMutationRpcPhase::Commit => "commit",
TierMutationRpcPhase::Abort => "abort",
_ => "unknown",
}
}
fn tier_mutation_control_status_error(phase: TierMutationRpcPhase, requested_version: u32, status: tonic::Status) -> Error {
let message = format!("peer tier mutation {} RPC failed: {status}", tier_mutation_phase_label(phase));
let legacy_rejection = format!("unsupported tier mutation peer protocol version: {requested_version}");
// RUSTFS_COMPAT_TODO(backlog-2097-tier-mutation-v4-error-text): retain this exact v3-server rejection classifier for mixed-version peers. Remove after every supported peer returns the signed v4 failure class.
if requested_version == rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION
&& status.code() == tonic::Code::FailedPrecondition
&& status.message().as_bytes() == legacy_rejection.as_bytes()
{
return tier_mutation_definitely_rejected_error(message);
}
Error::other(message)
}
fn tier_mutation_failed_response_error(version: u32, failure_class: i32, error_info: Option<String>) -> Error {
let message = error_info.unwrap_or_else(|| "peer tier mutation failed without an error".to_string());
if version == rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION
&& TierMutationFailureClass::try_from(failure_class).ok() == Some(TierMutationFailureClass::PreDispatchRejected)
{
return tier_mutation_definitely_rejected_error(message);
}
// Missing/zero, unknown, and explicit Ambiguous are deliberately the same
// fail-closed result: the coordinator must include this peer in Abort.
Error::other(message)
}
impl PeerRestClient {
fn recovery_monitor_span(grid_host: &str) -> tracing::Span {
tracing::info_span!(
"recovery-monitor",
component = "ecstore",
subsystem = "peer_rest_client",
kind = "peer_rest",
grid_host = %grid_host
)
}
pub fn new(host: XHost, grid_host: String) -> Self {
let topology_member = host.to_string();
Self {
host,
grid_host,
topology_member,
offline: Arc::new(AtomicBool::new(false)),
recovery_running: Arc::new(AtomicBool::new(false)),
}
}
fn parse_topology_host(peer_host_port: &str, grid_host: &str) -> Result<XHost> {
let url = url::Url::parse(grid_host).map_err(|_| Error::other("peer grid host is not a valid URL"))?;
if !matches!(url.scheme(), "http" | "https")
|| !url.username().is_empty()
|| url.password().is_some()
|| url.query().is_some()
|| url.fragment().is_some()
|| url.path() != "/"
{
return Err(Error::other("peer grid host has an invalid URL shape"));
}
let url_host = url.host().ok_or_else(|| Error::other("peer grid host is missing a host"))?;
let topology_host = match url.port() {
Some(port) => format!("{url_host}:{port}"),
None => url_host.to_string(),
};
let explicit_port = url.port();
let name = match url_host {
url::Host::Domain(domain) => domain.to_string(),
url::Host::Ipv4(address) => address.to_string(),
url::Host::Ipv6(address) if explicit_port.is_none() => format!("[{address}]"),
url::Host::Ipv6(address) => address.to_string(),
};
let port = url
.port_or_known_default()
.filter(|port| *port > 0)
.ok_or_else(|| Error::other("peer grid host is missing a valid port"))?;
let host = XHost {
name,
port,
is_port_set: explicit_port.is_some(),
};
if topology_host != peer_host_port {
return Err(Error::other("peer topology host does not match its grid URL"));
}
Ok(host)
}
fn build_clients_from_slots(
slots: Vec<(String, Option<String>, bool)>,
) -> (Vec<Option<Self>>, Vec<Option<Self>>, Vec<String>) {
let mut remote = Vec::with_capacity(slots.len().saturating_sub(1));
let mut all = vec![None; slots.len()];
let mut remote_topology_hosts = Vec::with_capacity(slots.len().saturating_sub(1));
for (idx, (peer_host_port, grid_host, is_local)) in slots.into_iter().enumerate() {
if is_local {
continue;
}
let client = match grid_host {
Some(grid_host) => match Self::parse_topology_host(&peer_host_port, &grid_host) {
Ok(host) => {
let mut client = PeerRestClient::new(host, grid_host);
client.topology_member = peer_host_port.clone();
Some(client)
}
Err(err) => {
warn!(peer = %peer_host_port, "peer topology host parse failed while constructing peer client: {err:?}");
None
}
},
None => {
warn!(peer = %peer_host_port, "grid host is missing while constructing peer client");
None
}
};
all[idx] = client.clone();
remote.push(client);
remote_topology_hosts.push(peer_host_port);
}
(remote, all, remote_topology_hosts)
}
pub async fn new_clients(eps: EndpointServerPools) -> (Vec<Option<Self>>, Vec<Option<Self>>) {
let (remote, all, _) = Self::new_clients_with_topology(eps).await;
(remote, all)
}
pub async fn new_clients_with_topology(eps: EndpointServerPools) -> (Vec<Option<Self>>, Vec<Option<Self>>, Vec<String>) {
if !runtime_sources::setup_is_dist_erasure().await {
return (Vec::new(), Vec::new(), Vec::new());
}
let (remote, all, remote_topology_hosts) = Self::build_clients_from_slots(eps.peer_grid_host_slots_sorted());
if all.len() != remote.len() + 1 {
warn!(
all_hosts = all.len(),
remote_slots = remote.len(),
"Expected number of all hosts to be remote slots + local node"
);
}
(remote, all, remote_topology_hosts)
}
pub async fn get_client(&self) -> Result<NodeServiceClient<InterceptedService<AuthenticatedChannel, TonicInterceptor>>> {
if self.offline.load(Ordering::Acquire) {
self.mark_offline_and_spawn_recovery();
return Err(Error::RemoteClientUnavailable(format!("peer {} is temporarily offline", self.grid_host)));
}
node_service_time_out_client(&self.grid_host, TonicInterceptor::Signature(gen_tonic_signature_interceptor()))
.await
.map_err(|err| {
let storage_err = Error::RemoteClientUnavailable(format!("can not get client, err: {err}"));
if Self::is_network_like_error(&storage_err) {
self.mark_offline_and_spawn_recovery();
}
storage_err
})
}
async fn get_heal_control_client(
&self,
) -> Result<
rustfs_protos::proto_gen::node_service::heal_control_service_client::HealControlServiceClient<
InterceptedService<AuthenticatedChannel, TonicInterceptor>,
>,
> {
if self.offline.load(Ordering::Acquire) {
self.mark_offline_and_spawn_recovery();
return Err(Error::RemoteClientUnavailable(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::RemoteClientUnavailable(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
})
}
async fn get_tier_mutation_control_client(
&self,
) -> Result<TierMutationControlServiceClient<InterceptedService<AuthenticatedChannel, TonicInterceptor>>> {
if self.offline.load(Ordering::Acquire) {
self.mark_offline_and_spawn_recovery();
return Err(Error::RemoteClientUnavailable(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::RemoteClientUnavailable(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) {
evict_failed_connection(&self.grid_host).await;
}
/// Prepare this client for an immediate fresh-connection retry.
///
/// On a network-like failure `finalize_result` both evicts the channel and
/// sets the offline gate, after which `get_client` fast-fails with
/// "temporarily offline" and only the async background recovery monitor
/// would clear the gate (not within this call). So a plain `evict_connection`
/// is not enough to make an in-call retry actually re-dial: the gate still
/// short-circuits it. This drops the cached channel AND clears the gate so
/// the very next `get_client` re-dials. See rustfs/backlog#1049 (P1-B).
pub async fn prepare_retry(&self) {
self.evict_connection().await;
self.offline.store(false, Ordering::Release);
}
/// Whether this failure means the peer is unreachable, so it should be
/// gated offline and its connection evicted.
///
/// RPC failures are classified by their typed gRPC code first
/// (`is_network_like_status`); an application error from a live peer must
/// never take it offline no matter what its message says. The substring
/// fallback only covers failures that exist purely as text, such as the
/// dial errors `get_client` wraps.
pub(crate) fn is_network_like_error(err: &Error) -> bool {
if let Error::Io(io_err) = err
&& let Some(status) = embedded_tonic_status(io_err)
{
return is_network_like_status(status);
}
message_has_network_needle(&err.to_string())
}
fn mark_offline_and_spawn_recovery(&self) {
self.offline.store(true, Ordering::Release);
if self
.recovery_running
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Relaxed)
.is_err()
{
return;
}
let grid_host = self.grid_host.clone();
let offline = Arc::downgrade(&self.offline);
let recovery_running = Arc::downgrade(&self.recovery_running);
// The offline flag and its recovery are the silent half of
// rustfs/backlog#888: log the monitor's start and its success so an
// "offline then back" episode leaves a trace on the observing node.
warn!(
event = "peer_connection_marked_offline",
grid_host = %self.grid_host,
"peer RPC connection marked offline after a network-like failure; starting background recovery monitor"
);
drop(Self::spawn_recovery_monitor(grid_host, offline, recovery_running));
}
fn spawn_recovery_monitor(
grid_host: String,
offline: Weak<AtomicBool>,
recovery_running: Weak<AtomicBool>,
) -> tokio::task::JoinHandle<()> {
let span = Self::recovery_monitor_span(&grid_host);
super::spawn_background_monitor(span, async move {
let mut delay = get_drive_active_check_interval();
let connect_timeout = get_drive_active_check_timeout();
for attempt in 1..=PEER_REST_RECOVERY_MAX_ATTEMPTS {
if offline.strong_count() == 0 || recovery_running.strong_count() == 0 {
return;
}
tokio::time::sleep(delay).await;
if offline.strong_count() == 0 || recovery_running.strong_count() == 0 {
return;
}
if Self::perform_connectivity_check(&grid_host, connect_timeout).await.is_ok() {
let Some(offline) = offline.upgrade() else {
return;
};
let Some(recovery_running) = recovery_running.upgrade() else {
return;
};
offline.store(false, Ordering::Release);
recovery_running.store(false, Ordering::Release);
info!(
event = "peer_connection_recovered",
grid_host = %grid_host,
attempts = attempt,
"peer connectivity restored by background recovery monitor"
);
return;
}
delay = std::cmp::min(delay.saturating_mul(2), PEER_REST_RECOVERY_MAX_BACKOFF);
}
warn!(
grid_host = %grid_host,
attempts = PEER_REST_RECOVERY_MAX_ATTEMPTS,
"peer recovery monitor reached max attempts; will retry on next request"
);
if let Some(recovery_running) = recovery_running.upgrade() {
recovery_running.store(false, Ordering::Release);
}
})
}
#[cfg(test)]
fn spawn_recovery_monitor_log_probe_for_test(&self) -> tokio::sync::oneshot::Receiver<()> {
let (tx, rx) = tokio::sync::oneshot::channel();
let grid_host = self.grid_host.clone();
let span = Self::recovery_monitor_span(&grid_host);
super::spawn_background_monitor(span, async move {
warn!(grid_host = %grid_host, "peer recovery monitor log probe");
let _ = tx.send(());
});
rx
}
async fn perform_connectivity_check(addr: &str, timeout_duration: Duration) -> Result<()> {
let url = url::Url::parse(addr).map_err(|e| Error::other(format!("Invalid URL: {e}")))?;
let Some(host) = url.host_str() else {
return Err(Error::other("No host in URL".to_string()));
};
let port = url.port_or_known_default().unwrap_or(80);
match tokio::time::timeout(timeout_duration, TcpStream::connect((host, port))).await {
Ok(Ok(stream)) => {
drop(stream);
Ok(())
}
_ => Err(Error::other(format!("Cannot connect to {host}:{port}"))),
}
}
async fn finalize_result<T>(&self, result: Result<T>) -> Result<T> {
if let Err(err) = &result
&& Self::is_network_like_error(err)
{
self.mark_offline_and_spawn_recovery();
self.evict_connection().await;
}
result
}
}
impl PeerRestClient {
pub async fn local_storage_info(&self) -> Result<rustfs_madmin::StorageInfo> {
self.finalize_result(self.local_storage_info_inner().await).await
}
async fn local_storage_info_inner(&self) -> Result<rustfs_madmin::StorageInfo> {
let mut client = self.get_client().await?;
let request = Request::new(LocalStorageInfoRequest { metrics: true });
let response = client.local_storage_info(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure(
"local_storage_info",
None,
response.error_code,
response.error_info,
));
}
let data = response.storage_info;
let mut buf = Deserializer::new(Cursor::new(data));
let storage_info: rustfs_madmin::StorageInfo = Deserialize::deserialize(&mut buf)?;
Ok(storage_info)
}
pub async fn server_info(&self) -> Result<ServerProperties> {
self.finalize_result(self.server_info_inner().await).await
}
async fn server_info_inner(&self) -> Result<ServerProperties> {
let mut client = self.get_client().await?;
let request = Request::new(ServerInfoRequest { metrics: true });
let response = client.server_info(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("server_info", None));
}
let data = response.server_properties;
let mut buf = Deserializer::new(Cursor::new(data));
let storage_properties: ServerProperties = Deserialize::deserialize(&mut buf)?;
Ok(storage_properties)
}
pub async fn get_cpus(&self) -> Result<Cpus> {
self.finalize_result(self.get_cpus_inner().await).await
}
async fn get_cpus_inner(&self) -> Result<Cpus> {
let mut client = self.get_client().await?;
let request = Request::new(GetCpusRequest {});
let response = client.get_cpus(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("get_cpus", None));
}
let data = response.cpus;
let mut buf = Deserializer::new(Cursor::new(data));
let cpus: Cpus = Deserialize::deserialize(&mut buf)?;
Ok(cpus)
}
/// This peer's own rolling-day transition counters, per remote tier.
///
/// The response is untrusted peer input: a ring of the wrong width or an
/// unrepresentable clock is rejected here rather than merged, so a corrupt
/// answer makes the node non-reporting instead of silently shifting a
/// cluster total.
pub async fn tier_daily_stats(&self) -> Result<DailyAllTierStats> {
self.finalize_result(self.tier_daily_stats_inner().await).await
}
async fn tier_daily_stats_inner(&self) -> Result<DailyAllTierStats> {
let mut client = self.get_client().await?;
let request = Request::new(TierDailyStatsRequest {});
let response = client.tier_daily_stats(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("tier_daily_stats", None));
}
let mut buf = Deserializer::new(Cursor::new(response.tier_daily_stats));
let wire: HashMap<String, TierDailyStatsWire> = Deserialize::deserialize(&mut buf)?;
wire.into_iter()
.map(|(tier, stats)| {
LastDayTierStats::from_wire(stats)
.map(|stats| (tier, stats))
.map_err(Error::from)
})
.collect()
}
pub async fn get_net_info(&self) -> Result<NetInfo> {
self.finalize_result(self.get_net_info_inner().await).await
}
async fn get_net_info_inner(&self) -> Result<NetInfo> {
let mut client = self.get_client().await?;
let request = Request::new(GetNetInfoRequest {});
let response = client.get_net_info(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("get_net_info", None));
}
let data = response.net_info;
let mut buf = Deserializer::new(Cursor::new(data));
let net_info: NetInfo = Deserialize::deserialize(&mut buf)?;
Ok(net_info)
}
pub async fn get_partitions(&self) -> Result<Partitions> {
self.finalize_result(self.get_partitions_inner().await).await
}
async fn get_partitions_inner(&self) -> Result<Partitions> {
let mut client = self.get_client().await?;
let request = Request::new(GetPartitionsRequest {});
let response = client.get_partitions(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("get_partitions", None));
}
let data = response.partitions;
let mut buf = Deserializer::new(Cursor::new(data));
let partitions: Partitions = Deserialize::deserialize(&mut buf)?;
Ok(partitions)
}
pub async fn get_os_info(&self) -> Result<OsInfo> {
self.finalize_result(self.get_os_info_inner().await).await
}
async fn get_os_info_inner(&self) -> Result<OsInfo> {
let mut client = self.get_client().await?;
let request = Request::new(GetOsInfoRequest {});
let response = client.get_os_info(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("get_os_info", None));
}
let data = response.os_info;
let mut buf = Deserializer::new(Cursor::new(data));
let os_info: OsInfo = Deserialize::deserialize(&mut buf)?;
Ok(os_info)
}
pub async fn get_se_linux_info(&self) -> Result<SysServices> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let request = Request::new(GetSeLinuxInfoRequest {});
let response = client.get_se_linux_info(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("get_se_linux_info", None));
}
let data = response.sys_services;
let mut buf = Deserializer::new(Cursor::new(data));
let sys_services: SysServices = Deserialize::deserialize(&mut buf)?;
Ok(sys_services)
}
.await,
)
.await
}
pub async fn get_sys_config(&self) -> Result<SysConfig> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let request = Request::new(GetSysConfigRequest {});
let response = client.get_sys_config(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("get_sys_config", None));
}
let data = response.sys_config;
let mut buf = Deserializer::new(Cursor::new(data));
let sys_config: SysConfig = Deserialize::deserialize(&mut buf)?;
Ok(sys_config)
}
.await,
)
.await
}
pub async fn get_sys_errors(&self) -> Result<SysErrors> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let request = Request::new(GetSysErrorsRequest {});
let response = client.get_sys_errors(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("get_sys_errors", None));
}
let data = response.sys_errors;
let mut buf = Deserializer::new(Cursor::new(data));
let sys_errors: SysErrors = Deserialize::deserialize(&mut buf)?;
Ok(sys_errors)
}
.await,
)
.await
}
pub async fn get_mem_info(&self) -> Result<MemInfo> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let request = Request::new(GetMemInfoRequest {});
let response = client.get_mem_info(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("get_mem_info", None));
}
let data = response.mem_info;
let mut buf = Deserializer::new(Cursor::new(data));
let mem_info: MemInfo = Deserialize::deserialize(&mut buf)?;
Ok(mem_info)
}
.await,
)
.await
}
pub async fn get_metrics(&self, t: MetricType, opts: &CollectMetricsOpts) -> Result<RealtimeMetrics> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut buf_t = Vec::new();
t.serialize(&mut Serializer::new(&mut buf_t))?;
let mut buf_o = Vec::new();
opts.serialize(&mut Serializer::new(&mut buf_o))?;
let request = Request::new(GetMetricsRequest {
metric_type: buf_t.into(),
opts: buf_o.into(),
});
let response = client.get_metrics(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("get_metrics", None));
}
let data = response.realtime_metrics;
let mut buf = Deserializer::new(Cursor::new(data));
let realtime_metrics: RealtimeMetrics = Deserialize::deserialize(&mut buf)?;
Ok(realtime_metrics)
}
.await,
)
.await
}
pub async fn get_live_events(&self, after_sequence: u64, limit: u32) -> Result<PeerLiveEventsBatch> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let request = Request::new(GetLiveEventsRequest { after_sequence, limit });
let response = client.get_live_events(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("get_live_events", None));
}
Ok(PeerLiveEventsBatch {
events: response.events.to_vec(),
next_sequence: response.next_sequence,
truncated: response.truncated,
})
}
.await,
)
.await
}
pub async fn get_proc_info(&self) -> Result<ProcInfo> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let request = Request::new(GetProcInfoRequest {});
let response = client.get_proc_info(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("get_proc_info", None));
}
let data = response.proc_info;
let mut buf = Deserializer::new(Cursor::new(data));
let proc_info: ProcInfo = Deserialize::deserialize(&mut buf)?;
Ok(proc_info)
}
.await,
)
.await
}
pub async fn start_profiling(&self, profiler: &str) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let request = Request::new(StartProfilingRequest {
profiler: profiler.to_string(),
});
let response = client.start_profiling(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("start_profiling", None));
}
Ok(())
}
.await,
)
.await
}
pub async fn download_profile_data(&self) -> Result<()> {
warn!("download_profile_data is not implemented in PeerRestClient");
Err(Error::NotImplemented)
}
pub async fn get_bucket_stats(&self, bucket: &str) -> Result<BucketStats> {
let response = self
.finalize_result(
async {
let mut client = self
.get_client()
.await?
.max_decoding_message_size(REPLICATION_STATS_MAX_MESSAGE_SIZE);
let response = client
.get_bucket_stats(Request::new(GetBucketStatsDataRequest {
bucket: bucket.to_string(),
}))
.await?
.into_inner();
Ok(response)
}
.await,
)
.await?;
decode_bucket_stats_response(response)
}
pub async fn get_sr_metrics(&self) -> Result<()> {
warn!("get_sr_metrics is not implemented in PeerRestClient");
Err(Error::NotImplemented)
}
pub async fn get_all_bucket_stats(&self) -> Result<()> {
warn!("get_all_bucket_stats is not implemented in PeerRestClient");
Err(Error::NotImplemented)
}
pub async fn background_heal_status(&self) -> Result<Option<Vec<u8>>> {
self.finalize_result(
async {
let mut client = self
.get_client()
.await?
.max_decoding_message_size(BACKGROUND_HEAL_STATUS_MAX_MESSAGE_SIZE);
let response = match client
.background_heal_status(Request::new(BackgroundHealStatusRequest {
protocol_version: rustfs_protos::BACKGROUND_HEAL_STATUS_PROTOCOL_VERSION,
}))
.await
{
Ok(response) => response.into_inner(),
Err(status) if status.code() == tonic::Code::Unimplemented => {
// RUSTFS_COMPAT_TODO(heal-status-rpc-v1): accept old peers without node heal snapshots. Remove after the minimum supported RustFS peer version implements BackgroundHealStatus.
return Ok(None);
}
Err(status) => return Err(status.into()),
};
if !response.success {
return Err(match (response.error_code, response.error_info) {
(None, None) => Error::other("peer background heal status failed without an error"),
(error_code, error_info) => control_plane_failure("background_heal_status", None, error_code, error_info),
});
}
Ok(Some(response.bg_heal_state.to_vec()))
}
.await,
)
.await
}
pub async fn replacement_recovery_status(&self) -> Result<Option<Vec<u8>>> {
self.finalize_result(
async {
let mut client = self
.get_client()
.await?
.max_decoding_message_size(REPLACEMENT_RECOVERY_STATUS_MAX_MESSAGE_SIZE);
let response = match client
.replacement_recovery_status(Request::new(ReplacementRecoveryStatusRequest::default()))
.await
{
Ok(response) => response.into_inner(),
Err(status) if status.code() == tonic::Code::Unimplemented => {
// RUSTFS_COMPAT_TODO(replacement-recovery-status-v1): old peers cannot prove replacement completion during rolling upgrades. Remove after the minimum supported RustFS peer version implements ReplacementRecoveryStatus.
return Ok(None);
}
Err(status) => return Err(status.into()),
};
if !response.success {
return Err(match (response.error_code, response.error_info) {
(None, None) => Error::other("peer replacement recovery status failed without an error"),
(error_code, error_info) => {
control_plane_failure("replacement_recovery_status", None, error_code, error_info)
}
});
}
Ok(Some(response.recovery_status.to_vec()))
}
.await,
)
.await
}
pub async fn prepare_tier_mutation(&self, mutation_id: Uuid, canonical_payload: Bytes) -> Result<PeerTierMutationOutcome> {
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<PeerTierMutationOutcome> {
self.tier_mutation_control(TierMutationRpcPhase::Commit, mutation_id, canonical_payload)
.await
}
pub async fn abort_tier_mutation(
&self,
mutation_id: Uuid,
canonical_prepare_payload: Bytes,
) -> Result<PeerTierMutationOutcome> {
self.tier_mutation_control(TierMutationRpcPhase::Abort, mutation_id, canonical_prepare_payload)
.await
}
async fn tier_mutation_control(
&self,
phase: TierMutationRpcPhase,
mutation_id: Uuid,
canonical_payload: Bytes,
) -> Result<PeerTierMutationOutcome> {
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
.map_err(|status| tier_mutation_control_status_error(phase, version, status))?
.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
.map_err(|status| tier_mutation_control_status_error(phase, version, status))?
.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
.map_err(|status| tier_mutation_control_status_error(phase, version, status))?
.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(tier_mutation_failed_response_error(version, response.failure_class, response.error_info));
}
if version == rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION
&& response.failure_class != TierMutationFailureClass::Unspecified as i32
{
return Err(Error::other("successful peer tier mutation response carried a failure class"));
}
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<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"));
}
let capability_probe = rustfs_protos::is_heal_control_capability_probe(&command);
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: topology_fingerprint.clone(),
command: command.clone().into(),
});
request.set_timeout(rustfs_protos::heal_control_execution_timeout());
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()),
));
}
if !capability_probe {
let canonical_response = rustfs_protos::canonical_heal_control_response_body(
version,
&topology_fingerprint,
&command,
&response.result,
)
.map_err(|_| Error::other("heal control response length cannot be represented"))?;
validate_heal_control_response_proof(&canonical_response, &response.response_proof)?;
}
Ok(response.result.to_vec())
}
.await,
)
.await
}
/// Confirms that a peer supports the current heal-control coordination
/// contract and has the same storage
/// topology. Every non-success response is an error so old or divergent
/// peers cannot be mistaken for compatible ones.
pub async fn probe_heal_control(&self, topology_fingerprint: String) -> Result<()> {
let nonce = uuid::Uuid::new_v4();
let probe = rustfs_protos::heal_control_capability_probe(nonce.as_bytes());
let canonical_ack = rustfs_protos::canonical_heal_control_capability_ack(
rustfs_protos::HEAL_CONTROL_PROTOCOL_VERSION,
&topology_fingerprint,
&probe,
)
.map_err(|_| Error::other("heal control capability acknowledgement length cannot be represented"))?;
let proof = self
.heal_control(rustfs_protos::HEAL_CONTROL_PROTOCOL_VERSION, topology_fingerprint, probe)
.await?;
validate_heal_control_capability_proof(&canonical_ack, &proof)
}
pub async fn probe_remote_version_state(&self, topology_fingerprint: String) -> Result<(String, Uuid)> {
let probe = rustfs_protos::remote_version_state_capability_probe(Uuid::new_v4().as_bytes());
let result = self
.heal_control(rustfs_protos::HEAL_CONTROL_PROTOCOL_VERSION, topology_fingerprint, probe)
.await?;
let epoch = decode_remote_version_state_capability(&self.topology_member, &result)?;
Ok((self.topology_member.clone(), epoch))
}
pub async fn probe_cross_pool_fence(&self, topology_fingerprint: String) -> Result<(String, u32, Uuid)> {
let mut probe = rustfs_protos::CROSS_POOL_FENCE_CAPABILITY_PROBE_PREFIX.to_vec();
probe.extend_from_slice(Uuid::new_v4().as_bytes());
let result = self
.heal_control(rustfs_protos::HEAL_CONTROL_PROTOCOL_VERSION, topology_fingerprint, probe)
.await?;
let (supported_version, epoch) = decode_cross_pool_fence_capability(&self.topology_member, &result)?;
Ok((self.topology_member.clone(), supported_version, epoch))
}
pub async fn load_bucket_metadata(&self, bucket: &str, scanner_maintenance_change: bool) -> Result<()> {
let result = tokio::time::timeout(BUCKET_METADATA_RELOAD_TIMEOUT, async {
let result = self.load_bucket_metadata_once(bucket, scanner_maintenance_change).await;
if let Err(err) = &result
&& Self::is_network_like_error(err)
{
self.prepare_retry().await;
return self.load_bucket_metadata_once(bucket, scanner_maintenance_change).await;
}
result
})
.await
.unwrap_or_else(|_| Err(Error::other(format!("load_bucket_metadata({bucket}) timed out"))));
self.finalize_result(result).await
}
async fn load_bucket_metadata_once(&self, bucket: &str, scanner_maintenance_change: bool) -> Result<()> {
let mut client = self.get_client().await?;
let mut request = Request::new(LoadBucketMetadataRequest {
bucket: bucket.to_string(),
scanner_maintenance_change,
});
set_tonic_mutation_body_digest(&mut request)?;
request.set_timeout(BUCKET_METADATA_RELOAD_TIMEOUT);
let response = client.load_bucket_metadata(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure(
"load_bucket_metadata",
Some(bucket),
response.error_code,
response.error_info,
));
}
Ok(())
}
pub async fn delete_bucket_metadata(&self, bucket: &str) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut request = Request::new(DeleteBucketMetadataRequest {
bucket: bucket.to_string(),
});
set_tonic_mutation_body_digest(&mut request)?;
let response = client.delete_bucket_metadata(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("delete_bucket_metadata", Some(bucket)));
}
Ok(())
}
.await,
)
.await
}
pub async fn delete_policy(&self, policy: &str) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut request = Request::new(DeletePolicyRequest {
policy_name: policy.to_string(),
});
set_tonic_mutation_body_digest(&mut request)?;
let response = client.delete_policy(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure("delete_policy", None, response.error_code, response.error_info));
}
Ok(())
}
.await,
)
.await
}
pub async fn load_policy(&self, policy: &str) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut request = Request::new(LoadPolicyRequest {
policy_name: policy.to_string(),
});
set_tonic_mutation_body_digest(&mut request)?;
let response = client.load_policy(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure("load_policy", None, response.error_code, response.error_info));
}
Ok(())
}
.await,
)
.await
}
pub async fn load_policy_mapping(&self, user_or_group: &str, user_type: u64, is_group: bool) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut request = Request::new(LoadPolicyMappingRequest {
user_or_group: user_or_group.to_string(),
user_type,
is_group,
});
set_tonic_mutation_body_digest(&mut request)?;
let response = client.load_policy_mapping(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure(
"load_policy_mapping",
None,
response.error_code,
response.error_info,
));
}
Ok(())
}
.await,
)
.await
}
pub async fn delete_user(&self, access_key: &str) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut request = Request::new(DeleteUserRequest {
access_key: access_key.to_string(),
});
set_tonic_mutation_body_digest(&mut request)?;
let response = client.delete_user(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure("delete_user", None, response.error_code, response.error_info));
}
Ok(())
}
.await,
)
.await
}
pub async fn delete_service_account(&self, access_key: &str) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut request = Request::new(DeleteServiceAccountRequest {
access_key: access_key.to_string(),
});
set_tonic_mutation_body_digest(&mut request)?;
let response = client.delete_service_account(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure(
"delete_service_account",
None,
response.error_code,
response.error_info,
));
}
Ok(())
}
.await,
)
.await
}
pub async fn load_user(&self, access_key: &str, temp: bool) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut request = Request::new(LoadUserRequest {
access_key: access_key.to_string(),
temp,
});
set_tonic_mutation_body_digest(&mut request)?;
let response = client.load_user(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure("load_user", None, response.error_code, response.error_info));
}
Ok(())
}
.await,
)
.await
}
pub async fn load_service_account(&self, access_key: &str) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut request = Request::new(LoadServiceAccountRequest {
access_key: access_key.to_string(),
});
set_tonic_mutation_body_digest(&mut request)?;
let response = client.load_service_account(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure(
"load_service_account",
None,
response.error_code,
response.error_info,
));
}
Ok(())
}
.await,
)
.await
}
pub async fn load_group(&self, group: &str) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut request = Request::new(LoadGroupRequest {
group: group.to_string(),
});
set_tonic_mutation_body_digest(&mut request)?;
let response = client.load_group(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure("load_group", None, response.error_code, response.error_info));
}
Ok(())
}
.await,
)
.await
}
pub async fn reload_site_replication_config(&self) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut request = Request::new(ReloadSiteReplicationConfigRequest {});
set_tonic_mutation_body_digest(&mut request)?;
let response = client.reload_site_replication_config(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure(
"reload_site_replication_config",
None,
response.error_code,
response.error_info,
));
}
Ok(())
}
.await,
)
.await
}
pub async fn signal_service(&self, sig: u64, sub_sys: &str, dry_run: bool, _exec_at: SystemTime) -> Result<()> {
self.signal_service_checked(sig, sub_sys, dry_run).await.map(|_| ())
}
/// Report the KMS configuration fingerprint the peer is currently running.
///
/// Sent as a dry-run reload signal so the peer answers without swapping its
/// own configuration. `None` means the peer has no KMS configuration. The
/// fingerprint is advisory and feeds cluster status reporting only, so the
/// response is not proof-signed.
pub async fn kms_config_fingerprint(&self) -> Result<Option<String>> {
self.signal_service_checked(SERVICE_SIGNAL_RELOAD_DYNAMIC, KMS_SIGNAL_SUBSYSTEM, true)
.await
.map(|response| response.config_fingerprint)
}
async fn signal_service_checked(&self, sig: u64, sub_sys: &str, dry_run: bool) -> Result<SignalServiceResponse> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut vars = HashMap::new();
vars.insert(PEER_RESTSIGNAL.to_string(), sig.to_string());
vars.insert(PEER_RESTSUB_SYS.to_string(), sub_sys.to_string());
vars.insert(PEER_RESTDRY_RUN.to_string(), dry_run.to_string());
let mut request = Request::new(SignalServiceRequest {
vars: Some(Mss { value: vars }),
});
set_tonic_mutation_body_digest(&mut request)?;
let response = client.signal_service(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("signal_service", None));
}
validate_signal_service_protocol(sig, sub_sys, response.protocol_version)?;
Ok(response)
}
.await,
)
.await
}
async fn scanner_activity_request_with_protocol(
&self,
acknowledge_instance_id: String,
acknowledge_dirty_usage_generation: u64,
protocol_version: u32,
) -> Result<ScannerPeerActivity> {
self.finalize_result(
async {
let challenge = Uuid::new_v4();
let mut client = self
.get_client()
.await?
.max_decoding_message_size(SCANNER_ACTIVITY_MAX_MESSAGE_SIZE)
.max_encoding_message_size(SCANNER_ACTIVITY_MAX_MESSAGE_SIZE);
let mut request = Request::new(ScannerActivityRequest {
challenge: challenge.as_bytes().to_vec().into(),
protocol_version,
acknowledge_instance_id,
acknowledge_dirty_usage_generation,
});
let canonical = rustfs_protos::canonical_scanner_activity_request_body(request.get_ref())
.map_err(|_| Error::other("scanner activity request is too large to authenticate"))?;
set_tonic_canonical_body_digest(&mut request, &canonical)?;
let response = client.scanner_activity(request).await?.into_inner();
decode_scanner_activity(response, challenge.as_bytes())
}
.await,
)
.await
}
pub async fn scanner_activity(&self) -> Result<ScannerPeerActivity> {
let result = self
.scanner_activity_request_with_protocol(String::new(), 0, SCANNER_ACTIVITY_PROTOCOL_VERSION)
.await;
if result.as_ref().err().is_some_and(scanner_activity_protocol_unsupported) {
// A v6 peer cannot parse the v7 marker. Its authenticated
// response is still decoded as untrusted terminal state, so the
// scanner will defer publication until every peer is v7.
self.scanner_activity_request_with_protocol(String::new(), 0, SCANNER_ACTIVITY_V6_PROTOCOL_VERSION)
.await
} else {
result
}
}
pub async fn scanner_dirty_usage_snapshot(&self) -> Result<ScannerPeerDirtyUsageSnapshot> {
self.finalize_result(
async {
let challenge = Uuid::new_v4();
let mut client = self
.get_client()
.await?
.max_decoding_message_size(SCANNER_DIRTY_USAGE_SNAPSHOT_RPC_MAX_MESSAGE_SIZE)
.max_encoding_message_size(SCANNER_DIRTY_USAGE_SNAPSHOT_RPC_MAX_MESSAGE_SIZE);
let mut request = Request::new(ScannerDirtyUsageSnapshotRequest {
challenge: challenge.as_bytes().to_vec().into(),
protocol_version: SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION,
});
let canonical = rustfs_protos::canonical_scanner_dirty_usage_snapshot_request_body(request.get_ref())
.map_err(|_| Error::other("scanner dirty usage snapshot request is too large to authenticate"))?;
set_tonic_canonical_body_digest(&mut request, &canonical)?;
let response = client.scanner_dirty_usage_snapshot(request).await?.into_inner();
decode_scanner_dirty_usage_snapshot(response, challenge.as_bytes())
}
.await,
)
.await
}
/// Probe only: scoped ACK production requires a durable per-bucket proof.
pub async fn scanner_scoped_dirty_usage_capability(
&self,
owner_id: String,
instance_id: String,
entries: Vec<rustfs_protos::proto_gen::node_service::ScannerScopedDirtyUsageEntry>,
) -> Result<bool> {
use rustfs_protos::scoped_dirty_usage::*;
let payload = rustfs_protos::proto_gen::node_service::ScannerScopedDirtyUsageAckRequest {
challenge: Uuid::new_v4().as_bytes().to_vec().into(),
protocol_version: SCOPED_DIRTY_USAGE_PROTOCOL_VERSION,
owner_id,
instance_id,
scope: SCOPED_DIRTY_USAGE_BUCKET_SCOPE,
probe_only: true,
entries,
};
let canonical = canonical_scoped_dirty_usage_request(&payload).map_err(|err| Error::other(err.to_string()))?;
self.finalize_result(
async {
let mut client = super::client::scanner_control_time_out_client(
&self.grid_host,
TonicInterceptor::Signature(gen_tonic_signature_interceptor()),
)
.await?;
let mut request = Request::new(payload.clone());
set_tonic_canonical_body_digest(&mut request, &canonical)?;
let response = client.scanner_scoped_dirty_usage_ack(request).await?.into_inner();
let body = canonical_scoped_dirty_usage_response(&canonical, &response)
.map_err(|_| Error::other("scoped dirty usage capability response is too large"))?;
verify_tonic_rpc_response_proof(&body, response.response_proof.as_ref())?;
if response.protocol_version != SCOPED_DIRTY_USAGE_PROTOCOL_VERSION
|| response.owner_id != payload.owner_id
|| response.instance_id != payload.instance_id
|| response.max_entries != SCOPED_DIRTY_USAGE_MAX_ENTRIES
|| response.max_request_bytes != SCOPED_DIRTY_USAGE_MAX_REQUEST_BYTES
|| response.cleared != 0
{
return Err(Error::other("scoped dirty usage capability response does not match request"));
}
Ok(response.supported)
}
.await,
)
.await
}
pub async fn acknowledge_scanner_dirty_usage(&self, instance_id: String, generation: u64) -> Result<ScannerPeerActivity> {
let result = self
.scanner_activity_request_with_protocol(instance_id.clone(), generation, SCANNER_ACTIVITY_PROTOCOL_VERSION)
.await;
if result.as_ref().err().is_some_and(scanner_activity_protocol_unsupported) {
self.scanner_activity_request_with_protocol(instance_id, generation, SCANNER_ACTIVITY_V6_PROTOCOL_VERSION)
.await
} else {
result
}
}
/// Acquire a bounded, storage-owned read admission on the peer that
/// produced the final activity generation. Older peers do not implement
/// the lease form and are rejected rather than downgraded.
pub async fn acquire_scanner_publication_lease(
&self,
expected_session_id: &str,
expected_generation: u64,
) -> Result<ScannerPublicationLease> {
let request_started = std::time::Instant::now();
self.finalize_result(
async {
let challenge = Uuid::new_v4();
let mut client = self
.get_client()
.await?
.max_decoding_message_size(SCANNER_ACTIVITY_MAX_MESSAGE_SIZE)
.max_encoding_message_size(SCANNER_ACTIVITY_MAX_MESSAGE_SIZE);
let mut request = Request::new(ScannerPublicationLeaseRequest {
challenge: challenge.as_bytes().to_vec().into(),
expected_movement_generation: expected_generation,
ttl_ms: crate::store::SCANNER_PUBLICATION_LEASE_TTL_MS,
expected_session_id: expected_session_id.to_string(),
token: Bytes::new(),
});
let canonical = rustfs_protos::canonical_scanner_publication_lease_request_body(request.get_ref())
.map_err(|_| Error::other("scanner publication lease request is too large to authenticate"))?;
set_tonic_canonical_body_digest(&mut request, &canonical)?;
let response = client.acquire_scanner_publication_lease(request).await?.into_inner();
let response_body =
rustfs_protos::canonical_scanner_publication_lease_response_body(challenge.as_bytes(), &response)
.map_err(|_| Error::other("scanner publication lease response is too large to authenticate"))?;
verify_tonic_rpc_response_proof(&response_body, &response.response_proof)
.map_err(|_| Error::other("peer returned an invalid scanner publication lease proof"))?;
let (token, owner_id) =
validate_scanner_publication_lease_response_fields(&response, expected_session_id, expected_generation)?;
Ok(ScannerPublicationLease {
token,
movement_generation: response.movement_generation,
owner_id,
session_id: response.session_id,
expires_at: scanner_publication_lease_deadline(
request_started,
std::time::Instant::now(),
response.lease_ttl_ms,
)?,
})
}
.await,
)
.await
}
/// Revalidate the exact token immediately before the coordinator commits
/// its final publication. The peer keeps the original movement read
/// guard in its token table; a restart drops that table and changes the
/// activity session, so this proof fails closed instead of accepting an
/// ABA generation value.
pub async fn validate_scanner_publication_lease(&self, lease: &ScannerPublicationLease) -> Result<()> {
self.finalize_result(
async {
let challenge = Uuid::new_v4();
let mut client = self
.get_client()
.await?
.max_decoding_message_size(SCANNER_ACTIVITY_MAX_MESSAGE_SIZE)
.max_encoding_message_size(SCANNER_ACTIVITY_MAX_MESSAGE_SIZE);
let mut request = Request::new(ScannerPublicationLeaseRequest {
challenge: challenge.as_bytes().to_vec().into(),
expected_movement_generation: lease.movement_generation,
ttl_ms: crate::store::SCANNER_PUBLICATION_LEASE_TTL_MS,
expected_session_id: lease.session_id.clone(),
token: lease.token.as_bytes().to_vec().into(),
});
let canonical = rustfs_protos::canonical_scanner_publication_lease_request_body(request.get_ref())
.map_err(|_| Error::other("scanner publication lease validation request is too large to authenticate"))?;
set_tonic_canonical_body_digest(&mut request, &canonical)?;
let response = client.acquire_scanner_publication_lease(request).await?.into_inner();
let response_body =
rustfs_protos::canonical_scanner_publication_lease_response_body(challenge.as_bytes(), &response).map_err(
|_| Error::other("scanner publication lease validation response is too large to authenticate"),
)?;
verify_tonic_rpc_response_proof(&response_body, &response.response_proof)
.map_err(|_| Error::other("peer returned an invalid scanner publication lease validation proof"))?;
let (token, owner_id) =
validate_scanner_publication_lease_response_fields(&response, &lease.session_id, lease.movement_generation)?;
if token != lease.token {
return Err(Error::other("peer returned a different scanner publication lease token"));
}
if owner_id != lease.owner_id {
return Err(Error::other("peer returned a different scanner publication lease owner"));
}
Ok(())
}
.await,
)
.await
}
pub async fn release_scanner_publication_lease(&self, lease: &ScannerPublicationLease) -> Result<()> {
self.finalize_result(
async {
let challenge = Uuid::new_v4();
let mut client = self.get_client().await?;
let mut request = Request::new(ScannerPublicationLeaseReleaseRequest {
challenge: challenge.as_bytes().to_vec().into(),
token: lease.token.as_bytes().to_vec().into(),
owner_id: lease.owner_id.clone(),
session_id: lease.session_id.clone(),
});
let canonical = rustfs_protos::canonical_scanner_publication_lease_release_request_body(request.get_ref())
.map_err(|_| Error::other("scanner publication lease release request is too large to authenticate"))?;
set_tonic_canonical_body_digest(&mut request, &canonical)?;
let request_body = request.get_ref().clone();
let response = client.release_scanner_publication_lease(request).await?.into_inner();
let response_body = rustfs_protos::canonical_scanner_publication_lease_release_response_body(
challenge.as_bytes(),
&request_body,
&response,
)
.map_err(|_| Error::other("scanner publication lease release response is too large to authenticate"))?;
verify_tonic_rpc_response_proof(&response_body, &response.response_proof)
.map_err(|_| Error::other("peer returned an invalid scanner publication lease release proof"))?;
if response.success {
Ok(())
} else {
Err(Error::other(
response
.error
.map(|error| error.error_info)
.unwrap_or_else(|| "peer rejected scanner publication lease release".to_string()),
))
}
}
.await,
)
.await
}
pub async fn get_metacache_listing(&self) -> Result<()> {
warn!("get_metacache_listing is not implemented in PeerRestClient");
Err(Error::NotImplemented)
}
pub async fn update_metacache_listing(&self) -> Result<()> {
warn!("update_metacache_listing is not implemented in PeerRestClient");
Err(Error::NotImplemented)
}
pub async fn reload_pool_meta(&self) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut request = Request::new(ReloadPoolMetaRequest {});
set_tonic_mutation_body_digest(&mut request)?;
let response = client.reload_pool_meta(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure("reload_pool_meta", None, response.error_code, response.error_info));
}
Ok(())
}
.await,
)
.await
}
pub async fn stop_rebalance(&self, expected_rebalance_id: Option<&str>) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut request = Request::new(StopRebalanceRequest {
expected_rebalance_id: expected_rebalance_id.unwrap_or_default().to_string(),
});
set_tonic_mutation_body_digest(&mut request)?;
let response = client.stop_rebalance(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure("stop_rebalance", None, response.error_code, response.error_info));
}
Ok(())
}
.await,
)
.await
}
pub async fn load_rebalance_meta(&self, start_rebalance: bool) -> Result<()> {
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut request = Request::new(LoadRebalanceMetaRequest { start_rebalance });
set_tonic_mutation_body_digest(&mut request)?;
let response = client.load_rebalance_meta(request).await?.into_inner();
debug!(
event = "peer_rebalance_meta",
component = "ecstore",
subsystem = "peer_rest_client",
action = "load_rebalance_meta",
result = "response_received",
peer = %self.grid_host,
success = response.success,
start_rebalance = start_rebalance,
"peer rebalance metadata response"
);
if !response.success {
return Err(control_plane_failure(
"load_rebalance_meta",
None,
response.error_code,
response.error_info,
));
}
Ok(())
}
.await,
)
.await
}
pub async fn start_decommission(&self, pool_indices: Vec<usize>) -> Result<()> {
self.finalize_result(
async {
let pool_indices = pool_indices
.into_iter()
.map(|idx| {
u32::try_from(idx).map_err(|_| Error::other(format!("decommission pool index {idx} exceeds RPC range")))
})
.collect::<Result<Vec<_>>>()?;
let mut client = self.get_client().await?;
let mut request = Request::new(StartDecommissionRequest { pool_indices });
set_tonic_mutation_body_digest(&mut request)?;
let response = client.start_decommission(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure(
"start_decommission",
None,
response.error_code,
response.error_info,
));
}
Ok(())
}
.await,
)
.await
}
pub async fn decommission_cancel(&self, pool_index: usize) -> Result<()> {
self.finalize_result(
async {
let pool_index = u32::try_from(pool_index)
.map_err(|_| Error::other(format!("decommission pool index {pool_index} exceeds RPC range")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(CancelDecommissionRequest { pool_index });
set_tonic_mutation_body_digest(&mut request)?;
let response = client.cancel_decommission(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure(
"decommission_cancel",
None,
response.error_code,
response.error_info,
));
}
Ok(())
}
.await,
)
.await
}
pub async fn clear_decommission(&self, pool_index: usize) -> Result<()> {
self.finalize_result(
async {
let pool_index = u32::try_from(pool_index)
.map_err(|_| Error::other(format!("decommission pool index {pool_index} exceeds RPC range")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(ClearDecommissionRequest { pool_index });
set_tonic_mutation_body_digest(&mut request)?;
let response = client.clear_decommission(request).await?.into_inner();
if !response.success {
return Err(control_plane_failure(
"clear_decommission",
None,
response.error_code,
response.error_info,
));
}
Ok(())
}
.await,
)
.await
}
pub async fn load_transition_tier_config(&self) -> Result<()> {
match self.load_transition_tier_config_outcome().await {
TierConfigReloadOutcome::Success => Ok(()),
// Only a reconnect-class failure says anything about the channel.
// `finalize_result` marks the peer offline and evicts its connection
// whenever the message looks network-like, and a peer that answered
// and rejected the apply can easily report one ("release RPC failed:
// transport error"). Routing those through here would gate a healthy,
// responding peer out of every unrelated RPC.
TierConfigReloadOutcome::TransientReconnect(err) => self.finalize_result(Err(err)).await,
TierConfigReloadOutcome::TransientRetrySameChannel(err) | TierConfigReloadOutcome::Terminal(err) => Err(err),
}
}
pub(crate) async fn load_transition_tier_config_outcome(&self) -> TierConfigReloadOutcome {
let outcome = self.load_transition_tier_config_single_attempt_outcome().await;
if outcome.is_transient() {
return self.load_transition_tier_config_once_outcome().await;
}
outcome
}
pub(crate) async fn load_transition_tier_config_single_attempt_outcome(&self) -> TierConfigReloadOutcome {
let outcome = self.load_transition_tier_config_once_outcome().await;
if outcome.requires_reconnect() {
self.prepare_retry().await;
}
outcome
}
pub(crate) async fn load_transition_tier_config_once_outcome(&self) -> TierConfigReloadOutcome {
let mut client = match self.get_client().await {
Ok(client) => client,
Err(err) => return tier_config_reload_connection_outcome(err),
};
let mut request = Request::new(LoadTransitionTierConfigRequest {});
if let Err(err) = set_tonic_mutation_body_digest(&mut request) {
return TierConfigReloadOutcome::Terminal(Error::other(err));
}
request.set_timeout(rustfs_protos::heal_control_execution_timeout());
let response = match client.load_transition_tier_config(request).await {
Ok(response) => response.into_inner(),
Err(status) => return tier_config_reload_status_outcome(status),
};
if !response.success {
return tier_config_reload_remote_failure(response.error_code, response.error_info);
}
TierConfigReloadOutcome::Success
}
}
pub(crate) enum TierConfigReloadOutcome {
Success,
TransientReconnect(Error),
TransientRetrySameChannel(Error),
Terminal(Error),
}
impl TierConfigReloadOutcome {
fn is_transient(&self) -> bool {
matches!(self, Self::TransientReconnect(_) | Self::TransientRetrySameChannel(_))
}
fn requires_reconnect(&self) -> bool {
matches!(self, Self::TransientReconnect(_))
}
}
fn tier_config_reload_connection_outcome(err: Error) -> TierConfigReloadOutcome {
if is_tier_config_reload_connection_failure(&err) {
TierConfigReloadOutcome::TransientReconnect(err)
} else {
TierConfigReloadOutcome::Terminal(err)
}
}
fn is_tier_config_reload_connection_failure(err: &Error) -> bool {
// A bare "unavailable" is only trusted inside the local dial failure from
// `get_client` (typed as RemoteClientUnavailable), never in application text.
if let Error::RemoteClientUnavailable(detail) = err
&& detail.to_ascii_lowercase().contains("unavailable")
{
return true;
}
message_has_network_needle(&err.to_string())
}
/// Classifies a reload the peer answered but refused to apply.
///
/// The peer replied, so the channel is healthy and only the remote apply
/// failed. Those failures are transient by nature: the reload reads the tier
/// mutation intents and takes the distributed tier-config lock, both of which
/// fail while any other node is restarting or while the lock quorum is briefly
/// disturbed. Retiring the worker on the first such rejection leaves that peer
/// pinned to the old configuration with nothing left to heal it, so it answers
/// `TierNotFound` for a tier the rest of the cluster already committed until a
/// second admin mutation happens to spawn a fresh worker.
///
/// Convergence is the whole point of this path, so a rejection is retried on
/// the same channel. The worker's exponential backoff caps the cost at one
/// reload every `TIER_CONFIG_RELOAD_RETRY_CAP`, and `Terminal` stays reachable
/// for transport and gRPC status failures, which is where a genuinely
/// unrecoverable peer surfaces.
fn tier_config_reload_remote_failure(error_code: Option<i32>, error_info: Option<String>) -> TierConfigReloadOutcome {
// Remote rejections are transient by design (see the doc comment above);
// the typed not-initialized code keeps the error typed for downstream
// classifiers instead of a bare string (backlog#1845).
if error_code == Some(rustfs_protos::proto_gen::node_service::ControlPlaneErrorCode::ControlPlaneErrorNotInitialized as i32) {
return TierConfigReloadOutcome::TransientRetrySameChannel(Error::RemoteNotInitialized);
}
TierConfigReloadOutcome::TransientRetrySameChannel(Error::other(error_info.unwrap_or_default()))
}
fn tier_config_reload_status_outcome(status: tonic::Status) -> TierConfigReloadOutcome {
use tonic::Code;
if matches!(status.code(), Code::Unavailable | Code::DeadlineExceeded) {
TierConfigReloadOutcome::TransientReconnect(status.into())
} else if status.code() == Code::Unknown && status.message().starts_with("Service was not ready:") {
TierConfigReloadOutcome::TransientRetrySameChannel(status.into())
} else if status.code() == Code::Unknown
&& is_tier_config_reload_connection_failure(&Error::other(status.message().to_string()))
{
// tonic reports a connection dropped mid-call as `Unknown` carrying the
// transport error text rather than as `Unavailable`, which is what a peer
// restarting under an active mutation produces. Reconnect and retry, so
// the restart does not permanently retire this peer's reload worker.
TierConfigReloadOutcome::TransientReconnect(status.into())
} else {
TierConfigReloadOutcome::Terminal(status.into())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::com::STORAGE_CLASS_SUB_SYS;
use crate::disk::error::DiskError;
use crate::disk::error_reduce::reduce_errs;
use crate::layout::{disks_layout::DisksLayout, endpoints::SetupType};
use rustfs_config::{ENV_KUBERNETES_SERVICE_HOST, ENV_LOCAL_ENDPOINT_HOST, ENV_STARTUP_TOPOLOGY_WAIT_MODE};
use serde_json::Value;
use serial_test::serial;
use std::io::{self, Write};
use std::sync::{Arc, Mutex};
use temp_env::async_with_vars;
use tracing_subscriber::{Registry, fmt::MakeWriter, layer::SubscriberExt};
#[test]
fn control_plane_failure_prefers_typed_not_initialized_code() {
use rustfs_protos::proto_gen::node_service::ControlPlaneErrorCode;
let code = Some(ControlPlaneErrorCode::ControlPlaneErrorNotInitialized as i32);
// Typed code wins even when the legacy string is present (dual-write).
let err = control_plane_failure("load_bucket_metadata", Some("b"), code, Some("errServerNotInitialized".to_string()));
assert!(matches!(err, Error::RemoteNotInitialized));
assert!(crate::error::is_err_not_initialized(&err), "typed variant must satisfy the predicate");
// Legacy peers: no code, string only — the substring fallback still classifies.
let err = control_plane_failure("load_bucket_metadata", Some("b"), None, Some("errServerNotInitialized".to_string()));
assert!(crate::error::is_err_not_initialized(&err), "legacy string form must keep classifying");
// No code, no string: detail-free per-op failure, not misread as not-initialized.
let err = control_plane_failure("load_bucket_metadata", Some("b"), None, None);
assert!(!crate::error::is_err_not_initialized(&err));
assert!(err.to_string().contains("load_bucket_metadata"));
// Unspecified code behaves like no code.
let err = control_plane_failure(
"load_bucket_metadata",
None,
Some(ControlPlaneErrorCode::ControlPlaneErrorUnspecified as i32),
Some("boom".to_string()),
);
assert!(!matches!(err, Error::RemoteNotInitialized));
assert_eq!(err.to_string(), "Io error: boom");
}
#[test]
fn control_plane_not_initialized_wire_value_is_pinned() {
// The discriminant is wire contract: old peers ignore it, but a renumber
// would silently flip classification on mixed-version clusters.
use rustfs_protos::proto_gen::node_service::ControlPlaneErrorCode;
assert_eq!(ControlPlaneErrorCode::ControlPlaneErrorUnspecified as i32, 0);
assert_eq!(ControlPlaneErrorCode::ControlPlaneErrorNotInitialized as i32, 1);
assert_eq!(ControlPlaneErrorCode::ControlPlaneErrorInvalidArgument as i32, 2);
}
#[test]
fn control_plane_failure_preserves_typed_invalid_argument_reason() {
use rustfs_protos::proto_gen::node_service::ControlPlaneErrorCode;
let reason = "durable unresolved-entry recovery requires pool metadata V2 or V3";
let err = control_plane_failure(
"start_decommission",
None,
Some(ControlPlaneErrorCode::ControlPlaneErrorInvalidArgument as i32),
Some(reason.to_string()),
);
assert!(
matches!(
err,
Error::InvalidArgument(ref scope, ref operation, ref actual_reason)
if scope == "control-plane" && operation == "start_decommission" && actual_reason == reason
),
"forwarded validation failures must remain typed and actionable"
);
}
#[test]
fn tier_config_reload_remote_failure_keeps_typed_not_initialized() {
use rustfs_protos::proto_gen::node_service::ControlPlaneErrorCode;
let outcome = tier_config_reload_remote_failure(
Some(ControlPlaneErrorCode::ControlPlaneErrorNotInitialized as i32),
Some("errServerNotInitialized".to_string()),
);
match outcome {
TierConfigReloadOutcome::TransientRetrySameChannel(err) => {
assert!(matches!(err, Error::RemoteNotInitialized));
}
TierConfigReloadOutcome::TransientReconnect(err) | TierConfigReloadOutcome::Terminal(err) => {
panic!("not-initialized must stay retry-same-channel, got {err}")
}
TierConfigReloadOutcome::Success => panic!("a rejection cannot classify as success"),
}
}
#[test]
fn scanner_publication_lease_response_rejects_stale_generation_and_session() {
let token = Uuid::new_v4();
let response = ScannerPublicationLeaseResponse {
success: true,
token: token.as_bytes().to_vec().into(),
movement_generation: 7,
lease_ttl_ms: crate::store::SCANNER_PUBLICATION_LEASE_TTL_MS,
error: None,
response_proof: Bytes::new(),
owner_id: Uuid::new_v4().to_string(),
session_id: "session-a".to_string(),
};
assert!(validate_scanner_publication_lease_response_fields(&response, "session-a", 7).is_ok());
let stale_generation = ScannerPublicationLeaseResponse {
movement_generation: 6,
..response.clone()
};
let error = validate_scanner_publication_lease_response_fields(&stale_generation, "session-a", 7)
.expect_err("a response from an older movement generation must be rejected");
assert!(error.to_string().contains("different scanner publication lease generation"));
let stale_session = ScannerPublicationLeaseResponse {
session_id: "session-b".to_string(),
..response
};
let error = validate_scanner_publication_lease_response_fields(&stale_session, "session-a", 7)
.expect_err("a response from an older scanner session must be rejected");
assert!(error.to_string().contains("different scanner publication lease session"));
}
#[test]
fn scanner_publication_lease_deadline_accounts_for_delayed_rpc_response() {
let started = std::time::Instant::now();
let expected_deadline = started + Duration::from_secs(55);
let deadline = scanner_publication_lease_deadline(started, started + Duration::from_secs(10), 60_000)
.expect("a response inside the safety window should retain the original deadline");
assert_eq!(deadline, expected_deadline);
let error = scanner_publication_lease_deadline(started, started + Duration::from_secs(55), 60_000)
.expect_err("a response arriving at the safety boundary must fail closed");
assert!(error.to_string().contains("after its safety window"));
}
#[test]
fn replication_stats_response_decodes_valid_empty_provider() {
let mut stats = BucketStats::default();
stats.replication_stats.provider_available = true;
let payload = rmp_serde::to_vec_named(&stats).expect("bucket statistics should encode");
let decoded = decode_bucket_stats_response(GetBucketStatsDataResponse {
success: true,
bucket_stats: payload.into(),
error_info: None,
})
.expect("valid bucket statistics should decode");
assert!(decoded.replication_stats.provider_available);
assert!(decoded.replication_stats.stats.is_empty());
}
#[test]
fn replication_stats_response_rejects_unavailable_malformed_and_oversized_payloads() {
let unavailable = decode_bucket_stats_response(GetBucketStatsDataResponse {
success: false,
bucket_stats: Bytes::new(),
error_info: Some("provider unavailable".to_string()),
})
.expect_err("unavailable provider must not become a zero snapshot");
assert!(unavailable.to_string().contains("provider unavailable"));
let malformed = decode_bucket_stats_response(GetBucketStatsDataResponse {
success: true,
bucket_stats: Bytes::from_static(b"not-msgpack"),
error_info: None,
})
.expect_err("malformed peer statistics must fail closed");
assert!(!malformed.to_string().is_empty());
let oversized = decode_bucket_stats_response(GetBucketStatsDataResponse {
success: true,
bucket_stats: Bytes::from(vec![0; REPLICATION_STATS_MAX_MESSAGE_SIZE + 1]),
error_info: None,
})
.expect_err("oversized peer statistics must fail closed");
assert!(oversized.to_string().contains("size limit"));
}
#[derive(Clone, Default)]
struct CapturedLogs {
buffer: Arc<Mutex<Vec<u8>>>,
}
struct CapturedLogWriter {
buffer: Arc<Mutex<Vec<u8>>>,
}
impl CapturedLogs {
fn lines(&self) -> Vec<Value> {
let buffer = self
.buffer
.lock()
.expect("captured logs mutex should not be poisoned")
.clone();
String::from_utf8(buffer)
.expect("captured logs should be valid UTF-8")
.lines()
.map(|line| serde_json::from_str::<Value>(line).expect("captured log line should be valid JSON"))
.collect()
}
}
impl Write for CapturedLogWriter {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
self.buffer
.lock()
.expect("captured logs mutex should not be poisoned")
.extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
impl<'a> MakeWriter<'a> for CapturedLogs {
type Writer = CapturedLogWriter;
fn make_writer(&'a self) -> Self::Writer {
CapturedLogWriter {
buffer: Arc::clone(&self.buffer),
}
}
}
fn test_peer_client() -> PeerRestClient {
PeerRestClient::new(
XHost {
name: "127.0.0.1".to_string(),
port: 9000,
is_port_set: true,
},
"http://127.0.0.1:9000".to_string(),
)
}
fn decode_test_scanner_activity(response: ScannerActivityResponse) -> Result<ScannerPeerActivity> {
decode_scanner_activity_with_verifier(response, &[9; 16], |_canonical, proof| {
(proof == b"proof")
.then_some(())
.ok_or_else(|| Error::other("peer returned an invalid scanner activity response proof"))
})
}
fn decode_test_scanner_dirty_usage_snapshot(
response: ScannerDirtyUsageSnapshotResponse,
) -> Result<ScannerPeerDirtyUsageSnapshot> {
decode_scanner_dirty_usage_snapshot_with_verifier(response, &[9; 16], |_canonical, proof| {
(proof == b"proof")
.then_some(())
.ok_or_else(|| Error::other("peer returned an invalid scanner dirty usage snapshot response proof"))
})
}
fn test_scanner_dirty_usage_snapshot_response() -> ScannerDirtyUsageSnapshotResponse {
ScannerDirtyUsageSnapshotResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
generation: 7,
pending_bucket_count: 2,
protocol_version: SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION,
complete: true,
buckets: vec![
rustfs_protos::proto_gen::node_service::ScannerDirtyUsageBucket {
bucket: "archive".to_string(),
generation: 3,
},
rustfs_protos::proto_gen::node_service::ScannerDirtyUsageBucket {
bucket: "photos".to_string(),
generation: 7,
},
],
response_proof: b"proof".to_vec().into(),
}
}
#[test]
fn scanner_dirty_usage_snapshot_requires_a_complete_authenticated_ordered_view() {
let decoded = decode_test_scanner_dirty_usage_snapshot(test_scanner_dirty_usage_snapshot_response())
.expect("a complete authenticated dirty usage snapshot should decode");
assert_eq!(decoded.instance_id, "0123456789abcdef0123456789abcdef");
assert_eq!(decoded.generation, 7);
assert_eq!(decoded.pending_bucket_count, 2);
assert_eq!(decoded.protocol_version, SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION);
assert!(decoded.complete);
assert_eq!(decoded.buckets.get("archive"), Some(&3));
assert_eq!(decoded.buckets.get("photos"), Some(&7));
let overflow_count =
u64::try_from(SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES + 1).expect("the test snapshot entry limit should fit in u64");
let overflow = decode_test_scanner_dirty_usage_snapshot(ScannerDirtyUsageSnapshotResponse {
pending_bucket_count: overflow_count,
complete: false,
buckets: Vec::new(),
..test_scanner_dirty_usage_snapshot_response()
})
.expect("an explicit all-or-nothing overflow snapshot should decode");
assert!(!overflow.complete);
assert!(overflow.buckets.is_empty());
}
#[test]
fn scanner_dirty_usage_snapshot_rejects_inconsistent_or_partial_peer_data() {
let mut cases = Vec::new();
let mut invalid_instance = test_scanner_dirty_usage_snapshot_response();
invalid_instance.instance_id = "ABCDEF0123456789ABCDEF0123456789".to_string();
cases.push((invalid_instance, "instance ID"));
let mut unsupported = test_scanner_dirty_usage_snapshot_response();
unsupported.protocol_version = SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION + 1;
cases.push((unsupported, "unsupported"));
let mut exhausted = test_scanner_dirty_usage_snapshot_response();
exhausted.generation = u64::MAX;
cases.push((exhausted, "exhausted"));
let mut inconsistent_count = test_scanner_dirty_usage_snapshot_response();
inconsistent_count.pending_bucket_count = 3;
cases.push((inconsistent_count, "bucket count"));
let mut unordered = test_scanner_dirty_usage_snapshot_response();
unordered.buckets.reverse();
cases.push((unordered, "strictly ordered"));
let mut future_bucket = test_scanner_dirty_usage_snapshot_response();
future_bucket.buckets[0].generation = 8;
cases.push((future_bucket, "bucket generation"));
let mut zero_generation = test_scanner_dirty_usage_snapshot_response();
zero_generation.buckets[0].generation = 0;
cases.push((zero_generation, "bucket generation"));
let mut empty_bucket = test_scanner_dirty_usage_snapshot_response();
empty_bucket.buckets[0].bucket.clear();
cases.push((empty_bucket, "empty bucket name"));
let mut partial = test_scanner_dirty_usage_snapshot_response();
partial.complete = false;
cases.push((partial, "entry-limit overflow"));
let too_many_buckets = ScannerDirtyUsageSnapshotResponse {
generation: 1,
pending_bucket_count: u64::try_from(SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES + 1)
.expect("the test snapshot entry limit should fit in u64"),
buckets: (0..=SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES)
.map(|index| rustfs_protos::proto_gen::node_service::ScannerDirtyUsageBucket {
bucket: format!("bucket-{index:04}"),
generation: 1,
})
.collect(),
..test_scanner_dirty_usage_snapshot_response()
};
cases.push((too_many_buckets, "exceeds the entry limit"));
let overflow_count =
u64::try_from(SCANNER_DIRTY_USAGE_SNAPSHOT_MAX_ENTRIES + 1).expect("the test snapshot entry limit should fit in u64");
let invalid_overflow = ScannerDirtyUsageSnapshotResponse {
generation: 0,
pending_bucket_count: overflow_count,
complete: false,
buckets: Vec::new(),
..test_scanner_dirty_usage_snapshot_response()
};
cases.push((invalid_overflow, "without a generation"));
for (response, expected) in cases {
let err =
decode_test_scanner_dirty_usage_snapshot(response).expect_err("malformed dirty usage snapshots must fail closed");
assert!(err.to_string().contains(expected), "expected {expected:?} in {err}");
}
let mut invalid_proof = test_scanner_dirty_usage_snapshot_response();
invalid_proof.protocol_version = SCANNER_DIRTY_USAGE_SNAPSHOT_PROTOCOL_VERSION + 1;
invalid_proof.response_proof = b"invalid".to_vec().into();
let err = decode_test_scanner_dirty_usage_snapshot(invalid_proof)
.expect_err("an invalid response proof must fail before peer fields are trusted");
assert!(err.to_string().contains("response proof"));
}
#[test]
fn build_clients_from_slots_preserves_missing_remote_topology_slots() {
let slots = vec![
("127.0.0.1:9000".to_string(), None, true),
(
"rustfs-1.invalid:9001".to_string(),
Some("http://rustfs-1.invalid:9001".to_string()),
false,
),
("rustfs-2.invalid".to_string(), Some("http://rustfs-2.invalid".to_string()), false),
("127.0.0.1:notaport".to_string(), Some("http://127.0.0.1:notaport".to_string()), false),
("127.0.0.1:9003".to_string(), None, false),
];
let (remote, all, remote_topology_hosts) = PeerRestClient::build_clients_from_slots(slots);
assert_eq!(remote.len(), 4, "local node is excluded but remote slots are not compacted away");
assert_eq!(all.len(), 5, "all slots preserve the sorted cluster topology shape");
assert_eq!(
remote_topology_hosts,
vec![
"rustfs-1.invalid:9001".to_string(),
"rustfs-2.invalid".to_string(),
"127.0.0.1:notaport".to_string(),
"127.0.0.1:9003".to_string()
]
);
let unresolved = remote[0]
.as_ref()
.expect("temporarily unresolved remote peer should retain a client");
assert_eq!(unresolved.host.to_string(), "rustfs-1.invalid:9001");
let default_port = remote[1]
.as_ref()
.expect("temporarily unresolved scheme-default remote peer should retain a client");
assert_eq!(default_port.host.to_string(), "rustfs-2.invalid");
assert_eq!(default_port.host.port, 80);
assert!(!default_port.host.is_port_set);
assert!(remote[2].is_none(), "unparseable remote peer should remain observable as a missing slot");
assert!(remote[3].is_none(), "missing grid host should remain observable as a missing slot");
assert!(all[0].is_none(), "local node is represented by the local server_info row");
assert!(all[1].is_some());
assert!(all[2].is_some());
assert!(all[3].is_none());
assert!(all[4].is_none());
}
#[test]
fn topology_host_parser_preserves_names_and_bracketed_ipv6() {
let domain = PeerRestClient::parse_topology_host("rustfs-1.invalid", "https://rustfs-1.invalid")
.expect("unresolved HTTPS topology host should parse without DNS");
assert_eq!(domain.to_string(), "rustfs-1.invalid");
assert_eq!(domain.port, 443);
assert!(!domain.is_port_set);
let ipv6 = PeerRestClient::parse_topology_host("[2001:db8::1]:9000", "http://[2001:db8::1]:9000")
.expect("bracketed IPv6 topology host should parse without changing its identity");
assert_eq!(ipv6.to_string(), "[2001:db8::1]:9000");
let default_port_ipv6 = PeerRestClient::parse_topology_host("[2001:db8::2]", "http://[2001:db8::2]")
.expect("scheme-default IPv6 topology host should parse without DNS");
assert_eq!(default_port_ipv6.to_string(), "[2001:db8::2]");
assert_eq!(default_port_ipv6.port, 80);
assert!(!default_port_ipv6.is_port_set);
assert!(PeerRestClient::parse_topology_host("peer.invalid:0", "http://peer.invalid:0").is_err());
assert!(PeerRestClient::parse_topology_host("peer-a.invalid:9000", "http://peer-b.invalid:9000").is_err());
assert!(PeerRestClient::parse_topology_host("peer.invalid:9000", "http://peer.invalid:9000/unexpected").is_err());
}
#[tokio::test]
#[serial]
async fn unresolved_default_port_endpoint_topology_retains_all_peer_clients() {
let volumes = (0..4)
.map(|index| format!("http://rustfs-{index}.invalid:80/data{index}"))
.collect::<Vec<_>>();
let layout = DisksLayout::from_volumes(&volumes).expect("distributed default-port topology should parse");
async_with_vars(
[
(ENV_STARTUP_TOPOLOGY_WAIT_MODE, Some("orchestrated")),
(ENV_LOCAL_ENDPOINT_HOST, Some("rustfs-0.invalid")),
(ENV_KUBERNETES_SERVICE_HOST, None),
],
async {
let (server_pools, setup_type) = EndpointServerPools::create_server_endpoints("0.0.0.0:80", &layout)
.await
.expect("explicit local identity should avoid peer DNS during endpoint construction");
assert_eq!(setup_type, SetupType::DistErasure);
let (remote, all, remote_topology_hosts) =
PeerRestClient::build_clients_from_slots(server_pools.peer_grid_host_slots_sorted());
assert_eq!(remote.len(), 3);
assert!(
remote.iter().all(Option::is_some),
"unresolved remote peers must retain reconnectable clients"
);
assert_eq!(all.len(), 4);
assert_eq!(all.iter().filter(|client| client.is_none()).count(), 1);
assert_eq!(remote_topology_hosts.len(), 3);
assert!(
remote_topology_hosts.iter().all(|host| !host.contains(':')),
"scheme-default ports must preserve the legacy topology identity"
);
assert!(
remote
.iter()
.flatten()
.all(|client| client.host.port == 80 && !client.host.is_port_set),
"scheme-default peers must retain the effective dial port"
);
},
)
.await;
}
#[test]
fn scanner_activity_requires_restart_safe_peer_identity() {
let legacy = decode_test_scanner_activity(ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION,
topology_digest: Vec::new().into(),
data_movement_active: false,
response_proof: Vec::new().into(),
dirty_usage_generation: 0,
dirty_usage_pending: false,
movement_generation: None,
publication_blocked: None,
})
.expect("legacy peers should retain their activity generations during a rolling upgrade");
assert_eq!(
legacy,
ScannerPeerActivity {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION,
topology_digest: None,
data_movement_active: None,
dirty_usage_generation: None,
dirty_usage_pending: None,
movement_generation: None,
publication_blocked: None,
}
);
let previous = decode_test_scanner_activity(ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: true,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 0,
dirty_usage_pending: false,
movement_generation: None,
publication_blocked: None,
})
.expect("protocol v4 peers should remain observable during a rolling upgrade");
assert_eq!(
previous,
ScannerPeerActivity {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION,
topology_digest: Some([7; 32]),
data_movement_active: Some(true),
dirty_usage_generation: None,
dirty_usage_pending: None,
movement_generation: None,
publication_blocked: None,
}
);
let v6 = decode_test_scanner_activity(ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_V6_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: true,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
movement_generation: None,
publication_blocked: None,
})
.expect("v6 peers should remain readable without a v7 publication proof");
assert_eq!(v6.movement_generation, None);
assert_eq!(v6.publication_blocked, None);
assert_eq!(v6.dirty_usage_generation, Some(11));
let malformed_topology = ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: vec![7; 31].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
movement_generation: Some(19),
publication_blocked: Some(false),
};
assert!(
decode_test_scanner_activity(malformed_topology)
.expect_err("activity topology digests must have the protocol-defined length")
.to_string()
.contains("topology digest")
);
let missing_instance = ScannerActivityResponse {
instance_id: String::new(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
movement_generation: Some(19),
publication_blocked: Some(false),
};
assert!(
decode_test_scanner_activity(missing_instance)
.expect_err("an empty instance ID is not restart safe")
.to_string()
.contains("instance ID")
);
let malformed_instance = ScannerActivityResponse {
instance_id: "ABCDEF0123456789ABCDEF0123456789".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
movement_generation: Some(19),
publication_blocked: Some(false),
};
assert!(
decode_test_scanner_activity(malformed_instance)
.expect_err("activity instance IDs must use the canonical lowercase hex form")
.to_string()
.contains("instance ID")
);
let activity = decode_test_scanner_activity(ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: true,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
movement_generation: Some(19),
publication_blocked: Some(false),
})
.expect("complete activity responses should be accepted");
assert_eq!(
activity,
ScannerPeerActivity {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: Some([7; 32]),
data_movement_active: Some(true),
dirty_usage_generation: Some(11),
dirty_usage_pending: Some(true),
movement_generation: Some(19),
publication_blocked: Some(false),
}
);
let missing_movement_generation = ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 11,
dirty_usage_pending: false,
movement_generation: None,
publication_blocked: Some(false),
};
assert!(
decode_test_scanner_activity(missing_movement_generation)
.expect_err("v7 activity must carry movement generation")
.to_string()
.contains("movement generation")
);
let pending_without_generation = ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 0,
dirty_usage_pending: true,
movement_generation: Some(19),
publication_blocked: Some(false),
};
assert!(
decode_test_scanner_activity(pending_without_generation)
.expect_err("pending dirty usage must carry a nonzero generation")
.to_string()
.contains("without a generation")
);
let previous_with_dirty_usage = ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
movement_generation: None,
publication_blocked: None,
};
assert!(
decode_test_scanner_activity(previous_with_dirty_usage)
.expect_err("protocol v4 responses must not claim unauthenticated dirty usage fields")
.to_string()
.contains("unauthenticated v5 fields")
);
let legacy_with_topology = ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 0,
dirty_usage_pending: false,
movement_generation: None,
publication_blocked: None,
};
assert!(
decode_test_scanner_activity(legacy_with_topology)
.expect_err("legacy protocol responses must not claim extended fields")
.to_string()
.contains("unexpected extended fields")
);
let unsupported_protocol = ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION + 1,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: b"proof".to_vec().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
movement_generation: None,
publication_blocked: None,
};
assert!(
decode_test_scanner_activity(unsupported_protocol)
.expect_err("unknown activity protocols must fail closed")
.to_string()
.contains("unsupported scanner activity protocol")
);
let missing_proof = ScannerActivityResponse {
instance_id: "0123456789abcdef0123456789abcdef".to_string(),
namespace_generation: 7,
maintenance_generation: 3,
protocol_version: SCANNER_ACTIVITY_PROTOCOL_VERSION,
topology_digest: vec![7; 32].into(),
data_movement_active: false,
response_proof: Vec::new().into(),
dirty_usage_generation: 11,
dirty_usage_pending: true,
movement_generation: Some(19),
publication_blocked: Some(false),
};
assert!(
decode_test_scanner_activity(missing_proof)
.expect_err("unsigned scanner activity responses must fail closed")
.to_string()
.contains("response proof")
);
}
#[test]
fn dynamic_scanner_config_requires_versioned_peer_acknowledgement() {
for sub_system in [SCANNER_SUB_SYS, HEAL_SUB_SYS] {
let err = validate_signal_service_protocol(SERVICE_SIGNAL_RELOAD_DYNAMIC, sub_system, 0)
.expect_err("an unversioned peer must not claim scanner config convergence");
assert!(err.to_string().contains("does not support dynamic"));
validate_signal_service_protocol(
SERVICE_SIGNAL_RELOAD_DYNAMIC,
sub_system,
rustfs_protos::DYNAMIC_CONFIG_PROTOCOL_VERSION,
)
.expect("a current peer should support dynamic scanner config");
}
validate_signal_service_protocol(SERVICE_SIGNAL_RELOAD_DYNAMIC, STORAGE_CLASS_SUB_SYS, 0)
.expect("unrelated dynamic config keeps its existing compatibility contract");
validate_signal_service_protocol(SERVICE_SIGNAL_REFRESH_CONFIG, SCANNER_SUB_SYS, 0)
.expect("full refresh compatibility is guarded by its scanner preflight");
}
#[test]
fn dynamic_kms_config_requires_versioned_peer_acknowledgement() {
let err = validate_signal_service_protocol(SERVICE_SIGNAL_RELOAD_DYNAMIC, KMS_SIGNAL_SUBSYSTEM, 0)
.expect_err("an unversioned peer must not claim KMS config convergence");
assert!(err.to_string().contains("does not support dynamic"));
validate_signal_service_protocol(
SERVICE_SIGNAL_RELOAD_DYNAMIC,
KMS_SIGNAL_SUBSYSTEM,
rustfs_protos::DYNAMIC_CONFIG_PROTOCOL_VERSION,
)
.expect("a current peer should support dynamic KMS config");
}
#[test]
fn peer_rest_client_marks_network_like_errors() {
assert!(PeerRestClient::is_network_like_error(&Error::other("transport error")));
assert!(PeerRestClient::is_network_like_error(&Error::other("connection refused")));
assert!(!PeerRestClient::is_network_like_error(&Error::NotImplemented));
}
#[test]
fn peer_rest_client_network_classifier_uses_typed_status_code() {
// The one code that means "nothing is answering on this channel".
assert!(PeerRestClient::is_network_like_error(&Error::from(tonic::Status::unavailable(
"storage layer is not initialized"
))));
// Application statuses from a live peer must not mark it offline,
// even when their message contains transport-sounding words.
assert!(!PeerRestClient::is_network_like_error(&Error::from(tonic::Status::internal(
"failed to reload metadata for bucket \"unavailable-logs\""
))));
assert!(!PeerRestClient::is_network_like_error(&Error::from(tonic::Status::unauthenticated(
"No valid auth token"
))));
// A request-budget expiry answered by a live peer is an application
// outcome, not a transport failure.
assert!(!PeerRestClient::is_network_like_error(&Error::from(tonic::Status::deadline_exceeded(
"heal control request expired"
))));
// Unknown is the transport's escape hatch for a cause it could not
// map, and our handlers never return it, so there the text decides.
assert!(PeerRestClient::is_network_like_error(&Error::from(tonic::Status::unknown(
"Service was not ready: transport error"
))));
assert!(!PeerRestClient::is_network_like_error(&Error::from(tonic::Status::unknown(
"peer response unknown"
))));
}
#[test]
fn peer_rest_client_network_classifier_ignores_transport_words_in_application_statuses() {
// The reason classification reads the code rather than the text: a
// peer that answers is reachable, even when what it says describes a
// connection failure of its own. A handler interpolating a local
// io::Error into Status::internal, or relaying trouble with its own
// downstream, must not cost us the channel to a healthy peer.
for status in [
tonic::Status::internal("connection refused while dialing downstream backend"),
tonic::Status::internal("write failed: broken pipe"),
tonic::Status::unauthenticated("connection reset while validating token"),
tonic::Status::failed_precondition("scanner lease timed out"),
tonic::Status::deadline_exceeded("heal control request timed out"),
] {
let rendered = status.to_string();
assert!(
!PeerRestClient::is_network_like_error(&Error::from(status)),
"an answered application status must not mark the peer offline: {rendered}"
);
}
}
#[test]
fn peer_rest_client_network_classifier_keeps_slow_peers_online() {
// The per-RPC channel deadline (RUSTFS_INTERNODE_RPC_TIMEOUT, 30s)
// surfaces as Cancelled "Timeout expired" carrying the transport
// cause as its source. A peer that is merely slow must stay online:
// gating it would spend a full recovery cycle fast-failing every RPC
// to a host that is still answering, turning load into a partition.
let timeout_status = tonic::Status::cancelled("Timeout expired");
assert!(!PeerRestClient::is_network_like_error(&Error::from(timeout_status)));
let sourced = tonic::Status::from_error(Box::new(std::io::Error::other("Timeout expired")));
assert!(
std::error::Error::source(&sourced).is_some(),
"the transport builds this status through Status::from_error, which attaches the cause"
);
assert!(!PeerRestClient::is_network_like_error(&Error::from(sourced)));
}
#[test]
fn rpc_status_errors_keep_their_rendering_and_hide_peer_metadata() {
let err = Error::from(tonic::Status::unavailable("peer gone"));
assert_eq!(
err.to_string(),
"Io error: code: 'The service is currently unavailable', message: \"peer gone\""
);
// tonic's own Debug prints the MetadataMap, i.e. every response header
// the peer sent; those must not reach a log through this error.
let mut status = tonic::Status::unavailable("peer gone");
status
.metadata_mut()
.insert("authorization", "Bearer secret".parse().expect("valid header value"));
let rendered = format!("{:?}", Error::from(status));
assert!(!rendered.contains("Bearer secret"), "{rendered}");
assert!(!rendered.contains("MetadataMap"), "{rendered}");
}
#[test]
fn peer_rest_client_network_classifier_ignores_application_text_containing_unavailable() {
// Regression: a bare "unavailable" needle used to match application
// strings like these and take a healthy peer offline.
assert!(!PeerRestClient::is_network_like_error(&Error::other(
"peer replication statistics provider is unavailable"
)));
assert!(!PeerRestClient::is_network_like_error(&Error::other(
"bucket \"unavailable-logs\" not found"
)));
// Anchored renderings of a flattened Unavailable status still match:
// tonic >= 0.14 form ...
assert!(PeerRestClient::is_network_like_error(&Error::other(
"peer tier mutation commit RPC failed: code: 'The service is currently unavailable', message: \"peer gone\""
)));
// ... which is only anchored as long as tonic renders Unavailable this
// way. A tonic bump that reworded it leaves the typed path correct but
// this needle stale, so pin the coupling rather than discover it in a
// partition.
assert!(
tonic::Status::unavailable("peer gone")
.to_string()
.to_ascii_lowercase()
.contains("code: 'the service is currently unavailable'"),
"tonic reworded Code::Unavailable; update the anchored needle"
);
// ... and the tonic <= 0.13 form peers may relay in error_info.
assert!(PeerRestClient::is_network_like_error(&Error::other(
"peer tier mutation commit RPC failed: status: Unavailable, message: \"peer gone\""
)));
}
#[test]
fn tier_config_reload_outcome_keeps_tonic_and_remote_errors_typed() {
assert!(matches!(
tier_config_reload_status_outcome(tonic::Status::unavailable("peer offline")),
TierConfigReloadOutcome::TransientReconnect(_)
));
assert!(matches!(
tier_config_reload_status_outcome(tonic::Status::deadline_exceeded("peer timeout")),
TierConfigReloadOutcome::TransientReconnect(_)
));
assert!(matches!(
tier_config_reload_status_outcome(tonic::Status::permission_denied("bad signature")),
TierConfigReloadOutcome::Terminal(_)
));
assert!(matches!(
tier_config_reload_status_outcome(tonic::Status::unknown("Service was not ready: test client")),
TierConfigReloadOutcome::TransientRetrySameChannel(_)
));
assert!(matches!(
tier_config_reload_status_outcome(tonic::Status::unknown("peer response unknown")),
TierConfigReloadOutcome::Terminal(_)
));
assert!(matches!(
tier_config_reload_status_outcome(tonic::Status::cancelled("request cancelled")),
TierConfigReloadOutcome::Terminal(_)
));
// A peer that answered and then refused the apply is retried rather than
// retired: the channel is healthy, so the rejection reflects remote state
// that the next attempt can find healed.
assert!(matches!(
tier_config_reload_remote_failure(None, Some("backend unavailable".to_string())),
TierConfigReloadOutcome::TransientRetrySameChannel(_)
));
assert!(matches!(
tier_config_reload_remote_failure(None, Some("errServerNotInitialized".to_string())),
TierConfigReloadOutcome::TransientRetrySameChannel(_)
));
assert!(matches!(
tier_config_reload_connection_outcome(Error::other("backend unavailable")),
TierConfigReloadOutcome::Terminal(_)
));
assert!(matches!(
tier_config_reload_connection_outcome(Error::RemoteClientUnavailable("connection unavailable".to_string())),
TierConfigReloadOutcome::TransientReconnect(_)
));
// The bare word is trusted only inside the typed local dial failure,
// not anywhere in application text — including text that mimics the
// old "can not get client" string form, which is retired.
assert!(matches!(
tier_config_reload_connection_outcome(Error::other(
"bucket unavailable-logs rejected it, then: can not get client, err: some other reason"
)),
TierConfigReloadOutcome::Terminal(_)
));
assert!(matches!(
tier_config_reload_connection_outcome(Error::other("can not get client, err: connection unavailable")),
TierConfigReloadOutcome::Terminal(_)
));
}
/// A tier mutation issued while another node restarts must still converge on
/// the nodes that stayed up. Those peers answer the reload RPC and reject the
/// apply, because reloading reads the tier mutation intents and takes the
/// distributed tier-config lock while the lock quorum is still disturbed.
/// Classifying those rejections as terminal retired the reload worker on its
/// first attempt and pinned the peer to the previous configuration, so it
/// served `TierNotFound` for an already-committed tier until an unrelated
/// second admin mutation spawned a new worker.
#[test]
fn tier_config_reload_retries_peers_that_reject_the_apply_mid_restart() {
for error_info in [
"Lock acquisition timeout for resource '.rustfs.sys/config/tier-config.bin.lock' after 5s",
"Resource '.rustfs.sys/config/tier-config.bin.lock' is already locked by node-3",
"Internal error: release RPC failed: transport error",
"save_config_with_opts: err: PreconditionFailed",
"erasure read quorum",
] {
assert!(
matches!(
tier_config_reload_remote_failure(None, Some(error_info.to_string())),
TierConfigReloadOutcome::TransientRetrySameChannel(_)
),
"a peer that rejected the apply must stay retryable so it converges: {error_info}"
);
}
// An absent error message is still a rejection, not a reason to stop.
assert!(matches!(
tier_config_reload_remote_failure(None, None),
TierConfigReloadOutcome::TransientRetrySameChannel(_)
));
// tonic surfaces a connection dropped mid-call as `Unknown`, not `Unavailable`.
assert!(matches!(
tier_config_reload_status_outcome(tonic::Status::unknown("transport error")),
TierConfigReloadOutcome::TransientReconnect(_)
));
// An `Unknown` that is not transport-shaped stays terminal.
assert!(matches!(
tier_config_reload_status_outcome(tonic::Status::unknown("peer response unknown")),
TierConfigReloadOutcome::Terminal(_)
));
}
#[tokio::test]
async fn tier_config_reload_single_attempt_clears_offline_gate_without_redial() {
let client = test_peer_client();
client.offline.store(true, Ordering::Release);
let outcome = client.load_transition_tier_config_single_attempt_outcome().await;
assert!(matches!(outcome, TierConfigReloadOutcome::TransientReconnect(_)));
assert!(!client.offline.load(Ordering::Acquire));
}
#[test]
fn tier_config_reload_readiness_retry_does_not_require_reconnect() {
let client = test_peer_client();
client.offline.store(true, Ordering::Release);
let outcome = tier_config_reload_status_outcome(tonic::Status::unknown("Service was not ready: startup"));
assert!(matches!(outcome, TierConfigReloadOutcome::TransientRetrySameChannel(_)));
assert!(!outcome.requires_reconnect());
assert!(client.offline.load(Ordering::Acquire));
}
#[tokio::test]
async fn peer_rest_client_fast_fails_when_marked_offline() {
let client = test_peer_client();
client.offline.store(true, Ordering::Release);
let err = client
.get_client()
.await
.expect_err("offline peer should fast-fail before dialing");
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));
}
#[test]
fn heal_control_capability_proof_must_authenticate_exact_ack() {
runtime_sources::ensure_test_rpc_secret();
let proof = crate::cluster::rpc::sign_tonic_rpc_response_proof(b"expected").expect("test proof should sign");
assert!(validate_heal_control_capability_proof(b"expected", &proof).is_ok());
let err = validate_heal_control_capability_proof(b"different", &proof)
.expect_err("a proof for a different acknowledgement must fail closed");
assert!(err.to_string().contains("invalid heal control capability proof"));
}
#[test]
fn heal_control_response_proof_binds_command_and_result() {
runtime_sources::ensure_test_rpc_secret();
let canonical = rustfs_protos::canonical_heal_control_response_body(2, "fingerprint", b"query", b"result")
.expect("small response should encode");
let proof = crate::cluster::rpc::sign_tonic_rpc_response_proof(&canonical).expect("test proof should sign");
assert!(validate_heal_control_response_proof(&canonical, &proof).is_ok());
for tampered in [
rustfs_protos::canonical_heal_control_response_body(2, "fingerprint", b"cancel", b"result").unwrap(),
rustfs_protos::canonical_heal_control_response_body(2, "fingerprint", b"query", b"tampered").unwrap(),
] {
let err = validate_heal_control_response_proof(&tampered, &proof)
.expect_err("proof must not authenticate a different command or result");
assert!(err.to_string().contains("invalid heal control response proof"));
}
}
#[test]
fn remote_version_state_capability_decoder_fails_closed() {
let epoch = Uuid::new_v4();
let result = rustfs_protos::encode_remote_version_state_capability("node-a:9000", epoch.as_bytes())
.expect("small capability response should encode");
assert_eq!(
decode_remote_version_state_capability("node-a:9000", &result).expect("valid epoch should decode"),
epoch
);
assert!(decode_remote_version_state_capability("node-b:9000", &result).is_err());
assert!(decode_remote_version_state_capability("node-a:9000", &result[..result.len() - 1]).is_err());
let nil = rustfs_protos::encode_remote_version_state_capability("node-a:9000", Uuid::nil().as_bytes())
.expect("small capability response should encode");
assert!(decode_remote_version_state_capability("node-a:9000", &nil).is_err());
}
#[test]
fn cross_pool_fence_capability_decoder_fails_closed() {
let epoch = Uuid::new_v4();
let result = rustfs_protos::encode_cross_pool_fence_capability(1, "node-a:9000", epoch.as_bytes())
.expect("small capability response should encode");
assert_eq!(
decode_cross_pool_fence_capability("node-a:9000", &result).expect("valid capability should decode"),
(1, epoch)
);
for malformed in [&[][..], &[0, 0, 0][..], &result[..result.len() - 1]] {
assert!(decode_cross_pool_fence_capability("node-a:9000", malformed).is_err());
}
assert!(decode_cross_pool_fence_capability("node-b:9000", &result).is_err());
let nil = rustfs_protos::encode_cross_pool_fence_capability(1, "node-a:9000", Uuid::nil().as_bytes())
.expect("small capability response should encode");
assert!(decode_cross_pool_fence_capability("node-a:9000", &nil).is_err());
}
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>,
failure_class: i32,
}
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,
failure_class: input.failure_class,
})
.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(),
failure_class: input.failure_class,
}
}
#[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,
failure_class: TierMutationFailureClass::Unspecified as i32,
});
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()
},
TierMutationControlResponse {
failure_class: TierMutationFailureClass::Ambiguous as i32,
..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_response_rejects_oversized_proof_and_error_before_verification() {
let mutation_id = Uuid::new_v4();
let payload = b"tier-mutation-prepare";
let oversized_proof = TierMutationControlResponse {
success: false,
state: TierMutationPeerState::Unspecified as i32,
applied: false,
error_info: None,
response_proof: vec![0; rustfs_protos::TIER_MUTATION_RPC_MAX_RESPONSE_PROOF_SIZE + 1].into(),
failure_class: TierMutationFailureClass::Ambiguous as i32,
};
let err = validate_tier_mutation_response_proof(
rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION,
TierMutationRpcPhase::Prepare,
mutation_id,
payload,
&oversized_proof,
)
.expect_err("oversized proof must fail before cryptographic verification");
assert!(err.to_string().contains("response proof exceeds size limit"));
let oversized_error = TierMutationControlResponse {
response_proof: Bytes::new(),
error_info: Some("e".repeat(rustfs_protos::TIER_MUTATION_RPC_MAX_ERROR_INFO_SIZE + 1)),
..oversized_proof
};
let err = validate_tier_mutation_response_proof(
rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION,
TierMutationRpcPhase::Prepare,
mutation_id,
payload,
&oversized_error,
)
.expect_err("oversized error detail must fail before proof construction");
assert!(err.to_string().contains("error response exceeds size limit"));
}
#[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()
);
assert!(validate_tier_mutation_payload_len(TierMutationRpcPhase::Abort, 0).is_err());
validate_tier_mutation_payload_len(TierMutationRpcPhase::Abort, 1).expect("non-empty abort payload should fit");
validate_tier_mutation_payload_len(TierMutationRpcPhase::Abort, rustfs_protos::TIER_MUTATION_RPC_MAX_ABORT_PAYLOAD_SIZE)
.expect("max abort payload should fit");
assert!(
validate_tier_mutation_payload_len(
TierMutationRpcPhase::Abort,
rustfs_protos::TIER_MUTATION_RPC_MAX_ABORT_PAYLOAD_SIZE + 1,
)
.is_err()
);
}
#[test]
fn tier_mutation_rpc_status_matrix_fails_closed_for_old_or_unresponsive_peers() {
for (phase, label) in [
(TierMutationRpcPhase::Prepare, "prepare"),
(TierMutationRpcPhase::Commit, "commit"),
(TierMutationRpcPhase::Abort, "abort"),
] {
for status in [
tonic::Status::unimplemented("old peer has no tier mutation control service"),
tonic::Status::deadline_exceeded("peer tier mutation control timed out"),
tonic::Status::unavailable("peer tier mutation control unavailable"),
] {
let err = tier_mutation_control_status_error(phase, rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, status);
let rendered = err.to_string();
assert!(rendered.contains(&format!("peer tier mutation {label} RPC failed")), "{rendered}");
assert!(
rendered.contains("old peer")
|| rendered.contains("timed out")
|| rendered.contains("unavailable")
|| rendered.contains("Unavailable"),
"{rendered}"
);
}
}
}
#[test]
fn tier_mutation_v4_to_v3_rejection_classification_requires_exact_status_and_message() {
let version = rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION;
let exact = format!("unsupported tier mutation peer protocol version: {version}");
let rejected = tier_mutation_control_status_error(
TierMutationRpcPhase::Prepare,
version,
tonic::Status::failed_precondition(exact.clone()),
);
assert!(tier_mutation_error_is_definitely_rejected(&rejected));
for status in [
tonic::Status::failed_precondition(format!("{exact}.")),
tonic::Status::failed_precondition(format!("unsupported tier mutation peer protocol version: {}", version - 1)),
tonic::Status::invalid_argument(exact.clone()),
tonic::Status::unimplemented(exact),
] {
let ambiguous = tier_mutation_control_status_error(TierMutationRpcPhase::Prepare, version, status);
assert!(
!tier_mutation_error_is_definitely_rejected(&ambiguous),
"near-text, wrong-code, and Unimplemented failures must remain ambiguous"
);
}
}
#[test]
fn tier_mutation_v4_failure_class_is_typed_and_fails_closed() {
let version = rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION;
let rejected = tier_mutation_failed_response_error(
version,
TierMutationFailureClass::PreDispatchRejected as i32,
Some("rejected".to_string()),
);
assert!(tier_mutation_error_is_definitely_rejected(&rejected));
for failure_class in [
TierMutationFailureClass::Unspecified as i32,
TierMutationFailureClass::Ambiguous as i32,
99,
] {
let ambiguous = tier_mutation_failed_response_error(version, failure_class, None);
assert!(
!tier_mutation_error_is_definitely_rejected(&ambiguous),
"missing, unknown, and explicit ambiguous classes must trigger Abort fanout"
);
}
let v3_ignores_v4_class = tier_mutation_failed_response_error(
rustfs_protos::TIER_MUTATION_RPC_PREVIOUS_PROTOCOL_VERSION,
TierMutationFailureClass::PreDispatchRejected as i32,
Some("legacy failure".to_string()),
);
assert!(!tier_mutation_error_is_definitely_rejected(&v3_ignores_v4_class));
}
#[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
// clearing it, an in-call retry would fast-fail on the gate instead of
// re-dialing (rustfs/backlog#1049 P1-B). prepare_retry must clear it.
let client = test_peer_client();
client.offline.store(true, Ordering::Release);
client.prepare_retry().await;
assert!(
!client.offline.load(Ordering::Acquire),
"prepare_retry must clear the offline gate so the next get_client re-dials"
);
}
#[tokio::test]
async fn peer_rest_client_finalize_result_marks_offline_for_network_errors() {
let client = test_peer_client();
let err = client
.finalize_result::<()>(Err(Error::other("transport error")))
.await
.expect_err("network error should still be returned");
assert!(err.to_string().contains("transport error"));
assert!(client.offline.load(Ordering::Acquire));
}
#[tokio::test]
async fn peer_rest_client_finalize_result_keeps_online_for_business_errors() {
let client = test_peer_client();
let err = client
.finalize_result::<()>(Err(Error::VolumeNotFound))
.await
.expect_err("business error should still be returned");
assert!(matches!(err, Error::VolumeNotFound));
assert!(!client.offline.load(Ordering::Acquire));
}
#[tokio::test(start_paused = true)]
async fn dropped_peer_client_releases_and_stops_its_recovery_monitor() {
let client = test_peer_client();
client.offline.store(true, Ordering::Release);
client.recovery_running.store(true, Ordering::Release);
let offline = Arc::downgrade(&client.offline);
let recovery_running = Arc::downgrade(&client.recovery_running);
let handle = PeerRestClient::spawn_recovery_monitor(client.grid_host.clone(), offline.clone(), recovery_running.clone());
let started = tokio::time::Instant::now();
drop(client);
assert!(offline.upgrade().is_none(), "detached recovery must not retain offline state");
assert!(
recovery_running.upgrade().is_none(),
"detached recovery must not retain its running state"
);
handle.await.expect("recovery monitor should not panic");
assert_eq!(
tokio::time::Instant::now(),
started,
"recovery monitor should stop before advancing to its first delayed probe"
);
}
#[tokio::test]
async fn peer_rest_client_finalize_result_keeps_online_for_app_errors_mentioning_unavailable() {
// Regression: application error text containing "unavailable" (a
// remote error_info payload, or a bucket named "unavailable-logs" in
// a typed application status) must not take a healthy peer offline.
let client = test_peer_client();
let err = client
.finalize_result::<()>(Err(Error::other("peer replication statistics provider is unavailable")))
.await
.expect_err("application error should still be returned");
assert!(err.to_string().contains("provider is unavailable"));
assert!(!client.offline.load(Ordering::Acquire));
let err = client
.finalize_result::<()>(Err(Error::from(tonic::Status::internal(
"failed to reload metadata for bucket \"unavailable-logs\"",
))))
.await
.expect_err("application status should still be returned");
assert!(err.to_string().contains("unavailable-logs"));
assert!(!client.offline.load(Ordering::Acquire));
}
#[tokio::test]
async fn peer_rest_client_finalize_result_marks_offline_for_typed_unavailable_status() {
let client = test_peer_client();
client
.finalize_result::<()>(Err(Error::from(tonic::Status::unavailable("storage layer is not initialized"))))
.await
.expect_err("network error should still be returned");
assert!(client.offline.load(Ordering::Acquire));
}
#[tokio::test(flavor = "current_thread")]
async fn peer_rest_recovery_probe_logs_keep_request_id_span_context() {
let logs = CapturedLogs::default();
let subscriber = Registry::default().with(
tracing_subscriber::fmt::layer()
.with_writer(logs.clone())
.with_ansi(false)
.without_time()
.json()
.flatten_event(true)
.with_current_span(true)
.with_span_list(true),
);
let _guard = tracing::subscriber::set_default(subscriber);
// The `recovery-monitor` callsite is shared with the production
// `mark_offline_and_spawn_recovery` path that sibling tests exercise from
// subscriber-less threads; without this the span can be cached as
// `Interest::never()` and silently degrade to `Span::none()`.
let _callsite_pin = crate::test_tracing::pin_callsite_interest_for_test();
let client = test_peer_client();
let span = tracing::info_span!("request-span", request_id = "req-peer-rest");
let _entered = span.enter();
let done = client.spawn_recovery_monitor_log_probe_for_test();
done.await.expect("recovery monitor probe should signal completion");
let log = logs
.lines()
.into_iter()
.find(|value| value.get("message").and_then(Value::as_str) == Some("peer recovery monitor log probe"))
.expect("expected peer recovery monitor probe log");
assert_eq!(log["span"]["name"], Value::String("recovery-monitor".to_string()));
assert_eq!(log["span"]["kind"], Value::String("peer_rest".to_string()));
let spans = log["spans"].as_array().expect("spans should be present");
assert!(spans.iter().any(|span| {
span.get("name").and_then(Value::as_str) == Some("request-span")
&& span.get("request_id").and_then(Value::as_str) == Some("req-peer-rest")
}));
}
/// Every operation name passed to `peer_failure_without_details` in this file.
const PEER_FAILURE_OPS: &[&str] = &[
"local_storage_info",
"server_info",
"get_cpus",
"get_net_info",
"get_partitions",
"get_os_info",
"get_se_linux_info",
"get_sys_config",
"get_sys_errors",
"get_mem_info",
"get_metrics",
"get_live_events",
"get_proc_info",
"start_profiling",
"load_bucket_metadata",
"delete_bucket_metadata",
"delete_policy",
"load_policy",
"load_policy_mapping",
"delete_user",
"delete_service_account",
"load_user",
"load_service_account",
"load_group",
"reload_site_replication_config",
"signal_service",
"reload_pool_meta",
"stop_rebalance",
"load_rebalance_meta",
"start_decommission",
"decommission_cancel",
"clear_decommission",
];
#[test]
fn peer_failure_without_details_names_operation_and_bucket() {
for op in PEER_FAILURE_OPS {
let message = peer_failure_without_details(op, None).to_string();
assert!(message.contains(op), "{op} message must name the operation: {message}");
}
for op in ["load_bucket_metadata", "delete_bucket_metadata"] {
let message = peer_failure_without_details(op, Some("ops-bucket")).to_string();
assert!(message.contains(op), "{op} message must name the operation: {message}");
assert!(message.contains("ops-bucket"), "{op} message must name the bucket: {message}");
}
}
#[test]
fn peer_failure_without_details_keeps_one_reduce_errs_bucket_per_operation() {
// reduce_errs groups Io errors by kind plus rendered message: peers failing the
// same operation must stay a single dominant error instead of one bucket per peer.
let per_peer_errs = (0..4)
.map(|_| {
Some(
peer_failure_without_details("load_bucket_metadata", Some("shared"))
.narrow_to_disk()
.unwrap_or_else(DiskError::other),
)
})
.collect::<Vec<_>>();
let (count, dominant) = reduce_errs(&per_peer_errs, &[]);
assert_eq!(count, 4, "one shared failure must not split into per-peer buckets");
assert_eq!(
dominant,
Some(
peer_failure_without_details("load_bucket_metadata", Some("shared"))
.narrow_to_disk()
.unwrap_or_else(DiskError::other)
)
);
assert_ne!(
peer_failure_without_details("load_bucket_metadata", Some("shared")).to_string(),
peer_failure_without_details("delete_bucket_metadata", Some("shared")).to_string()
);
assert_ne!(
peer_failure_without_details("load_bucket_metadata", Some("bucket-a")).to_string(),
peer_failure_without_details("load_bucket_metadata", Some("bucket-b")).to_string()
);
}
#[test]
fn peer_failure_without_details_never_reads_as_a_network_failure() {
// `finalize_result` marks the peer offline and evicts its connection whenever the
// message matches a network needle. A peer that answered `success = false` is alive,
// so no operation or bucket name may push this text over that classifier.
for op in PEER_FAILURE_OPS {
let err = peer_failure_without_details(op, None);
assert!(
!PeerRestClient::is_network_like_error(&err),
"{op} must not read as a transport failure: {err}"
);
let scoped = peer_failure_without_details(op, Some("bucket-name"));
assert!(
!PeerRestClient::is_network_like_error(&scoped),
"{op} must not read as a transport failure: {scoped}"
);
}
// The bucket name is caller-supplied. Every needle carries a space, which S3 bucket
// names cannot, and the name is closed by `)` before the literal text resumes, so no
// needle can straddle the boundary either.
for bucket in [
"timed-out",
"connection-reset",
"transport-error",
"broken-pipe",
"unavailable-logs",
] {
let err = peer_failure_without_details("load_bucket_metadata", Some(bucket));
assert!(
!PeerRestClient::is_network_like_error(&err),
"bucket {bucket} must not push the message over the network classifier: {err}"
);
}
}
}