Give up on an unreachable worker control stream and reap .ft- test daemons

A worker whose server is killed outright retried its control-stream
connect forever at a 5 s capped backoff, so it never finalized, never
stopped its sandbox, and lived until reboot. The worker now tracks the
start of each run of continuous connection failure and gives up after
60 s (10 s once its parent is pid 1), through the existing fatal
control-loss path that interrupts interviews and cancels the run. After
that fatal fires, the runner keeps driving the cancelled pipeline for a
bounded grace so `conclude` can stop the sandbox before the process
exits, instead of dropping the pipeline future mid-flight.

The test harness's stale-daemon reaper only matched `fabro server`
titles bound under `/tmp/.tmp*`, but `TestContext` roots are
`.ft-<label>-*` under `std::env::temp_dir()`, so daemons bound there
were never reaped. The regex now also matches sockets below a `.ft-`
directory component under any parent, and the unit test covers both
roots plus real-looking paths and `.ft-` outside a directory component.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-19 07:27:55 -04:00
parent 6e9be17e13
commit 1166caed0b
No known key found for this signature in database
2 changed files with 190 additions and 10 deletions

View file

@ -1,5 +1,6 @@
use std::collections::{HashSet, VecDeque};
use std::path::{Path, PathBuf};
use std::pin::pin;
use std::sync::Arc;
use std::time::Duration;
@ -27,6 +28,8 @@ use fabro_workflow::runtime_store::{RunStoreBackend, RunStoreHandle};
use fabro_workflow::services::FabroRunToolServices;
use futures::{SinkExt, StreamExt};
use jsonwebtoken::dangerous::insecure_decode;
#[cfg(unix)]
use nix::unistd;
#[cfg(test)]
use tokio::io::DuplexStream;
use tokio::net::TcpStream;
@ -172,13 +175,26 @@ pub(crate) async fn execute(
};
if let Some(mut control_manager) = control_manager {
let mut execution = pin!(execution);
tokio::select! {
result = execution => {
result = &mut execution => {
control_manager.finish();
result?;
}
fatal = control_manager.fatal_control_loss() => {
control_manager.finish();
// The fatal path already cancelled `cancel_token`. Keep
// driving the pipeline so it reaches its finalize step and
// stops the sandbox, instead of dropping it mid-flight.
if time::timeout(WORKER_FATAL_FINALIZE_GRACE, &mut execution)
.await
.is_err()
{
tracing::warn!(
grace = ?WORKER_FATAL_FINALIZE_GRACE,
"Pipeline did not finalize after worker control loss; exiting"
);
}
return Err(fatal);
}
}
@ -253,6 +269,21 @@ async fn load_worker_vault(storage_dir: &Path) -> Result<Arc<AsyncRwLock<Vault>>
const WORKER_CONTROL_RECONNECT_INITIAL_BACKOFF: Duration = Duration::from_millis(100);
const WORKER_CONTROL_RECONNECT_MAX_BACKOFF: Duration = Duration::from_secs(5);
/// How long the worker keeps retrying an unreachable control stream before it
/// gives up. Covers a `fabro server restart`: the 5 s shutdown grace, server
/// startup, and run reconciliation, with headroom for a loaded machine.
/// Longer buys nothing: the server has already marked the run failed once
/// its worker is unreachable, and the worker's only remaining job is to stop
/// its sandbox cleanly.
const WORKER_CONTROL_GIVE_UP: Duration = Duration::from_mins(1);
/// The shorter give-up used when the worker's parent is pid 1. Parent pid 1
/// is not proof that the server died, because the server daemonizes with
/// `setsid`, so this only shortens the window; it never triggers on its own.
const WORKER_CONTROL_ORPHAN_GIVE_UP: Duration = Duration::from_secs(10);
/// How long the worker keeps driving the cancelled pipeline after a fatal
/// control-stream loss, so `conclude` can stop the sandbox before the process
/// exits. Covers the cancelled-run diff timeout plus a container stop.
const WORKER_FATAL_FINALIZE_GRACE: Duration = Duration::from_secs(45);
const WORKER_CONTROL_APPLIED_ID_DEDUPE_CAPACITY: usize = 2048;
#[derive(Default)]
@ -436,6 +467,10 @@ async fn run_worker_control_manager(
let mut fatal_tx = Some(fatal_tx);
let mut backoff = WORKER_CONTROL_RECONNECT_INITIAL_BACKOFF;
let mut applied_ids = AppliedWorkerControlDeliveryIds::default();
// Start of the current run of continuous connection failure; cleared by
// every successful connect, so a connection that lands and later drops
// restarts the window.
let mut failing_since: Option<Instant> = None;
while !done.is_cancelled() {
let request = match build_worker_control_stream_request(
@ -464,6 +499,7 @@ async fn run_worker_control_manager(
let _ = first_tx.send(Ok(()));
}
backoff = WORKER_CONTROL_RECONNECT_INITIAL_BACKOFF;
failing_since = None;
match handle_worker_control_socket(
&mut socket,
&interviewer,
@ -505,6 +541,25 @@ async fn run_worker_control_manager(
}
Err(WorkerControlConnectError::Other(err)) => {
tracing::debug!(error = %err, "Worker control stream connection failed");
let now = Instant::now();
failing_since.get_or_insert(now);
if control_loss_exceeded(failing_since, now, worker_is_orphaned()) {
let elapsed =
failing_since.map_or(Duration::ZERO, |since| now.duration_since(since));
tracing::warn!(
elapsed = ?elapsed,
"Worker control stream unreachable; giving up"
);
report_fatal_control_loss(
&interviewer,
&cancel_token,
&mut first_tx,
&mut fatal_tx,
format!("worker control stream unreachable for {elapsed:?}; giving up"),
)
.await;
return;
}
}
}
@ -538,6 +593,33 @@ async fn sleep_or_done(done: &CancellationToken, delay: Duration) {
}
}
/// Whether the continuous connection-failure window that began at
/// `failing_since` has outlasted the applicable give-up limit. `None` means
/// no failure run is in progress, so it never fires.
fn control_loss_exceeded(failing_since: Option<Instant>, now: Instant, orphaned: bool) -> bool {
let Some(since) = failing_since else {
return false;
};
let limit = if orphaned {
WORKER_CONTROL_ORPHAN_GIVE_UP
} else {
WORKER_CONTROL_GIVE_UP
};
now.duration_since(since) > limit
}
/// Whether this worker has been reparented to pid 1.
fn worker_is_orphaned() -> bool {
#[cfg(unix)]
{
unistd::getppid().as_raw() == 1
}
#[cfg(not(unix))]
{
false
}
}
fn next_worker_control_reconnect_backoff(current: Duration) -> Duration {
current
.saturating_mul(2)
@ -1189,17 +1271,18 @@ mod tests {
use fabro_vault::{SecretType, Vault};
use fabro_workflow::event::RunEventSink;
use fabro_workflow::run_control::RunControlState;
use tokio::time;
use tokio::time::{self, Instant};
use tokio_tungstenite::tungstenite::protocol::{Message as TestWebSocketMessage, Role};
use tokio_util::sync::CancellationToken;
use super::{
AppliedWorkerControlDeliveryIds, WorkerControlConnectError, WorkerControlSocket,
AppliedWorkerControlDeliveryIds, WORKER_CONTROL_GIVE_UP, WORKER_CONTROL_ORPHAN_GIVE_UP,
WORKER_CONTROL_RECONNECT_MAX_BACKOFF, WorkerControlConnectError, WorkerControlSocket,
WorkerTitlePhase, apply_worker_control_delivery_frame, apply_worker_control_message,
build_worker_control_stream_request, connect_worker_control_stream,
build_worker_control_stream_request, connect_worker_control_stream, control_loss_exceeded,
handle_worker_control_socket, initial_worker_title_phase, load_worker_vault,
next_worker_control_reconnect_backoff, stamp_system_worker, worker_title,
worker_title_phase_for_event,
next_worker_control_reconnect_backoff, spawn_worker_control_manager, stamp_system_worker,
worker_title, worker_title_phase_for_event,
};
use crate::args::RunWorkerMode;
@ -1629,6 +1712,75 @@ mod tests {
);
}
#[test]
fn control_loss_exceeded_fires_only_past_the_applicable_limit() {
let start = Instant::now();
let at = |secs: u64| start + Duration::from_secs(secs);
assert!(!control_loss_exceeded(None, at(600), false));
assert!(!control_loss_exceeded(None, at(600), true));
assert!(!control_loss_exceeded(Some(start), at(5), false));
assert!(!control_loss_exceeded(Some(start), at(5), true));
assert!(!control_loss_exceeded(Some(start), at(11), false));
assert!(control_loss_exceeded(Some(start), at(11), true));
assert!(control_loss_exceeded(Some(start), at(61), false));
assert!(control_loss_exceeded(Some(start), at(61), true));
}
#[cfg(unix)]
#[tokio::test(start_paused = true)]
async fn worker_control_manager_gives_up_when_server_stays_unreachable() {
let temp = tempfile::tempdir().unwrap();
let socket_path = temp.path().join("missing.sock");
let interviewer = Arc::new(ControlInterviewer::new());
let cancel_token = CancellationToken::new();
let mut question = Question::new("Approve?", QuestionType::YesNo);
question.id = "q-1".to_string();
let ask_interviewer = Arc::clone(&interviewer);
let answer_task = tokio::spawn(async move { ask_interviewer.ask(question).await });
tokio::task::yield_now().await;
let started = Instant::now();
let mut handle = spawn_worker_control_manager(
ServerTarget::unix_socket_path(&socket_path).unwrap(),
fixtures::RUN_1,
"worker-token".to_string(),
Arc::clone(&interviewer),
cancel_token.clone(),
test_steering_hub(),
RunControlState::new(),
);
let first = handle.wait_for_first_connection().await;
let fatal = handle.fatal_control_loss().await;
let elapsed = started.elapsed();
handle.finish();
assert!(
first.is_err(),
"first connection should fail with the manager"
);
assert!(
fatal.to_string().contains("unreachable for"),
"unexpected fatal error: {fatal}"
);
assert!(cancel_token.is_cancelled());
assert_eq!(
answer_task.await.unwrap().answer.value,
AnswerValue::Interrupted
);
// The orphan limit is the earliest the manager may give up; the full
// limit plus one capped backoff sleep is the latest.
assert!(
elapsed > WORKER_CONTROL_ORPHAN_GIVE_UP
&& elapsed <= WORKER_CONTROL_GIVE_UP + WORKER_CONTROL_RECONNECT_MAX_BACKOFF,
"gave up after {elapsed:?}"
);
}
#[tokio::test(start_paused = true)]
async fn worker_control_socket_times_out_without_liveness() {
let (worker_io, _server_io) = tokio::io::duplex(1024);

View file

@ -1262,10 +1262,18 @@ fn parse_tmp_daemon_ps_line(line: &str) -> Option<(u32, u64, &str)> {
Some((pid, elapsed_secs, command))
}
/// Matches the process title of a test-spawned `fabro server` daemon: a Unix
/// socket under a `/tmp/.tmp*` tempdir, or anywhere below a
/// `TestContext` root (`.ft-<label>-<suffix>` under `std::env::temp_dir()`,
/// which is not `/tmp` on macOS). Real servers bind under `~/.fabro/` or a
/// configured storage dir, never a `.tmp` or `.ft-` directory.
fn tmp_daemon_socket_re() -> &'static Regex {
static TMP_DAEMON_SOCKET_RE: OnceLock<Regex> = OnceLock::new();
TMP_DAEMON_SOCKET_RE.get_or_init(|| {
Regex::new(r"^fabro server unix:/tmp/\.tmp[^/]+/[^/]+\.sock\s*$").expect("static regex")
Regex::new(
r"^fabro server unix:(?:/tmp/\.tmp[^/]+|(?:/[^/]+)+/\.ft-[^/]+-[^/]+(?:/[^/]+)*)/[^/]+\.sock\s*$",
)
.expect("static regex")
})
}
@ -2756,12 +2764,32 @@ mod tests {
fn tmp_daemon_socket_regex_matches_only_test_tmp_unix_daemons() {
let re = tmp_daemon_socket_re();
// `/tmp/.tmp*` roots.
assert!(re.is_match("fabro server unix:/tmp/.tmpAbC/fabro.sock"));
assert!(!re.is_match("fabro server unix:/tmp/.ft-foo-XYZ/test.sock"));
assert!(re.is_match("fabro server unix:/tmp/.tmp5vNC7E/configured.sock "));
// `TestContext` roots under any temp dir, including nested dirs.
assert!(re.is_match("fabro server unix:/tmp/.ft-foo-XYZ/test.sock"));
assert!(re.is_match(
"fabro server unix:/private/var/folders/nm/abc/T/.ft-a_crash_before_t-bMEw4g/install-storage.sock"
));
assert!(re.is_match(
"fabro server unix:/var/folders/nm/abc/T/.ft-auth_login_reje-Q1w2e3/temp/fabro.sock"
));
// Not a daemon title.
assert!(!re.is_match("sh -c fabro server unix:/tmp/.tmpAbC/fabro.sock"));
assert!(!re.is_match("fabro server unix:/Users/me/.fabro/fabro.sock"));
assert!(!re.is_match("fabro server unix:/tmp/notmatching/fabro.sock"));
assert!(!re.is_match("fabro server tcp:127.0.0.1:32276"));
// Real servers.
assert!(!re.is_match("fabro server unix:/Users/me/.fabro/fabro.sock"));
assert!(!re.is_match("fabro server unix:/var/lib/fabro/storage/fabro.sock"));
assert!(!re.is_match("fabro server unix:/tmp/notmatching/fabro.sock"));
assert!(!re.is_match("fabro server unix:/tmp/.tmpAbC/nested/fabro.sock"));
// `.ft-` anywhere other than as its own directory component.
assert!(!re.is_match("fabro server unix:/tmp/.ft-x-y.sock"));
assert!(!re.is_match("fabro server unix:/tmp/my.ft-foo-XYZ/fabro.sock"));
assert!(!re.is_match("fabro server unix:/Users/me/.fabro/.ft-foo-XYZ.sock"));
assert!(!re.is_match("fabro server unix:/tmp/.ft-noseparator/fabro.sock"));
assert!(!re.is_match("fabro server unix:.ft-foo-XYZ/fabro.sock"));
}
#[test]