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