mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 04:21:35 +00:00
fix(ecstore): preserve RPC status when cloning storage errors (#8207)
* fix(ecstore): preserve RPC status when cloning storage errors * test(heal): retain clone coverage without redundant ownership
This commit is contained in:
@@ -436,6 +436,11 @@ impl DiskError {
|
||||
.is_some_and(|error| error.status().code() == tonic::Code::Cancelled)
|
||||
}
|
||||
|
||||
pub(crate) fn clone_rpc_status_io_error(error: &io::Error) -> Option<io::Error> {
|
||||
let status = error.get_ref()?.downcast_ref::<RpcStatusError>()?;
|
||||
Some(io::Error::new(error.kind(), RpcStatusError(status.0.clone())))
|
||||
}
|
||||
|
||||
pub fn internode_http_error_kind(&self) -> Option<InternodeHttpErrorKind> {
|
||||
match self {
|
||||
DiskError::Io(io_error) => io_error
|
||||
@@ -709,8 +714,8 @@ impl Clone for DiskError {
|
||||
DiskError::conditional_file_not_committed(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::<RpcStatusError>()) {
|
||||
return DiskError::Io(io::Error::new(io_error.kind(), RpcStatusError(status.0.clone())));
|
||||
if let Some(error) = Self::clone_rpc_status_io_error(io_error) {
|
||||
return DiskError::Io(error);
|
||||
}
|
||||
DiskError::Io(
|
||||
Self::clone_dangling_delete_grace(io_error)
|
||||
@@ -908,6 +913,25 @@ mod tests {
|
||||
use super::*;
|
||||
use std::collections::HashMap;
|
||||
|
||||
#[test]
|
||||
fn rpc_status_survives_disk_and_storage_clones() {
|
||||
for code in [tonic::Code::Cancelled, tonic::Code::PermissionDenied] {
|
||||
let status = tonic::Status::new(code, "operation was canceled");
|
||||
let original = DiskError::Io(io::Error::new(io::ErrorKind::Interrupted, RpcStatusError::from(status)));
|
||||
let display = original.to_string();
|
||||
let disk = original.clone();
|
||||
let storage = crate::error::StorageError::from(disk);
|
||||
let cloned = storage.clone();
|
||||
for error in [io::Error::from(original), io::Error::from(storage), io::Error::from(cloned)] {
|
||||
assert_eq!(error.kind(), io::ErrorKind::Interrupted);
|
||||
let status = error.get_ref().unwrap().downcast_ref::<RpcStatusError>().unwrap().status();
|
||||
assert_eq!(status.code(), code);
|
||||
assert_eq!(status.message(), "operation was canceled");
|
||||
assert_eq!(DiskError::from(error).to_string(), display);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn retired_marker_deferral_survives_disk_and_storage_clones() {
|
||||
let original = DiskError::retired_marker_deferred("missing retirement record");
|
||||
|
||||
@@ -623,8 +623,9 @@ impl Clone for StorageError {
|
||||
fn clone(&self) -> Self {
|
||||
match self {
|
||||
StorageError::Io(e) => {
|
||||
if let Some(error) =
|
||||
DiskError::clone_dangling_delete_grace(e).or_else(|| DiskError::clone_retired_marker_deferred(e))
|
||||
if let Some(error) = DiskError::clone_rpc_status_io_error(e)
|
||||
.or_else(|| DiskError::clone_dangling_delete_grace(e))
|
||||
.or_else(|| DiskError::clone_retired_marker_deferred(e))
|
||||
{
|
||||
return StorageError::Io(error);
|
||||
}
|
||||
|
||||
@@ -267,10 +267,16 @@ mod tests {
|
||||
#[test]
|
||||
fn cancelled_rpc_is_recoverable_across_storage_error_conversions() {
|
||||
let status = tonic::Status::cancelled("operation was canceled");
|
||||
let disk = DiskError::from(status.clone());
|
||||
let storage = EcstoreError::from(status.clone());
|
||||
let disk_storage = EcstoreError::from(DiskError::from(status.clone()));
|
||||
for error in [
|
||||
Error::Disk(DiskError::from(status.clone())),
|
||||
Error::Storage(EcstoreError::from(status.clone())),
|
||||
Error::Storage(EcstoreError::from(DiskError::from(status.clone()))),
|
||||
Error::Disk(disk.clone()),
|
||||
Error::Disk(disk),
|
||||
Error::Storage(storage.clone()),
|
||||
Error::Storage(storage),
|
||||
Error::Storage(disk_storage.clone()),
|
||||
Error::Storage(disk_storage),
|
||||
Error::Io(std::io::Error::from(DiskError::from(status))),
|
||||
] {
|
||||
assert!(error.is_recoverable_heal(), "typed RPC cancellation must be recoverable: {error:?}");
|
||||
@@ -282,7 +288,7 @@ mod tests {
|
||||
tonic::Status::internal("operation was canceled"),
|
||||
] {
|
||||
assert!(
|
||||
!Error::Storage(EcstoreError::from(status)).is_recoverable_heal(),
|
||||
!Error::Storage(EcstoreError::from(status).clone()).is_recoverable_heal(),
|
||||
"cancellation text alone must not change application error classification"
|
||||
);
|
||||
}
|
||||
|
||||
@@ -2476,7 +2476,11 @@ mod resume_loop_tests {
|
||||
}
|
||||
HealOutcome::Transient => Ok((HealResultItem::default(), Some(Error::Storage(EcstoreError::DiskNotFound)))),
|
||||
HealOutcome::RpcCancelled(_) => {
|
||||
Err(Error::Storage(EcstoreError::from(tonic::Status::cancelled("injected peer cancellation"))))
|
||||
// Pool aggregation clones the selected error before returning it to heal.
|
||||
let error = EcstoreError::from(tonic::Status::cancelled("injected peer cancellation"));
|
||||
let cloned = error.clone();
|
||||
assert_eq!(cloned, error, "pool error clone must preserve its kind and message");
|
||||
Err(Error::Storage(cloned))
|
||||
}
|
||||
HealOutcome::Cancelled => Err(Error::TaskCancelled),
|
||||
HealOutcome::Timeout => Err(Error::TaskTimeout),
|
||||
|
||||
Reference in New Issue
Block a user