mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-08 03:10:26 +00:00
Fix attach test deadlines around workflow completion
This commit is contained in:
parent
e7b4859038
commit
3c942d8a79
1 changed files with 110 additions and 74 deletions
|
|
@ -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<Vec<u8>>,
|
||||
child: &mut Child,
|
||||
stdout_reader: JoinHandle<Vec<u8>>,
|
||||
stderr_reader: JoinHandle<Vec<u8>>,
|
||||
signal_rx: &mpsc::Receiver<()>,
|
||||
needle: &str,
|
||||
) -> std::thread::JoinHandle<Vec<u8>> {
|
||||
) -> (JoinHandle<Vec<u8>>, JoinHandle<Vec<u8>>) {
|
||||
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<Vec<u8>> {
|
||||
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<ExitStatus> {
|
||||
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::<RunProjection>()
|
||||
.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<Vec<u8>>,
|
||||
stderr_reader: JoinHandle<Vec<u8>>,
|
||||
) -> 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"]);
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue