1use crate::envelope::{MsgKind, Outcome, Principal};
8use crate::ops::SessionProjection;
9use chrono::{DateTime, Utc};
10use schemars::JsonSchema;
11use serde::{Deserialize, Serialize};
12use serde_json::Value;
13
14pub const TURN_END_WITHOUT_COMPLETE: &str = "turn_end_without_complete";
17
18pub const DELIVERY_BLOCKED: &str = "delivery_blocked";
21
22pub const HANDOFF: &str = "handoff";
25
26pub const CLIENT_EVENT_CLASSES: [&str; 3] = [TURN_END_WITHOUT_COMPLETE, DELIVERY_BLOCKED, HANDOFF];
33
34pub fn client_event(class: &str, payload: Value) -> Option<Event> {
38 Some(match class {
39 TURN_END_WITHOUT_COMPLETE => Event::TurnEndWithoutComplete(payload),
40 DELIVERY_BLOCKED => Event::DeliveryBlocked(payload),
41 HANDOFF => Event::Handoff(payload),
42 _ => return None,
43 })
44}
45
46#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
48#[serde(rename_all = "snake_case")]
49pub enum Presence {
50 Online,
51 Offline,
52 Draining,
55}
56
57impl Presence {
58 pub fn as_str(self) -> &'static str {
59 match self {
60 Presence::Online => "online",
61 Presence::Offline => "offline",
62 Presence::Draining => "draining",
63 }
64 }
65}
66
67#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
69#[serde(rename_all = "snake_case")]
70pub enum GatewayHealth {
71 Online,
72 Reconnecting,
73 Failed,
74}
75
76impl GatewayHealth {
77 pub fn as_str(self) -> &'static str {
78 match self {
79 GatewayHealth::Online => "online",
80 GatewayHealth::Reconnecting => "reconnecting",
81 GatewayHealth::Failed => "failed",
82 }
83 }
84}
85
86#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema, Default)]
89#[serde(rename_all = "snake_case")]
90pub enum Lifecycle {
91 #[default]
92 Created,
93 Working,
94 Idle,
95 Exited,
96}
97
98#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
100#[serde(rename_all = "snake_case")]
101pub enum LedgerState {
102 Queued,
104 InFlight,
106 Acked,
108 Rejected,
110 Expired,
112}
113
114impl LedgerState {
115 pub fn as_str(self) -> &'static str {
116 match self {
117 LedgerState::Queued => "queued",
118 LedgerState::InFlight => "in_flight",
119 LedgerState::Acked => "acked",
120 LedgerState::Rejected => "rejected",
121 LedgerState::Expired => "expired",
122 }
123 }
124}
125
126impl std::fmt::Display for LedgerState {
127 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
128 f.write_str(self.as_str())
129 }
130}
131
132#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
133#[serde(rename_all = "snake_case")]
134pub struct RolePresence {
135 pub role: String,
136 pub state: Presence,
137 #[serde(default, skip_serializing_if = "Option::is_none")]
139 pub aggregate: Option<String>,
140 pub sessions: u32,
142 #[serde(default, skip_serializing_if = "Option::is_none")]
143 pub detail: Option<String>,
144}
145
146#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
147#[serde(rename_all = "snake_case", default)]
148pub struct SessionStateEvent {
149 #[serde(skip_serializing_if = "Option::is_none")]
154 pub task_id: Option<String>,
155 pub role: String,
156 pub session_id: String,
157 pub generation: u64,
158 pub seq: u64,
159 pub projection: SessionProjection,
160 #[serde(skip_serializing_if = "Option::is_none")]
163 pub admin: Option<Principal>,
164}
165
166#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
167#[serde(rename_all = "snake_case")]
168pub struct LedgerStateEvent {
169 pub msg_id: String,
170 pub op_id: Option<String>,
171 pub kind: MsgKind,
172 pub from: Principal,
173 pub to: Principal,
174 #[serde(default, skip_serializing_if = "Option::is_none")]
175 pub task: Option<String>,
176 pub state: LedgerState,
177 #[serde(default, skip_serializing_if = "Option::is_none")]
178 pub outcome: Option<Outcome>,
179 #[serde(default, skip_serializing_if = "Option::is_none")]
180 pub reason: Option<String>,
181}
182
183#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
184#[serde(rename_all = "snake_case", default)]
185pub struct FaultEvent {
186 pub id: i64,
188 #[serde(default, skip_serializing_if = "Option::is_none")]
189 pub task_id: Option<String>,
190 #[serde(default, skip_serializing_if = "Option::is_none")]
191 pub role: Option<String>,
192 #[serde(default, skip_serializing_if = "Option::is_none")]
193 pub session_id: Option<String>,
194 #[serde(default, skip_serializing_if = "Option::is_none")]
195 pub generation: Option<u64>,
196 #[serde(default, skip_serializing_if = "Option::is_none")]
197 pub seq: Option<u64>,
198 pub kind: String,
201 pub reason: String,
202 #[serde(default, skip_serializing_if = "Option::is_none")]
203 pub desired: Option<Value>,
204 #[serde(default, skip_serializing_if = "Option::is_none")]
205 pub observed: Option<Value>,
206 #[serde(default, skip_serializing_if = "Option::is_none")]
207 pub intent: Option<String>,
208 #[serde(default, skip_serializing_if = "Option::is_none")]
209 pub attempt: Option<u64>,
210 #[serde(default, skip_serializing_if = "Option::is_none")]
211 pub backend_ref: Option<Value>,
212 #[serde(default, skip_serializing_if = "Option::is_none")]
213 pub state: Option<String>,
214 #[serde(default, skip_serializing_if = "Option::is_none")]
215 pub created_at: Option<i64>,
216}
217
218#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
219#[serde(rename_all = "snake_case", default)]
220pub struct SpecReloaded {
221 pub spec_hash: String,
222 pub roles: u32,
223 pub gateways: u32,
224 pub routes: u32,
225}
226
227#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
229#[serde(rename_all = "snake_case", tag = "type", content = "data")]
230pub enum Event {
231 RolePresence(RolePresence),
233 SessionState(SessionStateEvent),
235 LedgerState(LedgerStateEvent),
237 Fault(FaultEvent),
239 GatewayPresence {
241 gateway: String,
242 platform: String,
243 state: GatewayHealth,
244 #[serde(default, skip_serializing_if = "Option::is_none")]
245 detail: Option<String>,
246 },
247 SpecReloaded(SpecReloaded),
249 TurnEndWithoutComplete(Value),
256 DeliveryBlocked(Value),
259 Handoff(Value),
261}
262
263impl Event {
264 pub fn type_name(&self) -> &'static str {
265 match self {
266 Event::RolePresence(_) => "role_presence",
267 Event::SessionState(_) => "session_state",
268 Event::LedgerState(_) => "ledger_state",
269 Event::Fault(_) => "fault",
270 Event::GatewayPresence { .. } => "gateway_presence",
271 Event::SpecReloaded(_) => "spec_reloaded",
272 Event::TurnEndWithoutComplete(_) => TURN_END_WITHOUT_COMPLETE,
273 Event::DeliveryBlocked(_) => DELIVERY_BLOCKED,
274 Event::Handoff(_) => HANDOFF,
275 }
276 }
277
278 pub fn tier(&self) -> EventTier {
281 match self {
282 Event::LedgerState(_) | Event::SessionState(_) => EventTier::Durable,
283 Event::RolePresence(_)
284 | Event::Fault(_)
285 | Event::GatewayPresence { .. }
286 | Event::SpecReloaded(_) => EventTier::Advisory,
287 Event::TurnEndWithoutComplete(_) | Event::DeliveryBlocked(_) | Event::Handoff(_) => {
291 EventTier::Durable
292 }
293 }
294 }
295}
296
297#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
299#[serde(rename_all = "snake_case")]
300pub enum EventTier {
301 Durable,
303 Advisory,
305}
306
307#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
309#[serde(rename_all = "snake_case")]
310pub struct EventRow {
311 pub seq: u64,
312 pub created_at: DateTime<Utc>,
313 pub event: Event,
314}