Skip to main content

harn_cli/commands/run/
json_events.rs

1//! `harn run --json`: NDJSON event-stream emitter.
2//!
3//! Each line is a [`JsonEnvelope`] wrapping a [`RunEventWire`]. Wire
4//! events tag themselves with `event_type` for cheap discrimination
5//! by `jq`-style consumers and carry a strictly monotonic `seq`
6//! starting at `1`. See issue #1755 / epic #1753.
7
8use std::io::Write;
9use std::sync::atomic::{AtomicU64, Ordering};
10use std::sync::{Arc, Mutex};
11
12use harn_vm::run_events::{RunEvent, RunEventSink};
13use serde::Serialize;
14
15use crate::json_envelope::{JsonEnvelope, JsonError};
16
17/// Schema version for the `harn run --json` event stream. Bump on any
18/// breaking change to the wire shape; agents key off this to negotiate
19/// compatibility.
20pub const RUN_JSON_SCHEMA_VERSION: u32 = 1;
21
22/// Wire form of a single event emitted by `harn run --json`. The
23/// `event_type` tag is flat so consumers can `jq '.data.event_type'`.
24/// `seq` is monotonic and process-local — the first event in a run is
25/// `seq=1`.
26#[derive(Debug, Clone, Serialize)]
27#[serde(tag = "event_type", rename_all = "snake_case")]
28pub enum RunEventWire {
29    Stdout {
30        seq: u64,
31        payload: String,
32    },
33    Stderr {
34        seq: u64,
35        payload: String,
36    },
37    Transcript {
38        seq: u64,
39        #[serde(skip_serializing_if = "Option::is_none")]
40        agent_id: Option<String>,
41        kind: String,
42        payload: serde_json::Value,
43    },
44    ToolCall {
45        seq: u64,
46        call_id: String,
47        name: String,
48        args: serde_json::Value,
49        started_at: String,
50    },
51    ToolResult {
52        seq: u64,
53        call_id: String,
54        ok: bool,
55        result: serde_json::Value,
56    },
57    Hook {
58        seq: u64,
59        name: String,
60        phase: String,
61        #[serde(skip_serializing_if = "serde_json::Value::is_null")]
62        payload: serde_json::Value,
63    },
64    PersonaStage {
65        seq: u64,
66        persona: String,
67        stage: String,
68        transition: String,
69    },
70    PackRun {
71        seq: u64,
72        bundle_hash: String,
73        signature_verified: bool,
74        #[serde(skip_serializing_if = "Option::is_none")]
75        key_id: Option<String>,
76        cache_hit: bool,
77        dry_run_verify: bool,
78        execution_artifact_state: String,
79        #[serde(skip_serializing_if = "Option::is_none")]
80        fallback_reason: Option<String>,
81        artifact_decode_ms: u64,
82    },
83    Result {
84        seq: u64,
85        value: serde_json::Value,
86        exit_code: i32,
87    },
88    Error {
89        seq: u64,
90        error: JsonError,
91    },
92}
93
94impl RunEventWire {
95    /// The monotonic sequence number assigned at emission.
96    pub fn seq(&self) -> u64 {
97        match self {
98            Self::Stdout { seq, .. }
99            | Self::Stderr { seq, .. }
100            | Self::Transcript { seq, .. }
101            | Self::ToolCall { seq, .. }
102            | Self::ToolResult { seq, .. }
103            | Self::Hook { seq, .. }
104            | Self::PersonaStage { seq, .. }
105            | Self::PackRun { seq, .. }
106            | Self::Result { seq, .. }
107            | Self::Error { seq, .. } => *seq,
108        }
109    }
110}
111
112/// Writer that drains [`RunEvent`]s, assigns monotonic seq numbers,
113/// wraps them in [`JsonEnvelope`]s, and emits one NDJSON line per
114/// event. Lines are flushed per event so streaming consumers see them
115/// as the run progresses.
116pub struct NdjsonEmitter {
117    inner: Arc<NdjsonEmitterInner>,
118}
119
120struct NdjsonEmitterInner {
121    seq: AtomicU64,
122    quiet: bool,
123    /// Output sink. Behind a Mutex so concurrent emits stay
124    /// line-atomic; serde line writes are tiny so contention is
125    /// negligible.
126    out: Mutex<Box<dyn Write + Send>>,
127}
128
129impl NdjsonEmitter {
130    /// Build an emitter that writes to `out`. `quiet` suppresses
131    /// `Stdout` and `Stderr` events (transcript/tool/hook/persona/
132    /// result events still flow).
133    pub fn new(out: Box<dyn Write + Send>, quiet: bool) -> Self {
134        Self {
135            inner: Arc::new(NdjsonEmitterInner {
136                seq: AtomicU64::new(0),
137                quiet,
138                out: Mutex::new(out),
139            }),
140        }
141    }
142
143    /// Build a thread-safe sink that forwards [`RunEvent`]s into this
144    /// emitter, applying seq numbering and `quiet` filtering.
145    pub fn sink(&self) -> Arc<dyn RunEventSink> {
146        Arc::new(NdjsonSink {
147            inner: self.inner.clone(),
148        })
149    }
150
151    /// Next monotonic seq value. The first event in a run is `seq=1`.
152    fn next_seq(&self) -> u64 {
153        self.inner.seq.fetch_add(1, Ordering::SeqCst) + 1
154    }
155
156    fn write_envelope(inner: &NdjsonEmitterInner, event: RunEventWire) {
157        let envelope = JsonEnvelope::ok(RUN_JSON_SCHEMA_VERSION, event);
158        let line = serde_json::to_string(&envelope)
159            .unwrap_or_else(|_| r#"{"schemaVersion":1,"ok":false}"#.to_string());
160        if let Ok(mut out) = inner.out.lock() {
161            let _ = writeln!(out, "{line}");
162            let _ = out.flush();
163        }
164    }
165
166    /// Emit the terminal `Result` event for a run.
167    pub fn emit_result(&self, value: serde_json::Value, exit_code: i32) {
168        let event = RunEventWire::Result {
169            seq: self.next_seq(),
170            value,
171            exit_code,
172        };
173        Self::write_envelope(&self.inner, event);
174    }
175
176    /// Emit the terminal `Error` event for a fatal run failure (e.g.
177    /// compile error before the VM started).
178    pub fn emit_error(&self, code: impl Into<String>, message: impl Into<String>) {
179        let event = RunEventWire::Error {
180            seq: self.next_seq(),
181            error: JsonError {
182                code: code.into(),
183                message: message.into(),
184                details: serde_json::Value::Null,
185            },
186        };
187        Self::write_envelope(&self.inner, event);
188    }
189}
190
191struct NdjsonSink {
192    inner: Arc<NdjsonEmitterInner>,
193}
194
195impl RunEventSink for NdjsonSink {
196    fn emit(&self, event: RunEvent) {
197        let seq = self.inner.seq.fetch_add(1, Ordering::SeqCst) + 1;
198        let wire = match event {
199            RunEvent::Stdout { payload } => {
200                if self.inner.quiet {
201                    // Undo the seq bump so monotonicity stays tight.
202                    self.inner.seq.fetch_sub(1, Ordering::SeqCst);
203                    return;
204                }
205                RunEventWire::Stdout { seq, payload }
206            }
207            RunEvent::Stderr { payload } => {
208                if self.inner.quiet {
209                    self.inner.seq.fetch_sub(1, Ordering::SeqCst);
210                    return;
211                }
212                RunEventWire::Stderr { seq, payload }
213            }
214            RunEvent::Transcript {
215                agent_id,
216                kind,
217                payload,
218            } => RunEventWire::Transcript {
219                seq,
220                agent_id,
221                kind,
222                payload,
223            },
224            RunEvent::ToolCall {
225                call_id,
226                name,
227                args,
228                started_at,
229            } => RunEventWire::ToolCall {
230                seq,
231                call_id,
232                name,
233                args,
234                started_at,
235            },
236            RunEvent::ToolResult {
237                call_id,
238                ok,
239                result,
240            } => RunEventWire::ToolResult {
241                seq,
242                call_id,
243                ok,
244                result,
245            },
246            RunEvent::Hook {
247                name,
248                phase,
249                payload,
250            } => RunEventWire::Hook {
251                seq,
252                name,
253                phase,
254                payload,
255            },
256            RunEvent::PersonaStage {
257                persona,
258                stage,
259                transition,
260            } => RunEventWire::PersonaStage {
261                seq,
262                persona,
263                stage,
264                transition,
265            },
266            RunEvent::PackRun {
267                bundle_hash,
268                signature_verified,
269                key_id,
270                cache_hit,
271                dry_run_verify,
272                execution_artifact_state,
273                fallback_reason,
274                artifact_decode_ms,
275            } => RunEventWire::PackRun {
276                seq,
277                bundle_hash,
278                signature_verified,
279                key_id,
280                cache_hit,
281                dry_run_verify,
282                execution_artifact_state,
283                fallback_reason,
284                artifact_decode_ms,
285            },
286        };
287        NdjsonEmitter::write_envelope(&self.inner, wire);
288    }
289}
290
291#[cfg(test)]
292mod tests {
293    use super::*;
294
295    struct BufWriter(Arc<Mutex<Vec<u8>>>);
296    impl Write for BufWriter {
297        fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
298            self.0.lock().unwrap().extend_from_slice(buf);
299            Ok(buf.len())
300        }
301        fn flush(&mut self) -> std::io::Result<()> {
302            Ok(())
303        }
304    }
305
306    #[tokio::test]
307    async fn emits_monotonic_seq_across_events() {
308        let buf = Arc::new(Mutex::new(Vec::<u8>::new()));
309        let emitter = NdjsonEmitter::new(Box::new(BufWriter(buf.clone())), false);
310        let sink = emitter.sink();
311        harn_vm::run_events::scope(sink, async {
312            harn_vm::run_events::emit(RunEvent::Stdout {
313                payload: "hello\n".into(),
314            });
315            harn_vm::run_events::emit(RunEvent::Stderr {
316                payload: "warn\n".into(),
317            });
318        })
319        .await;
320        emitter.emit_result(serde_json::Value::Null, 0);
321
322        let raw = String::from_utf8(buf.lock().unwrap().clone()).expect("utf8");
323        let lines: Vec<&str> = raw.lines().filter(|line| !line.is_empty()).collect();
324        assert_eq!(lines.len(), 3, "expected 3 NDJSON lines, got:\n{raw}");
325        let seqs: Vec<u64> = lines
326            .iter()
327            .map(|line| {
328                let v: serde_json::Value = serde_json::from_str(line).expect("valid json");
329                v["data"]["seq"].as_u64().expect("seq present")
330            })
331            .collect();
332        assert_eq!(seqs, vec![1, 2, 3]);
333        let types: Vec<String> = lines
334            .iter()
335            .map(|line| {
336                let v: serde_json::Value = serde_json::from_str(line).expect("valid json");
337                v["data"]["event_type"].as_str().expect("type").to_string()
338            })
339            .collect();
340        assert_eq!(types, vec!["stdout", "stderr", "result"]);
341    }
342
343    #[tokio::test]
344    async fn quiet_drops_stdout_and_stderr_without_gaps() {
345        let buf = Arc::new(Mutex::new(Vec::<u8>::new()));
346        let emitter = NdjsonEmitter::new(Box::new(BufWriter(buf.clone())), true);
347        let sink = emitter.sink();
348        harn_vm::run_events::scope(sink, async {
349            harn_vm::run_events::emit(RunEvent::Stdout {
350                payload: "ignored\n".into(),
351            });
352            harn_vm::run_events::emit(RunEvent::Hook {
353                name: "PreRun".into(),
354                phase: "allow".into(),
355                payload: serde_json::Value::Null,
356            });
357        })
358        .await;
359        emitter.emit_result(serde_json::Value::Null, 0);
360
361        let raw = String::from_utf8(buf.lock().unwrap().clone()).expect("utf8");
362        let lines: Vec<&str> = raw.lines().filter(|line| !line.is_empty()).collect();
363        // stdout suppressed; hook + result remain.
364        assert_eq!(lines.len(), 2, "raw:\n{raw}");
365        let seqs: Vec<u64> = lines
366            .iter()
367            .map(|line| {
368                let v: serde_json::Value = serde_json::from_str(line).expect("valid json");
369                v["data"]["seq"].as_u64().expect("seq")
370            })
371            .collect();
372        assert_eq!(
373            seqs,
374            vec![1, 2],
375            "seq must stay contiguous after quiet filtering"
376        );
377    }
378
379    #[tokio::test(flavor = "current_thread")]
380    async fn explicit_exit_emits_stdio_and_one_result_event() {
381        harn_vm::reset_thread_local_state();
382        let temp = tempfile::TempDir::new().expect("temp dir");
383        let script = temp.path().join("main.harn");
384        std::fs::write(
385            &script,
386            r#"
387fn main(harness: Harness) {
388  harness.stdio.print("before ")
389  harness.stdio.println("exit")
390  harness.stdio.eprintln("diagnostic")
391  harness.runtime.exit(2)
392}
393"#,
394        )
395        .expect("write script");
396        let buffer = Arc::new(Mutex::new(Vec::<u8>::new()));
397
398        let outcome = super::super::execute_run_json(
399            &script.to_string_lossy(),
400            false,
401            std::collections::HashSet::new(),
402            Vec::new(),
403            Vec::new(),
404            super::super::CliLlmMockMode::Off,
405            None,
406            super::super::RunProfileOptions::default(),
407            Box::new(BufWriter(buffer.clone())),
408            super::super::RunJsonOptions::default(),
409        )
410        .await;
411
412        assert_eq!(outcome.exit_code, 2, "stderr:\n{}", outcome.stderr);
413        let events: Vec<serde_json::Value> = String::from_utf8(buffer.lock().unwrap().clone())
414            .expect("utf8")
415            .lines()
416            .filter(|line| !line.is_empty())
417            .map(|line| serde_json::from_str(line).expect("valid NDJSON event"))
418            .collect();
419        let stdout = events
420            .iter()
421            .filter(|event| event["data"]["event_type"] == "stdout")
422            .map(|event| event["data"]["payload"].as_str().expect("stdout payload"))
423            .collect::<String>();
424        let stderr = events
425            .iter()
426            .filter(|event| event["data"]["event_type"] == "stderr")
427            .map(|event| event["data"]["payload"].as_str().expect("stderr payload"))
428            .collect::<String>();
429        let terminal: Vec<&serde_json::Value> = events
430            .iter()
431            .filter(|event| event["data"]["event_type"] == "result")
432            .collect();
433
434        assert_eq!(stdout, "before exit\n");
435        assert_eq!(stderr, "diagnostic\n");
436        assert_eq!(terminal.len(), 1, "events: {events:#?}");
437        assert_eq!(terminal[0]["data"]["exit_code"], 2);
438        assert!(terminal[0]["data"]["value"].is_null());
439        assert!(events
440            .iter()
441            .all(|event| event["data"]["event_type"] != "error"));
442        harn_vm::reset_thread_local_state();
443    }
444
445    struct BarrierWriter {
446        bytes: Arc<Mutex<Vec<u8>>>,
447        first_event_barrier: Arc<std::sync::Barrier>,
448        reached_barrier: bool,
449    }
450
451    impl Write for BarrierWriter {
452        fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
453            self.bytes.lock().unwrap().extend_from_slice(buf);
454            if !self.reached_barrier && buf.contains(&b'\n') {
455                self.reached_barrier = true;
456                self.first_event_barrier.wait();
457            }
458            Ok(buf.len())
459        }
460
461        fn flush(&mut self) -> std::io::Result<()> {
462            Ok(())
463        }
464    }
465
466    #[test]
467    fn concurrent_json_runs_receive_only_their_own_ordered_events() {
468        let temp = tempfile::TempDir::new().expect("temp dir");
469        let alpha_path = temp.path().join("alpha.harn");
470        let beta_path = temp.path().join("beta.harn");
471        std::fs::write(
472            &alpha_path,
473            r#"
474fn main(harness: Harness) {
475  harness.stdio.println("alpha-out")
476  harness.stdio.eprintln("alpha-err")
477}
478"#,
479        )
480        .expect("write alpha script");
481        std::fs::write(
482            &beta_path,
483            r#"
484fn main(harness: Harness) {
485  harness.stdio.println("beta-out")
486  harness.stdio.eprintln("beta-err")
487}
488"#,
489        )
490        .expect("write beta script");
491
492        let barrier = Arc::new(std::sync::Barrier::new(2));
493        let alpha_bytes = Arc::new(Mutex::new(Vec::new()));
494        let beta_bytes = Arc::new(Mutex::new(Vec::new()));
495
496        let spawn_run = |path: std::path::PathBuf, bytes: Arc<Mutex<Vec<u8>>>| {
497            let first_event_barrier = barrier.clone();
498            std::thread::spawn(move || {
499                tokio::runtime::Builder::new_current_thread()
500                    .enable_all()
501                    .build()
502                    .expect("runtime")
503                    .block_on(super::super::execute_run_json(
504                        &path.to_string_lossy(),
505                        false,
506                        std::collections::HashSet::new(),
507                        Vec::new(),
508                        Vec::new(),
509                        super::super::CliLlmMockMode::Off,
510                        None,
511                        super::super::RunProfileOptions::default(),
512                        Box::new(BarrierWriter {
513                            bytes,
514                            first_event_barrier,
515                            reached_barrier: false,
516                        }),
517                        super::super::RunJsonOptions::default(),
518                    ))
519            })
520        };
521
522        let alpha = spawn_run(alpha_path, alpha_bytes.clone());
523        let beta = spawn_run(beta_path, beta_bytes.clone());
524        assert_eq!(alpha.join().expect("alpha thread").exit_code, 0);
525        assert_eq!(beta.join().expect("beta thread").exit_code, 0);
526
527        let parse = |bytes: &Arc<Mutex<Vec<u8>>>| {
528            String::from_utf8(bytes.lock().unwrap().clone())
529                .expect("utf8")
530                .lines()
531                .map(|line| serde_json::from_str::<serde_json::Value>(line).expect("valid event"))
532                .collect::<Vec<_>>()
533        };
534        let alpha_events = parse(&alpha_bytes);
535        let beta_events = parse(&beta_bytes);
536
537        let assert_stream = |events: &[serde_json::Value], stdout: &str, stderr: &str| {
538            assert_eq!(
539                events
540                    .iter()
541                    .map(|event| event["data"]["seq"].as_u64().expect("seq"))
542                    .collect::<Vec<_>>(),
543                [1, 2, 3]
544            );
545            assert_eq!(events[0]["data"]["event_type"], "stdout");
546            assert_eq!(events[0]["data"]["payload"], stdout);
547            assert_eq!(events[1]["data"]["event_type"], "stderr");
548            assert_eq!(events[1]["data"]["payload"], stderr);
549            assert_eq!(events[2]["data"]["event_type"], "result");
550        };
551        assert_stream(&alpha_events, "alpha-out\n", "alpha-err\n");
552        assert_stream(&beta_events, "beta-out\n", "beta-err\n");
553    }
554}