Skip to main content

foundry_evm_core/fork/
multi.rs

1//! Support for running multiple fork backends.
2//!
3//! The design is similar to the single `SharedBackend`, `BackendHandler` but supports multiple
4//! concurrently active pairs at once.
5
6use super::{CreateFork, Fork, bal};
7use crate::{FoundryBlock, opts::ForkContext};
8use alloy_eips::{BlockNumHash, eip7928::BlockAccessList};
9use alloy_evm::EvmEnv;
10use alloy_network::{AnyNetwork, Network};
11use alloy_primitives::map::HashMap;
12use foundry_config::Config;
13use foundry_fork_db::{
14    BackendHandler, BlockchainDb, ForkBlock, ForkBlockEnv, SharedBackend, cache::BlockchainDbMeta,
15};
16use futures::{
17    FutureExt, StreamExt,
18    channel::mpsc::{Receiver, Sender, channel},
19    stream::Fuse,
20    task::{Context, Poll},
21};
22use revm::primitives::hardfork::SpecId;
23use std::{
24    fmt::{self, Write},
25    pin::Pin,
26    sync::{
27        Arc,
28        atomic::{AtomicBool, AtomicUsize, Ordering},
29        mpsc::{Sender as OneshotSender, channel as oneshot_channel},
30    },
31    time::Duration,
32};
33
34/// The _unique_ identifier for a specific fork, this could be the name of the network a custom
35/// descriptive name.
36#[derive(Clone, Debug, PartialEq, Eq, Hash)]
37pub struct ForkId(pub String);
38
39impl ForkId {
40    /// Returns the identifier for a Fork from a URL and block number.
41    pub fn new(url: &str, num: Option<u64>) -> Self {
42        Self::new_with_context(url, num, None)
43    }
44
45    fn new_with_context(
46        url: &str,
47        num: Option<u64>,
48        context: Option<&crate::opts::ForkContext>,
49    ) -> Self {
50        let mut id = url.to_string();
51        if let Some(context) = context {
52            write!(
53                id,
54                "#{}:{}:{}:{}:{:?}:{:?}:{:?}:{:?}",
55                context.execution_chain_id,
56                context.source_chain_id,
57                context.network,
58                context.network_profile.execution_profile_name(),
59                context.hardfork,
60                context.instance_id,
61                context.source_fork_block_number,
62                context.source_fork_block_hash
63            )
64            .unwrap();
65        }
66        id.push('@');
67        match num {
68            Some(n) => write!(id, "{n:#x}").unwrap(),
69            None => id.push_str("latest"),
70        }
71        Self(id)
72    }
73
74    /// Returns the identifier for an exactly resolved fork.
75    fn resolved(url: &str, fork: &Fork) -> Self {
76        let mut id = Self::exact(url, fork, fork.block());
77        if fork.state_by_number {
78            id.0.push_str("#state-by-number");
79        }
80        id
81    }
82
83    /// Exact rolls always read state by hash, regardless of the parent fork's state mode.
84    fn exact(url: &str, fork: &Fork, block: BlockNumHash) -> Self {
85        let mut id = Self::new_with_context(url, Some(block.number), Some(&fork.context())).0;
86        write!(id, "#{}:{}", block.hash, fork.source_id()).unwrap();
87        Self(id)
88    }
89
90    /// Returns the identifier of the fork.
91    pub fn as_str(&self) -> &str {
92        &self.0
93    }
94}
95
96impl fmt::Display for ForkId {
97    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
98        self.0.fmt(f)
99    }
100}
101
102impl<T: Into<String>> From<T> for ForkId {
103    fn from(id: T) -> Self {
104        Self(id.into())
105    }
106}
107
108/// Backend, environment, and identity returned after creating or rolling a fork.
109pub struct ForkResult<N: Network, SPEC, BLOCK: ForkBlockEnv> {
110    /// Identifier assigned to the fork.
111    pub id: ForkId,
112    /// Backend pinned to the resolved fork block.
113    pub backend: SharedBackend<N, BLOCK>,
114    /// EVM environment reconstructed from the resolved fork block.
115    pub env: EvmEnv<SPEC, BLOCK>,
116    /// Exact source and block identity used to construct the backend.
117    pub fork: Fork,
118}
119
120/// The Sender half of multi fork pair.
121/// Can send requests to the `MultiForkHandler` to create forks.
122#[derive(Clone, Debug)]
123#[must_use]
124pub struct MultiFork<N: Network, SPEC, BLOCK: ForkBlockEnv> {
125    /// Channel to send `Request`s to the handler.
126    handler: Sender<Request<N, SPEC, BLOCK>>,
127    /// Ensures that all rpc resources get flushed properly.
128    _shutdown: Arc<ShutDownMultiFork<N, SPEC, BLOCK>>,
129}
130
131impl<
132    N: Network,
133    SPEC: Into<SpecId> + Default + Copy + Unpin + Send + 'static,
134    BLOCK: FoundryBlock + ForkBlockEnv + Default + Unpin,
135> MultiFork<N, SPEC, BLOCK>
136{
137    /// Creates a new pair and spawns the `MultiForkHandler` on a background thread.
138    pub fn spawn() -> Self {
139        Self::spawn_with_forks(HashMap::default(), None)
140    }
141
142    /// Starts an independent registry sharing the existing remote caches.
143    pub(crate) fn scoped(&self) -> eyre::Result<Self> {
144        let (sender, rx) = oneshot_channel();
145        self.handler
146            .clone()
147            .try_send(Request::CloneForks(sender))
148            .map_err(|e| eyre::eyre!("{e:?}"))?;
149        Ok(Self::spawn_with_forks(rx.recv()?, Some(self._shutdown.clone())))
150    }
151
152    fn spawn_with_forks(
153        forks: HashMap<ForkId, CreatedFork<N, SPEC, BLOCK>>,
154        parent: Option<Arc<ShutDownMultiFork<N, SPEC, BLOCK>>>,
155    ) -> Self {
156        trace!(target: "fork::multi", "spawning multifork");
157        let (mut fork, mut handler) = Self::new();
158        handler.forks = forks;
159        Arc::get_mut(&mut fork._shutdown).unwrap()._parent = parent;
160
161        // Spawn a light-weight thread just for sending and receiving data from the remote
162        // client(s).
163        let fut = async move {
164            // Flush cache every 60s, this ensures that long-running fork tests get their
165            // cache flushed from time to time.
166            // NOTE: we install the interval here because the `tokio::timer::Interval`
167            // requires a rt.
168            handler.set_flush_cache_interval(Duration::from_secs(60));
169            handler.await
170        };
171        match tokio::runtime::Handle::try_current() {
172            Ok(rt) => _ = rt.spawn(fut),
173            Err(_) => {
174                trace!(target: "fork::multi", "spawning multifork backend thread");
175                _ = std::thread::Builder::new()
176                    .name("multi-fork-backend".into())
177                    .spawn(move || {
178                        tokio::runtime::Builder::new_current_thread()
179                            .enable_all()
180                            .build()
181                            .expect("failed to build tokio runtime")
182                            .block_on(fut)
183                    })
184                    .expect("failed to spawn thread")
185            }
186        }
187
188        trace!(target: "fork::multi", "spawned MultiForkHandler thread");
189        fork
190    }
191
192    /// Creates a new pair multi fork pair.
193    ///
194    /// Use [`spawn`](Self::spawn) instead.
195    #[doc(hidden)]
196    pub fn new() -> (Self, MultiForkHandler<N, SPEC, BLOCK>) {
197        let (handler, handler_rx) = channel(1);
198        let _shutdown =
199            Arc::new(ShutDownMultiFork { handler: Some(handler.clone()), _parent: None });
200        (Self { handler, _shutdown }, MultiForkHandler::new(handler_rx))
201    }
202
203    /// Returns a fork backend.
204    ///
205    /// If no matching fork backend exists it will be created.
206    pub fn create_fork(&self, fork: CreateFork) -> eyre::Result<ForkResult<N, SPEC, BLOCK>> {
207        trace!("Creating new fork, url={}, block={:?}", fork.url, fork.evm_opts.fork_block_number);
208        let (sender, rx) = oneshot_channel();
209        let req = Request::CreateFork(Box::new(fork), sender);
210        self.handler.clone().try_send(req).map_err(|e| eyre::eyre!("{:?}", e))?;
211        rx.recv()?
212    }
213
214    /// Rolls the block of the fork.
215    ///
216    /// If no matching fork backend exists it will be created.
217    pub fn roll_fork(&self, fork: ForkId, block: u64) -> eyre::Result<ForkResult<N, SPEC, BLOCK>> {
218        trace!(?fork, ?block, "rolling fork");
219        let (sender, rx) = oneshot_channel();
220        let req = Request::RollFork(fork, block, sender);
221        self.handler.clone().try_send(req).map_err(|e| eyre::eyre!("{:?}", e))?;
222        rx.recv()?
223    }
224
225    /// Rolls a fork to an already resolved exact block.
226    pub fn roll_fork_exact(
227        &self,
228        fork: ForkId,
229        block: BlockNumHash,
230    ) -> eyre::Result<ForkResult<N, SPEC, BLOCK>> {
231        self.roll_fork_exact_with_bal(fork, block, false)
232    }
233
234    /// Rolls to an exact parent block, optionally warming its cache before transaction replay.
235    pub(crate) fn roll_fork_exact_with_bal(
236        &self,
237        fork: ForkId,
238        block: BlockNumHash,
239        prewarm_bal: bool,
240    ) -> eyre::Result<ForkResult<N, SPEC, BLOCK>> {
241        trace!(?fork, ?block, "rolling fork to exact block");
242        let (sender, rx) = oneshot_channel();
243        let req = Request::RollForkExact(fork, block, prewarm_bal, sender);
244        self.handler.clone().try_send(req).map_err(|e| eyre::eyre!("{:?}", e))?;
245        rx.recv()?
246    }
247
248    /// Returns the `EvmEnv` of the given fork, if any.
249    pub fn get_evm_env(&self, fork: ForkId) -> eyre::Result<Option<EvmEnv<SPEC, BLOCK>>> {
250        trace!(?fork, "getting env config");
251        let (sender, rx) = oneshot_channel();
252        let req = Request::GetEvmEnv(fork, sender);
253        self.handler.clone().try_send(req).map_err(|e| eyre::eyre!("{:?}", e))?;
254        Ok(rx.recv()?)
255    }
256
257    /// Returns the corresponding fork if it exists.
258    ///
259    /// Returns `None` if no matching fork backend is available.
260    pub fn get_fork(&self, id: impl Into<ForkId>) -> eyre::Result<Option<SharedBackend<N, BLOCK>>> {
261        let id = id.into();
262        trace!(?id, "get fork backend");
263        let (sender, rx) = oneshot_channel();
264        let req = Request::GetFork(id, sender);
265        self.handler.clone().try_send(req).map_err(|e| eyre::eyre!("{:?}", e))?;
266        Ok(rx.recv()?)
267    }
268
269    /// Returns the immutable remote fork owned by the selected backend.
270    pub fn get_fork_info(&self, id: impl Into<ForkId>) -> eyre::Result<Option<Fork>> {
271        let (sender, rx) = oneshot_channel();
272        self.handler
273            .clone()
274            .try_send(Request::GetForkInfo(id.into(), sender))
275            .map_err(|e| eyre::eyre!("{e:?}"))?;
276        Ok(rx.recv()?)
277    }
278
279    /// Returns the corresponding fork url if it exists.
280    ///
281    /// Returns `None` if no matching fork is available.
282    pub fn get_fork_url(&self, id: impl Into<ForkId>) -> eyre::Result<Option<String>> {
283        let (sender, rx) = oneshot_channel();
284        let req = Request::GetForkUrl(id.into(), sender);
285        self.handler.clone().try_send(req).map_err(|e| eyre::eyre!("{:?}", e))?;
286        Ok(rx.recv()?)
287    }
288
289    /// Returns the options used to create the corresponding fork if it exists.
290    ///
291    /// Returns `None` if no matching fork is available.
292    pub fn get_fork_options(&self, id: impl Into<ForkId>) -> eyre::Result<Option<CreateFork>> {
293        let (sender, rx) = oneshot_channel();
294        let req = Request::GetForkOptions(id.into(), sender);
295        self.handler.clone().try_send(req).map_err(|e| eyre::eyre!("{:?}", e))?;
296        Ok(rx.recv()?)
297    }
298}
299
300type CreateFuture<SPEC, BLOCK> = Pin<
301    Box<
302        dyn Future<
303                Output = eyre::Result<(
304                    ForkId,
305                    CreateFork,
306                    Fork,
307                    EvmEnv<SPEC, BLOCK>,
308                    Option<BlockAccessList>,
309                )>,
310            > + Send,
311    >,
312>;
313type CreateSender<N, SPEC, BLOCK> = OneshotSender<eyre::Result<ForkResult<N, SPEC, BLOCK>>>;
314type GetEvmEnvSender<SPEC, BLOCK> = OneshotSender<Option<EvmEnv<SPEC, BLOCK>>>;
315
316/// Request that's send to the handler.
317#[derive(Debug)]
318enum Request<N: Network, SPEC, BLOCK: ForkBlockEnv> {
319    CloneForks(OneshotSender<HashMap<ForkId, CreatedFork<N, SPEC, BLOCK>>>),
320    /// Creates a new ForkBackend.
321    CreateFork(Box<CreateFork>, CreateSender<N, SPEC, BLOCK>),
322    /// Returns the Fork backend for the `ForkId` if it exists.
323    GetFork(ForkId, OneshotSender<Option<SharedBackend<N, BLOCK>>>),
324    /// Adjusts the block that's being forked, by creating a new fork at the new block.
325    RollFork(ForkId, u64, CreateSender<N, SPEC, BLOCK>),
326    /// Adjusts the fork to an already resolved exact block.
327    RollForkExact(ForkId, BlockNumHash, bool, CreateSender<N, SPEC, BLOCK>),
328    /// Returns the environment of the fork.
329    GetEvmEnv(ForkId, GetEvmEnvSender<SPEC, BLOCK>),
330    /// Shutdowns the entire `MultiForkHandler`, see `ShutDownMultiFork`
331    ShutDown(OneshotSender<()>),
332    /// Returns the Fork Url for the `ForkId` if it exists.
333    GetForkUrl(ForkId, OneshotSender<Option<String>>),
334    GetForkInfo(ForkId, OneshotSender<Option<Fork>>),
335    /// Returns the options used to create the `ForkId` if it exists.
336    GetForkOptions(ForkId, OneshotSender<Option<CreateFork>>),
337}
338
339enum ForkTask<N: Network, SPEC, BLOCK: ForkBlockEnv> {
340    /// Contains the future that will establish a new fork.
341    Create {
342        future: CreateFuture<SPEC, BLOCK>,
343        id: ForkId,
344        prewarm_bal: bool,
345        no_fork_bal: bool,
346        sender: CreateSender<N, SPEC, BLOCK>,
347        additional_senders: Vec<CreateSender<N, SPEC, BLOCK>>,
348    },
349}
350
351/// The type that manages connections in the background.
352#[must_use = "futures do nothing unless polled"]
353pub struct MultiForkHandler<N: Network, SPEC, BLOCK: ForkBlockEnv> {
354    /// Incoming requests from the `MultiFork`.
355    incoming: Fuse<Receiver<Request<N, SPEC, BLOCK>>>,
356
357    /// All active handlers.
358    ///
359    /// It's expected that this list will be rather small (<10).
360    handlers: Vec<(ForkId, BackendHandler<N, BLOCK>)>,
361
362    // tasks currently in progress
363    pending_tasks: Vec<ForkTask<N, SPEC, BLOCK>>,
364
365    /// All _unique_ forkids mapped to their corresponding backend.
366    ///
367    /// Note: The backend can be shared by multiple ForkIds if the target the same provider and
368    /// block number.
369    forks: HashMap<ForkId, CreatedFork<N, SPEC, BLOCK>>,
370
371    /// Optional periodic interval to flush rpc cache.
372    flush_cache_interval: Option<tokio::time::Interval>,
373}
374
375impl<
376    N: Network,
377    SPEC: Into<SpecId> + Default + Copy + Send + 'static,
378    BLOCK: FoundryBlock + ForkBlockEnv + Default,
379> MultiForkHandler<N, SPEC, BLOCK>
380{
381    fn new(incoming: Receiver<Request<N, SPEC, BLOCK>>) -> Self {
382        Self {
383            incoming: incoming.fuse(),
384            handlers: Default::default(),
385            pending_tasks: Default::default(),
386            forks: Default::default(),
387            flush_cache_interval: None,
388        }
389    }
390
391    /// Sets the interval after which all rpc caches should be flushed periodically.
392    pub fn set_flush_cache_interval(&mut self, period: Duration) -> &mut Self {
393        self.flush_cache_interval =
394            Some(tokio::time::interval_at(tokio::time::Instant::now() + period, period));
395        self
396    }
397
398    /// Returns the list of additional senders of a matching task for the given id, if any.
399    fn find_in_progress_task(
400        &mut self,
401        id: &ForkId,
402        prewarm_bal: bool,
403        no_fork_bal: bool,
404    ) -> Option<&mut Vec<CreateSender<N, SPEC, BLOCK>>> {
405        for ForkTask::Create {
406            id: in_progress,
407            prewarm_bal: pending_prewarm,
408            no_fork_bal: pending_opt_out,
409            additional_senders,
410            ..
411        } in &mut self.pending_tasks
412        {
413            if in_progress == id
414                && *pending_prewarm == prewarm_bal
415                && *pending_opt_out == no_fork_bal
416            {
417                return Some(additional_senders);
418            }
419        }
420        None
421    }
422
423    fn create_fork(&mut self, fork: CreateFork, sender: CreateSender<N, SPEC, BLOCK>) {
424        self.create_fork_with_identity(fork, None, false, None, sender);
425    }
426
427    fn create_fork_with_identity(
428        &mut self,
429        fork: CreateFork,
430        expected_identity: Option<ForkContext>,
431        prewarm_bal: bool,
432        exact_block: Option<(Fork, BlockNumHash)>,
433        sender: CreateSender<N, SPEC, BLOCK>,
434    ) {
435        let no_fork_bal = fork.evm_opts.no_fork_bal;
436        let prewarm_bal = prewarm_bal && !no_fork_bal;
437        let resolved_id =
438            exact_block.as_ref().map(|(parent, block)| ForkId::exact(&fork.url, parent, *block));
439        trace!(?resolved_id, "creating fork");
440
441        // Only deduplicate requests that already carry an exact identity. Unresolved requests at
442        // the same URL and height can resolve to different blocks across a reorganization.
443        if let Some(fork_id) = &resolved_id
444            && let Some(in_progress) = self.find_in_progress_task(fork_id, prewarm_bal, no_fork_bal)
445        {
446            in_progress.push(sender);
447            return;
448        }
449
450        let already_prewarmed = resolved_id
451            .as_ref()
452            .and_then(|id| self.forks.get(id))
453            .is_some_and(|fork| fork.bal_prewarmed.load(Ordering::Relaxed));
454
455        // Need to create a new fork.
456        let task_id =
457            resolved_id.unwrap_or_else(|| ForkId::new(&fork.url, fork.evm_opts.fork_block_number));
458        let needs_bal = prewarm_bal && !already_prewarmed;
459        let future = Box::pin(create_fork(fork, expected_identity, needs_bal, exact_block));
460        self.pending_tasks.push(ForkTask::Create {
461            future,
462            id: task_id,
463            prewarm_bal,
464            no_fork_bal,
465            sender,
466            additional_senders: Vec::new(),
467        });
468    }
469
470    fn insert_new_fork(
471        &mut self,
472        fork_id: ForkId,
473        fork: CreatedFork<N, SPEC, BLOCK>,
474        sender: CreateSender<N, SPEC, BLOCK>,
475        additional_senders: Vec<CreateSender<N, SPEC, BLOCK>>,
476    ) {
477        self.forks.insert(fork_id.clone(), fork.clone());
478        let resolved = fork.fork.clone();
479        let _ = sender.send(Ok(ForkResult {
480            id: fork_id.clone(),
481            backend: fork.backend.clone(),
482            env: fork.evm_env.clone(),
483            fork: resolved.clone(),
484        }));
485
486        // Notify all additional senders and track unique forkIds.
487        for sender in additional_senders {
488            let next_fork_id = fork.inc_senders(fork_id.clone());
489            self.forks.insert(next_fork_id.clone(), fork.clone());
490            let _ = sender.send(Ok(ForkResult {
491                id: next_fork_id,
492                backend: fork.backend.clone(),
493                env: fork.evm_env.clone(),
494                fork: resolved.clone(),
495            }));
496        }
497    }
498
499    fn on_request(&mut self, req: Request<N, SPEC, BLOCK>) {
500        match req {
501            Request::CloneForks(sender) => {
502                let _ = sender.send(self.forks.clone());
503            }
504            Request::CreateFork(fork, sender) => self.create_fork(*fork, sender),
505            Request::GetFork(fork_id, sender) => {
506                let fork = self.forks.get(&fork_id).map(|f| f.backend.clone());
507                let _ = sender.send(fork);
508            }
509            Request::RollFork(fork_id, block, sender) => {
510                if let Some(fork) = self.forks.get(&fork_id) {
511                    trace!(target: "fork::multi", "rolling {} to {}", fork_id, block);
512                    let expected_identity = Some(fork.fork.context());
513                    let mut opts = fork.opts.clone();
514                    opts.evm_opts.fork_block_number = Some(block);
515                    opts.evm_opts.fork_block_number_is_inferred = false;
516                    self.create_fork_with_identity(opts, expected_identity, false, None, sender)
517                } else {
518                    let _ =
519                        sender.send(Err(eyre::eyre!("No matching fork exists for {}", fork_id)));
520                }
521            }
522            Request::RollForkExact(fork_id, block, prewarm_bal, sender) => {
523                if let Some(fork) = self.forks.get(&fork_id) {
524                    trace!(target: "fork::multi", "rolling {} to exact block {:?}", fork_id, block);
525                    let mut opts = fork.opts.clone();
526                    opts.evm_opts.fork_block_number = Some(block.number);
527                    opts.evm_opts.fork_state_by_number = false;
528                    opts.evm_opts.fork_block_number_is_inferred = false;
529                    self.create_fork_with_identity(
530                        opts,
531                        None,
532                        prewarm_bal,
533                        Some((fork.fork.clone(), block)),
534                        sender,
535                    )
536                } else {
537                    let _ =
538                        sender.send(Err(eyre::eyre!("No matching fork exists for {}", fork_id)));
539                }
540            }
541            Request::GetEvmEnv(fork_id, sender) => {
542                let _ = sender.send(self.forks.get(&fork_id).map(|fork| fork.evm_env.clone()));
543            }
544            Request::ShutDown(sender) => {
545                trace!(target: "fork::multi", "received shutdown signal");
546                // We're emptying all fork backends, this way we ensure all caches get flushed.
547                self.forks.clear();
548                self.handlers.clear();
549                let _ = sender.send(());
550            }
551            Request::GetForkInfo(fork_id, sender) => {
552                let _ = sender.send(self.forks.get(&fork_id).map(|fork| fork.fork.clone()));
553            }
554            Request::GetForkUrl(fork_id, sender) => {
555                let fork = self.forks.get(&fork_id).map(|f| f.opts.url.clone());
556                let _ = sender.send(fork);
557            }
558            Request::GetForkOptions(fork_id, sender) => {
559                let fork = self.forks.get(&fork_id).map(|f| f.opts.clone());
560                let _ = sender.send(fork);
561            }
562        }
563    }
564}
565
566// Drives all handler to completion.
567// This future will finish once all underlying BackendHandler are completed.
568impl<
569    N: Network,
570    SPEC: Into<SpecId> + Default + Copy + Unpin + Send + 'static,
571    BLOCK: FoundryBlock + ForkBlockEnv + Default + Unpin,
572> Future for MultiForkHandler<N, SPEC, BLOCK>
573{
574    type Output = ();
575
576    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
577        let this = self.get_mut();
578
579        // Receive new requests.
580        loop {
581            match this.incoming.poll_next_unpin(cx) {
582                Poll::Ready(Some(req)) => this.on_request(req),
583                Poll::Ready(None) => {
584                    // Channel closed, but we still need to drive the fork handlers to completion.
585                    trace!(target: "fork::multi", "request channel closed");
586                    break;
587                }
588                Poll::Pending => break,
589            }
590        }
591
592        // Advance all tasks.
593        for n in (0..this.pending_tasks.len()).rev() {
594            let task = this.pending_tasks.swap_remove(n);
595            match task {
596                ForkTask::Create {
597                    mut future,
598                    id,
599                    prewarm_bal,
600                    no_fork_bal,
601                    sender,
602                    additional_senders,
603                } => {
604                    if let Poll::Ready(resp) = future.poll_unpin(cx) {
605                        let resp = resp.and_then(|(fork_id, opts, resolved, evm_env, mut bal)| {
606                            let (fork_id, fork) = if let Some(mut cached) =
607                                this.forks.get(&fork_id).cloned()
608                            {
609                                if bal.is_some()
610                                    && resolved.fingerprint() != cached.fork.fingerprint()
611                                {
612                                    debug!(target: "backend::fork", "ignoring fork BAL for a different cache identity");
613                                    bal = None;
614                                }
615                                // Share the remote cache, RPC client, and block while retaining
616                                // this consumer's selector, execution environment, and BAL
617                                // policy. The fork ID includes the block hash, so the cached
618                                // block is the same block.
619                                cached.opts = opts;
620                                cached.fork = Fork {
621                                    client: cached.fork.client,
622                                    block: cached.fork.block,
623                                    ..resolved
624                                };
625                                cached.evm_env = evm_env;
626                                (cached.inc_senders(fork_id), cached)
627                            } else {
628                                // Only a new backend loads the disk cache.
629                                let (backend, handler) =
630                                    create_backend(&opts, &resolved, &evm_env.block_env)?;
631                                this.handlers.push((fork_id.clone(), handler));
632                                (fork_id, CreatedFork::new(opts, resolved, evm_env, backend))
633                            };
634                            Ok((fork_id, fork, bal))
635                        });
636                        match resp {
637                            Ok((fork_id, fork, bal)) => {
638                                // Apply only after choosing the backend, including an existing
639                                // cache.
640                                if let Some(bal) = bal {
641                                    bal::cache_bal(&fork.backend.data(), bal);
642                                    fork.bal_prewarmed.store(true, Ordering::Relaxed);
643                                }
644                                this.insert_new_fork(fork_id, fork, sender, additional_senders);
645                            }
646                            Err(err) => {
647                                let _ = sender.send(Err(eyre::eyre!("{err}")));
648                                for sender in additional_senders {
649                                    let _ = sender.send(Err(eyre::eyre!("{err}")));
650                                }
651                            }
652                        }
653                    } else {
654                        this.pending_tasks.push(ForkTask::Create {
655                            future,
656                            id,
657                            prewarm_bal,
658                            no_fork_bal,
659                            sender,
660                            additional_senders,
661                        });
662                    }
663                }
664            }
665        }
666
667        // Advance all handlers.
668        for n in (0..this.handlers.len()).rev() {
669            let (id, mut handler) = this.handlers.swap_remove(n);
670            match handler.poll_unpin(cx) {
671                Poll::Ready(_) => {
672                    trace!(target: "fork::multi", "fork {:?} completed", id);
673                }
674                Poll::Pending => {
675                    this.handlers.push((id, handler));
676                }
677            }
678        }
679
680        if this.handlers.is_empty() && this.incoming.is_done() {
681            trace!(target: "fork::multi", "completed");
682            return Poll::Ready(());
683        }
684
685        // Periodically flush cached RPC state.
686        if this
687            .flush_cache_interval
688            .as_mut()
689            .map(|interval| interval.poll_tick(cx).is_ready())
690            .unwrap_or_default()
691            && !this.forks.is_empty()
692        {
693            trace!(target: "fork::multi", "tick flushing caches");
694            // Only the registry owning a handler flushes its cache, once per backend.
695            let forks = this
696                .handlers
697                .iter()
698                .filter_map(|(id, _)| this.forks.get(id))
699                .map(|fork| fork.backend.clone())
700                .collect::<Vec<_>>();
701            // Flush this on new thread to not block here.
702            std::thread::Builder::new()
703                .name("flusher".into())
704                .spawn(move || {
705                    for fork in forks {
706                        fork.flush_cache();
707                    }
708                })
709                .expect("failed to spawn thread");
710        }
711
712        Poll::Pending
713    }
714}
715
716/// Tracks the created Fork
717#[derive(Debug, Clone)]
718struct CreatedFork<N: Network, SPEC, BLOCK: ForkBlockEnv> {
719    /// How the fork was initially created.
720    opts: CreateFork,
721    /// The immutable remote fork owned by this backend.
722    fork: Fork,
723    /// The resolved EVM environment (fetched from the provider).
724    evm_env: EvmEnv<SPEC, BLOCK>,
725    /// Copy of the sender.
726    backend: SharedBackend<N, BLOCK>,
727    /// How many consumers there are, since a `SharedBacked` can be used by multiple
728    /// consumers.
729    num_senders: Arc<AtomicUsize>,
730    /// Whether this exact shared backend has successfully received a validated parent BAL.
731    bal_prewarmed: Arc<AtomicBool>,
732}
733
734impl<N: Network, SPEC, BLOCK: ForkBlockEnv> CreatedFork<N, SPEC, BLOCK> {
735    pub fn new(
736        opts: CreateFork,
737        fork: Fork,
738        evm_env: EvmEnv<SPEC, BLOCK>,
739        backend: SharedBackend<N, BLOCK>,
740    ) -> Self {
741        Self {
742            opts,
743            fork,
744            evm_env,
745            backend,
746            num_senders: Arc::new(AtomicUsize::new(1)),
747            bal_prewarmed: Arc::default(),
748        }
749    }
750
751    /// Increment senders and return unique identifier of the fork.
752    fn inc_senders(&self, fork_id: ForkId) -> ForkId {
753        format!(
754            "{}-{}",
755            fork_id.as_str(),
756            self.num_senders.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
757        )
758        .into()
759    }
760}
761
762/// A type that's used to signaling the `MultiForkHandler` when it's time to shut down.
763///
764/// This is essentially a sync on drop, so that the `MultiForkHandler` can flush all rpc cashes.
765///
766/// This type intentionally does not implement `Clone` since it's intended that there's only once
767/// instance.
768#[derive(Debug)]
769struct ShutDownMultiFork<N: Network, SPEC, BLOCK: ForkBlockEnv> {
770    handler: Option<Sender<Request<N, SPEC, BLOCK>>>,
771    // Keep shared backend handlers alive until this registry has shut down.
772    _parent: Option<Arc<Self>>,
773}
774
775impl<N: Network, SPEC, BLOCK: ForkBlockEnv> Drop for ShutDownMultiFork<N, SPEC, BLOCK> {
776    fn drop(&mut self) {
777        trace!(target: "fork::multi", "initiating shutdown");
778        let (sender, rx) = oneshot_channel();
779        let req = Request::ShutDown(sender);
780        if let Some(mut handler) = self.handler.take()
781            && handler.try_send(req).is_ok()
782        {
783            let _ = rx.recv();
784            trace!(target: "fork::cache", "multifork backend shutdown");
785        }
786    }
787}
788
789/// Creates a new fork.
790///
791/// This resolves the fork block and environment through the endpoint. The handler creates the
792/// backend with [`create_backend`] only if no backend exists for the fork ID.
793async fn create_fork<
794    SPEC: Into<SpecId> + Default + Copy + Send,
795    BLOCK: FoundryBlock + ForkBlockEnv + Default,
796>(
797    mut fork: CreateFork,
798    expected_identity: Option<ForkContext>,
799    prewarm_bal: bool,
800    exact_block: Option<(Fork, BlockNumHash)>,
801) -> eyre::Result<(ForkId, CreateFork, Fork, EvmEnv<SPEC, BLOCK>, Option<BlockAccessList>)> {
802    // Ensure evm_opts reflects the fork URL (may differ from the resolved CreateFork url when
803    // created via cheatcodes, where evm_opts is cloned from the base config).
804    let execution_networks = fork.evm_opts.networks;
805    let require_endpoint_family_match =
806        fork.evm_opts.fork_network_is_inferred || !execution_networks.has_network_selection();
807    let targets_new_endpoint =
808        fork.evm_opts.fork_url.as_ref().is_some_and(|endpoint| endpoint != &fork.url)
809            || fork
810                .evm_opts
811                .fork_endpoint
812                .as_ref()
813                .is_some_and(|identity| identity.endpoint != fork.url);
814    if targets_new_endpoint {
815        // The EVM implementation is already fixed, so use its family as the fallback for a custom
816        // endpoint without metadata. Clear identity and chain values inferred from the old URL;
817        // authoritative metadata from the new endpoint is still checked below.
818        fork.evm_opts.fork_endpoint = None;
819        fork.evm_opts.expected_fork_endpoint = None;
820        fork.evm_opts.fork_network_is_inferred = false;
821        if fork.evm_opts.fork_chain_id_is_inferred {
822            fork.evm_opts.env.chain_id = None;
823            fork.evm_opts.fork_chain_id_is_inferred = false;
824        }
825        if fork.evm_opts.fork_block_number_is_inferred {
826            fork.evm_opts.fork_block_number = None;
827            fork.evm_opts.fork_block_number_is_inferred = false;
828        }
829    }
830    fork.evm_opts.fork_url = Some(fork.url.clone());
831
832    // Initialise the fork environment.
833    // Here we use [`AnyNetwork`] to maximize compatibility with custom chains, aligned with
834    // `EvmOpts::env` impl.
835    let resolved = if let Some((parent, block)) = exact_block {
836        let provider = parent.provider::<AnyNetwork>();
837        fork.evm_opts.check_fork_endpoint(&provider, &parent).await?;
838        let fork_state = parent.at_block(block).await?;
839        fork.evm_opts.check_fork_endpoint(&provider, &fork_state).await?;
840        fork_state
841    } else {
842        fork.evm_opts.prepare_fork().await?.expect("fork URL is configured")
843    };
844    let any_provider = resolved.provider::<AnyNetwork>();
845    let fork_context = resolved.context();
846    let evm_env = fork.evm_opts.fork_env_from_block::<SPEC, BLOCK, AnyNetwork>(
847        fork.evm_opts.chain_id_override().unwrap_or(fork_context.execution_chain_id),
848        fork_context.source_chain_id,
849        &resolved.block,
850    );
851    if require_endpoint_family_match
852        && !execution_networks.supports_fork_source(&fork_context.network_profile)
853    {
854        eyre::bail!(
855            "cannot create a `{}` fork with an EVM instantiated for `{}`; run the script with --rpc-url pointing to the forked chain to select its EVM",
856            fork_context.network,
857            execution_networks.execution_network()
858        );
859    }
860    if let Some(expected) = expected_identity {
861        eyre::ensure!(
862            fork_context.has_same_endpoint_identity(expected),
863            "fork endpoint identity changed while the fork was being rolled"
864        );
865    }
866    let bal = if prewarm_bal {
867        bal::prepare(&any_provider, &resolved, &resolved.block).await
868    } else {
869        None
870    };
871    let fork_id = ForkId::resolved(&fork.url, &resolved);
872
873    Ok((fork_id, fork, resolved, evm_env, bal))
874}
875
876/// Creates the backend of a resolved fork.
877///
878/// This loads the disk cache of the fork block, so it only runs for a new fork ID.
879fn create_backend<N: Network, BLOCK: FoundryBlock + ForkBlockEnv>(
880    fork: &CreateFork,
881    resolved: &Fork,
882    block_env: &BLOCK,
883) -> eyre::Result<(SharedBackend<N, BLOCK>, BackendHandler<N, BLOCK>)> {
884    let fork_context = resolved.context();
885    let number = resolved.number();
886    let account_fetch_policy = crate::backend::account_fetch_policy_for_source(
887        fork_context.source_chain_id,
888        fork_context.network_profile,
889    );
890    let meta = BlockchainDbMeta::new(block_env.clone(), fork.url.clone())
891        .with_fork_identity(resolved.hash(), resolved.source_id())
892        .with_account_fetch_policy(account_fetch_policy);
893
894    // Determine the cache path if caching is enabled.
895    // Number-addressed state can follow a replacement block, so keep it in memory only.
896    let cache_path = if fork.enable_caching && !resolved.state_by_number {
897        Config::foundry_block_cache_dir(fork_context.source_chain_id, number)
898    } else {
899        None
900    };
901
902    let provider = resolved.provider::<N>();
903    let db = BlockchainDb::new(meta, cache_path);
904    let anchor = ForkBlock::with_rpc_number(
905        block_env.number().saturating_to(),
906        resolved.number(),
907        resolved.hash(),
908    );
909    if resolved.state_by_number {
910        SharedBackend::new_with_anchor_by_number(provider, db, anchor)
911    } else {
912        SharedBackend::new_with_anchor(provider, db, anchor)
913    }
914}
915
916#[cfg(test)]
917mod tests {
918    use super::*;
919    use crate::opts::EvmOpts;
920    use alloy_chains::NamedChain;
921    use alloy_eips::eip7928::{AccountChanges, BlockAccessIndex, SlotChanges, StorageChange};
922    use alloy_network::TransactionBuilder;
923    use alloy_primitives::{Address, B256, U256, bytes};
924    use alloy_provider::{Provider, ProviderBuilder, mock::Asserter};
925    use alloy_rpc_types::TransactionRequest;
926    use alloy_serde::WithOtherFields;
927    use foundry_evm_networks::{NetworkConfigs, NetworkVariant};
928    use foundry_fork_db::AccountFetchPolicy;
929    use foundry_test_utils::rpc::{
930        spawn_rpc_proxy_method_not_found_before, spawn_rpc_proxy_recording_method,
931    };
932    use futures::{channel::oneshot, task::noop_waker_ref};
933    use revm::context::{BlockEnv, TxEnv};
934    use std::sync::mpsc::{Receiver as OneshotReceiver, TryRecvError};
935
936    fn context(block_number: u64) -> ForkContext {
937        ForkContext {
938            execution_chain_id: 1,
939            source_chain_id: 1,
940            network: NetworkVariant::Ethereum,
941            network_profile: NetworkConfigs::default(),
942            block_number,
943            hardfork: None,
944            instance_id: None,
945            source_fork_block_number: None,
946            source_fork_block_hash: None,
947        }
948    }
949
950    #[test]
951    fn resolved_fork_ids_include_hash_and_source_identity() {
952        let url = "http://localhost:8545";
953        let first = Fork::test(
954            url,
955            None,
956            None,
957            Some(1),
958            BlockNumHash::new(1, B256::with_last_byte(1)),
959            context(1),
960        );
961        let replacement = Fork::test(
962            url,
963            None,
964            None,
965            Some(1),
966            BlockNumHash::new(1, B256::with_last_byte(2)),
967            context(1),
968        );
969        let authenticated = Fork::test(
970            url,
971            Some(&["Authorization: secret".to_string()]),
972            None,
973            Some(1),
974            BlockNumHash::new(1, B256::with_last_byte(1)),
975            context(1),
976        );
977
978        assert_ne!(ForkId::resolved(url, &first), ForkId::resolved(url, &replacement));
979        assert_ne!(ForkId::resolved(url, &first), ForkId::resolved(url, &authenticated));
980        let mut numbered = first.clone();
981        numbered.state_by_number = true;
982        assert_ne!(ForkId::resolved(url, &first), ForkId::resolved(url, &numbered));
983        assert_eq!(ForkId::resolved(url, &first), ForkId::exact(url, &numbered, first.block()));
984    }
985
986    #[tokio::test(flavor = "multi_thread")]
987    async fn fork_state_by_number_does_not_load_or_overwrite_hash_cache() {
988        let chain_id = u64::from_be_bytes(B256::random()[..8].try_into().unwrap());
989        let (_api, handle) = anvil::spawn(
990            anvil::NodeConfig::test().with_chain_id(Some(chain_id)).with_no_mining(true),
991        )
992        .await;
993        let opts = EvmOpts {
994            fork_url: Some(handle.http_endpoint()),
995            fork_block_number: Some(0),
996            ..Default::default()
997        };
998        let resolved = opts.prepare_fork().await.unwrap().unwrap();
999        let (env, _) = opts.env_at_fork::<SpecId, BlockEnv, TxEnv>(Some(&resolved)).await.unwrap();
1000        let path = Config::foundry_block_cache_dir(chain_id, resolved.number()).unwrap();
1001        assert!(!path.exists());
1002        let meta = BlockchainDbMeta::new(env.block_env, handle.http_endpoint())
1003            .with_fork_identity(resolved.hash(), resolved.source_id());
1004        let cached = BlockchainDb::new(meta.clone(), Some(path.clone()));
1005        let address = Address::with_last_byte(0x42);
1006        cached.storage().write().entry(address).or_default().insert(U256::ZERO, U256::from(42));
1007        cached.cache().flush();
1008        assert!(path.exists());
1009
1010        for state_by_number in [true, false] {
1011            let fork = CreateFork {
1012                url: handle.http_endpoint(),
1013                enable_caching: true,
1014                evm_opts: EvmOpts { fork_state_by_number: state_by_number, ..opts.clone() },
1015            };
1016            let (_, fork, fork_state, evm_env, _) =
1017                create_fork::<SpecId, BlockEnv>(fork, None, false, None).await.unwrap();
1018            let (backend, handler) =
1019                create_backend::<AnyNetwork, _>(&fork, &fork_state, &evm_env.block_env).unwrap();
1020            assert_eq!(backend.data().storage.read().contains_key(&address), !state_by_number);
1021            if state_by_number {
1022                backend
1023                    .data()
1024                    .storage
1025                    .write()
1026                    .entry(address)
1027                    .or_default()
1028                    .insert(U256::ZERO, U256::from(99));
1029            }
1030            backend.flush_cache();
1031            drop(backend);
1032            drop(handler);
1033            let reloaded = BlockchainDb::new(meta.clone(), Some(path.clone()));
1034            assert_eq!(reloaded.storage().read()[&address][&U256::ZERO], U256::from(42));
1035        }
1036        std::fs::remove_file(&path).unwrap();
1037        std::fs::remove_dir(path.parent().unwrap()).unwrap();
1038    }
1039
1040    #[test]
1041    fn account_fetch_policy_follows_source_identity() {
1042        assert_eq!(
1043            crate::backend::account_fetch_policy_for_source(
1044                NamedChain::Tempo as u64,
1045                NetworkConfigs::with_ethereum(),
1046            ),
1047            AccountFetchPolicy::RequireAccountInfo,
1048        );
1049        assert_eq!(
1050            crate::backend::account_fetch_policy_for_source(
1051                NamedChain::Mainnet as u64,
1052                NetworkConfigs::with_tempo(),
1053            ),
1054            AccountFetchPolicy::Auto,
1055        );
1056        assert_eq!(
1057            crate::backend::account_fetch_policy_for_source(123_456, NetworkConfigs::with_tempo(),),
1058            AccountFetchPolicy::RequireAccountInfo,
1059        );
1060    }
1061
1062    #[test]
1063    fn fork_bal_pending_requests_preserve_opt_out() {
1064        let url = "http://localhost:8545";
1065        let resolved = Fork::test(
1066            url,
1067            None,
1068            None,
1069            Some(1),
1070            BlockNumHash::new(1, B256::with_last_byte(1)),
1071            context(1),
1072        );
1073        let (_, receiver) = channel(1);
1074        let mut handler =
1075            MultiForkHandler::<AnyNetwork, SpecId, revm::context::BlockEnv>::new(receiver);
1076        let mut fork = CreateFork {
1077            enable_caching: false,
1078            url: url.to_string(),
1079            evm_opts: Default::default(),
1080        };
1081        let (sender, _) = oneshot_channel();
1082        handler.create_fork_with_identity(
1083            fork.clone(),
1084            None,
1085            false,
1086            Some((resolved.clone(), resolved.block())),
1087            sender,
1088        );
1089        fork.evm_opts.no_fork_bal = true;
1090        let (sender, _) = oneshot_channel();
1091        handler.create_fork_with_identity(
1092            fork.clone(),
1093            None,
1094            false,
1095            Some((resolved.clone(), resolved.block())),
1096            sender,
1097        );
1098        assert_eq!(handler.pending_tasks.len(), 2, "each fork must retain its opt-out policy");
1099        let (sender, _) = oneshot_channel();
1100        handler.create_fork_with_identity(
1101            fork.clone(),
1102            None,
1103            false,
1104            Some((resolved.clone(), resolved.block())),
1105            sender,
1106        );
1107        assert_eq!(handler.pending_tasks.len(), 2, "identical policies can share creation");
1108        fork.evm_opts.no_fork_bal = false;
1109        let (sender, _) = oneshot_channel();
1110        handler.create_fork_with_identity(
1111            fork.clone(),
1112            None,
1113            true,
1114            Some((resolved.clone(), resolved.block())),
1115            sender,
1116        );
1117        assert_eq!(handler.pending_tasks.len(), 3, "prewarming must not join ordinary creation");
1118        let (sender, _) = oneshot_channel();
1119        handler.create_fork_with_identity(
1120            fork,
1121            None,
1122            true,
1123            Some((resolved.clone(), resolved.block())),
1124            sender,
1125        );
1126        assert_eq!(handler.pending_tasks.len(), 3, "matching prewarm requests can share creation");
1127    }
1128
1129    #[test]
1130    fn fork_bal_fills_selected_cache_for_equivalent_profiles() {
1131        let url = "http://localhost:8545";
1132        let block = BlockNumHash::new(1, B256::with_last_byte(1));
1133        let address = Address::with_last_byte(1);
1134        let create = |resolved: Fork| {
1135            let provider = ProviderBuilder::<_, _, AnyNetwork>::default()
1136                .connect_mocked_client(Asserter::new());
1137            let env = EvmEnv::<SpecId>::default();
1138            let meta = BlockchainDbMeta::new(env.block_env.clone(), url.to_string())
1139                .with_fork_identity(resolved.hash(), resolved.source_id());
1140            let (backend, handler) = SharedBackend::new_with_anchor(
1141                provider,
1142                BlockchainDb::new(meta, None),
1143                ForkBlock::with_rpc_number(1, 1, resolved.hash()),
1144            )
1145            .unwrap();
1146            let opts = CreateFork {
1147                enable_caching: false,
1148                url: url.to_string(),
1149                evm_opts: Default::default(),
1150            };
1151            (CreatedFork::new(opts, resolved, env, backend), handler)
1152        };
1153
1154        for profile in [NetworkConfigs::default(), NetworkConfigs::with_ethereum()] {
1155            let cached = Fork::test(url, None, None, None, block, context(1));
1156            let candidate = Fork::test(
1157                url,
1158                None,
1159                None,
1160                Some(1),
1161                block,
1162                ForkContext { network_profile: profile, ..context(1) },
1163            );
1164            let id = ForkId::resolved(url, &candidate);
1165            assert_eq!(id, ForkId::resolved(url, &cached));
1166            assert_eq!(candidate.fingerprint(), cached.fingerprint());
1167            let (cached, cached_handler) = create(cached);
1168            let cached_db = cached.backend.data();
1169            let bal = vec![AccountChanges::new(address).with_storage_change(SlotChanges::new(
1170                U256::ONE,
1171                vec![StorageChange::new(BlockAccessIndex::new(1), U256::from(42))],
1172            ))];
1173            let future = futures::future::ready(Ok((
1174                id.clone(),
1175                cached.opts.clone(),
1176                candidate,
1177                cached.evm_env.clone(),
1178                Some(bal),
1179            )));
1180            let (_incoming, receiver) = channel(1);
1181            let mut manager = MultiForkHandler::<AnyNetwork, SpecId, BlockEnv>::new(receiver);
1182            manager.forks.insert(id.clone(), cached);
1183            manager.handlers.push((id.clone(), cached_handler));
1184            let (sender, receiver) = oneshot_channel();
1185            manager.pending_tasks.push(ForkTask::Create {
1186                future: Box::pin(future),
1187                id: id.clone(),
1188                prewarm_bal: true,
1189                no_fork_bal: false,
1190                sender,
1191                additional_senders: Vec::new(),
1192            });
1193
1194            assert!(manager.poll_unpin(&mut Context::from_waker(noop_waker_ref())).is_pending());
1195            let result = receiver.try_recv().unwrap().unwrap();
1196            assert!(Arc::ptr_eq(&result.backend.data(), &cached_db));
1197            // Reusing the cached backend creates no second backend.
1198            assert_eq!(manager.handlers.len(), 1);
1199            assert_eq!(
1200                cached_db
1201                    .storage
1202                    .read()
1203                    .get(&address)
1204                    .and_then(|slots| slots.get(&U256::ONE))
1205                    .copied(),
1206                Some(U256::from(42)),
1207            );
1208            assert!(manager.forks[&id].bal_prewarmed.load(Ordering::Relaxed));
1209        }
1210    }
1211
1212    #[tokio::test(flavor = "multi_thread")]
1213    async fn fork_bal_reuses_prewarmed_cache_only_for_equivalent_identities() {
1214        let (api, handle) = anvil::spawn(
1215            anvil::NodeConfig::test()
1216                .with_chain_id(Some(1u64))
1217                .with_hardfork(Some(anvil::EthereumHardfork::Amsterdam.into()))
1218                .with_genesis_timestamp(Some(1_800_000_000u64))
1219                .with_no_mining(true),
1220        )
1221        .await;
1222        let address = Address::with_last_byte(0x42);
1223        api.anvil_set_code(address, bytes!("602a60015500")).await.unwrap();
1224        api.send_transaction(WithOtherFields::new(
1225            TransactionRequest::default()
1226                .with_from(handle.dev_accounts().next().unwrap())
1227                .with_to(address)
1228                .with_nonce(0)
1229                // Enough for the EIP-8037 state gas of creating slot one.
1230                .with_gas_limit(1_000_000)
1231                .with_gas_price(2_000_000_000),
1232        ))
1233        .await
1234        .unwrap();
1235        api.mine_one().await.unwrap();
1236        let block = handle.http_provider().get_block_by_number(1.into()).await.unwrap().unwrap();
1237        let block = BlockNumHash::new(1, block.header.hash);
1238        // Expose native BALs through an endpoint without mutable Anvil identity.
1239        let endpoint = spawn_rpc_proxy_method_not_found_before(
1240            handle.http_endpoint(),
1241            "anvil_nodeInfo",
1242            usize::MAX,
1243        )
1244        .await;
1245        let (endpoint, bal_requests) =
1246            spawn_rpc_proxy_recording_method(endpoint, "eth_getBlockAccessList").await;
1247        let (endpoint, probes) = spawn_rpc_proxy_recording_method(endpoint, "anvil_nodeInfo").await;
1248        let (_incoming, receiver) = channel(1);
1249        let mut manager = MultiForkHandler::<AnyNetwork, SpecId, BlockEnv>::new(receiver);
1250        let mut first = None;
1251        let mut probe_counts = Vec::new();
1252        let other_headers = ["X-Test-Source: other".to_string()];
1253
1254        for (profile, headers) in [
1255            (NetworkConfigs::default(), None),
1256            (NetworkConfigs::with_ethereum(), None),
1257            (NetworkConfigs::default(), None),
1258            (NetworkConfigs::default(), Some(other_headers.as_slice())),
1259        ] {
1260            let evm_opts = EvmOpts {
1261                fork_url: Some(endpoint.clone()),
1262                fork_block_number: Some(1),
1263                fork_headers: headers.map(<[_]>::to_vec),
1264                networks: profile,
1265                ..Default::default()
1266            };
1267            let resolved = evm_opts.prepare_fork().await.unwrap().unwrap();
1268            assert_eq!(resolved.block(), block);
1269            let fingerprint = resolved.fingerprint();
1270            let fork = CreateFork { enable_caching: false, url: endpoint.clone(), evm_opts };
1271            let before = probes.lock().unwrap().len();
1272            let (sender, receiver) = oneshot_channel();
1273            manager.create_fork_with_identity(
1274                fork,
1275                None,
1276                true,
1277                Some((resolved.clone(), resolved.block())),
1278                sender,
1279            );
1280            let result = tokio::time::timeout(
1281                Duration::from_secs(10),
1282                futures::future::poll_fn(|cx| {
1283                    assert!(manager.poll_unpin(cx).is_pending());
1284                    match receiver.try_recv() {
1285                        Ok(result) => Poll::Ready(result),
1286                        Err(std::sync::mpsc::TryRecvError::Empty) => Poll::Pending,
1287                        Err(error) => panic!("fork response channel closed: {error}"),
1288                    }
1289                }),
1290            )
1291            .await
1292            .unwrap()
1293            .unwrap();
1294
1295            let db = result.backend.data();
1296            assert_eq!(
1297                Arc::ptr_eq(first.get_or_insert_with(|| db.clone()), &db),
1298                headers.is_none(),
1299            );
1300            assert_eq!(result.fork.fingerprint(), fingerprint);
1301            assert_eq!(db.storage.read()[&address][&U256::ONE], U256::from(42));
1302            assert!(manager.forks[&result.id].bal_prewarmed.load(Ordering::Relaxed));
1303            assert_eq!(bal_requests.lock().unwrap().len(), if headers.is_some() { 2 } else { 1 });
1304            probe_counts.push(probes.lock().unwrap().len() - before);
1305        }
1306
1307        // Only cold caches need the two source probes surrounding BAL preparation.
1308        assert_eq!(probe_counts[0], probe_counts[1] + 2);
1309        assert_eq!(probe_counts[1], probe_counts[2]);
1310        assert_eq!(probe_counts[0], probe_counts[3]);
1311    }
1312
1313    #[tokio::test(flavor = "multi_thread")]
1314    async fn fork_bal_shared_backend_preserves_opposite_policies_on_roll() {
1315        let (api, handle) = anvil::spawn(
1316            anvil::NodeConfig::test()
1317                .with_chain_id(Some(1u64))
1318                .with_hardfork(Some(anvil::EthereumHardfork::Amsterdam.into()))
1319                .with_genesis_timestamp(Some(1_800_000_000u64))
1320                .with_no_mining(true),
1321        )
1322        .await;
1323        let mut blocks = Vec::new();
1324        for number in 1..=3 {
1325            api.mine_one().await.unwrap();
1326            let block =
1327                handle.http_provider().get_block_by_number(number.into()).await.unwrap().unwrap();
1328            blocks.push(BlockNumHash::new(number, block.header.hash));
1329        }
1330        // Expose native BALs through an endpoint without mutable Anvil identity.
1331        let endpoint = spawn_rpc_proxy_method_not_found_before(
1332            handle.http_endpoint(),
1333            "anvil_nodeInfo",
1334            usize::MAX,
1335        )
1336        .await;
1337        let (endpoint, bal_requests) =
1338            spawn_rpc_proxy_recording_method(endpoint, "eth_getBlockAccessList").await;
1339        let resolved = EvmOpts {
1340            fork_url: Some(endpoint.clone()),
1341            fork_block_number: Some(1),
1342            ..Default::default()
1343        }
1344        .prepare_fork()
1345        .await
1346        .unwrap()
1347        .unwrap();
1348        assert_eq!(resolved.block(), blocks[0]);
1349
1350        for disabled_first in [false, true] {
1351            let (_incoming, receiver) = channel(1);
1352            let mut manager = MultiForkHandler::<AnyNetwork, SpecId, BlockEnv>::new(receiver);
1353            bal_requests.lock().unwrap().clear();
1354            let mut requests = [false, true].map(|no_fork_bal| {
1355                let fork = CreateFork {
1356                    enable_caching: false,
1357                    url: endpoint.clone(),
1358                    evm_opts: EvmOpts {
1359                        fork_url: Some(endpoint.clone()),
1360                        fork_block_number: Some(1),
1361                        no_fork_bal,
1362                        ..Default::default()
1363                    },
1364                };
1365                let (sender, receiver) = oneshot_channel();
1366                manager.create_fork(fork, sender);
1367                let ForkTask::Create { future, .. } = manager.pending_tasks.last_mut().unwrap();
1368                let original = std::mem::replace(future, Box::pin(futures::future::pending()));
1369                let (release, gate) = oneshot::channel();
1370                // Run real creation, but explicitly control which result reaches backend reuse.
1371                *future = Box::pin(async move {
1372                    let result = original.await;
1373                    gate.await.unwrap();
1374                    result
1375                });
1376                (no_fork_bal, release, receiver)
1377            });
1378            assert_eq!(manager.pending_tasks.len(), 2);
1379            if disabled_first {
1380                requests.reverse();
1381            }
1382            let [
1383                (first_policy, first_gate, first_receiver),
1384                (second_policy, second_gate, second_receiver),
1385            ] = requests;
1386            first_gate.send(()).unwrap();
1387            let first = complete_fork(&mut manager, first_receiver).await;
1388            assert!(matches!(second_receiver.try_recv(), Err(TryRecvError::Empty)));
1389            assert_eq!(manager.pending_tasks.len(), 1);
1390            second_gate.send(()).unwrap();
1391            let second = complete_fork(&mut manager, second_receiver).await;
1392
1393            assert_ne!(first.id, second.id);
1394            assert!(Arc::ptr_eq(&first.backend.data(), &second.backend.data()));
1395            assert!(Arc::ptr_eq(first.fork.client.inner(), second.fork.client.inner()));
1396            assert!(Arc::ptr_eq(&first.fork.block, &second.fork.block));
1397            assert_eq!(manager.forks[&first.id].opts.evm_opts.no_fork_bal, first_policy);
1398            assert_eq!(manager.forks[&second.id].opts.evm_opts.no_fork_bal, second_policy);
1399            assert!(bal_requests.lock().unwrap().is_empty());
1400
1401            // Each roll targets a cold block so shared prewarming cannot mask a lost policy.
1402            for (no_fork_bal, fork, block) in
1403                [(first_policy, first, blocks[1]), (second_policy, second, blocks[2])]
1404            {
1405                bal_requests.lock().unwrap().clear();
1406                let (sender, receiver) = oneshot_channel();
1407                manager.on_request(Request::RollForkExact(fork.id, block, true, sender));
1408                let rolled = complete_fork(&mut manager, receiver).await;
1409                assert_eq!(rolled.fork.block(), block);
1410                assert_eq!(manager.forks[&rolled.id].opts.evm_opts.no_fork_bal, no_fork_bal);
1411                assert_eq!(
1412                    manager.forks[&rolled.id].bal_prewarmed.load(Ordering::Relaxed),
1413                    !no_fork_bal,
1414                );
1415                let expected =
1416                    if no_fork_bal { vec![] } else { vec![serde_json::json!([block.hash])] };
1417                assert_eq!(*bal_requests.lock().unwrap(), expected);
1418            }
1419        }
1420    }
1421
1422    async fn complete_fork(
1423        manager: &mut MultiForkHandler<AnyNetwork, SpecId, BlockEnv>,
1424        receiver: OneshotReceiver<eyre::Result<ForkResult<AnyNetwork, SpecId, BlockEnv>>>,
1425    ) -> ForkResult<AnyNetwork, SpecId, BlockEnv> {
1426        tokio::time::timeout(
1427            Duration::from_secs(10),
1428            futures::future::poll_fn(|cx| {
1429                assert!(manager.poll_unpin(cx).is_pending());
1430                match receiver.try_recv() {
1431                    Ok(result) => Poll::Ready(result),
1432                    Err(TryRecvError::Empty) => Poll::Pending,
1433                    Err(error) => panic!("fork response channel closed: {error}"),
1434                }
1435            }),
1436        )
1437        .await
1438        .unwrap()
1439        .unwrap()
1440    }
1441}