diff --git a/bin/dipper-service/src/admin_rpc_server/handlers/indexing_agreements.rs b/bin/dipper-service/src/admin_rpc_server/handlers/indexing_agreements.rs index 689b7415..827537ff 100644 --- a/bin/dipper-service/src/admin_rpc_server/handlers/indexing_agreements.rs +++ b/bin/dipper-service/src/admin_rpc_server/handlers/indexing_agreements.rs @@ -135,5 +135,6 @@ fn into_indexing_agreement_status( IndexingAgreementRecordStatus::AbandonedByIndexer => { IndexingAgreementStatus::AbandonedByIndexer } + IndexingAgreementRecordStatus::Cancelling => IndexingAgreementStatus::Cancelling, } } diff --git a/bin/dipper-service/src/cancel_dispatch.rs b/bin/dipper-service/src/cancel_dispatch.rs index 7452dbd6..e4c9c80f 100644 --- a/bin/dipper-service/src/cancel_dispatch.rs +++ b/bin/dipper-service/src/cancel_dispatch.rs @@ -1,12 +1,15 @@ //! On-chain cancel dispatch. Every cancel goes through //! [`cancel_agreement_on_chain`] so the manager-routed path lives in one place. +use dipper_core::time::now_secs; use thegraph_core::alloy::primitives::B256; use crate::{ chain_client::{ChainClient, ChainClientError}, config::IndexingAgreementConfig, - registry::IndexingAgreement, + registry::{ + AgreementRegistry, IndexingAgreement, IndexingAgreementStatus, Result as RegistryResult, + }, }; /// Pass both ACTIVE and PENDING; local status lags the chain, so let the @@ -59,8 +62,135 @@ pub async fn cancel_agreement_on_chain( Ok(outcome) } +/// What [`start_cancel`] left an agreement as. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum CancelStarted { + /// It was accepted and its cancel landed: now `CanceledByRequester`. + Ended, + /// Still `Cancelling`; the chain listener finishes it once it can't go live. + Cancelling, +} + +/// Start ending an agreement that may be live on-chain. It is marked `Cancelling` before +/// its cancel goes out, so an offer for it still in flight withdraws itself on landing. +/// Fails, sending nothing, when the mark can't be written. +pub async fn start_cancel( + registry: &R, + chain_client: &T, + agreement: &IndexingAgreement, + config: &IndexingAgreementConfig, +) -> RegistryResult +where + R: AgreementRegistry + Sync, + T: ChainClient, +{ + registry + .mark_indexing_agreement_as_cancelling(&agreement.id) + .await?; + let tx_hash = match cancel_agreement_on_chain(chain_client, agreement, config).await { + Ok(tx_hash) => tx_hash, + Err(err) => { + tracing::warn!( + agreement_id = %agreement.id, + error = %err, + "On-chain cancel failed; the chain listener retries it" + ); + return Ok(CancelStarted::Cancelling); + } + }; + tracing::info!( + agreement_id = %agreement.id, + tx_hash = ?tx_hash, + "Submitted on-chain cancellation" + ); + // An offer never accepted could still land and be accepted until its deadline. + if agreement.status != IndexingAgreementStatus::AcceptedOnChain { + return Ok(CancelStarted::Cancelling); + } + Ok(confirm_cancelled(registry, agreement, tx_hash, config).await) +} + +/// Mark an accepted agreement whose cancel landed `CanceledByRequester` and record the +/// cancel, so the `terminated` sweep announces it. +async fn confirm_cancelled( + registry: &R, + agreement: &IndexingAgreement, + tx_hash: Option, + config: &IndexingAgreementConfig, +) -> CancelStarted { + if let Err(err) = registry + .mark_indexing_agreement_as_canceled_by_requester(&agreement.id) + .await + { + tracing::warn!( + agreement_id = %agreement.id, + error = %err, + "Failed to mark a cancelled agreement; the chain listener finishes it" + ); + return CancelStarted::Cancelling; + } + record_cancel(registry, agreement, tx_hash, config).await; + CancelStarted::Ended +} + +/// Record dipper's own cancel of an accepted agreement, so the `terminated` sweep +/// announces it. +pub async fn record_cancel( + registry: &R, + agreement: &IndexingAgreement, + tx_hash: Option, + config: &IndexingAgreementConfig, +) { + let manager = config.recurring_agreement_manager().to_string(); + let tx = tx_hash.map(|hash| hash.to_string()); + if let Err(err) = registry + .record_cancel_audit(&agreement.id, now_secs(), &manager, tx.as_deref()) + .await + { + tracing::warn!( + agreement_id = %agreement.id, + error = %err, + "failed to record cancel audit; terminated event may emit with fallback fields" + ); + } +} + +/// What [`cancel_if_live`] found and did. +#[derive(Debug)] +pub enum LiveCancel { + /// The chain showed nothing live, so no cancel was sent. + NotLive, + /// A cancel went out and the chain confirmed the agreement ended. + Ended(Option), + /// The chain could not be read, so nothing was sent. + ReadFailed(ChainClientError), + /// The cancel failed or did not end the agreement. + CancelFailed(ChainClientError), +} + +/// Cancel an agreement on-chain only if the chain shows it live: a pending offer, or +/// accepted and not yet ended. Reading first saves a wasted transaction, since a cancel +/// of an agreement that already ended still mines. +pub async fn cancel_if_live( + chain_client: &T, + agreement: &IndexingAgreement, + config: &IndexingAgreementConfig, +) -> LiveCancel { + match chain_client + .agreement_still_active(agreement.id.as_bytes()) + .await + { + Err(err) => LiveCancel::ReadFailed(err), + Ok(false) => LiveCancel::NotLive, + Ok(true) => match cancel_agreement_on_chain(chain_client, agreement, config).await { + Ok(tx_hash) => LiveCancel::Ended(tx_hash), + Err(err) => LiveCancel::CancelFailed(err), + }, + } +} + #[cfg(test)] -mod tests { +pub(crate) mod tests { use std::sync::Mutex; use async_trait::async_trait; @@ -154,31 +284,16 @@ mod tests { fn manager_conf(collector: Address) -> IndexingAgreementConfig { IndexingAgreementConfig { - data_service: Address::ZERO, recurring_collector: collector, recurring_agreement_manager: Address::repeat_byte(0x33), - max_agreement_grt_per_30_days: 0.0, - max_seconds_per_collection: 0, - min_seconds_per_collection: 0, - duration_seconds: 0, - deadline_seconds: 0, - max_grt_per_30_days: std::collections::BTreeMap::new(), - max_grt_per_billion_entities_per_30_days: 0.0, - declined_indexer_lookback_days: 0, - price_rejection_lookback_days: 0, - transient_rejection_lookback_minutes: 0, - uncertain_rejection_lookback_days: 0, - unresponsive_indexer_lookback_days: 0, - mass_unresponsive_trip_fraction: 0.5, - mass_unresponsive_reset_fraction: 0.25, - dips_accepting_snapshot_max_age_hours: 48, - dips_accepting_cache_ttl_seconds: 300, - max_in_flight_offers_per_indexer: None, - max_in_flight_offers_total: None, + ..IndexingAgreementConfig::for_tests() } } - fn agreement(status: IndexingAgreementStatus, hash: Option>) -> IndexingAgreement { + pub(crate) fn agreement( + status: IndexingAgreementStatus, + hash: Option>, + ) -> IndexingAgreement { let deployment_id: DeploymentId = "QmTXzATwNfgGVukV1fX2T6xw9f6LAYRVWpsdXyRWzUR2H9" .parse() .unwrap(); diff --git a/bin/dipper-service/src/config.rs b/bin/dipper-service/src/config.rs index b480a33e..740974f3 100644 --- a/bin/dipper-service/src/config.rs +++ b/bin/dipper-service/src/config.rs @@ -1202,6 +1202,37 @@ pub struct IndexingAgreementConfig { pub max_in_flight_offers_total: Option, } +#[cfg(test)] +impl IndexingAgreementConfig { + /// Zero addresses and limits, with permissive breaker and cache settings, for tests + /// to adjust the fields they care about. + pub fn for_tests() -> Self { + Self { + data_service: Address::ZERO, + recurring_collector: Address::ZERO, + recurring_agreement_manager: Address::ZERO, + max_agreement_grt_per_30_days: 0.0, + max_seconds_per_collection: 0, + min_seconds_per_collection: 0, + duration_seconds: 0, + deadline_seconds: 0, + max_grt_per_30_days: BTreeMap::new(), + max_grt_per_billion_entities_per_30_days: 0.0, + declined_indexer_lookback_days: 0, + price_rejection_lookback_days: 0, + transient_rejection_lookback_minutes: 0, + uncertain_rejection_lookback_days: 0, + unresponsive_indexer_lookback_days: 0, + mass_unresponsive_trip_fraction: 0.5, + mass_unresponsive_reset_fraction: 0.25, + dips_accepting_snapshot_max_age_hours: 48, + dips_accepting_cache_ttl_seconds: 300, + max_in_flight_offers_per_indexer: None, + max_in_flight_offers_total: None, + } + } +} + /// Per-chain pricing for indexing agreements (runtime). #[derive(Debug)] pub struct IndexingAgreementChainPrices { diff --git a/bin/dipper-service/src/network/service.rs b/bin/dipper-service/src/network/service.rs index 86f9bb3e..3fc9d0f0 100644 --- a/bin/dipper-service/src/network/service.rs +++ b/bin/dipper-service/src/network/service.rs @@ -1,3 +1,4 @@ +pub mod cancel_retry; pub mod chain_events; pub mod chain_listener; pub mod domain_refresh; diff --git a/bin/dipper-service/src/network/service/cancel_retry.rs b/bin/dipper-service/src/network/service/cancel_retry.rs new file mode 100644 index 00000000..ca465fa9 --- /dev/null +++ b/bin/dipper-service/src/network/service/cancel_retry.rs @@ -0,0 +1,551 @@ +//! Finishes the cancels dipper starts. An agreement dipper wants ended is marked +//! `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 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}, +}; + +/// 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; + +/// 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. +const SWEEP_BUDGET: std::time::Duration = std::time::Duration::from_secs(30); + +/// Minutes an agreement stays out of the retry after it is marked, so the cancel sent +/// 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. +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. +pub async fn retry_cancelling_agreements( + registry: &R, + chain_client: &T, + config: &IndexingAgreementConfig, + chain_now: u64, +) where + R: AgreementRegistry + Sync, + T: ChainClient, +{ + let cancelling = match registry + .get_cancelling_agreements(BATCH_SIZE, MAX_CANCEL_ATTEMPTS, SETTLE_MINUTES) + .await + { + Ok(cancelling) => cancelling, + Err(err) => { + tracing::warn!(error = %err, "Failed to list agreements still being cancelled"); + return; + } + }; + let started = std::time::Instant::now(); + for (done, row) in cancelling.iter().enumerate() { + if started.elapsed() >= SWEEP_BUDGET { + tracing::info!( + left = cancelling.len() - done, + "Cancel retry ran out of time; the rest wait for the next sweep" + ); + break; + } + retry_cancel(registry, chain_client, config, row, chain_now).await; + } +} + +async fn retry_cancel( + registry: &R, + chain_client: &T, + config: &IndexingAgreementConfig, + row: &CancellingAgreement, + chain_now: u64, +) where + R: AgreementRegistry + Sync, + T: ChainClient, +{ + let agreement_id = row.agreement.id; + let (tx_hash, failure) = match cancel_if_live(chain_client, &row.agreement, config).await { + LiveCancel::ReadFailed(err) => { + tracing::warn!( + %agreement_id, + error = %err, + "Failed to read a cancelling agreement on-chain, will retry" + ); + // Unread, it may still be live, so it can't be confirmed ended. + return note_check(registry, row, None).await; + } + LiveCancel::NotLive => (None, None), + LiveCancel::Ended(tx_hash) => { + tracing::info!( + %agreement_id, + tx_hash = ?tx_hash, + "Cancelled an agreement still live on-chain" + ); + (tx_hash, None) + } + LiveCancel::CancelFailed(err) => (None, Some(err)), + }; + if failure.is_none() && confirm_if_over(registry, 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( + registry: &R, + config: &IndexingAgreementConfig, + row: &CancellingAgreement, + tx_hash: Option, + chain_now: u64, +) -> bool { + 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 { + return false; + } + if let Err(err) = registry + .mark_indexing_agreement_as_canceled_by_requester(&agreement.id) + .await + { + tracing::warn!( + agreement_id = %agreement.id, + error = %err, + "Failed to mark an ended agreement cancelled, will retry" + ); + return false; + } + tracing::info!( + agreement_id = %agreement.id, + indexing_request_id = %agreement.indexing_request_id, + old_status = "CANCELLING", + new_status = "CANCELED_BY_REQUESTER", + reason = "cancel_confirmed_on_chain", + "agreement state transition" + ); + // Without its own transaction, a late read by the chain listener fills in the cancel. + if row.accepted_on_chain && tx_hash.is_some() { + record_cancel(registry, agreement, tx_hash, config).await; + } + 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( + registry: &R, + row: &CancellingAgreement, + failure: Option<&ChainClientError>, +) { + let agreement = &row.agreement; + let failed_attempts = failure.map_or(0, failed_attempts); + if let Some(err) = failure + && failed_attempts == 0 + { + log_uncounted_failure(agreement, err); + } + match registry + .record_cancel_check(&agreement.id, failed_attempts) + .await + { + Ok(attempts) => { + if let Some(err) = failure.filter(|_| failed_attempts > 0) { + log_failed_cancel(agreement, attempts, err); + } + } + Err(err) => tracing::warn!( + agreement_id = %agreement.id, + error = %err, + "Failed to record a check of a cancelling agreement" + ), + } +} + +/// A revert before sending is logged as an ERROR: a paused or misconfigured manager makes +/// every cancel revert, so it is retried rather than counted against the agreement. +fn log_uncounted_failure(agreement: &IndexingAgreement, err: &ChainClientError) { + if matches!(err, ChainClientError::ContractRevert { .. }) { + tracing::error!( + agreement_id = %agreement.id, + error = %err, + "Cancel of an agreement reverted before it was sent, will retry" + ); + } else { + tracing::warn!( + agreement_id = %agreement.id, + error = %err, + "Cancel of an agreement failed or could not be confirmed, will retry" + ); + } +} + +fn log_failed_cancel(agreement: &IndexingAgreement, attempts: u32, err: &ChainClientError) { + if attempts < MAX_CANCEL_ATTEMPTS { + tracing::warn!( + agreement_id = %agreement.id, + attempts, + error = %err, + "Cancel did not end the agreement, will retry" + ); + return; + } + tracing::error!( + event = "agreement_cancel_abandoned", + agreement_id = %agreement.id, + indexer_id = %agreement.indexer.id, + indexing_request_id = %agreement.indexing_request_id, + attempts, + error = %err, + "Gave up cancelling an agreement on-chain; it may still be live" + ); +} + +/// How many of an agreement's cancel attempts a failure uses up. A cancel mined without +/// ending it, or reverted, counts, and one that can never be sent uses them all. An +/// unreachable chain or a revert before sending is retried freely. +fn failed_attempts(err: &ChainClientError) -> u32 { + match err { + ChainClientError::CancelNotConfirmed { .. } | ChainClientError::TxReverted { .. } => 1, + ChainClientError::MissingTermsVersionHash { .. } => MAX_CANCEL_ATTEMPTS, + _ => 0, + } +} + +#[cfg(test)] +mod tests { + use std::sync::{ + Mutex, + atomic::{AtomicBool, AtomicU32, Ordering}, + }; + + use async_trait::async_trait; + use dipper_core::ids::IndexingAgreementId; + use dipper_rpc::indexer::indexer_client::sol::RecurringCollectionAgreement; + use thegraph_core::alloy::primitives::Address; + + use super::*; + use crate::{ + cancel_dispatch::tests::agreement, + registry::{IndexingAgreementStatus, StubAgreementRegistry}, + }; + + const DEADLINE: u64 = 1_000; + + #[derive(Default)] + struct MockRegistry { + cancelling: Vec, + marked_cancelled: Mutex>, + audits: Mutex>>, + attempts: AtomicU32, + checks: AtomicU32, + } + + #[async_trait] + impl StubAgreementRegistry for MockRegistry { + async fn get_cancelling_agreements( + &self, + _batch_size: i64, + _max_attempts: u32, + _min_age_minutes: i32, + ) -> crate::registry::Result> { + Ok(self.cancelling.clone()) + } + async fn mark_indexing_agreement_as_canceled_by_requester( + &self, + id: &IndexingAgreementId, + ) -> crate::registry::Result<()> { + self.marked_cancelled.lock().unwrap().push(*id); + Ok(()) + } + async fn record_cancel_audit( + &self, + _id: &IndexingAgreementId, + _canceled_at: u64, + _canceled_by: &str, + canceled_tx: Option<&str>, + ) -> crate::registry::Result<()> { + self.audits + .lock() + .unwrap() + .push(canceled_tx.map(str::to_owned)); + Ok(()) + } + async fn record_cancel_check( + &self, + _id: &IndexingAgreementId, + failed_attempts: u32, + ) -> crate::registry::Result { + self.checks.fetch_add(1, Ordering::SeqCst); + Ok(self.attempts.fetch_add(failed_attempts, Ordering::SeqCst) + failed_attempts) + } + } + + /// An agreement live on-chain until a cancel ends it, unless set to ignore cancels. + #[derive(Default)] + struct MockChain { + live: AtomicBool, + read_fails: bool, + send_fails: bool, + mined_cancel_reverts: bool, + cancel_has_no_effect: bool, + cancels_sent: AtomicU32, + } + + #[async_trait] + impl ChainClient for MockChain { + async fn offer_via_manager( + &self, + _rca: &RecurringCollectionAgreement, + ) -> Result, ChainClientError> { + unimplemented!() + } + async fn cancel_via_manager( + &self, + _collector: Address, + _agreement_id: &[u8; 16], + _version_hash: B256, + _options: u16, + ) -> Result, ChainClientError> { + if self.send_fails { + return Err(ChainClientError::RpcError(anyhow::anyhow!("rpc down"))); + } + if self.mined_cancel_reverts { + return Err(ChainClientError::TxReverted { + tx_hash: B256::repeat_byte(0xee), + }); + } + self.cancels_sent.fetch_add(1, Ordering::SeqCst); + if !self.cancel_has_no_effect { + self.live.store(false, Ordering::SeqCst); + } + Ok(Some(B256::repeat_byte(0xcd))) + } + async fn reconcile_provider( + &self, + _collector: Address, + _provider: Address, + ) -> Result, ChainClientError> { + unimplemented!() + } + async fn reconcile_agreement( + &self, + _collector: Address, + _agreement_id: &[u8; 16], + ) -> Result, ChainClientError> { + unimplemented!() + } + async fn agreement_still_active( + &self, + _agreement_id: &[u8; 16], + ) -> Result { + if self.read_fails { + return Err(ChainClientError::RpcError(anyhow::anyhow!("rpc down"))); + } + Ok(self.live.load(Ordering::SeqCst)) + } + async fn latest_block_timestamp(&self) -> Result { + unimplemented!() + } + } + + fn registry_with_one(accepted_on_chain: bool) -> MockRegistry { + let mut cancelling = agreement(IndexingAgreementStatus::Cancelling, Some(vec![7u8; 32])); + cancelling.terms.deadline = DEADLINE; + MockRegistry { + cancelling: vec![CancellingAgreement { + agreement: cancelling, + accepted_on_chain, + }], + ..MockRegistry::default() + } + } + + fn live_chain() -> MockChain { + MockChain { + live: AtomicBool::new(true), + ..MockChain::default() + } + } + + async fn retry(registry: &MockRegistry, chain: &MockChain, chain_now: u64) { + let config = IndexingAgreementConfig::for_tests(); + retry_cancelling_agreements(registry, chain, &config, chain_now).await; + } + + #[tokio::test] + async fn cancels_a_live_accepted_agreement_and_records_the_cancel() { + let registry = registry_with_one(true); + let chain = live_chain(); + + retry(®istry, &chain, 0).await; + + assert_eq!(chain.cancels_sent.load(Ordering::SeqCst), 1); + assert_eq!(registry.marked_cancelled.lock().unwrap().len(), 1); + let tx = B256::repeat_byte(0xcd).to_string(); + assert_eq!(*registry.audits.lock().unwrap(), vec![Some(tx)]); + } + + #[tokio::test] + async fn leaves_an_accepted_agreement_that_already_ended_to_the_listener() { + // The indexer may have ended it, or an earlier cancel whose result went unread; + // the chain listener reads which, and records when and in which transaction. + let registry = registry_with_one(true); + let chain = MockChain::default(); + + retry(®istry, &chain, 0).await; + + assert_eq!(chain.cancels_sent.load(Ordering::SeqCst), 0); + assert!(registry.marked_cancelled.lock().unwrap().is_empty()); + assert!(registry.audits.lock().unwrap().is_empty()); + } + + #[tokio::test] + async fn marks_an_ended_accepted_agreement_itself_once_the_listener_has_had_long_enough() { + // In case the listener never reads its end, it would otherwise stay cancelling. + 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::default(); + + retry(®istry, &chain, 0).await; + + assert_eq!(registry.marked_cancelled.lock().unwrap().len(), 1); + assert!(registry.audits.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. + let registry = registry_with_one(false); + let chain = MockChain::default(); + + retry(®istry, &chain, DEADLINE).await; + assert!(registry.marked_cancelled.lock().unwrap().is_empty()); + + retry(®istry, &chain, DEADLINE + 1).await; + assert_eq!(registry.marked_cancelled.lock().unwrap().len(), 1); + assert!(registry.audits.lock().unwrap().is_empty()); + } + + #[tokio::test] + async fn withdraws_a_live_offer_before_its_deadline_and_keeps_watching() { + let registry = registry_with_one(false); + let chain = live_chain(); + + retry(®istry, &chain, 0).await; + + assert_eq!(chain.cancels_sent.load(Ordering::SeqCst), 1); + assert!(registry.marked_cancelled.lock().unwrap().is_empty()); + } + + #[tokio::test] + async fn notes_each_check_that_leaves_an_agreement_cancelling() { + // So the next sweep starts with the agreements checked longest ago. + let registry = registry_with_one(false); + let chain = MockChain::default(); + + retry(®istry, &chain, DEADLINE).await; + + assert_eq!(registry.checks.load(Ordering::SeqCst), 1); + assert_eq!(registry.attempts.load(Ordering::SeqCst), 0); + } + + #[tokio::test] + async fn never_confirms_an_agreement_it_could_not_read() { + // Past its deadline but unread, it may be live: an accept the listener hasn't + // recorded, or one from before accepts were recorded. + let registry = registry_with_one(false); + let chain = MockChain { + read_fails: true, + ..live_chain() + }; + + 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 an_unreachable_chain_neither_sends_nor_counts_an_attempt() { + for chain in [ + MockChain { + read_fails: true, + ..live_chain() + }, + MockChain { + send_fails: true, + ..live_chain() + }, + ] { + let registry = registry_with_one(true); + + retry(®istry, &chain, 0).await; + + assert_eq!(chain.cancels_sent.load(Ordering::SeqCst), 0); + assert_eq!(registry.attempts.load(Ordering::SeqCst), 0); + assert!(registry.marked_cancelled.lock().unwrap().is_empty()); + } + } + + #[tokio::test] + async fn gives_up_at_once_on_an_agreement_it_can_never_cancel() { + // Without a stored terms hash no cancel can be sent, so retrying only delays the alert. + let mut registry = registry_with_one(true); + registry.cancelling[0].agreement.terms_version_hash = None; + let chain = live_chain(); + + retry(®istry, &chain, 0).await; + + assert_eq!( + registry.attempts.load(Ordering::SeqCst), + MAX_CANCEL_ATTEMPTS + ); + assert_eq!(chain.cancels_sent.load(Ordering::SeqCst), 0); + } + + #[tokio::test] + async fn counts_a_cancel_that_is_mined_and_reverts() { + // Each one costs gas, so it can't be retried without limit. + let registry = registry_with_one(true); + let chain = MockChain { + mined_cancel_reverts: true, + ..live_chain() + }; + + retry(®istry, &chain, 0).await; + + assert_eq!(registry.attempts.load(Ordering::SeqCst), 1); + } + + #[tokio::test] + async fn counts_a_cancel_that_mines_without_ending_the_agreement() { + let registry = registry_with_one(true); + let chain = MockChain { + cancel_has_no_effect: true, + ..live_chain() + }; + + retry(®istry, &chain, 0).await; + + assert_eq!(registry.attempts.load(Ordering::SeqCst), 1); + assert!(registry.marked_cancelled.lock().unwrap().is_empty()); + } +} diff --git a/bin/dipper-service/src/network/service/chain_listener.rs b/bin/dipper-service/src/network/service/chain_listener.rs index f5ce990d..1b16dc77 100644 --- a/bin/dipper-service/src/network/service/chain_listener.rs +++ b/bin/dipper-service/src/network/service/chain_listener.rs @@ -83,6 +83,9 @@ const SWEEP_BATCH_SIZE: i64 = 1000; /// crash-recovery; the steady-state fan-out fires from finalize on a /// fresh accept, so per-poll execution is wasted DB work. const SWEEP_POLLS: u64 = 60; +/// How often agreements still being cancelled get their cancel retried, whichever rate the +/// listener polls at: just under the slow poll interval, so every slow poll retries. +const CANCEL_RETRY_INTERVAL: Duration = Duration::from_secs(290); /// Handle for controlling the chain listener service lifecycle #[derive(Clone)] @@ -228,6 +231,7 @@ where // Starts at SWEEP_POLLS so the first poll runs the sweep, // recovering any pre-startup orphans. let mut polls_since_sweep: u64 = SWEEP_POLLS; + let mut last_cancel_retry: Option = None; // Pause the event sweeps after a Kafka send failure, backing off from one // poll interval up to the idle interval, so a hung broker cannot stall the // poll loop on every iteration. @@ -293,6 +297,21 @@ where tokio::time::sleep(backoff).await; } + // Ahead of the drain, which ends the poll early while the subgraph is down: + // 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; + } + let outcome = match drain_once( &mut cursor, &mut last_persisted_timestamp, @@ -877,6 +896,7 @@ where } queue_cancel_if_cancelled_but_accepted(snapshot, &agreement, worker_queue).await?; + record_accept_of_cancelling(snapshot, &agreement, registry).await; // Both transitions are applied atomically downstream so the // Accept-then-Cancel-in-one-snapshot path can't leak an intermediate @@ -936,9 +956,10 @@ where } } -/// How long after an offer's deadline dipper's cancel can still have raced an -/// accept the listener hadn't caught up with. -const LISTENER_LAG_SLACK_SECS: i64 = 3_600; +/// When v0.1.10, the first release that announces lifecycle events, came out (2026-08-11 +/// UTC). Agreements dipper created before then are never announced, even when a replay of +/// the chain reads them again. +const LIFECYCLE_EVENTS_START: i64 = 1_786_406_400; /// Record the accept and cancel of an agreement dipper had already marked /// cancelled that went live on-chain first, so its accepted and terminated @@ -949,7 +970,7 @@ async fn record_accept_and_cancel_from_chain( agreement: &IndexingAgreement, registry: &R, ) { - if snapshot.accepted_at == 0 || !cancelled_while_offer_open(agreement) { + if snapshot.accepted_at == 0 || !created_after_events_started(agreement) { return; } let canceled_by = snapshot.canceled_by.to_string(); @@ -979,12 +1000,10 @@ async fn record_accept_and_cancel_from_chain( } } -/// Whether dipper cancelled the agreement while its offer could still be accepted, -/// so an accept on-chain may have slipped past it. One accepted before lifecycle -/// events existed was cancelled long after its deadline and is never announced. -fn cancelled_while_offer_open(agreement: &IndexingAgreement) -> bool { - let deadline = i64::try_from(agreement.terms.deadline).unwrap_or(i64::MAX); - agreement.updated_at.unix_timestamp() <= deadline.saturating_add(LISTENER_LAG_SLACK_SECS) +/// Whether dipper created the agreement once it announced lifecycle events. Its own clock +/// at creation, unlike any read of the chain, doesn't depend on how far the listener lags. +fn created_after_events_started(agreement: &IndexingAgreement) -> bool { + agreement.created_at.unix_timestamp() >= LIFECYCLE_EVENTS_START } /// Safety net for an agreement dipper cancelled whose offer the indexer accepted @@ -1012,6 +1031,33 @@ async fn queue_cancel_if_cancelled_but_accepted( Ok(()) } +/// An agreement dipper is cancelling stays `Cancelling` when the chain shows it accepted, +/// so its accept is recorded here; its end is then announced, along with the accept, +/// once the cancel lands. A withdrawn offer reads as cancelled with no accept time. +async fn record_accept_of_cancelling( + snapshot: &AgreementStateSnapshot, + agreement: &IndexingAgreement, + registry: &R, +) { + if agreement.status != IndexingAgreementStatus::Cancelling + || !snapshot.state.reached_accepted() + || snapshot.accepted_at == 0 + || !created_after_events_started(agreement) + { + return; + } + if let Err(err) = registry + .record_accepted_audit(&agreement.id, snapshot.accepted_at, &snapshot.accepted_tx) + .await + { + tracing::warn!( + agreement_id = %agreement.id, + error = %err, + "failed to record the on-chain accept of an agreement being cancelled" + ); + } +} + /// Log the transition that landed and, on fresh accepts, fan out the /// linked pending cancellations. #[expect( @@ -1115,16 +1161,9 @@ where /// Execute pending cancellations linked to a newly-accepted agreement. /// /// Called from the Created -> AcceptedOnChain and Expired -> AcceptedOnChain -/// transitions. For each pending cancellation, fires -/// `cancelIndexingAgreementByPayer` against the RecurringCollector contract and -/// flips the dipper DB row to `CanceledByRequester`, an unaccepted one first. Each pending row is -/// deleted individually after both steps succeed; transient failures retain -/// the record so the next reconcile pass can retry. -#[expect( - clippy::cognitive_complexity, - clippy::too_many_lines, - reason = "predates this lint; fix when next touched" -)] +/// transitions. Each replaced agreement is marked `Cancelling` and its on-chain cancel +/// sent once; the cancel retry in this listener finishes any that don't end at once. +/// A pending row is deleted once its agreement is marked; a failed mark keeps it for retry. async fn execute_pending_cancellations( agreement_id: &IndexingAgreementId, registry: &R, @@ -1139,178 +1178,134 @@ where .get_pending_cancellations_by_new_agreement(*agreement_id) .await?; - if pending.is_empty() { - return Ok(()); - } - let mut transient_failures: u32 = 0; - for cancellation in &pending { - let old_agreement = match registry - .get_indexing_agreement_by_id(&cancellation.old_agreement_id) - .await? - { - None => { - tracing::warn!( - old_agreement_id = %cancellation.old_agreement_id, - "Pending cancellation references non-existent agreement, cleaning up" - ); - registry - .delete_pending_cancellation(*agreement_id, cancellation.old_agreement_id) - .await?; - continue; - } - Some(a) => a, - }; - - // An unaccepted one is marked first. An offer for it still in flight then - // lands before this cancel, which withdraws it, or sees the mark and withdraws itself. - let unaccepted = old_agreement.status == IndexingAgreementStatus::Created; - if unaccepted { - match registry - .mark_indexing_agreement_as_canceled_by_requester(&cancellation.old_agreement_id) - .await - { - Ok(()) | Err(crate::registry::Error::NoRecordsUpdated) => {} - Err(err) => { - tracing::error!( - old_agreement_id = %cancellation.old_agreement_id, - error = %err, - "Failed to mark replaced agreement cancelled before its on-chain cancel, \ - retaining pending row" - ); - transient_failures += 1; - continue; - } - } - } - - let mut on_chain_cancel_tx: Option = None; - match crate::cancel_dispatch::cancel_agreement_on_chain( + let old_agreement_id = cancellation.old_agreement_id; + if start_replaced_cancel( + agreement_id, + &old_agreement_id, + registry, chain_client, - &old_agreement, config, ) - .await + .await? { - Ok(Some(tx_hash)) => { - tracing::info!( - new_agreement_id = %agreement_id, - old_agreement_id = %cancellation.old_agreement_id, - %tx_hash, - "Submitted on-chain cancellation for replaced agreement" - ); - on_chain_cancel_tx = Some(tx_hash.to_string()); - } - Ok(None) => { - tracing::info!( - new_agreement_id = %agreement_id, - old_agreement_id = %cancellation.old_agreement_id, - "Agreement already canceled on-chain; proceeding with local cleanup" - ); - } - Err(err) => { - tracing::warn!( - old_agreement_id = %cancellation.old_agreement_id, - error = %err, - "On-chain cancel failed, retaining pending cancellation for retry" - ); - transient_failures += 1; - continue; - } + registry + .delete_pending_cancellation(*agreement_id, old_agreement_id) + .await?; + } else { + transient_failures += 1; } + } - if !unaccepted { - match registry - .mark_indexing_agreement_as_canceled_by_requester(&cancellation.old_agreement_id) - .await - { - Ok(()) => {} - Err(crate::registry::Error::NoRecordsUpdated) => { - tracing::debug!( - old_agreement_id = %cancellation.old_agreement_id, - "Old agreement already in terminal state, skipping local cancel flip" - ); - registry - .delete_pending_cancellation(*agreement_id, cancellation.old_agreement_id) - .await?; - continue; - } - Err(err) => { - tracing::error!( - old_agreement_id = %cancellation.old_agreement_id, - error = %err, - "On-chain cancel succeeded but DB update failed, retaining pending row" - ); - transient_failures += 1; - continue; - } - } - } + if transient_failures > 0 { + anyhow::bail!( + "{transient_failures} pending cancellation(s) failed for agreement {}; \ + records retained for retry", + agreement_id, + ); + } - registry - .delete_pending_cancellation(*agreement_id, cancellation.old_agreement_id) - .await?; + Ok(()) +} - tracing::info!( - new_agreement_id = %agreement_id, - old_agreement_id = %cancellation.old_agreement_id, - "Canceled old agreement on-chain and in dipper DB after replacement confirmed" +/// Start cancelling an agreement its accepted replacement supersedes. Returns whether its +/// pending cancellation is done with: marked, already ended, or gone. +async fn start_replaced_cancel( + new_agreement_id: &IndexingAgreementId, + old_agreement_id: &IndexingAgreementId, + registry: &R, + chain_client: &T, + config: &crate::config::IndexingAgreementConfig, +) -> anyhow::Result +where + R: AgreementRegistry + Sync, + T: ChainClient, +{ + let Some(old_agreement) = registry + .get_indexing_agreement_by_id(old_agreement_id) + .await? + else { + tracing::warn!( + %old_agreement_id, + "Pending cancellation references non-existent agreement, cleaning up" ); - - // Record the cancel audit so `sweep_pending_terminated_events` can emit - // the `terminated` durably. Only for an accepted agreement: one that never - // was gets the chain's own cancel data if the indexer accepts it after all. - if old_agreement.status != IndexingAgreementStatus::AcceptedOnChain { - continue; + return Ok(true); + }; + if old_agreement.status == IndexingAgreementStatus::Expired { + match expired_but_live(chain_client, &old_agreement).await { + Some(true) => {} + Some(false) => return Ok(true), + None => return Ok(false), } - let manager = config.recurring_agreement_manager().to_string(); - if let Err(err) = registry - .record_cancel_audit( - &cancellation.old_agreement_id, - dipper_core::time::now_secs(), - &manager, - on_chain_cancel_tx.as_deref(), - ) - .await - { + } + let started = + crate::cancel_dispatch::start_cancel(registry, chain_client, &old_agreement, config).await; + Ok(note_replaced_cancel( + new_agreement_id, + &old_agreement, + started, + )) +} + +/// Whether an agreement marked `Expired` is live on-chain after all, accepted unseen by a +/// lagging listener; `None`, logged, when the chain can't be read. One that really expired +/// stays `Expired`. +async fn expired_but_live( + chain_client: &T, + agreement: &IndexingAgreement, +) -> Option { + match chain_client + .agreement_still_active(agreement.id.as_bytes()) + .await + { + Ok(live) => Some(live), + Err(err) => { tracing::warn!( - old_agreement_id = %cancellation.old_agreement_id, + old_agreement_id = %agreement.id, error = %err, - "failed to record cancel audit; terminated event may emit with fallback fields" + "Failed to read whether a replaced expired agreement is live, retaining pending row" ); + None } } +} - if transient_failures > 0 { - anyhow::bail!( - "{transient_failures} pending cancellation(s) failed for agreement {}; \ - records retained for retry", - agreement_id, - ); +/// Log how cancelling a replaced agreement started; false when it couldn't be marked. +fn note_replaced_cancel( + new_agreement_id: &IndexingAgreementId, + old_agreement: &IndexingAgreement, + started: crate::registry::Result, +) -> bool { + let old_agreement_id = old_agreement.id; + match started { + Ok(started) => tracing::info!( + %new_agreement_id, + %old_agreement_id, + old_status = %old_agreement.status, + ?started, + reason = "replacement_accepted", + "Cancelling replaced agreement" + ), + Err(crate::registry::Error::NoRecordsUpdated) => tracing::debug!( + %old_agreement_id, + "Replaced agreement already ended or being cancelled" + ), + Err(err) => { + tracing::error!( + %old_agreement_id, + error = %err, + "Failed to mark replaced agreement cancelling, retaining pending row" + ); + return false; + } } - - Ok(()) + true } -/// Retry on-chain cancels for agreements orphaned by a failed shrink-to-zero. -/// -/// When a `set_indexing_target_candidates(num_candidates = 0)` call flips the -/// request row to `Canceled`, reassessment fires `cancelIndexingAgreementByPayer` -/// for every agreement under it. A transient chain-client error during that -/// fan-out leaves the request row `Canceled` and at least one agreement still -/// `AcceptedOnChain` — the local intent and on-chain state disagree, and the -/// admin RPC has nothing left to trigger. -/// -/// This sweep runs periodically on the chain_listener tick and re-fires the -/// on-chain cancel for each such orphan. The chain-side cancel is idempotent -/// (the `Ok(None)` revert path handles already-canceled agreements), and the -/// DB transition is gated on chain success, so this is safe to run on every -/// sweep without coordination with reassessment. -#[expect( - clippy::cognitive_complexity, - reason = "predates this lint; fix when next touched" -)] +/// Start cancelling agreements orphaned by a failed shrink-to-zero: still `AcceptedOnChain` +/// although their request is `Canceled`, because reassessment couldn't mark them. Each starts +/// like any other cancel, so the cancel retry finishes it with the same limit. async fn sweep_orphan_canceled_agreements( registry: &R, chain_client: &T, @@ -1333,76 +1328,29 @@ async fn sweep_orphan_canceled_agreements( } }; - if orphans.is_empty() { - return; - } - - tracing::debug!( - count = orphans.len(), - "Sweeping orphan agreements whose parent request is Canceled" - ); - for agreement in orphans { - let mut on_chain_cancel_tx: Option = None; - match crate::cancel_dispatch::cancel_agreement_on_chain(chain_client, &agreement, config) - .await - { - Ok(Some(tx_hash)) => { - tracing::info!( - agreement_id = %agreement.id, - %tx_hash, - "Submitted on-chain cancel for orphan agreement" - ); - on_chain_cancel_tx = Some(tx_hash.to_string()); - } - Ok(None) => { - tracing::info!( - agreement_id = %agreement.id, - "Orphan agreement already canceled on-chain; cleaning up local state" - ); - } - Err(err) => { - tracing::warn!( - error = %err, - agreement_id = %agreement.id, - "Failed to cancel orphan agreement on-chain; will retry next sweep" - ); - continue; - } - } - - if let Err(err) = registry - .mark_indexing_agreement_as_canceled_by_requester(&agreement.id) - .await - { - tracing::error!( - error = %err, - agreement_id = %agreement.id, - "Failed to mark orphan agreement as canceled in local DB" - ); - continue; - } + let started = + crate::cancel_dispatch::start_cancel(registry, chain_client, &agreement, config).await; + log_orphan_cancel(&agreement, started); + } +} - // Orphan (previously accepted) agreement canceled on-chain by dipper. - // Record the cancel audit; `sweep_pending_terminated_events` emits the - // `terminated` durably (the row is `AcceptedOnChain` -> terminal, so it is - // sweep-eligible). - let manager = config.recurring_agreement_manager().to_string(); - if let Err(err) = registry - .record_cancel_audit( - &agreement.id, - dipper_core::time::now_secs(), - &manager, - on_chain_cancel_tx.as_deref(), - ) - .await - { - tracing::warn!( - agreement_id = %agreement.id, - error = %err, - "failed to record cancel audit; terminated event may emit with fallback fields" - ); - } +fn log_orphan_cancel( + agreement: &IndexingAgreement, + started: crate::registry::Result, +) { + match started { + Ok(started) => tracing::info!( + agreement_id = %agreement.id, + ?started, + reason = "request_canceled", + "Cancelling orphan agreement" + ), + Err(err) => tracing::warn!( + error = %err, + agreement_id = %agreement.id, + "Failed to mark orphan agreement cancelling; will retry next sweep" + ), } } @@ -1903,29 +1851,7 @@ mod tests { } fn test_agreement_conf() -> std::sync::Arc { - std::sync::Arc::new(crate::config::IndexingAgreementConfig { - data_service: thegraph_core::alloy::primitives::Address::ZERO, - recurring_collector: thegraph_core::alloy::primitives::Address::ZERO, - recurring_agreement_manager: thegraph_core::alloy::primitives::Address::ZERO, - max_agreement_grt_per_30_days: 0.0, - max_seconds_per_collection: 0, - min_seconds_per_collection: 0, - duration_seconds: 0, - deadline_seconds: 0, - max_grt_per_30_days: std::collections::BTreeMap::new(), - max_grt_per_billion_entities_per_30_days: 0.0, - declined_indexer_lookback_days: 0, - price_rejection_lookback_days: 0, - transient_rejection_lookback_minutes: 0, - uncertain_rejection_lookback_days: 0, - unresponsive_indexer_lookback_days: 0, - mass_unresponsive_trip_fraction: 0.5, - mass_unresponsive_reset_fraction: 0.25, - dips_accepting_snapshot_max_age_hours: 48, - dips_accepting_cache_ttl_seconds: 300, - max_in_flight_offers_per_indexer: None, - max_in_flight_offers_total: None, - }) + std::sync::Arc::new(crate::config::IndexingAgreementConfig::for_tests()) } // Deterministic audit values used by `make_snapshot` so emit tests can assert @@ -1975,6 +1901,7 @@ mod tests { agreements: std::collections::HashMap, marked_accepted_on_chain: Vec, marked_canceled_by_requester: Vec, + marked_cancelling: Vec, marked_canceled_by_indexer: Vec, /// Ids passed to `record_cancel_audit` -- the signal a cancel path drives /// the terminated event (the sweep emits from this audit). @@ -2057,6 +1984,12 @@ mod tests { } } + fn set_agreement_created_at(&self, agreement_id: IndexingAgreementId, unix: i64) { + if let Some(a) = self.state.lock().unwrap().agreements.get_mut(&agreement_id) { + a.created_at = OffsetDateTime::from_unix_timestamp(unix).unwrap(); + } + } + fn set_agreement_request_id( &self, agreement_id: IndexingAgreementId, @@ -2075,6 +2008,10 @@ mod tests { .contains(id) } + fn was_marked_cancelling(&self, id: &IndexingAgreementId) -> bool { + self.state.lock().unwrap().marked_cancelling.contains(id) + } + fn was_marked_canceled_by_requester(&self, id: &IndexingAgreementId) -> bool { self.state .lock() @@ -2250,6 +2187,39 @@ mod tests { Ok(()) } + async fn mark_indexing_agreement_as_cancelling( + &self, + id: &IndexingAgreementId, + ) -> RegistryResult<()> { + let mut state = self.state.lock().unwrap(); + if state.fail_cancel_for.contains(id) { + return Err(crate::registry::Error::BackendError( + dipper_pgregistry::Error::DbError(sqlx::Error::Protocol( + "simulated transient failure".into(), + )), + )); + } + state.marked_cancelling.push(*id); + Ok(()) + } + + async fn get_cancelling_agreements( + &self, + _batch_size: i64, + _max_attempts: u32, + _min_age_minutes: i32, + ) -> RegistryResult> { + Ok(Vec::new()) + } + + async fn record_cancel_check( + &self, + _id: &IndexingAgreementId, + failed_attempts: u32, + ) -> RegistryResult { + Ok(failed_attempts) + } + async fn record_cancel_audit( &self, agreement_id: &IndexingAgreementId, @@ -2322,12 +2292,16 @@ mod tests { Some( IndexingAgreementStatus::Created | IndexingAgreementStatus::AcceptedOnChain - | IndexingAgreementStatus::Rejected, + | IndexingAgreementStatus::Rejected + | IndexingAgreementStatus::Cancelling, ), ), Some(crate::registry::CancelKind::ByIndexer) => matches!( effective_status_for_cancel, - Some(IndexingAgreementStatus::AcceptedOnChain), + Some( + IndexingAgreementStatus::AcceptedOnChain + | IndexingAgreementStatus::Cancelling, + ), ), None => false, }; @@ -2594,6 +2568,9 @@ mod tests { /// When set, each cancel records whether its agreement was already marked. registry: Option, marked_at_cancel: Arc>>, + fail_cancels: bool, + /// When set, every agreement reads as live until a cancel is sent for it. + live_until_cancelled: bool, } impl MockChainClient { @@ -2638,12 +2615,17 @@ mod tests { Option, crate::chain_client::ChainClientError, > { + if self.fail_cancels { + return Err(crate::chain_client::ChainClientError::RpcError( + anyhow::anyhow!("rpc down"), + )); + } // Record manager-routed cancels through the same recorder so the // existing assertions hold. self.cancels.lock().unwrap().push(*agreement_id); if let Some(registry) = &self.registry { let id = IndexingAgreementId::from_bytes(*agreement_id); - let was_marked = registry.was_marked_canceled_by_requester(&id); + let was_marked = registry.was_marked_cancelling(&id); self.marked_at_cancel.lock().unwrap().push(was_marked); } // A manager-routed cancel has no "already canceled" result: the @@ -2675,11 +2657,11 @@ mod tests { async fn agreement_still_active( &self, - _agreement_id: &[u8; 16], + agreement_id: &[u8; 16], ) -> Result { // Cancel dispatch always reads back after a mined cancel; reporting // not-active here means "cancel confirmed", which these tests expect. - Ok(false) + Ok(self.live_until_cancelled && !self.cancels.lock().unwrap().contains(agreement_id)) } } @@ -2756,6 +2738,120 @@ mod tests { assert!(!worker_queue.was_cancellation_queued(&agreement_id)); } + #[tokio::test] + async fn reconcile_records_the_accept_of_an_agreement_being_cancelled() { + // It stays cancelling, and the recorded accept lets its end be announced. + 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, IndexingAgreementStatus::Cancelling); + + let snapshot = make_snapshot(agreement_id, AgreementState::Accepted, Address::ZERO); + 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.audit_writes(), vec![("accept", agreement_id)]); + assert!(!worker_queue.was_cancellation_queued(&agreement_id)); + } + + #[tokio::test] + async fn reconcile_records_no_accept_of_an_agreement_from_before_events_existed() { + 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, IndexingAgreementStatus::Cancelling); + registry.set_agreement_created_at(agreement_id, LIFECYCLE_EVENTS_START - 1); + + let snapshot = make_snapshot(agreement_id, AgreementState::Accepted, Address::ZERO); + reconcile_agreement( + &snapshot, + ®istry, + &worker_queue, + &chain_client, + test_agreement_conf().as_ref(), + ) + .await + .expect("reconcile ok"); + + assert!(registry.audit_writes().is_empty()); + } + + #[tokio::test] + async fn reconcile_records_no_accept_for_a_withdrawn_offer() { + // The subgraph reports a withdrawn offer as cancelled by the payer with no accept + // time; recording it would announce an agreement that was never live. + 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, IndexingAgreementStatus::Cancelling); + 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_canceled_by_requester(&agreement_id)); + 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 [ + (AgreementState::CanceledByPayer, true), + (AgreementState::CanceledByServiceProvider, false), + ] { + 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, IndexingAgreementStatus::Cancelling); + + let snapshot = make_snapshot(agreement_id, state, Address::ZERO); + reconcile_agreement( + &snapshot, + ®istry, + &worker_queue, + &chain_client, + test_agreement_conf().as_ref(), + ) + .await + .expect("reconcile ok"); + + assert_eq!( + registry.was_marked_canceled_by_requester(&agreement_id), + ended_by_dipper + ); + assert_eq!( + registry.was_marked_canceled_by_indexer(&agreement_id), + !ended_by_dipper + ); + } + } + #[tokio::test] async fn reconcile_accept_marks_accepted_without_emitting() { use crate::test_support::CapturingEventsProducer; @@ -3095,8 +3191,9 @@ mod tests { let worker_queue = MockWorkerQueue::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); registry.add_agreement(agreement_id, IndexingAgreementStatus::CanceledByRequester); - // Dipper cancelled it while its offer could still be accepted. - registry.set_agreement_deadline_from_now(agreement_id, 600); + // The listener lagged: dipper only cancelled it locally long after the offer's + // deadline, which once made this accept look too old to announce. + registry.set_agreement_deadline_from_now(agreement_id, -2 * 86_400); let snapshot = make_snapshot(agreement_id, AgreementState::CanceledByPayer, Address::ZERO); let result = reconcile_agreement( @@ -3124,7 +3221,6 @@ mod tests { let worker_queue = MockWorkerQueue::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); registry.add_agreement(agreement_id, IndexingAgreementStatus::CanceledByRequester); - registry.set_agreement_deadline_from_now(agreement_id, 600); registry.state.lock().unwrap().fail_cancel_audit = true; let snapshot = make_snapshot(agreement_id, AgreementState::CanceledByPayer, Address::ZERO); @@ -3143,15 +3239,14 @@ mod tests { #[tokio::test] async fn test_reconcile_does_not_announce_an_agreement_accepted_before_events_existed() { - // Agreements accepted before lifecycle events existed have no recorded - // accept and are never announced. Dipper cancelled this one long after its - // offer's deadline, so it can't be an accept dipper missed. + // Agreements created before lifecycle events existed have no recorded accept + // and are never announced, even when a replay of the chain reads them again. 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, IndexingAgreementStatus::CanceledByRequester); - registry.set_agreement_deadline_from_now(agreement_id, -2 * 86_400); + registry.set_agreement_created_at(agreement_id, LIFECYCLE_EVENTS_START - 1); let snapshot = make_snapshot(agreement_id, AgreementState::CanceledByPayer, Address::ZERO); let result = reconcile_agreement( @@ -3393,10 +3488,10 @@ mod tests { } #[tokio::test] - async fn test_pending_cancellations_records_no_audit_for_a_never_accepted_agreement() { - // A cancel record on an agreement that was never accepted would later win - // over the chain's own, if the indexer accepted it after all and it was - // then ended: the terminated event would report an end before the accept. + async fn test_pending_cancellations_leave_a_never_accepted_agreement_cancelling() { + // Its offer could still land and be accepted until the deadline, so the cancel + // retry finishes it then. No cancel is recorded: one would win over the chain's + // own if the indexer accepted it after all, reporting an end before the accept. let registry = MockRegistry::new(); let chain_client = MockChainClient::default(); let new_id = IndexingAgreementId::from_bytes(rand::random()); @@ -3415,12 +3510,13 @@ mod tests { .await; assert!(result.is_ok()); - assert!(registry.was_marked_canceled_by_requester(&old_id)); + assert!(registry.was_marked_cancelling(&old_id)); + assert!(!registry.was_marked_canceled_by_requester(&old_id)); assert!(!registry.was_cancel_audit_recorded(&old_id)); } /// Runs the pending cancellation of one old agreement in `status`, returning - /// whether it was already marked cancelled when its on-chain cancel went out. + /// whether it was already marked cancelling when its on-chain cancel went out. async fn marked_at_pending_cancel(status: IndexingAgreementStatus) -> Vec { let registry = MockRegistry::new(); let chain_client = MockChainClient { @@ -3442,7 +3538,6 @@ mod tests { .await; assert!(result.is_ok(), "got {result:?}"); - assert!(registry.was_marked_canceled_by_requester(&old_id)); assert!(registry.was_pending_cancellation_deleted(&new_id, &old_id)); chain_client.marked_at_cancel.lock().unwrap().clone() } @@ -3456,9 +3551,71 @@ mod tests { } #[tokio::test] - async fn test_pending_cancellations_mark_an_accepted_agreement_after_its_cancel() { + async fn test_pending_cancellations_cancel_an_expired_agreement_only_if_live() { + // A lagging listener can mark an agreement expired that was in fact accepted; one + // that really expired keeps its status, and no needless cancel is sent. + for live in [true, false] { + let registry = MockRegistry::new(); + let chain_client = MockChainClient { + live_until_cancelled: live, + ..MockChainClient::default() + }; + let new_id = IndexingAgreementId::from_bytes(rand::random()); + let old_id = IndexingAgreementId::from_bytes(rand::random()); + registry.add_agreement(new_id, IndexingAgreementStatus::AcceptedOnChain); + registry.add_agreement(old_id, IndexingAgreementStatus::Expired); + registry.add_pending_cancellation(new_id, old_id); + + let result = execute_pending_cancellations( + &new_id, + ®istry, + &chain_client, + test_agreement_conf().as_ref(), + ) + .await; + + assert!(result.is_ok(), "got {result:?}"); + assert_eq!(registry.was_marked_cancelling(&old_id), live); + assert_eq!(chain_client.was_on_chain_cancel_attempted(&old_id), live); + assert!(registry.was_pending_cancellation_deleted(&new_id, &old_id)); + } + } + + #[tokio::test] + async fn test_pending_cancellations_mark_an_accepted_agreement_before_its_cancel() { + // A failed cancel then leaves it cancelling, which the cancel retry picks up. let marked = marked_at_pending_cancel(IndexingAgreementStatus::AcceptedOnChain).await; - assert_eq!(marked, vec![false]); + assert_eq!(marked, vec![true]); + } + + #[tokio::test] + async fn test_pending_cancellations_leave_a_failed_cancel_to_the_cancel_retry() { + // Retried from here, a cancel that never works was resent on every sweep with no + // limit, for as long as the replacement stayed accepted. The cancel retry limits it. + let registry = MockRegistry::new(); + let chain_client = MockChainClient { + fail_cancels: true, + ..MockChainClient::default() + }; + let new_id = IndexingAgreementId::from_bytes(rand::random()); + let old_id = IndexingAgreementId::from_bytes(rand::random()); + registry.add_agreement(new_id, IndexingAgreementStatus::AcceptedOnChain); + registry.add_agreement(old_id, IndexingAgreementStatus::AcceptedOnChain); + registry.add_pending_cancellation(new_id, old_id); + + let result = execute_pending_cancellations( + &new_id, + ®istry, + &chain_client, + test_agreement_conf().as_ref(), + ) + .await; + + assert!(result.is_ok(), "got {result:?}"); + assert!(registry.was_marked_cancelling(&old_id)); + assert!(!registry.was_marked_canceled_by_requester(&old_id)); + assert!(!registry.was_cancel_audit_recorded(&old_id)); + assert!(registry.was_pending_cancellation_deleted(&new_id, &old_id)); } #[tokio::test] @@ -3481,9 +3638,10 @@ mod tests { ) .await; - // The chain cancel failed, the row is NOT marked canceled, and no cancel - // audit is recorded (so the sweep emits nothing). + // The row couldn't be marked, so no cancel went out, nothing is recorded (the + // sweep emits nothing), and the pending row stays for a retry. assert!(result.is_err()); + assert!(!chain_client.was_on_chain_cancel_attempted(&old_fail)); assert!(!registry.was_marked_canceled_by_requester(&old_fail)); assert!(!registry.was_cancel_audit_recorded(&old_fail)); } @@ -4355,6 +4513,27 @@ mod tests { ); } + #[tokio::test] + async fn test_orphan_sweep_leaves_a_failed_cancel_to_the_cancel_retry() { + // Retried by this sweep, a cancel that never works was resent with no limit. + let registry = MockRegistry::new(); + let chain_client = MockChainClient { + fail_cancels: true, + ..MockChainClient::default() + }; + let agreement_id = IndexingAgreementId::from_bytes(rand::random()); + let request_id = IndexingRequestId::new(); + registry.add_agreement(agreement_id, IndexingAgreementStatus::AcceptedOnChain); + registry.set_agreement_request_id(agreement_id, request_id); + registry.mark_request_canceled(request_id); + + sweep_orphan_canceled_agreements(®istry, &chain_client, test_agreement_conf().as_ref()) + .await; + + assert!(registry.was_marked_cancelling(&agreement_id)); + assert!(!registry.was_marked_canceled_by_requester(&agreement_id)); + } + #[tokio::test] async fn test_orphan_sweep_records_cancel_audit() { // The orphan sweep no longer emits `terminated` directly; it records the diff --git a/bin/dipper-service/src/network/service/liveness_checker.rs b/bin/dipper-service/src/network/service/liveness_checker.rs index 09716aba..34e43ac2 100644 --- a/bin/dipper-service/src/network/service/liveness_checker.rs +++ b/bin/dipper-service/src/network/service/liveness_checker.rs @@ -1216,29 +1216,7 @@ mod tests { /// Default agreement config for the cancel-path tests. fn test_agreement_conf() -> crate::config::IndexingAgreementConfig { - crate::config::IndexingAgreementConfig { - data_service: thegraph_core::alloy::primitives::Address::ZERO, - recurring_collector: thegraph_core::alloy::primitives::Address::ZERO, - recurring_agreement_manager: thegraph_core::alloy::primitives::Address::ZERO, - max_agreement_grt_per_30_days: 0.0, - max_seconds_per_collection: 0, - min_seconds_per_collection: 0, - duration_seconds: 0, - deadline_seconds: 0, - max_grt_per_30_days: std::collections::BTreeMap::new(), - max_grt_per_billion_entities_per_30_days: 0.0, - declined_indexer_lookback_days: 0, - price_rejection_lookback_days: 0, - transient_rejection_lookback_minutes: 0, - uncertain_rejection_lookback_days: 0, - unresponsive_indexer_lookback_days: 0, - mass_unresponsive_trip_fraction: 0.5, - mass_unresponsive_reset_fraction: 0.25, - dips_accepting_snapshot_max_age_hours: 48, - dips_accepting_cache_ttl_seconds: 300, - max_in_flight_offers_per_indexer: None, - max_in_flight_offers_total: None, - } + crate::config::IndexingAgreementConfig::for_tests() } // ---- Pure function tests ---- diff --git a/bin/dipper-service/src/registry.rs b/bin/dipper-service/src/registry.rs index 82d9baf6..7d30552c 100644 --- a/bin/dipper-service/src/registry.rs +++ b/bin/dipper-service/src/registry.rs @@ -23,8 +23,8 @@ pub use self::agreement_stub::StubAgreementRegistry; use self::result::Result as RegistryResult; pub use self::{ agreement::{ - AgreementFeeRate, AgreementRegistry, CancelKind, IndexingAgreement, NewAgreementParams, - ReconciliationAudit, ReconciliationItem, ReconciliationOutcome, + AgreementFeeRate, AgreementRegistry, CancelKind, CancellingAgreement, IndexingAgreement, + NewAgreementParams, ReconciliationAudit, ReconciliationItem, ReconciliationOutcome, Status as IndexingAgreementStatus, Terms as IndexingAgreementTerms, TermsMetadata as IndexingAgreementTermsMetadata, }, @@ -391,6 +391,43 @@ impl AgreementRegistry for RegistryProvider { .map_err(Into::into) } + async fn mark_indexing_agreement_as_cancelling( + &self, + id: &IndexingAgreementId, + ) -> RegistryResult<()> { + self.inner + .mark_indexing_agreement_as_cancelling(id) + .await + .map_err(Into::into) + } + + async fn get_cancelling_agreements( + &self, + batch_size: i64, + max_attempts: u32, + min_age_minutes: i32, + ) -> RegistryResult> { + Ok(self + .inner + .get_cancelling_agreements(batch_size, max_attempts, min_age_minutes) + .await? + .into_iter() + .map(CancellingAgreement::try_from) + .filter_map(filter_map_with_logging) + .collect()) + } + + async fn record_cancel_check( + &self, + id: &IndexingAgreementId, + failed_attempts: u32, + ) -> RegistryResult { + self.inner + .record_cancel_check(id, failed_attempts) + .await + .map_err(Into::into) + } + async fn apply_reconciliation( &self, id: &IndexingAgreementId, diff --git a/bin/dipper-service/src/registry/agreement.rs b/bin/dipper-service/src/registry/agreement.rs index 32227b2c..36d2c9ae 100644 --- a/bin/dipper-service/src/registry/agreement.rs +++ b/bin/dipper-service/src/registry/agreement.rs @@ -259,13 +259,37 @@ pub trait AgreementRegistry { /// Mark an indexing agreement as `CANCELED_BY_REQUESTER`. /// /// If there is no indexing agreement with the given ID, or if the agreement is not in the - /// `CREATED` or `ACCEPTED_ON_CHAIN` state, this method returns a + /// `CREATED`, `ACCEPTED_ON_CHAIN`, `REJECTED` or `CANCELLING` state, this method returns a /// [`NoRecordUpdated`](Error::NoRecordsUpdated) error. async fn mark_indexing_agreement_as_canceled_by_requester( &self, id: &IndexingAgreementId, ) -> RegistryResult<()>; + /// Mark a `CREATED`, `ACCEPTED_ON_CHAIN`, `REJECTED` or `EXPIRED` agreement `CANCELLING`, + /// before its cancel is sent; [`NoRecordUpdated`](Error::NoRecordsUpdated) otherwise. + async fn mark_indexing_agreement_as_cancelling( + &self, + id: &IndexingAgreementId, + ) -> RegistryResult<()>; + + /// `CANCELLING` agreements marked over `min_age_minutes` ago whose cancel has failed + /// fewer than `max_attempts` times, those checked longest ago first. + async fn get_cancelling_agreements( + &self, + batch_size: i64, + max_attempts: u32, + min_age_minutes: i32, + ) -> RegistryResult>; + + /// Record a check of a `CANCELLING` agreement that left it cancelling, adding + /// `failed_attempts` to its failed cancels and returning the new count. + async fn record_cancel_check( + &self, + id: &IndexingAgreementId, + failed_attempts: u32, + ) -> RegistryResult; + /// Apply a reconciliation-driven state transition atomically. /// /// Used by `chain_listener::reconcile_agreement` so the @@ -522,6 +546,25 @@ pub struct AgreementFeeRate { pub tokens_per_entity_per_second: f64, } +/// An agreement dipper is still cancelling on-chain. +#[derive(Debug, Clone)] +pub struct CancellingAgreement { + pub agreement: IndexingAgreement, + /// Whether dipper saw it accepted on-chain, so its end is announced. + pub accepted_on_chain: bool, +} + +impl TryFrom for CancellingAgreement { + type Error = anyhow::Error; + + fn try_from(value: dipper_pgregistry::CancellingAgreement) -> Result { + Ok(Self { + agreement: value.agreement.try_into()?, + accepted_on_chain: value.accepted_on_chain, + }) + } +} + /// An Indexing Agreement represents the contract between the DIPs Gateway (Dipper) and the indexer /// to index the data. /// @@ -689,6 +732,11 @@ pub enum Status { /// /// This is a terminal state. AbandonedByIndexer, + + /// Dipper decided to end the agreement and is cancelling it on-chain, where it may + /// still be live. It becomes `CanceledByRequester`, announced as ended, only once + /// the chain confirms the end. + Cancelling, } impl std::fmt::Display for Status { @@ -702,6 +750,7 @@ impl std::fmt::Display for Status { Status::AcceptedOnChain => "ACCEPTED_ON_CHAIN", Status::Rejected => "REJECTED", Status::AbandonedByIndexer => "ABANDONED_BY_INDEXER", + Status::Cancelling => "CANCELLING", }; f.write_str(status) } @@ -733,6 +782,7 @@ impl TryFrom for IndexingAgreement { dipper_pgregistry::IndexingAgreementStatus::AbandonedByIndexer => { Status::AbandonedByIndexer } + dipper_pgregistry::IndexingAgreementStatus::Cancelling => Status::Cancelling, _ => { return Err(anyhow::anyhow!("Invalid status: {:?}", value.status)); } diff --git a/bin/dipper-service/src/registry/agreement_stub.rs b/bin/dipper-service/src/registry/agreement_stub.rs index ec9bfd3b..3e7477be 100644 --- a/bin/dipper-service/src/registry/agreement_stub.rs +++ b/bin/dipper-service/src/registry/agreement_stub.rs @@ -9,9 +9,9 @@ use thegraph_core::{DeploymentId, IndexerId, alloy::primitives::ChainId}; use super::{ agreement::{ - AgreementFeeRate, AgreementRegistry, CancelKind, IndexingAgreement, NewAgreementParams, - PendingAcceptedEvent, PendingExpiredEvent, PendingTerminatedEvent, ReconciliationItem, - ReconciliationOutcome, + AgreementFeeRate, AgreementRegistry, CancelKind, CancellingAgreement, IndexingAgreement, + NewAgreementParams, PendingAcceptedEvent, PendingExpiredEvent, PendingTerminatedEvent, + ReconciliationItem, ReconciliationOutcome, }, result::Result, }; @@ -133,6 +133,27 @@ pub trait StubAgreementRegistry: Send + Sync { unimplemented!("mark_indexing_agreement_as_canceled_by_requester") } + async fn mark_indexing_agreement_as_cancelling(&self, _id: &IndexingAgreementId) -> Result<()> { + unimplemented!("mark_indexing_agreement_as_cancelling") + } + + async fn get_cancelling_agreements( + &self, + _batch_size: i64, + _max_attempts: u32, + _min_age_minutes: i32, + ) -> Result> { + Ok(Vec::new()) + } + + async fn record_cancel_check( + &self, + _id: &IndexingAgreementId, + failed_attempts: u32, + ) -> Result { + Ok(failed_attempts) + } + async fn apply_reconciliation( &self, _id: &IndexingAgreementId, @@ -415,6 +436,33 @@ impl AgreementRegistry for T { StubAgreementRegistry::mark_indexing_agreement_as_canceled_by_requester(self, id).await } + async fn mark_indexing_agreement_as_cancelling(&self, id: &IndexingAgreementId) -> Result<()> { + StubAgreementRegistry::mark_indexing_agreement_as_cancelling(self, id).await + } + + async fn get_cancelling_agreements( + &self, + batch_size: i64, + max_attempts: u32, + min_age_minutes: i32, + ) -> Result> { + StubAgreementRegistry::get_cancelling_agreements( + self, + batch_size, + max_attempts, + min_age_minutes, + ) + .await + } + + async fn record_cancel_check( + &self, + id: &IndexingAgreementId, + failed_attempts: u32, + ) -> Result { + StubAgreementRegistry::record_cancel_check(self, id, failed_attempts).await + } + async fn apply_reconciliation( &self, id: &IndexingAgreementId, 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 c135defb..46ae9be2 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 @@ -1,6 +1,6 @@ -//! Cancel on-chain, via the RecurringAgreementManager, an agreement dipper doesn't -//! want that is still live: one the indexer rejected off-chain, or one dipper had already -//! cancelled. The chain listener queues it, and a reassessment does when its own cancel fails. +//! Cancel on-chain, via the RecurringAgreementManager, an agreement dipper doesn't want +//! that was accepted anyway: one the indexer rejected off-chain, or one dipper had already +//! cancelled. The chain listener queues it. use std::{ collections::HashSet, @@ -11,7 +11,7 @@ use std::{ use dipper_core::ids::IndexingAgreementId; use crate::{ - cancel_dispatch::cancel_agreement_on_chain, + cancel_dispatch::{LiveCancel, cancel_agreement_on_chain, cancel_if_live}, chain_client::{ChainClient, ChainClientError}, config::IndexingAgreementConfig, registry::{AgreementRegistry, IndexingAgreement, IndexingAgreementStatus}, @@ -148,7 +148,7 @@ where } /// Agreements a job in this process is cancelling right now. The listener can -/// queue one twice before the chain shows it ended, and two jobs at once would +/// queue one twice before the chain shows it ended, and 2 jobs at once would /// both cancel it and both alert. static CANCELLING: LazyLock>> = LazyLock::new(Mutex::default); @@ -173,9 +173,9 @@ impl Drop for Cancelling { } } -/// Cancel on-chain an agreement dipper had already cancelled that is still live: accepted -/// anyway, or an offer a failed cancel left open. The row is already terminal. The chain is -/// read first, so a stale snapshot of one dipper has since ended raises no alert. +/// Cancel on-chain an agreement dipper had already cancelled that the indexer accepted +/// anyway; the row is already terminal. The chain is read first, so a stale snapshot of +/// one dipper has since ended raises no alert. async fn cancel_live_agreement_dipper_cancelled( ctx: &Ctx, agreement: &IndexingAgreement, @@ -190,25 +190,32 @@ where ); return Ok(()); }; - if !live_on_chain(&ctx.chain_client, agreement).await? { - tracing::info!( - agreement_id = %agreement.id, - "Cancelled agreement is no longer live on-chain; nothing to cancel" - ); - return Ok(()); - } - - match cancel_agreement_on_chain(&ctx.chain_client, agreement, &ctx.agreement_conf).await { - Ok(tx_hash) => { + match cancel_if_live(&ctx.chain_client, agreement, &ctx.agreement_conf).await { + LiveCancel::NotLive => { + tracing::info!( + agreement_id = %agreement.id, + "Cancelled agreement is no longer live on-chain; nothing to cancel" + ); + Ok(()) + } + LiveCancel::ReadFailed(err) => { + tracing::warn!( + agreement_id = %agreement.id, + error = %err, + "Failed to read whether a cancelled agreement is live on-chain, will retry" + ); + Err(JobError::Retryable(err.into(), Duration::from_secs(30))) + } + LiveCancel::Ended(tx_hash) => { let tx = tx_hash.map_or_else(|| "none".to_owned(), |hash| hash.to_string()); log_caught_live_agreement(agreement, "cancelled", &tx); Ok(()) } - Err(err @ ChainClientError::MissingTermsVersionHash { .. }) => { + LiveCancel::CancelFailed(err @ ChainClientError::MissingTermsVersionHash { .. }) => { log_caught_live_agreement(agreement, "cancel_impossible", &err.to_string()); Err(JobError::Fatal(err.into())) } - Err(err) => { + LiveCancel::CancelFailed(err) => { log_caught_live_agreement(agreement, "cancel_failed", &err.to_string()); Err(JobError::Retryable(err.into(), Duration::from_secs(30))) } @@ -230,24 +237,6 @@ fn log_caught_live_agreement(agreement: &IndexingAgreement, outcome: &str, detai ); } -/// Read the chain for whether the agreement is still live; a failed read retries. -async fn live_on_chain( - chain_client: &T, - agreement: &IndexingAgreement, -) -> JobResult { - chain_client - .agreement_still_active(agreement.id.as_bytes()) - .await - .map_err(|err| { - tracing::warn!( - agreement_id = %agreement.id, - error = %err, - "Failed to read whether a cancelled agreement is live on-chain, will retry" - ); - JobError::Retryable(err.into(), Duration::from_secs(30)) - }) -} - /// Flip the row to CanceledByRequester once the chain shows it cancelled. A failure /// is logged, not fatal: the chain is already right and the listener retries the /// DB update. Returns whether the row is now terminal, gating the cancel audit. @@ -474,6 +463,30 @@ mod tests { Ok(()) } + async fn mark_indexing_agreement_as_cancelling( + &self, + _id: &IndexingAgreementId, + ) -> crate::registry::Result<()> { + Ok(()) + } + + async fn get_cancelling_agreements( + &self, + _batch_size: i64, + _max_attempts: u32, + _min_age_minutes: i32, + ) -> crate::registry::Result> { + Ok(Vec::new()) + } + + async fn record_cancel_check( + &self, + _id: &IndexingAgreementId, + failed_attempts: u32, + ) -> crate::registry::Result { + Ok(failed_attempts) + } + async fn apply_reconciliation( &self, _id: &IndexingAgreementId, @@ -629,29 +642,7 @@ mod tests { } fn test_agreement_conf() -> Arc { - Arc::new(IndexingAgreementConfig { - data_service: Address::ZERO, - recurring_collector: Address::ZERO, - recurring_agreement_manager: Address::ZERO, - max_agreement_grt_per_30_days: 0.0, - max_seconds_per_collection: 0, - min_seconds_per_collection: 0, - duration_seconds: 0, - deadline_seconds: 0, - max_grt_per_30_days: std::collections::BTreeMap::new(), - max_grt_per_billion_entities_per_30_days: 0.0, - declined_indexer_lookback_days: 0, - price_rejection_lookback_days: 0, - transient_rejection_lookback_minutes: 0, - uncertain_rejection_lookback_days: 0, - unresponsive_indexer_lookback_days: 0, - mass_unresponsive_trip_fraction: 0.5, - mass_unresponsive_reset_fraction: 0.25, - dips_accepting_snapshot_max_age_hours: 48, - dips_accepting_cache_ttl_seconds: 300, - max_in_flight_offers_per_indexer: None, - max_in_flight_offers_total: None, - }) + Arc::new(IndexingAgreementConfig::for_tests()) } fn test_deployment_id() -> DeploymentId { @@ -848,7 +839,7 @@ mod tests { #[tokio::test] async fn a_second_job_for_the_same_agreement_leaves_it_to_the_first() { // The listener can queue an agreement twice before the chain shows it - // ended; two jobs at once would both cancel it and both alert. + // ended; 2 jobs at once would both cancel it and both alert. let agreement = make_agreement(IndexingAgreementStatus::CanceledByRequester); let agreement_id = agreement.id; let chain = MockChainClient::live(); 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 0fb5bfa7..340bf26b 100644 --- a/bin/dipper-service/src/worker/handlers/reassess_indexing_request.rs +++ b/bin/dipper-service/src/worker/handlers/reassess_indexing_request.rs @@ -634,122 +634,22 @@ where } // Cancel old agreements that have no replacement to pair with: these indexers - // leave the target group with nothing taking their place. An accepted one is - // cancelled on-chain before it is marked, so a failed cancel leaves it for a retry. + // leave the target group with nothing taking their place. let mut directly_cancelled = 0u32; let mut cancel_failures = 0u32; for old_agreement in old_iter { - let was_accepted = matches!( - old_agreement.status, - crate::registry::IndexingAgreementStatus::AcceptedOnChain - ); - // A Created agreement's offer may already be on-chain, so it is cancelled - // on-chain too; without a stored terms hash it can only be marked locally. - let needs_on_chain_cancel = was_accepted - || (old_agreement.status == crate::registry::IndexingAgreementStatus::Created - && old_agreement.terms_version_hash.is_some()); - // An unaccepted one is marked first. An offer in flight then lands before the - // cancel, which withdraws it, or after, when its job sees the mark and withdraws it. - if !was_accepted && !mark_unpaired_cancelled(&ctx.registry, old_agreement).await { - cancel_failures += 1; - continue; - } - - let mut on_chain_cancel_tx: Option = None; - if needs_on_chain_cancel { - match crate::cancel_dispatch::cancel_agreement_on_chain( - &ctx.chain_client, - old_agreement, - &ctx.agreement_conf, - ) - .await - { - Ok(Some(tx_hash)) => { - tracing::info!( - agreement_id = %old_agreement.id, - indexing_request_id = %indexing_request_id, - %tx_hash, - "Submitted on-chain cancellation for unpaired old agreement" - ); - on_chain_cancel_tx = Some(tx_hash.to_string()); - } - Ok(None) => { - tracing::info!( - agreement_id = %old_agreement.id, - indexing_request_id = %indexing_request_id, - "Unpaired old agreement already canceled on-chain; proceeding with local cleanup" - ); - } - Err(err) => { - tracing::warn!( - error = %err, - agreement_id = %old_agreement.id, - was_accepted, - "On-chain cancel failed; an accepted agreement is retried later, an \ - unaccepted one by a queued cancel job" - ); - cancel_failures += 1; - if was_accepted { - continue; - } - if let Err(err) = ctx - .queue - .cancel_rejected_agreement_on_chain( - old_agreement.id, - JobPriority::Background, - ) - .await - { - tracing::error!( - error = %err, - agreement_id = %old_agreement.id, - "Failed to queue a retry of the on-chain cancel; an offer already \ - on-chain stays open" - ); - } - } - } - } - - if was_accepted && !mark_unpaired_cancelled(&ctx.registry, old_agreement).await { + let Some(new_status) = cancel_unpaired(&ctx, old_agreement).await else { cancel_failures += 1; continue; - } - + }; tracing::info!( agreement_id = %old_agreement.id, indexing_request_id = %indexing_request_id, old_status = %old_agreement.status, - new_status = "CANCELED_BY_REQUESTER", + new_status, reason = "reassessment_not_in_target_group", "agreement state transition" ); - - // Record the cancel audit for the accepted-on-chain agreements dipper just - // cancelled, so the chain_listener's `terminated` sweep announces them - // durably. Never-accepted agreements (`!was_accepted`) were never live - // on-chain: they are not sweep-eligible (`accepted_at IS NULL`) and - // correctly emit nothing. - if was_accepted { - let manager = ctx.agreement_conf.recurring_agreement_manager().to_string(); - if let Err(err) = ctx - .registry - .record_cancel_audit( - &old_agreement.id, - now_secs(), - &manager, - on_chain_cancel_tx.as_deref(), - ) - .await - { - tracing::warn!( - agreement_id = %old_agreement.id, - error = %err, - "failed to record cancel audit; terminated event may emit with fallback fields" - ); - } - } - directly_cancelled += 1; } @@ -762,16 +662,15 @@ where } if cancel_failures > 0 { - // An unaccepted agreement's failed cancel was queued as its own job above. An - // accepted one stays AcceptedOnChain: the orphan-cancel sweep retries it once the - // request is Canceled, the periodic reassignment service (24 h) while it is Open. + // An agreement that couldn't be marked keeps its status: the orphan-cancel sweep + // retries it once the request is Canceled, the periodic reassignment service + // (24 h) while it is Open. tracing::warn!( indexing_request_id=%indexing_request_id, failures=cancel_failures, - "some agreement cancels failed during reassessment; retry will fire \ - via a queued cancel job (unaccepted agreements), the orphan-cancel \ - sweep (Canceled requests) or the periodic reassignment service \ - (Open requests over-target)" + "some agreements could not be marked for cancelling during reassessment; retry \ + will fire via the orphan-cancel sweep (Canceled requests) or the periodic \ + reassignment service (Open requests over-target)" ); } @@ -1062,8 +961,8 @@ mod lifecycle_event_tests { #[derive(Default, Clone)] struct MockChainClient { cancelled: Arc>>, - /// When set, each cancel records whether its agreement was already marked. - marked_cancelled: Option>>>, + /// When set, each cancel records whether its agreement was already marked cancelling. + marked_cancelling: Option>>>, marked_at_cancel: Arc>>, fail_cancel: bool, } @@ -1094,7 +993,7 @@ mod lifecycle_event_tests { return Err(ChainClientError::RpcError(anyhow::anyhow!("rpc down"))); } self.cancelled.lock().unwrap().push(*agreement_id); - if let Some(marked) = &self.marked_cancelled { + if let Some(marked) = &self.marked_cancelling { let id = IndexingAgreementId::from_bytes(*agreement_id); let was_marked = marked.lock().unwrap().contains(&id); self.marked_at_cancel.lock().unwrap().push(was_marked); @@ -1144,6 +1043,7 @@ mod lifecycle_event_tests { shortfall_active: std::sync::Mutex, /// Ids marked CanceledByRequester locally. marked_cancelled: Arc>>, + marked_cancelling: Arc>>, } #[async_trait] @@ -1295,6 +1195,30 @@ mod lifecycle_event_tests { self.marked_cancelled.lock().unwrap().push(*id); Ok(()) } + async fn mark_indexing_agreement_as_cancelling( + &self, + id: &IndexingAgreementId, + ) -> RegistryResult<()> { + self.marked_cancelling.lock().unwrap().push(*id); + Ok(()) + } + + async fn get_cancelling_agreements( + &self, + _batch_size: i64, + _max_attempts: u32, + _min_age_minutes: i32, + ) -> RegistryResult> { + Ok(Vec::new()) + } + + async fn record_cancel_check( + &self, + _id: &IndexingAgreementId, + failed_attempts: u32, + ) -> RegistryResult { + Ok(failed_attempts) + } async fn apply_reconciliation( &self, _id: &IndexingAgreementId, @@ -1449,19 +1373,11 @@ mod lifecycle_event_tests { min_seconds_per_collection: 60, duration_seconds: 86400, deadline_seconds: 3600, - max_grt_per_30_days: std::collections::BTreeMap::new(), - max_grt_per_billion_entities_per_30_days: 0.0, declined_indexer_lookback_days: 30, price_rejection_lookback_days: 1, transient_rejection_lookback_minutes: 30, uncertain_rejection_lookback_days: 1, - unresponsive_indexer_lookback_days: 0, - mass_unresponsive_trip_fraction: 0.5, - mass_unresponsive_reset_fraction: 0.25, - dips_accepting_snapshot_max_age_hours: 48, - dips_accepting_cache_ttl_seconds: 300, - max_in_flight_offers_per_indexer: None, - max_in_flight_offers_total: None, + ..IndexingAgreementConfig::for_tests() } } @@ -1700,6 +1616,7 @@ mod lifecycle_event_tests { // Already in shortfall. shortfall_active: std::sync::Mutex::new(true), marked_cancelled: Arc::default(), + marked_cancelling: Arc::default(), }; let ctx = build_ctx( registry, @@ -1737,6 +1654,7 @@ mod lifecycle_event_tests { chain_state_lookup_fails: false, shortfall_active: std::sync::Mutex::new(false), marked_cancelled: Arc::default(), + marked_cancelling: Arc::default(), }; let ctx = build_ctx( registry, @@ -1808,6 +1726,7 @@ mod lifecycle_event_tests { chain_state_lookup_fails: false, shortfall_active: std::sync::Mutex::new(false), marked_cancelled: Arc::default(), + marked_cancelling: Arc::default(), }; let ctx = build_ctx( registry, @@ -1862,6 +1781,7 @@ mod lifecycle_event_tests { chain_state_lookup_fails: false, shortfall_active: std::sync::Mutex::new(false), marked_cancelled: Arc::default(), + marked_cancelling: Arc::default(), }; let ctx = build_ctx( registry, @@ -1927,6 +1847,7 @@ mod lifecycle_event_tests { chain_state_lookup_fails: false, shortfall_active: std::sync::Mutex::new(false), marked_cancelled: Arc::default(), + marked_cancelling: Arc::default(), }; let ctx = build_ctx( registry, @@ -1950,8 +1871,8 @@ mod lifecycle_event_tests { assert!(matches!(captured[0], CapturedEvent::Proposed { .. })); } - /// Ctx whose only active agreement, in `status`, leaves the target group, with - /// the chain mock recording whether the row was already marked at each cancel. + /// Ctx whose only active agreement, in `status`, leaves the target group, with the + /// chain mock recording whether the row was already marked cancelling at each cancel. fn ctx_cancelling_one( status: IndexingAgreementStatus, ) -> ( @@ -1970,16 +1891,18 @@ mod lifecycle_event_tests { CapturingEventsProducer::new(), indexer_urls::Snapshot::new(), ); - ctx.chain_client.marked_cancelled = Some(ctx.registry.marked_cancelled.clone()); + ctx.chain_client.marked_cancelling = Some(ctx.registry.marked_cancelling.clone()); (ctx, leaving) } #[tokio::test] - async fn marks_an_unaccepted_agreement_cancelled_before_cancelling_it_on_chain() { + async fn marks_an_unaccepted_agreement_cancelling_before_cancelling_it_on_chain() { // Its offer may be in flight. An offer landing after the cancel can't be - // withdrawn by it, so its job must find the agreement already marked. + // withdrawn by it, so its job must find the agreement already marked. It + // stays cancelling: the offer could still be accepted until its deadline. let (ctx, leaving) = ctx_cancelling_one(IndexingAgreementStatus::Created); let chain = ctx.chain_client.clone(); + let cancelled = ctx.registry.marked_cancelled.clone(); let result = handle(ctx, &test_message(0)).await; @@ -1989,50 +1912,43 @@ mod lifecycle_event_tests { vec![*leaving.id.as_bytes()] ); assert_eq!(*chain.marked_at_cancel.lock().unwrap(), vec![true]); + assert!(cancelled.lock().unwrap().is_empty()); } #[tokio::test] - async fn cancels_an_accepted_agreement_on_chain_before_marking_it() { + async fn marks_an_accepted_agreement_cancelled_once_its_cancel_lands() { let (ctx, leaving) = ctx_cancelling_one(IndexingAgreementStatus::AcceptedOnChain); let chain = ctx.chain_client.clone(); - let marked = ctx.registry.marked_cancelled.clone(); + let cancelled = ctx.registry.marked_cancelled.clone(); let result = handle(ctx, &test_message(0)).await; assert!(result.is_ok(), "got {result:?}"); - assert_eq!(*chain.marked_at_cancel.lock().unwrap(), vec![false]); - assert_eq!(*marked.lock().unwrap(), vec![leaving.id]); - } - - #[tokio::test] - async fn an_unaccepted_agreement_whose_cancel_fails_is_marked_and_its_cancel_queued() { - // Nothing else would send its cancel again: the row is already terminal, - // and its offer job may have finished before the mark. - let (mut ctx, leaving) = ctx_cancelling_one(IndexingAgreementStatus::Created); - ctx.chain_client.fail_cancel = true; - let marked = ctx.registry.marked_cancelled.clone(); - let queue = ctx.queue.clone(); - - let result = handle(ctx, &test_message(0)).await; - - assert!(result.is_ok(), "got {result:?}"); - assert_eq!(*marked.lock().unwrap(), vec![leaving.id]); - assert_eq!(*queue.cancels_queued.lock().unwrap(), vec![leaving.id]); + assert_eq!(*chain.marked_at_cancel.lock().unwrap(), vec![true]); + assert_eq!(*cancelled.lock().unwrap(), vec![leaving.id]); } #[tokio::test] - async fn an_accepted_agreement_whose_cancel_fails_stays_for_a_retry() { - // Marking it cancelled would leave it live with nothing to end it. - let (mut ctx, _leaving) = ctx_cancelling_one(IndexingAgreementStatus::AcceptedOnChain); - ctx.chain_client.fail_cancel = true; - let marked = ctx.registry.marked_cancelled.clone(); - let queue = ctx.queue.clone(); + async fn an_agreement_whose_cancel_fails_is_left_cancelling_for_the_listener() { + // The chain listener retries the cancel of every cancelling agreement until + // the chain shows it ended, so nothing is queued here. + for status in [ + IndexingAgreementStatus::Created, + IndexingAgreementStatus::AcceptedOnChain, + ] { + let (mut ctx, leaving) = ctx_cancelling_one(status); + ctx.chain_client.fail_cancel = true; + let cancelling = ctx.registry.marked_cancelling.clone(); + let cancelled = ctx.registry.marked_cancelled.clone(); + let queue = ctx.queue.clone(); - let result = handle(ctx, &test_message(0)).await; + let result = handle(ctx, &test_message(0)).await; - assert!(result.is_ok(), "got {result:?}"); - assert!(marked.lock().unwrap().is_empty()); - assert!(queue.cancels_queued.lock().unwrap().is_empty()); + assert!(result.is_ok(), "got {result:?}"); + assert_eq!(*cancelling.lock().unwrap(), vec![leaving.id]); + assert!(cancelled.lock().unwrap().is_empty()); + assert!(queue.cancels_queued.lock().unwrap().is_empty()); + } } #[tokio::test] @@ -2049,6 +1965,7 @@ mod lifecycle_event_tests { chain_state_lookup_fails: false, shortfall_active: std::sync::Mutex::new(false), marked_cancelled: Arc::default(), + marked_cancelling: Arc::default(), }; let ctx = build_ctx( registry, @@ -2065,6 +1982,46 @@ mod lifecycle_event_tests { } } +/// Take an old agreement out of the target group, returning its new status, or `None` +/// (logged) when it couldn't be marked. One that may be live on-chain is cancelled there +/// too, which the chain listener retries until it ends. +async fn cancel_unpaired( + ctx: &Ctx, + agreement: &crate::registry::IndexingAgreement, +) -> Option<&'static str> +where + R: AgreementRegistry + Sync, + T: ChainClient, +{ + let may_be_live = agreement.status == crate::registry::IndexingAgreementStatus::AcceptedOnChain + || (agreement.status == crate::registry::IndexingAgreementStatus::Created + && agreement.terms_version_hash.is_some()); + if !may_be_live { + return mark_unpaired_cancelled(&ctx.registry, agreement) + .await + .then_some("CANCELED_BY_REQUESTER"); + } + match crate::cancel_dispatch::start_cancel( + &ctx.registry, + &ctx.chain_client, + agreement, + &ctx.agreement_conf, + ) + .await + { + Ok(crate::cancel_dispatch::CancelStarted::Ended) => Some("CANCELED_BY_REQUESTER"), + Ok(crate::cancel_dispatch::CancelStarted::Cancelling) => Some("CANCELLING"), + Err(err) => { + tracing::error!( + error = %err, + agreement_id = %agreement.id, + "Failed to mark unpaired old agreement as cancelling in local DB" + ); + None + } + } +} + /// Mark an unpaired old agreement CanceledByRequester; false, logged, if that fails. async fn mark_unpaired_cancelled( registry: &R, diff --git a/bin/dipper-service/src/worker/handlers/send_indexing_agreement_proposal.rs b/bin/dipper-service/src/worker/handlers/send_indexing_agreement_proposal.rs index 1aa582ac..af9cbf83 100644 --- a/bin/dipper-service/src/worker/handlers/send_indexing_agreement_proposal.rs +++ b/bin/dipper-service/src/worker/handlers/send_indexing_agreement_proposal.rs @@ -522,6 +522,30 @@ mod tests { Ok(()) } + async fn mark_indexing_agreement_as_cancelling( + &self, + _id: &IndexingAgreementId, + ) -> crate::registry::Result<()> { + Ok(()) + } + + async fn get_cancelling_agreements( + &self, + _batch_size: i64, + _max_attempts: u32, + _min_age_minutes: i32, + ) -> crate::registry::Result> { + Ok(Vec::new()) + } + + async fn record_cancel_check( + &self, + _id: &IndexingAgreementId, + failed_attempts: u32, + ) -> crate::registry::Result { + Ok(failed_attempts) + } + async fn apply_reconciliation( &self, _id: &IndexingAgreementId, diff --git a/bin/dipper-service/src/worker/handlers/submit_offer.rs b/bin/dipper-service/src/worker/handlers/submit_offer.rs index 172ddcdd..e27f3d74 100644 --- a/bin/dipper-service/src/worker/handlers/submit_offer.rs +++ b/bin/dipper-service/src/worker/handlers/submit_offer.rs @@ -23,7 +23,7 @@ use thegraph_core::{DeploymentId, alloy::primitives::ChainId}; use url::Url; use crate::{ - cancel_dispatch::cancel_agreement_on_chain, + cancel_dispatch::{LiveCancel, cancel_if_live}, chain_client::{ChainClient, ChainClientError, decode_revert_reason}, config::IndexingAgreementConfig, indexer_rpc_client::into_sol_rca, @@ -213,9 +213,7 @@ async fn next_step( NextStep::Skip } Some(a) if a.status == IndexingAgreementStatus::Created => NextStep::Offer(a), - Some(a) if a.status == IndexingAgreementStatus::CanceledByRequester => { - NextStep::Withdraw(a) - } + Some(a) if dipper_cancelled(a.status) => NextStep::Withdraw(a), Some(a) => { tracing::warn!( agreement_id = %agreement_id, @@ -234,16 +232,9 @@ async fn withdraw_offer_if_stored( ctx: &Ctx, agreement: &IndexingAgreement, ) -> JobResult<()> { - let stored = ctx - .chain_client - .agreement_still_active(agreement.id.as_bytes()) - .await - .map_err(|err| retry_withdraw(agreement, err))?; - if !stored { - return Ok(()); - } - match cancel_agreement_on_chain(&ctx.chain_client, agreement, &ctx.agreement_conf).await { - Ok(tx_hash) => { + match cancel_if_live(&ctx.chain_client, agreement, &ctx.agreement_conf).await { + LiveCancel::NotLive => Ok(()), + LiveCancel::Ended(tx_hash) => { tracing::info!( agreement_id = %agreement.id, tx_hash = ?tx_hash, @@ -251,7 +242,7 @@ async fn withdraw_offer_if_stored( ); Ok(()) } - Err(err @ ChainClientError::MissingTermsVersionHash { .. }) => { + LiveCancel::CancelFailed(err @ ChainClientError::MissingTermsVersionHash { .. }) => { tracing::error!( agreement_id = %agreement.id, error = %err, @@ -259,7 +250,9 @@ async fn withdraw_offer_if_stored( ); Err(JobError::Fatal(err.into())) } - Err(err) => Err(retry_withdraw(agreement, err)), + LiveCancel::ReadFailed(err) | LiveCancel::CancelFailed(err) => { + Err(retry_withdraw(agreement, err)) + } } } @@ -272,22 +265,30 @@ fn retry_withdraw(agreement: &IndexingAgreement, err: ChainClientError) -> JobEr JobError::Retryable(err.into(), TRANSIENT_RETRY_BASE) } -/// The agreement, if it was cancelled after this job's status check. +/// Whether dipper has cancelled the agreement, or started to. +fn dipper_cancelled(status: IndexingAgreementStatus) -> bool { + matches!( + status, + IndexingAgreementStatus::Cancelling | IndexingAgreementStatus::CanceledByRequester + ) +} + +/// The agreement, if it was cancelled after this job's status check. After a failed read +/// the job still finishes: retrying would send the offer again, and the chain listener's +/// cancel retry withdraws the offer of an agreement left cancelling. async fn cancelled_meanwhile( registry: &R, agreement_id: &IndexingAgreementId, ) -> Option { match registry.get_indexing_agreement_by_id(agreement_id).await { - Ok(Some(agreement)) if agreement.status == IndexingAgreementStatus::CanceledByRequester => { - Some(agreement) - } + Ok(Some(agreement)) if dipper_cancelled(agreement.status) => Some(agreement), Ok(_) => None, Err(err) => { tracing::warn!( agreement_id = %agreement_id, error = %err, - "Failed to re-read agreement after its offer landed; a cancel made meanwhile \ - would leave the offer open until its deadline" + "Failed to re-read agreement after its offer landed; the cancel retry \ + withdraws the offer if it was cancelled" ); None } @@ -298,7 +299,7 @@ async fn cancelled_meanwhile( mod tests { use std::sync::{ Arc, Mutex, - atomic::{AtomicBool, Ordering}, + atomic::{AtomicBool, AtomicU32, Ordering}, }; use async_trait::async_trait; @@ -321,6 +322,9 @@ mod tests { struct MockRegistry { agreement: SharedAgreement, + /// When set, every read after the first fails. + later_reads_fail: bool, + reads: AtomicU32, } #[async_trait] @@ -329,6 +333,9 @@ mod tests { &self, _id: &IndexingAgreementId, ) -> crate::registry::Result> { + if self.reads.fetch_add(1, Ordering::SeqCst) > 0 && self.later_reads_fail { + return Err(crate::registry::Error::NoRecordsUpdated); + } Ok(self.agreement.lock().unwrap().clone()) } async fn update_offer_tx_hash( @@ -364,7 +371,7 @@ mod tests { if let Some(agreement) = &self.cancel_mid_send && let Some(row) = agreement.lock().unwrap().as_mut() { - row.status = IndexingAgreementStatus::CanceledByRequester; + row.status = IndexingAgreementStatus::Cancelling; } let result = self .offer_result @@ -489,6 +496,8 @@ mod tests { Ctx { registry: MockRegistry { agreement: Arc::new(Mutex::new(Some(agreement))), + later_reads_fail: false, + reads: AtomicU32::new(0), }, chain_client: MockChainClient { offer_result: Mutex::new(offer_result), @@ -502,29 +511,7 @@ mod tests { } fn test_agreement_conf() -> crate::config::IndexingAgreementConfig { - crate::config::IndexingAgreementConfig { - data_service: Address::ZERO, - recurring_collector: Address::ZERO, - recurring_agreement_manager: Address::ZERO, - max_agreement_grt_per_30_days: 0.0, - max_seconds_per_collection: 0, - min_seconds_per_collection: 0, - duration_seconds: 0, - deadline_seconds: 0, - max_grt_per_30_days: std::collections::BTreeMap::new(), - max_grt_per_billion_entities_per_30_days: 0.0, - declined_indexer_lookback_days: 0, - price_rejection_lookback_days: 0, - transient_rejection_lookback_minutes: 0, - uncertain_rejection_lookback_days: 0, - unresponsive_indexer_lookback_days: 0, - mass_unresponsive_trip_fraction: 0.5, - mass_unresponsive_reset_fraction: 0.25, - dips_accepting_snapshot_max_age_hours: 48, - dips_accepting_cache_ttl_seconds: 300, - max_in_flight_offers_per_indexer: None, - max_in_flight_offers_total: None, - } + crate::config::IndexingAgreementConfig::for_tests() } #[tokio::test] @@ -547,6 +534,21 @@ mod tests { assert_eq!(*cancelled.lock().unwrap(), vec![*agreement_id.as_bytes()]); } + #[tokio::test] + async fn finishes_without_resending_when_it_cannot_recheck_the_agreement() { + //* Arrange - a retry would send the offer again; the mock panics on a second send + let agreement = make_test_agreement(); + let message = make_message(agreement.id); + let mut ctx = ctx_with_offer_result(agreement, Ok(Some(B256::repeat_byte(0xab)))); + ctx.registry.later_reads_fail = true; + + //* Act + let result = handle(ctx, &message).await; + + //* Assert + assert!(result.is_ok(), "got {result:?}"); + } + #[tokio::test] async fn keeps_its_offer_when_the_agreement_is_still_wanted() { //* Arrange @@ -583,23 +585,28 @@ mod tests { #[tokio::test] async fn withdraws_a_stored_offer_of_an_agreement_already_cancelled() { - //* Arrange - an earlier attempt sent the offer, then the agreement was - // cancelled before this retry; no offer result, so a send would panic - let mut agreement = make_test_agreement(); - agreement.status = IndexingAgreementStatus::CanceledByRequester; - agreement.terms_version_hash = Some(vec![7u8; 32]); - let agreement_id = agreement.id; - let message = make_message(agreement_id); - let ctx = ctx_with(agreement, None); - ctx.chain_client.on_chain.store(true, Ordering::SeqCst); - let cancelled = ctx.chain_client.cancelled.clone(); - - //* Act - let result = handle(ctx, &message).await; - - //* Assert - assert!(result.is_ok(), "got {result:?}"); - assert_eq!(*cancelled.lock().unwrap(), vec![*agreement_id.as_bytes()]); + for status in [ + IndexingAgreementStatus::Cancelling, + IndexingAgreementStatus::CanceledByRequester, + ] { + //* Arrange - an earlier attempt sent the offer, then the agreement was + // cancelled before this retry; no offer result, so a send would panic + let mut agreement = make_test_agreement(); + agreement.status = status; + agreement.terms_version_hash = Some(vec![7u8; 32]); + let agreement_id = agreement.id; + let message = make_message(agreement_id); + let ctx = ctx_with(agreement, None); + ctx.chain_client.on_chain.store(true, Ordering::SeqCst); + let cancelled = ctx.chain_client.cancelled.clone(); + + //* Act + let result = handle(ctx, &message).await; + + //* Assert + assert!(result.is_ok(), "got {result:?}"); + assert_eq!(*cancelled.lock().unwrap(), vec![*agreement_id.as_bytes()]); + } } #[tokio::test] diff --git a/dipper-pgregistry/migrations/20261002000000_add_cancelling_status.sql b/dipper-pgregistry/migrations/20261002000000_add_cancelling_status.sql new file mode 100644 index 00000000..e9c86d29 --- /dev/null +++ b/dipper-pgregistry/migrations/20261002000000_add_cancelling_status.sql @@ -0,0 +1,15 @@ +-- Cancelling (status = 9): dipper has decided to end the agreement and keeps sending the +-- on-chain cancel until the chain confirms it ended; only then does the row become +-- CanceledByRequester, which announces the end. +-- +-- cancel_attempts counts cancels that reached the chain without ending the agreement, +-- so one that can never work stops being retried and is left for an operator. +-- cancel_checked_at lets each retry sweep start with the agreements checked longest ago. +ALTER TABLE dipper_reg_indexing_agreements + ADD COLUMN cancel_attempts INTEGER NOT NULL DEFAULT 0, + ADD COLUMN cancel_checked_at TIMESTAMPTZ; + +-- Partial, so it only covers agreements still being cancelled (normally none). +CREATE INDEX idx_indexing_agreements_cancelling + ON dipper_reg_indexing_agreements (cancel_checked_at NULLS FIRST, updated_at) + WHERE status = 9; diff --git a/dipper-pgregistry/src/indexing_agreement.rs b/dipper-pgregistry/src/indexing_agreement.rs index 53c20c32..4a3b24fe 100644 --- a/dipper-pgregistry/src/indexing_agreement.rs +++ b/dipper-pgregistry/src/indexing_agreement.rs @@ -206,6 +206,11 @@ pub enum Status { /// This is a terminal state. AbandonedByIndexer = 8, + /// Dipper decided to end the agreement and is cancelling it on-chain, where it may + /// still be live. It becomes `CanceledByRequester`, announced as ended, only once + /// the chain confirms the end. + Cancelling = 9, + /// A fallback for unknown status values. Unknown = i32::MAX, } @@ -221,6 +226,7 @@ impl std::fmt::Display for Status { Status::AcceptedOnChain => "ACCEPTED_ON_CHAIN", Status::Rejected => "REJECTED", Status::AbandonedByIndexer => "ABANDONED_BY_INDEXER", + Status::Cancelling => "CANCELLING", Status::Unknown => "UNKNOWN", }; f.write_str(status) diff --git a/dipper-pgregistry/src/lib.rs b/dipper-pgregistry/src/lib.rs index 20b40ea6..8e0ea869 100644 --- a/dipper-pgregistry/src/lib.rs +++ b/dipper-pgregistry/src/lib.rs @@ -20,9 +20,9 @@ pub use indexing_request::{ Status as IndexingRequestStatus, }; pub use postgres::{ - CancelKind, ChainListenerStateRow, NewAgreementParams, PendingAcceptedEvent, - PendingExpiredEvent, PendingTerminatedEvent, PgRegistry, ReconciliationAudit, - ReconciliationItem, ReconciliationOutcome, + CancelKind, CancellingAgreement, ChainListenerStateRow, NewAgreementParams, + PendingAcceptedEvent, PendingExpiredEvent, PendingTerminatedEvent, PgRegistry, + ReconciliationAudit, ReconciliationItem, ReconciliationOutcome, }; pub use result::{Error, Result}; diff --git a/dipper-pgregistry/src/postgres.rs b/dipper-pgregistry/src/postgres.rs index 14c85160..d40439b3 100644 --- a/dipper-pgregistry/src/postgres.rs +++ b/dipper-pgregistry/src/postgres.rs @@ -187,6 +187,39 @@ impl sqlx::FromRow<'_, sqlx::postgres::PgRow> for PendingAcceptedEvent { } } +/// An agreement dipper is still cancelling on-chain. +#[derive(Debug, Clone)] +pub struct CancellingAgreement { + pub agreement: IndexingAgreement, + /// Whether dipper saw it accepted on-chain, so its end is announced. + pub accepted_on_chain: bool, +} + +impl sqlx::FromRow<'_, sqlx::postgres::PgRow> for CancellingAgreement { + fn from_row(row: &sqlx::postgres::PgRow) -> Result { + use sqlx::Row as _; + let accepted_at: Option = row.try_get("accepted_at")?; + Ok(Self { + agreement: IndexingAgreement::from_row(row)?, + accepted_on_chain: accepted_at.is_some(), + }) + } +} + +/// Statuses an on-chain cancel by dipper ends. +const CANCEL_BY_REQUESTER_FROM: &[IndexingAgreementStatus] = &[ + IndexingAgreementStatus::Created, + IndexingAgreementStatus::AcceptedOnChain, + IndexingAgreementStatus::Rejected, + IndexingAgreementStatus::Cancelling, +]; + +/// Statuses an on-chain cancel by the indexer ends. +const CANCEL_BY_INDEXER_FROM: &[IndexingAgreementStatus] = &[ + IndexingAgreementStatus::AcceptedOnChain, + IndexingAgreementStatus::Cancelling, +]; + /// A row that needs a `request.expired` lifecycle event emitted. Sourced from /// the agreement row alone; `request_expired_at` is the terms deadline (the true /// expiry instant), so the sweep needs no chain-time snapshot. @@ -911,29 +944,117 @@ impl PgRegistry { &self, agreement_id: &IndexingAgreementId, ) -> Result<(), Error> { - let record: Option<(IndexingAgreementId,)> = sqlx::query_as( + self.set_status_from( + agreement_id, + IndexingAgreementStatus::CanceledByRequester, + CANCEL_BY_REQUESTER_FROM, + ) + .await + } + + /// Mark an agreement that may be live on-chain `Cancelling`, before dipper sends its + /// on-chain cancel. One marked `Expired` may have been accepted unseen by a lagging listener. + pub async fn mark_indexing_agreement_as_cancelling( + &self, + agreement_id: &IndexingAgreementId, + ) -> Result<(), Error> { + self.set_status_from( + agreement_id, + IndexingAgreementStatus::Cancelling, + &[ + IndexingAgreementStatus::Created, + IndexingAgreementStatus::AcceptedOnChain, + IndexingAgreementStatus::Rejected, + IndexingAgreementStatus::Expired, + ], + ) + .await + } + + /// `Cancelling` agreements marked over `min_age_minutes` ago whose cancel has failed + /// fewer than `max_attempts` times, those checked longest ago first. + pub async fn get_cancelling_agreements( + &self, + batch_size: i64, + max_attempts: u32, + min_age_minutes: i32, + ) -> Result, Error> { + sqlx::query_as( + r#" + SELECT + id, + nonce_uuid, + created_at, + updated_at, + status, + indexing_request_id, + deployment_id, + indexer_id, + indexer_url, + terms, + last_block_height, + last_progress_at, + rejection_reason, + terms_version_hash, + accepted_at + FROM dipper_reg_indexing_agreements + 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 + LIMIT $3 + "#, + ) + .bind(IndexingAgreementStatus::Cancelling) + .bind(i32::try_from(max_attempts).unwrap_or(i32::MAX)) + .bind(batch_size) + .bind(min_age_minutes) + .fetch_all(&self.pool) + .await + .map_err(Into::into) + } + + /// Record a check of a `Cancelling` agreement that left it cancelling, adding + /// `failed_attempts` to its failed cancels and returning the new count. + pub async fn record_cancel_check( + &self, + agreement_id: &IndexingAgreementId, + failed_attempts: u32, + ) -> Result { + let record: Option<(i32,)> = sqlx::query_as( r#" UPDATE dipper_reg_indexing_agreements SET - status = $1, - updated_at = timezone('UTC', now()) - WHERE id = $2 AND status IN ($3, $4, $5) - RETURNING id + cancel_attempts = LEAST(cancel_attempts::BIGINT + $3, 2147483647)::INTEGER, + cancel_checked_at = timezone('UTC', now()) + WHERE id = $1 AND status = $2 + RETURNING cancel_attempts "#, ) - .bind(IndexingAgreementStatus::CanceledByRequester) .bind(agreement_id) - .bind(IndexingAgreementStatus::Created) - .bind(IndexingAgreementStatus::AcceptedOnChain) - .bind(IndexingAgreementStatus::Rejected) + .bind(IndexingAgreementStatus::Cancelling) + .bind(i64::from(failed_attempts)) .fetch_optional(&self.pool) .await?; + let (attempts,) = record.ok_or(Error::NoRecordsUpdated)?; + Ok(u32::try_from(attempts).unwrap_or_default()) + } - if record.is_none() { - return Err(Error::NoRecordsUpdated); + /// Move an agreement to `new_status` if it is in one of `allowed_from`. + async fn set_status_from( + &self, + agreement_id: &IndexingAgreementId, + new_status: IndexingAgreementStatus, + allowed_from: &[IndexingAgreementStatus], + ) -> Result<(), Error> { + let mut tx = self.pool.begin().await?; + let updated = update_status_from(&mut tx, agreement_id, new_status, allowed_from).await?; + tx.commit().await?; + if updated { + Ok(()) + } else { + Err(Error::NoRecordsUpdated) } - - Ok(()) } /// Atomically apply a reconciliation-driven state transition (accept @@ -986,15 +1107,11 @@ impl PgRegistry { let (new_status, allowed_from): (_, &[IndexingAgreementStatus]) = match kind { CancelKind::ByRequester => ( IndexingAgreementStatus::CanceledByRequester, - &[ - IndexingAgreementStatus::Created, - IndexingAgreementStatus::AcceptedOnChain, - IndexingAgreementStatus::Rejected, - ], + CANCEL_BY_REQUESTER_FROM, ), CancelKind::ByIndexer => ( IndexingAgreementStatus::CanceledByIndexer, - &[IndexingAgreementStatus::AcceptedOnChain], + CANCEL_BY_INDEXER_FROM, ), }; did_cancel = @@ -1084,15 +1201,11 @@ impl PgRegistry { let (new_status, allowed_from): (_, &[IndexingAgreementStatus]) = match cancel_kind { CancelKind::ByRequester => ( IndexingAgreementStatus::CanceledByRequester, - &[ - IndexingAgreementStatus::Created, - IndexingAgreementStatus::AcceptedOnChain, - IndexingAgreementStatus::Rejected, - ], + CANCEL_BY_REQUESTER_FROM, ), CancelKind::ByIndexer => ( IndexingAgreementStatus::CanceledByIndexer, - &[IndexingAgreementStatus::AcceptedOnChain], + CANCEL_BY_INDEXER_FROM, ), }; let did_cancel = @@ -1127,11 +1240,7 @@ impl PgRegistry { &mut tx, &cancel_by_requester, IndexingAgreementStatus::CanceledByRequester, - &[ - IndexingAgreementStatus::Created, - IndexingAgreementStatus::AcceptedOnChain, - IndexingAgreementStatus::Rejected, - ], + CANCEL_BY_REQUESTER_FROM, ) .await? { @@ -1142,7 +1251,7 @@ impl PgRegistry { &mut tx, &cancel_by_indexer, IndexingAgreementStatus::CanceledByIndexer, - &[IndexingAgreementStatus::AcceptedOnChain], + CANCEL_BY_INDEXER_FROM, ) .await? { @@ -1834,8 +1943,8 @@ impl PgRegistry { /// Returns (agreement_id, indexer_id, deployment_id, base_rate_wei, /// entity_rate_wei) per active agreement for optimistic fee estimation. /// - /// Queries all `Created` or `AcceptedOnChain` agreements and extracts - /// both rate fields from the terms metadata. + /// Queries all `Created`, `AcceptedOnChain` or `Cancelling` agreements, the last + /// still paid until their cancel lands, and extracts both rate fields from the terms. pub async fn get_agreement_fee_rates( &self, ) -> Result, Error> { @@ -1847,11 +1956,12 @@ impl PgRegistry { r#" SELECT id, indexer_id, terms FROM dipper_reg_indexing_agreements - WHERE status IN ($1, $2) + WHERE status IN ($1, $2, $3) "#, ) .bind(IndexingAgreementStatus::Created) .bind(IndexingAgreementStatus::AcceptedOnChain) + .bind(IndexingAgreementStatus::Cancelling) .fetch_all(&self.pool) .await?; diff --git a/dipper-pgregistry/tests/it_registry_postgres.rs b/dipper-pgregistry/tests/it_registry_postgres.rs index cd6f3e7a..18530ded 100644 --- a/dipper-pgregistry/tests/it_registry_postgres.rs +++ b/dipper-pgregistry/tests/it_registry_postgres.rs @@ -3344,3 +3344,181 @@ async fn count_created_agreements_by_indexer_counts_only_created() { ); assert_eq!(global, 3, "global counts only the 3 Created rows"); } + +fn fixture_agreement(prefix: u8) -> IndexingAgreementId { + let mut bytes = [0u8; 16]; + bytes[0] = prefix; + bytes[15] = 1; + IndexingAgreementId::from_bytes(bytes) +} + +#[tokio::test] +async fn cancelling_agreements_are_listed_until_their_cancel_fails_too_often() { + let (db, _temp_db) = temp_registry_db().await; + run_fixture( + &db, + include_str!("fixtures/0003_multi_indexer_agreements.sql"), + ) + .await + .expect("Failed to run fixture"); + let registry = PgRegistry::new(db); + let created = fixture_agreement(0xaa); + let accepted = + IndexingAgreementId::from_bytes([0xaa, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 2]); + let ended = + IndexingAgreementId::from_bytes([0xaa, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 3]); + let expired = fixture_agreement(0xcc); + + for id in [created, accepted, expired] { + registry + .mark_indexing_agreement_as_cancelling(&id) + .await + .expect("an agreement that may be live can be marked cancelling"); + } + let result = registry.mark_indexing_agreement_as_cancelling(&ended).await; + assert!( + matches!(result, Err(Error::NoRecordsUpdated)), + "got {result:?}" + ); + registry + .mark_indexing_agreement_as_canceled_by_requester(&expired) + .await + .expect("leave 2 cancelling"); + registry + .record_accepted_audit(&accepted, 1_700_000_000, "0xacc") + .await + .expect("accept record"); + + let just_marked = registry + .get_cancelling_agreements(100, 2, 5) + .await + .expect("cancelling query"); + assert!( + just_marked.is_empty(), + "one just marked waits for the cancel sent with the mark to be mined" + ); + + let listed = registry + .get_cancelling_agreements(100, 2, 0) + .await + .expect("cancelling query"); + let mut seen: Vec<_> = listed + .iter() + .map(|row| { + ( + row.agreement.id, + row.accepted_on_chain, + row.agreement.status, + ) + }) + .collect(); + seen.sort_by_key(|(id, ..)| *id); + assert_eq!( + seen, + vec![ + (created, false, IndexingAgreementStatus::Cancelling), + (accepted, true, IndexingAgreementStatus::Cancelling), + ] + ); + + 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![created, accepted], + "one checked longest ago goes first" + ); + + 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 + .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], + "one that failed too often is left alone" + ); + + let not_cancelling = registry + .record_cancel_check(&fixture_agreement(0xbb), 1) + .await; + assert!(matches!(not_cancelling, Err(Error::NoRecordsUpdated))); +} + +#[tokio::test] +async fn a_cancelling_agreement_stays_live_and_unannounced_until_it_ends() { + let (db, _temp_db) = temp_registry_db().await; + run_fixture( + &db, + include_str!("fixtures/0003_multi_indexer_agreements.sql"), + ) + .await + .expect("Failed to run fixture"); + let registry = PgRegistry::new(db); + let cancelling = + IndexingAgreementId::from_bytes([0xaa, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 2]); + let by_indexer = fixture_agreement(0xbb); + registry + .mark_indexing_agreement_as_canceled_by_requester(&fixture_agreement(0xaa)) + .await + .expect("cancel the other live agreement"); + for id in [cancelling, by_indexer] { + registry + .mark_indexing_agreement_as_cancelling(&id) + .await + .expect("mark cancelling"); + } + registry + .record_accepted_audit(&cancelling, 1_700_000_000, "0xacc") + .await + .expect("accept record"); + registry + .record_cancel_audit(&cancelling, 1_700_000_001, "0xmgr", Some("0xcxl")) + .await + .expect("cancel record"); + + assert!( + !registry.exists_active_agreements().await.unwrap(), + "agreements being cancelled don't keep the listener polling fast; their retry \ + runs on its own timer" + ); + let fee_rates = registry + .get_agreement_fee_rates() + .await + .expect("fee rates query"); + assert!( + fee_rates.iter().any(|(id, ..)| *id == cancelling), + "its fees still count until the cancel lands" + ); + let terminated = registry + .get_agreements_pending_terminated_emission(100) + .await + .expect("terminated query"); + assert!( + terminated.iter().all(|p| p.agreement_id != cancelling), + "not announced as ended while still cancelling" + ); + + let outcome = registry + .apply_reconciliation(&by_indexer, false, Some(CancelKind::ByIndexer)) + .await + .expect("indexer's cancel read from the chain"); + assert!(outcome.did_cancel); + registry + .mark_indexing_agreement_as_canceled_by_requester(&cancelling) + .await + .expect("dipper's cancel confirmed"); + + let terminated = registry + .get_agreements_pending_terminated_emission(100) + .await + .expect("terminated query"); + assert!(terminated.iter().any(|p| p.agreement_id == cancelling)); +} diff --git a/dipper-rpc/src/admin/indexing_agreements.rs b/dipper-rpc/src/admin/indexing_agreements.rs index a08c7e9a..7be13742 100644 --- a/dipper-rpc/src/admin/indexing_agreements.rs +++ b/dipper-rpc/src/admin/indexing_agreements.rs @@ -122,6 +122,9 @@ pub enum Status { /// This is a terminal state. AbandonedByIndexer, + /// Dipper is cancelling the agreement on-chain; it may still be live there. + Cancelling, + /// A fallback for unknown status values. Unknown, } @@ -140,6 +143,7 @@ impl serde::Serialize for Status { Status::AcceptedOnChain => "ACCEPTED_ON_CHAIN", Status::Rejected => "REJECTED", Status::AbandonedByIndexer => "ABANDONED_BY_INDEXER", + Status::Cancelling => "CANCELLING", Status::Unknown => "UNKNOWN", }; serializer.serialize_str(status) @@ -161,6 +165,7 @@ impl<'de> serde::Deserialize<'de> for Status { "ACCEPTED_ON_CHAIN" => Status::AcceptedOnChain, "REJECTED" => Status::Rejected, "ABANDONED_BY_INDEXER" => Status::AbandonedByIndexer, + "CANCELLING" => Status::Cancelling, _ => Status::Unknown, }; Ok(status)