From 16354186fa05f10283b725af317d984facfe53a8 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 17 Sep 2026 21:20:22 -0400 Subject: [PATCH] Run a Petri run in the worker process over the HTTP store When `fabro run __run-worker` finds its run's stored spec names Petri, the new `petri_worker` module executes it through `fabro_petri::engine` over `HttpRunStore`, leased for a launch id the worker mints and logs at start. `--mode start` loads the admitted graphs through the client's blob read; `--mode resume` continues the run from its records. The worker's existing services carry over: the control channel's cancel and SIGTERM/SIGINT cancel Petri's root invocation politely, a lost control channel cancels the run and is reported once it settles, and pause, unpause and steer are received and ignored with a warning until their adapters land. The model client comes from the worker's catalog and vault snapshot for the providers whose credentials resolve, and the lifecycle events (`run.starting`, `run.running`, then `run.completed` or `run.failed`) go through the client as the legacy worker's do. Scenario tests against the real binary: a command-only Petri run executes in the worker a foreground server launched, its records reach `petri_records` over the HTTP store and its lease ends with the worker; and a run whose server and worker are both killed mid-stage resumes in a new worker after the server restarts, with one `run.completed`. Co-Authored-By: Claude Fable 5.1 --- Cargo.lock | 3 + lib/apps/fabro-cli/Cargo.toml | 2 + lib/apps/fabro-cli/src/commands/run/mod.rs | 1 + .../src/commands/run/petri_worker.rs | 255 ++++++++ lib/apps/fabro-cli/src/commands/run/runner.rs | 35 +- lib/apps/fabro-cli/tests/it/scenario/mod.rs | 1 + lib/apps/fabro-cli/tests/it/scenario/petri.rs | 568 ++++++++++++++++++ 7 files changed, 855 insertions(+), 10 deletions(-) create mode 100644 lib/apps/fabro-cli/src/commands/run/petri_worker.rs create mode 100644 lib/apps/fabro-cli/tests/it/scenario/petri.rs diff --git a/Cargo.lock b/Cargo.lock index 2373be029..670f22e8b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2396,6 +2396,7 @@ dependencies = [ "fabro-checkpoint", "fabro-client", "fabro-config", + "fabro-db", "fabro-dump", "fabro-environment", "fabro-github", @@ -2410,6 +2411,7 @@ dependencies = [ "fabro-mcp", "fabro-mcp-server", "fabro-oauth", + "fabro-petri", "fabro-proc", "fabro-redact", "fabro-sandbox", @@ -2888,6 +2890,7 @@ version = "0.357.0-nightly.0" dependencies = [ "anyhow", "async-trait", + "bytes", "fabro-api", "fabro-auth", "fabro-client", diff --git a/lib/apps/fabro-cli/Cargo.toml b/lib/apps/fabro-cli/Cargo.toml index e953d0a26..881616e2f 100644 --- a/lib/apps/fabro-cli/Cargo.toml +++ b/lib/apps/fabro-cli/Cargo.toml @@ -34,6 +34,7 @@ fabro-install = { path = "../../components/fabro-install" } fabro-interview = { path = "../../components/fabro-interview" } fabro-mcp = { path = "../../components/fabro-mcp" } fabro-mcp-server = { path = "../fabro-mcp-server" } +fabro-petri = { path = "../../components/fabro-petri" } fabro-manifest = { path = "../../components/fabro-manifest" } fabro-proc = { path = "../../foundation/fabro-proc" } fabro-sandbox = { path = "../../components/fabro-sandbox" } @@ -117,6 +118,7 @@ chrono = { workspace = true } [dev-dependencies] assert_cmd = "2" +fabro-db = { path = "../../foundation/fabro-db" } walkdir.workspace = true fabro-acp = { path = "../../components/fabro-acp", features = ["test-support"] } fabro-mcp = { path = "../../components/fabro-mcp", features = ["test-support"] } diff --git a/lib/apps/fabro-cli/src/commands/run/mod.rs b/lib/apps/fabro-cli/src/commands/run/mod.rs index b904549b4..47ae4ad9f 100644 --- a/lib/apps/fabro-cli/src/commands/run/mod.rs +++ b/lib/apps/fabro-cli/src/commands/run/mod.rs @@ -20,6 +20,7 @@ pub(crate) mod fork; pub(crate) mod logs; pub(crate) mod output; pub(crate) mod overrides; +mod petri_worker; pub(crate) mod preview; mod remote_workflow; mod resolution; diff --git a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs new file mode 100644 index 000000000..2ec71756c --- /dev/null +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -0,0 +1,255 @@ +//! A Petri run in the worker process. +//! +//! When `fabro run __run-worker` finds that its run's stored spec names +//! Petri as the engine, the run executes here instead of through the legacy +//! executor, over the same worker services: the authenticated client, the +//! control channel the server pushes cancels through, the signal handlers, +//! the vault snapshot and the CLI catalog. The engine assembly itself is +//! `fabro_petri::engine`, shared with the server's in-process test path, so +//! the run gets the same runtime, options and interviewer either way. +//! +//! The run's record is [`HttpRunStore`] over the worker's client, leased +//! for this launch: the worker mints one owner id at start, logs it, and +//! every lease the run takes over the API names it. `--mode start` loads +//! the admitted graphs through the client's blob read and runs them; +//! `--mode resume` continues the run from its records. Either way the +//! worker appends the lifecycle events Fabro's read side needs +//! (`run.starting`, `run.running`, then `run.completed` or `run.failed`) +//! through the client, as the legacy worker does. +//! +//! Of the server's controls, cancel is wired: the control channel's cancel +//! and `SIGTERM`/`SIGINT` fire one token, which cancels Petri's root +//! invocation politely. Pause, unpause and steer are received and ignored +//! with a warning until their Petri adapters land. A control channel that +//! is lost for good cancels the run the same way, and the worker exits with +//! that loss as its error once the run has settled. +//! +//! The runtime's settings layer is left empty here: the run's graphs were +//! lowered and admitted at create time with the server's layer, and nothing +//! lowers again at execution. The model client is built from the worker's +//! catalog and vault snapshot for the providers whose credentials resolve, +//! the same eligible set the legacy worker's LLM backend uses. + +use std::path::{Path, PathBuf}; +use std::sync::Arc; +use std::time::Instant; + +use anyhow::{Context, Result, anyhow, bail}; +use fabro_auth::VaultCredentialSource; +use fabro_client::{Client, ServerTarget}; +use fabro_interview::ControlInterviewer; +use fabro_llm::credentials::{CredentialProvider, readiness}; +use fabro_petri::engine::{self, Conclusion, Execution, RunRequest}; +use fabro_petri::petri::OwnerId; +use fabro_petri::runtime::{self, RuntimeSpec}; +use fabro_petri::{HttpRunStore, admission}; +use fabro_store::RunProjection; +use fabro_types::settings::run::RunMode; +use fabro_types::{FailureReason, RunId, RunTiming, StageOutcome, SuccessReason}; +use fabro_workflow::Error as WorkflowError; +use fabro_workflow::event::{self as workflow_event, Emitter, Event, RunEventSink}; +use fabro_workflow::run_control::RunControlState; +use fabro_workflow::runtime_store::RunStoreHandle; +use tokio_util::sync::CancellationToken; +use tracing::{info, warn}; + +use super::runner::{self, WorkerTitlePhase}; +use crate::args::RunWorkerMode; +use crate::command_context; + +/// What the worker holds when it hands a run to Petri. +pub(super) struct PetriWorker<'a> { + pub(super) run_id: RunId, + pub(super) target: ServerTarget, + pub(super) client: Client, + /// The legacy run store over the same client, which carries the + /// lifecycle events to the server with its retries. + pub(super) run_store: RunStoreHandle, + pub(super) run_state: RunProjection, + pub(super) storage_dir: &'a Path, + pub(super) run_dir: PathBuf, + pub(super) mode: RunWorkerMode, + pub(super) worker_token: &'a str, +} + +/// Execute the run to its end. `Ok` when the record says it succeeded; +/// the failure otherwise, after the terminal event is appended, so the +/// worker exits as the legacy worker does for a failed run. +pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { + let run_id = worker.run_id; + let Some(admission) = worker.run_state.spec.engine.petri().cloned() else { + bail!("run {run_id} names Petri as its engine but carries no admission"); + }; + let owner = OwnerId::mint(); + info!( + run_id = %run_id, + owner = %owner, + mode = ?worker.mode, + "Petri worker starting; every lease of this run names this launch" + ); + let store = Arc::new(HttpRunStore::for_worker( + worker.client.clone_for_reuse(), + owner, + )); + + let cancel_token = CancellationToken::new(); + let run_control = RunControlState::new(); + runner::install_signal_handlers(Arc::clone(&run_control), cancel_token.clone())?; + let interviewer = Arc::new(ControlInterviewer::new()); + let emitter = Arc::new(Emitter::new(run_id)); + let steering_hub = Arc::new(fabro_workflow::SteeringHub::new(Arc::clone(&emitter))); + let mut control_manager = runner::spawn_worker_control_manager( + worker.target.clone(), + run_id, + worker.worker_token.to_owned(), + interviewer, + cancel_token.clone(), + steering_hub, + run_control, + ); + control_manager.wait_for_first_connection().await?; + warn!( + run_id = %run_id, + "a Petri run answers cancel only: pause, unpause and steer are not wired yet and are ignored" + ); + let sink = RunEventSink::map( + runner::stamp_system_worker, + RunEventSink::backend(worker.run_store.clone()), + ); + + let runtime = runtime_spec(worker.storage_dir, &worker.run_state).await?; + let execution = match worker.mode { + RunWorkerMode::Start => { + let client = worker.client.clone_for_reuse(); + let graphs = admission::load_with( + |blob| { + let client = client.clone_for_reuse(); + async move { client.read_run_blob(&run_id, &blob).await } + }, + &admission, + ) + .await + .context("loading the admitted graphs")?; + Execution::Start(graphs) + } + RunWorkerMode::Resume => Execution::Resume, + }; + + let started = Instant::now(); + for event in [Event::RunStarting, Event::RunRunning] { + workflow_event::append_event_to_sink(&sink, &run_id, &event).await?; + } + runner::set_worker_title(&run_id, WorkerTitlePhase::Running); + + let request = RunRequest { + run_id: run_id.to_string(), + run_dir: worker.run_dir.join("petri"), + execution, + store, + runtime, + provider: worker + .run_state + .spec + .settings + .run + .environment + .provider + .clone(), + cancel: cancel_token.clone(), + }; + let run = Box::pin(engine::run(request)); + tokio::pin!(run); + let mut control_lost = None; + let result = loop { + tokio::select! { + result = &mut run => break result, + lost = control_manager.fatal_control_loss(), if control_lost.is_none() => { + // The server can no longer reach this worker: end the run + // politely, let Petri record why, then report the loss. + warn!(run_id = %run_id, error = %lost, "worker control lost; cancelling the Petri run"); + control_lost = Some(lost); + cancel_token.cancel(); + } + } + }; + control_manager.finish(); + + let timing = RunTiming { + wall_time_ms: u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX), + ..RunTiming::default() + }; + let (event, phase, failure) = match engine::conclusion(&result) { + Conclusion::Succeeded => { + info!(run_id = %run_id, "Petri run completed"); + ( + Event::WorkflowRunCompleted { + timing, + artifact_count: 0, + status: StageOutcome::Succeeded.to_string(), + reason: SuccessReason::Completed, + final_git_commit_sha: None, + final_patch: None, + diff_summary: None, + usage: None, + }, + WorkerTitlePhase::Succeeded, + None, + ) + } + Conclusion::Failed { reason, message } => { + info!(run_id = %run_id, error = %message, "Petri run did not succeed"); + let error = match reason { + FailureReason::Cancelled => WorkflowError::Cancelled, + _ => WorkflowError::engine(message.clone()), + }; + let phase = if reason == FailureReason::Cancelled { + WorkerTitlePhase::Cancelled + } else { + WorkerTitlePhase::Failed + }; + ( + Event::workflow_run_failed_from_error( + &error, timing, reason, None, None, None, None, + ), + phase, + Some(message), + ) + } + }; + workflow_event::append_event_to_sink(&sink, &run_id, &event).await?; + runner::set_worker_title(&run_id, phase); + if let Some(lost) = control_lost { + return Err(lost); + } + match failure { + None => Ok(()), + Some(message) => Err(anyhow!("Petri run failed: {message}")), + } +} + +/// The runtime the worker hands Petri: no settings layer (nothing lowers +/// at execution), the model client over the worker's catalog and vault for +/// the providers whose credentials resolve, and the run's mode. +async fn runtime_spec(storage_dir: &Path, run_state: &RunProjection) -> Result { + let catalog = + command_context::load_cli_catalog().context("failed to build worker LLM catalog")?; + let vault = runner::load_worker_vault(storage_dir).await?; + let credentials: Arc = Arc::new(VaultCredentialSource::new(vault)); + let ready = readiness(catalog.enabled_providers(), credentials.as_ref()).await; + for (provider, issue) in &ready.issues { + warn!(provider = %provider, error = %issue, "model provider credentials unusable"); + } + let model_client = match runtime::model_client(catalog, credentials, None, &ready.ready) { + Ok(client) => client, + Err(err) => { + warn!(error = %err, "Petri model client unavailable; LLM nodes run without one"); + None + } + }; + Ok(RuntimeSpec { + settings_toml: None, + model_client, + dry_run: run_state.spec.settings.run.execution.mode == RunMode::DryRun, + fabro_home: None, + }) +} diff --git a/lib/apps/fabro-cli/src/commands/run/runner.rs b/lib/apps/fabro-cli/src/commands/run/runner.rs index b67fd873d..92c84ce2b 100644 --- a/lib/apps/fabro-cli/src/commands/run/runner.rs +++ b/lib/apps/fabro-cli/src/commands/run/runner.rs @@ -44,6 +44,7 @@ use tokio_tungstenite::tungstenite::protocol::{self, Message as WebSocketMessage use tokio_tungstenite::{MaybeTlsStream, WebSocketStream, connect_async, tungstenite}; use tokio_util::sync::CancellationToken; +use super::petri_worker::{self, PetriWorker}; use crate::args::RunWorkerMode; use crate::shared::github::build_github_credentials; use crate::{command_context, server_client}; @@ -55,7 +56,7 @@ const RUN_STORE_RETRY_DELAYS: [Duration; 3] = [ ]; #[derive(Clone, Copy, Debug, PartialEq, Eq)] -enum WorkerTitlePhase { +pub(super) enum WorkerTitlePhase { Start, Resume, Init, @@ -85,6 +86,20 @@ pub(crate) async fn execute( .state() .await .with_context(|| format!("failed to load run state for {run_id}"))?; + if run_state.spec.engine.is_petri() { + return Box::pin(petri_worker::execute(PetriWorker { + run_id, + target, + client, + run_store, + run_state, + storage_dir: &storage_dir, + run_dir, + mode, + worker_token, + })) + .await; + } let run_spec = &run_state.spec; let catalog = Arc::new( command_context::load_cli_catalog().context("failed to build worker LLM catalog")?, @@ -237,7 +252,7 @@ fn build_fabro_run_tool_services( /// /// A worker always receives the server storage root so it can load the same /// secret vault as the server. -async fn load_worker_vault(storage_dir: &Path) -> Result>> { +pub(super) async fn load_worker_vault(storage_dir: &Path) -> Result>> { let storage = Storage::new(storage_dir); let vault = SecretStore::open_snapshot(storage.sqlite_path(), storage.secrets_path()) .await @@ -286,7 +301,7 @@ impl AppliedWorkerControlDeliveryIds { } } -struct WorkerControlManagerHandle { +pub(super) struct WorkerControlManagerHandle { first_connection: Option>>, fatal: Option>, done: CancellationToken, @@ -294,7 +309,7 @@ struct WorkerControlManagerHandle { } impl WorkerControlManagerHandle { - async fn wait_for_first_connection(&mut self) -> Result<()> { + pub(super) async fn wait_for_first_connection(&mut self) -> Result<()> { let receiver = self .first_connection .take() @@ -304,7 +319,7 @@ impl WorkerControlManagerHandle { .context("worker control manager stopped before first connection")? } - async fn fatal_control_loss(&mut self) -> anyhow::Error { + pub(super) async fn fatal_control_loss(&mut self) -> anyhow::Error { let Some(receiver) = self.fatal.take() else { return anyhow!("worker control fatal receiver missing"); }; @@ -313,7 +328,7 @@ impl WorkerControlManagerHandle { .unwrap_or_else(|_| anyhow!("worker control manager stopped before workflow completed")) } - fn finish(&self) { + pub(super) fn finish(&self) { self.done.cancel(); self.task.abort(); } @@ -380,7 +395,7 @@ enum WorkerControlConnectError { Other(anyhow::Error), } -fn spawn_worker_control_manager( +pub(super) fn spawn_worker_control_manager( target: ServerTarget, run_id: RunId, worker_token: String, @@ -1012,7 +1027,7 @@ impl RunStoreBackend for HttpRunStore { } } -fn set_worker_title(run_id: &RunId, phase: WorkerTitlePhase) { +pub(super) fn set_worker_title(run_id: &RunId, phase: WorkerTitlePhase) { fabro_proc::title_set(&worker_title(run_id, phase)); } @@ -1064,7 +1079,7 @@ fn update_worker_title_from_event(event: &RunEvent) { } } -fn stamp_system_worker(mut event: RunEvent) -> RunEvent { +pub(super) fn stamp_system_worker(mut event: RunEvent) -> RunEvent { if event.actor.is_none() { event.actor = Some(Principal::Worker { run_id: event.run_id, @@ -1123,7 +1138,7 @@ fn requires_github_credentials(run: &RunNamespace, has_repo_origin: bool) -> boo && has_repo_origin } -fn install_signal_handlers( +pub(super) fn install_signal_handlers( run_control: Arc, cancel_token: CancellationToken, ) -> Result<()> { diff --git a/lib/apps/fabro-cli/tests/it/scenario/mod.rs b/lib/apps/fabro-cli/tests/it/scenario/mod.rs index ec7320e43..81388b806 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/mod.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/mod.rs @@ -8,6 +8,7 @@ mod artifacts; mod auth; mod exec; mod lifecycle; +mod petri; mod server_lifecycle; mod smoke; diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri.rs b/lib/apps/fabro-cli/tests/it/scenario/petri.rs new file mode 100644 index 000000000..ec95672b1 --- /dev/null +++ b/lib/apps/fabro-cli/tests/it/scenario/petri.rs @@ -0,0 +1,568 @@ +//! Runs on Petri through a real server and its worker subprocess: the run +//! is created and started with `fabro run --detach`, the server launches +//! `fabro run __run-worker` for it as it does for a legacy run, and the +//! worker executes it through Petri over the HTTP run store. +//! +//! Each test starts its own foreground server on disk storage, because the +//! session's shared daemon keeps its object store in memory and the resume +//! scenario restarts the server. The runs take their host scope through the +//! sandbox-driver host plugin, so the tests skip, and say why, when the +//! executable is not found, unless `FABRO_REQUIRE_SANDBOX_PLUGINS` is set. +//! The plugin's path override crosses into the server and its workers the +//! way `PATH` does. + +#![expect( + clippy::disallowed_methods, + reason = "these scenarios start a real server subprocess, locate the plugin through the process environment, and poll processes" +)] +#![expect( + clippy::disallowed_types, + reason = "the scenarios own the server Child so they can SIGKILL it mid-run" +)] +#![expect(clippy::print_stderr, reason = "a skipped test says why on its stderr")] + +use std::env; +use std::path::{Path, PathBuf}; +use std::process::{Child, Command, Stdio}; +use std::time::{Duration, Instant}; + +use fabro_client::ServerTarget; +use fabro_config::{Storage, envfile}; +use fabro_petri::SqliteRunStore; +use fabro_petri::engine::{self, RunStatus}; +use fabro_petri::petri::RunKey; +use fabro_static::EnvVars; +use fabro_store::EventEnvelope; +use fabro_test::{apply_test_isolation, expect_reqwest_json, isolated_storage_dir, test_context}; +use fabro_types::EventBody; + +use crate::cmd::support::created_run_id; +use crate::support::{ + TEST_DEV_TOKEN, TEST_SESSION_SECRET, parse_event_envelopes, seed_dev_token_auth, +}; + +const HOST_PLUGIN: &str = "sandbox-driver-host"; +const REQUIRE_ENV: &str = "FABRO_REQUIRE_SANDBOX_PLUGINS"; +const RUN_TIMEOUT: Duration = Duration::from_mins(1); +const POLL: Duration = Duration::from_millis(50); + +/// The host plugin as Petri's lookup finds it: the override variable, else +/// the executable on `PATH`. `None`, after saying so, when the test should +/// skip; a panic when the environment forbids a skip. +fn host_plugin() -> Option { + let found = env::var_os(EnvVars::PETRI_SANDBOX_HOST_PLUGIN) + .map(PathBuf::from) + .or_else(|| { + env::split_paths(&env::var_os(EnvVars::PATH)?) + .map(|dir| dir.join(HOST_PLUGIN)) + .find(|candidate| candidate.is_file()) + }); + if found.is_none() { + assert!( + env::var_os(REQUIRE_ENV).is_none(), + "{REQUIRE_ENV} is set, but {HOST_PLUGIN} is not on PATH and {} is unset", + EnvVars::PETRI_SANDBOX_HOST_PLUGIN + ); + eprintln!( + "skipping: {HOST_PLUGIN} is not on PATH and {} is unset", + EnvVars::PETRI_SANDBOX_HOST_PLUGIN + ); + } + found +} + +/// A foreground server on its own disk storage, dev-token auth, started +/// from the compiled `fabro` binary. Dropping it kills the process. +struct RunningServer { + child: Option, + home_root: tempfile::TempDir, + _storage_root: tempfile::TempDir, + storage_dir: PathBuf, + config_path: PathBuf, + port: u16, + api_base_url: String, +} + +impl RunningServer { + async fn start() -> Self { + 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, + "_version = 1\n\n[server.auth]\nmethods = [\"dev-token\"]\n", + ) + .expect("the server settings write"); + let runtime_directory = Storage::new(&storage_dir).runtime_directory(); + envfile::merge_env_file(&runtime_directory.env_path(), [ + ("SESSION_SECRET", TEST_SESSION_SECRET), + ("FABRO_DEV_TOKEN", TEST_DEV_TOKEN), + ]) + .expect("the server env writes"); + fabro_util::dev_token::write_dev_token(&runtime_directory.dev_token_path(), TEST_DEV_TOKEN) + .expect("the dev token writes"); + let mut server = Self { + child: None, + home_root, + _storage_root: storage_root, + storage_dir, + config_path, + port, + api_base_url: format!("http://127.0.0.1:{port}"), + }; + server.launch().await; + server + } + + /// Start the server process over this storage; the same call brings + /// it back after a kill. + async fn launch(&mut self) { + assert!(self.child.is_none(), "the server is already running"); + let mut cmd = Command::new(env!("CARGO_BIN_EXE_fabro")); + apply_test_isolation(&mut cmd, self.home_root.path()); + // The resume scenario restarts the server: its object store must + // outlive the process. + cmd.env(EnvVars::FABRO_TEST_IN_MEMORY_STORE, "0"); + cmd.env( + EnvVars::FABRO_HOME, + self.home_root.path().join("fabro-home"), + ); + cmd.args(["server", "start", "--foreground"]) + .arg("--storage-dir") + .arg(&self.storage_dir) + .arg("--bind") + .arg(format!("127.0.0.1:{}", self.port)) + .arg("--config") + .arg(&self.config_path) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(self.stderr_log()); + let mut child = cmd.spawn().expect("the server spawns"); + wait_for_http_ready(&self.api_base_url, &mut child).await; + self.child = Some(child); + } + + /// Where the server's stderr goes: a file beside its storage, so a + /// chatty server never blocks on a pipe nobody reads, and a failing + /// test can show it. + fn stderr_log(&self) -> Stdio { + let path = self.storage_dir.with_file_name("server.stderr.log"); + let file = std::fs::OpenOptions::new() + .create(true) + .append(true) + .open(path) + .expect("the server stderr log opens"); + Stdio::from(file) + } + + fn stderr_text(&self) -> String { + std::fs::read_to_string(self.storage_dir.with_file_name("server.stderr.log")) + .unwrap_or_default() + } + + /// The `--server` target a CLI command reaches this server at. + fn target(&self) -> String { + format!("{}/api/v1", self.api_base_url) + } + + /// Kill the server outright, as a crash would; its workers live on in + /// their own process groups. + fn kill(&mut self) { + let mut child = self.child.take().expect("the server is running"); + child.kill().expect("the server dies"); + let _ = child.wait(); + } + + fn shutdown(mut self) { + let mut stop = Command::new(env!("CARGO_BIN_EXE_fabro")); + apply_test_isolation(&mut stop, self.home_root.path()); + stop.args(["server", "stop"]) + .arg("--storage-dir") + .arg(&self.storage_dir); + let output = stop.output().expect("server stop runs"); + assert!( + output.status.success(), + "server stop failed\nstdout:\n{}\nstderr:\n{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + let status = self + .child + .take() + .expect("the server is running") + .wait() + .expect("the server exit status reads"); + assert!( + status.success(), + "the server exited unsuccessfully\nstderr:\n{}", + self.stderr_text() + ); + } + + /// Petri's store over the server's database, read beside the server: + /// what `petri inspect` would see. + async fn petri_store(&self) -> SqliteRunStore { + let database = fabro_db::Database::connect(Storage::new(&self.storage_dir).sqlite_path()) + .await + .expect("the server database opens"); + SqliteRunStore::new(database.clone_pool()) + } +} + +impl Drop for RunningServer { + fn drop(&mut self) { + if let Some(child) = self.child.as_mut() { + if child.try_wait().ok().flatten().is_none() { + let _ = child.kill(); + let _ = child.wait(); + } + } + } +} + +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) { + 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") { + panic!("the server exited before it was ready with status {status}"); + } + 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}"), + } + } +} + +/// 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 { + let workspace = context.temp_dir.join("petri-workspace"); + std::fs::create_dir_all(&workspace).expect("the workspace creates"); + std::fs::write( + workspace.join("workflow.fabro"), + format!( + "digraph Command {{\n graph [goal=\"Run one command\", default_max_retries=0]\n start \ + [shape=Mdiamond]\n exit [shape=Msquare]\n say [shape=parallelogram, \ + script=\"{script}\", max_retries=0]\n start -> say -> exit\n}}\n" + ), + ) + .expect("the workflow writes"); + std::fs::write( + workspace.join("workflow.toml"), + "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\nengine = \"petri\"\n\n[run]\ngoal \ + = \"Run one command\"\n", + ) + .expect("the settings write"); + workspace +} + +/// `fabro run --detach` against the server: the run is created and started, +/// and its id comes back. +fn run_detached( + context: &fabro_test::TestContext, + server: &RunningServer, + workspace: &Path, +) -> String { + let target = server.target(); + seed_dev_token_auth( + &context.home_dir, + &ServerTarget::http_url(&target).expect("the target parses"), + TEST_DEV_TOKEN, + ); + let output = context + .run_cmd() + .current_dir(workspace) + .args([ + "--server", + &target, + "--detach", + "--auto-approve", + "--environment", + "local", + "workflow.toml", + ]) + .output() + .expect("the detached run executes"); + assert!( + output.status.success(), + "detached run failed\nstdout:\n{}\nstderr:\n{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + created_run_id(&output) +} + +async fn run_json(server: &RunningServer, path: &str) -> serde_json::Value { + let response = fabro_test::test_http_client() + .get(format!("{}/api/v1/{path}", server.api_base_url)) + .bearer_auth(TEST_DEV_TOKEN) + .send() + .await + .expect("the request sends"); + expect_reqwest_json( + response, + fabro_http::StatusCode::OK, + format!("GET /api/v1/{path}"), + ) + .await +} + +async fn run_status(server: &RunningServer, run_id: &str) -> String { + run_json(server, &format!("runs/{run_id}")).await["lifecycle"]["status"]["kind"] + .as_str() + .expect("the run has a status kind") + .to_string() +} + +async fn wait_for_status(server: &RunningServer, run_id: &str, expected: &[&str]) -> String { + let deadline = Instant::now() + RUN_TIMEOUT; + loop { + let status = run_status(server, run_id).await; + if expected.contains(&status.as_str()) { + return status; + } + assert!( + Instant::now() < deadline, + "run {run_id} did not reach {expected:?}; last status {status}" + ); + tokio::time::sleep(POLL).await; + } +} + +async fn run_events(server: &RunningServer, run_id: &str) -> Vec { + parse_event_envelopes(&run_json(server, &format!("runs/{run_id}/events")).await) +} + +fn event_names(events: &[EventEnvelope]) -> Vec<&str> { + events + .iter() + .map(|envelope| envelope.event.event_name()) + .collect() +} + +/// The pid of the worker subprocess the server launched for the run: the +/// worker retitles itself `fabro `, so that +/// is what the process table shows. +fn worker_pid(run_id: &str) -> Option { + let short_id: String = run_id.chars().take(12).collect(); + let output = Command::new("pgrep") + .args(["-f", &format!("^fabro {short_id} ")]) + .output() + .expect("pgrep runs"); + String::from_utf8_lossy(&output.stdout) + .lines() + .find_map(|line| line.trim().parse().ok()) +} + +fn wait_for_worker(run_id: &str) -> u32 { + let deadline = Instant::now() + RUN_TIMEOUT; + loop { + if let Some(pid) = worker_pid(run_id) { + return pid; + } + assert!( + Instant::now() < deadline, + "no worker process appeared for run {run_id}" + ); + std::thread::sleep(POLL); + } +} + +/// Whether a process is waiting on the gate file: the stage is mid-flight. +fn gate_is_polled(gate: &Path) -> bool { + let output = Command::new("pgrep") + .args(["-f", &gate.display().to_string()]) + .output() + .expect("pgrep runs"); + output.status.success() +} + +fn wait_until_gate_is_polled(gate: &Path) { + let deadline = Instant::now() + RUN_TIMEOUT; + while !gate_is_polled(gate) { + assert!( + Instant::now() < deadline, + "the stage never started waiting on {}", + gate.display() + ); + std::thread::sleep(POLL); + } +} + +/// A command-only Petri bundle runs to completion in the worker the server +/// launched: Fabro reports the run succeeded, the worker wrote Petri's +/// records through the HTTP store, and its lease ended with it. +#[tokio::test(flavor = "multi_thread")] +async fn a_petri_run_executes_in_the_server_launched_worker() { + if host_plugin().is_none() { + return; + } + let context = test_context!(); + let server = RunningServer::start().await; + let workspace = write_petri_workspace(&context, "echo hello from petri"); + let run_id = run_detached(&context, &server, &workspace); + + let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await; + let run = run_json(&server, &format!("runs/{run_id}")).await; + assert_eq!(status, "succeeded", "run: {run}"); + let state = run_json(&server, &format!("runs/{run_id}/state")).await; + assert_eq!(state["spec"]["engine"]["kind"], "petri", "state: {state}"); + + let events = run_events(&server, &run_id).await; + let names = event_names(&events); + assert_eq!( + names + .iter() + .filter(|name| **name == "run.completed") + .count(), + 1, + "{names:?}" + ); + assert!( + names.contains(&"run.starting") && names.contains(&"run.running"), + "{names:?}" + ); + + let store = server.petri_store().await; + let outcome = engine::outcome_of(&store, &run_id) + .await + .expect("the run's Petri record inspects"); + assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}"); + assert!(outcome.complete, "{:?}", outcome.incomplete); + let key = RunKey::new(run_id.as_str()); + let deadline = Instant::now() + Duration::from_secs(10); + loop { + let holder = store.owner(&key).await.expect("the lease reads"); + if holder.is_none() { + break; + } + assert!( + Instant::now() < deadline, + "the worker's lease {holder:?} outlived the worker" + ); + tokio::time::sleep(POLL).await; + } + server.shutdown(); +} + +/// A Petri run whose worker and server both die mid-stage continues after +/// the server restarts: the new server releases the dead worker's lease, +/// asks the run to start again as a resume, and launches a worker in +/// resume mode, which finishes the run with one `run.completed`. +#[tokio::test(flavor = "multi_thread")] +async fn a_petri_run_resumes_in_a_new_worker_after_the_server_restarts() { + if host_plugin().is_none() { + return; + } + let context = test_context!(); + let mut server = RunningServer::start().await; + let gate = context.temp_dir.join("resume.gate"); + let script = format!("while [ ! -f {} ]; do sleep 0.05; done", gate.display()); + let workspace = write_petri_workspace(&context, &script); + let run_id = run_detached(&context, &server, &workspace); + + wait_for_status(&server, &run_id, &["running"]).await; + eprintln!("run {run_id} is running"); + let worker = wait_for_worker(&run_id); + eprintln!("worker {worker} launched"); + wait_until_gate_is_polled(&gate); + eprintln!("stage is waiting on the gate"); + + // The crash: the server first, so it never observes the worker exit, + // then the worker's whole process group, plugin and stage included. + server.kill(); + fabro_proc::sigkill_process_group(worker); + let deadline = Instant::now() + Duration::from_secs(10); + while fabro_proc::process_running(worker) { + assert!(Instant::now() < deadline, "the worker did not die"); + std::thread::sleep(POLL); + } + assert_eq!( + run_status_offline(&server).await, + None, + "the server is down" + ); + + server.launch().await; + eprintln!("server restarted"); + let status = wait_for_status(&server, &run_id, &["running", "succeeded", "failed"]).await; + eprintln!("run {run_id} is {status} after the restart"); + let resumed = wait_for_worker(&run_id); + assert_ne!(resumed, worker, "a new worker was launched"); + eprintln!("worker {resumed} launched for the resume"); + wait_until_gate_is_polled(&gate); + eprintln!("stage is waiting on the gate again"); + std::fs::write(&gate, "go").expect("the gate opens"); + + let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await; + let events = run_events(&server, &run_id).await; + let names = event_names(&events); + assert_eq!( + status, + "succeeded", + "events: {names:?}\nserver stderr:\n{}", + server.stderr_text() + ); + assert_eq!( + names + .iter() + .filter(|name| **name == "run.completed") + .count(), + 1, + "{names:?}" + ); + // `fabro run` asked for the first start; the restart asked for a + // resume, after the run had been running. + let first_running = names + .iter() + .position(|name| *name == "run.running") + .expect("the run ran before the crash"); + let resume_request = events + .iter() + .position(|envelope| { + matches!( + &envelope.event.body, + EventBody::RunStartRequested(props) if props.resume + ) + }) + .expect("the restart asked for a resume"); + assert!(resume_request > first_running, "{names:?}"); + assert_eq!( + names.iter().filter(|name| **name == "run.running").count(), + 2, + "the run ran once before and once after the restart: {names:?}" + ); + + let store = server.petri_store().await; + let outcome = engine::outcome_of(&store, &run_id) + .await + .expect("the run's Petri record inspects"); + assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}"); + assert!(outcome.complete, "{:?}", outcome.incomplete); + server.shutdown(); +} + +/// The run's status while the server may be down: `None` when it is. +async fn run_status_offline(server: &RunningServer) -> Option { + fabro_test::test_http_client() + .get(format!("{}/health", server.api_base_url)) + .send() + .await + .ok() + .map(|response| response.status().to_string()) +}