From 178320e7a56d270fefecc01d25eea414832cf26e Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Thu, 13 Aug 2026 14:53:32 -0400 Subject: [PATCH 1/8] Add immutable workflow version resource Add the WorkflowVersion domain resource with exactly entrypoint, files, and workflow_dependencies, plus strict WorkflowPath validation and deterministic canonical raw JSON. Semantic validation of graph imports, templates, file references, workflow.toml rules, Dockerfile paths, and exact child-workflow dependency bindings lives in the new fabro-workflow-version crate, which validates the complete stored dependency closure through the shared blob store before writing a root. The authenticated create-only POST /api/v1/workflow-versions endpoint ships with its OpenAPI contract, Rust type replacements, and generated TypeScript client. Squashed from the resource commits of the original combined branch; the walker unification this builds on landed separately. Co-Authored-By: Claude Fable 5 --- Cargo.lock | 17 + docs/public/api-reference/fabro-api.yaml | 130 +++++ lib/apps/fabro-server/Cargo.toml | 1 + .../fabro-server/src/server/handler/mod.rs | 2 + .../src/server/handler/workflow_versions.rs | 357 ++++++++++++ .../fabro-workflow-version/Cargo.toml | 27 + .../fabro-workflow-version/src/lib.rs | 507 ++++++++++++++++++ .../fabro-workflow-version/src/store.rs | 306 +++++++++++ lib/foundation/fabro-api/build.rs | 3 + lib/foundation/fabro-api/src/lib.rs | 3 +- .../tests/workflow_version_round_trip.rs | 54 ++ lib/foundation/fabro-types/src/lib.rs | 11 + .../fabro-types/src/workflow_path.rs | 301 +++++++++++ .../fabro-types/src/workflow_version.rs | 419 +++++++++++++++ .../fabro-types/src/workflow_version_id.rs | 96 ++++ .../src/.openapi-generator/FILES | 3 + lib/packages/fabro-api-client/src/api.ts | 1 + .../src/api/workflow-versions-api.ts | 134 +++++ .../create-workflow-version-response.ts | 25 + .../fabro-api-client/src/models/index.ts | 2 + .../src/models/workflow-version.ts | 33 ++ 21 files changed, 2431 insertions(+), 1 deletion(-) create mode 100644 lib/apps/fabro-server/src/server/handler/workflow_versions.rs create mode 100644 lib/components/fabro-workflow-version/Cargo.toml create mode 100644 lib/components/fabro-workflow-version/src/lib.rs create mode 100644 lib/components/fabro-workflow-version/src/store.rs create mode 100644 lib/foundation/fabro-api/tests/workflow_version_round_trip.rs create mode 100644 lib/foundation/fabro-types/src/workflow_path.rs create mode 100644 lib/foundation/fabro-types/src/workflow_version.rs create mode 100644 lib/foundation/fabro-types/src/workflow_version_id.rs create mode 100644 lib/packages/fabro-api-client/src/api/workflow-versions-api.ts create mode 100644 lib/packages/fabro-api-client/src/models/create-workflow-version-response.ts create mode 100644 lib/packages/fabro-api-client/src/models/workflow-version.ts diff --git a/Cargo.lock b/Cargo.lock index 661b4ab8f..93c745736 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3060,6 +3060,7 @@ dependencies = [ "fabro-variable", "fabro-vault", "fabro-workflow", + "fabro-workflow-version", "futures-util", "globset", "hex", @@ -3430,6 +3431,22 @@ dependencies = [ "walkdir", ] +[[package]] +name = "fabro-workflow-version" +version = "0.324.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]] name = "fail-parallel" version = "0.5.1" diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 261badbbb..274282e72 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -33,6 +33,8 @@ tags: description: Internal run details (stages, turns, context, configuration) - name: Workflows description: Workflow definitions and execution + - name: Workflow Versions + description: Immutable, content-addressed workflow packages - name: Billing description: Token counts and billed totals - name: Insights @@ -1065,6 +1067,69 @@ paths: schema: $ref: "#/components/schemas/ErrorResponse" + # ── Workflow Versions ───────────────────────────────────────────────── + + /api/v1/workflow-versions: + post: + operationId: createWorkflowVersion + tags: [Workflow Versions] + summary: Create Workflow Version + description: >- + Validates and stores an immutable workflow package in content-addressed + storage. Repeating the same canonical content returns the same identifier. + requestBody: + required: true + content: + application/json: + schema: + $ref: "#/components/schemas/WorkflowVersion" + responses: + "201": + description: Workflow version stored or already present + content: + application/json: + schema: + $ref: "#/components/schemas/CreateWorkflowVersionResponse" + "400": + description: Malformed JSON (`invalid_json`) + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + "413": + description: Request body exceeds 2 MiB (`workflow_version_too_large`) + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + "422": + description: >- + Invalid workflow content (`workflow_version_invalid`) or an absent, + invalid, or non-canonical dependency + (`workflow_version_dependency_not_found`) + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + "500": + description: Workflow version storage failed + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + # ── Runs ────────────────────────────────────────────────────────────── /api/v1/runs: @@ -9074,6 +9139,71 @@ components: detail: $ref: "#/components/schemas/FailureDetail" + WorkflowPath: + description: >- + Canonical portable path inside one workflow version. Paths are UTF-8, + relative, at most 240 bytes and 16 components, and cannot contain empty, + dot, parent, backslash, control, tilde-root, or drive-letter segments. + Map keys receive stricter byte and structural validation in the domain + model than OpenAPI can express. + type: string + minLength: 1 + maxLength: 240 + example: graphs/main.fabro + + WorkflowVersionId: + description: SHA-256 identity of validated canonical workflow-version bytes. + type: string + pattern: "^[0-9a-f]{64}$" + example: "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef" + + WorkflowVersion: + description: >- + Complete immutable package for one rooted workflow. It contains at most + 512 files and 512 workflow dependencies, each file is at most 512 KiB + of UTF-8 content, and its compact canonical JSON representation is at + most 2 MiB. + type: object + additionalProperties: false + required: + - entrypoint + - files + - workflow_dependencies + properties: + entrypoint: + $ref: "#/components/schemas/WorkflowPath" + files: + type: object + description: >- + Workflow-local text files keyed by canonical path. Keys receive + stricter domain validation than OpenAPI can express; each value is + limited to 512 KiB of UTF-8 bytes. + maxProperties: 512 + propertyNames: + $ref: "#/components/schemas/WorkflowPath" + additionalProperties: + type: string + workflow_dependencies: + type: object + description: >- + Exact stored workflow-version IDs keyed by resolved child-workflow + path. Keys receive stricter domain validation than OpenAPI can express. + maxProperties: 512 + propertyNames: + $ref: "#/components/schemas/WorkflowPath" + additionalProperties: + $ref: "#/components/schemas/WorkflowVersionId" + + CreateWorkflowVersionResponse: + description: Identity of the stored immutable workflow version. + type: object + additionalProperties: false + required: + - workflow_version_id + properties: + workflow_version_id: + $ref: "#/components/schemas/WorkflowVersionId" + RunManifest: description: Self-contained workflow run manifest. type: object diff --git a/lib/apps/fabro-server/Cargo.toml b/lib/apps/fabro-server/Cargo.toml index 8aea2bade..5ad20cfbb 100644 --- a/lib/apps/fabro-server/Cargo.toml +++ b/lib/apps/fabro-server/Cargo.toml @@ -32,6 +32,7 @@ fabro-hooks = { path = "../../components/fabro-hooks" } fabro-interview = { path = "../../components/fabro-interview" } fabro-slack = { path = "../../components/fabro-slack" } fabro-workflow = { path = "../../components/fabro-workflow" } +fabro-workflow-version = { path = "../../components/fabro-workflow-version" } fabro-validate = { path = "../../components/fabro-validate" } fabro-sandbox = { path = "../../components/fabro-sandbox", features = ["daytona", "docker"] } fabro-github = { path = "../../components/fabro-github" } diff --git a/lib/apps/fabro-server/src/server/handler/mod.rs b/lib/apps/fabro-server/src/server/handler/mod.rs index c186464d7..bcb7f9ef1 100644 --- a/lib/apps/fabro-server/src/server/handler/mod.rs +++ b/lib/apps/fabro-server/src/server/handler/mod.rs @@ -30,6 +30,7 @@ mod steer; pub(in crate::server) mod system; mod variables; mod worker_control; +mod workflow_versions; pub(super) use system::{health, openapi_spec}; @@ -226,6 +227,7 @@ pub(super) fn real_routes() -> Router> { .merge(secrets::routes()) .merge(variables::routes()) .merge(worker_control::routes()) + .merge(workflow_versions::routes()) .merge(sessions::routes()) .merge(system::routes()) .merge(completions::routes()) diff --git a/lib/apps/fabro-server/src/server/handler/workflow_versions.rs b/lib/apps/fabro-server/src/server/handler/workflow_versions.rs new file mode 100644 index 000000000..ebff58d82 --- /dev/null +++ b/lib/apps/fabro-server/src/server/handler/workflow_versions.rs @@ -0,0 +1,357 @@ +use std::sync::Arc; + +use axum::extract::DefaultBodyLimit; +use axum::extract::rejection::JsonRejection; +use fabro_api::types::{CreateWorkflowVersionResponse, WorkflowVersion}; +use fabro_types::MAX_WORKFLOW_VERSION_BYTES; +use fabro_util::error; +use fabro_workflow_version::{ + ValidatedWorkflowVersion, WorkflowVersionStore, WorkflowVersionStoreError, +}; + +use super::super::{ + ApiError, AppState, IntoResponse, Json, RequiredUser, Response, Router, State, StatusCode, post, +}; + +const INVALID_JSON_CODE: &str = "invalid_json"; +const INVALID_VERSION_CODE: &str = "workflow_version_invalid"; +const DEPENDENCY_NOT_FOUND_CODE: &str = "workflow_version_dependency_not_found"; +const VERSION_TOO_LARGE_CODE: &str = "workflow_version_too_large"; + +pub(super) fn routes() -> Router> { + Router::new().route( + "/workflow-versions", + post(create_workflow_version).layer(DefaultBodyLimit::max(MAX_WORKFLOW_VERSION_BYTES)), + ) +} + +async fn create_workflow_version( + _auth: RequiredUser, + State(state): State>, + 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, + error_chain = ?error::collect_chain(&err), + "Failed to open workflow version storage" + ); + internal_store_error() + })?; + let store = WorkflowVersionStore::new(blobs); + let workflow_version_id = store.put(&version).await.map_err(store_error)?; + + Ok(( + StatusCode::CREATED, + Json(CreateWorkflowVersionResponse { + workflow_version_id, + }), + ) + .into_response()) +} + +fn json_rejection(rejection: JsonRejection) -> ApiError { + if rejection.status() == StatusCode::PAYLOAD_TOO_LARGE { + return ApiError::with_code( + StatusCode::PAYLOAD_TOO_LARGE, + "workflow version request exceeds 2 MiB", + VERSION_TOO_LARGE_CODE, + ); + } + + match rejection { + JsonRejection::JsonDataError(err) => ApiError::with_code( + StatusCode::UNPROCESSABLE_ENTITY, + err.body_text(), + INVALID_VERSION_CODE, + ), + other => ApiError::with_code( + StatusCode::BAD_REQUEST, + other.body_text(), + INVALID_JSON_CODE, + ), + } +} + +fn store_error(err: WorkflowVersionStoreError) -> ApiError { + match err { + err @ WorkflowVersionStoreError::DependencyNotFound { .. } => ApiError::with_code( + StatusCode::UNPROCESSABLE_ENTITY, + err.to_string(), + DEPENDENCY_NOT_FOUND_CODE, + ), + WorkflowVersionStoreError::InvalidVersion(source) => ApiError::with_code( + StatusCode::UNPROCESSABLE_ENTITY, + 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, + error_chain = ?error::collect_chain(&err), + "Workflow version store operation failed" + ); + internal_store_error() + } + } +} + +fn internal_store_error() -> ApiError { + ApiError::new( + StatusCode::INTERNAL_SERVER_ERROR, + "workflow version store operation failed", + ) +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use axum::body::{Body, to_bytes}; + use axum::http::{Method, Request, StatusCode, header}; + use axum::response::IntoResponse; + use fabro_types::WorkflowVersionId; + use serde_json::{Value, json}; + use tower::ServiceExt; + + use super::{ + DEPENDENCY_NOT_FOUND_CODE, INVALID_JSON_CODE, INVALID_VERSION_CODE, + MAX_WORKFLOW_VERSION_BYTES, VERSION_TOO_LARGE_CODE, store_error, + }; + use crate::server; + use crate::test_support::{self, TestAppStateBuilder}; + + const GRAPH: &str = "digraph W { start [shape=Mdiamond] exit [shape=Msquare] start -> exit }"; + + fn request(body: impl Into) -> Request { + Request::builder() + .method(Method::POST) + .uri("/api/v1/workflow-versions") + .header(header::CONTENT_TYPE, "application/json") + .body(body.into()) + .unwrap() + } + + fn version(graph: &str) -> Value { + json!({ + "entrypoint": "workflow.fabro", + "files": { "workflow.fabro": graph }, + "workflow_dependencies": {} + }) + } + + async fn response_json(response: axum::response::Response) -> Value { + let bytes = to_bytes(response.into_body(), usize::MAX).await.unwrap(); + serde_json::from_slice(&bytes).unwrap() + } + + fn error_code(body: &Value) -> &str { + body["errors"][0]["code"].as_str().unwrap() + } + + #[tokio::test] + async fn create_requires_authenticated_user() { + let state = TestAppStateBuilder::new().build(); + let app = server::build_router(state, test_support::test_auth_mode()); + let response = app + .oneshot(request(serde_json::to_vec(&version(GRAPH)).unwrap())) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::UNAUTHORIZED); + } + + #[tokio::test] + async fn valid_and_equivalent_requests_return_the_same_id() { + let state = TestAppStateBuilder::new().build(); + let app = test_support::build_test_router(Arc::clone(&state)); + let first = app + .clone() + .oneshot(request(serde_json::to_vec(&version(GRAPH)).unwrap())) + .await + .unwrap(); + assert_eq!(first.status(), StatusCode::CREATED); + let first = response_json(first).await; + assert_eq!(first.as_object().unwrap().len(), 1); + + let reordered = format!( + r#"{{"workflow_dependencies":{{}},"files":{{"workflow.fabro":{}}},"entrypoint":"workflow.fabro"}}"#, + serde_json::to_string(GRAPH).unwrap() + ); + let second = app.oneshot(request(reordered)).await.unwrap(); + assert_eq!(second.status(), StatusCode::CREATED); + assert_eq!(response_json(second).await, first); + + let id = first["workflow_version_id"] + .as_str() + .unwrap() + .parse::() + .unwrap(); + assert!( + state + .store_ref() + .blobs() + .await + .unwrap() + .read(&id.into()) + .await + .unwrap() + .is_some() + ); + } + + #[tokio::test] + async fn invalid_json_and_domain_content_have_distinct_codes() { + let app = test_support::build_test_router(TestAppStateBuilder::new().build()); + let malformed = app.clone().oneshot(request("{")).await.unwrap(); + assert_eq!(malformed.status(), StatusCode::BAD_REQUEST); + assert_eq!( + error_code(&response_json(malformed).await), + INVALID_JSON_CODE + ); + + let unknown = json!({ + "entrypoint": "workflow.fabro", + "files": { "workflow.fabro": GRAPH }, + "workflow_dependencies": {}, + "metadata": {} + }); + let invalid = app + .oneshot(request(serde_json::to_vec(&unknown).unwrap())) + .await + .unwrap(); + assert_eq!(invalid.status(), StatusCode::UNPROCESSABLE_ENTITY); + assert_eq!( + error_code(&response_json(invalid).await), + INVALID_VERSION_CODE + ); + } + + #[tokio::test] + async fn unavailable_dependency_has_specific_code() { + let state = TestAppStateBuilder::new().build(); + let app = test_support::build_test_router(Arc::clone(&state)); + let missing_id = WorkflowVersionId::from(fabro_types::BlobHash::new(b"missing")); + let root = json!({ + "entrypoint": "workflow.fabro", + "files": { + "workflow.fabro": "digraph W { child [stack.child_workflow=\"child.fabro\"] }" + }, + "workflow_dependencies": { "child.fabro": missing_id } + }); + let response = app + .oneshot(request(serde_json::to_vec(&root).unwrap())) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::UNPROCESSABLE_ENTITY); + assert_eq!( + error_code(&response_json(response).await), + DEPENDENCY_NOT_FOUND_CODE + ); + } + + #[tokio::test] + async fn corrupt_stored_dependency_returns_curated_internal_error() { + let state = TestAppStateBuilder::new().build(); + let app = test_support::build_test_router(Arc::clone(&state)); + let dependency_id = WorkflowVersionId::from( + state + .store_ref() + .blobs() + .await + .unwrap() + .write(b"not a workflow version") + .await + .unwrap(), + ); + let root = json!({ + "entrypoint": "workflow.fabro", + "files": { + "workflow.fabro": "digraph W { child [stack.child_workflow=\"child.fabro\"] }" + }, + "workflow_dependencies": { "child.fabro": dependency_id } + }); + + let response = app + .oneshot(request(serde_json::to_vec(&root).unwrap())) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::INTERNAL_SERVER_ERROR); + let body = response_json(response).await; + assert_eq!( + body["errors"][0]["detail"], + "workflow version store operation failed" + ); + assert!(!body.to_string().contains("cannot be decoded")); + } + + #[tokio::test] + async fn stored_child_can_be_pinned_as_a_dependency() { + let app = test_support::build_test_router(TestAppStateBuilder::new().build()); + let child = app + .clone() + .oneshot(request(serde_json::to_vec(&version(GRAPH)).unwrap())) + .await + .unwrap(); + assert_eq!(child.status(), StatusCode::CREATED); + let child_id = response_json(child).await["workflow_version_id"].clone(); + let root = json!({ + "entrypoint": "workflow.fabro", + "files": { + "workflow.fabro": "digraph W { child [stack.child_workflow=\"child.fabro\"] }" + }, + "workflow_dependencies": { "child.fabro": child_id } + }); + + let response = app + .oneshot(request(serde_json::to_vec(&root).unwrap())) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::CREATED); + assert_eq!(response_json(response).await.as_object().unwrap().len(), 1); + } + + #[tokio::test] + async fn body_limit_has_specific_code() { + let app = test_support::build_test_router(TestAppStateBuilder::new().build()); + let response = app + .oneshot(request(vec![b' '; MAX_WORKFLOW_VERSION_BYTES + 1])) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::PAYLOAD_TOO_LARGE); + assert_eq!( + error_code(&response_json(response).await), + VERSION_TOO_LARGE_CODE + ); + } + + #[tokio::test] + async fn storage_fault_response_is_curated() { + let response = store_error(fabro_workflow_version::WorkflowVersionStoreError::Storage { + source: fabro_store::Error::Other("private persistence detail".to_string()), + }) + .into_response(); + assert_eq!(response.status(), StatusCode::INTERNAL_SERVER_ERROR); + + let body = response_json(response).await; + assert_eq!( + body["errors"][0]["detail"], + "workflow version store operation failed" + ); + assert!(!body.to_string().contains("private persistence detail")); + } +} diff --git a/lib/components/fabro-workflow-version/Cargo.toml b/lib/components/fabro-workflow-version/Cargo.toml new file mode 100644 index 000000000..bc91e3f89 --- /dev/null +++ b/lib/components/fabro-workflow-version/Cargo.toml @@ -0,0 +1,27 @@ +[package] +name = "fabro-workflow-version" +edition.workspace = true +version.workspace = true +publish = false +license.workspace = true +description = "Semantic validation and storage for immutable workflow versions" + +[lib] +doctest = false + +[lints] +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 new file mode 100644 index 000000000..320d09139 --- /dev/null +++ b/lib/components/fabro-workflow-version/src/lib.rs @@ -0,0 +1,507 @@ +//! 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::parse::{SettingsSource, validate_settings_source}; +use fabro_config::{EnvironmentDockerfileLayer, EnvironmentImageLayer, SettingsLayer}; +use fabro_graphviz::parser; +use fabro_template::{ + BundleTemplateStore, GraphReference, GraphReferenceError, StaticReferenceError, + TemplateDiscoveryError, TemplateSource, discover_static_dependency_closure, + validate_static_reference, visit_graph_references, +}; +use fabro_types::graph::ReferenceKind; +use fabro_types::{ManifestPath, WorkflowPath, WorkflowPathParseError, WorkflowVersion}; +use thiserror::Error; + +mod store; + +pub use store::{WorkflowVersionStore, WorkflowVersionStoreError}; + +#[derive(Debug, Error)] +pub enum WorkflowVersionError { + #[error("workflow graph `{path}` is invalid")] + GraphParse { + path: WorkflowPath, + #[source] + source: fabro_graphviz::Error, + }, + #[error("invalid {kind} in `{path}`: `{reference}`")] + InvalidReference { + path: WorkflowPath, + kind: ReferenceKind, + reference: String, + #[source] + source: WorkflowPathParseError, + }, + #[error("invalid static reference in `{path}`")] + StaticReference { + path: WorkflowPath, + #[source] + source: StaticReferenceError, + }, + #[error("{kind} in `{path}` references missing file `{target}`")] + MissingFile { + path: WorkflowPath, + kind: ReferenceKind, + target: WorkflowPath, + }, + #[error("template dependencies for `{path}` are invalid")] + Template { + path: WorkflowPath, + #[source] + source: Box, + }, + #[error("workflow.toml is invalid")] + Config { + #[source] + source: fabro_config::ParseError, + }, + #[error( + "workflow.toml selects graph `{configured}`, but the version entrypoint is `{entrypoint}`" + )] + ConfigEntrypointMismatch { + configured: WorkflowPath, + entrypoint: WorkflowPath, + }, + #[error("workflow dependencies do not match child workflow references")] + DependencyMismatch { + missing: Vec, + unused: Vec, + }, +} + +/// 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 + } +} + +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 })?; + validate_settings_source(&layer, SettingsSource::Workflow) + .map_err(|source| WorkflowVersionError::Config { source })?; + + 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(), + }); + } + } + + for image in layer.environment_images() { + validate_dockerfile(version, &config_path, image)?; + } + Ok(()) +} + +fn validate_dockerfile( + version: &WorkflowVersion, + config_path: &WorkflowPath, + image: &EnvironmentImageLayer, +) -> Result<(), WorkflowVersionError> { + let Some(EnvironmentDockerfileLayer::Path { path }) = image.dockerfile.as_ref() else { + return Ok(()); + }; + validate_static_reference(path, ReferenceKind::Dockerfile).map_err(|source| { + WorkflowVersionError::StaticReference { + path: config_path.clone(), + source, + } + })?; + 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; + } + 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, + })?; + + visit_graph_references(&graph, |reference| match reference { + GraphReference::GoalFile { reference } => { + let target = resolve_reference(&path, ReferenceKind::GraphGoalFile, reference)?; + let content = + require_file(version, &path, ReferenceKind::GraphGoalFile, target.clone())?; + validate_template(&target, content, &template_store, &template_root) + } + GraphReference::GoalInline { content } | GraphReference::InlinePrompt { content } => { + validate_template(&path, content, &template_store, &template_root) + } + GraphReference::Import { reference } => { + let target = resolve_reference(&path, ReferenceKind::Import, reference)?; + require_file(version, &path, ReferenceKind::Import, target.clone())?; + queue.push_back(target); + Ok(()) + } + GraphReference::ChildWorkflow { reference } => { + let target = resolve_reference(&path, ReferenceKind::ChildWorkflow, reference)?; + child_workflows.insert(target); + Ok(()) + } + GraphReference::FileInline { key, reference } => { + let target = resolve_reference(&path, ReferenceKind::FileInline, reference)?; + let content = + require_file(version, &path, ReferenceKind::FileInline, target.clone())?; + if key == "prompt" { + validate_template(&target, content, &template_store, &template_root)?; + } + Ok(()) + } + }) + .map_err(|error| match error { + GraphReferenceError::StaticReference(source) => WorkflowVersionError::StaticReference { + path: path.clone(), + source, + }, + GraphReferenceError::Visit(error) => error, + })?; + } + + let configured = version + .workflow_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_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 template_store(version: &WorkflowVersion) -> BundleTemplateStore { + BundleTemplateStore::new( + version + .files() + .iter() + .map(|(path, content)| (manifest_path(path), content.clone())) + .collect::>(), + ) +} + +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 { + ManifestPath::from_wire(path.as_str()) + .expect("validated workflow paths must also be valid manifest paths") +} + +#[cfg(test)] +mod tests { + use fabro_types::{BlobHash, WorkflowPath, WorkflowVersion, WorkflowVersionId}; + + use super::{ValidatedWorkflowVersion, WorkflowVersionError}; + + fn path(value: &str) -> WorkflowPath { + value.parse().unwrap() + } + + fn dependency_id(value: &[u8]) -> WorkflowVersionId { + BlobHash::new(value).into() + } + + fn version_with( + files: impl IntoIterator, + dependencies: impl IntoIterator, + ) -> 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 validates_imports_templates_file_refs_and_dependencies() { + let version = version_with( + [ + ( + "workflow.fabro", + r#"digraph W { + graph [goal="@prompts/goal.md"] + start [shape=Mdiamond] + imported [import="graphs/imported.fabro"] + child [stack.child_workflow="children/check.fabro"] + exit [shape=Msquare] + start -> imported -> child -> exit + }"#, + ), + ( + "graphs/imported.fabro", + r#"digraph I { step [prompt="{% include \"../prompts/partial.md\" %}"] }"#, + ), + ("prompts/goal.md", "{% include \"partial.md\" %}"), + ("prompts/partial.md", "Do the work"), + ], + [("children/check.fabro", dependency_id(b"child"))], + ) + .unwrap(); + + assert_eq!(version.version().workflow_dependencies().len(), 1); + } + + #[test] + fn rejects_missing_and_unused_dependencies() { + let error = version_with( + [( + "workflow.fabro", + r#"digraph W { child [stack.child_workflow="child.fabro"] }"#, + )], + [("unused.fabro", dependency_id(b"unused"))], + ) + .unwrap_err(); + + let WorkflowVersionError::DependencyMismatch { missing, unused } = error else { + panic!("expected dependency mismatch"); + }; + assert_eq!(missing, vec![path("child.fabro")]); + assert_eq!(unused, vec![path("unused.fabro")]); + } + + #[test] + fn rejects_config_entrypoint_and_missing_dockerfile() { + let error = version_with( + [ + ( + "workflow.fabro", + "digraph W { start [shape=Mdiamond] exit [shape=Msquare] start -> exit }", + ), + ( + "workflow.toml", + "_version = 1\n[workflow]\ngraph = \"other.fabro\"\n", + ), + ], + [], + ) + .unwrap_err(); + assert!(matches!( + error, + WorkflowVersionError::ConfigEntrypointMismatch { .. } + )); + + let missing_dockerfile = version_with( + [ + ("workflow.fabro", "digraph W {}"), + ( + "workflow.toml", + "_version = 1\n[run.environment.image]\ndockerfile = { path = \"docker/Dockerfile\" }\n", + ), + ], + [], + ) + .unwrap_err(); + assert!(matches!( + missing_dockerfile, + WorkflowVersionError::MissingFile { .. } + )); + + let invalid_config = version_with( + [ + ("workflow.fabro", "digraph W {}"), + ("workflow.toml", "not valid toml = ["), + ], + [], + ) + .unwrap_err(); + assert!(matches!( + invalid_config, + WorkflowVersionError::Config { .. } + )); + } + + #[test] + fn accepts_root_config_and_all_dockerfile_path_sources() { + let version = version_with( + [ + ("workflow.fabro", "digraph W {}"), + ( + "workflow.toml", + r#"_version = 1 +[workflow] +graph = "workflow.fabro" + +[environments.cloud] +provider = "daytona" + +[environments.cloud.image] +dockerfile = { path = "docker/named.Dockerfile" } + +[run.environment.image] +dockerfile = { path = "docker/run.Dockerfile" } +"#, + ), + ("docker/named.Dockerfile", "FROM alpine\n"), + ("docker/run.Dockerfile", "FROM ubuntu\n"), + ], + [], + ) + .unwrap(); + + assert_eq!(version.version().entrypoint(), &path("workflow.fabro")); + } + + #[test] + fn rejects_server_managed_environment_cwd_in_workflow_config() { + let error = version_with( + [ + ("workflow.fabro", "digraph W {}"), + ( + "workflow.toml", + "_version = 1\n[environments.local]\nprovider = \"local\"\ncwd = \"/tmp\"\n", + ), + ], + [], + ) + .unwrap_err(); + + assert!(matches!(error, WorkflowVersionError::Config { .. })); + assert!(error.to_string().contains("workflow.toml is invalid")); + } + + #[test] + fn rejects_escaping_and_dynamic_template_references() { + let escaping = version_with( + [( + "workflow.fabro", + r#"digraph W { imported [import="../outside.fabro"] }"#, + )], + [], + ) + .unwrap_err(); + assert!(matches!( + escaping, + WorkflowVersionError::InvalidReference { .. } + )); + + let dynamic = version_with( + [( + "workflow.fabro", + r#"digraph W { step [prompt="{% include template_name %}"] }"#, + )], + [], + ) + .unwrap_err(); + assert!(matches!(dynamic, WorkflowVersionError::Template { .. })); + } +} diff --git a/lib/components/fabro-workflow-version/src/store.rs b/lib/components/fabro-workflow-version/src/store.rs new file mode 100644 index 000000000..80dcebd53 --- /dev/null +++ b/lib/components/fabro-workflow-version/src/store.rs @@ -0,0 +1,306 @@ +use std::collections::{BTreeMap, HashSet, VecDeque}; +use std::sync::Arc; + +use fabro_store::BlobStore; +use fabro_types::{WorkflowPath, WorkflowVersion, WorkflowVersionId, WorkflowVersionShapeError}; +use thiserror::Error; + +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, + id: WorkflowVersionId, + }, + #[error("workflow-version dependency `{id}` at `{path}` is invalid")] + DependencyInvalid { + path: WorkflowPath, + id: WorkflowVersionId, + #[source] + source: Box, + }, + #[error("workflow-version blob `{id}` cannot be decoded as a valid workflow version")] + Decode { + id: WorkflowVersionId, + #[source] + source: serde_json::Error, + }, + #[error("workflow-version blob `{id}` is not canonical")] + NonCanonical { id: WorkflowVersionId }, + #[error("workflow-version storage operation failed")] + Storage { + #[source] + 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, +} + +impl WorkflowVersionStore { + #[must_use] + pub fn new(blobs: Arc) -> Self { + Self { blobs } + } + + pub async fn put( + &self, + version: &ValidatedWorkflowVersion, + ) -> Result { + let canonical = version.version().canonical_bytes()?; + self.validate_dependency_closure(version.version().workflow_dependencies()) + .await?; + self.blobs + .write(&canonical) + .await + .map(WorkflowVersionId::from) + .map_err(|source| WorkflowVersionStoreError::Storage { source }) + } + + pub async fn get( + &self, + id: &WorkflowVersionId, + ) -> Result, WorkflowVersionStoreError> { + let Some(version) = self.load_one(id).await? else { + return Ok(None); + }; + self.validate_dependency_closure(version.version().workflow_dependencies()) + .await?; + Ok(Some(version)) + } + + async fn load_one( + &self, + id: &WorkflowVersionId, + ) -> Result, WorkflowVersionStoreError> { + let blob_id = (*id).into(); + let Some(bytes) = self + .blobs + .read(&blob_id) + .await + .map_err(|source| WorkflowVersionStoreError::Storage { source })? + else { + return Ok(None); + }; + let version = serde_json::from_slice::(&bytes) + .map_err(|source| WorkflowVersionStoreError::Decode { id: *id, source })?; + 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(validated)) + } + + async fn validate_dependency_closure( + &self, + dependencies: &BTreeMap, + ) -> Result<(), WorkflowVersionStoreError> { + let mut pending = dependencies + .iter() + .map(|(path, id)| (path.clone(), *id)) + .collect::>(); + let mut visited = HashSet::new(); + + while let Some((path, id)) = pending.pop_front() { + if !visited.insert(id) { + continue; + } + match self.load_one(&id).await { + Ok(Some(dependency)) => { + pending.extend( + dependency + .version() + .workflow_dependencies() + .iter() + .map(|(path, id)| (path.clone(), *id)), + ); + } + Ok(None) => { + return Err(WorkflowVersionStoreError::DependencyNotFound { path, id }); + } + // Persistence failures are server faults, not evidence that + // the caller supplied an invalid dependency. + Err(source @ WorkflowVersionStoreError::Storage { .. }) => return Err(source), + Err(source) => { + return Err(WorkflowVersionStoreError::DependencyInvalid { + path, + id, + source: Box::new(source), + }); + } + } + } + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use std::collections::BTreeMap; + use std::sync::Arc; + use std::time::Duration; + + use fabro_store::{BlobStore, Database}; + use fabro_types::{WorkflowPath, WorkflowVersion, WorkflowVersionId}; + use object_store::memory::InMemory; + + use super::{WorkflowVersionStore, WorkflowVersionStoreError}; + use crate::ValidatedWorkflowVersion; + + fn path(value: &str) -> WorkflowPath { + value.parse().unwrap() + } + + fn version( + graph: &str, + dependencies: BTreeMap, + ) -> ValidatedWorkflowVersion { + ValidatedWorkflowVersion::new( + WorkflowVersion::new( + path("workflow.fabro"), + BTreeMap::from([(path("workflow.fabro"), graph.to_owned())]), + dependencies, + ) + .unwrap(), + ) + .unwrap() + } + + async fn stores() -> (Arc, WorkflowVersionStore) { + let database = Database::new( + Arc::new(InMemory::new()), + "", + Duration::from_millis(1), + None, + ); + let blobs = database.blobs().await.unwrap(); + let versions = WorkflowVersionStore::new(Arc::clone(&blobs)); + (blobs, versions) + } + + #[tokio::test] + async fn put_get_reuses_exact_blob_digest() { + let (blobs, store) = stores().await; + let version = version("digraph W {}", BTreeMap::new()); + let expected_bytes = version.version().canonical_bytes().unwrap(); + let expected_id = WorkflowVersionId::from(fabro_types::BlobHash::new(&expected_bytes)); + + let id = store.put(&version).await.unwrap(); + assert_eq!(id, expected_id); + let blob_id = id.into(); + assert_eq!(blobs.read(&blob_id).await.unwrap().unwrap(), expected_bytes); + assert_eq!(store.get(&id).await.unwrap(), Some(version)); + } + + #[tokio::test] + async fn identical_content_is_idempotent() { + let (_, store) = stores().await; + let original = version("digraph W {}", BTreeMap::new()); + + assert_eq!( + store.put(&original).await.unwrap(), + store.put(&original).await.unwrap() + ); + + let changed = version("digraph W { changed [label=\"yes\"] }", BTreeMap::new()); + assert_ne!( + store.put(&original).await.unwrap(), + store.put(&changed).await.unwrap() + ); + } + + #[tokio::test] + async fn dependency_must_be_stored_first() { + let (blobs, store) = stores().await; + let child = version("digraph Child {}", BTreeMap::new()); + let child_id = WorkflowVersionId::from(fabro_types::BlobHash::new( + &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::BlobHash::new( + &root.version().canonical_bytes().unwrap(), + )); + + let error = store.put(&root).await.unwrap_err(); + assert!(matches!( + error, + WorkflowVersionStoreError::DependencyNotFound { .. } + )); + assert!(!blobs.exists(&root_id.into()).await.unwrap()); + assert_eq!(store.put(&child).await.unwrap(), child_id); + assert!(store.put(&root).await.is_ok()); + } + + #[tokio::test] + async fn dependency_closure_must_be_complete_before_root_write() { + let (blobs, store) = stores().await; + let missing_grandchild_id = WorkflowVersionId::from(fabro_types::BlobHash::new(b"missing")); + let child = version( + r#"digraph Child { grandchild [stack.child_workflow="grandchild.fabro"] }"#, + BTreeMap::from([(path("grandchild.fabro"), missing_grandchild_id)]), + ); + let child_bytes = child.version().canonical_bytes().unwrap(); + let child_id = WorkflowVersionId::from(blobs.write(&child_bytes).await.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::BlobHash::new( + &root.version().canonical_bytes().unwrap(), + )); + + assert!(matches!( + store.put(&root).await.unwrap_err(), + WorkflowVersionStoreError::DependencyNotFound { id, .. } + if id == missing_grandchild_id + )); + assert!(!blobs.exists(&root_id.into()).await.unwrap()); + assert!(matches!( + store.get(&child_id).await.unwrap_err(), + WorkflowVersionStoreError::DependencyNotFound { id, .. } + if id == missing_grandchild_id + )); + } + + #[tokio::test] + async fn get_rejects_arbitrary_and_noncanonical_blobs() { + let (blobs, store) = stores().await; + let arbitrary = WorkflowVersionId::from(blobs.write(b"not json").await.unwrap()); + assert!(matches!( + store.get(&arbitrary).await.unwrap_err(), + WorkflowVersionStoreError::Decode { .. } + )); + + let invalid_bytes = br#"{"entrypoint":"missing.fabro","files":{"workflow.fabro":"digraph W {}"},"workflow_dependencies":{}}"#; + let invalid = WorkflowVersionId::from(blobs.write(invalid_bytes).await.unwrap()); + assert!(matches!( + store.get(&invalid).await.unwrap_err(), + WorkflowVersionStoreError::Decode { .. } + )); + + let version = version("digraph W {}", BTreeMap::new()); + 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(), + WorkflowVersionStoreError::NonCanonical { .. } + )); + } +} diff --git a/lib/foundation/fabro-api/build.rs b/lib/foundation/fabro-api/build.rs index 38811dd01..ea8ece1eb 100644 --- a/lib/foundation/fabro-api/build.rs +++ b/lib/foundation/fabro-api/build.rs @@ -722,6 +722,9 @@ fn main() { ("CompletionMessage", "fabro_types::Message", &[]), ("CompletionMessageRole", "fabro_types::Role", &[]), ("CompletionContentPart", "fabro_types::ContentPart", &[]), + ("WorkflowVersion", "fabro_types::WorkflowVersion", &[]), + ("WorkflowPath", "fabro_types::WorkflowPath", &[]), + ("WorkflowVersionId", "fabro_types::WorkflowVersionId", &[]), ("CostSource", "fabro_model::CostSource", &[]), ]; for (name, path, impls) in replacements { diff --git a/lib/foundation/fabro-api/src/lib.rs b/lib/foundation/fabro-api/src/lib.rs index d9d6aa025..b40087831 100644 --- a/lib/foundation/fabro-api/src/lib.rs +++ b/lib/foundation/fabro-api/src/lib.rs @@ -73,7 +73,8 @@ pub mod types { StageModelUsage, StageOutcome, StageProjection, StageState, StageToolBatchProjection, SubAgentProjection, SubAgentStatus, SystemActorKind, SystemIntegrationStatus, SystemIntegrationsResponse, TodoListProjection, TurnId, UpdateVariableRequest, - UserPrincipal, Variable, VariableListResponse, WorkflowSettings, + UserPrincipal, Variable, VariableListResponse, WorkflowPath, WorkflowSettings, + WorkflowVersion, WorkflowVersionId, }; 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 new file mode 100644 index 000000000..894c1c910 --- /dev/null +++ b/lib/foundation/fabro-api/tests/workflow_version_round_trip.rs @@ -0,0 +1,54 @@ +use std::any::{TypeId, type_name}; + +use fabro_api::types::{ + WorkflowPath as ApiWorkflowPath, WorkflowVersion as ApiWorkflowVersion, + WorkflowVersionId as ApiWorkflowVersionId, +}; +use fabro_types::{WorkflowPath, WorkflowVersion, WorkflowVersionId}; +use serde_json::json; + +#[test] +fn workflow_version_schemas_reuse_domain_types() { + assert_same_type::(); + assert_same_type::(); + assert_same_type::(); +} + +#[test] +fn workflow_version_round_trips_exact_wire_shape() { + let value = json!({ + "entrypoint": "workflow.fabro", + "files": { + "prompts/goal.md": "Ship it", + "workflow.fabro": "digraph W { start [shape=Mdiamond] exit [shape=Msquare] start -> exit }" + }, + "workflow_dependencies": {} + }); + + let version: ApiWorkflowVersion = serde_json::from_value(value.clone()).unwrap(); + assert_eq!(serde_json::to_value(version).unwrap(), value); +} + +#[test] +fn workflow_version_replacement_rejects_unknown_fields() { + let value = json!({ + "entrypoint": "workflow.fabro", + "files": { + "workflow.fabro": "digraph W {}" + }, + "workflow_dependencies": {}, + "metadata": {} + }); + + assert!(serde_json::from_value::(value).is_err()); +} + +fn assert_same_type() { + assert_eq!( + TypeId::of::(), + TypeId::of::(), + "{} and {} should be the same type", + type_name::(), + type_name::() + ); +} diff --git a/lib/foundation/fabro-types/src/lib.rs b/lib/foundation/fabro-types/src/lib.rs index e2b05384c..bb3c86620 100644 --- a/lib/foundation/fabro-types/src/lib.rs +++ b/lib/foundation/fabro-types/src/lib.rs @@ -54,6 +54,9 @@ pub mod timing; 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; pub use auth::{IdpIdentity, IdpIdentityError}; @@ -183,3 +186,11 @@ pub use transcript::{ pub use variable::{ CreateVariableRequest, UpdateVariableRequest, Variable, VariableListResponse, is_env_style_name, }; +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_DEPENDENCIES, 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_path.rs b/lib/foundation/fabro-types/src/workflow_path.rs new file mode 100644 index 000000000..32297d93e --- /dev/null +++ b/lib/foundation/fabro-types/src/workflow_path.rs @@ -0,0 +1,301 @@ +use std::fmt; +use std::str::FromStr; + +use serde::{Deserialize, Serialize}; +use thiserror::Error; + +pub const MAX_WORKFLOW_PATH_BYTES: usize = 240; +pub const MAX_WORKFLOW_PATH_COMPONENTS: usize = 16; + +#[derive(Clone, Debug, PartialEq, Eq, Error)] +#[error("invalid workflow path `{value}`: {reason}")] +pub struct WorkflowPathParseError { + value: String, + reason: &'static str, +} + +impl WorkflowPathParseError { + fn new(value: &str, reason: &'static str) -> Self { + Self { + value: value.to_owned(), + reason, + } + } + + #[must_use] + pub fn value(&self) -> &str { + &self.value + } + + #[must_use] + pub fn reason(&self) -> &'static str { + self.reason + } +} + +#[derive(Clone, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] +#[serde(into = "String", try_from = "String")] +pub struct WorkflowPath(String); + +impl WorkflowPath { + pub fn new(value: impl Into) -> Result { + let value = value.into(); + validate(&value)?; + Ok(Self(value)) + } + + #[must_use] + pub fn as_str(&self) -> &str { + &self.0 + } + + #[must_use] + pub fn parent(&self) -> Option { + self.0 + .rsplit_once('/') + .map(|(parent, _)| Self(parent.to_owned())) + } + + #[must_use] + pub fn is_ancestor_of(&self, other: &Self) -> bool { + other.0.len() > self.0.len() + && other.0.starts_with(self.0.as_str()) + && other.0.as_bytes()[self.0.len()] == b'/' + } + + pub fn resolve_reference(&self, reference: &str) -> Result { + validate_reference_shape(reference)?; + let mut components = self + .0 + .rsplit_once('/') + .map_or_else(Vec::new, |(parent, _)| { + parent.split('/').collect::>() + }); + + for component in reference.split('/') { + match component { + "" | "." => {} + ".." => { + if components.pop().is_none() { + return Err(WorkflowPathParseError::new( + reference, + "reference escapes the workflow root", + )); + } + } + value => components.push(value), + } + } + + Self::new(components.join("/")) + } +} + +impl FromStr for WorkflowPath { + type Err = WorkflowPathParseError; + + fn from_str(value: &str) -> Result { + Self::new(value) + } +} + +impl TryFrom for WorkflowPath { + type Error = WorkflowPathParseError; + + fn try_from(value: String) -> Result { + Self::new(value) + } +} + +impl From for String { + fn from(value: WorkflowPath) -> Self { + value.0 + } +} + +impl fmt::Display for WorkflowPath { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str(&self.0) + } +} + +fn validate(value: &str) -> Result<(), WorkflowPathParseError> { + validate_reference_shape(value)?; + if value + .split('/') + .any(|component| matches!(component, "." | "..")) + { + return Err(WorkflowPathParseError::new( + value, + "dot segments are not allowed in stored paths", + )); + } + if value.split('/').count() > MAX_WORKFLOW_PATH_COMPONENTS { + return Err(WorkflowPathParseError::new( + value, + "path has too many components", + )); + } + if value.len() > MAX_WORKFLOW_PATH_BYTES { + return Err(WorkflowPathParseError::new(value, "path is too long")); + } + Ok(()) +} + +fn validate_reference_shape(value: &str) -> Result<(), WorkflowPathParseError> { + if value.is_empty() { + return Err(WorkflowPathParseError::new(value, "path is empty")); + } + if value.starts_with('/') { + return Err(WorkflowPathParseError::new( + value, + "absolute paths are not allowed", + )); + } + if value.starts_with('~') { + return Err(WorkflowPathParseError::new( + value, + "tilde-prefixed paths are not allowed", + )); + } + if value.contains('\\') { + return Err(WorkflowPathParseError::new( + value, + "backslashes are not allowed", + )); + } + if value.ends_with('/') { + return Err(WorkflowPathParseError::new( + value, + "trailing slashes are not allowed", + )); + } + if value.contains("//") { + return Err(WorkflowPathParseError::new( + value, + "repeated slashes are not allowed", + )); + } + let bytes = value.as_bytes(); + if bytes.len() >= 2 && bytes[0].is_ascii_alphabetic() && bytes[1] == b':' { + return Err(WorkflowPathParseError::new( + value, + "Windows drive paths are not allowed", + )); + } + if value.bytes().any(|byte| byte.is_ascii_control()) { + return Err(WorkflowPathParseError::new( + value, + "control characters are not allowed", + )); + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use std::collections::BTreeMap; + + use serde_json::json; + + use super::{MAX_WORKFLOW_PATH_BYTES, MAX_WORKFLOW_PATH_COMPONENTS, WorkflowPath}; + + #[test] + fn accepts_canonical_portable_paths() { + for value in ["workflow.fabro", "graphs/main.fabro", "prompts/日本語.md"] { + let path: WorkflowPath = value.parse().expect("path should parse"); + assert_eq!(path.as_str(), value); + } + } + + #[test] + fn rejects_non_canonical_or_unsafe_paths() { + for value in [ + "", "/root", "root/", "a//b", "a\\b", "~/a", "C:/a", ".", "..", "a/./b", "a/../b", + "a\nb", + ] { + assert!(value.parse::().is_err(), "accepted {value:?}"); + } + } + + #[test] + fn enforces_byte_and_component_limits() { + assert!( + "a".repeat(MAX_WORKFLOW_PATH_BYTES) + .parse::() + .is_ok() + ); + assert!( + "a".repeat(MAX_WORKFLOW_PATH_BYTES + 1) + .parse::() + .is_err() + ); + assert!( + vec!["a"; MAX_WORKFLOW_PATH_COMPONENTS] + .join("/") + .parse::() + .is_ok() + ); + assert!( + vec!["a"; MAX_WORKFLOW_PATH_COMPONENTS + 1] + .join("/") + .parse::() + .is_err() + ); + } + + #[test] + fn resolves_references_without_escaping_root() { + let graph: WorkflowPath = "graphs/nested/main.fabro".parse().unwrap(); + assert_eq!( + graph.resolve_reference("../prompts/plan.md").unwrap(), + "graphs/prompts/plan.md".parse().unwrap() + ); + assert!(graph.resolve_reference("../../../outside.md").is_err()); + assert!(graph.resolve_reference("prompts//plan.md").is_err()); + assert!(graph.resolve_reference("prompts/").is_err()); + } + + #[test] + fn ancestor_checks_component_boundaries() { + let parent: WorkflowPath = "dir/file".parse().unwrap(); + assert!(parent.is_ancestor_of(&"dir/file/child".parse().unwrap())); + assert!(!parent.is_ancestor_of(&"dir/filename".parse().unwrap())); + } + + #[test] + fn serde_and_ordered_map_keys_preserve_canonical_text() { + let paths = BTreeMap::from([ + ("z/last.md".parse::().unwrap(), 2), + ("a/first.md".parse::().unwrap(), 1), + ]); + + assert_eq!( + serde_json::to_value(&paths).unwrap(), + json!({"a/first.md": 1, "z/last.md": 2}) + ); + assert_eq!( + serde_json::from_value::>(json!({ + "a/first.md": 1, + "z/last.md": 2 + })) + .unwrap(), + paths + ); + } + + #[test] + fn byte_limit_counts_utf8_bytes() { + assert!( + "é".repeat(MAX_WORKFLOW_PATH_BYTES / 2) + .parse::() + .is_ok() + ); + assert!( + "é".repeat(MAX_WORKFLOW_PATH_BYTES / 2 + 1) + .parse::() + .is_err() + ); + assert!("notes/\u{85}.md".parse::().is_ok()); + } +} 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..c0b35f8ad --- /dev/null +++ b/lib/foundation/fabro-types/src/workflow_version.rs @@ -0,0 +1,419 @@ +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_DEPENDENCIES: 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 version has {actual} workflow dependencies; maximum is {maximum}")] + TooManyWorkflowDependencies { 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, + workflow_dependencies: BTreeMap, +} + +impl WorkflowVersion { + pub fn new( + entrypoint: WorkflowPath, + files: BTreeMap, + workflow_dependencies: BTreeMap, + ) -> Result { + let version = Self { + entrypoint, + files, + workflow_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 workflow_dependencies(&self) -> &BTreeMap { + &self.workflow_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, + }); + } + if self.workflow_dependencies.len() > MAX_WORKFLOW_VERSION_DEPENDENCIES { + return Err(WorkflowVersionShapeError::TooManyWorkflowDependencies { + actual: self.workflow_dependencies.len(), + maximum: MAX_WORKFLOW_VERSION_DEPENDENCIES, + }); + } + 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 workflow dependencies. + let mut paths = self + .files + .keys() + .chain(self.workflow_dependencies.keys()) + .collect::>(); + paths.sort_unstable(); + for pair in paths.windows(2) { + let [first, second] = pair else { + unreachable!("a two-item window must contain two paths") + }; + if first == second || first.is_ancestor_of(second) { + 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, + workflow_dependencies: UniqueBTreeMap, + } + + let wire = Wire::deserialize(deserializer)?; + Self::new(wire.entrypoint, wire.files.0, wire.workflow_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_DEPENDENCIES, + MAX_WORKFLOW_VERSION_FILE_BYTES, MAX_WORKFLOW_VERSION_FILES, WorkflowVersion, + WorkflowVersionShapeError, + }; + use crate::{BlobHash, WorkflowPath, WorkflowVersionId}; + + 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"},"workflow_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 enforces_workflow_dependency_count_boundary() { + let dependencies = (0..MAX_WORKFLOW_VERSION_DEPENDENCIES) + .map(|index| { + ( + path(&format!("dependency-{index:03}.fabro")), + WorkflowVersionId::from(BlobHash::new(index.to_string().as_bytes())), + ) + }) + .collect::>(); + let files = BTreeMap::from([(path("workflow.fabro"), "digraph W {}".to_owned())]); + assert!( + WorkflowVersion::new(path("workflow.fabro"), files.clone(), dependencies.clone()) + .is_ok() + ); + + let mut oversized = dependencies; + oversized.insert( + path("too-many.fabro"), + WorkflowVersionId::from(BlobHash::new(b"too many")), + ); + assert!(matches!( + WorkflowVersion::new(path("workflow.fabro"), files, oversized).unwrap_err(), + WorkflowVersionShapeError::TooManyWorkflowDependencies { .. } + )); + } + + #[test] + fn deserialize_rejects_unknown_fields_and_duplicate_keys() { + let unknown = r#"{ + "entrypoint":"workflow.fabro", + "files":{"workflow.fabro":"digraph W {}"}, + "workflow_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 {}"}, + "workflow_dependencies":{} + }"#; + assert!(serde_json::from_str::(duplicate).is_err()); + } +} diff --git a/lib/foundation/fabro-types/src/workflow_version_id.rs b/lib/foundation/fabro-types/src/workflow_version_id.rs new file mode 100644 index 000000000..f3a618482 --- /dev/null +++ b/lib/foundation/fabro-types/src/workflow_version_id.rs @@ -0,0 +1,96 @@ +use std::fmt; +use std::str::FromStr; + +use serde::{Deserialize, Serialize}; +use thiserror::Error; + +use crate::BlobHash; + +#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] +#[serde(into = "String", try_from = "String")] +pub struct WorkflowVersionId(BlobHash); + +impl From for WorkflowVersionId { + fn from(value: BlobHash) -> Self { + Self(value) + } +} + +impl From for BlobHash { + fn from(value: WorkflowVersionId) -> Self { + value.0 + } +} + +impl fmt::Display for WorkflowVersionId { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + fmt::Display::fmt(&self.0, f) + } +} + +impl From for String { + fn from(value: WorkflowVersionId) -> Self { + value.to_string() + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq, Error)] +#[error("workflow version ID must be exactly 64 lowercase hexadecimal characters")] +pub struct WorkflowVersionIdParseError; + +impl FromStr for WorkflowVersionId { + type Err = WorkflowVersionIdParseError; + + fn from_str(value: &str) -> Result { + // `BlobHash` enforces length and hex charset but accepts uppercase digits; + // the canonical wire form is lowercase only. + if value.bytes().any(|byte| byte.is_ascii_uppercase()) { + return Err(WorkflowVersionIdParseError); + } + value + .parse::() + .map(Self) + .map_err(|_| WorkflowVersionIdParseError) + } +} + +impl TryFrom for WorkflowVersionId { + type Error = WorkflowVersionIdParseError; + + fn try_from(value: String) -> Result { + value.parse() + } +} + +#[cfg(test)] +mod tests { + use crate::{BlobHash, WorkflowVersionId}; + + #[test] + fn conversion_preserves_digest_and_display() { + let blob_id = BlobHash::new(b"workflow"); + let version_id = WorkflowVersionId::from(blob_id); + assert_eq!(version_id.to_string(), blob_id.to_string()); + assert_eq!(BlobHash::from(version_id), blob_id); + } + + #[test] + fn parse_and_serde_require_lowercase_hex() { + let value = BlobHash::new(b"workflow").to_string(); + let id: WorkflowVersionId = value.parse().unwrap(); + assert_eq!(serde_json::to_value(id).unwrap(), value); + assert!(value.to_uppercase().parse::().is_err()); + for invalid in [ + String::new(), + "0".repeat(63), + "0".repeat(65), + "g".repeat(64), + ] { + assert!(invalid.parse::().is_err()); + } + assert!( + serde_json::from_value::(serde_json::json!(value.to_uppercase())) + .is_err() + ); + } +} diff --git a/lib/packages/fabro-api-client/src/.openapi-generator/FILES b/lib/packages/fabro-api-client/src/.openapi-generator/FILES index 4f99a6d3d..4f8d5f871 100644 --- a/lib/packages/fabro-api-client/src/.openapi-generator/FILES +++ b/lib/packages/fabro-api-client/src/.openapi-generator/FILES @@ -22,6 +22,7 @@ api/sessions-api.ts api/settings-api.ts api/system-api.ts api/variables-api.ts +api/workflow-versions-api.ts api/workflows-api.ts base.ts common.ts @@ -105,6 +106,7 @@ models/create-run-pull-request-request.ts models/create-run-session-request.ts models/create-secret-request.ts models/create-variable-request.ts +models/create-workflow-version-response.ts models/delete-run-response.ts models/delete-run-sandbox.ts models/delete-secret-request.ts @@ -541,4 +543,5 @@ models/workflow-ref.ts models/workflow-reference.ts models/workflow-schedule-summary.ts models/workflow-settings.ts +models/workflow-version.ts models/write-blob-response.ts diff --git a/lib/packages/fabro-api-client/src/api.ts b/lib/packages/fabro-api-client/src/api.ts index 02fa11e7b..3d4ddb3c7 100644 --- a/lib/packages/fabro-api-client/src/api.ts +++ b/lib/packages/fabro-api-client/src/api.ts @@ -37,4 +37,5 @@ export * from './api/sessions-api'; export * from './api/settings-api'; export * from './api/system-api'; export * from './api/variables-api'; +export * from './api/workflow-versions-api'; export * from './api/workflows-api'; diff --git a/lib/packages/fabro-api-client/src/api/workflow-versions-api.ts b/lib/packages/fabro-api-client/src/api/workflow-versions-api.ts new file mode 100644 index 000000000..edd6a2ccf --- /dev/null +++ b/lib/packages/fabro-api-client/src/api/workflow-versions-api.ts @@ -0,0 +1,134 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.1.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + +import type { Configuration } from '../configuration'; +import type { AxiosPromise, AxiosInstance, RawAxiosRequestConfig } from 'axios'; +import globalAxios from 'axios'; +// Some imports not used depending on template conditions +// @ts-ignore +import { DUMMY_BASE_URL, assertParamExists, setApiKeyToObject, setBasicAuthToObject, setBearerAuthToObject, setOAuthToObject, setSearchParams, serializeDataIfNeeded, toPathString, createRequestFunction, replaceWithSerializableTypeIfNeeded } from '../common'; +// @ts-ignore +import { BASE_PATH, COLLECTION_FORMATS, type RequestArgs, BaseAPI, RequiredError, operationServerMap } from '../base'; +// @ts-ignore +import type { CreateWorkflowVersionResponse } from '../models'; +// @ts-ignore +import type { ErrorResponse } from '../models'; +// @ts-ignore +import type { WorkflowVersion } from '../models'; +/** + * WorkflowVersionsApi - axios parameter creator + */ +export const WorkflowVersionsApiAxiosParamCreator = function (configuration?: Configuration) { + return { + /** + * Validates and stores an immutable workflow package in content-addressed storage. Repeating the same canonical content returns the same identifier. + * @summary Create Workflow Version + * @param {WorkflowVersion} workflowVersion + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + createWorkflowVersion: async (workflowVersion: WorkflowVersion, options: RawAxiosRequestConfig = {}): Promise => { + // verify required parameter 'workflowVersion' is not null or undefined + assertParamExists('createWorkflowVersion', 'workflowVersion', workflowVersion) + const localVarPath = `/api/v1/workflow-versions`; + // use dummy base URL string because the URL constructor only accepts absolute URLs. + const localVarUrlObj = new URL(localVarPath, DUMMY_BASE_URL); + let baseOptions; + if (configuration) { + baseOptions = configuration.baseOptions; + } + + const localVarRequestOptions = { method: 'POST', ...baseOptions, ...options}; + const localVarHeaderParameter = {} as any; + const localVarQueryParameter = {} as any; + + // authentication SessionCookie required + + // authentication BearerAuth required + // http bearer authentication required + await setBearerAuthToObject(localVarHeaderParameter, configuration) + + localVarHeaderParameter['Content-Type'] = 'application/json'; + localVarHeaderParameter['Accept'] = 'application/json'; + + setSearchParams(localVarUrlObj, localVarQueryParameter); + let headersFromBaseOptions = baseOptions && baseOptions.headers ? baseOptions.headers : {}; + localVarRequestOptions.headers = {...localVarHeaderParameter, ...headersFromBaseOptions, ...options.headers}; + localVarRequestOptions.data = serializeDataIfNeeded(workflowVersion, localVarRequestOptions, configuration) + + return { + url: toPathString(localVarUrlObj), + options: localVarRequestOptions, + }; + }, + } +}; + +/** + * WorkflowVersionsApi - functional programming interface + */ +export const WorkflowVersionsApiFp = function(configuration?: Configuration) { + const localVarAxiosParamCreator = WorkflowVersionsApiAxiosParamCreator(configuration) + return { + /** + * Validates and stores an immutable workflow package in content-addressed storage. Repeating the same canonical content returns the same identifier. + * @summary Create Workflow Version + * @param {WorkflowVersion} workflowVersion + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + async createWorkflowVersion(workflowVersion: WorkflowVersion, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + const localVarAxiosArgs = await localVarAxiosParamCreator.createWorkflowVersion(workflowVersion, options); + const localVarOperationServerIndex = configuration?.serverIndex ?? 0; + const localVarOperationServerBasePath = operationServerMap['WorkflowVersionsApi.createWorkflowVersion']?.[localVarOperationServerIndex]?.url; + return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); + }, + } +}; + +/** + * WorkflowVersionsApi - factory interface + */ +export const WorkflowVersionsApiFactory = function (configuration?: Configuration, basePath?: string, axios?: AxiosInstance) { + const localVarFp = WorkflowVersionsApiFp(configuration) + return { + /** + * Validates and stores an immutable workflow package in content-addressed storage. Repeating the same canonical content returns the same identifier. + * @summary Create Workflow Version + * @param {WorkflowVersion} workflowVersion + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + createWorkflowVersion(workflowVersion: WorkflowVersion, options?: RawAxiosRequestConfig): AxiosPromise { + return localVarFp.createWorkflowVersion(workflowVersion, options).then((request) => request(axios, basePath)); + }, + }; +}; + +/** + * WorkflowVersionsApi - object-oriented interface + */ +export class WorkflowVersionsApi extends BaseAPI { + /** + * Validates and stores an immutable workflow package in content-addressed storage. Repeating the same canonical content returns the same identifier. + * @summary Create Workflow Version + * @param {WorkflowVersion} workflowVersion + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + public createWorkflowVersion(workflowVersion: WorkflowVersion, options?: RawAxiosRequestConfig) { + return WorkflowVersionsApiFp(this.configuration).createWorkflowVersion(workflowVersion, options).then((request) => request(this.axios, this.basePath)); + } +} diff --git a/lib/packages/fabro-api-client/src/models/create-workflow-version-response.ts b/lib/packages/fabro-api-client/src/models/create-workflow-version-response.ts new file mode 100644 index 000000000..79ede2d1b --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/create-workflow-version-response.ts @@ -0,0 +1,25 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.1.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + + +/** + * Identity of the stored immutable workflow version. + */ +export interface CreateWorkflowVersionResponse { + /** + * SHA-256 identity of validated canonical workflow-version bytes. + */ + 'workflow_version_id': string; +} diff --git a/lib/packages/fabro-api-client/src/models/index.ts b/lib/packages/fabro-api-client/src/models/index.ts index 2fdd2e9fe..04df237be 100644 --- a/lib/packages/fabro-api-client/src/models/index.ts +++ b/lib/packages/fabro-api-client/src/models/index.ts @@ -76,6 +76,7 @@ export * from './create-run-pull-request-request'; export * from './create-run-session-request'; export * from './create-secret-request'; export * from './create-variable-request'; +export * from './create-workflow-version-response'; export * from './delete-run-response'; export * from './delete-run-sandbox'; export * from './delete-secret-request'; @@ -511,4 +512,5 @@ export * from './workflow-ref'; export * from './workflow-reference'; export * from './workflow-schedule-summary'; export * from './workflow-settings'; +export * from './workflow-version'; export * from './write-blob-response'; diff --git a/lib/packages/fabro-api-client/src/models/workflow-version.ts b/lib/packages/fabro-api-client/src/models/workflow-version.ts new file mode 100644 index 000000000..d54f4064e --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/workflow-version.ts @@ -0,0 +1,33 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.1.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + + +/** + * Complete immutable package for one rooted workflow. It contains at most 512 files and 512 workflow dependencies, each file is at most 512 KiB of UTF-8 content, and its compact canonical JSON representation is at most 2 MiB. + */ +export interface WorkflowVersion { + /** + * Canonical portable path inside one workflow version. Paths are UTF-8, relative, at most 240 bytes and 16 components, and cannot contain empty, dot, parent, backslash, control, tilde-root, or drive-letter segments. Map keys receive stricter byte and structural validation in the domain model than OpenAPI can express. + */ + 'entrypoint': string; + /** + * Workflow-local text files keyed by canonical path. Keys receive stricter domain validation than OpenAPI can express; each value is limited to 512 KiB of UTF-8 bytes. + */ + 'files': { [key: string]: string; }; + /** + * Exact stored workflow-version IDs keyed by resolved child-workflow path. Keys receive stricter domain validation than OpenAPI can express. + */ + 'workflow_dependencies': { [key: string]: string; }; +} From 79e44262aa230b21fae212a5fe40eb1537c4caa2 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Thu, 13 Aug 2026 10:37:19 -0400 Subject: [PATCH 2/8] Detect workflow path collisions hidden by sort order The adjacent-pair scan over the byte-sorted path list missed file/directory collisions whenever a sibling path sorted between the ancestor and its descendant (any byte below '/' after the shared prefix, e.g. "assets.txt" between "assets" and "assets/item.txt"). Replace it with an exhaustive ancestor-prefix lookup over a path set, which also catches equal paths across files and workflow dependencies. Co-Authored-By: Claude Fable 5 --- .../fabro-types/src/workflow_version.rs | 104 +++++++++++++++--- 1 file changed, 90 insertions(+), 14 deletions(-) diff --git a/lib/foundation/fabro-types/src/workflow_version.rs b/lib/foundation/fabro-types/src/workflow_version.rs index c0b35f8ad..cba3a04c1 100644 --- a/lib/foundation/fabro-types/src/workflow_version.rs +++ b/lib/foundation/fabro-types/src/workflow_version.rs @@ -1,4 +1,4 @@ -use std::collections::BTreeMap; +use std::collections::{BTreeMap, HashMap}; use std::fmt; use std::marker::PhantomData; @@ -135,23 +135,27 @@ impl WorkflowVersion { fn validate_path_collisions(&self) -> Result<(), WorkflowVersionShapeError> { // Keys are unique within each map, so equality can only collide // across files and workflow dependencies. - let mut paths = self - .files - .keys() - .chain(self.workflow_dependencies.keys()) - .collect::>(); - paths.sort_unstable(); - for pair in paths.windows(2) { - let [first, second] = pair else { - unreachable!("a two-item window must contain two paths") - }; - if first == second || first.is_ancestor_of(second) { + let mut by_text = + HashMap::with_capacity(self.files.len() + self.workflow_dependencies.len()); + for path in self.files.keys().chain(self.workflow_dependencies.keys()) { + if let Some(existing) = by_text.insert(path.as_str(), path) { return Err(WorkflowVersionShapeError::PathCollision { - first: (*first).clone(), - second: (*second).clone(), + first: existing.clone(), + second: path.clone(), }); } } + for path in self.files.keys().chain(self.workflow_dependencies.keys()) { + let text = path.as_str(); + for (index, _) in text.match_indices('/') { + if let Some(ancestor) = by_text.get(&text[..index]) { + return Err(WorkflowVersionShapeError::PathCollision { + first: (*ancestor).clone(), + second: path.clone(), + }); + } + } + } Ok(()) } } @@ -295,6 +299,78 @@ mod tests { )); } + #[test] + fn rejects_ancestor_collisions_hidden_by_sort_order() { + // `assets.txt` sorts between `assets` and `assets/item.txt` because + // '.' precedes '/', so an adjacent-pair scan over the sorted list + // would miss this collision. + let error = WorkflowVersion::new( + path("workflow.fabro"), + BTreeMap::from([ + (path("workflow.fabro"), "digraph W {}".to_string()), + (path("assets"), "file".to_string()), + (path("assets.txt"), "sibling".to_string()), + (path("assets/item.txt"), "nested".to_string()), + ]), + BTreeMap::new(), + ) + .unwrap_err(); + assert!(matches!( + error, + WorkflowVersionShapeError::PathCollision { first, second } + if first.as_str() == "assets" && second.as_str() == "assets/item.txt" + )); + + assert!( + WorkflowVersion::new( + path("workflow.fabro"), + BTreeMap::from([ + (path("workflow.fabro"), "digraph W {}".to_string()), + (path("assets.txt"), "sibling".to_string()), + (path("assets/item.txt"), "nested".to_string()), + ]), + BTreeMap::new(), + ) + .is_ok() + ); + } + + #[test] + fn rejects_collisions_across_files_and_workflow_dependencies() { + let dependency_id = WorkflowVersionId::from(BlobHash::new(b"child")); + + let equal = WorkflowVersion::new( + path("workflow.fabro"), + BTreeMap::from([ + (path("workflow.fabro"), "digraph W {}".to_string()), + (path("child.fabro"), "digraph C {}".to_string()), + ]), + BTreeMap::from([(path("child.fabro"), dependency_id)]), + ) + .unwrap_err(); + assert!(matches!( + equal, + WorkflowVersionShapeError::PathCollision { first, second } + if first == second && first.as_str() == "child.fabro" + )); + + let ancestor = WorkflowVersion::new( + path("workflow.fabro"), + BTreeMap::from([ + (path("workflow.fabro"), "digraph W {}".to_string()), + (path("libs"), "file".to_string()), + (path("libs.md"), "sibling".to_string()), + ]), + BTreeMap::from([(path("libs/child.fabro"), dependency_id)]), + ) + .unwrap_err(); + assert!(matches!( + ancestor, + WorkflowVersionShapeError::PathCollision { first, second } + if first.as_str() == "libs" && second.as_str() == "libs/child.fabro" + )); + } + #[test] fn enforces_file_count_file_size_and_canonical_size_boundaries() { let mut files = BTreeMap::from([(path("workflow.fabro"), "digraph W {}".to_string())]); From e688bd98768575c74e0d3d0b890496a6decc1ad2 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Thu, 13 Aug 2026 10:37:19 -0400 Subject: [PATCH 3/8] Return 422 for invalid workflow-version dependencies DependencyInvalid fell through to the curated 500 even though the OpenAPI contract promises 422 workflow_version_dependency_not_found for an absent, invalid, or non-canonical dependency. Route it to that response alongside DependencyNotFound; the top-level message only names the caller-supplied path and id, so no internal chain leaks. Drop the InvalidVersion/InvalidShape arms, which were unreachable from the only call site. Co-Authored-By: Claude Fable 5 --- .../src/server/handler/workflow_versions.rs | 28 ++++++++----------- 1 file changed, 12 insertions(+), 16 deletions(-) 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 ebff58d82..beb5802f5 100644 --- a/lib/apps/fabro-server/src/server/handler/workflow_versions.rs +++ b/lib/apps/fabro-server/src/server/handler/workflow_versions.rs @@ -83,21 +83,14 @@ fn json_rejection(rejection: JsonRejection) -> ApiError { fn store_error(err: WorkflowVersionStoreError) -> ApiError { match err { - err @ WorkflowVersionStoreError::DependencyNotFound { .. } => ApiError::with_code( + // The top-level message names the offending dependency without its + // internal source chain, so it is safe to surface to the caller. + err @ (WorkflowVersionStoreError::DependencyNotFound { .. } + | WorkflowVersionStoreError::DependencyInvalid { .. }) => ApiError::with_code( StatusCode::UNPROCESSABLE_ENTITY, err.to_string(), DEPENDENCY_NOT_FOUND_CODE, ), - WorkflowVersionStoreError::InvalidVersion(source) => ApiError::with_code( - StatusCode::UNPROCESSABLE_ENTITY, - 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, @@ -264,7 +257,7 @@ mod tests { } #[tokio::test] - async fn corrupt_stored_dependency_returns_curated_internal_error() { + async fn invalid_stored_dependency_is_a_client_error_without_internals() { let state = TestAppStateBuilder::new().build(); let app = test_support::build_test_router(Arc::clone(&state)); let dependency_id = WorkflowVersionId::from( @@ -289,11 +282,14 @@ mod tests { .oneshot(request(serde_json::to_vec(&root).unwrap())) .await .unwrap(); - assert_eq!(response.status(), StatusCode::INTERNAL_SERVER_ERROR); + assert_eq!(response.status(), StatusCode::UNPROCESSABLE_ENTITY); let body = response_json(response).await; - assert_eq!( - body["errors"][0]["detail"], - "workflow version store operation failed" + assert_eq!(error_code(&body), DEPENDENCY_NOT_FOUND_CODE); + assert!( + body["errors"][0]["detail"] + .as_str() + .unwrap() + .contains("child.fabro") ); assert!(!body.to_string().contains("cannot be decoded")); } From 20c9fba0b13975c95d95d18e94684bb70dbda10e Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Thu, 13 Aug 2026 10:42:16 -0400 Subject: [PATCH 4/8] Serialize workflow-version canonical bytes once at construction WorkflowVersion::new serialized the whole version just to enforce the size limit and threw the bytes away, the store re-serialized them to write the blob, and every read re-serialized a third time for the canonicality comparison. Cache the canonical bytes on the struct at construction (skipped during serde) and expose them as an infallible borrow; the now-unconstructable InvalidShape store error variant goes away with it. Co-Authored-By: Claude Fable 5 --- .../fabro-workflow-version/src/store.rs | 28 +++++-------- .../fabro-types/src/workflow_version.rs | 40 ++++++++++--------- 2 files changed, 32 insertions(+), 36 deletions(-) diff --git a/lib/components/fabro-workflow-version/src/store.rs b/lib/components/fabro-workflow-version/src/store.rs index 80dcebd53..fd52cebe5 100644 --- a/lib/components/fabro-workflow-version/src/store.rs +++ b/lib/components/fabro-workflow-version/src/store.rs @@ -2,7 +2,7 @@ use std::collections::{BTreeMap, HashSet, VecDeque}; use std::sync::Arc; use fabro_store::BlobStore; -use fabro_types::{WorkflowPath, WorkflowVersion, WorkflowVersionId, WorkflowVersionShapeError}; +use fabro_types::{WorkflowPath, WorkflowVersion, WorkflowVersionId}; use thiserror::Error; use crate::{ValidatedWorkflowVersion, WorkflowVersionError}; @@ -11,8 +11,6 @@ use crate::{ValidatedWorkflowVersion, WorkflowVersionError}; 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, @@ -60,11 +58,10 @@ impl WorkflowVersionStore { &self, version: &ValidatedWorkflowVersion, ) -> Result { - let canonical = version.version().canonical_bytes()?; self.validate_dependency_closure(version.version().workflow_dependencies()) .await?; self.blobs - .write(&canonical) + .write(version.version().canonical_bytes()) .await .map(WorkflowVersionId::from) .map_err(|source| WorkflowVersionStoreError::Storage { source }) @@ -98,8 +95,7 @@ impl WorkflowVersionStore { let version = serde_json::from_slice::(&bytes) .map_err(|source| WorkflowVersionStoreError::Decode { id: *id, source })?; let validated = ValidatedWorkflowVersion::new(version)?; - let canonical = validated.version().canonical_bytes()?; - if canonical.as_slice() != bytes.as_ref() { + if validated.version().canonical_bytes() != bytes.as_ref() { return Err(WorkflowVersionStoreError::NonCanonical { id: *id }); } Ok(Some(validated)) @@ -196,7 +192,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.version().canonical_bytes().unwrap(); + let expected_bytes = version.version().canonical_bytes().to_vec(); let expected_id = WorkflowVersionId::from(fabro_types::BlobHash::new(&expected_bytes)); let id = store.put(&version).await.unwrap(); @@ -228,15 +224,14 @@ mod tests { let (blobs, store) = stores().await; let child = version("digraph Child {}", BTreeMap::new()); let child_id = WorkflowVersionId::from(fabro_types::BlobHash::new( - &child.version().canonical_bytes().unwrap(), + child.version().canonical_bytes(), )); 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::BlobHash::new( - &root.version().canonical_bytes().unwrap(), - )); + let root_id = + WorkflowVersionId::from(fabro_types::BlobHash::new(root.version().canonical_bytes())); let error = store.put(&root).await.unwrap_err(); assert!(matches!( @@ -256,15 +251,14 @@ mod tests { r#"digraph Child { grandchild [stack.child_workflow="grandchild.fabro"] }"#, BTreeMap::from([(path("grandchild.fabro"), missing_grandchild_id)]), ); - let child_bytes = child.version().canonical_bytes().unwrap(); - let child_id = WorkflowVersionId::from(blobs.write(&child_bytes).await.unwrap()); + let child_bytes = child.version().canonical_bytes(); + let child_id = WorkflowVersionId::from(blobs.write(child_bytes).await.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::BlobHash::new( - &root.version().canonical_bytes().unwrap(), - )); + let root_id = + WorkflowVersionId::from(fabro_types::BlobHash::new(root.version().canonical_bytes())); assert!(matches!( store.put(&root).await.unwrap_err(), diff --git a/lib/foundation/fabro-types/src/workflow_version.rs b/lib/foundation/fabro-types/src/workflow_version.rs index cba3a04c1..9b4570943 100644 --- a/lib/foundation/fabro-types/src/workflow_version.rs +++ b/lib/foundation/fabro-types/src/workflow_version.rs @@ -53,6 +53,8 @@ pub struct WorkflowVersion { entrypoint: WorkflowPath, files: BTreeMap, workflow_dependencies: BTreeMap, + #[serde(skip)] + canonical: Vec, } impl WorkflowVersion { @@ -61,13 +63,22 @@ impl WorkflowVersion { files: BTreeMap, workflow_dependencies: BTreeMap, ) -> Result { - let version = Self { + let mut version = Self { entrypoint, files, workflow_dependencies, + canonical: Vec::new(), }; version.validate_shape()?; - version.canonical_bytes()?; + let canonical = serde_json::to_vec(&version) + .map_err(|source| WorkflowVersionShapeError::Serialization { source })?; + if canonical.len() > MAX_WORKFLOW_VERSION_BYTES { + return Err(WorkflowVersionShapeError::VersionTooLarge { + actual: canonical.len(), + maximum: MAX_WORKFLOW_VERSION_BYTES, + }); + } + version.canonical = canonical; Ok(version) } @@ -86,20 +97,11 @@ impl WorkflowVersion { &self.workflow_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) + /// Canonical wire bytes, serialized and size-checked once at + /// construction. + #[must_use] + pub fn canonical_bytes(&self) -> &[u8] { + &self.canonical } fn validate_shape(&self) -> Result<(), WorkflowVersionShapeError> { @@ -250,7 +252,7 @@ mod tests { .unwrap(); assert_eq!( - String::from_utf8(version.canonical_bytes().unwrap()).unwrap(), + String::from_utf8(version.canonical_bytes().to_vec()).unwrap(), r#"{"entrypoint":"workflow.fabro","files":{"a.txt":"A","workflow.fabro":"digraph W {}","z.txt":"Z"},"workflow_dependencies":{}}"# ); } @@ -419,7 +421,7 @@ mod tests { BTreeMap::new(), ) .unwrap(); - let remaining = MAX_WORKFLOW_VERSION_BYTES - empty.canonical_bytes().unwrap().len(); + let remaining = MAX_WORKFLOW_VERSION_BYTES - empty.canonical_bytes().len(); let per_file = remaining / 4; let remainder = remaining % 4; for index in 0..4 { @@ -434,7 +436,7 @@ mod tests { ) .unwrap(); assert_eq!( - exact_version.canonical_bytes().unwrap().len(), + exact_version.canonical_bytes().len(), MAX_WORKFLOW_VERSION_BYTES ); exact_version_files From cc711027da7559c3e28687a3a2e5bf0988bb37f0 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Thu, 13 Aug 2026 10:43:26 -0400 Subject: [PATCH 5/8] Parse workflow version IDs case-insensitively MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit WorkflowVersionId bolted a lowercase-only byte scan onto BlobHash parsing, giving the same 64-hex concept two parse behaviors across entry points. Identity is the decoded 32-byte digest and canonical serialization always emits lowercase, so accepting either case on input is lossless — the stored-blob canonicality check still rejects non-canonical bytes independently. Delegate straight to BlobHash. Co-Authored-By: Claude Fable 5 --- .../fabro-types/src/workflow_version_id.rs | 20 +++++++++---------- 1 file changed, 10 insertions(+), 10 deletions(-) diff --git a/lib/foundation/fabro-types/src/workflow_version_id.rs b/lib/foundation/fabro-types/src/workflow_version_id.rs index f3a618482..4bf597f47 100644 --- a/lib/foundation/fabro-types/src/workflow_version_id.rs +++ b/lib/foundation/fabro-types/src/workflow_version_id.rs @@ -35,18 +35,13 @@ impl From for String { } #[derive(Clone, Copy, Debug, PartialEq, Eq, Error)] -#[error("workflow version ID must be exactly 64 lowercase hexadecimal characters")] +#[error("workflow version ID must be exactly 64 hexadecimal characters")] pub struct WorkflowVersionIdParseError; impl FromStr for WorkflowVersionId { type Err = WorkflowVersionIdParseError; fn from_str(value: &str) -> Result { - // `BlobHash` enforces length and hex charset but accepts uppercase digits; - // the canonical wire form is lowercase only. - if value.bytes().any(|byte| byte.is_ascii_uppercase()) { - return Err(WorkflowVersionIdParseError); - } value .parse::() .map(Self) @@ -75,11 +70,14 @@ mod tests { } #[test] - fn parse_and_serde_require_lowercase_hex() { + fn parse_accepts_any_case_and_serializes_lowercase() { let value = BlobHash::new(b"workflow").to_string(); let id: WorkflowVersionId = value.parse().unwrap(); assert_eq!(serde_json::to_value(id).unwrap(), value); - assert!(value.to_uppercase().parse::().is_err()); + assert_eq!( + value.to_uppercase().parse::().unwrap(), + id + ); for invalid in [ String::new(), "0".repeat(63), @@ -88,9 +86,11 @@ mod tests { ] { assert!(invalid.parse::().is_err()); } - assert!( + assert_eq!( serde_json::from_value::(serde_json::json!(value.to_uppercase())) - .is_err() + .unwrap() + .to_string(), + value ); } } From 75fb1e3a1d41aa36a1571738ec5db57282665dc9 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Thu, 13 Aug 2026 10:44:58 -0400 Subject: [PATCH 6/8] Prove workflow-version ID JSON parity with the OpenAPI schema The round-trip fixtures only used empty workflow_dependencies, so no WorkflowVersionId value ever appeared on the wire in a fabro-api assertion and CreateWorkflowVersionResponse had no coverage at all. Put a real 64-hex id in the fixture, round-trip the response type, and pin serialization to the schema's ^[0-9a-f]{64}$ pattern including lowercase normalization of case-insensitive input. Co-Authored-By: Claude Fable 5 --- .../tests/workflow_version_round_trip.rs | 33 +++++++++++++++++-- 1 file changed, 30 insertions(+), 3 deletions(-) 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 894c1c910..3bf0c8835 100644 --- a/lib/foundation/fabro-api/tests/workflow_version_round_trip.rs +++ b/lib/foundation/fabro-api/tests/workflow_version_round_trip.rs @@ -1,12 +1,14 @@ use std::any::{TypeId, type_name}; use fabro_api::types::{ - WorkflowPath as ApiWorkflowPath, WorkflowVersion as ApiWorkflowVersion, - WorkflowVersionId as ApiWorkflowVersionId, + CreateWorkflowVersionResponse, WorkflowPath as ApiWorkflowPath, + WorkflowVersion as ApiWorkflowVersion, WorkflowVersionId as ApiWorkflowVersionId, }; use fabro_types::{WorkflowPath, WorkflowVersion, WorkflowVersionId}; use serde_json::json; +const DEPENDENCY_ID: &str = "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"; + #[test] fn workflow_version_schemas_reuse_domain_types() { assert_same_type::(); @@ -22,13 +24,38 @@ fn workflow_version_round_trips_exact_wire_shape() { "prompts/goal.md": "Ship it", "workflow.fabro": "digraph W { start [shape=Mdiamond] exit [shape=Msquare] start -> exit }" }, - "workflow_dependencies": {} + "workflow_dependencies": { "children/check.fabro": DEPENDENCY_ID } }); let version: ApiWorkflowVersion = serde_json::from_value(value.clone()).unwrap(); assert_eq!(serde_json::to_value(version).unwrap(), value); } +#[test] +fn create_workflow_version_response_round_trips_exact_wire_shape() { + let value = json!({ "workflow_version_id": DEPENDENCY_ID }); + + let response: CreateWorkflowVersionResponse = serde_json::from_value(value.clone()).unwrap(); + assert_eq!(serde_json::to_value(&response).unwrap(), value); +} + +#[test] +fn workflow_version_id_emits_the_documented_lowercase_pattern() { + // Input is accepted case-insensitively, but serialization must match the + // OpenAPI schema pattern `^[0-9a-f]{64}$`. + let id = serde_json::from_value::(json!(DEPENDENCY_ID.to_uppercase())) + .unwrap(); + let emitted = serde_json::to_value(id).unwrap(); + assert_eq!(emitted, json!(DEPENDENCY_ID)); + + let text = emitted.as_str().unwrap(); + assert_eq!(text.len(), 64); + assert!( + text.bytes() + .all(|byte| matches!(byte, b'0'..=b'9' | b'a'..=b'f')) + ); +} + #[test] fn workflow_version_replacement_rejects_unknown_fields() { let value = json!({ From 180330c117423816f0a222cc8abbdf06bde9d334 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Thu, 13 Aug 2026 10:48:45 -0400 Subject: [PATCH 7/8] Remove unused WorkflowPath::parent and is_ancestor_of Neither method has callers anywhere in the workspace: resolve_reference splits on '/' directly, and the path-collision validator now checks ancestor prefixes against a path set. parent() also constructed Self without going through validate(), so dropping it removes an unvalidated construction path from the wire type's public API. Co-Authored-By: Claude Fable 5 --- .../fabro-types/src/workflow_path.rs | 21 ------------------- 1 file changed, 21 deletions(-) diff --git a/lib/foundation/fabro-types/src/workflow_path.rs b/lib/foundation/fabro-types/src/workflow_path.rs index 32297d93e..a489153a9 100644 --- a/lib/foundation/fabro-types/src/workflow_path.rs +++ b/lib/foundation/fabro-types/src/workflow_path.rs @@ -49,20 +49,6 @@ impl WorkflowPath { &self.0 } - #[must_use] - pub fn parent(&self) -> Option { - self.0 - .rsplit_once('/') - .map(|(parent, _)| Self(parent.to_owned())) - } - - #[must_use] - pub fn is_ancestor_of(&self, other: &Self) -> bool { - other.0.len() > self.0.len() - && other.0.starts_with(self.0.as_str()) - && other.0.as_bytes()[self.0.len()] == b'/' - } - pub fn resolve_reference(&self, reference: &str) -> Result { validate_reference_shape(reference)?; let mut components = self @@ -256,13 +242,6 @@ mod tests { assert!(graph.resolve_reference("prompts/").is_err()); } - #[test] - fn ancestor_checks_component_boundaries() { - let parent: WorkflowPath = "dir/file".parse().unwrap(); - assert!(parent.is_ancestor_of(&"dir/file/child".parse().unwrap())); - assert!(!parent.is_ancestor_of(&"dir/filename".parse().unwrap())); - } - #[test] fn serde_and_ordered_map_keys_preserve_canonical_text() { let paths = BTreeMap::from([ From 8c3ff6216cb4ef049fdb20b7dbaa0403b39d7b9b Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Thu, 13 Aug 2026 16:02:12 -0400 Subject: [PATCH 8/8] Revert "Serialize workflow-version canonical bytes once at construction" This reverts commit 8090d7030984862564a929ee9264e93911014e00. The cached canonical field was optimizing an unmeasured path: without the (deferred) O(closure) dependency re-validation multiplier, the repeated serialization is microseconds for realistic versions. Compute canonical bytes on demand like the environment, automation, and MCP stores do, rather than carrying a serde-skipped cache field, a construction bootstrap, and doubled memory for it. Purely in-memory: stored blobs and version IDs are unchanged. Co-Authored-By: Claude Fable 5 --- .../fabro-workflow-version/src/store.rs | 28 ++++++++----- .../fabro-types/src/workflow_version.rs | 40 +++++++++---------- 2 files changed, 36 insertions(+), 32 deletions(-) diff --git a/lib/components/fabro-workflow-version/src/store.rs b/lib/components/fabro-workflow-version/src/store.rs index fd52cebe5..80dcebd53 100644 --- a/lib/components/fabro-workflow-version/src/store.rs +++ b/lib/components/fabro-workflow-version/src/store.rs @@ -2,7 +2,7 @@ use std::collections::{BTreeMap, HashSet, VecDeque}; use std::sync::Arc; use fabro_store::BlobStore; -use fabro_types::{WorkflowPath, WorkflowVersion, WorkflowVersionId}; +use fabro_types::{WorkflowPath, WorkflowVersion, WorkflowVersionId, WorkflowVersionShapeError}; use thiserror::Error; use crate::{ValidatedWorkflowVersion, WorkflowVersionError}; @@ -11,6 +11,8 @@ use crate::{ValidatedWorkflowVersion, WorkflowVersionError}; 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, @@ -58,10 +60,11 @@ impl WorkflowVersionStore { &self, version: &ValidatedWorkflowVersion, ) -> Result { + let canonical = version.version().canonical_bytes()?; self.validate_dependency_closure(version.version().workflow_dependencies()) .await?; self.blobs - .write(version.version().canonical_bytes()) + .write(&canonical) .await .map(WorkflowVersionId::from) .map_err(|source| WorkflowVersionStoreError::Storage { source }) @@ -95,7 +98,8 @@ impl WorkflowVersionStore { let version = serde_json::from_slice::(&bytes) .map_err(|source| WorkflowVersionStoreError::Decode { id: *id, source })?; let validated = ValidatedWorkflowVersion::new(version)?; - if validated.version().canonical_bytes() != bytes.as_ref() { + let canonical = validated.version().canonical_bytes()?; + if canonical.as_slice() != bytes.as_ref() { return Err(WorkflowVersionStoreError::NonCanonical { id: *id }); } Ok(Some(validated)) @@ -192,7 +196,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.version().canonical_bytes().to_vec(); + let expected_bytes = version.version().canonical_bytes().unwrap(); let expected_id = WorkflowVersionId::from(fabro_types::BlobHash::new(&expected_bytes)); let id = store.put(&version).await.unwrap(); @@ -224,14 +228,15 @@ mod tests { let (blobs, store) = stores().await; let child = version("digraph Child {}", BTreeMap::new()); let child_id = WorkflowVersionId::from(fabro_types::BlobHash::new( - child.version().canonical_bytes(), + &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::BlobHash::new(root.version().canonical_bytes())); + let root_id = WorkflowVersionId::from(fabro_types::BlobHash::new( + &root.version().canonical_bytes().unwrap(), + )); let error = store.put(&root).await.unwrap_err(); assert!(matches!( @@ -251,14 +256,15 @@ mod tests { r#"digraph Child { grandchild [stack.child_workflow="grandchild.fabro"] }"#, BTreeMap::from([(path("grandchild.fabro"), missing_grandchild_id)]), ); - let child_bytes = child.version().canonical_bytes(); - let child_id = WorkflowVersionId::from(blobs.write(child_bytes).await.unwrap()); + let child_bytes = child.version().canonical_bytes().unwrap(); + let child_id = WorkflowVersionId::from(blobs.write(&child_bytes).await.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::BlobHash::new(root.version().canonical_bytes())); + let root_id = WorkflowVersionId::from(fabro_types::BlobHash::new( + &root.version().canonical_bytes().unwrap(), + )); assert!(matches!( store.put(&root).await.unwrap_err(), diff --git a/lib/foundation/fabro-types/src/workflow_version.rs b/lib/foundation/fabro-types/src/workflow_version.rs index 9b4570943..cba3a04c1 100644 --- a/lib/foundation/fabro-types/src/workflow_version.rs +++ b/lib/foundation/fabro-types/src/workflow_version.rs @@ -53,8 +53,6 @@ pub struct WorkflowVersion { entrypoint: WorkflowPath, files: BTreeMap, workflow_dependencies: BTreeMap, - #[serde(skip)] - canonical: Vec, } impl WorkflowVersion { @@ -63,22 +61,13 @@ impl WorkflowVersion { files: BTreeMap, workflow_dependencies: BTreeMap, ) -> Result { - let mut version = Self { + let version = Self { entrypoint, files, workflow_dependencies, - canonical: Vec::new(), }; version.validate_shape()?; - let canonical = serde_json::to_vec(&version) - .map_err(|source| WorkflowVersionShapeError::Serialization { source })?; - if canonical.len() > MAX_WORKFLOW_VERSION_BYTES { - return Err(WorkflowVersionShapeError::VersionTooLarge { - actual: canonical.len(), - maximum: MAX_WORKFLOW_VERSION_BYTES, - }); - } - version.canonical = canonical; + version.canonical_bytes()?; Ok(version) } @@ -97,11 +86,20 @@ impl WorkflowVersion { &self.workflow_dependencies } - /// Canonical wire bytes, serialized and size-checked once at - /// construction. - #[must_use] - pub fn canonical_bytes(&self) -> &[u8] { - &self.canonical + /// 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> { @@ -252,7 +250,7 @@ mod tests { .unwrap(); assert_eq!( - String::from_utf8(version.canonical_bytes().to_vec()).unwrap(), + String::from_utf8(version.canonical_bytes().unwrap()).unwrap(), r#"{"entrypoint":"workflow.fabro","files":{"a.txt":"A","workflow.fabro":"digraph W {}","z.txt":"Z"},"workflow_dependencies":{}}"# ); } @@ -421,7 +419,7 @@ mod tests { BTreeMap::new(), ) .unwrap(); - let remaining = MAX_WORKFLOW_VERSION_BYTES - empty.canonical_bytes().len(); + 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 { @@ -436,7 +434,7 @@ mod tests { ) .unwrap(); assert_eq!( - exact_version.canonical_bytes().len(), + exact_version.canonical_bytes().unwrap().len(), MAX_WORKFLOW_VERSION_BYTES ); exact_version_files