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
31 changes: 31 additions & 0 deletions bin/ethlambda/src/benchmark/import/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Expand Down Expand Up @@ -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);
}
}
23 changes: 20 additions & 3 deletions bin/ethlambda/src/benchmark/import/replay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -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();
Expand All @@ -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,
))
}
Expand Down Expand Up @@ -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(),
Expand All @@ -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;
}
}

Expand Down
20 changes: 18 additions & 2 deletions bin/ethlambda/src/benchmark/report/import.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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());
Expand Down
92 changes: 91 additions & 1 deletion crates/blockchain/src/import_timing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -79,6 +87,8 @@ pub struct StoreTimings {
pub verify_crypto_end: Option<Instant>,
pub stf_start: Option<Instant>,
pub stf_end: Option<Instant>,
pub writer_wait_start: Option<Instant>,
pub writer_wait_end: Option<Instant>,
pub db_write_start: Option<Instant>,
pub db_write_end: Option<Instant>,
pub fc_head_start: Option<Instant>,
Expand Down Expand Up @@ -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<Instant>,
pub stf_end: Option<Instant>,

/// 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<Instant>,
pub writer_wait_end: Option<Instant>,

/// 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
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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<Row> {
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),
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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<Duration> {
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
Expand Down
8 changes: 8 additions & 0 deletions crates/blockchain/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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()));
}
Expand Down
1 change: 1 addition & 0 deletions crates/blockchain/src/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@ pub const BLOCK_IMPORT_PHASES: &[&str] = &[
"verify_struct",
"verify_crypto",
"stf",
"writer_wait",
"db_write",
"fc_head",
"block_atts",
Expand Down
9 changes: 8 additions & 1 deletion crates/storage/src/state_writer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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},
Expand Down Expand Up @@ -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())
}
}

Expand Down
Loading
Loading