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
29 changes: 28 additions & 1 deletion bin/dipper-service/src/network/service/chain_listener.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
&registry,
&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();
Expand Down
11 changes: 10 additions & 1 deletion dipper-pgregistry/src/postgres.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand All @@ -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
"#,
Expand All @@ -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)
Expand Down
58 changes: 58 additions & 0 deletions dipper-pgregistry/tests/it_registry_postgres.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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"
);
}
Loading