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
94 changes: 77 additions & 17 deletions bin/dipper-service/src/cancel_dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -163,7 +163,8 @@ where

/// Mark an agreement the chain shows dipper ended as ended, 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.
/// fails; it stays `Cancelling` for the cancel retry. One the chain listener already marked
/// ended counts as ended.
pub async fn confirm_cancelled<R: AgreementRegistry + Sync>(
registry: &R,
agreement: &IndexingAgreement,
Expand All @@ -175,16 +176,26 @@ pub async fn confirm_cancelled<R: AgreementRegistry + Sync>(
if tx_hash.is_some() {
record_cancel(registry, agreement, tx_hash, config).await;
}
if let Err(err) = registry
match 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; the cancel retry tries again"
);
return false;
Ok(()) => {}
Err(crate::registry::Error::NoRecordsUpdated) => {
tracing::debug!(
agreement_id = %agreement.id,
"Agreement already marked ended, as the chain listener can do first"
);
return true;
}
Err(err) => {
tracing::warn!(
agreement_id = %agreement.id,
error = %err,
"Failed to mark an ended agreement cancelled; the cancel retry tries again"
);
return false;
}
}
tracing::info!(
agreement_id = %agreement.id,
Expand Down Expand Up @@ -232,20 +243,23 @@ where
R: AgreementRegistry + Sync,
T: ChainClient,
{
match chain_client
let seen_live = match chain_client
.agreement_on_chain(agreement.id.as_bytes())
.await
{
Ok(AgreementOnChain::Live) => {}
Ok(AgreementOnChain::Live) => true,
Ok(_) => return Ok(false),
Err(err) => tracing::warn!(
agreement_id = %agreement.id,
error = %err,
"Failed to read an ended agreement reported live; the cancel retry checks it"
),
}
Err(err) => {
tracing::warn!(
agreement_id = %agreement.id,
error = %err,
"Failed to read an ended agreement reported live; the cancel retry checks it"
);
false
}
};
match registry
.reopen_indexing_agreement_cancel(&agreement.id)
.reopen_indexing_agreement_cancel(&agreement.id, seen_live)
.await
{
Ok(()) => {}
Expand Down Expand Up @@ -632,4 +646,50 @@ pub(crate) mod tests {
}
));
}

/// Records what each reopen was told about the chain.
#[derive(Default)]
struct ReopenRegistry {
seen_live: Mutex<Vec<bool>>,
}

#[async_trait]
impl crate::registry::StubAgreementRegistry for ReopenRegistry {
async fn reopen_indexing_agreement_cancel(
&self,
_id: &IndexingAgreementId,
seen_live: bool,
) -> crate::registry::Result<()> {
self.seen_live.lock().unwrap().push(seen_live);
Ok(())
}
}

#[tokio::test]
async fn a_reopen_clears_the_end_on_record_only_when_the_chain_shows_it_live() {
// An unread chain reopens it all the same, but it may have ended, so its record stays.
let ag = agreement(
IndexingAgreementStatus::CanceledByRequester,
Some(vec![7u8; 32]),
);
let live = RecordingChainClient {
still_active_after_cancel: true,
..Default::default()
};
let unread = RecordingChainClient {
read_back_fails: true,
..Default::default()
};

for (client, seen_live) in [(live, true), (unread, false)] {
let registry = ReopenRegistry::default();

let reopened = super::reopen_if_live(&registry, &client, &ag)
.await
.expect("reopen");

assert!(reopened);
assert_eq!(*registry.seen_live.lock().unwrap(), vec![seen_live]);
}
}
}
18 changes: 18 additions & 0 deletions bin/dipper-service/src/network/service/cancel_retry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -337,6 +337,8 @@ mod tests {
attempts: AtomicU32,
checks: AtomicU32,
found_ended: Mutex<Vec<Option<bool>>>,
/// The chain listener marks it ended before the retry's own mark lands.
listener_ended_it: bool,
writes: Mutex<Vec<&'static str>>,
}

Expand All @@ -354,6 +356,9 @@ mod tests {
&self,
id: &IndexingAgreementId,
) -> crate::registry::Result<()> {
if self.listener_ended_it {
return Err(crate::registry::Error::NoRecordsUpdated);
}
self.marked_cancelled.lock().unwrap().push(*id);
self.writes.lock().unwrap().push("ended");
Ok(())
Expand Down Expand Up @@ -577,6 +582,19 @@ mod tests {
);
}

#[tokio::test]
async fn counts_an_agreement_the_listener_marked_ended_first_as_ended() {
// Not a failure: there is nothing left to retry.
let registry = MockRegistry {
listener_ended_it: true,
..registry_with_one(true)
};

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

assert_eq!(registry.checks.load(Ordering::SeqCst), 0);
}

#[tokio::test]
async fn leaves_an_accepted_agreement_that_already_ended_to_the_listener() {
// The indexer may have ended it, or an earlier cancel whose result went unread;
Expand Down
1 change: 1 addition & 0 deletions bin/dipper-service/src/network/service/chain_listener.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2217,6 +2217,7 @@ mod tests {
async fn reopen_indexing_agreement_cancel(
&self,
id: &IndexingAgreementId,
_seen_live: bool,
) -> RegistryResult<()> {
self.state.lock().unwrap().reopened.push(*id);
Ok(())
Expand Down
3 changes: 2 additions & 1 deletion bin/dipper-service/src/registry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -414,9 +414,10 @@ impl AgreementRegistry for RegistryProvider {
async fn reopen_indexing_agreement_cancel(
&self,
id: &IndexingAgreementId,
seen_live: bool,
) -> RegistryResult<()> {
self.inner
.reopen_indexing_agreement_cancel(id)
.reopen_indexing_agreement_cancel(id, seen_live)
.await
.map_err(Into::into)
}
Expand Down
4 changes: 3 additions & 1 deletion bin/dipper-service/src/registry/agreement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -284,10 +284,12 @@ pub trait AgreementRegistry {

/// Move a `CANCELED_BY_REQUESTER` or `REJECTED` agreement the chain shows live back to
/// `CANCELLING`, its cancel attempts reset; [`NoRecordUpdated`](Error::NoRecordsUpdated)
/// otherwise.
/// otherwise. `seen_live` when the chain was read and showed it live, which clears the end
/// on record and any announcement of it.
async fn reopen_indexing_agreement_cancel(
&self,
id: &IndexingAgreementId,
seen_live: bool,
) -> RegistryResult<()>;

/// `CANCELLING` agreements marked over `min_age_minutes` ago, those checked longest ago
Expand Down
14 changes: 11 additions & 3 deletions bin/dipper-service/src/registry/agreement_stub.rs
Original file line number Diff line number Diff line change
Expand Up @@ -141,7 +141,11 @@ pub trait StubAgreementRegistry: Send + Sync {
unimplemented!("mark_indexing_agreement_as_abandoning")
}

async fn reopen_indexing_agreement_cancel(&self, _id: &IndexingAgreementId) -> Result<()> {
async fn reopen_indexing_agreement_cancel(
&self,
_id: &IndexingAgreementId,
_seen_live: bool,
) -> Result<()> {
unimplemented!("reopen_indexing_agreement_cancel")
}

Expand Down Expand Up @@ -446,8 +450,12 @@ impl<T: StubAgreementRegistry> AgreementRegistry for T {
StubAgreementRegistry::mark_indexing_agreement_as_abandoning(self, id).await
}

async fn reopen_indexing_agreement_cancel(&self, id: &IndexingAgreementId) -> Result<()> {
StubAgreementRegistry::reopen_indexing_agreement_cancel(self, id).await
async fn reopen_indexing_agreement_cancel(
&self,
id: &IndexingAgreementId,
seen_live: bool,
) -> Result<()> {
StubAgreementRegistry::reopen_indexing_agreement_cancel(self, id, seen_live).await
}

async fn get_cancelling_agreements(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,7 @@ mod tests {
async fn reopen_indexing_agreement_cancel(
&self,
id: &IndexingAgreementId,
_seen_live: bool,
) -> crate::registry::Result<()> {
self.reopened.lock().unwrap().push(*id);
Ok(())
Expand Down
26 changes: 22 additions & 4 deletions dipper-pgregistry/src/postgres.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1008,10 +1008,13 @@ impl PgRegistry {
/// Move an agreement dipper had already ended, cancelled or rejected, back to `Cancelling`
/// once the chain shows it live after all, with its cancel attempts started afresh. It counts
/// as checked, since no cancel is sent with it, so the retry takes it on its next sweep rather
/// than waiting for one to be mined.
/// than waiting for one to be mined. When the chain was read and showed it live (`seen_live`),
/// the end on record, and any announcement of it, no longer stands, so both are cleared for
/// the end still to come; an unread chain leaves them, as the agreement may have ended.
pub async fn reopen_indexing_agreement_cancel(
&self,
agreement_id: &IndexingAgreementId,
seen_live: bool,
) -> Result<(), Error> {
let updated = sqlx::query(
r#"
Expand All @@ -1021,6 +1024,11 @@ impl PgRegistry {
cancel_attempts = 0,
cancel_checked_at = timezone('UTC', now()),
ended_seen_at = NULL,
canceled_at = CASE WHEN $5::BOOLEAN THEN NULL ELSE canceled_at END,
canceled_by = CASE WHEN $5 THEN NULL ELSE canceled_by END,
canceled_tx = CASE WHEN $5 THEN NULL ELSE canceled_tx END,
terminated_event_emitted_at =
CASE WHEN $5 THEN NULL ELSE terminated_event_emitted_at END,
updated_at = timezone('UTC', now())
WHERE id = $2 AND status IN ($3, $4)
"#,
Expand All @@ -1029,6 +1037,7 @@ impl PgRegistry {
.bind(agreement_id)
.bind(IndexingAgreementStatus::CanceledByRequester)
.bind(IndexingAgreementStatus::Rejected)
.bind(seen_live)
.execute(&self.pool)
.await?;
if updated.rows_affected() == 0 {
Expand Down Expand Up @@ -1544,6 +1553,8 @@ impl PgRegistry {
/// emission sweep can populate the `terminated` event's tx/by/at fields.
/// `COALESCE` keeps any value already observed on-chain. Best-effort
/// enrichment: the event still emits (with fallbacks) if never recorded.
/// An end recorded before the agreement's accept, such as its offer's withdrawal before the
/// offer landed after all, can't be its end, so a later one replaces it and is announced.
#[expect(
clippy::cast_possible_wrap,
reason = "predates this lint; fix when next touched"
Expand All @@ -1558,9 +1569,16 @@ impl PgRegistry {
sqlx::query(
r#"
UPDATE dipper_reg_indexing_agreements
SET canceled_at = COALESCE(canceled_at, $2),
canceled_by = COALESCE(canceled_by, $3),
canceled_tx = COALESCE(canceled_tx, $4)
SET canceled_at = CASE WHEN canceled_at < accepted_at AND $2 >= accepted_at
THEN $2 ELSE COALESCE(canceled_at, $2) END,
canceled_by = CASE WHEN canceled_at < accepted_at AND $2 >= accepted_at
THEN $3 ELSE COALESCE(canceled_by, $3) END,
canceled_tx = CASE WHEN canceled_at < accepted_at AND $2 >= accepted_at
THEN $4 ELSE COALESCE(canceled_tx, $4) END,
terminated_event_emitted_at = CASE
WHEN canceled_at < accepted_at AND $2 >= accepted_at THEN NULL
ELSE terminated_event_emitted_at
END
WHERE id = $1
"#,
)
Expand Down
Loading
Loading