Simplify automation workflow-version packaging

Share the RunIntent shape between the scheduler and the API trigger via
AutomationRunMaterialized::into_run_intent, drop the pass-through
packaging wrappers and the unreachable VersionIdMismatch error, and move
the config-path and version-ID derivations onto WorkflowVersion so the
server, validator, and collector stop re-deriving them.

The collector now owns the collected sources (moving file contents
instead of cloning them), resolves the workflow location once, and shares
the not-found probe with build_run_manifest. The bundler reads a goal
file once instead of twice.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-08-27 16:29:45 -04:00
parent e86dd3bea4
commit 077d020223
9 changed files with 147 additions and 231 deletions

View file

@ -3,9 +3,10 @@ use std::sync::Arc;
use async_trait::async_trait;
use fabro_automation::AutomationId;
use fabro_manifest::{CollectedWorkflowClosure, WorkflowVersionCollectError};
use fabro_manifest::WorkflowVersionCollectError;
use fabro_types::{
GitHubRepositorySlug, GitRunTarget, RunId, TargetValidationError, WorkflowVersionId,
GitHubRepositorySlug, GitRunTarget, RunId, RunIntent, RunIntentArgs, RunTarget,
TargetValidationError, WorkflowVersionId,
};
use fabro_workflow_version::{WorkflowVersionStore, WorkflowVersionStoreError};
use tokio::{fs, task};
@ -29,6 +30,22 @@ pub(crate) struct AutomationRunMaterialized {
pub target: GitRunTarget,
}
impl AutomationRunMaterialized {
/// The admission request for an automation run: the packaged workflow
/// version at the exact checked-out target, with no caller overrides.
pub(crate) fn into_run_intent(self) -> RunIntent {
RunIntent {
workflow_version_id: self.workflow_version_id,
target: RunTarget::Git(self.target),
args: RunIntentArgs::default(),
environment_id: None,
parent_id: None,
title: None,
goal: None,
}
}
}
#[derive(thiserror::Error, Debug)]
pub(crate) enum RunMaterializeError {
#[error("invalid automation Git target")]
@ -67,11 +84,6 @@ pub(crate) enum RunMaterializeError {
#[source]
source: WorkflowVersionStoreError,
},
#[error("stored workflow version ID `{actual}` did not match local ID `{expected}`")]
VersionIdMismatch {
expected: WorkflowVersionId,
actual: WorkflowVersionId,
},
#[error("failed to load GitHub credentials")]
Credentials {
#[source]
@ -165,59 +177,34 @@ impl AutomationRunMaterializer for ProductionAutomationRunMaterializer {
let mut exact_target = input.target;
exact_target.sha = Some(checked_out_sha);
let package_input = PackageFromCheckoutInput {
workflow: input.workflow,
checkout_dir,
target: exact_target,
};
let packaged = task::spawn_blocking(move || package_from_checkout(package_input))
.await
.map_err(|source| RunMaterializeError::PackageTask { source })??;
let workflow = PathBuf::from(input.workflow);
let closure = task::spawn_blocking(move || {
fabro_manifest::collect_workflow_versions(&workflow, &checkout_dir)
.map_err(package_error)
})
.await
.map_err(|source| RunMaterializeError::PackageTask { source })??;
let root_id = packaged.closure.root_id();
for (expected, version) in packaged.closure.versions() {
let actual = self
// Versions arrive dependency-first, and the store derives the same
// content hash the collector did, so the root ID is known up front.
for (expected, version) in closure.versions() {
let stored = self
.version_store
.put(version)
.await
.map_err(|source| RunMaterializeError::VersionStore { source })?;
if actual != expected {
return Err(RunMaterializeError::VersionIdMismatch { expected, actual });
}
debug_assert_eq!(
stored, expected,
"store and collector disagree on version ID"
);
}
Ok(AutomationRunMaterialized {
workflow_version_id: root_id,
target: packaged.target,
workflow_version_id: closure.root_id(),
target: exact_target,
})
}
}
#[derive(Debug)]
struct PackageFromCheckoutInput {
workflow: String,
checkout_dir: PathBuf,
target: GitRunTarget,
}
struct PackagedAutomationWorkflow {
closure: CollectedWorkflowClosure,
target: GitRunTarget,
}
fn package_from_checkout(
args: PackageFromCheckoutInput,
) -> Result<PackagedAutomationWorkflow, RunMaterializeError> {
let PackageFromCheckoutInput {
workflow,
checkout_dir,
target,
} = args;
let workflow = PathBuf::from(workflow);
let closure = fabro_manifest::collect_workflow_versions(&workflow, &checkout_dir)
.map_err(package_error)?;
Ok(PackagedAutomationWorkflow { closure, target })
}
fn package_error(error: WorkflowVersionCollectError) -> RunMaterializeError {
if matches!(&error, WorkflowVersionCollectError::WorkflowNotFound { .. }) {
RunMaterializeError::WorkflowNotFound { source: error }
@ -345,12 +332,11 @@ impl AutomationRunMaterializer for TestAutomationRunMaterializer {
.await
.map_err(|source| RunMaterializeError::VersionStore { source })?
} else {
let canonical = materialized
materialized
.version
.version()
.canonical_bytes()
.expect("validated test workflow version should serialize canonically");
WorkflowVersionId::from(fabro_types::BlobHash::new(&canonical))
.id()
.expect("validated test workflow version should serialize canonically")
};
Ok(AutomationRunMaterialized {
workflow_version_id,
@ -367,6 +353,7 @@ mod tests {
)]
use std::fs;
use std::path::Path;
use std::time::Duration;
use object_store::memory::InMemory;
@ -374,53 +361,6 @@ mod tests {
use super::*;
#[test]
fn package_builder_uses_checkout_relative_versions_and_exact_target() {
let temp = TempDir::new().unwrap();
let checkout = temp.path().join("checkout");
let workflow_dir = checkout.join(".fabro/workflows/demo");
fs::create_dir_all(&workflow_dir).unwrap();
fs::write(checkout.join(".fabro/project.toml"), "_version = 1\n").unwrap();
fs::write(
workflow_dir.join("workflow.fabro"),
r#"digraph Demo { graph [goal="Ship automation"] start [shape=Mdiamond] exit [shape=Msquare] start -> exit }"#,
)
.unwrap();
fs::write(
workflow_dir.join("workflow.toml"),
"_version = 1\n[workflow]\ngraph = \"workflow.fabro\"\n",
)
.unwrap();
let sha = "0123456789abcdef0123456789abcdef01234567".to_string();
let packaged = package_from_checkout(PackageFromCheckoutInput {
workflow: "demo".to_string(),
checkout_dir: checkout,
target: GitRunTarget {
repo: "workspace-org/app".to_string(),
branch: "release".to_string(),
tag: Some("v1".to_string()),
sha: Some(sha.clone()),
},
})
.expect("workflow versions should package from checkout");
assert_eq!(
packaged
.closure
.versions()
.last()
.unwrap()
.1
.version()
.entrypoint()
.as_str(),
".fabro/workflows/demo/workflow.fabro"
);
assert_eq!(packaged.target.tag.as_deref(), Some("v1"));
assert_eq!(packaged.target.sha.as_deref(), Some(sha.as_str()));
}
#[tokio::test]
async fn collected_closure_stores_dependency_first_and_idempotently() {
let temp = TempDir::new().unwrap();
@ -442,17 +382,8 @@ mod tests {
fs::create_dir_all(&child_dir).unwrap();
fs::write(child_dir.join("workflow.fabro"), "digraph Child {}").unwrap();
let packaged = package_from_checkout(PackageFromCheckoutInput {
workflow: "root".to_string(),
checkout_dir: checkout,
target: GitRunTarget {
repo: "workspace-org/app".to_string(),
branch: "main".to_string(),
tag: None,
sha: Some("0123456789abcdef0123456789abcdef01234567".to_string()),
},
})
.unwrap();
let closure =
fabro_manifest::collect_workflow_versions(Path::new("root"), &checkout).unwrap();
let database = fabro_store::test_support::test_database(
Arc::new(InMemory::new()),
"",
@ -462,16 +393,16 @@ mod tests {
let store = WorkflowVersionStore::new(database.blobs());
for _ in 0..2 {
for (expected, version) in packaged.closure.versions() {
for (expected, version) in closure.versions() {
assert_eq!(store.put(version).await.unwrap(), expected);
}
}
let loaded = store
.get_closure(&packaged.closure.root_id())
.get_closure(&closure.root_id())
.await
.unwrap()
.unwrap();
assert_eq!(loaded.root_id(), packaged.closure.root_id());
assert_eq!(loaded.root_id(), closure.root_id());
assert_eq!(loaded.versions().count(), 2);
}
}

View file

@ -303,10 +303,7 @@ fn mount_version(
.get(version.entrypoint())
.cloned()
.expect("validated workflow versions contain their entrypoint file");
let config_local = version
.entrypoint()
.resolve_reference("workflow.toml")
.expect("the static workflow config path should resolve beside a valid entrypoint");
let config_local = version.config_path();
let config_path = version.files().get(&config_local).map(|source| {
rebase_path(
version.entrypoint(),

View file

@ -8,9 +8,7 @@ use croner::errors::CronError;
use fabro_automation::{
Automation, AutomationId, AutomationRevision, AutomationTriggerId, parse_schedule_expression,
};
use fabro_types::{
AutomationRef, Principal, RunId, RunIntent, RunIntentArgs, RunTarget, SystemActorKind,
};
use fabro_types::{AutomationRef, Principal, RunId, SystemActorKind};
use tokio::time::sleep;
use tracing::{Instrument, error, info, info_span, warn};
@ -272,15 +270,7 @@ async fn fire_scheduled_automation_run(
let response = Box::pin(handler::runs::create_run_from_intent(
Arc::clone(&state),
handler::runs::CreateRunFromIntentRequest {
intent: RunIntent {
workflow_version_id: materialized.workflow_version_id,
target: RunTarget::Git(materialized.target),
args: RunIntentArgs::default(),
environment_id: None,
parent_id: None,
title: None,
goal: None,
},
intent: materialized.into_run_intent(),
explicit_run_id: Some(run_id),
actor: actor.clone(),
headers: HeaderMap::new(),
@ -351,7 +341,7 @@ mod tests {
use fabro_automation::{AutomationDraft, AutomationTrigger, ScheduleTrigger};
use fabro_static::EnvVars;
use fabro_store::ListRunsQuery;
use fabro_types::{GitRunTarget, RunStatus};
use fabro_types::{GitRunTarget, RunStatus, RunTarget};
use super::*;
use crate::test_support::{TestAppStateBuilder, TestAutomationRunMaterializer};

View file

@ -6,7 +6,7 @@ use fabro_automation::{
Automation, AutomationDraft, AutomationId, AutomationReplace, AutomationStoreError,
};
use fabro_store::{RunSummaryListQuery, RunSummaryVisibility};
use fabro_types::{AutomationRef, RunId, RunIntent, RunIntentArgs, RunTarget};
use fabro_types::{AutomationRef, RunId};
use fabro_util::error as error_util;
use serde::Serialize;
@ -151,15 +151,7 @@ async fn create_automation_run(
let response = Box::pin(runs::create_run_from_intent(
Arc::clone(&state),
runs::CreateRunFromIntentRequest {
intent: RunIntent {
workflow_version_id: materialized.workflow_version_id,
target: RunTarget::Git(materialized.target),
args: RunIntentArgs::default(),
environment_id: None,
parent_id: None,
title: None,
goal: None,
},
intent: materialized.into_run_intent(),
explicit_run_id: Some(run_id),
actor: actor.clone(),
headers,

View file

@ -130,13 +130,7 @@ pub fn build_sparse_run_overrides(input: RunOverrideInput<'_>) -> Option<RunLaye
}
pub fn build_run_manifest(input: ManifestBuildInput) -> Result<BuiltManifest> {
let root_location = WorkflowLocation::resolve(&input.workflow, &input.cwd)?;
if root_location.toml.is_none() && !root_location.graph.is_file() {
return Err(fabro_config::Error::WorkflowNotFound(
root_location.graph.display().to_string(),
)
.into());
}
let root_location = resolve_existing_workflow_location(&input.workflow, &input.cwd)?;
let project_config = discover_project_config(&root_location.dir)?;
let project_config_source = project_config
.as_ref()
@ -400,6 +394,19 @@ fn push_manifest_branch_best_effort(
let _ = push_branch_noninteractive(repo_path, "origin", branch);
}
/// Resolve a workflow reference and reject it when neither its config nor
/// its graph exists on disk.
/// A missing workflow surfaces as `fabro_config::Error::WorkflowNotFound`.
fn resolve_existing_workflow_location(workflow: &Path, cwd: &Path) -> Result<WorkflowLocation> {
let location = WorkflowLocation::resolve(workflow, cwd)?;
if location.toml.is_none() && !location.graph.is_file() {
return Err(
fabro_config::Error::WorkflowNotFound(location.graph.display().to_string()).into(),
);
}
Ok(location)
}
fn normalize_absolute_path(base_dir: &Path, reference: &str) -> Option<PathBuf> {
let path = Path::new(reference);
if path.is_absolute() || reference.starts_with('~') {

View file

@ -75,9 +75,12 @@ impl<'a> WorkflowBundler<'a> {
.collect())
}
pub(super) fn collect_versions(mut self, workflow: &Path) -> Result<CollectedWorkflowSources> {
pub(super) fn collect_versions(
mut self,
root: &WorkflowLocation,
) -> Result<CollectedWorkflowSources> {
self.workflow_version_projection = true;
let root_key = self.collect_workflow_entry(workflow, self.cwd)?;
let root_key = self.collect_workflow_location(root)?;
Ok(CollectedWorkflowSources {
root_key,
workflows: self.workflows,
@ -409,9 +412,11 @@ impl<'a> WorkflowBundler<'a> {
ReferenceKind::RunGoalFile,
Some(config_path.clone()),
)?;
std::fs::read_to_string(&bundled.absolute_path).with_context(|| {
format!("Failed to read {}", bundled.absolute_path.display())
})?
files
.get(&bundled.path.to_string())
.expect("collect_bundled_file inserts the goal file it returns")
.content
.clone()
}
};
self.collect_template_include_files(

View file

@ -1,9 +1,9 @@
use std::collections::{BTreeMap, HashMap, HashSet};
use std::path::{Path, PathBuf};
use fabro_config::project::WorkflowLocation;
use fabro_api::types;
use fabro_types::{
BlobHash, WorkflowPath, WorkflowPathParseError, WorkflowVersion, WorkflowVersionId,
WorkflowPath, WorkflowPathParseError, WorkflowVersion, WorkflowVersionId,
WorkflowVersionShapeError,
};
use fabro_workflow_version::{ValidatedWorkflowVersion, WorkflowVersionError};
@ -78,49 +78,56 @@ pub fn collect_workflow_versions(
workflow: &Path,
checkout_root: &Path,
) -> Result<CollectedWorkflowClosure, WorkflowVersionCollectError> {
let location = WorkflowLocation::resolve(workflow, checkout_root).map_err(|source| {
WorkflowVersionCollectError::Collect {
path: workflow.to_path_buf(),
source: anyhow::Error::new(source),
}
})?;
if location.toml.is_none() && !location.graph.is_file() {
return Err(WorkflowVersionCollectError::WorkflowNotFound {
path: workflow.to_path_buf(),
});
}
let location =
crate::resolve_existing_workflow_location(workflow, checkout_root).map_err(|source| {
if matches!(
source.downcast_ref::<fabro_config::Error>(),
Some(fabro_config::Error::WorkflowNotFound(_))
) {
WorkflowVersionCollectError::WorkflowNotFound {
path: workflow.to_path_buf(),
}
} else {
WorkflowVersionCollectError::Collect {
path: workflow.to_path_buf(),
source,
}
}
})?;
let inputs = HashMap::new();
let collected = WorkflowBundler::new(checkout_root, &inputs)
.collect_versions(workflow)
.collect_versions(&location)
.map_err(|source| WorkflowVersionCollectError::Collect {
path: workflow.to_path_buf(),
source,
})?;
VersionAssembler::new(&collected).assemble()
VersionAssembler::new(collected).assemble()
}
struct VersionAssembler<'a> {
collected: &'a CollectedWorkflowSources,
visiting: HashSet<String>,
ids: HashMap<String, WorkflowVersionId>,
emitted: HashSet<WorkflowVersionId>,
versions: Vec<(WorkflowVersionId, ValidatedWorkflowVersion)>,
struct VersionAssembler {
root_key: String,
/// Sources still waiting to be assembled; each is removed once visited.
pending: HashMap<String, CollectedWorkflowSource>,
visiting: HashSet<String>,
ids: HashMap<String, WorkflowVersionId>,
versions: Vec<(WorkflowVersionId, ValidatedWorkflowVersion)>,
}
impl<'a> VersionAssembler<'a> {
fn new(collected: &'a CollectedWorkflowSources) -> Self {
impl VersionAssembler {
fn new(collected: CollectedWorkflowSources) -> Self {
Self {
collected,
root_key: collected.root_key,
pending: collected.workflows,
visiting: HashSet::new(),
ids: HashMap::new(),
emitted: HashSet::new(),
ids: HashMap::new(),
versions: Vec::new(),
}
}
fn assemble(mut self) -> Result<CollectedWorkflowClosure, WorkflowVersionCollectError> {
let root_id = self.assemble_one(&self.collected.root_key)?;
let root_key = std::mem::take(&mut self.root_key);
let root_id = self.assemble_one(&root_key)?;
Ok(CollectedWorkflowClosure {
root_id,
versions: self.versions,
@ -139,31 +146,20 @@ impl<'a> VersionAssembler<'a> {
path: workflow_path(key)?,
});
}
let dependency_keys = self
.collected
.workflows
.get(key)
.ok_or_else(|| WorkflowVersionCollectError::MissingWorkflow {
let source = self.pending.remove(key).ok_or_else(|| {
WorkflowVersionCollectError::MissingWorkflow {
path: key.to_owned(),
})?
.dependency_keys
.iter()
.cloned()
.collect::<Vec<_>>();
}
})?;
let mut dependencies = BTreeMap::new();
for dependency_key in dependency_keys {
let dependency_id = self.assemble_one(&dependency_key)?;
dependencies.insert(workflow_path(&dependency_key)?, dependency_id);
for dependency_key in &source.dependency_keys {
let dependency_id = self.assemble_one(dependency_key)?;
dependencies.insert(workflow_path(dependency_key)?, dependency_id);
}
let source = self
.collected
.workflows
.get(key)
.expect("collected workflow must remain present while assembling");
let entrypoint = workflow_path(key)?;
let files = workflow_files(&entrypoint, source)?;
let files = workflow_files(&entrypoint, source.workflow)?;
let version =
WorkflowVersion::new(entrypoint.clone(), files, dependencies).map_err(|source| {
WorkflowVersionCollectError::InvalidShape {
@ -177,49 +173,38 @@ impl<'a> VersionAssembler<'a> {
source,
}
})?;
let canonical = validated.version().canonical_bytes().map_err(|source| {
let id = validated.version().id().map_err(|source| {
WorkflowVersionCollectError::InvalidShape {
entrypoint: entrypoint.clone(),
source,
}
})?;
let id = WorkflowVersionId::from(BlobHash::new(&canonical));
self.visiting.remove(key);
// Keys are distinct entrypoints and the entrypoint is part of the
// canonical bytes, so each key yields a distinct ID.
self.ids.insert(key.to_owned(), id);
if self.emitted.insert(id) {
self.versions.push((id, validated));
}
self.versions.push((id, validated));
Ok(id)
}
}
fn workflow_files(
entrypoint: &WorkflowPath,
source: &CollectedWorkflowSource,
workflow: types::ManifestWorkflow,
) -> Result<BTreeMap<WorkflowPath, String>, WorkflowVersionCollectError> {
let mut files = BTreeMap::new();
insert_file(
&mut files,
entrypoint,
entrypoint.clone(),
source.workflow.source.clone(),
)?;
if let Some(config) = source.workflow.config.as_ref() {
insert_file(&mut files, entrypoint, entrypoint.clone(), workflow.source)?;
if let Some(config) = workflow.config {
insert_file(
&mut files,
entrypoint,
workflow_path(&config.path)?,
config.source.clone(),
config.source,
)?;
}
for (path, file) in &source.workflow.files {
insert_file(
&mut files,
entrypoint,
workflow_path(path)?,
file.content.clone(),
)?;
for (path, file) in workflow.files {
insert_file(&mut files, entrypoint, workflow_path(&path)?, file.content)?;
}
Ok(files)
}

View file

@ -116,7 +116,7 @@ impl ValidatedWorkflowVersion {
/// run-admission time read the same bytes validation proved present.
#[must_use]
pub fn resolved_goal_file_content(&self) -> Option<&str> {
let config_path = workflow_config_path(&self.0);
let config_path = self.0.config_path();
let source = self.0.files().get(&config_path)?;
let layer: SettingsLayer = source.parse().expect("validated workflow.toml must parse");
let RunGoalLayer::File { file } = layer.run.as_ref().and_then(|run| run.goal.as_ref())?
@ -163,7 +163,7 @@ fn validate_config(
version: &WorkflowVersion,
template_roots: &mut TemplateRoots,
) -> Result<(), WorkflowVersionError> {
let config_path = workflow_config_path(version);
let config_path = version.config_path();
let Some(source) = version.files().get(&config_path) else {
return Ok(());
};
@ -212,13 +212,6 @@ fn validate_config(
Ok(())
}
fn workflow_config_path(version: &WorkflowVersion) -> WorkflowPath {
version
.entrypoint()
.resolve_reference("workflow.toml")
.expect("the static workflow config path must resolve beside a valid entrypoint")
}
#[expect(
clippy::disallowed_methods,
reason = "workflow-version validation preserves authored template source for dependency discovery"

View file

@ -6,7 +6,7 @@ use serde::de::{Error as _, MapAccess, Visitor};
use serde::{Deserialize, Deserializer, Serialize};
use thiserror::Error;
use crate::{WorkflowPath, WorkflowVersionId};
use crate::{BlobHash, WorkflowPath, WorkflowVersionId};
pub const MAX_WORKFLOW_VERSION_FILES: usize = 512;
pub const MAX_WORKFLOW_VERSION_DEPENDENCIES: usize = 512;
@ -86,6 +86,22 @@ impl WorkflowVersion {
&self.workflow_dependencies
}
/// Path of the optional `workflow.toml` that configures this version. It
/// always sits beside the entrypoint graph.
#[must_use]
pub fn config_path(&self) -> WorkflowPath {
self.entrypoint
.resolve_reference("workflow.toml")
.expect("the static workflow config path must resolve beside a valid entrypoint")
}
/// Content-addressed identity: the hash of the canonical wire form.
pub fn id(&self) -> Result<WorkflowVersionId, WorkflowVersionShapeError> {
Ok(WorkflowVersionId::from(BlobHash::new(
&self.canonical_bytes()?,
)))
}
/// Serialize to the canonical wire form.
///
/// Structural validity is guaranteed by construction, so this only