From bf8e59f81fdb0d40a95bb3161c8c69a6b9dc85b9 Mon Sep 17 00:00:00 2001 From: Michael Taylor Date: Thu, 10 Sep 2026 05:20:46 -0700 Subject: [PATCH 1/7] feat(rewards): wire the peer claim loop onto a cadence driver from real startup (#3268) The peer reward-claim engine shipped complete and tested in #594 but INERT: nothing constructed it, so the 86400s cadence never fired while `rewards_claim.enabled` defaulted to `true` -- a config asserting a subsystem is on while nothing runs. `rewards_claim/driver.rs` is a SCHEDULER, not a chain adapter: it derives this node's own payout puzzle hash, loads `RewardsClaimConfig`, builds a `ClaimEngine` against the only production port that exists (`UnavailableClaimChainPort`, until #3249 lands a real one) and drives `run_cycle` every `cadence_seconds + jitter`, jitter drawn from the OS CSPRNG. `server.rs`'s `serve_with_shutdown` makes exactly one call into it, beside `self_heal::spawn_driver_if_service()`. `enabled = true` now means: a background task exists, drives a counted cycle per interval, and its outcome is readable in-process as a NAMED state. With `UnavailableClaimChainPort` every cycle honestly reports `ChainSourceUnavailable` -- the gap is loud instead of silent. Anti-silence: `ClaimLoopHandle` carries a monotonic `cycles_driven` counter alongside the status, because `Idle` before the first cycle is correct and honest, so status alone cannot tell "scheduler never fired" from "nothing was claimable". The gate takes an INJECTED handle rather than reading the process-wide singleton, so `ClaimDriverRefusal::{Disabled, ChainSyncDisabled, NoOperatorWallet}` and "spawned but never ticked" are four pairwise-distinct readings a test asserts in-process. Nothing goes on the wire: no RPC method, dispatch row, handler or OpenRPC entry. `ClaimStatus` stays off the wire until #3249's real adapter lets the status surface be re-derived against it. Co-Authored-By: Claude Opus 5 (1M context) --- Cargo.lock | 1 + crates/dig-node-service/Cargo.toml | 6 + .../src/rewards_claim/driver.rs | 872 ++++++++++++++++++ .../dig-node-service/src/rewards_claim/mod.rs | 24 +- crates/dig-node-service/src/server.rs | 11 + 5 files changed, 904 insertions(+), 10 deletions(-) create mode 100644 crates/dig-node-service/src/rewards_claim/driver.rs diff --git a/Cargo.lock b/Cargo.lock index 058fad93..a8839fe2 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3079,6 +3079,7 @@ dependencies = [ "notify", "num-bigint", "reqwest", + "ring", "rpassword", "rustls", "serde", diff --git a/crates/dig-node-service/Cargo.toml b/crates/dig-node-service/Cargo.toml index e340a965..ed7e3f94 100644 --- a/crates/dig-node-service/Cargo.toml +++ b/crates/dig-node-service/Cargo.toml @@ -116,6 +116,12 @@ dig-mirror-coin = "0.9" # a dissenting THIRD source and did not distinguish source CLASSES from bare source strings. dig-stun = "0.2" +# The claim loop's cadence jitter (DIG-Network/dig_ecosystem#3268, SPEC §8.6) draws from the OS +# CSPRNG, never a global/thread RNG -- `ring::rand::SystemRandom` is the same primitive this +# crate's own signing paths already use for randomness (`seams::dig_peer::holdings`, in +# `dig-node-core`), so this is not a second RNG choice entering the tree. +ring = "0.17" + # The canonical `ChainSource` trait `dig-mirror-coin`'s census is generic over. Declared, not # implemented: `chia-query` already provides the implementation this node uses # (`ChiaQueryProvider`), reached through `dig-wallet`'s one shared chain transport. Caret-matched diff --git a/crates/dig-node-service/src/rewards_claim/driver.rs b/crates/dig-node-service/src/rewards_claim/driver.rs new file mode 100644 index 00000000..8f20e575 --- /dev/null +++ b/crates/dig-node-service/src/rewards_claim/driver.rs @@ -0,0 +1,872 @@ +//! Wires the claim engine onto a real background cadence, reachable from the node's actual +//! startup path (DIG-Network/dig_ecosystem#3268). Mirrors [`crate::self_heal`]'s split exactly: +//! a private injected-tick [`drive`] (testable under `#[tokio::test(start_paused = true)]`) behind +//! a pure, tested gate ([`decide_claim_driver`] / [`spawn_claim_driver_if`]) that `server.rs` calls +//! exactly once. +//! +//! # `enabled = true` must stop being a false statement +//! Before this module, nothing in the codebase ever constructed a [`super::ClaimEngine`] outside +//! its own tests (see [`super`]'s module doc, now updated). After it, a background task always +//! exists whenever `rewards_claim.enabled` and `enable_chain_sync` are both true, drives a cycle +//! every `cadence_seconds + jitter`, and its outcome is readable in-process via [`handle`] as a +//! NAMED [`super::ClaimLoopState`] — see [`ClaimLoopHandle`]. +//! +//! # Diverges from `self_heal::drive` on purpose: the FIRST pass waits for the interval +//! `self_heal::drive` fires its pass immediately, then once per fixed tick — right for a +//! maintenance sweep with no anti-silence surface. This driver's whole point (A2, the ticket's +//! headline acceptance item) is that "scheduler running, zero cycles ever fired" must be +//! DISTINGUISHABLE from "it ran" via a monotonic cycle counter that reads `0` before any interval +//! has elapsed. Running a pass at spawn, before the counter could ever read `0` under observation, +//! would defeat that on every startup. So [`drive`] sleeps `cadence_seconds + jitter` FIRST, then +//! runs a cycle, then repeats — the counter is genuinely `0` until the first interval elapses. +//! +//! # No RPC surface here (SCOPE) +//! [`handle`] is an IN-PROCESS accessor only — a future RPC (blocked on DIG-Network/dig_ecosystem#3249 +//! re-deriving the `ClaimStatus` wire semantics) can read it; this module puts nothing on the wire +//! and adds no RPC method, dispatch-table row or handler. +//! +//! # The only production adapter is [`super::UnavailableClaimChainPort`] +//! #3249 has not landed, so every real cycle this driver runs reports [`super::ClaimLoopState::ChainSourceUnavailable`] +//! and submits nothing — the honest state, not an invented adapter. + +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{Mutex, OnceLock}; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; + +use chia_protocol::Bytes32; + +use super::cadence::{next_interval_seconds, JitterSource}; +use super::config::RewardsClaimConfig; +use super::engine::ClaimEngine; +use super::hints::{DistributorHintSource, NoHintSource}; +use super::port::{ClaimChainPort, UnavailableClaimChainPort}; +use super::types::ClaimStatus; + +/// The in-process accessor onto the running claim loop (SCOPE: never exposed over the wire here). +/// Cheap to clone -- every field is an `Arc`-backed handle onto the same shared state. +/// +/// # A2, the anti-silence test +/// [`Self::cycles_driven`] is the counter a reader compares against [`Self::status`]'s +/// [`super::ClaimLoopState`] to tell "constructed and spawned but never drove a cycle" (`0`, +/// `Idle`) apart from "ran and reported a real outcome" (`> 0`, whatever [`super::ClaimEngine`] +/// computed). Neither field alone would do it: `status()` before the first cycle is already +/// `Idle` BY DESIGN (see [`super::types::ClaimLoopState::Idle`]'s doc, "no cycle has ever been +/// attempted yet") -- that is the correct, honest reading, not a defect, and a test that only +/// checked `status()` for `Idle` could not tell a scheduler that never fires apart from one that +/// correctly reports nothing pending. The count is the only thing here that is monotonic and can +/// never be read as "healthy" by a writer describing itself. +#[derive(Clone, Default)] +pub struct ClaimLoopHandle { + status: std::sync::Arc>, + cycles_driven: std::sync::Arc, + refusal: std::sync::Arc>>, +} + +/// Why the driver never reached [`drive`]'s loop at all -- distinct from anything +/// [`super::ClaimLoopState`] can say, because every one of ITS states presupposes an engine that +/// exists and a cycle that was at least attempted. Without this, "disabled", "chain sync is off" +/// and "no operator wallet, so there is nothing to build an engine with" all collapse into the +/// same reassuring `Idle` + zero-count reading -- three different truths about whether this peer +/// is being paid, indistinguishable to an operator or a future `dig.getRewardClaimStatus`. Kept on +/// the DRIVER's own handle, never added to [`super::types::ClaimStatus`] (read-only, and it is the +/// wrong home: it is a fact about whether an engine exists, not about a cycle one ran). +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ClaimDriverRefusal { + /// `rewards_claim.enabled = false` -- the ordinary, deliberate off state. + Disabled, + /// `enabled = true` but `enable_chain_sync = false` -- see [`ClaimDriverDecision::ChainSyncDisabled`]'s doc. + ChainSyncDisabled, + /// `enabled = true`, chain sync is on, but this node has no operator wallet to derive + /// [`own_payout_puzzle_hash`] from -- there is no puzzle hash to build an engine with at all. + NoOperatorWallet, +} + +impl ClaimLoopHandle { + /// The most recent [`ClaimStatus`] any cycle has produced, or [`ClaimStatus::default`]'s + /// `Idle` state before the first one ever runs. + #[must_use] + pub fn status(&self) -> ClaimStatus { + *self + .status + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + } + + /// How many times [`super::ClaimEngine::run_cycle`] has been invoked through this handle -- + /// incremented on EVERY invocation, whatever it returned (a refused, faulted or empty cycle + /// still counts: A1 requires an OBSERVED CYCLE COUNT, never "the task was spawned"). + #[must_use] + pub fn cycles_driven(&self) -> u64 { + self.cycles_driven.load(Ordering::SeqCst) + } + + /// Why no engine was ever built for this handle, or `None` when one was (whether or not it has + /// driven a cycle yet -- see [`ClaimDriverRefusal`]'s doc for the three-way collapse this + /// exists to prevent). + #[must_use] + pub fn refusal(&self) -> Option { + *self + .refusal + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + } + + /// Record why no engine will ever be built on this handle. Called only from + /// [`spawn_claim_driver_if`]'s non-`Spawn` branches and [`run_claim_driver`]'s + /// no-operator-wallet path. + fn set_refusal(&self, reason: ClaimDriverRefusal) { + *self + .refusal + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(reason); + } + + /// Record that one cycle was driven and publish its resulting status. Called only from + /// [`drive`], once per cycle, after `run_cycle` returns. + fn record(&self, status: ClaimStatus) { + *self + .status + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) = status; + self.cycles_driven.fetch_add(1, Ordering::SeqCst); + } +} + +/// The process-wide handle to the running (or never-spawned) claim loop -- one per node process, +/// mirroring how [`crate::state::state_dir`] and friends are process-wide singletons. Initialized +/// lazily to its `Idle`/zero default so a reader (a future RPC, a test) never has to handle +/// "not spawned yet" as a THIRD state distinct from `Idle` -- it is the same state, honestly. +static HANDLE: OnceLock = OnceLock::new(); + +/// The in-process accessor a future RPC (blocked on #3249) reads. Never wired onto the wire here. +#[must_use] +pub fn handle() -> ClaimLoopHandle { + HANDLE.get_or_init(ClaimLoopHandle::default).clone() +} + +/// Drive the claim cadence: sleep `cadence_seconds + jitter` (drawn from `jitter`), run one cycle, +/// record it on `handle`, repeat forever. `now` and `jitter` are injected -- never a global RNG or +/// the clock read directly here -- so the schedule is deterministic and falsifiable under +/// `#[tokio::test(start_paused = true)]` (see this module's doc for why the FIRST pass waits +/// rather than firing immediately, unlike [`crate::self_heal::drive`]). +async fn drive( + mut engine: ClaimEngine, + cadence_seconds: u64, + jitter_seconds: u64, + jitter: &dyn JitterSource, + mut now: impl FnMut() -> u64, + handle: ClaimLoopHandle, +) where + P: ClaimChainPort, + H: DistributorHintSource, +{ + loop { + let interval = next_interval_seconds(cadence_seconds, jitter_seconds, jitter); + tokio::time::sleep(Duration::from_secs(interval)).await; + let t = now(); + engine.run_cycle(t).await; + handle.record(engine.status()); + } +} + +/// A jitter source drawing from the OS CSPRNG (`ring::rand::SystemRandom`, the same primitive +/// [`crate::mirror`]'s signing paths use for randomness in this crate) -- never a global/thread +/// RNG. A CSPRNG failure (the underlying OS call erroring) fails to jitter `0` rather than +/// panicking the driver: the worst case is every node's cadence landing exactly on +/// `cadence_seconds` with no spread, not a crashed claim loop. +struct OsJitter; + +impl JitterSource for OsJitter { + fn jitter_seconds(&self, bound: u64) -> u64 { + if bound == 0 { + return 0; + } + use ring::rand::SecureRandom; + let rng = ring::rand::SystemRandom::new(); + let mut buf = [0u8; 8]; + if rng.fill(&mut buf).is_err() { + return 0; + } + u64::from_le_bytes(buf) % (bound + 1) + } +} + +fn unix_now_seconds() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|d| d.as_secs()) + .unwrap_or(0) +} + +/// A6, money-correctness: this node's own payout puzzle hash as the claim-chain entry slot +/// compares it. `$DIG` is a CAT, so the operator's coins -- and therefore the puzzle hash an +/// `InitiatePayout` should be admitted under -- sit at the canonical CAT wrapping of the owner's +/// inner puzzle hash, never the bare inner hash itself: the SAME derivation +/// [`crate::mirror::lifecycle`]'s `reclaimed_coin_id` and [`crate::mirror::funding::dig_cat_puzzle_hash`]'s +/// own doc use ("the operator's ordinary $DIG coins... sit at the canonical CAT wrapping... never +/// the bare owner puzzle hash"). Picking the unwrapped hash here would make every distributor +/// refuse this node's claims (`claims_refused_payout_mismatch`) while [`super::ClaimLoopState::compute_state`] +/// still reads `Nominal` when nothing is claimable at all -- exactly the misdirection this epic has +/// already measured. See this module's tests for the two-sided proof (a distributor paying to this +/// derivation claims; one paying the unwrapped hash is refused). +#[must_use] +pub fn own_payout_puzzle_hash(owner_inner_puzzle_hash: Bytes32) -> Bytes32 { + crate::mirror::funding::dig_cat_puzzle_hash(owner_inner_puzzle_hash) +} + +/// The real, detached claim-loop task: derive this node's own payout puzzle hash from its operator +/// wallet (the same public, no-unseal-required derivation [`crate::server::spawn_mirror_passes`] +/// falls back to), load [`RewardsClaimConfig`], build a [`ClaimEngine`] against the only +/// production adapter that exists ([`UnavailableClaimChainPort`] -- see this module's doc), and +/// drive it forever. +/// +/// Never called directly by `server.rs` -- see [`spawn_claim_driver_if`], the tested gate that +/// decides WHETHER to call this. `handle` is INJECTED (never the [`handle`] singleton read +/// directly) so a test can drive this against a private, non-shared handle instead of the +/// process-wide one. +async fn run_claim_driver(handle: ClaimLoopHandle) { + let paths = dig_wallet::autoseed::default_paths(); + let Some(owner_inner_puzzle_hash) = dig_wallet::operator_wallet::operator_puzzle_hash(&paths) + else { + tracing::warn!( + target: "rewards_claim", + "no operator wallet is available, so this node has no payout puzzle hash to claim \ + against; the claim loop is NOT started -- rewards_claim.enabled stays true but no \ + cycle will ever run until an operator wallet exists" + ); + handle.set_refusal(ClaimDriverRefusal::NoOperatorWallet); + return; + }; + let own_payout_puzzle_hash = own_payout_puzzle_hash(owner_inner_puzzle_hash); + + let state_dir = crate::state::state_dir(); + let cfg = RewardsClaimConfig::load_from(&state_dir); + // A4/F8: a corrupt config is not a reason to refuse to SPAWN -- `ClaimEngine::run_cycle` + // already fails closed and reports `PersistedStateCorrupt` by name on every cycle until an + // operator fixes or removes the file (see `engine.rs`'s `run_cycle` doc). Refusing to spawn + // here instead would report NOTHING at all, which is the exact silent failure this ticket + // exists to prevent -- a corrupt file must stay visible, not vanish into "never started". + + let engine = ClaimEngine::new( + UnavailableClaimChainPort, + NoHintSource, + own_payout_puzzle_hash, + cfg.max_fee_mojos, + cfg.max_cycle_fee_budget_mojos, + dig_mirror_coin::DIG_ASSET_ID, + ) + .with_rotation_cursor(cfg.rotation_cursor) + .with_persisted_fee_window(&state_dir, cfg.cadence_seconds); + + drive( + engine, + cfg.cadence_seconds, + cfg.jitter_seconds, + &OsJitter, + unix_now_seconds, + handle, + ) + .await; +} + +/// Spawn the real claim-loop task, detached, against `handle` -- injected, never the [`handle`] +/// singleton read from inside, so the only place the process-wide singleton is named is +/// [`spawn_claim_driver_from_config`]. +fn spawn_claim_driver(handle: ClaimLoopHandle) { + tokio::spawn(run_claim_driver(handle)); +} + +/// Why [`spawn_claim_driver_if`] declined to spawn -- named so the caller can log a reason instead +/// of silence (A3). `Disabled` is the ordinary, expected off state (`rewards_claim.enabled = +/// false`); `ChainSyncDisabled` is the one that matters most, because it is reachable with +/// `enabled = true` -- exactly the shape this ticket exists to close: an operator who reads +/// `enabled: true` and believes claims are running, on a node where `enable_chain_sync` is off. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum ClaimDriverDecision { + Spawn, + Disabled, + ChainSyncDisabled, +} + +/// The pure decision behind [`spawn_claim_driver_if`] -- no I/O, no logging, so a test can assert +/// every branch directly. `enabled` gates on its own (SPEC-level opt-out); `enable_chain_sync` is +/// gated the same way `spawn_collateral_census` and `mirror::bond_verify::spawn_bond_verifier_install` +/// already are in `server.rs` -- that flag already means "this node talks to the Chia network", and +/// an integration harness sets it false precisely so nothing dials. +fn decide_claim_driver(enabled: bool, enable_chain_sync: bool) -> ClaimDriverDecision { + if !enabled { + return ClaimDriverDecision::Disabled; + } + if !enable_chain_sync { + return ClaimDriverDecision::ChainSyncDisabled; + } + ClaimDriverDecision::Spawn +} + +/// The single wiring seam `server.rs`'s `serve_with_shutdown` calls (A1's "exact precedent": +/// `self_heal::spawn_driver_if`). `spawn` is invoked exactly when [`decide_claim_driver`] returns +/// `Spawn`; every other branch logs its reason instead of spawning silently (A3) and leaves +/// `handle` at its already-honest `Idle` default -- never a third, undocumented state. +/// +/// `handle` is INJECTED rather than read from the process-wide [`handle`] singleton, so a test can +/// assert the recorded [`ClaimDriverRefusal`] of each branch in-process, on a private handle, with +/// no cross-test interference from a `OnceLock` that outlives the test that touched it. +fn spawn_claim_driver_if( + enabled: bool, + enable_chain_sync: bool, + handle: &ClaimLoopHandle, + spawn: impl FnOnce(), +) { + match decide_claim_driver(enabled, enable_chain_sync) { + ClaimDriverDecision::Spawn => spawn(), + ClaimDriverDecision::Disabled => { + tracing::debug!( + target: "rewards_claim", + "rewards_claim.enabled=false; the claim loop is not started" + ); + handle.set_refusal(ClaimDriverRefusal::Disabled); + } + ClaimDriverDecision::ChainSyncDisabled => { + tracing::warn!( + target: "rewards_claim", + "rewards_claim.enabled=true but enable_chain_sync=false; the claim loop is NOT \ + started -- rewards_claim.enabled is a false statement on this node until chain \ + sync is enabled" + ); + handle.set_refusal(ClaimDriverRefusal::ChainSyncDisabled); + } + } +} + +/// Reads [`RewardsClaimConfig::load`] (the node's own state-dir config) and `enable_chain_sync`, +/// and calls [`spawn_claim_driver_if`] -- the exact one call `serve_with_shutdown` makes. +pub fn spawn_claim_driver_from_config(enable_chain_sync: bool) { + let cfg = RewardsClaimConfig::load(); + // The ONE place the process-wide singleton is read: everything below it takes an injected + // handle so it stays testable in-process. + let process_handle = handle(); + let driver_handle = process_handle.clone(); + spawn_claim_driver_if(cfg.enabled, enable_chain_sync, &process_handle, move || { + spawn_claim_driver(driver_handle); + }); +} + +#[cfg(test)] +mod tests { + use super::*; + use async_trait::async_trait; + use std::sync::atomic::AtomicUsize; + use std::sync::Arc; + + use super::super::port::ClaimPortError; + use super::super::types::{DiscoveredDistributor, OwnEntry}; + + // ---- decide_claim_driver / spawn_claim_driver_if (A3) ---------------------------------- + + #[test] + fn disabled_never_spawns() { + assert_eq!( + decide_claim_driver(false, true), + ClaimDriverDecision::Disabled + ); + assert_eq!( + decide_claim_driver(false, false), + ClaimDriverDecision::Disabled + ); + } + + #[test] + fn enabled_but_chain_sync_off_refuses_named() { + assert_eq!( + decide_claim_driver(true, false), + ClaimDriverDecision::ChainSyncDisabled + ); + } + + #[test] + fn enabled_and_chain_sync_on_spawns() { + assert_eq!(decide_claim_driver(true, true), ClaimDriverDecision::Spawn); + } + + #[test] + fn gate_invokes_spawn_only_on_the_spawn_decision() { + let spawned = Arc::new(AtomicUsize::new(0)); + let handle = ClaimLoopHandle::default(); + + let s = spawned.clone(); + spawn_claim_driver_if(false, true, &handle, || { + s.fetch_add(1, Ordering::SeqCst); + }); + assert_eq!(spawned.load(Ordering::SeqCst), 0, "enabled=false: no spawn"); + + let s = spawned.clone(); + spawn_claim_driver_if(true, false, &handle, || { + s.fetch_add(1, Ordering::SeqCst); + }); + assert_eq!( + spawned.load(Ordering::SeqCst), + 0, + "enabled=true, chain sync off: no spawn" + ); + + let s = spawned.clone(); + spawn_claim_driver_if(true, true, &handle, || { + s.fetch_add(1, Ordering::SeqCst); + }); + assert_eq!( + spawned.load(Ordering::SeqCst), + 1, + "enabled+chain sync: spawns" + ); + } + + /// ACCEPTANCE A3: the gate's two refusals are READABLE off the injected handle and distinct + /// from each other -- not both collapsed into the same zero-cycle `Idle` reading. The third + /// refusal (`NoOperatorWallet`) is proven in + /// `a_missing_operator_wallet_is_a_distinct_named_refusal` below; the fourth truth, + /// "spawned and running but never ticked", is `refusal() == None` with `cycles_driven() == 0`, + /// asserted here and driven past zero in + /// `zero_cycles_before_the_interval_elapses_then_a_counted_number_after`. + #[test] + fn each_refusal_is_readable_and_distinct_on_the_injected_handle() { + let disabled = ClaimLoopHandle::default(); + assert_eq!( + disabled.refusal(), + None, + "nothing refused before the gate runs" + ); + spawn_claim_driver_if(false, true, &disabled, || {}); + assert_eq!(disabled.refusal(), Some(ClaimDriverRefusal::Disabled)); + + let chain_off = ClaimLoopHandle::default(); + spawn_claim_driver_if(true, false, &chain_off, || {}); + assert_eq!( + chain_off.refusal(), + Some(ClaimDriverRefusal::ChainSyncDisabled) + ); + + let spawned = ClaimLoopHandle::default(); + spawn_claim_driver_if(true, true, &spawned, || {}); + assert_eq!( + spawned.refusal(), + None, + "a spawned driver has refused nothing: the fourth truth, `refusal() == None` with a zero cycle count" + ); + assert_eq!(spawned.cycles_driven(), 0); + + // The four readings are pairwise distinct, which is the whole point of A3: three refusals + // plus "running but never ticked" are four different answers to "is this peer being paid". + let readings = [ + disabled.refusal(), + chain_off.refusal(), + Some(ClaimDriverRefusal::NoOperatorWallet), + spawned.refusal(), + ]; + for (i, a) in readings.iter().enumerate() { + for b in &readings[i + 1..] { + assert_ne!(a, b, "two driver refusals must never read the same"); + } + } + } + + /// ACCEPTANCE A3 (third refusal): the no-operator-wallet path records its OWN named reason on + /// the handle rather than leaving `Idle` + zero cycles, and drives no cycle. `run_claim_driver` + /// reads the real default wallet paths, so this asserts the refusal only when this machine + /// genuinely has no operator wallet; where one exists the driver legitimately proceeds and the + /// refusal stays `None` -- either way the reading is a NAMED one, never a silent `Idle`. + #[tokio::test] + async fn a_missing_operator_wallet_is_a_distinct_named_refusal() { + let paths = dig_wallet::autoseed::default_paths(); + if dig_wallet::operator_wallet::operator_puzzle_hash(&paths).is_some() { + return; // this machine HAS an operator wallet; the refusal branch is unreachable here + } + let handle = ClaimLoopHandle::default(); + run_claim_driver(handle.clone()).await; + assert_eq!( + handle.refusal(), + Some(ClaimDriverRefusal::NoOperatorWallet), + "no operator wallet must be a named refusal, not a reassuring Idle" + ); + assert_eq!(handle.cycles_driven(), 0, "and it must drive no cycle"); + } + + // ---- A1 + A2: the anti-silence cycle counter through the real drive() loop ------------- + + /// A fake port whose every call succeeds with an empty/zero answer -- enough to let + /// `run_cycle` reach `Nominal` every time, so the driven-cycle counter is exercised against a + /// REAL completed cycle, not just an early `ChainSourceUnavailable` return. + struct EmptyPort; + + #[async_trait] + impl ClaimChainPort for EmptyPort { + async fn discover_distributors( + &self, + ) -> Result, ClaimPortError> { + Ok(Vec::new()) + } + async fn resolve_launch_comment( + &self, + _launcher_id: Bytes32, + ) -> Result, ClaimPortError> { + Ok(None) + } + async fn reserve_asset_id(&self, _launcher_id: Bytes32) -> Result { + Ok(Bytes32::from([0u8; 32])) + } + async fn payout_threshold(&self, _launcher_id: Bytes32) -> Result { + Ok(0) + } + async fn own_entry( + &self, + _launcher_id: Bytes32, + _payout_puzzle_hash: Bytes32, + ) -> Result, ClaimPortError> { + Ok(None) + } + async fn required_fee_mojos(&self, _launcher_id: Bytes32) -> Result { + Ok(0) + } + async fn submit_initiate_payout( + &self, + _launcher_id: Bytes32, + _payout_puzzle_hash: Bytes32, + _fee_mojos: u64, + ) -> Result<(), ClaimPortError> { + Ok(()) + } + } + + fn empty_engine() -> ClaimEngine { + ClaimEngine::new( + EmptyPort, + NoHintSource, + Bytes32::from([1u8; 32]), + 1, + 10, + Bytes32::from([2u8; 32]), + ) + } + + async fn settle() { + for _ in 0..8 { + tokio::task::yield_now().await; + } + } + + /// ACCEPTANCE A1 + A2 (the anti-silence test): a scheduler that is running but whose interval + /// has never elapsed must read `cycles_driven() == 0` -- NOT "spawn returned", an observed + /// count. Advancing the clock past `cadence_seconds + jitter` must then drive `run_cycle` a + /// counted number of times. + #[tokio::test(start_paused = true)] + async fn zero_cycles_before_the_interval_elapses_then_a_counted_number_after() { + let cadence = 100u64; + let handle = ClaimLoopHandle::default(); + let h = handle.clone(); + let driver = tokio::spawn(async move { + drive( + empty_engine(), + cadence, + 0, + &super::super::cadence::FixedJitter(0), + { + let mut t = 0u64; + move || { + t += cadence; + t + } + }, + h, + ) + .await; + }); + + settle().await; + assert_eq!( + handle.cycles_driven(), + 0, + "THE ANTI-SILENCE TEST: a scheduler that is running but has never fired a cycle must \ + report a driven-cycle count of 0, not silence and not a false 'ran' reading" + ); + + tokio::time::advance(Duration::from_secs(cadence)).await; + settle().await; + assert_eq!( + handle.cycles_driven(), + 1, + "one interval elapsed, one cycle driven" + ); + assert_eq!( + handle.status().state, + super::super::types::ClaimLoopState::Nominal, + "the driven cycle's real outcome is readable, not just its count" + ); + + tokio::time::advance(Duration::from_secs(cadence)).await; + settle().await; + assert_eq!( + handle.cycles_driven(), + 2, + "a second interval drives a second cycle" + ); + + driver.abort(); + } + + // ---- A3: enabled=false vs. enabled=true+gate-refused are both zero-cycle, named states - + + #[tokio::test(start_paused = true)] + async fn disabled_config_never_drives_a_cycle_via_the_configured_seam() { + let dir = tempfile::tempdir().unwrap(); + let cfg = RewardsClaimConfig { + enabled: false, + ..RewardsClaimConfig::default() + }; + cfg.save_to(dir.path()).unwrap(); + // The gate itself (not the full production seam, which reads the process-wide state dir) + // is what's under test here -- see `gate_invokes_spawn_only_on_the_spawn_decision` above + // for the direct proof that `enabled=false` never calls `spawn`. + assert_eq!( + decide_claim_driver(cfg.enabled, true), + ClaimDriverDecision::Disabled + ); + } + + // ---- A6: own_payout_puzzle_hash is the CAT-wrapped hash, proven against the engine ------ + + /// A distributor whose recorded entry is keyed to THIS node's own payout derivation is + /// claimable; one keyed to the bare, unwrapped owner puzzle hash is refused + /// (`PayoutPuzzleHashMismatch`) -- proving `own_payout_puzzle_hash` computes the CAT-wrapped + /// hash the engine's `own_entry` comparison expects, not the raw inner hash. + struct OneDistributorPort { + entry_keyed_to: Bytes32, + dig_asset_id: Bytes32, + } + + #[async_trait] + impl ClaimChainPort for OneDistributorPort { + async fn discover_distributors( + &self, + ) -> Result, ClaimPortError> { + Ok(vec![DiscoveredDistributor { + launcher_id: Bytes32::from([9u8; 32]), + store_id: Bytes32::from([0u8; 32]), + root: Bytes32::from([0u8; 32]), + }]) + } + async fn resolve_launch_comment( + &self, + _launcher_id: Bytes32, + ) -> Result, ClaimPortError> { + Ok(None) + } + async fn reserve_asset_id(&self, _launcher_id: Bytes32) -> Result { + Ok(self.dig_asset_id) + } + async fn payout_threshold(&self, _launcher_id: Bytes32) -> Result { + Ok(1) + } + async fn own_entry( + &self, + _launcher_id: Bytes32, + payout_puzzle_hash: Bytes32, + ) -> Result, ClaimPortError> { + Ok(Some(OwnEntry { + payout_puzzle_hash: self.entry_keyed_to, + counter: 0, + accrued_base_units: 1_000, + })) + // NOTE: `payout_puzzle_hash` (the argument the engine passed in, this node's own + // derivation) is ignored on purpose -- this fake always hands back the entry keyed to + // `entry_keyed_to`, so the engine's OWN comparison (`entry.payout_puzzle_hash != + // self.own_payout_puzzle_hash`) is what decides claimable vs. refused, exactly the + // real chain behaviour this proves against. + } + async fn required_fee_mojos(&self, _launcher_id: Bytes32) -> Result { + Ok(0) + } + async fn submit_initiate_payout( + &self, + _launcher_id: Bytes32, + _payout_puzzle_hash: Bytes32, + _fee_mojos: u64, + ) -> Result<(), ClaimPortError> { + Ok(()) + } + } + + #[tokio::test] + async fn a_distributor_paying_this_nodes_derivation_is_claimable() { + let owner_inner = Bytes32::from([7u8; 32]); + let wrapped = own_payout_puzzle_hash(owner_inner); + let asset_id = Bytes32::from([3u8; 32]); + let mut engine = ClaimEngine::new( + OneDistributorPort { + entry_keyed_to: wrapped, + dig_asset_id: asset_id, + }, + NoHintSource, + wrapped, + 1_000_000, + 10_000_000, + asset_id, + ); + let outcomes = engine.run_cycle(1).await; + assert_eq!(outcomes.len(), 1); + assert!( + matches!( + outcomes[0], + super::super::types::ClaimOutcome::Submitted { .. } + ), + "a distributor keyed to the CAT-wrapped derivation must be claimable, got {:?}", + outcomes[0] + ); + } + + #[tokio::test] + async fn a_distributor_paying_the_unwrapped_hash_is_refused() { + let owner_inner = Bytes32::from([7u8; 32]); + let asset_id = Bytes32::from([3u8; 32]); + let wrapped = own_payout_puzzle_hash(owner_inner); + let mut engine = ClaimEngine::new( + OneDistributorPort { + // Keyed to the RAW inner hash -- the wrong derivation -- not the wrapped one. + entry_keyed_to: owner_inner, + dig_asset_id: asset_id, + }, + NoHintSource, + wrapped, + 1_000_000, + 10_000_000, + asset_id, + ); + let outcomes = engine.run_cycle(1).await; + assert_eq!(outcomes.len(), 1); + assert!( + matches!( + outcomes[0], + super::super::types::ClaimOutcome::PayoutPuzzleHashMismatch { .. } + ), + "a distributor keyed to the unwrapped hash must be refused, got {:?}", + outcomes[0] + ); + } + + // ---- A5: restart safety + clock movement, through the real persisted-config path -------- + + #[tokio::test] + async fn restart_with_a_recent_completion_skips_via_cadence_not_elapsed() { + let dir = tempfile::tempdir().unwrap(); + let cfg = RewardsClaimConfig { + cadence_seconds: 1_000, + last_cycle_completed_at: Some(500), + ..RewardsClaimConfig::default() + }; + cfg.save_to(dir.path()).unwrap(); + + let mut engine = ClaimEngine::new( + EmptyPort, + NoHintSource, + Bytes32::from([1u8; 32]), + 1, + 10, + Bytes32::from([2u8; 32]), + ) + .with_persisted_fee_window(dir.path(), cfg.cadence_seconds); + + let outcomes = engine.run_cycle(900).await; // 900 - 500 = 400 < 1_000 + assert!(outcomes.is_empty()); + assert_eq!( + engine.status().state, + super::super::types::ClaimLoopState::CadenceNotElapsed, + "a crash-restart loop must not immediately re-run a cycle that already ran" + ); + } + + #[tokio::test] + async fn restart_with_an_elapsed_completion_runs_a_cycle() { + let dir = tempfile::tempdir().unwrap(); + let cfg = RewardsClaimConfig { + cadence_seconds: 1_000, + last_cycle_completed_at: Some(500), + ..RewardsClaimConfig::default() + }; + cfg.save_to(dir.path()).unwrap(); + + let mut engine = ClaimEngine::new( + EmptyPort, + NoHintSource, + Bytes32::from([1u8; 32]), + 1, + 10, + Bytes32::from([2u8; 32]), + ) + .with_persisted_fee_window(dir.path(), cfg.cadence_seconds); + + let outcomes = engine.run_cycle(2_000).await; // 2_000 - 500 = 1_500 >= 1_000 + assert!(outcomes.is_empty(), "nothing to claim, but the cycle RAN"); + assert_eq!( + engine.status().state, + super::super::types::ClaimLoopState::Nominal + ); + } + + #[tokio::test] + async fn a_future_dated_completion_fails_closed_not_underflowed() { + let dir = tempfile::tempdir().unwrap(); + let cfg = RewardsClaimConfig { + cadence_seconds: 1_000, + last_cycle_completed_at: Some(10_000), // in the future relative to `now` below + ..RewardsClaimConfig::default() + }; + cfg.save_to(dir.path()).unwrap(); + + let mut engine = ClaimEngine::new( + EmptyPort, + NoHintSource, + Bytes32::from([1u8; 32]), + 1, + 10, + Bytes32::from([2u8; 32]), + ) + .with_persisted_fee_window(dir.path(), cfg.cadence_seconds); + + let outcomes = engine.run_cycle(100).await; // now < last_cycle_completed_at + assert!(outcomes.is_empty()); + assert_eq!( + engine.status().state, + super::super::types::ClaimLoopState::PersistedStateCorrupt, + "a future-dated clock must fail CLOSED, never compute a negative/underflowed interval" + ); + } + + // ---- A4: corrupt = true (via a torn file through load_from), never engine::corrupt set -- + + #[tokio::test] + async fn a_torn_config_file_never_runs_a_cycle() { + let dir = tempfile::tempdir().unwrap(); + std::fs::write(dir.path().join("rewards-claim.json"), b"{ not json").unwrap(); + + let cfg = RewardsClaimConfig::load_from(dir.path()); + assert!( + cfg.corrupt, + "load_from must observe the torn file as corrupt" + ); + + let mut engine = ClaimEngine::new( + EmptyPort, + NoHintSource, + Bytes32::from([1u8; 32]), + 1, + 10, + Bytes32::from([2u8; 32]), + ) + .with_persisted_fee_window(dir.path(), 1_000); + + let outcomes = engine.run_cycle(1).await; + assert!(outcomes.is_empty()); + assert_eq!( + engine.status().state, + super::super::types::ClaimLoopState::PersistedStateCorrupt + ); + } +} diff --git a/crates/dig-node-service/src/rewards_claim/mod.rs b/crates/dig-node-service/src/rewards_claim/mod.rs index 5d5d4a57..7804be4c 100644 --- a/crates/dig-node-service/src/rewards_claim/mod.rs +++ b/crates/dig-node-service/src/rewards_claim/mod.rs @@ -32,19 +32,22 @@ //! prevent (SPEC §2.4): with the unavailable adapter wired, zero claims IS the true state, so the //! status surface must say so by name, not by omission. //! -//! # Not yet wired into node startup (Defect D — stated, not fixed here) -//! Nothing in this codebase constructs a [`ClaimEngine`] outside this module's own tests: there is -//! no scheduler that drives [`ClaimEngine::run_cycle`] on a cadence, and no RPC method exposes -//! [`ClaimStatus`] to an operator, even though [`RewardsClaimConfig::enabled`] defaults to `true`. -//! Wiring this into node startup — picking a concrete [`ClaimChainPort`] adapter, starting the -//! cadence loop, and exposing `ClaimStatus` over RPC — is a separate unit of work with its own -//! review surface, deferred out of this PR on purpose: the only production adapter available today -//! is [`UnavailableClaimChainPort`], and the real one arrives with -//! DIG-Network/dig_ecosystem#3249. Until that wiring lands, this module compiles, is fully tested -//! against the fake chain port, and does nothing in a running node. +//! # Wired into node startup (DIG-Network/dig_ecosystem#3268) +//! [`driver::spawn_claim_driver_from_config`] is the one call `dig-node-service::server`'s +//! `serve_with_shutdown` makes: it is gated on `RewardsClaimConfig::enabled` AND +//! `Config::enable_chain_sync` (the same flag `spawn_collateral_census` and +//! `mirror::bond_verify::spawn_bond_verifier_install` already gate on), and when both are true it +//! spawns a detached task that drives [`ClaimEngine::run_cycle`] on a jittered cadence forever. +//! [`driver::handle`] is the IN-PROCESS accessor a future RPC can read once DIG-Network/dig_ecosystem#3249 +//! lands a real [`ClaimChainPort`] adapter and the `ClaimStatus` wire semantics are re-derived +//! against it — this module puts nothing on the wire itself (see `driver`'s own module doc for +//! why). Until #3249 lands, the only production adapter is still [`UnavailableClaimChainPort`], so +//! every real cycle reports [`ClaimLoopState::ChainSourceUnavailable`] and submits nothing — the +//! honest state, not a silent no-op. mod cadence; mod config; +mod driver; mod engine; mod hints; mod parser; @@ -56,6 +59,7 @@ pub use config::{ RewardsClaimConfig, CLAIM_CADENCE_SECONDS_DEFAULT, CLAIM_CYCLE_FEE_BUDGET_MOJOS_DEFAULT, CLAIM_FEE_CEILING_MOJOS_DEFAULT, }; +pub use driver::{handle, spawn_claim_driver_from_config, ClaimDriverRefusal, ClaimLoopHandle}; pub use engine::ClaimEngine; pub use hints::{DistributorHint, DistributorHintSource, NoHintSource}; pub use parser::parse_launch_comment; diff --git a/crates/dig-node-service/src/server.rs b/crates/dig-node-service/src/server.rs index a646f962..cec18b25 100644 --- a/crates/dig-node-service/src/server.rs +++ b/crates/dig-node-service/src/server.rs @@ -2204,6 +2204,17 @@ where // (a tested unit, #1864) so it cannot be silently flipped to always- or never-spawn. crate::self_heal::spawn_driver_if_service(); + // The peer reward-claim loop (DIG-Network/dig_ecosystem#3268, #3251): drives + // `rewards_claim::ClaimEngine::run_cycle` on a jittered cadence so + // `RewardsClaimConfig::enabled = true` stops being a false statement. Gated the same way the + // census and bond-verifier spawns above are -- `enable_chain_sync` already means "this node + // talks to the Chia network", and a harness sets it false precisely so nothing dials. The + // service-gate lives inside the seam (a tested unit, mirroring `self_heal::spawn_driver_if`) + // so it cannot be silently flipped to always- or never-spawn. The only production chain + // adapter until DIG-Network/dig_ecosystem#3249 lands is `UnavailableClaimChainPort`, so every + // real cycle reports `ChainSourceUnavailable` and submits nothing -- the honest state. + crate::rewards_claim::spawn_claim_driver_from_config(config.enable_chain_sync); + // Best-effort wallet mTLS listener (#368, Sage byte-parity, node-class clients, §5.3). Binds // loopback only on [`DEFAULT_MTLS_PORT`], which is deliberately NOT Sage's own RPC port // (dig-node#260). A bind failure is NON-FATAL — the wallet stays reachable over the plain-HTTP From 2ebb3d9766ad68f220e431d10fc32be397fcc4da Mon Sep 17 00:00:00 2001 From: Michael Taylor Date: Thu, 10 Sep 2026 05:41:17 -0700 Subject: [PATCH 2/7] fix(rewards): silence the deliberately-ignored fake-port argument (#3268) `OneDistributorPort::own_entry` ignores the puzzle hash the engine passes in on purpose -- the fake always returns the entry keyed to `entry_keyed_to` so the ENGINE's own comparison is what decides claimable vs. refused. Named it `_payout_puzzle_hash` (clippy `-D unused-variables`) and moved the rationale onto the parameter, where the next reader meets it. Co-Authored-By: Claude Opus 5 (1M context) --- crates/dig-node-service/src/rewards_claim/driver.rs | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/crates/dig-node-service/src/rewards_claim/driver.rs b/crates/dig-node-service/src/rewards_claim/driver.rs index 8f20e575..ebba1b15 100644 --- a/crates/dig-node-service/src/rewards_claim/driver.rs +++ b/crates/dig-node-service/src/rewards_claim/driver.rs @@ -668,14 +668,16 @@ mod tests { async fn own_entry( &self, _launcher_id: Bytes32, - payout_puzzle_hash: Bytes32, + // Ignored ON PURPOSE (see the NOTE below): this fake always hands back the entry keyed + // to `entry_keyed_to`, so the ENGINE's own comparison decides claimable vs. refused. + _payout_puzzle_hash: Bytes32, ) -> Result, ClaimPortError> { Ok(Some(OwnEntry { payout_puzzle_hash: self.entry_keyed_to, counter: 0, accrued_base_units: 1_000, })) - // NOTE: `payout_puzzle_hash` (the argument the engine passed in, this node's own + // NOTE: `_payout_puzzle_hash` (the argument the engine passed in, this node's own // derivation) is ignored on purpose -- this fake always hands back the entry keyed to // `entry_keyed_to`, so the engine's OWN comparison (`entry.payout_puzzle_hash != // self.own_payout_puzzle_hash`) is what decides claimable vs. refused, exactly the From acc30f388909e4991838534c4e9d6788c0f90684 Mon Sep 17 00:00:00 2001 From: Michael Taylor Date: Thu, 10 Sep 2026 09:18:30 -0700 Subject: [PATCH 3/7] test(rewards): close the untested joint between the claim gate and the drive loop (#3268) `decide_claim_driver` was tested and `drive` was tested, but the production body joining them -- load the config from the state dir, derive the engine, reach `drive` -- was exercised by nothing. That is the exact shape of #594, which shipped a complete, fully-tested and entirely inert claim engine: had this body returned early, built the engine wrong, or never reached `drive`, every test on this change would still have passed and a real node would still never claim. Split `run_claim_driver` on the same `load` / `load_from` pattern the config itself uses: `run_claim_driver_in(state_dir, own_payout_puzzle_hash, port, handle)` holds the whole body and is generic over the port, and `run_claim_driver` is reduced to the wallet-derivation adapter that cannot be reached from a test. Adds two tests through the real body: counted cycles from a written config (zero before the interval, exactly one per interval after), and `UnavailableClaimChainPort` reporting `ChainSourceUnavailable` by name on a driven cycle -- proving the production adapter path is reached, not only a fake. No behaviour change: same config, same engine construction, same port. Co-Authored-By: Claude Opus 5 (1M context) --- .../src/rewards_claim/driver.rs | 141 +++++++++++++++++- 1 file changed, 137 insertions(+), 4 deletions(-) diff --git a/crates/dig-node-service/src/rewards_claim/driver.rs b/crates/dig-node-service/src/rewards_claim/driver.rs index ebba1b15..552c2d2a 100644 --- a/crates/dig-node-service/src/rewards_claim/driver.rs +++ b/crates/dig-node-service/src/rewards_claim/driver.rs @@ -29,6 +29,7 @@ //! #3249 has not landed, so every real cycle this driver runs reports [`super::ClaimLoopState::ChainSourceUnavailable`] //! and submits nothing — the honest state, not an invented adapter. +use std::path::Path; use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Mutex, OnceLock}; use std::time::{Duration, SystemTime, UNIX_EPOCH}; @@ -239,8 +240,34 @@ async fn run_claim_driver(handle: ClaimLoopHandle) { }; let own_payout_puzzle_hash = own_payout_puzzle_hash(owner_inner_puzzle_hash); - let state_dir = crate::state::state_dir(); - let cfg = RewardsClaimConfig::load_from(&state_dir); + run_claim_driver_in( + &crate::state::state_dir(), + own_payout_puzzle_hash, + UnavailableClaimChainPort, + handle, + ) + .await; +} + +/// The whole production body of the claim loop, with every process global it used to read taken as +/// an argument: the state directory it loads [`RewardsClaimConfig`] from, this node's own payout +/// puzzle hash, and the chain `port`. Split out of [`run_claim_driver`] on the same `load` / +/// `load_from` pattern [`RewardsClaimConfig`] itself uses, for one reason: the joint between the +/// tested gate and the tested [`drive`] loop was previously the only UNTESTED link in the chain, +/// and an untested joint is exactly how #594's claim engine shipped complete and inert. +/// +/// Generic over `P` so a test can drive this real body against a fake port; production always +/// passes [`UnavailableClaimChainPort`] (see the module doc -- there is deliberately no second +/// production adapter until #3249 lands). +async fn run_claim_driver_in

( + state_dir: &Path, + own_payout_puzzle_hash: Bytes32, + port: P, + handle: ClaimLoopHandle, +) where + P: ClaimChainPort, +{ + let cfg = RewardsClaimConfig::load_from(state_dir); // A4/F8: a corrupt config is not a reason to refuse to SPAWN -- `ClaimEngine::run_cycle` // already fails closed and reports `PersistedStateCorrupt` by name on every cycle until an // operator fixes or removes the file (see `engine.rs`'s `run_cycle` doc). Refusing to spawn @@ -248,7 +275,7 @@ async fn run_claim_driver(handle: ClaimLoopHandle) { // exists to prevent -- a corrupt file must stay visible, not vanish into "never started". let engine = ClaimEngine::new( - UnavailableClaimChainPort, + port, NoHintSource, own_payout_puzzle_hash, cfg.max_fee_mojos, @@ -256,7 +283,7 @@ async fn run_claim_driver(handle: ClaimLoopHandle) { dig_mirror_coin::DIG_ASSET_ID, ) .with_rotation_cursor(cfg.rotation_cursor) - .with_persisted_fee_window(&state_dir, cfg.cadence_seconds); + .with_persisted_fee_window(state_dir, cfg.cadence_seconds); drive( engine, @@ -871,4 +898,110 @@ mod tests { super::super::types::ClaimLoopState::PersistedStateCorrupt ); } + // ---- The JOINT: the production body itself, not the halves around it ------------------- + + /// Write a `rewards-claim.json` with a fixed cadence and NO jitter, so a composition test can + /// advance the clock by an exact number of seconds and know precisely how many cycles that + /// buys. + fn write_config(dir: &Path, cadence_seconds: u64) { + RewardsClaimConfig { + enabled: true, + cadence_seconds, + jitter_seconds: 0, + ..RewardsClaimConfig::default() + } + .save_to(dir) + .unwrap(); + } + + /// THE COMPOSITION TEST. `decide_claim_driver` was tested, `drive` was tested -- and the + /// production body that joins them (`run_claim_driver_in`: load the config from the state dir, + /// construct the engine, reach `drive`) was tested by NOTHING. That is the same shape as #594, + /// which shipped a complete, fully-tested, entirely INERT claim engine: if this body returned + /// early, built the engine wrong, or never reached `drive`, every other test on this change + /// would still pass and a real node would still never claim. + /// + /// So this drives the REAL body -- the one production calls -- and asserts the anti-silence + /// property through it: zero cycles before the configured interval elapses, then an exactly + /// COUNTED number after. + #[tokio::test(start_paused = true)] + async fn the_production_body_drives_counted_cycles_from_a_written_config() { + let cadence = 100u64; + let dir = tempfile::tempdir().unwrap(); + write_config(dir.path(), cadence); + + let handle = ClaimLoopHandle::default(); + let h = handle.clone(); + let state_dir = dir.path().to_path_buf(); + let driver = tokio::spawn(async move { + run_claim_driver_in(&state_dir, Bytes32::from([1u8; 32]), EmptyPort, h).await; + }); + + settle().await; + assert_eq!( + handle.cycles_driven(), + 0, + "the production body must honour the configured interval: no cycle before it elapses" + ); + + tokio::time::advance(Duration::from_secs(cadence)).await; + settle().await; + assert_eq!( + handle.cycles_driven(), + 1, + "one configured interval elapsed: the production body drove exactly one cycle, proving the joint between the tested gate and the tested drive loop is live" + ); + + tokio::time::advance(Duration::from_secs(cadence)).await; + settle().await; + assert_eq!( + handle.cycles_driven(), + 2, + "and it keeps driving, one cycle per configured interval" + ); + + driver.abort(); + } + + /// The same production body against the port production ACTUALLY passes it + /// ([`UnavailableClaimChainPort`], the only adapter until #3249) reports + /// [`ClaimLoopState::ChainSourceUnavailable`] by name once a cycle has been driven -- the + /// honest state of a real node today. Proves the real adapter path is reached, not only a fake + /// one: a counted cycle whose outcome names the missing chain source, never a reassuring + /// `Nominal` and never silence. + #[tokio::test(start_paused = true)] + async fn the_production_adapter_reports_chain_source_unavailable_by_name() { + let cadence = 100u64; + let dir = tempfile::tempdir().unwrap(); + write_config(dir.path(), cadence); + + let handle = ClaimLoopHandle::default(); + let h = handle.clone(); + let state_dir = dir.path().to_path_buf(); + let driver = tokio::spawn(async move { + run_claim_driver_in( + &state_dir, + Bytes32::from([1u8; 32]), + UnavailableClaimChainPort, + h, + ) + .await; + }); + + tokio::time::advance(Duration::from_secs(cadence)).await; + settle().await; + assert_eq!(handle.cycles_driven(), 1, "one cycle was driven"); + assert_eq!( + handle.status().state, + super::super::types::ClaimLoopState::ChainSourceUnavailable, + "with no chain adapter wired, the driven cycle must name ChainSourceUnavailable" + ); + assert_eq!( + handle.refusal(), + None, + "the loop RAN: an unavailable chain source is a cycle outcome, not a refusal to start" + ); + + driver.abort(); + } } From ac01370edc7e13310d89457a60342261c75d3e49 Mon Sep 17 00:00:00 2001 From: Michael Taylor Date: Thu, 10 Sep 2026 09:56:59 -0700 Subject: [PATCH 4/7] fix(rewards): settle before advancing, and keep the wrapped assertion out of rustfmt's reach (#3268) Two repairs to the new composition tests: - The `ChainSourceUnavailable` test advanced the paused clock before the spawned body had reached its first `sleep`, so the timer was not yet registered and the advance bought no cycle at all -- it read zero cycles, not a driven one. A `settle()` first, mirroring the counted-cycles test. - rustfmt rejoined a `\`-continued assertion message into one line, leaving 14 literal spaces mid-sentence and tripping the repo's own `continuation_guard`. `concat!` states the wrap explicitly, so no formatter pass can reintroduce the run. Co-Authored-By: Claude Opus 5 (1M context) --- crates/dig-node-service/src/rewards_claim/driver.rs | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/crates/dig-node-service/src/rewards_claim/driver.rs b/crates/dig-node-service/src/rewards_claim/driver.rs index 552c2d2a..2efe291e 100644 --- a/crates/dig-node-service/src/rewards_claim/driver.rs +++ b/crates/dig-node-service/src/rewards_claim/driver.rs @@ -949,7 +949,10 @@ mod tests { assert_eq!( handle.cycles_driven(), 1, - "one configured interval elapsed: the production body drove exactly one cycle, proving the joint between the tested gate and the tested drive loop is live" + concat!( + "one configured interval elapsed: the production body drove exactly one cycle, ", + "proving the joint between the tested gate and the tested drive loop is live" + ) ); tokio::time::advance(Duration::from_secs(cadence)).await; @@ -988,6 +991,9 @@ mod tests { .await; }); + // Let the spawned body reach its first `sleep` before advancing: under paused time an + // `advance` that lands before the timer is registered buys no cycle at all. + settle().await; tokio::time::advance(Duration::from_secs(cadence)).await; settle().await; assert_eq!(handle.cycles_driven(), 1, "one cycle was driven"); From e27f6d2d34e05da220fc28ea31275c7d30e1f412 Mon Sep 17 00:00:00 2001 From: Michael Taylor Date: Thu, 10 Sep 2026 14:02:24 -0700 Subject: [PATCH 5/7] feat(rewards): emit a per-cycle event so the claim loop has a reader (#3268) The adversarial gate blocked #605 on this: the PR justified itself by making an inert subsystem loud, but nothing in the shipped binary could hear it. ClaimLoopHandle had no caller outside driver.rs tests, drive() emitted no event, and all three tracing calls fired only on paths where the loop does NOT run -- so on the default path (enabled=true, chain sync on) the observable output was identical to before the PR: silence. Today that silence covers a permanent ChainSourceUnavailable; after #3249 it would also cover Faulted, PersistedStateCorrupt and ClaimableButNotClaiming. log_cycle() now names the state and the cycle count after every cycle -- info for Nominal, warn for everything else, because "this peer is earning nothing and here is why" is a warning, not routine chatter. Tested by capturing the subscriber output rather than asserting the call site exists, since this ticket exists because a guarantee that cannot be observed in a running node is not a guarantee. Refs DIG-Network/dig_ecosystem#3268 Co-Authored-By: Claude Opus 5 (1M context) --- .../src/rewards_claim/driver.rs | 195 +++++++++++++++++- 1 file changed, 193 insertions(+), 2 deletions(-) diff --git a/crates/dig-node-service/src/rewards_claim/driver.rs b/crates/dig-node-service/src/rewards_claim/driver.rs index 2efe291e..1c057c07 100644 --- a/crates/dig-node-service/src/rewards_claim/driver.rs +++ b/crates/dig-node-service/src/rewards_claim/driver.rs @@ -41,7 +41,7 @@ use super::config::RewardsClaimConfig; use super::engine::ClaimEngine; use super::hints::{DistributorHintSource, NoHintSource}; use super::port::{ClaimChainPort, UnavailableClaimChainPort}; -use super::types::ClaimStatus; +use super::types::{ClaimLoopState, ClaimStatus}; /// The in-process accessor onto the running claim loop (SCOPE: never exposed over the wire here). /// Cheap to clone -- every field is an `Arc`-backed handle onto the same shared state. @@ -166,7 +166,45 @@ async fn drive( tokio::time::sleep(Duration::from_secs(interval)).await; let t = now(); engine.run_cycle(t).await; - handle.record(engine.status()); + let status = engine.status(); + handle.record(status); + log_cycle(&status, handle.cycles_driven()); + } +} + +/// Emit the ONE record that makes a driven cycle observable in a running node. +/// +/// Without this, the whole status surface has no reader in a shipped binary: [`handle`] is +/// in-process only and deliberately carries no RPC (deferred to DIG-Network/dig_ecosystem#3249), +/// so a node whose claim loop can never claim a single reward would produce output IDENTICAL to a +/// healthy one -- silence. A status nobody can read is a doc claim, not a measurement. +/// +/// [`ClaimLoopState::Nominal`] is the routine case (`info`). Every other state means this peer is +/// earning nothing and names why, which on a money surface is a warning, not chatter. +fn log_cycle(status: &ClaimStatus, cycles_driven: u64) { + if status.state == ClaimLoopState::Nominal { + tracing::info!( + target: "rewards_claim", + state = ?status.state, + cycles_driven, + distributors_known = status.distributors_known, + distributors_claimable = status.distributors_claimable, + claims_submitted = status.claims_submitted, + "claim cycle complete" + ); + } else { + tracing::warn!( + target: "rewards_claim", + state = ?status.state, + cycles_driven, + distributors_known = status.distributors_known, + distributors_claimable = status.distributors_claimable, + claims_submitted = status.claims_submitted, + concat!( + "claim cycle complete but this node is NOT claiming rewards -- see the named ", + "state for why" + ) + ); } } @@ -1010,4 +1048,157 @@ mod tests { driver.abort(); } + + // ---- the cycle log: the only reader of the status surface in a shipped binary ---------- + + /// An in-memory sink a `tracing_subscriber::fmt` layer renders records into, so a test can + /// assert what a running node would actually print (the same pattern `never_log.rs` and + /// `server.rs` use for their log assertions). + #[derive(Clone)] + struct CapturedLogs(Arc>>); + + impl std::io::Write for CapturedLogs { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0 + .lock() + .expect("the capture buffer") + .extend_from_slice(buf); + Ok(buf.len()) + } + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + + impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for CapturedLogs { + type Writer = CapturedLogs; + fn make_writer(&'a self) -> Self::Writer { + self.clone() + } + } + + impl CapturedLogs { + fn rendered(&self) -> String { + String::from_utf8(self.0.lock().expect("the capture buffer").clone()) + .expect("the rendered lines are utf-8") + } + } + + /// Install a capturing subscriber for the duration of the returned guard. `set_default` is + /// thread-local, and `#[tokio::test]` runs a current-thread runtime, so the driver task + /// spawned below is polled on this very thread and its records land in the buffer. + fn capture_logs() -> (CapturedLogs, tracing::subscriber::DefaultGuard) { + let buffer = CapturedLogs(Arc::new(Mutex::new(Vec::new()))); + let subscriber = tracing_subscriber::fmt() + .with_writer(buffer.clone()) + .with_ansi(false) + .without_time() + .finish(); + let guard = tracing::subscriber::set_default(subscriber); + (buffer, guard) + } + + /// THE ACCEPTANCE BAR: a driven cycle is OBSERVABLE, not merely readable through an + /// in-process handle nothing in the shipped binary calls. A healthy cycle says so at `INFO`, + /// naming its state and its cycle count. + #[tokio::test(start_paused = true)] + async fn a_driven_cycle_emits_an_event_naming_its_state() { + let cadence = 100u64; + let (logs, _guard) = capture_logs(); + let handle = ClaimLoopHandle::default(); + let h = handle.clone(); + let driver = tokio::spawn(async move { + drive( + empty_engine(), + cadence, + 0, + &super::super::cadence::FixedJitter(0), + { + let mut t = 0u64; + move || { + t += cadence; + t + } + }, + h, + ) + .await; + }); + + settle().await; + assert_eq!( + logs.rendered(), + "", + "no interval has elapsed, so there is nothing to report yet" + ); + + tokio::time::advance(Duration::from_secs(cadence)).await; + settle().await; + driver.abort(); + + let rendered = logs.rendered(); + assert!( + rendered.contains("Nominal"), + "the event must NAME the state a reader has to act on; got: {rendered}" + ); + assert!( + rendered.contains("cycles_driven=1"), + "and the cycle count that distinguishes a running loop from a stalled one; got: {}", + rendered + ); + assert!( + rendered.contains("rewards_claim"), + "under the module's own target, so it can be filtered on; got: {rendered}" + ); + } + + /// A cycle that CANNOT claim -- today's real production path, with no chain adapter wired -- + /// must be a WARNING naming the state, not an `INFO` line that reads like health. This is the + /// defect the whole ticket exists to remove: silence covering a permanent inability to earn. + #[tokio::test(start_paused = true)] + async fn a_cycle_that_cannot_claim_warns_and_names_why() { + let cadence = 100u64; + let (logs, _guard) = capture_logs(); + let handle = ClaimLoopHandle::default(); + let h = handle.clone(); + let driver = tokio::spawn(async move { + drive( + ClaimEngine::new( + UnavailableClaimChainPort, + NoHintSource, + Bytes32::from([1u8; 32]), + 1, + 10, + Bytes32::from([2u8; 32]), + ), + cadence, + 0, + &super::super::cadence::FixedJitter(0), + { + let mut t = 0u64; + move || { + t += cadence; + t + } + }, + h, + ) + .await; + }); + + settle().await; + tokio::time::advance(Duration::from_secs(cadence)).await; + settle().await; + driver.abort(); + + let rendered = logs.rendered(); + assert!( + rendered.contains("WARN"), + "earning nothing is a warning, not routine chatter; got: {rendered}" + ); + assert!( + rendered.contains("ChainSourceUnavailable"), + "and it must name WHY this node is not claiming; got: {rendered}" + ); + } } From d1f9cac9cff7e963802a159e1220be154cddb63f Mon Sep 17 00:00:00 2001 From: Michael Taylor Date: Thu, 10 Sep 2026 14:19:07 -0700 Subject: [PATCH 6/7] fix(rewards): stop a u64::MAX jitter bound panicking the claim driver `OsJitter::jitter_seconds` computed `bound + 1` for its modulus. `jitter_seconds` comes from the node's persisted `rewards_claim` config and is not clamped, so a config carrying `u64::MAX` overflow-panicked inside the detached claim-driver task -- which has no restart and emits no further log output, so the claim loop would die silently for the rest of the process lifetime. `saturating_add(1)` keeps the draw within `0..=bound` for every input; the composed `next_interval_seconds` range is unchanged. Co-Authored-By: Claude Opus 5 (1M context) --- .../src/rewards_claim/driver.rs | 19 ++++++++++++++++++- 1 file changed, 18 insertions(+), 1 deletion(-) diff --git a/crates/dig-node-service/src/rewards_claim/driver.rs b/crates/dig-node-service/src/rewards_claim/driver.rs index 1c057c07..e562efca 100644 --- a/crates/dig-node-service/src/rewards_claim/driver.rs +++ b/crates/dig-node-service/src/rewards_claim/driver.rs @@ -226,7 +226,11 @@ impl JitterSource for OsJitter { if rng.fill(&mut buf).is_err() { return 0; } - u64::from_le_bytes(buf) % (bound + 1) + // `bound.saturating_add(1)` rather than `bound + 1`: `jitter_seconds` comes from the + // persisted config unclamped, so `u64::MAX` reaches here and `+ 1` would overflow-panic + // inside the detached driver task -- killing the claim loop silently for the process + // lifetime. Saturating keeps the draw in `0..=bound` for every input. + u64::from_le_bytes(buf) % bound.saturating_add(1) } } @@ -1201,4 +1205,17 @@ mod tests { "and it must name WHY this node is not claiming; got: {rendered}" ); } + + /// `jitter_seconds` is read from the persisted config WITHOUT a clamp, so the maximum `u64` + /// reaches `OsJitter`. An overflowing `bound + 1` there panics the detached driver task, + /// which never restarts -- the claim loop would die silently for the process lifetime. + #[test] + fn an_unclamped_max_jitter_bound_does_not_panic_the_driver() { + let bound = std::hint::black_box(u64::MAX); + let offset = OsJitter.jitter_seconds(bound); + assert!( + offset <= bound, + "the draw must stay within 0..=bound; got {offset}" + ); + } } From 5fce847770455862f39d022f0a9a9285fb8e1375 Mon Sep 17 00:00:00 2001 From: Michael Taylor Date: Thu, 10 Sep 2026 14:24:35 -0700 Subject: [PATCH 7/7] fix(rewards): sanitize the claim schedule so no config value silently disables the loop `next_interval_seconds` saturates instead of panicking, so a persisted `jitter_seconds = u64::MAX` no longer crashes the driver -- it schedules the next cycle ~585 billion years out. The claim loop then never fires again: no cycle, no `log_cycle` line, and a permanent, reassuring `0` cycle count. That is #594's inert-but-green shape reopened one level up, in the config file. `run_claim_driver_in` now sanitizes both schedule fields where it reads them, before either reaches the engine's fee window or `drive`: - `CLAIM_SCHEDULE_SECONDS_MAX = 31 * 24 * 60 * 60` (31 days) -- above every documented default (86,400s cadence, 3,600s jitter) and above "claim monthly", while excluding everything that means never. - out of range (or a zero cadence, which would busy-loop) substitutes the published default and emits `tracing::warn!` naming the field, the rejected value and the substituted one. Nothing is accepted silently. `config.rs` is untouched: it keeps reporting what is on disk. Co-Authored-By: Claude Opus 5 (1M context) --- .../src/rewards_claim/driver.rs | 170 +++++++++++++++++- 1 file changed, 165 insertions(+), 5 deletions(-) diff --git a/crates/dig-node-service/src/rewards_claim/driver.rs b/crates/dig-node-service/src/rewards_claim/driver.rs index e562efca..858b978f 100644 --- a/crates/dig-node-service/src/rewards_claim/driver.rs +++ b/crates/dig-node-service/src/rewards_claim/driver.rs @@ -36,8 +36,8 @@ use std::time::{Duration, SystemTime, UNIX_EPOCH}; use chia_protocol::Bytes32; -use super::cadence::{next_interval_seconds, JitterSource}; -use super::config::RewardsClaimConfig; +use super::cadence::{next_interval_seconds, JitterSource, CLAIM_JITTER_SECONDS_DEFAULT}; +use super::config::{RewardsClaimConfig, CLAIM_CADENCE_SECONDS_DEFAULT}; use super::engine::ClaimEngine; use super::hints::{DistributorHintSource, NoHintSource}; use super::port::{ClaimChainPort, UnavailableClaimChainPort}; @@ -301,6 +301,69 @@ async fn run_claim_driver(handle: ClaimLoopHandle) { /// Generic over `P` so a test can drive this real body against a fake port; production always /// passes [`UnavailableClaimChainPort`] (see the module doc -- there is deliberately no second /// production adapter until #3249 lands). +/// The largest schedule value this driver will honour, in seconds: 31 days. Chosen to sit +/// comfortably above every documented default -- [`CLAIM_CADENCE_SECONDS_DEFAULT`] is 86,400s +/// (1 day) and [`CLAIM_JITTER_SECONDS_DEFAULT`] is 3,600s (1 hour) -- and above any plausible +/// operator choice ("claim monthly" is 30 days), while excluding every value that means NEVER. +/// +/// # Why a ceiling exists at all +/// [`next_interval_seconds`] saturates rather than panicking, so a persisted +/// `jitter_seconds = u64::MAX` (or a cadence of the same shape) no longer crashes the node -- it +/// schedules the next cycle roughly 585 billion years out. That is strictly WORSE than a panic for +/// this ticket: the claim loop never fires again, so no cycle, no `log_cycle` line, and the +/// cycle counter reads a permanent, reassuring `0`. #594 shipped an engine that was inert and +/// green; a config value must not be able to put this driver back in that state silently. +const CLAIM_SCHEDULE_SECONDS_MAX: u64 = 31 * 24 * 60 * 60; + +/// Replace a schedule value that would switch the loop off (or spin it) with its documented +/// default, saying so at `WARN` -- never silently accept it, and never silently accept the +/// default either. Returns `(cadence_seconds, jitter_seconds)` fit to schedule with. +/// +/// A zero cadence is rejected for the opposite reason to a huge one: it would busy-loop the claim +/// engine as fast as the runtime can poll it. A zero JITTER is legitimate (it means "no jitter") +/// and is left alone. +fn sanitized_schedule(cadence_seconds: u64, jitter_seconds: u64) -> (u64, u64) { + let cadence = if cadence_seconds == 0 || cadence_seconds > CLAIM_SCHEDULE_SECONDS_MAX { + tracing::warn!( + target: "rewards_claim", + field = "cadence_seconds", + rejected = cadence_seconds, + substituted = CLAIM_CADENCE_SECONDS_DEFAULT, + max = CLAIM_SCHEDULE_SECONDS_MAX, + "{}", + concat!( + "rewards_claim.cadence_seconds is outside the honoured range and was IGNORED; ", + "the documented default is used instead -- a value that large would stop the ", + "claim loop from ever firing again, and zero would busy-loop it" + ) + ); + CLAIM_CADENCE_SECONDS_DEFAULT + } else { + cadence_seconds + }; + + let jitter = if jitter_seconds > CLAIM_SCHEDULE_SECONDS_MAX { + tracing::warn!( + target: "rewards_claim", + field = "jitter_seconds", + rejected = jitter_seconds, + substituted = CLAIM_JITTER_SECONDS_DEFAULT, + max = CLAIM_SCHEDULE_SECONDS_MAX, + "{}", + concat!( + "rewards_claim.jitter_seconds is outside the honoured range and was IGNORED; ", + "the documented default is used instead -- a value that large saturates the ", + "next interval and the claim loop would never fire again" + ) + ); + CLAIM_JITTER_SECONDS_DEFAULT + } else { + jitter_seconds + }; + + (cadence, jitter) +} + async fn run_claim_driver_in

( state_dir: &Path, own_payout_puzzle_hash: Bytes32, @@ -316,6 +379,11 @@ async fn run_claim_driver_in

( // here instead would report NOTHING at all, which is the exact silent failure this ticket // exists to prevent -- a corrupt file must stay visible, not vanish into "never started". + // The config is operator-writable and unclamped at rest (`config.rs` deliberately reports what + // is on disk). Sanitize HERE, at the read, before either value can reach the scheduler. + let (cadence_seconds, jitter_seconds) = + sanitized_schedule(cfg.cadence_seconds, cfg.jitter_seconds); + let engine = ClaimEngine::new( port, NoHintSource, @@ -325,12 +393,12 @@ async fn run_claim_driver_in

( dig_mirror_coin::DIG_ASSET_ID, ) .with_rotation_cursor(cfg.rotation_cursor) - .with_persisted_fee_window(state_dir, cfg.cadence_seconds); + .with_persisted_fee_window(state_dir, cadence_seconds); drive( engine, - cfg.cadence_seconds, - cfg.jitter_seconds, + cadence_seconds, + jitter_seconds, &OsJitter, unix_now_seconds, handle, @@ -1218,4 +1286,96 @@ mod tests { "the draw must stay within 0..=bound; got {offset}" ); } + + // ---- the config-read sanitizer: a value must not be able to switch the loop off --------- + + /// Write a config with BOTH schedule fields chosen by the caller, so a test can persist a + /// value production would otherwise honour to the letter. + fn write_schedule_config(dir: &Path, cadence_seconds: u64, jitter_seconds: u64) { + RewardsClaimConfig { + enabled: true, + cadence_seconds, + jitter_seconds, + ..RewardsClaimConfig::default() + } + .save_to(dir) + .unwrap(); + } + + /// An out-of-range `cadence_seconds` must be REPLACED by the documented default, and the + /// substitution must be visible: a value this large means "never fire again", and silently + /// honouring it reopens #594's inert-but-green shape one level up, in the config file. + #[test] + fn an_out_of_range_cadence_is_replaced_by_the_default_and_warned() { + let (logs, _guard) = capture_logs(); + + let rejected = std::hint::black_box(u64::MAX); + let (cadence, jitter) = sanitized_schedule(rejected, 0); + + assert_eq!( + cadence, CLAIM_CADENCE_SECONDS_DEFAULT, + "an out-of-range cadence must fall back to the documented default" + ); + assert_eq!(jitter, 0, "a legitimate zero jitter is left alone"); + + let rendered = logs.rendered(); + assert!( + rendered.contains("WARN"), + "ignoring a configured value is a warning, not routine chatter; got: {rendered}" + ); + assert!( + rendered.contains("cadence_seconds"), + "the warning must name the FIELD that was ignored; got: {rendered}" + ); + assert!( + rendered.contains(&rejected.to_string()) + && rendered.contains(&CLAIM_CADENCE_SECONDS_DEFAULT.to_string()), + "and both the rejected and the substituted value; got: {rendered}" + ); + } + + /// The same for `jitter_seconds` -- and, through the REAL production body, that the loop still + /// drives counted cycles instead of never firing again. Without the sanitizer, + /// `next_interval_seconds` saturates on this value and no cycle is ever driven: green, silent + /// and unpaid. + #[tokio::test(start_paused = true)] + async fn an_out_of_range_jitter_is_replaced_and_the_loop_still_drives_cycles() { + let (logs, _guard) = capture_logs(); + + let cadence = 100u64; + let dir = tempfile::tempdir().unwrap(); + write_schedule_config(dir.path(), cadence, std::hint::black_box(u64::MAX)); + + let handle = ClaimLoopHandle::default(); + let h = handle.clone(); + let state_dir = dir.path().to_path_buf(); + let driver = tokio::spawn(async move { + run_claim_driver_in(&state_dir, Bytes32::from([1u8; 32]), EmptyPort, h).await; + }); + + settle().await; + // The substituted jitter is the DEFAULT hour, so one interval is at most cadence + 3600s. + tokio::time::advance(Duration::from_secs(cadence + CLAIM_JITTER_SECONDS_DEFAULT)).await; + settle().await; + + assert!( + handle.cycles_driven() >= 1, + concat!( + "an out-of-range jitter must not switch the claim loop off: with the default ", + "substituted, at least one cycle is driven within cadence + the default jitter" + ) + ); + + let rendered = logs.rendered(); + assert!( + rendered.contains("WARN") && rendered.contains("jitter_seconds"), + "the ignored jitter field must be named at WARN; got: {rendered}" + ); + assert!( + rendered.contains(&CLAIM_JITTER_SECONDS_DEFAULT.to_string()), + "and the substituted default must be readable; got: {rendered}" + ); + + driver.abort(); + } }