diff --git a/docs/plans/2026-04-23-001-refactor-collapse-settings-resolve-indirection-plan.md b/docs/plans/2026-04-23-001-refactor-collapse-settings-resolve-indirection-plan.md index b94e5a739..c1977a80b 100644 --- a/docs/plans/2026-04-23-001-refactor-collapse-settings-resolve-indirection-plan.md +++ b/docs/plans/2026-04-23-001-refactor-collapse-settings-resolve-indirection-plan.md @@ -1,7 +1,7 @@ --- title: "refactor: collapse settings resolve indirection layer" type: refactor -status: active +status: completed date: 2026-04-23 --- @@ -105,7 +105,7 @@ None directly applicable. The original `Resolver` commit (52d24531d) explicitly ## Implementation Units -- [ ] **Unit 1: Add `WorkflowSettings` to `fabro-config::context`** +- [x] **Unit 1: Add `WorkflowSettings` to `fabro-config::context`** **Goal:** New consumer bundle for workflow ops, sibling to `ServerSettings`/`UserSettings`. Workflow-only namespaces. @@ -145,7 +145,7 @@ None directly applicable. The original `Resolver` commit (52d24531d) explicitly --- -- [ ] **Unit 2: Delete `Resolver`; rewrite `*Settings::from_layer` and `resolve_*_from_file` to not depend on it** +- [x] **Unit 2: Delete `Resolver`; rewrite `*Settings::from_layer` and `resolve_*_from_file` to not depend on it** **Goal:** Remove the `Resolver` indirection. Each consumer-bundle constructor and each standalone helper inlines `apply_builtin_defaults` + the per-namespace resolve. @@ -186,7 +186,7 @@ None directly applicable. The original `Resolver` commit (52d24531d) explicitly --- -- [ ] **Unit 3: Migrate `create_run` to `WorkflowSettings`; add `storage_root` parameter; delete `ResolvedSettingsTree` + `resolve_settings_tree` + `combined_labels` free fn** +- [x] **Unit 3: Migrate `create_run` to `WorkflowSettings`; add `storage_root` parameter; delete `ResolvedSettingsTree` + `resolve_settings_tree` + `combined_labels` free fn** **Goal:** `create_run` consumes `WorkflowSettings` for workflow-owned config and `storage_root` from the caller for deployment-owned config. `WorkflowSettings::from_layer` failures keep today's `; `-joined `render_resolve_errors` format (R6). Storage-root interpolation failures move out to the HTTP handler with new error attribution (R5). @@ -247,7 +247,7 @@ None directly applicable. The original `Resolver` commit (52d24531d) explicitly --- -- [ ] **Unit 4: Workspace build + lint sweep** +- [x] **Unit 4: Workspace build + lint sweep** **Goal:** Confirm no caller across the workspace still references deleted symbols. diff --git a/lib/crates/fabro-config/src/context.rs b/lib/crates/fabro-config/src/context.rs index 5cddd9e6a..42d77793f 100644 --- a/lib/crates/fabro-config/src/context.rs +++ b/lib/crates/fabro-config/src/context.rs @@ -1,9 +1,16 @@ -use fabro_types::settings::{CliNamespace, FeaturesNamespace, ServerNamespace, SettingsLayer}; +use std::collections::HashMap; + +use fabro_types::settings::{ + CliNamespace, FeaturesNamespace, ProjectNamespace, RunNamespace, ServerNamespace, + SettingsLayer, WorkflowNamespace, +}; use serde::{Deserialize, Serialize}; -use crate::resolve::Resolver; use crate::user::load_settings_config; -use crate::{Error, Result}; +use crate::{ + Error, ResolveError, Result, apply_builtin_defaults, resolve_cli, resolve_features, + resolve_project, resolve_run, resolve_server, resolve_workflow, +}; #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct ServerSettings { @@ -13,10 +20,10 @@ pub struct ServerSettings { impl ServerSettings { pub fn from_layer(layer: &SettingsLayer) -> Result { - let resolver = Resolver::from_layer(layer); + let layer = apply_builtin_defaults(layer.clone()); let mut errors = Vec::new(); - let server = resolver.server_into(&mut errors); - let features = resolver.features_into(&mut errors); + let server = resolve_server(&layer.server.clone().unwrap_or_default(), &mut errors); + let features = resolve_features(&layer.features.clone().unwrap_or_default(), &mut errors); if errors.is_empty() { Ok(Self { server, features }) } else { @@ -38,10 +45,10 @@ pub struct UserSettings { impl UserSettings { pub fn from_layer(layer: &SettingsLayer) -> Result { - let resolver = Resolver::from_layer(layer); + let layer = apply_builtin_defaults(layer.clone()); let mut errors = Vec::new(); - let cli = resolver.cli_into(&mut errors); - let features = resolver.features_into(&mut errors); + let cli = resolve_cli(&layer.cli.clone().unwrap_or_default(), &mut errors); + let features = resolve_features(&layer.features.clone().unwrap_or_default(), &mut errors); if errors.is_empty() { Ok(Self { cli, features }) } else { @@ -54,3 +61,36 @@ impl UserSettings { Self::from_layer(&layer) } } + +#[derive(Debug, Clone, PartialEq, Serialize)] +pub struct WorkflowSettings { + pub project: ProjectNamespace, + pub workflow: WorkflowNamespace, + pub run: RunNamespace, +} + +impl WorkflowSettings { + pub fn from_layer(layer: &SettingsLayer) -> std::result::Result> { + let layer = apply_builtin_defaults(layer.clone()); + let mut errors = Vec::new(); + let project = resolve_project(&layer.project.clone().unwrap_or_default(), &mut errors); + let workflow = resolve_workflow(&layer.workflow.clone().unwrap_or_default(), &mut errors); + let run = resolve_run(&layer.run.clone().unwrap_or_default(), &mut errors); + if errors.is_empty() { + Ok(Self { + project, + workflow, + run, + }) + } else { + Err(errors) + } + } + + pub fn combined_labels(&self) -> HashMap { + let mut labels = self.project.metadata.clone(); + labels.extend(self.workflow.metadata.clone()); + labels.extend(self.run.metadata.clone()); + labels + } +} diff --git a/lib/crates/fabro-config/src/lib.rs b/lib/crates/fabro-config/src/lib.rs index ed5772707..b5739ed9e 100644 --- a/lib/crates/fabro-config/src/lib.rs +++ b/lib/crates/fabro-config/src/lib.rs @@ -2,8 +2,9 @@ clippy::disallowed_methods, reason = "sync config loading utilities used at startup; not on a Tokio path" )] -//! Resolved settings entrypoints: [`ServerSettings`] for the running server and -//! [`UserSettings`] for the CLI/user perspective. +//! Resolved settings entrypoints: [`ServerSettings`] for the running server, +//! [`UserSettings`] for the CLI/user perspective, and [`WorkflowSettings`] for +//! workflow execution. extern crate self as fabro_config; @@ -27,7 +28,7 @@ pub mod user; use std::path::Path; -pub use context::{ServerSettings, UserSettings}; +pub use context::{ServerSettings, UserSettings, WorkflowSettings}; pub use defaults::{apply_builtin_defaults, defaults_layer}; pub use error::{Error, Result}; pub use fabro_util::path::expand_tilde; @@ -37,7 +38,7 @@ pub use load::{ }; pub use parse::{ParseError, parse_settings_layer}; pub use resolve::{ - ResolveError, Resolver, dev_token_auth_enabled, render_resolve_errors, resolve_cli, + ResolveError, dev_token_auth_enabled, render_resolve_errors, resolve_cli, resolve_cli_from_file, resolve_features, resolve_features_from_file, resolve_project, resolve_project_from_file, resolve_run, resolve_run_from_file, resolve_server, resolve_server_from_file, resolve_storage_root, resolve_workflow, resolve_workflow_from_file, diff --git a/lib/crates/fabro-config/src/resolve/mod.rs b/lib/crates/fabro-config/src/resolve/mod.rs index f35282ed7..754cba237 100644 --- a/lib/crates/fabro-config/src/resolve/mod.rs +++ b/lib/crates/fabro-config/src/resolve/mod.rs @@ -2,7 +2,6 @@ mod cli; mod error; mod features; mod project; -mod resolver; mod run; mod server; mod workflow; @@ -15,45 +14,71 @@ use fabro_types::settings::{ }; pub use features::resolve_features; pub use project::resolve_project; -pub use resolver::Resolver; pub use run::resolve_run; pub use server::{dev_token_auth_enabled, resolve_server}; pub use workflow::resolve_workflow; +use crate::apply_builtin_defaults; +use crate::user::default_storage_dir; + pub fn resolve_storage_root(file: &SettingsLayer) -> InterpString { - Resolver::from_layer(file).storage_root() + let layer = apply_builtin_defaults(file.clone()); + layer + .server + .as_ref() + .and_then(|server| server.storage.as_ref()) + .and_then(|storage| storage.root.clone()) + .unwrap_or_else(|| default_interp(default_storage_dir())) } pub fn resolve_cli_from_file(file: &SettingsLayer) -> Result> { - Resolver::from_layer(file).cli() + let layer = apply_builtin_defaults(file.clone()); + let mut errors = Vec::new(); + let value = resolve_cli(&layer.cli.clone().unwrap_or_default(), &mut errors); + finish(value, errors) } pub fn resolve_server_from_file( file: &SettingsLayer, ) -> Result> { - Resolver::from_layer(file).server() + let layer = apply_builtin_defaults(file.clone()); + let mut errors = Vec::new(); + let value = resolve_server(&layer.server.clone().unwrap_or_default(), &mut errors); + finish(value, errors) } pub fn resolve_project_from_file( file: &SettingsLayer, ) -> Result> { - Resolver::from_layer(file).project() + let layer = apply_builtin_defaults(file.clone()); + let mut errors = Vec::new(); + let value = resolve_project(&layer.project.clone().unwrap_or_default(), &mut errors); + finish(value, errors) } pub fn resolve_features_from_file( file: &SettingsLayer, ) -> Result> { - Resolver::from_layer(file).features() + let layer = apply_builtin_defaults(file.clone()); + let mut errors = Vec::new(); + let value = resolve_features(&layer.features.clone().unwrap_or_default(), &mut errors); + finish(value, errors) } pub fn resolve_run_from_file(file: &SettingsLayer) -> Result> { - Resolver::from_layer(file).run() + let layer = apply_builtin_defaults(file.clone()); + let mut errors = Vec::new(); + let value = resolve_run(&layer.run.clone().unwrap_or_default(), &mut errors); + finish(value, errors) } pub fn resolve_workflow_from_file( file: &SettingsLayer, ) -> Result> { - Resolver::from_layer(file).workflow() + let layer = apply_builtin_defaults(file.clone()); + let mut errors = Vec::new(); + let value = resolve_workflow(&layer.workflow.clone().unwrap_or_default(), &mut errors); + finish(value, errors) } /// Render a list of [`ResolveError`]s as a single semicolon-separated message @@ -101,6 +126,14 @@ pub(crate) fn default_interp(path: impl AsRef) -> InterpString InterpString::parse(&path.as_ref().to_string_lossy()) } +fn finish(value: T, errors: Vec) -> Result> { + if errors.is_empty() { + Ok(value) + } else { + Err(errors) + } +} + #[cfg(test)] mod tests { use std::collections::HashMap; diff --git a/lib/crates/fabro-config/src/resolve/resolver.rs b/lib/crates/fabro-config/src/resolve/resolver.rs deleted file mode 100644 index 893aa41ac..000000000 --- a/lib/crates/fabro-config/src/resolve/resolver.rs +++ /dev/null @@ -1,122 +0,0 @@ -//! Cache builtin defaults across multiple per-namespace resolutions. -//! -//! [`resolve_storage_root`] and the per-namespace `resolve_*_from_file` -//! helpers each call [`apply_builtin_defaults`], which clones both the input -//! layer and the embedded defaults layer before merging them. Callers that -//! need more than one namespace would otherwise pay that cost N times. -//! -//! [`Resolver`] applies defaults once on construction, then exposes per- -//! namespace methods that work against the materialized layer. It is the -//! shared backend for the standalone `resolve_*_from_file` helpers and the -//! preferred entrypoint when more than one namespace is needed. - -use fabro_types::settings::{ - CliNamespace, FeaturesNamespace, InterpString, ProjectNamespace, RunNamespace, ServerNamespace, - SettingsLayer, WorkflowNamespace, -}; - -use super::{ - ResolveError, default_interp, resolve_cli, resolve_features, resolve_project, resolve_run, - resolve_server, resolve_workflow, -}; -use crate::apply_builtin_defaults; -use crate::user::default_storage_dir; - -pub struct Resolver { - layer: SettingsLayer, -} - -impl Resolver { - #[must_use] - pub fn from_layer(layer: &SettingsLayer) -> Self { - Self { - layer: apply_builtin_defaults(layer.clone()), - } - } - - pub fn cli(&self) -> Result> { - let mut errors = Vec::new(); - let value = self.cli_into(&mut errors); - finish(value, errors) - } - - pub fn server(&self) -> Result> { - let mut errors = Vec::new(); - let value = self.server_into(&mut errors); - finish(value, errors) - } - - pub fn project(&self) -> Result> { - let mut errors = Vec::new(); - let value = self.project_into(&mut errors); - finish(value, errors) - } - - pub fn features(&self) -> Result> { - let mut errors = Vec::new(); - let value = self.features_into(&mut errors); - finish(value, errors) - } - - pub fn run(&self) -> Result> { - let mut errors = Vec::new(); - let value = self.run_into(&mut errors); - finish(value, errors) - } - - pub fn workflow(&self) -> Result> { - let mut errors = Vec::new(); - let value = self.workflow_into(&mut errors); - finish(value, errors) - } - - /// Resolved storage root, defaulting to [`default_storage_dir`] when the - /// input layer doesn't pin one. - #[must_use] - pub fn storage_root(&self) -> InterpString { - self.layer - .server - .as_ref() - .and_then(|server| server.storage.as_ref()) - .and_then(|storage| storage.root.clone()) - .unwrap_or_else(|| default_interp(default_storage_dir())) - } - - pub fn cli_into(&self, errors: &mut Vec) -> CliNamespace { - let layer = self.layer.cli.clone().unwrap_or_default(); - resolve_cli(&layer, errors) - } - - pub fn server_into(&self, errors: &mut Vec) -> ServerNamespace { - let layer = self.layer.server.clone().unwrap_or_default(); - resolve_server(&layer, errors) - } - - pub fn project_into(&self, errors: &mut Vec) -> ProjectNamespace { - let layer = self.layer.project.clone().unwrap_or_default(); - resolve_project(&layer, errors) - } - - pub fn features_into(&self, errors: &mut Vec) -> FeaturesNamespace { - let layer = self.layer.features.clone().unwrap_or_default(); - resolve_features(&layer, errors) - } - - pub fn run_into(&self, errors: &mut Vec) -> RunNamespace { - let layer = self.layer.run.clone().unwrap_or_default(); - resolve_run(&layer, errors) - } - - pub fn workflow_into(&self, errors: &mut Vec) -> WorkflowNamespace { - let layer = self.layer.workflow.clone().unwrap_or_default(); - resolve_workflow(&layer, errors) - } -} - -fn finish(value: T, errors: Vec) -> Result> { - if errors.is_empty() { - Ok(value) - } else { - Err(errors) - } -} diff --git a/lib/crates/fabro-config/tests/resolve_root.rs b/lib/crates/fabro-config/tests/resolve_root.rs index d4620053e..f0f79c275 100644 --- a/lib/crates/fabro-config/tests/resolve_root.rs +++ b/lib/crates/fabro-config/tests/resolve_root.rs @@ -1,4 +1,5 @@ use fabro_config::parse_settings_layer; +use fabro_types::settings::run::RunMode; use fabro_types::settings::{InterpString, SettingsLayer}; fn parse(source: &str) -> SettingsLayer { @@ -102,3 +103,92 @@ name = "gpt-5" Some("gpt-5".to_string()) ); } + +#[test] +fn workflow_settings_resolve_defaults_and_expose_fields() { + let settings = SettingsLayer::default(); + let resolved = + fabro_config::WorkflowSettings::from_layer(&settings).expect("defaults should resolve"); + + assert_eq!(resolved.project.directory, "."); + assert_eq!(resolved.workflow.graph, "workflow.fabro"); + assert_eq!(resolved.run.execution.mode, RunMode::Normal); +} + +#[test] +fn workflow_settings_combine_labels_with_later_namespaces_winning() { + let settings = parse( + r#" +_version = 1 + +[project.metadata] +project = "yes" +shared = "project" + +[workflow.metadata] +workflow = "yes" +shared = "workflow" + +[run.metadata] +run = "yes" +shared = "run" +"#, + ); + + let labels = fabro_config::WorkflowSettings::from_layer(&settings) + .expect("workflow settings should resolve") + .combined_labels(); + + assert_eq!(labels.get("project").map(String::as_str), Some("yes")); + assert_eq!(labels.get("workflow").map(String::as_str), Some("yes")); + assert_eq!(labels.get("run").map(String::as_str), Some("yes")); + assert_eq!(labels.get("shared").map(String::as_str), Some("run")); +} + +#[test] +fn workflow_settings_report_invalid_run_sandbox_provider() { + let settings = parse( + r#" +_version = 1 + +[run.sandbox] +provider = "not-a-provider" +"#, + ); + + let errors = fabro_config::WorkflowSettings::from_layer(&settings) + .expect_err("invalid workflow settings should fail"); + + assert!(errors.iter().any(|error| { + matches!( + error, + fabro_config::ResolveError::Invalid { path, .. } if path == "run.sandbox.provider" + ) + })); +} + +#[test] +fn workflow_settings_accumulate_multiple_run_errors() { + let settings = parse( + r#" +_version = 1 + +[run.sandbox] +provider = "not-a-provider" + +[[run.prepare.steps]] +script = "echo hi" +command = ["echo", "hi"] +"#, + ); + + let rendered = fabro_config::WorkflowSettings::from_layer(&settings) + .expect_err("invalid workflow settings should fail") + .into_iter() + .map(|error| error.to_string()) + .collect::>() + .join("\n"); + + assert!(rendered.contains("run.sandbox.provider")); + assert!(rendered.contains("run.prepare.steps[0]")); +} diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index fc3ee36f3..90267e99e 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -4016,7 +4016,23 @@ async fn create_run( create_input.provenance = Some(run_provenance(&headers, &subject)); create_input.submitted_manifest_bytes = Some(body.to_vec()); - let created = match Box::pin(operations::create(state.store.as_ref(), create_input)).await { + let storage_root = match resolve_interp_string(&state.server_settings().server.storage.root) { + Ok(path) => PathBuf::from(path), + Err(err) => { + return ApiError::new( + StatusCode::INTERNAL_SERVER_ERROR, + format!("Failed to resolve server storage root: {err}"), + ) + .into_response(); + } + }; + let created = match Box::pin(operations::create( + state.store.as_ref(), + create_input, + storage_root, + )) + .await + { Ok(created) => created, Err(WorkflowError::ValidationFailed { .. } | WorkflowError::Parse(_)) => { return ApiError::bad_request("Validation failed").into_response(); diff --git a/lib/crates/fabro-workflow/src/operations/create.rs b/lib/crates/fabro-workflow/src/operations/create.rs index 149e799d1..d2bd417e8 100644 --- a/lib/crates/fabro-workflow/src/operations/create.rs +++ b/lib/crates/fabro-workflow/src/operations/create.rs @@ -8,17 +8,15 @@ use std::collections::{BTreeMap, HashMap}; use std::path::{Path, PathBuf}; use std::sync::Arc; -use fabro_config::Storage; +use fabro_config::{Storage, WorkflowSettings}; use fabro_graphviz::graph::{AttrValue, Graph}; use fabro_model::{Catalog, Provider}; use fabro_sandbox::SandboxProvider; use fabro_sandbox::daytona::detect_repo_info; use fabro_store::Database; use fabro_template::{TemplateContext, render as render_template}; +use fabro_types::settings::SettingsLayer; use fabro_types::settings::run::RunMode; -use fabro_types::settings::{ - InterpString, ProjectNamespace, RunNamespace, SettingsLayer, WorkflowNamespace, -}; use fabro_types::{RunId, RunProvenance}; use fabro_util::json::normalize_json_value; use tokio::task::spawn_blocking; @@ -60,13 +58,6 @@ pub struct CreatedRun { pub dot_path: Option, } -struct ResolvedSettingsTree { - server_storage_root: InterpString, - project: ProjectNamespace, - workflow: WorkflowNamespace, - run: RunNamespace, -} - struct PersistCreateOptions { settings: SettingsLayer, run_id: Option, @@ -82,7 +73,11 @@ struct PersistCreateOptions { } /// Resolve workflow inputs, normalize settings, and persist a run directory. -pub async fn create(store: &Database, request: CreateRunInput) -> Result { +pub async fn create( + store: &Database, + request: CreateRunInput, + storage_root: PathBuf, +) -> Result { let resolved = resolve_workflow(ResolveWorkflowInput { workflow: request.workflow, settings: request.settings, @@ -113,18 +108,10 @@ pub async fn create(store: &Database, request: CreateRunInput) -> Result Result Error { Error::engine(err.to_string()) } -fn resolve_settings_tree(settings: &SettingsLayer) -> Result { - let resolver = fabro_config::Resolver::from_layer(settings); - let to_error = - |errors: Vec<_>| Error::Precondition(fabro_config::render_resolve_errors(&errors)); - Ok(ResolvedSettingsTree { - server_storage_root: resolver.storage_root(), - project: resolver.project().map_err(to_error)?, - workflow: resolver.workflow().map_err(to_error)?, - run: resolver.run().map_err(to_error)?, - }) -} - -fn combined_labels(settings: &ResolvedSettingsTree) -> HashMap { - let mut labels = settings.project.metadata.clone(); - labels.extend(settings.workflow.metadata.clone()); - labels.extend(settings.run.metadata.clone()); - labels -} - fn validate_sandbox_provider(settings: &SettingsLayer) -> Result<(), Error> { let resolved = fabro_config::resolve_run_from_file(settings) .map_err(|errors| Error::Precondition(fabro_config::render_resolve_errors(&errors)))?; @@ -735,25 +703,30 @@ mod tests { work [label="Work"] }"#; let dir = tempfile::tempdir().unwrap(); + let storage_root = dir.path().join("storage"); let store = memory_store(); - let err = create(&store, CreateRunInput { - workflow: WorkflowInput::DotSource { - source: dot.to_string(), - base_dir: None, + let err = create( + &store, + CreateRunInput { + workflow: WorkflowInput::DotSource { + source: dot.to_string(), + base_dir: None, + }, + settings: test_default_settings(), + cwd: dir.path().to_path_buf(), + workflow_slug: None, + workflow_path: None, + workflow_bundle: None, + submitted_manifest_bytes: None, + run_id: None, + host_repo_path: None, + repo_origin_url: None, + base_branch: None, + provenance: None, + configured_providers: Vec::new(), }, - settings: test_default_settings(), - cwd: dir.path().to_path_buf(), - workflow_slug: None, - workflow_path: None, - workflow_bundle: None, - submitted_manifest_bytes: None, - run_id: None, - host_repo_path: None, - repo_origin_url: None, - base_branch: None, - provenance: None, - configured_providers: Vec::new(), - }) + storage_root, + ) .await .unwrap_err(); @@ -765,57 +738,121 @@ mod tests { } } + #[tokio::test] + async fn create_reports_workflow_settings_errors_with_rendered_message() { + let dir = tempfile::tempdir().unwrap(); + let storage_root = dir.path().join("storage"); + let store = memory_store(); + let err = create( + &store, + CreateRunInput { + workflow: WorkflowInput::DotSource { + source: MINIMAL_DOT.to_string(), + base_dir: None, + }, + settings: { + use fabro_types::settings::run::{ + RunExecutionLayer, RunLayer, RunMode, RunSandboxLayer, + }; + let mut layer = SettingsLayer { + run: Some(RunLayer { + execution: Some(RunExecutionLayer { + mode: Some(RunMode::DryRun), + ..RunExecutionLayer::default() + }), + sandbox: Some(RunSandboxLayer { + provider: Some("not-a-provider".to_string()), + ..RunSandboxLayer::default() + }), + ..RunLayer::default() + }), + ..SettingsLayer::default() + }; + layer.ensure_test_auth_methods(); + layer + }, + cwd: dir.path().to_path_buf(), + workflow_slug: None, + workflow_path: None, + workflow_bundle: None, + submitted_manifest_bytes: None, + run_id: None, + host_repo_path: None, + repo_origin_url: None, + base_branch: None, + provenance: None, + configured_providers: Vec::new(), + }, + storage_root, + ) + .await + .unwrap_err(); + + match err { + Error::Precondition(message) => { + assert!(message.contains("run.sandbox.provider")); + assert!(!message.contains('\n')); + } + other => panic!("expected Precondition, got {other:?}"), + } + } + #[tokio::test] async fn create_persists_normalized_config_and_initial_state() { let dir = tempfile::tempdir().unwrap(); + let storage_root = dir.path().join("storage"); let store = memory_store(); - let created = create(&store, CreateRunInput { - workflow: WorkflowInput::DotSource { - source: MINIMAL_DOT.to_string(), - base_dir: None, + let created = create( + &store, + CreateRunInput { + workflow: WorkflowInput::DotSource { + source: MINIMAL_DOT.to_string(), + base_dir: None, + }, + settings: { + use fabro_types::settings::run::{ + RunExecutionLayer, RunGoalLayer, RunLayer, RunMode, RunModelLayer, + RunPullRequestLayer, + }; + let mut metadata = HashMap::new(); + metadata.insert("env".to_string(), "test".to_string()); + let mut layer = SettingsLayer { + run: Some(RunLayer { + goal: Some(RunGoalLayer::Inline(InterpString::parse("override goal"))), + metadata, + model: Some(RunModelLayer { + name: Some(InterpString::parse("sonnet")), + ..RunModelLayer::default() + }), + pull_request: Some(RunPullRequestLayer { + enabled: Some(false), + ..RunPullRequestLayer::default() + }), + execution: Some(RunExecutionLayer { + mode: Some(RunMode::DryRun), + ..RunExecutionLayer::default() + }), + ..RunLayer::default() + }), + ..SettingsLayer::default() + }; + layer.ensure_test_auth_methods(); + layer + }, + cwd: dir.path().to_path_buf(), + workflow_slug: Some("slug".to_string()), + workflow_path: None, + workflow_bundle: None, + submitted_manifest_bytes: None, + run_id: Some(fixtures::RUN_1), + host_repo_path: Some(dir.path().display().to_string()), + repo_origin_url: None, + base_branch: Some("main".to_string()), + provenance: None, + configured_providers: Vec::new(), }, - settings: { - use fabro_types::settings::run::{ - RunExecutionLayer, RunGoalLayer, RunLayer, RunMode, RunModelLayer, - RunPullRequestLayer, - }; - let mut metadata = HashMap::new(); - metadata.insert("env".to_string(), "test".to_string()); - let mut layer = SettingsLayer { - run: Some(RunLayer { - goal: Some(RunGoalLayer::Inline(InterpString::parse("override goal"))), - metadata, - model: Some(RunModelLayer { - name: Some(InterpString::parse("sonnet")), - ..RunModelLayer::default() - }), - pull_request: Some(RunPullRequestLayer { - enabled: Some(false), - ..RunPullRequestLayer::default() - }), - execution: Some(RunExecutionLayer { - mode: Some(RunMode::DryRun), - ..RunExecutionLayer::default() - }), - ..RunLayer::default() - }), - ..SettingsLayer::default() - }; - layer.ensure_test_auth_methods(); - layer - }, - cwd: dir.path().to_path_buf(), - workflow_slug: Some("slug".to_string()), - workflow_path: None, - workflow_bundle: None, - submitted_manifest_bytes: None, - run_id: Some(fixtures::RUN_1), - host_repo_path: Some(dir.path().display().to_string()), - repo_origin_url: None, - base_branch: Some("main".to_string()), - provenance: None, - configured_providers: Vec::new(), - }) + storage_root.clone(), + ) .await .unwrap(); @@ -869,7 +906,13 @@ mod tests { run_store.state().await.unwrap().status.unwrap(), crate::run_status::RunStatus::Submitted ); - assert_eq!(created.run_dir, default_run_dir(&fixtures::RUN_1)); + assert_eq!( + created.run_dir, + Storage::new(&storage_root) + .run_scratch(&fixtures::RUN_1) + .root() + .to_path_buf() + ); assert!(created.run_dir.is_dir()); } @@ -878,41 +921,46 @@ mod tests { let dir = tempfile::tempdir().unwrap(); let workspace = dir.path().join("workspace"); std::fs::create_dir_all(&workspace).unwrap(); + let storage_root = dir.path().join("storage"); let store = memory_store(); - let created = create(&store, CreateRunInput { - workflow: WorkflowInput::DotSource { - source: MINIMAL_DOT.to_string(), - base_dir: None, - }, - settings: { - use fabro_types::settings::run::{RunExecutionLayer, RunLayer, RunMode}; - let mut layer = SettingsLayer { - run: Some(RunLayer { - working_dir: Some(InterpString::parse("workspace")), - execution: Some(RunExecutionLayer { - mode: Some(RunMode::DryRun), - ..RunExecutionLayer::default() + let created = create( + &store, + CreateRunInput { + workflow: WorkflowInput::DotSource { + source: MINIMAL_DOT.to_string(), + base_dir: None, + }, + settings: { + use fabro_types::settings::run::{RunExecutionLayer, RunLayer, RunMode}; + let mut layer = SettingsLayer { + run: Some(RunLayer { + working_dir: Some(InterpString::parse("workspace")), + execution: Some(RunExecutionLayer { + mode: Some(RunMode::DryRun), + ..RunExecutionLayer::default() + }), + ..RunLayer::default() }), - ..RunLayer::default() - }), - ..SettingsLayer::default() - }; - layer.ensure_test_auth_methods(); - layer + ..SettingsLayer::default() + }; + layer.ensure_test_auth_methods(); + layer + }, + cwd: dir.path().to_path_buf(), + workflow_slug: None, + workflow_path: None, + workflow_bundle: None, + submitted_manifest_bytes: None, + run_id: Some(fixtures::RUN_2), + host_repo_path: None, + repo_origin_url: None, + base_branch: None, + provenance: None, + configured_providers: Vec::new(), }, - cwd: dir.path().to_path_buf(), - workflow_slug: None, - workflow_path: None, - workflow_bundle: None, - submitted_manifest_bytes: None, - run_id: Some(fixtures::RUN_2), - host_repo_path: None, - repo_origin_url: None, - base_branch: None, - provenance: None, - configured_providers: Vec::new(), - }) + storage_root, + ) .await .unwrap(); @@ -933,25 +981,30 @@ mod tests { #[tokio::test] async fn create_persists_repo_origin_url_from_request() { let dir = tempfile::tempdir().unwrap(); + let storage_root = dir.path().join("storage"); let store = memory_store(); - let created = create(&store, CreateRunInput { - workflow: WorkflowInput::DotSource { - source: MINIMAL_DOT.to_string(), - base_dir: None, + let created = create( + &store, + CreateRunInput { + workflow: WorkflowInput::DotSource { + source: MINIMAL_DOT.to_string(), + base_dir: None, + }, + settings: dry_run_only_settings(), + cwd: dir.path().to_path_buf(), + workflow_slug: None, + workflow_path: None, + workflow_bundle: None, + submitted_manifest_bytes: None, + run_id: Some(fixtures::RUN_2), + host_repo_path: None, + repo_origin_url: Some("https://github.com/acme/widgets".to_string()), + base_branch: None, + provenance: None, + configured_providers: Vec::new(), }, - settings: dry_run_only_settings(), - cwd: dir.path().to_path_buf(), - workflow_slug: None, - workflow_path: None, - workflow_bundle: None, - submitted_manifest_bytes: None, - run_id: Some(fixtures::RUN_2), - host_repo_path: None, - repo_origin_url: Some("https://github.com/acme/widgets".to_string()), - base_branch: None, - provenance: None, - configured_providers: Vec::new(), - }) + storage_root, + ) .await .unwrap(); @@ -1013,24 +1066,28 @@ mod tests { Duration::from_millis(1), None, )); - let created = create(store.as_ref(), CreateRunInput { - workflow: WorkflowInput::DotSource { - source: MINIMAL_DOT.to_string(), - base_dir: None, + let created = create( + store.as_ref(), + CreateRunInput { + workflow: WorkflowInput::DotSource { + source: MINIMAL_DOT.to_string(), + base_dir: None, + }, + settings: dry_run_with_storage(&storage_dir), + cwd: dir.path().to_path_buf(), + workflow_slug: Some("slug".to_string()), + workflow_path: None, + workflow_bundle: None, + submitted_manifest_bytes: None, + run_id: Some(fixtures::RUN_3), + host_repo_path: None, + repo_origin_url: None, + base_branch: None, + provenance: None, + configured_providers: Vec::new(), }, - settings: dry_run_with_storage(&storage_dir), - cwd: dir.path().to_path_buf(), - workflow_slug: Some("slug".to_string()), - workflow_path: None, - workflow_bundle: None, - submitted_manifest_bytes: None, - run_id: Some(fixtures::RUN_3), - host_repo_path: None, - repo_origin_url: None, - base_branch: None, - provenance: None, - configured_providers: Vec::new(), - }) + storage_dir.clone(), + ) .await .unwrap(); let run_store = store.open_run_reader(&created.run_id).await.unwrap(); @@ -1052,37 +1109,41 @@ mod tests { Duration::from_millis(1), None, )); - let created = create(store.as_ref(), CreateRunInput { - workflow: WorkflowInput::DotSource { - source: MINIMAL_DOT.to_string(), - base_dir: None, + let created = create( + store.as_ref(), + CreateRunInput { + workflow: WorkflowInput::DotSource { + source: MINIMAL_DOT.to_string(), + base_dir: None, + }, + settings: dry_run_with_storage(&storage_dir), + cwd: dir.path().to_path_buf(), + workflow_slug: Some("slug".to_string()), + workflow_path: None, + workflow_bundle: None, + submitted_manifest_bytes: None, + run_id: Some(fixtures::RUN_64), + host_repo_path: None, + repo_origin_url: None, + base_branch: None, + provenance: Some(fabro_types::RunProvenance { + server: Some(fabro_types::RunServerProvenance { + version: "0.9.0".to_string(), + }), + client: Some(fabro_types::RunClientProvenance { + user_agent: Some("fabro-cli/0.9.0".to_string()), + name: Some("fabro-cli".to_string()), + version: Some("0.9.0".to_string()), + }), + subject: Some(fabro_types::RunSubjectProvenance { + login: None, + auth_method: fabro_types::RunAuthMethod::Disabled, + }), + }), + configured_providers: Vec::new(), }, - settings: dry_run_with_storage(&storage_dir), - cwd: dir.path().to_path_buf(), - workflow_slug: Some("slug".to_string()), - workflow_path: None, - workflow_bundle: None, - submitted_manifest_bytes: None, - run_id: Some(fixtures::RUN_64), - host_repo_path: None, - repo_origin_url: None, - base_branch: None, - provenance: Some(fabro_types::RunProvenance { - server: Some(fabro_types::RunServerProvenance { - version: "0.9.0".to_string(), - }), - client: Some(fabro_types::RunClientProvenance { - user_agent: Some("fabro-cli/0.9.0".to_string()), - name: Some("fabro-cli".to_string()), - version: Some("0.9.0".to_string()), - }), - subject: Some(fabro_types::RunSubjectProvenance { - login: None, - auth_method: fabro_types::RunAuthMethod::Disabled, - }), - }), - configured_providers: Vec::new(), - }) + storage_dir, + ) .await .unwrap(); diff --git a/lib/crates/fabro-workflow/src/operations/start.rs b/lib/crates/fabro-workflow/src/operations/start.rs index 7a087b34a..3a5df431f 100644 --- a/lib/crates/fabro-workflow/src/operations/start.rs +++ b/lib/crates/fabro-workflow/src/operations/start.rs @@ -1003,42 +1003,55 @@ mod tests { )) } - async fn persisted_workflow(dot: &str, run_dir: &Path) -> (Persisted, Arc) { + fn storage_root_and_run_dir(temp: &tempfile::TempDir) -> (PathBuf, PathBuf) { + let storage_root = temp.path().join("storage"); + let run_dir = fabro_config::Storage::new(&storage_root) + .run_scratch(&fixtures::RUN_1) + .root() + .to_path_buf(); + (storage_root, run_dir) + } + + async fn persisted_workflow(dot: &str, storage_root: &Path) -> (Persisted, Arc) { let store = memory_store(); - let created = crate::operations::create(&store, crate::operations::CreateRunInput { - workflow: crate::operations::WorkflowInput::DotSource { - source: dot.to_string(), - base_dir: None, - }, - settings: { - let mut layer = SettingsLayer { - run: Some(RunLayer { - execution: Some(RunExecutionLayer { - mode: Some(RunMode::DryRun), - ..RunExecutionLayer::default() + let created = crate::operations::create( + &store, + crate::operations::CreateRunInput { + workflow: crate::operations::WorkflowInput::DotSource { + source: dot.to_string(), + base_dir: None, + }, + settings: { + let mut layer = SettingsLayer { + run: Some(RunLayer { + execution: Some(RunExecutionLayer { + mode: Some(RunMode::DryRun), + ..RunExecutionLayer::default() + }), + ..RunLayer::default() }), - ..RunLayer::default() - }), - ..SettingsLayer::default() - }; - layer.ensure_test_auth_methods(); - layer + ..SettingsLayer::default() + }; + layer.ensure_test_auth_methods(); + layer + }, + cwd: storage_root + .parent() + .unwrap_or_else(|| Path::new(".")) + .to_path_buf(), + workflow_slug: Some("test".to_string()), + workflow_path: None, + workflow_bundle: None, + submitted_manifest_bytes: None, + run_id: Some(fixtures::RUN_1), + host_repo_path: None, + repo_origin_url: None, + base_branch: None, + provenance: None, + configured_providers: Vec::new(), }, - cwd: run_dir - .parent() - .unwrap_or_else(|| Path::new(".")) - .to_path_buf(), - workflow_slug: Some("test".to_string()), - workflow_path: None, - workflow_bundle: None, - submitted_manifest_bytes: None, - run_id: Some(fixtures::RUN_1), - host_repo_path: None, - repo_origin_url: None, - base_branch: None, - provenance: None, - configured_providers: Vec::new(), - }) + storage_root.to_path_buf(), + ) .await .unwrap(); (created.persisted, store) @@ -1078,7 +1091,7 @@ mod tests { #[tokio::test] async fn start_captures_checkpoint_git_sha_in_conclusion() { let temp = tempfile::tempdir().unwrap(); - let run_dir = temp.path().join("run"); + let (storage_root, run_dir) = storage_root_and_run_dir(&temp); let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); let registry = Arc::new(test_registry()); let injected = Arc::new(AtomicBool::new(false)); @@ -1113,7 +1126,7 @@ mod tests { }); } - let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &run_dir).await; + let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &storage_root).await; let started = start( &run_dir, test_start_services(&store, &run_dir, emitter, registry).await, @@ -1132,11 +1145,11 @@ mod tests { #[tokio::test] async fn start_loads_persisted_from_run_dir() { let temp = tempfile::tempdir().unwrap(); - let run_dir = temp.path().join("run"); + let (storage_root, run_dir) = storage_root_and_run_dir(&temp); let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); let registry = Arc::new(test_registry()); - let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &run_dir).await; + let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &storage_root).await; let started = start( &run_dir, @@ -1153,7 +1166,7 @@ mod tests { #[tokio::test] async fn start_can_run_bundle_backed_child_workflow_without_workflow_bundle_json() { let temp = tempfile::tempdir().unwrap(); - let run_dir = temp.path().join("run"); + let (storage_root, run_dir) = storage_root_and_run_dir(&temp); let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); let registry = Arc::new(test_registry()); let store = memory_store(); @@ -1187,39 +1200,43 @@ mod tests { }), ])); - crate::operations::create(&store, crate::operations::CreateRunInput { - workflow: crate::operations::WorkflowInput::Bundled( - workflow_bundle - .workflow(Path::new("workflow.fabro")) - .unwrap() - .clone(), - ), - settings: { - let mut layer = SettingsLayer { - run: Some(RunLayer { - execution: Some(RunExecutionLayer { - mode: Some(RunMode::DryRun), - ..RunExecutionLayer::default() + crate::operations::create( + &store, + crate::operations::CreateRunInput { + workflow: crate::operations::WorkflowInput::Bundled( + workflow_bundle + .workflow(Path::new("workflow.fabro")) + .unwrap() + .clone(), + ), + settings: { + let mut layer = SettingsLayer { + run: Some(RunLayer { + execution: Some(RunExecutionLayer { + mode: Some(RunMode::DryRun), + ..RunExecutionLayer::default() + }), + ..RunLayer::default() }), - ..RunLayer::default() - }), - ..SettingsLayer::default() - }; - layer.ensure_test_auth_methods(); - layer + ..SettingsLayer::default() + }; + layer.ensure_test_auth_methods(); + layer + }, + cwd: temp.path().to_path_buf(), + workflow_slug: Some("bundle-child".to_string()), + workflow_path: Some(PathBuf::from("workflow.fabro")), + workflow_bundle: Some(workflow_bundle), + submitted_manifest_bytes: None, + run_id: Some(fixtures::RUN_1), + host_repo_path: None, + repo_origin_url: None, + base_branch: None, + provenance: None, + configured_providers: Vec::new(), }, - cwd: temp.path().to_path_buf(), - workflow_slug: Some("bundle-child".to_string()), - workflow_path: Some(PathBuf::from("workflow.fabro")), - workflow_bundle: Some(workflow_bundle), - submitted_manifest_bytes: None, - run_id: Some(fixtures::RUN_1), - host_repo_path: None, - repo_origin_url: None, - base_branch: None, - provenance: None, - configured_providers: Vec::new(), - }) + storage_root, + ) .await .unwrap(); @@ -1236,12 +1253,12 @@ mod tests { #[tokio::test] async fn start_invokes_on_node_callback_before_execution() { let temp = tempfile::tempdir().unwrap(); - let run_dir = temp.path().join("run"); + let (storage_root, run_dir) = storage_root_and_run_dir(&temp); let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); let registry = Arc::new(test_registry()); let visited = Arc::new(Mutex::new(Vec::new())); - let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &run_dir).await; + let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &storage_root).await; let started = start(&run_dir, StartServices { on_node: Some(Arc::new({ @@ -1262,11 +1279,11 @@ mod tests { #[tokio::test] async fn start_errors_when_checkpoint_exists() { let temp = tempfile::tempdir().unwrap(); - let run_dir = temp.path().join("run"); + let (storage_root, run_dir) = storage_root_and_run_dir(&temp); let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); let registry = Arc::new(test_registry()); - let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &run_dir).await; + let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &storage_root).await; let services = test_start_services(&store, &run_dir, emitter, registry).await; // Seed an authoritative checkpoint event so start() sees it @@ -1331,11 +1348,11 @@ mod tests { #[tokio::test] async fn resume_errors_when_checkpoint_missing() { let temp = tempfile::tempdir().unwrap(); - let run_dir = temp.path().join("run"); + let (storage_root, run_dir) = storage_root_and_run_dir(&temp); let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); let registry = Arc::new(test_registry()); - let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &run_dir).await; + let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &storage_root).await; let result = resume( &run_dir, @@ -1353,12 +1370,12 @@ mod tests { #[tokio::test] async fn resume_errors_when_run_already_finished_successfully() { let temp = tempfile::tempdir().unwrap(); - let run_dir = temp.path().join("run"); + let (storage_root, run_dir) = storage_root_and_run_dir(&temp); std::fs::create_dir_all(&run_dir).unwrap(); let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); let registry = Arc::new(test_registry()); - let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &run_dir).await; + let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &storage_root).await; let checkpoint = Checkpoint::from_context( &Context::new(),