diff --git a/lib/apps/fabro-server/Cargo.toml b/lib/apps/fabro-server/Cargo.toml index 811dec41f..6248f2056 100644 --- a/lib/apps/fabro-server/Cargo.toml +++ b/lib/apps/fabro-server/Cargo.toml @@ -42,6 +42,7 @@ pebble-coding-agent.workspace = true fabro-llm = { path = "../../components/fabro-llm" } fabro-manifest = { path = "../../components/fabro-manifest" } fabro-mcp-store = { path = "../../components/fabro-mcp-store" } +fabro-petri = { path = "../../components/fabro-petri" } fabro-proc = { path = "../../foundation/fabro-proc" } fabro-template = { path = "../../foundation/fabro-template" } fabro-tool = { path = "../../components/fabro-tool" } diff --git a/lib/apps/fabro-server/src/run_compiler.rs b/lib/apps/fabro-server/src/run_compiler.rs index f8e842e29..b4b5cad66 100644 --- a/lib/apps/fabro-server/src/run_compiler.rs +++ b/lib/apps/fabro-server/src/run_compiler.rs @@ -13,7 +13,10 @@ //! fabro-workflow pipeline. //! 3. Model pinning — materialize run-level model settings against the catalog //! and the configured provider set. Stages 2's graph compilation and stage 3 -//! share one blocking dispatch via [`compile_and_pin`]. +//! share one blocking dispatch via [`compile_and_pin`]. A run Petri admitted +//! takes [`compile_admitted`] instead: Petri compiled, linted and pinned +//! models at its own admission, so only the Fabro graph the read side +//! displays is parsed here. //! 4. [`assemble_run`] — purely assemble the complete persistence input; no //! field is mutated after assembly. //! @@ -37,8 +40,8 @@ use fabro_llm::lithos_catalog::Catalog; use fabro_types::settings::interp::{InterpString, ResolveError}; use fabro_types::settings::run::{McpServerSettings, RunGoal}; use fabro_types::{ - AutomationRef, GitContext, ManifestPath, RunId, RunProvenance, RunTarget, WorkflowSettings, - WorkflowVersionId, + AutomationRef, GitContext, ManifestPath, PetriAdmission, RunEngine, RunId, RunProvenance, + RunTarget, WorkflowSettings, WorkflowVersionId, }; use fabro_util::workspace_glob::{WorkspaceGlob, WorkspaceGlobError}; use fabro_workflow::Error as WorkflowError; @@ -136,6 +139,19 @@ impl PreparedRun { &self.layered.settings } + /// The acquired bundle, for an engine that compiles it itself. + pub(crate) fn workflow_bundle(&self) -> &WorkflowBundle { + &self.layered.workflow_bundle + } + + pub(crate) fn entrypoint(&self) -> &ManifestPath { + &self.layered.entrypoint + } + + pub(crate) fn target(&self) -> Option<&RunTarget> { + self.layered.metadata.target.as_ref() + } + pub(crate) fn with_target_and_git( mut self, target: RunTarget, @@ -173,6 +189,8 @@ struct GraphCompiledRun { pub(crate) struct PinnedRun { materialized: MaterializedRun, metadata: RunMetadata, + /// The engine the run was created for, with what it admitted. + engine: RunEngine, } #[derive(Debug, thiserror::Error)] @@ -391,6 +409,50 @@ pub(crate) async fn compile_and_pin( })? } +/// Stages two and three for a run Petri admitted: parse the Fabro graph +/// the read side displays, with no lint and no model pinning, and record +/// the admission on the run. +pub(crate) async fn compile_admitted( + prepared: PreparedRun, + admission: PetriAdmission, +) -> Result { + task::spawn_blocking(move || { + let PreparedRun { + layered: + LayeredRun { + workflow_bundle, + entrypoint, + workflow, + settings, + cwd, + metadata, + }, + vars, + } = prepared; + let compiled = operations::compile_admitted_run(CreateRunCompileInput { + workflow: WorkflowInput::Bundled(workflow), + settings, + vars, + cwd, + workflow_path: Some(entrypoint), + workflow_bundle: Some(workflow_bundle), + configured_providers: Vec::new(), + })?; + Ok(PinnedRun { + materialized: operations::materialize_admitted_run(compiled), + metadata, + engine: RunEngine::Petri(admission), + }) + }) + .await + .map_err(|source| { + RunCompilerError::Workflow(WorkflowError::engine_with_source( + "workflow create task failed", + source, + )) + })? +} + /// Stage two's graph compilation: parse, transform, and validate through the /// fabro-workflow pipeline, with undefined template variables promoted to /// hard errors. @@ -435,6 +497,7 @@ fn pin_models(compiled: GraphCompiledRun, catalog: &Catalog) -> Result CreateRunPersistenceInput { let PinnedRun { materialized, metadata, + engine, } = pinned; let RunMetadata { run_id, @@ -472,7 +536,7 @@ pub(crate) fn assemble_run(pinned: PinnedRun) -> CreateRunPersistenceInput { parent_id, provenance, web_url, - engine: fabro_types::RunEngine::Legacy, + engine, }) } diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 7c3a7c83e..c8f85c871 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -88,7 +88,7 @@ use fabro_types::settings::server::{ GithubIntegrationSettings, GithubIntegrationStrategy, LogDestination, }; use fabro_types::{ - AgentBackend, AskFabro, AskFabroUnavailableReason, BlobHash, EventBody, + AgentBackend, AskFabro, AskFabroUnavailableReason, BlobHash, Engine, EventBody, InterviewQuestionRecord, ModelRef, ModelTestMode, PairId, PairMessageId, PairTarget, PendingReason, Principal, PullRequestLink, QuestionType, RunControlAction, RunEvent, RunId, RunRunnableSource, RunStatusKind, SandboxProviderKind, ServerSettings, SessionCapability, @@ -166,6 +166,7 @@ use crate::{ mod automation_scheduler; mod handler; +mod petri_runs; mod pull_request_supervisor; pub(crate) mod resource_sampler; mod session_runtime; @@ -1126,6 +1127,9 @@ pub struct AppState { pub(super) server_secrets: ServerSecrets, pub(crate) llm_source: Arc, + /// The database pool the stores share, for the Petri run store a + /// server-process run writes its records through. + pub(crate) db_pool: DbPool, manifest_run_defaults: RwLock>, manifest_run_settings: RwLock>, pub(crate) server_settings: RwLock>, @@ -2430,6 +2434,7 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result anyhow::Result, run_id: RunId) { return; } + match run_engine(&state, run_id).await { + Ok(Engine::Petri) => { + Box::pin(petri_runs::execute(state, run_id)).await; + return; + } + Ok(Engine::Legacy) => {} + Err(err) => { + tracing::error!(run_id = %run_id, error = %err, "Failed to read the run's engine"); + fail_managed_run( + &state, + run_id, + FailureReason::WorkflowError, + format!("Failed to read the run's engine: {err}"), + ); + state.scheduler_notify.notify_one(); + return; + } + } + if state.registry_factory_override.is_some() { Box::pin(execute_run_in_process(state, run_id)).await; return; @@ -3965,6 +3990,13 @@ async fn execute_run(state: Arc, run_id: RunId) { Box::pin(execute_run_subprocess(state, run_id)).await; } +/// The engine the run was created for, from its stored spec. +async fn run_engine(state: &AppState, run_id: RunId) -> anyhow::Result { + let run_store = state.stores.runs.open_run(&run_id).await?; + let run_state = run_store.state().await?; + Ok(run_state.spec.engine.engine()) +} + async fn execute_run_in_process(state: Arc, run_id: RunId) { // Transition to Starting and set up cancel infrastructure let (cancel_rx, run_dir, event_tx, cancel_token, execution_mode) = { diff --git a/lib/apps/fabro-server/src/server/handler/runs.rs b/lib/apps/fabro-server/src/server/handler/runs.rs index 3437e9b0f..43030cb64 100644 --- a/lib/apps/fabro-server/src/server/handler/runs.rs +++ b/lib/apps/fabro-server/src/server/handler/runs.rs @@ -27,11 +27,11 @@ use fabro_store::{ RunSummaryListQuery, RunSummarySort, RunSummarySortDirection, RunSummaryVisibility, }; use fabro_types::{ - AutomationRef, ContextWindowStaleness, ManifestPath, Principal, Run, RunClientProvenance, - RunId, RunProvenance, RunServerProvenance, RunStatusKind, RunTarget, SandboxProviderKind, - StageContextWindow, StageContextWindowUnavailableReason, StageHandler, StageModelUsage, - StageProjection, SystemActorKind, ValidatedRunTarget, json_scalar_to_toml_value, - parse_blob_ref, + AutomationRef, ContextWindowStaleness, Engine, ManifestPath, Principal, Run, + RunClientProvenance, RunId, RunProvenance, RunServerProvenance, RunStatusKind, RunTarget, + SandboxProviderKind, StageContextWindow, StageContextWindowUnavailableReason, StageHandler, + StageModelUsage, StageProjection, SystemActorKind, ValidatedRunTarget, + json_scalar_to_toml_value, parse_blob_ref, }; use fabro_util::error as error_util; use fabro_util::version::FABRO_VERSION; @@ -48,7 +48,8 @@ use super::super::{ AppState, DeleteRunOutcome, ListResponse, RunExecutionMode, VariableError, answer_from_request, api_question_from_pending_interview, clamp_page_limit, clamp_page_offset, default_page_limit, delete_run_internal, load_pending_interview, managed_run, parse_run_id_path, - parse_stage_id_path, reject_if_archived, submit_pending_interview_answer, workflow_event, + parse_stage_id_path, petri_runs, reject_if_archived, submit_pending_interview_answer, + workflow_event, }; use crate::error::ApiError; use crate::principal_middleware::{ @@ -767,13 +768,27 @@ async fn finalize_created_run( ready_provider_ids.clone() } }; - let pinned = - match run_compiler::compile_and_pin(prepared, run_materialization_provider_ids, catalog) - .await - { - Ok(pinned) => pinned, - Err(error) => return run_intent_admission_error(error.into()), - }; + // Petri compiles a Petri run: the bundle goes to `Runtime::check`, its + // diagnostics come back in Fabro's shape, and the admitted graph is what + // the run executes. The legacy compile, lint and model pinning are + // skipped for it; Fabro's own settings resolution ran above as for any + // run. + let engine = petri_runs::engine_for(prepared.settings(), &state.server_settings()); + let pinned = match engine { + Engine::Legacy => { + run_compiler::compile_and_pin(prepared, run_materialization_provider_ids, catalog).await + } + Engine::Petri => { + match petri_runs::admit(&state, &prepared, &run_materialization_provider_ids).await { + Ok(admission) => run_compiler::compile_admitted(prepared, admission).await, + Err(error) => Err(error), + } + } + }; + let pinned = match pinned { + Ok(pinned) => pinned, + Err(error) => return run_intent_admission_error(error.into()), + }; let persistence_input = run_compiler::assemble_run(pinned); let created = match Box::pin(operations::persist_create_run( state.stores.runs.as_ref(), @@ -947,9 +962,14 @@ fn run_intent_admission_error(error: RunIntentAdmissionError) -> Response { ), }, // Return the curated compiler detail; retain its source chain in the log. + // A validation failure names its diagnostics, since the message alone + // ("Validation failed") tells the caller nothing to fix. RunIntentAdmissionError::Compiler(error) => intent_error( StatusCode::UNPROCESSABLE_ENTITY, - format!("run intent could not be compiled: {error}"), + format!( + "run intent could not be compiled: {}", + compiler_error_detail(&error) + ), "run_compile_invalid", ), RunIntentAdmissionError::VariableSnapshot { .. } => intent_error( @@ -970,6 +990,26 @@ fn run_intent_admission_error(error: RunIntentAdmissionError) -> Response { } } +/// The compiler error's text, with every error diagnostic of a validation +/// failure listed as `rule: message`. +fn compiler_error_detail(error: &run_compiler::RunCompilerError) -> String { + let run_compiler::RunCompilerError::Workflow(WorkflowError::ValidationFailed { diagnostics }) = + error + else { + return error.to_string(); + }; + let listed = diagnostics + .iter() + .filter(|diagnostic| diagnostic.severity == fabro_validate::Severity::Error) + .map(|diagnostic| format!("{}: {}", diagnostic.rule, diagnostic.message)) + .collect::>(); + if listed.is_empty() { + error.to_string() + } else { + format!("{error}: {}", listed.join("; ")) + } +} + async fn validate_intent_actor_target( state: &AppState, actor: &Principal, diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs new file mode 100644 index 000000000..b6944865a --- /dev/null +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -0,0 +1,440 @@ +//! Runs on Petri: what the server does at create time and at execution when +//! a run's engine is Petri. +//! +//! At create, [`admit`] hands the workflow version's bundle, the run's inputs +//! and the launch to Petri's `Runtime::check` through `fabro_petri::check`, +//! maps Petri's diagnostics onto Fabro's, and stores the admitted graphs in +//! 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. + +use std::collections::{BTreeMap, HashSet}; +use std::sync::Arc; +use std::time::Instant; + +use fabro_config::SettingsLayer; +use fabro_llm::selection; +use fabro_petri::check::{self, Bundle, CheckError, CheckRequest, Diagnostic, Launch}; +use fabro_petri::engine::{self, RunOutcome, RunRequest}; +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, +}; +use fabro_util::error as error_util; +use fabro_validate::{Diagnostic as FabroDiagnostic, Severity}; +use fabro_workflow::Error as WorkflowError; +use fabro_workflow::run_status::{FailureReason, RunStatus, SuccessReason}; +use lithos_llm::catalog::ProviderId; +use tokio::task; +use tokio_util::sync::CancellationToken; +use tracing::{error, info, warn}; + +use super::{AppState, clear_live_run_state, workflow_event}; +use crate::run_compiler::{PreparedRun, RunCompilerError}; + +/// The engine a run gets: the one its workflow version names, else the +/// server's default. +pub(crate) fn engine_for( + settings: &fabro_types::WorkflowSettings, + server: &ServerSettings, +) -> Engine { + settings + .workflow + .engine + .unwrap_or(server.server.execution.engine) +} + +/// The runtime Petri gets, at create and at execution: the server's run +/// defaults as the settings layer, the model client over the server's +/// catalog and credentials for the eligible providers, and the run mode. +pub(crate) fn runtime_spec( + state: &AppState, + eligible: &[ProviderId], + dry_run: bool, +) -> RuntimeSpec { + let settings_toml = settings_layer_toml(state); + let catalog = state.catalog(); + let model_client = match runtime::model_client( + (*catalog).clone(), + Arc::clone(&state.llm_source), + state.http_client.clone(), + eligible, + ) { + Ok(client) => client, + Err(err) => { + warn!(error = %err, "Petri model client unavailable; LLM nodes stay unpinned"); + None + } + }; + RuntimeSpec { + settings_toml, + model_client, + dry_run, + fabro_home: None, + } +} + +/// The server's `[run]` defaults, as the text of the operator settings +/// layer the Fabro frontend reads below `.fabro/project.toml` and +/// `workflow.toml`. +fn settings_layer_toml(state: &AppState) -> Option { + let layer = SettingsLayer { + version: Some(1), + run: Some((*state.manifest_run_defaults()).clone()), + ..SettingsLayer::default() + }; + match toml::to_string(&layer) { + Ok(text) => Some(text), + Err(err) => { + warn!(error = %err, "server run defaults do not serialize; Petri gets no settings layer"); + None + } + } +} + +/// Petri compiles the run: check the bundle, map the diagnostics, and +/// persist the admitted graphs. A refusal is the same validation error the +/// legacy compiler raises, carrying Petri's diagnostics. +pub(crate) async fn admit( + state: &AppState, + prepared: &PreparedRun, + eligible: &[ProviderId], +) -> Result { + let settings = prepared.settings(); + let mut files = BTreeMap::new(); + for workflow in prepared.workflow_bundle().workflows().values() { + for (path, text) in &workflow.files { + files.insert(path.to_string(), text.clone()); + } + files.insert(workflow.path.to_string(), workflow.source.clone()); + if let Some(config) = &workflow.config { + files.insert(config.path.to_string(), config.source.clone()); + } + } + let mut inputs = BTreeMap::new(); + for (name, value) in &settings.run.inputs { + let value = serde_json::to_value(value).map_err(|err| { + RunCompilerError::Workflow(WorkflowError::engine_with_source( + format!("run input `{name}` does not encode as JSON"), + err, + )) + })?; + inputs.insert(name.clone(), value); + } + // The launch: Fabro's resolved model and provider. When the settings + // name neither, the default offering of the eligible providers, as the + // legacy compiler picked it, is bound as the launch model alone: a node + // that names no model runs on it, and a node that names a model the + // catalog lacks stays unqualified, so Petri's admission refuses it. + let catalog = state.catalog(); + let model = settings.run.model.name.clone().or_else(|| { + if settings.run.model.provider.is_some() { + return None; + } + let eligible = eligible.iter().cloned().collect::>(); + selection::select_default(&catalog, &eligible) + .ok() + .map(|offering| offering.model.id().to_string()) + }); + let provider = settings.run.model.provider.clone(); + let repository = match prepared.target() { + Some(RunTarget::Folder { path }) => Some(path.into()), + Some(RunTarget::Git(_) | RunTarget::None {}) | None => None, + }; + let dry_run = settings.run.execution.mode == RunMode::DryRun; + let request = CheckRequest { + bundle: Bundle { + files, + entrypoint: prepared.entrypoint().to_string(), + project_toml: None, + }, + inputs, + launch: Launch { + model, + provider, + repository, + }, + runtime: runtime_spec(state, eligible, dry_run), + }; + let admitted = task::spawn_blocking(move || check::check(&request)) + .await + .map_err(|source| { + RunCompilerError::Workflow(WorkflowError::engine_with_source( + "Petri check task failed", + source, + )) + })? + .map_err(|err| match err { + CheckError::Rejected(diagnostics) => { + RunCompilerError::Workflow(WorkflowError::ValidationFailed { + diagnostics: diagnostics.iter().map(fabro_diagnostic).collect(), + }) + } + other => RunCompilerError::Workflow(WorkflowError::engine_with_source( + "Petri could not check the workflow", + other, + )), + })?; + for warning in &admitted.warnings { + info!(code = %warning.code, message = %warning.message, "Petri warned at admission"); + } + admission::persist(&state.store_ref().blobs(), &admitted) + .await + .map_err(|err| { + RunCompilerError::Workflow(WorkflowError::engine_with_source( + "the admitted graphs could not be stored", + err, + )) + }) +} + +/// Petri's diagnostic in Fabro's shape: the code is the rule, the hint is +/// the fix, the bundle-relative file and position are the source location. +fn fabro_diagnostic(diagnostic: &Diagnostic) -> FabroDiagnostic { + FabroDiagnostic { + rule: diagnostic.code.clone(), + severity: if diagnostic.is_error() { + Severity::Error + } else { + Severity::Warning + }, + message: diagnostic.message.clone(), + fix: diagnostic.hint.clone(), + source_path: Some(diagnostic.file.clone()), + line: diagnostic.line, + column: diagnostic.column, + ..FabroDiagnostic::default() + } +} + +/// Execute a Petri run in the server process: runnable → starting → running +/// → succeeded or failed, with the lifecycle events Fabro's read side needs. +pub(crate) async fn execute(state: Arc, run_id: RunId) { + let (run_dir, cancel) = { + 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, + _ => return, + }; + let Some(run_dir) = managed_run.run_dir.clone() else { + return; + }; + let cancel = CancellationToken::new(); + managed_run.status = RunStatus::Starting; + managed_run.cancel_token = Some(cancel.clone()); + (run_dir, cancel) + }; + + let run_store = match state.stores.runs.open_run(&run_id).await { + Ok(run_store) => run_store, + Err(err) => { + error!(run_id = %run_id, error = %err, "Failed to open run store"); + finish( + &state, + run_id, + RunStatus::Failed { + reason: FailureReason::WorkflowError, + }, + Some(format!("Failed to open run store: {err}")), + ); + return; + } + }; + tokio::spawn(super::forward_run_events_to_global( + Arc::clone(&state), + run_id, + run_store.subscribe(), + )); + let run_state = match run_store.state().await { + Ok(run_state) => run_state, + Err(err) => { + error!(run_id = %run_id, error = %err, "Failed to load run state"); + finish( + &state, + run_id, + RunStatus::Failed { + reason: FailureReason::WorkflowError, + }, + Some(format!("Failed to load run state: {err}")), + ); + return; + } + }; + let Some(admission) = run_state.spec.engine.petri().cloned() else { + fail_before_execution(&state, &run_store, run_id, "the run has no Petri admission").await; + return; + }; + let server_settings = state.server_settings(); + if super::reject_run_if_sandbox_provider_disabled( + &state, + &server_settings, + run_id, + &run_state.spec.settings.run, + ) + .await + { + return; + } + let started = Instant::now(); + for event in [ + workflow_event::Event::RunStarting, + workflow_event::Event::RunRunning, + ] { + 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 lifecycle event"); + finish( + &state, + run_id, + RunStatus::Failed { + reason: FailureReason::WorkflowError, + }, + Some(format!("Failed to persist run lifecycle event: {err}")), + ); + return; + } + } + { + let mut runs = state.runs.lock().expect("runs lock poisoned"); + if let Some(managed_run) = runs.get_mut(&run_id) { + if managed_run.status == RunStatus::Starting { + managed_run.status = RunStatus::Running; + } + } + } + let (_, eligible) = state.resolve_llm_client_with_ready_ids().await; + let dry_run = run_state.spec.settings.run.execution.mode == RunMode::DryRun; + let request = RunRequest { + run_id: run_id.to_string(), + run_dir: run_dir.join("petri"), + admission, + blobs: state.store_ref().blobs(), + 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 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, + .. + }) => { + info!(run_id = %run_id, "Petri run completed"); + ( + RunStatus::Succeeded { + reason: SuccessReason::Completed, + }, + None, + workflow_event::Event::WorkflowRunCompleted { + timing, + artifact_count: 0, + status: StageOutcome::Succeeded.to_string(), + reason: SuccessReason::Completed, + final_git_commit_sha: None, + final_patch: None, + diff_summary: None, + usage: None, + }, + ) + } + 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); + 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"); + } + 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(), + }; + if !outcome.complete { + message.push_str(" (record incomplete: "); + message.push_str(&outcome.incomplete.join("; ")); + message.push(')'); + } + message +} + +/// The failed status, its message, and the `run.failed` event for it. +fn failed( + reason: FailureReason, + message: String, + timing: RunTiming, +) -> (RunStatus, Option, workflow_event::Event) { + let error = match reason { + FailureReason::Cancelled => WorkflowError::Cancelled, + _ => WorkflowError::engine(message.clone()), + }; + ( + RunStatus::Failed { reason }, + Some(message), + workflow_event::Event::workflow_run_failed_from_error( + &error, timing, reason, None, None, None, None, + ), + ) +} + +/// Record a failure that happened before Petri ran, then finish the run. +async fn fail_before_execution( + state: &Arc, + run_store: &fabro_store::RunDatabase, + run_id: RunId, + message: &str, +) { + error!(run_id = %run_id, error = message, "Petri run cannot start"); + let (status, error, event) = failed( + FailureReason::WorkflowError, + message.to_string(), + RunTiming::default(), + ); + 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 failure status"); + } + finish(state, run_id, status, error); +} + +/// Settle the managed run and release its scheduler slot. +fn finish(state: &Arc, run_id: RunId, status: RunStatus, error: Option) { + let mut runs = state.runs.lock().expect("runs lock poisoned"); + if let Some(managed_run) = runs.get_mut(&run_id) { + managed_run.status = status; + managed_run.error = error; + clear_live_run_state(managed_run); + } + drop(runs); + state.scheduler_notify.notify_one(); +} diff --git a/lib/apps/fabro-server/src/test_support.rs b/lib/apps/fabro-server/src/test_support.rs index d3da04a44..4f8aeec7e 100644 --- a/lib/apps/fabro-server/src/test_support.rs +++ b/lib/apps/fabro-server/src/test_support.rs @@ -707,6 +707,13 @@ pub(crate) fn load_test_server_secrets( ServerSecrets::load(path, env).expect("test server secrets should load") } +/// The database pool the app state's stores share, for a test that reads +/// what a run wrote through another store over the same database. +#[must_use] +pub fn test_app_db_pool(state: &AppState) -> DbPool { + state.db_pool.clone() +} + pub fn test_secret_store_path() -> PathBuf { let dir = std::env::temp_dir().join(format!("fabro-test-{}", Ulid::new())); std::fs::create_dir_all(&dir).expect("test temp dir should be creatable"); diff --git a/lib/apps/fabro-server/tests/it/scenario/mod.rs b/lib/apps/fabro-server/tests/it/scenario/mod.rs index 08b3936b4..a58db82af 100644 --- a/lib/apps/fabro-server/tests/it/scenario/mod.rs +++ b/lib/apps/fabro-server/tests/it/scenario/mod.rs @@ -1,6 +1,7 @@ mod archive; mod dry_run; mod lifecycle; +mod petri; mod run_completion; mod sse; mod usage; diff --git a/lib/apps/fabro-server/tests/it/scenario/petri.rs b/lib/apps/fabro-server/tests/it/scenario/petri.rs new file mode 100644 index 000000000..4eeade6a4 --- /dev/null +++ b/lib/apps/fabro-server/tests/it/scenario/petri.rs @@ -0,0 +1,382 @@ +//! Runs on Petri through the server: a run goes to Petri when its workflow +//! version names `engine = "petri"` or when the server's +//! `[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 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 +//! refusals need no plugin and always run. + +#![expect( + clippy::disallowed_methods, + reason = "the tests locate the plugin executable through the process environment" +)] +#![expect(clippy::print_stderr, reason = "a skipped test says why on its stderr")] + +use std::collections::BTreeMap; +use std::env; +use std::path::PathBuf; +use std::sync::Arc; + +use axum::body::Body; +use axum::http::{Request, StatusCode}; +use fabro_petri::SqliteRunStore; +use fabro_petri::engine::{self, RunStatus}; +use fabro_server::server::AppState; +use fabro_server::test_support::{ + TestAppStateBuilder, llm_overlay_with_provider_base_url, test_app_db_pool, + test_register_workflow_version, +}; +use fabro_static::EnvVars; +use fabro_test::{TwinScenario, TwinScenarios, twin_openai}; +use fabro_types::{WorkflowPath, WorkflowVersion}; +use tower::ServiceExt; + +use crate::helpers::{ + api, create_and_start_run_from_intent, read_repo_file, response_json, run_json, + settings_from_toml, test_app_state_with_options, test_app_with_scheduler, test_settings, + wait_for_run_status, +}; + +const HOST_PLUGIN: &str = "sandbox-driver-host"; +const HOST_PLUGIN_OVERRIDE: &str = "PETRI_SANDBOX_HOST_PLUGIN"; +const REQUIRE_ENV: &str = "FABRO_REQUIRE_SANDBOX_PLUGINS"; + +const OPENAI_MODEL: &str = "gpt-5.4"; + +/// A command-only workflow: one script stage between start and exit. +const COMMAND_DOT: &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 +}"#; + +/// A workflow whose one stage carries an attribute the language does not +/// have. +const UNKNOWN_ATTRIBUTE_DOT: &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 +}"#; + +/// A workflow with an edge to a node nobody declared. +const UNDECLARED_NODE_DOT: &str = r#"digraph Bad { + graph [goal="Refuse me"] + start [shape=Mdiamond] + exit [shape=Msquare] + work [shape=box, prompt="Do the work"] + start -> work -> nowhere -> exit +}"#; + +/// A workflow whose stage names a model no catalog has. +const UNKNOWN_MODEL_DOT: &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 PLAIN_SETTINGS: &str = "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\n"; +const PETRI_SETTINGS: &str = + "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\nengine = \"petri\"\n"; + +/// The host plugin as Petri's lookup finds it: the override variable, else +/// the executable on `PATH`. `None`, after saying so, when the test should +/// skip; a panic when the environment forbids a skip. +fn host_plugin() -> Option { + let found = env::var_os(HOST_PLUGIN_OVERRIDE) + .map(PathBuf::from) + .or_else(|| { + env::split_paths(&env::var_os("PATH")?) + .map(|dir| dir.join(HOST_PLUGIN)) + .find(|candidate| candidate.is_file()) + }); + if found.is_none() { + assert!( + env::var_os(REQUIRE_ENV).is_none(), + "{REQUIRE_ENV} is set, but {HOST_PLUGIN} is not on PATH and {HOST_PLUGIN_OVERRIDE} is unset" + ); + eprintln!("skipping: {HOST_PLUGIN} is not on PATH and {HOST_PLUGIN_OVERRIDE} is unset"); + } + found +} + +/// Register a version whose entrypoint is `workflow.fabro`, with the given +/// files beside it. +async fn register_version(app: &axum::Router, files: &[(&str, &str)]) -> String { + let entrypoint = WorkflowPath::new("workflow.fabro").expect("entrypoint path is valid"); + let files = files + .iter() + .map(|(path, text)| { + ( + WorkflowPath::new(*path).expect("fixture path is valid"), + (*text).to_string(), + ) + }) + .collect::>(); + let version = + WorkflowVersion::new(entrypoint, files, BTreeMap::new()).expect("fixture version is valid"); + test_register_workflow_version(app, &version, None) + .await + .to_string() +} + +fn intent(version_id: &str, workspace: &std::path::Path) -> serde_json::Value { + serde_json::json!({ + "workflow_version_id": version_id, + "target": {"kind": "folder", "path": workspace}, + "environment_id": "local", + "args": {}, + }) +} + +/// The `hello` bundle checked into this repository, with `engine = "petri"` +/// added to its `[workflow]` table. +fn hello_files() -> [(&'static str, String); 2] { + let workflow = read_repo_file(".fabro/workflows/hello/workflow.fabro"); + let settings = read_repo_file(".fabro/workflows/hello/workflow.toml"); + assert!( + settings.trim_end().ends_with("graph = \"workflow.fabro\""), + "the hello settings end with the [workflow] table, so an engine key appends to it" + ); + [ + ("workflow.fabro", workflow), + ( + "workflow.toml", + format!("{}\nengine = \"petri\"\n", settings.trim_end()), + ), + ] +} + +/// The run's record in Petri's store, read through the same database the +/// server wrote it to. +async fn petri_outcome(state: &AppState, run_id: &str) -> engine::RunOutcome { + let store = SqliteRunStore::new(test_app_db_pool(state)); + engine::outcome_of(&store, run_id) + .await + .expect("the run's Petri record inspects") +} + +async fn run_engine(app: &axum::Router, run_id: &str) -> serde_json::Value { + let req = Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}/state"))) + .body(Body::empty()) + .expect("state request should build"); + let response = app + .clone() + .oneshot(req) + .await + .expect("state request routes"); + let body = response_json( + response, + StatusCode::OK, + format!("GET /api/v1/runs/{run_id}/state"), + ) + .await; + body["spec"]["engine"].clone() +} + +async fn create_run_response(app: &axum::Router, intent: serde_json::Value) -> serde_json::Value { + let req = Request::builder() + .method("POST") + .uri(api("/runs")) + .header("content-type", "application/json") + .body(Body::from( + serde_json::to_string(&intent).expect("intent serializes"), + )) + .expect("create-run request should build"); + let response = app + .clone() + .oneshot(req) + .await + .expect("create request routes"); + response_json( + response, + StatusCode::UNPROCESSABLE_ENTITY, + "POST /api/v1/runs", + ) + .await +} + +/// The `hello` bundle, whose one stage is a prompt, runs on Petri when its +/// version names the engine: the prompt reaches the twin through Petri's +/// model client, Fabro reports the run succeeded, and Petri's record of the +/// run says the same. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn the_hello_bundle_runs_on_petri_when_the_version_names_the_engine() { + if host_plugin().is_none() { + return; + } + let workspace = tempfile::tempdir().expect("workspace tempdir"); + let twin = twin_openai().await; + let namespace = format!("{}::{}", module_path!(), line!()); + TwinScenarios::new(&namespace) + .scenario(TwinScenario::responses(OPENAI_MODEL).text("A haiku, added.")) + .load(twin) + .await; + let settings = test_settings(); + let state = TestAppStateBuilder::new() + .runtime_settings(settings.server_settings, settings.manifest_run_defaults) + .max_concurrent_runs(5) + .llm_overlay(llm_overlay_with_provider_base_url( + "openai", + twin.base_url.clone(), + )) + .vault_entries([(EnvVars::OPENAI_API_KEY, namespace.clone())]) + .build(); + let app = test_app_with_scheduler(Arc::clone(&state)); + + let [(workflow_path, workflow), (settings_path, settings)] = hello_files(); + let version_id = register_version(&app, &[ + (workflow_path, &workflow), + (settings_path, &settings), + ]) + .await; + let mut intent = intent(&version_id, workspace.path()); + intent["args"]["model"] = serde_json::json!(OPENAI_MODEL); + let run_id = create_and_start_run_from_intent(&app, intent).await; + + let status = wait_for_run_status(&app, &run_id, &["succeeded", "failed"]).await; + let run = run_json(&app, &run_id).await; + assert_eq!(status, "succeeded", "run: {run}"); + assert_eq!(run_engine(&app, &run_id).await["kind"], "petri"); + let outcome = petri_outcome(&state, &run_id).await; + assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}"); + assert!(outcome.complete, "{:?}", outcome.incomplete); + let logs = twin.request_logs(&namespace).await; + let requests = logs["requests"] + .as_array() + .expect("twin request logs are an array"); + assert!( + requests + .iter() + .any(|request| request["model"] == OPENAI_MODEL), + "the prompt stage should have called the twin, got {logs}" + ); +} + +/// A command-only bundle runs on Petri when the server's setting names the +/// engine and the version names none, and Petri's record agrees. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_command_bundle_runs_on_petri_under_the_server_setting() { + if host_plugin().is_none() { + return; + } + let workspace = tempfile::tempdir().expect("workspace tempdir"); + let settings = settings_from_toml( + "_version = 1\n\n[run.environment]\nid = \"local\"\n\n[server.execution]\nengine = \ + \"petri\"\n", + ); + let state = test_app_state_with_options(settings, 5); + let app = test_app_with_scheduler(Arc::clone(&state)); + + let version_id = register_version(&app, &[ + ("workflow.fabro", COMMAND_DOT), + ("workflow.toml", PLAIN_SETTINGS), + ]) + .await; + let run_id = + create_and_start_run_from_intent(&app, intent(&version_id, workspace.path())).await; + + let status = wait_for_run_status(&app, &run_id, &["succeeded", "failed"]).await; + let run = run_json(&app, &run_id).await; + assert_eq!(status, "succeeded", "run: {run}"); + assert_eq!(run_engine(&app, &run_id).await["kind"], "petri"); + let outcome = petri_outcome(&state, &run_id).await; + assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}"); + assert!(outcome.complete, "{:?}", outcome.incomplete); +} + +/// A version that names no engine on a server whose setting is the default +/// keeps the legacy executor: the run's spec records no Petri admission. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_version_that_names_no_engine_stays_on_the_legacy_executor() { + let workspace = tempfile::tempdir().expect("workspace tempdir"); + let state = test_app_state_with_options(test_settings(), 5); + let app = test_app_with_scheduler(state); + + let version_id = register_version(&app, &[ + ("workflow.fabro", COMMAND_DOT), + ("workflow.toml", PLAIN_SETTINGS), + ]) + .await; + let mut intent = intent(&version_id, workspace.path()); + intent["args"]["dry_run"] = serde_json::json!(true); + let run_id = create_and_start_run_from_intent(&app, intent).await; + + assert_eq!(run_engine(&app, &run_id).await, serde_json::Value::Null); +} + +/// A workflow with an attribute the language does not have is refused at +/// create with Petri's code in Fabro's diagnostic shape. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn an_unknown_attribute_is_refused_at_create_with_petris_code() { + let workspace = tempfile::tempdir().expect("workspace tempdir"); + let state = test_app_state_with_options(test_settings(), 5); + let app = test_app_with_scheduler(state); + + let version_id = register_version(&app, &[ + ("workflow.fabro", UNKNOWN_ATTRIBUTE_DOT), + ("workflow.toml", PETRI_SETTINGS), + ]) + .await; + let body = create_run_response(&app, intent(&version_id, workspace.path())).await; + + let detail = body["errors"][0]["detail"].as_str().unwrap_or_default(); + assert_eq!(body["errors"][0]["code"], "run_compile_invalid", "{body}"); + assert!( + detail.contains("attractor.unknown_attribute") && detail.contains("bogus"), + "expected Petri's diagnostic in the detail, got {body}" + ); +} + +/// An edge to a node nobody declared is refused at create with Petri's code. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn an_edge_to_an_undeclared_node_is_refused_at_create_with_petris_code() { + let workspace = tempfile::tempdir().expect("workspace tempdir"); + let state = test_app_state_with_options(test_settings(), 5); + let app = test_app_with_scheduler(state); + + let version_id = register_version(&app, &[ + ("workflow.fabro", UNDECLARED_NODE_DOT), + ("workflow.toml", PETRI_SETTINGS), + ]) + .await; + let body = create_run_response(&app, intent(&version_id, workspace.path())).await; + + let detail = body["errors"][0]["detail"].as_str().unwrap_or_default(); + assert!( + detail.contains("attractor.undeclared_node") && detail.contains("nowhere"), + "expected Petri's diagnostic in the detail, got {body}" + ); +} + +/// A model selector the catalog cannot resolve is refused at create with +/// `attractor.model.unknown`: Petri's admission pass pins every model +/// against the server's catalog, so nothing is left for a run to discover. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn an_unknown_model_is_refused_at_create_with_attractor_model_unknown() { + let workspace = tempfile::tempdir().expect("workspace tempdir"); + let state = test_app_state_with_options(test_settings(), 5); + let app = test_app_with_scheduler(state); + + let version_id = register_version(&app, &[ + ("workflow.fabro", UNKNOWN_MODEL_DOT), + ("workflow.toml", PETRI_SETTINGS), + ]) + .await; + let body = create_run_response(&app, intent(&version_id, workspace.path())).await; + + let detail = body["errors"][0]["detail"].as_str().unwrap_or_default(); + assert!( + detail.contains("attractor.model.unknown") && detail.contains("no-such-model-9000"), + "expected the admission diagnostic in the detail, got {body}" + ); +} diff --git a/lib/components/fabro-workflow/src/operations/create.rs b/lib/components/fabro-workflow/src/operations/create.rs index d60e31ec8..555d900f5 100644 --- a/lib/components/fabro-workflow/src/operations/create.rs +++ b/lib/components/fabro-workflow/src/operations/create.rs @@ -274,6 +274,95 @@ pub async fn create( Box::pin(persist_create_run(store, persistence_input)).await } +/// Stage two for a run another engine admitted: the Fabro graph is parsed +/// and transformed for the read side (the goal, the node count, labels), with +/// no lint rule, no model resolution and no promotion of template +/// diagnostics. The engine that admitted the run judged the workflow; a +/// graph Fabro's own parser cannot read is still refused, since the read side +/// needs one. +pub fn compile_admitted_run(input: CreateRunCompileInput) -> Result { + let CreateRunCompileInput { + workflow, + settings, + vars, + cwd, + workflow_path, + workflow_bundle, + configured_providers, + } = input; + let resolved = resolve_workflow(ResolveWorkflowInput { + workflow, + settings, + cwd, + }) + .map_err(|err| Error::Parse(err.to_string()))?; + let settings = resolved.settings; + let labels = settings.combined_labels(); + let definition = match (workflow_path, workflow_bundle) { + (Some(workflow_path), Some(workflow_bundle)) => { + Some(RunDefinition::new(workflow_path, workflow_bundle)) + } + _ => None, + }; + let mut parsed = pipeline::parse(&resolved.raw_source)?; + apply_goal_override(&mut parsed.graph, resolved.goal_override.as_deref()); + let transformed = pipeline::transform(parsed, &TransformOptions { + current_dir: resolved.current_dir.clone(), + file_resolver: resolved.file_resolver.clone(), + template_context: template_context(Some(&settings), vars), + source_name: resolved + .dot_path + .as_ref() + .map(|path| path.display().to_string()), + render_mode: RenderMode::Structural, + custom_transforms: Vec::new(), + model_resolution: None, + })?; + let validated = Validated::new( + transformed.graph, + transformed.source, + transformed.diagnostics, + ); + Ok(CompiledRun { + validated, + settings, + raw_source: resolved.raw_source, + workflow_slug: resolved.workflow_slug, + dot_path: resolved.dot_path, + definition, + source_directory: resolved.working_directory.to_string_lossy().to_string(), + labels, + configured_providers, + }) +} + +/// Stage three for a run another engine admitted: no model pinning, since +/// the engine pinned every route at its own admission. +#[must_use] +pub fn materialize_admitted_run(compiled: CompiledRun) -> MaterializedRun { + let CompiledRun { + validated, + settings, + raw_source, + workflow_slug, + dot_path, + definition, + source_directory, + labels, + configured_providers: _, + } = compiled; + MaterializedRun { + validated, + settings, + raw_source, + workflow_slug, + dot_path, + definition, + source_directory, + labels, + } +} + /// Resolve, preprocess, validate, and promote a workflow for run creation. /// /// This stage is synchronous and may read workflow files. Async callers must diff --git a/lib/components/fabro-workflow/src/operations/mod.rs b/lib/components/fabro-workflow/src/operations/mod.rs index 57ed472f4..5de6be333 100644 --- a/lib/components/fabro-workflow/src/operations/mod.rs +++ b/lib/components/fabro-workflow/src/operations/mod.rs @@ -17,8 +17,8 @@ pub use archive::{ pub use create::{ CompiledRun, CreateRunCompileInput, CreateRunInput, CreateRunPersistenceInput, CreateRunPersistenceMetadata, CreatedRun, MaterializedRun, - assemble_create_run_persistence_input, compile_create_run, create, make_run_dir, - materialize_create_run, persist_create_run, + assemble_create_run_persistence_input, compile_admitted_run, compile_create_run, create, + make_run_dir, materialize_admitted_run, materialize_create_run, persist_create_run, }; pub use fork::{ForkOutcome, ForkRunInput, ResolvedForkTarget, fork_run}; pub use resume::resume;