pub mod events;
pub mod identity;
pub mod params_scope;
pub mod resume;
pub mod storage;
pub mod writer;
pub use events::CheckpointData;
pub use identity::{PathSegment, PhaseIdentity};
pub use resume::{ResumeAction, ResumePlan};
pub use storage::{Checkpoint, OpCounts, PhaseEntry, PhaseStatus};
pub use writer::CheckpointWriter;
pub fn declare_scene_tree_phases(
writer: &CheckpointWriter,
tree: &crate::scene_tree::SceneTree,
phases: &std::collections::HashMap<String, nmbrs_workload::model::WorkloadPhase>,
) {
for node in tree.dfs_phases() {
let identity = PhaseIdentity {
yaml_path: node.yaml_path.clone(),
coords: node.labels.clone(),
phase_hash: None,
};
let skip_eligible = phases
.get(&node.name)
.and_then(|p| p.checkpoint.as_ref())
.map(|c| c.idempotent)
.unwrap_or(false);
writer.declare_phase(identity, skip_eligible);
}
}
pub fn scene_tree_resume_candidates(
tree: &crate::scene_tree::SceneTree,
scope_tree: &crate::scope_tree::ScopeTree,
phases: &std::collections::HashMap<String, nmbrs_workload::model::WorkloadPhase>,
) -> Vec<(PhaseIdentity, bool)> {
tree.dfs_phases()
.map(|node| {
let phase = phases.get(&node.name);
let chain = ancestor_chain_hash(scope_tree, &node.name);
let config = phase.map(phase_config_hash).unwrap_or([0u8; 32]);
let identity = PhaseIdentity {
yaml_path: node.yaml_path.clone(),
coords: node.labels.clone(),
phase_hash: Some(compose_phase_hash(chain, config)),
};
let idempotent = phase
.and_then(|p| p.checkpoint.as_ref())
.map(|c| c.idempotent)
.unwrap_or(false);
(identity, idempotent)
})
.collect()
}
pub(crate) fn compose_phase_hash(chain: Option<[u8; 32]>, config: [u8; 32]) -> [u8; 32] {
use sha2::{Digest, Sha256};
let mut h = Sha256::new();
h.update(b"nmbrs-phase-identity-v2\n");
match chain {
Some(c) => {
h.update(b"chain:");
h.update(c);
}
None => h.update(b"chain:none"),
}
h.update(b"config:");
h.update(config);
h.finalize().into()
}
pub(crate) fn phase_config_canonical_text(phase: &nmbrs_workload::model::WorkloadPhase) -> String {
let mut value = serde_json::to_value(phase).unwrap_or(serde_json::Value::Null);
sort_json_keys(&mut value);
serde_json::to_string(&value).unwrap_or_default()
}
pub(crate) fn phase_config_hash(phase: &nmbrs_workload::model::WorkloadPhase) -> [u8; 32] {
config_text_hash(&phase_config_canonical_text(phase))
}
pub(crate) fn config_text_hash(text: &str) -> [u8; 32] {
use sha2::{Digest, Sha256};
let mut h = Sha256::new();
h.update(b"nmbrs-phase-config-v1\n");
h.update(text.as_bytes());
h.finalize().into()
}
fn sort_json_keys(v: &mut serde_json::Value) {
match v {
serde_json::Value::Object(map) => {
let mut entries: Vec<(String, serde_json::Value)> =
std::mem::take(map).into_iter().collect();
entries.sort_by(|a, b| a.0.cmp(&b.0));
for (_, val) in entries.iter_mut() {
sort_json_keys(val);
}
map.extend(entries);
}
serde_json::Value::Array(items) => {
for item in items {
sort_json_keys(item);
}
}
_ => {}
}
}
pub(crate) fn ancestor_chain_hash(
scope_tree: &crate::scope_tree::ScopeTree,
phase_name: &str,
) -> Option<[u8; 32]> {
let idx = scope_tree.phase_node_by_name(phase_name)?;
let (ancestors, _params_module) = scope_tree.ancestor_kernels_split(idx);
if ancestors.is_empty() {
return None;
}
let head = ancestors[0].program();
let tail: Vec<&polydat::kernel::PolydatProgram> = ancestors[1..]
.iter()
.map(|k| k.program().as_ref())
.collect();
Some(head.instance_hash(&tail))
}