From 3c942d8a79b432fe4714f5e1714cb0c318544232 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Thu, 1 Oct 2026 16:33:44 -0400 Subject: [PATCH] Fix attach test deadlines around workflow completion --- lib/apps/fabro-cli/tests/it/cmd/attach.rs | 184 +++++++++++++--------- 1 file changed, 110 insertions(+), 74 deletions(-) diff --git a/lib/apps/fabro-cli/tests/it/cmd/attach.rs b/lib/apps/fabro-cli/tests/it/cmd/attach.rs index 8c03ba3ca..c8978361e 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/attach.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/attach.rs @@ -5,15 +5,18 @@ use std::fmt::Write as _; use std::io::{BufRead, BufReader, Read, Write}; -use std::process::{Output, Stdio}; +use std::process::{Child, ExitStatus, Output, Stdio}; use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::mpsc; +use std::thread::JoinHandle; use std::time::{Duration, Instant}; +use anyhow::{Context as _, Result, ensure}; use fabro_test::{ apply_filters, assert_reqwest_status, expect_reqwest_json, fabro_json_snapshot, fabro_snapshot, test_context, }; +use fabro_types::RunProjection; use httpmock::{HttpMockResponse, MockServer}; use serde_json::Value; @@ -24,6 +27,7 @@ use super::support::{ use crate::support::run_output_filters; const SHARED_DAEMON_TIMEOUT: Duration = Duration::from_secs(30); +const ATTACH_EXIT_TIMEOUT: Duration = Duration::from_secs(5); async fn wait_for_server_question( client: &fabro_http::HttpClient, @@ -107,30 +111,30 @@ fn redact_volatile_fields(value: &mut Value) { } fn wait_for_output_signal( - child: &mut std::process::Child, - stdout: &mut impl Read, - stderr_reader: std::thread::JoinHandle>, + child: &mut Child, + stdout_reader: JoinHandle>, + stderr_reader: JoinHandle>, signal_rx: &mpsc::Receiver<()>, needle: &str, -) -> std::thread::JoinHandle> { +) -> (JoinHandle>, JoinHandle>) { let deadline = Instant::now() + SHARED_DAEMON_TIMEOUT; let mut stderr_reader = Some(stderr_reader); loop { match signal_rx.recv_timeout(Duration::from_millis(20)) { Ok(()) => { - return stderr_reader - .take() - .expect("stderr reader should still be available"); + return ( + stdout_reader, + stderr_reader + .take() + .expect("stderr reader should still be available"), + ); } Err(mpsc::RecvTimeoutError::Timeout | mpsc::RecvTimeoutError::Disconnected) => {} } if let Some(status) = child.try_wait().expect("attach should stay alive") { - let mut stdout_bytes = Vec::new(); - stdout - .read_to_end(&mut stdout_bytes) - .expect("attach stdout should be readable"); + let stdout_bytes = stdout_reader.join().expect("stdout reader should join"); let stderr_bytes = stderr_reader .take() .expect("stderr reader should still be available") @@ -146,10 +150,7 @@ fn wait_for_output_signal( if Instant::now() >= deadline { let _ = child.kill(); let status = child.wait().expect("attach should exit after kill"); - let mut stdout_bytes = Vec::new(); - stdout - .read_to_end(&mut stdout_bytes) - .expect("attach stdout should be readable"); + let stdout_bytes = stdout_reader.join().expect("stdout reader should join"); let stderr_bytes = stderr_reader .take() .expect("stderr reader should still be available") @@ -166,28 +167,92 @@ fn wait_for_output_signal( #[expect( clippy::disallowed_methods, - reason = "This sync integration helper polls a child process without a Tokio runtime." + reason = "A dedicated reader thread drains each child pipe while the test waits." )] -fn wait_for_child_exit(child: &mut std::process::Child, label: &str) -> std::process::ExitStatus { - let deadline = Instant::now() + Duration::from_secs(5); +fn read_output(mut reader: impl Read + Send + 'static) -> JoinHandle> { + std::thread::spawn(move || { + let mut bytes = Vec::new(); + reader + .read_to_end(&mut bytes) + .expect("attach output should be readable"); + bytes + }) +} + +async fn wait_for_attach_exit( + context: &fabro_test::TestContext, + run_id: &str, + child: &mut Child, +) -> Result { + let (client, base_url) = + server_endpoint(&context.storage_dir).context("server endpoint should exist")?; + let mut deadline = Instant::now() + SHARED_DAEMON_TIMEOUT; + let mut workflow_finished = false; + let mut last_status = None; loop { - if let Some(status) = child - .try_wait() - .unwrap_or_else(|err| panic!("{label} status should be readable: {err}")) - { - return status; + if let Some(status) = child.try_wait().context("read attach exit status")? { + return Ok(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}"); + ensure!( + Instant::now() < deadline, + "timed out waiting for {}; last run status: {last_status:?}", + if workflow_finished { + "attach to exit after workflow completion" + } else { + "workflow completion" + }, + ); + if !workflow_finished { + let state = client + .get(format!("{base_url}/api/v1/runs/{run_id}/state")) + .timeout(deadline.saturating_duration_since(Instant::now())) + .send() + .await? + .error_for_status()? + .json::() + .await + .context("read run state while waiting for workflow completion")?; + last_status = Some(state.status); + if state.status.is_terminal() { + // The workflow (including checkpoints) gets its own budget. Only + // a completed workflow starts the shorter attach shutdown budget. + workflow_finished = true; + deadline = Instant::now() + ATTACH_EXIT_TIMEOUT; + } } - std::thread::sleep(Duration::from_millis(20)); + tokio::time::sleep(Duration::from_millis(20)).await; } } +fn wait_for_attach_output( + context: &fabro_test::TestContext, + run_id: &str, + child: &mut Child, + stdout_reader: JoinHandle>, + stderr_reader: JoinHandle>, +) -> Output { + let result = tokio::runtime::Runtime::new() + .expect("test runtime should build") + .block_on(wait_for_attach_exit(context, run_id, child)); + if result.is_err() { + let _ = child.kill(); + } + let status = child.wait().expect("attach should be reaped"); + let output = Output { + status, + stdout: stdout_reader.join().expect("stdout reader should join"), + stderr: stderr_reader.join().expect("stderr reader should join"), + }; + if let Err(error) = result { + panic!( + "attach for run {run_id}: {error:#}\nstatus: {status}\nstdout:\n{}\nstderr:\n{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr), + ); + } + output +} + fn start_detached_human_run( context: &fabro_test::TestContext, filename: &str, @@ -404,8 +469,8 @@ fn attach_with_stdin(context: &fabro_test::TestContext, run_id: &str, input: &[u attach_cmd.stderr(Stdio::piped()); let mut child = attach_cmd.spawn().expect("attach should spawn"); - let mut stdout = child.stdout.take().expect("attach stdout should be piped"); - let mut stderr = child.stderr.take().expect("attach stderr should be piped"); + let stdout_reader = read_output(child.stdout.take().expect("attach stdout should be piped")); + let stderr_reader = read_output(child.stderr.take().expect("attach stderr should be piped")); { let mut stdin = child.stdin.take().expect("attach stdin should be piped"); stdin @@ -413,20 +478,7 @@ fn attach_with_stdin(context: &fabro_test::TestContext, run_id: &str, input: &[u .expect("scripted attach input should be writable"); } - 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 mut stderr_bytes = Vec::new(); - stderr - .read_to_end(&mut stderr_bytes) - .expect("attach stderr should be readable"); - Output { - status, - stdout: stdout_bytes, - stderr: stderr_bytes, - } + wait_for_attach_output(context, run_id, &mut child, stdout_reader, stderr_reader) } #[test] @@ -631,7 +683,7 @@ fn attach_advances_when_pending_question_is_answered_elsewhere() { 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 stdout_reader = read_output(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 || { @@ -658,9 +710,9 @@ fn attach_advances_when_pending_question_is_answered_elsewhere() { stderr_bytes }); - let stderr_reader = wait_for_output_signal( + let (stdout_reader, stderr_reader) = wait_for_output_signal( &mut child, - &mut stdout, + stdout_reader, stderr_reader, &signal_rx, "Approve?", @@ -684,18 +736,10 @@ fn attach_advances_when_pending_question_is_answered_elsewhere() { .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"), - }; + let output = + wait_for_attach_output(&context, &run_id, &mut child, stdout_reader, stderr_reader); assert!( - status.success(), + output.status.success(), "attach failed after external answer:\nstdout:\n{}\nstderr:\n{}", String::from_utf8_lossy(&output.stdout), String::from_utf8_lossy(&output.stderr) @@ -746,7 +790,7 @@ fn attach_before_completion_streams_to_finished_state() { attach_cmd.stdout(Stdio::piped()); attach_cmd.stderr(Stdio::piped()); let mut child = attach_cmd.spawn().expect("attach should spawn"); - let mut stdout = child.stdout.take().expect("attach stdout should be piped"); + let stdout_reader = read_output(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 || { @@ -773,24 +817,16 @@ fn attach_before_completion_streams_to_finished_state() { stderr_bytes }); - let stderr_reader = wait_for_output_signal( + let (stdout_reader, stderr_reader) = wait_for_output_signal( &mut child, - &mut stdout, + stdout_reader, stderr_reader, &signal_rx, "✓ start", ); gate.release(); - let status = child.wait().expect("attach should exit"); - 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"), - }; + let output = + wait_for_attach_output(&context, &run_id, &mut child, stdout_reader, stderr_reader); let snapshot = format_output_snapshot(&output, &filters); wait_for_status(&run.run_dir, &["succeeded"]);