From 3630933e93e0c9a6939db4ebb47304d5ff8acdf6 Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Fri, 2 Oct 2026 22:11:52 +0300 Subject: [PATCH 1/5] fix(selection): keep indexers dipper is cancelling out of selection An agreement dipper is still cancelling may be live and paying its indexer, yet IISA wasn't told to skip that indexer, so it could be picked again and paid twice. It now goes on the deployment's declined list, which blocks a pick without counting it as part of the group. --- .../src/worker/handlers/selection_context.rs | 79 +++++++++++++++++-- 1 file changed, 73 insertions(+), 6 deletions(-) diff --git a/bin/dipper-service/src/worker/handlers/selection_context.rs b/bin/dipper-service/src/worker/handlers/selection_context.rs index 0f65da62..3485463d 100644 --- a/bin/dipper-service/src/worker/handlers/selection_context.rs +++ b/bin/dipper-service/src/worker/handlers/selection_context.rs @@ -7,7 +7,9 @@ use thegraph_core::{DeploymentId, IndexerId, alloy::primitives::ChainId}; use crate::{ network::service::entity_count_cache::EntityCountCache, - registry::{AgreementRegistry, IndexerDenylistRegistry, IndexingAgreementStatus}, + registry::{ + AgreementRegistry, IndexerDenylistRegistry, IndexingAgreement, IndexingAgreementStatus, + }, worker::result::{JobError, JobResult}, }; @@ -42,11 +44,12 @@ where R: AgreementRegistry + IndexerDenylistRegistry, { // Get indexers that already have active agreements for this deployment - let existing_indexers = registry + let agreements = registry .get_indexing_agreements_by_deployment_id(deployment_id) .await - .map_err(|err| JobError::Fatal(err.into()))? - .into_iter() + .map_err(|err| JobError::Fatal(err.into()))?; + let existing_indexers = agreements + .iter() .filter(|a| is_active_agreement(&a.status)) .map(|a| a.indexer.id) .collect::>(); @@ -58,7 +61,7 @@ where .map_err(|err| JobError::Fatal(err.into()))?; // Get indexers that declined within their respective lookback periods - let declined_indexers = registry + let mut declined_indexers = registry .get_declined_indexers_by_deployment( declined_indexer_lookback_days, price_rejection_lookback_days, @@ -67,6 +70,7 @@ where ) .await .map_err(|err| JobError::Fatal(err.into()))?; + exclude_cancelling_indexers(&mut declined_indexers, *deployment_id, &agreements); // Get denied indexers that should be excluded from selection let indexer_denylist = registry @@ -173,6 +177,30 @@ fn wei_per_second_to_grt_per_28d(wei_per_second: f64) -> f64 { wei_per_second * SECONDS_PER_28_DAYS / WEI_PER_GRT } +/// Add to the deployment's declined list the indexers whose agreement dipper is still +/// cancelling. That agreement may still be live and paid on-chain, so its indexer must not +/// be picked again, but it no longer counts towards the group IISA sizes. +fn exclude_cancelling_indexers( + declined: &mut HashMap>, + deployment_id: DeploymentId, + agreements: &[IndexingAgreement], +) { + let cancelling = agreements + .iter() + .filter(|a| a.status == IndexingAgreementStatus::Cancelling) + .map(|a| a.indexer.id) + .collect::>(); + if cancelling.is_empty() { + return; + } + let excluded = declined.entry(deployment_id).or_default(); + for indexer in cancelling { + if !excluded.contains(&indexer) { + excluded.push(indexer); + } + } +} + /// Check if an agreement status represents an active agreement. fn is_active_agreement(status: &IndexingAgreementStatus) -> bool { matches!( @@ -184,7 +212,46 @@ fn is_active_agreement(status: &IndexingAgreementStatus) -> bool { #[cfg(test)] mod tests { use super::*; - use crate::registry::AgreementFeeRate; + use crate::{cancel_dispatch::tests::agreement, registry::AgreementFeeRate}; + + fn indexer(hex_digit: char) -> IndexerId { + format!("0x{}", hex_digit.to_string().repeat(40)) + .parse() + .unwrap() + } + + #[test] + fn an_indexer_still_being_cancelled_cannot_be_picked_again_for_the_deployment() { + let deployment: DeploymentId = "QmTXzATwNfgGVukV1fX2T6xw9f6LAYRVWpsdXyRWzUR2H9" + .parse() + .unwrap(); + let mut cancelling = agreement(IndexingAgreementStatus::Cancelling, None); + cancelling.indexer.id = indexer('b'); + let mut accepted = agreement(IndexingAgreementStatus::AcceptedOnChain, None); + accepted.indexer.id = indexer('c'); + let mut declined = HashMap::from([(deployment, vec![indexer('a')])]); + + exclude_cancelling_indexers(&mut declined, deployment, &[cancelling, accepted]); + + assert_eq!(declined[&deployment], vec![indexer('a'), indexer('b')]); + } + + #[test] + fn a_cancelling_indexer_already_declined_is_listed_once_and_none_adds_no_entry() { + let deployment: DeploymentId = "QmTXzATwNfgGVukV1fX2T6xw9f6LAYRVWpsdXyRWzUR2H9" + .parse() + .unwrap(); + let mut cancelling = agreement(IndexingAgreementStatus::Cancelling, None); + cancelling.indexer.id = indexer('a'); + let mut declined = HashMap::from([(deployment, vec![indexer('a')])]); + exclude_cancelling_indexers(&mut declined, deployment, &[cancelling]); + assert_eq!(declined[&deployment], vec![indexer('a')]); + + let mut none_declined = HashMap::new(); + let accepted = agreement(IndexingAgreementStatus::AcceptedOnChain, None); + exclude_cancelling_indexers(&mut none_declined, deployment, &[accepted]); + assert!(none_declined.is_empty()); + } #[test] fn test_wei_per_second_to_grt_per_28d() { From 5eb995ba7dc442be971e95c44bcc8f91da3d5845 Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Fri, 2 Oct 2026 22:11:57 +0300 Subject: [PATCH 2/5] fix(cancel): never record an indexer's cancel as dipper's The cancel retry marked an accepted agreement it found ended as cancelled by dipper after an hour, even when the indexer had ended it, so the wrong end was announced. It now reads who ended it from the contract, and leaves an end by the indexer for the chain listener to record as theirs. --- bin/dipper-service/src/cancel_dispatch.rs | 6 ++ bin/dipper-service/src/chain_client.rs | 14 +++ bin/dipper-service/src/chain_client/client.rs | 55 +++++------ .../src/network/service/cancel_retry.rs | 95 +++++++++++++++++-- .../src/network/service/chain_listener.rs | 6 ++ .../src/network/service/escrow_reconciler.rs | 6 ++ .../src/network/service/liveness_checker.rs | 6 ++ .../cancel_rejected_agreement_on_chain.rs | 6 ++ .../handlers/reassess_indexing_request.rs | 12 +++ .../src/worker/handlers/submit_offer.rs | 6 ++ 10 files changed, 173 insertions(+), 39 deletions(-) diff --git a/bin/dipper-service/src/cancel_dispatch.rs b/bin/dipper-service/src/cancel_dispatch.rs index e4c9c80f..4f6734a2 100644 --- a/bin/dipper-service/src/cancel_dispatch.rs +++ b/bin/dipper-service/src/cancel_dispatch.rs @@ -280,6 +280,12 @@ pub(crate) mod tests { *self.active_reads.lock().unwrap() += 1; Ok(self.still_active_after_cancel) } + async fn agreement_ended_by_indexer( + &self, + _agreement_id: &[u8; 16], + ) -> Result { + Ok(false) + } } fn manager_conf(collector: Address) -> IndexingAgreementConfig { diff --git a/bin/dipper-service/src/chain_client.rs b/bin/dipper-service/src/chain_client.rs index 5c734b26..f9f1583f 100644 --- a/bin/dipper-service/src/chain_client.rs +++ b/bin/dipper-service/src/chain_client.rs @@ -128,6 +128,13 @@ pub trait ChainClient { agreement_id: &[u8; 16], ) -> Result; + /// Read whether the indexer ended the agreement on-chain, which `getAgreementDetails` + /// reports with its BY_PROVIDER flag. False for one still live or ended by dipper. + async fn agreement_ended_by_indexer( + &self, + agreement_id: &[u8; 16], + ) -> Result; + /// Read the latest block's unix timestamp from the chain. Lets agreement /// deadlines be stamped from live chain time when the chain-clock bypass is /// on, instead of a cached listener timestamp that can lag a fast chain. @@ -182,6 +189,13 @@ impl ChainClient for Arc { (**self).agreement_still_active(agreement_id).await } + async fn agreement_ended_by_indexer( + &self, + agreement_id: &[u8; 16], + ) -> Result { + (**self).agreement_ended_by_indexer(agreement_id).await + } + async fn latest_block_timestamp(&self) -> Result { (**self).latest_block_timestamp().await } diff --git a/bin/dipper-service/src/chain_client/client.rs b/bin/dipper-service/src/chain_client/client.rs index e759edb4..8aae2df7 100644 --- a/bin/dipper-service/src/chain_client/client.rs +++ b/bin/dipper-service/src/chain_client/client.rs @@ -58,12 +58,13 @@ 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). `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. +/// 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. const STATE_REGISTERED: u16 = 1; const STATE_ACCEPTED: u16 = 2; const STATE_NOTICE_GIVEN: u16 = 4; +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 @@ -602,6 +603,19 @@ impl AlloyChainClient { } } + /// The agreement's `AgreementDetails.state` flags for its current terms. + async fn agreement_state(&self, agreement_id: &[u8; 16]) -> Result { + let call = IRecurringCollector::getAgreementDetailsCall { + agreementId: FixedBytes::<16>::from_slice(agreement_id), + index: thegraph_core::alloy::primitives::U256::from(VERSION_CURRENT), + }; + let collector = self.inner.recurring_collector_address; + Ok(self + .view(collector, call, "get_agreement_details") + .await? + .state) + } + /// Run a read-only contract call and decode its return value. async fn view( &self, @@ -789,35 +803,14 @@ impl ChainClient for AlloyChainClient { &self, agreement_id: &[u8; 16], ) -> Result { - let calldata = IRecurringCollector::getAgreementDetailsCall { - agreementId: FixedBytes::<16>::from_slice(agreement_id), - index: thegraph_core::alloy::primitives::U256::from(VERSION_CURRENT), - } - .abi_encode(); - - let collector = self.inner.recurring_collector_address; - let output = self - .inner - .rpc_pool - .execute("get_agreement_details", |provider| { - let calldata = calldata.clone(); - async move { - let tx = TransactionRequest::default() - .to(collector) - .input(calldata.into()); - provider.call(tx).await - } - }) - .await?; - - let details = IRecurringCollector::getAgreementDetailsCall::abi_decode_returns(&output) - .map_err(|err| { - ChainClientError::RpcError(anyhow::anyhow!( - "undecodable getAgreementDetails from {collector}: {err}" - )) - })?; + Ok(still_live(self.agreement_state(agreement_id).await?)) + } - Ok(still_live(details.state)) + async fn agreement_ended_by_indexer( + &self, + agreement_id: &[u8; 16], + ) -> Result { + Ok(self.agreement_state(agreement_id).await? & STATE_BY_PROVIDER != 0) } async fn reconcile_provider( diff --git a/bin/dipper-service/src/network/service/cancel_retry.rs b/bin/dipper-service/src/network/service/cancel_retry.rs index ca465fa9..ae44e32b 100644 --- a/bin/dipper-service/src/network/service/cancel_retry.rs +++ b/bin/dipper-service/src/network/service/cancel_retry.rs @@ -26,8 +26,8 @@ const SWEEP_BUDGET: std::time::Duration = std::time::Duration::from_secs(30); /// when it was marked can be mined first instead of being sent again. const SETTLE_MINUTES: i32 = 2; -/// How long the chain listener gets to record who ended an accepted agreement, and when, -/// before the retry marks it ended without those details. +/// How long the chain listener gets to record when, and in which transaction, dipper's own +/// cancel ended an accepted agreement, 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, @@ -96,29 +96,36 @@ async fn retry_cancel( } LiveCancel::CancelFailed(err) => (None, Some(err)), }; - if failure.is_none() && confirm_if_over(registry, config, row, tx_hash, chain_now).await { + if failure.is_none() + && confirm_if_over(registry, chain_client, config, row, tx_hash, chain_now).await + { return; } note_check(registry, row, failure.as_ref()).await; } /// Mark the agreement `CanceledByRequester` once it can't go live again: this sweep's cancel -/// ended it, or nobody accepted its offer before the deadline to. One accepted that ended -/// otherwise is left to the chain listener, which reads who ended it and when, for a while. -async fn confirm_if_over( +/// ended it, or nobody accepted its offer before the deadline to. One the indexer ended is +/// left to the chain listener, and so, for a while, is one dipper ended in an earlier send. +async fn confirm_if_over( registry: &R, + chain_client: &T, config: &IndexingAgreementConfig, row: &CancellingAgreement, tx_hash: Option, chain_now: u64, -) -> bool { +) -> bool +where + R: AgreementRegistry + Sync, + T: ChainClient, +{ let agreement = &row.agreement; let can_confirm = if row.accepted_on_chain { tx_hash.is_some() || agreement.updated_at < time::OffsetDateTime::now_utc() - LISTENER_GRACE } else { chain_now > agreement.terms.deadline }; - if !can_confirm { + if !can_confirm || (tx_hash.is_none() && ended_by_indexer(chain_client, agreement).await) { return false; } if let Err(err) = registry @@ -147,6 +154,32 @@ async fn confirm_if_over( true } +/// Whether the chain shows the indexer ended the agreement, which the chain listener records +/// as theirs. An unread answer counts as yes, so an end is never wrongly put down to dipper. +async fn ended_by_indexer(chain_client: &T, agreement: &IndexingAgreement) -> bool { + match chain_client + .agreement_ended_by_indexer(agreement.id.as_bytes()) + .await + { + Ok(false) => false, + Ok(true) => { + tracing::info!( + agreement_id = %agreement.id, + "The indexer ended an agreement dipper was cancelling; the chain listener records it" + ); + true + } + Err(err) => { + tracing::warn!( + agreement_id = %agreement.id, + error = %err, + "Failed to read who ended a cancelling agreement, will retry" + ); + true + } + } +} + /// Record that the agreement was checked and is still cancelling, counting a cancel the /// chain answered without ending it; past the limit, dipper gives up with an ERROR. async fn note_check( @@ -305,6 +338,8 @@ mod tests { send_fails: bool, mined_cancel_reverts: bool, cancel_has_no_effect: bool, + ended_by_indexer: bool, + who_read_fails: bool, cancels_sent: AtomicU32, } @@ -360,6 +395,15 @@ mod tests { } Ok(self.live.load(Ordering::SeqCst)) } + async fn agreement_ended_by_indexer( + &self, + _agreement_id: &[u8; 16], + ) -> Result { + if self.who_read_fails { + return Err(ChainClientError::RpcError(anyhow::anyhow!("rpc down"))); + } + Ok(self.ended_by_indexer) + } async fn latest_block_timestamp(&self) -> Result { unimplemented!() } @@ -430,6 +474,41 @@ mod tests { assert!(registry.audits.lock().unwrap().is_empty()); } + #[tokio::test] + async fn never_puts_an_end_by_the_indexer_down_to_dipper() { + // The listener records it as the indexer's. An accepted agreement can lack an accept + // time, if it was accepted before accepts were recorded, so both kinds are checked. + for accepted_on_chain in [true, false] { + let mut registry = registry_with_one(accepted_on_chain); + registry.cancelling[0].agreement.updated_at = + time::OffsetDateTime::now_utc() - LISTENER_GRACE - time::Duration::MINUTE; + let chain = MockChain { + ended_by_indexer: true, + ..MockChain::default() + }; + + retry(®istry, &chain, DEADLINE + 1).await; + + assert!(registry.marked_cancelled.lock().unwrap().is_empty()); + assert_eq!(registry.checks.load(Ordering::SeqCst), 1); + } + } + + #[tokio::test] + async fn does_not_confirm_an_end_when_who_ended_it_cannot_be_read() { + let mut registry = registry_with_one(true); + registry.cancelling[0].agreement.updated_at = + time::OffsetDateTime::now_utc() - LISTENER_GRACE - time::Duration::MINUTE; + let chain = MockChain { + who_read_fails: true, + ..MockChain::default() + }; + + retry(®istry, &chain, 0).await; + + assert!(registry.marked_cancelled.lock().unwrap().is_empty()); + } + #[tokio::test] async fn keeps_an_unaccepted_agreement_cancelling_until_its_deadline() { // An offer still in flight could land and be accepted until then. diff --git a/bin/dipper-service/src/network/service/chain_listener.rs b/bin/dipper-service/src/network/service/chain_listener.rs index 1b16dc77..145f8cb2 100644 --- a/bin/dipper-service/src/network/service/chain_listener.rs +++ b/bin/dipper-service/src/network/service/chain_listener.rs @@ -2663,6 +2663,12 @@ mod tests { // not-active here means "cancel confirmed", which these tests expect. Ok(self.live_until_cancelled && !self.cancels.lock().unwrap().contains(agreement_id)) } + async fn agreement_ended_by_indexer( + &self, + _agreement_id: &[u8; 16], + ) -> Result { + Ok(false) + } } #[async_trait::async_trait] diff --git a/bin/dipper-service/src/network/service/escrow_reconciler.rs b/bin/dipper-service/src/network/service/escrow_reconciler.rs index f24fbaf2..5fba0f1a 100644 --- a/bin/dipper-service/src/network/service/escrow_reconciler.rs +++ b/bin/dipper-service/src/network/service/escrow_reconciler.rs @@ -698,6 +698,12 @@ mod tests { ) -> Result { unimplemented!() } + async fn agreement_ended_by_indexer( + &self, + _agreement_id: &[u8; 16], + ) -> Result { + Ok(false) + } async fn reconcile_agreement( &self, _collector: Address, diff --git a/bin/dipper-service/src/network/service/liveness_checker.rs b/bin/dipper-service/src/network/service/liveness_checker.rs index 34e43ac2..9e520241 100644 --- a/bin/dipper-service/src/network/service/liveness_checker.rs +++ b/bin/dipper-service/src/network/service/liveness_checker.rs @@ -1209,6 +1209,12 @@ mod tests { // not-active means "cancel confirmed", which these tests expect. Ok(false) } + async fn agreement_ended_by_indexer( + &self, + _agreement_id: &[u8; 16], + ) -> Result { + Ok(false) + } } const DB_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5); diff --git a/bin/dipper-service/src/worker/handlers/cancel_rejected_agreement_on_chain.rs b/bin/dipper-service/src/worker/handlers/cancel_rejected_agreement_on_chain.rs index 46ae9be2..651245c0 100644 --- a/bin/dipper-service/src/worker/handlers/cancel_rejected_agreement_on_chain.rs +++ b/bin/dipper-service/src/worker/handlers/cancel_rejected_agreement_on_chain.rs @@ -639,6 +639,12 @@ mod tests { } Ok(self.live.load(Ordering::SeqCst)) } + async fn agreement_ended_by_indexer( + &self, + _agreement_id: &[u8; 16], + ) -> Result { + Ok(false) + } } fn test_agreement_conf() -> Arc { diff --git a/bin/dipper-service/src/worker/handlers/reassess_indexing_request.rs b/bin/dipper-service/src/worker/handlers/reassess_indexing_request.rs index 340bf26b..51b46811 100644 --- a/bin/dipper-service/src/worker/handlers/reassess_indexing_request.rs +++ b/bin/dipper-service/src/worker/handlers/reassess_indexing_request.rs @@ -1023,6 +1023,12 @@ mod lifecycle_event_tests { // Cancel confirmed: agreement is no longer active on-chain. Ok(false) } + async fn agreement_ended_by_indexer( + &self, + _agreement_id: &[u8; 16], + ) -> Result { + Ok(false) + } } // ---- Mock: registry (all five traits) ----------------------------------- @@ -2454,6 +2460,12 @@ mod deadline_clock_tests { ) -> Result { Ok(false) } + async fn agreement_ended_by_indexer( + &self, + _agreement_id: &[u8; 16], + ) -> Result { + Ok(false) + } } /// `fail: true` makes the state lookup itself error, exercising the diff --git a/bin/dipper-service/src/worker/handlers/submit_offer.rs b/bin/dipper-service/src/worker/handlers/submit_offer.rs index e27f3d74..2a24a1a4 100644 --- a/bin/dipper-service/src/worker/handlers/submit_offer.rs +++ b/bin/dipper-service/src/worker/handlers/submit_offer.rs @@ -418,6 +418,12 @@ mod tests { ) -> Result { Ok(self.on_chain.load(Ordering::SeqCst)) } + async fn agreement_ended_by_indexer( + &self, + _agreement_id: &[u8; 16], + ) -> Result { + Ok(false) + } async fn latest_block_timestamp(&self) -> Result { unimplemented!() } From 65fc1328ff26c6a02c250e8a65e4d70b24dece97 Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Fri, 2 Oct 2026 22:13:53 +0300 Subject: [PATCH 3/5] fix(cancel): retry cancels of agreements that may be paying first The cancel retry took 10 agreements per run, oldest check first, so offers nobody had accepted could hold up a live agreement for hours. It now takes up to 50, putting first those accepted or past their offer deadline, and its 30 second time limit decides how many it gets through. --- .../src/network/service/cancel_retry.rs | 5 +++-- bin/dipper-service/src/registry/agreement.rs | 3 ++- dipper-pgregistry/src/postgres.rs | 9 +++++++-- dipper-pgregistry/tests/it_registry_postgres.rs | 16 +++++++++++++--- 4 files changed, 25 insertions(+), 8 deletions(-) diff --git a/bin/dipper-service/src/network/service/cancel_retry.rs b/bin/dipper-service/src/network/service/cancel_retry.rs index ae44e32b..a62f044d 100644 --- a/bin/dipper-service/src/network/service/cancel_retry.rs +++ b/bin/dipper-service/src/network/service/cancel_retry.rs @@ -15,8 +15,9 @@ use crate::{ /// leaves it to an operator. Other failures don't count (see `failed_attempts`). pub const MAX_CANCEL_ATTEMPTS: u32 = 10; -/// Agreements checked per sweep, those checked longest ago first. -const BATCH_SIZE: i64 = 10; +/// Agreements a sweep takes on, those that may be paying an indexer first; the time budget +/// below decides how many it gets through. +const BATCH_SIZE: i64 = 50; /// Time a sweep may take before leaving the rest to the next one: it holds up the chain /// listener while it runs, and each cancel can wait up to 15 s to be mined. diff --git a/bin/dipper-service/src/registry/agreement.rs b/bin/dipper-service/src/registry/agreement.rs index 36d2c9ae..406204b0 100644 --- a/bin/dipper-service/src/registry/agreement.rs +++ b/bin/dipper-service/src/registry/agreement.rs @@ -274,7 +274,8 @@ pub trait AgreementRegistry { ) -> RegistryResult<()>; /// `CANCELLING` agreements marked over `min_age_minutes` ago whose cancel has failed - /// fewer than `max_attempts` times, those checked longest ago first. + /// fewer than `max_attempts` times. Those that may be paying an indexer come first (accepted, + /// or past the offer deadline, which only an accepted one outlives), then those checked longest ago. async fn get_cancelling_agreements( &self, batch_size: i64, diff --git a/dipper-pgregistry/src/postgres.rs b/dipper-pgregistry/src/postgres.rs index d40439b3..8111f80d 100644 --- a/dipper-pgregistry/src/postgres.rs +++ b/dipper-pgregistry/src/postgres.rs @@ -972,7 +972,8 @@ impl PgRegistry { } /// `Cancelling` agreements marked over `min_age_minutes` ago whose cancel has failed - /// fewer than `max_attempts` times, those checked longest ago first. + /// fewer than `max_attempts` times. Those that may be paying an indexer come first (accepted, + /// or past the offer deadline, which only an accepted one outlives), then those checked longest ago. pub async fn get_cancelling_agreements( &self, batch_size: i64, @@ -1001,7 +1002,11 @@ impl PgRegistry { WHERE status = $1 AND cancel_attempts < $2 AND updated_at < timezone('UTC', now()) - make_interval(mins => $4) - ORDER BY cancel_checked_at ASC NULLS FIRST, updated_at ASC + ORDER BY + (accepted_at IS NOT NULL + OR CAST(terms->>'deadline' AS bigint) < EXTRACT(EPOCH FROM now())) DESC, + cancel_checked_at ASC NULLS FIRST, + updated_at ASC LIMIT $3 "#, ) diff --git a/dipper-pgregistry/tests/it_registry_postgres.rs b/dipper-pgregistry/tests/it_registry_postgres.rs index 18530ded..81e5d456 100644 --- a/dipper-pgregistry/tests/it_registry_postgres.rs +++ b/dipper-pgregistry/tests/it_registry_postgres.rs @@ -3361,8 +3361,18 @@ async fn cancelling_agreements_are_listed_until_their_cancel_fails_too_often() { ) .await .expect("Failed to run fixture"); - let registry = PgRegistry::new(db); let created = fixture_agreement(0xaa); + // An offer still open to acceptance, so only the accepted agreement can be paying. + sqlx::query( + "UPDATE dipper_reg_indexing_agreements \ + SET terms = jsonb_set(terms::jsonb, '{deadline}', to_jsonb(4102444800::bigint)) \ + WHERE id = $1", + ) + .bind(created) + .execute(&db) + .await + .expect("Failed to update deadline"); + let registry = PgRegistry::new(db); let accepted = IndexingAgreementId::from_bytes([0xaa, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 2]); let ended = @@ -3429,8 +3439,8 @@ async fn cancelling_agreements_are_listed_until_their_cancel_fails_too_often() { let ids: Vec<_> = listed.iter().map(|row| row.agreement.id).collect(); assert_eq!( ids, - vec![created, accepted], - "one checked longest ago goes first" + vec![accepted, created], + "one accepted on-chain may be paying its indexer, so it goes first" ); assert_eq!(registry.record_cancel_check(&created, 1).await.unwrap(), 1); From f973d46374e0a700ed1144107d444e14ed3767ce Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Fri, 2 Oct 2026 22:19:01 +0300 Subject: [PATCH 4/5] fix(cancel): close an indexer's cancel the listener never recorded An agreement the indexer ended was left for the chain listener to record, so if the listener never did, it stayed cancelling for good and its indexer was kept out of selection. After an hour, the cancel retry now marks it cancelled by the indexer itself, naming them as the one who ended it. --- .../src/network/service/cancel_retry.rs | 133 +++++++++++++++--- 1 file changed, 113 insertions(+), 20 deletions(-) diff --git a/bin/dipper-service/src/network/service/cancel_retry.rs b/bin/dipper-service/src/network/service/cancel_retry.rs index a62f044d..dc4874b2 100644 --- a/bin/dipper-service/src/network/service/cancel_retry.rs +++ b/bin/dipper-service/src/network/service/cancel_retry.rs @@ -2,13 +2,14 @@ //! `Cancelling` before its on-chain cancel goes out; this sweep re-sends the cancel while //! the chain shows it live, and marks it `CanceledByRequester` once it can no longer be. +use dipper_core::time::now_secs; use thegraph_core::alloy::primitives::B256; use crate::{ cancel_dispatch::{LiveCancel, cancel_if_live, record_cancel}, chain_client::{ChainClient, ChainClientError}, config::IndexingAgreementConfig, - registry::{AgreementRegistry, CancellingAgreement, IndexingAgreement}, + registry::{AgreementRegistry, CancelKind, CancellingAgreement, IndexingAgreement}, }; /// Retried cancels mined without ending an agreement before dipper stops retrying it and @@ -27,8 +28,8 @@ const SWEEP_BUDGET: std::time::Duration = std::time::Duration::from_secs(30); /// when it was marked can be mined first instead of being sent again. const SETTLE_MINUTES: i32 = 2; -/// How long the chain listener gets to record when, and in which transaction, dipper's own -/// cancel ended an accepted agreement, before the retry marks it ended without those details. +/// How long the chain listener gets to record when, and in which transaction, an accepted +/// 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, @@ -106,8 +107,8 @@ async fn retry_cancel( } /// Mark the agreement `CanceledByRequester` once it can't go live again: this sweep's cancel -/// ended it, or nobody accepted its offer before the deadline to. One the indexer ended is -/// left to the chain listener, and so, for a while, is one dipper ended in an earlier send. +/// ended it, or nobody accepted its offer before the deadline to. One ended otherwise is left +/// to the chain listener for a while; one the indexer ended then becomes `CanceledByIndexer`. async fn confirm_if_over( registry: &R, chain_client: &T, @@ -121,14 +122,22 @@ where T: ChainClient, { let agreement = &row.agreement; + let past_grace = agreement.updated_at < time::OffsetDateTime::now_utc() - LISTENER_GRACE; let can_confirm = if row.accepted_on_chain { - tx_hash.is_some() || agreement.updated_at < time::OffsetDateTime::now_utc() - LISTENER_GRACE + tx_hash.is_some() || past_grace } else { chain_now > agreement.terms.deadline }; - if !can_confirm || (tx_hash.is_none() && ended_by_indexer(chain_client, agreement).await) { + if !can_confirm { return false; } + if tx_hash.is_none() { + match ended_by_indexer(chain_client, agreement).await { + None => return false, + Some(true) => return past_grace && record_end_by_indexer(registry, agreement).await, + Some(false) => {} + } + } if let Err(err) = registry .mark_indexing_agreement_as_canceled_by_requester(&agreement.id) .await @@ -155,18 +164,65 @@ where true } -/// Whether the chain shows the indexer ended the agreement, which the chain listener records -/// as theirs. An unread answer counts as yes, so an end is never wrongly put down to dipper. -async fn ended_by_indexer(chain_client: &T, agreement: &IndexingAgreement) -> bool { +/// Whether the chain shows the indexer ended the agreement, or `None` when it can't be read, +/// so an end is never wrongly put down to dipper. +async fn ended_by_indexer( + chain_client: &T, + agreement: &IndexingAgreement, +) -> Option { match chain_client .agreement_ended_by_indexer(agreement.id.as_bytes()) .await { - Ok(false) => false, - Ok(true) => { + Ok(by_indexer) => { + if by_indexer { + tracing::info!( + agreement_id = %agreement.id, + "The indexer ended an agreement dipper was cancelling" + ); + } + Some(by_indexer) + } + Err(err) => { + tracing::warn!( + agreement_id = %agreement.id, + error = %err, + "Failed to read who ended a cancelling agreement, will retry" + ); + None + } + } +} + +/// Mark an agreement the indexer ended `CanceledByIndexer` when the chain listener hasn't in +/// time, so it doesn't stay cancelling for good. The indexer is recorded as ending it first, +/// so its announcement names them; the time recorded is when dipper noticed. +async fn record_end_by_indexer( + registry: &R, + agreement: &IndexingAgreement, +) -> bool { + let indexer = agreement.indexer.id.to_string(); + let marked = match registry + .record_cancel_audit(&agreement.id, now_secs(), &indexer, None) + .await + { + Ok(()) => { + registry + .apply_reconciliation(&agreement.id, false, Some(CancelKind::ByIndexer)) + .await + } + Err(err) => Err(err), + }; + match marked { + Ok(outcome) => { tracing::info!( agreement_id = %agreement.id, - "The indexer ended an agreement dipper was cancelling; the chain listener records it" + indexing_request_id = %agreement.indexing_request_id, + old_status = "CANCELLING", + new_status = "CANCELED_BY_INDEXER", + applied = outcome.did_cancel, + reason = "indexer_cancel_seen_on_chain", + "agreement state transition" ); true } @@ -174,9 +230,9 @@ async fn ended_by_indexer(chain_client: &T, agreement: &Indexing tracing::warn!( agreement_id = %agreement.id, error = %err, - "Failed to read who ended a cancelling agreement, will retry" + "Failed to mark an agreement the indexer ended, will retry" ); - true + false } } } @@ -286,7 +342,9 @@ mod tests { struct MockRegistry { cancelling: Vec, marked_cancelled: Mutex>, + marked_by_indexer: Mutex>, audits: Mutex>>, + audited_by: Mutex>, attempts: AtomicU32, checks: AtomicU32, } @@ -312,15 +370,29 @@ mod tests { &self, _id: &IndexingAgreementId, _canceled_at: u64, - _canceled_by: &str, + canceled_by: &str, canceled_tx: Option<&str>, ) -> crate::registry::Result<()> { + self.audited_by.lock().unwrap().push(canceled_by.to_owned()); self.audits .lock() .unwrap() .push(canceled_tx.map(str::to_owned)); Ok(()) } + async fn apply_reconciliation( + &self, + id: &IndexingAgreementId, + _apply_accept: bool, + cancel: Option, + ) -> crate::registry::Result { + assert_eq!(cancel, Some(CancelKind::ByIndexer)); + self.marked_by_indexer.lock().unwrap().push(*id); + Ok(crate::registry::ReconciliationOutcome { + did_accept: false, + did_cancel: true, + }) + } async fn record_cancel_check( &self, _id: &IndexingAgreementId, @@ -476,9 +548,26 @@ mod tests { } #[tokio::test] - async fn never_puts_an_end_by_the_indexer_down_to_dipper() { - // The listener records it as the indexer's. An accepted agreement can lack an accept - // time, if it was accepted before accepts were recorded, so both kinds are checked. + async fn leaves_an_end_by_the_indexer_to_the_listener_for_a_while() { + // The listener records when and in which transaction. An accepted agreement can lack + // an accept time, if accepted before accepts were recorded, so both kinds are checked. + for accepted_on_chain in [true, false] { + let registry = registry_with_one(accepted_on_chain); + let chain = MockChain { + ended_by_indexer: true, + ..MockChain::default() + }; + + retry(®istry, &chain, DEADLINE + 1).await; + + assert!(registry.marked_cancelled.lock().unwrap().is_empty()); + assert!(registry.marked_by_indexer.lock().unwrap().is_empty()); + assert_eq!(registry.checks.load(Ordering::SeqCst), 1); + } + } + + #[tokio::test] + async fn marks_an_end_by_the_indexer_as_theirs_once_the_listener_has_had_long_enough() { for accepted_on_chain in [true, false] { let mut registry = registry_with_one(accepted_on_chain); registry.cancelling[0].agreement.updated_at = @@ -491,7 +580,10 @@ mod tests { retry(®istry, &chain, DEADLINE + 1).await; assert!(registry.marked_cancelled.lock().unwrap().is_empty()); - assert_eq!(registry.checks.load(Ordering::SeqCst), 1); + assert_eq!(registry.marked_by_indexer.lock().unwrap().len(), 1); + let indexer = registry.cancelling[0].agreement.indexer.id.to_string(); + assert_eq!(*registry.audited_by.lock().unwrap(), vec![indexer]); + assert_eq!(*registry.audits.lock().unwrap(), vec![None]); } } @@ -508,6 +600,7 @@ mod tests { retry(®istry, &chain, 0).await; assert!(registry.marked_cancelled.lock().unwrap().is_empty()); + assert!(registry.marked_by_indexer.lock().unwrap().is_empty()); } #[tokio::test] From 578afdbfb4a7be973811886ba3538b6b1eb31965 Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Fri, 2 Oct 2026 22:19:06 +0300 Subject: [PATCH 5/5] fix(cancel): stop open offers waiting forever behind paying agreements Agreements that may be paying an indexer were always retried first, so if enough of them kept failing, offers an indexer could still accept were never retried. Those agreements now get an hour's head start instead, so an offer left unchecked for over an hour still gets its turn. --- bin/dipper-service/src/registry/agreement.rs | 5 +++-- dipper-pgregistry/src/postgres.rs | 14 +++++++----- .../tests/it_registry_postgres.rs | 22 ++++++++++++++++++- 3 files changed, 33 insertions(+), 8 deletions(-) diff --git a/bin/dipper-service/src/registry/agreement.rs b/bin/dipper-service/src/registry/agreement.rs index 406204b0..1309d1d3 100644 --- a/bin/dipper-service/src/registry/agreement.rs +++ b/bin/dipper-service/src/registry/agreement.rs @@ -274,8 +274,9 @@ pub trait AgreementRegistry { ) -> RegistryResult<()>; /// `CANCELLING` agreements marked over `min_age_minutes` ago whose cancel has failed - /// fewer than `max_attempts` times. Those that may be paying an indexer come first (accepted, - /// or past the offer deadline, which only an accepted one outlives), then those checked longest ago. + /// fewer than `max_attempts` times, those checked longest ago first. One that may be paying + /// an indexer (accepted, or past the offer deadline, which only an accepted one outlives) + /// counts as checked an hour earlier, so it goes first without holding the rest back. async fn get_cancelling_agreements( &self, batch_size: i64, diff --git a/dipper-pgregistry/src/postgres.rs b/dipper-pgregistry/src/postgres.rs index 8111f80d..fb36f95b 100644 --- a/dipper-pgregistry/src/postgres.rs +++ b/dipper-pgregistry/src/postgres.rs @@ -972,8 +972,9 @@ impl PgRegistry { } /// `Cancelling` agreements marked over `min_age_minutes` ago whose cancel has failed - /// fewer than `max_attempts` times. Those that may be paying an indexer come first (accepted, - /// or past the offer deadline, which only an accepted one outlives), then those checked longest ago. + /// fewer than `max_attempts` times, those checked longest ago first. One that may be paying + /// an indexer (accepted, or past the offer deadline, which only an accepted one outlives) + /// counts as checked an hour earlier, so it goes first without holding the rest back. pub async fn get_cancelling_agreements( &self, batch_size: i64, @@ -1003,9 +1004,12 @@ impl PgRegistry { AND cancel_attempts < $2 AND updated_at < timezone('UTC', now()) - make_interval(mins => $4) ORDER BY - (accepted_at IS NOT NULL - OR CAST(terms->>'deadline' AS bigint) < EXTRACT(EPOCH FROM now())) DESC, - cancel_checked_at ASC NULLS FIRST, + cancel_checked_at - CASE + WHEN accepted_at IS NOT NULL + OR CAST(terms->>'deadline' AS bigint) < EXTRACT(EPOCH FROM now()) + THEN INTERVAL '1 hour' + ELSE INTERVAL '0 seconds' + END ASC NULLS FIRST, updated_at ASC LIMIT $3 "#, diff --git a/dipper-pgregistry/tests/it_registry_postgres.rs b/dipper-pgregistry/tests/it_registry_postgres.rs index 81e5d456..ff1de13f 100644 --- a/dipper-pgregistry/tests/it_registry_postgres.rs +++ b/dipper-pgregistry/tests/it_registry_postgres.rs @@ -3372,7 +3372,7 @@ async fn cancelling_agreements_are_listed_until_their_cancel_fails_too_often() { .execute(&db) .await .expect("Failed to update deadline"); - let registry = PgRegistry::new(db); + let registry = PgRegistry::new(db.clone()); let accepted = IndexingAgreementId::from_bytes([0xaa, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 2]); let ended = @@ -3431,6 +3431,7 @@ async fn cancelling_agreements_are_listed_until_their_cancel_fails_too_often() { ] ); + assert_eq!(registry.record_cancel_check(&created, 0).await.unwrap(), 0); assert_eq!(registry.record_cancel_check(&accepted, 0).await.unwrap(), 0); let listed = registry .get_cancelling_agreements(100, 2, 0) @@ -3443,6 +3444,25 @@ async fn cancelling_agreements_are_listed_until_their_cancel_fails_too_often() { "one accepted on-chain may be paying its indexer, so it goes first" ); + sqlx::query( + "UPDATE dipper_reg_indexing_agreements \ + SET cancel_checked_at = cancel_checked_at - INTERVAL '2 hours' WHERE id = $1", + ) + .bind(created) + .execute(&db) + .await + .expect("Failed to age the check"); + let listed = registry + .get_cancelling_agreements(100, 2, 0) + .await + .expect("cancelling query"); + let ids: Vec<_> = listed.iter().map(|row| row.agreement.id).collect(); + assert_eq!( + ids, + vec![created, accepted], + "an offer unchecked for over an hour isn't held back for ever" + ); + assert_eq!(registry.record_cancel_check(&created, 1).await.unwrap(), 1); assert_eq!(registry.record_cancel_check(&created, 1).await.unwrap(), 2); let listed = registry