From 837286f95942934250ccef9d9f9c402521571200 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=A9=AC=E7=99=BB=E5=B1=B1?= Date: Tue, 28 Jul 2026 13:37:43 +0800 Subject: [PATCH] feat(tiering): persist opaque remote versions --- crates/ecstore/src/object_api/types.rs | 10 +- crates/ecstore/src/services/tier/test_util.rs | 5 +- .../src/set_disk/core/io_primitives.rs | 1 + crates/ecstore/src/set_disk/metadata.rs | 1 + crates/ecstore/src/set_disk/ops/object.rs | 67 +++++----- crates/filemeta/src/fileinfo.rs | 23 ++++ crates/filemeta/src/filemeta/version.rs | 125 ++++++++++++------ 7 files changed, 158 insertions(+), 74 deletions(-) diff --git a/crates/ecstore/src/object_api/types.rs b/crates/ecstore/src/object_api/types.rs index c9a0e560b..33f53347b 100644 --- a/crates/ecstore/src/object_api/types.rs +++ b/crates/ecstore/src/object_api/types.rs @@ -464,11 +464,11 @@ impl ObjectInfo { let transitioned_object = TransitionedObject { name: fi.transitioned_objname.clone(), - version_id: if let Some(transition_version_id) = fi.transition_version_id { - transition_version_id.to_string() - } else { - "".to_string() - }, + version_id: fi + .transition_version + .clone() + .or_else(|| fi.transition_version_id.map(|version_id| version_id.to_string())) + .unwrap_or_default(), status: fi.transition_status.clone(), free_version: fi.tier_free_version(), tier: fi.transition_tier.clone(), diff --git a/crates/ecstore/src/services/tier/test_util.rs b/crates/ecstore/src/services/tier/test_util.rs index 1a2b37d26..751f2e92b 100644 --- a/crates/ecstore/src/services/tier/test_util.rs +++ b/crates/ecstore/src/services/tier/test_util.rs @@ -931,7 +931,10 @@ pub async fn read_transition_meta(disk_path: &Path, bucket: &str, object: &str) status: fi.transition_status.clone(), tier: fi.transition_tier.clone(), remote_object: fi.transitioned_objname.clone(), - remote_version_id: fi.transition_version_id.map(|id| id.to_string()), + remote_version_id: fi + .transition_version + .clone() + .or_else(|| fi.transition_version_id.map(|id| id.to_string())), free_version_count, }) } diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 9d07d8310..ceaff1996 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -544,6 +544,7 @@ pub(in crate::set_disk) fn metadata_early_stop_candidate_matches(left: &FileInfo && left.transitioned_objname == right.transitioned_objname && left.transition_tier == right.transition_tier && left.transition_version_id == right.transition_version_id + && left.transition_version == right.transition_version && left.expire_restored == right.expire_restored && left.size == right.size && left.mod_time == right.mod_time diff --git a/crates/ecstore/src/set_disk/metadata.rs b/crates/ecstore/src/set_disk/metadata.rs index 136a6cf7e..419055f7a 100644 --- a/crates/ecstore/src/set_disk/metadata.rs +++ b/crates/ecstore/src/set_disk/metadata.rs @@ -578,6 +578,7 @@ impl SetDisks { Self::update_hash_str(hasher, &meta.transition_tier); Self::update_hash_str(hasher, &meta.transitioned_objname); Self::update_hash_optional_uuid(hasher, meta.transition_version_id); + Self::update_hash_optional_str(hasher, meta.transition_version.as_deref()); Self::update_hash_optional_u32(hasher, meta.mode); Self::update_hash_optional_u64(hasher, meta.written_by_version); diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 757af70b4..3adf2ffa8 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -2107,11 +2107,12 @@ async fn pause_transition_commit(bucket: &str, object: &str, pause: TransitionCo } } -fn parse_transition_version_id(remote_version: &str) -> std::result::Result, uuid::Error> { - if remote_version.is_empty() { - return Ok(None); +fn parse_transition_version_id(remote_version: &str) -> Option { + if remote_version.is_empty() || Uuid::parse_str(remote_version).is_ok_and(|version_id| version_id.is_nil()) { + None + } else { + Some(remote_version.to_string()) } - Uuid::parse_str(remote_version).map(|version_id| (!version_id.is_nil()).then_some(version_id)) } #[cfg(test)] @@ -2278,11 +2279,8 @@ mod transition_version_id_tests { #[test] fn normalizes_persisted_unversioned_ids_and_preserves_put_constraints() { - assert_eq!(parse_transition_version_id("").expect("empty remote version should be valid"), None); - assert_eq!( - parse_transition_version_id(&Uuid::nil().to_string()).expect("nil remote version should be valid"), - None - ); + assert_eq!(parse_transition_version_id(""), None); + assert_eq!(parse_transition_version_id(&Uuid::nil().to_string()), None); let nil_put_response = Uuid::nil().to_string(); let nil_candidate = TransitionUploadCandidate::from_put_response(nil_put_response.clone()); assert_eq!(nil_candidate.cleanup_version(), nil_put_response); @@ -2294,11 +2292,12 @@ mod transition_version_id_tests { } #[test] - fn preserves_valid_remote_id_and_rejects_invalid_text() { + fn preserves_uuid_and_opaque_remote_ids() { let version_id = Uuid::new_v4(); + assert_eq!(parse_transition_version_id(&version_id.to_string()), Some(version_id.to_string())); assert_eq!( - parse_transition_version_id(&version_id.to_string()).expect("UUID remote version should be valid"), - Some(version_id) + parse_transition_version_id("opaque-version-token"), + Some("opaque-version-token".to_string()) ); assert_eq!( TransitionUploadCandidate::from_put_response(version_id.to_string()).cleanup_version(), @@ -2308,7 +2307,6 @@ mod transition_version_id_tests { TransitionUploadCandidate::from_put_response("opaque-version-token".to_string()).cleanup_version(), "opaque-version-token" ); - assert!(parse_transition_version_id("not-a-uuid").is_err()); } } @@ -3554,16 +3552,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object).await; return Err(err); } - let transition_version_id = match parse_transition_version_id(candidate.remote_version()) { - Ok(version_id) => version_id, - Err(err) => { - if upload_cleanup.cleanup().await.is_ok() { - delete_transition_transaction_after_remote_cleanup(transaction_api.as_ref(), transaction_id, bucket, object) - .await; - } - return Err(err.into()); - } - }; + let transition_version_id = parse_transition_version_id(candidate.remote_version()); let mut commit_opts = opts.clone(); commit_opts.no_lock = true; @@ -3621,7 +3610,10 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { current_fi.transition_status = TRANSITION_COMPLETE.to_string(); current_fi.transitioned_objname = dest_obj; current_fi.transition_tier = opts.transition.tier.clone(); - current_fi.transition_version_id = transition_version_id; + current_fi.transition_version_id = transition_version_id + .as_deref() + .and_then(|version_id| Uuid::parse_str(version_id).ok()); + current_fi.transition_version = transition_version_id; rustfs_utils::http::metadata_compat::insert_str( &mut current_fi.metadata, rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, @@ -6242,7 +6234,7 @@ mod transition_upload_integrity_tests { #[tokio::test] #[serial_test::serial] - async fn opaque_remote_version_is_cleaned_before_parse_failure() { + async fn opaque_remote_version_is_persisted_exactly() { let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await; let bucket = "transition-unknown-version-bucket"; let object = "object.bin"; @@ -6255,12 +6247,25 @@ mod transition_upload_integrity_tests { set_disks .transition_object(bucket, object, &transition_options(&original, tier_name)) .await - .expect_err("an unparseable remote version must fail closed"); - let removed_versions = backend.remove_versions().await; - assert_eq!(removed_versions.len(), 1); - assert_eq!(removed_versions[0].1, "opaque-version-token"); - assert_eq!(backend.object_count().await, 0); - assert_local_source_intact(&set_disks, bucket, object, &payload).await; + .expect("an accepted opaque remote version must commit"); + let (fi, _, _) = set_disks + .get_object_fileinfo( + bucket, + object, + &ObjectOptions { + no_lock: true, + metadata_cache_safe: false, + ..Default::default() + }, + true, + false, + ) + .await + .expect("committed opaque transition metadata should be readable"); + assert_eq!(fi.transition_version_id, None); + assert_eq!(fi.transition_version.as_deref(), Some("opaque-version-token")); + assert_eq!(backend.remove_count().await, 0); + assert_eq!(backend.object_count().await, 1); } #[tokio::test] diff --git a/crates/filemeta/src/fileinfo.rs b/crates/filemeta/src/fileinfo.rs index 98a39c8a9..012adccb2 100644 --- a/crates/filemeta/src/fileinfo.rs +++ b/crates/filemeta/src/fileinfo.rs @@ -230,6 +230,8 @@ pub struct FileInfo { pub transitioned_objname: String, pub transition_tier: String, pub transition_version_id: Option, + #[serde(default)] + pub transition_version: Option, pub expire_restored: bool, pub data_dir: Option, pub mod_time: Option, @@ -455,6 +457,10 @@ impl FileInfo { if self.mod_time.is_none_or(|mod_time| mod_time <= OffsetDateTime::UNIX_EPOCH) || (!allow_nil_version_id && self.version_id.is_some_and(|version_id| version_id.is_nil())) || self.transition_version_id.is_some_and(|version_id| version_id.is_nil()) + || self + .transition_version + .as_ref() + .is_some_and(|version_id| version_id.is_empty()) || self.size != 0 || self.data_dir.is_some() || self.mode.is_some() @@ -488,6 +494,7 @@ impl FileInfo { || !self.transitioned_objname.is_empty() || !self.transition_tier.is_empty() || self.transition_version_id.is_some() + || self.transition_version.is_some() || self.expire_restored || self.size != 0 || self.data_dir.is_some() @@ -532,6 +539,11 @@ impl FileInfo { /// return `None`. pub fn validate(&self, mode: ValidationMode) -> Result> { self.validate_collection_bounds()?; + if let (Some(version), Some(version_id)) = (&self.transition_version, self.transition_version_id) + && Uuid::parse_str(version).ok() != Some(version_id) + { + return Err(Error::FileCorrupt); + } let erasure_layout = match mode { ValidationMode::RequireErasure => Some(self.validate_erasure_geometry()?), @@ -824,6 +836,7 @@ impl FileInfo { && self.transition_tier == other.transition_tier && self.transitioned_objname == other.transitioned_objname && self.transition_version_id == other.transition_version_id + && self.transition_version == other.transition_version } /// Check if metadata maps are equal @@ -1328,6 +1341,15 @@ mod tests { assert_file_corrupt(&fi, ValidationMode::DeleteOnly); } + #[test] + fn metadata_read_validation_rejects_conflicting_transition_versions() { + let mut fi = one_shard_validation_fileinfo(1); + fi.transition_version_id = Some(Uuid::new_v4()); + fi.transition_version = Some(Uuid::new_v4().to_string()); + + assert_file_corrupt(&fi, ValidationMode::RequireErasure); + } + #[test] fn metadata_read_validation_requires_canonical_delete_marker_shape() { let marker = FileInfo { @@ -1699,6 +1721,7 @@ mod tests { transitioned_objname, transition_tier, transition_version_id, + transition_version: transition_version_id.map(|version_id| version_id.to_string()), expire_restored, data_dir, mod_time, diff --git a/crates/filemeta/src/filemeta/version.rs b/crates/filemeta/src/filemeta/version.rs index f491e75b9..fa3cf6c0d 100644 --- a/crates/filemeta/src/filemeta/version.rs +++ b/crates/filemeta/src/filemeta/version.rs @@ -43,6 +43,7 @@ const MSGPACK_FIXEXT8: u8 = 0xd7; const MSGPACK_TIME_EXT_LEGACY: i8 = 5; const MSGPACK_TIME_EXT_OFFICIAL: i8 = -1; const MSGPACK_TIME_LEN: u8 = 12; +const MAX_TRANSITION_VERSION_LEN: usize = 1024; /// Sentinel signature returned when a version has no computable body (invalid / /// missing inner object). Mirrors MinIO's `signatureErr` so such versions never @@ -251,23 +252,38 @@ fn parse_legacy_uuid_bytes(bytes: &[u8], field: &str) -> Result> { /// Decode a stored transitioned-version-id from a version's `meta_sys`. /// -/// RustFS writes it as 16 raw UUID bytes; MinIO-migrated tiered objects store -/// the remote tier's version id as a UUID *string*. Accept both, and treat any -/// absent / nil / otherwise-unparseable value as "no tier version" (matching the -/// tolerant pre-hardening behavior) rather than failing the whole object read — -/// a malformed tier id must not make an otherwise-readable object unreadable. -fn transitioned_version_id_from_meta_sys(meta_sys: &HashMap>) -> Option { +/// Legacy RustFS writes used 16 raw UUID bytes. New writes and MinIO-migrated +/// records use the provider's exact UTF-8 version text. Empty, nil UUID, and +/// malformed bytes are not usable remote versions. +fn transitioned_version_from_meta_sys(meta_sys: &HashMap>) -> Option { let value = get_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_ID)?; if value.is_empty() { return None; } if let Ok(id) = Uuid::from_slice(&value) { - return (!id.is_nil()).then_some(id); + return (!id.is_nil()).then(|| id.to_string()); } - std::str::from_utf8(&value) - .ok() - .and_then(|s| Uuid::parse_str(s.trim()).ok()) - .filter(|id| !id.is_nil()) + let value = String::from_utf8(value).ok()?; + if value.is_empty() + || value.len() > MAX_TRANSITION_VERSION_LEN + || value.chars().any(char::is_control) + || Uuid::parse_str(&value).is_ok_and(|id| id.is_nil()) + { + None + } else { + Some(value) + } +} + +fn legacy_transitioned_version_id_from_meta_sys(meta_sys: &HashMap>) -> Option { + transitioned_version_from_meta_sys(meta_sys).and_then(|value| Uuid::parse_str(&value).ok()) +} + +fn transitioned_version_bytes(fi: &FileInfo) -> Option> { + fi.transition_version + .as_ref() + .map(|version| version.as_bytes().to_vec()) + .or_else(|| fi.transition_version_id.map(|version_id| version_id.as_bytes().to_vec())) } fn parse_legacy_erasure_algo(value: &str) -> ErasureAlgo { @@ -2398,7 +2414,8 @@ impl MetaObject { let transitioned_objname = get_bytes(&self.meta_sys, SUFFIX_TRANSITIONED_OBJECTNAME) .map(|v| String::from_utf8_lossy(&v).to_string()) .unwrap_or_default(); - let transition_version_id = transitioned_version_id_from_meta_sys(&self.meta_sys); + let transition_version = transitioned_version_from_meta_sys(&self.meta_sys); + let transition_version_id = transition_version.as_deref().and_then(|value| Uuid::parse_str(value).ok()); let transition_tier = get_bytes(&self.meta_sys, SUFFIX_TRANSITION_TIER) .map(|v| String::from_utf8_lossy(&v).to_string()) .unwrap_or_default(); @@ -2419,6 +2436,7 @@ impl MetaObject { transition_status, transitioned_objname, transition_version_id, + transition_version, transition_tier, ..Default::default() }) @@ -2431,12 +2449,8 @@ impl MetaObject { SUFFIX_TRANSITIONED_OBJECTNAME, fi.transitioned_objname.as_bytes().to_vec(), ); - if let Some(transition_version_id) = fi.transition_version_id.as_ref() { - insert_bytes( - &mut self.meta_sys, - SUFFIX_TRANSITIONED_VERSION_ID, - transition_version_id.as_bytes().to_vec(), - ); + if let Some(transition_version) = transitioned_version_bytes(fi) { + insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version); } insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER, fi.transition_tier.as_bytes().to_vec()); if let Some(destination_id) = get_str(&fi.metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID) { @@ -2562,8 +2576,8 @@ impl From for MetaObject { ); } - if let Some(vid) = &value.transition_version_id { - insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, vid.as_bytes().to_vec()); + if let Some(transition_version) = transitioned_version_bytes(&value) { + insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version); } if !value.transition_tier.is_empty() { @@ -2706,7 +2720,8 @@ impl MetaDeleteMarker { .map(|v| String::from_utf8_lossy(&v).to_string()) .unwrap_or_default(); - fi.transition_version_id = transitioned_version_id_from_meta_sys(&self.meta_sys); + fi.transition_version = transitioned_version_from_meta_sys(&self.meta_sys); + fi.transition_version_id = legacy_transitioned_version_id_from_meta_sys(&self.meta_sys); } fi @@ -2859,8 +2874,8 @@ impl From for MetaDeleteMarker { value.transitioned_objname.as_bytes().to_vec(), ); } - if let Some(version_id) = value.transition_version_id { - insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, version_id.as_bytes().to_vec()); + if let Some(transition_version) = transitioned_version_bytes(&value) { + insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version); } if !value.transition_tier.is_empty() { insert_bytes(&mut meta_sys, SUFFIX_TRANSITION_TIER, value.transition_tier.as_bytes().to_vec()); @@ -3412,7 +3427,7 @@ mod tests { .insert("x-rustfs-internal-healing".to_string(), "true".to_string()); marker.metadata.insert("content-type".to_string(), "text/plain".to_string()); let remote_version_id = Uuid::new_v4(); - marker.transition_version_id = Some(remote_version_id); + marker.transition_version = Some(remote_version_id.to_string()); let converted = MetaDeleteMarker::from(marker); @@ -3420,7 +3435,19 @@ mod tests { assert_eq!(converted.meta_sys.get("x-minio-internal-purgestatus"), Some(&b"pending".to_vec())); assert_eq!( get_bytes(&converted.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID), - Some(remote_version_id.as_bytes().to_vec()) + Some(remote_version_id.to_string().into_bytes()) + ); + assert_eq!( + converted + .meta_sys + .get(&format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_TRANSITIONED_VERSION_ID}")), + Some(&remote_version_id.to_string().into_bytes()) + ); + assert_eq!( + converted + .meta_sys + .get(&format!("{}{SUFFIX_TRANSITIONED_VERSION_ID}", rustfs_utils::http::MINIO_INTERNAL_PREFIX)), + Some(&remote_version_id.to_string().into_bytes()) ); assert!(!converted.meta_sys.contains_key("x-rustfs-internal-healing")); assert!(!converted.meta_sys.contains_key("content-type")); @@ -4097,19 +4124,42 @@ mod tests { .into_fileinfo("b", "k", false) .expect("into_fileinfo"); assert_eq!(fi.transition_version_id, Some(id)); + assert_eq!(fi.transition_version, Some(id.to_string())); } #[test] - fn meta_object_transition_version_id_unparseable_stays_readable_as_none() { - // A non-UUID / non-16-byte tier version id must NOT make the object - // unreadable; it is tolerated as "no tier version" (compat with - // pre-hardening behavior and foreign/edge metadata). + fn meta_object_transition_version_id_opaque_text_is_preserved() { let mut sys = HashMap::new(); - insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, b"not-a-uuid".to_vec()); + insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, b"opaque-generation-42".to_vec()); let fi = make_meta_object_with_sys(sys) .into_fileinfo("b", "k", false) - .expect("unparseable transition version id must not fail the object read"); + .expect("opaque transition version id must decode"); assert_eq!(fi.transition_version_id, None); + assert_eq!(fi.transition_version.as_deref(), Some("opaque-generation-42")); + } + + #[test] + fn meta_object_transition_version_id_invalid_utf8_yields_none() { + let mut sys = HashMap::new(); + insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, vec![0xff]); + let fi = make_meta_object_with_sys(sys) + .into_fileinfo("b", "k", false) + .expect("invalid transition version bytes must not fail the object read"); + assert_eq!(fi.transition_version_id, None); + assert_eq!(fi.transition_version, None); + } + + #[test] + fn meta_object_transition_version_id_unsafe_text_yields_none() { + for value in [b"opaque\0version".to_vec(), vec![b'x'; MAX_TRANSITION_VERSION_LEN + 1]] { + let mut sys = HashMap::new(); + insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, value); + let fi = make_meta_object_with_sys(sys) + .into_fileinfo("b", "k", false) + .expect("unsafe transition version text must not fail the object read"); + assert_eq!(fi.transition_version_id, None); + assert_eq!(fi.transition_version, None); + } } #[test] @@ -4123,6 +4173,7 @@ mod tests { .into_fileinfo("b", "k", false) .expect("string-form transition version id must decode"); assert_eq!(fi.transition_version_id, Some(id)); + assert_eq!(fi.transition_version, Some(id.to_string())); } #[test] @@ -4152,16 +4203,14 @@ mod tests { } .into_fileinfo("b", "k", false); assert_eq!(fi.transition_version_id, Some(id)); + assert_eq!(fi.transition_version, Some(id.to_string())); } #[test] - fn delete_marker_free_version_transition_version_id_unparseable_stays_readable() { - // A malformed tier version id must not make a free-version record corrupt: - // it decodes to None and stays readable. Otherwise free-version expiry - // fails and the remote-tier object leaks. + fn delete_marker_free_version_transition_version_id_opaque_text_is_preserved() { let mut sys = HashMap::new(); insert_bytes(&mut sys, SUFFIX_FREE_VERSION, vec![]); - insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, b"not-a-uuid".to_vec()); + insert_bytes(&mut sys, SUFFIX_TRANSITIONED_VERSION_ID, b"opaque-generation-42".to_vec()); insert_bytes(&mut sys, SUFFIX_TRANSITION_TIER, b"WARM".to_vec()); insert_bytes(&mut sys, SUFFIX_TRANSITIONED_OBJECTNAME, b"remote-object".to_vec()); let fi = MetaDeleteMarker { @@ -4172,8 +4221,9 @@ mod tests { .into_fileinfo("b", "k", false); assert_eq!(fi.transition_version_id, None); + assert_eq!(fi.transition_version.as_deref(), Some("opaque-generation-42")); fi.validate_for_metadata_read() - .expect("free-version record with an unparseable tier id must remain readable"); + .expect("free-version record with an opaque tier id must remain readable"); } #[test] @@ -4193,6 +4243,7 @@ mod tests { .into_fileinfo("b", "k", false); assert_eq!(fi.transition_version_id, Some(id)); + assert_eq!(fi.transition_version, Some(id.to_string())); } #[test]