Skip to main content

turnframe_eval/
runner.rs

1//! Running a corpus against a real orchestrator (spec §27.6).
2//!
3//! # What the runner does, and what it refuses to do
4//!
5//! For every selected item it asks the harness for a fresh world, builds the
6//! turn the item describes, runs it through a real [`Orchestrator`], reads back
7//! what happened from the stores, and checks the item's deterministic
8//! expectations. Then it does it again, `samples_per_item` times.
9//!
10//! It never turns a failing item into an error. A scenario that passes seven
11//! times out of ten is a *measurement* — the most valuable one in the whole
12//! harness, because it is the one a single run cannot see. So the sample
13//! records its failures, the item records its variance, and the run continues.
14//! The only thing that stops a run early is an explicit
15//! [`stop_after_failures`](crate::config::ExecutionConfig::stop_after_failures).
16//!
17//! # The seam
18//!
19//! [`EvalHarness`] is the one thing an application implements. It owns the
20//! domain: it knows how to turn the item's seeded JSON state into a trip or a
21//! traveler, which providers to configure, and how much autonomy to grant.
22//! The runner owns everything that must not vary between applications — the
23//! order, the sampling, the assertions and the report.
24//!
25//! The `sample` index is handed to [`EvalHarness::prepare`] on purpose: a
26//! harness that wants to exercise model variance without a network can script a
27//! different answer per sample, and a harness talking to a real endpoint can
28//! simply ignore it.
29
30use 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/// Which execution of an item this is, zero-based.
51#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
52pub struct SampleIndex(pub u32);
53
54impl SampleIndex {
55    /// The zero-based index.
56    #[must_use]
57    pub const fn index(self) -> u32 {
58        self.0
59    }
60
61    /// The 1-based number shown in reports.
62    #[must_use]
63    pub const fn number(self) -> u32 {
64        self.0 + 1
65    }
66
67    /// Returns `true` for the first sample of an item.
68    #[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
80/// One item's world, freshly built for one sample.
81///
82/// Every sample gets its own: two samples that shared a store would be
83/// measuring the second one against the first one's effects.
84pub struct PreparedRun {
85    /// The runtime under test.
86    pub orchestrator: Arc<Orchestrator>,
87    /// The workflows it was built with, so the runner can read a case's
88    /// revision back without knowing the domain types.
89    pub workflows: Arc<WorkflowRegistry>,
90    /// Who takes the turn.
91    pub actor: ActorContext,
92    /// The conversation the turn belongs to. The harness must have created it.
93    pub conversation_id: ConversationId,
94    /// The identifier the turn will carry.
95    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/// Builds one world per sample.
108#[async_trait]
109pub trait EvalHarness: Send + Sync {
110    /// Seeds the item's starting state and returns a runtime to run it against.
111    ///
112    /// # Errors
113    ///
114    /// [`HarnessError`] when the world could not be built — an unknown
115    /// workflow, an unparseable seeded state, a provider that could not be
116    /// configured. The sample is recorded as unmeasured rather than as failing.
117    async fn prepare(
118        &self,
119        item: &EvalItem,
120        sample: SampleIndex,
121    ) -> Result<PreparedRun, HarnessError>;
122}
123
124/// Why a sample could not be set up or started.
125#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
126#[non_exhaustive]
127pub enum HarnessError {
128    /// The world could not be built.
129    #[error("the harness could not prepare the run: {message}")]
130    Setup {
131        /// What went wrong.
132        message: String,
133    },
134    /// The item names a workflow the harness did not register.
135    #[error("the item names workflow `{workflow}`, which the harness did not register")]
136    UnknownWorkflow {
137        /// The workflow key.
138        workflow: String,
139    },
140    /// The item answers a card, and the case has none.
141    #[error("case {case} has no blocking card for the item's reply")]
142    NoBlockingCard {
143        /// The case, as `workflow/case_id`.
144        case: String,
145    },
146    /// The card store could not be read.
147    #[error("the interaction store could not be read: {message}")]
148    Store {
149        /// What the store said.
150        message: String,
151    },
152}
153
154impl HarnessError {
155    /// A setup failure with a message.
156    #[must_use]
157    pub fn setup(message: impl Into<String>) -> Self {
158        Self::Setup {
159            message: message.into(),
160        }
161    }
162}
163
164/// Runs a corpus.
165pub 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    /// A runner with this configuration and no judge, which is all a corpus of
186    /// deterministic assertions needs.
187    #[must_use]
188    pub fn new(config: EvalConfig) -> Self {
189        Self {
190            config,
191            judge: None,
192        }
193    }
194
195    /// Attaches a judge, used only for the criteria an item asks for.
196    #[must_use]
197    pub fn with_judge(mut self, judge: Arc<Judge>) -> Self {
198        self.judge = Some(judge);
199        self
200    }
201
202    /// The configuration in force.
203    #[must_use]
204    pub const fn config(&self) -> &EvalConfig {
205        &self.config
206    }
207
208    /// Runs every selected item of `suite`.
209    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    /// Runs the same suite twice, against the same harness, and hands back both
235    /// reports as a [`ControlRun`].
236    ///
237    /// This is how the noise floor stops being an assumption. Nothing changes
238    /// between the two passes — same corpus, same harness, same code — so
239    /// whatever difference [`ControlRun::noise_floor`] finds is the harness's
240    /// own variation, and a later comparison can say whether a difference
241    /// exceeds it instead of leaving a reader to guess.
242    ///
243    /// It costs exactly twice a run, which is the honest price of knowing
244    /// whether the first one meant anything.
245    ///
246    /// ```no_run
247    /// # use std::sync::Arc;
248    /// # use turnframe_eval::config::EvalConfig;
249    /// # use turnframe_eval::corpus::Suite;
250    /// # use turnframe_eval::runner::{EvalHarness, Runner};
251    /// # async fn run(suite: &Suite, harness: &dyn EvalHarness) {
252    /// let control = Runner::new(EvalConfig::default().with_samples_per_item(10))
253    ///     .run_control(suite, harness)
254    ///     .await;
255    /// let floor = control.noise_floor();
256    /// println!("{}", floor.summary());
257    /// # }
258    /// ```
259    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    /// Runs one item, `samples_per_item` times, at most
266    /// [`sample_concurrency`](crate::config::ExecutionConfig::sample_concurrency)
267    /// of them at once.
268    ///
269    /// The default concurrency of one reproduces the strictly serial behaviour
270    /// exactly. Above one the samples *execute* interleaved — which is the
271    /// point, when each of them is a call to a real endpoint — but the report
272    /// does not move: the samples come back in index order either way, and the
273    /// same set of results comes back whatever the setting is.
274    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`, not `buffer_unordered`: the samples run together and are
282        // still handed back in index order, so a report never depends on which
283        // endpoint answered first.
284        .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            // Recorded on every report, always: a comparison that cannot see
292            // what the item contained cannot tell an intended projection change
293            // from an unrelated edit, and would report an unpaired difference
294            // as a regression.
295            fingerprint: item.fingerprint(),
296            samples,
297        }
298    }
299
300    /// Runs one sample of one item.
301    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        // The turns before the observed one play out as a person would take them: one
312        // that fails leaves the conversation where it stands, and the next is taken.
313        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        // "Abandoned" is a turn that produced nothing a person could act on:
336        // it failed, or it came back with no blocks at all (spec §26.3).
337        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    /// Polls the judge for the criteria this item asks for.
396    ///
397    /// Nothing from [`Observation`] reaches the judge except the model-authored
398    /// text and — for
399    /// [`OperationalClaimIntegrity`](crate::judge::JudgeCriterion::OperationalClaimIntegrity)
400    /// — the event types the turn committed, which are the ledger's answer and
401    /// not something the judge is asked to establish.
402    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
438/// The sample record of a run that never happened.
439/// Applies a change from outside the conversation as the record's own system would: the
440/// domain's command, at the revision the record is at, from an external origin.
441async 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
512/// What the person asked, for the judge's benefit.
513fn 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
529/// Turns the item's description of a turn into a real [`TurnInput`].
530async 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
560/// Finds the blocking card an item's reply answers.
561///
562/// An item file cannot name an [`InteractionId`](turnframe_core::ids::InteractionId):
563/// the identifier is minted while the corpus is running. What it names instead
564/// is the case, which is how a person would describe the click anyway — and the
565/// revision the reply carries is read from the card itself, so the corpus never
566/// has to know which revision the previous turn left behind.
567async 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}