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(),