diff --git a/bin/dipper-service/src/cancel_dispatch.rs b/bin/dipper-service/src/cancel_dispatch.rs index 17c0eed2..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, @@ -232,20 +243,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 +646,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/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; 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..23faca70 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 { @@ -1544,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" @@ -1558,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 c0b5de0f..a246b2a3 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,13 +3661,63 @@ 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:?}" ); } +#[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;