Guides

FlowDSL to Service

Take a FlowDSL diagram, validate it, generate Go models, and run it with Redelay.

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.

What you'll build: A webhook alert pipeline — ingests HTTP webhooks, filters by severity, enriches the payload, then dispatches alerts to PagerDuty and Slack.

1. Understand the FlowDSL vocabulary

A FlowDSL document has three sections:

SectionWhat it declares
nodesBusiness logic units — each maps to one handler function
edgesHow packets flow between nodes — transport, durability, retries
components.packetsTyped schemas for the data on each edge

Node kinds

KindRole in RedelayGo pattern
sourceEntry point — no inputsHTTP handler / Kafka consumer that publishes a typed.EventDefinition
transformMaps one packet type to anotherPure function: input → output, no side effects
routerRoutes to one of several outputs based on contentSwitch/map on a field value, returns the port name
llmCalls a language modelHTTP call to OpenAI/Anthropic + JSON extraction
actionSide effect in an external systemCalls external API, sends message, creates ticket
terminalEnd of path — no outputsArchive to DB, emit audit event, discard
publishEmits to the event busCalls eventBus.Publish() with a typed message
checkpointSaves pipeline state and passes throughUpsert to MongoDB before yielding
integrationBridges to another FlowDSL flowCalls the other flow's source node via the runtime API

Delivery modes (edge mode)

ModeTransportWhen to use
directIn-process, synchronousCheap deterministic transforms — field extraction, validation
ephemeralRedis Streams / NATSBurst absorption, worker pool smoothing
checkpointMongoDB per stageLong ETL pipelines — resume from last successful stage
durableMongoDB packet storeBusiness-critical — payments, email/SMS sends, LLM calls
streamKafka / Redis / NATSFan-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:

yaml
# 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

shell
# 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.

shell
# 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

text
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

go
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

go
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.

go
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.

go
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.

go
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.

go
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.

go
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

go
// 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

go
// 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))
}
shell
# Start infrastructure
cd infra && make up

# Run the service
cd backend && go run ./cmd/api/

The service exposes:

EndpointWhat it is
POST /api/v1/webhooksSource node — webhook ingest
GET /asyncapi.jsonAuto-generated AsyncAPI spec from your events
GET /referenceOpenAPI docs for the HTTP actions
GET /healthHealth check

8. Test without Kafka

Use the in-memory transport to test the full node chain in unit tests — no Docker required.

go
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