mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-06 03:59:14 +00:00
fix(replication): address a replicated marker purge by the target's version id (backlog#2290)
A source-side DELETE ?versionId=<marker> replicates as a version purge, and replicate_delete_to_target addressed it by the SOURCE marker id on every target. A generic S3 target answers a DELETE of an unknown versionId with 204 and keeps its marker, so the purge reported success and the marker stayed; the same event also spawned a second delayed-purge watcher that journaled a duplicate intent. Real VMs (R6.1 in backlog#2080) failed on the persisted-id fix alone because this path never consulted the mapping. Resolve the target version through the recorded mapping for marker purges (a corrupt record refuses, as the watcher does; nothing recorded keeps the source-derived id for id-mirroring peers), and do not spawn the delayed watcher for a version purge — that purge is the replication itself and its failures reach the journal as a purge entry.
This commit is contained in:
@@ -20,9 +20,9 @@ pub use rustfs_replication::{
|
||||
pub(crate) use rustfs_replication::{
|
||||
ReplicationDeleteSource, ReplicationMultipartPartInput, ReplicationResyncTargetObject, delete_marker_purge_mrf_entry,
|
||||
delete_marker_purge_version_id, delete_replication_creates_marker, delete_replication_missing_source_decision,
|
||||
delete_replication_object_opts, heal_uses_delete_replication_path, is_object_lock_denied_delete,
|
||||
is_retryable_delete_replication_head_error, is_version_delete_replication, replicate_delete_outcome, replication_etags_match,
|
||||
replication_multipart_complete_actual_size, replication_multipart_part_plan, replication_single_put_size_error,
|
||||
resync_existing_delete_replication_info, resync_target_for_object, should_retry_delete_marker_purge,
|
||||
single_part_replica_etag_mismatch, target_delete_version_id,
|
||||
delete_replication_object_opts, delete_replication_target_version_id, heal_uses_delete_replication_path,
|
||||
is_object_lock_denied_delete, is_retryable_delete_replication_head_error, is_version_delete_replication,
|
||||
replicate_delete_outcome, replication_etags_match, replication_multipart_complete_actual_size,
|
||||
replication_multipart_part_plan, replication_single_put_size_error, resync_existing_delete_replication_info,
|
||||
resync_target_for_object, should_retry_delete_marker_purge, single_part_replica_etag_mismatch,
|
||||
};
|
||||
|
||||
@@ -32,11 +32,11 @@ use super::replication_msgp_boundary::ReplicationMsgpCodec;
|
||||
use super::replication_object_config::{ReplicationConfig, get_replication_config, must_replicate};
|
||||
use super::replication_object_decision_boundary::{
|
||||
MustReplicateOptions, ReplicationMultipartPartInput, delete_marker_purge_mrf_entry, delete_marker_purge_version_id,
|
||||
delete_replication_creates_marker, heal_uses_delete_replication_path, is_object_lock_denied_delete,
|
||||
is_retryable_delete_replication_head_error, is_version_delete_replication, replicate_delete_outcome, replication_etags_match,
|
||||
replication_multipart_complete_actual_size, replication_multipart_part_plan, replication_single_put_size_error,
|
||||
resync_existing_delete_replication_info, should_retry_delete_marker_purge, single_part_replica_etag_mismatch,
|
||||
target_delete_version_id,
|
||||
delete_replication_creates_marker, delete_replication_target_version_id, heal_uses_delete_replication_path,
|
||||
is_object_lock_denied_delete, is_retryable_delete_replication_head_error, is_version_delete_replication,
|
||||
replicate_delete_outcome, replication_etags_match, replication_multipart_complete_actual_size,
|
||||
replication_multipart_part_plan, replication_single_put_size_error, resync_existing_delete_replication_info,
|
||||
should_retry_delete_marker_purge, single_part_replica_etag_mismatch,
|
||||
};
|
||||
use super::replication_queue_boundary::{DeletedObjectReplicationInfo, ReplicationQueueAdmission};
|
||||
use super::replication_resync_boundary::ResyncStatusType;
|
||||
@@ -2051,7 +2051,11 @@ pub(crate) async fn replicate_delete_with_outcome<S: ReplicationStorage>(
|
||||
|
||||
let is_version_purge = is_version_delete_replication(&dobj.delete_object);
|
||||
|
||||
let requires_delayed_purge = should_retry_delete_marker_purge(&dobj.delete_object);
|
||||
// The watcher exists to purge a replicated marker once the SOURCE marker
|
||||
// vanishes. A version purge is that purge already (its failures reach the
|
||||
// journal as a purge entry), so it must not spawn a second watcher that
|
||||
// journals a duplicate intent (backlog#2290).
|
||||
let requires_delayed_purge = should_retry_delete_marker_purge(&dobj.delete_object) && !is_version_purge;
|
||||
|
||||
let (replication_status, prev_status) = if !is_version_purge {
|
||||
(
|
||||
@@ -2761,12 +2765,6 @@ fn unavailable_delete_target_info(dobj: &DeletedObjectReplicationInfo, arn: &str
|
||||
}
|
||||
|
||||
async fn replicate_delete_to_target(dobj: &DeletedObjectReplicationInfo, tgt_client: Arc<TargetClient>) -> ReplicatedTargetInfo {
|
||||
let version_id = if let Some(version_id) = &dobj.delete_object.delete_marker_version_id {
|
||||
version_id.to_owned()
|
||||
} else {
|
||||
dobj.delete_object.version_id.unwrap_or_default()
|
||||
};
|
||||
|
||||
let mut rinfo = dobj
|
||||
.delete_object
|
||||
.replication_state
|
||||
@@ -2799,7 +2797,25 @@ async fn replicate_delete_to_target(dobj: &DeletedObjectReplicationInfo, tgt_cli
|
||||
return rinfo;
|
||||
}
|
||||
|
||||
let version_id = target_delete_version_id(version_id, is_version_purge);
|
||||
// Purging a replicated delete marker addresses the version the target
|
||||
// assigned (recorded when the marker was created there); see
|
||||
// `delete_replication_target_version_id`. A corrupt record is a failure,
|
||||
// not a guess: the entry stays visible until the metadata is repaired.
|
||||
let Some(version_id) = delete_replication_target_version_id(&dobj.delete_object, &tgt_client.arn) else {
|
||||
warn!(
|
||||
event = EVENT_DELETE_MARKER_PURGE_FAILED,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC,
|
||||
bucket = tgt_client.bucket,
|
||||
object = dobj.delete_object.object_name,
|
||||
arn = %tgt_client.arn,
|
||||
reason = "recorded_target_version_inconsistent",
|
||||
"Replicated version purge refused: recorded target delete-marker version metadata is inconsistent"
|
||||
);
|
||||
rinfo.version_purge_status = VersionPurgeStatusType::Failed;
|
||||
rinfo.error = Some("recorded target delete-marker version metadata is inconsistent".to_string());
|
||||
return rinfo;
|
||||
};
|
||||
|
||||
if dobj.delete_object.delete_marker && dobj.delete_object.delete_marker_version_id.is_some() {
|
||||
match head_object_for_worker(
|
||||
|
||||
@@ -253,6 +253,28 @@ pub fn delete_marker_purge_version_id(
|
||||
})
|
||||
}
|
||||
|
||||
/// The version a delete replication addresses on `arn`, or `None` to refuse.
|
||||
///
|
||||
/// A version purge whose purged version is a delete marker must address the
|
||||
/// marker version the TARGET assigned — the recorded mapping, exactly as the
|
||||
/// delayed-purge watcher does. The source-side `DELETE ?versionId=<marker>`
|
||||
/// replicates as such a purge, and a generic S3 target answers a DELETE of an
|
||||
/// unknown versionId with 204 while keeping its marker, so addressing it by
|
||||
/// the source id reported success and left the marker behind (backlog#2290,
|
||||
/// R6.1 on the VMs). Nothing recorded falls back to the source-derived id
|
||||
/// (id-mirroring peers); a corrupt record refuses, as the watcher does.
|
||||
pub fn delete_replication_target_version_id(dobj: &DeletedObject, arn: &str) -> Option<Option<String>> {
|
||||
let is_version_purge = is_version_delete_replication(dobj);
|
||||
if is_version_purge
|
||||
&& !dobj.delete_marker
|
||||
&& let Some(marker) = dobj.delete_marker_version_id
|
||||
{
|
||||
return delete_marker_purge_version_id(dobj.replication_state.as_ref(), arn, marker);
|
||||
}
|
||||
let source_version = dobj.delete_marker_version_id.or(dobj.version_id).unwrap_or_default();
|
||||
Some(target_delete_version_id(source_version, is_version_purge))
|
||||
}
|
||||
|
||||
/// Shape an exhausted purge intent as a marker-creation delete entry. Replay
|
||||
/// reconstructs it with `delete_marker: true`, finds the source marker gone,
|
||||
/// and funnels into the stale-marker branch of `replicate_delete_with_outcome`
|
||||
@@ -273,9 +295,9 @@ mod tests {
|
||||
|
||||
use super::{
|
||||
DeletedObjectReplicationInfo, delete_marker_purge_mrf_entry, delete_marker_purge_version_id,
|
||||
delete_replication_creates_marker, is_object_lock_denied_delete, is_retryable_delete_replication_head_error,
|
||||
is_version_delete_replication, replicate_delete_outcome, resync_existing_delete_replication_info,
|
||||
should_retry_delete_marker_purge, target_delete_version_id,
|
||||
delete_replication_creates_marker, delete_replication_target_version_id, is_object_lock_denied_delete,
|
||||
is_retryable_delete_replication_head_error, is_version_delete_replication, replicate_delete_outcome,
|
||||
resync_existing_delete_replication_info, should_retry_delete_marker_purge, target_delete_version_id,
|
||||
};
|
||||
use crate::storage_api::DeletedObject;
|
||||
use crate::{
|
||||
@@ -741,4 +763,57 @@ mod tests {
|
||||
assert!(!is_object_lock_denied_delete(Some("InternalError"), Some("retention lookup failed")));
|
||||
assert!(!is_object_lock_denied_delete(None, Some("legal hold")));
|
||||
}
|
||||
|
||||
fn purge_of_marker(marker: Uuid, state: Option<ReplicationState>) -> DeletedObject {
|
||||
DeletedObject {
|
||||
object_name: "obj".to_string(),
|
||||
delete_marker: false,
|
||||
delete_marker_version_id: Some(marker),
|
||||
version_id: None,
|
||||
replication_state: state,
|
||||
..Default::default()
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn delete_replication_target_version_id_addresses_recorded_marker_for_purges() {
|
||||
let arn = "arn:minio:replication::generic:photos";
|
||||
let marker = Uuid::new_v4();
|
||||
let mut state = ReplicationState::default();
|
||||
state
|
||||
.target_delete_marker_version_ids
|
||||
.insert(arn.to_string(), "remote-marker".to_string());
|
||||
|
||||
// purge of a replicated marker: the target's own version
|
||||
assert_eq!(
|
||||
delete_replication_target_version_id(&purge_of_marker(marker, Some(state.clone())), arn),
|
||||
Some(Some("remote-marker".to_string()))
|
||||
);
|
||||
// nothing recorded for this arn: the source-derived id (id-mirroring peers)
|
||||
assert_eq!(
|
||||
delete_replication_target_version_id(&purge_of_marker(marker, None), arn),
|
||||
Some(Some(marker.to_string()))
|
||||
);
|
||||
// corrupt record: refuse instead of guessing
|
||||
state.target_delete_marker_version_ids_corrupt = true;
|
||||
assert_eq!(delete_replication_target_version_id(&purge_of_marker(marker, Some(state)), arn), None);
|
||||
|
||||
// marker creation keeps the source id (the target mints its own on a
|
||||
// versionless DELETE; the id only travels in the source header)
|
||||
let creation = DeletedObject {
|
||||
object_name: "obj".to_string(),
|
||||
delete_marker: true,
|
||||
delete_marker_version_id: Some(marker),
|
||||
..Default::default()
|
||||
};
|
||||
assert_eq!(delete_replication_target_version_id(&creation, arn), Some(Some(marker.to_string())));
|
||||
// plain version purge: the source version id
|
||||
let version = Uuid::new_v4();
|
||||
let purge = DeletedObject {
|
||||
object_name: "obj".to_string(),
|
||||
version_id: Some(version),
|
||||
..Default::default()
|
||||
};
|
||||
assert_eq!(delete_replication_target_version_id(&purge, arn), Some(Some(version.to_string())));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -41,9 +41,9 @@ pub use config::{
|
||||
};
|
||||
pub use delete::{
|
||||
DeletedObjectReplicationInfo, delete_marker_purge_mrf_entry, delete_marker_purge_version_id,
|
||||
delete_replication_creates_marker, is_object_lock_denied_delete, is_retryable_delete_replication_head_error,
|
||||
is_version_delete_replication, replicate_delete_outcome, resync_existing_delete_replication_info,
|
||||
should_retry_delete_marker_purge, target_delete_version_id,
|
||||
delete_replication_creates_marker, delete_replication_target_version_id, is_object_lock_denied_delete,
|
||||
is_retryable_delete_replication_head_error, is_version_delete_replication, replicate_delete_outcome,
|
||||
resync_existing_delete_replication_info, should_retry_delete_marker_purge, target_delete_version_id,
|
||||
};
|
||||
pub use filemeta::{
|
||||
NULL_VERSION_ID, REPLICATE_EXISTING, REPLICATE_EXISTING_DELETE, REPLICATE_HEAL, REPLICATE_HEAL_DELETE, REPLICATE_INCOMING,
|
||||
|
||||
Reference in New Issue
Block a user