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()); + } +}