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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-17 21:20:22 -04:00
parent b5ebb472ae
commit 16354186fa
No known key found for this signature in database
7 changed files with 855 additions and 10 deletions

3
Cargo.lock generated
View file

@ -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",

View file

@ -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"] }

View file

@ -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;

View file

@ -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<RuntimeSpec> {
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<dyn CredentialProvider> = 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,
})
}

View file

@ -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<Arc<AsyncRwLock<Vault>>> {
pub(super) async fn load_worker_vault(storage_dir: &Path) -> Result<Arc<AsyncRwLock<Vault>>> {
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<oneshot::Receiver<Result<()>>>,
fatal: Option<oneshot::Receiver<anyhow::Error>>,
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<RunControlState>,
cancel_token: CancellationToken,
) -> Result<()> {

View file

@ -8,6 +8,7 @@ mod artifacts;
mod auth;
mod exec;
mod lifecycle;
mod petri;
mod server_lifecycle;
mod smoke;

View file

@ -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<PathBuf> {
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<Child>,
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<EventEnvelope> {
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 <first 12 of the run id> <phase>`, so that
/// is what the process table shows.
fn worker_pid(run_id: &str) -> Option<u32> {
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<String> {
fabro_test::test_http_client()
.get(format!("{}/health", server.api_base_url))
.send()
.await
.ok()
.map(|response| response.status().to_string())
}