Skip to main content

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}