1use std::sync::Arc;
31
32use async_trait::async_trait;
33use futures::StreamExt as _;
34use turnframe_core::case::CaseKey;
35use turnframe_core::flow::WorkflowRegistry;
36use turnframe_core::ids::{CaseId, ConversationId, OriginToken, TurnId};
37use turnframe_core::locale::Locale;
38use turnframe_core::turn::{ActorContext, InteractionResponse, OriginRef, TurnInput};
39use turnframe_runtime::orchestrator::{Orchestrator, error_code};
40use turnframe_store::interaction::InteractionReader;
41
42use crate::assertions::check;
43use crate::config::EvalConfig;
44use crate::control::ControlRun;
45use crate::corpus::{CardReplySpec, EvalItem, ExternalSpec, Suite, TurnSpec};
46use crate::judge::{CriterionOutcome, Judge, JudgeInput};
47use crate::observation::Observation;
48use crate::report::{EvalReport, ItemReport, SampleReport};
49
50#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
52pub struct SampleIndex(pub u32);
53
54impl SampleIndex {
55 #[must_use]
57 pub const fn index(self) -> u32 {
58 self.0
59 }
60
61 #[must_use]
63 pub const fn number(self) -> u32 {
64 self.0 + 1
65 }
66
67 #[must_use]
69 pub const fn is_first(self) -> bool {
70 self.0 == 0
71 }
72}
73
74impl std::fmt::Display for SampleIndex {
75 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
76 write!(f, "{}", self.number())
77 }
78}
79
80pub struct PreparedRun {
85 pub orchestrator: Arc<Orchestrator>,
87 pub workflows: Arc<WorkflowRegistry>,
90 pub actor: ActorContext,
92 pub conversation_id: ConversationId,
94 pub turn_id: TurnId,
96}
97
98impl std::fmt::Debug for PreparedRun {
99 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
100 f.debug_struct("PreparedRun")
101 .field("account", &self.actor.account_id)
102 .field("turn_id", &self.turn_id)
103 .finish_non_exhaustive()
104 }
105}
106
107#[async_trait]
109pub trait EvalHarness: Send + Sync {
110 async fn prepare(
118 &self,
119 item: &EvalItem,
120 sample: SampleIndex,
121 ) -> Result<PreparedRun, HarnessError>;
122}
123
124#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
126#[non_exhaustive]
127pub enum HarnessError {
128 #[error("the harness could not prepare the run: {message}")]
130 Setup {
131 message: String,
133 },
134 #[error("the item names workflow `{workflow}`, which the harness did not register")]
136 UnknownWorkflow {
137 workflow: String,
139 },
140 #[error("case {case} has no blocking card for the item's reply")]
142 NoBlockingCard {
143 case: String,
145 },
146 #[error("the interaction store could not be read: {message}")]
148 Store {
149 message: String,
151 },
152}
153
154impl HarnessError {
155 #[must_use]
157 pub fn setup(message: impl Into<String>) -> Self {
158 Self::Setup {
159 message: message.into(),
160 }
161 }
162}
163
164pub struct Runner {
166 config: EvalConfig,
167 judge: Option<Arc<Judge>>,
168}
169
170impl std::fmt::Debug for Runner {
171 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
172 f.debug_struct("Runner")
173 .field("samples_per_item", &self.config.execution.samples_per_item)
174 .field(
175 "sample_concurrency",
176 &self.config.execution.sample_concurrency,
177 )
178 .field("votes_per_sample", &self.config.judging.votes_per_sample)
179 .field("judge", &self.judge.is_some())
180 .finish()
181 }
182}
183
184impl Runner {
185 #[must_use]
188 pub fn new(config: EvalConfig) -> Self {
189 Self {
190 config,
191 judge: None,
192 }
193 }
194
195 #[must_use]
197 pub fn with_judge(mut self, judge: Arc<Judge>) -> Self {
198 self.judge = Some(judge);
199 self
200 }
201
202 #[must_use]
204 pub const fn config(&self) -> &EvalConfig {
205 &self.config
206 }
207
208 pub async fn run(&self, suite: &Suite, harness: &dyn EvalHarness) -> EvalReport {
210 let mut items = Vec::new();
211 let mut failed_samples = 0_u32;
212 for item in suite.select(&self.config.selection) {
213 let report = self.run_item(item, harness).await;
214 failed_samples +=
215 u32::try_from(report.total_samples() - report.samples_passed()).unwrap_or(u32::MAX);
216 items.push(report);
217 if self
218 .config
219 .execution
220 .stop_after_failures
221 .is_some_and(|budget| failed_samples >= budget)
222 {
223 break;
224 }
225 }
226 EvalReport::new(
227 suite.name.clone(),
228 chrono::Utc::now(),
229 self.config.clone(),
230 items,
231 )
232 }
233
234 pub async fn run_control(&self, suite: &Suite, harness: &dyn EvalHarness) -> ControlRun {
260 let first = self.run(suite, harness).await;
261 let second = self.run(suite, harness).await;
262 ControlRun::new(first, second)
263 }
264
265 pub async fn run_item(&self, item: &EvalItem, harness: &dyn EvalHarness) -> ItemReport {
275 let count = self.config.execution.samples_per_item.max(1);
276 let in_flight =
277 usize::try_from(self.config.execution.sample_concurrency.max(1)).unwrap_or(usize::MAX);
278 let samples = futures::stream::iter(
279 (0..count).map(|index| self.run_sample(item, harness, SampleIndex(index))),
280 )
281 .buffered(in_flight)
285 .collect::<Vec<_>>()
286 .await;
287 ItemReport {
288 id: item.id.clone(),
289 name: item.name.clone(),
290 tags: item.tags.clone(),
291 fingerprint: item.fingerprint(),
296 samples,
297 }
298 }
299
300 pub async fn run_sample(
302 &self,
303 item: &EvalItem,
304 harness: &dyn EvalHarness,
305 sample: SampleIndex,
306 ) -> SampleReport {
307 let prepared = match harness.prepare(item, sample).await {
308 Ok(prepared) => prepared,
309 Err(error) => return unmeasured(sample, &error.to_string()),
310 };
311 let mut earlier = Vec::new();
314 for spec in &item.before {
315 if let Some(external) = &spec.external {
316 if let Err(error) = outside(&prepared, external).await {
317 return unmeasured(sample, &error);
318 }
319 continue;
320 }
321 let turn_id = TurnId::new();
322 let input = match build_input(&prepared, spec, turn_id).await {
323 Ok(input) => input,
324 Err(error) => return unmeasured(sample, &error.to_string()),
325 };
326 let _ = prepared.orchestrator.handle_turn(input).await;
327 earlier.push(turn_id);
328 }
329 let input = match build_input(&prepared, &item.turn, prepared.turn_id).await {
330 Ok(input) => input,
331 Err(error) => return unmeasured(sample, &error.to_string()),
332 };
333
334 let outcome = prepared.orchestrator.handle_turn(input).await;
335 let abandoned = outcome.as_ref().map_or(true, |turn| turn.blocks.is_empty());
338 let observed = Observation::collect_bounded(
339 prepared.orchestrator.stores(),
340 prepared.workflows.as_ref(),
341 &prepared.actor.account_id,
342 prepared.turn_id,
343 &item.setup.cases,
344 outcome.as_ref().map_err(error_code),
345 self.config.execution.max_observed_events,
346 )
347 .await
348 .with_conversation_cases(
349 prepared.orchestrator.stores(),
350 prepared.workflows.as_ref(),
351 &prepared.actor.account_id,
352 &earlier,
353 )
354 .await;
355
356 let failures = check(&item.expect, &observed);
357 let judge = self.judge_sample(item, sample, &observed).await;
358
359 SampleReport {
360 sample: sample.number(),
361 failures,
362 harness_error: None,
363 signature: observed.signature(),
364 judge,
365 acts_proposed: observed.acts.len(),
366 acts_refused: observed
367 .acts
368 .iter()
369 .filter(|act| act.outcome.as_deref() == Some("rejected"))
370 .count(),
371 commands_journaled: observed.commands.len(),
372 provider_failures: observed.provider_failures,
373 cards_created: observed.cards_created,
374 abandoned,
375 discarded_answers: observed
376 .discard_codes()
377 .into_iter()
378 .map(str::to_owned)
379 .collect(),
380 answer: if self.config.execution.record_answers {
381 observed.answer.clone()
382 } else {
383 String::new()
384 },
385 tasks: item
386 .expect
387 .understanding
388 .as_ref()
389 .zip(item.turn.text.as_deref())
390 .map(|(expected, text)| expected.score(text, &observed))
391 .unwrap_or_default(),
392 }
393 }
394
395 async fn judge_sample(
403 &self,
404 item: &EvalItem,
405 sample: SampleIndex,
406 observed: &Observation,
407 ) -> Vec<CriterionOutcome> {
408 let Some(judge) = self.judge.as_ref() else {
409 return Vec::new();
410 };
411 if item.judge.is_empty() {
412 return Vec::new();
413 }
414 if !self.config.judging.judge_every_sample && !sample.is_first() {
415 return Vec::new();
416 }
417 let input = JudgeInput::new(question_of(&item.turn), observed.answer.clone())
418 .with_committed(observed.events.clone());
419 if input.is_empty() {
420 return item
421 .judge
422 .iter()
423 .map(|criterion| CriterionOutcome::empty(*criterion))
424 .collect();
425 }
426 let mut outcomes = Vec::new();
427 for criterion in &item.judge {
428 outcomes.push(
429 judge
430 .poll(*criterion, &input, self.config.judging.votes_per_sample)
431 .await,
432 );
433 }
434 outcomes
435 }
436}
437
438async fn outside(prepared: &PreparedRun, external: &ExternalSpec) -> Result<(), String> {
442 use turnframe_core::case::CaseRef;
443 use turnframe_core::command::{
444 AtomicityScope, CommandBatch, CommandEnvelope, CommandOrigin, IdempotencyKey,
445 };
446 use turnframe_core::ids::{BatchId, CommandId};
447 use turnframe_core::understanding::{ActId, UnitId};
448
449 let registered = prepared
450 .workflows
451 .require(&external.workflow)
452 .map_err(|error| error.to_string())?;
453 let account = &prepared.actor.account_id;
454 let loaded = registered
455 .executor
456 .load(account, &external.case_id)
457 .await
458 .map_err(|error| error.to_string())?;
459 let case_ref = CaseRef::new(
460 external.workflow.clone(),
461 external.case_id.clone(),
462 loaded.revision,
463 );
464 let turn_id = TurnId::new();
465 let origin = CommandOrigin::ExternalCallback {
466 callback_id: "eval.external".to_owned(),
467 signature_verified: true,
468 };
469 let idempotency_key =
470 IdempotencyKey::derive(account, &turn_id, &case_ref, &origin, &external.command)
471 .map_err(|error| error.to_string())?;
472 let batch = CommandBatch {
473 batch_id: BatchId::derive(&turn_id, &case_ref.key(), &AtomicityScope::PerCase),
474 scope: AtomicityScope::PerCase,
475 envelopes: vec![CommandEnvelope {
476 command_id: CommandId::derive(&turn_id, ActId::new(UnitId(1), 1), 0),
477 turn_id,
478 actor: ActorContext::new(account.clone(), "external"),
479 case_ref,
480 idempotency_key,
481 origin,
482 command: external.command.clone(),
483 }],
484 };
485 registered
486 .executor
487 .execute(batch)
488 .await
489 .map(|_| ())
490 .map_err(|error| error.to_string())
491}
492
493fn unmeasured(sample: SampleIndex, message: &str) -> SampleReport {
494 SampleReport {
495 sample: sample.number(),
496 failures: Vec::new(),
497 harness_error: Some(message.to_owned()),
498 signature: format!("harness_error={message}"),
499 judge: Vec::new(),
500 acts_proposed: 0,
501 acts_refused: 0,
502 commands_journaled: 0,
503 provider_failures: 0,
504 cards_created: 0,
505 abandoned: true,
506 discarded_answers: Vec::new(),
507 answer: String::new(),
508 tasks: Default::default(),
509 }
510}
511
512fn question_of(spec: &TurnSpec) -> String {
514 match (&spec.text, &spec.reply) {
515 (Some(text), _) => text.clone(),
516 (None, Some(reply)) => format!(
517 "(the user chose `{}` on {}/{})",
518 reply.option,
519 reply.workflow,
520 reply
521 .case_id
522 .as_ref()
523 .map_or("the open card", CaseId::as_str)
524 ),
525 (None, None) => String::new(),
526 }
527}
528
529async fn build_input(
531 prepared: &PreparedRun,
532 spec: &TurnSpec,
533 turn_id: TurnId,
534) -> Result<TurnInput, HarnessError> {
535 let mut actor = prepared.actor.clone();
536 if let Some(user_id) = &spec.user_id {
537 actor.user_id = turnframe_core::ids::UserId::new(user_id.clone());
538 }
539 let interaction_response = match &spec.reply {
540 Some(reply) => Some(resolve_card(prepared, reply).await?),
541 None => None,
542 };
543 Ok(TurnInput {
544 turn_id,
545 conversation_id: prepared.conversation_id,
546 actor,
547 text: spec.text.clone(),
548 interaction_response,
549 attachments: Vec::new(),
550 origin: spec.origin.as_ref().map(|origin| OriginRef {
551 origin_token: OriginToken::from(origin.token.as_str()),
552 signature: None,
553 surface: origin.surface.clone(),
554 }),
555 locale: spec.locale.clone().unwrap_or_else(|| Locale::from("en")),
556 effort: None,
557 })
558}
559
560async fn resolve_card(
568 prepared: &PreparedRun,
569 reply: &CardReplySpec,
570) -> Result<InteractionResponse, HarnessError> {
571 let store = prepared.orchestrator.stores().interactions();
572 let open = match &reply.case_id {
573 Some(case_id) => {
574 let case = CaseKey::new(reply.workflow.clone(), case_id.clone());
575 InteractionReader::list_open_for_case(store.as_ref(), &prepared.actor.account_id, &case)
576 .await
577 }
578 None => {
579 InteractionReader::list_open_for_conversation(
580 store.as_ref(),
581 &prepared.actor.account_id,
582 &prepared.conversation_id,
583 )
584 .await
585 }
586 }
587 .map_err(|error| HarnessError::Store {
588 message: error.to_string(),
589 })?;
590 let card = open
591 .into_iter()
592 .find(|card| card.blocking && card.case_ref.workflow == reply.workflow)
593 .ok_or_else(|| HarnessError::NoBlockingCard {
594 case: format!(
595 "{}/{}",
596 reply.workflow,
597 reply.case_id.as_ref().map_or("*", CaseId::as_str)
598 ),
599 })?;
600 Ok(InteractionResponse {
601 interaction_id: card.id,
602 option_id: reply.option.clone(),
603 expected_case_revision: card.case_ref.expected_revision,
604 freeform_input: reply.freeform.clone(),
605 })
606}