mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-10 03:30:59 +00:00
refactor(config): collapse settings resolver
This commit is contained in:
parent
5041669d6d
commit
f5758d99d0
9 changed files with 565 additions and 429 deletions
|
|
@ -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.
|
||||
|
||||
|
|
|
|||
|
|
@ -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<Self> {
|
||||
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<Self> {
|
||||
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<Self, Vec<ResolveError>> {
|
||||
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<String, String> {
|
||||
let mut labels = self.project.metadata.clone();
|
||||
labels.extend(self.workflow.metadata.clone());
|
||||
labels.extend(self.run.metadata.clone());
|
||||
labels
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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<CliNamespace, Vec<ResolveError>> {
|
||||
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<ServerNamespace, Vec<ResolveError>> {
|
||||
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<ProjectNamespace, Vec<ResolveError>> {
|
||||
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<FeaturesNamespace, Vec<ResolveError>> {
|
||||
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<RunNamespace, Vec<ResolveError>> {
|
||||
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<WorkflowNamespace, Vec<ResolveError>> {
|
||||
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<std::path::Path>) -> InterpString
|
|||
InterpString::parse(&path.as_ref().to_string_lossy())
|
||||
}
|
||||
|
||||
fn finish<T>(value: T, errors: Vec<ResolveError>) -> Result<T, Vec<ResolveError>> {
|
||||
if errors.is_empty() {
|
||||
Ok(value)
|
||||
} else {
|
||||
Err(errors)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::collections::HashMap;
|
||||
|
|
|
|||
|
|
@ -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<CliNamespace, Vec<ResolveError>> {
|
||||
let mut errors = Vec::new();
|
||||
let value = self.cli_into(&mut errors);
|
||||
finish(value, errors)
|
||||
}
|
||||
|
||||
pub fn server(&self) -> Result<ServerNamespace, Vec<ResolveError>> {
|
||||
let mut errors = Vec::new();
|
||||
let value = self.server_into(&mut errors);
|
||||
finish(value, errors)
|
||||
}
|
||||
|
||||
pub fn project(&self) -> Result<ProjectNamespace, Vec<ResolveError>> {
|
||||
let mut errors = Vec::new();
|
||||
let value = self.project_into(&mut errors);
|
||||
finish(value, errors)
|
||||
}
|
||||
|
||||
pub fn features(&self) -> Result<FeaturesNamespace, Vec<ResolveError>> {
|
||||
let mut errors = Vec::new();
|
||||
let value = self.features_into(&mut errors);
|
||||
finish(value, errors)
|
||||
}
|
||||
|
||||
pub fn run(&self) -> Result<RunNamespace, Vec<ResolveError>> {
|
||||
let mut errors = Vec::new();
|
||||
let value = self.run_into(&mut errors);
|
||||
finish(value, errors)
|
||||
}
|
||||
|
||||
pub fn workflow(&self) -> Result<WorkflowNamespace, Vec<ResolveError>> {
|
||||
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<ResolveError>) -> CliNamespace {
|
||||
let layer = self.layer.cli.clone().unwrap_or_default();
|
||||
resolve_cli(&layer, errors)
|
||||
}
|
||||
|
||||
pub fn server_into(&self, errors: &mut Vec<ResolveError>) -> ServerNamespace {
|
||||
let layer = self.layer.server.clone().unwrap_or_default();
|
||||
resolve_server(&layer, errors)
|
||||
}
|
||||
|
||||
pub fn project_into(&self, errors: &mut Vec<ResolveError>) -> ProjectNamespace {
|
||||
let layer = self.layer.project.clone().unwrap_or_default();
|
||||
resolve_project(&layer, errors)
|
||||
}
|
||||
|
||||
pub fn features_into(&self, errors: &mut Vec<ResolveError>) -> FeaturesNamespace {
|
||||
let layer = self.layer.features.clone().unwrap_or_default();
|
||||
resolve_features(&layer, errors)
|
||||
}
|
||||
|
||||
pub fn run_into(&self, errors: &mut Vec<ResolveError>) -> RunNamespace {
|
||||
let layer = self.layer.run.clone().unwrap_or_default();
|
||||
resolve_run(&layer, errors)
|
||||
}
|
||||
|
||||
pub fn workflow_into(&self, errors: &mut Vec<ResolveError>) -> WorkflowNamespace {
|
||||
let layer = self.layer.workflow.clone().unwrap_or_default();
|
||||
resolve_workflow(&layer, errors)
|
||||
}
|
||||
}
|
||||
|
||||
fn finish<T>(value: T, errors: Vec<ResolveError>) -> Result<T, Vec<ResolveError>> {
|
||||
if errors.is_empty() {
|
||||
Ok(value)
|
||||
} else {
|
||||
Err(errors)
|
||||
}
|
||||
}
|
||||
|
|
@ -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::<Vec<_>>()
|
||||
.join("\n");
|
||||
|
||||
assert!(rendered.contains("run.sandbox.provider"));
|
||||
assert!(rendered.contains("run.prepare.steps[0]"));
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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<PathBuf>,
|
||||
}
|
||||
|
||||
struct ResolvedSettingsTree {
|
||||
server_storage_root: InterpString,
|
||||
project: ProjectNamespace,
|
||||
workflow: WorkflowNamespace,
|
||||
run: RunNamespace,
|
||||
}
|
||||
|
||||
struct PersistCreateOptions {
|
||||
settings: SettingsLayer,
|
||||
run_id: Option<RunId>,
|
||||
|
|
@ -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<CreatedRun, Error> {
|
||||
pub async fn create(
|
||||
store: &Database,
|
||||
request: CreateRunInput,
|
||||
storage_root: PathBuf,
|
||||
) -> Result<CreatedRun, Error> {
|
||||
let resolved = resolve_workflow(ResolveWorkflowInput {
|
||||
workflow: request.workflow,
|
||||
settings: request.settings,
|
||||
|
|
@ -113,18 +108,10 @@ pub async fn create(store: &Database, request: CreateRunInput) -> Result<Created
|
|||
} = request;
|
||||
|
||||
let settings = resolved.settings.clone();
|
||||
let resolved_settings = resolve_settings_tree(&settings)?;
|
||||
let resolved_settings = WorkflowSettings::from_layer(&settings)
|
||||
.map_err(|errors| Error::Precondition(fabro_config::render_resolve_errors(&errors)))?;
|
||||
let run_id = run_id.unwrap_or_else(RunId::new);
|
||||
let storage_root = resolved_settings
|
||||
.server_storage_root
|
||||
.resolve(|name| std::env::var(name).ok())
|
||||
.map_err(|err| {
|
||||
Error::Precondition(format!(
|
||||
"failed to resolve {}: {err}",
|
||||
resolved_settings.server_storage_root.as_source()
|
||||
))
|
||||
})?;
|
||||
let storage = Storage::new(storage_root.value);
|
||||
let storage = Storage::new(storage_root);
|
||||
let run_dir = storage.run_scratch(&run_id).root().to_path_buf();
|
||||
let working_directory = resolved.working_directory.clone();
|
||||
let host_repo_path =
|
||||
|
|
@ -163,7 +150,7 @@ pub async fn create(store: &Database, request: CreateRunInput) -> Result<Created
|
|||
run_id: Some(run_id),
|
||||
run_dir: Some(persisted_run_dir),
|
||||
workflow_slug: workflow_slug.or(resolved_workflow_slug),
|
||||
labels: combined_labels(&resolved_settings),
|
||||
labels: resolved_settings.combined_labels(),
|
||||
base_branch,
|
||||
working_directory,
|
||||
host_repo_path,
|
||||
|
|
@ -284,25 +271,6 @@ fn store_error(err: impl std::fmt::Display) -> Error {
|
|||
Error::engine(err.to_string())
|
||||
}
|
||||
|
||||
fn resolve_settings_tree(settings: &SettingsLayer) -> Result<ResolvedSettingsTree, Error> {
|
||||
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<String, String> {
|
||||
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();
|
||||
|
||||
|
|
|
|||
|
|
@ -1003,42 +1003,55 @@ mod tests {
|
|||
))
|
||||
}
|
||||
|
||||
async fn persisted_workflow(dot: &str, run_dir: &Path) -> (Persisted, Arc<Database>) {
|
||||
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<Database>) {
|
||||
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(),
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue