1use 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#[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}