1use 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#[derive(Debug, Clone, Copy)]
20pub struct EvalRunOptions {
21 pub concurrency: usize,
23 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#[async_trait]
40pub trait EvalExecutor: Send + Sync {
41 async fn execute_task(&self, task: &EvalTask) -> Result<EvalRunResult>;
43
44 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
53pub async fn run_suite(executor: &dyn EvalExecutor, suite: &EvalSuite) -> Result<EvalReport> {
60 run_suite_with_options(executor, suite, EvalRunOptions::default()).await
61}
62
63pub 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 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(®_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 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 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 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}