From 03309d4240a252f4e6ab0bb499e3abb133db7698 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 17 Sep 2026 20:27:44 -0400 Subject: [PATCH] Let Petri compile and execute a Fabro run through fabro-petri `check` materializes a workflow version's bundle into a temporary directory (`Runtime::check` reads files from disk), lowers it with the run's inputs and launch, and returns the admitted graphs or Petri's diagnostics in a shape the server maps onto Fabro's. `admission` keeps the admitted graphs in the blob store, named on the run spec and verified by digest on load. `runtime` assembles the same Petri runtime at create and at execution: the Fabro frontend with the server's settings layer, the Attractor step kinds, the model client as the PebbleClient capability so admission pins every model. `engine` runs the admitted graph in the server process over SqliteRunStore under the Fabro run id, with the standalone defaults, an interviewer that fails any question, and cancel on a token, and derives the outcome from inspect_run over the run's record. Co-Authored-By: Claude Fable 5.1 --- Cargo.lock | 8 + lib/components/fabro-petri/Cargo.toml | 9 +- lib/components/fabro-petri/README.md | 34 ++- lib/components/fabro-petri/src/admission.rs | 104 ++++++++ lib/components/fabro-petri/src/check.rs | 245 ++++++++++++++++++ lib/components/fabro-petri/src/engine.rs | 226 ++++++++++++++++ lib/components/fabro-petri/src/interviewer.rs | 24 ++ lib/components/fabro-petri/src/lib.rs | 18 +- lib/components/fabro-petri/src/runtime.rs | 99 +++++++ lib/components/fabro-petri/tests/check.rs | 224 ++++++++++++++++ 10 files changed, 986 insertions(+), 5 deletions(-) create mode 100644 lib/components/fabro-petri/src/admission.rs create mode 100644 lib/components/fabro-petri/src/check.rs create mode 100644 lib/components/fabro-petri/src/engine.rs create mode 100644 lib/components/fabro-petri/src/interviewer.rs create mode 100644 lib/components/fabro-petri/src/runtime.rs create mode 100644 lib/components/fabro-petri/tests/check.rs diff --git a/Cargo.lock b/Cargo.lock index 70342d49b..6be4ffd77 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2887,9 +2887,13 @@ name = "fabro-petri" version = "0.357.0-nightly.0" dependencies = [ "async-trait", + "fabro-auth", "fabro-db", + "fabro-http", + "fabro-llm", "fabro-store", "fabro-types", + "lithos-llm", "petri-attractor-steps", "petri-execution", "petri-frontend-attractor", @@ -2897,10 +2901,13 @@ dependencies = [ "petri-runtime", "petri-store", "petri-testkit", + "serde", "serde_json", "sqlx", "tempfile", + "thiserror 2.0.18", "tokio", + "tokio-util", "tracing", ] @@ -3005,6 +3012,7 @@ dependencies = [ "fabro-macros", "fabro-manifest", "fabro-mcp-store", + "fabro-petri", "fabro-proc", "fabro-redact", "fabro-sandbox", diff --git a/lib/components/fabro-petri/Cargo.toml b/lib/components/fabro-petri/Cargo.toml index 9398c5481..dda58aacd 100644 --- a/lib/components/fabro-petri/Cargo.toml +++ b/lib/components/fabro-petri/Cargo.toml @@ -14,6 +14,7 @@ workspace = true [dependencies] fabro-db = { path = "../../foundation/fabro-db" } +fabro-http.workspace = true fabro-store = { path = "../fabro-store" } fabro-types = { path = "../../foundation/fabro-types" } petri_runtime.workspace = true @@ -22,14 +23,20 @@ petri_store.workspace = true petri_attractor_steps.workspace = true petri_frontend_attractor.workspace = true petri_frontend_fabro.workspace = true +lithos-llm = { workspace = true, features = ["runtime"] } async-trait.workspace = true +serde.workspace = true serde_json.workspace = true sqlx.workspace = true +tempfile = "3" +thiserror.workspace = true tokio.workspace = true +tokio-util.workspace = true tracing.workspace = true [dev-dependencies] +fabro-auth = { path = "../../foundation/fabro-auth", features = ["test-support"] } +fabro-llm = { path = "../fabro-llm", features = ["test-support"] } fabro-store = { path = "../fabro-store", features = ["test-support"] } petri_testkit.workspace = true -tempfile = "3" tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } diff --git a/lib/components/fabro-petri/README.md b/lib/components/fabro-petri/README.md index 727afb430..908af3662 100644 --- a/lib/components/fabro-petri/README.md +++ b/lib/components/fabro-petri/README.md @@ -19,8 +19,28 @@ Every adapter the integration plan describes lands here. run and its writer lease, `petri_records` for every record of every log, and the shared `blobs` table). The module docs state the lease and append rules. -- The platform adapters the plan adds after it: hooks, interviews, secrets, - output storage, run tools, the event projection. +- `runtime`: the Petri runtime Fabro assembles, the same way at create time + and at execution: the Fabro frontend with the server's settings layer, the + Attractor step kinds (real, or simulated for a dry run), the model client + as the `PebbleClient` capability, the Fabro home. +- `check`: Petri compiles at create time. The workflow version's bundle is + materialized into a temporary directory (`Runtime::check` reads files from + disk), lowered with the run's inputs and launch, and the admitted graphs or + 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 + interview adapter lands. +- The platform adapters the plan adds after it: hooks, interviews over + Fabro's API, secrets, output storage, run tools, the event projection. + +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`. ## How it is tested @@ -32,6 +52,10 @@ Integration tests live under `tests/`: Both skip, and say why, when the `sandbox-driver-host` plugin executable is not on `PATH` (every run takes its scope's environment through it); the sandbox-plugins CI job requires them. +- `check.rs` admits the `hello` bundle and round-trips its graph through + the blob store, binds the launch, and refuses an unknown attribute and, + with a model client over the test catalog, an unknown model + (`attractor.model.unknown`). No plugin is needed. - `sqlite_store.rs` runs Petri's store conformance suite (`petri_testkit::run_store::conformance`) against `SqliteRunStore`, plus the operator release, lease exclusivity, a crash between appends, and blob @@ -42,3 +66,9 @@ Run them with: ```sh 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. diff --git a/lib/components/fabro-petri/src/admission.rs b/lib/components/fabro-petri/src/admission.rs new file mode 100644 index 000000000..8ba4397f7 --- /dev/null +++ b/lib/components/fabro-petri/src/admission.rs @@ -0,0 +1,104 @@ +//! The admitted graphs in Fabro's blob store. +//! +//! What `Runtime::check` admitted is what the run executes and resumes from, +//! so the root graph and every pre-lowered child are serialized into the +//! blob store at create time and named on the run spec as a +//! [`PetriAdmission`]: the blob by hash, and Petri's own content digest, +//! 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. + +use fabro_store::BlobStore; +use fabro_types::{PetriAdmission, PetriGraphRef}; +use petri_runtime::frontend::graph_digest; +use petri_runtime::ir::Graph; + +use crate::check::Admitted; + +/// 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("graph `{digest}` is not in the blob store")] + Missing { digest: String }, + #[error("graph `{digest}` does not encode as JSON")] + Encode { + digest: String, + #[source] + source: serde_json::Error, + }, + #[error("blob `{digest}` does not decode as a graph")] + Decode { + digest: String, + #[source] + source: serde_json::Error, + }, + #[error("blob `{blob}` decodes to graph `{found}`, not `{digest}`")] + DigestMismatch { + blob: String, + digest: String, + found: String, + }, +} + +/// Serialize the admitted graphs into `blobs` and name them. +pub async fn persist( + blobs: &BlobStore, + admitted: &Admitted, +) -> Result { + let graph = persist_graph(blobs, &admitted.graph).await?; + let mut children = Vec::with_capacity(admitted.children.len()); + for child in &admitted.children { + children.push(persist_graph(blobs, child).await?); + } + Ok(PetriAdmission { graph, children }) +} + +/// The root graph and its children, read back from `blobs` and checked +/// against their digests. +pub async fn load( + blobs: &BlobStore, + admission: &PetriAdmission, +) -> Result<(Graph, Vec), AdmissionError> { + let graph = load_graph(blobs, &admission.graph).await?; + let mut children = Vec::with_capacity(admission.children.len()); + for child in &admission.children { + children.push(load_graph(blobs, child).await?); + } + Ok((graph, children)) +} + +async fn persist_graph(blobs: &BlobStore, graph: &Graph) -> Result { + let digest = graph_digest(graph); + let bytes = serde_json::to_vec(graph).map_err(|source| AdmissionError::Encode { + digest: digest.clone(), + source, + })?; + let blob = blobs.write(&bytes).await.map_err(AdmissionError::Store)?; + Ok(PetriGraphRef { blob, digest }) +} + +async fn load_graph(blobs: &BlobStore, graph: &PetriGraphRef) -> Result { + let bytes = blobs + .read(&graph.blob) + .await + .map_err(AdmissionError::Store)? + .ok_or_else(|| AdmissionError::Missing { + digest: graph.digest.clone(), + })?; + let decoded: Graph = + serde_json::from_slice(&bytes).map_err(|source| AdmissionError::Decode { + digest: graph.digest.clone(), + source, + })?; + let found = graph_digest(&decoded); + if found != graph.digest { + return Err(AdmissionError::DigestMismatch { + blob: graph.blob.to_string(), + digest: graph.digest.clone(), + found, + }); + } + Ok(decoded) +} diff --git a/lib/components/fabro-petri/src/check.rs b/lib/components/fabro-petri/src/check.rs new file mode 100644 index 000000000..5b4022f09 --- /dev/null +++ b/lib/components/fabro-petri/src/check.rs @@ -0,0 +1,245 @@ +//! Petri compiles: the create handler hands a workflow version's files, the +//! run's inputs and the launch to `Runtime::check`, and gets back either the +//! admitted graphs or Petri's diagnostics. +//! +//! `Runtime::check` reads the workflow and its settings files from disk, so +//! the bundle is materialized into a temporary directory first, laid out the +//! way the Fabro frontend expects: the workflow file with `workflow.toml` +//! beside it under a bundle root that holds a `.fabro` directory (with +//! `.fabro/project.toml` when the caller has one). The directory is removed +//! when the check returns. An in-memory `FileSource` entry point on +//! `Runtime` would remove the round trip; that is a Petri follow-up. +//! +//! The launch binds the compile variables the Fabro frontend reads: +//! `petri.launch_model` and `petri.launch_provider` as the model default +//! below every file layer, and `petri.repository` as the repository the root +//! `start` stage checks out. A caller with no local repository binds `null`, +//! and the run starts from an empty workspace. + +use std::collections::BTreeMap; +use std::io; +use std::path::{Path, PathBuf}; + +use petri_runtime::LoadError; +use petri_runtime::frontend::{ + self, CompileInputs, LAUNCH_MODEL_VAR, LAUNCH_PROVIDER_VAR, REPOSITORY_VAR, Severity, +}; +use petri_runtime::ir::Graph; +use serde::{Deserialize, Serialize}; +use serde_json::Value; + +use crate::runtime::RuntimeSpec; + +/// The directory under the temporary bundle root the version's files land +/// in. Its parent holds `.fabro`, so the Fabro frontend takes the parent as +/// the bundle root. +const BUNDLE_DIR: &str = "bundle"; + +/// The project settings file the Fabro frontend reads at the bundle root. +const PROJECT_FILE: &str = ".fabro/project.toml"; + +/// One workflow bundle to check: its files by bundle-relative path. +#[derive(Clone, Debug, Default)] +pub struct Bundle { + /// Every file of the version closure, keyed by its path relative to + /// the bundle (`workflow.fabro`, `workflow.toml`, `prompts/goal.md`, + /// `children/check.fabro`), with `/` separators. + pub files: BTreeMap, + /// The workflow file to check, one of `files`. + pub entrypoint: String, + /// `.fabro/project.toml` at the bundle root, when the caller has one. + pub project_toml: Option, +} + +/// What the launch binds below the file layers. +#[derive(Clone, Debug, Default)] +pub struct Launch { + pub model: Option, + pub provider: Option, + /// The local repository the root `start` stage checks out into the + /// workspace; `None` starts the run from an empty workspace. + pub repository: Option, +} + +/// One check: the bundle, the run's inputs, the launch and the runtime. +#[derive(Clone, Default)] +pub struct CheckRequest { + pub bundle: Bundle, + /// The intent's inputs, under which `[run.inputs]` defaults fill in. + pub inputs: BTreeMap, + pub launch: Launch, + pub runtime: RuntimeSpec, +} + +/// Petri's diagnostic, in the shape Fabro's create handler maps onto its +/// own: the stable code, the text, the hint, and the position in the +/// bundle when known. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct Diagnostic { + pub severity: DiagnosticSeverity, + /// Petri's stable code: `attractor.model.unknown`, `fabro.hooks.toml`, + /// `unsupported.workflow_toml.key`. + pub code: String, + pub message: String, + pub hint: Option, + /// The bundle-relative file the diagnostic names. + pub file: String, + /// 1-based; `None` for a whole-file problem. + pub line: Option, + pub column: Option, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum DiagnosticSeverity { + Error, + Warning, +} + +/// What Petri admitted: the lowered root graph, the pre-lowered child +/// graphs, and the warnings the lowering raised. +pub struct Admitted { + pub graph: Graph, + pub children: Vec, + pub warnings: Vec, +} + +/// Why a check produced no graph. +#[derive(Debug, thiserror::Error)] +pub enum CheckError { + /// Petri refused the workflow. Every diagnostic is here, warnings + /// included; at least one is an error. + #[error("Petri refused the workflow with {} diagnostic(s)", .0.len())] + Rejected(Vec), + /// The bundle could not be materialized for the check. + #[error("could not materialize the workflow bundle at `{path}`")] + Materialize { + path: PathBuf, + #[source] + source: io::Error, + }, + /// The bundle's entrypoint is not one of its files, or no frontend + /// claims it. + #[error("the workflow could not be loaded")] + Load(#[source] LoadError), +} + +/// Materialize the bundle, run `Runtime::check`, and hand back the admitted +/// graphs or the diagnostics. Blocking: it reads and writes files and +/// lowers the graph, so a server calls it from its blocking pool. +pub fn check(request: &CheckRequest) -> Result { + let root = tempfile::tempdir().map_err(|source| CheckError::Materialize { + path: std::env::temp_dir(), + source, + })?; + let workflow = materialize(root.path(), &request.bundle)?; + let runtime = request.runtime.runtime(false); + let inputs = compile_inputs(&request.inputs, &request.launch); + let lowered = runtime + .check(&workflow, None, None, &inputs) + .map_err(CheckError::Load)?; + let diagnostics: Vec = lowered + .diagnostics + .iter() + .map(|diagnostic| convert(diagnostic, root.path())) + .collect(); + match lowered.graph { + Some(graph) => Ok(Admitted { + graph, + children: lowered.children, + warnings: diagnostics, + }), + None => Err(CheckError::Rejected(diagnostics)), + } +} + +/// Write the bundle under `root/bundle/`, with `root/.fabro` beside it so +/// the frontend takes `root` as the bundle root. Returns the entrypoint's +/// path. +#[expect( + clippy::disallowed_methods, + reason = "the check is a blocking function; its caller runs it on the blocking pool" +)] +fn materialize(root: &Path, bundle: &Bundle) -> Result { + let write = |relative: &str, text: &str| -> Result<(), CheckError> { + let path = root.join(relative); + let materialize = |source| CheckError::Materialize { + path: path.clone(), + source, + }; + if let Some(parent) = path.parent() { + std::fs::create_dir_all(parent).map_err(materialize)?; + } + std::fs::write(&path, text).map_err(materialize) + }; + let fabro_dir = root.join(".fabro"); + std::fs::create_dir_all(&fabro_dir).map_err(|source| CheckError::Materialize { + path: fabro_dir, + source, + })?; + if let Some(project) = &bundle.project_toml { + write(PROJECT_FILE, project)?; + } + for (relative, text) in &bundle.files { + write(&format!("{BUNDLE_DIR}/{relative}"), text)?; + } + Ok(root.join(BUNDLE_DIR).join(&bundle.entrypoint)) +} + +/// The compile inputs: the intent's inputs, and the launch variables. +fn compile_inputs(inputs: &BTreeMap, launch: &Launch) -> CompileInputs { + let mut compile = CompileInputs::new(); + for (name, value) in inputs { + compile.inputs.insert(name.as_str().into(), value.clone()); + } + let text = |value: &Option| match value { + Some(text) if !text.trim().is_empty() => Value::String(text.clone()), + _ => Value::Null, + }; + compile + .vars + .insert(LAUNCH_MODEL_VAR.into(), text(&launch.model)); + compile + .vars + .insert(LAUNCH_PROVIDER_VAR.into(), text(&launch.provider)); + // Bound even when absent: `Runtime::lower` would otherwise bind the + // temporary bundle root, which is gone by the time the run starts. + let repository = launch.repository.as_ref().map_or(Value::Null, |path| { + Value::String(path.to_string_lossy().into_owned()) + }); + compile.vars.insert(REPOSITORY_VAR.into(), repository); + compile +} + +/// Petri's diagnostic in Fabro's shape, with the file made relative to the +/// bundle. +fn convert(diagnostic: &frontend::Diagnostic, root: &Path) -> Diagnostic { + let file = diagnostic.span.file.as_str(); + let prefix = format!("{BUNDLE_DIR}/"); + let file = Path::new(file) + .strip_prefix(root) + .map_or(file, |relative| relative.to_str().unwrap_or(file)) + .to_string(); + let file = file + .strip_prefix(&prefix) + .map_or(file.as_str(), |relative| relative) + .to_string(); + Diagnostic { + severity: match diagnostic.severity { + Severity::Error => DiagnosticSeverity::Error, + Severity::Warning => DiagnosticSeverity::Warning, + }, + code: diagnostic.code.to_string(), + message: diagnostic.message.clone(), + hint: diagnostic.hint.clone(), + file, + line: (diagnostic.span.line > 0).then_some(diagnostic.span.line), + column: (diagnostic.span.column > 0).then_some(diagnostic.span.column), + } +} + +impl Diagnostic { + pub fn is_error(&self) -> bool { + self.severity == DiagnosticSeverity::Error + } +} diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs new file mode 100644 index 000000000..9b16f0880 --- /dev/null +++ b/lib/components/fabro-petri/src/engine.rs @@ -0,0 +1,226 @@ +//! A Fabro run executed by Petri, in the server process. +//! +//! 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 +//! `inspect_run` over a read handle of the same store, so what the caller +//! reports is what the durable record says. +//! +//! What the standalone runner's defaults give the run: Petri's local hook +//! service for `[[run.hooks]]`, no `ExecutionHooks` of Fabro's own, the +//! [`Unattended`] interviewer that fails any question, no host tools, and +//! `Retention::Always` for every workspace, Fabro's default. Cancellation +//! rides the caller's token: when it fires, the root invocation is cancelled +//! politely and Petri records why. +//! +//! 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. + +use std::path::PathBuf; +use std::sync::Arc; + +use fabro_store::BlobStore; +use fabro_types::{PetriAdmission, 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_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::interviewer::Unattended; +use crate::run_store::SqliteRunStore; +use crate::runtime::RuntimeSpec; + +/// One run to execute. +pub struct RunRequest { + /// The Fabro run id, which becomes Petri's run key: the run's identity + /// in the store and the label on every sandbox of the run. + 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 runtime: RuntimeSpec, + /// The sandbox provider Fabro resolved for the run's environment. + pub provider: SandboxProviderKind, + /// Fires to cancel the run. + pub cancel: CancellationToken, +} + +/// The recorded status of a finished run. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RunStatus { + Success, + Failed, + Cancelled, +} + +/// What the durable record says about the run once it ended. +#[derive(Clone, Debug)] +pub struct RunOutcome { + pub status: RunStatus, + /// The root invocation's failure message, when it failed. + pub failure: Option, + /// Whether the record is whole: the run recorded its finish and every + /// log replays byte for byte. + pub complete: bool, + /// Every reason `complete` is false. + pub incomplete: Vec, +} + +/// Why the run could not be executed or its outcome read. +#[derive(Debug, thiserror::Error)] +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 inspected")] + Inspect(#[source] InspectError), + #[error("the run ended without recording a status; the record says: {}", .0.join("; "))] + Unfinished(Vec), +} + +/// 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 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| { + dispatcher.wire(handle.clone(), secrets); + cancel_task = Some(tokio::spawn(async move { + cancel.cancelled().await; + info!("cancelling the Petri run"); + 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; + if let Some(task) = cancel_task { + task.abort(); + } + let receipt = dispatcher.shutdown().await; + write_receipt(&request.run_dir, &receipt).await; + match &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?; + 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 { + let inspection = inspect(store, &RunKey::new(run_id)).await?; + outcome(inspection, None) +} + +/// The sandbox backend for Fabro's provider kind. +fn backend(provider: &SandboxProviderKind) -> Result { + if *provider == SandboxProviderKind::LOCAL { + Ok(SandboxBackend::Host) + } else if *provider == SandboxProviderKind::DOCKER { + Ok(SandboxBackend::Docker) + } else if *provider == SandboxProviderKind::DAYTONA { + Ok(SandboxBackend::Daytona) + } else { + Err(RunError::UnsupportedProvider { + provider: provider.clone(), + }) + } +} + +/// Read the run back through a handle that holds no lease. +async fn inspect(store: &SqliteRunStore, key: &RunKey) -> Result { + let logs = store + .open(key, Access::Read) + .await + .map_err(RunError::Open)?; + inspect::inspect_run(&*logs) + .await + .map_err(RunError::Inspect) +} + +/// The outcome the record supports. A run whose record has no status is +/// unfinished: the host error, when there is one, says why. +fn outcome( + inspection: RunInspection, + host_error: Option, +) -> Result { + let status = match inspection.status.as_deref() { + Some("success") => RunStatus::Success, + Some("failed") => RunStatus::Failed, + Some("cancelled") => RunStatus::Cancelled, + _ => { + let mut reasons = inspection.incomplete.clone(); + if let Some(error) = host_error { + reasons.push(error.to_string()); + } + return Err(RunError::Unfinished(reasons)); + } + }; + let failure = inspection + .invocations + .iter() + .find(|invocation| invocation.invocation == inspection.root.invocation) + .and_then(|root| root.result.as_ref()) + .and_then(|result| result.failure.as_ref()) + .map(|failure| failure.message.clone()); + Ok(RunOutcome { + status, + failure, + complete: inspection.complete, + incomplete: inspection.incomplete, + }) +} + +/// The interview receipt beside the run, as the standalone runner writes +/// it. A receipt that cannot be written is logged: the run's record does +/// not depend on it. +async fn write_receipt(run_dir: &std::path::Path, receipt: &petri_execution::InterviewReceipt) { + let path = run_dir.join(RECEIPT_FILE); + let bytes = match serde_json::to_vec_pretty(receipt) { + Ok(bytes) => bytes, + Err(error) => { + warn!(error = %error, "could not encode the interview receipt"); + return; + } + }; + if let Err(error) = fs::create_dir_all(run_dir).await { + warn!(path = %run_dir.display(), error = %error, "could not create the run directory"); + return; + } + if let Err(error) = fs::write(&path, bytes).await { + warn!(path = %path.display(), error = %error, "could not write the interview receipt"); + } +} diff --git a/lib/components/fabro-petri/src/interviewer.rs b/lib/components/fabro-petri/src/interviewer.rs new file mode 100644 index 000000000..594843ba0 --- /dev/null +++ b/lib/components/fabro-petri/src/interviewer.rs @@ -0,0 +1,24 @@ +//! The interviewer of a run nobody is watching. +//! +//! Until the questions adapter over Fabro's API lands (F3.2), a Petri run in +//! the server has no way to reach a person. A human gate that asks anyway +//! gets a failure that says so, the gate fails closed, and the reason +//! reaches the interview receipt, instead of a question that waits forever. + +use petri_execution::{InterviewError, InterviewReply, InterviewRequest, Interviewer}; +use tokio_util::sync::CancellationToken; + +/// Fails every question with a clear error. +#[derive(Clone, Copy, Debug, Default)] +pub struct Unattended; + +#[async_trait::async_trait] +impl Interviewer for Unattended { + async fn reply(&self, request: InterviewRequest, _cancel: CancellationToken) -> InterviewReply { + InterviewReply::Failed(InterviewError::new(format!( + "node `{}` asked a question, but a Petri run has no interviewer yet: questions reach \ + nobody until the interview adapter lands", + request.node + ))) + } +} diff --git a/lib/components/fabro-petri/src/lib.rs b/lib/components/fabro-petri/src/lib.rs index dd3f77139..0d06cc428 100644 --- a/lib/components/fabro-petri/src/lib.rs +++ b/lib/components/fabro-petri/src/lib.rs @@ -10,12 +10,26 @@ //! //! - [`SqliteRunStore`]: Petri's run store over Fabro's SQLite database, so a //! run's records are its source of truth in Fabro's tables; -//! - the platform adapters: hooks, interviews, secrets, output storage, the run -//! tools, the event projection. +//! - [`runtime`]: the Petri runtime Fabro assembles, at create time and at +//! execution; +//! - [`check`]: Petri compiles a workflow version's bundle at create time, and +//! 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; +//! - [`interviewer`]: the interviewer of a run nobody is watching; +//! - the platform adapters still to come: hooks, interviews over Fabro's API, +//! secrets, output storage, the run tools, the event projection. //! //! The Petri packages are pinned by revision in the workspace `Cargo.toml` //! under `petri_*` keys. +pub mod admission; +pub mod check; +pub mod engine; +pub mod interviewer; pub mod run_store; +pub mod runtime; pub use run_store::SqliteRunStore; diff --git a/lib/components/fabro-petri/src/runtime.rs b/lib/components/fabro-petri/src/runtime.rs new file mode 100644 index 000000000..b63509c84 --- /dev/null +++ b/lib/components/fabro-petri/src/runtime.rs @@ -0,0 +1,99 @@ +//! The Petri runtime Fabro runs its workflows on, assembled the same way at +//! create time (for `Runtime::check`) and at execution. +//! +//! The pieces are Petri's own: [`Runtime::standard`] with the Fabro frontend +//! carrying the server's settings layer, the Attractor step kinds (the real +//! ones, or the simulated registry for a dry run), the model client as the +//! `PebbleClient` capability so Petri's admission pass pins every LLM node's +//! route, and the Fabro home for the skills step. Nothing here knows about a +//! run: the store and the run options are added by the caller. + +use std::path::PathBuf; +use std::sync::Arc; + +use fabro_http::HttpClient; +use lithos_llm::Client; +use lithos_llm::catalog::{Catalog, ProviderId}; +use lithos_llm::client::ClientBuildError; +use lithos_llm::credentials::CredentialProvider; +use petri_attractor_steps::pebble::PebbleClient; +use petri_attractor_steps::skills::FabroHome; +use petri_frontend_fabro::Fabro; +use petri_runtime::Runtime; +use tracing::debug; + +/// What every Petri runtime Fabro builds is configured with. +#[derive(Clone, Default)] +pub struct RuntimeSpec { + /// The operator's settings layer, as `~/.fabro/settings.toml` text: the + /// lowest of the three layers the Fabro frontend reads (`[run.model]` + /// defaults, `[[run.hooks]]`, `[run.agent.mcps]`). + pub settings_toml: Option, + /// The model client the native agent and prompt steps call, and the + /// catalog the admission pass resolves model selectors against. `None` + /// leaves every LLM node unpinned and every model call unconfigured. + pub model_client: Option, + /// Run the simulated step registry (Fabro's `--dry-run` handlers) + /// instead of the real one. + pub dry_run: bool, + /// The Fabro home the skills step reads; `None` leaves it to Petri's + /// own lookup (`FABRO_HOME`, else `$HOME/.fabro`). + pub fabro_home: Option, +} + +impl RuntimeSpec { + /// Assemble the runtime. The admission pass that pins models is part of + /// the real registry, so a dry run's `check` still uses the real + /// registry: only execution swaps in the stubs. + #[must_use] + pub fn runtime(&self, for_execution: bool) -> Runtime { + let mut runtime = Runtime::standard() + .frontend(Fabro::new().with_settings_toml(self.settings_toml.clone())); + if let Some(client) = &self.model_client { + runtime = runtime.capability(PebbleClient(client.clone())); + } + let home = self + .fabro_home + .clone() + .map(FabroHome) + .or_else(FabroHome::from_env); + if let Some(home) = home { + runtime = runtime.capability(home); + } + if for_execution && self.dry_run { + petri_attractor_steps::register_stubs(runtime) + } else { + petri_attractor_steps::register(runtime) + } + } +} + +/// The model client Fabro hands Petri: the server's catalog, its credential +/// provider, its HTTP client (so a test's loopback client and a server's +/// proxy policy carry over), and only the providers whose credentials are +/// ready, the same eligible set the legacy compiler pinned models against. +/// `None` when no provider is eligible, so Petri's admission pass leaves the +/// graph alone rather than refusing every model. +pub fn model_client( + catalog: Catalog, + credentials: Arc, + http: Option, + eligible: &[ProviderId], +) -> Result, ClientBuildError> { + if eligible.is_empty() { + debug!("no eligible model provider; the Petri runtime gets no model client"); + return Ok(None); + } + let mut builder = Client::builder() + .catalog(catalog) + .credentials_arc(credentials) + .enabled_providers(eligible.iter().cloned()); + if let Some(http) = http { + builder = builder.http(http); + } + let build = builder.build()?; + for issue in &build.issues { + debug!(provider = %issue.provider, cause = %issue.cause, "model provider unavailable"); + } + Ok(Some(build.client)) +} diff --git a/lib/components/fabro-petri/tests/check.rs b/lib/components/fabro-petri/tests/check.rs new file mode 100644 index 000000000..351166810 --- /dev/null +++ b/lib/components/fabro-petri/tests/check.rs @@ -0,0 +1,224 @@ +//! Petri compiles at create time: `fabro_petri::check` materializes a bundle, +//! hands it to `Runtime::check`, and returns the admitted graphs or Petri's +//! diagnostics in Fabro's shape; `fabro_petri::admission` round-trips the +//! admitted graphs through Fabro's blob store. +//! +//! No sandbox plugin is needed: nothing here runs a graph. + +#![expect( + clippy::disallowed_methods, + reason = "the tests read checked-in fixture files synchronously before any run" +)] + +use std::collections::BTreeMap; +use std::path::{Path, PathBuf}; + +use fabro_auth::test_support::env_credential_source; +use fabro_llm::test_support::test_catalog; +use fabro_petri::admission; +use fabro_petri::check::{self, Bundle, CheckError, CheckRequest, DiagnosticSeverity, Launch}; +use fabro_petri::runtime::{self, RuntimeSpec}; +use fabro_store::{BlobStore, test_support}; +use lithos_llm::catalog::ProviderId; + +const COMMAND_WORKFLOW: &str = r#"digraph Command { + graph [goal="Run one command"] + start [shape=Mdiamond] + exit [shape=Msquare] + say [shape=parallelogram, script="echo hello from petri"] + start -> say -> exit +}"#; + +const UNKNOWN_ATTRIBUTE_WORKFLOW: &str = r#"digraph Bad { + graph [goal="Refuse me"] + start [shape=Mdiamond] + exit [shape=Msquare] + work [shape=box, prompt="Do the work", bogus="yes"] + start -> work -> exit +}"#; + +const UNKNOWN_MODEL_WORKFLOW: &str = r#"digraph Bad { + graph [goal="Refuse me"] + start [shape=Mdiamond] + exit [shape=Msquare] + work [shape=box, prompt="Do the work", model="no-such-model-9000"] + start -> work -> exit +}"#; + +const SETTINGS: &str = "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\n"; + +/// The `.fabro/workflows/hello` bundle checked into this repository. +fn hello_bundle() -> PathBuf { + Path::new(env!("CARGO_MANIFEST_DIR")).join("../../../.fabro/workflows/hello") +} + +fn bundle(files: &[(&str, &str)]) -> Bundle { + Bundle { + files: files + .iter() + .map(|(path, text)| ((*path).to_string(), (*text).to_string())) + .collect(), + entrypoint: "workflow.fabro".to_string(), + project_toml: None, + } +} + +fn request(bundle: Bundle, runtime: RuntimeSpec) -> CheckRequest { + CheckRequest { + bundle, + inputs: BTreeMap::new(), + launch: Launch::default(), + runtime, + } +} + +/// A runtime with a model client over the test catalog, with `openai` +/// eligible, as a server with an OpenAI key configured builds it. +fn runtime_with_openai() -> RuntimeSpec { + let credentials = env_credential_source(|name| match name { + "OPENAI_API_KEY" => Some("test-key".to_string()), + _ => None, + }); + let client = runtime::model_client(test_catalog(), credentials, None, &[ProviderId::new( + "openai", + )]) + .expect("the model client builds") + .expect("openai is eligible"); + RuntimeSpec { + model_client: Some(client), + ..RuntimeSpec::default() + } +} + +#[tokio::test] +async fn the_hello_bundle_is_admitted_and_round_trips_through_the_blob_store() { + let workflow = std::fs::read_to_string(hello_bundle().join("workflow.fabro")) + .expect("the hello workflow is checked in"); + let settings = std::fs::read_to_string(hello_bundle().join("workflow.toml")) + .expect("the hello settings are checked in"); + let request = request( + bundle(&[("workflow.fabro", &workflow), ("workflow.toml", &settings)]), + RuntimeSpec::default(), + ); + + let admitted = check::check(&request).expect("the hello bundle is admitted"); + + assert!( + admitted + .warnings + .iter() + .all(|w| w.severity == DiagnosticSeverity::Warning), + "{:?}", + admitted.warnings + ); + let blobs = BlobStore::new(test_support::in_memory_pool_with(&[ + fabro_db::BLOBS_MIGRATION_SQL, + ])); + let record = admission::persist(&blobs, &admitted) + .await + .expect("the graphs persist"); + assert!(record.children.is_empty()); + let (graph, children) = admission::load(&blobs, &record) + .await + .expect("the graphs load"); + assert_eq!(graph, admitted.graph); + assert!(children.is_empty()); +} + +#[tokio::test] +async fn a_launch_binds_the_repository_and_the_model_default() { + let repository = tempfile::tempdir().expect("a temp dir"); + let request = CheckRequest { + bundle: bundle(&[ + ("workflow.fabro", COMMAND_WORKFLOW), + ("workflow.toml", SETTINGS), + ]), + inputs: BTreeMap::new(), + launch: Launch { + model: Some("gpt-5.4".to_string()), + provider: None, + repository: Some(repository.path().to_path_buf()), + }, + runtime: RuntimeSpec::default(), + }; + + let admitted = check::check(&request).expect("the command bundle is admitted"); + + let launch = &admitted.graph.params["fabro.launch"]; + assert_eq!(launch["model"], "gpt-5.4"); + assert_eq!( + launch["clone"]["repository"], + repository.path().to_string_lossy().as_ref() + ); +} + +#[tokio::test] +async fn an_unknown_attribute_is_refused_with_petris_code() { + let request = request( + bundle(&[ + ("workflow.fabro", UNKNOWN_ATTRIBUTE_WORKFLOW), + ("workflow.toml", SETTINGS), + ]), + RuntimeSpec::default(), + ); + + let Err(CheckError::Rejected(diagnostics)) = check::check(&request) else { + panic!("an unknown attribute should be refused"); + }; + + let error = diagnostics + .iter() + .find(|d| d.code == "attractor.unknown_attribute") + .unwrap_or_else(|| panic!("no unknown-attribute diagnostic in {diagnostics:?}")); + assert!(error.is_error()); + assert!(error.message.contains("bogus"), "{error:?}"); + assert_eq!(error.file, "workflow.fabro"); + assert!(error.line.is_some(), "{error:?}"); +} + +#[tokio::test] +async fn an_unknown_model_is_refused_at_admission_when_a_catalog_is_installed() { + let request = request( + bundle(&[ + ("workflow.fabro", UNKNOWN_MODEL_WORKFLOW), + ("workflow.toml", SETTINGS), + ]), + runtime_with_openai(), + ); + + let Err(CheckError::Rejected(diagnostics)) = check::check(&request) else { + panic!("an unknown model should be refused when the runtime has a catalog"); + }; + + let error = diagnostics + .iter() + .find(|d| d.code == "attractor.model.unknown") + .unwrap_or_else(|| panic!("no model diagnostic in {diagnostics:?}")); + assert!(error.message.contains("no-such-model-9000"), "{error:?}"); +} + +#[tokio::test] +async fn a_known_model_is_pinned_at_admission() { + let workflow = UNKNOWN_MODEL_WORKFLOW.replace("no-such-model-9000", "gpt-5.4"); + let request = request( + bundle(&[("workflow.fabro", &workflow), ("workflow.toml", SETTINGS)]), + runtime_with_openai(), + ); + + let admitted = check::check(&request).expect("a catalog model is admitted"); + + let work = admitted + .graph + .body + .nodes + .iter() + .find(|node| node.name == "work") + .expect("the work node is in the graph"); + assert_eq!(work.step.config["provider"], "openai"); + assert_eq!(work.step.config["model"], "gpt-5.4"); + assert!( + work.step.config.get("plan").is_some(), + "{:?}", + work.step.config + ); +}