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
74 changes: 48 additions & 26 deletions bin/dipper-service/src/network/service/cancel_retry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,8 +12,8 @@ use crate::{
registry::{AgreementRegistry, CancelKind, CancellingAgreement, IndexingAgreement},
};

/// Retried cancels mined without ending an agreement before dipper stops retrying it and
/// leaves it to an operator. Other failures don't count (see `failed_attempts`).
/// Failed cancels before dipper alerts an operator and retries the agreement only hourly, so a
/// paused manager recovers once unpaused. Outages don't count (see `failed_attempts`).
pub const MAX_CANCEL_ATTEMPTS: u32 = 10;

/// Agreements a sweep takes on, those that may be paying an indexer first; the time budget
Expand Down Expand Up @@ -250,7 +250,7 @@ async fn note_check<R: AgreementRegistry + Sync>(
{
Ok(attempts) => {
if let Some(err) = failure.filter(|_| failed_attempts > 0) {
log_failed_cancel(agreement, attempts, err);
log_failed_cancel(agreement, attempts, failed_attempts, err);
}
}
Err(err) => tracing::warn!(
Expand All @@ -261,26 +261,25 @@ async fn note_check<R: AgreementRegistry + Sync>(
}
}

/// A revert before sending is logged as an ERROR: a paused or misconfigured manager makes
/// every cancel revert, so it is retried rather than counted against the agreement.
fn log_uncounted_failure(agreement: &IndexingAgreement, err: &ChainClientError) {
if matches!(err, ChainClientError::ContractRevert { .. }) {
tracing::error!(
agreement_id = %agreement.id,
error = %err,
"Cancel of an agreement reverted before it was sent, will retry"
);
} else {
tracing::warn!(
agreement_id = %agreement.id,
error = %err,
"Cancel of an agreement failed or could not be confirmed, will retry"
);
}
tracing::warn!(
agreement_id = %agreement.id,
error = %err,
"Cancel of an agreement failed or could not be confirmed, will retry"
);
}

fn log_failed_cancel(agreement: &IndexingAgreement, attempts: u32, err: &ChainClientError) {
if attempts < MAX_CANCEL_ATTEMPTS {
/// One ERROR as an agreement reaches the limit, for an operator to look into; a WARN for
/// every other failed cancel.
fn log_failed_cancel(
agreement: &IndexingAgreement,
attempts: u32,
failed: u32,
err: &ChainClientError,
) {
let reached_limit =
attempts >= MAX_CANCEL_ATTEMPTS && attempts.saturating_sub(failed) < MAX_CANCEL_ATTEMPTS;
if !reached_limit {
tracing::warn!(
agreement_id = %agreement.id,
attempts,
Expand All @@ -290,22 +289,24 @@ fn log_failed_cancel(agreement: &IndexingAgreement, attempts: u32, err: &ChainCl
return;
}
tracing::error!(
event = "agreement_cancel_abandoned",
event = "agreement_cancel_stuck",
agreement_id = %agreement.id,
indexer_id = %agreement.indexer.id,
indexing_request_id = %agreement.indexing_request_id,
attempts,
error = %err,
"Gave up cancelling an agreement on-chain; it may still be live"
"Cancelling an agreement keeps failing; it may still be live. Dipper now retries it hourly"
);
}

/// How many of an agreement's cancel attempts a failure uses up. A cancel mined without
/// ending it, or reverted, counts, and one that can never be sent uses them all. An
/// unreachable chain or a revert before sending is retried freely.
/// How many of an agreement's cancel attempts a failure uses up. A cancel the contract
/// refused, before sending or once mined, or that mined without ending the agreement counts,
/// and one that can never be sent uses them all. An unreachable chain is retried freely.
fn failed_attempts(err: &ChainClientError) -> u32 {
match err {
ChainClientError::CancelNotConfirmed { .. } | ChainClientError::TxReverted { .. } => 1,
ChainClientError::CancelNotConfirmed { .. }
| ChainClientError::TxReverted { .. }
| ChainClientError::ContractRevert { .. } => 1,
ChainClientError::MissingTermsVersionHash { .. } => MAX_CANCEL_ATTEMPTS,
_ => 0,
}
Expand Down Expand Up @@ -403,6 +404,7 @@ mod tests {
read_fails: bool,
send_fails: bool,
mined_cancel_reverts: bool,
reverts_before_sending: bool,
cancel_has_no_effect: bool,
clock_fails: bool,
now: AtomicU64,
Expand All @@ -429,6 +431,12 @@ mod tests {
if self.send_fails {
return Err(ChainClientError::RpcError(anyhow::anyhow!("rpc down")));
}
if self.reverts_before_sending {
return Err(ChainClientError::ContractRevert {
selector: [0xde, 0xad, 0xbe, 0xef],
data: Default::default(),
});
}
if self.mined_cancel_reverts {
return Err(ChainClientError::TxReverted {
tx_hash: B256::repeat_byte(0xee),
Expand Down Expand Up @@ -721,6 +729,20 @@ mod tests {
assert_eq!(registry.attempts.load(Ordering::SeqCst), 1);
}

#[tokio::test]
async fn counts_a_cancel_the_contract_refuses_before_it_is_sent() {
// Otherwise one that always reverts is retried, and alerted on, for ever.
let registry = registry_with_one(true);
let chain = MockChain {
reverts_before_sending: true,
..live_chain()
};

retry(&registry, &chain, 0).await;

assert_eq!(registry.attempts.load(Ordering::SeqCst), 1);
}

#[tokio::test]
async fn counts_a_cancel_that_mines_without_ending_the_agreement() {
let registry = registry_with_one(true);
Expand Down
4 changes: 2 additions & 2 deletions bin/dipper-service/src/registry/agreement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -281,8 +281,8 @@ pub trait AgreementRegistry {
id: &IndexingAgreementId,
) -> RegistryResult<()>;

/// `CANCELLING` agreements marked over `min_age_minutes` ago whose cancel has failed
/// fewer than `max_attempts` times, those checked longest ago first. One that may be paying
/// `CANCELLING` agreements marked over `min_age_minutes` ago, those checked longest ago
/// first; one whose cancel has failed `max_attempts` times only once an hour. One that may be paying
/// an indexer (accepted, or past the offer deadline, which only an accepted one outlives)
/// counts as checked an hour earlier, so it goes first without holding the rest back.
async fn get_cancelling_agreements(
Expand Down
9 changes: 6 additions & 3 deletions dipper-pgregistry/src/postgres.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1000,8 +1000,8 @@ impl PgRegistry {
Ok(())
}

/// `Cancelling` agreements marked over `min_age_minutes` ago whose cancel has failed
/// fewer than `max_attempts` times, those checked longest ago first. One that may be paying
/// `Cancelling` agreements marked over `min_age_minutes` ago, those checked longest ago
/// first; one whose cancel has failed `max_attempts` times only once an hour. One that may be paying
/// an indexer (accepted, or past the offer deadline, which only an accepted one outlives)
/// counts as checked an hour earlier, so it goes first without holding the rest back.
pub async fn get_cancelling_agreements(
Expand Down Expand Up @@ -1030,7 +1030,10 @@ impl PgRegistry {
accepted_at
FROM dipper_reg_indexing_agreements
WHERE status = $1
AND cancel_attempts < $2
AND (
cancel_attempts < $2
OR cancel_checked_at < timezone('UTC', now()) - INTERVAL '1 hour'
)
AND updated_at < timezone('UTC', now()) - make_interval(mins => $4)
ORDER BY
cancel_checked_at - CASE
Expand Down
19 changes: 18 additions & 1 deletion dipper-pgregistry/tests/it_registry_postgres.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3473,7 +3473,24 @@ async fn cancelling_agreements_are_listed_until_their_cancel_fails_too_often() {
assert_eq!(
ids,
vec![accepted],
"one that failed too often is left alone"
"one that failed too often waits an hour between checks"
);

sqlx::query(
"UPDATE dipper_reg_indexing_agreements \
SET cancel_checked_at = cancel_checked_at - INTERVAL '2 hours' WHERE id = $1",
)
.bind(created)
.execute(&db)
.await
.expect("Failed to age the check");
let listed = registry
.get_cancelling_agreements(100, 2, 0)
.await
.expect("cancelling query");
assert!(
listed.iter().any(|row| row.agreement.id == created),
"and is tried again after it"
);

let not_cancelling = registry
Expand Down
Loading