Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion bin/dipper-service/src/chain_client/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 \
Expand Down
43 changes: 39 additions & 4 deletions bin/dipper-service/src/chain_client/rpc_provider.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<u64> {
/// 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<u64> {
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 {
Expand Down Expand Up @@ -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.
Expand Down
Loading