diff --git a/bin/dipper-service/src/network/service/chain_listener.rs b/bin/dipper-service/src/network/service/chain_listener.rs index aa8b761e..bfc83eb4 100644 --- a/bin/dipper-service/src/network/service/chain_listener.rs +++ b/bin/dipper-service/src/network/service/chain_listener.rs @@ -931,10 +931,9 @@ where /// Record the accept and cancel of an agreement dipper had already marked /// cancelled that went live on-chain first, so its accepted and terminated -/// events go out. Cancel first: the terminated sweep waits only for the accept. +/// events go out. Both go in 1 write, as the terminated sweep waits only for the accept. /// Existing values win, so an agreement dipper already recorded is unchanged, -/// except an end recorded before the accept, such as its offer's withdrawal: the -/// cancel is recorded again once the accept is, so that end gives way to this one. +/// except an end recorded before the accept, such as its offer's withdrawal, which gives way. async fn record_accept_and_cancel_from_chain( snapshot: &AgreementStateSnapshot, agreement: &IndexingAgreement, @@ -943,36 +942,16 @@ async fn record_accept_and_cancel_from_chain( if snapshot.accepted_at == 0 { return; } - let canceled_by = snapshot.canceled_by.to_string(); - let recorded = match registry - .record_cancel_audit( + let recorded = registry + .record_accept_and_cancel_audit( &agreement.id, + snapshot.accepted_at, + &snapshot.accepted_tx, snapshot.canceled_at, - &canceled_by, + &snapshot.canceled_by.to_string(), Some(&snapshot.canceled_tx), ) - .await - { - Ok(()) => { - registry - .record_accepted_audit(&agreement.id, snapshot.accepted_at, &snapshot.accepted_tx) - .await - } - Err(err) => Err(err), - }; - let recorded = match recorded { - Ok(()) => { - registry - .record_cancel_audit( - &agreement.id, - snapshot.canceled_at, - &canceled_by, - Some(&snapshot.canceled_tx), - ) - .await - } - Err(err) => Err(err), - }; + .await; if let Err(err) = recorded { tracing::warn!( agreement_id = %agreement.id, @@ -1899,9 +1878,9 @@ mod tests { /// Ids passed to `record_cancel_audit` -- the signal a cancel path drives /// the terminated event (the sweep emits from this audit). recorded_cancel_audit: Vec, - /// Every audit write in order, as ("cancel" | "accept", id). + /// Every audit write in order, as ("cancel" | "accept" | "accept and cancel", id). audit_writes: Vec<(&'static str, IndexingAgreementId)>, - /// When true, `record_cancel_audit` fails. + /// When true, writes that record a cancel fail. fail_cancel_audit: bool, pending_cancellations: std::collections::HashMap< IndexingAgreementId, @@ -2219,6 +2198,25 @@ mod tests { Ok(()) } + async fn record_accept_and_cancel_audit( + &self, + agreement_id: &IndexingAgreementId, + _accepted_at: u64, + _accepted_tx: &str, + _canceled_at: u64, + _canceled_by: &str, + _canceled_tx: Option<&str>, + ) -> RegistryResult<()> { + let mut state = self.state.lock().unwrap(); + if state.fail_cancel_audit { + return Err(crate::registry::Error::NoRecordsUpdated); + } + state + .audit_writes + .push(("accept and cancel", *agreement_id)); + Ok(()) + } + async fn record_accepted_audit( &self, agreement_id: &IndexingAgreementId, @@ -3128,8 +3126,7 @@ mod tests { async fn test_reconcile_records_accept_and_cancel_of_cancelled_agreement_that_went_live() { // Dipper had marked the agreement cancelled, but it was accepted on-chain // before being ended there. Recording both lets the accepted and terminated - // events go out; the cancel goes first because the terminated sweep only - // waits for the accept, and again after it, to replace an end from before it. + // events go out, in 1 write, as the terminated sweep only waits for the accept. let registry = MockRegistry::new(); let chain_client = MockChainClient::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); @@ -3150,18 +3147,12 @@ mod tests { assert!(result.is_ok()); assert_eq!( registry.audit_writes(), - vec![ - ("cancel", agreement_id), - ("accept", agreement_id), - ("cancel", agreement_id) - ] + vec![("accept and cancel", agreement_id)] ); } #[tokio::test] - async fn test_reconcile_records_no_accept_when_the_cancel_record_fails() { - // An accept recorded without its cancel would let the terminated event go - // out with fallback cancel fields. + async fn test_reconcile_survives_a_failed_accept_and_cancel_record() { let registry = MockRegistry::new(); let chain_client = MockChainClient::default(); let agreement_id = IndexingAgreementId::from_bytes(rand::random()); diff --git a/bin/dipper-service/src/registry.rs b/bin/dipper-service/src/registry.rs index 86ac2a7f..71470cda 100644 --- a/bin/dipper-service/src/registry.rs +++ b/bin/dipper-service/src/registry.rs @@ -626,6 +626,28 @@ impl AgreementRegistry for RegistryProvider { Ok(()) } + async fn record_accept_and_cancel_audit( + &self, + agreement_id: &IndexingAgreementId, + accepted_at: u64, + accepted_tx: &str, + canceled_at: u64, + canceled_by: &str, + canceled_tx: Option<&str>, + ) -> RegistryResult<()> { + self.inner + .record_accept_and_cancel_audit( + agreement_id, + accepted_at, + accepted_tx, + canceled_at, + canceled_by, + canceled_tx, + ) + .await?; + Ok(()) + } + async fn get_expired_created_agreements( &self, batch_size: i64, diff --git a/bin/dipper-service/src/registry/agreement.rs b/bin/dipper-service/src/registry/agreement.rs index 10cc02f8..7f8618b1 100644 --- a/bin/dipper-service/src/registry/agreement.rs +++ b/bin/dipper-service/src/registry/agreement.rs @@ -427,6 +427,21 @@ pub trait AgreementRegistry { Ok(()) } + /// Record an agreement's accept and its end together, in 1 write, so nothing reads one + /// without the other. Default no-op so mocks need not override. + #[allow(clippy::too_many_arguments)] + async fn record_accept_and_cancel_audit( + &self, + _agreement_id: &IndexingAgreementId, + _accepted_at: u64, + _accepted_tx: &str, + _canceled_at: u64, + _canceled_by: &str, + _canceled_tx: Option<&str>, + ) -> RegistryResult<()> { + Ok(()) + } + /// Record the cancel audit payload for a dipper-initiated cancel so the /// emission sweep can populate the `terminated` event fields. Default no-op /// so mocks need not override. diff --git a/bin/dipper-service/src/registry/agreement_stub.rs b/bin/dipper-service/src/registry/agreement_stub.rs index cc3c717f..dd2950b4 100644 --- a/bin/dipper-service/src/registry/agreement_stub.rs +++ b/bin/dipper-service/src/registry/agreement_stub.rs @@ -327,6 +327,19 @@ pub trait StubAgreementRegistry: Send + Sync { ) -> Result<()> { Ok(()) } + + #[allow(clippy::too_many_arguments)] + async fn record_accept_and_cancel_audit( + &self, + _agreement_id: &IndexingAgreementId, + _accepted_at: u64, + _accepted_tx: &str, + _canceled_at: u64, + _canceled_by: &str, + _canceled_tx: Option<&str>, + ) -> Result<()> { + Ok(()) + } } // Every stub is a full AgreementRegistry: each method delegates to the stub @@ -650,4 +663,25 @@ impl AgreementRegistry for T { ) .await } + + async fn record_accept_and_cancel_audit( + &self, + agreement_id: &IndexingAgreementId, + accepted_at: u64, + accepted_tx: &str, + canceled_at: u64, + canceled_by: &str, + canceled_tx: Option<&str>, + ) -> Result<()> { + StubAgreementRegistry::record_accept_and_cancel_audit( + self, + agreement_id, + accepted_at, + accepted_tx, + canceled_at, + canceled_by, + canceled_tx, + ) + .await + } } diff --git a/dipper-pgregistry/src/postgres.rs b/dipper-pgregistry/src/postgres.rs index d2bce6bd..3c3fdfd2 100644 --- a/dipper-pgregistry/src/postgres.rs +++ b/dipper-pgregistry/src/postgres.rs @@ -1591,6 +1591,53 @@ impl PgRegistry { Ok(()) } + /// Record an agreement's accept and its end together, as the chain shows them, in 1 write, + /// with the rules of [`Self::record_accepted_audit`] and [`Self::record_cancel_audit`]. An + /// end recorded before the accept is judged against the accept being recorded with it. + #[expect( + clippy::cast_possible_wrap, + reason = "chain timestamps are far below i64::MAX" + )] + pub async fn record_accept_and_cancel_audit( + &self, + agreement_id: &IndexingAgreementId, + accepted_at: u64, + accepted_tx: &str, + canceled_at: u64, + canceled_by: &str, + canceled_tx: Option<&str>, + ) -> Result<(), Error> { + sqlx::query( + r#" + UPDATE dipper_reg_indexing_agreements + SET accepted_at = COALESCE(accepted_at, $2), + accepted_tx = COALESCE(accepted_tx, $3), + canceled_at = CASE + WHEN canceled_at < COALESCE(accepted_at, $2) AND $4 >= COALESCE(accepted_at, $2) + THEN $4 ELSE COALESCE(canceled_at, $4) END, + canceled_by = CASE + WHEN canceled_at < COALESCE(accepted_at, $2) AND $4 >= COALESCE(accepted_at, $2) + THEN $5 ELSE COALESCE(canceled_by, $5) END, + canceled_tx = CASE + WHEN canceled_at < COALESCE(accepted_at, $2) AND $4 >= COALESCE(accepted_at, $2) + THEN $6 ELSE COALESCE(canceled_tx, $6) END, + terminated_event_emitted_at = CASE + WHEN canceled_at < COALESCE(accepted_at, $2) AND $4 >= COALESCE(accepted_at, $2) + THEN NULL ELSE terminated_event_emitted_at END + WHERE id = $1 + "#, + ) + .bind(agreement_id) + .bind(accepted_at as i64) + .bind(accepted_tx) + .bind(canceled_at as i64) + .bind(canceled_by) + .bind(canceled_tx) + .execute(&self.pool) + .await?; + Ok(()) + } + // ========================================================================= // Reassignment operations // ========================================================================= diff --git a/dipper-pgregistry/tests/it_registry_postgres.rs b/dipper-pgregistry/tests/it_registry_postgres.rs index 88b993ef..8937d3f1 100644 --- a/dipper-pgregistry/tests/it_registry_postgres.rs +++ b/dipper-pgregistry/tests/it_registry_postgres.rs @@ -3934,3 +3934,64 @@ async fn an_abandoned_agreement_awaits_replacement_once_ended_until_queued() { .expect("note it queued"); assert!(awaiting().await.is_empty()); } + +/// The accept and end of an agreement replayed from the chain go in together: an end already +/// known stays, unless it came before the accept, and then the replayed end is announced. +#[tokio::test] +async fn an_accept_and_end_from_the_chain_are_recorded_in_1_write() { + 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()); + let record = async |accepted_at: u64, canceled_at: u64, tx: &str| { + registry + .record_accept_and_cancel_audit( + &id, + accepted_at, + "0xacc", + canceled_at, + "payer", + Some(tx), + ) + .await + .expect("record"); + }; + let recorded = async || { + sqlx::query_as::<_, (Option, Option, Option, bool)>( + "SELECT accepted_at, canceled_at, canceled_tx, terminated_event_emitted_at IS NULL \ + FROM dipper_reg_indexing_agreements WHERE id = $1", + ) + .bind(id) + .fetch_one(&db) + .await + .expect("read") + }; + + // Its offer was withdrawn at 50, already announced, before an accept at 100 came to light. + sqlx::query( + "UPDATE dipper_reg_indexing_agreements SET canceled_at = 50, canceled_tx = '0xwd', \ + terminated_event_emitted_at = now() WHERE id = $1", + ) + .bind(id) + .execute(&db) + .await + .expect("withdraw"); + record(100, 200, "0xend").await; + assert_eq!( + recorded().await, + (Some(100), Some(200), Some("0xend".to_owned()), true), + "the end before the accept gives way, to be announced again" + ); + + record(300, 400, "0xlater").await; + assert_eq!( + recorded().await, + (Some(100), Some(200), Some("0xend".to_owned()), true), + "what is already known stays" + ); +}