fix(site-replication): bound retry coordination

This commit is contained in:
overtrue
2026-09-05 11:25:52 +08:00
parent 41ceb044b0
commit bde7c58892
5 changed files with 213 additions and 60 deletions
@@ -1824,7 +1824,7 @@ async fn site_replication_reconcile_prerequisites_ready() -> bool {
fn reconcile_site_replication_retry_drain() -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send>> {
Box::pin(async {
let Some(_lifecycle) = SiteReplicationLifecycleGuard::try_acquire() else {
let Some(lifecycle) = SiteReplicationLifecycleGuard::try_acquire() else {
return;
};
if !site_replication_reconcile_prerequisites_ready().await {
@@ -1839,6 +1839,10 @@ fn reconcile_site_replication_retry_drain() -> std::pin::Pin<Box<dyn std::future
}
Err(_) => return,
}
// Retry sends re-check membership under the bucket-op read lock for
// each bounded peer request. Do not hold the lifecycle guard across
// an arbitrarily large snapshot replay.
drop(lifecycle);
drain_site_replication_retry_queue().await;
})
}
@@ -9573,10 +9577,6 @@ mod tests {
SITE_REPLICATION_LIFECYCLE_LOCK_TIMEOUT >= SITE_REPLICATION_PEER_REQUEST_TIMEOUT,
"a waiter must not give up before the holder's single wedged peer probe can finish"
);
assert!(
crate::site_replication::SITE_REPLICATION_RETRY_DRAIN_BATCH_TIMEOUT < SITE_REPLICATION_LIFECYCLE_LOCK_TIMEOUT,
"the retry drainer must release the lifecycle guard before operator waiters time out"
);
}
#[tokio::test(start_paused = true)]
+3 -1
View File
@@ -75,7 +75,7 @@ use crate::auth::get_condition_values_with_client_info;
use crate::error::ApiError;
use crate::shared_types::RemoteAddr;
use crate::site_replication::{
cancel_site_replication_delete_bucket_hook, finish_site_replication_delete_bucket_hook,
SITE_REPLICATION_BUCKET_OP_LOCK, cancel_site_replication_delete_bucket_hook, finish_site_replication_delete_bucket_hook,
prepare_site_replication_delete_bucket_hook, site_replication_bucket_meta_hook, site_replication_make_bucket_hook,
};
use crate::storage::storage_api::lock_bucket_targets_metadata;
@@ -1398,6 +1398,7 @@ impl DefaultBucketUsecase {
authorize_request(&mut req, Action::S3Action(S3Action::ForceDeleteBucketAction)).await?;
}
let bucket_op_guard = SITE_REPLICATION_BUCKET_OP_LOCK.read().await;
let replication_delete_intent = prepare_site_replication_delete_bucket_hook(&input.bucket, force).await?;
let delete_result = store
.delete_bucket(
@@ -1415,6 +1416,7 @@ impl DefaultBucketUsecase {
}
return Err(err.into());
}
drop(bucket_op_guard);
// Drop every cached object body for the now-deleted bucket so dead
// bytes do not sit resident until TTL. Covers both the normal and the
+2 -14
View File
@@ -409,7 +409,7 @@ async fn broadcast_site_replication_destructive_bucket_op(
.await;
match result {
Ok(_) => {
dequeue_site_replication_retry_event(peer, path).await;
dequeue_site_replication_destructive_retry_events(peer, path).await;
None
}
Err(err) => {
@@ -456,8 +456,6 @@ pub(crate) async fn prepare_site_replication_delete_bucket_hook(
let Some(runtime) = runtime_site_replication_targets().await? else {
return Ok(None);
};
let store =
current_object_store_handle().ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))?;
let peers = runtime
.state
.peers
@@ -468,17 +466,7 @@ pub(crate) async fn prepare_site_replication_delete_bucket_hook(
})
.cloned()
.collect::<Vec<_>>();
let reserve_peers = peers.clone();
let reserve_path = path.clone();
let result =
with_config_object_write_lock(store, SITE_REPLICATION_REPAIR_EXECUTION_LOCK_PATH.to_string(), move || async move {
prequeue_site_replication_destructive_events(&reserve_peers, &reserve_path).await
})
.await;
let reserved_peers = match result {
Ok(result) => result?,
Err(err) => return Err(ApiError::from(err).into()),
};
let reserved_peers = prequeue_site_replication_destructive_events(&peers, &path).await?;
Ok(Some(SiteReplicationDeleteBucketIntent {
runtime,
peers,
+155 -40
View File
@@ -16,11 +16,6 @@ use super::*;
pub(crate) const SITE_REPLICATION_RETRY_QUEUE_LIMIT: usize = 256;
/// Keep the reconcile tick's lifecycle guard below the 30-second operator
/// wait bound. Successfully settled events are persisted as the batch runs,
/// so a later tick continues with the remaining work.
pub(crate) const SITE_REPLICATION_RETRY_DRAIN_BATCH_TIMEOUT: Duration = Duration::from_secs(25);
/// Attempts before an entry reports as `failed` in retryStats. Visibility
/// only: a `failed` entry stays drain-eligible, and the reachability probe
/// short-circuits its backoff once the peer answers again — so an early
@@ -66,6 +61,13 @@ pub(crate) fn retry_event_matches(event: &SiteReplicationRetryEvent, peer: &Peer
(event.peer_deployment_id == peer.deployment_id || event.peer_endpoint == peer.endpoint) && event.path == path
}
pub(crate) fn retry_event_is_destructive_bucket_op(event: &SiteReplicationRetryEvent) -> bool {
matches!(
retry_bucket_operation(&event.path).as_deref(),
Some("delete-bucket" | "force-delete-bucket" | "purge-deleted-bucket")
)
}
pub(crate) const SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH: &str = "internal:retry-snapshot:iam";
pub(crate) const SITE_REPLICATION_RETRY_BUCKET_METADATA_SNAPSHOT_PATH: &str = "internal:retry-snapshot:bucket-metadata";
@@ -244,8 +246,12 @@ pub(crate) fn upsert_site_replication_retry_event(
deletions_recorded: false,
});
if queue.len() > SITE_REPLICATION_RETRY_QUEUE_LIMIT {
let overflow = queue.len() - SITE_REPLICATION_RETRY_QUEUE_LIMIT;
queue.drain(0..overflow);
while queue.len() > SITE_REPLICATION_RETRY_QUEUE_LIMIT {
let Some(index) = queue.iter().position(|event| !retry_event_is_destructive_bucket_op(event)) else {
break;
};
queue.remove(index);
}
}
}
@@ -307,6 +313,52 @@ pub(crate) fn prequeue_site_replication_destructive_events_in_state(
Ok(missing)
}
pub(crate) fn settle_site_replication_destructive_retry_events(
queue: &mut Vec<SiteReplicationRetryEvent>,
peer: &PeerInfo,
path: &str,
) -> usize {
let Some(bucket) = retry_bucket_name(path) else {
return 0;
};
let before = queue.len();
queue.retain(|event| {
!(retry_event_matches_peer(event, peer)
&& retry_event_is_destructive_bucket_op(event)
&& retry_bucket_name(&event.path).as_deref() == Some(bucket.as_str()))
});
before.saturating_sub(queue.len())
}
fn retry_event_matches_peer(event: &SiteReplicationRetryEvent, peer: &PeerInfo) -> bool {
event.peer_deployment_id == peer.deployment_id || event.peer_endpoint == peer.endpoint
}
pub(crate) async fn dequeue_site_replication_destructive_retry_events(peer: &PeerInfo, path: &str) {
let peer_owned = peer.clone();
let path_owned = path.to_string();
let result = update_site_replication_state_when_changed(move |state| {
let removed = settle_site_replication_destructive_retry_events(&mut state.retry_queue, &peer_owned, &path_owned);
Ok(if removed == 0 {
StateCommit::Unchanged(())
} else {
StateCommit::Changed(())
})
})
.await;
if let Err(err) = result {
warn!(
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
peer = %peer.endpoint,
path,
error = ?err,
"failed to settle destructive site replication retry events"
);
}
}
pub(crate) async fn prequeue_site_replication_destructive_events(peers: &[PeerInfo], path: &str) -> S3Result<Vec<PeerInfo>> {
if peers.is_empty() {
return Ok(Vec::new());
@@ -783,27 +835,69 @@ impl RetrySnapshot {
}
}
pub(crate) async fn send(&self, transport: &PeerTransport, access_key: &str, secret_key: &str) -> S3Result<()> {
pub(crate) async fn send(
&self,
peer: &PeerInfo,
transport: &PeerTransport,
access_key: &str,
secret_key: &str,
) -> S3Result<bool> {
match self {
Self::Iam(items) => {
for item in items {
SiteReplicationRepairTask::Iam(item)
.send(transport, access_key, secret_key)
.await?;
if !send_retry_task_if_peer_current(
peer,
&SiteReplicationRepairTask::Iam(item),
transport,
access_key,
secret_key,
)
.await?
{
return Ok(false);
}
}
}
Self::BucketMetadata(items) => {
for item in items {
SiteReplicationRepairTask::BucketMetadata(item)
.send(transport, access_key, secret_key)
.await?;
if !send_retry_task_if_peer_current(
peer,
&SiteReplicationRepairTask::BucketMetadata(item),
transport,
access_key,
secret_key,
)
.await?
{
return Ok(false);
}
}
}
}
Ok(())
Ok(true)
}
}
pub(crate) async fn send_retry_task_if_peer_current(
peer: &PeerInfo,
task: &SiteReplicationRepairTask<'_>,
transport: &PeerTransport,
access_key: &str,
secret_key: &str,
) -> S3Result<bool> {
let _bucket_op = SITE_REPLICATION_BUCKET_OP_LOCK.read().await;
let state = load_site_replication_state().await?;
let current = state
.peers
.get(&peer.deployment_id)
.is_some_and(|current| same_identity_endpoint(&current.endpoint, &peer.endpoint));
if !current {
return Ok(false);
}
task.send(transport, access_key, secret_key).await?;
Ok(true)
}
#[derive(Hash, PartialEq, Eq)]
pub(crate) enum IamSnapshotKey {
Policy(String),
@@ -1316,13 +1410,14 @@ pub(crate) async fn drain_site_replication_retry_queue_inner() -> S3Result<()> {
// guard) may have started since. Re-check on the fresh state.
return Ok(());
}
// Serialize against operator repair execution. This does NOT close the
// Serialize against operator repair execution. Peer membership is
// re-checked under the bucket-op read lock immediately before each
// network request, so the caller need not hold the lifecycle guard while
// a large snapshot is replayed. This does NOT close the
// dry-run -> execute window (dry-run takes no lock): a drain settling a
// replayable bucket-op entry in that window changes the preflight token
// and execute fails safe with "preflight is stale" — the operator
// re-runs the dry-run. Lock order matches repair: lifecycle guard (held
// by the reconcile tick) -> repair execution lock -> state object lock
// inside the send bookkeeping. The lock also elects one server to probe
// re-runs the dry-run. The lock also elects one server to probe
// and replay the queue; after acquiring it, reload state so a settled
// event or deleted bucket cannot be replayed from this admission snapshot.
with_config_object_write_lock(store, SITE_REPLICATION_REPAIR_EXECUTION_LOCK_PATH.to_string(), move || async move {
@@ -1341,17 +1436,11 @@ pub(crate) async fn drain_site_replication_retry_queue_inner() -> S3Result<()> {
let now = OffsetDateTime::now_utc();
let actionable = actionable_site_replication_retry_events(&runtime.state, now);
let deferred = deferred_site_replication_retry_events(&runtime.state, now);
let drain = async {
promote_reachable_deferred_retry_events(&runtime, &actionable, deferred).await?;
if actionable.is_empty() {
return Ok(());
}
drain_site_replication_retry_queue_locked(runtime, actionable).await
};
match tokio::time::timeout(SITE_REPLICATION_RETRY_DRAIN_BATCH_TIMEOUT, drain).await {
Ok(result) => result,
Err(_) => Ok(()),
promote_reachable_deferred_retry_events(&runtime, &actionable, deferred).await?;
if actionable.is_empty() {
return Ok(());
}
drain_site_replication_retry_queue_locked(runtime, actionable).await
})
.await
.map_err(ApiError::from)?
@@ -1472,12 +1561,21 @@ pub(crate) async fn drain_one_site_replication_retry_event(
drop_corrupt_iam_deletion_replay(peer, &record.id).await;
continue;
};
if let Err(err) = SiteReplicationRepairTask::Iam(&item)
.send(transport, access_key, secret_key)
.await
match send_retry_task_if_peer_current(
peer,
&SiteReplicationRepairTask::Iam(&item),
transport,
access_key,
secret_key,
)
.await
{
enqueue_site_replication_retry_event(peer, &event.path, &err).await;
return Err(err);
Ok(true) => {}
Ok(false) => return Ok(false),
Err(err) => {
enqueue_site_replication_retry_event(peer, &event.path, &err).await;
return Err(err);
}
}
replayed_record_ids.push(record.id.clone());
}
@@ -1486,9 +1584,13 @@ pub(crate) async fn drain_one_site_replication_retry_event(
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);
match replay.send(peer, transport, access_key, secret_key).await {
Ok(true) => {}
Ok(false) => return Ok(false),
Err(err) => {
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)?;
@@ -1533,9 +1635,13 @@ pub(crate) async fn drain_one_site_replication_retry_event(
return Ok(true);
}
for task in &tasks {
if let Err(err) = task.send(transport, access_key, secret_key).await {
enqueue_site_replication_retry_event(peer, &event.path, &err).await;
return Err(err);
match send_retry_task_if_peer_current(peer, task, transport, access_key, secret_key).await {
Ok(true) => {}
Ok(false) => return Ok(false),
Err(err) => {
enqueue_site_replication_retry_event(peer, &event.path, &err).await;
return Err(err);
}
}
}
dequeue_site_replication_retry_event(peer, &event.path).await;
@@ -1564,6 +1670,15 @@ pub(crate) async fn drain_one_site_replication_retry_event(
let edit_path = peer_edit_path_with_fence(local_deployment_id, generation);
let delivery_fence = local_deployment_id.is_some().then_some(generation);
for body in &bodies {
let _bucket_op = SITE_REPLICATION_BUCKET_OP_LOCK.read().await;
let current_state = load_site_replication_state().await?;
let current = current_state
.peers
.get(&peer.deployment_id)
.is_some_and(|current| same_identity_endpoint(&current.endpoint, &peer.endpoint));
if !current {
return Ok(false);
}
if let Err(err) = PeerAdminRequest::put(&transport.connection, &edit_path, access_key)
.with_client(&transport.client)
.send(secret_key, body)
+48
View File
@@ -862,6 +862,54 @@ fn test_destructive_bucket_op_prequeue_does_not_evict_existing_liability() {
assert_eq!(state.retry_queue.iter().map(|event| event.id.clone()).collect::<Vec<_>>(), existing_ids);
}
#[test]
fn test_destructive_success_settles_equivalent_bucket_intents() {
let peer = peer("remote", "https://remote.example.com");
let failed = drain_event(
"remote",
"/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=delete-bucket&retryIntent=a",
1,
None,
);
let current = drain_event(
"remote",
"/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=force-delete-bucket&retryIntent=b",
1,
None,
);
let unrelated = drain_event(
"remote",
"/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=videos&operation=delete-bucket&retryIntent=c",
1,
None,
);
let mut queue = vec![failed, current.clone(), unrelated.clone()];
assert_eq!(settle_site_replication_destructive_retry_events(&mut queue, &peer, &current.path), 2);
assert_eq!(queue.len(), 1);
assert_eq!(queue[0].id, unrelated.id);
}
#[test]
fn test_ordinary_retry_does_not_evict_destructive_liability() {
let mut queue = (0..SITE_REPLICATION_RETRY_QUEUE_LIMIT)
.map(|index| {
drain_event(
&format!("peer-{index}"),
&format!("/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=bucket-{index}&operation=delete-bucket"),
1,
None,
)
})
.collect::<Vec<_>>();
let ordinary = peer("ordinary", "https://ordinary.example.com");
upsert_site_replication_retry_event(&mut queue, &ordinary, SITE_REPLICATION_PEER_EDIT_PATH, "failed (connect)", None);
assert_eq!(queue.len(), SITE_REPLICATION_RETRY_QUEUE_LIMIT);
assert!(queue.iter().all(retry_event_is_destructive_bucket_op));
}
#[test]
fn test_retry_snapshot_fingerprint_detects_concurrent_iam_change() {
let old = SRIAMItem {