fix(replication): stop sending versionId query once a target rejects it

Merge the approved fix from pull request #8109.
This commit is contained in:
Nikita Bakun
2026-09-28 13:40:20 +03:00
committed by GitHub
parent d556d4dc94
commit cc7d5e3f5f
+146 -5
View File
@@ -409,6 +409,11 @@ pub struct BucketTargetSys {
/// lock (never held across an await) so the replication worker can record
/// a verdict from inside its synchronous PUT-response audit.
version_identity_map: Arc<std::sync::RwLock<HashMap<String, VersionIdentityCapability>>>,
/// Target ARNs whose remote rejected a PUT / CreateMultipartUpload
/// `versionId` query with `InvalidArgument`; later requests to the same
/// ARN omit the query instead of failing again. Reset alongside
/// `version_identity_map` — same lifetime, same reasoning.
put_version_id_query_rejected: Arc<std::sync::RwLock<HashSet<String>>>,
pub targets_map: Arc<RwLock<HashMap<String, Vec<BucketTarget>>>>,
/// Buckets whose persisted `bucket-targets.json` exists but cannot be
/// decoded (rustfs/backlog#2282). Written under the bucket's update mutex
@@ -458,6 +463,7 @@ impl BucketTargetSys {
arn_remotes_map: Arc::new(RwLock::new(HashMap::new())),
ssec_passthrough_map: Arc::new(RwLock::new(HashMap::new())),
version_identity_map: Arc::new(std::sync::RwLock::new(HashMap::new())),
put_version_id_query_rejected: Arc::new(std::sync::RwLock::new(HashSet::new())),
targets_map: Arc::new(RwLock::new(HashMap::new())),
unreadable_targets: Arc::new(RwLock::new(HashSet::new())),
h_mutex: Arc::new(RwLock::new(HashMap::new())),
@@ -782,10 +788,36 @@ impl BucketTargetSys {
health_map.remove(&target.arn);
ssec_map.remove(&target.arn);
self.forget_version_identity_capability(&target.arn);
self.forget_put_version_id_query_rejected(&target.arn);
}
}
}
/// Whether this target ARN is already known to reject a `versionId`
/// query on PUT / CreateMultipartUpload.
pub fn put_version_id_query_rejected(&self, arn: &str) -> bool {
self.put_version_id_query_rejected
.read()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.contains(arn)
}
/// Record that this target ARN rejected a `versionId` query on PUT /
/// CreateMultipartUpload, so later requests to it omit the query.
pub fn record_put_version_id_query_rejected(&self, arn: &str) {
self.put_version_id_query_rejected
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.insert(arn.to_string());
}
fn forget_put_version_id_query_rejected(&self, arn: &str) {
self.put_version_id_query_rejected
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.remove(arn);
}
/// Cached version-identity verdict for a target ARN; `Unknown` until a
/// replication write or a replication-check VersionFidelity probe judged
/// it since the target was built.
@@ -1252,6 +1284,7 @@ impl BucketTargetSys {
.is_none_or(|edited| !same_replication_service(edited, &target))
{
self.forget_version_identity_capability(&target.arn);
self.forget_put_version_id_query_rejected(&target.arn);
}
self.update_bandwidth_limit(bucket, &target.arn, 0);
}
@@ -1449,8 +1482,13 @@ pub fn resolve_delete_api_version_id(version_id: Option<String>, opts: &RemoveOb
/// "null" (the delete path already maps it via `target_delete_version_id`),
/// and an empty id means the source object carries no version: send no query
/// so an unversioned target stays valid.
fn resolve_put_api_version_id(source_version_id: &str) -> Option<&str> {
if source_version_id.is_empty() {
///
/// A standard S3 target (AWS S3 included) instead rejects this query on PUT /
/// CreateMultipartUpload outright with `InvalidArgument: This operation does
/// not accept a version-id`. Once a target ARN has been observed to reject
/// it (`arn_rejected_query`), it is never sent to that ARN again.
fn resolve_put_api_version_id(source_version_id: &str, arn_rejected_query: bool) -> Option<&str> {
if arn_rejected_query || source_version_id.is_empty() {
None
} else if Uuid::parse_str(source_version_id).is_ok_and(|uuid| uuid.is_nil()) {
Some(rustfs_filemeta::NULL_VERSION_ID)
@@ -1459,6 +1497,13 @@ fn resolve_put_api_version_id(source_version_id: &str) -> Option<&str> {
}
}
/// True when a PUT / CreateMultipartUpload service error is the standard S3
/// rejection of a `versionId` query on those operations (AWS S3 and other
/// strict S3-compatible targets; MinIO/RustFS peers accept the query).
fn is_put_version_id_query_rejected(code: Option<&str>, message: Option<&str>) -> bool {
code == Some("InvalidArgument") && message.is_some_and(|m| m.to_ascii_lowercase().contains("version-id"))
}
/// 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
@@ -2251,7 +2296,9 @@ impl TargetClient {
if !version_id.is_empty() {
insert_header(&mut headers, SUFFIX_SOURCE_VERSION_ID, &version_id);
}
let api_version_id = resolve_put_api_version_id(&version_id).map(ToOwned::to_owned);
let api_version_id =
resolve_put_api_version_id(&version_id, BucketTargetSys::get().put_version_id_query_rejected(&self.arn))
.map(ToOwned::to_owned);
// A PUT carrying Object Lock parameters must also carry Content-MD5 or
// an x-amz-checksum-* header on AWS-compatible targets (rustfs#7082).
@@ -2333,6 +2380,9 @@ impl TargetClient {
SdkError::ServiceError(service_err) => {
let err = service_err.into_err();
let meta = err.meta();
if is_put_version_id_query_rejected(meta.code(), meta.message()) {
BucketTargetSys::get().record_put_version_id_query_rejected(&self.arn);
}
let error = match (meta.code(), meta.message()) {
(Some(code), Some(message)) => format!("put_object failed: {code}: {message}"),
(Some(code), None) => format!("put_object failed: {code}"),
@@ -2369,7 +2419,9 @@ impl TargetClient {
let version_id = opts.internal.source_version_id.clone();
// The remote version of a multipart replication is decided at initiate
// time; CompleteMultipartUpload does not read a versionId.
let api_version_id = resolve_put_api_version_id(&version_id).map(ToOwned::to_owned);
let api_version_id =
resolve_put_api_version_id(&version_id, BucketTargetSys::get().put_version_id_query_rejected(&self.arn))
.map(ToOwned::to_owned);
match self
.client
@@ -2395,7 +2447,15 @@ impl TargetClient {
.await
{
Ok(res) => Ok(res.upload_id.unwrap_or_default()),
Err(e) => Err(e.into()),
Err(e) => {
if let SdkError::ServiceError(service_err) = &e {
let meta = service_err.err().meta();
if is_put_version_id_query_rejected(meta.code(), meta.message()) {
BucketTargetSys::get().record_put_version_id_query_rejected(&self.arn);
}
}
Err(e.into())
}
}
}
@@ -3597,6 +3657,87 @@ mod tests {
assert_eq!(sys.version_identity_capability(arn), VersionIdentityCapability::Unknown);
}
#[test]
fn put_version_id_query_rejection_is_per_arn_and_forgotten_with_the_target() {
let sys = BucketTargetSys::default();
let arn = "arn:rustfs:replication:us-east-1:bucket:versionid";
assert!(!sys.put_version_id_query_rejected(arn));
sys.record_put_version_id_query_rejected(arn);
assert!(sys.put_version_id_query_rejected(arn));
assert!(!sys.put_version_id_query_rejected("other"));
// A rebuilt target may point at a different service.
sys.forget_put_version_id_query_rejected(arn);
assert!(!sys.put_version_id_query_rejected(arn));
}
#[tokio::test]
async fn update_all_targets_forgets_put_version_id_rejection_when_endpoint_changes() {
// Editing a target's ARN to point at a different endpoint must not
// carry over a stale "this remote rejects versionId" verdict from the
// service it used to point at (rustfs review follow-up on the
// versionId-on-PUT fix): otherwise a peer that legitimately needs the
// query silently loses version-id fidelity until the process restarts.
let sys = BucketTargetSys::default();
let arn = "arn:rustfs:replication:us-east-1:bucket:versionid";
let target = |endpoint: &str| BucketTarget {
arn: arn.to_string(),
endpoint: endpoint.to_string(),
target_bucket: "target-bucket".to_string(),
region: "us-east-1".to_string(),
credentials: Some(Credentials {
access_key: "access".to_string(),
secret_key: "secret".to_string(),
session_token: None,
expiration: None,
}),
..Default::default()
};
sys.update_all_targets(
"bucket",
Some(&BucketTargets {
targets: vec![target("aws-s3.example:9000")],
}),
)
.await;
sys.record_put_version_id_query_rejected(arn);
assert!(sys.put_version_id_query_rejected(arn));
// Same ARN, different endpoint: the update path (not `delete`) is the
// only thing that runs here, so the reset must come from
// `update_all_targets_locked` itself.
sys.update_all_targets(
"bucket",
Some(&BucketTargets {
targets: vec![target("minio-peer.example:9000")],
}),
)
.await;
assert!(
!sys.put_version_id_query_rejected(arn),
"editing the target to a different endpoint must forget the stale rejection verdict"
);
}
#[test]
fn resolve_put_api_version_id_omits_query_once_rejected() {
let version_id = Uuid::new_v4().to_string();
assert_eq!(resolve_put_api_version_id(&version_id, false), Some(version_id.as_str()));
assert_eq!(resolve_put_api_version_id(&version_id, true), None);
}
#[test]
fn is_put_version_id_query_rejected_matches_aws_invalid_argument() {
assert!(is_put_version_id_query_rejected(
Some("InvalidArgument"),
Some("This operation does not accept a version-id.")
));
assert!(!is_put_version_id_query_rejected(Some("AccessDenied"), Some("access denied")));
assert!(!is_put_version_id_query_rejected(Some("InvalidArgument"), Some("bucket name is invalid")));
assert!(!is_put_version_id_query_rejected(None, None));
}
#[test]
fn endpoint_health_key_preserves_explicit_port() {
let url = Url::parse("https://remote.example:9443").expect("url should parse");