1use 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
40pub const CHECKPOINT_VERSION: u32 = 1;
44
45pub struct CheckpointWriter {
50 path: PathBuf,
52 enabled: bool,
58 inner: Mutex<Inner>,
59 run_reached_end: std::sync::atomic::AtomicBool,
69 _lock_fd: Option<LockHandle>,
75}
76
77struct LockHandle(#[allow(dead_code)] File);
82
83struct Inner {
84 doc: Checkpoint,
90 index: HashMap<String, usize>,
94 file: File,
99}
100
101impl CheckpointWriter {
102 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 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 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 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 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 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 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 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 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 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 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 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 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 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 pub fn flush(&self) -> Result<(), String> {
437 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 pub fn snapshot(&self) -> Checkpoint {
452 self.inner.lock().unwrap().doc.clone()
453 }
454
455 pub fn mark_run_reached_end(&self) {
461 self.run_reached_end
462 .store(true, std::sync::atomic::Ordering::Relaxed);
463 }
464
465 pub fn resume_hint(&self) -> Option<String> {
469 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 PhaseStatus::Failed | PhaseStatus::Running => true,
484 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 pub fn path(&self) -> &std::path::Path {
517 &self.path
518 }
519
520 fn append_event(&self, event: CheckpointData) {
525 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 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 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
563fn 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
586pub(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
594fn 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 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 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 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 assert!(
735 w.resume_hint().is_some(),
736 "an interrupted run must advise resuming pending phases"
737 );
738
739 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 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 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); let snap = w.snapshot();
793 assert_eq!(snap.phases.len(), 1);
794 assert!(snap.phases[0].skip_eligible, "first declare wins");
795
796 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 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 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 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 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 w.emit_scope_enter("for_each", outer_coord(x), Vec::new());
874 for y in ["a", "b"] {
875 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 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 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 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 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 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 assert_eq!(y_at(&scope_events[3], "/coords/y"), Some("b".to_string()));
951 assert_eq!(kind(&scope_events[4]), "scope_exit");
952
953 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 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 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 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}