Go FlowDSL
Module: github.com/redelay/go-flowdsl
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:
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
go get github.com/redelay/go-flowdsl
go.mod
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.
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:
| Type | Description |
|---|---|
Document | Root container for an entire IR |
Module | Business/domain unit with entities, events, actions, workflows, config, settings, CRUD, consumers |
Entity | Business entity or domain aggregate with typed fields |
Field | Single field — type, constraints, nested properties, $ref |
FieldType | Enum: string, integer, float, boolean, array, object, date, datetime, uuid, binary, ref, enum, any |
Event | Business event — entity_type + action, optional topic and payload |
Command | State-changing intent with input/output schemas and emitted events |
Query | Read operation with input/output schemas |
Action | Top-level operation — kind (command/query/mutation/subscription), method, path, I/O, emits |
Workflow | Multi-step process with Node graph and Edge connections |
Node | Workflow step — kind (action/event/gateway/start/end/wait/condition/parallel) |
Edge | Directed connection between nodes with optional condition |
Channel | Message channel with publish/subscribe events, protocol, bindings |
Schema | Reusable data structure definition |
Packet | Data transfer structure (event payload, message body) |
Schedule | Time-based trigger — cron expression or interval |
Migration | Database migration — backend, version, up/down SQL |
Consumer | Event consumer — topic, group ID, handler reference, batch size |
Provenance | Origin tracking — source, URI, source ref, mappings |
ConfigDefinition | Module runtime configuration (entries with env vars, types, defaults) |
SettingsDefinition | Admin-UI-editable settings (groups of fields with components) |
CRUDDefinition | CRUD capabilities for an entity (create/read/update/delete/list) |
FlowDSLNode | Module-exported FlowDSL node — kind, ports, settings, refs to events/actions/workflows |
FlowDSLPort | Named 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.
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:
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:
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):
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.
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.
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.
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.
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.
import "github.com/redelay/go-flowdsl/resolver"
sorted, diag := resolver.Resolve(doc.Modules)
// sorted — []*ir.Module in dependency order (dependencies first)
// diag — *diagnostics.Collector
| Code | Severity | Rule |
|---|---|---|
RES001 | error | Duplicate module ID |
RES002 | error | Depends on unknown module |
RES003 | error | Circular dependency detected |
RES004 | error | Module is part of a dependency cycle |
validate — Structural validator
Checks an IR Document for structural correctness and referential integrity.
import "github.com/redelay/go-flowdsl/validate"
diag := validate.Validate(doc) // *diagnostics.Collector
| Code | Severity | Rule |
|---|---|---|
VAL001 | error | Document is nil |
VAL002 | error | Document ID is required |
VAL003 | warning | Document version is empty |
VAL010 | error | Duplicate ID |
VAL020 | error | Module depends on unknown module |
VAL021 | warning | Event references unknown payload |
VAL022 | warning | Action references unknown input |
VAL023 | warning | Action references unknown output |
VAL030 | error | Workflow edge references unknown source node |
VAL031 | error | Workflow edge references unknown target node |
VAL040 | error | Module has no name |
VAL041 | warning | Module ID should start with module. |
flowdsl — FlowDSL YAML parser
Parses, validates, normalizes, and exports the FlowDSL YAML application definition format.
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:
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:
| Code | Severity | Rule |
|---|---|---|
FDL001 | error | Document is nil |
FDL002 | warning | Version is empty |
FDL010 | error | Event missing entity_type |
FDL011 | error | Event missing action |
FDL012 | warning | Event payload field has invalid type |
FDL020 | error | Workflow step missing ID |
FDL021 | warning | Step next references unknown step |
FDL022 | warning | Step on_error references unknown step |
FDL023 | warning | Step action references unknown action |
FDL024 | warning | Step event references unknown event |
FDL025 | warning | Step kind is not a recognised value |
FDL030 | warning | Module depends_on unknown module |
FDL040 | warning | Schedule action references unknown action |
FDL050 | warning | Channel protocol not recognised |
FDL060 | warning | Action field has invalid type |
FDL061 | warning | Action kind not recognised |
FDL062 | warning | Action emits references unknown event |
modfile — Module file I/O
Loads, saves, and exports ir.Module and ir.Document values as YAML or JSON.
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.
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.
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:
| File | Contents |
|---|---|
manifests/<slug>.flowdsl-node.json | Node manifest for each event, action, command, and workflow |
nodes_generated.go | Go handler stubs implementing node.NodeHandler interface |
nodes_registry.go | FlowDSLNodes() function returning []*ir.FlowDSLNode for module wiring |
Node generation rules:
| IR element | Node kind | Notes |
|---|---|---|
Event | source | One output port named after the event |
Action | action or transform | Queries become transform; mutations become action |
Command | action | Input/output ports |
Workflow | subworkflow | Ports 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.
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.
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.
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.
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:
| Constant | Value | Import | Export |
|---|---|---|---|
compile.FormatFlowDSL | "flowdsl" | yes | yes |
compile.FormatBusinessModel | "businessmodel" | yes | — |
Additional formats (asyncapi, openapi) are registered by go-framework packages.
Compilation stages:
- Validate — structural correctness (IDs, references, module naming conventions)
- Resolve — topological sort of modules by dependencies (Kahn's algorithm)
- 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
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:
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:
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
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
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
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.
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:
// 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).
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
| Type | Description |
|---|---|
Document | Top-level spec document — flowdsl version, info, flows, components |
Info | Document metadata — title, version, description |
Flow | Single flow — nodes (map of $ref pointers), edges |
NodeRef | $ref pointer into #/components/nodes/... |
Edge | Directed connection — from, to, when, packet, delivery |
DeliveryPolicy | Edge delivery mode — mode (direct, ephemeral, durable, stream, checkpoint) |
Components | Shared definitions — nodes, packets, policies |
NodeDef | Node component — operationId, title, kind, settings, x-ui |
XUI | Canvas visual metadata — position, icon, color, registryId, group |
XYPos | Canvas coordinates — x, y |
FromWorkflow — IR to spec
doc := spec.FromWorkflow(wf) // *ir.Workflow → *spec.Document
- Each
ir.Nodebecomes an entry incomponents.nodesand a$refpointer 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 becomedescription. - Returns
nilfornilinput.
ToWorkflow — spec to IR
wf := spec.ToWorkflow(doc, "main") // *spec.Document, flowID → *ir.Workflow
- Resolves
$refpointers to component node definitions. - Spec kinds are mapped back:
source→start,terminal→end,router→gateway,transform/llm/checkpoint/publish/integration/subworkflow→action. - If
flowIDis empty, uses the first flow in the document. - Returns
nilfornilinput.
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:
- An executor that wraps
Engine.Startand emits lifecycle events for every run, node, and edge. - An immutable flow store — flow documents saved as append-only versions identified by a content hash, with head/published pointers.
- 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
noopfor zero external deps).
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:
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.
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):
| Method | Path | Purpose |
|---|---|---|
| GET | /deployments | List (cursor pagination) |
| POST | /deployments | Create |
| GET | /deployments/{id} | Read |
| PATCH | /deployments/{id} | Replace variants / stickyBy / labels |
| DELETE | /deployments/{id} | Delete (flows untouched) |
| POST | /deployments/{id}/variants/from-template | Provision 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
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
FlowRunrow written at start (Status: "running") and updated at end (completedorfailed). - Kind-
run.started→ per-nodenode.started/node.done/node.failed→ kind-run.completedorrun.failedevents. Node events are derived fromexec.StepsafterEngine.Startreturns (the engine lacks hooks today; node events will move to real-time emission once hooks land).
sink.FlowEventSink interface
type FlowEventSink interface {
RecordFlowEvent(ctx context.Context, ev FlowEvent) error
RecordFlowRun(ctx context.Context, run FlowRun) error
}
Event kinds:
| Kind | When |
|---|---|
run.started | Immediately after the executor accepts a run |
run.completed / run.failed | After Engine.Start returns |
node.started / node.done / node.failed | One pair per step in exec.Steps |
packet.emit | Reserved for future packet-emission events |
edge.deliver | Reserved 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:
| Package | Purpose |
|---|---|
sink/noop | Drops everything. Default when Mongo is not wired. |
sink/mongots | Writes flow_events as a time-series collection and flow_runs as a regular collection. Implements Reader for historical queries. |
store.Store — immutable versioning
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:
Flowis mutable. It holdsName,Description,Labels, and two pointers:HeadVersionID(newest saved) andPublishedVersionID(what runs resolve to).FlowVersionis append-only. It holds the fullir.Workflowdocument plusVersionHash = 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.
| Backend | Dependencies | Use for |
|---|---|---|
memstore | stdlib only | tests, ephemeral demos |
mongostore | mongo-driver | production |
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:
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.
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.
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:
| Package | Purpose |
|---|---|
metrics/inproc | Always-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/clickhouse | Batched 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/rules | Server-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.
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:
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:
| Component | File | Role |
|---|---|---|
coordFrameSink | flowexec/module/coord_forward.go | Wraps inproc.Sink; writes every flushed metrics.Frame to the coord stream live:<flowID> via Coordinator.AppendStream. |
frameHasSignal | same file | Producer-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. |
handleFlowLiveSSE | flowexec/module/admin/admin.go | SSE 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:
- Admin-api maintains a small in-memory map
(flowID, tickSecond) → aggregatedFrame. - 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).
- When a frame from every known
InstanceIDhas arrived for a tick (or a short timeout expires), the aggregate is flushed to the SSE client. - 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:
// 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:
| Method | Path | Purpose |
|---|---|---|
| GET | /flows | List flows (cursor pagination) |
| POST | /flows | Create flow |
| GET | /flows/{id} | Read flow |
| PATCH | /flows/{id} | Update name / description / labels |
| DELETE | /flows/{id} | Delete flow |
| GET | /flows/{id}/versions | List versions |
| POST | /flows/{id}/versions | Save a new version (supports ?format=spec for Studio-native spec.Document bodies) |
| GET | /flows/{id}/versions/{vid} | Read a version |
| POST | /flows/{id}/publish | Publish a version or set weighted variants |
| POST | /flows/{id}/rollback | Rollback to a prior published version |
| GET | /flows/{id}/route | Read the currently-published route / variants |
| GET | /flows/{id}/publish-history | Audit log of publish actions |
| GET | /flows/templates | List embedded + module-contributed templates |
| GET | /flows/templates/* | Read one template (id may contain slashes) |
| POST | /flows/{id}/runs | Launch a run |
| GET | /flows/{id}/runs | List runs for the flow (cursor pagination) |
| GET | /runs/{runID} | Read one run |
| GET | /runs/{runID}/events | SSE stream of run events |
| GET | /runs/{runID}/events/history | Historical event replay |
| GET | /flows/{id}/live | SSE stream of live metrics.Frame rollups |
| GET | /flows/{id}/metrics/history | Historical frames |
| GET | /deployments | List deployments |
| POST | /deployments | Create 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-template | Provision + 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.jsonand the module browser UI - Modules can override auto-generated nodes by providing an explicit
FlowDSLNodewith the sameWorkflowRef - 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:
| File | Purpose |
|---|---|
module.yaml | Node definitions (embedded via go:embed) |
ir.go | //go:embed module.yaml + var moduleYAML []byte |
module.go | init() factory + Module struct embedding modules.IRBase |
module_test.go | Factory registration, manifest, node catalog, kind checks |
Opt-in via blank import alongside the parent module:
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:
| Subpackage | Factory | Nodes | Kinds |
|---|---|---|---|
modules/auth/flowdsl | auth-flowdsl | 7 | 1 action, 1 transform, 1 router, 4 source |
modules/users/flowdsl | users-flowdsl | 8 | 3 action, 2 transform, 3 source |
modules/groups/flowdsl | groups-flowdsl | 7 | 3 action, 1 transform, 3 source |
go-modules:
| Subpackage | Factory | Nodes | Kinds |
|---|---|---|---|
scheduler/flowdsl | scheduler-flowdsl | 4 | 2 source (with settings_schema), 2 action |
settings/flowdsl | settings-flowdsl | 3 | 1 transform (with NotFound), 2 action |
storage/flowdsl | storage-flowdsl | 4 | 2 action, 2 transform (1 NotFound, 1 settings_schema) |
notifications/flowdsl | notifications-flowdsl | 3 | 2 action, 1 source |
verification/flowdsl | verification-flowdsl | 2 | 1 action, 1 router (Valid/Invalid) |
links/flowdsl | links-flowdsl | 3 | 2 action, 1 transform (with NotFound) |
go-module-email:
| Subpackage | Factory | Nodes | Kinds |
|---|---|---|---|
flowdsl | email-flowdsl | 4 | 2 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.
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:
go install github.com/redelay/go-flowdsl/cmd/flowdsl@latest
Commands
| Command | Description |
|---|---|
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
flowdsl validate app.flow.yaml
Parses the file and runs flowdsl.Validate. Prints all diagnostics with [E]/[W]/[I]
prefixes. On success:
ok — no issues found
On failure (exit code 1):
[E] FDL010: event "order_event" has no entity_type
[W] FDL002: version is empty
gen
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
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
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.
| Flag | Default | Description |
|---|---|---|
-f, --format | flowdsl | Input format: flowdsl or businessmodel |
--out | json | Output format: json or yaml |
For the heavier CLI with asyncapi/openapi import, scaffold, audit, diff, and convert commands, see the redelayctl reference.