Skip to main content

nmbrs_runtime/checkpoint/
storage.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Checkpoint storage — fold-state types + JSONL event-log
5//! reader.
6//!
7//! See SRD-44a §"Reader behaviour" for the fold algorithm.
8//! `Checkpoint` is the in-memory representation produced by
9//! folding the on-disk `checkpoint.jsonl` event stream; each
10//! per-phase `PhaseEntry` records the entry's identity, status,
11//! and any cursor state for in-flight Tier 2 resume.
12
13use serde::{Deserialize, Serialize};
14use std::io::{BufRead, BufReader};
15use std::path::Path;
16
17use super::events::{CheckpointData, hex_to_hash};
18use super::identity::PhaseIdentity;
19
20/// In-memory checkpoint state — the fold of an append-only
21/// JSONL event stream at `logs/<session>/checkpoint.jsonl`.
22/// One per session. Per SRD-44a, the on-disk format is the
23/// event log; this struct is what consumers (resume planner,
24/// summary report) see after reading and folding.
25#[derive(Clone, Debug, Serialize, Deserialize)]
26pub struct Checkpoint {
27    /// File-format version. `1` until we ship a `2`. Each
28    /// resume invocation refuses to read a checkpoint whose
29    /// version it doesn't recognise — fail fast over silent
30    /// schema drift.
31    pub version: u32,
32    /// Session identifier — same string used in
33    /// `logs/<session>/`. The resume CLI's auto-detect picks
34    /// the most-recent session; explicit `--resume <id>`
35    /// names this directly.
36    pub session: String,
37    /// RFC 3339 timestamp of session start — first
38    /// invocation's `nmbrs run` start.
39    pub started_at: String,
40    /// RFC 3339 timestamp of this flush — updated on every
41    /// successful write.
42    pub checkpoint_at: String,
43    /// 1-based invocation counter. First `nmbrs run` is `1`;
44    /// each `--resume` increments by 1. Used in `session.log`
45    /// separator lines (`--- RESUMED <ts> [#N] ---`) and in
46    /// post-run summary diagnostics if needed.
47    pub invocation: u32,
48    /// One entry per pre-mapped phase. Order matches the
49    /// scenario tree's DFS — same order the post-run summary
50    /// uses, same order the resume planner walks.
51    pub phases: Vec<PhaseEntry>,
52}
53
54/// Per-phase entry in the checkpoint document.
55#[derive(Clone, Debug, Serialize, Deserialize)]
56pub struct PhaseEntry {
57    /// Phase identity tuple. The resume planner uses
58    /// `(yaml_path, coords)` as the structural match key and
59    /// (when present) the `phase_hash` as the sufficiency
60    /// check. See SRD-44 §"Phase identity".
61    #[serde(flatten)]
62    pub identity: PhaseIdentity,
63    /// Whether this phase is *eligible to skip* on resume
64    /// per its `checkpoint:` declaration. `true` means the
65    /// resume planner may classify the phase as Skip when
66    /// the saved status is Completed and identity matches.
67    /// `false` (operator declared `checkpoint: none` or no
68    /// declaration at all) means the phase always re-runs,
69    /// regardless of saved status.
70    pub skip_eligible: bool,
71    /// SRD-107 — the consumed-params map as canonical JSON
72    /// (`{"name":"<value sha256 hex>",…}`); the per-param leg of
73    /// resume skip validity. `None` on entries written before
74    /// the field existed (their base hash never matches current
75    /// formulas, so they re-run regardless).
76    #[serde(default, skip_serializing_if = "Option::is_none")]
77    pub params_consumed: Option<String>,
78    /// Lifecycle status at the time of the most recent flush.
79    pub status: PhaseStatus,
80    /// Wall-clock duration of the *successful* execution.
81    /// Set when status transitions to Completed; preserved on
82    /// subsequent flushes. None for in-flight (`Running`) or
83    /// terminally-failed phases.
84    #[serde(default)]
85    pub duration_secs: Option<f64>,
86    /// Op counts from the live activity. Captured on every
87    /// flush; useful for ETA and post-run summary in resumed
88    /// sessions.
89    #[serde(default)]
90    pub op_counts: Option<OpCounts>,
91    /// Tier 2 only. Opaque cursor-state snapshot from the
92    /// active source factory, captured on each flush while
93    /// the phase is `Running`. Resume planner restores this
94    /// to a freshly-constructed cursor source so the phase
95    /// continues from where it left off.
96    #[serde(default)]
97    pub cursor_state: Option<serde_json::Value>,
98    /// Per-phase error message recorded when the phase
99    /// transitioned to `Failed` — preserved across flushes
100    /// so resume diagnostics can reference the original
101    /// failure mode.
102    #[serde(default)]
103    pub error: Option<String>,
104}
105
106/// Lifecycle status as recorded in the checkpoint file.
107/// Mirrors `crate::scene_tree::PhaseStatus` semantically; kept
108/// as a separate type so the on-disk vocabulary doesn't
109/// coupled-evolve with scene-tree internal status (e.g. if we
110/// ever add a transient state for a runtime invariant that
111/// has no on-disk meaning, the storage type stays clean).
112#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
113#[serde(rename_all = "lowercase")]
114pub enum PhaseStatus {
115    Pending,
116    Running,
117    Completed,
118    Failed,
119}
120
121impl From<crate::scene_tree::PhaseStatus> for PhaseStatus {
122    fn from(s: crate::scene_tree::PhaseStatus) -> Self {
123        match s {
124            crate::scene_tree::PhaseStatus::Pending => Self::Pending,
125            crate::scene_tree::PhaseStatus::Running => Self::Running,
126            crate::scene_tree::PhaseStatus::Completed => Self::Completed,
127            crate::scene_tree::PhaseStatus::Failed(_) => Self::Failed,
128        }
129    }
130}
131
132/// Op-execution counts captured on every flush.
133#[derive(Clone, Debug, Default, Serialize, Deserialize)]
134pub struct OpCounts {
135    pub started: u64,
136    pub finished: u64,
137    pub errors: u64,
138}
139
140/// Format the current wall-clock time as an RFC 3339 UTC
141/// timestamp, e.g. `2026-01-01T00:00:00Z`. Used by the writer
142/// for `started_at` / `checkpoint_at` and by the runner when
143/// stamping a fresh session. Local implementation to avoid
144/// dragging chrono in for one call site.
145pub fn now_rfc3339() -> String {
146    let dur = std::time::SystemTime::now()
147        .duration_since(std::time::UNIX_EPOCH)
148        .unwrap_or_default();
149    let secs = dur.as_secs();
150    let days = secs / 86400;
151    let time_of_day = secs % 86400;
152    let hours = time_of_day / 3600;
153    let minutes = (time_of_day % 3600) / 60;
154    let seconds = time_of_day % 60;
155    let (year, month, day) = days_to_ymd(days);
156    format!("{year:04}-{month:02}-{day:02}T{hours:02}:{minutes:02}:{seconds:02}Z")
157}
158
159fn days_to_ymd(days: u64) -> (u64, u64, u64) {
160    let z = days + 719468;
161    let era = z / 146097;
162    let doe = z - era * 146097;
163    let yoe = (doe - doe / 1460 + doe / 36524 - doe / 146096) / 365;
164    let y = yoe + era * 400;
165    let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
166    let mp = (5 * doy + 2) / 153;
167    let d = doy - (153 * mp + 2) / 5 + 1;
168    let m = if mp < 10 { mp + 3 } else { mp - 9 };
169    let y = if m <= 2 { y + 1 } else { y };
170    (y, m, d)
171}
172
173/// Stream events from the JSONL log at `path`. Returns an
174/// iterator over `CheckpointData` records; malformed lines
175/// surface as `Err` items so the caller decides whether to
176/// stop or continue. Truncated-tail recovery is handled by
177/// [`read`]'s fold; this function is the lower-level building
178/// block for diagnostics tools that want raw event streams.
179pub fn iter_events(path: &Path) -> Result<Option<EventIter>, String> {
180    match std::fs::File::open(path) {
181        Ok(f) => {
182            let reader = BufReader::new(f);
183            Ok(Some(EventIter {
184                lines: reader.lines(),
185                path: path.to_path_buf(),
186            }))
187        }
188        Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
189        Err(e) => Err(format!("read checkpoint log {}: {e}", path.display())),
190    }
191}
192
193/// Iterator over [`CheckpointData`] records from a JSONL log.
194pub struct EventIter {
195    lines: std::io::Lines<BufReader<std::fs::File>>,
196    path: std::path::PathBuf,
197}
198
199impl Iterator for EventIter {
200    type Item = Result<CheckpointData, String>;
201
202    fn next(&mut self) -> Option<Self::Item> {
203        loop {
204            let line = match self.lines.next()? {
205                Ok(l) => l,
206                Err(e) => return Some(Err(format!("read line from {}: {e}", self.path.display()))),
207            };
208            if line.trim().is_empty() {
209                continue;
210            }
211            return Some(
212                serde_json::from_str(&line)
213                    .map_err(|e| format!("parse line in {}: {e}", self.path.display())),
214            );
215        }
216    }
217}
218
219/// Read and fold every event in `path` into a [`Checkpoint`].
220/// Returns `Ok(None)` when the file doesn't exist (fresh
221/// session); returns `Err` only on hard parse failures
222/// mid-stream. A truncated last line is recovered with a Warn
223/// per SRD-44a §"Truncated-tail recovery".
224pub fn read(path: &Path) -> Result<Option<Checkpoint>, String> {
225    use std::collections::HashMap;
226
227    let raw = match std::fs::read(path) {
228        Ok(b) => b,
229        Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
230        Err(e) => return Err(format!("read checkpoint log {}: {e}", path.display())),
231    };
232
233    // Trim a partial trailing line (one that doesn't end in
234    // `\n`) — this is the truncated-tail recovery. Anything
235    // before the last `\n` is a complete record per the
236    // append-mode write guarantee.
237    let cutoff = raw
238        .iter()
239        .rposition(|&b| b == b'\n')
240        .map(|i| i + 1)
241        .unwrap_or(0);
242    if cutoff < raw.len() {
243        eprintln!(
244            "warning: checkpoint {}: truncated tail (last {} bytes lacked newline), dropping",
245            path.display(),
246            raw.len() - cutoff,
247        );
248    }
249    let body = &raw[..cutoff];
250    let body_str = std::str::from_utf8(body)
251        .map_err(|e| format!("checkpoint log {}: invalid UTF-8: {e}", path.display()))?;
252
253    // First record MUST be a session_start per SRD-44a.
254    let mut lines = body_str.lines().filter(|l| !l.trim().is_empty());
255    let first_line = match lines.next() {
256        Some(l) => l,
257        None => return Ok(None), // empty log — treat as fresh
258    };
259    let first_event: CheckpointData = serde_json::from_str(first_line).map_err(|e| {
260        format!(
261            "checkpoint log {}: malformed first record: {e}",
262            path.display()
263        )
264    })?;
265
266    let mut doc = match first_event {
267        CheckpointData::SessionStart {
268            version,
269            session,
270            started_at,
271            invocation,
272            at,
273            ..
274        } => {
275            if version != 1 {
276                return Err(format!(
277                    "checkpoint {}: unsupported version {version} (this build supports v1)",
278                    path.display(),
279                ));
280            }
281            Checkpoint {
282                version,
283                session,
284                started_at,
285                checkpoint_at: at,
286                invocation,
287                phases: Vec::new(),
288            }
289        }
290        other => {
291            return Err(format!(
292                "checkpoint {}: first record must be session_start, got {:?}",
293                path.display(),
294                discriminator(&other),
295            ));
296        }
297    };
298
299    let mut index: HashMap<String, usize> = HashMap::new();
300
301    for line in lines {
302        let event: CheckpointData = match serde_json::from_str(line) {
303            Ok(e) => e,
304            Err(e) => {
305                // Per SRD-44a forward-compat: unknown event
306                // types (which serde rejects since we use
307                // tagged enum) are downgraded to a Debug log.
308                // Same fate for malformed mid-file records —
309                // they're a strong signal of corruption but
310                // not load-bearing for fold correctness as
311                // long as we keep going. Switching to Warn for
312                // visibility.
313                eprintln!(
314                    "warning: checkpoint {}: ignoring unparseable line: {e}",
315                    path.display(),
316                );
317                continue;
318            }
319        };
320        apply_event(&mut doc, &mut index, event);
321    }
322
323    Ok(Some(doc))
324}
325
326fn discriminator(e: &CheckpointData) -> &'static str {
327    match e {
328        CheckpointData::SessionStart { .. } => "session_start",
329        CheckpointData::SessionEnd { .. } => "session_end",
330        CheckpointData::PhaseDeclared { .. } => "phase_declared",
331        CheckpointData::PhaseStarted { .. } => "phase_started",
332        CheckpointData::PhaseProgress { .. } => "phase_progress",
333        CheckpointData::PhaseCompleted { .. } => "phase_completed",
334        CheckpointData::PhaseFailed { .. } => "phase_failed",
335        CheckpointData::PhaseHash { .. } => "phase_hash",
336        CheckpointData::ScopeEnter { .. } => "scope_enter",
337        CheckpointData::ScopeExit { .. } => "scope_exit",
338    }
339}
340
341fn apply_event(
342    doc: &mut Checkpoint,
343    index: &mut std::collections::HashMap<String, usize>,
344    event: CheckpointData,
345) {
346    match event {
347        CheckpointData::SessionStart {
348            invocation,
349            at,
350            started_at,
351            session,
352            ..
353        } => {
354            // Resume continuation: bump the invocation and
355            // refresh the per-flush timestamps. The phase list
356            // built so far stays as-is (fold semantics).
357            doc.invocation = invocation;
358            doc.checkpoint_at = at;
359            doc.started_at = started_at;
360            doc.session = session;
361        }
362        CheckpointData::SessionEnd { at, .. } => {
363            doc.checkpoint_at = at;
364        }
365        CheckpointData::PhaseDeclared {
366            at,
367            identity,
368            skip_eligible,
369        } => {
370            let key = super::writer::identity_key(&identity);
371            if let std::collections::hash_map::Entry::Vacant(e) = index.entry(key) {
372                doc.phases.push(PhaseEntry {
373                    identity,
374                    skip_eligible,
375                    params_consumed: None,
376                    status: PhaseStatus::Pending,
377                    duration_secs: None,
378                    op_counts: None,
379                    cursor_state: None,
380                    error: None,
381                });
382                e.insert(doc.phases.len() - 1);
383            }
384            doc.checkpoint_at = at;
385        }
386        CheckpointData::PhaseStarted { at, identity } => {
387            if let Some(entry) = lookup_mut(doc, index, &identity) {
388                entry.status = PhaseStatus::Running;
389                entry.error = None;
390            }
391            doc.checkpoint_at = at;
392        }
393        CheckpointData::PhaseProgress {
394            at,
395            identity,
396            op_counts,
397            cursor_state,
398        } => {
399            if let Some(entry) = lookup_mut(doc, index, &identity) {
400                entry.op_counts = Some(op_counts);
401                if cursor_state.is_some() {
402                    entry.cursor_state = cursor_state;
403                }
404            }
405            doc.checkpoint_at = at;
406        }
407        CheckpointData::PhaseCompleted {
408            at,
409            identity,
410            duration_secs,
411            op_counts,
412        } => {
413            if let Some(entry) = lookup_mut(doc, index, &identity) {
414                entry.status = PhaseStatus::Completed;
415                entry.duration_secs = Some(duration_secs);
416                entry.op_counts = Some(op_counts);
417                entry.cursor_state = None;
418                entry.error = None;
419            }
420            doc.checkpoint_at = at;
421        }
422        CheckpointData::PhaseFailed {
423            at,
424            identity,
425            error,
426            op_counts,
427        } => {
428            if let Some(entry) = lookup_mut(doc, index, &identity) {
429                entry.status = PhaseStatus::Failed;
430                entry.error = Some(error);
431                if let Some(c) = op_counts {
432                    entry.op_counts = Some(c);
433                }
434                entry.cursor_state = None;
435            }
436            doc.checkpoint_at = at;
437        }
438        CheckpointData::PhaseHash {
439            at,
440            identity,
441            hash_hex,
442            params_consumed,
443        } => {
444            if let Some(entry) = lookup_mut(doc, index, &identity)
445                && let Some(h) = hex_to_hash(&hash_hex)
446            {
447                entry.identity.phase_hash = Some(h);
448                entry.params_consumed = params_consumed;
449            }
450            doc.checkpoint_at = at;
451        }
452        CheckpointData::ScopeEnter { at, .. } | CheckpointData::ScopeExit { at, .. } => {
453            // Push 3 territory — fold to nothing today; the
454            // event lives on disk for forensic replay.
455            doc.checkpoint_at = at;
456        }
457    }
458}
459
460fn lookup_mut<'a>(
461    doc: &'a mut Checkpoint,
462    index: &std::collections::HashMap<String, usize>,
463    identity: &PhaseIdentity,
464) -> Option<&'a mut PhaseEntry> {
465    let key = super::writer::identity_key(identity);
466    let idx = *index.get(&key)?;
467    Some(&mut doc.phases[idx])
468}
469
470#[cfg(test)]
471mod tests {
472    use super::*;
473    use crate::checkpoint::CheckpointWriter;
474    use crate::checkpoint::PathSegment;
475    use crate::checkpoint::PhaseIdentity;
476
477    fn ident(name: &str) -> PhaseIdentity {
478        PhaseIdentity {
479            yaml_path: vec![
480                PathSegment::Scenario("test".into()),
481                PathSegment::Phase(name.into()),
482            ],
483            coords: String::new(),
484            phase_hash: None,
485        }
486    }
487
488    #[test]
489    fn read_missing_file_yields_none() {
490        let dir = tempdir();
491        let path = dir.join("nonexistent.jsonl");
492        let result = read(&path).expect("read should not error on missing file");
493        assert!(result.is_none());
494    }
495
496    #[test]
497    fn read_empty_file_yields_none() {
498        let dir = tempdir();
499        let path = dir.join("empty.jsonl");
500        std::fs::write(&path, "").expect("write");
501        let result = read(&path).expect("read should not error on empty file");
502        assert!(result.is_none(), "empty log = fresh session");
503    }
504
505    #[test]
506    fn fold_full_lifecycle_matches_in_memory_snapshot() {
507        // Write events via the writer, then read them back via
508        // the fold algorithm; the two views must agree on every
509        // field a resume planner observes.
510        let dir = tempdir();
511        let path = dir.join("checkpoint.jsonl");
512        let snap_in_memory = {
513            let w = CheckpointWriter::new(
514                path.clone(),
515                "sess".into(),
516                "2026-01-01T00:00:00Z".into(),
517                1,
518            );
519            let id1 = ident("schema");
520            let id2 = ident("rampup");
521            w.declare_phase(id1.clone(), true);
522            w.declare_phase(id2.clone(), true);
523            w.phase_started(&id1);
524            w.phase_completed(&id1, 1.5);
525            w.phase_started(&id2);
526            w.update_op_counts(
527                &id2,
528                OpCounts {
529                    started: 100,
530                    finished: 99,
531                    errors: 1,
532                },
533            );
534            w.flush().expect("flush");
535            w.snapshot()
536        };
537        let folded = read(&path).expect("read").expect("present");
538        assert_eq!(folded.session, snap_in_memory.session);
539        assert_eq!(folded.invocation, snap_in_memory.invocation);
540        assert_eq!(folded.phases.len(), snap_in_memory.phases.len());
541        for (i, phase) in folded.phases.iter().enumerate() {
542            assert_eq!(
543                phase.status, snap_in_memory.phases[i].status,
544                "status mismatch on phase {i}"
545            );
546            assert_eq!(phase.skip_eligible, snap_in_memory.phases[i].skip_eligible);
547            assert_eq!(phase.duration_secs, snap_in_memory.phases[i].duration_secs);
548            assert_eq!(
549                phase.op_counts.as_ref().map(|c| c.started),
550                snap_in_memory.phases[i]
551                    .op_counts
552                    .as_ref()
553                    .map(|c| c.started)
554            );
555        }
556    }
557
558    #[test]
559    fn truncated_tail_is_recovered() {
560        let dir = tempdir();
561        let path = dir.join("checkpoint.jsonl");
562        {
563            let w =
564                CheckpointWriter::new(path.clone(), "s".into(), "2026-01-01T00:00:00Z".into(), 1);
565            w.declare_phase(ident("p"), true);
566            w.flush().expect("flush");
567        }
568        // Append a partial line (no trailing newline) — simulates
569        // a crash mid-write.
570        use std::io::Write;
571        let mut f = std::fs::OpenOptions::new()
572            .append(true)
573            .open(&path)
574            .unwrap();
575        f.write_all(b"{\"type\":\"phase_started\",\"at\":\"2026-")
576            .unwrap();
577        drop(f);
578
579        let folded = read(&path).expect("read should recover").expect("present");
580        // Truncated tail dropped, prior records intact.
581        assert_eq!(folded.phases.len(), 1);
582    }
583
584    #[test]
585    fn first_record_must_be_session_start() {
586        let dir = tempdir();
587        let path = dir.join("checkpoint.jsonl");
588        // A valid `phase_started` record shaped correctly but
589        // appearing as the first line — the reader rejects.
590        let body = r#"{"type":"phase_started","at":"x","identity":{"yaml_path":[],"coords":""}}"#;
591        std::fs::write(&path, format!("{body}\n")).expect("write");
592        let err = read(&path).expect_err("first-record check must error");
593        assert!(
594            err.contains("first record must be session_start"),
595            "got: {err}"
596        );
597    }
598
599    #[test]
600    fn unsupported_version_in_session_start_errors() {
601        let dir = tempdir();
602        let path = dir.join("checkpoint.jsonl");
603        let body = r#"{"type":"session_start","at":"t","version":99,"session":"x","started_at":"t","invocation":1}"#;
604        std::fs::write(&path, format!("{body}\n")).expect("write");
605        let err = read(&path).expect_err("expected version-mismatch error");
606        assert!(err.contains("version 99"), "got: {err}");
607    }
608
609    fn tempdir() -> std::path::PathBuf {
610        let d = std::env::temp_dir().join(format!("nmbrs-checkpoint-test-{}", rand_suffix()));
611        std::fs::create_dir_all(&d).unwrap();
612        d
613    }
614
615    fn rand_suffix() -> String {
616        crate::scratch_suffix()
617    }
618}