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, Default, PartialEq)]
741pub struct SessionProjection {
742 pub metadata: SessionMetadata,
744 pub session_id: Option<SessionId>,
746 pub profile_revision_id: Option<ProfileRevisionId>,
749 pub deleted: bool,
751 pub last_seq: u64,
753 pub active_run_id: Option<RunId>,
755 pub waiting_interaction_id: Option<InteractionId>,
757 pub messages: Vec<ProjectedMessage>,
759 pub injected_context: Vec<ProjectedContext>,
761 pub run_status: BTreeMap<RunId, RunState>,
763 pub open_tool_calls: BTreeMap<ToolCallId, OpenToolCall>,
765 pub started_tool_calls: BTreeSet<ToolCallId>,
767 pub queued_inputs: BTreeMap<InputId, (RunId, DeliveryMode)>,
769 pub claimed_inputs: BTreeMap<InputId, RunId>,
771 pub open_turn: Option<u32>,
773 pub open_steps: BTreeSet<u32>,
775 pub next_step: u32,
777 pub open_compaction: Option<(String, u64)>,
779 pub summary: Option<String>,
781 corrected_usage: BTreeMap<(RunId, String), OperationUsage>,
782 metering_corrections: BTreeMap<af_context::MeteringCorrectionId, AppliedMeteringCorrection>,
783 usage_operations: BTreeMap<(RunId, String), RecordedUsage>,
784 pending_usage_operations: BTreeMap<(RunId, String), PreparedUsage>,
785 seen_tool_calls: BTreeSet<ToolCallId>,
786}
787
788type PreparedUsage = RecordedUsage;
789
790type RecordedUsage = (u64, u64, u64, Option<MeteringDetails>);
791
792#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
794#[serde(rename_all = "snake_case")]
795pub enum RunState {
796 Running,
798 WaitingForInput,
800 Terminal(RunStatus),
802}
803
804impl RunState {
805 pub const fn as_str(self) -> &'static str {
807 match self {
808 Self::Running => "running",
809 Self::WaitingForInput => "waiting_for_input",
810 Self::Terminal(status) => status.as_str(),
811 }
812 }
813
814 pub const fn is_terminal(self) -> bool {
816 matches!(self, Self::Terminal(_))
817 }
818}
819
820#[derive(Debug, Clone, PartialEq)]
822pub struct ProjectedContext {
823 pub run_id: RunId,
825 pub step: u32,
827 pub contribution_id: String,
829 pub source: String,
831 pub version: String,
833 pub authority: String,
835 pub form: String,
837 pub content: Vec<ContentBlock>,
839}
840
841#[derive(Debug, Clone, PartialEq)]
843pub struct ProjectedMessage {
844 pub role: &'static str,
846 pub run_id: RunId,
848 pub content: Vec<ContentBlock>,
850}
851
852#[derive(Debug, Clone, PartialEq)]
854pub struct OpenToolCall {
855 pub run_id: RunId,
857 pub step: u32,
859 pub tool: String,
861 pub arguments: Value,
863 pub source_event_seq: u64,
865}
866
867impl SessionProjection {
868 pub fn replay(events: &[SessionEvent]) -> Result<Self, EventError> {
870 let mut projection = Self::default();
871 for event in events {
872 projection.apply(event)?;
873 }
874 Ok(projection)
875 }
876
877 pub fn apply(&mut self, envelope: &SessionEvent) -> Result<(), EventError> {
879 self.apply_internal(envelope, true)
880 }
881
882 pub fn apply_facts(&mut self, envelope: &SessionEvent) -> Result<(), EventError> {
885 self.messages.clear();
886 self.injected_context.clear();
887 self.apply_internal(envelope, false)
888 }
889
890 pub fn replay_facts(events: &[SessionEvent]) -> Result<Self, EventError> {
892 let mut projection = Self::default();
893 for event in events {
894 projection.apply_facts(event)?;
895 }
896 Ok(projection)
897 }
898
899 fn apply_internal(
900 &mut self,
901 envelope: &SessionEvent,
902 include_content: bool,
903 ) -> Result<(), EventError> {
904 if envelope.seq != self.last_seq + 1 {
905 return Err(EventError::Sequence {
906 expected: self.last_seq + 1,
907 actual: envelope.seq,
908 });
909 }
910 let session_id = self
911 .session_id
912 .get_or_insert_with(|| envelope.session_id.clone());
913 if *session_id != envelope.session_id {
914 return Err(EventError::SessionMismatch);
915 }
916 if envelope.format_version() != SESSION_EVENT_FORMAT_VERSION {
917 if envelope.ignorable() {
918 self.last_seq = envelope.seq;
919 return Ok(());
920 }
921 return Err(EventError::UnsupportedFormat(envelope.format_version()));
922 }
923 if self.deleted && !matches!(envelope.event, Event::UsageCorrected { .. }) {
924 return Err(EventError::SessionClosed);
925 }
926 match &envelope.event {
927 Event::UsageCorrected {
928 correction,
929 actor_id,
930 } => self.apply_metering_correction(correction, actor_id, envelope.seq)?,
931 Event::SessionCreated {
932 profile_revision_id,
933 } => {
934 if envelope.seq != 1 || self.profile_revision_id.is_some() {
935 return Err(EventError::DuplicateSession);
936 }
937 self.profile_revision_id = Some(profile_revision_id.clone());
938 }
939 Event::SessionMetadataUpdated { metadata } => {
940 metadata.validate()?;
941 if self.profile_revision_id.is_none()
942 || self.metadata.version.checked_add(1) != Some(metadata.version)
943 {
944 return Err(EventError::InvalidSessionMetadata);
945 }
946 self.metadata = metadata.clone();
947 }
948 Event::SessionDeleted { .. } => {
949 if let Some(run_id) = self.active_run_id.take() {
950 self.run_status
951 .insert(run_id, RunState::Terminal(RunStatus::Cancelled));
952 }
953 for (_, (run_id, _)) in std::mem::take(&mut self.queued_inputs) {
954 self.run_status
955 .insert(run_id, RunState::Terminal(RunStatus::Cancelled));
956 }
957 self.waiting_interaction_id = None;
958 self.open_turn = None;
959 self.open_steps.clear();
960 self.open_tool_calls.clear();
961 self.started_tool_calls.clear();
962 self.open_compaction = None;
963 self.deleted = true;
964 }
965 Event::RunStarted { run_id, input_id } => {
966 if self.active_run_id.is_some() {
967 return Err(EventError::ConcurrentRun);
968 }
969 if self.claimed_inputs.get(input_id) != Some(run_id) {
970 return Err(EventError::UnclaimedInput(input_id.clone()));
971 }
972 self.active_run_id = Some(run_id.clone());
973 self.next_step = 1;
974 self.run_status.insert(run_id.clone(), RunState::Running);
975 }
976 Event::RunWaiting {
977 run_id,
978 interaction_id,
979 } => {
980 self.require_active(run_id)?;
981 self.waiting_interaction_id = Some(interaction_id.clone());
982 self.run_status
983 .insert(run_id.clone(), RunState::WaitingForInput);
984 }
985 Event::RunResumed {
986 run_id,
987 interaction_id,
988 } => {
989 self.require_active(run_id)?;
990 if self.waiting_interaction_id.as_ref() != Some(interaction_id) {
991 return Err(EventError::InteractionMismatch);
992 }
993 self.waiting_interaction_id = None;
994 self.run_status.insert(run_id.clone(), RunState::Running);
995 }
996 Event::RunFinished { run_id, status, .. } => {
997 self.require_active(run_id)?;
998 if !self.open_tool_calls.is_empty()
999 || !self.open_steps.is_empty()
1000 || self.open_turn.is_some()
1001 || self.open_compaction.is_some()
1002 || self
1003 .queued_inputs
1004 .values()
1005 .any(|(target_run_id, _)| target_run_id == run_id)
1006 {
1007 return Err(EventError::OpenLifecycle);
1008 }
1009 self.run_status
1010 .insert(run_id.clone(), RunState::Terminal(*status));
1011 self.active_run_id = None;
1012 self.waiting_interaction_id = None;
1013 }
1014 Event::UserMessage { run_id, content } if include_content => {
1015 self.messages.push(ProjectedMessage {
1016 role: "user",
1017 run_id: run_id.clone(),
1018 content: content.clone(),
1019 })
1020 }
1021 Event::AssistantMessage {
1022 run_id, content, ..
1023 } if include_content => self.messages.push(ProjectedMessage {
1024 role: "assistant",
1025 run_id: run_id.clone(),
1026 content: content.clone(),
1027 }),
1028 Event::ContextInjected {
1029 run_id,
1030 step,
1031 contribution_id,
1032 source,
1033 version,
1034 authority,
1035 form,
1036 content,
1037 } => {
1038 self.require_active(run_id)?;
1039 if !self.open_steps.contains(step) {
1040 return Err(EventError::LifecycleMismatch);
1041 }
1042 if include_content {
1043 self.injected_context.push(ProjectedContext {
1044 run_id: run_id.clone(),
1045 step: *step,
1046 contribution_id: contribution_id.clone(),
1047 source: source.clone(),
1048 version: version.clone(),
1049 authority: authority.clone(),
1050 form: form.clone(),
1051 content: content.clone(),
1052 });
1053 }
1054 }
1055 Event::InputQueued {
1056 input_id,
1057 run_id,
1058 mode,
1059 ..
1060 } => {
1061 if self.metadata.archived {
1062 return Err(EventError::SessionArchived);
1063 }
1064 if *mode != DeliveryMode::Followup {
1065 self.require_active(run_id)?;
1066 }
1067 if self
1068 .queued_inputs
1069 .insert(input_id.clone(), (run_id.clone(), *mode))
1070 .is_some()
1071 {
1072 return Err(EventError::DuplicateInput(input_id.clone()));
1073 }
1074 }
1075 Event::InputClaimed { input_id, run_id } => {
1076 if self
1077 .queued_inputs
1078 .remove(input_id)
1079 .map(|value| value.0)
1080 .as_ref()
1081 != Some(run_id)
1082 || self
1083 .claimed_inputs
1084 .insert(input_id.clone(), run_id.clone())
1085 .is_some()
1086 {
1087 return Err(EventError::UnqueuedInput(input_id.clone()));
1088 }
1089 }
1090 Event::InputCancelled {
1091 input_id,
1092 run_id,
1093 error_code,
1094 } => {
1095 if self
1096 .queued_inputs
1097 .remove(input_id)
1098 .map(|value| value.0)
1099 .as_ref()
1100 != Some(run_id)
1101 {
1102 return Err(EventError::UnqueuedInput(input_id.clone()));
1103 }
1104 self.run_status.insert(
1105 run_id.clone(),
1106 RunState::Terminal(if error_code == "cancelled" {
1107 RunStatus::Cancelled
1108 } else {
1109 RunStatus::Failed
1110 }),
1111 );
1112 }
1113 Event::TurnStarted { run_id, turn } => {
1114 self.require_active(run_id)?;
1115 if self.open_turn.replace(*turn).is_some() {
1116 return Err(EventError::ConcurrentTurn);
1117 }
1118 }
1119 Event::TurnFinished { run_id, turn } => {
1120 self.require_active(run_id)?;
1121 if self.open_turn != Some(*turn)
1122 || !self.open_steps.is_empty()
1123 || !self.open_tool_calls.is_empty()
1124 {
1125 return Err(EventError::LifecycleMismatch);
1126 }
1127 self.open_turn = None;
1128 }
1129 Event::StepStarted { run_id, step } => {
1130 self.require_active(run_id)?;
1131 if self.open_turn.is_none()
1132 || !self.open_steps.is_empty()
1133 || *step != self.next_step
1134 || !self.open_steps.insert(*step)
1135 {
1136 return Err(EventError::ConcurrentStep);
1137 }
1138 }
1139 Event::StepFinished { run_id, step } => {
1140 self.require_active(run_id)?;
1141 if self
1142 .open_tool_calls
1143 .values()
1144 .any(|call| call.run_id == *run_id && call.step == *step)
1145 || !self.open_steps.remove(step)
1146 {
1147 return Err(EventError::LifecycleMismatch);
1148 }
1149 self.next_step = step.saturating_add(1);
1150 }
1151 Event::ToolCall {
1152 run_id,
1153 step,
1154 call_id,
1155 tool,
1156 arguments,
1157 } => {
1158 self.require_active(run_id)?;
1159 if !self.open_steps.contains(step) {
1160 return Err(EventError::LifecycleMismatch);
1161 }
1162 if !self.seen_tool_calls.insert(call_id.clone())
1163 || self
1164 .open_tool_calls
1165 .insert(
1166 call_id.clone(),
1167 OpenToolCall {
1168 run_id: run_id.clone(),
1169 step: *step,
1170 tool: tool.clone(),
1171 arguments: arguments.clone(),
1172 source_event_seq: envelope.seq,
1173 },
1174 )
1175 .is_some()
1176 {
1177 return Err(EventError::DuplicateToolCall(call_id.clone()));
1178 }
1179 }
1180 Event::ToolResult {
1181 run_id,
1182 step,
1183 call_id,
1184 ..
1185 } => {
1186 self.require_active(run_id)?;
1187 if !self.open_steps.contains(step) {
1188 return Err(EventError::LifecycleMismatch);
1189 }
1190 let Some(call) = self.open_tool_calls.get(call_id) else {
1191 return Err(EventError::OrphanToolResult(call_id.clone()));
1192 };
1193 if call.run_id != *run_id || call.step != *step {
1194 return Err(EventError::ToolResultMismatch(call_id.clone()));
1195 }
1196 self.open_tool_calls.remove(call_id);
1197 self.started_tool_calls.remove(call_id);
1198 }
1199 Event::ToolExecutionStarted {
1200 run_id,
1201 step,
1202 call_id,
1203 metering,
1204 reserved_cost_units,
1205 } => {
1206 self.require_active(run_id)?;
1207 let Some(call) = self.open_tool_calls.get(call_id) else {
1208 return Err(EventError::OrphanToolResult(call_id.clone()));
1209 };
1210 if call.run_id != *run_id
1211 || call.step != *step
1212 || self.started_tool_calls.contains(call_id)
1213 {
1214 return Err(EventError::ToolResultMismatch(call_id.clone()));
1215 }
1216 if let Some(details) = metering {
1217 details.validate()?;
1218 if !matches!(details, MeteringDetails::Tool { call_id: id, name, source: MeteringSource::Estimated, outcome: MeteringOutcome::Unknown } if id == call_id && name == &call.tool)
1219 {
1220 return Err(EventError::InvalidMetering);
1221 }
1222 }
1223 self.started_tool_calls.insert(call_id.clone());
1224 self.pending_usage_operations.insert(
1225 (run_id.clone(), format!("tool:{call_id}")),
1226 (0, 0, *reserved_cost_units, metering.clone()),
1227 );
1228 }
1229 Event::ModelRequestPrepared {
1230 run_id,
1231 operation_id,
1232 reserved_prompt_tokens,
1233 reserved_completion_tokens,
1234 metering,
1235 ..
1236 } if !operation_id.is_empty() => {
1237 self.require_active(run_id)?;
1238 let key = (run_id.clone(), operation_id.clone());
1239 if !self.usage_operations.contains_key(&key) {
1240 if let Some(details) = metering {
1241 details.validate()?;
1242 if !matches!(
1243 details,
1244 MeteringDetails::Model {
1245 source: MeteringSource::Estimated,
1246 outcome: MeteringOutcome::Unknown,
1247 ..
1248 }
1249 ) {
1250 return Err(EventError::InvalidMetering);
1251 }
1252 }
1253 let reservation = (
1254 *reserved_prompt_tokens,
1255 *reserved_completion_tokens,
1256 0,
1257 metering.clone(),
1258 );
1259 match self.pending_usage_operations.get(&key) {
1260 Some(existing) if *existing != reservation => {
1261 return Err(EventError::UsageConflict(operation_id.clone()));
1262 }
1263 Some(_) => {}
1264 None => {
1265 self.pending_usage_operations.insert(key, reservation);
1266 }
1267 }
1268 }
1269 }
1270 Event::UsageRecorded {
1271 run_id,
1272 operation_id,
1273 prompt_tokens,
1274 completion_tokens,
1275 cost_units,
1276 metering,
1277 } => {
1278 self.require_active(run_id)?;
1279 if let Some(details) = metering {
1280 details.validate()?;
1281 }
1282 let key = (run_id.clone(), operation_id.clone());
1283 if let Some((_, _, _, Some(prepared))) = self.pending_usage_operations.get(&key) {
1284 if !metering
1285 .as_ref()
1286 .is_some_and(|details| details.completes(prepared))
1287 {
1288 return Err(EventError::UsageConflict(operation_id.clone()));
1289 }
1290 }
1291 let usage = (
1292 *prompt_tokens,
1293 *completion_tokens,
1294 *cost_units,
1295 metering.clone(),
1296 );
1297 match self.usage_operations.get(&key) {
1298 Some(existing) if *existing != usage => {
1299 return Err(EventError::UsageConflict(operation_id.clone()));
1300 }
1301 Some(_) => {}
1302 None => {
1303 self.pending_usage_operations.remove(&key);
1304 self.usage_operations.insert(key, usage);
1305 }
1306 }
1307 }
1308 Event::CompactionStarted {
1309 run_id,
1310 compaction_id,
1311 source_through_seq,
1312 } => {
1313 self.require_active(run_id)?;
1314 if self.open_compaction.is_some() {
1315 return Err(EventError::ConcurrentCompaction);
1316 }
1317 self.open_compaction = Some((compaction_id.clone(), *source_through_seq));
1318 }
1319 Event::CompactionFinished {
1320 run_id,
1321 compaction_id,
1322 ..
1323 } => {
1324 self.require_active(run_id)?;
1325 if self.open_compaction.as_ref().map(|value| value.0.as_str())
1326 != Some(compaction_id.as_str())
1327 {
1328 return Err(EventError::CompactionMismatch);
1329 }
1330 self.open_compaction = None;
1331 }
1332 Event::SummaryReplaced { summary, .. } => self.summary = Some(summary.clone()),
1333 Event::Opaque {
1334 event_type,
1335 ignorable: false,
1336 ..
1337 } => return Err(EventError::UnknownRequired(event_type.clone())),
1338 _ => {}
1339 }
1340 self.last_seq = envelope.seq;
1341 Ok(())
1342 }
1343
1344 fn require_active(&self, run_id: &RunId) -> Result<(), EventError> {
1345 if self.active_run_id.as_ref() == Some(run_id) {
1346 Ok(())
1347 } else {
1348 Err(EventError::RunMismatch)
1349 }
1350 }
1351
1352 pub fn usage_for(&self, run_id: &str) -> (u64, u64) {
1354 self.usage_operations
1355 .iter()
1356 .filter(|((recorded_run_id, _), _)| recorded_run_id.as_str() == run_id)
1357 .fold((0, 0), |total, (_, usage)| {
1358 (total.0 + usage.0, total.1 + usage.1)
1359 })
1360 }
1361
1362 pub fn billable_units_for(&self, run_id: &str) -> u64 {
1364 let recorded = self
1365 .usage_operations
1366 .iter()
1367 .filter(|((recorded_run_id, _), _)| recorded_run_id.as_str() == run_id)
1368 .map(|(_, usage)| usage.0 + usage.1 + usage.2)
1369 .sum::<u64>();
1370 recorded
1371 + self
1372 .pending_usage_operations
1373 .iter()
1374 .filter(|((recorded_run_id, _), _)| recorded_run_id.as_str() == run_id)
1375 .map(|(_, usage)| usage.0 + usage.1 + usage.2)
1376 .sum::<u64>()
1377 }
1378}
1379
1380#[derive(Debug, thiserror::Error, PartialEq, Eq)]
1382pub enum EventError {
1383 #[error("invalid Session metadata or metadata version conflict")]
1385 InvalidSessionMetadata,
1386 #[error("invalid metering attribution")]
1388 InvalidMetering,
1389 #[error("Session is archived")]
1391 SessionArchived,
1392 #[error("unsupported session event format version {0}")]
1394 UnsupportedFormat(u32),
1395 #[error("unknown required session event type {0}")]
1397 UnknownRequired(String),
1398 #[error("event sequence mismatch: expected {expected}, got {actual}")]
1400 Sequence {
1401 expected: u64,
1403 actual: u64,
1405 },
1406 #[error("event belongs to another session")]
1408 SessionMismatch,
1409 #[error("session creation must be the first and only creation event")]
1411 DuplicateSession,
1412 #[error("session already has an active run")]
1414 ConcurrentRun,
1415 #[error("session is closed")]
1417 SessionClosed,
1418 #[error("event does not match the active run")]
1420 RunMismatch,
1421 #[error("interaction does not match the waiting run")]
1423 InteractionMismatch,
1424 #[error("run cannot finish with an open turn, step or tool call")]
1426 OpenLifecycle,
1427 #[error("input was queued twice: {0}")]
1429 DuplicateInput(InputId),
1430 #[error("input was claimed before it was queued: {0}")]
1432 UnqueuedInput(InputId),
1433 #[error("run started from an unclaimed input: {0}")]
1435 UnclaimedInput(InputId),
1436 #[error("session already has an active turn")]
1438 ConcurrentTurn,
1439 #[error("turn already has this active step")]
1441 ConcurrentStep,
1442 #[error("turn or step lifecycle does not pair")]
1444 LifecycleMismatch,
1445 #[error("duplicate tool call {0}")]
1447 DuplicateToolCall(ToolCallId),
1448 #[error("tool result has no matching call {0}")]
1450 OrphanToolResult(ToolCallId),
1451 #[error("tool result does not match the call run and step: {0}")]
1453 ToolResultMismatch(ToolCallId),
1454 #[error("usage operation was recorded with different totals: {0}")]
1456 UsageConflict(String),
1457 #[error("session already has an active compaction")]
1459 ConcurrentCompaction,
1460 #[error("compaction lifecycle does not pair")]
1462 CompactionMismatch,
1463 #[error("event store conflict: {0}")]
1465 Conflict(String),
1466 #[error("event store unavailable: {0}")]
1468 Unavailable(String),
1469}
1470
1471#[async_trait]
1473pub trait SessionEventStore: Send + Sync {
1474 async fn append(
1476 &self,
1477 tenant_id: &str,
1478 session_id: &str,
1479 expected_seq: u64,
1480 events: Vec<Event>,
1481 ) -> Result<Vec<SessionEvent>, EventError>;
1482 async fn load(
1484 &self,
1485 tenant_id: &str,
1486 session_id: &str,
1487 after_seq: u64,
1488 ) -> Result<Vec<SessionEvent>, EventError>;
1489}
1490
1491pub fn text(value: impl Into<String>) -> Vec<ContentBlock> {
1493 vec![ContentBlock::Text { text: value.into() }]
1494}
1495
1496pub fn recovery_events(projection: &SessionProjection) -> Vec<Event> {
1498 failure_events(projection, "worker_restarted")
1499}
1500
1501pub fn failure_events(projection: &SessionProjection, error_code: &str) -> Vec<Event> {
1503 termination_events(projection, RunStatus::Failed, error_code)
1504}
1505
1506pub fn cancel_events(projection: &SessionProjection) -> Vec<Event> {
1508 termination_events(projection, RunStatus::Cancelled, "cancelled")
1509}
1510
1511pub fn session_deletion_events(projection: &SessionProjection, reason: &str) -> Vec<Event> {
1513 let active = projection.active_run_id.as_deref();
1514 let mut events = cancel_events(projection);
1515 events.extend(
1516 projection
1517 .queued_inputs
1518 .iter()
1519 .filter(|(_, (run_id, _))| Some(run_id.as_str()) != active)
1520 .map(|(input_id, (run_id, _))| Event::InputCancelled {
1521 input_id: input_id.clone(),
1522 run_id: run_id.clone(),
1523 error_code: "cancelled".into(),
1524 }),
1525 );
1526 events.push(Event::SessionDeleted {
1527 reason: reason.into(),
1528 });
1529 events
1530}
1531
1532fn termination_events(
1533 projection: &SessionProjection,
1534 status: RunStatus,
1535 error_code: &str,
1536) -> Vec<Event> {
1537 let Some(run_id) = &projection.active_run_id else {
1538 return Vec::new();
1539 };
1540 let mut events = projection
1541 .queued_inputs
1542 .iter()
1543 .filter(|(_, (target_run_id, _))| target_run_id == run_id)
1544 .map(|(input_id, _)| Event::InputCancelled {
1545 input_id: input_id.clone(),
1546 run_id: run_id.clone(),
1547 error_code: error_code.into(),
1548 })
1549 .collect::<Vec<_>>();
1550 events.extend(
1551 projection
1552 .open_tool_calls
1553 .iter()
1554 .map(|(call_id, call)| Event::ToolResult {
1555 run_id: run_id.clone(),
1556 step: call.step,
1557 call_id: call_id.clone(),
1558 result: termination_result(error_code),
1559 is_error: true,
1560 }),
1561 );
1562 if let Some((compaction_id, _)) = &projection.open_compaction {
1563 events.push(Event::CompactionFinished {
1564 run_id: run_id.clone(),
1565 compaction_id: compaction_id.clone(),
1566 status: "failed".into(),
1567 error: Some(error_code.into()),
1568 });
1569 }
1570 events.extend(
1571 projection
1572 .open_steps
1573 .iter()
1574 .map(|step| Event::StepFinished {
1575 run_id: run_id.clone(),
1576 step: *step,
1577 }),
1578 );
1579 if let Some(turn) = projection.open_turn {
1580 events.push(Event::TurnFinished {
1581 run_id: run_id.clone(),
1582 turn,
1583 });
1584 }
1585 events.push(Event::RunFinished {
1586 run_id: run_id.clone(),
1587 status,
1588 error_code: Some(error_code.into()),
1589 });
1590 events
1591}
1592
1593fn termination_result(error_code: &str) -> Value {
1594 if error_code == "tool_outcome_unknown" || error_code == "worker_restarted" {
1595 serde_json::json!({
1596 "error": error_code,
1597 "guidance": "The tool outcome is unknown. Verify external state before retrying any operation with side effects; ask the user when verification is unavailable."
1598 })
1599 } else {
1600 serde_json::json!({"error":error_code})
1601 }
1602}
1603
1604#[cfg(test)]
1605mod tests {
1606 use super::*;
1607
1608 fn event(seq: u64, event: Event) -> SessionEvent {
1609 SessionEvent {
1610 session_id: "s".parse().unwrap(),
1611 seq,
1612 occurred_at: Utc::now(),
1613 event,
1614 }
1615 }
1616
1617 #[test]
1618 fn replay_enforces_single_run_and_tool_pairs() {
1619 let events = vec![
1620 event(
1621 1,
1622 Event::SessionCreated {
1623 profile_revision_id: "p1".parse().unwrap(),
1624 },
1625 ),
1626 event(
1627 2,
1628 Event::InputQueued {
1629 input_id: "i1".parse().unwrap(),
1630 run_id: "r1".parse().unwrap(),
1631 mode: DeliveryMode::Followup,
1632 content: text("hi"),
1633 explicit_skill: None,
1634 },
1635 ),
1636 event(
1637 3,
1638 Event::InputClaimed {
1639 input_id: "i1".parse().unwrap(),
1640 run_id: "r1".parse().unwrap(),
1641 },
1642 ),
1643 event(
1644 4,
1645 Event::RunStarted {
1646 run_id: "r1".parse().unwrap(),
1647 input_id: "i1".parse().unwrap(),
1648 },
1649 ),
1650 event(
1651 5,
1652 Event::TurnStarted {
1653 run_id: "r1".parse().unwrap(),
1654 turn: 1,
1655 },
1656 ),
1657 event(
1658 6,
1659 Event::StepStarted {
1660 run_id: "r1".parse().unwrap(),
1661 step: 1,
1662 },
1663 ),
1664 event(
1665 7,
1666 Event::ToolCall {
1667 run_id: "r1".parse().unwrap(),
1668 step: 1,
1669 call_id: "c1".parse().unwrap(),
1670 tool: "echo".into(),
1671 arguments: serde_json::json!({"x":1}),
1672 },
1673 ),
1674 event(
1675 8,
1676 Event::ToolResult {
1677 run_id: "r1".parse().unwrap(),
1678 step: 1,
1679 call_id: "c1".parse().unwrap(),
1680 result: serde_json::json!({"x":1}),
1681 is_error: false,
1682 },
1683 ),
1684 event(
1685 9,
1686 Event::StepFinished {
1687 run_id: "r1".parse().unwrap(),
1688 step: 1,
1689 },
1690 ),
1691 event(
1692 10,
1693 Event::TurnFinished {
1694 run_id: "r1".parse().unwrap(),
1695 turn: 1,
1696 },
1697 ),
1698 event(
1699 11,
1700 Event::RunFinished {
1701 run_id: "r1".parse().unwrap(),
1702 status: RunStatus::Completed,
1703 error_code: None,
1704 },
1705 ),
1706 ];
1707 let projection = SessionProjection::replay(&events).unwrap();
1708 assert_eq!(projection.last_seq, 11);
1709 assert!(projection.active_run_id.is_none());
1710 }
1711
1712 #[test]
1713 fn replay_rejects_orphan_tool_result() {
1714 let events = vec![
1715 event(
1716 1,
1717 Event::SessionCreated {
1718 profile_revision_id: "p1".parse().unwrap(),
1719 },
1720 ),
1721 event(
1722 2,
1723 Event::InputQueued {
1724 input_id: "i1".parse().unwrap(),
1725 run_id: "r1".parse().unwrap(),
1726 mode: DeliveryMode::Followup,
1727 content: text("hi"),
1728 explicit_skill: None,
1729 },
1730 ),
1731 event(
1732 3,
1733 Event::InputClaimed {
1734 input_id: "i1".parse().unwrap(),
1735 run_id: "r1".parse().unwrap(),
1736 },
1737 ),
1738 event(
1739 4,
1740 Event::RunStarted {
1741 run_id: "r1".parse().unwrap(),
1742 input_id: "i1".parse().unwrap(),
1743 },
1744 ),
1745 event(
1746 5,
1747 Event::TurnStarted {
1748 run_id: "r1".parse().unwrap(),
1749 turn: 1,
1750 },
1751 ),
1752 event(
1753 6,
1754 Event::StepStarted {
1755 run_id: "r1".parse().unwrap(),
1756 step: 1,
1757 },
1758 ),
1759 event(
1760 7,
1761 Event::ToolResult {
1762 run_id: "r1".parse().unwrap(),
1763 step: 1,
1764 call_id: "missing".parse().unwrap(),
1765 result: Value::Null,
1766 is_error: true,
1767 },
1768 ),
1769 ];
1770 assert_eq!(
1771 SessionProjection::replay(&events).unwrap_err(),
1772 EventError::OrphanToolResult("missing".parse().unwrap())
1773 );
1774 }
1775
1776 #[test]
1777 fn queued_input_failure_is_not_projected_as_cancellation() {
1778 let events = vec![
1779 event(
1780 1,
1781 Event::SessionCreated {
1782 profile_revision_id: "p1".parse().unwrap(),
1783 },
1784 ),
1785 event(
1786 2,
1787 Event::InputQueued {
1788 input_id: "i1".parse().unwrap(),
1789 run_id: "r1".parse().unwrap(),
1790 mode: DeliveryMode::Followup,
1791 content: text("hi"),
1792 explicit_skill: None,
1793 },
1794 ),
1795 event(
1796 3,
1797 Event::InputCancelled {
1798 input_id: "i1".parse().unwrap(),
1799 run_id: "r1".parse().unwrap(),
1800 error_code: "profile_not_found".into(),
1801 },
1802 ),
1803 ];
1804 let projection = SessionProjection::replay(&events).unwrap();
1805 assert_eq!(
1806 projection.run_status.get("r1").map(|state| state.as_str()),
1807 Some("failed")
1808 );
1809 }
1810
1811 #[test]
1812 fn session_deletion_closes_active_and_queued_runs_before_tombstone() {
1813 let mut events = vec![
1814 event(
1815 1,
1816 Event::SessionCreated {
1817 profile_revision_id: "p1".parse().unwrap(),
1818 },
1819 ),
1820 event(
1821 2,
1822 Event::InputQueued {
1823 input_id: "i1".parse().unwrap(),
1824 run_id: "r1".parse().unwrap(),
1825 mode: DeliveryMode::Followup,
1826 content: text("start"),
1827 explicit_skill: None,
1828 },
1829 ),
1830 event(
1831 3,
1832 Event::InputClaimed {
1833 input_id: "i1".parse().unwrap(),
1834 run_id: "r1".parse().unwrap(),
1835 },
1836 ),
1837 event(
1838 4,
1839 Event::RunStarted {
1840 run_id: "r1".parse().unwrap(),
1841 input_id: "i1".parse().unwrap(),
1842 },
1843 ),
1844 event(
1845 5,
1846 Event::TurnStarted {
1847 run_id: "r1".parse().unwrap(),
1848 turn: 1,
1849 },
1850 ),
1851 event(
1852 6,
1853 Event::InputQueued {
1854 input_id: "i2".parse().unwrap(),
1855 run_id: "r2".parse().unwrap(),
1856 mode: DeliveryMode::Followup,
1857 content: text("later"),
1858 explicit_skill: None,
1859 },
1860 ),
1861 ];
1862 let projection = SessionProjection::replay(&events).unwrap();
1863 for event_value in session_deletion_events(&projection, "api_deleted") {
1864 let seq = events.len() as u64 + 1;
1865 events.push(event(seq, event_value));
1866 }
1867 let deleted = SessionProjection::replay(&events).unwrap();
1868 assert!(deleted.deleted);
1869 assert_eq!(
1870 deleted.run_status.get("r1").map(|state| state.as_str()),
1871 Some("cancelled")
1872 );
1873 assert_eq!(
1874 deleted.run_status.get("r2").map(|state| state.as_str()),
1875 Some("cancelled")
1876 );
1877 assert!(events.iter().any(|event| matches!(
1878 &event.event,
1879 Event::RunFinished { run_id, status: RunStatus::Cancelled, .. } if run_id == "r1"
1880 )));
1881 assert!(matches!(
1882 events.last().map(|event| &event.event),
1883 Some(Event::SessionDeleted { .. })
1884 ));
1885 }
1886
1887 #[test]
1888 fn usage_operations_are_idempotent_and_conflicts_fail_replay() {
1889 let mut events = vec![
1890 event(
1891 1,
1892 Event::SessionCreated {
1893 profile_revision_id: "p1".parse().unwrap(),
1894 },
1895 ),
1896 event(
1897 2,
1898 Event::InputQueued {
1899 input_id: "i1".parse().unwrap(),
1900 run_id: "r1".parse().unwrap(),
1901 mode: DeliveryMode::Followup,
1902 content: text("hi"),
1903 explicit_skill: None,
1904 },
1905 ),
1906 event(
1907 3,
1908 Event::InputClaimed {
1909 input_id: "i1".parse().unwrap(),
1910 run_id: "r1".parse().unwrap(),
1911 },
1912 ),
1913 event(
1914 4,
1915 Event::RunStarted {
1916 run_id: "r1".parse().unwrap(),
1917 input_id: "i1".parse().unwrap(),
1918 },
1919 ),
1920 event(
1921 5,
1922 Event::UsageRecorded {
1923 metering: None,
1924 run_id: "r1".parse().unwrap(),
1925 operation_id: "model:1:attempt:1".into(),
1926 prompt_tokens: 7,
1927 completion_tokens: 3,
1928 cost_units: 5,
1929 },
1930 ),
1931 event(
1932 6,
1933 Event::UsageRecorded {
1934 metering: None,
1935 run_id: "r1".parse().unwrap(),
1936 operation_id: "model:1:attempt:1".into(),
1937 prompt_tokens: 7,
1938 completion_tokens: 3,
1939 cost_units: 5,
1940 },
1941 ),
1942 ];
1943 assert_eq!(
1944 SessionProjection::replay(&events).unwrap().usage_for("r1"),
1945 (7, 3)
1946 );
1947 assert_eq!(
1948 SessionProjection::replay(&events)
1949 .unwrap()
1950 .billable_units_for("r1"),
1951 15
1952 );
1953
1954 let mut attributed = events.clone();
1955 let details = MeteringDetails::Model {
1956 model: "m".into(),
1957 provider: Some("p".into()),
1958 provider_attempt_id: "r1:model:1:attempt:1".parse().unwrap(),
1959 source: MeteringSource::Estimated,
1960 outcome: MeteringOutcome::Unknown,
1961 };
1962 for entry in &mut attributed[4..] {
1963 if let Event::UsageRecorded { metering, .. } = &mut entry.event {
1964 *metering = Some(details.clone());
1965 }
1966 }
1967 assert_eq!(
1968 SessionProjection::replay(&attributed)
1969 .unwrap()
1970 .billable_units_for("r1"),
1971 15
1972 );
1973 if let Event::UsageRecorded {
1974 metering: Some(MeteringDetails::Model { source, .. }),
1975 ..
1976 } = &mut attributed[5].event
1977 {
1978 *source = MeteringSource::Reported;
1979 }
1980 assert!(matches!(
1981 SessionProjection::replay(&attributed),
1982 Err(EventError::UsageConflict(_))
1983 ));
1984 let mut prepared = attributed[..4].to_vec();
1985 prepared.push(event(
1986 5,
1987 Event::ModelRequestPrepared {
1988 metering: Some(details.clone()),
1989 run_id: "r1".parse().unwrap(),
1990 step: 1,
1991 attempt: 1,
1992 provider_attempt_id: "r1:model:1:attempt:1".into(),
1993 operation_id: "model:1:attempt:1".into(),
1994 reserved_prompt_tokens: 10,
1995 reserved_completion_tokens: 20,
1996 request: Value::Null,
1997 prompt_sections: Value::Null,
1998 },
1999 ));
2000 prepared.push(attributed[5].clone());
2001 assert_eq!(
2002 SessionProjection::replay(&prepared)
2003 .unwrap()
2004 .billable_units_for("r1"),
2005 15
2006 );
2007 if let Event::UsageRecorded {
2008 metering: Some(MeteringDetails::Model { model, .. }),
2009 ..
2010 } = &mut prepared[5].event
2011 {
2012 *model = "wrong".into();
2013 }
2014 assert!(matches!(
2015 SessionProjection::replay(&prepared),
2016 Err(EventError::UsageConflict(_))
2017 ));
2018 if let Event::UsageRecorded { metering, .. } = &mut prepared[5].event {
2019 *metering = None;
2020 }
2021 assert!(matches!(
2022 SessionProjection::replay(&prepared),
2023 Err(EventError::UsageConflict(_))
2024 ));
2025 let serialized = serde_json::to_value(&events[4]).unwrap();
2026 assert!(serialized["event"].get("metering").is_none());
2027 assert_eq!(
2028 serde_json::from_value::<SessionEvent>(serialized).unwrap(),
2029 events[4]
2030 );
2031
2032 events.push(event(
2033 7,
2034 Event::UsageRecorded {
2035 metering: None,
2036 run_id: "r1".parse().unwrap(),
2037 operation_id: "model:1:attempt:1".into(),
2038 prompt_tokens: 8,
2039 completion_tokens: 3,
2040 cost_units: 0,
2041 },
2042 ));
2043 assert_eq!(
2044 SessionProjection::replay(&events).unwrap_err(),
2045 EventError::UsageConflict("model:1:attempt:1".into())
2046 );
2047 }
2048
2049 #[test]
2050 fn unresolved_provider_attempt_is_conservatively_billable_and_reconcilable() {
2051 let mut events = vec![
2052 event(
2053 1,
2054 Event::SessionCreated {
2055 profile_revision_id: "p1".parse().unwrap(),
2056 },
2057 ),
2058 event(
2059 2,
2060 Event::InputQueued {
2061 input_id: "i1".parse().unwrap(),
2062 run_id: "r1".parse().unwrap(),
2063 mode: DeliveryMode::Followup,
2064 content: text("hi"),
2065 explicit_skill: None,
2066 },
2067 ),
2068 event(
2069 3,
2070 Event::InputClaimed {
2071 input_id: "i1".parse().unwrap(),
2072 run_id: "r1".parse().unwrap(),
2073 },
2074 ),
2075 event(
2076 4,
2077 Event::RunStarted {
2078 run_id: "r1".parse().unwrap(),
2079 input_id: "i1".parse().unwrap(),
2080 },
2081 ),
2082 event(
2083 5,
2084 Event::ModelRequestPrepared {
2085 metering: None,
2086 run_id: "r1".parse().unwrap(),
2087 step: 1,
2088 attempt: 1,
2089 provider_attempt_id: "r1:model:1:attempt:1".into(),
2090 operation_id: "model:1:attempt:1".into(),
2091 reserved_prompt_tokens: 7,
2092 reserved_completion_tokens: 11,
2093 request: Value::Null,
2094 prompt_sections: Value::Null,
2095 },
2096 ),
2097 ];
2098 assert_eq!(
2099 SessionProjection::replay(&events)
2100 .unwrap()
2101 .billable_units_for("r1"),
2102 18
2103 );
2104 events.push(event(
2105 6,
2106 Event::UsageRecorded {
2107 metering: None,
2108 run_id: "r1".parse().unwrap(),
2109 operation_id: "model:1:attempt:1".into(),
2110 prompt_tokens: 6,
2111 completion_tokens: 2,
2112 cost_units: 0,
2113 },
2114 ));
2115 assert_eq!(
2116 SessionProjection::replay(&events)
2117 .unwrap()
2118 .billable_units_for("r1"),
2119 8
2120 );
2121 }
2122
2123 #[test]
2124 fn a3_incremental_facts_preserve_every_lifecycle_and_usage_invariant() {
2125 let mut events = open_step_events();
2126 for fact in [
2127 Event::UserMessage {
2128 run_id: "r1".parse().unwrap(),
2129 content: text("history".repeat(10_000)),
2130 },
2131 Event::UsageRecorded {
2132 metering: None,
2133 run_id: "r1".parse().unwrap(),
2134 operation_id: "attempt".into(),
2135 prompt_tokens: 5,
2136 completion_tokens: 3,
2137 cost_units: 0,
2138 },
2139 Event::ToolCall {
2140 run_id: "r1".parse().unwrap(),
2141 step: 1,
2142 call_id: "pending".parse().unwrap(),
2143 tool: "echo".into(),
2144 arguments: Value::Null,
2145 },
2146 ] {
2147 events.push(event(events.len() as u64 + 1, fact));
2148 }
2149 let mut full = SessionProjection::replay(&events).unwrap();
2150 assert!(!full.messages.is_empty());
2151 full.messages.clear();
2152 full.injected_context.clear();
2153 let mut folded = SessionProjection::default();
2154 for page in events.chunks(3) {
2155 for fact in page {
2156 folded.apply_facts(fact).unwrap();
2157 }
2158 assert!(folded.messages.is_empty());
2159 }
2160 assert_eq!(folded, full);
2161 assert_eq!(SessionProjection::replay_facts(&events).unwrap(), full);
2162 assert_eq!(cancel_events(&folded), cancel_events(&full));
2163 assert_eq!(folded.usage_for("r1"), (5, 3));
2164 let invalid = event(
2165 events.len() as u64 + 1,
2166 Event::ToolResult {
2167 run_id: "r1".parse().unwrap(),
2168 step: 1,
2169 call_id: "missing".parse().unwrap(),
2170 result: Value::Null,
2171 is_error: true,
2172 },
2173 );
2174 assert_eq!(folded.apply_facts(&invalid), full.apply(&invalid));
2175 }
2176
2177 fn open_step_events() -> Vec<SessionEvent> {
2178 vec![
2179 event(
2180 1,
2181 Event::SessionCreated {
2182 profile_revision_id: "p1".parse().unwrap(),
2183 },
2184 ),
2185 event(
2186 2,
2187 Event::InputQueued {
2188 input_id: "i1".parse().unwrap(),
2189 run_id: "r1".parse().unwrap(),
2190 mode: DeliveryMode::Followup,
2191 content: text("hi"),
2192 explicit_skill: None,
2193 },
2194 ),
2195 event(
2196 3,
2197 Event::InputClaimed {
2198 input_id: "i1".parse().unwrap(),
2199 run_id: "r1".parse().unwrap(),
2200 },
2201 ),
2202 event(
2203 4,
2204 Event::RunStarted {
2205 run_id: "r1".parse().unwrap(),
2206 input_id: "i1".parse().unwrap(),
2207 },
2208 ),
2209 event(
2210 5,
2211 Event::TurnStarted {
2212 run_id: "r1".parse().unwrap(),
2213 turn: 1,
2214 },
2215 ),
2216 event(
2217 6,
2218 Event::StepStarted {
2219 run_id: "r1".parse().unwrap(),
2220 step: 1,
2221 },
2222 ),
2223 ]
2224 }
2225
2226 #[test]
2227 fn replay_rejects_overlapping_or_out_of_order_steps() {
2228 let mut overlapping = open_step_events();
2229 overlapping.push(event(
2230 7,
2231 Event::StepStarted {
2232 run_id: "r1".parse().unwrap(),
2233 step: 2,
2234 },
2235 ));
2236 assert_eq!(
2237 SessionProjection::replay(&overlapping).unwrap_err(),
2238 EventError::ConcurrentStep
2239 );
2240
2241 let mut skipped = open_step_events();
2242 skipped[5] = event(
2243 6,
2244 Event::StepStarted {
2245 run_id: "r1".parse().unwrap(),
2246 step: 2,
2247 },
2248 );
2249 assert_eq!(
2250 SessionProjection::replay(&skipped).unwrap_err(),
2251 EventError::ConcurrentStep
2252 );
2253 }
2254
2255 #[test]
2256 fn replay_rejects_cross_step_results_and_dangling_calls() {
2257 let mut events = open_step_events();
2258 events.push(event(
2259 7,
2260 Event::ToolCall {
2261 run_id: "r1".parse().unwrap(),
2262 step: 1,
2263 call_id: "c1".parse().unwrap(),
2264 tool: "echo".into(),
2265 arguments: Value::Null,
2266 },
2267 ));
2268 events.push(event(
2269 8,
2270 Event::ToolResult {
2271 run_id: "r1".parse().unwrap(),
2272 step: 2,
2273 call_id: "c1".parse().unwrap(),
2274 result: Value::Null,
2275 is_error: false,
2276 },
2277 ));
2278 assert_eq!(
2279 SessionProjection::replay(&events).unwrap_err(),
2280 EventError::LifecycleMismatch
2281 );
2282
2283 let mut dangling = open_step_events();
2284 dangling.push(event(
2285 7,
2286 Event::ToolCall {
2287 run_id: "r1".parse().unwrap(),
2288 step: 1,
2289 call_id: "c1".parse().unwrap(),
2290 tool: "echo".into(),
2291 arguments: Value::Null,
2292 },
2293 ));
2294 dangling.push(event(
2295 8,
2296 Event::StepFinished {
2297 run_id: "r1".parse().unwrap(),
2298 step: 1,
2299 },
2300 ));
2301 assert_eq!(
2302 SessionProjection::replay(&dangling).unwrap_err(),
2303 EventError::LifecycleMismatch
2304 );
2305 }
2306
2307 #[test]
2308 fn replay_rejects_reused_tool_call_ids() {
2309 let mut events = open_step_events();
2310 events.extend([
2311 event(
2312 7,
2313 Event::ToolCall {
2314 run_id: "r1".parse().unwrap(),
2315 step: 1,
2316 call_id: "c1".parse().unwrap(),
2317 tool: "echo".into(),
2318 arguments: Value::Null,
2319 },
2320 ),
2321 event(
2322 8,
2323 Event::ToolResult {
2324 run_id: "r1".parse().unwrap(),
2325 step: 1,
2326 call_id: "c1".parse().unwrap(),
2327 result: Value::Null,
2328 is_error: false,
2329 },
2330 ),
2331 event(
2332 9,
2333 Event::StepFinished {
2334 run_id: "r1".parse().unwrap(),
2335 step: 1,
2336 },
2337 ),
2338 event(
2339 10,
2340 Event::StepStarted {
2341 run_id: "r1".parse().unwrap(),
2342 step: 2,
2343 },
2344 ),
2345 event(
2346 11,
2347 Event::ToolCall {
2348 run_id: "r1".parse().unwrap(),
2349 step: 2,
2350 call_id: "c1".parse().unwrap(),
2351 tool: "echo".into(),
2352 arguments: Value::Null,
2353 },
2354 ),
2355 ]);
2356 assert_eq!(
2357 SessionProjection::replay(&events).unwrap_err(),
2358 EventError::DuplicateToolCall("c1".parse().unwrap())
2359 );
2360 }
2361
2362 #[test]
2363 fn replay_keeps_injected_context_without_changing_run_state() {
2364 let mut events = open_step_events();
2365 events.push(event(
2366 7,
2367 Event::ContextInjected {
2368 run_id: "r1".parse().unwrap(),
2369 step: 1,
2370 contribution_id: "memory:1".into(),
2371 source: "memory".into(),
2372 version: "v1".into(),
2373 authority: "tenant".into(),
2374 form: "message".into(),
2375 content: text("tenant context"),
2376 },
2377 ));
2378
2379 let projection = SessionProjection::replay(&events).unwrap();
2380 assert_eq!(projection.active_run_id.as_deref(), Some("r1"));
2381 assert_eq!(projection.open_steps, BTreeSet::from([1]));
2382 assert_eq!(
2383 projection.injected_context,
2384 vec![ProjectedContext {
2385 run_id: "r1".parse().unwrap(),
2386 step: 1,
2387 contribution_id: "memory:1".into(),
2388 source: "memory".into(),
2389 version: "v1".into(),
2390 authority: "tenant".into(),
2391 form: "message".into(),
2392 content: text("tenant context"),
2393 }]
2394 );
2395 }
2396
2397 #[test]
2398 fn replay_skips_unknown_ignorable_formats_and_rejects_required_events() {
2399 let optional = event(
2400 1,
2401 Event::Opaque {
2402 format_version: SESSION_EVENT_FORMAT_VERSION,
2403 event_type: "future_optional".into(),
2404 ignorable: true,
2405 payload: serde_json::json!({"answer": 42}),
2406 },
2407 );
2408 let projection = SessionProjection::replay(&[optional]).unwrap();
2409 assert_eq!(projection.last_seq, 1);
2410
2411 let required = event(
2412 1,
2413 Event::Opaque {
2414 format_version: SESSION_EVENT_FORMAT_VERSION,
2415 event_type: "future_required".into(),
2416 ignorable: false,
2417 payload: Value::Null,
2418 },
2419 );
2420 assert_eq!(
2421 SessionProjection::replay(&[required]).unwrap_err(),
2422 EventError::UnknownRequired("future_required".into())
2423 );
2424
2425 let unsupported = event(
2426 1,
2427 Event::Opaque {
2428 format_version: SESSION_EVENT_FORMAT_VERSION + 1,
2429 event_type: "future_optional".into(),
2430 ignorable: true,
2431 payload: Value::Null,
2432 },
2433 );
2434 assert_eq!(
2435 SessionProjection::replay(&[unsupported]).unwrap().last_seq,
2436 1
2437 );
2438
2439 let required_new_format = event(
2440 1,
2441 Event::Opaque {
2442 format_version: SESSION_EVENT_FORMAT_VERSION + 1,
2443 event_type: "future_required".into(),
2444 ignorable: false,
2445 payload: Value::Null,
2446 },
2447 );
2448 assert_eq!(
2449 SessionProjection::replay(&[required_new_format]).unwrap_err(),
2450 EventError::UnsupportedFormat(SESSION_EVENT_FORMAT_VERSION + 1)
2451 );
2452 }
2453}