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