Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
343 changes: 343 additions & 0 deletions crates/blockchain/src/epoch_precompute.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,343 @@
//! Beacon epoch-transition precompute.
//!
//! The first block of an epoch runs `process_epoch` inline on the import path,
//! then pays the post-epoch rehash for its state-root check. Both can be done
//! before the block arrives, since the parent is known by the end of the last
//! slot of the previous epoch: clone the head state, advance it to the epoch's
//! first slot on a blocking worker, hash it, and park it in the store's state
//! cache under [`CacheKey::CheckpointState`]. `fork_choice::on_block` resumes
//! from that entry when it finds one, and falls back to the inline path when it
//! does not, so the precompute is a pure speed trade.
//!
//! Two triggers feed the same worker:
//!
//! | Trigger | When | Covers |
//! |---|---|---|
//! | [`Trigger::Head`] | an import leaves a last-slot block as head | a late last-slot block |
//! | [`Trigger::Timer`] | three quarters into a last slot | a skipped last slot |
//!
//! At most one worker runs at a time, and a key already in the cache is never
//! recomputed. A block that arrives while the worker is still running imports
//! inline; the late result is stored and simply goes unused by that block.

use std::sync::Arc;
use std::time::Duration;

use ethlambda_state_transition::beacon::fork_choice;
use ethlambda_state_transition::beacon::helpers::misc::compute_start_slot_at_epoch;
use ethlambda_state_transition::metrics as stf_metrics;
use ethlambda_storage::{CacheKey, Chain};
use ethlambda_types::ShortRoot;
use ethlambda_types::beacon::containers::BeaconState;
use ethlambda_types::beacon::preset;
use ethlambda_types::primitives::H256;
use spawned_concurrency::message::Message;
use spawned_concurrency::tasks::{Context, Handler, send_after};
use tracing::{debug, info, warn};

use crate::BlockChainServer;

/// What caused a precompute, and the label it is counted under.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Trigger {
/// An import made a last-slot block the head.
Head,
/// The three-quarter point of a last slot, carrying that slot.
Timer { slot: u64 },
}

impl Trigger {
fn label(self) -> &'static str {
match self {
Trigger::Head => "head",
Trigger::Timer { .. } => "timer",
}
}
}

/// The state a precompute produces: the head's post-state advanced to the first
/// slot of `epoch`. The same identity as `CacheKey::CheckpointState`.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct PrecomputeKey {
pub(crate) epoch: u64,
pub(crate) root: H256,
}

/// Everything [`decide`] reads, gathered so the decision is a pure function.
pub(crate) struct DecisionInputs {
pub(crate) trigger: Trigger,
pub(crate) syncing: bool,
pub(crate) head_slot: u64,
pub(crate) head_root: H256,
/// The wall-clock slot now.
pub(crate) current_slot: u64,
/// The worker running right now, if any.
pub(crate) in_flight: Option<PrecomputeKey>,
/// Whether the store already holds a state for the candidate key.
pub(crate) already_cached: bool,
}

fn is_last_slot_of_epoch(slot: u64) -> bool {
(slot + 1).is_multiple_of(preset::SLOTS_PER_EPOCH)
}

/// Whether to start a precompute, and for which key.
///
/// `None` when the node is syncing, when the trigger does not sit on the last
/// slot of an epoch, when a worker is already running (at most one), or when
/// the store already holds the state. A head trigger also requires the head to
/// be fresh (at most one slot behind the wall clock): the sync tracker moves
/// once per slot, so a follower importing history would otherwise precompute
/// every epoch it replays.
pub(crate) fn decide(inputs: &DecisionInputs) -> Option<PrecomputeKey> {
if inputs.syncing || inputs.in_flight.is_some() {
return None;
}
let epoch = match inputs.trigger {
Trigger::Head => {
if !is_last_slot_of_epoch(inputs.head_slot)
|| inputs.current_slot > inputs.head_slot.saturating_add(1)
{
return None;
}
inputs.head_slot / preset::SLOTS_PER_EPOCH + 1
}
Trigger::Timer { slot } => {
// A head past the slot would be a stale message racing a newer
// import; one already at the boundary has nothing left to advance.
if !is_last_slot_of_epoch(slot) || inputs.head_slot > slot {
return None;
}
slot / preset::SLOTS_PER_EPOCH + 1
}
};
if inputs.head_slot >= compute_start_slot_at_epoch(epoch) {
return None;
}
let key = PrecomputeKey {
epoch,
root: inputs.head_root,
};
(!inputs.already_cached).then_some(key)
}

/// Delay from the start of a slot to its three-quarter point.
pub(crate) fn three_quarter_slot(slot_duration_ms: u64) -> Duration {
Duration::from_millis(slot_duration_ms / 4 * 3)
}

/// Self-message armed at a last slot's start, delivered at its three-quarter point.
pub(crate) struct EpochPrecomputeCheck {
pub(crate) slot: u64,
}
impl Message for EpochPrecomputeCheck {
type Result = ();
}

/// A worker's result, sent back to the actor.
pub(crate) struct EpochPrecomputed {
pub(crate) key: PrecomputeKey,
pub(crate) result: Result<Arc<BeaconState>, String>,
}
impl Message for EpochPrecomputed {
type Result = ();
}

impl BlockChainServer {
/// Arm the three-quarter check if `slot` is the last of its epoch.
///
/// Called from the once-per-slot tick, which fires at the slot boundary, so
/// the delay is the three-quarter point less whatever of the slot the tick
/// has already consumed (a tick delayed by a long import).
pub(crate) fn arm_epoch_precompute_check(&self, slot: u64, ctx: &Context<Self>) {
if self.store.chain() != Chain::Beacon || !is_last_slot_of_epoch(slot) {
return;
}
let slot_duration_ms = self.store.config().time_grid().milliseconds_per_slot;
let into_slot = Duration::from_millis(self.ms_into_slot(slot).max(0) as u64);
let delay = three_quarter_slot(slot_duration_ms).saturating_sub(into_slot);
send_after(delay, ctx.clone(), EpochPrecomputeCheck { slot });
}

/// Start a precompute for the current head if [`decide`] allows one.
pub(crate) fn maybe_start_epoch_precompute(&mut self, trigger: Trigger, ctx: &Context<Self>) {
if self.store.chain() != Chain::Beacon {
return;
}
let Some((head_slot, head_root)) = self.store.beacon_head() else {
return;
};
let candidate_epoch = match trigger {
Trigger::Head => head_slot / preset::SLOTS_PER_EPOCH + 1,
Trigger::Timer { slot } => slot / preset::SLOTS_PER_EPOCH + 1,
};
let already_cached = self
.store
.cached_state(CacheKey::CheckpointState {
epoch: candidate_epoch,
root: head_root,
})
.is_some();
let inputs = DecisionInputs {
trigger,
syncing: self.sync_status.is_syncing(),
head_slot,
head_root,
current_slot: self.wall_clock_slot(),
in_flight: self.epoch_precompute_in_flight,
already_cached,
};
let Some(key) = decide(&inputs) else {
return;
};
// The head's post-state is resident (it was just imported); a read that
// cannot find it means the head moved to a block this node has no state
// for, which a later trigger can retry.
let Ok(Some(head_state)) = self.store.get_state(&head_root) else {
return;
};

self.epoch_precompute_in_flight = Some(key);
stf_metrics::inc_epoch_precompute_started(trigger.label());
debug!(
epoch = key.epoch,
head_slot,
head_root = %ShortRoot(&head_root.0),
trigger = trigger.label(),
"Starting epoch precompute"
);

let config = self.store.config();
let actor = ctx.actor_ref();
tokio::task::spawn_blocking(move || {
let result = {
let _timing = stf_metrics::time_epoch_precompute();
fork_choice::advance_to_epoch_start(&head_state, key.epoch, &config)
};
let result = result.map(Arc::new).map_err(|err| err.to_string());
let _ = actor.send(EpochPrecomputed { key, result });
});
}
}

impl Handler<EpochPrecomputeCheck> for BlockChainServer {
async fn handle(&mut self, msg: EpochPrecomputeCheck, ctx: &Context<Self>) {
self.maybe_start_epoch_precompute(Trigger::Timer { slot: msg.slot }, ctx);
}
}

impl Handler<EpochPrecomputed> for BlockChainServer {
async fn handle(&mut self, msg: EpochPrecomputed, _ctx: &Context<Self>) {
self.epoch_precompute_in_flight = None;
match msg.result {
Ok(state) => {
// Stored even if the block it was meant for already imported
// inline: attestation targets for the epoch read the same key.
self.store.cache_state(
CacheKey::CheckpointState {
epoch: msg.key.epoch,
root: msg.key.root,
},
state,
);
info!(
epoch = msg.key.epoch,
head_root = %ShortRoot(&msg.key.root.0),
"Epoch precompute stored"
);
}
Err(err) => warn!(
epoch = msg.key.epoch,
head_root = %ShortRoot(&msg.key.root.0),
%err,
"Epoch precompute failed"
),
}
}
}

#[cfg(test)]
mod tests {
use super::*;

const SPE: u64 = preset::SLOTS_PER_EPOCH;

fn inputs(trigger: Trigger, head_slot: u64) -> DecisionInputs {
DecisionInputs {
trigger,
syncing: false,
head_slot,
head_root: H256([7; 32]),
current_slot: head_slot,
in_flight: None,
already_cached: false,
}
}

#[test]
fn a_last_slot_head_precomputes_the_next_epoch() {
let key = decide(&inputs(Trigger::Head, 2 * SPE - 1)).expect("starts");
assert_eq!(key.epoch, 2);
assert_eq!(key.root, H256([7; 32]));
}

#[test]
fn a_head_elsewhere_in_the_epoch_starts_nothing() {
assert!(decide(&inputs(Trigger::Head, SPE)).is_none());
assert!(decide(&inputs(Trigger::Head, 2 * SPE - 2)).is_none());
}

#[test]
fn a_syncing_node_starts_nothing() {
let mut i = inputs(Trigger::Head, 2 * SPE - 1);
i.syncing = true;
assert!(decide(&i).is_none());
}

#[test]
fn a_stale_head_trigger_starts_nothing() {
let mut i = inputs(Trigger::Head, 2 * SPE - 1);
i.current_slot = 2 * SPE + 5;
assert!(decide(&i).is_none());
}

#[test]
fn a_running_worker_blocks_a_second_one() {
let mut i = inputs(Trigger::Head, 2 * SPE - 1);
i.in_flight = Some(PrecomputeKey {
epoch: 2,
root: H256([7; 32]),
});
assert!(decide(&i).is_none());
}

#[test]
fn a_state_already_cached_is_not_recomputed() {
let mut i = inputs(Trigger::Timer { slot: 2 * SPE - 1 }, 2 * SPE - 3);
i.already_cached = true;
assert!(decide(&i).is_none());
}

#[test]
fn the_timer_covers_a_skipped_last_slot() {
let key =
decide(&inputs(Trigger::Timer { slot: 2 * SPE - 1 }, 2 * SPE - 3)).expect("starts");
assert_eq!(key.epoch, 2);
}

#[test]
fn the_timer_ignores_slots_that_are_not_the_last() {
assert!(decide(&inputs(Trigger::Timer { slot: SPE }, SPE - 1)).is_none());
}

#[test]
fn the_timer_skips_a_head_already_at_the_boundary() {
assert!(decide(&inputs(Trigger::Timer { slot: 2 * SPE - 1 }, 2 * SPE)).is_none());
}

#[test]
fn three_quarters_follows_the_configured_slot_duration() {
assert_eq!(three_quarter_slot(12_000), Duration::from_secs(9));
assert_eq!(three_quarter_slot(6_000), Duration::from_millis(4_500));
}
}
Loading
Loading