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, Eq, Serialize, Deserialize)]
130pub struct ToolResultArtifact {
131 pub result_ref: String,
133 pub sha256: String,
135 pub original_bytes: u64,
137 pub view_bytes: u64,
139 #[serde(default, skip_serializing_if = "Option::is_none")]
141 pub expires_at: Option<DateTime<Utc>>,
142 #[serde(default, skip_serializing_if = "Vec::is_empty")]
144 pub omitted: Vec<String>,
145}
146
147#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
149#[serde(tag = "type", rename_all = "snake_case")]
150pub enum Event {
151 SessionCreated {
153 profile_revision_id: ProfileRevisionId,
155 },
156 SessionForked {
158 parent_session_id: SessionId,
160 parent_seq: u64,
162 },
163 SessionMetadataUpdated {
165 metadata: SessionMetadata,
167 },
168 SessionDeleted {
170 reason: String,
172 },
173 InputQueued {
175 input_id: InputId,
177 run_id: RunId,
179 mode: DeliveryMode,
181 content: Vec<ContentBlock>,
183 explicit_skill: Option<String>,
185 },
186 InputClaimed {
188 input_id: InputId,
190 run_id: RunId,
192 },
193 InputCancelled {
195 input_id: InputId,
197 run_id: RunId,
199 error_code: String,
201 },
202 RunStarted {
204 run_id: RunId,
206 input_id: InputId,
208 },
209 RunWaiting {
211 run_id: RunId,
213 interaction_id: InteractionId,
215 },
216 RunResumed {
218 run_id: RunId,
220 interaction_id: InteractionId,
222 },
223 RunFinished {
225 run_id: RunId,
227 status: RunStatus,
229 error_code: Option<String>,
231 },
232 TurnStarted {
234 run_id: RunId,
236 turn: u32,
238 },
239 TurnFinished {
241 run_id: RunId,
243 turn: u32,
245 },
246 StepStarted {
248 run_id: RunId,
250 step: u32,
252 },
253 StepFinished {
255 run_id: RunId,
257 step: u32,
259 },
260 UserMessage {
262 run_id: RunId,
264 content: Vec<ContentBlock>,
266 },
267 AssistantDelta {
269 run_id: RunId,
271 step: u32,
273 attempt: u32,
275 content: String,
277 },
278 AssistantMessage {
280 run_id: RunId,
282 step: u32,
284 attempt: u32,
286 content: Vec<ContentBlock>,
288 },
289 AssistantToolCalls {
291 run_id: RunId,
293 step: u32,
295 content: Option<String>,
297 calls: Vec<RecordedToolCall>,
299 },
300 ToolCall {
302 run_id: RunId,
304 step: u32,
306 call_id: ToolCallId,
308 tool: String,
310 arguments: Value,
312 },
313 ToolAuthorization {
315 run_id: RunId,
317 step: u32,
319 call_id: ToolCallId,
321 status: ToolAuthorizationStatus,
323 reason: Option<String>,
325 },
326 ToolExecutionStarted {
328 #[serde(default, skip_serializing_if = "Option::is_none")]
330 metering: Option<MeteringDetails>,
331 #[serde(default)]
333 reserved_cost_units: u64,
334 run_id: RunId,
336 step: u32,
338 call_id: ToolCallId,
340 },
341 ToolResult {
343 run_id: RunId,
345 step: u32,
347 call_id: ToolCallId,
349 result: Value,
351 is_error: bool,
353 },
354 UsageCorrected {
356 correction: MeteringCorrection,
358 actor_id: af_context::SubjectId,
360 },
361 UsageRecorded {
363 #[serde(default, skip_serializing_if = "Option::is_none")]
365 metering: Option<MeteringDetails>,
366 run_id: RunId,
368 operation_id: String,
370 prompt_tokens: u64,
372 completion_tokens: u64,
374 #[serde(default)]
376 cost_units: u64,
377 },
378 RetryScheduled {
380 run_id: RunId,
382 attempt: u32,
384 delay_ms: u64,
386 reason: String,
388 },
389 ModelRequestPrepared {
391 #[serde(default, skip_serializing_if = "Option::is_none")]
393 metering: Option<MeteringDetails>,
394 run_id: RunId,
396 step: u32,
398 attempt: u32,
400 #[serde(default)]
402 provider_attempt_id: String,
403 #[serde(default)]
405 operation_id: String,
406 #[serde(default)]
408 reserved_prompt_tokens: u64,
409 #[serde(default)]
411 reserved_completion_tokens: u64,
412 request: Value,
414 prompt_sections: Value,
416 },
417 ContextInjected {
419 run_id: RunId,
421 step: u32,
423 contribution_id: String,
425 source: String,
427 version: String,
429 authority: String,
431 form: String,
433 content: Vec<ContentBlock>,
435 },
436 ModelAttemptFailed {
438 run_id: RunId,
440 step: u32,
442 attempt: u32,
444 error: String,
446 retryable: bool,
448 },
449 CompactionStarted {
451 run_id: RunId,
453 compaction_id: String,
455 source_through_seq: u64,
457 },
458 ToolResultsPruned {
460 run_id: RunId,
462 call_ids: Vec<ToolCallId>,
464 },
465 SummaryReplaced {
467 run_id: RunId,
469 through_seq: u64,
471 summary: String,
473 compactor: String,
475 model: String,
477 },
478 CompactionFinished {
480 run_id: RunId,
482 compaction_id: String,
484 status: String,
486 error: Option<String>,
488 },
489 InteractionRequested {
491 run_id: RunId,
493 interaction_id: InteractionId,
495 kind: InteractionKind,
497 payload: Value,
499 },
500 InteractionResolved {
502 run_id: RunId,
504 interaction_id: InteractionId,
506 resolution: InteractionResolution,
508 payload: Value,
510 },
511 ChildSessionLinked {
513 run_id: RunId,
515 child_session_id: SessionId,
517 provider: String,
519 },
520 Extension {
522 run_id: RunId,
524 plugin_id: String,
526 event_type: String,
528 payload: Value,
530 },
531 #[serde(skip)]
533 Opaque {
534 format_version: u32,
536 event_type: String,
538 ignorable: bool,
540 payload: Value,
542 },
543}
544
545impl Event {
546 pub fn format_version(&self) -> u32 {
548 match self {
549 Self::Opaque { format_version, .. } => *format_version,
550 _ => SESSION_EVENT_FORMAT_VERSION,
551 }
552 }
553
554 pub fn event_type(&self) -> &str {
556 match self {
557 Self::SessionCreated { .. } => "session_created",
558 Self::SessionMetadataUpdated { .. } => "session_metadata_updated",
559 Self::SessionForked { .. } => "session_forked",
560 Self::SessionDeleted { .. } => "session_deleted",
561 Self::InputQueued { .. } => "input_queued",
562 Self::InputClaimed { .. } => "input_claimed",
563 Self::InputCancelled { .. } => "input_cancelled",
564 Self::RunStarted { .. } => "run_started",
565 Self::RunWaiting { .. } => "run_waiting",
566 Self::RunResumed { .. } => "run_resumed",
567 Self::RunFinished { .. } => "run_finished",
568 Self::TurnStarted { .. } => "turn_started",
569 Self::TurnFinished { .. } => "turn_finished",
570 Self::StepStarted { .. } => "step_started",
571 Self::StepFinished { .. } => "step_finished",
572 Self::UserMessage { .. } => "user_message",
573 Self::AssistantDelta { .. } => "assistant_delta",
574 Self::AssistantMessage { .. } => "assistant_message",
575 Self::AssistantToolCalls { .. } => "assistant_tool_calls",
576 Self::ToolCall { .. } => "tool_call",
577 Self::ToolAuthorization { .. } => "tool_authorization",
578 Self::ToolExecutionStarted { .. } => "tool_execution_started",
579 Self::ToolResult { .. } => "tool_result",
580 Self::UsageRecorded { .. } => "usage_recorded",
581 Self::UsageCorrected { .. } => "usage_corrected",
582 Self::RetryScheduled { .. } => "retry_scheduled",
583 Self::ModelRequestPrepared { .. } => "model_request_prepared",
584 Self::ContextInjected { .. } => "context_injected",
585 Self::ModelAttemptFailed { .. } => "model_attempt_failed",
586 Self::CompactionStarted { .. } => "compaction_started",
587 Self::ToolResultsPruned { .. } => "tool_results_pruned",
588 Self::SummaryReplaced { .. } => "summary_replaced",
589 Self::CompactionFinished { .. } => "compaction_finished",
590 Self::InteractionRequested { .. } => "interaction_requested",
591 Self::InteractionResolved { .. } => "interaction_resolved",
592 Self::ChildSessionLinked { .. } => "child_session_linked",
593 Self::Extension { .. } => "extension",
594 Self::Opaque { event_type, .. } => event_type,
595 }
596 }
597
598 pub fn ignorable(&self) -> bool {
600 match self {
601 Self::Opaque { ignorable, .. } => *ignorable,
602 _ => false,
603 }
604 }
605}
606
607#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
609#[serde(rename_all = "snake_case")]
610pub enum DeliveryMode {
611 Followup,
613 Steer,
615 Inject,
617}
618
619#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
621#[serde(rename_all = "snake_case")]
622pub enum RunStatus {
623 Completed,
625 Failed,
627 Cancelled,
629 MaxStepsReached,
631}
632
633impl RunStatus {
634 pub const fn as_str(self) -> &'static str {
636 match self {
637 Self::Completed => "completed",
638 Self::Failed => "failed",
639 Self::Cancelled => "cancelled",
640 Self::MaxStepsReached => "max_steps_reached",
641 }
642 }
643}
644
645#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
647#[serde(rename_all = "snake_case")]
648pub enum InteractionKind {
649 Action,
651 UserQuestion,
653}
654
655impl InteractionKind {
656 pub const fn as_str(self) -> &'static str {
658 match self {
659 Self::Action => "action",
660 Self::UserQuestion => "user_question",
661 }
662 }
663}
664
665#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
667#[serde(rename_all = "snake_case")]
668pub enum InteractionResolution {
669 Confirmed,
671 Rejected,
673 Answered,
675}
676
677impl InteractionResolution {
678 pub const fn as_str(self) -> &'static str {
680 match self {
681 Self::Confirmed => "confirmed",
682 Self::Rejected => "rejected",
683 Self::Answered => "answered",
684 }
685 }
686}
687
688#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
690#[serde(rename_all = "snake_case")]
691pub enum ToolAuthorizationStatus {
692 Allowed,
694 Waiting,
696 Denied,
698}
699
700#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
702pub struct RecordedToolCall {
703 pub call_id: ToolCallId,
705 pub tool: String,
707 pub arguments: Value,
709}
710
711impl ToolAuthorizationStatus {
712 pub const fn as_str(self) -> &'static str {
714 match self {
715 Self::Allowed => "allowed",
716 Self::Waiting => "waiting",
717 Self::Denied => "denied",
718 }
719 }
720}
721
722#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
724#[serde(tag = "type", rename_all = "snake_case")]
725pub enum ContentBlock {
726 Text {
728 text: String,
730 },
731 Resource {
733 resource_id: String,
735 media_type: String,
737 },
738 Data {
740 slot: String,
742 value: Value,
744 },
745 Citation {
747 resource_id: String,
749 label: String,
751 uri: String,
753 excerpt: Option<String>,
755 },
756}
757
758#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
760pub struct UsageCharge {
761 pub prompt_tokens: u64,
763 pub completion_tokens: u64,
765 pub cost_units: u64,
767}
768
769impl UsageCharge {
770 pub fn total(self) -> u64 {
772 self.prompt_tokens
773 .saturating_add(self.completion_tokens)
774 .saturating_add(self.cost_units)
775 }
776}
777
778#[derive(Debug, Clone, Default, PartialEq)]
780pub struct SessionProjection {
781 pub metadata: SessionMetadata,
783 pub session_id: Option<SessionId>,
785 pub profile_revision_id: Option<ProfileRevisionId>,
788 pub deleted: bool,
790 pub last_seq: u64,
792 pub active_run_id: Option<RunId>,
794 pub waiting_interaction_id: Option<InteractionId>,
796 pub messages: Vec<ProjectedMessage>,
798 pub injected_context: Vec<ProjectedContext>,
800 pub run_status: BTreeMap<RunId, RunState>,
802 pub open_tool_calls: BTreeMap<ToolCallId, OpenToolCall>,
804 pub started_tool_calls: BTreeSet<ToolCallId>,
806 pub queued_inputs: BTreeMap<InputId, (RunId, DeliveryMode)>,
808 pub claimed_inputs: BTreeMap<InputId, RunId>,
810 pub open_turn: Option<u32>,
812 pub open_steps: BTreeSet<u32>,
814 pub next_step: u32,
816 pub open_compaction: Option<(String, u64)>,
818 pub summary: Option<String>,
820 corrected_usage: BTreeMap<(RunId, String), OperationUsage>,
821 metering_corrections: BTreeMap<af_context::MeteringCorrectionId, AppliedMeteringCorrection>,
822 usage_operations: BTreeMap<(RunId, String), RecordedUsage>,
823 pending_usage_operations: BTreeMap<(RunId, String), PreparedUsage>,
824 seen_tool_calls: BTreeSet<ToolCallId>,
825}
826
827type PreparedUsage = RecordedUsage;
828
829type RecordedUsage = (u64, u64, u64, Option<MeteringDetails>);
830
831#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
833#[serde(rename_all = "snake_case")]
834pub enum RunState {
835 Running,
837 WaitingForInput,
839 Terminal(RunStatus),
841}
842
843impl RunState {
844 pub const fn as_str(self) -> &'static str {
846 match self {
847 Self::Running => "running",
848 Self::WaitingForInput => "waiting_for_input",
849 Self::Terminal(status) => status.as_str(),
850 }
851 }
852
853 pub const fn is_terminal(self) -> bool {
855 matches!(self, Self::Terminal(_))
856 }
857}
858
859#[derive(Debug, Clone, PartialEq)]
861pub struct ProjectedContext {
862 pub run_id: RunId,
864 pub step: u32,
866 pub contribution_id: String,
868 pub source: String,
870 pub version: String,
872 pub authority: String,
874 pub form: String,
876 pub content: Vec<ContentBlock>,
878}
879
880#[derive(Debug, Clone, PartialEq)]
882pub struct ProjectedMessage {
883 pub role: &'static str,
885 pub run_id: RunId,
887 pub content: Vec<ContentBlock>,
889}
890
891#[derive(Debug, Clone, PartialEq)]
893pub struct OpenToolCall {
894 pub run_id: RunId,
896 pub step: u32,
898 pub tool: String,
900 pub arguments: Value,
902 pub source_event_seq: u64,
904}
905
906impl SessionProjection {
907 pub fn replay(events: &[SessionEvent]) -> Result<Self, EventError> {
909 let mut projection = Self::default();
910 for event in events {
911 projection.apply(event)?;
912 }
913 Ok(projection)
914 }
915
916 pub fn apply(&mut self, envelope: &SessionEvent) -> Result<(), EventError> {
918 self.apply_internal(envelope, true)
919 }
920
921 pub fn apply_facts(&mut self, envelope: &SessionEvent) -> Result<(), EventError> {
924 self.messages.clear();
925 self.injected_context.clear();
926 self.apply_internal(envelope, false)
927 }
928
929 pub fn replay_facts(events: &[SessionEvent]) -> Result<Self, EventError> {
931 let mut projection = Self::default();
932 for event in events {
933 projection.apply_facts(event)?;
934 }
935 Ok(projection)
936 }
937
938 fn apply_internal(
939 &mut self,
940 envelope: &SessionEvent,
941 include_content: bool,
942 ) -> Result<(), EventError> {
943 if envelope.seq != self.last_seq + 1 {
944 return Err(EventError::Sequence {
945 expected: self.last_seq + 1,
946 actual: envelope.seq,
947 });
948 }
949 let session_id = self
950 .session_id
951 .get_or_insert_with(|| envelope.session_id.clone());
952 if *session_id != envelope.session_id {
953 return Err(EventError::SessionMismatch);
954 }
955 if envelope.format_version() != SESSION_EVENT_FORMAT_VERSION {
956 if envelope.ignorable() {
957 self.last_seq = envelope.seq;
958 return Ok(());
959 }
960 return Err(EventError::UnsupportedFormat(envelope.format_version()));
961 }
962 if self.deleted && !matches!(envelope.event, Event::UsageCorrected { .. }) {
963 return Err(EventError::SessionClosed);
964 }
965 match &envelope.event {
966 Event::UsageCorrected {
967 correction,
968 actor_id,
969 } => self.apply_metering_correction(correction, actor_id, envelope.seq)?,
970 Event::SessionCreated {
971 profile_revision_id,
972 } => {
973 if envelope.seq != 1 || self.profile_revision_id.is_some() {
974 return Err(EventError::DuplicateSession);
975 }
976 self.profile_revision_id = Some(profile_revision_id.clone());
977 }
978 Event::SessionMetadataUpdated { metadata } => {
979 metadata.validate()?;
980 if self.profile_revision_id.is_none()
981 || self.metadata.version.checked_add(1) != Some(metadata.version)
982 {
983 return Err(EventError::InvalidSessionMetadata);
984 }
985 self.metadata = metadata.clone();
986 }
987 Event::SessionDeleted { .. } => {
988 if let Some(run_id) = self.active_run_id.take() {
989 self.run_status
990 .insert(run_id, RunState::Terminal(RunStatus::Cancelled));
991 }
992 for (_, (run_id, _)) in std::mem::take(&mut self.queued_inputs) {
993 self.run_status
994 .insert(run_id, RunState::Terminal(RunStatus::Cancelled));
995 }
996 self.waiting_interaction_id = None;
997 self.open_turn = None;
998 self.open_steps.clear();
999 self.open_tool_calls.clear();
1000 self.started_tool_calls.clear();
1001 self.open_compaction = None;
1002 self.deleted = true;
1003 }
1004 Event::RunStarted { run_id, input_id } => {
1005 if self.active_run_id.is_some() {
1006 return Err(EventError::ConcurrentRun);
1007 }
1008 if self.claimed_inputs.get(input_id) != Some(run_id) {
1009 return Err(EventError::UnclaimedInput(input_id.clone()));
1010 }
1011 self.active_run_id = Some(run_id.clone());
1012 self.next_step = 1;
1013 self.run_status.insert(run_id.clone(), RunState::Running);
1014 }
1015 Event::RunWaiting {
1016 run_id,
1017 interaction_id,
1018 } => {
1019 self.require_active(run_id)?;
1020 self.waiting_interaction_id = Some(interaction_id.clone());
1021 self.run_status
1022 .insert(run_id.clone(), RunState::WaitingForInput);
1023 }
1024 Event::RunResumed {
1025 run_id,
1026 interaction_id,
1027 } => {
1028 self.require_active(run_id)?;
1029 if self.waiting_interaction_id.as_ref() != Some(interaction_id) {
1030 return Err(EventError::InteractionMismatch);
1031 }
1032 self.waiting_interaction_id = None;
1033 self.run_status.insert(run_id.clone(), RunState::Running);
1034 }
1035 Event::RunFinished { run_id, status, .. } => {
1036 self.require_active(run_id)?;
1037 if !self.open_tool_calls.is_empty()
1038 || !self.open_steps.is_empty()
1039 || self.open_turn.is_some()
1040 || self.open_compaction.is_some()
1041 || self
1042 .queued_inputs
1043 .values()
1044 .any(|(target_run_id, _)| target_run_id == run_id)
1045 {
1046 return Err(EventError::OpenLifecycle);
1047 }
1048 self.run_status
1049 .insert(run_id.clone(), RunState::Terminal(*status));
1050 self.active_run_id = None;
1051 self.waiting_interaction_id = None;
1052 }
1053 Event::UserMessage { run_id, content } if include_content => {
1054 self.messages.push(ProjectedMessage {
1055 role: "user",
1056 run_id: run_id.clone(),
1057 content: content.clone(),
1058 })
1059 }
1060 Event::AssistantMessage {
1061 run_id, content, ..
1062 } if include_content => self.messages.push(ProjectedMessage {
1063 role: "assistant",
1064 run_id: run_id.clone(),
1065 content: content.clone(),
1066 }),
1067 Event::ContextInjected {
1068 run_id,
1069 step,
1070 contribution_id,
1071 source,
1072 version,
1073 authority,
1074 form,
1075 content,
1076 } => {
1077 self.require_active(run_id)?;
1078 if !self.open_steps.contains(step) {
1079 return Err(EventError::LifecycleMismatch);
1080 }
1081 if include_content {
1082 self.injected_context.push(ProjectedContext {
1083 run_id: run_id.clone(),
1084 step: *step,
1085 contribution_id: contribution_id.clone(),
1086 source: source.clone(),
1087 version: version.clone(),
1088 authority: authority.clone(),
1089 form: form.clone(),
1090 content: content.clone(),
1091 });
1092 }
1093 }
1094 Event::InputQueued {
1095 input_id,
1096 run_id,
1097 mode,
1098 ..
1099 } => {
1100 if self.metadata.archived {
1101 return Err(EventError::SessionArchived);
1102 }
1103 if *mode != DeliveryMode::Followup {
1104 self.require_active(run_id)?;
1105 }
1106 if self
1107 .queued_inputs
1108 .insert(input_id.clone(), (run_id.clone(), *mode))
1109 .is_some()
1110 {
1111 return Err(EventError::DuplicateInput(input_id.clone()));
1112 }
1113 }
1114 Event::InputClaimed { input_id, run_id } => {
1115 if self
1116 .queued_inputs
1117 .remove(input_id)
1118 .map(|value| value.0)
1119 .as_ref()
1120 != Some(run_id)
1121 || self
1122 .claimed_inputs
1123 .insert(input_id.clone(), run_id.clone())
1124 .is_some()
1125 {
1126 return Err(EventError::UnqueuedInput(input_id.clone()));
1127 }
1128 }
1129 Event::InputCancelled {
1130 input_id,
1131 run_id,
1132 error_code,
1133 } => {
1134 if self
1135 .queued_inputs
1136 .remove(input_id)
1137 .map(|value| value.0)
1138 .as_ref()
1139 != Some(run_id)
1140 {
1141 return Err(EventError::UnqueuedInput(input_id.clone()));
1142 }
1143 self.run_status.insert(
1144 run_id.clone(),
1145 RunState::Terminal(if error_code == "cancelled" {
1146 RunStatus::Cancelled
1147 } else {
1148 RunStatus::Failed
1149 }),
1150 );
1151 }
1152 Event::TurnStarted { run_id, turn } => {
1153 self.require_active(run_id)?;
1154 if self.open_turn.replace(*turn).is_some() {
1155 return Err(EventError::ConcurrentTurn);
1156 }
1157 }
1158 Event::TurnFinished { run_id, turn } => {
1159 self.require_active(run_id)?;
1160 if self.open_turn != Some(*turn)
1161 || !self.open_steps.is_empty()
1162 || !self.open_tool_calls.is_empty()
1163 {
1164 return Err(EventError::LifecycleMismatch);
1165 }
1166 self.open_turn = None;
1167 }
1168 Event::StepStarted { run_id, step } => {
1169 self.require_active(run_id)?;
1170 if self.open_turn.is_none()
1171 || !self.open_steps.is_empty()
1172 || *step != self.next_step
1173 || !self.open_steps.insert(*step)
1174 {
1175 return Err(EventError::ConcurrentStep);
1176 }
1177 }
1178 Event::StepFinished { run_id, step } => {
1179 self.require_active(run_id)?;
1180 if self
1181 .open_tool_calls
1182 .values()
1183 .any(|call| call.run_id == *run_id && call.step == *step)
1184 || !self.open_steps.remove(step)
1185 {
1186 return Err(EventError::LifecycleMismatch);
1187 }
1188 self.next_step = step.saturating_add(1);
1189 }
1190 Event::ToolCall {
1191 run_id,
1192 step,
1193 call_id,
1194 tool,
1195 arguments,
1196 } => {
1197 self.require_active(run_id)?;
1198 if !self.open_steps.contains(step) {
1199 return Err(EventError::LifecycleMismatch);
1200 }
1201 if !self.seen_tool_calls.insert(call_id.clone())
1202 || self
1203 .open_tool_calls
1204 .insert(
1205 call_id.clone(),
1206 OpenToolCall {
1207 run_id: run_id.clone(),
1208 step: *step,
1209 tool: tool.clone(),
1210 arguments: arguments.clone(),
1211 source_event_seq: envelope.seq,
1212 },
1213 )
1214 .is_some()
1215 {
1216 return Err(EventError::DuplicateToolCall(call_id.clone()));
1217 }
1218 }
1219 Event::ToolResult {
1220 run_id,
1221 step,
1222 call_id,
1223 ..
1224 } => {
1225 self.require_active(run_id)?;
1226 if !self.open_steps.contains(step) {
1227 return Err(EventError::LifecycleMismatch);
1228 }
1229 let Some(call) = self.open_tool_calls.get(call_id) else {
1230 return Err(EventError::OrphanToolResult(call_id.clone()));
1231 };
1232 if call.run_id != *run_id || call.step != *step {
1233 return Err(EventError::ToolResultMismatch(call_id.clone()));
1234 }
1235 self.open_tool_calls.remove(call_id);
1236 self.started_tool_calls.remove(call_id);
1237 }
1238 Event::ToolExecutionStarted {
1239 run_id,
1240 step,
1241 call_id,
1242 metering,
1243 reserved_cost_units,
1244 } => {
1245 self.require_active(run_id)?;
1246 let Some(call) = self.open_tool_calls.get(call_id) else {
1247 return Err(EventError::OrphanToolResult(call_id.clone()));
1248 };
1249 if call.run_id != *run_id
1250 || call.step != *step
1251 || self.started_tool_calls.contains(call_id)
1252 {
1253 return Err(EventError::ToolResultMismatch(call_id.clone()));
1254 }
1255 if let Some(details) = metering {
1256 details.validate()?;
1257 if !matches!(details, MeteringDetails::Tool { call_id: id, name, source: MeteringSource::Estimated, outcome: MeteringOutcome::Unknown } if id == call_id && name == &call.tool)
1258 {
1259 return Err(EventError::InvalidMetering);
1260 }
1261 }
1262 self.started_tool_calls.insert(call_id.clone());
1263 self.pending_usage_operations.insert(
1264 (run_id.clone(), format!("tool:{call_id}")),
1265 (0, 0, *reserved_cost_units, metering.clone()),
1266 );
1267 }
1268 Event::ModelRequestPrepared {
1269 run_id,
1270 operation_id,
1271 reserved_prompt_tokens,
1272 reserved_completion_tokens,
1273 metering,
1274 ..
1275 } if !operation_id.is_empty() => {
1276 self.require_active(run_id)?;
1277 let key = (run_id.clone(), operation_id.clone());
1278 if !self.usage_operations.contains_key(&key) {
1279 if let Some(details) = metering {
1280 details.validate()?;
1281 if !matches!(
1282 details,
1283 MeteringDetails::Model {
1284 source: MeteringSource::Estimated,
1285 outcome: MeteringOutcome::Unknown,
1286 ..
1287 }
1288 ) {
1289 return Err(EventError::InvalidMetering);
1290 }
1291 }
1292 let reservation = (
1293 *reserved_prompt_tokens,
1294 *reserved_completion_tokens,
1295 0,
1296 metering.clone(),
1297 );
1298 match self.pending_usage_operations.get(&key) {
1299 Some(existing) if *existing != reservation => {
1300 return Err(EventError::UsageConflict(operation_id.clone()));
1301 }
1302 Some(_) => {}
1303 None => {
1304 self.pending_usage_operations.insert(key, reservation);
1305 }
1306 }
1307 }
1308 }
1309 Event::UsageRecorded {
1310 run_id,
1311 operation_id,
1312 prompt_tokens,
1313 completion_tokens,
1314 cost_units,
1315 metering,
1316 } => {
1317 self.require_active(run_id)?;
1318 if let Some(details) = metering {
1319 details.validate()?;
1320 }
1321 let key = (run_id.clone(), operation_id.clone());
1322 if let Some((_, _, _, Some(prepared))) = self.pending_usage_operations.get(&key) {
1323 if !metering
1324 .as_ref()
1325 .is_some_and(|details| details.completes(prepared))
1326 {
1327 return Err(EventError::UsageConflict(operation_id.clone()));
1328 }
1329 }
1330 let usage = (
1331 *prompt_tokens,
1332 *completion_tokens,
1333 *cost_units,
1334 metering.clone(),
1335 );
1336 match self.usage_operations.get(&key) {
1337 Some(existing) if *existing != usage => {
1338 return Err(EventError::UsageConflict(operation_id.clone()));
1339 }
1340 Some(_) => {}
1341 None => {
1342 self.pending_usage_operations.remove(&key);
1343 self.usage_operations.insert(key, usage);
1344 }
1345 }
1346 }
1347 Event::CompactionStarted {
1348 run_id,
1349 compaction_id,
1350 source_through_seq,
1351 } => {
1352 self.require_active(run_id)?;
1353 if self.open_compaction.is_some() {
1354 return Err(EventError::ConcurrentCompaction);
1355 }
1356 self.open_compaction = Some((compaction_id.clone(), *source_through_seq));
1357 }
1358 Event::CompactionFinished {
1359 run_id,
1360 compaction_id,
1361 ..
1362 } => {
1363 self.require_active(run_id)?;
1364 if self.open_compaction.as_ref().map(|value| value.0.as_str())
1365 != Some(compaction_id.as_str())
1366 {
1367 return Err(EventError::CompactionMismatch);
1368 }
1369 self.open_compaction = None;
1370 }
1371 Event::SummaryReplaced { summary, .. } => self.summary = Some(summary.clone()),
1372 Event::Opaque {
1373 event_type,
1374 ignorable: false,
1375 ..
1376 } => return Err(EventError::UnknownRequired(event_type.clone())),
1377 _ => {}
1378 }
1379 self.last_seq = envelope.seq;
1380 Ok(())
1381 }
1382
1383 fn require_active(&self, run_id: &RunId) -> Result<(), EventError> {
1384 if self.active_run_id.as_ref() == Some(run_id) {
1385 Ok(())
1386 } else {
1387 Err(EventError::RunMismatch)
1388 }
1389 }
1390
1391 pub fn usage_for(&self, run_id: &str) -> (u64, u64) {
1393 self.usage_operations
1394 .iter()
1395 .filter(|((recorded_run_id, _), _)| recorded_run_id.as_str() == run_id)
1396 .fold((0, 0), |total, (_, usage)| {
1397 (total.0 + usage.0, total.1 + usage.1)
1398 })
1399 }
1400
1401 pub fn usage_charge_for(&self, run_id: &str) -> UsageCharge {
1403 let sum = |pending: bool| {
1404 let operations = if pending {
1405 &self.pending_usage_operations
1406 } else {
1407 &self.usage_operations
1408 };
1409 operations
1410 .iter()
1411 .filter(|((recorded_run_id, _), _)| recorded_run_id.as_str() == run_id)
1412 .fold(UsageCharge::default(), |mut total, (_, usage)| {
1413 total.prompt_tokens = total.prompt_tokens.saturating_add(usage.0);
1414 total.completion_tokens = total.completion_tokens.saturating_add(usage.1);
1415 total.cost_units = total.cost_units.saturating_add(usage.2);
1416 total
1417 })
1418 };
1419 let recorded = sum(false);
1420 let pending = sum(true);
1421 UsageCharge {
1422 prompt_tokens: recorded.prompt_tokens.saturating_add(pending.prompt_tokens),
1423 completion_tokens: recorded
1424 .completion_tokens
1425 .saturating_add(pending.completion_tokens),
1426 cost_units: recorded.cost_units.saturating_add(pending.cost_units),
1427 }
1428 }
1429
1430 pub fn billable_units_for(&self, run_id: &str) -> u64 {
1432 self.usage_charge_for(run_id).total()
1433 }
1434}
1435
1436#[derive(Debug, thiserror::Error, PartialEq, Eq)]
1438pub enum EventError {
1439 #[error("invalid Session metadata or metadata version conflict")]
1441 InvalidSessionMetadata,
1442 #[error("invalid metering attribution")]
1444 InvalidMetering,
1445 #[error("Session is archived")]
1447 SessionArchived,
1448 #[error("unsupported session event format version {0}")]
1450 UnsupportedFormat(u32),
1451 #[error("unknown required session event type {0}")]
1453 UnknownRequired(String),
1454 #[error("event sequence mismatch: expected {expected}, got {actual}")]
1456 Sequence {
1457 expected: u64,
1459 actual: u64,
1461 },
1462 #[error("event belongs to another session")]
1464 SessionMismatch,
1465 #[error("session creation must be the first and only creation event")]
1467 DuplicateSession,
1468 #[error("session already has an active run")]
1470 ConcurrentRun,
1471 #[error("session is closed")]
1473 SessionClosed,
1474 #[error("event does not match the active run")]
1476 RunMismatch,
1477 #[error("interaction does not match the waiting run")]
1479 InteractionMismatch,
1480 #[error("run cannot finish with an open turn, step or tool call")]
1482 OpenLifecycle,
1483 #[error("input was queued twice: {0}")]
1485 DuplicateInput(InputId),
1486 #[error("input was claimed before it was queued: {0}")]
1488 UnqueuedInput(InputId),
1489 #[error("run started from an unclaimed input: {0}")]
1491 UnclaimedInput(InputId),
1492 #[error("session already has an active turn")]
1494 ConcurrentTurn,
1495 #[error("turn already has this active step")]
1497 ConcurrentStep,
1498 #[error("turn or step lifecycle does not pair")]
1500 LifecycleMismatch,
1501 #[error("duplicate tool call {0}")]
1503 DuplicateToolCall(ToolCallId),
1504 #[error("tool result has no matching call {0}")]
1506 OrphanToolResult(ToolCallId),
1507 #[error("tool result does not match the call run and step: {0}")]
1509 ToolResultMismatch(ToolCallId),
1510 #[error("usage operation was recorded with different totals: {0}")]
1512 UsageConflict(String),
1513 #[error("session already has an active compaction")]
1515 ConcurrentCompaction,
1516 #[error("compaction lifecycle does not pair")]
1518 CompactionMismatch,
1519 #[error("event store conflict: {0}")]
1521 Conflict(String),
1522 #[error("event store unavailable: {0}")]
1524 Unavailable(String),
1525}
1526
1527#[async_trait]
1529pub trait SessionEventStore: Send + Sync {
1530 async fn append(
1532 &self,
1533 tenant_id: &str,
1534 session_id: &str,
1535 expected_seq: u64,
1536 events: Vec<Event>,
1537 ) -> Result<Vec<SessionEvent>, EventError>;
1538 async fn load(
1540 &self,
1541 tenant_id: &str,
1542 session_id: &str,
1543 after_seq: u64,
1544 ) -> Result<Vec<SessionEvent>, EventError>;
1545}
1546
1547pub fn text(value: impl Into<String>) -> Vec<ContentBlock> {
1549 vec![ContentBlock::Text { text: value.into() }]
1550}
1551
1552pub fn recovery_events(projection: &SessionProjection) -> Vec<Event> {
1554 failure_events(projection, "worker_restarted")
1555}
1556
1557pub fn failure_events(projection: &SessionProjection, error_code: &str) -> Vec<Event> {
1559 termination_events(projection, RunStatus::Failed, error_code)
1560}
1561
1562pub fn cancel_events(projection: &SessionProjection) -> Vec<Event> {
1564 termination_events(projection, RunStatus::Cancelled, "cancelled")
1565}
1566
1567pub fn session_deletion_events(projection: &SessionProjection, reason: &str) -> Vec<Event> {
1569 let active = projection.active_run_id.as_deref();
1570 let mut events = cancel_events(projection);
1571 events.extend(
1572 projection
1573 .queued_inputs
1574 .iter()
1575 .filter(|(_, (run_id, _))| Some(run_id.as_str()) != active)
1576 .map(|(input_id, (run_id, _))| Event::InputCancelled {
1577 input_id: input_id.clone(),
1578 run_id: run_id.clone(),
1579 error_code: "cancelled".into(),
1580 }),
1581 );
1582 events.push(Event::SessionDeleted {
1583 reason: reason.into(),
1584 });
1585 events
1586}
1587
1588fn termination_events(
1589 projection: &SessionProjection,
1590 status: RunStatus,
1591 error_code: &str,
1592) -> Vec<Event> {
1593 let Some(run_id) = &projection.active_run_id else {
1594 return Vec::new();
1595 };
1596 let mut events = projection
1597 .queued_inputs
1598 .iter()
1599 .filter(|(_, (target_run_id, _))| target_run_id == run_id)
1600 .map(|(input_id, _)| Event::InputCancelled {
1601 input_id: input_id.clone(),
1602 run_id: run_id.clone(),
1603 error_code: error_code.into(),
1604 })
1605 .collect::<Vec<_>>();
1606 events.extend(
1607 projection
1608 .open_tool_calls
1609 .iter()
1610 .map(|(call_id, call)| Event::ToolResult {
1611 run_id: run_id.clone(),
1612 step: call.step,
1613 call_id: call_id.clone(),
1614 result: termination_result(error_code),
1615 is_error: true,
1616 }),
1617 );
1618 if let Some((compaction_id, _)) = &projection.open_compaction {
1619 events.push(Event::CompactionFinished {
1620 run_id: run_id.clone(),
1621 compaction_id: compaction_id.clone(),
1622 status: "failed".into(),
1623 error: Some(error_code.into()),
1624 });
1625 }
1626 events.extend(
1627 projection
1628 .open_steps
1629 .iter()
1630 .map(|step| Event::StepFinished {
1631 run_id: run_id.clone(),
1632 step: *step,
1633 }),
1634 );
1635 if let Some(turn) = projection.open_turn {
1636 events.push(Event::TurnFinished {
1637 run_id: run_id.clone(),
1638 turn,
1639 });
1640 }
1641 events.push(Event::RunFinished {
1642 run_id: run_id.clone(),
1643 status,
1644 error_code: Some(error_code.into()),
1645 });
1646 events
1647}
1648
1649fn termination_result(error_code: &str) -> Value {
1650 if error_code == "tool_outcome_unknown" || error_code == "worker_restarted" {
1651 serde_json::json!({
1652 "error": error_code,
1653 "guidance": "The tool outcome is unknown. Verify external state before retrying any operation with side effects; ask the user when verification is unavailable."
1654 })
1655 } else {
1656 serde_json::json!({"error":error_code})
1657 }
1658}
1659
1660#[cfg(test)]
1661mod tests {
1662 use super::*;
1663
1664 #[test]
1665 fn legacy_tool_and_compaction_events_still_decode() {
1666 let tool: Event = serde_json::from_value(serde_json::json!({
1667 "type": "tool_result",
1668 "run_id": "run",
1669 "step": 1,
1670 "call_id": "call",
1671 "result": {"ok": true},
1672 "is_error": false
1673 }))
1674 .unwrap();
1675 assert!(matches!(tool, Event::ToolResult { .. }));
1676
1677 let compaction: Event = serde_json::from_value(serde_json::json!({
1678 "type": "compaction_finished",
1679 "run_id": "run",
1680 "compaction_id": "compaction:1",
1681 "status": "completed",
1682 "error": null
1683 }))
1684 .unwrap();
1685 assert!(matches!(compaction, Event::CompactionFinished { .. }));
1686 }
1687
1688 fn event(seq: u64, event: Event) -> SessionEvent {
1689 SessionEvent {
1690 session_id: "s".parse().unwrap(),
1691 seq,
1692 occurred_at: Utc::now(),
1693 event,
1694 }
1695 }
1696
1697 #[test]
1698 fn replay_enforces_single_run_and_tool_pairs() {
1699 let events = vec![
1700 event(
1701 1,
1702 Event::SessionCreated {
1703 profile_revision_id: "p1".parse().unwrap(),
1704 },
1705 ),
1706 event(
1707 2,
1708 Event::InputQueued {
1709 input_id: "i1".parse().unwrap(),
1710 run_id: "r1".parse().unwrap(),
1711 mode: DeliveryMode::Followup,
1712 content: text("hi"),
1713 explicit_skill: None,
1714 },
1715 ),
1716 event(
1717 3,
1718 Event::InputClaimed {
1719 input_id: "i1".parse().unwrap(),
1720 run_id: "r1".parse().unwrap(),
1721 },
1722 ),
1723 event(
1724 4,
1725 Event::RunStarted {
1726 run_id: "r1".parse().unwrap(),
1727 input_id: "i1".parse().unwrap(),
1728 },
1729 ),
1730 event(
1731 5,
1732 Event::TurnStarted {
1733 run_id: "r1".parse().unwrap(),
1734 turn: 1,
1735 },
1736 ),
1737 event(
1738 6,
1739 Event::StepStarted {
1740 run_id: "r1".parse().unwrap(),
1741 step: 1,
1742 },
1743 ),
1744 event(
1745 7,
1746 Event::ToolCall {
1747 run_id: "r1".parse().unwrap(),
1748 step: 1,
1749 call_id: "c1".parse().unwrap(),
1750 tool: "echo".into(),
1751 arguments: serde_json::json!({"x":1}),
1752 },
1753 ),
1754 event(
1755 8,
1756 Event::ToolResult {
1757 run_id: "r1".parse().unwrap(),
1758 step: 1,
1759 call_id: "c1".parse().unwrap(),
1760 result: serde_json::json!({"x":1}),
1761 is_error: false,
1762 },
1763 ),
1764 event(
1765 9,
1766 Event::StepFinished {
1767 run_id: "r1".parse().unwrap(),
1768 step: 1,
1769 },
1770 ),
1771 event(
1772 10,
1773 Event::TurnFinished {
1774 run_id: "r1".parse().unwrap(),
1775 turn: 1,
1776 },
1777 ),
1778 event(
1779 11,
1780 Event::RunFinished {
1781 run_id: "r1".parse().unwrap(),
1782 status: RunStatus::Completed,
1783 error_code: None,
1784 },
1785 ),
1786 ];
1787 let projection = SessionProjection::replay(&events).unwrap();
1788 assert_eq!(projection.last_seq, 11);
1789 assert!(projection.active_run_id.is_none());
1790 }
1791
1792 #[test]
1793 fn replay_rejects_orphan_tool_result() {
1794 let events = vec![
1795 event(
1796 1,
1797 Event::SessionCreated {
1798 profile_revision_id: "p1".parse().unwrap(),
1799 },
1800 ),
1801 event(
1802 2,
1803 Event::InputQueued {
1804 input_id: "i1".parse().unwrap(),
1805 run_id: "r1".parse().unwrap(),
1806 mode: DeliveryMode::Followup,
1807 content: text("hi"),
1808 explicit_skill: None,
1809 },
1810 ),
1811 event(
1812 3,
1813 Event::InputClaimed {
1814 input_id: "i1".parse().unwrap(),
1815 run_id: "r1".parse().unwrap(),
1816 },
1817 ),
1818 event(
1819 4,
1820 Event::RunStarted {
1821 run_id: "r1".parse().unwrap(),
1822 input_id: "i1".parse().unwrap(),
1823 },
1824 ),
1825 event(
1826 5,
1827 Event::TurnStarted {
1828 run_id: "r1".parse().unwrap(),
1829 turn: 1,
1830 },
1831 ),
1832 event(
1833 6,
1834 Event::StepStarted {
1835 run_id: "r1".parse().unwrap(),
1836 step: 1,
1837 },
1838 ),
1839 event(
1840 7,
1841 Event::ToolResult {
1842 run_id: "r1".parse().unwrap(),
1843 step: 1,
1844 call_id: "missing".parse().unwrap(),
1845 result: Value::Null,
1846 is_error: true,
1847 },
1848 ),
1849 ];
1850 assert_eq!(
1851 SessionProjection::replay(&events).unwrap_err(),
1852 EventError::OrphanToolResult("missing".parse().unwrap())
1853 );
1854 }
1855
1856 #[test]
1857 fn queued_input_failure_is_not_projected_as_cancellation() {
1858 let events = vec![
1859 event(
1860 1,
1861 Event::SessionCreated {
1862 profile_revision_id: "p1".parse().unwrap(),
1863 },
1864 ),
1865 event(
1866 2,
1867 Event::InputQueued {
1868 input_id: "i1".parse().unwrap(),
1869 run_id: "r1".parse().unwrap(),
1870 mode: DeliveryMode::Followup,
1871 content: text("hi"),
1872 explicit_skill: None,
1873 },
1874 ),
1875 event(
1876 3,
1877 Event::InputCancelled {
1878 input_id: "i1".parse().unwrap(),
1879 run_id: "r1".parse().unwrap(),
1880 error_code: "profile_not_found".into(),
1881 },
1882 ),
1883 ];
1884 let projection = SessionProjection::replay(&events).unwrap();
1885 assert_eq!(
1886 projection.run_status.get("r1").map(|state| state.as_str()),
1887 Some("failed")
1888 );
1889 }
1890
1891 #[test]
1892 fn session_deletion_closes_active_and_queued_runs_before_tombstone() {
1893 let mut events = vec![
1894 event(
1895 1,
1896 Event::SessionCreated {
1897 profile_revision_id: "p1".parse().unwrap(),
1898 },
1899 ),
1900 event(
1901 2,
1902 Event::InputQueued {
1903 input_id: "i1".parse().unwrap(),
1904 run_id: "r1".parse().unwrap(),
1905 mode: DeliveryMode::Followup,
1906 content: text("start"),
1907 explicit_skill: None,
1908 },
1909 ),
1910 event(
1911 3,
1912 Event::InputClaimed {
1913 input_id: "i1".parse().unwrap(),
1914 run_id: "r1".parse().unwrap(),
1915 },
1916 ),
1917 event(
1918 4,
1919 Event::RunStarted {
1920 run_id: "r1".parse().unwrap(),
1921 input_id: "i1".parse().unwrap(),
1922 },
1923 ),
1924 event(
1925 5,
1926 Event::TurnStarted {
1927 run_id: "r1".parse().unwrap(),
1928 turn: 1,
1929 },
1930 ),
1931 event(
1932 6,
1933 Event::InputQueued {
1934 input_id: "i2".parse().unwrap(),
1935 run_id: "r2".parse().unwrap(),
1936 mode: DeliveryMode::Followup,
1937 content: text("later"),
1938 explicit_skill: None,
1939 },
1940 ),
1941 ];
1942 let projection = SessionProjection::replay(&events).unwrap();
1943 for event_value in session_deletion_events(&projection, "api_deleted") {
1944 let seq = events.len() as u64 + 1;
1945 events.push(event(seq, event_value));
1946 }
1947 let deleted = SessionProjection::replay(&events).unwrap();
1948 assert!(deleted.deleted);
1949 assert_eq!(
1950 deleted.run_status.get("r1").map(|state| state.as_str()),
1951 Some("cancelled")
1952 );
1953 assert_eq!(
1954 deleted.run_status.get("r2").map(|state| state.as_str()),
1955 Some("cancelled")
1956 );
1957 assert!(events.iter().any(|event| matches!(
1958 &event.event,
1959 Event::RunFinished { run_id, status: RunStatus::Cancelled, .. } if run_id == "r1"
1960 )));
1961 assert!(matches!(
1962 events.last().map(|event| &event.event),
1963 Some(Event::SessionDeleted { .. })
1964 ));
1965 }
1966
1967 #[test]
1968 fn usage_operations_are_idempotent_and_conflicts_fail_replay() {
1969 let mut events = vec![
1970 event(
1971 1,
1972 Event::SessionCreated {
1973 profile_revision_id: "p1".parse().unwrap(),
1974 },
1975 ),
1976 event(
1977 2,
1978 Event::InputQueued {
1979 input_id: "i1".parse().unwrap(),
1980 run_id: "r1".parse().unwrap(),
1981 mode: DeliveryMode::Followup,
1982 content: text("hi"),
1983 explicit_skill: None,
1984 },
1985 ),
1986 event(
1987 3,
1988 Event::InputClaimed {
1989 input_id: "i1".parse().unwrap(),
1990 run_id: "r1".parse().unwrap(),
1991 },
1992 ),
1993 event(
1994 4,
1995 Event::RunStarted {
1996 run_id: "r1".parse().unwrap(),
1997 input_id: "i1".parse().unwrap(),
1998 },
1999 ),
2000 event(
2001 5,
2002 Event::UsageRecorded {
2003 metering: None,
2004 run_id: "r1".parse().unwrap(),
2005 operation_id: "model:1:attempt:1".into(),
2006 prompt_tokens: 7,
2007 completion_tokens: 3,
2008 cost_units: 5,
2009 },
2010 ),
2011 event(
2012 6,
2013 Event::UsageRecorded {
2014 metering: None,
2015 run_id: "r1".parse().unwrap(),
2016 operation_id: "model:1:attempt:1".into(),
2017 prompt_tokens: 7,
2018 completion_tokens: 3,
2019 cost_units: 5,
2020 },
2021 ),
2022 ];
2023 assert_eq!(
2024 SessionProjection::replay(&events).unwrap().usage_for("r1"),
2025 (7, 3)
2026 );
2027 assert_eq!(
2028 SessionProjection::replay(&events)
2029 .unwrap()
2030 .billable_units_for("r1"),
2031 15
2032 );
2033 assert_eq!(
2034 SessionProjection::replay(&events)
2035 .unwrap()
2036 .usage_charge_for("r1"),
2037 UsageCharge {
2038 prompt_tokens: 7,
2039 completion_tokens: 3,
2040 cost_units: 5,
2041 }
2042 );
2043
2044 let mut attributed = events.clone();
2045 let details = MeteringDetails::Model {
2046 model: "m".into(),
2047 provider: Some("p".into()),
2048 provider_attempt_id: "r1:model:1:attempt:1".parse().unwrap(),
2049 source: MeteringSource::Estimated,
2050 outcome: MeteringOutcome::Unknown,
2051 };
2052 for entry in &mut attributed[4..] {
2053 if let Event::UsageRecorded { metering, .. } = &mut entry.event {
2054 *metering = Some(details.clone());
2055 }
2056 }
2057 assert_eq!(
2058 SessionProjection::replay(&attributed)
2059 .unwrap()
2060 .billable_units_for("r1"),
2061 15
2062 );
2063 if let Event::UsageRecorded {
2064 metering: Some(MeteringDetails::Model { source, .. }),
2065 ..
2066 } = &mut attributed[5].event
2067 {
2068 *source = MeteringSource::Reported;
2069 }
2070 assert!(matches!(
2071 SessionProjection::replay(&attributed),
2072 Err(EventError::UsageConflict(_))
2073 ));
2074 let mut prepared = attributed[..4].to_vec();
2075 prepared.push(event(
2076 5,
2077 Event::ModelRequestPrepared {
2078 metering: Some(details.clone()),
2079 run_id: "r1".parse().unwrap(),
2080 step: 1,
2081 attempt: 1,
2082 provider_attempt_id: "r1:model:1:attempt:1".into(),
2083 operation_id: "model:1:attempt:1".into(),
2084 reserved_prompt_tokens: 10,
2085 reserved_completion_tokens: 20,
2086 request: Value::Null,
2087 prompt_sections: Value::Null,
2088 },
2089 ));
2090 prepared.push(attributed[5].clone());
2091 assert_eq!(
2092 SessionProjection::replay(&prepared)
2093 .unwrap()
2094 .billable_units_for("r1"),
2095 15
2096 );
2097 if let Event::UsageRecorded {
2098 metering: Some(MeteringDetails::Model { model, .. }),
2099 ..
2100 } = &mut prepared[5].event
2101 {
2102 *model = "wrong".into();
2103 }
2104 assert!(matches!(
2105 SessionProjection::replay(&prepared),
2106 Err(EventError::UsageConflict(_))
2107 ));
2108 if let Event::UsageRecorded { metering, .. } = &mut prepared[5].event {
2109 *metering = None;
2110 }
2111 assert!(matches!(
2112 SessionProjection::replay(&prepared),
2113 Err(EventError::UsageConflict(_))
2114 ));
2115 let serialized = serde_json::to_value(&events[4]).unwrap();
2116 assert!(serialized["event"].get("metering").is_none());
2117 assert_eq!(
2118 serde_json::from_value::<SessionEvent>(serialized).unwrap(),
2119 events[4]
2120 );
2121
2122 events.push(event(
2123 7,
2124 Event::UsageRecorded {
2125 metering: None,
2126 run_id: "r1".parse().unwrap(),
2127 operation_id: "model:1:attempt:1".into(),
2128 prompt_tokens: 8,
2129 completion_tokens: 3,
2130 cost_units: 0,
2131 },
2132 ));
2133 assert_eq!(
2134 SessionProjection::replay(&events).unwrap_err(),
2135 EventError::UsageConflict("model:1:attempt:1".into())
2136 );
2137 }
2138
2139 #[test]
2140 fn unresolved_provider_attempt_is_conservatively_billable_and_reconcilable() {
2141 let mut events = vec![
2142 event(
2143 1,
2144 Event::SessionCreated {
2145 profile_revision_id: "p1".parse().unwrap(),
2146 },
2147 ),
2148 event(
2149 2,
2150 Event::InputQueued {
2151 input_id: "i1".parse().unwrap(),
2152 run_id: "r1".parse().unwrap(),
2153 mode: DeliveryMode::Followup,
2154 content: text("hi"),
2155 explicit_skill: None,
2156 },
2157 ),
2158 event(
2159 3,
2160 Event::InputClaimed {
2161 input_id: "i1".parse().unwrap(),
2162 run_id: "r1".parse().unwrap(),
2163 },
2164 ),
2165 event(
2166 4,
2167 Event::RunStarted {
2168 run_id: "r1".parse().unwrap(),
2169 input_id: "i1".parse().unwrap(),
2170 },
2171 ),
2172 event(
2173 5,
2174 Event::ModelRequestPrepared {
2175 metering: None,
2176 run_id: "r1".parse().unwrap(),
2177 step: 1,
2178 attempt: 1,
2179 provider_attempt_id: "r1:model:1:attempt:1".into(),
2180 operation_id: "model:1:attempt:1".into(),
2181 reserved_prompt_tokens: 7,
2182 reserved_completion_tokens: 11,
2183 request: Value::Null,
2184 prompt_sections: Value::Null,
2185 },
2186 ),
2187 ];
2188 assert_eq!(
2189 SessionProjection::replay(&events)
2190 .unwrap()
2191 .billable_units_for("r1"),
2192 18
2193 );
2194 events.push(event(
2195 6,
2196 Event::UsageRecorded {
2197 metering: None,
2198 run_id: "r1".parse().unwrap(),
2199 operation_id: "model:1:attempt:1".into(),
2200 prompt_tokens: 6,
2201 completion_tokens: 2,
2202 cost_units: 0,
2203 },
2204 ));
2205 assert_eq!(
2206 SessionProjection::replay(&events)
2207 .unwrap()
2208 .billable_units_for("r1"),
2209 8
2210 );
2211 }
2212
2213 #[test]
2214 fn a3_incremental_facts_preserve_every_lifecycle_and_usage_invariant() {
2215 let mut events = open_step_events();
2216 for fact in [
2217 Event::UserMessage {
2218 run_id: "r1".parse().unwrap(),
2219 content: text("history".repeat(10_000)),
2220 },
2221 Event::UsageRecorded {
2222 metering: None,
2223 run_id: "r1".parse().unwrap(),
2224 operation_id: "attempt".into(),
2225 prompt_tokens: 5,
2226 completion_tokens: 3,
2227 cost_units: 0,
2228 },
2229 Event::ToolCall {
2230 run_id: "r1".parse().unwrap(),
2231 step: 1,
2232 call_id: "pending".parse().unwrap(),
2233 tool: "echo".into(),
2234 arguments: Value::Null,
2235 },
2236 ] {
2237 events.push(event(events.len() as u64 + 1, fact));
2238 }
2239 let mut full = SessionProjection::replay(&events).unwrap();
2240 assert!(!full.messages.is_empty());
2241 full.messages.clear();
2242 full.injected_context.clear();
2243 let mut folded = SessionProjection::default();
2244 for page in events.chunks(3) {
2245 for fact in page {
2246 folded.apply_facts(fact).unwrap();
2247 }
2248 assert!(folded.messages.is_empty());
2249 }
2250 assert_eq!(folded, full);
2251 assert_eq!(SessionProjection::replay_facts(&events).unwrap(), full);
2252 assert_eq!(cancel_events(&folded), cancel_events(&full));
2253 assert_eq!(folded.usage_for("r1"), (5, 3));
2254 let invalid = event(
2255 events.len() as u64 + 1,
2256 Event::ToolResult {
2257 run_id: "r1".parse().unwrap(),
2258 step: 1,
2259 call_id: "missing".parse().unwrap(),
2260 result: Value::Null,
2261 is_error: true,
2262 },
2263 );
2264 assert_eq!(folded.apply_facts(&invalid), full.apply(&invalid));
2265 }
2266
2267 fn open_step_events() -> Vec<SessionEvent> {
2268 vec![
2269 event(
2270 1,
2271 Event::SessionCreated {
2272 profile_revision_id: "p1".parse().unwrap(),
2273 },
2274 ),
2275 event(
2276 2,
2277 Event::InputQueued {
2278 input_id: "i1".parse().unwrap(),
2279 run_id: "r1".parse().unwrap(),
2280 mode: DeliveryMode::Followup,
2281 content: text("hi"),
2282 explicit_skill: None,
2283 },
2284 ),
2285 event(
2286 3,
2287 Event::InputClaimed {
2288 input_id: "i1".parse().unwrap(),
2289 run_id: "r1".parse().unwrap(),
2290 },
2291 ),
2292 event(
2293 4,
2294 Event::RunStarted {
2295 run_id: "r1".parse().unwrap(),
2296 input_id: "i1".parse().unwrap(),
2297 },
2298 ),
2299 event(
2300 5,
2301 Event::TurnStarted {
2302 run_id: "r1".parse().unwrap(),
2303 turn: 1,
2304 },
2305 ),
2306 event(
2307 6,
2308 Event::StepStarted {
2309 run_id: "r1".parse().unwrap(),
2310 step: 1,
2311 },
2312 ),
2313 ]
2314 }
2315
2316 #[test]
2317 fn replay_rejects_overlapping_or_out_of_order_steps() {
2318 let mut overlapping = open_step_events();
2319 overlapping.push(event(
2320 7,
2321 Event::StepStarted {
2322 run_id: "r1".parse().unwrap(),
2323 step: 2,
2324 },
2325 ));
2326 assert_eq!(
2327 SessionProjection::replay(&overlapping).unwrap_err(),
2328 EventError::ConcurrentStep
2329 );
2330
2331 let mut skipped = open_step_events();
2332 skipped[5] = event(
2333 6,
2334 Event::StepStarted {
2335 run_id: "r1".parse().unwrap(),
2336 step: 2,
2337 },
2338 );
2339 assert_eq!(
2340 SessionProjection::replay(&skipped).unwrap_err(),
2341 EventError::ConcurrentStep
2342 );
2343 }
2344
2345 #[test]
2346 fn replay_rejects_cross_step_results_and_dangling_calls() {
2347 let mut events = open_step_events();
2348 events.push(event(
2349 7,
2350 Event::ToolCall {
2351 run_id: "r1".parse().unwrap(),
2352 step: 1,
2353 call_id: "c1".parse().unwrap(),
2354 tool: "echo".into(),
2355 arguments: Value::Null,
2356 },
2357 ));
2358 events.push(event(
2359 8,
2360 Event::ToolResult {
2361 run_id: "r1".parse().unwrap(),
2362 step: 2,
2363 call_id: "c1".parse().unwrap(),
2364 result: Value::Null,
2365 is_error: false,
2366 },
2367 ));
2368 assert_eq!(
2369 SessionProjection::replay(&events).unwrap_err(),
2370 EventError::LifecycleMismatch
2371 );
2372
2373 let mut dangling = open_step_events();
2374 dangling.push(event(
2375 7,
2376 Event::ToolCall {
2377 run_id: "r1".parse().unwrap(),
2378 step: 1,
2379 call_id: "c1".parse().unwrap(),
2380 tool: "echo".into(),
2381 arguments: Value::Null,
2382 },
2383 ));
2384 dangling.push(event(
2385 8,
2386 Event::StepFinished {
2387 run_id: "r1".parse().unwrap(),
2388 step: 1,
2389 },
2390 ));
2391 assert_eq!(
2392 SessionProjection::replay(&dangling).unwrap_err(),
2393 EventError::LifecycleMismatch
2394 );
2395 }
2396
2397 #[test]
2398 fn replay_rejects_reused_tool_call_ids() {
2399 let mut events = open_step_events();
2400 events.extend([
2401 event(
2402 7,
2403 Event::ToolCall {
2404 run_id: "r1".parse().unwrap(),
2405 step: 1,
2406 call_id: "c1".parse().unwrap(),
2407 tool: "echo".into(),
2408 arguments: Value::Null,
2409 },
2410 ),
2411 event(
2412 8,
2413 Event::ToolResult {
2414 run_id: "r1".parse().unwrap(),
2415 step: 1,
2416 call_id: "c1".parse().unwrap(),
2417 result: Value::Null,
2418 is_error: false,
2419 },
2420 ),
2421 event(
2422 9,
2423 Event::StepFinished {
2424 run_id: "r1".parse().unwrap(),
2425 step: 1,
2426 },
2427 ),
2428 event(
2429 10,
2430 Event::StepStarted {
2431 run_id: "r1".parse().unwrap(),
2432 step: 2,
2433 },
2434 ),
2435 event(
2436 11,
2437 Event::ToolCall {
2438 run_id: "r1".parse().unwrap(),
2439 step: 2,
2440 call_id: "c1".parse().unwrap(),
2441 tool: "echo".into(),
2442 arguments: Value::Null,
2443 },
2444 ),
2445 ]);
2446 assert_eq!(
2447 SessionProjection::replay(&events).unwrap_err(),
2448 EventError::DuplicateToolCall("c1".parse().unwrap())
2449 );
2450 }
2451
2452 #[test]
2453 fn replay_keeps_injected_context_without_changing_run_state() {
2454 let mut events = open_step_events();
2455 events.push(event(
2456 7,
2457 Event::ContextInjected {
2458 run_id: "r1".parse().unwrap(),
2459 step: 1,
2460 contribution_id: "memory:1".into(),
2461 source: "memory".into(),
2462 version: "v1".into(),
2463 authority: "tenant".into(),
2464 form: "message".into(),
2465 content: text("tenant context"),
2466 },
2467 ));
2468
2469 let projection = SessionProjection::replay(&events).unwrap();
2470 assert_eq!(projection.active_run_id.as_deref(), Some("r1"));
2471 assert_eq!(projection.open_steps, BTreeSet::from([1]));
2472 assert_eq!(
2473 projection.injected_context,
2474 vec![ProjectedContext {
2475 run_id: "r1".parse().unwrap(),
2476 step: 1,
2477 contribution_id: "memory:1".into(),
2478 source: "memory".into(),
2479 version: "v1".into(),
2480 authority: "tenant".into(),
2481 form: "message".into(),
2482 content: text("tenant context"),
2483 }]
2484 );
2485 }
2486
2487 #[test]
2488 fn replay_skips_unknown_ignorable_formats_and_rejects_required_events() {
2489 let optional = event(
2490 1,
2491 Event::Opaque {
2492 format_version: SESSION_EVENT_FORMAT_VERSION,
2493 event_type: "future_optional".into(),
2494 ignorable: true,
2495 payload: serde_json::json!({"answer": 42}),
2496 },
2497 );
2498 let projection = SessionProjection::replay(&[optional]).unwrap();
2499 assert_eq!(projection.last_seq, 1);
2500
2501 let required = event(
2502 1,
2503 Event::Opaque {
2504 format_version: SESSION_EVENT_FORMAT_VERSION,
2505 event_type: "future_required".into(),
2506 ignorable: false,
2507 payload: Value::Null,
2508 },
2509 );
2510 assert_eq!(
2511 SessionProjection::replay(&[required]).unwrap_err(),
2512 EventError::UnknownRequired("future_required".into())
2513 );
2514
2515 let unsupported = event(
2516 1,
2517 Event::Opaque {
2518 format_version: SESSION_EVENT_FORMAT_VERSION + 1,
2519 event_type: "future_optional".into(),
2520 ignorable: true,
2521 payload: Value::Null,
2522 },
2523 );
2524 assert_eq!(
2525 SessionProjection::replay(&[unsupported]).unwrap().last_seq,
2526 1
2527 );
2528
2529 let required_new_format = event(
2530 1,
2531 Event::Opaque {
2532 format_version: SESSION_EVENT_FORMAT_VERSION + 1,
2533 event_type: "future_required".into(),
2534 ignorable: false,
2535 payload: Value::Null,
2536 },
2537 );
2538 assert_eq!(
2539 SessionProjection::replay(&[required_new_format]).unwrap_err(),
2540 EventError::UnsupportedFormat(SESSION_EVENT_FORMAT_VERSION + 1)
2541 );
2542 }
2543
2544 #[test]
2545 fn legacy_tool_artifact_metadata_defaults_without_changing_old_events() {
2546 let artifact: ToolResultArtifact = serde_json::from_value(serde_json::json!({
2547 "result_ref": "ref",
2548 "sha256": "digest",
2549 "original_bytes": 100,
2550 "view_bytes": 20
2551 }))
2552 .unwrap();
2553 assert_eq!(artifact.expires_at, None);
2554 assert!(artifact.omitted.is_empty());
2555 }
2556}