feat(replication): proxy GET/HEAD/Tagging for unreplicated objects to replication targets

Implements the MinIO active-active read-proxy protocol (P1-5 of the
replication compatibility review): when a GET/HEAD/GetObjectTagging/
PutObjectTagging/DeleteObjectTagging request fails locally with
not-found and the bucket has replication targets, the request is proxied
to the targets in rule order, mirroring bucket-replication.go
proxyGetToReplicationTarget/proxyHeadToRepTarget/proxyTaggingToRepTarget.

Protocol surface:
- Anti-loop: inbound {x-rustfs-,x-minio-}source-proxy-request is parsed
  into ObjectOptions (proxy_request + proxy_header_set, matching MinIO
  ProxyRequest/ProxyHeaderSet); a request carrying the marker with ANY
  value is never re-proxied. Outbound client proxy calls send the marker
  as "true"; replication worker convergence HEADs send it as "false" so
  a peer's proxy layer cannot answer a convergence check by proxying
  back to the source (which would fake Completed without a PUT).
- Target selection: new replication_proxy.rs get_proxy_targets — empty
  when the marker is set, versioning is suspended, or no replication
  config; otherwise filter_target_arns -> TargetClient lookup, skipping
  targets with proxying disabled.
- TargetClient gains head_object_for_proxy/get_object (streaming) and
  the three tagging calls. Proxy calls never send the replication-check
  SSE-C exemption header; customer SSE-C keys are forwarded verbatim so
  the target performs real decryption. Conditional (If-*) headers are
  not forwarded (MinIO parity); Range and part_number are, with
  parts_count/tag_count/storage_class/expiration passed through.
- Metrics: proxy counters now count only real client proxy traffic,
  MinIO-aligned (one total per proxied request, one failed when no
  target served it). The previous misattributed counters — replication
  worker HEAD/PUT (#2672) and local tagging operations (#2682) — are
  removed; ReplProxyMetric now maps the tagging counters instead of
  dropping them.

e2e (fake_s3_target extended with tagging + header journaling): proxied
GET body + outbound header contract (marker present, no
replication-check, SSE-C passthrough), HEAD, anti-loop 404 with zero
outbound requests, GetObjectTagging, and metric mapping unit tests.

Rolling note: proxying only activates for buckets with replication
targets; requests carrying the marker keep pre-upgrade behavior.

Refs rustfs/backlog#1675 (P1-5)
This commit is contained in:
唐小鸭
2026-08-17 18:33:42 +08:00
parent 21c2fb42bb
commit 9baa92563a
19 changed files with 1471 additions and 160 deletions
+2 -2
View File
@@ -4,8 +4,8 @@ This module is the shared failure-injection boundary for replication end-to-end
`FakeS3Target::start()` creates the listener. Add target buckets with `create_bucket`, point a RustFS remote target at `address()`, use `FAKE_ACCESS_KEY` / `FAKE_SECRET_KEY`, then enqueue per-operation faults with `inject`. Faults for one operation are consumed in FIFO order and do not consume faults queued for another operation. A fault is consumed only after `s3s` verifies the full request signature, so anonymous, other-access-key, and bad-signature traffic cannot disturb a script.
Supported data operations are HeadBucket, GetBucketVersioning, PUT/GET/HEAD/DELETE Object, and create/upload/complete/abort multipart upload. `create_bucket` models general-purpose buckets in S3's shared global namespace; account-regional namespace buckets and their `-an` names are intentionally out of scope. Buckets are versioned: PUT creates a version, DELETE without `versionId` creates a delete marker, and DELETE with `versionId` removes exactly that version. Internal source version IDs must be UUIDs and are stored canonically. Source mtime is honored only for source-replication PUT/DELETE requests; absent or invalid values use receipt time, matching RustFS, while multipart completion always uses receipt time. Replicated versions are ordered newest-first by source mtime so late older versions and delete markers do not become current. Equal mtimes prefer objects over delete markers, then canonical UUID order; RustFS's internal FileMeta signature tie-break is intentionally out of scope because it is not part of the target S3 protocol. Multipart part numbers follow S3's `1..=10000` range, and every completed part except the final part must be at least 5 MiB.
Supported data operations are HeadBucket, GetBucketVersioning, PUT/GET/HEAD/DELETE Object, Get/Put/Delete ObjectTagging (tags live per version; Put replaces the whole set, Delete clears it), and create/upload/complete/abort multipart upload. `create_bucket` models general-purpose buckets in S3's shared global namespace; account-regional namespace buckets and their `-an` names are intentionally out of scope. Buckets are versioned: PUT creates a version, DELETE without `versionId` creates a delete marker, and DELETE with `versionId` removes exactly that version. Internal source version IDs must be UUIDs and are stored canonically. Source mtime is honored only for source-replication PUT/DELETE requests; absent or invalid values use receipt time, matching RustFS, while multipart completion always uses receipt time. Replicated versions are ordered newest-first by source mtime so late older versions and delete markers do not become current. Equal mtimes prefer objects over delete markers, then canonical UUID order; RustFS's internal FileMeta signature tie-break is intentionally out of scope because it is not part of the target S3 protocol. Multipart part numbers follow S3's `1..=10000` range, and every completed part except the final part must be at least 5 MiB.
Fault actions cover HTTP 401/403/503 responses, pre-dispatch delay, connection abort when a logical request-body threshold is reached, streaming slow drain, and a deliberately wrong response ETag (including multipart-complete XML). `requests()` returns the ordered, credential-free request journal for assertions.
Fault actions cover HTTP 401/403/503 responses, pre-dispatch delay, connection abort when a logical request-body threshold is reached, streaming slow drain, and a deliberately wrong response ETag (including multipart-complete XML). `requests()` returns the ordered, credential-free request journal for assertions. Each record also journals a `ProxyHeaderSnapshot` — the read-proxy anti-loop marker (`x-{rustfs,minio}-source-proxy-request`), the replication-check exemption header, and the client SSE-C header family (algorithm and key-MD5 values; for the key itself only its presence) — so proxy tests can pin the exact wire contract.
The listener is loopback-only. It admits at most 64 active connections and two concurrently buffered request bodies; authenticated multipart-complete XML collection and assembly take both body permits. Keep-alive is disabled, request-header reads are bounded to 30 seconds, a parsed request is bounded to 65 seconds, and the complete connection lifetime is bounded to 100 seconds. It retains at most 256 buckets, 4,096 journal entries, 4,096 scripted faults, 4,096 object versions, 256 multipart uploads, and 10,000 multipart parts. Retained identifiers are capped at 1 KiB, user metadata at 2 KiB, and content type at 1 KiB. A PUT or uploaded part is capped at 64 MiB; a completed multipart object and all stored object/part data are capped at 128 MiB. Body drain, body-permit waits, delay, and slow-drain execution are bounded to 30 seconds; each slow-drain slice delay must be below that bound.
+155 -3
View File
@@ -30,10 +30,12 @@ use s3s::access::{S3Access, S3AccessContext};
use s3s::auth::SimpleAuth;
use s3s::dto::{
AbortMultipartUploadInput, AbortMultipartUploadOutput, CompleteMultipartUploadInput, CompleteMultipartUploadOutput,
CreateMultipartUploadInput, CreateMultipartUploadOutput, DeleteMarkerEntry, DeleteObjectInput, DeleteObjectOutput, ETag,
GetBucketVersioningInput, GetBucketVersioningOutput, GetObjectInput, GetObjectOutput, HeadBucketInput, HeadBucketOutput,
CreateMultipartUploadInput, CreateMultipartUploadOutput, DeleteMarkerEntry, DeleteObjectInput, DeleteObjectOutput,
DeleteObjectTaggingInput, DeleteObjectTaggingOutput, ETag, GetBucketVersioningInput, GetBucketVersioningOutput,
GetObjectInput, GetObjectOutput, GetObjectTaggingInput, GetObjectTaggingOutput, HeadBucketInput, HeadBucketOutput,
HeadObjectInput, HeadObjectOutput, ListObjectVersionsInput, ListObjectVersionsOutput, ObjectVersionId, PutObjectInput,
PutObjectOutput, StreamingBlob, Timestamp, TimestampFormat, UploadPartInput, UploadPartOutput,
PutObjectOutput, PutObjectTaggingInput, PutObjectTaggingOutput, StreamingBlob, Tag, TagSet, Timestamp, TimestampFormat,
UploadPartInput, UploadPartOutput,
};
use s3s::service::{S3Service, S3ServiceBuilder};
use s3s::validation::{AwsNameValidation, NameValidation};
@@ -103,6 +105,9 @@ pub enum Operation {
GetObject,
HeadObject,
DeleteObject,
GetObjectTagging,
PutObjectTagging,
DeleteObjectTagging,
ListObjectVersions,
CreateMultipartUpload,
UploadPart,
@@ -149,6 +154,35 @@ impl ReplicationTimestampHeaders {
}
}
/// Read-proxy related headers observed on a request, journaled so proxy
/// tests can assert the exact wire contract: the anti-loop marker present,
/// the replication-check exemption absent, and the client SSE-C key family
/// forwarded verbatim. The SSE-C key value itself is never retained — only
/// its presence.
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ProxyHeaderSnapshot {
pub source_proxy_request: Option<String>,
pub replication_check: Option<String>,
pub ssec_algorithm: Option<String>,
pub ssec_key_present: bool,
pub ssec_key_md5: Option<String>,
}
impl ProxyHeaderSnapshot {
fn from_headers(headers: &HeaderMap) -> Self {
Self {
source_proxy_request: header_value(headers, &["x-rustfs-source-proxy-request", "x-minio-source-proxy-request"])
.map(bounded_journal_value),
replication_check: header_value(headers, &["x-rustfs-source-replication-check", "x-minio-source-replication-check"])
.map(bounded_journal_value),
ssec_algorithm: header_value(headers, &["x-amz-server-side-encryption-customer-algorithm"])
.map(bounded_journal_value),
ssec_key_present: headers.contains_key("x-amz-server-side-encryption-customer-key"),
ssec_key_md5: header_value(headers, &["x-amz-server-side-encryption-customer-key-md5"]).map(bounded_journal_value),
}
}
}
/// Credential-free request metadata retained for deterministic assertions.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RequestRecord {
@@ -163,6 +197,7 @@ pub struct RequestRecord {
pub content_length: Option<u64>,
pub consumed_bytes: Option<usize>,
pub replication_timestamps: ReplicationTimestampHeaders,
pub proxy_headers: ProxyHeaderSnapshot,
pub fault: Option<FaultAction>,
}
@@ -199,6 +234,9 @@ struct ObjectVersion {
delete_marker: bool,
content_type: Option<String>,
metadata: Option<HashMap<String, String>>,
/// Object tags as ordered key/value pairs (PutObjectTagging replaces the
/// whole set, DeleteObjectTagging clears it).
tags: Vec<(String, String)>,
}
#[derive(Clone)]
@@ -569,6 +607,7 @@ impl S3Access for FaultAccess {
.and_then(|value| value.to_str().ok())
.and_then(|value| value.parse().ok());
let replication_timestamps = ReplicationTimestampHeaders::from_headers(context.headers());
let proxy_headers = ProxyHeaderSnapshot::from_headers(context.headers());
let fault = record_request(
&self.control,
operation,
@@ -576,6 +615,7 @@ impl S3Access for FaultAccess {
parsed,
content_length,
replication_timestamps,
proxy_headers,
);
if let Some(RequestFault {
action: FaultAction::Status(status),
@@ -615,6 +655,9 @@ fn operation_from_s3_name(name: &str) -> Operation {
"GetObject" => Operation::GetObject,
"HeadObject" => Operation::HeadObject,
"DeleteObject" => Operation::DeleteObject,
"GetObjectTagging" => Operation::GetObjectTagging,
"PutObjectTagging" => Operation::PutObjectTagging,
"DeleteObjectTagging" => Operation::DeleteObjectTagging,
"CreateMultipartUpload" => Operation::CreateMultipartUpload,
"UploadPart" => Operation::UploadPart,
"CompleteMultipartUpload" => Operation::CompleteMultipartUpload,
@@ -630,6 +673,7 @@ fn record_request(
parsed: ParsedRequest,
content_length: Option<u64>,
replication_timestamps: ReplicationTimestampHeaders,
proxy_headers: ProxyHeaderSnapshot,
) -> Option<RequestFault> {
let mut state = lock(control);
let action = parsed
@@ -655,6 +699,7 @@ fn record_request(
content_length,
consumed_bytes: None,
replication_timestamps,
proxy_headers,
fault: action.clone(),
});
action.map(|action| RequestFault { sequence, action })
@@ -721,6 +766,15 @@ fn parse_request(method: &Method, uri: &Uri) -> ParsedRequest {
(&Method::POST, true) if query.contains_key("uploads") => Operation::CreateMultipartUpload,
(&Method::POST, true) if upload_id.is_some() => Operation::CompleteMultipartUpload,
(&Method::DELETE, true) if upload_id.is_some() => Operation::AbortMultipartUpload,
(&Method::GET, true) if query.contains_key("tagging") && only_query_keys(&["tagging", "versionId"]) => {
Operation::GetObjectTagging
}
(&Method::PUT, true) if query.contains_key("tagging") && only_query_keys(&["tagging", "versionId"]) => {
Operation::PutObjectTagging
}
(&Method::DELETE, true) if query.contains_key("tagging") && only_query_keys(&["tagging", "versionId"]) => {
Operation::DeleteObjectTagging
}
// A replication PUT addresses the source version via `?versionId=`.
(&Method::PUT, true) if only_query_keys(&["versionId"]) => Operation::PutObject,
(&Method::GET, true) if only_query_keys(&["versionId"]) => Operation::GetObject,
@@ -1135,6 +1189,33 @@ fn find_version(state: &StoreState, bucket: &str, key: &str, version_id: Option<
Ok(version.clone())
}
/// Replace (or clear, with an empty vec) the tag set of the addressed
/// version, returning its version id. Mirrors `find_version` addressing:
/// explicit version id or the latest version, delete markers rejected.
fn set_version_tags(
state: &mut StoreState,
bucket: &str,
key: &str,
version_id: Option<&str>,
tags: Vec<(String, String)>,
) -> S3Result<String> {
// Resolve first (immutable) so the error paths match find_version.
let resolved = find_version(state, bucket, key, version_id)?.version_id;
let versions = state
.buckets
.get_mut(bucket)
.expect("bucket existence checked by find_version")
.objects
.get_mut(key)
.expect("key existence checked by find_version");
let version = versions
.iter_mut()
.find(|version| version.version_id == resolved)
.expect("version existence checked by find_version");
version.tags = tags;
Ok(resolved)
}
#[async_trait]
impl S3 for FakeBackend {
async fn head_bucket(&self, req: S3Request<HeadBucketInput>) -> S3Result<S3Response<HeadBucketOutput>> {
@@ -1248,6 +1329,7 @@ impl S3 for FakeBackend {
delete_marker: false,
content_type: input.content_type,
metadata: input.metadata,
tags: Vec::new(),
};
upsert_version(&mut lock(&self.store), &input.bucket, input.key, version)?;
Ok(apply_response_fault(
@@ -1305,6 +1387,72 @@ impl S3 for FakeBackend {
))
}
async fn get_object_tagging(&self, req: S3Request<GetObjectTaggingInput>) -> S3Result<S3Response<GetObjectTaggingOutput>> {
let fault = request_fault(&req);
apply_non_body_fault(fault.as_ref(), &self.control).await?;
let input = req.input;
let version = {
let state = lock(&self.store);
find_version(&state, &input.bucket, &input.key, input.version_id.as_deref())?
};
let tag_set: TagSet = version
.tags
.into_iter()
.map(|(key, value)| Tag {
key: Some(key),
value: Some(value),
})
.collect();
Ok(apply_response_fault(
S3Response::new(GetObjectTaggingOutput {
tag_set,
version_id: Some(ObjectVersionId::from(version.version_id)),
}),
fault.as_ref(),
))
}
async fn put_object_tagging(&self, req: S3Request<PutObjectTaggingInput>) -> S3Result<S3Response<PutObjectTaggingOutput>> {
let fault = request_fault(&req);
apply_non_body_fault(fault.as_ref(), &self.control).await?;
let input = req.input;
let tags = input
.tagging
.tag_set
.into_iter()
.map(|tag| (tag.key.unwrap_or_default(), tag.value.unwrap_or_default()))
.collect();
let version_id = {
let mut state = lock(&self.store);
set_version_tags(&mut state, &input.bucket, &input.key, input.version_id.as_deref(), tags)?
};
Ok(apply_response_fault(
S3Response::new(PutObjectTaggingOutput {
version_id: Some(ObjectVersionId::from(version_id)),
}),
fault.as_ref(),
))
}
async fn delete_object_tagging(
&self,
req: S3Request<DeleteObjectTaggingInput>,
) -> S3Result<S3Response<DeleteObjectTaggingOutput>> {
let fault = request_fault(&req);
apply_non_body_fault(fault.as_ref(), &self.control).await?;
let input = req.input;
let version_id = {
let mut state = lock(&self.store);
set_version_tags(&mut state, &input.bucket, &input.key, input.version_id.as_deref(), Vec::new())?
};
Ok(apply_response_fault(
S3Response::new(DeleteObjectTaggingOutput {
version_id: Some(ObjectVersionId::from(version_id)),
}),
fault.as_ref(),
))
}
async fn delete_object(&self, req: S3Request<DeleteObjectInput>) -> S3Result<S3Response<DeleteObjectOutput>> {
let fault = request_fault(&req);
apply_non_body_fault(fault.as_ref(), &self.control).await?;
@@ -1381,6 +1529,7 @@ impl S3 for FakeBackend {
delete_marker: true,
content_type: None,
metadata: None,
tags: Vec::new(),
},
)?;
Ok(apply_response_fault(
@@ -1583,6 +1732,7 @@ impl S3 for FakeBackend {
delete_marker: false,
content_type: upload.content_type,
metadata: upload.metadata,
tags: Vec::new(),
};
let mut state = lock(&self.store);
let current = state
@@ -3074,6 +3224,7 @@ mod tests {
},
Some(0),
ReplicationTimestampHeaders::default(),
ProxyHeaderSnapshot::default(),
);
}
let records = lock(&control).requests.clone();
@@ -3096,6 +3247,7 @@ mod tests {
},
None,
ReplicationTimestampHeaders::default(),
ProxyHeaderSnapshot::default(),
);
{
let bounded_records = lock(&bounded_control);
@@ -8417,3 +8417,304 @@ async fn test_scanner_never_compensates_when_existing_object_replication_disable
Ok(())
}
/// Shared setup for the P1-5 read-proxy scenarios (backlog#1675): a RustFS
/// source with an enabled replication rule pointing at the fake target, and
/// an object seeded DIRECTLY on the target — it exists remotely but not
/// locally, exactly the active-active replication-lag window the read proxy
/// serves.
async fn start_read_proxy_lab(
source_bucket: &str,
target_bucket: &str,
) -> Result<(FakeS3Target, RustFSTestEnvironment, Client, Client), Box<dyn Error + Send + Sync>> {
let target = FakeS3Target::start().await?;
target.create_bucket(target_bucket);
target.assign_own_version_ids(true);
let mut source_env = RustFSTestEnvironment::new().await?;
let mut process_env = replication_fast_env();
process_env.extend_from_slice(LOOPBACK_REPLICATION_TARGET_ENV);
process_env.extend_from_slice(&[
("NO_PROXY", "127.0.0.1,localhost"),
("HTTP_PROXY", ""),
("HTTPS_PROXY", ""),
("RUST_LOG", "error"),
]);
source_env.start_rustfs_server_with_env(vec![], &process_env).await?;
let source_client = source_env.create_s3_client();
source_client.create_bucket().bucket(source_bucket).send().await?;
enable_bucket_versioning(&source_env, source_bucket).await?;
let target_arn = set_replication_target_with_options(
&source_env,
source_bucket,
ReplicationTargetOptions {
endpoint: &target.address(),
access_key: FAKE_ACCESS_KEY,
secret_key: FAKE_SECRET_KEY,
target_bucket,
secure: false,
skip_tls_verify: false,
ca_cert_pem: None,
},
)
.await?;
put_bucket_replication(&source_env, source_bucket, &target_arn).await?;
let target_client = Client::from_conf(crate::common::build_test_s3_config(
target.endpoint(),
FAKE_ACCESS_KEY,
FAKE_SECRET_KEY,
None,
"read-proxy-e2e",
));
Ok((target, source_env, source_client, target_client))
}
/// P1-5 (backlog#1675): during the active-active replication lag window a
/// GET/HEAD for an object the local site does not have yet is proxied to the
/// replication target. Pins the wire contract: the anti-loop
/// `source-proxy-request` marker is sent, the replication worker's
/// `source-replication-check` SSE-C exemption is NEVER sent, client SSE-C
/// headers are forwarded verbatim, and an inbound request that was itself
/// proxied is answered locally (404) without touching the target.
#[tokio::test]
#[serial]
async fn test_get_and_head_proxy_unreplicated_object_to_replication_target() -> TestResult {
init_logging();
let source_bucket = "proxy-read-src";
let target_bucket = "proxy-read-dst";
let (target, source_env, source_client, target_client) = start_read_proxy_lab(source_bucket, target_bucket).await?;
let payload = b"proxy payload".to_vec();
target_client
.put_object()
.bucket(target_bucket)
.key("proxy-only")
.body(ByteStream::from(payload.clone()))
.send()
.await?;
target.take_requests();
// a. GET of the locally-missing object is served through the proxy.
let got = source_client
.get_object()
.bucket(source_bucket)
.key("proxy-only")
.send()
.await
.map_err(|err| format!("proxied GET failed: {}", err.into_service_error()))?;
assert_eq!(got.content_length, Some(payload.len() as i64));
let body = got.body.collect().await?.into_bytes();
assert_eq!(body.as_ref(), payload.as_slice(), "proxied GET must stream the target's body");
let get_record = target
.requests()
.into_iter()
.find(|record| record.operation == FakeTargetOperation::GetObject && record.key.as_deref() == Some("proxy-only"))
.ok_or("fake target never received the proxied GET")?;
assert_eq!(
get_record.proxy_headers.source_proxy_request.as_deref(),
Some("true"),
"proxied GET must carry the anti-loop source-proxy-request marker"
);
assert!(
get_record.proxy_headers.replication_check.is_none(),
"proxied GET must never carry the replication worker's source-replication-check exemption"
);
assert!(
get_record.proxy_headers.ssec_algorithm.is_none() && !get_record.proxy_headers.ssec_key_present,
"no client SSE-C headers were sent, so none may be forwarded"
);
// a2. Client SSE-C headers travel verbatim to the target (the target owns
// the real SSE-C decryption; the plaintext fake simply ignores them).
target.take_requests();
let ssec_key = "01234567890123456789012345678901";
let ssec_key_b64 = BASE64_STANDARD.encode(ssec_key);
let ssec_key_md5 = sse_customer_key_md5_base64(ssec_key);
let _ = source_client
.get_object()
.bucket(source_bucket)
.key("proxy-only")
.sse_customer_algorithm("AES256")
.sse_customer_key(&ssec_key_b64)
.sse_customer_key_md5(&ssec_key_md5)
.send()
.await
.map_err(|err| format!("proxied SSE-C GET failed: {}", err.into_service_error()))?;
let ssec_record = target
.requests()
.into_iter()
.find(|record| record.operation == FakeTargetOperation::GetObject && record.key.as_deref() == Some("proxy-only"))
.ok_or("fake target never received the proxied SSE-C GET")?;
assert_eq!(ssec_record.proxy_headers.ssec_algorithm.as_deref(), Some("AES256"));
assert!(ssec_record.proxy_headers.ssec_key_present, "SSE-C key header must be forwarded verbatim");
assert_eq!(ssec_record.proxy_headers.ssec_key_md5.as_deref(), Some(ssec_key_md5.as_str()));
assert!(ssec_record.proxy_headers.replication_check.is_none());
// b. HEAD of the locally-missing object is served through the proxy.
target.take_requests();
let head = source_client
.head_object()
.bucket(source_bucket)
.key("proxy-only")
.send()
.await
.map_err(|err| format!("proxied HEAD failed: {}", err.into_service_error()))?;
assert_eq!(head.content_length, Some(payload.len() as i64));
let head_record = target
.requests()
.into_iter()
.find(|record| record.operation == FakeTargetOperation::HeadObject && record.key.as_deref() == Some("proxy-only"))
.ok_or("fake target never received the proxied HEAD")?;
assert_eq!(head_record.proxy_headers.source_proxy_request.as_deref(), Some("true"));
assert!(head_record.proxy_headers.replication_check.is_none());
// c. Anti-loop: an inbound request that already carries the proxy marker
// is answered locally with 404 and never forwarded to the target.
target.take_requests();
let err = source_client
.get_object()
.bucket(source_bucket)
.key("proxy-only")
.customize()
.mutate_request(|req| {
req.headers_mut().insert("x-minio-source-proxy-request", "true");
})
.send()
.await
.expect_err("anti-loop GET must fail locally instead of proxying");
let service_err = err.into_service_error();
assert!(service_err.is_no_such_key(), "anti-loop GET must 404, got: {service_err}");
assert!(
!target
.requests()
.iter()
.any(|record| record.operation == FakeTargetOperation::GetObject),
"anti-loop GET must not reach the replication target; journal: {:?}",
target.requests()
);
// c2. MinIO ProxyHeaderSet parity: the header's mere PRESENCE disables
// proxying — "false" is exactly what a peer's replication worker sends on
// its convergence HEADs, and proxying that miss back would fake
// convergence.
target.take_requests();
let err = source_client
.get_object()
.bucket(source_bucket)
.key("proxy-only")
.customize()
.mutate_request(|req| {
req.headers_mut().insert("x-minio-source-proxy-request", "false");
})
.send()
.await
.expect_err("proxy-header-set GET must fail locally instead of proxying");
let service_err = err.into_service_error();
assert!(service_err.is_no_such_key(), "proxy-header-set GET must 404, got: {service_err}");
assert!(
!target
.requests()
.iter()
.any(|record| record.operation == FakeTargetOperation::GetObject),
"proxy-header-set GET must not reach the replication target; journal: {:?}",
target.requests()
);
// d. The replication worker's own convergence HEAD against the target
// must carry `source-proxy-request: false` (never proxied back) and the
// replication-check exemption. Trigger real replication and inspect the
// fake journal.
target.take_requests();
source_client
.put_object()
.bucket(source_bucket)
.key("worker-replicated")
.body(ByteStream::from_static(b"worker payload"))
.send()
.await?;
wait_for_target_request_version_id(&target, FakeTargetOperation::PutObject, "worker-replicated").await?;
let worker_head = target
.requests()
.into_iter()
.find(|record| record.operation == FakeTargetOperation::HeadObject && record.key.as_deref() == Some("worker-replicated"))
.ok_or_else(|| format!("replication worker never HEAD-ed the target; journal: {:?}", target.requests()))?;
assert_eq!(
worker_head.proxy_headers.source_proxy_request.as_deref(),
Some("false"),
"worker convergence HEAD must send source-proxy-request: false so the target answers locally"
);
assert_eq!(
worker_head.proxy_headers.replication_check.as_deref(),
Some("true"),
"worker convergence HEAD keeps the replication-check exemption"
);
drop(source_env);
target.shutdown().await;
Ok(())
}
/// P1-5 (backlog#1675): GetObjectTagging for an object missing locally is
/// proxied to the replication target with the anti-loop marker, mirroring
/// MinIO `proxyGetTaggingToRepTarget`.
#[tokio::test]
#[serial]
async fn test_get_object_tagging_proxies_unreplicated_object_to_replication_target() -> TestResult {
init_logging();
let source_bucket = "proxy-tag-src";
let target_bucket = "proxy-tag-dst";
let (target, source_env, source_client, target_client) = start_read_proxy_lab(source_bucket, target_bucket).await?;
target_client
.put_object()
.bucket(target_bucket)
.key("proxy-tagged")
.body(ByteStream::from_static(b"tagged payload"))
.send()
.await?;
target_client
.put_object_tagging()
.bucket(target_bucket)
.key("proxy-tagged")
.tagging(
aws_sdk_s3::types::Tagging::builder()
.tag_set(aws_sdk_s3::types::Tag::builder().key("team").value("storage").build()?)
.build()?,
)
.send()
.await?;
target.take_requests();
let tags = source_client
.get_object_tagging()
.bucket(source_bucket)
.key("proxy-tagged")
.send()
.await
.map_err(|err| format!("proxied GetObjectTagging failed: {}", err.into_service_error()))?;
assert_eq!(tags.tag_set.len(), 1, "proxied tagging read must return the target's tags");
assert_eq!(tags.tag_set[0].key.as_str(), "team");
assert_eq!(tags.tag_set[0].value.as_str(), "storage");
let record = target
.requests()
.into_iter()
.find(|record| record.operation == FakeTargetOperation::GetObjectTagging && record.key.as_deref() == Some("proxy-tagged"))
.ok_or("fake target never received the proxied GetObjectTagging")?;
assert_eq!(
record.proxy_headers.source_proxy_request.as_deref(),
Some("true"),
"proxied tagging read must carry the anti-loop marker"
);
assert!(record.proxy_headers.replication_check.is_none());
drop(source_env);
target.shutdown().await;
Ok(())
}
+7 -6
View File
@@ -198,12 +198,13 @@ pub mod bucket {
ReplicationType, ResyncOpts, ResyncStatusType, RuntimeReplicationTargetBacklog, TargetReplicationResyncStatus,
VersionPurgeStatusType, XferStats, commit_force_delete_intent, complete_force_delete_intent,
delete_replication_state_from_config, delete_replication_version_id, get_global_replication_pool,
get_global_replication_stats, init_background_replication, invalid_replication_config_status_field,
persist_force_delete_intent, read_durable_mrf_backlog, replication_state_to_filemeta, replication_status_to_filemeta,
replication_statuses_map, replication_target_arns, resync_start_conflict_id, should_remove_replication_target,
should_schedule_delete_replication, should_use_existing_delete_replication_info,
should_use_existing_delete_replication_source, unsupported_replication_config_field,
validate_replication_config_structure, validate_replication_config_target_arns, version_purge_status_to_filemeta,
get_global_replication_stats, get_proxy_targets, init_background_replication,
invalid_replication_config_status_field, persist_force_delete_intent, read_durable_mrf_backlog,
replication_state_to_filemeta, replication_status_to_filemeta, replication_statuses_map, replication_target_arns,
resync_start_conflict_id, should_remove_replication_target, should_schedule_delete_replication,
should_use_existing_delete_replication_info, should_use_existing_delete_replication_source,
unsupported_replication_config_field, validate_replication_config_structure, validate_replication_config_target_arns,
version_purge_status_to_filemeta,
};
}
+174 -2
View File
@@ -27,10 +27,15 @@ use aws_sdk_s3::config::SharedHttpClient;
use aws_sdk_s3::error::ProvideErrorMetadata;
use aws_sdk_s3::error::SdkError;
use aws_sdk_s3::operation::complete_multipart_upload::CompleteMultipartUploadOutput;
use aws_sdk_s3::operation::delete_object_tagging::{DeleteObjectTaggingError, DeleteObjectTaggingOutput};
use aws_sdk_s3::operation::get_object::{GetObjectError, GetObjectOutput};
use aws_sdk_s3::operation::get_object_tagging::{GetObjectTaggingError, GetObjectTaggingOutput};
use aws_sdk_s3::operation::head_bucket::HeadBucketError;
use aws_sdk_s3::operation::head_object::HeadObjectError;
use aws_sdk_s3::operation::put_object_tagging::{PutObjectTaggingError, PutObjectTaggingOutput};
use aws_sdk_s3::operation::upload_part::UploadPartOutput;
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::Tagging as SdkTagging;
use aws_sdk_s3::types::{
ChecksumMode, CompletedMultipartUpload, CompletedPart, ObjectLockLegalHoldStatus, ObjectLockRetentionMode,
};
@@ -57,8 +62,8 @@ use rustfs_utils::http::{
is_rustfs_header, is_standard_header, is_storageclass_header,
};
use rustfs_utils::http::{
SUFFIX_FORCE_DELETE, SUFFIX_SOURCE_DELETEMARKER, SUFFIX_SOURCE_ETAG, SUFFIX_SOURCE_MTIME, SUFFIX_SOURCE_REPLICATION_CHECK,
SUFFIX_SOURCE_REPLICATION_LEGALHOLD_TIMESTAMP, SUFFIX_SOURCE_REPLICATION_REQUEST,
SUFFIX_FORCE_DELETE, SUFFIX_SOURCE_DELETEMARKER, SUFFIX_SOURCE_ETAG, SUFFIX_SOURCE_MTIME, SUFFIX_SOURCE_PROXY_REQUEST,
SUFFIX_SOURCE_REPLICATION_CHECK, SUFFIX_SOURCE_REPLICATION_LEGALHOLD_TIMESTAMP, SUFFIX_SOURCE_REPLICATION_REQUEST,
SUFFIX_SOURCE_REPLICATION_RETENTION_TIMESTAMP, SUFFIX_SOURCE_REPLICATION_TAGGING_TIMESTAMP, SUFFIX_SOURCE_VERSION_ID,
insert_header,
};
@@ -1446,6 +1451,43 @@ fn resolve_put_api_version_id(source_version_id: &str) -> Option<&str> {
}
}
/// Resolve the S3 `versionId` for a proxied read against a remote target.
/// RustFS represents the null version internally as the nil UUID while the S3
/// API addresses it as the literal "null" (same mapping as
/// [`resolve_put_api_version_id`]); empty means "no version requested".
fn resolve_read_api_version_id(version_id: Option<String>) -> Option<String> {
let version_id = version_id?;
let trimmed = version_id.trim();
if trimmed.is_empty() {
None
} else if Uuid::parse_str(trimmed).is_ok_and(|uuid| uuid.is_nil()) {
Some(rustfs_filemeta::NULL_VERSION_ID.to_string())
} else {
Some(trimmed.to_string())
}
}
/// Outbound header set for a proxied read: the caller-provided passthrough
/// headers (client SSE-C key family, conditional headers) plus the anti-loop
/// `source-proxy-request` marker in both the x-rustfs- and x-minio- prefixes
/// (a MinIO target only understands the latter). Never adds
/// `source-replication-check`: that exemption channel belongs exclusively to
/// the replication worker's HEAD.
fn proxy_outbound_headers(mut extra_headers: HeaderMap) -> HeaderMap {
insert_header(&mut extra_headers, SUFFIX_SOURCE_PROXY_REQUEST, "true");
extra_headers
}
/// Copy `headers` onto an SDK request inside `customize().map_request` (runs
/// before signing, so the headers join the SigV4 canonical request).
fn apply_extra_headers(mut req: HttpRequest, headers: &HeaderMap) -> Result<HttpRequest, std::convert::Infallible> {
for (k, v) in headers.iter() {
req.headers_mut()
.insert(k.as_str().to_string(), v.to_str().unwrap_or("").to_string());
}
Ok(req)
}
/// Append `versionId=<id>` to an already-built request URI. aws-sdk-s3's
/// `PutObjectInput` / `CreateMultipartUploadInput` expose no version id
/// member, so the query is spliced in via `map_request`, which runs at
@@ -1851,6 +1893,13 @@ impl TargetClient {
// worker cannot hold; otherwise SSE-C replicas never converge on HEAD.
let mut headers = HeaderMap::new();
insert_header(&mut headers, SUFFIX_SOURCE_REPLICATION_CHECK, "true");
// `source-proxy-request: false` (MinIO `ProxyHeaderSet` semantics):
// the header's mere presence tells the receiver to answer LOCALLY
// instead of proxying the miss back to us. Without it, a not-found on
// the target gets read-proxied back to this source, echoes the source
// object with an identical ETag, and the worker concludes the object
// already converged — so it never actually replicates it.
insert_header(&mut headers, SUFFIX_SOURCE_PROXY_REQUEST, "false");
match self
.client
.head_object()
@@ -1875,6 +1924,129 @@ impl TargetClient {
}
}
/// HEAD used by the read-proxy path (GET/HEAD of an object not yet
/// replicated locally, MinIO `proxyHeadToRepTarget`).
///
/// Deliberately different from [`TargetClient::head_object`]: it must NOT
/// send `source-replication-check` — that header is the replication
/// worker's SSE-C metadata exemption channel. A proxied client request
/// instead forwards the client's own SSE-C headers (`extra_headers`) so
/// the target performs the real SSE-C validation/decryption. The
/// `source-proxy-request` marker is always added so the target does not
/// proxy the request onward (anti-loop).
pub async fn head_object_for_proxy(
&self,
bucket: &str,
object: &str,
version_id: Option<String>,
range: Option<String>,
part_number: Option<i32>,
extra_headers: HeaderMap,
) -> Result<HeadObjectOutput, SdkError<HeadObjectError>> {
let headers = proxy_outbound_headers(extra_headers);
self.client
.head_object()
.bucket(bucket)
.key(object)
.set_version_id(resolve_read_api_version_id(version_id))
.set_range(range)
.set_part_number(part_number)
.customize()
.map_request(move |req| apply_extra_headers(req, &headers))
.send()
.await
}
/// GET used by the read-proxy path (MinIO `proxyGetToReplicationTarget`).
/// Returns the streaming SDK output; callers must forward the body without
/// buffering it. Same header contract as [`Self::head_object_for_proxy`]:
/// anti-loop marker on, replication-check never sent, client SSE-C /
/// conditional headers forwarded verbatim via `extra_headers`.
pub async fn get_object(
&self,
bucket: &str,
object: &str,
version_id: Option<String>,
range: Option<String>,
part_number: Option<i32>,
extra_headers: HeaderMap,
) -> Result<GetObjectOutput, SdkError<GetObjectError>> {
let headers = proxy_outbound_headers(extra_headers);
self.client
.get_object()
.bucket(bucket)
.key(object)
.set_version_id(resolve_read_api_version_id(version_id))
.set_range(range)
.set_part_number(part_number)
.customize()
.map_request(move |req| apply_extra_headers(req, &headers))
.send()
.await
}
/// GetObjectTagging for the tagging read-proxy path
/// (MinIO `proxyGetTaggingToRepTarget`). Anti-loop marker always added.
pub async fn get_object_tagging(
&self,
bucket: &str,
object: &str,
version_id: Option<String>,
) -> Result<GetObjectTaggingOutput, SdkError<GetObjectTaggingError>> {
let headers = proxy_outbound_headers(HeaderMap::new());
self.client
.get_object_tagging()
.bucket(bucket)
.key(object)
.set_version_id(resolve_read_api_version_id(version_id))
.customize()
.map_request(move |req| apply_extra_headers(req, &headers))
.send()
.await
}
/// PutObjectTagging for the tagging proxy path
/// (MinIO `proxyTaggingToRepTarget`). Anti-loop marker always added.
pub async fn put_object_tagging(
&self,
bucket: &str,
object: &str,
version_id: Option<String>,
tagging: SdkTagging,
) -> Result<PutObjectTaggingOutput, SdkError<PutObjectTaggingError>> {
let headers = proxy_outbound_headers(HeaderMap::new());
self.client
.put_object_tagging()
.bucket(bucket)
.key(object)
.set_version_id(resolve_read_api_version_id(version_id))
.tagging(tagging)
.customize()
.map_request(move |req| apply_extra_headers(req, &headers))
.send()
.await
}
/// DeleteObjectTagging for the tagging proxy path
/// (MinIO `proxyTaggingToRepTarget`). Anti-loop marker always added.
pub async fn delete_object_tagging(
&self,
bucket: &str,
object: &str,
version_id: Option<String>,
) -> Result<DeleteObjectTaggingOutput, SdkError<DeleteObjectTaggingError>> {
let headers = proxy_outbound_headers(HeaderMap::new());
self.client
.delete_object_tagging()
.bucket(bucket)
.key(object)
.set_version_id(resolve_read_api_version_id(version_id))
.customize()
.map_request(move |req| apply_extra_headers(req, &headers))
.send()
.await
}
/// On success returns the version id the target assigned (from
/// `x-amz-version-id`), letting callers audit the version-identity
/// contract — a target that adopts the source version echoes it back.
@@ -14,6 +14,7 @@ paths.
| `datatypes.rs` | ECStore compatibility re-export for resync status enums. | Re-exports `rustfs-replication` contracts while downstream facade consumers migrate. |
| `replication_object_decision_boundary.rs` | Object replication option DTOs, resync target projection, delete replication decisions, and multipart planning helpers. | Keeps ECStore runtime modules from importing object decision contracts directly from `rustfs-replication`. |
| `replication_pool.rs` | Replication queue, worker pool, MRF persistence, bucket stats, and delete/object scheduling. | Depends on bucket target sys, bucket metadata sys, metadata paths, queue contracts through the queue boundary, file metadata replication contracts through local boundaries, config storage, storage contracts through the replication storage boundary, runtime sources, and notification state. |
| `replication_proxy.rs` | Proxy-target selection for GET/HEAD/Tagging reads of objects not yet replicated locally (MinIO `getProxyTargets` parity: anti-loop, version-suspended, and no-config empty branches). | Uses replication config lookup, rule matching, and target clients through local boundaries. |
| `replication_queue_boundary.rs` | Queue/admission DTOs, heal queue DTOs, worker sizing, and backpressure helpers. | Keeps ECStore runtime modules from importing queue/backpressure contracts directly from `rustfs-replication`. |
| `replication_resync_boundary.rs` | Resync DTOs, status classifiers, persisted resync/MRF codec wrappers, and ECStore error mapping. | Keeps ECStore runtime modules from importing resync contract helpers directly from `rustfs-replication`. |
| `replication_resyncer.rs` | Object replication, delete replication, resync execution, target calls, and multipart target upload paths. | Depends on target calls and target config types through the replication target boundary, metadata paths and metadata systems through the replication metadata boundary, file metadata replication contracts through the filemeta boundary, object decisions and multipart planning through the object decision boundary, resync contracts through the resync boundary, queue DTOs through the queue boundary, error contracts through the error boundary, versioning systems, storage contracts through the replication storage boundary, config-derived storage class labels through the config store, runtime sources, notification events and local event host selection through the event sink, bandwidth reader wrapping, and SetDisks lock timing. |
@@ -29,6 +29,7 @@ mod replication_object_bridge;
mod replication_object_config;
mod replication_object_decision_boundary;
pub(crate) mod replication_pool;
mod replication_proxy;
mod replication_queue_boundary;
mod replication_resync_boundary;
mod replication_resyncer;
@@ -74,6 +75,7 @@ pub use replication_pool::{
get_global_replication_pool, get_global_replication_stats, init_background_replication, persist_force_delete_intent,
read_durable_mrf_backlog, resync_start_conflict_id,
};
pub use replication_proxy::get_proxy_targets;
pub use replication_queue_boundary::{
DeletedObjectReplicationInfo, ReplicationBatchAdmission, ReplicationHealQueueResult, ReplicationOperation,
ReplicationPriority, ReplicationQueueAdmission,
@@ -0,0 +1,150 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! Proxy-target selection for reads of objects not yet replicated locally
//! (MinIO `getProxyTargets`, bucket-replication.go).
//!
//! During the active-active replication lag window a GET/HEAD/Tagging request
//! for an object the local site does not have yet may be served by proxying to
//! a replication target. This module only *selects* the candidate targets; the
//! request-path callers perform the remote calls and response translation.
use std::sync::Arc;
use tracing::debug;
use super::replication_config_boundary::{ObjectOpts, ReplicationConfigurationExt as _};
use super::replication_object_config::get_replication_config;
use super::replication_storage_boundary::ObjectOptions;
use super::replication_target_boundary::{ReplicationTargetStore, TargetClient};
/// Returns the replication-target clients eligible to serve a proxied read of
/// `bucket/object`, in rule order. Mirrors MinIO's `getProxyTargets`:
///
/// - the `source-proxy-request` header family was present at all
/// (`opts.proxy_request` / `opts.proxy_header_set`, MinIO `ProxyRequest` /
/// `ProxyHeaderSet`) -> empty. "true" is the anti-loop marker of an
/// already-proxied client read; "false" is what a peer's replication
/// worker sends on convergence HEADs so the receiver answers locally —
/// proxying that miss back would echo the source object and fake
/// convergence, permanently skipping replication;
/// - the bucket's versioning is suspended for the object -> empty;
/// - no replication configuration / no matching rule -> empty;
/// - otherwise every distinct target ARN whose rules match the object,
/// resolved through the bucket target system, skipping targets that opted
/// out of proxying (`disable_proxy`).
pub async fn get_proxy_targets(bucket: &str, object: &str, opts: &ObjectOptions) -> Vec<Arc<TargetClient>> {
if opts.proxy_request || opts.proxy_header_set {
return Vec::new();
}
if opts.version_suspended {
return Vec::new();
}
let cfg = match get_replication_config(bucket).await {
Ok(Some(cfg)) => cfg,
Ok(None) => return Vec::new(),
Err(err) => {
debug!(bucket, object, error = %err, "read proxy: failed to load replication config; not proxying");
return Vec::new();
}
};
let arns = cfg.filter_target_arns(&ObjectOpts {
name: object.to_string(),
..Default::default()
});
let mut targets = Vec::with_capacity(arns.len());
for arn in arns {
let Some(client) = ReplicationTargetStore::remote_target_client(bucket, &arn).await else {
debug!(bucket, object, arn, "read proxy: no client for replication target ARN");
continue;
};
if client.disable_proxy {
continue;
}
targets.push(client);
}
targets
}
#[cfg(test)]
mod tests {
use super::*;
fn opts() -> ObjectOptions {
ObjectOptions::default()
}
/// Anti-loop: a request that was already proxied by a peer must never be
/// proxied onward, regardless of replication configuration.
#[tokio::test]
async fn proxy_request_yields_no_targets() {
let targets = get_proxy_targets(
"bucket",
"object",
&ObjectOptions {
proxy_request: true,
..opts()
},
)
.await;
assert!(targets.is_empty());
}
/// MinIO `ProxyHeaderSet` parity: the header family being present at all
/// disables proxying, even with the value "false" — that is what a
/// peer's replication worker sends on convergence HEADs.
#[tokio::test]
async fn proxy_header_set_yields_no_targets() {
let targets = get_proxy_targets(
"bucket",
"object",
&ObjectOptions {
proxy_header_set: true,
proxy_request: false,
..opts()
},
)
.await;
assert!(targets.is_empty());
}
/// Suspended versioning disables proxying (MinIO parity): the local null
/// version is authoritative and a remote read could resurrect data.
#[tokio::test]
async fn version_suspended_yields_no_targets() {
let targets = get_proxy_targets(
"bucket",
"object",
&ObjectOptions {
version_suspended: true,
..opts()
},
)
.await;
assert!(targets.is_empty());
}
/// A bucket without replication configuration has nothing to proxy to.
/// (No metadata system is running in unit tests, so the config lookup
/// resolves to "no configuration" — the same empty-result contract.)
#[tokio::test]
async fn missing_replication_config_yields_no_targets() {
let targets = get_proxy_targets("bucket-without-replication", "object", &opts()).await;
assert!(targets.is_empty());
}
}
@@ -15,10 +15,14 @@
use super::replication_error_boundary::{Error, Result};
use super::replication_filemeta_boundary::MrfReplicateEntry;
/// Kept test-only: the runtime consumer was the worker HEAD's fake proxy
/// counting (removed in backlog#1675 P1-5); the resyncer tests still pin the
/// classifier's semantics for the real client read-proxy failure accounting.
#[cfg(test)]
pub(crate) use rustfs_replication::should_count_head_proxy_failure;
pub use rustfs_replication::{BucketReplicationResyncStatus, ResyncOpts, ResyncStatusType, TargetReplicationResyncStatus};
pub(crate) use rustfs_replication::{
is_version_id_mismatch, resync_state_accepts_update, sanitize_resync_error_detail, should_auto_resume_resync,
should_count_head_proxy_failure,
};
#[allow(
@@ -36,6 +36,7 @@ use super::replication_object_decision_boundary::{
};
use super::replication_queue_boundary::{DeletedObjectReplicationInfo, ReplicationQueueAdmission};
use super::replication_resync_boundary::ResyncStatusType;
#[cfg(test)]
use super::replication_resync_boundary::should_count_head_proxy_failure;
use super::replication_resync_boundary::{
BucketReplicationResyncStatus, ResyncOpts, TargetReplicationResyncStatus, encode_resync_file, is_version_id_mismatch,
@@ -234,32 +235,18 @@ fn audit_target_version_identity(tgt_client: &TargetClient, source_version_id: &
}
}
fn is_head_proxy_failure(err: &SdkError<HeadObjectError>) -> bool {
let (is_not_found, code) = err
.as_service_error()
.map(|service_err| (service_err.is_not_found(), service_err.code()))
.unwrap_or((false, None));
let raw_status = err.raw_response().map(|resp| resp.status().as_u16());
should_count_head_proxy_failure(is_not_found, code, raw_status)
}
async fn record_proxy_request(bucket: &str, api: &str, is_err: bool) {
if let Some(stats) = runtime_sources::replication_stats() {
stats.inc_proxy(bucket, api, is_err).await;
}
}
async fn head_object_with_proxy_stats(
source_bucket: &str,
/// HEAD against a replication target on behalf of the replication worker
/// (resync/heal/delete convergence checks). This is NOT a client read proxy:
/// it must not touch the proxy metrics — those count only real GET/HEAD/
/// Tagging requests proxied for clients (see `replication_proxy.rs` /
/// `TargetClient::head_object_for_proxy`).
async fn head_object_for_worker(
target_client: &TargetClient,
target_bucket: &str,
object: &str,
version_id: Option<String>,
) -> std::result::Result<HeadObjectOutput, SdkError<HeadObjectError>> {
let result = target_client.head_object(target_bucket, object, version_id).await;
let is_err = result.as_ref().err().is_some_and(is_head_proxy_failure);
record_proxy_request(source_bucket, "HeadObject", is_err).await;
result
target_client.head_object(target_bucket, object, version_id).await
}
fn is_version_id_format_mismatch(err: &SdkError<HeadObjectError>) -> bool {
@@ -282,11 +269,10 @@ async fn mark_replication_target_offline_if_needed(target_client: &Arc<TargetCli
}
async fn head_object_fallback(
source_bucket: &str,
tgt_client: &TargetClient,
object: &str,
) -> std::result::Result<Option<HeadObjectOutput>, SdkError<HeadObjectError>> {
match head_object_with_proxy_stats(source_bucket, tgt_client, &tgt_client.bucket, object, None).await {
match head_object_for_worker(tgt_client, &tgt_client.bucket, object, None).await {
Ok(oi) => Ok(Some(oi)),
Err(e) if e.as_service_error().is_some_and(|se| se.is_not_found()) || has_raw_status(&e, 404) => Ok(None),
Err(e) => Err(e),
@@ -1186,7 +1172,7 @@ async fn verify_resync_head_result(
// (400). Re-verify without the versionId before
// concluding the object failed to replicate, instead
// of counting a well-replicated object as failed.
match head_object_fallback(&roi.bucket, target_client.as_ref(), &roi.name).await {
match head_object_fallback(target_client.as_ref(), &roi.name).await {
Ok(Some(_)) => {
st.replicated_count += 1;
st.replicated_size += roi.size;
@@ -1236,8 +1222,7 @@ async fn resync_worker_process_object<S: ReplicationStorage>(
let reset_id = target_client.reset_id.clone();
let head_result = head_object_with_proxy_stats(
bucket_name,
let head_result = head_object_for_worker(
target_client.as_ref(),
&target_client.bucket,
&roi.name,
@@ -2521,8 +2506,7 @@ async fn replicate_delete_to_target(dobj: &DeletedObjectReplicationInfo, tgt_cli
let version_id = target_delete_version_id(version_id, is_version_purge);
if dobj.delete_object.delete_marker && dobj.delete_object.delete_marker_version_id.is_some() {
match head_object_with_proxy_stats(
&dobj.bucket,
match head_object_for_worker(
tgt_client.as_ref(),
&tgt_client.bucket,
&dobj.delete_object.object_name,
@@ -2985,14 +2969,8 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
}
let mut replication_action = replication_action;
match head_object_with_proxy_stats(
&bucket,
tgt_client.as_ref(),
&tgt_client.bucket,
&object,
self.version_id.map(|v| v.to_string()),
)
.await
match head_object_for_worker(tgt_client.as_ref(), &tgt_client.bucket, &object, self.version_id.map(|v| v.to_string()))
.await
{
Ok(oi) => {
replication_action = replication_action_for_target_head(&object_info, &oi, self.op_type);
@@ -3009,7 +2987,7 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
// Object not on target yet → fall through to PUT.
} else if is_version_id_format_mismatch(&e) {
// Version-ID format mismatch: retry without versionId and compare ETags.
match head_object_fallback(&bucket, &tgt_client, &object).await {
match head_object_fallback(&tgt_client, &object).await {
Ok(Some(oi)) if replication_etags_match(object_info.etag.as_deref(), oi.e_tag.as_deref()) => {
rinfo.replication_status = ReplicationStatusType::Completed;
rinfo.replication_resynced = true;
@@ -3085,7 +3063,6 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
}
};
let has_tagging_replication = !put_opts.user_tags.is_empty();
if let Some(err) = if is_multipart {
drop(gr);
let result = replicate_object_with_multipart(MultipartReplicationContext {
@@ -3100,10 +3077,6 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
put_opts,
})
.await;
record_proxy_request(&bucket, "PutObject", result.is_err()).await;
if has_tagging_replication {
record_proxy_request(&bucket, "PutObjectTagging", result.is_err()).await;
}
result.err()
} else {
gr.stream = wrap_with_bandwidth_monitor(gr.stream, &put_opts, &bucket, &rinfo.arn);
@@ -3119,10 +3092,6 @@ impl ReplicateObjectInfoExt for ReplicateObjectInfo {
)
})
.map_err(|e| std::io::Error::other(e.to_string()));
record_proxy_request(&bucket, "PutObject", result.is_err()).await;
if has_tagging_replication {
record_proxy_request(&bucket, "PutObjectTagging", result.is_err()).await;
}
result.err()
} {
rinfo.replication_status = ReplicationStatusType::Failed;
@@ -3486,15 +3455,7 @@ async fn resolve_replicate_all_action(
rinfo: &mut ReplicatedTargetInfo,
) -> Option<(ReplicationAction, ObjectInfo)> {
let replication_action;
match head_object_with_proxy_stats(
bucket,
tgt_client.as_ref(),
&tgt_client.bucket,
object,
roi.version_id.map(|v| v.to_string()),
)
.await
{
match head_object_for_worker(tgt_client.as_ref(), &tgt_client.bucket, object, roi.version_id.map(|v| v.to_string())).await {
Ok(oi) => {
replication_action = replication_action_for_target_head(&object_info, &oi, roi.op_type);
rinfo.replication_status = ReplicationStatusType::Completed;
@@ -3545,7 +3506,7 @@ async fn resolve_replicate_all_action(
Err(e) => {
if is_version_id_format_mismatch(&e) {
// Version-ID format mismatch: retry without versionId and compare ETags.
match head_object_fallback(bucket, tgt_client, object).await {
match head_object_fallback(tgt_client, object).await {
Ok(Some(oi)) => {
replication_action = if replication_etags_match(object_info.etag.as_deref(), oi.e_tag.as_deref()) {
ReplicationAction::None
@@ -3668,7 +3629,6 @@ async fn replicate_all_payload_to_target<S: ReplicationObjectIO>(
ctx: ReplicateAllPayloadContext<'_, S>,
mut gr: GetObjectReader,
) -> Option<std::io::Error> {
let has_tagging_replication = !ctx.put_opts.user_tags.is_empty();
if ctx.is_multipart {
drop(gr);
let result = replicate_object_with_multipart(MultipartReplicationContext {
@@ -3683,10 +3643,6 @@ async fn replicate_all_payload_to_target<S: ReplicationObjectIO>(
put_opts: ctx.put_opts,
})
.await;
record_proxy_request(ctx.bucket, "PutObject", result.is_err()).await;
if has_tagging_replication {
record_proxy_request(ctx.bucket, "PutObjectTagging", result.is_err()).await;
}
result.err()
} else {
gr.stream = wrap_with_bandwidth_monitor(gr.stream, &ctx.put_opts, ctx.bucket, ctx.arn);
@@ -3703,10 +3659,6 @@ async fn replicate_all_payload_to_target<S: ReplicationObjectIO>(
)
})
.map_err(|e| std::io::Error::other(e.to_string()));
record_proxy_request(ctx.bucket, "PutObject", result.is_err()).await;
if has_tagging_replication {
record_proxy_request(ctx.bucket, "PutObjectTagging", result.is_err()).await;
}
result.err()
}
}
@@ -1161,6 +1161,31 @@ mod tests {
assert!(all.contains_key("proxy-only-bucket"));
}
/// Pins the read-proxy metric contract (backlog#1675 P1-5): the API
/// strings the GET/HEAD/Tagging proxy paths record map onto the
/// get/head/tagging totals, and only unexpected failures raise the
/// failed counters.
#[tokio::test]
async fn test_proxy_stats_map_read_proxy_apis_to_totals() {
let stats = ReplicationStats::new();
stats.inc_proxy("proxy-bucket", "GetObject", false).await;
stats.inc_proxy("proxy-bucket", "GetObject", true).await;
stats.inc_proxy("proxy-bucket", "HeadObject", false).await;
stats.inc_proxy("proxy-bucket", "GetObjectTagging", false).await;
stats.inc_proxy("proxy-bucket", "PutObjectTagging", false).await;
stats.inc_proxy("proxy-bucket", "DeleteObjectTagging", true).await;
let metric = stats.get_proxy_stats("proxy-bucket").await;
assert_eq!(metric.get_total, 2);
assert_eq!(metric.get_failed, 1);
assert_eq!(metric.head_total, 1);
assert_eq!(metric.head_failed, 0);
assert_eq!(metric.get_tag_total, 1);
assert_eq!(metric.put_tag_total, 1);
assert_eq!(metric.delete_tag_total, 1);
assert_eq!(metric.delete_tag_failed, 1);
}
#[tokio::test]
async fn test_calculate_bucket_replication_stats_merges_resync_metrics() {
let stats = ReplicationStats::new();
+14
View File
@@ -277,6 +277,20 @@ pub struct ObjectOptions {
/// fence avoids recursively acquiring the read lock behind a queued writer.
pub bucket_lifecycle_lock_fence: Option<NamespaceLockFence>,
pub replication_request: bool,
/// True when the inbound request carried the
/// `{x-rustfs-,x-minio-}source-proxy-request` header family with the
/// value "true": the request was already proxied by a replication peer,
/// so this server must not proxy a local miss onward (anti-loop,
/// MinIO-compatible). The header only disables proxying — it grants no
/// capability — so no authorization gate is required to honor it.
pub proxy_request: bool,
/// True when the `source-proxy-request` header family was present at
/// all, regardless of value (MinIO's `ProxyHeaderSet`). A replication
/// peer sends `source-proxy-request: false` on its worker convergence
/// HEADs precisely so the receiver answers locally instead of proxying
/// back — otherwise a proxied 404->200 echo makes the worker believe the
/// object already converged and it never replicates it.
pub proxy_header_set: bool,
/// Source-cluster LWW timestamps carried by an authorized replication
/// request; None when the source never modified the category. Only the
/// replication-authorized options builders may set these.
@@ -4586,7 +4586,10 @@ async fn build_metrics_summary(local_peer: &PeerInfo) -> SRMetricsSummary {
head_failed_total: non_negative_u64(node.proxy_head_failed),
put_tag_total: non_negative_u64(node.proxy_put_tag_total),
put_tag_failed_total: non_negative_u64(node.proxy_put_tag_failed),
..Default::default()
get_tag_total: non_negative_u64(node.proxy_get_tag_total),
get_tag_failed_total: non_negative_u64(node.proxy_get_tag_failed),
remove_tag_total: non_negative_u64(node.proxy_delete_tag_total),
remove_tag_failed_total: non_negative_u64(node.proxy_delete_tag_failed),
},
metrics,
uptime: node.uptime,
+248 -3
View File
@@ -46,9 +46,10 @@ use super::storage_api::object_usecase::bucket::{
replication::{
DeleteReplicationConfigSnapshot, REPLICATE_INCOMING_DELETE, ReplicationStatusType, commit_force_delete_intent,
delete_replication_state_from_config, delete_replication_version_id, deleted_object_has_pending_replication_delete,
force_delete_target_set, has_active_delete_rule, load_delete_config_snapshot, must_replicate_object,
persist_force_delete_intent, schedule_object_replication, schedule_replication_delete, schedule_replication_deletes,
set_deleted_object_replication_state, should_schedule_delete_replication, should_use_existing_delete_replication_info,
force_delete_target_set, get_read_proxy_targets, has_active_delete_rule, load_delete_config_snapshot,
must_replicate_object, persist_force_delete_intent, record_replication_proxy, schedule_object_replication,
schedule_replication_delete, schedule_replication_deletes, set_deleted_object_replication_state,
should_schedule_delete_replication, should_use_existing_delete_replication_info,
},
tagging::decode_tags,
validate_restore_request,
@@ -6598,6 +6599,226 @@ impl DefaultObjectUsecase {
})
}
/// Headers a proxied read forwards verbatim to the replication target:
/// only the client's SSE-C key family, so the target performs the real
/// SSE-C decryption (never the replication-check exemption). HTTP
/// conditional headers (If-Match & co.) are deliberately NOT forwarded —
/// MinIO does not forward them either, and a remote 304/412 would leak a
/// conditional evaluation against a replica the local site never saw.
/// Range and part-number travel as typed SDK parameters instead.
fn proxy_read_passthrough_headers(headers: &HeaderMap) -> HeaderMap {
const FORWARDED: &[&str] = &[
"x-amz-server-side-encryption-customer-algorithm",
"x-amz-server-side-encryption-customer-key",
"x-amz-server-side-encryption-customer-key-md5",
];
let mut forwarded = HeaderMap::new();
for name in FORWARDED {
if let Ok(header_name) = http::HeaderName::from_str(name)
&& let Some(value) = headers.get(&header_name)
{
forwarded.insert(header_name, value.clone());
}
}
forwarded
}
/// True when a proxied SDK call failed because the target does not have
/// the object either (service-level not-found or a raw 404, which also
/// covers NoSuchVersion): the caller tries the next target silently.
fn proxy_sdk_error_is_not_found<E>(err: &aws_sdk_s3::error::SdkError<E>) -> bool {
err.raw_response().is_some_and(|resp| resp.status().as_u16() == 404)
}
/// Serve a GET whose local read failed with not-found by proxying to the
/// bucket's replication targets (MinIO `proxyGetToReplicationTarget`,
/// backlog#1675 P1-5). Returns None when no target can serve the object;
/// the caller then returns the original local error.
async fn proxy_get_object_to_replication_targets(
req: &S3Request<GetObjectInput>,
bucket: &str,
key: &str,
opts: &ObjectOptions,
) -> Option<GetObjectOutput> {
let targets = get_read_proxy_targets(bucket, key, opts).await;
if targets.is_empty() {
return None;
}
let extra_headers = Self::proxy_read_passthrough_headers(&req.headers);
let range = req
.headers
.get(http::header::RANGE)
.and_then(|value| value.to_str().ok())
.map(str::to_owned);
let part_number = req.input.part_number;
for target in targets {
match target
.get_object(
&target.bucket,
key,
opts.version_id.clone(),
range.clone(),
part_number,
extra_headers.clone(),
)
.await
{
Ok(remote) => {
// MinIO-aligned accounting: one total per proxy attempt
// (targets were available), one failed when no target
// served it — never per target.
record_replication_proxy(bucket, "GetObject", false).await;
return Some(Self::proxy_sdk_get_output_to_s3s(remote));
}
Err(err) if Self::proxy_sdk_error_is_not_found(&err) => {
debug!(bucket, key, arn = %target.arn, "read proxy: target does not have the object");
}
Err(err) => {
warn!(bucket, key, arn = %target.arn, error = %err, "read proxy: GET against replication target failed");
}
}
}
record_replication_proxy(bucket, "GetObject", true).await;
None
}
/// Serve a HEAD whose local lookup failed with not-found by proxying to
/// the bucket's replication targets (MinIO `proxyHeadToRepTarget`).
async fn proxy_head_object_to_replication_targets(
req: &S3Request<HeadObjectInput>,
bucket: &str,
key: &str,
opts: &ObjectOptions,
) -> Option<HeadObjectOutput> {
let targets = get_read_proxy_targets(bucket, key, opts).await;
if targets.is_empty() {
return None;
}
let extra_headers = Self::proxy_read_passthrough_headers(&req.headers);
let range = req
.headers
.get(http::header::RANGE)
.and_then(|value| value.to_str().ok())
.map(str::to_owned);
let part_number = req.input.part_number;
for target in targets {
match target
.head_object_for_proxy(
&target.bucket,
key,
opts.version_id.clone(),
range.clone(),
part_number,
extra_headers.clone(),
)
.await
{
Ok(remote) => {
// MinIO-aligned accounting: one total per proxy attempt,
// one failed when no target served it.
record_replication_proxy(bucket, "HeadObject", false).await;
return Some(Self::proxy_sdk_head_output_to_s3s(remote));
}
Err(err) if Self::proxy_sdk_error_is_not_found(&err) => {
debug!(bucket, key, arn = %target.arn, "read proxy: target does not have the object");
}
Err(err) => {
warn!(bucket, key, arn = %target.arn, error = %err, "read proxy: HEAD against replication target failed");
}
}
}
record_replication_proxy(bucket, "HeadObject", true).await;
None
}
/// Translate a proxied SDK GET response into the s3s output, forwarding
/// the body as a stream (no buffering, no local persistence).
fn proxy_sdk_get_output_to_s3s(remote: aws_sdk_s3::operation::get_object::GetObjectOutput) -> GetObjectOutput {
let body = remote.body;
let body_stream = tokio_util::io::ReaderStream::with_capacity(body.into_async_read(), 64 * 1024);
GetObjectOutput {
body: Some(StreamingBlob::wrap(body_stream)),
content_length: remote.content_length,
content_range: remote.content_range,
content_type: remote.content_type.as_deref().and_then(|v| ContentType::from_str(v).ok()),
content_encoding: remote.content_encoding,
content_disposition: remote.content_disposition,
content_language: remote.content_language,
cache_control: remote.cache_control,
accept_ranges: Some(ACCEPT_RANGES_BYTES.to_string()),
e_tag: remote.e_tag.as_deref().and_then(|v| ETag::from_str(v).ok()),
last_modified: remote
.last_modified
.and_then(|dt| OffsetDateTime::from_unix_timestamp_nanos(dt.as_nanos()).ok())
.map(Timestamp::from),
metadata: remote.metadata,
version_id: remote.version_id,
server_side_encryption: remote
.server_side_encryption
.map(|sse| ServerSideEncryption::from(sse.as_str().to_string())),
sse_customer_algorithm: remote.sse_customer_algorithm,
sse_customer_key_md5: remote.sse_customer_key_md5,
ssekms_key_id: remote.ssekms_key_id,
parts_count: remote.parts_count,
tag_count: remote.tag_count,
storage_class: remote.storage_class.map(|sc| StorageClass::from(sc.as_str().to_string())),
expiration: remote.expiration,
restore: remote.restore,
checksum_crc32: remote.checksum_crc32,
checksum_crc32c: remote.checksum_crc32_c,
checksum_crc64nvme: remote.checksum_crc64_nvme,
checksum_sha1: remote.checksum_sha1,
checksum_sha256: remote.checksum_sha256,
checksum_type: remote.checksum_type.map(|ct| ChecksumType::from(ct.as_str().to_string())),
..Default::default()
}
}
/// Translate a proxied SDK HEAD response into the s3s output.
///
/// Known gaps: the SDK's HeadObjectOutput does not model 206/Content-Range
/// for a ranged HEAD (the SDK exposes no content_range member on HEAD),
/// and s3s' typed HeadObjectOutput has no tag_count field (the local path
/// injects x-amz-tagging-count as a raw header) — both are dropped for
/// proxied HEADs.
fn proxy_sdk_head_output_to_s3s(remote: aws_sdk_s3::operation::head_object::HeadObjectOutput) -> HeadObjectOutput {
HeadObjectOutput {
content_length: remote.content_length,
content_type: remote.content_type.as_deref().and_then(|v| ContentType::from_str(v).ok()),
content_encoding: remote.content_encoding,
content_disposition: remote.content_disposition,
content_language: remote.content_language,
cache_control: remote.cache_control,
accept_ranges: Some(ACCEPT_RANGES_BYTES.to_string()),
e_tag: remote.e_tag.as_deref().and_then(|v| ETag::from_str(v).ok()),
last_modified: remote
.last_modified
.and_then(|dt| OffsetDateTime::from_unix_timestamp_nanos(dt.as_nanos()).ok())
.map(Timestamp::from),
metadata: remote.metadata,
version_id: remote.version_id,
server_side_encryption: remote
.server_side_encryption
.map(|sse| ServerSideEncryption::from(sse.as_str().to_string())),
sse_customer_algorithm: remote.sse_customer_algorithm,
sse_customer_key_md5: remote.sse_customer_key_md5,
ssekms_key_id: remote.ssekms_key_id,
parts_count: remote.parts_count,
storage_class: remote.storage_class.map(|sc| StorageClass::from(sc.as_str().to_string())),
expiration: remote.expiration,
restore: remote.restore,
checksum_crc32: remote.checksum_crc32,
checksum_crc32c: remote.checksum_crc32_c,
checksum_crc64nvme: remote.checksum_crc64_nvme,
checksum_sha1: remote.checksum_sha1,
checksum_sha256: remote.checksum_sha256,
checksum_type: remote.checksum_type.map(|ct| ChecksumType::from(ct.as_str().to_string())),
..Default::default()
}
}
#[instrument(name = "execute_get_object", level = "trace", skip(self, req))]
pub async fn execute_get_object(&self, req: S3Request<GetObjectInput>) -> S3Result<S3Response<GetObjectOutput>> {
self.execute_get_object_boxed(req).await
@@ -6723,6 +6944,19 @@ impl DefaultObjectUsecase {
{
Ok(prepared_read) => prepared_read,
Err(err) => {
// Active-active replication lag window: an object missing
// locally (and only missing — other errors keep their
// semantics) may still be served by proxying the GET to a
// replication target (backlog#1675 P1-5).
if matches!(*err.code(), S3ErrorCode::NoSuchKey | S3ErrorCode::NoSuchVersion)
&& let Some(output) = Self::proxy_get_object_to_replication_targets(&req, &bucket, &key, &opts).await
{
lifecycle.finish_ok();
let response = wrap_response_with_cors(&bucket, &req.method, &req.headers, output).await;
let result = Ok(response);
let _ = helper.version_id(version_id_for_event).complete(&result);
return result;
}
lifecycle.finish_err();
return Err(err);
}
@@ -8632,6 +8866,17 @@ impl DefaultObjectUsecase {
let msg = head_prefix_not_found_message(&bucket, &key, has_children);
return Err(S3Error::with_message(S3ErrorCode::NoSuchKey, msg));
}
// Active-active replication lag window: an object missing
// locally may still be served by proxying the HEAD to a
// replication target (backlog#1675 P1-5).
if let Some(output) = Self::proxy_head_object_to_replication_targets(&req, &bucket, &key, &opts).await {
let response = wrap_response_with_cors(&bucket, &req.method, &req.headers, output).await;
let result = Ok(response);
let _ = helper
.version_id(req.input.version_id.clone().unwrap_or_default())
.complete(&result);
return result;
}
return Err(S3Error::new(S3ErrorCode::NoSuchKey));
}
// Other errors, such as insufficient permissions, still return the original error
+18
View File
@@ -627,6 +627,24 @@ pub(crate) mod bucket {
#[cfg(test)]
pub(crate) use replication_contracts::replication_statuses_map;
/// Remote replication-target client used by the read-proxy path.
pub(crate) type ProxyTargetClient = crate::storage::storage_api::ecstore_bucket::bucket_target_sys::TargetClient;
/// Proxy-request metric recorder (get/head/tagging totals + failures).
pub(crate) use crate::storage::storage_api::record_replication_proxy;
/// Replication targets eligible to serve a proxied GET/HEAD/Tagging of
/// an object not present locally (MinIO `getProxyTargets`; empty when
/// the request was itself proxied, versioning is suspended, or no
/// replication rule matches). backlog#1675 P1-5.
pub(crate) async fn get_read_proxy_targets(
bucket: &str,
object: &str,
opts: &crate::storage::storage_api::StorageObjectOptions,
) -> Vec<Arc<ProxyTargetClient>> {
replication_contracts::get_proxy_targets(bucket, object, opts).await
}
pub(crate) async fn persist_force_delete_intent(
store: Arc<crate::storage::storage_api::ECStore>,
bucket: String,
+229 -38
View File
@@ -12,15 +12,15 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use super::StorageVersioningConfigExt as _;
use super::{
BUCKET_ACCELERATE_CONFIG, BUCKET_LOGGING_CONFIG, BUCKET_REQUEST_PAYMENT_CONFIG, BUCKET_VERSIONING_CONFIG,
BUCKET_WEBSITE_CONFIG, BucketVersioningSys, OBJECT_LOCK_CONFIG, StorageError, check_retention_for_modification, decode_tags,
decode_tags_to_map, delete_bucket_metadata_config_if_incarnation, encode_tags, get_bucket_accelerate_config,
get_bucket_logging_config, get_bucket_object_lock_config, get_bucket_replication_config, get_bucket_request_payment_config,
get_bucket_website_config, is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found,
record_replication_proxy, serialize, update_bucket_metadata_config_if_incarnation,
get_bucket_logging_config, get_bucket_object_lock_config, get_bucket_request_payment_config, get_bucket_website_config,
is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found, record_replication_proxy, serialize,
update_bucket_metadata_config_if_incarnation,
};
use super::{StorageReplicationConfigExt as _, StorageVersioningConfigExt as _};
use crate::admin::handlers::site_replication::site_replication_bucket_meta_hook;
use crate::error::ApiError;
use crate::storage::access::{apply_bucket_generation_guard, bucket_config_mutation_incarnation, has_bypass_governance_header};
@@ -59,7 +59,7 @@ const LOG_SUBSYSTEM_OBJECT_LOCK: &str = "object_lock";
const LOG_SUBSYSTEM_TAGGING: &str = "tagging";
use crate::app::storage_api::object_usecase::bucket::replication::{
ReplicateDecision, must_replicate_metadata, schedule_metadata_replication,
ReplicateDecision, get_read_proxy_targets, must_replicate_metadata, schedule_metadata_replication,
};
use crate::storage::storage_api::ecfs_consumer::StorageObjectOptions as ObjectOptions;
@@ -105,18 +105,152 @@ impl FS {
&self.server_ctx
}
async fn replication_tagging_enabled(bucket: &str, object: &str) -> bool {
get_bucket_replication_config(bucket)
.await
.map(|(cfg, _)| cfg.has_active_rules(object, true))
.unwrap_or(false)
/// Not-found classifier for proxied SDK tagging calls: a raw 404 covers
/// NoSuchKey and NoSuchVersion alike; the caller silently tries the next
/// replication target.
fn proxy_sdk_error_is_not_found<E>(err: &aws_sdk_s3::error::SdkError<E>) -> bool {
err.raw_response().is_some_and(|resp| resp.status().as_u16() == 404)
}
async fn record_replication_tagging_metric(bucket: &str, object: &str, api: &str, is_err: bool) {
if !Self::replication_tagging_enabled(bucket, object).await {
return;
/// Selector options for a tagging proxy. Reuses `get_opts` so the
/// anti-loop `source-proxy-request` header family and the bucket's
/// version-suspension state gate proxying exactly like GET/HEAD.
async fn tagging_proxy_opts(
bucket: &str,
object: &str,
version_id: Option<String>,
headers: &http::HeaderMap,
) -> Option<ObjectOptions> {
get_opts(bucket, object, version_id, None, headers).await.ok()
}
/// Serve a GetObjectTagging for an object missing locally by proxying to
/// the bucket's replication targets (MinIO `proxyGetTaggingToRepTarget`,
/// backlog#1675 P1-5). None means no target had the object.
async fn proxy_get_object_tagging(
bucket: &str,
object: &str,
version_id: Option<String>,
headers: &http::HeaderMap,
) -> Option<TagSet> {
let opts = Self::tagging_proxy_opts(bucket, object, version_id, headers).await?;
let targets = get_read_proxy_targets(bucket, object, &opts).await;
if targets.is_empty() {
return None;
}
record_replication_proxy(bucket, api, is_err).await;
for target in targets {
match target
.get_object_tagging(&target.bucket, object, opts.version_id.clone())
.await
{
Ok(remote) => {
// MinIO-aligned accounting: one total per proxy attempt,
// one failed when no target served it.
record_replication_proxy(bucket, "GetObjectTagging", false).await;
return Some(
remote
.tag_set
.into_iter()
.map(|tag| Tag {
key: Some(tag.key),
value: Some(tag.value),
})
.collect(),
);
}
Err(err) if Self::proxy_sdk_error_is_not_found(&err) => {
debug!(bucket, object, arn = %target.arn, "tagging proxy: target does not have the object");
}
Err(err) => {
warn!(bucket, object, arn = %target.arn, error = %err, "tagging proxy: GetObjectTagging against replication target failed");
}
}
}
record_replication_proxy(bucket, "GetObjectTagging", true).await;
None
}
/// Apply a PutObjectTagging for an object missing locally on a
/// replication target (MinIO `proxyTaggingToRepTarget`).
async fn proxy_put_object_tagging(
bucket: &str,
object: &str,
version_id: Option<String>,
headers: &http::HeaderMap,
tag_set: &TagSet,
) -> Option<()> {
let opts = Self::tagging_proxy_opts(bucket, object, version_id, headers).await?;
let mut tagging = aws_sdk_s3::types::Tagging::builder();
for tag in tag_set {
let sdk_tag = aws_sdk_s3::types::Tag::builder()
.key(tag.key.clone().unwrap_or_default())
.value(tag.value.clone().unwrap_or_default())
.build()
.ok()?;
tagging = tagging.tag_set(sdk_tag);
}
let tagging = tagging.build().ok()?;
let targets = get_read_proxy_targets(bucket, object, &opts).await;
if targets.is_empty() {
return None;
}
for target in targets {
match target
.put_object_tagging(&target.bucket, object, opts.version_id.clone(), tagging.clone())
.await
{
Ok(_) => {
// MinIO-aligned accounting: one total per proxy attempt,
// one failed when no target served it.
record_replication_proxy(bucket, "PutObjectTagging", false).await;
return Some(());
}
Err(err) if Self::proxy_sdk_error_is_not_found(&err) => {
debug!(bucket, object, arn = %target.arn, "tagging proxy: target does not have the object");
}
Err(err) => {
warn!(bucket, object, arn = %target.arn, error = %err, "tagging proxy: PutObjectTagging against replication target failed");
}
}
}
record_replication_proxy(bucket, "PutObjectTagging", true).await;
None
}
/// Apply a DeleteObjectTagging for an object missing locally on a
/// replication target (MinIO `proxyTaggingToRepTarget`).
async fn proxy_delete_object_tagging(
bucket: &str,
object: &str,
version_id: Option<String>,
headers: &http::HeaderMap,
) -> Option<()> {
let opts = Self::tagging_proxy_opts(bucket, object, version_id, headers).await?;
let targets = get_read_proxy_targets(bucket, object, &opts).await;
if targets.is_empty() {
return None;
}
for target in targets {
match target
.delete_object_tagging(&target.bucket, object, opts.version_id.clone())
.await
{
Ok(_) => {
// MinIO-aligned accounting: one total per proxy attempt,
// one failed when no target served it.
record_replication_proxy(bucket, "DeleteObjectTagging", false).await;
return Some(());
}
Err(err) if Self::proxy_sdk_error_is_not_found(&err) => {
debug!(bucket, object, arn = %target.arn, "tagging proxy: target does not have the object");
}
Err(err) => {
warn!(bucket, object, arn = %target.arn, error = %err, "tagging proxy: DeleteObjectTagging against replication target failed");
}
}
}
record_replication_proxy(bucket, "DeleteObjectTagging", true).await;
None
}
pub async fn get_object_tag_conditions_for_policy(
@@ -447,7 +581,27 @@ impl S3 for FS {
let mut opts = get_opts(&bucket, &object, version_id.clone(), None, &req.headers)
.await
.map_err(ApiError::from)?;
let existing_object_info = store.get_object_info(&bucket, &object, &opts).await.map_err(ApiError::from)?;
let existing_object_info = match store.get_object_info(&bucket, &object, &opts).await {
Ok(info) => info,
Err(e) => {
// Replication lag window: apply the tagging delete on a
// replication target that already has the object
// (backlog#1675 P1-5). No local object exists, so no bucket
// notification event is emitted for the proxied write.
if (is_err_object_not_found(&e) || is_err_version_not_found(&e))
&& Self::proxy_delete_object_tagging(&bucket, &object, version_id.clone(), &req.headers)
.await
.is_some()
{
counter!("rustfs_delete_object_tagging_success").increment(1);
let duration = start_time.elapsed();
histogram!("rustfs_object_tagging_operation_duration_seconds", "operation" => "delete")
.record(duration.as_secs_f64());
return Ok(S3Response::new(DeleteObjectTaggingOutput { version_id }));
}
return Err(ApiError::from(e).into());
}
};
let dsc = must_replicate_metadata(
&bucket,
&object,
@@ -470,7 +624,6 @@ impl S3 for FS {
}
let delete_tags_result = store.delete_object_tags(&bucket, &object, &opts).await;
Self::record_replication_tagging_metric(&bucket, &object, "DeleteObjectTagging", delete_tags_result.is_err()).await;
let object_info = delete_tags_result.map_err(|e| {
error!(
component = LOG_COMPONENT_STORAGE,
@@ -928,32 +1081,49 @@ impl S3 for FS {
..Default::default()
};
let tags_result = store.get_object_tags(bucket, object, &opts).await;
Self::record_replication_tagging_metric(bucket, object, "GetObjectTagging", tags_result.is_err()).await;
let tags = tags_result.map_err(|e| {
if is_err_object_not_found(&e) {
debug!(
let tags = match store.get_object_tags(bucket, object, &opts).await {
Ok(tags) => tags,
Err(e) => {
// Replication lag window: the object may exist on a
// replication target even though it is missing locally —
// proxy the tagging read there (backlog#1675 P1-5).
if (is_err_object_not_found(&e) || is_err_version_not_found(&e))
&& let Some(tag_set) =
Self::proxy_get_object_tagging(bucket, object, req.input.version_id.clone(), &req.headers).await
{
counter!("rustfs_get_object_tagging_success").increment(1);
let duration = start_time.elapsed();
histogram!("rustfs_object_tagging_operation_duration_seconds", "operation" => "get")
.record(duration.as_secs_f64());
return Ok(S3Response::new(GetObjectTaggingOutput {
tag_set,
version_id: req.input.version_id.clone(),
}));
}
if is_err_object_not_found(&e) {
debug!(
component = LOG_COMPONENT_STORAGE,
subsystem = LOG_SUBSYSTEM_TAGGING,
event = "object_tagging_not_found",
bucket = %bucket,
object = %object,
error = %e,
"Object tags not found"
);
return Err(s3_error!(NoSuchKey));
}
error!(
component = LOG_COMPONENT_STORAGE,
subsystem = LOG_SUBSYSTEM_TAGGING,
event = "object_tagging_not_found",
event = "object_tagging_get_failed",
bucket = %bucket,
object = %object,
error = %e,
"Object tags not found"
"Failed to load object tags"
);
return s3_error!(NoSuchKey);
return Err(ApiError::from(e).into());
}
error!(
component = LOG_COMPONENT_STORAGE,
subsystem = LOG_SUBSYSTEM_TAGGING,
event = "object_tagging_get_failed",
bucket = %bucket,
object = %object,
error = %e,
"Failed to load object tags"
);
ApiError::from(e).into()
})?;
};
let tag_set = decode_tags(tags.as_str());
debug!(
@@ -1629,14 +1799,36 @@ impl S3 for FS {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
};
let tags = encode_tags(tagging.tag_set);
let tags = encode_tags(tagging.tag_set.clone());
debug!("Encoded tags: {}", tags);
let version_id = req.input.version_id.clone();
let mut opts = get_opts(&bucket, &object, version_id.clone(), None, &req.headers)
.await
.map_err(ApiError::from)?;
let existing_object_info = store.get_object_info(&bucket, &object, &opts).await.map_err(ApiError::from)?;
let existing_object_info = match store.get_object_info(&bucket, &object, &opts).await {
Ok(info) => info,
Err(e) => {
// Replication lag window: apply the tagging update on a
// replication target that already has the object
// (backlog#1675 P1-5). No local object exists, so no bucket
// notification event is emitted for the proxied write.
if (is_err_object_not_found(&e) || is_err_version_not_found(&e))
&& Self::proxy_put_object_tagging(&bucket, &object, version_id.clone(), &req.headers, &tagging.tag_set)
.await
.is_some()
{
counter!("rustfs_put_object_tagging_success").increment(1);
let duration = start_time.elapsed();
histogram!("rustfs_object_tagging_operation_duration_seconds", "operation" => "put")
.record(duration.as_secs_f64());
return Ok(S3Response::new(PutObjectTaggingOutput {
version_id: req.input.version_id.clone(),
}));
}
return Err(ApiError::from(e).into());
}
};
let dsc = must_replicate_metadata(
&bucket,
&object,
@@ -1659,7 +1851,6 @@ impl S3 for FS {
}
let put_tags_result = store.put_object_tags(&bucket, &object, &tags, &opts).await;
Self::record_replication_tagging_metric(&bucket, &object, "PutObjectTagging", put_tags_result.is_err()).await;
let object_info = put_tags_result.map_err(|e| {
error!("Failed to put object tags: {}", e);
counter!("rustfs_put_object_tagging_failure").increment(1);
+18 -18
View File
@@ -57,24 +57,24 @@ pub(crate) use storage_api::{
QuotaError, RUSTFS_META_BUCKET, RawFileInfo, ReadMultipleReq, ReadMultipleResp, ReadOptions, RenameDataResp,
ReplicationStats, ReplicationStatusType, Result, SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC,
StorageDeletedObject, StorageDiskRpcExt, StorageError, StorageGetObjectReader, StorageObjectInfo, StorageObjectOptions,
StorageObjectToDelete, StoragePeerS3ClientExt, StoragePutObjReader, StorageReplicationConfigExt, StorageVersioningConfigExt,
TONIC_RPC_PREFIX, TierConfigMgr, UpdateMetadataOpts, VolumeInfo, WalkDirOptions, WorkloadAdmissionSnapshotProviderRef,
WriteEncryption, WritePlan, access_consumer, add_object_lock_years, all_local_disk, all_local_disk_path,
check_retention_for_modification, collect_local_metrics, compression_metadata_value, contract, decode_tags,
decode_tags_to_map, delete_bucket_metadata_config, delete_bucket_metadata_config_if_incarnation, disk_drive_path,
disk_endpoint, ecfs_consumer, ecfs_extend_consumer, ecstore_admin, ecstore_bucket, ecstore_capacity, ecstore_client,
ecstore_cluster, ecstore_compression, ecstore_config, ecstore_data_usage, ecstore_disk, ecstore_error, ecstore_event,
ecstore_layout, ecstore_metrics, ecstore_notification, ecstore_rebalance, ecstore_rio, ecstore_rpc, ecstore_set_disk,
ecstore_storage, ecstore_tier, encode_tags, find_local_disk_by_ref, get_bucket_accelerate_config, get_bucket_cors_config,
get_bucket_logging_config, get_bucket_metadata, get_bucket_notification_config, get_bucket_object_lock_config,
get_bucket_replication_config, get_bucket_request_payment_config, get_bucket_sse_config, get_bucket_website_config,
get_local_server_property, get_lock_acquire_timeout, head_prefix_consumer, helper_consumer, init_background_replication,
init_bucket_metadata_sys, init_ecstore_config, init_local_disks_with_instance_ctx, init_lock_clients,
is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found, is_valid_storage_class, options_consumer,
prewarm_local_disk_id_map_with_instance_ctx, read_config, record_replication_proxy, rpc_consumer, runtime_sources_consumer,
s3_api_consumer, serialize, table_catalog_path_hash, to_s3s_etag, topology_snapshot_from_endpoint_pools_with_capabilities,
try_migrate_bucket_metadata, try_migrate_iam_config, try_migrate_server_config, update_bucket_metadata_config,
update_bucket_metadata_config_if_incarnation, verify_rpc_signature, wrap_reader,
StorageObjectToDelete, StoragePeerS3ClientExt, StoragePutObjReader, StorageVersioningConfigExt, TONIC_RPC_PREFIX,
TierConfigMgr, UpdateMetadataOpts, VolumeInfo, WalkDirOptions, WorkloadAdmissionSnapshotProviderRef, WriteEncryption,
WritePlan, access_consumer, add_object_lock_years, all_local_disk, all_local_disk_path, check_retention_for_modification,
collect_local_metrics, compression_metadata_value, contract, decode_tags, decode_tags_to_map, delete_bucket_metadata_config,
delete_bucket_metadata_config_if_incarnation, disk_drive_path, disk_endpoint, ecfs_consumer, ecfs_extend_consumer,
ecstore_admin, ecstore_bucket, ecstore_capacity, ecstore_client, ecstore_cluster, ecstore_compression, ecstore_config,
ecstore_data_usage, ecstore_disk, ecstore_error, ecstore_event, ecstore_layout, ecstore_metrics, ecstore_notification,
ecstore_rebalance, ecstore_rio, ecstore_rpc, ecstore_set_disk, ecstore_storage, ecstore_tier, encode_tags,
find_local_disk_by_ref, get_bucket_accelerate_config, get_bucket_cors_config, get_bucket_logging_config, get_bucket_metadata,
get_bucket_notification_config, get_bucket_object_lock_config, get_bucket_request_payment_config, get_bucket_sse_config,
get_bucket_website_config, get_local_server_property, get_lock_acquire_timeout, head_prefix_consumer, helper_consumer,
init_background_replication, init_bucket_metadata_sys, init_ecstore_config, init_local_disks_with_instance_ctx,
init_lock_clients, is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found, is_valid_storage_class,
options_consumer, prewarm_local_disk_id_map_with_instance_ctx, read_config, record_replication_proxy, rpc_consumer,
runtime_sources_consumer, s3_api_consumer, serialize, table_catalog_path_hash, to_s3s_etag,
topology_snapshot_from_endpoint_pools_with_capabilities, try_migrate_bucket_metadata, try_migrate_iam_config,
try_migrate_server_config, update_bucket_metadata_config, update_bucket_metadata_config_if_incarnation, verify_rpc_signature,
wrap_reader,
};
#[cfg(test)]
+93 -3
View File
@@ -19,9 +19,10 @@ use http::{HeaderMap, HeaderValue};
use rustfs_utils::http::{
AMZ_BUCKET_REPLICATION_STATUS, SUFFIX_FORCE_DELETE, SUFFIX_OBJECTLOCK_LEGALHOLD_TIMESTAMP,
SUFFIX_OBJECTLOCK_RETENTION_TIMESTAMP, SUFFIX_REPLICATION_ACTUAL_OBJECT_SIZE, SUFFIX_REPLICATION_SSEC_CRC,
SUFFIX_SOURCE_DELETEMARKER, SUFFIX_SOURCE_ETAG, SUFFIX_SOURCE_MTIME, SUFFIX_SOURCE_REPLICATION_LEGALHOLD_TIMESTAMP,
SUFFIX_SOURCE_REPLICATION_REQUEST, SUFFIX_SOURCE_REPLICATION_RETENTION_TIMESTAMP,
SUFFIX_SOURCE_REPLICATION_TAGGING_TIMESTAMP, SUFFIX_SOURCE_VERSION_ID, SUFFIX_TAGGING_TIMESTAMP, get_header,
SUFFIX_SOURCE_DELETEMARKER, SUFFIX_SOURCE_ETAG, SUFFIX_SOURCE_MTIME, SUFFIX_SOURCE_PROXY_REQUEST,
SUFFIX_SOURCE_REPLICATION_LEGALHOLD_TIMESTAMP, SUFFIX_SOURCE_REPLICATION_REQUEST,
SUFFIX_SOURCE_REPLICATION_RETENTION_TIMESTAMP, SUFFIX_SOURCE_REPLICATION_TAGGING_TIMESTAMP, SUFFIX_SOURCE_VERSION_ID,
SUFFIX_TAGGING_TIMESTAMP, get_header,
header_compat::{MINIO_ENCRYPTION_PREFIX, RUSTFS_ENCRYPTION_PREFIX},
insert_header_map, insert_str,
metadata_compat::{MINIO_INTERNAL_PREFIX, RUSTFS_INTERNAL_PREFIX},
@@ -276,6 +277,19 @@ pub async fn get_opts(
// Background scanner still performs full integrity checks asynchronously.
opts.skip_verify_bitrot = get_skip_verify_bitrot();
// Anti-loop markers for the replication read proxy
// (`{x-rustfs-,x-minio-}source-proxy-request` header family).
// MinIO semantics: the header being PRESENT at all (`ProxyHeaderSet`)
// disables proxying, whatever its value — a peer's replication worker
// sends "false" on its convergence HEADs so the receiver answers locally
// instead of proxying the miss back (a proxied echo would fake
// convergence and the object would never replicate). Deliberately not
// gated on replication authorization: the header only disables proxying
// (it grants nothing).
let proxy_header = get_header(headers, SUFFIX_SOURCE_PROXY_REQUEST);
opts.proxy_header_set = proxy_header.is_some();
opts.proxy_request = proxy_header.map(|v| v.as_ref() == "true").unwrap_or_default();
fill_conditional_writes_opts_from_header(headers, &mut opts)?;
Ok(opts)
@@ -2544,4 +2558,80 @@ mod tests {
}
}
}
/// The replication read-proxy anti-loop markers must be honored under
/// both interop prefixes (a MinIO peer sends x-minio-, a RustFS peer
/// sends both). `proxy_request` is set only for the literal value
/// "true", while `proxy_header_set` (MinIO `ProxyHeaderSet`) is set by
/// the header's mere presence — "false" (the replication worker's
/// convergence-HEAD marker) and arbitrary values included — so the
/// selector refuses to proxy either way.
#[tokio::test]
async fn test_get_opts_parses_source_proxy_request_under_both_prefixes() {
for header_name in ["x-rustfs-source-proxy-request", "x-minio-source-proxy-request"] {
let mut headers = HeaderMap::new();
headers.insert(header_name, HeaderValue::from_static("true"));
let opts = get_opts("test-bucket", "test-object", None, None, &headers)
.await
.expect("get_opts should succeed");
assert!(opts.proxy_request, "{header_name} must set opts.proxy_request");
assert!(opts.proxy_header_set, "{header_name} must set opts.proxy_header_set");
}
let opts = get_opts("test-bucket", "test-object", None, None, &HeaderMap::new())
.await
.expect("get_opts should succeed");
assert!(!opts.proxy_request, "absent header must leave proxy_request off");
assert!(!opts.proxy_header_set, "absent header must leave proxy_header_set off");
for (header_name, value) in [
("x-minio-source-proxy-request", "false"),
("x-rustfs-source-proxy-request", "false"),
("x-minio-source-proxy-request", "anything-else"),
] {
let mut headers = HeaderMap::new();
headers.insert(header_name, HeaderValue::from_static(value));
let opts = get_opts("test-bucket", "test-object", None, None, &headers)
.await
.expect("get_opts should succeed");
assert!(!opts.proxy_request, "{header_name}: non-'true' value must leave proxy_request off");
assert!(
opts.proxy_header_set,
"{header_name}: value {value:?} must still set proxy_header_set (presence disables proxying)"
);
}
}
/// Pin that the source-proxy-request transport family cannot be
/// materialized as bare stored metadata via an `x-*-meta-` disguise: the
/// reserved-key namespacing (`x-rustfs-source-` / `x-minio-source-`
/// prefixes in `is_reserved_user_metadata_key`) must keep covering it.
#[test]
fn test_source_proxy_request_family_is_reserved_user_metadata() {
let mut headers = HeaderMap::new();
headers.insert("x-amz-meta-x-minio-source-proxy-request", HeaderValue::from_static("true"));
headers.insert("x-rustfs-meta-x-rustfs-source-proxy-request", HeaderValue::from_static("true"));
// The bare transport header itself is not a user-metadata prefix and
// must never land in stored metadata at all.
headers.insert("x-minio-source-proxy-request", HeaderValue::from_static("true"));
let metadata = extract_metadata(&headers);
assert!(
!metadata.contains_key("x-minio-source-proxy-request"),
"bare source-proxy-request key must not be storable: {metadata:?}"
);
assert!(
!metadata.contains_key("x-rustfs-source-proxy-request"),
"bare source-proxy-request key must not be storable: {metadata:?}"
);
assert!(
metadata.contains_key("x-amz-meta-x-minio-source-proxy-request"),
"disguised key must be namespaced back under x-amz-meta-: {metadata:?}"
);
assert!(
metadata.contains_key("x-amz-meta-x-rustfs-source-proxy-request"),
"disguised key must be namespaced back under x-amz-meta-: {metadata:?}"
);
}
}
+8 -18
View File
@@ -805,6 +805,10 @@ impl StorageReplicationStatsHandle {
proxy_head_failed: metrics.proxied.head_failed,
proxy_put_tag_total: metrics.proxied.put_tag_total,
proxy_put_tag_failed: metrics.proxied.put_tag_failed,
proxy_get_tag_total: metrics.proxied.get_tag_total,
proxy_get_tag_failed: metrics.proxied.get_tag_failed,
proxy_delete_tag_total: metrics.proxied.delete_tag_total,
proxy_delete_tag_failed: metrics.proxied.delete_tag_failed,
replica_size: metrics.replica_size,
replica_count: metrics.replica_count,
}
@@ -841,6 +845,10 @@ pub(crate) struct ReplicationSiteMetricsSnapshot {
pub(crate) proxy_head_failed: i64,
pub(crate) proxy_put_tag_total: i64,
pub(crate) proxy_put_tag_failed: i64,
pub(crate) proxy_get_tag_total: i64,
pub(crate) proxy_get_tag_failed: i64,
pub(crate) proxy_delete_tag_total: i64,
pub(crate) proxy_delete_tag_failed: i64,
pub(crate) replica_size: i64,
pub(crate) replica_count: i64,
}
@@ -1487,12 +1495,6 @@ pub(crate) async fn get_bucket_object_lock_config(
ecstore_bucket::metadata_sys::get_object_lock_config(bucket).await
}
pub(crate) async fn get_bucket_replication_config(
bucket: &str,
) -> Result<(s3s::dto::ReplicationConfiguration, time::OffsetDateTime)> {
ecstore_bucket::metadata_sys::get_replication_config(bucket).await
}
pub(crate) async fn persist_force_delete_intent(
api: Arc<ECStore>,
entry: ecstore_bucket::replication::MrfReplicateEntry,
@@ -1838,18 +1840,6 @@ pub(crate) async fn find_local_disk_by_ref(disk_ref: &str) -> Option<DiskStore>
ecstore_storage::find_local_disk_by_ref(disk_ref).await
}
pub(crate) trait StorageReplicationConfigExt {
fn has_active_rules(&self, prefix: &str, recursive: bool) -> bool;
}
impl StorageReplicationConfigExt for s3s::dto::ReplicationConfiguration {
fn has_active_rules(&self, prefix: &str, recursive: bool) -> bool {
<s3s::dto::ReplicationConfiguration as ecstore_bucket::replication::ReplicationConfigurationExt>::has_active_rules(
self, prefix, recursive,
)
}
}
pub(crate) trait StorageVersioningConfigExt {
fn enabled(&self) -> bool;
}