fix(tier): model R2 as unversioned storage (#8004)

This commit is contained in:
cxymds
2026-09-18 18:12:51 +08:00
committed by GitHub
parent e9da1fa074
commit b0f4e66c92
2 changed files with 125 additions and 21 deletions
@@ -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<TransitionCandidateProbe, std::io::Error> {
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<String, String> {
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<String>) -> Option<(WarmBackendS3, tokio::task::JoinHandle<Vec<String>>)> {
scripted_probe_fixture_for_tier(responses, "s3").await
}
async fn scripted_probe_fixture_for_tier(
responses: Vec<String>,
tier_type: &str,
) -> Option<(WarmBackendS3, tokio::task::JoinHandle<Vec<String>>)> {
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<TransitionCandidateProbe, std::io::Error> {
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
+33 -6
View File
@@ -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),