use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use crate::case::CaseRef;
use crate::command::{CommandOrigin, IdempotencyKey};
use crate::error::RejectionCode;
use crate::hash::Digest;
use crate::ids::{
AccountId, AttemptId, BlockId, CaseRevision, ConversationId, EventId, InteractionId, ModelKey,
OutboxId, ProviderKey, TurnId, WorkflowKey, WorkflowVersion,
};
use crate::policy::PolicyDecision;
use crate::prompt::PromptRef;
use crate::reduce::CommandRef;
use crate::target::TargetResolution;
use crate::understanding::Understanding;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
#[non_exhaustive]
pub enum TurnPhase {
Received,
Interpreted,
Reduced,
Executing,
Committed,
Composed,
Delivered,
Failed,
}
impl TurnPhase {
#[must_use]
pub fn is_terminal(self) -> bool {
matches!(self, Self::Delivered | Self::Failed)
}
#[must_use]
pub fn effects_may_exist(self) -> bool {
matches!(
self,
Self::Executing | Self::Committed | Self::Composed | Self::Delivered
)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct WorkflowVersionRecord {
pub key: WorkflowKey,
pub version: WorkflowVersion,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct TargetResolutionRecord {
pub act: crate::understanding::ActId,
pub resolution: TargetResolution,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
#[non_exhaustive]
pub enum CommandOutcome {
Committed {
new_revision: CaseRevision,
event_ids: Vec<EventId>,
},
IdempotentReplay,
RevisionConflict {
current_revision: CaseRevision,
},
Rejected {
code: RejectionCode,
},
Failed {
code: String,
},
OutcomeUnknown {
attempt_id: AttemptId,
},
AwaitingConfirmation {
interaction_id: InteractionId,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CommandOutcomeRecord {
pub command_ref: CommandRef,
pub idempotency_key: IdempotencyKey,
pub case_ref: CaseRef,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub origin: Option<CommandOrigin>,
pub outcome: CommandOutcome,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
#[non_exhaustive]
pub enum ProviderAttemptOutcome {
Succeeded,
Failed {
code: String,
},
FellBack {
code: String,
},
Cancelled,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ProviderAttemptRecord {
pub attempt: u32,
pub purpose: String,
pub provider_key: ProviderKey,
pub model_key: ModelKey,
pub request_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub prompt_version: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub prompt_ref: Option<PromptRef>,
pub outcome: ProviderAttemptOutcome,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub latency_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub input_tokens: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_tokens: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub temperature: Option<f32>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub finish_reasons: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct DiscardedAnswer {
pub purpose: String,
pub round: u32,
pub code: String,
pub reason: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ReplayRecord {
pub turn_id: TurnId,
pub conversation_id: ConversationId,
pub account_id: AccountId,
pub phase: TurnPhase,
pub workflow_versions: Vec<WorkflowVersionRecord>,
pub loaded_cases: Vec<CaseRef>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub understanding: Option<Understanding>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub plan_hash: Option<Digest>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub act_outcomes: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reduction_plan_hash: Option<Digest>,
pub target_resolutions: Vec<TargetResolutionRecord>,
pub policy_decisions: Vec<PolicyDecision>,
pub command_outcomes: Vec<CommandOutcomeRecord>,
pub interactions_created: Vec<InteractionId>,
pub event_ids: Vec<EventId>,
pub response_block_ids: Vec<BlockId>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub outbox_ids: Vec<OutboxId>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub reconciliation_attempt_ids: Vec<AttemptId>,
pub provider_attempts: Vec<ProviderAttemptRecord>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub prompt_refs: Vec<PromptRef>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub tasks: Vec<TaskRecord>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub budget: Option<BudgetReport>,
#[serde(default)]
pub effort: crate::effort::Effort,
pub recorded_at: DateTime<Utc>,
}
impl ReplayRecord {
#[must_use]
pub fn received(
turn_id: TurnId,
conversation_id: ConversationId,
account_id: AccountId,
now: DateTime<Utc>,
) -> Self {
Self {
turn_id,
conversation_id,
account_id,
phase: TurnPhase::Received,
workflow_versions: Vec::new(),
loaded_cases: Vec::new(),
understanding: None,
plan_hash: None,
act_outcomes: Vec::new(),
reduction_plan_hash: None,
target_resolutions: Vec::new(),
policy_decisions: Vec::new(),
command_outcomes: Vec::new(),
interactions_created: Vec::new(),
event_ids: Vec::new(),
response_block_ids: Vec::new(),
outbox_ids: Vec::new(),
reconciliation_attempt_ids: Vec::new(),
provider_attempts: Vec::new(),
prompt_refs: Vec::new(),
tasks: Vec::new(),
budget: None,
effort: crate::effort::Effort::Medium,
recorded_at: now,
}
}
#[must_use]
pub fn discarded_answers(&self) -> Vec<DiscardedAnswer> {
let mut rounds: std::collections::BTreeMap<(&str, &str), u32> =
std::collections::BTreeMap::new();
self.tasks
.iter()
.filter_map(|task| match &task.verdict {
TaskVerdict::Rejected { code, reason } => {
let owner = task.task_id.split('/').next().unwrap_or_default();
let round = rounds.entry((owner, task.kind.as_str())).or_default();
*round += 1;
Some(DiscardedAnswer {
purpose: task.kind.clone(),
round: *round,
code: code.clone(),
reason: reason.clone(),
})
}
_ => None,
})
.collect()
}
#[must_use]
pub fn pending_reconciliations(&self) -> Vec<AttemptId> {
let from_commands =
self.command_outcomes
.iter()
.filter_map(|record| match &record.outcome {
CommandOutcome::OutcomeUnknown { attempt_id } => Some(attempt_id),
_ => None,
});
let mut found: Vec<AttemptId> = Vec::new();
for attempt in from_commands.chain(self.reconciliation_attempt_ids.iter()) {
if !found.contains(attempt) {
found.push(attempt.clone());
}
}
found
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[non_exhaustive]
pub struct TaskRecord {
pub task_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parent: Option<String>,
pub depth: u8,
pub kind: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub prompt_ref: Option<PromptRef>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub provider_key: Option<ProviderKey>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model_key: Option<ModelKey>,
pub params: TaskParams,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub input_digest: Option<Digest>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub rendered: Option<serde_json::Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub raw_output: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parsed: Option<serde_json::Value>,
pub verdict: TaskVerdict,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub input_tokens: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_tokens: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub latency_ms: Option<u64>,
}
impl TaskRecord {
#[must_use]
pub fn new(task_id: impl Into<String>, kind: impl Into<String>, verdict: TaskVerdict) -> Self {
Self {
task_id: task_id.into(),
parent: None,
depth: 0,
kind: kind.into(),
prompt_ref: None,
provider_key: None,
model_key: None,
params: TaskParams::default(),
input_digest: None,
rendered: None,
raw_output: None,
parsed: None,
verdict,
input_tokens: None,
output_tokens: None,
latency_ms: None,
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
#[non_exhaustive]
pub struct TaskParams {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub temperature: Option<f32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_output_tokens: Option<u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reasoning_effort: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub seed: Option<u64>,
pub timeout_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
#[non_exhaustive]
pub enum TaskVerdict {
Accepted,
Rejected {
code: String,
reason: String,
},
Outvoted,
Failed {
code: String,
},
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
#[non_exhaustive]
pub struct BudgetReport {
pub model_calls: u32,
pub prompt_tokens: u64,
pub max_depth: u8,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub exhausted: Option<String>,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn phases() {
assert!(TurnPhase::Failed.is_terminal());
assert!(!TurnPhase::Reduced.effects_may_exist());
assert!(TurnPhase::Executing.effects_may_exist());
}
#[test]
fn a_refused_answer_is_counted_by_its_task_chain() {
let rejected = |code: &str| TaskVerdict::Rejected {
code: code.to_owned(),
reason: format!("{code} in words"),
};
let mut record = full_record();
record.tasks = vec![
TaskRecord::new("u1/extract", "extract", rejected("out_of_range")),
TaskRecord::new("u1/extract#repair1", "extract", rejected("wrong_kind")),
TaskRecord::new("u1/extract#repair2", "extract", TaskVerdict::Accepted),
TaskRecord::new("u1.a2/extract", "extract", rejected("out_of_range")),
];
let discarded = record.discarded_answers();
let rounds: Vec<(&str, u32)> = discarded
.iter()
.map(|answer| (answer.code.as_str(), answer.round))
.collect();
assert_eq!(
rounds,
vec![("out_of_range", 1), ("wrong_kind", 2), ("out_of_range", 1)]
);
assert_eq!(discarded[0].purpose, "extract");
}
fn attempt(temperature: Option<f32>, finish_reasons: &[&str]) -> ProviderAttemptRecord {
ProviderAttemptRecord {
attempt: 1,
purpose: "extract".to_owned(),
provider_key: ProviderKey::from("openai"),
model_key: ModelKey::from("gpt-5.4"),
request_id: "req-1".to_owned(),
prompt_version: None,
prompt_ref: None,
outcome: ProviderAttemptOutcome::Succeeded,
latency_ms: Some(120),
input_tokens: Some(10),
output_tokens: Some(20),
temperature,
finish_reasons: finish_reasons.iter().map(|s| (*s).to_owned()).collect(),
}
}
fn full_record() -> ReplayRecord {
let now = DateTime::from_timestamp(1_700_000_000, 0).unwrap();
let mut record = ReplayRecord::received(
TurnId::nil(),
ConversationId::nil(),
AccountId::from("a"),
now,
);
record.outbox_ids = vec![OutboxId::nil()];
record.reconciliation_attempt_ids = vec![AttemptId::from("attempt-dispatch-2")];
record.provider_attempts = vec![attempt(Some(0.2), &["stop", "length"])];
record
}
#[test]
fn record_round_trips() {
let now = DateTime::from_timestamp(1_700_000_000, 0).unwrap();
let r = ReplayRecord::received(
TurnId::nil(),
ConversationId::nil(),
AccountId::from("a"),
now,
);
let json = serde_json::to_string(&r).unwrap();
assert_eq!(serde_json::from_str::<ReplayRecord>(&json).unwrap(), r);
}
#[test]
fn external_effect_identifiers_round_trip() {
let record = full_record();
let json = serde_json::to_value(&record).unwrap();
assert_eq!(json["outbox_ids"][0], serde_json::json!(OutboxId::nil()));
assert_eq!(json["reconciliation_attempt_ids"][0], "attempt-dispatch-2");
assert_eq!(
serde_json::from_value::<ReplayRecord>(json).unwrap(),
record
);
}
#[test]
fn provider_attempt_carries_temperature_and_finish_reasons() {
let with_sampling = attempt(Some(0.2), &["stop", "length"]);
let json = serde_json::to_value(&with_sampling).unwrap();
assert_eq!(json["temperature"], serde_json::json!(0.2_f32));
assert_eq!(
json["finish_reasons"],
serde_json::json!(["stop", "length"])
);
assert_eq!(
serde_json::from_value::<ProviderAttemptRecord>(json).unwrap(),
with_sampling
);
let default_sampling = attempt(None, &[]);
let json = serde_json::to_value(&default_sampling).unwrap();
assert!(json.get("temperature").is_none());
assert!(json.get("finish_reasons").is_none());
assert_ne!(default_sampling, attempt(Some(0.0), &[]));
}
#[test]
fn records_written_before_the_new_fields_still_load() {
let legacy = serde_json::json!({
"turn_id": TurnId::nil(),
"conversation_id": ConversationId::nil(),
"account_id": "a",
"phase": "received",
"workflow_versions": [],
"loaded_cases": [],
"target_resolutions": [],
"policy_decisions": [],
"command_outcomes": [],
"interactions_created": [],
"event_ids": [],
"response_block_ids": [],
"provider_attempts": [{
"attempt": 1,
"purpose": "extract",
"provider_key": "openai",
"model_key": "gpt-5.4",
"request_id": "req-1",
"outcome": { "kind": "succeeded" }
}],
"recorded_at": "2023-11-14T22:13:20Z",
});
let loaded: ReplayRecord = serde_json::from_value(legacy).unwrap();
assert!(loaded.outbox_ids.is_empty());
assert!(loaded.reconciliation_attempt_ids.is_empty());
assert_eq!(loaded.provider_attempts[0].temperature, None);
assert!(loaded.provider_attempts[0].finish_reasons.is_empty());
assert!(loaded.prompt_refs.is_empty());
assert_eq!(loaded.provider_attempts[0].prompt_ref, None);
}
#[test]
fn a_prompt_reference_reaches_both_the_turn_and_the_attempt_that_used_it() {
let reference = crate::prompt::PromptRef::of_text(
"interpret.system",
"9f2a1c",
"Answer with the plan only.",
);
let mut record = full_record();
record.prompt_refs = vec![reference.clone()];
record.provider_attempts[0].prompt_ref = Some(reference.clone());
let json = serde_json::to_value(&record).unwrap();
assert_eq!(json["prompt_refs"][0]["name"], "interpret.system");
assert_eq!(json["prompt_refs"][0]["version"], "9f2a1c");
assert_eq!(
json["provider_attempts"][0]["prompt_ref"]["hash"],
serde_json::json!(reference.hash.as_str())
);
assert_eq!(
serde_json::from_value::<ReplayRecord>(json).unwrap(),
record
);
let quiet = full_record();
let json = serde_json::to_value(&quiet).unwrap();
assert!(json.get("prompt_refs").is_none());
assert!(json["provider_attempts"][0].get("prompt_ref").is_none());
}
#[test]
fn a_record_written_before_effort_existed_reads_as_medium() {
let mut value = serde_json::to_value(full_record()).unwrap();
value.as_object_mut().unwrap().remove("effort");
let read: ReplayRecord = serde_json::from_value(value).unwrap();
assert_eq!(read.effort, crate::effort::Effort::Medium);
}
#[test]
fn pending_reconciliations_unions_both_sources_without_repeats() {
let mut record = full_record();
let from_command = AttemptId::from("attempt-command-1");
record.command_outcomes = vec![
CommandOutcomeRecord {
command_ref: CommandRef {
batch_id: crate::ids::BatchId::nil(),
command_id: crate::ids::CommandId::nil(),
},
idempotency_key: IdempotencyKey::new("k1"),
case_ref: CaseRef::new("w", "c", CaseRevision(1)),
origin: None,
outcome: CommandOutcome::OutcomeUnknown {
attempt_id: from_command.clone(),
},
},
CommandOutcomeRecord {
command_ref: CommandRef {
batch_id: crate::ids::BatchId::nil(),
command_id: crate::ids::CommandId::nil(),
},
idempotency_key: IdempotencyKey::new("k2"),
case_ref: CaseRef::new("w", "c", CaseRevision(1)),
origin: None,
outcome: CommandOutcome::IdempotentReplay,
},
];
record.reconciliation_attempt_ids =
vec![from_command.clone(), AttemptId::from("attempt-dispatch-2")];
assert_eq!(
record.pending_reconciliations(),
vec![from_command, AttemptId::from("attempt-dispatch-2")],
"the command outcome comes first and nothing is listed twice"
);
let quiet = ReplayRecord::received(
TurnId::nil(),
ConversationId::nil(),
AccountId::from("a"),
DateTime::from_timestamp(1_700_000_000, 0).unwrap(),
);
assert!(quiet.pending_reconciliations().is_empty());
}
}