From 597cf631985d5b8069ca89bf1fc6fbaa6af423e7 Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Mon, 5 Oct 2026 13:03:35 +0100 Subject: [PATCH 1/3] fix(cancel): clear the old end when an agreement is found live again An agreement reopened because the chain shows it live kept its earlier end on record, such as a withdrawn offer, so its later real end was announced with the old transaction and time, or never if one had gone out. That record is now cleared when the chain showed it live. --- bin/dipper-service/src/cancel_dispatch.rs | 67 ++++++++++++++++--- .../src/network/service/chain_listener.rs | 1 + bin/dipper-service/src/registry.rs | 3 +- bin/dipper-service/src/registry/agreement.rs | 4 +- .../src/registry/agreement_stub.rs | 14 +++- .../cancel_rejected_agreement_on_chain.rs | 1 + dipper-pgregistry/src/postgres.rs | 11 ++- .../tests/it_registry_postgres.rs | 23 ++++++- 8 files changed, 106 insertions(+), 18 deletions(-) diff --git a/bin/dipper-service/src/cancel_dispatch.rs b/bin/dipper-service/src/cancel_dispatch.rs index 17c0eed2..0a3e8e83 100644 --- a/bin/dipper-service/src/cancel_dispatch.rs +++ b/bin/dipper-service/src/cancel_dispatch.rs @@ -232,20 +232,23 @@ where R: AgreementRegistry + Sync, T: ChainClient, { - match chain_client + let seen_live = match chain_client .agreement_on_chain(agreement.id.as_bytes()) .await { - Ok(AgreementOnChain::Live) => {} + Ok(AgreementOnChain::Live) => true, Ok(_) => return Ok(false), - Err(err) => tracing::warn!( - agreement_id = %agreement.id, - error = %err, - "Failed to read an ended agreement reported live; the cancel retry checks it" - ), - } + Err(err) => { + tracing::warn!( + agreement_id = %agreement.id, + error = %err, + "Failed to read an ended agreement reported live; the cancel retry checks it" + ); + false + } + }; match registry - .reopen_indexing_agreement_cancel(&agreement.id) + .reopen_indexing_agreement_cancel(&agreement.id, seen_live) .await { Ok(()) => {} @@ -632,4 +635,50 @@ pub(crate) mod tests { } )); } + + /// Records what each reopen was told about the chain. + #[derive(Default)] + struct ReopenRegistry { + seen_live: Mutex>, + } + + #[async_trait] + impl crate::registry::StubAgreementRegistry for ReopenRegistry { + async fn reopen_indexing_agreement_cancel( + &self, + _id: &IndexingAgreementId, + seen_live: bool, + ) -> crate::registry::Result<()> { + self.seen_live.lock().unwrap().push(seen_live); + Ok(()) + } + } + + #[tokio::test] + async fn a_reopen_clears_the_end_on_record_only_when_the_chain_shows_it_live() { + // An unread chain reopens it all the same, but it may have ended, so its record stays. + let ag = agreement( + IndexingAgreementStatus::CanceledByRequester, + Some(vec![7u8; 32]), + ); + let live = RecordingChainClient { + still_active_after_cancel: true, + ..Default::default() + }; + let unread = RecordingChainClient { + read_back_fails: true, + ..Default::default() + }; + + for (client, seen_live) in [(live, true), (unread, false)] { + let registry = ReopenRegistry::default(); + + let reopened = super::reopen_if_live(®istry, &client, &ag) + .await + .expect("reopen"); + + assert!(reopened); + assert_eq!(*registry.seen_live.lock().unwrap(), vec![seen_live]); + } + } } diff --git a/bin/dipper-service/src/network/service/chain_listener.rs b/bin/dipper-service/src/network/service/chain_listener.rs index 64947665..eacc5ebb 100644 --- a/bin/dipper-service/src/network/service/chain_listener.rs +++ b/bin/dipper-service/src/network/service/chain_listener.rs @@ -2217,6 +2217,7 @@ mod tests { async fn reopen_indexing_agreement_cancel( &self, id: &IndexingAgreementId, + _seen_live: bool, ) -> RegistryResult<()> { self.state.lock().unwrap().reopened.push(*id); Ok(()) diff --git a/bin/dipper-service/src/registry.rs b/bin/dipper-service/src/registry.rs index 239ff563..e49d8c4c 100644 --- a/bin/dipper-service/src/registry.rs +++ b/bin/dipper-service/src/registry.rs @@ -414,9 +414,10 @@ impl AgreementRegistry for RegistryProvider { async fn reopen_indexing_agreement_cancel( &self, id: &IndexingAgreementId, + seen_live: bool, ) -> RegistryResult<()> { self.inner - .reopen_indexing_agreement_cancel(id) + .reopen_indexing_agreement_cancel(id, seen_live) .await .map_err(Into::into) } diff --git a/bin/dipper-service/src/registry/agreement.rs b/bin/dipper-service/src/registry/agreement.rs index 54ac69f0..c7cfeaab 100644 --- a/bin/dipper-service/src/registry/agreement.rs +++ b/bin/dipper-service/src/registry/agreement.rs @@ -284,10 +284,12 @@ pub trait AgreementRegistry { /// Move a `CANCELED_BY_REQUESTER` or `REJECTED` agreement the chain shows live back to /// `CANCELLING`, its cancel attempts reset; [`NoRecordUpdated`](Error::NoRecordsUpdated) - /// otherwise. + /// otherwise. `seen_live` when the chain was read and showed it live, which clears the end + /// on record and any announcement of it. async fn reopen_indexing_agreement_cancel( &self, id: &IndexingAgreementId, + seen_live: bool, ) -> RegistryResult<()>; /// `CANCELLING` agreements marked over `min_age_minutes` ago, those checked longest ago diff --git a/bin/dipper-service/src/registry/agreement_stub.rs b/bin/dipper-service/src/registry/agreement_stub.rs index cbb62dd6..d9092a6d 100644 --- a/bin/dipper-service/src/registry/agreement_stub.rs +++ b/bin/dipper-service/src/registry/agreement_stub.rs @@ -141,7 +141,11 @@ pub trait StubAgreementRegistry: Send + Sync { unimplemented!("mark_indexing_agreement_as_abandoning") } - async fn reopen_indexing_agreement_cancel(&self, _id: &IndexingAgreementId) -> Result<()> { + async fn reopen_indexing_agreement_cancel( + &self, + _id: &IndexingAgreementId, + _seen_live: bool, + ) -> Result<()> { unimplemented!("reopen_indexing_agreement_cancel") } @@ -446,8 +450,12 @@ impl AgreementRegistry for T { StubAgreementRegistry::mark_indexing_agreement_as_abandoning(self, id).await } - async fn reopen_indexing_agreement_cancel(&self, id: &IndexingAgreementId) -> Result<()> { - StubAgreementRegistry::reopen_indexing_agreement_cancel(self, id).await + async fn reopen_indexing_agreement_cancel( + &self, + id: &IndexingAgreementId, + seen_live: bool, + ) -> Result<()> { + StubAgreementRegistry::reopen_indexing_agreement_cancel(self, id, seen_live).await } async fn get_cancelling_agreements( 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 cde9677a..c6d5735c 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 @@ -82,6 +82,7 @@ mod tests { async fn reopen_indexing_agreement_cancel( &self, id: &IndexingAgreementId, + _seen_live: bool, ) -> crate::registry::Result<()> { self.reopened.lock().unwrap().push(*id); Ok(()) diff --git a/dipper-pgregistry/src/postgres.rs b/dipper-pgregistry/src/postgres.rs index 2a553614..4120268b 100644 --- a/dipper-pgregistry/src/postgres.rs +++ b/dipper-pgregistry/src/postgres.rs @@ -1008,10 +1008,13 @@ impl PgRegistry { /// 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. It counts /// as checked, since no cancel is sent with it, so the retry takes it on its next sweep rather - /// than waiting for one to be mined. + /// than waiting for one to be mined. When the chain was read and showed it live (`seen_live`), + /// the end on record, and any announcement of it, no longer stands, so both are cleared for + /// the end still to come; an unread chain leaves them, as the agreement may have ended. pub async fn reopen_indexing_agreement_cancel( &self, agreement_id: &IndexingAgreementId, + seen_live: bool, ) -> Result<(), Error> { let updated = sqlx::query( r#" @@ -1021,6 +1024,11 @@ impl PgRegistry { cancel_attempts = 0, cancel_checked_at = timezone('UTC', now()), ended_seen_at = NULL, + canceled_at = CASE WHEN $5::BOOLEAN THEN NULL ELSE canceled_at END, + canceled_by = CASE WHEN $5 THEN NULL ELSE canceled_by END, + canceled_tx = CASE WHEN $5 THEN NULL ELSE canceled_tx END, + terminated_event_emitted_at = + CASE WHEN $5 THEN NULL ELSE terminated_event_emitted_at END, updated_at = timezone('UTC', now()) WHERE id = $2 AND status IN ($3, $4) "#, @@ -1029,6 +1037,7 @@ impl PgRegistry { .bind(agreement_id) .bind(IndexingAgreementStatus::CanceledByRequester) .bind(IndexingAgreementStatus::Rejected) + .bind(seen_live) .execute(&self.pool) .await?; if updated.rows_affected() == 0 { diff --git a/dipper-pgregistry/tests/it_registry_postgres.rs b/dipper-pgregistry/tests/it_registry_postgres.rs index c0b5de0f..29aee3c7 100644 --- a/dipper-pgregistry/tests/it_registry_postgres.rs +++ b/dipper-pgregistry/tests/it_registry_postgres.rs @@ -3617,7 +3617,7 @@ async fn an_ended_agreement_found_live_on_chain_goes_back_to_cancelling() { ) .await .expect("Failed to run fixture"); - let registry = PgRegistry::new(db); + let registry = PgRegistry::new(db.clone()); 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]); @@ -3625,6 +3625,11 @@ async fn an_ended_agreement_found_live_on_chain_goes_back_to_cancelling() { .mark_indexing_agreement_as_cancelling(&ended) .await .expect("mark cancelling"); + // The withdrawal of its offer, recorded as its end before the offer landed after all. + registry + .record_cancel_audit(&ended, 1_700_000_000, "0xpayer", Some("0xwithdrawal")) + .await + .expect("cancel record"); assert_eq!( registry.record_cancel_check(&ended, 2, None).await.unwrap(), 2 @@ -3635,9 +3640,19 @@ async fn an_ended_agreement_found_live_on_chain_goes_back_to_cancelling() { .expect("mark ended"); registry - .reopen_indexing_agreement_cancel(&ended) + .reopen_indexing_agreement_cancel(&ended, true) .await .expect("an ended agreement can be reopened"); + let (canceled_tx,): (Option,) = + sqlx::query_as("SELECT canceled_tx FROM dipper_reg_indexing_agreements WHERE id = $1") + .bind(ended) + .fetch_one(&db) + .await + .expect("cancel record query"); + assert_eq!( + canceled_tx, None, + "the end on record no longer stands once the chain shows it live" + ); // Reopening sends no cancel, so there is none to wait on being mined. let listed = registry @@ -3646,7 +3661,9 @@ async fn an_ended_agreement_found_live_on_chain_goes_back_to_cancelling() { .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; + let still_wanted = registry + .reopen_indexing_agreement_cancel(&accepted, true) + .await; assert!( matches!(still_wanted, Err(Error::NoRecordsUpdated)), "got {still_wanted:?}" From 3483c253ed25679cc26cfa9785035ea8d3549e16 Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Mon, 5 Oct 2026 13:04:16 +0100 Subject: [PATCH 2/3] fix(cancel): count an agreement the listener already ended as ended If the chain listener marked an agreement ended between dipper sending its cancel and confirming it, the confirm found nothing to update, warned that it failed, and reported the agreement as still cancelling. It now counts it as ended and logs that at debug level. --- bin/dipper-service/src/cancel_dispatch.rs | 27 +++++++++++++------ .../src/network/service/cancel_retry.rs | 18 +++++++++++++ 2 files changed, 37 insertions(+), 8 deletions(-) diff --git a/bin/dipper-service/src/cancel_dispatch.rs b/bin/dipper-service/src/cancel_dispatch.rs index 0a3e8e83..54e50521 100644 --- a/bin/dipper-service/src/cancel_dispatch.rs +++ b/bin/dipper-service/src/cancel_dispatch.rs @@ -163,7 +163,8 @@ where /// Mark an agreement the chain shows dipper ended as ended, 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. +/// fails; it stays `Cancelling` for the cancel retry. One the chain listener already marked +/// ended counts as ended. pub async fn confirm_cancelled( registry: &R, agreement: &IndexingAgreement, @@ -175,16 +176,26 @@ pub async fn confirm_cancelled( if tx_hash.is_some() { record_cancel(registry, agreement, tx_hash, config).await; } - if let Err(err) = registry + match 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; the cancel retry tries again" - ); - return false; + Ok(()) => {} + Err(crate::registry::Error::NoRecordsUpdated) => { + tracing::debug!( + agreement_id = %agreement.id, + "Agreement already marked ended, as the chain listener can do first" + ); + return true; + } + Err(err) => { + tracing::warn!( + agreement_id = %agreement.id, + error = %err, + "Failed to mark an ended agreement cancelled; the cancel retry tries again" + ); + return false; + } } tracing::info!( agreement_id = %agreement.id, diff --git a/bin/dipper-service/src/network/service/cancel_retry.rs b/bin/dipper-service/src/network/service/cancel_retry.rs index ec9e0e63..63cf5d50 100644 --- a/bin/dipper-service/src/network/service/cancel_retry.rs +++ b/bin/dipper-service/src/network/service/cancel_retry.rs @@ -337,6 +337,8 @@ mod tests { attempts: AtomicU32, checks: AtomicU32, found_ended: Mutex>>, + /// The chain listener marks it ended before the retry's own mark lands. + listener_ended_it: bool, writes: Mutex>, } @@ -354,6 +356,9 @@ mod tests { &self, id: &IndexingAgreementId, ) -> crate::registry::Result<()> { + if self.listener_ended_it { + return Err(crate::registry::Error::NoRecordsUpdated); + } self.marked_cancelled.lock().unwrap().push(*id); self.writes.lock().unwrap().push("ended"); Ok(()) @@ -577,6 +582,19 @@ mod tests { ); } + #[tokio::test] + async fn counts_an_agreement_the_listener_marked_ended_first_as_ended() { + // Not a failure: there is nothing left to retry. + let registry = MockRegistry { + listener_ended_it: true, + ..registry_with_one(true) + }; + + retry(®istry, &live_chain(), 0).await; + + assert_eq!(registry.checks.load(Ordering::SeqCst), 0); + } + #[tokio::test] async fn leaves_an_accepted_agreement_that_already_ended_to_the_listener() { // The indexer may have ended it, or an earlier cancel whose result went unread; From 6b4b2b5f810798022d2b0b1d3cddb0bbb6bb85c4 Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Mon, 5 Oct 2026 13:07:21 +0100 Subject: [PATCH 3/3] fix(registry): replace an end recorded before the agreement's accept A reopen after a failed chain read keeps the end on record, and later cancels only filled blank fields, so a withdrawn offer's end could still be announced for an agreement accepted after it. An end recorded before the accept can't be the real one, so a later end now replaces it. --- dipper-pgregistry/src/postgres.rs | 15 ++++-- .../tests/it_registry_postgres.rs | 48 +++++++++++++++++++ 2 files changed, 60 insertions(+), 3 deletions(-) diff --git a/dipper-pgregistry/src/postgres.rs b/dipper-pgregistry/src/postgres.rs index 4120268b..23faca70 100644 --- a/dipper-pgregistry/src/postgres.rs +++ b/dipper-pgregistry/src/postgres.rs @@ -1553,6 +1553,8 @@ impl PgRegistry { /// emission sweep can populate the `terminated` event's tx/by/at fields. /// `COALESCE` keeps any value already observed on-chain. Best-effort /// enrichment: the event still emits (with fallbacks) if never recorded. + /// An end recorded before the agreement's accept, such as its offer's withdrawal before the + /// offer landed after all, can't be its end, so a later one replaces it and is announced. #[expect( clippy::cast_possible_wrap, reason = "predates this lint; fix when next touched" @@ -1567,9 +1569,16 @@ impl PgRegistry { sqlx::query( r#" UPDATE dipper_reg_indexing_agreements - SET canceled_at = COALESCE(canceled_at, $2), - canceled_by = COALESCE(canceled_by, $3), - canceled_tx = COALESCE(canceled_tx, $4) + SET canceled_at = CASE WHEN canceled_at < accepted_at AND $2 >= accepted_at + THEN $2 ELSE COALESCE(canceled_at, $2) END, + canceled_by = CASE WHEN canceled_at < accepted_at AND $2 >= accepted_at + THEN $3 ELSE COALESCE(canceled_by, $3) END, + canceled_tx = CASE WHEN canceled_at < accepted_at AND $2 >= accepted_at + THEN $4 ELSE COALESCE(canceled_tx, $4) END, + terminated_event_emitted_at = CASE + WHEN canceled_at < accepted_at AND $2 >= accepted_at THEN NULL + ELSE terminated_event_emitted_at + END WHERE id = $1 "#, ) diff --git a/dipper-pgregistry/tests/it_registry_postgres.rs b/dipper-pgregistry/tests/it_registry_postgres.rs index 29aee3c7..a246b2a3 100644 --- a/dipper-pgregistry/tests/it_registry_postgres.rs +++ b/dipper-pgregistry/tests/it_registry_postgres.rs @@ -3670,6 +3670,54 @@ async fn an_ended_agreement_found_live_on_chain_goes_back_to_cancelling() { ); } +#[tokio::test] +async fn an_end_recorded_before_the_accept_gives_way_to_the_real_one() { + // An offer withdrawn, then accepted when it landed after all: the withdrawal can't be the + // end of an agreement accepted later, however the agreement came back to be cancelled. + 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.clone()); + let id = fixture_agreement(0xaa); + let canceled_tx = async || -> Option { + let (tx,): (Option,) = + sqlx::query_as("SELECT canceled_tx FROM dipper_reg_indexing_agreements WHERE id = $1") + .bind(id) + .fetch_one(&db) + .await + .expect("cancel record query"); + tx + }; + registry + .record_cancel_audit(&id, 1_000, "0xpayer", Some("0xwithdrawal")) + .await + .expect("cancel record"); + registry + .record_accepted_audit(&id, 2_000, "0xaccept") + .await + .expect("accept record"); + + registry + .record_cancel_audit(&id, 3_000, "0xpayer", Some("0xcancel")) + .await + .expect("cancel record"); + assert_eq!(canceled_tx().await.as_deref(), Some("0xcancel")); + + registry + .record_cancel_audit(&id, 4_000, "0xpayer", Some("0xlater")) + .await + .expect("cancel record"); + assert_eq!( + canceled_tx().await.as_deref(), + Some("0xcancel"), + "a real end, after the accept, is kept" + ); +} + #[tokio::test] async fn a_cancelling_agreement_stays_live_and_unannounced_until_it_ends() { let (db, _temp_db) = temp_registry_db().await;