commit bcbe1f9edfb82ca4caf2c490771134ef176d8b49
parent 0c9daaf165e7d546c82932a87394a3d77f20ec2e
Author: Joris Hartog <jorishartog@hotmail.com>
Date: Mon, 17 Aug 2026 00:27:05 +0200
Harden p2p block import and mempool sizing
Diffstat:
8 files changed, 142 insertions(+), 92 deletions(-)
diff --git a/src/adapters/p2p/fetch.rs b/src/adapters/p2p/fetch.rs
@@ -152,11 +152,9 @@ pub(super) async fn validate_snapshot_extension(
.await
.context("chain snapshot adoption worker failed")?;
}
- let missing_blocks = ledger.missing_snapshot_blocks(&snapshot)?;
- verify_blocks_vdf(missing_blocks).await?;
tokio::task::spawn_blocking(move || {
- ledger.extend_from_preverified_snapshot_at(snapshot, now_ms)?;
+ ledger.extend_from_snapshot_at(snapshot, now_ms)?;
Ok(ledger)
})
.await
@@ -171,11 +169,10 @@ pub(super) async fn validate_blocks_extension(
if blocks.is_empty() {
return Ok(ledger);
}
- verify_blocks_vdf(blocks.clone()).await?;
tokio::task::spawn_blocking(move || {
for block in blocks {
- ledger.apply_preverified_block_at(block, now_ms)?;
+ ledger.apply_block_at(block, now_ms)?;
}
Ok(ledger)
})
@@ -207,30 +204,15 @@ pub(super) async fn verify_block_vdf(block: Block) -> Result<Block> {
Ok(block)
}
-async fn verify_blocks_vdf(blocks: Vec<Block>) -> Result<()> {
- let mut tasks = tokio::task::JoinSet::new();
- for block in blocks {
- tasks.spawn_blocking(move || {
- if !verify_vdf(&block.vdf_seed(), block.vdf_rounds, &block.vdf_output) {
- anyhow::bail!("block {} VDF output is invalid", block.height);
- }
- Ok::<(), anyhow::Error>(())
- });
- }
-
- while let Some(result) = tasks.join_next().await {
- result.context("VDF verification worker failed")??;
- }
-
- Ok(())
-}
-
#[cfg(test)]
mod tests {
- use crate::{app::GossipEnvelope, domain::Wallet};
+ use crate::{
+ app::GossipEnvelope,
+ domain::{Block, FinalizerMode, RevealBundleSection, Wallet},
+ };
use super::super::test_support::{allocations, node};
- use super::join_snapshot_response;
+ use super::{join_snapshot_response, validate_blocks_extension};
#[test]
fn join_snapshot_response_ignores_status_noise_before_snapshot() {
@@ -262,4 +244,40 @@ mod tests {
assert_eq!(parsed, Some(snapshot));
}
+
+ #[tokio::test]
+ async fn block_batch_prechecks_before_vdf_verification() {
+ let alice = Wallet::from_seed("batch-precheck-alice");
+ let test_node = node(
+ "alice",
+ alice.clone(),
+ allocations(std::slice::from_ref(&alice), 1_000),
+ );
+ let tip = test_node.chain_snapshot().blocks.last().unwrap().clone();
+ let ledger = test_node.clone_ledger();
+ let invalid_height_block = Block {
+ height: ledger.height() + 2,
+ prev_hash: tip.hash,
+ timestamp_ms: tip.timestamp_ms + 1,
+ miner: alice.address().to_string(),
+ finalizer_mode: FinalizerMode::Ticket,
+ finalizer_rank: 0,
+ reward: 0,
+ vdf_rounds: ledger.vdf_rounds(),
+ vdf_output: "not-a-vdf-solution".to_string(),
+ leader_proof: None,
+ blinded_transactions: Vec::new(),
+ reveal_bundle_section: RevealBundleSection::default(),
+ transactions: Vec::new(),
+ hash: "invalid-hash".to_string(),
+ };
+
+ let error = validate_blocks_extension(ledger, vec![invalid_height_block], u64::MAX)
+ .await
+ .unwrap_err();
+ let message = format!("{error:#}");
+
+ assert!(message.contains("expected block height"));
+ assert!(!message.contains("VDF output is invalid"));
+ }
}
diff --git a/src/domain/ledger_apply.rs b/src/domain/ledger_apply.rs
@@ -227,6 +227,7 @@ impl Ledger {
&& self.pending_reveal_transaction(reveal).is_ok()
})
.collect();
+ self.refresh_pending_pool_byte_counters()?;
self.promote_orphan_transactions()?;
self.vdf_rounds = self.next_vdf_rounds_after_tip();
Ok(())
diff --git a/src/domain/ledger_chain.rs b/src/domain/ledger_chain.rs
@@ -7,7 +7,7 @@ use super::genesis::{build_genesis_block, utxos_after_genesis, validate_genesis_
use super::ledger_ops::validate_genesis_allocations;
use super::ticket::genesis_tickets;
use super::{
- Amount, Block, ChainSnapshot, GenesisBurn, LaunchProfile, Ledger, MINE_REWARD, Transaction,
+ Amount, ChainSnapshot, GenesisBurn, LaunchProfile, Ledger, MINE_REWARD, Transaction,
unix_now_ms,
};
@@ -54,6 +54,10 @@ impl Ledger {
orphans: Vec::new(),
pending_blinded: Vec::new(),
pending_reveals: Vec::new(),
+ pending_bytes: 0,
+ orphan_bytes: 0,
+ pending_blinded_bytes: 0,
+ pending_reveal_bytes: 0,
active_blinded: BTreeMap::new(),
mine_reward: MINE_REWARD,
initial_vdf_rounds: vdf_rounds,
@@ -109,6 +113,10 @@ impl Ledger {
orphans: Vec::new(),
pending_blinded: Vec::new(),
pending_reveals: Vec::new(),
+ pending_bytes: 0,
+ orphan_bytes: 0,
+ pending_blinded_bytes: 0,
+ pending_reveal_bytes: 0,
active_blinded: BTreeMap::new(),
mine_reward: MINE_REWARD,
initial_vdf_rounds: vdf_rounds,
@@ -135,27 +143,12 @@ impl Ledger {
self.extend_from_snapshot_with_vdf_policy(snapshot, true, unix_now_ms())
}
- pub(crate) fn extend_from_preverified_snapshot_at(
+ pub(crate) fn extend_from_snapshot_at(
&mut self,
snapshot: ChainSnapshot,
now_ms: u64,
) -> Result<bool> {
- self.extend_from_snapshot_with_vdf_policy(snapshot, false, now_ms)
- }
-
- pub(crate) fn missing_snapshot_blocks(&self, snapshot: &ChainSnapshot) -> Result<Vec<Block>> {
- let remote_height = self.validate_snapshot_identity(snapshot)?;
- if remote_height <= self.height() {
- return Ok(Vec::new());
- }
- let common_ancestor_height = self.common_ancestor_height(snapshot)?;
-
- Ok(snapshot
- .blocks
- .iter()
- .skip(common_ancestor_height as usize + 1)
- .cloned()
- .collect())
+ self.extend_from_snapshot_with_vdf_policy(snapshot, true, now_ms)
}
fn extend_from_snapshot_with_vdf_policy(
@@ -203,20 +196,6 @@ impl Ledger {
Ok(remote_height)
}
- fn common_ancestor_height(&self, snapshot: &ChainSnapshot) -> Result<u64> {
- self.validate_snapshot_identity(snapshot)?;
- let max_common_index = self.chain.len().min(snapshot.blocks.len()) - 1;
- for index in 0..=max_common_index {
- if self.chain[index] != snapshot.blocks[index] {
- if index == 0 {
- bail!("chain snapshot has no common genesis block");
- }
- return Ok(index as u64 - 1);
- }
- }
- Ok(max_common_index as u64)
- }
-
fn fork_point_with_candidate(&self, candidate: &Ledger) -> Result<ForkPoint> {
if candidate.genesis_hash() != self.genesis_hash() {
bail!("candidate chain has no common genesis block");
diff --git a/src/domain/ledger_mempool.rs b/src/domain/ledger_mempool.rs
@@ -48,13 +48,14 @@ impl Ledger {
if self.pending_blinded.len() >= MAX_PENDING_TRANSACTIONS {
bail!("blinded mempool is full");
}
- ensure_pending_pool_bytes(
+ let candidate_bytes = ensure_pending_pool_bytes(
"blinded mempool",
- &self.pending_blinded,
+ self.pending_blinded_bytes,
&transaction,
MAX_PENDING_POOL_BYTES,
)?;
self.pending_blinded.push(transaction);
+ self.pending_blinded_bytes = self.pending_blinded_bytes.saturating_add(candidate_bytes);
Ok(true)
}
@@ -65,30 +66,34 @@ impl Ledger {
self.validate_blinded_reveal_terms(&reveal)?;
self.pending_reveal_transaction(&reveal)?;
if self.pending_reveals.len() >= MAX_PENDING_TRANSACTIONS
- && !self.drop_one_invalid_pending_blinded_reveal()
+ && !self.drop_one_invalid_pending_blinded_reveal()?
{
bail!("blinded reveal pool is full");
}
- ensure_pending_pool_bytes(
+ let candidate_bytes = ensure_pending_pool_bytes(
"blinded reveal pool",
- &self.pending_reveals,
+ self.pending_reveal_bytes,
&reveal,
MAX_PENDING_POOL_BYTES,
)?;
self.pending_reveals.push(reveal);
+ self.pending_reveal_bytes = self.pending_reveal_bytes.saturating_add(candidate_bytes);
Ok(true)
}
- fn drop_one_invalid_pending_blinded_reveal(&mut self) -> bool {
+ fn drop_one_invalid_pending_blinded_reveal(&mut self) -> Result<bool> {
let Some(index) = self
.pending_reveals
.iter()
.position(|reveal| self.pending_reveal_transaction(reveal).is_err())
else {
- return false;
+ return Ok(false);
};
- self.pending_reveals.remove(index);
- true
+ let removed = self.pending_reveals.remove(index);
+ self.pending_reveal_bytes = self
+ .pending_reveal_bytes
+ .saturating_sub(serialized_len(&removed)?);
+ Ok(true)
}
pub fn submit_transaction_with_outcome(
@@ -120,50 +125,68 @@ impl Ledger {
if self.orphans.len() >= MAX_ORPHAN_TRANSACTIONS {
bail!("orphan transaction pool is full");
}
- ensure_pending_pool_bytes(
+ let candidate_bytes = ensure_pending_pool_bytes(
"orphan transaction pool",
- &self.orphans,
+ self.orphan_bytes,
&transaction,
MAX_PENDING_POOL_BYTES,
)?;
self.orphans.push(transaction);
+ self.orphan_bytes = self.orphan_bytes.saturating_add(candidate_bytes);
return Ok(TransactionSubmitOutcome::Added);
}
apply_transaction(&transaction, &mut utxos)?;
- ensure_pending_pool_bytes(
+ let candidate_bytes = ensure_pending_pool_bytes(
"mempool",
- &self.pending,
+ self.pending_bytes,
&transaction,
MAX_PENDING_POOL_BYTES,
)?;
self.pending.push(transaction);
+ self.pending_bytes = self.pending_bytes.saturating_add(candidate_bytes);
self.promote_orphan_transactions()?;
Ok(TransactionSubmitOutcome::Added)
}
+
+ pub(super) fn refresh_pending_pool_byte_counters(&mut self) -> Result<()> {
+ self.pending_bytes = serialized_pool_len(&self.pending)?;
+ self.orphan_bytes = serialized_pool_len(&self.orphans)?;
+ self.pending_blinded_bytes = serialized_pool_len(&self.pending_blinded)?;
+ self.pending_reveal_bytes = serialized_pool_len(&self.pending_reveals)?;
+ Ok(())
+ }
}
fn ensure_pending_pool_bytes<T: Serialize>(
label: &str,
- existing: &[T],
+ existing_bytes: usize,
candidate: &T,
max_bytes: usize,
-) -> Result<()> {
- let existing_bytes = existing.iter().try_fold(0usize, |total, item| {
- let bytes = serde_json::to_vec(item)
- .context("failed to serialize pending item for size check")?
- .len();
- total
- .checked_add(bytes)
- .context("pending pool byte size overflow")
- })?;
- let candidate_bytes = serde_json::to_vec(candidate)
- .context("failed to serialize pending item for size check")?
- .len();
+) -> Result<usize> {
+ let candidate_bytes = serialized_len(candidate)?;
let total_bytes = existing_bytes
.checked_add(candidate_bytes)
.context("pending pool byte size overflow")?;
if total_bytes > max_bytes {
bail!("{label} byte limit exceeded");
}
- Ok(())
+ Ok(candidate_bytes)
+}
+
+fn serialized_pool_len<T: Serialize>(items: &[T]) -> Result<usize> {
+ items.iter().try_fold(0usize, |total, item| {
+ total
+ .checked_add(serialized_len(item)?)
+ .context("pending pool byte size overflow")
+ })
+}
+
+pub(super) fn pending_pool_item_bytes<T: Serialize>(item: &T) -> Result<usize> {
+ serialized_len(item)
+}
+
+fn serialized_len<T: Serialize>(item: &T) -> Result<usize> {
+ serde_json::to_vec(item)
+ .context("failed to serialize pending item for size check")
+ .map(|bytes| bytes.len())
}
diff --git a/src/domain/ledger_pending.rs b/src/domain/ledger_pending.rs
@@ -7,6 +7,7 @@ use super::blinded::{
blinded_reveal_inputs_match, blinded_transaction_commitment, decrypt_blinded_transaction,
verify_blinded_input_signatures,
};
+use super::ledger_mempool::pending_pool_item_bytes;
use super::ledger_ops::{
apply_spendable_pending_transaction, apply_transaction, best_selectable_blinded_index,
best_selectable_burn_from_index, best_selectable_transaction_index,
@@ -31,7 +32,7 @@ use super::validation::{
use super::{
Amount, BLINDED_KEY_BYTES, BLINDED_NONCE_BYTES, BLINDED_VISIBLE_INPUTS_REQUIRED_HEIGHT,
BLOCK_ITEM_FEES_REQUIRED_HEIGHT, BlindedReveal, BlindedTransaction, Ledger,
- MAX_BLINDED_TRANSACTION_EXPIRY_HEIGHTS, MAX_PENDING_TRANSACTIONS,
+ MAX_BLINDED_TRANSACTION_EXPIRY_HEIGHTS, MAX_PENDING_POOL_BYTES, MAX_PENDING_TRANSACTIONS,
MINE_ACTIONS_PER_ANCHOR_LIMIT, OutPoint, RevealBundleSection, Transaction, TxOutput,
decode_hex, decode_hex_array,
};
@@ -369,7 +370,7 @@ impl Ledger {
if self.pending.len() >= MAX_PENDING_TRANSACTIONS {
return Ok(());
}
- let mut promoted_index = None;
+ let mut promoted = None;
let mut utxos = self.utxos_after_valid_pending_and_blinded()?;
for (index, transaction) in self.orphans.iter().enumerate() {
if transaction_inputs_spent_by(transaction, &self.pending) {
@@ -381,15 +382,26 @@ impl Ledger {
if self.validate_new_transaction(transaction).is_ok()
&& apply_transaction(transaction, &mut utxos).is_ok()
{
- promoted_index = Some(index);
+ let transaction_bytes = pending_pool_item_bytes(transaction)?;
+ let promoted_bytes = self
+ .pending_bytes
+ .checked_add(transaction_bytes)
+ .context("pending pool byte size overflow")?;
+ if promoted_bytes > MAX_PENDING_POOL_BYTES {
+ continue;
+ }
+ promoted = Some((index, transaction_bytes));
break;
}
}
- let Some(index) = promoted_index else {
+ let Some((index, transaction_bytes)) = promoted else {
return Ok(());
};
- self.pending.push(self.orphans.remove(index));
+ let transaction = self.orphans.remove(index);
+ self.orphan_bytes = self.orphan_bytes.saturating_sub(transaction_bytes);
+ self.pending.push(transaction);
+ self.pending_bytes = self.pending_bytes.saturating_add(transaction_bytes);
}
}
diff --git a/src/domain/ledger_queries.rs b/src/domain/ledger_queries.rs
@@ -332,14 +332,17 @@ impl Ledger {
.iter()
.any(|input| spent.contains(&input.outpoint))
});
+ let _ = self.refresh_pending_pool_byte_counters();
}
pub(crate) fn clear_pending_blinded_transactions(&mut self) {
self.pending_blinded.clear();
+ self.pending_blinded_bytes = 0;
}
pub(crate) fn clear_pending_transactions(&mut self) {
self.pending.clear();
+ self.pending_bytes = 0;
}
pub fn orphan_transactions(&self) -> &[Transaction] {
diff --git a/src/domain/ledger_state.rs b/src/domain/ledger_state.rs
@@ -18,6 +18,10 @@ pub struct Ledger {
pub(super) orphans: Vec<Transaction>,
pub(super) pending_blinded: Vec<BlindedTransaction>,
pub(super) pending_reveals: Vec<BlindedReveal>,
+ pub(super) pending_bytes: usize,
+ pub(super) orphan_bytes: usize,
+ pub(super) pending_blinded_bytes: usize,
+ pub(super) pending_reveal_bytes: usize,
pub(super) active_blinded: BTreeMap<String, ActiveBlindedTransaction>,
pub(super) mine_reward: Amount,
pub(super) initial_vdf_rounds: u64,
diff --git a/src/domain/tests.rs b/src/domain/tests.rs
@@ -214,6 +214,14 @@ fn mine_preverified_as_next_leader(
block
}
+fn mine_valid_as_next_leader(ledger: &mut Ledger, wallets: &[Wallet], timestamp_ms: u64) -> Block {
+ let leader = ledger.expected_leader_for_next_block().unwrap();
+ let wallet = wallet_for_address(wallets, &leader);
+ let block = ledger.mine_next_block(wallet, timestamp_ms).unwrap();
+ ledger.apply_block_at(block.clone(), u64::MAX).unwrap();
+ block
+}
+
fn mine_preverified_as_next_leader_with_reveal_bundles(
ledger: &mut Ledger,
wallets: &[Wallet],
@@ -2120,6 +2128,7 @@ fn blinded_mempool_rejects_byte_limit_even_before_height_750() {
ledger.pending_blinded = (0..existing_count)
.map(large_inputless_zero_fee_blinded_spam)
.collect();
+ ledger.refresh_pending_pool_byte_counters().unwrap();
let error = ledger
.submit_blinded_transaction(large_inputless_zero_fee_blinded_spam(existing_count))
@@ -2516,6 +2525,7 @@ fn active_blinded_reveal_displaces_invalid_legacy_reveal_when_pool_is_full() {
key: "00".repeat(BLINDED_KEY_BYTES),
})
.collect();
+ ledger.refresh_pending_pool_byte_counters().unwrap();
assert_eq!(
ledger.pending_blinded_reveals().len(),
MAX_PENDING_TRANSACTIONS
@@ -2852,19 +2862,19 @@ fn abandoned_fork_blinded_transactions_return_to_mempool() {
.submit_blinded_transaction(blinded.transaction.clone())
.unwrap();
queue_next_leader_burn(&mut local, &finalizers);
- mine_preverified_as_next_leader(&mut local, &finalizers, 1);
+ mine_valid_as_next_leader(&mut local, &finalizers, 1);
for timestamp_ms in [1, 2] {
let leader = remote.expected_leader_for_next_block().unwrap();
let wallet = wallet_for_address(&finalizers, &leader);
let burn = remote.build_burn(wallet, 1, 0).unwrap();
remote.submit_transaction(burn).unwrap();
- mine_preverified_as_next_leader(&mut remote, &finalizers, timestamp_ms);
+ mine_valid_as_next_leader(&mut remote, &finalizers, timestamp_ms);
}
assert!(
local
- .extend_from_preverified_snapshot_at(remote.snapshot(), u64::MAX)
+ .extend_from_snapshot_at(remote.snapshot(), u64::MAX)
.unwrap()
);
assert!(local.has_blinded_transaction(&blinded.transaction.commitment));