Merge pull request #884 from fabro-sh/worker-control-give-up

Give up on an unreachable worker control stream and reap .ft- test daemons
This commit is contained in:
Bryan Helmkamp 2026-09-19 09:52:34 -04:00 • committed by GitHub
commit 1caf37cc31
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
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]