nmbrs-runtime 0.4.0

Workload execution runtime for nmbrs
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
// Copyright 2024-2026 Jonathan Shook
// SPDX-License-Identifier: Apache-2.0

//! Resume planner — classify each freshly-pre-mapped phase
//! against the saved checkpoint document and emit a
//! `ResumePlan` the executor consults before dispatch.
//!
//! Per SRD-44 §"Resume protocol", the planner walks the
//! pre-map phase list (already in DFS-of-the-scenario-tree
//! order) and for each phase asks:
//!
//! 1. Is there a structural match (same `yaml_path` and
//!    `coords`) in the saved doc? If not → `ReRun` (this is a
//!    new phase that wasn't in the last invocation).
//! 2. If matched and the saved status is `Completed`:
//!    - phase declared `checkpoint: idempotent` (or long form
//!      with `idempotent: true`) → `Skip` when hashes match,
//!      `IdentityMismatch` (and re-run) when they don't.
//!    - phase declared `checkpoint: none` (or absent) → `ReRun`
//!      regardless of saved status (operator opted out).
//! 3. If matched and the saved status is `Running`:
//!    - the phase has cursor state recorded → `CursorResume`
//!      with that opaque snapshot (Tier 2).
//!    - no cursor state → `ReRun` (Tier 1 inflight crash).
//! 4. If matched and the saved status is `Failed` →
//!    `ReRun` (the errors cascade decides whether to surface
//!    the failure again — see SRD-44 §"Error handling is
//!    invocation-agnostic"). The planner doesn't pre-decide.
//! 5. If matched and the saved status is `Pending` → `ReRun`
//!    (the previous invocation never started this phase).

use std::collections::HashMap;

use super::identity::PhaseIdentity;
use super::storage::{Checkpoint, PhaseStatus};

/// What the executor should do with one pre-mapped phase. The
/// planner emits one of these per phase keyed by identity-key
/// (see [`super::writer`] for the key shape).
#[derive(Clone, Debug)]
pub enum ResumeAction {
    /// Phase already completed successfully on a prior
    /// invocation, identity matches, and the operator declared
    /// it idempotent. Executor emits the post-run summary line
    /// for this phase as `[skipped]` and moves on; the writer
    /// keeps the saved entry intact.
    Skip,
    /// Phase was in-flight (`Running`) and recorded a cursor
    /// snapshot the source factory understands. Executor
    /// reconstructs the source via the factory, calls
    /// `restore_cursor(snapshot)` on it, and runs the remaining
    /// cycles per SRD-44 §"Tier 2".
    CursorResume { cursor_state: serde_json::Value },
    /// Default — the phase runs from scratch this invocation.
    /// Triggered by: no saved entry, declared `checkpoint:
    /// none`, saved status `Pending`/`Failed`, or `Running`
    /// without cursor state.
    ReRun,
    /// Saved entry exists, structural match holds, but the
    /// program hash differs (or some other identity mismatch).
    /// Executor re-runs the phase and the writer overwrites
    /// the stale entry. Carries a human-readable reason for
    /// the resume diagnostic banner.
    IdentityMismatch { reason: String },
}

/// The resume plan for one pre-mapped scenario tree. Maps each
/// phase's identity-key to its [`ResumeAction`]. A "fresh" plan
/// (no saved checkpoint) maps every phase to `ReRun`.
#[derive(Clone, Debug, Default)]
pub struct ResumePlan {
    actions: HashMap<String, ResumeAction>,
    /// `true` when at least one phase was loaded from a
    /// previously-flushed checkpoint document — i.e. this is a
    /// resume run, not a fresh start. Drives the resume banner
    /// in the runner header and the `--- RESUMED ---` separator
    /// in `session.log`.
    pub is_resume: bool,
}

impl ResumePlan {
    /// A plan with no saved state — every phase re-runs.
    pub fn fresh() -> Self {
        Self::default()
    }

    /// Build a plan from a saved checkpoint and a list of
    /// freshly-pre-mapped phase candidates. `candidates` is
    /// `(identity, declared_idempotent)`; the second entry
    /// reflects the workload's `checkpoint:` declaration for
    /// that phase (`true` for idempotent, `false` for
    /// none/absent).
    pub fn from_checkpoint(
        saved: &Checkpoint,
        candidates: &[(PhaseIdentity, bool)],
        current_params: &HashMap<String, String>,
    ) -> Self {
        let saved_index: HashMap<String, &super::storage::PhaseEntry> = saved
            .phases
            .iter()
            .map(|e| (identity_key(&e.identity), e))
            .collect();

        let mut actions = HashMap::with_capacity(candidates.len());
        for (cand, declared_idempotent) in candidates {
            let key = identity_key(cand);
            let action = match saved_index.get(&key) {
                None => ResumeAction::ReRun,
                Some(saved_entry) => {
                    classify(cand, saved_entry, *declared_idempotent, current_params)
                }
            };
            actions.insert(key, action);
        }
        Self {
            actions,
            is_resume: true,
        }
    }

    /// Look up the action for a given phase identity. Returns
    /// `ReRun` for unknown phases — the conservative default so
    /// pre-map drift can never cause a phase to silently skip.
    pub fn action_for(&self, identity: &PhaseIdentity) -> ResumeAction {
        let key = identity_key(identity);
        self.actions
            .get(&key)
            .cloned()
            .unwrap_or(ResumeAction::ReRun)
    }

    /// Number of phases the plan classifies as `Skip`.
    pub fn skip_count(&self) -> usize {
        self.actions
            .values()
            .filter(|a| matches!(a, ResumeAction::Skip))
            .count()
    }

    /// Number of phases the plan classifies as `CursorResume`.
    pub fn cursor_resume_count(&self) -> usize {
        self.actions
            .values()
            .filter(|a| matches!(a, ResumeAction::CursorResume { .. }))
            .count()
    }

    /// Number of phases the plan classifies as
    /// `IdentityMismatch`. Surfaced in the resume banner so
    /// operators can tell at a glance whether their YAML edit
    /// invalidated phases they thought were stable.
    pub fn mismatch_count(&self) -> usize {
        self.actions
            .values()
            .filter(|a| matches!(a, ResumeAction::IdentityMismatch { .. }))
            .count()
    }
}

fn classify(
    candidate: &PhaseIdentity,
    saved: &super::storage::PhaseEntry,
    declared_idempotent: bool,
    current_params: &HashMap<String, String>,
) -> ResumeAction {
    // The saved entry might have been written before the
    // operator changed `checkpoint:` for this phase. Honour the
    // *current* declaration: if the workload now says `none`,
    // re-run regardless of saved status.
    if !declared_idempotent {
        return ResumeAction::ReRun;
    }

    match saved.status {
        PhaseStatus::Completed => {
            // Structural match has been established by the
            // identity_key lookup; check sufficiency (base hash)
            // only when both sides carry one.
            if !candidate.matches_full(&saved.identity) {
                return ResumeAction::IdentityMismatch {
                    reason: format!(
                        "phase '{}': scope or phase config changed since \
                         this phase last ran (base hash differs)",
                        phase_label(&candidate.yaml_path),
                    ),
                };
            }
            // SRD-107 — the per-param leg: every param the saved
            // run consumed must digest to the same value NOW. The
            // blocker names the param. A base-matching entry
            // without the stored map is incomparable (mid-upgrade
            // oddity) → conservative re-run.
            if let Err(mismatch) =
                saved_params_still_valid(saved.params_consumed.as_deref(), current_params)
            {
                return ResumeAction::IdentityMismatch {
                    reason: format!(
                        "phase '{}': {mismatch} since this phase last ran",
                        phase_label(&candidate.yaml_path),
                    ),
                };
            }
            if !saved.skip_eligible {
                // Saved entry was written when the phase was
                // declared `none`; even though the current
                // declaration is `idempotent`, the saved
                // execution may not have been idempotent. Be
                // conservative and re-run.
                return ResumeAction::ReRun;
            }
            ResumeAction::Skip
        }
        PhaseStatus::Running => match &saved.cursor_state {
            Some(cs) => ResumeAction::CursorResume {
                cursor_state: cs.clone(),
            },
            None => ResumeAction::ReRun,
        },
        // Pending: never started in the prior invocation.
        // Failed: errors cascade decides whether to retry; from
        //   the planner's perspective, a Failed phase always
        //   re-runs (the cascade may then mark it Failed again
        //   or surface a different outcome).
        PhaseStatus::Pending | PhaseStatus::Failed => ResumeAction::ReRun,
    }
}

/// SRD-107 — evaluate a saved consumed-params map against the
/// CURRENT param values. `Ok(())` when every stored name digests
/// identically now; `Err(description)` naming the first changed
/// or missing param, or the comparability gap.
fn saved_params_still_valid(
    stored_json: Option<&str>,
    current_params: &HashMap<String, String>,
) -> Result<(), String> {
    let Some(json) = stored_json else {
        return Err("saved entry carries no consumed-params record".into());
    };
    let Ok(stored) = serde_json::from_str::<std::collections::BTreeMap<String, String>>(json)
    else {
        return Err("saved consumed-params record is unreadable".into());
    };
    for (name, stored_digest) in stored {
        let current = current_params
            .get(&name)
            .map(|v| crate::checkpoint::params_scope::value_digest(v));
        if current.as_deref() != Some(stored_digest.as_str()) {
            return Err(format!("param '{name}' changed"));
        }
    }
    Ok(())
}

/// Reuse the writer's identity-key shape so both sides agree on
/// equality semantics.
fn identity_key(identity: &PhaseIdentity) -> String {
    let path_json = serde_json::to_string(&identity.yaml_path).unwrap_or_else(|_| String::new());
    format!("{path_json}\x1f{}", identity.coords)
}

fn phase_label(yaml_path: &[super::identity::PathSegment]) -> String {
    use super::identity::PathSegment;
    yaml_path
        .iter()
        .filter_map(|seg| match seg {
            PathSegment::Phase(n) => Some(n.clone()),
            _ => None,
        })
        .next_back()
        .unwrap_or_else(|| "<unknown>".to_string())
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::checkpoint::{
        PathSegment,
        storage::{OpCounts, PhaseEntry},
    };

    fn ident_with_hash(name: &str, coords: &str, hash: Option<[u8; 32]>) -> PhaseIdentity {
        PhaseIdentity {
            yaml_path: vec![
                PathSegment::Scenario("s".into()),
                PathSegment::Phase(name.into()),
            ],
            coords: coords.into(),
            phase_hash: hash,
        }
    }

    fn entry(identity: PhaseIdentity, status: PhaseStatus, skip_eligible: bool) -> PhaseEntry {
        PhaseEntry {
            identity,
            skip_eligible,
            // Comparable empty consumed-set: these fixtures pin the
            // status/hash legs; the params leg has its own tests.
            params_consumed: Some("{}".into()),
            status,
            duration_secs: Some(1.0),
            op_counts: Some(OpCounts::default()),
            cursor_state: None,
            error: None,
        }
    }

    fn checkpoint_with(phases: Vec<PhaseEntry>) -> Checkpoint {
        Checkpoint {
            version: 1,
            session: "s".into(),
            started_at: "t".into(),
            checkpoint_at: "t".into(),
            invocation: 1,
            phases,
        }
    }

    #[test]
    fn fresh_plan_rerun_for_everything() {
        let plan = ResumePlan::fresh();
        assert!(matches!(
            plan.action_for(&ident_with_hash("p", "", None)),
            ResumeAction::ReRun
        ));
        assert!(!plan.is_resume);
    }

    /// SRD-107 — the per-param leg: same base hash, but a param
    /// the saved run consumed now digests differently → the
    /// mismatch names the param.
    #[test]
    fn consumed_param_change_invalidates_and_names_the_param() {
        let h = [0xab; 32];
        let id = ident_with_hash("load", "", Some(h));
        let mut e = entry(id.clone(), PhaseStatus::Completed, true);
        e.params_consumed = Some(format!(
            r#"{{"dataset":"{}"}}"#,
            crate::checkpoint::params_scope::value_digest("example"),
        ));
        let saved = checkpoint_with(vec![e]);

        let mut params = HashMap::new();
        params.insert("dataset".to_string(), "sift10m".to_string());
        let plan = ResumePlan::from_checkpoint(&saved, &[(id.clone(), true)], &params);
        match plan.action_for(&id) {
            ResumeAction::IdentityMismatch { reason } => {
                assert!(
                    reason.contains("param 'dataset' changed"),
                    "reason must name the param: {reason}"
                );
            }
            other => panic!("expected IdentityMismatch, got {other:?}"),
        }
    }

    /// SRD-107 — unrelated params may change freely: only the
    /// STORED names are checked, so a phase that consumed nothing
    /// (or different params) still skips.
    #[test]
    fn unrelated_param_change_still_skips() {
        let h = [0xab; 32];
        let id = ident_with_hash("load", "", Some(h));
        let mut e = entry(id.clone(), PhaseStatus::Completed, true);
        e.params_consumed = Some(format!(
            r#"{{"dataset":"{}"}}"#,
            crate::checkpoint::params_scope::value_digest("example"),
        ));
        let saved = checkpoint_with(vec![e]);

        let mut params = HashMap::new();
        params.insert("dataset".to_string(), "example".to_string());
        params.insert("suite_k".to_string(), "100".to_string());
        let plan = ResumePlan::from_checkpoint(&saved, &[(id.clone(), true)], &params);
        assert!(
            matches!(plan.action_for(&id), ResumeAction::Skip),
            "an unconsumed param's change must not invalidate"
        );
    }

    #[test]
    fn completed_idempotent_with_matching_hash_skips() {
        let h = [0xab; 32];
        let id = ident_with_hash("schema", "", Some(h));
        let saved = checkpoint_with(vec![entry(id.clone(), PhaseStatus::Completed, true)]);
        let plan = ResumePlan::from_checkpoint(&saved, &[(id.clone(), true)], &HashMap::new());
        assert!(matches!(plan.action_for(&id), ResumeAction::Skip));
        assert_eq!(plan.skip_count(), 1);
    }

    #[test]
    fn completed_with_hash_mismatch_invalidates() {
        let h_old = [0x01; 32];
        let h_new = [0x02; 32];
        let id_saved = ident_with_hash("schema", "", Some(h_old));
        let id_now = ident_with_hash("schema", "", Some(h_new));
        let saved = checkpoint_with(vec![entry(id_saved.clone(), PhaseStatus::Completed, true)]);
        let plan = ResumePlan::from_checkpoint(&saved, &[(id_now.clone(), true)], &HashMap::new());
        match plan.action_for(&id_now) {
            ResumeAction::IdentityMismatch { reason } => {
                assert!(reason.contains("schema"), "reason: {reason}");
            }
            other => panic!("expected IdentityMismatch, got {other:?}"),
        }
        assert_eq!(plan.mismatch_count(), 1);
    }

    #[test]
    fn declared_none_always_reruns_even_if_completed() {
        let id = ident_with_hash("schema", "", None);
        let saved = checkpoint_with(vec![entry(id.clone(), PhaseStatus::Completed, false)]);
        let plan = ResumePlan::from_checkpoint(&saved, &[(id.clone(), false)], &HashMap::new());
        assert!(matches!(plan.action_for(&id), ResumeAction::ReRun));
    }

    #[test]
    fn running_with_cursor_state_yields_cursor_resume() {
        let id = ident_with_hash("rampup", "", None);
        let mut e = entry(id.clone(), PhaseStatus::Running, true);
        e.cursor_state = Some(serde_json::json!({"next_cycle": 12345}));
        let saved = checkpoint_with(vec![e]);
        let plan = ResumePlan::from_checkpoint(&saved, &[(id.clone(), true)], &HashMap::new());
        match plan.action_for(&id) {
            ResumeAction::CursorResume { cursor_state } => {
                assert_eq!(cursor_state["next_cycle"], 12345);
            }
            other => panic!("expected CursorResume, got {other:?}"),
        }
        assert_eq!(plan.cursor_resume_count(), 1);
    }

    #[test]
    fn running_without_cursor_state_reruns() {
        let id = ident_with_hash("rampup", "", None);
        let saved = checkpoint_with(vec![entry(id.clone(), PhaseStatus::Running, true)]);
        let plan = ResumePlan::from_checkpoint(&saved, &[(id.clone(), true)], &HashMap::new());
        assert!(matches!(plan.action_for(&id), ResumeAction::ReRun));
    }

    #[test]
    fn unknown_candidate_reruns() {
        let saved = checkpoint_with(vec![]);
        let id = ident_with_hash("brand_new", "", None);
        let plan = ResumePlan::from_checkpoint(&saved, &[(id.clone(), true)], &HashMap::new());
        assert!(matches!(plan.action_for(&id), ResumeAction::ReRun));
    }

    #[test]
    fn failed_phase_reruns() {
        let id = ident_with_hash("flaky", "", None);
        let mut e = entry(id.clone(), PhaseStatus::Failed, true);
        e.error = Some("boom".into());
        let saved = checkpoint_with(vec![e]);
        let plan = ResumePlan::from_checkpoint(&saved, &[(id.clone(), true)], &HashMap::new());
        assert!(matches!(plan.action_for(&id), ResumeAction::ReRun));
    }
}