1#![deny(missing_docs)]
5#![deny(rustdoc::broken_intra_doc_links)]
6
7mod corrections;
8pub use corrections::{
9 AppliedMeteringCorrection, MeteringCorrection, OperationUsage, RunOperationUsage,
10};
11
12mod metering;
13pub use metering::{MeteringDetails, MeteringOutcome, MeteringSource};
14
15use af_context::{InputId, InteractionId, ProfileRevisionId, RunId, SessionId, ToolCallId};
16use std::collections::{BTreeMap, BTreeSet};
17
18use async_trait::async_trait;
19use chrono::{DateTime, Utc};
20use serde::{Deserialize, Serialize};
21use serde_json::Value;
22
23pub const SESSION_EVENT_FORMAT_VERSION: u32 = 1;
25
26#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
28#[serde(deny_unknown_fields)]
29pub struct SessionMetadata {
30 pub title: String,
32 pub archived: bool,
34 #[serde(default)]
36 pub approval_mode: ApprovalMode,
37 pub version: u64,
39}
40
41#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
43#[serde(rename_all = "snake_case")]
44pub enum ApprovalMode {
45 ReadOnly,
47 #[default]
49 ConfirmChanges,
50 LowRiskAuto,
52}
53
54impl ApprovalMode {
55 pub const fn as_str(self) -> &'static str {
57 match self {
58 Self::ReadOnly => "read_only",
59 Self::ConfirmChanges => "confirm_changes",
60 Self::LowRiskAuto => "low_risk_auto",
61 }
62 }
63}
64
65impl std::str::FromStr for ApprovalMode {
66 type Err = EventError;
67
68 fn from_str(value: &str) -> Result<Self, Self::Err> {
69 match value {
70 "read_only" => Ok(Self::ReadOnly),
71 "confirm_changes" => Ok(Self::ConfirmChanges),
72 "low_risk_auto" => Ok(Self::LowRiskAuto),
73 _ => Err(EventError::InvalidSessionMetadata),
74 }
75 }
76}
77
78impl SessionMetadata {
79 pub fn validate(&self) -> Result<(), EventError> {
81 if self.title.chars().count() > 200 || self.title.chars().any(char::is_control) {
82 return Err(EventError::InvalidSessionMetadata);
83 }
84 Ok(())
85 }
86}
87
88#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
90pub struct SessionEvent {
91 pub session_id: SessionId,
93 pub seq: u64,
95 pub occurred_at: DateTime<Utc>,
97 pub event: Event,
99}
100
101impl SessionEvent {
102 pub fn pending(session_id: impl Into<SessionId>, event: Event) -> Self {
104 Self {
105 session_id: session_id.into(),
106 seq: 0,
107 occurred_at: Utc::now(),
108 event,
109 }
110 }
111
112 pub fn format_version(&self) -> u32 {
114 self.event.format_version()
115 }
116
117 pub fn event_type(&self) -> &str {
119 self.event.event_type()
120 }
121
122 pub fn ignorable(&self) -> bool {
124 self.event.ignorable()
125 }
126}
127
128#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
130#[serde(tag = "type", rename_all = "snake_case")]
131pub enum Event {
132 SessionCreated {
134 profile_revision_id: ProfileRevisionId,
136 },
137 SessionForked {
139 parent_session_id: SessionId,
141 parent_seq: u64,
143 },
144 SessionMetadataUpdated {
146 metadata: SessionMetadata,
148 },
149 SessionDeleted {
151 reason: String,
153 },
154 InputQueued {
156 input_id: InputId,
158 run_id: RunId,
160 mode: DeliveryMode,
162 content: Vec<ContentBlock>,
164 explicit_skill: Option<String>,
166 },
167 InputClaimed {
169 input_id: InputId,
171 run_id: RunId,
173 },
174 InputCancelled {
176 input_id: InputId,
178 run_id: RunId,
180 error_code: String,
182 },
183 RunStarted {
185 run_id: RunId,
187 input_id: InputId,
189 },
190 RunWaiting {
192 run_id: RunId,
194 interaction_id: InteractionId,
196 },
197 RunResumed {
199 run_id: RunId,
201 interaction_id: InteractionId,
203 },
204 RunFinished {
206 run_id: RunId,
208 status: RunStatus,
210 error_code: Option<String>,
212 },
213 TurnStarted {
215 run_id: RunId,
217 turn: u32,
219 },
220 TurnFinished {
222 run_id: RunId,
224 turn: u32,
226 },
227 StepStarted {
229 run_id: RunId,
231 step: u32,
233 },
234 StepFinished {
236 run_id: RunId,
238 step: u32,
240 },
241 UserMessage {
243 run_id: RunId,
245 content: Vec<ContentBlock>,
247 },
248 AssistantDelta {
250 run_id: RunId,
252 step: u32,
254 attempt: u32,
256 content: String,
258 },
259 AssistantMessage {
261 run_id: RunId,
263 step: u32,
265 attempt: u32,
267 content: Vec<ContentBlock>,
269 },
270 AssistantToolCalls {
272 run_id: RunId,
274 step: u32,
276 content: Option<String>,
278 calls: Vec<RecordedToolCall>,
280 },
281 ToolCall {
283 run_id: RunId,
285 step: u32,
287 call_id: ToolCallId,
289 tool: String,
291 arguments: Value,
293 },
294 ToolAuthorization {
296 run_id: RunId,
298 step: u32,
300 call_id: ToolCallId,
302 status: ToolAuthorizationStatus,
304 reason: Option<String>,
306 },
307 ToolExecutionStarted {
309 #[serde(default, skip_serializing_if = "Option::is_none")]
311 metering: Option<MeteringDetails>,
312 #[serde(default)]
314 reserved_cost_units: u64,
315 run_id: RunId,
317 step: u32,
319 call_id: ToolCallId,
321 },
322 ToolResult {
324 run_id: RunId,
326 step: u32,
328 call_id: ToolCallId,
330 result: Value,
332 is_error: bool,
334 },
335 UsageCorrected {
337 correction: MeteringCorrection,
339 actor_id: af_context::SubjectId,
341 },
342 UsageRecorded {
344 #[serde(default, skip_serializing_if = "Option::is_none")]
346 metering: Option<MeteringDetails>,
347 run_id: RunId,
349 operation_id: String,
351 prompt_tokens: u64,
353 completion_tokens: u64,
355 #[serde(default)]
357 cost_units: u64,
358 },
359 RetryScheduled {
361 run_id: RunId,
363 attempt: u32,
365 delay_ms: u64,
367 reason: String,
369 },
370 ModelRequestPrepared {
372 #[serde(default, skip_serializing_if = "Option::is_none")]
374 metering: Option<MeteringDetails>,
375 run_id: RunId,
377 step: u32,
379 attempt: u32,
381 #[serde(default)]
383 provider_attempt_id: String,
384 #[serde(default)]
386 operation_id: String,
387 #[serde(default)]
389 reserved_prompt_tokens: u64,
390 #[serde(default)]
392 reserved_completion_tokens: u64,
393 request: Value,
395 prompt_sections: Value,
397 },
398 ContextInjected {
400 run_id: RunId,
402 step: u32,
404 contribution_id: String,
406 source: String,
408 version: String,
410 authority: String,
412 form: String,
414 content: Vec<ContentBlock>,
416 },
417 ModelAttemptFailed {
419 run_id: RunId,
421 step: u32,
423 attempt: u32,
425 error: String,
427 retryable: bool,
429 },
430 CompactionStarted {
432 run_id: RunId,
434 compaction_id: String,
436 source_through_seq: u64,
438 },
439 ToolResultsPruned {
441 run_id: RunId,
443 call_ids: Vec<ToolCallId>,
445 },
446 SummaryReplaced {
448 run_id: RunId,
450 through_seq: u64,
452 summary: String,
454 compactor: String,
456 model: String,
458 },
459 CompactionFinished {
461 run_id: RunId,
463 compaction_id: String,
465 status: String,
467 error: Option<String>,
469 },
470 InteractionRequested {
472 run_id: RunId,
474 interaction_id: InteractionId,
476 kind: InteractionKind,
478 payload: Value,
480 },
481 InteractionResolved {
483 run_id: RunId,
485 interaction_id: InteractionId,
487 resolution: InteractionResolution,
489 payload: Value,
491 },
492 ChildSessionLinked {
494 run_id: RunId,
496 child_session_id: SessionId,
498 provider: String,
500 },
501 Extension {
503 run_id: RunId,
505 plugin_id: String,
507 event_type: String,
509 payload: Value,
511 },
512 #[serde(skip)]
514 Opaque {
515 format_version: u32,
517 event_type: String,
519 ignorable: bool,
521 payload: Value,
523 },
524}
525
526impl Event {
527 pub fn format_version(&self) -> u32 {
529 match self {
530 Self::Opaque { format_version, .. } => *format_version,
531 _ => SESSION_EVENT_FORMAT_VERSION,
532 }
533 }
534
535 pub fn event_type(&self) -> &str {
537 match self {
538 Self::SessionCreated { .. } => "session_created",
539 Self::SessionMetadataUpdated { .. } => "session_metadata_updated",
540 Self::SessionForked { .. } => "session_forked",
541 Self::SessionDeleted { .. } => "session_deleted",
542 Self::InputQueued { .. } => "input_queued",
543 Self::InputClaimed { .. } => "input_claimed",
544 Self::InputCancelled { .. } => "input_cancelled",
545 Self::RunStarted { .. } => "run_started",
546 Self::RunWaiting { .. } => "run_waiting",
547 Self::RunResumed { .. } => "run_resumed",
548 Self::RunFinished { .. } => "run_finished",
549 Self::TurnStarted { .. } => "turn_started",
550 Self::TurnFinished { .. } => "turn_finished",
551 Self::StepStarted { .. } => "step_started",
552 Self::StepFinished { .. } => "step_finished",
553 Self::UserMessage { .. } => "user_message",
554 Self::AssistantDelta { .. } => "assistant_delta",
555 Self::AssistantMessage { .. } => "assistant_message",
556 Self::AssistantToolCalls { .. } => "assistant_tool_calls",
557 Self::ToolCall { .. } => "tool_call",
558 Self::ToolAuthorization { .. } => "tool_authorization",
559 Self::ToolExecutionStarted { .. } => "tool_execution_started",
560 Self::ToolResult { .. } => "tool_result",
561 Self::UsageRecorded { .. } => "usage_recorded",
562 Self::UsageCorrected { .. } => "usage_corrected",
563 Self::RetryScheduled { .. } => "retry_scheduled",
564 Self::ModelRequestPrepared { .. } => "model_request_prepared",
565 Self::ContextInjected { .. } => "context_injected",
566 Self::ModelAttemptFailed { .. } => "model_attempt_failed",
567 Self::CompactionStarted { .. } => "compaction_started",
568 Self::ToolResultsPruned { .. } => "tool_results_pruned",
569 Self::SummaryReplaced { .. } => "summary_replaced",
570 Self::CompactionFinished { .. } => "compaction_finished",
571 Self::InteractionRequested { .. } => "interaction_requested",
572 Self::InteractionResolved { .. } => "interaction_resolved",
573 Self::ChildSessionLinked { .. } => "child_session_linked",
574 Self::Extension { .. } => "extension",
575 Self::Opaque { event_type, .. } => event_type,
576 }
577 }
578
579 pub fn ignorable(&self) -> bool {
581 match self {
582 Self::Opaque { ignorable, .. } => *ignorable,
583 _ => false,
584 }
585 }
586}
587
588#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
590#[serde(rename_all = "snake_case")]
591pub enum DeliveryMode {
592 Followup,
594 Steer,
596 Inject,
598}
599
600#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
602#[serde(rename_all = "snake_case")]
603pub enum RunStatus {
604 Completed,
606 Failed,
608 Cancelled,
610 MaxStepsReached,
612}
613
614impl RunStatus {
615 pub const fn as_str(self) -> &'static str {
617 match self {
618 Self::Completed => "completed",
619 Self::Failed => "failed",
620 Self::Cancelled => "cancelled",
621 Self::MaxStepsReached => "max_steps_reached",
622 }
623 }
624}
625
626#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
628#[serde(rename_all = "snake_case")]
629pub enum InteractionKind {
630 Action,
632 UserQuestion,
634}
635
636impl InteractionKind {
637 pub const fn as_str(self) -> &'static str {
639 match self {
640 Self::Action => "action",
641 Self::UserQuestion => "user_question",
642 }
643 }
644}
645
646#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
648#[serde(rename_all = "snake_case")]
649pub enum InteractionResolution {
650 Confirmed,
652 Rejected,
654 Answered,
656}
657
658impl InteractionResolution {
659 pub const fn as_str(self) -> &'static str {
661 match self {
662 Self::Confirmed => "confirmed",
663 Self::Rejected => "rejected",
664 Self::Answered => "answered",
665 }
666 }
667}
668
669#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
671#[serde(rename_all = "snake_case")]
672pub enum ToolAuthorizationStatus {
673 Allowed,
675 Waiting,
677 Denied,
679}
680
681#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
683pub struct RecordedToolCall {
684 pub call_id: ToolCallId,
686 pub tool: String,
688 pub arguments: Value,
690}
691
692impl ToolAuthorizationStatus {
693 pub const fn as_str(self) -> &'static str {
695 match self {
696 Self::Allowed => "allowed",
697 Self::Waiting => "waiting",
698 Self::Denied => "denied",
699 }
700 }
701}
702
703#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
705#[serde(tag = "type", rename_all = "snake_case")]
706pub enum ContentBlock {
707 Text {
709 text: String,
711 },
712 Resource {
714 resource_id: String,
716 media_type: String,
718 },
719 Data {
721 slot: String,
723 value: Value,
725 },
726 Citation {
728 resource_id: String,
730 label: String,
732 uri: String,
734 excerpt: Option<String>,
736 },
737}
738
739#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
741pub struct UsageCharge {
742 pub prompt_tokens: u64,
744 pub completion_tokens: u64,
746 pub cost_units: u64,
748}
749
750impl UsageCharge {
751 pub fn total(self) -> u64 {
753 self.prompt_tokens
754 .saturating_add(self.completion_tokens)
755 .saturating_add(self.cost_units)
756 }
757}
758
759#[derive(Debug, Clone, Default, PartialEq)]
761pub struct SessionProjection {
762 pub metadata: SessionMetadata,
764 pub session_id: Option<SessionId>,
766 pub profile_revision_id: Option<ProfileRevisionId>,
769 pub deleted: bool,
771 pub last_seq: u64,
773 pub active_run_id: Option<RunId>,
775 pub waiting_interaction_id: Option<InteractionId>,
777 pub messages: Vec<ProjectedMessage>,
779 pub injected_context: Vec<ProjectedContext>,
781 pub run_status: BTreeMap<RunId, RunState>,
783 pub open_tool_calls: BTreeMap<ToolCallId, OpenToolCall>,
785 pub started_tool_calls: BTreeSet<ToolCallId>,
787 pub queued_inputs: BTreeMap<InputId, (RunId, DeliveryMode)>,
789 pub claimed_inputs: BTreeMap<InputId, RunId>,
791 pub open_turn: Option<u32>,
793 pub open_steps: BTreeSet<u32>,
795 pub next_step: u32,
797 pub open_compaction: Option<(String, u64)>,
799 pub summary: Option<String>,
801 corrected_usage: BTreeMap<(RunId, String), OperationUsage>,
802 metering_corrections: BTreeMap<af_context::MeteringCorrectionId, AppliedMeteringCorrection>,
803 usage_operations: BTreeMap<(RunId, String), RecordedUsage>,
804 pending_usage_operations: BTreeMap<(RunId, String), PreparedUsage>,
805 seen_tool_calls: BTreeSet<ToolCallId>,
806}
807
808type PreparedUsage = RecordedUsage;
809
810type RecordedUsage = (u64, u64, u64, Option<MeteringDetails>);
811
812#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
814#[serde(rename_all = "snake_case")]
815pub enum RunState {
816 Running,
818 WaitingForInput,
820 Terminal(RunStatus),
822}
823
824impl RunState {
825 pub const fn as_str(self) -> &'static str {
827 match self {
828 Self::Running => "running",
829 Self::WaitingForInput => "waiting_for_input",
830 Self::Terminal(status) => status.as_str(),
831 }
832 }
833
834 pub const fn is_terminal(self) -> bool {
836 matches!(self, Self::Terminal(_))
837 }
838}
839
840#[derive(Debug, Clone, PartialEq)]
842pub struct ProjectedContext {
843 pub run_id: RunId,
845 pub step: u32,
847 pub contribution_id: String,
849 pub source: String,
851 pub version: String,
853 pub authority: String,
855 pub form: String,
857 pub content: Vec<ContentBlock>,
859}
860
861#[derive(Debug, Clone, PartialEq)]
863pub struct ProjectedMessage {
864 pub role: &'static str,
866 pub run_id: RunId,
868 pub content: Vec<ContentBlock>,
870}
871
872#[derive(Debug, Clone, PartialEq)]
874pub struct OpenToolCall {
875 pub run_id: RunId,
877 pub step: u32,
879 pub tool: String,
881 pub arguments: Value,
883 pub source_event_seq: u64,
885}
886
887impl SessionProjection {
888 pub fn replay(events: &[SessionEvent]) -> Result<Self, EventError> {
890 let mut projection = Self::default();
891 for event in events {
892 projection.apply(event)?;
893 }
894 Ok(projection)
895 }
896
897 pub fn apply(&mut self, envelope: &SessionEvent) -> Result<(), EventError> {
899 self.apply_internal(envelope, true)
900 }
901
902 pub fn apply_facts(&mut self, envelope: &SessionEvent) -> Result<(), EventError> {
905 self.messages.clear();
906 self.injected_context.clear();
907 self.apply_internal(envelope, false)
908 }
909
910 pub fn replay_facts(events: &[SessionEvent]) -> Result<Self, EventError> {
912 let mut projection = Self::default();
913 for event in events {
914 projection.apply_facts(event)?;
915 }
916 Ok(projection)
917 }
918
919 fn apply_internal(
920 &mut self,
921 envelope: &SessionEvent,
922 include_content: bool,
923 ) -> Result<(), EventError> {
924 if envelope.seq != self.last_seq + 1 {
925 return Err(EventError::Sequence {
926 expected: self.last_seq + 1,
927 actual: envelope.seq,
928 });
929 }
930 let session_id = self
931 .session_id
932 .get_or_insert_with(|| envelope.session_id.clone());
933 if *session_id != envelope.session_id {
934 return Err(EventError::SessionMismatch);
935 }
936 if envelope.format_version() != SESSION_EVENT_FORMAT_VERSION {
937 if envelope.ignorable() {
938 self.last_seq = envelope.seq;
939 return Ok(());
940 }
941 return Err(EventError::UnsupportedFormat(envelope.format_version()));
942 }
943 if self.deleted && !matches!(envelope.event, Event::UsageCorrected { .. }) {
944 return Err(EventError::SessionClosed);
945 }
946 match &envelope.event {
947 Event::UsageCorrected {
948 correction,
949 actor_id,
950 } => self.apply_metering_correction(correction, actor_id, envelope.seq)?,
951 Event::SessionCreated {
952 profile_revision_id,
953 } => {
954 if envelope.seq != 1 || self.profile_revision_id.is_some() {
955 return Err(EventError::DuplicateSession);
956 }
957 self.profile_revision_id = Some(profile_revision_id.clone());
958 }
959 Event::SessionMetadataUpdated { metadata } => {
960 metadata.validate()?;
961 if self.profile_revision_id.is_none()
962 || self.metadata.version.checked_add(1) != Some(metadata.version)
963 {
964 return Err(EventError::InvalidSessionMetadata);
965 }
966 self.metadata = metadata.clone();
967 }
968 Event::SessionDeleted { .. } => {
969 if let Some(run_id) = self.active_run_id.take() {
970 self.run_status
971 .insert(run_id, RunState::Terminal(RunStatus::Cancelled));
972 }
973 for (_, (run_id, _)) in std::mem::take(&mut self.queued_inputs) {
974 self.run_status
975 .insert(run_id, RunState::Terminal(RunStatus::Cancelled));
976 }
977 self.waiting_interaction_id = None;
978 self.open_turn = None;
979 self.open_steps.clear();
980 self.open_tool_calls.clear();
981 self.started_tool_calls.clear();
982 self.open_compaction = None;
983 self.deleted = true;
984 }
985 Event::RunStarted { run_id, input_id } => {
986 if self.active_run_id.is_some() {
987 return Err(EventError::ConcurrentRun);
988 }
989 if self.claimed_inputs.get(input_id) != Some(run_id) {
990 return Err(EventError::UnclaimedInput(input_id.clone()));
991 }
992 self.active_run_id = Some(run_id.clone());
993 self.next_step = 1;
994 self.run_status.insert(run_id.clone(), RunState::Running);
995 }
996 Event::RunWaiting {
997 run_id,
998 interaction_id,
999 } => {
1000 self.require_active(run_id)?;
1001 self.waiting_interaction_id = Some(interaction_id.clone());
1002 self.run_status
1003 .insert(run_id.clone(), RunState::WaitingForInput);
1004 }
1005 Event::RunResumed {
1006 run_id,
1007 interaction_id,
1008 } => {
1009 self.require_active(run_id)?;
1010 if self.waiting_interaction_id.as_ref() != Some(interaction_id) {
1011 return Err(EventError::InteractionMismatch);
1012 }
1013 self.waiting_interaction_id = None;
1014 self.run_status.insert(run_id.clone(), RunState::Running);
1015 }
1016 Event::RunFinished { run_id, status, .. } => {
1017 self.require_active(run_id)?;
1018 if !self.open_tool_calls.is_empty()
1019 || !self.open_steps.is_empty()
1020 || self.open_turn.is_some()
1021 || self.open_compaction.is_some()
1022 || self
1023 .queued_inputs
1024 .values()
1025 .any(|(target_run_id, _)| target_run_id == run_id)
1026 {
1027 return Err(EventError::OpenLifecycle);
1028 }
1029 self.run_status
1030 .insert(run_id.clone(), RunState::Terminal(*status));
1031 self.active_run_id = None;
1032 self.waiting_interaction_id = None;
1033 }
1034 Event::UserMessage { run_id, content } if include_content => {
1035 self.messages.push(ProjectedMessage {
1036 role: "user",
1037 run_id: run_id.clone(),
1038 content: content.clone(),
1039 })
1040 }
1041 Event::AssistantMessage {
1042 run_id, content, ..
1043 } if include_content => self.messages.push(ProjectedMessage {
1044 role: "assistant",
1045 run_id: run_id.clone(),
1046 content: content.clone(),
1047 }),
1048 Event::ContextInjected {
1049 run_id,
1050 step,
1051 contribution_id,
1052 source,
1053 version,
1054 authority,
1055 form,
1056 content,
1057 } => {
1058 self.require_active(run_id)?;
1059 if !self.open_steps.contains(step) {
1060 return Err(EventError::LifecycleMismatch);
1061 }
1062 if include_content {
1063 self.injected_context.push(ProjectedContext {
1064 run_id: run_id.clone(),
1065 step: *step,
1066 contribution_id: contribution_id.clone(),
1067 source: source.clone(),
1068 version: version.clone(),
1069 authority: authority.clone(),
1070 form: form.clone(),
1071 content: content.clone(),
1072 });
1073 }
1074 }
1075 Event::InputQueued {
1076 input_id,
1077 run_id,
1078 mode,
1079 ..
1080 } => {
1081 if self.metadata.archived {
1082 return Err(EventError::SessionArchived);
1083 }
1084 if *mode != DeliveryMode::Followup {
1085 self.require_active(run_id)?;
1086 }
1087 if self
1088 .queued_inputs
1089 .insert(input_id.clone(), (run_id.clone(), *mode))
1090 .is_some()
1091 {
1092 return Err(EventError::DuplicateInput(input_id.clone()));
1093 }
1094 }
1095 Event::InputClaimed { input_id, run_id } => {
1096 if self
1097 .queued_inputs
1098 .remove(input_id)
1099 .map(|value| value.0)
1100 .as_ref()
1101 != Some(run_id)
1102 || self
1103 .claimed_inputs
1104 .insert(input_id.clone(), run_id.clone())
1105 .is_some()
1106 {
1107 return Err(EventError::UnqueuedInput(input_id.clone()));
1108 }
1109 }
1110 Event::InputCancelled {
1111 input_id,
1112 run_id,
1113 error_code,
1114 } => {
1115 if self
1116 .queued_inputs
1117 .remove(input_id)
1118 .map(|value| value.0)
1119 .as_ref()
1120 != Some(run_id)
1121 {
1122 return Err(EventError::UnqueuedInput(input_id.clone()));
1123 }
1124 self.run_status.insert(
1125 run_id.clone(),
1126 RunState::Terminal(if error_code == "cancelled" {
1127 RunStatus::Cancelled
1128 } else {
1129 RunStatus::Failed
1130 }),
1131 );
1132 }
1133 Event::TurnStarted { run_id, turn } => {
1134 self.require_active(run_id)?;
1135 if self.open_turn.replace(*turn).is_some() {
1136 return Err(EventError::ConcurrentTurn);
1137 }
1138 }
1139 Event::TurnFinished { run_id, turn } => {
1140 self.require_active(run_id)?;
1141 if self.open_turn != Some(*turn)
1142 || !self.open_steps.is_empty()
1143 || !self.open_tool_calls.is_empty()
1144 {
1145 return Err(EventError::LifecycleMismatch);
1146 }
1147 self.open_turn = None;
1148 }
1149 Event::StepStarted { run_id, step } => {
1150 self.require_active(run_id)?;
1151 if self.open_turn.is_none()
1152 || !self.open_steps.is_empty()
1153 || *step != self.next_step
1154 || !self.open_steps.insert(*step)
1155 {
1156 return Err(EventError::ConcurrentStep);
1157 }
1158 }
1159 Event::StepFinished { run_id, step } => {
1160 self.require_active(run_id)?;
1161 if self
1162 .open_tool_calls
1163 .values()
1164 .any(|call| call.run_id == *run_id && call.step == *step)
1165 || !self.open_steps.remove(step)
1166 {
1167 return Err(EventError::LifecycleMismatch);
1168 }
1169 self.next_step = step.saturating_add(1);
1170 }
1171 Event::ToolCall {
1172 run_id,
1173 step,
1174 call_id,
1175 tool,
1176 arguments,
1177 } => {
1178 self.require_active(run_id)?;
1179 if !self.open_steps.contains(step) {
1180 return Err(EventError::LifecycleMismatch);
1181 }
1182 if !self.seen_tool_calls.insert(call_id.clone())
1183 || self
1184 .open_tool_calls
1185 .insert(
1186 call_id.clone(),
1187 OpenToolCall {
1188 run_id: run_id.clone(),
1189 step: *step,
1190 tool: tool.clone(),
1191 arguments: arguments.clone(),
1192 source_event_seq: envelope.seq,
1193 },
1194 )
1195 .is_some()
1196 {
1197 return Err(EventError::DuplicateToolCall(call_id.clone()));
1198 }
1199 }
1200 Event::ToolResult {
1201 run_id,
1202 step,
1203 call_id,
1204 ..
1205 } => {
1206 self.require_active(run_id)?;
1207 if !self.open_steps.contains(step) {
1208 return Err(EventError::LifecycleMismatch);
1209 }
1210 let Some(call) = self.open_tool_calls.get(call_id) else {
1211 return Err(EventError::OrphanToolResult(call_id.clone()));
1212 };
1213 if call.run_id != *run_id || call.step != *step {
1214 return Err(EventError::ToolResultMismatch(call_id.clone()));
1215 }
1216 self.open_tool_calls.remove(call_id);
1217 self.started_tool_calls.remove(call_id);
1218 }
1219 Event::ToolExecutionStarted {
1220 run_id,
1221 step,
1222 call_id,
1223 metering,
1224 reserved_cost_units,
1225 } => {
1226 self.require_active(run_id)?;
1227 let Some(call) = self.open_tool_calls.get(call_id) else {
1228 return Err(EventError::OrphanToolResult(call_id.clone()));
1229 };
1230 if call.run_id != *run_id
1231 || call.step != *step
1232 || self.started_tool_calls.contains(call_id)
1233 {
1234 return Err(EventError::ToolResultMismatch(call_id.clone()));
1235 }
1236 if let Some(details) = metering {
1237 details.validate()?;
1238 if !matches!(details, MeteringDetails::Tool { call_id: id, name, source: MeteringSource::Estimated, outcome: MeteringOutcome::Unknown } if id == call_id && name == &call.tool)
1239 {
1240 return Err(EventError::InvalidMetering);
1241 }
1242 }
1243 self.started_tool_calls.insert(call_id.clone());
1244 self.pending_usage_operations.insert(
1245 (run_id.clone(), format!("tool:{call_id}")),
1246 (0, 0, *reserved_cost_units, metering.clone()),
1247 );
1248 }
1249 Event::ModelRequestPrepared {
1250 run_id,
1251 operation_id,
1252 reserved_prompt_tokens,
1253 reserved_completion_tokens,
1254 metering,
1255 ..
1256 } if !operation_id.is_empty() => {
1257 self.require_active(run_id)?;
1258 let key = (run_id.clone(), operation_id.clone());
1259 if !self.usage_operations.contains_key(&key) {
1260 if let Some(details) = metering {
1261 details.validate()?;
1262 if !matches!(
1263 details,
1264 MeteringDetails::Model {
1265 source: MeteringSource::Estimated,
1266 outcome: MeteringOutcome::Unknown,
1267 ..
1268 }
1269 ) {
1270 return Err(EventError::InvalidMetering);
1271 }
1272 }
1273 let reservation = (
1274 *reserved_prompt_tokens,
1275 *reserved_completion_tokens,
1276 0,
1277 metering.clone(),
1278 );
1279 match self.pending_usage_operations.get(&key) {
1280 Some(existing) if *existing != reservation => {
1281 return Err(EventError::UsageConflict(operation_id.clone()));
1282 }
1283 Some(_) => {}
1284 None => {
1285 self.pending_usage_operations.insert(key, reservation);
1286 }
1287 }
1288 }
1289 }
1290 Event::UsageRecorded {
1291 run_id,
1292 operation_id,
1293 prompt_tokens,
1294 completion_tokens,
1295 cost_units,
1296 metering,
1297 } => {
1298 self.require_active(run_id)?;
1299 if let Some(details) = metering {
1300 details.validate()?;
1301 }
1302 let key = (run_id.clone(), operation_id.clone());
1303 if let Some((_, _, _, Some(prepared))) = self.pending_usage_operations.get(&key) {
1304 if !metering
1305 .as_ref()
1306 .is_some_and(|details| details.completes(prepared))
1307 {
1308 return Err(EventError::UsageConflict(operation_id.clone()));
1309 }
1310 }
1311 let usage = (
1312 *prompt_tokens,
1313 *completion_tokens,
1314 *cost_units,
1315 metering.clone(),
1316 );
1317 match self.usage_operations.get(&key) {
1318 Some(existing) if *existing != usage => {
1319 return Err(EventError::UsageConflict(operation_id.clone()));
1320 }
1321 Some(_) => {}
1322 None => {
1323 self.pending_usage_operations.remove(&key);
1324 self.usage_operations.insert(key, usage);
1325 }
1326 }
1327 }
1328 Event::CompactionStarted {
1329 run_id,
1330 compaction_id,
1331 source_through_seq,
1332 } => {
1333 self.require_active(run_id)?;
1334 if self.open_compaction.is_some() {
1335 return Err(EventError::ConcurrentCompaction);
1336 }
1337 self.open_compaction = Some((compaction_id.clone(), *source_through_seq));
1338 }
1339 Event::CompactionFinished {
1340 run_id,
1341 compaction_id,
1342 ..
1343 } => {
1344 self.require_active(run_id)?;
1345 if self.open_compaction.as_ref().map(|value| value.0.as_str())
1346 != Some(compaction_id.as_str())
1347 {
1348 return Err(EventError::CompactionMismatch);
1349 }
1350 self.open_compaction = None;
1351 }
1352 Event::SummaryReplaced { summary, .. } => self.summary = Some(summary.clone()),
1353 Event::Opaque {
1354 event_type,
1355 ignorable: false,
1356 ..
1357 } => return Err(EventError::UnknownRequired(event_type.clone())),
1358 _ => {}
1359 }
1360 self.last_seq = envelope.seq;
1361 Ok(())
1362 }
1363
1364 fn require_active(&self, run_id: &RunId) -> Result<(), EventError> {
1365 if self.active_run_id.as_ref() == Some(run_id) {
1366 Ok(())
1367 } else {
1368 Err(EventError::RunMismatch)
1369 }
1370 }
1371
1372 pub fn usage_for(&self, run_id: &str) -> (u64, u64) {
1374 self.usage_operations
1375 .iter()
1376 .filter(|((recorded_run_id, _), _)| recorded_run_id.as_str() == run_id)
1377 .fold((0, 0), |total, (_, usage)| {
1378 (total.0 + usage.0, total.1 + usage.1)
1379 })
1380 }
1381
1382 pub fn usage_charge_for(&self, run_id: &str) -> UsageCharge {
1384 let sum = |pending: bool| {
1385 let operations = if pending {
1386 &self.pending_usage_operations
1387 } else {
1388 &self.usage_operations
1389 };
1390 operations
1391 .iter()
1392 .filter(|((recorded_run_id, _), _)| recorded_run_id.as_str() == run_id)
1393 .fold(UsageCharge::default(), |mut total, (_, usage)| {
1394 total.prompt_tokens = total.prompt_tokens.saturating_add(usage.0);
1395 total.completion_tokens = total.completion_tokens.saturating_add(usage.1);
1396 total.cost_units = total.cost_units.saturating_add(usage.2);
1397 total
1398 })
1399 };
1400 let recorded = sum(false);
1401 let pending = sum(true);
1402 UsageCharge {
1403 prompt_tokens: recorded.prompt_tokens.saturating_add(pending.prompt_tokens),
1404 completion_tokens: recorded
1405 .completion_tokens
1406 .saturating_add(pending.completion_tokens),
1407 cost_units: recorded.cost_units.saturating_add(pending.cost_units),
1408 }
1409 }
1410
1411 pub fn billable_units_for(&self, run_id: &str) -> u64 {
1413 self.usage_charge_for(run_id).total()
1414 }
1415}
1416
1417#[derive(Debug, thiserror::Error, PartialEq, Eq)]
1419pub enum EventError {
1420 #[error("invalid Session metadata or metadata version conflict")]
1422 InvalidSessionMetadata,
1423 #[error("invalid metering attribution")]
1425 InvalidMetering,
1426 #[error("Session is archived")]
1428 SessionArchived,
1429 #[error("unsupported session event format version {0}")]
1431 UnsupportedFormat(u32),
1432 #[error("unknown required session event type {0}")]
1434 UnknownRequired(String),
1435 #[error("event sequence mismatch: expected {expected}, got {actual}")]
1437 Sequence {
1438 expected: u64,
1440 actual: u64,
1442 },
1443 #[error("event belongs to another session")]
1445 SessionMismatch,
1446 #[error("session creation must be the first and only creation event")]
1448 DuplicateSession,
1449 #[error("session already has an active run")]
1451 ConcurrentRun,
1452 #[error("session is closed")]
1454 SessionClosed,
1455 #[error("event does not match the active run")]
1457 RunMismatch,
1458 #[error("interaction does not match the waiting run")]
1460 InteractionMismatch,
1461 #[error("run cannot finish with an open turn, step or tool call")]
1463 OpenLifecycle,
1464 #[error("input was queued twice: {0}")]
1466 DuplicateInput(InputId),
1467 #[error("input was claimed before it was queued: {0}")]
1469 UnqueuedInput(InputId),
1470 #[error("run started from an unclaimed input: {0}")]
1472 UnclaimedInput(InputId),
1473 #[error("session already has an active turn")]
1475 ConcurrentTurn,
1476 #[error("turn already has this active step")]
1478 ConcurrentStep,
1479 #[error("turn or step lifecycle does not pair")]
1481 LifecycleMismatch,
1482 #[error("duplicate tool call {0}")]
1484 DuplicateToolCall(ToolCallId),
1485 #[error("tool result has no matching call {0}")]
1487 OrphanToolResult(ToolCallId),
1488 #[error("tool result does not match the call run and step: {0}")]
1490 ToolResultMismatch(ToolCallId),
1491 #[error("usage operation was recorded with different totals: {0}")]
1493 UsageConflict(String),
1494 #[error("session already has an active compaction")]
1496 ConcurrentCompaction,
1497 #[error("compaction lifecycle does not pair")]
1499 CompactionMismatch,
1500 #[error("event store conflict: {0}")]
1502 Conflict(String),
1503 #[error("event store unavailable: {0}")]
1505 Unavailable(String),
1506}
1507
1508#[async_trait]
1510pub trait SessionEventStore: Send + Sync {
1511 async fn append(
1513 &self,
1514 tenant_id: &str,
1515 session_id: &str,
1516 expected_seq: u64,
1517 events: Vec<Event>,
1518 ) -> Result<Vec<SessionEvent>, EventError>;
1519 async fn load(
1521 &self,
1522 tenant_id: &str,
1523 session_id: &str,
1524 after_seq: u64,
1525 ) -> Result<Vec<SessionEvent>, EventError>;
1526}
1527
1528pub fn text(value: impl Into<String>) -> Vec<ContentBlock> {
1530 vec![ContentBlock::Text { text: value.into() }]
1531}
1532
1533pub fn recovery_events(projection: &SessionProjection) -> Vec<Event> {
1535 failure_events(projection, "worker_restarted")
1536}
1537
1538pub fn failure_events(projection: &SessionProjection, error_code: &str) -> Vec<Event> {
1540 termination_events(projection, RunStatus::Failed, error_code)
1541}
1542
1543pub fn cancel_events(projection: &SessionProjection) -> Vec<Event> {
1545 termination_events(projection, RunStatus::Cancelled, "cancelled")
1546}
1547
1548pub fn session_deletion_events(projection: &SessionProjection, reason: &str) -> Vec<Event> {
1550 let active = projection.active_run_id.as_deref();
1551 let mut events = cancel_events(projection);
1552 events.extend(
1553 projection
1554 .queued_inputs
1555 .iter()
1556 .filter(|(_, (run_id, _))| Some(run_id.as_str()) != active)
1557 .map(|(input_id, (run_id, _))| Event::InputCancelled {
1558 input_id: input_id.clone(),
1559 run_id: run_id.clone(),
1560 error_code: "cancelled".into(),
1561 }),
1562 );
1563 events.push(Event::SessionDeleted {
1564 reason: reason.into(),
1565 });
1566 events
1567}
1568
1569fn termination_events(
1570 projection: &SessionProjection,
1571 status: RunStatus,
1572 error_code: &str,
1573) -> Vec<Event> {
1574 let Some(run_id) = &projection.active_run_id else {
1575 return Vec::new();
1576 };
1577 let mut events = projection
1578 .queued_inputs
1579 .iter()
1580 .filter(|(_, (target_run_id, _))| target_run_id == run_id)
1581 .map(|(input_id, _)| Event::InputCancelled {
1582 input_id: input_id.clone(),
1583 run_id: run_id.clone(),
1584 error_code: error_code.into(),
1585 })
1586 .collect::<Vec<_>>();
1587 events.extend(
1588 projection
1589 .open_tool_calls
1590 .iter()
1591 .map(|(call_id, call)| Event::ToolResult {
1592 run_id: run_id.clone(),
1593 step: call.step,
1594 call_id: call_id.clone(),
1595 result: termination_result(error_code),
1596 is_error: true,
1597 }),
1598 );
1599 if let Some((compaction_id, _)) = &projection.open_compaction {
1600 events.push(Event::CompactionFinished {
1601 run_id: run_id.clone(),
1602 compaction_id: compaction_id.clone(),
1603 status: "failed".into(),
1604 error: Some(error_code.into()),
1605 });
1606 }
1607 events.extend(
1608 projection
1609 .open_steps
1610 .iter()
1611 .map(|step| Event::StepFinished {
1612 run_id: run_id.clone(),
1613 step: *step,
1614 }),
1615 );
1616 if let Some(turn) = projection.open_turn {
1617 events.push(Event::TurnFinished {
1618 run_id: run_id.clone(),
1619 turn,
1620 });
1621 }
1622 events.push(Event::RunFinished {
1623 run_id: run_id.clone(),
1624 status,
1625 error_code: Some(error_code.into()),
1626 });
1627 events
1628}
1629
1630fn termination_result(error_code: &str) -> Value {
1631 if error_code == "tool_outcome_unknown" || error_code == "worker_restarted" {
1632 serde_json::json!({
1633 "error": error_code,
1634 "guidance": "The tool outcome is unknown. Verify external state before retrying any operation with side effects; ask the user when verification is unavailable."
1635 })
1636 } else {
1637 serde_json::json!({"error":error_code})
1638 }
1639}
1640
1641#[cfg(test)]
1642mod tests {
1643 use super::*;
1644
1645 fn event(seq: u64, event: Event) -> SessionEvent {
1646 SessionEvent {
1647 session_id: "s".parse().unwrap(),
1648 seq,
1649 occurred_at: Utc::now(),
1650 event,
1651 }
1652 }
1653
1654 #[test]
1655 fn replay_enforces_single_run_and_tool_pairs() {
1656 let events = vec![
1657 event(
1658 1,
1659 Event::SessionCreated {
1660 profile_revision_id: "p1".parse().unwrap(),
1661 },
1662 ),
1663 event(
1664 2,
1665 Event::InputQueued {
1666 input_id: "i1".parse().unwrap(),
1667 run_id: "r1".parse().unwrap(),
1668 mode: DeliveryMode::Followup,
1669 content: text("hi"),
1670 explicit_skill: None,
1671 },
1672 ),
1673 event(
1674 3,
1675 Event::InputClaimed {
1676 input_id: "i1".parse().unwrap(),
1677 run_id: "r1".parse().unwrap(),
1678 },
1679 ),
1680 event(
1681 4,
1682 Event::RunStarted {
1683 run_id: "r1".parse().unwrap(),
1684 input_id: "i1".parse().unwrap(),
1685 },
1686 ),
1687 event(
1688 5,
1689 Event::TurnStarted {
1690 run_id: "r1".parse().unwrap(),
1691 turn: 1,
1692 },
1693 ),
1694 event(
1695 6,
1696 Event::StepStarted {
1697 run_id: "r1".parse().unwrap(),
1698 step: 1,
1699 },
1700 ),
1701 event(
1702 7,
1703 Event::ToolCall {
1704 run_id: "r1".parse().unwrap(),
1705 step: 1,
1706 call_id: "c1".parse().unwrap(),
1707 tool: "echo".into(),
1708 arguments: serde_json::json!({"x":1}),
1709 },
1710 ),
1711 event(
1712 8,
1713 Event::ToolResult {
1714 run_id: "r1".parse().unwrap(),
1715 step: 1,
1716 call_id: "c1".parse().unwrap(),
1717 result: serde_json::json!({"x":1}),
1718 is_error: false,
1719 },
1720 ),
1721 event(
1722 9,
1723 Event::StepFinished {
1724 run_id: "r1".parse().unwrap(),
1725 step: 1,
1726 },
1727 ),
1728 event(
1729 10,
1730 Event::TurnFinished {
1731 run_id: "r1".parse().unwrap(),
1732 turn: 1,
1733 },
1734 ),
1735 event(
1736 11,
1737 Event::RunFinished {
1738 run_id: "r1".parse().unwrap(),
1739 status: RunStatus::Completed,
1740 error_code: None,
1741 },
1742 ),
1743 ];
1744 let projection = SessionProjection::replay(&events).unwrap();
1745 assert_eq!(projection.last_seq, 11);
1746 assert!(projection.active_run_id.is_none());
1747 }
1748
1749 #[test]
1750 fn replay_rejects_orphan_tool_result() {
1751 let events = vec![
1752 event(
1753 1,
1754 Event::SessionCreated {
1755 profile_revision_id: "p1".parse().unwrap(),
1756 },
1757 ),
1758 event(
1759 2,
1760 Event::InputQueued {
1761 input_id: "i1".parse().unwrap(),
1762 run_id: "r1".parse().unwrap(),
1763 mode: DeliveryMode::Followup,
1764 content: text("hi"),
1765 explicit_skill: None,
1766 },
1767 ),
1768 event(
1769 3,
1770 Event::InputClaimed {
1771 input_id: "i1".parse().unwrap(),
1772 run_id: "r1".parse().unwrap(),
1773 },
1774 ),
1775 event(
1776 4,
1777 Event::RunStarted {
1778 run_id: "r1".parse().unwrap(),
1779 input_id: "i1".parse().unwrap(),
1780 },
1781 ),
1782 event(
1783 5,
1784 Event::TurnStarted {
1785 run_id: "r1".parse().unwrap(),
1786 turn: 1,
1787 },
1788 ),
1789 event(
1790 6,
1791 Event::StepStarted {
1792 run_id: "r1".parse().unwrap(),
1793 step: 1,
1794 },
1795 ),
1796 event(
1797 7,
1798 Event::ToolResult {
1799 run_id: "r1".parse().unwrap(),
1800 step: 1,
1801 call_id: "missing".parse().unwrap(),
1802 result: Value::Null,
1803 is_error: true,
1804 },
1805 ),
1806 ];
1807 assert_eq!(
1808 SessionProjection::replay(&events).unwrap_err(),
1809 EventError::OrphanToolResult("missing".parse().unwrap())
1810 );
1811 }
1812
1813 #[test]
1814 fn queued_input_failure_is_not_projected_as_cancellation() {
1815 let events = vec![
1816 event(
1817 1,
1818 Event::SessionCreated {
1819 profile_revision_id: "p1".parse().unwrap(),
1820 },
1821 ),
1822 event(
1823 2,
1824 Event::InputQueued {
1825 input_id: "i1".parse().unwrap(),
1826 run_id: "r1".parse().unwrap(),
1827 mode: DeliveryMode::Followup,
1828 content: text("hi"),
1829 explicit_skill: None,
1830 },
1831 ),
1832 event(
1833 3,
1834 Event::InputCancelled {
1835 input_id: "i1".parse().unwrap(),
1836 run_id: "r1".parse().unwrap(),
1837 error_code: "profile_not_found".into(),
1838 },
1839 ),
1840 ];
1841 let projection = SessionProjection::replay(&events).unwrap();
1842 assert_eq!(
1843 projection.run_status.get("r1").map(|state| state.as_str()),
1844 Some("failed")
1845 );
1846 }
1847
1848 #[test]
1849 fn session_deletion_closes_active_and_queued_runs_before_tombstone() {
1850 let mut events = vec![
1851 event(
1852 1,
1853 Event::SessionCreated {
1854 profile_revision_id: "p1".parse().unwrap(),
1855 },
1856 ),
1857 event(
1858 2,
1859 Event::InputQueued {
1860 input_id: "i1".parse().unwrap(),
1861 run_id: "r1".parse().unwrap(),
1862 mode: DeliveryMode::Followup,
1863 content: text("start"),
1864 explicit_skill: None,
1865 },
1866 ),
1867 event(
1868 3,
1869 Event::InputClaimed {
1870 input_id: "i1".parse().unwrap(),
1871 run_id: "r1".parse().unwrap(),
1872 },
1873 ),
1874 event(
1875 4,
1876 Event::RunStarted {
1877 run_id: "r1".parse().unwrap(),
1878 input_id: "i1".parse().unwrap(),
1879 },
1880 ),
1881 event(
1882 5,
1883 Event::TurnStarted {
1884 run_id: "r1".parse().unwrap(),
1885 turn: 1,
1886 },
1887 ),
1888 event(
1889 6,
1890 Event::InputQueued {
1891 input_id: "i2".parse().unwrap(),
1892 run_id: "r2".parse().unwrap(),
1893 mode: DeliveryMode::Followup,
1894 content: text("later"),
1895 explicit_skill: None,
1896 },
1897 ),
1898 ];
1899 let projection = SessionProjection::replay(&events).unwrap();
1900 for event_value in session_deletion_events(&projection, "api_deleted") {
1901 let seq = events.len() as u64 + 1;
1902 events.push(event(seq, event_value));
1903 }
1904 let deleted = SessionProjection::replay(&events).unwrap();
1905 assert!(deleted.deleted);
1906 assert_eq!(
1907 deleted.run_status.get("r1").map(|state| state.as_str()),
1908 Some("cancelled")
1909 );
1910 assert_eq!(
1911 deleted.run_status.get("r2").map(|state| state.as_str()),
1912 Some("cancelled")
1913 );
1914 assert!(events.iter().any(|event| matches!(
1915 &event.event,
1916 Event::RunFinished { run_id, status: RunStatus::Cancelled, .. } if run_id == "r1"
1917 )));
1918 assert!(matches!(
1919 events.last().map(|event| &event.event),
1920 Some(Event::SessionDeleted { .. })
1921 ));
1922 }
1923
1924 #[test]
1925 fn usage_operations_are_idempotent_and_conflicts_fail_replay() {
1926 let mut events = vec![
1927 event(
1928 1,
1929 Event::SessionCreated {
1930 profile_revision_id: "p1".parse().unwrap(),
1931 },
1932 ),
1933 event(
1934 2,
1935 Event::InputQueued {
1936 input_id: "i1".parse().unwrap(),
1937 run_id: "r1".parse().unwrap(),
1938 mode: DeliveryMode::Followup,
1939 content: text("hi"),
1940 explicit_skill: None,
1941 },
1942 ),
1943 event(
1944 3,
1945 Event::InputClaimed {
1946 input_id: "i1".parse().unwrap(),
1947 run_id: "r1".parse().unwrap(),
1948 },
1949 ),
1950 event(
1951 4,
1952 Event::RunStarted {
1953 run_id: "r1".parse().unwrap(),
1954 input_id: "i1".parse().unwrap(),
1955 },
1956 ),
1957 event(
1958 5,
1959 Event::UsageRecorded {
1960 metering: None,
1961 run_id: "r1".parse().unwrap(),
1962 operation_id: "model:1:attempt:1".into(),
1963 prompt_tokens: 7,
1964 completion_tokens: 3,
1965 cost_units: 5,
1966 },
1967 ),
1968 event(
1969 6,
1970 Event::UsageRecorded {
1971 metering: None,
1972 run_id: "r1".parse().unwrap(),
1973 operation_id: "model:1:attempt:1".into(),
1974 prompt_tokens: 7,
1975 completion_tokens: 3,
1976 cost_units: 5,
1977 },
1978 ),
1979 ];
1980 assert_eq!(
1981 SessionProjection::replay(&events).unwrap().usage_for("r1"),
1982 (7, 3)
1983 );
1984 assert_eq!(
1985 SessionProjection::replay(&events)
1986 .unwrap()
1987 .billable_units_for("r1"),
1988 15
1989 );
1990 assert_eq!(
1991 SessionProjection::replay(&events)
1992 .unwrap()
1993 .usage_charge_for("r1"),
1994 UsageCharge {
1995 prompt_tokens: 7,
1996 completion_tokens: 3,
1997 cost_units: 5,
1998 }
1999 );
2000
2001 let mut attributed = events.clone();
2002 let details = MeteringDetails::Model {
2003 model: "m".into(),
2004 provider: Some("p".into()),
2005 provider_attempt_id: "r1:model:1:attempt:1".parse().unwrap(),
2006 source: MeteringSource::Estimated,
2007 outcome: MeteringOutcome::Unknown,
2008 };
2009 for entry in &mut attributed[4..] {
2010 if let Event::UsageRecorded { metering, .. } = &mut entry.event {
2011 *metering = Some(details.clone());
2012 }
2013 }
2014 assert_eq!(
2015 SessionProjection::replay(&attributed)
2016 .unwrap()
2017 .billable_units_for("r1"),
2018 15
2019 );
2020 if let Event::UsageRecorded {
2021 metering: Some(MeteringDetails::Model { source, .. }),
2022 ..
2023 } = &mut attributed[5].event
2024 {
2025 *source = MeteringSource::Reported;
2026 }
2027 assert!(matches!(
2028 SessionProjection::replay(&attributed),
2029 Err(EventError::UsageConflict(_))
2030 ));
2031 let mut prepared = attributed[..4].to_vec();
2032 prepared.push(event(
2033 5,
2034 Event::ModelRequestPrepared {
2035 metering: Some(details.clone()),
2036 run_id: "r1".parse().unwrap(),
2037 step: 1,
2038 attempt: 1,
2039 provider_attempt_id: "r1:model:1:attempt:1".into(),
2040 operation_id: "model:1:attempt:1".into(),
2041 reserved_prompt_tokens: 10,
2042 reserved_completion_tokens: 20,
2043 request: Value::Null,
2044 prompt_sections: Value::Null,
2045 },
2046 ));
2047 prepared.push(attributed[5].clone());
2048 assert_eq!(
2049 SessionProjection::replay(&prepared)
2050 .unwrap()
2051 .billable_units_for("r1"),
2052 15
2053 );
2054 if let Event::UsageRecorded {
2055 metering: Some(MeteringDetails::Model { model, .. }),
2056 ..
2057 } = &mut prepared[5].event
2058 {
2059 *model = "wrong".into();
2060 }
2061 assert!(matches!(
2062 SessionProjection::replay(&prepared),
2063 Err(EventError::UsageConflict(_))
2064 ));
2065 if let Event::UsageRecorded { metering, .. } = &mut prepared[5].event {
2066 *metering = None;
2067 }
2068 assert!(matches!(
2069 SessionProjection::replay(&prepared),
2070 Err(EventError::UsageConflict(_))
2071 ));
2072 let serialized = serde_json::to_value(&events[4]).unwrap();
2073 assert!(serialized["event"].get("metering").is_none());
2074 assert_eq!(
2075 serde_json::from_value::<SessionEvent>(serialized).unwrap(),
2076 events[4]
2077 );
2078
2079 events.push(event(
2080 7,
2081 Event::UsageRecorded {
2082 metering: None,
2083 run_id: "r1".parse().unwrap(),
2084 operation_id: "model:1:attempt:1".into(),
2085 prompt_tokens: 8,
2086 completion_tokens: 3,
2087 cost_units: 0,
2088 },
2089 ));
2090 assert_eq!(
2091 SessionProjection::replay(&events).unwrap_err(),
2092 EventError::UsageConflict("model:1:attempt:1".into())
2093 );
2094 }
2095
2096 #[test]
2097 fn unresolved_provider_attempt_is_conservatively_billable_and_reconcilable() {
2098 let mut events = vec![
2099 event(
2100 1,
2101 Event::SessionCreated {
2102 profile_revision_id: "p1".parse().unwrap(),
2103 },
2104 ),
2105 event(
2106 2,
2107 Event::InputQueued {
2108 input_id: "i1".parse().unwrap(),
2109 run_id: "r1".parse().unwrap(),
2110 mode: DeliveryMode::Followup,
2111 content: text("hi"),
2112 explicit_skill: None,
2113 },
2114 ),
2115 event(
2116 3,
2117 Event::InputClaimed {
2118 input_id: "i1".parse().unwrap(),
2119 run_id: "r1".parse().unwrap(),
2120 },
2121 ),
2122 event(
2123 4,
2124 Event::RunStarted {
2125 run_id: "r1".parse().unwrap(),
2126 input_id: "i1".parse().unwrap(),
2127 },
2128 ),
2129 event(
2130 5,
2131 Event::ModelRequestPrepared {
2132 metering: None,
2133 run_id: "r1".parse().unwrap(),
2134 step: 1,
2135 attempt: 1,
2136 provider_attempt_id: "r1:model:1:attempt:1".into(),
2137 operation_id: "model:1:attempt:1".into(),
2138 reserved_prompt_tokens: 7,
2139 reserved_completion_tokens: 11,
2140 request: Value::Null,
2141 prompt_sections: Value::Null,
2142 },
2143 ),
2144 ];
2145 assert_eq!(
2146 SessionProjection::replay(&events)
2147 .unwrap()
2148 .billable_units_for("r1"),
2149 18
2150 );
2151 events.push(event(
2152 6,
2153 Event::UsageRecorded {
2154 metering: None,
2155 run_id: "r1".parse().unwrap(),
2156 operation_id: "model:1:attempt:1".into(),
2157 prompt_tokens: 6,
2158 completion_tokens: 2,
2159 cost_units: 0,
2160 },
2161 ));
2162 assert_eq!(
2163 SessionProjection::replay(&events)
2164 .unwrap()
2165 .billable_units_for("r1"),
2166 8
2167 );
2168 }
2169
2170 #[test]
2171 fn a3_incremental_facts_preserve_every_lifecycle_and_usage_invariant() {
2172 let mut events = open_step_events();
2173 for fact in [
2174 Event::UserMessage {
2175 run_id: "r1".parse().unwrap(),
2176 content: text("history".repeat(10_000)),
2177 },
2178 Event::UsageRecorded {
2179 metering: None,
2180 run_id: "r1".parse().unwrap(),
2181 operation_id: "attempt".into(),
2182 prompt_tokens: 5,
2183 completion_tokens: 3,
2184 cost_units: 0,
2185 },
2186 Event::ToolCall {
2187 run_id: "r1".parse().unwrap(),
2188 step: 1,
2189 call_id: "pending".parse().unwrap(),
2190 tool: "echo".into(),
2191 arguments: Value::Null,
2192 },
2193 ] {
2194 events.push(event(events.len() as u64 + 1, fact));
2195 }
2196 let mut full = SessionProjection::replay(&events).unwrap();
2197 assert!(!full.messages.is_empty());
2198 full.messages.clear();
2199 full.injected_context.clear();
2200 let mut folded = SessionProjection::default();
2201 for page in events.chunks(3) {
2202 for fact in page {
2203 folded.apply_facts(fact).unwrap();
2204 }
2205 assert!(folded.messages.is_empty());
2206 }
2207 assert_eq!(folded, full);
2208 assert_eq!(SessionProjection::replay_facts(&events).unwrap(), full);
2209 assert_eq!(cancel_events(&folded), cancel_events(&full));
2210 assert_eq!(folded.usage_for("r1"), (5, 3));
2211 let invalid = event(
2212 events.len() as u64 + 1,
2213 Event::ToolResult {
2214 run_id: "r1".parse().unwrap(),
2215 step: 1,
2216 call_id: "missing".parse().unwrap(),
2217 result: Value::Null,
2218 is_error: true,
2219 },
2220 );
2221 assert_eq!(folded.apply_facts(&invalid), full.apply(&invalid));
2222 }
2223
2224 fn open_step_events() -> Vec<SessionEvent> {
2225 vec![
2226 event(
2227 1,
2228 Event::SessionCreated {
2229 profile_revision_id: "p1".parse().unwrap(),
2230 },
2231 ),
2232 event(
2233 2,
2234 Event::InputQueued {
2235 input_id: "i1".parse().unwrap(),
2236 run_id: "r1".parse().unwrap(),
2237 mode: DeliveryMode::Followup,
2238 content: text("hi"),
2239 explicit_skill: None,
2240 },
2241 ),
2242 event(
2243 3,
2244 Event::InputClaimed {
2245 input_id: "i1".parse().unwrap(),
2246 run_id: "r1".parse().unwrap(),
2247 },
2248 ),
2249 event(
2250 4,
2251 Event::RunStarted {
2252 run_id: "r1".parse().unwrap(),
2253 input_id: "i1".parse().unwrap(),
2254 },
2255 ),
2256 event(
2257 5,
2258 Event::TurnStarted {
2259 run_id: "r1".parse().unwrap(),
2260 turn: 1,
2261 },
2262 ),
2263 event(
2264 6,
2265 Event::StepStarted {
2266 run_id: "r1".parse().unwrap(),
2267 step: 1,
2268 },
2269 ),
2270 ]
2271 }
2272
2273 #[test]
2274 fn replay_rejects_overlapping_or_out_of_order_steps() {
2275 let mut overlapping = open_step_events();
2276 overlapping.push(event(
2277 7,
2278 Event::StepStarted {
2279 run_id: "r1".parse().unwrap(),
2280 step: 2,
2281 },
2282 ));
2283 assert_eq!(
2284 SessionProjection::replay(&overlapping).unwrap_err(),
2285 EventError::ConcurrentStep
2286 );
2287
2288 let mut skipped = open_step_events();
2289 skipped[5] = event(
2290 6,
2291 Event::StepStarted {
2292 run_id: "r1".parse().unwrap(),
2293 step: 2,
2294 },
2295 );
2296 assert_eq!(
2297 SessionProjection::replay(&skipped).unwrap_err(),
2298 EventError::ConcurrentStep
2299 );
2300 }
2301
2302 #[test]
2303 fn replay_rejects_cross_step_results_and_dangling_calls() {
2304 let mut events = open_step_events();
2305 events.push(event(
2306 7,
2307 Event::ToolCall {
2308 run_id: "r1".parse().unwrap(),
2309 step: 1,
2310 call_id: "c1".parse().unwrap(),
2311 tool: "echo".into(),
2312 arguments: Value::Null,
2313 },
2314 ));
2315 events.push(event(
2316 8,
2317 Event::ToolResult {
2318 run_id: "r1".parse().unwrap(),
2319 step: 2,
2320 call_id: "c1".parse().unwrap(),
2321 result: Value::Null,
2322 is_error: false,
2323 },
2324 ));
2325 assert_eq!(
2326 SessionProjection::replay(&events).unwrap_err(),
2327 EventError::LifecycleMismatch
2328 );
2329
2330 let mut dangling = open_step_events();
2331 dangling.push(event(
2332 7,
2333 Event::ToolCall {
2334 run_id: "r1".parse().unwrap(),
2335 step: 1,
2336 call_id: "c1".parse().unwrap(),
2337 tool: "echo".into(),
2338 arguments: Value::Null,
2339 },
2340 ));
2341 dangling.push(event(
2342 8,
2343 Event::StepFinished {
2344 run_id: "r1".parse().unwrap(),
2345 step: 1,
2346 },
2347 ));
2348 assert_eq!(
2349 SessionProjection::replay(&dangling).unwrap_err(),
2350 EventError::LifecycleMismatch
2351 );
2352 }
2353
2354 #[test]
2355 fn replay_rejects_reused_tool_call_ids() {
2356 let mut events = open_step_events();
2357 events.extend([
2358 event(
2359 7,
2360 Event::ToolCall {
2361 run_id: "r1".parse().unwrap(),
2362 step: 1,
2363 call_id: "c1".parse().unwrap(),
2364 tool: "echo".into(),
2365 arguments: Value::Null,
2366 },
2367 ),
2368 event(
2369 8,
2370 Event::ToolResult {
2371 run_id: "r1".parse().unwrap(),
2372 step: 1,
2373 call_id: "c1".parse().unwrap(),
2374 result: Value::Null,
2375 is_error: false,
2376 },
2377 ),
2378 event(
2379 9,
2380 Event::StepFinished {
2381 run_id: "r1".parse().unwrap(),
2382 step: 1,
2383 },
2384 ),
2385 event(
2386 10,
2387 Event::StepStarted {
2388 run_id: "r1".parse().unwrap(),
2389 step: 2,
2390 },
2391 ),
2392 event(
2393 11,
2394 Event::ToolCall {
2395 run_id: "r1".parse().unwrap(),
2396 step: 2,
2397 call_id: "c1".parse().unwrap(),
2398 tool: "echo".into(),
2399 arguments: Value::Null,
2400 },
2401 ),
2402 ]);
2403 assert_eq!(
2404 SessionProjection::replay(&events).unwrap_err(),
2405 EventError::DuplicateToolCall("c1".parse().unwrap())
2406 );
2407 }
2408
2409 #[test]
2410 fn replay_keeps_injected_context_without_changing_run_state() {
2411 let mut events = open_step_events();
2412 events.push(event(
2413 7,
2414 Event::ContextInjected {
2415 run_id: "r1".parse().unwrap(),
2416 step: 1,
2417 contribution_id: "memory:1".into(),
2418 source: "memory".into(),
2419 version: "v1".into(),
2420 authority: "tenant".into(),
2421 form: "message".into(),
2422 content: text("tenant context"),
2423 },
2424 ));
2425
2426 let projection = SessionProjection::replay(&events).unwrap();
2427 assert_eq!(projection.active_run_id.as_deref(), Some("r1"));
2428 assert_eq!(projection.open_steps, BTreeSet::from([1]));
2429 assert_eq!(
2430 projection.injected_context,
2431 vec![ProjectedContext {
2432 run_id: "r1".parse().unwrap(),
2433 step: 1,
2434 contribution_id: "memory:1".into(),
2435 source: "memory".into(),
2436 version: "v1".into(),
2437 authority: "tenant".into(),
2438 form: "message".into(),
2439 content: text("tenant context"),
2440 }]
2441 );
2442 }
2443
2444 #[test]
2445 fn replay_skips_unknown_ignorable_formats_and_rejects_required_events() {
2446 let optional = event(
2447 1,
2448 Event::Opaque {
2449 format_version: SESSION_EVENT_FORMAT_VERSION,
2450 event_type: "future_optional".into(),
2451 ignorable: true,
2452 payload: serde_json::json!({"answer": 42}),
2453 },
2454 );
2455 let projection = SessionProjection::replay(&[optional]).unwrap();
2456 assert_eq!(projection.last_seq, 1);
2457
2458 let required = event(
2459 1,
2460 Event::Opaque {
2461 format_version: SESSION_EVENT_FORMAT_VERSION,
2462 event_type: "future_required".into(),
2463 ignorable: false,
2464 payload: Value::Null,
2465 },
2466 );
2467 assert_eq!(
2468 SessionProjection::replay(&[required]).unwrap_err(),
2469 EventError::UnknownRequired("future_required".into())
2470 );
2471
2472 let unsupported = event(
2473 1,
2474 Event::Opaque {
2475 format_version: SESSION_EVENT_FORMAT_VERSION + 1,
2476 event_type: "future_optional".into(),
2477 ignorable: true,
2478 payload: Value::Null,
2479 },
2480 );
2481 assert_eq!(
2482 SessionProjection::replay(&[unsupported]).unwrap().last_seq,
2483 1
2484 );
2485
2486 let required_new_format = event(
2487 1,
2488 Event::Opaque {
2489 format_version: SESSION_EVENT_FORMAT_VERSION + 1,
2490 event_type: "future_required".into(),
2491 ignorable: false,
2492 payload: Value::Null,
2493 },
2494 );
2495 assert_eq!(
2496 SessionProjection::replay(&[required_new_format]).unwrap_err(),
2497 EventError::UnsupportedFormat(SESSION_EVENT_FORMAT_VERSION + 1)
2498 );
2499 }
2500}