1use super::corpus_io::{
38 CorpusDirEntry, canonical_replay_dirs, read_corpus_dir, read_corpus_dir_strict,
39};
40use crate::{
41 executors::{Executor, RawCallResult, invariant::execute_tx},
42 inspectors::{CmpOperands, EdgeIndexMap, MAX_EDGE_COUNT},
43};
44use alloy_dyn_abi::JsonAbiExt;
45use alloy_json_abi::Function;
46use alloy_primitives::{Address, Bytes, I256, U256};
47use eyre::{Result, eyre};
48use foundry_common::{ContractsByAddress, ContractsByArtifact, TestFunctionExt, sh_warn};
49use foundry_config::FuzzCorpusConfig;
50use foundry_evm_core::{constants::CALLER, evm::FoundryEvmNetwork, utils::StateChangeset};
51use foundry_evm_fuzz::{
52 BasicTxDetails, CallDetails, ObservedCall,
53 invariant::{
54 ArtifactFilters, FuzzRunIdentifiedContracts, InvariantContract, TargetedContracts,
55 },
56 sequence::{ComparisonHint, CorpusEntryView, SequenceGenerator, SequencePlan},
57};
58use proptest::test_runner::TestRunner;
59use serde::{Deserialize, Serialize};
60use std::{
61 collections::HashSet,
62 fmt,
63 path::{Path, PathBuf},
64 sync::{
65 Arc,
66 atomic::{AtomicBool, AtomicUsize, Ordering},
67 },
68 time::{Duration, Instant, SystemTime, UNIX_EPOCH},
69};
70use uuid::Uuid;
71
72const WORKER: &str = "worker";
73const CORPUS_DIR: &str = "corpus";
74const SYNC_DIR: &str = "sync";
75const OPTIMIZATION_BEST_FILE: &str = "optimization_best.json";
76
77const FAVORABILITY_THRESHOLD: f64 = 0.3;
78
79const GZIP_THRESHOLD: usize = 4 * 1024;
82
83#[derive(Clone, Serialize, Deserialize)]
85struct OptimizationState {
86 best_value: I256,
87 best_sequence: Vec<BasicTxDetails>,
88}
89
90#[derive(Clone, Serialize)]
92struct CorpusEntry {
93 uuid: Uuid,
95 total_mutations: usize,
97 new_finds_produced: usize,
99 #[serde(skip_serializing)]
101 tx_seq: Vec<BasicTxDetails>,
102 #[serde(skip_serializing)]
105 cmp_seq: Vec<Vec<ComparisonHint>>,
106 is_favored: bool,
109 #[serde(skip_serializing)]
111 timestamp: u64,
112 #[serde(skip_serializing)]
114 persisted_file_name: Option<String>,
115}
116
117impl CorpusEntry {
118 pub fn new(tx_seq: Vec<BasicTxDetails>) -> Self {
120 Self::new_with_cmp(tx_seq, Vec::new(), Uuid::new_v4())
121 }
122
123 pub fn new_with_cmp(
125 tx_seq: Vec<BasicTxDetails>,
126 cmp_seq: Vec<Vec<ComparisonHint>>,
127 uuid: Uuid,
128 ) -> Self {
129 Self {
130 uuid,
131 total_mutations: 0,
132 new_finds_produced: 0,
133 tx_seq,
134 cmp_seq,
135 is_favored: false,
136 timestamp: SystemTime::now()
137 .duration_since(UNIX_EPOCH)
138 .expect("time went backwards")
139 .as_secs(),
140 persisted_file_name: None,
141 }
142 }
143
144 fn write_to_disk_in(&self, dir: &Path, can_gzip: bool) -> foundry_common::fs::Result<PathBuf> {
145 let should_gzip = self.should_gzip(can_gzip);
146 let file_name = self.file_name(should_gzip);
147 let path = dir.join(&file_name);
148 let temp_path = dir.join(format!(".{file_name}.{}.tmp", Uuid::new_v4()));
149
150 let write_result = if should_gzip {
151 foundry_common::fs::write_json_gzip_file(&temp_path, &self.tx_seq)
152 } else {
153 foundry_common::fs::write_json_file(&temp_path, &self.tx_seq)
154 };
155 if let Err(err) = write_result {
156 let _ = foundry_common::fs::remove_file(&temp_path);
157 return Err(err);
158 }
159
160 if let Err(err) = std::fs::rename(&temp_path, &path) {
161 let _ = foundry_common::fs::remove_file(&temp_path);
162 return Err(foundry_common::errors::FsPathError::write(err, &path));
163 }
164
165 Ok(path)
166 }
167
168 fn file_name(&self, gzip: bool) -> String {
169 if let Some(name) = &self.persisted_file_name {
170 return name.clone();
171 }
172 let ext = if gzip { ".json.gz" } else { ".json" };
173 format!("{}-{}{ext}", self.uuid, self.timestamp)
174 }
175
176 fn should_gzip(&self, can_gzip: bool) -> bool {
177 if !can_gzip {
178 return false;
179 }
180 let size: usize = self.tx_seq.iter().map(|tx| tx.estimate_serialized_size()).sum();
181 size > GZIP_THRESHOLD
182 }
183}
184
185pub fn persist_corpus_seed(
187 config: &FuzzCorpusConfig,
188 tx_seq: Vec<BasicTxDetails>,
189) -> foundry_common::fs::Result<Option<PathBuf>> {
190 let Some(root) = &config.corpus_dir else {
191 return Ok(None);
192 };
193 for dir in canonical_replay_dirs(root) {
194 for entry in read_corpus_dir(&dir) {
195 match entry.read_tx_seq() {
196 Ok(existing) if same_tx_sequence(&existing, &tx_seq) => {
197 return Ok(Some(entry.path));
198 }
199 Ok(_) => {}
200 Err(err) => debug!(%err, path = ?entry.path, "failed to read corpus seed"),
201 }
202 }
203 }
204 let corpus_dir = root.join(format!("{WORKER}0")).join(CORPUS_DIR);
205 foundry_common::fs::create_dir_all(&corpus_dir)?;
206 CorpusEntry::new(tx_seq).write_to_disk_in(&corpus_dir, config.corpus_gzip).map(Some)
207}
208
209fn same_tx_sequence(left: &[BasicTxDetails], right: &[BasicTxDetails]) -> bool {
210 left.len() == right.len()
211 && left.iter().zip(right).all(|(left, right)| {
212 left.warp == right.warp
213 && left.roll == right.roll
214 && left.sender == right.sender
215 && left.call_details.target == right.call_details.target
216 && left.call_details.calldata == right.call_details.calldata
217 && left.call_details.value == right.call_details.value
218 })
219}
220
221fn link_corpus_file(from: &Path, to: &Path) -> bool {
222 match std::fs::hard_link(from, to) {
223 Ok(()) => true,
224 Err(err) if err.kind() == std::io::ErrorKind::AlreadyExists => {
225 if !std::fs::symlink_metadata(to).is_ok_and(|metadata| metadata.file_type().is_file()) {
226 debug!(target: "corpus", from=?from, to=?to, "corpus destination is not a regular file");
227 return false;
228 }
229 let entry = |path: &Path| CorpusDirEntry {
230 path: path.to_path_buf(),
231 uuid: Uuid::nil(),
232 timestamp: 0,
233 };
234 let matches = entry(from)
235 .read_tx_seq()
236 .ok()
237 .zip(entry(to).read_tx_seq().ok())
238 .is_some_and(|(from, to)| same_tx_sequence(&from, &to));
239 if !matches {
240 debug!(target: "corpus", from=?from, to=?to, "conflicting corpus file already exists");
241 }
242 matches
243 }
244 Err(err) => {
245 debug!(target: "corpus", %err, from=?from, to=?to, "failed to link corpus file");
246 false
247 }
248 }
249}
250
251fn accept_synced_corpus_file(
252 entry: &CorpusDirEntry,
253 tx_seq: &[BasicTxDetails],
254 corpus_path: &Path,
255) -> bool {
256 if corpus_path.is_file() {
257 let existing = CorpusDirEntry {
258 path: corpus_path.to_path_buf(),
259 uuid: entry.uuid,
260 timestamp: entry.timestamp,
261 };
262 if existing.read_tx_seq().is_ok_and(|existing| same_tx_sequence(&existing, tx_seq)) {
263 if let Err(err) = std::fs::remove_file(&entry.path) {
264 debug!(target: "corpus", %err, "failed to remove synced corpus link {}", entry.path.display());
265 return false;
266 }
267 return true;
268 }
269
270 warn!(target: "corpus", "not overwriting conflicting corpus file {}", corpus_path.display());
271 let quarantine_path =
272 entry.path.with_file_name(format!("{}.{}.invalid", entry.name(), Uuid::new_v4()));
273 if let Err(err) = std::fs::rename(&entry.path, &quarantine_path) {
274 debug!(target: "corpus", %err, "failed to quarantine conflicting corpus file {}", entry.path.display());
275 }
276 return false;
277 }
278
279 if let Err(err) = std::fs::rename(&entry.path, corpus_path) {
280 debug!(target: "corpus", %err, "failed to move synced corpus from {:?} to {corpus_path:?} dir", entry.path);
281 return false;
282 }
283 true
284}
285
286#[derive(Clone, Copy, Debug, PartialEq, Eq)]
287pub(crate) enum CorpusInsertionMode {
288 Live,
289 MemoryOnly,
290}
291
292struct ReplayOutcome {
293 keep_entry: bool,
294 new_coverage: bool,
295 new_edge: bool,
297 cmp_seq: Vec<Vec<ComparisonHint>>,
298 failed_replays: usize,
299}
300
301#[derive(Clone, Copy)]
302pub struct StatelessReplayTarget<'a> {
303 pub function: &'a Function,
304 pub address: Address,
305}
306
307impl StatelessReplayTarget<'_> {
308 fn can_replay(self, tx: &BasicTxDetails) -> bool {
309 tx.call_details.target == self.address
310 && tx
311 .call_details
312 .calldata
313 .get(..4)
314 .is_some_and(|selector| self.function.selector() == selector)
315 }
316}
317
318#[derive(Clone, Copy)]
319pub(crate) struct ReplayTarget<'a> {
320 pub(crate) stateless: Option<StatelessReplayTarget<'a>>,
321 pub(crate) fuzzed_contracts: Option<&'a FuzzRunIdentifiedContracts>,
322 pub(crate) dynamic: Option<&'a DynamicTargetCtx<'a>>,
323}
324
325struct ReplayCoverage<'a> {
326 history_map: &'a mut Vec<u8>,
327 edge_indices: &'a mut EdgeIndexMap,
328 sancov_history_map: &'a mut Vec<u8>,
329 metrics: Option<&'a mut CorpusMetrics>,
330}
331
332#[derive(Clone, Default)]
338pub(crate) struct WorkerCorpusSeed {
339 in_memory_corpus: Vec<CorpusEntry>,
340 history_map: Vec<u8>,
341 edge_indices: EdgeIndexMap,
342 sancov_history_map: Vec<u8>,
343 metrics: CorpusMetrics,
344 replay_dirs: Option<Vec<PathBuf>>,
345 failed_replays: usize,
346 optimization_best_value: Option<I256>,
347 optimization_best_sequence: Vec<BasicTxDetails>,
348 last_new_edge_at: Option<Instant>,
351}
352
353impl WorkerCorpusSeed {
354 fn empty(config: &FuzzCorpusConfig) -> Self {
355 let history_map =
359 if config.collect_evm_edge_coverage() && !config.evm_edge_coverage_collision_free() {
360 vec![0u8; MAX_EDGE_COUNT]
361 } else {
362 Vec::new()
363 };
364 Self { history_map, ..Default::default() }
365 }
366
367 fn with_optimization_state(mut self, config: &FuzzCorpusConfig) -> Self {
368 if let Some((value, sequence)) = load_optimization_state(config) {
369 self.optimization_best_value = Some(value);
370 self.optimization_best_sequence = sequence;
371 }
372 self
373 }
374
375 pub(crate) fn clone_for_worker(
376 &self,
377 worker_id: usize,
378 worker_count: usize,
379 include_cmp_seq: bool,
380 ) -> Self {
381 let in_memory_corpus = self
382 .in_memory_corpus
383 .iter()
384 .enumerate()
385 .filter(|(idx, _)| idx % worker_count == worker_id)
386 .map(|(_, entry)| {
387 let mut entry = entry.clone();
388 if !include_cmp_seq {
389 entry.cmp_seq.clear();
390 }
391 entry
392 })
393 .collect::<Vec<_>>();
394
395 let mut metrics = self.metrics.clone();
396 metrics.corpus_count = in_memory_corpus.len();
397 metrics.favored_items = in_memory_corpus.iter().filter(|entry| entry.is_favored).count();
398
399 Self {
400 in_memory_corpus,
401 history_map: self.history_map.clone(),
402 edge_indices: self.edge_indices.clone(),
403 sancov_history_map: self.sancov_history_map.clone(),
404 metrics,
405 replay_dirs: self.replay_dirs.clone(),
406 failed_replays: self.failed_replays,
407 optimization_best_value: self.optimization_best_value,
408 optimization_best_sequence: self.optimization_best_sequence.clone(),
409 last_new_edge_at: self.last_new_edge_at,
410 }
411 }
412
413 pub(crate) fn retain_replayable(&mut self, targeted_contracts: &TargetedContracts) {
414 let is_replayable =
415 |tx_seq: &[BasicTxDetails]| tx_seq.iter().all(|tx| targeted_contracts.can_replay(tx));
416 self.in_memory_corpus.retain(|entry| is_replayable(&entry.tx_seq));
417 self.metrics.corpus_count = self.in_memory_corpus.len();
418 self.metrics.favored_items =
419 self.in_memory_corpus.iter().filter(|entry| entry.is_favored).count();
420
421 if !self.optimization_best_sequence.is_empty()
422 && !is_replayable(&self.optimization_best_sequence)
423 {
424 self.optimization_best_value = None;
425 self.optimization_best_sequence.clear();
426 }
427 }
428
429 pub(crate) fn load_from_disk<FEN: FoundryEvmNetwork>(
430 config: &FuzzCorpusConfig,
431 replay_root: Option<&Path>,
432 executor: Option<&Executor<FEN>>,
433 target: ReplayTarget<'_>,
434 ) -> Result<Self> {
435 let mut seed = Self::empty(config).with_optimization_state(config);
436 let Some(corpus_dir) = &config.corpus_dir else {
437 return Ok(seed);
438 };
439 let replay_dirs = canonical_replay_dirs(replay_root.unwrap_or(corpus_dir));
440 seed.replay_dirs = Some(replay_dirs.clone());
441
442 if !seed.optimization_best_sequence.is_empty() {
445 seed.in_memory_corpus.push(CorpusEntry::new(seed.optimization_best_sequence.clone()));
446 seed.metrics.corpus_count += 1;
447 }
448
449 if target.fuzzed_contracts.is_some() && has_legacy_invariant_corpus_dirs(corpus_dir) {
450 let _ = sh_warn!(
451 "Ignoring legacy invariant corpus directories under {}; new corpus entries are persisted under the contract-level corpus directory.",
452 corpus_dir.display(),
453 );
454 }
455
456 let Some(executor) = executor else {
457 return Ok(seed);
458 };
459 let mut seen_entries =
460 seed.in_memory_corpus.iter().map(|entry| entry.uuid).collect::<HashSet<_>>();
461 for entry in unique_corpus_entries(&replay_dirs, &mut seen_entries) {
462 let tx_seq = match entry.read_tx_seq() {
466 Ok(tx_seq) => tx_seq,
467 Err(err) => {
468 let _ =
469 sh_warn!("Skipping unreadable corpus file {}: {err}", entry.path.display());
470 continue;
471 }
472 };
473 if tx_seq.is_empty() {
474 continue;
475 }
476
477 let coverage = ReplayCoverage {
478 history_map: &mut seed.history_map,
479 edge_indices: &mut seed.edge_indices,
480 sancov_history_map: &mut seed.sancov_history_map,
481 metrics: Some(&mut seed.metrics),
482 };
483 let ReplayOutcome { keep_entry, new_edge, cmp_seq, failed_replays, .. } =
484 replay_corpus_sequence(&tx_seq, executor, target, coverage)?;
485 seed.failed_replays += failed_replays;
486 if new_edge {
488 seed.last_new_edge_at = Some(Instant::now());
489 }
490 if !keep_entry {
491 continue;
492 }
493
494 seed.metrics.corpus_count += 1;
495 debug!(
496 target: "corpus",
497 "load sequence with len {} from corpus file {}",
498 tx_seq.len(),
499 entry.path.display()
500 );
501 seed.in_memory_corpus.push(CorpusEntry::new_with_cmp(tx_seq, cmp_seq, entry.uuid));
502 }
503
504 Ok(seed)
505 }
506}
507
508#[derive(Default)]
509pub(crate) struct GlobalCorpusMetrics {
510 cumulative_edges_seen: AtomicUsize,
512 cumulative_features_seen: AtomicUsize,
514 corpus_count: AtomicUsize,
516 favored_items: AtomicUsize,
518}
519
520pub(crate) struct CorpusSyncCoordinator {
521 workers: usize,
522 arrived: AtomicUsize,
523 phase: AtomicUsize,
524 aborted: AtomicBool,
525}
526
527impl CorpusSyncCoordinator {
528 pub(crate) const fn new(workers: usize) -> Self {
529 Self {
530 workers,
531 arrived: AtomicUsize::new(0),
532 phase: AtomicUsize::new(0),
533 aborted: AtomicBool::new(false),
534 }
535 }
536
537 pub(crate) fn abort(&self) {
538 self.aborted.store(true, Ordering::Release);
539 }
540
541 fn wait(&self) -> bool {
542 if self.aborted.load(Ordering::Acquire) {
543 return false;
544 }
545
546 let phase = self.phase.load(Ordering::Acquire);
547 if self.arrived.fetch_add(1, Ordering::AcqRel) + 1 == self.workers {
548 self.arrived.store(0, Ordering::Release);
549 self.phase.fetch_add(1, Ordering::AcqRel);
550 return true;
551 }
552
553 while self.phase.load(Ordering::Acquire) == phase && !self.aborted.load(Ordering::Acquire) {
554 rayon::yield_now();
555 }
556 !self.aborted.load(Ordering::Acquire)
557 }
558}
559
560impl fmt::Display for GlobalCorpusMetrics {
561 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
562 self.load().fmt(f)
563 }
564}
565
566impl GlobalCorpusMetrics {
567 pub(crate) fn load(&self) -> CorpusMetrics {
568 CorpusMetrics {
569 cumulative_edges_seen: self.cumulative_edges_seen.load(Ordering::Relaxed),
570 cumulative_features_seen: self.cumulative_features_seen.load(Ordering::Relaxed),
571 corpus_count: self.corpus_count.load(Ordering::Relaxed),
572 favored_items: self.favored_items.load(Ordering::Relaxed),
573 }
574 }
575}
576
577#[derive(Serialize, Default, Clone)]
578pub(crate) struct CorpusMetrics {
579 cumulative_edges_seen: usize,
581 cumulative_features_seen: usize,
583 corpus_count: usize,
585 favored_items: usize,
587}
588
589impl fmt::Display for CorpusMetrics {
590 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
591 writeln!(f)?;
592 writeln!(f, " Edge coverage metrics:")?;
593 writeln!(f, " - cumulative edges seen: {}", self.cumulative_edges_seen)?;
594 writeln!(f, " - cumulative features seen: {}", self.cumulative_features_seen)?;
595 writeln!(f, " - corpus count: {}", self.corpus_count)?;
596 write!(f, " - favored items: {}", self.favored_items)?;
597 Ok(())
598 }
599}
600
601impl CorpusMetrics {
602 pub const fn update_seen(&mut self, is_edge: bool) {
604 if is_edge {
605 self.cumulative_edges_seen += 1;
606 } else {
607 self.cumulative_features_seen += 1;
608 }
609 }
610
611 pub const fn update_favored(&mut self, is_favored: bool, corpus_favored: bool) {
613 if is_favored && !corpus_favored {
614 self.favored_items += 1;
615 } else if !is_favored && corpus_favored {
616 self.favored_items -= 1;
617 }
618 }
619}
620
621pub struct WorkerCorpus {
623 id: usize,
625 in_memory_corpus: Vec<CorpusEntry>,
628 history_map: Vec<u8>,
630 edge_indices: EdgeIndexMap,
632 sancov_history_map: Vec<u8>,
634 pub(crate) failed_replays: usize,
636 pub(crate) metrics: CorpusMetrics,
638 sequence_generator: SequenceGenerator,
640 current_mutated_index: Option<usize>,
642 config: Arc<FuzzCorpusConfig>,
644 worker_sync_enabled: bool,
646 new_entry_indices: Vec<usize>,
648 initial_export_dirs: Option<Vec<PathBuf>>,
650 worker_dir: Option<PathBuf>,
653 warned_persistence_failure: bool,
655 last_sync_metrics: CorpusMetrics,
657 optimization_best_value: Option<I256>,
659 optimization_best_sequence: Vec<BasicTxDetails>,
661 last_new_edge_at: Option<Instant>,
667}
668
669#[derive(Clone, Copy)]
672pub struct DynamicTargetCtx<'a> {
673 pub project_contracts: &'a ContractsByArtifact,
674 pub setup_contracts: &'a ContractsByAddress,
675 pub artifact_filters: &'a ArtifactFilters,
676}
677
678pub(crate) fn register_replay_created(
681 state_changeset: &StateChangeset,
682 dynamic: Option<&DynamicTargetCtx<'_>>,
683 fuzzed_contracts: Option<&FuzzRunIdentifiedContracts>,
684 created: &mut Vec<Address>,
685) {
686 let (Some(dynamic), Some(fuzzed_contracts)) = (dynamic, fuzzed_contracts) else {
687 return;
688 };
689 if let Err(error) = fuzzed_contracts.collect_created_contracts(
690 state_changeset,
691 dynamic.project_contracts,
692 dynamic.setup_contracts,
693 dynamic.artifact_filters,
694 created,
695 ) {
696 warn!(target: "corpus", "{error}");
697 }
698}
699
700pub(crate) fn rollback_replay_created(
702 fuzzed_contracts: Option<&FuzzRunIdentifiedContracts>,
703 created: Vec<Address>,
704) {
705 if !created.is_empty()
706 && let Some(fuzzed_contracts) = fuzzed_contracts
707 {
708 fuzzed_contracts.clear_created_contracts(created);
709 }
710}
711
712fn load_optimization_state(config: &FuzzCorpusConfig) -> Option<(I256, Vec<BasicTxDetails>)> {
713 let corpus_dir = config.corpus_dir.as_ref()?;
714 let opt_path = corpus_dir.join(OPTIMIZATION_BEST_FILE);
715 if !opt_path.is_file() {
716 return None;
717 }
718
719 match foundry_common::fs::read_json_file::<OptimizationState>(&opt_path) {
720 Ok(state) => {
721 debug!(
722 target: "corpus",
723 "loaded optimization best value {} with sequence len {}",
724 state.best_value,
725 state.best_sequence.len()
726 );
727 Some((state.best_value, state.best_sequence))
728 }
729 Err(err) => {
730 let _ = sh_warn!(
731 "failed to load optimization state from {}: {err}; starting without persisted optimization seed",
732 opt_path.display()
733 );
734 None
735 }
736 }
737}
738
739fn replay_corpus_sequence<FEN: FoundryEvmNetwork>(
740 tx_seq: &[BasicTxDetails],
741 executor: &Executor<FEN>,
742 target: ReplayTarget<'_>,
743 coverage: ReplayCoverage<'_>,
744) -> Result<ReplayOutcome> {
745 let mut executor = executor.clone();
746 replay_corpus_sequence_with_executor(tx_seq, &mut executor, target, coverage, false, true)
747}
748
749fn replay_corpus_sequence_with_executor<FEN: FoundryEvmNetwork>(
750 tx_seq: &[BasicTxDetails],
751 executor: &mut Executor<FEN>,
752 target: ReplayTarget<'_>,
753 mut coverage: ReplayCoverage<'_>,
754 trace_sync: bool,
755 reject_unmatched_function: bool,
756) -> Result<ReplayOutcome> {
757 let mut cmp_seq = Vec::with_capacity(tx_seq.len());
758 let mut failed_replays = 0;
759 let mut new_coverage_for_entry = false;
760 let mut new_edge_for_entry = false;
761 let mut created: Vec<Address> = Vec::new();
762
763 for tx in tx_seq {
764 if WorkerCorpus::can_replay_tx(tx, target.stateless, target.fuzzed_contracts) {
765 let mut call_result = execute_tx(executor, tx)?;
766 cmp_seq.push(
767 call_result
768 .evm_cmp_values
769 .take()
770 .unwrap_or_default()
771 .into_iter()
772 .map(|cmp| ComparisonHint { lhs: cmp.op1, rhs: cmp.op2 })
773 .collect(),
774 );
775 let (new_coverage, is_edge) = call_result.merge_all_coverage(
776 coverage.history_map,
777 coverage.edge_indices,
778 coverage.sancov_history_map,
779 );
780 if new_coverage {
781 new_coverage_for_entry = true;
782 new_edge_for_entry |= is_edge;
783 if let Some(metrics) = coverage.metrics.as_deref_mut() {
784 metrics.update_seen(is_edge);
785 }
786 }
787
788 register_replay_created(
789 &call_result.state_changeset,
790 target.dynamic,
791 target.fuzzed_contracts,
792 &mut created,
793 );
794
795 if target.fuzzed_contracts.is_some() {
797 executor.commit(&mut call_result);
798 }
799
800 if trace_sync {
801 trace!(
802 target: "corpus",
803 %new_coverage,
804 ?tx,
805 "replayed tx for syncing",
806 );
807 }
808 } else {
809 cmp_seq.push(Vec::new());
810 failed_replays += 1;
811
812 if reject_unmatched_function && target.stateless.is_some() {
813 rollback_replay_created(target.fuzzed_contracts, created);
814 return Ok(ReplayOutcome {
815 keep_entry: false,
816 new_coverage: new_coverage_for_entry,
817 new_edge: new_edge_for_entry,
818 cmp_seq,
819 failed_replays,
820 });
821 }
822 }
823 }
824 rollback_replay_created(target.fuzzed_contracts, created);
825
826 Ok(ReplayOutcome {
827 keep_entry: true,
828 new_coverage: new_coverage_for_entry,
829 new_edge: new_edge_for_entry,
830 cmp_seq,
831 failed_replays,
832 })
833}
834
835impl WorkerCorpus {
836 pub fn new_sequence(&mut self, test_runner: &mut TestRunner) -> Result<SequencePlan> {
838 if self.config.is_coverage_guided() && !self.in_memory_corpus.is_empty() {
839 self.evict_oldest_corpus()?;
840 }
841 let corpus_len = self.in_memory_corpus.len();
842 let plan = self.sequence_generator.start(
843 test_runner,
844 corpus_len,
845 |index| {
846 let entry = self
847 .in_memory_corpus
848 .get(index)
849 .ok_or_else(|| eyre::eyre!("corpus index {index} is out of bounds"))?;
850 CorpusEntryView::new(&entry.tx_seq, &entry.cmp_seq)
851 },
852 self.config.is_coverage_guided(),
853 )?;
854 self.current_mutated_index = plan.source();
855 Ok(plan)
856 }
857
858 pub fn new<FEN: FoundryEvmNetwork>(
859 id: usize,
860 config: FuzzCorpusConfig,
861 sequence_generator: SequenceGenerator,
862 replay_root: Option<&Path>,
863 executor: Option<&Executor<FEN>>,
865 target: ReplayTarget<'_>,
866 ) -> Result<Self> {
867 let seed = if id == 0 {
868 WorkerCorpusSeed::load_from_disk(&config, replay_root, executor, target)?
869 } else {
870 WorkerCorpusSeed::empty(&config).with_optimization_state(&config)
871 };
872 let mut corpus = Self::from_seed(id, config, sequence_generator, seed)?;
873 corpus.worker_sync_enabled = true;
874 Ok(corpus)
875 }
876
877 pub(crate) fn from_seed(
878 id: usize,
879 config: FuzzCorpusConfig,
880 sequence_generator: SequenceGenerator,
881 mut seed: WorkerCorpusSeed,
882 ) -> Result<Self> {
883 let initial_export_dirs = if id == 0 { seed.replay_dirs.take() } else { None };
884 let worker_dir = if let Some(corpus_dir) = &config.corpus_dir {
885 let worker_dir = corpus_dir.join(format!("{WORKER}{id}"));
886 let worker_corpus = worker_dir.join(CORPUS_DIR);
887 let sync_dir = worker_dir.join(SYNC_DIR);
888
889 foundry_common::fs::create_dir_all(&worker_corpus)?;
891 foundry_common::fs::create_dir_all(&sync_dir)?;
892
893 Some(worker_dir)
894 } else {
895 None
896 };
897
898 Ok(Self {
899 id,
900 in_memory_corpus: seed.in_memory_corpus,
901 history_map: seed.history_map,
902 edge_indices: seed.edge_indices,
903 sancov_history_map: seed.sancov_history_map,
904 failed_replays: seed.failed_replays,
905 metrics: seed.metrics,
906 sequence_generator,
907 current_mutated_index: None,
908 config: config.into(),
909 worker_sync_enabled: false,
910 new_entry_indices: Default::default(),
911 initial_export_dirs,
912 worker_dir,
913 warned_persistence_failure: false,
914 last_sync_metrics: Default::default(),
915 optimization_best_value: seed.optimization_best_value,
916 optimization_best_sequence: seed.optimization_best_sequence,
917 last_new_edge_at: seed.last_new_edge_at,
918 })
919 }
920
921 #[instrument(skip_all)]
925 pub fn process_inputs(
926 &mut self,
927 inputs: &[BasicTxDetails],
928 cmp_seq: &[Vec<CmpOperands>],
929 new_coverage: bool,
930 optimization: Option<(I256, Vec<BasicTxDetails>)>,
931 ) {
932 self.process_inputs_inner(
933 inputs,
934 cmp_seq,
935 new_coverage,
936 optimization,
937 CorpusInsertionMode::Live,
938 true,
939 );
940 }
941
942 #[instrument(skip_all)]
947 pub fn process_inputs_for_campaign(
948 &mut self,
949 inputs: &[BasicTxDetails],
950 cmp_seq: &[Vec<CmpOperands>],
951 new_coverage: bool,
952 optimization: Option<(I256, Vec<BasicTxDetails>)>,
953 ) {
954 self.process_inputs_inner(
955 inputs,
956 cmp_seq,
957 new_coverage,
958 optimization,
959 CorpusInsertionMode::Live,
960 false,
961 );
962 }
963
964 fn process_inputs_inner(
965 &mut self,
966 inputs: &[BasicTxDetails],
967 cmp_seq: &[Vec<CmpOperands>],
968 new_coverage: bool,
969 optimization: Option<(I256, Vec<BasicTxDetails>)>,
970 insertion_mode: CorpusInsertionMode,
971 persist_optimization: bool,
972 ) {
973 let improved_optimization = optimization.as_ref().is_some_and(|(value, _)| {
975 self.optimization_best_value.is_none_or(|best| *value > best)
976 });
977
978 if let Some(index) = self.current_mutated_index.take() {
980 let should_credit = new_coverage || improved_optimization;
981 if let Some(corpus) = self.in_memory_corpus.get_mut(index) {
982 corpus.total_mutations += 1;
983 if should_credit {
984 corpus.new_finds_produced += 1
985 }
986 let is_favored = (corpus.new_finds_produced as f64 / corpus.total_mutations as f64)
987 > FAVORABILITY_THRESHOLD;
988 self.metrics.update_favored(is_favored, corpus.is_favored);
989 corpus.is_favored = is_favored;
990
991 trace!(
992 target: "corpus",
993 "updated corpus {}, total mutations: {}, new finds: {}",
994 corpus.uuid, corpus.total_mutations, corpus.new_finds_produced
995 );
996 }
997 }
998 if let Some((value, best_seq)) = optimization
999 && improved_optimization
1000 {
1001 self.optimization_best_value = Some(value);
1002 self.optimization_best_sequence = best_seq;
1003 if persist_optimization {
1004 self.persist_optimization_state();
1005 }
1006 }
1007
1008 if !self.config.is_coverage_guided() {
1009 return;
1010 }
1011
1012 if !new_coverage && !improved_optimization {
1014 return;
1015 }
1016
1017 let corpus_inputs = if improved_optimization && (!new_coverage || inputs.is_empty()) {
1024 self.optimization_best_sequence.clone()
1025 } else {
1026 inputs.to_vec()
1027 };
1028 if corpus_inputs.is_empty() {
1029 return;
1030 }
1031 let corpus_cmp_seq = cmp_seq
1032 .iter()
1033 .take(corpus_inputs.len())
1034 .map(|values| {
1035 values.iter().map(|cmp| ComparisonHint { lhs: cmp.op1, rhs: cmp.op2 }).collect()
1036 })
1037 .collect();
1038 let corpus = CorpusEntry::new_with_cmp(corpus_inputs, corpus_cmp_seq, Uuid::new_v4());
1039
1040 self.insert_corpus_entry(corpus, insertion_mode)
1041 }
1042
1043 fn insert_corpus_entry(&mut self, corpus: CorpusEntry, insertion_mode: CorpusInsertionMode) {
1044 if matches!(insertion_mode, CorpusInsertionMode::Live)
1045 && let Some(worker_dir) = &self.worker_dir
1046 {
1047 let worker_corpus = worker_dir.join(CORPUS_DIR);
1048 let write_result = corpus.write_to_disk_in(&worker_corpus, self.config.corpus_gzip);
1049 if let Err(err) = write_result {
1050 if !self.warned_persistence_failure {
1051 let _ = sh_warn!(
1052 "Failed to persist coverage corpus entries for worker {} in {}: {err}",
1053 self.id,
1054 worker_corpus.display()
1055 );
1056 self.warned_persistence_failure = true;
1057 }
1058 debug!(target: "corpus", %err, "failed to record call sequence {:?}", corpus.tx_seq);
1059 } else {
1060 trace!(
1061 target: "corpus",
1062 "persisted {} inputs for new coverage for {} corpus",
1063 corpus.tx_seq.len(),
1064 corpus.uuid,
1065 );
1066 }
1067 }
1068
1069 self.push_corpus_entry(corpus);
1070 }
1071
1072 fn push_corpus_entry(&mut self, corpus: CorpusEntry) {
1073 let new_index = self.in_memory_corpus.len();
1074 if self.worker_sync_enabled {
1075 self.new_entry_indices.push(new_index);
1076 }
1077 self.metrics.corpus_count += 1;
1078 self.in_memory_corpus.push(corpus);
1079 }
1080
1081 pub fn optimization_initial_state(&self) -> (Option<I256>, Vec<BasicTxDetails>) {
1083 (self.optimization_best_value, self.optimization_best_sequence.clone())
1084 }
1085
1086 fn persist_optimization_state(&self) {
1088 let optimization_best = self
1089 .optimization_best_value
1090 .map(|value| (value, self.optimization_best_sequence.as_slice()));
1091 persist_optimization_output(&self.config, optimization_best);
1092 }
1093
1094 pub fn merge_edge_coverage<FEN: FoundryEvmNetwork>(
1096 &mut self,
1097 call_result: &mut RawCallResult<FEN>,
1098 ) -> bool {
1099 if !self.config.collect_edge_coverage() {
1100 return false;
1101 }
1102
1103 let (new_coverage, is_edge) = call_result.merge_all_coverage(
1104 &mut self.history_map,
1105 &mut self.edge_indices,
1106 &mut self.sancov_history_map,
1107 );
1108 if new_coverage {
1109 self.metrics.update_seen(is_edge);
1110 if is_edge {
1113 self.last_new_edge_at = Some(Instant::now());
1114 }
1115 }
1116 new_coverage
1117 }
1118
1119 pub(crate) fn time_since_new_edge(&self) -> Option<Duration> {
1122 self.last_new_edge_at.map(|at| at.elapsed())
1123 }
1124 pub fn hoist_observed_calls(
1129 &mut self,
1130 observed: &[ObservedCall],
1131 parent_tx: &BasicTxDetails,
1132 targeted_contracts: &FuzzRunIdentifiedContracts,
1133 insertion_mode: CorpusInsertionMode,
1134 ) {
1135 if !self.config.is_coverage_guided() || observed.is_empty() {
1136 return;
1137 }
1138
1139 let tx_seq = {
1140 let targets = targeted_contracts.targets();
1141 sequence_from_observed(
1142 observed,
1143 &targets,
1144 ObservedCallDepth::All,
1145 Some((parent_tx.warp, parent_tx.roll)),
1146 )
1147 };
1148
1149 self.push_observed_sequence(tx_seq, insertion_mode)
1150 }
1151
1152 pub fn seed_from_test_traces<FEN: FoundryEvmNetwork>(
1157 &mut self,
1158 invariant_contract: &InvariantContract<'_>,
1159 targeted_contracts: &FuzzRunIdentifiedContracts,
1160 executor: &Executor<FEN>,
1161 ) -> Result<usize> {
1162 if !self.config.is_coverage_guided() {
1163 return Ok(0);
1164 }
1165
1166 let mut added = 0;
1167
1168 for func in invariant_contract.abi.functions() {
1169 if !func.is_unit_test() {
1170 continue;
1171 }
1172 if invariant_contract
1173 .invariant_fns
1174 .iter()
1175 .any(|(invariant_fn, _)| func.selector() == invariant_fn.selector())
1176 {
1177 continue;
1178 }
1179
1180 let calldata = match func.abi_encode_input(&[]) {
1181 Ok(calldata) => Bytes::from(calldata),
1182 Err(_) => continue,
1183 };
1184
1185 let exec = executor.clone();
1186
1187 let raw = match exec.call_raw(CALLER, invariant_contract.address, calldata, U256::ZERO)
1188 {
1189 Ok(raw) => raw,
1190 Err(_) => continue,
1191 };
1192 if raw.reverted {
1193 continue;
1194 }
1195
1196 let observed = raw.observed_calls;
1197 if observed.is_empty() {
1198 continue;
1199 }
1200
1201 let seq = {
1202 let targets = targeted_contracts.targets();
1203 sequence_from_observed(&observed, &targets, ObservedCallDepth::DirectOnly, None)
1204 };
1205
1206 let insertion_mode = if self.id == 0 {
1207 CorpusInsertionMode::Live
1208 } else {
1209 CorpusInsertionMode::MemoryOnly
1210 };
1211 let len_before = self.in_memory_corpus.len();
1212 self.push_observed_sequence(seq, insertion_mode);
1213 if self.in_memory_corpus.len() > len_before {
1214 debug!(target: "corpus", test = %func.name, "seeded corpus sequence from test trace");
1215 added += 1;
1216 }
1217 }
1218
1219 Ok(added)
1220 }
1221
1222 fn push_observed_sequence(
1223 &mut self,
1224 tx_seq: Vec<BasicTxDetails>,
1225 insertion_mode: CorpusInsertionMode,
1226 ) {
1227 if !self.config.is_coverage_guided() || tx_seq.is_empty() {
1228 return;
1229 }
1230
1231 let corpus = CorpusEntry::new(tx_seq);
1232
1233 self.insert_corpus_entry(corpus, insertion_mode)
1234 }
1235 fn evict_oldest_corpus(&mut self) -> Result<()> {
1238 if self.in_memory_corpus.len() > self.config.corpus_min_size.max(1)
1239 && let Some(index) =
1240 self.in_memory_corpus.iter().enumerate().position(|(index, corpus)| {
1241 self.new_entry_indices.binary_search(&index).is_err()
1242 && corpus.total_mutations > self.config.corpus_min_mutations
1243 && !corpus.is_favored
1244 })
1245 {
1246 let corpus = &self.in_memory_corpus[index];
1247
1248 trace!(target: "corpus", corpus=%serde_json::to_string(&corpus).unwrap(), "evict corpus");
1249
1250 self.in_memory_corpus.remove(index);
1252
1253 self.new_entry_indices.retain_mut(|i| {
1255 if *i > index {
1256 *i -= 1; true } else {
1259 *i != index }
1261 });
1262 }
1263 Ok(())
1264 }
1265 fn load_sync_corpus(&self, strict: bool) -> Result<Vec<(CorpusDirEntry, Vec<BasicTxDetails>)>> {
1270 let Some(worker_dir) = &self.worker_dir else {
1271 return Ok(vec![]);
1272 };
1273
1274 let sync_dir = worker_dir.join(SYNC_DIR);
1275 if !strict && !sync_dir.is_dir() {
1276 return Ok(vec![]);
1277 }
1278
1279 let mut imports = vec![];
1280 let entries = if strict {
1281 read_corpus_dir_strict(&sync_dir)?
1282 } else {
1283 read_corpus_dir(&sync_dir).collect()
1284 };
1285 for entry in entries {
1286 let tx_seq = match entry.read_tx_seq() {
1288 Ok(tx_seq) => tx_seq,
1289 Err(err) if strict => {
1290 return Err(eyre!(
1291 "failed to read final corpus entry {}: {err}",
1292 entry.path.display()
1293 ));
1294 }
1295 Err(err) => {
1296 warn!(target: "corpus", "skipping unreadable corpus file {}: {err}", entry.path.display());
1297 let quarantine_path = entry.path.with_file_name(format!(
1298 "{}.{}.invalid",
1299 entry.name(),
1300 Uuid::new_v4()
1301 ));
1302 if let Err(err) = std::fs::rename(&entry.path, &quarantine_path) {
1303 debug!(target: "corpus", %err, "failed to quarantine unreadable corpus file {}", entry.path.display());
1304 }
1305 continue;
1306 }
1307 };
1308 if tx_seq.is_empty() {
1309 warn!(target: "corpus", "skipping empty corpus entry: {}", entry.path.display());
1310 if let Err(err) = std::fs::remove_file(&entry.path) {
1311 if strict {
1312 return Err(err.into());
1313 }
1314 debug!(target: "corpus", %err, "failed to remove empty corpus file {}", entry.path.display());
1315 }
1316 continue;
1317 }
1318 imports.push((entry, tx_seq));
1319 }
1320
1321 if !imports.is_empty() {
1322 debug!(target: "corpus", "imported {} new corpus entries", imports.len());
1323 }
1324
1325 Ok(imports)
1326 }
1327
1328 fn push_synced_corpus_entry(
1330 &mut self,
1331 mut corpus: CorpusEntry,
1332 timestamp: u64,
1333 file_name: String,
1334 ) {
1335 corpus.timestamp = timestamp;
1336 corpus.persisted_file_name = Some(file_name);
1337 if self.worker_sync_enabled && self.id == 0 {
1338 self.new_entry_indices.push(self.in_memory_corpus.len());
1339 }
1340 self.in_memory_corpus.push(corpus);
1341 }
1342
1343 #[instrument(skip_all)]
1346 fn calibrate<FEN: FoundryEvmNetwork>(
1347 &mut self,
1348 executor: &Executor<FEN>,
1349 target: ReplayTarget<'_>,
1350 strict: bool,
1351 ) -> Result<()> {
1352 let Some(worker_dir) = &self.worker_dir else {
1353 return Ok(());
1354 };
1355 let corpus_dir = worker_dir.join(CORPUS_DIR);
1356
1357 for (entry, tx_seq) in self.load_sync_corpus(strict)? {
1358 let mut history_map = self.history_map.clone();
1359 let mut edge_indices = self.edge_indices.clone();
1360 let mut sancov_history_map = self.sancov_history_map.clone();
1361 let mut metrics = self.metrics.clone();
1362 let coverage = ReplayCoverage {
1363 history_map: &mut history_map,
1364 edge_indices: &mut edge_indices,
1365 sancov_history_map: &mut sancov_history_map,
1366 metrics: Some(&mut metrics),
1367 };
1368 let mut replay_executor = executor.clone();
1369 let ReplayOutcome { keep_entry, new_coverage, new_edge, cmp_seq, .. } =
1370 replay_corpus_sequence_with_executor(
1371 &tx_seq,
1372 &mut replay_executor,
1373 target,
1374 coverage,
1375 true,
1376 false,
1377 )?;
1378
1379 let sync_path = &entry.path;
1380 if keep_entry && new_coverage {
1381 let corpus_path = corpus_dir.join(sync_path.components().next_back().unwrap());
1383 if !accept_synced_corpus_file(&entry, &tx_seq, &corpus_path) {
1384 if strict {
1385 return Err(eyre!("failed to accept final corpus entry {}", entry.name()));
1386 }
1387 continue;
1388 }
1389
1390 self.history_map = history_map;
1391 self.edge_indices = edge_indices;
1392 self.sancov_history_map = sancov_history_map;
1393 self.metrics = metrics;
1394 if new_edge {
1396 self.last_new_edge_at = Some(Instant::now());
1397 }
1398
1399 debug!(
1400 target: "corpus",
1401 name=%entry.name(),
1402 "moved synced corpus to corpus dir",
1403 );
1404
1405 let corpus_entry = CorpusEntry::new_with_cmp(tx_seq.clone(), cmp_seq, entry.uuid);
1406 self.push_synced_corpus_entry(
1407 corpus_entry,
1408 entry.timestamp,
1409 entry.name().to_owned(),
1410 );
1411 } else {
1412 if let Err(err) = std::fs::remove_file(&entry.path) {
1414 if strict {
1415 return Err(err.into());
1416 }
1417 debug!(target: "corpus", %err, "failed to remove synced corpus from {sync_path:?}");
1418 continue;
1419 }
1420 trace!(target: "corpus", "removed synced corpus from {sync_path:?}");
1421 }
1422 }
1423
1424 Ok(())
1425 }
1426
1427 #[instrument(skip_all)]
1429 fn export_to_master(&mut self) -> Result<()> {
1430 assert_ne!(self.id, 0, "non-master only");
1432
1433 if self.new_entry_indices.is_empty() || self.worker_dir.is_none() {
1435 return Ok(());
1436 }
1437
1438 let worker_dir = self.worker_dir.as_ref().unwrap();
1439 let Some(master_sync_dir) = self
1440 .config
1441 .corpus_dir
1442 .as_ref()
1443 .map(|dir| dir.join(format!("{WORKER}0")).join(SYNC_DIR))
1444 else {
1445 return Ok(());
1446 };
1447
1448 let mut exported = 0;
1449 let corpus_dir = worker_dir.join(CORPUS_DIR);
1450 let mut delivered = HashSet::new();
1451
1452 for &index in &self.new_entry_indices {
1453 let Some(corpus) = self.in_memory_corpus.get(index) else {
1454 delivered.insert(index);
1455 continue;
1456 };
1457 let file_name = corpus.file_name(corpus.should_gzip(self.config.corpus_gzip));
1458 let file_path = corpus_dir.join(&file_name);
1459 if !file_path.is_file()
1460 && let Err(err) = corpus.write_to_disk_in(&corpus_dir, self.config.corpus_gzip)
1461 {
1462 debug!(target: "corpus", %err, "failed to persist corpus {} for export", corpus.uuid);
1463 continue;
1464 }
1465 let sync_path = master_sync_dir.join(&file_name);
1466 if !link_corpus_file(&file_path, &sync_path) {
1467 continue;
1468 }
1469 exported += 1;
1470 delivered.insert(index);
1471 }
1472 self.new_entry_indices.retain(|index| !delivered.contains(index));
1473
1474 debug!(target: "corpus", "exported {exported} new corpus entries");
1475
1476 Ok(())
1477 }
1478
1479 #[instrument(skip_all)]
1481 fn export_to_workers(&mut self, num_workers: usize) -> Result<()> {
1482 assert_eq!(self.id, 0, "master worker only");
1483 if self.worker_dir.is_none() {
1484 return Ok(());
1485 }
1486
1487 let worker_dir = self.worker_dir.as_ref().unwrap();
1488 let master_corpus_dir = worker_dir.join(CORPUS_DIR);
1489 let startup_entries = if let Some(replay_dirs) = &self.initial_export_dirs {
1490 let mut seen_entries = HashSet::new();
1491 unique_corpus_entries(replay_dirs, &mut seen_entries)
1492 .map(|entry| entry.path)
1493 .collect::<Vec<_>>()
1494 } else {
1495 Vec::new()
1496 };
1497 let mut pending_entries = Vec::new();
1498 let mut delivered = HashSet::new();
1499 for &index in &self.new_entry_indices {
1500 let Some(corpus) = self.in_memory_corpus.get(index) else {
1501 delivered.insert(index);
1502 continue;
1503 };
1504 let path = master_corpus_dir
1505 .join(corpus.file_name(corpus.should_gzip(self.config.corpus_gzip)));
1506 if !path.is_file()
1507 && let Err(err) =
1508 corpus.write_to_disk_in(&master_corpus_dir, self.config.corpus_gzip)
1509 {
1510 debug!(target: "corpus", %err, "failed to persist corpus {} for fan-out", corpus.uuid);
1511 continue;
1512 }
1513 pending_entries.push((index, path));
1514 }
1515
1516 let mut target_dirs = Vec::new();
1517 for target_worker in 1..num_workers {
1518 let target_dir = self
1519 .config
1520 .corpus_dir
1521 .as_ref()
1522 .unwrap()
1523 .join(format!("{WORKER}{target_worker}"))
1524 .join(SYNC_DIR);
1525 if !target_dir.is_dir() {
1526 foundry_common::fs::create_dir_all(&target_dir)?;
1527 }
1528 target_dirs.push(target_dir);
1529 }
1530
1531 let mut any_distributed = false;
1532 let mut startup_delivered = true;
1533 for path in &startup_entries {
1534 let Some(name) = path.file_name() else {
1535 startup_delivered = false;
1536 continue;
1537 };
1538 let mut delivered_to_all = true;
1539 for target_dir in &target_dirs {
1540 let sync_path = target_dir.join(name);
1541 if link_corpus_file(path, &sync_path) {
1542 any_distributed = true;
1543 trace!(target: "corpus", name=%name.to_string_lossy(), ?target_dir, "distributed corpus");
1544 } else {
1545 delivered_to_all = false;
1546 }
1547 }
1548 startup_delivered &= delivered_to_all;
1549 }
1550
1551 for (index, path) in pending_entries {
1552 let Some(name) = path.file_name() else { continue };
1553 let mut delivered_to_all = true;
1554 for target_dir in &target_dirs {
1555 let sync_path = target_dir.join(name);
1556 if link_corpus_file(&path, &sync_path) {
1557 any_distributed = true;
1558 trace!(target: "corpus", name=%name.to_string_lossy(), ?target_dir, "distributed corpus");
1559 } else {
1560 delivered_to_all = false;
1561 }
1562 }
1563 if delivered_to_all {
1564 delivered.insert(index);
1565 }
1566 }
1567
1568 self.new_entry_indices.retain(|index| !delivered.contains(index));
1569 if startup_delivered {
1570 self.initial_export_dirs = None;
1571 }
1572 debug!(target: "corpus", %any_distributed, "distributed master corpus to all workers");
1573
1574 Ok(())
1575 }
1576
1577 pub(crate) fn sync_metrics(&mut self, global_corpus_metrics: &GlobalCorpusMetrics) {
1580 let edges_delta = self
1582 .metrics
1583 .cumulative_edges_seen
1584 .saturating_sub(self.last_sync_metrics.cumulative_edges_seen);
1585 let features_delta = self
1586 .metrics
1587 .cumulative_features_seen
1588 .saturating_sub(self.last_sync_metrics.cumulative_features_seen);
1589 let corpus_count_delta =
1591 self.metrics.corpus_count as isize - self.last_sync_metrics.corpus_count as isize;
1592 let favored_delta =
1593 self.metrics.favored_items as isize - self.last_sync_metrics.favored_items as isize;
1594
1595 if edges_delta > 0 {
1598 global_corpus_metrics.cumulative_edges_seen.fetch_add(edges_delta, Ordering::Relaxed);
1599 }
1600 if features_delta > 0 {
1601 global_corpus_metrics
1602 .cumulative_features_seen
1603 .fetch_add(features_delta, Ordering::Relaxed);
1604 }
1605
1606 if corpus_count_delta > 0 {
1607 global_corpus_metrics
1608 .corpus_count
1609 .fetch_add(corpus_count_delta as usize, Ordering::Relaxed);
1610 } else if corpus_count_delta < 0 {
1611 global_corpus_metrics
1612 .corpus_count
1613 .fetch_sub((-corpus_count_delta) as usize, Ordering::Relaxed);
1614 }
1615
1616 if favored_delta > 0 {
1617 global_corpus_metrics
1618 .favored_items
1619 .fetch_add(favored_delta as usize, Ordering::Relaxed);
1620 } else if favored_delta < 0 {
1621 global_corpus_metrics
1622 .favored_items
1623 .fetch_sub((-favored_delta) as usize, Ordering::Relaxed);
1624 }
1625
1626 self.last_sync_metrics = self.metrics.clone();
1628 }
1629
1630 #[instrument(skip_all)]
1632 pub fn sync<FEN: FoundryEvmNetwork>(
1633 &mut self,
1634 num_workers: usize,
1635 executor: &Executor<FEN>,
1636 target: ReplayTarget<'_>,
1637 global_corpus_metrics: &GlobalCorpusMetrics,
1638 ) -> Result<()> {
1639 trace!(target: "corpus", "syncing");
1640
1641 self.sync_metrics(global_corpus_metrics);
1642
1643 self.calibrate(executor, target, false)?;
1644 if self.id == 0 {
1645 self.export_to_workers(num_workers)?;
1646 } else {
1647 self.export_to_master()?;
1648 }
1649
1650 debug!(target: "corpus", "synced");
1651
1652 Ok(())
1653 }
1654
1655 pub(crate) fn finalize_sync<FEN: FoundryEvmNetwork>(
1657 &mut self,
1658 executor: &Executor<FEN>,
1659 target: ReplayTarget<'_>,
1660 coordinator: &CorpusSyncCoordinator,
1661 ) -> Result<()> {
1662 if self.id != 0 {
1663 self.export_to_master()?;
1664 if self.worker_dir.is_some() && !self.new_entry_indices.is_empty() {
1665 return Err(eyre!("worker {} failed to complete final corpus export", self.id));
1666 }
1667 }
1668 if !coordinator.wait() {
1669 return Ok(());
1670 }
1671
1672 if self.id == 0 {
1673 self.calibrate(executor, target, true)?;
1674 self.export_to_workers(coordinator.workers)?;
1675 if self.worker_dir.is_some()
1676 && (!self.new_entry_indices.is_empty() || self.initial_export_dirs.is_some())
1677 {
1678 return Err(eyre!("master failed to complete final corpus fan-out"));
1679 }
1680 }
1681 if !coordinator.wait() {
1682 return Ok(());
1683 }
1684
1685 if self.id != 0 {
1686 self.calibrate(executor, target, true)?;
1687 }
1688 if !coordinator.wait() {
1689 return Ok(());
1690 }
1691
1692 Ok(())
1693 }
1694
1695 pub(crate) fn can_replay_tx(
1697 tx: &BasicTxDetails,
1698 stateless: Option<StatelessReplayTarget<'_>>,
1699 fuzzed_contracts: Option<&FuzzRunIdentifiedContracts>,
1700 ) -> bool {
1701 fuzzed_contracts.is_some_and(|contracts| contracts.targets().can_replay(tx))
1702 || stateless.is_some_and(|target| target.can_replay(tx))
1703 }
1704}
1705
1706#[derive(Clone, Copy)]
1707enum ObservedCallDepth {
1708 DirectOnly,
1709 All,
1710}
1711
1712fn sequence_from_observed(
1713 observed: &[ObservedCall],
1714 targets: &TargetedContracts,
1715 depth: ObservedCallDepth,
1716 first_delay: Option<(Option<U256>, Option<U256>)>,
1717) -> Vec<BasicTxDetails> {
1718 let mut first_delay = first_delay;
1719 observed
1720 .iter()
1721 .filter(|call| matches!(depth, ObservedCallDepth::All) || call.depth == 1)
1722 .filter_map(|call| {
1723 let mut tx = BasicTxDetails {
1724 warp: None,
1725 roll: None,
1726 sender: call.caller,
1727 call_details: CallDetails {
1728 target: call.target,
1729 calldata: call.calldata.clone(),
1730 value: call.value,
1731 },
1732 };
1733 targets.can_replay(&tx).then(|| {
1734 let (warp, roll) = first_delay.take().unwrap_or((None, None));
1735 tx.warp = warp;
1736 tx.roll = roll;
1737 tx
1738 })
1739 })
1740 .collect()
1741}
1742
1743fn persist_optimization_output(
1744 config: &FuzzCorpusConfig,
1745 optimization_best: Option<(I256, &[BasicTxDetails])>,
1746) {
1747 let Some(root) = &config.corpus_dir else {
1748 return;
1749 };
1750 let Some((value, sequence)) = optimization_best else {
1751 return;
1752 };
1753 let state = OptimizationState { best_value: value, best_sequence: sequence.to_vec() };
1754 let path = root.join(OPTIMIZATION_BEST_FILE);
1755 if let Err(err) = foundry_common::fs::write_json_file(&path, &state) {
1756 debug!(target: "corpus", %err, "failed to persist optimization state");
1757 } else {
1758 trace!(
1759 target: "corpus",
1760 "persisted optimization best value {} with sequence len {}",
1761 value,
1762 sequence.len()
1763 );
1764 }
1765}
1766
1767pub(crate) fn persist_campaign_optimization(
1768 config: &FuzzCorpusConfig,
1769 value: Option<I256>,
1770 sequence: &[BasicTxDetails],
1771) {
1772 persist_optimization_output(config, value.map(|value| (value, sequence)));
1773}
1774
1775fn has_legacy_invariant_corpus_dirs(path: &Path) -> bool {
1776 std::fs::read_dir(path).is_ok_and(|entries| {
1777 entries.flatten().any(|entry| {
1778 let path = entry.path();
1779 path.is_dir()
1780 && entry.file_name().to_str().is_some_and(|name| !name.starts_with(WORKER))
1781 && !path.join(OPTIMIZATION_BEST_FILE).is_file()
1782 })
1783 })
1784}
1785
1786fn unique_corpus_entries<'a>(
1787 replay_dirs: &'a [PathBuf],
1788 seen_entries: &'a mut HashSet<Uuid>,
1789) -> impl Iterator<Item = CorpusDirEntry> + 'a {
1790 replay_dirs.iter().flat_map(|replay_dir| read_corpus_dir(replay_dir)).filter(|entry| {
1791 let is_new = seen_entries.insert(entry.uuid);
1792 if !is_new {
1793 trace!(target: "corpus", "skipping duplicate corpus entry {}", entry.uuid);
1794 }
1795 is_new
1796 })
1797}
1798
1799#[cfg(test)]
1800mod tests {
1801 use super::*;
1802 use crate::{
1803 executors::ExecutorBuilder,
1804 inspectors::{EdgeCovHit, EdgeCoverage, EdgeKey},
1805 };
1806 use alloy_dyn_abi::DynSolValue;
1807 use foundry_config::FuzzDictionaryConfig;
1808 use foundry_evm_core::{
1809 backend::Backend,
1810 evm::{EthEvmNetwork, EvmEnvFor, TxEnvFor},
1811 };
1812 use foundry_evm_fuzz::strategies::{EvmFuzzState, TxGenerator};
1813 use proptest::prelude::{Just, Strategy};
1814 use rayon::prelude::*;
1815 use revm::{
1816 bytecode::Bytecode,
1817 database::{CacheDB, EmptyDB},
1818 };
1819 use std::fs;
1820
1821 fn basic_tx() -> BasicTxDetails {
1822 BasicTxDetails {
1823 warp: None,
1824 roll: None,
1825 sender: Address::ZERO,
1826 call_details: foundry_evm_fuzz::CallDetails {
1827 target: Address::ZERO,
1828 calldata: Bytes::new(),
1829 value: None,
1830 },
1831 }
1832 }
1833
1834 fn basic_tx_with_calldata(calldata: impl Into<Bytes>) -> BasicTxDetails {
1835 let mut tx = basic_tx();
1836 tx.call_details.calldata = calldata.into();
1837 tx
1838 }
1839
1840 fn tx_for_function(
1841 target: Address,
1842 function: &Function,
1843 args: &[DynSolValue],
1844 ) -> BasicTxDetails {
1845 BasicTxDetails {
1846 warp: None,
1847 roll: None,
1848 sender: Address::ZERO,
1849 call_details: foundry_evm_fuzz::CallDetails {
1850 target,
1851 calldata: Bytes::from(function.abi_encode_input(args).unwrap()),
1852 value: None,
1853 },
1854 }
1855 }
1856
1857 fn empty_fuzz_state() -> EvmFuzzState {
1858 EvmFuzzState::new(
1859 &[],
1860 &CacheDB::<EmptyDB>::default(),
1861 FuzzDictionaryConfig::default(),
1862 None,
1863 )
1864 }
1865
1866 fn temp_corpus_dir() -> PathBuf {
1867 let dir = std::env::temp_dir().join(format!("foundry-corpus-tests-{}", Uuid::new_v4()));
1868 let _ = fs::create_dir_all(&dir);
1869 dir
1870 }
1871
1872 fn corpus_config(corpus_dir: PathBuf) -> FuzzCorpusConfig {
1873 FuzzCorpusConfig {
1874 corpus_dir: Some(corpus_dir),
1875 corpus_gzip: false,
1876 corpus_min_mutations: 0,
1877 corpus_min_size: 0,
1878 ..Default::default()
1879 }
1880 }
1881
1882 fn test_sequence(config: &FuzzCorpusConfig, tx: TxGenerator) -> SequenceGenerator {
1883 SequenceGenerator::stateless(
1884 tx,
1885 empty_fuzz_state().stateless_worker(),
1886 Function::parse("test(uint256)").unwrap(),
1887 config,
1888 )
1889 .unwrap()
1890 }
1891
1892 fn worker_corpus(id: usize, corpus_root: PathBuf, seed: WorkerCorpusSeed) -> WorkerCorpus {
1893 let config = corpus_config(corpus_root);
1894 let generator =
1895 test_sequence(&config, TxGenerator::from_strategy(Just(basic_tx()).boxed()));
1896 let mut corpus = WorkerCorpus::from_seed(id, config, generator, seed).unwrap();
1897 corpus.worker_sync_enabled = true;
1898 corpus
1899 }
1900
1901 fn empty_worker_corpus(id: usize, corpus_root: PathBuf) -> WorkerCorpus {
1902 worker_corpus(id, corpus_root, WorkerCorpusSeed::default())
1903 }
1904
1905 #[test]
1906 fn worker_initialization_fails_when_corpus_directories_cannot_be_created() {
1907 let corpus_root = temp_corpus_dir().join("not-a-directory");
1908 fs::write(&corpus_root, b"blocked").unwrap();
1909 let config = corpus_config(corpus_root);
1910 let generator =
1911 test_sequence(&config, TxGenerator::from_strategy(Just(basic_tx()).boxed()));
1912
1913 assert!(
1914 WorkerCorpus::from_seed(0, config, generator, WorkerCorpusSeed::default()).is_err()
1915 );
1916 }
1917
1918 fn sync_test_executor(corpus_root: PathBuf, target: Address) -> Executor<EthEvmNetwork> {
1919 let mut executor = ExecutorBuilder::<EthEvmNetwork>::default().gas_limit(1 << 24).build(
1920 EvmEnvFor::<EthEvmNetwork>::default(),
1921 TxEnvFor::<EthEvmNetwork>::default(),
1922 Backend::spawn(None).unwrap(),
1923 Default::default(),
1924 );
1925 executor.inspector_mut().collect_edge_coverage_with_config(&corpus_config(corpus_root));
1926 executor
1928 .set_code(
1929 target,
1930 Bytecode::new_raw(Bytes::from_static(&[
1931 0x60, 0x04, 0x35, 0x60, 0x08, 0x57, 0x00, 0x00, 0x5b, 0x00,
1932 ])),
1933 )
1934 .unwrap();
1935 executor
1936 }
1937
1938 fn finalize_test_worker(
1939 worker: &mut WorkerCorpus,
1940 corpus_root: PathBuf,
1941 target_address: Address,
1942 coordinator: &CorpusSyncCoordinator,
1943 ) {
1944 let function = Function::parse("test(uint256)").unwrap();
1945 let executor = sync_test_executor(corpus_root, target_address);
1946 worker
1947 .finalize_sync(
1948 &executor,
1949 ReplayTarget {
1950 stateless: Some(StatelessReplayTarget {
1951 function: &function,
1952 address: target_address,
1953 }),
1954 fuzzed_contracts: None,
1955 dynamic: None,
1956 },
1957 coordinator,
1958 )
1959 .unwrap();
1960 }
1961
1962 fn seeded_worker_corpus(
1963 id: usize,
1964 corpus_root: PathBuf,
1965 entries: Vec<CorpusEntry>,
1966 ) -> WorkerCorpus {
1967 worker_corpus(
1968 id,
1969 corpus_root,
1970 WorkerCorpusSeed { in_memory_corpus: entries, ..Default::default() },
1971 )
1972 }
1973 fn new_manager_with_single_corpus() -> (WorkerCorpus, Uuid) {
1974 let corpus = CorpusEntry::new(vec![basic_tx()]);
1975 let seed_uuid = corpus.uuid;
1976 let mut manager = seeded_worker_corpus(0, temp_corpus_dir(), vec![corpus]);
1977 manager.current_mutated_index = Some(0);
1978
1979 (manager, seed_uuid)
1980 }
1981
1982 fn targeted_contracts_with_selective_functions(
1983 target: Address,
1984 functions: Vec<Function>,
1985 targeted_selectors: impl IntoIterator<Item = alloy_primitives::Selector>,
1986 ) -> FuzzRunIdentifiedContracts {
1987 use alloy_json_abi::JsonAbi;
1988 use foundry_evm_fuzz::invariant::TargetedContract;
1989
1990 let mut abi = JsonAbi::new();
1991 for function in functions {
1992 abi.functions.entry(function.name.clone()).or_default().push(function);
1993 }
1994
1995 let mut contract = TargetedContract::new("Target".to_string(), abi);
1996 contract.add_selectors(targeted_selectors, false).unwrap();
1997
1998 let mut targets = TargetedContracts::new();
1999 targets.inner.insert(target, contract);
2000 FuzzRunIdentifiedContracts::new(targets, false)
2001 }
2002
2003 #[test]
2007 fn corrupt_corpus_file_surfaces_as_error_for_load_to_skip() {
2008 let dir = temp_corpus_dir();
2009
2010 let valid = CorpusEntry::new(vec![basic_tx()]);
2012 valid.write_to_disk_in(&dir, false).unwrap();
2013
2014 let corrupt_path = dir.join(format!("{}-123.json", Uuid::new_v4()));
2016 fs::write(&corrupt_path, b"{ not valid json").unwrap();
2017
2018 let entries = read_corpus_dir(&dir).collect::<Vec<_>>();
2019 assert_eq!(entries.len(), 2, "directory scan should surface both files");
2020
2021 let (mut ok, mut err) = (0u32, 0u32);
2022 for entry in &entries {
2023 match entry.read_tx_seq() {
2024 Ok(seq) => {
2025 ok += 1;
2026 assert_eq!(seq.len(), 1);
2027 }
2028 Err(_) => err += 1,
2029 }
2030 }
2031 assert_eq!((ok, err), (1, 1), "the corrupt file must read as Err, the valid one as Ok");
2032 }
2033
2034 #[test]
2035 fn sync_inbox_loads_entries_regardless_of_timestamp() {
2036 let corpus_root = temp_corpus_dir();
2037 let manager = empty_worker_corpus(0, corpus_root.clone());
2038 let sync_dir = corpus_root.join("worker0").join(SYNC_DIR);
2039 let mut corpus = CorpusEntry::new(vec![basic_tx()]);
2040 corpus.timestamp = 0;
2041 corpus.write_to_disk_in(&sync_dir, false).unwrap();
2042 let mut empty = CorpusEntry::new(vec![]);
2043 empty.timestamp = 0;
2044 let empty_path = empty.write_to_disk_in(&sync_dir, false).unwrap();
2045 let corrupt_path = sync_dir.join(format!("{}-0.json", Uuid::new_v4()));
2046 fs::write(&corrupt_path, b"{ not valid json").unwrap();
2047
2048 let imports = manager.load_sync_corpus(false).unwrap();
2049
2050 assert_eq!(imports.len(), 1);
2051 assert_eq!(imports[0].0.timestamp, 0);
2052 assert!(!empty_path.exists());
2053 assert!(!corrupt_path.exists());
2054 assert!(fs::read_dir(&sync_dir).unwrap().flatten().any(|entry| {
2055 entry.file_name().to_string_lossy().ends_with(".invalid") && entry.path().is_file()
2056 }));
2057 }
2058
2059 #[test]
2060 fn synced_entries_are_queued_for_fanout_only_on_master() {
2061 let corpus_root = temp_corpus_dir();
2062 let mut master = empty_worker_corpus(0, corpus_root.clone());
2063 let mut worker = empty_worker_corpus(1, corpus_root);
2064 let master_entry = CorpusEntry::new(vec![basic_tx()]);
2065 let worker_entry = CorpusEntry::new(vec![basic_tx()]);
2066
2067 master.push_synced_corpus_entry(master_entry, 1, "master-1.json".to_string());
2068 worker.push_synced_corpus_entry(worker_entry, 2, "worker-2.json".to_string());
2069
2070 assert_eq!(master.new_entry_indices, [0]);
2071 assert!(worker.new_entry_indices.is_empty());
2072 assert_eq!(master.in_memory_corpus[0].timestamp, 1);
2073 assert_eq!(worker.in_memory_corpus[0].timestamp, 2);
2074 assert_eq!(master.metrics.corpus_count, 0);
2075 assert_eq!(worker.metrics.corpus_count, 0);
2076 }
2077
2078 #[test]
2079 fn master_distributes_old_synced_entries() {
2080 let corpus_root = temp_corpus_dir();
2081 let mut master = empty_worker_corpus(0, corpus_root.clone());
2082 master.initial_export_dirs = None;
2083 let mut corpus =
2084 CorpusEntry::new(vec![basic_tx_with_calldata(vec![0; GZIP_THRESHOLD * 2])]);
2085 corpus.timestamp = 0;
2086 let path =
2087 corpus.write_to_disk_in(&corpus_root.join("worker0").join(CORPUS_DIR), true).unwrap();
2088 let name = path.file_name().unwrap().to_str().unwrap().to_owned();
2089 assert!(name.ends_with(".json.gz"));
2090 master.push_synced_corpus_entry(corpus, 0, name.clone());
2091
2092 master.export_to_workers(3).unwrap();
2093
2094 assert!(corpus_root.join("worker1").join(SYNC_DIR).join(&name).is_file());
2095 assert!(corpus_root.join("worker2").join(SYNC_DIR).join(name).is_file());
2096 assert!(master.new_entry_indices.is_empty());
2097 }
2098
2099 #[test]
2100 fn worker_retries_failed_export() {
2101 let corpus_root = temp_corpus_dir();
2102 let mut worker = empty_worker_corpus(1, corpus_root.clone());
2103 let corpus = CorpusEntry::new(vec![basic_tx_with_calldata([1])]);
2104 let name = corpus.file_name(false);
2105 worker.push_corpus_entry(corpus);
2106
2107 let master_sync = corpus_root.join("worker0").join(SYNC_DIR);
2108 fs::create_dir_all(&master_sync).unwrap();
2109 let destination = master_sync.join(&name);
2110 foundry_common::fs::write_json_file(&destination, &vec![basic_tx_with_calldata([2])])
2111 .unwrap();
2112
2113 worker.export_to_master().unwrap();
2114 assert_eq!(worker.new_entry_indices, [0]);
2115 assert!(corpus_root.join("worker1").join(CORPUS_DIR).join(&name).is_file());
2116
2117 fs::remove_file(&destination).unwrap();
2118 worker.export_to_master().unwrap();
2119
2120 assert!(worker.new_entry_indices.is_empty());
2121 let exported = read_corpus_dir(&master_sync).next().unwrap().read_tx_seq().unwrap();
2122 assert!(same_tx_sequence(&exported, &[basic_tx_with_calldata([1])]));
2123 }
2124
2125 #[test]
2126 fn master_retries_partial_fanout() {
2127 let corpus_root = temp_corpus_dir();
2128 let mut master = empty_worker_corpus(0, corpus_root.clone());
2129 master.initial_export_dirs = None;
2130 let corpus = CorpusEntry::new(vec![basic_tx_with_calldata([1])]);
2131 let name = corpus.file_name(false);
2132 corpus.write_to_disk_in(&corpus_root.join("worker0").join(CORPUS_DIR), false).unwrap();
2133 master.push_corpus_entry(corpus);
2134
2135 let worker2_sync = corpus_root.join("worker2").join(SYNC_DIR);
2136 fs::create_dir_all(&worker2_sync).unwrap();
2137 let worker2_destination = worker2_sync.join(&name);
2138 foundry_common::fs::write_json_file(
2139 &worker2_destination,
2140 &vec![basic_tx_with_calldata([2])],
2141 )
2142 .unwrap();
2143
2144 master.export_to_workers(3).unwrap();
2145 assert_eq!(master.new_entry_indices, [0]);
2146 assert!(corpus_root.join("worker1").join(SYNC_DIR).join(&name).is_file());
2147
2148 fs::remove_file(&worker2_destination).unwrap();
2149 master.export_to_workers(3).unwrap();
2150
2151 assert!(master.new_entry_indices.is_empty());
2152 assert!(worker2_destination.is_file());
2153 }
2154
2155 #[test]
2156 fn final_sync_coordinator_yields_to_nested_rayon_workers() {
2157 rayon::ThreadPoolBuilder::new().num_threads(2).build().unwrap().install(|| {
2158 (0..2usize).into_par_iter().for_each(|_| {
2159 let coordinator = CorpusSyncCoordinator::new(2);
2160 (0..2usize).into_par_iter().for_each(|_| assert!(coordinator.wait()));
2161 });
2162 });
2163 }
2164
2165 #[test]
2166 fn final_calibration_rejects_non_file_corpus_entry() {
2167 let corpus_root = temp_corpus_dir();
2168 let mut worker = empty_worker_corpus(1, corpus_root.clone());
2169 let sync_dir = corpus_root.join("worker1").join(SYNC_DIR);
2170 fs::create_dir(sync_dir.join("00000000-0000-0000-0000-000000000001-1.json")).unwrap();
2171 let target = Address::repeat_byte(0x11);
2172 let executor = sync_test_executor(corpus_root, target);
2173 let function = Function::parse("test(uint256)").unwrap();
2174
2175 let err = worker
2176 .calibrate(
2177 &executor,
2178 ReplayTarget {
2179 stateless: Some(StatelessReplayTarget { function: &function, address: target }),
2180 fuzzed_contracts: None,
2181 dynamic: None,
2182 },
2183 true,
2184 )
2185 .unwrap_err()
2186 .to_string();
2187
2188 assert!(err.contains("not a regular file"), "{err}");
2189 }
2190
2191 #[test]
2192 fn final_sync_completes_interleaved_corpus_lifecycle() {
2193 let corpus_root = temp_corpus_dir();
2194 let function = Function::parse("test(uint256)").unwrap();
2195 let target_address = Address::repeat_byte(0x11);
2196 let startup = CorpusEntry::new(vec![tx_for_function(
2197 target_address,
2198 &function,
2199 &[DynSolValue::Uint(U256::ZERO, 256)],
2200 )]);
2201 let startup_name = startup.file_name(false);
2202 fs::create_dir_all(corpus_root.join("worker1").join(CORPUS_DIR)).unwrap();
2203 startup.write_to_disk_in(&corpus_root.join("worker1").join(CORPUS_DIR), false).unwrap();
2204 let coordinator = Arc::new(CorpusSyncCoordinator::new(3));
2205 let setup = Arc::new(std::sync::Barrier::new(3));
2206
2207 std::thread::scope(|scope| {
2208 let root = corpus_root.clone();
2209 let master_startup_name = startup_name.clone();
2210 let master_coordinator = coordinator.clone();
2211 let master_setup = setup.clone();
2212 scope.spawn(move || {
2213 let function = Function::parse("test(uint256)").unwrap();
2214 let seed = WorkerCorpusSeed {
2215 replay_dirs: Some(canonical_replay_dirs(&root)),
2216 ..Default::default()
2217 };
2218 let mut master = worker_corpus(0, root.clone(), seed);
2219 let periodic = CorpusEntry::new(vec![tx_for_function(
2220 target_address,
2221 &function,
2222 &[DynSolValue::Uint(U256::ZERO, 256)],
2223 )]);
2224 let periodic_name = periodic.file_name(false);
2225 periodic.write_to_disk_in(&root.join("worker0").join(CORPUS_DIR), false).unwrap();
2226 master.push_corpus_entry(periodic);
2227
2228 let blocked = root.join("worker2").join(SYNC_DIR).join(&periodic_name);
2230 fs::create_dir_all(&blocked).unwrap();
2231 master.export_to_workers(3).unwrap();
2232 assert_eq!(master.new_entry_indices, [0]);
2233 assert!(root.join("worker2").join(SYNC_DIR).join(master_startup_name).is_file());
2234 fs::remove_dir(blocked).unwrap();
2235 master_setup.wait();
2236
2237 finalize_test_worker(&mut master, root, target_address, &master_coordinator);
2238 });
2239
2240 for id in 1..=2 {
2241 let root = corpus_root.clone();
2242 let coordinator = coordinator.clone();
2243 let setup = setup.clone();
2244 scope.spawn(move || {
2245 let mut worker = empty_worker_corpus(id, root.clone());
2246 setup.wait();
2247
2248 if id == 1 {
2249 let function = Function::parse("test(uint256)").unwrap();
2251 let mut late = CorpusEntry::new(vec![tx_for_function(
2252 target_address,
2253 &function,
2254 &[DynSolValue::Uint(U256::from(1), 256)],
2255 )]);
2256 late.timestamp = 0;
2257 worker.push_corpus_entry(late);
2258 }
2259
2260 finalize_test_worker(&mut worker, root, target_address, &coordinator);
2261 });
2262 }
2263 });
2264
2265 assert_eq!(read_corpus_dir(&corpus_root.join("worker2").join(CORPUS_DIR)).count(), 2);
2266 assert!(read_corpus_dir(&corpus_root.join("worker2").join(SYNC_DIR)).next().is_none());
2267 }
2268
2269 #[test]
2270 fn master_distributes_startup_corpus_only_once() {
2271 let corpus_root = temp_corpus_dir();
2272 let corpus = CorpusEntry::new(vec![basic_tx()]);
2273 let name = corpus.file_name(false);
2274 let non_master_corpus = corpus_root.join("worker1").join(CORPUS_DIR);
2275 fs::create_dir_all(&non_master_corpus).unwrap();
2276 corpus.write_to_disk_in(&non_master_corpus, false).unwrap();
2277 let seed = WorkerCorpusSeed {
2278 replay_dirs: Some(canonical_replay_dirs(&corpus_root)),
2279 ..Default::default()
2280 };
2281 let mut master = worker_corpus(0, corpus_root.clone(), seed);
2282
2283 let other_worker_sync_dir = corpus_root.join("worker2").join(SYNC_DIR);
2284 fs::create_dir_all(&other_worker_sync_dir).unwrap();
2285 let other_worker_sync = other_worker_sync_dir.join(&name);
2286 foundry_common::fs::write_json_file(&other_worker_sync, &vec![basic_tx_with_calldata([1])])
2287 .unwrap();
2288
2289 master.export_to_workers(3).unwrap();
2290 let source_worker_sync = corpus_root.join("worker1").join(SYNC_DIR).join(&name);
2291 assert!(source_worker_sync.is_file());
2292 assert!(other_worker_sync.is_file());
2293 assert!(master.initial_export_dirs.is_some());
2294
2295 fs::remove_file(&other_worker_sync).unwrap();
2296 master.export_to_workers(3).unwrap();
2297 assert!(other_worker_sync.is_file());
2298 assert!(master.initial_export_dirs.is_none());
2299
2300 let source_entry = read_corpus_dir(source_worker_sync.parent().unwrap()).next().unwrap();
2301 assert!(accept_synced_corpus_file(
2302 &source_entry,
2303 &source_entry.read_tx_seq().unwrap(),
2304 &non_master_corpus.join(&name),
2305 ));
2306 assert!(!source_worker_sync.exists());
2307 assert!(non_master_corpus.join(&name).is_file());
2308
2309 fs::remove_file(&other_worker_sync).unwrap();
2310 master.export_to_workers(3).unwrap();
2311 assert!(!other_worker_sync.exists());
2312
2313 let flat_root = temp_corpus_dir();
2314 let flat_corpus = CorpusEntry::new(vec![basic_tx()]);
2315 let flat_name = flat_corpus.file_name(false);
2316 flat_corpus.write_to_disk_in(&flat_root, false).unwrap();
2317 let config = corpus_config(flat_root.clone());
2318 let seed = WorkerCorpusSeed::load_from_disk::<foundry_evm_core::evm::EthEvmNetwork>(
2319 &config,
2320 None,
2321 None,
2322 ReplayTarget { stateless: None, fuzzed_contracts: None, dynamic: None },
2323 )
2324 .unwrap();
2325 fs::create_dir_all(flat_root.join("worker1").join(CORPUS_DIR)).unwrap();
2326 let generator =
2327 test_sequence(&config, TxGenerator::from_strategy(Just(basic_tx()).boxed()));
2328 let mut flat_master = WorkerCorpus::from_seed(0, config, generator, seed).unwrap();
2329 flat_master.worker_sync_enabled = true;
2330 let pending = CorpusEntry::new(vec![basic_tx_with_calldata([1])]);
2331 let pending_name = pending.file_name(false);
2332 pending.write_to_disk_in(&flat_root.join("worker0").join(CORPUS_DIR), false).unwrap();
2333 flat_master.push_corpus_entry(pending);
2334
2335 flat_master.export_to_workers(2).unwrap();
2336
2337 assert!(flat_root.join("worker1").join(SYNC_DIR).join(flat_name).is_file());
2338 assert!(flat_root.join("worker1").join(SYNC_DIR).join(pending_name).is_file());
2339 assert!(flat_master.new_entry_indices.is_empty());
2340 }
2341
2342 #[test]
2343 fn pending_master_fanout_entries_are_not_evicted() {
2344 let corpus_root = temp_corpus_dir();
2345 let mut master = empty_worker_corpus(0, corpus_root.clone());
2346 master.initial_export_dirs = None;
2347 let mut pending = CorpusEntry::new(vec![basic_tx_with_calldata([1])]);
2348 pending.total_mutations = 1;
2349 let retained = CorpusEntry::new(vec![basic_tx_with_calldata([2])]);
2350 let pending_name = pending.file_name(false);
2351 let retained_name = retained.file_name(false);
2352 let pending_timestamp = pending.timestamp;
2353 let retained_timestamp = retained.timestamp;
2354 let master_corpus = corpus_root.join("worker0").join(CORPUS_DIR);
2355 pending.write_to_disk_in(&master_corpus, false).unwrap();
2356 retained.write_to_disk_in(&master_corpus, false).unwrap();
2357 master.push_synced_corpus_entry(pending, pending_timestamp, pending_name.clone());
2358 master.push_synced_corpus_entry(retained, retained_timestamp, retained_name.clone());
2359
2360 master.evict_oldest_corpus().unwrap();
2361 master.export_to_workers(2).unwrap();
2362
2363 let worker_sync = corpus_root.join("worker1").join(SYNC_DIR);
2364 assert!(worker_sync.join(pending_name).is_file());
2365 assert!(worker_sync.join(retained_name).is_file());
2366 }
2367
2368 #[test]
2369 fn campaign_processing_writes_worker_file_immediately() {
2370 let corpus_root = temp_corpus_dir();
2371 let worker_subdir = corpus_root.join("worker1");
2372 let mut manager = empty_worker_corpus(1, corpus_root);
2373
2374 manager.process_inputs_for_campaign(&[basic_tx()], &[], true, None);
2375
2376 assert_eq!(manager.in_memory_corpus.len(), 1);
2377 assert_eq!(manager.metrics.corpus_count, 1);
2378 assert_eq!(read_corpus_dir(&worker_subdir.join(CORPUS_DIR)).count(), 1);
2379 }
2380
2381 fn edge_call(edge: EdgeKey, count: u8) -> RawCallResult {
2383 RawCallResult {
2384 edge_coverage: Some(EdgeCoverage::CollisionFree(vec![EdgeCovHit { edge, count }])),
2385 ..Default::default()
2386 }
2387 }
2388
2389 #[test]
2390 fn merge_edge_coverage_advances_timer_only_for_new_edges() {
2391 let corpus_root = temp_corpus_dir();
2392 let mut manager = empty_worker_corpus(1, corpus_root);
2393
2394 assert!(manager.time_since_new_edge().is_none());
2396 assert_eq!(manager.metrics.cumulative_edges_seen, 0);
2397
2398 let edge =
2399 EdgeKey { address: Address::ZERO, depth: None, pc: 0, jump_dest: U256::from(10) };
2400
2401 assert!(manager.merge_edge_coverage(&mut edge_call(edge, 1)));
2403 let first = manager.last_new_edge_at.expect("timer set after first new edge");
2404 assert_eq!(manager.metrics.cumulative_edges_seen, 1);
2405
2406 assert!(manager.merge_edge_coverage(&mut edge_call(edge, 8)));
2408 assert_eq!(manager.last_new_edge_at, Some(first));
2409 assert_eq!(manager.metrics.cumulative_edges_seen, 1);
2410 assert_eq!(manager.metrics.cumulative_features_seen, 1);
2411
2412 let other =
2414 EdgeKey { address: Address::ZERO, depth: None, pc: 1, jump_dest: U256::from(20) };
2415 assert!(manager.merge_edge_coverage(&mut edge_call(other, 1)));
2416 let second = manager.last_new_edge_at.expect("timer present");
2417 assert!(second >= first);
2418 assert_eq!(manager.metrics.cumulative_edges_seen, 2);
2419 assert!(manager.time_since_new_edge().is_some());
2420 }
2421
2422 #[test]
2423 fn empty_input_sequence_with_new_coverage_does_not_panic_or_insert() {
2424 let corpus_root = temp_corpus_dir();
2429 let worker_subdir = corpus_root.join("worker1");
2430 let mut manager = empty_worker_corpus(1, corpus_root);
2431
2432 manager.process_inputs_for_campaign(&[], &[], true, None);
2433
2434 assert_eq!(manager.in_memory_corpus.len(), 0);
2435 assert_eq!(manager.metrics.corpus_count, 0);
2436 assert_eq!(read_corpus_dir(&worker_subdir.join(CORPUS_DIR)).count(), 0);
2437
2438 manager.process_inputs(&[], &[], true, None);
2440 assert_eq!(manager.in_memory_corpus.len(), 0);
2441 assert_eq!(read_corpus_dir(&worker_subdir.join(CORPUS_DIR)).count(), 0);
2442 }
2443
2444 #[test]
2445 fn campaign_processing_defers_only_optimization_persistence() {
2446 let corpus_root = temp_corpus_dir();
2447 let mut manager = empty_worker_corpus(1, corpus_root.clone());
2448 let sequence = vec![basic_tx()];
2449 manager.process_inputs_for_campaign(
2450 &sequence,
2451 &[],
2452 false,
2453 Some((I256::try_from(7).unwrap(), sequence.clone())),
2454 );
2455
2456 let worker_corpus_dir = corpus_root.join("worker1").join(CORPUS_DIR);
2457 let entries = read_corpus_dir(&worker_corpus_dir).collect::<Vec<_>>();
2458 assert_eq!(entries.len(), 1);
2459 let persisted_sequence = entries[0].read_tx_seq().unwrap();
2460 assert_eq!(persisted_sequence.len(), sequence.len());
2461 assert_eq!(persisted_sequence[0].sender, sequence[0].sender);
2462 assert_eq!(persisted_sequence[0].call_details.target, sequence[0].call_details.target);
2463 assert_eq!(persisted_sequence[0].call_details.calldata, sequence[0].call_details.calldata);
2464 assert!(!corpus_root.join(OPTIMIZATION_BEST_FILE).exists());
2465
2466 persist_campaign_optimization(
2467 &corpus_config(corpus_root.clone()),
2468 Some(I256::try_from(7).unwrap()),
2469 &sequence,
2470 );
2471
2472 let state: OptimizationState =
2473 foundry_common::fs::read_json_file(&corpus_root.join(OPTIMIZATION_BEST_FILE)).unwrap();
2474 assert_eq!(state.best_value, I256::try_from(7).unwrap());
2475 assert_eq!(state.best_sequence.len(), sequence.len());
2476 assert_eq!(state.best_sequence[0].sender, sequence[0].sender);
2477 assert_eq!(state.best_sequence[0].call_details.target, sequence[0].call_details.target);
2478 assert_eq!(state.best_sequence[0].call_details.calldata, sequence[0].call_details.calldata);
2479 }
2480
2481 #[test]
2482 fn persisted_worker_corpus_entries_are_deduped_by_uuid() {
2483 let corpus_root = temp_corpus_dir();
2484 let corpus = CorpusEntry::new(vec![basic_tx()]);
2485 let duplicate = corpus.clone();
2486
2487 let worker0_corpus = corpus_root.join("worker0").join(CORPUS_DIR);
2488 let worker1_corpus = corpus_root.join("worker1").join(CORPUS_DIR);
2489 fs::create_dir_all(&worker0_corpus).unwrap();
2490 fs::create_dir_all(&worker1_corpus).unwrap();
2491 corpus.write_to_disk_in(&worker0_corpus, false).unwrap();
2492 duplicate.write_to_disk_in(&worker1_corpus, false).unwrap();
2493
2494 let mut seen = HashSet::new();
2495 let entries = unique_corpus_entries(&canonical_replay_dirs(&corpus_root), &mut seen)
2496 .collect::<Vec<_>>();
2497
2498 assert_eq!(entries.len(), 1);
2499 assert_eq!(entries[0].uuid, corpus.uuid);
2500 }
2501
2502 #[test]
2503 fn corpus_entry_write_uses_unparsable_temp_file() {
2504 let corpus_dir = temp_corpus_dir();
2505 let corpus = CorpusEntry::new(vec![basic_tx()]);
2506 let temp_path =
2507 corpus_dir.join(format!(".{}.{}.tmp", corpus.file_name(false), Uuid::new_v4()));
2508 fs::write(&temp_path, b"{").unwrap();
2509
2510 let path = corpus.write_to_disk_in(&corpus_dir, false).unwrap();
2511 let entries = read_corpus_dir(&corpus_dir).collect::<Vec<_>>();
2512
2513 assert_eq!(entries.len(), 1);
2514 assert_eq!(entries[0].path, path);
2515 assert!(temp_path.exists());
2516 }
2517
2518 #[test]
2519 fn persist_corpus_seed_skips_duplicate_sequence() {
2520 let corpus_root = temp_corpus_dir();
2521 let config = corpus_config(corpus_root.clone());
2522 let sequence = vec![basic_tx_with_calldata(vec![0x12, 0x34])];
2523
2524 let first = persist_corpus_seed(&config, sequence.clone()).unwrap().unwrap();
2525 let second = persist_corpus_seed(&config, sequence).unwrap().unwrap();
2526 let entries =
2527 read_corpus_dir(&corpus_root.join("worker0").join(CORPUS_DIR)).collect::<Vec<_>>();
2528
2529 assert_eq!(first, second);
2530 assert_eq!(entries.len(), 1);
2531 }
2532
2533 #[test]
2534 fn non_master_campaign_worker_uses_persisted_optimization_baseline() {
2535 let corpus_root = temp_corpus_dir();
2536 let persisted_sequence = vec![basic_tx()];
2537 let persisted_state = OptimizationState {
2538 best_value: I256::try_from(100).unwrap(),
2539 best_sequence: persisted_sequence,
2540 };
2541 foundry_common::fs::write_json_file(
2542 &corpus_root.join(OPTIMIZATION_BEST_FILE),
2543 &persisted_state,
2544 )
2545 .unwrap();
2546 let config = corpus_config(corpus_root);
2547 let generator =
2548 test_sequence(&config, TxGenerator::from_strategy(Just(basic_tx()).boxed()));
2549 let mut manager = WorkerCorpus::new::<foundry_evm_core::evm::EthEvmNetwork>(
2550 1,
2551 config,
2552 generator,
2553 None,
2554 None,
2555 ReplayTarget { stateless: None, fuzzed_contracts: None, dynamic: None },
2556 )
2557 .unwrap();
2558
2559 let worse_sequence = vec![basic_tx()];
2560 manager.process_inputs_for_campaign(
2561 &worse_sequence,
2562 &[],
2563 false,
2564 Some((I256::try_from(50).unwrap(), worse_sequence.clone())),
2565 );
2566
2567 let better_sequence = vec![basic_tx()];
2568 manager.process_inputs_for_campaign(
2569 &better_sequence,
2570 &[],
2571 false,
2572 Some((I256::try_from(150).unwrap(), better_sequence.clone())),
2573 );
2574 assert_eq!(manager.optimization_best_value, Some(I256::try_from(150).unwrap()));
2575 }
2576
2577 #[test]
2578 fn worker_can_initialize_from_warmed_seed() {
2579 let corpus_root = temp_corpus_dir();
2580 let tx_seq = vec![basic_tx()];
2581 let seed = WorkerCorpusSeed {
2582 in_memory_corpus: vec![CorpusEntry::new(tx_seq.clone())],
2583 history_map: vec![1, 2, 3],
2584 edge_indices: EdgeIndexMap::default(),
2585 sancov_history_map: vec![4, 5],
2586 metrics: CorpusMetrics {
2587 cumulative_edges_seen: 7,
2588 cumulative_features_seen: 11,
2589 corpus_count: 1,
2590 favored_items: 0,
2591 },
2592 replay_dirs: None,
2593 failed_replays: 13,
2594 optimization_best_value: Some(I256::try_from(17).unwrap()),
2595 optimization_best_sequence: tx_seq,
2596 last_new_edge_at: None,
2597 };
2598
2599 let config = corpus_config(corpus_root);
2600 let generator =
2601 test_sequence(&config, TxGenerator::from_strategy(Just(basic_tx()).boxed()));
2602 let manager = WorkerCorpus::from_seed(1, config, generator, seed).unwrap();
2603
2604 assert_eq!(manager.in_memory_corpus.len(), 1);
2605 assert_eq!(manager.history_map, vec![1, 2, 3]);
2606 assert_eq!(manager.sancov_history_map, vec![4, 5]);
2607 assert_eq!(manager.metrics.cumulative_edges_seen, 7);
2608 assert_eq!(manager.metrics.cumulative_features_seen, 11);
2609 assert_eq!(manager.metrics.corpus_count, 1);
2610 assert_eq!(manager.failed_replays, 13);
2611 let (value, sequence) = manager.optimization_initial_state();
2612 assert_eq!(value, Some(I256::try_from(17).unwrap()));
2613 assert_eq!(sequence.len(), 1);
2614 }
2615
2616 #[test]
2617 fn clone_for_worker_shards_warmed_corpus_and_recomputes_metrics() {
2618 let entries = (0..10)
2619 .map(|idx| {
2620 let mut entry = CorpusEntry::new(vec![basic_tx()]);
2621 entry.is_favored = idx % 2 == 0;
2622 entry
2623 })
2624 .collect::<Vec<_>>();
2625 let entry_ids = entries.iter().map(|entry| entry.uuid).collect::<Vec<_>>();
2626 let seed = WorkerCorpusSeed {
2627 in_memory_corpus: entries,
2628 history_map: vec![1, 2, 3],
2629 edge_indices: EdgeIndexMap::default(),
2630 sancov_history_map: vec![4, 5],
2631 metrics: CorpusMetrics {
2632 cumulative_edges_seen: 7,
2633 cumulative_features_seen: 11,
2634 corpus_count: 10,
2635 favored_items: 5,
2636 },
2637 replay_dirs: None,
2638 failed_replays: 13,
2639 optimization_best_value: Some(I256::try_from(17).unwrap()),
2640 optimization_best_sequence: vec![basic_tx()],
2641 last_new_edge_at: None,
2642 };
2643
2644 let worker_count = 3;
2645 let shards = (0..worker_count)
2646 .map(|worker_id| seed.clone_for_worker(worker_id, worker_count, true))
2647 .collect::<Vec<_>>();
2648 let mut sharded_ids = shards
2649 .iter()
2650 .flat_map(|shard| shard.in_memory_corpus.iter().map(|entry| entry.uuid))
2651 .collect::<Vec<_>>();
2652 let mut expected_ids = entry_ids.clone();
2653 sharded_ids.sort_unstable();
2654 expected_ids.sort_unstable();
2655
2656 assert_eq!(sharded_ids, expected_ids);
2657 assert_eq!(
2658 shards[0].in_memory_corpus.iter().map(|entry| entry.uuid).collect::<Vec<_>>(),
2659 [entry_ids[0], entry_ids[3], entry_ids[6], entry_ids[9]]
2660 );
2661 assert_eq!(
2662 shards[1].in_memory_corpus.iter().map(|entry| entry.uuid).collect::<Vec<_>>(),
2663 [entry_ids[1], entry_ids[4], entry_ids[7]]
2664 );
2665 assert_eq!(
2666 shards[2].in_memory_corpus.iter().map(|entry| entry.uuid).collect::<Vec<_>>(),
2667 [entry_ids[2], entry_ids[5], entry_ids[8]]
2668 );
2669 assert_eq!(
2670 shards.iter().map(|shard| shard.in_memory_corpus.len()).collect::<Vec<_>>(),
2671 [4, 3, 3]
2672 );
2673 assert_eq!(
2674 shards.iter().map(|shard| shard.metrics.corpus_count).collect::<Vec<_>>(),
2675 [4, 3, 3]
2676 );
2677 assert_eq!(
2678 shards.iter().map(|shard| shard.metrics.favored_items).collect::<Vec<_>>(),
2679 [2, 1, 2]
2680 );
2681 assert!(shards.iter().all(|shard| shard.history_map == seed.history_map));
2682 assert!(shards.iter().all(|shard| shard.sancov_history_map == seed.sancov_history_map));
2683 assert!(shards.iter().all(|shard| shard.metrics.cumulative_edges_seen == 7));
2684 assert!(shards.iter().all(|shard| shard.metrics.cumulative_features_seen == 11));
2685 }
2686
2687 #[test]
2688 fn clone_for_worker_can_strip_cmp_sequences() {
2689 let cmp = CmpOperands {
2690 op1: U256::from(1),
2691 op2: U256::from(2),
2692 pc: 3,
2693 address: Address::ZERO,
2694 opcode: 0,
2695 };
2696 let entries = (0..2)
2697 .map(|_| {
2698 CorpusEntry::new_with_cmp(
2699 vec![basic_tx()],
2700 vec![vec![ComparisonHint { lhs: cmp.op1, rhs: cmp.op2 }]],
2701 Uuid::new_v4(),
2702 )
2703 })
2704 .collect::<Vec<_>>();
2705 let seed = WorkerCorpusSeed { in_memory_corpus: entries, ..Default::default() };
2706
2707 let with_cmp = seed.clone_for_worker(0, 1, true);
2708 let without_cmp = seed.clone_for_worker(0, 1, false);
2709
2710 assert!(with_cmp.in_memory_corpus.iter().all(|entry| !entry.cmp_seq[0].is_empty()));
2711 assert!(without_cmp.in_memory_corpus.iter().all(|entry| entry.cmp_seq.is_empty()));
2712 }
2713
2714 #[test]
2715 fn retain_replayable_removes_off_target_corpus_entries() {
2716 let target = Address::from([0x11; 20]);
2717 let foo = Function::parse("foo()").unwrap();
2718 let bar = Function::parse("bar()").unwrap();
2719 let foo_selector = foo.selector();
2720 let foo_tx = tx_for_function(target, &foo, &[]);
2721 let bar_tx = tx_for_function(target, &bar, &[]);
2722 let mut foo_entry = CorpusEntry::new(vec![foo_tx.clone()]);
2723 foo_entry.is_favored = true;
2724 let mut bar_entry = CorpusEntry::new(vec![bar_tx.clone()]);
2725 bar_entry.is_favored = true;
2726 let mut seed = WorkerCorpusSeed {
2727 in_memory_corpus: vec![foo_entry, bar_entry],
2728 metrics: CorpusMetrics { corpus_count: 2, favored_items: 2, ..Default::default() },
2729 optimization_best_value: Some(I256::try_from(17).unwrap()),
2730 optimization_best_sequence: vec![bar_tx],
2731 ..Default::default()
2732 };
2733 let targeted_contracts =
2734 targeted_contracts_with_selective_functions(target, vec![foo, bar], [foo_selector]);
2735 let targets = targeted_contracts.targets();
2736
2737 seed.retain_replayable(&targets);
2738
2739 assert_eq!(seed.in_memory_corpus.len(), 1);
2740 assert_eq!(seed.in_memory_corpus[0].tx_seq.len(), 1);
2741 assert_eq!(
2742 seed.in_memory_corpus[0].tx_seq[0].call_details.target,
2743 foo_tx.call_details.target
2744 );
2745 assert_eq!(
2746 seed.in_memory_corpus[0].tx_seq[0].call_details.calldata,
2747 foo_tx.call_details.calldata
2748 );
2749 assert_eq!(seed.metrics.corpus_count, 1);
2750 assert_eq!(seed.metrics.favored_items, 1);
2751 assert!(seed.optimization_best_value.is_none());
2752 assert!(seed.optimization_best_sequence.is_empty());
2753 }
2754
2755 #[test]
2756 fn hoist_observed_calls_bundles_replayable_subcalls_into_one_corpus_entry() {
2757 let target = Address::from([0x42; 20]);
2758 let other = Address::from([0x43; 20]);
2759 let sender = Address::from([0xaa; 20]);
2760 let observed_caller = Address::from([0xbb; 20]);
2761 let foo = Function::parse("foo(uint256)").unwrap();
2762 let bar = Function::parse("bar()").unwrap();
2763 let foo_selector = foo.selector();
2764 let bar_selector = bar.selector();
2765 let targeted_contracts = targeted_contracts_with_selective_functions(
2766 target,
2767 vec![foo, bar],
2768 [foo_selector, bar_selector],
2769 );
2770
2771 let mut foo_calldata = vec![0u8; 36];
2772 foo_calldata[..4].copy_from_slice(&foo_selector[..]);
2773 let bar_calldata = bar_selector.to_vec();
2774 let mut unknown_selector = vec![0u8; 36];
2775 unknown_selector[..4].copy_from_slice(&[0xde, 0xad, 0xbe, 0xef]);
2776 let value = U256::from(1);
2777
2778 let observed = vec![
2779 ObservedCall {
2780 depth: 1,
2781 caller: observed_caller,
2782 target: other,
2783 calldata: Bytes::from(foo_calldata.clone()),
2784 value: Some(value),
2785 },
2786 ObservedCall {
2787 depth: 1,
2788 caller: observed_caller,
2789 target,
2790 calldata: Bytes::from(foo_calldata),
2791 value: None,
2792 },
2793 ObservedCall {
2794 depth: 2,
2795 caller: observed_caller,
2796 target,
2797 calldata: Bytes::from(bar_calldata),
2798 value: None,
2799 },
2800 ObservedCall {
2801 depth: 1,
2802 caller: observed_caller,
2803 target,
2804 calldata: Bytes::from(unknown_selector),
2805 value: None,
2806 },
2807 ObservedCall {
2808 depth: 1,
2809 caller: observed_caller,
2810 target,
2811 calldata: Bytes::from(vec![0u8; 3]),
2812 value: None,
2813 },
2814 ];
2815 let parent_tx = BasicTxDetails {
2816 warp: Some(U256::from(123)),
2817 roll: Some(U256::from(456)),
2818 sender,
2819 call_details: CallDetails {
2820 target: Address::from([0x99; 20]),
2821 calldata: Bytes::new(),
2822 value: None,
2823 },
2824 };
2825 let mut manager = empty_worker_corpus(0, temp_corpus_dir());
2826
2827 manager.hoist_observed_calls(
2828 &observed,
2829 &parent_tx,
2830 &targeted_contracts,
2831 CorpusInsertionMode::Live,
2832 );
2833
2834 assert_eq!(manager.in_memory_corpus.len(), 1);
2835 assert_eq!(manager.metrics.corpus_count, 1);
2836
2837 let entry = manager.in_memory_corpus.last().unwrap();
2838 assert_eq!(entry.tx_seq.len(), 2);
2839 let tx = &entry.tx_seq[0];
2840 assert_eq!(tx.warp, parent_tx.warp);
2841 assert_eq!(tx.roll, parent_tx.roll);
2842 assert_eq!(tx.sender, observed_caller);
2843 assert_eq!(tx.call_details.target, target);
2844 assert_eq!(&tx.call_details.calldata[..4], &foo_selector[..]);
2845 assert_eq!(tx.call_details.value, None);
2846
2847 let tx = &entry.tx_seq[1];
2848 assert_eq!(tx.warp, None);
2849 assert_eq!(tx.roll, None);
2850 assert_eq!(tx.sender, observed_caller);
2851 assert_eq!(tx.call_details.target, target);
2852 assert_eq!(&tx.call_details.calldata[..4], &bar_selector[..]);
2853 assert_eq!(tx.call_details.value, None);
2854 }
2855
2856 #[test]
2857 fn hoist_observed_calls_persists_immediately() {
2858 let target = Address::from([0x42; 20]);
2859 let foo = Function::parse("foo()").unwrap();
2860 let selector = foo.selector();
2861 let targeted_contracts = targeted_contracts_with_selective_functions(target, vec![foo], []);
2862 let observed = vec![ObservedCall {
2863 depth: 1,
2864 caller: Address::from([0xaa; 20]),
2865 target,
2866 calldata: Bytes::from(selector.to_vec()),
2867 value: None,
2868 }];
2869 let corpus_root = temp_corpus_dir();
2870 let worker_corpus_dir = corpus_root.join("worker1").join(CORPUS_DIR);
2871 let mut manager = empty_worker_corpus(1, corpus_root);
2872
2873 manager.hoist_observed_calls(
2874 &observed,
2875 &basic_tx(),
2876 &targeted_contracts,
2877 CorpusInsertionMode::Live,
2878 );
2879
2880 assert_eq!(manager.in_memory_corpus.len(), 1);
2881 assert_eq!(read_corpus_dir(&worker_corpus_dir).count(), 1);
2882 }
2883
2884 #[test]
2885 fn hoist_observed_calls_skips_empty_or_non_coverage_guided_inputs() {
2886 let target = Address::from([0x42; 20]);
2887 let foo = Function::parse("foo()").unwrap();
2888 let selector = foo.selector();
2889 let targeted_contracts = targeted_contracts_with_selective_functions(target, vec![foo], []);
2890 let observed = vec![ObservedCall {
2891 depth: 1,
2892 caller: Address::from([0xaa; 20]),
2893 target,
2894 calldata: Bytes::from(selector.to_vec()),
2895 value: None,
2896 }];
2897
2898 let mut no_corpus_config = corpus_config(temp_corpus_dir());
2899 no_corpus_config.corpus_dir = None;
2900 let generator =
2901 test_sequence(&no_corpus_config, TxGenerator::from_strategy(Just(basic_tx()).boxed()));
2902 let mut manager =
2903 WorkerCorpus::from_seed(0, no_corpus_config, generator, WorkerCorpusSeed::default())
2904 .unwrap();
2905 manager.hoist_observed_calls(
2906 &observed,
2907 &basic_tx(),
2908 &targeted_contracts,
2909 CorpusInsertionMode::Live,
2910 );
2911 assert!(manager.in_memory_corpus.is_empty());
2912
2913 let mut manager = empty_worker_corpus(0, temp_corpus_dir());
2914 manager.hoist_observed_calls(
2915 &[],
2916 &basic_tx(),
2917 &targeted_contracts,
2918 CorpusInsertionMode::Live,
2919 );
2920 assert!(manager.in_memory_corpus.is_empty());
2921 }
2922
2923 #[test]
2924 fn sequence_from_observed_keeps_only_direct_replayable_calls() {
2925 let target = Address::from([0x42; 20]);
2926 let other = Address::from([0x43; 20]);
2927 let sender = Address::from([0xaa; 20]);
2928 let nested_caller = Address::from([0xbb; 20]);
2929 let foo = Function::parse("foo(uint256)").unwrap();
2930 let bar = Function::parse("bar()").unwrap();
2931 let foo_selector = foo.selector();
2932 let bar_selector = bar.selector();
2933 let targeted_contracts =
2934 targeted_contracts_with_selective_functions(target, vec![foo, bar], [foo_selector]);
2935 let targets = targeted_contracts.targets();
2936
2937 let mut foo_calldata = vec![0u8; 36];
2938 foo_calldata[..4].copy_from_slice(&foo_selector[..]);
2939 let bar_calldata = bar_selector.to_vec();
2940 let observed = vec![
2941 ObservedCall {
2942 depth: 1,
2943 caller: sender,
2944 target,
2945 calldata: Bytes::from(foo_calldata.clone()),
2946 value: None,
2947 },
2948 ObservedCall {
2949 depth: 2,
2950 caller: nested_caller,
2951 target,
2952 calldata: Bytes::from(foo_calldata),
2953 value: None,
2954 },
2955 ObservedCall {
2956 depth: 1,
2957 caller: sender,
2958 target,
2959 calldata: Bytes::from(bar_calldata),
2960 value: None,
2961 },
2962 ObservedCall {
2963 depth: 1,
2964 caller: sender,
2965 target: other,
2966 calldata: Bytes::from(foo_selector.to_vec()),
2967 value: None,
2968 },
2969 ];
2970
2971 let seq = sequence_from_observed(&observed, &targets, ObservedCallDepth::DirectOnly, None);
2972
2973 assert_eq!(seq.len(), 1);
2974 assert_eq!(seq[0].sender, sender);
2975 assert_eq!(seq[0].call_details.target, target);
2976 assert_eq!(&seq[0].call_details.calldata[..4], &foo_selector[..]);
2977 }
2978
2979 #[test]
2980 fn push_observed_sequence_live_persists_and_memory_only_does_not() {
2981 let corpus_root = temp_corpus_dir();
2982 let worker0_corpus_dir = corpus_root.join("worker0").join(CORPUS_DIR);
2983 let mut manager = empty_worker_corpus(0, corpus_root.clone());
2984
2985 manager.push_observed_sequence(vec![basic_tx()], CorpusInsertionMode::Live);
2986 assert_eq!(manager.in_memory_corpus.len(), 1);
2987 assert_eq!(read_corpus_dir(&worker0_corpus_dir).count(), 1);
2988
2989 let mut manager = empty_worker_corpus(1, corpus_root.clone());
2990 let worker1_corpus_dir = corpus_root.join("worker1").join(CORPUS_DIR);
2991 manager.push_observed_sequence(vec![basic_tx()], CorpusInsertionMode::MemoryOnly);
2992 assert_eq!(manager.in_memory_corpus.len(), 1);
2993 assert_eq!(read_corpus_dir(&worker1_corpus_dir).count(), 0);
2994 }
2995
2996 #[test]
2997 fn detects_legacy_invariant_corpus_dirs_without_matching_worker_dirs() {
2998 let corpus_root = temp_corpus_dir();
2999 fs::create_dir_all(corpus_root.join("worker0")).unwrap();
3000 assert!(!has_legacy_invariant_corpus_dirs(&corpus_root));
3001
3002 fs::create_dir_all(corpus_root.join("invariant_a")).unwrap();
3003 assert!(has_legacy_invariant_corpus_dirs(&corpus_root));
3004 }
3005
3006 #[test]
3007 fn ignores_optimization_invariant_corpus_dirs_when_detecting_legacy_dirs() {
3008 let corpus_root = temp_corpus_dir();
3009 fs::create_dir_all(corpus_root.join("worker0")).unwrap();
3010 let optimization_dir = corpus_root.join("invariant_optimize");
3011 fs::create_dir_all(optimization_dir.join("worker0")).unwrap();
3012 fs::write(optimization_dir.join(OPTIMIZATION_BEST_FILE), "{}").unwrap();
3013
3014 assert!(!has_legacy_invariant_corpus_dirs(&corpus_root));
3015
3016 fs::create_dir_all(corpus_root.join("invariant_legacy").join("worker0")).unwrap();
3017 assert!(has_legacy_invariant_corpus_dirs(&corpus_root));
3018 }
3019
3020 #[test]
3021 fn favored_sets_true_and_metrics_increment_when_ratio_gt_threshold() {
3022 let (mut manager, uuid) = new_manager_with_single_corpus();
3023 let corpus = manager.in_memory_corpus.iter_mut().find(|c| c.uuid == uuid).unwrap();
3024 corpus.total_mutations = 4;
3025 corpus.new_finds_produced = 2; corpus.is_favored = false;
3027
3028 assert_eq!(manager.metrics.favored_items, 0);
3030
3031 manager.current_mutated_index = Some(0);
3033 manager.process_inputs(&[basic_tx()], &[], true, None);
3034
3035 let corpus = manager.in_memory_corpus.iter().find(|c| c.uuid == uuid).unwrap();
3036 assert!(corpus.is_favored, "expected favored to be true when ratio > threshold");
3037 assert_eq!(
3038 manager.metrics.favored_items, 1,
3039 "favored_items should increment on false→true"
3040 );
3041 }
3042
3043 #[test]
3044 fn favored_sets_false_and_metrics_decrement_when_ratio_lt_threshold() {
3045 let (mut manager, uuid) = new_manager_with_single_corpus();
3046 let corpus = manager.in_memory_corpus.iter_mut().find(|c| c.uuid == uuid).unwrap();
3047 corpus.total_mutations = 9;
3048 corpus.new_finds_produced = 3; corpus.is_favored = true; manager.metrics.favored_items = 1;
3052
3053 manager.current_mutated_index = Some(0);
3055 manager.process_inputs(&[basic_tx()], &[], false, None);
3056
3057 let corpus = manager.in_memory_corpus.iter().find(|c| c.uuid == uuid).unwrap();
3058 assert!(!corpus.is_favored, "expected favored to be false when ratio < threshold");
3059 assert_eq!(
3060 manager.metrics.favored_items, 0,
3061 "favored_items should decrement on true→false"
3062 );
3063 }
3064
3065 #[test]
3066 fn favored_is_false_on_ratio_equal_threshold() {
3067 let (mut manager, uuid) = new_manager_with_single_corpus();
3068 let corpus = manager.in_memory_corpus.iter_mut().find(|c| c.uuid == uuid).unwrap();
3069 corpus.total_mutations = 9;
3071 corpus.new_finds_produced = 2;
3072 corpus.is_favored = false;
3073
3074 manager.current_mutated_index = Some(0);
3075 manager.process_inputs(&[basic_tx()], &[], true, None);
3076
3077 let corpus = manager.in_memory_corpus.iter().find(|c| c.uuid == uuid).unwrap();
3078 assert!(
3079 !(corpus.is_favored),
3080 "with strict '>' comparison, favored must be false when ratio == threshold"
3081 );
3082 }
3083
3084 #[test]
3085 fn eviction_skips_favored_and_evicts_non_favored() {
3086 let mut favored = CorpusEntry::new(vec![basic_tx()]);
3088 favored.total_mutations = 2;
3089 favored.is_favored = true;
3090
3091 let mut non_favored = CorpusEntry::new(vec![basic_tx()]);
3092 non_favored.total_mutations = 2;
3093 non_favored.is_favored = false;
3094 let non_favored_uuid = non_favored.uuid;
3095
3096 let mut manager = seeded_worker_corpus(0, temp_corpus_dir(), vec![favored, non_favored]);
3097
3098 manager.evict_oldest_corpus().unwrap();
3100 assert_eq!(manager.in_memory_corpus.len(), 1);
3101 assert!(manager.in_memory_corpus.iter().all(|c| c.is_favored));
3102
3103 manager.evict_oldest_corpus().unwrap();
3105 assert_eq!(manager.in_memory_corpus.len(), 1, "favored corpus must not be evicted");
3106
3107 assert!(manager.in_memory_corpus.iter().all(|c| c.uuid != non_favored_uuid));
3109 }
3110
3111 #[test]
3112 fn non_synchronizing_entries_remain_evictable() {
3113 let corpus_root = temp_corpus_dir();
3114 let config = corpus_config(corpus_root);
3115 let generator =
3116 test_sequence(&config, TxGenerator::from_strategy(Just(basic_tx()).boxed()));
3117 let mut manager =
3118 WorkerCorpus::from_seed(0, config, generator, WorkerCorpusSeed::default()).unwrap();
3119 let mut evictable = CorpusEntry::new(vec![basic_tx_with_calldata([1])]);
3120 evictable.total_mutations = 1;
3121 let evictable_uuid = evictable.uuid;
3122 manager.push_corpus_entry(evictable);
3123 manager.push_corpus_entry(CorpusEntry::new(vec![basic_tx_with_calldata([2])]));
3124
3125 manager.evict_oldest_corpus().unwrap();
3126
3127 assert_eq!(manager.in_memory_corpus.len(), 1);
3128 assert!(manager.in_memory_corpus.iter().all(|entry| entry.uuid != evictable_uuid));
3129 assert!(manager.new_entry_indices.is_empty());
3130 }
3131}