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;
#[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 {
pub task_id: String,
pub role: String,
pub session_id: String,
pub generation: u64,
pub seq: u64,
pub projection: SessionProjection,
}
#[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),
}
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",
}
}
pub fn tier(&self) -> EventTier {
match self {
Event::LedgerState(_) | Event::SessionState(_) => EventTier::Durable,
Event::RolePresence(_)
| Event::Fault(_)
| Event::GatewayPresence { .. }
| Event::SpecReloaded(_) => EventTier::Advisory,
}
}
}
#[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,
}
#[cfg(test)]
mod tests {
use super::*;
use crate::envelope::{Causality, new_task_id};
use crate::ops::SessionProjection;
#[test]
fn events_carry_a_wire_type_name_and_payload() {
let event = Event::RolePresence(RolePresence {
role: "builder".into(),
state: Presence::Offline,
aggregate: None,
sessions: 0,
detail: None,
});
let value = serde_json::to_value(&event).expect("encode");
assert_eq!(value["type"], "role_presence");
assert_eq!(value["data"]["state"], "offline");
let back: Event = serde_json::from_value(value).expect("decode");
assert_eq!(back, event);
}
#[test]
fn ledger_state_names_are_the_documented_strings() {
assert_eq!(LedgerState::InFlight.as_str(), "in_flight");
assert_eq!(LedgerState::Queued.as_str(), "queued");
assert_eq!(LedgerState::Expired.as_str(), "expired");
assert_eq!(Presence::Draining.as_str(), "draining");
assert_eq!(GatewayHealth::Reconnecting.as_str(), "reconnecting");
}
#[test]
fn durable_tiers_cover_ledger_and_session_only() {
let session = Event::SessionState(SessionStateEvent {
task_id: new_task_id(),
role: "builder".into(),
session_id: "s".into(),
generation: 1,
seq: 2,
projection: SessionProjection::default_working(),
});
assert_eq!(session.tier(), EventTier::Durable);
let fault = Event::Fault(FaultEvent {
id: 1,
task_id: Some(Causality::root(new_task_id()).task),
role: None,
session_id: None,
generation: None,
seq: None,
kind: "idle_fault".into(),
reason: "no heartbeat".into(),
desired: None,
observed: None,
intent: None,
attempt: None,
backend_ref: None,
state: None,
created_at: None,
});
assert_eq!(fault.tier(), EventTier::Advisory);
let value = serde_json::to_value(&fault).expect("encode");
assert_eq!(value["type"], "fault");
assert_eq!(value["data"]["kind"], "idle_fault");
}
#[test]
fn lifecycle_and_outcome_serialise_snake_case() {
assert_eq!(
serde_json::to_value(Lifecycle::Exited).expect("encode"),
Value::String("exited".into())
);
assert_eq!(
serde_json::to_value(Outcome::Cancelled).expect("encode"),
Value::String("cancelled".into())
);
}
}