Adding Events
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
- Go module scaffold in place (see Creating a Go Module)
go-eventsadded to your app's dependenciesEventBuswired at bootstrap (see Creating a Go App)
1. Define event types
Create events.go alongside module.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:
| Field | Convention |
|---|---|
Name | {entity}.{action} — matches Topic for simplicity |
EntityType | Lowercase singular noun (order, user, invoice) |
Action | Past-tense verb (created, updated, deleted, completed) |
Topic | Kafka/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:
// modules.AsyncAPISchemasProvider — reflects the struct into
// /asyncapi.json#/components/schemas/OrderCreatedPayload
func (m *Module) AsyncAPISchemas() []any {
return []any{OrderCreatedPayload{}}
}
# 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:
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.
// 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:
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:
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:
{
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:
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
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:
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
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
| Step | What to do |
|---|---|
| 1 | Create events.go with typed.EventDefinition variables |
| 2 | Pass deps.EventBus from factory → service |
| 3 | Call eventBus.Publish(ctx, msg) after successful writes; guard != nil |
| 4 | Implement Events() and Consumers() on the module for IR + runtime wiring |
| 5 | Test 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.