From 3d8e17aedf94aea5c2a178be7836d0824f168dc1 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Tue, 24 Feb 2026 09:50:19 -0500 Subject: [PATCH] Write progress.ndjson and live.json to logs dir during pipeline runs MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Every PipelineEvent is now logged as a JSON envelope with timestamp, run_id, and event fields — matching the Kilroy reference implementation. Co-Authored-By: Claude Opus 4.6 --- crates/attractor/src/cli/run.rs | 35 +++++++++++++++++++++- crates/attractor/tests/cli.rs | 51 +++++++++++++++++++++++++++++++++ 2 files changed, 85 insertions(+), 1 deletion(-) diff --git a/crates/attractor/src/cli/run.rs b/crates/attractor/src/cli/run.rs index 86cc03cc6..254c1aca4 100644 --- a/crates/attractor/src/cli/run.rs +++ b/crates/attractor/src/cli/run.rs @@ -3,7 +3,7 @@ use std::sync::{Arc, Mutex}; use std::time::Instant; use anyhow::bail; -use chrono::Local; +use chrono::{Local, Utc}; use terminal::Styles; use crate::checkpoint::Checkpoint; @@ -93,6 +93,39 @@ pub async fn run_command(args: RunArgs, styles: &'static Styles) -> anyhow::Resu } }); + // NDJSON progress log + live.json snapshot + { + let ndjson_path = logs_dir.join("progress.ndjson"); + let live_path = logs_dir.join("live.json"); + let run_id = Arc::new(Mutex::new(String::new())); + let run_id_clone = Arc::clone(&run_id); + emitter.on_event(move |event| { + if let crate::event::PipelineEvent::PipelineStarted { id, .. } = event { + *run_id_clone.lock().unwrap() = id.clone(); + } + let envelope = serde_json::json!({ + "timestamp": Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Millis, true), + "run_id": *run_id_clone.lock().unwrap(), + "event": event, + }); + // Append to progress.ndjson + if let Ok(line) = serde_json::to_string(&envelope) { + use std::io::Write; + if let Ok(mut f) = std::fs::OpenOptions::new() + .create(true) + .append(true) + .open(&ndjson_path) + { + let _ = writeln!(f, "{line}"); + } + } + // Overwrite live.json + if let Ok(pretty) = serde_json::to_string_pretty(&envelope) { + let _ = std::fs::write(&live_path, pretty); + } + }); + } + if args.verbose >= 2 { emitter.on_event(move |event| { eprint!("{}", format_event_detail(event, styles)); diff --git a/crates/attractor/tests/cli.rs b/crates/attractor/tests/cli.rs index 4a964d686..e12dbdd9c 100644 --- a/crates/attractor/tests/cli.rs +++ b/crates/attractor/tests/cli.rs @@ -117,3 +117,54 @@ fn dry_run_styled() { .assert() .success(); } + +// -- NDJSON logging ---------------------------------------------------------- + +#[test] +fn dry_run_writes_ndjson_and_live_json() { + let tmp = tempfile::tempdir().unwrap(); + let logs_dir = tmp.path().join("logs"); + + attractor() + .args([ + "run", + "--dry-run", + "--auto-approve", + "--logs-dir", + logs_dir.to_str().unwrap(), + "../../test/simple.dot", + ]) + .assert() + .success(); + + // progress.ndjson must exist and contain valid JSON lines + let ndjson_path = logs_dir.join("progress.ndjson"); + assert!(ndjson_path.exists(), "progress.ndjson should exist"); + let ndjson_content = std::fs::read_to_string(&ndjson_path).unwrap(); + let lines: Vec<&str> = ndjson_content.lines().collect(); + assert!(!lines.is_empty(), "progress.ndjson should have at least one line"); + + // Every line must be valid JSON with timestamp, run_id, and event keys + let first_line: serde_json::Value = serde_json::from_str(lines[0]).unwrap(); + assert!(first_line.get("timestamp").is_some(), "line should have timestamp"); + assert!(first_line.get("run_id").is_some(), "line should have run_id"); + assert!(first_line.get("event").is_some(), "line should have event"); + + // First event should be PipelineStarted + let first_event = &first_line["event"]; + assert!(first_event.get("PipelineStarted").is_some(), "first event should be PipelineStarted"); + + // run_id should be non-empty after PipelineStarted + let last_line: serde_json::Value = serde_json::from_str(lines[lines.len() - 1]).unwrap(); + let run_id = last_line["run_id"].as_str().unwrap(); + assert!(!run_id.is_empty(), "run_id should be non-empty"); + + // live.json must exist and contain valid JSON matching the last NDJSON line + let live_path = logs_dir.join("live.json"); + assert!(live_path.exists(), "live.json should exist"); + let live_content: serde_json::Value = + serde_json::from_str(&std::fs::read_to_string(&live_path).unwrap()).unwrap(); + assert!(live_content.get("timestamp").is_some()); + assert!(live_content.get("run_id").is_some()); + assert!(live_content.get("event").is_some()); +}