nmbrs_runtime/checkpoint/mod.rs
1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Workload checkpointing — see SRD-44.
5//!
6//! Submodules:
7//!
8//! - [`identity`] — `PathSegment`, `PhaseIdentity`, and the
9//! per-phase canonical-program hash. Identity is per-phase
10//! (no workload-level identity tuple) and `(yaml_path,
11//! coords)` is necessary; the program hash is sufficiency.
12//! - [`storage`] — JSON file format + atomic-rename writer.
13//! - [`writer`] — `CheckpointWriter` actor: subscribes to
14//! phase-lifecycle events, flushes on the metrics-tick
15//! cadence with sqlite-fsync-then-checkpoint-fsync ordering.
16//! - [`resume`] — resume planner: loads checkpoint, classifies
17//! each freshly-pre-mapped phase per the resume protocol,
18//! produces a `ResumePlan` the executor consults before
19//! dispatch.
20
21pub mod events;
22pub mod identity;
23pub mod params_scope;
24pub mod resume;
25pub mod storage;
26pub mod writer;
27
28pub use events::CheckpointData;
29pub use identity::{PathSegment, PhaseIdentity};
30pub use resume::{ResumeAction, ResumePlan};
31pub use storage::{Checkpoint, OpCounts, PhaseEntry, PhaseStatus};
32pub use writer::CheckpointWriter;
33
34/// Declare every phase node in a freshly-pre-mapped scene tree
35/// to the writer. Called once at session bootstrap, immediately
36/// after [`crate::executor::pre_map_tree`] returns. Each phase
37/// gets a `Pending` entry with no hash; the runtime updates the
38/// hash via [`CheckpointWriter::update_phase_hash`] when the
39/// phase compiles.
40///
41/// The `phases` map carries each phase's `checkpoint:`
42/// declaration (parsed by `nmbrs-workload`); entries with
43/// `checkpoint: idempotent` set `skip_eligible = true`, all
44/// others set `false`.
45pub fn declare_scene_tree_phases(
46 writer: &CheckpointWriter,
47 tree: &crate::scene_tree::SceneTree,
48 phases: &std::collections::HashMap<String, nmbrs_workload::model::WorkloadPhase>,
49) {
50 for node in tree.dfs_phases() {
51 let identity = PhaseIdentity {
52 yaml_path: node.yaml_path.clone(),
53 coords: node.labels.clone(),
54 phase_hash: None,
55 };
56 let skip_eligible = phases
57 .get(&node.name)
58 .and_then(|p| p.checkpoint.as_ref())
59 .map(|c| c.idempotent)
60 .unwrap_or(false);
61 writer.declare_phase(identity, skip_eligible);
62 }
63}
64
65/// Build a list of `(identity, declared_idempotent)` pairs for
66/// every phase the scene tree will execute. The runner feeds
67/// this to [`ResumePlan::from_checkpoint`] when resuming, so the
68/// planner can classify each freshly-pre-mapped phase against
69/// the saved document.
70///
71/// Each candidate's `phase_hash` is [`compose_phase_hash`] over
72/// two pre-map-computable digests (SRD-106 D2 — this is THE
73/// skip-validity anchor, shared verbatim with the executor's
74/// stamped value so saved and fresh compare directly):
75///
76/// - the **ancestor-chain instance hash** — SHA-256 over the
77/// canonical_hash of every installed ancestor kernel, from
78/// the immediate parent scope up through the workload root
79/// AND the session-level workload-params module. Catches
80/// upstream binding edits and any param-value change (params
81/// live as const slots on the params module).
82/// - the **phase-config digest** ([`phase_config_hash`]) — a
83/// canonical serialization of the phase's full declared
84/// configuration: ops (statement templates included),
85/// bindings, cycles, concurrency, rate, stop conditions —
86/// every `WorkloadPhase` field. Catches edits the compiled
87/// program chain cannot see (an op's statement text, a
88/// cycle-count change).
89pub fn scene_tree_resume_candidates(
90 tree: &crate::scene_tree::SceneTree,
91 scope_tree: &crate::scope_tree::ScopeTree,
92 phases: &std::collections::HashMap<String, nmbrs_workload::model::WorkloadPhase>,
93) -> Vec<(PhaseIdentity, bool)> {
94 tree.dfs_phases()
95 .map(|node| {
96 let phase = phases.get(&node.name);
97 let chain = ancestor_chain_hash(scope_tree, &node.name);
98 let config = phase.map(phase_config_hash).unwrap_or([0u8; 32]);
99 let identity = PhaseIdentity {
100 yaml_path: node.yaml_path.clone(),
101 coords: node.labels.clone(),
102 phase_hash: Some(compose_phase_hash(chain, config)),
103 };
104 let idempotent = phase
105 .and_then(|p| p.checkpoint.as_ref())
106 .map(|c| c.idempotent)
107 .unwrap_or(false);
108 (identity, idempotent)
109 })
110 .collect()
111}
112
113/// Compose the two provenance digests into the one phase hash
114/// that every store and gate carries: the checkpoint document
115/// (SRD-44 resume classification), the persisted phase-outcome
116/// row (SRD-77 refine hash gate), and the resume planner's
117/// candidates all use this same formula, so a saved hash and a
118/// freshly computed one compare directly.
119pub(crate) fn compose_phase_hash(chain: Option<[u8; 32]>, config: [u8; 32]) -> [u8; 32] {
120 use sha2::{Digest, Sha256};
121 let mut h = Sha256::new();
122 h.update(b"nmbrs-phase-identity-v2\n");
123 match chain {
124 Some(c) => {
125 h.update(b"chain:");
126 h.update(c);
127 }
128 None => h.update(b"chain:none"),
129 }
130 h.update(b"config:");
131 h.update(config);
132 h.finalize().into()
133}
134
135/// Canonical serialization of a phase's full declared
136/// configuration — every `WorkloadPhase` field, ops and bindings
137/// included — through `serde_json::Value` with recursively sorted
138/// object keys so HashMap-backed fields serialize stably across
139/// processes (the workspace enables serde_json's `preserve_order`,
140/// so insertion order alone is NOT stable). One surface, two
141/// consumers: the config digest below and SRD-107's textual
142/// `{name}` interpolation scan.
143pub(crate) fn phase_config_canonical_text(phase: &nmbrs_workload::model::WorkloadPhase) -> String {
144 let mut value = serde_json::to_value(phase).unwrap_or(serde_json::Value::Null);
145 sort_json_keys(&mut value);
146 serde_json::to_string(&value).unwrap_or_default()
147}
148
149/// Canonical SHA-256 over [`phase_config_canonical_text`].
150pub(crate) fn phase_config_hash(phase: &nmbrs_workload::model::WorkloadPhase) -> [u8; 32] {
151 config_text_hash(&phase_config_canonical_text(phase))
152}
153
154/// Hash an already-canonicalized config text (callers that also
155/// feed the text to the SRD-107 scan avoid serializing twice).
156pub(crate) fn config_text_hash(text: &str) -> [u8; 32] {
157 use sha2::{Digest, Sha256};
158 let mut h = Sha256::new();
159 h.update(b"nmbrs-phase-config-v1\n");
160 h.update(text.as_bytes());
161 h.finalize().into()
162}
163
164fn sort_json_keys(v: &mut serde_json::Value) {
165 match v {
166 serde_json::Value::Object(map) => {
167 let mut entries: Vec<(String, serde_json::Value)> =
168 std::mem::take(map).into_iter().collect();
169 entries.sort_by(|a, b| a.0.cmp(&b.0));
170 for (_, val) in entries.iter_mut() {
171 sort_json_keys(val);
172 }
173 map.extend(entries);
174 }
175 serde_json::Value::Array(items) => {
176 for item in items {
177 sort_json_keys(item);
178 }
179 }
180 _ => {}
181 }
182}
183
184/// Compute a phase's ancestor-chain instance hash by looking up
185/// the scope-tree node and walking its installed ancestor
186/// kernels — immediate parent first, up through the workload
187/// root. The session node's kernel (the workload-params module)
188/// is EXCLUDED (SRD-107): param values participate per-phase via
189/// the consumed-params digest, not the chain, so an unrelated
190/// param change cannot flip this hash. Returns `None` if the
191/// scope tree has no installed kernels below the session
192/// (defensive — the workload root always has one in production).
193///
194/// Shared by the resume planner's candidates and the executor's
195/// stamped hash ([`compose_phase_hash`] composes it with the
196/// phase-config digest at both sites) — one formula, one walk.
197pub(crate) fn ancestor_chain_hash(
198 scope_tree: &crate::scope_tree::ScopeTree,
199 phase_name: &str,
200) -> Option<[u8; 32]> {
201 let idx = scope_tree.phase_node_by_name(phase_name)?;
202 let (ancestors, _params_module) = scope_tree.ancestor_kernels_split(idx);
203 if ancestors.is_empty() {
204 return None;
205 }
206 // The chain hash uses PolydatProgram::instance_hash with the
207 // first ancestor as the "self" anchor and the rest as
208 // ancestors-of-ancestor. The phase's OWN program is
209 // deliberately absent (it compiles lazily; its declared
210 // matter is covered by the phase-config digest instead), so
211 // this value is computable at pre-map time and identical at
212 // both compute sites.
213 let head = ancestors[0].program();
214 let tail: Vec<&polydat::kernel::PolydatProgram> = ancestors[1..]
215 .iter()
216 .map(|k| k.program().as_ref())
217 .collect();
218 Some(head.instance_hash(&tail))
219}