mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-07 03:00:29 +00:00
fabro(01KYQN78K19NY7PNSCDYP6CG9G): implement (succeeded)
Fabro-Run: 01KYQN78K19NY7PNSCDYP6CG9G
Fabro-Completed: 5
Fabro-Checkpoint: 879008d9d5
⚒️ Generated with [Fabro](https://fabro.sh)
This commit is contained in:
parent
0815f9fa8f
commit
ac6e3ced6a
7 changed files with 2257 additions and 245 deletions
|
|
@ -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;
|
||||
|
|
|
|||
1105
lib/apps/fabro-server/src/run_compiler.rs
Normal file
1105
lib/apps/fabro-server/src/run_compiler.rs
Normal file
File diff suppressed because it is too large
Load diff
|
|
@ -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<types::GitContext>,
|
||||
pub root_source: String,
|
||||
#[allow(
|
||||
dead_code,
|
||||
reason = "create now resolves identity in the run compiler adapter"
|
||||
)]
|
||||
pub run_id: Option<RunId>,
|
||||
#[allow(
|
||||
dead_code,
|
||||
reason = "create now resolves lineage in the run compiler adapter"
|
||||
)]
|
||||
pub parent_id: Option<RunId>,
|
||||
#[allow(
|
||||
dead_code,
|
||||
reason = "create now normalizes titles in the run compiler adapter"
|
||||
)]
|
||||
pub title: Option<String>,
|
||||
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<RunLayer>,
|
||||
cli: Option<CliLayer>,
|
||||
input_overrides: HashMap<String, toml::Value>,
|
||||
pub(crate) struct ManifestSettingsOverrides {
|
||||
pub(crate) run: Option<RunLayer>,
|
||||
pub(crate) cli: Option<CliLayer>,
|
||||
pub(crate) input_overrides: HashMap<String, toml::Value>,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
|
|
@ -235,34 +247,6 @@ fn manifest_validate_input(
|
|||
}
|
||||
}
|
||||
|
||||
pub(crate) fn create_run_input(
|
||||
prepared: PreparedManifest,
|
||||
configured_providers: Vec<ProviderId>,
|
||||
provenance: RunProvenance,
|
||||
web_url: Option<String>,
|
||||
vars: HashMap<String, String>,
|
||||
) -> 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<ManifestSettingsOverrides> {
|
||||
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<ManifestPath> {
|
||||
|
|
|
|||
|
|
@ -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<AutomationRef>,
|
||||
}
|
||||
|
||||
struct ManifestRunCompilerAdapter {
|
||||
workflow_bundle: WorkflowBundle,
|
||||
entrypoint: ManifestPath,
|
||||
cwd: PathBuf,
|
||||
project_settings: Vec<ProjectSettingsSource>,
|
||||
user_toml: Vec<String>,
|
||||
run_overrides: Option<fabro_config::RunLayer>,
|
||||
cli_overrides: Option<fabro_config::CliLayer>,
|
||||
input_overrides: HashMap<String, toml::Value>,
|
||||
inline_goal_override: Option<String>,
|
||||
run_id: Option<RunId>,
|
||||
parent_id: Option<RunId>,
|
||||
title: Option<String>,
|
||||
git: Option<fabro_types::GitContext>,
|
||||
}
|
||||
|
||||
fn adapt_manifest_for_run_compiler(
|
||||
manifest: &RunManifest,
|
||||
explicit_run_id: Option<RunId>,
|
||||
) -> anyhow::Result<ManifestRunCompilerAdapter> {
|
||||
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::<anyhow::Result<Vec<_>>>()?;
|
||||
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::<RunId>)
|
||||
.transpose()
|
||||
.context("invalid run ID")?;
|
||||
let parent_id = manifest
|
||||
.parent_id
|
||||
.as_deref()
|
||||
.map(str::parse::<RunId>)
|
||||
.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<AppState>,
|
||||
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();
|
||||
|
|
|
|||
|
|
@ -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<_>>(),
|
||||
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::<RunBlobId>()
|
||||
.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 =
|
||||
|
|
|
|||
|
|
@ -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<String>,
|
||||
}
|
||||
|
||||
/// 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<String, String>,
|
||||
pub cwd: PathBuf,
|
||||
pub workflow_path: Option<ManifestPath>,
|
||||
pub workflow_bundle: Option<WorkflowBundle>,
|
||||
pub configured_providers: Vec<ProviderId>,
|
||||
}
|
||||
|
||||
/// 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<String>,
|
||||
pub submitted_manifest_bytes: Option<Vec<u8>>,
|
||||
pub title: Option<String>,
|
||||
pub automation: Option<AutomationRef>,
|
||||
pub git: Option<GitContext>,
|
||||
pub fork_source_ref: Option<ForkSourceRef>,
|
||||
pub parent_id: Option<RunId>,
|
||||
pub provenance: RunProvenance,
|
||||
pub web_url: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct CreatedRun {
|
||||
pub persisted: Persisted,
|
||||
|
|
@ -66,20 +95,132 @@ pub struct CreatedRun {
|
|||
pub dot_path: Option<PathBuf>,
|
||||
}
|
||||
|
||||
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<RunId>,
|
||||
run_dir: Option<PathBuf>,
|
||||
raw_source: String,
|
||||
workflow_slug: Option<String>,
|
||||
source_name: Option<String>,
|
||||
workflow_config: Option<String>,
|
||||
dot_path: Option<PathBuf>,
|
||||
current_dir: Option<PathBuf>,
|
||||
file_resolver: Option<Arc<dyn FileResolver>>,
|
||||
definition: Option<RunDefinition>,
|
||||
source_directory: String,
|
||||
labels: HashMap<String, String>,
|
||||
source_directory: Option<String>,
|
||||
automation: Option<AutomationRef>,
|
||||
git: Option<GitContext>,
|
||||
fork_source_ref: Option<ForkSourceRef>,
|
||||
provenance: RunProvenance,
|
||||
configured_providers: Vec<ProviderId>,
|
||||
catalog: Arc<Catalog>,
|
||||
}
|
||||
|
||||
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<Arc<dyn FileResolver>> {
|
||||
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<String, String> {
|
||||
&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<String>,
|
||||
submitted_manifest_bytes: Option<Vec<u8>>,
|
||||
title: Option<String>,
|
||||
automation: Option<AutomationRef>,
|
||||
git: Option<GitContext>,
|
||||
fork_source_ref: Option<ForkSourceRef>,
|
||||
parent_id: Option<RunId>,
|
||||
provenance: RunProvenance,
|
||||
web_url: Option<String>,
|
||||
}
|
||||
|
||||
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<Catalog>,
|
||||
) -> Result<CreatedRun, Error> {
|
||||
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<Catalog>,
|
||||
) -> Result<CompiledRun, Error> {
|
||||
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<MaterializedRun, Error> {
|
||||
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<CreatedRun, Error> {
|
||||
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<String, String>,
|
||||
options: PersistCreateOptions,
|
||||
current_dir: Option<PathBuf>,
|
||||
file_resolver: Option<Arc<dyn FileResolver>>,
|
||||
goal_override: Option<&str>,
|
||||
) -> Result<Persisted, Error> {
|
||||
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<Persisted, Error> {
|
||||
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<_>>(),
|
||||
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 {
|
||||
|
|
|
|||
|
|
@ -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};
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue