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)]
91pub enum ConversationEvent {
92 SnapshotUpdated(SnapshotUpdatedEvent),
93 CatalogChanged(CatalogChangedEvent),
94 QueueChanged(QueueChangedEvent),
95 SessionPolicyChanged(SessionPolicyChangedEvent),
96 Live(crate::LiveEvent),
97}
98
99impl<'de> Deserialize<'de> for ConversationEvent {
100 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
101 where
102 D: Deserializer<'de>,
103 {
104 let value = Value::deserialize(deserializer)?;
105 let event_type = value
106 .get("type")
107 .and_then(Value::as_str)
108 .ok_or_else(|| D::Error::custom("conversation event type is required"))?;
109 match event_type {
110 "snapshot.updated" => serde_json::from_value(value)
111 .map(Self::SnapshotUpdated)
112 .map_err(D::Error::custom),
113 "thread.catalog_changed" => serde_json::from_value(value)
114 .map(Self::CatalogChanged)
115 .map_err(D::Error::custom),
116 "queue.changed" => serde_json::from_value(value)
117 .map(Self::QueueChanged)
118 .map_err(D::Error::custom),
119 "session.policy_changed" => serde_json::from_value(value)
120 .map(Self::SessionPolicyChanged)
121 .map_err(D::Error::custom),
122 _ => serde_json::from_value(value)
123 .map(Self::Live)
124 .map_err(D::Error::custom),
125 }
126 }
127}
128
129impl ConversationEvent {
130 pub fn surface_id(&self) -> &str {
131 match self {
132 Self::SnapshotUpdated(event) => &event.surface_id,
133 Self::CatalogChanged(event) => &event.surface_id,
134 Self::QueueChanged(event) => &event.surface_id,
135 Self::SessionPolicyChanged(event) => &event.surface_id,
136 Self::Live(event) => &event.surface_id,
137 }
138 }
139
140 pub fn event_type(&self) -> &str {
141 match self {
142 Self::SnapshotUpdated(event) => &event.event_type,
143 Self::CatalogChanged(event) => &event.event_type,
144 Self::QueueChanged(event) => &event.event_type,
145 Self::SessionPolicyChanged(event) => &event.event_type,
146 Self::Live(event) => &event.event_type,
147 }
148 }
149}
150
151#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
152#[serde(rename_all = "camelCase")]
153pub struct DebugEvent {
154 #[serde(rename = "type")]
155 pub event_type: String,
156 pub surface_id: String,
157 pub thread_id: String,
158 pub request_id: String,
159 pub occurred_at: String,
160 pub data: DebugEventData,
161}
162
163#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
164pub struct DebugEventData {
165 pub scope: String,
166 pub payload: Value,
167}