diff --git a/crates/contributor-rewards/CHANGELOG.md b/crates/contributor-rewards/CHANGELOG.md index fd61d020..88309abd 100644 --- a/crates/contributor-rewards/CHANGELOG.md +++ b/crates/contributor-rewards/CHANGELOG.md @@ -7,6 +7,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +- fix(contributor-rewards): resolve the Solana epoch for a timestamp from real block times instead of dividing wall clock by a hardcoded 400ms slot duration. The old estimate drifted about 30k slots per day of lookback and picked the wrong epoch near a boundary, and no fixed constant survives the SIMD-0525 rollout. That epoch selects the leader schedule rewards are computed against, so the search now errors rather than returning a wrong answer: a backfill older than the endpoint's ledger retention fails on the `ingestor::demand` path instead of silently mis-estimating (malbeclabs/infra#2317) +- fix(contributor-rewards): `snapshot` validates before writing. It warns and continues when the leader schedule cannot be fetched, but every consumer rejects a snapshot without one, so the command exited 0 having written an unusable file under the canonical name and a `snapshot` then `export-shapley` chain failed a step late. Pre-existing, but reachable now that resolving the Solana epoch depends on block-time reads (malbeclabs/infra#2317) - migrate to Solana 3.0: workspace `solana-*` crates and `solana-sdk` move to the 3.0 line, `solana-program-test` to 3.0.12, and the doublezero SDK git-deps repin from `client/v0.27.1` to the malbeclabs/doublezero#3830 merge revision (malbeclabs/infra#1853) - release artifact now builds as a static `x86_64-unknown-linux-musl` binary (malbeclabs/infra#1853) - TLS for HTTP clients moves from openssl to rustls; trust roots are the bundled webpki Mozilla set plus the host OS certificate store, so OS-installed private CAs remain trusted (malbeclabs/infra#1853) diff --git a/crates/contributor-rewards/src/cli/snapshot.rs b/crates/contributor-rewards/src/cli/snapshot.rs index 45b280b6..6e8f7c46 100644 --- a/crates/contributor-rewards/src/cli/snapshot.rs +++ b/crates/contributor-rewards/src/cli/snapshot.rs @@ -252,31 +252,22 @@ pub async fn create_snapshot( fetcher.dz_rpc_client.clone(), fetcher.solana_read_client.clone(), ); - let solana_epoch = match epoch_finder - .find_epoch_at_timestamp(fetch_data.start_us) + // fetch_leader_schedule resolves the Solana epoch itself and reports which + // one it used, so taking the epoch from its result avoids running the + // chain-verified epoch search twice over the same timestamp. + let leader_schedule = match epoch_finder + .fetch_leader_schedule(fetch_epoch, fetch_data.start_us) .await { - Ok(epoch) => Some(epoch), + Ok(schedule) => Some(schedule), Err(e) => { - warn!("Failed to determine Solana epoch: {}", e); + warn!("Failed to get leader schedule: {}", e); None } }; - - let leader_schedule = if solana_epoch.is_some() { - match epoch_finder - .fetch_leader_schedule(fetch_epoch, fetch_data.start_us) - .await - { - Ok(schedule) => Some(schedule), - Err(e) => { - warn!("Failed to get leader schedule: {}", e); - None - } - } - } else { - None - }; + let solana_epoch = leader_schedule + .as_ref() + .map(|schedule| schedule.solana_epoch); // Create metadata let metadata = SnapshotMetadata { @@ -327,6 +318,11 @@ pub async fn create_snapshot( Network::Devnet => "dn", }; + // Refuse to write a snapshot no consumer can read. Without this the command + // exits 0 having left an unusable file under the canonical name, and the + // failure surfaces a step later in whatever reads it next. + snapshot.validate()?; + // Export: local override or configured storage if local_file.is_some() || local_dir.is_some() { // Save to local filesystem (ignores storage backend config) diff --git a/crates/contributor-rewards/src/ingestor/epoch.rs b/crates/contributor-rewards/src/ingestor/epoch.rs index b6992d0c..a136580f 100644 --- a/crates/contributor-rewards/src/ingestor/epoch.rs +++ b/crates/contributor-rewards/src/ingestor/epoch.rs @@ -1,19 +1,24 @@ //! Epoch calculation utilities for mapping timestamps to Solana epochs //! //! This module provides functionality to: -//! - Calculate Solana epochs from slots //! - Estimate slots from timestamps //! - Find epochs corresponding to specific timestamps use std::{collections::BTreeMap, sync::Arc, time::Duration}; -use anyhow::{Result, anyhow, bail}; +use anyhow::{Context, Result, anyhow, bail, ensure}; use backon::{ExponentialBuilder, Retryable}; use chrono::Utc; use doublezero_solana_client_tools::rpc::DoubleZeroLedgerConnection; use serde::{Deserialize, Serialize}; use solana_client::{ - client_error::ClientError as SolanaClientError, nonblocking::rpc_client::RpcClient, + client_error::{ClientError as SolanaClientError, ClientErrorKind}, + nonblocking::rpc_client::RpcClient, + rpc_custom_error::{ + JSON_RPC_SERVER_ERROR_BLOCK_CLEANED_UP, JSON_RPC_SERVER_ERROR_BLOCK_NOT_AVAILABLE, + JSON_RPC_SERVER_ERROR_LONG_TERM_STORAGE_SLOT_SKIPPED, JSON_RPC_SERVER_ERROR_SLOT_SKIPPED, + }, + rpc_request::RpcError, }; use solana_sdk::epoch_schedule::EpochSchedule; use tracing::{debug, info}; @@ -23,8 +28,47 @@ use crate::cli::{ traits::Exportable, }; -/// Approximate slot duration in microseconds (400ms) -pub const SLOT_DURATION_US: u64 = 400_000; +// Seed slot duration for the epoch search in `find_epoch_at_timestamp`. What +// matters is the direction of the error, not its size, so this is deliberately +// slower than any real cluster rate rather than close to one. +// +// Dividing elapsed wall clock by a value above the real rate under-counts the +// slots elapsed, seeding at or after the epoch being looked for, so the search +// walks backward only and never probes older than the answer's own epoch. A +// value below the real rate probes older than the target instead, which can fall +// outside the endpoint's retention and fail the search for a target that is +// itself readable. The margin costs search steps, roughly one extra epoch of +// backward walk per two days of lookback. +const SEED_SLOT_DURATION_US: u64 = 500_000; + +// 400_000 is mainnet-beta's rate until epoch 1020 (2026-08-21) and the slowest +// any cluster runs; every step of the SIMD-0525 rollout only lowers it. +const _: () = assert!(SEED_SLOT_DURATION_US > 400_000); + +// `getSlot` can name a slot that has no block yet, so the chain tip lookup walks +// backward from it. A tip that needs more than this many slots to find a block is +// a stalled cluster rather than a run of skipped slots. +const MAX_CHAIN_TIP_SEARCH_SLOTS: u64 = 128; + +// How far past an epoch's first slot that epoch's first block may be and still be +// trusted to date the epoch. +// +// `getBlocksWithLimit` steps over a gap in the endpoint's own block history +// without saying that it did, so the slot it returns cannot by itself tell a +// routine run of skipped boundary slots from a node restored from a snapshot +// partway through the epoch. Distance is the tell: dating an epoch by a block +// that far in overstates when the epoch began by the length of the gap, and +// timestamps inside the gap then resolve to the previous epoch and pick up its +// leader schedule. +// +// 432 slots is 0.1% of a mainnet epoch, around two and a half minutes, chosen to +// sit far above any routine skip run and far below a gap worth guessing through. +const MAX_BOUNDARY_SKIP_SLOTS: u64 = 432; + +// Each search step moves the candidate by one epoch. The seed is normally within +// an epoch or two of correct, so this cap exists only to bound a pathological +// seed rather than loop forever. +const MAX_EPOCH_SEARCH_STEPS: usize = 16; // key: validator_pk, val: slot count pub type LeaderScheduleMap = BTreeMap; @@ -48,12 +92,109 @@ impl Exportable for LeaderSchedule { } } -/// Calculate the epoch for a given slot using the epoch schedule +/// Report whether an RPC error means the slot has no block to report a time for, +/// as opposed to the request itself failing. /// -/// This handles normal epochs & ignores warmup period (that's relevant only in genesis) -pub fn calculate_epoch_from_slot(slot: u64, schedule: &EpochSchedule) -> u64 { - // Normal epoch calculation - ((slot - schedule.first_normal_slot) / schedule.slots_per_epoch) + schedule.first_normal_epoch +/// `getBlockTime` says this in three ways depending on what the endpoint has +/// behind it, and missing any one of them fails the search on endpoints of that +/// shape: a validator with long term storage answers with a coded skipped-slot +/// error, one without it answers with a JSON `null` that +/// `RpcClient::get_block_time` turns into `RpcError::ForUser("Block Not +/// Found: ...")`, and a slot the endpoint has not rooted yet is reported as +/// block-not-available. +/// +/// `JSON_RPC_SERVER_ERROR_BLOCK_CLEANED_UP` is deliberately excluded, because a +/// pruned ledger means the answer is unknowable rather than absent. See +/// [`is_block_cleaned_up`]. +fn is_block_unavailable(err: &SolanaClientError) -> bool { + match err.kind() { + ClientErrorKind::RpcError(RpcError::RpcResponseError { code, .. }) => matches!( + *code, + JSON_RPC_SERVER_ERROR_SLOT_SKIPPED + | JSON_RPC_SERVER_ERROR_LONG_TERM_STORAGE_SLOT_SKIPPED + | JSON_RPC_SERVER_ERROR_BLOCK_NOT_AVAILABLE + ), + ClientErrorKind::RpcError(RpcError::ForUser(message)) => { + message.starts_with("Block Not Found") + } + _ => false, + } +} + +/// Report whether an RPC error means the ledger no longer holds the slot. +/// +/// This fails a lookup rather than counting as an absent block, since the block +/// time is gone rather than nonexistent. It is just as settled, so it is not +/// worth retrying either. +fn is_block_cleaned_up(err: &SolanaClientError) -> bool { + matches!( + err.kind(), + ClientErrorKind::RpcError(RpcError::RpcResponseError { + code: JSON_RPC_SERVER_ERROR_BLOCK_CLEANED_UP, + .. + }) + ) +} + +/// Report whether an RPC error is settled, meaning a retry would sleep through +/// the backoff schedule only to be told the same thing. +fn is_settled_block_error(err: &SolanaClientError) -> bool { + is_block_unavailable(err) || is_block_cleaned_up(err) +} + +/// Report whether a block at `first_block_slot` is close enough to `first_slot`, +/// the first slot of an epoch, to date when that epoch began. +/// +/// `MAX_BOUNDARY_SKIP_SLOTS` is the real bound. The next epoch's first slot only +/// binds first on a cluster whose epochs are shorter than that budget, such as a +/// local test validator. +fn can_date_epoch_start( + first_slot: u64, + first_block_slot: u64, + next_epoch_first_slot: u64, +) -> bool { + let last_datable_slot = (first_slot + MAX_BOUNDARY_SKIP_SLOTS).min(next_epoch_first_slot); + + first_block_slot < last_datable_slot +} + +#[derive(Debug, PartialEq, Eq)] +enum EpochSearchStep { + Earlier, + Later, + Found, +} + +/// Decide which way the epoch search should move from its current candidate. +/// +/// `epoch_start_time` and `next_epoch_start_time` are the block times of the +/// first block in the candidate epoch and in the epoch after it, in seconds. +/// `None` means that epoch has produced no block yet and so bounds nothing: it +/// rules the candidate out, or, for the next epoch, means the candidate has no +/// upper bound. Accepting the candidate unbounded is only sound because the +/// caller has already established that the target is at or before the chain tip. +/// +/// The lower bound is inclusive and the upper bound is exclusive, so a timestamp +/// falling exactly on an epoch's first block time belongs to that epoch. +fn decide_epoch_search_step( + target_time: i64, + epoch_start_time: Option, + next_epoch_start_time: Option, +) -> EpochSearchStep { + let Some(epoch_start_time) = epoch_start_time else { + return EpochSearchStep::Earlier; + }; + + if target_time < epoch_start_time { + return EpochSearchStep::Earlier; + } + + match next_epoch_start_time { + Some(next_epoch_start_time) if target_time >= next_epoch_start_time => { + EpochSearchStep::Later + } + _ => EpochSearchStep::Found, + } } /// Estimate the slot at a given timestamp based on current slot and time @@ -70,7 +211,7 @@ pub fn estimate_slot_from_timestamp( // Calculate approximate slot at the given timestamp let time_diff_us = current_time_us - timestamp_us; - let slots_ago = time_diff_us / SLOT_DURATION_US; + let slots_ago = time_diff_us / SEED_SLOT_DURATION_US; if slots_ago > current_slot { bail!("Timestamp {timestamp_us} is too far in the past"); @@ -158,9 +299,160 @@ impl EpochFinder { .expect("solana_schedule cannot be none")) } + /// Get a slot's block time in seconds, or `Ok(None)` when the slot has no + /// block to report a time for. + /// + /// Transport failures are retried; a slot that produced no block and a pruned + /// ledger are not, since both are settled and retrying only burns the backoff + /// schedule before arriving at the same answer. + async fn try_get_block_time(&self, slot: u64) -> Result> { + let block_time = (|| async { self.solana_read_client.get_block_time(slot).await }) + .retry(&ExponentialBuilder::default().with_jitter()) + .when(|err: &SolanaClientError| !is_settled_block_error(err)) + .notify(|err: &SolanaClientError, dur: Duration| { + info!( + "retrying get_block_time error: {:?} with sleeping {:?}", + err, dur + ) + }) + .await; + + match block_time { + Ok(block_time) => Ok(Some(block_time)), + Err(err) if is_block_unavailable(&err) => Ok(None), + Err(err) => { + Err(err).with_context(|| format!("Failed to get block time for Solana slot {slot}")) + } + } + } + + /// Find the block time in seconds of the first block in `epoch`. + /// + /// Returns `Ok(None)` only when the epoch has produced no block, meaning its + /// first slot is past `chain_tip_slot`. Measuring against the chain tip rather + /// than `getSlot` matters because `getSlot` can name a blockless slot, which + /// would make a just-started epoch whose opening slots were all skipped look + /// like an endpoint missing history and fail the search. + /// + /// Every other way of failing to date the epoch is an error rather than a + /// `None`, and `chain_tip_slot` is what makes that sound: it has a block, so + /// an epoch starting at or before it must have one too, leaving missing + /// history as the only reading of an empty answer. Returning `None` there + /// would walk the search backward past the right answer. + /// + /// [`can_date_epoch_start`] guards the one thing `getBlocksWithLimit` will not + /// report, which is that it stepped over a gap in the endpoint's own history. + /// + /// The block time is returned as is, with no back estimation of the skipped + /// slots before it: the first block's time is the epoch's effective start for + /// this purpose, so subtracting an estimate would only add error. This is a + /// deliberate difference from `estimate_block_time_for_skipped_slot` in + /// `validator-debt/src/rpc.rs`, which does subtract one. + async fn try_epoch_start_block_time( + &self, + schedule: &EpochSchedule, + epoch: u64, + chain_tip_slot: u64, + ) -> Result> { + let first_slot = schedule.get_first_slot_in_epoch(epoch); + if first_slot > chain_tip_slot { + return Ok(None); + } + + let first_block_slot = (|| async { + self.solana_read_client + .get_blocks_with_limit(first_slot, 1) + .await + }) + .retry(&ExponentialBuilder::default().with_jitter()) + .when(|err: &SolanaClientError| !is_settled_block_error(err)) + .notify(|err: &SolanaClientError, dur: Duration| { + info!( + "retrying get_blocks_with_limit error: {:?} with sleeping {:?}", + err, dur + ) + }) + .await + .with_context(|| format!("Failed to find the first block of Solana epoch {epoch}"))? + .first() + .copied() + .with_context(|| { + format!( + "Solana endpoint reports no block at or after slot {first_slot}, the first slot \ + of epoch {epoch}, even though slot {chain_tip_slot} has one. The endpoint is \ + most likely missing block history for that range" + ) + })?; + + ensure!( + can_date_epoch_start( + first_slot, + first_block_slot, + schedule.get_first_slot_in_epoch(epoch + 1) + ), + "The first block at or after slot {first_slot} is slot {first_block_slot}, more than \ + {MAX_BOUNDARY_SKIP_SLOTS} slots into Solana epoch {epoch}. The endpoint is most \ + likely missing block history there, and dating the epoch by that block would place \ + its start late enough that timestamps inside the gap resolve to the previous epoch" + ); + + let block_time = self + .try_get_block_time(first_block_slot) + .await? + .with_context(|| { + format!( + "Solana slot {first_block_slot} was reported as the first block of epoch \ + {epoch} but has no block time" + ) + })?; + + Ok(Some(block_time)) + } + + /// Find the newest slot at or before `current_slot` that has a block, and + /// return it with its block time in seconds. + /// + /// The walk runs backward because `getSlot` can name a slot with no block yet, + /// whether skipped or not yet caught up to. It resolves on the first probe in + /// the ordinary case. + /// + /// The time bounds how recent a timestamp the search accepts; the slot is what + /// [`Self::try_epoch_start_block_time`] measures against to tell "no block + /// yet" apart from "endpoint is missing history". + async fn try_chain_tip_block(&self, current_slot: u64) -> Result<(u64, i64)> { + let oldest_slot_to_search = current_slot.saturating_sub(MAX_CHAIN_TIP_SEARCH_SLOTS - 1); + + for slot in (oldest_slot_to_search..=current_slot).rev() { + if let Some(block_time) = self.try_get_block_time(slot).await? { + return Ok((slot, block_time)); + } + } + + bail!( + "No block within {MAX_CHAIN_TIP_SEARCH_SLOTS} slots at or before the current slot \ + {current_slot}" + ) + } + /// Find the Solana epoch that was active at a given timestamp /// - /// This uses the Solana network to map timestamps to Solana epochs + /// The timestamp seeds a slot estimate, and the epoch that seed lands in is + /// then verified against real block times. The verification is what makes the + /// answer correct: the seed drifts by thousands of slots over a day of + /// lookback and no fixed slot duration survives the SIMD-0525 rollout, so a + /// seeded guess alone picks the wrong epoch near a boundary. That epoch + /// chooses the leader schedule contributor rewards are computed against, so a + /// wrong answer corrupts rewards and an error is the better outcome. + /// + /// The forward step exists despite the backward-biased seed because the local + /// clock, not the chain, decides where the seed lands, so a clock running + /// ahead can still overshoot. + /// + /// A second chain-verified epoch search lives in `validator-debt/src/rpc.rs` + /// (`find_solana_epoch_before_timestamp`). It threads a `leaky-bucket` rate + /// limiter and searches one direction only, so the two are not yet worth + /// unifying, but a fix to the boundary handling here probably belongs there + /// too. pub async fn find_epoch_at_timestamp(&mut self, timestamp_us: u64) -> Result { // Get current slot from Solana let current_slot = (|| async { self.solana_read_client.get_slot().await }) @@ -172,19 +464,72 @@ impl EpochFinder { let current_time_us = Utc::now().timestamp_micros() as u64; - // Estimate the slot at the given timestamp - let target_slot = + // Also rejects a future or unreachably old timestamp before any RPC + // calls are spent on it. + let estimated_slot = estimate_slot_from_timestamp(timestamp_us, current_slot, current_time_us)?; - // Get SOLANA epoch schedule and calculate epoch - let schedule = self.get_solana_schedule().await?; - let epoch = calculate_epoch_from_slot(target_slot, schedule); + // Copied rather than borrowed: the borrow is tied to the &mut self that + // the block time lookups below also need. + let schedule = self.get_solana_schedule().await?.clone(); + + let mut candidate_epoch = schedule.get_epoch(estimated_slot); + let target_time = (timestamp_us / 1_000_000) as i64; + + // The seed was only checked against the local clock, which says nothing + // about how far the endpoint has caught up. Since the search accepts an + // unbounded candidate as the answer, a timestamp past the chain tip would + // otherwise resolve to whatever epoch a lagging endpoint sits in. + let (chain_tip_slot, chain_tip_time) = self.try_chain_tip_block(current_slot).await?; + if target_time > chain_tip_time { + bail!( + "Timestamp {timestamp_us} is ahead of the Solana chain tip at slot \ + {chain_tip_slot} (block time {chain_tip_time}), so the epoch containing it is \ + not yet determined" + ); + } - debug!( - "Mapped timestamp {} to Solana epoch {}", - timestamp_us, epoch - ); - Ok(epoch) + // Each step reuses the bound it already resolved, so only the far side of + // the move needs a lookup. Halves the round trips of a multi-step walk. + let mut epoch_start_time = self + .try_epoch_start_block_time(&schedule, candidate_epoch, chain_tip_slot) + .await?; + let mut next_epoch_start_time = self + .try_epoch_start_block_time(&schedule, candidate_epoch + 1, chain_tip_slot) + .await?; + + for _ in 0..MAX_EPOCH_SEARCH_STEPS { + match decide_epoch_search_step(target_time, epoch_start_time, next_epoch_start_time) { + EpochSearchStep::Found => { + debug!( + "Mapped timestamp {} to Solana epoch {}", + timestamp_us, candidate_epoch + ); + return Ok(candidate_epoch); + } + EpochSearchStep::Earlier => { + candidate_epoch = candidate_epoch.checked_sub(1).with_context(|| { + format!("Timestamp {timestamp_us} precedes the first Solana epoch") + })?; + next_epoch_start_time = epoch_start_time; + epoch_start_time = self + .try_epoch_start_block_time(&schedule, candidate_epoch, chain_tip_slot) + .await?; + } + EpochSearchStep::Later => { + candidate_epoch += 1; + epoch_start_time = next_epoch_start_time; + next_epoch_start_time = self + .try_epoch_start_block_time(&schedule, candidate_epoch + 1, chain_tip_slot) + .await?; + } + } + } + + bail!( + "Could not resolve timestamp {timestamp_us} to a Solana epoch within \ + {MAX_EPOCH_SEARCH_STEPS} steps, last candidate was epoch {candidate_epoch}" + ) } /// Fetch leader schedule for a DZ epoch @@ -258,35 +603,24 @@ impl EpochFinder { #[cfg(test)] mod tests { - use super::*; + use solana_client::rpc_request::RpcResponseErrorData; - #[test] - fn test_calculate_epoch_from_slot_normal() { - let schedule = EpochSchedule { - slots_per_epoch: 432000, - leader_schedule_slot_offset: 432000, - warmup: false, - first_normal_epoch: 0, - first_normal_slot: 0, - }; - - assert_eq!(calculate_epoch_from_slot(0, &schedule), 0); - assert_eq!(calculate_epoch_from_slot(432000, &schedule), 1); - assert_eq!(calculate_epoch_from_slot(864000, &schedule), 2); - assert_eq!(calculate_epoch_from_slot(431999, &schedule), 0); - } + use super::*; #[test] fn test_estimate_slot_from_timestamp() { let current_slot = 1000000; let current_time_us = 1_000_000_000_000; // 1 million seconds in microseconds - // Test normal case - 400 seconds ago (1000 slots) - let timestamp_us = current_time_us - 400_000_000; + // Test normal case - 500 seconds ago (500_000_000 us / 500_000 us per + // slot = 1000 slots) + let timestamp_us = current_time_us - 500_000_000; let result = estimate_slot_from_timestamp(timestamp_us, current_slot, current_time_us); assert_eq!(result.unwrap(), 999000); - // Test future timestamp + // Test future timestamp. find_epoch_at_timestamp seeds its search with + // this call, so this guard is what keeps a future timestamp out of the + // search entirely. let future_timestamp = current_time_us + 1000; let result = estimate_slot_from_timestamp(future_timestamp, current_slot, current_time_us); assert!(result.is_err()); @@ -296,4 +630,174 @@ mod tests { let result = estimate_slot_from_timestamp(ancient_timestamp, current_slot, current_time_us); assert!(result.is_err()); } + + // The epoch's first confirmed block time is the inclusive lower bound, so a + // timestamp landing exactly on it belongs to the candidate epoch. + #[test] + fn test_decide_epoch_search_step_at_epoch_start_is_found() { + assert_eq!( + decide_epoch_search_step(1_700_000_000, Some(1_700_000_000), Some(1_700_100_000)), + EpochSearchStep::Found + ); + } + + // One second earlier belongs to the previous epoch. This is the boundary case + // that a seeded estimate alone got wrong. + #[test] + fn test_decide_epoch_search_step_before_epoch_start_steps_earlier() { + assert_eq!( + decide_epoch_search_step(1_699_999_999, Some(1_700_000_000), Some(1_700_100_000)), + EpochSearchStep::Earlier + ); + } + + // The next epoch's first confirmed block time is the exclusive upper bound, so + // a timestamp landing exactly on it belongs to the next epoch, not this one. + #[test] + fn test_decide_epoch_search_step_at_next_epoch_start_steps_later() { + assert_eq!( + decide_epoch_search_step(1_700_100_000, Some(1_700_000_000), Some(1_700_100_000)), + EpochSearchStep::Later + ); + } + + // Without this the search would fail on every recent timestamp. + #[test] + fn test_decide_epoch_search_step_current_epoch_has_no_upper_bound() { + assert_eq!( + decide_epoch_search_step(1_700_100_000, Some(1_700_000_000), None), + EpochSearchStep::Found + ); + } + + // The caller feeds this step into a checked_sub, so a timestamp older than + // the earliest available block errors rather than underflowing or silently + // returning epoch 0. + #[test] + fn test_decide_epoch_search_step_unstarted_epoch_steps_earlier() { + assert_eq!( + decide_epoch_search_step(1_700_000_000, None, None), + EpochSearchStep::Earlier + ); + } + + fn rpc_response_error(code: i64) -> SolanaClientError { + RpcError::RpcResponseError { + code, + message: "test".to_string(), + data: RpcResponseErrorData::Empty, + } + .into() + } + + // Which shape arrives depends on what the endpoint has behind it, so missing + // any one of them fails every lookup on endpoints of that shape. + #[test] + fn test_is_block_unavailable_covers_every_absent_block_shape() { + assert!(is_block_unavailable(&rpc_response_error( + JSON_RPC_SERVER_ERROR_SLOT_SKIPPED + ))); + assert!(is_block_unavailable(&rpc_response_error( + JSON_RPC_SERVER_ERROR_LONG_TERM_STORAGE_SLOT_SKIPPED + ))); + assert!(is_block_unavailable(&rpc_response_error( + JSON_RPC_SERVER_ERROR_BLOCK_NOT_AVAILABLE + ))); + // What RpcClient::get_block_time synthesizes from a JSON null response. + assert!(is_block_unavailable( + &RpcError::ForUser("Block Not Found: slot=123".to_string()).into() + )); + } + + // A pruned ledger has to fail the lookup rather than read as an absent block, + // and is not worth retrying either. + #[test] + fn test_is_block_cleaned_up_is_not_an_absent_block() { + let err = rpc_response_error(JSON_RPC_SERVER_ERROR_BLOCK_CLEANED_UP); + assert!(!is_block_unavailable(&err)); + assert!(is_block_cleaned_up(&err)); + // Settled either way, so neither is worth retrying. + assert!(is_settled_block_error(&err)); + } + + // A transport failure is neither, so it stays retryable. + #[test] + fn test_transport_error_is_retryable() { + let err = RpcError::RpcRequestError("connection reset".to_string()).into(); + assert!(!is_block_unavailable(&err)); + assert!(!is_block_cleaned_up(&err)); + assert!(!is_settled_block_error(&err)); + } + + // The ordinary cases: the epoch's own first slot produced a block, or a short + // run of skipped slots pushed the first block a few slots in. + #[test] + fn test_can_date_epoch_start_accepts_a_short_skip_run() { + let first_slot = 432_000; + let next_epoch_first_slot = 864_000; + + assert!(can_date_epoch_start( + first_slot, + first_slot, + next_epoch_first_slot + )); + assert!(can_date_epoch_start( + first_slot, + first_slot + 4, + next_epoch_first_slot + )); + } + + // The budget is exclusive, so the last accepted slot is one below it. + #[test] + fn test_can_date_epoch_start_budget_edge() { + let first_slot = 432_000; + let next_epoch_first_slot = 864_000; + + assert!(can_date_epoch_start( + first_slot, + first_slot + MAX_BOUNDARY_SKIP_SLOTS - 1, + next_epoch_first_slot + )); + assert!(!can_date_epoch_start( + first_slot, + first_slot + MAX_BOUNDARY_SKIP_SLOTS, + next_epoch_first_slot + )); + } + + // A block this far in means getBlocksWithLimit stepped over a history gap. + // Dating the epoch by it would resolve timestamps inside the gap to the + // previous epoch and its leader schedule. + #[test] + fn test_can_date_epoch_start_rejects_a_history_gap() { + let first_slot = 432_000; + let next_epoch_first_slot = 864_000; + + assert!(!can_date_epoch_start( + first_slot, + first_slot + 100_000, + next_epoch_first_slot + )); + } + + // On a cluster whose epochs are shorter than the budget, such as a local test + // validator, the epoch's own end is what binds. Without the `min` the budget + // would accept a block belonging to a later epoch. + #[test] + fn test_can_date_epoch_start_short_epoch_binds_before_the_budget() { + let first_slot = 64; + let next_epoch_first_slot = 96; + + assert!(can_date_epoch_start( + first_slot, + next_epoch_first_slot - 1, + next_epoch_first_slot + )); + assert!(!can_date_epoch_start( + first_slot, + next_epoch_first_slot, + next_epoch_first_slot + )); + } } diff --git a/crates/solana-cli/CHANGELOG.md b/crates/solana-cli/CHANGELOG.md index 4a32ae47..46185e56 100644 --- a/crates/solana-cli/CHANGELOG.md +++ b/crates/solana-cli/CHANGELOG.md @@ -7,6 +7,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +- `shreds`: collapse `pay`'s `SLOT_DURATION_SECS` and `prepare-offchain-message`'s `SLOT_DURATION_MS` into one `NOMINAL_SLOT_DURATION` at 350ms, matching mainnet-beta from epoch 1020 (2026-08-21). Deliberately cluster-independent, because a `~` prefixed estimate and a deadline slot the CLI and operator must both compute want reproducibility over accuracy. `--valid-for 1h` now resolves to 10,285 slots rather than 9,000, and the epoch-remaining estimates shrink by an eighth. Testnet runs at 200ms, so its estimates stay wrong in the other direction, and SIMD-0525 will need one more bump here (malbeclabs/infra#2317) - `shreds validator-client-rewards`: read the `ValidatorClientRewards` account through the SDK's zero-copy mirror instead of the hand-written byte-offset parser. Every subcommand now requires at least 184 bytes of account data rather than 116 - `shreds pay`: print the fastshreds.com transition notice, and require a y/N confirmation when the operator can actually answer it. The notice fires before the wallet is built and before the first RPC call, so declining needs neither a loadable keypair nor a reachable cluster and signs nothing. The prompt requires *both* stdin and stdout to be a terminal: the prompt is written to stdout, so under a redirect (`> pay.log`, `| tee`) the question would land in the file while `read_line` blocked on a terminal showing nothing. Dry runs print the notice without prompting, matching the epoch-remaining prompt, because a simulation signs and sends nothing. A non-terminal stdin also prints without prompting — `read_line` there returns EOF, which the helper reads as "no", and keypair-from-stdin requires a non-terminal stdin, so a prompt would consume the keypair bytes. `--accept-deprecation-notice` skips the prompt for interactive batch and multi-seat workflows; the notice still prints. Declining exits non-zero, matching the epoch-remaining prompt, so a `shreds pay ... && ...` chain does not treat a decline as a completed payment (malbeclabs/infra#2164) - `shreds pay`: check the `--amount` floor against the price the program actually charges. A new instant seat allocation is priced from the metro/device ring entries at the execution controller's `last_settled_epoch` (the seat covers the remainder of the epoch currently being served), not from the newest entries — those two agree only while prices are static, so the epoch a metro repriced the preflight passed an underfunded amount through to an opaque `invalid account data for instruction`. A pure escrow top-up for an already-active seat submits no `RequestInstantSeatAllocation` and keeps its floor at the newest entry's price. When either ring has no entry for `last_settled_epoch` the command refuses to submit instead of falling back to the newest entry, mirroring the program's own error. The rejection message now names both prices and both epochs (#405) diff --git a/crates/solana-cli/src/command/shreds/mod.rs b/crates/solana-cli/src/command/shreds/mod.rs index 7c2d8c1e..d2199c0d 100644 --- a/crates/solana-cli/src/command/shreds/mod.rs +++ b/crates/solana-cli/src/command/shreds/mod.rs @@ -6,7 +6,7 @@ pub mod publisher_rewards; pub mod validator_client_rewards; pub mod withdraw; -use std::io::Write; +use std::{io::Write, time::Duration}; use anyhow::{Result, bail}; use clap::{Args, Subcommand}; @@ -23,6 +23,15 @@ use solana_client::{ }; use solana_sdk::pubkey::Pubkey; +// Solana's nominal slot duration. Mainnet-beta moves 400ms to 350ms at the +// start of epoch 1020 (2026-08-21) and SIMD-0525 continues stepping it down to +// 200ms, so this needs one more update when that rollout completes. Testnet is +// already at 200ms. This value deliberately does not vary by cluster, because +// both users want a deterministic, reproducible number more than an accurate +// one: one prints a "~" prefixed estimate and the other computes a deadline +// slot that the CLI and the operator must agree on. +pub(in crate::command::shreds) const NOMINAL_SLOT_DURATION: Duration = Duration::from_millis(350); + #[derive(Debug, Args)] pub struct ShredsCommand { /// Override the DZ Ledger RPC URL. When omitted, the URL is derived from diff --git a/crates/solana-cli/src/command/shreds/pay.rs b/crates/solana-cli/src/command/shreds/pay.rs index 6335fbc6..64290acb 100644 --- a/crates/solana-cli/src/command/shreds/pay.rs +++ b/crates/solana-cli/src/command/shreds/pay.rs @@ -33,7 +33,7 @@ use solana_compute_budget_interface::ComputeBudgetInstruction; use solana_sdk::pubkey::Pubkey; use spl_associated_token_account_interface::address::get_associated_token_address; -use super::{make_dz_connection, serviceability_program_id}; +use super::{NOMINAL_SLOT_DURATION, make_dz_connection, serviceability_program_id}; /// Warn if less than 10% of the current Solana epoch remains. const EPOCH_REMAINING_WARNING_THRESHOLD: f64 = 0.10; @@ -46,10 +46,6 @@ const DEPRECATION_NOTICE: &str = "The doublezero-solana client will soon be depr are available. You can then proceed to withdraw funds in the CLI and close out this account. \ This command will fund your seat now"; -/// Solana's target slot duration in seconds. Actual varies with network -/// conditions; used only for approximate display estimates (prefixed with `~`). -const SLOT_DURATION_SECS: f64 = 0.4; - /// Inputs for the epoch-remaining warning check, separated from I/O for testability. struct EpochWarningInput { accept_partial_epoch: bool, @@ -98,8 +94,12 @@ fn epoch_warning_prompt(input: &EpochWarningInput) -> Option { return None; } - let remaining_secs = (input.slots_in_epoch - input.slot_index) as f64 * SLOT_DURATION_SECS; - let total_secs = input.slots_in_epoch as f64 * SLOT_DURATION_SECS; + // The nominal slot duration, not the observed one. Actual slot times vary + // with network conditions, so both numbers are approximate and are printed + // with a "~" prefix. + let remaining_secs = + (input.slots_in_epoch - input.slot_index) as f64 * NOMINAL_SLOT_DURATION.as_secs_f64(); + let total_secs = input.slots_in_epoch as f64 * NOMINAL_SLOT_DURATION.as_secs_f64(); Some(format!( "Only {:.1}% of the current Solana epoch remains ({} of {}).\n \ @@ -861,10 +861,10 @@ mod tests { fn warning_message_contains_time_estimates() { let input = make_input(0.05); let prompt = epoch_warning_prompt(&input).expect("should warn"); - // 5% of 432000 slots = 21600 slots * 0.4s = 8640s = 2.4 hours - assert!(prompt.contains("~2.4 hours")); - // Total epoch: 432000 * 0.4s = 172800s = 48.0 hours - assert!(prompt.contains("~48.0 hours")); + // 5% of 432,000 slots = 21,600 slots * 0.35 s = 7,560 s = 2.1 hours + assert!(prompt.contains("~2.1 hours")); + // Total epoch: 432,000 * 0.35 s = 151,200 s = 42.0 hours + assert!(prompt.contains("~42.0 hours")); } #[test] @@ -898,8 +898,11 @@ mod tests { // rather than panicking. #[test] fn format_duration_exactly_zero_slots_remaining() { - // 0 slots * 0.4 s = 0.0 s - assert_eq!(format_duration(0.0 * SLOT_DURATION_SECS), "~0 minutes"); + // 0 slots * 0.35 s = 0.0 s + assert_eq!( + format_duration(0.0 * NOMINAL_SLOT_DURATION.as_secs_f64()), + "~0 minutes" + ); } // --- Additional coverage: epoch_warning_prompt edge cases --- @@ -982,14 +985,14 @@ mod tests { dry_run: false, seat_active_this_epoch: false, prorated_service_enabled: false, - slot_index: slots_in_epoch - 1, // 1 slot = 0.4 s remaining + slot_index: slots_in_epoch - 1, // 1 slot = 0.35 s remaining slots_in_epoch, }; let prompt = epoch_warning_prompt(&input).expect("should warn"); - // 1 slot * 0.4 s = 0.4 s → rounds to "~0 minutes" + // 1 slot * 0.35 s = 0.35 s → rounds to "~0 minutes" assert!(prompt.contains("~0 minutes")); // Total epoch is hours, not minutes - assert!(prompt.contains("~48.0 hours")); + assert!(prompt.contains("~42.0 hours")); } // Spec row: "passed | any | any | any → Proceed silently, no warning". @@ -1026,20 +1029,20 @@ mod tests { // Verify the warning content uses minutes (not hours) for the remaining-time // fields when less than an hour remains, even though the total epoch time is // still rendered in hours. - // 1% of a 432,000-slot epoch = 4,320 slots * 0.4 s = 1,728 s ≈ 28.8 minutes + // 1% of a 432,000-slot epoch = 4,320 slots * 0.35 s = 1,512 s = 25.2 minutes #[test] fn warning_near_epoch_end_shows_minutes_for_remaining_time() { - let input = make_input(0.01); // 1% remaining → ~29 minutes + let input = make_input(0.01); // 1% remaining → ~25 minutes let prompt = epoch_warning_prompt(&input).expect("should warn"); // The remaining-time estimate must appear as minutes. assert!( - prompt.contains("~29 minutes"), - "expected '~29 minutes' for remaining time, got: {prompt}" + prompt.contains("~25 minutes"), + "expected '~25 minutes' for remaining time, got: {prompt}" ); // The total epoch estimate is still rendered in hours. assert!( - prompt.contains("~48.0 hours"), - "expected '~48.0 hours' for total epoch, got: {prompt}" + prompt.contains("~42.0 hours"), + "expected '~42.0 hours' for total epoch, got: {prompt}" ); } @@ -1052,11 +1055,11 @@ mod tests { let prompt = epoch_warning_prompt(&input).unwrap(); assert!(prompt.contains("5.0%"), "missing percentage"); assert!( - prompt.contains("~2.4 hours"), + prompt.contains("~2.1 hours"), "missing remaining time estimate" ); assert!( - prompt.contains("~48.0 hours"), + prompt.contains("~42.0 hours"), "missing total epoch time estimate" ); assert!( diff --git a/crates/solana-cli/src/command/shreds/publisher_rewards/prepare_offchain_message.rs b/crates/solana-cli/src/command/shreds/publisher_rewards/prepare_offchain_message.rs index 23c1bca9..d2849e94 100644 --- a/crates/solana-cli/src/command/shreds/publisher_rewards/prepare_offchain_message.rs +++ b/crates/solana-cli/src/command/shreds/publisher_rewards/prepare_offchain_message.rs @@ -12,7 +12,7 @@ use doublezero_solana_sdk::{ }, }; -use super::rewards_mint_arg::RewardsMintArg; +use super::{super::NOMINAL_SLOT_DURATION, rewards_mint_arg::RewardsMintArg}; /* doublezero-solana shreds publisher-rewards prepare-offchain-message \ @@ -21,7 +21,6 @@ use super::rewards_mint_arg::RewardsMintArg; [--deadline-slot | --valid-for ] [--json] */ -const SLOT_DURATION_MS: u64 = 400; const DEFAULT_VALID_FOR: Duration = Duration::from_secs(60 * 60); #[derive(Debug, Args)] @@ -42,7 +41,9 @@ pub struct PrepareOffchainMessageCommand { pub deadline_slot: Option, /// Duration the authorization remains valid (e.g. `1h`, `30m`, `7200s`). - /// Default: `1h`. Parsed via `humantime`. + /// Default: `1h`. Parsed via `humantime`. The program enforces an absolute + /// slot, so this is converted at a nominal 350ms per slot and the real + /// elapsed time varies with the cluster's slot rate. #[arg(long, value_parser = parse_valid_for)] pub valid_for: Option, @@ -137,9 +138,9 @@ impl PrepareOffchainMessageCommand { /// Pure helper: resolve the absolute deadline slot from CLI inputs. /// -/// `--deadline-slot` always wins. If absent, `--valid-for` is divided by the -/// nominal slot duration (400ms) and added to `current_slot`. If both are -/// `None`, defaults to 1h. Both supplied is an error. +/// `--deadline-slot` always wins. If absent, `--valid-for` is divided by +/// `NOMINAL_SLOT_DURATION` and added to `current_slot`. If both are `None`, +/// defaults to 1h. Both supplied is an error. pub(crate) fn resolve_deadline_slot( current_slot: u64, deadline_slot: Option, @@ -155,7 +156,7 @@ pub(crate) fn resolve_deadline_slot( let duration = valid_for.unwrap_or(DEFAULT_VALID_FOR); let slots = duration .as_millis() - .checked_div(SLOT_DURATION_MS as u128) + .checked_div(NOMINAL_SLOT_DURATION.as_millis()) .context("invalid slot duration")?; let slots: u64 = slots .try_into() @@ -187,14 +188,18 @@ mod tests { #[test] fn valid_for_default_one_hour() { let resolved = resolve_deadline_slot(100, None, None).unwrap(); - assert_eq!(resolved, 100 + 9_000); + // 3,600,000 ms / 350 ms = 10,285 slots (integer division truncates the + // remainder of 250 ms). + assert_eq!(resolved, 100 + 10_285); } #[test] fn valid_for_explicit_30m() { let resolved = resolve_deadline_slot(100, None, Some(Duration::from_secs(30 * 60))).unwrap(); - assert_eq!(resolved, 100 + 4_500); + // 1,800,000 ms / 350 ms = 5,142 slots (integer division truncates the + // remainder of 300 ms). + assert_eq!(resolved, 100 + 5_142); } #[test] diff --git a/crates/validator-debt/CHANGELOG.md b/crates/validator-debt/CHANGELOG.md index 2c83d2cb..060526f1 100644 --- a/crates/validator-debt/CHANGELOG.md +++ b/crates/validator-debt/CHANGELOG.md @@ -7,6 +7,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +- refactor(validator-debt): delete the two dead timestamp to Solana epoch code paths that hardcoded a 0.4s slot duration, both `SLOT_TIME_DURATION_SECONDS` constants, and the two `ValidatorRewards` trait methods only they reached. Nothing called them: `rpc.rs` already maps DZ epochs to Solana epochs from real block times for the production path, and the worker calls `get_total_rewards` with an explicit epoch. Also removes an unsigned subtraction that panicked in debug and wrapped silently in release. Note that `estimate_block_time_for_skipped_slot` in `rpc.rs` still assumes 0.4s per slot implicitly (malbeclabs/infra#2317) - migrate to Solana 3.0: workspace `solana-*` crates and `solana-sdk` move to the 3.0 line, `solana-program-test` to 3.0.12, and the doublezero SDK git-deps repin from `client/v0.27.1` to the malbeclabs/doublezero#3830 merge revision (malbeclabs/infra#1853) - release artifact now builds as a static `x86_64-unknown-linux-musl` binary (malbeclabs/infra#1853) - TLS for HTTP clients moves from openssl to rustls; trust roots are the bundled webpki Mozilla set plus the host OS certificate store, so OS-installed private CAs remain trusted (malbeclabs/infra#1853) diff --git a/crates/validator-debt/src/ledger.rs b/crates/validator-debt/src/ledger.rs index 00bcdf36..3071edb4 100644 --- a/crates/validator-debt/src/ledger.rs +++ b/crates/validator-debt/src/ledger.rs @@ -5,7 +5,6 @@ use doublezero_solana_client_tools::rpc::DoubleZeroLedgerConnection; use solana_client::{nonblocking::rpc_client::RpcClient, rpc_config::RpcSendTransactionConfig}; use solana_commitment_config::CommitmentConfig; use solana_sdk::{ - clock::Epoch, hash::Hash, pubkey::Pubkey, signer::{Signer, keypair::Keypair}, @@ -13,36 +12,9 @@ use solana_sdk::{ use crate::validator_debt::ComputedSolanaValidatorDebts; -const SLOT_TIME_DURATION_SECONDS: f64 = 0.4; - pub const DOUBLEZERO_LEDGER_MAINNET_BETA_GENESIS_HASH: Pubkey = solana_sdk::pubkey!("5wVUvkFcFGYiKRUZ8Jp8Wc5swjhDEqT7hTdyssxDpC7P"); -pub async fn get_solana_epoch_from_dz_epoch( - solana_client: &RpcClient, - ledger_client: &RpcClient, - dz_epoch: Epoch, -) -> Result<(u64, u64)> { - let epoch_info = ledger_client.get_epoch_info().await?; - - let first_slot_in_current_epoch = epoch_info.absolute_slot - epoch_info.slot_index; - - let epoch_diff = epoch_info.epoch - dz_epoch; - - let first_slot = first_slot_in_current_epoch - (epoch_info.slots_in_epoch * epoch_diff) - 1; - let last_slot = first_slot + (epoch_info.slots_in_epoch - 1); - - let solana_epoch_from_first_dz_epoch_slot = - get_solana_epoch_from_dz_slot(solana_client, ledger_client, first_slot).await?; - let solana_epoch_from_last_dz_epoch_slot = - get_solana_epoch_from_dz_slot(solana_client, ledger_client, last_slot).await?; - - Ok(( - solana_epoch_from_first_dz_epoch_slot + 1, - solana_epoch_from_last_dz_epoch_slot, - )) -} - pub async fn create_record_on_ledger( rpc_client: &RpcClient, recent_blockhash: Hash, @@ -119,34 +91,6 @@ pub async fn try_fetch_debt_record( Ok((debt_record.header, debt_record.data)) } -async fn get_solana_epoch_from_dz_slot( - solana_client: &RpcClient, - ledger_client: &RpcClient, - slot: u64, -) -> Result { - let block = ledger_client.get_block(slot).await?; - - let dz_block_time = block.block_time.unwrap(); - let dz_block_time: u64 = dz_block_time as u64; - - let solana_epoch_info = solana_client.get_epoch_info().await?; - - let first_slot_in_current_solana_epoch = - solana_epoch_info.absolute_slot - solana_epoch_info.slot_index; - - let block_time = solana_client - .get_block_time(first_slot_in_current_solana_epoch) - .await?; - let block_time: u64 = block_time as u64; - - let num_slots: u64 = ((block_time - dz_block_time) as f64 / SLOT_TIME_DURATION_SECONDS) as u64; - - Ok( - (solana_epoch_info.epoch * solana_epoch_info.slots_in_epoch - num_slots) - / solana_epoch_info.slots_in_epoch, - ) -} - pub async fn ensure_same_network_environment( dz_ledger_rpc: &RpcClient, is_mainnet: bool, diff --git a/crates/validator-debt/src/rewards.rs b/crates/validator-debt/src/rewards.rs index 68fdd900..c996f522 100644 --- a/crates/validator-debt/src/rewards.rs +++ b/crates/validator-debt/src/rewards.rs @@ -5,17 +5,12 @@ //! - JITO rewards per epoch //! //! The rewards from all sources for an epoch are summed and associated with a validator_id -use std::collections::HashMap; - -use anyhow::{Result, anyhow}; +use anyhow::Result; use borsh::{BorshDeserialize, BorshSerialize}; use serde::Deserialize; -use solana_sdk::clock::DEFAULT_SLOTS_PER_EPOCH; use crate::{block, inflation, jito, solana_debt_calculator::ValidatorRewards}; -const SLOT_TIME_DURATION_SECONDS: f64 = 0.4; - #[derive(Deserialize, Debug, BorshDeserialize, BorshSerialize)] pub struct EpochRewards { pub epoch: u64, @@ -33,26 +28,6 @@ pub struct Reward { pub block_base: u64, } -pub async fn get_rewards_between_timestamps( - solana_debt_calculator: &impl ValidatorRewards, - start_timestamp: u64, - end_timestamp: u64, - validator_ids: &[String], -) -> Result>> { - let mut rewards: HashMap> = HashMap::new(); - let current_slot = solana_debt_calculator.get_slot().await?; - let block_time = solana_debt_calculator.get_block_time(current_slot).await?; - let block_time: u64 = block_time as u64; - - let start_epoch = epoch_from_timestamp(block_time, current_slot, start_timestamp)?; - let end_epoch = epoch_from_timestamp(block_time, current_slot, end_timestamp)?; - for epoch in start_epoch..=end_epoch { - let reward = get_total_rewards(solana_debt_calculator, validator_ids, epoch).await?; - rewards.insert(epoch, reward.rewards); - } - Ok(rewards) -} - // this function will return a hashmap of total rewards keyed by validator pubkey pub async fn get_total_rewards( solana_debt_calculator: &impl ValidatorRewards, @@ -104,24 +79,10 @@ pub async fn get_total_rewards( Ok(rewards) } -// get the number of slots by subtracting the timestamp from the block time and dividing it by the time per slot -// get the desired slot by subtracting the num_slots from the current_slot -// then get the epoch by dividing the desired_slot by the DEFAULT_SLOTS_PER_EPOCH -// NOTE: This can change if solana changes -fn epoch_from_timestamp(block_time: u64, current_slot: u64, timestamp: u64) -> Result { - if timestamp > block_time { - return Err(anyhow!( - "timestamp cannot be greater than block_time: {timestamp}, {block_time}" - )); - } - let num_slots: u64 = ((block_time - timestamp) as f64 / SLOT_TIME_DURATION_SECONDS) as u64; - let desired_slot = current_slot - num_slots; - // epoch - Ok(desired_slot / DEFAULT_SLOTS_PER_EPOCH) -} - #[cfg(test)] mod tests { + use std::collections::HashMap; + use solana_client::rpc_response::{ RpcInflationReward, RpcVoteAccountInfo, RpcVoteAccountStatus, }; @@ -135,167 +96,6 @@ mod tests { solana_debt_calculator::MockValidatorRewards, }; - #[tokio::test] - async fn test_get_rewards_between_timestamps() { - // Set up test variables and mock data. - let validator_id = "6WgdYhhGE53WrZ7ywJA15hBVkw7CRbQ8yDBBTwmBtAHN"; - let validator_ids: &[String] = &[String::from(validator_id)]; - let epoch = 824; - let block_reward: u64 = 40000; - let inflation_reward = 2500; - let jito_reward = 10000; - - let start_timestamp = 1752727180; - let end_timestamp = 1752727280; - - let mut mock_solana_debt_calculator = MockValidatorRewards::new(); - - // Set up mock expectations for the ValidatorRewards trait. - // These mocks simulate the behavior of external dependencies. - mock_solana_debt_calculator - .expect_get_slot() - .times(1) - .returning(move || Ok(356170122)); - - mock_solana_debt_calculator - .expect_get_block_time() - .times(1) - .returning(move |_| Ok(1752728180)); - - let signatures = vec![ - "One".to_string(), - "Two".to_string(), - "Three".to_string(), - "Four".to_string(), - "Five".to_string(), - "Six".to_string(), - "Seven".to_string(), - "Eight".to_string(), - "Nine".to_string(), - "Ten".to_string(), - "Eleven".to_string(), - "Twelve".to_string(), - ]; - let priority_fees = signatures.len() as u64 * 2500; - let mock_block = UiConfirmedBlock { - num_reward_partitions: Some(1), - signatures: Some(signatures), - rewards: Some(vec![solana_transaction_status_client_types::Reward { - pubkey: validator_id.to_string(), - lamports: block_reward as i64, - post_balance: block_reward, - reward_type: Some(Fee), - commission: None, - }]), - previous_blockhash: "".to_string(), - blockhash: "".to_string(), - parent_slot: 0, - transactions: None, - block_time: None, - block_height: None, - }; - - let slot_index: usize = 10; - - let mock_epoch_info = EpochInfo { - epoch, - slot_index: 100000, - absolute_slot: 10000000, - block_height: 103030003, - slots_in_epoch: 5000000, - transaction_count: Some(1000), - }; - - let first_slot = 9900010; - mock_solana_debt_calculator - .expect_get::() - .withf(move |url| url.contains(&format!("epoch={epoch}"))) - .times(1) - .returning(move |_| { - Ok(JitoRewards { - total_count: 1000, - rewards: vec![JitoReward { - vote_account: validator_id.to_string(), - mev_revenue: jito_reward, - }], - }) - }); - - mock_solana_debt_calculator - .expect_get_epoch_info() - .times(1) - .returning(move || Ok(mock_epoch_info.clone())); - - mock_solana_debt_calculator - .expect_get_block_with_config() - .withf(move |s| *s == first_slot) - .times(1) - .returning(move |_| Ok(mock_block.clone())); - - let mock_rpc_vote_account_status = RpcVoteAccountStatus { - current: vec![RpcVoteAccountInfo { - vote_pubkey: "6WgdYhhGE53WrZ7ywJA15hBVkw7CRbQ8yDBBTwmBtBBN".to_string(), - node_pubkey: validator_id.to_string(), - activated_stake: 4_200_000_000_000, - epoch_vote_account: true, - epoch_credits: vec![(812, 256, 128), (811, 128, 64)], - commission: 10, - last_vote: 123456789, - root_slot: 123456700, - }], - delinquent: vec![], - }; - - mock_solana_debt_calculator - .expect_get_vote_accounts_with_config() - .withf(move || true) - .times(1) - .returning(move || Ok(mock_rpc_vote_account_status.clone())); - - let mock_rpc_inflation_reward = vec![Some(RpcInflationReward { - epoch, - effective_slot: 123456789, - amount: inflation_reward, - post_balance: 1_500_002_500, - commission: Some(1), - })]; - - mock_solana_debt_calculator - .expect_get_inflation_reward() - .times(1) - .returning(move |_, _| Ok(mock_rpc_inflation_reward.clone())); - - let mut leader_schedule = HashMap::new(); - leader_schedule.insert(validator_id.to_string(), vec![slot_index]); - - mock_solana_debt_calculator - .expect_get_leader_schedule() - .times(1) - .returning(move |_| Ok(leader_schedule.clone())); - - // Call the function under test with the prepared data and mocks. - let rewards = get_rewards_between_timestamps( - &mock_solana_debt_calculator, - start_timestamp, - end_timestamp, - validator_ids, - ) - .await - .unwrap(); - - let epoch_rewards = rewards.get(&epoch).unwrap(); - let reward = epoch_rewards - .iter() - .find(|reward| reward.validator_id == validator_id) - .unwrap(); - assert_eq!(reward.block_base + reward.block_priority, block_reward); - assert_eq!(reward.epoch, epoch); - assert_eq!(reward.inflation, inflation_reward); - assert_eq!(reward.jito, jito_reward); - assert_eq!(reward.total, block_reward + inflation_reward + jito_reward); - assert_eq!(reward.block_priority, block_reward - priority_fees); - } - #[tokio::test] async fn test_get_total_rewards() { // Set up test variables and mock data. diff --git a/crates/validator-debt/src/rpc.rs b/crates/validator-debt/src/rpc.rs index df7167b7..f8ccd6c7 100644 --- a/crates/validator-debt/src/rpc.rs +++ b/crates/validator-debt/src/rpc.rs @@ -193,6 +193,11 @@ impl JoinedSolanaEpochs { } } + // A second chain-verified epoch search lives in the contributor-rewards crate + // (`ingestor::epoch::EpochFinder::find_epoch_at_timestamp`). That one steps + // from an estimated seed slot and threads no rate limiter, so the two are not + // yet worth unifying. A fix to the skipped-slot or boundary handling here + // probably belongs there too. async fn find_solana_epoch_before_timestamp( solana_client: &RpcClient, rate_limiter: &RateLimiter, diff --git a/crates/validator-debt/src/solana_debt_calculator.rs b/crates/validator-debt/src/solana_debt_calculator.rs index ad624902..a92b4ee9 100644 --- a/crates/validator-debt/src/solana_debt_calculator.rs +++ b/crates/validator-debt/src/solana_debt_calculator.rs @@ -51,8 +51,6 @@ pub trait ValidatorRewards { vote_keys: Vec, epoch: u64, ) -> Result>, ClientError>; - async fn get_slot(&self) -> Result; - async fn get_block_time(&self, slot: u64) -> Result; } pub struct SolanaDebtCalculator { @@ -131,11 +129,4 @@ impl ValidatorRewards for SolanaDebtCalculator { .get_inflation_reward(&vote_keys, Some(epoch)) .await } - async fn get_slot(&self) -> Result { - self.solana_rpc_client.get_slot().await - } - - async fn get_block_time(&self, slot: u64) -> Result { - self.solana_rpc_client.get_block_time(slot).await - } }