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    AssistantPreviewEventType,
37    AssistantPreview,
38    "assistant.preview"
39);
40event_type!(RunProgressEventType, RunProgress, "run.progress");
41event_type!(
42    RunSystemFailedEventType,
43    RunSystemFailed,
44    "run.system_failed"
45);
46event_type!(DebugEventType, RunDebug, "run.debug");
47
48#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
49#[serde(rename_all = "camelCase", deny_unknown_fields)]
50pub struct SnapshotUpdatedEvent {
51    #[serde(rename = "type")]
52    pub event_type: SnapshotUpdatedEventType,
53    pub surface_id: ConversationSurface,
54    pub thread_id: String,
55    pub request_id: String,
56    pub snapshot_version: u64,
57    pub occurred_at: String,
58    pub data: SnapshotUpdatedData,
59}
60
61#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
62#[serde(rename_all = "camelCase", deny_unknown_fields)]
63pub struct SnapshotUpdatedData {
64    #[serde(default)]
65    pub session_activity_at_unix_ms: i64,
66    #[serde(default)]
67    pub messages_added: Vec<ConversationMessage>,
68    pub active_run: Option<ActiveRun>,
69    pub pending_interaction: Option<PendingInteraction>,
70    pub run_outcome: Option<RunOutcome>,
71}
72
73#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
74#[serde(rename_all = "camelCase", deny_unknown_fields)]
75pub struct QueueChangedEvent {
76    #[serde(rename = "type")]
77    pub event_type: QueueChangedEventType,
78    pub surface_id: ConversationSurface,
79    pub thread_id: String,
80    pub queue_version: u64,
81    pub occurred_at: String,
82    pub data: QueueChangedData,
83}
84
85#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
86#[serde(rename_all = "camelCase", deny_unknown_fields)]
87pub struct QueueChangedData {
88    pub items: Vec<QueuedMessage>,
89    pub active_run: Option<ActiveRun>,
90}
91
92#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
93#[serde(rename_all = "camelCase", deny_unknown_fields)]
94pub struct CatalogChangedEvent {
95    #[serde(rename = "type")]
96    pub event_type: CatalogChangedEventType,
97    pub surface_id: ConversationSurface,
98    pub catalog_version: u64,
99    pub occurred_at: String,
100    pub data: CatalogChangedData,
101}
102
103#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
104#[serde(rename_all = "camelCase", deny_unknown_fields)]
105pub struct CatalogChangedData {
106    pub active_thread_id: Option<String>,
107    pub threads: Vec<ThreadSummary>,
108}
109
110#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
111#[serde(rename_all = "camelCase", deny_unknown_fields)]
112pub struct SessionPolicyChangedEvent {
113    #[serde(rename = "type")]
114    pub event_type: SessionPolicyChangedEventType,
115    pub version: u64,
116    pub surface_id: ConversationSurface,
117    pub occurred_at: String,
118    pub data: SessionPolicyChangedData,
119}
120
121#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
122#[serde(rename_all = "camelCase", deny_unknown_fields)]
123pub struct SessionPolicyChangedData {
124    pub idle_timeout_minutes: u32,
125    pub updated_at_ms: Option<i64>,
126}
127
128#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
129#[serde(rename_all = "camelCase", deny_unknown_fields)]
130pub struct AssistantPreviewEvent {
131    #[serde(rename = "type")]
132    pub event_type: AssistantPreviewEventType,
133    pub surface_id: ConversationSurface,
134    pub thread_id: String,
135    pub request_id: String,
136    pub base_snapshot_version: u64,
137    pub offset: u64,
138    pub occurred_at: String,
139    pub data: AssistantPreviewData,
140}
141
142#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
143#[serde(rename_all = "camelCase", deny_unknown_fields)]
144pub struct AssistantPreviewData {
145    pub text: String,
146    pub append: bool,
147}
148
149#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
150#[serde(rename_all = "camelCase", deny_unknown_fields)]
151pub struct RunProgressEvent {
152    #[serde(rename = "type")]
153    pub event_type: RunProgressEventType,
154    pub surface_id: ConversationSurface,
155    pub thread_id: String,
156    pub request_id: String,
157    pub base_snapshot_version: u64,
158    pub offset: u64,
159    pub occurred_at: String,
160    pub data: RunProgressData,
161}
162
163#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
164#[serde(
165    tag = "kind",
166    rename_all = "snake_case",
167    rename_all_fields = "camelCase",
168    deny_unknown_fields
169)]
170pub enum RunProgressData {
171    Thinking {},
172    Planning {},
173    WaitingExternal {},
174    ToolScheduled { tool_names: Vec<String> },
175    ToolStarted { tool_name: String },
176}
177
178#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
179#[serde(rename_all = "camelCase", deny_unknown_fields)]
180pub struct RunSystemFailedEvent {
181    #[serde(rename = "type")]
182    pub event_type: RunSystemFailedEventType,
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: RunSystemFailedData,
190}
191
192#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
193#[serde(rename_all = "camelCase", deny_unknown_fields)]
194pub struct RunSystemFailedData {
195    pub code: String,
196    pub message: String,
197}
198
199#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
200#[serde(untagged)]
201pub enum LiveConversationEvent {
202    AssistantPreview(AssistantPreviewEvent),
203    RunProgress(RunProgressEvent),
204    RunSystemFailed(RunSystemFailedEvent),
205}
206
207#[derive(Debug, Clone, PartialEq, Eq)]
208pub struct LiveEventMetadata {
209    pub surface_id: ConversationSurface,
210    pub thread_id: String,
211    pub request_id: String,
212    pub base_snapshot_version: u64,
213    pub offset: u64,
214    pub occurred_at: String,
215}
216
217#[derive(Debug, Clone, PartialEq, Eq)]
218pub enum LiveConversationEventData {
219    AssistantPreview(AssistantPreviewData),
220    RunProgress(RunProgressData),
221    RunSystemFailed(RunSystemFailedData),
222}
223
224impl LiveConversationEvent {
225    pub fn new(metadata: LiveEventMetadata, data: LiveConversationEventData) -> Self {
226        match data {
227            LiveConversationEventData::AssistantPreview(data) => {
228                Self::AssistantPreview(AssistantPreviewEvent {
229                    event_type: AssistantPreviewEventType::default(),
230                    surface_id: metadata.surface_id,
231                    thread_id: metadata.thread_id,
232                    request_id: metadata.request_id,
233                    base_snapshot_version: metadata.base_snapshot_version,
234                    offset: metadata.offset,
235                    occurred_at: metadata.occurred_at,
236                    data,
237                })
238            }
239            LiveConversationEventData::RunProgress(data) => Self::RunProgress(RunProgressEvent {
240                event_type: RunProgressEventType::default(),
241                surface_id: metadata.surface_id,
242                thread_id: metadata.thread_id,
243                request_id: metadata.request_id,
244                base_snapshot_version: metadata.base_snapshot_version,
245                offset: metadata.offset,
246                occurred_at: metadata.occurred_at,
247                data,
248            }),
249            LiveConversationEventData::RunSystemFailed(data) => {
250                Self::RunSystemFailed(RunSystemFailedEvent {
251                    event_type: RunSystemFailedEventType::default(),
252                    surface_id: metadata.surface_id,
253                    thread_id: metadata.thread_id,
254                    request_id: metadata.request_id,
255                    base_snapshot_version: metadata.base_snapshot_version,
256                    offset: metadata.offset,
257                    occurred_at: metadata.occurred_at,
258                    data,
259                })
260            }
261        }
262    }
263
264    pub fn surface(&self) -> &ConversationSurface {
265        match self {
266            Self::AssistantPreview(event) => &event.surface_id,
267            Self::RunProgress(event) => &event.surface_id,
268            Self::RunSystemFailed(event) => &event.surface_id,
269        }
270    }
271
272    pub fn event_type(&self) -> &'static str {
273        match self {
274            Self::AssistantPreview(_) => "assistant.preview",
275            Self::RunProgress(_) => "run.progress",
276            Self::RunSystemFailed(_) => "run.system_failed",
277        }
278    }
279
280    pub fn base_snapshot_version(&self) -> u64 {
281        match self {
282            Self::AssistantPreview(event) => event.base_snapshot_version,
283            Self::RunProgress(event) => event.base_snapshot_version,
284            Self::RunSystemFailed(event) => event.base_snapshot_version,
285        }
286    }
287
288    pub fn offset(&self) -> u64 {
289        match self {
290            Self::AssistantPreview(event) => event.offset,
291            Self::RunProgress(event) => event.offset,
292            Self::RunSystemFailed(event) => event.offset,
293        }
294    }
295}
296
297#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
298#[serde(untagged)]
299#[allow(clippy::large_enum_variant)]
300pub enum ConversationEvent {
301    SnapshotUpdated(SnapshotUpdatedEvent),
302    CatalogChanged(CatalogChangedEvent),
303    QueueChanged(QueueChangedEvent),
304    SessionPolicyChanged(SessionPolicyChangedEvent),
305    AssistantPreview(AssistantPreviewEvent),
306    RunProgress(RunProgressEvent),
307    RunSystemFailed(RunSystemFailedEvent),
308}
309
310impl ConversationEvent {
311    pub fn surface(&self) -> &ConversationSurface {
312        match self {
313            Self::SnapshotUpdated(event) => &event.surface_id,
314            Self::CatalogChanged(event) => &event.surface_id,
315            Self::QueueChanged(event) => &event.surface_id,
316            Self::SessionPolicyChanged(event) => &event.surface_id,
317            Self::AssistantPreview(event) => &event.surface_id,
318            Self::RunProgress(event) => &event.surface_id,
319            Self::RunSystemFailed(event) => &event.surface_id,
320        }
321    }
322
323    pub fn event_type(&self) -> &'static str {
324        match self {
325            Self::SnapshotUpdated(_) => "snapshot.updated",
326            Self::CatalogChanged(_) => "thread.catalog_changed",
327            Self::QueueChanged(_) => "queue.changed",
328            Self::SessionPolicyChanged(_) => "session.policy_changed",
329            Self::AssistantPreview(_) => "assistant.preview",
330            Self::RunProgress(_) => "run.progress",
331            Self::RunSystemFailed(_) => "run.system_failed",
332        }
333    }
334}
335
336#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
337#[serde(rename_all = "camelCase", deny_unknown_fields)]
338pub struct DebugEvent {
339    #[serde(rename = "type")]
340    pub event_type: DebugEventType,
341    pub surface_id: ConversationSurface,
342    pub thread_id: String,
343    pub request_id: String,
344    pub occurred_at: String,
345    pub data: DebugEventData,
346}
347
348#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
349#[serde(deny_unknown_fields)]
350pub struct DebugEventData {
351    pub scope: String,
352    pub payload: Value,
353}