commit 7e1a99500eaf12eb018cf7f447b956b08e4cc880
parent 6886cbd060443c45f7e93bcc5e8b1b32bda3ce4b
Author: Joris Hartog <jorishartog@hotmail.com>
Date: Mon, 21 Sep 2026 16:59:18 +0200
perf(p2p): parallelize sync VDF verification
Diffstat:
4 files changed, 151 insertions(+), 25 deletions(-)
diff --git a/src/adapters/p2p.rs b/src/adapters/p2p.rs
@@ -29,7 +29,9 @@ mod sync;
mod test_support;
mod writer;
use error::SyncError;
-pub use fetch::{fetch_peer_height, fetch_snapshot, fetch_snapshot_with_announcement};
+pub use fetch::{
+ fetch_peer_height, fetch_snapshot, fetch_snapshot_with_announcement, validate_chain_snapshot,
+};
use fetch::{
network_adjusted_time_ms, validate_blocks_extension, validate_chain_bootstrap, verify_block_vdf,
};
diff --git a/src/adapters/p2p/fetch.rs b/src/adapters/p2p/fetch.rs
@@ -1,4 +1,4 @@
-use std::{collections::BTreeMap, net::SocketAddr};
+use std::{collections::BTreeMap, net::SocketAddr, time::Instant};
use anyhow::{Context, Result};
use tokio::{
@@ -8,9 +8,9 @@ use tokio::{
use crate::{
app::{
- ChainBootstrap, GossipEnvelope, NETWORK_ID, PROTOCOL_VERSION, ProtocolHello, now_ms,
- protocol_capabilities, validate_network_genesis, validate_protocol_capabilities,
- validate_transaction_v2_peer_capability,
+ ChainBootstrap, GossipEnvelope, NETWORK_ID, PROTOCOL_VERSION, ProtocolHello,
+ debug_logging_enabled, now_ms, protocol_capabilities, validate_network_genesis,
+ validate_protocol_capabilities, validate_transaction_v2_peer_capability,
},
domain::{Block, ChainSnapshot, LaunchProfile, Ledger, verify_vdf},
};
@@ -28,6 +28,32 @@ pub async fn fetch_peer_height(peer: &str) -> Result<u64> {
fetch_peer_status(peer).await.map(|status| status.height)
}
+/// Validates a downloaded snapshot without making VDF verification a serial part of replay.
+/// State-dependent consensus rules are checked first, then independent VDF proofs are checked
+/// across a bounded number of worker threads before the candidate ledger is returned.
+pub async fn validate_chain_snapshot(snapshot: ChainSnapshot) -> Result<Ledger> {
+ let validation_time_ms = now_ms();
+ tokio::task::spawn_blocking(move || {
+ let started = Instant::now();
+ let ledger = Ledger::from_preverified_snapshot_at(snapshot, validation_time_ms)?;
+ let replay_elapsed = started.elapsed();
+
+ let vdf_started = Instant::now();
+ verify_block_vdfs_parallel(ledger.chain().iter().skip(1))?;
+ if debug_logging_enabled() {
+ eprintln!(
+ "initial chain validation: state={:.3}s vdf={:.3}s blocks={}",
+ replay_elapsed.as_secs_f64(),
+ vdf_started.elapsed().as_secs_f64(),
+ ledger.chain().len().saturating_sub(1),
+ );
+ }
+ Ok(ledger)
+ })
+ .await
+ .context("chain snapshot validation worker failed")?
+}
+
async fn fetch_peer_status(peer: &str) -> Result<PeerStatus> {
let stream = TcpStream::connect(peer)
.await
@@ -253,6 +279,7 @@ pub(super) async fn validate_blocks_extension(
}
tokio::task::spawn_blocking(move || {
+ let state_started = Instant::now();
let target_height = blocks.last().map(|block| block.height);
if blocks[0].prev_hash != ledger.tip_hash() {
let mut candidate = ledger.snapshot();
@@ -261,27 +288,27 @@ pub(super) async fn validate_blocks_extension(
.iter()
.position(|block| block.hash == blocks[0].prev_hash)
.ok_or(super::SyncError::BlockPageHasNoCommonAncestor)?;
- #[cfg(feature = "e2e")]
- for block in &blocks {
- if !verify_vdf(&block.vdf_seed(), block.vdf_rounds, &block.vdf_output) {
- anyhow::bail!("block VDF output is invalid");
- }
- }
candidate.blocks.truncate(ancestor + 1);
- candidate.blocks.extend(blocks);
- #[cfg(feature = "e2e")]
- ledger.extend_from_preverified_snapshot_for_e2e(candidate)?;
- #[cfg(not(feature = "e2e"))]
- ledger.extend_from_snapshot_at(candidate, now_ms)?;
+ candidate.blocks.extend(blocks.iter().cloned());
+ ledger.extend_from_preverified_snapshot_at(candidate, now_ms)?;
+ let state_elapsed = state_started.elapsed();
+ let vdf_started = Instant::now();
+ verify_block_vdfs_parallel(&blocks)?;
+ log_batch_validation_timing(blocks.len(), state_elapsed, vdf_started.elapsed(), true);
if let Some(target_height) = target_height {
on_progress(target_height);
}
return Ok(ledger);
}
- for block in blocks {
- let height = block.height;
- ledger.apply_block_at(block, now_ms)?;
- on_progress(height);
+ for block in blocks.iter().cloned() {
+ ledger.apply_preverified_block_at(block, now_ms)?;
+ }
+ let state_elapsed = state_started.elapsed();
+ let vdf_started = Instant::now();
+ verify_block_vdfs_parallel(&blocks)?;
+ log_batch_validation_timing(blocks.len(), state_elapsed, vdf_started.elapsed(), false);
+ for block in &blocks {
+ on_progress(block.height);
}
Ok(ledger)
})
@@ -289,6 +316,67 @@ pub(super) async fn validate_blocks_extension(
.context("block batch extension worker failed")?
}
+fn verify_block_vdfs_parallel<'a>(blocks: impl IntoIterator<Item = &'a Block>) -> Result<()> {
+ let blocks = blocks.into_iter().collect::<Vec<_>>();
+ if blocks.is_empty() {
+ return Ok(());
+ }
+
+ // Two chain candidates may be validated concurrently by the network coordinator. Giving
+ // each validation at most half the available CPUs prevents the pair from oversubscribing the
+ // machine, while the upper bound keeps untrusted batches from creating excessive threads.
+ let available = std::thread::available_parallelism()
+ .map(usize::from)
+ .unwrap_or(1);
+ let workers = available.div_ceil(2).clamp(1, 8).min(blocks.len());
+ let chunk_size = blocks.len().div_ceil(workers);
+ let invalid_height = std::thread::scope(|scope| -> Result<Option<u64>> {
+ let handles = blocks
+ .chunks(chunk_size)
+ .map(|chunk| {
+ scope.spawn(move || {
+ chunk.iter().find_map(|block| {
+ (!verify_vdf(&block.vdf_seed(), block.vdf_rounds, &block.vdf_output))
+ .then_some(block.height)
+ })
+ })
+ })
+ .collect::<Vec<_>>();
+
+ let mut invalid_height = None;
+ for handle in handles {
+ let height = handle
+ .join()
+ .map_err(|_| anyhow::anyhow!("VDF verification worker panicked"))?;
+ invalid_height = match (invalid_height, height) {
+ (Some(left), Some(right)) => Some(left.min(right)),
+ (height @ Some(_), None) | (None, height @ Some(_)) => height,
+ (None, None) => None,
+ };
+ }
+ Ok(invalid_height)
+ })?;
+ if invalid_height.is_some() {
+ anyhow::bail!("block VDF output is invalid");
+ }
+ Ok(())
+}
+
+fn log_batch_validation_timing(
+ blocks: usize,
+ state_elapsed: std::time::Duration,
+ vdf_elapsed: std::time::Duration,
+ fork: bool,
+) {
+ if debug_logging_enabled() {
+ eprintln!(
+ "block batch validation: state={:.3}s vdf={:.3}s blocks={blocks} fork={fork}",
+ state_elapsed.as_secs_f64(),
+ vdf_elapsed.as_secs_f64(),
+ );
+ }
+}
+
pub(super) async fn network_adjusted_time_ms(network: &GossipNetwork) -> u64 {
let local_time_ms = now_ms();
network
@@ -315,8 +403,11 @@ pub(super) async fn verify_block_vdf(block: Block) -> Result<Block> {
#[cfg(test)]
mod tests {
- use super::join_client_hello;
+ use std::collections::BTreeMap;
+
+ use super::{join_client_hello, verify_block_vdfs_parallel};
use crate::app::{GossipEnvelope, NETWORK_ID, PROTOCOL_VERSION};
+ use crate::domain::Ledger;
#[test]
fn snapshot_join_identifies_as_an_unannounced_setup_placeholder() {
@@ -331,4 +422,19 @@ mod tests {
assert!(hello.listen_addr.is_none());
assert!(hello.node_id.is_none());
}
+
+ #[test]
+ fn parallel_vdf_verification_rejects_an_invalid_proof() {
+ let genesis = Ledger::new(BTreeMap::new(), 1).chain()[0].clone();
+ let mut later = genesis.clone();
+ later.height = 9;
+ later.vdf_output = "invalid-vdf".to_string();
+ let mut earlier = genesis;
+ earlier.height = 3;
+ earlier.vdf_output = "also-invalid".to_string();
+
+ let error = verify_block_vdfs_parallel([&later, &earlier]).unwrap_err();
+
+ assert_eq!(error.to_string(), "block VDF output is invalid");
+ }
}
diff --git a/src/domain/ledger_chain.rs b/src/domain/ledger_chain.rs
@@ -143,7 +143,14 @@ impl Ledger {
}
pub(crate) fn from_preverified_snapshot(snapshot: ChainSnapshot) -> Result<Self> {
- Self::from_snapshot_with_vdf_policy(snapshot, false, u64::MAX)
+ Self::from_preverified_snapshot_at(snapshot, u64::MAX)
+ }
+
+ pub(crate) fn from_preverified_snapshot_at(
+ snapshot: ChainSnapshot,
+ now_ms: u64,
+ ) -> Result<Self> {
+ Self::from_snapshot_with_vdf_policy(snapshot, false, now_ms)
}
pub(crate) fn from_snapshot_at(snapshot: ChainSnapshot, now_ms: u64) -> Result<Self> {
@@ -256,10 +263,21 @@ impl Ledger {
&mut self,
snapshot: ChainSnapshot,
) -> Result<bool> {
- self.extend_from_snapshot_with_vdf_policy(snapshot, false, u64::MAX)
+ self.extend_from_preverified_snapshot_at(snapshot, u64::MAX)
+ }
+
+ /// Applies all normal snapshot and fork-choice validation while skipping only VDF proof
+ /// verification. Callers must verify every untrusted block proof separately before adopting
+ /// the returned ledger.
+ pub(crate) fn extend_from_preverified_snapshot_at(
+ &mut self,
+ snapshot: ChainSnapshot,
+ now_ms: u64,
+ ) -> Result<bool> {
+ self.extend_from_snapshot_with_vdf_policy(snapshot, false, now_ms)
}
- #[cfg_attr(feature = "e2e", allow(dead_code))]
+ #[allow(dead_code)]
pub(crate) fn extend_from_snapshot_at(
&mut self,
snapshot: ChainSnapshot,
diff --git a/src/main.rs b/src/main.rs
@@ -811,7 +811,7 @@ async fn join_chain_ledger(
.map(|block| block.height)
.unwrap_or(0);
println!("joined chain from {peer} at height {height}");
- return Ledger::from_snapshot(snapshot);
+ return p2p::validate_chain_snapshot(snapshot).await;
}
Err(error) => {
errors.push(format!("{peer}: {error:#}"));