Skip to main content

foundry_evm/executors/
corpus.rs

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