use crate::envelope::{MsgKind, Outcome, Principal};
use crate::ops::SessionProjection;
use chrono::{DateTime, Utc};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use serde_json::Value;
pub const TURN_END_WITHOUT_COMPLETE: &str = "turn_end_without_complete";
pub const DELIVERY_BLOCKED: &str = "delivery_blocked";
pub const HANDOFF: &str = "handoff";
pub const CLIENT_EVENT_CLASSES: [&str; 3] = [TURN_END_WITHOUT_COMPLETE, DELIVERY_BLOCKED, HANDOFF];
pub fn client_event(class: &str, payload: Value) -> Option<Event> {
Some(match class {
TURN_END_WITHOUT_COMPLETE => Event::TurnEndWithoutComplete(payload),
DELIVERY_BLOCKED => Event::DeliveryBlocked(payload),
HANDOFF => Event::Handoff(payload),
_ => return None,
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum Presence {
Online,
Offline,
Draining,
}
impl Presence {
pub fn as_str(self) -> &'static str {
match self {
Presence::Online => "online",
Presence::Offline => "offline",
Presence::Draining => "draining",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum GatewayHealth {
Online,
Reconnecting,
Failed,
}
impl GatewayHealth {
pub fn as_str(self) -> &'static str {
match self {
GatewayHealth::Online => "online",
GatewayHealth::Reconnecting => "reconnecting",
GatewayHealth::Failed => "failed",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case")]
pub enum Lifecycle {
#[default]
Created,
Working,
Idle,
Exited,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum LedgerState {
Queued,
InFlight,
Acked,
Rejected,
Expired,
}
impl LedgerState {
pub fn as_str(self) -> &'static str {
match self {
LedgerState::Queued => "queued",
LedgerState::InFlight => "in_flight",
LedgerState::Acked => "acked",
LedgerState::Rejected => "rejected",
LedgerState::Expired => "expired",
}
}
}
impl std::fmt::Display for LedgerState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub struct RolePresence {
pub role: String,
pub state: Presence,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub aggregate: Option<String>,
pub sessions: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub detail: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", default)]
pub struct SessionStateEvent {
#[serde(skip_serializing_if = "Option::is_none")]
pub task_id: Option<String>,
pub role: String,
pub session_id: String,
pub generation: u64,
pub seq: u64,
pub projection: SessionProjection,
#[serde(skip_serializing_if = "Option::is_none")]
pub admin: Option<Principal>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub struct LedgerStateEvent {
pub msg_id: String,
pub op_id: Option<String>,
pub kind: MsgKind,
pub from: Principal,
pub to: Principal,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub task: Option<String>,
pub state: LedgerState,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub outcome: Option<Outcome>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reason: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", default)]
pub struct FaultEvent {
pub id: i64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub task_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub role: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub session_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub generation: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub seq: Option<u64>,
pub kind: String,
pub reason: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub desired: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub observed: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub intent: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub attempt: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub backend_ref: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub state: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub created_at: Option<i64>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", default)]
pub struct SpecReloaded {
pub spec_hash: String,
pub roles: u32,
pub gateways: u32,
pub routes: u32,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case", tag = "type", content = "data")]
pub enum Event {
RolePresence(RolePresence),
SessionState(SessionStateEvent),
LedgerState(LedgerStateEvent),
Fault(FaultEvent),
GatewayPresence {
gateway: String,
platform: String,
state: GatewayHealth,
#[serde(default, skip_serializing_if = "Option::is_none")]
detail: Option<String>,
},
SpecReloaded(SpecReloaded),
TurnEndWithoutComplete(Value),
DeliveryBlocked(Value),
Handoff(Value),
}
impl Event {
pub fn type_name(&self) -> &'static str {
match self {
Event::RolePresence(_) => "role_presence",
Event::SessionState(_) => "session_state",
Event::LedgerState(_) => "ledger_state",
Event::Fault(_) => "fault",
Event::GatewayPresence { .. } => "gateway_presence",
Event::SpecReloaded(_) => "spec_reloaded",
Event::TurnEndWithoutComplete(_) => TURN_END_WITHOUT_COMPLETE,
Event::DeliveryBlocked(_) => DELIVERY_BLOCKED,
Event::Handoff(_) => HANDOFF,
}
}
pub fn tier(&self) -> EventTier {
match self {
Event::LedgerState(_) | Event::SessionState(_) => EventTier::Durable,
Event::RolePresence(_)
| Event::Fault(_)
| Event::GatewayPresence { .. }
| Event::SpecReloaded(_) => EventTier::Advisory,
Event::TurnEndWithoutComplete(_) | Event::DeliveryBlocked(_) | Event::Handoff(_) => {
EventTier::Durable
}
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum EventTier {
Durable,
Advisory,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub struct EventRow {
pub seq: u64,
pub created_at: DateTime<Utc>,
pub event: Event,
}