Skip to main content

conversation_api/
event.rs

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}