From 3a041e0e3c3e3d40da137f4a2d0bb572054077a6 Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 5 Sep 2026 08:17:38 +0800 Subject: [PATCH] fix(tier): probe exact legacy transition versions --- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 100 ++++++++---------- crates/ecstore/src/services/tier/tier.rs | 13 +++ .../ecstore/src/services/tier/warm_backend.rs | 53 +++++++++- .../src/services/tier/warm_backend_s3.rs | 43 +++++++- 4 files changed, 154 insertions(+), 55 deletions(-) diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index f52c9ea3f..dbb5dd154 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -660,7 +660,6 @@ fn transition_remote_version_delete_plan(oi: &ObjectInfo) -> Result Ok(ResolvedTransitionDeleteVersion { version_id_exact, - verify_missing_after_delete: false, remote_already_missing: false, }), TransitionDeleteVersionPlan::ProbeLegacyUnknown => { - let probe = lease.probe_transition_candidate(&oi.transitioned_object.name).await?; - match (oi.transitioned_object.version_id.as_str(), probe) { - ("", crate::services::tier::warm_backend::TransitionCandidateProbe::UnversionedPresent) => { - Ok(ResolvedTransitionDeleteVersion { - version_id_exact: false, - verify_missing_after_delete: true, - remote_already_missing: false, - }) - } + 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.is_empty() && expected == actual => + if expected == actual => { lease.validate_remote_version_id(expected)?; Ok(ResolvedTransitionDeleteVersion { version_id_exact: true, - verify_missing_after_delete: false, remote_already_missing: false, }) } (_, crate::services::tier::warm_backend::TransitionCandidateProbe::Missing) => { Ok(ResolvedTransitionDeleteVersion { version_id_exact: false, - verify_missing_after_delete: false, remote_already_missing: true, }) } @@ -743,15 +741,6 @@ async fn execute_resolved_transition_delete( ) .await?; } - if resolved.verify_missing_after_delete - && lease.probe_transition_candidate(&oi.transitioned_object.name).await? - != crate::services::tier::warm_backend::TransitionCandidateProbe::Missing - { - return Err(std::io::Error::new( - std::io::ErrorKind::WouldBlock, - "remote tier could not confirm legacy unversioned deletion", - )); - } Ok(()) } @@ -6936,7 +6925,7 @@ mod tests { #[cfg(feature = "test-util")] #[tokio::test] - async fn free_version_delete_probes_legacy_unknown_exact_version_before_remove() { + 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; @@ -6953,6 +6942,17 @@ mod tests { .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 { @@ -6976,13 +6976,13 @@ mod tests { assert_eq!( backend.op_log().await, vec![ - MockWarmOp::Probe { + MockWarmOp::Get { object: remote_object.clone() }, MockWarmOp::Remove { object: remote_object.clone() }, - MockWarmOp::Probe { + MockWarmOp::Get { object: remote_object.clone() }, ] @@ -6995,7 +6995,7 @@ mod tests { #[cfg(feature = "test-util")] #[tokio::test] - async fn free_version_delete_probes_legacy_unknown_unversioned_before_remove() { + 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; @@ -7027,40 +7027,34 @@ mod tests { ..Default::default() }; - super::delete_free_version_remote_object(&object_info, &manager) + let err = super::delete_free_version_remote_object(&object_info, &manager) .await - .expect("probe-proven legacy unversioned cleanup should delete the remote object"); + .expect_err("legacy unversioned cleanup cannot exclude a versioning-state race"); - assert_eq!( - backend.op_log().await, - vec![ - MockWarmOp::Probe { - object: remote_object.clone() - }, - MockWarmOp::Remove { - object: remote_object.clone() - }, - MockWarmOp::Probe { - object: remote_object.clone() - }, - ] - ); - assert_eq!(backend.remove_versions().await, vec![(remote_object, String::new())]); + 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_retains_legacy_unknown_when_probe_disagrees() { + 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 - .set_transition_candidate_probe_override(Some(TransitionCandidateProbe::VersionedPresent( - "different-version".to_string(), - ))) - .await; + .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 { @@ -7075,13 +7069,13 @@ mod tests { ..Default::default() }; - let err = super::delete_free_version_remote_object(&object_info, &manager) + super::delete_free_version_remote_object(&object_info, &manager) .await - .expect_err("legacy unknown cleanup must not delete when the probe disagrees"); + .expect("a missing exact legacy version should be an idempotent cleanup success"); - assert_eq!(err.kind(), std::io::ErrorKind::WouldBlock); - assert_eq!(backend.op_log().await, vec![MockWarmOp::Probe { object: remote_object }]); + 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")] 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 {