From 7bf40373f95a886ce45da78911aaf2f5b6ab04e8 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Tue, 11 Aug 2026 12:32:22 -0400 Subject: [PATCH] Split workflow version wire type from semantic validation Move the WorkflowVersion wire type and its structural invariants (file count and size limits, entrypoint presence, unique keys, path collisions, canonical form) into fabro-types, so fabro-api replaces the generated schema type without pulling graph parsing or the template engine into every API consumer. The wire shape is unchanged. fabro-workflow-version keeps the expensive graph/config/template validation behind a ValidatedWorkflowVersion newtype and now owns WorkflowVersionStore: put only accepts validated versions and get still re-validates blobs read from shared storage. fabro-store goes back to being domain-agnostic persistence. Move the static-reference attribute vocabulary (ReferenceKind, AttributeScope, reference_kind_for_attribute) to fabro_types::graph and template-syntax validation to fabro-template, so fabro-graphviz no longer depends on the template engine. Co-Authored-By: Claude Fable 5 --- Cargo.lock | 6 +- .../src/server/handler/workflow_versions.rs | 20 +- lib/components/fabro-graphviz/Cargo.toml | 1 - lib/components/fabro-graphviz/src/lib.rs | 1 - .../fabro-graphviz/src/static_reference.rs | 133 --- lib/components/fabro-manifest/src/lib.rs | 6 +- .../fabro-manifest/src/workflow_bundler.rs | 23 +- lib/components/fabro-store/Cargo.toml | 1 - lib/components/fabro-store/src/lib.rs | 2 - .../fabro-workflow-version/Cargo.toml | 7 +- .../fabro-workflow-version/src/lib.rs | 899 ++++++------------ .../src/store.rs} | 57 +- .../src/handler/manager_loop.rs | 3 +- .../fabro-workflow/src/transforms/import.rs | 4 +- .../src/transforms/importable_field.rs | 3 +- .../src/transforms/variable_expansion.rs | 6 +- lib/foundation/fabro-api/Cargo.toml | 1 - lib/foundation/fabro-api/build.rs | 6 +- lib/foundation/fabro-api/src/lib.rs | 3 +- .../tests/workflow_version_round_trip.rs | 3 +- lib/foundation/fabro-template/src/lib.rs | 62 ++ lib/foundation/fabro-types/src/graph.rs | 66 ++ lib/foundation/fabro-types/src/lib.rs | 5 + .../fabro-types/src/workflow_version.rs | 379 ++++++++ 24 files changed, 894 insertions(+), 803 deletions(-) delete mode 100644 lib/components/fabro-graphviz/src/static_reference.rs rename lib/components/{fabro-store/src/workflow_version_store.rs => fabro-workflow-version/src/store.rs} (78%) create mode 100644 lib/foundation/fabro-types/src/workflow_version.rs diff --git a/Cargo.lock b/Cargo.lock index 996ddec8f..e3006a26d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2328,7 +2328,6 @@ dependencies = [ "fabro-environment", "fabro-model", "fabro-types", - "fabro-workflow-version", "openapiv3", "prettyplease", "progenitor", @@ -2677,7 +2676,6 @@ name = "fabro-graphviz" version = "0.324.0-nightly.0" dependencies = [ "anyhow", - "fabro-template", "fabro-types", "graphviz-sys", "nom", @@ -3149,7 +3147,6 @@ dependencies = [ "fabro-db", "fabro-types", "fabro-util", - "fabro-workflow-version", "futures", "hex", "insta", @@ -3440,11 +3437,14 @@ version = "0.316.0-nightly.0" dependencies = [ "fabro-config", "fabro-graphviz", + "fabro-store", "fabro-template", "fabro-types", + "object_store", "serde", "serde_json", "thiserror 2.0.18", + "tokio", ] [[package]] diff --git a/lib/apps/fabro-server/src/server/handler/workflow_versions.rs b/lib/apps/fabro-server/src/server/handler/workflow_versions.rs index 4f8a2762a..6b4a98f9a 100644 --- a/lib/apps/fabro-server/src/server/handler/workflow_versions.rs +++ b/lib/apps/fabro-server/src/server/handler/workflow_versions.rs @@ -3,9 +3,11 @@ use std::sync::Arc; use axum::extract::DefaultBodyLimit; use axum::extract::rejection::JsonRejection; use fabro_api::types::{CreateWorkflowVersionResponse, WorkflowVersion}; -use fabro_store::{WorkflowVersionStore, WorkflowVersionStoreError}; +use fabro_types::MAX_WORKFLOW_VERSION_BYTES; use fabro_util::error; -use fabro_workflow_version::MAX_WORKFLOW_VERSION_BYTES; +use fabro_workflow_version::{ + ValidatedWorkflowVersion, WorkflowVersionStore, WorkflowVersionStoreError, +}; use super::super::{ ApiError, AppState, IntoResponse, Json, RequiredUser, Response, Router, State, StatusCode, post, @@ -29,6 +31,13 @@ async fn create_workflow_version( payload: Result, JsonRejection>, ) -> Result { let Json(version) = payload.map_err(json_rejection)?; + let version = ValidatedWorkflowVersion::new(version).map_err(|err| { + ApiError::with_code( + StatusCode::UNPROCESSABLE_ENTITY, + err.to_string(), + INVALID_VERSION_CODE, + ) + })?; let blobs = state.store_ref().blobs().await.map_err(|err| { tracing::error!( error = %err, @@ -85,6 +94,11 @@ fn store_error(err: WorkflowVersionStoreError) -> ApiError { source.to_string(), INVALID_VERSION_CODE, ), + WorkflowVersionStoreError::InvalidShape(source) => ApiError::with_code( + StatusCode::UNPROCESSABLE_ENTITY, + source.to_string(), + INVALID_VERSION_CODE, + ), err => { tracing::error!( error = %err, @@ -307,7 +321,7 @@ mod tests { #[tokio::test] async fn storage_fault_response_is_curated() { - let response = store_error(fabro_store::WorkflowVersionStoreError::Storage { + let response = store_error(fabro_workflow_version::WorkflowVersionStoreError::Storage { source: fabro_store::Error::Other("private persistence detail".to_string()), }) .into_response(); diff --git a/lib/components/fabro-graphviz/Cargo.toml b/lib/components/fabro-graphviz/Cargo.toml index 0df998065..f78b3b9c2 100644 --- a/lib/components/fabro-graphviz/Cargo.toml +++ b/lib/components/fabro-graphviz/Cargo.toml @@ -16,7 +16,6 @@ workspace = true anyhow.workspace = true graphviz-sys.workspace = true fabro-types = { path = "../../foundation/fabro-types" } -fabro-template = { path = "../../foundation/fabro-template" } nom = "7" regex = { workspace = true } serde = { workspace = true } diff --git a/lib/components/fabro-graphviz/src/lib.rs b/lib/components/fabro-graphviz/src/lib.rs index 3efef3593..ac4a2562f 100644 --- a/lib/components/fabro-graphviz/src/lib.rs +++ b/lib/components/fabro-graphviz/src/lib.rs @@ -4,7 +4,6 @@ pub mod fidelity; pub mod graph; pub mod parser; pub mod render; -pub mod static_reference; pub mod stylesheet; pub use error::{Error, Result}; diff --git a/lib/components/fabro-graphviz/src/static_reference.rs b/lib/components/fabro-graphviz/src/static_reference.rs deleted file mode 100644 index 55ec7013b..000000000 --- a/lib/components/fabro-graphviz/src/static_reference.rs +++ /dev/null @@ -1,133 +0,0 @@ -use fabro_template::contains_template_syntax; -use thiserror::Error; - -#[derive(Clone, Copy, Debug, Eq, PartialEq, strum::Display)] -pub enum ReferenceKind { - #[strum(to_string = "file inline reference")] - FileInline, - #[strum(to_string = "import reference")] - Import, - #[strum(to_string = "child workflow reference")] - ChildWorkflow, - #[strum(to_string = "Dockerfile reference")] - Dockerfile, - #[strum(to_string = "graph goal file reference")] - GraphGoalFile, -} - -impl ReferenceKind { - pub fn validate(self, value: &str) -> Result<(), StaticReferenceError> { - validate_static_reference(value, self) - } -} - -#[derive(Clone, Copy, Debug, Eq, PartialEq)] -pub enum AttributeScope { - Graph, - Node, - Edge, -} - -#[derive(Debug, Error)] -#[error("templates are not supported in {kind}s: {value}")] -pub struct StaticReferenceError { - kind: ReferenceKind, - value: String, -} - -impl StaticReferenceError { - #[must_use] - pub fn new(kind: ReferenceKind, value: impl Into) -> Self { - Self { - kind, - value: value.into(), - } - } - - #[must_use] - pub fn kind(&self) -> ReferenceKind { - self.kind - } - - #[must_use] - pub fn value(&self) -> &str { - &self.value - } -} - -pub fn validate_static_reference( - value: &str, - kind: ReferenceKind, -) -> Result<(), StaticReferenceError> { - if contains_template_syntax(value) { - return Err(StaticReferenceError::new(kind, value)); - } - Ok(()) -} - -#[must_use] -pub fn reference_kind_for_attribute( - scope: AttributeScope, - key: &str, - value: &str, -) -> Option { - match key { - "import" => Some(ReferenceKind::Import), - "stack.child_workflow" | "stack.child_dotfile" => Some(ReferenceKind::ChildWorkflow), - "goal" if matches!(scope, AttributeScope::Graph) && value.starts_with('@') => { - Some(ReferenceKind::GraphGoalFile) - } - "prompt" | "output_schema" - if matches!(scope, AttributeScope::Node) && value.starts_with('@') => - { - Some(ReferenceKind::FileInline) - } - _ => None, - } -} - -#[cfg(test)] -mod tests { - use super::{AttributeScope, ReferenceKind, reference_kind_for_attribute}; - - #[test] - fn output_schema_at_value_is_file_inline_reference() { - assert_eq!( - reference_kind_for_attribute( - AttributeScope::Node, - "output_schema", - "@schemas/result.schema.json", - ), - Some(ReferenceKind::FileInline), - ); - } - - #[test] - fn output_schema_builtin_keyword_is_not_file_inline_reference() { - assert_eq!( - reference_kind_for_attribute(AttributeScope::Node, "output_schema", "routing"), - None, - ); - } - - #[test] - fn output_schema_reference_rejects_template_syntax() { - let error = reference_kind_for_attribute( - AttributeScope::Node, - "output_schema", - "@schemas/{{ inputs.schema }}.json", - ) - .expect("output_schema @ references should be static references") - .validate("@schemas/{{ inputs.schema }}.json") - .unwrap_err(); - - assert_eq!(error.kind(), ReferenceKind::FileInline); - assert_eq!(error.value(), "@schemas/{{ inputs.schema }}.json"); - assert!( - error - .to_string() - .contains("templates are not supported in file inline references"), - "unexpected error: {error}", - ); - } -} diff --git a/lib/components/fabro-manifest/src/lib.rs b/lib/components/fabro-manifest/src/lib.rs index dc9a0d2cf..d66e2e21a 100644 --- a/lib/components/fabro-manifest/src/lib.rs +++ b/lib/components/fabro-manifest/src/lib.rs @@ -19,7 +19,8 @@ use fabro_config::{ }; use fabro_graphviz::graph::AttrValue; use fabro_graphviz::parser; -use fabro_graphviz::static_reference::ReferenceKind; +use fabro_template::validate_static_reference; +use fabro_types::graph::ReferenceKind; use fabro_types::settings::interp::InterpString; use fabro_types::settings::run::{ApprovalMode, ResolvedGoalSource, ResolvedRunGoal, RunMode}; use fabro_types::{DirtyStatus, GitContext, ManifestPath, WorkflowSettings}; @@ -264,8 +265,7 @@ fn resolve_manifest_goal( return Ok(None); }; if let Some(reference) = goal.strip_prefix('@') { - ReferenceKind::GraphGoalFile - .validate(reference) + validate_static_reference(reference, ReferenceKind::GraphGoalFile) .map_err(anyhow::Error::new)?; let goal_path = normalize_absolute_path( root_dot_path.parent().unwrap_or_else(|| Path::new(".")), diff --git a/lib/components/fabro-manifest/src/workflow_bundler.rs b/lib/components/fabro-manifest/src/workflow_bundler.rs index 7b150e629..b26d966b7 100644 --- a/lib/components/fabro-manifest/src/workflow_bundler.rs +++ b/lib/components/fabro-manifest/src/workflow_bundler.rs @@ -8,11 +8,11 @@ use fabro_config::project::WorkflowLocation; use fabro_config::{EnvironmentDockerfileLayer, EnvironmentImageLayer, SettingsLayer}; use fabro_graphviz::graph::AttrValue; use fabro_graphviz::parser; -use fabro_graphviz::static_reference::{self, AttributeScope, ReferenceKind}; use fabro_template::{ BundleTemplateStore, FilesystemTemplateStore, RecordingTemplateStore, TemplateContext, - TemplateDependencyClosure, TemplateRenderMode, TemplateSource, + TemplateDependencyClosure, TemplateRenderMode, TemplateSource, validate_static_reference, }; +use fabro_types::graph::{AttributeScope, ReferenceKind, reference_kind_for_attribute}; use fabro_types::ManifestPath; use crate::{manifest_path_from_absolute, normalize_absolute_path}; @@ -176,11 +176,7 @@ impl<'a> WorkflowBundler<'a> { continue; }; let Some(ReferenceKind::FileInline) = - static_reference::reference_kind_for_attribute( - AttributeScope::Node, - name, - value, - ) + reference_kind_for_attribute(AttributeScope::Node, name, value) else { continue; }; @@ -234,13 +230,12 @@ impl<'a> WorkflowBundler<'a> { .get("stack.child_workflow") .and_then(AttrValue::as_str) { - manifest_attr_reference_kind( + let kind = manifest_attr_reference_kind( AttributeScope::Node, "stack.child_workflow", child_ref, - )? - .validate(child_ref) - .map_err(anyhow::Error::new)?; + )?; + validate_static_reference(child_ref, kind).map_err(anyhow::Error::new)?; self.collect_workflow_entry(Path::new(child_ref), workflow_base_dir)?; } } @@ -397,9 +392,7 @@ impl<'a> WorkflowBundler<'a> { reference_kind: ReferenceKind, from: Option, ) -> Result { - reference_kind - .validate(reference) - .map_err(anyhow::Error::new)?; + validate_static_reference(reference, reference_kind).map_err(anyhow::Error::new)?; let absolute_path = normalize_absolute_path(base_dir, reference) .ok_or_else(|| anyhow!("unsupported manifest reference: {reference}"))?; @@ -468,7 +461,7 @@ fn manifest_attr_reference_kind( key: &str, value: &str, ) -> Result { - static_reference::reference_kind_for_attribute(scope, key, value) + reference_kind_for_attribute(scope, key, value) .ok_or_else(|| anyhow!("unsupported manifest reference attribute: {key}={value}")) } diff --git a/lib/components/fabro-store/Cargo.toml b/lib/components/fabro-store/Cargo.toml index 3e73f9237..588e285a0 100644 --- a/lib/components/fabro-store/Cargo.toml +++ b/lib/components/fabro-store/Cargo.toml @@ -16,7 +16,6 @@ test-support = [] [dependencies] fabro-types = { path = "../../foundation/fabro-types" } -fabro-workflow-version = { path = "../fabro-workflow-version" } fabro-util = { path = "../../foundation/fabro-util" } hex.workspace = true slatedb.workspace = true diff --git a/lib/components/fabro-store/src/lib.rs b/lib/components/fabro-store/src/lib.rs index 229936b45..1b514a8e0 100644 --- a/lib/components/fabro-store/src/lib.rs +++ b/lib/components/fabro-store/src/lib.rs @@ -13,7 +13,6 @@ mod slate; #[cfg(any(test, feature = "test-support"))] pub mod test_support; mod types; -mod workflow_version_store; pub use artifact_store::{ ArtifactKey, ArtifactStore, NodeArtifact, StageArtifactEntry, retry_storage_segment, @@ -39,7 +38,6 @@ pub use slate::{ RefreshToken, RefreshTokenStore, RunCatalogIndex, RunDatabase, Runs, UnreadableRun, }; pub use types::EventPayload; -pub use workflow_version_store::{WorkflowVersionStore, WorkflowVersionStoreError}; #[derive(Debug, Default, Clone, PartialEq, Eq)] pub struct ListRunsQuery { diff --git a/lib/components/fabro-workflow-version/Cargo.toml b/lib/components/fabro-workflow-version/Cargo.toml index 84c315479..bc91e3f89 100644 --- a/lib/components/fabro-workflow-version/Cargo.toml +++ b/lib/components/fabro-workflow-version/Cargo.toml @@ -4,7 +4,7 @@ edition.workspace = true version.workspace = true publish = false license.workspace = true -description = "Immutable workflow-version package model and validation" +description = "Semantic validation and storage for immutable workflow versions" [lib] doctest = false @@ -15,8 +15,13 @@ workspace = true [dependencies] fabro-config = { path = "../../foundation/fabro-config" } fabro-graphviz = { path = "../fabro-graphviz" } +fabro-store = { path = "../fabro-store" } fabro-template = { path = "../../foundation/fabro-template" } fabro-types = { path = "../../foundation/fabro-types" } serde.workspace = true serde_json.workspace = true thiserror.workspace = true + +[dev-dependencies] +object_store.workspace = true +tokio = { workspace = true, features = ["full"] } diff --git a/lib/components/fabro-workflow-version/src/lib.rs b/lib/components/fabro-workflow-version/src/lib.rs index 61976595d..8de855c85 100644 --- a/lib/components/fabro-workflow-version/src/lib.rs +++ b/lib/components/fabro-workflow-version/src/lib.rs @@ -1,44 +1,30 @@ -use std::collections::{BTreeMap, BTreeSet, HashMap, VecDeque}; -use std::fmt; -use std::marker::PhantomData; +//! Semantic validation for immutable workflow versions. +//! +//! The wire type ([`fabro_types::WorkflowVersion`]) enforces structural +//! invariants at construction. This crate owns the expensive semantic +//! validation — graph closure, config, and template checks — behind the +//! [`ValidatedWorkflowVersion`] newtype, and the content-addressed +//! [`WorkflowVersionStore`] that only accepts and returns validated versions. + +use std::collections::{BTreeSet, HashMap, VecDeque}; use fabro_config::{EnvironmentDockerfileLayer, EnvironmentImageLayer, SettingsLayer}; use fabro_graphviz::graph::Graph; use fabro_graphviz::parser; -use fabro_graphviz::static_reference::{ - AttributeScope, ReferenceKind, StaticReferenceError, reference_kind_for_attribute, -}; use fabro_template::{ - BundleTemplateStore, TemplateDiscoveryError, TemplateSource, discover_static_dependency_closure, + BundleTemplateStore, StaticReferenceError, TemplateDiscoveryError, TemplateSource, + discover_static_dependency_closure, validate_static_reference, }; -use fabro_types::{ManifestPath, WorkflowPath, WorkflowPathParseError, WorkflowVersionId}; -use serde::de::{Error as _, MapAccess, Visitor}; -use serde::{Deserialize, Deserializer, Serialize}; +use fabro_types::graph::{AttributeScope, ReferenceKind, reference_kind_for_attribute}; +use fabro_types::{ManifestPath, WorkflowPath, WorkflowPathParseError, WorkflowVersion}; use thiserror::Error; -pub const MAX_WORKFLOW_VERSION_FILES: usize = 512; -pub const MAX_WORKFLOW_VERSION_FILE_BYTES: usize = 512 * 1024; -pub const MAX_WORKFLOW_VERSION_BYTES: usize = 2 * 1024 * 1024; +mod store; + +pub use store::{WorkflowVersionStore, WorkflowVersionStoreError}; #[derive(Debug, Error)] pub enum WorkflowVersionError { - #[error("workflow version has {actual} files; maximum is {maximum}")] - TooManyFiles { actual: usize, maximum: usize }, - #[error("workflow file `{path}` is {actual} bytes; maximum is {maximum}")] - FileTooLarge { - path: WorkflowPath, - actual: usize, - maximum: usize, - }, - #[error("workflow version is {actual} canonical bytes; maximum is {maximum}")] - VersionTooLarge { actual: usize, maximum: usize }, - #[error("entrypoint `{path}` is not present in workflow files")] - MissingEntrypoint { path: WorkflowPath }, - #[error("workflow paths collide: `{first}` and `{second}`")] - PathCollision { - first: WorkflowPath, - second: WorkflowPath, - }, #[error("workflow graph `{path}` is invalid")] GraphParse { path: WorkflowPath, @@ -88,372 +74,295 @@ pub enum WorkflowVersionError { missing: Vec, unused: Vec, }, - #[error("failed to serialize canonical workflow version")] - Serialization { - #[source] - source: serde_json::Error, - }, } -#[derive(Clone, Debug, PartialEq, Eq, Serialize)] -pub struct WorkflowVersion { - entrypoint: WorkflowPath, - files: BTreeMap, - dependencies: BTreeMap, +/// A workflow version whose graph, config, and template content passed +/// semantic validation. +/// +/// This is the only door: functions that require a semantically valid +/// version take this type, and the only way to obtain one is [`Self::new`] +/// (or loading through [`WorkflowVersionStore`], which validates on read). +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct ValidatedWorkflowVersion(WorkflowVersion); + +impl ValidatedWorkflowVersion { + pub fn new(version: WorkflowVersion) -> Result { + validate_config(&version)?; + validate_graph_closure(&version)?; + Ok(Self(version)) + } + + #[must_use] + pub fn version(&self) -> &WorkflowVersion { + &self.0 + } + + #[must_use] + pub fn into_version(self) -> WorkflowVersion { + self.0 + } } -impl WorkflowVersion { - pub fn new( - entrypoint: WorkflowPath, - files: BTreeMap, - dependencies: BTreeMap, - ) -> Result { - let version = Self { - entrypoint, - files, - dependencies, - }; - version.validate_structure()?; - version.canonical_bytes()?; - Ok(version) - } +fn validate_config(version: &WorkflowVersion) -> Result<(), WorkflowVersionError> { + let config_path = + WorkflowPath::new("workflow.toml").expect("the static workflow config path must be valid"); + let Some(source) = version.files().get(&config_path) else { + return Ok(()); + }; + let layer = source + .parse::() + .map_err(|source| WorkflowVersionError::Config { source })?; - #[must_use] - pub fn entrypoint(&self) -> &WorkflowPath { - &self.entrypoint - } - - #[must_use] - pub fn files(&self) -> &BTreeMap { - &self.files - } - - #[must_use] - pub fn dependencies(&self) -> &BTreeMap { - &self.dependencies - } - - /// Serialize to the canonical wire form. - /// - /// Structural validity is guaranteed by construction (`new` and - /// `Deserialize` both validate), so this only serializes and enforces - /// the canonical size limit. - pub fn canonical_bytes(&self) -> Result, WorkflowVersionError> { - let bytes = serde_json::to_vec(self) - .map_err(|source| WorkflowVersionError::Serialization { source })?; - if bytes.len() > MAX_WORKFLOW_VERSION_BYTES { - return Err(WorkflowVersionError::VersionTooLarge { - actual: bytes.len(), - maximum: MAX_WORKFLOW_VERSION_BYTES, + if let Some(configured) = layer + .workflow + .as_ref() + .and_then(|workflow| workflow.graph.as_deref()) + { + let configured = resolve_reference(&config_path, ReferenceKind::FileInline, configured)?; + if configured != *version.entrypoint() { + return Err(WorkflowVersionError::ConfigEntrypointMismatch { + configured, + entrypoint: version.entrypoint().clone(), }); } - Ok(bytes) } - fn validate_structure(&self) -> Result<(), WorkflowVersionError> { - if self.files.len() > MAX_WORKFLOW_VERSION_FILES { - return Err(WorkflowVersionError::TooManyFiles { - actual: self.files.len(), - maximum: MAX_WORKFLOW_VERSION_FILES, - }); + for environment in layer.environments.values() { + validate_dockerfile(version, &config_path, environment.image.as_ref())?; + } + if let Some(image) = layer + .run + .as_ref() + .and_then(|run| run.environment.as_ref()) + .and_then(|environment| environment.image.as_ref()) + { + validate_dockerfile(version, &config_path, Some(image))?; + } + Ok(()) +} + +fn validate_dockerfile( + version: &WorkflowVersion, + config_path: &WorkflowPath, + image: Option<&EnvironmentImageLayer>, +) -> Result<(), WorkflowVersionError> { + let Some(EnvironmentDockerfileLayer::Path { path }) = + image.and_then(|image| image.dockerfile.as_ref()) + else { + return Ok(()); + }; + validate_static_reference(path, ReferenceKind::Dockerfile).map_err(|source| { + WorkflowVersionError::StaticReference { + path: config_path.clone(), + source, } - for (path, content) in &self.files { - if content.len() > MAX_WORKFLOW_VERSION_FILE_BYTES { - return Err(WorkflowVersionError::FileTooLarge { - path: path.clone(), - actual: content.len(), - maximum: MAX_WORKFLOW_VERSION_FILE_BYTES, - }); - } + })?; + let target = resolve_reference(config_path, ReferenceKind::Dockerfile, path)?; + require_file(version, config_path, ReferenceKind::Dockerfile, target).map(|_| ()) +} + +fn validate_graph_closure(version: &WorkflowVersion) -> Result<(), WorkflowVersionError> { + let template_store = template_store(version); + let template_root = ManifestPath::from_wire(".") + .expect("the template package root must be a valid manifest path"); + let mut queue = VecDeque::from([version.entrypoint().clone()]); + let mut visited = BTreeSet::new(); + let mut child_workflows = BTreeSet::new(); + + while let Some(path) = queue.pop_front() { + if !visited.insert(path.clone()) { + continue; } - if !self.files.contains_key(&self.entrypoint) { - return Err(WorkflowVersionError::MissingEntrypoint { - path: self.entrypoint.clone(), - }); - } - self.validate_path_collisions()?; - self.validate_config()?; - self.validate_graph_closure() + let source = + version + .files() + .get(&path) + .ok_or_else(|| WorkflowVersionError::MissingFile { + path: path.clone(), + kind: ReferenceKind::Import, + target: path.clone(), + })?; + let graph = parser::parse(source).map_err(|source| WorkflowVersionError::GraphParse { + path: path.clone(), + source, + })?; + + validate_graph_goal(version, &path, &graph, &template_store, &template_root)?; + validate_graph_nodes( + version, + &path, + &graph, + &template_store, + &template_root, + &mut queue, + &mut child_workflows, + )?; } - fn validate_path_collisions(&self) -> Result<(), WorkflowVersionError> { - // Keys are unique within each map, so equality can only collide - // across files and dependencies. - let paths = self - .files - .keys() - .chain(self.dependencies.keys()) - .collect::>(); - for (index, first) in paths.iter().enumerate() { - for second in &paths[index + 1..] { - if first == second || first.is_ancestor_of(second) || second.is_ancestor_of(first) { - return Err(WorkflowVersionError::PathCollision { - first: (*first).clone(), - second: (*second).clone(), - }); - } - } - } - Ok(()) + let configured = version + .dependencies() + .keys() + .cloned() + .collect::>(); + if child_workflows != configured { + return Err(WorkflowVersionError::DependencyMismatch { + missing: child_workflows.difference(&configured).cloned().collect(), + unused: configured.difference(&child_workflows).cloned().collect(), + }); } + Ok(()) +} - fn validate_config(&self) -> Result<(), WorkflowVersionError> { - let config_path = WorkflowPath::new("workflow.toml") - .expect("the static workflow config path must be valid"); - let Some(source) = self.files.get(&config_path) else { - return Ok(()); - }; - let layer = source - .parse::() - .map_err(|source| WorkflowVersionError::Config { source })?; - - if let Some(configured) = layer - .workflow - .as_ref() - .and_then(|workflow| workflow.graph.as_deref()) - { - let configured = - Self::resolve_reference(&config_path, ReferenceKind::FileInline, configured)?; - if configured != self.entrypoint { - return Err(WorkflowVersionError::ConfigEntrypointMismatch { - configured, - entrypoint: self.entrypoint.clone(), - }); - } - } - - for environment in layer.environments.values() { - self.validate_dockerfile(&config_path, environment.image.as_ref())?; - } - if let Some(image) = layer - .run - .as_ref() - .and_then(|run| run.environment.as_ref()) - .and_then(|environment| environment.image.as_ref()) - { - self.validate_dockerfile(&config_path, Some(image))?; - } - Ok(()) +fn validate_graph_goal( + version: &WorkflowVersion, + graph_path: &WorkflowPath, + graph: &Graph, + template_store: &BundleTemplateStore, + template_root: &ManifestPath, +) -> Result<(), WorkflowVersionError> { + let goal = graph.goal(); + if goal.is_empty() { + return Ok(()); } - - fn validate_dockerfile( - &self, - config_path: &WorkflowPath, - image: Option<&EnvironmentImageLayer>, - ) -> Result<(), WorkflowVersionError> { - let Some(EnvironmentDockerfileLayer::Path { path }) = - image.and_then(|image| image.dockerfile.as_ref()) - else { - return Ok(()); - }; - ReferenceKind::Dockerfile.validate(path).map_err(|source| { + if let Some(reference) = goal.strip_prefix('@') { + validate_static_reference(reference, ReferenceKind::GraphGoalFile).map_err(|source| { WorkflowVersionError::StaticReference { - path: config_path.clone(), + path: graph_path.clone(), source, } })?; - let target = Self::resolve_reference(config_path, ReferenceKind::Dockerfile, path)?; - self.require_file(config_path, ReferenceKind::Dockerfile, target) - .map(|_| ()) + let target = resolve_reference(graph_path, ReferenceKind::GraphGoalFile, reference)?; + let content = require_file( + version, + graph_path, + ReferenceKind::GraphGoalFile, + target.clone(), + )?; + return validate_template(&target, content, template_store, template_root); } + validate_template(graph_path, goal, template_store, template_root) +} - fn validate_graph_closure(&self) -> Result<(), WorkflowVersionError> { - let template_store = self.template_store(); - let template_root = ManifestPath::from_wire(".") - .expect("the template package root must be a valid manifest path"); - let mut queue = VecDeque::from([self.entrypoint.clone()]); - let mut visited = BTreeSet::new(); - let mut child_workflows = BTreeSet::new(); - - while let Some(path) = queue.pop_front() { - if !visited.insert(path.clone()) { +#[allow( + clippy::too_many_arguments, + reason = "Graph validation threads one explicit closure accumulator through node attributes." +)] +fn validate_graph_nodes( + version: &WorkflowVersion, + graph_path: &WorkflowPath, + graph: &Graph, + template_store: &BundleTemplateStore, + template_root: &ManifestPath, + imports: &mut VecDeque, + child_workflows: &mut BTreeSet, +) -> Result<(), WorkflowVersionError> { + for node in graph.nodes.values() { + for (key, value) in &node.attrs { + let Some(value) = value.as_str() else { continue; - } - let source = - self.files - .get(&path) - .ok_or_else(|| WorkflowVersionError::MissingFile { - path: path.clone(), - kind: ReferenceKind::Import, - target: path.clone(), - })?; - let graph = - parser::parse(source).map_err(|source| WorkflowVersionError::GraphParse { - path: path.clone(), - source, - })?; - - self.validate_graph_goal(&path, &graph, &template_store, &template_root)?; - self.validate_graph_nodes( - &path, - &graph, - &template_store, - &template_root, - &mut queue, - &mut child_workflows, - )?; - } - - let configured = self.dependencies.keys().cloned().collect::>(); - if child_workflows != configured { - return Err(WorkflowVersionError::DependencyMismatch { - missing: child_workflows.difference(&configured).cloned().collect(), - unused: configured.difference(&child_workflows).cloned().collect(), - }); - } - Ok(()) - } - - fn validate_graph_goal( - &self, - graph_path: &WorkflowPath, - graph: &Graph, - template_store: &BundleTemplateStore, - template_root: &ManifestPath, - ) -> Result<(), WorkflowVersionError> { - let goal = graph.goal(); - if goal.is_empty() { - return Ok(()); - } - if let Some(reference) = goal.strip_prefix('@') { - ReferenceKind::GraphGoalFile - .validate(reference) - .map_err(|source| WorkflowVersionError::StaticReference { + }; + let Some(kind) = reference_kind_for_attribute(AttributeScope::Node, key, value) else { + continue; + }; + validate_static_reference(value, kind).map_err(|source| { + WorkflowVersionError::StaticReference { path: graph_path.clone(), source, - })?; - let target = - Self::resolve_reference(graph_path, ReferenceKind::GraphGoalFile, reference)?; - let content = - self.require_file(graph_path, ReferenceKind::GraphGoalFile, target.clone())?; - return Self::validate_template(&target, content, template_store, template_root); - } - Self::validate_template(graph_path, goal, template_store, template_root) - } - - #[allow( - clippy::too_many_arguments, - reason = "Graph validation threads one explicit closure accumulator through node attributes." - )] - fn validate_graph_nodes( - &self, - graph_path: &WorkflowPath, - graph: &Graph, - template_store: &BundleTemplateStore, - template_root: &ManifestPath, - imports: &mut VecDeque, - child_workflows: &mut BTreeSet, - ) -> Result<(), WorkflowVersionError> { - for node in graph.nodes.values() { - for (key, value) in &node.attrs { - let Some(value) = value.as_str() else { - continue; - }; - let Some(kind) = reference_kind_for_attribute(AttributeScope::Node, key, value) - else { - continue; - }; - kind.validate(value) - .map_err(|source| WorkflowVersionError::StaticReference { - path: graph_path.clone(), - source, - })?; - - match kind { - ReferenceKind::Import => { - let target = Self::resolve_reference(graph_path, kind, value)?; - self.require_file(graph_path, kind, target.clone())?; - imports.push_back(target); - } - ReferenceKind::ChildWorkflow => { - let target = Self::resolve_reference(graph_path, kind, value)?; - child_workflows.insert(target); - } - ReferenceKind::FileInline => { - let reference = value.strip_prefix('@').ok_or_else(|| { - WorkflowVersionError::MissingFile { - path: graph_path.clone(), - kind, - target: graph_path.clone(), - } - })?; - let target = Self::resolve_reference(graph_path, kind, reference)?; - let content = self.require_file(graph_path, kind, target.clone())?; - if key == "prompt" { - Self::validate_template( - &target, - content, - template_store, - template_root, - )?; - } - } - ReferenceKind::Dockerfile | ReferenceKind::GraphGoalFile => {} } - } + })?; - if let Some(prompt) = node.prompt().filter(|prompt| !prompt.starts_with('@')) { - Self::validate_template(graph_path, prompt, template_store, template_root)?; + match kind { + ReferenceKind::Import => { + let target = resolve_reference(graph_path, kind, value)?; + require_file(version, graph_path, kind, target.clone())?; + imports.push_back(target); + } + ReferenceKind::ChildWorkflow => { + let target = resolve_reference(graph_path, kind, value)?; + child_workflows.insert(target); + } + ReferenceKind::FileInline => { + let reference = value.strip_prefix('@').ok_or_else(|| { + WorkflowVersionError::MissingFile { + path: graph_path.clone(), + kind, + target: graph_path.clone(), + } + })?; + let target = resolve_reference(graph_path, kind, reference)?; + let content = require_file(version, graph_path, kind, target.clone())?; + if key == "prompt" { + validate_template(&target, content, template_store, template_root)?; + } + } + ReferenceKind::Dockerfile | ReferenceKind::GraphGoalFile => {} } } - Ok(()) - } - fn validate_template( - path: &WorkflowPath, - content: &str, - store: &BundleTemplateStore, - root: &ManifestPath, - ) -> Result<(), WorkflowVersionError> { - let manifest_path = manifest_path(path); - discover_static_dependency_closure( - [TemplateSource::new(manifest_path, root.clone(), content)], - store, - ) - .map_err(|source| WorkflowVersionError::Template { - path: path.clone(), - source: Box::new(source), - })?; - Ok(()) + if let Some(prompt) = node.prompt().filter(|prompt| !prompt.starts_with('@')) { + validate_template(graph_path, prompt, template_store, template_root)?; + } } + Ok(()) +} - fn template_store(&self) -> BundleTemplateStore { - BundleTemplateStore::new( - self.files - .iter() - .map(|(path, content)| (manifest_path(path), content.clone())) - .collect::>(), - ) - } +fn validate_template( + path: &WorkflowPath, + content: &str, + store: &BundleTemplateStore, + root: &ManifestPath, +) -> Result<(), WorkflowVersionError> { + let manifest_path = manifest_path(path); + discover_static_dependency_closure( + [TemplateSource::new(manifest_path, root.clone(), content)], + store, + ) + .map_err(|source| WorkflowVersionError::Template { + path: path.clone(), + source: Box::new(source), + })?; + Ok(()) +} - fn resolve_reference( - path: &WorkflowPath, - kind: ReferenceKind, - reference: &str, - ) -> Result { - path.resolve_reference(reference) - .map_err(|source| WorkflowVersionError::InvalidReference { - path: path.clone(), - kind, - reference: reference.to_owned(), - source, - }) - } +fn template_store(version: &WorkflowVersion) -> BundleTemplateStore { + BundleTemplateStore::new( + version + .files() + .iter() + .map(|(path, content)| (manifest_path(path), content.clone())) + .collect::>(), + ) +} - fn require_file( - &self, - path: &WorkflowPath, - kind: ReferenceKind, - target: WorkflowPath, - ) -> Result<&str, WorkflowVersionError> { - self.files.get(&target).map(String::as_str).ok_or_else(|| { - WorkflowVersionError::MissingFile { - path: path.clone(), - kind, - target, - } +fn resolve_reference( + path: &WorkflowPath, + kind: ReferenceKind, + reference: &str, +) -> Result { + path.resolve_reference(reference) + .map_err(|source| WorkflowVersionError::InvalidReference { + path: path.clone(), + kind, + reference: reference.to_owned(), + source, + }) +} + +fn require_file<'version>( + version: &'version WorkflowVersion, + path: &WorkflowPath, + kind: ReferenceKind, + target: WorkflowPath, +) -> Result<&'version str, WorkflowVersionError> { + version + .files() + .get(&target) + .map(String::as_str) + .ok_or_else(|| WorkflowVersionError::MissingFile { + path: path.clone(), + kind, + target, }) - } } fn manifest_path(path: &WorkflowPath) -> ManifestPath { @@ -461,76 +370,11 @@ fn manifest_path(path: &WorkflowPath) -> ManifestPath { .expect("validated workflow paths must also be valid manifest paths") } -impl<'de> Deserialize<'de> for WorkflowVersion { - fn deserialize(deserializer: D) -> Result - where - D: Deserializer<'de>, - { - #[derive(Deserialize)] - #[serde(deny_unknown_fields)] - struct Wire { - entrypoint: WorkflowPath, - files: UniqueBTreeMap, - dependencies: UniqueBTreeMap, - } - - let wire = Wire::deserialize(deserializer)?; - Self::new(wire.entrypoint, wire.files.0, wire.dependencies.0).map_err(D::Error::custom) - } -} - -struct UniqueBTreeMap(BTreeMap); - -impl<'de, K, V> Deserialize<'de> for UniqueBTreeMap -where - K: Deserialize<'de> + Ord + fmt::Display, - V: Deserialize<'de>, -{ - fn deserialize(deserializer: D) -> Result - where - D: Deserializer<'de>, - { - struct MapVisitor(PhantomData<(K, V)>); - - impl<'de, K, V> Visitor<'de> for MapVisitor - where - K: Deserialize<'de> + Ord + fmt::Display, - V: Deserialize<'de>, - { - type Value = UniqueBTreeMap; - - fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { - formatter.write_str("a map with unique keys") - } - - fn visit_map(self, mut access: A) -> Result - where - A: MapAccess<'de>, - { - let mut values = BTreeMap::new(); - while let Some((key, value)) = access.next_entry::()? { - if values.insert(key, value).is_some() { - return Err(A::Error::custom("duplicate workflow map key")); - } - } - Ok(UniqueBTreeMap(values)) - } - } - - deserializer.deserialize_map(MapVisitor(PhantomData)) - } -} - #[cfg(test)] mod tests { - use std::collections::BTreeMap; + use fabro_types::{RunBlobId, WorkflowPath, WorkflowVersion, WorkflowVersionId}; - use fabro_types::{RunBlobId, WorkflowPath, WorkflowVersionId}; - - use super::{ - MAX_WORKFLOW_VERSION_BYTES, MAX_WORKFLOW_VERSION_FILE_BYTES, MAX_WORKFLOW_VERSION_FILES, - WorkflowVersion, WorkflowVersionError, - }; + use super::{ValidatedWorkflowVersion, WorkflowVersionError}; fn path(value: &str) -> WorkflowPath { value.parse().unwrap() @@ -543,43 +387,23 @@ mod tests { fn version_with( files: impl IntoIterator, dependencies: impl IntoIterator, - ) -> Result { - WorkflowVersion::new( - path("workflow.fabro"), - files - .into_iter() - .map(|(path_value, content)| (path(path_value), content.to_owned())) - .collect(), - dependencies - .into_iter() - .map(|(path_value, id)| (path(path_value), id)) - .collect(), + ) -> Result { + ValidatedWorkflowVersion::new( + WorkflowVersion::new( + path("workflow.fabro"), + files + .into_iter() + .map(|(path_value, content)| (path(path_value), content.to_owned())) + .collect(), + dependencies + .into_iter() + .map(|(path_value, id)| (path(path_value), id)) + .collect(), + ) + .expect("test fixtures must be structurally valid"), ) } - #[test] - fn canonical_bytes_have_fixed_field_and_map_order() { - let version = WorkflowVersion::new( - path("workflow.fabro"), - BTreeMap::from([ - (path("z.txt"), "Z".to_string()), - ( - path("workflow.fabro"), - "digraph W { start [shape=Mdiamond] exit [shape=Msquare] start -> exit }" - .to_string(), - ), - (path("a.txt"), "A".to_string()), - ]), - BTreeMap::new(), - ) - .unwrap(); - - assert_eq!( - String::from_utf8(version.canonical_bytes().unwrap()).unwrap(), - r#"{"entrypoint":"workflow.fabro","files":{"a.txt":"A","workflow.fabro":"digraph W { start [shape=Mdiamond] exit [shape=Msquare] start -> exit }","z.txt":"Z"},"dependencies":{}}"# - ); - } - #[test] fn validates_imports_templates_file_refs_and_dependencies() { let version = version_with( @@ -606,7 +430,7 @@ mod tests { ) .unwrap(); - assert_eq!(version.dependencies().len(), 1); + assert_eq!(version.version().dependencies().len(), 1); } #[test] @@ -706,121 +530,17 @@ dockerfile = { path = "docker/run.Dockerfile" } ) .unwrap(); - assert_eq!(version.entrypoint(), &path("workflow.fabro")); - } - - #[test] - fn rejects_path_collisions_and_large_files() { - let collision = version_with( - [ - ("workflow.fabro", "digraph W {}"), - ("assets", "file"), - ("assets/item.txt", "nested"), - ], - [], - ) - .unwrap_err(); - assert!(matches!( - collision, - WorkflowVersionError::PathCollision { .. } - )); - - let mut files = BTreeMap::from([(path("workflow.fabro"), "digraph W {}".to_string())]); - files.insert( - path("large.txt"), - "x".repeat(MAX_WORKFLOW_VERSION_FILE_BYTES + 1), - ); - let large = - WorkflowVersion::new(path("workflow.fabro"), files, BTreeMap::new()).unwrap_err(); - assert!(matches!(large, WorkflowVersionError::FileTooLarge { .. })); - } - - #[test] - fn enforces_file_count_file_size_and_canonical_size_boundaries() { - let mut files = BTreeMap::from([(path("workflow.fabro"), "digraph W {}".to_string())]); - for index in 0..MAX_WORKFLOW_VERSION_FILES - 1 { - files.insert(path(&format!("file-{index:03}.txt")), String::new()); - } - assert!( - WorkflowVersion::new(path("workflow.fabro"), files.clone(), BTreeMap::new()).is_ok() - ); - files.insert(path("too-many.txt"), String::new()); - assert!(matches!( - WorkflowVersion::new(path("workflow.fabro"), files, BTreeMap::new()).unwrap_err(), - WorkflowVersionError::TooManyFiles { .. } - )); - - let exact_file = BTreeMap::from([ - (path("workflow.fabro"), "digraph W {}".to_string()), - ( - path("payload.txt"), - "x".repeat(MAX_WORKFLOW_VERSION_FILE_BYTES), - ), - ]); - assert!( - WorkflowVersion::new(path("workflow.fabro"), exact_file.clone(), BTreeMap::new()) - .is_ok() - ); - let mut oversized_file = exact_file; - oversized_file - .get_mut(&path("payload.txt")) - .unwrap() - .push('x'); - assert!(matches!( - WorkflowVersion::new(path("workflow.fabro"), oversized_file, BTreeMap::new()) - .unwrap_err(), - WorkflowVersionError::FileTooLarge { .. } - )); - - let mut exact_version_files = - BTreeMap::from([(path("workflow.fabro"), "digraph W {}".to_string())]); - for index in 0..4 { - exact_version_files.insert(path(&format!("payload-{index}.txt")), String::new()); - } - let empty = WorkflowVersion::new( - path("workflow.fabro"), - exact_version_files.clone(), - BTreeMap::new(), - ) - .unwrap(); - let remaining = MAX_WORKFLOW_VERSION_BYTES - empty.canonical_bytes().unwrap().len(); - let per_file = remaining / 4; - let remainder = remaining % 4; - for index in 0..4 { - let length = per_file + usize::from(index < remainder); - assert!(length <= MAX_WORKFLOW_VERSION_FILE_BYTES); - exact_version_files.insert(path(&format!("payload-{index}.txt")), "x".repeat(length)); - } - let exact_version = WorkflowVersion::new( - path("workflow.fabro"), - exact_version_files.clone(), - BTreeMap::new(), - ) - .unwrap(); - assert_eq!( - exact_version.canonical_bytes().unwrap().len(), - MAX_WORKFLOW_VERSION_BYTES - ); - exact_version_files - .get_mut(&path("payload-0.txt")) - .unwrap() - .push('x'); - assert!(matches!( - WorkflowVersion::new(path("workflow.fabro"), exact_version_files, BTreeMap::new()) - .unwrap_err(), - WorkflowVersionError::VersionTooLarge { .. } - )); + assert_eq!(version.version().entrypoint(), &path("workflow.fabro")); } #[test] fn rejects_escaping_and_dynamic_template_references() { - let escaping = WorkflowVersion::new( - path("workflow.fabro"), - BTreeMap::from([( - path("workflow.fabro"), - r#"digraph W { imported [import="../outside.fabro"] }"#.to_string(), - )]), - BTreeMap::new(), + let escaping = version_with( + [( + "workflow.fabro", + r#"digraph W { imported [import="../outside.fabro"] }"#, + )], + [], ) .unwrap_err(); assert!(matches!( @@ -828,33 +548,14 @@ dockerfile = { path = "docker/run.Dockerfile" } WorkflowVersionError::InvalidReference { .. } )); - let dynamic = WorkflowVersion::new( - path("workflow.fabro"), - BTreeMap::from([( - path("workflow.fabro"), - r#"digraph W { step [prompt="{% include template_name %}"] }"#.to_string(), - )]), - BTreeMap::new(), + let dynamic = version_with( + [( + "workflow.fabro", + r#"digraph W { step [prompt="{% include template_name %}"] }"#, + )], + [], ) .unwrap_err(); assert!(matches!(dynamic, WorkflowVersionError::Template { .. })); } - - #[test] - fn deserialize_rejects_unknown_fields_and_duplicate_keys() { - let unknown = r#"{ - "entrypoint":"workflow.fabro", - "files":{"workflow.fabro":"digraph W {}"}, - "dependencies":{}, - "metadata":{} - }"#; - assert!(serde_json::from_str::(unknown).is_err()); - - let duplicate = r#"{ - "entrypoint":"workflow.fabro", - "files":{"workflow.fabro":"digraph W {}","workflow.fabro":"digraph X {}"}, - "dependencies":{} - }"#; - assert!(serde_json::from_str::(duplicate).is_err()); - } } diff --git a/lib/components/fabro-store/src/workflow_version_store.rs b/lib/components/fabro-workflow-version/src/store.rs similarity index 78% rename from lib/components/fabro-store/src/workflow_version_store.rs rename to lib/components/fabro-workflow-version/src/store.rs index 4101487d6..314b88388 100644 --- a/lib/components/fabro-store/src/workflow_version_store.rs +++ b/lib/components/fabro-workflow-version/src/store.rs @@ -1,15 +1,17 @@ use std::sync::Arc; -use fabro_types::{WorkflowPath, WorkflowVersionId}; -use fabro_workflow_version::{WorkflowVersion, WorkflowVersionError}; +use fabro_store::BlobStore; +use fabro_types::{WorkflowPath, WorkflowVersion, WorkflowVersionId, WorkflowVersionShapeError}; use thiserror::Error; -use crate::BlobStore; +use crate::{ValidatedWorkflowVersion, WorkflowVersionError}; #[derive(Debug, Error)] pub enum WorkflowVersionStoreError { #[error(transparent)] InvalidVersion(#[from] WorkflowVersionError), + #[error(transparent)] + InvalidShape(#[from] WorkflowVersionShapeError), #[error("workflow-version dependency `{id}` at `{path}` is not stored")] DependencyNotFound { path: WorkflowPath, @@ -33,10 +35,15 @@ pub enum WorkflowVersionStoreError { #[error("workflow-version storage operation failed")] Storage { #[source] - source: crate::Error, + source: fabro_store::Error, }, } +/// Content-addressed storage for validated workflow versions. +/// +/// `put` only accepts semantically validated versions; `get` re-validates +/// blobs on read because the blob namespace is shared and storage is not +/// trusted to contain only canonical versions. #[derive(Clone, Debug)] pub struct WorkflowVersionStore { blobs: Arc, @@ -50,10 +57,10 @@ impl WorkflowVersionStore { pub async fn put( &self, - version: &WorkflowVersion, + version: &ValidatedWorkflowVersion, ) -> Result { - let canonical = version.canonical_bytes()?; - for (path, id) in version.dependencies() { + let canonical = version.version().canonical_bytes()?; + for (path, id) in version.version().dependencies() { match self.get(id).await { Ok(Some(_)) => {} Ok(None) => { @@ -81,7 +88,7 @@ impl WorkflowVersionStore { pub async fn get( &self, id: &WorkflowVersionId, - ) -> Result, WorkflowVersionStoreError> { + ) -> Result, WorkflowVersionStoreError> { let blob_id = (*id).into(); let Some(bytes) = self .blobs @@ -93,11 +100,12 @@ impl WorkflowVersionStore { }; let version = serde_json::from_slice::(&bytes) .map_err(|source| WorkflowVersionStoreError::Decode { id: *id, source })?; - let canonical = version.canonical_bytes()?; + let validated = ValidatedWorkflowVersion::new(version)?; + let canonical = validated.version().canonical_bytes()?; if canonical.as_slice() != bytes.as_ref() { return Err(WorkflowVersionStoreError::NonCanonical { id: *id }); } - Ok(Some(version)) + Ok(Some(validated)) } } @@ -107,12 +115,12 @@ mod tests { use std::sync::Arc; use std::time::Duration; - use fabro_types::{WorkflowPath, WorkflowVersionId}; - use fabro_workflow_version::WorkflowVersion; + use fabro_store::{BlobStore, Database}; + use fabro_types::{WorkflowPath, WorkflowVersion, WorkflowVersionId}; use object_store::memory::InMemory; use super::{WorkflowVersionStore, WorkflowVersionStoreError}; - use crate::Database; + use crate::ValidatedWorkflowVersion; fn path(value: &str) -> WorkflowPath { value.parse().unwrap() @@ -121,16 +129,19 @@ mod tests { fn version( graph: &str, dependencies: BTreeMap, - ) -> WorkflowVersion { - WorkflowVersion::new( - path("workflow.fabro"), - BTreeMap::from([(path("workflow.fabro"), graph.to_owned())]), - dependencies, + ) -> ValidatedWorkflowVersion { + ValidatedWorkflowVersion::new( + WorkflowVersion::new( + path("workflow.fabro"), + BTreeMap::from([(path("workflow.fabro"), graph.to_owned())]), + dependencies, + ) + .unwrap(), ) .unwrap() } - async fn stores() -> (Arc, WorkflowVersionStore) { + async fn stores() -> (Arc, WorkflowVersionStore) { let database = Database::new( Arc::new(InMemory::new()), "", @@ -146,7 +157,7 @@ mod tests { async fn put_get_reuses_exact_blob_digest() { let (blobs, store) = stores().await; let version = version("digraph W {}", BTreeMap::new()); - let expected_bytes = version.canonical_bytes().unwrap(); + let expected_bytes = version.version().canonical_bytes().unwrap(); let expected_id = WorkflowVersionId::from(fabro_types::RunBlobId::new(&expected_bytes)); let id = store.put(&version).await.unwrap(); @@ -178,14 +189,14 @@ mod tests { let (blobs, store) = stores().await; let child = version("digraph Child {}", BTreeMap::new()); let child_id = WorkflowVersionId::from(fabro_types::RunBlobId::new( - &child.canonical_bytes().unwrap(), + &child.version().canonical_bytes().unwrap(), )); let root = version( r#"digraph Root { child [stack.child_workflow="child.fabro"] }"#, BTreeMap::from([(path("child.fabro"), child_id)]), ); let root_id = WorkflowVersionId::from(fabro_types::RunBlobId::new( - &root.canonical_bytes().unwrap(), + &root.version().canonical_bytes().unwrap(), )); let error = store.put(&root).await.unwrap_err(); @@ -215,7 +226,7 @@ mod tests { )); let version = version("digraph W {}", BTreeMap::new()); - let pretty = serde_json::to_vec_pretty(&version).unwrap(); + let pretty = serde_json::to_vec_pretty(version.version()).unwrap(); let noncanonical = WorkflowVersionId::from(blobs.write(&pretty).await.unwrap()); assert!(matches!( store.get(&noncanonical).await.unwrap_err(), diff --git a/lib/components/fabro-workflow/src/handler/manager_loop.rs b/lib/components/fabro-workflow/src/handler/manager_loop.rs index 6d58f4cc5..8cc82f634 100644 --- a/lib/components/fabro-workflow/src/handler/manager_loop.rs +++ b/lib/components/fabro-workflow/src/handler/manager_loop.rs @@ -5,9 +5,10 @@ use std::time::Duration; use async_trait::async_trait; use fabro_graphviz::graph::{AttrValue, Graph, Node}; -use fabro_graphviz::static_reference::{ReferenceKind, validate_static_reference}; use fabro_store::{ArtifactStore, Database}; +use fabro_template::validate_static_reference; use fabro_types::WorkflowSettings; +use fabro_types::graph::ReferenceKind; use object_store::memory::InMemory; use tokio::fs; use tokio::time::{sleep, timeout}; diff --git a/lib/components/fabro-workflow/src/transforms/import.rs b/lib/components/fabro-workflow/src/transforms/import.rs index fe42079dd..bab7ee1b1 100644 --- a/lib/components/fabro-workflow/src/transforms/import.rs +++ b/lib/components/fabro-workflow/src/transforms/import.rs @@ -4,8 +4,8 @@ use std::sync::Arc; use fabro_graphviz::graph::{AttrValue, Edge, Graph, Node}; use fabro_graphviz::parser; -use fabro_graphviz::static_reference::{ReferenceKind, validate_static_reference}; -use fabro_template::TemplateContext; +use fabro_template::{TemplateContext, validate_static_reference}; +use fabro_types::graph::ReferenceKind; use fabro_validate::Diagnostic; use super::file_inlining::template_render_store; diff --git a/lib/components/fabro-workflow/src/transforms/importable_field.rs b/lib/components/fabro-workflow/src/transforms/importable_field.rs index 51978e845..cf4b5c2d1 100644 --- a/lib/components/fabro-workflow/src/transforms/importable_field.rs +++ b/lib/components/fabro-workflow/src/transforms/importable_field.rs @@ -13,7 +13,8 @@ //! [`super::file_inlining`], where the `FileResolver` and current-dir context //! live. -use fabro_graphviz::static_reference::{ReferenceKind, validate_static_reference}; +use fabro_template::validate_static_reference; +use fabro_types::graph::ReferenceKind; use crate::error::Error; diff --git a/lib/components/fabro-workflow/src/transforms/variable_expansion.rs b/lib/components/fabro-workflow/src/transforms/variable_expansion.rs index 642a1557c..da7aba76d 100644 --- a/lib/components/fabro-workflow/src/transforms/variable_expansion.rs +++ b/lib/components/fabro-workflow/src/transforms/variable_expansion.rs @@ -4,13 +4,11 @@ use std::fmt::Write as _; use std::sync::Arc; use fabro_graphviz::graph::{AttrValue, Graph, Node}; -use fabro_graphviz::static_reference::{ - AttributeScope, ReferenceKind, reference_kind_for_attribute, validate_static_reference, -}; use fabro_template::{ TemplateContext, TemplateError, TemplateRenderMode, TemplateSource, TemplateSourceOrigin, - TemplateStore, + TemplateStore, validate_static_reference, }; +use fabro_types::graph::{AttributeScope, ReferenceKind, reference_kind_for_attribute}; use fabro_types::settings::interp::Namespace; use fabro_types::settings::{InterpString, ResolveCtx, ResolveError, ResolveErrorKind}; use fabro_util::error::collect_chain; diff --git a/lib/foundation/fabro-api/Cargo.toml b/lib/foundation/fabro-api/Cargo.toml index 70294db9a..bc28670af 100644 --- a/lib/foundation/fabro-api/Cargo.toml +++ b/lib/foundation/fabro-api/Cargo.toml @@ -20,7 +20,6 @@ fabro-config = { path = "../fabro-config" } fabro-environment.workspace = true fabro-model = { path = "../fabro-model" } fabro-types = { path = "../fabro-types" } -fabro-workflow-version = { path = "../../components/fabro-workflow-version" } progenitor-client = "0.13" regress = "0.10" reqwest.workspace = true diff --git a/lib/foundation/fabro-api/build.rs b/lib/foundation/fabro-api/build.rs index 9f1a0c2ca..ea8ece1eb 100644 --- a/lib/foundation/fabro-api/build.rs +++ b/lib/foundation/fabro-api/build.rs @@ -722,11 +722,7 @@ fn main() { ("CompletionMessage", "fabro_types::Message", &[]), ("CompletionMessageRole", "fabro_types::Role", &[]), ("CompletionContentPart", "fabro_types::ContentPart", &[]), - ( - "WorkflowVersion", - "fabro_workflow_version::WorkflowVersion", - &[], - ), + ("WorkflowVersion", "fabro_types::WorkflowVersion", &[]), ("WorkflowPath", "fabro_types::WorkflowPath", &[]), ("WorkflowVersionId", "fabro_types::WorkflowVersionId", &[]), ("CostSource", "fabro_model::CostSource", &[]), diff --git a/lib/foundation/fabro-api/src/lib.rs b/lib/foundation/fabro-api/src/lib.rs index 0a26ed039..b40087831 100644 --- a/lib/foundation/fabro-api/src/lib.rs +++ b/lib/foundation/fabro-api/src/lib.rs @@ -74,9 +74,8 @@ pub mod types { SubAgentProjection, SubAgentStatus, SystemActorKind, SystemIntegrationStatus, SystemIntegrationsResponse, TodoListProjection, TurnId, UpdateVariableRequest, UserPrincipal, Variable, VariableListResponse, WorkflowPath, WorkflowSettings, - WorkflowVersionId, + WorkflowVersion, WorkflowVersionId, }; - pub use fabro_workflow_version::WorkflowVersion; pub use crate::generated::types::*; } diff --git a/lib/foundation/fabro-api/tests/workflow_version_round_trip.rs b/lib/foundation/fabro-api/tests/workflow_version_round_trip.rs index 2101bba73..42e9ef5f0 100644 --- a/lib/foundation/fabro-api/tests/workflow_version_round_trip.rs +++ b/lib/foundation/fabro-api/tests/workflow_version_round_trip.rs @@ -4,8 +4,7 @@ use fabro_api::types::{ WorkflowPath as ApiWorkflowPath, WorkflowVersion as ApiWorkflowVersion, WorkflowVersionId as ApiWorkflowVersionId, }; -use fabro_types::{WorkflowPath, WorkflowVersionId}; -use fabro_workflow_version::WorkflowVersion; +use fabro_types::{WorkflowPath, WorkflowVersion, WorkflowVersionId}; use serde_json::json; #[test] diff --git a/lib/foundation/fabro-template/src/lib.rs b/lib/foundation/fabro-template/src/lib.rs index f312e9416..0a77f18e6 100644 --- a/lib/foundation/fabro-template/src/lib.rs +++ b/lib/foundation/fabro-template/src/lib.rs @@ -3,6 +3,7 @@ use std::fmt; use std::sync::{Arc, Mutex}; use fabro_types::ManifestPath; +use fabro_types::graph::ReferenceKind; use miette::{LabeledSpan, NamedSource, SourceCode, SourceSpan}; use minijinja::value::{Object, Value}; use minijinja::{AutoEscape, Environment, ErrorKind, UndefinedBehavior}; @@ -524,6 +525,46 @@ pub fn contains_template_syntax(template: &str) -> bool { template.contains("{{") || template.contains("{%") || template.contains("{#") } +/// A static file reference that unexpectedly contains template syntax. +#[derive(Debug, thiserror::Error)] +#[error("templates are not supported in {kind}s: {value}")] +pub struct StaticReferenceError { + kind: ReferenceKind, + value: String, +} + +impl StaticReferenceError { + #[must_use] + pub fn new(kind: ReferenceKind, value: impl Into) -> Self { + Self { + kind, + value: value.into(), + } + } + + #[must_use] + pub fn kind(&self) -> ReferenceKind { + self.kind + } + + #[must_use] + pub fn value(&self) -> &str { + &self.value + } +} + +/// Reject static file references (imports, child workflows, `@` file values) +/// that contain template syntax. +pub fn validate_static_reference( + value: &str, + kind: ReferenceKind, +) -> Result<(), StaticReferenceError> { + if contains_template_syntax(value) { + return Err(StaticReferenceError::new(kind, value)); + } + Ok(()) +} + /// Whether `template` references `name` as a top-level variable. /// /// Used by the goal self-reference lint: a graph `goal` may not reference @@ -793,6 +834,27 @@ mod tests { use super::*; + #[test] + fn static_reference_rejects_template_syntax() { + let error = validate_static_reference( + "@schemas/{{ inputs.schema }}.json", + ReferenceKind::FileInline, + ) + .unwrap_err(); + + assert_eq!(error.kind(), ReferenceKind::FileInline); + assert_eq!(error.value(), "@schemas/{{ inputs.schema }}.json"); + assert!( + error + .to_string() + .contains("templates are not supported in file inline references"), + "unexpected error: {error}", + ); + assert!( + validate_static_reference("@schemas/result.json", ReferenceKind::FileInline).is_ok() + ); + } + fn manifest_path(value: &str) -> ManifestPath { ManifestPath::from_wire(value).expect("path should parse") } diff --git a/lib/foundation/fabro-types/src/graph.rs b/lib/foundation/fabro-types/src/graph.rs index e2ed3fab4..7ed575ac5 100644 --- a/lib/foundation/fabro-types/src/graph.rs +++ b/lib/foundation/fabro-types/src/graph.rs @@ -590,6 +590,52 @@ impl Graph { } } +/// Where an attribute appears in a workflow graph. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum AttributeScope { + Graph, + Node, + Edge, +} + +/// Kinds of static (non-templated) file references a graph attribute can +/// carry. +#[derive(Clone, Copy, Debug, Eq, PartialEq, strum::Display)] +pub enum ReferenceKind { + #[strum(to_string = "file inline reference")] + FileInline, + #[strum(to_string = "import reference")] + Import, + #[strum(to_string = "child workflow reference")] + ChildWorkflow, + #[strum(to_string = "Dockerfile reference")] + Dockerfile, + #[strum(to_string = "graph goal file reference")] + GraphGoalFile, +} + +/// Classify a graph attribute as a static file reference, if it is one. +#[must_use] +pub fn reference_kind_for_attribute( + scope: AttributeScope, + key: &str, + value: &str, +) -> Option { + match key { + "import" => Some(ReferenceKind::Import), + "stack.child_workflow" | "stack.child_dotfile" => Some(ReferenceKind::ChildWorkflow), + "goal" if matches!(scope, AttributeScope::Graph) && value.starts_with('@') => { + Some(ReferenceKind::GraphGoalFile) + } + "prompt" | "output_schema" + if matches!(scope, AttributeScope::Node) && value.starts_with('@') => + { + Some(ReferenceKind::FileInline) + } + _ => None, + } +} + #[cfg(test)] mod tests { use super::*; @@ -1111,4 +1157,24 @@ mod tests { ); assert_eq!(g.loop_restart_signature_limit(), 3); } + + #[test] + fn output_schema_at_value_is_file_inline_reference() { + assert_eq!( + reference_kind_for_attribute( + AttributeScope::Node, + "output_schema", + "@schemas/result.schema.json", + ), + Some(ReferenceKind::FileInline), + ); + } + + #[test] + fn output_schema_builtin_keyword_is_not_file_inline_reference() { + assert_eq!( + reference_kind_for_attribute(AttributeScope::Node, "output_schema", "routing"), + None, + ); + } } diff --git a/lib/foundation/fabro-types/src/lib.rs b/lib/foundation/fabro-types/src/lib.rs index 28dc93996..c7324a1d9 100644 --- a/lib/foundation/fabro-types/src/lib.rs +++ b/lib/foundation/fabro-types/src/lib.rs @@ -55,6 +55,7 @@ pub mod todo; pub mod transcript; pub mod variable; pub mod workflow_path; +pub mod workflow_version; pub mod workflow_version_id; pub use artifact::ArtifactUpload; @@ -188,4 +189,8 @@ pub use variable::{ pub use workflow_path::{ MAX_WORKFLOW_PATH_BYTES, MAX_WORKFLOW_PATH_COMPONENTS, WorkflowPath, WorkflowPathParseError, }; +pub use workflow_version::{ + MAX_WORKFLOW_VERSION_BYTES, MAX_WORKFLOW_VERSION_FILE_BYTES, MAX_WORKFLOW_VERSION_FILES, + WorkflowVersion, WorkflowVersionShapeError, +}; pub use workflow_version_id::{WorkflowVersionId, WorkflowVersionIdParseError}; diff --git a/lib/foundation/fabro-types/src/workflow_version.rs b/lib/foundation/fabro-types/src/workflow_version.rs new file mode 100644 index 000000000..ff376bbab --- /dev/null +++ b/lib/foundation/fabro-types/src/workflow_version.rs @@ -0,0 +1,379 @@ +use std::collections::BTreeMap; +use std::fmt; +use std::marker::PhantomData; + +use serde::de::{Error as _, MapAccess, Visitor}; +use serde::{Deserialize, Deserializer, Serialize}; +use thiserror::Error; + +use crate::{WorkflowPath, WorkflowVersionId}; + +pub const MAX_WORKFLOW_VERSION_FILES: usize = 512; +pub const MAX_WORKFLOW_VERSION_FILE_BYTES: usize = 512 * 1024; +pub const MAX_WORKFLOW_VERSION_BYTES: usize = 2 * 1024 * 1024; + +#[derive(Debug, Error)] +pub enum WorkflowVersionShapeError { + #[error("workflow version has {actual} files; maximum is {maximum}")] + TooManyFiles { actual: usize, maximum: usize }, + #[error("workflow file `{path}` is {actual} bytes; maximum is {maximum}")] + FileTooLarge { + path: WorkflowPath, + actual: usize, + maximum: usize, + }, + #[error("workflow version is {actual} canonical bytes; maximum is {maximum}")] + VersionTooLarge { actual: usize, maximum: usize }, + #[error("entrypoint `{path}` is not present in workflow files")] + MissingEntrypoint { path: WorkflowPath }, + #[error("workflow paths collide: `{first}` and `{second}`")] + PathCollision { + first: WorkflowPath, + second: WorkflowPath, + }, + #[error("failed to serialize canonical workflow version")] + Serialization { + #[source] + source: serde_json::Error, + }, +} + +/// Canonical wire form of an immutable workflow version. +/// +/// Construction (and therefore deserialization) enforces the structural +/// invariants: file-count and byte-size limits, entrypoint presence, unique +/// map keys, and collision-free paths. Semantic validation of graph, config, +/// and template content is a separate concern owned by +/// `fabro-workflow-version`. +#[derive(Clone, Debug, PartialEq, Eq, Serialize)] +pub struct WorkflowVersion { + entrypoint: WorkflowPath, + files: BTreeMap, + dependencies: BTreeMap, +} + +impl WorkflowVersion { + pub fn new( + entrypoint: WorkflowPath, + files: BTreeMap, + dependencies: BTreeMap, + ) -> Result { + let version = Self { + entrypoint, + files, + dependencies, + }; + version.validate_shape()?; + version.canonical_bytes()?; + Ok(version) + } + + #[must_use] + pub fn entrypoint(&self) -> &WorkflowPath { + &self.entrypoint + } + + #[must_use] + pub fn files(&self) -> &BTreeMap { + &self.files + } + + #[must_use] + pub fn dependencies(&self) -> &BTreeMap { + &self.dependencies + } + + /// Serialize to the canonical wire form. + /// + /// Structural validity is guaranteed by construction, so this only + /// serializes and enforces the canonical size limit. + pub fn canonical_bytes(&self) -> Result, WorkflowVersionShapeError> { + let bytes = serde_json::to_vec(self) + .map_err(|source| WorkflowVersionShapeError::Serialization { source })?; + if bytes.len() > MAX_WORKFLOW_VERSION_BYTES { + return Err(WorkflowVersionShapeError::VersionTooLarge { + actual: bytes.len(), + maximum: MAX_WORKFLOW_VERSION_BYTES, + }); + } + Ok(bytes) + } + + fn validate_shape(&self) -> Result<(), WorkflowVersionShapeError> { + if self.files.len() > MAX_WORKFLOW_VERSION_FILES { + return Err(WorkflowVersionShapeError::TooManyFiles { + actual: self.files.len(), + maximum: MAX_WORKFLOW_VERSION_FILES, + }); + } + for (path, content) in &self.files { + if content.len() > MAX_WORKFLOW_VERSION_FILE_BYTES { + return Err(WorkflowVersionShapeError::FileTooLarge { + path: path.clone(), + actual: content.len(), + maximum: MAX_WORKFLOW_VERSION_FILE_BYTES, + }); + } + } + if !self.files.contains_key(&self.entrypoint) { + return Err(WorkflowVersionShapeError::MissingEntrypoint { + path: self.entrypoint.clone(), + }); + } + self.validate_path_collisions() + } + + fn validate_path_collisions(&self) -> Result<(), WorkflowVersionShapeError> { + // Keys are unique within each map, so equality can only collide + // across files and dependencies. + let paths = self + .files + .keys() + .chain(self.dependencies.keys()) + .collect::>(); + for (index, first) in paths.iter().enumerate() { + for second in &paths[index + 1..] { + if first == second || first.is_ancestor_of(second) || second.is_ancestor_of(first) { + return Err(WorkflowVersionShapeError::PathCollision { + first: (*first).clone(), + second: (*second).clone(), + }); + } + } + } + Ok(()) + } +} + +impl<'de> Deserialize<'de> for WorkflowVersion { + fn deserialize(deserializer: D) -> Result + where + D: Deserializer<'de>, + { + #[derive(Deserialize)] + #[serde(deny_unknown_fields)] + struct Wire { + entrypoint: WorkflowPath, + files: UniqueBTreeMap, + dependencies: UniqueBTreeMap, + } + + let wire = Wire::deserialize(deserializer)?; + Self::new(wire.entrypoint, wire.files.0, wire.dependencies.0).map_err(D::Error::custom) + } +} + +struct UniqueBTreeMap(BTreeMap); + +impl<'de, K, V> Deserialize<'de> for UniqueBTreeMap +where + K: Deserialize<'de> + Ord + fmt::Display, + V: Deserialize<'de>, +{ + fn deserialize(deserializer: D) -> Result + where + D: Deserializer<'de>, + { + struct MapVisitor(PhantomData<(K, V)>); + + impl<'de, K, V> Visitor<'de> for MapVisitor + where + K: Deserialize<'de> + Ord + fmt::Display, + V: Deserialize<'de>, + { + type Value = UniqueBTreeMap; + + fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("a map with unique keys") + } + + fn visit_map(self, mut access: A) -> Result + where + A: MapAccess<'de>, + { + let mut values = BTreeMap::new(); + while let Some((key, value)) = access.next_entry::()? { + if values.insert(key, value).is_some() { + return Err(A::Error::custom("duplicate workflow map key")); + } + } + Ok(UniqueBTreeMap(values)) + } + } + + deserializer.deserialize_map(MapVisitor(PhantomData)) + } +} + +#[cfg(test)] +mod tests { + use std::collections::BTreeMap; + + use super::{ + MAX_WORKFLOW_VERSION_BYTES, MAX_WORKFLOW_VERSION_FILE_BYTES, MAX_WORKFLOW_VERSION_FILES, + WorkflowVersion, WorkflowVersionShapeError, + }; + use crate::WorkflowPath; + + fn path(value: &str) -> WorkflowPath { + value.parse().unwrap() + } + + #[test] + fn canonical_bytes_have_fixed_field_and_map_order() { + let version = WorkflowVersion::new( + path("workflow.fabro"), + BTreeMap::from([ + (path("z.txt"), "Z".to_string()), + (path("workflow.fabro"), "digraph W {}".to_string()), + (path("a.txt"), "A".to_string()), + ]), + BTreeMap::new(), + ) + .unwrap(); + + assert_eq!( + String::from_utf8(version.canonical_bytes().unwrap()).unwrap(), + r#"{"entrypoint":"workflow.fabro","files":{"a.txt":"A","workflow.fabro":"digraph W {}","z.txt":"Z"},"dependencies":{}}"# + ); + } + + #[test] + fn rejects_missing_entrypoint() { + let error = WorkflowVersion::new( + path("missing.fabro"), + BTreeMap::from([(path("workflow.fabro"), "digraph W {}".to_string())]), + BTreeMap::new(), + ) + .unwrap_err(); + assert!(matches!( + error, + WorkflowVersionShapeError::MissingEntrypoint { .. } + )); + } + + #[test] + fn rejects_path_collisions_and_large_files() { + let collision = WorkflowVersion::new( + path("workflow.fabro"), + BTreeMap::from([ + (path("workflow.fabro"), "digraph W {}".to_string()), + (path("assets"), "file".to_string()), + (path("assets/item.txt"), "nested".to_string()), + ]), + BTreeMap::new(), + ) + .unwrap_err(); + assert!(matches!( + collision, + WorkflowVersionShapeError::PathCollision { .. } + )); + + let mut files = BTreeMap::from([(path("workflow.fabro"), "digraph W {}".to_string())]); + files.insert( + path("large.txt"), + "x".repeat(MAX_WORKFLOW_VERSION_FILE_BYTES + 1), + ); + let large = + WorkflowVersion::new(path("workflow.fabro"), files, BTreeMap::new()).unwrap_err(); + assert!(matches!( + large, + WorkflowVersionShapeError::FileTooLarge { .. } + )); + } + + #[test] + fn enforces_file_count_file_size_and_canonical_size_boundaries() { + let mut files = BTreeMap::from([(path("workflow.fabro"), "digraph W {}".to_string())]); + for index in 0..MAX_WORKFLOW_VERSION_FILES - 1 { + files.insert(path(&format!("file-{index:03}.txt")), String::new()); + } + assert!( + WorkflowVersion::new(path("workflow.fabro"), files.clone(), BTreeMap::new()).is_ok() + ); + files.insert(path("too-many.txt"), String::new()); + assert!(matches!( + WorkflowVersion::new(path("workflow.fabro"), files, BTreeMap::new()).unwrap_err(), + WorkflowVersionShapeError::TooManyFiles { .. } + )); + + let exact_file = BTreeMap::from([ + (path("workflow.fabro"), "digraph W {}".to_string()), + ( + path("payload.txt"), + "x".repeat(MAX_WORKFLOW_VERSION_FILE_BYTES), + ), + ]); + assert!( + WorkflowVersion::new(path("workflow.fabro"), exact_file.clone(), BTreeMap::new()) + .is_ok() + ); + let mut oversized_file = exact_file; + oversized_file + .get_mut(&path("payload.txt")) + .unwrap() + .push('x'); + assert!(matches!( + WorkflowVersion::new(path("workflow.fabro"), oversized_file, BTreeMap::new()) + .unwrap_err(), + WorkflowVersionShapeError::FileTooLarge { .. } + )); + + let mut exact_version_files = + BTreeMap::from([(path("workflow.fabro"), "digraph W {}".to_string())]); + for index in 0..4 { + exact_version_files.insert(path(&format!("payload-{index}.txt")), String::new()); + } + let empty = WorkflowVersion::new( + path("workflow.fabro"), + exact_version_files.clone(), + BTreeMap::new(), + ) + .unwrap(); + let remaining = MAX_WORKFLOW_VERSION_BYTES - empty.canonical_bytes().unwrap().len(); + let per_file = remaining / 4; + let remainder = remaining % 4; + for index in 0..4 { + let length = per_file + usize::from(index < remainder); + assert!(length <= MAX_WORKFLOW_VERSION_FILE_BYTES); + exact_version_files.insert(path(&format!("payload-{index}.txt")), "x".repeat(length)); + } + let exact_version = WorkflowVersion::new( + path("workflow.fabro"), + exact_version_files.clone(), + BTreeMap::new(), + ) + .unwrap(); + assert_eq!( + exact_version.canonical_bytes().unwrap().len(), + MAX_WORKFLOW_VERSION_BYTES + ); + exact_version_files + .get_mut(&path("payload-0.txt")) + .unwrap() + .push('x'); + assert!(matches!( + WorkflowVersion::new(path("workflow.fabro"), exact_version_files, BTreeMap::new()) + .unwrap_err(), + WorkflowVersionShapeError::VersionTooLarge { .. } + )); + } + + #[test] + fn deserialize_rejects_unknown_fields_and_duplicate_keys() { + let unknown = r#"{ + "entrypoint":"workflow.fabro", + "files":{"workflow.fabro":"digraph W {}"}, + "dependencies":{}, + "metadata":{} + }"#; + assert!(serde_json::from_str::(unknown).is_err()); + + let duplicate = r#"{ + "entrypoint":"workflow.fabro", + "files":{"workflow.fabro":"digraph W {}","workflow.fabro":"digraph X {}"}, + "dependencies":{} + }"#; + assert!(serde_json::from_str::(duplicate).is_err()); + } +}