mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-02 02:13:49 +00:00
Run a workflow through Petri in the server, behind the engine flag
When a run's engine is Petri, the create handler hands the bundle, inputs and launch to Petri's check instead of the legacy compile, lint and model pinning, refuses the run with the validation error the legacy validator uses (Petri's codes as the rules, listed in the API detail), and records the admission on the run spec. The Fabro graph the read side displays is parsed without validation. The scheduler executes a Petri run in the server process through fabro_petri::engine, appending only the run lifecycle events the read side needs (run.starting, run.running, run.completed or run.failed); no stage or agent event is projected yet. Scenario tests run the hello bundle on the OpenAI twin under the version flag and a command-only bundle under the server setting, check Petri's record agrees, and cover the refusals for an unknown attribute, an undeclared node and an unknown model. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
03309d4240
commit
d2f0e70dca
10 changed files with 1077 additions and 21 deletions
|
|
@ -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" }
|
||||
|
|
|
|||
|
|
@ -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<PinnedRun> {
|
||||
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<PinnedRun
|
|||
Ok(PinnedRun {
|
||||
materialized,
|
||||
metadata,
|
||||
engine: RunEngine::Legacy,
|
||||
})
|
||||
}
|
||||
|
||||
|
|
@ -445,6 +508,7 @@ pub(crate) fn assemble_run(pinned: PinnedRun) -> 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,
|
||||
})
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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<dyn CredentialProvider>,
|
||||
/// 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<Arc<RunLayer>>,
|
||||
manifest_run_settings: RwLock<std::result::Result<RunNamespace, SharedError>>,
|
||||
pub(crate) server_settings: RwLock<Arc<ServerSettings>>,
|
||||
|
|
@ -2430,6 +2434,7 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result<Arc<AppS
|
|||
automation_materializer_override,
|
||||
} = config;
|
||||
|
||||
let store_pool = db_pool.clone();
|
||||
let automation_migration_pool = db_pool.clone();
|
||||
load_store_blocking("automation environment migration", move || async move {
|
||||
fabro_automation::backfill_environment_selectors(&automation_migration_pool)
|
||||
|
|
@ -2595,6 +2600,7 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result<Arc<AppS
|
|||
parent_link_lock: AsyncMutex::new(()),
|
||||
server_secrets,
|
||||
llm_source,
|
||||
db_pool: store_pool,
|
||||
manifest_run_defaults: RwLock::new(current_manifest_run_defaults),
|
||||
manifest_run_settings: RwLock::new(current_manifest_run_settings),
|
||||
server_settings: RwLock::new(current_server_settings),
|
||||
|
|
@ -3957,6 +3963,25 @@ async fn execute_run(state: Arc<AppState>, 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<AppState>, 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<Engine> {
|
||||
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<AppState>, run_id: RunId) {
|
||||
// Transition to Starting and set up cancel infrastructure
|
||||
let (cancel_rx, run_dir, event_tx, cancel_token, execution_mode) = {
|
||||
|
|
|
|||
|
|
@ -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::<Vec<_>>();
|
||||
if listed.is_empty() {
|
||||
error.to_string()
|
||||
} else {
|
||||
format!("{error}: {}", listed.join("; "))
|
||||
}
|
||||
}
|
||||
|
||||
async fn validate_intent_actor_target(
|
||||
state: &AppState,
|
||||
actor: &Principal,
|
||||
|
|
|
|||
440
lib/apps/fabro-server/src/server/petri_runs.rs
Normal file
440
lib/apps/fabro-server/src/server/petri_runs.rs
Normal file
|
|
@ -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<String> {
|
||||
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<PetriAdmission, RunCompilerError> {
|
||||
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::<HashSet<_>>();
|
||||
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<AppState>, 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<String>, 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<AppState>,
|
||||
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<AppState>, run_id: RunId, status: RunStatus, error: Option<String>) {
|
||||
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();
|
||||
}
|
||||
|
|
@ -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");
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
mod archive;
|
||||
mod dry_run;
|
||||
mod lifecycle;
|
||||
mod petri;
|
||||
mod run_completion;
|
||||
mod sse;
|
||||
mod usage;
|
||||
|
|
|
|||
382
lib/apps/fabro-server/tests/it/scenario/petri.rs
Normal file
382
lib/apps/fabro-server/tests/it/scenario/petri.rs
Normal file
|
|
@ -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<PathBuf> {
|
||||
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::<BTreeMap<_, _>>();
|
||||
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}"
|
||||
);
|
||||
}
|
||||
|
|
@ -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<CompiledRun, Error> {
|
||||
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
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue