1use 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#[derive(Clone, Debug, PartialEq, Eq, Hash)]
37pub struct ForkId(pub String);
38
39impl ForkId {
40 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 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 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 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
108pub struct ForkResult<N: Network, SPEC, BLOCK: ForkBlockEnv> {
110 pub id: ForkId,
112 pub backend: SharedBackend<N, BLOCK>,
114 pub env: EvmEnv<SPEC, BLOCK>,
116 pub fork: Fork,
118}
119
120#[derive(Clone, Debug)]
123#[must_use]
124pub struct MultiFork<N: Network, SPEC, BLOCK: ForkBlockEnv> {
125 handler: Sender<Request<N, SPEC, BLOCK>>,
127 _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 pub fn spawn() -> Self {
139 Self::spawn_with_forks(HashMap::default(), None)
140 }
141
142 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 let fut = async move {
164 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 #[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 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 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 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 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 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 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 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 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 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#[derive(Debug)]
318enum Request<N: Network, SPEC, BLOCK: ForkBlockEnv> {
319 CloneForks(OneshotSender<HashMap<ForkId, CreatedFork<N, SPEC, BLOCK>>>),
320 CreateFork(Box<CreateFork>, CreateSender<N, SPEC, BLOCK>),
322 GetFork(ForkId, OneshotSender<Option<SharedBackend<N, BLOCK>>>),
324 RollFork(ForkId, u64, CreateSender<N, SPEC, BLOCK>),
326 RollForkExact(ForkId, BlockNumHash, bool, CreateSender<N, SPEC, BLOCK>),
328 GetEvmEnv(ForkId, GetEvmEnvSender<SPEC, BLOCK>),
330 ShutDown(OneshotSender<()>),
332 GetForkUrl(ForkId, OneshotSender<Option<String>>),
334 GetForkInfo(ForkId, OneshotSender<Option<Fork>>),
335 GetForkOptions(ForkId, OneshotSender<Option<CreateFork>>),
337}
338
339enum ForkTask<N: Network, SPEC, BLOCK: ForkBlockEnv> {
340 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#[must_use = "futures do nothing unless polled"]
353pub struct MultiForkHandler<N: Network, SPEC, BLOCK: ForkBlockEnv> {
354 incoming: Fuse<Receiver<Request<N, SPEC, BLOCK>>>,
356
357 handlers: Vec<(ForkId, BackendHandler<N, BLOCK>)>,
361
362 pending_tasks: Vec<ForkTask<N, SPEC, BLOCK>>,
364
365 forks: HashMap<ForkId, CreatedFork<N, SPEC, BLOCK>>,
370
371 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 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 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 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 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 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 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
566impl<
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 loop {
581 match this.incoming.poll_next_unpin(cx) {
582 Poll::Ready(Some(req)) => this.on_request(req),
583 Poll::Ready(None) => {
584 trace!(target: "fork::multi", "request channel closed");
586 break;
587 }
588 Poll::Pending => break,
589 }
590 }
591
592 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 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 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 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 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 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 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 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#[derive(Debug, Clone)]
718struct CreatedFork<N: Network, SPEC, BLOCK: ForkBlockEnv> {
719 opts: CreateFork,
721 fork: Fork,
723 evm_env: EvmEnv<SPEC, BLOCK>,
725 backend: SharedBackend<N, BLOCK>,
727 num_senders: Arc<AtomicUsize>,
730 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 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#[derive(Debug)]
769struct ShutDownMultiFork<N: Network, SPEC, BLOCK: ForkBlockEnv> {
770 handler: Option<Sender<Request<N, SPEC, BLOCK>>>,
771 _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
789async 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 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 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 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
876fn 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 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 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 .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 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 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 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 *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 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}