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
198 changes: 186 additions & 12 deletions bin/dipper-service/src/chain_client/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,21 +14,21 @@ use std::{
use async_trait::async_trait;
use dipper_rpc::indexer::indexer_client::sol::RecurringCollectionAgreement;
use thegraph_core::alloy::{
eips::{BlockNumberOrTag, eip2718::Encodable2718},
eips::{BlockId, BlockNumberOrTag, eip2718::Encodable2718},
network::{EthereumWallet, TransactionBuilder},
primitives::{Address, B256, FixedBytes, U256},
providers::Provider,
rpc::types::TransactionRequest,
signers::local::PrivateKeySigner,
sol_types::{SolCall, SolValue},
transports::TransportError,
transports::{TransportError, TransportErrorKind},
};
use tokio::sync::Mutex;

use super::{
abi::{IRecurringAgreementManager, IRecurringCollector},
gas::{GasEstimator, calculate_max_fee, exceeds_max_gas_price, get_gas_prices},
rpc_provider::RpcProviderPool,
rpc_provider::{BEHIND_A_SEEN_BLOCK, RpcProviderPool},
};
use crate::{
chain_client::{
Expand Down Expand Up @@ -58,21 +58,23 @@ const RECEIPT_POLL_INTERVAL: Duration = Duration::from_millis(500);
const VERSION_CURRENT: u64 = 0;

/// `AgreementDetails.state` flags from `IAgreementCollector.sol` (REGISTERED=1,
/// ACCEPTED=2, NOTICE_GIVEN=4, BY_PROVIDER=32). `getAgreementDetails` keeps ACCEPTED set on
/// a canceled agreement and ORs in NOTICE_GIVEN, so a cancel must clear it, not just lack
/// it. REGISTERED without ACCEPTED is an offer still waiting. BY_PROVIDER: the indexer cancelled.
/// ACCEPTED=2, NOTICE_GIVEN=4, SETTLED=8, BY_PROVIDER=32). `getAgreementDetails` keeps
/// ACCEPTED set on a canceled agreement and ORs in NOTICE_GIVEN, so a cancel must clear it,
/// not just lack it. SETTLED: nothing left to claim. BY_PROVIDER: the indexer cancelled.
const STATE_REGISTERED: u16 = 1;
const STATE_ACCEPTED: u16 = 2;
const STATE_NOTICE_GIVEN: u16 = 4;
const STATE_SETTLED: u16 = 8;
const STATE_BY_PROVIDER: u16 = 32;

/// Live iff the terms are accepted and no cancellation notice exists, or an
/// offer is still stored for the indexer to accept. A cancel sets NOTICE_GIVEN
/// while ACCEPTED stays set, so the notice bit tells a live agreement from a
/// cancelled one; a revoked offer reads as an empty state.
/// cancelled one; a revoked offer reads as an empty state. An offer past its
/// deadline stays stored but is SETTLED, since it can no longer be accepted.
fn still_live(state: u16) -> bool {
let accepted = state & STATE_ACCEPTED != 0;
let pending_offer = state & STATE_REGISTERED != 0 && !accepted;
let pending_offer = state & STATE_REGISTERED != 0 && !accepted && state & STATE_SETTLED == 0;
pending_offer || (accepted && state & STATE_NOTICE_GIVEN == 0)
}

Expand Down Expand Up @@ -238,6 +240,9 @@ struct AlloyChainClientInner {
submit_lock: Mutex<()>,
/// How long one submission may hold `submit_lock`; see `derive_submit_deadline`.
submit_deadline: Duration,
/// The newest block dipper has seen, from reads and receipts. An agreement's state is
/// never read from an endpoint behind it, so a lagging endpoint can't undo what dipper saw.
seen_block: AtomicU64,
}

impl AlloyChainClient {
Expand Down Expand Up @@ -288,6 +293,7 @@ impl AlloyChainClient {
nonce: AtomicU64::new(NONCE_UNINITIALIZED),
submit_lock: Mutex::new(()),
submit_deadline,
seen_block: AtomicU64::new(0),
}),
})
}
Expand Down Expand Up @@ -611,11 +617,51 @@ impl AlloyChainClient {
};
let collector = self.inner.recurring_collector_address;
Ok(self
.view(collector, call, "get_agreement_details")
.view_at_seen_block(collector, call, "get_agreement_details")
.await?
.state)
}

/// Run a read-only contract call at the endpoint's latest block, refusing an endpoint
/// whose latest block is older than one dipper has already seen; the pool moves on to
/// the next endpoint instead.
async fn view_at_seen_block<C: SolCall>(
&self,
to: Address,
call: C,
operation: &'static str,
) -> Result<C::Return, ChainClientError> {
let calldata = call.abi_encode();
let seen = self.inner.seen_block.load(Ordering::Relaxed);
let (head, output) = self
.inner
.rpc_pool
.execute(operation, |provider| {
let calldata = calldata.clone();
async move {
let head = provider.get_block_number().await?;
if head < seen {
return Err(TransportErrorKind::custom_str(&format!(
"endpoint is at block {head}, {BEHIND_A_SEEN_BLOCK} ({seen})"
)));
}
let tx = TransactionRequest::default().to(to).input(calldata.into());
let output = provider.call(tx).block(BlockId::number(head)).await?;
Ok((head, output))
}
})
.await?;
self.note_block(head);
C::abi_decode_returns(&output).map_err(|err| {
ChainClientError::RpcError(anyhow::anyhow!("undecodable {operation} from {to}: {err}"))
})
}

/// Remember a block dipper has seen, so later reads are never older.
fn note_block(&self, block: u64) {
self.inner.seen_block.fetch_max(block, Ordering::Relaxed);
}

/// Run a read-only contract call and decode its return value.
async fn view<C: SolCall>(
&self,
Expand Down Expand Up @@ -659,7 +705,12 @@ impl AlloyChainClient {
.await;

match receipt {
Ok(Some(r)) => return Ok(Some(r.status())),
Ok(Some(r)) => {
if let Some(block) = r.block_number {
self.note_block(block);
}
return Ok(Some(r.status()));
}
Ok(None) => {} // not mined yet
Err(e) => {
// Transient RPC error: log and keep polling. If it persists, the outer
Expand Down Expand Up @@ -693,6 +744,7 @@ impl ChainClient for AlloyChainClient {
.ok_or_else(|| {
ChainClientError::RpcError(anyhow::anyhow!("no latest block returned"))
})?;
self.note_block(block.header.number);
Ok(block.header.timestamp)
}

Expand Down Expand Up @@ -1016,15 +1068,22 @@ mod tests {
/// or an id it never saw, as an empty state.
#[test]
fn still_live_covers_accepted_agreements_and_pending_offers() {
const SETTLED: u16 = 8;
const BY_PAYER: u16 = 16;
assert!(still_live(STATE_REGISTERED | STATE_ACCEPTED));
assert!(
still_live(STATE_REGISTERED),
"a pending offer can still be accepted"
);
assert!(
!still_live(STATE_REGISTERED | STATE_SETTLED),
"an offer past its deadline can't be"
);
assert!(
still_live(STATE_REGISTERED | STATE_ACCEPTED | STATE_SETTLED),
"an accepted agreement just collected from is still live"
);
assert!(!still_live(
STATE_REGISTERED | STATE_ACCEPTED | STATE_NOTICE_GIVEN | BY_PAYER | SETTLED
STATE_REGISTERED | STATE_ACCEPTED | STATE_NOTICE_GIVEN | BY_PAYER | STATE_SETTLED
));
assert!(!still_live(0), "revoked or never offered");
}
Expand Down Expand Up @@ -2021,6 +2080,121 @@ mod tests {
}
}

/// Answers as an endpoint whose latest block is `head`, reporting `state` for any
/// agreement read at that block.
struct AgreementStateResponder {
head: AtomicU64,
/// Blocks the endpoint gains each time it reports its head.
catch_up: u64,
state: u16,
}

impl Respond for AgreementStateResponder {
fn respond(&self, request: &Request) -> ResponseTemplate {
let body: serde_json::Value =
serde_json::from_slice(&request.body).expect("JSON-RPC request body");
let result = match body["method"].as_str().unwrap_or_default() {
"eth_blockNumber" => {
let head = self.head.fetch_add(self.catch_up, Ordering::SeqCst);
format!("{head:#x}")
}
"eth_call" => {
let at = body["params"][1].as_str().expect("a block number");
let at = u64::from_str_radix(at.trim_start_matches("0x"), 16).expect("hex");
assert!(
at <= self.head.load(Ordering::SeqCst),
"read at a block the endpoint has"
);
let details = IRecurringCollector::AgreementDetails {
agreementId: FixedBytes::<16>::ZERO,
payer: Address::ZERO,
dataService: Address::ZERO,
serviceProvider: Address::ZERO,
versionHash: B256::ZERO,
state: self.state,
};
let output =
IRecurringCollector::getAgreementDetailsCall::abi_encode_returns(&details);
format!(
"0x{}",
thegraph_core::alloy::primitives::hex::encode(output)
)
}
other => panic!("unexpected call {other}"),
};
ResponseTemplate::new(200).set_body_json(serde_json::json!({
"jsonrpc": "2.0",
"id": body["id"],
"result": result,
}))
}
}

async fn server_at_block(head: u64, state: u16) -> MockServer {
server_catching_up(head, 0, state).await
}

async fn server_catching_up(head: u64, catch_up: u64, state: u16) -> MockServer {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(AgreementStateResponder {
head: AtomicU64::new(head),
catch_up,
state,
})
.mount(&server)
.await;
server
}

#[tokio::test]
async fn waits_for_an_endpoint_a_block_behind_to_catch_up() {
// Hosted endpoints spread calls across nodes, so one a block behind is routine
// rather than a reason to give up on the endpoint.
let endpoint = server_catching_up(94, 1, STATE_REGISTERED | STATE_ACCEPTED).await;
let client = client_over_retrying(vec![endpoint.uri().parse().expect("provider URL")], 1);
client.note_block(95);

let live = client
.agreement_still_active(&[0xab; 16])
.await
.expect("read once the endpoint caught up");

assert!(live);
}

#[tokio::test]
async fn never_reads_an_agreement_from_an_endpoint_behind_a_block_already_seen() {
// The lagging endpoint still shows the offer it hasn't seen accepted and cancelled.
let lagging = server_at_block(90, STATE_REGISTERED).await;
let current =
server_at_block(100, STATE_REGISTERED | STATE_ACCEPTED | STATE_NOTICE_GIVEN).await;
let client = client_over(vec![
lagging.uri().parse().expect("provider URL"),
current.uri().parse().expect("provider URL"),
]);
client.note_block(95);

let live = client
.agreement_still_active(&[0xab; 16])
.await
.expect("read");

assert!(!live, "read from the endpoint that has reached block 95");
assert_eq!(client.inner.seen_block.load(Ordering::Relaxed), 100);
}

#[tokio::test]
async fn fails_a_read_rather_than_answer_from_endpoints_all_behind() {
let lagging = server_at_block(90, STATE_REGISTERED).await;
let client = client_over(vec![lagging.uri().parse().expect("provider URL")]);
client.note_block(95);

let read = client.agreement_still_active(&[0xab; 16]).await;

assert!(read.is_err(), "got {read:?}");
}

async fn client_over_manager(
responder: ManagerViewsResponder,
) -> (AlloyChainClient, MockServer) {
Expand Down
16 changes: 16 additions & 0 deletions bin/dipper-service/src/chain_client/rpc_provider.rs
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,10 @@ fn describe_failure(url: &Url, error: &TransportError) -> String {
/// Error text that indicates a transient failure worth retrying, used only for faults
/// that arrive as prose rather than as a status code or JSON-RPC error object.
const RETRYABLE_ERROR_PATTERNS: &[&str] = &[
BEHIND_A_SEEN_BLOCK,
// A node behind the rest of its provider's fleet, asked for a block it hasn't reached.
"header not found",
"unknown block",
"connection refused",
"connection reset",
"connection closed",
Expand All @@ -68,6 +72,10 @@ const RETRYABLE_ERROR_PATTERNS: &[&str] = &[
"temporary internal error",
];

/// How a read refused by an endpoint behind a block dipper has already seen describes it. It
/// is retried, since an endpoint a block or so behind catches up within a second or two.
pub(super) const BEHIND_A_SEEN_BLOCK: &str = "behind a block already seen";

/// Type alias for the provider with default fillers.
pub type HttpProvider = FillProvider<
JoinFill<
Expand Down Expand Up @@ -633,6 +641,14 @@ mod tests {
);
}

#[test]
fn a_node_that_has_not_reached_a_block_yet_is_retryable() {
let payload = serde_json::from_str(r#"{"code":-32000,"message":"header not found"}"#)
.expect("JSON-RPC error payload");
let err: TransportError = RpcError::ErrorResp(payload);
assert!(RpcProviderPool::is_retryable(&err));
}

#[test]
fn each_retry_waits_twice_as_long_up_to_a_ceiling() {
// 1s, 2s, 4s, 8s, 16s, 32s->30s
Expand Down
Loading
Loading