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}