Skip to main content

nmbrs_runtime/checkpoint/
writer.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! `CheckpointWriter` — append-only event-log owner for the
5//! per-session checkpoint document. SRD-44a §"Writer behaviour".
6//!
7//! Per SRD-44a, every state-changing observation is one event
8//! line written to `logs/<session>/checkpoint.jsonl`. The writer:
9//!
10//! - is created at session bootstrap (or restored from a saved
11//!   document on resume) — emits a `session_start` event,
12//! - has phases *declared* into it during pre-map — one
13//!   `phase_declared` event each,
14//! - receives phase-lifecycle calls (`phase_started`,
15//!   `phase_completed`, `phase_failed`) — one event per call,
16//! - receives op-count and cursor-state updates from the
17//!   metrics tick callback — one `phase_progress` event per
18//!   tick.
19//!
20//! No whole-document rewrite. The file is opened in append mode
21//! (`O_APPEND`); each event is a single `\n`-terminated line.
22//! In-memory state is the same `Checkpoint` document as before
23//! so `snapshot()` and `resume_hint()` keep working without
24//! re-folding from disk.
25//!
26//! The writer is `Send + Sync` (interior mutex) so the executor,
27//! the tick callback, and the cursor-state collector can all
28//! hold one `Arc<CheckpointWriter>`.
29
30use std::collections::{BTreeMap, HashMap};
31use std::fs::{File, OpenOptions};
32use std::io::Write;
33use std::path::PathBuf;
34use std::sync::Mutex;
35
36use super::events::{CheckpointData, hash_to_hex};
37use super::identity::PhaseIdentity;
38use super::storage::{Checkpoint, OpCounts, PhaseEntry, PhaseStatus, now_rfc3339};
39
40/// On-disk checkpoint version this build emits / accepts. Bump
41/// only when an incompatible schema change ships; resume against
42/// a different version is rejected at read time.
43pub const CHECKPOINT_VERSION: u32 = 1;
44
45/// Writer-side handle to the per-session checkpoint event log.
46/// One instance per session; held as `Arc<CheckpointWriter>` and
47/// shared between the executor (lifecycle calls) and the metrics
48/// tick (count + cursor + flush calls).
49pub struct CheckpointWriter {
50    /// Absolute path to `logs/<session>/checkpoint.jsonl`.
51    path: PathBuf,
52    /// When `false` (a dry-run writer), every event is dropped before
53    /// it touches disk and no `checkpoint.jsonl` is created — a dry-run
54    /// records no resumable state (SRD-44). The in-memory mirror is
55    /// still maintained so `snapshot` stays coherent for tests, but
56    /// `resume_hint` stays silent.
57    enabled: bool,
58    inner: Mutex<Inner>,
59    /// Set once by [`Self::mark_run_reached_end`] when the runner
60    /// reaches its session-end boundary (both run shapes converge
61    /// there; early error returns never set it). Read by
62    /// [`Self::resume_hint`]: a run that ended cleanly must not
63    /// read declared-but-never-started phases as resumable work —
64    /// the runner declares EVERY pre-mapped phase up front (so a
65    /// resume can tell "didn't run yet" from "wasn't planned"),
66    /// and runtime predicates (`continue_if`, for-loop bounds)
67    /// legitimately leave some of those entries Pending forever.
68    run_reached_end: std::sync::atomic::AtomicBool,
69    /// Open handle on the held lockfile
70    /// (`logs/<session>/checkpoint.lock`). The handle is owned
71    /// for the lifetime of the writer; closing it (Drop)
72    /// releases the advisory lock automatically. `None` when
73    /// lock acquisition failed soft.
74    _lock_fd: Option<LockHandle>,
75}
76
77/// Owning wrapper around the open lockfile that releases the
78/// advisory lock when dropped — unlock is implicit on close
79/// (`flock(LOCK_UN)` on Unix, `UnlockFile` on Windows; both are
80/// what `std`'s file locking maps to).
81struct LockHandle(#[allow(dead_code)] File);
82
83struct Inner {
84    /// In-memory fold mirror. Always reflects the most recent
85    /// state declared / observed; consumers (`snapshot`,
86    /// `resume_hint`) read this. The on-disk view is the
87    /// append-only event stream — these two converge after a
88    /// reader folds the log.
89    doc: Checkpoint,
90    /// Map from `identity_key(&PhaseIdentity)` to the index of
91    /// the matching entry in `doc.phases`. Avoids an O(n) linear
92    /// scan on every lifecycle call.
93    index: HashMap<String, usize>,
94    /// Open append-mode handle to the JSONL log. Keeping the
95    /// file open avoids one `open(2)` per event; the `O_APPEND`
96    /// flag guarantees writes are atomic up to `PIPE_BUF`
97    /// (4 KB on Linux) without explicit locking.
98    file: File,
99}
100
101impl CheckpointWriter {
102    /// Construct a fresh writer for a brand-new session. Emits
103    /// the leading `session_start` record before returning.
104    /// Subsequent `declare_phase` / lifecycle calls append events
105    /// onto the same log.
106    pub fn new(path: PathBuf, session: String, started_at: String, invocation: u32) -> Self {
107        let doc = Checkpoint {
108            version: CHECKPOINT_VERSION,
109            session: session.clone(),
110            started_at: started_at.clone(),
111            checkpoint_at: started_at.clone(),
112            invocation,
113            phases: Vec::new(),
114        };
115        let lock = acquire_flock(&path);
116        let file = open_append(&path);
117        let writer = Self {
118            path,
119            enabled: true,
120            inner: Mutex::new(Inner {
121                doc,
122                index: HashMap::new(),
123                file,
124            }),
125            run_reached_end: std::sync::atomic::AtomicBool::new(false),
126            _lock_fd: lock,
127        };
128        // Per SRD-44a §"File location and format" — the first
129        // line of a fresh log MUST be a `session_start` record;
130        // the reader rejects logs whose first record is anything
131        // else. Resume continues the same log by appending its
132        // own `session_start`, so this stays correct across
133        // invocations too.
134        writer.append_event(CheckpointData::SessionStart {
135            at: now_rfc3339(),
136            version: CHECKPOINT_VERSION,
137            session,
138            started_at,
139            invocation,
140        });
141        writer
142    }
143
144    /// Restore a writer from a previously-written document on
145    /// resume. The caller is responsible for having parsed and
146    /// version-checked the document via [`super::storage::read`].
147    /// The restored writer keeps the existing phase entries and
148    /// emits a fresh `session_start` event with the new
149    /// invocation counter — appending continues onto the same
150    /// JSONL log per SRD-44a.
151    pub fn from_existing(
152        path: PathBuf,
153        mut doc: Checkpoint,
154        new_checkpoint_at: String,
155        new_invocation: u32,
156    ) -> Self {
157        doc.checkpoint_at = new_checkpoint_at;
158        doc.invocation = new_invocation;
159        let session = doc.session.clone();
160        let started_at = doc.started_at.clone();
161        let index = build_index(&doc.phases);
162        let lock = acquire_flock(&path);
163        let file = open_append(&path);
164        let writer = Self {
165            path,
166            enabled: true,
167            inner: Mutex::new(Inner { doc, index, file }),
168            run_reached_end: std::sync::atomic::AtomicBool::new(false),
169            _lock_fd: lock,
170        };
171        writer.append_event(CheckpointData::SessionStart {
172            at: now_rfc3339(),
173            version: CHECKPOINT_VERSION,
174            session,
175            started_at,
176            invocation: new_invocation,
177        });
178        writer
179    }
180
181    /// A resume-inert writer for dry-runs: no `checkpoint.jsonl` is
182    /// created, no flock is taken, and every lifecycle event is dropped
183    /// before it reaches disk. A dry-run short-circuits ops and may run
184    /// against placeholder params, so persisting its phases as
185    /// "completed" would poison a later `--resume-latest` (SRD-44) —
186    /// this constructor guarantees it can't. The file handle targets
187    /// `/dev/null` purely to keep the struct total; `append_event`
188    /// never writes to it.
189    pub fn disabled(path: PathBuf) -> Self {
190        let doc = Checkpoint {
191            version: CHECKPOINT_VERSION,
192            session: String::new(),
193            started_at: String::new(),
194            checkpoint_at: String::new(),
195            invocation: 0,
196            phases: Vec::new(),
197        };
198        let file = open_append(std::path::Path::new("/dev/null"));
199        Self {
200            path,
201            enabled: false,
202            inner: Mutex::new(Inner {
203                doc,
204                index: HashMap::new(),
205                file,
206            }),
207            run_reached_end: std::sync::atomic::AtomicBool::new(false),
208            _lock_fd: None,
209        }
210    }
211
212    /// Declare a phase the run plans to execute. Called during
213    /// pre-map for every phase; idempotent (re-declaration is a
214    /// no-op so the resume path can declare the same phases the
215    /// saved doc already lists).
216    pub fn declare_phase(&self, identity: PhaseIdentity, skip_eligible: bool) {
217        let event = {
218            let mut g = self.inner.lock().unwrap();
219            let key = identity_key(&identity);
220            if g.index.contains_key(&key) {
221                return;
222            }
223            let entry = PhaseEntry {
224                identity: identity.clone(),
225                skip_eligible,
226                params_consumed: None,
227                status: PhaseStatus::Pending,
228                duration_secs: None,
229                op_counts: None,
230                cursor_state: None,
231                error: None,
232            };
233            g.doc.phases.push(entry);
234            let idx = g.doc.phases.len() - 1;
235            g.index.insert(key, idx);
236            CheckpointData::PhaseDeclared {
237                at: now_rfc3339(),
238                identity,
239                skip_eligible,
240            }
241        };
242        self.append_event(event);
243    }
244
245    /// Mark a declared phase as `Running`.
246    pub fn phase_started(&self, identity: &PhaseIdentity) {
247        let updated = self.with_entry(identity, |e| {
248            e.status = PhaseStatus::Running;
249            e.error = None;
250        });
251        if updated {
252            self.append_event(CheckpointData::PhaseStarted {
253                at: now_rfc3339(),
254                identity: identity.clone(),
255            });
256        }
257    }
258
259    /// Mark a declared phase as `Completed` with the given
260    /// wall-clock duration.
261    pub fn phase_completed(&self, identity: &PhaseIdentity, duration_secs: f64) {
262        let final_counts = {
263            let mut g = self.inner.lock().unwrap();
264            let key = identity_key(identity);
265            if let Some(&idx) = g.index.get(&key) {
266                let entry = &mut g.doc.phases[idx];
267                let counts = entry.op_counts.clone().unwrap_or_default();
268                entry.status = PhaseStatus::Completed;
269                entry.duration_secs = Some(duration_secs);
270                // Pin the final counts on the in-memory mirror
271                // so a subsequent `snapshot()` matches what the
272                // event records — and what the reader-fold will
273                // produce. Leaving it `None` would diverge from
274                // the disk view.
275                entry.op_counts = Some(counts.clone());
276                entry.cursor_state = None;
277                entry.error = None;
278                counts
279            } else {
280                return;
281            }
282        };
283        self.append_event(CheckpointData::PhaseCompleted {
284            at: now_rfc3339(),
285            identity: identity.clone(),
286            duration_secs,
287            op_counts: final_counts,
288        });
289    }
290
291    /// Mark a declared phase as `Failed`. The error message is
292    /// preserved for resume diagnostics.
293    pub fn phase_failed(&self, identity: &PhaseIdentity, error: &str) {
294        let counts = {
295            let err_owned = error.to_string();
296            let mut g = self.inner.lock().unwrap();
297            let key = identity_key(identity);
298            if let Some(&idx) = g.index.get(&key) {
299                let entry = &mut g.doc.phases[idx];
300                entry.status = PhaseStatus::Failed;
301                entry.error = Some(err_owned);
302                entry.cursor_state = None;
303                entry.op_counts.clone()
304            } else {
305                return;
306            }
307        };
308        self.append_event(CheckpointData::PhaseFailed {
309            at: now_rfc3339(),
310            identity: identity.clone(),
311            error: error.to_string(),
312            op_counts: counts,
313        });
314    }
315
316    /// Record op-execution counts from the live activity. Called
317    /// from the metrics tick callback for the currently-running
318    /// phase. Folds into the matching `PhaseEntry` and emits one
319    /// `phase_progress` event per call — the reader keeps only
320    /// the most recent per identity when folding.
321    pub fn update_op_counts(&self, identity: &PhaseIdentity, counts: OpCounts) {
322        let cursor_state = {
323            let mut g = self.inner.lock().unwrap();
324            let key = identity_key(identity);
325            if let Some(&idx) = g.index.get(&key) {
326                let entry = &mut g.doc.phases[idx];
327                entry.op_counts = Some(counts.clone());
328                entry.cursor_state.clone()
329            } else {
330                return;
331            }
332        };
333        self.append_event(CheckpointData::PhaseProgress {
334            at: now_rfc3339(),
335            identity: identity.clone(),
336            op_counts: counts,
337            cursor_state,
338        });
339    }
340
341    /// Set the program-canonical hash on a declared phase.
342    pub fn update_phase_hash(
343        &self,
344        identity: &PhaseIdentity,
345        hash: [u8; 32],
346        params_consumed: Option<String>,
347    ) {
348        let updated = self.with_entry(identity, |e| {
349            e.identity.phase_hash = Some(hash);
350            e.params_consumed = params_consumed.clone();
351        });
352        if updated {
353            self.append_event(CheckpointData::PhaseHash {
354                at: now_rfc3339(),
355                identity: identity.clone(),
356                hash_hex: hash_to_hex(&hash),
357                params_consumed,
358            });
359        }
360    }
361
362    /// Record the latest cursor-state snapshot for a Tier 2
363    /// phase.
364    pub fn update_cursor(&self, identity: &PhaseIdentity, cursor_state: serde_json::Value) {
365        let counts = {
366            let mut g = self.inner.lock().unwrap();
367            let key = identity_key(identity);
368            if let Some(&idx) = g.index.get(&key) {
369                let entry = &mut g.doc.phases[idx];
370                entry.cursor_state = Some(cursor_state.clone());
371                entry.op_counts.clone().unwrap_or_default()
372            } else {
373                return;
374            }
375        };
376        self.append_event(CheckpointData::PhaseProgress {
377            at: now_rfc3339(),
378            identity: identity.clone(),
379            op_counts: counts,
380            cursor_state: Some(cursor_state),
381        });
382    }
383
384    /// Emit a `scope_enter` event marking entry into a
385    /// `for_each` / `for_combinations` / `do_while` / `do_until`
386    /// iteration. Per SRD-44a §"Event taxonomy", `coords` is the
387    /// `{var: value}` map for THIS scope's own bindings and
388    /// `path` is the leaf-first chain of enclosing scopes' coords
389    /// — together they pin the executor's position in the
390    /// scenario tree at iteration time. Scope events are
391    /// write-and-go: the in-memory mirror has no slot to fold
392    /// them into (the reader's fold is also a no-op today), so
393    /// this is one direct `append_event`.
394    pub fn emit_scope_enter(
395        &self,
396        kind: &str,
397        coords: BTreeMap<String, serde_json::Value>,
398        path: Vec<BTreeMap<String, serde_json::Value>>,
399    ) {
400        self.append_event(CheckpointData::ScopeEnter {
401            at: now_rfc3339(),
402            kind: kind.to_string(),
403            coords,
404            path,
405        });
406    }
407
408    /// Emit a `scope_exit` event marking the end of one
409    /// iteration. `outcome` is `"completed"` when the iteration's
410    /// terminal action returned `Ok`, `"interrupted"` when it
411    /// returned an error or the executor unwound through a stop
412    /// signal. Same write-and-go shape as
413    /// [`emit_scope_enter`](Self::emit_scope_enter).
414    pub fn emit_scope_exit(
415        &self,
416        kind: &str,
417        coords: BTreeMap<String, serde_json::Value>,
418        path: Vec<BTreeMap<String, serde_json::Value>>,
419        outcome: &str,
420    ) {
421        self.append_event(CheckpointData::ScopeExit {
422            at: now_rfc3339(),
423            kind: kind.to_string(),
424            coords,
425            path,
426            outcome: outcome.to_string(),
427        });
428    }
429
430    /// Force a `fdatasync(2)` on the underlying log. Per SRD-44a
431    /// §"Writer behaviour", lifecycle records (start / completed
432    /// / failed) deserve a per-event sync; the periodic
433    /// progress tick can batch. The runtime calls this at every
434    /// phase-lifecycle boundary so a crash between ticks loses
435    /// at most one tick's worth of progress.
436    pub fn flush(&self) -> Result<(), String> {
437        // Nothing is ever written on a dry-run writer, and its handle is
438        // `/dev/null` (which rejects fdatasync) — so flush is a no-op.
439        if !self.enabled {
440            return Ok(());
441        }
442        let g = self.inner.lock().unwrap();
443        match g.file.sync_data() {
444            Ok(()) => Ok(()),
445            Err(e) => Err(format!("fdatasync {}: {e}", self.path.display())),
446        }
447    }
448
449    /// Read-only snapshot of the current in-memory document.
450    /// Useful for diagnostics and tests.
451    pub fn snapshot(&self) -> Checkpoint {
452        self.inner.lock().unwrap().doc.clone()
453    }
454
455    /// Record that the runner reached its session-end boundary
456    /// (the convergence point right before `run_finished()`).
457    /// Early error returns and interrupts never get here, so the
458    /// flag cleanly separates "the run ended" from "the run was
459    /// cut short" for [`Self::resume_hint`].
460    pub fn mark_run_reached_end(&self) {
461        self.run_reached_end
462            .store(true, std::sync::atomic::Ordering::Relaxed);
463    }
464
465    /// If the workload has incomplete phases declared
466    /// `checkpoint: idempotent`, return a multi-line hint string
467    /// the runtime can show the operator on exit.
468    pub fn resume_hint(&self) -> Option<String> {
469        // A dry-run persisted nothing — never advise resuming it.
470        if !self.enabled {
471            return None;
472        }
473        let ended = self
474            .run_reached_end
475            .load(std::sync::atomic::Ordering::Relaxed);
476        let cp = self.snapshot();
477        let recoverable = cp.phases.iter().any(|e| {
478            e.skip_eligible
479                && match e.status {
480                    PhaseStatus::Completed => false,
481                    // Failed (and the anomalous started-never-finished)
482                    // are actionable regardless of how the run ended.
483                    PhaseStatus::Failed | PhaseStatus::Running => true,
484                    // Declared-but-never-started is only evidence of
485                    // interruption when the run was cut short. Every
486                    // pre-mapped phase is declared up front (so resume
487                    // can tell "didn't run yet" from "wasn't planned"),
488                    // and runtime predicates (`continue_if`, for-loop
489                    // bounds) legitimately leave some entries Pending
490                    // forever — a clean run end must not read those as
491                    // work left behind. (Observed 2026-08-05: a fully
492                    // successful 78/78 adaptive run advised resuming
493                    // its 28 continue_if-excluded tiers.)
494                    PhaseStatus::Pending => !ended,
495                }
496        });
497        if !recoverable {
498            return None;
499        }
500        Some(format!(
501            "This session has resumable phases that didn't complete.\n  \
502             To continue from where it stopped:\n    \
503             nmbrs run <workload> --session-dir {} (already set if you exported \
504             SESSION_DIRECTORY) --resume\n  \
505             To pin the session name for repeatable resumes:\n    \
506             nmbrs run <workload> --session {} (then add --resume next time)",
507            self.path
508                .parent()
509                .map(|p| p.display().to_string())
510                .unwrap_or_default(),
511            cp.session,
512        ))
513    }
514
515    /// Path the writer flushes to.
516    pub fn path(&self) -> &std::path::Path {
517        &self.path
518    }
519
520    /// Append one event line to the log. Mutates the file, but
521    /// not the in-memory mirror — the caller is responsible for
522    /// updating the mirror first (so a reader-fold and the
523    /// in-memory snapshot stay equivalent).
524    fn append_event(&self, event: CheckpointData) {
525        // Dry-run writer: drop every event before it touches disk so no
526        // resumable state is ever persisted (SRD-44).
527        if !self.enabled {
528            return;
529        }
530        let mut g = self.inner.lock().unwrap();
531        let mut line = match serde_json::to_string(&event) {
532            Ok(s) => s,
533            Err(e) => {
534                // Serialisation failure is a programming bug —
535                // log loudly and drop the event rather than
536                // panicking the whole session.
537                eprintln!("checkpoint: serialise event failed: {e}; dropping record",);
538                return;
539            }
540        };
541        line.push('\n');
542        if let Err(e) = g.file.write_all(line.as_bytes()) {
543            eprintln!("checkpoint: append to {}: {e}", self.path.display(),);
544        }
545    }
546
547    /// Apply a mutation closure to the entry matching `identity`,
548    /// returning `true` when the entry was found and updated. The
549    /// caller emits the event after a successful mutation so the
550    /// in-memory mirror and the on-disk log stay in lock-step.
551    fn with_entry<F: FnOnce(&mut PhaseEntry)>(&self, identity: &PhaseIdentity, f: F) -> bool {
552        let mut g = self.inner.lock().unwrap();
553        let key = identity_key(identity);
554        if let Some(&idx) = g.index.get(&key) {
555            f(&mut g.doc.phases[idx]);
556            true
557        } else {
558            false
559        }
560    }
561}
562
563/// Open `path` in append mode, creating it (and the parent
564/// directory) if needed. Returns the file or panics — the
565/// writer can't function without it, so a hard failure here is
566/// the correct response. Production callers wrap construction
567/// in a fallible bootstrap path; the panic surfaces as a
568/// session-startup error rather than a silent skip.
569fn open_append(path: &std::path::Path) -> File {
570    if let Some(parent) = path.parent()
571        && let Err(e) = std::fs::create_dir_all(parent)
572    {
573        panic!(
574            "checkpoint: create parent dir {} for {}: {e}",
575            parent.display(),
576            path.display(),
577        );
578    }
579    OpenOptions::new()
580        .create(true)
581        .append(true)
582        .open(path)
583        .unwrap_or_else(|e| panic!("checkpoint: open append {} failed: {e}", path.display(),))
584}
585
586/// Build a lookup key for the index map. Identity equality is
587/// `(yaml_path, coords)` per SRD-44; the hash is sufficiency,
588/// not identity, so it's deliberately excluded from the key.
589pub(crate) fn identity_key(identity: &PhaseIdentity) -> String {
590    let path_json = serde_json::to_string(&identity.yaml_path).unwrap_or_else(|_| String::new());
591    format!("{path_json}\x1f{}", identity.coords)
592}
593
594/// Take a non-blocking exclusive advisory lock on a sibling
595/// `checkpoint.lock` file alongside the checkpoint document.
596fn acquire_flock(checkpoint_path: &std::path::Path) -> Option<LockHandle> {
597    let parent = checkpoint_path.parent()?;
598    if let Err(e) = std::fs::create_dir_all(parent) {
599        eprintln!(
600            "warning: could not create checkpoint dir {}: {e} (concurrent-resume protection skipped)",
601            parent.display(),
602        );
603        return None;
604    }
605    let lock_path = checkpoint_path.with_extension("lock");
606    let file = match OpenOptions::new()
607        .read(true)
608        .write(true)
609        .create(true)
610        .open(&lock_path)
611    {
612        Ok(f) => f,
613        Err(e) => {
614            eprintln!(
615                "warning: could not open lockfile {}: {e} (concurrent-resume protection skipped)",
616                lock_path.display(),
617            );
618            return None;
619        }
620    };
621    match file.try_lock() {
622        Ok(()) => Some(LockHandle(file)),
623        Err(std::fs::TryLockError::WouldBlock) => {
624            panic!(
625                "checkpoint: another process holds the resume lock at {} \
626                 (concurrent `nmbrs run --resume` against the same session?). \
627                 If you're certain no other process is running, remove the \
628                 lockfile and retry.",
629                lock_path.display(),
630            );
631        }
632        Err(std::fs::TryLockError::Error(e)) => {
633            eprintln!(
634                "warning: lock on {} failed: {e} (concurrent-resume protection skipped)",
635                lock_path.display(),
636            );
637            None
638        }
639    }
640}
641
642fn build_index(phases: &[PhaseEntry]) -> HashMap<String, usize> {
643    let mut m = HashMap::with_capacity(phases.len());
644    for (i, e) in phases.iter().enumerate() {
645        m.insert(identity_key(&e.identity), i);
646    }
647    m
648}
649
650#[cfg(test)]
651mod tests {
652    use super::*;
653    use crate::checkpoint::{PathSegment, PhaseIdentity};
654
655    fn ident(name: &str, coords: &str) -> PhaseIdentity {
656        PhaseIdentity {
657            yaml_path: vec![
658                PathSegment::Scenario("s".into()),
659                PathSegment::Phase(name.into()),
660            ],
661            coords: coords.into(),
662            phase_hash: Some([0xcd; 32]),
663        }
664    }
665
666    fn tempdir() -> std::path::PathBuf {
667        let d = std::env::temp_dir().join(format!(
668            "nmbrs-checkpoint-writer-{}",
669            crate::scratch_suffix()
670        ));
671        std::fs::create_dir_all(&d).unwrap();
672        d
673    }
674
675    #[test]
676    fn declare_then_complete_then_flush() {
677        let dir = tempdir();
678        let path = dir.join("checkpoint.jsonl");
679        let w = CheckpointWriter::new(
680            path.clone(),
681            "sess".into(),
682            "2026-01-01T00:00:00Z".into(),
683            1,
684        );
685        let id = ident("schema", "");
686        w.declare_phase(id.clone(), true);
687        w.phase_started(&id);
688        w.phase_completed(&id, 1.5);
689        w.flush().expect("flush");
690
691        // Verify in-memory mirror reflects the lifecycle.
692        let snap = w.snapshot();
693        assert_eq!(snap.phases.len(), 1);
694        assert_eq!(snap.phases[0].status, PhaseStatus::Completed);
695        assert_eq!(snap.phases[0].duration_secs, Some(1.5));
696
697        // Verify the on-disk log carries each event line in order.
698        let raw = std::fs::read_to_string(&path).expect("read log");
699        let lines: Vec<&str> = raw.lines().collect();
700        assert_eq!(
701            lines.len(),
702            4,
703            "expected session_start, phase_declared, phase_started, phase_completed"
704        );
705        assert!(lines[0].contains("\"type\":\"session_start\""));
706        assert!(lines[1].contains("\"type\":\"phase_declared\""));
707        assert!(lines[2].contains("\"type\":\"phase_started\""));
708        assert!(lines[3].contains("\"type\":\"phase_completed\""));
709    }
710
711    #[test]
712    fn resume_hint_respects_run_end_boundary() {
713        // The 2026-08-05 false positive: every pre-mapped phase is
714        // declared up front, and runtime predicates (continue_if,
715        // for-loop bounds) leave excluded ones Pending forever — a
716        // fully successful 78/78 run advised resuming its 28
717        // excluded tiers. Pending counts as resumable ONLY when the
718        // run never reached its end boundary; Failed counts always.
719        let dir = tempdir();
720        let w = CheckpointWriter::new(
721            dir.join("checkpoint.jsonl"),
722            "sess".into(),
723            "2026-01-01T00:00:00Z".into(),
724            1,
725        );
726        let ran = ident("tier", "(part=0)");
727        let excluded = ident("tier", "(part=17)");
728        w.declare_phase(ran.clone(), true);
729        w.declare_phase(excluded.clone(), true);
730        w.phase_started(&ran);
731        w.phase_completed(&ran, 1.0);
732
733        // Cut short (no end mark): the Pending entry is resumable.
734        assert!(
735            w.resume_hint().is_some(),
736            "an interrupted run must advise resuming pending phases"
737        );
738
739        // Clean end: the same Pending entry was excluded by a
740        // runtime predicate, not left behind.
741        w.mark_run_reached_end();
742        assert!(
743            w.resume_hint().is_none(),
744            "a run that reached its end must not advise resuming \
745             predicate-excluded phases"
746        );
747
748        // A failure stays actionable even after a clean end.
749        let failed = ident("tier", "(part=3)");
750        w.declare_phase(failed.clone(), true);
751        w.phase_started(&failed);
752        w.phase_failed(&failed, "boom");
753        assert!(
754            w.resume_hint().is_some(),
755            "failed phases must keep the hint even on a clean end"
756        );
757    }
758
759    #[test]
760    fn disabled_writer_persists_nothing() {
761        // SRD-44 dry-run gate: a disabled writer must create NO
762        // checkpoint.jsonl and drop every lifecycle event, so a later
763        // `--resume-latest` can never pick a dry-run's phases up.
764        let dir = tempdir();
765        let path = dir.join("checkpoint.jsonl");
766        let w = CheckpointWriter::disabled(path.clone());
767        let id = ident("teardown", "(table=changeme_default)");
768        w.declare_phase(id.clone(), true);
769        w.phase_started(&id);
770        w.phase_completed(&id, 1.0);
771        w.flush().expect("flush is a harmless no-op");
772
773        assert!(
774            !path.exists(),
775            "dry-run must not create a checkpoint file at {}",
776            path.display()
777        );
778        assert!(
779            w.resume_hint().is_none(),
780            "dry-run must never advertise a resumable session"
781        );
782    }
783
784    #[test]
785    fn redundant_declare_is_idempotent() {
786        let dir = tempdir();
787        let path = dir.join("c.jsonl");
788        let w = CheckpointWriter::new(path.clone(), "s".into(), "t".into(), 1);
789        let id = ident("p", "(k=1)");
790        w.declare_phase(id.clone(), true);
791        w.declare_phase(id.clone(), false); // re-declare ignored
792        let snap = w.snapshot();
793        assert_eq!(snap.phases.len(), 1);
794        assert!(snap.phases[0].skip_eligible, "first declare wins");
795
796        // Only one phase_declared in the log (plus the leading
797        // session_start).
798        let raw = std::fs::read_to_string(&path).expect("read");
799        let count = raw
800            .lines()
801            .filter(|l| l.contains("\"type\":\"phase_declared\""))
802            .count();
803        assert_eq!(count, 1, "second declare must not emit a duplicate event");
804    }
805
806    #[test]
807    fn from_existing_emits_fresh_session_start() {
808        let dir = tempdir();
809        let path = dir.join("c.jsonl");
810        // First session writes some events, then drops to release
811        // the flock. Mirrors the production lifecycle: prior
812        // process exits before resume starts.
813        let saved = {
814            let w =
815                CheckpointWriter::new(path.clone(), "s".into(), "2026-01-01T00:00:00Z".into(), 1);
816            let id = ident("schema", "");
817            w.declare_phase(id.clone(), true);
818            w.phase_completed(&id, 0.5);
819            w.flush().expect("flush");
820            w.snapshot()
821        };
822        let w2 =
823            CheckpointWriter::from_existing(path.clone(), saved, "2026-01-01T00:01:00Z".into(), 2);
824        let snap = w2.snapshot();
825        assert_eq!(snap.invocation, 2);
826        assert_eq!(snap.phases.len(), 1);
827        assert_eq!(snap.phases[0].status, PhaseStatus::Completed);
828
829        // Log carries TWO session_start events, marking the
830        // invocation boundary.
831        let raw = std::fs::read_to_string(&path).expect("read");
832        let count = raw
833            .lines()
834            .filter(|l| l.contains("\"type\":\"session_start\""))
835            .count();
836        assert_eq!(
837            count, 2,
838            "resume must append a fresh session_start, not rewrite"
839        );
840    }
841
842    #[test]
843    fn scope_enter_exit_pairs_for_two_deep_for_each() {
844        // SRD-44a Push 3: drive the writer over the event
845        // sequence the executor's scope walker produces for a
846        // 2-deep `for_each` workload — outer iterates `x in
847        // [1, 2]`, inner iterates `y in ["a", "b"]`. Two
848        // outer iterations × two inner iterations = four leaf
849        // bracket pairs, plus the two outer-loop bracket pairs
850        // wrapping each inner sub-walk.
851        let dir = tempdir();
852        let path = dir.join("c.jsonl");
853        let w = CheckpointWriter::new(path.clone(), "s".into(), "2026-01-01T00:00:00Z".into(), 1);
854
855        // Scope coords as the executor would synthesize per
856        // iteration. `coord_path` is root-first; the writer's
857        // event takes the leaf as `coords` and the prefix
858        // (reversed to leaf-first) as `path`. We model that
859        // shape here directly.
860        let outer_coord = |xv: u64| -> BTreeMap<String, serde_json::Value> {
861            let mut m = BTreeMap::new();
862            m.insert("x".into(), serde_json::Value::from(xv));
863            m
864        };
865        let inner_coord = |yv: &str| -> BTreeMap<String, serde_json::Value> {
866            let mut m = BTreeMap::new();
867            m.insert("y".into(), serde_json::Value::from(yv));
868            m
869        };
870
871        for x in [1u64, 2u64] {
872            // Outer enter: coords={x=…}, path=[]
873            w.emit_scope_enter("for_each", outer_coord(x), Vec::new());
874            for y in ["a", "b"] {
875                // Inner enter: coords={y=…}, path=[{x=…}]
876                w.emit_scope_enter("for_each", inner_coord(y), vec![outer_coord(x)]);
877                w.emit_scope_exit(
878                    "for_each",
879                    inner_coord(y),
880                    vec![outer_coord(x)],
881                    "completed",
882                );
883            }
884            w.emit_scope_exit("for_each", outer_coord(x), Vec::new(), "completed");
885        }
886        w.flush().expect("flush");
887
888        // Parse the JSONL log line-by-line, dropping the
889        // leading session_start. Confirm the bracket sequence
890        // is well-formed and the path/coords carry the
891        // expected x/y values for each iteration.
892        let raw = std::fs::read_to_string(&path).expect("read log");
893        let scope_events: Vec<serde_json::Value> = raw
894            .lines()
895            .map(|l| serde_json::from_str::<serde_json::Value>(l).expect("parse line"))
896            .filter(|v| {
897                let t = v.get("type").and_then(|t| t.as_str()).unwrap_or("");
898                t == "scope_enter" || t == "scope_exit"
899            })
900            .collect();
901        // 2 outer (enter+exit) + 2*2 inner (enter+exit) = 12 events.
902        assert_eq!(
903            scope_events.len(),
904            12,
905            "expected 12 scope events, got {}",
906            scope_events.len()
907        );
908
909        let kind = |v: &serde_json::Value| {
910            v.get("type")
911                .and_then(|t| t.as_str())
912                .unwrap_or("")
913                .to_string()
914        };
915        let x_at = |v: &serde_json::Value, idx: &str| -> Option<u64> {
916            v.pointer(idx).and_then(|n| n.as_u64())
917        };
918        let y_at = |v: &serde_json::Value, idx: &str| -> Option<String> {
919            v.pointer(idx)
920                .and_then(|n| n.as_str())
921                .map(|s| s.to_string())
922        };
923
924        // Outer entry first: coords={x=1}, path=[].
925        assert_eq!(kind(&scope_events[0]), "scope_enter");
926        assert_eq!(x_at(&scope_events[0], "/coords/x"), Some(1));
927        assert!(
928            scope_events[0]
929                .pointer("/path")
930                .and_then(|p| p.as_array())
931                .map(|a| a.is_empty())
932                .unwrap_or(false),
933            "outer enter must have empty path"
934        );
935
936        // First inner enter under x=1, y=a. path leaf-first is
937        // [{x=1}].
938        assert_eq!(kind(&scope_events[1]), "scope_enter");
939        assert_eq!(y_at(&scope_events[1], "/coords/y"), Some("a".to_string()));
940        assert_eq!(x_at(&scope_events[1], "/path/0/x"), Some(1));
941
942        // First inner exit, completed.
943        assert_eq!(kind(&scope_events[2]), "scope_exit");
944        assert_eq!(
945            scope_events[2].pointer("/outcome").and_then(|s| s.as_str()),
946            Some("completed")
947        );
948
949        // Second inner under x=1, y=b.
950        assert_eq!(y_at(&scope_events[3], "/coords/y"), Some("b".to_string()));
951        assert_eq!(kind(&scope_events[4]), "scope_exit");
952
953        // Outer exit for x=1.
954        assert_eq!(kind(&scope_events[5]), "scope_exit");
955        assert_eq!(x_at(&scope_events[5], "/coords/x"), Some(1));
956        assert_eq!(
957            scope_events[5].pointer("/outcome").and_then(|s| s.as_str()),
958            Some("completed")
959        );
960
961        // Second outer (x=2) entry, plus its inner pairs and exit.
962        assert_eq!(kind(&scope_events[6]), "scope_enter");
963        assert_eq!(x_at(&scope_events[6], "/coords/x"), Some(2));
964        assert_eq!(x_at(&scope_events[7], "/path/0/x"), Some(2));
965        assert_eq!(y_at(&scope_events[7], "/coords/y"), Some("a".to_string()));
966        assert_eq!(kind(&scope_events[11]), "scope_exit");
967        assert_eq!(x_at(&scope_events[11], "/coords/x"), Some(2));
968
969        // The reader's fold treats scope events as no-ops in v1
970        // (no scope mirror in the in-memory document), so a
971        // round-trip read should succeed and leave the phases
972        // list untouched.
973        let folded = super::super::storage::read(&path)
974            .expect("read folds")
975            .expect("non-empty");
976        assert!(
977            folded.phases.is_empty(),
978            "no phases declared, fold should be empty"
979        );
980    }
981
982    #[test]
983    fn scope_exit_outcome_distinguishes_interrupted_from_completed() {
984        // The kind/outcome surface is what distinguishes a clean
985        // bracket from one closed by a stop signal. Verify both
986        // outcomes round-trip through the JSONL.
987        let dir = tempdir();
988        let path = dir.join("c.jsonl");
989        let w = CheckpointWriter::new(path.clone(), "s".into(), "t".into(), 1);
990        let mut coords = BTreeMap::new();
991        coords.insert("k".into(), serde_json::Value::from(7u64));
992        w.emit_scope_enter("do_while", coords.clone(), Vec::new());
993        w.emit_scope_exit("do_while", coords, Vec::new(), "interrupted");
994        w.flush().expect("flush");
995
996        let raw = std::fs::read_to_string(&path).expect("read");
997        let exit_line = raw
998            .lines()
999            .find(|l| l.contains("\"type\":\"scope_exit\""))
1000            .expect("scope_exit line");
1001        let ev: serde_json::Value = serde_json::from_str(exit_line).unwrap();
1002        assert_eq!(
1003            ev.pointer("/kind").and_then(|s| s.as_str()),
1004            Some("do_while")
1005        );
1006        assert_eq!(
1007            ev.pointer("/outcome").and_then(|s| s.as_str()),
1008            Some("interrupted")
1009        );
1010    }
1011
1012    #[test]
1013    fn flock_blocks_concurrent_writer_on_same_path() {
1014        let dir = tempdir();
1015        let path = dir.join("c.jsonl");
1016        let _w = CheckpointWriter::new(path.clone(), "s".into(), "t".into(), 1);
1017        let result = std::panic::catch_unwind(|| {
1018            let _w2 = CheckpointWriter::new(path.clone(), "s".into(), "t".into(), 1);
1019        });
1020        assert!(
1021            result.is_err(),
1022            "second writer should panic on flock contention"
1023        );
1024    }
1025}