From ada72f0aeef5849861e0adbbdc9f194a992cf4e8 Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Fri, 2 Oct 2026 22:46:18 +0300 Subject: [PATCH 1/3] refactor(cancel): share one step for marking a cancel dipper confirmed Both the first cancel attempt and the cancel retry marked an agreement cancelled by dipper, logged it and recorded the cancel, each with its own copy that could drift apart. They now share one step, which records the cancel whenever its transaction is known. --- bin/dipper-service/src/cancel_dispatch.rs | 37 ++++++++++++++----- .../src/network/service/cancel_retry.rs | 27 +------------- 2 files changed, 29 insertions(+), 35 deletions(-) diff --git a/bin/dipper-service/src/cancel_dispatch.rs b/bin/dipper-service/src/cancel_dispatch.rs index 4f6734a2..1e091884 100644 --- a/bin/dipper-service/src/cancel_dispatch.rs +++ b/bin/dipper-service/src/cancel_dispatch.rs @@ -107,17 +107,24 @@ where if agreement.status != IndexingAgreementStatus::AcceptedOnChain { return Ok(CancelStarted::Cancelling); } - Ok(confirm_cancelled(registry, agreement, tx_hash, config).await) + Ok( + if confirm_cancelled(registry, agreement, tx_hash, config).await { + CancelStarted::Ended + } else { + CancelStarted::Cancelling + }, + ) } -/// Mark an accepted agreement whose cancel landed `CanceledByRequester` and record the -/// cancel, so the `terminated` sweep announces it. -async fn confirm_cancelled( +/// Mark an agreement the chain shows dipper ended `CanceledByRequester`, recording the cancel +/// when its transaction is known, so the `terminated` sweep announces it. False, logged, when +/// the mark fails; it stays `Cancelling` for the cancel retry. +pub async fn confirm_cancelled( registry: &R, agreement: &IndexingAgreement, tx_hash: Option, config: &IndexingAgreementConfig, -) -> CancelStarted { +) -> bool { if let Err(err) = registry .mark_indexing_agreement_as_canceled_by_requester(&agreement.id) .await @@ -125,17 +132,27 @@ async fn confirm_cancelled( tracing::warn!( agreement_id = %agreement.id, error = %err, - "Failed to mark a cancelled agreement; the chain listener finishes it" + "Failed to mark an ended agreement cancelled; the cancel retry tries again" ); - return CancelStarted::Cancelling; + 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" + ); + if tx_hash.is_some() { + record_cancel(registry, agreement, tx_hash, config).await; } - record_cancel(registry, agreement, tx_hash, config).await; - CancelStarted::Ended + true } /// Record dipper's own cancel of an accepted agreement, so the `terminated` sweep /// announces it. -pub async fn record_cancel( +async fn record_cancel( registry: &R, agreement: &IndexingAgreement, tx_hash: Option, diff --git a/bin/dipper-service/src/network/service/cancel_retry.rs b/bin/dipper-service/src/network/service/cancel_retry.rs index 4ea06248..5b37a2dd 100644 --- a/bin/dipper-service/src/network/service/cancel_retry.rs +++ b/bin/dipper-service/src/network/service/cancel_retry.rs @@ -6,7 +6,7 @@ use dipper_core::time::now_secs; use thegraph_core::alloy::primitives::B256; use crate::{ - cancel_dispatch::{LiveCancel, cancel_if_live, record_cancel}, + cancel_dispatch::{LiveCancel, cancel_if_live, confirm_cancelled}, chain_client::{ChainClient, ChainClientError}, config::IndexingAgreementConfig, registry::{AgreementRegistry, CancelKind, CancellingAgreement, IndexingAgreement}, @@ -154,30 +154,7 @@ where Some(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 + confirm_cancelled(registry, agreement, tx_hash, config).await } /// Whether the chain shows the indexer ended the agreement, or `None` when it can't be read, From d1e56486373ca6a67b6912ce9417322bc931f72c Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Fri, 2 Oct 2026 22:47:02 +0300 Subject: [PATCH 2/3] fix(cancel): read the chain before sending a cancel Starting a cancel sent it blind, so cancelling offers that never reached the chain cost gas each, and the reassessment waited on every receipt while holding its lock. The agreement is still marked first, so an offer that lands later withdraws itself, but a cancel now goes out only if it is live. --- bin/dipper-service/src/cancel_dispatch.rs | 13 ++--- .../src/network/service/chain_listener.rs | 51 +++++++++++-------- .../handlers/reassess_indexing_request.rs | 26 ++++++++-- 3 files changed, 57 insertions(+), 33 deletions(-) diff --git a/bin/dipper-service/src/cancel_dispatch.rs b/bin/dipper-service/src/cancel_dispatch.rs index 1e091884..3b4c47b1 100644 --- a/bin/dipper-service/src/cancel_dispatch.rs +++ b/bin/dipper-service/src/cancel_dispatch.rs @@ -71,9 +71,9 @@ pub enum CancelStarted { 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. +/// Start ending an agreement that may be live on-chain. It is marked `Cancelling` first, so an +/// offer for it still in flight withdraws itself on landing, then cancelled only if the chain +/// shows it live. Fails, sending nothing, when the mark can't be written. pub async fn start_cancel( registry: &R, chain_client: &T, @@ -87,9 +87,10 @@ where 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) => { + let tx_hash = match cancel_if_live(chain_client, agreement, config).await { + LiveCancel::Ended(tx_hash) => tx_hash, + LiveCancel::NotLive => return Ok(CancelStarted::Cancelling), + LiveCancel::ReadFailed(err) | LiveCancel::CancelFailed(err) => { tracing::warn!( agreement_id = %agreement.id, error = %err, diff --git a/bin/dipper-service/src/network/service/chain_listener.rs b/bin/dipper-service/src/network/service/chain_listener.rs index 278db264..29a1fce1 100644 --- a/bin/dipper-service/src/network/service/chain_listener.rs +++ b/bin/dipper-service/src/network/service/chain_listener.rs @@ -2573,6 +2573,14 @@ mod tests { live_until_cancelled: bool, } + /// A chain on which every agreement is live until a cancel is sent for it. + fn live_chain() -> MockChainClient { + MockChainClient { + live_until_cancelled: true, + ..MockChainClient::default() + } + } + impl MockChainClient { fn was_on_chain_cancel_attempted(&self, id: &IndexingAgreementId) -> bool { self.cancels.lock().unwrap().contains(id.as_bytes()) @@ -3058,7 +3066,7 @@ mod tests { #[tokio::test] async fn test_reconcile_recovers_expired_agreement() { let registry = MockRegistry::new(); - let chain_client = MockChainClient::default(); + let chain_client = live_chain(); let worker_queue = MockWorkerQueue::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); let old_agreement_id = IndexingAgreementId::from_bytes(rand::random()); @@ -3434,7 +3442,7 @@ mod tests { // We should run the acceptance-side bookkeeping (pending cancellations) // AND mark the agreement as CanceledByRequester. let registry = MockRegistry::new(); - let chain_client = MockChainClient::default(); + let chain_client = live_chain(); let worker_queue = MockWorkerQueue::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); let old_agreement_id = IndexingAgreementId::from_bytes(rand::random()); @@ -3474,7 +3482,7 @@ mod tests { #[tokio::test] async fn test_pending_cancellations_all_succeed_records_deleted() { let registry = MockRegistry::new(); - let chain_client = MockChainClient::default(); + let chain_client = live_chain(); let new_id = IndexingAgreementId::from_bytes(rand::random()); let old_id_1 = IndexingAgreementId::from_bytes(rand::random()); let old_id_2 = IndexingAgreementId::from_bytes(rand::random()); @@ -3510,7 +3518,7 @@ mod tests { // execute_pending no longer emits `terminated` directly; it records the // cancel audit and the chain_listener sweep announces it durably. let registry = MockRegistry::new(); - let chain_client = MockChainClient::default(); + let chain_client = live_chain(); let new_id = IndexingAgreementId::from_bytes(rand::random()); let old_id = IndexingAgreementId::from_bytes(rand::random()); @@ -3568,6 +3576,7 @@ mod tests { let registry = MockRegistry::new(); let chain_client = MockChainClient { registry: Some(registry.clone()), + live_until_cancelled: true, ..MockChainClient::default() }; let new_id = IndexingAgreementId::from_bytes(rand::random()); @@ -3696,7 +3705,7 @@ mod tests { #[tokio::test] async fn test_pending_cancellations_transient_failure_retains_record() { let registry = MockRegistry::new(); - let chain_client = MockChainClient::default(); + let chain_client = live_chain(); let new_id = IndexingAgreementId::from_bytes(rand::random()); let old_ok = IndexingAgreementId::from_bytes(rand::random()); let old_fail = IndexingAgreementId::from_bytes(rand::random()); @@ -3782,13 +3791,10 @@ mod tests { #[tokio::test] async fn test_pending_cancellations_already_canceled_on_chain_succeeds() { - // Crash-recovery edge case: the cancel tx confirmed on-chain on a - // prior pass, but dipper crashed before deleting the pending row. - // On the next sweep the chain call surfaces as Ok(None) (the - // SubgraphService contract reverts with IndexingAgreementNotActive; - // the chain client translates that into "already canceled"). The - // handler must still flip the local row to CanceledByRequester and - // delete the pending row, not loop forever. + // Crash-recovery edge case: the cancel landed on a prior pass, but dipper + // crashed before deleting the pending row. The chain shows nothing live, so no + // cancel is sent; the row is left cancelling for the cancel retry to confirm + // and the pending row is deleted rather than retried forever. let registry = MockRegistry::new(); let chain_client = MockChainClient::default(); let new_id = IndexingAgreementId::from_bytes(rand::random()); @@ -3811,8 +3817,8 @@ mod tests { result.is_ok(), "expected idempotent success, got {result:?}" ); - assert!(chain_client.was_on_chain_cancel_attempted(&old_id)); - assert!(registry.was_marked_canceled_by_requester(&old_id)); + assert!(!chain_client.was_on_chain_cancel_attempted(&old_id)); + assert!(registry.was_marked_cancelling(&old_id)); assert!(registry.was_pending_cancellation_deleted(&new_id, &old_id)); } @@ -3838,7 +3844,7 @@ mod tests { ) .await; - assert!(registry.was_marked_canceled_by_requester(&old_id)); + assert!(registry.was_marked_cancelling(&old_id)); assert!(registry.was_pending_cancellation_deleted(&new_id, &old_id)); let remaining = registry .get_pending_cancellations_by_new_agreement(new_id) @@ -3857,7 +3863,7 @@ mod tests { // and the old agreement is still alive. The sweep must complete // the cancellation without needing another snapshot to arrive. let registry = MockRegistry::new(); - let chain_client = MockChainClient::default(); + let chain_client = live_chain(); let new_id = IndexingAgreementId::from_bytes(rand::random()); let old_id = IndexingAgreementId::from_bytes(rand::random()); @@ -4539,7 +4545,7 @@ mod tests { // Canceled is the orphan signature: reassessment fired the chain // cancel and failed, then bailed out. The sweep must pick it up. let registry = MockRegistry::new(); - let chain_client = MockChainClient::default(); + let chain_client = live_chain(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); let request_id = IndexingRequestId::new(); @@ -4586,7 +4592,7 @@ mod tests { // The orphan sweep no longer emits `terminated` directly; it records the // cancel audit and `sweep_pending_terminated_events` announces it. let registry = MockRegistry::new(); - let chain_client = MockChainClient::default(); + let chain_client = live_chain(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); let request_id = IndexingRequestId::new(); @@ -4606,8 +4612,8 @@ mod tests { #[tokio::test] async fn test_orphan_sweep_handles_already_canceled_on_chain() { - // Idempotency check: the chain reports the agreement is already - // canceled (Ok(None)). The sweep must still clean up the local row. + // The chain shows the agreement already ended, so no cancel is sent; it is + // left cancelling for the cancel retry to confirm who ended it. let registry = MockRegistry::new(); let chain_client = MockChainClient::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); @@ -4621,8 +4627,9 @@ mod tests { sweep_orphan_canceled_agreements(®istry, &chain_client, test_agreement_conf().as_ref()) .await; - assert!(chain_client.was_on_chain_cancel_attempted(&agreement_id)); - assert!(registry.was_marked_canceled_by_requester(&agreement_id)); + assert!(!chain_client.was_on_chain_cancel_attempted(&agreement_id)); + assert!(registry.was_marked_cancelling(&agreement_id)); + assert!(!registry.was_marked_canceled_by_requester(&agreement_id)); } #[tokio::test] 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 51b46811..e59d3207 100644 --- a/bin/dipper-service/src/worker/handlers/reassess_indexing_request.rs +++ b/bin/dipper-service/src/worker/handlers/reassess_indexing_request.rs @@ -956,11 +956,12 @@ mod lifecycle_event_tests { // ---- Mock: chain client -------------------------------------------------- - /// Always reports a successful cancel that the post-cancel read confirms, and - /// records the id of every agreement it was asked to cancel. + /// Shows every agreement live until a cancel is sent for it, unless nothing is on + /// chain, and records the id of every agreement it was asked to cancel. #[derive(Default, Clone)] struct MockChainClient { cancelled: Arc>>, + nothing_on_chain: bool, /// When set, each cancel records whether its agreement was already marked cancelling. marked_cancelling: Option>>>, marked_at_cancel: Arc>>, @@ -1018,10 +1019,9 @@ mod lifecycle_event_tests { async fn agreement_still_active( &self, - _agreement_id: &[u8; 16], + agreement_id: &[u8; 16], ) -> std::result::Result { - // Cancel confirmed: agreement is no longer active on-chain. - Ok(false) + Ok(!self.nothing_on_chain && !self.cancelled.lock().unwrap().contains(agreement_id)) } async fn agreement_ended_by_indexer( &self, @@ -1921,6 +1921,22 @@ mod lifecycle_event_tests { assert!(cancelled.lock().unwrap().is_empty()); } + #[tokio::test] + async fn sends_no_cancel_for_an_offer_that_never_reached_the_chain() { + // A cancel of nothing still mines and costs gas. The mark made first means an + // offer still in flight withdraws itself when it lands. + let (mut ctx, leaving) = ctx_cancelling_one(IndexingAgreementStatus::Created); + ctx.chain_client.nothing_on_chain = true; + let chain = ctx.chain_client.clone(); + let marked = ctx.registry.marked_cancelling.clone(); + + let result = handle(ctx, &test_message(0)).await; + + assert!(result.is_ok(), "got {result:?}"); + assert!(chain.cancelled.lock().unwrap().is_empty()); + assert_eq!(*marked.lock().unwrap(), vec![leaving.id]); + } + #[tokio::test] async fn marks_an_accepted_agreement_cancelled_once_its_cancel_lands() { let (ctx, leaving) = ctx_cancelling_one(IndexingAgreementStatus::AcceptedOnChain); From 2ce3c39c8f68a6b3d05b9e1f75d1d9af515458ad Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Fri, 2 Oct 2026 22:47:10 +0300 Subject: [PATCH 3/3] fix(cancel): finish a live cancelled agreement through the cancel retry An agreement dipper had rejected or cancelled that the subgraph showed accepted went to a separate job with its own 40 minute retry budget, queued again on every poll. The listener now reads the chain and, if it is live, moves it back to cancelling so the one cancel retry ends it. --- .../handlers/indexing_requests.rs | 8 - bin/dipper-service/src/cancel_dispatch.rs | 45 + bin/dipper-service/src/main.rs | 1 - .../src/network/service/chain_listener.rs | 246 ++--- .../src/network/service/expiration.rs | 7 - .../src/network/service/liveness_checker.rs | 7 - bin/dipper-service/src/registry.rs | 10 + bin/dipper-service/src/registry/agreement.rs | 8 + .../src/registry/agreement_stub.rs | 8 + bin/dipper-service/src/worker/context.rs | 1 - .../cancel_rejected_agreement_on_chain.rs | 951 ++---------------- .../handlers/reassess_indexing_request.rs | 16 +- .../send_indexing_agreement_proposal.rs | 15 +- .../src/worker/service_queue.rs | 58 +- dipper-pgregistry/src/postgres.rs | 29 + .../tests/it_registry_postgres.rs | 41 + 16 files changed, 302 insertions(+), 1149 deletions(-) diff --git a/bin/dipper-service/src/admin_rpc_server/handlers/indexing_requests.rs b/bin/dipper-service/src/admin_rpc_server/handlers/indexing_requests.rs index 17666cbd..b846dd8e 100644 --- a/bin/dipper-service/src/admin_rpc_server/handlers/indexing_requests.rs +++ b/bin/dipper-service/src/admin_rpc_server/handlers/indexing_requests.rs @@ -364,14 +364,6 @@ mod tests { Ok(JobId::default()) } - async fn cancel_rejected_agreement_on_chain( - &self, - _agreement_id: IndexingAgreementId, - _priority: crate::worker::queue::JobPriority, - ) -> anyhow::Result { - unimplemented!() - } - async fn submit_offer( &self, _agreement_id: IndexingAgreementId, diff --git a/bin/dipper-service/src/cancel_dispatch.rs b/bin/dipper-service/src/cancel_dispatch.rs index 3b4c47b1..b35e8a20 100644 --- a/bin/dipper-service/src/cancel_dispatch.rs +++ b/bin/dipper-service/src/cancel_dispatch.rs @@ -173,6 +173,51 @@ async fn record_cancel( } } +/// Move an agreement dipper had already rejected or cancelled back into `Cancelling` when the +/// chain shows it live after all, so the cancel retry ends it. The chain is read first, so a +/// subgraph report from before dipper's cancel landed reopens nothing; an unreadable chain +/// reopens it anyway, as the retry reads again before sending. True if it was reopened. +pub async fn reopen_if_live( + registry: &R, + chain_client: &T, + agreement: &IndexingAgreement, +) -> RegistryResult +where + R: AgreementRegistry + Sync, + T: ChainClient, +{ + match chain_client + .agreement_still_active(agreement.id.as_bytes()) + .await + { + Ok(false) => return Ok(false), + Ok(true) => {} + Err(err) => tracing::warn!( + agreement_id = %agreement.id, + error = %err, + "Failed to read an ended agreement reported live; the cancel retry checks it" + ), + } + match registry + .reopen_indexing_agreement_cancel(&agreement.id) + .await + { + Ok(()) => {} + Err(crate::registry::Error::NoRecordsUpdated) => return Ok(false), + Err(err) => return Err(err), + } + tracing::warn!( + agreement_id = %agreement.id, + indexer_id = %agreement.indexer.id, + indexing_request_id = %agreement.indexing_request_id, + old_status = %agreement.status, + new_status = "CANCELLING", + reason = "live_on_chain_after_end", + "agreement state transition" + ); + Ok(true) +} + /// What [`cancel_if_live`] found and did. #[derive(Debug)] pub enum LiveCancel { diff --git a/bin/dipper-service/src/main.rs b/bin/dipper-service/src/main.rs index 8b22c649..e6a243f1 100644 --- a/bin/dipper-service/src/main.rs +++ b/bin/dipper-service/src/main.rs @@ -522,7 +522,6 @@ pub async fn main() -> anyhow::Result<()> { let ctx = network::service::chain_listener::Ctx { registry: registry.clone(), - worker_queue: worker_handle.queue().clone(), event_source, chain_client: chain_client.clone(), agreement_conf: chain_listener_agreement_conf.clone(), diff --git a/bin/dipper-service/src/network/service/chain_listener.rs b/bin/dipper-service/src/network/service/chain_listener.rs index 29a1fce1..585626ca 100644 --- a/bin/dipper-service/src/network/service/chain_listener.rs +++ b/bin/dipper-service/src/network/service/chain_listener.rs @@ -56,7 +56,6 @@ use crate::{ AgreementRegistry, CancelKind, IndexingAgreement, IndexingAgreementStatus, PendingCancellationRegistry, ReconciliationItem, }, - worker::service::{JobPriority, WorkerQueue}, }; /// Idle interval used when no `Created` agreements are awaiting acceptance. @@ -106,11 +105,9 @@ impl Handle { } /// Context required by the chain listener service -pub struct Ctx { +pub struct Ctx { /// Registry for querying and updating agreements pub registry: R, - /// Worker queue (still used by reconciliation paths that hand work back to the worker) - pub worker_queue: W, /// Chain event source (subgraph) pub event_source: E, /// Chain client used to cancel on-chain via the RecurringAgreementManager @@ -150,10 +147,9 @@ pub struct ChainListenerState { clippy::too_many_lines, reason = "predates this lint; fix when next touched" )] -pub fn new(ctx: Ctx) -> (Handle, impl Future>) +pub fn new(ctx: Ctx) -> (Handle, impl Future>) where R: AgreementRegistry + ChainListenerStateRegistry + PendingCancellationRegistry + Send + Sync, - W: WorkerQueue + Send + Sync, E: ChainEventSource, T: ChainClient + Send + Sync, { @@ -161,7 +157,6 @@ where let Ctx { registry, - worker_queue, event_source, chain_client, agreement_conf, @@ -321,7 +316,6 @@ where chain_ts_drift_tolerance_secs, bypass_chain_clock_defenses, ®istry, - &worker_queue, &chain_client, &event_source, &mut rx_stop, @@ -440,7 +434,7 @@ struct DrainOutcome { clippy::too_many_lines, reason = "predates this lint; fix when next touched" )] -async fn drain_once( +async fn drain_once( cursor: &mut Cursor, last_persisted_timestamp: &mut Option, last_chain_ts_persist_wall: &mut std::time::Instant, @@ -452,7 +446,6 @@ async fn drain_once( chain_ts_drift_tolerance_secs: u64, bypass_chain_clock_defenses: bool, registry: &R, - worker_queue: &W, chain_client: &T, event_source: &E, rx_stop: &mut mpsc::Receiver<()>, @@ -460,7 +453,6 @@ async fn drain_once( ) -> Result where R: AgreementRegistry + ChainListenerStateRegistry + PendingCancellationRegistry + Send + Sync, - W: WorkerQueue + Send + Sync, E: ChainEventSource, T: ChainClient + Send + Sync, { @@ -603,7 +595,7 @@ where } let agreement = agreements_by_id.remove(&snapshot.agreement_id); - match prepare_reconciliation(&snapshot, agreement, registry, worker_queue).await { + match prepare_reconciliation(&snapshot, agreement, registry, chain_client).await { Ok(Some(prep)) => prepared.push(prep), Ok(None) => {} Err(err) => { @@ -823,15 +815,15 @@ fn apply_chain_ts_drift_cap( clippy::cognitive_complexity, reason = "predates this lint; fix when next touched" )] -async fn prepare_reconciliation( +async fn prepare_reconciliation( snapshot: &AgreementStateSnapshot, agreement: Option, registry: &R, - worker_queue: &W, + chain_client: &T, ) -> anyhow::Result> where R: AgreementRegistry + Sync, - W: WorkerQueue, + T: ChainClient, { tracing::debug!( agreement_id = %snapshot.agreement_id, @@ -868,12 +860,9 @@ where tracing::warn!( agreement_id = %agreement.id, indexer = %snapshot.indexer, - "Rejected agreement accepted on-chain, queuing cancellation" + "Rejected agreement accepted on-chain, cancelling it" ); - worker_queue - // Background: on-chain cleanup of a rejected-then-accepted agreement. - .cancel_rejected_agreement_on_chain(agreement.id, JobPriority::Background) - .await?; + crate::cancel_dispatch::reopen_if_live(registry, chain_client, &agreement).await?; // This row goes Rejected -> Canceled without ever transiting // AcceptedOnChain, so `apply_reconciliation` never records the accept. @@ -892,7 +881,7 @@ where } } - queue_cancel_if_cancelled_but_accepted(snapshot, &agreement, worker_queue).await?; + reopen_if_cancelled_but_accepted(snapshot, &agreement, registry, chain_client).await?; record_accept_of_cancelling(snapshot, &agreement, registry).await; // Both transitions are applied atomically downstream so the @@ -1008,25 +997,23 @@ fn created_after_events_started(agreement: &IndexingAgreement) -> bool { /// Safety net for an agreement dipper cancelled whose offer the indexer accepted /// anyway, such as one that landed after dipper's cancel. Nothing else would end -/// it: reconciliation ignores an accept on a cancelled row. The job reads the -/// chain before acting, so a stale snapshot of an agreement already ended is free. -async fn queue_cancel_if_cancelled_but_accepted( +/// it: reconciliation ignores an accept on a cancelled row. It goes back to +/// `Cancelling` for the cancel retry, unless the chain shows it already ended. +async fn reopen_if_cancelled_but_accepted( snapshot: &AgreementStateSnapshot, agreement: &IndexingAgreement, - worker_queue: &W, -) -> anyhow::Result<()> { + registry: &R, + chain_client: &T, +) -> anyhow::Result<()> +where + R: AgreementRegistry + Sync, + T: ChainClient, +{ if agreement.status == IndexingAgreementStatus::CanceledByRequester && snapshot.state.reached_accepted() && !snapshot.state.is_canceled() { - tracing::warn!( - agreement_id = %agreement.id, - indexer = %snapshot.indexer, - "Cancelled agreement accepted on-chain, queuing cancellation" - ); - worker_queue - .cancel_rejected_agreement_on_chain(agreement.id, JobPriority::Background) - .await?; + crate::cancel_dispatch::reopen_if_live(registry, chain_client, agreement).await?; } Ok(()) } @@ -1131,22 +1118,20 @@ where /// and applies whatever transitions the diff implies. See the module-level /// transition table for the full mapping. #[cfg(test)] -async fn reconcile_agreement( +async fn reconcile_agreement( snapshot: &AgreementStateSnapshot, registry: &R, - worker_queue: &W, chain_client: &T, config: &crate::config::IndexingAgreementConfig, ) -> anyhow::Result<()> where R: AgreementRegistry + PendingCancellationRegistry + Sync, - W: WorkerQueue, T: ChainClient, { let agreement = registry .get_indexing_agreement_by_id(&snapshot.agreement_id) .await?; - let Some(prep) = prepare_reconciliation(snapshot, agreement, registry, worker_queue).await? + let Some(prep) = prepare_reconciliation(snapshot, agreement, registry, chain_client).await? else { return Ok(()); }; @@ -1817,7 +1802,6 @@ mod tests { use dipper_core::ids::{IndexingAgreementId, IndexingRequestId}; use thegraph_core::{DeploymentId, IndexerId, alloy::primitives::ChainId}; use time::OffsetDateTime; - use url::Url; use super::{super::chain_events::AgreementState, *}; use crate::registry::{ @@ -1902,6 +1886,7 @@ mod tests { marked_accepted_on_chain: Vec, marked_canceled_by_requester: Vec, marked_cancelling: Vec, + reopened: 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). @@ -2012,6 +1997,10 @@ mod tests { self.state.lock().unwrap().marked_cancelling.contains(id) } + fn was_reopened(&self, id: &IndexingAgreementId) -> bool { + self.state.lock().unwrap().reopened.contains(id) + } + fn was_marked_canceled_by_requester(&self, id: &IndexingAgreementId) -> bool { self.state .lock() @@ -2203,6 +2192,14 @@ mod tests { Ok(()) } + async fn reopen_indexing_agreement_cancel( + &self, + id: &IndexingAgreementId, + ) -> RegistryResult<()> { + self.state.lock().unwrap().reopened.push(*id); + Ok(()) + } + async fn get_cancelling_agreements( &self, _batch_size: i64, @@ -2545,18 +2542,6 @@ mod tests { } } - // Mock worker queue - #[derive(Clone, Default)] - struct MockWorkerQueue { - cancel_jobs: Arc>>, - } - - impl MockWorkerQueue { - fn was_cancellation_queued(&self, id: &IndexingAgreementId) -> bool { - self.cancel_jobs.lock().unwrap().contains(id) - } - } - /// Minimal `ChainClient` mock for chain_listener tests. Records every /// on-chain cancel attempt. Tests can mark specific agreements as /// already-canceled-on-chain (cancel returns `Ok(None)`); unmarked @@ -2679,60 +2664,12 @@ mod tests { } } - #[async_trait::async_trait] - impl crate::worker::service::WorkerQueue for MockWorkerQueue { - async fn send_indexing_agreement_proposal( - &self, - _candidate_url: Url, - _agreement_id: IndexingAgreementId, - _indexing_request_id: IndexingRequestId, - _deployment_id: DeploymentId, - _deployment_chain_id: ChainId, - _priority: JobPriority, - ) -> anyhow::Result { - Ok(dipper_pgmq::JobId::default()) - } - - async fn reassess_indexing_request( - &self, - _indexing_request_id: IndexingRequestId, - _deployment_id: DeploymentId, - _deployment_chain_id: ChainId, - _num_candidates: usize, - _priority: JobPriority, - ) -> anyhow::Result { - Ok(dipper_pgmq::JobId::default()) - } - - async fn cancel_rejected_agreement_on_chain( - &self, - agreement_id: IndexingAgreementId, - _priority: JobPriority, - ) -> anyhow::Result { - self.cancel_jobs.lock().unwrap().push(agreement_id); - Ok(dipper_pgmq::JobId::default()) - } - - async fn submit_offer( - &self, - _agreement_id: IndexingAgreementId, - _indexing_request_id: IndexingRequestId, - _indexer_url: Url, - _deployment_id: DeploymentId, - _deployment_chain_id: ChainId, - _priority: JobPriority, - ) -> anyhow::Result { - Ok(dipper_pgmq::JobId::default()) - } - } - // -- reconcile_agreement tests -- #[tokio::test] async fn test_reconcile_transitions_created_to_accepted() { 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::Created); @@ -2741,7 +2678,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2749,7 +2685,7 @@ mod tests { assert!(result.is_ok()); assert!(registry.was_marked_accepted_on_chain(&agreement_id)); - assert!(!worker_queue.was_cancellation_queued(&agreement_id)); + assert!(!registry.was_reopened(&agreement_id)); } #[tokio::test] @@ -2757,7 +2693,6 @@ mod tests { // 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); @@ -2765,7 +2700,6 @@ mod tests { reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2774,14 +2708,13 @@ mod tests { 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)); + assert!(!registry.was_reopened(&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); @@ -2790,7 +2723,6 @@ mod tests { reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2806,7 +2738,6 @@ mod tests { // 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 = @@ -2816,7 +2747,6 @@ mod tests { reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2842,7 +2772,6 @@ mod tests { ] { let registry = MockRegistry::new(); let chain_client = MockChainClient::default(); - let worker_queue = MockWorkerQueue::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); registry.add_agreement(agreement_id, status); let mut snapshot = @@ -2852,7 +2781,6 @@ mod tests { reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2881,7 +2809,6 @@ mod tests { ] { 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); @@ -2889,7 +2816,6 @@ mod tests { reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2913,7 +2839,6 @@ mod tests { let registry = MockRegistry::new(); let chain_client = MockChainClient::default(); - let worker_queue = MockWorkerQueue::default(); let events = CapturingEventsProducer::new(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); registry.add_agreement(agreement_id, IndexingAgreementStatus::Created); @@ -2922,7 +2847,6 @@ mod tests { reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2945,7 +2869,6 @@ mod tests { let registry = MockRegistry::new(); let chain_client = MockChainClient::default(); - let worker_queue = MockWorkerQueue::default(); let events = CapturingEventsProducer::new(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); registry.add_agreement(agreement_id, IndexingAgreementStatus::AcceptedOnChain); @@ -2962,7 +2885,6 @@ mod tests { reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2980,10 +2902,12 @@ mod tests { } #[tokio::test] - async fn test_reconcile_queues_cancellation_for_rejected() { + async fn test_reconcile_reopens_the_cancel_of_a_rejected_agreement_live_on_chain() { let registry = MockRegistry::new(); - let chain_client = MockChainClient::default(); - let worker_queue = MockWorkerQueue::default(); + let chain_client = MockChainClient { + live_until_cancelled: true, + ..MockChainClient::default() + }; let agreement_id = IndexingAgreementId::from_bytes(rand::random()); registry.add_agreement(agreement_id, IndexingAgreementStatus::Rejected); @@ -2992,7 +2916,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3000,7 +2923,7 @@ mod tests { assert!(result.is_ok()); assert!(!registry.was_marked_accepted_on_chain(&agreement_id)); - assert!(worker_queue.was_cancellation_queued(&agreement_id)); + assert!(registry.was_reopened(&agreement_id)); } #[tokio::test] @@ -3012,7 +2935,6 @@ mod tests { // the canceler address. let registry = MockRegistry::new(); let chain_client = MockChainClient::default(); - let worker_queue = MockWorkerQueue::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); let signer_address: Address = "0xaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" .parse() @@ -3028,14 +2950,13 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) .await; assert!(result.is_ok()); - assert!(!worker_queue.was_cancellation_queued(&agreement_id)); + assert!(!registry.was_reopened(&agreement_id)); assert!(registry.was_marked_canceled_by_requester(&agreement_id)); assert!(!registry.was_marked_accepted_on_chain(&agreement_id)); } @@ -3044,7 +2965,6 @@ mod tests { async fn test_reconcile_ignores_unknown_agreement() { let registry = MockRegistry::new(); let chain_client = MockChainClient::default(); - let worker_queue = MockWorkerQueue::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); // Don't add the agreement to the registry @@ -3052,7 +2972,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3060,14 +2979,13 @@ mod tests { assert!(result.is_ok()); assert!(!registry.was_marked_accepted_on_chain(&agreement_id)); - assert!(!worker_queue.was_cancellation_queued(&agreement_id)); + assert!(!registry.was_reopened(&agreement_id)); } #[tokio::test] async fn test_reconcile_recovers_expired_agreement() { let registry = MockRegistry::new(); let chain_client = live_chain(); - let worker_queue = MockWorkerQueue::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); let old_agreement_id = IndexingAgreementId::from_bytes(rand::random()); @@ -3079,7 +2997,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3095,7 +3012,6 @@ mod tests { async fn test_reconcile_marks_canceled_by_indexer() { let registry = MockRegistry::new(); let chain_client = MockChainClient::default(); - let worker_queue = MockWorkerQueue::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); let indexer_address: Address = "0x1234567890123456789012345678901234567890" .parse() @@ -3111,7 +3027,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3126,7 +3041,6 @@ mod tests { async fn test_reconcile_marks_canceled_by_requester() { let registry = MockRegistry::new(); let chain_client = MockChainClient::default(); - let worker_queue = MockWorkerQueue::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); let signer_address: Address = "0xaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" .parse() @@ -3142,7 +3056,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3160,7 +3073,6 @@ mod tests { // the state: CanceledByPayer -> ByRequester, the only kind allowed here. let registry = MockRegistry::new(); let chain_client = MockChainClient::default(); - let worker_queue = MockWorkerQueue::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); // A payer address deliberately distinct from dipper's signer key. let payer_address: Address = "0xcccccccccccccccccccccccccccccccccccccccc" @@ -3173,7 +3085,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3184,17 +3095,19 @@ mod tests { assert!(!registry.was_marked_canceled_by_indexer(&agreement_id)); assert!(!registry.was_marked_accepted_on_chain(&agreement_id)); // Already canceled on-chain: must not queue a fresh cancel job. - assert!(!worker_queue.was_cancellation_queued(&agreement_id)); + assert!(!registry.was_reopened(&agreement_id)); } #[tokio::test] - async fn test_reconcile_queues_cancel_for_cancelled_agreement_accepted_on_chain() { + async fn test_reconcile_reopens_the_cancel_of_a_cancelled_agreement_live_on_chain() { // Dipper cancelled the agreement locally, but the indexer accepted its offer // (for example one that landed after dipper's cancel). Nothing else would end - // it, so the listener queues an on-chain cancel and leaves the row cancelled. + // it, so it goes back to cancelling for the cancel retry to end. let registry = MockRegistry::new(); - let chain_client = MockChainClient::default(); - let worker_queue = MockWorkerQueue::default(); + let chain_client = MockChainClient { + live_until_cancelled: true, + ..MockChainClient::default() + }; let agreement_id = IndexingAgreementId::from_bytes(rand::random()); registry.add_agreement(agreement_id, IndexingAgreementStatus::CanceledByRequester); @@ -3202,22 +3115,42 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) .await; assert!(result.is_ok()); - assert!(worker_queue.was_cancellation_queued(&agreement_id)); + assert!(registry.was_reopened(&agreement_id)); assert!(!registry.was_marked_accepted_on_chain(&agreement_id)); } + #[tokio::test] + async fn test_reconcile_leaves_a_cancelled_agreement_the_chain_shows_ended() { + // The subgraph can still report an agreement accepted for a few polls after + // dipper's cancel lands; reopening it would only bring it back to cancelling. + let registry = MockRegistry::new(); + let chain_client = MockChainClient::default(); + let agreement_id = IndexingAgreementId::from_bytes(rand::random()); + registry.add_agreement(agreement_id, IndexingAgreementStatus::CanceledByRequester); + + let snapshot = make_snapshot(agreement_id, AgreementState::Accepted, Address::ZERO); + reconcile_agreement( + &snapshot, + ®istry, + &chain_client, + test_agreement_conf().as_ref(), + ) + .await + .expect("reconcile ok"); + + assert!(!registry.was_reopened(&agreement_id)); + } + #[tokio::test] async fn test_reconcile_cancelled_agreement_already_cancelled_on_chain_queues_nothing() { 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); @@ -3225,14 +3158,13 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) .await; assert!(result.is_ok()); - assert!(!worker_queue.was_cancellation_queued(&agreement_id)); + assert!(!registry.was_reopened(&agreement_id)); } #[tokio::test] @@ -3243,7 +3175,6 @@ mod tests { // waits for the accept. 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); // The listener lagged: dipper only cancelled it locally long after the offer's @@ -3254,7 +3185,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3273,7 +3203,6 @@ mod tests { // out with fallback cancel fields. 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.state.lock().unwrap().fail_cancel_audit = true; @@ -3282,7 +3211,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3298,7 +3226,6 @@ mod tests { // 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_created_at(agreement_id, LIFECYCLE_EVENTS_START - 1); @@ -3307,7 +3234,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3322,7 +3248,6 @@ mod tests { // A withdrawn offer was never accepted, so there is nothing to announce. 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); @@ -3332,7 +3257,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3348,7 +3272,6 @@ mod tests { // agreement has actually ended on-chain. 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); @@ -3356,7 +3279,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3370,7 +3292,6 @@ mod tests { async fn test_reconcile_ignores_already_canceled() { 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::CanceledByIndexer); @@ -3383,7 +3304,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3407,7 +3327,6 @@ mod tests { // rather than incrementing its `errors` counter. let registry = MockRegistry::new(); let chain_client = MockChainClient::default(); - let worker_queue = MockWorkerQueue::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); let signer_address: Address = "0xaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" .parse() @@ -3423,7 +3342,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3443,7 +3361,6 @@ mod tests { // AND mark the agreement as CanceledByRequester. let registry = MockRegistry::new(); let chain_client = live_chain(); - let worker_queue = MockWorkerQueue::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); let old_agreement_id = IndexingAgreementId::from_bytes(rand::random()); let signer_address: Address = "0xaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" @@ -3462,7 +3379,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -4018,7 +3934,6 @@ mod tests { let ctx = Ctx { registry: registry.clone(), - worker_queue: MockWorkerQueue::default(), chain_client: MockChainClient::default(), event_source, config, @@ -4102,7 +4017,6 @@ mod tests { let ctx = Ctx { registry: registry.clone(), - worker_queue: MockWorkerQueue::default(), chain_client: MockChainClient::default(), event_source, config, @@ -4183,7 +4097,6 @@ mod tests { let ctx = Ctx { registry: registry.clone(), - worker_queue: MockWorkerQueue::default(), chain_client: MockChainClient::default(), event_source, config, @@ -4268,7 +4181,6 @@ mod tests { 10, false, ®istry, - &MockWorkerQueue::default(), &chain_client, &event_source, &mut rx_stop, @@ -4333,7 +4245,6 @@ mod tests { 10, false, ®istry, - &MockWorkerQueue::default(), &chain_client, &event_source, &mut rx_stop, @@ -4361,7 +4272,6 @@ mod tests { 10, false, ®istry, - &MockWorkerQueue::default(), &chain_client, &event_source, &mut rx_stop, @@ -4465,7 +4375,6 @@ mod tests { 10, false, ®istry, - &MockWorkerQueue::default(), &chain_client, &event_source, &mut rx_stop, @@ -4506,7 +4415,6 @@ mod tests { let ctx = Ctx { registry: MockRegistry::new(), - worker_queue: MockWorkerQueue::default(), chain_client: MockChainClient::default(), event_source: TimingEventSource { poll_times: poll_times.clone(), diff --git a/bin/dipper-service/src/network/service/expiration.rs b/bin/dipper-service/src/network/service/expiration.rs index 14652149..12802fcf 100644 --- a/bin/dipper-service/src/network/service/expiration.rs +++ b/bin/dipper-service/src/network/service/expiration.rs @@ -538,13 +538,6 @@ mod tests { ) -> anyhow::Result { Ok(JobId::default()) } - async fn cancel_rejected_agreement_on_chain( - &self, - _agreement_id: IndexingAgreementId, - _priority: JobPriority, - ) -> anyhow::Result { - unimplemented!() - } async fn submit_offer( &self, _agreement_id: IndexingAgreementId, diff --git a/bin/dipper-service/src/network/service/liveness_checker.rs b/bin/dipper-service/src/network/service/liveness_checker.rs index 9e520241..367f3a98 100644 --- a/bin/dipper-service/src/network/service/liveness_checker.rs +++ b/bin/dipper-service/src/network/service/liveness_checker.rs @@ -1098,13 +1098,6 @@ mod tests { self.calls.reassessments.lock().unwrap().push(req_id); Ok(JobId::default()) } - async fn cancel_rejected_agreement_on_chain( - &self, - _agr_id: IndexingAgreementId, - _priority: JobPriority, - ) -> anyhow::Result { - unimplemented!() - } async fn submit_offer( &self, _agreement_id: IndexingAgreementId, diff --git a/bin/dipper-service/src/registry.rs b/bin/dipper-service/src/registry.rs index 7d30552c..eb5c7185 100644 --- a/bin/dipper-service/src/registry.rs +++ b/bin/dipper-service/src/registry.rs @@ -401,6 +401,16 @@ impl AgreementRegistry for RegistryProvider { .map_err(Into::into) } + async fn reopen_indexing_agreement_cancel( + &self, + id: &IndexingAgreementId, + ) -> RegistryResult<()> { + self.inner + .reopen_indexing_agreement_cancel(id) + .await + .map_err(Into::into) + } + async fn get_cancelling_agreements( &self, batch_size: i64, diff --git a/bin/dipper-service/src/registry/agreement.rs b/bin/dipper-service/src/registry/agreement.rs index 1309d1d3..762848a8 100644 --- a/bin/dipper-service/src/registry/agreement.rs +++ b/bin/dipper-service/src/registry/agreement.rs @@ -273,6 +273,14 @@ pub trait AgreementRegistry { id: &IndexingAgreementId, ) -> RegistryResult<()>; + /// Move a `CANCELED_BY_REQUESTER` or `REJECTED` agreement the chain shows live back to + /// `CANCELLING`, its cancel attempts reset; [`NoRecordUpdated`](Error::NoRecordsUpdated) + /// otherwise. + async fn reopen_indexing_agreement_cancel( + &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. One that may be paying /// an indexer (accepted, or past the offer deadline, which only an accepted one outlives) diff --git a/bin/dipper-service/src/registry/agreement_stub.rs b/bin/dipper-service/src/registry/agreement_stub.rs index 3e7477be..9e516c4b 100644 --- a/bin/dipper-service/src/registry/agreement_stub.rs +++ b/bin/dipper-service/src/registry/agreement_stub.rs @@ -137,6 +137,10 @@ pub trait StubAgreementRegistry: Send + Sync { unimplemented!("mark_indexing_agreement_as_cancelling") } + async fn reopen_indexing_agreement_cancel(&self, _id: &IndexingAgreementId) -> Result<()> { + unimplemented!("reopen_indexing_agreement_cancel") + } + async fn get_cancelling_agreements( &self, _batch_size: i64, @@ -440,6 +444,10 @@ impl AgreementRegistry for T { StubAgreementRegistry::mark_indexing_agreement_as_cancelling(self, id).await } + async fn reopen_indexing_agreement_cancel(&self, id: &IndexingAgreementId) -> Result<()> { + StubAgreementRegistry::reopen_indexing_agreement_cancel(self, id).await + } + async fn get_cancelling_agreements( &self, batch_size: i64, diff --git a/bin/dipper-service/src/worker/context.rs b/bin/dipper-service/src/worker/context.rs index 5e1b7e60..73175fc5 100644 --- a/bin/dipper-service/src/worker/context.rs +++ b/bin/dipper-service/src/worker/context.rs @@ -234,7 +234,6 @@ impl_from_state!(SendIndexingAgreementProposalCtx { impl_from_state!(CancelRejectedAgreementOnChainCtx { registry, chain_client, - agreement_conf, }); impl_from_state!(SubmitOfferCtx { 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 651245c0..700cb26a 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,27 +1,21 @@ -//! 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. +//! Kept so jobs queued before an upgrade still run. The chain listener now moves an +//! agreement dipper rejected or cancelled that went live on-chain anyway back into +//! `Cancelling` itself, and the cancel retry ends it; this job does the same. -use std::{ - collections::HashSet, - sync::{Arc, LazyLock, Mutex, PoisonError}, - time::Duration, -}; +use std::time::Duration; use dipper_core::ids::IndexingAgreementId; use crate::{ - cancel_dispatch::{LiveCancel, cancel_agreement_on_chain, cancel_if_live}, - chain_client::{ChainClient, ChainClientError}, - config::IndexingAgreementConfig, - registry::{AgreementRegistry, IndexingAgreement, IndexingAgreementStatus}, + cancel_dispatch::reopen_if_live, + chain_client::ChainClient, + registry::{AgreementRegistry, IndexingAgreementStatus}, worker::result::{JobError, JobResult}, }; pub struct Ctx { pub registry: R, pub chain_client: T, - pub agreement_conf: Arc, } /// Cancel on-chain an agreement dipper rejected or cancelled that was accepted anyway. @@ -30,597 +24,97 @@ pub struct Message { pub agreement_id: IndexingAgreementId, } -/// Cancel on-chain an agreement dipper rejected or cancelled that was accepted -/// anyway, via `cancelIndexingAgreementByPayer`, so the indexer isn't paid for -/// work dipper didn't want. -#[expect( - clippy::cognitive_complexity, - reason = "predates this lint; fix when next touched" -)] +/// Hand an agreement dipper rejected or cancelled that is live on-chain to the cancel retry. pub async fn handle(ctx: Ctx, Message { agreement_id }: &Message) -> JobResult<()> where R: AgreementRegistry + Sync, T: ChainClient, { - // Look up the agreement - let agreement = ctx + let Some(agreement) = ctx .registry .get_indexing_agreement_by_id(agreement_id) .await - .map_err(|err| JobError::Fatal(err.into()))?; - - let agreement = match agreement { - Some(a) => a, - None => { - tracing::error!( - agreement_id = %agreement_id, - "Agreement not found for on-chain cancellation" - ); - return Ok(()); - } - }; - - // This job is only queued for these 2 statuses. - match agreement.status { - IndexingAgreementStatus::Rejected => {} - IndexingAgreementStatus::CanceledByRequester => { - return cancel_live_agreement_dipper_cancelled(&ctx, &agreement).await; - } - status => { - tracing::warn!( - agreement_id = %agreement_id, - status = %status, - "Agreement neither Rejected nor CanceledByRequester, skipping on-chain cancellation" - ); - return Ok(()); - } - } - - tracing::info!( - agreement_id = %agreement_id, - indexer_id = %agreement.indexer.id, - "Canceling rejected agreement on-chain" - ); - - // Send the cancellation transaction (mode-aware dispatch). - let on_chain_cancel_tx: Option = - match cancel_agreement_on_chain(&ctx.chain_client, &agreement, &ctx.agreement_conf).await { - Ok(Some(tx_hash)) => { - tracing::info!( - agreement_id = %agreement_id, - tx_hash = %tx_hash, - "Successfully submitted on-chain cancellation" - ); - Some(tx_hash.to_string()) - } - Ok(None) => { - tracing::info!( - agreement_id = %agreement_id, - "Rejected agreement already canceled on-chain; reconciling local state" - ); - None - } - Err(err @ ChainClientError::MissingTermsVersionHash { .. }) => { - // Permanent: the hash never appears, so retrying can't help. Fail - // terminally and leave the live agreement for operator action. - tracing::error!( - agreement_id = %agreement_id, - error = %err, - "Cannot cancel rejected agreement: missing terms_version_hash" - ); - return Err(JobError::Fatal(err.into())); - } - Err(err) => { - tracing::warn!( - agreement_id = %agreement_id, - error = %err, - "Failed to cancel agreement on-chain, will retry" - ); - // Retry with backoff - on-chain transactions can fail due to gas issues, nonce, etc. - return Err(JobError::Retryable(err.into(), Duration::from_secs(30))); - } - }; - - // Once the row is terminal, the cancel audit lets the `terminated` sweep - // announce it (the accept was recorded when the listener queued this job). - // If the mark failed, the listener sees the on-chain cancel and flips it. - if mark_cancellation_complete(&ctx.registry, agreement_id).await { - let manager = ctx.agreement_conf.recurring_agreement_manager().to_string(); - if let Err(err) = ctx - .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" - ); - } - } - - Ok(()) -} - -/// Agreements a job in this process is cancelling right now. The listener can -/// 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); - -/// A job's claim on cancelling one agreement, released when the job ends. -struct Cancelling(IndexingAgreementId); - -impl Cancelling { - fn claim(agreement_id: IndexingAgreementId) -> Option { - // Unlock before a claim exists: dropping one locks the set again. - let claimed = CANCELLING - .lock() - .unwrap_or_else(PoisonError::into_inner) - .insert(agreement_id); - claimed.then(|| Self(agreement_id)) - } -} - -impl Drop for Cancelling { - fn drop(&mut self) { - let mut cancelling = CANCELLING.lock().unwrap_or_else(PoisonError::into_inner); - cancelling.remove(&self.0); - } -} - -/// 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, -) -> JobResult<()> -where - T: ChainClient, -{ - let Some(_claim) = Cancelling::claim(agreement.id) else { - tracing::info!( - agreement_id = %agreement.id, - "Another job is already cancelling this agreement" - ); + .map_err(|err| JobError::Fatal(err.into()))? + else { + tracing::warn!(%agreement_id, "Agreement not found for on-chain cancellation"); return Ok(()); }; - 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(()) - } - LiveCancel::CancelFailed(err @ ChainClientError::MissingTermsVersionHash { .. }) => { - log_caught_live_agreement(agreement, "cancel_impossible", &err.to_string()); - Err(JobError::Fatal(err.into())) - } - LiveCancel::CancelFailed(err) => { - log_caught_live_agreement(agreement, "cancel_failed", &err.to_string()); - Err(JobError::Retryable(err.into(), Duration::from_secs(30))) - } + if !matches!( + agreement.status, + IndexingAgreementStatus::Rejected | IndexingAgreementStatus::CanceledByRequester + ) { + return Ok(()); } -} - -/// ERROR line with stable `event` and `outcome` for alerting: `cancelled` once per -/// agreement ended, `cancel_failed` per failed attempt (so running out of retries -/// is never silent), `cancel_impossible` when it never can be. -fn log_caught_live_agreement(agreement: &IndexingAgreement, outcome: &str, detail: &str) { - tracing::error!( - event = "cancelled_agreement_live_on_chain", - outcome, - agreement_id = %agreement.id, - indexer_id = %agreement.indexer.id, - indexing_request_id = %agreement.indexing_request_id, - detail, - "Agreement dipper had cancelled is live on-chain" - ); -} - -/// 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. -async fn mark_cancellation_complete(registry: &R, agreement_id: &IndexingAgreementId) -> bool -where - R: AgreementRegistry + Sync, -{ - match registry - .mark_indexing_agreement_as_canceled_by_requester(agreement_id) + reopen_if_live(&ctx.registry, &ctx.chain_client, &agreement) .await - { - Ok(()) => { - tracing::info!( - agreement_id = %agreement_id, - old_status = "REJECTED", - new_status = "CANCELED_BY_REQUESTER", - reason = "canceled_on_chain_after_rejection", - "agreement state transition" - ); - true - } - Err(err) => { - tracing::error!( - agreement_id = %agreement_id, - error = %err, - "Failed to update agreement status after on-chain cancellation" - ); - false - } - } + .map(|_| ()) + .map_err(|err| JobError::Retryable(err.into(), Duration::from_secs(30))) } #[cfg(test)] mod tests { - use std::sync::{ - Mutex, - atomic::{AtomicBool, Ordering}, - }; + use std::sync::{Arc, Mutex}; use async_trait::async_trait; - use dipper_core::ids::IndexingRequestId; use dipper_rpc::indexer::indexer_client::sol::RecurringCollectionAgreement; - use thegraph_core::{ - DeploymentId, IndexerId, - alloy::primitives::{Address, B256, U256}, - }; - use time::OffsetDateTime; - use url::Url; + use thegraph_core::alloy::primitives::{Address, B256}; use super::*; use crate::{ - chain_client::{ChainClient, ChainClientError}, - registry::{ - IndexingAgreement, IndexingAgreementStatus, IndexingAgreementTerms, - IndexingAgreementTermsMetadata, - }, + cancel_dispatch::tests::agreement, + chain_client::ChainClientError, + registry::{IndexingAgreement, StubAgreementRegistry}, }; - // ========================================================================= - // Mock implementations - // ========================================================================= - - /// Returns one configurable agreement and records the terminal-cancel and - /// cancel-audit calls. Clones share state, so a test can still assert on its - /// copy after `handle` consumes the ctx. - #[derive(Clone)] struct MockRegistry { - agreement: Arc>>, - marked_canceled: Arc>>, - /// Ids passed to `record_cancel_audit` -- the signal the handler drives - /// the terminated event (the chain_listener sweep emits from this audit). - recorded_cancel_audit: Arc>>, - /// When true, `mark_indexing_agreement_as_canceled_by_requester` errors. - fail_mark: bool, - } - - impl MockRegistry { - fn new(agreement: IndexingAgreement) -> Self { - Self { - agreement: Arc::new(Mutex::new(Some(agreement))), - marked_canceled: Arc::new(Mutex::new(Vec::new())), - recorded_cancel_audit: Arc::new(Mutex::new(Vec::new())), - fail_mark: false, - } - } - - fn with_mark_failure(agreement: IndexingAgreement) -> Self { - Self { - fail_mark: true, - ..Self::new(agreement) - } - } + agreement: IndexingAgreement, + reopened: Arc>>, } #[async_trait] - impl AgreementRegistry for MockRegistry { + impl StubAgreementRegistry for MockRegistry { async fn get_indexing_agreement_by_id( &self, _id: &IndexingAgreementId, ) -> crate::registry::Result> { - Ok(self.agreement.lock().unwrap().clone()) - } - - async fn get_indexing_agreements_by_deployment_id( - &self, - _deployment_id: &DeploymentId, - ) -> crate::registry::Result> { - Ok(vec![]) - } - - async fn get_indexing_agreements_by_indexer_id( - &self, - _indexer_id: &IndexerId, - ) -> crate::registry::Result> { - Ok(vec![]) - } - - async fn get_pending_agreement_indexers_by_deployment( - &self, - _indexer_ids: &[IndexerId], - ) -> crate::registry::Result>> - { - Ok(std::collections::HashMap::new()) + Ok(Some(self.agreement.clone())) } - - async fn get_declined_indexers_by_deployment( - &self, - _default_lookback_days: i32, - _price_lookback_days: i32, - _transient_lookback_minutes: i32, - _uncertain_lookback_days: i32, - ) -> crate::registry::Result>> - { - Ok(std::collections::HashMap::new()) - } - - async fn get_indexing_agreements_by_indexing_request_id( - &self, - _request_id: &IndexingRequestId, - ) -> crate::registry::Result> { - Ok(vec![]) - } - - async fn get_active_indexing_agreements_by_indexing_request_id( - &self, - _request_id: &IndexingRequestId, - ) -> crate::registry::Result> { - Ok(vec![]) - } - - async fn count_accepted_agreements_by_deployment( - &self, - _deployment_id: &DeploymentId, - ) -> crate::registry::Result { - Ok(0) - } - - async fn record_cancel_audit( - &self, - agreement_id: &IndexingAgreementId, - _canceled_at: u64, - _canceled_by: &str, - _canceled_tx: Option<&str>, - ) -> crate::registry::Result<()> { - self.recorded_cancel_audit - .lock() - .unwrap() - .push(*agreement_id); - Ok(()) - } - - async fn register_new_indexing_agreement( - &self, - _params: crate::registry::NewAgreementParams, - ) -> crate::registry::Result { - Ok(IndexingAgreementId::from_bytes(rand::random())) - } - - async fn register_agreement_with_pending_cancellation( - &self, - _params: crate::registry::NewAgreementParams, - _old_agreement_id: IndexingAgreementId, - ) -> crate::registry::Result { - Ok(IndexingAgreementId::from_bytes(rand::random())) - } - - async fn get_unresponsive_indexers( - &self, - _lookback_days: i32, - _chain_id: thegraph_core::alloy::primitives::ChainId, - ) -> crate::registry::Result> { - unimplemented!() - } - - async fn mark_indexing_agreement_as_unresponsive( - &self, - _id: &IndexingAgreementId, - ) -> crate::registry::Result<()> { - unimplemented!() - } - - async fn count_created_agreements_by_indexer( - &self, - ) -> crate::registry::Result<(std::collections::HashMap, u64)> { - unimplemented!() - } - - async fn update_offer_tx_hash( - &self, - _id: &IndexingAgreementId, - _tx_hash: &[u8; 32], - ) -> crate::registry::Result<()> { - Ok(()) - } - - async fn mark_indexing_agreement_as_canceled_by_requester( + async fn reopen_indexing_agreement_cancel( &self, id: &IndexingAgreementId, ) -> crate::registry::Result<()> { - if self.fail_mark { - return Err(crate::registry::Error::NoRecordsUpdated); - } - self.marked_canceled.lock().unwrap().push(*id); + self.reopened.lock().unwrap().push(*id); 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, - _apply_accept: bool, - _cancel: Option, - ) -> crate::registry::Result { - Ok(crate::registry::ReconciliationOutcome { - did_accept: false, - did_cancel: false, - }) - } - - async fn get_expired_created_agreements( - &self, - _batch_size: i64, - _chain_timestamp: u64, - ) -> crate::registry::Result> { - Ok(vec![]) - } - - async fn mark_indexing_agreement_as_expired( - &self, - _id: &IndexingAgreementId, - ) -> crate::registry::Result<()> { - Ok(()) - } - - async fn mark_indexing_agreement_as_rejected( - &self, - _id: &IndexingAgreementId, - _rejection_reason: Option<&str>, - ) -> crate::registry::Result<()> { - Ok(()) - } - - async fn get_accepted_on_chain_agreements( - &self, - _batch_size: i64, - ) -> crate::registry::Result> { - Ok(vec![]) - } - - async fn get_agreements_pending_chain_cancel( - &self, - _batch_size: i64, - ) -> crate::registry::Result> { - Ok(vec![]) - } - - async fn update_agreement_sync_progress( - &self, - _id: &IndexingAgreementId, - _block_height: u64, - _progress_at: time::OffsetDateTime, - ) -> crate::registry::Result<()> { - Ok(()) - } - - async fn count_active_agreements_by_deployment( - &self, - ) -> crate::registry::Result> { - Ok(std::collections::HashMap::new()) - } - - async fn mark_indexing_agreement_as_abandoned( - &self, - _id: &IndexingAgreementId, - ) -> crate::registry::Result { - Err(crate::registry::Error::NoRecordsUpdated) - } - - async fn get_agreement_fee_rates( - &self, - ) -> crate::registry::Result> { - Ok(vec![]) - } } - /// Chain client whose manager cancel always mines (returns a tx hash) and - /// ends the agreement, so the post-cancel liveness read confirms it. `live` - /// is whether the agreement is live on-chain before any cancel lands. - #[derive(Default, Clone)] - struct MockChainClient { - live: Arc, - cancelled: Arc>>, - fail_liveness_read: bool, - fail_cancel: bool, - } - - impl MockChainClient { - fn live() -> Self { - Self { - live: Arc::new(AtomicBool::new(true)), - ..Self::default() - } - } + struct MockChain { + live: bool, } #[async_trait] - impl ChainClient for MockChainClient { - async fn latest_block_timestamp(&self) -> Result { - Err(ChainClientError::RpcError(anyhow::anyhow!( - "latest_block_timestamp not mocked" - ))) - } - + impl ChainClient for MockChain { async fn offer_via_manager( &self, _rca: &RecurringCollectionAgreement, ) -> Result, ChainClientError> { - Ok(None) + unimplemented!() } - async fn cancel_via_manager( &self, _collector: Address, - agreement_id: &[u8; 16], + _agreement_id: &[u8; 16], _version_hash: B256, _options: u16, ) -> Result, ChainClientError> { - if self.fail_cancel { - return Err(ChainClientError::RpcError(anyhow::anyhow!("rpc down"))); - } - self.cancelled.lock().unwrap().push(*agreement_id); - self.live.store(false, Ordering::SeqCst); - Ok(Some(B256::ZERO)) + panic!("the cancel retry sends cancels, not this job") } - async fn reconcile_provider( &self, _collector: Address, _provider: Address, ) -> Result, ChainClientError> { - Ok(None) + unimplemented!() } async fn reconcile_agreement( &self, @@ -629,370 +123,63 @@ mod tests { ) -> Result, ChainClientError> { unimplemented!() } - async fn agreement_still_active( &self, _agreement_id: &[u8; 16], ) -> Result { - if self.fail_liveness_read { - return Err(ChainClientError::RpcError(anyhow::anyhow!("rpc down"))); - } - Ok(self.live.load(Ordering::SeqCst)) + Ok(self.live) } async fn agreement_ended_by_indexer( &self, _agreement_id: &[u8; 16], ) -> Result { - Ok(false) - } - } - - fn test_agreement_conf() -> Arc { - Arc::new(IndexingAgreementConfig::for_tests()) - } - - fn test_deployment_id() -> DeploymentId { - "QmTXzATwNfgGVukV1fX2T6xw9f6LAYRVWpsdXyRWzUR2H9" - .parse() - .unwrap() - } - - fn make_agreement(status: IndexingAgreementStatus) -> IndexingAgreement { - IndexingAgreement { - id: IndexingAgreementId::from_bytes(rand::random()), - nonce_uuid: uuid::Uuid::now_v7(), - created_at: OffsetDateTime::now_utc(), - updated_at: OffsetDateTime::now_utc(), - status, - indexing_request_id: IndexingRequestId::new(), - indexer: crate::registry::Indexer { - id: IndexerId::from(Address::ZERO), - url: Url::parse("https://indexer.example").unwrap(), - }, - terms: IndexingAgreementTerms { - payer: Address::ZERO, - service_provider: Address::ZERO, - data_service: Address::ZERO, - deadline: 0, - ends_at: 0, - max_initial_tokens: U256::ZERO, - max_ongoing_tokens_per_second: U256::ZERO, - min_seconds_per_collection: 0, - max_seconds_per_collection: 0, - conditions: 0, - metadata: IndexingAgreementTermsMetadata { - tokens_per_second: U256::ZERO, - tokens_per_entity_per_second: U256::ZERO, - subgraph_deployment_id: test_deployment_id(), - protocol_network: 1u64, - chain_id: 1u64, - proposed_at: 0, - }, - }, - last_block_height: None, - last_progress_at: None, - rejection_reason: None, - // 32-byte hash so the on-chain cancel path is exercised. - terms_version_hash: Some(vec![0u8; 32]), - } - } - - // ========================================================================= - // Tests - // ========================================================================= - - fn ctx_for(registry: MockRegistry) -> Ctx { - ctx_with_chain(registry, MockChainClient::default()) - } - - fn ctx_with_chain( - registry: MockRegistry, - chain_client: MockChainClient, - ) -> Ctx { - Ctx { - registry, - chain_client, - agreement_conf: test_agreement_conf(), - } - } - - /// Records the `outcome` of every ERROR line carrying the alert's `event`. - #[derive(Clone, Default)] - struct AlertLines(Arc>>); - - impl AlertLines { - fn outcomes(&self) -> Vec { - self.0.lock().unwrap().clone() + unimplemented!() } - } - - impl tracing_subscriber::Layer for AlertLines { - fn on_event( - &self, - event: &tracing::Event<'_>, - _ctx: tracing_subscriber::layer::Context<'_, S>, - ) { - #[derive(Default)] - struct Fields { - event: Option, - outcome: Option, - } - impl tracing::field::Visit for Fields { - fn record_str(&mut self, field: &tracing::field::Field, value: &str) { - match field.name() { - "event" => self.event = Some(value.to_owned()), - "outcome" => self.outcome = Some(value.to_owned()), - _ => {} - } - } - fn record_debug( - &mut self, - _field: &tracing::field::Field, - _value: &dyn std::fmt::Debug, - ) { - } - } - if *event.metadata().level() != tracing::Level::ERROR { - return; - } - let mut fields = Fields::default(); - event.record(&mut fields); - if fields.event.as_deref() == Some("cancelled_agreement_live_on_chain") { - self.0 - .lock() - .unwrap() - .push(fields.outcome.unwrap_or_default()); - } + async fn latest_block_timestamp(&self) -> Result { + unimplemented!() } } - /// Run `handle` with alert lines captured; tokio tests run on one thread. - async fn handle_capturing_alerts( - ctx: Ctx, - agreement_id: IndexingAgreementId, - ) -> (JobResult<()>, Vec) { - use tracing_subscriber::layer::SubscriberExt; - let alerts = AlertLines::default(); - let _guard = - tracing::subscriber::set_default(tracing_subscriber::registry().with(alerts.clone())); - let result = handle(ctx, &Message { agreement_id }).await; - (result, alerts.outcomes()) - } - - #[tokio::test] - async fn logs_one_alert_line_for_a_live_agreement_it_ends() { - let agreement = make_agreement(IndexingAgreementStatus::CanceledByRequester); - let agreement_id = agreement.id; - let ctx = ctx_with_chain(MockRegistry::new(agreement), MockChainClient::live()); - - let (result, alerts) = handle_capturing_alerts(ctx, agreement_id).await; - - assert!(result.is_ok(), "got {result:?}"); - assert_eq!(alerts, vec!["cancelled"]); - } - - #[tokio::test] - async fn logs_no_alert_line_for_an_agreement_that_already_ended() { - let agreement = make_agreement(IndexingAgreementStatus::CanceledByRequester); - let agreement_id = agreement.id; - let ctx = ctx_with_chain(MockRegistry::new(agreement), MockChainClient::default()); - - let (result, alerts) = handle_capturing_alerts(ctx, agreement_id).await; - - assert!(result.is_ok(), "got {result:?}"); - assert!( - alerts.is_empty(), - "a stale snapshot must not alert: {alerts:?}" - ); - } - - #[tokio::test] - async fn logs_an_alert_line_for_each_failed_attempt_at_a_live_agreement() { - // The job can run out of retries; each failure is visible, so a live - // agreement dipper couldn't end is never silent. - let agreement = make_agreement(IndexingAgreementStatus::CanceledByRequester); - let agreement_id = agreement.id; - let chain = MockChainClient { - fail_cancel: true, - ..MockChainClient::live() - }; - let ctx = ctx_with_chain(MockRegistry::new(agreement), chain); - - let (result, alerts) = handle_capturing_alerts(ctx, agreement_id).await; - - assert!( - matches!(result, Err(JobError::Retryable(_, _))), - "got {result:?}" - ); - assert_eq!(alerts, vec!["cancel_failed"]); - } - - #[tokio::test] - async fn fails_with_an_alert_line_when_a_live_agreement_cannot_be_cancelled() { - let mut agreement = make_agreement(IndexingAgreementStatus::CanceledByRequester); - agreement.terms_version_hash = None; - let agreement_id = agreement.id; - let chain = MockChainClient::live(); - let ctx = ctx_with_chain(MockRegistry::new(agreement), chain.clone()); - - let (result, alerts) = handle_capturing_alerts(ctx, agreement_id).await; - - assert!(matches!(result, Err(JobError::Fatal(_))), "got {result:?}"); - assert_eq!(alerts, vec!["cancel_impossible"]); - assert!(chain.cancelled.lock().unwrap().is_empty()); - } - - #[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; 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(); - let ctx = ctx_with_chain(MockRegistry::new(agreement), chain.clone()); - let _first_job = Cancelling::claim(agreement_id).expect("not yet claimed"); - - let (result, alerts) = handle_capturing_alerts(ctx, agreement_id).await; - - assert!(result.is_ok(), "got {result:?}"); - assert!(alerts.is_empty(), "{alerts:?}"); - assert!(chain.cancelled.lock().unwrap().is_empty()); - } - - #[tokio::test] - async fn cancels_a_live_agreement_dipper_had_cancelled() { - // Dipper cancelled the agreement locally, but the indexer accepted its offer - // anyway. The job ends it on-chain and leaves the already-terminal row alone. - let agreement = make_agreement(IndexingAgreementStatus::CanceledByRequester); - let agreement_id = agreement.id; - let registry = MockRegistry::new(agreement); - let chain = MockChainClient::live(); - - let result = handle( - ctx_with_chain(registry.clone(), chain.clone()), - &Message { agreement_id }, - ) - .await; - - assert!(result.is_ok(), "got {result:?}"); - assert_eq!( - *chain.cancelled.lock().unwrap(), - vec![*agreement_id.as_bytes()] - ); - assert!( - registry.marked_canceled.lock().unwrap().is_empty(), - "the row is already cancelled; it must not be marked again" - ); - } - - #[tokio::test] - async fn leaves_alone_an_agreement_dipper_cancelled_that_already_ended() { - // A stale snapshot can report an accept after dipper's own cancel already - // ended the agreement. Reading the chain first avoids a pointless cancel. - let agreement = make_agreement(IndexingAgreementStatus::CanceledByRequester); - let agreement_id = agreement.id; - let chain = MockChainClient::default(); - - let result = handle( - ctx_with_chain(MockRegistry::new(agreement), chain.clone()), - &Message { agreement_id }, - ) - .await; - - assert!(result.is_ok(), "got {result:?}"); - assert!( - chain.cancelled.lock().unwrap().is_empty(), - "nothing to cancel" - ); - } - - #[tokio::test] - async fn retries_without_cancelling_when_the_chain_cannot_be_read() { - let agreement = make_agreement(IndexingAgreementStatus::CanceledByRequester); - let agreement_id = agreement.id; - let chain = MockChainClient { - fail_liveness_read: true, - ..MockChainClient::live() + async fn run(status: IndexingAgreementStatus, live: bool) -> Vec { + let reopened = Arc::new(Mutex::new(Vec::new())); + let agreement = agreement(status, Some(vec![7u8; 32])); + let message = Message { + agreement_id: agreement.id, }; - - let result = handle( - ctx_with_chain(MockRegistry::new(agreement), chain.clone()), - &Message { agreement_id }, - ) - .await; - - assert!( - matches!(result, Err(JobError::Retryable(_, _))), - "got {result:?}" - ); - assert!(chain.cancelled.lock().unwrap().is_empty()); - } - - #[tokio::test] - async fn retries_when_cancelling_a_live_agreement_fails() { - let agreement = make_agreement(IndexingAgreementStatus::CanceledByRequester); - let agreement_id = agreement.id; - let chain = MockChainClient { - fail_cancel: true, - ..MockChainClient::live() + let ctx = Ctx { + registry: MockRegistry { + agreement, + reopened: Arc::clone(&reopened), + }, + chain_client: MockChain { live }, }; - let result = handle( - ctx_with_chain(MockRegistry::new(agreement), chain), - &Message { agreement_id }, - ) - .await; + handle(ctx, &message).await.expect("job ok"); - assert!( - matches!(result, Err(JobError::Retryable(_, _))), - "got {result:?}" - ); + reopened.lock().unwrap().clone() } #[tokio::test] - async fn rejected_agreement_records_cancel_audit_once() { - // The handler no longer emits `terminated` directly: it records the cancel - // audit, and the chain_listener sweep announces it durably. Assert exactly - // one audit was recorded for the agreement. - let agreement = make_agreement(IndexingAgreementStatus::Rejected); - let agreement_id = agreement.id; - let registry = MockRegistry::new(agreement); - - let result = handle(ctx_for(registry.clone()), &Message { agreement_id }).await; - assert!(result.is_ok(), "handle should succeed: {result:?}"); - - let recorded = registry.recorded_cancel_audit.lock().unwrap().clone(); - assert_eq!(recorded, vec![agreement_id], "exactly one cancel audit"); + async fn hands_a_live_agreement_dipper_ended_to_the_cancel_retry() { + for status in [ + IndexingAgreementStatus::Rejected, + IndexingAgreementStatus::CanceledByRequester, + ] { + assert_eq!(run(status, true).await.len(), 1, "{status}"); + } } #[tokio::test] - async fn failed_local_mark_records_no_cancel_audit() { - // The on-chain cancel succeeds but the DB mark fails, so the row stays - // non-terminal and no cancel audit is recorded: the listener sees the - // on-chain cancel, flips the row, and the sweep emits from there. - let agreement = make_agreement(IndexingAgreementStatus::Rejected); - let agreement_id = agreement.id; - let registry = MockRegistry::with_mark_failure(agreement); - - let result = handle(ctx_for(registry.clone()), &Message { agreement_id }).await; - assert!(result.is_ok(), "handle should still return Ok: {result:?}"); + async fn leaves_an_agreement_that_already_ended_or_is_still_wanted() { assert!( - registry.recorded_cancel_audit.lock().unwrap().is_empty(), - "no cancel audit recorded when the local mark failed" + run(IndexingAgreementStatus::CanceledByRequester, false) + .await + .is_empty() ); - } - - #[tokio::test] - async fn non_rejected_agreement_records_nothing() { - let agreement = make_agreement(IndexingAgreementStatus::Created); - let agreement_id = agreement.id; - let registry = MockRegistry::new(agreement); - - let result = handle(ctx_for(registry.clone()), &Message { agreement_id }).await; - assert!(result.is_ok(), "handle should return Ok: {result:?}"); assert!( - registry.recorded_cancel_audit.lock().unwrap().is_empty(), - "no cancel audit recorded for a non-Rejected agreement" + run(IndexingAgreementStatus::AcceptedOnChain, true) + .await + .is_empty() ); } } 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 e59d3207..b9897a4f 100644 --- a/bin/dipper-service/src/worker/handlers/reassess_indexing_request.rs +++ b/bin/dipper-service/src/worker/handlers/reassess_indexing_request.rs @@ -932,15 +932,6 @@ mod lifecycle_event_tests { unimplemented!("not exercised by reassess handler") } - async fn cancel_rejected_agreement_on_chain( - &self, - agreement_id: IndexingAgreementId, - _priority: crate::worker::queue::JobPriority, - ) -> anyhow::Result { - self.cancels_queued.lock().unwrap().push(agreement_id); - Ok(crate::worker::queue::JobId::default()) - } - async fn submit_offer( &self, _agreement_id: IndexingAgreementId, @@ -1209,6 +1200,13 @@ mod lifecycle_event_tests { Ok(()) } + async fn reopen_indexing_agreement_cancel( + &self, + _id: &IndexingAgreementId, + ) -> crate::registry::Result<()> { + unimplemented!() + } + async fn get_cancelling_agreements( &self, _batch_size: i64, 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 af9cbf83..a3139803 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 @@ -529,6 +529,13 @@ mod tests { Ok(()) } + async fn reopen_indexing_agreement_cancel( + &self, + _id: &IndexingAgreementId, + ) -> crate::registry::Result<()> { + unimplemented!() + } + async fn get_cancelling_agreements( &self, _batch_size: i64, @@ -753,14 +760,6 @@ mod tests { Ok(JobId::default()) } - async fn cancel_rejected_agreement_on_chain( - &self, - _agreement_id: IndexingAgreementId, - _priority: JobPriority, - ) -> anyhow::Result { - Ok(JobId::default()) - } - async fn submit_offer( &self, _agreement_id: IndexingAgreementId, diff --git a/bin/dipper-service/src/worker/service_queue.rs b/bin/dipper-service/src/worker/service_queue.rs index 2a0195cd..49ad37e1 100644 --- a/bin/dipper-service/src/worker/service_queue.rs +++ b/bin/dipper-service/src/worker/service_queue.rs @@ -4,19 +4,11 @@ use thegraph_core::{DeploymentId, alloy::primitives::ChainId}; use url::Url; use super::{ - handlers::{ - CancelRejectedAgreementOnChain, ReassessIndexingRequest, SendIndexingAgreementProposal, - SubmitOffer, - }, + handlers::{ReassessIndexingRequest, SendIndexingAgreementProposal, SubmitOffer}, messages::Message, queue::{JobId, JobPriority, Queue}, }; -/// Retries for the job that cancels on-chain an agreement dipper doesn't want but -/// that went live anyway: about 40 minutes of attempts at its 30 s backoff base, -/// since giving up leaves the indexer paid. -const CANCEL_ON_CHAIN_MAX_RETRIES: u32 = 10; - #[async_trait] pub trait WorkerQueue { async fn send_indexing_agreement_proposal( @@ -38,15 +30,6 @@ pub trait WorkerQueue { priority: JobPriority, ) -> anyhow::Result; - /// Cancel a rejected agreement on-chain. When an indexer rejected off-chain - /// but accepted on-chain, this cancels the agreement via - /// `cancelIndexingAgreementByPayer`. - async fn cancel_rejected_agreement_on_chain( - &self, - agreement_id: IndexingAgreementId, - priority: JobPriority, - ) -> anyhow::Result; - /// Submit an RCA offer on-chain as the first step of a new proposal. The /// job retries until the indexer's window to accept has closed, so a /// provider outage costs one offer only if it outlasts that window. @@ -129,22 +112,6 @@ where .await } - async fn cancel_rejected_agreement_on_chain( - &self, - agreement_id: IndexingAgreementId, - priority: JobPriority, - ) -> anyhow::Result { - self.queue - .push_with_max_retries( - Message::CancelRejectedAgreementOnChain(CancelRejectedAgreementOnChain { - agreement_id, - }), - priority, - CANCEL_ON_CHAIN_MAX_RETRIES, - ) - .await - } - async fn submit_offer( &self, agreement_id: IndexingAgreementId, @@ -275,27 +242,4 @@ mod tests { //* Assert assert_eq!(*queue.queue.pushes.lock().unwrap(), vec![None]); } - - /// Giving up on cancelling a live agreement dipper doesn't want leaves the - /// indexer paid, so that job keeps trying well past the queue default. - #[tokio::test] - async fn an_on_chain_cancel_carries_its_longer_retry_budget() { - //* Arrange - let queue = handle(4); - - //* Act - queue - .cancel_rejected_agreement_on_chain( - IndexingAgreementId::from_bytes([0; 16]), - JobPriority::Background, - ) - .await - .unwrap(); - - //* Assert - assert_eq!( - *queue.queue.pushes.lock().unwrap(), - vec![Some(CANCEL_ON_CHAIN_MAX_RETRIES)] - ); - } } diff --git a/dipper-pgregistry/src/postgres.rs b/dipper-pgregistry/src/postgres.rs index fb36f95b..485ab10e 100644 --- a/dipper-pgregistry/src/postgres.rs +++ b/dipper-pgregistry/src/postgres.rs @@ -971,6 +971,35 @@ impl PgRegistry { .await } + /// Move an agreement dipper had already ended, cancelled or rejected, back to `Cancelling` + /// once the chain shows it live after all, with its cancel attempts started afresh. + pub async fn reopen_indexing_agreement_cancel( + &self, + agreement_id: &IndexingAgreementId, + ) -> Result<(), Error> { + let updated = sqlx::query( + r#" + UPDATE dipper_reg_indexing_agreements + SET + status = $1, + cancel_attempts = 0, + cancel_checked_at = NULL, + updated_at = timezone('UTC', now()) + WHERE id = $2 AND status IN ($3, $4) + "#, + ) + .bind(IndexingAgreementStatus::Cancelling) + .bind(agreement_id) + .bind(IndexingAgreementStatus::CanceledByRequester) + .bind(IndexingAgreementStatus::Rejected) + .execute(&self.pool) + .await?; + if updated.rows_affected() == 0 { + return Err(Error::NoRecordsUpdated); + } + 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 /// an indexer (accepted, or past the offer deadline, which only an accepted one outlives) diff --git a/dipper-pgregistry/tests/it_registry_postgres.rs b/dipper-pgregistry/tests/it_registry_postgres.rs index ff1de13f..e82c5617 100644 --- a/dipper-pgregistry/tests/it_registry_postgres.rs +++ b/dipper-pgregistry/tests/it_registry_postgres.rs @@ -3482,6 +3482,47 @@ async fn cancelling_agreements_are_listed_until_their_cancel_fails_too_often() { assert!(matches!(not_cancelling, Err(Error::NoRecordsUpdated))); } +#[tokio::test] +async fn an_ended_agreement_found_live_on_chain_goes_back_to_cancelling() { + 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 ended = fixture_agreement(0xaa); + let accepted = + IndexingAgreementId::from_bytes([0xaa, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 2]); + registry + .mark_indexing_agreement_as_cancelling(&ended) + .await + .expect("mark cancelling"); + assert_eq!(registry.record_cancel_check(&ended, 2).await.unwrap(), 2); + registry + .mark_indexing_agreement_as_canceled_by_requester(&ended) + .await + .expect("mark ended"); + + registry + .reopen_indexing_agreement_cancel(&ended) + .await + .expect("an ended agreement can be reopened"); + + let listed = registry + .get_cancelling_agreements(100, 1, 0) + .await + .expect("cancelling query"); + let ids: Vec<_> = listed.iter().map(|row| row.agreement.id).collect(); + assert_eq!(ids, vec![ended], "its cancel attempts start afresh"); + let still_wanted = registry.reopen_indexing_agreement_cancel(&accepted).await; + assert!( + matches!(still_wanted, Err(Error::NoRecordsUpdated)), + "got {still_wanted:?}" + ); +} + #[tokio::test] async fn a_cancelling_agreement_stays_live_and_unannounced_until_it_ends() { let (db, _temp_db) = temp_registry_db().await;