diff --git a/Cargo.lock b/Cargo.lock index 85ed3184..4c66901e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -630,6 +630,7 @@ name = "construct-adapter-common" version = "0.17.4" dependencies = [ "construct-protocol", + "serde_json", "tempfile", "tokio", ] diff --git a/crates/adapter-claude/src/lib.rs b/crates/adapter-claude/src/lib.rs index b0850a93..427dacc4 100644 --- a/crates/adapter-claude/src/lib.rs +++ b/crates/adapter-claude/src/lib.rs @@ -25,7 +25,8 @@ use construct_adapter_common::context_breakdown::{ estimate_tokens_from_chars, BreakdownGate, FixedOverheadPin, }; use construct_adapter_common::{ - claude_transcript_path, drive_turn, next_native_seq, spawn_stderr_log, TurnOutcome, + claude_transcript_path, drive_turn, emit_launch_failure_if_silent, next_native_seq, + spawn_stderr_tail, TurnOutcome, }; use construct_protocol::adapter::pty::{run_session as run_pty, PtySpec}; use construct_protocol::adapter::{ @@ -1110,6 +1111,7 @@ async fn run_session(params: SessionStartParams, ctx: AdapterContext) { } command.env("CONSTRUCT_SESSION_ID", &agentd_session_id); + let events_before_spawn = emit.events_emitted(); let mut child = match command.spawn() { Ok(c) => c, Err(e) => { @@ -1129,7 +1131,7 @@ async fn run_session(params: SessionStartParams, ctx: AdapterContext) { // Write the user message, then close stdin so claude knows we're done. let writer_task = spawn_writer(child_stdin, user_text.clone()); - let stderr_task = spawn_stderr_log(child_stderr, emit.clone()); + let (stderr_task, stderr_tail) = spawn_stderr_tail(child_stderr, emit.clone()); let captured_sid = Arc::new(StdMutex::new(None::)); let parser_task = spawn_parser(child_stdout, emit.clone(), captured_sid.clone()); @@ -1140,7 +1142,15 @@ async fn run_session(params: SessionStartParams, ctx: AdapterContext) { let _ = parser_task.await; let _ = stderr_task.await; // Make sure the child is fully reaped. - let _ = child.wait().await; + let exit_status = child.wait().await.ok(); + if matches!(outcome, TurnOutcome::Completed) { + emit_launch_failure_if_silent( + &emit, + events_before_spawn, + exit_status.as_ref(), + &stderr_tail.snapshot(), + ); + } // Always adopt the latest native id. A turn can report a new id (e.g. // after the interactive equivalent of a context reset); subsequent diff --git a/crates/adapter-codex/src/lib.rs b/crates/adapter-codex/src/lib.rs index 054c2361..1207ba85 100644 --- a/crates/adapter-codex/src/lib.rs +++ b/crates/adapter-codex/src/lib.rs @@ -21,7 +21,8 @@ use construct_adapter_common::context_breakdown::{ estimate_tokens_from_chars, BreakdownGate, FixedOverheadPin, }; use construct_adapter_common::{ - codex_sessions_root, drive_turn, next_native_seq, short, TurnOutcome, + codex_sessions_root, drive_turn, emit_launch_failure_if_silent, next_native_seq, short, + StderrTail, TurnOutcome, }; use construct_protocol::adapter::pty::{run_session as run_pty, PtySpec}; use construct_protocol::adapter::{ @@ -1337,6 +1338,7 @@ async fn run_session(params: SessionStartParams, ctx: AdapterContext) { } command.env("CONSTRUCT_SESSION_ID", &agentd_session_id); + let events_before_spawn = emit.events_emitted(); let mut child = match command.spawn() { Ok(c) => c, Err(e) => { @@ -1353,14 +1355,20 @@ async fn run_session(params: SessionStartParams, ctx: AdapterContext) { let child_stdout = child.stdout.take().expect("piped"); let child_stderr = child.stderr.take().expect("piped"); let diagnostics = Arc::new(StdMutex::new(HeadlessTurnDiagnostics::default())); + let stderr_tail = StderrTail::default(); let stdout_task = spawn_stdout(child_stdout, emit.clone(), diagnostics.clone()); - let stderr_task = spawn_headless_stderr(child_stderr, emit.clone(), diagnostics.clone()); + let stderr_task = spawn_headless_stderr( + child_stderr, + emit.clone(), + diagnostics.clone(), + stderr_tail.clone(), + ); let outcome = drive_turn(&mut child, &mut inbox, &emit, &mut pending).await; let _ = stdout_task.await; let stderr_error = stderr_task.await.ok().flatten(); - let _ = child.wait().await; + let exit_status = child.wait().await.ok(); let diagnostics = diagnostics.lock().unwrap().clone(); @@ -1369,6 +1377,15 @@ async fn run_session(params: SessionStartParams, ctx: AdapterContext) { break 1; } + if matches!(outcome, TurnOutcome::Completed) { + emit_launch_failure_if_silent( + &emit, + events_before_spawn, + exit_status.as_ref(), + &stderr_tail.snapshot(), + ); + } + // Always adopt the latest native id so a mid-run reset is honored // on subsequent turns (and written for daemon resume). if let Some(sid) = diagnostics.session_id.clone() { @@ -1412,6 +1429,7 @@ fn spawn_headless_stderr( reader: R, emit: EventEmitter, diagnostics: Arc>, + tail: StderrTail, ) -> tokio::task::JoinHandle> where R: tokio::io::AsyncRead + Unpin + Send + 'static, @@ -1432,6 +1450,7 @@ where } } emit.log(format!("stderr: {line}")); + tail.push(line); } error }) @@ -1856,7 +1875,12 @@ tokens used let (emit, _rx) = EventEmitter::channel("session"); let diagnostics = Arc::new(StdMutex::new(HeadlessTurnDiagnostics::default())); - spawn_headless_stderr(REAL_REFUSAL_STDERR.as_bytes(), emit, diagnostics.clone()) + spawn_headless_stderr( + REAL_REFUSAL_STDERR.as_bytes(), + emit, + diagnostics.clone(), + StderrTail::default(), + ) .await .expect("stderr task should finish"); @@ -1876,7 +1900,12 @@ tokens used let (emit, _rx) = EventEmitter::channel("session"); let diagnostics = Arc::new(StdMutex::new(HeadlessTurnDiagnostics::default())); - spawn_headless_stderr(OUTPUT.as_bytes(), emit, diagnostics.clone()) + spawn_headless_stderr( + OUTPUT.as_bytes(), + emit, + diagnostics.clone(), + StderrTail::default(), + ) .await .expect("stderr task should finish"); diff --git a/crates/adapter-common/Cargo.toml b/crates/adapter-common/Cargo.toml index 68d2f3d0..6fa0e81a 100644 --- a/crates/adapter-common/Cargo.toml +++ b/crates/adapter-common/Cargo.toml @@ -17,5 +17,6 @@ construct-protocol = { workspace = true } tokio.workspace = true [dev-dependencies] +serde_json.workspace = true tempfile = "3" diff --git a/crates/adapter-common/src/lib.rs b/crates/adapter-common/src/lib.rs index 8ebcfe19..0130b9e2 100644 --- a/crates/adapter-common/src/lib.rs +++ b/crates/adapter-common/src/lib.rs @@ -1,6 +1,8 @@ use std::collections::VecDeque; +use std::sync::{Arc, Mutex}; use construct_protocol::adapter::{AdapterInboxMsg, EventEmitter}; +use construct_protocol::SessionEvent; use tokio::io::{AsyncBufReadExt, AsyncRead, BufReader}; use tokio::sync::mpsc; @@ -80,6 +82,104 @@ where }) } +/// How many trailing stderr lines [`StderrTail`] retains for launch-failure +/// reporting. A fatal CLI/config error is almost always in the last few +/// lines; the bound keeps a chatty harness from growing the buffer unbounded. +const STDERR_TAIL_LINES: usize = 20; + +/// Bounded tail of a harness child's stderr, shared between the stderr +/// logging task and the adapter's turn loop so a silent launch failure can be +/// surfaced with its actual cause instead of only landing in the daemon log. +#[derive(Clone, Default)] +pub struct StderrTail { + lines: Arc>>, +} + +impl StderrTail { + pub fn push(&self, line: String) { + let mut lines = self.lines.lock().unwrap(); + if lines.len() >= STDERR_TAIL_LINES { + lines.pop_front(); + } + lines.push_back(line); + } + + pub fn snapshot(&self) -> Vec { + self.lines.lock().unwrap().iter().cloned().collect() + } +} + +/// Like [`spawn_stderr_log`], but also retains the last few stderr lines so +/// the turn loop can report them if the child turns out to have died at +/// launch. +pub fn spawn_stderr_tail( + reader: R, + emit: EventEmitter, +) -> (tokio::task::JoinHandle<()>, StderrTail) +where + R: AsyncRead + Unpin + Send + 'static, +{ + let tail = StderrTail::default(); + let task_tail = tail.clone(); + let task = tokio::spawn(async move { + let mut lines = BufReader::new(reader).lines(); + while let Ok(Some(line)) = lines.next_line().await { + emit.log(format!("stderr: {line}")); + task_tail.push(line); + } + }); + (task, tail) +} + +/// Human-readable message for a harness child that exited abnormally without +/// producing any session output. +pub fn launch_failure_message( + status: Option<&std::process::ExitStatus>, + stderr_tail: &[String], +) -> String { + let status = match status { + Some(s) => s.to_string(), + None => "unknown exit status".to_string(), + }; + let mut message = format!("harness exited during launch ({status}) without producing output"); + if stderr_tail.is_empty() { + message.push_str("; no stderr captured (see daemon log)"); + } else { + message.push_str(":\n"); + message.push_str(&stderr_tail.join("\n")); + } + message +} + +/// Surface a harness child that failed at launch as a session error. +/// +/// A wrapper harness handed a flag or config it rejects exits non-zero +/// before emitting a single session event; without this the session just +/// flips back to awaiting-input with an empty transcript and looks hung, +/// while the real error sits in the adapter's stderr log. Call this after +/// the turn's child has been reaped, with the emitter's +/// [`EventEmitter::events_emitted`] count snapshotted just before spawn: +/// if the child failed and nothing was emitted since, this emits a +/// [`SessionEvent::Error`] carrying the captured stderr tail. Returns +/// whether it emitted. Only call it for turns that ran to completion — +/// an interrupted or stopped turn is killed by us and would misreport +/// the kill as a launch failure. +pub fn emit_launch_failure_if_silent( + emit: &EventEmitter, + events_before_spawn: u64, + status: Option<&std::process::ExitStatus>, + stderr_tail: &[String], +) -> bool { + let failed = status.map(|s| !s.success()).unwrap_or(false); + if !failed || emit.events_emitted() > events_before_spawn { + return false; + } + emit.emit(SessionEvent::Error { + message: launch_failure_message(status, stderr_tail), + }); + true +} + /// Post-incrementing counter for native-subagent emission ordinals /// (`SessionEvent::NativeSubagent::seq`): returns the current ordinal and /// advances it. Adapters number every emission derived from a child's own @@ -92,3 +192,146 @@ pub fn next_native_seq(ord: &mut u64) -> u64 { *ord += 1; v } + +#[cfg(test)] +mod launch_failure_tests { + use super::*; + use std::process::Stdio; + use tokio::process::Command; + + /// Run a shell one-liner the way headless adapters run a harness child: + /// stderr tailed, turn driven to completion, exit status captured. + async fn run_child( + script: &str, + emit: &EventEmitter, + ) -> (TurnOutcome, Option, StderrTail) { + let mut child = Command::new("sh") + .arg("-c") + .arg(script) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::piped()) + .spawn() + .expect("spawn sh"); + let (stderr_task, tail) = + spawn_stderr_tail(child.stderr.take().expect("piped"), emit.clone()); + // Keep the sender alive so the inbox pends instead of signalling Stop. + let (_inbox_tx, mut inbox) = mpsc::channel::(8); + let mut pending = VecDeque::new(); + let outcome = drive_turn(&mut child, &mut inbox, emit, &mut pending).await; + let _ = stderr_task.await; + let status = child.wait().await.ok(); + (outcome, status, tail) + } + + fn error_messages(rx: &mut mpsc::UnboundedReceiver) -> Vec { + let mut errors = Vec::new(); + while let Ok(v) = rx.try_recv() { + let Some(event) = v.pointer("/params/event") else { + continue; + }; + if event.get("type").and_then(|t| t.as_str()) == Some("error") { + let msg = event + .get("message") + .and_then(|m| m.as_str()) + .unwrap_or_default(); + errors.push(msg.to_string()); + } + } + errors + } + + /// Regression test for construct-worlds/construct#1208: a harness child + /// that rejects its CLI flags and dies before emitting anything used to + /// leave the transcript empty (the error only reached the daemon log via + /// `emit.log`), so the session looked hung at awaiting-input. The failure + /// must surface as a `SessionEvent::Error` carrying the child's stderr. + #[tokio::test] + async fn silent_launch_failure_surfaces_stderr_as_session_error() { + let (emit, mut rx) = EventEmitter::channel("session"); + let baseline = emit.events_emitted(); + let (outcome, status, tail) = run_child( + "echo 'Error: --allow \"MultiEdit(...)\": unknown tool prefix: MultiEdit' >&2; exit 2", + &emit, + ) + .await; + + assert!(matches!(outcome, TurnOutcome::Completed)); + assert!(emit_launch_failure_if_silent( + &emit, + baseline, + status.as_ref(), + &tail.snapshot(), + )); + + let errors = error_messages(&mut rx); + assert_eq!(errors.len(), 1, "exactly one error event, got {errors:?}"); + assert!( + errors[0].contains("unknown tool prefix: MultiEdit"), + "error must carry the child's stderr: {}", + errors[0] + ); + assert!( + errors[0].contains("harness exited during launch"), + "error must name the failure mode: {}", + errors[0] + ); + } + + #[tokio::test] + async fn successful_exit_is_not_a_launch_failure() { + let (emit, mut rx) = EventEmitter::channel("session"); + let baseline = emit.events_emitted(); + let (_outcome, status, tail) = run_child("exit 0", &emit).await; + + assert!(!emit_launch_failure_if_silent( + &emit, + baseline, + status.as_ref(), + &tail.snapshot(), + )); + assert!(error_messages(&mut rx).is_empty()); + } + + #[tokio::test] + async fn failure_after_real_output_is_not_a_launch_failure() { + // A harness that produced session output and then exited non-zero is + // a failed turn, not a launch failure — the transcript already shows + // what happened, so no synthetic error is added. + let (emit, mut rx) = EventEmitter::channel("session"); + let baseline = emit.events_emitted(); + emit.emit(SessionEvent::Message { + role: construct_protocol::MessageRole::Assistant, + text: "partial output".into(), + }); + let (_outcome, status, tail) = run_child("echo boom >&2; exit 1", &emit).await; + + assert!(!emit_launch_failure_if_silent( + &emit, + baseline, + status.as_ref(), + &tail.snapshot(), + )); + assert!(error_messages(&mut rx) + .iter() + .all(|m| !m.contains("harness exited during launch"))); + } + + #[test] + fn launch_failure_message_without_stderr_points_at_the_daemon_log() { + let message = launch_failure_message(None, &[]); + assert!(message.contains("no stderr captured")); + assert!(message.contains("daemon log")); + } + + #[test] + fn stderr_tail_is_bounded_and_keeps_the_last_lines() { + let tail = StderrTail::default(); + for i in 0..100 { + tail.push(format!("line {i}")); + } + let snapshot = tail.snapshot(); + assert_eq!(snapshot.len(), STDERR_TAIL_LINES); + assert_eq!(snapshot.last().map(String::as_str), Some("line 99")); + } +} diff --git a/crates/adapter-grok/src/lib.rs b/crates/adapter-grok/src/lib.rs index 7b5a95ae..6d310b4c 100644 --- a/crates/adapter-grok/src/lib.rs +++ b/crates/adapter-grok/src/lib.rs @@ -14,7 +14,7 @@ use construct_adapter_common::{ context_breakdown::{estimate_tokens_from_chars, BreakdownGate, FixedOverheadPin}, - drive_turn, next_native_seq, spawn_stderr_log, TurnOutcome, + drive_turn, emit_launch_failure_if_silent, next_native_seq, spawn_stderr_tail, TurnOutcome, }; use construct_protocol::adapter::pty::{run_session as run_pty, PtySpec}; use construct_protocol::adapter::{ @@ -1463,6 +1463,7 @@ async fn run_session(params: SessionStartParams, ctx: AdapterContext) { } command.env("CONSTRUCT_SESSION_ID", &agentd_session_id); + let events_before_spawn = emit.events_emitted(); let mut child = match command.spawn() { Ok(c) => c, Err(e) => { @@ -1479,7 +1480,7 @@ async fn run_session(params: SessionStartParams, ctx: AdapterContext) { let child_stdout = child.stdout.take().expect("piped"); let child_stderr = child.stderr.take().expect("piped"); - let stderr_task = spawn_stderr_log(child_stderr, emit.clone()); + let (stderr_task, stderr_tail) = spawn_stderr_tail(child_stderr, emit.clone()); let captured_sid = Arc::new(StdMutex::new(None::)); let parser_task = spawn_parser(child_stdout, emit.clone(), captured_sid.clone()); @@ -1487,7 +1488,7 @@ async fn run_session(params: SessionStartParams, ctx: AdapterContext) { let _ = parser_task.await; let _ = stderr_task.await; - let _ = child.wait().await; + let exit_status = child.wait().await.ok(); // Always adopt the latest native id so a mid-run reset is honored // on subsequent turns (and written for daemon resume). @@ -1499,7 +1500,15 @@ async fn run_session(params: SessionStartParams, ctx: AdapterContext) { } match outcome { - TurnOutcome::Completed => continue, + TurnOutcome::Completed => { + emit_launch_failure_if_silent( + &emit, + events_before_spawn, + exit_status.as_ref(), + &stderr_tail.snapshot(), + ); + continue; + } TurnOutcome::Interrupted => { emit.log("turn interrupted; awaiting next input"); continue; diff --git a/crates/adapter-hermes/src/lib.rs b/crates/adapter-hermes/src/lib.rs index 48719003..cad76427 100644 --- a/crates/adapter-hermes/src/lib.rs +++ b/crates/adapter-hermes/src/lib.rs @@ -13,7 +13,9 @@ //! location at `~/.local/bin/hermes`. `CONSTRUCT_HERMES_HOME` can point the //! adapter and child at a non-default Hermes home. -use construct_adapter_common::{drive_turn, spawn_stderr_log, TurnOutcome}; +use construct_adapter_common::{ + drive_turn, emit_launch_failure_if_silent, spawn_stderr_tail, StderrTail, TurnOutcome, +}; use construct_protocol::adapter::pty::{run_session as run_pty, PtySpec}; use construct_protocol::adapter::{ run as adapter_run, AdapterContext, AdapterInboxMsg, EventEmitter, @@ -335,8 +337,10 @@ async fn run_headless(params: SessionStartParams, ctx: AdapterContext) { break 127; }; let stdout = child.stdout.take(); + let mut stderr_tail = StderrTail::default(); if let Some(stderr) = child.stderr.take() { - spawn_stderr_log(stderr, emit.clone()); + let (_task, tail) = spawn_stderr_tail(stderr, emit.clone()); + stderr_tail = tail; } let stdout_task = tokio::spawn(async move { let mut bytes = Vec::new(); @@ -350,6 +354,7 @@ async fn run_headless(params: SessionStartParams, ctx: AdapterContext) { state: SessionState::Running, detail: Some(format!("{} --oneshot", command.argv_preview())), }); + let events_before_turn = emit.events_emitted(); let outcome = drive_turn(&mut child, &mut inbox, &emit, &mut pending).await; let output = stdout_task.await.unwrap_or_default(); match outcome { @@ -361,6 +366,13 @@ async fn run_headless(params: SessionStartParams, ctx: AdapterContext) { let message = String::from_utf8_lossy(&output).trim().to_string(); if !message.is_empty() { emit.emit(SessionEvent::Error { message }); + } else { + emit_launch_failure_if_silent( + &emit, + events_before_turn, + Some(&status), + &stderr_tail.snapshot(), + ); } } } diff --git a/crates/adapter-muse/src/lib.rs b/crates/adapter-muse/src/lib.rs index 3bf2fedc..73e8ba1a 100644 --- a/crates/adapter-muse/src/lib.rs +++ b/crates/adapter-muse/src/lib.rs @@ -13,7 +13,9 @@ use std::process::Stdio; use std::sync::{Arc, Mutex as StdMutex}; use std::time::{Duration, SystemTime}; -use construct_adapter_common::{drive_turn, spawn_stderr_log, TurnOutcome}; +use construct_adapter_common::{ + drive_turn, emit_launch_failure_if_silent, spawn_stderr_tail, TurnOutcome, +}; use construct_protocol::adapter::pty::{run_session_with_pid as run_pty_with_pid, PtySpec}; use construct_protocol::adapter::{ run as adapter_run, AdapterContext, AdapterInboxMsg, EventEmitter, @@ -280,11 +282,13 @@ async fn run_headless(params: SessionStartParams, mut ctx: AdapterContext) { ctx.emit.clone(), meta.clone(), ); - let stderr_task = spawn_stderr_log(child.stderr.take().expect("piped"), ctx.emit.clone()); + let (stderr_task, stderr_tail) = + spawn_stderr_tail(child.stderr.take().expect("piped"), ctx.emit.clone()); ctx.emit.emit(SessionEvent::Status { state: SessionState::Running, detail: Some(format!("{} exec", command.argv_preview())), }); + let events_before_turn = ctx.emit.events_emitted(); let outcome = drive_turn(&mut child, &mut ctx.inbox, &ctx.emit, &mut pending).await; let _ = stdout_task.await; @@ -296,10 +300,18 @@ async fn run_headless(params: SessionStartParams, mut ctx: AdapterContext) { ctx.emit.log("muse turn interrupted; awaiting next input"); } TurnOutcome::Completed => { - if let Some(status) = status.filter(|status| !status.success()) { - ctx.emit.emit(SessionEvent::Error { - message: format!("muse exec exited with {status}"), - }); + let reported = emit_launch_failure_if_silent( + &ctx.emit, + events_before_turn, + status.as_ref(), + &stderr_tail.snapshot(), + ); + if !reported { + if let Some(status) = status.filter(|status| !status.success()) { + ctx.emit.emit(SessionEvent::Error { + message: format!("muse exec exited with {status}"), + }); + } } } } diff --git a/crates/adapter-pi/src/lib.rs b/crates/adapter-pi/src/lib.rs index 08f923ca..3352b505 100644 --- a/crates/adapter-pi/src/lib.rs +++ b/crates/adapter-pi/src/lib.rs @@ -42,7 +42,9 @@ use std::time::Duration; use construct_adapter_common::context_breakdown::{ estimate_tokens_from_chars, BreakdownGate, FixedOverheadPin, }; -use construct_adapter_common::{drive_turn, spawn_stderr_log, TurnOutcome}; +use construct_adapter_common::{ + drive_turn, emit_launch_failure_if_silent, spawn_stderr_tail, TurnOutcome, +}; use construct_protocol::adapter::pty::{run_session as run_pty, PtySpec}; use construct_protocol::adapter::{ run as adapter_run, AdapterContext, AdapterInboxMsg, EventEmitter, @@ -996,6 +998,7 @@ async fn run_headless(params: SessionStartParams, mut ctx: AdapterContext, flavo } cmd.env("CONSTRUCT_SESSION_ID", &ctx.session_id); + let events_before_spawn = emit.events_emitted(); let mut child = match cmd.spawn() { Ok(c) => c, Err(e) => { @@ -1019,13 +1022,21 @@ async fn run_headless(params: SessionStartParams, mut ctx: AdapterContext, flavo captured_sid.clone(), flavor, ); - let stderr_task = spawn_stderr_log(child_stderr, emit.clone()); + let (stderr_task, stderr_tail) = spawn_stderr_tail(child_stderr, emit.clone()); let outcome = drive_turn(&mut child, &mut ctx.inbox, &emit, &mut pending).await; let _ = stdout_task.await; let _ = stderr_task.await; - let _ = child.wait().await; + let exit_status = child.wait().await.ok(); + if matches!(outcome, TurnOutcome::Completed) { + emit_launch_failure_if_silent( + &emit, + events_before_spawn, + exit_status.as_ref(), + &stderr_tail.snapshot(), + ); + } // Adopt the session the turn actually ran as (fresh spawn, fork, or // a continue that re-minted the uuid) so the next turn and a daemon diff --git a/crates/protocol/src/adapter.rs b/crates/protocol/src/adapter.rs index 8021836e..693a47ad 100644 --- a/crates/protocol/src/adapter.rs +++ b/crates/protocol/src/adapter.rs @@ -497,6 +497,7 @@ pub struct AdapterContext { pub struct EventEmitter { out_tx: mpsc::UnboundedSender, session_id: String, + events_emitted: std::sync::Arc, } impl EventEmitter { @@ -512,12 +513,22 @@ impl EventEmitter { Self { out_tx, session_id: session_id.into(), + events_emitted: Default::default(), }, out_rx, ) } + /// Number of [`SessionEvent`]s emitted so far, shared across clones. + /// Adapters snapshot this before spawning a harness child to tell a turn + /// that produced output apart from one that died silently at launch. + pub fn events_emitted(&self) -> u64 { + self.events_emitted.load(std::sync::atomic::Ordering::Relaxed) + } + pub fn emit(&self, event: SessionEvent) { + self.events_emitted + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); let env = EventEnvelope { session_id: self.session_id.clone(), event, @@ -795,6 +806,7 @@ where emit: EventEmitter { out_tx: out_tx.clone(), session_id: params.session_id.clone(), + events_emitted: Default::default(), }, inbox: rx, }; @@ -1043,6 +1055,7 @@ where emit: EventEmitter { out_tx: out_tx.clone(), session_id: params.session_id.clone(), + events_emitted: Default::default(), }, inbox: rx, };