diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri.rs b/lib/apps/fabro-cli/tests/it/scenario/petri.rs index c641054fa..0662cf70f 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/petri.rs @@ -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::>(); - panic!( - "the server exited before it was ready with status {status}; its log ends \ - with:\n{}", - tail.into_iter().rev().collect::>().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::>(); + panic!( + "the server exited before it was ready with status {status}; its log ends \ + with:\n{}", + tail.into_iter().rev().collect::>().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 {