Skip to main content

foundry_evm/executors/invariant/
campaign.rs

1use super::{
2    FailureKey, InvariantFailureMetrics, InvariantFailures, InvariantFuzzError,
3    InvariantFuzzTestResult, InvariantMetrics,
4};
5use crate::executors::{EarlyExit, EvmExecutionCancellation, corpus::CampaignCorpusEntry};
6use alloy_primitives::{Address, I256, Selector};
7use eyre::{Result, ensure};
8use foundry_evm_coverage::HitMaps;
9use foundry_evm_fuzz::BasicTxDetails;
10use std::{
11    collections::{HashMap, HashSet},
12    sync::{
13        Arc, Mutex,
14        atomic::{AtomicBool, AtomicU32, AtomicU64, Ordering},
15    },
16    time::{Duration, Instant},
17};
18
19/// Immutable plan-level description for an invariant campaign.
20///
21/// This is only a planning contract for splitting one logical campaign into worker ranges. It does
22/// not start workers, choose worker counts, or decide corpus/failure persistence.
23#[derive(Clone, Copy, Debug, PartialEq, Eq)]
24pub struct InvariantCampaignSpec {
25    /// Total logical runs configured for the campaign.
26    pub total_runs: u32,
27}
28
29impl InvariantCampaignSpec {
30    pub const fn new(total_runs: u32) -> Self {
31        Self { total_runs }
32    }
33
34    /// Partitions the logical campaign into contiguous worker run ranges.
35    ///
36    /// This only describes work assignment. It does not start worker execution and does not
37    /// attribute failures to worker/run origins.
38    pub fn worker_plans(self, workers: usize) -> Result<Vec<InvariantWorkerPlan>> {
39        ensure!(workers > 0, "invariant campaign requires at least one worker");
40
41        if self.total_runs == 0 {
42            return Ok(vec![InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 0 }]);
43        }
44
45        let worker_count = workers.min(self.total_runs as usize) as u32;
46        let base_runs = self.total_runs / worker_count;
47        let extra_runs = self.total_runs % worker_count;
48
49        let mut first_global_run = 0;
50        let mut plans = Vec::with_capacity(worker_count as usize);
51        for worker_id in 0..worker_count {
52            let runs = base_runs + u32::from(worker_id < extra_runs);
53            plans.push(InvariantWorkerPlan { worker_id, first_global_run, runs });
54            first_global_run += runs;
55        }
56
57        debug_assert_eq!(first_global_run, self.total_runs);
58        Ok(plans)
59    }
60}
61
62/// Static assignment of a contiguous logical run range to one worker.
63///
64/// The assigned range is `[first_global_run, first_global_run + runs)`.
65/// Worker `0` is the master worker for master-only artifacts such as persisted corpus replay
66/// counts.
67#[derive(Clone, Copy, Debug, PartialEq, Eq)]
68pub struct InvariantWorkerPlan {
69    pub worker_id: u32,
70    pub first_global_run: u32,
71    pub runs: u32,
72}
73
74/// Shared state used only to coordinate invariant worker execution.
75pub struct InvariantCampaignState {
76    started_at: Instant,
77    timed: bool,
78    total_runs: AtomicU32,
79    total_txs: AtomicU64,
80    total_gas: AtomicU64,
81    cancellation: EvmExecutionCancellation,
82    last_metrics_report: Mutex<Instant>,
83    failure_metrics: Mutex<CampaignFailureMetrics>,
84}
85
86#[derive(Default)]
87struct CampaignFailureMetrics {
88    metrics: InvariantFailureMetrics,
89    handler_sites: HashSet<(Address, Selector)>,
90}
91
92impl InvariantCampaignState {
93    pub fn new(early_exit: EarlyExit, timeout: Option<u32>) -> Self {
94        let started_at = Instant::now();
95        let deadline = timeout
96            .map(|timeout| Duration::from_secs(timeout.into()))
97            .and_then(|timeout| started_at.checked_add(timeout));
98        Self {
99            started_at,
100            timed: timeout.is_some(),
101            total_runs: AtomicU32::new(0),
102            total_txs: AtomicU64::new(0),
103            total_gas: AtomicU64::new(0),
104            cancellation: EvmExecutionCancellation::campaign(
105                early_exit,
106                Arc::new(AtomicBool::new(false)),
107                deadline,
108            ),
109            last_metrics_report: Mutex::new(started_at),
110            failure_metrics: Mutex::new(CampaignFailureMetrics::default()),
111        }
112    }
113
114    pub fn increment_runs(&self) -> u32 {
115        self.total_runs.fetch_add(1, Ordering::Relaxed) + 1
116    }
117
118    #[cfg(test)]
119    pub fn total_runs(&self) -> u32 {
120        self.total_runs.load(Ordering::Relaxed)
121    }
122
123    pub fn record_call(&self, gas_used: u64) {
124        self.total_txs.fetch_add(1, Ordering::Relaxed);
125        self.total_gas.fetch_add(gas_used, Ordering::Relaxed);
126    }
127
128    pub fn throughput_totals(&self) -> (u64, u64) {
129        (self.total_txs.load(Ordering::Relaxed), self.total_gas.load(Ordering::Relaxed))
130    }
131
132    pub fn elapsed(&self) -> Duration {
133        self.started_at.elapsed()
134    }
135
136    pub const fn is_timed_campaign(&self) -> bool {
137        self.timed
138    }
139
140    pub fn should_stop(&self) -> bool {
141        self.cancellation.should_stop(true)
142    }
143
144    pub fn request_terminal_stop(&self) {
145        self.cancellation.request_stop();
146    }
147
148    pub fn should_emit_metrics_report(&self, interval: Duration) -> bool {
149        let mut last_report =
150            self.last_metrics_report.lock().expect("metrics report lock poisoned");
151        if last_report.elapsed() <= interval {
152            return false;
153        }
154
155        *last_report = Instant::now();
156        true
157    }
158
159    pub(super) fn record_invariant_failure(
160        &self,
161        invariant_name: &str,
162        target: &str,
163        reason: &str,
164    ) {
165        let mut failure_metrics =
166            self.failure_metrics.lock().expect("failure metrics lock poisoned");
167        if !failure_metrics.metrics.unique_failures.contains(invariant_name) {
168            failure_metrics.metrics.record_failure(invariant_name, target, reason);
169        }
170    }
171
172    pub(super) fn sync_handler_failures(&self, failures: &InvariantFailures) {
173        let mut failure_metrics =
174            self.failure_metrics.lock().expect("failure metrics lock poisoned");
175        for (key, error) in &failures.failures {
176            let FailureKey::Handler(addr, selector) = key else { continue };
177            if failure_metrics.handler_sites.insert((*addr, *selector)) {
178                let reason = error.revert_reason().unwrap_or_default();
179                failure_metrics.metrics.record_handler_failure(*addr, *selector, &reason);
180            }
181        }
182        debug_assert_eq!(
183            failure_metrics.metrics.broken_handlers,
184            failure_metrics.handler_sites.len()
185        );
186    }
187
188    pub(super) fn failure_metrics(&self) -> InvariantFailureMetrics {
189        self.failure_metrics.lock().expect("failure metrics lock poisoned").metrics.clone()
190    }
191
192    pub const fn early_exit(&self) -> &EarlyExit {
193        self.cancellation.early_exit_ref()
194    }
195
196    pub const fn cancellation(&self) -> &EvmExecutionCancellation {
197        &self.cancellation
198    }
199}
200
201/// Output produced by one invariant worker.
202///
203/// This is a data envelope for aggregation only. It does not imply that this module executed the
204/// worker, shrank failures, or wrote any persisted corpus/failure files.
205#[derive(Debug)]
206pub struct InvariantWorkerOutput {
207    pub plan: InvariantWorkerPlan,
208    pub result: InvariantFuzzTestResult,
209    pub corpus_entries: Vec<CampaignCorpusEntry>,
210}
211
212impl InvariantWorkerOutput {
213    #[cfg(test)]
214    pub const fn new(plan: InvariantWorkerPlan, result: InvariantFuzzTestResult) -> Self {
215        Self { plan, result, corpus_entries: Vec::new() }
216    }
217}
218
219/// Merges worker outputs back into one logical invariant campaign result.
220///
221/// Merge policy:
222/// - outputs are folded in `first_global_run` order;
223/// - predicate failures keep the first failure in logical run order;
224/// - handler assertion failures keep the shorter reproducer, with equal lengths preserving the
225///   earlier logical worker;
226/// - optimization mode keeps the maximum value, with ties preserving the earlier logical worker;
227/// - `failed_corpus_replays` is a master-worker-only value from worker `0`;
228/// - run/call counts, reverts, gas traces, selector metrics, and line coverage accumulate into the
229///   logical campaign result.
230#[derive(Debug)]
231pub struct InvariantCampaignAggregator {
232    spec: InvariantCampaignSpec,
233    outputs: Vec<InvariantWorkerOutput>,
234}
235
236impl InvariantCampaignAggregator {
237    pub const fn new(spec: InvariantCampaignSpec) -> Self {
238        Self { spec, outputs: Vec::new() }
239    }
240
241    pub fn push(&mut self, output: InvariantWorkerOutput) {
242        self.outputs.push(output);
243    }
244
245    /// Validates the collected worker ranges and folds them into one logical campaign result.
246    #[cfg(test)]
247    pub fn finish(self) -> Result<InvariantFuzzTestResult> {
248        Ok(self.finish_with_corpus_entries()?.0)
249    }
250
251    /// Validates the collected worker ranges and folds them into one logical campaign result with
252    /// corpus artifacts selected in logical worker order.
253    pub fn finish_with_corpus_entries(
254        mut self,
255    ) -> Result<(InvariantFuzzTestResult, Vec<CampaignCorpusEntry>)> {
256        ensure!(!self.outputs.is_empty(), "missing invariant worker output");
257
258        self.outputs.sort_by_key(|output| output.plan.first_global_run);
259        ensure_outputs_cover_campaign(self.spec, &self.outputs)?;
260        fold_outputs(self.outputs)
261    }
262
263    /// Folds timeout worker outputs without requiring full logical campaign coverage.
264    ///
265    /// Timeout campaigns share a wall-clock deadline across workers. When the deadline hits, any
266    /// worker may have completed fewer than its assigned runs, so the original static ranges can
267    /// contain gaps. The merge still validates worker identity and preserves deterministic worker
268    /// order, but final run count is derived from the completed worker counters.
269    pub fn finish_partial_with_corpus_entries(
270        mut self,
271    ) -> Result<(InvariantFuzzTestResult, Vec<CampaignCorpusEntry>)> {
272        ensure!(!self.outputs.is_empty(), "missing invariant worker output");
273
274        self.outputs.sort_by_key(|output| output.plan.first_global_run);
275        ensure_worker_ids_are_valid(&self.outputs)?;
276        fold_outputs(self.outputs)
277    }
278}
279
280fn fold_outputs(
281    outputs: Vec<InvariantWorkerOutput>,
282) -> Result<(InvariantFuzzTestResult, Vec<CampaignCorpusEntry>)> {
283    let workers = outputs.len();
284    let mut errors = HashMap::default();
285    let mut handler_errors = HashMap::default();
286    let mut runs = 0;
287    let mut calls = 0;
288    let mut reverts = 0;
289    let mut last_run_inputs = Vec::new();
290    let mut gas_report_traces = Vec::new();
291    let mut line_coverage = None;
292    let mut metrics = HashMap::default();
293    let mut corpus_entries = Vec::new();
294    let mut failed_corpus_replays = 0;
295    let mut optimization_best = None;
296
297    for InvariantWorkerOutput { plan, result, corpus_entries: worker_entries } in outputs {
298        if plan.worker_id == 0 {
299            failed_corpus_replays = result.failed_corpus_replays;
300        }
301        for (invariant, error) in result.errors {
302            errors.entry(invariant).or_insert(error);
303        }
304        merge_handler_errors(&mut handler_errors, result.handler_errors);
305        corpus_entries.extend(worker_entries);
306        runs += result.runs;
307        calls += result.calls;
308        reverts += result.reverts;
309        if !result.last_run_inputs.is_empty() {
310            last_run_inputs = result.last_run_inputs;
311        }
312        gas_report_traces.extend(result.gas_report_traces);
313        HitMaps::merge_opt(&mut line_coverage, result.line_coverage);
314        merge_metrics(&mut metrics, result.metrics);
315        merge_optimization(
316            &mut optimization_best,
317            result.optimization_best_value,
318            result.optimization_best_sequence,
319        );
320    }
321    let (optimization_best_value, optimization_best_sequence) =
322        optimization_best.map(|(value, sequence)| (Some(value), sequence)).unwrap_or_default();
323    Ok((
324        InvariantFuzzTestResult::new(
325            errors,
326            handler_errors,
327            runs,
328            calls,
329            reverts,
330            last_run_inputs,
331            gas_report_traces,
332            line_coverage,
333            metrics,
334            failed_corpus_replays,
335            workers,
336            optimization_best_value,
337            optimization_best_sequence,
338        ),
339        corpus_entries,
340    ))
341}
342
343fn ensure_outputs_cover_campaign(
344    spec: InvariantCampaignSpec,
345    outputs: &[InvariantWorkerOutput],
346) -> Result<()> {
347    ensure_worker_ids_are_valid(outputs)?;
348
349    if spec.total_runs == 0 {
350        ensure!(
351            outputs.len() == 1
352                && outputs[0].plan.first_global_run == 0
353                && outputs[0].plan.runs == 0,
354            "invariant worker outputs do not cover the logical campaign"
355        );
356        return Ok(());
357    }
358
359    let mut next_global_run = 0;
360    for output in outputs {
361        ensure!(output.plan.runs > 0, "invariant worker outputs do not cover the logical campaign");
362        ensure!(
363            output.plan.first_global_run == next_global_run,
364            "invariant worker outputs do not cover the logical campaign"
365        );
366        next_global_run = next_global_run
367            .checked_add(output.plan.runs)
368            .ok_or_else(|| eyre::eyre!("invariant worker output range overflows"))?;
369    }
370
371    ensure!(
372        next_global_run == spec.total_runs,
373        "invariant worker outputs do not cover the logical campaign"
374    );
375    Ok(())
376}
377
378fn ensure_worker_ids_are_valid(outputs: &[InvariantWorkerOutput]) -> Result<()> {
379    let mut seen = HashSet::with_capacity(outputs.len());
380    for output in outputs {
381        ensure!(
382            seen.insert(output.plan.worker_id),
383            "duplicate invariant worker output for worker {}",
384            output.plan.worker_id
385        );
386    }
387
388    ensure!(seen.contains(&0), "missing invariant master worker output");
389    Ok(())
390}
391
392/// Deduplicates handler assertion failures by site, keeping the shorter reproducer.
393/// Equal-length reproducers keep the one already inserted, which is the earlier logical worker
394/// because the caller folds worker outputs in `first_global_run` order.
395fn merge_handler_errors(
396    merged: &mut HashMap<(Address, Selector), InvariantFuzzError>,
397    worker_errors: HashMap<(Address, Selector), InvariantFuzzError>,
398) {
399    for (site, error) in worker_errors {
400        let candidate_len = handler_error_sequence_len(&error);
401        if merged
402            .get(&site)
403            .is_none_or(|existing| handler_error_sequence_len(existing) > candidate_len)
404        {
405            merged.insert(site, error);
406        }
407    }
408}
409
410/// Adds worker-local selector metrics into the logical campaign totals.
411fn merge_metrics(
412    merged: &mut HashMap<String, InvariantMetrics>,
413    worker_metrics: HashMap<String, InvariantMetrics>,
414) {
415    for (selector, metrics) in worker_metrics {
416        let entry = merged.entry(selector).or_default();
417        entry.calls += metrics.calls;
418        entry.reverts += metrics.reverts;
419        entry.discards += metrics.discards;
420    }
421}
422
423/// Keeps the best optimization value, using logical run order to break ties.
424fn merge_optimization(
425    best: &mut Option<(I256, Vec<BasicTxDetails>)>,
426    candidate_value: Option<I256>,
427    candidate_sequence: Vec<BasicTxDetails>,
428) {
429    let Some(candidate_value) = candidate_value else {
430        return;
431    };
432
433    if best.as_ref().is_none_or(|(best, _)| candidate_value > *best) {
434        *best = Some((candidate_value, candidate_sequence));
435    }
436}
437
438fn handler_error_sequence_len(error: &InvariantFuzzError) -> usize {
439    error.as_handler_assertion().map_or(usize::MAX, |failure| failure.call_sequence.len())
440}
441
442#[cfg(test)]
443mod tests {
444    use super::{
445        super::error::{FailedInvariantCaseData, HandlerAssertionFailure},
446        *,
447    };
448    use alloy_primitives::{B256, Bytes};
449    use foundry_evm_coverage::HitMap;
450    use foundry_evm_fuzz::CallDetails;
451    use proptest::test_runner::TestError;
452    use revm_inspectors::tracing::CallTraceArena;
453
454    fn empty_result(reverts: usize, failed_corpus_replays: usize) -> InvariantFuzzTestResult {
455        InvariantFuzzTestResult::new(
456            HashMap::default(),
457            HashMap::default(),
458            0,
459            0,
460            reverts,
461            Vec::new(),
462            Vec::new(),
463            None,
464            HashMap::default(),
465            failed_corpus_replays,
466            1,
467            None,
468            Vec::new(),
469        )
470    }
471
472    fn basic_tx(sender: u8) -> BasicTxDetails {
473        BasicTxDetails {
474            warp: None,
475            roll: None,
476            sender: Address::repeat_byte(sender),
477            call_details: CallDetails {
478                target: Address::repeat_byte(sender.wrapping_add(1)),
479                calldata: Bytes::from(vec![0, 0, 0, sender]),
480                value: None,
481            },
482        }
483    }
484
485    fn hit_maps(pc: u32, hits: u32) -> HitMaps {
486        let mut hit_map = HitMap::new(Bytes::from_static(&[0]));
487        hit_map.hits(pc, hits);
488
489        let mut maps = HitMaps::default();
490        maps.insert(B256::ZERO, hit_map);
491        maps
492    }
493
494    /// Builds a worker-local result fixture with the fields merged by the aggregator.
495    fn worker_result(
496        reverts: usize,
497        last_input_sender: u8,
498        metric_name: &str,
499        metrics: InvariantMetrics,
500        coverage_hits: u32,
501        failed_corpus_replays: usize,
502    ) -> InvariantFuzzTestResult {
503        let mut result = empty_result(reverts, failed_corpus_replays);
504        result.runs = 1;
505        result.calls = metrics.calls;
506        result.last_run_inputs = vec![basic_tx(last_input_sender)];
507        result.gas_report_traces.push(vec![CallTraceArena::default()]);
508        result.line_coverage = Some(hit_maps(7, coverage_hits));
509        result.metrics.insert(metric_name.to_string(), metrics);
510        result
511    }
512
513    fn sequence(len: usize, first_sender: u8) -> Vec<BasicTxDetails> {
514        (0..len).map(|idx| basic_tx(first_sender.wrapping_add(idx as u8))).collect()
515    }
516
517    /// Builds a predicate failure fixture with a reproducible call sequence.
518    fn predicate_error(reason: &str, sequence_len: usize) -> InvariantFuzzError {
519        InvariantFuzzError::BrokenInvariant(FailedInvariantCaseData {
520            test_error: TestError::Fail(reason.to_string().into(), sequence(sequence_len, 0x80)),
521            return_reason: reason.to_string().into(),
522            revert_reason: reason.to_string(),
523            addr: Address::repeat_byte(0x70),
524            calldata: Bytes::new(),
525            inner_sequence: Vec::new(),
526            shrink_run_limit: 0,
527            fail_on_revert: false,
528            assertion_failure: false,
529        })
530    }
531
532    /// Builds a handler assertion fixture with a reproducible call sequence.
533    fn handler_error(
534        reverter: Address,
535        selector: Selector,
536        sequence_len: usize,
537        reason: &str,
538    ) -> InvariantFuzzError {
539        InvariantFuzzError::HandlerAssertion(HandlerAssertionFailure {
540            reverter,
541            selector,
542            call_sequence: sequence(sequence_len, 0x90),
543            original_sequence_len: sequence_len,
544            revert_reason: reason.to_string(),
545            edge_fingerprint: B256::ZERO,
546        })
547    }
548
549    fn one_worker_plan(total_runs: u32) -> InvariantWorkerPlan {
550        let mut plans = InvariantCampaignSpec::new(total_runs).worker_plans(1).unwrap();
551        assert_eq!(plans.len(), 1);
552        plans.pop().unwrap()
553    }
554
555    #[test]
556    fn worker_plans_cover_logical_campaign_with_one_worker() {
557        let plan = one_worker_plan(3);
558
559        assert_eq!(plan.worker_id, 0);
560        assert_eq!(plan.first_global_run, 0);
561        assert_eq!(plan.runs, 3);
562    }
563
564    #[test]
565    fn worker_plans_split_runs_evenly() {
566        let plans = InvariantCampaignSpec::new(100).worker_plans(4).unwrap();
567
568        assert_eq!(
569            plans,
570            vec![
571                InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 25 },
572                InvariantWorkerPlan { worker_id: 1, first_global_run: 25, runs: 25 },
573                InvariantWorkerPlan { worker_id: 2, first_global_run: 50, runs: 25 },
574                InvariantWorkerPlan { worker_id: 3, first_global_run: 75, runs: 25 },
575            ]
576        );
577    }
578
579    #[test]
580    fn worker_plans_distribute_remainder_to_earlier_workers() {
581        let plans = InvariantCampaignSpec::new(10).worker_plans(3).unwrap();
582
583        assert_eq!(
584            plans,
585            vec![
586                InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 4 },
587                InvariantWorkerPlan { worker_id: 1, first_global_run: 4, runs: 3 },
588                InvariantWorkerPlan { worker_id: 2, first_global_run: 7, runs: 3 },
589            ]
590        );
591    }
592
593    #[test]
594    fn worker_plans_do_not_create_empty_workers_when_runs_are_available() {
595        let plans = InvariantCampaignSpec::new(2).worker_plans(8).unwrap();
596
597        assert_eq!(
598            plans,
599            vec![
600                InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 1 },
601                InvariantWorkerPlan { worker_id: 1, first_global_run: 1, runs: 1 },
602            ]
603        );
604    }
605
606    #[test]
607    fn worker_plans_keep_zero_run_campaign_as_single_empty_plan() {
608        let plans = InvariantCampaignSpec::new(0).worker_plans(4).unwrap();
609
610        assert_eq!(plans, vec![InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 0 }]);
611    }
612
613    #[test]
614    fn worker_plans_reject_zero_workers() {
615        let err = InvariantCampaignSpec::new(1).worker_plans(0).unwrap_err();
616        assert!(err.to_string().contains("requires at least one worker"));
617
618        let err = InvariantCampaignSpec::new(0).worker_plans(0).unwrap_err();
619        assert!(err.to_string().contains("requires at least one worker"));
620    }
621
622    #[test]
623    fn campaign_state_stops_after_terminal_request() {
624        let state = InvariantCampaignState::new(EarlyExit::new(false), None);
625        assert!(!state.should_stop());
626
627        state.request_terminal_stop();
628
629        assert!(state.should_stop());
630    }
631
632    #[test]
633    fn campaign_state_uses_shared_timeout_and_global_throughput() {
634        let state = InvariantCampaignState::new(EarlyExit::new(false), Some(0));
635        std::thread::sleep(Duration::from_millis(1));
636
637        assert!(state.is_timed_campaign());
638        assert!(state.should_stop());
639
640        state.record_call(20);
641        state.record_call(30);
642        assert_eq!(state.throughput_totals(), (2, 50));
643        assert_eq!(state.increment_runs(), 1);
644        assert_eq!(state.total_runs(), 1);
645    }
646
647    #[test]
648    fn campaign_state_deduplicates_handler_failure_events_across_workers() {
649        let state = InvariantCampaignState::new(EarlyExit::new(false), None);
650        let target = Address::repeat_byte(0x11);
651        let selector = Selector::from([0xde, 0xad, 0xbe, 0xef]);
652        let mut first_worker = InvariantFailures::new();
653        first_worker.seed_handler_failure(
654            target,
655            selector,
656            handler_error(target, selector, 2, "assertion failed"),
657        );
658        let mut second_worker = InvariantFailures::new();
659        second_worker.seed_handler_failure(
660            target,
661            selector,
662            handler_error(target, selector, 1, "assertion failed"),
663        );
664
665        state.sync_handler_failures(&first_worker);
666        state.sync_handler_failures(&second_worker);
667
668        assert_eq!(state.failure_metrics().broken_handlers, 1);
669    }
670
671    #[test]
672    fn aggregator_returns_single_worker_result_without_rewriting() {
673        let spec = InvariantCampaignSpec::new(1);
674        let worker = InvariantWorkerOutput::new(one_worker_plan(1), empty_result(2, 3));
675
676        let mut aggregator = InvariantCampaignAggregator::new(spec);
677        aggregator.push(worker);
678        let result = aggregator.finish().unwrap();
679
680        assert_eq!(result.reverts, 2);
681        assert_eq!(result.failed_corpus_replays, 3);
682    }
683
684    #[test]
685    fn aggregator_accepts_single_worker_output_for_zero_run_campaign() {
686        let spec = InvariantCampaignSpec::new(0);
687        let worker = InvariantWorkerOutput::new(
688            InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 0 },
689            empty_result(0, 0),
690        );
691
692        let mut aggregator = InvariantCampaignAggregator::new(spec);
693        aggregator.push(worker);
694        let result = aggregator.finish().unwrap();
695
696        assert_eq!(result.reverts, 0);
697    }
698
699    #[test]
700    fn aggregator_merges_multiple_worker_outputs_in_logical_run_order() {
701        let spec = InvariantCampaignSpec::new(3);
702        let plans = [
703            InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 1 },
704            InvariantWorkerPlan { worker_id: 1, first_global_run: 1, runs: 1 },
705            InvariantWorkerPlan { worker_id: 2, first_global_run: 2, runs: 1 },
706        ];
707
708        let mut aggregator = InvariantCampaignAggregator::new(spec);
709        aggregator.push(InvariantWorkerOutput::new(
710            plans[2],
711            worker_result(
712                3,
713                0x30,
714                "transfer(address)",
715                InvariantMetrics { calls: 3, reverts: 1, discards: 0 },
716                3,
717                0,
718            ),
719        ));
720        aggregator.push(InvariantWorkerOutput::new(
721            plans[0],
722            worker_result(
723                1,
724                0x10,
725                "transfer(address)",
726                InvariantMetrics { calls: 1, reverts: 0, discards: 2 },
727                1,
728                4,
729            ),
730        ));
731        aggregator.push(InvariantWorkerOutput::new(
732            plans[1],
733            worker_result(
734                2,
735                0x20,
736                "approve(address)",
737                InvariantMetrics { calls: 2, reverts: 1, discards: 1 },
738                2,
739                0,
740            ),
741        ));
742
743        let result = aggregator.finish().unwrap();
744
745        assert_eq!(result.runs, 3);
746        assert_eq!(result.calls, 6);
747        assert_eq!(result.reverts, 6);
748        assert_eq!(result.gas_report_traces.len(), 3);
749        assert_eq!(result.last_run_inputs[0].sender, Address::repeat_byte(0x30));
750
751        let transfer_metrics = result.metrics.get("transfer(address)").unwrap();
752        assert_eq!(transfer_metrics, &InvariantMetrics { calls: 4, reverts: 1, discards: 2 });
753        let approve_metrics = result.metrics.get("approve(address)").unwrap();
754        assert_eq!(approve_metrics, &InvariantMetrics { calls: 2, reverts: 1, discards: 1 });
755
756        let coverage = result.line_coverage.unwrap();
757        assert_eq!(coverage.get(&B256::ZERO).unwrap().get(7).unwrap().get(), 6);
758        assert_eq!(result.failed_corpus_replays, 4);
759    }
760
761    #[test]
762    fn aggregator_preserves_run_and_call_counts() {
763        let spec = InvariantCampaignSpec::new(3);
764        let plans = [
765            InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 1 },
766            InvariantWorkerPlan { worker_id: 1, first_global_run: 1, runs: 2 },
767        ];
768        let mut first = empty_result(0, 0);
769        first.runs = 1;
770        first.calls = 1000;
771        let mut second = empty_result(0, 0);
772        second.runs = 2;
773        second.calls = 2000;
774
775        let mut aggregator = InvariantCampaignAggregator::new(spec);
776        aggregator.push(InvariantWorkerOutput::new(plans[1], second));
777        aggregator.push(InvariantWorkerOutput::new(plans[0], first));
778        let result = aggregator.finish().unwrap();
779
780        assert_eq!(result.runs, 3);
781        assert_eq!(result.calls, 3000);
782    }
783
784    #[test]
785    fn timeout_aggregator_accepts_partial_outputs_with_range_gaps() {
786        fn result_with_counts(
787            runs: usize,
788            calls: usize,
789            has_last_run: bool,
790            failed_corpus_replays: usize,
791        ) -> InvariantFuzzTestResult {
792            let mut result = empty_result(0, failed_corpus_replays);
793            result.runs = runs;
794            result.calls = calls;
795            result.last_run_inputs = if has_last_run { vec![basic_tx(0x44)] } else { Vec::new() };
796            result
797        }
798
799        let spec = InvariantCampaignSpec::new(10);
800        let outputs = [
801            InvariantWorkerOutput::new(
802                InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 2 },
803                result_with_counts(2, 20, true, 5),
804            ),
805            InvariantWorkerOutput::new(
806                InvariantWorkerPlan { worker_id: 1, first_global_run: 4, runs: 0 },
807                result_with_counts(0, 0, false, 0),
808            ),
809            InvariantWorkerOutput::new(
810                InvariantWorkerPlan { worker_id: 2, first_global_run: 7, runs: 1 },
811                result_with_counts(1, 10, true, 0),
812            ),
813        ];
814
815        let mut strict = InvariantCampaignAggregator::new(spec);
816        for output in outputs {
817            strict.push(output);
818        }
819        let err = strict.finish().unwrap_err();
820        assert!(err.to_string().contains("do not cover the logical campaign"));
821
822        let mut partial = InvariantCampaignAggregator::new(spec);
823        partial.push(InvariantWorkerOutput::new(
824            InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 2 },
825            result_with_counts(2, 20, true, 5),
826        ));
827        partial.push(InvariantWorkerOutput::new(
828            InvariantWorkerPlan { worker_id: 1, first_global_run: 4, runs: 0 },
829            result_with_counts(0, 0, false, 0),
830        ));
831        partial.push(InvariantWorkerOutput::new(
832            InvariantWorkerPlan { worker_id: 2, first_global_run: 7, runs: 1 },
833            result_with_counts(1, 10, true, 0),
834        ));
835
836        let (result, corpus_entries) = partial.finish_partial_with_corpus_entries().unwrap();
837
838        assert_eq!(result.runs, 3);
839        assert_eq!(result.calls, 30);
840        assert_eq!(result.failed_corpus_replays, 5);
841        assert!(corpus_entries.is_empty());
842    }
843
844    #[test]
845    fn aggregator_keeps_earlier_predicate_failure_for_each_invariant() {
846        let spec = InvariantCampaignSpec::new(2);
847        let plans = [
848            InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 1 },
849            InvariantWorkerPlan { worker_id: 1, first_global_run: 1, runs: 1 },
850        ];
851        let mut earlier = empty_result(0, 0);
852        earlier.errors.insert("invariant_balance".to_string(), predicate_error("earlier", 3));
853        let mut later = empty_result(0, 0);
854        later.errors.insert("invariant_balance".to_string(), predicate_error("later", 1));
855
856        let mut aggregator = InvariantCampaignAggregator::new(spec);
857        aggregator.push(InvariantWorkerOutput::new(plans[1], later));
858        aggregator.push(InvariantWorkerOutput::new(plans[0], earlier));
859        let result = aggregator.finish().unwrap();
860
861        assert_eq!(result.errors.len(), 1);
862        assert_eq!(result.errors["invariant_balance"].revert_reason().as_deref(), Some("earlier"));
863    }
864
865    #[test]
866    fn aggregator_dedupes_handler_assertions_by_site_and_keeps_shorter_sequence() {
867        let spec = InvariantCampaignSpec::new(2);
868        let plans = [
869            InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 1 },
870            InvariantWorkerPlan { worker_id: 1, first_global_run: 1, runs: 1 },
871        ];
872        let site = (Address::repeat_byte(0xaa), Selector::from([1, 2, 3, 4]));
873        let mut longer = empty_result(0, 0);
874        longer.handler_errors.insert(site, handler_error(site.0, site.1, 4, "longer"));
875        let mut shorter = empty_result(0, 0);
876        shorter.handler_errors.insert(site, handler_error(site.0, site.1, 2, "shorter"));
877
878        let mut aggregator = InvariantCampaignAggregator::new(spec);
879        aggregator.push(InvariantWorkerOutput::new(plans[1], shorter));
880        aggregator.push(InvariantWorkerOutput::new(plans[0], longer));
881        let result = aggregator.finish().unwrap();
882
883        let failure = result.handler_errors[&site].as_handler_assertion().unwrap();
884        assert_eq!(result.handler_errors.len(), 1);
885        assert_eq!(failure.call_sequence.len(), 2);
886        assert_eq!(failure.revert_reason, "shorter");
887    }
888
889    #[test]
890    fn aggregator_keeps_earlier_handler_assertion_when_lengths_tie() {
891        let spec = InvariantCampaignSpec::new(2);
892        let plans = [
893            InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 1 },
894            InvariantWorkerPlan { worker_id: 1, first_global_run: 1, runs: 1 },
895        ];
896        let site = (Address::repeat_byte(0xaa), Selector::from([1, 2, 3, 4]));
897        let mut earlier = empty_result(0, 0);
898        earlier.handler_errors.insert(site, handler_error(site.0, site.1, 2, "earlier"));
899        let mut later = empty_result(0, 0);
900        later.handler_errors.insert(site, handler_error(site.0, site.1, 2, "later"));
901
902        let mut aggregator = InvariantCampaignAggregator::new(spec);
903        aggregator.push(InvariantWorkerOutput::new(plans[1], later));
904        aggregator.push(InvariantWorkerOutput::new(plans[0], earlier));
905        let result = aggregator.finish().unwrap();
906
907        let failure = result.handler_errors[&site].as_handler_assertion().unwrap();
908        assert_eq!(result.handler_errors.len(), 1);
909        assert_eq!(failure.call_sequence.len(), 2);
910        assert_eq!(failure.revert_reason, "earlier");
911    }
912
913    #[test]
914    fn aggregator_keeps_distinct_predicate_failures() {
915        let spec = InvariantCampaignSpec::new(2);
916        let plans = [
917            InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 1 },
918            InvariantWorkerPlan { worker_id: 1, first_global_run: 1, runs: 1 },
919        ];
920        let mut earlier = empty_result(0, 0);
921        earlier.errors.insert("invariant_a".to_string(), predicate_error("a", 3));
922        let mut later = empty_result(0, 0);
923        later.errors.insert("invariant_b".to_string(), predicate_error("b", 2));
924
925        let mut aggregator = InvariantCampaignAggregator::new(spec);
926        aggregator.push(InvariantWorkerOutput::new(plans[1], later));
927        aggregator.push(InvariantWorkerOutput::new(plans[0], earlier));
928        let result = aggregator.finish().unwrap();
929
930        assert_eq!(result.errors.len(), 2);
931        assert_eq!(result.errors["invariant_a"].revert_reason().as_deref(), Some("a"));
932        assert_eq!(result.errors["invariant_b"].revert_reason().as_deref(), Some("b"));
933    }
934
935    #[test]
936    fn aggregator_keeps_first_max_optimization_value_on_tie() {
937        let spec = InvariantCampaignSpec::new(3);
938        let plans = [
939            InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 1 },
940            InvariantWorkerPlan { worker_id: 1, first_global_run: 1, runs: 1 },
941            InvariantWorkerPlan { worker_id: 2, first_global_run: 2, runs: 1 },
942        ];
943        let mut first = empty_result(0, 0);
944        first.optimization_best_value = Some(I256::try_from(7).unwrap());
945        first.optimization_best_sequence = sequence(1, 0x10);
946        let mut earlier_best = empty_result(0, 0);
947        earlier_best.optimization_best_value = Some(I256::try_from(9).unwrap());
948        earlier_best.optimization_best_sequence = sequence(1, 0x20);
949        let mut later_tie = empty_result(0, 0);
950        later_tie.optimization_best_value = Some(I256::try_from(9).unwrap());
951        later_tie.optimization_best_sequence = sequence(1, 0x30);
952
953        let mut aggregator = InvariantCampaignAggregator::new(spec);
954        aggregator.push(InvariantWorkerOutput::new(plans[2], later_tie));
955        aggregator.push(InvariantWorkerOutput::new(plans[0], first));
956        aggregator.push(InvariantWorkerOutput::new(plans[1], earlier_best));
957        let result = aggregator.finish().unwrap();
958
959        assert_eq!(result.optimization_best_value, Some(I256::try_from(9).unwrap()));
960        assert_eq!(result.optimization_best_sequence[0].sender, Address::repeat_byte(0x20));
961    }
962
963    #[test]
964    fn aggregator_rejects_overlapping_outputs() {
965        let spec = InvariantCampaignSpec::new(1);
966        let mut aggregator = InvariantCampaignAggregator::new(spec);
967
968        aggregator.push(InvariantWorkerOutput::new(
969            InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 1 },
970            empty_result(0, 0),
971        ));
972        aggregator.push(InvariantWorkerOutput::new(
973            InvariantWorkerPlan { worker_id: 1, first_global_run: 0, runs: 1 },
974            empty_result(0, 0),
975        ));
976        let err = aggregator.finish().unwrap_err();
977
978        assert!(err.to_string().contains("do not cover the logical campaign"));
979    }
980
981    #[test]
982    fn aggregator_rejects_duplicate_worker_ids() {
983        let spec = InvariantCampaignSpec::new(2);
984        let mut aggregator = InvariantCampaignAggregator::new(spec);
985
986        aggregator.push(InvariantWorkerOutput::new(
987            InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 1 },
988            empty_result(0, 0),
989        ));
990        aggregator.push(InvariantWorkerOutput::new(
991            InvariantWorkerPlan { worker_id: 0, first_global_run: 1, runs: 1 },
992            empty_result(0, 0),
993        ));
994        let err = aggregator.finish().unwrap_err();
995
996        assert!(err.to_string().contains("duplicate invariant worker output"));
997    }
998
999    #[test]
1000    fn aggregator_allows_non_dense_worker_ids_with_contiguous_ranges() {
1001        let spec = InvariantCampaignSpec::new(2);
1002        let mut aggregator = InvariantCampaignAggregator::new(spec);
1003
1004        aggregator.push(InvariantWorkerOutput::new(
1005            InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 1 },
1006            empty_result(0, 0),
1007        ));
1008        aggregator.push(InvariantWorkerOutput::new(
1009            InvariantWorkerPlan { worker_id: 2, first_global_run: 1, runs: 1 },
1010            empty_result(2, 0),
1011        ));
1012        let result = aggregator.finish().unwrap();
1013
1014        assert_eq!(result.reverts, 2);
1015    }
1016
1017    #[test]
1018    fn aggregator_rejects_missing_master_worker() {
1019        let spec = InvariantCampaignSpec::new(2);
1020        let mut aggregator = InvariantCampaignAggregator::new(spec);
1021
1022        aggregator.push(InvariantWorkerOutput::new(
1023            InvariantWorkerPlan { worker_id: 1, first_global_run: 0, runs: 1 },
1024            empty_result(0, 0),
1025        ));
1026        aggregator.push(InvariantWorkerOutput::new(
1027            InvariantWorkerPlan { worker_id: 2, first_global_run: 1, runs: 1 },
1028            empty_result(0, 0),
1029        ));
1030        let err = aggregator.finish().unwrap_err();
1031
1032        assert!(err.to_string().contains("missing invariant master worker output"));
1033    }
1034
1035    #[test]
1036    fn aggregator_uses_master_failed_corpus_replays() {
1037        let spec = InvariantCampaignSpec::new(2);
1038        let plans = [
1039            InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 1 },
1040            InvariantWorkerPlan { worker_id: 1, first_global_run: 1, runs: 1 },
1041        ];
1042
1043        let mut aggregator = InvariantCampaignAggregator::new(spec);
1044        aggregator.push(InvariantWorkerOutput::new(plans[0], empty_result(0, 7)));
1045        aggregator.push(InvariantWorkerOutput::new(plans[1], empty_result(0, 1)));
1046        let result = aggregator.finish().unwrap();
1047
1048        assert_eq!(result.failed_corpus_replays, 7);
1049    }
1050
1051    #[test]
1052    fn aggregator_uses_master_failed_corpus_replays_independent_of_output_order() {
1053        let spec = InvariantCampaignSpec::new(2);
1054        let plans = [
1055            InvariantWorkerPlan { worker_id: 1, first_global_run: 0, runs: 1 },
1056            InvariantWorkerPlan { worker_id: 0, first_global_run: 1, runs: 1 },
1057        ];
1058
1059        let mut aggregator = InvariantCampaignAggregator::new(spec);
1060        aggregator.push(InvariantWorkerOutput::new(plans[0], empty_result(0, 0)));
1061        aggregator.push(InvariantWorkerOutput::new(plans[1], empty_result(0, 7)));
1062        let result = aggregator.finish().unwrap();
1063
1064        assert_eq!(result.failed_corpus_replays, 7);
1065    }
1066
1067    #[test]
1068    fn aggregator_rejects_plan_that_does_not_cover_campaign() {
1069        let spec = InvariantCampaignSpec::new(2);
1070        let plan = InvariantWorkerPlan { worker_id: 0, first_global_run: 0, runs: 1 };
1071        let worker = InvariantWorkerOutput::new(plan, empty_result(0, 0));
1072
1073        let mut aggregator = InvariantCampaignAggregator::new(spec);
1074        aggregator.push(worker);
1075        let err = aggregator.finish().unwrap_err();
1076
1077        assert!(err.to_string().contains("do not cover the logical campaign"));
1078    }
1079
1080    #[test]
1081    fn aggregator_rejects_missing_output() {
1082        let aggregator = InvariantCampaignAggregator::new(InvariantCampaignSpec::new(1));
1083        let err = aggregator.finish().unwrap_err();
1084
1085        assert!(err.to_string().contains("missing invariant worker output"));
1086    }
1087}