From 82019d356f97916af05c82bbc8533f636371f367 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 8 May 2026 09:35:53 -0700 Subject: [PATCH] fix(cli): unblock attach on external interview answers Keep attach reading run events while a local interview prompt is active so answers from the web UI or API can resolve the prompt and let the CLI advance. --- Cargo.lock | 1 + lib/crates/fabro-cli/Cargo.toml | 1 + .../fabro-cli/src/commands/run/attach.rs | 313 +++++++++++++++++- lib/crates/fabro-cli/tests/it/cmd/attach.rs | 163 +++++++++ 4 files changed, 467 insertions(+), 11 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 955644a77..31aced142 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1698,6 +1698,7 @@ dependencies = [ "jsonwebtoken", "libc", "miette", + "nix 0.30.1", "object_store", "openssl", "paste", diff --git a/lib/crates/fabro-cli/Cargo.toml b/lib/crates/fabro-cli/Cargo.toml index d727f4aff..864ef3fb4 100644 --- a/lib/crates/fabro-cli/Cargo.toml +++ b/lib/crates/fabro-cli/Cargo.toml @@ -96,6 +96,7 @@ object_store.workspace = true bytes.workspace = true tokio-util.workspace = true libc = "0.2" +nix = { version = "0.30", features = ["fs"] } [target.'cfg(target_os = "macos")'.dependencies] core-foundation = { version = "0.9", optional = true } diff --git a/lib/crates/fabro-cli/src/commands/run/attach.rs b/lib/crates/fabro-cli/src/commands/run/attach.rs index 62af35a32..63fc71004 100644 --- a/lib/crates/fabro-cli/src/commands/run/attach.rs +++ b/lib/crates/fabro-cli/src/commands/run/attach.rs @@ -8,6 +8,8 @@ )] use std::io::{IsTerminal, Write}; +#[cfg(unix)] +use std::os::fd::AsFd; #[cfg(test)] use std::path::Path; #[cfg(test)] @@ -17,17 +19,17 @@ use std::time::Duration; use anyhow::Result; use fabro_api::types; -use fabro_interview::{AnswerValue, ConsoleInterviewer, Question}; +use fabro_interview::{Answer, AnswerValue, Question}; use fabro_store::EventEnvelope; use fabro_types::settings::run::ApprovalMode; -use fabro_types::{EventBody, InterviewOption, RunId}; +use fabro_types::{EventBody, InterviewOption, QuestionType, RunId}; use fabro_util::json::normalize_json_value; use fabro_util::printer::Printer; use fabro_util::terminal::Styles; use fabro_workflow::outcome::StageOutcome; use fabro_workflow::run_status::RunStatus; use tokio::signal::ctrl_c; -use tokio::time::sleep; +use tokio::time::{Duration as TokioDuration, sleep}; use super::run_progress; use crate::server_client; @@ -36,6 +38,87 @@ const INTERVIEW_UNANSWERED_MESSAGE: &str = "Interview ended without an answer. The run is still waiting for input; reattach to answer it."; const JSON_INTERVIEW_MESSAGE: &str = "This run is waiting for human input, but --json is non-interactive. Reattach without --json to answer it."; const ATTACH_PREMATURE_EOF_MESSAGE: &str = "Attach stream ended before terminal run event."; +const PROMPT_READ_POLL_INTERVAL: TokioDuration = TokioDuration::from_millis(50); + +enum PromptRead { + Line(String), + Eof, + Error, +} + +#[cfg(unix)] +use nix::errno::Errno; +#[cfg(unix)] +use nix::fcntl::{FcntlArg, OFlag, fcntl}; +#[cfg(unix)] +use nix::unistd; +#[cfg(unix)] +enum LineRead { + Pending, + Complete(String), + Eof, + Error, +} + +#[cfg(unix)] +struct NonblockingStdin { + stdin: std::io::Stdin, + original_flags: OFlag, +} + +#[cfg(unix)] +impl NonblockingStdin { + fn new() -> Option { + let stdin = std::io::stdin(); + let original_flags = + OFlag::from_bits_truncate(fcntl(stdin.as_fd(), FcntlArg::F_GETFL).ok()?); + fcntl( + stdin.as_fd(), + FcntlArg::F_SETFL(original_flags | OFlag::O_NONBLOCK), + ) + .ok()?; + Some(Self { + stdin, + original_flags, + }) + } + + fn read_line(&self, buffer: &mut Vec) -> LineRead { + let mut chunk = [0_u8; 256]; + loop { + match unistd::read(self.stdin.as_fd(), &mut chunk) { + Ok(0) => { + return if buffer.is_empty() { + LineRead::Eof + } else { + let line = std::mem::take(buffer); + LineRead::Complete(String::from_utf8_lossy(&line).to_string()) + }; + } + Ok(read) => { + buffer.extend_from_slice(&chunk[..read]); + if let Some(newline) = buffer.iter().position(|byte| *byte == b'\n') { + let line = buffer.drain(..=newline).collect::>(); + return LineRead::Complete( + String::from_utf8_lossy(&line) + .trim_end_matches(['\r', '\n']) + .to_string(), + ); + } + } + Err(Errno::EAGAIN) => return LineRead::Pending, + Err(_) => return LineRead::Error, + } + } + } +} + +#[cfg(unix)] +impl Drop for NonblockingStdin { + fn drop(&mut self) { + let _ = fcntl(self.stdin.as_fd(), FcntlArg::F_SETFL(self.original_flags)); + } +} /// Attach to a running (or finished) workflow run, rendering progress live. /// @@ -168,15 +251,17 @@ async fn attach_live_run_with_client( emit_progress_line(&mut progress_ui, &line, opts.json_output)?; } - if let Some(exit_code) = handle_pending_server_interview( + if let Some(exit_code) = Box::pin(handle_pending_server_interview( client, run_id, + &mut stream, opts.auto_approve, &mut progress_ui, styles, opts.json_output, + opts.kill_on_detach, printer, - ) + )) .await? { return Ok(exit_code); @@ -206,15 +291,17 @@ async fn attach_live_run_with_client( } if event_starts_interview(&event) { - if let Some(exit_code) = handle_pending_server_interview( + if let Some(exit_code) = Box::pin(handle_pending_server_interview( client, run_id, + &mut stream, opts.auto_approve, &mut progress_ui, styles, opts.json_output, + opts.kill_on_detach, printer, - ) + )) .await? { return Ok(exit_code); @@ -226,10 +313,12 @@ async fn attach_live_run_with_client( async fn handle_pending_server_interview( client: &server_client::Client, run_id: &RunId, + stream: &mut server_client::RunEventStream, auto_approve: bool, progress_ui: &mut run_progress::ProgressUI, styles: &'static Styles, json_output: bool, + kill_on_detach: bool, printer: Printer, ) -> Result> { let Some(question) = client.list_run_questions(run_id).await?.into_iter().next() else { @@ -245,10 +334,42 @@ async fn handle_pending_server_interview( } hide_progress(progress_ui, json_output); - let interviewer = ConsoleInterviewer::new(styles, fabro_types::Principal::Anonymous); - let submission = - fabro_interview::Interviewer::ask(&interviewer, api_question_to_question(&question)).await; - let answer = submission.answer; + let ask = ask_attach_question(api_question_to_question(&question), styles); + tokio::pin!(ask); + let ctrl_c_signal = ctrl_c(); + tokio::pin!(ctrl_c_signal); + + let answer = loop { + let next_event = tokio::select! { + answer = &mut ask => { + break answer; + } + _ = &mut ctrl_c_signal => { + handle_detach_signal(client, run_id, kill_on_detach, printer).await; + show_progress(progress_ui, json_output); + return Ok(Some(ExitCode::from(1))); + } + result = stream.next_event() => result?, + }; + + let Some(event) = next_event else { + show_progress(progress_ui, json_output); + return Err(anyhow::anyhow!(ATTACH_PREMATURE_EOF_MESSAGE)); + }; + + let line = event_payload_line(&event)?; + emit_progress_line(progress_ui, &line, json_output)?; + + if let Some(exit_code) = event_exit_code(&event) { + show_progress(progress_ui, json_output); + return Ok(Some(exit_code)); + } + + if event_resolves_interview(&event, &question.id) { + show_progress(progress_ui, json_output); + return Ok(None); + } + }; show_progress(progress_ui, json_output); if answer_requires_reattach(&answer) { @@ -307,6 +428,167 @@ fn api_question_to_question(question: &types::ApiQuestion) -> Question { converted } +#[allow( + clippy::print_stderr, + reason = "Interactive questions and options belong on stderr, not captured stdout." +)] +async fn ask_attach_question(question: Question, styles: &'static Styles) -> Answer { + if let Some(ref context_text) = question.context_display { + let rendered = styles.render_markdown(context_text); + eprint!("{rendered}"); + } + eprintln!("{} {}", styles.bold_cyan.apply_to("?"), question.text); + + match question.question_type { + QuestionType::MultipleChoice | QuestionType::MultiSelect => { + for (i, opt) in question.options.iter().enumerate() { + eprintln!( + " {}{}{} {} - {}", + styles.dim.apply_to("["), + styles.bold.apply_to(i + 1), + styles.dim.apply_to("]"), + opt.key, + opt.label, + ); + } + if question.allow_freeform { + eprintln!(" Or type a free-text response"); + } + parse_choice_response(&question, read_attach_line("Select: ").await) + } + QuestionType::YesNo | QuestionType::Confirmation => { + parse_confirm_response(read_attach_line("[Y/N]: ").await) + } + QuestionType::Freeform => parse_freeform_response(read_attach_line("> ").await), + } +} + +#[allow( + clippy::print_stderr, + reason = "Prompts go to stderr so piped stdout stays machine-readable." +)] +async fn read_attach_line(prompt: &str) -> PromptRead { + eprint!("{prompt}"); + let _ = std::io::stderr().flush(); + read_attach_line_after_prompt().await +} + +#[cfg(unix)] +async fn read_attach_line_after_prompt() -> PromptRead { + let Some(stdin) = NonblockingStdin::new() else { + return PromptRead::Error; + }; + let mut buffer = Vec::new(); + loop { + match stdin.read_line(&mut buffer) { + LineRead::Pending => sleep(PROMPT_READ_POLL_INTERVAL).await, + LineRead::Complete(line) => return PromptRead::Line(line), + LineRead::Eof => return PromptRead::Eof, + LineRead::Error => return PromptRead::Error, + } + } +} + +#[cfg(not(unix))] +async fn read_attach_line_after_prompt() -> PromptRead { + use tokio::io::{self, AsyncBufReadExt, BufReader}; + + let stdin = io::stdin(); + let mut reader = BufReader::new(stdin); + let mut line = String::new(); + match reader.read_line(&mut line).await { + Ok(0) => PromptRead::Eof, + Ok(_) => PromptRead::Line(line.trim_end().to_string()), + Err(_) => PromptRead::Error, + } +} + +fn parse_choice_response(question: &Question, prompt_read: PromptRead) -> Answer { + let PromptRead::Line(response) = prompt_read else { + return Answer::interrupted(); + }; + if response.trim().is_empty() { + return Answer::interrupted(); + } + if question.question_type == QuestionType::MultiSelect { + let selected = response + .split([',', ' ']) + .filter(|part| !part.trim().is_empty()) + .map(str::trim) + .map(|part| { + question + .options + .iter() + .find(|option| option.key.eq_ignore_ascii_case(part)) + .map(|option| option.key.clone()) + .or_else(|| { + part.parse::().ok().and_then(|idx| { + idx.checked_sub(1) + .and_then(|zero_idx| question.options.get(zero_idx)) + .map(|option| option.key.clone()) + }) + }) + }) + .collect::>>(); + if let Some(selected) = selected.filter(|keys| !keys.is_empty()) { + return Answer::multi_selected(selected); + } + } + if let Some(answer) = find_matching_option(&response, &question.options) { + return answer; + } + if question.allow_freeform { + return Answer::text(response); + } + Answer::interrupted() +} + +fn parse_confirm_response(prompt_read: PromptRead) -> Answer { + let PromptRead::Line(response) = prompt_read else { + return Answer::interrupted(); + }; + match response.trim().to_lowercase().as_str() { + "y" | "yes" => Answer::yes(), + "n" | "no" => Answer::no(), + _ => Answer::interrupted(), + } +} + +fn parse_freeform_response(prompt_read: PromptRead) -> Answer { + let PromptRead::Line(response) = prompt_read else { + return Answer::interrupted(); + }; + if response.trim().is_empty() { + Answer::interrupted() + } else { + Answer::text(response) + } +} + +fn find_matching_option(response: &str, options: &[InterviewOption]) -> Option { + let trimmed = response.trim(); + for opt in options { + if opt.key.eq_ignore_ascii_case(trimmed) { + return Some(Answer { + value: AnswerValue::Selected(opt.key.clone()), + selected_option: Some(opt.clone()), + text: None, + }); + } + } + if let Ok(idx) = trimmed.parse::() { + if idx >= 1 && idx <= options.len() { + let opt = &options[idx - 1]; + return Some(Answer { + value: AnswerValue::Selected(opt.key.clone()), + selected_option: Some(opt.clone()), + text: None, + }); + } + } + None +} + async fn submit_server_interview_answer( client: &server_client::Client, run_id: &RunId, @@ -481,6 +763,15 @@ fn event_starts_interview(event: &EventEnvelope) -> bool { matches!(event.event.body, EventBody::InterviewStarted(_)) } +fn event_resolves_interview(event: &EventEnvelope, question_id: &str) -> bool { + match &event.event.body { + EventBody::InterviewCompleted(props) => props.question_id == question_id, + EventBody::InterviewInterrupted(props) => props.question_id == question_id, + EventBody::InterviewTimeout(props) => props.question_id == question_id, + _ => false, + } +} + #[cfg(test)] mod tests { #![allow( diff --git a/lib/crates/fabro-cli/tests/it/cmd/attach.rs b/lib/crates/fabro-cli/tests/it/cmd/attach.rs index 37dd18b57..9104592f3 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/attach.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/attach.rs @@ -122,6 +122,30 @@ fn wait_for_output_signal( } } +#[expect( + clippy::disallowed_methods, + reason = "This sync integration helper polls a child process without a Tokio runtime." +)] +fn wait_for_child_exit(child: &mut std::process::Child, label: &str) -> std::process::ExitStatus { + let deadline = Instant::now() + Duration::from_secs(5); + loop { + if let Some(status) = child + .try_wait() + .unwrap_or_else(|err| panic!("{label} status should be readable: {err}")) + { + return status; + } + if Instant::now() >= deadline { + let _ = child.kill(); + let status = child + .wait() + .unwrap_or_else(|err| panic!("{label} should exit after kill: {err}")); + panic!("{label} did not exit before timeout; killed with status {status}"); + } + std::thread::sleep(Duration::from_millis(20)); + } +} + #[test] fn attach_replays_completed_detached_run() { let context = test_context!(); @@ -169,6 +193,145 @@ fn attach_replays_completed_detached_run() { "); } +#[test] +#[expect( + clippy::disallowed_methods, + reason = "This sync integration test keeps a child stdin pipe open to reproduce attach waiting on input while the API answers the same question." +)] +fn attach_advances_when_pending_question_is_answered_elsewhere() { + let context = test_context!(); + context.ensure_home_server_auth_methods(); + let workflow = context.temp_dir.join("human-gate.fabro"); + context.write_temp( + "human-gate.fabro", + r#"digraph HumanGate { + graph [goal="Wait for approval"] + start [shape=Mdiamond, label="Start"] + exit [shape=Msquare, label="Exit"] + approve [shape=hexagon, label="Approve?"] + ship [shape=parallelogram, script="echo shipped"] + start -> approve + approve -> ship [label="[A] Approve"] + ship -> exit +} +"#, + ); + + let run_output = context + .command() + .env("OPENAI_API_KEY", "test") + .args([ + "run", + "--detach", + "--no-retro", + "--sandbox", + "local", + "--provider", + "openai", + workflow.to_str().unwrap(), + ]) + .output() + .expect("detached run should execute"); + assert!( + run_output.status.success(), + "detached run failed:\nstdout:\n{}\nstderr:\n{}", + String::from_utf8_lossy(&run_output.stdout), + String::from_utf8_lossy(&run_output.stderr) + ); + let run_id = output_stdout(&run_output).trim().to_string(); + let cleanup_run_id = run_id.clone(); + scopeguard::defer! { + let _ = context.command().args(["rm", "--force", &cleanup_run_id]).output(); + } + + let runtime = tokio::runtime::Runtime::new().expect("test runtime should build"); + let (client, base_url) = + server_endpoint(&context.storage_dir).expect("server endpoint should exist"); + let question = runtime.block_on(wait_for_server_question(&client, &base_url, &run_id)); + let question_id = question["id"] + .as_str() + .expect("question id should be present") + .to_string(); + + let mut attach_cmd = std::process::Command::new(env!("CARGO_BIN_EXE_fabro")); + fabro_test::apply_test_isolation(&mut attach_cmd, &context.home_dir); + attach_cmd.current_dir(&context.temp_dir); + attach_cmd.args(["attach", &run_id]); + attach_cmd.stdin(Stdio::piped()); + attach_cmd.stdout(Stdio::piped()); + attach_cmd.stderr(Stdio::piped()); + let mut child = attach_cmd.spawn().expect("attach should spawn"); + let _stdin = child.stdin.take().expect("attach stdin should be piped"); + let mut stdout = child.stdout.take().expect("attach stdout should be piped"); + let stderr = child.stderr.take().expect("attach stderr should be piped"); + let (signal_tx, signal_rx) = mpsc::channel(); + let stderr_reader = std::thread::spawn(move || { + let mut reader = BufReader::new(stderr); + let mut stderr_bytes = Vec::new(); + let mut line = Vec::new(); + + loop { + line.clear(); + let read = reader + .read_until(b'\n', &mut line) + .expect("attach stderr should be readable"); + if read == 0 { + break; + } + if line + .windows("Approve?".len()) + .any(|window| window == "Approve?".as_bytes()) + { + let _ = signal_tx.send(()); + } + stderr_bytes.extend_from_slice(&line); + } + + stderr_bytes + }); + let stderr_reader = wait_for_output_signal( + &mut child, + &mut stdout, + stderr_reader, + &signal_rx, + "Approve?", + ); + + runtime.block_on(async { + let response = client + .post(format!( + "{base_url}/api/v1/runs/{run_id}/questions/{question_id}/answer" + )) + .json(&serde_json::json!({ "kind": "selected", "option_key": "A" })) + .send() + .await + .expect("answer submission should succeed"); + assert_reqwest_status( + response, + fabro_http::StatusCode::NO_CONTENT, + format!("POST /api/v1/runs/{run_id}/questions/{question_id}/answer"), + ) + .await; + }); + + let status = wait_for_child_exit(&mut child, "attach"); + let mut stdout_bytes = Vec::new(); + stdout + .read_to_end(&mut stdout_bytes) + .expect("attach stdout should be readable"); + let output = Output { + status, + stdout: stdout_bytes, + stderr: stderr_reader.join().expect("stderr reader should join"), + }; + assert!( + status.success(), + "attach failed after external answer:\nstdout:\n{}\nstderr:\n{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); +} + #[test] #[expect( clippy::disallowed_methods,