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..dc4874b2 100644 --- a/bin/dipper-service/src/network/service/cancel_retry.rs +++ b/bin/dipper-service/src/network/service/cancel_retry.rs @@ -2,21 +2,23 @@ //! `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 /// 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. @@ -26,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 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, 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, @@ -96,31 +98,46 @@ 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 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, config: &IndexingAgreementConfig, row: &CancellingAgreement, tx_hash: Option, chain_now: u64, -) -> bool { +) -> bool +where + R: AgreementRegistry + Sync, + 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 { 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 @@ -147,6 +164,79 @@ async fn confirm_if_over( true } +/// 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(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, + 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 + } + Err(err) => { + tracing::warn!( + agreement_id = %agreement.id, + error = %err, + "Failed to mark an agreement the indexer ended, will retry" + ); + false + } + } +} + /// 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( @@ -252,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, } @@ -278,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, @@ -305,6 +411,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 +468,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 +547,62 @@ mod tests { assert!(registry.audits.lock().unwrap().is_empty()); } + #[tokio::test] + 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 = + 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.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]); + } + } + + #[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()); + assert!(registry.marked_by_indexer.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/registry/agreement.rs b/bin/dipper-service/src/registry/agreement.rs index 36d2c9ae..1309d1d3 100644 --- a/bin/dipper-service/src/registry/agreement.rs +++ b/bin/dipper-service/src/registry/agreement.rs @@ -274,7 +274,9 @@ 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 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/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/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() { 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!() } diff --git a/dipper-pgregistry/src/postgres.rs b/dipper-pgregistry/src/postgres.rs index d40439b3..fb36f95b 100644 --- a/dipper-pgregistry/src/postgres.rs +++ b/dipper-pgregistry/src/postgres.rs @@ -972,7 +972,9 @@ 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 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, @@ -1001,7 +1003,14 @@ 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 + 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 18530ded..ff1de13f 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.clone()); let accepted = IndexingAgreementId::from_bytes([0xaa, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 2]); let ended = @@ -3421,16 +3431,36 @@ 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) .await .expect("cancelling query"); let ids: Vec<_> = listed.iter().map(|row| row.agreement.id).collect(); + assert_eq!( + ids, + vec![accepted, created], + "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], - "one checked longest ago goes first" + "an offer unchecked for over an hour isn't held back for ever" ); assert_eq!(registry.record_cancel_check(&created, 1).await.unwrap(), 1);