From cfce7bd9b12011ccf4d3c0a7d50c599c88334117 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Tue, 4 Aug 2026 17:25:28 +0800 Subject: [PATCH] fix(filemeta): preserve FileInfo wire compatibility (#5689) --- Cargo.lock | 1 + crates/ecstore/src/cluster/rpc/remote_disk.rs | 59 +- crates/filemeta/Cargo.toml | 1 + crates/filemeta/src/fileinfo.rs | 565 +++++++++++++++++- docs/architecture/compat-cleanup-register.md | 1 + .../minio-rustfs-router-compatibility.md | 2 +- rustfs/src/storage/rpc/node_service/disk.rs | 28 +- scripts/check_logging_guardrails.sh | 10 +- 8 files changed, 625 insertions(+), 42 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index c6da4be88..ed8771c70 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9417,6 +9417,7 @@ dependencies = [ "rustfs-utils", "s3s", "serde", + "serde_json", "tempfile", "thiserror 2.0.19", "time", diff --git a/crates/ecstore/src/cluster/rpc/remote_disk.rs b/crates/ecstore/src/cluster/rpc/remote_disk.rs index 0d8d507e4..b319b2083 100644 --- a/crates/ecstore/src/cluster/rpc/remote_disk.rs +++ b/crates/ecstore/src/cluster/rpc/remote_disk.rs @@ -1117,15 +1117,31 @@ impl RemoteDisk { } /// Initial capacity hint (bytes) for msgpack encode buffers, sized to cover a typical single- -/// version `FileInfo` without repeated growth reallocations. Larger payloads still grow as needed. +/// request without repeated growth reallocations. Larger payloads still grow as needed. const MSGPACK_ENCODE_CAPACITY_HINT: usize = 512; +const FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT: usize = 1024; -fn encode_msgpack(value: &T) -> Result> { - let mut serializer = rmp_serde::Serializer::new(Vec::with_capacity(MSGPACK_ENCODE_CAPACITY_HINT)); +fn encode_msgpack_with_capacity(value: &T, capacity: usize) -> Result> { + let mut serializer = rmp_serde::Serializer::new(Vec::with_capacity(capacity)); value.serialize(&mut serializer)?; Ok(serializer.into_inner()) } +fn encode_msgpack(value: &T) -> Result> { + encode_msgpack_with_capacity(value, MSGPACK_ENCODE_CAPACITY_HINT) +} + +fn encode_file_info_msgpack(value: &FileInfo) -> Result> { + encode_msgpack_with_capacity(value, FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT) +} + +fn encode_file_info_versions_msgpack(value: &FileInfoVersions) -> Result> { + let version_count = value.versions.len().saturating_add(value.free_versions.len()); + let capacity = + MSGPACK_ENCODE_CAPACITY_HINT.saturating_add(FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT.saturating_mul(version_count)); + encode_msgpack_with_capacity(value, capacity) +} + /// JSON compatibility string for a dual-encoded (`_bin` + text) request field. Returns an empty /// string only when msgpack-only mode and its explicit fleet confirmation guard are both enabled; /// otherwise the legacy JSON encoding is retained for old peers. @@ -1136,12 +1152,6 @@ fn compat_json(value: &T) -> Result { Ok(serde_json::to_string(value)?) } -fn encode_msgpack_named(value: &T) -> Result> { - let mut serializer = rmp_serde::Serializer::new(Vec::with_capacity(MSGPACK_ENCODE_CAPACITY_HINT)).with_struct_map(); - value.serialize(&mut serializer)?; - Ok(serializer.into_inner()) -} - fn decode_msgpack_or_json(binary: &[u8], json: &str, value_name: &'static str) -> Result { if !binary.is_empty() { let mut deserializer = rmp_serde::Deserializer::new(Cursor::new(binary)); @@ -1580,7 +1590,7 @@ impl DiskAPI for RemoteDisk { || async { // `_bin` support for DeleteVersion is new (grpc-optimization P2); always dual-write // JSON + msgpack until its fallback counter has read zero across a release window. - let file_info_bin = encode_msgpack(&fi)?; + let file_info_bin = encode_file_info_msgpack(&fi)?; let opts_bin = encode_msgpack(&opts)?; let file_info = serde_json::to_string(&fi)?; let opts = serde_json::to_string(&opts)?; @@ -1670,7 +1680,7 @@ impl DiskAPI for RemoteDisk { return errors; } }); - versions_bin.push(match encode_msgpack(file_info_versions) { + versions_bin.push(match encode_file_info_versions_msgpack(file_info_versions) { Ok(versions_bin) => Bytes::from(versions_bin), Err(err) => { let mut errors = Vec::with_capacity(versions.len()); @@ -1886,7 +1896,7 @@ impl DiskAPI for RemoteDisk { "Remote disk RPC started" ); let file_info = compat_json(&fi)?; - let file_info_bin = encode_msgpack(&fi)?; + let file_info_bin = encode_file_info_msgpack(&fi)?; self.execute_with_timeout_for_op( "write_metadata", @@ -1965,7 +1975,7 @@ impl DiskAPI for RemoteDisk { ); let file_info = compat_json(&fi)?; let opts_str = compat_json(&opts)?; - let file_info_bin = encode_msgpack(&fi)?; + let file_info_bin = encode_file_info_msgpack(&fi)?; let opts_bin = encode_msgpack(opts)?; self.execute_with_timeout_for_op( @@ -2229,7 +2239,7 @@ impl DiskAPI for RemoteDisk { "rename_data", || async { let file_info = compat_json(&fi)?; - let file_info_bin = encode_msgpack_named(&fi)?; + let file_info_bin = encode_file_info_msgpack(&fi)?; let mut client = self .get_client() .await @@ -3733,26 +3743,13 @@ mod tests { fn rename_data_file_info_named_msgpack_is_smaller_than_json() { let file_info = sample_rename_data_file_info(); let json = serde_json::to_vec(&file_info).expect("file info json should encode"); - let named_msgpack = encode_msgpack_named(&file_info).expect("file info named msgpack should encode"); + let named_msgpack = encode_file_info_msgpack(&file_info).expect("file info named msgpack should encode"); assert!( - named_msgpack.len() < json.len(), - "expected named msgpack payload to be smaller than json (msgpack={}, json={})", - named_msgpack.len(), - json.len() + named_msgpack.len() <= FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT, + "typical FileInfo should fit the msgpack capacity hint (msgpack={}, hint={FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT})", + named_msgpack.len() ); - } - - #[test] - fn rename_data_resp_named_msgpack_is_smaller_than_json() { - let response = RenameDataResp { - old_data_dir: Some(Uuid::new_v4()), - sign: Some(vec![1_u8; 32]), - old_current_size: Some(crate::disk::OldCurrentSize::Present(4096)), - }; - let json = serde_json::to_vec(&response).expect("rename data response json should encode"); - let named_msgpack = encode_msgpack_named(&response).expect("rename data response named msgpack should encode"); - assert!( named_msgpack.len() < json.len(), "expected named msgpack payload to be smaller than json (msgpack={}, json={})", diff --git a/crates/filemeta/Cargo.toml b/crates/filemeta/Cargo.toml index 78d6e1c1c..c2e7ee95f 100644 --- a/crates/filemeta/Cargo.toml +++ b/crates/filemeta/Cargo.toml @@ -54,6 +54,7 @@ arc-swap.workspace = true criterion = { workspace = true, features = ["html_reports"] } tempfile = { workspace = true } proptest = "1" +serde_json.workspace = true [[bench]] name = "xl_meta_bench" diff --git a/crates/filemeta/src/fileinfo.rs b/crates/filemeta/src/fileinfo.rs index e1f1f2e1b..bdfbb8348 100644 --- a/crates/filemeta/src/fileinfo.rs +++ b/crates/filemeta/src/fileinfo.rs @@ -22,6 +22,8 @@ use rustfs_utils::http::{ }; use s3s::dto::{RestoreStatus, Timestamp}; use s3s::header::X_AMZ_RESTORE; +use serde::de::{self, MapAccess, SeqAccess, Visitor, value::MapAccessDeserializer}; +use serde::ser::SerializeMap; use serde::{Deserialize, Serialize}; use std::collections::HashMap; use time::{OffsetDateTime, format_description::well_known::Rfc3339}; @@ -229,7 +231,7 @@ pub enum TransitionVersionState { Exact, } -#[derive(Serialize, Deserialize, Debug, PartialEq, Clone, Default)] +#[derive(Debug, PartialEq, Clone, Default)] pub struct FileInfo { pub volume: String, pub name: String, @@ -240,9 +242,7 @@ pub struct FileInfo { pub transitioned_objname: String, pub transition_tier: String, pub transition_version_id: Option, - #[serde(default)] pub transition_version: Option, - #[serde(default)] pub transition_version_state: TransitionVersionState, pub expire_restored: bool, pub data_dir: Option, @@ -271,6 +271,246 @@ pub struct FileInfo { pub uses_legacy_checksum: bool, } +#[derive(Deserialize)] +#[serde(remote = "FileInfo")] +struct FileInfoMapDef { + volume: String, + name: String, + version_id: Option, + is_latest: bool, + deleted: bool, + transition_status: String, + transitioned_objname: String, + transition_tier: String, + transition_version_id: Option, + #[serde(default)] + transition_version: Option, + #[serde(default)] + transition_version_state: TransitionVersionState, + expire_restored: bool, + data_dir: Option, + mod_time: Option, + size: i64, + mode: Option, + written_by_version: Option, + metadata: HashMap, + parts: Vec, + erasure: ErasureInfo, + mark_deleted: bool, + replication_state_internal: Option, + data: Option, + num_versions: usize, + successor_mod_time: Option, + fresh: bool, + idx: usize, + checksum: Option, + versioned: bool, + uses_legacy_checksum: bool, +} + +#[derive(Deserialize)] +#[serde(untagged)] +enum TransitionVersionOrExpireRestored { + TransitionVersion(Option), + ExpireRestored(bool), +} + +const FILE_INFO_FIELDS: &[&str] = &[ + "volume", + "name", + "version_id", + "is_latest", + "deleted", + "transition_status", + "transitioned_objname", + "transition_tier", + "transition_version_id", + "transition_version", + "transition_version_state", + "expire_restored", + "data_dir", + "mod_time", + "size", + "mode", + "written_by_version", + "metadata", + "parts", + "erasure", + "mark_deleted", + "replication_state_internal", + "data", + "num_versions", + "successor_mod_time", + "fresh", + "idx", + "checksum", + "versioned", + "uses_legacy_checksum", +]; + +impl Serialize for FileInfo { + fn serialize(&self, serializer: S) -> std::result::Result + where + S: serde::Serializer, + { + let mut map = serializer.serialize_map(Some(FILE_INFO_FIELDS.len()))?; + map.serialize_entry("volume", &self.volume)?; + map.serialize_entry("name", &self.name)?; + map.serialize_entry("version_id", &self.version_id)?; + map.serialize_entry("is_latest", &self.is_latest)?; + map.serialize_entry("deleted", &self.deleted)?; + map.serialize_entry("transition_status", &self.transition_status)?; + map.serialize_entry("transitioned_objname", &self.transitioned_objname)?; + map.serialize_entry("transition_tier", &self.transition_tier)?; + map.serialize_entry("transition_version_id", &self.transition_version_id)?; + map.serialize_entry("transition_version", &self.transition_version)?; + map.serialize_entry("transition_version_state", &self.transition_version_state)?; + map.serialize_entry("expire_restored", &self.expire_restored)?; + map.serialize_entry("data_dir", &self.data_dir)?; + map.serialize_entry("mod_time", &self.mod_time)?; + map.serialize_entry("size", &self.size)?; + map.serialize_entry("mode", &self.mode)?; + map.serialize_entry("written_by_version", &self.written_by_version)?; + map.serialize_entry("metadata", &self.metadata)?; + map.serialize_entry("parts", &self.parts)?; + map.serialize_entry("erasure", &self.erasure)?; + map.serialize_entry("mark_deleted", &self.mark_deleted)?; + map.serialize_entry("replication_state_internal", &self.replication_state_internal)?; + map.serialize_entry("data", &self.data)?; + map.serialize_entry("num_versions", &self.num_versions)?; + map.serialize_entry("successor_mod_time", &self.successor_mod_time)?; + map.serialize_entry("fresh", &self.fresh)?; + map.serialize_entry("idx", &self.idx)?; + map.serialize_entry("checksum", &self.checksum)?; + map.serialize_entry("versioned", &self.versioned)?; + map.serialize_entry("uses_legacy_checksum", &self.uses_legacy_checksum)?; + map.end() + } +} + +impl<'de> Deserialize<'de> for FileInfo { + fn deserialize(deserializer: D) -> std::result::Result + where + D: serde::Deserializer<'de>, + { + struct FileInfoVisitor; + + impl<'de> Visitor<'de> for FileInfoVisitor { + type Value = FileInfo; + + fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter.write_str("a FileInfo map or supported positional array") + } + + fn visit_map(self, map: A) -> std::result::Result + where + A: MapAccess<'de>, + { + FileInfoMapDef::deserialize(MapAccessDeserializer::new(map)) + } + + fn visit_seq(self, mut seq: A) -> std::result::Result + where + A: SeqAccess<'de>, + { + // RUSTFS_COMPAT_TODO(rustfs-5509): beta.11 and beta.12 wrote incompatible positional arrays. Remove after every supported direct-upgrade release writes named maps and retained RPC payloads cannot contain either array. + let declared_len = seq.size_hint(); + if let Some(len) = declared_len + && len != 28 + && len != 30 + { + return Err(de::Error::invalid_length(len, &self)); + } + + macro_rules! next_field { + ($field:literal) => { + seq.next_element()? + .ok_or_else(|| de::Error::missing_field($field))? + }; + } + + let volume = next_field!("volume"); + let name = next_field!("name"); + let version_id = next_field!("version_id"); + let is_latest = next_field!("is_latest"); + let deleted = next_field!("deleted"); + let transition_status = next_field!("transition_status"); + let transitioned_objname = next_field!("transitioned_objname"); + let transition_tier = next_field!("transition_tier"); + let transition_version_id = next_field!("transition_version_id"); + let transition_or_expire = next_field!("transition_version or expire_restored"); + let (transition_version, transition_version_state, expire_restored) = match transition_or_expire { + TransitionVersionOrExpireRestored::TransitionVersion(transition_version) => ( + transition_version, + next_field!("transition_version_state"), + next_field!("expire_restored"), + ), + TransitionVersionOrExpireRestored::ExpireRestored(expire_restored) => { + (None, TransitionVersionState::Unknown, expire_restored) + } + }; + let data_dir = next_field!("data_dir"); + let mod_time = next_field!("mod_time"); + let size = next_field!("size"); + let mode = next_field!("mode"); + let written_by_version = next_field!("written_by_version"); + let metadata = next_field!("metadata"); + let parts = next_field!("parts"); + let erasure = next_field!("erasure"); + let mark_deleted = next_field!("mark_deleted"); + let replication_state_internal = next_field!("replication_state_internal"); + let data = next_field!("data"); + let num_versions = next_field!("num_versions"); + let successor_mod_time = next_field!("successor_mod_time"); + let fresh = next_field!("fresh"); + let idx = next_field!("idx"); + let checksum = next_field!("checksum"); + let versioned = next_field!("versioned"); + let uses_legacy_checksum = next_field!("uses_legacy_checksum"); + + if seq.next_element::()?.is_some() { + return Err(de::Error::invalid_length(declared_len.unwrap_or(29), &self)); + } + + Ok(FileInfo { + volume, + name, + version_id, + is_latest, + deleted, + transition_status, + transitioned_objname, + transition_tier, + transition_version_id, + transition_version, + transition_version_state, + expire_restored, + data_dir, + mod_time, + size, + mode, + written_by_version, + metadata, + parts, + erasure, + mark_deleted, + replication_state_internal, + data, + num_versions, + successor_mod_time, + fresh, + idx, + checksum, + versioned, + uses_legacy_checksum, + }) + } + } + + deserializer.deserialize_struct("FileInfo", FILE_INFO_FIELDS, FileInfoVisitor) + } +} + /// Selects the validation policy for a trusted operation boundary. /// /// This mode is deliberately caller-selected and is never inferred from @@ -1800,6 +2040,325 @@ mod tests { }) } + fn positional_compat_file_info() -> FileInfo { + FileInfo { + volume: "bucket".to_string(), + name: "object".to_string(), + transition_version_id: Some(Uuid::from_u128(1)), + transition_version: Some(Uuid::from_u128(1).to_string()), + transition_version_state: TransitionVersionState::Exact, + expire_restored: true, + size: -1, + fresh: true, + idx: 7, + versioned: true, + uses_legacy_checksum: true, + ..Default::default() + } + } + + #[derive(Clone, Copy)] + enum HistoricalFileInfoLayout { + Beta11, + Beta12, + } + + fn encode_historical_file_info(value: &FileInfo, layout: HistoricalFileInfoLayout) -> Vec { + let mut encoded = Vec::new(); + let field_count = match layout { + HistoricalFileInfoLayout::Beta11 => 28, + HistoricalFileInfoLayout::Beta12 => 30, + }; + rmp::encode::write_array_len(&mut encoded, field_count).expect("historical FileInfo array header should encode"); + + macro_rules! encode_fields { + ($($field:expr),+ $(,)?) => { + $($field + .serialize(&mut Serializer::new(&mut encoded)) + .expect("historical FileInfo field should encode");)+ + }; + } + + encode_fields!( + &value.volume, + &value.name, + &value.version_id, + &value.is_latest, + &value.deleted, + &value.transition_status, + &value.transitioned_objname, + &value.transition_tier, + &value.transition_version_id, + ); + if matches!(layout, HistoricalFileInfoLayout::Beta12) { + encode_fields!(&value.transition_version, &value.transition_version_state); + } + encode_fields!( + &value.expire_restored, + &value.data_dir, + &value.mod_time, + &value.size, + &value.mode, + &value.written_by_version, + &value.metadata, + &value.parts, + &value.erasure, + &value.mark_deleted, + &value.replication_state_internal, + &value.data, + &value.num_versions, + &value.successor_mod_time, + &value.fresh, + &value.idx, + &value.checksum, + &value.versioned, + &value.uses_legacy_checksum, + ); + encoded + } + + fn field_order_file_info() -> FileInfo { + let transition_version_id = Uuid::from_u128(12); + FileInfo { + volume: "volume-sentinel".to_string(), + name: "name-sentinel".to_string(), + version_id: Some(Uuid::from_u128(11)), + is_latest: true, + deleted: false, + transition_status: "transition-status".to_string(), + transitioned_objname: "transitioned-object".to_string(), + transition_tier: "transition-tier".to_string(), + transition_version_id: Some(transition_version_id), + transition_version: Some(transition_version_id.to_string()), + transition_version_state: TransitionVersionState::Exact, + expire_restored: true, + data_dir: Some(Uuid::from_u128(13)), + mod_time: Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(14)), + size: 15, + mode: Some(16), + written_by_version: Some(17), + metadata: [("metadata-key".to_string(), "metadata-value".to_string())] + .into_iter() + .collect(), + parts: vec![ObjectPartInfo { + etag: "part-etag".to_string(), + number: 18, + size: 19, + actual_size: 20, + mod_time: Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(21)), + index: Some(Bytes::from_static(b"part-index")), + checksums: Some( + [("part-checksum".to_string(), "checksum-value".to_string())] + .into_iter() + .collect(), + ), + error: Some("part-error".to_string()), + }], + erasure: ErasureInfo { + algorithm: "erasure-algorithm".to_string(), + data_blocks: 2, + parity_blocks: 1, + block_size: 1024, + index: 3, + distribution: vec![2, 1, 3], + checksums: vec![ChecksumInfo { + part_number: 22, + algorithm: HashAlgorithm::SHA256, + hash: Bytes::from_static(b"erasure-hash"), + }], + }, + mark_deleted: false, + replication_state_internal: Some(ReplicationState { + replicate_decision_str: "replication-decision".to_string(), + delete_marker: true, + ..Default::default() + }), + data: Some(Bytes::from_static(b"inline-data")), + num_versions: 23, + successor_mod_time: Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(24)), + fresh: true, + idx: 25, + checksum: Some(Bytes::from_static(b"combined-checksum")), + versioned: false, + uses_legacy_checksum: true, + } + } + + const BETA11_FILEINFO_FIXTURE: &[u8] = &[ + 220, 0, 28, 166, 98, 117, 99, 107, 101, 116, 166, 111, 98, 106, 101, 99, 116, 192, 194, 194, 160, 160, 160, 196, 16, 0, + 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 1, 195, 192, 192, 255, 192, 192, 128, 144, 151, 160, 0, 0, 0, 0, 144, 144, 194, + 192, 192, 0, 192, 195, 7, 192, 195, 195, + ]; + + const BETA12_FILEINFO_FIXTURE: &[u8] = &[ + 220, 0, 30, 166, 98, 117, 99, 107, 101, 116, 166, 111, 98, 106, 101, 99, 116, 192, 194, 194, 160, 160, 160, 196, 16, 0, + 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 1, 217, 36, 48, 48, 48, 48, 48, 48, 48, 48, 45, 48, 48, 48, 48, 45, 48, 48, 48, + 48, 45, 48, 48, 48, 48, 45, 48, 48, 48, 48, 48, 48, 48, 48, 48, 48, 48, 49, 165, 101, 120, 97, 99, 116, 195, 192, 192, + 255, 192, 192, 128, 144, 151, 160, 0, 0, 0, 0, 144, 144, 194, 192, 192, 0, 192, 195, 7, 192, 195, 195, + ]; + + #[test] + fn fileinfo_decodes_beta11_positional_layout() { + let expected = positional_compat_file_info(); + let decoded = FileInfo::unmarshal(BETA11_FILEINFO_FIXTURE).expect("beta.11 positional FileInfo should decode"); + + assert_eq!(decoded.volume, expected.volume); + assert_eq!(decoded.expire_restored, expected.expire_restored); + assert_eq!(decoded.idx, expected.idx); + assert_eq!(decoded.transition_version, None); + assert_eq!(decoded.transition_version_state, TransitionVersionState::Unknown); + } + + #[test] + fn fileinfo_decodes_beta12_positional_layout() { + let expected = positional_compat_file_info(); + let decoded = FileInfo::unmarshal(BETA12_FILEINFO_FIXTURE).expect("beta.12 positional FileInfo should decode"); + + assert_eq!(decoded, expected); + } + + #[test] + fn fileinfo_historical_layouts_preserve_all_field_positions() { + let expected = field_order_file_info(); + let beta12 = encode_historical_file_info(&expected, HistoricalFileInfoLayout::Beta12); + assert_eq!(FileInfo::unmarshal(&beta12).expect("beta.12 full-field FileInfo should decode"), expected); + + let beta11 = encode_historical_file_info(&expected, HistoricalFileInfoLayout::Beta11); + let mut beta11_expected = expected; + beta11_expected.transition_version = None; + beta11_expected.transition_version_state = TransitionVersionState::Unknown; + assert_eq!( + FileInfo::unmarshal(&beta11).expect("beta.11 full-field FileInfo should decode"), + beta11_expected + ); + } + + #[test] + fn fileinfo_rejects_unsupported_positional_lengths() { + let mut fields_29 = BETA11_FILEINFO_FIXTURE.to_vec(); + fields_29[2] = 29; + fields_29.push(0xc0); + let error = FileInfo::unmarshal(&fields_29).expect_err("29-field FileInfo must fail closed"); + assert!(matches!(error, Error::RmpSerdeDecode(message) if message.contains("invalid length 29"))); + + let mut fields_31 = BETA12_FILEINFO_FIXTURE.to_vec(); + fields_31[2] = 31; + fields_31.push(0xc0); + let error = FileInfo::unmarshal(&fields_31).expect_err("31-field FileInfo must fail closed"); + assert!(matches!(error, Error::RmpSerdeDecode(message) if message.contains("invalid length 31"))); + } + + #[test] + fn fileinfo_rejects_truncated_historical_layouts() { + let beta11 = &BETA11_FILEINFO_FIXTURE[..BETA11_FILEINFO_FIXTURE.len() - 1]; + assert!(matches!(FileInfo::unmarshal(beta11), Err(Error::RmpSerdeDecode(_)))); + + let state_offset = BETA12_FILEINFO_FIXTURE + .windows(6) + .position(|window| window == [0xa5, b'e', b'x', b'a', b'c', b't']) + .expect("beta.12 fixture should contain the exact transition state"); + assert!(matches!( + FileInfo::unmarshal(&BETA12_FILEINFO_FIXTURE[..state_offset]), + Err(Error::RmpSerdeDecode(_)) + )); + } + + fn wrap_historical_file_info(fixture: &[u8]) -> Vec { + let mut encoded = Vec::new(); + rmp::encode::write_array_len(&mut encoded, 5).expect("FileInfoVersions array header should encode"); + "bucket" + .serialize(&mut Serializer::new(&mut encoded)) + .expect("FileInfoVersions volume should encode"); + "object" + .serialize(&mut Serializer::new(&mut encoded)) + .expect("FileInfoVersions name should encode"); + Option::::None + .serialize(&mut Serializer::new(&mut encoded)) + .expect("FileInfoVersions mod time should encode"); + rmp::encode::write_array_len(&mut encoded, 1).expect("FileInfoVersions versions header should encode"); + encoded.extend_from_slice(fixture); + rmp::encode::write_array_len(&mut encoded, 0).expect("FileInfoVersions free versions header should encode"); + encoded + } + + #[test] + fn fileinfo_versions_decodes_nested_historical_layouts() { + for fixture in [BETA11_FILEINFO_FIXTURE, BETA12_FILEINFO_FIXTURE] { + let encoded = wrap_historical_file_info(fixture); + let decoded: FileInfoVersions = + rmp_serde::from_slice(&encoded).expect("nested historical FileInfo should decode through FileInfoVersions"); + assert_eq!(decoded.versions.len(), 1); + assert_eq!(decoded.versions[0].name, "object"); + assert!(decoded.free_versions.is_empty()); + } + } + + #[derive(Deserialize)] + struct Beta11MapProbe { + volume: String, + expire_restored: bool, + idx: usize, + } + + #[derive(Deserialize)] + struct Beta12MapProbe { + transition_version: Option, + transition_version_state: TransitionVersionState, + expire_restored: bool, + } + + #[derive(Deserialize)] + struct Beta11NestedMapProbe { + file_info: Beta11MapProbe, + } + + #[derive(Serialize)] + struct NestedFileInfo<'a> { + file_info: &'a FileInfo, + } + + #[test] + fn fileinfo_serializes_as_map_readable_by_beta11_and_beta12_shapes() { + let expected = positional_compat_file_info(); + let encoded = expected.marshal_msg().expect("current FileInfo map should encode"); + let mut cursor = encoded.as_slice(); + let field_count = usize::try_from(rmp::decode::read_map_len(&mut cursor).expect("FileInfo should start with a map")) + .expect("FileInfo map field count should fit usize"); + assert_eq!(field_count, FILE_INFO_FIELDS.len()); + + let beta11: Beta11MapProbe = rmp_serde::from_slice(&encoded).expect("beta.11 field shape should read current map"); + assert_eq!(beta11.volume, expected.volume); + assert_eq!(beta11.expire_restored, expected.expire_restored); + assert_eq!(beta11.idx, expected.idx); + + let beta12: Beta12MapProbe = rmp_serde::from_slice(&encoded).expect("beta.12 field shape should read current map"); + assert_eq!(beta12.transition_version, expected.transition_version); + assert_eq!(beta12.transition_version_state, expected.transition_version_state); + assert_eq!(beta12.expire_restored, expected.expire_restored); + + let nested = NestedFileInfo { file_info: &expected }; + let nested_encoded = rmp_serde::to_vec(&nested).expect("nested FileInfo should encode"); + let nested_beta11: Beta11NestedMapProbe = + rmp_serde::from_slice(&nested_encoded).expect("beta.11 field shape should read nested current map"); + assert_eq!(nested_beta11.file_info.volume, expected.volume); + assert_eq!(nested_beta11.file_info.expire_restored, expected.expire_restored); + } + + #[test] + fn fileinfo_deserializes_beta11_map_without_transition_fields() { + let expected = positional_compat_file_info(); + let mut value = serde_json::to_value(&expected).expect("current FileInfo should serialize to JSON"); + let object = value.as_object_mut().expect("FileInfo JSON should be an object"); + object.remove("transition_version"); + object.remove("transition_version_state"); + + let decoded: FileInfo = serde_json::from_value(value).expect("beta.11 FileInfo JSON should decode"); + assert_eq!(decoded.volume, expected.volume); + assert_eq!(decoded.expire_restored, expected.expire_restored); + assert_eq!(decoded.transition_version, None); + assert_eq!(decoded.transition_version_state, TransitionVersionState::Unknown); + } + proptest! { #[test] fn fileinfo_msgpack_round_trips(value in file_info_strategy()) { diff --git a/docs/architecture/compat-cleanup-register.md b/docs/architecture/compat-cleanup-register.md index 776773cb2..db8a8fcac 100644 --- a/docs/architecture/compat-cleanup-register.md +++ b/docs/architecture/compat-cleanup-register.md @@ -12,6 +12,7 @@ for later deletion. ## Open Items +- `rustfs-5509` FileInfo positional MessagePack decoding: beta.11 serialized 28 fields, while beta.12 inserted transition-version fields in the middle and serialized an incompatible 30-field array. New releases write named maps and retain readers for both shipped array layouts so direct and rolling upgrades can read either release. Remove the positional-array readers after every supported direct-upgrade release writes named maps and no retained RPC payload can contain a pre-map FileInfo array. - `rustfs-5416` Helm distributed startup wait setting: charts that predate explicit local endpoint identity expose startupWaitTimeoutSeconds for their peer DNS/TCP init gate. The new chart keeps the value accepted but ignores it after moving startup convergence into RustFS. Remove the value and its documentation after the minimum supported direct-upgrade chart includes localEndpointHost.autoInject and no longer renders the peer gate. - `rustfs-5416-wait-mode` startup wait-mode validation: releases before explicit local endpoint identity treat an unknown RUSTFS_STARTUP_TOPOLOGY_WAIT_MODE value as auto. New servers retain that fallback only when no explicit local endpoint host is configured; an anchor requires a recognized mode so a typo cannot bypass DNS locality. Remove the fallback and reject every unknown value after every supported direct-upgrade chart validates this setting before rollout. - `rustfs-5416-kubernetes-alias-dns` Kubernetes endpoint identity fallback: deployments created before explicit local endpoint identity may use resolvable aliases that do not match the Pod hostname. An implicit auto-mode zero match retains legacy DNS locality with a bounded deadline, while ambiguous matches and invalid explicit anchors still fail closed. Remove the fallback after every supported direct-upgrade chart and deployment manifest provides a canonical RUSTFS_LOCAL_ENDPOINT_HOST for domain-based distributed topologies. diff --git a/docs/architecture/minio-rustfs-router-compatibility.md b/docs/architecture/minio-rustfs-router-compatibility.md index bf8971a24..00c9ff68c 100644 --- a/docs/architecture/minio-rustfs-router-compatibility.md +++ b/docs/architecture/minio-rustfs-router-compatibility.md @@ -157,7 +157,7 @@ maps MinIO admin route families to RustFS status. | TLS debug (`/debug/tls/status`) | 已实现 | `rustfs/src/admin/handlers/tls_debug.rs`, `profile.rs` | | KMS management / dynamic / keys | 已实现 | `rustfs/src/admin/handlers/kms_management.rs`, `kms_dynamic.rs`, `kms_keys.rs` | | OIDC public + config | 已实现 | `rustfs/src/admin/handlers/oidc.rs` | -| Table catalog (Iceberg) | 已实现 | `rustfs/src/admin/handlers/table_catalog.rs` | +| Table catalog (Iceberg) | 已实现 | `rustfs/src/admin/handlers/table_catalog/mod.rs` | ### Registered-but-incomplete diff --git a/rustfs/src/storage/rpc/node_service/disk.rs b/rustfs/src/storage/rpc/node_service/disk.rs index d5fa44815..31d63bd97 100644 --- a/rustfs/src/storage/rpc/node_service/disk.rs +++ b/rustfs/src/storage/rpc/node_service/disk.rs @@ -35,9 +35,9 @@ use tokio_util::time::DelayQueue; use tonic::{Request, Response, Status}; use tracing::debug; -/// Initial capacity hint (bytes) for msgpack encode buffers, sized to cover a typical single- -/// version `FileInfo` without repeated growth reallocations. Larger payloads still grow as needed. +/// Initial capacity hint (bytes) for typical small msgpack requests and responses. const MSGPACK_ENCODE_CAPACITY_HINT: usize = 512; +const FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT: usize = 1024; const SNAPSHOT_LEASE_PROTOCOL_VERSION: u32 = 1; const SNAPSHOT_LEASE_MIN_TTL: Duration = Duration::from_secs(5); const SNAPSHOT_LEASE_MAX_TTL: Duration = Duration::from_secs(5 * 60); @@ -170,14 +170,26 @@ fn decode_msgpack_or_json( } } -fn encode_msgpack(value: &T, value_name: &str) -> std::result::Result, DiskError> { - let mut serializer = rmp_serde::Serializer::new(Vec::with_capacity(MSGPACK_ENCODE_CAPACITY_HINT)); +fn encode_msgpack_with_capacity( + value: &T, + value_name: &str, + capacity: usize, +) -> std::result::Result, DiskError> { + let mut serializer = rmp_serde::Serializer::new(Vec::with_capacity(capacity)); value .serialize(&mut serializer) .map_err(|err| DiskError::other(format!("encode {value_name} msgpack failed: {err}")))?; Ok(serializer.into_inner()) } +fn encode_msgpack(value: &T, value_name: &str) -> std::result::Result, DiskError> { + encode_msgpack_with_capacity(value, value_name, MSGPACK_ENCODE_CAPACITY_HINT) +} + +fn encode_file_info_msgpack(value: &FileInfo) -> std::result::Result, DiskError> { + encode_msgpack_with_capacity(value, "FileInfo", FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT) +} + fn encode_msgpack_named(value: &T, value_name: &str) -> std::result::Result, DiskError> { let mut serializer = rmp_serde::Serializer::new(Vec::with_capacity(MSGPACK_ENCODE_CAPACITY_HINT)).with_struct_map(); value @@ -242,7 +254,11 @@ fn encode_batch_read_version_response_payloads( compat_response_json(batch_read_version_resp) .map_err(|err| DiskError::other(format!("encode BatchReadVersionResp json failed: {err}")))?, ); - batch_read_version_resps_bin.push(Bytes::from(encode_msgpack(batch_read_version_resp, "BatchReadVersionResp")?)); + batch_read_version_resps_bin.push(Bytes::from(encode_msgpack_with_capacity( + batch_read_version_resp, + "BatchReadVersionResp", + FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT, + )?)); } Ok((batch_read_version_resps_json, batch_read_version_resps_bin)) @@ -779,7 +795,7 @@ impl NodeService { { Ok(file_info) => { let file_info_json = compat_response_json(&file_info); - let file_info_bin = encode_msgpack(&file_info, "FileInfo"); + let file_info_bin = encode_file_info_msgpack(&file_info); match (file_info_json, file_info_bin) { (Ok(file_info), Ok(file_info_bin)) => Ok(Response::new(ReadVersionResponse { success: true, diff --git a/scripts/check_logging_guardrails.sh b/scripts/check_logging_guardrails.sh index 12fd1fa4d..e1f00af66 100755 --- a/scripts/check_logging_guardrails.sh +++ b/scripts/check_logging_guardrails.sh @@ -12,7 +12,15 @@ checked_files=( "rustfs/src/auth.rs" "rustfs/src/protocols/client.rs" "rustfs/src/admin/router.rs" - "rustfs/src/admin/handlers/table_catalog.rs" + "rustfs/src/admin/handlers/table_catalog/config.rs" + "rustfs/src/admin/handlers/table_catalog/credentials.rs" + "rustfs/src/admin/handlers/table_catalog/maintenance.rs" + "rustfs/src/admin/handlers/table_catalog/mod.rs" + "rustfs/src/admin/handlers/table_catalog/namespace.rs" + "rustfs/src/admin/handlers/table_catalog/refs.rs" + "rustfs/src/admin/handlers/table_catalog/routes.rs" + "rustfs/src/admin/handlers/table_catalog/table.rs" + "rustfs/src/admin/handlers/table_catalog/view.rs" "rustfs/src/admin/handlers/service_account.rs" "rustfs/src/admin/handlers/kms_audit.rs" "rustfs/src/admin/handlers/kms_dynamic.rs"