Guides

Running Flows

End-to-end guide for executing FlowDSL flows in a Redelay service — create a flow, save and publish a version, launch a run, and subscribe to live events over SSE.

Running Flows

This guide walks through running FlowDSL flows in a Redelay Go service using the flowexec addon module. You will:

  1. Wire flowexec into your app and register a node handler.
  2. Create a flow, save an immutable version, publish it.
  3. Launch a run and stream its events live over SSE.
  4. Understand how immutable versioning keeps replays deterministic.

The underlying orchestration, storage, and sink abstractions live in go-flowdsl/flowexec. go-flowdsl/flowexec/module is the HTTP wrapper that exposes them through your Redelay server.

Why this module exists

runtime.Engine in go-flowdsl executes a single ir.Workflow in memory. Real applications need more:

  • Persistence — flows that survive restarts, with a clean edit/publish cycle.
  • Replay determinism — past runs must resolve to the exact workflow bytes they ran on, even after the flow is edited.
  • Observability — per-node, per-run lifecycle events for UI dashboards, audit logs, and alerting.
  • Live UX — subscribers want progress without polling.

flowexec adds all four without coupling the engine to a database or transport.

Prerequisites

  • A working Redelay Go app (see Creating a Go App).
  • MongoDB reachable at MONGO_URI for production; none required for tests.

1. Register the module

go
package main

import (
    _ "github.com/redelay/go-framework/modules/health"
    _ "github.com/redelay/go-framework/modules/openapi"
    _ "github.com/redelay/go-flowdsl/flowexec/module"   // the module
)

func main() {
    server.Default().Run()
}

On startup flowexec auto-wires:

  • mongostore when deps.DB is non-nil, otherwise memstore.
  • mongots time-series sink when Mongo is available, otherwise noop.
  • A live-fanout broadcaster between the executor and sink so SSE subscribers receive events in real time while the sink still gets the canonical record.

2. Register node handlers

Handlers are registered on the underlying runtime.Engine exposed by the module. The cleanest spot is right after app.Bootstrap:

go
import (
    "context"
    "github.com/redelay/go-framework/app"
    "github.com/redelay/go-flowdsl/flowexec/module"
    "github.com/redelay/go-flowdsl/runtime"
)

inst, _ := app.Bootstrap()
m := inst.Registry.MustGet("flowexec").(*flowexec.Module)

m.Engine().RegisterHandler("send_welcome_email", func(ctx context.Context, step *runtime.Step) error {
    step.Output["sent"] = true
    return nil
})

3. Create a flow and save a version

Flows are mutable metadata. Versions are immutable documents identified by a content hash. You edit the flow by creating new versions; you promote a version by publishing it.

shell
# Create the flow
curl -sX POST $API/flows \
  -H 'Content-Type: application/json' \
  -d '{"name": "Welcome sequence"}'
# → { "id": "flow.018f...", "name": "Welcome sequence", "headVersionID": "", "publishedVersionID": "" }

FLOW=flow.018f...

# Save an immutable version
curl -sX POST $API/flows/$FLOW/versions \
  -H 'Content-Type: application/json' \
  -d @welcome.flow.json
# → { "id": "v.018f...", "flowID": "flow.018f...", "versionHash": "f3e1...", ... }

welcome.flow.json is a full ir.Workflow document — either hand-written or compiled from a FlowDSL source file via flowdsl compile.

Using the Studio spec format

If you are building flows in the FlowDSL Studio editor, the document is in spec format (nested flows, components.nodes, x-ui metadata) rather than raw ir.Workflow format.

Append ?format=spec to the version endpoints and the server handles conversion automatically:

shell
# Save a spec-format document as a version
curl -sX POST "$API/flows/$FLOW/versions?format=spec" \
  -H 'Content-Type: application/json' \
  -d @my-flow.spec.json

# Read a version back in spec format (for the Studio editor)
curl -s "$API/flows/$FLOW/versions/$VID?format=spec"

The server uses go-flowdsl/spec to convert between spec documents and ir.Workflow transparently. Without ?format=spec, endpoints use raw ir.Workflow format.

4. Publish the version

Publishing flips the flow's publishedVersionID pointer. It is atomic, cheap, and non-destructive — older versions stay in the store and remain addressable by ID.

shell
curl -sX POST $API/flows/$FLOW/publish \
  -H 'Content-Type: application/json' \
  -d '{"versionID": "v.018f..."}'

Calling POST /flows/$FLOW/runs before publishing returns 409 Conflict with flow has no published version.

5. Launch a run

shell
curl -sX POST $API/flows/$FLOW/runs \
  -H 'Content-Type: application/json' \
  -d '{"input": {"user_id": "u123"}}'
# → { "runID": "run.018f...", "versionID": "v.018f...", "versionHash": "f3e1..." }

The response is 202 Accepted. The run executes asynchronously against the published version — not whatever happens to be the latest save — so a publish-then-edit race can never land a partially-saved draft in production.

6. Subscribe to live events

GET /runs/{runID}/events is a text/event-stream. Subscribe immediately after launching, or any time before the run completes, to see progress live:

shell
curl -N $API/runs/run.018f.../events
text
: connected

event: run.started
data: {"runID":"run.018f...","kind":"run.started","versionHash":"f3e1...",...}

event: node.started
data: {"runID":"run.018f...","kind":"node.started","nodeID":"send_welcome_email",...}

event: node.done
data: {"runID":"run.018f...","kind":"node.done","nodeID":"send_welcome_email",...}

event: run.completed
data: {"runID":"run.018f...","kind":"run.completed","durationMs":42,...}

The stream closes automatically on run.completed or run.failed. : keep-alive lines are emitted every 15 seconds so reverse proxies do not drop idle connections.

Browser-side example

ts
const es = new EventSource(`/api/runs/${runID}/events`);
es.addEventListener("node.started", (ev) => {
  const payload = JSON.parse(ev.data);
  console.log("node started:", payload.nodeID);
});
es.addEventListener("run.completed", () => es.close());
es.addEventListener("run.failed", () => es.close());

Subscribe-before-launch race

RunAsync pre-generates the runID and stamps it onto RunInput before the goroutine fires run.started. This means the HTTP response returns the runID to the client immediately while the executor is still about to emit its first event. Most clients will attach to the SSE stream before run.started arrives; if not, the broadcaster fans subsequent events out normally. The SSE stream is live-only by design — use the sink's Reader interface (implemented by mongots) if you need a full replayable history.

7. Immutable versioning — why replay is deterministic

Every event emitted during a run carries versionHash = sha256(workflow JSON) of the version it resolved. Six months later you can audit a failed run and be certain you are looking at the same document the executor ran on, even if the flow has been edited a hundred times since.

The pointer-flip publish means:

  • Editing a flow appends a new row to flow_versions and advances Flow.HeadVersionID. Runs are unaffected.
  • Publishing updates Flow.PublishedVersionID atomically. The next run sees the new version; in-flight runs keep running against the version they started on.
  • Deleting the flow removes metadata and versions, but any already-written runs and events in the sink remain addressable by runID + versionHash.

Deployment without MongoDB

For tests or ephemeral worker processes, omit MONGO_URI. The module falls back to memstore

  • noop, which keeps the HTTP surface usable without any external dependencies. Flows and versions live in memory for the lifetime of the process.

Next steps

  • Trace + Live Observability — once a flow is running, use Flow Studio's Trace and Live modes to inspect individual runs and watch aggregate health (rate, p95, error rate, saturation) in real time with server-issued recommendations.
  • Flow Templates — let your module ship starter flows that show up in the Studio's "Create flow → From template" picker via the TemplateProvider interface.
  • Addon Modules → flowexec — route-by-route reference.
  • go-flowdsl → flowexec — executor, sinks, store, and metrics internals.
  • FlowDSL Examples — sample workflows ready to POST as versions.