diff --git a/bin/dipper-service/src/chain_client/client.rs b/bin/dipper-service/src/chain_client/client.rs index 8aae2df7..9cc58a9d 100644 --- a/bin/dipper-service/src/chain_client/client.rs +++ b/bin/dipper-service/src/chain_client/client.rs @@ -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::{ @@ -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) } @@ -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 { @@ -288,6 +293,7 @@ impl AlloyChainClient { nonce: AtomicU64::new(NONCE_UNINITIALIZED), submit_lock: Mutex::new(()), submit_deadline, + seen_block: AtomicU64::new(0), }), }) } @@ -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( + &self, + to: Address, + call: C, + operation: &'static str, + ) -> Result { + 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( &self, @@ -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 @@ -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) } @@ -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"); } @@ -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) { diff --git a/bin/dipper-service/src/chain_client/rpc_provider.rs b/bin/dipper-service/src/chain_client/rpc_provider.rs index 98954d46..b2f280b7 100644 --- a/bin/dipper-service/src/chain_client/rpc_provider.rs +++ b/bin/dipper-service/src/chain_client/rpc_provider.rs @@ -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", @@ -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< @@ -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 diff --git a/bin/dipper-service/src/network/service/cancel_retry.rs b/bin/dipper-service/src/network/service/cancel_retry.rs index dc4874b2..4ea06248 100644 --- a/bin/dipper-service/src/network/service/cancel_retry.rs +++ b/bin/dipper-service/src/network/service/cancel_retry.rs @@ -32,17 +32,20 @@ const SETTLE_MINUTES: i32 = 2; /// agreement ended, before the retry marks it ended without those details. const LISTENER_GRACE: time::Duration = time::Duration::HOUR; -/// Retry the cancel of agreements still `Cancelling`. `chain_now`, in chain seconds, -/// decides when an offer that was never accepted no longer can be. +/// Retry the cancel of agreements still `Cancelling`. The chain's own latest block time +/// decides when an offer that was never accepted no longer can be, so a subgraph that has +/// fallen behind doesn't hold that up. pub async fn retry_cancelling_agreements( registry: &R, chain_client: &T, config: &IndexingAgreementConfig, - chain_now: u64, ) where R: AgreementRegistry + Sync, T: ChainClient, { + let Some(chain_now) = chain_time(chain_client).await else { + return; + }; let cancelling = match registry .get_cancelling_agreements(BATCH_SIZE, MAX_CANCEL_ATTEMPTS, SETTLE_MINUTES) .await @@ -66,6 +69,19 @@ pub async fn retry_cancelling_agreements( } } +async fn chain_time(chain_client: &T) -> Option { + match chain_client.latest_block_timestamp().await { + Ok(chain_now) => Some(chain_now), + Err(err) => { + tracing::warn!( + error = %err, + "Failed to read the chain's time; cancels are retried next sweep" + ); + None + } + } +} + async fn retry_cancel( registry: &R, chain_client: &T, @@ -322,7 +338,7 @@ fn failed_attempts(err: &ChainClientError) -> u32 { mod tests { use std::sync::{ Mutex, - atomic::{AtomicBool, AtomicU32, Ordering}, + atomic::{AtomicBool, AtomicU32, AtomicU64, Ordering}, }; use async_trait::async_trait; @@ -411,6 +427,8 @@ mod tests { send_fails: bool, mined_cancel_reverts: bool, cancel_has_no_effect: bool, + clock_fails: bool, + now: AtomicU64, ended_by_indexer: bool, who_read_fails: bool, cancels_sent: AtomicU32, @@ -478,7 +496,10 @@ mod tests { Ok(self.ended_by_indexer) } async fn latest_block_timestamp(&self) -> Result { - unimplemented!() + if self.clock_fails { + return Err(ChainClientError::RpcError(anyhow::anyhow!("rpc down"))); + } + Ok(self.now.load(Ordering::SeqCst)) } } @@ -503,7 +524,22 @@ mod tests { async fn retry(registry: &MockRegistry, chain: &MockChain, chain_now: u64) { let config = IndexingAgreementConfig::for_tests(); - retry_cancelling_agreements(registry, chain, &config, chain_now).await; + chain.now.store(chain_now, Ordering::SeqCst); + retry_cancelling_agreements(registry, chain, &config).await; + } + + #[tokio::test] + async fn waits_for_the_next_sweep_when_the_chain_time_cannot_be_read() { + let registry = registry_with_one(true); + let chain = MockChain { + clock_fails: true, + ..live_chain() + }; + + retry(®istry, &chain, 0).await; + + assert_eq!(chain.cancels_sent.load(Ordering::SeqCst), 0); + assert_eq!(registry.checks.load(Ordering::SeqCst), 0); } #[tokio::test] diff --git a/bin/dipper-service/src/network/service/chain_listener.rs b/bin/dipper-service/src/network/service/chain_listener.rs index 145f8cb2..278db264 100644 --- a/bin/dipper-service/src/network/service/chain_listener.rs +++ b/bin/dipper-service/src/network/service/chain_listener.rs @@ -301,13 +301,10 @@ where // finishing a cancel needs only the chain. if last_cancel_retry.is_none_or(|at| at.elapsed() >= CANCEL_RETRY_INTERVAL) { last_cancel_retry = Some(Instant::now()); - let chain_now = - last_persisted_timestamp.unwrap_or_else(dipper_core::time::now_secs); super::cancel_retry::retry_cancelling_agreements( ®istry, &chain_client, &agreement_conf, - chain_now, ) .await; } @@ -901,10 +898,13 @@ where // Both transitions are applied atomically downstream so the // Accept-then-Cancel-in-one-snapshot path can't leak an intermediate // AcceptedOnChain to concurrent readers. + // A withdrawn offer reads as cancelled with no accept time: it never went live. + let withdrawn_offer = snapshot.state.is_canceled() && snapshot.accepted_at == 0; let apply_accept = matches!( agreement.status, IndexingAgreementStatus::Created | IndexingAgreementStatus::Expired, - ) && snapshot.state.reached_accepted(); + ) && snapshot.state.reached_accepted() + && !withdrawn_offer; let already_terminal_cancel = matches!( agreement.status, @@ -2824,6 +2824,47 @@ mod tests { ); } + #[tokio::test] + async fn reconcile_does_not_accept_an_offer_withdrawn_before_anyone_accepted_it() { + // Announcing it would send accepted and terminated events for an agreement that + // was never live; an expired one stays expired. + for status in [ + IndexingAgreementStatus::Created, + IndexingAgreementStatus::Expired, + ] { + let registry = MockRegistry::new(); + let chain_client = MockChainClient::default(); + let worker_queue = MockWorkerQueue::default(); + let agreement_id = IndexingAgreementId::from_bytes(rand::random()); + registry.add_agreement(agreement_id, status); + let mut snapshot = + make_snapshot(agreement_id, AgreementState::CanceledByPayer, Address::ZERO); + snapshot.accepted_at = 0; + + reconcile_agreement( + &snapshot, + ®istry, + &worker_queue, + &chain_client, + test_agreement_conf().as_ref(), + ) + .await + .expect("reconcile ok"); + + assert!(!registry.was_marked_accepted_on_chain(&agreement_id)); + assert_eq!( + registry.was_marked_canceled_by_requester(&agreement_id), + status == IndexingAgreementStatus::Created + ); + assert!( + registry + .audit_writes() + .iter() + .all(|(kind, _)| *kind != "accept") + ); + } + } + #[tokio::test] async fn reconcile_marks_an_agreement_being_cancelled_once_the_chain_shows_it_ended() { for (state, ended_by_dipper) in [