use crate::envelope::{MsgKind, Outcome, Principal};
use crate::event::{
Event, FaultEvent, LedgerState, LedgerStateEvent, Lifecycle, Presence, SessionStateEvent,
};
use crate::lifecycle::{AgentPhase, DeliveryPhase, RecoveryPhase, ResourcePhase};
use crate::ops::{LedgerEntry, RoleInfo, SessionRow};
use chrono::{DateTime, Utc};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::BTreeMap;
pub const RESYNC_LAG_KIND: &str = "resync_lag";
pub const FAULT_STATE_OPEN: &str = "open";
pub const EVENT_TAIL_LIMIT: usize = 128;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum DeliveryAxis {
Queued,
InFlight,
Settled,
}
impl DeliveryAxis {
pub fn as_str(self) -> &'static str {
match self {
DeliveryAxis::Queued => "queued",
DeliveryAxis::InFlight => "in_flight",
DeliveryAxis::Settled => "settled",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum SessionState {
Opening,
Busy,
Idle,
Suspended,
Closed,
}
#[derive(
Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
)]
#[serde(rename_all = "snake_case")]
pub enum BoardColumn {
Queued,
Running,
Waiting,
Done,
FailedOrBlocked,
}
impl BoardColumn {
pub fn as_str(self) -> &'static str {
match self {
BoardColumn::Queued => "queued",
BoardColumn::Running => "running",
BoardColumn::Waiting => "waiting",
BoardColumn::Done => "done",
BoardColumn::FailedOrBlocked => "failed_or_blocked",
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub struct DeliveryView {
pub msg_id: String,
pub op_id: Option<String>,
pub kind: MsgKind,
pub from: Principal,
pub to: Principal,
pub task_id: Option<String>,
pub family: Option<String>,
pub hop: Option<u32>,
pub origin: Option<String>,
pub attempt: Option<u32>,
pub state: LedgerState,
pub outcome: Option<Outcome>,
pub reason: Option<String>,
pub out_head: Option<String>,
pub enqueued_at: Option<DateTime<Utc>>,
pub acked_at: Option<DateTime<Utc>>,
}
impl DeliveryView {
pub fn axis(&self) -> DeliveryAxis {
match self.state {
LedgerState::Queued => DeliveryAxis::Queued,
LedgerState::InFlight => DeliveryAxis::InFlight,
LedgerState::Acked | LedgerState::Rejected | LedgerState::Expired => {
DeliveryAxis::Settled
}
}
}
pub fn from_entry(entry: &LedgerEntry) -> Self {
Self {
msg_id: entry.msg_id.clone(),
op_id: entry.op_id.clone(),
kind: entry.kind,
from: entry.from.clone(),
to: entry.to.clone(),
task_id: entry.task.clone(),
family: entry.family.clone(),
hop: Some(entry.hop),
origin: entry.origin.clone(),
attempt: Some(entry.attempt),
state: entry.state,
outcome: None,
reason: entry.reason.clone(),
out_head: entry.out_head.clone(),
enqueued_at: Some(entry.enqueued_at),
acked_at: entry.acked_at,
}
}
fn from_event(reported: &LedgerStateEvent) -> Self {
Self {
msg_id: reported.msg_id.clone(),
op_id: reported.op_id.clone(),
kind: reported.kind,
from: reported.from.clone(),
to: reported.to.clone(),
task_id: reported.task.clone(),
family: None,
hop: None,
origin: None,
attempt: None,
state: reported.state,
outcome: reported.outcome,
reason: reported.reason.clone(),
out_head: None,
enqueued_at: None,
acked_at: None,
}
}
fn absorb(&mut self, reported: &LedgerStateEvent) {
self.op_id = reported.op_id.clone();
self.kind = reported.kind;
self.from = reported.from.clone();
self.to = reported.to.clone();
self.task_id = reported.task.clone();
self.state = reported.state;
self.outcome = reported.outcome;
self.reason = reported.reason.clone();
}
fn absorb_causality(&mut self, entry: &LedgerEntry) -> bool {
let mut filled = false;
filled |= fill_absent(&mut self.family, entry.family.clone());
filled |= fill_absent(&mut self.hop, Some(entry.hop));
filled |= fill_absent(&mut self.origin, entry.origin.clone());
filled |= fill_absent(&mut self.attempt, Some(entry.attempt));
filled |= fill_absent(&mut self.out_head, entry.out_head.clone());
filled |= fill_absent(&mut self.enqueued_at, Some(entry.enqueued_at));
filled |= fill_absent(&mut self.acked_at, entry.acked_at);
filled |= fill_absent(&mut self.reason, entry.reason.clone());
filled
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub struct SessionView {
pub session_id: String,
pub role: Option<String>,
pub task_id: Option<String>,
pub generation: u64,
pub seq: u64,
pub lifecycle: Lifecycle,
pub agent: AgentPhase,
pub delivery: DeliveryPhase,
pub resource: ResourcePhase,
pub recovery: RecoveryPhase,
pub outcome: Option<Outcome>,
pub observed: Option<Value>,
pub admin: Option<Principal>,
pub updated_at: Option<String>,
pub last_seen: Option<String>,
pub heartbeat_stale: bool,
}
impl SessionView {
pub fn state(&self) -> SessionState {
if self.lifecycle == Lifecycle::Exited || self.agent == AgentPhase::Gone {
return SessionState::Closed;
}
if self.lifecycle == Lifecycle::Working {
return SessionState::Busy;
}
if self.resource == ResourcePhase::Closed {
return SessionState::Suspended;
}
if self.lifecycle == Lifecycle::Created {
return SessionState::Opening;
}
SessionState::Idle
}
pub fn is_blocked(&self) -> bool {
self.outcome == Some(Outcome::Blocked)
}
pub fn from_row(row: &SessionRow) -> Self {
Self {
session_id: row.session_id.clone(),
role: row.role.clone(),
task_id: row.task_id.clone(),
generation: row.generation,
seq: row.seq,
lifecycle: row.projection.lifecycle,
agent: row.projection.agent,
delivery: row.projection.delivery,
resource: row.projection.resource,
recovery: row.projection.recovery,
outcome: row.projection.outcome,
observed: row.projection.observed.clone(),
admin: None,
updated_at: row.updated_at.clone(),
last_seen: row.last_seen.clone(),
heartbeat_stale: row.heartbeat_stale,
}
}
fn from_event(reported: &SessionStateEvent) -> Self {
Self {
session_id: reported.session_id.clone(),
role: Some(reported.role.clone()),
task_id: reported.task_id.clone(),
generation: reported.generation,
seq: reported.seq,
lifecycle: reported.projection.lifecycle,
agent: reported.projection.agent,
delivery: reported.projection.delivery,
resource: reported.projection.resource,
recovery: reported.projection.recovery,
outcome: reported.projection.outcome,
observed: reported.projection.observed.clone(),
admin: reported.admin.clone(),
updated_at: None,
last_seen: None,
heartbeat_stale: false,
}
}
fn absorb(&mut self, reported: &SessionStateEvent) {
self.role = Some(reported.role.clone());
self.task_id = reported.task_id.clone();
self.generation = reported.generation;
self.seq = reported.seq;
self.lifecycle = reported.projection.lifecycle;
self.agent = reported.projection.agent;
self.delivery = reported.projection.delivery;
self.resource = reported.projection.resource;
self.recovery = reported.projection.recovery;
self.outcome = reported.projection.outcome;
self.observed = reported.projection.observed.clone();
self.admin = reported.admin.clone();
self.updated_at = None;
self.last_seen = None;
self.heartbeat_stale = false;
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case", default)]
pub struct ClusterSummary {
pub cluster: Option<String>,
pub version: Option<String>,
pub spec_hash: Option<String>,
pub role_count: Option<u64>,
pub gateway_count: Option<u64>,
pub routes: Option<u64>,
pub channels: Option<u64>,
pub connected_roles: Option<u64>,
pub connected_gateways: Option<u64>,
pub event_head: Option<u64>,
pub uptime_s: Option<i64>,
}
impl ClusterSummary {
pub fn from_status(status: &Value) -> Self {
fn text(status: &Value, key: &str) -> Option<String> {
status.get(key)?.as_str().map(str::to_string)
}
fn number(status: &Value, key: &str) -> Option<u64> {
status.get(key)?.as_u64()
}
Self {
cluster: text(status, "cluster"),
version: text(status, "version"),
spec_hash: text(status, "spec_hash"),
role_count: number(status, "role_count"),
gateway_count: number(status, "gateway_count"),
routes: number(status, "routes"),
channels: number(status, "channels"),
connected_roles: number(status, "connected_roles"),
connected_gateways: number(status, "connected_gateways"),
event_head: number(status, "event_head"),
uptime_s: status.get("uptime_s").and_then(Value::as_i64),
}
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case", default)]
pub struct SessionCounts {
pub busy: u32,
pub idle: u32,
pub suspended: u32,
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct Card<'a> {
pub delivery: &'a DeliveryView,
pub session: Option<&'a SessionView>,
}
impl Card<'_> {
pub fn column(&self) -> BoardColumn {
match self.delivery.axis() {
DeliveryAxis::Queued => BoardColumn::Queued,
DeliveryAxis::InFlight => match self.session.map(SessionView::state) {
Some(SessionState::Busy) => BoardColumn::Running,
_ => BoardColumn::Waiting,
},
DeliveryAxis::Settled => match self.delivery.state {
LedgerState::Acked => match self.session.and_then(|session| session.outcome) {
Some(Outcome::Blocked) => BoardColumn::FailedOrBlocked,
_ => BoardColumn::Done,
},
_ => BoardColumn::FailedOrBlocked,
},
}
}
pub fn waits(&self) -> bool {
self.column() == BoardColumn::Waiting || self.session.is_some_and(SessionView::is_blocked)
}
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case", default)]
pub struct Snapshot {
pub status: Option<Value>,
pub roles: Vec<RoleInfo>,
pub sessions: Vec<SessionRow>,
pub ledger: Vec<LedgerEntry>,
pub faults: Vec<FaultEvent>,
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case", default)]
pub struct View {
pub cluster: ClusterSummary,
pub roles: BTreeMap<String, RoleInfo>,
pub sessions: BTreeMap<String, SessionView>,
pub deliveries: BTreeMap<String, DeliveryView>,
pub faults: BTreeMap<i64, FaultEvent>,
pub event_tail: Vec<Event>,
pub stale: bool,
}
impl View {
pub fn session_for(&self, task_id: &str) -> Option<&SessionView> {
self.sessions
.values()
.find(|session| session.task_id.as_deref() == Some(task_id))
}
pub fn cards<'a>(&'a self, role: &'a str) -> impl Iterator<Item = Card<'a>> + 'a {
self.deliveries
.values()
.filter_map(move |delivery| match principal_role(&delivery.to) {
Some(to) if to == role => Some(Card {
delivery,
session: delivery
.task_id
.as_deref()
.and_then(|task_id| self.session_for(task_id)),
}),
_ => None,
})
}
pub fn counts(&self, role: &str) -> SessionCounts {
let mut counts = SessionCounts::default();
for session in self.sessions.values() {
if session.role.as_deref() != Some(role) {
continue;
}
match session.state() {
SessionState::Busy => counts.busy += 1,
SessionState::Idle => counts.idle += 1,
SessionState::Suspended => counts.suspended += 1,
SessionState::Opening | SessionState::Closed => {}
}
}
counts
}
pub fn family<'a>(&'a self, family: &'a str) -> impl Iterator<Item = &'a DeliveryView> + 'a {
self.deliveries
.values()
.filter(move |delivery| delivery.family.as_deref() == Some(family))
}
pub fn open_faults(&self) -> impl Iterator<Item = &FaultEvent> {
self.faults.values().filter(|fault| fault_is_open(fault))
}
pub fn session_log<'a>(&'a self, session_id: &'a str) -> impl Iterator<Item = &'a Event> + 'a {
self.event_tail.iter().filter(move |event| match event {
Event::SessionState(reported) => reported.session_id == session_id,
_ => false,
})
}
pub fn merge_ledger(&mut self, ledger: &[LedgerEntry]) -> usize {
ledger
.iter()
.filter(|entry| {
self.deliveries
.get_mut(&entry.msg_id)
.is_some_and(|delivery| delivery.absorb_causality(entry))
})
.count()
}
}
pub fn fault_is_open(fault: &FaultEvent) -> bool {
fault
.state
.as_deref()
.is_none_or(|state| state == FAULT_STATE_OPEN)
}
pub fn is_resync_lag(event: &Event) -> bool {
matches!(event, Event::Fault(fault) if fault.kind == RESYNC_LAG_KIND)
}
pub fn snapshot_to_view(snapshot: &Snapshot) -> View {
let sessions: BTreeMap<String, SessionView> = snapshot
.sessions
.iter()
.map(|row| {
let session = SessionView::from_row(row);
(session.session_id.clone(), session)
})
.collect();
let deliveries: BTreeMap<String, DeliveryView> = snapshot
.ledger
.iter()
.map(|entry| {
let delivery = DeliveryView::from_entry(entry);
(delivery.msg_id.clone(), delivery)
})
.collect();
let faults: BTreeMap<i64, FaultEvent> = snapshot
.faults
.iter()
.map(|fault| (fault.id, fault.clone()))
.collect();
let roles: BTreeMap<String, RoleInfo> = snapshot
.roles
.iter()
.map(|role| (role.name.clone(), role.clone()))
.collect();
View {
cluster: snapshot
.status
.as_ref()
.map(ClusterSummary::from_status)
.unwrap_or_default(),
roles,
sessions,
deliveries,
faults,
event_tail: Vec::new(),
stale: false,
}
}
pub fn update(mut view: View, event: &Event) -> View {
let news = match event {
Event::RolePresence(reported) => {
if let Some(role) = view.roles.get_mut(&reported.role) {
role.state = reported.state;
role.sessions = reported.sessions;
role.aggregate = reported.aggregate.clone();
role.detail = reported.detail.clone();
}
view.cluster.connected_roles = Some(
view.roles
.values()
.filter(|role| role.state != Presence::Offline)
.count() as u64,
);
true
}
Event::SessionState(reported) => {
match view.sessions.get_mut(&reported.session_id) {
Some(stored)
if (reported.generation, reported.seq) <= (stored.generation, stored.seq) =>
{
false
}
Some(stored) => {
stored.absorb(reported);
true
}
None => {
let session = SessionView::from_event(reported);
view.sessions.insert(session.session_id.clone(), session);
true
}
}
}
Event::LedgerState(reported) => {
match view.deliveries.get_mut(&reported.msg_id) {
Some(stored) => stored.absorb(reported),
None => {
let delivery = DeliveryView::from_event(reported);
view.deliveries.insert(delivery.msg_id.clone(), delivery);
}
}
true
}
Event::Fault(reported) if reported.kind == RESYNC_LAG_KIND => {
view.stale = true;
false
}
Event::Fault(reported) => {
view.faults.insert(reported.id, reported.clone());
true
}
Event::GatewayPresence { .. } => true,
Event::SpecReloaded(_) => true,
Event::DeliveryBlocked(_) | Event::TurnEndWithoutComplete(_) | Event::Handoff(_) => true,
};
if news {
view.event_tail.insert(0, event.clone());
view.event_tail.truncate(EVENT_TAIL_LIMIT);
}
view
}
pub fn update_class(view: View, class: &str, payload: Value) -> View {
let encoded = serde_json::json!({ "type": class, "data": payload });
match serde_json::from_value::<Event>(encoded) {
Ok(event) => update(view, &event),
Err(_) => view,
}
}
fn principal_role(principal: &Principal) -> Option<&str> {
match principal {
Principal::Role { role, .. } => Some(role.as_str()),
Principal::Gateway { .. } | Principal::Cluster { .. } => None,
}
}
fn fill_absent<T>(slot: &mut Option<T>, value: Option<T>) -> bool {
match (slot.is_none(), value) {
(true, Some(value)) => {
*slot = Some(value);
true
}
_ => false,
}
}
#[cfg(test)]
mod tests {
use super::{LedgerEntry, Principal, Snapshot, View, snapshot_to_view, update};
use crate::envelope::MsgKind;
use crate::event::{Event, LedgerState, LedgerStateEvent};
use chrono::{DateTime, Utc};
#[derive(Clone, Copy)]
enum Born {
Stream,
Snapshot,
}
struct Case {
name: &'static str,
born: Born,
read: &'static [&'static str],
filled: usize,
}
fn fixed_time() -> DateTime<Utc> {
DateTime::parse_from_rfc3339("2026-09-29T00:00:00Z")
.expect("a fixed timestamp")
.with_timezone(&Utc)
}
fn entry(msg_id: &str) -> LedgerEntry {
LedgerEntry {
msg_id: msg_id.to_string(),
op_id: Some("o-1".into()),
kind: MsgKind::Task,
from: Principal::role("_supervisor"),
to: Principal::role("planner"),
task: Some("t1".into()),
parent_task: None,
hop: 2,
family: Some("t1".into()),
hop_budget: Some(8),
origin: Some("_supervisor".into()),
deadline: None,
labels: None,
attempt: 1,
state: LedgerState::Queued,
reason: None,
out_head: Some("h-1".into()),
body_json: None,
enqueued_at: fixed_time(),
acked_at: None,
}
}
fn transition(msg_id: &str) -> Event {
Event::LedgerState(LedgerStateEvent {
msg_id: msg_id.to_string(),
op_id: Some("o-1".into()),
kind: MsgKind::Task,
from: Principal::role("_supervisor"),
to: Principal::role("planner"),
task: Some("t1".into()),
state: LedgerState::InFlight,
outcome: None,
reason: None,
})
}
fn view_with_read_row(msg_id: &str) -> View {
snapshot_to_view(&Snapshot {
ledger: vec![entry(msg_id)],
..Snapshot::default()
})
}
#[test]
fn a_ledger_read_fills_only_what_a_stream_born_row_is_missing() {
let cases = [
Case {
name: "a row born on the stream gains its causality and its clock",
born: Born::Stream,
read: &["m1"],
filled: 1,
},
Case {
name: "a row the snapshot already read whole is left alone",
born: Born::Snapshot,
read: &["m1"],
filled: 0,
},
Case {
name: "a row this read does not name is left alone",
born: Born::Snapshot,
read: &["m9"],
filled: 0,
},
Case {
name: "the count is the rows filled, not the rows read",
born: Born::Stream,
read: &["m1", "m9"],
filled: 1,
},
];
for case in cases {
let mut view = match case.born {
Born::Stream => update(View::default(), &transition("m1")),
Born::Snapshot => view_with_read_row("m1"),
};
let ledger: Vec<LedgerEntry> = case.read.iter().map(|id| entry(id)).collect();
let before = view.clone();
let filled = view.merge_ledger(&ledger);
assert_eq!(
filled, case.filled,
"{}: the count is the rows a column moved on",
case.name
);
let delivery = view.deliveries.get("m1").expect("the view still holds m1");
match case.born {
Born::Stream => {
assert_eq!(delivery.family.as_deref(), Some("t1"), "{}", case.name);
assert_eq!(delivery.hop, Some(2), "{}", case.name);
assert_eq!(
delivery.origin.as_deref(),
Some("_supervisor"),
"{}",
case.name
);
assert_eq!(delivery.attempt, Some(1), "{}", case.name);
assert_eq!(delivery.out_head.as_deref(), Some("h-1"), "{}", case.name);
assert_eq!(delivery.enqueued_at, Some(fixed_time()), "{}", case.name);
assert_eq!(delivery.state, LedgerState::InFlight, "{}", case.name);
assert_eq!(delivery.acked_at, None, "{}", case.name);
assert_eq!(delivery.reason, None, "{}", case.name);
}
Born::Snapshot => assert_eq!(
delivery,
before.deliveries.get("m1").expect("the row was there"),
"{}: a row with nothing missing does not move",
case.name
),
}
assert_eq!(view.deliveries.len(), 1, "{}", case.name);
assert_eq!(view.event_tail, before.event_tail, "{}", case.name);
assert_eq!(view.stale, before.stale, "{}", case.name);
}
}
}