Reference

Go FlowDSL

API reference for go-flowdsl — FlowDSL parsing, validation, compilation, and code generation.

Module: github.com/redelay/go-flowdsl

Browse the full FlowDSL Node Catalog — 147 nodes across every framework, addon, and infrastructure module, each with a settings table and example snippet.

Overview

go-flowdsl is a standalone Go library for everything FlowDSL and IR: parsing, validation, code generation, SVG diagrams, drift detection, and a flow execution engine.

Key property: the only external dependencies are github.com/google/uuid and gopkg.in/yaml.v3. There is no dependency on go-framework, MongoDB, chi, JWT, zap, or any transport backend. This makes it suitable for lightweight workers, build tools, and CLIs that need IR tooling without the full framework stack.

The dependency flows from go-framework into go-flowdsl, not the reverse:

text
go-flowdsl                    ← standalone (uuid + yaml.v3 only)
       ↑
go-framework/core/asyncapi    ← registers AsyncAPI importer into go-flowdsl/compile
go-framework/core/openapi     ← registers OpenAPI importer/exporter into go-flowdsl/compile
go-framework/core/compile     ← wraps go-flowdsl/compile; adds asyncapi + openapi formats

Install

shell
go get github.com/redelay/go-flowdsl

go.mod

text
module github.com/redelay/go-flowdsl
go 1.24.0
toolchain go1.24.4

require (
    github.com/flowdsl/flowdsl-go v0.1.0
    github.com/google/uuid v1.6.0
    gopkg.in/yaml.v3 v3.0.1
)

Package reference

ir — Canonical IR types

Defines the canonical Intermediate Representation that all import and export formats convert to and from. No external dependencies.

go
import "github.com/redelay/go-flowdsl/ir"

doc := ir.NewDocument("My App", "1.0.0")
// doc.ID        — auto-generated UUID
// doc.Modules   — business/domain modules
// doc.Entities  — domain aggregates
// doc.Events    — business events
// doc.Commands  — state-changing intents
// doc.Queries   — read operations
// doc.Actions   — top-level operations
// doc.Workflows — multi-step processes with nodes + edges
// doc.Channels  — message/event channels
// doc.Schemas   — reusable data structure definitions
// doc.Packets   — data transfer payloads
// doc.Schedules — time-based triggers (cron/interval)
// doc.Migrations — database migration definitions
// doc.Metadata  — arbitrary key-value metadata
// doc.Provenance — origin and mapping history

Key types:

TypeDescription
DocumentRoot container for an entire IR
ModuleBusiness/domain unit with entities, events, actions, workflows, config, settings, CRUD, consumers
EntityBusiness entity or domain aggregate with typed fields
FieldSingle field — type, constraints, nested properties, $ref
FieldTypeEnum: string, integer, float, boolean, array, object, date, datetime, uuid, binary, ref, enum, any
EventBusiness event — entity_type + action, optional topic and payload
CommandState-changing intent with input/output schemas and emitted events
QueryRead operation with input/output schemas
ActionTop-level operation — kind (command/query/mutation/subscription), method, path, I/O, emits
WorkflowMulti-step process with Node graph and Edge connections
NodeWorkflow step — kind (action/event/gateway/start/end/wait/condition/parallel)
EdgeDirected connection between nodes with optional condition
ChannelMessage channel with publish/subscribe events, protocol, bindings
SchemaReusable data structure definition
PacketData transfer structure (event payload, message body)
ScheduleTime-based trigger — cron expression or interval
MigrationDatabase migration — backend, version, up/down SQL
ConsumerEvent consumer — topic, group ID, handler reference, batch size
ProvenanceOrigin tracking — source, URI, source ref, mappings
ConfigDefinitionModule runtime configuration (entries with env vars, types, defaults)
SettingsDefinitionAdmin-UI-editable settings (groups of fields with components)
CRUDDefinitionCRUD capabilities for an entity (create/read/update/delete/list)
FlowDSLNodeModule-exported FlowDSL node — kind, ports, settings, refs to events/actions/workflows
FlowDSLPortNamed input or output port on a FlowDSLNode

FlowDSLNode — Module node exports

FlowDSLNode describes a FlowDSL node that a module exports. It captures the node identity, kind, ports, and settings — enough for the module browser to list nodes and for the generator to produce .flowdsl-node.json manifests.

go
type FlowDSLNode struct {
    ID             string            // unique node identifier (namespace/slug)
    Name           string            // human-readable display name
    Version        string            // semver version
    Summary        string            // one-line description
    Description    string            // full markdown description
    Kind           string            // functional category (source, transform, action, subworkflow, etc.)
    Language       string            // implementation language (go, python, nodejs)
    Published      bool              // publicly visible?
    Icon           string            // emoji or icon name for Studio
    Color          string            // hex color for Studio card
    Tags           []string          // search/filter tags
    Handler        string            // fully-qualified handler path
    Protocols      []string          // supported invocation protocols
    Inputs         []FlowDSLPort     // named input ports
    Outputs        []FlowDSLPort     // named output ports
    SettingsSchema map[string]any    // inline JSON Schema for node configuration
    ModuleRef      string            // owning redelay module
    EventRef       string            // ties node to a module event
    ActionRef      string            // ties node to a module action
    WorkflowRef    string            // ties node to a module workflow (for kind=subworkflow)
    Emits          []string          // domain event topics this node may publish
}

WorkflowRef and subworkflow nodes: When Kind is "subworkflow", WorkflowRef links the node to an ir.Workflow. The runtime resolves the reference and executes the internal graph. The node's ports are the public interface; the workflow is the private implementation.

Emits and lifecycle composition: Emits is a static declaration of the event topics an action/router node may publish when it runs — the counterpart to the consume side (a source node's configured eventName). Declared in module.yaml:

yaml
flowdsl_nodes:
  - id: orders/apply-payment
    kind: router
    module_ref: orders
    emits: [order.paid, order.updated]
    outputs:
      - name: Applied
      - name: NotFound

It powers lifecycle composition: a resolver links otherwise-decoupled event-driven flows into one graph by matching what each flow emits against what other flows consume (their source node's eventName). Nodes whose emitted topic is dynamic (chosen from config, e.g. a generic publish node) leave Emits empty and are resolved from config at graph-build time.

Subworkflow node generation

SubworkflowNodeFromWorkflow creates a subworkflow FlowDSLNode from an ir.Workflow:

go
node := ir.SubworkflowNodeFromWorkflow(wf, "redelay", "orders")
// node.Kind        = "subworkflow"
// node.WorkflowRef = wf.ID
// node.Inputs      — derived from workflow triggers
// node.Outputs     — derived from terminal nodes (kind=end or leaf nodes)

Port derivation rules:

  • Inputs: one port per workflow trigger; defaults to a single "Input" port if no triggers.
  • Outputs: one port per terminal node (kind=end) or leaf node with no outgoing edges; defaults to "Output" if none found.
  • Error output: if any edge has condition: "error" or "on_error", an "Error" output port is appended.

The framework registry auto-generates subworkflow nodes for any module that provides workflows via FlowDSLProvider without explicitly defining nodes for them. See the auto-generation section below.

Manifest conversion

FlowDSLNode has bidirectional conversion to the canonical FlowDSL SDK node.Manifest type (from github.com/flowdsl/flowdsl-go/pkg/node):

go
manifest := irNode.ToManifest()     // IR → SDK Manifest (camelCase JSON)
irNode = ir.FlowDSLNodeFromManifest(manifest) // SDK Manifest → IR (snake_case JSON)

This bridges the IR layer (snake_case, redelay-specific fields) with the FlowDSL spec layer (camelCase, schema-compliant). The SDK dependency is github.com/flowdsl/flowdsl-go.

diagnostics — Structured error/warning collector

Used throughout all validation and compilation stages.

go
import "github.com/redelay/go-flowdsl/diagnostics"

c := diagnostics.NewCollector()
c.Error("source", "CODE", "message")
c.Warning("source", "CODE", "message")
c.Info("source", "CODE", "message")

c.HasErrors()  // bool
c.Errors()     // []*Diagnostic  (error-level only)
c.Warnings()   // []*Diagnostic  (warning-level only)
c.All()        // []*Diagnostic  (all levels)

Each Diagnostic carries: Severity, Code, Message, Source, and optional File, Line, Column, Pointer, Related.

schema — JSON Schema generator

Generates JSON Schema (Draft 2020-12) from Go types via reflection. Used by go-framework/openapi for struct → JSON Schema conversion.

go
import "github.com/redelay/go-flowdsl/schema"

s := schema.From(MyStruct{})
// returns map[string]any compatible with JSON Schema Draft 2020-12

Handles json tags, validate tags (required, email, min, max), pointer nullability, time.Time, slices, maps, and embedded structs.

metadata — IR metadata helpers

Extracts summary statistics from a document and stamps compiler metadata.

go
import "github.com/redelay/go-flowdsl/metadata"

info := metadata.Extract(doc)
// info.Title, info.Version, info.Compiler
// info.ModuleCount, info.EntityCount, info.EventCount
// info.ActionCount, info.WorkflowCount, info.ChannelCount, info.ScheduleCount

metadata.Stamp(doc)              // adds compiler version + compiled_at to doc.Metadata
names := metadata.ModuleNames(doc)   // []string
names  = metadata.ActionNames(doc)   // []string
names  = metadata.EventNames(doc)    // []string

businessmodel — Business model parser

Parses, normalizes, and exports the Redelay business model YAML format.

go
import "github.com/redelay/go-flowdsl/businessmodel"

spec, err := businessmodel.ParseYAML(r)  // io.Reader → *BusinessModel
doc := businessmodel.Normalize(spec)     // *BusinessModel → *ir.Document

resolver — Module dependency resolver

Topological sort of modules by their depends_on declarations using Kahn's algorithm with cycle detection.

go
import "github.com/redelay/go-flowdsl/resolver"

sorted, diag := resolver.Resolve(doc.Modules)
// sorted — []*ir.Module in dependency order (dependencies first)
// diag   — *diagnostics.Collector
CodeSeverityRule
RES001errorDuplicate module ID
RES002errorDepends on unknown module
RES003errorCircular dependency detected
RES004errorModule is part of a dependency cycle

validate — Structural validator

Checks an IR Document for structural correctness and referential integrity.

go
import "github.com/redelay/go-flowdsl/validate"

diag := validate.Validate(doc) // *diagnostics.Collector
CodeSeverityRule
VAL001errorDocument is nil
VAL002errorDocument ID is required
VAL003warningDocument version is empty
VAL010errorDuplicate ID
VAL020errorModule depends on unknown module
VAL021warningEvent references unknown payload
VAL022warningAction references unknown input
VAL023warningAction references unknown output
VAL030errorWorkflow edge references unknown source node
VAL031errorWorkflow edge references unknown target node
VAL040errorModule has no name
VAL041warningModule ID should start with module.

flowdsl — FlowDSL YAML parser

Parses, validates, normalizes, and exports the FlowDSL YAML application definition format.

go
import "github.com/redelay/go-flowdsl/flowdsl"

parsed, err := flowdsl.Parse(r)       // io.Reader → *FlowDSL
diag := flowdsl.Validate(parsed)      // *FlowDSL → *diagnostics.Collector
doc := flowdsl.Normalize(parsed)      // *FlowDSL → *ir.Document
err = flowdsl.Export(doc, w)          // *ir.Document → FlowDSL YAML

FlowDSL YAML format:

yaml
version: "1.0"
title: My Application
modules:
  billing:
    name: billing
    description: Billing domain
    depends_on: [users]
    provides: [invoices]
entities:
  invoice:
    name: invoice
    fields:
      amount: { type: float, required: true }
      currency: { type: string, default: USD }
events:
  invoice_created:
    name: invoice_created
    entity_type: invoice
    action: created
    topic: billing.invoices
actions:
  create_invoice:
    name: create_invoice
    kind: command
    method: POST
    path: /invoices
workflows:
  payment_flow:
    name: payment_flow
    triggers: [invoice_created]
    steps:
      - id: validate
        kind: action
        action: validate_payment
        next: charge
      - id: charge
        kind: action
        action: charge_card
channels:
  notifications:
    name: notifications
    protocol: kafka
schedules:
  daily_report:
    name: daily_report
    cron: "0 9 * * *"
    action: generate_report

Validation error codes:

CodeSeverityRule
FDL001errorDocument is nil
FDL002warningVersion is empty
FDL010errorEvent missing entity_type
FDL011errorEvent missing action
FDL012warningEvent payload field has invalid type
FDL020errorWorkflow step missing ID
FDL021warningStep next references unknown step
FDL022warningStep on_error references unknown step
FDL023warningStep action references unknown action
FDL024warningStep event references unknown event
FDL025warningStep kind is not a recognised value
FDL030warningModule depends_on unknown module
FDL040warningSchedule action references unknown action
FDL050warningChannel protocol not recognised
FDL060warningAction field has invalid type
FDL061warningAction kind not recognised
FDL062warningAction emits references unknown event

modfile — Module file I/O

Loads, saves, and exports ir.Module and ir.Document values as YAML or JSON.

go
import "github.com/redelay/go-flowdsl/modfile"

// Load a module from any reader — DetectFormat infers YAML/JSON from filename.
m, err := modfile.LoadModule(r, modfile.DetectFormat(filename))

// Save a module.
err = modfile.SaveModule(m, w, modfile.FormatYAML)

// Save a full document.
err = modfile.SaveDocument(doc, w, modfile.FormatYAML)

modgen — Go module scaffolder

Generates Go source file stubs from an ir.Module. Produces module.go, schemas.go, and handler.go ready for business logic implementation. Existing files are not overwritten.

go
import "github.com/redelay/go-flowdsl/modgen"

result, err := modgen.Scaffold(m, modgen.Options{OutputDir: "./modules/orders"})
for _, path := range result.Files {
    fmt.Println("wrote", path)
}

nodegen — FlowDSL node generator

Generates FlowDSL node manifests (.flowdsl-node.json) and Go handler scaffolds from ir.Module definitions. Inspects events, actions, commands, and workflows to produce appropriate node definitions.

go
import "github.com/redelay/go-flowdsl/nodegen"

result, err := nodegen.FromIRModule(m, nodegen.Options{
    Namespace:   "myapp",
    GoModule:    "github.com/myorg/myapp",
    GoPackage:   "nodes",
    OutputDir:   "./nodes/",
    ManifestDir: "./nodes/manifests/",
})
for _, path := range result.Manifests {
    fmt.Println("manifest:", path)
}
for _, path := range result.Handlers {
    fmt.Println("handler:", path)
}

Generated output per module:

FileContents
manifests/<slug>.flowdsl-node.jsonNode manifest for each event, action, command, and workflow
nodes_generated.goGo handler stubs implementing node.NodeHandler interface
nodes_registry.goFlowDSLNodes() function returning []*ir.FlowDSLNode for module wiring

Node generation rules:

IR elementNode kindNotes
EventsourceOne output port named after the event
Actionaction or transformQueries become transform; mutations become action
CommandactionInput/output ports
WorkflowsubworkflowPorts derived via SubworkflowNodeFromWorkflow

FromFlowDocument additionally extracts nodes from workflow step graphs in a FlowDSL document.

modsvg — SVG diagram generator

Generates an SVG visual diagram from an ir.Module showing entities, events, config, CRUD, settings, commands, queries, workflows, schedules, consumers, dependencies, and provides. Supports dark mode via prefers-color-scheme CSS media query.

go
import "github.com/redelay/go-flowdsl/modsvg"

err := modsvg.Generate(w, m) // io.Writer, *ir.Module

modsync — Drift detection

Compares a declared IR module against another (e.g. a runtime export) and reports structural drift.

go
import "github.com/redelay/go-flowdsl/modsync"

report := modsync.Compare(declared, actual) // both *ir.Module
for _, diff := range report.Diffs {
    fmt.Printf("[%s] %s: %s\n", diff.Kind, diff.Path, diff.Message)
}

Diff kinds: added, removed, changed.

modval — Module YAML validator

Validates naming conventions, reference integrity, and structural rules for ir.Module files — the same checks run by redelayctl validate-module.

go
import "github.com/redelay/go-flowdsl/modval"

diag := modval.Validate(m) // *ir.Module → *diagnostics.Collector

Rules include: required fields (ID, name), semver version, entity IDs prefixed with module ID, event names containing a dot separator, config keys UPPER_SNAKE_CASE, settings keys lower_snake_case, CRUD entity refs match declared entities, no duplicate IDs, no self-dependencies.

compile — Base compile pipeline

Plugin-based compilation pipeline. Built-in importers handle flowdsl and businessmodel. go-framework/core/asyncapi and go-framework/core/openapi register additional importers and exporters via RegisterImporter/RegisterExporter.

go
import "github.com/redelay/go-flowdsl/compile"

// Register a custom importer (done by go-framework core packages at init time).
compile.RegisterImporter("asyncapi", myAsyncAPIImporter)
compile.RegisterExporter("openapi", myOpenAPIExporter)

// Import a source format into IR.
doc, err := compile.Import(compile.FormatFlowDSL, r)

// Validate + resolve + enrich.
result := compile.Compile(doc)
result.Document    // *ir.Document — enriched
result.Diagnostics // *diagnostics.Collector
result.Modules     // []*ir.Module in dependency order
result.HasErrors() // bool

// One-step import + compile.
result, err = compile.ImportAndCompile(compile.FormatFlowDSL, r)

// Export to a format.
err = compile.Export(doc, compile.FormatFlowDSL, w)

Built-in format constants:

ConstantValueImportExport
compile.FormatFlowDSL"flowdsl"yesyes
compile.FormatBusinessModel"businessmodel"yes—

Additional formats (asyncapi, openapi) are registered by go-framework packages.

Compilation stages:

  1. Validate — structural correctness (IDs, references, module naming conventions)
  2. Resolve — topological sort of modules by dependencies (Kahn's algorithm)
  3. Enrich — promote module-level entities/events/actions/workflows to document level; warn on unresolved provides

runtime — Flow execution engine

Executes FlowDSL workflows. Nodes are dispatched to registered NodeHandler functions by node ID, node name, or kind. Execution state can be persisted via a Checkpoint implementation for resumable long-running runs.

NodeHandler

go
type NodeHandler func(ctx context.Context, step *Step) error

This is a plain function type — not an interface.

Step

Per-node execution context passed to every NodeHandler:

go
type Step struct {
    Node        *ir.Node       // IR definition of the current node
    Workflow    *ir.Workflow   // IR definition of the owning workflow
    ExecutionID string         // unique run ID
    Attempt     int            // 1-based attempt count (1 on first try)
    Input       map[string]any // data from the previous step or trigger
    Output      map[string]any // written by the handler; passed to the next step
}

Node I/O binding ($inputMap / $outputMap)

The runtime is an accumulating-packet model: each node's Input is the previous node's whole Output, and edges route but never transform. To reuse a generic node in a flow whose packet uses different field names — without editing the node — set two reserved keys in the node's Config:

yaml
config:
  $inputMap:  { address: "customer.shipping_address" }   # node field ← packet path
  $outputMap: { Valid: "shipping_valid" }                # node output → flow key

Before the handler runs, $inputMap pre-populates step.Input["address"] from the dotted path (additive — the passthrough packet is untouched). After it runs, $outputMap copies step.Output["Valid"] to step.Output["shipping_valid"] (non-destructive). Both are no-ops when absent, so existing flows are unaffected; there is no ir.Node struct change. Applied in both Start and Resume. In Flow Studio these are edited via the Inspector's I/O Binding section.

RetryPolicy

go
type RetryPolicy struct {
    MaxAttempts   int           // 1 = no retries (default)
    InitialDelay  time.Duration
    BackoffFactor float64       // 2.0 = exponential doubling
    MaxDelay      time.Duration // caps delay between retries
}

DefaultRetryPolicy is RetryPolicy{MaxAttempts: 1}.

Checkpoint interface

go
type Checkpoint interface {
    Save(ctx context.Context, exec *ExecutionRecord) error
    Load(ctx context.Context, executionID string) (*ExecutionRecord, error)
    SaveStep(ctx context.Context, step *StepRecord) error
}

MemoryCheckpoint is the built-in in-process implementation used by default when EngineConfig.Checkpoint is nil.

ExecutionRecord and StepRecord

ExecutionRecord is the top-level persisted record for a workflow run. It contains an ID, WorkflowID, Status, Input, Output, Error, Steps []*StepRecord, and timestamps.

StepRecord records one node's execution: ExecutionID, NodeID, NodeName, Status, Attempt, Input, Output, Error, StartedAt, CompletedAt.

Status values: pending, running, completed, failed, retrying, skipped.

Engine API

go
import "github.com/redelay/go-flowdsl/runtime"

// Create engine — MemoryCheckpoint used when no Checkpoint is supplied.
eng := runtime.NewEngine(runtime.EngineConfig{})

// Register a handler for a specific node ID or name.
eng.RegisterHandler("validate_payment", func(ctx context.Context, step *runtime.Step) error {
    step.Output["valid"] = true
    return nil
})

// Register a fallback handler for all nodes of a given kind.
eng.RegisterKindHandler(ir.NodeKindAction, fallbackHandler)

// Set a per-node retry policy (keyed by node ID or name).
eng.SetRetryPolicy("charge_card", runtime.RetryPolicy{
    MaxAttempts:   3,
    InitialDelay:  2 * time.Second,
    BackoffFactor: 2.0,
    MaxDelay:      30 * time.Second,
})

// Start a workflow. Runs synchronously — call in a goroutine for non-blocking behaviour.
exec, err := eng.Start(ctx, wf, map[string]any{"user_id": "abc"})
// exec.Status — completed or failed
// exec.Output — output from the last node

// Resume a previously started execution, skipping already-completed steps.
exec, err = eng.Resume(ctx, wf, exec.ID)

Handler resolution order for a node: node ID → node name → node ActionRef → kind fallback. An on-error edge (edge with condition: "error" or "on_error") routes execution to an error-handling node instead of aborting the workflow.

StepObserver — real-time step lifecycle hooks

A StepObserver receives step events inline with handler execution, not in a post-run batch. Observers are the hook flowexec uses to emit node.started / node.done events with true wall-clock timestamps while the workflow is still running.

go
type StepEventKind string

const (
    StepEventStarted StepEventKind = "started"
    StepEventDone    StepEventKind = "done"
    StepEventFailed  StepEventKind = "failed"
    StepEventRetried StepEventKind = "retried"
)

type StepEvent struct {
    Kind     StepEventKind
    Step     *Step       // same pointer the handler saw
    Attempt  int         // 1-based
    Err      error       // non-nil on Failed / Retried
    Duration time.Duration
}

type StepObserver func(ev StepEvent)

Two registration paths, per-run wins over engine-wide:

go
// Engine-wide observer — fires for every run unless overridden.
eng.SetStepObserver(func(ev runtime.StepEvent) { ... })

// Per-run observer via context. Concurrent runs do not interfere —
// each carries its own observer through context.
ctx = runtime.WithStepObserver(ctx, myObserver)
exec, _ := eng.Start(ctx, wf, input)

Firing order is Started → (Retried × N) → Done|Failed per step. Observers run synchronously on the executor goroutine — keep them fast (send to a channel, append to a ring buffer, etc.). Panics inside an observer are recovered but surface as run.failed with the panic message.

Tests: go-flowdsl/runtime/observer_test.go.

spec — Studio spec document converters

Bidirectional converters between the canonical ir.Workflow (flat nodes + edges used at runtime) and the Studio spec document format (nested flows, components, visual metadata).

go
import "github.com/redelay/go-flowdsl/spec"

The ir.Workflow is the runtime representation stored by flowexec. The spec Document is the editing representation used by the FlowDSL Studio frontend — it carries flows, components.nodes, and x-ui visual metadata (positions, icons, colors).

Types

TypeDescription
DocumentTop-level spec document — flowdsl version, info, flows, components
InfoDocument metadata — title, version, description
FlowSingle flow — nodes (map of $ref pointers), edges
NodeRef$ref pointer into #/components/nodes/...
EdgeDirected connection — from, to, when, packet, delivery
DeliveryPolicyEdge delivery mode — mode (direct, ephemeral, durable, stream, checkpoint)
ComponentsShared definitions — nodes, packets, policies
NodeDefNode component — operationId, title, kind, settings, x-ui
XUICanvas visual metadata — position, icon, color, registryId, group
XYPosCanvas coordinates — x, y

FromWorkflow — IR to spec

go
doc := spec.FromWorkflow(wf) // *ir.Workflow → *spec.Document
  • Each ir.Node becomes an entry in components.nodes and a $ref pointer in the flow.
  • Node kinds are mapped: start→source, end→terminal, gateway→router, condition→router, wait→checkpoint, event→source, action→action.
  • Edge conditions become when, labels become description.
  • Returns nil for nil input.

ToWorkflow — spec to IR

go
wf := spec.ToWorkflow(doc, "main") // *spec.Document, flowID → *ir.Workflow
  • Resolves $ref pointers to component node definitions.
  • Spec kinds are mapped back: source→start, terminal→end, router→gateway, transform/llm/checkpoint/publish/integration/subworkflow→action.
  • If flowID is empty, uses the first flow in the document.
  • Returns nil for nil input.

Round-trip guarantee

FromWorkflow(wf) → ToWorkflow(doc, flowID) preserves all structural information (node IDs, names, kinds, edges, conditions). Visual metadata (x-ui) is generated on FromWorkflow with default values and round-trips through ToWorkflow without loss.

flowexec — orchestration, storage, sinks {#flowexec}

flowexec is a subpackage of go-flowdsl that sits on top of runtime.Engine and adds the three things real applications need but the raw engine does not provide:

  1. An executor that wraps Engine.Start and emits lifecycle events for every run, node, and edge.
  2. An immutable flow store — flow documents saved as append-only versions identified by a content hash, with head/published pointers.
  3. A pluggable sink abstraction for persisting events and run summaries (Mongo time-series, ClickHouse, no-op).

flowexec has no dependency on go-framework; it builds on runtime + ir plus the Go MongoDB driver (only for the mongostore and mongots sink implementations — pick memstore

  • noop for zero external deps).
text
go-flowdsl/
  flowexec/
    executor.go        Executor, Options, RunInput, RunAsync
    sink/              FlowEventSink, UsageSink, Reader, Combined, FlowEvent, FlowRun
    sink/noop/         No-op sink for tests and Mongo-less deployments
    sink/mongots/      Mongo time-series sink (flow_events + flow_runs)
    store/             Store interface + Flow / FlowVersion / FlowDeployment types
    store/memstore/    In-memory store for tests
    store/mongostore/  MongoDB-backed store with versionHash indexes
    module/            HTTP module — flow + deployment CRUD routes

FlowDeployment — the caller-facing routing table

Flow + FlowVersion give you immutable diagrams. FlowDeployment sits above both as the business-level binding every consumer module references:

text
Module config → DeploymentID → FlowDeployment.variants:
                                  - stable       → flow.X (or pinned ver.Y)   90%
                                  - canary       → flow.Z (or pinned ver.W)   10%
                                  - experiment_a → flow.Q                      0%  (staged)

One deployment, many flows. Weights are session-sticky (FNV-1a hash of the configured StickyBy key bucketed into a CDF), so A/B tests between radically different flow shapes live in a single config knob without flapping users between variants on reloads.

go
import flowstore "github.com/redelay/go-flowdsl/flowexec/store"

// Once at boot: make sure the deployment exists. Safe to call every
// startup — existing deployments are returned untouched.
dep, _ := flowstore.EnsureDeployment(ctx, s, flowstore.CreateDeploymentInput{
    ID:   "assistant",
    Name: "Production Assistant",
    Variants: []flowstore.FlowDeploymentVariant{
        {Label: "stable", FlowID: "assistant-v1", Weight: 100},
    },
    StickyBy: flowstore.StickySession,
})

// On every incoming request: resolve which flow + version to run for
// this session. Same (deployment, session_id) always → same variant.
resolved, err := flowstore.ResolveDeployment(ctx, s, "assistant", map[string]any{
    "session_id": sessionID,
}, nil)
// resolved.Label == "stable", resolved.FlowID, resolved.Version ready for the executor.

Admin HTTP surface (registered by go-flowdsl/flowexec/module, admin-gated via admin:access):

MethodPathPurpose
GET/deploymentsList (cursor pagination)
POST/deploymentsCreate
GET/deployments/{id}Read
PATCH/deployments/{id}Replace variants / stickyBy / labels
DELETE/deployments/{id}Delete (flows untouched)
POST/deployments/{id}/variants/from-templateProvision a new flow from a template, publish v1, append as variant — all in one call

Consumer modules (assistant, community, any future module) stay thin: they hold a DeploymentID in config and call ResolveDeployment per request. Adding a canary flow is a curl against /deployments/{id}/variants/from-template; no module redeploy needed.

Executor

go
import "github.com/redelay/go-flowdsl/flowexec"

ex := flowexec.New(flowexec.Options{
    Engine: eng,     // a *runtime.Engine with handlers registered
    Sink:   sink,    // a sink.FlowEventSink — mongots, noop, or a custom decorator
})

exec, runID, err := ex.Run(ctx, flowexec.RunInput{
    Workflow:    wf,
    VersionID:   version.ID,
    VersionHash: version.VersionHash,
    Input:       map[string]any{"user": "u1"},
    Labels:      map[string]string{"source": "api"},
})

// Fire and forget; subscribe to SSE for progress.
runID, errCh, err := ex.RunAsync(ctx, flowexec.RunInput{Workflow: wf})

RunInput.RunID can be pre-populated by the caller; if non-empty it is used verbatim instead of a fresh UUID. RunAsync always pre-generates the ID and stamps it back onto the input, so subscribers can attach to the event stream before the run begins emitting events.

Every run produces:

  • One FlowRun row written at start (Status: "running") and updated at end (completed or failed).
  • Kind-run.started → per-node node.started / node.done / node.failed → kind-run.completed or run.failed events. Node events are derived from exec.Steps after Engine.Start returns (the engine lacks hooks today; node events will move to real-time emission once hooks land).

sink.FlowEventSink interface

go
type FlowEventSink interface {
    RecordFlowEvent(ctx context.Context, ev FlowEvent) error
    RecordFlowRun(ctx context.Context, run FlowRun) error
}

Event kinds:

KindWhen
run.startedImmediately after the executor accepts a run
run.completed / run.failedAfter Engine.Start returns
node.started / node.done / node.failedOne pair per step in exec.Steps
packet.emitReserved for future packet-emission events
edge.deliverReserved for future edge-delivery events

FlowEvent carries RunID, FlowID, VersionHash, NodeID, Timestamp, Seq, DurationMs, ErrorMsg, Payload, and a Meta block (FlowID, RunID, AssistantID, UserID). Every event is stamped with VersionHash so past runs remain replay-correct even after the flow is edited and re-published.

Strict emission order via Seq. Seq is a per-run monotonically-increasing counter assigned by the Executor in the exact order events are produced. It serves as the tie-breaker whenever multiple events share a millisecond timestamp — which happens for pass-through entry/exit nodes whose StartedAt and CompletedAt resolve to the same wall-clock ms. mongots.ListRunEvents sorts by (ts asc, seq asc) so consumers see a strict total order; clients (Studio, custom dashboards) MUST treat that order as canonical and not re-sort.

Available implementations:

PackagePurpose
sink/noopDrops everything. Default when Mongo is not wired.
sink/mongotsWrites flow_events as a time-series collection and flow_runs as a regular collection. Implements Reader for historical queries.

store.Store — immutable versioning

go
type Store interface {
    CreateFlow(ctx, CreateFlowInput) (*Flow, error)
    GetFlow(ctx, id) (*Flow, error)
    UpdateFlow(ctx, id, name, description, labels) (*Flow, error)
    DeleteFlow(ctx, id) error
    ListFlows(ctx, ListFlowsOptions) ([]*Flow, nextCursor, error)

    SaveVersion(ctx, SaveVersionInput) (*FlowVersion, error)
    GetVersion(ctx, id) (*FlowVersion, error)
    ListVersions(ctx, flowID, ListVersionsOptions) ([]*FlowVersion, nextCursor, error)

    PublishVersion(ctx, flowID, versionID) (*Flow, error)
    GetPublishedVersion(ctx, flowID) (*FlowVersion, error)
}

Two collections, two invariants:

  • Flow is mutable. It holds Name, Description, Labels, and two pointers: HeadVersionID (newest saved) and PublishedVersionID (what runs resolve to).
  • FlowVersion is append-only. It holds the full ir.Workflow document plus VersionHash = sha256(workflow JSON). Once saved, versions never change.

SaveVersion hashes the incoming document, inserts a new flow_versions row, and atomically advances Flow.HeadVersionID. PublishVersion is a pointer flip — no document rewrites, no cascading updates. Past runs reference versions by ID and hash, so replay is deterministic.

BackendDependenciesUse for
memstorestdlib onlytests, ephemeral demos
mongostoremongo-driverproduction

mongostore.EnsureIndexes creates indexes on flow_versions.flowID+createdAt (desc), flow_versions.versionHash, and flows.updatedAt (desc).

Zero-dependency deployments

For workers that should run flows without a database, pair memstore + sinknoop:

go
st := memstore.New()
snk := sinknoop.New()
ex := flowexec.New(flowexec.Options{Engine: eng, Sink: snk})

This keeps flowexec usable in serverless or test scenarios where MongoDB is not available.

For HTTP endpoints, SSE event streaming, and framework integration, see the flowexec addon module.

metrics — live rollups, SLOs, recommendations {#metrics}

flowexec/metrics is the aggregation layer that powers Flow Studio's Live mode. The runtime produces a Frame per flow per tick (default 1 s); subscribers receive the current frame over SSE and an optional archival sink persists frames to ClickHouse.

go
type Frame struct {
    Ts          time.Time
    WindowMs    int64
    Meta        Meta                          // tenantId, flowId, versionHash, assistantId, userId
    Nodes       map[string]NodeMetrics        // per-node rollup, keyed by nodeId
    Edges       map[string]EdgeMetrics        // per-edge rollup, keyed by "${from}__${to}"
    Channels    map[string]ChannelMetrics     // per-transport backend (kafka lag, redis depth, …)
    Suggestions []Suggestion                  // server-evaluated rule output
}

NodeMetrics carries two cadences side-by-side: Completed / Failed reset every frame (drives rate + error-rate), while CompletedTotal / FailedTotal are monotonically-increasing counters since aggregator boot. The Studio canvas chips show both — the rate first, then a Σ total so users can tell at a glance how many packets a node has handled since Live mode started. EdgeMetrics follows the same pattern with Delivered (window) and DeliveredTotal (cumulative).

Identity rides on Meta, never on routing metadata. Topic names, collection names, and ClickHouse table names stay content-agnostic; multi-tenant deployments filter on meta.tenantId on the read side. Future Kafka/NATS exporters carry Meta as message headers.

go
type MetricsSink interface {
    RecordFrame(ctx context.Context, frame Frame) error
    ListFrames(ctx context.Context, flowID string, since time.Time, limit int) ([]Frame, error)
    Subscribe(flowID string, buf int) (<-chan Frame, func())
}

Implementations:

PackagePurpose
metrics/inprocAlways-on aggregator. Per-flow mutex on the hot path, 5-minute ring buffer, non-blocking fan-out to SSE subscribers. Default tick 1 s, default ring size 300 frames.
metrics/clickhouseBatched archival sink. Frames fan out to flat Rows (one per node/edge/channel); Driver interface isolates the network call so drivers ship separately. Safe to wire with empty URL (returns nil sink).
metrics/rulesServer-side recommendation engine. Stateless Apply(frame, workflow) []Suggestion. Seven baseline rules covering capacity, reliability, SLOs, transport upgrades, backpressure.

Executor integration. flowexec.Options gained a Metrics MetricsRecorder field — any sink that implements RecordEvent(ev), RegisterEdges(flowID, meta, keys), and RegisterEdgeMode(flowID, key, mode). When set, every FlowEvent the Executor emits is also handed to the aggregator (same emit closure, so Seq ordering is preserved). Before a run fires, RegisterEdges pre-seeds the flow's edge keys so empty edges render at 0/s and the rule engine has each edge's DeliveryMode to reason about.

go
agg := inproc.New(inproc.Options{TickInterval: time.Second, Downstream: chSink})
agg.Start(ctx)
ex := flowexec.New(flowexec.Options{Engine: eng, Sink: bcast, Metrics: agg})

SLOs as first-class IR fields. ir.Node.SLO declares per-node service-level objectives that the rules package enforces:

go
type NodeSLO struct {
    P95LatencyMs  int64   // hard ceiling on p95 latency (ms)
    MaxErrorRate  float64 // upper bound on failures / total (0..1)
    MaxSaturation float64 // upper bound on (inflight / capacity)
}

Zero on any field means "no SLO on that dimension". Breaches fire a slo-kind Suggestion with severity critical.

See Trace + Live Observability for the user-facing walkthrough, wire format, and canvas rendering.

Cross-container live-metrics fan-out

In multi-container deployments, any api container can run a flow but only the admin-api container serves the Studio Live SSE. The two are bridged through the Coordination layer:

ComponentFileRole
coordFrameSinkflowexec/module/coord_forward.goWraps inproc.Sink; writes every flushed metrics.Frame to the coord stream live:<flowID> via Coordinator.AppendStream.
frameHasSignalsame fileProducer-side filter. Drops frames with no node activity, no edge traffic, and no versionHash so idle aggregators in containers not running the flow don't pollute the stream.
handleFlowLiveSSEflowexec/module/admin/admin.goSSE handler on GET /flows/{id}/live. Reads from Coordinator.ReadStream(ctx, "live:"+flowID, StreamEnd) when a Coordinator is wired; falls back to the local in-proc aggregator otherwise.

The module is safe on a single-container setup: with no Coordinator the sink is a pass-through and the SSE handler reads directly from its own aggregator.

Instant first paint. handleFlowLiveSSE now emits a topology-only seed frame the moment the SSE connection opens — zeroed counters for every node and edge pulled straight from the published workflow. The Studio canvas therefore renders every badge immediately instead of waiting up to one flush tick plus coord propagation (~1-3 s). Subsequent real frames overwrite the seed.

Staleness watchdog. Each live connection runs a forwarder goroutine (runCoordForwarder) that re-opens the coord reader after coordIdleTimeout seconds of silence (default 45 s). This is the defense against NATS OrderedConsumer stalls, transient broker hiccups, or any issue that silently stops delivering messages without closing the subscription. The HTTP response stream stays up the whole time — only the coord reader is swapped — so the Studio widget never sees a disconnect.

Per-producer attribution. metrics.Meta.InstanceID is stamped onto every emitted frame with the value of HOSTNAME (Docker sets this to the container id) or os.Hostname as a fallback. Admin-api can use this to tell "how many containers are reporting for this flow" and is the hook for the planned multi-container aggregator that will sum window counters across producers per tick.

Multi-container aggregation (planned)

Today the admin-api SSE handler forwards frames from every container as-is. In a deployment where multiple api containers run the same flow, the Studio canvas sees frames from each container in turn — each carries only that container's window counters. A canvas that wants the global rate must sum the window counters from all containers for the same tick.

Planned design:

  1. Admin-api maintains a small in-memory map (flowID, tickSecond) → aggregatedFrame.
  2. Incoming coord frames are matched to the map by rounded timestamp (1 s granularity) and added: window counters sum, totals sum, inflight sums, latency percentiles merge via weighted mean (approximate — exact quantile merging needs t-digest or similar).
  3. When a frame from every known InstanceID has arrived for a tick (or a short timeout expires), the aggregate is flushed to the SSE client.
  4. InstanceID list is discovered from the frames themselves — first-seen per (flowID, 5-minute window) — so scale-out happens automatically.

Groundwork in this release: InstanceID on metrics.Meta. Aggregator layer itself is not yet implemented — follow-up work.

Module HTTP surface — public vs admin

flowexec/module/module.go registers the executor, metrics aggregator, template provider, and usage sink — but no HTTP routes. Every flow / run / deployment endpoint lives in flowexec/module/admin (id: flowexec-admin, depends on flowexec). Blank-import the admin companion only in admin-api:

go
// cmd/admin-api admin_imports.go
import _ "github.com/redelay/go-flowdsl/flowexec/module/admin"

This means cmd/api can run flows (the engine + sinks are registered by the core module) but its public HTTP surface carries zero /flows, /runs, or /deployments paths. The regression is locked in by TestOpenAPIHasNoFlowAdminPaths.

All admin routes are gated on auth.RequireActiveUser + auth.RequirePermission("admin:access") and sit under modules.AdminPrefix() (default /admin, override via ADMIN_ROUTE_PREFIX).

Full route table — flows + versions + runs + live + deployments, all in the admin companion:

MethodPathPurpose
GET/flowsList flows (cursor pagination)
POST/flowsCreate flow
GET/flows/{id}Read flow
PATCH/flows/{id}Update name / description / labels
DELETE/flows/{id}Delete flow
GET/flows/{id}/versionsList versions
POST/flows/{id}/versionsSave a new version (supports ?format=spec for Studio-native spec.Document bodies)
GET/flows/{id}/versions/{vid}Read a version
POST/flows/{id}/publishPublish a version or set weighted variants
POST/flows/{id}/rollbackRollback to a prior published version
GET/flows/{id}/routeRead the currently-published route / variants
GET/flows/{id}/publish-historyAudit log of publish actions
GET/flows/templatesList embedded + module-contributed templates
GET/flows/templates/*Read one template (id may contain slashes)
POST/flows/{id}/runsLaunch a run
GET/flows/{id}/runsList runs for the flow (cursor pagination)
GET/runs/{runID}Read one run
GET/runs/{runID}/eventsSSE stream of run events
GET/runs/{runID}/events/historyHistorical event replay
GET/flows/{id}/liveSSE stream of live metrics.Frame rollups
GET/flows/{id}/metrics/historyHistorical frames
GET/deploymentsList deployments
POST/deploymentsCreate deployment
GET/deployments/{id}Read deployment
PATCH/deployments/{id}Replace variants / stickyBy / labels
DELETE/deployments/{id}Delete deployment (flows untouched)
POST/deployments/{id}/variants/from-templateProvision + publish + bind in one call

Auto-generation of subworkflow nodes {#auto-generation-of-subworkflow-nodes}

The go-framework module registry automatically generates subworkflow FlowDSLNode entries for any module that provides workflows via FlowDSLProvider. This means modules only need to implement FlowDSLFragments() — the corresponding node definitions are created at runtime.

The auto-generation runs inside Registry.IRModules() after collecting both workflows and explicit nodes. For each workflow, it checks whether a FlowDSLNode with a matching WorkflowRef already exists. If not, it calls ir.SubworkflowNodeFromWorkflow() to generate one automatically.

This ensures:

  • Every module workflow is visible in /modules.json and the module browser UI
  • Modules can override auto-generated nodes by providing an explicit FlowDSLNode with the same WorkflowRef
  • The module browser shows step count and workflow reference for subworkflow nodes

FlowDSL node subpackage pattern {#flowdsl-node-subpackages}

Modules that export FlowDSL nodes use a flowdsl/ subdirectory that registers as a separate companion module. This keeps the FlowDSL node YAML definitions and registration code isolated from the core module logic.

Each subpackage has four files:

FilePurpose
module.yamlNode definitions (embedded via go:embed)
ir.go//go:embed module.yaml + var moduleYAML []byte
module.goinit() factory + Module struct embedding modules.IRBase
module_test.goFactory registration, manifest, node catalog, kind checks

Opt-in via blank import alongside the parent module:

go
import (
    _ "github.com/redelay/go-framework/modules/auth"         // core module
    _ "github.com/redelay/go-framework/modules/auth/flowdsl"  // FlowDSL nodes
)

Node catalog by repository:

go-framework:

SubpackageFactoryNodesKinds
modules/auth/flowdslauth-flowdsl71 action, 1 transform, 1 router, 4 source
modules/users/flowdslusers-flowdsl83 action, 2 transform, 3 source
modules/groups/flowdslgroups-flowdsl73 action, 1 transform, 3 source

go-modules:

SubpackageFactoryNodesKinds
scheduler/flowdslscheduler-flowdsl42 source (with settings_schema), 2 action
settings/flowdslsettings-flowdsl31 transform (with NotFound), 2 action
storage/flowdslstorage-flowdsl42 action, 2 transform (1 NotFound, 1 settings_schema)
notifications/flowdslnotifications-flowdsl32 action, 1 source
verification/flowdslverification-flowdsl21 action, 1 router (Valid/Invalid)
links/flowdsllinks-flowdsl32 action, 1 transform (with NotFound)

go-module-email:

SubpackageFactoryNodesKinds
flowdslemail-flowdsl42 action, 1 transform, 1 source

FlowDSL SDK dependency

go-flowdsl depends on the FlowDSL Go SDK (github.com/flowdsl/flowdsl-go) for canonical node manifest types. The SDK provides zero-dependency types in pkg/node (manifest, ports, settings schema) that both ir and nodegen packages use for manifest generation and conversion.

text
github.com/flowdsl/flowdsl-go    ← zero-dep canonical types
       ↑
github.com/redelay/go-flowdsl    ← IR layer + code generation
       ↑
github.com/redelay/go-framework  ← module registry + modspec browser

cmd/flowdsl — CLI

The standalone flowdsl CLI provides portable toolchain commands with no framework dependencies.

Install:

shell
go install github.com/redelay/go-flowdsl/cmd/flowdsl@latest

Commands

CommandDescription
flowdsl validate <file>Validate a FlowDSL YAML file; exit 1 on errors
flowdsl gen <module.yaml> [-o outdir]Scaffold a Go module from IR module YAML
flowdsl svg <module.yaml> [-o out.svg]Generate SVG diagram from IR module YAML
flowdsl compile <file> [-f format] [--out json|yaml]Compile to enriched IR document

validate

shell
flowdsl validate app.flow.yaml

Parses the file and runs flowdsl.Validate. Prints all diagnostics with [E]/[W]/[I] prefixes. On success:

text
ok — no issues found

On failure (exit code 1):

text
[E] FDL010: event "order_event" has no entity_type
[W] FDL002: version is empty

gen

shell
flowdsl gen modules/orders.module.yaml -o ./modules/orders/

Loads the ir.Module YAML, validates it with modval, then scaffolds Go source files via modgen.Scaffold. Prints each written file path. Existing files are not overwritten.

svg

shell
flowdsl svg modules/auth.module.yaml -o auth.svg

Generates an SVG diagram and writes it to the specified file. Omit -o to write to stdout.

compile

shell
flowdsl compile app.flow.yaml                        # FlowDSL → IR as JSON (default)
flowdsl compile app.flow.yaml --out yaml             # FlowDSL → IR as YAML
flowdsl compile model.yaml -f businessmodel          # businessmodel → IR as JSON

Runs the full Import + Compile pipeline and writes the enriched ir.Document to stdout. Diagnostics go to stderr. Exit code 1 if compilation errors occur.

FlagDefaultDescription
-f, --formatflowdslInput format: flowdsl or businessmodel
--outjsonOutput format: json or yaml

For the heavier CLI with asyncapi/openapi import, scaffold, audit, diff, and convert commands, see the redelayctl reference.