Concepts

Transports

How Redelay decouples event publishing and consuming from any specific broker.

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

Rendering diagram...

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:

go
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

TransportPackageDeliveryCluster supportUse case
Memorygo-events/transport/memoryIn-process—Unit tests, no infrastructure
Kafkago-events/transport/kafkaDurable streamMulti-broker + SASL/TLSProduction; highest throughput
NATS corego-events/transport/natsEphemeral pub/subCluster routesLow-latency, fire-and-forget
NATS JetStreamgo-events/transport/natsDurable streamCluster + super-clusterDurable with built-in replay
Redis Streamsgo-events/transport/redisDurable streamRedis Cluster / SentinelLightweight; no separate broker

Delivery mode comparison

CharacteristicMemoryKafkaNATS coreNATS JetStreamRedis Streams
Persistence✗✓✗✓✓
Consumer groups✓✓✗✓✓
At-least-once✓*✓✗✓✓
Replay✗✓✗✓✓
Backpressure✗✓✗✓✗
Infra requiredNoneKafkaNATSNATS + JetStreamRedis

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

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

goKafka
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))

Explicit configuration

Pass a Config struct directly for programmatic control:

goKafka
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",
    },
})

Cluster configuration

All transports support multi-node clusters via environment variables.

Kafka cluster

shell
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

shell
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

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

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

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

go
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_id back into context.Context via WithCorrelationID.
go
// 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:

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