Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
e8beaba
feat(agreements): keep retrying cancels until the chain confirms them
MoonBoi9001 Oct 2, 2026
443641e
refactor(worker): share the read-then-cancel step between cancel paths
MoonBoi9001 Oct 2, 2026
3456774
fix(worker): retry an offer job that can't recheck its agreement
MoonBoi9001 Oct 2, 2026
e2a82b9
fix(listener): announce a missed accept however late it is read
MoonBoi9001 Oct 2, 2026
103a974
test(config): build every test's agreement config from one helper
MoonBoi9001 Oct 2, 2026
075f8d8
docs(worker): write job counts as numerals in the cancel job's comments
MoonBoi9001 Oct 2, 2026
3374a09
fix(listener): record no accept for a withdrawn offer being cancelled
MoonBoi9001 Oct 2, 2026
ac6470f
fix(listener): let the listener confirm an accepted agreement that ended
MoonBoi9001 Oct 2, 2026
4d75ad9
fix(listener): retry cancels fairly, briefly and never twice at once
MoonBoi9001 Oct 2, 2026
cd84305
fix(listener): count only cancels that are mined without effect
MoonBoi9001 Oct 2, 2026
340e399
fix(worker): finish an offer job whose recheck fails without resending
MoonBoi9001 Oct 2, 2026
b8daeac
fix(listener): retry cancels on a timer, not by keeping polls fast
MoonBoi9001 Oct 2, 2026
556d300
fix(registry): count fees of agreements still being cancelled
MoonBoi9001 Oct 2, 2026
0de1b2b
fix(registry): cancel a replaced agreement marked expired as well
MoonBoi9001 Oct 2, 2026
f106463
fix(listener): start orphan cancels like any other, with a retry limit
MoonBoi9001 Oct 2, 2026
ddb651c
refactor(registry): require every registry to implement the cancel retry
MoonBoi9001 Oct 2, 2026
2774efb
refactor(registry): stop reading a cancel count nothing uses
MoonBoi9001 Oct 2, 2026
f7632df
fix(listener): never confirm a cancel the chain could not be read for
MoonBoi9001 Oct 2, 2026
8c8c78a
fix(listener): record no accept of an old agreement being cancelled
MoonBoi9001 Oct 2, 2026
3ad96af
fix(listener): cancel a replaced expired agreement only if it is live
MoonBoi9001 Oct 2, 2026
735ea06
fix(listener): limit retries of a cancel that is mined and reverts
MoonBoi9001 Oct 2, 2026
280f0c1
fix(listener): mark an ended accepted agreement after an hour's wait
MoonBoi9001 Oct 2, 2026
c3e5bda
fix(listener): stop a cancel retry sweep after 30 seconds
MoonBoi9001 Oct 2, 2026
7d662d0
docs(worker): say only the chain listener queues the on-chain cancel job
MoonBoi9001 Oct 2, 2026
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
Original file line number Diff line number Diff line change
Expand Up @@ -135,5 +135,6 @@ fn into_indexing_agreement_status(
IndexingAgreementRecordStatus::AbandonedByIndexer => {
IndexingAgreementStatus::AbandonedByIndexer
}
IndexingAgreementRecordStatus::Cancelling => IndexingAgreementStatus::Cancelling,
}
}
159 changes: 137 additions & 22 deletions bin/dipper-service/src/cancel_dispatch.rs
Original file line number Diff line number Diff line change
@@ -1,12 +1,15 @@
//! On-chain cancel dispatch. Every cancel goes through
//! [`cancel_agreement_on_chain`] so the manager-routed path lives in one place.

use dipper_core::time::now_secs;
use thegraph_core::alloy::primitives::B256;

use crate::{
chain_client::{ChainClient, ChainClientError},
config::IndexingAgreementConfig,
registry::IndexingAgreement,
registry::{
AgreementRegistry, IndexingAgreement, IndexingAgreementStatus, Result as RegistryResult,
},
};

/// Pass both ACTIVE and PENDING; local status lags the chain, so let the
Expand Down Expand Up @@ -59,8 +62,135 @@ pub async fn cancel_agreement_on_chain<T: ChainClient>(
Ok(outcome)
}

/// What [`start_cancel`] left an agreement as.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CancelStarted {
/// It was accepted and its cancel landed: now `CanceledByRequester`.
Ended,
/// Still `Cancelling`; the chain listener finishes it once it can't go live.
Cancelling,
}

/// Start ending an agreement that may be live on-chain. It is marked `Cancelling` before
/// its cancel goes out, so an offer for it still in flight withdraws itself on landing.
/// Fails, sending nothing, when the mark can't be written.
pub async fn start_cancel<R, T>(
registry: &R,
chain_client: &T,
agreement: &IndexingAgreement,
config: &IndexingAgreementConfig,
) -> RegistryResult<CancelStarted>
where
R: AgreementRegistry + Sync,
T: ChainClient,
{
registry
.mark_indexing_agreement_as_cancelling(&agreement.id)
.await?;
let tx_hash = match cancel_agreement_on_chain(chain_client, agreement, config).await {
Ok(tx_hash) => tx_hash,
Err(err) => {
tracing::warn!(
agreement_id = %agreement.id,
error = %err,
"On-chain cancel failed; the chain listener retries it"
);
return Ok(CancelStarted::Cancelling);
}
};
tracing::info!(
agreement_id = %agreement.id,
tx_hash = ?tx_hash,
"Submitted on-chain cancellation"
);
// An offer never accepted could still land and be accepted until its deadline.
if agreement.status != IndexingAgreementStatus::AcceptedOnChain {
return Ok(CancelStarted::Cancelling);
}
Ok(confirm_cancelled(registry, agreement, tx_hash, config).await)
}

/// Mark an accepted agreement whose cancel landed `CanceledByRequester` and record the
/// cancel, so the `terminated` sweep announces it.
async fn confirm_cancelled<R: AgreementRegistry + Sync>(
registry: &R,
agreement: &IndexingAgreement,
tx_hash: Option<B256>,
config: &IndexingAgreementConfig,
) -> CancelStarted {
if let Err(err) = registry
.mark_indexing_agreement_as_canceled_by_requester(&agreement.id)
.await
{
tracing::warn!(
agreement_id = %agreement.id,
error = %err,
"Failed to mark a cancelled agreement; the chain listener finishes it"
);
return CancelStarted::Cancelling;
}
record_cancel(registry, agreement, tx_hash, config).await;
CancelStarted::Ended
}

/// Record dipper's own cancel of an accepted agreement, so the `terminated` sweep
/// announces it.
pub async fn record_cancel<R: AgreementRegistry + Sync>(
registry: &R,
agreement: &IndexingAgreement,
tx_hash: Option<B256>,
config: &IndexingAgreementConfig,
) {
let manager = config.recurring_agreement_manager().to_string();
let tx = tx_hash.map(|hash| hash.to_string());
if let Err(err) = registry
.record_cancel_audit(&agreement.id, now_secs(), &manager, tx.as_deref())
.await
{
tracing::warn!(
agreement_id = %agreement.id,
error = %err,
"failed to record cancel audit; terminated event may emit with fallback fields"
);
}
}

/// What [`cancel_if_live`] found and did.
#[derive(Debug)]
pub enum LiveCancel {
/// The chain showed nothing live, so no cancel was sent.
NotLive,
/// A cancel went out and the chain confirmed the agreement ended.
Ended(Option<B256>),
/// The chain could not be read, so nothing was sent.
ReadFailed(ChainClientError),
/// The cancel failed or did not end the agreement.
CancelFailed(ChainClientError),
}

/// Cancel an agreement on-chain only if the chain shows it live: a pending offer, or
/// accepted and not yet ended. Reading first saves a wasted transaction, since a cancel
/// of an agreement that already ended still mines.
pub async fn cancel_if_live<T: ChainClient>(
chain_client: &T,
agreement: &IndexingAgreement,
config: &IndexingAgreementConfig,
) -> LiveCancel {
match chain_client
.agreement_still_active(agreement.id.as_bytes())
.await
{
Err(err) => LiveCancel::ReadFailed(err),
Ok(false) => LiveCancel::NotLive,
Ok(true) => match cancel_agreement_on_chain(chain_client, agreement, config).await {
Ok(tx_hash) => LiveCancel::Ended(tx_hash),
Err(err) => LiveCancel::CancelFailed(err),
},
}
}

#[cfg(test)]
mod tests {
pub(crate) mod tests {
use std::sync::Mutex;

use async_trait::async_trait;
Expand Down Expand Up @@ -154,31 +284,16 @@ mod tests {

fn manager_conf(collector: Address) -> IndexingAgreementConfig {
IndexingAgreementConfig {
data_service: Address::ZERO,
recurring_collector: collector,
recurring_agreement_manager: Address::repeat_byte(0x33),
max_agreement_grt_per_30_days: 0.0,
max_seconds_per_collection: 0,
min_seconds_per_collection: 0,
duration_seconds: 0,
deadline_seconds: 0,
max_grt_per_30_days: std::collections::BTreeMap::new(),
max_grt_per_billion_entities_per_30_days: 0.0,
declined_indexer_lookback_days: 0,
price_rejection_lookback_days: 0,
transient_rejection_lookback_minutes: 0,
uncertain_rejection_lookback_days: 0,
unresponsive_indexer_lookback_days: 0,
mass_unresponsive_trip_fraction: 0.5,
mass_unresponsive_reset_fraction: 0.25,
dips_accepting_snapshot_max_age_hours: 48,
dips_accepting_cache_ttl_seconds: 300,
max_in_flight_offers_per_indexer: None,
max_in_flight_offers_total: None,
..IndexingAgreementConfig::for_tests()
}
}

fn agreement(status: IndexingAgreementStatus, hash: Option<Vec<u8>>) -> IndexingAgreement {
pub(crate) fn agreement(
status: IndexingAgreementStatus,
hash: Option<Vec<u8>>,
) -> IndexingAgreement {
let deployment_id: DeploymentId = "QmTXzATwNfgGVukV1fX2T6xw9f6LAYRVWpsdXyRWzUR2H9"
.parse()
.unwrap();
Expand Down
31 changes: 31 additions & 0 deletions bin/dipper-service/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1202,6 +1202,37 @@ pub struct IndexingAgreementConfig {
pub max_in_flight_offers_total: Option<u32>,
}

#[cfg(test)]
impl IndexingAgreementConfig {
/// Zero addresses and limits, with permissive breaker and cache settings, for tests
/// to adjust the fields they care about.
pub fn for_tests() -> Self {
Self {
data_service: Address::ZERO,
recurring_collector: Address::ZERO,
recurring_agreement_manager: Address::ZERO,
max_agreement_grt_per_30_days: 0.0,
max_seconds_per_collection: 0,
min_seconds_per_collection: 0,
duration_seconds: 0,
deadline_seconds: 0,
max_grt_per_30_days: BTreeMap::new(),
max_grt_per_billion_entities_per_30_days: 0.0,
declined_indexer_lookback_days: 0,
price_rejection_lookback_days: 0,
transient_rejection_lookback_minutes: 0,
uncertain_rejection_lookback_days: 0,
unresponsive_indexer_lookback_days: 0,
mass_unresponsive_trip_fraction: 0.5,
mass_unresponsive_reset_fraction: 0.25,
dips_accepting_snapshot_max_age_hours: 48,
dips_accepting_cache_ttl_seconds: 300,
max_in_flight_offers_per_indexer: None,
max_in_flight_offers_total: None,
}
}
}

/// Per-chain pricing for indexing agreements (runtime).
#[derive(Debug)]
pub struct IndexingAgreementChainPrices {
Expand Down
1 change: 1 addition & 0 deletions bin/dipper-service/src/network/service.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
pub mod cancel_retry;
pub mod chain_events;
pub mod chain_listener;
pub mod domain_refresh;
Expand Down
Loading
Loading