Skip to main content

onlyne_proto/
event.rs

1//! The observation plane (§9, decision D11).
2//!
3//! Events are at-most-once pushes with a monotonic per-server `seq`. A
4//! subscriber that falls behind re-subscribes with `since_seq` and rebuilds from
5//! the ledger queries; nothing here replays by itself.
6
7use 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
14/// A turn ended with its task still open. Published by the client that
15/// witnessed the ending (`docs/v2-CONTRACT.md` §"Slice 7").
16pub const TURN_END_WITHOUT_COMPLETE: &str = "turn_end_without_complete";
17
18/// A delivery settled blocked. Published by the client that witnessed the
19/// settlement (`docs/v2-CONTRACT.md` §"Slice 7").
20pub const DELIVERY_BLOCKED: &str = "delivery_blocked";
21
22/// One turn handed work on instead of finishing. Published by the client
23/// that witnessed the handoff (`docs/v2-CONTRACT.md` §"Slice 7").
24pub const HANDOFF: &str = "handoff";
25
26/// The closed set of classes a client may publish (`ClientOp::PublishEvent`):
27/// the turn-end family and nothing else. A client owns these facts, so the
28/// server carries a client's word for them and never invents one. A class
29/// outside this set is refused by name, so a peer cannot invent an event
30/// name and a hook cannot bind to a spelling nobody publishes
31/// (`docs/v2-CONTRACT.md` §"Slice 7").
32pub const CLIENT_EVENT_CLASSES: [&str; 3] = [TURN_END_WITHOUT_COMPLETE, DELIVERY_BLOCKED, HANDOFF];
33
34/// Build the settlement event a client published, when `class` is in the
35/// closed set. `None` for a class the client may not publish, so the server's
36/// one check refuses by name rather than decode a free string into a variant.
37pub 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/// Liveness of a registered role's client connection.
47#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
48#[serde(rename_all = "snake_case")]
49pub enum Presence {
50    Online,
51    Offline,
52    /// Connected and refusing new deliveries while it drains running sessions
53    /// (decision D3).
54    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/// Health of a connected gateway process.
68#[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/// Public lifecycle projection of a session, as published by the client that
87/// owns it.
88#[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/// Ledger row state for one envelope (decision D10, D11).
99#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
100#[serde(rename_all = "snake_case")]
101pub enum LedgerState {
102    /// Accepted for a recipient that is not connected yet.
103    Queued,
104    /// Handed to the recipient's client, awaiting `ack`.
105    InFlight,
106    /// The recipient settled it.
107    Acked,
108    /// Refused at the gate; the sender got an error frame.
109    Rejected,
110    /// A note whose `ttl_ms` elapsed before delivery.
111    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    /// Aggregate cluster this role represents, when the spec declares one.
138    #[serde(default, skip_serializing_if = "Option::is_none")]
139    pub aggregate: Option<String>,
140    /// Live session count reported by the client.
141    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    /// The delivery this session was serving when the write landed, read off
150    /// the row's binding. Absent for a session no delivery is bound to, and the
151    /// field keeps its place so a write that names a delivery still encodes to
152    /// the same bytes it always did.
153    #[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    /// The operator who filed this write on the session's behalf over the
161    /// admin surface. Absent when the session reported on itself.
162    #[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    /// `faults.id` in the ledger.
187    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    /// Stable fault taxonomy, e.g. `intent_exhausted`, `idle_fault`,
199    /// `gateway_unconfigured`, `spec_reload_failed`.
200    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/// Observation-plane event. `seq` lives on the enclosing frame.
228#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
229#[serde(rename_all = "snake_case", tag = "type", content = "data")]
230pub enum Event {
231    /// A role's client connected, drained, or dropped.
232    RolePresence(RolePresence),
233    /// A session projection moved.
234    SessionState(SessionStateEvent),
235    /// A ledger row changed state.
236    LedgerState(LedgerStateEvent),
237    /// The server recorded a fault that needs a supervisor decision.
238    Fault(FaultEvent),
239    /// A gateway process changed health.
240    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    /// `spec.toml` was reloaded successfully.
248    SpecReloaded(SpecReloaded),
249    /// §3c's turn-end family, published by the client that witnessed it. The
250    /// server carries these on its stream so a hook bound to the class can
251    /// serve the plan's own example, and the client keeps no second copy
252    /// (`docs/v2-CONTRACT.md` §"Slice 7"). `type_name` is the class a hook
253    /// binds to, so the row's `type` column is the class and a hook's `on`
254    /// list filters by it exactly as a subscriber does.
255    TurnEndWithoutComplete(Value),
256    /// A delivery settled blocked: the work waits on something outside the
257    /// delivery, and a board reads it as waiting rather than as failed.
258    DeliveryBlocked(Value),
259    /// One turn handed work on instead of finishing.
260    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    /// Tiers drive subscription filtering: control-plane events are durable and
279    /// always replayable, observation events are ring-buffered.
280    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            // The settlement family is persisted in `events` and replayable
288            // from a cursor exactly as the durable class is, so a subscriber
289            // that reconnects with its last `seq` resumes with no gap.
290            Event::TurnEndWithoutComplete(_) | Event::DeliveryBlocked(_) | Event::Handoff(_) => {
291                EventTier::Durable
292            }
293        }
294    }
295}
296
297/// Subscription filter classes (§4 `subscribe.tiers`).
298#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
299#[serde(rename_all = "snake_case")]
300pub enum EventTier {
301    /// Persisted in `events`, replayable from a cursor.
302    Durable,
303    /// Best effort; a lagging subscriber resyncs by querying.
304    Advisory,
305}
306
307/// A timestamped event row, as returned by `watch` and `history`.
308#[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}