Skip to main content

turnframe_core/
replay.rs

1//! Replay records and turn phases (spec §23.1, I20).
2//!
3//! A [`ReplayRecord`] holds enough to reconstruct why a turn produced its
4//! commands and response: versions, revisions, the normalized plan and its
5//! hash, target resolutions, policy decisions, command outcomes, event ids,
6//! block ids, the outbox rows and reconciliation handles of its external
7//! effects, and every provider attempt.
8
9use chrono::{DateTime, Utc};
10use serde::{Deserialize, Serialize};
11
12use crate::case::CaseRef;
13use crate::command::{CommandOrigin, IdempotencyKey};
14use crate::error::RejectionCode;
15use crate::hash::Digest;
16use crate::ids::{
17    AccountId, AttemptId, BlockId, CaseRevision, ConversationId, EventId, InteractionId, ModelKey,
18    OutboxId, ProviderKey, TurnId, WorkflowKey, WorkflowVersion,
19};
20use crate::policy::PolicyDecision;
21use crate::prompt::PromptRef;
22use crate::reduce::CommandRef;
23use crate::target::TargetResolution;
24use crate::understanding::Understanding;
25
26/// Persisted phase marker of a turn, for crash recovery (spec §23.1).
27///
28/// The pipeline grows phases, so downstream matches need a wildcard arm.
29#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
30#[serde(rename_all = "snake_case")]
31#[non_exhaustive]
32pub enum TurnPhase {
33    /// Input accepted and persisted.
34    Received,
35    /// A plan was accepted.
36    Interpreted,
37    /// A reduction plan exists.
38    Reduced,
39    /// Commands are executing.
40    Executing,
41    /// Commands committed.
42    Committed,
43    /// Response blocks composed and persisted.
44    Composed,
45    /// Response delivered.
46    Delivered,
47    /// The turn failed.
48    Failed,
49}
50
51impl TurnPhase {
52    /// Returns `true` for `Delivered` and `Failed`.
53    #[must_use]
54    pub fn is_terminal(self) -> bool {
55        matches!(self, Self::Delivered | Self::Failed)
56    }
57
58    /// Returns `true` once commands may have executed; recovery must resume by
59    /// idempotency key instead of re-interpreting (spec §23.1).
60    #[must_use]
61    pub fn effects_may_exist(self) -> bool {
62        matches!(
63            self,
64            Self::Executing | Self::Committed | Self::Composed | Self::Delivered
65        )
66    }
67}
68
69/// Workflow version in force for a turn.
70#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
71pub struct WorkflowVersionRecord {
72    /// The workflow.
73    pub key: WorkflowKey,
74    /// Its version.
75    pub version: WorkflowVersion,
76}
77
78/// Target resolution of one act.
79#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
80pub struct TargetResolutionRecord {
81    /// The act.
82    pub act: crate::understanding::ActId,
83    /// How it resolved.
84    pub resolution: TargetResolution,
85}
86
87/// What happened to one command.
88///
89/// New outcomes appear as execution learns to say more, so downstream matches
90/// need a wildcard arm.
91#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
92#[serde(tag = "kind", rename_all = "snake_case")]
93#[non_exhaustive]
94pub enum CommandOutcome {
95    /// Committed.
96    Committed {
97        /// Revision after the commit.
98        new_revision: CaseRevision,
99        /// Events produced.
100        event_ids: Vec<EventId>,
101    },
102    /// The journal already had the key; the original outcome was returned.
103    IdempotentReplay,
104    /// The expected revision was stale.
105    RevisionConflict {
106        /// Revision found.
107        current_revision: CaseRevision,
108    },
109    /// The domain rejected.
110    Rejected {
111        /// Code.
112        code: RejectionCode,
113    },
114    /// Execution failed.
115    Failed {
116        /// Stable code.
117        code: String,
118    },
119    /// An external effect has an unknown outcome.
120    OutcomeUnknown {
121        /// Attempt id for reconciliation.
122        attempt_id: AttemptId,
123    },
124    /// Waiting for a confirmation interaction.
125    AwaitingConfirmation {
126        /// The interaction.
127        interaction_id: InteractionId,
128    },
129}
130
131/// One command and its outcome.
132#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
133pub struct CommandOutcomeRecord {
134    /// The command.
135    pub command_ref: CommandRef,
136    /// Its idempotency key.
137    pub idempotency_key: IdempotencyKey,
138    /// Case and expected revision.
139    pub case_ref: CaseRef,
140    /// What authorized the command.
141    ///
142    /// Recorded so the audit answers "which interaction authorized this
143    /// command" from the record alone, without reloading the batches that
144    /// produced it (spec §26.4, I20). Absent on records written before the
145    /// field existed, which is why it is optional rather than required.
146    #[serde(default, skip_serializing_if = "Option::is_none")]
147    pub origin: Option<CommandOrigin>,
148    /// The outcome.
149    pub outcome: CommandOutcome,
150}
151
152/// How a provider attempt ended.
153///
154/// Routing strategies grow, so downstream matches need a wildcard arm.
155#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
156#[serde(tag = "kind", rename_all = "snake_case")]
157#[non_exhaustive]
158pub enum ProviderAttemptOutcome {
159    /// Usable output.
160    Succeeded,
161    /// Failed with a normalized code.
162    Failed {
163        /// Stable code.
164        code: String,
165    },
166    /// Failed and routing moved to another candidate.
167    FellBack {
168        /// Stable code.
169        code: String,
170    },
171    /// Cancelled.
172    Cancelled,
173}
174
175/// One model call (spec §20.7: "record every provider attempt").
176///
177/// # Why `attempt` is a number and [`AttemptId`] is an identifier
178///
179/// A model call leaves nothing behind outside the process. When one fails or
180/// times out, the only thing anyone needs to know is *which try it was* inside
181/// a stage this record already identifies — the turn, the `purpose`, the
182/// provider and the model — so a position in that sequence says everything, and
183/// nothing else ever refers to it.
184///
185/// An external effect attempt is the opposite: it may have happened even though
186/// the answer never arrived (spec §16.5), and settling it means naming that
187/// exact attempt to a remote system. That is what [`AttemptId`] is for, why
188/// [`CommandOutcome::OutcomeUnknown`] carries one, and why
189/// [`ReplayRecord::reconciliation_attempt_ids`] lists them. The asymmetry is
190/// the difference between counting retries and naming an effect, not an
191/// oversight.
192///
193/// Note that this record is compared with [`PartialEq`] only: `temperature` is
194/// a float, so `Eq` would be a promise about `NaN` that the type cannot keep.
195#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
196pub struct ProviderAttemptRecord {
197    /// 1-based attempt number within the stage. A plain ordinal on purpose;
198    /// see the type documentation.
199    pub attempt: u32,
200    /// Stage purpose (e.g. `"extract"`).
201    pub purpose: String,
202    /// Provider key.
203    pub provider_key: ProviderKey,
204    /// Model key.
205    pub model_key: ModelKey,
206    /// Stable request id sent to the provider.
207    pub request_id: String,
208    /// Prompt or template version, when known.
209    #[serde(default, skip_serializing_if = "Option::is_none")]
210    pub prompt_version: Option<String>,
211    /// The exact prompt text this call ran under, when a prompt source supplied
212    /// it (roadmap: "prompt management, with an optional Langfuse prompt
213    /// source").
214    ///
215    /// `None` for a stage whose instructions are compiled into the library —
216    /// which is every stage until an application configures a
217    /// `turnframe-prompt` source — and also for a stage whose configured source
218    /// failed and fell back to those built-in instructions. That is deliberate:
219    /// the absence of a reference is the audit signal that the text was not the
220    /// text the source was asked for.
221    ///
222    /// Unlike [`prompt_version`](Self::prompt_version), which is a bare label,
223    /// this carries the digest of the text, so the record can be falsified
224    /// rather than merely believed.
225    #[serde(default, skip_serializing_if = "Option::is_none")]
226    pub prompt_ref: Option<PromptRef>,
227    /// Outcome.
228    pub outcome: ProviderAttemptOutcome,
229    /// Latency in milliseconds.
230    #[serde(default, skip_serializing_if = "Option::is_none")]
231    pub latency_ms: Option<u64>,
232    /// Input tokens, when reported.
233    #[serde(default, skip_serializing_if = "Option::is_none")]
234    pub input_tokens: Option<u64>,
235    /// Output tokens, when reported.
236    #[serde(default, skip_serializing_if = "Option::is_none")]
237    pub output_tokens: Option<u64>,
238    /// Sampling temperature the request carried, when the stage set one.
239    ///
240    /// `None` means the call left the provider's default in place, which is not
241    /// the same as `Some(0.0)`. An observability bridge reports it as
242    /// `gen_ai.request.temperature`; kept at the request's own `f32` precision
243    /// so the audit value is the value that was sent, not a widened copy of it.
244    #[serde(default, skip_serializing_if = "Option::is_none")]
245    pub temperature: Option<f32>,
246    /// Finish reasons the provider reported for the response, in the order it
247    /// reported them, verbatim.
248    ///
249    /// Empty when the attempt produced none — a transport failure, a
250    /// cancellation, or a provider that does not report them. Providers spell
251    /// these differently (`"stop"`, `"length"`, `"max_tokens"`, ...) and the
252    /// strings are kept as received: normalizing them here would erase the
253    /// distinction an audit is being read for. An observability bridge reports
254    /// them as `gen_ai.response.finish_reasons`.
255    #[serde(default, skip_serializing_if = "Vec::is_empty")]
256    pub finish_reasons: Vec<String>,
257}
258
259/// One answer a model produced and a check refused whole.
260///
261/// A call can succeed and the turn still lose its answer, to a pointer outside the
262/// message or a document of the wrong shape. The repair round usually recovers and
263/// nobody counts the cost; when it does not, the turn does nothing and asks nothing,
264/// and this is the only place that says why. Derived from [`ReplayRecord::tasks`].
265#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
266pub struct DiscardedAnswer {
267    /// The task's purpose, such as `"extract"` or `"answer"`.
268    pub purpose: String,
269    /// Which refusal of its task this was, from 1.
270    pub round: u32,
271    /// Stable code of the failed check: the field to group by.
272    pub code: String,
273    /// The refusal in words, as the repair round was told it.
274    pub reason: String,
275}
276
277/// Everything needed to replay a turn (I20).
278///
279/// Spec §26.1 asks that one turn be traceable end to end, which includes the
280/// external effects it started: the outbox rows it enqueued and the attempts
281/// whose outcome is not yet known. Those live in [`outbox_ids`] and
282/// [`reconciliation_attempt_ids`], so an auditor holding nothing but this
283/// record can ask the outbox what became of the turn's side effects instead of
284/// inferring it from the transcript.
285///
286/// Compared with [`PartialEq`] only, because [`ProviderAttemptRecord`] carries
287/// a float; see that type.
288///
289/// [`outbox_ids`]: ReplayRecord::outbox_ids
290/// [`reconciliation_attempt_ids`]: ReplayRecord::reconciliation_attempt_ids
291#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
292pub struct ReplayRecord {
293    /// The turn.
294    pub turn_id: TurnId,
295    /// The conversation.
296    pub conversation_id: ConversationId,
297    /// The tenant.
298    pub account_id: AccountId,
299    /// Last persisted phase.
300    pub phase: TurnPhase,
301    /// Workflow versions in force.
302    pub workflow_versions: Vec<WorkflowVersionRecord>,
303    /// Cases and revisions loaded at the start of the turn.
304    pub loaded_cases: Vec<CaseRef>,
305    /// What the turn was understood to say: the reducer's input.
306    #[serde(default, skip_serializing_if = "Option::is_none")]
307    pub understanding: Option<Understanding>,
308    /// Hash of the understanding.
309    #[serde(default, skip_serializing_if = "Option::is_none")]
310    pub plan_hash: Option<Digest>,
311    /// What the reduction decided about each act, by index, as
312    /// [`crate::reduce::PlannedActResult::name`].
313    ///
314    /// The hash below identifies a plan; this says what happened in it. A
315    /// measurement that wants to know how often the structure refused a
316    /// reading cannot get there from the command count: an act that changes
317    /// nothing and one that is waiting for a confirmation both journal zero
318    /// commands and neither was refused, while a plan holding one refusal
319    /// beside one write journals a command and hides the refusal entirely.
320    #[serde(default, skip_serializing_if = "Vec::is_empty")]
321    pub act_outcomes: Vec<String>,
322    /// Hash of the reduction plan.
323    #[serde(default, skip_serializing_if = "Option::is_none")]
324    pub reduction_plan_hash: Option<Digest>,
325    /// Target resolutions per act.
326    pub target_resolutions: Vec<TargetResolutionRecord>,
327    /// Policy decisions per command.
328    pub policy_decisions: Vec<PolicyDecision>,
329    /// Command outcomes.
330    pub command_outcomes: Vec<CommandOutcomeRecord>,
331    /// Interactions created by the turn.
332    pub interactions_created: Vec<InteractionId>,
333    /// Events committed by the turn.
334    pub event_ids: Vec<EventId>,
335    /// Block ids of the assistant turn, in order.
336    pub response_block_ids: Vec<BlockId>,
337    /// Outbox rows the turn enqueued, in the order they were enqueued
338    /// (spec §16.4, §26.1).
339    ///
340    /// Each identifies a row in the dispatch queue, so the record names the
341    /// turn's external effects even before any of them has been dispatched.
342    /// Empty for a turn with no external effect, which is the common case; the
343    /// field defaults on deserialization, so records written before it existed
344    /// still load.
345    #[serde(default, skip_serializing_if = "Vec::is_empty")]
346    pub outbox_ids: Vec<OutboxId>,
347    /// Attempts at external effects that the turn left unsettled, in the order
348    /// they were made (spec §16.5, I15).
349    ///
350    /// An attempt lands here when the request was transmitted and the outcome
351    /// never came back: it may or may not have taken effect, so it is neither
352    /// a success to claim nor a failure to retry blindly, and a poller or a
353    /// callback has to settle it against the remote system. Reading them
354    /// together with the ones implied by
355    /// [`CommandOutcome::OutcomeUnknown`] is
356    /// what [`pending_reconciliations`](ReplayRecord::pending_reconciliations)
357    /// is for. Defaults on deserialization.
358    #[serde(default, skip_serializing_if = "Vec::is_empty")]
359    pub reconciliation_attempt_ids: Vec<AttemptId>,
360    /// Every provider attempt.
361    pub provider_attempts: Vec<ProviderAttemptRecord>,
362    /// Every prompt a configured prompt source supplied for this turn, in the
363    /// order the stages asked for them and without repeats.
364    ///
365    /// [`ProviderAttemptRecord::prompt_ref`] answers "which prompt produced
366    /// *this call*"; this list answers "which prompts were in force for the
367    /// turn at all", which is the question an audit of a released prompt
368    /// version asks. It stays empty for an application that has configured no
369    /// prompt source, which is the recommended default, and it defaults on
370    /// deserialization so records written before it existed still load.
371    #[serde(default, skip_serializing_if = "Vec::is_empty")]
372    pub prompt_refs: Vec<PromptRef>,
373    /// Every model task the turn ran, repairs, votes and escalations included.
374    #[serde(default, skip_serializing_if = "Vec::is_empty")]
375    pub tasks: Vec<TaskRecord>,
376    /// What the turn's model calls spent, and the bound that stopped them if one did.
377    #[serde(default, skip_serializing_if = "Option::is_none")]
378    pub budget: Option<BudgetReport>,
379    /// The effort the turn ran at.
380    #[serde(default)]
381    pub effort: crate::effort::Effort,
382    /// When the record was last written.
383    pub recorded_at: DateTime<Utc>,
384}
385
386impl ReplayRecord {
387    /// A record for a turn that was just received.
388    #[must_use]
389    pub fn received(
390        turn_id: TurnId,
391        conversation_id: ConversationId,
392        account_id: AccountId,
393        now: DateTime<Utc>,
394    ) -> Self {
395        Self {
396            turn_id,
397            conversation_id,
398            account_id,
399            phase: TurnPhase::Received,
400            workflow_versions: Vec::new(),
401            loaded_cases: Vec::new(),
402            understanding: None,
403            plan_hash: None,
404            act_outcomes: Vec::new(),
405            reduction_plan_hash: None,
406            target_resolutions: Vec::new(),
407            policy_decisions: Vec::new(),
408            command_outcomes: Vec::new(),
409            interactions_created: Vec::new(),
410            event_ids: Vec::new(),
411            response_block_ids: Vec::new(),
412            outbox_ids: Vec::new(),
413            reconciliation_attempt_ids: Vec::new(),
414            provider_attempts: Vec::new(),
415            prompt_refs: Vec::new(),
416            tasks: Vec::new(),
417            budget: None,
418            effort: crate::effort::Effort::Medium,
419            recorded_at: now,
420        }
421    }
422
423    /// Every answer a check refused whole, in the order the tasks ran.
424    #[must_use]
425    pub fn discarded_answers(&self) -> Vec<DiscardedAnswer> {
426        let mut rounds: std::collections::BTreeMap<(&str, &str), u32> =
427            std::collections::BTreeMap::new();
428        self.tasks
429            .iter()
430            .filter_map(|task| match &task.verdict {
431                TaskVerdict::Rejected { code, reason } => {
432                    let owner = task.task_id.split('/').next().unwrap_or_default();
433                    let round = rounds.entry((owner, task.kind.as_str())).or_default();
434                    *round += 1;
435                    Some(DiscardedAnswer {
436                        purpose: task.kind.clone(),
437                        round: *round,
438                        code: code.clone(),
439                        reason: reason.clone(),
440                    })
441                }
442                _ => None,
443            })
444            .collect()
445    }
446
447    /// Every attempt at an external effect this record says is unsettled, in
448    /// record order and without repeats.
449    ///
450    /// Two places name one: the explicit
451    /// [`reconciliation_attempt_ids`](ReplayRecord::reconciliation_attempt_ids),
452    /// which a dispatcher adds to as it retries, and
453    /// [`CommandOutcome::OutcomeUnknown`], which execution writes against the
454    /// command that caused it. A reconciler wants the union of the two and
455    /// wants it once each, so it does not query the remote system twice for the
456    /// same attempt.
457    ///
458    /// The command outcomes come first, because they carry the attempt the turn
459    /// itself observed.
460    #[must_use]
461    pub fn pending_reconciliations(&self) -> Vec<AttemptId> {
462        let from_commands =
463            self.command_outcomes
464                .iter()
465                .filter_map(|record| match &record.outcome {
466                    CommandOutcome::OutcomeUnknown { attempt_id } => Some(attempt_id),
467                    _ => None,
468                });
469        let mut found: Vec<AttemptId> = Vec::new();
470        for attempt in from_commands.chain(self.reconciliation_attempt_ids.iter()) {
471            if !found.contains(attempt) {
472                found.push(attempt.clone());
473            }
474        }
475        found
476    }
477}
478
479/// One call of one model task: what it was asked under, and what became of the answer.
480///
481/// Identifiers are paths (`u2/extract`, `u2/extract#repair1`), so the calls of one
482/// turn read as the graph they ran as. The rendered prompt and the raw answer are
483/// kept only when the deployment asks for them.
484#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
485#[non_exhaustive]
486pub struct TaskRecord {
487    /// Where the call sits in the turn.
488    pub task_id: String,
489    /// The task that asked for this one.
490    #[serde(default, skip_serializing_if = "Option::is_none")]
491    pub parent: Option<String>,
492    /// How many calls ran before it in its chain.
493    pub depth: u8,
494    /// The task kind, as its routing purpose names it.
495    pub kind: String,
496    /// The instructions it ran under.
497    #[serde(default, skip_serializing_if = "Option::is_none")]
498    pub prompt_ref: Option<PromptRef>,
499    /// The provider that answered, when one did.
500    #[serde(default, skip_serializing_if = "Option::is_none")]
501    pub provider_key: Option<ProviderKey>,
502    /// The model that answered, when one did.
503    #[serde(default, skip_serializing_if = "Option::is_none")]
504    pub model_key: Option<ModelKey>,
505    /// The settings the call was sent with.
506    pub params: TaskParams,
507    /// Digest of the request as sent.
508    #[serde(default, skip_serializing_if = "Option::is_none")]
509    pub input_digest: Option<Digest>,
510    /// The request as sent, when the deployment keeps prompts.
511    #[serde(default, skip_serializing_if = "Option::is_none")]
512    pub rendered: Option<serde_json::Value>,
513    /// The answer as received, when the deployment keeps it.
514    #[serde(default, skip_serializing_if = "Option::is_none")]
515    pub raw_output: Option<String>,
516    /// The answer as parsed, when it parsed.
517    #[serde(default, skip_serializing_if = "Option::is_none")]
518    pub parsed: Option<serde_json::Value>,
519    /// What the runtime made of it.
520    pub verdict: TaskVerdict,
521    /// Prompt tokens the provider reported.
522    #[serde(default, skip_serializing_if = "Option::is_none")]
523    pub input_tokens: Option<u64>,
524    /// Output tokens the provider reported.
525    #[serde(default, skip_serializing_if = "Option::is_none")]
526    pub output_tokens: Option<u64>,
527    /// Wall-clock time of the call.
528    #[serde(default, skip_serializing_if = "Option::is_none")]
529    pub latency_ms: Option<u64>,
530}
531
532impl TaskRecord {
533    /// A record of `kind` at `task_id`, with nothing learnt yet.
534    #[must_use]
535    pub fn new(task_id: impl Into<String>, kind: impl Into<String>, verdict: TaskVerdict) -> Self {
536        Self {
537            task_id: task_id.into(),
538            parent: None,
539            depth: 0,
540            kind: kind.into(),
541            prompt_ref: None,
542            provider_key: None,
543            model_key: None,
544            params: TaskParams::default(),
545            input_digest: None,
546            rendered: None,
547            raw_output: None,
548            parsed: None,
549            verdict,
550            input_tokens: None,
551            output_tokens: None,
552            latency_ms: None,
553        }
554    }
555}
556
557/// The settings one task call was sent with.
558#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
559#[non_exhaustive]
560pub struct TaskParams {
561    /// Sampling temperature, when one was sent.
562    #[serde(default, skip_serializing_if = "Option::is_none")]
563    pub temperature: Option<f32>,
564    /// Output cap, when one was set.
565    #[serde(default, skip_serializing_if = "Option::is_none")]
566    pub max_output_tokens: Option<u32>,
567    /// Reasoning effort, when one was sent.
568    #[serde(default, skip_serializing_if = "Option::is_none")]
569    pub reasoning_effort: Option<String>,
570    /// Sampling seed, when one was sent.
571    #[serde(default, skip_serializing_if = "Option::is_none")]
572    pub seed: Option<u64>,
573    /// The call's deadline.
574    pub timeout_ms: u64,
575}
576
577/// What became of one task call.
578#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
579#[serde(tag = "kind", rename_all = "snake_case")]
580#[non_exhaustive]
581pub enum TaskVerdict {
582    /// The answer was used.
583    Accepted,
584    /// The answer failed a check; a repair or an escalation may follow.
585    Rejected {
586        /// Stable code of the failed check.
587        code: String,
588        /// The failure in words, as the repair round was told it.
589        reason: String,
590    },
591    /// Another answer of the same vote was used.
592    Outvoted,
593    /// No answer came back.
594    Failed {
595        /// Stable code of the failure.
596        code: String,
597    },
598}
599
600/// What a turn's model calls spent, and which bound stopped them if one did.
601#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
602#[non_exhaustive]
603pub struct BudgetReport {
604    /// Calls reserved.
605    pub model_calls: u32,
606    /// Prompt tokens the providers reported.
607    pub prompt_tokens: u64,
608    /// The longest chain of dependent calls.
609    pub max_depth: u8,
610    /// The bound that stopped the turn's model calls, when one did.
611    #[serde(default, skip_serializing_if = "Option::is_none")]
612    pub exhausted: Option<String>,
613}
614
615#[cfg(test)]
616mod tests {
617    use super::*;
618
619    #[test]
620    fn phases() {
621        assert!(TurnPhase::Failed.is_terminal());
622        assert!(!TurnPhase::Reduced.effects_may_exist());
623        assert!(TurnPhase::Executing.effects_may_exist());
624    }
625
626    #[test]
627    fn a_refused_answer_is_counted_by_its_task_chain() {
628        let rejected = |code: &str| TaskVerdict::Rejected {
629            code: code.to_owned(),
630            reason: format!("{code} in words"),
631        };
632        let mut record = full_record();
633        record.tasks = vec![
634            TaskRecord::new("u1/extract", "extract", rejected("out_of_range")),
635            TaskRecord::new("u1/extract#repair1", "extract", rejected("wrong_kind")),
636            TaskRecord::new("u1/extract#repair2", "extract", TaskVerdict::Accepted),
637            TaskRecord::new("u1.a2/extract", "extract", rejected("out_of_range")),
638        ];
639        let discarded = record.discarded_answers();
640        let rounds: Vec<(&str, u32)> = discarded
641            .iter()
642            .map(|answer| (answer.code.as_str(), answer.round))
643            .collect();
644        assert_eq!(
645            rounds,
646            vec![("out_of_range", 1), ("wrong_kind", 2), ("out_of_range", 1)]
647        );
648        assert_eq!(discarded[0].purpose, "extract");
649    }
650
651    fn attempt(temperature: Option<f32>, finish_reasons: &[&str]) -> ProviderAttemptRecord {
652        ProviderAttemptRecord {
653            attempt: 1,
654            purpose: "extract".to_owned(),
655            provider_key: ProviderKey::from("openai"),
656            model_key: ModelKey::from("gpt-5.4"),
657            request_id: "req-1".to_owned(),
658            prompt_version: None,
659            prompt_ref: None,
660            outcome: ProviderAttemptOutcome::Succeeded,
661            latency_ms: Some(120),
662            input_tokens: Some(10),
663            output_tokens: Some(20),
664            temperature,
665            finish_reasons: finish_reasons.iter().map(|s| (*s).to_owned()).collect(),
666        }
667    }
668
669    fn full_record() -> ReplayRecord {
670        let now = DateTime::from_timestamp(1_700_000_000, 0).unwrap();
671        let mut record = ReplayRecord::received(
672            TurnId::nil(),
673            ConversationId::nil(),
674            AccountId::from("a"),
675            now,
676        );
677        record.outbox_ids = vec![OutboxId::nil()];
678        record.reconciliation_attempt_ids = vec![AttemptId::from("attempt-dispatch-2")];
679        record.provider_attempts = vec![attempt(Some(0.2), &["stop", "length"])];
680        record
681    }
682
683    #[test]
684    fn record_round_trips() {
685        let now = DateTime::from_timestamp(1_700_000_000, 0).unwrap();
686        let r = ReplayRecord::received(
687            TurnId::nil(),
688            ConversationId::nil(),
689            AccountId::from("a"),
690            now,
691        );
692        let json = serde_json::to_string(&r).unwrap();
693        assert_eq!(serde_json::from_str::<ReplayRecord>(&json).unwrap(), r);
694    }
695
696    #[test]
697    fn external_effect_identifiers_round_trip() {
698        let record = full_record();
699        let json = serde_json::to_value(&record).unwrap();
700        assert_eq!(json["outbox_ids"][0], serde_json::json!(OutboxId::nil()));
701        assert_eq!(json["reconciliation_attempt_ids"][0], "attempt-dispatch-2");
702        assert_eq!(
703            serde_json::from_value::<ReplayRecord>(json).unwrap(),
704            record
705        );
706    }
707
708    #[test]
709    fn provider_attempt_carries_temperature_and_finish_reasons() {
710        let with_sampling = attempt(Some(0.2), &["stop", "length"]);
711        let json = serde_json::to_value(&with_sampling).unwrap();
712        assert_eq!(json["temperature"], serde_json::json!(0.2_f32));
713        assert_eq!(
714            json["finish_reasons"],
715            serde_json::json!(["stop", "length"])
716        );
717        assert_eq!(
718            serde_json::from_value::<ProviderAttemptRecord>(json).unwrap(),
719            with_sampling
720        );
721
722        // Absent is not zero, and an empty list of reasons stays out of the JSON.
723        let default_sampling = attempt(None, &[]);
724        let json = serde_json::to_value(&default_sampling).unwrap();
725        assert!(json.get("temperature").is_none());
726        assert!(json.get("finish_reasons").is_none());
727        assert_ne!(default_sampling, attempt(Some(0.0), &[]));
728    }
729
730    #[test]
731    fn records_written_before_the_new_fields_still_load() {
732        // Exactly what a record serialized by an older build looks like: no
733        // outbox ids, no reconciliation ids, an attempt with no sampling data.
734        let legacy = serde_json::json!({
735            "turn_id": TurnId::nil(),
736            "conversation_id": ConversationId::nil(),
737            "account_id": "a",
738            "phase": "received",
739            "workflow_versions": [],
740            "loaded_cases": [],
741            "target_resolutions": [],
742            "policy_decisions": [],
743            "command_outcomes": [],
744            "interactions_created": [],
745            "event_ids": [],
746            "response_block_ids": [],
747            "provider_attempts": [{
748                "attempt": 1,
749                "purpose": "extract",
750                "provider_key": "openai",
751                "model_key": "gpt-5.4",
752                "request_id": "req-1",
753                "outcome": { "kind": "succeeded" }
754            }],
755            "recorded_at": "2023-11-14T22:13:20Z",
756        });
757        let loaded: ReplayRecord = serde_json::from_value(legacy).unwrap();
758        assert!(loaded.outbox_ids.is_empty());
759        assert!(loaded.reconciliation_attempt_ids.is_empty());
760        assert_eq!(loaded.provider_attempts[0].temperature, None);
761        assert!(loaded.provider_attempts[0].finish_reasons.is_empty());
762        assert!(loaded.prompt_refs.is_empty());
763        assert_eq!(loaded.provider_attempts[0].prompt_ref, None);
764    }
765
766    #[test]
767    fn a_prompt_reference_reaches_both_the_turn_and_the_attempt_that_used_it() {
768        let reference = crate::prompt::PromptRef::of_text(
769            "interpret.system",
770            "9f2a1c",
771            "Answer with the plan only.",
772        );
773        let mut record = full_record();
774        record.prompt_refs = vec![reference.clone()];
775        record.provider_attempts[0].prompt_ref = Some(reference.clone());
776
777        let json = serde_json::to_value(&record).unwrap();
778        assert_eq!(json["prompt_refs"][0]["name"], "interpret.system");
779        assert_eq!(json["prompt_refs"][0]["version"], "9f2a1c");
780        assert_eq!(
781            json["provider_attempts"][0]["prompt_ref"]["hash"],
782            serde_json::json!(reference.hash.as_str())
783        );
784        assert_eq!(
785            serde_json::from_value::<ReplayRecord>(json).unwrap(),
786            record
787        );
788
789        // A turn with no prompt source keeps both out of the JSON entirely.
790        let quiet = full_record();
791        let json = serde_json::to_value(&quiet).unwrap();
792        assert!(json.get("prompt_refs").is_none());
793        assert!(json["provider_attempts"][0].get("prompt_ref").is_none());
794    }
795
796    #[test]
797    fn a_record_written_before_effort_existed_reads_as_medium() {
798        let mut value = serde_json::to_value(full_record()).unwrap();
799        value.as_object_mut().unwrap().remove("effort");
800        let read: ReplayRecord = serde_json::from_value(value).unwrap();
801        assert_eq!(read.effort, crate::effort::Effort::Medium);
802    }
803
804    #[test]
805    fn pending_reconciliations_unions_both_sources_without_repeats() {
806        let mut record = full_record();
807        let from_command = AttemptId::from("attempt-command-1");
808        record.command_outcomes = vec![
809            CommandOutcomeRecord {
810                command_ref: CommandRef {
811                    batch_id: crate::ids::BatchId::nil(),
812                    command_id: crate::ids::CommandId::nil(),
813                },
814                idempotency_key: IdempotencyKey::new("k1"),
815                case_ref: CaseRef::new("w", "c", CaseRevision(1)),
816                origin: None,
817                outcome: CommandOutcome::OutcomeUnknown {
818                    attempt_id: from_command.clone(),
819                },
820            },
821            CommandOutcomeRecord {
822                command_ref: CommandRef {
823                    batch_id: crate::ids::BatchId::nil(),
824                    command_id: crate::ids::CommandId::nil(),
825                },
826                idempotency_key: IdempotencyKey::new("k2"),
827                case_ref: CaseRef::new("w", "c", CaseRevision(1)),
828                origin: None,
829                outcome: CommandOutcome::IdempotentReplay,
830            },
831        ];
832        // The dispatcher also recorded the attempt the turn observed, plus one
833        // of its own retries.
834        record.reconciliation_attempt_ids =
835            vec![from_command.clone(), AttemptId::from("attempt-dispatch-2")];
836
837        assert_eq!(
838            record.pending_reconciliations(),
839            vec![from_command, AttemptId::from("attempt-dispatch-2")],
840            "the command outcome comes first and nothing is listed twice"
841        );
842
843        let quiet = ReplayRecord::received(
844            TurnId::nil(),
845            ConversationId::nil(),
846            AccountId::from("a"),
847            DateTime::from_timestamp(1_700_000_000, 0).unwrap(),
848        );
849        assert!(quiet.pending_reconciliations().is_empty());
850    }
851}