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 #[tokio::test]
297 async fn emits_monotonic_seq_across_events() {
298 let buf = Arc::new(Mutex::new(Vec::<u8>::new()));
299 let emitter = NdjsonEmitter::new(Box::new(BufWriter(buf.clone())), false);
300 let sink = emitter.sink();
301 harn_vm::run_events::scope(sink, async {
302 harn_vm::run_events::emit(RunEvent::Stdout {
303 payload: "hello\n".into(),
304 });
305 harn_vm::run_events::emit(RunEvent::Stderr {
306 payload: "warn\n".into(),
307 });
308 })
309 .await;
310 emitter.emit_result(serde_json::Value::Null, 0);
311
312 let raw = String::from_utf8(buf.lock().unwrap().clone()).expect("utf8");
313 let lines: Vec<&str> = raw.lines().filter(|line| !line.is_empty()).collect();
314 assert_eq!(lines.len(), 3, "expected 3 NDJSON lines, got:\n{raw}");
315 let seqs: Vec<u64> = lines
316 .iter()
317 .map(|line| {
318 let v: serde_json::Value = serde_json::from_str(line).expect("valid json");
319 v["data"]["seq"].as_u64().expect("seq present")
320 })
321 .collect();
322 assert_eq!(seqs, vec![1, 2, 3]);
323 let types: Vec<String> = lines
324 .iter()
325 .map(|line| {
326 let v: serde_json::Value = serde_json::from_str(line).expect("valid json");
327 v["data"]["event_type"].as_str().expect("type").to_string()
328 })
329 .collect();
330 assert_eq!(types, vec!["stdout", "stderr", "result"]);
331 }
332
333 #[tokio::test]
334 async fn quiet_drops_stdout_and_stderr_without_gaps() {
335 let buf = Arc::new(Mutex::new(Vec::<u8>::new()));
336 let emitter = NdjsonEmitter::new(Box::new(BufWriter(buf.clone())), true);
337 let sink = emitter.sink();
338 harn_vm::run_events::scope(sink, async {
339 harn_vm::run_events::emit(RunEvent::Stdout {
340 payload: "ignored\n".into(),
341 });
342 harn_vm::run_events::emit(RunEvent::Hook {
343 name: "PreRun".into(),
344 phase: "allow".into(),
345 payload: serde_json::Value::Null,
346 });
347 })
348 .await;
349 emitter.emit_result(serde_json::Value::Null, 0);
350
351 let raw = String::from_utf8(buf.lock().unwrap().clone()).expect("utf8");
352 let lines: Vec<&str> = raw.lines().filter(|line| !line.is_empty()).collect();
353 assert_eq!(lines.len(), 2, "raw:\n{raw}");
355 let seqs: Vec<u64> = lines
356 .iter()
357 .map(|line| {
358 let v: serde_json::Value = serde_json::from_str(line).expect("valid json");
359 v["data"]["seq"].as_u64().expect("seq")
360 })
361 .collect();
362 assert_eq!(
363 seqs,
364 vec![1, 2],
365 "seq must stay contiguous after quiet filtering"
366 );
367 }
368
369 #[tokio::test(flavor = "current_thread")]
370 async fn explicit_exit_emits_stdio_and_one_result_event() {
371 harn_vm::reset_thread_local_state();
372 let temp = tempfile::TempDir::new().expect("temp dir");
373 let script = temp.path().join("main.harn");
374 std::fs::write(
375 &script,
376 r#"
377fn main(harness: Harness) {
378 harness.stdio.print("before ")
379 harness.stdio.println("exit")
380 harness.stdio.eprintln("diagnostic")
381 exit(2)
382}
383"#,
384 )
385 .expect("write script");
386 let buffer = Arc::new(Mutex::new(Vec::<u8>::new()));
387
388 let outcome = super::super::execute_run_json(
389 &script.to_string_lossy(),
390 false,
391 std::collections::HashSet::new(),
392 Vec::new(),
393 Vec::new(),
394 super::super::CliLlmMockMode::Off,
395 None,
396 super::super::RunProfileOptions::default(),
397 Box::new(BufWriter(buffer.clone())),
398 super::super::RunJsonOptions::default(),
399 )
400 .await;
401
402 assert_eq!(outcome.exit_code, 2, "stderr:\n{}", outcome.stderr);
403 let events: Vec<serde_json::Value> = String::from_utf8(buffer.lock().unwrap().clone())
404 .expect("utf8")
405 .lines()
406 .filter(|line| !line.is_empty())
407 .map(|line| serde_json::from_str(line).expect("valid NDJSON event"))
408 .collect();
409 let stdout = events
410 .iter()
411 .filter(|event| event["data"]["event_type"] == "stdout")
412 .map(|event| event["data"]["payload"].as_str().expect("stdout payload"))
413 .collect::<String>();
414 let stderr = events
415 .iter()
416 .filter(|event| event["data"]["event_type"] == "stderr")
417 .map(|event| event["data"]["payload"].as_str().expect("stderr payload"))
418 .collect::<String>();
419 let terminal: Vec<&serde_json::Value> = events
420 .iter()
421 .filter(|event| event["data"]["event_type"] == "result")
422 .collect();
423
424 assert_eq!(stdout, "before exit\n");
425 assert_eq!(stderr, "diagnostic\n");
426 assert_eq!(terminal.len(), 1, "events: {events:#?}");
427 assert_eq!(terminal[0]["data"]["exit_code"], 2);
428 assert!(terminal[0]["data"]["value"].is_null());
429 assert!(events
430 .iter()
431 .all(|event| event["data"]["event_type"] != "error"));
432 harn_vm::reset_thread_local_state();
433 }
434
435 struct BarrierWriter {
436 bytes: Arc<Mutex<Vec<u8>>>,
437 first_event_barrier: Arc<std::sync::Barrier>,
438 reached_barrier: bool,
439 }
440
441 impl Write for BarrierWriter {
442 fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
443 self.bytes.lock().unwrap().extend_from_slice(buf);
444 if !self.reached_barrier && buf.contains(&b'\n') {
445 self.reached_barrier = true;
446 self.first_event_barrier.wait();
447 }
448 Ok(buf.len())
449 }
450
451 fn flush(&mut self) -> std::io::Result<()> {
452 Ok(())
453 }
454 }
455
456 #[test]
457 fn concurrent_json_runs_receive_only_their_own_ordered_events() {
458 let temp = tempfile::TempDir::new().expect("temp dir");
459 let alpha_path = temp.path().join("alpha.harn");
460 let beta_path = temp.path().join("beta.harn");
461 std::fs::write(
462 &alpha_path,
463 r#"
464fn main(harness: Harness) {
465 harness.stdio.println("alpha-out")
466 harness.stdio.eprintln("alpha-err")
467}
468"#,
469 )
470 .expect("write alpha script");
471 std::fs::write(
472 &beta_path,
473 r#"
474fn main(harness: Harness) {
475 harness.stdio.println("beta-out")
476 harness.stdio.eprintln("beta-err")
477}
478"#,
479 )
480 .expect("write beta script");
481
482 let barrier = Arc::new(std::sync::Barrier::new(2));
483 let alpha_bytes = Arc::new(Mutex::new(Vec::new()));
484 let beta_bytes = Arc::new(Mutex::new(Vec::new()));
485
486 let spawn_run = |path: std::path::PathBuf, bytes: Arc<Mutex<Vec<u8>>>| {
487 let first_event_barrier = barrier.clone();
488 std::thread::spawn(move || {
489 tokio::runtime::Builder::new_current_thread()
490 .enable_all()
491 .build()
492 .expect("runtime")
493 .block_on(super::super::execute_run_json(
494 &path.to_string_lossy(),
495 false,
496 std::collections::HashSet::new(),
497 Vec::new(),
498 Vec::new(),
499 super::super::CliLlmMockMode::Off,
500 None,
501 super::super::RunProfileOptions::default(),
502 Box::new(BarrierWriter {
503 bytes,
504 first_event_barrier,
505 reached_barrier: false,
506 }),
507 super::super::RunJsonOptions::default(),
508 ))
509 })
510 };
511
512 let alpha = spawn_run(alpha_path, alpha_bytes.clone());
513 let beta = spawn_run(beta_path, beta_bytes.clone());
514 assert_eq!(alpha.join().expect("alpha thread").exit_code, 0);
515 assert_eq!(beta.join().expect("beta thread").exit_code, 0);
516
517 let parse = |bytes: &Arc<Mutex<Vec<u8>>>| {
518 String::from_utf8(bytes.lock().unwrap().clone())
519 .expect("utf8")
520 .lines()
521 .map(|line| serde_json::from_str::<serde_json::Value>(line).expect("valid event"))
522 .collect::<Vec<_>>()
523 };
524 let alpha_events = parse(&alpha_bytes);
525 let beta_events = parse(&beta_bytes);
526
527 let assert_stream = |events: &[serde_json::Value], stdout: &str, stderr: &str| {
528 assert_eq!(
529 events
530 .iter()
531 .map(|event| event["data"]["seq"].as_u64().expect("seq"))
532 .collect::<Vec<_>>(),
533 [1, 2, 3]
534 );
535 assert_eq!(events[0]["data"]["event_type"], "stdout");
536 assert_eq!(events[0]["data"]["payload"], stdout);
537 assert_eq!(events[1]["data"]["event_type"], "stderr");
538 assert_eq!(events[1]["data"]["payload"], stderr);
539 assert_eq!(events[2]["data"]["event_type"], "result");
540 };
541 assert_stream(&alpha_events, "alpha-out\n", "alpha-err\n");
542 assert_stream(&beta_events, "beta-out\n", "beta-err\n");
543 }
544}