From e9e12d52eb56866647f2781f12e0ae9a4f9f34ae Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Fri, 2 Oct 2026 22:49:06 +0300 Subject: [PATCH 1/3] fix(cancel): count a cancel the contract refuses before it is sent A cancel the contract rejected while it was being prepared was never counted as a failed attempt, so one that always failed that way was retried, with an error logged, every 5 minutes for ever. It now counts like a cancel that reverts once mined, so it reaches the give-up limit and alert. --- .../src/network/service/cancel_retry.rs | 51 ++++++++++++------- 1 file changed, 32 insertions(+), 19 deletions(-) diff --git a/bin/dipper-service/src/network/service/cancel_retry.rs b/bin/dipper-service/src/network/service/cancel_retry.rs index 5b37a2dd..bce65e6c 100644 --- a/bin/dipper-service/src/network/service/cancel_retry.rs +++ b/bin/dipper-service/src/network/service/cancel_retry.rs @@ -261,22 +261,12 @@ async fn note_check( } } -/// 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" - ); - } + 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) { @@ -300,12 +290,14 @@ fn log_failed_cancel(agreement: &IndexingAgreement, attempts: u32, err: &ChainCl ); } -/// 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. +/// How many of an agreement's cancel attempts a failure uses up. A cancel the contract +/// refused, before sending or once mined, or that mined without ending the agreement counts, +/// and one that can never be sent uses them all. An unreachable chain is retried freely. fn failed_attempts(err: &ChainClientError) -> u32 { match err { - ChainClientError::CancelNotConfirmed { .. } | ChainClientError::TxReverted { .. } => 1, + ChainClientError::CancelNotConfirmed { .. } + | ChainClientError::TxReverted { .. } + | ChainClientError::ContractRevert { .. } => 1, ChainClientError::MissingTermsVersionHash { .. } => MAX_CANCEL_ATTEMPTS, _ => 0, } @@ -403,6 +395,7 @@ mod tests { read_fails: bool, send_fails: bool, mined_cancel_reverts: bool, + reverts_before_sending: bool, cancel_has_no_effect: bool, clock_fails: bool, now: AtomicU64, @@ -429,6 +422,12 @@ mod tests { if self.send_fails { return Err(ChainClientError::RpcError(anyhow::anyhow!("rpc down"))); } + if self.reverts_before_sending { + return Err(ChainClientError::ContractRevert { + selector: [0xde, 0xad, 0xbe, 0xef], + data: Default::default(), + }); + } if self.mined_cancel_reverts { return Err(ChainClientError::TxReverted { tx_hash: B256::repeat_byte(0xee), @@ -721,6 +720,20 @@ mod tests { assert_eq!(registry.attempts.load(Ordering::SeqCst), 1); } + #[tokio::test] + async fn counts_a_cancel_the_contract_refuses_before_it_is_sent() { + // Otherwise one that always reverts is retried, and alerted on, for ever. + let registry = registry_with_one(true); + let chain = MockChain { + reverts_before_sending: 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); From ef2a9faf302b73ba8da67be6b6dd333a0affe9c7 Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Fri, 2 Oct 2026 22:50:00 +0300 Subject: [PATCH 2/3] fix(cancel): keep checking cancels dipper gave up on until they end Once a cancel failed 10 times dipper never looked at the agreement again, so if it later ended it stayed cancelling for good, its fees counted and its indexer kept out of selection. It is now read hourly, without sending more cancels, and closed once the chain shows it ended. --- .../src/network/service/cancel_retry.rs | 58 +++++++++++++++++-- bin/dipper-service/src/registry/agreement.rs | 7 ++- dipper-pgregistry/src/postgres.rs | 16 +++-- .../tests/it_registry_postgres.rs | 21 ++++++- 4 files changed, 91 insertions(+), 11 deletions(-) diff --git a/bin/dipper-service/src/network/service/cancel_retry.rs b/bin/dipper-service/src/network/service/cancel_retry.rs index bce65e6c..b0967311 100644 --- a/bin/dipper-service/src/network/service/cancel_retry.rs +++ b/bin/dipper-service/src/network/service/cancel_retry.rs @@ -12,8 +12,8 @@ use crate::{ registry::{AgreementRegistry, CancelKind, CancellingAgreement, IndexingAgreement}, }; -/// Retried cancels mined without ending an agreement before dipper stops retrying it and -/// leaves it to an operator. Other failures don't count (see `failed_attempts`). +/// Failed cancels before dipper stops sending them and alerts an operator; it still reads the +/// agreement hourly and closes it once it ends. Outages don't count (see `failed_attempts`). pub const MAX_CANCEL_ATTEMPTS: u32 = 10; /// Agreements a sweep takes on, those that may be paying an indexer first; the time budget @@ -93,7 +93,10 @@ async fn retry_cancel( T: ChainClient, { let agreement_id = row.agreement.id; - let (tx_hash, failure) = match cancel_if_live(chain_client, &row.agreement, config).await { + let Some(outcome) = cancel_unless_given_up(chain_client, config, row).await else { + return note_check(registry, row, None).await; + }; + let (tx_hash, failure) = match outcome { LiveCancel::ReadFailed(err) => { tracing::warn!( %agreement_id, @@ -122,6 +125,26 @@ async fn retry_cancel( note_check(registry, row, failure.as_ref()).await; } +/// Cancel the agreement if the chain shows it live. One dipper gave up on is only read, so it +/// is still closed once it ends without paying for more cancels; `None` if it is still live. +async fn cancel_unless_given_up( + chain_client: &T, + config: &IndexingAgreementConfig, + row: &CancellingAgreement, +) -> Option { + if row.cancel_attempts < MAX_CANCEL_ATTEMPTS { + return Some(cancel_if_live(chain_client, &row.agreement, config).await); + } + match chain_client + .agreement_still_active(row.agreement.id.as_bytes()) + .await + { + Ok(true) => None, + Ok(false) => Some(LiveCancel::NotLive), + Err(err) => Some(LiveCancel::ReadFailed(err)), + } +} + /// 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 ended otherwise is left /// to the chain listener for a while; one the indexer ended then becomes `CanceledByIndexer`. @@ -286,7 +309,8 @@ fn log_failed_cancel(agreement: &IndexingAgreement, attempts: u32, err: &ChainCl indexing_request_id = %agreement.indexing_request_id, attempts, error = %err, - "Gave up cancelling an agreement on-chain; it may still be live" + "Gave up cancelling an agreement on-chain; it may still be live. Dipper checks it hourly \ + and closes it once it ends; set its cancel_attempts to 0 to send cancels again" ); } @@ -486,6 +510,7 @@ mod tests { cancelling: vec![CancellingAgreement { agreement: cancelling, accepted_on_chain, + cancel_attempts: 0, }], ..MockRegistry::default() } @@ -734,6 +759,31 @@ mod tests { assert_eq!(registry.attempts.load(Ordering::SeqCst), 1); } + #[tokio::test] + async fn only_reads_an_agreement_it_gave_up_cancelling() { + let mut registry = registry_with_one(true); + registry.cancelling[0].cancel_attempts = MAX_CANCEL_ATTEMPTS; + let chain = live_chain(); + + retry(®istry, &chain, 0).await; + + assert_eq!(chain.cancels_sent.load(Ordering::SeqCst), 0); + assert_eq!(registry.checks.load(Ordering::SeqCst), 1); + assert!(registry.marked_cancelled.lock().unwrap().is_empty()); + } + + #[tokio::test] + async fn closes_an_agreement_it_gave_up_cancelling_once_it_ends() { + // Otherwise it stays cancelling, its indexer kept out of selection, for ever. + let mut registry = registry_with_one(false); + registry.cancelling[0].cancel_attempts = MAX_CANCEL_ATTEMPTS; + let chain = MockChain::default(); + + retry(®istry, &chain, DEADLINE + 1).await; + + assert_eq!(registry.marked_cancelled.lock().unwrap().len(), 1); + } + #[tokio::test] async fn counts_a_cancel_that_mines_without_ending_the_agreement() { let registry = registry_with_one(true); diff --git a/bin/dipper-service/src/registry/agreement.rs b/bin/dipper-service/src/registry/agreement.rs index 762848a8..c4d2e660 100644 --- a/bin/dipper-service/src/registry/agreement.rs +++ b/bin/dipper-service/src/registry/agreement.rs @@ -281,8 +281,8 @@ pub trait AgreementRegistry { 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. One that may be paying + /// `CANCELLING` agreements marked over `min_age_minutes` ago, those checked longest ago + /// first; one whose cancel has failed `max_attempts` times only once an hour. One that may be paying /// an indexer (accepted, or past the offer deadline, which only an accepted one outlives) /// counts as checked an hour earlier, so it goes first without holding the rest back. async fn get_cancelling_agreements( @@ -562,6 +562,8 @@ pub struct CancellingAgreement { pub agreement: IndexingAgreement, /// Whether dipper saw it accepted on-chain, so its end is announced. pub accepted_on_chain: bool, + /// Cancels that failed in a way retrying may not fix. + pub cancel_attempts: u32, } impl TryFrom for CancellingAgreement { @@ -571,6 +573,7 @@ impl TryFrom for CancellingAgreement { Ok(Self { agreement: value.agreement.try_into()?, accepted_on_chain: value.accepted_on_chain, + cancel_attempts: value.cancel_attempts, }) } } diff --git a/dipper-pgregistry/src/postgres.rs b/dipper-pgregistry/src/postgres.rs index 485ab10e..2a76bd74 100644 --- a/dipper-pgregistry/src/postgres.rs +++ b/dipper-pgregistry/src/postgres.rs @@ -193,15 +193,19 @@ pub struct CancellingAgreement { pub agreement: IndexingAgreement, /// Whether dipper saw it accepted on-chain, so its end is announced. pub accepted_on_chain: bool, + /// Cancels that failed in a way retrying may not fix. + pub cancel_attempts: u32, } 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")?; + let cancel_attempts: i32 = row.try_get("cancel_attempts")?; Ok(Self { agreement: IndexingAgreement::from_row(row)?, accepted_on_chain: accepted_at.is_some(), + cancel_attempts: u32::try_from(cancel_attempts).unwrap_or_default(), }) } } @@ -1000,8 +1004,8 @@ impl PgRegistry { Ok(()) } - /// `Cancelling` agreements marked over `min_age_minutes` ago whose cancel has failed - /// fewer than `max_attempts` times, those checked longest ago first. One that may be paying + /// `Cancelling` agreements marked over `min_age_minutes` ago, those checked longest ago + /// first; one whose cancel has failed `max_attempts` times only once an hour. One that may be paying /// an indexer (accepted, or past the offer deadline, which only an accepted one outlives) /// counts as checked an hour earlier, so it goes first without holding the rest back. pub async fn get_cancelling_agreements( @@ -1027,10 +1031,14 @@ impl PgRegistry { last_progress_at, rejection_reason, terms_version_hash, - accepted_at + accepted_at, + cancel_attempts FROM dipper_reg_indexing_agreements WHERE status = $1 - AND cancel_attempts < $2 + AND ( + cancel_attempts < $2 + OR cancel_checked_at < timezone('UTC', now()) - INTERVAL '1 hour' + ) AND updated_at < timezone('UTC', now()) - make_interval(mins => $4) ORDER BY cancel_checked_at - CASE diff --git a/dipper-pgregistry/tests/it_registry_postgres.rs b/dipper-pgregistry/tests/it_registry_postgres.rs index e82c5617..dfb91dbc 100644 --- a/dipper-pgregistry/tests/it_registry_postgres.rs +++ b/dipper-pgregistry/tests/it_registry_postgres.rs @@ -3473,7 +3473,26 @@ async fn cancelling_agreements_are_listed_until_their_cancel_fails_too_often() { assert_eq!( ids, vec![accepted], - "one that failed too often is left alone" + "one that failed too often waits an hour between checks" + ); + + sqlx::query( + "UPDATE dipper_reg_indexing_agreements \ + SET cancel_checked_at = cancel_checked_at - INTERVAL '2 hours' WHERE id = $1", + ) + .bind(created) + .execute(&db) + .await + .expect("Failed to age the check"); + let listed = registry + .get_cancelling_agreements(100, 2, 0) + .await + .expect("cancelling query"); + let given_up = listed.iter().find(|row| row.agreement.id == created); + assert_eq!( + given_up.map(|row| row.cancel_attempts), + Some(2), + "and is checked again after it" ); let not_cancelling = registry From 083b8fdd69d9b668bd3c5c24414123ea2a0dfcaf Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Fri, 2 Oct 2026 22:52:45 +0300 Subject: [PATCH 3/3] fix(cancel): slow a failing cancel to hourly instead of stopping it Counting refusals before sending meant a paused manager used up every agreement's 10 attempts in under an hour, and dipper then never sent their cancels again once it was unpaused. It now alerts once at the limit and keeps retrying hourly, so those cancels resume by themselves. --- .../src/network/service/cancel_retry.rs | 75 +++++-------------- bin/dipper-service/src/registry/agreement.rs | 3 - dipper-pgregistry/src/postgres.rs | 7 +- .../tests/it_registry_postgres.rs | 8 +- 4 files changed, 21 insertions(+), 72 deletions(-) diff --git a/bin/dipper-service/src/network/service/cancel_retry.rs b/bin/dipper-service/src/network/service/cancel_retry.rs index b0967311..49a651a6 100644 --- a/bin/dipper-service/src/network/service/cancel_retry.rs +++ b/bin/dipper-service/src/network/service/cancel_retry.rs @@ -12,8 +12,8 @@ use crate::{ registry::{AgreementRegistry, CancelKind, CancellingAgreement, IndexingAgreement}, }; -/// Failed cancels before dipper stops sending them and alerts an operator; it still reads the -/// agreement hourly and closes it once it ends. Outages don't count (see `failed_attempts`). +/// Failed cancels before dipper alerts an operator and retries the agreement only hourly, so a +/// paused manager recovers once unpaused. Outages don't count (see `failed_attempts`). pub const MAX_CANCEL_ATTEMPTS: u32 = 10; /// Agreements a sweep takes on, those that may be paying an indexer first; the time budget @@ -93,10 +93,7 @@ async fn retry_cancel( T: ChainClient, { let agreement_id = row.agreement.id; - let Some(outcome) = cancel_unless_given_up(chain_client, config, row).await else { - return note_check(registry, row, None).await; - }; - let (tx_hash, failure) = match outcome { + let (tx_hash, failure) = match cancel_if_live(chain_client, &row.agreement, config).await { LiveCancel::ReadFailed(err) => { tracing::warn!( %agreement_id, @@ -125,26 +122,6 @@ async fn retry_cancel( note_check(registry, row, failure.as_ref()).await; } -/// Cancel the agreement if the chain shows it live. One dipper gave up on is only read, so it -/// is still closed once it ends without paying for more cancels; `None` if it is still live. -async fn cancel_unless_given_up( - chain_client: &T, - config: &IndexingAgreementConfig, - row: &CancellingAgreement, -) -> Option { - if row.cancel_attempts < MAX_CANCEL_ATTEMPTS { - return Some(cancel_if_live(chain_client, &row.agreement, config).await); - } - match chain_client - .agreement_still_active(row.agreement.id.as_bytes()) - .await - { - Ok(true) => None, - Ok(false) => Some(LiveCancel::NotLive), - Err(err) => Some(LiveCancel::ReadFailed(err)), - } -} - /// 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 ended otherwise is left /// to the chain listener for a while; one the indexer ended then becomes `CanceledByIndexer`. @@ -273,7 +250,7 @@ async fn note_check( { Ok(attempts) => { if let Some(err) = failure.filter(|_| failed_attempts > 0) { - log_failed_cancel(agreement, attempts, err); + log_failed_cancel(agreement, attempts, failed_attempts, err); } } Err(err) => tracing::warn!( @@ -292,8 +269,17 @@ fn log_uncounted_failure(agreement: &IndexingAgreement, err: &ChainClientError) ); } -fn log_failed_cancel(agreement: &IndexingAgreement, attempts: u32, err: &ChainClientError) { - if attempts < MAX_CANCEL_ATTEMPTS { +/// One ERROR as an agreement reaches the limit, for an operator to look into; a WARN for +/// every other failed cancel. +fn log_failed_cancel( + agreement: &IndexingAgreement, + attempts: u32, + failed: u32, + err: &ChainClientError, +) { + let reached_limit = + attempts >= MAX_CANCEL_ATTEMPTS && attempts.saturating_sub(failed) < MAX_CANCEL_ATTEMPTS; + if !reached_limit { tracing::warn!( agreement_id = %agreement.id, attempts, @@ -303,14 +289,13 @@ fn log_failed_cancel(agreement: &IndexingAgreement, attempts: u32, err: &ChainCl return; } tracing::error!( - event = "agreement_cancel_abandoned", + event = "agreement_cancel_stuck", 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. Dipper checks it hourly \ - and closes it once it ends; set its cancel_attempts to 0 to send cancels again" + "Cancelling an agreement keeps failing; it may still be live. Dipper now retries it hourly" ); } @@ -510,7 +495,6 @@ mod tests { cancelling: vec![CancellingAgreement { agreement: cancelling, accepted_on_chain, - cancel_attempts: 0, }], ..MockRegistry::default() } @@ -759,31 +743,6 @@ mod tests { assert_eq!(registry.attempts.load(Ordering::SeqCst), 1); } - #[tokio::test] - async fn only_reads_an_agreement_it_gave_up_cancelling() { - let mut registry = registry_with_one(true); - registry.cancelling[0].cancel_attempts = MAX_CANCEL_ATTEMPTS; - let chain = live_chain(); - - retry(®istry, &chain, 0).await; - - assert_eq!(chain.cancels_sent.load(Ordering::SeqCst), 0); - assert_eq!(registry.checks.load(Ordering::SeqCst), 1); - assert!(registry.marked_cancelled.lock().unwrap().is_empty()); - } - - #[tokio::test] - async fn closes_an_agreement_it_gave_up_cancelling_once_it_ends() { - // Otherwise it stays cancelling, its indexer kept out of selection, for ever. - let mut registry = registry_with_one(false); - registry.cancelling[0].cancel_attempts = MAX_CANCEL_ATTEMPTS; - let chain = MockChain::default(); - - retry(®istry, &chain, DEADLINE + 1).await; - - assert_eq!(registry.marked_cancelled.lock().unwrap().len(), 1); - } - #[tokio::test] async fn counts_a_cancel_that_mines_without_ending_the_agreement() { let registry = registry_with_one(true); diff --git a/bin/dipper-service/src/registry/agreement.rs b/bin/dipper-service/src/registry/agreement.rs index c4d2e660..24122084 100644 --- a/bin/dipper-service/src/registry/agreement.rs +++ b/bin/dipper-service/src/registry/agreement.rs @@ -562,8 +562,6 @@ pub struct CancellingAgreement { pub agreement: IndexingAgreement, /// Whether dipper saw it accepted on-chain, so its end is announced. pub accepted_on_chain: bool, - /// Cancels that failed in a way retrying may not fix. - pub cancel_attempts: u32, } impl TryFrom for CancellingAgreement { @@ -573,7 +571,6 @@ impl TryFrom for CancellingAgreement { Ok(Self { agreement: value.agreement.try_into()?, accepted_on_chain: value.accepted_on_chain, - cancel_attempts: value.cancel_attempts, }) } } diff --git a/dipper-pgregistry/src/postgres.rs b/dipper-pgregistry/src/postgres.rs index 2a76bd74..3e7d062b 100644 --- a/dipper-pgregistry/src/postgres.rs +++ b/dipper-pgregistry/src/postgres.rs @@ -193,19 +193,15 @@ pub struct CancellingAgreement { pub agreement: IndexingAgreement, /// Whether dipper saw it accepted on-chain, so its end is announced. pub accepted_on_chain: bool, - /// Cancels that failed in a way retrying may not fix. - pub cancel_attempts: u32, } 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")?; - let cancel_attempts: i32 = row.try_get("cancel_attempts")?; Ok(Self { agreement: IndexingAgreement::from_row(row)?, accepted_on_chain: accepted_at.is_some(), - cancel_attempts: u32::try_from(cancel_attempts).unwrap_or_default(), }) } } @@ -1031,8 +1027,7 @@ impl PgRegistry { last_progress_at, rejection_reason, terms_version_hash, - accepted_at, - cancel_attempts + accepted_at FROM dipper_reg_indexing_agreements WHERE status = $1 AND ( diff --git a/dipper-pgregistry/tests/it_registry_postgres.rs b/dipper-pgregistry/tests/it_registry_postgres.rs index dfb91dbc..215e946c 100644 --- a/dipper-pgregistry/tests/it_registry_postgres.rs +++ b/dipper-pgregistry/tests/it_registry_postgres.rs @@ -3488,11 +3488,9 @@ async fn cancelling_agreements_are_listed_until_their_cancel_fails_too_often() { .get_cancelling_agreements(100, 2, 0) .await .expect("cancelling query"); - let given_up = listed.iter().find(|row| row.agreement.id == created); - assert_eq!( - given_up.map(|row| row.cancel_attempts), - Some(2), - "and is checked again after it" + assert!( + listed.iter().any(|row| row.agreement.id == created), + "and is tried again after it" ); let not_cancelling = registry