mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-10 03:30:59 +00:00
Merge pull request #941 from fabro-sh/codex/petri-test-server-port
Some checks are pending
Rust / Test (macOS) (push) Waiting to run
Rust / Process titles (musl) (push) Waiting to run
Rust / Format (push) Waiting to run
Rust / Clippy (push) Waiting to run
Rust / Rustdoc (push) Waiting to run
Rust / Generated Docs (push) Waiting to run
Rust / Test (Linux) (push) Waiting to run
Rust / Sandbox providers (Docker) (push) Waiting to run
Some checks are pending
Rust / Test (macOS) (push) Waiting to run
Rust / Process titles (musl) (push) Waiting to run
Rust / Format (push) Waiting to run
Rust / Clippy (push) Waiting to run
Rust / Rustdoc (push) Waiting to run
Rust / Generated Docs (push) Waiting to run
Rust / Test (Linux) (push) Waiting to run
Rust / Sandbox providers (Docker) (push) Waiting to run
Bind Petri scenario servers to OS-assigned ports
This commit is contained in:
commit
36f8b61b60
1 changed files with 98 additions and 34 deletions
|
|
@ -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 {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue