From b0f4e66c92a53bc9fe3cfba541a161353a6e8579 Mon Sep 17 00:00:00 2001 From: cxymds Date: Fri, 18 Sep 2026 18:12:51 +0800 Subject: [PATCH] fix(tier): model R2 as unversioned storage (#8004) --- .../src/services/tier/warm_backend_s3.rs | 107 +++++++++++++++--- crates/s3-client/src/provider_versions.rs | 39 ++++++- 2 files changed, 125 insertions(+), 21 deletions(-) diff --git a/crates/ecstore/src/services/tier/warm_backend_s3.rs b/crates/ecstore/src/services/tier/warm_backend_s3.rs index 4e790b560..8303c1105 100644 --- a/crates/ecstore/src/services/tier/warm_backend_s3.rs +++ b/crates/ecstore/src/services/tier/warm_backend_s3.rs @@ -37,7 +37,7 @@ use rustfs_s3_client::{ api_s3_datatypes::ListVersionsResult, credentials::{Credentials, SignatureType, Static, Value}, provider_versions::validate_remote_version_id, - transition_api::{BucketLookupType, Options, TransitionClient, TransitionCore}, + transition_api::{BucketLookupType, ObjectInfo, Options, TransitionClient, TransitionCore}, transition_api::{ReadCloser, ReaderImpl}, }; use rustfs_utils::egress::validate_outbound_url; @@ -278,19 +278,7 @@ impl WarmBackendS3 { let mut stat_opts = GetObjectOptions::default(); stat_opts.version_id.clone_from(&version.version_id); let info = self.client.stat_object(&self.bucket, &remote_object, &stat_opts).await?; - let mut metadata = info.user_metadata; - for (name, value) in &info.metadata { - if (name - .as_str() - .starts_with(rustfs_utils::http::metadata_compat::RUSTFS_INTERNAL_PREFIX) - || name - .as_str() - .starts_with(rustfs_utils::http::metadata_compat::MINIO_INTERNAL_PREFIX)) - && let Ok(value) = value.to_str() - { - metadata.insert(name.as_str().to_string(), value.to_string()); - } - } + let metadata = transition_candidate_metadata(info); if transition_candidate_metadata_matches(&metadata, identity)? { if matched_version.is_some() { return Ok(TransitionCandidateProbe::Ambiguous); @@ -313,6 +301,53 @@ impl WarmBackendS3 { advance_version_markers(&mut key_marker, &mut version_id_marker, &versions)?; } } + + async fn probe_unversioned_transition_candidate_identity( + &self, + object: &str, + identity: TransitionCandidateIdentity, + ) -> Result { + let remote_object = self.get_dest(object); + match self + .client + .stat_object(&self.bucket, &remote_object, &GetObjectOptions::default()) + .await + { + Ok(info) => { + let metadata = transition_candidate_metadata(info); + if transition_candidate_metadata_matches(&metadata, identity)? { + Ok(TransitionCandidateProbe::UnversionedPresent) + } else { + Ok(TransitionCandidateProbe::Unsupported) + } + } + Err(err) => { + let response = to_error_response(&err); + if response.code == S3ErrorCode::NoSuchKey { + Ok(TransitionCandidateProbe::Missing) + } else { + Err(err) + } + } + } + } +} + +fn transition_candidate_metadata(info: ObjectInfo) -> HashMap { + let mut metadata = info.user_metadata; + for (name, value) in &info.metadata { + if (name + .as_str() + .starts_with(rustfs_utils::http::metadata_compat::RUSTFS_INTERNAL_PREFIX) + || name + .as_str() + .starts_with(rustfs_utils::http::metadata_compat::MINIO_INTERNAL_PREFIX)) + && let Ok(value) = value.to_str() + { + metadata.insert(name.as_str().to_string(), value.to_string()); + } + } + metadata } fn transition_candidate_metadata_matches( @@ -525,6 +560,13 @@ mod tests { } async fn scripted_probe_fixture(responses: Vec) -> Option<(WarmBackendS3, tokio::task::JoinHandle>)> { + scripted_probe_fixture_for_tier(responses, "s3").await + } + + async fn scripted_probe_fixture_for_tier( + responses: Vec, + tier_type: &str, + ) -> Option<(WarmBackendS3, tokio::task::JoinHandle>)> { let listener = match tokio::net::TcpListener::bind("127.0.0.1:0").await { Ok(listener) => listener, Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return None, @@ -571,7 +613,7 @@ mod tests { max_retries: 1, ..Default::default() }, - "s3", + tier_type, ) .await .expect("fixture client should build"), @@ -675,6 +717,37 @@ mod tests { assert!(requests[8].to_ascii_lowercase().contains("?versionid=historical-version")); } + #[tokio::test] + async fn unversioned_candidate_reconciler_uses_head_without_version_apis() { + let identity = candidate_identity(); + let transaction_id = identity.transaction_id; + let destination_id = rustfs_utils::crypto::hex(identity.destination_id); + let response = format!( + "HTTP/1.1 200 OK\r\nContent-Length: 0\r\nx-amz-meta-x-rustfs-internal-transition-transaction-id: {transaction_id}\r\nx-amz-meta-x-rustfs-internal-transition-tier-destination-id: {destination_id}\r\nConnection: close\r\n\r\n" + ); + let Some((backend, fixture)) = scripted_probe_fixture_for_tier(vec![response], "r2").await else { + return; + }; + + assert_eq!( + crate::services::tier::warm_backend::TransitionCandidateReconciler::probe_transition_candidate_for( + &backend, + "archive/object", + identity, + ) + .await + .expect("R2 candidate probe should use the current unversioned object"), + TransitionCandidateProbe::UnversionedPresent + ); + let requests = fixture.await.expect("R2 candidate fixture should finish"); + assert_eq!(requests.len(), 1); + let request = requests[0].to_ascii_lowercase(); + assert!(request.starts_with("head /bucket/archive/object")); + assert!(!request.contains("versioning")); + assert!(!request.contains("versions")); + assert!(!request.contains("versionid=")); + } + fn legacy_probe_xml_response(body: &str) -> String { format!( "HTTP/1.1 200 OK\r\nContent-Type: application/xml\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}", @@ -1089,6 +1162,10 @@ impl TransitionCandidateReconciler for WarmBackendS3 { object: &str, identity: TransitionCandidateIdentity, ) -> Result { + let capabilities = self.client.provider_version_capabilities(); + if !capabilities.bucket_versioning_state && !capabilities.list_object_versions { + return self.probe_unversioned_transition_candidate_identity(object, identity).await; + } let bucket_versioning = self.remote_bucket_versioning().await?; self.probe_transition_candidate_identity(object, identity, bucket_versioning) .await diff --git a/crates/s3-client/src/provider_versions.rs b/crates/s3-client/src/provider_versions.rs index 8b36de5fc..ac6dea21c 100644 --- a/crates/s3-client/src/provider_versions.rs +++ b/crates/s3-client/src/provider_versions.rs @@ -83,13 +83,11 @@ impl ProviderVersionCapabilities { if tier_type.eq_ignore_ascii_case("s3") || tier_type.eq_ignore_ascii_case("rustfs") || tier_type.eq_ignore_ascii_case("minio") - || tier_type.eq_ignore_ascii_case("r2") || tier_type.eq_ignore_ascii_case("wasabi") { let list_object_versions = tier_type.eq_ignore_ascii_case("s3") || tier_type.eq_ignore_ascii_case("rustfs") - || tier_type.eq_ignore_ascii_case("minio") - || tier_type.eq_ignore_ascii_case("r2"); + || tier_type.eq_ignore_ascii_case("minio"); Self { raw_version_header: Some(X_AMZ_VERSION_ID), bucket_versioning_state: list_object_versions, @@ -101,6 +99,17 @@ impl ProviderVersionCapabilities { }, exact_get_delete: true, } + } else if tier_type.eq_ignore_ascii_case("r2") { + // R2 exposes an x-amz-version-id on some PUT responses, but it does + // not provide routable S3 object versions. GET/DELETE must therefore + // use the unversioned object path and must not call version APIs. + Self { + raw_version_header: None, + bucket_versioning_state: false, + list_object_versions: false, + conditional_create: ConditionalCreateCapability::IfNoneMatchStar, + exact_get_delete: false, + } } else if tier_type.eq_ignore_ascii_case("aliyun") { Self { raw_version_header: Some(X_OSS_VERSION_ID), @@ -210,8 +219,6 @@ mod tests { ("RustFS", "x-amz-version-id"), ("minio", "x-amz-version-id"), ("MinIO", "x-amz-version-id"), - ("r2", "x-amz-version-id"), - ("R2", "x-amz-version-id"), ("wasabi", "x-amz-version-id"), ("Wasabi", "x-amz-version-id"), ("aliyun", "x-oss-version-id"), @@ -255,6 +262,26 @@ mod tests { ); } + #[test] + fn r2_ignores_non_routable_put_version_header() { + let mut headers = HeaderMap::new(); + headers.insert("x-amz-version-id", HeaderValue::from_static("r2-response-token")); + let capabilities = ProviderVersionCapabilities::for_tier_type("r2"); + + assert_eq!( + capabilities + .raw_version_id(&headers) + .expect("R2 header parsing should succeed"), + None + ); + assert_eq!( + capabilities + .remote_version(&headers, BucketVersioningState::Disabled) + .expect("R2 uses unversioned routing"), + RemoteVersion::Disabled + ); + } + #[test] fn provider_version_missing_header_is_unknown_until_bucket_state_is_known() { let headers = HeaderMap::new(); @@ -280,7 +307,7 @@ mod tests { ("s3", true, true, ConditionalCreateCapability::IfNoneMatchStar, true), ("rustfs", true, true, ConditionalCreateCapability::Unsupported, true), ("minio", true, true, ConditionalCreateCapability::Unsupported, true), - ("r2", true, true, ConditionalCreateCapability::IfNoneMatchStar, true), + ("r2", false, false, ConditionalCreateCapability::IfNoneMatchStar, false), ("wasabi", false, false, ConditionalCreateCapability::Unsupported, true), ("aliyun", false, false, ConditionalCreateCapability::Unsupported, true), ("tencent", false, false, ConditionalCreateCapability::Unsupported, true),