Skip to main content

vtcode_eval/
executor.rs

1//! Executor trait boundary and pure eval orchestration.
2//!
3//! The [`EvalExecutor`] trait isolates the harness orchestration from the
4//! concrete agent runner. This is the "interface guard rail": `run_suite`
5//! depends only on the trait, so it can be unit-tested with a fake executor
6//! and the production executor implementation can be swapped
7//! without touching the orchestration logic.
8
9use anyhow::Result;
10use async_trait::async_trait;
11use futures::{StreamExt, stream};
12
13use crate::task::EvalTask;
14use crate::{
15    EvalReport, EvalRunResult, EvalSuite, SuiteReport, aggregate_metrics, build_task_report, compute_metric_with_k,
16};
17
18/// Controls bounded eval scheduling and the sampling value used by metrics.
19#[derive(Debug, Clone, Copy)]
20pub struct EvalRunOptions {
21    /// Maximum number of attempts executing concurrently.
22    pub concurrency: usize,
23    /// Sampling value for pass@k/pass^k. Defaults to the suite attempt count.
24    pub metric_k: Option<u32>,
25}
26
27impl Default for EvalRunOptions {
28    fn default() -> Self {
29        Self { concurrency: 2, metric_k: None }
30    }
31}
32
33/// Executes a single eval task attempt and returns the verified outcome.
34///
35/// Implementations own the full "execute this task" semantics — running the
36/// agent and applying any environment verification (probes). Keeping this
37/// behind a trait decouples the orchestration in [`run_suite`] from the
38/// concrete runner, which makes the harness independently testable.
39#[async_trait]
40pub trait EvalExecutor: Send + Sync {
41    /// Run one attempt of `task` and return the outcome.
42    async fn execute_task(&self, task: &EvalTask) -> Result<EvalRunResult>;
43
44    /// Run a numbered attempt. Existing executors remain source-compatible;
45    /// production executors override this to isolate each attempt.
46    async fn execute_task_attempt(&self, task: &EvalTask, attempt: u32) -> Result<EvalRunResult> {
47        let mut result = self.execute_task(task).await?;
48        result.attempt = attempt;
49        Ok(result)
50    }
51}
52
53/// Pure orchestration core: loop tasks × attempts through the executor,
54/// compute per-task metrics, and assemble the report.
55///
56/// This function performs no file I/O, no config reads, and no trust checks —
57/// those belong to the caller. It depends only on [`EvalExecutor`], which makes
58/// it fully unit-testable with an in-memory fake (see `executor::tests`).
59pub async fn run_suite(executor: &dyn EvalExecutor, suite: &EvalSuite) -> Result<EvalReport> {
60    run_suite_with_options(executor, suite, EvalRunOptions::default()).await
61}
62
63/// Execute a suite with bounded concurrency and deterministic report order.
64pub async fn run_suite_with_options(
65    executor: &dyn EvalExecutor,
66    suite: &EvalSuite,
67    options: EvalRunOptions,
68) -> Result<EvalReport> {
69    suite.validate()?;
70    let metric_k = options.metric_k.unwrap_or(suite.attempts);
71    anyhow::ensure!(metric_k >= 1 && metric_k <= suite.attempts, "evaluation metric k must be between 1 and attempts");
72    let concurrency = options.concurrency.max(1);
73
74    let jobs =
75        suite.tasks.iter().enumerate().flat_map(|(task_index, task)| {
76            (1..=suite.attempts).map(move |attempt| (task_index, task.clone(), attempt))
77        });
78    // Drain every in-flight attempt before propagating an error. A production
79    // executor may own an isolated worktree or another external resource that
80    // is cleaned up after its future completes; `try_collect` would drop the
81    // remaining futures on the first error and strand those resources.
82    let mut results = stream::iter(jobs)
83        .map(|(task_index, task, attempt)| async move {
84            let result = executor.execute_task_attempt(&task, attempt).await?;
85            Ok::<_, anyhow::Error>((task_index, attempt, result))
86        })
87        .buffer_unordered(concurrency)
88        .collect::<Vec<_>>()
89        .await
90        .into_iter()
91        .collect::<Result<Vec<_>>>()?;
92    for (task_index, expected_attempt, result) in &results {
93        let expected_task = suite
94            .tasks
95            .get(*task_index)
96            .ok_or_else(|| anyhow::anyhow!("executor returned an unknown task index {task_index}"))?;
97        anyhow::ensure!(
98            result.task_id == expected_task.id,
99            "executor returned task id '{}' for task '{}' attempt {}",
100            result.task_id,
101            expected_task.id,
102            expected_attempt
103        );
104        anyhow::ensure!(
105            result.attempt == *expected_attempt,
106            "executor returned attempt {} for task '{}' attempt {}",
107            result.attempt,
108            expected_task.id,
109            expected_attempt
110        );
111    }
112    results.sort_by_key(|(task_index, attempt, _)| (*task_index, *attempt));
113
114    let mut all_task_reports = Vec::new();
115    let mut all_cost_usd = 0.0;
116    let mut known_cost_runs = 0u32;
117    let mut unpriced_runs = 0u32;
118    let mut duration_secs = 0.0;
119    let mut trace_summary: Option<crate::trace_analyzer::HarnessTraceSummary> = None;
120
121    for (task_index, task) in suite.tasks.iter().enumerate() {
122        let run_results: Vec<EvalRunResult> = results
123            .iter()
124            .filter(|(index, _, _)| *index == task_index)
125            .map(|(_, _, result)| result.clone())
126            .collect();
127        for result in &run_results {
128            duration_secs += result.duration_secs;
129            if let Some(cost) = crate::report::priced_cost(result.cost_usd) {
130                all_cost_usd += cost;
131                known_cost_runs = known_cost_runs.saturating_add(1);
132            } else {
133                unpriced_runs = unpriced_runs.saturating_add(1);
134            }
135            if let Some(summary) = &result.trace_summary {
136                trace_summary.get_or_insert_with(Default::default).merge(summary);
137            }
138        }
139        let metric = compute_metric_with_k(&task.id, &run_results, metric_k)?;
140        all_task_reports.push(build_task_report(&task.id, &task.name, task.category, metric));
141    }
142
143    let all_metrics: Vec<_> = all_task_reports.iter().map(|r| r.metric.clone()).collect();
144    let cap_metrics: Vec<_> = all_task_reports
145        .iter()
146        .filter(|r| r.category == "Capability")
147        .map(|r| r.metric.clone())
148        .collect();
149    let reg_metrics: Vec<_> = all_task_reports
150        .iter()
151        .filter(|r| r.category == "Regression")
152        .map(|r| r.metric.clone())
153        .collect();
154
155    let efficiency = crate::report::CostEfficiency::from_runs(results.iter().map(|(_, _, result)| result));
156
157    Ok(EvalReport {
158        generated_at: chrono::Utc::now().to_rfc3339(),
159        suites: vec![SuiteReport {
160            suite_id: suite.id.clone(),
161            suite_name: suite.name.clone(),
162            task_reports: all_task_reports,
163            aggregate: aggregate_metrics(&all_metrics),
164            capability_metrics: aggregate_metrics(&cap_metrics),
165            regression_metrics: aggregate_metrics(&reg_metrics),
166            cost_usd: (known_cost_runs > 0).then_some(all_cost_usd),
167            unpriced_runs,
168            duration_secs,
169            trace_summary,
170            mean_cost_per_attempt: efficiency.mean_cost_per_attempt,
171            cost_per_solve: efficiency.cost_per_solve,
172            mean_tokens_per_attempt: efficiency.mean_tokens_per_attempt,
173            mean_turns_per_attempt: efficiency.mean_turns_per_attempt,
174        }],
175    })
176}
177
178#[cfg(test)]
179mod tests {
180    use super::*;
181    use crate::task::{EvalCategory, EvalRunResult, EvalTask, RunOutcome};
182    use std::sync::{
183        Arc,
184        atomic::{AtomicUsize, Ordering},
185    };
186    use tokio::time::{Duration, sleep};
187
188    /// In-memory fake executor for isolating `run_suite`.
189    struct FakeExecutor {
190        outcomes: Vec<RunOutcome>,
191        calls: AtomicUsize,
192    }
193
194    #[async_trait]
195    impl EvalExecutor for FakeExecutor {
196        async fn execute_task(&self, _task: &EvalTask) -> Result<EvalRunResult> {
197            let i = self.calls.fetch_add(1, Ordering::SeqCst);
198            let outcome = self.outcomes[i % self.outcomes.len()];
199            Ok(EvalRunResult {
200                task_id: _task.id.clone(),
201                outcome,
202                error_message: None,
203                duration_secs: 0.0,
204                attempt: (i + 1) as u32,
205                cost_usd: None,
206                transcript_path: None,
207                trace_summary: None,
208            })
209        }
210    }
211
212    fn suite(attempts: u32) -> EvalSuite {
213        EvalSuite {
214            id: "s1".into(),
215            name: "demo".into(),
216            tasks: vec![
217                EvalTask {
218                    id: "t1".into(),
219                    name: "t1".into(),
220                    category: EvalCategory::Capability,
221                    prompt: "p".into(),
222                    verify_commands: vec![],
223                    timeout_secs: None,
224                },
225                EvalTask {
226                    id: "t2".into(),
227                    name: "t2".into(),
228                    category: EvalCategory::Regression,
229                    prompt: "p".into(),
230                    verify_commands: vec![],
231                    timeout_secs: None,
232                },
233            ],
234            attempts,
235        }
236    }
237
238    #[tokio::test]
239    async fn run_suite_aggregates_capability_and_regression() {
240        // t1 gets [Pass, Fail]; t2 gets [Pass, Pass]
241        let exec = FakeExecutor {
242            outcomes: vec![RunOutcome::Pass, RunOutcome::Fail, RunOutcome::Pass, RunOutcome::Pass],
243            calls: AtomicUsize::new(0),
244        };
245        let report = run_suite(&exec, &suite(2)).await.unwrap();
246        let s = &report.suites[0];
247        // With k=attempts, t1's two-sample pass@2 is 1.0 and t2 is 1.0.
248        assert_eq!(s.aggregate.passed_runs, 3);
249        assert_eq!(s.aggregate.total_runs, 4);
250        assert!((s.capability_metrics.pass_at_k - 1.0).abs() < 1e-9);
251        assert!((s.regression_metrics.pass_at_k - 1.0).abs() < 1e-9);
252        assert!(s.capability_metrics.pass_all_k < 1.0);
253        assert_eq!(s.regression_metrics.pass_all_k, 1.0);
254    }
255
256    struct BoundedExecutor {
257        active: AtomicUsize,
258        max_active: AtomicUsize,
259    }
260
261    #[async_trait]
262    impl EvalExecutor for BoundedExecutor {
263        async fn execute_task(&self, task: &EvalTask) -> Result<EvalRunResult> {
264            self.execute_task_attempt(task, 1).await
265        }
266
267        async fn execute_task_attempt(&self, task: &EvalTask, attempt: u32) -> Result<EvalRunResult> {
268            let active = self.active.fetch_add(1, Ordering::SeqCst) + 1;
269            self.max_active.fetch_max(active, Ordering::SeqCst);
270            sleep(Duration::from_millis(5)).await;
271            self.active.fetch_sub(1, Ordering::SeqCst);
272            Ok(EvalRunResult {
273                task_id: task.id.clone(),
274                outcome: RunOutcome::Pass,
275                error_message: None,
276                duration_secs: 0.005,
277                attempt,
278                cost_usd: Some(0.0),
279                transcript_path: None,
280                trace_summary: None,
281            })
282        }
283    }
284
285    #[tokio::test]
286    async fn run_suite_bounds_attempts_and_keeps_report_order() {
287        let executor = Arc::new(BoundedExecutor {
288            active: AtomicUsize::new(0),
289            max_active: AtomicUsize::new(0),
290        });
291        let report =
292            run_suite_with_options(executor.as_ref(), &suite(4), EvalRunOptions { concurrency: 2, metric_k: Some(2) })
293                .await
294                .expect("bounded suite should complete");
295
296        assert!(executor.max_active.load(Ordering::SeqCst) <= 2);
297        let task_reports = &report.suites[0].task_reports;
298        assert_eq!(
299            task_reports
300                .iter()
301                .map(|report| report.metric.task_id.as_str())
302                .collect::<Vec<_>>(),
303            ["t1", "t2"]
304        );
305        assert_eq!(report.suites[0].unpriced_runs, 0);
306        assert_eq!(report.suites[0].cost_usd, Some(0.0));
307    }
308
309    struct IdentityViolatingExecutor;
310
311    #[async_trait]
312    impl EvalExecutor for IdentityViolatingExecutor {
313        async fn execute_task(&self, task: &EvalTask) -> Result<EvalRunResult> {
314            Ok(EvalRunResult {
315                task_id: task.id.clone(),
316                outcome: RunOutcome::Pass,
317                error_message: None,
318                duration_secs: 0.0,
319                attempt: 1,
320                cost_usd: None,
321                transcript_path: None,
322                trace_summary: None,
323            })
324        }
325
326        async fn execute_task_attempt(&self, task: &EvalTask, _attempt: u32) -> Result<EvalRunResult> {
327            let mut result = self.execute_task(task).await?;
328            result.task_id = "wrong-task".into();
329            result.attempt = 99;
330            Ok(result)
331        }
332    }
333
334    #[tokio::test]
335    async fn run_suite_rejects_executor_identity_mismatch() {
336        let error = run_suite(&IdentityViolatingExecutor, &suite(1))
337            .await
338            .expect_err("executor identity mismatch must fail the suite");
339        assert!(error.to_string().contains("executor returned task id"));
340    }
341}