1use std::sync::{
7 Arc,
8 atomic::{AtomicBool, AtomicU64, Ordering},
9};
10
11use alloy_consensus::{
12 Block, BlockBody, BlockHeader, EMPTY_OMMER_ROOT_HASH, Header, TxReceipt, proofs,
13 transaction::{SignerRecoverable, TxHashRef},
14};
15use alloy_eips::eip2718::Decodable2718;
16use alloy_evm::{
17 EvmFactory,
18 block::{BlockExecutor, BlockExecutorFactory},
19};
20use alloy_primitives::{Address, B64, B256, Bytes, U256};
21use alloy_rpc_types_eth::BlockNumberOrTag;
22use arb_evm::config::{ArbEvmConfig, arbos_version_from_mix_hash, l1_block_number_from_mix_hash};
23use arb_primitives::{ArbPrimitives, signed_tx::ArbTransactionSigned, tx_types::ArbInternalTx};
24use arb_rpc::block_producer::{
25 BlockProducer, BlockProducerError, BlockProductionInput, ProducedBlock,
26};
27use arbos::{
28 header::{ArbHeaderInfo, derive_arb_header_info},
29 internal_tx,
30 parse_l2::{ParsedTransaction, parse_l2_transactions, parsed_tx_to_signed},
31 types::parse_init_message,
32};
33use parking_lot::Mutex;
34use reth_chain_state::{CanonicalInMemoryState, ExecutedBlock, NewCanonicalChain};
35use reth_chainspec::ChainSpec;
36use reth_evm::ConfigureEvm;
37use reth_metrics::{
38 Metrics,
39 metrics::{self, Counter, Gauge, Histogram},
40};
41use reth_primitives_traits::{NodePrimitives, SealedHeader, logs_bloom};
42use reth_provider::{BlockNumReader, BlockReaderIdExt, HeaderProvider, StateProviderFactory};
43use reth_revm::database::StateProviderDatabase;
44use reth_storage_api::{StateProvider, StateProviderBox};
45use reth_trie_common::{HashedPostState, TrieInputSorted};
46use revm::database::{BundleState, StateBuilder};
47use revm_database::states::bundle_state::BundleRetention;
48use tracing::{debug, info, warn};
49
50use crate::genesis;
51
52pub trait InMemoryStateAccess {
58 type Primitives: NodePrimitives;
59 fn canonical_in_memory_state(&self) -> CanonicalInMemoryState<Self::Primitives>;
60}
61
62impl<N> InMemoryStateAccess for reth_provider::providers::BlockchainProvider<N>
64where
65 N: reth_provider::providers::ProviderNodeTypes,
66{
67 type Primitives = N::Primitives;
68 fn canonical_in_memory_state(&self) -> CanonicalInMemoryState<Self::Primitives> {
69 self.canonical_in_memory_state()
70 }
71}
72
73pub const DEFAULT_FLUSH_INTERVAL: u64 = 128;
74const DEFAULT_MAX_INFLIGHT: usize = 512;
75
76fn max_inflight() -> usize {
77 static MAX: std::sync::OnceLock<usize> = std::sync::OnceLock::new();
78 *MAX.get_or_init(|| {
79 std::env::var("ARB_RETH_MAX_INFLIGHT")
80 .ok()
81 .and_then(|s| s.parse::<usize>().ok())
82 .filter(|n| *n > 0)
83 .unwrap_or(DEFAULT_MAX_INFLIGHT)
84 })
85}
86
87pub struct FlushScheduler {
90 interval: u64,
91 ema_commit_latency_ms: u64,
92}
93
94impl FlushScheduler {
95 pub fn new(interval: u64) -> Self {
96 Self {
97 interval,
98 ema_commit_latency_ms: 0,
99 }
100 }
101
102 pub fn should_flush(&self, since_last: u64) -> bool {
103 since_last >= self.interval
104 }
105
106 pub fn observe(&mut self, commit_latency_ms: u64) {
107 self.ema_commit_latency_ms = (self.ema_commit_latency_ms * 7 + commit_latency_ms * 3) / 10;
108 }
109
110 pub fn current_interval(&self) -> u64 {
111 self.interval
112 }
113}
114
115#[cfg(target_os = "linux")]
116fn read_dirty_pages_mb() -> Option<u64> {
117 let content = std::fs::read_to_string("/proc/meminfo").ok()?;
118 for line in content.lines() {
119 if let Some(rest) = line.strip_prefix("Dirty:") {
120 let kb: u64 = rest.trim().trim_end_matches(" kB").trim().parse().ok()?;
121 return Some(kb / 1024);
122 }
123 }
124 None
125}
126
127#[cfg(not(target_os = "linux"))]
128fn read_dirty_pages_mb() -> Option<u64> {
129 None
130}
131
132#[derive(Metrics)]
134#[metrics(scope = "arb_block_producer")]
135struct ArbBlockProducerMetrics {
136 head_block: Gauge,
138 blocks_produced_total: Counter,
140 gas_processed_total: Counter,
142 transactions_processed_total: Counter,
144 flush_commit_duration_seconds: Histogram,
146 backpressure_stall_seconds: Histogram,
148}
149
150pub struct ArbBlockProducer<Provider> {
152 provider: Provider,
153 chain_spec: Arc<ChainSpec>,
154 evm_config: ArbEvmConfig,
155 in_memory_state: CanonicalInMemoryState<ArbPrimitives>,
156 head_block_num: AtomicU64,
157 blocks_since_flush: AtomicU64,
158 scheduler: Mutex<FlushScheduler>,
159 accumulated_trie_input: Mutex<Arc<TrieInputSorted>>,
160 flushing_trie_input: Mutex<Option<Arc<TrieInputSorted>>>,
161 pending_flush: AtomicBool,
162 produce_lock: tokio::sync::Mutex<()>,
163 cached_init: Mutex<Option<arbos::types::ParsedInitMessage>>,
164 finality: Mutex<FinalityMarkers>,
166 validated_watcher: Mutex<Option<Arc<parking_lot::RwLock<alloy_primitives::B256>>>>,
170 cached_overlay: Mutex<Option<CachedOverlay>>,
174 cached_prestate: Mutex<Option<CachedPrestate>>,
175 metrics: ArbBlockProducerMetrics,
176}
177
178#[derive(Debug, Default, Clone)]
179struct FinalityMarkers {
180 safe: Option<alloy_primitives::B256>,
181 finalized: Option<alloy_primitives::B256>,
182 validated: Option<alloy_primitives::B256>,
183}
184
185struct CachedOverlay {
186 parent_hash: B256,
187 overlay: Arc<crate::coalesced_state::CoalescedOverlay>,
188}
189
190struct CachedPrestate {
191 parent_hash: B256,
192 contracts: Arc<alloy_primitives::map::B256Map<revm::bytecode::Bytecode>>,
193}
194
195impl<Provider> ArbBlockProducer<Provider>
196where
197 Provider: BlockNumReader,
198{
199 pub fn new(
200 provider: Provider,
201 chain_spec: Arc<ChainSpec>,
202 evm_config: ArbEvmConfig,
203 in_memory_state: CanonicalInMemoryState<ArbPrimitives>,
204 flush_interval: u64,
205 ) -> Self {
206 let head = provider.last_block_number().unwrap_or(0);
207 Self {
208 provider,
209 chain_spec,
210 evm_config,
211 in_memory_state,
212 head_block_num: AtomicU64::new(head),
213 blocks_since_flush: AtomicU64::new(0),
214 scheduler: Mutex::new(FlushScheduler::new(flush_interval)),
215 accumulated_trie_input: Mutex::new(Arc::new(TrieInputSorted::default())),
216 flushing_trie_input: Mutex::new(None),
217 pending_flush: AtomicBool::new(false),
218 produce_lock: tokio::sync::Mutex::new(()),
219 cached_init: Mutex::new(None),
220 finality: Mutex::new(FinalityMarkers::default()),
221 validated_watcher: Mutex::new(None),
222 cached_overlay: Mutex::new(None),
223 cached_prestate: Mutex::new(None),
224 metrics: ArbBlockProducerMetrics::default(),
225 }
226 }
227
228 fn get_or_build_overlay(
229 &self,
230 parent_hash: B256,
231 head_state: &reth_chain_state::BlockState<ArbPrimitives>,
232 ) -> Arc<crate::coalesced_state::CoalescedOverlay> {
233 let mut cache = self.cached_overlay.lock();
234 if let Some(c) = cache.as_ref()
235 && c.parent_hash == parent_hash
236 {
237 return c.overlay.clone();
238 }
239 let overlay = Arc::new(crate::coalesced_state::CoalescedOverlay::from_chain(
240 head_state,
241 ));
242 *cache = Some(CachedOverlay {
243 parent_hash,
244 overlay: overlay.clone(),
245 });
246 overlay
247 }
248
249 fn extend_cached_overlay(&self, new_block_hash: B256, bundle: &BundleState) {
250 let mut cache = self.cached_overlay.lock();
251 let mut overlay = match cache.take() {
252 Some(c) => match Arc::try_unwrap(c.overlay) {
253 Ok(o) => o,
254 Err(arc) => (*arc).clone(),
255 },
256 None => crate::coalesced_state::CoalescedOverlay::default(),
257 };
258 overlay.extend_with_block(bundle);
259 *cache = Some(CachedOverlay {
260 parent_hash: new_block_hash,
261 overlay: Arc::new(overlay),
262 });
263 }
264
265 fn invalidate_cached_overlay(&self) {
266 *self.cached_overlay.lock() = None;
267 }
268
269 fn get_or_build_prestate(
270 &self,
271 parent_hash: B256,
272 head_state: Option<&reth_chain_state::BlockState<ArbPrimitives>>,
273 ) -> Arc<alloy_primitives::map::B256Map<revm::bytecode::Bytecode>> {
274 let mut cache = self.cached_prestate.lock();
275 if let Some(c) = cache.as_ref()
276 && c.parent_hash == parent_hash
277 {
278 return c.contracts.clone();
279 }
280 let mut contracts: alloy_primitives::map::B256Map<revm::bytecode::Bytecode> =
281 Default::default();
282 if let Some(head_state) = head_state {
283 for block_state in head_state.chain() {
284 let exec_output = &block_state.block().execution_output;
285 for (hash, code) in &exec_output.state.contracts {
286 contracts.entry(*hash).or_insert_with(|| code.clone());
287 }
288 }
289 }
290 let arc = Arc::new(contracts);
291 *cache = Some(CachedPrestate {
292 parent_hash,
293 contracts: arc.clone(),
294 });
295 arc
296 }
297
298 fn extend_cached_prestate(&self, new_block_hash: B256, bundle: &BundleState) {
299 let mut cache = self.cached_prestate.lock();
300 let mut contracts = match cache.take() {
301 Some(c) => match Arc::try_unwrap(c.contracts) {
302 Ok(map) => map,
303 Err(arc) => (*arc).clone(),
304 },
305 None => Default::default(),
306 };
307 for (hash, code) in &bundle.contracts {
308 contracts.entry(*hash).or_insert_with(|| code.clone());
309 }
310 *cache = Some(CachedPrestate {
311 parent_hash: new_block_hash,
312 contracts: Arc::new(contracts),
313 });
314 }
315
316 fn invalidate_cached_prestate(&self) {
317 *self.cached_prestate.lock() = None;
318 }
319
320 pub fn finality_markers(
322 &self,
323 ) -> (
324 Option<alloy_primitives::B256>,
325 Option<alloy_primitives::B256>,
326 Option<alloy_primitives::B256>,
327 ) {
328 let f = self.finality.lock();
329 (f.safe, f.finalized, f.validated)
330 }
331}
332
333impl<Provider> ArbBlockProducer<Provider>
334where
335 Provider: BlockNumReader
336 + BlockReaderIdExt
337 + HeaderProvider<Header = Header>
338 + StateProviderFactory
339 + Send
340 + Sync
341 + 'static,
342{
343 fn head_block_number(&self) -> Result<u64, BlockProducerError> {
345 let head = self.head_block_num.load(Ordering::SeqCst);
346 if head > 0 {
347 Ok(head)
348 } else {
349 self.provider
350 .last_block_number()
351 .map_err(|e| BlockProducerError::StateAccess(e.to_string()))
352 }
353 }
354
355 fn parent_header(&self, head_num: u64) -> Result<SealedHeader<Header>, BlockProducerError> {
357 self.provider
358 .sealed_header_by_number_or_tag(BlockNumberOrTag::Number(head_num))
359 .map_err(|e| BlockProducerError::StateAccess(e.to_string()))?
360 .ok_or_else(|| {
361 BlockProducerError::StateAccess(format!("Parent block {head_num} not found"))
362 })
363 }
364
365 fn drain_completed_flush(&self) -> bool {
366 if !self.pending_flush.load(Ordering::SeqCst) {
367 return false;
368 }
369 let Some(result) = crate::launcher::try_flush_result() else {
370 return false;
371 };
372 self.in_memory_state
373 .remove_persisted_blocks(result.last_num_hash);
374 *self.flushing_trie_input.lock() = None;
375 self.pending_flush.store(false, Ordering::SeqCst);
376 self.invalidate_cached_overlay();
377 self.invalidate_cached_prestate();
378 let commit_latency_ms = result.duration.as_millis() as u64;
379 self.metrics
380 .flush_commit_duration_seconds
381 .record(result.duration.as_secs_f64());
382 let flush_interval_current = {
383 let mut sched = self.scheduler.lock();
384 sched.observe(commit_latency_ms);
385 sched.current_interval()
386 };
387 let dirty_pages_mb = read_dirty_pages_mb().unwrap_or(0);
388 let chain_len_unflushed = self
389 .in_memory_state
390 .head_state()
391 .map(|s| s.chain().count())
392 .unwrap_or(0) as u64;
393 info!(
394 target: "block_producer",
395 flushed = result.count,
396 last_block = result.last_num_hash.number,
397 mdbx_commit_latency_ms = commit_latency_ms,
398 dirty_pages_mb,
399 flush_interval_current,
400 chain_len_unflushed,
401 "block flush"
402 );
403 true
404 }
405
406 async fn apply_backpressure(&self) {
407 let chain_len = self
408 .in_memory_state
409 .head_state()
410 .map(|s| s.chain().count())
411 .unwrap_or(0);
412 let limit = max_inflight();
413 if chain_len <= limit {
414 return;
415 }
416 if !self.pending_flush.load(Ordering::SeqCst) {
417 self.start_async_flush();
418 }
419 let start = std::time::Instant::now();
420 let notifier = crate::launcher::flush_notifier();
421 loop {
422 if let Some(n) = notifier.as_ref() {
423 let notified = n.notified();
426 if self.drain_completed_flush() {
427 break;
428 }
429 let waited = tokio::time::timeout(std::time::Duration::from_secs(30), notified)
430 .await
431 .is_ok();
432 if !waited {
433 warn!(
434 target: "block_producer",
435 chain_len,
436 waited_ms = start.elapsed().as_millis() as u64,
437 "Backpressure: flush notification timed out, polling once"
438 );
439 }
440 } else {
441 if self.drain_completed_flush() {
442 break;
443 }
444 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
445 }
446 }
447 self.metrics
448 .backpressure_stall_seconds
449 .record(start.elapsed().as_secs_f64());
450 warn!(
451 target: "block_producer",
452 chain_len,
453 limit,
454 waited_ms = start.elapsed().as_millis() as u64,
455 "Backpressure: drained pending flush"
456 );
457 }
458
459 fn produce_block_with_execution(
460 &self,
461 input: &BlockProductionInput,
462 parsed_txs: Vec<ParsedTransaction>,
463 ) -> Result<ProducedBlock, BlockProducerError> {
464 self.drain_completed_flush();
465
466 let head_num = self.head_block_number()?;
467 let l2_block_number = head_num + 1;
468 let parent_header = self.parent_header(head_num)?;
469
470 let timestamp = input.l1_timestamp.max(parent_header.timestamp());
471 let time_passed = timestamp.saturating_sub(parent_header.timestamp());
472
473 let parent_mix_hash = parent_header.mix_hash().unwrap_or_default();
474 let parent_arbos_version = arbos_version_from_mix_hash(&parent_mix_hash);
475
476 let l1_block_number = input.l1_block_number;
479 let block_l1_block_number = monotonic_l1_block_number(l1_block_number, &parent_mix_hash);
480 let arbos_version = parent_arbos_version; let send_count = {
484 let mut buf = [0u8; 8];
485 buf.copy_from_slice(&parent_mix_hash.0[0..8]);
486 u64::from_be_bytes(buf)
487 };
488 let provisional_mix_hash =
489 compute_mix_hash(send_count, block_l1_block_number, arbos_version);
490
491 let raw_state_provider = self
493 .provider
494 .state_by_block_hash(parent_header.hash())
495 .map_err(|e| BlockProducerError::StateAccess(e.to_string()))?;
496
497 let state_provider: StateProviderBox = match self
498 .in_memory_state
499 .state_by_hash(parent_header.hash())
500 {
501 Some(head_state) => {
502 let overlay = self.get_or_build_overlay(parent_header.hash(), &head_state);
503 if overlay.is_empty() {
504 raw_state_provider
505 } else {
506 crate::coalesced_state::CoalescedStateProvider::new(raw_state_provider, overlay)
507 .boxed()
508 }
509 }
510 _ => raw_state_provider,
511 };
512
513 let l2_base_fee = {
515 let read_slot = |addr: Address, slot: B256| state_provider.storage(addr, slot);
516 arbos::header::read_l2_base_fee(&read_slot)
517 .map_err(|e| BlockProducerError::Storage(e.to_string()))?
518 .or(parent_header.base_fee_per_gas())
519 };
520
521 let provisional_header = Header {
523 parent_hash: parent_header.hash(),
524 ommers_hash: EMPTY_OMMER_ROOT_HASH,
525 beneficiary: input.sender,
526 state_root: B256::ZERO, transactions_root: B256::ZERO,
528 receipts_root: B256::ZERO,
529 withdrawals_root: None,
530 logs_bloom: Default::default(),
531 timestamp,
532 mix_hash: provisional_mix_hash,
533 nonce: B64::from(input.delayed_messages_read.to_be_bytes()),
534 base_fee_per_gas: l2_base_fee,
535 number: l2_block_number,
536 gas_limit: parent_header.gas_limit(),
537 difficulty: U256::from(1),
538 gas_used: 0,
539 extra_data: Default::default(),
540 parent_beacon_block_root: None,
541 blob_gas_used: None,
542 excess_blob_gas: None,
543 requests_hash: None,
544 };
545
546 let evm_env = self
547 .evm_config
548 .evm_env(&provisional_header)
549 .map_err(|_| BlockProducerError::Execution("evm_env construction failed".into()))?;
550
551 let prestate = {
558 let head_state_opt = self.in_memory_state.state_by_hash(parent_header.hash());
559 let contracts =
560 self.get_or_build_prestate(parent_header.hash(), head_state_opt.as_deref());
561 BundleState {
562 contracts: (*contracts).clone(),
563 ..Default::default()
564 }
565 };
566
567 let mut db = StateBuilder::new()
568 .with_database(StateProviderDatabase::new(state_provider.as_ref()))
569 .with_bundle_prestate(prestate)
570 .with_bundle_update()
571 .build();
572
573 let chain_id = self.chain_spec.chain().id();
574
575 if let Some(init_msg) = self.cached_init.lock().take() {
582 if !genesis::is_arbos_initialized(&mut db) {
583 let initial_version = std::env::var("ARB_INITIAL_ARBOS_VERSION")
589 .ok()
590 .and_then(|v| v.parse::<u64>().ok())
591 .unwrap_or({
592 if parent_arbos_version > 0 {
593 parent_arbos_version
594 } else {
595 genesis::INITIAL_ARBOS_VERSION
596 }
597 });
598 info!(
599 target: "block_producer",
600 initial_version,
601 "Applying cached ArbOS Init during block {} execution",
602 l2_block_number
603 );
604 genesis::initialize_arbos_state(
605 &mut db,
606 &init_msg,
607 chain_id,
608 initial_version,
609 genesis::DEFAULT_CHAIN_OWNER,
610 genesis::ArbOSInit::default(),
611 )
612 .map_err(|e| BlockProducerError::Execution(e.to_string()))?;
613 } else {
614 use arbos::{arbos_state::ArbosState, burn::SystemBurner};
615 info!(
616 target: "block_producer",
617 initial_l1_base_fee = %init_msg.initial_l1_base_fee,
618 "ArbOS already initialized; overriding L1 price_per_unit from Init message"
619 );
620 let state_ptr: *mut _ = &mut db;
626 let mut arb_state =
627 ArbosState::open(unsafe { &mut *state_ptr }, SystemBurner::new(None, false))
628 .map_err(|e| BlockProducerError::Execution(e.to_string()))?;
629 let _ = arb_state
630 .l1_pricing_state
631 .set_price_per_unit(unsafe { &mut *state_ptr }, init_msg.initial_l1_base_fee);
632 if let Ok(target) = std::env::var("ARB_INITIAL_ARBOS_VERSION")
633 && let Ok(target_version) = target.parse::<u64>()
634 {
635 let current = arb_state.arbos_version();
636 if target_version > current {
637 match arb_state.upgrade_arbos_version(
638 unsafe { &mut *state_ptr },
639 target_version,
640 true,
641 ) {
642 Err(e) => {
643 info!(target: "block_producer", err = ?e, target_version, "ArbOS upgrade via env var failed");
644 }
645 _ => {
646 info!(
647 target: "block_producer",
648 from = current,
649 to = target_version,
650 "ArbOS upgraded via ARB_INITIAL_ARBOS_VERSION"
651 );
652 }
653 }
654 }
655 }
656 }
657 }
658
659 let parent_extra = parent_header.extra_data().to_vec();
660 let mut exec_extra = parent_extra.clone();
661 exec_extra.resize(32, 0);
662 exec_extra.extend_from_slice(&input.delayed_messages_read.to_be_bytes());
663
664 let exec_ctx = alloy_evm::eth::EthBlockExecutionCtx {
665 tx_count_hint: Some(parsed_txs.len() + 2), parent_hash: parent_header.hash(),
667 parent_beacon_block_root: None,
668 ommers: &[],
669 withdrawals: None,
670 extra_data: exec_extra.into(),
671 };
672
673 let multi_gas_sink = arb_evm::multi_gas::MultiGasSink::default();
678 let evm = self
679 .evm_config
680 .block_executor_factory()
681 .evm_factory()
682 .create_evm_with_inspector(
683 &mut db,
684 evm_env.clone(),
685 arb_evm::multi_gas::MultiGasInspector::with_sink(multi_gas_sink.clone()),
686 );
687 let mut executor = self
688 .evm_config
689 .block_executor_factory()
690 .create_arb_executor(evm, exec_ctx, chain_id);
691 executor.set_multi_gas_sink(multi_gas_sink);
692 executor.arb_ctx.l2_block_number = l2_block_number;
693 executor.arb_ctx.l1_block_number = block_l1_block_number;
694
695 let l2_hash_entries = {
697 let mut entries = Vec::new();
698 let parent_num = l2_block_number.saturating_sub(1);
699 entries.push((parent_num, parent_header.hash()));
700 let cache_cold = parent_num > 1
701 && self
702 .evm_config
703 .executor_factory
704 .arb_evm_factory()
705 .chain_caches()
706 .l2_block_hashes
707 .lock()
708 .get(&parent_num.saturating_sub(1))
709 .is_none();
710 if cache_cold {
711 let mut hash = parent_header.parent_hash();
712 for i in 2..=256u64 {
713 let Some(n) = l2_block_number.checked_sub(i) else {
714 break;
715 };
716 entries.push((n, hash));
717 match self
718 .provider
719 .sealed_header_by_number_or_tag(BlockNumberOrTag::Number(n))
720 {
721 Ok(Some(h)) => hash = h.parent_hash(),
722 _ => break,
723 }
724 }
725 }
726 entries
727 };
728
729 executor
731 .apply_pre_execution_changes()
732 .map_err(|e| BlockProducerError::Execution(format!("pre-exec: {e}")))?;
733
734 for (l2_num, hash) in l2_hash_entries {
735 executor
736 .precompile_ctx
737 .block
738 .cache_l2_block_hash(l2_num, hash);
739 }
740
741 let mut all_txs: Vec<ArbTransactionSigned> = Vec::new();
742
743 let l1_base_fee = input.l1_base_fee.unwrap_or(U256::ZERO);
745 let start_block_data = internal_tx::encode_start_block(
746 l1_base_fee,
747 l1_block_number,
748 l2_block_number,
749 time_passed,
750 );
751
752 let start_block_tx = create_internal_tx(chain_id, &start_block_data);
753 execute_and_commit_tx(&mut executor, &start_block_tx, "StartBlock")?;
754 all_txs.push(start_block_tx);
755
756 let pre_recovered: Vec<Option<ArbTransactionSigned>> = {
758 use rayon::prelude::*;
759 parsed_txs
760 .par_iter()
761 .map(|parsed| match parsed {
762 ParsedTransaction::InternalStartBlock { .. }
763 | ParsedTransaction::BatchPostingReport { .. } => None,
764 other => {
765 let signed = parsed_tx_to_signed(other, chain_id)?;
766 let _ = signed.recover_signer();
767 Some(signed)
768 }
769 })
770 .collect()
771 };
772
773 for (idx, parsed) in parsed_txs.iter().enumerate() {
775 match parsed {
776 ParsedTransaction::InternalStartBlock { .. } => {
777 continue;
779 }
780 ParsedTransaction::BatchPostingReport {
781 batch_timestamp,
782 batch_poster,
783 batch_number,
784 l1_base_fee_estimate,
785 extra_gas,
786 ..
787 } => {
788 let report_data =
791 if parent_arbos_version >= arb_chainspec::arbos_version::ARBOS_VERSION_50 {
792 let (length, non_zeros) = input.batch_data_stats.unwrap_or((0, 0));
794 internal_tx::encode_batch_posting_report_v2(
795 *batch_timestamp,
796 *batch_poster,
797 *batch_number,
798 length,
799 non_zeros,
800 *extra_gas,
801 *l1_base_fee_estimate,
802 )
803 } else {
804 let legacy_gas = input.batch_gas_cost.unwrap_or(0);
806 let batch_data_gas = legacy_gas.saturating_add(*extra_gas);
807 internal_tx::encode_batch_posting_report(
808 *batch_timestamp,
809 *batch_poster,
810 *batch_number,
811 batch_data_gas,
812 *l1_base_fee_estimate,
813 )
814 };
815 let report_tx = create_internal_tx(chain_id, &report_data);
816 execute_and_commit_tx(&mut executor, &report_tx, "BatchPostingReport")?;
817 all_txs.push(report_tx);
818 continue;
819 }
820 _ => {}
821 }
822
823 let signed_tx = match pre_recovered.get(idx).and_then(|s| s.clone()) {
824 Some(tx) => tx,
825 None => {
826 debug!(target: "block_producer", ?parsed, "Skipping unparseable transaction");
827 continue;
828 }
829 };
830
831 let recovered = match signed_tx.clone().try_into_recovered() {
832 Ok(r) => r,
833 Err(e) => {
834 warn!(target: "block_producer", error = %e, "Failed to recover tx sender, skipping");
835 continue;
836 }
837 };
838 let tx_hash = *signed_tx.tx_hash();
839 let (exec_outcome, hostio_records) = arb_rpc::stylus_tracer::with_trace_buffer(|| {
840 executor.execute_transaction_without_commit(recovered)
841 });
842 match exec_outcome {
843 Ok(result) => {
844 match executor.commit_transaction(result) {
845 Ok(_gas_used) => {
846 all_txs.push(signed_tx);
847 if !hostio_records.is_empty() {
848 arb_rpc::stylus_tracer::cache_trace(tx_hash, hostio_records);
849 }
850
851 loop {
856 let scheduled = executor.drain_scheduled_txs();
857 debug!(
858 target: "block_producer",
859 count = scheduled.len(),
860 "Drained scheduled txs"
861 );
862 if scheduled.is_empty() {
863 break;
864 }
865 for encoded in scheduled {
866 let retry_tx: Option<ArbTransactionSigned> =
867 ArbTransactionSigned::decode_2718(&mut &encoded[..]).ok();
868 if let Some(retry_tx) = retry_tx {
869 let retry_signed = retry_tx.clone();
870 let retry_hash = *retry_signed.tx_hash();
871 match retry_tx.try_into_recovered() {
872 Ok(recovered_retry) => {
873 let (retry_outcome, retry_records) =
874 arb_rpc::stylus_tracer::with_trace_buffer(
875 || {
876 executor
877 .execute_transaction_without_commit(
878 recovered_retry,
879 )
880 },
881 );
882 match retry_outcome {
883 Ok(retry_result) => {
884 match executor
885 .commit_transaction(retry_result)
886 {
887 Ok(_) => {
888 all_txs.push(retry_signed);
889 if !retry_records.is_empty() {
890 arb_rpc::stylus_tracer::cache_trace(
891 retry_hash,
892 retry_records,
893 );
894 }
895 }
896 Err(e) => {
897 warn!(
898 target: "block_producer",
899 error = %e,
900 "Failed to commit auto-redeem tx"
901 );
902 }
903 }
904 }
905 Err(e) => {
906 warn!(
907 target: "block_producer",
908 error = %e,
909 "Auto-redeem tx execution failed"
910 );
911 }
912 }
913 }
914 Err(e) => {
915 warn!(
916 target: "block_producer",
917 error = %e,
918 "Failed to recover auto-redeem tx sender"
919 );
920 }
921 }
922 }
923 }
924 }
925 }
926 Err(e) => {
927 warn!(target: "block_producer", error = %e, "Failed to commit transaction");
928 }
929 }
930 }
931 Err(ref e) if e.to_string().contains("block gas limit reached") => {
932 break;
933 }
934 Err(e) => {
935 warn!(target: "block_producer", error = %e, "Transaction execution failed, skipping");
936 }
937 }
938 }
939
940 let zombie_accounts = executor.zombie_accounts().clone();
941 let finalise_deleted = executor.finalise_deleted().clone();
942
943 let (_, exec_result) = executor
944 .finish()
945 .map_err(|e| BlockProducerError::Execution(format!("finish: {e}")))?;
946
947 let receipts: Vec<arb_primitives::ArbReceipt> = exec_result.receipts;
948
949 db.merge_transitions(BundleRetention::Reverts);
950 let mut bundle = db.take_bundle();
951
952 augment_bundle_from_cache(&mut bundle, &db.cache, &*state_provider)?;
953
954 let keccak_empty_hash = alloy_primitives::B256::from(alloy_primitives::keccak256([]));
956 for addr in &finalise_deleted {
957 if zombie_accounts.contains(addr) {
958 continue;
959 }
960 if bundle.state.contains_key(addr) {
961 let existed_before = state_provider.basic_account(addr).ok().flatten().is_some();
962 if existed_before {
963 let still_empty = bundle
967 .state
968 .get(addr)
969 .and_then(|a| a.info.as_ref())
970 .is_none_or(|info| {
971 info.nonce == 0
972 && info.balance.is_zero()
973 && info.code_hash == keccak_empty_hash
974 });
975 if still_empty && let Some(bundle_acct) = bundle.state.get_mut(addr) {
976 bundle_acct.info = None;
977 }
978 } else {
979 let still_empty = bundle
980 .state
981 .get(addr)
982 .and_then(|a| a.info.as_ref())
983 .is_none_or(|info| {
984 info.nonce == 0
985 && info.balance.is_zero()
986 && info.code_hash == keccak_empty_hash
987 });
988 if still_empty {
989 bundle.state.remove(addr);
990 }
991 }
992 continue;
993 }
994 if let Ok(Some(acct)) = state_provider.basic_account(addr) {
995 let was_originally_empty = acct.balance.is_zero()
996 && acct.nonce == 0
997 && acct.bytecode_hash.is_none_or(|h| h == keccak_empty_hash);
998 if was_originally_empty {
999 continue;
1000 }
1001 bundle.state.insert(
1002 *addr,
1003 revm_database::BundleAccount {
1004 info: None, original_info: None,
1006 storage: Default::default(),
1007 status: revm_database::AccountStatus::Changed,
1008 },
1009 );
1010 }
1011 }
1012
1013 filter_unchanged_storage(&mut bundle);
1014 delete_empty_accounts(&mut bundle, &zombie_accounts, &*state_provider);
1015
1016 let hashed_state =
1017 HashedPostState::from_bundle_state::<reth_trie_common::KeccakKeyHasher>(bundle.state());
1018
1019 let (state_root, trie_updates) = {
1020 let acc_arc = self.accumulated_trie_input.lock().clone();
1021 let flushing_arc = self.flushing_trie_input.lock().clone();
1022
1023 let block_state_sorted = hashed_state.clone().into_sorted();
1024 let prefix_sets = block_state_sorted.construct_prefix_sets().freeze();
1025
1026 let mut new_acc_state = (*acc_arc.state).clone();
1027 new_acc_state.extend_ref_and_sort(&block_state_sorted);
1028 let new_acc_state_arc = Arc::new(new_acc_state);
1029
1030 let (overlay_state_arc, overlay_nodes_arc) = if let Some(f) = &flushing_arc {
1031 let mut s = (*f.state).clone();
1032 s.extend_ref_and_sort(&new_acc_state_arc);
1033 let mut n = (*f.nodes).clone();
1034 n.extend_ref_and_sort(&acc_arc.nodes);
1035 (Arc::new(s), Arc::new(n))
1036 } else {
1037 (Arc::clone(&new_acc_state_arc), Arc::clone(&acc_arc.nodes))
1038 };
1039
1040 let overlay = Arc::new(TrieInputSorted::new(
1041 overlay_nodes_arc,
1042 overlay_state_arc,
1043 Default::default(),
1044 ));
1045
1046 let (root, updates) =
1047 crate::launcher::compute_parallel_state_root(overlay, prefix_sets)
1048 .map_err(|e| BlockProducerError::Execution(format!("state root: {e}")))?;
1049
1050 let mut new_acc_nodes = (*acc_arc.nodes).clone();
1051 new_acc_nodes.extend_ref_and_sort(&updates.clone_into_sorted());
1052 *self.accumulated_trie_input.lock() = Arc::new(TrieInputSorted::new(
1053 Arc::new(new_acc_nodes),
1054 new_acc_state_arc,
1055 Default::default(),
1056 ));
1057
1058 (root, updates)
1059 };
1060
1061 let arb_info =
1063 derive_header_info_from_state(state_provider.as_ref(), &bundle, input.sender)?;
1064
1065 let final_mix_hash = arb_info
1066 .as_ref()
1067 .map(|info| info.compute_mix_hash())
1068 .unwrap_or(provisional_mix_hash);
1069
1070 let extra_data: Bytes = arb_info
1071 .as_ref()
1072 .map(|info| {
1073 let mut data = info.send_root.to_vec();
1074 data.resize(32, 0);
1075 data.into()
1076 })
1077 .unwrap_or_else(|| {
1078 let mut data = parent_extra.clone();
1079 data.resize(32, 0);
1080 data.into()
1081 });
1082
1083 let send_root = arb_info
1084 .as_ref()
1085 .map(|info| info.send_root)
1086 .unwrap_or_else(|| {
1087 if parent_extra.len() >= 32 {
1088 B256::from_slice(&parent_extra[..32])
1089 } else {
1090 B256::ZERO
1091 }
1092 });
1093
1094 let gas_used = exec_result.gas_used;
1096 let logs_bloom_val = logs_bloom(receipts.iter().flat_map(|r| r.logs()));
1097
1098 let transactions_root =
1099 proofs::calculate_transaction_root::<ArbTransactionSigned>(&all_txs);
1100 let receipts_root = proofs::calculate_receipt_root(
1101 &receipts
1102 .iter()
1103 .map(|r| r.with_bloom_ref())
1104 .collect::<Vec<_>>(),
1105 );
1106
1107 let header = Header {
1108 parent_hash: parent_header.hash(),
1109 ommers_hash: EMPTY_OMMER_ROOT_HASH,
1110 beneficiary: input.sender,
1111 state_root,
1112 transactions_root,
1113 receipts_root,
1114 withdrawals_root: None,
1115 logs_bloom: logs_bloom_val,
1116 timestamp,
1117 mix_hash: final_mix_hash,
1118 nonce: B64::from(input.delayed_messages_read.to_be_bytes()),
1119 base_fee_per_gas: l2_base_fee,
1120 number: l2_block_number,
1121 gas_limit: parent_header.gas_limit(),
1122 difficulty: U256::from(1),
1123 gas_used,
1124 extra_data,
1125 parent_beacon_block_root: None,
1126 blob_gas_used: None,
1127 excess_blob_gas: None,
1128 requests_hash: None,
1129 };
1130
1131 let block = Block::<ArbTransactionSigned> {
1132 header,
1133 body: BlockBody {
1134 transactions: all_txs,
1135 ommers: Default::default(),
1136 withdrawals: None,
1137 },
1138 };
1139
1140 let sealed = reth_primitives_traits::SealedBlock::seal_slow(block);
1141 let block_hash = sealed.hash();
1142
1143 self.extend_cached_overlay(block_hash, &bundle);
1144 self.extend_cached_prestate(block_hash, &bundle);
1145
1146 {
1148 use alloy_evm::block::BlockExecutionResult;
1149 use reth_chain_state::ComputedTrieData;
1150 use reth_execution_types::BlockExecutionOutput;
1151 use reth_primitives_traits::RecoveredBlock;
1152
1153 let recovered = Arc::new(RecoveredBlock::new_sealed(sealed.clone(), vec![]));
1154 let exec_output = Arc::new(BlockExecutionOutput {
1155 state: bundle,
1156 result: BlockExecutionResult {
1157 receipts,
1158 requests: Default::default(),
1159 gas_used,
1160 blob_gas_used: 0,
1161 },
1162 });
1163 let computed = ComputedTrieData {
1164 hashed_state: Arc::new(hashed_state.into_sorted()),
1165 trie_updates: Arc::new(trie_updates.into_sorted()),
1166 anchored_trie_input: None,
1167 };
1168 let executed = ExecutedBlock::new(recovered, exec_output, computed);
1169
1170 self.in_memory_state
1171 .update_chain(NewCanonicalChain::Commit {
1172 new: vec![executed],
1173 });
1174
1175 let sealed_header = SealedHeader::new(sealed.header().clone(), sealed.hash());
1176 self.in_memory_state.set_canonical_head(sealed_header);
1177 }
1178
1179 self.head_block_num.store(l2_block_number, Ordering::SeqCst);
1180
1181 let num_txs = sealed.body().transactions.len();
1182 {
1184 self.metrics.head_block.set(l2_block_number as f64);
1185 self.metrics.blocks_produced_total.increment(1);
1186 self.metrics.gas_processed_total.increment(gas_used);
1187 self.metrics
1188 .transactions_processed_total
1189 .increment(num_txs as u64);
1190 }
1191
1192 let since_flush = self.blocks_since_flush.fetch_add(1, Ordering::SeqCst) + 1;
1193 let should_flush = self.scheduler.lock().should_flush(since_flush);
1194 if should_flush && !self.pending_flush.load(Ordering::SeqCst) {
1195 self.start_async_flush();
1196 }
1197
1198 info!(
1199 target: "block_producer",
1200 block_num = l2_block_number,
1201 ?block_hash,
1202 ?send_root,
1203 ?state_root,
1204 num_txs,
1205 gas_used,
1206 "Produced block"
1207 );
1208
1209 Ok(ProducedBlock {
1210 block_hash,
1211 send_root,
1212 })
1213 }
1214
1215 fn start_async_flush(&self) {
1217 let mut blocks: Vec<ExecutedBlock<ArbPrimitives>> = Vec::new();
1218 if let Some(head_state) = self.in_memory_state.head_state() {
1219 for block_state in head_state.chain() {
1220 blocks.push(block_state.block().clone());
1221 }
1222 }
1223 blocks.reverse();
1224
1225 let Some(last) = blocks.last() else {
1226 return;
1227 };
1228 let last_num_hash = alloy_eips::BlockNumHash::new(
1229 last.recovered_block().number(),
1230 last.recovered_block().hash(),
1231 );
1232
1233 let current = std::mem::take(&mut *self.accumulated_trie_input.lock());
1235 *self.flushing_trie_input.lock() = Some(current);
1236
1237 self.blocks_since_flush.store(0, Ordering::SeqCst);
1238 self.pending_flush.store(true, Ordering::SeqCst);
1239
1240 let count = blocks.len();
1241 crate::launcher::start_flush(crate::launcher::FlushRequest {
1242 blocks,
1243 last_num_hash,
1244 });
1245
1246 debug!(
1247 target: "block_producer",
1248 count,
1249 last_block = last_num_hash.number,
1250 "Started async flush"
1251 );
1252 }
1253}
1254
1255#[async_trait::async_trait]
1256impl<Provider> BlockProducer for ArbBlockProducer<Provider>
1257where
1258 Provider: BlockNumReader
1259 + BlockReaderIdExt
1260 + HeaderProvider<Header = Header>
1261 + StateProviderFactory
1262 + Send
1263 + Sync
1264 + 'static,
1265{
1266 fn cache_init_message(&self, l2_msg: &[u8]) -> Result<(), BlockProducerError> {
1267 let init_msg = parse_init_message(l2_msg)
1268 .map_err(|e| BlockProducerError::Parse(format!("init message: {e}")))?;
1269
1270 info!(
1271 target: "block_producer",
1272 chain_id = %init_msg.chain_id,
1273 initial_l1_base_fee = %init_msg.initial_l1_base_fee,
1274 "Cached Init message params"
1275 );
1276
1277 *self.cached_init.lock() = Some(init_msg);
1278 Ok(())
1279 }
1280
1281 async fn produce_block(
1282 &self,
1283 msg_idx: u64,
1284 input: BlockProductionInput,
1285 ) -> Result<ProducedBlock, BlockProducerError> {
1286 let _lock = self.produce_lock.lock().await;
1287
1288 let head_num = self.head_block_number()?;
1290 let expected_block = head_num + 1;
1291 let actual_block = msg_idx;
1292
1293 if expected_block != actual_block {
1294 return Err(BlockProducerError::Unexpected(format!(
1295 "Expected block {expected_block} but got msg_idx {msg_idx} (block {actual_block})"
1296 )));
1297 }
1298
1299 let chain_id = self.chain_spec.chain().id();
1301
1302 let parsed_txs = parse_l2_transactions(
1303 input.kind,
1304 input.sender,
1305 &input.l2_msg,
1306 input.request_id,
1307 input.l1_base_fee,
1308 chain_id,
1309 )
1310 .unwrap_or_else(|e| {
1311 warn!(target: "block_producer", error=%e, "Error parsing L2 message, treating as empty");
1312 vec![]
1313 });
1314
1315 debug!(
1316 target: "block_producer",
1317 msg_idx,
1318 kind = input.kind,
1319 num_txs = parsed_txs.len(),
1320 "Parsed L1 message"
1321 );
1322
1323 self.apply_backpressure().await;
1324 self.produce_block_with_execution(&input, parsed_txs)
1325 }
1326
1327 async fn reset_to_block(&self, target_block_number: u64) -> Result<(), BlockProducerError> {
1328 let _lock = self.produce_lock.lock().await;
1329 let current = self.head_block_num.load(Ordering::SeqCst);
1330 if target_block_number > current {
1331 return Err(BlockProducerError::Unexpected(format!(
1332 "reset target {target_block_number} > current head {current}"
1333 )));
1334 }
1335 if target_block_number == current {
1336 return Ok(());
1337 }
1338
1339 let header = self
1340 .provider
1341 .sealed_header_by_number_or_tag(BlockNumberOrTag::Number(target_block_number))
1342 .map_err(|e| BlockProducerError::StateAccess(e.to_string()))?
1343 .ok_or_else(|| {
1344 BlockProducerError::Unexpected(format!(
1345 "reset target block {target_block_number} not found"
1346 ))
1347 })?;
1348
1349 if self.pending_flush.load(Ordering::SeqCst)
1351 && let Some(result) = crate::launcher::try_flush_result()
1352 {
1353 self.in_memory_state
1354 .remove_persisted_blocks(result.last_num_hash);
1355 *self.flushing_trie_input.lock() = None;
1356 self.pending_flush.store(false, Ordering::SeqCst);
1357 }
1358
1359 let mut old_blocks: Vec<reth_chain_state::ExecutedBlock<ArbPrimitives>> = Vec::new();
1364 for bn in (target_block_number + 1)..=current {
1365 if let Some(state) = self.in_memory_state.state_by_number(bn) {
1366 old_blocks.push(state.block());
1367 }
1368 }
1369
1370 if !old_blocks.is_empty() {
1372 self.in_memory_state
1373 .update_chain(reth_chain_state::NewCanonicalChain::Reorg {
1374 new: Vec::new(),
1375 old: old_blocks,
1376 });
1377 }
1378
1379 self.invalidate_cached_overlay();
1380 self.invalidate_cached_prestate();
1381
1382 self.in_memory_state.set_canonical_head(header.clone());
1385
1386 self.head_block_num
1389 .store(target_block_number, Ordering::SeqCst);
1390
1391 if let Some(rx) = crate::launcher::start_unwind(target_block_number) {
1394 match rx.recv() {
1395 Ok(Ok(())) => {}
1396 Ok(Err(e)) => {
1397 return Err(BlockProducerError::Storage(format!(
1398 "unwind above {target_block_number}: {e}"
1399 )));
1400 }
1401 Err(e) => {
1402 return Err(BlockProducerError::Storage(format!(
1403 "unwind channel closed: {e}"
1404 )));
1405 }
1406 }
1407 }
1408
1409 *self.accumulated_trie_input.lock() = Arc::new(TrieInputSorted::default());
1411
1412 info!(
1413 target: "block_producer",
1414 target = target_block_number,
1415 hash = %header.hash(),
1416 old_count = current - target_block_number,
1417 "reset head"
1418 );
1419 Ok(())
1420 }
1421
1422 fn set_finality(
1423 &self,
1424 safe: Option<alloy_primitives::B256>,
1425 finalized: Option<alloy_primitives::B256>,
1426 validated: Option<alloy_primitives::B256>,
1427 ) -> Result<(), BlockProducerError> {
1428 let mut f = self.finality.lock();
1429 if safe.is_some() {
1430 f.safe = safe;
1431 }
1432 if finalized.is_some() {
1433 f.finalized = finalized;
1434 }
1435 if validated.is_some() {
1436 f.validated = validated;
1437 }
1438 drop(f);
1439
1440 if let Some(h) = safe
1444 && let Ok(Some(sealed)) = self.provider.sealed_header_by_hash(h)
1445 {
1446 self.in_memory_state.set_safe(sealed);
1447 }
1448 if let Some(h) = finalized
1449 && let Ok(Some(sealed)) = self.provider.sealed_header_by_hash(h)
1450 {
1451 self.in_memory_state.set_finalized(sealed);
1452 }
1453 if let Some(h) = validated
1457 && let Some(w) = self.validated_watcher.lock().as_ref()
1458 {
1459 *w.write() = h;
1460 }
1461 Ok(())
1462 }
1463
1464 fn attach_validated_watcher(&self, watcher: Arc<parking_lot::RwLock<alloy_primitives::B256>>) {
1465 *self.validated_watcher.lock() = Some(watcher);
1466 }
1467}
1468
1469fn create_internal_tx(chain_id: u64, data: &[u8]) -> ArbTransactionSigned {
1475 use arb_primitives::signed_tx::ArbTypedTransaction;
1476 let tx = ArbTypedTransaction::Internal(ArbInternalTx {
1477 chain_id: U256::from(chain_id),
1478 data: Bytes::copy_from_slice(data),
1479 });
1480 let sig = alloy_primitives::Signature::new(U256::ZERO, U256::ZERO, false);
1481 ArbTransactionSigned::new_unhashed(tx, sig)
1482}
1483
1484fn execute_and_commit_tx<E>(
1486 executor: &mut E,
1487 tx: &ArbTransactionSigned,
1488 label: &str,
1489) -> Result<(), BlockProducerError>
1490where
1491 E: BlockExecutor<Transaction = ArbTransactionSigned>,
1492{
1493 let recovered = tx
1494 .clone()
1495 .try_into_recovered()
1496 .map_err(|e| BlockProducerError::Execution(format!("{label} recovery: {e}")))?;
1497
1498 let result = executor
1499 .execute_transaction_without_commit(recovered)
1500 .map_err(|e| BlockProducerError::Execution(format!("{label} execution: {e}")))?;
1501
1502 executor
1503 .commit_transaction(result)
1504 .map_err(|e| BlockProducerError::Execution(format!("{label} commit: {e}")))?;
1505
1506 Ok(())
1507}
1508
1509fn compute_mix_hash(send_count: u64, l1_block_number: u64, arbos_version: u64) -> B256 {
1510 arbos::header::compute_arbos_mixhash(send_count, l1_block_number, arbos_version, false)
1511}
1512
1513fn monotonic_l1_block_number(reported: u64, parent_mix_hash: &B256) -> u64 {
1516 reported.max(l1_block_number_from_mix_hash(parent_mix_hash))
1517}
1518
1519fn delete_empty_accounts(
1521 bundle: &mut BundleState,
1522 zombie_accounts: &rustc_hash::FxHashSet<Address>,
1523 state_provider: &dyn StateProvider,
1524) {
1525 let keccak_empty = alloy_primitives::B256::from(alloy_primitives::keccak256([]));
1526 let mut to_remove = Vec::new();
1527 for (addr, account) in bundle.state.iter_mut() {
1528 if let Some(ref info) = account.info {
1529 let is_empty =
1530 info.nonce == 0 && info.balance.is_zero() && info.code_hash == keccak_empty;
1531 if is_empty && !zombie_accounts.contains(addr) {
1532 let existed_before = state_provider.basic_account(addr).ok().flatten().is_some();
1533 if existed_before {
1534 debug!(
1535 target: "block_producer",
1536 addr = ?addr,
1537 "EIP-161: deleting empty account from state"
1538 );
1539 account.info = None;
1540 } else {
1541 to_remove.push(*addr);
1542 }
1543 }
1544 }
1545 }
1546 for addr in to_remove {
1547 bundle.state.remove(&addr);
1548 }
1549}
1550
1551fn filter_unchanged_storage(bundle: &mut BundleState) {
1553 for (_addr, account) in bundle.state.iter_mut() {
1554 account
1555 .storage
1556 .retain(|_key, slot| slot.present_value != slot.previous_or_original_value);
1557 }
1558}
1559
1560fn derive_header_info_from_state(
1562 state_provider: &dyn StateProvider,
1563 bundle_state: &BundleState,
1564 coinbase: Address,
1565) -> Result<Option<ArbHeaderInfo>, BlockProducerError> {
1566 let read_slot = |addr: Address, slot: B256| {
1567 if let Some(account) = bundle_state.state.get(&addr) {
1568 let slot_u256 = U256::from_be_bytes(slot.0);
1569 if let Some(storage_slot) = account.storage.get(&slot_u256) {
1570 return Ok(Some(storage_slot.present_value));
1571 }
1572 }
1573 state_provider.storage(addr, slot)
1574 };
1575
1576 derive_arb_header_info(&read_slot, coinbase)
1577 .map_err(|e| BlockProducerError::Storage(e.to_string()))
1578}
1579
1580fn augment_bundle_from_cache(
1582 bundle: &mut BundleState,
1583 cache: &revm_database::CacheState,
1584 state_provider: &dyn StateProvider,
1585) -> Result<(), BlockProducerError> {
1586 use revm_database::states::plain_account::StorageSlot;
1587
1588 for (addr, cache_acct) in &cache.accounts {
1589 let current_info = cache_acct.account.as_ref().map(|a| a.info.clone());
1590 let current_storage = cache_acct
1591 .account
1592 .as_ref()
1593 .map(|a| &a.storage)
1594 .cloned()
1595 .unwrap_or_default();
1596
1597 if let Some(bundle_acct) = bundle.state.get_mut(addr) {
1598 bundle_acct.info = current_info;
1600
1601 for (key, value) in ¤t_storage {
1602 if let Some(slot) = bundle_acct.storage.get_mut(key) {
1603 slot.present_value = *value;
1604 } else {
1605 let original_value = state_provider
1607 .storage(*addr, B256::from(*key))
1608 .map_err(|e| BlockProducerError::Storage(e.to_string()))?
1609 .unwrap_or(U256::ZERO);
1610 if *value != original_value {
1611 bundle_acct.storage.insert(
1612 *key,
1613 StorageSlot {
1614 previous_or_original_value: original_value,
1615 present_value: *value,
1616 },
1617 );
1618 }
1619 }
1620 }
1621 } else {
1622 let original = state_provider
1624 .basic_account(addr)
1625 .map_err(|e| BlockProducerError::Storage(e.to_string()))?;
1626
1627 let info_changed = match (&original, ¤t_info) {
1628 (None, None) => false,
1629 (Some(_), None) | (None, Some(_)) => true,
1630 (Some(orig), Some(curr)) => {
1631 orig.balance != curr.balance
1632 || orig.nonce != curr.nonce
1633 || orig
1634 .bytecode_hash
1635 .unwrap_or(alloy_primitives::KECCAK256_EMPTY)
1636 != curr.code_hash
1637 }
1638 };
1639
1640 let mut storage_changes: revm_database::StorageWithOriginalValues = Default::default();
1641 for (key, value) in ¤t_storage {
1642 let original_value = state_provider
1643 .storage(*addr, B256::from(*key))
1644 .map_err(|e| BlockProducerError::Storage(e.to_string()))?
1645 .unwrap_or(U256::ZERO);
1646 if original_value != *value {
1647 storage_changes.insert(
1648 *key,
1649 StorageSlot {
1650 previous_or_original_value: original_value,
1651 present_value: *value,
1652 },
1653 );
1654 }
1655 }
1656
1657 if info_changed || !storage_changes.is_empty() {
1658 let original_info = original.as_ref().map(|a| revm::state::AccountInfo {
1659 balance: a.balance,
1660 nonce: a.nonce,
1661 code_hash: a.bytecode_hash.unwrap_or(alloy_primitives::KECCAK256_EMPTY),
1662 code: None,
1663 account_id: None,
1664 });
1665
1666 let status = if original.is_some() {
1667 revm_database::AccountStatus::Changed
1668 } else {
1669 revm_database::AccountStatus::InMemoryChange
1670 };
1671
1672 bundle.state.insert(
1673 *addr,
1674 revm_database::BundleAccount {
1675 info: current_info,
1676 original_info,
1677 storage: storage_changes,
1678 status,
1679 },
1680 );
1681 }
1682 }
1683 }
1684 Ok(())
1685}
1686
1687#[cfg(test)]
1688mod tests {
1689 use arbos::header::compute_arbos_mixhash;
1690
1691 use super::*;
1692
1693 #[test]
1694 fn l1_block_number_clamps_to_parent() {
1695 let parent = compute_arbos_mixhash(0, 10_538_022, 51, false);
1696 assert_eq!(monotonic_l1_block_number(10_537_967, &parent), 10_538_022);
1698 assert_eq!(monotonic_l1_block_number(10_538_099, &parent), 10_538_099);
1700 assert_eq!(monotonic_l1_block_number(10_538_022, &parent), 10_538_022);
1702 }
1703}