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
Original file line number Diff line number Diff line change
Expand Up @@ -364,14 +364,6 @@ mod tests {
Ok(JobId::default())
}

async fn cancel_rejected_agreement_on_chain(
&self,
_agreement_id: IndexingAgreementId,
_priority: crate::worker::queue::JobPriority,
) -> anyhow::Result<JobId> {
unimplemented!()
}

async fn submit_offer(
&self,
_agreement_id: IndexingAgreementId,
Expand Down
95 changes: 79 additions & 16 deletions bin/dipper-service/src/cancel_dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -71,9 +71,9 @@ pub enum CancelStarted {
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.
/// Start ending an agreement that may be live on-chain. It is marked `Cancelling` first, so an
/// offer for it still in flight withdraws itself on landing, then cancelled only if the chain
/// shows it live. Fails, sending nothing, when the mark can't be written.
pub async fn start_cancel<R, T>(
registry: &R,
chain_client: &T,
Expand All @@ -87,9 +87,10 @@ where
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) => {
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::ReadFailed(err) | LiveCancel::CancelFailed(err) => {
tracing::warn!(
agreement_id = %agreement.id,
error = %err,
Expand All @@ -107,35 +108,52 @@ where
if agreement.status != IndexingAgreementStatus::AcceptedOnChain {
return Ok(CancelStarted::Cancelling);
}
Ok(confirm_cancelled(registry, agreement, tx_hash, config).await)
Ok(
if confirm_cancelled(registry, agreement, tx_hash, config).await {
CancelStarted::Ended
} else {
CancelStarted::Cancelling
},
)
}

/// 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>(
/// Mark an agreement the chain shows dipper ended `CanceledByRequester`, recording the cancel
/// when its transaction is known, so the `terminated` sweep announces it. False, logged, when
/// the mark fails; it stays `Cancelling` for the cancel retry.
pub async fn confirm_cancelled<R: AgreementRegistry + Sync>(
registry: &R,
agreement: &IndexingAgreement,
tx_hash: Option<B256>,
config: &IndexingAgreementConfig,
) -> CancelStarted {
) -> bool {
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"
"Failed to mark an ended agreement cancelled; the cancel retry tries again"
);
return CancelStarted::Cancelling;
return false;
}
record_cancel(registry, agreement, tx_hash, config).await;
CancelStarted::Ended
tracing::info!(
agreement_id = %agreement.id,
indexing_request_id = %agreement.indexing_request_id,
old_status = "CANCELLING",
new_status = "CANCELED_BY_REQUESTER",
reason = "cancel_confirmed_on_chain",
"agreement state transition"
);
if tx_hash.is_some() {
record_cancel(registry, agreement, tx_hash, config).await;
}
true
}

/// Record dipper's own cancel of an accepted agreement, so the `terminated` sweep
/// announces it.
pub async fn record_cancel<R: AgreementRegistry + Sync>(
async fn record_cancel<R: AgreementRegistry + Sync>(
registry: &R,
agreement: &IndexingAgreement,
tx_hash: Option<B256>,
Expand All @@ -155,6 +173,51 @@ pub async fn record_cancel<R: AgreementRegistry + Sync>(
}
}

/// Move an agreement dipper had already rejected or cancelled back into `Cancelling` when the
/// chain shows it live after all, so the cancel retry ends it. The chain is read first, so a
/// subgraph report from before dipper's cancel landed reopens nothing; an unreadable chain
/// reopens it anyway, as the retry reads again before sending. True if it was reopened.
pub async fn reopen_if_live<R, T>(
registry: &R,
chain_client: &T,
agreement: &IndexingAgreement,
) -> RegistryResult<bool>
where
R: AgreementRegistry + Sync,
T: ChainClient,
{
match chain_client
.agreement_still_active(agreement.id.as_bytes())
.await
{
Ok(false) => return Ok(false),
Ok(true) => {}
Err(err) => tracing::warn!(
agreement_id = %agreement.id,
error = %err,
"Failed to read an ended agreement reported live; the cancel retry checks it"
),
}
match registry
.reopen_indexing_agreement_cancel(&agreement.id)
.await
{
Ok(()) => {}
Err(crate::registry::Error::NoRecordsUpdated) => return Ok(false),
Err(err) => return Err(err),
}
tracing::warn!(
agreement_id = %agreement.id,
indexer_id = %agreement.indexer.id,
indexing_request_id = %agreement.indexing_request_id,
old_status = %agreement.status,
new_status = "CANCELLING",
reason = "live_on_chain_after_end",
"agreement state transition"
);
Ok(true)
}

/// What [`cancel_if_live`] found and did.
#[derive(Debug)]
pub enum LiveCancel {
Expand Down
1 change: 0 additions & 1 deletion bin/dipper-service/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -522,7 +522,6 @@ pub async fn main() -> anyhow::Result<()> {

let ctx = network::service::chain_listener::Ctx {
registry: registry.clone(),
worker_queue: worker_handle.queue().clone(),
event_source,
chain_client: chain_client.clone(),
agreement_conf: chain_listener_agreement_conf.clone(),
Expand Down
27 changes: 2 additions & 25 deletions bin/dipper-service/src/network/service/cancel_retry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ use dipper_core::time::now_secs;
use thegraph_core::alloy::primitives::B256;

use crate::{
cancel_dispatch::{LiveCancel, cancel_if_live, record_cancel},
cancel_dispatch::{LiveCancel, cancel_if_live, confirm_cancelled},
chain_client::{ChainClient, ChainClientError},
config::IndexingAgreementConfig,
registry::{AgreementRegistry, CancelKind, CancellingAgreement, IndexingAgreement},
Expand Down Expand Up @@ -154,30 +154,7 @@ where
Some(false) => {}
}
}
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 an ended agreement cancelled, will retry"
);
return false;
}
tracing::info!(
agreement_id = %agreement.id,
indexing_request_id = %agreement.indexing_request_id,
old_status = "CANCELLING",
new_status = "CANCELED_BY_REQUESTER",
reason = "cancel_confirmed_on_chain",
"agreement state transition"
);
// Without its own transaction, a late read by the chain listener fills in the cancel.
if row.accepted_on_chain && tx_hash.is_some() {
record_cancel(registry, agreement, tx_hash, config).await;
}
true
confirm_cancelled(registry, agreement, tx_hash, config).await
}

/// Whether the chain shows the indexer ended the agreement, or `None` when it can't be read,
Expand Down
Loading
Loading