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