From ac6e3ced6adf252f7aea435d5286c9853d8a4ea8 Mon Sep 17 00:00:00 2001 From: Fabro Date: Wed, 29 Jul 2026 22:27:05 +0000 Subject: [PATCH] fabro(01KYQN78K19NY7PNSCDYP6CG9G): implement (succeeded) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Fabro-Run: 01KYQN78K19NY7PNSCDYP6CG9G Fabro-Completed: 5 Fabro-Checkpoint: 879008d9d587ff62ea2633980529d51587e11dbc ⚒️ Generated with [Fabro](https://fabro.sh) --- lib/apps/fabro-server/src/lib.rs | 1 + lib/apps/fabro-server/src/run_compiler.rs | 1105 +++++++++++++++++ lib/apps/fabro-server/src/run_manifest.rs | 56 +- .../fabro-server/src/server/handler/runs.rs | 255 +++- lib/apps/fabro-server/src/server/tests.rs | 243 ++++ .../fabro-workflow/src/operations/create.rs | 835 ++++++++++--- .../fabro-workflow/src/operations/mod.rs | 7 +- 7 files changed, 2257 insertions(+), 245 deletions(-) create mode 100644 lib/apps/fabro-server/src/run_compiler.rs diff --git a/lib/apps/fabro-server/src/lib.rs b/lib/apps/fabro-server/src/lib.rs index 5fc735ba0..9a7e12978 100644 --- a/lib/apps/fabro-server/src/lib.rs +++ b/lib/apps/fabro-server/src/lib.rs @@ -34,6 +34,7 @@ pub mod manifest_validation; mod migrations; mod principal_middleware; mod request_id; +mod run_compiler; mod run_files; mod run_files_security; mod run_manifest; diff --git a/lib/apps/fabro-server/src/run_compiler.rs b/lib/apps/fabro-server/src/run_compiler.rs new file mode 100644 index 000000000..341c4565c --- /dev/null +++ b/lib/apps/fabro-server/src/run_compiler.rs @@ -0,0 +1,1105 @@ +use std::collections::HashMap; +use std::error::Error as StdError; +use std::path::PathBuf; +use std::sync::Arc; + +use fabro_config::parse::{self, SettingsSource}; +use fabro_config::{ + CliLayer, EnvironmentDockerfileLayer, EnvironmentImageLayer, EnvironmentLayer, MergeMap, + RunLayer, SettingsLayer, WorkflowSettingsBuilder, +}; +use fabro_model::{Catalog, ModelSelectionError, ProviderId}; +use fabro_types::settings::AmbiguousModelRef; +use fabro_types::settings::interp::{InterpString, ResolveError}; +use fabro_types::settings::run::{McpServerSettings, RunGoal}; +use fabro_types::{ + AutomationRef, GitContext, ManifestPath, RunId, RunProvenance, WorkflowSettings, +}; +use fabro_util::workspace_glob::{WorkspaceGlob, WorkspaceGlobError}; +use fabro_workflow::Error as WorkflowError; +use fabro_workflow::operations::{ + self, CompiledRun, CreateRunCompileInput, CreateRunPersistenceInput, + CreateRunPersistenceMetadata, MaterializedRun, WorkflowInput, +}; +use fabro_workflow::workflow_bundle::{BundledWorkflow, WorkflowBundle}; +use tokio::task; + +/// One project settings source already normalized into the manifest path +/// namespace used by the workflow bundle. +#[derive(Clone, Debug)] +pub(crate) struct ProjectSettingsSource { + pub(crate) path: ManifestPath, + pub(crate) toml: String, +} + +/// Transport-neutral inputs for compiling one submitted run. +/// +/// IDs, title, git metadata, and provenance are resolved by the caller. This +/// boundary owns only source normalization, settings resolution, workflow +/// compilation, materialization, and persistence-input assembly. +#[derive(Clone, Debug)] +pub(crate) struct RawRunCompilerInput { + pub(crate) workflow_bundle: WorkflowBundle, + pub(crate) entrypoint: ManifestPath, + pub(crate) cwd: PathBuf, + pub(crate) server_run_defaults: RunLayer, + pub(crate) server_environment_defaults: MergeMap, + pub(crate) server_mcp_catalog: HashMap, + pub(crate) project_settings: Vec, + pub(crate) user_toml: Vec, + pub(crate) run_overrides: Option, + pub(crate) cli_overrides: Option, + pub(crate) input_overrides: HashMap, + pub(crate) inline_goal_override: Option, + pub(crate) vars: HashMap, + pub(crate) run_id: Option, + pub(crate) title: Option, + pub(crate) parent_id: Option, + pub(crate) git: Option, + pub(crate) storage_root: PathBuf, + pub(crate) configured_providers: Vec, + pub(crate) workflow_slug: Option, + pub(crate) provenance: RunProvenance, + pub(crate) web_url: Option, + pub(crate) submitted_manifest_bytes: Option>, + pub(crate) automation: Option, +} + +#[derive(Clone, Debug)] +struct RunMetadata { + run_id: Option, + title: Option, + parent_id: Option, + git: Option, + storage_root: PathBuf, + workflow_slug: Option, + provenance: RunProvenance, + web_url: Option, + submitted_manifest_bytes: Option>, + automation: Option, +} + +/// Stage-one output: the selected bundled workflow and all client settings +/// sources have been parsed and normalized, but no settings have been layered. +pub(crate) struct NormalizedRun { + workflow_bundle: WorkflowBundle, + entrypoint: ManifestPath, + workflow: BundledWorkflow, + workflow_layer: Option, + project_layers: Vec, + user_toml: Vec, + cwd: PathBuf, + server_run_defaults: RunLayer, + server_environment_defaults: MergeMap, + server_mcp_catalog: HashMap, + run_overrides: Option, + cli_overrides: Option, + input_overrides: HashMap, + inline_goal_override: Option, + vars: HashMap, + configured_providers: Vec, + metadata: RunMetadata, +} + +/// Settings-layered output. Variable substitution is intentionally separate +/// so callers can snapshot variables after source/settings preparation, as the +/// create handler historically does. +pub(crate) struct LayeredRun { + workflow_bundle: WorkflowBundle, + entrypoint: ManifestPath, + workflow: BundledWorkflow, + settings: WorkflowSettings, + cwd: PathBuf, + vars: HashMap, + configured_providers: Vec, + metadata: RunMetadata, +} + +impl LayeredRun { + pub(crate) fn with_vars(mut self, vars: HashMap) -> Self { + self.vars = vars; + self + } +} + +/// Settings-resolved stage output. Callers may inspect this before policy +/// checks, then move it into [`compile_graph`] after those checks pass. +pub(crate) struct PreparedRun { + workflow_bundle: WorkflowBundle, + entrypoint: ManifestPath, + workflow: BundledWorkflow, + settings: WorkflowSettings, + cwd: PathBuf, + vars: HashMap, + configured_providers: Vec, + metadata: RunMetadata, +} + +impl PreparedRun { + pub(crate) fn resolve_run_id(mut self) -> (Self, RunId) { + let run_id = self.metadata.run_id.unwrap_or_default(); + self.metadata.run_id = Some(run_id); + (self, run_id) + } + + pub(crate) fn with_web_url(mut self, web_url: Option) -> Self { + self.metadata.web_url = web_url; + self + } + + pub(crate) fn with_configured_providers( + mut self, + configured_providers: Vec, + ) -> Self { + self.configured_providers = configured_providers; + self + } + + pub(crate) fn settings(&self) -> &WorkflowSettings { + &self.settings + } + + pub(crate) fn parent_id(&self) -> Option { + self.metadata.parent_id + } +} + +/// Graph-compiled stage output, retaining the metadata needed by later pure +/// assembly. +pub(crate) struct GraphCompiledRun { + compiled: CompiledRun, + entrypoint: ManifestPath, + metadata: RunMetadata, +} + +impl GraphCompiledRun { + #[cfg(test)] + pub(crate) fn compiled(&self) -> &CompiledRun { + &self.compiled + } +} + +/// Materialized stage output ready for pure persistence-input assembly. +pub(crate) struct PersistenceReadyRun { + materialized: MaterializedRun, + entrypoint: ManifestPath, + metadata: RunMetadata, +} + +/// Complete output of this boundary. +pub(crate) struct RunCompilerOutput { + persistence_input: CreateRunPersistenceInput, + entrypoint: ManifestPath, +} + +impl RunCompilerOutput { + #[cfg(test)] + pub(crate) fn persistence_input(&self) -> &CreateRunPersistenceInput { + &self.persistence_input + } + + #[cfg(test)] + pub(crate) fn entrypoint(&self) -> &ManifestPath { + &self.entrypoint + } + + pub(crate) fn into_parts(self) -> (CreateRunPersistenceInput, ManifestPath) { + (self.persistence_input, self.entrypoint) + } +} + +#[derive(Debug, thiserror::Error)] +pub(crate) enum RunCompilerError { + #[error("invalid run source: {source}")] + InvalidSource { + #[source] + source: InvalidSourceError, + }, + + #[error("invalid run settings: {source}")] + InvalidSettings { + #[source] + source: Box, + }, + + #[error("run config variable interpolation failed: {source}")] + VariableInterpolation { + #[source] + source: VariableInterpolationError, + }, + + #[error("workflow validation or parse failed: {source}")] + ValidationOrParse { + #[source] + source: WorkflowError, + }, + + #[error("model selection failed: {source}")] + ModelSelection { + #[source] + source: ModelSelectionError, + }, + + #[error("model reference failed: {source}")] + ModelReference { + #[source] + source: AmbiguousModelRef, + }, + + #[error("{context}")] + Internal { + context: &'static str, + #[source] + source: Box, + }, +} + +#[derive(Debug, thiserror::Error)] +pub(crate) enum InvalidSourceError { + #[error("bundle entrypoint {entrypoint} is missing from the workflow bundle")] + MissingEntrypoint { entrypoint: ManifestPath }, + + #[error("unsupported dockerfile reference {reference:?} in {config_path}")] + UnsupportedDockerfileReference { + config_path: ManifestPath, + reference: String, + }, + + #[error("bundled dockerfile {dockerfile_path} referenced by {config_path} is missing")] + MissingDockerfile { + config_path: ManifestPath, + dockerfile_path: ManifestPath, + }, +} + +#[derive(Debug, thiserror::Error)] +pub(crate) enum InvalidSettingsError { + #[error("failed to parse {kind} settings at {path}")] + Parse { + kind: &'static str, + path: ManifestPath, + #[source] + source: Box, + }, + + #[error("failed to parse user settings: {source}")] + User { + #[source] + source: fabro_config::Error, + }, + + #[error("failed to resolve layered workflow settings: {source}")] + Resolve { + #[source] + source: fabro_config::ResolveErrors, + }, +} + +#[derive(Debug, thiserror::Error)] +pub(crate) enum VariableInterpolationError { + #[error(transparent)] + Interpolation(#[from] ResolveError), + + #[error("run.artifacts.include[{index}]: {source}")] + ArtifactGlob { + index: usize, + #[source] + source: WorkspaceGlobError, + }, +} + +pub(crate) type Result = std::result::Result; + +/// Normalize the bundle entrypoint and parse workflow/project settings while +/// resolving Dockerfile references against the selected workflow's files. +pub(crate) fn normalize_source(input: RawRunCompilerInput) -> Result { + let RawRunCompilerInput { + workflow_bundle, + entrypoint, + cwd, + server_run_defaults, + server_environment_defaults, + server_mcp_catalog, + project_settings, + user_toml, + run_overrides, + cli_overrides, + input_overrides, + inline_goal_override, + vars, + run_id, + title, + parent_id, + git, + storage_root, + configured_providers, + workflow_slug, + provenance, + web_url, + submitted_manifest_bytes, + automation, + } = input; + let mut workflow = workflow_bundle + .workflow(&entrypoint) + .cloned() + .ok_or_else(|| RunCompilerError::InvalidSource { + source: InvalidSourceError::MissingEntrypoint { + entrypoint: entrypoint.clone(), + }, + })?; + workflow.path = entrypoint.clone(); + + let workflow_layer = workflow + .config + .as_ref() + .map(|config| { + parse_settings_layer( + &config.source, + &config.path, + &workflow.files, + SettingsSource::Workflow, + "workflow", + ) + }) + .transpose()?; + let project_layers = project_settings + .into_iter() + .map(|project| { + parse_settings_layer( + &project.toml, + &project.path, + &workflow.files, + SettingsSource::Project, + "project", + ) + }) + .collect::>>()?; + + Ok(NormalizedRun { + workflow_bundle, + entrypoint, + workflow, + workflow_layer, + project_layers, + user_toml, + cwd, + server_run_defaults, + server_environment_defaults, + server_mcp_catalog, + run_overrides, + cli_overrides, + input_overrides, + inline_goal_override, + vars, + configured_providers, + metadata: RunMetadata { + run_id, + title, + parent_id, + git, + storage_root, + workflow_slug, + provenance, + web_url, + submitted_manifest_bytes, + automation, + }, + }) +} + +/// Layer settings and apply the submitted input and goal overrides. +pub(crate) fn layer_settings(normalized: NormalizedRun) -> Result { + let NormalizedRun { + workflow_bundle, + entrypoint, + workflow, + workflow_layer, + project_layers, + user_toml, + cwd, + server_run_defaults, + server_environment_defaults, + server_mcp_catalog, + run_overrides, + cli_overrides, + input_overrides, + inline_goal_override, + vars, + configured_providers, + metadata, + } = normalized; + let mut builder = WorkflowSettingsBuilder::new() + .server_manifest_defaults(server_run_defaults, server_environment_defaults) + .server_mcp_catalog(server_mcp_catalog); + if let Some(run) = run_overrides { + builder = builder.run_overrides(run); + } + if let Some(cli) = cli_overrides { + builder = builder.cli_overrides(cli); + } + if let Some(layer) = workflow_layer { + builder = builder.workflow_layer(layer); + } + for layer in project_layers { + builder = builder.project_layer(layer); + } + for source in user_toml { + builder = + builder + .user_toml(&source) + .map_err(|source| RunCompilerError::InvalidSettings { + source: Box::new(InvalidSettingsError::User { source }), + })?; + } + let mut settings = builder + .build() + .map_err(|source| RunCompilerError::InvalidSettings { + source: Box::new(InvalidSettingsError::Resolve { source }), + })?; + settings.run.inputs.extend(input_overrides); + if let Some(goal) = inline_goal_override { + settings.run.goal = Some(RunGoal::Inline(InterpString::parse(&goal))); + } + + Ok(LayeredRun { + workflow_bundle, + entrypoint, + workflow, + settings, + cwd, + vars, + configured_providers, + metadata, + }) +} + +/// Apply a run-variable snapshot and validate the resulting artifact globs. +pub(crate) fn apply_run_variables(mut layered: LayeredRun) -> Result { + substitute_variables(&layered.vars, &mut layered.settings)?; + Ok(PreparedRun { + workflow_bundle: layered.workflow_bundle, + entrypoint: layered.entrypoint, + workflow: layered.workflow, + settings: layered.settings, + cwd: layered.cwd, + vars: layered.vars, + configured_providers: layered.configured_providers, + metadata: layered.metadata, + }) +} + +/// Run stage one and the settings/variables portion of stage two. This is the +/// convenient boundary for callers that already own a variable snapshot. +#[cfg(test)] +pub(crate) fn prepare_run(input: RawRunCompilerInput) -> Result { + apply_run_variables(layer_settings(normalize_source(input)?)?) +} + +/// Compile and validate the graph on Tokio's blocking pool. +pub(crate) async fn compile_graph( + prepared: PreparedRun, + catalog: Arc, +) -> Result { + let PreparedRun { + workflow_bundle, + entrypoint, + workflow, + settings, + cwd, + vars, + configured_providers, + metadata, + } = prepared; + let compile_input = CreateRunCompileInput { + workflow: WorkflowInput::Bundled(workflow), + settings, + vars, + cwd, + workflow_path: Some(entrypoint.clone()), + workflow_bundle: Some(workflow_bundle), + configured_providers, + }; + let compiled = + task::spawn_blocking(move || operations::compile_create_run(compile_input, catalog)) + .await + .map_err(|source| RunCompilerError::Internal { + context: "workflow compilation failed", + source: Box::new(WorkflowError::engine_with_source( + "workflow create task failed", + source, + )), + })? + .map_err(classify_workflow_error)?; + + Ok(GraphCompiledRun { + compiled, + entrypoint, + metadata, + }) +} + +/// Materialize run-level model settings on Tokio's blocking pool. +pub(crate) async fn materialize_run( + compiled: GraphCompiledRun, + catalog: Arc, +) -> Result { + let GraphCompiledRun { + compiled, + entrypoint, + metadata, + } = compiled; + let materialized = task::spawn_blocking(move || { + operations::materialize_create_run(compiled, catalog.as_ref()) + }) + .await + .map_err(|source| RunCompilerError::Internal { + context: "workflow compilation failed", + source: Box::new(WorkflowError::engine_with_source( + "workflow create task failed", + source, + )), + })? + .map_err(classify_workflow_error)?; + + Ok(PersistenceReadyRun { + materialized, + entrypoint, + metadata, + }) +} + +/// Purely assemble the complete workflow persistence input. +pub(crate) fn assemble_run(ready: PersistenceReadyRun) -> RunCompilerOutput { + let PersistenceReadyRun { + materialized, + entrypoint, + metadata, + } = ready; + let RunMetadata { + run_id, + title, + parent_id, + git, + storage_root, + workflow_slug, + provenance, + web_url, + submitted_manifest_bytes, + automation, + } = metadata; + let persistence_input = operations::assemble_create_run_persistence_input( + materialized, + CreateRunPersistenceMetadata { + run_id: run_id.unwrap_or_default(), + storage_root, + workflow_slug, + submitted_manifest_bytes, + title, + automation, + git, + fork_source_ref: None, + parent_id, + provenance, + web_url, + }, + ); + + RunCompilerOutput { + persistence_input, + entrypoint, + } +} + +/// Compile a raw run all the way to a complete persistence input. +#[cfg(test)] +pub(crate) async fn compile_run( + input: RawRunCompilerInput, + catalog: Arc, +) -> Result { + let prepared = prepare_run(input)?; + let compiled = compile_graph(prepared, Arc::clone(&catalog)).await?; + let materialized = materialize_run(compiled, catalog).await?; + Ok(assemble_run(materialized)) +} + +fn parse_settings_layer( + source: &str, + config_path: &ManifestPath, + files: &HashMap, + settings_source: SettingsSource, + kind: &'static str, +) -> Result { + let mut layer = + source + .parse::() + .map_err(|source| RunCompilerError::InvalidSettings { + source: Box::new(InvalidSettingsError::Parse { + kind, + path: config_path.clone(), + source: Box::new(source), + }), + })?; + parse::validate_settings_source(&layer, settings_source).map_err(|source| { + RunCompilerError::InvalidSettings { + source: Box::new(InvalidSettingsError::Parse { + kind, + path: config_path.clone(), + source: Box::new(source), + }), + } + })?; + resolve_dockerfiles(&mut layer, config_path, files)?; + Ok(layer) +} + +fn resolve_dockerfiles( + layer: &mut SettingsLayer, + config_path: &ManifestPath, + files: &HashMap, +) -> Result<()> { + for environment in layer.environments.values_mut() { + if let Some(image) = environment.image.as_mut() { + resolve_dockerfile(image, config_path, files)?; + } + } + if let Some(image) = layer + .run + .as_mut() + .and_then(|run| run.environment.as_mut()) + .and_then(|environment| environment.image.as_mut()) + { + resolve_dockerfile(image, config_path, files)?; + } + Ok(()) +} + +fn resolve_dockerfile( + image: &mut EnvironmentImageLayer, + config_path: &ManifestPath, + files: &HashMap, +) -> Result<()> { + let Some(EnvironmentDockerfileLayer::Path { path }) = image.dockerfile.as_ref() else { + return Ok(()); + }; + let reference = path.clone(); + let dockerfile_path = ManifestPath::from_reference(config_path.parent_or_dot(), &reference) + .ok_or_else(|| RunCompilerError::InvalidSource { + source: InvalidSourceError::UnsupportedDockerfileReference { + config_path: config_path.clone(), + reference: reference.clone(), + }, + })?; + let content = + files + .get(&dockerfile_path) + .cloned() + .ok_or_else(|| RunCompilerError::InvalidSource { + source: InvalidSourceError::MissingDockerfile { + config_path: config_path.clone(), + dockerfile_path: dockerfile_path.clone(), + }, + })?; + image.dockerfile = Some(EnvironmentDockerfileLayer::Inline(content)); + Ok(()) +} + +fn substitute_variables( + variables: &HashMap, + settings: &mut WorkflowSettings, +) -> Result<()> { + settings + .run + .substitute_variables(|name| variables.get(name).cloned()) + .map_err(|source| RunCompilerError::VariableInterpolation { + source: VariableInterpolationError::Interpolation(source), + })?; + for (index, pattern) in settings.run.artifacts.include.iter().enumerate() { + WorkspaceGlob::try_new(pattern).map_err(|source| { + RunCompilerError::VariableInterpolation { + source: VariableInterpolationError::ArtifactGlob { index, source }, + } + })?; + } + Ok(()) +} + +fn classify_workflow_error(error: WorkflowError) -> RunCompilerError { + match error { + WorkflowError::ModelSelection(source) => RunCompilerError::ModelSelection { source }, + WorkflowError::ModelReference(source) => RunCompilerError::ModelReference { source }, + source @ (WorkflowError::Parse(_) | WorkflowError::ValidationFailed { .. }) => { + RunCompilerError::ValidationOrParse { source } + } + source => RunCompilerError::Internal { + context: "workflow compilation failed", + source: Box::new(source), + }, + } +} + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + use std::error::Error as _; + use std::sync::Arc; + + use fabro_config::EnvironmentDockerfileLayer; + use fabro_graphviz::graph::AttrValue; + use fabro_model::Catalog; + use fabro_types::settings::interp::ResolveCtx; + use fabro_types::settings::run::RunGoal; + use fabro_types::{AutomationRef, Principal, RunProvenance, SystemActorKind}; + use fabro_workflow::workflow_bundle::ParsedWorkflowConfig; + + use super::*; + + const DOT: &str = r#"digraph Test { + graph [goal="Graph goal"] + start [shape=Mdiamond] + work [prompt="Ship {{ inputs.target }} for {{ vars.owner }}", model="gpt-5.4"] + exit [shape=Msquare] + start -> work -> exit + }"#; + + fn manifest_path(value: &str) -> ManifestPath { + ManifestPath::from_wire(value).expect("fixture manifest path should be valid") + } + + fn provenance() -> RunProvenance { + RunProvenance { + server: None, + client: None, + subject: Principal::System { + system_kind: SystemActorKind::Engine, + }, + } + } + + fn workflow( + entrypoint: &ManifestPath, + workflow_toml: Option<&str>, + files: HashMap, + ) -> BundledWorkflow { + BundledWorkflow { + path: entrypoint.clone(), + source: DOT.to_string(), + config: workflow_toml.map(|source| ParsedWorkflowConfig { + path: manifest_path("flows/workflow.toml"), + source: source.to_string(), + }), + files, + } + } + + fn raw_input( + workflow_toml: Option<&str>, + files: HashMap, + ) -> RawRunCompilerInput { + let entrypoint = manifest_path("flows/workflow.fabro"); + let workflow = workflow(&entrypoint, workflow_toml, files); + RawRunCompilerInput { + workflow_bundle: WorkflowBundle::new(HashMap::from([(entrypoint.clone(), workflow)])), + entrypoint, + cwd: PathBuf::from("/workspace"), + server_run_defaults: RunLayer::default(), + server_environment_defaults: fabro_environment::seeded_catalog_layer(), + server_mcp_catalog: HashMap::new(), + project_settings: Vec::new(), + user_toml: Vec::new(), + run_overrides: None, + cli_overrides: None, + input_overrides: HashMap::new(), + inline_goal_override: None, + vars: HashMap::new(), + run_id: Some(RunId::new()), + title: None, + parent_id: None, + git: None, + storage_root: PathBuf::from("/tmp/fabro-storage"), + configured_providers: Catalog::builtin().all_provider_ids().into_iter().collect(), + workflow_slug: None, + provenance: provenance(), + web_url: None, + submitted_manifest_bytes: None, + automation: None, + } + } + + #[test] + fn stage_one_rejects_missing_entrypoint() { + let mut input = raw_input(None, HashMap::new()); + input.entrypoint = manifest_path("flows/missing.fabro"); + + let Err(error) = normalize_source(input) else { + panic!("missing entrypoint should fail"); + }; + + assert!(matches!(error, RunCompilerError::InvalidSource { + source: InvalidSourceError::MissingEntrypoint { .. }, + })); + } + + #[test] + fn unresolved_run_id_is_allocated_only_after_variables_are_applied() { + let mut input = raw_input(None, HashMap::new()); + input.run_id = None; + let normalized = normalize_source(input).expect("source should normalize"); + let layered = layer_settings(normalized).expect("settings should layer"); + let prepared = apply_run_variables(layered).expect("variables should apply"); + + let (prepared, run_id) = prepared.resolve_run_id(); + + assert_eq!(prepared.metadata.run_id, Some(run_id)); + } + + #[test] + fn stage_one_rejects_missing_dockerfile_and_preserves_source_chain() { + let workflow_toml = r#" +_version = 1 + +[run.environment.image] +dockerfile = { path = "Dockerfile" } +"#; + + let Err(error) = normalize_source(raw_input(Some(workflow_toml), HashMap::new())) else { + panic!("missing dockerfile should fail"); + }; + + assert!(matches!(error, RunCompilerError::InvalidSource { + source: InvalidSourceError::MissingDockerfile { .. }, + })); + let source = error + .source() + .expect("top-level error should retain source"); + assert!(source.to_string().contains("Dockerfile")); + } + + #[test] + fn stage_one_resolves_bundled_dockerfile() { + let workflow_toml = r#" +_version = 1 + +[run.environment.image] +dockerfile = { path = "Dockerfile" } +"#; + let normalized = normalize_source(raw_input( + Some(workflow_toml), + HashMap::from([( + manifest_path("flows/Dockerfile"), + "FROM ubuntu:24.04\n".to_string(), + )]), + )) + .expect("bundled dockerfile should resolve"); + let dockerfile = normalized + .workflow_layer + .as_ref() + .and_then(|layer| layer.run.as_ref()) + .and_then(|run| run.environment.as_ref()) + .and_then(|environment| environment.image.as_ref()) + .and_then(|image| image.dockerfile.as_ref()); + + assert_eq!( + dockerfile, + Some(&EnvironmentDockerfileLayer::Inline( + "FROM ubuntu:24.04\n".to_string() + )) + ); + } + + #[test] + fn settings_apply_precedence_vars_inputs_and_safe_artifact_globs() { + let workflow_toml = r#" +_version = 1 + +[run.metadata] +layer = "workflow" +owner = "{{ vars.owner }}" + +[run.inputs] +target = "workflow" + +[run.artifacts] +include = ["reports/{{ vars.owner }}/*.json"] +"#; + let mut input = raw_input(Some(workflow_toml), HashMap::new()); + input.project_settings.push(ProjectSettingsSource { + path: manifest_path(".fabro/project.toml"), + toml: r#" +_version = 1 + +[run.metadata] +layer = "project" +"# + .to_string(), + }); + input.user_toml = vec![ + r#" +_version = 1 + +[run.metadata] +layer = "user" +"# + .to_string(), + ]; + input.run_overrides = Some( + toml::from_str::( + r#" +_version = 1 + +[run.metadata] +layer = "args" +owner = "{{ vars.owner }}" +"#, + ) + .expect("args settings should parse") + .run + .expect("args run layer should exist"), + ); + input.input_overrides.insert( + "target".to_string(), + toml::Value::String("override".to_string()), + ); + input.inline_goal_override = Some("Ship {{ vars.owner }}".to_string()); + input + .vars + .insert("owner".to_string(), "payments".to_string()); + + let prepared = prepare_run(input).expect("settings should prepare"); + let settings = prepared.settings(); + + assert_eq!( + settings.run.metadata.get("layer").map(String::as_str), + Some("args") + ); + assert_eq!( + settings.run.metadata.get("owner").map(String::as_str), + Some("payments") + ); + assert_eq!( + settings.run.inputs.get("target"), + Some(&toml::Value::String("override".to_string())) + ); + assert_eq!(settings.run.artifacts.include, vec![ + "reports/payments/*.json" + ]); + let Some(RunGoal::Inline(goal)) = settings.run.goal.as_ref() else { + panic!("inline goal override should win"); + }; + assert_eq!( + goal.resolve_with(&mut ResolveCtx::default()).unwrap(), + "Ship payments" + ); + } + + #[test] + fn settings_reject_artifact_glob_made_unsafe_by_variable() { + let workflow_toml = r#" +_version = 1 + +[run.artifacts] +include = ["reports/{{ vars.path }}/*.json"] +"#; + let mut input = raw_input(Some(workflow_toml), HashMap::new()); + input + .vars + .insert("path".to_string(), "../secrets".to_string()); + + let Err(error) = prepare_run(input) else { + panic!("unsafe artifact glob should fail"); + }; + + assert!(matches!(error, RunCompilerError::VariableInterpolation { + source: VariableInterpolationError::ArtifactGlob { .. }, + })); + } + + #[tokio::test] + async fn graph_vars_are_hard_errors_and_successfully_render_when_present() { + let catalog = Arc::new(Catalog::from_builtin().unwrap()); + let missing = prepare_run(raw_input(None, HashMap::new())) + .expect("settings preparation should not compile graph vars"); + let Err(error) = compile_graph(missing, Arc::clone(&catalog)).await else { + panic!("missing graph variable should be a hard error"); + }; + assert!(matches!(error, RunCompilerError::ValidationOrParse { + source: WorkflowError::ValidationFailed { .. }, + })); + + let mut input = raw_input(None, HashMap::new()); + input + .vars + .insert("owner".to_string(), "payments".to_string()); + input.input_overrides.insert( + "target".to_string(), + toml::Value::String("checkout".to_string()), + ); + let compiled = compile_graph( + prepare_run(input).expect("settings should prepare"), + catalog, + ) + .await + .expect("graph variables should render"); + let work = &compiled.compiled().validated().graph().nodes["work"]; + + assert_eq!( + work.attrs.get("prompt").and_then(AttrValue::as_str), + Some("Ship checkout for payments") + ); + assert_eq!( + work.attrs.get("provider").and_then(AttrValue::as_str), + Some("openai") + ); + } + + #[tokio::test] + async fn assembly_retains_entrypoint_and_run_metadata() { + let run_id = RunId::new(); + let parent_id = RunId::new(); + let automation = AutomationRef { + id: "nightly".to_string(), + name: Some("Nightly".to_string()), + trigger_id: Some("schedule".to_string()), + }; + let submitted = b"submitted manifest".to_vec(); + let mut input = raw_input(None, HashMap::new()); + input.run_id = Some(run_id); + input.parent_id = Some(parent_id); + input.title = Some("Compiler boundary".to_string()); + input.workflow_slug = Some("compiler-boundary".to_string()); + input.web_url = Some(format!("https://fabro.test/runs/{run_id}")); + input.submitted_manifest_bytes = Some(submitted.clone()); + input.automation = Some(automation.clone()); + input + .vars + .insert("owner".to_string(), "payments".to_string()); + input.input_overrides.insert( + "target".to_string(), + toml::Value::String("checkout".to_string()), + ); + let expected_entrypoint = input.entrypoint.clone(); + + let output = compile_run(input, Arc::new(Catalog::from_builtin().unwrap())) + .await + .expect("run should compile"); + let persistence = output.persistence_input(); + + assert_eq!(output.entrypoint(), &expected_entrypoint); + assert_eq!(persistence.run_id(), run_id); + assert_eq!(persistence.workflow_slug(), Some("compiler-boundary")); + assert_eq!( + persistence.submitted_manifest_bytes(), + Some(submitted.as_slice()) + ); + assert_eq!(persistence.automation(), Some(&automation)); + assert_eq!( + persistence + .definition() + .map(|definition| &definition.workflow_path), + Some(&expected_entrypoint) + ); + assert_eq!( + persistence.materialized().settings().run.goal.as_ref(), + Some(&RunGoal::Inline(InterpString::parse("Graph goal"))) + ); + } +} diff --git a/lib/apps/fabro-server/src/run_manifest.rs b/lib/apps/fabro-server/src/run_manifest.rs index fd9a1f50c..707b63931 100644 --- a/lib/apps/fabro-server/src/run_manifest.rs +++ b/lib/apps/fabro-server/src/run_manifest.rs @@ -38,8 +38,7 @@ use fabro_validate::Severity; use fabro_workflow::Error as WorkflowError; use fabro_workflow::model_fallback::resolve_model_fallbacks; use fabro_workflow::operations::{ - CreateRunInput, ValidateInput, WorkflowInput, validate, validate_with_catalog, - validate_with_ready_providers, + ValidateInput, WorkflowInput, validate, validate_with_catalog, validate_with_ready_providers, }; use fabro_workflow::pipeline::Validated; use fabro_workflow::run_materialization::materialize_run_with_ready_providers; @@ -56,21 +55,34 @@ pub(crate) struct PreparedManifest { pub cwd: PathBuf, pub git: Option, pub root_source: String, + #[allow( + dead_code, + reason = "create now resolves identity in the run compiler adapter" + )] pub run_id: Option, + #[allow( + dead_code, + reason = "create now resolves lineage in the run compiler adapter" + )] pub parent_id: Option, + #[allow( + dead_code, + reason = "create now normalizes titles in the run compiler adapter" + )] pub title: Option, pub settings: WorkflowSettings, pub target_path: ManifestPath, + #[allow(dead_code, reason = "create now owns the bundle through run_compiler")] pub workflow_bundle: WorkflowBundle, pub workflow_input: BundledWorkflow, pub source_directory: PathBuf, } #[derive(Clone, Debug, Default)] -struct ManifestSettingsOverrides { - run: Option, - cli: Option, - input_overrides: HashMap, +pub(crate) struct ManifestSettingsOverrides { + pub(crate) run: Option, + pub(crate) cli: Option, + pub(crate) input_overrides: HashMap, } #[cfg(test)] @@ -235,34 +247,6 @@ fn manifest_validate_input( } } -pub(crate) fn create_run_input( - prepared: PreparedManifest, - configured_providers: Vec, - provenance: RunProvenance, - web_url: Option, - vars: HashMap, -) -> CreateRunInput { - CreateRunInput { - workflow: WorkflowInput::Bundled(prepared.workflow_input), - settings: prepared.settings, - vars, - cwd: prepared.cwd, - workflow_slug: None, - workflow_path: Some(prepared.target_path), - workflow_bundle: Some(prepared.workflow_bundle), - submitted_manifest_bytes: None, - run_id: prepared.run_id, - title: prepared.title, - automation: None, - git: prepared.git, - fork_source_ref: None, - parent_id: prepared.parent_id, - provenance, - configured_providers, - web_url, - } -} - pub(crate) async fn run_preflight( state: &AppState, prepared: &PreparedManifest, @@ -387,7 +371,7 @@ fn settings_layer_with_resolved_dockerfiles( Ok(layer) } -fn manifest_args_overrides( +pub(crate) fn manifest_args_overrides( args: Option<&types::ManifestArgs>, ) -> Result { let Some(args) = args else { @@ -468,7 +452,7 @@ fn resolve_manifest_dockerfile( Ok(()) } -fn manifest_project_config_path( +pub(crate) fn manifest_project_config_path( config: &types::ManifestConfig, cwd: &Path, ) -> Result { diff --git a/lib/apps/fabro-server/src/server/handler/runs.rs b/lib/apps/fabro-server/src/server/handler/runs.rs index 0327d0cda..21451adf8 100644 --- a/lib/apps/fabro-server/src/server/handler/runs.rs +++ b/lib/apps/fabro-server/src/server/handler/runs.rs @@ -1,5 +1,6 @@ use std::collections::{HashMap, HashSet}; use std::io::ErrorKind; +use std::path::PathBuf; use std::sync::Arc; use axum::extract::{Path, Query, State}; @@ -13,7 +14,8 @@ use base64::engine::general_purpose::STANDARD as BASE64_STANDARD; use bytes::Bytes; use chrono::{DateTime, Utc}; use fabro_api::types::{ - BoardColumn, RunManifest, SubmitAnswerRequest, UpdateRunParentRequest, UpdateRunRequest, + BoardColumn, ManifestGoalType, RunManifest, SubmitAnswerRequest, UpdateRunParentRequest, + UpdateRunRequest, }; use fabro_config::Storage; use fabro_interview::AnswerSubmission; @@ -23,8 +25,8 @@ use fabro_store::{ }; use fabro_types::settings::ResolveError; use fabro_types::{ - AutomationRef, Principal, Run, RunClientProvenance, RunId, RunProvenance, RunServerProvenance, - RunStatusKind, StageContextWindow, StageContextWindowStaleness, + AutomationRef, ManifestPath, Principal, Run, RunClientProvenance, RunId, RunProvenance, + RunServerProvenance, RunStatusKind, StageContextWindow, StageContextWindowStaleness, StageContextWindowUnavailableReason, StageHandler, StageModelUsage, StageProjection, SystemActorKind, WorkflowSettings, parse_blob_ref, }; @@ -32,6 +34,7 @@ use fabro_util::version::FABRO_VERSION; use fabro_util::workspace_glob::{WorkspaceGlob, WorkspaceGlobError}; use fabro_workflow::command_log::{command_log_path, read_json_string_blob, read_log_slice}; use fabro_workflow::run_status::RunStatus; +use fabro_workflow::workflow_bundle::WorkflowBundle; use fabro_workflow::{Error as WorkflowError, operations}; use strum::VariantArray as _; use tokio::fs; @@ -48,6 +51,7 @@ use crate::principal_middleware::{ RequireCommandLog, RequireRunManagementTarget, RequireRunScoped, RequireRunStageScoped, RequiredRunManagementActor, RequiredUser, }; +use crate::run_compiler::{self, ProjectSettingsSource, RawRunCompilerInput, RunCompilerError}; use crate::run_files::{list_run_commits, list_run_files}; use crate::run_manifest; use crate::run_selector::{ResolveRunError, resolve_run_by_selector}; @@ -552,6 +556,151 @@ pub(crate) struct CreateRunFromManifestRequest { pub(crate) automation: Option, } +struct ManifestRunCompilerAdapter { + workflow_bundle: WorkflowBundle, + entrypoint: ManifestPath, + cwd: PathBuf, + project_settings: Vec, + user_toml: Vec, + run_overrides: Option, + cli_overrides: Option, + input_overrides: HashMap, + inline_goal_override: Option, + run_id: Option, + parent_id: Option, + title: Option, + git: Option, +} + +fn adapt_manifest_for_run_compiler( + manifest: &RunManifest, + explicit_run_id: Option, +) -> anyhow::Result { + use anyhow::{Context as _, anyhow, bail}; + use fabro_api::types::ManifestConfigType; + + if manifest.version != 1 { + bail!("unsupported manifest version {}", manifest.version); + } + let cwd = PathBuf::from(&manifest.cwd); + let entrypoint = ManifestPath::from_wire(&manifest.target.path) + .ok_or_else(|| anyhow!("invalid manifest target path: {}", manifest.target.path))?; + let workflow_bundle = run_manifest::workflow_bundle_from_manifest(&manifest.workflows)?; + if workflow_bundle.workflow(&entrypoint).is_none() { + return Err(anyhow!( + "manifest target path is missing from workflows map" + )); + } + let overrides = run_manifest::manifest_args_overrides(manifest.args.as_ref()) + .context("failed to parse manifest args")?; + let project_settings = manifest + .configs + .iter() + .filter(|config| config.type_ == ManifestConfigType::Project) + .filter_map(|config| config.source.as_ref().map(|source| (config, source))) + .map(|(config, source)| { + Ok(ProjectSettingsSource { + path: run_manifest::manifest_project_config_path(config, &cwd)?, + toml: source.clone(), + }) + }) + .collect::>>()?; + let user_toml = manifest + .configs + .iter() + .filter(|config| config.type_ == ManifestConfigType::User) + .filter_map(|config| config.source.clone()) + .collect(); + let inline_goal_override = manifest + .goal + .as_ref() + .filter(|goal| goal.type_ != ManifestGoalType::Graph) + .map(|goal| goal.text.clone()); + let title = manifest + .title + .as_ref() + .map(|title| fabro_types::normalize_explicit_run_title(title.as_str())) + .transpose()?; + let manifest_run_id = manifest + .run_id + .as_deref() + .map(str::parse::) + .transpose() + .context("invalid run ID")?; + let parent_id = manifest + .parent_id + .as_deref() + .map(str::parse::) + .transpose() + .context("invalid parent run ID")?; + + Ok(ManifestRunCompilerAdapter { + workflow_bundle, + entrypoint, + cwd, + project_settings, + user_toml, + run_overrides: overrides.run, + cli_overrides: overrides.cli, + input_overrides: overrides.input_overrides, + inline_goal_override, + run_id: explicit_run_id.or(manifest_run_id), + parent_id, + title, + git: manifest.git.clone(), + }) +} + +fn compiler_preparation_error_response(error: RunCompilerError) -> Response { + use run_compiler::{InvalidSettingsError, InvalidSourceError}; + + let detail = match error { + RunCompilerError::InvalidSource { source } => match source { + InvalidSourceError::MissingEntrypoint { .. } => { + "manifest target path is missing from workflows map".to_string() + } + InvalidSourceError::UnsupportedDockerfileReference { reference, .. } => { + format!("unsupported dockerfile reference: {reference}") + } + InvalidSourceError::MissingDockerfile { + dockerfile_path, .. + } => format!("missing bundled dockerfile: {dockerfile_path}"), + }, + RunCompilerError::InvalidSettings { source } => match *source { + InvalidSettingsError::Parse { .. } => "Failed to parse run config TOML".to_string(), + InvalidSettingsError::User { source } => source.to_string(), + InvalidSettingsError::Resolve { .. } => { + "failed to resolve manifest settings".to_string() + } + }, + RunCompilerError::VariableInterpolation { source } => { + format!("Run config variable interpolation failed: {source}") + } + other => return compiler_execution_error_response(other), + }; + ApiError::bad_request(detail).into_response() +} + +fn compiler_execution_error_response(error: RunCompilerError) -> Response { + match error { + RunCompilerError::ValidationOrParse { .. } => { + ApiError::bad_request("Validation failed").into_response() + } + RunCompilerError::ModelSelection { source } => { + ApiError::bad_request(format!("Model selection failed: {source}")).into_response() + } + RunCompilerError::ModelReference { source } => { + ApiError::bad_request(format!("Model reference failed: {source}")).into_response() + } + RunCompilerError::Internal { source, .. } => ApiError::new( + StatusCode::INTERNAL_SERVER_ERROR, + format!("Failed to persist run state: {source}"), + ) + .into_response(), + other => compiler_preparation_error_response(other), + } +} + pub(crate) async fn create_run_from_manifest( state: Arc, request: CreateRunFromManifestRequest, @@ -568,15 +717,45 @@ pub(crate) async fn create_run_from_manifest( let manifest_run_defaults = state.manifest_run_defaults(); let manifest_environment_defaults = state.environment_store().catalog_layer(); let manifest_mcp_server_catalog = state.mcp_server_store().catalog_settings(); - let mut prepared = match run_manifest::prepare_manifest_with_environment_defaults( - manifest_run_defaults.as_ref(), - manifest_environment_defaults.as_ref(), - &manifest_mcp_server_catalog, - &manifest, - ) { - Ok(prepared) => prepared, + let manifest_adapter = match adapt_manifest_for_run_compiler(&manifest, explicit_run_id) { + Ok(adapter) => adapter, Err(err) => return ApiError::bad_request(err.to_string()).into_response(), }; + let provenance = run_provenance(&headers, &actor); + let raw_compiler_input = RawRunCompilerInput { + workflow_bundle: manifest_adapter.workflow_bundle, + entrypoint: manifest_adapter.entrypoint, + cwd: manifest_adapter.cwd, + server_run_defaults: manifest_run_defaults.as_ref().clone(), + server_environment_defaults: manifest_environment_defaults.as_ref().clone(), + server_mcp_catalog: manifest_mcp_server_catalog, + project_settings: manifest_adapter.project_settings, + user_toml: manifest_adapter.user_toml, + run_overrides: manifest_adapter.run_overrides, + cli_overrides: manifest_adapter.cli_overrides, + input_overrides: manifest_adapter.input_overrides, + inline_goal_override: manifest_adapter.inline_goal_override, + vars: HashMap::new(), + run_id: manifest_adapter.run_id, + title: manifest_adapter.title, + parent_id: manifest_adapter.parent_id, + git: manifest_adapter.git, + storage_root: state.server_storage_dir(), + configured_providers: Vec::new(), + workflow_slug: None, + provenance, + web_url: None, + submitted_manifest_bytes: Some(submitted_manifest_bytes), + automation, + }; + let normalized = match run_compiler::normalize_source(raw_compiler_input) { + Ok(normalized) => normalized, + Err(err) => return compiler_preparation_error_response(err), + }; + let layered = match run_compiler::layer_settings(normalized) { + Ok(layered) => layered, + Err(err) => return compiler_preparation_error_response(err), + }; let vars = match snapshot_run_variables(&state).await { Ok(vars) => vars, Err(err) => { @@ -584,20 +763,19 @@ pub(crate) async fn create_run_from_manifest( .into_response(); } }; - if let Err(err) = substitute_run_variables(&vars, &mut prepared.settings) { - return ApiError::bad_request(format!("Run config variable interpolation failed: {err}")) - .into_response(); - } - let run_id = explicit_run_id - .or(prepared.run_id) - .unwrap_or_else(RunId::new); - let provider = run_manifest::effective_sandbox_provider(&prepared.settings.run); + let prepared = match run_compiler::apply_run_variables(layered.with_vars(vars)) { + Ok(prepared) => prepared, + Err(err) => return compiler_preparation_error_response(err), + }; + let (prepared, run_id) = prepared.resolve_run_id(); + let prepared = prepared.with_web_url(state.run_web_url(&run_id)); + let provider = run_manifest::effective_sandbox_provider(&prepared.settings().run); if let Some(error) = run_manifest::sandbox_provider_policy_error(&state.server_settings(), provider) { return ApiError::bad_request(error).into_response(); } - if let Some(parent_id) = prepared.parent_id { + if let Some(parent_id) = prepared.parent_id() { if parent_id == run_id { return ApiError::bad_request("A run cannot be its own parent.").into_response(); } @@ -607,7 +785,6 @@ pub(crate) async fn create_run_from_manifest( } info!(run_id = %run_id, "Run created"); - let web_url = state.run_web_url(&run_id); let catalog = state.catalog(); // Resolve once: we need both the provider IDs (for the run create input // and ask-fabro-readiness) and the LLM client itself (for the spawned @@ -628,34 +805,24 @@ pub(crate) async fn create_run_from_manifest( ready_provider_ids.clone() } }; - let provenance = run_provenance(&headers, &actor); - let mut create_input = run_manifest::create_run_input( - prepared.clone(), - run_materialization_provider_ids, - provenance, - web_url.clone(), - vars, - ); - create_input.run_id = Some(run_id); - create_input.submitted_manifest_bytes = Some(submitted_manifest_bytes); - create_input.automation = automation; - - let storage_root = state.server_storage_dir(); - let created = match Box::pin(operations::create( + let prepared = prepared.with_configured_providers(run_materialization_provider_ids); + let compiled = match run_compiler::compile_graph(prepared, Arc::clone(&catalog)).await { + Ok(compiled) => compiled, + Err(err) => return compiler_execution_error_response(err), + }; + let materialized = match run_compiler::materialize_run(compiled, catalog).await { + Ok(materialized) => materialized, + Err(err) => return compiler_execution_error_response(err), + }; + let compiler_output = run_compiler::assemble_run(materialized); + let (persistence_input, title_generation_target) = compiler_output.into_parts(); + let created = match Box::pin(operations::persist_create_run( state.stores.runs.as_ref(), - create_input, - storage_root, - catalog, + persistence_input, )) .await { Ok(created) => created, - Err(WorkflowError::ValidationFailed { .. } | WorkflowError::Parse(_)) => { - return ApiError::bad_request("Validation failed").into_response(); - } - Err(err @ (WorkflowError::ModelSelection(_) | WorkflowError::ModelReference(_))) => { - return ApiError::bad_request(err.to_string()).into_response(); - } Err(err) => { return ApiError::new( StatusCode::INTERNAL_SERVER_ERROR, @@ -699,7 +866,7 @@ pub(crate) async fn create_run_from_manifest( let run_spec = created.persisted.run_spec(); let workflow = run_title_generation::workflow_summary(&run_spec.graph); let run_inputs = run_spec.settings.run.inputs.clone(); - let workflow_target = prepared.target_path.to_string(); + let workflow_target = title_generation_target.to_string(); let title_catalog = state.catalog(); let title_model = title_catalog.small_default_for_configured_ids(&ready_provider_ids); let title_model_id = title_model.id.clone(); diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index b3da5be63..fd920c042 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -3660,6 +3660,249 @@ async fn create_run_from_manifest_helper_persists_automation_metadata() { assert_eq!(summary.automation, Some(automation)); } +#[tokio::test] +async fn create_run_from_manifest_pins_compiled_and_persisted_behavior() { + let state = TestAppStateBuilder::new() + .runtime_settings( + default_test_server_settings(), + manifest_run_defaults_from_toml( + r#" +[run.metadata] +server-label = "server" +layer = "server" +"#, + ), + ) + .env_lookup(|_| None) + .vault_entries([(EnvVars::OPENAI_API_KEY, "test-openai-api-key")]) + .build(); + let run_id = RunId::new(); + let dot = r#"digraph CompilePin { + graph [goal="Graph goal", target="{{ inputs.target }}"] + start [shape=Mdiamond] + work [prompt="Ship {{ inputs.target }}", model="gpt-5.4"] + exit [shape=Msquare] + start -> work -> exit + }"#; + let mut manifest_json = minimal_manifest_json(dot); + manifest_json["title"] = json!(" Pinned create "); + manifest_json["goal"] = json!({ + "type": "value", + "text": "Inline release goal" + }); + manifest_json["args"] = json!({ + "model": "gpt-5.4", + "input": ["target=payments"] + }); + manifest_json["configs"] = json!([{ + "type": "project", + "path": "/tmp/project/.fabro/project.toml", + "source": r#" +_version = 1 + +[project] +name = "payments-project" + +[run.metadata] +project-label = "project" +layer = "project" +"# + }]); + manifest_json["cwd"] = json!("/tmp/project"); + manifest_json["git"] = json!({ + "origin_url": "https://github.com/acme/payments.git", + "branch": "feature/compiler", + "sha": "0123456789abcdef", + "dirty": "clean", + "push_outcome": { "type": "not_attempted" } + }); + let manifest: RunManifest = serde_json::from_value(manifest_json).unwrap(); + let submitted_manifest_bytes = serde_json::to_vec(&manifest).unwrap(); + let mut headers = HeaderMap::new(); + headers.insert( + header::USER_AGENT, + "fabro-cli/9.8.7".parse().expect("user agent should parse"), + ); + + let response = Box::pin(handler::runs::create_run_from_manifest( + Arc::clone(&state), + handler::runs::CreateRunFromManifestRequest { + manifest, + submitted_manifest_bytes: submitted_manifest_bytes.clone(), + explicit_run_id: Some(run_id), + explicit_title_supplied: true, + actor: Principal::System { + system_kind: SystemActorKind::Engine, + }, + headers, + automation: None, + }, + )) + .await; + + let body = response_json!(response, StatusCode::CREATED).await; + assert_eq!(body["id"], run_id.to_string()); + assert_eq!(body["title"], "Pinned create"); + assert_eq!(body["lifecycle"]["status"]["kind"], "submitted"); + + let run_store = state.stores.runs.open_run_reader(&run_id).await.unwrap(); + let events = run_store.list_events().await.unwrap(); + assert_eq!( + events + .iter() + .map(|envelope| envelope.event.event_name()) + .collect::>(), + vec!["run.created", "run.submitted"] + ); + let run_state = run_store.state().await.unwrap(); + let spec = &run_state.spec; + assert_eq!(spec.run_id, run_id); + assert_eq!(spec.graph.goal(), "Inline release goal"); + assert_eq!( + spec.graph.attrs.get("target").and_then(AttrValue::as_str), + Some("{{ inputs.target }}") + ); + assert_eq!( + spec.graph.nodes["work"] + .attrs + .get("prompt") + .and_then(AttrValue::as_str), + Some("Ship payments") + ); + assert_eq!( + spec.graph.nodes["work"] + .attrs + .get("model") + .and_then(AttrValue::as_str), + Some("gpt-5.4") + ); + assert_eq!( + spec.graph.nodes["work"] + .attrs + .get("provider") + .and_then(AttrValue::as_str), + Some("openai") + ); + assert_eq!(spec.settings.run.model.name.as_deref(), Some("gpt-5.4")); + assert_eq!(spec.settings.run.model.provider.as_deref(), Some("openai")); + assert_eq!( + spec.settings.run.inputs.get("target"), + Some(&toml::Value::String("payments".to_string())) + ); + assert_eq!( + spec.settings.project.name.as_deref(), + Some("payments-project") + ); + assert_eq!( + spec.labels.get("project-label").map(String::as_str), + Some("project") + ); + assert_eq!( + spec.labels.get("layer").map(String::as_str), + Some("project") + ); + assert_eq!( + spec.git.as_ref().map(|git| git.origin_url.as_str()), + Some("https://github.com/acme/payments.git") + ); + + let created = events[0].event.to_value().unwrap(); + assert_eq!(created["properties"]["title"], "Pinned create"); + assert_eq!(created["properties"]["labels"]["project-label"], "project"); + assert_eq!( + created["properties"]["provenance"]["client"]["user_agent"], + "fabro-cli/9.8.7" + ); + assert_eq!( + created["properties"]["provenance"]["subject"]["kind"], + "system" + ); + let manifest_blob = created["properties"]["manifest_blob"] + .as_str() + .expect("run.created should carry the submitted source blob") + .parse::() + .unwrap(); + let persisted_manifest = run_store + .read_blob(&manifest_blob) + .await + .unwrap() + .expect("submitted source blob should exist"); + assert_eq!(persisted_manifest.as_ref(), submitted_manifest_bytes); +} + +#[tokio::test] +async fn create_run_from_manifest_pins_compiler_http_error_mappings() { + let cases = [ + ( + { + let mut manifest = minimal_manifest_json(MINIMAL_DOT); + manifest["version"] = json!(2); + manifest + }, + "unsupported manifest version 2", + ), + ( + minimal_manifest_json( + r#"digraph Test { + graph [goal="Test"] + start [shape=Mdiamond] + work [prompt="Use {{ vars.MISSING }}"] + exit [shape=Msquare] + start -> work -> exit + }"#, + ), + "Validation failed", + ), + ( + { + let mut manifest = minimal_manifest_json( + r#"digraph Test { + graph [goal="Test"] + start [shape=Mdiamond] + work [prompt="Do work", model="gpt-5.4", provider="missing-provider"] + exit [shape=Msquare] + start -> work -> exit + }"#, + ); + manifest["args"] = json!({ + "model": "gpt-5.4", + "provider": "missing-provider" + }); + manifest + }, + "Model selection failed: unknown model provider 'missing-provider'", + ), + ]; + + for (manifest_json, expected_detail) in cases { + let state = TestAppStateBuilder::new() + .env_lookup(|_| None) + .vault_entries([(EnvVars::OPENAI_API_KEY, "test-openai-api-key")]) + .build(); + let manifest: RunManifest = serde_json::from_value(manifest_json).unwrap(); + let submitted_manifest_bytes = serde_json::to_vec(&manifest).unwrap(); + + let response = Box::pin(handler::runs::create_run_from_manifest( + state, + handler::runs::CreateRunFromManifestRequest { + manifest, + submitted_manifest_bytes, + explicit_run_id: Some(RunId::new()), + explicit_title_supplied: true, + actor: Principal::System { + system_kind: SystemActorKind::Engine, + }, + headers: HeaderMap::new(), + automation: None, + }, + )) + .await; + + let body = response_json!(response, StatusCode::BAD_REQUEST).await; + assert_eq!(body["errors"][0]["detail"], expected_detail); + } +} + #[tokio::test] async fn fake_automation_materializer_injection_captures_input_and_returns_manifest() { let materialized_manifest: RunManifest = diff --git a/lib/components/fabro-workflow/src/operations/create.rs b/lib/components/fabro-workflow/src/operations/create.rs index 8451ef1df..f238fbc40 100644 --- a/lib/components/fabro-workflow/src/operations/create.rs +++ b/lib/components/fabro-workflow/src/operations/create.rs @@ -26,8 +26,7 @@ use crate::file_resolver::FileResolver; use crate::pipeline::types::PersistOptions; use crate::pipeline::{self, Persisted, TransformOptions, Validated}; use crate::records::RunSpec; -use crate::run_lookup::default_scratch_base; -use crate::run_materialization::materialize_run; +use crate::run_materialization; use crate::transforms::{ModelResolutionTransform, RenderMode}; use crate::workflow_bundle::{RunDefinition, WorkflowBundle}; @@ -58,6 +57,36 @@ pub struct CreateRunInput { pub web_url: Option, } +/// Inputs needed to resolve and compile a workflow for run creation. +#[derive(Clone, Debug)] +pub struct CreateRunCompileInput { + pub workflow: WorkflowInput, + pub settings: WorkflowSettings, + pub vars: HashMap, + pub cwd: PathBuf, + pub workflow_path: Option, + pub workflow_bundle: Option, + pub configured_providers: Vec, +} + +/// Durable metadata joined to a materialized workflow before persistence. +/// `run_id` is already resolved, and `storage_root` is used to derive the +/// run's scratch directory during pure input assembly. +#[derive(Clone, Debug)] +pub struct CreateRunPersistenceMetadata { + pub run_id: RunId, + pub storage_root: PathBuf, + pub workflow_slug: Option, + pub submitted_manifest_bytes: Option>, + pub title: Option, + pub automation: Option, + pub git: Option, + pub fork_source_ref: Option, + pub parent_id: Option, + pub provenance: RunProvenance, + pub web_url: Option, +} + #[derive(Debug)] pub struct CreatedRun { pub persisted: Persisted, @@ -66,20 +95,132 @@ pub struct CreatedRun { pub dot_path: Option, } -struct PersistCreateOptions { +/// Result of resolving, preprocessing, validating, and promoting a workflow +/// for run creation. Model selectors in the graph are resolved, while the run +/// settings still reflect the compiled source and have not been materialized. +pub struct CompiledRun { + validated: Validated, settings: WorkflowSettings, - run_id: Option, - run_dir: Option, + raw_source: String, workflow_slug: Option, - source_name: Option, + workflow_config: Option, + dot_path: Option, + current_dir: Option, + file_resolver: Option>, + definition: Option, + source_directory: String, labels: HashMap, - source_directory: Option, - automation: Option, - git: Option, - fork_source_ref: Option, - provenance: RunProvenance, configured_providers: Vec, - catalog: Arc, +} + +impl CompiledRun { + pub fn validated(&self) -> &Validated { + &self.validated + } + + pub fn settings(&self) -> &WorkflowSettings { + &self.settings + } + + pub fn resolved_source(&self) -> &str { + &self.raw_source + } + + pub fn workflow_slug(&self) -> Option<&str> { + self.workflow_slug.as_deref() + } + + pub fn workflow_config(&self) -> Option<&str> { + self.workflow_config.as_deref() + } + + pub fn dot_path(&self) -> Option<&Path> { + self.dot_path.as_deref() + } + + pub fn current_dir(&self) -> Option<&Path> { + self.current_dir.as_deref() + } + + pub fn file_resolver(&self) -> Option> { + self.file_resolver.clone() + } + + pub fn definition(&self) -> Option<&RunDefinition> { + self.definition.as_ref() + } + + pub fn source_directory(&self) -> &str { + &self.source_directory + } + + pub fn labels(&self) -> &HashMap { + &self.labels + } +} + +/// Compiled workflow with its run-level model settings materialized against +/// the same provider snapshot used during compilation. +pub struct MaterializedRun { + compiled: CompiledRun, + settings: WorkflowSettings, +} + +impl MaterializedRun { + pub fn compiled(&self) -> &CompiledRun { + &self.compiled + } + + pub fn settings(&self) -> &WorkflowSettings { + &self.settings + } +} + +/// Complete input for creating a durable run. The run ID and run directory +/// are resolved during assembly, before persistence begins. +pub struct CreateRunPersistenceInput { + materialized: MaterializedRun, + run_id: RunId, + run_dir: PathBuf, + workflow_slug: Option, + submitted_manifest_bytes: Option>, + title: Option, + automation: Option, + git: Option, + fork_source_ref: Option, + parent_id: Option, + provenance: RunProvenance, + web_url: Option, +} + +impl CreateRunPersistenceInput { + pub fn materialized(&self) -> &MaterializedRun { + &self.materialized + } + + pub fn run_id(&self) -> RunId { + self.run_id + } + + pub fn run_dir(&self) -> &Path { + &self.run_dir + } + + pub fn workflow_slug(&self) -> Option<&str> { + self.workflow_slug.as_deref() + } + + pub fn submitted_manifest_bytes(&self) -> Option<&[u8]> { + self.submitted_manifest_bytes.as_deref() + } + + pub fn automation(&self) -> Option<&AutomationRef> { + self.automation.as_ref() + } + + pub fn definition(&self) -> Option<&RunDefinition> { + self.materialized.compiled.definition.as_ref() + } } /// Resolve workflow inputs, normalize settings using the caller-provided @@ -90,96 +231,278 @@ pub async fn create( storage_root: PathBuf, catalog: Arc, ) -> Result { - let resolved = resolve_workflow(ResolveWorkflowInput { - workflow: request.workflow, - settings: request.settings, - cwd: request.cwd, + let run_id = request.run_id.unwrap_or_default(); + let persistence_input = spawn_blocking(move || { + let CreateRunInput { + workflow, + settings, + vars, + cwd, + workflow_slug, + workflow_path, + workflow_bundle, + submitted_manifest_bytes, + run_id: _, + title, + automation, + git, + fork_source_ref, + parent_id, + provenance, + configured_providers, + web_url, + } = request; + let compiled = compile_create_run( + CreateRunCompileInput { + workflow, + settings, + vars, + cwd, + workflow_path, + workflow_bundle, + configured_providers, + }, + Arc::clone(&catalog), + )?; + let materialized = materialize_create_run(compiled, catalog.as_ref())?; + Ok::<_, Error>(assemble_create_run_persistence_input( + materialized, + CreateRunPersistenceMetadata { + run_id, + storage_root, + workflow_slug, + submitted_manifest_bytes, + title, + automation, + git, + fork_source_ref, + parent_id, + provenance, + web_url, + }, + )) }) - .map_err(|err| Error::Parse(err.to_string()))?; - let labels = resolved.settings.combined_labels(); - let settings = resolved.settings.clone(); + .await + .map_err(|err| Error::engine_with_source("workflow create task failed", err))??; - let CreateRunInput { - workflow: _, - settings: _, + Box::pin(persist_create_run(store, persistence_input)).await +} + +/// Resolve, preprocess, validate, and promote a workflow for run creation. +/// +/// This stage is synchronous and may read workflow files. Async callers must +/// run it on a blocking thread. +pub fn compile_create_run( + input: CreateRunCompileInput, + catalog: Arc, +) -> Result { + let CreateRunCompileInput { + workflow, + settings, vars, - cwd: _, - workflow_slug, + cwd, workflow_path, workflow_bundle, - submitted_manifest_bytes, + 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 workflow_config = resolved + .workflow_toml_path + .as_deref() + .and_then(|path| std::fs::read_to_string(path).ok()); + let source_name = resolved + .dot_path + .as_ref() + .map(|path| path.display().to_string()); + let definition = match (workflow_path, workflow_bundle) { + (Some(workflow_path), Some(workflow_bundle)) => { + let bundled = workflow_bundle.workflow(&workflow_path).ok_or_else(|| { + Error::Parse("workflow path is missing from workflow bundle".to_string()) + })?; + if bundled.source != resolved.raw_source { + return Err(Error::Parse( + "resolved workflow does not match workflow bundle entrypoint".to_string(), + )); + } + Some(RunDefinition::new(workflow_path, workflow_bundle)) + } + (None, None) => None, + _ => { + return Err(Error::Parse( + "workflow path and workflow bundle must be provided together".to_string(), + )); + } + }; + let mut validated = preprocess_and_validate( + &resolved.raw_source, + resolved.goal_override.as_deref(), + &TransformOptions { + current_dir: resolved.current_dir.clone(), + file_resolver: resolved.file_resolver.clone(), + template_context: template_context(Some(&settings), vars), + source_name, + render_mode: RenderMode::Structural, + custom_transforms: Vec::new(), + model_resolution: Some( + ModelResolutionTransform::for_eligible( + catalog, + configured_providers.iter().cloned().collect(), + ) + .with_default_provider(configured_default_provider(&settings)), + ), + }, + )?; + + validated.promote_template_undefined_variables_to_errors(); + if validated.has_errors() { + return Err(Error::ValidationFailed { + diagnostics: validated.diagnostics().to_vec(), + }); + } + + Ok(CompiledRun { + validated, + settings, + raw_source: resolved.raw_source, + workflow_slug: resolved.workflow_slug, + workflow_config, + dot_path: resolved.dot_path, + current_dir: resolved.current_dir, + file_resolver: resolved.file_resolver, + definition, + source_directory: resolved.working_directory.to_string_lossy().to_string(), + labels, + configured_providers, + }) +} + +/// Materialize run-level model settings from a compiled workflow. +pub fn materialize_create_run( + compiled: CompiledRun, + catalog: &Catalog, +) -> Result { + let settings = run_materialization::materialize_run( + compiled.settings.clone(), + compiled.validated.graph(), + catalog, + &compiled.configured_providers, + )?; + Ok(MaterializedRun { compiled, settings }) +} + +/// Assemble all inputs needed for persistence without I/O or recompilation. +pub fn assemble_create_run_persistence_input( + materialized: MaterializedRun, + metadata: CreateRunPersistenceMetadata, +) -> CreateRunPersistenceInput { + let CreateRunPersistenceMetadata { run_id, + storage_root, + workflow_slug, + submitted_manifest_bytes, title, automation, git, fork_source_ref, parent_id, provenance, - configured_providers, web_url, - } = request; + } = metadata; + let run_dir = Storage::new(storage_root) + .run_scratch(&run_id) + .root() + .to_path_buf(); + let workflow_slug = workflow_slug.or_else(|| materialized.compiled.workflow_slug.clone()); - let run_id = run_id.unwrap_or_else(RunId::new); - let storage = Storage::new(storage_root); - let run_dir = storage.run_scratch(&run_id).root().to_path_buf(); - let source_directory = Some(resolved.working_directory.to_string_lossy().to_string()); + CreateRunPersistenceInput { + materialized, + run_id, + run_dir, + workflow_slug, + submitted_manifest_bytes, + title, + automation, + git, + fork_source_ref, + parent_id, + provenance, + web_url, + } +} - let goal_override = resolved.goal_override.clone(); - let current_dir = resolved.current_dir.clone(); - let file_resolver = resolved.file_resolver.clone(); - let resolved_workflow_slug = resolved.workflow_slug.clone(); +/// Persist one already-compiled and materialized run without recompiling it. +pub async fn persist_create_run( + store: &Database, + input: CreateRunPersistenceInput, +) -> Result { + let CreateRunPersistenceInput { + materialized, + run_id, + run_dir, + workflow_slug, + submitted_manifest_bytes, + title, + automation, + git, + fork_source_ref, + parent_id, + provenance, + web_url, + } = input; + let MaterializedRun { compiled, settings } = materialized; + let CompiledRun { + validated, + settings: _, + raw_source, + workflow_slug: _, + workflow_config, + dot_path, + current_dir: _, + file_resolver: _, + definition, + source_directory, + labels, + configured_providers: _, + } = compiled; let persisted_run_dir = run_dir.clone(); - let accepted_definition = match (&workflow_path, &workflow_bundle) { - (Some(workflow_path), Some(workflow_bundle)) => Some(RunDefinition::new( - workflow_path.clone(), - workflow_bundle.clone(), - )), - _ => None, - }; - - let raw_source = resolved.raw_source.clone(); - let source_name = resolved - .dot_path - .as_ref() - .map(|path| path.display().to_string()); let persisted = spawn_blocking(move || { - create_from_source( - &raw_source, - vars, - PersistCreateOptions { - settings, - run_id: Some(run_id), - run_dir: Some(persisted_run_dir), - workflow_slug: workflow_slug.or(resolved_workflow_slug), - source_name, - labels, - source_directory, - automation, - git, - fork_source_ref, - provenance, - configured_providers, - catalog, - }, - current_dir, - file_resolver, - goal_override.as_deref(), - ) + let run_spec = RunSpec { + run_id, + settings, + graph: validated.graph().clone(), + graph_source: Some(validated.source().to_string()), + workflow_slug, + automation, + source_directory: Some(source_directory), + labels, + provenance, + manifest_blob: None, + definition_blob: None, + git, + fork_source_ref, + }; + pipeline::persist(validated, PersistOptions { + run_dir: persisted_run_dir, + run_spec, + }) }) .await .map_err(|err| Error::engine_with_source("workflow create task failed", err))??; - let workflow_config = resolved - .workflow_toml_path - .as_deref() - .and_then(|path| std::fs::read_to_string(path).ok()); persist_created_run( store, &persisted, - &resolved.raw_source, + &raw_source, workflow_config, submitted_manifest_bytes.as_deref(), - accepted_definition.as_ref(), + definition.as_ref(), title, parent_id, web_url, @@ -190,7 +513,7 @@ pub async fn create( persisted, run_id, run_dir, - dot_path: resolved.dot_path, + dot_path, }) } @@ -285,40 +608,6 @@ fn store_error(err: impl std::fmt::Display) -> Error { Error::engine(err.to_string()) } -fn create_from_source( - dot_source: &str, - vars: HashMap, - options: PersistCreateOptions, - current_dir: Option, - file_resolver: Option>, - goal_override: Option<&str>, -) -> Result { - let mut validated = preprocess_and_validate(dot_source, goal_override, &TransformOptions { - current_dir, - file_resolver, - template_context: template_context(Some(&options.settings), vars), - source_name: options.source_name.clone(), - render_mode: RenderMode::Structural, - custom_transforms: Vec::new(), - model_resolution: Some( - ModelResolutionTransform::for_eligible( - Arc::clone(&options.catalog), - options.configured_providers.iter().cloned().collect(), - ) - .with_default_provider(configured_default_provider(&options.settings)), - ), - })?; - - validated.promote_template_undefined_variables_to_errors(); - if validated.has_errors() { - return Err(Error::ValidationFailed { - diagnostics: validated.diagnostics().to_vec(), - }); - } - - persist_validated(validated, options) -} - /// Parse, transform, and validate `dot_source`. /// /// `options.model_resolution` drives both halves of catalog awareness: it @@ -376,59 +665,6 @@ fn apply_goal_override(graph: &mut Graph, goal_override: Option<&str>) { } } -fn persist_validated( - validated: Validated, - options: PersistCreateOptions, -) -> Result { - let PersistCreateOptions { - settings, - run_id, - run_dir, - workflow_slug, - source_name: _, - labels, - source_directory, - automation, - git, - fork_source_ref, - provenance, - configured_providers, - catalog, - } = options; - - let settings = materialize_run( - settings, - validated.graph(), - catalog.as_ref(), - &configured_providers, - )?; - - let run_id = run_id.unwrap_or_else(RunId::new); - let run_dir = run_dir.unwrap_or_else(|| default_run_dir(&run_id)); - - let run_spec = RunSpec { - run_id, - settings, - graph: validated.graph().clone(), - graph_source: Some(validated.source().to_string()), - workflow_slug, - automation, - source_directory, - labels, - provenance, - manifest_blob: None, - definition_blob: None, - git, - fork_source_ref, - }; - - pipeline::persist(validated, PersistOptions { run_dir, run_spec }) -} - -pub(crate) fn default_run_dir(run_id: &RunId) -> PathBuf { - make_run_dir(&default_scratch_base(), run_id) -} - pub fn make_run_dir(scratch_base: &Path, run_id: &RunId) -> PathBuf { fabro_config::RunScratch::for_run(scratch_base, run_id) .root() @@ -456,6 +692,7 @@ mod tests { use object_store::memory::InMemory; use super::*; + use crate::file_resolver::FileResolver; use crate::operations::{ValidateInput, validate, validate_with_catalog}; use crate::pipeline::types::{GOAL_SELF_REFERENCE_RULE, TEMPLATE_UNDEFINED_VARIABLE_RULE}; use crate::transforms::Transform; @@ -547,6 +784,38 @@ reasoning = false Catalog::builtin().all_provider_ids().into_iter().collect() } + fn compile_input(request: &CreateRunInput) -> CreateRunCompileInput { + CreateRunCompileInput { + workflow: request.workflow.clone(), + settings: request.settings.clone(), + vars: request.vars.clone(), + cwd: request.cwd.clone(), + workflow_path: request.workflow_path.clone(), + workflow_bundle: request.workflow_bundle.clone(), + configured_providers: request.configured_providers.clone(), + } + } + + fn persistence_metadata( + request: &CreateRunInput, + run_id: RunId, + storage_root: &Path, + ) -> CreateRunPersistenceMetadata { + CreateRunPersistenceMetadata { + run_id, + storage_root: storage_root.to_path_buf(), + workflow_slug: request.workflow_slug.clone(), + submitted_manifest_bytes: request.submitted_manifest_bytes.clone(), + title: request.title.clone(), + automation: request.automation.clone(), + git: request.git.clone(), + fork_source_ref: request.fork_source_ref.clone(), + parent_id: request.parent_id, + provenance: request.provenance.clone(), + web_url: request.web_url.clone(), + } + } + fn validate_dot(dot_source: &str, settings: WorkflowSettings) -> Validated { validate_with_catalog( ValidateInput { @@ -1337,6 +1606,244 @@ reasoning = false ); } + #[test] + fn assemble_create_run_persistence_input_resolves_complete_durable_identity() { + let dir = tempfile::tempdir().unwrap(); + let storage_root = dir.path().join("storage"); + let automation = AutomationRef { + id: "nightly".to_string(), + name: Some("Nightly".to_string()), + trigger_id: Some("schedule_1".to_string()), + }; + let request = CreateRunInput { + workflow: WorkflowInput::DotSource { + source: MINIMAL_DOT.to_string(), + base_dir: None, + }, + settings: test_default_settings(), + vars: HashMap::new(), + cwd: dir.path().to_path_buf(), + workflow_slug: Some("request-slug".to_string()), + workflow_path: None, + workflow_bundle: None, + submitted_manifest_bytes: Some(b"submitted manifest".to_vec()), + run_id: Some(fixtures::RUN_1), + title: Some("Assembled run".to_string()), + automation: Some(automation.clone()), + git: None, + fork_source_ref: None, + parent_id: Some(fixtures::RUN_2), + provenance: test_support::test_run_provenance(), + configured_providers: test_provider_ids(), + web_url: Some("https://fabro.test/runs/1".to_string()), + }; + let catalog = test_catalog(); + let resolved_run_id = fixtures::RUN_64; + + let compiled = compile_create_run(compile_input(&request), Arc::clone(&catalog)).unwrap(); + let materialized = materialize_create_run(compiled, catalog.as_ref()).unwrap(); + let metadata = persistence_metadata(&request, resolved_run_id, &storage_root); + let input = assemble_create_run_persistence_input(materialized, metadata); + + assert_eq!(input.run_id(), resolved_run_id); + assert_eq!( + input.run_dir(), + Storage::new(&storage_root) + .run_scratch(&resolved_run_id) + .root() + ); + assert_eq!(input.workflow_slug(), Some("request-slug")); + assert_eq!( + input.submitted_manifest_bytes(), + Some(b"submitted manifest".as_slice()) + ); + assert_eq!(input.automation(), Some(&automation)); + assert_eq!( + input.materialized().settings().run.model.name.as_deref(), + Some("claude-sonnet-5") + ); + } + + #[test] + fn compile_create_run_rejects_mismatched_bundle_definition() { + let workflow_path = ManifestPath::from_wire("workflows/main.fabro").unwrap(); + let compiled_workflow = BundledWorkflow { + path: workflow_path.clone(), + source: MINIMAL_DOT.to_string(), + config: None, + files: HashMap::new(), + }; + let mismatched_bundle = + WorkflowBundle::new(HashMap::from([(workflow_path.clone(), BundledWorkflow { + source: MINIMAL_DOT.replace("Build feature", "Different goal"), + ..compiled_workflow.clone() + })])); + + let Err(error) = compile_create_run( + CreateRunCompileInput { + workflow: WorkflowInput::Bundled(compiled_workflow), + settings: test_default_settings(), + vars: HashMap::new(), + cwd: PathBuf::from("/tmp/project"), + workflow_path: Some(workflow_path), + workflow_bundle: Some(mismatched_bundle), + configured_providers: test_provider_ids(), + }, + test_catalog(), + ) else { + panic!("mismatched accepted definition should fail"); + }; + + assert!(matches!(error, Error::Parse(message) if message == + "resolved workflow does not match workflow bundle entrypoint")); + } + + #[test] + fn compile_create_run_exposes_resolved_metadata_and_definition() { + let workflow_path = ManifestPath::from_wire("workflows/main.fabro").unwrap(); + let bundled = BundledWorkflow { + path: workflow_path.clone(), + source: MINIMAL_DOT.to_string(), + config: None, + files: HashMap::new(), + }; + let bundle = WorkflowBundle::new(HashMap::from([(workflow_path.clone(), bundled.clone())])); + let compiled = compile_create_run( + CreateRunCompileInput { + workflow: WorkflowInput::Bundled(bundled), + settings: test_default_settings(), + vars: HashMap::new(), + cwd: PathBuf::from("/tmp/project"), + workflow_path: Some(workflow_path.clone()), + workflow_bundle: Some(bundle), + configured_providers: test_provider_ids(), + }, + test_catalog(), + ) + .unwrap(); + + assert_eq!(compiled.resolved_source(), MINIMAL_DOT); + assert_eq!(compiled.current_dir(), Some(Path::new("workflows"))); + assert_eq!(compiled.dot_path(), Some(workflow_path.as_path())); + assert!(compiled.file_resolver().is_some()); + assert_eq!(compiled.labels(), &compiled.settings().combined_labels()); + let materialized = materialize_create_run(compiled, test_catalog().as_ref()).unwrap(); + let input = + assemble_create_run_persistence_input(materialized, CreateRunPersistenceMetadata { + run_id: fixtures::RUN_1, + storage_root: PathBuf::from("/tmp/storage"), + workflow_slug: None, + submitted_manifest_bytes: None, + title: None, + automation: None, + git: None, + fork_source_ref: None, + parent_id: None, + provenance: test_support::test_run_provenance(), + web_url: None, + }); + let definition = input + .definition() + .expect("bundled create input should retain a run definition"); + assert_eq!(definition.workflow_path, workflow_path); + } + + #[tokio::test] + async fn persist_create_run_uses_compiled_graph_without_recompiling_source() { + let dir = tempfile::tempdir().unwrap(); + let storage_root = dir.path().join("storage"); + let dot_path = dir.path().join("workflow.fabro"); + let compiled_source = MINIMAL_DOT.replace("Build feature", "Compiled goal"); + std::fs::write(&dot_path, &compiled_source).unwrap(); + let automation = AutomationRef { + id: "nightly".to_string(), + name: Some("Nightly".to_string()), + trigger_id: Some("schedule_1".to_string()), + }; + let request = CreateRunInput { + workflow: WorkflowInput::Path(dot_path.clone()), + settings: test_default_settings(), + vars: HashMap::new(), + cwd: dir.path().to_path_buf(), + workflow_slug: Some("compiled-slug".to_string()), + workflow_path: None, + workflow_bundle: None, + submitted_manifest_bytes: Some(b"submitted manifest".to_vec()), + run_id: Some(fixtures::RUN_2), + title: Some("Compiled run".to_string()), + automation: Some(automation.clone()), + git: None, + fork_source_ref: None, + parent_id: None, + provenance: test_support::test_run_provenance(), + configured_providers: test_provider_ids(), + web_url: None, + }; + let catalog = test_catalog(); + let workflow_config_path = dir.path().join("workflow.toml"); + std::fs::write( + &workflow_config_path, + "_version = 1\n[workflow]\ngraph = \"workflow.fabro\"\n", + ) + .unwrap(); + let compiled = compile_create_run(compile_input(&request), Arc::clone(&catalog)).unwrap(); + assert_eq!( + compiled.workflow_config(), + Some("_version = 1\n[workflow]\ngraph = \"workflow.fabro\"\n") + ); + + std::fs::write(&dot_path, "this is no longer a graph").unwrap(); + std::fs::write(&workflow_config_path, "changed after compilation").unwrap(); + + let materialized = materialize_create_run(compiled, catalog.as_ref()).unwrap(); + let metadata = persistence_metadata(&request, fixtures::RUN_2, &storage_root); + let input = assemble_create_run_persistence_input(materialized, metadata); + let store = memory_store(); + let created = persist_create_run(store.as_ref(), input).await.unwrap(); + + assert_eq!(created.run_id, fixtures::RUN_2); + assert_eq!(created.dot_path.as_deref(), Some(dot_path.as_path())); + assert_eq!(created.persisted.graph().goal(), "Compiled goal"); + assert_eq!(created.persisted.source(), compiled_source); + + let run_store = store.open_run_reader(&fixtures::RUN_2).await.unwrap(); + let state = run_store.state().await.unwrap(); + assert_eq!(state.spec.graph.goal(), "Compiled goal"); + assert_eq!(state.spec.automation, Some(automation)); + let events = run_store.list_events().await.unwrap(); + assert_eq!( + events + .iter() + .map(|event| event.event.event_name()) + .collect::>(), + vec!["run.created", "run.submitted"] + ); + let EventBody::RunCreated(created) = &events[0].event.body else { + panic!("first durable event should be run.created"); + }; + assert_eq!( + created.workflow_source.as_deref(), + Some(compiled_source.as_str()) + ); + assert_eq!( + created.workflow_config.as_deref(), + Some("_version = 1\n[workflow]\ngraph = \"workflow.fabro\"\n") + ); + let manifest_blob = created + .manifest_blob + .as_ref() + .expect("submitted manifest should be persisted"); + assert_eq!( + run_store + .read_blob(manifest_blob) + .await + .unwrap() + .expect("submitted manifest blob should exist") + .as_ref(), + b"submitted manifest" + ); + } + #[tokio::test] async fn create_returns_validation_failed_with_diagnostics() { let dot = r#"digraph Test { diff --git a/lib/components/fabro-workflow/src/operations/mod.rs b/lib/components/fabro-workflow/src/operations/mod.rs index 9523c5ff5..57ed472f4 100644 --- a/lib/components/fabro-workflow/src/operations/mod.rs +++ b/lib/components/fabro-workflow/src/operations/mod.rs @@ -14,7 +14,12 @@ pub use archive::{ ArchiveOutcome, UnarchiveOutcome, archive, archived_rejection_message, ensure_not_archived, unarchive, }; -pub use create::{CreateRunInput, CreatedRun, create, make_run_dir}; +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, +}; pub use fork::{ForkOutcome, ForkRunInput, ResolvedForkTarget, fork_run}; pub use resume::resume; pub use retry::{RetryOutcome, RetryRunInput, retry_run};