Skip to main content

conversation_api/
event.rs

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