Write progress.ndjson and live.json to logs dir during pipeline runs

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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-02-24 09:50:19 -05:00
parent fa828dd8b1
commit 3d8e17aedf
2 changed files with 85 additions and 1 deletions

View file

@ -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));

View file

@ -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());
}