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 offset: u64,
167 pub occurred_at: String,
168 pub data: AssistantPreviewData,
169}
170
171#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
172#[serde(rename_all = "camelCase", deny_unknown_fields)]
173pub struct AssistantPreviewData {
174 pub text: String,
175 pub append: bool,
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 offset: u64,
188 pub occurred_at: String,
189 pub data: RunProgressData,
190}
191
192#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
193#[serde(
194 tag = "kind",
195 rename_all = "snake_case",
196 rename_all_fields = "camelCase",
197 deny_unknown_fields
198)]
199pub enum RunProgressData {
200 Thinking {},
201 Planning {},
202 WaitingExternal {},
203 ToolScheduled { tool_names: Vec<String> },
204 ToolStarted { tool_name: String },
205}
206
207#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
208#[serde(rename_all = "camelCase", deny_unknown_fields)]
209pub struct RunSystemFailedEvent {
210 #[serde(rename = "type")]
211 pub event_type: RunSystemFailedEventType,
212 pub surface_id: ConversationSurface,
213 pub thread_id: String,
214 pub request_id: String,
215 pub base_snapshot_version: u64,
216 pub offset: u64,
217 pub occurred_at: String,
218 pub data: RunSystemFailedData,
219}
220
221#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
222#[serde(rename_all = "camelCase", deny_unknown_fields)]
223pub struct RunSystemFailedData {
224 pub code: String,
225 pub message: String,
226}
227
228#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
229#[serde(untagged)]
230pub enum LiveConversationEvent {
231 AssistantPreview(AssistantPreviewEvent),
232 RunProgress(RunProgressEvent),
233 RunSystemFailed(RunSystemFailedEvent),
234}
235
236#[derive(Debug, Clone, PartialEq, Eq)]
237pub struct LiveEventMetadata {
238 pub surface_id: ConversationSurface,
239 pub thread_id: String,
240 pub request_id: String,
241 pub base_snapshot_version: u64,
242 pub offset: u64,
243 pub occurred_at: String,
244}
245
246#[derive(Debug, Clone, PartialEq, Eq)]
247pub enum LiveConversationEventData {
248 AssistantPreview(AssistantPreviewData),
249 RunProgress(RunProgressData),
250 RunSystemFailed(RunSystemFailedData),
251}
252
253impl LiveConversationEvent {
254 pub fn new(metadata: LiveEventMetadata, data: LiveConversationEventData) -> Self {
255 match data {
256 LiveConversationEventData::AssistantPreview(data) => {
257 Self::AssistantPreview(AssistantPreviewEvent {
258 event_type: AssistantPreviewEventType::default(),
259 surface_id: metadata.surface_id,
260 thread_id: metadata.thread_id,
261 request_id: metadata.request_id,
262 base_snapshot_version: metadata.base_snapshot_version,
263 offset: metadata.offset,
264 occurred_at: metadata.occurred_at,
265 data,
266 })
267 }
268 LiveConversationEventData::RunProgress(data) => Self::RunProgress(RunProgressEvent {
269 event_type: RunProgressEventType::default(),
270 surface_id: metadata.surface_id,
271 thread_id: metadata.thread_id,
272 request_id: metadata.request_id,
273 base_snapshot_version: metadata.base_snapshot_version,
274 offset: metadata.offset,
275 occurred_at: metadata.occurred_at,
276 data,
277 }),
278 LiveConversationEventData::RunSystemFailed(data) => {
279 Self::RunSystemFailed(RunSystemFailedEvent {
280 event_type: RunSystemFailedEventType::default(),
281 surface_id: metadata.surface_id,
282 thread_id: metadata.thread_id,
283 request_id: metadata.request_id,
284 base_snapshot_version: metadata.base_snapshot_version,
285 offset: metadata.offset,
286 occurred_at: metadata.occurred_at,
287 data,
288 })
289 }
290 }
291 }
292
293 pub fn surface(&self) -> &ConversationSurface {
294 match self {
295 Self::AssistantPreview(event) => &event.surface_id,
296 Self::RunProgress(event) => &event.surface_id,
297 Self::RunSystemFailed(event) => &event.surface_id,
298 }
299 }
300
301 pub fn event_type(&self) -> &'static str {
302 match self {
303 Self::AssistantPreview(_) => "assistant.preview",
304 Self::RunProgress(_) => "run.progress",
305 Self::RunSystemFailed(_) => "run.system_failed",
306 }
307 }
308
309 pub fn base_snapshot_version(&self) -> u64 {
310 match self {
311 Self::AssistantPreview(event) => event.base_snapshot_version,
312 Self::RunProgress(event) => event.base_snapshot_version,
313 Self::RunSystemFailed(event) => event.base_snapshot_version,
314 }
315 }
316
317 pub fn offset(&self) -> u64 {
318 match self {
319 Self::AssistantPreview(event) => event.offset,
320 Self::RunProgress(event) => event.offset,
321 Self::RunSystemFailed(event) => event.offset,
322 }
323 }
324}
325
326#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
327#[serde(untagged)]
328#[allow(clippy::large_enum_variant)]
329pub enum ConversationEvent {
330 SnapshotUpdated(SnapshotUpdatedEvent),
331 CatalogChanged(CatalogChangedEvent),
332 QueueChanged(QueueChangedEvent),
333 SessionPolicyChanged(SessionPolicyChangedEvent),
334 AdmissionChanged(AdmissionChangedEvent),
335 AssistantPreview(AssistantPreviewEvent),
336 RunProgress(RunProgressEvent),
337 RunSystemFailed(RunSystemFailedEvent),
338}
339
340impl ConversationEvent {
341 pub fn surface(&self) -> &ConversationSurface {
342 match self {
343 Self::SnapshotUpdated(event) => &event.surface_id,
344 Self::CatalogChanged(event) => &event.surface_id,
345 Self::QueueChanged(event) => &event.surface_id,
346 Self::SessionPolicyChanged(event) => &event.surface_id,
347 Self::AdmissionChanged(event) => &event.surface_id,
348 Self::AssistantPreview(event) => &event.surface_id,
349 Self::RunProgress(event) => &event.surface_id,
350 Self::RunSystemFailed(event) => &event.surface_id,
351 }
352 }
353
354 pub fn event_type(&self) -> &'static str {
355 match self {
356 Self::SnapshotUpdated(_) => "snapshot.updated",
357 Self::CatalogChanged(_) => "thread.catalog_changed",
358 Self::QueueChanged(_) => "queue.changed",
359 Self::SessionPolicyChanged(_) => "session.policy_changed",
360 Self::AdmissionChanged(_) => "admission.changed",
361 Self::AssistantPreview(_) => "assistant.preview",
362 Self::RunProgress(_) => "run.progress",
363 Self::RunSystemFailed(_) => "run.system_failed",
364 }
365 }
366}
367
368#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
369#[serde(rename_all = "camelCase", deny_unknown_fields)]
370pub struct DebugEvent {
371 #[serde(rename = "type")]
372 pub event_type: DebugEventType,
373 pub surface_id: ConversationSurface,
374 pub thread_id: String,
375 pub request_id: String,
376 pub occurred_at: String,
377 pub data: DebugEventData,
378}
379
380#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
381#[serde(deny_unknown_fields)]
382pub struct DebugEventData {
383 pub scope: String,
384 pub payload: Value,
385}