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.
This commit is contained in:
Bryan Helmkamp 2026-05-08 09:35:53 -07:00
parent b0dfd6b3b4
commit 82019d356f
No known key found for this signature in database
4 changed files with 467 additions and 11 deletions

1
Cargo.lock generated
View file

@ -1698,6 +1698,7 @@ dependencies = [
"jsonwebtoken",
"libc",
"miette",
"nix 0.30.1",
"object_store",
"openssl",
"paste",

View file

@ -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 }

View file

@ -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<Self> {
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<u8>) -> 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::<Vec<_>>();
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<Option<ExitCode>> {
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::<usize>().ok().and_then(|idx| {
idx.checked_sub(1)
.and_then(|zero_idx| question.options.get(zero_idx))
.map(|option| option.key.clone())
})
})
})
.collect::<Option<Vec<_>>>();
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<Answer> {
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::<usize>() {
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(

View file

@ -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,