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