1use std::collections::{HashMap, HashSet};
2use std::path::{Path, PathBuf};
3use std::sync::{Arc, Mutex, OnceLock, RwLock, RwLockReadGuard, RwLockWriteGuard};
4use std::time::Instant;
5
6use ai_agents_hooks::AgentHooks;
7use ai_agents_observability::config::{ExportConfig, ExportFormat};
8use ai_agents_observability::{CostEstimator, ObservabilityConfig};
9use ai_agents_runtime::spec::{AgentSpec, LLMConfigOrSelector, StorageConfig};
10use ai_agents_runtime::{Agent, AgentBuilder, RuntimeAgent, StreamChunk};
11use async_trait::async_trait;
12use futures::{StreamExt, stream};
13use serde_json::{Value, json};
14use tokio::time::{Duration, Instant as TokioInstant, timeout, timeout_at};
15
16use crate::assertion::{
17 Assertion, AssertionEvalContext, AssertionOutcome, AssertionResultDetail, evaluate_assertion,
18};
19use crate::budget::{BudgetProviderConfig, ScenarioBudgetTracker};
20use crate::compatibility::suite_from_jsonl;
21use crate::evidence::{
22 ApprovalEvidence, LlmRequestEvidence, collect_turn_evidence, relationship_snapshot,
23};
24use crate::fixtures::{
25 AttemptFixtureContext, AttemptWorkspace, LlmFixtureMode, RecordingToolLog,
26 WorkspacePolicyFixtureConfig, build_approval_handler, build_llm_registry, build_tool_registry,
27 resolve_fixture_context, start_mock_server,
28};
29use crate::judge::{JudgeConfig, JudgeResolver};
30use crate::metrics::compute_metrics;
31use crate::redaction::{redact_text, redact_value};
32use crate::suite::{
33 AttemptResult, EvalResult, EvalSuite, FailureCategory, IsolationMode, ResetStepConfig,
34 Scenario, ScenarioResult, ScenarioStatus, ScenarioStep, Turn, TurnResult, turn_expected_error,
35 turn_runtime_context,
36};
37use crate::{EvalError, Result};
38
39struct EvalRecordHooks {
41 tool_log: RecordingToolLog,
42 approval_log: RecordingApprovalLog,
43 llm_log: RecordingLlmLog,
44}
45
46#[derive(Clone, Default)]
47struct RecordingLlmLog {
48 records: Arc<Mutex<Vec<LlmRequestEvidence>>>,
49}
50
51impl RecordingLlmLog {
52 fn len(&self) -> usize {
53 self.records
54 .lock()
55 .unwrap_or_else(|poisoned| poisoned.into_inner())
56 .len()
57 }
58
59 fn push_messages(&self, messages: &[ai_agents_core::ChatMessage]) {
60 self.records
61 .lock()
62 .unwrap_or_else(|poisoned| poisoned.into_inner())
63 .push(LlmRequestEvidence::from_messages(messages));
64 }
65
66 fn records_since(&self, start: usize) -> Vec<LlmRequestEvidence> {
67 self.records
68 .lock()
69 .unwrap_or_else(|poisoned| poisoned.into_inner())
70 .get(start..)
71 .unwrap_or_default()
72 .to_vec()
73 }
74}
75
76#[derive(Clone, Default)]
77struct RecordingApprovalLog {
78 records: Arc<Mutex<Vec<ApprovalEvidence>>>,
79}
80
81impl RecordingApprovalLog {
82 fn len(&self) -> usize {
83 self.records
84 .lock()
85 .unwrap_or_else(|poisoned| poisoned.into_inner())
86 .len()
87 }
88
89 fn push(&self, evidence: ApprovalEvidence) {
90 self.records
91 .lock()
92 .unwrap_or_else(|poisoned| poisoned.into_inner())
93 .push(evidence);
94 }
95
96 fn records_since(&self, start: usize) -> Vec<ApprovalEvidence> {
97 self.records
98 .lock()
99 .unwrap_or_else(|poisoned| poisoned.into_inner())
100 .get(start..)
101 .unwrap_or_default()
102 .to_vec()
103 }
104}
105
106#[async_trait]
107impl AgentHooks for EvalRecordHooks {
108 async fn on_llm_start(&self, messages: &[ai_agents_core::ChatMessage]) {
109 self.llm_log.push_messages(messages);
110 }
111
112 async fn on_tool_execution_record(&self, record: &ai_agents_core::ToolExecutionRecord) {
113 self.tool_log.push_executor_record(record);
114 }
115
116 async fn on_approval_resolved(
117 &self,
118 request: &ai_agents_hitl::ApprovalRequest,
119 raw_result: &ai_agents_hitl::ApprovalResult,
120 outcome: &ai_agents_hitl::ApprovalResolvedOutcome,
121 ) {
122 self.approval_log.push(ApprovalEvidence::from_resolution(
123 request, raw_result, outcome,
124 ));
125 }
126}
127
128#[derive(Debug, Clone, Default)]
130pub struct EvalRunnerOptions {
131 pub agent: Option<PathBuf>,
133 pub scenarios: Option<PathBuf>,
135 pub output: PathBuf,
137 pub ids: Vec<String>,
139 pub tags: Vec<String>,
141 pub tag_mode_all: bool,
143 pub languages: Vec<String>,
145 pub retries: Option<u32>,
147 pub timeout_ms: Option<u64>,
149 pub parallel: Option<usize>,
151 pub fail_fast: bool,
153 pub observability: bool,
155 pub llm_mode: Option<LlmFixtureMode>,
157 pub cassette: Option<PathBuf>,
159}
160
161pub struct EvalRunner {
163 suite_path: PathBuf,
165 suite: EvalSuite,
167 options: EvalRunnerOptions,
169}
170
171struct BuildAgentParams<'a> {
173 agent_path: &'a Path,
174 base_dir: &'a Path,
175 attempt_context: &'a AttemptFixtureContext,
176 tool_log: RecordingToolLog,
177 approval_log: RecordingApprovalLog,
178 llm_log: RecordingLlmLog,
179 approval_handler: Option<Arc<dyn ai_agents_hitl::ApprovalHandler>>,
180 budget: Option<ScenarioBudgetTracker>,
181}
182
183struct RunTurnParams<'a> {
185 agent: &'a RuntimeAgent,
186 scenario: &'a Scenario,
187 turn: &'a Turn,
188 index: usize,
189 tool_log: &'a RecordingToolLog,
190 approval_log: &'a RecordingApprovalLog,
191 llm_log: &'a RecordingLlmLog,
192}
193
194fn load_eval_suite(path: &Path) -> Result<EvalSuite> {
195 let content = std::fs::read_to_string(path)?;
196 if path.extension().and_then(|extension| extension.to_str()) == Some("jsonl") {
197 suite_from_jsonl(
198 path.file_stem()
199 .and_then(|stem| stem.to_str())
200 .unwrap_or("eval")
201 .to_string(),
202 &content,
203 )
204 } else {
205 parse_eval_suite_yaml(&content)
206 }
207}
208
209fn parse_eval_suite_yaml(content: &str) -> Result<EvalSuite> {
210 let mut unknown_fields = Vec::new();
211 let deserializer = serde_yaml::Deserializer::from_str(content);
212 let suite = serde_ignored::deserialize(deserializer, |path| {
213 unknown_fields.push(path.to_string());
214 })?;
215 if !unknown_fields.is_empty() {
216 unknown_fields.sort();
217 unknown_fields.dedup();
218 return Err(EvalError::Config(format!(
219 "unknown eval configuration field(s): {}",
220 unknown_fields.join(", ")
221 )));
222 }
223 Ok(suite)
224}
225
226fn authorize_llm_mode(
227 suite_mode: LlmFixtureMode,
228 override_mode: Option<LlmFixtureMode>,
229) -> Result<()> {
230 if override_mode.is_some()
231 || matches!(suite_mode, LlmFixtureMode::Mock | LlmFixtureMode::Replay)
232 {
233 return Ok(());
234 }
235 match suite_mode {
236 LlmFixtureMode::Real => Err(EvalError::Config(
237 "suite-declared fixtures.llm.mode real requires --real-llm or EvalRunnerOptions.llm_mode = Some(LlmFixtureMode::Real)"
238 .to_string(),
239 )),
240 LlmFixtureMode::Record => Err(EvalError::Config(
241 "suite-declared fixtures.llm.mode record requires --record or EvalRunnerOptions.llm_mode = Some(LlmFixtureMode::Record)"
242 .to_string(),
243 )),
244 LlmFixtureMode::Mock | LlmFixtureMode::Replay => Ok(()),
245 }
246}
247
248impl EvalRunner {
249 pub fn from_file(path: impl AsRef<Path>, options: EvalRunnerOptions) -> Result<Self> {
250 let path = path.as_ref().to_path_buf();
251 let mut suite = load_eval_suite(&path)?;
252 if let Some(agent) = &options.agent {
253 suite.agent = Some(agent.clone());
254 }
255 if let Some(retries) = options.retries {
256 suite.settings.retries = retries;
257 }
258 if let Some(timeout_ms) = options.timeout_ms {
259 suite.settings.timeout_per_turn_ms = timeout_ms;
260 }
261 if let Some(parallel) = options.parallel {
262 suite.settings.parallel = parallel > 1;
263 suite.settings.max_concurrent = parallel.max(1);
264 }
265 if options.fail_fast {
266 suite.settings.fail_fast = true;
267 }
268 authorize_llm_mode(suite.fixtures.llm.mode, options.llm_mode)?;
269 if let Some(mode) = options.llm_mode {
270 suite.fixtures.llm.mode = mode;
271 }
272 if let Some(cassette) = &options.cassette {
273 suite.fixtures.llm.cassette = Some(cassette.clone());
274 }
275 suite.validate(options.agent.as_ref())?;
276 Ok(Self {
277 suite_path: path,
278 suite,
279 options,
280 })
281 }
282
283 pub fn validate_file(path: impl AsRef<Path>, agent_override: Option<PathBuf>) -> Result<()> {
284 let path = path.as_ref();
285 let mut suite = load_eval_suite(path)?;
286 if let Some(agent) = &agent_override {
287 suite.agent = Some(agent.clone());
288 }
289 suite.validate(agent_override.as_ref())?;
290
291 let base_dir = path.parent().unwrap_or_else(|| Path::new("."));
292 let agent_path = if let Some(agent) = agent_override {
293 agent
294 } else {
295 let agent = suite.agent.ok_or_else(|| {
296 EvalError::Config("agent path is required in suite or CLI".into())
297 })?;
298 if agent.is_absolute() {
299 agent
300 } else {
301 base_dir.join(agent)
302 }
303 };
304 let content = std::fs::read_to_string(&agent_path)?;
305 let spec = AgentSpec::from_yaml_strict(&content).map_err(|error| {
306 EvalError::Config(format!(
307 "invalid agent configuration '{}': {}",
308 agent_path.display(),
309 error
310 ))
311 })?;
312 spec.validate().map_err(|error| {
313 EvalError::Config(format!(
314 "invalid agent configuration '{}': {}",
315 agent_path.display(),
316 error
317 ))
318 })
319 }
320
321 pub async fn run(&self) -> Result<EvalResult> {
322 let start = Instant::now();
323 let base_dir = self.suite_path.parent().unwrap_or_else(|| Path::new("."));
324 let scenarios = self.filtered_scenarios();
325 if scenarios.is_empty() {
326 return Err(EvalError::Config(
327 "scenario selection matched zero scenarios".to_string(),
328 ));
329 }
330 let agent_path = self.resolve_agent_path(base_dir)?;
331 let results = if self.suite.settings.parallel && !self.suite.settings.fail_fast {
332 self.run_scenarios_parallel(&agent_path, base_dir, scenarios)
333 .await
334 } else {
335 self.run_scenarios_serial(&agent_path, base_dir, scenarios)
336 .await
337 };
338
339 let total = results.len();
340 let passed = results.iter().filter(|r| r.status.is_passed()).count();
341 let failed = results
342 .iter()
343 .filter(|r| r.status.is_failed() || r.status.is_error())
344 .count();
345 let skipped = results
346 .iter()
347 .filter(|r| matches!(r.status, ScenarioStatus::Skipped { .. }))
348 .count();
349 let metrics = compute_metrics(&results);
350
351 let observability = final_observability_report(&results);
352
353 Ok(EvalResult {
354 schema_version: 1,
355 suite: self.suite.name.clone(),
356 agent: agent_path.display().to_string(),
357 total,
358 passed,
359 failed,
360 skipped,
361 duration_ms: start.elapsed().as_millis() as u64,
362 scenarios: results,
363 metrics,
364 observability,
365 })
366 }
367
368 async fn run_scenarios_serial(
369 &self,
370 agent_path: &Path,
371 base_dir: &Path,
372 scenarios: Vec<&Scenario>,
373 ) -> Vec<ScenarioResult> {
374 let mut results = Vec::new();
375 for scenario in scenarios {
376 let result = self.run_scenario(agent_path, base_dir, scenario).await;
377 match result {
378 Ok(result) => {
379 let stop = self.suite.settings.fail_fast
380 && (result.status.is_failed() || result.status.is_error());
381 results.push(result);
382 if stop {
383 break;
384 }
385 }
386 Err(error) => {
387 results.push(error_result(scenario, error, FailureCategory::RuntimeError));
388 if self.suite.settings.fail_fast {
389 break;
390 }
391 }
392 }
393 }
394 results
395 }
396
397 async fn run_scenarios_parallel(
398 &self,
399 agent_path: &Path,
400 base_dir: &Path,
401 scenarios: Vec<&Scenario>,
402 ) -> Vec<ScenarioResult> {
403 let max_concurrent = self.suite.settings.max_concurrent.max(1);
404 let mut indexed = stream::iter(scenarios.into_iter().enumerate())
405 .map(|(idx, scenario)| async move {
406 let result = self.run_scenario(agent_path, base_dir, scenario).await;
407 let result = result.unwrap_or_else(|error| {
408 error_result(scenario, error, FailureCategory::RuntimeError)
409 });
410 (idx, result)
411 })
412 .buffer_unordered(max_concurrent)
413 .collect::<Vec<_>>()
414 .await;
415 indexed.sort_by_key(|(idx, _)| *idx);
416 indexed.into_iter().map(|(_, result)| result).collect()
417 }
418
419 fn resolve_agent_path(&self, base_dir: &Path) -> Result<PathBuf> {
420 if let Some(agent) = &self.options.agent {
421 return Ok(agent.clone());
422 }
423 let agent =
424 self.suite.agent.clone().ok_or_else(|| {
425 EvalError::Config("agent path is required in suite or CLI".into())
426 })?;
427 Ok(if agent.is_absolute() {
428 agent
429 } else {
430 base_dir.join(agent)
431 })
432 }
433
434 fn filtered_scenarios(&self) -> Vec<&Scenario> {
435 let ids: HashSet<_> = self.options.ids.iter().collect();
436 let tags: HashSet<_> = self.options.tags.iter().collect();
437 let languages: HashSet<_> = self.options.languages.iter().collect();
438 self.suite
439 .scenarios
440 .iter()
441 .filter(|scenario| {
442 if !ids.is_empty() && !ids.contains(&scenario.id) {
443 return false;
444 }
445 if !languages.is_empty() {
446 let Some(language) = &scenario.language else {
447 return false;
448 };
449 if !languages.contains(language) {
450 return false;
451 }
452 }
453 if !tags.is_empty() {
454 let scenario_tags: HashSet<_> = scenario.tags.iter().collect();
455 if self.options.tag_mode_all {
456 if !tags.iter().all(|tag| scenario_tags.contains(*tag)) {
457 return false;
458 }
459 } else if !tags.iter().any(|tag| scenario_tags.contains(*tag)) {
460 return false;
461 }
462 }
463 true
464 })
465 .collect()
466 }
467
468 async fn run_scenario(
469 &self,
470 agent_path: &Path,
471 base_dir: &Path,
472 scenario: &Scenario,
473 ) -> Result<ScenarioResult> {
474 let start = Instant::now();
475 if scenario.skip.is_skipped() {
476 return Ok(ScenarioResult {
477 id: scenario.id.clone(),
478 name: scenario.name.clone(),
479 tags: scenario.tags.clone(),
480 language: scenario.language.clone(),
481 status: ScenarioStatus::Skipped {
482 reason: scenario.skip.reason(),
483 },
484 failure_category: None,
485 flaky: false,
486 attempts: Vec::new(),
487 duration_ms: 0,
488 retries_used: 0,
489 });
490 }
491
492 let mut attempts = Vec::new();
493 let mut final_status = ScenarioStatus::Failed {
494 reason: "not run".to_string(),
495 };
496 let mut category = Some(FailureCategory::AssertionFailed);
497 let max_attempt = self.suite.settings.retries + 1;
498 let budget = if scenario.budget.is_configured() {
499 Some(ScenarioBudgetTracker::new(
500 scenario.budget.clone(),
501 self.budget_cost_estimator(base_dir, scenario)?,
502 ))
503 } else {
504 None
505 };
506
507 for attempt_idx in 0..max_attempt {
508 let attempt_future =
509 self.run_attempt(agent_path, base_dir, scenario, attempt_idx, budget.clone());
510 let attempt = if let Some(timeout_ms) = self.suite.settings.timeout_per_scenario_ms {
511 match timeout(Duration::from_millis(timeout_ms), attempt_future).await {
512 Ok(result) => result,
513 Err(_) => Err(EvalError::Runtime(format!(
514 "scenario '{}' attempt {} timed out after {}ms",
515 scenario.id, attempt_idx, timeout_ms
516 ))),
517 }
518 } else {
519 attempt_future.await
520 };
521 match attempt {
522 Ok(attempt_result) => {
523 final_status = attempt_result.status.clone();
524 if final_status.is_passed() {
525 attempts.push(attempt_result);
526 category = if attempt_idx > 0 {
527 Some(FailureCategory::FlakyPass)
528 } else {
529 None
530 };
531 break;
532 }
533 category = Some(if final_status.is_error() {
534 FailureCategory::RuntimeError
535 } else {
536 failure_category_for_attempt(&attempt_result)
537 });
538 attempts.push(attempt_result);
539 }
540 Err(error) => {
541 final_status = ScenarioStatus::Error {
542 message: error.to_string(),
543 };
544 category = Some(FailureCategory::RuntimeError);
545 attempts.push(AttemptResult {
546 attempt: attempt_idx,
547 turns: Vec::new(),
548 status: final_status.clone(),
549 duration_ms: 0,
550 });
551 }
552 }
553 if budget
554 .as_ref()
555 .is_some_and(ScenarioBudgetTracker::has_failed)
556 {
557 break;
558 }
559 if attempt_idx + 1 < max_attempt {
560 tokio::time::sleep(Duration::from_millis(self.suite.settings.retry_delay_ms)).await;
561 }
562 }
563
564 let flaky = final_status.is_passed() && attempts.len() > 1;
565 Ok(ScenarioResult {
566 id: scenario.id.clone(),
567 name: scenario.name.clone(),
568 tags: scenario.tags.clone(),
569 language: scenario.language.clone(),
570 status: final_status,
571 failure_category: category,
572 flaky,
573 duration_ms: start.elapsed().as_millis() as u64,
574 retries_used: attempts.len().saturating_sub(1) as u32,
575 attempts,
576 })
577 }
578
579 async fn run_attempt(
580 &self,
581 agent_path: &Path,
582 base_dir: &Path,
583 scenario: &Scenario,
584 attempt: u32,
585 budget: Option<ScenarioBudgetTracker>,
586 ) -> Result<AttemptResult> {
587 let start = Instant::now();
588 let _env_guard = EnvGuard::apply(&scenario.env)?;
589 let mock_server = start_mock_server(self.suite.fixtures.mock_server.as_ref()).await?;
590 let attempt_context = AttemptWorkspace::create(mock_server.as_ref())?;
591 let tool_log = RecordingToolLog::new();
592 let approval_log = RecordingApprovalLog::default();
593 let llm_log = RecordingLlmLog::default();
594 let approval_handler = self
595 .suite
596 .fixtures
597 .approvals
598 .as_ref()
599 .map(build_approval_handler);
600 let mut agent = self
601 .build_agent(BuildAgentParams {
602 agent_path,
603 base_dir,
604 attempt_context: &attempt_context,
605 tool_log: tool_log.clone(),
606 approval_log: approval_log.clone(),
607 llm_log: llm_log.clone(),
608 approval_handler: approval_handler.clone(),
609 budget: budget.clone(),
610 })
611 .await?;
612 apply_base_context(&agent, &self.suite, base_dir, scenario, &attempt_context)?;
613 let mut turns = Vec::new();
614 let mut status = ScenarioStatus::Passed;
615
616 if !scenario.turns.is_empty() {
617 for (idx, turn) in scenario.turns.iter().enumerate() {
618 let turn_execution = self
619 .run_turn(RunTurnParams {
620 agent: &agent,
621 scenario,
622 turn,
623 index: idx,
624 tool_log: &tool_log,
625 approval_log: &approval_log,
626 llm_log: &llm_log,
627 })
628 .await?;
629 let turn_failed = turn_execution
630 .result
631 .assertion_results
632 .iter()
633 .any(|result| !result.passed);
634 let runtime_error = turn_execution.unhandled_runtime_error;
635 turns.push(turn_execution.result);
636 if let Some(message) = runtime_error {
637 status = ScenarioStatus::Error { message };
638 break;
639 }
640 if turn_failed {
641 status = ScenarioStatus::Failed {
642 reason: format!("turn {} assertion failed", idx + 1),
643 };
644 break;
645 }
646 if self.suite.settings.isolation == IsolationMode::Turn
647 && idx + 1 < scenario.turns.len()
648 {
649 agent.reset().await?;
650 apply_base_context(&agent, &self.suite, base_dir, scenario, &attempt_context)?;
651 }
652 }
653 }
654
655 for step in &scenario.steps {
656 if !status.is_passed() {
657 break;
658 }
659 match step {
660 ScenarioStep::Run(run) => {
661 for turn in &run.turns {
662 let idx = turns.len();
663 let turn_execution = self
664 .run_turn(RunTurnParams {
665 agent: &agent,
666 scenario,
667 turn,
668 index: idx,
669 tool_log: &tool_log,
670 approval_log: &approval_log,
671 llm_log: &llm_log,
672 })
673 .await?;
674 let turn_failed = turn_execution
675 .result
676 .assertion_results
677 .iter()
678 .any(|result| !result.passed);
679 let runtime_error = turn_execution.unhandled_runtime_error;
680 turns.push(turn_execution.result);
681 if let Some(message) = runtime_error {
682 status = ScenarioStatus::Error { message };
683 break;
684 }
685 if turn_failed {
686 status = ScenarioStatus::Failed {
687 reason: format!("turn {} assertion failed", idx + 1),
688 };
689 break;
690 }
691 }
692 if status.is_passed()
693 && let Some(session) = &run.save_session
694 {
695 agent.save_session(session).await?;
696 }
697 }
698 ScenarioStep::ResetAgent(reset) => {
699 if let Some(options) = reset_options(reset) {
700 if options.delete_persistence || !options.preserve_storage {
701 let _ = std::fs::remove_dir_all(&attempt_context.workspace);
702 std::fs::create_dir_all(&attempt_context.workspace)?;
703 }
704 let preserved_actor = options
705 .preserve_actor_id
706 .then(|| agent.actor_id())
707 .flatten();
708 if matches!(options.profile, crate::reset::ResetProfile::Conversation)
709 && !options.delete_persistence
710 {
711 agent.reset().await?;
712 } else {
713 agent = self
714 .build_agent(BuildAgentParams {
715 agent_path,
716 base_dir,
717 attempt_context: &attempt_context,
718 tool_log: tool_log.clone(),
719 approval_log: approval_log.clone(),
720 llm_log: llm_log.clone(),
721 approval_handler: approval_handler.clone(),
722 budget: budget.clone(),
723 })
724 .await?;
725 }
726 if options.preserve_host_context {
727 apply_base_context(
728 &agent,
729 &self.suite,
730 base_dir,
731 scenario,
732 &attempt_context,
733 )?;
734 } else {
735 apply_context_map(&agent, attempt_context.runtime_context())?;
736 }
737 if let Some(actor) = preserved_actor.or_else(|| scenario.actor.clone()) {
738 agent.set_actor_id(&actor)?;
739 agent.load_actor_memory().await?;
740 agent.load_actor_relationship().await?;
741 }
742 }
743 }
744 ScenarioStep::SaveSession(name) => {
745 agent.save_session(name).await?;
746 }
747 ScenarioStep::LoadSession(name) => {
748 let _ = agent.load_session(name).await?;
749 }
750 ScenarioStep::SetContext { values } => {
751 apply_context_value(&agent, values)?;
752 }
753 ScenarioStep::SetActor { actor } => {
754 agent.set_actor_id(actor)?;
755 agent.load_actor_memory().await?;
756 agent.load_actor_relationship().await?;
757 }
758 ScenarioStep::CleanupExpired => {
759 let _ = agent.cleanup_expired_sessions().await?;
760 }
761 }
762 }
763
764 Ok(AttemptResult {
765 attempt,
766 turns,
767 status,
768 duration_ms: start.elapsed().as_millis() as u64,
769 })
770 }
771
772 async fn build_agent(&self, params: BuildAgentParams<'_>) -> Result<RuntimeAgent> {
773 let BuildAgentParams {
774 agent_path,
775 base_dir,
776 attempt_context,
777 tool_log,
778 approval_log,
779 llm_log,
780 approval_handler,
781 budget,
782 } = params;
783 let content = std::fs::read_to_string(agent_path)?;
784 let mut spec = AgentSpec::from_yaml_strict(&content)?;
785 apply_eval_llm_settings(&mut spec, &self.suite.settings);
786 isolate_spec_storage(&mut spec, attempt_context);
787 apply_workspace_policy(
788 &mut spec,
789 self.suite.fixtures.workspace_policy.as_ref(),
790 attempt_context,
791 )?;
792 spec.validate()
793 .map_err(|error| EvalError::Config(error.to_string()))?;
794 let llm_fixture = attempt_context.interpolate_llm_fixture(&self.suite.fixtures.llm)?;
795 let provider_configs = budget_provider_configs(&spec);
796 let (mut llm_registry, _judge_llm) = build_llm_registry(&spec, &llm_fixture, base_dir)?;
797 if let Some(budget) = budget {
798 llm_registry =
799 llm_registry.map_providers(|alias, provider| {
800 let config = provider_configs.get(alias).cloned().unwrap_or_else(|| {
801 BudgetProviderConfig {
802 provider: provider.provider_name().to_string(),
803 model: alias.to_string(),
804 max_output_tokens: 2_000,
805 }
806 });
807 budget.wrap(provider, config)
808 });
809 }
810 let tool_registry = build_tool_registry(&self.suite.fixtures, tool_log.clone())?;
811 let agent_base_dir = agent_path.parent().unwrap_or_else(|| Path::new("."));
812 let mut builder = AgentBuilder::from_spec_with_base_dir(spec, agent_base_dir)
813 .llm_registry(llm_registry)
814 .tools(tool_registry)
815 .hooks(Arc::new(EvalRecordHooks {
816 tool_log,
817 approval_log,
818 llm_log,
819 }))
820 .auto_configure_features()
821 .map_err(|error| EvalError::Config(error.to_string()))?
822 .auto_configure_mcp()
823 .await
824 .map_err(|error| EvalError::Config(error.to_string()))?;
825
826 if let Some(approval_handler) = approval_handler {
827 builder = builder.approval_handler(approval_handler);
828 }
829
830 if let Some(observability) = self.observability_config(base_dir)? {
831 let manager = ai_agents_observability::ObservabilityManager::new(observability);
832 builder = builder.observability(manager);
833 }
834 builder = builder
835 .auto_configure_spawner()
836 .await
837 .map_err(|error| EvalError::Config(error.to_string()))?;
838 let agent = builder
839 .build()
840 .map_err(|error| EvalError::Config(error.to_string()))?;
841 agent.init_storage().await?;
842 Ok(agent)
843 }
844
845 fn observability_config(&self, base_dir: &Path) -> Result<Option<ObservabilityConfig>> {
846 let mut config = if let Some(config) = self.suite.observability.clone() {
847 config
848 } else if self.options.observability {
849 ObservabilityConfig {
850 enabled: true,
851 export: ExportConfig {
852 formats: vec![ExportFormat::Json],
853 path: self
854 .options
855 .output
856 .join("observability")
857 .display()
858 .to_string(),
859 write_report: true,
860 ..Default::default()
861 },
862 ..Default::default()
863 }
864 } else {
865 return Ok(None);
866 };
867 if !config.enabled {
868 return Ok(None);
869 }
870 config = config
871 .with_pricing_file_loaded(Some(base_dir))
872 .map_err(|error| EvalError::Config(error.to_string()))?;
873 config
874 .validate()
875 .map_err(|error| EvalError::Config(error.to_string()))?;
876 Ok(Some(config))
877 }
878
879 fn budget_cost_estimator(
880 &self,
881 base_dir: &Path,
882 scenario: &Scenario,
883 ) -> Result<Option<CostEstimator>> {
884 if scenario.budget.max_cost_usd.is_none() {
885 return Ok(None);
886 }
887 let config = self.suite.observability.clone().ok_or_else(|| {
888 EvalError::Config(format!(
889 "scenario '{}' budget.max_cost_usd requires suite observability.cost pricing",
890 scenario.id
891 ))
892 })?;
893 let config = config
894 .with_pricing_file_loaded(Some(base_dir))
895 .map_err(|error| EvalError::Config(error.to_string()))?;
896 if !config.cost.enabled {
897 return Err(EvalError::Config(format!(
898 "scenario '{}' budget.max_cost_usd requires observability.cost.enabled: true",
899 scenario.id
900 )));
901 }
902 Ok(Some(CostEstimator::new(config.cost)))
903 }
904
905 async fn run_turn(&self, params: RunTurnParams<'_>) -> Result<TurnExecution> {
906 let RunTurnParams {
907 agent,
908 scenario,
909 turn,
910 index,
911 tool_log,
912 approval_log,
913 llm_log,
914 } = params;
915 apply_context_value(agent, &turn_runtime_context(turn))?;
916 if let Some(actor) = &turn.actor {
917 agent.set_actor_id(actor)?;
918 }
919 let before_relationship = relationship_snapshot(agent);
920 let tool_start = tool_log.len();
921 let approval_start = approval_log.len();
922 let llm_start = llm_log.len();
923 let start = Instant::now();
924 let timeout_ms = turn
925 .timeout_ms
926 .unwrap_or(self.suite.settings.timeout_per_turn_ms);
927 let mut operation = if turn.stream.unwrap_or(false) {
928 collect_stream_response(agent, &turn.input, timeout_ms).await
929 } else {
930 match timeout(Duration::from_millis(timeout_ms), agent.chat(&turn.input)).await {
931 Ok(Ok(response)) => TurnOperation {
932 response_content: response.content,
933 response_metadata: response.metadata,
934 response_present: true,
935 runtime_error: None,
936 },
937 Ok(Err(error)) => TurnOperation::error(error.to_string()),
938 Err(_) => TurnOperation::error(format!("turn timed out after {}ms", timeout_ms)),
939 }
940 };
941 if let Err(error) = agent.flush_background_tasks().await
942 && operation.runtime_error.is_none()
943 {
944 operation.runtime_error = Some(error.to_string());
945 }
946 let latency_ms = start.elapsed().as_millis() as u64;
947 let mut evidence = collect_turn_evidence(
948 agent,
949 operation.response_metadata.clone(),
950 tool_log,
951 tool_start,
952 before_relationship,
953 );
954 evidence.approvals = approval_log.records_since(approval_start);
955 evidence.llm_requests = llm_log.records_since(llm_start);
956 let judge = self.build_judge(agent);
957 let mut assertion_results = if let Some(assertion) = &turn.assertions {
958 match evaluate_assertion(
959 assertion,
960 AssertionEvalContext {
961 evidence: &evidence,
962 response: &operation.response_content,
963 user_input: Some(&turn.input),
964 scenario_id: Some(&scenario.id),
965 language: scenario.language.as_deref(),
966 judge_resolver: Some(&judge),
967 },
968 )
969 .await
970 {
971 AssertionOutcome::Passed(details) | AssertionOutcome::Failed(details) => details,
972 AssertionOutcome::Error(message) => return Err(EvalError::Assertion(message)),
973 }
974 } else {
975 Vec::new()
976 };
977 if !operation.response_present
978 && turn
979 .assertions
980 .as_ref()
981 .is_some_and(assertion_uses_response)
982 {
983 assertion_results.push(AssertionResultDetail {
984 assertion: "response_present".to_string(),
985 passed: false,
986 actual: json!(false),
987 expected: json!(true),
988 message: Some("response assertions require a runtime response".to_string()),
989 });
990 }
991
992 let expected_error = turn_expected_error(turn);
993 let unhandled_runtime_error = match (&expected_error, &operation.runtime_error) {
994 (Some(expected), Some(error)) if expected.matches(error) => {
995 assertion_results.push(AssertionResultDetail {
996 assertion: "expect_error".to_string(),
997 passed: true,
998 actual: json!(error),
999 expected: json!(expected.items()),
1000 message: None,
1001 });
1002 None
1003 }
1004 (Some(expected), Some(error)) => {
1005 assertion_results.push(AssertionResultDetail {
1006 assertion: "expect_error".to_string(),
1007 passed: false,
1008 actual: json!(error),
1009 expected: json!(expected.items()),
1010 message: Some("runtime error did not match any expected substring".to_string()),
1011 });
1012 Some(format!(
1013 "runtime error did not match expect_error: {}",
1014 error
1015 ))
1016 }
1017 (Some(expected), None) => {
1018 assertion_results.push(AssertionResultDetail {
1019 assertion: "expect_error".to_string(),
1020 passed: false,
1021 actual: Value::Null,
1022 expected: json!(expected.items()),
1023 message: Some("expected a runtime error but the turn completed".to_string()),
1024 });
1025 None
1026 }
1027 (None, Some(error)) => Some(error.clone()),
1028 (None, None) => None,
1029 };
1030 if self.suite.settings.redact_outputs {
1031 redact_assertion_details(&mut assertion_results);
1032 }
1033 let observability_span_id = evidence
1034 .observability
1035 .as_ref()
1036 .and_then(|obs| obs.span_ids.last().cloned());
1037 let runtime_error = operation
1038 .runtime_error
1039 .as_deref()
1040 .map(|error| redact_text(error, self.suite.settings.redact_outputs, 0));
1041 let unhandled_runtime_error = unhandled_runtime_error
1042 .map(|error| redact_text(&error, self.suite.settings.redact_outputs, 0).value);
1043 Ok(TurnExecution {
1044 result: TurnResult {
1045 index,
1046 input: redact_text(&turn.input, self.suite.settings.redact_outputs, 0),
1047 response: if operation.response_present {
1048 redact_text(
1049 &operation.response_content,
1050 self.suite.settings.redact_outputs,
1051 0,
1052 )
1053 } else {
1054 crate::redaction::RedactedString::plain("")
1055 },
1056 response_present: operation.response_present,
1057 runtime_error,
1058 state: evidence.state.clone(),
1059 metadata: if self.suite.settings.redact_outputs {
1060 None
1061 } else {
1062 operation
1063 .response_metadata
1064 .and_then(|metadata| serde_json::to_value(metadata).ok())
1065 },
1066 evidence,
1067 assertion_results,
1068 latency_ms,
1069 observability_span_id,
1070 },
1071 unhandled_runtime_error,
1072 })
1073 }
1074
1075 fn build_judge(&self, agent: &RuntimeAgent) -> JudgeResolver {
1076 JudgeResolver::new(Arc::clone(agent.llm_registry()), JudgeConfig::default())
1077 }
1078}
1079
1080struct TurnExecution {
1081 result: TurnResult,
1082 unhandled_runtime_error: Option<String>,
1083}
1084
1085struct TurnOperation {
1086 response_content: String,
1087 response_metadata: Option<HashMap<String, Value>>,
1088 response_present: bool,
1089 runtime_error: Option<String>,
1090}
1091
1092impl TurnOperation {
1093 fn error(message: String) -> Self {
1094 Self {
1095 response_content: String::new(),
1096 response_metadata: None,
1097 response_present: false,
1098 runtime_error: Some(message),
1099 }
1100 }
1101}
1102
1103async fn collect_stream_response(
1104 agent: &RuntimeAgent,
1105 input: &str,
1106 timeout_ms: u64,
1107) -> TurnOperation {
1108 let deadline = TokioInstant::now() + Duration::from_millis(timeout_ms);
1109 let stream = match timeout_at(deadline, agent.chat_stream(input)).await {
1110 Ok(Ok(stream)) => stream,
1111 Ok(Err(error)) => return TurnOperation::error(error.to_string()),
1112 Err(_) => return TurnOperation::error(format!("turn timed out after {}ms", timeout_ms)),
1113 };
1114 consume_stream_response(stream, deadline, timeout_ms).await
1115}
1116
1117async fn consume_stream_response<S>(
1118 mut stream: S,
1119 deadline: TokioInstant,
1120 timeout_ms: u64,
1121) -> TurnOperation
1122where
1123 S: futures::Stream<Item = StreamChunk> + Unpin,
1124{
1125 let mut content = String::new();
1126 let mut content_seen = false;
1127 let mut runtime_error = None;
1128 let mut done = false;
1129 loop {
1130 match timeout_at(deadline, stream.next()).await {
1131 Ok(Some(StreamChunk::Content { text })) => {
1132 content_seen = true;
1133 content.push_str(&text);
1134 }
1135 Ok(Some(StreamChunk::Done {})) => {
1136 done = true;
1137 break;
1138 }
1139 Ok(Some(StreamChunk::Error { message })) => {
1140 if runtime_error.is_none() {
1141 runtime_error = Some(message);
1142 }
1143 }
1144 Ok(Some(_)) => {}
1145 Ok(None) => break,
1146 Err(_) => {
1147 if runtime_error.is_none() {
1148 runtime_error = Some(format!("turn timed out after {}ms", timeout_ms));
1149 }
1150 break;
1151 }
1152 }
1153 }
1154 if !done && runtime_error.is_none() {
1155 runtime_error = Some("stream ended before Done".to_string());
1156 }
1157 TurnOperation {
1158 response_content: content,
1159 response_metadata: None,
1160 response_present: content_seen || (done && runtime_error.is_none()),
1161 runtime_error,
1162 }
1163}
1164
1165fn assertion_uses_response(assertion: &Assertion) -> bool {
1166 assertion.response_contains.is_some()
1167 || assertion.response_contains_any.is_some()
1168 || assertion.response_not_contains.is_some()
1169 || assertion.response_not_empty.is_some()
1170 || assertion.response_semantic.is_some()
1171 || assertion.judge.is_some()
1172 || assertion
1173 .all
1174 .as_ref()
1175 .is_some_and(|children| children.iter().any(assertion_uses_response))
1176 || assertion
1177 .any
1178 .as_ref()
1179 .is_some_and(|children| children.iter().any(assertion_uses_response))
1180 || assertion
1181 .not
1182 .as_deref()
1183 .is_some_and(assertion_uses_response)
1184}
1185
1186fn apply_workspace_policy(
1187 spec: &mut AgentSpec,
1188 config: Option<&WorkspacePolicyFixtureConfig>,
1189 context: &AttemptFixtureContext,
1190) -> Result<()> {
1191 let Some(config) = config else {
1192 return Ok(());
1193 };
1194 if !context.workspace.is_absolute() {
1195 return Err(EvalError::Config(
1196 "eval attempt workspace must be absolute".to_string(),
1197 ));
1198 }
1199
1200 for tool_id in config.read_tools.iter().chain(&config.write_tools) {
1201 if !spec.tool_security.tools.contains_key(tool_id) {
1202 return Err(EvalError::Config(format!(
1203 "fixtures.workspace_policy names tool '{}' without an existing tool policy",
1204 tool_id
1205 )));
1206 }
1207 }
1208
1209 let workspace = context.workspace.display().to_string();
1210 for tool_id in &config.read_tools {
1211 spec.tool_security
1212 .tools
1213 .get_mut(tool_id)
1214 .expect("workspace policy tool was validated")
1215 .read_paths
1216 .push(workspace.clone());
1217 }
1218 for tool_id in &config.write_tools {
1219 spec.tool_security
1220 .tools
1221 .get_mut(tool_id)
1222 .expect("workspace policy tool was validated")
1223 .write_paths
1224 .push(workspace.clone());
1225 }
1226 Ok(())
1227}
1228
1229fn isolate_spec_storage(spec: &mut AgentSpec, context: &AttemptFixtureContext) {
1230 isolate_storage_config(
1231 &mut spec.storage,
1232 context,
1233 "parent-storage",
1234 "parent-storage.db",
1235 );
1236 if let Some(shared_storage) = spec
1237 .spawner
1238 .as_mut()
1239 .and_then(|spawner| spawner.shared_storage.as_mut())
1240 {
1241 isolate_storage_config(
1242 shared_storage,
1243 context,
1244 "spawner-shared-storage",
1245 "spawner-shared-storage.db",
1246 );
1247 }
1248}
1249
1250fn isolate_storage_config(
1251 storage: &mut StorageConfig,
1252 context: &AttemptFixtureContext,
1253 file_name: &str,
1254 sqlite_name: &str,
1255) {
1256 match storage {
1257 StorageConfig::File(config) => {
1258 config.path = context.workspace.join(file_name).display().to_string();
1259 }
1260 StorageConfig::Sqlite(config) => {
1261 config.path = context.workspace.join(sqlite_name).display().to_string();
1262 }
1263 StorageConfig::Redis(config) => {
1264 let prefix = config.prefix.as_deref().unwrap_or("agent:");
1265 config.prefix = Some(format!("{}eval:{}:", prefix, context.isolation_id));
1266 }
1267 StorageConfig::None => {}
1268 }
1269}
1270
1271fn apply_context_map(agent: &RuntimeAgent, values: HashMap<String, Value>) -> Result<()> {
1272 for (key, value) in values {
1273 agent.set_context(&key, value)?;
1274 }
1275 Ok(())
1276}
1277
1278fn apply_context_value(agent: &RuntimeAgent, value: &Value) -> Result<()> {
1279 let Value::Object(map) = value else {
1280 return Ok(());
1281 };
1282 for (key, value) in map {
1283 agent.set_context(key, value.clone())?;
1284 }
1285 Ok(())
1286}
1287
1288fn apply_base_context(
1289 agent: &RuntimeAgent,
1290 suite: &EvalSuite,
1291 base_dir: &Path,
1292 scenario: &Scenario,
1293 attempt_context: &AttemptFixtureContext,
1294) -> Result<()> {
1295 apply_context_map(agent, resolve_fixture_context(&suite.fixtures, base_dir)?)?;
1296 apply_context_map(agent, attempt_context.runtime_context())?;
1297 apply_context_value(agent, &scenario.context)?;
1298 if let Some(actor) = &scenario.actor {
1299 agent.set_actor_id(actor)?;
1300 }
1301 Ok(())
1302}
1303
1304fn reset_options(config: &ResetStepConfig) -> Option<crate::reset::ResetOptions> {
1305 match config {
1306 ResetStepConfig::Bool(false) => None,
1307 ResetStepConfig::Bool(true) => Some(crate::reset::ResetOptions::default()),
1308 ResetStepConfig::Options(options) => Some(options.clone()),
1309 }
1310}
1311
1312fn redact_assertion_details(details: &mut [crate::assertion::AssertionResultDetail]) {
1313 for detail in details {
1314 detail.actual = redact_value(std::mem::take(&mut detail.actual), true, 0);
1315 detail.expected = redact_value(std::mem::take(&mut detail.expected), true, 0);
1316 }
1317}
1318
1319fn error_result(
1320 scenario: &Scenario,
1321 error: EvalError,
1322 category: FailureCategory,
1323) -> ScenarioResult {
1324 ScenarioResult {
1325 id: scenario.id.clone(),
1326 name: scenario.name.clone(),
1327 tags: scenario.tags.clone(),
1328 language: scenario.language.clone(),
1329 status: ScenarioStatus::Error {
1330 message: error.to_string(),
1331 },
1332 failure_category: Some(category),
1333 flaky: false,
1334 attempts: Vec::new(),
1335 duration_ms: 0,
1336 retries_used: 0,
1337 }
1338}
1339
1340fn failure_category_for_attempt(attempt: &AttemptResult) -> FailureCategory {
1341 let judge_failed = attempt.turns.iter().any(|turn| {
1342 turn.assertion_results
1343 .iter()
1344 .any(|detail| !detail.passed && detail.assertion == "judge")
1345 });
1346 if judge_failed {
1347 FailureCategory::JudgeError
1348 } else {
1349 FailureCategory::AssertionFailed
1350 }
1351}
1352
1353fn final_observability_report(
1354 results: &[ScenarioResult],
1355) -> Option<ai_agents_observability::ObservabilityReport> {
1356 results
1357 .iter()
1358 .rev()
1359 .flat_map(|scenario| scenario.attempts.iter().rev())
1360 .flat_map(|attempt| attempt.turns.iter().rev())
1361 .find_map(|turn| {
1362 turn.evidence
1363 .observability
1364 .as_ref()
1365 .and_then(|obs| obs.report.clone())
1366 })
1367}
1368
1369fn apply_eval_llm_settings(spec: &mut AgentSpec, settings: &crate::suite::EvalSettings) {
1370 if let LLMConfigOrSelector::Config(config) = &mut spec.llm {
1371 apply_llm_config_settings(config, settings);
1372 }
1373 for config in spec.llms.values_mut() {
1374 apply_llm_config_settings(config, settings);
1375 }
1376}
1377
1378fn budget_provider_configs(spec: &AgentSpec) -> HashMap<String, BudgetProviderConfig> {
1379 if spec.llms.is_empty() {
1380 let config = spec.llm.as_config().cloned().unwrap_or_default();
1381 return HashMap::from([(
1382 "default".to_string(),
1383 BudgetProviderConfig {
1384 provider: config.provider,
1385 model: config.model,
1386 max_output_tokens: config.max_tokens,
1387 },
1388 )]);
1389 }
1390 spec.llms
1391 .iter()
1392 .map(|(alias, config)| {
1393 (
1394 alias.clone(),
1395 BudgetProviderConfig {
1396 provider: config.provider.clone(),
1397 model: config.model.clone(),
1398 max_output_tokens: config.max_tokens,
1399 },
1400 )
1401 })
1402 .collect()
1403}
1404
1405fn apply_llm_config_settings(
1406 config: &mut ai_agents_runtime::spec::LLMConfig,
1407 settings: &crate::suite::EvalSettings,
1408) {
1409 if let Some(temperature) = settings.temperature {
1410 config.temperature = temperature;
1411 }
1412 if let Some(seed) = settings.seed {
1413 config.extra.insert("seed".to_string(), json!(seed));
1414 }
1415}
1416
1417enum EnvExclusionGuard {
1419 Read {
1420 _guard: RwLockReadGuard<'static, ()>,
1421 },
1422 Write {
1423 _guard: RwLockWriteGuard<'static, ()>,
1424 },
1425}
1426
1427struct EnvGuard {
1429 previous: Vec<(String, Option<String>)>,
1431 _guard: EnvExclusionGuard,
1433}
1434
1435impl EnvGuard {
1436 fn apply(values: &HashMap<String, String>) -> Result<Self> {
1437 static ENV_LOCK: OnceLock<RwLock<()>> = OnceLock::new();
1438 let lock = ENV_LOCK.get_or_init(|| RwLock::new(()));
1439 if values.is_empty() {
1440 let guard = lock.read().map_err(|_| {
1441 EvalError::Runtime("failed to lock eval environment guard".to_string())
1442 })?;
1443 return Ok(Self {
1444 previous: Vec::new(),
1445 _guard: EnvExclusionGuard::Read { _guard: guard },
1446 });
1447 }
1448 let guard = lock
1449 .write()
1450 .map_err(|_| EvalError::Runtime("failed to lock eval environment guard".to_string()))?;
1451 let mut previous = Vec::new();
1452 for (key, value) in values {
1453 previous.push((key.clone(), std::env::var(key).ok()));
1454 unsafe {
1455 std::env::set_var(key, value);
1456 }
1457 }
1458 Ok(Self {
1459 previous,
1460 _guard: EnvExclusionGuard::Write { _guard: guard },
1461 })
1462 }
1463}
1464
1465impl Drop for EnvGuard {
1466 fn drop(&mut self) {
1467 for (key, value) in self.previous.drain(..).rev() {
1468 unsafe {
1469 if let Some(value) = value {
1470 std::env::set_var(key, value);
1471 } else {
1472 std::env::remove_var(key);
1473 }
1474 }
1475 }
1476 }
1477}
1478
1479#[cfg(test)]
1480mod tests {
1481 use super::*;
1482
1483 fn attempt_workspace(id: &str) -> PathBuf {
1484 std::env::temp_dir().join(format!("ai-agents-eval-{id}"))
1485 }
1486
1487 #[test]
1488 fn strict_suite_loader_rejects_nested_observability_typos() {
1489 let error = parse_eval_suite_yaml(
1490 r#"
1491name: strict
1492agent: agent.yaml
1493observability:
1494 enabeld: true
1495scenarios:
1496 - id: scenario
1497 turns:
1498 - input: hello
1499"#,
1500 )
1501 .unwrap_err()
1502 .to_string();
1503
1504 assert!(error.contains("enabeld"), "{error}");
1505 }
1506
1507 #[test]
1508 fn storage_isolation_rewrites_parent_and_spawner_backends() {
1509 let workspace = attempt_workspace("attempt-a");
1510 let context = AttemptFixtureContext {
1511 isolation_id: "attempt-a".to_string(),
1512 workspace: workspace.clone(),
1513 mock_server_base_url: None,
1514 };
1515 let mut file_spec: AgentSpec = serde_yaml::from_str(
1516 r#"
1517name: FileAgent
1518system_prompt: test
1519storage: { type: file, path: ./parent }
1520spawner:
1521 shared_storage: { type: sqlite, path: ./shared.db, table: shared_sessions }
1522"#,
1523 )
1524 .unwrap();
1525 isolate_spec_storage(&mut file_spec, &context);
1526 assert_eq!(
1527 file_spec.storage.get_path().map(PathBuf::from),
1528 Some(workspace.join("parent-storage"))
1529 );
1530 let shared = file_spec
1531 .spawner
1532 .as_ref()
1533 .unwrap()
1534 .shared_storage
1535 .as_ref()
1536 .unwrap();
1537 assert_eq!(
1538 shared.get_path().map(PathBuf::from),
1539 Some(workspace.join("spawner-shared-storage.db"))
1540 );
1541 assert_eq!(shared.get_table(), Some("shared_sessions"));
1542
1543 let mut redis_spec: AgentSpec = serde_yaml::from_str(
1544 r#"
1545name: RedisAgent
1546system_prompt: test
1547storage: { type: redis, url: redis://localhost, prefix: "parent:" }
1548spawner:
1549 shared_storage: { type: redis, url: redis://localhost }
1550"#,
1551 )
1552 .unwrap();
1553 isolate_spec_storage(&mut redis_spec, &context);
1554 assert_eq!(redis_spec.storage.get_prefix(), "parent:eval:attempt-a:");
1555 assert_eq!(
1556 redis_spec
1557 .spawner
1558 .as_ref()
1559 .unwrap()
1560 .shared_storage
1561 .as_ref()
1562 .unwrap()
1563 .get_prefix(),
1564 "agent:eval:attempt-a:"
1565 );
1566
1567 let other_context = AttemptFixtureContext {
1568 isolation_id: "attempt-b".to_string(),
1569 workspace: attempt_workspace("attempt-b"),
1570 mock_server_base_url: None,
1571 };
1572 let mut other_redis: AgentSpec = serde_yaml::from_str(
1573 r#"
1574name: RedisAgent
1575system_prompt: test
1576storage: { type: redis, url: redis://localhost, prefix: "parent:" }
1577"#,
1578 )
1579 .unwrap();
1580 isolate_spec_storage(&mut other_redis, &other_context);
1581 assert_ne!(
1582 redis_spec.storage.get_prefix(),
1583 other_redis.storage.get_prefix()
1584 );
1585 }
1586
1587 #[test]
1588 fn workspace_policy_is_narrow_and_isolated_per_attempt() {
1589 let source: AgentSpec = serde_yaml::from_str(
1590 r#"
1591name: PolicyAgent
1592system_prompt: test
1593tool_security:
1594 enabled: true
1595 fail_closed: true
1596 tools:
1597 file_read:
1598 read_paths: [./source]
1599 blocked_paths: [./blocked]
1600 file_write:
1601 write_paths: [./output]
1602 blocked_paths: [./blocked]
1603 grep:
1604 read_paths: [./repository]
1605"#,
1606 )
1607 .unwrap();
1608 let config = WorkspacePolicyFixtureConfig {
1609 read_tools: vec!["file_read".to_string()],
1610 write_tools: vec!["file_write".to_string()],
1611 };
1612 let first_workspace = attempt_workspace("attempt-a");
1613 let second_workspace = attempt_workspace("attempt-b");
1614 let first_context = AttemptFixtureContext {
1615 isolation_id: "attempt-a".to_string(),
1616 workspace: first_workspace.clone(),
1617 mock_server_base_url: None,
1618 };
1619 let second_context = AttemptFixtureContext {
1620 isolation_id: "attempt-b".to_string(),
1621 workspace: second_workspace.clone(),
1622 mock_server_base_url: None,
1623 };
1624
1625 let mut first = source.clone();
1626 apply_workspace_policy(&mut first, Some(&config), &first_context).unwrap();
1627 let mut second = source.clone();
1628 apply_workspace_policy(&mut second, Some(&config), &second_context).unwrap();
1629
1630 assert!(first.tool_security.fail_closed);
1631 assert_eq!(
1632 first.tool_security.tools["file_read"].read_paths,
1633 vec![
1634 "./source".to_string(),
1635 first_workspace.display().to_string()
1636 ]
1637 );
1638 assert_eq!(
1639 first.tool_security.tools["file_write"].write_paths,
1640 vec![
1641 "./output".to_string(),
1642 first_workspace.display().to_string()
1643 ]
1644 );
1645 assert_eq!(
1646 first.tool_security.tools["file_read"].blocked_paths,
1647 vec!["./blocked"]
1648 );
1649 assert_eq!(
1650 first.tool_security.tools["file_write"].blocked_paths,
1651 vec!["./blocked"]
1652 );
1653 assert_eq!(
1654 first.tool_security.tools["grep"].read_paths,
1655 vec!["./repository"]
1656 );
1657 assert_eq!(
1658 second.tool_security.tools["file_read"].read_paths,
1659 vec![
1660 "./source".to_string(),
1661 second_workspace.display().to_string()
1662 ]
1663 );
1664 assert_eq!(
1665 source.tool_security.tools["file_read"].read_paths,
1666 vec!["./source"]
1667 );
1668 assert_eq!(
1669 source.tool_security.tools["file_write"].write_paths,
1670 vec!["./output"]
1671 );
1672 }
1673
1674 #[test]
1675 fn workspace_policy_rejects_unknown_tools_without_partial_mutation() {
1676 let mut spec: AgentSpec = serde_yaml::from_str(
1677 r#"
1678name: PolicyAgent
1679system_prompt: test
1680tool_security:
1681 enabled: true
1682 fail_closed: true
1683 tools:
1684 file_read:
1685 read_paths: [./source]
1686"#,
1687 )
1688 .unwrap();
1689 let original = spec.tool_security.tools["file_read"].read_paths.clone();
1690 let config = WorkspacePolicyFixtureConfig {
1691 read_tools: vec!["file_read".to_string(), "missing_tool".to_string()],
1692 write_tools: Vec::new(),
1693 };
1694 let context = AttemptFixtureContext {
1695 isolation_id: "attempt-a".to_string(),
1696 workspace: attempt_workspace("attempt-a"),
1697 mock_server_base_url: None,
1698 };
1699
1700 let error = apply_workspace_policy(&mut spec, Some(&config), &context).unwrap_err();
1701
1702 assert!(error.to_string().contains("missing_tool"));
1703 assert!(
1704 error
1705 .to_string()
1706 .contains("without an existing tool policy")
1707 );
1708 assert_eq!(spec.tool_security.tools["file_read"].read_paths, original);
1709 assert!(spec.tool_security.fail_closed);
1710 }
1711
1712 #[tokio::test]
1713 async fn generated_attempt_context_has_stable_reset_precedence() {
1714 let dir = std::env::temp_dir().join(format!(
1715 "ai_agents_eval_attempt_context_test_{}",
1716 uuid::Uuid::new_v4()
1717 ));
1718 std::fs::create_dir_all(&dir).unwrap();
1719 write_test_agent(&dir);
1720 let suite_path = dir.join("suite.yaml");
1721 std::fs::write(
1722 &suite_path,
1723 r#"
1724name: Attempt Context
1725agent: agent.yaml
1726fixtures:
1727 context:
1728 eval: { workspace: fixture-value }
1729 mock_server: { base_url: fixture-value }
1730 fixture_only: true
1731 llm:
1732 mode: mock
1733 responses: [ok]
1734scenarios:
1735 - id: context
1736 context: { scenario_only: true, precedence: scenario }
1737 turns:
1738 - input: test
1739 context: { precedence: turn }
1740"#,
1741 )
1742 .unwrap();
1743 let runner = EvalRunner::from_file(
1744 &suite_path,
1745 EvalRunnerOptions {
1746 output: dir.join("out"),
1747 ..Default::default()
1748 },
1749 )
1750 .unwrap();
1751 let attempt_context = AttemptFixtureContext {
1752 isolation_id: "stable-attempt".to_string(),
1753 workspace: dir
1754 .join("workspace")
1755 .canonicalize()
1756 .unwrap_or_else(|_| dir.join("workspace")),
1757 mock_server_base_url: Some("http://127.0.0.1:40000".to_string()),
1758 };
1759 std::fs::create_dir_all(&attempt_context.workspace).unwrap();
1760 let scenario = &runner.suite.scenarios[0];
1761 let tool_log = RecordingToolLog::new();
1762 let approval_log = RecordingApprovalLog::default();
1763 let llm_log = RecordingLlmLog::default();
1764 let first = runner
1765 .build_agent(BuildAgentParams {
1766 agent_path: &dir.join("agent.yaml"),
1767 base_dir: &dir,
1768 attempt_context: &attempt_context,
1769 tool_log: tool_log.clone(),
1770 approval_log: approval_log.clone(),
1771 llm_log: llm_log.clone(),
1772 approval_handler: None,
1773 budget: None,
1774 })
1775 .await
1776 .unwrap();
1777 apply_base_context(&first, &runner.suite, &dir, scenario, &attempt_context).unwrap();
1778 let first_context = first.get_context();
1779 assert_eq!(
1780 first_context["eval"]["workspace"],
1781 json!(attempt_context.workspace.display().to_string())
1782 );
1783 assert_eq!(
1784 first_context["mock_server"]["base_url"],
1785 "http://127.0.0.1:40000"
1786 );
1787 assert_eq!(first_context["precedence"], "scenario");
1788 apply_context_value(&first, &turn_runtime_context(&scenario.turns[0])).unwrap();
1789 assert_eq!(first.get_context()["precedence"], "turn");
1790
1791 let reset = runner
1792 .build_agent(BuildAgentParams {
1793 agent_path: &dir.join("agent.yaml"),
1794 base_dir: &dir,
1795 attempt_context: &attempt_context,
1796 tool_log,
1797 approval_log,
1798 llm_log,
1799 approval_handler: None,
1800 budget: None,
1801 })
1802 .await
1803 .unwrap();
1804 apply_base_context(&reset, &runner.suite, &dir, scenario, &attempt_context).unwrap();
1805 assert_eq!(reset.get_context()["eval"], first_context["eval"]);
1806 assert_eq!(
1807 reset.get_context()["mock_server"],
1808 first_context["mock_server"]
1809 );
1810 let _ = std::fs::remove_dir_all(dir);
1811 }
1812
1813 #[tokio::test]
1814 async fn streaming_error_retains_partial_content_and_consumes_until_done() {
1815 let chunks = stream::iter(vec![
1816 StreamChunk::Content {
1817 text: "before ".to_string(),
1818 },
1819 StreamChunk::Error {
1820 message: "stream failed".to_string(),
1821 },
1822 StreamChunk::Content {
1823 text: "after".to_string(),
1824 },
1825 StreamChunk::Done {},
1826 ]);
1827 let operation =
1828 consume_stream_response(chunks, TokioInstant::now() + Duration::from_secs(1), 1_000)
1829 .await;
1830 assert_eq!(operation.response_content, "before after");
1831 assert!(operation.response_present);
1832 assert_eq!(operation.runtime_error.as_deref(), Some("stream failed"));
1833 }
1834
1835 #[test]
1836 fn dry_config_check_validates_real_suite_and_agent_without_authorization() {
1837 let dir = std::env::temp_dir().join(format!(
1838 "ai_agents_eval_dry_config_test_{}",
1839 uuid::Uuid::new_v4()
1840 ));
1841 std::fs::create_dir_all(&dir).unwrap();
1842 let suite_path = dir.join("suite.yaml");
1843 std::fs::write(
1844 &suite_path,
1845 r#"
1846name: Dry Config
1847agent: agent.yaml
1848fixtures:
1849 llm:
1850 mode: real
1851scenarios:
1852 - id: live
1853 turns:
1854 - input: hello
1855"#,
1856 )
1857 .unwrap();
1858 std::fs::write(
1859 dir.join("agent.yaml"),
1860 "name: TestAgent\nsystem_prompt: test\n",
1861 )
1862 .unwrap();
1863
1864 EvalRunner::validate_file(&suite_path, None).unwrap();
1865
1866 std::fs::write(
1867 dir.join("agent.yaml"),
1868 "name: TestAgent\nsystem_prompt: test\nmax_iteratons: 3\n",
1869 )
1870 .unwrap();
1871 let error = EvalRunner::validate_file(&suite_path, None).unwrap_err();
1872 assert!(error.to_string().contains("max_iteratons"));
1873 let _ = std::fs::remove_dir_all(dir);
1874 }
1875
1876 #[test]
1877 fn real_and_record_modes_require_explicit_authorization() {
1878 let dir = std::env::temp_dir().join(format!(
1879 "ai_agents_eval_authorization_test_{}",
1880 uuid::Uuid::new_v4()
1881 ));
1882 std::fs::create_dir_all(&dir).unwrap();
1883 let suite_path = dir.join("suite.yaml");
1884 std::fs::write(
1885 &suite_path,
1886 r#"
1887name: Authorization
1888agent: agent.yaml
1889fixtures:
1890 llm:
1891 mode: real
1892scenarios:
1893 - id: authorized
1894 turns:
1895 - input: hello
1896"#,
1897 )
1898 .unwrap();
1899
1900 let error = EvalRunner::from_file(&suite_path, EvalRunnerOptions::default())
1901 .err()
1902 .expect("real mode should require authorization");
1903 assert!(error.to_string().contains("--real-llm"));
1904 assert!(
1905 EvalRunner::from_file(
1906 &suite_path,
1907 EvalRunnerOptions {
1908 llm_mode: Some(LlmFixtureMode::Real),
1909 ..Default::default()
1910 },
1911 )
1912 .is_ok()
1913 );
1914
1915 let record_suite = std::fs::read_to_string(&suite_path)
1916 .unwrap()
1917 .replace("mode: real", "mode: record");
1918 std::fs::write(&suite_path, record_suite).unwrap();
1919 let error = EvalRunner::from_file(&suite_path, EvalRunnerOptions::default())
1920 .err()
1921 .expect("record mode should require authorization");
1922 assert!(error.to_string().contains("--record"));
1923 assert!(
1924 EvalRunner::from_file(
1925 &suite_path,
1926 EvalRunnerOptions {
1927 llm_mode: Some(LlmFixtureMode::Record),
1928 ..Default::default()
1929 },
1930 )
1931 .is_ok()
1932 );
1933 let _ = std::fs::remove_dir_all(dir);
1934 }
1935
1936 #[tokio::test]
1937 async fn zero_selected_scenarios_is_an_error() {
1938 let dir = std::env::temp_dir().join(format!(
1939 "ai_agents_eval_zero_selection_test_{}",
1940 uuid::Uuid::new_v4()
1941 ));
1942 std::fs::create_dir_all(&dir).unwrap();
1943 let suite_path = dir.join("suite.yaml");
1944 std::fs::write(
1945 &suite_path,
1946 r#"
1947name: Selection
1948agent: agent.yaml
1949fixtures:
1950 llm:
1951 mode: mock
1952 responses: [ok]
1953scenarios:
1954 - id: present
1955 turns:
1956 - input: hello
1957"#,
1958 )
1959 .unwrap();
1960 let runner = EvalRunner::from_file(
1961 &suite_path,
1962 EvalRunnerOptions {
1963 ids: vec!["missing".to_string()],
1964 ..Default::default()
1965 },
1966 )
1967 .unwrap();
1968
1969 let error = runner.run().await.unwrap_err();
1970
1971 assert!(error.to_string().contains("matched zero scenarios"));
1972 let _ = std::fs::remove_dir_all(dir);
1973 }
1974
1975 #[test]
1976 fn env_guard_allows_readers_and_excludes_writer() {
1977 use std::sync::{Barrier, mpsc};
1978 use std::thread;
1979
1980 let start = Arc::new(Barrier::new(3));
1981 let (acquired_tx, acquired_rx) = mpsc::channel();
1982 let mut releases = Vec::new();
1983 let mut readers = Vec::new();
1984 for index in 0..2 {
1985 let start = Arc::clone(&start);
1986 let acquired_tx = acquired_tx.clone();
1987 let (release_tx, release_rx) = mpsc::channel();
1988 releases.push(release_tx);
1989 readers.push(thread::spawn(move || {
1990 start.wait();
1991 let guard = EnvGuard::apply(&HashMap::new()).unwrap();
1992 acquired_tx.send(index).unwrap();
1993 release_rx.recv().unwrap();
1994 drop(guard);
1995 }));
1996 }
1997 start.wait();
1998 let first = acquired_rx.recv_timeout(Duration::from_secs(1));
1999 let second = acquired_rx.recv_timeout(Duration::from_secs(1));
2000 for release in releases {
2001 release.send(()).unwrap();
2002 }
2003 for reader in readers {
2004 reader.join().unwrap();
2005 }
2006 assert_ne!(first.unwrap(), second.unwrap());
2007
2008 let key = format!("AI_AGENTS_EVAL_ENV_TEST_{}", uuid::Uuid::new_v4());
2009 unsafe {
2010 std::env::remove_var(&key);
2011 }
2012 let reader = EnvGuard::apply(&HashMap::new()).unwrap();
2013 assert!(matches!(&reader._guard, EnvExclusionGuard::Read { .. }));
2014 let writer_start = Arc::new(Barrier::new(2));
2015 let writer_start_thread = Arc::clone(&writer_start);
2016 let (writer_tx, writer_rx) = mpsc::channel();
2017 let writer_key = key.clone();
2018 let writer = thread::spawn(move || {
2019 writer_start_thread.wait();
2020 let guard = EnvGuard::apply(&HashMap::from([(
2021 writer_key.clone(),
2022 "temporary".to_string(),
2023 )]))
2024 .unwrap();
2025 writer_tx.send(()).unwrap();
2026 assert_eq!(std::env::var(&writer_key).as_deref(), Ok("temporary"));
2027 drop(guard);
2028 });
2029 writer_start.wait();
2030 let writer_was_blocked = writer_rx.recv_timeout(Duration::from_millis(100)).is_err();
2031 drop(reader);
2032 if writer_was_blocked {
2033 writer_rx.recv_timeout(Duration::from_secs(1)).unwrap();
2034 }
2035 writer.join().unwrap();
2036 assert!(writer_was_blocked);
2037 assert!(std::env::var(&key).is_err());
2038 }
2039
2040 fn write_test_agent(dir: &Path) {
2041 std::fs::write(
2042 dir.join("agent.yaml"),
2043 r#"
2044name: TestAgent
2045system_prompt: "You are helpful."
2046llm:
2047 provider: openai
2048 model: gpt-4.1-nano
2049"#,
2050 )
2051 .unwrap();
2052 }
2053
2054 async fn run_test_suite(dir: &Path, name: &str, yaml: &str) -> EvalResult {
2055 let suite_path = dir.join(name);
2056 std::fs::write(&suite_path, yaml).unwrap();
2057 let options = EvalRunnerOptions {
2058 output: dir.join("out"),
2059 ..Default::default()
2060 };
2061 EvalRunner::from_file(&suite_path, options)
2062 .unwrap()
2063 .run()
2064 .await
2065 .unwrap()
2066 }
2067
2068 #[test]
2069 fn runtime_error_expectations_retain_turns_and_control_retries() {
2070 std::thread::Builder::new()
2071 .name("eval-runtime-error-test".to_string())
2072 .stack_size(16 * 1024 * 1024)
2073 .spawn(|| {
2074 let runtime = tokio::runtime::Runtime::new().unwrap();
2075 runtime.block_on(async {
2076 let dir = std::env::temp_dir().join(format!(
2077 "ai_agents_eval_runtime_error_test_{}",
2078 uuid::Uuid::new_v4()
2079 ));
2080 std::fs::create_dir_all(&dir).unwrap();
2081 write_test_agent(&dir);
2082 let errors = run_test_suite(
2083 &dir,
2084 "errors.yaml",
2085 r#"
2086name: Runtime Errors
2087agent: agent.yaml
2088settings:
2089 retries: 1
2090 retry_delay_ms: 0
2091 redact_outputs: false
2092fixtures:
2093 llm:
2094 mode: mock
2095 errors_by_alias:
2096 default: provider exploded
2097scenarios:
2098 - id: expected
2099 turns:
2100 - input: Hello
2101 expect_error: [timeout, provider exploded]
2102 - id: mismatched
2103 turns:
2104 - input: Hello
2105 expect_error: permission denied
2106 - id: unexpected
2107 turns:
2108 - input: Hello
2109"#,
2110 )
2111 .await;
2112
2113 let expected = &errors.scenarios[0];
2114 assert!(expected.status.is_passed());
2115 assert_eq!(expected.attempts.len(), 1);
2116 let expected_turn = &expected.attempts[0].turns[0];
2117 assert!(!expected_turn.response_present);
2118 assert!(expected_turn.runtime_error.is_some());
2119 assert!(
2120 expected_turn
2121 .assertion_results
2122 .iter()
2123 .any(|detail| detail.assertion == "expect_error" && detail.passed)
2124 );
2125
2126 for scenario in &errors.scenarios[1..] {
2127 assert!(scenario.status.is_error());
2128 assert_eq!(scenario.attempts.len(), 2);
2129 assert!(
2130 scenario
2131 .attempts
2132 .iter()
2133 .all(|attempt| attempt.turns.len() == 1)
2134 );
2135 assert!(
2136 scenario
2137 .attempts
2138 .iter()
2139 .all(|attempt| attempt.turns[0].runtime_error.is_some())
2140 );
2141 }
2142
2143 let missing = run_test_suite(
2144 &dir,
2145 "missing.yaml",
2146 r#"
2147name: Missing Runtime Error
2148agent: agent.yaml
2149settings:
2150 retry_delay_ms: 0
2151 redact_outputs: false
2152fixtures:
2153 llm:
2154 mode: mock
2155 responses: [ok]
2156scenarios:
2157 - id: missing
2158 turns:
2159 - input: Hello
2160 expect_error: timeout
2161"#,
2162 )
2163 .await;
2164 let missing = &missing.scenarios[0];
2165 assert!(missing.status.is_failed());
2166 let turn = &missing.attempts[0].turns[0];
2167 assert!(turn.response_present);
2168 assert!(turn.runtime_error.is_none());
2169 assert!(
2170 turn.assertion_results
2171 .iter()
2172 .any(|detail| detail.assertion == "expect_error" && !detail.passed)
2173 );
2174 let _ = std::fs::remove_dir_all(dir);
2175 });
2176 })
2177 .unwrap()
2178 .join()
2179 .unwrap();
2180 }
2181
2182 #[test]
2183 fn scenario_budget_is_shared_across_retries_and_agent_resets() {
2184 std::thread::Builder::new()
2185 .name("eval-budget-lifecycle-test".to_string())
2186 .stack_size(16 * 1024 * 1024)
2187 .spawn(|| {
2188 let runtime = tokio::runtime::Runtime::new().unwrap();
2189 runtime.block_on(async {
2190 let dir = std::env::temp_dir().join(format!(
2191 "ai_agents_eval_budget_lifecycle_test_{}",
2192 uuid::Uuid::new_v4()
2193 ));
2194 std::fs::create_dir_all(&dir).unwrap();
2195 write_test_agent(&dir);
2196
2197 let retried = run_test_suite(
2198 &dir,
2199 "budget-retry.yaml",
2200 r#"
2201name: Retry Budget
2202agent: agent.yaml
2203settings:
2204 retries: 1
2205 retry_delay_ms: 0
2206 redact_outputs: false
2207fixtures:
2208 llm:
2209 mode: mock
2210 responses: [wrong]
2211scenarios:
2212 - id: retry
2213 budget:
2214 max_llm_calls: 1
2215 turns:
2216 - input: Hello
2217 assert:
2218 response_contains: right
2219"#,
2220 )
2221 .await;
2222 let retried = &retried.scenarios[0];
2223 assert!(retried.status.is_error());
2224 assert_eq!(retried.attempts.len(), 2);
2225 assert!(
2226 retried.attempts[1].turns[0]
2227 .runtime_error
2228 .as_ref()
2229 .is_some_and(|error| error.value.contains("max_llm_calls=1"))
2230 );
2231
2232 let reset = run_test_suite(
2233 &dir,
2234 "budget-reset.yaml",
2235 r#"
2236name: Reset Budget
2237agent: agent.yaml
2238settings:
2239 retries: 0
2240 redact_outputs: false
2241fixtures:
2242 llm:
2243 mode: mock
2244 responses: [ok]
2245scenarios:
2246 - id: reset
2247 budget:
2248 max_llm_calls: 1
2249 steps:
2250 - !run
2251 turns:
2252 - input: First
2253 - !reset_agent true
2254 - !run
2255 turns:
2256 - input: Second
2257"#,
2258 )
2259 .await;
2260 let reset = &reset.scenarios[0];
2261 assert!(reset.status.is_error());
2262 assert_eq!(reset.attempts.len(), 1);
2263 assert_eq!(reset.attempts[0].turns.len(), 2);
2264 assert!(
2265 reset.attempts[0].turns[1]
2266 .runtime_error
2267 .as_ref()
2268 .is_some_and(|error| error.value.contains("max_llm_calls=1"))
2269 );
2270 let _ = std::fs::remove_dir_all(dir);
2271 });
2272 })
2273 .unwrap()
2274 .join()
2275 .unwrap();
2276 }
2277
2278 #[test]
2279 fn captures_composed_llm_requests_per_turn_across_reset() {
2280 std::thread::Builder::new()
2281 .name("eval-llm-evidence-test".to_string())
2282 .stack_size(16 * 1024 * 1024)
2283 .spawn(|| {
2284 let runtime = tokio::runtime::Runtime::new().unwrap();
2285 runtime.block_on(async {
2286 let dir = std::env::temp_dir().join(format!(
2287 "ai_agents_eval_llm_evidence_test_{}",
2288 uuid::Uuid::new_v4()
2289 ));
2290 std::fs::create_dir_all(&dir).unwrap();
2291 std::fs::write(
2292 dir.join("agent.yaml"),
2293 r#"
2294name: EvidenceAgent
2295system_prompt: "Base instruction marker."
2296llm:
2297 provider: openai
2298 model: gpt-4.1-nano
2299persona:
2300 identity:
2301 name: Evidence Guide
2302 role: Prompt Inspector
2303reasoning:
2304 mode: cot
2305 output: tagged
2306"#,
2307 )
2308 .unwrap();
2309 let result = run_test_suite(
2310 &dir,
2311 "llm-evidence.yaml",
2312 r#"
2313name: LLM Evidence
2314agent: agent.yaml
2315settings:
2316 redact_outputs: false
2317fixtures:
2318 llm:
2319 mode: mock
2320 responses: [first answer, second answer]
2321scenarios:
2322 - id: composed
2323 steps:
2324 - !run
2325 turns:
2326 - input: first question
2327 assert:
2328 llm_request:
2329 system_contains:
2330 - "You are Evidence Guide, Prompt Inspector."
2331 - "Base instruction marker."
2332 - "Think through this step by step"
2333 - "<instruction>"
2334 user_contains: first question
2335 count: 1
2336 same_request: true
2337 - input: second question
2338 assert:
2339 llm_request:
2340 user_contains: [first question, second question]
2341 assistant_contains: first answer
2342 count: 1
2343 same_request: true
2344 - !reset_agent true
2345 - !run
2346 turns:
2347 - input: after reset
2348 assert:
2349 llm_request:
2350 system_contains: "Base instruction marker."
2351 user_contains: after reset
2352 count: 1
2353 same_request: true
2354"#,
2355 )
2356 .await;
2357
2358 assert_eq!(result.passed, 1);
2359 let turns = &result.scenarios[0].attempts[0].turns;
2360 assert_eq!(turns.len(), 3);
2361 assert!(
2362 turns
2363 .iter()
2364 .all(|turn| turn.evidence.llm_requests.len() == 1)
2365 );
2366 let second_messages = &turns[1].evidence.llm_requests[0].messages;
2367 assert!(second_messages.iter().any(|message| {
2368 message.role == ai_agents_core::Role::Assistant
2369 && message.content.contains("first answer")
2370 }));
2371 let reset_messages = &turns[2].evidence.llm_requests[0].messages;
2372 assert!(
2373 reset_messages
2374 .iter()
2375 .all(|message| !message.content.contains("first question"))
2376 );
2377 assert!(
2378 reset_messages
2379 .iter()
2380 .all(|message| !message.content.contains("first answer"))
2381 );
2382
2383 let serialized = serde_json::to_string(&result).unwrap();
2384 let serialized_value: Value = serde_json::from_str(&serialized).unwrap();
2385 assert!(
2386 serialized_value["scenarios"][0]["attempts"][0]["turns"][0]
2387 .get("evidence")
2388 .is_none()
2389 );
2390 assert!(!serialized.contains("Base instruction marker"));
2391 assert!(!serialized.contains("Evidence Guide"));
2392 assert!(!serialized.contains("Think through this step by step"));
2393 let _ = std::fs::remove_dir_all(dir);
2394 });
2395 })
2396 .unwrap()
2397 .join()
2398 .unwrap();
2399 }
2400
2401 #[test]
2402 fn runner_executes_mocked_suite_and_redacts_outputs() {
2403 std::thread::Builder::new()
2404 .name("eval-runner-test".to_string())
2405 .stack_size(16 * 1024 * 1024)
2406 .spawn(|| {
2407 let runtime = tokio::runtime::Runtime::new().unwrap();
2408 runtime.block_on(async {
2409 let dir = std::env::temp_dir().join(format!(
2410 "ai_agents_eval_runner_test_{}",
2411 uuid::Uuid::new_v4()
2412 ));
2413 std::fs::create_dir_all(&dir).unwrap();
2414 write_test_agent(&dir);
2415 let suite_path = dir.join("suite.yaml");
2416 std::fs::write(
2417 &suite_path,
2418 r#"
2419name: Runner Suite
2420agent: agent.yaml
2421fixtures:
2422 llm:
2423 mode: mock
2424 responses:
2425 - "Hello from mock"
2426scenarios:
2427 - id: smoke
2428 turns:
2429 - input: Hello
2430 assert:
2431 response_contains: "Hello"
2432"#,
2433 )
2434 .unwrap();
2435 let options = EvalRunnerOptions {
2436 output: dir.join("out"),
2437 ..Default::default()
2438 };
2439 let runner = EvalRunner::from_file(&suite_path, options).unwrap();
2440 let result = runner.run().await.unwrap();
2441 assert_eq!(result.passed, 1);
2442 let turn = &result.scenarios[0].attempts[0].turns[0];
2443 assert_eq!(turn.input.value, "[redacted]");
2444 assert_eq!(turn.response.value, "[redacted]");
2445 let json = serde_json::to_string(&result).unwrap();
2446 assert!(!json.contains("Hello from mock"));
2447 let _ = std::fs::remove_dir_all(dir);
2448 });
2449 })
2450 .unwrap()
2451 .join()
2452 .unwrap();
2453 }
2454}