FlowDSL to Service
This guide takes you from a FlowDSL diagram — as seen on flowdsl.com/examples — through every step to a running Redelay service with typed, testable node implementations.
1. Understand the FlowDSL vocabulary
A FlowDSL document has three sections:
| Section | What it declares |
|---|---|
nodes | Business logic units — each maps to one handler function |
edges | How packets flow between nodes — transport, durability, retries |
components.packets | Typed schemas for the data on each edge |
Node kinds
| Kind | Role in Redelay | Go pattern |
|---|---|---|
source | Entry point — no inputs | HTTP handler / Kafka consumer that publishes a typed.EventDefinition |
transform | Maps one packet type to another | Pure function: input → output, no side effects |
router | Routes to one of several outputs based on content | Switch/map on a field value, returns the port name |
llm | Calls a language model | HTTP call to OpenAI/Anthropic + JSON extraction |
action | Side effect in an external system | Calls external API, sends message, creates ticket |
terminal | End of path — no outputs | Archive to DB, emit audit event, discard |
publish | Emits to the event bus | Calls eventBus.Publish() with a typed message |
checkpoint | Saves pipeline state and passes through | Upsert to MongoDB before yielding |
integration | Bridges to another FlowDSL flow | Calls the other flow's source node via the runtime API |
Delivery modes (edge mode)
| Mode | Transport | When to use |
|---|---|---|
direct | In-process, synchronous | Cheap deterministic transforms — field extraction, validation |
ephemeral | Redis Streams / NATS | Burst absorption, worker pool smoothing |
checkpoint | MongoDB per stage | Long ETL pipelines — resume from last successful stage |
durable | MongoDB packet store | Business-critical — payments, email/SMS sends, LLM calls |
stream | Kafka / Redis / NATS | Fan-out, external integration, event sourcing |
2. Pick an example from flowdsl.com/examples
Go to flowdsl.com/examples and pick a flow. For this guide we'll build the Webhook Alert Pipeline pattern:
# webhook-alert.flowdsl.yaml
version: "1.0"
title: Webhook Alert Pipeline
description: >
Ingest HTTP webhooks, filter by severity, enrich with metadata,
then dispatch to PagerDuty and Slack.
nodes:
IngestWebhook:
operationId: ingest_webhook
kind: source
summary: Receives raw webhook payloads over HTTP POST /webhooks
FilterBySeverity:
operationId: filter_by_severity
kind: router
summary: Routes critical/high to the alert path, low to the archive
inputs:
in: { packet: RawWebhook }
outputs:
alert: { packet: RawWebhook }
archive: { packet: RawWebhook }
EnrichPayload:
operationId: enrich_webhook
kind: transform
summary: Adds team ownership and service metadata from the config service
inputs:
in: { packet: RawWebhook }
outputs:
out: { packet: EnrichedWebhook }
DispatchPagerDuty:
operationId: dispatch_pagerduty
kind: action
summary: Opens a PagerDuty incident
inputs:
in: { packet: EnrichedWebhook }
outputs:
out: { packet: IncidentResult }
DispatchSlack:
operationId: dispatch_slack
kind: action
summary: Posts an alert message to the team Slack channel
inputs:
in: { packet: EnrichedWebhook }
ArchiveLow:
operationId: archive_low_severity
kind: terminal
summary: Writes low-severity webhooks to the archive collection
inputs:
in: { packet: RawWebhook }
edges:
- from: IngestWebhook
to: FilterBySeverity
delivery:
mode: direct
packet: RawWebhook
- from: FilterBySeverity.alert
to: EnrichPayload
delivery:
mode: durable
packet: RawWebhook
- from: FilterBySeverity.archive
to: ArchiveLow
delivery:
mode: ephemeral
packet: RawWebhook
stream: archive-queue
- from: EnrichPayload
to: DispatchPagerDuty
delivery:
mode: durable
packet: EnrichedWebhook
retryPolicy:
maxAttempts: 3
backoff: exponential
initialDelay: PT2S
idempotencyKey: "{{payload.webhookId}}-pagerduty"
- from: EnrichPayload
to: DispatchSlack
delivery:
mode: durable
packet: EnrichedWebhook
idempotencyKey: "{{payload.webhookId}}-slack"
components:
packets:
RawWebhook:
type: object
properties:
webhookId: { type: string }
source: { type: string }
severity:
type: string
enum: [critical, high, medium, low]
payload: { type: object }
receivedAt: { type: string, format: date-time }
required: [webhookId, source, severity]
EnrichedWebhook:
type: object
properties:
webhookId: { type: string }
source: { type: string }
severity: { type: string }
payload: { type: object }
teamOwner: { type: string }
serviceSlug: { type: string }
receivedAt: { type: string, format: date-time }
required: [webhookId, source, severity, teamOwner]
IncidentResult:
type: object
properties:
incidentId: { type: string }
url: { type: string }
status: { type: string }
3. Validate the flow
# Validate FlowDSL syntax and cross-references
redelayctl validate webhook-alert.flowdsl.yaml
# Output on success:
# OK — Webhook Alert Pipeline: 0 errors, 0 warnings
Validation checks: entity_type / action on every event, step IDs, depends_on refs, field types, channel protocols, schedule action refs.
4. Generate IR and module files
Redelay's compile pipeline converts FlowDSL → IR (Intermediate Representation) → Go module scaffolds.
# Step 1: Convert FlowDSL to IR JSON
redelayctl convert flowdsl json webhook-alert.flowdsl.yaml > webhook-alert.ir.json
# Step 2: Inspect the IR (optional — useful for debugging)
redelayctl ir webhook-alert.flowdsl.yaml
# Step 3: Generate module YAML files from IR
redelayctl generate webhook-alert.ir.json ./modules/
# Step 4: Export to AsyncAPI (for API docs and schema contracts)
redelayctl convert flowdsl asyncapi webhook-alert.flowdsl.yaml > asyncapi.yaml
# Step 5: Export to OpenAPI (for HTTP action docs)
redelayctl convert flowdsl openapi webhook-alert.flowdsl.yaml > openapi.yaml
# Step 6: Generate SVG diagram for each module
redelayctl diagram ./modules/alerts/module.yaml diagram.svg
Tip: Add the export steps to your CI pipeline. The AsyncAPI and OpenAPI outputs stay in sync with your flow automatically.
5. Translate nodes to Go handlers
Each node operationId maps to a function or struct in Go. The pattern is the same for every kind:
Project layout
modules/alerts/
├── module.go # Module wrapper + init() factory registration
├── config.go # Env-based configuration
├── packets.go # Go structs matching components.packets
├── nodes/
│ ├── source.go # ingest_webhook — HTTP handler → event publish
│ ├── router.go # filter_by_severity
│ ├── transform.go # enrich_webhook
│ ├── action_pd.go # dispatch_pagerduty
│ ├── action_sl.go # dispatch_slack
│ └── terminal.go # archive_low_severity
└── events.go # typed.EventDefinition declarations
packets.go — typed packet structs
package alerts
import "time"
// RawWebhook matches components.packets.RawWebhook in the FlowDSL document.
type RawWebhook struct {
WebhookID string `json:"webhookId"`
Source string `json:"source"`
Severity string `json:"severity"` // critical | high | medium | low
Payload map[string]any `json:"payload"`
ReceivedAt time.Time `json:"receivedAt"`
}
// EnrichedWebhook is the output of the enrich_webhook transform node.
type EnrichedWebhook struct {
WebhookID string `json:"webhookId"`
Source string `json:"source"`
Severity string `json:"severity"`
Payload map[string]any `json:"payload"`
TeamOwner string `json:"teamOwner"`
ServiceSlug string `json:"serviceSlug"`
ReceivedAt time.Time `json:"receivedAt"`
}
// IncidentResult is the output of the dispatch_pagerduty action node.
type IncidentResult struct {
IncidentID string `json:"incidentId"`
URL string `json:"url"`
Status string `json:"status"`
}
events.go — event definitions
package alerts
import "github.com/redelay/go-events/typed"
// WebhookReceivedEvent corresponds to the IngestWebhook source node.
// Publishing this event is what "triggers" the flow.
var WebhookReceivedEvent = typed.EventDefinition[RawWebhook]{
Name: "webhook.received",
EntityType: "webhook",
Action: "received",
Topic: "webhooks.received",
}
// WebhookAlertedEvent marks the end of a successful alert dispatch.
var WebhookAlertedEvent = typed.EventDefinition[EnrichedWebhook]{
Name: "webhook.alerted",
EntityType: "webhook",
Action: "alerted",
Topic: "webhooks.alerted",
}
nodes/source.go — ingest_webhook (kind: source)
A source node is an entry point. In Redelay it's an HTTP handler that validates input and publishes the typed event.
package nodes
import (
"net/http"
"time"
"github.com/redelay/go-events/typed"
"github.com/redelay/go-framework/server/httputil"
"github.com/redelay/go-framework/modules"
"github.com/myorg/myapp/modules/alerts"
)
type IngestWebhookHandler struct {
eventBus modules.EventBus
}
func NewIngestWebhookHandler(bus modules.EventBus) *IngestWebhookHandler {
return &IngestWebhookHandler{eventBus: bus}
}
// Handle is the HTTP POST /webhooks endpoint — the source node entry.
func (h *IngestWebhookHandler) Handle(w http.ResponseWriter, r *http.Request) {
var raw alerts.RawWebhook
if err := httputil.ReadJSON(r, &raw); err != nil {
httputil.WriteError(w, http.StatusBadRequest, "invalid payload", err)
return
}
if raw.WebhookID == "" || raw.Source == "" || raw.Severity == "" {
httputil.WriteError(w, http.StatusUnprocessableEntity, "missing required fields", nil)
return
}
raw.ReceivedAt = time.Now().UTC()
// Publishing the event is all this node does — Redelay routes it.
msg, err := alerts.WebhookReceivedEvent.NewMessage(
raw.WebhookID,
typed.Actor{Type: typed.ActorTypeSystem, ID: "alerts"},
raw,
)
if err != nil {
httputil.WriteError(w, http.StatusInternalServerError, "event creation failed", err)
return
}
if err := h.eventBus.Publish(r.Context(), msg); err != nil {
httputil.WriteError(w, http.StatusInternalServerError, "publish failed", err)
return
}
httputil.WriteJSON(w, http.StatusAccepted, map[string]string{
"webhookId": raw.WebhookID,
"status": "accepted",
})
}
nodes/router.go — filter_by_severity (kind: router)
A router node decides which downstream path receives the packet. No side effects.
package nodes
import (
"context"
"github.com/redelay/go-events/eventbus"
"github.com/redelay/go-events/typed"
"github.com/redelay/go-framework/modules"
"github.com/myorg/myapp/modules/alerts"
)
// FilterBySeverityNode subscribes to WebhookReceivedEvent and re-publishes
// to AlertPath or ArchivePath based on severity.
type FilterBySeverityNode struct {
alertEvent typed.EventDefinition[alerts.RawWebhook]
archiveEvent typed.EventDefinition[alerts.RawWebhook]
eventBus modules.EventBus
}
// Route returns the output port name — "alert" or "archive".
// This is the pure routing logic; publishing happens in the subscriber wrapper.
func (n *FilterBySeverityNode) Route(ctx context.Context, webhook *alerts.RawWebhook) string {
switch webhook.Severity {
case "critical", "high":
return "alert"
default:
return "archive"
}
}
nodes/transform.go — enrich_webhook (kind: transform)
A transform node is a pure mapping function — takes one packet type, returns another.
package nodes
import (
"context"
"github.com/myorg/myapp/modules/alerts"
)
// EnrichWebhookNode enriches a RawWebhook with team ownership metadata.
// It is a pure transform — no side effects, no external calls beyond the
// config service lookup (which should be in-process or low-latency Redis).
type EnrichWebhookNode struct {
ownershipMap map[string]string // source → teamOwner
slugMap map[string]string // source → serviceSlug
}
// Transform is the core function: RawWebhook → EnrichedWebhook.
// If called via the event bus, wrap this in a typed subscriber.
func (n *EnrichWebhookNode) Transform(ctx context.Context, raw *alerts.RawWebhook) (*alerts.EnrichedWebhook, error) {
teamOwner, ok := n.ownershipMap[raw.Source]
if !ok {
teamOwner = "platform"
}
serviceSlug, ok := n.slugMap[raw.Source]
if !ok {
serviceSlug = raw.Source
}
return &alerts.EnrichedWebhook{
WebhookID: raw.WebhookID,
Source: raw.Source,
Severity: raw.Severity,
Payload: raw.Payload,
TeamOwner: teamOwner,
ServiceSlug: serviceSlug,
ReceivedAt: raw.ReceivedAt,
}, nil
}
nodes/action_pd.go — dispatch_pagerduty (kind: action)
An action node performs an external side effect. It must be idempotent — Redelay can retry on failure.
package nodes
import (
"context"
"fmt"
"net/http"
"github.com/myorg/myapp/modules/alerts"
)
// DispatchPagerDutyNode opens a PagerDuty incident.
// The edge declares idempotencyKey = "{{payload.webhookId}}-pagerduty",
// so duplicate calls are safe — PagerDuty deduplicates by dedup_key.
type DispatchPagerDutyNode struct {
routingKey string
httpClient *http.Client
}
func (n *DispatchPagerDutyNode) Handle(ctx context.Context, enriched *alerts.EnrichedWebhook) (*alerts.IncidentResult, error) {
dedupKey := fmt.Sprintf("%s-pagerduty", enriched.WebhookID)
body := map[string]any{
"routing_key": n.routingKey,
"dedup_key": dedupKey,
"event_action": "trigger",
"payload": map[string]any{
"summary": fmt.Sprintf("[%s] %s — %s", enriched.Severity, enriched.TeamOwner, enriched.Source),
"severity": enriched.Severity,
"source": enriched.ServiceSlug,
"custom_details": map[string]any{
"webhookId": enriched.WebhookID,
"teamOwner": enriched.TeamOwner,
"serviceSlug": enriched.ServiceSlug,
},
},
}
_ = body // call PagerDuty Events API v2 here
return &alerts.IncidentResult{
IncidentID: dedupKey,
Status: "triggered",
}, nil
}
nodes/terminal.go — archive_low_severity (kind: terminal)
A terminal node has no outputs. It logs, archives, or discards — and returns.
package nodes
import (
"context"
"go.mongodb.org/mongo-driver/mongo"
"github.com/myorg/myapp/modules/alerts"
)
// ArchiveLowSeverityNode writes the webhook to the archive collection.
// It is a terminal — nothing is published after this.
type ArchiveLowSeverityNode struct {
coll *mongo.Collection
}
func (n *ArchiveLowSeverityNode) Handle(ctx context.Context, webhook *alerts.RawWebhook) error {
_, err := n.coll.InsertOne(ctx, map[string]any{
"webhookId": webhook.WebhookID,
"source": webhook.Source,
"severity": webhook.Severity,
"payload": webhook.Payload,
"receivedAt": webhook.ReceivedAt,
"archivedBy": "archive_low_severity",
})
return err
}
6. Wire everything into the module
// modules/alerts/module.go
package alerts
import (
"context"
"github.com/go-chi/chi/v5"
"github.com/redelay/go-framework/modules"
"github.com/myorg/myapp/modules/alerts/nodes"
)
func init() {
modules.RegisterFactory("alerts", func(deps modules.ModuleDeps) modules.Module {
return &Module{deps: deps}
})
}
type Module struct {
deps modules.ModuleDeps
// nodes
source *nodes.IngestWebhookHandler
enricher *nodes.EnrichWebhookNode
pdNode *nodes.DispatchPagerDutyNode
slNode *nodes.DispatchSlackNode
archive *nodes.ArchiveLowSeverityNode
}
func (m *Module) Manifest() modules.Manifest {
return modules.Manifest{
ID: "alerts",
Name: "Webhook Alert Pipeline",
Description: "FlowDSL-defined webhook ingest → filter → enrich → dispatch.",
Version: "1.0.0",
}
}
func (m *Module) Configure(r *modules.Registry) {}
func (m *Module) Startup(ctx context.Context) error {
cfg := DefaultConfig()
// Instantiate nodes
m.source = nodes.NewIngestWebhookHandler(m.deps.EventBus)
m.enricher = nodes.NewEnrichWebhookNode(cfg.OwnershipMapPath)
m.pdNode = nodes.NewDispatchPagerDutyNode(cfg.PagerDutyRoutingKey)
m.slNode = nodes.NewDispatchSlackNode(cfg.SlackWebhookURL)
db := m.deps.MongoDB.Database(m.deps.Config.MongoDatabase)
m.archive = nodes.NewArchiveLowSeverityNode(db.Collection(cfg.ArchiveCollection))
// Wire event subscriptions (the "edges" in Go)
// edge: IngestWebhook → FilterBySeverity (direct — handled inline by the router)
// edge: FilterBySeverity.archive → ArchiveLow (ephemeral via Redis stream)
// edge: FilterBySeverity.alert → EnrichPayload → PagerDuty + Slack (durable via EventBus)
return nil
}
func (m *Module) Shutdown(ctx context.Context) error { return nil }
func (m *Module) Routes(prefix string) func(chi.Router) {
return func(r chi.Router) {
r.Post("/webhooks", m.source.Handle)
}
}
7. Bootstrap and run
// cmd/api/main.go
package main
import (
"github.com/redelay/go-framework/app"
"github.com/redelay/go-framework/server"
"github.com/redelay/go-events/eventbus"
"github.com/redelay/go-events/transport/kafka"
_ "github.com/redelay/go-framework/modules/health"
_ "github.com/redelay/go-framework/modules/openapi"
_ "github.com/redelay/go-framework/modules/asyncapi" // serves /asyncapi.json
_ "github.com/myorg/myapp/modules/alerts"
)
func main() {
// Wire Kafka transport
t, err := kafka.NewFromEnv()
if err != nil {
panic(err)
}
bus := eventbus.New(t)
// Bootstrap Redelay
inst, err := app.Bootstrap(app.WithEventBus(bus))
if err != nil {
panic(err)
}
// Run HTTP server
app.Run(inst, server.Default(inst))
}
# Start infrastructure
cd infra && make up
# Run the service
cd backend && go run ./cmd/api/
The service exposes:
| Endpoint | What it is |
|---|---|
POST /api/v1/webhooks | Source node — webhook ingest |
GET /asyncapi.json | Auto-generated AsyncAPI spec from your events |
GET /reference | OpenAPI docs for the HTTP actions |
GET /health | Health check |
8. Test without Kafka
Use the in-memory transport to test the full node chain in unit tests — no Docker required.
package alerts_test
import (
"context"
"testing"
"github.com/redelay/go-events/eventbus"
"github.com/redelay/go-events/transport/memory"
"github.com/myorg/myapp/modules/alerts"
"github.com/myorg/myapp/modules/alerts/nodes"
)
func TestFilterBySeverity_RoutesCorrectly(t *testing.T) {
tr := memory.New()
bus := eventbus.New(tr)
_ = bus
router := &nodes.FilterBySeverityNode{}
critical := &alerts.RawWebhook{WebhookID: "w1", Source: "api", Severity: "critical"}
low := &alerts.RawWebhook{WebhookID: "w2", Source: "api", Severity: "low"}
if got := router.Route(context.Background(), critical); got != "alert" {
t.Errorf("expected alert, got %q", got)
}
if got := router.Route(context.Background(), low); got != "archive" {
t.Errorf("expected archive, got %q", got)
}
}
func TestEnrichWebhook_AddsOwnership(t *testing.T) {
enricher := nodes.NewEnrichWebhookNode("")
// inject test ownership map
raw := &alerts.RawWebhook{WebhookID: "w1", Source: "checkout-api", Severity: "high"}
enriched, err := enricher.Transform(context.Background(), raw)
if err != nil {
t.Fatal(err)
}
if enriched.WebhookID != raw.WebhookID {
t.Errorf("webhookId not preserved")
}
}
Next steps
- FlowDSL concepts — Nodes — all nine node kinds in depth
- FlowDSL concepts — Delivery Modes — choose the right mode for each edge
- FlowDSL Studio — visual canvas to draw and validate flows before writing YAML
- FlowDSL Examples — copy real flow patterns (smart email triage, DB anomaly detection, sales pipeline)
- Adding Events to a Go Module — full guide to typed event publishing in Redelay
- Choosing a Transport — Kafka vs NATS vs Redis Streams for your edges
- FlowDSL Examples (this site) — interactive examples with YAML viewer and export commands