diff --git a/bin/dipper-service/src/network/service/chain_listener.rs b/bin/dipper-service/src/network/service/chain_listener.rs index 4da1555c..64947665 100644 --- a/bin/dipper-service/src/network/service/chain_listener.rs +++ b/bin/dipper-service/src/network/service/chain_listener.rs @@ -881,8 +881,9 @@ where } } - reopen_if_cancelled_but_accepted(snapshot, &agreement, registry, chain_client).await?; - record_accept_of_cancelling(snapshot, &agreement, registry).await; + let reopened = + reopen_if_cancelled_but_accepted(snapshot, &agreement, registry, chain_client).await?; + record_accept_of_cancelling(snapshot, &agreement, reopened, registry).await; // Both transitions are applied atomically downstream so the // Accept-then-Cancel-in-one-snapshot path can't leak an intermediate @@ -998,13 +999,14 @@ 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. It goes back to -/// `Cancelling` for the cancel retry, unless the chain shows it already ended. +/// `Cancelling` for the cancel retry, unless the chain shows it already ended. Returns whether it +/// moved it back. async fn reopen_if_cancelled_but_accepted( snapshot: &AgreementStateSnapshot, agreement: &IndexingAgreement, registry: &R, chain_client: &T, -) -> anyhow::Result<()> +) -> anyhow::Result where R: AgreementRegistry + Sync, T: ChainClient, @@ -1013,20 +1015,24 @@ where && snapshot.state.reached_accepted() && !snapshot.state.is_canceled() { - crate::cancel_dispatch::reopen_if_live(registry, chain_client, agreement).await?; + return Ok( + crate::cancel_dispatch::reopen_if_live(registry, chain_client, agreement).await?, + ); } - Ok(()) + Ok(false) } -/// An agreement dipper is cancelling stays `Cancelling` when the chain shows it accepted, -/// so its accept is recorded here; its end is then announced, along with the accept, -/// once the cancel lands. A withdrawn offer reads as cancelled with no accept time. +/// An agreement dipper is cancelling, or has just moved back to cancelling, stays `Cancelling` +/// when the chain shows it accepted, so its accept is recorded here; its end is then announced, +/// along with the accept, once the cancel lands. A withdrawn offer reads as cancelled with no +/// accept time. async fn record_accept_of_cancelling( snapshot: &AgreementStateSnapshot, agreement: &IndexingAgreement, + reopened: bool, registry: &R, ) { - if agreement.status != IndexingAgreementStatus::Cancelling + if (agreement.status != IndexingAgreementStatus::Cancelling && !reopened) || !snapshot.state.reached_accepted() || snapshot.accepted_at == 0 || !created_after_events_started(agreement) @@ -1343,6 +1349,10 @@ fn log_orphan_cancel( reason = "request_canceled", "Cancelling orphan agreement" ), + Err(crate::registry::Error::NoRecordsUpdated) => tracing::debug!( + agreement_id = %agreement.id, + "Orphan agreement already ended or being cancelled" + ), Err(err) => tracing::warn!( error = %err, agreement_id = %agreement.id, @@ -2534,14 +2544,12 @@ mod tests { } } - /// 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 - /// agreements get a successful `Ok(Some(zero))`. + /// Minimal `ChainClient` mock for chain_listener tests. Records every on-chain cancel + /// attempt. Every agreement reads as not live, so nothing is sent, unless + /// `live_until_cancelled` is set. #[derive(Clone, Default)] struct MockChainClient { cancels: Arc>>, - already_canceled: Arc>>, /// When set, each cancel records whether its agreement was already marked. registry: Option, marked_at_cancel: Arc>>, @@ -2562,10 +2570,6 @@ mod tests { fn was_on_chain_cancel_attempted(&self, id: &IndexingAgreementId) -> bool { self.cancels.lock().unwrap().contains(id.as_bytes()) } - - fn mark_already_canceled_on_chain(&self, id: &IndexingAgreementId) { - self.already_canceled.lock().unwrap().push(*id.as_bytes()); - } } #[async_trait::async_trait] @@ -3112,6 +3116,10 @@ mod tests { assert!(result.is_ok()); assert!(registry.was_reopened(&agreement_id)); assert!(!registry.was_marked_accepted_on_chain(&agreement_id)); + assert!( + registry.audit_writes().contains(&("accept", agreement_id)), + "its accept is recorded, so the cancel retry treats it as paying" + ); } #[tokio::test] @@ -3708,7 +3716,6 @@ mod tests { registry.add_agreement(new_id, IndexingAgreementStatus::AcceptedOnChain); registry.add_agreement(old_id, IndexingAgreementStatus::AcceptedOnChain); registry.add_pending_cancellation(new_id, old_id); - chain_client.mark_already_canceled_on_chain(&old_id); let result = execute_pending_cancellations( &new_id, @@ -3740,7 +3747,6 @@ mod tests { registry.add_agreement(new_id, IndexingAgreementStatus::AcceptedOnChain); registry.add_agreement(old_id, IndexingAgreementStatus::AcceptedOnChain); registry.add_pending_cancellation(new_id, old_id); - chain_client.mark_already_canceled_on_chain(&old_id); sweep_executable_pending_cancellations( ®istry, @@ -4519,7 +4525,6 @@ mod tests { registry.add_agreement(agreement_id, IndexingAgreementStatus::AcceptedOnChain); registry.set_agreement_request_id(agreement_id, request_id); registry.mark_request_canceled(request_id); - chain_client.mark_already_canceled_on_chain(&agreement_id); sweep_orphan_canceled_agreements(®istry, &chain_client, test_agreement_conf().as_ref()) .await; 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 9fcc916b..e07f5ba5 100644 --- a/bin/dipper-service/src/worker/handlers/reassess_indexing_request.rs +++ b/bin/dipper-service/src/worker/handlers/reassess_indexing_request.rs @@ -638,9 +638,13 @@ where let mut directly_cancelled = 0u32; let mut cancel_failures = 0u32; for old_agreement in old_iter { - let Some(new_status) = cancel_unpaired(&ctx, old_agreement).await else { - cancel_failures += 1; - continue; + let new_status = match cancel_unpaired(&ctx, old_agreement).await { + Unpaired::Moved(new_status) => new_status, + Unpaired::AlreadyEnding => continue, + Unpaired::Failed => { + cancel_failures += 1; + continue; + } }; tracing::info!( agreement_id = %old_agreement.id, @@ -897,12 +901,11 @@ mod lifecycle_event_tests { // ---- Mock: worker queue -------------------------------------------------- - /// Records every `send_indexing_agreement_proposal` call's indexer URL and every - /// queued on-chain cancel. Clone shares the buffers for inspection after `handle`. + /// Records every `send_indexing_agreement_proposal` call's indexer URL. Clone shares the + /// buffer for inspection after `handle`. #[derive(Default, Clone)] struct MockQueue { proposals: Arc>>, - cancels_queued: Arc>>, } #[async_trait] @@ -1707,9 +1710,8 @@ mod lifecycle_event_tests { async fn never_accepted_unpaired_cancel_does_not_emit_terminated() { // Both old agreements were never accepted on-chain (Created). One add // pairs with the first old agreement; the second, unpaired old agreement - // reaches the cancel loop but, being never-accepted - // (`!was_accepted`), must NOT emit `terminated`. The add still - // emits `proposed`. Net: exactly one event, a `proposed`. + // reaches the cancel loop but, never accepted, must NOT emit `terminated`. + // The add still emits `proposed`. Net: exactly one event, a `proposed`. let new_idx = indexer_id(0x44); let old_paired = indexer_id(0x55); let old_unpaired = indexer_id(0x56); @@ -1836,14 +1838,12 @@ mod lifecycle_event_tests { ctx.chain_client.fail_cancel = true; let cancelling = ctx.registry.marked_cancelling.clone(); let cancelled = ctx.registry.marked_cancelled.clone(); - let queue = ctx.queue.clone(); let result = handle(ctx, &test_message(0)).await; assert!(result.is_ok(), "got {result:?}"); assert_eq!(*cancelling.lock().unwrap(), vec![leaving.id]); assert!(cancelled.lock().unwrap().is_empty()); - assert!(queue.cancels_queued.lock().unwrap().is_empty()); } } @@ -1876,15 +1876,45 @@ mod lifecycle_event_tests { assert_eq!(*chain_client.cancelled.lock().unwrap(), vec![leaving_id]); } + + #[test] + fn an_agreement_that_ended_since_it_was_listed_is_not_a_failure() { + let agreement = crate::cancel_dispatch::tests::agreement( + crate::registry::IndexingAgreementStatus::AcceptedOnChain, + None, + ); + let backend_down = crate::registry::Error::BackendError(dipper_pgregistry::Error::DbError( + sqlx::Error::PoolTimedOut, + )); + + assert_eq!( + super::unmarked(&agreement, &crate::registry::Error::NoRecordsUpdated), + super::Unpaired::AlreadyEnding + ); + assert_eq!( + super::unmarked(&agreement, &backend_down), + super::Unpaired::Failed + ); + } } -/// Take an old agreement out of the target group, returning its new status, or `None` -/// (logged) when it couldn't be marked. One that may be live on-chain is cancelled there -/// too, which the chain listener retries until it ends. +/// What became of an old agreement taken out of the target group. +#[derive(Debug, PartialEq, Eq)] +enum Unpaired { + /// Moved to the status named. + Moved(&'static str), + /// Already ended or being cancelled, as a race with another path can leave it. + AlreadyEnding, + /// Couldn't be marked; logged. + Failed, +} + +/// Take an old agreement out of the target group. One that may be live on-chain is cancelled +/// there too, which the chain listener retries until it ends. async fn cancel_unpaired( ctx: &Ctx, agreement: &crate::registry::IndexingAgreement, -) -> Option<&'static str> +) -> Unpaired where R: AgreementRegistry + Sync, T: ChainClient, @@ -1893,9 +1923,14 @@ where || (agreement.status == crate::registry::IndexingAgreementStatus::Created && agreement.terms_version_hash.is_some()); if !may_be_live { - return mark_unpaired_cancelled(&ctx.registry, agreement) + return match ctx + .registry + .mark_indexing_agreement_as_canceled_by_requester(&agreement.id) .await - .then_some("CANCELED_BY_REQUESTER"); + { + Ok(()) => Unpaired::Moved("CANCELED_BY_REQUESTER"), + Err(err) => unmarked(agreement, &err), + }; } match crate::cancel_dispatch::start_cancel( &ctx.registry, @@ -1906,38 +1941,33 @@ where ) .await { - Ok(crate::cancel_dispatch::CancelStarted::Ended) => Some("CANCELED_BY_REQUESTER"), - Ok(crate::cancel_dispatch::CancelStarted::Cancelling) => Some("CANCELLING"), - Err(err) => { - tracing::error!( - error = %err, - agreement_id = %agreement.id, - "Failed to mark unpaired old agreement as cancelling in local DB" - ); - None + Ok(crate::cancel_dispatch::CancelStarted::Ended) => { + Unpaired::Moved("CANCELED_BY_REQUESTER") } + Ok(crate::cancel_dispatch::CancelStarted::Cancelling) => Unpaired::Moved("CANCELLING"), + Err(err) => unmarked(agreement, &err), } } -/// Mark an unpaired old agreement CanceledByRequester; false, logged, if that fails. -async fn mark_unpaired_cancelled( - registry: &R, +/// Why an unpaired old agreement couldn't be marked, logged at the level it deserves: one that +/// ended, or started cancelling, since it was listed is expected and needs nothing more. +fn unmarked( agreement: &crate::registry::IndexingAgreement, -) -> bool { - match registry - .mark_indexing_agreement_as_canceled_by_requester(&agreement.id) - .await - { - Ok(()) => true, - Err(err) => { - tracing::error!( - error = %err, - agreement_id = %agreement.id, - "Failed to mark unpaired old agreement as canceled in local DB" - ); - false - } + err: &crate::registry::Error, +) -> Unpaired { + if matches!(err, crate::registry::Error::NoRecordsUpdated) { + tracing::debug!( + agreement_id = %agreement.id, + "Unpaired old agreement already ended or being cancelled" + ); + return Unpaired::AlreadyEnding; } + tracing::error!( + error = %err, + agreement_id = %agreement.id, + "Failed to mark unpaired old agreement as ended in local DB" + ); + Unpaired::Failed } /// Olds reserved from cancellation: one per add-cancel pairing lost to a