Reference

Go Events

Full API reference for the go-events library — transports, EventBus, filters, typed messages.

Module: github.com/redelay/go-events

Overview

go-events provides a high-level EventBus over pluggable transport backends (Kafka, NATS, Redis Streams, in-memory). The canonical event types (EventMessage, ConsumerRegistration, EventBus interface) are defined in go-framework/runtime/transport (stdlib-only, zero external deps) and re-exported as type aliases from both go-events/core and go-framework/modules. This means all three import paths refer to identical types with no circular module dependency.

Dependency graph

text
go-framework/runtime/transport  ← canonical zero-dep type home (stdlib only)
        ↑                ↑
go-events/core      go-framework/modules  (type aliases — identical types)
        ↑
go-events/eventbus  ← no mongo/chi/JWT/zap
go-events/module    ← imports go-framework/modules (lifecycle only)

Lightweight workers (Kafka consumers, scheduled jobs) that only need event types and no HTTP server can import go-events/core without pulling in the full framework stack.

FlowDSL integration — go-flowdsl/nodes

The redelay/event-source node subscribes to any event on the bus. It replaces the per-event source nodes that each module used to emit.

The node ships from go-flowdsl/nodes, not from go-events (it lived at go-events/flowdsl until the 2026-04-16 flowdsl refactor).

Why it lives there: the node pack programs against the EventBus interface in go-framework/runtime/transport and imports go-events nowhere. Had it stayed under go-events, every flow using redelay/event-source would have pulled in the event library — and transitively sarama or nats.go — regardless of which transport the project actually ran. Keeping it in go-flowdsl/nodes means FlowDSL depends on the transport interface only, so any implementation (kafka, nats, redis, memory) satisfies it.

It is packaged with the other transport-entry nodes (json-stream-source, json-to-event, kafka-produce, nats-publish, redis-stream-publish), so one import gets the whole entry-node pack.

Enable with a blank import:

go
import _ "github.com/redelay/go-flowdsl/nodes"      // event/stream sources + publishers
import _ "github.com/redelay/go-flowdsl/nodes/core" // stdlib primitives (transforms, control flow)

Configurable event-source node

One node type serves every event in the system. The specific event is picked via the node's eventName setting, which Studio renders as a dropdown.

yaml
on_user_signup:
  kind: source
  nodeType: redelay/event-source
  settings:
    eventName: users.user_created  # chosen from live dropdown
    groupID:   onboarding-welcome
    filter:    "payload.email_verified == true"  # optional, server-side
  outputs:
    - name: Event
      message:
        $ref: "asyncapi:default#/components/schemas/EventMessage"

Settings:

SettingTypeDescription
eventNameenum (required)Event topic to subscribe to — enum populated at runtime from the live module registry
groupIDstringConsumer group — defaults to a flow-id-scoped auto-generated value
filterstringOptional expression evaluated by the event bus filter chain before the flow wakes

Output port is always the generic EventMessage envelope. Payload shape varies per event — refer to the live AsyncAPI document for per-event payload schemas. In edge transforms, project payload fields via {{.payload.*}}:

yaml
edges:
  - from: on_user_signup
    to:   send_welcome_email
    when: "output.name == 'Event'"
    transform:
      to:   "{{.payload.email}}"
      name: "{{.payload.name}}"

How the live enum works

go-flowdsl/nodes/module.go overrides FlowDSLNodes() on its Module. At spec-build time it walks registry.IRModules(), collects every declared event name across all loaded modules, and injects them as the enum on settings.eventName in the node manifest. The module YAML ships with enum: [] — the list is always live.

New events show up in Studio automatically as soon as the corresponding module is loaded; no manifest regeneration required.

JSON stream ingestion nodes

Two companion nodes let a flow ingest arbitrary JSON from external sources (webhooks, NDJSON streams, raw message buses) and map the fields into one of the system's declared events.

Typical pipeline:

text
redelay/json-stream-source ──► redelay/json-to-event ──► publish / downstream action

The split lets you either chain both nodes for a full receive-map-publish flow, or use redelay/json-to-event alone with JSON from any upstream node.

redelay/json-stream-source (source)

Subscribes to an external JSON message stream and emits each decoded document as a packet. Does not interpret the payload.

Supported transports (transport setting):

ValueDescription
http-webhookExposes an HTTP POST endpoint at the configured path
ndjsonTails an NDJSON file/URL
kafka-rawReads a Kafka topic where the value is raw JSON (no envelope)
natsSubscribes to a NATS subject (core or JetStream); optional natsQueue for queue groups
redis-streamReads entries from a Redis Stream

Common settings: path, url, topic, natsQueue, groupID, authHeader, authSecret (secret), maxBodyBytes. Emits on the message output port with fields raw, received_at, source_transport, headers.

redelay/json-to-event (transform)

Consumes a JSON document and transforms it into a typed Redelay event, optionally publishing it to the bus. Target event is picked via the eventName dropdown (live enum, same mechanism as redelay/event-source).

Mapping rules are JSONPath-style key/value pairs:

json
{
  "entity_id":      "$.customer.id",
  "payload.email":  "$.customer.email_address",
  "payload.status": "$.status || \"pending\""
}

Settings: eventName (required), mapping (required), actorType, actorID, publish, strict. Two outputs: event (the constructed envelope) and error (emitted when mapping fails — required fields missing, invalid JSONPath, etc.).

When publish=true, the transformed event is published to the bus in addition to being emitted on the event output port.


Consumer auto-wiring

When using the auto-wire pattern (import _ "github.com/redelay/go-events/module"), the framework automatically discovers and registers all consumer subscriptions declared by other modules. No manual bus.Register() call is needed anywhere in your code.

How it works:

  1. Each module that wants to consume events implements Consumers() from the EventsProvider interface:
    go
    func (m *Module) Consumers() []*modules.ConsumerRegistration {
        return []*modules.ConsumerRegistration{{
            Topic:   "email.send",
            GroupID: "email-sender",
            Handler: m.handleEmailSend,
        }}
    }
    
  2. During the Configure lifecycle phase (before Startup), the events module iterates all registered modules, finds every EventsProvider, and calls bus.Register(reg) for each declared consumer.
  3. bus.Start() is called in Startup — all consumers are already registered and wired to the transport.

This means adding a new consumer to any module is a one-step change: return it from Consumers(). The rest is automatic.


Bootstrap patterns

go
import (
    _ "github.com/redelay/go-events/module"  // registers factory
)

inst, _ := app.Bootstrap() // EventBus auto-created from TRANSPORT env
app.Run(inst, server.Default(inst))

The TRANSPORT environment variable selects the backend (kafka by default). An explicit app.WithEventBus(bus) always takes priority over the auto-wired bus.

Manual wiring

go
import (
    "github.com/redelay/go-events/eventbus"
    "github.com/redelay/go-events/transport/kafka"
)

t, err := kafka.NewFromEnv()
bus := eventbus.New(t)
inst, _ := app.Bootstrap(app.WithEventBus(bus))

Standalone worker (no full framework)

go
import (
    "github.com/redelay/go-events/core"
    "github.com/redelay/go-events/eventbus"
    "github.com/redelay/go-events/transport/kafka"
)

// Only core + eventbus + transport — no mongo-driver, chi, JWT, zap.
t, _ := kafka.NewFromEnv()
bus := eventbus.New(t)
bus.Register(&core.ConsumerRegistration{
    Topic:   "user.created",
    GroupID: "my-worker",
    Handler: func(ctx context.Context, msg *core.EventMessage) error { ... },
})
bus.Start(ctx)

Transport backends

Every transport exposes a NewFromEnv() constructor that reads all settings from environment variables:

go
// Pick one:
t, err := kafka.NewFromEnv()  // reads KAFKA_BOOTSTRAP_SERVERS, KAFKA_SASL_*, ...
t, err := nats.NewFromEnv()   // reads NATS_URL, NATS_JETSTREAM, NATS_CREDENTIALS_FILE, ...
t, err := redis.NewFromEnv()  // reads REDIS_STREAMS_URL, REDIS_STREAMS_BATCH_SIZE, ...

bus := eventbus.New(t)

See Transports for the full environment variable reference and cluster configuration.

Kafka

go
import "github.com/redelay/go-events/transport/kafka"

t, err := kafka.New(kafka.Config{
    Brokers:  []string{"kafka:9092"},
    GroupID:  "my-service",
    ClientID: "my-service",
    // SASL for managed Kafka (Confluent Cloud, MSK, Aiven):
    // SASL: &kafka.SASLConfig{
    //     Mechanism: kafka.SASLMechanismSCRAMSHA256,
    //     Username:  "user",
    //     Password:  "secret",
    // },
})

NATS

go
import "github.com/redelay/go-events/transport/nats"

// Core pub/sub (ephemeral — no persistence):
t, err := nats.New(nats.Config{
    URL:          "nats://localhost:4222",
    UseJetStream: false,
})

// JetStream (durable — persistent, consumer groups, replay):
t, err := nats.New(nats.Config{
    URL:          "nats://nats1:4222,nats://nats2:4222",  // cluster
    UseJetStream: true,
    AckWait:      30 * time.Second,
    MaxDeliver:   3,
})

JetStream auto-creates streams on demand — no manual NATS admin needed. A stream is provisioned before the first QueueSubscribe and before the first Publish to a subject. The publish-side provisioning matters: without it, JetStream rejects a publish to a subject that has no stream, so publishing only worked for topics the same process happened to subscribe to first — an event consumed by another service, or by nothing yet, failed with "no response from stream". The stream-existence check is cached per subject, so it costs one round-trip per new topic, not one per message.

Redis Streams

go
import "github.com/redelay/go-events/transport/redis"

t, err := redis.New(redis.Config{
    Addr:          "redis:6379",
    BlockDuration: 200 * time.Millisecond,
    BatchSize:     10,
    MaxLen:        10_000,
})

Uses XADD / XREADGROUP under the hood. Consumer group semantics: each GroupID receives exactly one delivery across all replicas.

Memory (for tests)

go
import "github.com/redelay/go-events/transport/memory"

tr := memory.New()
// Full consumer-group semantics: each group gets the message,
// round-robin dispatch within a group.

No infrastructure required. Use for unit tests and fast feedback loops.

EventBus

go
import "github.com/redelay/go-events/eventbus"

bus := eventbus.New(t)          // t is any transport.Transport
bus.UsePublish(eventbus.LoggingPublishMiddleware)
bus.UseConsume(eventbus.LoggingConsumeMiddleware)

Key methods

MethodDescription
New(t)Create EventBus wrapping transport t
UsePublish(mw...)Append publish middleware
UseConsume(mw...)Append consume middleware
Register(reg)Register a ConsumerRegistration
Start(ctx)Wire registrations to transport and start consuming
Stop(ctx)Stop consuming
Publish(ctx, msg)Publish a modules.EventMessage
Close()Close the underlying transport

Built-in middleware

MiddlewareSideEffect
CorrelationPublishMiddlewarePublishSets correlation_id header from context (default)
CorrelationConsumeMiddlewareConsumeInjects correlation_id from message into context (default)
LoggingPublishMiddlewarePublishLogs entity_type, action, latency
LoggingConsumeMiddlewareConsumeLogs entity_type, action, latency

Context helpers

go
ctx = eventbus.WithCorrelationID(ctx, "req-abc")
cid := eventbus.CorrelationIDFromContext(ctx) // "req-abc"

Typed event definitions

Import github.com/redelay/go-events/typed to declare strongly-typed event definitions with JSON-serialised payloads.

go
import "github.com/redelay/go-events/typed"

type UserCreatedPayload struct {
    ID    string `json:"id"`
    Email string `json:"email"`
}

var UserCreatedEvent = typed.EventDefinition[UserCreatedPayload]{
    Name:       "user.created",
    EntityType: "user",
    Action:     "created",
    Topic:      "user.created",
}

msg, err := UserCreatedEvent.NewMessage(
    user.ID,
    typed.Actor{Type: typed.ActorTypeUser, ID: actorID},
    UserCreatedPayload{ID: user.ID, Email: user.Email},
)
bus.Publish(ctx, msg)

MustNewMessage panics instead of returning an error — tests only.

Where the pieces live

TypePackage
EventDefinition[T], Actor, ActorType*, SystemActorgo-events/typed
EmailSendPayload + EmailSendEventgo-module-email/events

Implementation sits in go-framework/runtime/events (stdlib + uuid, sibling to runtime/transport). go-events/typed re-exports from there so neither direction introduces a module cycle.

Filters

go
import "github.com/redelay/go-events/filter"
FilterDescription
HeaderFilter{Key, Value}Match headerKey == Value
ActionFilter{EntityType, Action}Match entity_type + action
RegexFilter{Key, Pattern}Match headerKey against regexp
CompositeFilter{Filters, Op}AND / OR combinator
PrefixFilter{Prefix}EntityType starts with prefix
FuncFilter{Fn}Custom predicate
All(filters...)Shorthand for AND composite
Any(filters...)Shorthand for OR composite
Func(fn)Shorthand for FuncFilter

Using a filter in a ConsumerRegistration:

go
bus.Register(&modules.ConsumerRegistration{
    Topic:   "events",
    GroupID: "g",
    Filter:  filter.All(&filter.ActionFilter{EntityType: "user", Action: "created"}).Matches,
    Handler: ...,
})

EventMessage

EventMessage is defined in go-framework/runtime/transport and aliased in both go-events/core and go-framework/modules — all three import paths refer to the same type.

go
type EventMessage struct {
    ID            string            `json:"event_id"`
    CorrelationID string            `json:"correlation_id,omitempty"`
    Timestamp     string            `json:"timestamp"`
    EntityType    string            `json:"entity_type"`
    EntityID      string            `json:"entity_id"`
    Action        string            `json:"action"`
    ActorType     string            `json:"actor_type,omitempty"`
    ActorID       string            `json:"actor_id,omitempty"`
    Payload       []byte            `json:"payload"`
    Headers       map[string]string `json:"headers,omitempty"`
}

Extract a typed payload:

go
var p OrderCreatedPayload
json.Unmarshal(msg.Payload, &p)

Legacy packages (backward compat)

PackageDescription
consumer/EventConsumer — direct sarama consumer group with handler registry
producer/EventProducer — direct sarama sync producer
messages/BaseEventMessage, TypedEventMessage[T], FlexTime, EventDefinition
schema/AsyncAPI schema registry + JSON Schema validator
client.goEventClient — legacy top-level API

These remain for backward compatibility. New code should use eventbus/ and transport/.

Configuration (env vars)

Kafka

VariableDescriptionDefault
KAFKA_BOOTSTRAP_SERVERSComma-separated broker listlocalhost:9092
KAFKA_GROUP_IDConsumer group ID""
KAFKA_SECURITY_PROTOCOLplaintext | sasl_plaintext | sasl_ssl | sslplaintext
KAFKA_SASL_MECHANISMplain | scram-sha-256 | scram-sha-512""
KAFKA_SASL_USERNAMESASL username""
KAFKA_SASL_PASSWORDSASL password""
ASYNCAPI_URLAsyncAPI schema endpoint for validation""

NATS

VariableDescriptionDefault
NATS_URLServer URL or comma-separated cluster URLsnats://localhost:4222
NATS_JETSTREAMtrue to enable durable JetStream deliveryfalse
NATS_CREDENTIALS_FILEPath to .creds file (NGS/Synadia accounts)""
NATS_NKEY_FILEPath to NKey seed file""
NATS_TLStrue to require TLSfalse
NATS_ACK_WAIT_SECONDSJetStream ack wait before redelivery30
NATS_MAX_DELIVERJetStream max redelivery attempts (-1 = unlimited)3

Redis Streams

VariableDescriptionDefault
REDIS_STREAMS_URLRedis URL (redis://host:port); falls back to REDIS_URLredis://localhost:6379
REDIS_STREAMS_BLOCK_MSXREADGROUP block duration in ms200
REDIS_STREAMS_BATCH_SIZEMessages per XREADGROUP fetch10
REDIS_STREAMS_MAX_LENMax stream length (0 = unlimited)10000

Build and test

shell
# Build all packages
cd go-events && go build ./...

# Unit tests (no infrastructure needed)
cd go-events && go test -count=1 -race ./...
# or via Makefile:
cd go-events && make test

# Integration tests (requires Docker — starts real Kafka, NATS, Redis containers)
cd go-events && go test -tags integration -v -timeout 180s ./transport/integration/
# or via Makefile:
cd go-events && make test-integration

# Both
cd go-events && make test-all

The integration test suite (transport/integration/) uses testcontainers-go to spin up:

TestContainer
TestMemoryTransport*None — in-process
TestNATSCoreTransport*nats:2-alpine
TestNATSJetStreamTransport*nats:2-alpine (JetStream enabled)
TestRedisStreamsTransport*redis:7-alpine
TestKafkaTransport*confluentinc/confluent-local:7.5.0

Each transport runs the same 5 scenarios: publish/subscribe, consumer groups, filtered delivery, correlation ID propagation, graceful shutdown.