diff --git a/lib/components/fabro-petri/Cargo.toml b/lib/components/fabro-petri/Cargo.toml index 039ac85b8..1933209e3 100644 --- a/lib/components/fabro-petri/Cargo.toml +++ b/lib/components/fabro-petri/Cargo.toml @@ -34,6 +34,7 @@ petri_frontend_fabro.workspace = true lithos-llm = { workspace = true, features = ["runtime"] } petri_testkit = { workspace = true, optional = true } anyhow.workspace = true +bytes.workspace = true async-trait.workspace = true serde.workspace = true serde_json.workspace = true diff --git a/lib/components/fabro-petri/README.md b/lib/components/fabro-petri/README.md index 5de85436a..9b8134b11 100644 --- a/lib/components/fabro-petri/README.md +++ b/lib/components/fabro-petri/README.md @@ -29,16 +29,20 @@ Every adapter the integration plan describes lands here. Petri's diagnostics come back in a shape the server maps onto Fabro's. - `admission`: the admitted graphs in Fabro's blob store, named on the run spec as `RunEngine::Petri(PetriAdmission)`, verified by digest on load. -- `engine`: a run executed by Petri in the server process over - `SqliteRunStore`, with the outcome read from the run's record through - `inspect_run`; `interviewer::Unattended` fails any question until the +- `engine`: a run executed by Petri, started from its admitted graphs or + resumed from its records, with the outcome read from the run's record + through `inspect_run` and mapped to the conclusion Fabro's read side + records. The run's worker process runs it over `HttpRunStore`; the server + runs it in its own process only under its test override, over + `SqliteRunStore`. `interviewer::Unattended` fails any question until the interview adapter lands. - `HttpRunStore`: the same store as a run's worker process reaches it, over the server's `/api/v1/runs/{id}/petri/*` endpoints with the worker's token. The server answers from its `SqliteRunStore`, so the lease and the `(log, seq)` rule are the store's; this layer carries requests, resends a - request whose reply was lost, and maps the server's error codes back to - `StoreError`. The module docs state the rules. + request whose reply was lost, maps the server's error codes back to + `StoreError`, and, for a worker, takes every lease for the worker's launch + id. The module docs state the rules. - `petri`: the Petri store vocabulary re-exported for the server, which answers the worker endpoints from a `SqliteRunStore` without naming a Petri package in its own manifest. @@ -49,7 +53,12 @@ A run goes to Petri when its workflow version's `workflow.toml` names `engine = "petri"` in `[workflow]`, or when the server's `[server.execution] engine` (`FABRO_SERVER_ENGINE`, `fabro server start --engine`) says so for versions that name none. The server side of both -halves is `fabro-server`'s `server::petri_runs`. +halves is `fabro-server`'s `server::petri_runs`; the worker side is +`fabro-cli`'s `commands::run::petri_worker`, which `fabro run __run-worker` +takes when the run's stored spec names Petri. After a server restart, a +Petri run left in flight goes back to a worker in `--mode resume`: the run +continues from its records, as Petri's own resume does, and full recovery +of the workspace to a durable snapshot is the plan's F3.5. ## How it is tested @@ -83,6 +92,15 @@ ulimit -n 4096 && cargo nextest run -p fabro-petri The server's end-to-end coverage is `lib/apps/fabro-server/tests/it/scenario/petri.rs`: the `hello` bundle on the OpenAI twin and a command-only bundle run to -completion through the create handler and the scheduler, under the version -flag and under the server setting, and Petri's diagnostics refuse a run at -create. +completion through the create handler and the scheduler, in the server +process under its test override, under the version flag and under the +server setting, and Petri's diagnostics refuse a run at create. The +server's `petri_runs` unit tests cover the lease ending at worker exit and +the restart reconcile that relaunches a worker in resume mode. + +The worker path is covered with the real binary in +`lib/apps/fabro-cli/tests/it/scenario/petri.rs`: 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`. diff --git a/lib/components/fabro-petri/src/admission.rs b/lib/components/fabro-petri/src/admission.rs index 8ba4397f7..1a7c78dfa 100644 --- a/lib/components/fabro-petri/src/admission.rs +++ b/lib/components/fabro-petri/src/admission.rs @@ -7,19 +7,38 @@ //! which is the key the coordinator registers the graph under and the name a //! nested-workflow step invokes its child by. Loading verifies the digest, //! so a blob that does not decode to the graph it claims is refused. +//! +//! The server loads through its [`BlobStore`]; a run's worker loads through +//! its client's blob read with [`load_with`], since the run's blobs are the +//! blob store the server answers `GET /runs/{id}/blobs/{hash}` from. + +use std::future::Future; use fabro_store::BlobStore; -use fabro_types::{PetriAdmission, PetriGraphRef}; +use fabro_types::{BlobHash, PetriAdmission, PetriGraphRef}; use petri_runtime::frontend::graph_digest; use petri_runtime::ir::Graph; use crate::check::Admitted; +/// The graphs a run starts from: the admitted root and its pre-lowered +/// children, loaded and verified. +pub struct AdmittedGraphs { + pub graph: Graph, + pub children: Vec, +} + /// Why an admission could not be stored or loaded. #[derive(Debug, thiserror::Error)] pub enum AdmissionError { #[error("the blob store failed")] Store(#[source] fabro_store::Error), + #[error("blob `{blob}` could not be read")] + Read { + blob: String, + #[source] + source: anyhow::Error, + }, #[error("graph `{digest}` is not in the blob store")] Missing { digest: String }, #[error("graph `{digest}` does not encode as JSON")] @@ -60,13 +79,30 @@ pub async fn persist( pub async fn load( blobs: &BlobStore, admission: &PetriAdmission, -) -> Result<(Graph, Vec), AdmissionError> { - let graph = load_graph(blobs, &admission.graph).await?; +) -> Result { + load_with( + |blob| async move { blobs.read(&blob).await.map_err(anyhow::Error::from) }, + admission, + ) + .await +} + +/// [`load`] over any blob read: `read` answers a hash with the blob's +/// bytes, or `None` when the store lacks it. +pub async fn load_with( + read: F, + admission: &PetriAdmission, +) -> Result +where + F: Fn(BlobHash) -> Fut, + Fut: Future>>, +{ + let graph = load_graph(&read, &admission.graph).await?; let mut children = Vec::with_capacity(admission.children.len()); for child in &admission.children { - children.push(load_graph(blobs, child).await?); + children.push(load_graph(&read, child).await?); } - Ok((graph, children)) + Ok(AdmittedGraphs { graph, children }) } async fn persist_graph(blobs: &BlobStore, graph: &Graph) -> Result { @@ -79,11 +115,17 @@ async fn persist_graph(blobs: &BlobStore, graph: &Graph) -> Result Result { - let bytes = blobs - .read(&graph.blob) +async fn load_graph(read: &F, graph: &PetriGraphRef) -> Result +where + F: Fn(BlobHash) -> Fut, + Fut: Future>>, +{ + let bytes = read(graph.blob) .await - .map_err(AdmissionError::Store)? + .map_err(|source| AdmissionError::Read { + blob: graph.blob.to_string(), + source, + })? .ok_or_else(|| AdmissionError::Missing { digest: graph.digest.clone(), })?; diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index 9b16f0880..0287ec699 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -1,11 +1,17 @@ -//! A Fabro run executed by Petri, in the server process. +//! A Fabro run executed by Petri: the one assembly the run's worker process +//! and the server share. //! -//! Until the worker's HTTP run store lands, a Petri run executes where the -//! server is: the runtime is assembled the same way the create handler -//! assembled it for `Runtime::check`, the run's records go to -//! [`SqliteRunStore`] under the Fabro run id as the run key, the admitted -//! graphs are loaded from the blob store, and `execution::host::run_configured` -//! runs the root invocation to its end. The outcome is then derived from +//! The worker is where a Petri run executes, as a legacy run does: it +//! reaches the run's record through [`HttpRunStore`](crate::HttpRunStore) +//! with its token, and everything else here is the same as in the server. +//! The server itself executes a run only under its test override, over +//! [`SqliteRunStore`](crate::SqliteRunStore) in its own process. Both build +//! the runtime the same way the create handler built it for +//! `Runtime::check`, name the Fabro run id as the run key, and hand the run +//! to `execution::host`: [`Execution::Start`] runs the admitted graphs +//! through `run_configured`; [`Execution::Resume`] continues the run from +//! its records through `resume_configured`, with the same observers a start +//! installs, as the host's docs require. The outcome is then derived from //! `inspect_run` over a read handle of the same store, so what the caller //! reports is what the durable record says. //! @@ -16,30 +22,46 @@ //! rides the caller's token: when it fires, the root invocation is cancelled //! politely and Petri records why. //! +//! A resume here is Petri's own: the run continues from its records, and +//! sandbox leases are reconciled by label. 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 and lands after this. +//! //! No stage or agent event is projected into Fabro's tables here; the //! caller appends only the run lifecycle events Fabro's read side needs to -//! finish the run. The projection over Petri's records is the read-side -//! item that follows. +//! finish the run, from the [`Conclusion`] this module derives. The +//! projection over Petri's records is the read-side item that follows. use std::path::PathBuf; use std::sync::Arc; -use fabro_store::BlobStore; -use fabro_types::{PetriAdmission, SandboxProviderKind}; +use fabro_types::{FailureReason, SandboxProviderKind}; use petri_execution::host::{self, HostError, HostRun}; use petri_execution::inspect::{self, InspectError, RunInspection}; -use petri_execution::{Access, CancelReason, InterviewDispatcher, RECEIPT_FILE, RunKey, RunStore}; +use petri_execution::{ + Access, CancelReason, InterviewDispatcher, InvocationId, RECEIPT_FILE, RunKey, RunStore, +}; use petri_runtime::executor::Retention; use petri_runtime::{RunOptions, SandboxBackend}; use tokio::fs; use tokio_util::sync::CancellationToken; use tracing::{debug, info, warn}; -use crate::admission::{self, AdmissionError}; +use crate::admission::AdmittedGraphs; use crate::interviewer::Unattended; -use crate::run_store::SqliteRunStore; use crate::runtime::RuntimeSpec; +/// How the run is entered: fresh, from the admitted graphs, or continued +/// from its records. +pub enum Execution { + /// Run the admitted graphs from the start; the run must not exist in + /// the store yet. + Start(AdmittedGraphs), + /// Continue the run from its records; the run must exist in the store + /// with its root invocation declared. + Resume, +} + /// One run to execute. pub struct RunRequest { /// The Fabro run id, which becomes Petri's run key: the run's identity @@ -47,12 +69,10 @@ pub struct RunRequest { pub run_id: String, /// Where the run's workspaces, step output and blobs live. pub run_dir: PathBuf, - /// What the create handler admitted. - pub admission: PetriAdmission, - /// The blob store the admitted graphs are read from. - pub blobs: Arc, - /// The run's durable record. - pub store: Arc, + pub execution: Execution, + /// The run's durable record: the worker's HTTP store, or the server's + /// SQLite store under the test override. + pub store: Arc, pub runtime: RuntimeSpec, /// The sandbox provider Fabro resolved for the run's environment. pub provider: SandboxProviderKind, @@ -86,32 +106,47 @@ pub struct RunOutcome { pub enum RunError { #[error("the run's sandbox provider `{provider}` is not one Petri serves")] UnsupportedProvider { provider: SandboxProviderKind }, - #[error("the admitted graphs could not be loaded")] - Admission(#[from] AdmissionError), #[error("the run's record could not be opened")] Open(#[source] petri_store::StoreError), + #[error("the run's record could not be read")] + Read(#[source] HostError), + #[error("the run's record has no root invocation, so there is nothing to resume")] + NothingToResume, #[error("the run's record could not be inspected")] Inspect(#[source] InspectError), #[error("the run ended without recording a status; the record says: {}", .0.join("; "))] Unfinished(Vec), } +/// How Fabro reports the run: what its read side records as the run's +/// terminal event. +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum Conclusion { + /// The record says the run succeeded and is whole. + Succeeded, + /// Anything else: the record says the run failed or was cancelled, the + /// record is incomplete, or the run could not be executed at all. + Failed { + reason: FailureReason, + message: String, + }, +} + /// Execute the run to its end and report what the record says. pub async fn run(request: RunRequest) -> Result { let backend = backend(&request.provider)?; - let (graph, children) = admission::load(&request.blobs, &request.admission).await?; let key = RunKey::new(request.run_id.as_str()); let mut options = RunOptions::new(&request.run_dir); options.run_key = Some(key.clone()); options.retention = Retention::Always; options.sandbox.backend = backend; - let store: Arc = request.store.clone(); - let runtime = request.runtime.runtime(true).store(store).options(options); + let runtime = request + .runtime + .runtime(true) + .store(Arc::clone(&request.store)) + .options(options); let dispatcher = InterviewDispatcher::new(Arc::new(Unattended)); - let host_run = HostRun::new(graph) - .with_children(children) - .observe(Arc::new(dispatcher.clone())); let cancel = request.cancel.clone(); let mut cancel_task = None; let with_handle = |handle: petri_execution::CoordinatorHandle, secrets| { @@ -122,8 +157,26 @@ pub async fn run(request: RunRequest) -> Result { handle.cancel_root_for(CancelReason::Control); })); }; - info!(run_id = %request.run_id, backend = %backend, "Starting Petri run"); - let result = Box::pin(host::run_configured(&runtime, host_run, with_handle)).await; + let result = match request.execution { + Execution::Start(graphs) => { + info!(run_id = %request.run_id, backend = %backend, "Starting Petri run"); + let host_run = HostRun::new(graphs.graph) + .with_children(graphs.children) + .observe(Arc::new(dispatcher.clone())); + Box::pin(host::run_configured(&runtime, host_run, with_handle)).await + } + Execution::Resume => { + check_resumable(request.store.as_ref(), &key).await?; + info!(run_id = %request.run_id, backend = %backend, "Resuming Petri run"); + Box::pin(host::resume_configured( + &runtime, + Vec::new(), + vec![Arc::new(dispatcher.clone())], + with_handle, + )) + .await + } + }; if let Some(task) = cancel_task { task.abort(); } @@ -133,18 +186,74 @@ pub async fn run(request: RunRequest) -> Result { Ok(report) => debug!(status = %report.status, "Petri run ended"), Err(error) => warn!(error = %error, "Petri run ended with a host error"), } - let inspection = inspect(&request.store, &key).await?; + let inspection = inspect(request.store.as_ref(), &key).await?; outcome(inspection, result.err()) } /// What the run's record says, read through a handle that holds no lease: /// the same derivation [`run`] ends with, for a caller that only holds the /// store, such as a test checking a finished run. -pub async fn outcome_of(store: &SqliteRunStore, run_id: &str) -> Result { +pub async fn outcome_of(store: &dyn RunStore, run_id: &str) -> Result { let inspection = inspect(store, &RunKey::new(run_id)).await?; outcome(inspection, None) } +/// How Fabro reports what [`run`] returned. A cancelled run is a failure +/// with the cancelled reason, as the legacy executor reports one; every +/// other shortfall is a workflow error whose message says what the record, +/// or the host, said. +#[must_use] +pub fn conclusion(result: &Result) -> Conclusion { + match result { + Ok(RunOutcome { + status: RunStatus::Success, + complete: true, + .. + }) => Conclusion::Succeeded, + Ok(outcome) => { + let reason = match outcome.status { + RunStatus::Cancelled => FailureReason::Cancelled, + RunStatus::Success | RunStatus::Failed => FailureReason::WorkflowError, + }; + Conclusion::Failed { + reason, + message: failure_message(outcome), + } + } + Err(error) => Conclusion::Failed { + reason: FailureReason::WorkflowError, + message: error_chain(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) { + (RunStatus::Cancelled, _) => "the run was cancelled".to_string(), + (_, Some(failure)) => failure.clone(), + (RunStatus::Failed, None) => "the run failed".to_string(), + (RunStatus::Success, None) => "the run's record is incomplete".to_string(), + }; + if !outcome.complete { + message.push_str(" (record incomplete: "); + message.push_str(&outcome.incomplete.join("; ")); + message.push(')'); + } + message +} + +/// The error and every cause under it, as one line. +fn error_chain(error: &RunError) -> String { + let mut parts = vec![error.to_string()]; + let mut cause = std::error::Error::source(error); + while let Some(next) = cause { + parts.push(next.to_string()); + cause = next.source(); + } + parts.join(": ") +} + /// The sandbox backend for Fabro's provider kind. fn backend(provider: &SandboxProviderKind) -> Result { if *provider == SandboxProviderKind::LOCAL { @@ -160,8 +269,25 @@ fn backend(provider: &SandboxProviderKind) -> Result { } } +/// Refuse a resume the host would not survive: `resume_configured` indexes +/// the root invocation of the stored state, so a record with none (the run +/// was created in the store and nothing more) is refused here with a named +/// error instead. +async fn check_resumable(store: &dyn RunStore, key: &RunKey) -> Result<(), RunError> { + let logs = store + .open(key, Access::Read) + .await + .map_err(RunError::Open)?; + let state = host::stored_state(&*logs).await.map_err(RunError::Read)?; + if state.invocations.contains_key(&InvocationId::ROOT) { + Ok(()) + } else { + Err(RunError::NothingToResume) + } +} + /// Read the run back through a handle that holds no lease. -async fn inspect(store: &SqliteRunStore, key: &RunKey) -> Result { +async fn inspect(store: &dyn RunStore, key: &RunKey) -> Result { let logs = store .open(key, Access::Read) .await @@ -224,3 +350,66 @@ async fn write_receipt(run_dir: &std::path::Path, receipt: &petri_execution::Int warn!(path = %path.display(), error = %error, "could not write the interview receipt"); } } + +#[cfg(test)] +mod tests { + use super::*; + + fn outcome_with(status: RunStatus, failure: Option<&str>, complete: bool) -> RunOutcome { + RunOutcome { + status, + failure: failure.map(ToOwned::to_owned), + complete, + incomplete: if complete { + Vec::new() + } else { + vec!["execution 0 did not finish".to_string()] + }, + } + } + + #[test] + fn a_whole_successful_record_concludes_succeeded() { + assert_eq!( + conclusion(&Ok(outcome_with(RunStatus::Success, None, true))), + Conclusion::Succeeded + ); + } + + #[test] + fn a_cancelled_record_concludes_cancelled() { + assert_eq!( + conclusion(&Ok(outcome_with(RunStatus::Cancelled, None, true))), + Conclusion::Failed { + reason: FailureReason::Cancelled, + message: "the run was cancelled".to_string(), + } + ); + } + + #[test] + fn a_failed_record_carries_the_root_failure_and_the_incomplete_reasons() { + assert_eq!( + conclusion(&Ok(outcome_with( + RunStatus::Failed, + Some("step `say` failed"), + false + ))), + Conclusion::Failed { + reason: FailureReason::WorkflowError, + message: "step `say` failed (record incomplete: execution 0 did not finish)" + .to_string(), + } + ); + } + + #[test] + fn a_host_error_concludes_with_its_chain() { + let error = RunError::Unfinished(vec!["no status".to_string()]); + assert_eq!(conclusion(&Err(error)), Conclusion::Failed { + reason: FailureReason::WorkflowError, + message: "the run ended without recording a status; the record says: no status" + .to_string(), + }); + } +} diff --git a/lib/components/fabro-petri/src/http_store.rs b/lib/components/fabro-petri/src/http_store.rs index d213706ed..76ee9b69c 100644 --- a/lib/components/fabro-petri/src/http_store.rs +++ b/lib/components/fabro-petri/src/http_store.rs @@ -26,6 +26,16 @@ //! runtime at drop, the server's worker-exit release is the backstop, and //! the drop says so in the log. //! +//! # The owner +//! +//! A store built with [`HttpRunStore::for_worker`] names one owner for the +//! whole process: every `Create` and `Write` takes the lease for the +//! worker's launch id, whatever owner Petri minted for the run runtime that +//! asked. One worker process executes one run, so the lease is the +//! launch's, the worker logs it once at start, and the server's lease row +//! names the launch that holds it. A store built with [`HttpRunStore::new`] +//! passes Petri's owner through unchanged. +//! //! # Lost replies //! //! Every call is one request. A reply that never arrives (a transport error, @@ -92,6 +102,9 @@ impl fmt::Debug for HttpRunStore { /// What the store and every handle it opens share. struct Shared { client: Client, + /// The owner every writer open takes the lease for, when the store is + /// a worker's; `None` passes Petri's owner through. + owner: Option, /// The writer handle alive in this process per run and owner, so a /// same-owner reopen shares it and the lease lasts while any handle /// does. @@ -101,12 +114,25 @@ struct Shared { } impl HttpRunStore { - /// A store over a client that carries the worker's token. + /// A store over a client that carries the worker's token, taking each + /// lease for the owner Petri names. #[must_use] pub fn new(client: Client) -> Self { + Self::build(client, None) + } + + /// A worker's store: every lease is taken for `owner`, the worker's + /// launch id, whatever owner Petri names. + #[must_use] + pub fn for_worker(client: Client, owner: OwnerId) -> Self { + Self::build(client, Some(owner)) + } + + fn build(client: Client, owner: Option) -> Self { Self { shared: Arc::new(Shared { client, + owner, live: Mutex::default(), releases: Mutex::default(), }), @@ -290,6 +316,9 @@ impl RunStore for HttpRunStore { Access::Write { owner } => (PetriAccess::Write, Some(owner)), Access::Read => (PetriAccess::Read, None), }; + // A worker's store leases for its launch, not for the owner Petri + // minted for this run runtime. + let owner = owner.map(|named| shared.owner.as_ref().unwrap_or(named)); let request = PetriOpenRequest { access: api_access, owner: owner.map(|owner| owner.as_str().to_string()), @@ -326,7 +355,12 @@ impl RunStore for HttpRunStore { }; let opened = opened.map_err(|error| shared.store_error(key, "open the run", None, error))?; - debug!(run_id = %key, access = ?request.access, "Petri run opened over the API"); + debug!( + run_id = %key, + access = ?request.access, + owner = owner.map(OwnerId::as_str), + "Petri run opened over the API" + ); match owner { Some(owner) => Ok(self.writer(key, run_id, owner.clone(), opened.locator)), None => Ok(Arc::new(HttpRunLogs { diff --git a/lib/components/fabro-petri/src/lib.rs b/lib/components/fabro-petri/src/lib.rs index 8e9c6c78b..443032ac8 100644 --- a/lib/components/fabro-petri/src/lib.rs +++ b/lib/components/fabro-petri/src/lib.rs @@ -16,11 +16,13 @@ //! its diagnostics come back in a shape Fabro maps onto its own; //! - [`admission`]: the admitted graphs in Fabro's blob store, named on the run //! spec; -//! - [`engine`]: a run executed by Petri in the server process, with the -//! outcome read from its record; +//! - [`engine`]: a run executed by Petri, started or resumed, in the run's +//! worker process over the HTTP store (or in the server process under its +//! test override), with the outcome read from its record; //! - [`interviewer`]: the interviewer of a run nobody is watching; //! - [`HttpRunStore`]: the same store as a run's worker process reaches it, -//! over the server's API with the worker's token; +//! over the server's API with the worker's token and its launch id as the +//! lease owner; //! - the platform adapters still to come: hooks, interviews over Fabro's API, //! secrets, output storage, the run tools, the event projection. //! diff --git a/lib/components/fabro-petri/tests/check.rs b/lib/components/fabro-petri/tests/check.rs index 351166810..e936f716a 100644 --- a/lib/components/fabro-petri/tests/check.rs +++ b/lib/components/fabro-petri/tests/check.rs @@ -118,11 +118,11 @@ async fn the_hello_bundle_is_admitted_and_round_trips_through_the_blob_store() { .await .expect("the graphs persist"); assert!(record.children.is_empty()); - let (graph, children) = admission::load(&blobs, &record) + let graphs = admission::load(&blobs, &record) .await .expect("the graphs load"); - assert_eq!(graph, admitted.graph); - assert!(children.is_empty()); + assert_eq!(graphs.graph, admitted.graph); + assert!(graphs.children.is_empty()); } #[tokio::test]