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