Skip to main content

foundry_evm/executors/
corpus.rs

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