Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
73 changes: 32 additions & 41 deletions bin/dipper-service/src/network/service/chain_listener.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<R: AgreementRegistry + Sync>(
snapshot: &AgreementStateSnapshot,
agreement: &IndexingAgreement,
Expand All @@ -943,36 +942,16 @@ async fn record_accept_and_cancel_from_chain<R: AgreementRegistry + Sync>(
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,
Expand Down Expand Up @@ -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<IndexingAgreementId>,
/// 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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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());
Expand All @@ -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());
Expand Down
22 changes: 22 additions & 0 deletions bin/dipper-service/src/registry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
15 changes: 15 additions & 0 deletions bin/dipper-service/src/registry/agreement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
34 changes: 34 additions & 0 deletions bin/dipper-service/src/registry/agreement_stub.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -650,4 +663,25 @@ impl<T: StubAgreementRegistry> 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
}
}
47 changes: 47 additions & 0 deletions dipper-pgregistry/src/postgres.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
// =========================================================================
Expand Down
61 changes: 61 additions & 0 deletions dipper-pgregistry/tests/it_registry_postgres.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<i64>, Option<i64>, Option<String>, 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"
);
}
Loading