1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
30#[serde(rename_all = "snake_case")]
31#[non_exhaustive]
32pub enum TurnPhase {
33 Received,
35 Interpreted,
37 Reduced,
39 Executing,
41 Committed,
43 Composed,
45 Delivered,
47 Failed,
49}
50
51impl TurnPhase {
52 #[must_use]
54 pub fn is_terminal(self) -> bool {
55 matches!(self, Self::Delivered | Self::Failed)
56 }
57
58 #[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#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
71pub struct WorkflowVersionRecord {
72 pub key: WorkflowKey,
74 pub version: WorkflowVersion,
76}
77
78#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
80pub struct TargetResolutionRecord {
81 pub act: crate::understanding::ActId,
83 pub resolution: TargetResolution,
85}
86
87#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
92#[serde(tag = "kind", rename_all = "snake_case")]
93#[non_exhaustive]
94pub enum CommandOutcome {
95 Committed {
97 new_revision: CaseRevision,
99 event_ids: Vec<EventId>,
101 },
102 IdempotentReplay,
104 RevisionConflict {
106 current_revision: CaseRevision,
108 },
109 Rejected {
111 code: RejectionCode,
113 },
114 Failed {
116 code: String,
118 },
119 OutcomeUnknown {
121 attempt_id: AttemptId,
123 },
124 AwaitingConfirmation {
126 interaction_id: InteractionId,
128 },
129}
130
131#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
133pub struct CommandOutcomeRecord {
134 pub command_ref: CommandRef,
136 pub idempotency_key: IdempotencyKey,
138 pub case_ref: CaseRef,
140 #[serde(default, skip_serializing_if = "Option::is_none")]
147 pub origin: Option<CommandOrigin>,
148 pub outcome: CommandOutcome,
150}
151
152#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
156#[serde(tag = "kind", rename_all = "snake_case")]
157#[non_exhaustive]
158pub enum ProviderAttemptOutcome {
159 Succeeded,
161 Failed {
163 code: String,
165 },
166 FellBack {
168 code: String,
170 },
171 Cancelled,
173}
174
175#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
196pub struct ProviderAttemptRecord {
197 pub attempt: u32,
200 pub purpose: String,
202 pub provider_key: ProviderKey,
204 pub model_key: ModelKey,
206 pub request_id: String,
208 #[serde(default, skip_serializing_if = "Option::is_none")]
210 pub prompt_version: Option<String>,
211 #[serde(default, skip_serializing_if = "Option::is_none")]
226 pub prompt_ref: Option<PromptRef>,
227 pub outcome: ProviderAttemptOutcome,
229 #[serde(default, skip_serializing_if = "Option::is_none")]
231 pub latency_ms: Option<u64>,
232 #[serde(default, skip_serializing_if = "Option::is_none")]
234 pub input_tokens: Option<u64>,
235 #[serde(default, skip_serializing_if = "Option::is_none")]
237 pub output_tokens: Option<u64>,
238 #[serde(default, skip_serializing_if = "Option::is_none")]
245 pub temperature: Option<f32>,
246 #[serde(default, skip_serializing_if = "Vec::is_empty")]
256 pub finish_reasons: Vec<String>,
257}
258
259#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
266pub struct DiscardedAnswer {
267 pub purpose: String,
269 pub round: u32,
271 pub code: String,
273 pub reason: String,
275}
276
277#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
292pub struct ReplayRecord {
293 pub turn_id: TurnId,
295 pub conversation_id: ConversationId,
297 pub account_id: AccountId,
299 pub phase: TurnPhase,
301 pub workflow_versions: Vec<WorkflowVersionRecord>,
303 pub loaded_cases: Vec<CaseRef>,
305 #[serde(default, skip_serializing_if = "Option::is_none")]
307 pub understanding: Option<Understanding>,
308 #[serde(default, skip_serializing_if = "Option::is_none")]
310 pub plan_hash: Option<Digest>,
311 #[serde(default, skip_serializing_if = "Vec::is_empty")]
321 pub act_outcomes: Vec<String>,
322 #[serde(default, skip_serializing_if = "Option::is_none")]
324 pub reduction_plan_hash: Option<Digest>,
325 pub target_resolutions: Vec<TargetResolutionRecord>,
327 pub policy_decisions: Vec<PolicyDecision>,
329 pub command_outcomes: Vec<CommandOutcomeRecord>,
331 pub interactions_created: Vec<InteractionId>,
333 pub event_ids: Vec<EventId>,
335 pub response_block_ids: Vec<BlockId>,
337 #[serde(default, skip_serializing_if = "Vec::is_empty")]
346 pub outbox_ids: Vec<OutboxId>,
347 #[serde(default, skip_serializing_if = "Vec::is_empty")]
359 pub reconciliation_attempt_ids: Vec<AttemptId>,
360 pub provider_attempts: Vec<ProviderAttemptRecord>,
362 #[serde(default, skip_serializing_if = "Vec::is_empty")]
372 pub prompt_refs: Vec<PromptRef>,
373 #[serde(default, skip_serializing_if = "Vec::is_empty")]
375 pub tasks: Vec<TaskRecord>,
376 #[serde(default, skip_serializing_if = "Option::is_none")]
378 pub budget: Option<BudgetReport>,
379 #[serde(default)]
381 pub effort: crate::effort::Effort,
382 pub recorded_at: DateTime<Utc>,
384}
385
386impl ReplayRecord {
387 #[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 #[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 #[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#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
485#[non_exhaustive]
486pub struct TaskRecord {
487 pub task_id: String,
489 #[serde(default, skip_serializing_if = "Option::is_none")]
491 pub parent: Option<String>,
492 pub depth: u8,
494 pub kind: String,
496 #[serde(default, skip_serializing_if = "Option::is_none")]
498 pub prompt_ref: Option<PromptRef>,
499 #[serde(default, skip_serializing_if = "Option::is_none")]
501 pub provider_key: Option<ProviderKey>,
502 #[serde(default, skip_serializing_if = "Option::is_none")]
504 pub model_key: Option<ModelKey>,
505 pub params: TaskParams,
507 #[serde(default, skip_serializing_if = "Option::is_none")]
509 pub input_digest: Option<Digest>,
510 #[serde(default, skip_serializing_if = "Option::is_none")]
512 pub rendered: Option<serde_json::Value>,
513 #[serde(default, skip_serializing_if = "Option::is_none")]
515 pub raw_output: Option<String>,
516 #[serde(default, skip_serializing_if = "Option::is_none")]
518 pub parsed: Option<serde_json::Value>,
519 pub verdict: TaskVerdict,
521 #[serde(default, skip_serializing_if = "Option::is_none")]
523 pub input_tokens: Option<u64>,
524 #[serde(default, skip_serializing_if = "Option::is_none")]
526 pub output_tokens: Option<u64>,
527 #[serde(default, skip_serializing_if = "Option::is_none")]
529 pub latency_ms: Option<u64>,
530}
531
532impl TaskRecord {
533 #[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#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
559#[non_exhaustive]
560pub struct TaskParams {
561 #[serde(default, skip_serializing_if = "Option::is_none")]
563 pub temperature: Option<f32>,
564 #[serde(default, skip_serializing_if = "Option::is_none")]
566 pub max_output_tokens: Option<u32>,
567 #[serde(default, skip_serializing_if = "Option::is_none")]
569 pub reasoning_effort: Option<String>,
570 #[serde(default, skip_serializing_if = "Option::is_none")]
572 pub seed: Option<u64>,
573 pub timeout_ms: u64,
575}
576
577#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
579#[serde(tag = "kind", rename_all = "snake_case")]
580#[non_exhaustive]
581pub enum TaskVerdict {
582 Accepted,
584 Rejected {
586 code: String,
588 reason: String,
590 },
591 Outvoted,
593 Failed {
595 code: String,
597 },
598}
599
600#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
602#[non_exhaustive]
603pub struct BudgetReport {
604 pub model_calls: u32,
606 pub prompt_tokens: u64,
608 pub max_depth: u8,
610 #[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 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 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 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 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}