feat(ecstore): type internode client-acquisition failures for quorum buckets (#6619)

* feat(ecstore): type internode client-acquisition failures for stable quorum buckets

Backlog#1845 step 3, first typed family. The largest other(format!) message family in ecstore was 'can not get client, err: {detail}' (~50 production sites): every internode RPC that fails to acquire a client wrapped the dial/auth error with per-peer detail into DiskError::other / StorageError::other, whose Io equality compares the rendered message. N disks failing for this same cause therefore counted as N distinct errors in reduce_errs, starving quorum aggregation, and remote_disk call sites double-wrapped the message on top of get_client's own wrap.

Introduce DiskError::RemoteClientUnavailable(String) (wire code 0x2B) and its StorageError twin (StorageErrorCode 0x54): equality and hashing use the wire code alone, so same-cause failures land in one quorum bucket regardless of per-peer detail, while Display keeps the detail so substring classifiers (network needles, heal recoverability) keep reading it unchanged. Wire encoding carries the rendered detail in error_info and decode restores the typed variant; old peers fall back to the legacy string form gracefully.

Call sites: remote_disk get_client/get_bulk_client/offline-bypass/recovery-probe now construct the typed variant and the ~60 redundant double-wrap map_errs are gone; peer_rest_client's three client getters and offline gates, peer_s3_client, and admin_server_info follow. The tier-config-reload connection classifier's anchored 'can not get client' substring check becomes a typed match on the variant (the string form is retired and now classifies as Terminal, pinned by test).

Ref rustfs/backlog#1845

* chore(ci): refresh error other ratchet baseline

* fix(ecstore): classify typed client network failures
This commit is contained in:
Zhengchao An
2026-08-26 11:24:56 +08:00
committed by GitHub
parent 7cac528de3
commit c0c208d89a
9 changed files with 176 additions and 147 deletions
+11
View File
@@ -198,6 +198,7 @@ pub(crate) fn message_has_network_needle(message: &str) -> bool {
pub(crate) fn is_network_like_disk_error(err: &DiskErrorType) -> bool {
match err {
DiskError::Timeout => true,
DiskError::RemoteClientUnavailable(detail) => message_has_network_needle(detail),
DiskError::Io(io_err) => {
if let Some(status) = embedded_tonic_status(io_err) {
return is_network_like_status(status);
@@ -719,6 +720,16 @@ mod tests {
assert!(!is_network_like_disk_error(&DiskError::FileNotFound));
}
#[test]
fn network_like_disk_error_keeps_typed_client_failures_classified() {
assert!(is_network_like_disk_error(&DiskError::RemoteClientUnavailable(
"transport error: connection refused".to_string()
)));
assert!(!is_network_like_disk_error(&DiskError::RemoteClientUnavailable(
"invalid client credentials".to_string()
)));
}
#[test]
fn test_signature_interceptor_keeps_auth_headers() {
ensure_test_rpc_secret();
@@ -626,13 +626,13 @@ impl PeerRestClient {
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::other(format!("peer {} is temporarily offline", self.grid_host)));
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::other(format!("can not get client, 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();
}
@@ -649,13 +649,13 @@ impl PeerRestClient {
> {
if self.offline.load(Ordering::Acquire) {
self.mark_offline_and_spawn_recovery();
return Err(Error::other(format!("peer {} is temporarily offline", self.grid_host)));
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::other(format!("can not get heal control client, 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();
}
@@ -668,13 +668,13 @@ impl PeerRestClient {
) -> Result<TierMutationControlServiceClient<InterceptedService<AuthenticatedChannel, TonicInterceptor>>> {
if self.offline.load(Ordering::Acquire) {
self.mark_offline_and_spawn_recovery();
return Err(Error::other(format!("peer {} is temporarily offline", self.grid_host)));
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::other(format!("can not get tier mutation control client, 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();
}
@@ -2213,17 +2213,14 @@ fn tier_config_reload_connection_outcome(err: Error) -> TierConfigReloadOutcome
}
fn is_tier_config_reload_connection_failure(err: &Error) -> bool {
let message = err.to_string();
// A bare "unavailable" is only trusted inside the local dial-failure
// wrapper from `get_client`, never in application text.
if message
.to_ascii_lowercase()
.split_once("can not get client, err:")
.is_some_and(|(_, local_error)| local_error.contains("unavailable"))
// 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(&message)
message_has_network_needle(&err.to_string())
}
/// Classifies a reload the peer answered but refused to apply.
@@ -3046,17 +3043,22 @@ mod tests {
TierConfigReloadOutcome::Terminal(_)
));
assert!(matches!(
tier_config_reload_connection_outcome(Error::other("can not get client, err: connection unavailable")),
tier_config_reload_connection_outcome(Error::RemoteClientUnavailable("connection unavailable".to_string())),
TierConfigReloadOutcome::TransientReconnect(_)
));
// The bare word is trusted only to the right of the dial-failure
// prefix, not anywhere in the message.
// 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
@@ -1090,7 +1090,7 @@ impl RemotePeerS3Client {
pub async fn get_client(&self) -> Result<NodeServiceClient<InterceptedService<AuthenticatedChannel, TonicInterceptor>>> {
node_service_time_out_client(&self.addr, TonicInterceptor::Signature(gen_tonic_signature_interceptor()))
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))
.map_err(|err| Error::RemoteClientUnavailable(err.to_string()))
}
/// Start health monitoring for the remote peer
+34 -120
View File
@@ -1379,7 +1379,7 @@ impl RemoteDisk {
let addr = addr.to_string();
let mut client = node_service_time_out_client(&addr, TonicInterceptor::Signature(gen_tonic_signature_interceptor()))
.await
.map_err(|err| (Error::other(format!("can not get client, err: {err}")), true))?;
.map_err(|err| (Error::RemoteClientUnavailable(err.to_string()), true))?;
let request = Request::new(DiskInfoRequest {
disk: endpoint.to_string(),
opts,
@@ -1719,7 +1719,7 @@ impl RemoteDisk {
/// recovers even without a background monitor. The recovery monitor's own probe path calls the
/// client directly and is unaffected.
fn offline_bypass_error(&self) -> Option<Error> {
internode_offline_bypass_reason(&self.addr).map(Error::other)
internode_offline_bypass_reason(&self.addr).map(Error::RemoteClientUnavailable)
}
async fn get_client(&self) -> Result<NodeServiceClient<InterceptedService<AuthenticatedChannel, TonicInterceptor>>> {
@@ -1728,7 +1728,7 @@ impl RemoteDisk {
}
node_service_time_out_client(&self.addr, TonicInterceptor::Signature(gen_tonic_signature_interceptor()))
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))
.map_err(|err| Error::RemoteClientUnavailable(err.to_string()))
}
/// Client for large `bytes`-carrying RPCs (ReadAll/WriteAll/ReadMultiple/BatchReadVersion).
@@ -1745,7 +1745,7 @@ impl RemoteDisk {
ChannelClass::Bulk,
)
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))
.map_err(|err| Error::RemoteClientUnavailable(err.to_string()))
}
async fn disk_ref(&self) -> String {
@@ -2015,10 +2015,7 @@ impl RemoteDisk {
|| async {
let file_info = compat_json(fi)?;
let file_info_bin = encode_file_info_msgpack(fi)?;
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(RenameDataRequest {
disk: self.endpoint.to_string(),
src_volume: src_volume.to_string(),
@@ -2088,10 +2085,7 @@ impl RemoteDisk {
self.execute_with_timeout(
|| async {
let options = serde_json::to_string(&opt)?;
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(DeleteRequest {
disk: self.endpoint.to_string(),
volume: volume.to_string(),
@@ -2214,10 +2208,7 @@ impl DiskAPI for RemoteDisk {
self.execute_with_timeout(
|| async {
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(MakeVolumeRequest {
disk: self.endpoint.to_string(),
volume: volume.to_string(),
@@ -2253,10 +2244,7 @@ impl DiskAPI for RemoteDisk {
self.execute_with_timeout(
|| async {
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(MakeVolumesRequest {
disk: self.endpoint.to_string(),
volumes: volumes.iter().map(|s| (*s).to_string()).collect(),
@@ -2291,10 +2279,7 @@ impl DiskAPI for RemoteDisk {
self.execute_with_timeout(
|| async {
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let request = Request::new(ListVolumesRequest {
disk: self.endpoint.to_string(),
});
@@ -2329,10 +2314,7 @@ impl DiskAPI for RemoteDisk {
self.execute_with_timeout(
|| async {
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let request = Request::new(StatVolumeRequest {
disk: self.endpoint.to_string(),
volume: volume.to_string(),
@@ -2368,10 +2350,7 @@ impl DiskAPI for RemoteDisk {
self.execute_with_timeout(
|| async {
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(DeleteVolumeRequest {
disk: self.endpoint.to_string(),
volume: volume.to_string(),
@@ -2423,10 +2402,7 @@ impl DiskAPI for RemoteDisk {
let file_info = serde_json::to_string(&fi)?;
let opts = serde_json::to_string(&opts)?;
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(DeleteVersionRequest {
disk: self.endpoint.to_string(),
volume: volume.to_string(),
@@ -2598,10 +2574,7 @@ impl DiskAPI for RemoteDisk {
self.execute_with_timeout(
|| async {
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(DeletePathsRequest {
disk: self.endpoint.to_string(),
volume: volume.to_string(),
@@ -2626,10 +2599,7 @@ impl DiskAPI for RemoteDisk {
async fn acquire_snapshot_lease(&self, volume: &str, path: &str) -> Result<SnapshotLeaseToken> {
self.execute_with_timeout(
|| async {
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(SnapshotLeaseRequest {
disk: self.endpoint.to_string(),
volume: volume.to_string(),
@@ -2650,10 +2620,7 @@ impl DiskAPI for RemoteDisk {
async fn renew_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> Result<SnapshotLeaseToken> {
self.execute_with_timeout(
|| async {
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(SnapshotLeaseRenewRequest {
disk: self.endpoint.to_string(),
volume: volume.to_string(),
@@ -2675,10 +2642,7 @@ impl DiskAPI for RemoteDisk {
async fn release_snapshot_lease(&self, volume: &str, path: &str, token: SnapshotLeaseToken) -> Result<()> {
self.execute_with_timeout(
|| async {
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(SnapshotLeaseReleaseRequest {
disk: self.endpoint.to_string(),
volume: volume.to_string(),
@@ -2718,10 +2682,7 @@ impl DiskAPI for RemoteDisk {
"write_metadata",
move || async move {
let disk = self.disk_ref().await;
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(WriteMetadataRequest {
disk,
volume: volume.to_string(),
@@ -2753,10 +2714,7 @@ impl DiskAPI for RemoteDisk {
"read_metadata",
|| async {
let disk = self.disk_ref().await;
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let request = Request::new(ReadMetadataRequest {
volume: volume.to_string(),
path: path.to_string(),
@@ -2798,10 +2756,7 @@ impl DiskAPI for RemoteDisk {
"update_metadata",
move || async move {
let disk = self.disk_ref().await;
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(UpdateMetadataRequest {
disk,
volume: volume.to_string(),
@@ -2864,10 +2819,7 @@ impl DiskAPI for RemoteDisk {
let opts_str = opts_str.clone();
let opts_bin = opts_bin.clone();
let disk = self.disk_ref().await;
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let request_payload_bytes = read_version_attribution_enabled.then(|| {
disk.len()
.saturating_add(volume.len())
@@ -2969,10 +2921,7 @@ impl DiskAPI for RemoteDisk {
move || async move {
let disk = self.disk_ref().await;
let disk_len = disk.len();
let mut client = self
.get_bulk_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_bulk_client().await?;
let request = Request::new(BatchReadVersionRequest {
disk,
batch_read_version_req,
@@ -3081,10 +3030,7 @@ impl DiskAPI for RemoteDisk {
self.execute_with_timeout(
|| async {
let disk = self.disk_ref().await;
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let request = Request::new(ReadXlRequest {
disk,
volume: volume.to_string(),
@@ -3129,10 +3075,7 @@ impl DiskAPI for RemoteDisk {
|| async {
let disk = self.disk_ref().await;
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let request = Request::new(ListDirRequest {
disk,
volume: volume.to_string(),
@@ -3393,10 +3336,7 @@ impl DiskAPI for RemoteDisk {
self.execute_with_timeout(
|| async {
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(RenameFileRequest {
disk: self.endpoint.to_string(),
src_volume: src_volume.to_string(),
@@ -3438,10 +3378,7 @@ impl DiskAPI for RemoteDisk {
self.execute_with_timeout(
|| async {
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(RenamePartRequest {
disk: self.endpoint.to_string(),
src_volume: src_volume.to_string(),
@@ -3477,10 +3414,7 @@ impl DiskAPI for RemoteDisk {
) -> Result<()> {
self.execute_with_timeout(
|| async {
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(PreparePartTransactionRequest {
disk: self.endpoint.to_string(),
src_volume: src_volume.to_string(),
@@ -3507,10 +3441,7 @@ impl DiskAPI for RemoteDisk {
async fn settle_part_transaction(&self, volume: &str, path: &str, action: PartTransactionAction) -> Result<()> {
self.execute_with_timeout(
|| async {
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let mut request = Request::new(SettlePartTransactionRequest {
disk: self.endpoint.to_string(),
volume: volume.to_string(),
@@ -3553,10 +3484,7 @@ impl DiskAPI for RemoteDisk {
self.execute_with_timeout(
|| async {
let file_info = serde_json::to_string(&fi)?;
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let request = Request::new(VerifyFileRequest {
disk: self.endpoint.to_string(),
volume: volume.to_string(),
@@ -3594,10 +3522,7 @@ impl DiskAPI for RemoteDisk {
);
self.execute_with_timeout(
|| async {
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let request = Request::new(ReadPartsRequest {
disk: self.endpoint.to_string(),
bucket: bucket.to_string(),
@@ -3635,10 +3560,7 @@ impl DiskAPI for RemoteDisk {
self.execute_with_timeout(
|| async {
let file_info = serde_json::to_string(&fi)?;
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let request = Request::new(CheckPartsRequest {
disk: self.endpoint.to_string(),
volume: volume.to_string(),
@@ -3680,10 +3602,7 @@ impl DiskAPI for RemoteDisk {
let read_multiple_req = compat_json(&req)?;
let read_multiple_req_bin = encode_msgpack(&req)?;
let disk = self.disk_ref().await;
let mut client = self
.get_bulk_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_bulk_client().await?;
let request = Request::new(ReadMultipleRequest {
disk,
read_multiple_req,
@@ -3728,9 +3647,8 @@ impl DiskAPI for RemoteDisk {
|| async {
let data_len = data.len();
let disk = self.disk_ref().await;
let mut client = self.get_bulk_client().await.map_err(|err| {
let mut client = self.get_bulk_client().await.inspect_err(|_| {
crate::cluster::rpc::runtime_sources::record_remote_disk_grpc_write_all_error();
Error::other(format!("can not get client, err: {err}"))
})?;
let mut request = Request::new(WriteAllRequest {
disk,
@@ -3785,9 +3703,8 @@ impl DiskAPI for RemoteDisk {
"read_all",
|| async {
let disk = self.disk_ref().await;
let mut client = self.get_bulk_client().await.map_err(|err| {
let mut client = self.get_bulk_client().await.inspect_err(|_| {
crate::cluster::rpc::runtime_sources::record_remote_disk_grpc_read_all_error();
Error::other(format!("can not get client, err: {err}"))
})?;
let request = Request::new(ReadAllRequest {
disk,
@@ -3824,10 +3741,7 @@ impl DiskAPI for RemoteDisk {
"disk_info",
|| async {
let opts = serde_json::to_string(&opts)?;
let mut client = self
.get_client()
.await
.map_err(|err| Error::other(format!("can not get client, err: {err}")))?;
let mut client = self.get_client().await?;
let request = Request::new(DiskInfoRequest {
disk: self.endpoint.to_string(),
opts,