Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
cd1b224
fix(registry): make every registry implement its event and audit writes
MoonBoi9001 Oct 9, 2026
ef40e47
fix(k8s): make the example Slack webhook placeholder a valid URL
MoonBoi9001 Oct 9, 2026
735c065
fix(liveness): stop a hung database from stalling the liveness check
MoonBoi9001 Oct 9, 2026
db50eca
fix(cancel): wait longer before retrying a cancel that was just sent
MoonBoi9001 Oct 9, 2026
f642a58
fix(listener): record no accept for a rejected agreement already ended
MoonBoi9001 Oct 9, 2026
90b8fb6
fix(cancel): retry accepted cancels even when the chain time is unknown
MoonBoi9001 Oct 9, 2026
7aa00d1
fix(pacing): count offers still being withdrawn as in flight
MoonBoi9001 Oct 9, 2026
b553c27
fix(listener): re-read a page of chain changes that failed to apply
MoonBoi9001 Oct 9, 2026
6ed2d5b
fix(chain): let 2 RPC endpoints, not 1, move the newest block far ahead
MoonBoi9001 Oct 9, 2026
ee98b94
docs(registry): say the test stub alone keeps the no-op audit defaults
MoonBoi9001 Oct 9, 2026
15c1e51
fix(reassess): keep an abandoned agreement's slot while it may be paid
MoonBoi9001 Oct 9, 2026
21abb6a
refactor(liveness): split marking a stale agreement from cancelling it
MoonBoi9001 Oct 9, 2026
b7c1d18
style(fmt): lay out this branch's changes the way rustfmt wants
MoonBoi9001 Oct 9, 2026
6424686
fix(reassess): never cancel an agreement because another holds a slot
MoonBoi9001 Oct 9, 2026
abb43e5
fix(liveness): stop waiting on a hung database when noting an end
MoonBoi9001 Oct 9, 2026
af2b710
fix(listener): move past a page of changes that keeps failing to apply
MoonBoi9001 Oct 9, 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
48 changes: 35 additions & 13 deletions bin/dipper-service/src/cancel_dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -119,32 +119,56 @@ where
R: AgreementRegistry + Sync,
T: ChainClient,
{
mark_cancelling(registry, agreement, reason).await?;
Ok(send_marked_cancel(registry, chain_client, agreement, reason, config).await)
}

/// The first half of [`start_cancel`]: mark the agreement `Cancelling`, or abandoning.
pub async fn mark_cancelling<R: AgreementRegistry + Sync>(
registry: &R,
agreement: &IndexingAgreement,
reason: CancelReason,
) -> RegistryResult<()> {
match reason {
CancelReason::NotWanted => {
registry
.mark_indexing_agreement_as_cancelling(&agreement.id)
.await?
.await
}
CancelReason::Abandoned => {
registry
.mark_indexing_agreement_as_abandoning(&agreement.id)
.await?
.await
}
}
}

/// The second half of [`start_cancel`]: cancel a marked agreement if the chain shows it live.
pub async fn send_marked_cancel<R, T>(
registry: &R,
chain_client: &T,
agreement: &IndexingAgreement,
reason: CancelReason,
config: &IndexingAgreementConfig,
) -> CancelStarted
where
R: AgreementRegistry + Sync,
T: ChainClient,
{
let tx_hash = match cancel_if_live(chain_client, agreement, config).await {
LiveCancel::Ended(tx_hash) => tx_hash,
LiveCancel::NotLive { .. } => return Ok(CancelStarted::NotLive),
LiveCancel::NotLive { .. } => return CancelStarted::NotLive,
LiveCancel::ReadFailed(err) | LiveCancel::CancelFailed(err) => {
tracing::warn!(
agreement_id = %agreement.id,
error = %err,
"On-chain cancel failed; the cancel retry sends it again"
);
return Ok(CancelStarted::MayBeLive);
return CancelStarted::MayBeLive;
}
LiveCancel::Unconfirmed { tx_hash, err } => {
log_unconfirmed(agreement, tx_hash, &err);
return Ok(CancelStarted::MayBeLive);
return CancelStarted::MayBeLive;
}
};
tracing::info!(
Expand All @@ -154,15 +178,13 @@ where
);
// An offer never accepted could still land and be accepted until its deadline.
if agreement.status != IndexingAgreementStatus::AcceptedOnChain {
return Ok(CancelStarted::NotLive);
return CancelStarted::NotLive;
}
if confirm_cancelled(registry, agreement, reason, tx_hash, config).await {
CancelStarted::Ended
} else {
CancelStarted::NotLive
}
Ok(
if confirm_cancelled(registry, agreement, reason, tx_hash, config).await {
CancelStarted::Ended
} else {
CancelStarted::NotLive
},
)
}

/// Mark an agreement the chain shows dipper ended as ended, recording the cancel when its
Expand Down
134 changes: 111 additions & 23 deletions bin/dipper-service/src/chain_client/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -92,12 +92,26 @@ const CROSS_CHECK_DEADLINE: Duration = Duration::from_secs(3);
/// Blocks too far ahead of it are refused only once it is confirmed, by the receipt for one of
/// dipper's transactions or by 2 endpoints agreeing, so an endpoint stuck far behind, or far
/// ahead, that answers first after a restart can't shut out the endpoints that are right.
/// Once confirmed, 1 endpoint alone can't move it far ahead either, for the same reason.
#[derive(Debug)]
struct SeenBlock {
number: u64,
moved_at: Instant,
confirmed: bool,
cross_checked_at: Option<Instant>,
/// An endpoint reported a block further ahead than 1 endpoint alone may move it, so the
/// endpoints are asked whether they agree.
jump_reported: bool,
}

/// What became of a block 1 endpoint reported.
#[derive(Debug, PartialEq, Eq)]
enum Noted {
Taken,
/// Too far ahead for 1 endpoint alone; it waits for 2 to agree.
Unagreed,
/// Too far ahead to be real.
TooFarAhead,
}

impl SeenBlock {
Expand All @@ -107,6 +121,7 @@ impl SeenBlock {
moved_at: Instant::now(),
confirmed: false,
cross_checked_at: None,
jump_reported: false,
}
}

Expand All @@ -130,24 +145,47 @@ impl SeenBlock {
(self.number, highest)
}

/// Take `block` as seen, unless it is too far ahead of a confirmed one to be real; false
/// when it is.
fn advance(&mut self, block: u64, now: Instant) -> bool {
if self.confirmed && block > self.believable_limit(now) {
return false;
/// The highest block 1 endpoint alone may move a confirmed newest block to: an hour of
/// blocks past it, plus what the chain can have added since.
fn lone_limit(&self, now: Instant) -> u64 {
let since = now.saturating_duration_since(self.moved_at).as_secs();
self.number
.saturating_add(ALERT_GAP_BLOCKS)
.saturating_add(since.saturating_mul(BLOCKS_PER_SECOND))
}

/// Take `block`, reported by 1 endpoint, as seen. Once the newest block is confirmed, one
/// too far ahead to be real is refused, and one past [`Self::lone_limit`] waits for 2
/// endpoints to agree, unless this is the `only_endpoint` there is to ask.
fn advance(&mut self, block: u64, now: Instant, only_endpoint: bool) -> Noted {
if self.confirmed {
if block > self.believable_limit(now) {
return Noted::TooFarAhead;
}
if !only_endpoint && block > self.lone_limit(now) {
self.jump_reported = true;
return Noted::Unagreed;
}
}
self.take(block, now);
Noted::Taken
}

fn take(&mut self, block: u64, now: Instant) {
if block > self.number {
self.number = block;
self.moved_at = now;
}
true
}

/// Take `block` as one the chain is known to have reached. An unconfirmed newest block more
/// than an hour of blocks past it came from a faulty endpoint, so it is replaced.
fn confirm(&mut self, block: u64, now: Instant) {
self.jump_reported = false;
if self.confirmed {
self.advance(block, now);
if block <= self.believable_limit(now) {
self.take(block, now);
}
return;
}
if block > self.number || self.number > block.saturating_add(ALERT_GAP_BLOCKS) {
Expand All @@ -157,13 +195,13 @@ impl SeenBlock {
self.confirmed = true;
}

/// Whether to ask every endpoint for its latest block: only while unconfirmed, and at most
/// once every [`CROSS_CHECK_INTERVAL`].
/// Whether to ask every endpoint for its latest block: only while unconfirmed or after a
/// jump 1 endpoint reported, and at most once every [`CROSS_CHECK_INTERVAL`].
fn cross_check_due(&mut self, now: Instant) -> bool {
let checked_lately = self
.cross_checked_at
.is_some_and(|at| now.saturating_duration_since(at) < CROSS_CHECK_INTERVAL);
if self.confirmed || checked_lately {
if (self.confirmed && !self.jump_reported) || checked_lately {
return false;
}
self.cross_checked_at = Some(now);
Expand Down Expand Up @@ -899,24 +937,37 @@ impl AlloyChainClient {
Some(head) => self.seen_block().confirm(head, Instant::now()),
None => tracing::warn!(
"No 2 RPC endpoints agree on the chain's latest block; reads aren't checked \
against blocks too far ahead until they do"
against blocks too far ahead, nor the newest block seen moved far ahead, until \
they do"
),
}
}

/// Remember a block dipper has seen, so later reads are never older. One too far ahead to
/// be real is ignored, so it can't refuse every read after it.
/// be real is ignored, so it can't refuse every read after it, and one over an hour ahead
/// waits for 2 endpoints to agree on it.
fn note_block(&self, block: u64) {
let (taken, seen) = {
let only_endpoint = self.inner.rpc_pool.endpoint_count() < 2;
let (noted, seen) = {
let mut seen = self.seen_block();
(seen.advance(block, Instant::now()), seen.number)
(
seen.advance(block, Instant::now(), only_endpoint),
seen.number,
)
};
if !taken {
tracing::warn!(
match noted {
Noted::Taken => {}
Noted::Unagreed => tracing::warn!(
block,
seen_block = seen,
"Not moving to a block over an hour ahead of the newest block dipper has seen \
until 2 RPC endpoints agree on it"
),
Noted::TooFarAhead => tracing::warn!(
block,
seen_block = seen,
"Ignoring a block too far ahead of the newest block dipper has seen"
);
),
}
}

Expand Down Expand Up @@ -2549,6 +2600,7 @@ mod tests {
moved_at: now.checked_sub(Duration::from_secs(600)).expect("instant"),
confirmed: true,
cross_checked_at: None,
jump_reported: false,
};

assert_eq!(seen.bounds(now), (100, 100 + WEEK_OF_BLOCKS + 2_400));
Expand Down Expand Up @@ -2590,8 +2642,8 @@ mod tests {
// So a read just after dipper's cancel mined can't come from an endpoint behind it.
let now = Instant::now();
let mut seen = SeenBlock::new();
assert!(seen.advance(100, now));
assert!(seen.advance(90, now));
assert_eq!(seen.advance(100, now, false), Noted::Taken);
assert_eq!(seen.advance(90, now, false), Noted::Taken);

assert_eq!(seen.bounds(now), (100, u64::MAX));
}
Expand All @@ -2603,21 +2655,57 @@ mod tests {
let now = Instant::now();
let mut seen = SeenBlock::new();
let real = 1_000 + WEEK_OF_BLOCKS * 3;
assert!(seen.advance(1_000, now));
assert!(seen.advance(real, now), "taken while unconfirmed");
assert_eq!(seen.advance(1_000, now, false), Noted::Taken);
assert_eq!(
seen.advance(real, now, false),
Noted::Taken,
"taken while unconfirmed"
);

seen.confirm(real, now);

assert!(!seen.advance(real + WEEK_OF_BLOCKS * 2, now));
assert_eq!(
seen.advance(real + WEEK_OF_BLOCKS * 2, now, false),
Noted::TooFarAhead
);
assert_eq!(seen.bounds(now).0, real);
}

#[test]
fn waits_for_2_endpoints_to_agree_before_a_jump_past_a_confirmed_block() {
// 1 faulty endpoint moving it far ahead would have every endpoint that is right
// refused as behind, for as long as the real chain takes to catch up.
let now = Instant::now();
let mut seen = SeenBlock::new();
let real = 5_000_000;
seen.confirm(real, now);
let jump = real + ALERT_GAP_BLOCKS + 1_000;

assert_eq!(seen.advance(real + 10, now, false), Noted::Taken);
assert_eq!(seen.advance(jump, now, false), Noted::Unagreed);
assert_eq!(seen.bounds(now).0, real + 10);
assert!(seen.cross_check_due(now), "the endpoints are asked");

seen.confirm(jump, now);

assert_eq!(seen.bounds(now).0, jump, "taken once 2 agree");
assert!(!seen.cross_check_due(now + CROSS_CHECK_INTERVAL));
assert_eq!(
seen.advance(jump + ALERT_GAP_BLOCKS + 1, now, true),
Noted::Taken,
"the only endpoint has no other to agree with"
);
}

#[test]
fn a_confirmed_block_replaces_one_far_ahead_that_was_never_confirmed() {
let now = Instant::now();
let mut seen = SeenBlock::new();
let real = 5_000_000;
assert!(seen.advance(real + WEEK_OF_BLOCKS * 10, now));
assert_eq!(
seen.advance(real + WEEK_OF_BLOCKS * 10, now, false),
Noted::Taken
);

seen.confirm(real, now);

Expand Down
17 changes: 17 additions & 0 deletions bin/dipper-service/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2118,6 +2118,23 @@ mod tests {
);
}

#[test]
fn the_example_configmap_alerts_section_parses() {
let example = include_str!("../../../k8s/configmap-example.yaml");
let json: String = example
.split_once("config.json: |\n")
.expect("the example holds a config.json")
.1
.lines()
.map(|line| format!("{}\n", line.strip_prefix(" ").unwrap_or(line)))
.collect();
let config: serde_json::Value = serde_json::from_str(&json).expect("config.json is JSON");

let alerts = serde_json::from_value::<AlertsConfig>(config["alerts"].clone())
.expect("the example's alerts section is a valid alerts config");
assert!(alerts.slack_webhook_url.is_some());
}

#[test]
fn alerts_config_takes_a_webhook_and_keeps_the_default_events() {
let alerts = serde_json::from_str::<AlertsConfig>(
Expand Down
Loading
Loading