From a7deb27bcf36604221f48059ea8ed00ae74c3378 Mon Sep 17 00:00:00 2001 From: overtrue Date: Tue, 15 Sep 2026 01:34:58 +0800 Subject: [PATCH] fix(heal): retain RPC transport failures across metadata probes (cherry picked from commit f48ec9f0c6e37b1f941620c666a9357c4b892d68) --- crates/ecstore/src/disk/error.rs | 23 ++++++++------ crates/ecstore/src/set_disk/ops/heal.rs | 41 +++++++++++++++++++++++++ 2 files changed, 55 insertions(+), 9 deletions(-) diff --git a/crates/ecstore/src/disk/error.rs b/crates/ecstore/src/disk/error.rs index 3f124780f..881f93778 100644 --- a/crates/ecstore/src/disk/error.rs +++ b/crates/ecstore/src/disk/error.rs @@ -700,15 +700,20 @@ impl Clone for DiskError { DiskError::Io(io_error) if self.is_conditional_file_not_committed() => DiskError::Io( DiskError::conditional_file_not_committed(io::Error::new(io_error.kind(), io_error.to_string())), ), - DiskError::Io(io_error) => DiskError::Io( - Self::clone_dangling_delete_grace(io_error) - .or_else(|| Self::clone_retired_marker_deferred(io_error)) - .or_else(|| rustfs_rio::clone_internode_http_io_error(io_error)) - .and_then(std::io::Error::into_inner) - // The helper derives a kind from the source; Clone must retain the original outer kind. - .map(|source| std::io::Error::new(io_error.kind(), source)) - .unwrap_or_else(|| std::io::Error::new(io_error.kind(), io_error.to_string())), - ), + DiskError::Io(io_error) => { + if let Some(status) = io_error.get_ref().and_then(|source| source.downcast_ref::()) { + return DiskError::Io(io::Error::new(io_error.kind(), RpcStatusError(status.0.clone()))); + } + DiskError::Io( + Self::clone_dangling_delete_grace(io_error) + .or_else(|| Self::clone_retired_marker_deferred(io_error)) + .or_else(|| rustfs_rio::clone_internode_http_io_error(io_error)) + .and_then(std::io::Error::into_inner) + // The helper derives a kind from the source; Clone must retain the original outer kind. + .map(|source| std::io::Error::new(io_error.kind(), source)) + .unwrap_or_else(|| std::io::Error::new(io_error.kind(), io_error.to_string())), + ) + } DiskError::MaxVersionsExceeded => DiskError::MaxVersionsExceeded, DiskError::Unexpected => DiskError::Unexpected, DiskError::CorruptedFormat => DiskError::CorruptedFormat, diff --git a/crates/ecstore/src/set_disk/ops/heal.rs b/crates/ecstore/src/set_disk/ops/heal.rs index e187c4139..3a53057dd 100644 --- a/crates/ecstore/src/set_disk/ops/heal.rs +++ b/crates/ecstore/src/set_disk/ops/heal.rs @@ -55,6 +55,20 @@ pub(crate) struct HealedObjectAbsence { pub removed: bool, } +fn is_unavailable_heal_rpc(error: &DiskError) -> bool { + let DiskError::Io(error) = error else { + return false; + }; + let Some(status) = crate::cluster::rpc::client::embedded_tonic_status(error) else { + return false; + }; + // Tonic can report a locally observed HTTP/2 GOAWAY as Internal. Only a + // typed transport source identifies that case; peer-supplied text cannot. + status.code() == tonic::Code::Unavailable + || (matches!(status.code(), tonic::Code::Internal | tonic::Code::Unknown) + && std::error::Error::source(status).is_some_and(|source| source.is::())) +} + fn heal_drive_state_for_error(error: &DiskError) -> DriveState { match error { DiskError::DiskNotFound | DiskError::RemoteClientUnavailable(_) => DriveState::Offline, @@ -65,6 +79,7 @@ fn heal_drive_state_for_error(error: &DiskError) -> DriveState { | DiskError::PartMissingOrCorrupt | DiskError::OutdatedXLMeta => DriveState::Missing, DiskError::FileCorrupt => DriveState::Corrupt, + _ if is_unavailable_heal_rpc(error) => DriveState::Offline, _ => DriveState::Unknown(error.to_string()), } } @@ -3168,6 +3183,32 @@ mod heal_result_report_tests { ); } + #[tokio::test] + async fn heal_preserves_rpc_transport_failure_through_metadata_classification() { + let transport = tonic::transport::Endpoint::from_static("http://127.0.0.1:0") + .connect() + .await + .expect_err("port zero cannot serve an internode connection"); + let mut interrupted = tonic::Status::internal("h2 protocol error: http2 error"); + interrupted.set_source(Arc::new(transport)); + for (status, offline) in [ + (interrupted, true), + (tonic::Status::unavailable("peer restarting"), true), + (tonic::Status::internal("h2 protocol error: http2 error"), false), + (tonic::Status::permission_denied("connection reset"), false), + (tonic::Status::data_loss("broken pipe"), false), + ] { + let error = DiskError::from(status); + let (_, _, reason) = super::should_heal_object_on_disk(&Some(error), &[], &FileInfo::default(), &FileInfo::default()); + let reason = reason.expect("metadata probe failure must survive classification"); + assert_eq!( + matches!(super::heal_drive_state_for_error(&reason), DriveState::Offline), + offline, + "{reason:?}" + ); + } + } + #[test] fn read_repair_commit_fingerprint_tracks_commit_identity_only() { let data_dir = Uuid::parse_str("11111111-1111-1111-1111-111111111111").expect("data dir should parse");