fix(site-replication): preserve retry accountability

This commit is contained in:
overtrue
2026-08-16 04:54:11 +08:00
parent af802756c5
commit 96a2f9c67f
2 changed files with 438 additions and 64 deletions
+3 -3
View File
@@ -258,7 +258,7 @@ pub struct SRLDAPUser {
pub api_version: Option<String>,
}
#[derive(Debug, Serialize, Deserialize, Default)]
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct SRIAMUser {
#[serde(rename = "accessKey", default)]
pub access_key: String,
@@ -270,7 +270,7 @@ pub struct SRIAMUser {
pub api_version: Option<String>,
}
#[derive(Debug, Serialize, Deserialize, Default)]
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct SRGroupInfo {
#[serde(rename = "updateReq", default)]
pub update_req: GroupAddRemove,
@@ -346,7 +346,7 @@ pub struct SRCredInfo {
pub api_version: Option<String>,
}
#[derive(Debug, Serialize, Deserialize, Default)]
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct SRIAMItem {
#[serde(default)]
pub r#type: String,
+435 -61
View File
@@ -2973,6 +2973,18 @@ fn reconcile_site_replication_wiring() -> std::pin::Pin<Box<dyn std::future::Fut
return;
};
if let Err(err) = migrate_collapsed_retry_queue_paths().await {
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
result = "retry_queue_migration_failed",
error = ?err,
"admin site replication state"
);
return;
}
match load_site_replication_state().await {
Ok(state) => {
if state.pending_endpoint_refresh.is_some() || state.pending_remove.is_some() || state.pending_rotation.is_some()
@@ -6040,6 +6052,73 @@ fn retry_event_matches(event: &SiteReplicationRetryEvent, peer: &PeerInfo, path:
(event.peer_deployment_id == peer.deployment_id || event.peer_endpoint == peer.endpoint) && event.path == path
}
const SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH: &str = "internal:retry-snapshot:iam";
const SITE_REPLICATION_RETRY_BUCKET_METADATA_SNAPSHOT_PATH: &str = "internal:retry-snapshot:bucket-metadata";
fn collapsed_retry_queue_path(path: &str) -> Option<&'static str> {
let base_path = path.split_once('?').map(|(base, _)| base).unwrap_or(path);
match base_path {
"/rustfs/admin/v3/site-replication/peer/iam-item" | SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH => {
Some(SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH)
}
"/rustfs/admin/v3/site-replication/peer/bucket-meta" | SITE_REPLICATION_RETRY_BUCKET_METADATA_SNAPSHOT_PATH => {
Some(SITE_REPLICATION_RETRY_BUCKET_METADATA_SNAPSHOT_PATH)
}
_ => None,
}
}
fn normalize_collapsed_retry_queue_paths(queue: &mut Vec<SiteReplicationRetryEvent>) -> bool {
let mut changed = false;
let mut normalized: Vec<SiteReplicationRetryEvent> = Vec::with_capacity(queue.len());
for mut event in queue.drain(..) {
if let Some(path) = collapsed_retry_queue_path(&event.path)
&& event.path != path
{
event.path = path.to_string();
changed = true;
}
let duplicate = normalized.iter().position(|existing| {
existing.path == event.path
&& (existing.peer_deployment_id == event.peer_deployment_id || existing.peer_endpoint == event.peer_endpoint)
});
let Some(index) = duplicate else {
normalized.push(event);
continue;
};
changed = true;
let existing = &mut normalized[index];
let event_is_newer = match (event.updated_at, existing.updated_at) {
(Some(event), Some(existing)) => event >= existing,
(Some(_), None) => true,
_ => false,
};
if event_is_newer {
let retry_count = existing.retry_count.max(event.retry_count);
*existing = event;
existing.retry_count = retry_count;
} else {
existing.retry_count = existing.retry_count.max(event.retry_count);
}
existing.failed = existing.retry_count >= SITE_REPLICATION_RETRY_FAILED_AFTER;
}
*queue = normalized;
changed
}
async fn migrate_collapsed_retry_queue_paths() -> S3Result<()> {
update_site_replication_state_when_changed(|state| {
Ok(if normalize_collapsed_retry_queue_paths(&mut state.retry_queue) {
StateCommit::Changed(())
} else {
StateCommit::Unchanged(())
})
})
.await
}
fn dequeue_site_replication_retry_events(queue: &mut Vec<SiteReplicationRetryEvent>, peer: &PeerInfo, path: &str) -> usize {
settle_site_replication_retry_events(queue, peer, path, None)
}
@@ -6054,7 +6133,11 @@ fn dequeue_site_replication_retry_events_including_escalated(
path: &str,
) -> usize {
let before = queue.len();
queue.retain(|event| !retry_event_matches(event, peer, path));
let collapsed_path = collapsed_retry_queue_path(path);
queue.retain(|event| {
!retry_event_matches(event, peer, path)
&& !collapsed_path.is_some_and(|collapsed_path| retry_event_matches(event, peer, collapsed_path))
});
before.saturating_sub(queue.len())
}
@@ -6073,10 +6156,17 @@ fn settle_site_replication_retry_events(
generation: Option<u64>,
) -> usize {
let before = queue.len();
let collapsed_path = collapsed_retry_queue_path(path);
queue.retain(|event| {
if !retry_event_matches(event, peer, path) {
return true;
}
// A wire-path success identifies no IAM or bucket-metadata entity.
// This also protects legacy rows until the startup migration moves
// them under their internal snapshot path.
if collapsed_path.is_some() {
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 —
@@ -6099,6 +6189,7 @@ fn upsert_site_replication_retry_event(
error: &str,
generation: Option<u64>,
) {
let path = collapsed_retry_queue_path(path).unwrap_or(path);
let now = OffsetDateTime::now_utc();
let detail = summarize_peer_error_detail(error);
if let Some(event) = queue.iter_mut().find(|event| retry_event_matches(event, peer, path)) {
@@ -6224,10 +6315,170 @@ enum RetryDrainAction {
PeerEdit,
}
#[derive(Clone)]
enum RetrySnapshot {
Iam(Vec<SRIAMItem>),
BucketMetadata(Vec<SRBucketMeta>),
}
impl RetrySnapshot {
fn from_plan(action: &RetryDrainAction, plan: &SiteReplicationBootstrapPlan) -> Option<Self> {
match action {
RetryDrainAction::IamSnapshot => Some(Self::Iam(plan.iam_items.clone())),
RetryDrainAction::BucketMetadataSnapshot => Some(Self::BucketMetadata(plan.bucket_items.clone())),
_ => None,
}
}
fn fingerprint(&self) -> S3Result<Vec<Vec<u8>>> {
let mut payloads = match self {
Self::Iam(items) => items.iter().map(serde_json::to_vec).collect::<Result<Vec<_>, _>>(),
Self::BucketMetadata(items) => items.iter().map(serde_json::to_vec).collect::<Result<Vec<_>, _>>(),
}
.map_err(|err| S3Error::with_message(S3ErrorCode::InternalError, format!("serialize retry snapshot failed: {err}")))?;
payloads.sort_unstable();
Ok(payloads)
}
fn replay_after_change(previous: &Self, fresh: &Self, observed_at: OffsetDateTime) -> Self {
match (previous, fresh) {
(Self::Iam(previous), Self::Iam(fresh)) => {
let fresh_keys: HashSet<IamSnapshotKey> = fresh.iter().filter_map(iam_snapshot_key).collect();
let mut replay = fresh.clone();
for item in previous {
if iam_snapshot_key(item).is_some_and(|key| !fresh_keys.contains(&key)) {
replay.extend(iam_snapshot_tombstones(item, observed_at));
}
}
Self::Iam(replay)
}
(Self::BucketMetadata(previous), Self::BucketMetadata(fresh)) => {
let fresh_keys: HashSet<(&str, &str)> = fresh
.iter()
.map(|item| (item.bucket.as_str(), item.r#type.as_str()))
.collect();
let mut replay = fresh.clone();
for item in previous {
if !fresh_keys.contains(&(item.bucket.as_str(), item.r#type.as_str())) {
replay.push(bucket_metadata_snapshot_tombstone(item, observed_at));
}
}
Self::BucketMetadata(replay)
}
_ => fresh.clone(),
}
}
async fn send(&self, transport: &PeerTransport, access_key: &str, secret_key: &str) -> S3Result<()> {
match self {
Self::Iam(items) => {
for item in items {
SiteReplicationRepairTask::Iam(item)
.send(transport, access_key, secret_key)
.await?;
}
}
Self::BucketMetadata(items) => {
for item in items {
SiteReplicationRepairTask::BucketMetadata(item)
.send(transport, access_key, secret_key)
.await?;
}
}
}
Ok(())
}
}
#[derive(Hash, PartialEq, Eq)]
enum IamSnapshotKey {
Policy(String),
User(String),
Group(String),
PolicyMapping { target: String, user_type: i64, is_group: bool },
}
fn iam_snapshot_key(item: &SRIAMItem) -> Option<IamSnapshotKey> {
match item.r#type.as_str() {
"policy" => Some(IamSnapshotKey::Policy(item.name.clone())),
"iam-user" => item
.iam_user
.as_ref()
.map(|user| IamSnapshotKey::User(user.access_key.clone())),
"group-info" => item
.group_info
.as_ref()
.map(|group| IamSnapshotKey::Group(group.update_req.group.clone())),
"policy-mapping" => item.policy_mapping.as_ref().map(|mapping| IamSnapshotKey::PolicyMapping {
target: mapping.user_or_group.clone(),
user_type: mapping.user_type,
is_group: mapping.is_group,
}),
_ => None,
}
}
fn iam_snapshot_tombstones(item: &SRIAMItem, observed_at: OffsetDateTime) -> Vec<SRIAMItem> {
let mut tombstone = item.clone();
tombstone.updated_at = Some(observed_at);
match item.r#type.as_str() {
"policy" => tombstone.policy = None,
"iam-user" => {
if let Some(user) = tombstone.iam_user.as_mut() {
user.is_delete_req = true;
user.user_req = None;
}
}
"group-info" => {
let Some(group) = tombstone.group_info.as_mut() else {
return Vec::new();
};
group.update_req.is_remove = true;
if group.update_req.members.is_empty() {
return vec![tombstone];
}
let mut delete = tombstone.clone();
if let Some(group) = delete.group_info.as_mut() {
group.update_req.members.clear();
}
return vec![tombstone, delete];
}
"policy-mapping" => {
if let Some(mapping) = tombstone.policy_mapping.as_mut() {
mapping.policy.clear();
}
}
_ => return Vec::new(),
}
vec![tombstone]
}
fn bucket_metadata_snapshot_tombstone(item: &SRBucketMeta, observed_at: OffsetDateTime) -> SRBucketMeta {
SRBucketMeta {
r#type: item.r#type.clone(),
bucket: item.bucket.clone(),
updated_at: Some(observed_at),
expiry_updated_at: Some(observed_at),
api_version: item.api_version.clone(),
..Default::default()
}
}
const SITE_REPLICATION_RETRY_SNAPSHOT_STABILITY_ATTEMPTS: usize = 3;
fn classify_site_replication_retry_event(event: &SiteReplicationRetryEvent) -> Option<RetryDrainAction> {
let snapshot_action = match event.path.as_str() {
SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH => Some(RetryDrainAction::IamSnapshot),
SITE_REPLICATION_RETRY_BUCKET_METADATA_SNAPSHOT_PATH => Some(RetryDrainAction::BucketMetadataSnapshot),
_ => None,
};
if snapshot_action.is_some() && event.last_error != SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER {
return snapshot_action;
}
if event.path.starts_with("internal:") {
// Marker records store payloads in `last_error` (legacy
// pending-endpoint-refresh backup); they are not delivery failures.
// pending-endpoint-refresh backup and snapshot liabilities); they are
// not drainable delivery failures.
return None;
}
if event.last_error == SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER {
@@ -6265,14 +6516,11 @@ fn retry_bucket_name(path: &str) -> Option<String> {
.find_map(|(key, value)| (key == "bucket" && !value.is_empty()).then(|| value.into_owned()))
}
/// A collapsed (constant-path) retry event after a successful snapshot
/// resend is escalated with this marker instead of being cleared: the
/// snapshot replays every entity that still exists, but a failed *deletion*
/// leaves no task in the plan, so remote absence is unproven and the entry
/// must stay operator-visible until a later full delivery or a manual repair
/// settles it. The drain skips marked entries so the once-per-episode
/// snapshot is not re-sent daily; a new hook failure overwrites the marker
/// and re-arms the drain.
/// A collapsed retry event after a stable snapshot resend is escalated with
/// this marker instead of being cleared: the snapshot contains no task for a
/// failed deletion, so remote absence remains operator-visible. Collapsed
/// failures use an internal queue path so ordinary successes and older nodes
/// cannot settle an unrelated entity's liability.
const SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER: &str = "snapshot replayed; a failed deletion cannot be replayed from a snapshot — run site replication repair or re-deliver to settle";
/// Escalate a collapsed retry event after its snapshot resend succeeded,
@@ -6280,30 +6528,51 @@ const SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER: &str = "snapshot replayed
/// failure belongs to a newer local commit the snapshot did not contain and
/// must keep the entry drain-eligible).
fn escalate_site_replication_retry_events_up_to(
queue: &mut [SiteReplicationRetryEvent],
queue: &mut Vec<SiteReplicationRetryEvent>,
peer: &PeerInfo,
path: &str,
snapshot_updated_at: Option<OffsetDateTime>,
) -> usize {
let mut escalated = 0usize;
for event in queue.iter_mut() {
if !retry_event_matches(event, peer, path) {
continue;
}
let newer_failure_recorded = match (event.updated_at, snapshot_updated_at) {
(Some(current), Some(seen)) => current > seen,
(Some(_), None) => true,
(None, _) => false,
};
if newer_failure_recorded {
continue;
}
event.failed = true;
event.retry_count = event.retry_count.max(SITE_REPLICATION_RETRY_FAILED_AFTER);
event.last_error = SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER.to_string();
escalated += 1;
let Some(marker_path) = collapsed_retry_queue_path(path) else {
return 0;
};
if path != marker_path {
queue.retain(|event| {
if !retry_event_matches(event, peer, path) {
return true;
}
matches!((event.updated_at, snapshot_updated_at), (Some(current), Some(seen)) if current > seen)
|| matches!((event.updated_at, snapshot_updated_at), (Some(_), None))
});
}
escalated
let marker_index = queue.iter().position(|event| retry_event_matches(event, peer, marker_path));
let marker_index = marker_index.unwrap_or_else(|| {
queue.push(SiteReplicationRetryEvent {
id: Uuid::new_v4().to_string(),
peer_deployment_id: peer.deployment_id.clone(),
peer_endpoint: peer.endpoint.clone(),
path: marker_path.to_string(),
updated_at: snapshot_updated_at,
..Default::default()
});
queue.len() - 1
});
let event = &mut queue[marker_index];
let newer_failure_recorded = match (event.updated_at, snapshot_updated_at) {
(Some(current), Some(seen)) => current > seen,
(Some(_), None) => true,
(None, _) => false,
};
if newer_failure_recorded && event.last_error != SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER {
return 0;
}
event.failed = true;
event.retry_count = event.retry_count.max(SITE_REPLICATION_RETRY_FAILED_AFTER);
event.last_error = SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER.to_string();
event.updated_at = Some(OffsetDateTime::now_utc());
1
}
async fn escalate_site_replication_retry_event_up_to(peer: &PeerInfo, path: &str, snapshot_updated_at: Option<OffsetDateTime>) {
@@ -6498,33 +6767,30 @@ async fn drain_one_site_replication_retry_event(
) -> S3Result<bool> {
let access_key = &runtime.state.service_account_access_key;
let secret_key = &runtime.service_account_secret_key;
match action {
match action.clone() {
RetryDrainAction::IamSnapshot | RetryDrainAction::BucketMetadataSnapshot => {
let Some(plan) = plan else {
return Ok(false);
};
let tasks: Vec<SiteReplicationRepairTask<'_>> = match action {
RetryDrainAction::IamSnapshot => plan.iam_items.iter().map(SiteReplicationRepairTask::Iam).collect(),
_ => plan
.bucket_items
.iter()
.map(SiteReplicationRepairTask::BucketMetadata)
.collect(),
};
for task in &tasks {
if let Err(err) = task.send(transport, access_key, secret_key).await {
let mut current_snapshot = RetrySnapshot::from_plan(&action, plan).expect("snapshot action has a snapshot");
let mut replay = current_snapshot.clone();
for _ in 0..SITE_REPLICATION_RETRY_SNAPSHOT_STABILITY_ATTEMPTS {
let current_fingerprint = current_snapshot.fingerprint()?;
if let Err(err) = replay.send(transport, access_key, secret_key).await {
enqueue_site_replication_retry_event(peer, &event.path, &err).await;
return Err(err);
}
let fresh_info = build_sr_info(&runtime.state, &runtime.local_peer).await?;
let fresh_plan = site_replication_bootstrap_plan(&fresh_info)?;
let fresh_snapshot = RetrySnapshot::from_plan(&action, &fresh_plan).expect("snapshot action has a snapshot");
if fresh_snapshot.fingerprint()? == current_fingerprint {
escalate_site_replication_retry_event_up_to(peer, &event.path, event.updated_at).await;
return Ok(true);
}
replay = RetrySnapshot::replay_after_change(&current_snapshot, &fresh_snapshot, OffsetDateTime::now_utc());
current_snapshot = fresh_snapshot;
}
// The snapshot replays every entity that still exists, but a
// failed *deletion* leaves no task in the plan — remote absence
// is unproven, so escalate (operator-visible, drain-idle) instead
// of clearing. Conditional on the snapshot timestamp: a hook
// failure recorded while this snapshot was in flight belongs to a
// newer commit and keeps the entry drain-eligible.
escalate_site_replication_retry_event_up_to(peer, &event.path, event.updated_at).await;
Ok(true)
Ok(false)
}
RetryDrainAction::BucketOpReplay { operation, bucket } => {
let Some(plan) = plan else {
@@ -11888,8 +12154,7 @@ mod tests {
/// replay faithfully. IAM / bucket-meta entries collapse per (peer, path)
/// with no body persisted — only a snapshot resend is truthful; bucket
/// makes/replication configs are re-derivable; destructive bucket ops and
/// `internal:` marker records (the pending-endpoint-refresh backup store)
/// are never background-replayed.
/// unrelated `internal:` marker records are never background-replayed.
#[test]
fn test_classify_site_replication_retry_event_actions() {
let now = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp");
@@ -11903,6 +12168,11 @@ mod tests {
classify("/rustfs/admin/v3/site-replication/peer/bucket-meta"),
Some(RetryDrainAction::BucketMetadataSnapshot)
);
assert_eq!(classify(SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH), Some(RetryDrainAction::IamSnapshot));
assert_eq!(
classify(SITE_REPLICATION_RETRY_BUCKET_METADATA_SNAPSHOT_PATH),
Some(RetryDrainAction::BucketMetadataSnapshot)
);
assert_eq!(classify(SITE_REPLICATION_PEER_EDIT_PATH), Some(RetryDrainAction::PeerEdit));
assert_eq!(
classify(
@@ -11936,6 +12206,62 @@ mod tests {
assert_eq!(classify("/rustfs/admin/v3/site-replication/peer/unknown"), None);
}
#[test]
fn test_retry_snapshot_fingerprint_detects_concurrent_iam_change() {
let old = SRIAMItem {
r#type: "policy".to_string(),
name: "readwrite".to_string(),
updated_at: Some(OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp")),
..Default::default()
};
let mut new = old.clone();
new.updated_at = Some(OffsetDateTime::from_unix_timestamp(1_700_000_001).expect("timestamp"));
let sent = RetrySnapshot::Iam(vec![old]);
let changed = RetrySnapshot::Iam(vec![new]);
assert_ne!(sent.fingerprint().unwrap(), changed.fingerprint().unwrap());
}
#[test]
fn test_retry_snapshot_replays_a_concurrent_deletion_as_a_tombstone() {
let observed_at = OffsetDateTime::from_unix_timestamp(1_700_000_010).expect("timestamp");
let policy = SRIAMItem {
r#type: "policy".to_string(),
name: "readwrite".to_string(),
policy: Some(serde_json::json!({"Version": "2012-10-17"})),
..Default::default()
};
let replay =
RetrySnapshot::replay_after_change(&RetrySnapshot::Iam(vec![policy]), &RetrySnapshot::Iam(Vec::new()), observed_at);
let RetrySnapshot::Iam(items) = replay else {
panic!("IAM snapshot expected");
};
assert_eq!(items.len(), 1);
assert_eq!(items[0].name, "readwrite");
assert!(items[0].policy.is_none());
assert_eq!(items[0].updated_at, Some(observed_at));
let bucket = SRBucketMeta {
r#type: "tags".to_string(),
bucket: "photos".to_string(),
tags: Some("encoded-tags".to_string()),
..Default::default()
};
let replay = RetrySnapshot::replay_after_change(
&RetrySnapshot::BucketMetadata(vec![bucket]),
&RetrySnapshot::BucketMetadata(Vec::new()),
observed_at,
);
let RetrySnapshot::BucketMetadata(items) = replay else {
panic!("bucket metadata snapshot expected");
};
assert_eq!(items.len(), 1);
assert_eq!(items[0].bucket, "photos");
assert_eq!(items[0].r#type, "tags");
assert!(items[0].tags.is_none());
assert_eq!(items[0].updated_at, Some(observed_at));
}
/// Exponential backoff gates every attempt: without it a dead peer's
/// entries hit `failed` (retry_count >= 3) within 30 minutes of reconcile
/// ticks and the retry stats lose their signal.
@@ -11973,7 +12299,7 @@ mod tests {
state.retry_queue = vec![
// Eligible: known peer, replayable, past backoff.
drain_event("remote", "/rustfs/admin/v3/site-replication/peer/iam-item", 1, old),
drain_event("remote", SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH, 1, old),
// Not yet due.
drain_event("remote", "/rustfs/admin/v3/site-replication/peer/bucket-meta", 2, Some(now)),
// Unknown peer (removed since the failure was recorded).
@@ -11991,7 +12317,7 @@ mod tests {
let actionable = actionable_site_replication_retry_events(&state, now);
assert_eq!(actionable.len(), 1, "only the due, replayable, known-peer event is actionable");
assert_eq!(actionable[0].path, "/rustfs/admin/v3/site-replication/peer/iam-item");
assert_eq!(actionable[0].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH);
}
/// The drain settles a peer-edit success under a freshly allocated
@@ -12021,9 +12347,19 @@ mod tests {
let snapshot_at = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp");
// Failure re-stamped after the snapshot: untouched, still eligible.
let mut queue = vec![drain_event("remote", path, 2, Some(snapshot_at + time::Duration::seconds(5)))];
let mut queue = vec![drain_event(
"remote",
SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH,
2,
Some(snapshot_at + time::Duration::seconds(5)),
)];
assert_eq!(
escalate_site_replication_retry_events_up_to(&mut queue, &target, path, Some(snapshot_at)),
escalate_site_replication_retry_events_up_to(
&mut queue,
&target,
SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH,
Some(snapshot_at),
),
0
);
assert!(!queue[0].failed);
@@ -12033,12 +12369,18 @@ mod tests {
);
// Unchanged since the snapshot: escalated, kept, drain-idle.
let mut queue = vec![drain_event("remote", path, 2, Some(snapshot_at))];
let mut queue = vec![drain_event(
"remote",
SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH,
2,
Some(snapshot_at),
)];
assert_eq!(
escalate_site_replication_retry_events_up_to(&mut queue, &target, path, Some(snapshot_at)),
1
);
assert_eq!(queue.len(), 1, "the entry must survive until remote absence is proven");
assert_eq!(queue[0].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH);
assert!(queue[0].failed);
assert_eq!(queue[0].last_error, SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER);
assert!(
@@ -12055,6 +12397,16 @@ mod tests {
assert_eq!(dequeue_site_replication_retry_events_including_escalated(&mut queue, &target, path), 1);
assert!(queue.is_empty());
// A failed Alice deletion is stored under the internal path, so a
// successful Bob update on the shared wire path cannot erase it even
// before the drain runs.
let mut queue = Vec::new();
upsert_site_replication_retry_event(&mut queue, &target, path, "alice delete failed", None);
assert_eq!(queue[0].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH);
assert_eq!(dequeue_site_replication_retry_events(&mut queue, &target, path), 0);
assert_eq!(queue.len(), 1);
assert_eq!(queue[0].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH);
// 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));
@@ -12068,13 +12420,34 @@ mod tests {
1
);
// Other (peer, path) entries are untouched.
// A cloned event can disappear during replay; escalation recreates
// the internal liability while leaving another peer's row untouched.
let mut queue = vec![drain_event("other", path, 2, Some(snapshot_at))];
assert_eq!(
escalate_site_replication_retry_events_up_to(&mut queue, &target, path, Some(snapshot_at)),
0
1
);
assert!(!queue[0].failed);
assert_eq!(queue.len(), 2);
assert_eq!(queue[1].peer_deployment_id, target.deployment_id);
assert_eq!(queue[1].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH);
}
#[test]
fn test_collapsed_retry_queue_migration_preserves_legacy_liability() {
let peer = PeerInfo {
deployment_id: "remote-dep".to_string(),
..peer("remote", "https://remote.example.com")
};
let wire_path = "/rustfs/admin/v3/site-replication/peer/iam-item";
let now = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp");
let mut queue = vec![drain_event("remote-dep", wire_path, 2, Some(now))];
assert_eq!(dequeue_site_replication_retry_events(&mut queue, &peer, wire_path), 0);
assert!(normalize_collapsed_retry_queue_paths(&mut queue));
assert_eq!(queue.len(), 1);
assert_eq!(queue[0].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH);
assert!(!normalize_collapsed_retry_queue_paths(&mut queue));
}
#[test]
@@ -13702,6 +14075,7 @@ mod tests {
upsert_site_replication_retry_event(&mut queue, &peer, "/rustfs/admin/v3/site-replication/peer/iam-item", "third", None);
assert_eq!(queue.len(), 1);
assert_eq!(queue[0].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH);
assert_eq!(queue[0].retry_count, SITE_REPLICATION_RETRY_FAILED_AFTER);
assert!(queue[0].failed);
assert_eq!(queue[0].last_error, "third");
@@ -13743,12 +14117,12 @@ mod tests {
);
assert!(queue.is_empty());
// Broadcast paths carry no generation and keep settling unconditionally
// — their retry events live under their own path and never collide
// with a peer-edit delivery.
// Collapsed broadcast failures live under an internal snapshot path;
// an unrelated success on their shared wire path cannot settle them.
let iam_path = "/rustfs/admin/v3/site-replication/peer/iam-item";
upsert_site_replication_retry_event(&mut queue, &peer, iam_path, "peer offline", None);
assert_eq!(dequeue_site_replication_retry_events(&mut queue, &peer, iam_path), 1);
assert_eq!(dequeue_site_replication_retry_events(&mut queue, &peer, iam_path), 0);
assert_eq!(queue[0].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH);
}
/// P1-15 review follow-up: the receiving side of the ordering fence. Two