diff --git a/Cargo.lock b/Cargo.lock index 4d4e38b57..472410069 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6476,6 +6476,7 @@ dependencies = [ "helix-tcp-types", "hex", "rand 0.9.5", + "rand_xorshift 0.4.0", "rayon", "rustc-hash", "secp256k1 0.30.0", diff --git a/crates/builder/Cargo.toml b/crates/builder/Cargo.toml index 0c88f4c68..608a82e7c 100644 --- a/crates/builder/Cargo.toml +++ b/crates/builder/Cargo.toml @@ -52,3 +52,4 @@ uuid = { workspace = true, features = ["serde"] } zstd.workspace = true [dev-dependencies] +rand_xorshift.workspace = true diff --git a/crates/builder/src/engine/mod.rs b/crates/builder/src/engine/mod.rs index 669f30b13..52a62fe4e 100644 --- a/crates/builder/src/engine/mod.rs +++ b/crates/builder/src/engine/mod.rs @@ -26,7 +26,7 @@ use helix_tcp_types::merging::{ MAX_BLOCK_TXS, MAX_ORDERS_PER_BLOCK, MAX_TX_BYTES, MergeOrderRef, OrderMeta, bundle_order_hash, is_tx_hash_ref, order_id, }, - relay_to_builder::{MergeableBlockV1, RevokeOrderV1, SlotStartV1}, + relay_to_builder::{MergeableBlockV1, SlotStartV1}, }; use ssz::Decode; use tokio::sync::watch; @@ -72,14 +72,6 @@ pub enum EngineEvent { recv_ns: u64, generation: u64, }, - /// A previously-pooled `latest_only` order was dropped from the - /// originating builder's newest submission this slot; the relay tells us - /// rather than us inferring it, since only the relay sees every - /// submission. - RevokeOrder { - msg: RevokeOrderV1, - generation: u64, - }, } /// Outputs to the server tile. `response_id` on `MergedBlockV1` is left 0; the @@ -307,37 +299,6 @@ impl MergeEngine { state.pending_activation = Some((block_hash, recv_ns)); true } - EngineEvent::RevokeOrder { msg, generation } => { - if generation != self.generation { - return false; - } - let Some(state) = self.slot.as_mut() else { return false }; - if msg.slot != state.slot { - return false; - } - let revoked_id = order_id(msg.order_hash, &msg.builder_pubkey); - if !state.remove_order(revoked_id, &msg.builder_pubkey) { - return false; - } - // A parked session may carry the revoked order forward if - // it's later resumed; drop them all rather than tracking - // which one(s) actually applied it. - for parked in state.parked.drain(..) { - parked.log_stats("revoked_parked_drop"); - } - let Some(session) = state.session.as_ref() else { return false }; - if !session.has_applied(&revoked_id) { - return false; - } - warn!( - order_hash = %msg.order_hash, - "live session already applied a just-revoked order, rebuilding" - ); - let base_block_hash = session.base_block_hash; - state.session = None; - state.pending_activation = Some((base_block_hash, crate::utils::utcnow_ns())); - true - } } } @@ -463,10 +424,10 @@ impl MergeEngine { } } - let SlotState { slot, proposer_fee_recipient, orders, session, .. } = state; + let SlotState { slot, proposer_fee_recipient, orders, session, excluded, .. } = state; let Some(session) = session.as_mut() else { return }; - let changed = session.try_extend(orders); + let changed = session.try_extend(orders, excluded); if !changed && !session.pending_emission && !session.has_pending_revenue() { return; } @@ -534,10 +495,24 @@ impl MergeEngine { .map_err(|err| fail(MergeError::InvalidOrder(err)))?; let txs = Arc::new(decoded); + let prepared_orders: Vec = msg + .merge_orders + .iter() + .map(|order_ref| prepare_order(&msg, order_ref, &txs, block_hash)) + .collect(); + + state.update_latest_only( + msg.builder_pubkey, + prepared_orders + .iter() + .filter(|order| order.latest_only) + .map(|order| order.order_hash) + .collect(), + ); + // Budget counts distinct pooled orders, as the relay's `orders_sent` does. let mut pool_full = false; - for order_ref in &msg.merge_orders { - let prepared = prepare_order(&msg, order_ref, &txs, block_hash); + for prepared in prepared_orders { match state.order_ids.get(&prepared.order_id) { Some(&existing_ix) => { // Duplicate order: attribution goes to the highest-value @@ -713,6 +688,8 @@ fn prepare_order( PreparedOrder { order_id: meta.order_id(), + order_hash, + latest_only: matches!(order_ref, MergeOrderRef::Bundle(b) if b.latest_only), origin: msg.builder_address, builder_pubkey: msg.builder_pubkey, source_block_hash: block_hash, diff --git a/crates/builder/src/engine/session.rs b/crates/builder/src/engine/session.rs index 65cbae026..31cff133d 100644 --- a/crates/builder/src/engine/session.rs +++ b/crates/builder/src/engine/session.rs @@ -66,6 +66,7 @@ pub struct MergeStats { pub emit_not_improved: u64, pub emit_no_revenue: u64, pub emit_throttled: u64, + pub orders_excluded_skipped: u64, } impl MergeStats { @@ -424,14 +425,18 @@ impl MergeSession { /// Presimulates candidate orders in parallel, then greedily applies them /// best-payment-first to the live context. Returns whether the block /// changed. Port of `append_greedily_until_gas_limit`. - pub fn try_extend(&mut self, orders: &[PreparedOrder]) -> bool { + pub fn try_extend(&mut self, orders: &[PreparedOrder], excluded: &FxHashSet) -> bool { self.trace.sim_start_ns = utcnow_ns(); let header = self.ctx.payload.header.clone(); + self.stats.orders_excluded_skipped = + orders.iter().filter(|order| excluded.contains(&order.order_hash)).count() as u64; + let candidates: Vec = (0..orders.len()) .filter(|&ix| { let order = &orders[ix]; - !self.applied_orders.contains(&order.order_id) && + !excluded.contains(&order.order_hash) && + !self.applied_orders.contains(&order.order_id) && order.source_block_hash != self.base_block_hash && simulate::gate_order( order, diff --git a/crates/builder/src/engine/tests.rs b/crates/builder/src/engine/tests.rs index ca3b1c9f7..812c2a870 100644 --- a/crates/builder/src/engine/tests.rs +++ b/crates/builder/src/engine/tests.rs @@ -1,6 +1,3 @@ -//! Merge-engine tests on an in-memory ethrex store: the end-to-end flow -//! (spawned worker), session parking, throttled-emission retry, and stats. - use std::{str::FromStr, sync::Arc, time::Duration}; use alloy_consensus::{SignableTransaction, TxEip1559}; @@ -19,7 +16,7 @@ use helix_tcp_types::merging::{ builder_to_relay::RejectCode, control::{BuilderCollateral, RelayConfigV1}, order::{MergeOrderRef, TxOrderRef}, - relay_to_builder::{MergeableBlockV1, RevokeOrderV1, SlotStartV1}, + relay_to_builder::{MergeableBlockV1, SlotStartV1}, }; use ssz::Encode; use tokio::sync::watch; @@ -33,8 +30,6 @@ use crate::{ node::HeadInfo, }; -/// Dev keys from ethrex's `fixtures/keys/private_keys_l1.txt`; the test picks -/// the ones the `LocalDevnet` genesis actually funds. const KEYS: [&str; 20] = [ "0x941e103320615d394a55708be13e45994c7d93b932b064dbcb2b511fe3254e2e", "0xbcdf20249abf0ed6d944c0288fad489e33f66b3960d9e6229c1cd214ed3bbe31", @@ -87,10 +82,6 @@ fn signed_transfer( alloy_consensus::TxEnvelope::from(tx.into_signed(signature)).encoded_2718() } -/// Shared in-memory ethrex chain plus the merge participants. Signer roles: -/// 0 = winning builder (base coinbase + payment sender), 1 = base user tx, -/// 2 = donor origin coinbase, 3/7 = order senders, 4 = relay fee recipient, -/// 5 = collateral safe (EOA stand-in), 6 = relay signer. struct Fixture { store: Store, blockchain: Arc, @@ -173,7 +164,6 @@ impl Fixture { timestamp: self.genesis_timestamp, is_synced: true, }); - // Keep the sender alive for the test duration. std::mem::forget(tx); rx } @@ -187,16 +177,10 @@ impl Fixture { } } - /// Builds a valid base block ([user tx, proposer payment]) on genesis; - /// vary `user_value` to get distinct block hashes. fn build_base(&self, user_value: U256) -> (MergeableBlockV1, B256) { self.build_base_with_payment(user_value, self.block_value) } - /// Like [`build_base`](Self::build_base), but with an independently - /// chosen payment value — lets a test hold the user tx fixed (so the - /// non-payment prefix is byte-identical) while only the trailing payment - /// tx changes, as in a real bid-ratchet resubmission. fn build_base_with_payment( &self, user_value: U256, @@ -250,9 +234,7 @@ impl Fixture { (msg, block_hash) } - /// A synthetic (never activated) donor block carrying one order: a - /// transfer paying the winning builder's coinbase. - fn donor( + fn mergeable_tx( &self, template: &MergeableBlockV1, order_sender_ix: usize, @@ -286,8 +268,6 @@ impl Fixture { } } - /// A direct-drive engine (no worker thread) for tests that assert on - /// internal state. fn direct_engine( &self, min_emission_interval: Duration, @@ -306,7 +286,6 @@ impl Fixture { (engine, output_rx) } - /// A direct-drive engine with a small per-slot order budget. fn direct_engine_with_order_cap( &self, max_orders_per_slot: usize, @@ -315,6 +294,121 @@ impl Fixture { engine.config.max_orders_per_slot = max_orders_per_slot; (engine, output_rx) } + + fn mergeable_orders( + &self, + template: &MergeableBlockV1, + pubkey: alloy_rpc_types::beacon::BlsPublicKey, + senders: &[usize], + flagged: bool, + hash_byte: u8, + ) -> (MergeableBlockV1, Vec) { + let mut txs = Vec::with_capacity(senders.len()); + let mut order_hashes = Vec::with_capacity(senders.len()); + let mut merge_orders = Vec::with_capacity(senders.len()); + for (i, &sender_ix) in senders.iter().enumerate() { + let tx = signed_transfer( + &self.signers[sender_ix], + self.chain_id, + 0, + self.signers[0].address(), + U256::from(ETH / 5), + 100 * GWEI, + GWEI, + ); + let tx_hash = alloy_primitives::keccak256(&tx); + order_hashes.push(if flagged { + helix_tcp_types::merging::order::bundle_order_hash(&[tx_hash]) + } else { + tx_hash + }); + merge_orders.push(if flagged { + MergeOrderRef::Bundle(helix_tcp_types::merging::order::BundleOrderRef { + txs: vec![i as u16], + reverting_txs: vec![], + dropping_txs: vec![], + latest_only: true, + }) + } else { + MergeOrderRef::Tx(TxOrderRef { index: i as u16, can_revert: false }) + }); + txs.push(tx.into()); + } + + let mut payload = template.execution_payload.clone(); + payload.payload_inner.payload_inner.transactions = txs; + payload.payload_inner.payload_inner.block_hash = B256::repeat_byte(hash_byte); + payload.payload_inner.payload_inner.fee_recipient = self.signers[2].address(); + + let msg = MergeableBlockV1 { + slot: SLOT, + builder_pubkey: pubkey, + block_value: U256::from(ETH / 10), + builder_address: self.signers[2].address(), + proposer_fee_recipient: self.proposer, + parent_beacon_block_root: B256::ZERO, + allow_appending: false, + merge_orders, + execution_payload: payload, + }; + (msg, order_hashes) + } + + fn mergeable_bundle( + &self, + template: &MergeableBlockV1, + pubkey: alloy_rpc_types::beacon::BlsPublicKey, + senders: &[usize], + hash_byte: u8, + ) -> (MergeableBlockV1, B256) { + let mut txs = Vec::with_capacity(senders.len()); + let mut tx_hashes = Vec::with_capacity(senders.len()); + for &sender_ix in senders { + let tx = signed_transfer( + &self.signers[sender_ix], + self.chain_id, + 0, + self.signers[0].address(), + U256::from(ETH / 5), + 100 * GWEI, + GWEI, + ); + tx_hashes.push(alloy_primitives::keccak256(&tx)); + txs.push(tx.into()); + } + let mut payload = template.execution_payload.clone(); + payload.payload_inner.payload_inner.transactions = txs; + payload.payload_inner.payload_inner.block_hash = B256::repeat_byte(hash_byte); + payload.payload_inner.payload_inner.fee_recipient = self.signers[2].address(); + + let msg = MergeableBlockV1 { + slot: SLOT, + builder_pubkey: pubkey, + block_value: U256::from(ETH / 10), + builder_address: self.signers[2].address(), + proposer_fee_recipient: self.proposer, + parent_beacon_block_root: B256::ZERO, + allow_appending: false, + merge_orders: vec![MergeOrderRef::Bundle( + helix_tcp_types::merging::order::BundleOrderRef { + txs: (0..senders.len() as u16).collect(), + reverting_txs: vec![], + dropping_txs: vec![], + latest_only: true, + }, + )], + execution_payload: payload, + }; + let hash = helix_tcp_types::merging::order::bundle_order_hash(&tx_hashes); + (msg, hash) + } + + fn started_engine(&self) -> (MergeEngine, crossbeam_channel::Receiver) { + let (mut engine, output_rx) = self.direct_engine(Duration::ZERO); + engine.handle_event(EngineEvent::RelayConfig(self.relay_config.clone())); + engine.handle_event(EngineEvent::SlotStart(self.slot_start())); + (engine, output_rx) + } } fn mergeable_event(msg: &MergeableBlockV1, recv_ns: u64) -> EngineEvent { @@ -337,12 +431,12 @@ fn expect_merged( } #[tokio::test(flavor = "multi_thread")] -async fn merges_donor_order_into_activated_base_block() { +async fn merges_order_into_activated_base_block() { let _ = tracing_subscriber::fmt().with_env_filter("debug").try_init(); let fixture = Fixture::new().await; let (base_msg, base_block_hash) = fixture.build_base(U256::from(ETH)); let order_value = U256::from(ETH / 5); - let donor_msg = fixture.donor(&base_msg, 3, order_value, 0xdd); + let mergeable_msg = fixture.mergeable_tx(&base_msg, 3, order_value, 0xdd); let base_tx_count = base_msg.execution_payload.payload_inner.payload_inner.transactions.len(); let (event_tx, event_rx) = crossbeam_channel::bounded(1024); @@ -359,7 +453,7 @@ async fn merges_donor_order_into_activated_base_block() { event_tx.send(EngineEvent::RelayConfig(fixture.relay_config.clone())).unwrap(); event_tx.send(EngineEvent::SlotStart(fixture.slot_start())).unwrap(); event_tx.send(mergeable_event(&base_msg, 1)).unwrap(); - event_tx.send(mergeable_event(&donor_msg, 2)).unwrap(); + event_tx.send(mergeable_event(&mergeable_msg, 2)).unwrap(); event_tx.send(activate_event(base_block_hash)).unwrap(); let output = output_rx.recv_timeout(Duration::from_secs(60)).expect("engine produced nothing"); @@ -375,7 +469,6 @@ async fn merges_donor_order_into_activated_base_block() { ); let merged_v1 = &merged.execution_payload.payload_inner.payload_inner; - // base txs + the merged order tx + the distribution tx assert_eq!(merged_v1.transactions.len(), base_tx_count + 2); assert_eq!(merged_v1.block_number, 1); assert_eq!(merged_v1.parent_hash, b256(fixture.genesis_hash)); @@ -393,10 +486,8 @@ async fn merges_donor_order_into_activated_base_block() { ); assert!(merged.appended_blobs.is_empty()); - // No further emission without new orders (nothing improved). assert!(output_rx.recv_timeout(Duration::from_millis(500)).is_err()); - // 2500 bps of distributable revenue -> proposer share is at most contribution/4. let proposer_added = merged.proposer_value - fixture.block_value; assert!(proposer_added > U256::ZERO); assert!(proposer_added <= inclusion.contribution / U256::from(4) + U256::from(1)); @@ -404,13 +495,6 @@ async fn merges_donor_order_into_activated_base_block() { let _ = aaddr(eaddr(fixture.proposer)); // keep converters exercised both ways } -/// A resubmission from the same builder that keeps the same non-payment -/// prefix and only changes the trailing payment tx (a bid ratchet — the -/// common case: the relay's own comments call this "the common case, since -/// builders resubmit near-identical blocks as their bid ratchets") hits the -/// replay checkpoint instead of re-replaying the shared prefix from the -/// parent, and the resulting session is still fully functional (base value -/// correct, still able to merge a fresh order and emit). #[tokio::test(flavor = "multi_thread")] async fn checkpoint_hit_reuses_shared_prefix_on_resubmission() { let fixture = Fixture::new().await; @@ -418,7 +502,6 @@ async fn checkpoint_hit_reuses_shared_prefix_on_resubmission() { let (base_b, hash_b) = fixture .build_base_with_payment(U256::from(ETH), fixture.block_value + U256::from(ETH / 100)); assert_ne!(hash_a, hash_b); - // Same non-payment prefix: the user tx is byte-identical. assert_eq!( base_a.execution_payload.payload_inner.payload_inner.transactions[0], base_b.execution_payload.payload_inner.payload_inner.transactions[0], @@ -449,10 +532,8 @@ async fn checkpoint_hit_reuses_shared_prefix_on_resubmission() { assert_eq!(state.session.as_ref().unwrap().base_block_hash, hash_b); } - // The checkpoint-hit session must still be fully functional: a fresh - // order applies and merges correctly on top of it. - let donor_msg = fixture.donor(&base_b, 3, U256::from(ETH / 5), 0xdd); - engine.handle_event(mergeable_event(&donor_msg, 3)); + let mergeable_msg = fixture.mergeable_tx(&base_b, 3, U256::from(ETH / 5), 0xdd); + engine.handle_event(mergeable_event(&mergeable_msg, 3)); engine.merge_pass(); let merged = @@ -465,272 +546,408 @@ async fn checkpoint_hit_reuses_shared_prefix_on_resubmission() { ); } -fn revoke_event( - order_hash: B256, - builder_pubkey: alloy_rpc_types::beacon::BlsPublicKey, -) -> EngineEvent { - EngineEvent::RevokeOrder { - msg: RevokeOrderV1 { slot: SLOT, order_hash, builder_pubkey }, - generation: 0, - } +fn pubkey(byte: u8) -> alloy_rpc_types::beacon::BlsPublicKey { + alloy_rpc_types::beacon::BlsPublicKey::repeat_byte(byte) } -/// Revoking an order that's pooled but never got applied just drops it from -/// the pool — there's no session state to unwind. #[tokio::test(flavor = "multi_thread")] -async fn revoke_removes_pooled_order_before_it_applies() { +async fn exclusion_holds_across_later_sessions() { let fixture = Fixture::new().await; - let (base_msg, base_block_hash) = fixture.build_base(U256::from(ETH)); - let donor_msg = fixture.donor(&base_msg, 3, U256::from(ETH / 5), 0xdd); - let order_tx_bytes = - donor_msg.execution_payload.payload_inner.payload_inner.transactions[0].clone(); - let order_hash = alloy_primitives::keccak256(order_tx_bytes.as_ref()); + let (base_a, base_a_hash) = fixture.build_base(U256::from(ETH)); + let (base_b, base_b_hash) = fixture.build_base(U256::from(ETH)); + let a = pubkey(0xaa); + let (with_order, hashes) = fixture.mergeable_orders(&base_a, a, &[3], true, 0xd1); + let (without_order, _) = fixture.mergeable_orders(&base_a, a, &[], true, 0xd2); - let (mut engine, output_rx) = fixture.direct_engine(Duration::ZERO); - engine.handle_event(EngineEvent::RelayConfig(fixture.relay_config.clone())); - engine.handle_event(EngineEvent::SlotStart(fixture.slot_start())); - engine.handle_event(mergeable_event(&base_msg, 1)); - engine.handle_event(mergeable_event(&donor_msg, 2)); + let (mut engine, output_rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base_a, 1)); + engine.handle_event(mergeable_event(&base_b, 2)); + engine.handle_event(mergeable_event(&with_order, 3)); + engine.handle_event(mergeable_event(&without_order, 4)); + + engine.handle_event(activate_event(base_a_hash)); + engine.merge_pass(); + while output_rx.try_recv().is_ok() {} - engine.handle_event(revoke_event(order_hash, Default::default())); - assert!(engine.slot.as_ref().unwrap().orders.is_empty(), "revoked order must leave the pool"); + engine.handle_event(activate_event(base_b_hash)); + engine.merge_pass(); - engine.handle_event(activate_event(base_block_hash)); + assert!(engine.slot.as_ref().unwrap().is_excluded(&hashes[0])); + while let Ok(out) = output_rx.try_recv() { + let merged = expect_merged(out); + let order_id = helix_tcp_types::merging::order::order_id(hashes[0], &a); + assert!( + !merged.included_order_ids.contains(&order_id), + "an excluded order must not appear in a later session's block" + ); + } +} + +#[tokio::test(flavor = "multi_thread")] +async fn exclusion_applies_within_the_live_session() { + let fixture = Fixture::new().await; + let (base, base_hash) = fixture.build_base(U256::from(ETH)); + let a = pubkey(0xaa); + let (with_order, hashes) = fixture.mergeable_orders(&base, a, &[3], true, 0xd1); + let (without_order, _) = fixture.mergeable_orders(&base, a, &[], true, 0xd2); + + let (mut engine, output_rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base, 1)); + engine.handle_event(mergeable_event(&with_order, 2)); + engine.handle_event(activate_event(base_hash)); + engine.handle_event(mergeable_event(&without_order, 3)); engine.merge_pass(); + assert!(engine.slot.as_ref().unwrap().is_excluded(&hashes[0])); + let order_id = helix_tcp_types::merging::order::order_id(hashes[0], &a); + while let Ok(out) = output_rx.try_recv() { + assert!(!expect_merged(out).included_order_ids.contains(&order_id)); + } +} + +#[tokio::test(flavor = "multi_thread")] +async fn exclusion_covers_identical_content_from_another_builder() { + let fixture = Fixture::new().await; + let (base, _) = fixture.build_base(U256::from(ETH)); + let (a, b) = (pubkey(0xaa), pubkey(0xbb)); + let (a_with, hashes) = fixture.mergeable_orders(&base, a, &[3], true, 0xd1); + let (b_with, b_hashes) = fixture.mergeable_orders(&base, b, &[3], true, 0xd2); + assert_eq!(hashes[0], b_hashes[0], "both builders must send the same content"); + let (a_without, _) = fixture.mergeable_orders(&base, a, &[], true, 0xd3); + + let (mut engine, _rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base, 1)); + engine.handle_event(mergeable_event(&a_with, 2)); + engine.handle_event(mergeable_event(&b_with, 3)); + engine.handle_event(mergeable_event(&a_without, 4)); + + let state = engine.slot.as_ref().unwrap(); + assert!(state.is_excluded(&hashes[0])); assert!( - output_rx.try_recv().is_err(), - "base alone with its only order revoked before activation must not emit" + state.orders.iter().filter(|o| o.order_hash == hashes[0]).count() >= 2, + "both contributors' entries must still be pooled" ); } -/// Revoking an order the live session already applied can't be reflected in -/// place (no way to un-apply a committed tx), so it forces a full rebuild -/// from the same base — which then excludes the revoked order since it's -/// gone from the pool. #[tokio::test(flavor = "multi_thread")] -async fn revoke_of_an_applied_order_rebuilds_the_session() { +async fn a_later_block_from_another_builder_does_not_restore() { let fixture = Fixture::new().await; - let (base_msg, base_block_hash) = fixture.build_base(U256::from(ETH)); - let donor_msg = fixture.donor(&base_msg, 3, U256::from(ETH / 5), 0xdd); - let order_tx_bytes = - donor_msg.execution_payload.payload_inner.payload_inner.transactions[0].clone(); - let order_hash = alloy_primitives::keccak256(order_tx_bytes.as_ref()); + let (base, _) = fixture.build_base(U256::from(ETH)); + let (a, b) = (pubkey(0xaa), pubkey(0xbb)); + let (a_with, hashes) = fixture.mergeable_orders(&base, a, &[3], true, 0xd1); + let (a_without, _) = fixture.mergeable_orders(&base, a, &[], true, 0xd2); + let (b_with, _) = fixture.mergeable_orders(&base, b, &[3], true, 0xd3); + + let (mut engine, _rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base, 1)); + engine.handle_event(mergeable_event(&a_with, 2)); + engine.handle_event(mergeable_event(&a_without, 3)); + engine.handle_event(mergeable_event(&b_with, 4)); - let (mut engine, output_rx) = fixture.direct_engine(Duration::ZERO); - engine.handle_event(EngineEvent::RelayConfig(fixture.relay_config.clone())); - engine.handle_event(EngineEvent::SlotStart(fixture.slot_start())); - engine.handle_event(mergeable_event(&base_msg, 1)); - engine.handle_event(mergeable_event(&donor_msg, 2)); - engine.handle_event(activate_event(base_block_hash)); - engine.merge_pass(); + assert!( + engine.slot.as_ref().unwrap().is_excluded(&hashes[0]), + "a later block must not lift an exclusion" + ); +} - let merged = expect_merged(output_rx.try_recv().expect("first emission")); - assert_eq!(merged.included_order_ids.len(), 1); +#[tokio::test(flavor = "multi_thread")] +async fn exclusion_applies_to_an_order_that_arrives_later() { + let fixture = Fixture::new().await; + let (base, _) = fixture.build_base(U256::from(ETH)); + let (a, b) = (pubkey(0xaa), pubkey(0xbb)); + let (a_with, hashes) = fixture.mergeable_orders(&base, a, &[3], true, 0xd1); + let (a_without, _) = fixture.mergeable_orders(&base, a, &[], true, 0xd2); + let (b_with, _) = fixture.mergeable_orders(&base, b, &[3], true, 0xd3); + + let (mut engine, _rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base, 1)); + engine.handle_event(mergeable_event(&a_with, 2)); + engine.handle_event(mergeable_event(&a_without, 3)); + engine.handle_event(mergeable_event(&b_with, 4)); - engine.handle_event(revoke_event(order_hash, Default::default())); - engine.merge_pass(); + let state = engine.slot.as_ref().unwrap(); + assert!(state.is_excluded(&hashes[0])); + assert!(state.orders.iter().any(|o| o.order_hash == hashes[0] && o.builder_pubkey == b)); +} + +#[tokio::test(flavor = "multi_thread")] +async fn unexcluded_orders_remain_candidates() { + let fixture = Fixture::new().await; + let (base, _) = fixture.build_base(U256::from(ETH)); + let a = pubkey(0xaa); + let (full, hashes) = fixture.mergeable_orders(&base, a, &[3, 4, 5], true, 0xd1); + let (partial, _) = fixture.mergeable_orders(&base, a, &[3, 5], true, 0xd2); + + let (mut engine, _rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base, 1)); + engine.handle_event(mergeable_event(&full, 2)); + engine.handle_event(mergeable_event(&partial, 3)); let state = engine.slot.as_ref().unwrap(); - assert!(state.orders.is_empty(), "revoked order must leave the pool"); - let session = state.session.as_ref().expect("session rebuilt from the same base"); - assert_eq!(session.base_block_hash, base_block_hash); - assert_eq!( - session.included_order_count(), - 0, - "rebuilt session must not carry the revoked order forward" - ); + assert!(!state.is_excluded(&hashes[0]), "still sent"); + assert!(state.is_excluded(&hashes[1]), "dropped from the newest block"); + assert!(!state.is_excluded(&hashes[2]), "still sent"); } -/// A→B→A activation flip resumes the parked session instead of re-replaying, -/// preserving its emission bookkeeping; stats counters reflect the work. #[tokio::test(flavor = "multi_thread")] -async fn base_flip_back_resumes_parked_session() { +async fn unflagged_content_is_never_excluded_by_absence() { let fixture = Fixture::new().await; - let (base_a, hash_a) = fixture.build_base(U256::from(ETH)); - let (base_b, hash_b) = fixture.build_base(U256::from(2 * ETH)); - assert_ne!(hash_a, hash_b); - let donor_msg = fixture.donor(&base_a, 3, U256::from(ETH / 5), 0xdd); + let (base, _) = fixture.build_base(U256::from(ETH)); + let a = pubkey(0xaa); + let (with_order, hashes) = fixture.mergeable_orders(&base, a, &[3], false, 0xd1); + let (gone, _) = fixture.mergeable_orders(&base, a, &[], false, 0xd2); - let (mut engine, output_rx) = fixture.direct_engine(Duration::ZERO); - engine.handle_event(EngineEvent::RelayConfig(fixture.relay_config.clone())); - engine.handle_event(EngineEvent::SlotStart(fixture.slot_start())); - engine.handle_event(mergeable_event(&base_a, 1)); - engine.handle_event(mergeable_event(&base_b, 2)); - engine.handle_event(mergeable_event(&donor_msg, 3)); + let (mut engine, _rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base, 1)); + engine.handle_event(mergeable_event(&with_order, 2)); + engine.handle_event(mergeable_event(&gone, 3)); - // Activate A: fresh session, order applied, emission. - engine.handle_event(activate_event(hash_a)); + assert!(!engine.slot.as_ref().unwrap().is_excluded(&hashes[0])); +} + +#[tokio::test(flavor = "multi_thread")] +async fn silence_is_not_exclusion() { + let fixture = Fixture::new().await; + let (base, base_hash) = fixture.build_base(U256::from(ETH)); + let a = pubkey(0xaa); + let (with_order, hashes) = fixture.mergeable_orders(&base, a, &[3], true, 0xd1); + + let (mut engine, _rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base, 1)); + engine.handle_event(mergeable_event(&with_order, 2)); + engine.handle_event(activate_event(base_hash)); + engine.merge_pass(); engine.merge_pass(); - let merged_a = expect_merged(output_rx.try_recv().expect("no emission for base A")); - assert_eq!(merged_a.base_block_hash, hash_a); - { - let state = engine.slot.as_ref().unwrap(); - let session = state.session.as_ref().unwrap(); - assert_eq!(session.base_block_hash, hash_a); - assert!(state.parked.is_empty()); - assert_eq!(session.stats().orders_applied, 1); - assert_eq!(session.stats().emissions, 1); - assert!(session.stats().candidates_screened >= 1); - } + assert!(!engine.slot.as_ref().unwrap().is_excluded(&hashes[0])); +} - // Activate B: A is parked, B builds fresh and emits its own merge. - engine.handle_event(activate_event(hash_b)); - engine.merge_pass(); - let merged_b = expect_merged(output_rx.try_recv().expect("no emission for base B")); - assert_eq!(merged_b.base_block_hash, hash_b); - { - let state = engine.slot.as_ref().unwrap(); - assert_eq!(state.session.as_ref().unwrap().base_block_hash, hash_b); - assert_eq!(state.parked.len(), 1); - assert_eq!(state.parked[0].base_block_hash, hash_a); - } +#[tokio::test(flavor = "multi_thread")] +async fn exclusion_does_not_change_pool_membership() { + let fixture = Fixture::new().await; + let (base, _) = fixture.build_base(U256::from(ETH)); + let a = pubkey(0xaa); + let (with_order, hashes) = fixture.mergeable_orders(&base, a, &[3], true, 0xd1); + let (without_order, _) = fixture.mergeable_orders(&base, a, &[], true, 0xd2); + + let (mut engine, _rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base, 1)); + engine.handle_event(mergeable_event(&with_order, 2)); + let pooled_before = engine.slot.as_ref().unwrap().orders.len(); + engine.handle_event(mergeable_event(&without_order, 3)); + let state = engine.slot.as_ref().unwrap(); - // Flip back to A: the parked session resumes with its bookkeeping intact - // (a rebuilt session would have best_emitted == 0 and would re-emit). - engine.handle_event(activate_event(hash_a)); - engine.merge_pass(); - { - let state = engine.slot.as_ref().unwrap(); - let session = state.session.as_ref().unwrap(); - assert_eq!(session.base_block_hash, hash_a); - assert_eq!(state.parked.len(), 1); - assert_eq!(state.parked[0].base_block_hash, hash_b); - assert_eq!(session.stats().emissions, 1, "resumed session kept its stats"); - } - // Nothing new to emit for the resumed session: same orders, same value. - assert!(output_rx.try_recv().is_err(), "resume must not re-emit a non-improving block"); + assert!(state.is_excluded(&hashes[0])); + assert_eq!(state.orders.len(), pooled_before, "nothing leaves the pool mid-slot"); + assert!(state.orders.iter().any(|o| o.order_hash == hashes[0])); } -/// A resubmission that dehydrates an already-seen order tx to a -/// `order::TX_HASH_REF_LEN`-byte hash reference (the relay's connection-scoped -/// dehydration, see `MergeableBlockV1`'s doc comment) must resolve against -/// the tx sent whole earlier this slot, not be rejected as an undecodable tx. #[tokio::test(flavor = "multi_thread")] -async fn resolves_dehydrated_tx_hash_reference() { +async fn exclusions_do_not_cross_a_slot() { let fixture = Fixture::new().await; - let (base_msg, _base_block_hash) = fixture.build_base(U256::from(ETH)); - let donor_full = fixture.donor(&base_msg, 3, U256::from(ETH / 5), 0xdd); - - let order_tx_bytes = - donor_full.execution_payload.payload_inner.payload_inner.transactions[0].clone(); - let order_tx_hash = alloy_primitives::keccak256(order_tx_bytes.as_ref()); - let dehydrated_block_hash = B256::repeat_byte(0xee); - let mut donor_dehydrated = donor_full.clone(); - donor_dehydrated.execution_payload.payload_inner.payload_inner.transactions = - vec![order_tx_hash.as_slice().to_vec().into()]; - donor_dehydrated.execution_payload.payload_inner.payload_inner.block_hash = - dehydrated_block_hash; + let (base, _) = fixture.build_base(U256::from(ETH)); + let a = pubkey(0xaa); + let (with_order, hashes) = fixture.mergeable_orders(&base, a, &[3], true, 0xd1); + let (without_order, _) = fixture.mergeable_orders(&base, a, &[], true, 0xd2); - let (mut engine, output_rx) = fixture.direct_engine(Duration::ZERO); - engine.handle_event(EngineEvent::RelayConfig(fixture.relay_config.clone())); - engine.handle_event(EngineEvent::SlotStart(fixture.slot_start())); - engine.handle_event(mergeable_event(&donor_full, 1)); - engine.handle_event(mergeable_event(&donor_dehydrated, 2)); + let (mut engine, _rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base, 1)); + engine.handle_event(mergeable_event(&with_order, 2)); + engine.handle_event(mergeable_event(&without_order, 3)); + assert!(engine.slot.as_ref().unwrap().is_excluded(&hashes[0])); + + let mut next = fixture.slot_start(); + next.slot = SLOT + 1; + engine.handle_event(EngineEvent::SlotStart(next)); let state = engine.slot.as_ref().unwrap(); - assert!( - state.blocks.contains_key(&dehydrated_block_hash), - "dehydrated resubmission with a resolvable tx-hash reference must be pooled, not rejected" - ); - assert!(output_rx.try_recv().is_err(), "no reject should have been emitted"); + assert!(!state.is_excluded(&hashes[0]), "exclusions must not cross a slot"); + assert!(state.latest_only.is_empty(), "latest_only state must not cross a slot either"); } -/// An improvement blocked by the emission-spacing gate is retried by the -/// worker after the window passes, with no further inbound events. #[tokio::test(flavor = "multi_thread")] -async fn throttled_emission_is_retried() { +async fn the_first_block_of_a_slot_excludes_nothing() { let fixture = Fixture::new().await; - let (base_msg, base_block_hash) = fixture.build_base(U256::from(ETH)); - let donor_one = fixture.donor(&base_msg, 3, U256::from(ETH / 5), 0xdd); - let donor_two = fixture.donor(&base_msg, 7, U256::from(ETH / 4), 0xee); + let (base, _) = fixture.build_base(U256::from(ETH)); + let a = pubkey(0xaa); + let (with_order, hashes) = fixture.mergeable_orders(&base, a, &[3, 4], true, 0xd1); - let (event_tx, event_rx) = crossbeam_channel::bounded(1024); - let (output_tx, output_rx) = crossbeam_channel::bounded(64); - let _engine = MergeEngine::spawn( - fixture.engine_config(Duration::from_millis(300)), - fixture.store.clone(), - fixture.blockchain.clone(), - fixture.head(), - event_rx, - output_tx, - ); + let (mut engine, _rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base, 1)); + engine.handle_event(mergeable_event(&with_order, 2)); - event_tx.send(EngineEvent::RelayConfig(fixture.relay_config.clone())).unwrap(); - event_tx.send(EngineEvent::SlotStart(fixture.slot_start())).unwrap(); - event_tx.send(mergeable_event(&base_msg, 1)).unwrap(); - event_tx.send(mergeable_event(&donor_one, 2)).unwrap(); - event_tx.send(activate_event(base_block_hash)).unwrap(); + let mut next = fixture.slot_start(); + next.slot = SLOT + 1; + engine.handle_event(EngineEvent::SlotStart(next)); - let first = - expect_merged(output_rx.recv_timeout(Duration::from_secs(60)).expect("no first emission")); + let mut base2 = base.clone(); + base2.slot = SLOT + 1; + let (mut first, _) = fixture.mergeable_orders(&base2, a, &[3], true, 0xe1); + first.slot = SLOT + 1; + engine.handle_event(mergeable_event(&base2, 3)); + engine.handle_event(mergeable_event(&first, 4)); - // The second improvement lands right after the first emission and hits the - // 300ms spacing gate; only the timeout wake-up can emit it. - event_tx.send(mergeable_event(&donor_two, 3)).unwrap(); - let second = expect_merged( - output_rx.recv_timeout(Duration::from_secs(10)).expect("throttled emission never retried"), - ); - assert!( - second.proposer_value > first.proposer_value, - "retried emission {} must improve on {}", - second.proposer_value, - first.proposer_value + let state = engine.slot.as_ref().unwrap(); + assert!(!state.is_excluded(&hashes[0])); + assert!(!state.is_excluded(&hashes[1]), "the previous slot must not exclude here"); +} + +#[tokio::test(flavor = "multi_thread")] +async fn a_block_for_another_slot_alters_nothing() { + let fixture = Fixture::new().await; + let (base, _) = fixture.build_base(U256::from(ETH)); + let a = pubkey(0xaa); + let (with_order, _) = fixture.mergeable_orders(&base, a, &[3], true, 0xd1); + let (mut stale, _) = fixture.mergeable_orders(&base, a, &[], true, 0xd2); + stale.slot = SLOT + 5; + + let (mut engine, _rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base, 1)); + engine.handle_event(mergeable_event(&with_order, 2)); + let flagged_before = engine.slot.as_ref().unwrap().latest_only.clone(); + engine.handle_event(mergeable_event(&stale, 3)); + + assert_eq!( + engine.slot.as_ref().unwrap().latest_only, + flagged_before, + "a block for another slot must not alter this slot's latest_only state" ); - assert_eq!(second.included_order_ids.len(), 2); } -/// A repeat announcement of a pooled order must not consume the per-slot budget. #[tokio::test(flavor = "multi_thread")] -async fn repeat_announcements_do_not_consume_the_order_budget() { +async fn exclusion_does_not_disturb_an_applied_bundle() { let fixture = Fixture::new().await; - let (base_msg, _base_block_hash) = fixture.build_base(U256::from(ETH)); + let (base, base_hash) = fixture.build_base(U256::from(ETH)); + let a = pubkey(0xaa); + let (with_order, order_hash) = fixture.mergeable_bundle(&base, a, &[3, 4], 0xd1); + let (without_order, _) = fixture.mergeable_orders(&base, a, &[], true, 0xd2); + let bundle_txs: Vec<_> = + with_order.execution_payload.payload_inner.payload_inner.transactions.clone(); + + let (mut engine, output_rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base, 1)); + engine.handle_event(mergeable_event(&with_order, 2)); + engine.handle_event(activate_event(base_hash)); + engine.merge_pass(); - // One order, re-announced in six blocks with distinct block hashes. - let donors: Vec = - (0u8..6).map(|i| fixture.donor(&base_msg, 3, U256::from(ETH / 5), 0xd0 + i)).collect(); + let merged = expect_merged(output_rx.try_recv().expect("bundle should merge while still sent")); + let order_id = helix_tcp_types::merging::order::order_id(order_hash, &a); + assert!(merged.included_order_ids.contains(&order_id), "bundle applied before exclusion"); + let applied_before = merged.execution_payload.payload_inner.payload_inner.transactions.clone(); - let (mut engine, output_rx) = fixture.direct_engine_with_order_cap(3); - engine.handle_event(EngineEvent::RelayConfig(fixture.relay_config.clone())); - engine.handle_event(EngineEvent::SlotStart(fixture.slot_start())); - for (i, donor) in donors.iter().enumerate() { - engine.handle_event(mergeable_event(donor, i as u64 + 1)); - } + engine.handle_event(mergeable_event(&without_order, 3)); - assert!(output_rx.try_recv().is_err(), "a repeat order must not trip the per-slot cap"); - let state = engine.slot.as_ref().unwrap(); - assert_eq!(state.blocks.len(), donors.len(), "every resubmission must be pooled"); - assert_eq!(state.orders.len(), 1, "the six announcements are one distinct order"); + let session = engine.slot.as_ref().unwrap().session.as_ref().expect("session still live"); + assert!(session.has_applied(&order_id), "an applied order stays applied"); + let positions: Vec = bundle_txs + .iter() + .map(|tx| applied_before.iter().position(|t| t == tx).expect("bundle tx present")) + .collect(); + assert!( + positions.windows(2).all(|w| w[1] == w[0] + 1), + "an applied bundle must stay contiguous and ordered" + ); } -/// The per-slot budget still binds on distinct orders. #[tokio::test(flavor = "multi_thread")] -async fn distinct_orders_still_hit_the_per_slot_cap() { +async fn emitted_block_stays_valid_across_an_exclusion() { let fixture = Fixture::new().await; - let (base_msg, _base_block_hash) = fixture.build_base(U256::from(ETH)); + let (base, base_hash) = fixture.build_base(U256::from(ETH)); + let a = pubkey(0xaa); + let (with_order, _) = fixture.mergeable_orders(&base, a, &[3, 4], true, 0xd1); + let (partial, _) = fixture.mergeable_orders(&base, a, &[3], true, 0xd2); + + let (mut engine, output_rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base, 1)); + engine.handle_event(mergeable_event(&with_order, 2)); + engine.handle_event(activate_event(base_hash)); + engine.merge_pass(); + while output_rx.try_recv().is_ok() {} - // Four distinct orders against a budget of three. - let donors: Vec = (0u8..4) - .map(|i| fixture.donor(&base_msg, 3, U256::from(ETH / 5 + u128::from(i)), 0xd0 + i)) - .collect(); + engine.handle_event(mergeable_event(&partial, 3)); + engine.merge_pass(); - let (mut engine, output_rx) = fixture.direct_engine_with_order_cap(3); - engine.handle_event(EngineEvent::RelayConfig(fixture.relay_config.clone())); - engine.handle_event(EngineEvent::SlotStart(fixture.slot_start())); - for (i, donor) in donors.iter().enumerate() { - engine.handle_event(mergeable_event(donor, i as u64 + 1)); + while let Ok(out) = output_rx.try_recv() { + let merged = expect_merged(out); + assert!( + merged.proposer_value > U256::ZERO, + "an emission after an exclusion must still pay the proposer" + ); + assert!( + !merged.execution_payload.payload_inner.payload_inner.transactions.is_empty(), + "an emission after an exclusion must still carry a payload" + ); } +} + +#[tokio::test(flavor = "multi_thread")] +async fn exclusion_state_is_independent_of_event_interleaving() { + use rand::{RngCore, SeedableRng}; + use rand_xorshift::XorShiftRng; - match output_rx.try_recv().expect("the fourth distinct order must be rejected") { - EngineOutput::Reject { msg, .. } => assert_eq!(msg.code, RejectCode::LimitExceeded), - EngineOutput::Merged { .. } => panic!("expected a reject, got a merged block"), + let fixture = Fixture::new().await; + let (base, base_hash) = fixture.build_base(U256::from(ETH)); + let (a, b) = (pubkey(0xaa), pubkey(0xbb)); + let (a_full, a_hashes) = fixture.mergeable_orders(&base, a, &[3, 4, 5], true, 0xd1); + let (a_partial, _) = fixture.mergeable_orders(&base, a, &[3, 5], true, 0xd2); + let (b_full, _) = fixture.mergeable_orders(&base, b, &[4], true, 0xd3); + let (a_final, _) = fixture.mergeable_orders(&base, a, &[3], true, 0xd4); + let blocks = [&base, &a_full, &a_partial, &b_full, &a_final]; + + let mut expected: Option> = None; + for seed in 0u64..8 { + let mut rng = XorShiftRng::seed_from_u64(seed); + let (mut engine, output_rx) = fixture.started_engine(); + for msg in blocks { + engine.handle_event(mergeable_event(msg, 1)); + match rng.next_u32() % 3 { + 0 => { + engine.handle_event(activate_event(base_hash)); + } + 1 => engine.merge_pass(), + _ => {} + } + } + while output_rx.try_recv().is_ok() {} + + let mut got: Vec = engine.slot.as_ref().unwrap().excluded.iter().copied().collect(); + got.sort(); + match &expected { + None => expected = Some(got), + Some(first) => { + assert_eq!(&got, first, "interleaving changed the excluded set (seed {seed})") + } + } } + + let excluded = expected.expect("at least one run"); + assert!(excluded.contains(&a_hashes[1]), "A dropped its second order"); + assert!(excluded.contains(&a_hashes[2]), "A dropped its third order"); + assert!(!excluded.contains(&a_hashes[0]), "A never dropped its first order"); +} + +#[tokio::test(flavor = "multi_thread")] +async fn every_pooled_order_is_accounted_for() { + let fixture = Fixture::new().await; + let (base, base_hash) = fixture.build_base(U256::from(ETH)); + let a = pubkey(0xaa); + let (full, _) = fixture.mergeable_orders(&base, a, &[3, 4, 5], true, 0xd1); + let (partial, _) = fixture.mergeable_orders(&base, a, &[3], true, 0xd2); + + let (mut engine, _rx) = fixture.started_engine(); + engine.handle_event(mergeable_event(&base, 1)); + engine.handle_event(mergeable_event(&full, 2)); + engine.handle_event(mergeable_event(&partial, 3)); + engine.handle_event(activate_event(base_hash)); + engine.merge_pass(); + let state = engine.slot.as_ref().unwrap(); - assert_eq!(state.orders.len(), 3, "the pool must stop at the cap"); + let excluded_in_pool = + state.orders.iter().filter(|o| state.is_excluded(&o.order_hash)).count() as u64; + let session = state.session.as_ref().expect("session live"); assert_eq!( - state.blocks.len(), - donors.len(), - "a block whose orders hit the cap is still stored as a base candidate" + session.stats().orders_excluded_skipped, + excluded_in_pool, + "the skip counter must account for exactly the excluded pool entries" ); } diff --git a/crates/builder/src/engine/types.rs b/crates/builder/src/engine/types.rs index ff05a3058..996aa5e83 100644 --- a/crates/builder/src/engine/types.rs +++ b/crates/builder/src/engine/types.rs @@ -4,7 +4,7 @@ use alloy_primitives::{Address, B256, U256}; use alloy_rpc_types::beacon::BlsPublicKey; use alloy_signer_local::PrivateKeySigner; use helix_tcp_types::merging::control::RelayConfigV1; -use rustc_hash::FxHashMap; +use rustc_hash::{FxHashMap, FxHashSet}; use crate::engine::session::{MergeSession, ReplayCheckpoint}; @@ -67,6 +67,8 @@ pub struct PreparedBlock { /// One mergeable order drawn from a prepared block's `merge_orders`. pub struct PreparedOrder { pub order_id: B256, + pub order_hash: B256, + pub latest_only: bool, pub origin: Address, pub builder_pubkey: BlsPublicKey, pub source_block_hash: B256, @@ -139,25 +141,22 @@ pub struct SlotState { /// scratch, this slot. pub checkpoint_hits: usize, pub checkpoint_misses: usize, + pub excluded: FxHashSet, + pub latest_only: FxHashMap>, } impl SlotState { - /// Removes `order_id` from the pool if `builder_pubkey` still holds its - /// current attribution — a higher-value duplicate from another builder - /// may have since superseded it, in which case it's left untouched (see - /// `ingest_mergeable_block`'s dedup). Returns whether anything changed. - pub fn remove_order(&mut self, order_id: B256, builder_pubkey: &BlsPublicKey) -> bool { - let Some(&ix) = self.order_ids.get(&order_id) else { return false }; - if &self.orders[ix].builder_pubkey != builder_pubkey { - return false; + pub fn update_latest_only(&mut self, pubkey: BlsPublicKey, current: FxHashSet) { + if let Some(previous) = self.latest_only.get(&pubkey) { + for dropped in previous.difference(¤t) { + self.excluded.insert(*dropped); + } } - self.order_ids.remove(&order_id); - self.orders.swap_remove(ix); - // The element `swap_remove` moved into `ix` needs its index updated. - if let Some(moved) = self.orders.get(ix) { - self.order_ids.insert(moved.order_id, ix); - } - true + self.latest_only.insert(pubkey, current); + } + + pub fn is_excluded(&self, order_hash: &B256) -> bool { + self.excluded.contains(order_hash) } pub fn new(msg: &helix_tcp_types::merging::relay_to_builder::SlotStartV1) -> Self { @@ -178,6 +177,8 @@ impl SlotState { replay_checkpoints: FxHashMap::default(), checkpoint_hits: 0, checkpoint_misses: 0, + excluded: FxHashSet::default(), + latest_only: FxHashMap::default(), } } } diff --git a/crates/builder/src/server/mod.rs b/crates/builder/src/server/mod.rs index b37825aa1..8c6da091f 100644 --- a/crates/builder/src/server/mod.rs +++ b/crates/builder/src/server/mod.rs @@ -23,7 +23,7 @@ use helix_tcp_types::{ MERGING_HEADER_SIZE, MergingFrameHeader, MergingHeaderError, MergingMsgId, builder_to_relay::{FatalV1, RejectCode, RejectSubject, RejectV1}, control::{MergerAckV1, MergerRegistrationV1, PingV1, PongV1, RelayConfigV1}, - relay_to_builder::{ActivateBaseBlockV1, RevokeOrderV1, SlotEndV1, SlotStartV1}, + relay_to_builder::{ActivateBaseBlockV1, SlotEndV1, SlotStartV1}, }, }; use rustc_hash::FxHashMap; @@ -529,14 +529,6 @@ fn handle_active_message( generation, }); } - MergingMsgId::RevokeOrderV1 => { - let Ok(msg) = RevokeOrderV1::from_ssz_bytes(body) else { - fatal(replies, to_disconnect, RejectCode::InvalidOrder, "undecodable revoke order"); - return; - }; - debug!(slot = msg.slot, order_hash = %msg.order_hash, "revoke order"); - send_control(engine_tx, EngineEvent::RevokeOrder { msg, generation }); - } MergingMsgId::MergeableBlockV1 => { // Not decoded here: the engine does SSZ + tx decoding off the tile // thread. Only the slot for a possible Busy reject is peeked. diff --git a/crates/relay/src/block_merging/tile.rs b/crates/relay/src/block_merging/tile.rs index 6e78df88a..34cae6710 100644 --- a/crates/relay/src/block_merging/tile.rs +++ b/crates/relay/src/block_merging/tile.rs @@ -3,7 +3,7 @@ use std::sync::{ atomic::{AtomicBool, Ordering}, }; -use alloy_primitives::{B256, Bytes, U256, keccak256}; +use alloy_primitives::{B256, Bytes, keccak256}; use flux::{ spine::SpineProducers, tile::Tile, @@ -28,10 +28,10 @@ use helix_tcp_types::merging::{ control::{ BuilderCollateral, MergerAckV1, MergerRegistrationV1, PingV1, PongV1, RelayConfigV1, }, - order::{MergeOrderRef, order_id}, - relay_to_builder::{ActivateBaseBlockV1, MergeableBlockV1, RevokeOrderV1, SlotStartV1}, + order::MergeOrderRef, + relay_to_builder::{ActivateBaseBlockV1, MergeableBlockV1, SlotStartV1}, }; -use helix_types::{BlobWithMetadata, BlsPublicKeyBytes, HydrationCache, Submission, payload_to_v3}; +use helix_types::{BlobWithMetadata, HydrationCache, Submission, payload_to_v3}; use rustc_hash::{FxHashMap, FxHashSet}; use ssz::Decode; use tracing::{debug, error, info, trace, warn}; @@ -120,27 +120,14 @@ struct SlotState { appendable: FxHashSet, /// Events replayed on re-handshake. replay_log: Vec, - /// base_block_hash -> best proposer_value and the order_ids it was built from. - best_merged: FxHashMap, /// Orders actually sent to the merge builder this slot, for the /// unbundling check on incoming merged blocks. order_txs: Vec, - /// Per-builder map of order_id -> order_hash for their latest_only bundles this - /// slot, diffed on each new submission to detect revocations. - latest_only_ids: FxHashMap>, } #[derive(Clone, Copy)] enum ReplayEvent { Forward(usize), - Revoke { order_hash: B256, builder_pubkey: BlsPublicKeyBytes }, -} - -#[derive(Default)] -struct BestMergedFloor { - value: U256, - /// order_ids of the merged block currently holding this floor. - order_ids: FxHashSet, } /// Per-slot counters, logged and reset on slot transition. @@ -160,12 +147,12 @@ struct SlotStats { skipped_over_limits: usize, /// Merge orders dropped for an out of range tx index. orders_dropped: usize, + orders_forwarded: usize, + orders_forwarded_latest_only: usize, /// Mergeable frames replayed on re-handshake. replayed: usize, merged_blocks: usize, merged_stale: usize, - /// Merged blocks discarded because a better one was already stored. - merged_regressed: usize, /// Merged blocks dropped because an appended blob's sidecar wasn't in our cache. merged_blob_missing: usize, /// Merged blocks whose simulation was skipped because a required piece of this @@ -308,20 +295,6 @@ fn handle_merged_block( ); return None; } - // Builders only guarantee monotonicity within a connection; filter - // so the stored merged bid never regresses. - if slot - .best_merged - .get(&merged.base_block_hash) - .is_some_and(|floor| merged.proposer_value <= floor.value) - { - stats.merged_regressed += 1; - return None; - } - slot.best_merged.insert(merged.base_block_hash, BestMergedFloor { - value: merged.proposer_value, - order_ids: merged.included_order_ids.iter().copied().collect(), - }); let Some(response) = merged_block_to_response(merged, blob_sidecars, max_blobs_per_block) else { stats.merged_blob_missing += 1; @@ -423,30 +396,6 @@ fn merge_sim_disable_check( Some((block_hash, err.clone())) } -/// order_id -> order_hash for every `latest_only` bundle in this submission. -fn latest_only_ids( - builder_pubkey: BlsPublicKeyBytes, - merge_orders: &[MergeOrderRef], - order_hashes: &[B256], -) -> FxHashMap { - merge_orders - .iter() - .zip(order_hashes) - .filter(|(order_ref, _)| matches!(order_ref, MergeOrderRef::Bundle(b) if b.latest_only)) - .map(|(_, &hash)| (order_id(hash, &builder_pubkey), hash)) - .collect() -} - -/// (order_id, order_hash) pairs present in `prev` but missing from `new` — -/// i.e. flagged latest_only before, dropped from the newest submission. -fn revoked_ids( - prev: Option<&FxHashMap>, - new: &FxHashMap, -) -> Vec<(B256, B256)> { - let Some(prev) = prev else { return Vec::new() }; - prev.iter().filter(|(id, _)| !new.contains_key(*id)).map(|(&id, &hash)| (id, hash)).collect() -} - impl BlockMergingTile { #[allow(clippy::too_many_arguments)] pub fn new( @@ -576,13 +525,6 @@ impl BlockMergingTile { for event in self.slot.replay_log.clone() { match event { ReplayEvent::Forward(ix) => self.forward_decoded(ix, Some(token)), - ReplayEvent::Revoke { order_hash, builder_pubkey } => { - let msg = - RevokeOrderV1 { slot: self.slot.bid_slot, order_hash, builder_pubkey }; - self.connector.write_or_enqueue_with(SendBehavior::Single(token), |buf| { - append_frame(buf, MergingMsgId::RevokeOrderV1, &msg); - }); - } } } } @@ -880,10 +822,11 @@ impl BlockMergingTile { skipped_builder_ahead = stats.skipped_builder_ahead, skipped_over_limits = stats.skipped_over_limits, orders_dropped = stats.orders_dropped, + orders_forwarded = stats.orders_forwarded, + orders_forwarded_latest_only = stats.orders_forwarded_latest_only, replayed = stats.replayed, merged_blocks = stats.merged_blocks, merged_stale = stats.merged_stale, - merged_regressed = stats.merged_regressed, merged_blob_missing = stats.merged_blob_missing, merged_slot_data_missing = stats.merged_slot_data_missing, merged_unbundled = stats.merged_unbundled, @@ -942,10 +885,6 @@ impl BlockMergingTile { } let Some(merging) = &sub.merging_data else { self.stats.skipped_no_merging_data += 1; - // No merging data revokes everything this builder previously flagged. - if !is_replay { - self.diff_latest_only(*sub.submission.builder_pubkey(), &[], &[]); - } self.feed_cache(&sub.submission); return; }; @@ -1053,10 +992,6 @@ impl BlockMergingTile { .map(|order_ref| order_ref_hash(order_ref, &tx_hashes)) .collect(); - if !is_replay { - self.diff_latest_only(signed.message.builder_pubkey, &msg.merge_orders, &order_hashes); - } - self.encode_buf.clear(); append_frame(&mut self.encode_buf, MergingMsgId::MergeableBlockV1, &msg); @@ -1083,6 +1018,14 @@ impl BlockMergingTile { debug!(?token, %block_hash, "skipping mergeable block over builder limits"); return; } + if !is_replay { + self.stats.orders_forwarded += msg.merge_orders.len(); + self.stats.orders_forwarded_latest_only += msg + .merge_orders + .iter() + .filter(|o| matches!(o, MergeOrderRef::Bundle(b) if b.latest_only)) + .count(); + } self.conn.orders_sent.extend(order_hashes); self.slot .order_txs @@ -1095,45 +1038,6 @@ impl BlockMergingTile { }); } - /// Diffs `builder_pubkey`'s latest_only bundles against their previous - /// submission this slot to detect revocations. - fn diff_latest_only( - &mut self, - builder_pubkey: BlsPublicKeyBytes, - merge_orders: &[MergeOrderRef], - order_hashes: &[B256], - ) { - let new_ids = latest_only_ids(builder_pubkey, merge_orders, order_hashes); - let prev_ids = self.slot.latest_only_ids.get(&builder_pubkey); - for (revoked_id, revoked_hash) in revoked_ids(prev_ids, &new_ids) { - self.revoke_order(revoked_id, revoked_hash, builder_pubkey); - } - self.slot.latest_only_ids.insert(builder_pubkey, new_ids); - } - - /// Relaxes the floor for any base block that depended on `order_id`, notifies - /// the auctioneer to evict cached bids, and tells the merge builder to drop it. - fn revoke_order( - &mut self, - order_id: B256, - order_hash: B256, - builder_pubkey: BlsPublicKeyBytes, - ) { - self.slot.best_merged.retain(|_, floor| !floor.order_ids.contains(&order_id)); - - let bid_slot = self.slot.bid_slot; - - self.slot.replay_log.push(ReplayEvent::Revoke { order_hash, builder_pubkey }); - if let Some(token) = self.token && - self.conn.active - { - let msg = RevokeOrderV1 { slot: bid_slot, order_hash, builder_pubkey }; - self.connector.write_or_enqueue_with(SendBehavior::Single(token), |buf| { - append_frame(buf, MergingMsgId::RevokeOrderV1, &msg); - }); - } - } - fn on_top_bid(&mut self, top_bid: TopBidUpdate) { if top_bid.slot != self.slot.bid_slot { return; @@ -1242,32 +1146,18 @@ mod tests { MergingBuilderCollateral, MergingBuilderEndpoint, SubmissionTrace, decoder::{Encoding, SubmissionDecoderParams}, }; - use helix_tcp_types::{ - MergeType, - merging::{ - builder_to_relay::MergeTraceV1, - order::{BundleOrderRef, TxOrderRef}, - }, - }; + use helix_tcp_types::{MergeType, merging::builder_to_relay::MergeTraceV1}; use helix_types::{ - BlobsBundle, BlockMergingData, BuilderInclusionResult, Compression, ExecutionPayload, - ExecutionRequests, ForkName, MergedBlockTrace, SignedBidSubmission, SubmissionVersion, - TestRandom, TestRandomSeed, dehydrated_submission_with_txs_for_test, full_tx_for_test, + BlobsBundle, BlockMergingData, BlsPublicKeyBytes, BuilderInclusionResult, Compression, + ExecutionPayload, ExecutionRequests, ForkName, MergedBlockTrace, SignedBidSubmission, + SubmissionVersion, TestRandom, TestRandomSeed, dehydrated_submission_with_txs_for_test, + full_tx_for_test, }; use rand::{SeedableRng, rngs::SmallRng}; use super::*; use crate::{SubmissionRef, auctioneer::SubmissionData}; - fn bundle(latest_only: bool) -> MergeOrderRef { - MergeOrderRef::Bundle(BundleOrderRef { - txs: vec![0], - reverting_txs: vec![], - dropping_txs: vec![], - latest_only, - }) - } - fn merge_response( payload: ExecutionPayload, proposer_value: U256, @@ -1291,67 +1181,6 @@ mod tests { } } - #[test] - fn latest_only_ids_ignores_non_flagged_and_tx_orders() { - let builder_pubkey = BlsPublicKeyBytes::default(); - let orders = [ - bundle(true), - bundle(false), - MergeOrderRef::Tx(TxOrderRef { index: 0, can_revert: false }), - ]; - let hashes = [B256::repeat_byte(1), B256::repeat_byte(2), B256::repeat_byte(3)]; - - let ids = latest_only_ids(builder_pubkey, &orders, &hashes); - - assert_eq!(ids.len(), 1); - assert_eq!(*ids.values().next().unwrap(), hashes[0]); - } - - #[test] - fn latest_only_ids_disambiguates_by_builder() { - let hash = [B256::repeat_byte(9)]; - let orders = [bundle(true)]; - let mut pubkey_b = [0u8; 48]; - pubkey_b[0] = 1; - - let ids_a = latest_only_ids(BlsPublicKeyBytes::default(), &orders, &hash); - let ids_b = latest_only_ids(BlsPublicKeyBytes::from(pubkey_b), &orders, &hash); - - assert_ne!(ids_a.keys().next(), ids_b.keys().next()); - } - - #[test] - fn revoked_ids_detects_a_dropped_flag() { - let builder_pubkey = BlsPublicKeyBytes::default(); - let hash = B256::repeat_byte(5); - let prev = latest_only_ids(builder_pubkey, &[bundle(true)], &[hash]); - - // Resubmission without the bundle at all. - let new = latest_only_ids(builder_pubkey, &[], &[]); - - let revoked = revoked_ids(Some(&prev), &new); - assert_eq!(revoked, vec![(*prev.keys().next().unwrap(), hash)]); - } - - #[test] - fn revoked_ids_empty_when_still_present() { - let builder_pubkey = BlsPublicKeyBytes::default(); - let hash = [B256::repeat_byte(6)]; - let orders = [bundle(true)]; - - let prev = latest_only_ids(builder_pubkey, &orders, &hash); - let new = latest_only_ids(builder_pubkey, &orders, &hash); - - assert!(revoked_ids(Some(&prev), &new).is_empty()); - } - - #[test] - fn revoked_ids_empty_on_first_submission() { - let new = - latest_only_ids(BlsPublicKeyBytes::default(), &[bundle(true)], &[B256::repeat_byte(7)]); - assert!(revoked_ids(None, &new).is_empty()); - } - fn inclusion(txs: Vec) -> BuilderInclusionResult { BuilderInclusionResult { contribution: U256::ZERO, revenue: U256::ZERO, txs } } @@ -1887,7 +1716,6 @@ mod tests { ); assert!(result.is_none()); - assert!(slot.best_merged.is_empty()); assert_eq!(stats.merged_blocks, 0); } @@ -1916,7 +1744,6 @@ mod tests { ); assert!(result.is_some()); - assert!(slot.best_merged.contains_key(&base_block_hash)); assert_eq!(stats.merged_blocks, 1); }