Go Events
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
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:
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.
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:
| Setting | Type | Description |
|---|---|---|
eventName | enum (required) | Event topic to subscribe to — enum populated at runtime from the live module registry |
groupID | string | Consumer group — defaults to a flow-id-scoped auto-generated value |
filter | string | Optional 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.*}}:
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:
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):
| Value | Description |
|---|---|
http-webhook | Exposes an HTTP POST endpoint at the configured path |
ndjson | Tails an NDJSON file/URL |
kafka-raw | Reads a Kafka topic where the value is raw JSON (no envelope) |
nats | Subscribes to a NATS subject (core or JetStream); optional natsQueue for queue groups |
redis-stream | Reads 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:
{
"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:
- Each module that wants to consume events implements
Consumers()from theEventsProviderinterface:gofunc (m *Module) Consumers() []*modules.ConsumerRegistration { return []*modules.ConsumerRegistration{{ Topic: "email.send", GroupID: "email-sender", Handler: m.handleEmailSend, }} } - During the
Configurelifecycle phase (beforeStartup), the events module iterates all registered modules, finds everyEventsProvider, and callsbus.Register(reg)for each declared consumer. bus.Start()is called inStartup— 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
Auto-wire (recommended)
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
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)
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
NewFromEnv() (recommended)
Every transport exposes a NewFromEnv() constructor that reads all settings from environment variables:
// 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
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
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
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)
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
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
| Method | Description |
|---|---|
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
| Middleware | Side | Effect |
|---|---|---|
CorrelationPublishMiddleware | Publish | Sets correlation_id header from context (default) |
CorrelationConsumeMiddleware | Consume | Injects correlation_id from message into context (default) |
LoggingPublishMiddleware | Publish | Logs entity_type, action, latency |
LoggingConsumeMiddleware | Consume | Logs entity_type, action, latency |
Context helpers
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.
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
| Type | Package |
|---|---|
EventDefinition[T], Actor, ActorType*, SystemActor | go-events/typed |
EmailSendPayload + EmailSendEvent | go-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
import "github.com/redelay/go-events/filter"
| Filter | Description |
|---|---|
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:
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.
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:
var p OrderCreatedPayload
json.Unmarshal(msg.Payload, &p)
Legacy packages (backward compat)
| Package | Description |
|---|---|
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.go | EventClient — legacy top-level API |
These remain for backward compatibility. New code should use eventbus/ and transport/.
Configuration (env vars)
Kafka
| Variable | Description | Default |
|---|---|---|
KAFKA_BOOTSTRAP_SERVERS | Comma-separated broker list | localhost:9092 |
KAFKA_GROUP_ID | Consumer group ID | "" |
KAFKA_SECURITY_PROTOCOL | plaintext | sasl_plaintext | sasl_ssl | ssl | plaintext |
KAFKA_SASL_MECHANISM | plain | scram-sha-256 | scram-sha-512 | "" |
KAFKA_SASL_USERNAME | SASL username | "" |
KAFKA_SASL_PASSWORD | SASL password | "" |
ASYNCAPI_URL | AsyncAPI schema endpoint for validation | "" |
NATS
| Variable | Description | Default |
|---|---|---|
NATS_URL | Server URL or comma-separated cluster URLs | nats://localhost:4222 |
NATS_JETSTREAM | true to enable durable JetStream delivery | false |
NATS_CREDENTIALS_FILE | Path to .creds file (NGS/Synadia accounts) | "" |
NATS_NKEY_FILE | Path to NKey seed file | "" |
NATS_TLS | true to require TLS | false |
NATS_ACK_WAIT_SECONDS | JetStream ack wait before redelivery | 30 |
NATS_MAX_DELIVER | JetStream max redelivery attempts (-1 = unlimited) | 3 |
Redis Streams
| Variable | Description | Default |
|---|---|---|
REDIS_STREAMS_URL | Redis URL (redis://host:port); falls back to REDIS_URL | redis://localhost:6379 |
REDIS_STREAMS_BLOCK_MS | XREADGROUP block duration in ms | 200 |
REDIS_STREAMS_BATCH_SIZE | Messages per XREADGROUP fetch | 10 |
REDIS_STREAMS_MAX_LEN | Max stream length (0 = unlimited) | 10000 |
Build and test
# 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:
| Test | Container |
|---|---|
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.
Domain Addon Modules
Domain addon modules: scheduler, settings, storage, notifications, verification, links, metrics, search, i18n, apilog, content, workers, coordination, websockets, and flowexec.
Email Module
Transactional email delivery for Redelay: SMTP, SendGrid, Mailgun, Mailjet, and SES backends, HTML templating, event-driven delivery.