From ee6663f50301fda9ab79d63961a95f220a585ed1 Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Fri, 9 Oct 2026 18:56:18 +0100 Subject: [PATCH 1/2] fix(registry): hold an end's announcement until its transaction is known Dipper can mark an agreement ended before the chain listener records the transaction that ended it, and the terminated event then went out without it. The event now waits for that transaction, for up to an hour, after which it goes out without one as before. --- dipper-pgregistry/src/postgres.rs | 11 +++- .../tests/it_registry_postgres.rs | 58 +++++++++++++++++++ 2 files changed, 68 insertions(+), 1 deletion(-) diff --git a/dipper-pgregistry/src/postgres.rs b/dipper-pgregistry/src/postgres.rs index 3c3fdfd2..e9b7411e 100644 --- a/dipper-pgregistry/src/postgres.rs +++ b/dipper-pgregistry/src/postgres.rs @@ -212,6 +212,11 @@ impl sqlx::FromRow<'_, sqlx::postgres::PgRow> for CancellingAgreement { } } +/// How long an ended agreement's `terminated` event waits for the transaction that ended it. +/// Dipper can mark an agreement ended before the chain listener records that transaction, so +/// the wait lets the event carry it; after this the event goes out without one. +const TERMINATED_TX_WAIT_MINUTES: i32 = 60; + /// Statuses an on-chain cancel by dipper ends. const CANCEL_BY_REQUESTER_FROM: &[IndexingAgreementStatus] = &[ IndexingAgreementStatus::Created, @@ -1385,7 +1390,8 @@ impl PgRegistry { /// Fetch a batch of agreements awaiting a `terminated` event: in a /// terminal-cancel state, genuinely accepted on-chain (`accepted_at IS NOT /// NULL`, so a never-accepted local cancel is excluded), and not yet - /// emitted. Oldest-marked first so the backlog drains in order. + /// emitted, once the transaction that ended it is known or + /// [`TERMINATED_TX_WAIT_MINUTES`] have passed. Oldest-marked first. pub async fn get_agreements_pending_terminated_emission( &self, limit: i64, @@ -1397,6 +1403,8 @@ impl PgRegistry { WHERE status IN ($1, $2, $3) AND accepted_at IS NOT NULL AND terminated_event_emitted_at IS NULL + AND (canceled_tx IS NOT NULL + OR updated_at < timezone('UTC', now()) - make_interval(mins => $5)) ORDER BY updated_at ASC LIMIT $4 "#, @@ -1405,6 +1413,7 @@ impl PgRegistry { .bind(IndexingAgreementStatus::CanceledByIndexer) .bind(IndexingAgreementStatus::AbandonedByIndexer) .bind(limit) + .bind(TERMINATED_TX_WAIT_MINUTES) .fetch_all(&self.pool) .await?; Ok(rows) diff --git a/dipper-pgregistry/tests/it_registry_postgres.rs b/dipper-pgregistry/tests/it_registry_postgres.rs index 8937d3f1..d5fd3cef 100644 --- a/dipper-pgregistry/tests/it_registry_postgres.rs +++ b/dipper-pgregistry/tests/it_registry_postgres.rs @@ -3995,3 +3995,61 @@ async fn an_accept_and_end_from_the_chain_are_recorded_in_1_write() { "what is already known stays" ); } + +/// An ended agreement's `terminated` event waits for the transaction that ended it, which the +/// chain listener can record after dipper marks it ended, but not for ever. +#[tokio::test] +async fn an_end_is_announced_once_its_transaction_is_known_or_after_an_hour() { + 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 id = fixture_agreement(0xaa); + let registry = PgRegistry::new(db.clone()); + registry + .mark_indexing_agreement_as_canceled_by_requester(&id) + .await + .expect("dipper ends it"); + registry + .record_accepted_audit(&id, 1_700_000_000, "0xacc") + .await + .expect("its accept is known"); + let pending = async || { + registry + .get_agreements_pending_terminated_emission(100) + .await + .expect("terminated query") + .iter() + .any(|p| p.agreement_id == id) + }; + assert!(!pending().await, "waits for its transaction"); + + sqlx::query( + "UPDATE dipper_reg_indexing_agreements \ + SET updated_at = timezone('UTC', now()) - interval '61 minutes' WHERE id = $1", + ) + .bind(id) + .execute(&db) + .await + .expect("age it"); + assert!(pending().await, "goes out without one after an hour"); + + sqlx::query( + "UPDATE dipper_reg_indexing_agreements SET updated_at = timezone('UTC', now()) WHERE id = $1", + ) + .bind(id) + .execute(&db) + .await + .expect("make it recent again"); + registry + .record_cancel_audit(&id, 1_700_000_100, "0xmgr", Some("0xend")) + .await + .expect("its transaction is recorded"); + assert!( + pending().await, + "goes out as soon as its transaction is known" + ); +} From 311cfe8c6d47272a18ffbd5f316bc6c000b89b42 Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Fri, 9 Oct 2026 19:00:09 +0100 Subject: [PATCH 2/2] fix(listener): record the ending transaction of an abandoned agreement The listener backfilled the accept and ending transaction of agreements dipper had cancelled, but skipped ones it had ended as abandoned, so their terminated event waited an hour for a transaction that never arrived. Abandoned agreements now get the same backfill. --- .../src/network/service/chain_listener.rs | 29 ++++++++++++++++++- 1 file changed, 28 insertions(+), 1 deletion(-) diff --git a/bin/dipper-service/src/network/service/chain_listener.rs b/bin/dipper-service/src/network/service/chain_listener.rs index bfc83eb4..a3ce327c 100644 --- a/bin/dipper-service/src/network/service/chain_listener.rs +++ b/bin/dipper-service/src/network/service/chain_listener.rs @@ -881,7 +881,9 @@ where let already_terminal_cancel = matches!( agreement.status, - IndexingAgreementStatus::CanceledByRequester | IndexingAgreementStatus::CanceledByIndexer, + IndexingAgreementStatus::CanceledByRequester + | IndexingAgreementStatus::CanceledByIndexer + | IndexingAgreementStatus::AbandonedByIndexer, ); // Classify off the on-chain state, which carries the contract's own // "canceled by" flag. The canceler address is the payer, not dipper's @@ -3151,6 +3153,31 @@ mod tests { ); } + #[tokio::test] + async fn test_reconcile_records_the_end_of_an_abandoned_agreement() { + // Dipper can mark it ended without the transaction that ended it, and its terminated + // event waits for that transaction, which only this record supplies. + let registry = MockRegistry::new(); + let chain_client = MockChainClient::default(); + let agreement_id = IndexingAgreementId::from_bytes(rand::random()); + registry.add_agreement(agreement_id, IndexingAgreementStatus::AbandonedByIndexer); + + let snapshot = make_snapshot(agreement_id, AgreementState::CanceledByPayer, Address::ZERO); + reconcile_agreement( + &snapshot, + ®istry, + &chain_client, + test_agreement_conf().as_ref(), + ) + .await + .expect("reconcile ok"); + + assert_eq!( + registry.audit_writes(), + vec![("accept and cancel", agreement_id)] + ); + } + #[tokio::test] async fn test_reconcile_survives_a_failed_accept_and_cancel_record() { let registry = MockRegistry::new();