1use 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
17pub const RUN_JSON_SCHEMA_VERSION: u32 = 1;
21
22#[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 },
79 Result {
80 seq: u64,
81 value: serde_json::Value,
82 exit_code: i32,
83 },
84 Error {
85 seq: u64,
86 error: JsonError,
87 },
88}
89
90impl RunEventWire {
91 pub fn seq(&self) -> u64 {
93 match self {
94 Self::Stdout { seq, .. }
95 | Self::Stderr { seq, .. }
96 | Self::Transcript { seq, .. }
97 | Self::ToolCall { seq, .. }
98 | Self::ToolResult { seq, .. }
99 | Self::Hook { seq, .. }
100 | Self::PersonaStage { seq, .. }
101 | Self::PackRun { seq, .. }
102 | Self::Result { seq, .. }
103 | Self::Error { seq, .. } => *seq,
104 }
105 }
106}
107
108pub struct NdjsonEmitter {
113 inner: Arc<NdjsonEmitterInner>,
114}
115
116struct NdjsonEmitterInner {
117 seq: AtomicU64,
118 quiet: bool,
119 out: Mutex<Box<dyn Write + Send>>,
123}
124
125impl NdjsonEmitter {
126 pub fn new(out: Box<dyn Write + Send>, quiet: bool) -> Self {
130 Self {
131 inner: Arc::new(NdjsonEmitterInner {
132 seq: AtomicU64::new(0),
133 quiet,
134 out: Mutex::new(out),
135 }),
136 }
137 }
138
139 pub fn sink(&self) -> Arc<dyn RunEventSink> {
142 Arc::new(NdjsonSink {
143 inner: self.inner.clone(),
144 })
145 }
146
147 fn next_seq(&self) -> u64 {
149 self.inner.seq.fetch_add(1, Ordering::SeqCst) + 1
150 }
151
152 fn write_envelope(inner: &NdjsonEmitterInner, event: RunEventWire) {
153 let envelope = JsonEnvelope::ok(RUN_JSON_SCHEMA_VERSION, event);
154 let line = serde_json::to_string(&envelope)
155 .unwrap_or_else(|_| r#"{"schemaVersion":1,"ok":false}"#.to_string());
156 if let Ok(mut out) = inner.out.lock() {
157 let _ = writeln!(out, "{line}");
158 let _ = out.flush();
159 }
160 }
161
162 pub fn emit_result(&self, value: serde_json::Value, exit_code: i32) {
164 let event = RunEventWire::Result {
165 seq: self.next_seq(),
166 value,
167 exit_code,
168 };
169 Self::write_envelope(&self.inner, event);
170 }
171
172 pub fn emit_error(&self, code: impl Into<String>, message: impl Into<String>) {
175 let event = RunEventWire::Error {
176 seq: self.next_seq(),
177 error: JsonError {
178 code: code.into(),
179 message: message.into(),
180 details: serde_json::Value::Null,
181 },
182 };
183 Self::write_envelope(&self.inner, event);
184 }
185}
186
187struct NdjsonSink {
188 inner: Arc<NdjsonEmitterInner>,
189}
190
191impl RunEventSink for NdjsonSink {
192 fn emit(&self, event: RunEvent) {
193 let seq = self.inner.seq.fetch_add(1, Ordering::SeqCst) + 1;
194 let wire = match event {
195 RunEvent::Stdout { payload } => {
196 if self.inner.quiet {
197 self.inner.seq.fetch_sub(1, Ordering::SeqCst);
199 return;
200 }
201 RunEventWire::Stdout { seq, payload }
202 }
203 RunEvent::Stderr { payload } => {
204 if self.inner.quiet {
205 self.inner.seq.fetch_sub(1, Ordering::SeqCst);
206 return;
207 }
208 RunEventWire::Stderr { seq, payload }
209 }
210 RunEvent::Transcript {
211 agent_id,
212 kind,
213 payload,
214 } => RunEventWire::Transcript {
215 seq,
216 agent_id,
217 kind,
218 payload,
219 },
220 RunEvent::ToolCall {
221 call_id,
222 name,
223 args,
224 started_at,
225 } => RunEventWire::ToolCall {
226 seq,
227 call_id,
228 name,
229 args,
230 started_at,
231 },
232 RunEvent::ToolResult {
233 call_id,
234 ok,
235 result,
236 } => RunEventWire::ToolResult {
237 seq,
238 call_id,
239 ok,
240 result,
241 },
242 RunEvent::Hook {
243 name,
244 phase,
245 payload,
246 } => RunEventWire::Hook {
247 seq,
248 name,
249 phase,
250 payload,
251 },
252 RunEvent::PersonaStage {
253 persona,
254 stage,
255 transition,
256 } => RunEventWire::PersonaStage {
257 seq,
258 persona,
259 stage,
260 transition,
261 },
262 RunEvent::PackRun {
263 bundle_hash,
264 signature_verified,
265 key_id,
266 cache_hit,
267 dry_run_verify,
268 } => RunEventWire::PackRun {
269 seq,
270 bundle_hash,
271 signature_verified,
272 key_id,
273 cache_hit,
274 dry_run_verify,
275 },
276 };
277 NdjsonEmitter::write_envelope(&self.inner, wire);
278 }
279}
280
281#[cfg(test)]
282mod tests {
283 use super::*;
284
285 struct BufWriter(Arc<Mutex<Vec<u8>>>);
286 impl Write for BufWriter {
287 fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
288 self.0.lock().unwrap().extend_from_slice(buf);
289 Ok(buf.len())
290 }
291 fn flush(&mut self) -> std::io::Result<()> {
292 Ok(())
293 }
294 }
295
296 #[test]
297 fn emits_monotonic_seq_across_events() {
298 let _guard = crate::tests::common::run_event_sink_lock::lock_run_event_sink();
299 let buf = Arc::new(Mutex::new(Vec::<u8>::new()));
300 let emitter = NdjsonEmitter::new(Box::new(BufWriter(buf.clone())), false);
301 let sink = emitter.sink();
302 let prior = harn_vm::run_events::install_sink(sink);
303 harn_vm::run_events::emit(RunEvent::Stdout {
304 payload: "hello\n".into(),
305 });
306 harn_vm::run_events::emit(RunEvent::Stderr {
307 payload: "warn\n".into(),
308 });
309 harn_vm::run_events::clear_sink();
310 emitter.emit_result(serde_json::Value::Null, 0);
311 if let Some(prior) = prior {
312 harn_vm::run_events::install_sink(prior);
313 }
314
315 let raw = String::from_utf8(buf.lock().unwrap().clone()).expect("utf8");
316 let lines: Vec<&str> = raw.lines().filter(|line| !line.is_empty()).collect();
317 assert_eq!(lines.len(), 3, "expected 3 NDJSON lines, got:\n{raw}");
318 let seqs: Vec<u64> = lines
319 .iter()
320 .map(|line| {
321 let v: serde_json::Value = serde_json::from_str(line).expect("valid json");
322 v["data"]["seq"].as_u64().expect("seq present")
323 })
324 .collect();
325 assert_eq!(seqs, vec![1, 2, 3]);
326 let types: Vec<String> = lines
327 .iter()
328 .map(|line| {
329 let v: serde_json::Value = serde_json::from_str(line).expect("valid json");
330 v["data"]["event_type"].as_str().expect("type").to_string()
331 })
332 .collect();
333 assert_eq!(types, vec!["stdout", "stderr", "result"]);
334 }
335
336 #[test]
337 fn quiet_drops_stdout_and_stderr_without_gaps() {
338 let _guard = crate::tests::common::run_event_sink_lock::lock_run_event_sink();
339 let buf = Arc::new(Mutex::new(Vec::<u8>::new()));
340 let emitter = NdjsonEmitter::new(Box::new(BufWriter(buf.clone())), true);
341 let sink = emitter.sink();
342 let prior = harn_vm::run_events::install_sink(sink);
343 harn_vm::run_events::emit(RunEvent::Stdout {
344 payload: "ignored\n".into(),
345 });
346 harn_vm::run_events::emit(RunEvent::Hook {
347 name: "PreRun".into(),
348 phase: "allow".into(),
349 payload: serde_json::Value::Null,
350 });
351 harn_vm::run_events::clear_sink();
352 emitter.emit_result(serde_json::Value::Null, 0);
353 if let Some(prior) = prior {
354 harn_vm::run_events::install_sink(prior);
355 }
356
357 let raw = String::from_utf8(buf.lock().unwrap().clone()).expect("utf8");
358 let lines: Vec<&str> = raw.lines().filter(|line| !line.is_empty()).collect();
359 assert_eq!(lines.len(), 2, "raw:\n{raw}");
361 let seqs: Vec<u64> = lines
362 .iter()
363 .map(|line| {
364 let v: serde_json::Value = serde_json::from_str(line).expect("valid json");
365 v["data"]["seq"].as_u64().expect("seq")
366 })
367 .collect();
368 assert_eq!(
369 seqs,
370 vec![1, 2],
371 "seq must stay contiguous after quiet filtering"
372 );
373 }
374
375 #[tokio::test(flavor = "current_thread")]
380 async fn explicit_exit_emits_stdio_and_one_result_event() {
381 let _guard = crate::tests::common::run_event_sink_lock::lock_run_event_sink_async().await;
382 harn_vm::reset_thread_local_state();
383 let temp = tempfile::TempDir::new().expect("temp dir");
384 let script = temp.path().join("main.harn");
385 std::fs::write(
386 &script,
387 r#"
388fn main(harness: Harness) {
389 harness.stdio.print("before ")
390 harness.stdio.println("exit")
391 harness.stdio.eprintln("diagnostic")
392 exit(2)
393}
394"#,
395 )
396 .expect("write script");
397 let buffer = Arc::new(Mutex::new(Vec::<u8>::new()));
398
399 let outcome = super::super::execute_run_json(
400 &script.to_string_lossy(),
401 false,
402 std::collections::HashSet::new(),
403 Vec::new(),
404 Vec::new(),
405 super::super::CliLlmMockMode::Off,
406 None,
407 super::super::RunProfileOptions::default(),
408 Box::new(BufWriter(buffer.clone())),
409 super::super::RunJsonOptions::default(),
410 )
411 .await;
412
413 assert_eq!(outcome.exit_code, 2, "stderr:\n{}", outcome.stderr);
414 let events: Vec<serde_json::Value> = String::from_utf8(buffer.lock().unwrap().clone())
415 .expect("utf8")
416 .lines()
417 .filter(|line| !line.is_empty())
418 .map(|line| serde_json::from_str(line).expect("valid NDJSON event"))
419 .collect();
420 let stdout = events
421 .iter()
422 .filter(|event| event["data"]["event_type"] == "stdout")
423 .map(|event| event["data"]["payload"].as_str().expect("stdout payload"))
424 .collect::<String>();
425 let stderr = events
426 .iter()
427 .filter(|event| event["data"]["event_type"] == "stderr")
428 .map(|event| event["data"]["payload"].as_str().expect("stderr payload"))
429 .collect::<String>();
430 let terminal: Vec<&serde_json::Value> = events
431 .iter()
432 .filter(|event| event["data"]["event_type"] == "result")
433 .collect();
434
435 assert_eq!(stdout, "before exit\n");
436 assert_eq!(stderr, "diagnostic\n");
437 assert_eq!(terminal.len(), 1, "events: {events:#?}");
438 assert_eq!(terminal[0]["data"]["exit_code"], 2);
439 assert!(terminal[0]["data"]["value"].is_null());
440 assert!(events
441 .iter()
442 .all(|event| event["data"]["event_type"] != "error"));
443 harn_vm::reset_thread_local_state();
444 }
445}