From a9c7d6b91aab3f482ca07dca9b897c03f2beffff Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Tom=C3=A1s=20Gr=C3=BCner?= <47506558+MegaRedHand@users.noreply.github.com> Date: Fri, 2 Oct 2026 14:44:03 -0300 Subject: [PATCH 1/2] feat(beacon): time the state-writer hand-off as its own import phase Beacon post-states are persisted by one background thread behind a bounded queue, and Store::insert_state blocks on it while the queue is full. That call sits inside fork_choice::on_block, which the chain actor times as the stf phase, so when blocks import faster than the writer persists them the importer's wait was reported as state-transition time. StateWriterHandle::send now returns the instants around its blocking send, insert_state keeps the span (first start, last end) and Store::take_state_handoff hands it to the chain actor, which records it as writer_wait_start/end. The report gains a writer_wait row after stf, and the phase label is added to lean_block_import_phase_seconds. Semantics change: beacon stf now excludes the writer wait (saturating subtraction in the report), so the rows still sum to the end-to-end time. Lean never records a wait, so its report is unchanged. --- crates/blockchain/src/import_timing.rs | 92 +++++++++++++++++++++++++- crates/blockchain/src/lib.rs | 8 +++ crates/blockchain/src/metrics.rs | 1 + crates/storage/src/state_writer.rs | 9 ++- crates/storage/src/store.rs | 39 ++++++++++- docs/benchmarking.md | 11 ++- docs/metrics.md | 3 +- 7 files changed, 157 insertions(+), 6 deletions(-) diff --git a/crates/blockchain/src/import_timing.rs b/crates/blockchain/src/import_timing.rs index 57181f7e..6625a472 100644 --- a/crates/blockchain/src/import_timing.rs +++ b/crates/blockchain/src/import_timing.rs @@ -50,6 +50,14 @@ //! locally built one) starts at the pull, since no earlier moment is //! knowable. //! +//! The beacon `stf` span is `fork_choice::on_block`, which ends by handing the +//! post-state to the storage crate's background writer and blocks while that +//! writer's queue is full. That wait is nested inside the `stf` pair and is +//! recorded as its own `writer_wait` pair. The `stf` row reports its pair +//! minus the wait (saturating), so `stf` is state-transition work and the +//! rows still sum to the end-to-end time. A block with no recorded wait, which +//! includes every lean block, reports `stf` unchanged. +//! //! One consequence worth keeping in mind when reading the numbers: a gossip //! block's total and a fetched block's total do not start at the same point in //! the block's life, and a fetched block's total excludes the request round @@ -79,6 +87,8 @@ pub struct StoreTimings { pub verify_crypto_end: Option, pub stf_start: Option, pub stf_end: Option, + pub writer_wait_start: Option, + pub writer_wait_end: Option, pub db_write_start: Option, pub db_write_end: Option, pub fc_head_start: Option, @@ -181,9 +191,18 @@ pub struct ImportTimings { /// The state transition. On beacon this is `fork_choice::on_block`, which /// bundles the transition, the state root and the state write, so /// `db_write` stays `None` there. + /// + /// This pair brackets the whole call, writer wait included. The `stf` row + /// reports it minus [`Self::writer_wait_start`]..[`Self::writer_wait_end`]. pub stf_start: Option, pub stf_end: Option, + /// Beacon: the blocking hand-off of the post-state to the storage crate's + /// background writer, which waits while the writer's queue is full. It + /// lies inside the `stf` pair; only the report separates the two. + pub writer_wait_start: Option, + pub writer_wait_end: Option, + /// Lean: the block write (`insert_signed_block`) and the state hand-off /// (`insert_state`, which since the storage crate moved state writes to a /// background thread only enqueues the state — cache, buffer and a @@ -311,6 +330,8 @@ impl ImportTimings { self.verify_crypto_end = store.verify_crypto_end; self.stf_start = store.stf_start; self.stf_end = store.stf_end; + self.writer_wait_start = store.writer_wait_start; + self.writer_wait_end = store.writer_wait_end; self.db_write_start = store.db_write_start; self.db_write_end = store.db_write_end; self.fc_head_start = store.fc_head_start; @@ -345,7 +366,16 @@ impl ImportTimings { /// Every section, in the order a block crosses them, whether or not it /// crossed this one. + /// + /// The `stf` row excludes `writer_wait` when one was recorded, since the + /// wait is nested inside the `stf` pair; that keeps the rows summing to + /// the end-to-end time and `stf` meaning state-transition work. fn rows(&self) -> Vec { + let mut stf = Row::new("stf", self.stf_start, self.stf_end); + let writer_wait = Row::new("writer_wait", self.writer_wait_start, self.writer_wait_end); + if let (Some(total), Some(wait)) = (stf.elapsed, writer_wait.elapsed) { + stf.elapsed = Some(total.saturating_sub(wait)); + } vec![ Row::new("decode", self.decode_start, self.decode_end), Row::new("queue", self.queue_start, self.queue_end), @@ -376,7 +406,8 @@ impl ImportTimings { self.verify_crypto_start, self.verify_crypto_end, ), - Row::new("stf", self.stf_start, self.stf_end), + stf, + writer_wait, Row::new("db_write", self.db_write_start, self.db_write_end), Row::new("fc_head", self.fc_head_start, self.fc_head_end), Row::new("block_atts", self.block_atts_start, self.block_atts_end), @@ -821,6 +852,65 @@ mod tests { assert_eq!(beacon.stf_end, Some(base + Duration::from_millis(90))); } + fn row_elapsed(timings: &ImportTimings, name: &str) -> Option { + timings + .rows() + .into_iter() + .find(|row| row.name == name) + .and_then(|row| row.elapsed) + } + + #[test] + fn the_stf_row_excludes_the_writer_wait() { + let base = Instant::now(); + let mut timings = ImportTimings::default(); + (timings.stf_start, timings.stf_end) = pair(base, 0, 90); + (timings.writer_wait_start, timings.writer_wait_end) = pair(base, 60, 30); + + assert_eq!( + row_elapsed(&timings, "stf"), + Some(Duration::from_millis(60)) + ); + assert_eq!( + row_elapsed(&timings, "writer_wait"), + Some(Duration::from_millis(30)) + ); + let present: Vec<&str> = timings + .rows() + .into_iter() + .filter(|row| row.elapsed.is_some()) + .map(|row| row.name) + .collect(); + assert_eq!(present, vec!["stf", "writer_wait"]); + } + + #[test] + fn an_absent_writer_wait_leaves_stf_unchanged() { + let base = Instant::now(); + let mut timings = ImportTimings::default(); + (timings.stf_start, timings.stf_end) = pair(base, 0, 90); + + assert_eq!( + row_elapsed(&timings, "stf"), + Some(Duration::from_millis(90)) + ); + assert_eq!(row_elapsed(&timings, "writer_wait"), None); + } + + #[test] + fn a_writer_wait_longer_than_stf_saturates_stf_to_zero() { + let base = Instant::now(); + let mut timings = ImportTimings::default(); + (timings.stf_start, timings.stf_end) = pair(base, 0, 10); + (timings.writer_wait_start, timings.writer_wait_end) = pair(base, 0, 25); + + assert_eq!(row_elapsed(&timings, "stf"), Some(Duration::ZERO)); + assert_eq!( + row_elapsed(&timings, "writer_wait"), + Some(Duration::from_millis(25)) + ); + } + #[test] fn the_tail_of_an_import_that_ran_nothing_lands_nowhere() { // A re-delivered block whose post-state the store already held runs no diff --git a/crates/blockchain/src/lib.rs b/crates/blockchain/src/lib.rs index 0044889d..07c029dd 100644 --- a/crates/blockchain/src/lib.rs +++ b/crates/blockchain/src/lib.rs @@ -1784,6 +1784,9 @@ impl BlockChainServer { } } + // Drop any hand-off span left by an insert outside this call, + // so the one taken below is this block's alone. + let _ = self.store.take_state_handoff(); timings.stf_start = Some(Instant::now()); let imported = fork_choice::on_block( &mut self.store, @@ -1794,6 +1797,11 @@ impl BlockChainServer { &committees, ); timings.stf_end = Some(Instant::now()); + (timings.writer_wait_start, timings.writer_wait_end) = + match self.store.take_state_handoff() { + Some((start, end)) => (Some(start), Some(end)), + None => (None, None), + }; if let Err(err) = imported { return (timings, Err(err.into())); } diff --git a/crates/blockchain/src/metrics.rs b/crates/blockchain/src/metrics.rs index 3a40d0d7..b5ad80f5 100644 --- a/crates/blockchain/src/metrics.rs +++ b/crates/blockchain/src/metrics.rs @@ -59,6 +59,7 @@ pub const BLOCK_IMPORT_PHASES: &[&str] = &[ "verify_struct", "verify_crypto", "stf", + "writer_wait", "db_write", "fc_head", "block_atts", diff --git a/crates/storage/src/state_writer.rs b/crates/storage/src/state_writer.rs index d8279626..e2de2e0d 100644 --- a/crates/storage/src/state_writer.rs +++ b/crates/storage/src/state_writer.rs @@ -43,6 +43,7 @@ use std::collections::HashMap; use std::sync::mpsc::{Receiver, SyncSender, sync_channel}; use std::sync::{Arc, Mutex}; use std::thread::JoinHandle; +use std::time::Instant; use ethlambda_types::{ beacon::{containers::BeaconState, fork::ForkName}, @@ -611,17 +612,23 @@ impl StateWriterHandle { /// Hands a state to the writer, blocking while the queue is full. /// + /// Returns the instants just before and just after the blocking `send`, + /// so a caller can report the hand-off wait apart from the work around + /// it. The span is near zero unless the queue was full. + /// /// # Panics /// /// If the writer thread is gone, which only happens after it panicked. Its /// own panic message is already on the default hook; this is the importer /// learning about it, one state later. - pub(crate) fn send(&self, request: StateWriteRequest) { + pub(crate) fn send(&self, request: StateWriteRequest) -> (Instant, Instant) { + let start = Instant::now(); self.tx .as_ref() .expect("the sender is taken only in Drop") .send(request) .expect("the state writer thread died; see the logged panic above"); + (start, Instant::now()) } } diff --git a/crates/storage/src/store.rs b/crates/storage/src/store.rs index 96addfbd..610ec68e 100644 --- a/crates/storage/src/store.rs +++ b/crates/storage/src/store.rs @@ -1,6 +1,7 @@ use std::collections::{BTreeMap, HashMap, HashSet, VecDeque}; use std::num::NonZeroUsize; use std::sync::{Arc, LazyLock, Mutex}; +use std::time::Instant; use lru::LruCache; @@ -887,6 +888,10 @@ pub struct Store { /// The background writer, joined when the last clone of this `Store` /// drops. See [`StateWriterHandle`]. state_writer: Arc, + /// The span [`Self::insert_state`] spent handing states to the writer + /// since the last [`Self::take_state_handoff`]: the first start and the + /// last end. Shared across clones, like the writer it describes. + state_handoff: Arc>>, } /// Build an empty state cache sized to [`STATE_CACHE_CAPACITY`]. @@ -1514,6 +1519,7 @@ impl Store { committee_cache: Arc::new(CommitteeCache::default()), beacon: Default::default(), state_writer, + state_handoff: Default::default(), } } @@ -2739,10 +2745,27 @@ impl Store { self.cache_state(CacheKey::BlockState(root), state.clone()); self.pending_states.insert(root, state.clone()); crate::metrics::inc_state_write_queue_depth(); - self.state_writer.send(StateWriteRequest { root, state }); + let (start, end) = self.state_writer.send(StateWriteRequest { root, state }); + // Repeated hand-offs keep the first start and the last end, so the + // span covers every wait since the caller last took it. + let mut handoff = self.state_handoff.lock().unwrap(); + *handoff = Some(match *handoff { + Some((first, _)) => (first, end), + None => (start, end), + }); Ok(()) } + /// Takes the span `insert_state` spent handing states to the writer since + /// the last call, clearing it. `None` when no state was inserted. + /// + /// The span is the blocking `send` only, so it is the importer's wait for + /// the writer (a full queue), not the encode/diff/commit work, which runs + /// on the writer thread. + pub fn take_state_handoff(&self) -> Option<(Instant, Instant)> { + self.state_handoff.lock().unwrap().take() + } + // ============ Attestation Extraction ============ fn record_vote( @@ -5437,6 +5460,20 @@ mod tests { assert_eq!(read.slot(), 7); } + #[test] + fn insert_state_records_the_writer_handoff_and_taking_it_clears_it() { + let mut store = beacon_test_store(Arc::new(InMemoryBackend::new())); + assert!(store.take_state_handoff().is_none()); + + store + .insert_state(H256::from([1u8; 32]), beacon_test_state(7)) + .expect("insert beacon state"); + + let (start, end) = store.take_state_handoff().expect("a hand-off was recorded"); + assert!(start <= end); + assert!(store.take_state_handoff().is_none()); + } + #[test] fn a_beacon_state_with_no_known_parent_block_is_its_own_snapshot() { // A bootstrap or checkpoint-sync anchor is the store's first-ever diff --git a/docs/benchmarking.md b/docs/benchmarking.md index a01f2f1a..19c743a9 100644 --- a/docs/benchmarking.md +++ b/docs/benchmarking.md @@ -253,11 +253,18 @@ workload reads its own histogram. A replayed block reports under without passing for a gossip or sync arrival. The phases are the `BLOCK_IMPORT_PHASES` labels: `decode`, `queue`, `defer`, `admit`, `guards`, `preamble`, `parent_wait`, `cascade_wait`, `da_check`, `columns_wait`, `engine`, -`verify_struct`, `verify_crypto`, `stf`, `db_write`, `fc_head`, `block_atts`; -plus the per-arrival sections that are not spans around the others, `prune`, +`verify_struct`, `verify_crypto`, `stf`, `writer_wait`, `db_write`, `fc_head`, +`block_atts`; plus the per-arrival sections that are not spans around the others, `prune`, `get_head` and `fcu`, since each `import_block` call is one arrival. `get_head` is where the head recomputation after every import is charged. +On beacon, `writer_wait` is the importer's blocking hand-off of the post-state +to the storage crate's background writer, which waits while the writer's queue +is full. It happens inside `fork_choice::on_block`, so `stf` excludes it: `stf` +is state-transition work and the two phases do not overlap. When blocks import +faster than the writer persists them, the wait shows up here rather than in +`stf`. + ### What is excluded, and why Replay runs nothing that a corpus already answers for: diff --git a/docs/metrics.md b/docs/metrics.md index bf215a3c..382bfd96 100644 --- a/docs/metrics.md +++ b/docs/metrics.md @@ -280,7 +280,8 @@ Per-block phases, in the order a block crosses them, plus `total` for a complete | `engine` | The `engine_newPayload` round trip, including its retry ladder | beacon | | `verify_struct` | Participant bounds checks and pubkey resolution | lean | | `verify_crypto` | The leanVM multi-message aggregate verification | lean | -| `stf` | The state transition. On beacon this bundles the transition, the state root and the state write | both | +| `stf` | The state transition. On beacon this bundles the transition and the state root, and excludes `writer_wait`, which is nested inside it | both | +| `writer_wait` | Beacon only: the importer's blocking hand-off of the post-state to the storage crate's background writer, which waits while the writer's queue is full. Subtracted from `stf`, so `stf` and `writer_wait` do not overlap | beacon | | `db_write` | The block write and handing the post-state off to the storage crate's background writer; the state's own encode/diff/commit cost is `lean_state_write_seconds` instead, not this row | lean | | `fc_head` | `update_head` | lean | | `block_atts` | Replaying the block's own attestations and slashings into fork choice | beacon | From 7be479153a7a94eceeb70f6f7e261a42be67f13a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Tom=C3=A1s=20Gr=C3=BCner?= <47506558+MegaRedHand@users.noreply.github.com> Date: Fri, 2 Oct 2026 16:01:23 -0300 Subject: [PATCH 2/2] feat(bench): add --block-delay to the import replay The replay feeds blocks back to back, while a live node gets one per slot. When the build imports faster than the background state writer persists post-states, the writer's queue fills and the importer blocks in writer_wait, so the replay measures the writer instead of the import (one merged build waited 326 ms per block). --block-delay (default 0, today's behavior) sleeps after each block's import, warm-up included, once its phases and wall time are recorded, so the sleep sits in no per-block phase, no wall_seconds and no summary row. The delay is recorded as params.block_delay_ms in the JSON report and printed in the human header. --- bin/ethlambda/src/benchmark/import/mod.rs | 31 ++++++++++++++++++++ bin/ethlambda/src/benchmark/import/replay.rs | 23 +++++++++++++-- bin/ethlambda/src/benchmark/report/import.rs | 20 +++++++++++-- docs/benchmarking.md | 12 ++++++++ 4 files changed, 81 insertions(+), 5 deletions(-) diff --git a/bin/ethlambda/src/benchmark/import/mod.rs b/bin/ethlambda/src/benchmark/import/mod.rs index 79dde168..f8344b80 100644 --- a/bin/ethlambda/src/benchmark/import/mod.rs +++ b/bin/ethlambda/src/benchmark/import/mod.rs @@ -86,6 +86,11 @@ pub(crate) struct ReplayOptions { /// of the same name. #[arg(long, default_value_t = ethlambda_types::beacon::constants::SAFE_SLOTS_TO_IMPORT_OPTIMISTICALLY)] pub safe_slots_to_import_optimistically: u64, + /// Milliseconds to sleep after each block's import before feeding the + /// next, so the background state writer can drain between blocks. The + /// sleep is outside every measured span. 0 feeds blocks back to back. + #[arg(long, value_name = "MILLISECONDS", default_value_t = 0)] + pub block_delay: u64, /// Report format printed to stdout. Logs go to stderr, so JSON output can /// be piped directly (e.g. into jq). #[arg(long, value_enum, default_value_t = OutputFormat::Human)] @@ -169,3 +174,29 @@ async fn run_replay(options: ReplayOptions) -> eyre::Result<()> { Ok(()) } + +#[cfg(test)] +mod tests { + use clap::Parser as _; + + use super::*; + + #[derive(Debug, clap::Parser)] + struct Harness { + #[command(flatten)] + replay: ReplayOptions, + } + + fn parse(extra: &[&str]) -> ReplayOptions { + let base = ["replay", "--corpus", "/tmp/c", "--data-dir", "/tmp/d"]; + Harness::try_parse_from(base.iter().chain(extra)) + .expect("replay flags parse") + .replay + } + + #[test] + fn block_delay_defaults_to_zero_and_parses_milliseconds() { + assert_eq!(parse(&[]).block_delay, 0); + assert_eq!(parse(&["--block-delay", "1500"]).block_delay, 1500); + } +} diff --git a/bin/ethlambda/src/benchmark/import/replay.rs b/bin/ethlambda/src/benchmark/import/replay.rs index 104e8dea..5cf16beb 100644 --- a/bin/ethlambda/src/benchmark/import/replay.rs +++ b/bin/ethlambda/src/benchmark/import/replay.rs @@ -18,7 +18,7 @@ use std::path::Path; use std::sync::Arc; -use std::time::Instant; +use std::time::{Duration, Instant}; use ethlambda_blockchain::metrics::{BLOCK_ARRIVAL_PHASES, BLOCK_IMPORT_PHASES}; use ethlambda_blockchain::{BlockChainServer, ImportOutcome}; @@ -130,6 +130,7 @@ pub(crate) async fn replay_corpus(dir: &Path, options: &ReplayOptions) -> eyre:: position + 1, format_ms(wall_seconds) ); + pause_between_blocks(options.block_delay).await; } let total = manifest.slots.len(); @@ -153,11 +154,14 @@ pub(crate) async fn replay_corpus(dir: &Path, options: &ReplayOptions) -> eyre:: phases, outcome: "imported", }); + // After `finish_at_most_once` and the sample are recorded, so the + // sleep is in no phase delta and no `wall_seconds`. + pause_between_blocks(options.block_delay).await; } Ok(Report::new( Environment::collect(), - params_from(&manifest, dir), + params_from(&manifest, dir, options.block_delay), samples, )) } @@ -233,7 +237,7 @@ fn read_anchor(dir: &Path, config: &Config) -> eyre::Result<(BeaconState, Signed /// The report parameters a manifest already carries, so the loop above never /// has to reconstruct them from samples. -fn params_from(manifest: &Manifest, dir: &Path) -> Params { +fn params_from(manifest: &Manifest, dir: &Path, block_delay_ms: u64) -> Params { Params { mode: "import", corpus: dir.display().to_string(), @@ -244,6 +248,19 @@ fn params_from(manifest: &Manifest, dir: &Path) -> Params { range_start: manifest.range_start, range_end: manifest.range_end, blocks: manifest.slots.len(), + block_delay_ms, + } +} + +/// Sleep `delay_ms` so the state writer can drain before the next block. +/// +/// Called only between measured spans: a replay feeds blocks back to back, +/// faster than a live node's one block per slot, so without a pause the +/// importer can wait on the writer's queue and measure the writer instead of +/// the import. +async fn pause_between_blocks(delay_ms: u64) { + if delay_ms > 0 { + tokio::time::sleep(Duration::from_millis(delay_ms)).await; } } diff --git a/bin/ethlambda/src/benchmark/report/import.rs b/bin/ethlambda/src/benchmark/report/import.rs index 86c976ba..74bb20e8 100644 --- a/bin/ethlambda/src/benchmark/report/import.rs +++ b/bin/ethlambda/src/benchmark/report/import.rs @@ -21,6 +21,9 @@ pub(crate) struct Params { /// Inclusive, as `fetch --to` is. pub range_end: u64, pub blocks: usize, + /// Milliseconds slept after each block's import, outside every measured + /// span. 0 means blocks were fed back to back. + pub block_delay_ms: u64, } #[derive(Debug, Serialize)] @@ -104,14 +107,15 @@ impl Report { let _ = writeln!(out, "Block-import benchmark — {} workload", params.mode); let _ = writeln!( out, - " corpus={} network={} anchor_slot={} warmup_blocks={} range=[{}, {}] blocks={}", + " corpus={} network={} anchor_slot={} warmup_blocks={} range=[{}, {}] blocks={} block_delay_ms={}", params.corpus, params.network, params.anchor_slot, params.warmup_blocks, params.range_start, params.range_end, - params.blocks + params.blocks, + params.block_delay_ms ); let _ = writeln!( out, @@ -212,9 +216,21 @@ mod tests { range_start: 1, range_end: 3, blocks: 2, + block_delay_ms: 0, } } + #[test] + fn the_block_delay_is_recorded_in_the_json_and_the_header() { + let mut params = params(); + params.block_delay_ms = 750; + let report = Report::new(Environment::collect(), params, Vec::new()); + + let json: serde_json::Value = serde_json::from_str(&report.to_json().unwrap()).unwrap(); + assert_eq!(json["params"]["block_delay_ms"], 750); + assert!(report.human_table().contains("block_delay_ms=750")); + } + #[test] fn the_range_prints_inclusive_as_fetch_takes_it() { let report = Report::new(Environment::collect(), params(), Vec::new()); diff --git a/docs/benchmarking.md b/docs/benchmarking.md index 19c743a9..7a69b858 100644 --- a/docs/benchmarking.md +++ b/docs/benchmarking.md @@ -237,6 +237,18 @@ ethlambda benchmark import replay \ --format json --output report.json ``` +**`--block-delay `** (default 0) sleeps after each block's import +before the next block is fed. Use it to measure import work when the build +outpaces the state writer: the replay feeds blocks back to back, while a live +node gets one per slot, so a fast build can fill the writer's queue and +`writer_wait` then measures the writer rather than the import. A delay longer +than one state write (about a second on mainnet) lets the writer drain fully +between blocks. The sleep runs after the block's phases and `wall_seconds` are +recorded, so it appears in no per-block number and no summary row; the report +has no end-to-end figure to adjust. The value is recorded as +`params.block_delay_ms` in the JSON report and printed in the human header, so +compare only reports made with the same delay. + ### What is measured The measured span is one `BlockChainServer::import_block` call per block, the