From af802756c5d539972c4770f8899d594f26685765 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=94=90=E5=B0=8F=E9=B8=AD?= Date: Sun, 16 Aug 2026 01:02:42 +0800 Subject: [PATCH] fix(site-replication): only a repair settles snapshot-escalated retry entries MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Second review round: every iam-item / bucket-meta delivery shares a constant path, so any later successful single-item delivery (a Bob update) dequeued the escalated marker recording a possibly-unreplayed deletion (a failed Alice delete) while the entity still existed remotely. Ordinary settlement now skips escalated entries; only the repair path — the operator's explicit accountability transfer — clears them via dequeue_..._including_escalated. A new hook failure still overwrites the marker and re-arms the drain. Regression covers survive-ordinary-dequeue and repair-clears. --- rustfs/src/admin/handlers/site_replication.rs | 35 ++++++++++++++++++- 1 file changed, 34 insertions(+), 1 deletion(-) diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index 8f305e185..533c7ebd3 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -3956,7 +3956,7 @@ async fn persist_site_replication_repair_task( match failure.as_deref() { Some(error) => upsert_site_replication_retry_event(&mut state.retry_queue, &peer, &path, error, None), None => { - dequeue_site_replication_retry_events(&mut state.retry_queue, &peer, &path); + dequeue_site_replication_retry_events_including_escalated(&mut state.retry_queue, &peer, &path); } } Ok(()) @@ -6044,6 +6044,20 @@ fn dequeue_site_replication_retry_events(queue: &mut Vec, + peer: &PeerInfo, + path: &str, +) -> usize { + let before = queue.len(); + queue.retain(|event| !retry_event_matches(event, peer, path)); + before.saturating_sub(queue.len()) +} + /// Remove the retry events for (peer, path) that `generation` is entitled to /// settle. A successful delivery only proves the peer reached the state the /// delivery carried: while it was in flight another edit can commit, fail its @@ -6063,6 +6077,13 @@ fn settle_site_replication_retry_events( if !retry_event_matches(event, peer, path) { return true; } + // A snapshot-escalated entry records a possibly-unreplayed deletion. + // Collapsed paths are shared by every entity, so a later successful + // delivery of a DIFFERENT item proves nothing about the deleted one — + // only a repair settles it (dequeue_..._including_escalated). + if event.last_error == SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER { + return true; + } match (generation, event.edit_generation) { (Some(settled), Some(failed)) => failed > settled, _ => false, @@ -12024,7 +12045,19 @@ mod tests { classify_site_replication_retry_event(&queue[0]).is_none(), "a snapshot-replayed entry must not be re-sent daily" ); + // Ordinary success dequeues must not clear the marker: collapsed + // paths are shared by every entity, so a successful Bob update + // proves nothing about a failed Alice deletion (second review + // round). + assert_eq!(dequeue_site_replication_retry_events(&mut queue, &target, path), 0); + assert_eq!(queue.len(), 1, "an escalated entry must survive an ordinary delivery success"); + // Only a repair — the operator's accountability transfer — settles it. + assert_eq!(dequeue_site_replication_retry_events_including_escalated(&mut queue, &target, path), 1); + assert!(queue.is_empty()); + // A later hook failure overwrites the marker and re-arms the drain. + let mut queue = vec![drain_event("remote", path, 2, Some(snapshot_at))]; + escalate_site_replication_retry_events_up_to(&mut queue, &target, path, Some(snapshot_at)); upsert_site_replication_retry_event(&mut queue, &target, path, "peer offline", None); assert!(classify_site_replication_retry_event(&queue[0]).is_some());