Merge remote-tracking branch 'origin/main' into petri-integration

# Conflicts:
#	lib/apps/fabro-cli/src/commands/run/runner.rs
This commit is contained in:
Bryan Helmkamp 2026-09-19 15:25:20 -04:00
commit ec5faab116
No known key found for this signature in database
4 changed files with 270 additions and 52 deletions

117
Cargo.lock generated
View file

@ -2006,7 +2006,7 @@ dependencies = [
"progenitor-client",
"regress",
"reqwest 0.13.4",
"sandbox-driver",
"sandbox-driver 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"serde",
"serde_json",
"serde_yaml",
@ -2145,7 +2145,7 @@ dependencies = [
"ring",
"rmcp",
"rustls",
"sandbox-driver",
"sandbox-driver 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"sandbox-driver-host",
"scopeguard",
"semver",
@ -2523,7 +2523,7 @@ dependencies = [
"fabro-types",
"fabro-util",
"pebble-coding-agent",
"sandbox-driver",
"sandbox-driver 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"sandbox-driver-host",
"sandbox-driver-testing",
"serde_json",
@ -2674,11 +2674,11 @@ dependencies = [
"percent-encoding",
"rand 0.9.4",
"reqwest 0.12.28",
"sandbox-driver",
"sandbox-driver 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"sandbox-driver-daytona",
"sandbox-driver-docker",
"sandbox-driver-host",
"sandbox-driver-protocol",
"sandbox-driver-protocol 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"sandbox-driver-testing",
"serde",
"serde_json",
@ -2880,7 +2880,7 @@ dependencies = [
"hex",
"lithos-llm",
"pebble-coding-agent",
"sandbox-driver",
"sandbox-driver 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"serde",
"serde_json",
"sha2 0.10.9",
@ -2987,7 +2987,7 @@ dependencies = [
"httpmock",
"lithos-llm",
"pebble-coding-agent",
"sandbox-driver",
"sandbox-driver 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"sandbox-driver-host",
"scopeguard",
"serde",
@ -5284,10 +5284,10 @@ dependencies = [
"async-trait",
"petri-executor",
"petri-ir",
"sandbox-driver",
"sandbox-driver-daytona-config",
"sandbox-driver-docker-config",
"sandbox-driver-protocol",
"sandbox-driver 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver.git?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d)",
"sandbox-driver-daytona-config 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver.git?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d)",
"sandbox-driver-docker-config 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver.git?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d)",
"sandbox-driver-protocol 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver.git?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d)",
"serde",
"serde_json",
"sha2 0.10.9",
@ -5434,7 +5434,7 @@ dependencies = [
"petri-ir",
"petri-steps",
"petri-store",
"sandbox-driver",
"sandbox-driver 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver.git?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d)",
"serde",
"serde_json",
"smol_str",
@ -6303,7 +6303,24 @@ dependencies = [
[[package]]
name = "sandbox-driver"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d#64c14b89d078a4b34d1555092ad01d41541f7a7d"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c#07600aa5c6695ec4c999da93c05d2cb78fe11b0c"
dependencies = [
"async-trait",
"globset",
"humantime",
"rand 0.10.1",
"serde",
"serde_json",
"thiserror 2.0.18",
"tokio",
"tokio-util",
"tracing",
]
[[package]]
name = "sandbox-driver"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/sandbox-driver.git?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d#64c14b89d078a4b34d1555092ad01d41541f7a7d"
dependencies = [
"async-trait",
"globset",
@ -6320,7 +6337,7 @@ dependencies = [
[[package]]
name = "sandbox-driver-daytona"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d#64c14b89d078a4b34d1555092ad01d41541f7a7d"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c#07600aa5c6695ec4c999da93c05d2cb78fe11b0c"
dependencies = [
"anyhow",
"async-trait",
@ -6330,11 +6347,11 @@ dependencies = [
"hmac 0.12.1",
"rand 0.10.1",
"reqwest 0.13.4",
"sandbox-driver",
"sandbox-driver-daytona-config",
"sandbox-driver 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"sandbox-driver-daytona-config 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"sandbox-driver-docker",
"sandbox-driver-docker-config",
"sandbox-driver-protocol",
"sandbox-driver-docker-config 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"sandbox-driver-protocol 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"serde",
"serde_json",
"sha2 0.10.9",
@ -6347,9 +6364,19 @@ dependencies = [
[[package]]
name = "sandbox-driver-daytona-config"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d#64c14b89d078a4b34d1555092ad01d41541f7a7d"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c#07600aa5c6695ec4c999da93c05d2cb78fe11b0c"
dependencies = [
"sandbox-driver-docker-config",
"sandbox-driver-docker-config 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"serde",
"serde_json",
]
[[package]]
name = "sandbox-driver-daytona-config"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/sandbox-driver.git?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d#64c14b89d078a4b34d1555092ad01d41541f7a7d"
dependencies = [
"sandbox-driver-docker-config 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver.git?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d)",
"serde",
"serde_json",
]
@ -6357,15 +6384,15 @@ dependencies = [
[[package]]
name = "sandbox-driver-docker"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d#64c14b89d078a4b34d1555092ad01d41541f7a7d"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c#07600aa5c6695ec4c999da93c05d2cb78fe11b0c"
dependencies = [
"anyhow",
"async-trait",
"bollard",
"futures-util",
"sandbox-driver",
"sandbox-driver-docker-config",
"sandbox-driver-protocol",
"sandbox-driver 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"sandbox-driver-docker-config 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"sandbox-driver-protocol 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"serde",
"serde_json",
"tar",
@ -6378,7 +6405,16 @@ dependencies = [
[[package]]
name = "sandbox-driver-docker-config"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d#64c14b89d078a4b34d1555092ad01d41541f7a7d"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c#07600aa5c6695ec4c999da93c05d2cb78fe11b0c"
dependencies = [
"serde",
"serde_json",
]
[[package]]
name = "sandbox-driver-docker-config"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/sandbox-driver.git?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d#64c14b89d078a4b34d1555092ad01d41541f7a7d"
dependencies = [
"serde",
"serde_json",
@ -6387,13 +6423,13 @@ dependencies = [
[[package]]
name = "sandbox-driver-host"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d#64c14b89d078a4b34d1555092ad01d41541f7a7d"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c#07600aa5c6695ec4c999da93c05d2cb78fe11b0c"
dependencies = [
"anyhow",
"async-trait",
"nix 0.30.1",
"sandbox-driver",
"sandbox-driver-protocol",
"sandbox-driver 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"sandbox-driver-protocol 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"serde",
"serde_json",
"tokio",
@ -6405,12 +6441,29 @@ dependencies = [
[[package]]
name = "sandbox-driver-protocol"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d#64c14b89d078a4b34d1555092ad01d41541f7a7d"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c#07600aa5c6695ec4c999da93c05d2cb78fe11b0c"
dependencies = [
"async-trait",
"base64",
"rand 0.10.1",
"sandbox-driver",
"sandbox-driver 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"serde",
"serde_json",
"sha2 0.10.9",
"tokio",
"tokio-util",
"tracing",
]
[[package]]
name = "sandbox-driver-protocol"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/sandbox-driver.git?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d#64c14b89d078a4b34d1555092ad01d41541f7a7d"
dependencies = [
"async-trait",
"base64",
"rand 0.10.1",
"sandbox-driver 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver.git?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d)",
"serde",
"serde_json",
"sha2 0.10.9",
@ -6422,10 +6475,10 @@ dependencies = [
[[package]]
name = "sandbox-driver-testing"
version = "0.1.0"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=64c14b89d078a4b34d1555092ad01d41541f7a7d#64c14b89d078a4b34d1555092ad01d41541f7a7d"
source = "git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c#07600aa5c6695ec4c999da93c05d2cb78fe11b0c"
dependencies = [
"async-trait",
"sandbox-driver",
"sandbox-driver 0.1.0 (git+https://github.com/lithoscomputer/sandbox-driver?rev=07600aa5c6695ec4c999da93c05d2cb78fe11b0c)",
"tokio",
]

View file

@ -90,16 +90,19 @@ futures-util = "0.3"
# before a generation serves, and a refusal of a replacement that reports
# another resource namespace; #22 made the Daytona exec decoder resync on a
# torn frame and report what it dropped as `ExecStreamingResult::output_loss`
# instead of failing the command). The CI plugin job installs the driver
# executables at the same rev, read from this file.
sandbox-driver = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "64c14b89d078a4b34d1555092ad01d41541f7a7d" }
sandbox-driver-protocol = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "64c14b89d078a4b34d1555092ad01d41541f7a7d" }
sandbox-driver-host = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "64c14b89d078a4b34d1555092ad01d41541f7a7d" }
sandbox-driver-docker = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "64c14b89d078a4b34d1555092ad01d41541f7a7d" }
sandbox-driver-docker-config = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "64c14b89d078a4b34d1555092ad01d41541f7a7d" }
sandbox-driver-daytona = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "64c14b89d078a4b34d1555092ad01d41541f7a7d" }
sandbox-driver-daytona-config = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "64c14b89d078a4b34d1555092ad01d41541f7a7d" }
sandbox-driver-testing = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "64c14b89d078a4b34d1555092ad01d41541f7a7d" }
# instead of failing the command; #23 made an idle Host sentinel reap its
# own process group once its owning provider pid is gone, so a provider
# that is killed rather than stopped leaves no sentinel behind). The CI
# plugin job installs the driver executables at the same rev, read from
# this file.
sandbox-driver = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "07600aa5c6695ec4c999da93c05d2cb78fe11b0c" }
sandbox-driver-protocol = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "07600aa5c6695ec4c999da93c05d2cb78fe11b0c" }
sandbox-driver-host = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "07600aa5c6695ec4c999da93c05d2cb78fe11b0c" }
sandbox-driver-docker = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "07600aa5c6695ec4c999da93c05d2cb78fe11b0c" }
sandbox-driver-docker-config = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "07600aa5c6695ec4c999da93c05d2cb78fe11b0c" }
sandbox-driver-daytona = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "07600aa5c6695ec4c999da93c05d2cb78fe11b0c" }
sandbox-driver-daytona-config = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "07600aa5c6695ec4c999da93c05d2cb78fe11b0c" }
sandbox-driver-testing = { git = "https://github.com/lithoscomputer/sandbox-driver", rev = "07600aa5c6695ec4c999da93c05d2cb78fe11b0c" }
# pebble: the coding agent loop fabro runs its agent stages, Ask Fabro
# sessions, hook evaluators, and `fabro exec` on. Pinned by rev to pebble
# `main`. Pebble owns the `PortRoutes` trait fabro implements over its run

View file

@ -20,6 +20,8 @@ use fabro_vault::{SecretStore, Vault};
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;
@ -155,6 +157,17 @@ pub(super) async fn load_worker_vault(storage_dir: &Path) -> Result<Arc<AsyncRwL
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);
const WORKER_CONTROL_APPLIED_ID_DEDUPE_CAPACITY: usize = 2048;
#[derive(Default)]
@ -339,6 +352,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(
@ -367,6 +384,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,
@ -407,6 +425,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;
}
}
}
@ -440,6 +477,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)
@ -780,17 +844,19 @@ mod tests {
use fabro_store::PlatformRecord;
use fabro_types::{Principal, QuestionType, SystemActorKind, fixtures};
use fabro_vault::{SecretType, Vault};
use tokio::time;
use tokio::time::{self, Instant};
use tokio_tungstenite::tungstenite::protocol::{Message as TestWebSocketMessage, Role};
use tokio_util::sync::CancellationToken;
use super::super::petri_worker::PetriControls;
use super::{
AppliedWorkerControlDeliveryIds, RunControls, WorkerControlConnectError,
WorkerControlSocket, WorkerControls, WorkerTitlePhase, apply_worker_control_delivery_frame,
apply_worker_control_message, build_worker_control_stream_request,
connect_worker_control_stream, handle_worker_control_socket, initial_worker_title_phase,
load_worker_vault, next_worker_control_reconnect_backoff, worker_title,
AppliedWorkerControlDeliveryIds, RunControls, WORKER_CONTROL_GIVE_UP,
WORKER_CONTROL_ORPHAN_GIVE_UP, WORKER_CONTROL_RECONNECT_MAX_BACKOFF,
WorkerControlConnectError, WorkerControlSocket, WorkerControls, WorkerTitlePhase,
apply_worker_control_delivery_frame, apply_worker_control_message,
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, spawn_worker_control_manager, worker_title,
};
use crate::args::RunWorkerMode;
@ -1205,6 +1271,74 @@ 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_controls(),
);
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

@ -1286,10 +1286,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")
})
}
@ -2780,12 +2788,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]