From 64fd198b239fd7ba0aa311e5a8f241e684584476 Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Fri, 9 Oct 2026 16:51:41 +0100 Subject: [PATCH] fix(rpc): stop a hung endpoint from holding up reads Until 2 endpoints agree on the chain's newest block, a read first asks every endpoint for it. 1 endpoint that never answered held that read for the full request timeout, once a minute. Each ask now gives up after 3 seconds and the check goes on with the answers it has. --- bin/dipper-service/src/chain_client/client.rs | 5 ++- .../src/chain_client/rpc_provider.rs | 43 +++++++++++++++++-- 2 files changed, 43 insertions(+), 5 deletions(-) diff --git a/bin/dipper-service/src/chain_client/client.rs b/bin/dipper-service/src/chain_client/client.rs index 69122dea..9376c63c 100644 --- a/bin/dipper-service/src/chain_client/client.rs +++ b/bin/dipper-service/src/chain_client/client.rs @@ -84,6 +84,9 @@ const FAR_AHEAD_OF_A_SEEN_BLOCK: &str = "too far ahead of the newest block seen" /// block it has seen, while that block is unconfirmed. const CROSS_CHECK_INTERVAL: Duration = Duration::from_secs(60); +/// An answer takes well under a second, so this only bounds how long a hung endpoint holds a read. +const CROSS_CHECK_DEADLINE: Duration = Duration::from_secs(3); + /// The newest block dipper has seen, from reads, receipts and latest-block lookups, and when it /// last moved. It never goes backwards, so a lagging endpoint can't show state from before it. /// Blocks too far ahead of it are refused only once it is confirmed, by the receipt for one of @@ -892,7 +895,7 @@ impl AlloyChainClient { if pool.endpoint_count() < 2 || !self.seen_block().cross_check_due(Instant::now()) { return; } - match agreed_head(pool.latest_blocks().await) { + match agreed_head(pool.latest_blocks(CROSS_CHECK_DEADLINE).await) { Some(head) => self.seen_block().confirm(head, Instant::now()), None => tracing::warn!( "No 2 RPC endpoints agree on the chain's latest block; reads aren't checked \ diff --git a/bin/dipper-service/src/chain_client/rpc_provider.rs b/bin/dipper-service/src/chain_client/rpc_provider.rs index d6e8bf56..f05bd1f7 100644 --- a/bin/dipper-service/src/chain_client/rpc_provider.rs +++ b/bin/dipper-service/src/chain_client/rpc_provider.rs @@ -176,13 +176,18 @@ impl RpcProviderPool { self.providers.len() } - /// Each endpoint's latest block, all asked at once with no retries. Endpoints that fail are - /// left out, so dipper can see whether the ones that answer agree. - pub async fn latest_blocks(&self) -> Vec { + /// Each endpoint's latest block, all asked at once with no retries. Endpoints that fail, or + /// don't answer within `deadline`, are left out, so dipper can see whether the rest agree. + pub async fn latest_blocks(&self, deadline: Duration) -> Vec { let mut asks = tokio::task::JoinSet::new(); for url in &self.providers { let (http, url) = (self.http.clone(), url.clone()); - asks.spawn(async move { (endpoint_name(&url), latest_block(http, &url).await) }); + asks.spawn(async move { + let head = tokio::time::timeout(deadline, latest_block(http, &url)) + .await + .unwrap_or_else(|_| Err(format!("no answer within {deadline:?}"))); + (endpoint_name(&url), head) + }); } let mut heads = Vec::with_capacity(self.providers.len()); while let Some(answer) = asks.join_next().await { @@ -618,6 +623,36 @@ mod tests { ); } + /// Reads wait on the cross-check, so 1 endpoint that never answers must not hold it for the + /// whole request timeout. The healthy endpoint's head still counts. + #[tokio::test] + async fn a_cross_check_stops_waiting_for_a_hung_endpoint() { + let healthy = server_answering_block(0x2a).await; + let hung = MockServer::start().await; + Mock::given(method("POST")) + .respond_with(ResponseTemplate::new(200).set_delay(Duration::from_secs(30))) + .mount(&hung) + .await; + let pool = RpcProviderPool::new( + vec![ + healthy.uri().parse().expect("healthy URL"), + hung.uri().parse().expect("hung URL"), + ], + Duration::from_secs(60), + 0, + ) + .expect("pool"); + + let heads = tokio::time::timeout( + Duration::from_secs(2), + pool.latest_blocks(Duration::from_millis(200)), + ) + .await + .expect("the cross-check should give up on the hung endpoint at its deadline"); + + assert_eq!(heads, vec![0x2a]); + } + /// An endpoint that answers can only describe its own refusal, so it never repeats the /// URL. One that never answers is described by the HTTP client instead, which says which /// URL it was reaching for, and that is where the key sits. Nothing listens on port 1.