Skip to main content

lash_remote_protocol/
observations.rs

1//! Session observation: cursors, resumable observation events, and live
2//! replay gap envelopes.
3
4use schemars::JsonSchema;
5use serde::{Deserialize, Serialize};
6
7use crate::registry_errors::{RemoteProtocolError, require_non_empty};
8use crate::usage_activity::{RemoteTurnActivity, RemoteUsage};
9use crate::{REMOTE_PROTOCOL_VERSION, ensure_protocol_version};
10
11/// Stable identity proving an admitted input became canonical turn input.
12///
13/// This wire DTO intentionally contains no display text. `input_id` and
14/// `source_key` correlate to the admission receipt, while `turn_id` and
15/// `committed_message_id` identify the canonical conversation application.
16#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
17pub struct RemoteTurnInputApplication {
18    pub input_id: String,
19    #[serde(default, skip_serializing_if = "Option::is_none")]
20    pub source_key: Option<String>,
21    pub turn_id: String,
22    pub committed_message_id: String,
23    #[serde(default, skip_serializing_if = "Option::is_none")]
24    pub checkpoint: Option<RemoteTurnInputCheckpoint>,
25}
26
27impl RemoteTurnInputApplication {
28    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
29        require_non_empty("RemoteTurnInputApplication", "input_id", &self.input_id)?;
30        require_non_empty("RemoteTurnInputApplication", "turn_id", &self.turn_id)?;
31        require_non_empty(
32            "RemoteTurnInputApplication",
33            "committed_message_id",
34            &self.committed_message_id,
35        )
36    }
37}
38
39#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
40#[serde(rename_all = "snake_case")]
41pub enum RemoteTurnInputCheckpoint {
42    AfterWork,
43    BeforeCompletion,
44}
45
46#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
47pub struct RemoteSessionCursor {
48    pub protocol_version: u32,
49    pub cursor: String,
50}
51
52impl RemoteSessionCursor {
53    pub fn new(cursor: impl Into<String>) -> Self {
54        Self {
55            protocol_version: REMOTE_PROTOCOL_VERSION,
56            cursor: cursor.into(),
57        }
58    }
59
60    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
61        ensure_protocol_version(self.protocol_version)?;
62        require_non_empty("RemoteSessionCursor", "cursor", &self.cursor)
63    }
64}
65
66#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
67pub struct RemoteSessionObservation {
68    pub protocol_version: u32,
69    pub session_id: String,
70    pub cursor: String,
71    pub turn_index: u64,
72    pub usage: RemoteUsage,
73}
74
75impl RemoteSessionObservation {
76    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
77        ensure_protocol_version(self.protocol_version)?;
78        require_non_empty("RemoteSessionObservation", "session_id", &self.session_id)?;
79        require_non_empty("RemoteSessionObservation", "cursor", &self.cursor)
80    }
81}
82
83#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
84pub struct RemoteSessionObservationEvent {
85    pub protocol_version: u32,
86    pub session_id: String,
87    pub revision: u64,
88    pub cursor: String,
89    #[serde(flatten)]
90    pub event: RemoteSessionObservationEventPayload,
91}
92
93impl RemoteSessionObservationEvent {
94    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
95        ensure_protocol_version(self.protocol_version)?;
96        require_non_empty(
97            "RemoteSessionObservationEvent",
98            "session_id",
99            &self.session_id,
100        )?;
101        require_non_empty("RemoteSessionObservationEvent", "cursor", &self.cursor)?;
102        if let RemoteSessionObservationEventPayload::TurnActivity { activity } = &self.event {
103            activity.validate()?;
104            if activity.protocol_version != self.protocol_version {
105                return Err(RemoteProtocolError::MismatchedNestedProtocolVersion {
106                    parent: "RemoteSessionObservationEvent",
107                    child: "activity",
108                    parent_version: self.protocol_version,
109                    child_version: activity.protocol_version,
110                });
111            }
112            if let crate::usage_activity::RemoteTurnEvent::TurnInputApplied { applications } =
113                &activity.event
114            {
115                for application in applications {
116                    application.validate()?;
117                }
118            }
119        }
120        Ok(())
121    }
122}
123
124#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
125#[serde(tag = "type", rename_all = "snake_case")]
126pub enum RemoteSessionObservationEventPayload {
127    TurnActivity {
128        activity: Box<RemoteTurnActivity>,
129    },
130    Committed,
131    AgentFrameSwitched {
132        frame_id: String,
133    },
134    QueueChanged {
135        kind: RemoteSessionQueueEventKind,
136        batch_ids: Vec<String>,
137    },
138    ProcessChanged {
139        kind: RemoteSessionProcessEventKind,
140        process_ids: Vec<String>,
141    },
142}
143
144#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
145#[serde(rename_all = "snake_case")]
146pub enum RemoteSessionQueueEventKind {
147    Enqueued,
148    Cancelled,
149}
150
151#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
152#[serde(rename_all = "snake_case")]
153pub enum RemoteSessionProcessEventKind {
154    Started,
155    Cancelled,
156}
157
158#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
159pub struct RemoteLiveReplayGap {
160    pub protocol_version: u32,
161    pub session_id: String,
162    pub requested_cursor: String,
163    pub latest_cursor: String,
164    pub latest_revision: u64,
165    pub reason: RemoteLiveReplayGapReason,
166}
167
168impl RemoteLiveReplayGap {
169    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
170        ensure_protocol_version(self.protocol_version)?;
171        require_non_empty("RemoteLiveReplayGap", "session_id", &self.session_id)?;
172        require_non_empty(
173            "RemoteLiveReplayGap",
174            "requested_cursor",
175            &self.requested_cursor,
176        )?;
177        require_non_empty("RemoteLiveReplayGap", "latest_cursor", &self.latest_cursor)
178    }
179}
180
181#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
182#[serde(rename_all = "snake_case")]
183pub enum RemoteLiveReplayGapReason {
184    Trimmed,
185    Unavailable,
186}