Guides

Adding Events

How to publish events, consume events, and test event-driven behaviour in a Redelay Go module.

Adding Events to a Go Module

This guide walks through the full event integration for a Go module: defining typed events, publishing from handlers, consuming events from other modules, and testing without Kafka.

Prerequisites

1. Define event types

Create events.go alongside module.go:

go
package orders

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

// OrderCreatedPayload is the canonical payload for order.created events.
// Defined here so any module can reference it without importing this package.
type OrderCreatedPayload struct {
    OrderID  string  `json:"order_id"`
    Total    float64 `json:"total"`
    Currency string  `json:"currency"`
}

// OrderCreatedEvent is the compile-time-safe event descriptor.
// Call NewMessage() to produce a ready-to-publish EventMessage.
var OrderCreatedEvent = typed.EventDefinition[OrderCreatedPayload]{
    Name:       "order.created",
    EntityType: "order",
    Action:     "created",
    Topic:      "order.created",
}

Rules for typed.EventDefinition:

FieldConvention
Name{entity}.{action} — matches Topic for simplicity
EntityTypeLowercase singular noun (order, user, invoice)
ActionPast-tense verb (created, updated, deleted, completed)
TopicKafka/Redis/NATS topic name. Mirror Name unless fan-out is needed

Surface the payload in AsyncAPI

So reaction flows and the AsyncAPI Studio see the real fields of your event (not an opaque envelope), publish the payload struct and link it to the event:

go
// modules.AsyncAPISchemasProvider — reflects the struct into
// /asyncapi.json#/components/schemas/OrderCreatedPayload
func (m *Module) AsyncAPISchemas() []any {
    return []any{OrderCreatedPayload{}}
}
yaml
# module.yaml — link the event to the reflected schema
events:
  - topic: order.created
    description: An order was created.
    payload_ref: OrderCreatedPayload

See AsyncAPI → Reflected Go payload structs for the full contract.

2. Publish events from service methods

Pass modules.EventBus into the service:

go
type Service struct {
    crud     *crud.MongoCRUD[*Order]
    eventBus modules.EventBus
    logger   *zap.Logger
}

func NewService(col *mongo.Collection, bus modules.EventBus, logger *zap.Logger) *Service {
    return &Service{crud: crud.NewMongoCRUD[*Order](col), eventBus: bus, logger: logger}
}

func (s *Service) Create(ctx context.Context, input *CreateOrderInput) (*Order, error) {
    order, err := s.crud.Create(ctx, &Order{
        Total:    input.Total,
        Currency: input.Currency,
    })
    if err != nil {
        return nil, err
    }

    if s.eventBus != nil {
        msg, err := OrderCreatedEvent.NewMessage(
            order.ID.Hex(),
            typed.Actor{Type: typed.ActorTypeUser, ID: input.UserID},
            OrderCreatedPayload{OrderID: order.ID.Hex(), Total: order.Total, Currency: order.Currency},
        )
        if err != nil {
            s.logger.Error("failed to build order.created message", zap.Error(err))
        } else if err := s.eventBus.Publish(ctx, msg); err != nil {
            s.logger.Error("failed to publish order.created", zap.Error(err))
        }
    }

    return order, nil
}

Always guard eventBus != nil — modules must work in apps without a transport configured.

Correlation ID propagation

Pass the request ctx directly to Publish. The CorrelationPublishMiddleware (installed by default in eventbus.New()) reads the correlation ID from context and injects it into the message header automatically.

go
// Set a correlation ID on an incoming HTTP request context:
ctx = eventbus.WithCorrelationID(ctx, r.Header.Get("X-Request-ID"))
// Now any Publish call within this ctx will carry the same correlation_id.

3. Consume events from other modules

Implement EventsProvider in module.go:

go
import (
    "encoding/json"

    "github.com/redelay/go-framework/core/ir"
    "github.com/redelay/go-framework/modules"
)

// Events returns IR metadata for modspec and AsyncAPI spec generation.
func (m *Module) Events() []*ir.Event {
    return []*ir.Event{
        {Name: "order.created", EntityType: "order", Action: "created"},
        // Consumed events can also be listed here for documentation:
        {Name: "payment.completed", EntityType: "payment", Action: "completed"},
    }
}

// Consumers returns runtime subscriptions. Called before bus.Start().
func (m *Module) Consumers() []*modules.ConsumerRegistration {
    return []*modules.ConsumerRegistration{
        {
            EventName: "payment.completed",
            Topic:     "payment.completed",
            GroupID:   "orders",
            Handler:   m.handlePaymentCompleted,
        },
    }
}

func (m *Module) handlePaymentCompleted(ctx context.Context, msg *modules.EventMessage) error {
    var p PaymentCompletedPayload
    if err := json.Unmarshal(msg.Payload, &p); err != nil {
        return fmt.Errorf("orders: unmarshal payment.completed: %w", err)
    }
    return m.Service.MarkAsPaid(ctx, p.OrderID)
}

Consumer group naming

Use your module name as GroupID. This ensures exactly-one delivery per module across scaled-out replicas:

text
GroupID: "orders"   // all replicas of the orders module share this group
GroupID: "notifications"
GroupID: "invoices"

Filtering

Register a Filter to skip messages that don't apply to this handler:

go
{
    Topic:   "entity.updated",
    GroupID: "orders",
    Filter: func(msg *modules.EventMessage) bool {
        return msg.EntityType == "order"
    },
    Handler: m.handleOrderUpdated,
}

Or use the filter package from go-events for reusable filters:

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

f := filter.All(
    &filter.ActionFilter{EntityType: "order", Action: "updated"},
    &filter.HeaderFilter{Key: "env", Value: "prod"},
)

{
    Topic:   "entity.updated",
    GroupID: "orders",
    Filter:  f.Matches,
    Handler: m.handleOrderUpdated,
}

4. Wire EventBus into the module factory

go
func init() {
    modules.RegisterFactory("orders", func(deps modules.ModuleDeps) (modules.Module, error) {
        cfg := DefaultConfig()
        return NewModule(deps.DB, cfg, deps.Logger, deps.EventBus), nil
    })
}

func NewModule(db *mongo.Database, cfg *Config, logger *zap.Logger, bus modules.EventBus) *Module {
    return &Module{db: db, cfg: cfg, logger: logger, eventBus: bus}
}

func (m *Module) Startup(ctx context.Context) error {
    col := m.db.Collection(m.cfg.Collection)
    m.Service = NewService(col, m.eventBus, m.logger)
    return m.Service.EnsureIndexes(ctx)
}

5. Testing without Kafka

Use memory.New() from go-events/transport/memory — no broker required:

go
package orders_test

import (
    "context"
    "testing"

    "github.com/redelay/go-events/eventbus"
    "github.com/redelay/go-events/transport/memory"
    "github.com/redelay/go-framework/modules"
    "yourorg/yourapp/modules/orders"
)

func TestOrderCreatedEventPublished(t *testing.T) {
    tr := memory.New()
    bus := eventbus.New(tr)

    var received *modules.EventMessage
    bus.Register(&modules.ConsumerRegistration{
        Topic:   "order.created",
        GroupID: "test",
        Handler: func(_ context.Context, msg *modules.EventMessage) error {
            received = msg
            return nil
        },
    })
    if err := bus.Start(context.Background()); err != nil {
        t.Fatal(err)
    }

    svc := orders.NewService(testCollection(t), bus, zaptest.NewLogger(t))
    _, err := svc.Create(context.Background(), &orders.CreateOrderInput{
        Total:    99.00,
        Currency: "USD",
        UserID:   "user-1",
    })
    if err != nil {
        t.Fatalf("Create: %v", err)
    }

    if received == nil {
        t.Fatal("expected order.created event to be published")
    }
    if received.EntityType != "order" || received.Action != "created" {
        t.Errorf("unexpected event: %s.%s", received.EntityType, received.Action)
    }
}

Testing consumers

go
func TestHandlePaymentCompleted(t *testing.T) {
    tr := memory.New()
    bus := eventbus.New(tr)

    mod := orders.NewModule(testDB(t), orders.DefaultConfig(), zaptest.NewLogger(t), bus)
    // Register module consumers
    for _, reg := range mod.Consumers() {
        bus.Register(reg)
    }
    if err := bus.Start(context.Background()); err != nil {
        t.Fatal(err)
    }
    if err := mod.Startup(context.Background()); err != nil {
        t.Fatal(err)
    }

    // Publish a payment.completed event
    _ = bus.Publish(context.Background(), &modules.EventMessage{
        EntityType: "payment",
        Action:     "completed",
        Headers:    map[string]string{"topic": "payment.completed"},
        Payload:    []byte(`{"order_id":"order-123","amount":99.00}`),
    })

    order, _ := mod.Service.GetByID(context.Background(), "order-123")
    if order.Status != "paid" {
        t.Errorf("expected order to be marked as paid")
    }
}

Summary

StepWhat to do
1Create events.go with typed.EventDefinition variables
2Pass deps.EventBus from factory → service
3Call eventBus.Publish(ctx, msg) after successful writes; guard != nil
4Implement Events() and Consumers() on the module for IR + runtime wiring
5Test with memory.New() + eventbus.New(tr) — no Kafka needed

See Choosing a Transport for the full transport API, and go-events reference for middleware and filter documentation.