Transports
Redelay uses a transport abstraction so module code never imports Kafka, Redis, or any specific broker directly. The same module works with any backend just by changing bootstrap configuration.
Architecture
Module code only ever calls modules.EventBus. The transport is wired once at bootstrap and is invisible to modules.
Transport interface
Defined in go-framework/runtime/transport:
type Transport interface {
Publish(ctx context.Context, msg *TransportMessage) error
Subscribe(ctx context.Context, topics []string, groupID string, handler MessageHandler) error
Start(ctx context.Context) error
Stop(ctx context.Context) error
Close() error
}
type TransportMessage struct {
Topic string
Key []byte
Value []byte
Headers map[string][]byte
}
This package has zero external dependencies — only stdlib. Implementations live in go-events.
Available transports
| Transport | Package | Delivery | Cluster support | Use case |
|---|---|---|---|---|
| Memory | go-events/transport/memory | In-process | — | Unit tests, no infrastructure |
| Kafka | go-events/transport/kafka | Durable stream | Multi-broker + SASL/TLS | Production; highest throughput |
| NATS core | go-events/transport/nats | Ephemeral pub/sub | Cluster routes | Low-latency, fire-and-forget |
| NATS JetStream | go-events/transport/nats | Durable stream | Cluster + super-cluster | Durable with built-in replay |
| Redis Streams | go-events/transport/redis | Durable stream | Redis Cluster / Sentinel | Lightweight; no separate broker |
Delivery mode comparison
| Characteristic | Memory | Kafka | NATS core | NATS JetStream | Redis Streams |
|---|---|---|---|---|---|
| Persistence | ✗ | ✓ | ✗ | ✓ | ✓ |
| Consumer groups | ✓ | ✓ | ✗ | ✓ | ✓ |
| At-least-once | ✓* | ✓ | ✗ | ✓ | ✓ |
| Replay | ✗ | ✓ | ✗ | ✓ | ✓ |
| Backpressure | ✗ | ✓ | ✗ | ✓ | ✗ |
| Infra required | None | Kafka | NATS | NATS + JetStream | Redis |
*Memory transport gives at-least-once within a single process only.
EventBus
The EventBus is the module-facing publishing interface. It wraps a Transport and adds:
- JSON serialisation of
modules.EventMessage - Correlation ID propagation via context
- Composable publish/consume middleware
- Per-subscription event filters
// modules.EventBus (go-framework) — what modules receive in ModuleDeps
type EventBus interface {
Publish(ctx context.Context, msg *modules.EventMessage) error
Close() error
}
Implemented by eventbus.EventBus in go-events.
Bootstrap wiring
Auto-wire via blank import (simplest)
Import go-events/module to register an EventBus factory. app.Bootstrap() will call it automatically when no explicit bus is provided:
import (
_ "github.com/redelay/go-events/module" // registers factory
)
inst, _ := app.Bootstrap() // EventBus auto-created from TRANSPORT env
Set TRANSPORT=kafka|nats|redis|memory (defaults to kafka).
An explicit app.WithEventBus(bus) always takes priority.
Using NewFromEnv() (manual)
Every transport has a NewFromEnv() constructor that reads all settings from environment variables:
import (
"github.com/redelay/go-events/eventbus"
"github.com/redelay/go-events/transport/kafka"
)
t, err := kafka.NewFromEnv() // reads KAFKA_BOOTSTRAP_SERVERS, KAFKA_SASL_*, etc.
bus := eventbus.New(t)
inst, err := app.Bootstrap(app.WithEventBus(bus))
import (
"github.com/redelay/go-events/eventbus"
"github.com/redelay/go-events/transport/nats"
)
t, err := nats.NewFromEnv() // reads NATS_URL, NATS_JETSTREAM, NATS_CREDENTIALS_FILE, etc.
bus := eventbus.New(t)
inst, err := app.Bootstrap(app.WithEventBus(bus))
import (
"github.com/redelay/go-events/eventbus"
"github.com/redelay/go-events/transport/redis"
)
t, err := redis.NewFromEnv() // reads REDIS_STREAMS_URL, REDIS_STREAMS_BATCH_SIZE, etc.
bus := eventbus.New(t)
inst, err := app.Bootstrap(app.WithEventBus(bus))
import (
"github.com/redelay/go-events/eventbus"
"github.com/redelay/go-events/transport/memory"
)
tr := memory.New()
bus := eventbus.New(tr)
inst, err := app.Bootstrap(app.WithEventBus(bus))
Explicit configuration
Pass a Config struct directly for programmatic control:
t, err := kafka.New(kafka.Config{
Brokers: []string{"kafka1:9092", "kafka2:9092", "kafka3:9092"},
ClientID: "my-service",
GroupID: "my-service",
SASL: &kafka.SASLConfig{
Mechanism: kafka.SASLMechanismSCRAMSHA256,
Username: "user",
Password: "secret",
},
})
t, err := nats.New(nats.Config{
URL: "nats://nats1:4222,nats://nats2:4222,nats://nats3:4222",
UseJetStream: false,
})
t, err := nats.New(nats.Config{
URL: "nats://nats1:4222,nats://nats2:4222",
UseJetStream: true,
AckWait: 60 * time.Second,
MaxDeliver: 5,
})
t, err := redis.New(redis.Config{
Addr: "redis:6379",
BlockDuration: 200 * time.Millisecond,
BatchSize: 10,
MaxLen: 10_000,
})
Cluster configuration
All transports support multi-node clusters via environment variables.
Kafka cluster
KAFKA_BOOTSTRAP_SERVERS=kafka1:9092,kafka2:9092,kafka3:9092
KAFKA_SECURITY_PROTOCOL=sasl_ssl # plaintext | sasl_plaintext | sasl_ssl | ssl
KAFKA_SASL_MECHANISM=scram-sha-256 # plain | scram-sha-256 | scram-sha-512
KAFKA_SASL_USERNAME=myuser
KAFKA_SASL_PASSWORD=mysecret
NATS cluster
NATS_URL=nats://nats1:4222,nats://nats2:4222,nats://nats3:4222
NATS_JETSTREAM=true # enable durable JetStream delivery
NATS_CREDENTIALS_FILE=/etc/nats/app.creds # NGS / Synadia accounts
NATS_TLS=true # TLS (nats+tls:// connections)
NATS_ACK_WAIT_SECONDS=30
NATS_MAX_DELIVER=3 # -1 = unlimited
Redis Streams
REDIS_STREAMS_URL=redis://redis:6379 # falls back to REDIS_URL if unset
REDIS_STREAMS_BLOCK_MS=200 # XREADGROUP block duration
REDIS_STREAMS_BATCH_SIZE=10 # messages per XREADGROUP fetch
REDIS_STREAMS_MAX_LEN=10000 # max stream length (0=unlimited)
Publishing events
Modules receive the EventBus via ModuleDeps.EventBus. Declare typed event definitions once and publish anywhere:
// Declare once in your module (import "github.com/redelay/go-events/typed")
var OrderCreatedEvent = typed.EventDefinition[OrderCreatedPayload]{
Name: "order.created",
EntityType: "order",
Action: "created",
Topic: "order.created",
}
// Publish from a handler — deps.EventBus injected at bootstrap
func (m *OrdersModule) createOrder(ctx context.Context, ...) {
// ... create order in DB ...
msg, err := OrderCreatedEvent.NewMessage(
order.ID,
typed.Actor{Type: typed.ActorTypeUser, ID: userID},
OrderCreatedPayload{OrderID: order.ID, Total: order.Total},
)
if err != nil { ... }
m.bus.Publish(ctx, msg)
}
Always check deps.EventBus != nil before publishing — the EventBus is nil when no transport is configured (e.g. offline IR generation tools).
Consuming events
Register consumer handlers before calling bus.Start:
bus.Register(&modules.ConsumerRegistration{
EventName: "order.created",
Topic: "order.created",
GroupID: "notifications-service",
Handler: func(ctx context.Context, msg *modules.EventMessage) error {
var p OrderCreatedPayload
if err := json.Unmarshal(msg.Payload, &p); err != nil { return err }
return sendConfirmationEmail(ctx, p)
},
})
bus.Start(ctx)
Event filters
Filters let a handler opt out of specific messages without affecting other subscribers in the same group:
import "github.com/redelay/go-events/filter"
bus.Register(&modules.ConsumerRegistration{
Topic: "events",
GroupID: "billing",
Filter: filter.All(
&filter.ActionFilter{EntityType: "invoice", Action: "created"},
filter.Func(func(msg *modules.EventMessage) bool {
return msg.ActorType != "system" // ignore system-generated invoices
}),
).Matches,
Handler: ...,
})
Available filters: HeaderFilter, ActionFilter, RegexFilter, CompositeFilter (All/Any), PrefixFilter, FuncFilter.
Correlation IDs
Correlation IDs chain events across services. The EventBus middleware propagates them automatically:
- Publish side: reads
CorrelationIDFromContext(ctx)and sets it on the outgoing message. - Consume side: injects the message's
correlation_idback intocontext.ContextviaWithCorrelationID.
// In an HTTP handler — inject a trace ID from the request
ctx = eventbus.WithCorrelationID(r.Context(), requestID)
bus.Publish(ctx, msg) // correlation_id header is set automatically
// In a consumer — read the propagated ID
cid := eventbus.CorrelationIDFromContext(ctx) // "req-abc-123"
Email events (built-in definition)
EmailSendPayload and EmailSendEvent are defined in the lightweight go-module-email/events package (it depends only on go-events/typed), so any module can trigger emails without importing the email module itself:
import (
"github.com/redelay/go-events/typed"
emailevents "github.com/redelay/go-module-email/events"
)
msg, _ := emailevents.EmailSendEvent.NewMessage("", typed.SystemActor, emailevents.EmailSendPayload{
To: []string{user.Email},
Template: "welcome",
Data: map[string]any{"name": user.Name},
})
bus.Publish(ctx, msg)
The go-module-email module subscribes to email.send and handles rendering + delivery.