Skip to main content

nmbrs_runtime/checkpoint/
events.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! SRD-44a — Checkpoint event taxonomy.
5//!
6//! `CheckpointData` is the on-disk record type for the
7//! append-only `checkpoint.jsonl` event log. Every state-
8//! changing observation is one variant; serde tags each line
9//! with a `"type"` discriminator so a reader can drop
10//! unrecognised types without poisoning the rest of the
11//! stream.
12//!
13//! See SRD-44a §"Event taxonomy" for the schema.
14
15use serde::{Deserialize, Serialize};
16
17use super::identity::PhaseIdentity;
18use super::storage::OpCounts;
19
20/// One record in the JSONL event log. Tagged on `type` so the
21/// stream is forward-extensible: adding a new variant is a
22/// no-op for older readers (they ignore unknown types per
23/// SRD-44a §"Reader behaviour").
24#[derive(Debug, Clone, Serialize, Deserialize)]
25#[serde(tag = "type", rename_all = "snake_case")]
26pub enum CheckpointData {
27    /// First line of every fresh invocation's section. Resume
28    /// increments `invocation` and writes a fresh `session_start`
29    /// to **continue** the same JSONL — no separate file rotation.
30    SessionStart {
31        /// RFC 3339 UTC timestamp of when the event was written.
32        at: String,
33        /// Format version. `1` until we ship a `2`.
34        version: u32,
35        /// Session id (matches `logs/<session>/`).
36        session: String,
37        /// RFC 3339 first-invocation start.
38        started_at: String,
39        /// 1-based invocation counter.
40        invocation: u32,
41    },
42
43    /// Written when the workload completes. Optional — its
44    /// absence means the invocation was interrupted, which the
45    /// resume planner uses to distinguish clean exit from crash.
46    SessionEnd {
47        at: String,
48        /// `"completed"` / `"errored"` / `"stopped"`.
49        outcome: String,
50        /// Top-level error message when `outcome == "errored"`.
51        #[serde(default, skip_serializing_if = "Option::is_none")]
52        error: Option<String>,
53    },
54
55    /// One per phase per invocation, written during pre-map.
56    /// Carries identity and eligibility flags so the reader can
57    /// build the planned phase index without reading the
58    /// workload YAML.
59    PhaseDeclared {
60        at: String,
61        identity: PhaseIdentity,
62        skip_eligible: bool,
63    },
64
65    /// Mark a declared phase as Running.
66    PhaseStarted { at: String, identity: PhaseIdentity },
67
68    /// Periodic op-count + cursor-state update for a Running
69    /// phase. Replaces the in-place mutation of
70    /// `PhaseEntry::op_counts` and `cursor_state` from SRD-44.
71    /// The reader keeps **only the most recent** progress
72    /// record per (identity, current invocation) when folding.
73    PhaseProgress {
74        at: String,
75        identity: PhaseIdentity,
76        op_counts: OpCounts,
77        /// Tier-2 opaque snapshot from the active source factory.
78        #[serde(default, skip_serializing_if = "Option::is_none")]
79        cursor_state: Option<serde_json::Value>,
80    },
81
82    /// Phase reached the Completed state. Per SRD-44 §"Status
83    /// Completed is load-bearing", this is the call that makes
84    /// the phase eligible to skip on a future resume.
85    PhaseCompleted {
86        at: String,
87        identity: PhaseIdentity,
88        duration_secs: f64,
89        op_counts: OpCounts,
90    },
91
92    /// Phase failed terminally. The error message is preserved
93    /// for resume diagnostics.
94    PhaseFailed {
95        at: String,
96        identity: PhaseIdentity,
97        error: String,
98        #[serde(default, skip_serializing_if = "Option::is_none")]
99        op_counts: Option<OpCounts>,
100    },
101
102    /// Update on a previously-declared phase: bind the
103    /// program-canonical hash. Called once per phase the first
104    /// time it compiles. Folds into the existing entry's
105    /// identity so a future resume can detect program drift via
106    /// `PhaseIdentity::matches_full`.
107    PhaseHash {
108        at: String,
109        identity: PhaseIdentity,
110        /// 32-byte program hash, hex-encoded for human
111        /// readability of the on-disk records. SRD-107: the
112        /// BASE hash (chain below the session node + config
113        /// digest); param values ride in `params_consumed`.
114        hash_hex: String,
115        /// SRD-107 — the consumed-params map as canonical JSON
116        /// (`{"name":"<value sha256 hex>",…}`). Absent on
117        /// records written before the field existed.
118        #[serde(default, skip_serializing_if = "Option::is_none")]
119        params_consumed: Option<String>,
120    },
121
122    /// Reserved for SRD-44a Push 3. The reader already
123    /// understands these variants so the writer can emit them
124    /// in a future patch without a schema bump.
125    ScopeEnter {
126        at: String,
127        kind: String,
128        coords: std::collections::BTreeMap<String, serde_json::Value>,
129        path: Vec<std::collections::BTreeMap<String, serde_json::Value>>,
130    },
131
132    /// Reserved for SRD-44a Push 3.
133    ScopeExit {
134        at: String,
135        kind: String,
136        coords: std::collections::BTreeMap<String, serde_json::Value>,
137        path: Vec<std::collections::BTreeMap<String, serde_json::Value>>,
138        outcome: String,
139    },
140}
141
142impl CheckpointData {
143    /// RFC 3339 timestamp this event was tagged with. All
144    /// variants carry one; the helper avoids a match in the
145    /// fold loop.
146    pub fn at(&self) -> &str {
147        match self {
148            Self::SessionStart { at, .. }
149            | Self::SessionEnd { at, .. }
150            | Self::PhaseDeclared { at, .. }
151            | Self::PhaseStarted { at, .. }
152            | Self::PhaseProgress { at, .. }
153            | Self::PhaseCompleted { at, .. }
154            | Self::PhaseFailed { at, .. }
155            | Self::PhaseHash { at, .. }
156            | Self::ScopeEnter { at, .. }
157            | Self::ScopeExit { at, .. } => at,
158        }
159    }
160
161    /// The lifecycle kind-tag this durable record corresponds
162    /// to, tying an on-disk [`CheckpointData`] record back to
163    /// the in-memory [`crate::lifecycle::EventType`] that the
164    /// readout binder fires on. `None` for records that have no
165    /// lifecycle fire point — structural / metadata records
166    /// (a pre-map `PhaseDeclared` declaration, a `PhaseHash`
167    /// binding) are written to the log but never fire a readout
168    /// slot, so they have no kind-tag.
169    pub fn event_type(&self) -> Option<crate::lifecycle::EventType> {
170        use crate::lifecycle::EventType;
171        match self {
172            Self::SessionStart { .. } => Some(EventType::SessionStart),
173            Self::SessionEnd { .. } => Some(EventType::SessionEnd),
174            // Pre-map declaration record — no lifecycle fire point.
175            Self::PhaseDeclared { .. } => None,
176            Self::PhaseStarted { .. } => Some(EventType::PhaseStart),
177            Self::PhaseProgress { .. } => Some(EventType::Update),
178            // Both terminal states fire the phase-end slot
179            // (`phase_outcome` / `error_readout` bind there).
180            Self::PhaseCompleted { .. } | Self::PhaseFailed { .. } => Some(EventType::PhaseEnd),
181            // Metadata update — binds the program hash, no fire.
182            Self::PhaseHash { .. } => None,
183            Self::ScopeEnter { .. } => Some(EventType::ScopeStart),
184            Self::ScopeExit { .. } => Some(EventType::ScopeEnd),
185        }
186    }
187}
188
189/// Hex-encode a 32-byte program hash for storage in
190/// `PhaseHash` events. Lowercase, no separators.
191pub fn hash_to_hex(bytes: &[u8; 32]) -> String {
192    let mut s = String::with_capacity(64);
193    for b in bytes {
194        s.push_str(&format!("{b:02x}"));
195    }
196    s
197}
198
199/// Decode a 32-byte program hash from its hex form. Returns
200/// `None` on any malformed input — the reader logs a Warn and
201/// continues folding.
202pub fn hex_to_hash(hex: &str) -> Option<[u8; 32]> {
203    if hex.len() != 64 {
204        return None;
205    }
206    let mut out = [0u8; 32];
207    for (i, byte) in out.iter_mut().enumerate() {
208        let pair = &hex[i * 2..i * 2 + 2];
209        *byte = u8::from_str_radix(pair, 16).ok()?;
210    }
211    Some(out)
212}
213
214#[cfg(test)]
215mod tests {
216    use super::*;
217
218    #[test]
219    fn hash_round_trip() {
220        let h = [0xabu8; 32];
221        let hex = hash_to_hex(&h);
222        assert_eq!(hex.len(), 64);
223        assert_eq!(hex_to_hash(&hex), Some(h));
224    }
225
226    #[test]
227    fn hex_to_hash_rejects_malformed() {
228        assert!(hex_to_hash("").is_none());
229        assert!(hex_to_hash("ab").is_none()); // wrong length
230        assert!(hex_to_hash(&"z".repeat(64)).is_none()); // non-hex
231    }
232
233    #[test]
234    fn session_start_round_trip_via_serde() {
235        let e = CheckpointData::SessionStart {
236            at: "2026-05-07T12:00:00Z".into(),
237            version: 1,
238            session: "test".into(),
239            started_at: "2026-05-07T12:00:00Z".into(),
240            invocation: 1,
241        };
242        let line = serde_json::to_string(&e).unwrap();
243        assert!(line.contains("\"type\":\"session_start\""), "line: {line}");
244        let parsed: CheckpointData = serde_json::from_str(&line).unwrap();
245        match parsed {
246            CheckpointData::SessionStart { invocation, .. } => {
247                assert_eq!(invocation, 1);
248            }
249            _ => panic!("expected SessionStart"),
250        }
251    }
252
253    #[test]
254    fn unknown_type_fails_to_parse() {
255        // Confirms serde rejects unknown discriminators —
256        // the reader catches the Err and logs Debug per the
257        // forward-compat policy. (The "ignore unknown types"
258        // semantics live in the reader, not in serde.)
259        let line = r#"{"type":"future_event","at":"2026-05-07T12:00:00Z"}"#;
260        let r: Result<CheckpointData, _> = serde_json::from_str(line);
261        assert!(r.is_err(), "unknown type must error at parse");
262    }
263}