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 4f6734a2..b35e8a20 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, @@ -107,17 +108,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 +133,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; } - record_cancel(registry, agreement, tx_hash, config).await; - CancelStarted::Ended + 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; + } + 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, @@ -155,6 +173,51 @@ pub 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/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, diff --git a/bin/dipper-service/src/network/service/chain_listener.rs b/bin/dipper-service/src/network/service/chain_listener.rs index 278db264..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 @@ -2573,6 +2558,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()) @@ -2671,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); @@ -2733,7 +2678,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2741,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] @@ -2749,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); @@ -2757,7 +2700,6 @@ mod tests { reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2766,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); @@ -2782,7 +2723,6 @@ mod tests { reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2798,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 = @@ -2808,7 +2747,6 @@ mod tests { reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2834,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 = @@ -2844,7 +2781,6 @@ mod tests { reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2873,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); @@ -2881,7 +2816,6 @@ mod tests { reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2905,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); @@ -2914,7 +2847,6 @@ mod tests { reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2937,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); @@ -2954,7 +2885,6 @@ mod tests { reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2972,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); @@ -2984,7 +2916,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -2992,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] @@ -3004,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() @@ -3020,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)); } @@ -3036,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 @@ -3044,7 +2972,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3052,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 = MockChainClient::default(); - let worker_queue = MockWorkerQueue::default(); + let chain_client = live_chain(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); let old_agreement_id = IndexingAgreementId::from_bytes(rand::random()); @@ -3071,7 +2997,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3087,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() @@ -3103,7 +3027,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3118,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() @@ -3134,7 +3056,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3152,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" @@ -3165,7 +3085,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3176,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); @@ -3194,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); @@ -3217,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] @@ -3235,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 @@ -3246,7 +3185,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3265,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; @@ -3274,7 +3211,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3290,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); @@ -3299,7 +3234,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3314,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); @@ -3324,7 +3257,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3340,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); @@ -3348,7 +3279,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3362,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); @@ -3375,7 +3304,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3399,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() @@ -3415,7 +3342,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3434,8 +3360,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 worker_queue = MockWorkerQueue::default(); + let chain_client = live_chain(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); let old_agreement_id = IndexingAgreementId::from_bytes(rand::random()); let signer_address: Address = "0xaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" @@ -3454,7 +3379,6 @@ mod tests { let result = reconcile_agreement( &snapshot, ®istry, - &worker_queue, &chain_client, test_agreement_conf().as_ref(), ) @@ -3474,7 +3398,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 +3434,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 +3492,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 +3621,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 +3707,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 +3733,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 +3760,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 +3779,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()); @@ -4012,7 +3934,6 @@ mod tests { let ctx = Ctx { registry: registry.clone(), - worker_queue: MockWorkerQueue::default(), chain_client: MockChainClient::default(), event_source, config, @@ -4096,7 +4017,6 @@ mod tests { let ctx = Ctx { registry: registry.clone(), - worker_queue: MockWorkerQueue::default(), chain_client: MockChainClient::default(), event_source, config, @@ -4177,7 +4097,6 @@ mod tests { let ctx = Ctx { registry: registry.clone(), - worker_queue: MockWorkerQueue::default(), chain_client: MockChainClient::default(), event_source, config, @@ -4262,7 +4181,6 @@ mod tests { 10, false, ®istry, - &MockWorkerQueue::default(), &chain_client, &event_source, &mut rx_stop, @@ -4327,7 +4245,6 @@ mod tests { 10, false, ®istry, - &MockWorkerQueue::default(), &chain_client, &event_source, &mut rx_stop, @@ -4355,7 +4272,6 @@ mod tests { 10, false, ®istry, - &MockWorkerQueue::default(), &chain_client, &event_source, &mut rx_stop, @@ -4459,7 +4375,6 @@ mod tests { 10, false, ®istry, - &MockWorkerQueue::default(), &chain_client, &event_source, &mut rx_stop, @@ -4500,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(), @@ -4539,7 +4453,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 +4500,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 +4520,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 +4535,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/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 51b46811..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, @@ -956,11 +947,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 +1010,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, @@ -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, @@ -1921,6 +1919,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); 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;