From 4dbc58887afe4c0e38dc11704f403479a9f7c79e Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 5 Sep 2026 14:00:14 +0800 Subject: [PATCH] fix(tier): probe legacy transition version state (#7138) --- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 782 +++++++++++++++++- crates/ecstore/src/services/tier/test_util.rs | 5 +- crates/ecstore/src/services/tier/tier.rs | 13 + .../ecstore/src/services/tier/warm_backend.rs | 53 +- .../src/services/tier/warm_backend_s3.rs | 43 +- crates/ecstore/src/store/init.rs | 160 +++- crates/filemeta/src/filemeta/version.rs | 125 ++- 7 files changed, 1133 insertions(+), 48 deletions(-) diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 7649e3ac9..dbb5dd154 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -584,33 +584,173 @@ impl ExpiryOp for FreeVersionTask { } } +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum TransitionDeleteVersionPlan { + Direct { version_id_exact: bool }, + ProbeLegacyUnknown, +} + +fn legacy_transition_version_state_missing(oi: &ObjectInfo) -> Result { + use rustfs_utils::http::metadata_compat::{ + SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, contains_key_str, get_consistent_str, + }; + + if !contains_key_str(&oi.user_defined, SUFFIX_TRANSITIONED_VERSION_STATE) { + let version_key_present = contains_key_str(&oi.user_defined, SUFFIX_TRANSITIONED_VERSION_ID); + if version_key_present { + if oi.transitioned_object.version_id.is_empty() { + let has_non_empty_version = oi.user_defined.iter().any(|(key, value)| { + rustfs_utils::http::metadata_compat::strip_internal_prefix_preserving_case(key) + .is_some_and(|suffix| suffix.eq_ignore_ascii_case(SUFFIX_TRANSITIONED_VERSION_ID)) + && !value.is_empty() + }); + if !has_non_empty_version { + // MinIO writes the transitioned-versionID key with an empty value + // for unversioned tier objects. The backend probe remains the proof. + return Ok(true); + } + } else if get_consistent_str(&oi.user_defined, SUFFIX_TRANSITIONED_VERSION_ID) + == Some(oi.transitioned_object.version_id.as_str()) + { + return Ok(true); + } + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "legacy remote tier version metadata is conflicting or malformed", + )); + } + if !oi.transitioned_object.version_id.is_empty() { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "legacy remote tier version metadata is missing or inconsistent", + )); + } + return Ok(true); + } + let persisted = get_consistent_str(&oi.user_defined, SUFFIX_TRANSITIONED_VERSION_STATE).ok_or_else(|| { + std::io::Error::new( + std::io::ErrorKind::InvalidData, + "remote tier object has conflicting transition version state metadata", + ) + })?; + if persisted != oi.transition_version_state.as_str() { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "remote tier object transition version state metadata changed during decoding", + )); + } + Ok(false) +} + +fn transition_remote_version_delete_plan(oi: &ObjectInfo) -> Result { + match oi.transition_version_state { + rustfs_filemeta::TransitionVersionState::Unknown => { + if legacy_transition_version_state_missing(oi)? { + Ok(TransitionDeleteVersionPlan::ProbeLegacyUnknown) + } else { + validate_transition_remote_version(oi) + .map(|version_id_exact| TransitionDeleteVersionPlan::Direct { version_id_exact }) + } + } + _ => validate_transition_remote_version(oi) + .map(|version_id_exact| TransitionDeleteVersionPlan::Direct { version_id_exact }), + } +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +struct ResolvedTransitionDeleteVersion { + version_id_exact: bool, + remote_already_missing: bool, +} + async fn acquire_free_version_tier_lease( oi: &ObjectInfo, tier_config_mgr: &Arc>, -) -> Result<(TierOperationLease, bool), std::io::Error> { - let version_id_exact = validate_transition_remote_version(oi)?; +) -> Result<(TierOperationLease, TransitionDeleteVersionPlan), std::io::Error> { + let delete_plan = transition_remote_version_delete_plan(oi)?; let identity = tier_destination_id_from_metadata(&oi.user_defined)? .ok_or_else(|| std::io::Error::other("tier free-version has no durable backend identity"))?; let lease = TierConfigMgr::acquire_operation_lease_for_backend_identity(tier_config_mgr, &oi.transitioned_object.tier, identity) .await .map_err(std::io::Error::other)?; - Ok((lease, version_id_exact)) + Ok((lease, delete_plan)) +} + +async fn resolve_transition_delete_version_plan( + oi: &ObjectInfo, + lease: &TierOperationLease, + delete_plan: TransitionDeleteVersionPlan, +) -> Result { + match delete_plan { + TransitionDeleteVersionPlan::Direct { version_id_exact } => Ok(ResolvedTransitionDeleteVersion { + version_id_exact, + remote_already_missing: false, + }), + TransitionDeleteVersionPlan::ProbeLegacyUnknown => { + let expected_version = oi.transitioned_object.version_id.as_str(); + if expected_version.is_empty() { + return Err(std::io::Error::new( + std::io::ErrorKind::WouldBlock, + "remote tier cannot safely delete a legacy object without an exact version ID", + )); + } + let probe = lease + .probe_transition_version(&oi.transitioned_object.name, expected_version) + .await?; + match (expected_version, probe) { + (expected, crate::services::tier::warm_backend::TransitionCandidateProbe::VersionedPresent(actual)) + if expected == actual => + { + lease.validate_remote_version_id(expected)?; + Ok(ResolvedTransitionDeleteVersion { + version_id_exact: true, + remote_already_missing: false, + }) + } + (_, crate::services::tier::warm_backend::TransitionCandidateProbe::Missing) => { + Ok(ResolvedTransitionDeleteVersion { + version_id_exact: false, + remote_already_missing: true, + }) + } + (_, crate::services::tier::warm_backend::TransitionCandidateProbe::Unsupported) => Err(std::io::Error::new( + std::io::ErrorKind::Unsupported, + "remote tier cannot prove legacy transition delete state", + )), + _ => Err(std::io::Error::new( + std::io::ErrorKind::WouldBlock, + "remote tier object version state is unknown", + )), + } + } + } +} + +async fn execute_resolved_transition_delete( + oi: &ObjectInfo, + lease: &TierOperationLease, + resolved: ResolvedTransitionDeleteVersion, +) -> Result<(), std::io::Error> { + if !resolved.remote_already_missing { + delete_object_from_remote_tier_with_lease_idempotent( + &oi.transitioned_object.name, + &oi.transitioned_object.version_id, + lease, + resolved.version_id_exact, + ) + .await?; + } + Ok(()) } async fn delete_free_version_remote_object_with_lease( oi: &ObjectInfo, lease: &TierOperationLease, - version_id_exact: bool, + delete_plan: TransitionDeleteVersionPlan, ) -> Result<(), std::io::Error> { - delete_object_from_remote_tier_with_lease_idempotent( - &oi.transitioned_object.name, - &oi.transitioned_object.version_id, - lease, - version_id_exact, - ) - .await?; - Ok(()) + let resolved = resolve_transition_delete_version_plan(oi, lease, delete_plan).await?; + execute_resolved_transition_delete(oi, lease, resolved).await } fn free_version_physical_topology_generation(api: &ECStore) -> String { @@ -641,6 +781,16 @@ fn free_version_remote_tuple_matches(candidate: &ObjectInfo, expected: &ObjectIn if candidate.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown || expected.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown { + let candidate_legacy_missing = legacy_transition_version_state_missing(candidate)?; + let expected_legacy_missing = legacy_transition_version_state_missing(expected)?; + if candidate.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown + && expected.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown + && candidate_legacy_missing + && expected_legacy_missing + && candidate.transitioned_object.version_id == expected.transitioned_object.version_id + { + return Ok(true); + } return Err(std::io::Error::new( std::io::ErrorKind::WouldBlock, "tier free-version remote version state is unknown", @@ -716,7 +866,7 @@ async fn cleanup_free_version_exact(api: Arc, oi: &ObjectInfo, cancel: .acquire_bucket_lifecycle_read_lock(&oi.bucket) .await .map_err(std::io::Error::other)?; - let (lease, version_id_exact) = acquire_free_version_tier_lease(oi, &api.tier_config_mgr()).await?; + let (lease, delete_plan) = acquire_free_version_tier_lease(oi, &api.tier_config_mgr()).await?; let local_object = encode_dir_object(&oi.name); let object_guards = api .acquire_all_physical_object_write_locks("tier_free_version_cleanup", &oi.bucket, &local_object) @@ -734,16 +884,30 @@ async fn cleanup_free_version_exact(api: Arc, oi: &ObjectInfo, cancel: "tier free-version cleanup fence is invalid before remote delete", )); } + let resolved = tokio::select! { + _ = cancel.cancelled() => { + return Err(std::io::Error::new(std::io::ErrorKind::Interrupted, "tier free-version cleanup was cancelled")); + } + result = tokio::time::timeout_at(deadline, resolve_transition_delete_version_plan(oi, &lease, delete_plan)) => { + result.map_err(|_| { + std::io::Error::new(std::io::ErrorKind::TimedOut, "tier free-version remote probe timed out") + })?? + } + }; + if !free_version_cleanup_fences_current(&topology_generation, &api, &bucket_guard, &object_guards, &lease, cancel, deadline) { + return Err(std::io::Error::new( + std::io::ErrorKind::WouldBlock, + "tier free-version cleanup fence changed after remote probe", + )); + } tokio::select! { _ = cancel.cancelled() => { return Err(std::io::Error::new(std::io::ErrorKind::Interrupted, "tier free-version cleanup was cancelled")); } - result = tokio::time::timeout_at( - deadline, - delete_free_version_remote_object_with_lease(oi, &lease, version_id_exact), - ) => { - result - .map_err(|_| std::io::Error::new(std::io::ErrorKind::TimedOut, "tier free-version remote delete timed out"))??; + result = tokio::time::timeout_at(deadline, execute_resolved_transition_delete(oi, &lease, resolved)) => { + result.map_err(|_| { + std::io::Error::new(std::io::ErrorKind::TimedOut, "tier free-version remote delete timed out") + })??; } } if !free_version_cleanup_fences_current(&topology_generation, &api, &bucket_guard, &object_guards, &lease, cancel, deadline) { @@ -791,8 +955,8 @@ async fn delete_free_version_remote_object( oi: &ObjectInfo, tier_config_mgr: &Arc>, ) -> Result<(), std::io::Error> { - let (lease, version_id_exact) = acquire_free_version_tier_lease(oi, tier_config_mgr).await?; - delete_free_version_remote_object_with_lease(oi, &lease, version_id_exact).await + let (lease, delete_plan) = acquire_free_version_tier_lease(oi, tier_config_mgr).await?; + delete_free_version_remote_object_with_lease(oi, &lease, delete_plan).await } #[allow( @@ -808,8 +972,8 @@ where F: FnOnce() -> Fut, Fut: std::future::Future, { - let (lease, version_id_exact) = acquire_free_version_tier_lease(oi, tier_config_mgr).await?; - delete_free_version_remote_object_with_lease(oi, &lease, version_id_exact).await?; + let (lease, delete_plan) = acquire_free_version_tier_lease(oi, tier_config_mgr).await?; + delete_free_version_remote_object_with_lease(oi, &lease, delete_plan).await?; let result = delete_local().await; drop(lease); Ok(result) @@ -4688,6 +4852,39 @@ fn validate_transition_remote_version(oi: &ObjectInfo) -> Result Result { + let version = oi.transitioned_object.version_id.as_str(); + match oi.transition_version_state { + rustfs_filemeta::TransitionVersionState::Unknown => { + if !legacy_transition_version_state_missing(oi)? { + return validate_transition_remote_version(oi).map(|_| TransitionReadVersionPlan::Direct); + } + if version.is_empty() { + Ok(TransitionReadVersionPlan::ProbeLegacyUnversioned) + } else { + Ok(TransitionReadVersionPlan::Direct) + } + } + rustfs_filemeta::TransitionVersionState::KnownDisabled if version.is_empty() => Ok(TransitionReadVersionPlan::Direct), + rustfs_filemeta::TransitionVersionState::SuspendedNull if version == "null" => Ok(TransitionReadVersionPlan::Direct), + rustfs_filemeta::TransitionVersionState::Exact if !version.is_empty() && version != "null" => { + Ok(TransitionReadVersionPlan::Direct) + } + _ => Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "remote tier object version state conflicts with its version ID", + )), + } +} + // The resolver joins the tier manager as the second injected port this read // needs; grouping the request half into a struct would churn every call site of // a bug fix. @@ -4702,7 +4899,12 @@ pub(crate) async fn get_transitioned_object_reader_with_tier_manager( tier_config_mgr: &Arc>, resolver: Option<&dyn ObjectEncryptionResolver>, ) -> Result { - validate_transition_remote_version(oi)?; + let read_plan = transition_remote_version_read_plan(oi)?; + // Reject invalid ranges and encryption requests before a compatibility + // probe can amplify them into remote listing work. + let plan = ReadPlan::build_for_request(rs.clone(), oi, opts, h, resolver) + .await + .map_err(|err| std::io::Error::other(format!("building the read plan for {bucket}/{object} failed: {err}")))?; let expected_identity = tier_destination_id_from_metadata(&oi.user_defined)?; let lease = match expected_identity { Some(identity) => { @@ -4716,7 +4918,36 @@ pub(crate) async fn get_transitioned_object_reader_with_tier_manager( Err(err) => return Err(std::io::Error::other(err)), }; - tgt_client.validate_remote_version_id(&oi.transitioned_object.version_id)?; + match read_plan { + TransitionReadVersionPlan::Direct => { + tgt_client.validate_remote_version_id(&oi.transitioned_object.version_id)?; + } + TransitionReadVersionPlan::ProbeLegacyUnversioned => { + // RUSTFS_COMPAT_TODO(backlog#2203): remove operation-time probing + // after an admin reconcile can persist every proven legacy state. + let probe = tokio::time::timeout( + LEGACY_TRANSITION_READ_PROBE_TIMEOUT, + tgt_client.probe_transition_candidate(&oi.transitioned_object.name), + ) + .await + .map_err(|_| std::io::Error::new(std::io::ErrorKind::TimedOut, "legacy remote tier version probe timed out"))??; + match probe { + crate::services::tier::warm_backend::TransitionCandidateProbe::UnversionedPresent => {} + crate::services::tier::warm_backend::TransitionCandidateProbe::Unsupported => { + return Err(std::io::Error::new( + std::io::ErrorKind::Unsupported, + "remote tier cannot prove legacy unversioned transition state", + )); + } + _ => { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "remote tier object version state is unknown", + )); + } + } + } + } // The same read plan the local path uses, so the tier fetch is positioned in // the object's *stored* coordinate system and the stream is handed the same @@ -4724,9 +4955,6 @@ pub(crate) async fn get_transitioned_object_reader_with_tier_manager( // through a plaintext-coordinate range and skipping the transform is how a // transitioned SSE object used to come back as silently corrupt bytes of the // right length (rustfs/rustfs#6025). - let plan = ReadPlan::build_for_request(rs.clone(), oi, opts, h, resolver) - .await - .map_err(|err| std::io::Error::other(format!("building the read plan for {bucket}/{object} failed: {err}")))?; let (off, length) = (plan.storage_offset() as i64, plan.storage_length()); let mut gopts = WarmBackendGetOpts::default(); @@ -5599,11 +5827,13 @@ mod tests { use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}; use crate::object_api::{ObjectInfo, ObjectOptions, PutObjReader}; #[cfg(feature = "test-util")] + use crate::services::tier::test_util::MockWarmOp; + #[cfg(feature = "test-util")] use crate::services::tier::test_util::register_mock_tier; #[cfg(feature = "test-util")] use crate::services::tier::tier::TierConfigMgr; #[cfg(feature = "test-util")] - use crate::services::tier::warm_backend::WarmBackend as _; + use crate::services::tier::warm_backend::{TransitionCandidateProbe, WarmBackend as _}; use crate::set_disk::{MultipartCommitBarrier, MultipartCommitPause}; use crate::set_disk::{RUSTFS_MULTIPART_BUCKET_KEY, RUSTFS_MULTIPART_OBJECT_KEY}; use crate::storage_api_contracts::namespace::NamespaceLocking as _; @@ -6299,7 +6529,75 @@ mod tests { #[cfg(feature = "test-util")] #[tokio::test] - async fn transitioned_get_rejects_unknown_version_state_before_backend_io() { + async fn transitioned_get_allows_legacy_unknown_exact_version_for_non_destructive_read() { + let manager = TierConfigMgr::new(); + let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + let backend = register_mock_tier(&manager, &tier).await; + let remote_object = format!("remote/{}", Uuid::new_v4()); + let body = Bytes::from_static(b"legacy transitioned object body"); + let remote_version = backend + .put( + &remote_object, + ReaderImpl::Body(body.clone()), + i64::try_from(body.len()).expect("body length should fit"), + ) + .await + .expect("mock remote object should be stored"); + let mut user_defined = HashMap::new(); + insert_legacy_transition_version_id(&mut user_defined, &remote_version); + let object_info = ObjectInfo { + bucket: "bucket".to_string(), + name: "object".to_string(), + size: i64::try_from(body.len()).expect("body length should fit"), + transitioned_object: TransitionedObject { + name: remote_object, + version_id: remote_version, + status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(), + tier: tier.clone(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined.into(), + ..Default::default() + }; + + let range = Some(crate::storage_api_contracts::range::HTTPRangeSpec { + is_suffix_length: false, + start: 7, + end: 18, + }); + let mut reader = get_transitioned_object_reader_with_tier_manager( + &object_info.bucket, + &object_info.name, + &range, + &HeaderMap::new(), + &object_info, + &ObjectOptions::default(), + &manager, + None, + ) + .await + .expect("legacy unknown state should still allow a non-destructive read"); + let mut got = Vec::new(); + reader + .stream + .read_to_end(&mut got) + .await + .expect("transitioned reader should drain"); + + assert_eq!(got, &body.as_ref()[7..=18]); + assert_eq!(backend.get_count().await, 1); + assert_eq!(backend.remove_count().await, 0); + assert_eq!( + TierConfigMgr::active_operation_lease_count(&manager, &tier).await, + 0, + "tier generation lease should release after EOF" + ); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn transitioned_get_rejects_explicit_unknown_version_state_before_backend_io() { let manager = TierConfigMgr::new(); let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); let backend = register_mock_tier(&manager, &tier).await; @@ -6315,6 +6613,181 @@ mod tests { ..Default::default() }, transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined_with_transition_version_state(rustfs_filemeta::TransitionVersionState::Unknown).into(), + ..Default::default() + }; + + let err = match get_transitioned_object_reader_with_tier_manager( + &object_info.bucket, + &object_info.name, + &None, + &HeaderMap::new(), + &object_info, + &ObjectOptions::default(), + &manager, + None, + ) + .await + { + Ok(_) => panic!("explicit unknown remote version state must fail before backend IO"), + Err(err) => err, + }; + + assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + assert_eq!(backend.op_log().await, Vec::::new()); + assert_eq!(backend.get_count().await, 0); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn transitioned_get_rejects_present_but_invalid_legacy_version_metadata() { + let manager = TierConfigMgr::new(); + let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + let backend = register_mock_tier(&manager, &tier).await; + + for persisted_version in [ + Uuid::nil().to_string(), + "\u{fffd}".to_string(), + "bad\u{0001}version".to_string(), + ] { + let mut user_defined = HashMap::new(); + insert_legacy_transition_version_id(&mut user_defined, &persisted_version); + let object_info = ObjectInfo { + bucket: "bucket".to_string(), + name: "object".to_string(), + size: 1, + transitioned_object: TransitionedObject { + name: "remote/object".to_string(), + version_id: String::new(), + status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(), + tier: tier.clone(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined.into(), + ..Default::default() + }; + + let err = match get_transitioned_object_reader_with_tier_manager( + &object_info.bucket, + &object_info.name, + &None, + &HeaderMap::new(), + &object_info, + &ObjectOptions::default(), + &manager, + None, + ) + .await + { + Ok(_) => panic!("present but invalid legacy version metadata must fail before backend IO"), + Err(err) => err, + }; + + assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + } + + assert_eq!(backend.op_log().await, Vec::::new()); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn transitioned_get_probes_legacy_empty_unknown_state_before_unversioned_read() { + let manager = TierConfigMgr::new(); + let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + let backend = register_mock_tier(&manager, &tier).await; + backend.set_put_remote_version(Some(String::new())).await; + let remote_object = format!("remote/{}", Uuid::new_v4()); + let body = Bytes::from_static(b"legacy unversioned transitioned object body"); + let remote_version = backend + .put( + &remote_object, + ReaderImpl::Body(body.clone()), + i64::try_from(body.len()).expect("body length should fit"), + ) + .await + .expect("mock remote object should be stored"); + assert!(remote_version.is_empty()); + let object_info = ObjectInfo { + bucket: "bucket".to_string(), + name: "object".to_string(), + size: i64::try_from(body.len()).expect("body length should fit"), + transitioned_object: TransitionedObject { + name: remote_object.clone(), + version_id: String::new(), + status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(), + tier: tier.clone(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: HashMap::from([("x-minio-internal-transitioned-versionID".to_string(), String::new())]).into(), + ..Default::default() + }; + + let mut reader = get_transitioned_object_reader_with_tier_manager( + &object_info.bucket, + &object_info.name, + &None, + &HeaderMap::new(), + &object_info, + &ObjectOptions::default(), + &manager, + None, + ) + .await + .expect("probe-proven legacy unversioned state should allow a non-destructive read"); + let mut got = Vec::new(); + reader + .stream + .read_to_end(&mut got) + .await + .expect("transitioned reader should drain"); + + assert_eq!(got, body.as_ref()); + assert_eq!(backend.remove_count().await, 0); + assert_eq!( + backend.op_log().await, + vec![ + MockWarmOp::Put { + object: remote_object.clone() + }, + MockWarmOp::Probe { + object: remote_object.clone() + }, + MockWarmOp::Get { object: remote_object }, + ] + ); + assert_eq!( + TierConfigMgr::active_operation_lease_count(&manager, &tier).await, + 0, + "tier generation lease should release after EOF" + ); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn transitioned_get_rejects_ambiguous_empty_unknown_state_without_backend_get() { + let manager = TierConfigMgr::new(); + let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + let backend = register_mock_tier(&manager, &tier).await; + let remote_object = format!("remote/{}", Uuid::new_v4()); + backend + .set_transition_candidate_probe_override(Some(TransitionCandidateProbe::VersionedPresent( + "versioned-candidate".to_string(), + ))) + .await; + let object_info = ObjectInfo { + bucket: "bucket".to_string(), + name: "object".to_string(), + size: 1, + transitioned_object: TransitionedObject { + name: remote_object.clone(), + version_id: String::new(), + status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(), + tier, + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, ..Default::default() }; @@ -6330,19 +6803,28 @@ mod tests { ) .await { - Ok(_) => panic!("unknown remote version state must fail before backend IO"), + Ok(_) => panic!("versioned legacy unknown state without stored version must fail before backend GET"), Err(err) => err, }; assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + assert_eq!(backend.op_log().await, vec![MockWarmOp::Probe { object: remote_object }]); assert_eq!(backend.get_count().await, 0); + assert_eq!(backend.remove_count().await, 0); } #[cfg(feature = "test-util")] #[tokio::test] - async fn free_version_delete_rejects_unknown_version_state_before_backend_io() { + async fn free_version_delete_rejects_explicit_unknown_before_backend_io() { let manager = TierConfigMgr::new(); let backend = register_mock_tier(&manager, "WARM").await; + let identity = test_tier_destination_identity(&manager, "WARM").await; + let mut user_defined = user_defined_with_tier_destination_identity(identity); + rustfs_utils::http::metadata_compat::insert_str( + &mut user_defined, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE, + rustfs_filemeta::TransitionVersionState::Unknown.as_str().to_string(), + ); let object_info = ObjectInfo { transitioned_object: TransitionedObject { name: "remote/object".to_string(), @@ -6351,17 +6833,251 @@ mod tests { ..Default::default() }, transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined.into(), ..Default::default() }; let err = super::delete_free_version_remote_object(&object_info, &manager) .await - .expect_err("unknown remote version state must fail before backend IO"); + .expect_err("explicit unknown cleanup must fail before backend IO"); assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + assert!(err.to_string().contains("version state is unknown")); + assert_eq!(backend.op_log().await, Vec::::new()); assert_eq!(backend.remove_count().await, 0); } + #[cfg(feature = "test-util")] + async fn test_tier_destination_identity( + manager: &Arc>, + tier: &str, + ) -> crate::services::tier::tier::TierDestinationId { + TierConfigMgr::acquire_operation_lease(manager, tier) + .await + .expect("test tier lease should be available") + .backend_identity() + } + + #[cfg(feature = "test-util")] + fn user_defined_with_tier_destination_identity( + identity: crate::services::tier::tier::TierDestinationId, + ) -> HashMap { + let mut user_defined = HashMap::new(); + rustfs_utils::http::metadata_compat::insert_str( + &mut user_defined, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + rustfs_utils::crypto::hex(identity), + ); + user_defined + } + + #[cfg(feature = "test-util")] + fn user_defined_with_transition_version_state(state: rustfs_filemeta::TransitionVersionState) -> HashMap { + let mut user_defined = HashMap::new(); + rustfs_utils::http::metadata_compat::insert_str( + &mut user_defined, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE, + state.as_str().to_string(), + ); + user_defined + } + + #[cfg(feature = "test-util")] + fn insert_legacy_transition_version_id(user_defined: &mut HashMap, version_id: &str) { + rustfs_utils::http::metadata_compat::insert_str( + user_defined, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_ID, + version_id.to_string(), + ); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn free_version_tuple_rejects_mixed_legacy_missing_and_explicit_unknown() { + let manager = TierConfigMgr::new(); + register_mock_tier(&manager, "WARM").await; + let identity = test_tier_destination_identity(&manager, "WARM").await; + let mut legacy_metadata = user_defined_with_tier_destination_identity(identity); + insert_legacy_transition_version_id(&mut legacy_metadata, "legacy-version"); + let mut explicit_metadata = legacy_metadata.clone(); + rustfs_utils::http::metadata_compat::insert_str( + &mut explicit_metadata, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE, + rustfs_filemeta::TransitionVersionState::Unknown.as_str().to_string(), + ); + let make_info = |user_defined: HashMap| ObjectInfo { + transitioned_object: TransitionedObject { + name: "remote/object".to_string(), + version_id: "legacy-version".to_string(), + tier: "WARM".to_string(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined.into(), + ..Default::default() + }; + + let err = super::free_version_remote_tuple_matches(&make_info(legacy_metadata), &make_info(explicit_metadata)) + .expect_err("mixed legacy-missing and explicit unknown provenance must fail closed"); + + assert_eq!(err.kind(), std::io::ErrorKind::WouldBlock); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn free_version_delete_probes_exact_version_hidden_by_current_delete_marker() { + let manager = TierConfigMgr::new(); + let tier = "WARM"; + let backend = register_mock_tier(&manager, tier).await; + let identity = test_tier_destination_identity(&manager, tier).await; + let remote_object = format!("remote/{}", Uuid::new_v4()); + let body = Bytes::from_static(b"legacy exact cleanup body"); + let remote_version = backend + .put( + &remote_object, + ReaderImpl::Body(body), + i64::try_from(b"legacy exact cleanup body".len()).expect("body length should fit"), + ) + .await + .expect("mock remote object should be stored"); + let mut user_defined = user_defined_with_tier_destination_identity(identity); + insert_legacy_transition_version_id(&mut user_defined, &remote_version); + backend + .set_transition_candidate_probe_override(Some(TransitionCandidateProbe::Missing)) + .await; + assert_eq!( + backend + .probe_transition_candidate_state(&remote_object) + .await + .expect("current remote view should be readable"), + TransitionCandidateProbe::Missing, + "a current delete marker must hide the historical data version from an unversioned probe" + ); + backend.clear_op_log().await; + let object_info = ObjectInfo { + transitioned_object: TransitionedObject { + name: remote_object.clone(), + version_id: remote_version, + tier: tier.to_string(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined.into(), + ..Default::default() + }; + + super::delete_free_version_remote_object(&object_info, &manager) + .await + .expect("probe-proven legacy exact cleanup should delete the remote version"); + super::delete_free_version_remote_object(&object_info, &manager) + .await + .expect("a retry after the exact remote version is already missing should be idempotent"); + + assert_eq!( + backend.op_log().await, + vec![ + MockWarmOp::Get { + object: remote_object.clone() + }, + MockWarmOp::Remove { + object: remote_object.clone() + }, + MockWarmOp::Get { + object: remote_object.clone() + }, + ] + ); + assert_eq!( + backend.remove_versions().await, + vec![(remote_object, object_info.transitioned_object.version_id)] + ); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn free_version_delete_retains_legacy_unknown_unversioned_object() { + let manager = TierConfigMgr::new(); + let tier = "WARM"; + let backend = register_mock_tier(&manager, tier).await; + backend.set_put_remote_version(Some(String::new())).await; + let identity = test_tier_destination_identity(&manager, tier).await; + let remote_object = format!("remote/{}", Uuid::new_v4()); + let body = Bytes::from_static(b"legacy unversioned cleanup body"); + let remote_version = backend + .put( + &remote_object, + ReaderImpl::Body(body), + i64::try_from(b"legacy unversioned cleanup body".len()).expect("body length should fit"), + ) + .await + .expect("mock remote object should be stored"); + assert!(remote_version.is_empty()); + backend.clear_op_log().await; + let mut user_defined = user_defined_with_tier_destination_identity(identity); + user_defined.insert("x-minio-internal-transitioned-versionID".to_string(), String::new()); + let object_info = ObjectInfo { + transitioned_object: TransitionedObject { + name: remote_object.clone(), + version_id: String::new(), + tier: tier.to_string(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined.into(), + ..Default::default() + }; + + let err = super::delete_free_version_remote_object(&object_info, &manager) + .await + .expect_err("legacy unversioned cleanup cannot exclude a versioning-state race"); + + assert_eq!(err.kind(), std::io::ErrorKind::WouldBlock); + assert!(backend.op_log().await.is_empty()); + assert_eq!(backend.remove_count().await, 0); + assert!(backend.remove_versions().await.is_empty()); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn free_version_delete_does_not_remove_a_different_remote_version() { + let manager = TierConfigMgr::new(); + let tier = "WARM"; + let backend = register_mock_tier(&manager, tier).await; + let identity = test_tier_destination_identity(&manager, tier).await; + let remote_object = format!("remote/{}", Uuid::new_v4()); + backend.set_put_remote_version(Some("different-version".to_string())).await; + backend + .put( + &remote_object, + ReaderImpl::Body(Bytes::from_static(b"different remote version")), + i64::try_from(b"different remote version".len()).expect("body length should fit"), + ) + .await + .expect("different remote version should be stored"); + backend.clear_op_log().await; + let mut user_defined = user_defined_with_tier_destination_identity(identity); + insert_legacy_transition_version_id(&mut user_defined, "legacy-version"); + let object_info = ObjectInfo { + transitioned_object: TransitionedObject { + name: remote_object.clone(), + version_id: "legacy-version".to_string(), + tier: tier.to_string(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined.into(), + ..Default::default() + }; + + super::delete_free_version_remote_object(&object_info, &manager) + .await + .expect("a missing exact legacy version should be an idempotent cleanup success"); + + assert_eq!(backend.op_log().await, vec![MockWarmOp::Get { object: remote_object }]); + assert_eq!(backend.remove_count().await, 0); + assert!(backend.remove_versions().await.is_empty()); + } + #[cfg(feature = "test-util")] #[tokio::test] async fn free_version_remote_delete_requires_persisted_destination_identity() { diff --git a/crates/ecstore/src/services/tier/test_util.rs b/crates/ecstore/src/services/tier/test_util.rs index f2c4ddeb6..59b4d6f4e 100644 --- a/crates/ecstore/src/services/tier/test_util.rs +++ b/crates/ecstore/src/services/tier/test_util.rs @@ -701,7 +701,7 @@ impl WarmBackend for MockWarmBackend { Ok(version) } - async fn get(&self, object: &str, _rv: &str, opts: WarmBackendGetOpts) -> Result { + async fn get(&self, object: &str, rv: &str, opts: WarmBackendGetOpts) -> Result { self.precondition().await?; let barrier = self.inner.get_barrier.lock().await.take(); if let Some(barrier) = barrier { @@ -719,6 +719,9 @@ impl WarmBackend for MockWarmBackend { let Some(stored) = objects.get(object) else { return Err(std::io::Error::new(std::io::ErrorKind::NotFound, "mock object not found")); }; + if !rv.is_empty() && stored.remote_version_id != rv { + return Err(std::io::Error::new(std::io::ErrorKind::NotFound, "NoSuchVersion")); + } let bytes = &stored.bytes; let start = opts.start_offset.max(0) as usize; diff --git a/crates/ecstore/src/services/tier/tier.rs b/crates/ecstore/src/services/tier/tier.rs index 887015a1e..af24423fa 100644 --- a/crates/ecstore/src/services/tier/tier.rs +++ b/crates/ecstore/src/services/tier/tier.rs @@ -2346,6 +2346,10 @@ impl WarmBackend for SharedWarmBackendProxy { self.0.probe_transition_candidate(object).await } + async fn probe_transition_version(&self, object: &str, remote_version_id: &str) -> io::Result { + self.0.probe_transition_version(object, remote_version_id).await + } + async fn in_use(&self) -> io::Result { self.0.in_use().await } @@ -2458,6 +2462,15 @@ impl TierOperationLease { Ok(()) } + pub(crate) async fn probe_transition_version( + &self, + object: &str, + remote_version_id: &str, + ) -> io::Result { + self.validate_remote_version_id(remote_version_id)?; + self.inner.driver.probe_transition_version(object, remote_version_id).await + } + pub(crate) fn is_current_generation(&self) -> bool { lock_unpoisoned(&self.runtime) .generations diff --git a/crates/ecstore/src/services/tier/warm_backend.rs b/crates/ecstore/src/services/tier/warm_backend.rs index ee48c116c..b5cf4ab38 100644 --- a/crates/ecstore/src/services/tier/warm_backend.rs +++ b/crates/ecstore/src/services/tier/warm_backend.rs @@ -40,6 +40,7 @@ use rustfs_s3_client::credentials::{Credentials, SignatureType, Static, Value}; use rustfs_s3_client::transition_api::{BucketLookupType, Options, TransitionClient, TransitionCore}; use rustfs_s3_client::{ admin_handler_utils::AdminError, + api_error_response::to_error_response, api_put_object::{AdvancedPutOptions, PutObjectOptions}, transition_api::{ReadCloser, ReaderImpl}, }; @@ -48,11 +49,14 @@ use rustfs_utils::egress::validate_outbound_url; use rustfs_utils::http::headers::{ CACHE_CONTROL, CONTENT_DISPOSITION, CONTENT_ENCODING, CONTENT_LANGUAGE, CONTENT_TYPE, EXPIRES, HeaderExt as _, }; -use s3s::dto::{ObjectLockLegalHoldStatus, ObjectLockRetentionMode, ReplicationStatus}; use s3s::header::{ X_AMZ_OBJECT_LOCK_LEGAL_HOLD, X_AMZ_OBJECT_LOCK_MODE, X_AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE, X_AMZ_REPLICATION_STATUS, X_AMZ_STORAGE_CLASS, }; +use s3s::{ + S3ErrorCode, + dto::{ObjectLockLegalHoldStatus, ObjectLockRetentionMode, ReplicationStatus}, +}; use std::collections::HashMap; use std::sync::Arc; use std::time::Duration; @@ -141,6 +145,42 @@ pub trait WarmBackend { async fn probe_transition_candidate(&self, _object: &str) -> Result { Ok(TransitionCandidateProbe::Unsupported) } + async fn probe_transition_version( + &self, + object: &str, + remote_version_id: &str, + ) -> Result { + if remote_version_id.is_empty() { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "an exact tier probe requires a remote version ID", + )); + } + self.validate_remote_version_id(remote_version_id)?; + match self + .get( + object, + remote_version_id, + WarmBackendGetOpts { + start_offset: 0, + length: 1, + }, + ) + .await + { + Ok(_) => Ok(TransitionCandidateProbe::VersionedPresent(remote_version_id.to_string())), + Err(err) if matches!(to_error_response(&err).code, S3ErrorCode::InvalidRange) => { + Ok(TransitionCandidateProbe::VersionedPresent(remote_version_id.to_string())) + } + Err(err) + if err.kind() == std::io::ErrorKind::NotFound + || matches!(to_error_response(&err).code, S3ErrorCode::NoSuchKey | S3ErrorCode::NoSuchVersion) => + { + Ok(TransitionCandidateProbe::Missing) + } + Err(err) => Err(err), + } + } async fn in_use(&self) -> Result; } @@ -437,6 +477,17 @@ impl WarmBackend for MeteredWarmBackend { Self::record(TierRequestOperation::Probe, result) } + async fn probe_transition_version( + &self, + object: &str, + remote_version_id: &str, + ) -> Result { + Self::record( + TierRequestOperation::Probe, + self.inner.probe_transition_version(object, remote_version_id).await, + ) + } + async fn in_use(&self) -> Result { Self::record(TierRequestOperation::InUse, self.inner.in_use().await) } diff --git a/crates/ecstore/src/services/tier/warm_backend_s3.rs b/crates/ecstore/src/services/tier/warm_backend_s3.rs index 5462fc52c..b830ea7f2 100644 --- a/crates/ecstore/src/services/tier/warm_backend_s3.rs +++ b/crates/ecstore/src/services/tier/warm_backend_s3.rs @@ -529,6 +529,10 @@ mod tests { "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 63\r\nConnection: close\r\n\r\nNoSuchKeymissing", "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 66\r\nConnection: close\r\n\r\nNoSuchObjectmissing", "HTTP/1.1 403 Forbidden\r\nContent-Type: application/xml\r\nContent-Length: 65\r\nConnection: close\r\n\r\nAccessDenieddenied", + "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 63\r\nConnection: close\r\n\r\nNoSuchKeymissing", + "HTTP/1.1 416 Range Not Satisfiable\r\nContent-Type: application/xml\r\nContent-Length: 72\r\nConnection: close\r\n\r\nInvalidRangeempty version", + "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 67\r\nConnection: close\r\n\r\nNoSuchVersionmissing", + "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 63\r\nConnection: close\r\n\r\nNoSuchKeymissing", ]; let mut requests = Vec::new(); for response in responses { @@ -622,15 +626,52 @@ mod tests { .await .expect_err("an authorization failure must not be mistaken for a missing key"); assert_eq!(to_error_response(&err).code, S3ErrorCode::AccessDenied); + assert_eq!( + backend + .probe_transition_candidate("delete-marker-hidden") + .await + .expect("a current delete marker should hide the data version"), + TransitionCandidateProbe::Missing + ); + assert_eq!( + backend + .probe_transition_version("delete-marker-hidden", "historical-version") + .await + .expect("the stored historical version should be probed exactly"), + TransitionCandidateProbe::VersionedPresent("historical-version".to_string()) + ); + assert_eq!( + backend + .probe_transition_version("delete-marker-hidden", "missing-version") + .await + .expect("a missing exact version should be classified"), + TransitionCandidateProbe::Missing + ); + assert_eq!( + backend + .probe_transition_version("missing-object", "historical-version") + .await + .expect("a missing key for an exact version probe should be classified"), + TransitionCandidateProbe::Missing + ); let requests = fixture.await.expect("candidate fixture should join"); - for request in requests { + for request in &requests[..6] { let request = request.to_ascii_lowercase(); assert!(request.starts_with("get /bucket/"), "candidate discovery must use object GET"); assert!(request.contains("\r\nrange: bytes=0-0\r\n")); assert!(!request.contains("?versioning")); assert!(!request.contains("?versions")); } + for request in &requests[6..] { + let request = request.to_ascii_lowercase(); + assert!(request.starts_with("get /bucket/"), "exact discovery must use object GET"); + assert!(request.contains("\r\nrange: bytes=0-0\r\n")); + } + assert!(!requests[5].to_ascii_lowercase().contains("versionid=")); + assert!(requests[6].to_ascii_lowercase().contains("?versionid=historical-version")); + assert!(requests[7].to_ascii_lowercase().contains("?versionid=missing-version")); + assert!(requests[8].to_ascii_lowercase().contains("?versionid=historical-version")); } fn list_versions(versions: &[(&str, &str)], delete_markers: &[(&str, &str)], is_truncated: bool) -> ListVersionsResult { diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 22d7e37c3..463f50a82 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -11575,6 +11575,7 @@ mod tests { pool_index: usize, bucket: &str, object: &str, + minio_unversioned: bool, ) { for disk_index in 0..4 { let metadata_path = @@ -11608,6 +11609,11 @@ mod tests { ] { rustfs_utils::http::metadata_compat::remove_bytes(&mut object_meta.meta_sys, suffix); } + if minio_unversioned { + object_meta + .meta_sys + .insert("x-minio-internal-transitioned-versionID".to_string(), Vec::new()); + } *shallow = rustfs_filemeta::FileMetaShallowVersion::try_from(version) .expect("legacy transitioned version should re-encode"); } @@ -11618,6 +11624,152 @@ mod tests { } } + #[cfg(feature = "test-util")] + async fn read_store_body( + store: &Arc, + bucket: &str, + object: &str, + range: Option, + opts: &ObjectOptions, + ) -> Vec { + let mut reader = store + .get_object_reader(bucket, object, range, HeaderMap::new(), opts) + .await + .expect("object reader should open"); + let mut body = Vec::new(); + reader.stream.read_to_end(&mut body).await.expect("object body should drain"); + body + } + + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn legacy_unknown_unversioned_transition_supports_head_get_and_range_without_backfill() { + let temp_dir = tempfile::tempdir().expect("create legacy unknown unversioned store dir"); + let (ctx, store, _shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "legacy-unknown-unversioned-read", &[4])).await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + let tier_name = "LEGACY-UNKNOWN-UNVERSIONED-READ"; + let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await; + backend.set_put_remote_version(Some(String::new())).await; + let bucket = "legacy-unknown-unversioned-read-bucket"; + let object = "object.bin"; + let payload = b"legacy unversioned remote tier object remains readable".repeat(1024); + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("legacy source bucket should be created"); + let mut reader = PutObjReader::from_vec(payload.clone()); + let source = store + .put_object(bucket, object, &mut reader, &ObjectOptions::default()) + .await + .expect("legacy source should be written"); + store + .transition_object( + bucket, + object, + &ObjectOptions { + transition: TransitionOptions { + status: TRANSITION_PENDING.to_string(), + tier: tier_name.to_string(), + etag: source.etag.clone().expect("legacy source should have an etag"), + ..Default::default() + }, + mod_time: source.mod_time, + ..Default::default() + }, + ) + .await + .expect("legacy source should transition"); + rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 0, bucket, object, true).await; + backend.clear_op_log().await; + + let opts = ObjectOptions { + metadata_cache_safe: false, + ..Default::default() + }; + let head = store + .get_object_info(bucket, object, &opts) + .await + .expect("legacy transitioned HEAD should use local metadata"); + assert_eq!(head.transition_version_state, rustfs_filemeta::TransitionVersionState::Unknown); + assert!(head.transitioned_object.version_id.is_empty()); + assert_eq!( + head.user_defined + .get("x-minio-internal-transitioned-versionID") + .map(String::as_str), + Some(""), + "the MinIO empty version-key provenance must survive xl.meta decoding" + ); + assert!( + !rustfs_utils::http::metadata_compat::contains_key_str( + &head.user_defined, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE, + ), + "the compatibility read must not synthesize version-state metadata" + ); + + let full_body = read_store_body(&store, bucket, object, None, &opts).await; + assert_eq!(full_body, payload); + + let range = HTTPRangeSpec { + is_suffix_length: false, + start: 7, + end: 38, + }; + let ranged_body = read_store_body(&store, bucket, object, Some(range), &opts).await; + assert_eq!(ranged_body, &payload[7..=38]); + + let after_read = store.pools[0] + .get_disks_by_key(object) + .load_file_info_versions_exact(bucket, object) + .await + .expect("legacy metadata should remain readable after GET") + .expect("legacy object metadata should remain on disk") + .versions + .into_iter() + .find(|version| version.transition_status == rustfs_filemeta::TRANSITION_COMPLETE) + .expect("legacy transitioned source should remain visible after GET"); + assert_eq!(after_read.transition_version_state, rustfs_filemeta::TransitionVersionState::Unknown); + assert!(after_read.transition_version.is_none()); + assert!(after_read.transition_version_id.is_none()); + assert_eq!( + after_read + .metadata + .get("x-minio-internal-transitioned-versionID") + .map(String::as_str), + Some(""), + "the MinIO empty version-key provenance must remain after GET and Range GET" + ); + assert!( + !rustfs_utils::http::metadata_compat::contains_key_str( + &after_read.metadata, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE, + ), + "the compatibility read must remain side-effect free" + ); + + assert_eq!( + backend.op_log().await, + vec![ + MockWarmOp::Probe { + object: after_read.transitioned_objname.clone(), + }, + MockWarmOp::Get { + object: after_read.transitioned_objname.clone(), + }, + MockWarmOp::Probe { + object: after_read.transitioned_objname.clone(), + }, + MockWarmOp::Get { + object: after_read.transitioned_objname, + }, + ], + "legacy reads should probe before each unversioned GET and never mutate local metadata" + ); + assert_eq!(backend.remove_count().await, 0); + } + #[cfg(feature = "test-util")] #[tokio::test] #[serial_test::serial(storage_class_env)] @@ -11658,7 +11810,7 @@ mod tests { ) .await .expect("legacy source should transition"); - rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 0, bucket, object).await; + rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 0, bucket, object, false).await; let legacy = store.pools[0] .get_disks_by_key(object) .load_file_info_versions_exact(bucket, object) @@ -12799,7 +12951,7 @@ mod tests { .expect("merge-loser source should transition"); copy_test_xlmeta_between_pools(temp_dir.path(), 0, 1, bucket, object).await; } - rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 1, bucket, "legacy/item.bin").await; + rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 1, bucket, "legacy/item.bin", false).await; backend.set_remove_failure(true); store.pools[1] .delete_object(bucket, "hidden/item.bin", ObjectOptions::default()) @@ -16866,6 +17018,10 @@ mod tests { .find(|version| version.version_id == history.version_id) .expect("transitioned history should exist"); transitioned.transition_version_state = rustfs_filemeta::TransitionVersionState::Unknown; + rustfs_utils::http::metadata_compat::remove_str( + &mut transitioned.metadata, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE, + ); metadata .add_version(transitioned) .expect("unknown state should replace the transitioned version"); diff --git a/crates/filemeta/src/filemeta/version.rs b/crates/filemeta/src/filemeta/version.rs index f0969c3cf..42aa5a74d 100644 --- a/crates/filemeta/src/filemeta/version.rs +++ b/crates/filemeta/src/filemeta/version.rs @@ -297,6 +297,20 @@ fn transitioned_version_from_bytes(value: Option<&[u8]>, state: TransitionVersio } } +fn transition_version_metadata_value(raw: &[u8], decoded: Option<&str>) -> String { + decoded.map(str::to_owned).unwrap_or_else(|| { + if raw.is_empty() { + String::new() + } else { + String::from_utf8_lossy(raw).into_owned() + } + }) +} + +fn is_transition_version_metadata_key(key: &str) -> bool { + strip_internal_prefix_preserving_case(key).is_some_and(|suffix| suffix.eq_ignore_ascii_case(SUFFIX_TRANSITIONED_VERSION_ID)) +} + fn validate_transition_version_state(state: TransitionVersionState, version: Option<&str>) -> Result<()> { let valid = match state { TransitionVersionState::Unknown | TransitionVersionState::KnownDisabled => version.is_none(), @@ -366,14 +380,26 @@ impl<'a> DerivedInternalMetadata<'a> { } *slot = Some(value.as_slice()); } + fn merge_consistent<'a>(canonical: Option<&'a [u8]>, legacy: Option<&'a [u8]>) -> Result> { + if let (Some(canonical), Some(legacy)) = (canonical, legacy) + && canonical != legacy + { + return Err(Error::FileCorrupt); + } + Ok(canonical.or(legacy)) + } + Ok(Self { checksum: canonical.checksum.or(legacy.checksum), part_checksums: canonical.part_checksums.or(legacy.part_checksums), - transition_status: canonical.transition_status.or(legacy.transition_status), - transitioned_object: canonical.transitioned_object.or(legacy.transitioned_object), - transitioned_version: canonical.transitioned_version.or(legacy.transitioned_version), - transitioned_version_state: canonical.transitioned_version_state.or(legacy.transitioned_version_state), - transition_tier: canonical.transition_tier.or(legacy.transition_tier), + transition_status: merge_consistent(canonical.transition_status, legacy.transition_status)?, + transitioned_object: merge_consistent(canonical.transitioned_object, legacy.transitioned_object)?, + transitioned_version: merge_consistent(canonical.transitioned_version, legacy.transitioned_version)?, + transitioned_version_state: merge_consistent( + canonical.transitioned_version_state, + legacy.transitioned_version_state, + )?, + transition_tier: merge_consistent(canonical.transition_tier, legacy.transition_tier)?, }) } } @@ -438,8 +464,14 @@ impl FileInfo { } } -fn set_transition_version_state(meta_sys: &mut HashMap>, state: TransitionVersionState) { - if state == TransitionVersionState::Unknown { +fn set_transition_version_state( + meta_sys: &mut HashMap>, + state: TransitionVersionState, + source_metadata: &HashMap, +) { + if state == TransitionVersionState::Unknown + && !rustfs_utils::http::metadata_compat::contains_key_str(source_metadata, SUFFIX_TRANSITIONED_VERSION_STATE) + { remove_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE); } else { insert_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE, state.as_str().as_bytes().to_vec()); @@ -2643,6 +2675,11 @@ impl MetaObject { if derived_metadata.transitioned_version_state.is_some() { validate_transition_version_state(transition_version_state, transition_version.as_deref())?; } + for (key, value) in &self.meta_sys { + if is_transition_version_metadata_key(key) { + metadata.insert(key.to_owned(), transition_version_metadata_value(value, transition_version.as_deref())); + } + } let transition_version_id = transition_version.as_deref().and_then(|value| Uuid::parse_str(value).ok()); let transition_tier = derived_metadata .transition_tier @@ -2689,7 +2726,7 @@ impl MetaObject { } else { remove_bytes(&mut self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID); } - set_transition_version_state(&mut self.meta_sys, fi.transition_version_state); + set_transition_version_state(&mut self.meta_sys, fi.transition_version_state, &fi.metadata); 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) { insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID, destination_id.into_bytes()); @@ -2830,7 +2867,7 @@ impl From for MetaObject { insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version); } if !value.transition_status.is_empty() { - set_transition_version_state(&mut meta_sys, value.transition_version_state); + set_transition_version_state(&mut meta_sys, value.transition_version_state, &value.metadata); } if !value.transition_tier.is_empty() { @@ -2985,6 +3022,12 @@ impl MetaDeleteMarker { fi.transition_version_state = transition_version_state_from_bytes(derived_metadata.transitioned_version_state)?; fi.transition_version = transitioned_version_from_bytes(derived_metadata.transitioned_version, fi.transition_version_state); + for (key, value) in &self.meta_sys { + if is_transition_version_metadata_key(key) { + fi.metadata + .insert(key.to_owned(), transition_version_metadata_value(value, fi.transition_version.as_deref())); + } + } fi.transition_version_id = fi.transition_version.as_deref().and_then(|value| Uuid::parse_str(value).ok()); if derived_metadata.transitioned_version_state.is_some() { validate_transition_version_state(fi.transition_version_state, fi.transition_version.as_deref())?; @@ -3152,7 +3195,7 @@ impl From for MetaDeleteMarker { insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version); } if !value.transition_status.is_empty() || value.tier_free_version() { - set_transition_version_state(&mut meta_sys, value.transition_version_state); + set_transition_version_state(&mut meta_sys, value.transition_version_state, &value.metadata); } if !value.transition_tier.is_empty() { insert_bytes(&mut meta_sys, SUFFIX_TRANSITION_TIER, value.transition_tier.as_bytes().to_vec()); @@ -4574,6 +4617,7 @@ mod tests { .into_fileinfo("b", "k", false) .expect("into_fileinfo"); assert_eq!(fi.transition_version_id, None); + assert_eq!(get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID), Some(String::new())); } #[test] @@ -4585,6 +4629,10 @@ mod tests { .into_fileinfo("b", "k", false) .expect("into_fileinfo"); assert_eq!(fi.transition_version_id, None); + assert!( + get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID).is_some_and(|value| !value.is_empty()), + "nil UUID bytes must remain distinguishable from an empty MinIO version" + ); } #[test] @@ -4598,6 +4646,7 @@ mod tests { assert_eq!(fi.transition_version_id, Some(id)); assert_eq!(fi.transition_version, Some(id.to_string())); assert_eq!(fi.transition_version_state, TransitionVersionState::Unknown); + assert_eq!(get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID), Some(id.to_string())); } #[test] @@ -4637,6 +4686,36 @@ mod tests { assert_eq!(fi.transition_version_state, TransitionVersionState::Unknown); } + #[test] + fn meta_object_transition_version_state_explicit_unknown_is_not_legacy_missing() { + let mut metadata = HashMap::new(); + rustfs_utils::http::metadata_compat::insert_str( + &mut metadata, + SUFFIX_TRANSITIONED_VERSION_STATE, + TransitionVersionState::Unknown.as_str().to_string(), + ); + let fi = FileInfo { + transition_status: "complete".to_string(), + transition_version_state: TransitionVersionState::Unknown, + metadata, + ..Default::default() + }; + + let object = MetaObject::from(fi); + assert_eq!( + get_consistent_bytes(&object.meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE), + Some(b"unknown".as_slice()) + ); + let decoded = object + .into_fileinfo("b", "k", false) + .expect("explicit unknown state should decode"); + assert_eq!(decoded.transition_version_state, TransitionVersionState::Unknown); + assert_eq!( + rustfs_utils::http::metadata_compat::get_consistent_str(&decoded.metadata, SUFFIX_TRANSITIONED_VERSION_STATE,), + Some("unknown") + ); + } + #[test] fn meta_object_transition_version_state_exact_round_trips_dual_keys() { let id = sample_version_id(); @@ -4753,6 +4832,10 @@ mod tests { .expect("invalid transition version bytes must not fail the object read"); assert_eq!(fi.transition_version_id, None); assert_eq!(fi.transition_version, None); + assert!( + get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID).is_some_and(|value| !value.is_empty()), + "invalid raw bytes must remain distinguishable from an empty MinIO version" + ); } #[test] @@ -4795,6 +4878,10 @@ mod tests { .into_fileinfo("b", "k", false) .expect("nil tier version should remain an absent remote version"); assert_eq!(fi.transition_version_id, None); + assert!( + get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID).is_some_and(|value| !value.is_empty()), + "nil UUID bytes must remain distinguishable from an empty MinIO version" + ); } #[test] @@ -4812,6 +4899,7 @@ mod tests { .expect("legacy binary UUID tier version should decode"); assert_eq!(fi.transition_version_id, Some(id)); assert_eq!(fi.transition_version, Some(id.to_string())); + assert_eq!(get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID), Some(id.to_string())); } #[test] @@ -4910,6 +4998,23 @@ mod tests { assert_eq!(err, Error::FileCorrupt); } + #[test] + fn meta_object_transition_version_state_mixed_case_alias_conflict_fails_closed() { + let sys = HashMap::from([ + ( + format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_TRANSITIONED_VERSION_STATE}"), + b"unknown".to_vec(), + ), + ("X-Minio-Internal-transitioned-version-state".to_string(), b"exact".to_vec()), + ]); + + let err = make_meta_object_with_sys(sys) + .into_fileinfo("b", "k", false) + .expect_err("mixed-case transition state aliases must agree"); + + assert_eq!(err, Error::FileCorrupt); + } + #[test] fn version_header_sorts_before_prefers_object_over_delete_marker_on_equal_mod_time() { let object = FileMetaVersionHeader {