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}