mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-06 12:09:12 +00:00
fix(site-replication): broadcast bucket ops to every peer and record each failure (backlog#2293)
The generic JSON broadcast (make/delete bucket, bucket-meta hook, bucket ops) returned at the first failing peer, so peers later in deployment-id order never received the request and got no retry event; a transport construction failure recorded nothing at all. Attempt every remote peer like the IAM change hook does: a success settles the peer/path retry event, a failure (transport construction included) enqueues one under the request path, and the first error is returned after all peers were attempted. (cherry picked from commit ce8f73bfd74bac61c383434780f27d17ca16d75e)
This commit is contained in:
@@ -3367,3 +3367,130 @@ fn test_retry_snapshot_tombstones_removed_service_accounts() {
|
||||
assert!(change.create.is_none());
|
||||
assert_eq!(tombstone.updated_at, Some(observed_at));
|
||||
}
|
||||
|
||||
/// Spawns a one-shot HTTP peer that answers 200 and flips the returned flag
|
||||
/// once a request head has arrived.
|
||||
async fn spawn_reached_probe_peer() -> (String, Arc<AtomicBool>, tokio::task::JoinHandle<()>) {
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind healthy peer");
|
||||
let endpoint = format!("http://{}", listener.local_addr().expect("healthy peer address"));
|
||||
let reached = Arc::new(AtomicBool::new(false));
|
||||
let reached_by_server = reached.clone();
|
||||
let server = tokio::spawn(async move {
|
||||
let Ok((mut stream, _)) = listener.accept().await else {
|
||||
return;
|
||||
};
|
||||
let mut request = Vec::new();
|
||||
let mut buffer = [0_u8; 1024];
|
||||
loop {
|
||||
let Ok(read) = stream.read(&mut buffer).await else {
|
||||
return;
|
||||
};
|
||||
if read == 0 {
|
||||
return;
|
||||
}
|
||||
request.extend_from_slice(&buffer[..read]);
|
||||
if request.windows(4).any(|window| window == b"\r\n\r\n") {
|
||||
break;
|
||||
}
|
||||
}
|
||||
reached_by_server.store(true, Ordering::SeqCst);
|
||||
let _ = stream
|
||||
.write_all(b"HTTP/1.1 200 OK\r\ncontent-length: 2\r\nconnection: close\r\n\r\nok")
|
||||
.await;
|
||||
});
|
||||
(endpoint, reached, server)
|
||||
}
|
||||
|
||||
/// Three-peer runtime whose local peer is `local`; BTreeMap order visits the
|
||||
/// failing peer `b` before the healthy peer `c`.
|
||||
fn broadcast_runtime_with_failing_peer_before_healthy(failing_endpoint: &str, healthy_endpoint: &str) -> SiteReplicationRuntime {
|
||||
let local_peer = PeerInfo {
|
||||
deployment_id: "local".to_string(),
|
||||
..peer("local", "http://127.0.0.1:9")
|
||||
};
|
||||
let mut state = SiteReplicationState {
|
||||
name: "local".to_string(),
|
||||
service_account_access_key: "site-replicator-0".to_string(),
|
||||
..Default::default()
|
||||
};
|
||||
state.peers.insert("local".to_string(), local_peer.clone());
|
||||
state.peers.insert(
|
||||
"b".to_string(),
|
||||
PeerInfo {
|
||||
deployment_id: "b".to_string(),
|
||||
..peer("b", failing_endpoint)
|
||||
},
|
||||
);
|
||||
state.peers.insert(
|
||||
"c".to_string(),
|
||||
PeerInfo {
|
||||
deployment_id: "c".to_string(),
|
||||
..peer("c", healthy_endpoint)
|
||||
},
|
||||
);
|
||||
SiteReplicationRuntime {
|
||||
state,
|
||||
local_peer,
|
||||
service_account_secret_key: "site-replicator-secret".to_string(),
|
||||
}
|
||||
}
|
||||
|
||||
const BROADCAST_PROBE_DELETE_BUCKET_PATH: &str =
|
||||
"/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=delete-bucket";
|
||||
|
||||
/// The generic JSON broadcast (bucket make/delete, bucket-meta hook, bucket
|
||||
/// ops) attempts every remote peer: a peer whose request fails must not stop
|
||||
/// delivery to the peers that follow it in deployment-id order, and the
|
||||
/// failure is still reported to the caller (backlog#2293).
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_broadcast_json_reaches_healthy_peers_after_a_failed_peer() {
|
||||
// Peer "b": nothing listens on the port, so the connect is refused.
|
||||
let refused = TcpListener::bind("127.0.0.1:0").await.expect("bind refused-peer probe");
|
||||
let refused_endpoint = format!("http://{}", refused.local_addr().expect("refused-peer address"));
|
||||
drop(refused);
|
||||
|
||||
let (healthy_endpoint, reached, server) = spawn_reached_probe_peer().await;
|
||||
let runtime = broadcast_runtime_with_failing_peer_before_healthy(&refused_endpoint, &healthy_endpoint);
|
||||
|
||||
let result = temp_env::async_with_vars([(ALLOW_LOOPBACK_REPLICATION_TARGET_ENV, Some("true"))], async {
|
||||
broadcast_site_replication_json_with_runtime(&runtime, BROADCAST_PROBE_DELETE_BUCKET_PATH, &serde_json::json!({})).await
|
||||
})
|
||||
.await;
|
||||
|
||||
let err = result.expect_err("peer b refuses connections, the broadcast must report it");
|
||||
assert!(
|
||||
reached.load(Ordering::SeqCst),
|
||||
"peer c never received the broadcast once peer b failed: {err}"
|
||||
);
|
||||
server.abort();
|
||||
}
|
||||
|
||||
/// Same guarantee when the failing peer never gets a transport: an endpoint
|
||||
/// that `PeerTransport::for_runtime_peer` rejects must be skipped past (and
|
||||
/// reported), not abort the broadcast before the healthy peers (backlog#2293).
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn test_broadcast_json_reaches_healthy_peers_after_a_peer_without_transport() {
|
||||
// Peer "b": a scheme the peer connection validator refuses outright.
|
||||
let forbidden_endpoint = "ftp://peer-b.example.com";
|
||||
|
||||
let (healthy_endpoint, reached, server) = spawn_reached_probe_peer().await;
|
||||
let runtime = broadcast_runtime_with_failing_peer_before_healthy(forbidden_endpoint, &healthy_endpoint);
|
||||
|
||||
let result = temp_env::async_with_vars([(ALLOW_LOOPBACK_REPLICATION_TARGET_ENV, Some("true"))], async {
|
||||
broadcast_site_replication_json_with_runtime(&runtime, BROADCAST_PROBE_DELETE_BUCKET_PATH, &serde_json::json!({})).await
|
||||
})
|
||||
.await;
|
||||
|
||||
let err = result.expect_err("peer b has no usable transport, the broadcast must report it");
|
||||
assert!(
|
||||
err.to_string().contains("invalid persisted site replication peer"),
|
||||
"the reported error must be peer b's transport failure: {err}"
|
||||
);
|
||||
assert!(
|
||||
reached.load(Ordering::SeqCst),
|
||||
"peer c never received the broadcast once peer b failed to get a transport: {err}"
|
||||
);
|
||||
server.abort();
|
||||
}
|
||||
|
||||
@@ -876,6 +876,14 @@ pub(crate) async fn broadcast_site_replication_json<T: Serialize>(path: &str, bo
|
||||
broadcast_site_replication_json_with_runtime(&runtime, path, body).await
|
||||
}
|
||||
|
||||
/// PUT `body` to `path` on every remote peer of the runtime.
|
||||
///
|
||||
/// Every peer is attempted: one peer's failure — transport construction
|
||||
/// included — must not skip the peers that follow it in deployment-id order,
|
||||
/// or they silently miss the change with no retry record (backlog#2293). A
|
||||
/// success settles the peer/path's queued retry event, a failure enqueues one
|
||||
/// under the request `path` (so the drain classifies it as today), and the
|
||||
/// first error is returned once all peers were attempted.
|
||||
pub(crate) async fn broadcast_site_replication_json_with_runtime<T: Serialize>(
|
||||
runtime: &SiteReplicationRuntime,
|
||||
path: &str,
|
||||
@@ -883,20 +891,30 @@ pub(crate) async fn broadcast_site_replication_json_with_runtime<T: Serialize>(
|
||||
) -> S3Result<()> {
|
||||
let state = &runtime.state;
|
||||
let local_peer = &runtime.local_peer;
|
||||
let mut first_error: Option<S3Error> = None;
|
||||
|
||||
for peer in state.peers.values() {
|
||||
if peer.deployment_id == local_peer.deployment_id || same_identity_endpoint(&peer.endpoint, &local_peer.endpoint) {
|
||||
continue;
|
||||
}
|
||||
|
||||
let transport = PeerTransport::for_runtime_peer(peer).await?;
|
||||
PeerAdminRequest::put(&transport.connection, path, &state.service_account_access_key)
|
||||
.with_client(&transport.client)
|
||||
.send_with_retry_event(peer, &runtime.service_account_secret_key, body)
|
||||
.await?;
|
||||
let sent = match PeerTransport::for_runtime_peer(peer).await {
|
||||
Ok(transport) => PeerAdminRequest::put(&transport.connection, path, &state.service_account_access_key)
|
||||
.with_client(&transport.client)
|
||||
.send_with_retry_event(peer, &runtime.service_account_secret_key, body)
|
||||
.await
|
||||
.map(|_| ()),
|
||||
Err(err) => {
|
||||
enqueue_site_replication_retry_event(peer, path, &err).await;
|
||||
Err(err)
|
||||
}
|
||||
};
|
||||
if let Err(err) = sent {
|
||||
first_error.get_or_insert(err);
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
first_error.map_or(Ok(()), Err)
|
||||
}
|
||||
|
||||
pub(crate) fn parse_endpoint_refresh_status(peer: &PeerInfo, body: &[u8]) -> S3Result<()> {
|
||||
|
||||
Reference in New Issue
Block a user