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