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    /// Source-owned ordering within this snapshot generation. Missing numbers are allowed.
167    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    /// Source-owned ordering within this snapshot generation. Missing numbers are allowed.
188    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    /// Source-owned ordering within this snapshot generation. Missing numbers are allowed.
218    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    /// Source-owned ordering within this snapshot generation. Missing numbers are allowed.
245    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}