1use crate::{
2 ActiveRun, ConversationMessage, PendingInteraction, QueuedMessage, RunOutcome, ThreadSummary,
3};
4use serde::de::Error as _;
5use serde::{Deserialize, Deserializer, Serialize};
6use serde_json::Value;
7
8#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
9#[serde(rename_all = "camelCase")]
10pub struct SnapshotUpdatedEvent {
11 #[serde(rename = "type")]
12 pub event_type: String,
13 pub surface_id: String,
14 pub thread_id: String,
15 pub request_id: String,
16 pub snapshot_version: u64,
17 pub occurred_at: String,
18 pub data: SnapshotUpdatedData,
19}
20
21#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
22#[serde(rename_all = "camelCase")]
23pub struct SnapshotUpdatedData {
24 pub reason: String,
25 #[serde(default)]
26 pub session_activity_at_unix_ms: i64,
27 #[serde(default)]
28 pub messages_added: Vec<ConversationMessage>,
29 pub active_run: Option<ActiveRun>,
30 pub pending_interaction: Option<PendingInteraction>,
31 pub run_outcome: Option<RunOutcome>,
32}
33
34#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
35#[serde(rename_all = "camelCase")]
36pub struct QueueChangedEvent {
37 #[serde(rename = "type")]
38 pub event_type: String,
39 pub surface_id: String,
40 pub thread_id: String,
41 pub queue_version: u64,
42 pub occurred_at: String,
43 pub data: QueueChangedData,
44}
45
46#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
47#[serde(rename_all = "camelCase")]
48pub struct QueueChangedData {
49 pub items: Vec<QueuedMessage>,
50 pub active_run: Option<ActiveRun>,
51}
52
53#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
54#[serde(rename_all = "camelCase")]
55pub struct CatalogChangedEvent {
56 #[serde(rename = "type")]
57 pub event_type: String,
58 pub surface_id: String,
59 pub catalog_version: u64,
60 pub occurred_at: String,
61 pub data: CatalogChangedData,
62}
63
64#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
65#[serde(rename_all = "camelCase")]
66pub struct CatalogChangedData {
67 pub active_thread_id: Option<String>,
68 pub threads: Vec<ThreadSummary>,
69}
70
71#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
72#[serde(rename_all = "camelCase")]
73pub struct SessionPolicyChangedEvent {
74 #[serde(rename = "type")]
75 pub event_type: String,
76 pub version: u64,
77 pub surface_id: String,
78 pub occurred_at: String,
79 pub data: SessionPolicyChangedData,
80}
81
82#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
83#[serde(rename_all = "camelCase")]
84pub struct SessionPolicyChangedData {
85 pub idle_timeout_minutes: u32,
86 pub updated_at_ms: Option<i64>,
87}
88
89#[derive(Debug, Clone, Serialize, PartialEq)]
90#[serde(untagged)]
91#[allow(clippy::large_enum_variant)]
92pub enum ConversationEvent {
93 SnapshotUpdated(SnapshotUpdatedEvent),
94 CatalogChanged(CatalogChangedEvent),
95 QueueChanged(QueueChangedEvent),
96 SessionPolicyChanged(SessionPolicyChangedEvent),
97 Live(crate::LiveEvent),
98}
99
100impl<'de> Deserialize<'de> for ConversationEvent {
101 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
102 where
103 D: Deserializer<'de>,
104 {
105 let value = Value::deserialize(deserializer)?;
106 let event_type = value
107 .get("type")
108 .and_then(Value::as_str)
109 .ok_or_else(|| D::Error::custom("conversation event type is required"))?;
110 match event_type {
111 "snapshot.updated" => serde_json::from_value(value)
112 .map(Self::SnapshotUpdated)
113 .map_err(D::Error::custom),
114 "thread.catalog_changed" => serde_json::from_value(value)
115 .map(Self::CatalogChanged)
116 .map_err(D::Error::custom),
117 "queue.changed" => serde_json::from_value(value)
118 .map(Self::QueueChanged)
119 .map_err(D::Error::custom),
120 "session.policy_changed" => serde_json::from_value(value)
121 .map(Self::SessionPolicyChanged)
122 .map_err(D::Error::custom),
123 _ => serde_json::from_value(value)
124 .map(Self::Live)
125 .map_err(D::Error::custom),
126 }
127 }
128}
129
130impl ConversationEvent {
131 pub fn surface_id(&self) -> &str {
132 match self {
133 Self::SnapshotUpdated(event) => &event.surface_id,
134 Self::CatalogChanged(event) => &event.surface_id,
135 Self::QueueChanged(event) => &event.surface_id,
136 Self::SessionPolicyChanged(event) => &event.surface_id,
137 Self::Live(event) => &event.surface_id,
138 }
139 }
140
141 pub fn event_type(&self) -> &str {
142 match self {
143 Self::SnapshotUpdated(event) => &event.event_type,
144 Self::CatalogChanged(event) => &event.event_type,
145 Self::QueueChanged(event) => &event.event_type,
146 Self::SessionPolicyChanged(event) => &event.event_type,
147 Self::Live(event) => &event.event_type,
148 }
149 }
150}
151
152#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
153#[serde(rename_all = "camelCase")]
154pub struct DebugEvent {
155 #[serde(rename = "type")]
156 pub event_type: String,
157 pub surface_id: String,
158 pub thread_id: String,
159 pub request_id: String,
160 pub occurred_at: String,
161 pub data: DebugEventData,
162}
163
164#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
165pub struct DebugEventData {
166 pub scope: String,
167 pub payload: Value,
168}