fix(replication): harden delete marker HEAD convergence (#7993)

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
Hauser
2026-09-18 11:05:47 +08:00
committed by GitHub
parent 415bd2310c
commit 785b44a696
@@ -199,6 +199,38 @@ fn has_raw_status(err: &SdkError<HeadObjectError>, status: u16) -> bool {
err.raw_response().is_some_and(|r| r.status().as_u16() == status)
}
fn raw_header_value(err: &SdkError<HeadObjectError>, name: &str) -> Option<String> {
err.raw_response()
.and_then(|response| response.headers().get(name))
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum DeleteMarkerHeadOutcome {
Missing,
Exists,
Ambiguous,
}
fn classify_delete_marker_head_error(err: &SdkError<HeadObjectError>) -> DeleteMarkerHeadOutcome {
let (is_not_found, code) = err
.as_service_error()
.map(|service_err| (service_err.is_not_found(), service_err.code()))
.unwrap_or((false, None));
// Prefer the raw HTTP status when present. HEAD responses commonly have no
// body, and a contradictory XML code must not turn a 403/5xx into success.
match err.raw_response().map(|response| response.status().as_u16()) {
Some(404) => DeleteMarkerHeadOutcome::Missing,
Some(405) => DeleteMarkerHeadOutcome::Exists,
Some(_) => DeleteMarkerHeadOutcome::Ambiguous,
None if is_not_found => DeleteMarkerHeadOutcome::Missing,
None if !is_retryable_delete_replication_head_error(false, code) => DeleteMarkerHeadOutcome::Exists,
None => DeleteMarkerHeadOutcome::Ambiguous,
}
}
fn metadata_requires_existing_target(op_type: ReplicationType, object_info: &ObjectInfo) -> bool {
op_type == ReplicationType::Metadata
&& object_info
@@ -1809,32 +1841,14 @@ async fn verify_resync_head_result(
(roi.size, None)
}
Err(err) if roi.delete_marker => {
// Verifying a replicated delete marker: only a
// definitive 404/NoSuchKey or 405/MethodNotAllowed
// confirms the marker propagated. Any other
// (retryable/ambiguous) HEAD error leaves the outcome
// unverified, so it must count as failed — not as a
// blanket success (backlog#862 / #799 B13).
// A HEAD carries no body, so the SDK only synthesizes an error
// code for 404; a target answering the delete marker's version
// with 405 (RustFS, MinIO, AWS) leaves `code()` empty and the
// raw status is the only evidence (rustfs/backlog#2479).
// A marker *version purge* (`version_purge_status` set) asks
// for the opposite outcome: the marker must be gone from the
// target, so only absence verifies it and a 405 proves the
// marker is still there.
let retryable = {
let (is_not_found, code) = err
.as_service_error()
.map(|se| (se.is_not_found(), se.code()))
.unwrap_or((false, None));
let not_found = is_not_found || has_raw_status(&err, 404);
if roi.version_purge_status.is_empty() {
let code = code.or_else(|| has_raw_status(&err, 405).then_some("405"));
is_retryable_delete_replication_head_error(not_found, code)
} else {
!not_found
}
// A marker creation is verified by either absence (404) or
// presence (405). A marker-version purge has the opposite goal:
// only absence proves completion. All other outcomes remain
// failed (backlog#862 / #799 B13).
let retryable = if roi.version_purge_status.is_empty() {
matches!(classify_delete_marker_head_error(&err), DeleteMarkerHeadOutcome::Ambiguous)
} else {
!matches!(classify_delete_marker_head_error(&err), DeleteMarkerHeadOutcome::Missing)
};
if retryable {
st.failed_count += 1;
@@ -3553,7 +3567,7 @@ async fn replicate_delete_to_target(
}
}
if dobj.delete_object.delete_marker && dobj.delete_object.delete_marker_version_id.is_some() {
if !is_version_purge && dobj.delete_object.delete_marker && dobj.delete_object.delete_marker_version_id.is_some() {
match head_object_for_worker(
tgt_client.as_ref(),
&tgt_client.bucket,
@@ -3562,20 +3576,48 @@ async fn replicate_delete_to_target(
)
.await
{
Ok(_) => {}
Ok(head) if head.delete_marker == Some(true) => {
// Some targets expose an existing marker as a successful HEAD
// with the marker bit instead of the usual 405. Treat that as
// positive convergence evidence and retain the target version
// id for any later marker purge.
rinfo.target_delete_marker_version_id = head.version_id.filter(|version_id| !version_id.is_empty());
rinfo.replication_status = ReplicationStatusType::Completed;
return rinfo;
}
Ok(_) => {
// A successful HEAD without the delete-marker bit proves that
// the addressed version is a regular object (or that the
// target omitted the marker contract). Creating a marker in
// that state would hide a live version, so fail closed.
rinfo.replication_status = ReplicationStatusType::Failed;
rinfo.error = Some("target version is not a delete marker".to_string());
return rinfo;
}
Err(e) => {
let non_retryable = matches!(
e.as_ref(),
SdkError::ServiceError(service_err)
if is_retryable_delete_replication_head_error(
service_err.err().is_not_found(),
service_err.err().code(),
)
);
if non_retryable {
rinfo.replication_status = ReplicationStatusType::Failed;
rinfo.error = Some(e.to_string());
return rinfo;
match classify_delete_marker_head_error(e.as_ref()) {
DeleteMarkerHeadOutcome::Exists => {
// An explicit HEAD of a delete-marker version is expected
// to answer 405 (often without an SDK error code). That
// is positive evidence the marker already exists, so the
// creation is complete and must not issue another
// versionless DELETE that would mint a second marker.
if let Some(version_id) = raw_header_value(e.as_ref(), "x-amz-version-id") {
rinfo.target_delete_marker_version_id = Some(version_id);
}
rinfo.replication_status = ReplicationStatusType::Completed;
return rinfo;
}
DeleteMarkerHeadOutcome::Missing => {}
DeleteMarkerHeadOutcome::Ambiguous => {
// Only a definitive not-found means the marker is absent
// and may be created. Keep transport and ambiguous
// failures fail-closed instead of creating blindly.
rinfo.replication_status = ReplicationStatusType::Failed;
rinfo.error = Some(e.to_string());
mark_replication_target_offline_if_needed(&tgt_client, &e).await;
return rinfo;
}
}
}
}
@@ -6529,6 +6571,110 @@ mod tests {
}
}
fn marker_creation_dobj(target_arn: &str) -> DeletedObjectReplicationInfo {
let mut state = ReplicationState::default();
state.targets.insert(target_arn.to_string(), ReplicationStatusType::Pending);
DeletedObjectReplicationInfo {
bucket: "source".to_string(),
target_arn: target_arn.to_string(),
op_type: ReplicationType::Heal,
delete_object: ReplicationDeletedObject {
object_name: "gone.txt".to_string(),
delete_marker: true,
delete_marker_version_id: Some(Uuid::new_v4()),
replication_state: Some(state),
..Default::default()
},
..Default::default()
}
}
/// A target that already has the marker answers a versioned HEAD with a
/// bodiless 405. The live worker must settle the marker as completed and
/// avoid minting another marker with a versionless DELETE.
#[tokio::test]
async fn live_delete_marker_worker_accepts_bodiless_405_without_delete() {
let (endpoint, server) = spawn_scripted_target_server(1, |_| {
"HTTP/1.1 405 Method Not Allowed\r\nX-Amz-Version-Id: target-marker\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
.to_string()
});
let target = test_target_client(endpoint);
register_test_target(&target).await;
let dobj = marker_creation_dobj(&target.arn);
let rinfo = replicate_delete_to_target(&dobj, target.clone(), None).await;
assert_eq!(rinfo.replication_status, ReplicationStatusType::Completed, "{rinfo:?}");
assert!(rinfo.error.is_none(), "{rinfo:?}");
assert_eq!(rinfo.target_delete_marker_version_id.as_deref(), Some("target-marker"));
let seen = server.join().expect("test HTTP server should finish");
assert_eq!(seen.len(), 1, "an existing marker must not be re-created: {seen:?}");
assert!(seen[0].starts_with("HEAD "), "{seen:?}");
}
/// A successful HEAD that does not identify a delete marker is not
/// convergence evidence: creating a marker would hide a live version.
#[tokio::test]
async fn live_delete_marker_worker_rejects_successful_regular_head_without_delete() {
let (endpoint, server) = spawn_scripted_target_server(1, |_| empty_response("200 OK"));
let target = test_target_client(endpoint);
register_test_target(&target).await;
let dobj = marker_creation_dobj(&target.arn);
let rinfo = replicate_delete_to_target(&dobj, target.clone(), None).await;
assert_eq!(rinfo.replication_status, ReplicationStatusType::Failed, "{rinfo:?}");
assert_eq!(rinfo.error.as_deref(), Some("target version is not a delete marker"));
let seen = server.join().expect("test HTTP server should finish");
assert_eq!(seen.len(), 1, "a regular object must not be hidden by a new marker: {seen:?}");
assert!(seen[0].starts_with("HEAD "), "{seen:?}");
}
/// A definitive 404 means the marker is absent, so exactly one marker
/// creation DELETE is still required.
#[tokio::test]
async fn live_delete_marker_worker_creates_marker_after_404() {
let (endpoint, server) = spawn_scripted_target_server(2, |line| {
if line.starts_with("HEAD ") {
empty_response("404 Not Found")
} else if line.starts_with("DELETE ") {
empty_response("204 No Content")
} else {
empty_response("500 Unexpected")
}
});
let target = test_target_client(endpoint);
register_test_target(&target).await;
let dobj = marker_creation_dobj(&target.arn);
let rinfo = replicate_delete_to_target(&dobj, target.clone(), None).await;
assert_eq!(rinfo.replication_status, ReplicationStatusType::Completed, "{rinfo:?}");
let seen = server.join().expect("test HTTP server should finish");
assert_eq!(seen.len(), 2, "a missing marker should be created once: {seen:?}");
assert!(seen[0].starts_with("HEAD "), "{seen:?}");
assert!(seen[1].starts_with("DELETE "), "{seen:?}");
assert!(!seen[1].contains("versionId="), "marker creation must remain versionless: {seen:?}");
}
/// Ambiguous service errors must remain failed and must not fall through
/// to a blind marker-creation DELETE.
#[tokio::test]
async fn live_delete_marker_worker_keeps_ambiguous_head_failure_closed() {
let (endpoint, server) = spawn_scripted_target_server(1, |_| empty_response("503 Service Unavailable"));
let target = test_target_client(endpoint);
register_test_target(&target).await;
let dobj = marker_creation_dobj(&target.arn);
let rinfo = replicate_delete_to_target(&dobj, target.clone(), None).await;
assert_eq!(rinfo.replication_status, ReplicationStatusType::Failed, "{rinfo:?}");
assert!(rinfo.error.is_some(), "{rinfo:?}");
let seen = server.join().expect("test HTTP server should finish");
assert_eq!(seen.len(), 1, "an ambiguous HEAD must not trigger DELETE: {seen:?}");
assert!(seen[0].starts_with("HEAD "), "{seen:?}");
}
/// A target answers `HEAD ?versionId=<delete marker>` with a bodiless
/// 405, so the SDK reports no error code. That is the marker having
/// propagated, not a failure (rustfs/backlog#2479).