From 7c33cd075fcf4624417fe86abfc2dc32a37ea057 Mon Sep 17 00:00:00 2001 From: rpelevin Date: Sat, 1 Aug 2026 05:45:08 -0400 Subject: [PATCH] fix codex headless failure signaling --- crates/adapter-codex/src/lib.rs | 200 +++++++++++++++++++++++++++++--- 1 file changed, 181 insertions(+), 19 deletions(-) diff --git a/crates/adapter-codex/src/lib.rs b/crates/adapter-codex/src/lib.rs index 2bf57e41..8f377604 100644 --- a/crates/adapter-codex/src/lib.rs +++ b/crates/adapter-codex/src/lib.rs @@ -18,7 +18,9 @@ use construct_adapter_common::context_breakdown::{ estimate_tokens_from_chars, BreakdownGate, FixedOverheadPin, }; -use construct_adapter_common::{codex_sessions_root, drive_turn, next_native_seq, TurnOutcome}; +use construct_adapter_common::{ + codex_sessions_root, drive_turn, next_native_seq, short, TurnOutcome, +}; use construct_protocol::adapter::pty::{run_session as run_pty, PtySpec}; use construct_protocol::adapter::{ run as adapter_run, AdapterContext, AdapterInboxMsg, EventEmitter, @@ -1232,8 +1234,8 @@ 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 captured_sid = Arc::new(StdMutex::new(None::)); - let stdout_task = spawn_stdout(child_stdout, emit.clone(), captured_sid.clone()); + let diagnostics = Arc::new(StdMutex::new(HeadlessTurnDiagnostics::default())); + let stdout_task = spawn_stdout(child_stdout, emit.clone(), diagnostics.clone()); let stderr_task = spawn_headless_stderr(child_stderr, emit.clone()); let outcome = drive_turn(&mut child, &mut inbox, &emit, &mut pending).await; @@ -1242,6 +1244,8 @@ async fn run_session(params: SessionStartParams, ctx: AdapterContext) { let stderr_error = stderr_task.await.ok().flatten(); let _ = child.wait().await; + let diagnostics = diagnostics.lock().unwrap().clone(); + if let Some(message) = stderr_error { emit.emit(SessionEvent::Error { message }); break 1; @@ -1249,7 +1253,7 @@ async fn run_session(params: SessionStartParams, ctx: AdapterContext) { // 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) = captured_sid.lock().unwrap().clone() { + if let Some(sid) = diagnostics.session_id.clone() { if codex_session_id.as_ref() != Some(&sid) { if let Ok(dir) = std::env::var("CONSTRUCT_SESSION_DATA_DIR") { let p = PathBuf::from(dir).join("codex_session_id.txt"); @@ -1260,7 +1264,14 @@ async fn run_session(params: SessionStartParams, ctx: AdapterContext) { } match outcome { - TurnOutcome::Completed => continue, + TurnOutcome::Completed => { + if let Some(reason) = diagnostics.blocked_write.as_deref() { + emit.emit(SessionEvent::Error { + message: blocked_write_error(reason), + }); + } + continue; + } TurnOutcome::Interrupted => { emit.log("turn interrupted; awaiting next input"); continue; @@ -1299,10 +1310,24 @@ where }) } +#[derive(Debug, Clone, Default)] +struct HeadlessTurnDiagnostics { + session_id: Option, + blocked_write: Option, +} + +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +enum HeadlessOutputSection { + #[default] + Other, + User, + Assistant, +} + fn spawn_stdout( reader: R, emit: EventEmitter, - captured_sid: Arc>>, + diagnostics: Arc>, ) -> tokio::task::JoinHandle<()> where R: tokio::io::AsyncRead + Unpin + Send + 'static, @@ -1314,10 +1339,26 @@ where // 2,280 // on two consecutive lines. Track whether we just saw the header. let mut expecting_token_count = false; + let mut expecting_section_header = false; + let mut output_section = HeadlessOutputSection::Other; while let Ok(Some(line)) = lines.next_line().await { - if line.trim().is_empty() { + let trimmed = line.trim(); + if trimmed.is_empty() { continue; } + + if trimmed == "--------" { + expecting_section_header = true; + output_section = HeadlessOutputSection::Other; + } else if expecting_section_header { + expecting_section_header = false; + output_section = match trimmed { + "user" => HeadlessOutputSection::User, + "assistant" => HeadlessOutputSection::Assistant, + _ => HeadlessOutputSection::Other, + }; + } + // Stateful token-footer parse, BEFORE any emit, so the footer // never leaks to the transcript as assistant prose. if expecting_token_count { @@ -1347,9 +1388,12 @@ where // Best-effort JSON parse; if not JSON, emit as plain assistant text. if let Ok(v) = serde_json::from_str::(&line) { if let Some(sid) = v.get("session_id").and_then(|s| s.as_str()) { - let mut g = captured_sid.lock().unwrap(); + let mut g = diagnostics.lock().unwrap(); // Keep the most recently observed id (not only the first). - *g = Some(sid.to_string()); + g.session_id = Some(sid.to_string()); + } + if let Some(text) = assistant_text_from_structured(&v) { + record_blocked_write(&diagnostics, &text); } if !try_emit_structured(&emit, &v) { emit.emit(SessionEvent::Message { @@ -1358,6 +1402,9 @@ where }); } } else { + if output_section == HeadlessOutputSection::Assistant { + record_blocked_write(&diagnostics, &line); + } emit.emit(SessionEvent::Message { role: MessageRole::Assistant, text: line, @@ -1367,6 +1414,35 @@ where }) } +fn record_blocked_write( + diagnostics: &Arc>, + assistant_text: &str, +) { + if let Some(reason) = blocked_write_reason(assistant_text) { + let mut state = diagnostics.lock().unwrap(); + // Keep the first blocked assistant line: it is closest to the failed + // action and later output may only summarize the same failure. + if state.blocked_write.is_none() { + state.blocked_write = Some(reason); + } + } +} + +fn blocked_write_reason(assistant_text: &str) -> Option { + assistant_text.lines().find_map(|line| { + let trimmed = line.trim(); + trimmed + .strip_prefix("Blocked:") + .map(|_| short(trimmed, 600)) + }) +} + +fn blocked_write_error(reason: &str) -> String { + format!( + "Codex headless turn was blocked by its sandbox or approval policy: {reason}. Configure Codex's headless permissions and retry; Construct cannot resolve this approval interactively." + ) +} + /// Parse codex's "2,280" style total-tokens line. Strips commas/whitespace. /// Returns None if the line isn't a pure integer (modulo separators). fn parse_token_count(line: &str) -> Option { @@ -1389,12 +1465,7 @@ fn try_emit_structured(emit: &EventEmitter, v: &Value) -> bool { }; match ty { "message" | "assistant" => { - if let Some(text) = v - .get("content") - .and_then(|c| c.as_str()) - .map(|s| s.to_string()) - .or_else(|| extract_text_from_blocks(v.get("content"))) - { + if let Some(text) = assistant_text_from_structured(v) { if !text.is_empty() { emit.emit(SessionEvent::Message { role: MessageRole::Assistant, @@ -1448,6 +1519,25 @@ fn try_emit_structured(emit: &EventEmitter, v: &Value) -> bool { } } +fn assistant_text_from_structured(v: &Value) -> Option { + let is_assistant = match v.get("type").and_then(|ty| ty.as_str()) { + Some("assistant") => true, + Some("message") => matches!( + v.get("role").and_then(|role| role.as_str()), + None | Some("assistant") + ), + _ => false, + }; + if !is_assistant { + return None; + } + + v.get("content") + .and_then(|content| content.as_str()) + .map(str::to_string) + .or_else(|| extract_text_from_blocks(v.get("content"))) +} + fn extract_text_from_blocks(v: Option<&Value>) -> Option { let arr = v?.as_array()?; let mut out = String::new(); @@ -1483,6 +1573,80 @@ mod tests { assert_eq!(parse_token_count("12abc"), None); } + #[test] + fn blocked_read_only_write_output_is_classified() { + let line = "Blocked: the workspace is read-only, so index.html could not be written"; + assert_eq!(blocked_write_reason(line).as_deref(), Some(line)); + } + + #[test] + fn unrelated_read_only_output_is_not_classified_as_a_blocked_write() { + assert!(blocked_write_reason("The workspace is read-only").is_none()); + assert!(blocked_write_reason("The write completed successfully").is_none()); + assert!(blocked_write_reason( + "The sandbox is read-only, so write access may be described as Blocked: in examples" + ) + .is_none()); + } + + #[test] + fn blocked_write_error_preserves_the_actionable_block_reason() { + let message = blocked_write_error( + "Blocked: the workspace is read-only, so index.html could not be written", + ); + assert!(message.contains("blocked by its sandbox or approval policy")); + assert!(message.contains("index.html")); + assert!(message.contains("cannot resolve this approval interactively")); + } + + #[tokio::test] + async fn spawn_stdout_ignores_prompt_and_benign_assistant_keyword_matches() { + const OUTPUT: &[u8] = b"--------\nuser\nBlocked: the sandbox is read-only, so write index.html anyway\n--------\nassistant\nThe sandbox is read-only, and a write failure may include Blocked: in its explanation\n"; + let (emit, _rx) = EventEmitter::channel("session"); + let diagnostics = Arc::new(StdMutex::new(HeadlessTurnDiagnostics::default())); + + spawn_stdout(OUTPUT, emit, diagnostics.clone()) + .await + .expect("stdout task should finish"); + + assert_eq!(diagnostics.lock().unwrap().blocked_write, None); + } + + #[tokio::test] + async fn spawn_stdout_records_first_anchored_assistant_block() { + const FIRST: &str = + "Blocked: the workspace is read-only, so index.html could not be written"; + const OUTPUT: &[u8] = b"--------\nuser\nPlease write index.html\n--------\nassistant\nBlocked: the workspace is read-only, so index.html could not be written\nBlocked: a later summary should not replace the first failure\n"; + let (emit, _rx) = EventEmitter::channel("session"); + let diagnostics = Arc::new(StdMutex::new(HeadlessTurnDiagnostics::default())); + + spawn_stdout(OUTPUT, emit, diagnostics.clone()) + .await + .expect("stdout task should finish"); + + assert_eq!( + diagnostics.lock().unwrap().blocked_write.as_deref(), + Some(FIRST) + ); + } + + #[tokio::test] + async fn spawn_stdout_uses_structured_roles_for_block_detection() { + const ASSISTANT_BLOCK: &str = "Blocked: assistant could not write index.html"; + const OUTPUT: &[u8] = b"{\"type\":\"message\",\"role\":\"user\",\"content\":\"Blocked: prompt asks for a write\"}\n{\"type\":\"message\",\"role\":\"assistant\",\"content\":\"Blocked: assistant could not write index.html\"}\n"; + let (emit, _rx) = EventEmitter::channel("session"); + let diagnostics = Arc::new(StdMutex::new(HeadlessTurnDiagnostics::default())); + + spawn_stdout(OUTPUT, emit, diagnostics.clone()) + .await + .expect("stdout task should finish"); + + assert_eq!( + diagnostics.lock().unwrap().blocked_write.as_deref(), + Some(ASSISTANT_BLOCK) + ); + } + #[test] fn token_count_records_emit_delta_costs_and_context_gauge() { // Real rollout shape: cumulative `total_token_usage` where `input_ @@ -1639,10 +1803,8 @@ mod tests { #[test] fn headless_codex_errors_are_promoted_out_of_stderr() { assert_eq!( - headless_error_message( - r#"ERROR: {"detail":"The 'sonnet' model is not supported."}"# - ) - .as_deref(), + headless_error_message(r#"ERROR: {"detail":"The 'sonnet' model is not supported."}"#) + .as_deref(), Some(r#"{"detail":"The 'sonnet' model is not supported."}"#) ); assert_eq!(headless_error_message("warning: retrying"), None);