Skip to main content

af_agent_session/
lib.rs

1//! Event-sourced Agent sessions. The append-only log is the conversation SSOT;
2//! messages, runs, interactions and UI state are deterministic projections.
3
4#![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
23/// Wire/storage format of [`Event`]; only structural changes bump it.
24pub const SESSION_EVENT_FORMAT_VERSION: u32 = 1;
25
26/// Versioned organizational metadata projected from the Session log.
27#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
28#[serde(deny_unknown_fields)]
29pub struct SessionMetadata {
30    /// User title; empty means the consumer derives an automatic title.
31    pub title: String,
32    /// Archived Sessions reject new input; already accepted work may converge.
33    pub archived: bool,
34    /// Conversation-wide action approval policy inherited by new Runs.
35    #[serde(default)]
36    pub approval_mode: ApprovalMode,
37    /// Monotonic metadata version independent of running event traffic.
38    pub version: u64,
39}
40
41/// Conversation-wide policy for tool and completion-action execution.
42#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
43#[serde(rename_all = "snake_case")]
44pub enum ApprovalMode {
45    /// Permit reads only; every write is denied.
46    ReadOnly,
47    /// Permit reads and require confirmation for every write.
48    #[default]
49    ConfirmChanges,
50    /// Permit reads and explicitly classified low-risk writes without confirmation.
51    LowRiskAuto,
52}
53
54impl ApprovalMode {
55    /// Stable snake_case wire/storage value.
56    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    /// Validate display metadata at every public write and replay boundary.
80    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/// One committed Session event with its position in the log.
89#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
90pub struct SessionEvent {
91    /// Session this record belongs to.
92    pub session_id: SessionId,
93    /// Strictly increasing sequence inside the Session.
94    pub seq: u64,
95    /// When the event happened.
96    pub occurred_at: DateTime<Utc>,
97    /// The event payload.
98    pub event: Event,
99}
100
101impl SessionEvent {
102    /// An event awaiting append: sequence 0 until the store assigns it.
103    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    /// Format version of the payload.
113    pub fn format_version(&self) -> u32 {
114        self.event.format_version()
115    }
116
117    /// Stable snake_case event type.
118    pub fn event_type(&self) -> &str {
119        self.event.event_type()
120    }
121
122    /// Whether a reader that does not know this event may skip it.
123    pub fn ignorable(&self) -> bool {
124        self.event.ignorable()
125    }
126}
127
128/// Every model-visible fact and lifecycle transition of a Session. The log is the source of truth; everything else is a projection.
129#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
130#[serde(tag = "type", rename_all = "snake_case")]
131pub enum Event {
132    /// First event of a Session, pinning its Profile revision.
133    SessionCreated {
134        /// Profile revision the Session pins for its lifetime.
135        profile_revision_id: ProfileRevisionId,
136    },
137    /// First event of a forked Session, citing the parent position it copied through.
138    SessionForked {
139        /// Session the fork was cut from.
140        parent_session_id: SessionId,
141        /// Last parent sequence included in the fork.
142        parent_seq: u64,
143    },
144    /// Atomically replace organizational metadata at the next metadata version.
145    SessionMetadataUpdated {
146        /// Full immutable metadata fact.
147        metadata: SessionMetadata,
148    },
149    /// Tombstone: active and queued Runs are closed, listing hides the Session, history stays readable.
150    SessionDeleted {
151        /// Why the Session was deleted.
152        reason: String,
153    },
154    /// Input entered the inbox for `run_id`.
155    InputQueued {
156        /// Queued input.
157        input_id: InputId,
158        /// Run that will consume the input (a new Run for `followup`).
159        run_id: RunId,
160        /// Delivery semantics.
161        mode: DeliveryMode,
162        /// Submitted content.
163        content: Vec<ContentBlock>,
164        /// Skill slug invoked explicitly with this input.
165        explicit_skill: Option<String>,
166    },
167    /// The Run took the input out of the inbox.
168    InputClaimed {
169        /// Claimed input.
170        input_id: InputId,
171        /// Claiming Run.
172        run_id: RunId,
173    },
174    /// The input was dropped before a Run consumed it.
175    InputCancelled {
176        /// Cancelled input.
177        input_id: InputId,
178        /// Run that would have consumed it.
179        run_id: RunId,
180        /// Machine-readable cancellation code.
181        error_code: String,
182    },
183    /// A Run began executing.
184    RunStarted {
185        /// The Run.
186        run_id: RunId,
187        /// Input that started it.
188        input_id: InputId,
189    },
190    /// The Run parked on an interaction.
191    RunWaiting {
192        /// The Run.
193        run_id: RunId,
194        /// Interaction it waits for.
195        interaction_id: InteractionId,
196    },
197    /// The Run resumed after its interaction resolved.
198    RunResumed {
199        /// The Run.
200        run_id: RunId,
201        /// Resolved interaction.
202        interaction_id: InteractionId,
203    },
204    /// The Run reached a terminal status.
205    RunFinished {
206        /// The Run.
207        run_id: RunId,
208        /// Terminal status.
209        status: RunStatus,
210        /// Machine-readable failure code.
211        error_code: Option<String>,
212    },
213    /// A Turn opened.
214    TurnStarted {
215        /// The Run.
216        run_id: RunId,
217        /// Turn number.
218        turn: u32,
219    },
220    /// A Turn closed.
221    TurnFinished {
222        /// The Run.
223        run_id: RunId,
224        /// Turn number.
225        turn: u32,
226    },
227    /// A Step opened.
228    StepStarted {
229        /// The Run.
230        run_id: RunId,
231        /// Step number.
232        step: u32,
233    },
234    /// A Step closed.
235    StepFinished {
236        /// The Run.
237        run_id: RunId,
238        /// Step number.
239        step: u32,
240    },
241    /// User content entered the transcript.
242    UserMessage {
243        /// The Run.
244        run_id: RunId,
245        /// Content blocks.
246        content: Vec<ContentBlock>,
247    },
248    /// Streamed assistant text fragment.
249    AssistantDelta {
250        /// The Run.
251        run_id: RunId,
252        /// Step producing it.
253        step: u32,
254        /// Provider attempt producing it.
255        attempt: u32,
256        /// Text fragment.
257        content: String,
258    },
259    /// Final assistant content for a step.
260    AssistantMessage {
261        /// The Run.
262        run_id: RunId,
263        /// Step.
264        step: u32,
265        /// Provider attempt that succeeded.
266        attempt: u32,
267        /// Content blocks.
268        content: Vec<ContentBlock>,
269    },
270    /// The model requested tool calls.
271    AssistantToolCalls {
272        /// The Run.
273        run_id: RunId,
274        /// Step.
275        step: u32,
276        /// Text accompanying the calls.
277        content: Option<String>,
278        /// Requested calls in model order.
279        calls: Vec<RecordedToolCall>,
280    },
281    /// One tool call was admitted with immutable arguments.
282    ToolCall {
283        /// The Run.
284        run_id: RunId,
285        /// Step.
286        step: u32,
287        /// Call identity.
288        call_id: ToolCallId,
289        /// Tool name.
290        tool: String,
291        /// Immutable arguments.
292        arguments: Value,
293    },
294    /// Authorization outcome for a tool call.
295    ToolAuthorization {
296        /// The Run.
297        run_id: RunId,
298        /// Step.
299        step: u32,
300        /// Call identity.
301        call_id: ToolCallId,
302        /// Outcome.
303        status: ToolAuthorizationStatus,
304        /// Why it was denied or parked.
305        reason: Option<String>,
306    },
307    /// Tool execution began.
308    ToolExecutionStarted {
309        /// Immutable tool attribution; absent only in legacy logs.
310        #[serde(default, skip_serializing_if = "Option::is_none")]
311        metering: Option<MeteringDetails>,
312        /// Declared resource weight frozen before dispatch.
313        #[serde(default)]
314        reserved_cost_units: u64,
315        /// The Run.
316        run_id: RunId,
317        /// Step.
318        step: u32,
319        /// Call identity.
320        call_id: ToolCallId,
321    },
322    /// Tool execution finished.
323    ToolResult {
324        /// The Run.
325        run_id: RunId,
326        /// Step.
327        step: u32,
328        /// Call identity.
329        call_id: ToolCallId,
330        /// Result value.
331        result: Value,
332        /// Whether the result is an error.
333        is_error: bool,
334    },
335    /// Privileged correction; does not reopen execution or its reservation settlement.
336    UsageCorrected {
337        /// Immutable reported measurement and idempotency key.
338        correction: MeteringCorrection,
339        /// Authenticated reporting subject.
340        actor_id: af_context::SubjectId,
341    },
342    /// Usage settled for one idempotent operation.
343    UsageRecorded {
344        /// Immutable attribution; absent only on legacy events.
345        #[serde(default, skip_serializing_if = "Option::is_none")]
346        metering: Option<MeteringDetails>,
347        /// The Run.
348        run_id: RunId,
349        /// Idempotent operation identity.
350        operation_id: String,
351        /// Prompt tokens.
352        prompt_tokens: u64,
353        /// Completion tokens.
354        completion_tokens: u64,
355        /// Provider-neutral cost units.
356        #[serde(default)]
357        cost_units: u64,
358    },
359    /// A model retry was scheduled.
360    RetryScheduled {
361        /// The Run.
362        run_id: RunId,
363        /// Failed attempt number.
364        attempt: u32,
365        /// Backoff before the next attempt.
366        delay_ms: u64,
367        /// Failure reason.
368        reason: String,
369    },
370    /// The frozen model request before it is sent; the durable attempt identity makes crashes reconcilable.
371    ModelRequestPrepared {
372        /// Attribution frozen before dispatch; absent on legacy events.
373        #[serde(default, skip_serializing_if = "Option::is_none")]
374        metering: Option<MeteringDetails>,
375        /// The Run.
376        run_id: RunId,
377        /// Step.
378        step: u32,
379        /// Attempt number.
380        attempt: u32,
381        /// Durable provider attempt identity.
382        #[serde(default)]
383        provider_attempt_id: String,
384        /// Usage operation identity.
385        #[serde(default)]
386        operation_id: String,
387        /// Prompt tokens reserved.
388        #[serde(default)]
389        reserved_prompt_tokens: u64,
390        /// Completion tokens reserved.
391        #[serde(default)]
392        reserved_completion_tokens: u64,
393        /// Redacted request as sent.
394        request: Value,
395        /// Prompt sections the request was composed from.
396        prompt_sections: Value,
397    },
398    /// Context a contributor injected into a step.
399    ContextInjected {
400        /// The Run.
401        run_id: RunId,
402        /// Step.
403        step: u32,
404        /// Contribution identity.
405        contribution_id: String,
406        /// Contributor source.
407        source: String,
408        /// Contributor version.
409        version: String,
410        /// `trusted` or `untrusted`.
411        authority: String,
412        /// Rendering form.
413        form: String,
414        /// Content blocks.
415        content: Vec<ContentBlock>,
416    },
417    /// A model attempt failed.
418    ModelAttemptFailed {
419        /// The Run.
420        run_id: RunId,
421        /// Step.
422        step: u32,
423        /// Attempt number.
424        attempt: u32,
425        /// Failure.
426        error: String,
427        /// Whether a retry may help.
428        retryable: bool,
429    },
430    /// Compaction began.
431    CompactionStarted {
432        /// The Run.
433        run_id: RunId,
434        /// Compaction identity.
435        compaction_id: String,
436        /// Last sequence the summary covers.
437        source_through_seq: u64,
438    },
439    /// Large tool results were truncated in the transcript.
440    ToolResultsPruned {
441        /// The Run.
442        run_id: RunId,
443        /// Truncated calls.
444        call_ids: Vec<ToolCallId>,
445    },
446    /// A summary replaced transcript through `through_seq`.
447    SummaryReplaced {
448        /// The Run.
449        run_id: RunId,
450        /// Last sequence replaced.
451        through_seq: u64,
452        /// Summary text.
453        summary: String,
454        /// Compactor name.
455        compactor: String,
456        /// Model used.
457        model: String,
458    },
459    /// Compaction finished.
460    CompactionFinished {
461        /// The Run.
462        run_id: RunId,
463        /// Compaction identity.
464        compaction_id: String,
465        /// `completed` or `failed`.
466        status: String,
467        /// Failure, when failed.
468        error: Option<String>,
469    },
470    /// The Run asked the user for an approval or answer.
471    InteractionRequested {
472        /// The Run.
473        run_id: RunId,
474        /// Interaction identity.
475        interaction_id: InteractionId,
476        /// Approval or question.
477        kind: InteractionKind,
478        /// What is being asked.
479        payload: Value,
480    },
481    /// The user resolved an interaction.
482    InteractionResolved {
483        /// The Run.
484        run_id: RunId,
485        /// Interaction identity.
486        interaction_id: InteractionId,
487        /// Outcome.
488        resolution: InteractionResolution,
489        /// Answer payload.
490        payload: Value,
491    },
492    /// A child Session was created for this Run.
493    ChildSessionLinked {
494        /// The Run.
495        run_id: RunId,
496        /// Child Session.
497        child_session_id: SessionId,
498        /// Subagent provider.
499        provider: String,
500    },
501    /// Plugin-defined event; readers that do not know `event_type` treat it as data.
502    Extension {
503        /// The Run.
504        run_id: RunId,
505        /// Plugin that emitted it.
506        plugin_id: String,
507        /// Plugin event type.
508        event_type: String,
509        /// Plugin payload.
510        payload: Value,
511    },
512    /// An event this build cannot decode; `ignorable` decides whether replay may skip it.
513    #[serde(skip)]
514    Opaque {
515        /// Producer format version.
516        format_version: u32,
517        /// Producer event type.
518        event_type: String,
519        /// Whether readers may skip it.
520        ignorable: bool,
521        /// Raw payload.
522        payload: Value,
523    },
524}
525
526impl Event {
527    /// Format version of this event.
528    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    /// Stable snake_case event type.
536    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    /// Whether an unknown reader may skip this event.
580    pub fn ignorable(&self) -> bool {
581        match self {
582            Self::Opaque { ignorable, .. } => *ignorable,
583            _ => false,
584        }
585    }
586}
587
588/// How queued input reaches a Run.
589#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
590#[serde(rename_all = "snake_case")]
591pub enum DeliveryMode {
592    /// Runs as its own Turn after the current Run.
593    Followup,
594    /// Interrupts the current Run at the next step boundary.
595    Steer,
596    /// Added as context to the next step without waking the loop.
597    Inject,
598}
599
600/// Terminal Run status.
601#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
602#[serde(rename_all = "snake_case")]
603pub enum RunStatus {
604    /// Finished with an answer.
605    Completed,
606    /// Failed with an error code.
607    Failed,
608    /// Cancelled by the user or platform.
609    Cancelled,
610    /// Stopped at a step, tool-call or token limit.
611    MaxStepsReached,
612}
613
614impl RunStatus {
615    /// Stable snake_case name.
616    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/// What the Run asks the user for.
627#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
628#[serde(rename_all = "snake_case")]
629pub enum InteractionKind {
630    /// Approve or reject a side effect.
631    Action,
632    /// Answer a question.
633    UserQuestion,
634}
635
636impl InteractionKind {
637    /// Stable snake_case name.
638    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/// How the user resolved an interaction.
647#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
648#[serde(rename_all = "snake_case")]
649pub enum InteractionResolution {
650    /// Approved.
651    Confirmed,
652    /// Rejected.
653    Rejected,
654    /// Answered a question.
655    Answered,
656}
657
658impl InteractionResolution {
659    /// Stable snake_case name.
660    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/// Authorization outcome of a tool call.
670#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
671#[serde(rename_all = "snake_case")]
672pub enum ToolAuthorizationStatus {
673    /// May execute.
674    Allowed,
675    /// Parked until the interaction resolves.
676    Waiting,
677    /// Refused.
678    Denied,
679}
680
681/// A tool call as recorded in `AssistantToolCalls`.
682#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
683pub struct RecordedToolCall {
684    /// Tool call this record refers to.
685    pub call_id: ToolCallId,
686    /// Tool name.
687    pub tool: String,
688    /// JSON arguments passed to the tool.
689    pub arguments: Value,
690}
691
692impl ToolAuthorizationStatus {
693    /// Stable snake_case name.
694    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/// Typed content unit; UI and models parse types, never prose conventions.
704#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
705#[serde(tag = "type", rename_all = "snake_case")]
706pub enum ContentBlock {
707    /// Plain text.
708    Text {
709        /// The text.
710        text: String,
711    },
712    /// Attached resource resolved by the host.
713    Resource {
714        /// Resource identity.
715        resource_id: String,
716        /// MIME type.
717        media_type: String,
718    },
719    /// Product-typed data for a registered UI slot.
720    Data {
721        /// Slot name.
722        slot: String,
723        /// Slot payload.
724        value: Value,
725    },
726    /// Citation of a retrieved resource.
727    Citation {
728        /// Resource identity.
729        resource_id: String,
730        /// Display label.
731        label: String,
732        /// Where to open it.
733        uri: String,
734        /// Supporting excerpt.
735        excerpt: Option<String>,
736    },
737}
738
739/// Provider usage preserved through product settlement.
740#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
741pub struct UsageCharge {
742    /// Input tokens consumed by model calls.
743    pub prompt_tokens: u64,
744    /// Output tokens consumed by model calls.
745    pub completion_tokens: u64,
746    /// Provider-neutral tool units.
747    pub cost_units: u64,
748}
749
750impl UsageCharge {
751    /// Legacy scalar representation used by meters that do not price components separately.
752    pub fn total(self) -> u64 {
753        self.prompt_tokens
754            .saturating_add(self.completion_tokens)
755            .saturating_add(self.cost_units)
756    }
757}
758
759/// Deterministic projection of a Session log used by the runtime, recovery and UI.
760#[derive(Debug, Clone, Default, PartialEq)]
761pub struct SessionProjection {
762    /// Latest versioned Session title and archival state.
763    pub metadata: SessionMetadata,
764    /// Session this record belongs to; `None` until the first event is applied.
765    pub session_id: Option<SessionId>,
766    /// Immutable Agent Profile revision pinned by the Session; `None` before
767    /// `SessionCreated`.
768    pub profile_revision_id: Option<ProfileRevisionId>,
769    /// Whether the Session was deleted.
770    pub deleted: bool,
771    /// Sequence of the last applied event.
772    pub last_seq: u64,
773    /// Run currently executing, if any.
774    pub active_run_id: Option<RunId>,
775    /// Interaction the active Run waits for, if any.
776    pub waiting_interaction_id: Option<InteractionId>,
777    /// Conversation messages in request order.
778    pub messages: Vec<ProjectedMessage>,
779    /// Context injected so far, in order.
780    pub injected_context: Vec<ProjectedContext>,
781    /// State of every Run seen.
782    pub run_status: BTreeMap<RunId, RunState>,
783    /// Tool calls admitted but not yet resulted.
784    pub open_tool_calls: BTreeMap<ToolCallId, OpenToolCall>,
785    /// Open tool calls whose execution started.
786    pub started_tool_calls: BTreeSet<ToolCallId>,
787    /// Inputs queued but not claimed.
788    pub queued_inputs: BTreeMap<InputId, (RunId, DeliveryMode)>,
789    /// Inputs claimed by a Run.
790    pub claimed_inputs: BTreeMap<InputId, RunId>,
791    /// Open Turn, if any.
792    pub open_turn: Option<u32>,
793    /// Open steps (at most one).
794    pub open_steps: BTreeSet<u32>,
795    /// Next step number.
796    pub next_step: u32,
797    /// Open compaction `(id, through_seq)`, if any.
798    pub open_compaction: Option<(String, u64)>,
799    /// Short human-readable summary.
800    pub summary: Option<String>,
801    corrected_usage: BTreeMap<(RunId, String), OperationUsage>,
802    metering_corrections: BTreeMap<af_context::MeteringCorrectionId, AppliedMeteringCorrection>,
803    usage_operations: BTreeMap<(RunId, String), RecordedUsage>,
804    pending_usage_operations: BTreeMap<(RunId, String), PreparedUsage>,
805    seen_tool_calls: BTreeSet<ToolCallId>,
806}
807
808type PreparedUsage = RecordedUsage;
809
810type RecordedUsage = (u64, u64, u64, Option<MeteringDetails>);
811
812/// Projected lifecycle state of one Run: live, parked on an interaction, or terminal.
813#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
814#[serde(rename_all = "snake_case")]
815pub enum RunState {
816    /// Executing.
817    Running,
818    /// Parked on an interaction.
819    WaitingForInput,
820    /// Finished.
821    Terminal(RunStatus),
822}
823
824impl RunState {
825    /// Stable snake_case name.
826    pub const fn as_str(self) -> &'static str {
827        match self {
828            Self::Running => "running",
829            Self::WaitingForInput => "waiting_for_input",
830            Self::Terminal(status) => status.as_str(),
831        }
832    }
833
834    /// Whether the Run has finished.
835    pub const fn is_terminal(self) -> bool {
836        matches!(self, Self::Terminal(_))
837    }
838}
839
840/// Injected context as projected for UI and replay.
841#[derive(Debug, Clone, PartialEq)]
842pub struct ProjectedContext {
843    /// Run this record belongs to.
844    pub run_id: RunId,
845    /// 1-based step number inside the Turn.
846    pub step: u32,
847    /// Contribution identity.
848    pub contribution_id: String,
849    /// Contributor source.
850    pub source: String,
851    /// Semantic version string.
852    pub version: String,
853    /// `trusted` or `untrusted`.
854    pub authority: String,
855    /// Rendering form.
856    pub form: String,
857    /// Content blocks carried by this record.
858    pub content: Vec<ContentBlock>,
859}
860
861/// Message as projected for UI and replay.
862#[derive(Debug, Clone, PartialEq)]
863pub struct ProjectedMessage {
864    /// `user` or `assistant`.
865    pub role: &'static str,
866    /// Run this record belongs to.
867    pub run_id: RunId,
868    /// Content blocks carried by this record.
869    pub content: Vec<ContentBlock>,
870}
871
872/// A tool call awaiting its result.
873#[derive(Debug, Clone, PartialEq)]
874pub struct OpenToolCall {
875    /// Run this record belongs to.
876    pub run_id: RunId,
877    /// 1-based step number inside the Turn.
878    pub step: u32,
879    /// Tool name.
880    pub tool: String,
881    /// JSON arguments passed to the tool.
882    pub arguments: Value,
883    /// Sequence of the Session event this was derived from.
884    pub source_event_seq: u64,
885}
886
887impl SessionProjection {
888    /// Rebuild the projection from a committed log; fails on any invariant violation.
889    pub fn replay(events: &[SessionEvent]) -> Result<Self, EventError> {
890        let mut projection = Self::default();
891        for event in events {
892            projection.apply(event)?;
893        }
894        Ok(projection)
895    }
896
897    /// Apply one committed event, enforcing sequence, lifecycle and pairing invariants.
898    pub fn apply(&mut self, envelope: &SessionEvent) -> Result<(), EventError> {
899        self.apply_internal(envelope, true)
900    }
901
902    /// Reduce execution/authorization/usage facts without retaining UI messages
903    /// or injected prompt content. Suitable for incremental replay of bounded pages.
904    pub fn apply_facts(&mut self, envelope: &SessionEvent) -> Result<(), EventError> {
905        self.messages.clear();
906        self.injected_context.clear();
907        self.apply_internal(envelope, false)
908    }
909
910    /// Fold the same lifecycle invariants without building a transcript copy.
911    pub fn replay_facts(events: &[SessionEvent]) -> Result<Self, EventError> {
912        let mut projection = Self::default();
913        for event in events {
914            projection.apply_facts(event)?;
915        }
916        Ok(projection)
917    }
918
919    fn apply_internal(
920        &mut self,
921        envelope: &SessionEvent,
922        include_content: bool,
923    ) -> Result<(), EventError> {
924        if envelope.seq != self.last_seq + 1 {
925            return Err(EventError::Sequence {
926                expected: self.last_seq + 1,
927                actual: envelope.seq,
928            });
929        }
930        let session_id = self
931            .session_id
932            .get_or_insert_with(|| envelope.session_id.clone());
933        if *session_id != envelope.session_id {
934            return Err(EventError::SessionMismatch);
935        }
936        if envelope.format_version() != SESSION_EVENT_FORMAT_VERSION {
937            if envelope.ignorable() {
938                self.last_seq = envelope.seq;
939                return Ok(());
940            }
941            return Err(EventError::UnsupportedFormat(envelope.format_version()));
942        }
943        if self.deleted && !matches!(envelope.event, Event::UsageCorrected { .. }) {
944            return Err(EventError::SessionClosed);
945        }
946        match &envelope.event {
947            Event::UsageCorrected {
948                correction,
949                actor_id,
950            } => self.apply_metering_correction(correction, actor_id, envelope.seq)?,
951            Event::SessionCreated {
952                profile_revision_id,
953            } => {
954                if envelope.seq != 1 || self.profile_revision_id.is_some() {
955                    return Err(EventError::DuplicateSession);
956                }
957                self.profile_revision_id = Some(profile_revision_id.clone());
958            }
959            Event::SessionMetadataUpdated { metadata } => {
960                metadata.validate()?;
961                if self.profile_revision_id.is_none()
962                    || self.metadata.version.checked_add(1) != Some(metadata.version)
963                {
964                    return Err(EventError::InvalidSessionMetadata);
965                }
966                self.metadata = metadata.clone();
967            }
968            Event::SessionDeleted { .. } => {
969                if let Some(run_id) = self.active_run_id.take() {
970                    self.run_status
971                        .insert(run_id, RunState::Terminal(RunStatus::Cancelled));
972                }
973                for (_, (run_id, _)) in std::mem::take(&mut self.queued_inputs) {
974                    self.run_status
975                        .insert(run_id, RunState::Terminal(RunStatus::Cancelled));
976                }
977                self.waiting_interaction_id = None;
978                self.open_turn = None;
979                self.open_steps.clear();
980                self.open_tool_calls.clear();
981                self.started_tool_calls.clear();
982                self.open_compaction = None;
983                self.deleted = true;
984            }
985            Event::RunStarted { run_id, input_id } => {
986                if self.active_run_id.is_some() {
987                    return Err(EventError::ConcurrentRun);
988                }
989                if self.claimed_inputs.get(input_id) != Some(run_id) {
990                    return Err(EventError::UnclaimedInput(input_id.clone()));
991                }
992                self.active_run_id = Some(run_id.clone());
993                self.next_step = 1;
994                self.run_status.insert(run_id.clone(), RunState::Running);
995            }
996            Event::RunWaiting {
997                run_id,
998                interaction_id,
999            } => {
1000                self.require_active(run_id)?;
1001                self.waiting_interaction_id = Some(interaction_id.clone());
1002                self.run_status
1003                    .insert(run_id.clone(), RunState::WaitingForInput);
1004            }
1005            Event::RunResumed {
1006                run_id,
1007                interaction_id,
1008            } => {
1009                self.require_active(run_id)?;
1010                if self.waiting_interaction_id.as_ref() != Some(interaction_id) {
1011                    return Err(EventError::InteractionMismatch);
1012                }
1013                self.waiting_interaction_id = None;
1014                self.run_status.insert(run_id.clone(), RunState::Running);
1015            }
1016            Event::RunFinished { run_id, status, .. } => {
1017                self.require_active(run_id)?;
1018                if !self.open_tool_calls.is_empty()
1019                    || !self.open_steps.is_empty()
1020                    || self.open_turn.is_some()
1021                    || self.open_compaction.is_some()
1022                    || self
1023                        .queued_inputs
1024                        .values()
1025                        .any(|(target_run_id, _)| target_run_id == run_id)
1026                {
1027                    return Err(EventError::OpenLifecycle);
1028                }
1029                self.run_status
1030                    .insert(run_id.clone(), RunState::Terminal(*status));
1031                self.active_run_id = None;
1032                self.waiting_interaction_id = None;
1033            }
1034            Event::UserMessage { run_id, content } if include_content => {
1035                self.messages.push(ProjectedMessage {
1036                    role: "user",
1037                    run_id: run_id.clone(),
1038                    content: content.clone(),
1039                })
1040            }
1041            Event::AssistantMessage {
1042                run_id, content, ..
1043            } if include_content => self.messages.push(ProjectedMessage {
1044                role: "assistant",
1045                run_id: run_id.clone(),
1046                content: content.clone(),
1047            }),
1048            Event::ContextInjected {
1049                run_id,
1050                step,
1051                contribution_id,
1052                source,
1053                version,
1054                authority,
1055                form,
1056                content,
1057            } => {
1058                self.require_active(run_id)?;
1059                if !self.open_steps.contains(step) {
1060                    return Err(EventError::LifecycleMismatch);
1061                }
1062                if include_content {
1063                    self.injected_context.push(ProjectedContext {
1064                        run_id: run_id.clone(),
1065                        step: *step,
1066                        contribution_id: contribution_id.clone(),
1067                        source: source.clone(),
1068                        version: version.clone(),
1069                        authority: authority.clone(),
1070                        form: form.clone(),
1071                        content: content.clone(),
1072                    });
1073                }
1074            }
1075            Event::InputQueued {
1076                input_id,
1077                run_id,
1078                mode,
1079                ..
1080            } => {
1081                if self.metadata.archived {
1082                    return Err(EventError::SessionArchived);
1083                }
1084                if *mode != DeliveryMode::Followup {
1085                    self.require_active(run_id)?;
1086                }
1087                if self
1088                    .queued_inputs
1089                    .insert(input_id.clone(), (run_id.clone(), *mode))
1090                    .is_some()
1091                {
1092                    return Err(EventError::DuplicateInput(input_id.clone()));
1093                }
1094            }
1095            Event::InputClaimed { input_id, run_id } => {
1096                if self
1097                    .queued_inputs
1098                    .remove(input_id)
1099                    .map(|value| value.0)
1100                    .as_ref()
1101                    != Some(run_id)
1102                    || self
1103                        .claimed_inputs
1104                        .insert(input_id.clone(), run_id.clone())
1105                        .is_some()
1106                {
1107                    return Err(EventError::UnqueuedInput(input_id.clone()));
1108                }
1109            }
1110            Event::InputCancelled {
1111                input_id,
1112                run_id,
1113                error_code,
1114            } => {
1115                if self
1116                    .queued_inputs
1117                    .remove(input_id)
1118                    .map(|value| value.0)
1119                    .as_ref()
1120                    != Some(run_id)
1121                {
1122                    return Err(EventError::UnqueuedInput(input_id.clone()));
1123                }
1124                self.run_status.insert(
1125                    run_id.clone(),
1126                    RunState::Terminal(if error_code == "cancelled" {
1127                        RunStatus::Cancelled
1128                    } else {
1129                        RunStatus::Failed
1130                    }),
1131                );
1132            }
1133            Event::TurnStarted { run_id, turn } => {
1134                self.require_active(run_id)?;
1135                if self.open_turn.replace(*turn).is_some() {
1136                    return Err(EventError::ConcurrentTurn);
1137                }
1138            }
1139            Event::TurnFinished { run_id, turn } => {
1140                self.require_active(run_id)?;
1141                if self.open_turn != Some(*turn)
1142                    || !self.open_steps.is_empty()
1143                    || !self.open_tool_calls.is_empty()
1144                {
1145                    return Err(EventError::LifecycleMismatch);
1146                }
1147                self.open_turn = None;
1148            }
1149            Event::StepStarted { run_id, step } => {
1150                self.require_active(run_id)?;
1151                if self.open_turn.is_none()
1152                    || !self.open_steps.is_empty()
1153                    || *step != self.next_step
1154                    || !self.open_steps.insert(*step)
1155                {
1156                    return Err(EventError::ConcurrentStep);
1157                }
1158            }
1159            Event::StepFinished { run_id, step } => {
1160                self.require_active(run_id)?;
1161                if self
1162                    .open_tool_calls
1163                    .values()
1164                    .any(|call| call.run_id == *run_id && call.step == *step)
1165                    || !self.open_steps.remove(step)
1166                {
1167                    return Err(EventError::LifecycleMismatch);
1168                }
1169                self.next_step = step.saturating_add(1);
1170            }
1171            Event::ToolCall {
1172                run_id,
1173                step,
1174                call_id,
1175                tool,
1176                arguments,
1177            } => {
1178                self.require_active(run_id)?;
1179                if !self.open_steps.contains(step) {
1180                    return Err(EventError::LifecycleMismatch);
1181                }
1182                if !self.seen_tool_calls.insert(call_id.clone())
1183                    || self
1184                        .open_tool_calls
1185                        .insert(
1186                            call_id.clone(),
1187                            OpenToolCall {
1188                                run_id: run_id.clone(),
1189                                step: *step,
1190                                tool: tool.clone(),
1191                                arguments: arguments.clone(),
1192                                source_event_seq: envelope.seq,
1193                            },
1194                        )
1195                        .is_some()
1196                {
1197                    return Err(EventError::DuplicateToolCall(call_id.clone()));
1198                }
1199            }
1200            Event::ToolResult {
1201                run_id,
1202                step,
1203                call_id,
1204                ..
1205            } => {
1206                self.require_active(run_id)?;
1207                if !self.open_steps.contains(step) {
1208                    return Err(EventError::LifecycleMismatch);
1209                }
1210                let Some(call) = self.open_tool_calls.get(call_id) else {
1211                    return Err(EventError::OrphanToolResult(call_id.clone()));
1212                };
1213                if call.run_id != *run_id || call.step != *step {
1214                    return Err(EventError::ToolResultMismatch(call_id.clone()));
1215                }
1216                self.open_tool_calls.remove(call_id);
1217                self.started_tool_calls.remove(call_id);
1218            }
1219            Event::ToolExecutionStarted {
1220                run_id,
1221                step,
1222                call_id,
1223                metering,
1224                reserved_cost_units,
1225            } => {
1226                self.require_active(run_id)?;
1227                let Some(call) = self.open_tool_calls.get(call_id) else {
1228                    return Err(EventError::OrphanToolResult(call_id.clone()));
1229                };
1230                if call.run_id != *run_id
1231                    || call.step != *step
1232                    || self.started_tool_calls.contains(call_id)
1233                {
1234                    return Err(EventError::ToolResultMismatch(call_id.clone()));
1235                }
1236                if let Some(details) = metering {
1237                    details.validate()?;
1238                    if !matches!(details, MeteringDetails::Tool { call_id: id, name, source: MeteringSource::Estimated, outcome: MeteringOutcome::Unknown } if id == call_id && name == &call.tool)
1239                    {
1240                        return Err(EventError::InvalidMetering);
1241                    }
1242                }
1243                self.started_tool_calls.insert(call_id.clone());
1244                self.pending_usage_operations.insert(
1245                    (run_id.clone(), format!("tool:{call_id}")),
1246                    (0, 0, *reserved_cost_units, metering.clone()),
1247                );
1248            }
1249            Event::ModelRequestPrepared {
1250                run_id,
1251                operation_id,
1252                reserved_prompt_tokens,
1253                reserved_completion_tokens,
1254                metering,
1255                ..
1256            } if !operation_id.is_empty() => {
1257                self.require_active(run_id)?;
1258                let key = (run_id.clone(), operation_id.clone());
1259                if !self.usage_operations.contains_key(&key) {
1260                    if let Some(details) = metering {
1261                        details.validate()?;
1262                        if !matches!(
1263                            details,
1264                            MeteringDetails::Model {
1265                                source: MeteringSource::Estimated,
1266                                outcome: MeteringOutcome::Unknown,
1267                                ..
1268                            }
1269                        ) {
1270                            return Err(EventError::InvalidMetering);
1271                        }
1272                    }
1273                    let reservation = (
1274                        *reserved_prompt_tokens,
1275                        *reserved_completion_tokens,
1276                        0,
1277                        metering.clone(),
1278                    );
1279                    match self.pending_usage_operations.get(&key) {
1280                        Some(existing) if *existing != reservation => {
1281                            return Err(EventError::UsageConflict(operation_id.clone()));
1282                        }
1283                        Some(_) => {}
1284                        None => {
1285                            self.pending_usage_operations.insert(key, reservation);
1286                        }
1287                    }
1288                }
1289            }
1290            Event::UsageRecorded {
1291                run_id,
1292                operation_id,
1293                prompt_tokens,
1294                completion_tokens,
1295                cost_units,
1296                metering,
1297            } => {
1298                self.require_active(run_id)?;
1299                if let Some(details) = metering {
1300                    details.validate()?;
1301                }
1302                let key = (run_id.clone(), operation_id.clone());
1303                if let Some((_, _, _, Some(prepared))) = self.pending_usage_operations.get(&key) {
1304                    if !metering
1305                        .as_ref()
1306                        .is_some_and(|details| details.completes(prepared))
1307                    {
1308                        return Err(EventError::UsageConflict(operation_id.clone()));
1309                    }
1310                }
1311                let usage = (
1312                    *prompt_tokens,
1313                    *completion_tokens,
1314                    *cost_units,
1315                    metering.clone(),
1316                );
1317                match self.usage_operations.get(&key) {
1318                    Some(existing) if *existing != usage => {
1319                        return Err(EventError::UsageConflict(operation_id.clone()));
1320                    }
1321                    Some(_) => {}
1322                    None => {
1323                        self.pending_usage_operations.remove(&key);
1324                        self.usage_operations.insert(key, usage);
1325                    }
1326                }
1327            }
1328            Event::CompactionStarted {
1329                run_id,
1330                compaction_id,
1331                source_through_seq,
1332            } => {
1333                self.require_active(run_id)?;
1334                if self.open_compaction.is_some() {
1335                    return Err(EventError::ConcurrentCompaction);
1336                }
1337                self.open_compaction = Some((compaction_id.clone(), *source_through_seq));
1338            }
1339            Event::CompactionFinished {
1340                run_id,
1341                compaction_id,
1342                ..
1343            } => {
1344                self.require_active(run_id)?;
1345                if self.open_compaction.as_ref().map(|value| value.0.as_str())
1346                    != Some(compaction_id.as_str())
1347                {
1348                    return Err(EventError::CompactionMismatch);
1349                }
1350                self.open_compaction = None;
1351            }
1352            Event::SummaryReplaced { summary, .. } => self.summary = Some(summary.clone()),
1353            Event::Opaque {
1354                event_type,
1355                ignorable: false,
1356                ..
1357            } => return Err(EventError::UnknownRequired(event_type.clone())),
1358            _ => {}
1359        }
1360        self.last_seq = envelope.seq;
1361        Ok(())
1362    }
1363
1364    fn require_active(&self, run_id: &RunId) -> Result<(), EventError> {
1365        if self.active_run_id.as_ref() == Some(run_id) {
1366            Ok(())
1367        } else {
1368            Err(EventError::RunMismatch)
1369        }
1370    }
1371
1372    /// `(prompt, completion)` tokens settled for `run_id`.
1373    pub fn usage_for(&self, run_id: &str) -> (u64, u64) {
1374        self.usage_operations
1375            .iter()
1376            .filter(|((recorded_run_id, _), _)| recorded_run_id.as_str() == run_id)
1377            .fold((0, 0), |total, (_, usage)| {
1378                (total.0 + usage.0, total.1 + usage.1)
1379            })
1380    }
1381
1382    /// Provider usage preserved for product settlement.
1383    pub fn usage_charge_for(&self, run_id: &str) -> UsageCharge {
1384        let sum = |pending: bool| {
1385            let operations = if pending {
1386                &self.pending_usage_operations
1387            } else {
1388                &self.usage_operations
1389            };
1390            operations
1391                .iter()
1392                .filter(|((recorded_run_id, _), _)| recorded_run_id.as_str() == run_id)
1393                .fold(UsageCharge::default(), |mut total, (_, usage)| {
1394                    total.prompt_tokens = total.prompt_tokens.saturating_add(usage.0);
1395                    total.completion_tokens = total.completion_tokens.saturating_add(usage.1);
1396                    total.cost_units = total.cost_units.saturating_add(usage.2);
1397                    total
1398                })
1399        };
1400        let recorded = sum(false);
1401        let pending = sum(true);
1402        UsageCharge {
1403            prompt_tokens: recorded.prompt_tokens.saturating_add(pending.prompt_tokens),
1404            completion_tokens: recorded
1405                .completion_tokens
1406                .saturating_add(pending.completion_tokens),
1407            cost_units: recorded.cost_units.saturating_add(pending.cost_units),
1408        }
1409    }
1410
1411    /// Settled plus reserved-but-unsettled units for `run_id`, for conservative billing.
1412    pub fn billable_units_for(&self, run_id: &str) -> u64 {
1413        self.usage_charge_for(run_id).total()
1414    }
1415}
1416
1417/// Invariant violated while applying an event.
1418#[derive(Debug, thiserror::Error, PartialEq, Eq)]
1419pub enum EventError {
1420    /// Metadata syntax or version does not match the prior state.
1421    #[error("invalid Session metadata or metadata version conflict")]
1422    InvalidSessionMetadata,
1423    /// Metering attribution is malformed.
1424    #[error("invalid metering attribution")]
1425    InvalidMetering,
1426    /// New input requires an unarchived Session.
1427    #[error("Session is archived")]
1428    SessionArchived,
1429    /// Unsupported session event format version.
1430    #[error("unsupported session event format version {0}")]
1431    UnsupportedFormat(u32),
1432    /// Unknown required session event type.
1433    #[error("unknown required session event type {0}")]
1434    UnknownRequired(String),
1435    /// Event sequence mismatch: expected `expected`, got `actual`.
1436    #[error("event sequence mismatch: expected {expected}, got {actual}")]
1437    Sequence {
1438        /// Sequence the projection expected next.
1439        expected: u64,
1440        /// Sequence the event carried.
1441        actual: u64,
1442    },
1443    /// Event belongs to another session.
1444    #[error("event belongs to another session")]
1445    SessionMismatch,
1446    /// Session creation must be the first and only creation event.
1447    #[error("session creation must be the first and only creation event")]
1448    DuplicateSession,
1449    /// Session already has an active run.
1450    #[error("session already has an active run")]
1451    ConcurrentRun,
1452    /// Session is closed.
1453    #[error("session is closed")]
1454    SessionClosed,
1455    /// Event does not match the active run.
1456    #[error("event does not match the active run")]
1457    RunMismatch,
1458    /// Interaction does not match the waiting run.
1459    #[error("interaction does not match the waiting run")]
1460    InteractionMismatch,
1461    /// Run cannot finish with an open turn, step or tool call.
1462    #[error("run cannot finish with an open turn, step or tool call")]
1463    OpenLifecycle,
1464    /// Input was queued twice.
1465    #[error("input was queued twice: {0}")]
1466    DuplicateInput(InputId),
1467    /// Input was claimed before it was queued.
1468    #[error("input was claimed before it was queued: {0}")]
1469    UnqueuedInput(InputId),
1470    /// Run started from an unclaimed input.
1471    #[error("run started from an unclaimed input: {0}")]
1472    UnclaimedInput(InputId),
1473    /// Session already has an active turn.
1474    #[error("session already has an active turn")]
1475    ConcurrentTurn,
1476    /// Turn already has this active step.
1477    #[error("turn already has this active step")]
1478    ConcurrentStep,
1479    /// Turn or step lifecycle does not pair.
1480    #[error("turn or step lifecycle does not pair")]
1481    LifecycleMismatch,
1482    /// Duplicate tool call.
1483    #[error("duplicate tool call {0}")]
1484    DuplicateToolCall(ToolCallId),
1485    /// Tool result has no matching call.
1486    #[error("tool result has no matching call {0}")]
1487    OrphanToolResult(ToolCallId),
1488    /// Tool result does not match the call run and step.
1489    #[error("tool result does not match the call run and step: {0}")]
1490    ToolResultMismatch(ToolCallId),
1491    /// Usage operation was recorded with different totals.
1492    #[error("usage operation was recorded with different totals: {0}")]
1493    UsageConflict(String),
1494    /// Session already has an active compaction.
1495    #[error("session already has an active compaction")]
1496    ConcurrentCompaction,
1497    /// Compaction lifecycle does not pair.
1498    #[error("compaction lifecycle does not pair")]
1499    CompactionMismatch,
1500    /// Event store conflict.
1501    #[error("event store conflict: {0}")]
1502    Conflict(String),
1503    /// Event store unavailable.
1504    #[error("event store unavailable: {0}")]
1505    Unavailable(String),
1506}
1507
1508/// Append-only Session event storage.
1509#[async_trait]
1510pub trait SessionEventStore: Send + Sync {
1511    /// Append events atomically after `expected_seq`; returns them with sequences.
1512    async fn append(
1513        &self,
1514        tenant_id: &str,
1515        session_id: &str,
1516        expected_seq: u64,
1517        events: Vec<Event>,
1518    ) -> Result<Vec<SessionEvent>, EventError>;
1519    /// Events after `after_seq`.
1520    async fn load(
1521        &self,
1522        tenant_id: &str,
1523        session_id: &str,
1524        after_seq: u64,
1525    ) -> Result<Vec<SessionEvent>, EventError>;
1526}
1527
1528/// One text content block.
1529pub fn text(value: impl Into<String>) -> Vec<ContentBlock> {
1530    vec![ContentBlock::Text { text: value.into() }]
1531}
1532
1533/// Events that close an interrupted Run as failed with `worker_restarted`, so recovery starts from a clean lifecycle.
1534pub fn recovery_events(projection: &SessionProjection) -> Vec<Event> {
1535    failure_events(projection, "worker_restarted")
1536}
1537
1538/// Events that close an interrupted Run as failed with `error_code`.
1539pub fn failure_events(projection: &SessionProjection, error_code: &str) -> Vec<Event> {
1540    termination_events(projection, RunStatus::Failed, error_code)
1541}
1542
1543/// Events that close an interrupted Run as cancelled.
1544pub fn cancel_events(projection: &SessionProjection) -> Vec<Event> {
1545    termination_events(projection, RunStatus::Cancelled, "cancelled")
1546}
1547
1548/// Events that tombstone a Session and close its Runs.
1549pub fn session_deletion_events(projection: &SessionProjection, reason: &str) -> Vec<Event> {
1550    let active = projection.active_run_id.as_deref();
1551    let mut events = cancel_events(projection);
1552    events.extend(
1553        projection
1554            .queued_inputs
1555            .iter()
1556            .filter(|(_, (run_id, _))| Some(run_id.as_str()) != active)
1557            .map(|(input_id, (run_id, _))| Event::InputCancelled {
1558                input_id: input_id.clone(),
1559                run_id: run_id.clone(),
1560                error_code: "cancelled".into(),
1561            }),
1562    );
1563    events.push(Event::SessionDeleted {
1564        reason: reason.into(),
1565    });
1566    events
1567}
1568
1569fn termination_events(
1570    projection: &SessionProjection,
1571    status: RunStatus,
1572    error_code: &str,
1573) -> Vec<Event> {
1574    let Some(run_id) = &projection.active_run_id else {
1575        return Vec::new();
1576    };
1577    let mut events = projection
1578        .queued_inputs
1579        .iter()
1580        .filter(|(_, (target_run_id, _))| target_run_id == run_id)
1581        .map(|(input_id, _)| Event::InputCancelled {
1582            input_id: input_id.clone(),
1583            run_id: run_id.clone(),
1584            error_code: error_code.into(),
1585        })
1586        .collect::<Vec<_>>();
1587    events.extend(
1588        projection
1589            .open_tool_calls
1590            .iter()
1591            .map(|(call_id, call)| Event::ToolResult {
1592                run_id: run_id.clone(),
1593                step: call.step,
1594                call_id: call_id.clone(),
1595                result: termination_result(error_code),
1596                is_error: true,
1597            }),
1598    );
1599    if let Some((compaction_id, _)) = &projection.open_compaction {
1600        events.push(Event::CompactionFinished {
1601            run_id: run_id.clone(),
1602            compaction_id: compaction_id.clone(),
1603            status: "failed".into(),
1604            error: Some(error_code.into()),
1605        });
1606    }
1607    events.extend(
1608        projection
1609            .open_steps
1610            .iter()
1611            .map(|step| Event::StepFinished {
1612                run_id: run_id.clone(),
1613                step: *step,
1614            }),
1615    );
1616    if let Some(turn) = projection.open_turn {
1617        events.push(Event::TurnFinished {
1618            run_id: run_id.clone(),
1619            turn,
1620        });
1621    }
1622    events.push(Event::RunFinished {
1623        run_id: run_id.clone(),
1624        status,
1625        error_code: Some(error_code.into()),
1626    });
1627    events
1628}
1629
1630fn termination_result(error_code: &str) -> Value {
1631    if error_code == "tool_outcome_unknown" || error_code == "worker_restarted" {
1632        serde_json::json!({
1633            "error": error_code,
1634            "guidance": "The tool outcome is unknown. Verify external state before retrying any operation with side effects; ask the user when verification is unavailable."
1635        })
1636    } else {
1637        serde_json::json!({"error":error_code})
1638    }
1639}
1640
1641#[cfg(test)]
1642mod tests {
1643    use super::*;
1644
1645    fn event(seq: u64, event: Event) -> SessionEvent {
1646        SessionEvent {
1647            session_id: "s".parse().unwrap(),
1648            seq,
1649            occurred_at: Utc::now(),
1650            event,
1651        }
1652    }
1653
1654    #[test]
1655    fn replay_enforces_single_run_and_tool_pairs() {
1656        let events = vec![
1657            event(
1658                1,
1659                Event::SessionCreated {
1660                    profile_revision_id: "p1".parse().unwrap(),
1661                },
1662            ),
1663            event(
1664                2,
1665                Event::InputQueued {
1666                    input_id: "i1".parse().unwrap(),
1667                    run_id: "r1".parse().unwrap(),
1668                    mode: DeliveryMode::Followup,
1669                    content: text("hi"),
1670                    explicit_skill: None,
1671                },
1672            ),
1673            event(
1674                3,
1675                Event::InputClaimed {
1676                    input_id: "i1".parse().unwrap(),
1677                    run_id: "r1".parse().unwrap(),
1678                },
1679            ),
1680            event(
1681                4,
1682                Event::RunStarted {
1683                    run_id: "r1".parse().unwrap(),
1684                    input_id: "i1".parse().unwrap(),
1685                },
1686            ),
1687            event(
1688                5,
1689                Event::TurnStarted {
1690                    run_id: "r1".parse().unwrap(),
1691                    turn: 1,
1692                },
1693            ),
1694            event(
1695                6,
1696                Event::StepStarted {
1697                    run_id: "r1".parse().unwrap(),
1698                    step: 1,
1699                },
1700            ),
1701            event(
1702                7,
1703                Event::ToolCall {
1704                    run_id: "r1".parse().unwrap(),
1705                    step: 1,
1706                    call_id: "c1".parse().unwrap(),
1707                    tool: "echo".into(),
1708                    arguments: serde_json::json!({"x":1}),
1709                },
1710            ),
1711            event(
1712                8,
1713                Event::ToolResult {
1714                    run_id: "r1".parse().unwrap(),
1715                    step: 1,
1716                    call_id: "c1".parse().unwrap(),
1717                    result: serde_json::json!({"x":1}),
1718                    is_error: false,
1719                },
1720            ),
1721            event(
1722                9,
1723                Event::StepFinished {
1724                    run_id: "r1".parse().unwrap(),
1725                    step: 1,
1726                },
1727            ),
1728            event(
1729                10,
1730                Event::TurnFinished {
1731                    run_id: "r1".parse().unwrap(),
1732                    turn: 1,
1733                },
1734            ),
1735            event(
1736                11,
1737                Event::RunFinished {
1738                    run_id: "r1".parse().unwrap(),
1739                    status: RunStatus::Completed,
1740                    error_code: None,
1741                },
1742            ),
1743        ];
1744        let projection = SessionProjection::replay(&events).unwrap();
1745        assert_eq!(projection.last_seq, 11);
1746        assert!(projection.active_run_id.is_none());
1747    }
1748
1749    #[test]
1750    fn replay_rejects_orphan_tool_result() {
1751        let events = vec![
1752            event(
1753                1,
1754                Event::SessionCreated {
1755                    profile_revision_id: "p1".parse().unwrap(),
1756                },
1757            ),
1758            event(
1759                2,
1760                Event::InputQueued {
1761                    input_id: "i1".parse().unwrap(),
1762                    run_id: "r1".parse().unwrap(),
1763                    mode: DeliveryMode::Followup,
1764                    content: text("hi"),
1765                    explicit_skill: None,
1766                },
1767            ),
1768            event(
1769                3,
1770                Event::InputClaimed {
1771                    input_id: "i1".parse().unwrap(),
1772                    run_id: "r1".parse().unwrap(),
1773                },
1774            ),
1775            event(
1776                4,
1777                Event::RunStarted {
1778                    run_id: "r1".parse().unwrap(),
1779                    input_id: "i1".parse().unwrap(),
1780                },
1781            ),
1782            event(
1783                5,
1784                Event::TurnStarted {
1785                    run_id: "r1".parse().unwrap(),
1786                    turn: 1,
1787                },
1788            ),
1789            event(
1790                6,
1791                Event::StepStarted {
1792                    run_id: "r1".parse().unwrap(),
1793                    step: 1,
1794                },
1795            ),
1796            event(
1797                7,
1798                Event::ToolResult {
1799                    run_id: "r1".parse().unwrap(),
1800                    step: 1,
1801                    call_id: "missing".parse().unwrap(),
1802                    result: Value::Null,
1803                    is_error: true,
1804                },
1805            ),
1806        ];
1807        assert_eq!(
1808            SessionProjection::replay(&events).unwrap_err(),
1809            EventError::OrphanToolResult("missing".parse().unwrap())
1810        );
1811    }
1812
1813    #[test]
1814    fn queued_input_failure_is_not_projected_as_cancellation() {
1815        let events = vec![
1816            event(
1817                1,
1818                Event::SessionCreated {
1819                    profile_revision_id: "p1".parse().unwrap(),
1820                },
1821            ),
1822            event(
1823                2,
1824                Event::InputQueued {
1825                    input_id: "i1".parse().unwrap(),
1826                    run_id: "r1".parse().unwrap(),
1827                    mode: DeliveryMode::Followup,
1828                    content: text("hi"),
1829                    explicit_skill: None,
1830                },
1831            ),
1832            event(
1833                3,
1834                Event::InputCancelled {
1835                    input_id: "i1".parse().unwrap(),
1836                    run_id: "r1".parse().unwrap(),
1837                    error_code: "profile_not_found".into(),
1838                },
1839            ),
1840        ];
1841        let projection = SessionProjection::replay(&events).unwrap();
1842        assert_eq!(
1843            projection.run_status.get("r1").map(|state| state.as_str()),
1844            Some("failed")
1845        );
1846    }
1847
1848    #[test]
1849    fn session_deletion_closes_active_and_queued_runs_before_tombstone() {
1850        let mut events = vec![
1851            event(
1852                1,
1853                Event::SessionCreated {
1854                    profile_revision_id: "p1".parse().unwrap(),
1855                },
1856            ),
1857            event(
1858                2,
1859                Event::InputQueued {
1860                    input_id: "i1".parse().unwrap(),
1861                    run_id: "r1".parse().unwrap(),
1862                    mode: DeliveryMode::Followup,
1863                    content: text("start"),
1864                    explicit_skill: None,
1865                },
1866            ),
1867            event(
1868                3,
1869                Event::InputClaimed {
1870                    input_id: "i1".parse().unwrap(),
1871                    run_id: "r1".parse().unwrap(),
1872                },
1873            ),
1874            event(
1875                4,
1876                Event::RunStarted {
1877                    run_id: "r1".parse().unwrap(),
1878                    input_id: "i1".parse().unwrap(),
1879                },
1880            ),
1881            event(
1882                5,
1883                Event::TurnStarted {
1884                    run_id: "r1".parse().unwrap(),
1885                    turn: 1,
1886                },
1887            ),
1888            event(
1889                6,
1890                Event::InputQueued {
1891                    input_id: "i2".parse().unwrap(),
1892                    run_id: "r2".parse().unwrap(),
1893                    mode: DeliveryMode::Followup,
1894                    content: text("later"),
1895                    explicit_skill: None,
1896                },
1897            ),
1898        ];
1899        let projection = SessionProjection::replay(&events).unwrap();
1900        for event_value in session_deletion_events(&projection, "api_deleted") {
1901            let seq = events.len() as u64 + 1;
1902            events.push(event(seq, event_value));
1903        }
1904        let deleted = SessionProjection::replay(&events).unwrap();
1905        assert!(deleted.deleted);
1906        assert_eq!(
1907            deleted.run_status.get("r1").map(|state| state.as_str()),
1908            Some("cancelled")
1909        );
1910        assert_eq!(
1911            deleted.run_status.get("r2").map(|state| state.as_str()),
1912            Some("cancelled")
1913        );
1914        assert!(events.iter().any(|event| matches!(
1915            &event.event,
1916            Event::RunFinished { run_id, status: RunStatus::Cancelled, .. } if run_id == "r1"
1917        )));
1918        assert!(matches!(
1919            events.last().map(|event| &event.event),
1920            Some(Event::SessionDeleted { .. })
1921        ));
1922    }
1923
1924    #[test]
1925    fn usage_operations_are_idempotent_and_conflicts_fail_replay() {
1926        let mut events = vec![
1927            event(
1928                1,
1929                Event::SessionCreated {
1930                    profile_revision_id: "p1".parse().unwrap(),
1931                },
1932            ),
1933            event(
1934                2,
1935                Event::InputQueued {
1936                    input_id: "i1".parse().unwrap(),
1937                    run_id: "r1".parse().unwrap(),
1938                    mode: DeliveryMode::Followup,
1939                    content: text("hi"),
1940                    explicit_skill: None,
1941                },
1942            ),
1943            event(
1944                3,
1945                Event::InputClaimed {
1946                    input_id: "i1".parse().unwrap(),
1947                    run_id: "r1".parse().unwrap(),
1948                },
1949            ),
1950            event(
1951                4,
1952                Event::RunStarted {
1953                    run_id: "r1".parse().unwrap(),
1954                    input_id: "i1".parse().unwrap(),
1955                },
1956            ),
1957            event(
1958                5,
1959                Event::UsageRecorded {
1960                    metering: None,
1961                    run_id: "r1".parse().unwrap(),
1962                    operation_id: "model:1:attempt:1".into(),
1963                    prompt_tokens: 7,
1964                    completion_tokens: 3,
1965                    cost_units: 5,
1966                },
1967            ),
1968            event(
1969                6,
1970                Event::UsageRecorded {
1971                    metering: None,
1972                    run_id: "r1".parse().unwrap(),
1973                    operation_id: "model:1:attempt:1".into(),
1974                    prompt_tokens: 7,
1975                    completion_tokens: 3,
1976                    cost_units: 5,
1977                },
1978            ),
1979        ];
1980        assert_eq!(
1981            SessionProjection::replay(&events).unwrap().usage_for("r1"),
1982            (7, 3)
1983        );
1984        assert_eq!(
1985            SessionProjection::replay(&events)
1986                .unwrap()
1987                .billable_units_for("r1"),
1988            15
1989        );
1990        assert_eq!(
1991            SessionProjection::replay(&events)
1992                .unwrap()
1993                .usage_charge_for("r1"),
1994            UsageCharge {
1995                prompt_tokens: 7,
1996                completion_tokens: 3,
1997                cost_units: 5,
1998            }
1999        );
2000
2001        let mut attributed = events.clone();
2002        let details = MeteringDetails::Model {
2003            model: "m".into(),
2004            provider: Some("p".into()),
2005            provider_attempt_id: "r1:model:1:attempt:1".parse().unwrap(),
2006            source: MeteringSource::Estimated,
2007            outcome: MeteringOutcome::Unknown,
2008        };
2009        for entry in &mut attributed[4..] {
2010            if let Event::UsageRecorded { metering, .. } = &mut entry.event {
2011                *metering = Some(details.clone());
2012            }
2013        }
2014        assert_eq!(
2015            SessionProjection::replay(&attributed)
2016                .unwrap()
2017                .billable_units_for("r1"),
2018            15
2019        );
2020        if let Event::UsageRecorded {
2021            metering: Some(MeteringDetails::Model { source, .. }),
2022            ..
2023        } = &mut attributed[5].event
2024        {
2025            *source = MeteringSource::Reported;
2026        }
2027        assert!(matches!(
2028            SessionProjection::replay(&attributed),
2029            Err(EventError::UsageConflict(_))
2030        ));
2031        let mut prepared = attributed[..4].to_vec();
2032        prepared.push(event(
2033            5,
2034            Event::ModelRequestPrepared {
2035                metering: Some(details.clone()),
2036                run_id: "r1".parse().unwrap(),
2037                step: 1,
2038                attempt: 1,
2039                provider_attempt_id: "r1:model:1:attempt:1".into(),
2040                operation_id: "model:1:attempt:1".into(),
2041                reserved_prompt_tokens: 10,
2042                reserved_completion_tokens: 20,
2043                request: Value::Null,
2044                prompt_sections: Value::Null,
2045            },
2046        ));
2047        prepared.push(attributed[5].clone());
2048        assert_eq!(
2049            SessionProjection::replay(&prepared)
2050                .unwrap()
2051                .billable_units_for("r1"),
2052            15
2053        );
2054        if let Event::UsageRecorded {
2055            metering: Some(MeteringDetails::Model { model, .. }),
2056            ..
2057        } = &mut prepared[5].event
2058        {
2059            *model = "wrong".into();
2060        }
2061        assert!(matches!(
2062            SessionProjection::replay(&prepared),
2063            Err(EventError::UsageConflict(_))
2064        ));
2065        if let Event::UsageRecorded { metering, .. } = &mut prepared[5].event {
2066            *metering = None;
2067        }
2068        assert!(matches!(
2069            SessionProjection::replay(&prepared),
2070            Err(EventError::UsageConflict(_))
2071        ));
2072        let serialized = serde_json::to_value(&events[4]).unwrap();
2073        assert!(serialized["event"].get("metering").is_none());
2074        assert_eq!(
2075            serde_json::from_value::<SessionEvent>(serialized).unwrap(),
2076            events[4]
2077        );
2078
2079        events.push(event(
2080            7,
2081            Event::UsageRecorded {
2082                metering: None,
2083                run_id: "r1".parse().unwrap(),
2084                operation_id: "model:1:attempt:1".into(),
2085                prompt_tokens: 8,
2086                completion_tokens: 3,
2087                cost_units: 0,
2088            },
2089        ));
2090        assert_eq!(
2091            SessionProjection::replay(&events).unwrap_err(),
2092            EventError::UsageConflict("model:1:attempt:1".into())
2093        );
2094    }
2095
2096    #[test]
2097    fn unresolved_provider_attempt_is_conservatively_billable_and_reconcilable() {
2098        let mut events = vec![
2099            event(
2100                1,
2101                Event::SessionCreated {
2102                    profile_revision_id: "p1".parse().unwrap(),
2103                },
2104            ),
2105            event(
2106                2,
2107                Event::InputQueued {
2108                    input_id: "i1".parse().unwrap(),
2109                    run_id: "r1".parse().unwrap(),
2110                    mode: DeliveryMode::Followup,
2111                    content: text("hi"),
2112                    explicit_skill: None,
2113                },
2114            ),
2115            event(
2116                3,
2117                Event::InputClaimed {
2118                    input_id: "i1".parse().unwrap(),
2119                    run_id: "r1".parse().unwrap(),
2120                },
2121            ),
2122            event(
2123                4,
2124                Event::RunStarted {
2125                    run_id: "r1".parse().unwrap(),
2126                    input_id: "i1".parse().unwrap(),
2127                },
2128            ),
2129            event(
2130                5,
2131                Event::ModelRequestPrepared {
2132                    metering: None,
2133                    run_id: "r1".parse().unwrap(),
2134                    step: 1,
2135                    attempt: 1,
2136                    provider_attempt_id: "r1:model:1:attempt:1".into(),
2137                    operation_id: "model:1:attempt:1".into(),
2138                    reserved_prompt_tokens: 7,
2139                    reserved_completion_tokens: 11,
2140                    request: Value::Null,
2141                    prompt_sections: Value::Null,
2142                },
2143            ),
2144        ];
2145        assert_eq!(
2146            SessionProjection::replay(&events)
2147                .unwrap()
2148                .billable_units_for("r1"),
2149            18
2150        );
2151        events.push(event(
2152            6,
2153            Event::UsageRecorded {
2154                metering: None,
2155                run_id: "r1".parse().unwrap(),
2156                operation_id: "model:1:attempt:1".into(),
2157                prompt_tokens: 6,
2158                completion_tokens: 2,
2159                cost_units: 0,
2160            },
2161        ));
2162        assert_eq!(
2163            SessionProjection::replay(&events)
2164                .unwrap()
2165                .billable_units_for("r1"),
2166            8
2167        );
2168    }
2169
2170    #[test]
2171    fn a3_incremental_facts_preserve_every_lifecycle_and_usage_invariant() {
2172        let mut events = open_step_events();
2173        for fact in [
2174            Event::UserMessage {
2175                run_id: "r1".parse().unwrap(),
2176                content: text("history".repeat(10_000)),
2177            },
2178            Event::UsageRecorded {
2179                metering: None,
2180                run_id: "r1".parse().unwrap(),
2181                operation_id: "attempt".into(),
2182                prompt_tokens: 5,
2183                completion_tokens: 3,
2184                cost_units: 0,
2185            },
2186            Event::ToolCall {
2187                run_id: "r1".parse().unwrap(),
2188                step: 1,
2189                call_id: "pending".parse().unwrap(),
2190                tool: "echo".into(),
2191                arguments: Value::Null,
2192            },
2193        ] {
2194            events.push(event(events.len() as u64 + 1, fact));
2195        }
2196        let mut full = SessionProjection::replay(&events).unwrap();
2197        assert!(!full.messages.is_empty());
2198        full.messages.clear();
2199        full.injected_context.clear();
2200        let mut folded = SessionProjection::default();
2201        for page in events.chunks(3) {
2202            for fact in page {
2203                folded.apply_facts(fact).unwrap();
2204            }
2205            assert!(folded.messages.is_empty());
2206        }
2207        assert_eq!(folded, full);
2208        assert_eq!(SessionProjection::replay_facts(&events).unwrap(), full);
2209        assert_eq!(cancel_events(&folded), cancel_events(&full));
2210        assert_eq!(folded.usage_for("r1"), (5, 3));
2211        let invalid = event(
2212            events.len() as u64 + 1,
2213            Event::ToolResult {
2214                run_id: "r1".parse().unwrap(),
2215                step: 1,
2216                call_id: "missing".parse().unwrap(),
2217                result: Value::Null,
2218                is_error: true,
2219            },
2220        );
2221        assert_eq!(folded.apply_facts(&invalid), full.apply(&invalid));
2222    }
2223
2224    fn open_step_events() -> Vec<SessionEvent> {
2225        vec![
2226            event(
2227                1,
2228                Event::SessionCreated {
2229                    profile_revision_id: "p1".parse().unwrap(),
2230                },
2231            ),
2232            event(
2233                2,
2234                Event::InputQueued {
2235                    input_id: "i1".parse().unwrap(),
2236                    run_id: "r1".parse().unwrap(),
2237                    mode: DeliveryMode::Followup,
2238                    content: text("hi"),
2239                    explicit_skill: None,
2240                },
2241            ),
2242            event(
2243                3,
2244                Event::InputClaimed {
2245                    input_id: "i1".parse().unwrap(),
2246                    run_id: "r1".parse().unwrap(),
2247                },
2248            ),
2249            event(
2250                4,
2251                Event::RunStarted {
2252                    run_id: "r1".parse().unwrap(),
2253                    input_id: "i1".parse().unwrap(),
2254                },
2255            ),
2256            event(
2257                5,
2258                Event::TurnStarted {
2259                    run_id: "r1".parse().unwrap(),
2260                    turn: 1,
2261                },
2262            ),
2263            event(
2264                6,
2265                Event::StepStarted {
2266                    run_id: "r1".parse().unwrap(),
2267                    step: 1,
2268                },
2269            ),
2270        ]
2271    }
2272
2273    #[test]
2274    fn replay_rejects_overlapping_or_out_of_order_steps() {
2275        let mut overlapping = open_step_events();
2276        overlapping.push(event(
2277            7,
2278            Event::StepStarted {
2279                run_id: "r1".parse().unwrap(),
2280                step: 2,
2281            },
2282        ));
2283        assert_eq!(
2284            SessionProjection::replay(&overlapping).unwrap_err(),
2285            EventError::ConcurrentStep
2286        );
2287
2288        let mut skipped = open_step_events();
2289        skipped[5] = event(
2290            6,
2291            Event::StepStarted {
2292                run_id: "r1".parse().unwrap(),
2293                step: 2,
2294            },
2295        );
2296        assert_eq!(
2297            SessionProjection::replay(&skipped).unwrap_err(),
2298            EventError::ConcurrentStep
2299        );
2300    }
2301
2302    #[test]
2303    fn replay_rejects_cross_step_results_and_dangling_calls() {
2304        let mut events = open_step_events();
2305        events.push(event(
2306            7,
2307            Event::ToolCall {
2308                run_id: "r1".parse().unwrap(),
2309                step: 1,
2310                call_id: "c1".parse().unwrap(),
2311                tool: "echo".into(),
2312                arguments: Value::Null,
2313            },
2314        ));
2315        events.push(event(
2316            8,
2317            Event::ToolResult {
2318                run_id: "r1".parse().unwrap(),
2319                step: 2,
2320                call_id: "c1".parse().unwrap(),
2321                result: Value::Null,
2322                is_error: false,
2323            },
2324        ));
2325        assert_eq!(
2326            SessionProjection::replay(&events).unwrap_err(),
2327            EventError::LifecycleMismatch
2328        );
2329
2330        let mut dangling = open_step_events();
2331        dangling.push(event(
2332            7,
2333            Event::ToolCall {
2334                run_id: "r1".parse().unwrap(),
2335                step: 1,
2336                call_id: "c1".parse().unwrap(),
2337                tool: "echo".into(),
2338                arguments: Value::Null,
2339            },
2340        ));
2341        dangling.push(event(
2342            8,
2343            Event::StepFinished {
2344                run_id: "r1".parse().unwrap(),
2345                step: 1,
2346            },
2347        ));
2348        assert_eq!(
2349            SessionProjection::replay(&dangling).unwrap_err(),
2350            EventError::LifecycleMismatch
2351        );
2352    }
2353
2354    #[test]
2355    fn replay_rejects_reused_tool_call_ids() {
2356        let mut events = open_step_events();
2357        events.extend([
2358            event(
2359                7,
2360                Event::ToolCall {
2361                    run_id: "r1".parse().unwrap(),
2362                    step: 1,
2363                    call_id: "c1".parse().unwrap(),
2364                    tool: "echo".into(),
2365                    arguments: Value::Null,
2366                },
2367            ),
2368            event(
2369                8,
2370                Event::ToolResult {
2371                    run_id: "r1".parse().unwrap(),
2372                    step: 1,
2373                    call_id: "c1".parse().unwrap(),
2374                    result: Value::Null,
2375                    is_error: false,
2376                },
2377            ),
2378            event(
2379                9,
2380                Event::StepFinished {
2381                    run_id: "r1".parse().unwrap(),
2382                    step: 1,
2383                },
2384            ),
2385            event(
2386                10,
2387                Event::StepStarted {
2388                    run_id: "r1".parse().unwrap(),
2389                    step: 2,
2390                },
2391            ),
2392            event(
2393                11,
2394                Event::ToolCall {
2395                    run_id: "r1".parse().unwrap(),
2396                    step: 2,
2397                    call_id: "c1".parse().unwrap(),
2398                    tool: "echo".into(),
2399                    arguments: Value::Null,
2400                },
2401            ),
2402        ]);
2403        assert_eq!(
2404            SessionProjection::replay(&events).unwrap_err(),
2405            EventError::DuplicateToolCall("c1".parse().unwrap())
2406        );
2407    }
2408
2409    #[test]
2410    fn replay_keeps_injected_context_without_changing_run_state() {
2411        let mut events = open_step_events();
2412        events.push(event(
2413            7,
2414            Event::ContextInjected {
2415                run_id: "r1".parse().unwrap(),
2416                step: 1,
2417                contribution_id: "memory:1".into(),
2418                source: "memory".into(),
2419                version: "v1".into(),
2420                authority: "tenant".into(),
2421                form: "message".into(),
2422                content: text("tenant context"),
2423            },
2424        ));
2425
2426        let projection = SessionProjection::replay(&events).unwrap();
2427        assert_eq!(projection.active_run_id.as_deref(), Some("r1"));
2428        assert_eq!(projection.open_steps, BTreeSet::from([1]));
2429        assert_eq!(
2430            projection.injected_context,
2431            vec![ProjectedContext {
2432                run_id: "r1".parse().unwrap(),
2433                step: 1,
2434                contribution_id: "memory:1".into(),
2435                source: "memory".into(),
2436                version: "v1".into(),
2437                authority: "tenant".into(),
2438                form: "message".into(),
2439                content: text("tenant context"),
2440            }]
2441        );
2442    }
2443
2444    #[test]
2445    fn replay_skips_unknown_ignorable_formats_and_rejects_required_events() {
2446        let optional = event(
2447            1,
2448            Event::Opaque {
2449                format_version: SESSION_EVENT_FORMAT_VERSION,
2450                event_type: "future_optional".into(),
2451                ignorable: true,
2452                payload: serde_json::json!({"answer": 42}),
2453            },
2454        );
2455        let projection = SessionProjection::replay(&[optional]).unwrap();
2456        assert_eq!(projection.last_seq, 1);
2457
2458        let required = event(
2459            1,
2460            Event::Opaque {
2461                format_version: SESSION_EVENT_FORMAT_VERSION,
2462                event_type: "future_required".into(),
2463                ignorable: false,
2464                payload: Value::Null,
2465            },
2466        );
2467        assert_eq!(
2468            SessionProjection::replay(&[required]).unwrap_err(),
2469            EventError::UnknownRequired("future_required".into())
2470        );
2471
2472        let unsupported = event(
2473            1,
2474            Event::Opaque {
2475                format_version: SESSION_EVENT_FORMAT_VERSION + 1,
2476                event_type: "future_optional".into(),
2477                ignorable: true,
2478                payload: Value::Null,
2479            },
2480        );
2481        assert_eq!(
2482            SessionProjection::replay(&[unsupported]).unwrap().last_seq,
2483            1
2484        );
2485
2486        let required_new_format = event(
2487            1,
2488            Event::Opaque {
2489                format_version: SESSION_EVENT_FORMAT_VERSION + 1,
2490                event_type: "future_required".into(),
2491                ignorable: false,
2492                payload: Value::Null,
2493            },
2494        );
2495        assert_eq!(
2496            SessionProjection::replay(&[required_new_format]).unwrap_err(),
2497            EventError::UnsupportedFormat(SESSION_EVENT_FORMAT_VERSION + 1)
2498        );
2499    }
2500}