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 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 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
112pub struct NdjsonEmitter {
117 inner: Arc<NdjsonEmitterInner>,
118}
119
120struct NdjsonEmitterInner {
121 seq: AtomicU64,
122 quiet: bool,
123 out: Mutex<Box<dyn Write + Send>>,
127}
128
129impl NdjsonEmitter {
130 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 pub fn sink(&self) -> Arc<dyn RunEventSink> {
146 Arc::new(NdjsonSink {
147 inner: self.inner.clone(),
148 })
149 }
150
151 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 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 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 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 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}