From 65a7cc9cd4ea423a3741ad91b969d9200d7d9c5f Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Wed, 26 Aug 2026 12:32:53 +0800 Subject: [PATCH] refactor(replication): name the resync state error and keep io failures typed (#6628) Backlog#1845 step 7. The replication crate's hand-rolled, crate-generic Error type actually describes one thing: failures of the persisted resync/MRF state files. Rename it to ResyncStateError so the name says so, and stop collapsing io::Error into Other(String): a new Io(std::io::Error) variant keeps the kind and source chain, Display renders identically, and the ecstore boundary maps it to StorageError::Io so the kind survives into store-layer classification instead of degrading into a stringified other(). No thiserror introduced - the crate keeps its zero-internal-deps posture and hand-written impls. Ref rustfs/backlog#1845 --- .../replication_resync_boundary.rs | 9 +- crates/replication/src/lib.rs | 2 +- crates/replication/src/mrf.rs | 14 +- crates/replication/src/resync.rs | 129 +++++++++++------- 4 files changed, 92 insertions(+), 62 deletions(-) diff --git a/crates/ecstore/src/bucket/replication/replication_resync_boundary.rs b/crates/ecstore/src/bucket/replication/replication_resync_boundary.rs index 32e3f6c61..53db55c2d 100644 --- a/crates/ecstore/src/bucket/replication/replication_resync_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_resync_boundary.rs @@ -53,10 +53,13 @@ pub(crate) const MRF_META_FORMAT: u16 = rustfs_replication::mrf::MRF_META_FORMAT )] pub(crate) const MRF_META_VERSION: u16 = rustfs_replication::mrf::MRF_META_VERSION; -fn map_replication_error(err: rustfs_replication::Error) -> Error { +fn map_replication_error(err: rustfs_replication::ResyncStateError) -> Error { match err { - rustfs_replication::Error::CorruptedFormat => Error::CorruptedFormat, - rustfs_replication::Error::Other(err) => Error::other(err), + rustfs_replication::ResyncStateError::CorruptedFormat => Error::CorruptedFormat, + // Keep the io error typed so its kind survives into StorageError::Io + // instead of degrading to a stringified other() (backlog#1845). + rustfs_replication::ResyncStateError::Io(err) => Error::Io(err), + rustfs_replication::ResyncStateError::Other(err) => Error::other(err), } } diff --git a/crates/replication/src/lib.rs b/crates/replication/src/lib.rs index def314f95..e35280749 100644 --- a/crates/replication/src/lib.rs +++ b/crates/replication/src/lib.rs @@ -78,7 +78,7 @@ pub use queue::{ replication_heal_queue_action, worker_queue_for_replication_type, }; pub use resync::{ - BucketReplicationResyncStatus, Error, RESYNC_FILE_MAX_BYTES, Result, ResyncOpts, ResyncStatusType, + BucketReplicationResyncStatus, RESYNC_FILE_MAX_BYTES, Result, ResyncOpts, ResyncStateError, ResyncStatusType, TargetReplicationResyncStatus, decode_resync_file, encode_resync_file, is_version_id_mismatch, resync_state_accepts_update, resync_status_duration, sanitize_resync_error_detail, should_auto_resume_resync, should_count_head_proxy_failure, }; diff --git a/crates/replication/src/mrf.rs b/crates/replication/src/mrf.rs index f488cf1a4..8698e7c2a 100644 --- a/crates/replication/src/mrf.rs +++ b/crates/replication/src/mrf.rs @@ -15,7 +15,7 @@ use byteorder::{ByteOrder, LittleEndian}; use std::fmt; -use crate::{Error, Result}; +use crate::{Result, ResyncStateError}; pub use crate::filemeta::{MrfOpKind, MrfReplicateEntry}; @@ -569,7 +569,7 @@ impl MrfV2Envelope { } } pub fn encode_mrf_file(entries: &[MrfReplicateEntry]) -> Result> { - let payload = rmp_serde::to_vec_named(entries).map_err(|e| Error::Other(e.to_string()))?; + let payload = rmp_serde::to_vec_named(entries).map_err(|e| ResyncStateError::Other(e.to_string()))?; let mut data = Vec::with_capacity(4 + payload.len()); let mut fmt = [0u8; 2]; LittleEndian::write_u16(&mut fmt, MRF_META_FORMAT); @@ -583,19 +583,19 @@ pub fn encode_mrf_file(entries: &[MrfReplicateEntry]) -> Result> { pub fn decode_mrf_file(data: &[u8]) -> Result> { if data.len() <= 4 { - return Err(Error::CorruptedFormat); + return Err(ResyncStateError::CorruptedFormat); } let mut fmt = [0u8; 2]; fmt.copy_from_slice(&data[0..2]); if LittleEndian::read_u16(&fmt) != MRF_META_FORMAT { - return Err(Error::CorruptedFormat); + return Err(ResyncStateError::CorruptedFormat); } let mut ver = [0u8; 2]; ver.copy_from_slice(&data[2..4]); if LittleEndian::read_u16(&ver) != MRF_META_VERSION { - return Err(Error::CorruptedFormat); + return Err(ResyncStateError::CorruptedFormat); } - rmp_serde::from_slice(&data[4..]).map_err(|e| Error::Other(e.to_string())) + rmp_serde::from_slice(&data[4..]).map_err(|e| ResyncStateError::Other(e.to_string())) } #[cfg(test)] @@ -755,7 +755,7 @@ mod tests { data.extend_from_slice(&MRF_META_VERSION.to_le_bytes()); data.push(0x90); - assert!(matches!(decode_mrf_file(&data), Err(Error::CorruptedFormat))); + assert!(matches!(decode_mrf_file(&data), Err(ResyncStateError::CorruptedFormat))); } #[test] diff --git a/crates/replication/src/resync.rs b/crates/replication/src/resync.rs index 96d0e6b3f..12425941c 100644 --- a/crates/replication/src/resync.rs +++ b/crates/replication/src/resync.rs @@ -49,62 +49,74 @@ const RESYNC_ERROR_SENSITIVE_MARKERS: &[&str] = &[ "password", ]; -pub type Result = std::result::Result; +pub type Result = std::result::Result; +/// Error type for the persisted resync/MRF state files (backlog#1845 step 7: +/// renamed from the crate-generic `Error`, with `Io` kept typed instead of +/// collapsing into a string so callers can still classify the io failure). #[derive(Debug)] -pub enum Error { +pub enum ResyncStateError { CorruptedFormat, + Io(std::io::Error), Other(String), } -impl Error { +impl ResyncStateError { fn other(err: impl Into) -> Self { Self::Other(err.into()) } } -impl fmt::Display for Error { +impl fmt::Display for ResyncStateError { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { match self { Self::CorruptedFormat => write!(f, "corrupted format"), + Self::Io(err) => write!(f, "{err}"), Self::Other(err) => write!(f, "{err}"), } } } -impl std::error::Error for Error {} - -impl From for Error { - fn from(err: std::io::Error) -> Self { - Self::other(err.to_string()) +impl std::error::Error for ResyncStateError { + fn source(&self) -> Option<&(dyn std::error::Error + 'static)> { + match self { + Self::Io(err) => Some(err), + _ => None, + } } } -impl From for Error { +impl From for ResyncStateError { + fn from(err: std::io::Error) -> Self { + Self::Io(err) + } +} + +impl From for ResyncStateError { fn from(err: std::string::FromUtf8Error) -> Self { Self::other(err.to_string()) } } -impl From for Error { +impl From for ResyncStateError { fn from(err: rmp::encode::ValueWriteError) -> Self { Self::other(err.to_string()) } } -impl From for Error { +impl From for ResyncStateError { fn from(err: rmp::decode::ValueReadError) -> Self { Self::other(err.to_string()) } } -impl From for Error { +impl From for ResyncStateError { fn from(err: rmp::decode::NumValueReadError) -> Self { Self::other(err.to_string()) } } -impl From for Error { +impl From for ResyncStateError { fn from(err: rmp_serde::decode::Error) -> Self { Self::other(err.to_string()) } @@ -360,7 +372,7 @@ impl BucketReplicationResyncStatus { rmp::encode::write_str(&mut wr, "v")?; rmp::encode::write_i32(&mut wr, i32::from(self.version))?; rmp::encode::write_str(&mut wr, "brs")?; - let target_count = u32::try_from(self.targets_map.len()).map_err(|_| Error::CorruptedFormat)?; + let target_count = u32::try_from(self.targets_map.len()).map_err(|_| ResyncStateError::CorruptedFormat)?; rmp::encode::write_map_len(&mut wr, target_count)?; for (arn, status) in &self.targets_map { rmp::encode::write_str(&mut wr, arn)?; @@ -386,11 +398,11 @@ impl BucketReplicationResyncStatus { match key.as_str() { "v" => { let v: i32 = rmp::decode::read_int(&mut rd)?; - out.version = u16::try_from(v).map_err(|_| Error::other("invalid resync version"))?; + out.version = u16::try_from(v).map_err(|_| ResyncStateError::other("invalid resync version"))?; } "brs" => { let map_len = rmp::decode::read_map_len(&mut rd)?; - let target_count = usize::try_from(map_len).map_err(|_| Error::CorruptedFormat)?; + let target_count = usize::try_from(map_len).map_err(|_| ResyncStateError::CorruptedFormat)?; let mut targets = HashMap::with_capacity(target_count); for _ in 0..map_len { let arn = read_msgp_str(&mut rd)?; @@ -436,19 +448,19 @@ pub fn encode_resync_file(status: &BucketReplicationResyncStatus) -> Result Result { if data.len() <= RESYNC_FILE_HEADER_LEN || data.len() > RESYNC_FILE_MAX_BYTES { - return Err(Error::CorruptedFormat); + return Err(ResyncStateError::CorruptedFormat); } let mut major = [0u8; 2]; major.copy_from_slice(&data[0..2]); if LittleEndian::read_u16(&major) != RESYNC_META_FORMAT { - return Err(Error::CorruptedFormat); + return Err(ResyncStateError::CorruptedFormat); } let mut minor = [0u8; 2]; minor.copy_from_slice(&data[2..4]); if LittleEndian::read_u16(&minor) != RESYNC_META_VERSION { - return Err(Error::CorruptedFormat); + return Err(ResyncStateError::CorruptedFormat); } let status = match BucketReplicationResyncStatus::unmarshal_msg(&data[4..]) { @@ -456,7 +468,7 @@ pub fn decode_resync_file(data: &[u8]) -> Result Err(_) => BucketReplicationResyncStatus::unmarshal_legacy_msg(&data[4..])?, }; if status.version != RESYNC_META_VERSION { - return Err(Error::CorruptedFormat); + return Err(ResyncStateError::CorruptedFormat); } Ok(status) } @@ -478,28 +490,28 @@ fn wire_zero_time() -> OffsetDateTime { fn validate_msgp_payload(data: &[u8]) -> Result<()> { if data.len() > RESYNC_MSGP_MAX_BYTES { - return Err(Error::CorruptedFormat); + return Err(ResyncStateError::CorruptedFormat); } let mut rd = Cursor::new(data); let mut values = 0usize; validate_msgp_value(&mut rd, 0, &mut values)?; if usize::try_from(rd.position()).ok() != Some(data.len()) { - return Err(Error::CorruptedFormat); + return Err(ResyncStateError::CorruptedFormat); } Ok(()) } fn validate_msgp_value(rd: &mut R, depth: usize, values: &mut usize) -> Result<()> { if depth > RESYNC_MSGP_MAX_DEPTH { - return Err(Error::CorruptedFormat); + return Err(ResyncStateError::CorruptedFormat); } - *values = values.checked_add(1).ok_or(Error::CorruptedFormat)?; + *values = values.checked_add(1).ok_or(ResyncStateError::CorruptedFormat)?; if *values > RESYNC_MSGP_MAX_VALUES { - return Err(Error::CorruptedFormat); + return Err(ResyncStateError::CorruptedFormat); } - let marker = rmp::decode::read_marker(rd).map_err(|e| Error::other(format!("{e:?}")))?; + let marker = rmp::decode::read_marker(rd).map_err(|e| ResyncStateError::other(format!("{e:?}")))?; let skip_len = match marker { Marker::Null | Marker::False | Marker::True | Marker::FixPos(_) | Marker::FixNeg(_) => 0, Marker::U8 | Marker::I8 => 1, @@ -558,7 +570,7 @@ fn validate_msgp_value(rd: &mut R, depth: usize, values: &mut usize) -> skip_exact(rd, 1)?; return skip_exact(rd, len); } - Marker::Reserved => return Err(Error::CorruptedFormat), + Marker::Reserved => return Err(ResyncStateError::CorruptedFormat), }; let skip_len = validate_msgp_element_len(skip_len)?; skip_exact(rd, skip_len) @@ -566,10 +578,10 @@ fn validate_msgp_value(rd: &mut R, depth: usize, values: &mut usize) -> fn validate_msgp_collection(rd: &mut R, len: usize, depth: usize, values: &mut usize, is_map: bool) -> Result<()> { if len > RESYNC_MSGP_MAX_COLLECTION_ITEMS { - return Err(Error::CorruptedFormat); + return Err(ResyncStateError::CorruptedFormat); } let values_per_item = if is_map { 2 } else { 1 }; - let child_count = len.checked_mul(values_per_item).ok_or(Error::CorruptedFormat)?; + let child_count = len.checked_mul(values_per_item).ok_or(ResyncStateError::CorruptedFormat)?; for _ in 0..child_count { validate_msgp_value(rd, depth + 1, values)?; } @@ -578,7 +590,7 @@ fn validate_msgp_collection(rd: &mut R, len: usize, depth: usize, value fn validate_msgp_element_len(len: usize) -> Result { if len > RESYNC_MSGP_MAX_ELEMENT_BYTES { - return Err(Error::CorruptedFormat); + return Err(ResyncStateError::CorruptedFormat); } Ok(len) } @@ -591,11 +603,11 @@ fn read_msgp_str(rd: &mut R) -> Result { } fn read_msgp_time_or_nil(rd: &mut R) -> Result> { - let marker = rmp::decode::read_marker(rd).map_err(|e| Error::other(format!("{e:?}")))?; + let marker = rmp::decode::read_marker(rd).map_err(|e| ResyncStateError::other(format!("{e:?}")))?; match marker { Marker::Null => Ok(None), Marker::Ext8 => Ok(Some(read_msgp_ext8_time(rd)?)), - other => Err(Error::other(format!("expected time ext or nil, got marker: {other:?}"))), + other => Err(ResyncStateError::other(format!("expected time ext or nil, got marker: {other:?}"))), } } @@ -604,21 +616,21 @@ fn read_msgp_ext8_time(rd: &mut R) -> Result { rd.read_exact(&mut len_buf)?; let len = len_buf[0] as usize; if len != MSGP_TIME_LEN as usize { - return Err(Error::other(format!("invalid msgp time len: {len}"))); + return Err(ResyncStateError::other(format!("invalid msgp time len: {len}"))); } let mut type_buf = [0u8; 1]; rd.read_exact(&mut type_buf)?; if type_buf[0] != MSGP_TIME_EXT_TYPE as u8 { - return Err(Error::other(format!("invalid msgp time type: {}", type_buf[0]))); + return Err(ResyncStateError::other(format!("invalid msgp time type: {}", type_buf[0]))); } let mut buf = [0u8; 12]; rd.read_exact(&mut buf)?; let sec = BigEndian::read_i64(&buf[0..8]); let nsec = BigEndian::read_u32(&buf[8..12]); OffsetDateTime::from_unix_timestamp(sec) - .map_err(|_| Error::other("invalid timestamp"))? + .map_err(|_| ResyncStateError::other("invalid timestamp"))? .replace_nanosecond(nsec) - .map_err(|_| Error::other("invalid nanosecond")) + .map_err(|_| ResyncStateError::other("invalid nanosecond")) } fn write_msgp_time(wr: &mut W, time: OffsetDateTime) -> Result<()> { @@ -660,9 +672,9 @@ fn read_skip_len(rd: &mut R, bytes: usize) -> Result { 4 => { let mut buf = [0u8; 4]; rd.read_exact(&mut buf)?; - usize::try_from(u32::from_be_bytes(buf)).map_err(|_| Error::CorruptedFormat) + usize::try_from(u32::from_be_bytes(buf)).map_err(|_| ResyncStateError::CorruptedFormat) } - _ => Err(Error::other("invalid MessagePack length width")), + _ => Err(ResyncStateError::other("invalid MessagePack length width")), } } @@ -685,7 +697,7 @@ fn resync_status_from_i32(code: i32) -> Result { 3 => Ok(ResyncStatusType::ResyncStarted), 4 => Ok(ResyncStatusType::ResyncCompleted), 5 => Ok(ResyncStatusType::ResyncFailed), - _ => Err(Error::other(format!("invalid resync status code: {code}"))), + _ => Err(ResyncStateError::other(format!("invalid resync status code: {code}"))), } } @@ -847,14 +859,14 @@ mod tests { let mut data = encode_resync_file(&status).expect("resync status should encode"); data.push(0xc0); - assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + assert!(matches!(decode_resync_file(&data), Err(ResyncStateError::CorruptedFormat))); } #[test] fn resync_file_rejects_reserved_messagepack_marker() { let data = resync_file_with_unknown_value(&[0xc1]); - assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + assert!(matches!(decode_resync_file(&data), Err(ResyncStateError::CorruptedFormat))); } #[test] @@ -866,7 +878,7 @@ mod tests { let value = msgp_str32(RESYNC_MSGP_MAX_ELEMENT_BYTES + 1); let data = resync_file_with_unknown_value(&value); - assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + assert!(matches!(decode_resync_file(&data), Err(ResyncStateError::CorruptedFormat))); } #[test] @@ -899,7 +911,7 @@ mod tests { for value in [array, map] { let data = resync_file_with_unknown_value(&value); - assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + assert!(matches!(decode_resync_file(&data), Err(ResyncStateError::CorruptedFormat))); } } @@ -921,7 +933,7 @@ mod tests { let data = resync_file_with_unknown_value(&extension(RESYNC_MSGP_MAX_ELEMENT_BYTES + 1)); - assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + assert!(matches!(decode_resync_file(&data), Err(ResyncStateError::CorruptedFormat))); } #[test] @@ -935,7 +947,7 @@ mod tests { let mut rejected = vec![0x91; RESYNC_MSGP_MAX_DEPTH]; rejected.push(0xc0); let data = resync_file_with_unknown_value(&rejected); - assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + assert!(matches!(decode_resync_file(&data), Err(ResyncStateError::CorruptedFormat))); } #[test] @@ -973,7 +985,7 @@ mod tests { let data = resync_file_with_unknown_value(&nested_arrays(accepted_leaf_count + 1)); - assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + assert!(matches!(decode_resync_file(&data), Err(ResyncStateError::CorruptedFormat))); } #[test] @@ -1004,7 +1016,7 @@ mod tests { data.push(0xc0); assert_eq!(data.len(), RESYNC_FILE_MAX_BYTES + 1); - assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + assert!(matches!(decode_resync_file(&data), Err(ResyncStateError::CorruptedFormat))); } #[test] @@ -1030,7 +1042,7 @@ mod tests { ); } - assert!(matches!(encode_resync_file(&status), Err(Error::CorruptedFormat))); + assert!(matches!(encode_resync_file(&status), Err(ResyncStateError::CorruptedFormat))); } #[test] @@ -1044,7 +1056,7 @@ mod tests { }, ); - assert!(matches!(encode_resync_file(&status), Err(Error::CorruptedFormat))); + assert!(matches!(encode_resync_file(&status), Err(ResyncStateError::CorruptedFormat))); } #[test] @@ -1115,7 +1127,7 @@ mod tests { assert!(matches!( BucketReplicationResyncStatus::unmarshal_legacy_msg(&payload), - Err(Error::CorruptedFormat) + Err(ResyncStateError::CorruptedFormat) )); } @@ -1173,4 +1185,19 @@ mod tests { assert!(matches!(rmp::decode::read_marker(&mut cursor), Ok(Marker::FixPos(1)))); } } + + #[test] + fn io_errors_stay_typed_in_resync_state_error() { + let io_err = std::io::Error::new(std::io::ErrorKind::PermissionDenied, "state file locked"); + let err: ResyncStateError = io_err.into(); + match &err { + ResyncStateError::Io(inner) => { + assert_eq!(inner.kind(), std::io::ErrorKind::PermissionDenied); + assert_eq!(inner.to_string(), "state file locked"); + } + other => panic!("io::Error must stay typed, got {other:?}"), + } + assert_eq!(err.to_string(), "state file locked"); + assert!(std::error::Error::source(&err).is_some(), "Io must expose its source"); + } }