Skip to main content

conversation_api/
event.rs

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}