From b5ebb472aec4b12529df2c5ff48e8d4469e0ea1d Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 17 Sep 2026 21:20:22 -0400 Subject: [PATCH] Launch a worker for a Petri run and resume it after a server restart `execute_run` no longer runs a Petri run in the server process by default: it takes the subprocess path a legacy run takes, and `worker_exited` still releases the worker's lease when the process ends. The in-process path stays under the handler-registry test override, so the scenario tests need no worker binary; it now honours the managed run's execution mode. At startup, `reconcile_incomplete_runs_on_startup` hands a Petri run the previous server left in flight (runnable, starting, running, blocked or paused, with no cancel pending) back to a worker instead of failing it: `PetriRuns::release_for_restart` ends the dead worker's lease from outside, which fences it should it still be alive, the run is asked to start again as a resume (`run.start_requested` with `resume`, then `run.runnable`, the pair the API's resume appends), and the managed run is registered in resume mode when Petri's store holds the run, else in start mode. Full workspace recovery is the plan's F3.5 and is noted in the module docs. Tests: the restart reconcile releases the lease, rewrites the history, and launches the worker with `--mode resume`; a worker's HTTP store leases for its launch id over the loopback server. Co-Authored-By: Claude Fable 5.1 --- lib/apps/fabro-server/src/petri_runs.rs | 163 ++++++++++++++++-- lib/apps/fabro-server/src/server.rs | 33 +++- .../fabro-server/src/server/petri_runs.rs | 160 +++++++++++------ .../fabro-server/tests/it/api/petri_store.rs | 41 +++++ .../fabro-server/tests/it/scenario/petri.rs | 10 ++ 5 files changed, 346 insertions(+), 61 deletions(-) diff --git a/lib/apps/fabro-server/src/petri_runs.rs b/lib/apps/fabro-server/src/petri_runs.rs index a3b22ae0a..9cacaefce 100644 --- a/lib/apps/fabro-server/src/petri_runs.rs +++ b/lib/apps/fabro-server/src/petri_runs.rs @@ -108,6 +108,16 @@ impl PetriRuns { drop(handle); } + /// End whatever lease the run's previous worker held, from outside: + /// what the server does for a run it finds in flight at startup, before + /// it launches a new worker for it. The previous worker, should it still + /// be alive, finds its handles stale on its next write. `NotFound` when + /// the store never held the run. + pub(crate) async fn release_for_restart(&self, run_id: RunId) -> Result<(), StoreError> { + self.worker_exited(run_id); + self.store.release_lease(&Self::key(&run_id)).await + } + /// Drop every handle held on the run: what the server does when it /// observes the run's worker exit, so a worker that died without /// releasing does not keep the lease. @@ -141,6 +151,7 @@ fn lock(mutex: &Mutex) -> MutexGuard<'_, T> { #[cfg(test)] mod tests { + use std::collections::BTreeMap; use std::pin::Pin; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; @@ -153,7 +164,8 @@ mod tests { use fabro_config::daemon::ServerDaemon; use fabro_petri::petri::RunStore as _; use fabro_static::EnvVars; - use fabro_types::{RunId, WorkflowPath, WorkflowVersion}; + use fabro_types::{RunId, RunStatus, WorkflowPath, WorkflowVersion}; + use fabro_workflow::event::{Event, append_event}; use serde_json::json; use tokio::io::AsyncRead; use tokio::sync::Notify; @@ -161,9 +173,10 @@ mod tests { use tower::ServiceExt as _; use super::*; - use crate::server::{AppState, spawn_scheduler}; + use crate::server::{AppState, reconcile_incomplete_runs_on_startup, spawn_scheduler}; use crate::test_support::{ TestAppStateBuilder, build_test_router, test_register_workflow_version, + test_secret_store_path, test_store_bundle, }; use crate::worker_runtime::{ StartedWorker, WorkerExit, WorkerLaunchSpec, WorkerRef, WorkerRuntime, @@ -176,13 +189,18 @@ mod tests { start -> exit }"#; + const PETRI_SETTINGS: &str = + "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\nengine = \"petri\"\n"; + /// A worker runtime whose one worker runs until the test ends it, so - /// the test can act while the server waits on the worker. + /// the test can act while the server waits on the worker. It keeps the + /// mode the server launched the worker with. #[derive(Default)] struct HeldWorkerRuntime { started: Notify, running: AtomicBool, exit: Arc, + mode: Mutex>, } impl HeldWorkerRuntime { @@ -196,11 +214,16 @@ mod tests { self.running.store(false, Ordering::SeqCst); self.exit.notify_one(); } + + fn launched_mode(&self) -> Option<&'static str> { + *lock(&self.mode) + } } #[async_trait::async_trait] impl WorkerRuntime for HeldWorkerRuntime { - async fn start(&self, _spec: WorkerLaunchSpec) -> anyhow::Result { + async fn start(&self, spec: WorkerLaunchSpec) -> anyhow::Result { + *lock(&self.mode) = Some(spec.mode); self.running.store(true, Ordering::SeqCst); let exit = Arc::clone(&self.exit); let stderr: Pin> = Box::pin(tokio::io::empty()); @@ -248,15 +271,27 @@ mod tests { .expect("the test server record writes"); } - /// A run created and started through the API, as a client would. + /// A legacy run created and started through the API, as a client would. async fn create_and_start_run(app: &axum::Router) -> RunId { + create_and_start_run_with(app, &[]).await + } + + /// A Petri run: its version's `workflow.toml` names the engine. + async fn create_and_start_petri_run(app: &axum::Router) -> RunId { + create_and_start_run_with(app, &[("workflow.toml", PETRI_SETTINGS)]).await + } + + async fn create_and_start_run_with(app: &axum::Router, extra: &[(&str, &str)]) -> RunId { let path = WorkflowPath::new("workflow.fabro").expect("a workflow path"); - let version = WorkflowVersion::new( - path.clone(), - std::collections::BTreeMap::from([(path, MINIMAL_DOT.to_string())]), - std::collections::BTreeMap::new(), - ) - .expect("a workflow version"); + let mut files = BTreeMap::from([(path.clone(), MINIMAL_DOT.to_string())]); + for (name, text) in extra { + files.insert( + WorkflowPath::new(*name).expect("a workflow path"), + (*text).to_string(), + ); + } + let version = + WorkflowVersion::new(path, files, BTreeMap::new()).expect("a workflow version"); let version_id = test_register_workflow_version(app, &version, None).await; let intent = json!({ "workflow_version_id": version_id, @@ -360,4 +395,110 @@ mod tests { .expect("the next owner takes the run"); drop(resumed); } + + /// After a restart, a Petri run the previous server left running goes + /// back to a worker in resume mode: the lease its worker held is + /// released from outside, the run is asked to start again as a resume, + /// and the scheduler launches the worker with `--mode resume`. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn a_petri_run_left_running_by_a_restart_goes_back_to_a_worker_in_resume_mode() { + let (store, artifact_store) = test_store_bundle(); + let vault_path = test_secret_store_path(); + let before = TestAppStateBuilder::new() + .store_bundle(Arc::clone(&store), artifact_store.clone()) + .vault_path(vault_path.clone()) + .vault_entries([(EnvVars::OPENAI_API_KEY, "test-openai-api-key")]) + .build(); + let app = build_test_router(Arc::clone(&before)); + let run_id = create_and_start_petri_run(&app).await; + let key = PetriRuns::key(&run_id); + + // The worker took the run as far as running and holds its lease; + // then the server died, so nothing released it. + let run_store = before + .stores + .runs + .open_run(&run_id) + .await + .expect("the run opens"); + for event in [Event::RunStarting, Event::RunRunning] { + append_event(&run_store, &run_id, &event) + .await + .expect("the lifecycle event appends"); + } + let held = before + .petri_runs + .open(run_id, Access::Create { + owner: OwnerId::new("worker-1"), + }) + .await + .expect("the worker takes the run"); + drop(held); + assert_eq!( + before + .petri_runs + .store() + .owner(&key) + .await + .expect("reads the lease"), + Some(OwnerId::new("worker-1")) + ); + + let runtime = Arc::new(HeldWorkerRuntime::default()); + let after = TestAppStateBuilder::new() + .store_bundle(store, artifact_store) + .vault_path(vault_path) + .vault_entries([(EnvVars::OPENAI_API_KEY, "test-openai-api-key")]) + .worker_runtime(Arc::clone(&runtime) as Arc) + .build(); + let reconciled = reconcile_incomplete_runs_on_startup(&after) + .await + .expect("the restart reconciles"); + assert_eq!(reconciled, 1); + + assert_eq!( + after + .petri_runs + .store() + .owner(&key) + .await + .expect("reads the lease"), + None, + "the previous worker's lease is released" + ); + let reader = after + .stores + .runs + .open_run_reader(&run_id) + .await + .expect("the run opens for reading"); + let run_state = reader.state().await.expect("the run state loads"); + assert_eq!(run_state.status, RunStatus::Runnable); + let names = reader + .list_events() + .await + .expect("the history lists") + .into_iter() + .map(|envelope| envelope.event.event_name().to_string()) + .collect::>(); + assert_eq!( + &names[names.len() - 4..], + [ + "run.starting", + "run.running", + "run.start_requested", + "run.runnable" + ], + "{names:?}" + ); + + write_test_server_record(&after); + spawn_scheduler(Arc::clone(&after)); + runtime.wait_for_start().await; + assert_eq!(runtime.launched_mode(), Some("resume")); + runtime.end_worker(); + // The first server's handles must outlive the check above: a real + // crash releases nothing, and dropping them here would. + drop(before); + } } diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index ea58db9c5..21a57cf24 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -3143,6 +3143,17 @@ pub(crate) async fn reconcile_incomplete_runs_on_startup( for summary in summaries { let run_store = state.stores.runs.open_run(&summary.id).await?; + // A Petri run continues from its records in a new worker, unless a + // cancel was pending or the run was being removed: those end as a + // legacy run's do. + if petri_run_resumes_on_restart(&summary) { + let run_state = run_store.state().await?; + if run_state.spec.engine.is_petri() { + petri_runs::reconcile_on_startup(state, summary.id, &run_store, &run_state).await?; + reconciled += 1; + continue; + } + } let (error, reason) = failure_for_incomplete_run( summary.lifecycle.pending_control, "Fabro server restarted before the run reached a terminal state.".to_string(), @@ -3163,6 +3174,21 @@ pub(crate) async fn reconcile_incomplete_runs_on_startup( Ok(reconciled) } +/// Whether a run the server finds in flight at startup is one a Petri +/// worker can continue: it was runnable or running (blocked or paused +/// count), no cancel was pending, and it was not being removed. +fn petri_run_resumes_on_restart(summary: &fabro_types::Run) -> bool { + summary.lifecycle.pending_control != Some(RunControlAction::Cancel) + && matches!( + summary.lifecycle.status, + RunStatus::Runnable + | RunStatus::Starting + | RunStatus::Running + | RunStatus::Blocked { .. } + | RunStatus::Paused { .. } + ) +} + fn live_worker_processes(state: &AppState) -> Vec { let runs = state.runs.lock().expect("runs lock poisoned"); runs.iter() @@ -3982,12 +4008,15 @@ async fn execute_run(state: Arc, run_id: RunId) { return; } + // A Petri run takes the worker path a legacy run takes. Under the test + // override it executes in this process instead, so the scenario tests + // need no worker binary. match run_engine(&state, run_id).await { - Ok(Engine::Petri) => { + Ok(Engine::Petri) if state.registry_factory_override.is_some() => { Box::pin(petri_runs::execute(state, run_id)).await; return; } - Ok(Engine::Legacy) => {} + Ok(Engine::Petri | Engine::Legacy) => {} Err(err) => { tracing::error!(run_id = %run_id, error = %err, "Failed to read the run's engine"); fail_managed_run( diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index b6944865a..931a8a459 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -7,26 +7,36 @@ //! the blob store so the run executes and resumes from what was admitted. //! Petri compiled, linted and pinned models; the legacy compile is skipped. //! -//! At execution, [`execute`] runs the admitted graph through -//! `fabro_petri::engine` in the server process, over the run store in the -//! server's database, until the worker's HTTP run store lands. Only the run -//! lifecycle events Fabro's read side needs are appended (`run.starting`, -//! `run.running`, then `run.completed` or `run.failed`); no stage or agent -//! event is projected, which is the read-side item that follows. +//! At execution, a Petri run takes the same path a legacy run does: the +//! scheduler launches `fabro run __run-worker` with the worker's token, and +//! the worker executes the run through `fabro_petri::engine` over the HTTP +//! run store, appending the run lifecycle events Fabro's read side needs +//! (`run.starting`, `run.running`, then `run.completed` or `run.failed`). +//! The server keeps the worker's lease for as long as the worker lives +//! (`crate::petri_runs`). Under the test override that replaces the handler +//! registry, [`execute`] runs the same engine in the server process over the +//! run store in the server's database, so the scenario tests need no +//! worker binary. No stage or agent event is projected either way, which is +//! the read-side item that follows. +//! +//! After a server restart, [`reconcile_on_startup`] hands a Petri run the +//! previous server left in flight back to a worker in resume mode. use std::collections::{BTreeMap, HashSet}; use std::sync::Arc; use std::time::Instant; -use fabro_config::SettingsLayer; +use fabro_config::{SettingsLayer, Storage}; use fabro_llm::selection; use fabro_petri::check::{self, Bundle, CheckError, CheckRequest, Diagnostic, Launch}; -use fabro_petri::engine::{self, RunOutcome, RunRequest}; +use fabro_petri::engine::{self, Conclusion, Execution, RunRequest}; +use fabro_petri::petri::StoreError; use fabro_petri::runtime::{self, RuntimeSpec}; use fabro_petri::{SqliteRunStore, admission}; use fabro_types::settings::run::RunMode; use fabro_types::{ - Engine, PetriAdmission, RunId, RunTarget, RunTiming, ServerSettings, StageOutcome, + Engine, PetriAdmission, RunId, RunRunnableSource, RunTarget, RunTiming, ServerSettings, + StageOutcome, }; use fabro_util::error as error_util; use fabro_validate::{Diagnostic as FabroDiagnostic, Severity}; @@ -37,7 +47,8 @@ use tokio::task; use tokio_util::sync::CancellationToken; use tracing::{error, info, warn}; -use super::{AppState, clear_live_run_state, workflow_event}; +use super::{AppState, RunExecutionMode, clear_live_run_state, workflow_event}; +use crate::petri_runs::PetriRuns; use crate::run_compiler::{PreparedRun, RunCompilerError}; /// The engine a run gets: the one its workflow version names, else the @@ -215,10 +226,12 @@ fn fabro_diagnostic(diagnostic: &Diagnostic) -> FabroDiagnostic { } } -/// Execute a Petri run in the server process: runnable → starting → running -/// → succeeded or failed, with the lifecycle events Fabro's read side needs. +/// Execute a Petri run in the server process, under the test override: +/// runnable → starting → running → succeeded or failed, with the lifecycle +/// events Fabro's read side needs. Outside tests a Petri run executes in +/// its worker process, launched as a legacy run's worker is. pub(crate) async fn execute(state: Arc, run_id: RunId) { - let (run_dir, cancel) = { + let (run_dir, cancel, mode) = { let mut runs = state.runs.lock().expect("runs lock poisoned"); let managed_run = match runs.get_mut(&run_id) { Some(run) if run.status == RunStatus::Runnable => run, @@ -230,7 +243,7 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { let cancel = CancellationToken::new(); managed_run.status = RunStatus::Starting; managed_run.cancel_token = Some(cancel.clone()); - (run_dir, cancel) + (run_dir, cancel, managed_run.execution_mode) }; let run_store = match state.stores.runs.open_run(&run_id).await { @@ -283,6 +296,19 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { { return; } + let execution = match mode { + RunExecutionMode::Start => { + match admission::load(&state.store_ref().blobs(), &admission).await { + Ok(graphs) => Execution::Start(graphs), + Err(err) => { + let message = error_util::collect_chain(&err).join(": "); + fail_before_execution(&state, &run_store, run_id, &message).await; + return; + } + } + } + RunExecutionMode::Resume => Execution::Resume, + }; let started = Instant::now(); for event in [ workflow_event::Event::RunStarting, @@ -314,24 +340,19 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { let request = RunRequest { run_id: run_id.to_string(), run_dir: run_dir.join("petri"), - admission, - blobs: state.store_ref().blobs(), + execution, store: Arc::new(SqliteRunStore::new(state.db_pool.clone())), runtime: runtime_spec(&state, &eligible, dry_run), provider: run_state.spec.settings.run.environment.provider.clone(), cancel, }; - let outcome = Box::pin(engine::run(request)).await; + let result = Box::pin(engine::run(request)).await; let timing = RunTiming { wall_time_ms: u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX), ..RunTiming::default() }; - let (status, error, event) = match outcome { - Ok(RunOutcome { - status: engine::RunStatus::Success, - complete: true, - .. - }) => { + let (status, error, event) = match engine::conclusion(&result) { + Conclusion::Succeeded => { info!(run_id = %run_id, "Petri run completed"); ( RunStatus::Succeeded { @@ -350,22 +371,10 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { }, ) } - Ok(outcome) => { - let reason = match outcome.status { - engine::RunStatus::Cancelled => FailureReason::Cancelled, - engine::RunStatus::Success | engine::RunStatus::Failed => { - FailureReason::WorkflowError - } - }; - let message = failure_message(&outcome); + Conclusion::Failed { reason, message } => { info!(run_id = %run_id, error = %message, "Petri run did not succeed"); failed(reason, message, timing) } - Err(err) => { - let message = error_util::collect_chain(&err).join(": "); - error!(run_id = %run_id, error = %message, "Petri run failed"); - failed(FailureReason::WorkflowError, message, timing) - } }; if let Err(err) = workflow_event::append_event(&run_store, &run_id, &event).await { error!(run_id = %run_id, error = %err, "Failed to persist run outcome"); @@ -373,20 +382,75 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { finish(&state, run_id, status, error); } -/// The failure of a run whose record says it did not succeed. -fn failure_message(outcome: &RunOutcome) -> String { - let mut message = match (&outcome.status, &outcome.failure) { - (engine::RunStatus::Cancelled, _) => "the run was cancelled".to_string(), - (_, Some(failure)) => failure.clone(), - (engine::RunStatus::Failed, None) => "the run failed".to_string(), - (engine::RunStatus::Success, None) => "the run's record is incomplete".to_string(), +/// Bring a Petri run the server left in flight back to its worker after a +/// restart: the run continues from its records, as Petri's own resume does. +/// +/// The lease the previous worker held is released from outside, which +/// fences that worker should it still be alive; then the run is asked to +/// start again as a resume (`run.start_requested` with `resume`, then +/// `run.runnable`, the same pair the API's resume appends), and a managed +/// run is registered for the scheduler in resume mode when Petri's store +/// holds the run, else in start mode: a worker that died before it created +/// the run's record left nothing to continue from, so the run starts from +/// its admitted graphs. +/// +/// Full recovery, where the workspace a resumed stage sees is restored to +/// the snapshot its durable state names, is the integration plan's F3.5. +/// Until it lands, a retained workspace is used as the previous worker left +/// it. +pub(crate) async fn reconcile_on_startup( + state: &Arc, + run_id: RunId, + run_store: &fabro_store::RunDatabase, + run_state: &fabro_store::RunProjection, +) -> anyhow::Result<()> { + let key = PetriRuns::key(&run_id); + let held = match state.petri_runs.release_for_restart(run_id).await { + Ok(()) => true, + Err(StoreError::NotFound { .. }) => false, + Err(err) => { + return Err(anyhow::Error::new(err).context("releasing the Petri run's lease")); + } }; - if !outcome.complete { - message.push_str(" (record incomplete: "); - message.push_str(&outcome.incomplete.join("; ")); - message.push(')'); + let mode = if held { + RunExecutionMode::Resume + } else { + RunExecutionMode::Start + }; + info!( + run_id = %run_id, + petri_key = %key, + mode = super::worker_mode_arg(mode), + "Petri run left in flight by the previous server; relaunching its worker" + ); + for event in [ + workflow_event::Event::RunStartRequested { + resume: true, + actor: None, + }, + workflow_event::Event::RunRunnable { + source: RunRunnableSource::StartRequested, + actor: None, + }, + ] { + workflow_event::append_event(run_store, &run_id, &event).await?; } - message + let run_dir = Storage::new(state.server_storage_dir()) + .run_scratch(&run_id) + .root() + .to_path_buf(); + let mut runs = state.runs.lock().expect("runs lock poisoned"); + runs.insert( + run_id, + super::managed_run( + run_state.spec.graph_source.clone().unwrap_or_default(), + RunStatus::Runnable, + run_id.created_at(), + run_dir, + mode, + ), + ); + Ok(()) } /// The failed status, its message, and the `run.failed` event for it. diff --git a/lib/apps/fabro-server/tests/it/api/petri_store.rs b/lib/apps/fabro-server/tests/it/api/petri_store.rs index 9c8d6d157..7701dd156 100644 --- a/lib/apps/fabro-server/tests/it/api/petri_store.rs +++ b/lib/apps/fabro-server/tests/it/api/petri_store.rs @@ -305,3 +305,44 @@ async fn two_workers_cannot_both_hold_a_run_lease() { drop(taken); wait_until_released(server_store, &key).await; } + +/// A worker's store leases for the worker's launch id, whatever owner Petri +/// minted for the run runtime that opened the run: the server's lease row +/// names the launch, a reopen for the same launch shares the lease, and the +/// launch's release ends it. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_worker_store_leases_for_its_launch_not_for_petris_owner() { + let state = test_app_state(); + let base_url = serve(Arc::clone(&state), |router| router).await; + let run_id = RunId::new(); + let key = petri_key(run_id); + let token = state.test_issue_worker_token(&run_id); + let launch = OwnerId::new("launch-1"); + let store = HttpRunStore::for_worker( + worker_client(&base_url, &token, Duration::from_secs(5)).await, + launch.clone(), + ); + + let created = store + .open(&key, Access::Create { + owner: OwnerId::mint(), + }) + .await + .expect("the worker creates the run"); + let server_store = state.test_petri_run_store(); + assert_eq!( + server_store.owner(&key).await.expect("reads the lease"), + Some(launch.clone()), + "the lease names the launch" + ); + let reopened = store + .open(&key, Access::Write { + owner: OwnerId::mint(), + }) + .await + .expect("the same launch reopens the run"); + assert_eq!(reopened.locator(), created.locator()); + drop(reopened); + drop(created); + wait_until_released(server_store, &key).await; +} diff --git a/lib/apps/fabro-server/tests/it/scenario/petri.rs b/lib/apps/fabro-server/tests/it/scenario/petri.rs index 4eeade6a4..ae2dae66c 100644 --- a/lib/apps/fabro-server/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-server/tests/it/scenario/petri.rs @@ -3,6 +3,11 @@ //! `[server.execution] engine` says so, Petri's record of the run agrees //! with Fabro's status, and Petri's diagnostics refuse a run at create. //! +//! The runs here execute in the server process under the handler-registry +//! test override; outside it the scheduler launches a worker for a Petri +//! run, which the CLI's scenario tests cover with the real binary +//! (`lib/apps/fabro-cli/tests/it/scenario/petri.rs`). +//! //! The runs that execute take their host scope through the sandbox-driver //! host plugin, so those tests skip, and say why, when the executable is not //! found, unless `FABRO_REQUIRE_SANDBOX_PLUGINS` is set. The create-time @@ -222,9 +227,14 @@ async fn the_hello_bundle_runs_on_petri_when_the_version_names_the_engine() { .load(twin) .await; let settings = test_settings(); + // The handler-registry override is the test switch that keeps a Petri + // run in this process; without it the scheduler launches a worker. let state = TestAppStateBuilder::new() .runtime_settings(settings.server_settings, settings.manifest_run_defaults) .max_concurrent_runs(5) + .registry_factory(|interviewer| { + fabro_workflow::handler::default_registry(interviewer, || None) + }) .llm_overlay(llm_overlay_with_provider_base_url( "openai", twin.base_url.clone(),