Bind Petri scenario servers to OS-assigned ports

This commit is contained in:
Scott Werner 2026-10-08 10:33:38 -04:00
parent b609f6e015
commit d3d8893bf9

View file

@ -27,7 +27,8 @@ use std::process::{Child, Command, Output, Stdio};
use std::time::{Duration, Instant};
use fabro_client::ServerTarget;
use fabro_config::{Storage, envfile};
use fabro_config::daemon::ServerDaemon;
use fabro_config::{RuntimeDirectory, Storage, envfile};
use fabro_petri::SqliteRunStore;
use fabro_petri::checkpoint::CheckpointKey;
use fabro_petri::engine::{self, RunStatus};
@ -54,7 +55,6 @@ pub(super) struct RunningServer {
_storage_root: tempfile::TempDir,
pub(super) storage_dir: PathBuf,
config_path: PathBuf,
port: u16,
pub(super) api_base_url: String,
/// The checkpoint gate directory the server forwards to its workers.
gates_dir: PathBuf,
@ -85,7 +85,6 @@ impl RunningServer {
let home_root = tempfile::tempdir_in("/tmp").expect("home tempdir");
let storage_root = isolated_storage_dir();
let storage_dir = storage_root.path().join("storage");
let port = reserve_port();
let config_path = home_root.path().join("settings.toml");
std::fs::write(
&config_path,
@ -117,8 +116,7 @@ impl RunningServer {
_storage_root: storage_root,
storage_dir,
config_path,
port,
api_base_url: format!("http://127.0.0.1:{port}"),
api_base_url: String::new(),
gates_dir,
env: env
.iter()
@ -150,16 +148,21 @@ impl RunningServer {
.arg("--storage-dir")
.arg(&self.storage_dir)
.arg("--bind")
.arg(format!("127.0.0.1:{}", self.port))
.arg("127.0.0.1:0")
.arg("--config")
.arg(&self.config_path)
.stdin(Stdio::null())
.stdout(self.stderr_log())
.stderr(self.stderr_log());
let mut child = cmd.spawn().expect("the server spawns");
// Keep ownership during readiness so a failed launch is killed on Drop.
self.child = Some(cmd.spawn().expect("the server spawns"));
let log_path = self.storage_dir.with_file_name("server.stderr.log");
wait_for_http_ready(&self.api_base_url, &mut child, &log_path).await;
self.child = Some(child);
self.api_base_url = wait_for_http_ready(
&Storage::new(&self.storage_dir).runtime_directory(),
self.child.as_mut().expect("the server was spawned"),
&log_path,
)
.await;
}
/// Where the server's stdout and stderr go: a file beside its storage,
@ -351,38 +354,99 @@ impl Drop for RunningServer {
}
}
fn reserve_port() -> u16 {
std::net::TcpListener::bind("127.0.0.1:0")
.expect("a port binds")
.local_addr()
.expect("the listener has an address")
.port()
}
async fn wait_for_http_ready(base_url: &str, child: &mut Child, log_path: &Path) {
async fn wait_for_http_ready(
runtime_directory: &RuntimeDirectory,
child: &mut Child,
log_path: &Path,
) -> String {
let client = fabro_test::test_http_client();
let deadline = Instant::now() + Duration::from_secs(10);
loop {
match client.get(format!("{base_url}/health")).send().await {
Ok(response) if response.status().is_success() => return,
Ok(_) | Err(_) if Instant::now() < deadline => {
if let Some(status) = child.try_wait().expect("the server polls") {
let log = std::fs::read_to_string(log_path).unwrap_or_default();
let tail = log.lines().rev().take(20).collect::<Vec<_>>();
panic!(
"the server exited before it was ready with status {status}; its log ends \
with:\n{}",
tail.into_iter().rev().collect::<Vec<_>>().join("\n")
);
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
Ok(response) => panic!("server at {base_url} was not ready: {}", response.status()),
Err(err) => panic!("server at {base_url} was not ready: {err}"),
if let Some(status) = child.try_wait().expect("the server polls") {
let log = std::fs::read_to_string(log_path).unwrap_or_default();
let tail = log.lines().rev().take(20).collect::<Vec<_>>();
panic!(
"the server exited before it was ready with status {status}; its log ends \
with:\n{}",
tail.into_iter().rev().collect::<Vec<_>>().join("\n")
);
}
// The server binds port zero itself and publishes the actual address.
// A SIGKILL leaves the previous record behind: only trust this child.
if let Some(daemon) = ServerDaemon::read(runtime_directory)
.expect("the server record reads")
.filter(|daemon| daemon.pid == child.id())
{
let base_url = daemon.bind.to_target();
match client
.get(format!("{base_url}/health"))
.timeout(Duration::from_millis(250))
.send()
.await
{
Ok(response) if response.status().is_success() => return base_url,
Ok(_) | Err(_) => {}
}
}
assert!(
Instant::now() < deadline,
"the server did not become ready; record: {:?}; log:\n{}",
ServerDaemon::read(runtime_directory),
std::fs::read_to_string(log_path).unwrap_or_default(),
);
tokio::time::sleep(Duration::from_millis(25)).await;
}
}
/// A crashed server's port may already belong to a healthy unrelated
/// service. Relaunch must discover the new child, not trust the stale record.
#[tokio::test(flavor = "multi_thread")]
async fn a_server_relaunch_uses_its_own_port_when_the_old_port_is_occupied() {
let _context = test_context!();
let mut server = RunningServer::start().await;
let old_url = server.api_base_url.clone();
server.kill();
let listener = tokio::net::TcpListener::bind(
old_url
.strip_prefix("http://")
.expect("the server uses HTTP"),
)
.await
.expect("the crashed server's port is occupied");
let impostor = tokio::spawn(async move {
axum::serve(
listener,
axum::Router::new().route("/health", axum::routing::get(|| async { "ok" })),
)
.await
.expect("the unrelated health server serves");
});
let _impostor_guard = scopeguard::guard(impostor, |task| task.abort());
fabro_test::expect_reqwest_status(
fabro_test::test_http_client()
.get(format!("{old_url}/health"))
.send()
.await
.expect("the unrelated health server responds"),
axum::http::StatusCode::OK,
"GET /health on the occupied old port",
)
.await;
server.launch().await;
assert_ne!(server.api_base_url, old_url);
let daemon = ServerDaemon::read(&Storage::new(&server.storage_dir).runtime_directory())
.expect("the server record reads")
.expect("the new server wrote its record");
assert_eq!(
daemon.pid,
server.child.as_ref().expect("the new server runs").id()
);
assert_eq!(daemon.bind.to_target(), server.api_base_url);
server.shutdown();
}
/// A workspace holding a command-only bundle whose `workflow.toml` names
/// Petri, with the given stage script.
fn write_petri_workspace(context: &fabro_test::TestContext, script: &str) -> PathBuf {