Skip to content
49 changes: 27 additions & 22 deletions bin/dipper-service/src/network/service/chain_listener.rs
Original file line number Diff line number Diff line change
Expand Up @@ -881,8 +881,9 @@ where
}
}

reopen_if_cancelled_but_accepted(snapshot, &agreement, registry, chain_client).await?;
record_accept_of_cancelling(snapshot, &agreement, registry).await;
let reopened =
reopen_if_cancelled_but_accepted(snapshot, &agreement, registry, chain_client).await?;
record_accept_of_cancelling(snapshot, &agreement, reopened, registry).await;

// Both transitions are applied atomically downstream so the
// Accept-then-Cancel-in-one-snapshot path can't leak an intermediate
Expand Down Expand Up @@ -998,13 +999,14 @@ fn created_after_events_started(agreement: &IndexingAgreement) -> bool {
/// Safety net for an agreement dipper cancelled whose offer the indexer accepted
/// anyway, such as one that landed after dipper's cancel. Nothing else would end
/// it: reconciliation ignores an accept on a cancelled row. It goes back to
/// `Cancelling` for the cancel retry, unless the chain shows it already ended.
/// `Cancelling` for the cancel retry, unless the chain shows it already ended. Returns whether it
/// moved it back.
async fn reopen_if_cancelled_but_accepted<R, T>(
snapshot: &AgreementStateSnapshot,
agreement: &IndexingAgreement,
registry: &R,
chain_client: &T,
) -> anyhow::Result<()>
) -> anyhow::Result<bool>
where
R: AgreementRegistry + Sync,
T: ChainClient,
Expand All @@ -1013,20 +1015,24 @@ where
&& snapshot.state.reached_accepted()
&& !snapshot.state.is_canceled()
{
crate::cancel_dispatch::reopen_if_live(registry, chain_client, agreement).await?;
return Ok(
crate::cancel_dispatch::reopen_if_live(registry, chain_client, agreement).await?,
);
}
Ok(())
Ok(false)
}

/// An agreement dipper is cancelling stays `Cancelling` when the chain shows it accepted,
/// so its accept is recorded here; its end is then announced, along with the accept,
/// once the cancel lands. A withdrawn offer reads as cancelled with no accept time.
/// An agreement dipper is cancelling, or has just moved back to cancelling, stays `Cancelling`
/// when the chain shows it accepted, so its accept is recorded here; its end is then announced,
/// along with the accept, once the cancel lands. A withdrawn offer reads as cancelled with no
/// accept time.
async fn record_accept_of_cancelling<R: AgreementRegistry + Sync>(
snapshot: &AgreementStateSnapshot,
agreement: &IndexingAgreement,
reopened: bool,
registry: &R,
) {
if agreement.status != IndexingAgreementStatus::Cancelling
if (agreement.status != IndexingAgreementStatus::Cancelling && !reopened)
|| !snapshot.state.reached_accepted()
|| snapshot.accepted_at == 0
|| !created_after_events_started(agreement)
Expand Down Expand Up @@ -1343,6 +1349,10 @@ fn log_orphan_cancel(
reason = "request_canceled",
"Cancelling orphan agreement"
),
Err(crate::registry::Error::NoRecordsUpdated) => tracing::debug!(
agreement_id = %agreement.id,
"Orphan agreement already ended or being cancelled"
),
Err(err) => tracing::warn!(
error = %err,
agreement_id = %agreement.id,
Expand Down Expand Up @@ -2534,14 +2544,12 @@ mod tests {
}
}

/// Minimal `ChainClient` mock for chain_listener tests. Records every
/// on-chain cancel attempt. Tests can mark specific agreements as
/// already-canceled-on-chain (cancel returns `Ok(None)`); unmarked
/// agreements get a successful `Ok(Some(zero))`.
/// Minimal `ChainClient` mock for chain_listener tests. Records every on-chain cancel
/// attempt. Every agreement reads as not live, so nothing is sent, unless
/// `live_until_cancelled` is set.
#[derive(Clone, Default)]
struct MockChainClient {
cancels: Arc<Mutex<Vec<[u8; 16]>>>,
already_canceled: Arc<Mutex<Vec<[u8; 16]>>>,
/// When set, each cancel records whether its agreement was already marked.
registry: Option<MockRegistry>,
marked_at_cancel: Arc<Mutex<Vec<bool>>>,
Expand All @@ -2562,10 +2570,6 @@ mod tests {
fn was_on_chain_cancel_attempted(&self, id: &IndexingAgreementId) -> bool {
self.cancels.lock().unwrap().contains(id.as_bytes())
}

fn mark_already_canceled_on_chain(&self, id: &IndexingAgreementId) {
self.already_canceled.lock().unwrap().push(*id.as_bytes());
}
}

#[async_trait::async_trait]
Expand Down Expand Up @@ -3112,6 +3116,10 @@ mod tests {
assert!(result.is_ok());
assert!(registry.was_reopened(&agreement_id));
assert!(!registry.was_marked_accepted_on_chain(&agreement_id));
assert!(
registry.audit_writes().contains(&("accept", agreement_id)),
"its accept is recorded, so the cancel retry treats it as paying"
);
}

#[tokio::test]
Expand Down Expand Up @@ -3708,7 +3716,6 @@ mod tests {
registry.add_agreement(new_id, IndexingAgreementStatus::AcceptedOnChain);
registry.add_agreement(old_id, IndexingAgreementStatus::AcceptedOnChain);
registry.add_pending_cancellation(new_id, old_id);
chain_client.mark_already_canceled_on_chain(&old_id);

let result = execute_pending_cancellations(
&new_id,
Expand Down Expand Up @@ -3740,7 +3747,6 @@ mod tests {
registry.add_agreement(new_id, IndexingAgreementStatus::AcceptedOnChain);
registry.add_agreement(old_id, IndexingAgreementStatus::AcceptedOnChain);
registry.add_pending_cancellation(new_id, old_id);
chain_client.mark_already_canceled_on_chain(&old_id);

sweep_executable_pending_cancellations(
&registry,
Expand Down Expand Up @@ -4519,7 +4525,6 @@ mod tests {
registry.add_agreement(agreement_id, IndexingAgreementStatus::AcceptedOnChain);
registry.set_agreement_request_id(agreement_id, request_id);
registry.mark_request_canceled(request_id);
chain_client.mark_already_canceled_on_chain(&agreement_id);

sweep_orphan_canceled_agreements(&registry, &chain_client, test_agreement_conf().as_ref())
.await;
Expand Down
116 changes: 73 additions & 43 deletions bin/dipper-service/src/worker/handlers/reassess_indexing_request.rs
Original file line number Diff line number Diff line change
Expand Up @@ -638,9 +638,13 @@ where
let mut directly_cancelled = 0u32;
let mut cancel_failures = 0u32;
for old_agreement in old_iter {
let Some(new_status) = cancel_unpaired(&ctx, old_agreement).await else {
cancel_failures += 1;
continue;
let new_status = match cancel_unpaired(&ctx, old_agreement).await {
Unpaired::Moved(new_status) => new_status,
Unpaired::AlreadyEnding => continue,
Unpaired::Failed => {
cancel_failures += 1;
continue;
}
};
tracing::info!(
agreement_id = %old_agreement.id,
Expand Down Expand Up @@ -897,12 +901,11 @@ mod lifecycle_event_tests {

// ---- Mock: worker queue --------------------------------------------------

/// Records every `send_indexing_agreement_proposal` call's indexer URL and every
/// queued on-chain cancel. Clone shares the buffers for inspection after `handle`.
/// Records every `send_indexing_agreement_proposal` call's indexer URL. Clone shares the
/// buffer for inspection after `handle`.
#[derive(Default, Clone)]
struct MockQueue {
proposals: Arc<Mutex<Vec<Url>>>,
cancels_queued: Arc<Mutex<Vec<IndexingAgreementId>>>,
}

#[async_trait]
Expand Down Expand Up @@ -1707,9 +1710,8 @@ mod lifecycle_event_tests {
async fn never_accepted_unpaired_cancel_does_not_emit_terminated() {
// Both old agreements were never accepted on-chain (Created). One add
// pairs with the first old agreement; the second, unpaired old agreement
// reaches the cancel loop but, being never-accepted
// (`!was_accepted`), must NOT emit `terminated`. The add still
// emits `proposed`. Net: exactly one event, a `proposed`.
// reaches the cancel loop but, never accepted, must NOT emit `terminated`.
// The add still emits `proposed`. Net: exactly one event, a `proposed`.
let new_idx = indexer_id(0x44);
let old_paired = indexer_id(0x55);
let old_unpaired = indexer_id(0x56);
Expand Down Expand Up @@ -1836,14 +1838,12 @@ mod lifecycle_event_tests {
ctx.chain_client.fail_cancel = true;
let cancelling = ctx.registry.marked_cancelling.clone();
let cancelled = ctx.registry.marked_cancelled.clone();
let queue = ctx.queue.clone();

let result = handle(ctx, &test_message(0)).await;

assert!(result.is_ok(), "got {result:?}");
assert_eq!(*cancelling.lock().unwrap(), vec![leaving.id]);
assert!(cancelled.lock().unwrap().is_empty());
assert!(queue.cancels_queued.lock().unwrap().is_empty());
}
}

Expand Down Expand Up @@ -1876,15 +1876,45 @@ mod lifecycle_event_tests {

assert_eq!(*chain_client.cancelled.lock().unwrap(), vec![leaving_id]);
}

#[test]
fn an_agreement_that_ended_since_it_was_listed_is_not_a_failure() {
let agreement = crate::cancel_dispatch::tests::agreement(
crate::registry::IndexingAgreementStatus::AcceptedOnChain,
None,
);
let backend_down = crate::registry::Error::BackendError(dipper_pgregistry::Error::DbError(
sqlx::Error::PoolTimedOut,
));

assert_eq!(
super::unmarked(&agreement, &crate::registry::Error::NoRecordsUpdated),
super::Unpaired::AlreadyEnding
);
assert_eq!(
super::unmarked(&agreement, &backend_down),
super::Unpaired::Failed
);
}
}

/// Take an old agreement out of the target group, returning its new status, or `None`
/// (logged) when it couldn't be marked. One that may be live on-chain is cancelled there
/// too, which the chain listener retries until it ends.
/// What became of an old agreement taken out of the target group.
#[derive(Debug, PartialEq, Eq)]
enum Unpaired {
/// Moved to the status named.
Moved(&'static str),
/// Already ended or being cancelled, as a race with another path can leave it.
AlreadyEnding,
/// Couldn't be marked; logged.
Failed,
}

/// Take an old agreement out of the target group. One that may be live on-chain is cancelled
/// there too, which the chain listener retries until it ends.
async fn cancel_unpaired<R, W, I, T>(
ctx: &Ctx<R, W, I, T>,
agreement: &crate::registry::IndexingAgreement,
) -> Option<&'static str>
) -> Unpaired
where
R: AgreementRegistry + Sync,
T: ChainClient,
Expand All @@ -1893,9 +1923,14 @@ where
|| (agreement.status == crate::registry::IndexingAgreementStatus::Created
&& agreement.terms_version_hash.is_some());
if !may_be_live {
return mark_unpaired_cancelled(&ctx.registry, agreement)
return match ctx
.registry
.mark_indexing_agreement_as_canceled_by_requester(&agreement.id)
.await
.then_some("CANCELED_BY_REQUESTER");
{
Ok(()) => Unpaired::Moved("CANCELED_BY_REQUESTER"),
Err(err) => unmarked(agreement, &err),
};
}
match crate::cancel_dispatch::start_cancel(
&ctx.registry,
Expand All @@ -1906,38 +1941,33 @@ where
)
.await
{
Ok(crate::cancel_dispatch::CancelStarted::Ended) => Some("CANCELED_BY_REQUESTER"),
Ok(crate::cancel_dispatch::CancelStarted::Cancelling) => Some("CANCELLING"),
Err(err) => {
tracing::error!(
error = %err,
agreement_id = %agreement.id,
"Failed to mark unpaired old agreement as cancelling in local DB"
);
None
Ok(crate::cancel_dispatch::CancelStarted::Ended) => {
Unpaired::Moved("CANCELED_BY_REQUESTER")
}
Ok(crate::cancel_dispatch::CancelStarted::Cancelling) => Unpaired::Moved("CANCELLING"),
Err(err) => unmarked(agreement, &err),
}
}

/// Mark an unpaired old agreement CanceledByRequester; false, logged, if that fails.
async fn mark_unpaired_cancelled<R: AgreementRegistry + Sync>(
registry: &R,
/// Why an unpaired old agreement couldn't be marked, logged at the level it deserves: one that
/// ended, or started cancelling, since it was listed is expected and needs nothing more.
fn unmarked(
agreement: &crate::registry::IndexingAgreement,
) -> bool {
match registry
.mark_indexing_agreement_as_canceled_by_requester(&agreement.id)
.await
{
Ok(()) => true,
Err(err) => {
tracing::error!(
error = %err,
agreement_id = %agreement.id,
"Failed to mark unpaired old agreement as canceled in local DB"
);
false
}
err: &crate::registry::Error,
) -> Unpaired {
if matches!(err, crate::registry::Error::NoRecordsUpdated) {
tracing::debug!(
agreement_id = %agreement.id,
"Unpaired old agreement already ended or being cancelled"
);
return Unpaired::AlreadyEnding;
}
tracing::error!(
error = %err,
agreement_id = %agreement.id,
"Failed to mark unpaired old agreement as ended in local DB"
);
Unpaired::Failed
}

/// Olds reserved from cancellation: one per add-cancel pairing lost to a
Expand Down
Loading