From 31dc78eab08a9c41154242877abc4da079f4c9e4 Mon Sep 17 00:00:00 2001 From: houseme Date: Wed, 22 Jul 2026 14:13:32 +0800 Subject: [PATCH] feat(tier): probe transition candidates from providers (#5112) * feat(tier): add transition candidate probe contract Add a fail-closed WarmBackend probe contract for provider-authoritative transition candidate state. Default providers report Unsupported, while the shared mock backend can now model missing, unversioned, and exact-version candidates for follow-up recovery tests. This is a forward-compatible foundation for #1352/#1358 recovery work and does not change production cleanup behavior. Co-Authored-By: heihutu * feat(tier): probe transition candidates from providers Implement provider-authoritative transition candidate probing for S3-family warm backends by querying ListObjectVersions with exact-key filtering and fail-closed classification for delete markers, multiple versions, truncation, and unknown versioning state. This keeps non-S3 providers on the default Unsupported probe result and forwards MinIO, RustFS, and R2 through the S3 probe implementation. Co-Authored-By: heihutu --------- Co-authored-by: heihutu --- crates/ecstore/src/client/api_list.rs | 198 +++++++++++------- crates/ecstore/src/client/api_s3_datatypes.rs | 52 ++--- crates/ecstore/src/services/tier/test_util.rs | 97 ++++++++- crates/ecstore/src/services/tier/tier.rs | 6 +- .../ecstore/src/services/tier/warm_backend.rs | 28 +++ .../src/services/tier/warm_backend_minio.rs | 6 +- .../src/services/tier/warm_backend_r2.rs | 6 +- .../src/services/tier/warm_backend_rustfs.rs | 6 +- .../src/services/tier/warm_backend_s3.rs | 192 ++++++++++++++++- 9 files changed, 484 insertions(+), 107 deletions(-) diff --git a/crates/ecstore/src/client/api_list.rs b/crates/ecstore/src/client/api_list.rs index b1eaf0e56..6bd8591c3 100644 --- a/crates/ecstore/src/client/api_list.rs +++ b/crates/ecstore/src/client/api_list.rs @@ -24,7 +24,7 @@ use crate::client::{ ListBucketResult, ListBucketV2Result, ListMultipartUploadsResult, ListObjectPartsResult, ListVersionsResult, ObjectPart, }, credentials, - transition_api::{ReaderImpl, RequestMetadata, TransitionClient}, + transition_api::{ReaderImpl, RequestMetadata, TransitionClient, collect_response_body}, }; use crate::storage_api_contracts::bucket::BucketInfo; use http::{HeaderMap, StatusCode}; @@ -88,7 +88,7 @@ impl TransitionClient { url_values.insert("max-keys".to_string(), max_keys.to_string()); } - let mut resp = self + let resp = self .execute_method( http::Method::GET, &mut RequestMetadata { @@ -164,7 +164,7 @@ impl TransitionClient { Ok(list_bucket_result) } - pub fn list_object_versions_query( + pub async fn list_object_versions_query( &self, bucket_name: &str, opts: &ListObjectsOptions, @@ -172,93 +172,88 @@ impl TransitionClient { version_id_marker: &str, delimiter: &str, ) -> Result { - /*if err := s3utils.CheckValidBucketName(bucketName); err != nil { - return ListVersionsResult{}, err + let mut url_values = HashMap::new(); + url_values.insert("versions".to_string(), "".to_string()); + url_values.insert("prefix".to_string(), opts.prefix.clone()); + url_values.insert("delimiter".to_string(), delimiter.to_string()); + url_values.insert("encoding-type".to_string(), "url".to_string()); + + if !key_marker.is_empty() { + url_values.insert("key-marker".to_string(), key_marker.to_string()); } - if err := s3utils.CheckValidObjectNamePrefix(opts.Prefix); err != nil { - return ListVersionsResult{}, err - } - urlValues := make(url.Values) - - urlValues.Set("versions", "") - - urlValues.Set("prefix", opts.Prefix) - - urlValues.Set("delimiter", delimiter) - - if keyMarker != "" { - urlValues.Set("key-marker", keyMarker) - } - if opts.max_keys > 0 { - urlValues.Set("max-keys", fmt.Sprintf("%d", opts.max_keys)) + url_values.insert("max-keys".to_string(), opts.max_keys.to_string()); + } + if !version_id_marker.is_empty() { + url_values.insert("version-id-marker".to_string(), version_id_marker.to_string()); + } + if opts.with_metadata { + url_values.insert("metadata".to_string(), "true".to_string()); } - if versionIDMarker != "" { - urlValues.Set("version-id-marker", versionIDMarker) + let mut resp = self + .execute_method( + http::Method::GET, + &mut RequestMetadata { + bucket_name: bucket_name.to_string(), + object_name: "".to_string(), + query_values: url_values, + content_sha256_hex: EMPTY_STRING_SHA256_HASH.to_string(), + custom_header: opts.headers.clone(), + content_body: ReaderImpl::Body(Bytes::new()), + content_length: 0, + content_md5_base64: "".to_string(), + stream_sha256: false, + trailer: HeaderMap::new(), + pre_sign_url: Default::default(), + add_crc: Default::default(), + extra_pre_sign_header: Default::default(), + bucket_location: Default::default(), + expires: Default::default(), + }, + ) + .await?; + + let resp_status = resp.status(); + let headers = resp.headers().clone(); + let body = collect_response_body(resp.into_body(), MAX_S3_CLIENT_RESPONSE_SIZE).await?; + if resp_status != StatusCode::OK { + return Err(std::io::Error::other(http_resp_to_error_response( + resp_status, + &headers, + body, + bucket_name, + "", + ))); } - if opts.WithMetadata { - urlValues.Set("metadata", "true") + let mut versions = quick_xml::de::from_reader::<_, ListVersionsResult>(body.as_slice()) + .map_err(|err| std::io::Error::new(std::io::ErrorKind::InvalidData, err))?; + for version in &mut versions.versions { + version.key = decode_s3_name(&version.key, &versions.encoding_type)?; + } + for marker in &mut versions.delete_markers { + marker.key = decode_s3_name(&marker.key, &versions.encoding_type)?; + } + for prefix in &mut versions.common_prefixes { + prefix.prefix = decode_s3_name(&prefix.prefix, &versions.encoding_type)?; + } + if !versions.next_key_marker.is_empty() { + versions.next_key_marker = decode_s3_name(&versions.next_key_marker, &versions.encoding_type)?; } - urlValues.Set("encoding-type", "url") - - let resp = self.executeMethod(http::Method::GET, &mut RequestMetadata{ - bucketName: bucketName, - queryValues: urlValues, - contentSHA256Hex: emptySHA256Hex, - customHeader: opts.headers, - }).await?; - defer closeResponse(resp) - if err != nil { - return ListVersionsResult{}, err - } - if resp != nil { - if resp.StatusCode != http.StatusOK { - return ListVersionsResult{}, httpRespToErrorResponse(resp, bucketName, "") - } - } - - listObjectVersionsOutput := ListVersionsResult{} - err = xml_decoder(resp.Body, &listObjectVersionsOutput) - if err != nil { - return ListVersionsResult{}, err - } - - for i, obj := range listObjectVersionsOutput.Versions { - listObjectVersionsOutput.Versions[i].Key, err = decode_s3_name(obj.Key, listObjectVersionsOutput.EncodingType) - if err != nil { - return listObjectVersionsOutput, err - } - } - - for i, obj := range listObjectVersionsOutput.CommonPrefixes { - listObjectVersionsOutput.CommonPrefixes[i].Prefix, err = decode_s3_name(obj.Prefix, listObjectVersionsOutput.EncodingType) - if err != nil { - return listObjectVersionsOutput, err - } - } - - if listObjectVersionsOutput.NextKeyMarker != "" { - listObjectVersionsOutput.NextKeyMarker, err = decode_s3_name(listObjectVersionsOutput.NextKeyMarker, listObjectVersionsOutput.EncodingType) - if err != nil { - return listObjectVersionsOutput, err - } - } - - Ok(listObjectVersionsOutput)*/ - Err(std::io::Error::new( - ErrorKind::Unsupported, - credentials::ErrorResponse { + if versions.is_truncated && versions.next_key_marker.is_empty() { + return Err(std::io::Error::other(credentials::ErrorResponse { sts_error: credentials::STSError { r#type: "".to_string(), code: "NotImplemented".to_string(), - message: format!("list_object_versions_query is not implemented for bucket {bucket_name}"), + message: "Truncated ListObjectVersions response should have next key marker set".to_string(), }, request_id: "".to_string(), - }, - )) + })); + } + + Ok(versions) } pub fn list_objects_query( @@ -364,6 +359,7 @@ impl TransitionClient { } } +#[derive(Default)] #[allow(dead_code)] pub struct ListObjectsOptions { reverse_versions: bool, @@ -431,3 +427,51 @@ fn decode_s3_name(name: &str, encoding_type: &str) -> Result + tier-bucket + archive/object + + + 2 + true + archive/object + version-a + + archive/object + version-a + true + 2026-07-22T00:00:00Z + "etag-a" + 5 + STANDARD + + + archive/object + marker-a + false + 2026-07-22T00:00:01Z + + + "#; + + let parsed = + quick_xml::de::from_reader::<_, ListVersionsResult>(xml.as_slice()).expect("ListObjectVersions XML should parse"); + + assert!(parsed.is_truncated); + assert_eq!(parsed.next_key_marker, "archive/object"); + assert_eq!(parsed.next_version_id_marker, "version-a"); + assert_eq!(parsed.versions.len(), 1); + assert_eq!(parsed.versions[0].key, "archive/object"); + assert_eq!(parsed.versions[0].version_id, "version-a"); + assert_eq!(parsed.delete_markers.len(), 1); + assert_eq!(parsed.delete_markers[0].version_id, "marker-a"); + } +} diff --git a/crates/ecstore/src/client/api_s3_datatypes.rs b/crates/ecstore/src/client/api_s3_datatypes.rs index dea2918a7..158637cfd 100644 --- a/crates/ecstore/src/client/api_s3_datatypes.rs +++ b/crates/ecstore/src/client/api_s3_datatypes.rs @@ -55,34 +55,38 @@ pub struct ListBucketV2Result { pub start_after: String, } -#[allow(dead_code)] +#[derive(Debug, Default, Serialize, Deserialize)] +#[serde(default, rename_all = "PascalCase")] pub struct Version { - etag: String, - is_latest: bool, - key: String, - last_modified: OffsetDateTime, - owner: Owner, - size: i64, - storage_class: String, - version_id: String, - user_metadata: HashMap, - user_tags: HashMap, - is_delete_marker: bool, + #[serde(rename = "ETag")] + pub etag: String, + pub is_latest: bool, + pub key: String, + pub size: i64, + pub storage_class: String, + pub version_id: String, + pub user_metadata: HashMap, + pub user_tags: HashMap, } +#[derive(Debug, Default, Serialize, Deserialize)] +#[serde(default, rename_all = "PascalCase")] pub struct ListVersionsResult { - versions: Vec, - common_prefixes: Vec, - name: String, - prefix: String, - delimiter: String, - max_keys: i64, - encoding_type: String, - is_truncated: bool, - key_marker: String, - version_id_marker: String, - next_key_marker: String, - next_version_id_marker: String, + #[serde(rename = "Version")] + pub versions: Vec, + #[serde(rename = "DeleteMarker")] + pub delete_markers: Vec, + pub common_prefixes: Vec, + pub name: String, + pub prefix: String, + pub delimiter: String, + pub max_keys: i64, + pub encoding_type: String, + pub is_truncated: bool, + pub key_marker: String, + pub version_id_marker: String, + pub next_key_marker: String, + pub next_version_id_marker: String, } pub struct ListBucketResult { diff --git a/crates/ecstore/src/services/tier/test_util.rs b/crates/ecstore/src/services/tier/test_util.rs index c69758fe0..d560c12f5 100644 --- a/crates/ecstore/src/services/tier/test_util.rs +++ b/crates/ecstore/src/services/tier/test_util.rs @@ -74,7 +74,9 @@ use crate::disk::format::FormatV3; use crate::disk::{DiskAPI, DiskOption, FORMAT_CONFIG_FILE, RUSTFS_META_BUCKET, STORAGE_FORMAT_FILE, new_disk}; use crate::services::tier::tier::TierConfigMgr; use crate::services::tier::tier_config::{TierConfig, TierMinIO, TierType}; -use crate::services::tier::warm_backend::{WarmBackend, WarmBackendGetOpts, build_transition_put_options}; +use crate::services::tier::warm_backend::{ + TransitionCandidateProbe, WarmBackend, WarmBackendGetOpts, build_transition_put_options, +}; use rustfs_filemeta::FileMeta; use rustfs_utils::path::path_join_buf; @@ -142,6 +144,7 @@ pub enum MockWarmOp { Put { object: String }, Get { object: String }, Remove { object: String }, + Probe { object: String }, ExternalRemove { object: String }, InUse, } @@ -444,6 +447,11 @@ impl MockWarmBackend { self.inner.remove_versions.lock().await.clone() } + /// Return the provider-authoritative view of a remote transition candidate. + pub async fn probe_transition_candidate_state(&self, object: &str) -> Result { + self.probe_transition_candidate(object).await + } + /// Number of `get` calls recorded — useful to assert restore reads hit the /// local copy rather than the remote tier. pub async fn get_count(&self) -> usize { @@ -725,6 +733,23 @@ impl WarmBackend for MockWarmBackend { self.remove(object, rv).await } + async fn probe_transition_candidate(&self, object: &str) -> Result { + self.precondition().await?; + self.record(MockWarmOp::Probe { + object: object.to_string(), + }) + .await; + let objects = self.inner.objects.lock().await; + let Some(stored) = objects.get(object) else { + return Ok(TransitionCandidateProbe::Missing); + }; + if stored.remote_version_id.is_empty() { + Ok(TransitionCandidateProbe::UnversionedPresent) + } else { + Ok(TransitionCandidateProbe::VersionedPresent(stored.remote_version_id.clone())) + } + } + async fn in_use(&self) -> Result { self.precondition().await?; self.record(MockWarmOp::InUse).await; @@ -909,3 +934,73 @@ pub async fn wait_for_free_version_absence(disk_path: &Path, bucket: &str, objec tokio::time::sleep(POLL_INTERVAL).await; } } + +#[cfg(test)] +mod tests { + use super::*; + use bytes::Bytes; + + #[tokio::test] + async fn mock_probe_distinguishes_missing_unversioned_and_versioned_candidates() { + let backend = MockWarmBackend::new(); + + assert_eq!( + backend + .probe_transition_candidate_state("missing") + .await + .expect("probe missing candidate"), + TransitionCandidateProbe::Missing + ); + + backend.set_put_remote_version(Some(String::new())).await; + backend + .put("unversioned", ReaderImpl::Body(Bytes::new()), 0) + .await + .expect("put unversioned candidate"); + assert_eq!( + backend + .probe_transition_candidate_state("unversioned") + .await + .expect("probe unversioned candidate"), + TransitionCandidateProbe::UnversionedPresent + ); + + let remote_version = Uuid::new_v4().to_string(); + backend.set_put_remote_version(Some(remote_version.clone())).await; + backend + .put("versioned", ReaderImpl::Body(Bytes::new()), 0) + .await + .expect("put versioned candidate"); + assert_eq!( + backend + .probe_transition_candidate_state("versioned") + .await + .expect("probe versioned candidate"), + TransitionCandidateProbe::VersionedPresent(remote_version) + ); + + assert_eq!( + backend + .op_log() + .await + .into_iter() + .filter(|op| matches!(op, MockWarmOp::Probe { .. })) + .count(), + 3 + ); + } + + #[tokio::test] + async fn mock_probe_preserves_fault_fail_closed_behavior() { + let backend = MockWarmBackend::new(); + backend.set_reject_credentials(true).await; + + let err = backend + .probe_transition_candidate("remote-object") + .await + .expect_err("credential rejection must fail the authoritative probe"); + + assert_eq!(err.kind(), std::io::ErrorKind::PermissionDenied); + assert!(backend.op_log().await.is_empty()); + } +} diff --git a/crates/ecstore/src/services/tier/tier.rs b/crates/ecstore/src/services/tier/tier.rs index 670ec9d1e..c0df19feb 100644 --- a/crates/ecstore/src/services/tier/tier.rs +++ b/crates/ecstore/src/services/tier/tier.rs @@ -53,7 +53,7 @@ use crate::services::tier::{ tier_admin::TierCreds, tier_config::{TierConfig, TierType, TierWasabi}, tier_handlers::{ERR_TIER_ALREADY_EXISTS, ERR_TIER_NAME_NOT_UPPERCASE, ERR_TIER_NOT_FOUND, ERR_TIER_RESERVED_NAME}, - warm_backend::{WarmBackend, check_warm_backend, new_warm_backend}, + warm_backend::{TransitionCandidateProbe, WarmBackend, check_warm_backend, new_warm_backend}, }; use crate::storage_api_contracts::{ bucket::BucketOperations, @@ -1113,6 +1113,10 @@ impl WarmBackend for SharedWarmBackendProxy { self.0.remove_exact(object, rv).await } + async fn probe_transition_candidate(&self, object: &str) -> io::Result { + self.0.probe_transition_candidate(object).await + } + async fn in_use(&self) -> io::Result { self.0.in_use().await } diff --git a/crates/ecstore/src/services/tier/warm_backend.rs b/crates/ecstore/src/services/tier/warm_backend.rs index 62947cf24..1c667c843 100644 --- a/crates/ecstore/src/services/tier/warm_backend.rs +++ b/crates/ecstore/src/services/tier/warm_backend.rs @@ -64,6 +64,15 @@ pub struct WarmBackendGetOpts { pub length: i64, } +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum TransitionCandidateProbe { + Missing, + UnversionedPresent, + VersionedPresent(String), + Ambiguous, + Unsupported, +} + #[async_trait::async_trait] pub trait WarmBackend { async fn validate(&self) -> Result<(), std::io::Error> { @@ -101,6 +110,9 @@ pub trait WarmBackend { } self.remove(object, rv).await } + async fn probe_transition_candidate(&self, _object: &str) -> Result { + Ok(TransitionCandidateProbe::Unsupported) + } async fn in_use(&self) -> Result; } @@ -615,6 +627,22 @@ mod tests { assert_eq!(removes.load(Ordering::SeqCst), 0); } + #[tokio::test] + async fn default_transition_candidate_probe_is_unsupported() { + let backend = RejectingValidationBackend { + validations: Arc::new(AtomicUsize::new(0)), + puts: Arc::new(AtomicUsize::new(0)), + removes: Arc::new(AtomicUsize::new(0)), + }; + + let probe = backend + .probe_transition_candidate("remote-object") + .await + .expect("default candidate probe should be a safe capability response"); + + assert_eq!(probe, TransitionCandidateProbe::Unsupported); + } + #[tokio::test] async fn check_warm_backend_removes_exact_probe_when_versioning_drifts() { let gets = Arc::new(AtomicUsize::new(0)); diff --git a/crates/ecstore/src/services/tier/warm_backend_minio.rs b/crates/ecstore/src/services/tier/warm_backend_minio.rs index 5d1908e03..fa62df7cb 100644 --- a/crates/ecstore/src/services/tier/warm_backend_minio.rs +++ b/crates/ecstore/src/services/tier/warm_backend_minio.rs @@ -29,7 +29,7 @@ use crate::client::{ }; use crate::services::tier::{ tier_config::TierMinIO, - warm_backend::{WarmBackend, WarmBackendGetOpts, build_transition_put_options}, + warm_backend::{TransitionCandidateProbe, WarmBackend, WarmBackendGetOpts, build_transition_put_options}, warm_backend_s3::WarmBackendS3, }; use tracing::warn; @@ -126,6 +126,10 @@ impl WarmBackend for WarmBackendMinIO { self.0.remove(object, rv).await } + async fn probe_transition_candidate(&self, object: &str) -> Result { + self.0.probe_transition_candidate(object).await + } + async fn in_use(&self) -> Result { self.0.in_use().await } diff --git a/crates/ecstore/src/services/tier/warm_backend_r2.rs b/crates/ecstore/src/services/tier/warm_backend_r2.rs index 3ced55197..e93c1ad9c 100644 --- a/crates/ecstore/src/services/tier/warm_backend_r2.rs +++ b/crates/ecstore/src/services/tier/warm_backend_r2.rs @@ -29,7 +29,7 @@ use crate::client::{ }; use crate::services::tier::{ tier_config::TierR2, - warm_backend::{WarmBackend, WarmBackendGetOpts, build_transition_put_options}, + warm_backend::{TransitionCandidateProbe, WarmBackend, WarmBackendGetOpts, build_transition_put_options}, warm_backend_s3::WarmBackendS3, }; use tracing::warn; @@ -126,6 +126,10 @@ impl WarmBackend for WarmBackendR2 { self.0.remove(object, rv).await } + async fn probe_transition_candidate(&self, object: &str) -> Result { + self.0.probe_transition_candidate(object).await + } + async fn in_use(&self) -> Result { self.0.in_use().await } diff --git a/crates/ecstore/src/services/tier/warm_backend_rustfs.rs b/crates/ecstore/src/services/tier/warm_backend_rustfs.rs index 65c751968..db03c9041 100644 --- a/crates/ecstore/src/services/tier/warm_backend_rustfs.rs +++ b/crates/ecstore/src/services/tier/warm_backend_rustfs.rs @@ -29,7 +29,7 @@ use crate::client::{ }; use crate::services::tier::{ tier_config::TierRustFS, - warm_backend::{WarmBackend, WarmBackendGetOpts, build_transition_put_options}, + warm_backend::{TransitionCandidateProbe, WarmBackend, WarmBackendGetOpts, build_transition_put_options}, warm_backend_s3::WarmBackendS3, }; @@ -123,6 +123,10 @@ impl WarmBackend for WarmBackendRustFS { self.0.remove(object, rv).await } + async fn probe_transition_candidate(&self, object: &str) -> Result { + self.0.probe_transition_candidate(object).await + } + async fn in_use(&self) -> Result { self.0.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 9b464ea8e..e37a00b7c 100644 --- a/crates/ecstore/src/services/tier/warm_backend_s3.rs +++ b/crates/ecstore/src/services/tier/warm_backend_s3.rs @@ -24,8 +24,10 @@ use url::Url; use crate::client::{ api_get_options::GetObjectOptions, + api_list::ListObjectsOptions, api_put_object::PutObjectOptions, api_remove::{RemoveObjectOptions, RemoveObjectResult}, + api_s3_datatypes::ListVersionsResult, credentials::{Credentials, SignatureType, Static, Value}, transition_api::{BucketLookupType, Options, TransitionClient, TransitionCore}, transition_api::{ReadCloser, ReaderImpl}, @@ -34,11 +36,12 @@ use crate::error::ErrorResponse; use crate::error::error_resp_to_object_err; use crate::services::tier::{ tier_config::TierS3, - warm_backend::{WarmBackend, WarmBackendGetOpts, build_transition_put_options}, + warm_backend::{TransitionCandidateProbe, WarmBackend, WarmBackendGetOpts, build_transition_put_options}, }; use http::HeaderMap; use rustfs_utils::egress::validate_outbound_url; use rustfs_utils::path::SLASH_SEPARATOR; +use s3s::dto::BucketVersioningStatus; pub struct WarmBackendS3 { pub client: Arc, @@ -48,6 +51,27 @@ pub struct WarmBackendS3 { pub storage_class: String, } +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum RemoteBucketVersioning { + Disabled, + Suspended, + Enabled, +} + +fn remote_bucket_versioning_from_status(status: Option<&str>) -> Result { + Ok(match status { + Some(BucketVersioningStatus::ENABLED) => RemoteBucketVersioning::Enabled, + Some(BucketVersioningStatus::SUSPENDED) => RemoteBucketVersioning::Suspended, + Some(status) => { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + format!("remote tier bucket returned unsupported versioning status {status}"), + )); + } + None => RemoteBucketVersioning::Disabled, + }) +} + impl WarmBackendS3 { pub async fn new(conf: &TierS3, _tier: &str) -> Result { Self::new_with_bucket_lookup(conf, BucketLookupType::BucketLookupAuto, "s3").await @@ -159,11 +183,57 @@ impl WarmBackendS3 { let (_, headers, reader) = self.core.get_object(&self.bucket, &self.get_dest(object), &gopts).await?; Ok((headers, reader)) } + + async fn remote_bucket_versioning(&self) -> Result { + let config = self.client.get_bucket_versioning(&self.bucket).await?; + remote_bucket_versioning_from_status(config.status.as_ref().map(|status| status.as_str())) + } + + async fn list_transition_candidate_versions(&self, object: &str) -> Result { + let mut opts = ListObjectsOptions::default(); + opts.set("prefix", &self.get_dest(object)); + opts.set("max-keys", "2"); + self.client.list_object_versions_query(&self.bucket, &opts, "", "", "").await + } +} + +fn classify_transition_candidate_versions( + remote_object: &str, + bucket_versioning: RemoteBucketVersioning, + versions: &ListVersionsResult, +) -> TransitionCandidateProbe { + if versions.is_truncated { + return TransitionCandidateProbe::Ambiguous; + } + + if versions.delete_markers.iter().any(|marker| marker.key == remote_object) { + return TransitionCandidateProbe::Ambiguous; + } + + let mut exact_versions = versions.versions.iter().filter(|version| version.key == remote_object); + let Some(version) = exact_versions.next() else { + return TransitionCandidateProbe::Missing; + }; + if exact_versions.next().is_some() { + return TransitionCandidateProbe::Ambiguous; + } + + match bucket_versioning { + RemoteBucketVersioning::Disabled => TransitionCandidateProbe::UnversionedPresent, + RemoteBucketVersioning::Suspended if version.version_id == "null" => { + TransitionCandidateProbe::VersionedPresent(version.version_id.clone()) + } + RemoteBucketVersioning::Suspended | RemoteBucketVersioning::Enabled if !version.version_id.is_empty() => { + TransitionCandidateProbe::VersionedPresent(version.version_id.clone()) + } + RemoteBucketVersioning::Suspended | RemoteBucketVersioning::Enabled => TransitionCandidateProbe::Ambiguous, + } } #[cfg(test)] mod tests { use super::*; + use crate::client::api_s3_datatypes::{ListVersionsResult, Version}; #[tokio::test] async fn new_rejects_loopback_endpoint_before_network_setup() { @@ -181,6 +251,116 @@ mod tests { Err(err) => assert!(err.to_string().contains("not allowed")), } } + + fn list_versions(versions: &[(&str, &str)], delete_markers: &[(&str, &str)], is_truncated: bool) -> ListVersionsResult { + ListVersionsResult { + versions: versions + .iter() + .map(|(key, version_id)| Version { + key: (*key).to_string(), + version_id: (*version_id).to_string(), + ..Default::default() + }) + .collect(), + delete_markers: delete_markers + .iter() + .map(|(key, version_id)| Version { + key: (*key).to_string(), + version_id: (*version_id).to_string(), + ..Default::default() + }) + .collect(), + is_truncated, + ..Default::default() + } + } + + #[test] + fn transition_candidate_probe_classifier_is_fail_closed() { + assert_eq!( + classify_transition_candidate_versions( + "archive/object", + RemoteBucketVersioning::Disabled, + &list_versions(&[], &[], false), + ), + TransitionCandidateProbe::Missing + ); + assert_eq!( + classify_transition_candidate_versions( + "archive/object", + RemoteBucketVersioning::Disabled, + &list_versions(&[("archive/object", "")], &[], false), + ), + TransitionCandidateProbe::UnversionedPresent + ); + assert_eq!( + classify_transition_candidate_versions( + "archive/object", + RemoteBucketVersioning::Enabled, + &list_versions(&[("archive/object", "version-a")], &[], false), + ), + TransitionCandidateProbe::VersionedPresent("version-a".to_string()) + ); + assert_eq!( + classify_transition_candidate_versions( + "archive/object", + RemoteBucketVersioning::Suspended, + &list_versions(&[("archive/object", "null")], &[], false), + ), + TransitionCandidateProbe::VersionedPresent("null".to_string()) + ); + assert_eq!( + classify_transition_candidate_versions( + "archive/object", + RemoteBucketVersioning::Enabled, + &list_versions(&[("archive/object", "")], &[], false), + ), + TransitionCandidateProbe::Ambiguous + ); + assert_eq!( + classify_transition_candidate_versions( + "archive/object", + RemoteBucketVersioning::Enabled, + &list_versions(&[("archive/object", "version-a"), ("archive/object", "version-b")], &[], false), + ), + TransitionCandidateProbe::Ambiguous + ); + assert_eq!( + classify_transition_candidate_versions( + "archive/object", + RemoteBucketVersioning::Enabled, + &list_versions(&[("archive/object", "version-a")], &[("archive/object", "marker-a")], false), + ), + TransitionCandidateProbe::Ambiguous + ); + assert_eq!( + classify_transition_candidate_versions( + "archive/object", + RemoteBucketVersioning::Enabled, + &list_versions(&[("archive/object", "version-a")], &[], true), + ), + TransitionCandidateProbe::Ambiguous + ); + } + + #[test] + fn remote_bucket_versioning_status_parser_fails_closed() { + assert_eq!( + remote_bucket_versioning_from_status(None).expect("absent status means disabled"), + RemoteBucketVersioning::Disabled + ); + assert_eq!( + remote_bucket_versioning_from_status(Some(BucketVersioningStatus::ENABLED)).expect("enabled status should parse"), + RemoteBucketVersioning::Enabled + ); + assert_eq!( + remote_bucket_versioning_from_status(Some(BucketVersioningStatus::SUSPENDED)).expect("suspended status should parse"), + RemoteBucketVersioning::Suspended + ); + let err = remote_bucket_versioning_from_status(Some("UnexpectedStatus")) + .expect_err("unknown versioning status must fail closed"); + assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + } } #[async_trait::async_trait] @@ -215,6 +395,16 @@ impl WarmBackend for WarmBackendS3 { self.remove_with_result(object, rv).await.map(|_| ()) } + async fn probe_transition_candidate(&self, object: &str) -> Result { + let bucket_versioning = self.remote_bucket_versioning().await?; + let versions = self.list_transition_candidate_versions(object).await?; + Ok(classify_transition_candidate_versions( + &self.get_dest(object), + bucket_versioning, + &versions, + )) + } + async fn in_use(&self) -> Result { let result = self .core