Skip to content
22 changes: 13 additions & 9 deletions bin/dipper-service/src/cancel_dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -92,13 +92,17 @@ impl CancelReason {
}
}

/// What [`start_cancel`] left an agreement as.
/// What [`start_cancel`] left an agreement as. Unless `Ended`, it stays `Cancelling` for the
/// cancel retry to finish.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CancelStarted {
/// It was accepted and its cancel landed: now ended, as its [`CancelReason`] says.
Ended,
/// Still `Cancelling`; the chain listener finishes it once it can't go live.
Cancelling,
/// The chain shows nothing live, so nobody is being paid for it.
NotLive,
/// The chain couldn't be read, or its cancel failed or couldn't be confirmed, so it may
/// still be live and paid.
MayBeLive,
}

/// Start ending an agreement that may be live on-chain. It is marked `Cancelling` first, so an
Expand Down Expand Up @@ -129,18 +133,18 @@ where
}
let tx_hash = match cancel_if_live(chain_client, agreement, config).await {
LiveCancel::Ended(tx_hash) => tx_hash,
LiveCancel::NotLive { .. } => return Ok(CancelStarted::Cancelling),
LiveCancel::NotLive { .. } => return Ok(CancelStarted::NotLive),
LiveCancel::ReadFailed(err) | LiveCancel::CancelFailed(err) => {
tracing::warn!(
agreement_id = %agreement.id,
error = %err,
"On-chain cancel failed; the chain listener retries it"
"On-chain cancel failed; the cancel retry sends it again"
);
return Ok(CancelStarted::Cancelling);
return Ok(CancelStarted::MayBeLive);
}
LiveCancel::Unconfirmed { tx_hash, err } => {
log_unconfirmed(agreement, tx_hash, &err);
return Ok(CancelStarted::Cancelling);
return Ok(CancelStarted::MayBeLive);
}
};
tracing::info!(
Expand All @@ -150,13 +154,13 @@ where
);
// An offer never accepted could still land and be accepted until its deadline.
if agreement.status != IndexingAgreementStatus::AcceptedOnChain {
return Ok(CancelStarted::Cancelling);
return Ok(CancelStarted::NotLive);
}
Ok(
if confirm_cancelled(registry, agreement, reason, tx_hash, config).await {
CancelStarted::Ended
} else {
CancelStarted::Cancelling
CancelStarted::NotLive
},
)
}
Expand Down
16 changes: 16 additions & 0 deletions bin/dipper-service/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,7 @@ pub async fn main() -> anyhow::Result<()> {
let chain_listener_agreement_conf = agreement_conf.clone();
let liveness_agreement_conf = agreement_conf.clone();
let escrow_reconciler_agreement_conf = agreement_conf.clone();
let cancel_retry_agreement_conf = agreement_conf.clone();

// Canonical chain id and RecurringCollector address, read once and shared by the
// admin signer, the gRPC proposal signer, and the on-chain chain client so their
Expand Down Expand Up @@ -545,6 +546,15 @@ pub async fn main() -> anyhow::Result<()> {
_ => None,
};

//- The cancel retry, always on: it alone finishes the cancels dipper starts
let (cancel_retry_handle, cancel_retry_service) =
network::service::cancel_retry::new(network::service::cancel_retry::Ctx {
registry: registry.clone(),
chain_client: chain_client.clone(),
agreement_conf: cancel_retry_agreement_conf,
worker_queue: worker_handle.queue().clone(),
});

//- The liveness checker service (optional, enabled by config)
// Detects indexers who silently stop indexing active AcceptedOnChain agreements
let liveness_checker_handle = match conf.liveness_checker {
Expand Down Expand Up @@ -683,6 +693,9 @@ pub async fn main() -> anyhow::Result<()> {
None
};

let cancel_retry_task_handle = task_tree.spawn(cancel_retry_service);
tracing::debug!(task_id=%cancel_retry_task_handle.id(), "Cancel retry service started");

// Spawn the escrow reconciler service if enabled
let escrow_reconciler_stop_handle = if let Some((handle, service)) = escrow_reconciler_handle {
let task_handle = task_tree.spawn(service);
Expand Down Expand Up @@ -763,6 +776,9 @@ pub async fn main() -> anyhow::Result<()> {
all_stopped &= stop_service("Chain listener", handle.stop()).await;
}

// Stop the cancel retry before worker (it queues replacements)
all_stopped &= stop_service("Cancel retry", cancel_retry_handle.stop()).await;

// Stop escrow reconciler service before the DB pool closes
if let Some(handle) = escrow_reconciler_stop_handle {
all_stopped &= stop_service("Escrow reconciler", handle.stop()).await;
Expand Down
Loading
Loading