1use crate::{
2 ActiveRun, ConversationMessage, ConversationSurface, PendingInteraction, QueuedMessage,
3 RunOutcome, ThreadSummary,
4};
5use serde::{Deserialize, Serialize};
6use serde_json::Value;
7
8macro_rules! event_type {
9 ($name:ident, $variant:ident, $wire:literal) => {
10 #[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
11 pub enum $name {
12 #[default]
13 #[serde(rename = $wire)]
14 $variant,
15 }
16 };
17}
18
19event_type!(
20 SnapshotUpdatedEventType,
21 SnapshotUpdated,
22 "snapshot.updated"
23);
24event_type!(
25 CatalogChangedEventType,
26 CatalogChanged,
27 "thread.catalog_changed"
28);
29event_type!(QueueChangedEventType, QueueChanged, "queue.changed");
30event_type!(
31 SessionPolicyChangedEventType,
32 SessionPolicyChanged,
33 "session.policy_changed"
34);
35event_type!(
36 AdmissionChangedEventType,
37 AdmissionChanged,
38 "admission.changed"
39);
40event_type!(
41 AssistantPreviewEventType,
42 AssistantPreview,
43 "assistant.preview"
44);
45event_type!(RunProgressEventType, RunProgress, "run.progress");
46event_type!(
47 RunSystemFailedEventType,
48 RunSystemFailed,
49 "run.system_failed"
50);
51event_type!(DebugEventType, RunDebug, "run.debug");
52
53#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
54#[serde(rename_all = "camelCase", deny_unknown_fields)]
55pub struct SnapshotUpdatedEvent {
56 #[serde(rename = "type")]
57 pub event_type: SnapshotUpdatedEventType,
58 pub surface_id: ConversationSurface,
59 pub thread_id: String,
60 pub request_id: String,
61 pub snapshot_version: u64,
62 pub occurred_at: String,
63 pub data: SnapshotUpdatedData,
64}
65
66#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
67#[serde(rename_all = "camelCase", deny_unknown_fields)]
68pub struct SnapshotUpdatedData {
69 #[serde(default)]
70 pub session_activity_at_unix_ms: i64,
71 #[serde(default)]
72 pub messages_added: Vec<ConversationMessage>,
73 pub active_run: Option<ActiveRun>,
74 pub pending_interaction: Option<PendingInteraction>,
75 pub run_outcome: Option<RunOutcome>,
76}
77
78#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
79#[serde(rename_all = "camelCase", deny_unknown_fields)]
80pub struct QueueChangedEvent {
81 #[serde(rename = "type")]
82 pub event_type: QueueChangedEventType,
83 pub surface_id: ConversationSurface,
84 pub thread_id: String,
85 pub queue_version: u64,
86 pub occurred_at: String,
87 pub data: QueueChangedData,
88}
89
90#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
91#[serde(rename_all = "camelCase", deny_unknown_fields)]
92pub struct QueueChangedData {
93 pub items: Vec<QueuedMessage>,
94 pub active_run: Option<ActiveRun>,
95}
96
97#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
98#[serde(rename_all = "camelCase", deny_unknown_fields)]
99pub struct CatalogChangedEvent {
100 #[serde(rename = "type")]
101 pub event_type: CatalogChangedEventType,
102 pub surface_id: ConversationSurface,
103 pub catalog_version: u64,
104 pub occurred_at: String,
105 pub data: CatalogChangedData,
106}
107
108#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
109#[serde(rename_all = "camelCase", deny_unknown_fields)]
110pub struct CatalogChangedData {
111 pub active_thread_id: Option<String>,
112 pub threads: Vec<ThreadSummary>,
113}
114
115#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
116#[serde(rename_all = "camelCase", deny_unknown_fields)]
117pub struct SessionPolicyChangedEvent {
118 #[serde(rename = "type")]
119 pub event_type: SessionPolicyChangedEventType,
120 pub version: u64,
121 pub surface_id: ConversationSurface,
122 pub occurred_at: String,
123 pub data: SessionPolicyChangedData,
124}
125
126#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
127#[serde(rename_all = "camelCase", deny_unknown_fields)]
128pub struct SessionPolicyChangedData {
129 pub idle_timeout_minutes: u32,
130 pub updated_at_ms: Option<i64>,
131}
132
133#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
134#[serde(rename_all = "camelCase", deny_unknown_fields)]
135pub struct AdmissionChangedEvent {
136 #[serde(rename = "type")]
137 pub event_type: AdmissionChangedEventType,
138 pub surface_id: ConversationSurface,
139 pub request_id: String,
140 pub admission_version: u64,
141 pub occurred_at: String,
142 pub data: AdmissionChangedData,
143}
144
145#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
146#[serde(rename_all = "camelCase", deny_unknown_fields)]
147pub struct AdmissionChangedData {
148 pub state: crate::RequestAdmissionState,
149 #[serde(default, skip_serializing_if = "Option::is_none")]
150 pub reason: Option<crate::RequestAdmissionReason>,
151 #[serde(default, skip_serializing_if = "Option::is_none")]
152 pub thread_id: Option<String>,
153 #[serde(default, skip_serializing_if = "Option::is_none")]
154 pub queue_version: Option<u64>,
155}
156
157#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
158#[serde(rename_all = "camelCase", deny_unknown_fields)]
159pub struct AssistantPreviewEvent {
160 #[serde(rename = "type")]
161 pub event_type: AssistantPreviewEventType,
162 pub surface_id: ConversationSurface,
163 pub thread_id: String,
164 pub request_id: String,
165 pub base_snapshot_version: u64,
166 pub sequence: u64,
168 pub occurred_at: String,
169 pub data: AssistantPreviewData,
170}
171
172#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
173#[serde(rename_all = "camelCase", deny_unknown_fields)]
174pub struct AssistantPreviewData {
175 pub text: String,
176}
177
178#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
179#[serde(rename_all = "camelCase", deny_unknown_fields)]
180pub struct RunProgressEvent {
181 #[serde(rename = "type")]
182 pub event_type: RunProgressEventType,
183 pub surface_id: ConversationSurface,
184 pub thread_id: String,
185 pub request_id: String,
186 pub base_snapshot_version: u64,
187 pub sequence: u64,
189 pub occurred_at: String,
190 pub data: RunProgressData,
191}
192
193#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
194#[serde(
195 tag = "kind",
196 rename_all = "snake_case",
197 rename_all_fields = "camelCase",
198 deny_unknown_fields
199)]
200pub enum RunProgressData {
201 Thinking {},
202 Planning {},
203 WaitingExternal {},
204 ToolScheduled { tool_names: Vec<String> },
205 ToolStarted { tool_name: String },
206}
207
208#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
209#[serde(rename_all = "camelCase", deny_unknown_fields)]
210pub struct RunSystemFailedEvent {
211 #[serde(rename = "type")]
212 pub event_type: RunSystemFailedEventType,
213 pub surface_id: ConversationSurface,
214 pub thread_id: String,
215 pub request_id: String,
216 pub base_snapshot_version: u64,
217 pub sequence: u64,
219 pub occurred_at: String,
220 pub data: RunSystemFailedData,
221}
222
223#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
224#[serde(rename_all = "camelCase", deny_unknown_fields)]
225pub struct RunSystemFailedData {
226 pub code: String,
227 pub message: String,
228}
229
230#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
231#[serde(untagged)]
232pub enum LiveConversationEvent {
233 AssistantPreview(AssistantPreviewEvent),
234 RunProgress(RunProgressEvent),
235 RunSystemFailed(RunSystemFailedEvent),
236}
237
238#[derive(Debug, Clone, PartialEq, Eq)]
239pub struct LiveEventMetadata {
240 pub surface_id: ConversationSurface,
241 pub thread_id: String,
242 pub request_id: String,
243 pub base_snapshot_version: u64,
244 pub sequence: u64,
246 pub occurred_at: String,
247}
248
249#[derive(Debug, Clone, PartialEq, Eq)]
250pub enum LiveConversationEventData {
251 AssistantPreview(AssistantPreviewData),
252 RunProgress(RunProgressData),
253 RunSystemFailed(RunSystemFailedData),
254}
255
256impl LiveConversationEvent {
257 pub fn new(metadata: LiveEventMetadata, data: LiveConversationEventData) -> Self {
258 match data {
259 LiveConversationEventData::AssistantPreview(data) => {
260 Self::AssistantPreview(AssistantPreviewEvent {
261 event_type: AssistantPreviewEventType::default(),
262 surface_id: metadata.surface_id,
263 thread_id: metadata.thread_id,
264 request_id: metadata.request_id,
265 base_snapshot_version: metadata.base_snapshot_version,
266 sequence: metadata.sequence,
267 occurred_at: metadata.occurred_at,
268 data,
269 })
270 }
271 LiveConversationEventData::RunProgress(data) => Self::RunProgress(RunProgressEvent {
272 event_type: RunProgressEventType::default(),
273 surface_id: metadata.surface_id,
274 thread_id: metadata.thread_id,
275 request_id: metadata.request_id,
276 base_snapshot_version: metadata.base_snapshot_version,
277 sequence: metadata.sequence,
278 occurred_at: metadata.occurred_at,
279 data,
280 }),
281 LiveConversationEventData::RunSystemFailed(data) => {
282 Self::RunSystemFailed(RunSystemFailedEvent {
283 event_type: RunSystemFailedEventType::default(),
284 surface_id: metadata.surface_id,
285 thread_id: metadata.thread_id,
286 request_id: metadata.request_id,
287 base_snapshot_version: metadata.base_snapshot_version,
288 sequence: metadata.sequence,
289 occurred_at: metadata.occurred_at,
290 data,
291 })
292 }
293 }
294 }
295
296 pub fn surface(&self) -> &ConversationSurface {
297 match self {
298 Self::AssistantPreview(event) => &event.surface_id,
299 Self::RunProgress(event) => &event.surface_id,
300 Self::RunSystemFailed(event) => &event.surface_id,
301 }
302 }
303
304 pub fn event_type(&self) -> &'static str {
305 match self {
306 Self::AssistantPreview(_) => "assistant.preview",
307 Self::RunProgress(_) => "run.progress",
308 Self::RunSystemFailed(_) => "run.system_failed",
309 }
310 }
311
312 pub fn base_snapshot_version(&self) -> u64 {
313 match self {
314 Self::AssistantPreview(event) => event.base_snapshot_version,
315 Self::RunProgress(event) => event.base_snapshot_version,
316 Self::RunSystemFailed(event) => event.base_snapshot_version,
317 }
318 }
319
320 pub fn sequence(&self) -> u64 {
321 match self {
322 Self::AssistantPreview(event) => event.sequence,
323 Self::RunProgress(event) => event.sequence,
324 Self::RunSystemFailed(event) => event.sequence,
325 }
326 }
327}
328
329#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
330#[serde(untagged)]
331#[allow(clippy::large_enum_variant)]
332pub enum ConversationEvent {
333 ControlUpdated(crate::ControlUpdatedEvent),
334 SnapshotUpdated(SnapshotUpdatedEvent),
335 CatalogChanged(CatalogChangedEvent),
336 QueueChanged(QueueChangedEvent),
337 SessionPolicyChanged(SessionPolicyChangedEvent),
338 AdmissionChanged(AdmissionChangedEvent),
339 AssistantPreview(AssistantPreviewEvent),
340 RunProgress(RunProgressEvent),
341 RunSystemFailed(RunSystemFailedEvent),
342}
343
344impl ConversationEvent {
345 pub fn surface(&self) -> &ConversationSurface {
346 match self {
347 Self::ControlUpdated(event) => &event.surface_id,
348 Self::SnapshotUpdated(event) => &event.surface_id,
349 Self::CatalogChanged(event) => &event.surface_id,
350 Self::QueueChanged(event) => &event.surface_id,
351 Self::SessionPolicyChanged(event) => &event.surface_id,
352 Self::AdmissionChanged(event) => &event.surface_id,
353 Self::AssistantPreview(event) => &event.surface_id,
354 Self::RunProgress(event) => &event.surface_id,
355 Self::RunSystemFailed(event) => &event.surface_id,
356 }
357 }
358
359 pub fn event_type(&self) -> &'static str {
360 match self {
361 Self::ControlUpdated(_) => "control.updated",
362 Self::SnapshotUpdated(_) => "snapshot.updated",
363 Self::CatalogChanged(_) => "thread.catalog_changed",
364 Self::QueueChanged(_) => "queue.changed",
365 Self::SessionPolicyChanged(_) => "session.policy_changed",
366 Self::AdmissionChanged(_) => "admission.changed",
367 Self::AssistantPreview(_) => "assistant.preview",
368 Self::RunProgress(_) => "run.progress",
369 Self::RunSystemFailed(_) => "run.system_failed",
370 }
371 }
372}
373
374#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
375#[serde(rename_all = "camelCase", deny_unknown_fields)]
376pub struct DebugEvent {
377 #[serde(rename = "type")]
378 pub event_type: DebugEventType,
379 pub surface_id: ConversationSurface,
380 pub thread_id: String,
381 pub request_id: String,
382 pub occurred_at: String,
383 pub data: DebugEventData,
384}
385
386#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
387#[serde(deny_unknown_fields)]
388pub struct DebugEventData {
389 pub scope: String,
390 pub payload: Value,
391}