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}