AsyncAPI
The asyncapi module is the event-side counterpart to the openapi module. It walks every
registered module at startup, collects all declared events and consumers, and emits a complete
AsyncAPI 2.6 specification — with typed payload schemas, Kafka bindings, and channel
descriptions. An interactive AsyncAPI Studio UI is served alongside the raw YAML.
How it works
Module.Events() → channels (publish side)
Module.Consumers() → channels (subscribe side)
ir.Event.Payload → components/schemas (inline ir.Packet)
Module.AsyncAPISchemas() → components/schemas (reflected Go structs)
ir.Event.PayloadRef → links a channel message to a named schema
↓
asyncapi.Export()
↓
/asyncapi.json ← spec (YAML, application/yaml)
/asyncapi ← AsyncAPI Studio UI
Events declared on any module via EventsProvider are automatically included — no
additional wiring needed. The spec is generated lazily on first request and cached.
Quick start
import (
_ "github.com/redelay/go-framework/modules/asyncapi"
)
Blank-import registers the module. On startup your service gains:
| Path | Description |
|---|---|
/asyncapi.json | AsyncAPI 2.6 spec (YAML) |
/asyncapi | AsyncAPI Studio — interactive event browser |
Declaring events
Implement EventsProvider on your module. The same Events() return value drives both
the runtime consumer routing and the AsyncAPI spec:
type Module struct{}
func (m *Module) Name() string { return "orders" }
// Events — published by this module
func (m *Module) Events() []*ir.Event {
return []*ir.Event{
{
Name: "order.created",
EntityType: "order",
Action: "created",
Topic: "order.created",
Description: "Fired when a new order is placed.",
Payload: &ir.Packet{
ID: "OrderCreatedPayload",
Name: "OrderCreatedPayload",
Fields: []*ir.Field{
{Name: "order_id", Type: ir.FieldTypeUUID, Required: true},
{Name: "customer_id", Type: ir.FieldTypeUUID, Required: true},
{Name: "total", Type: ir.FieldTypeFloat, Required: true},
{Name: "currency", Type: ir.FieldTypeString, Required: true},
{Name: "items_count", Type: ir.FieldTypeInt, Required: true},
},
},
},
{
Name: "order.cancelled",
EntityType: "order",
Action: "cancelled",
Topic: "order.cancelled",
Description: "Fired when an order is cancelled by the customer or system.",
},
}
}
// Consumers — events this module subscribes to
func (m *Module) Consumers() []*modules.ConsumerRegistration {
return []*modules.ConsumerRegistration{
{
EventName: "payment.completed",
Topic: "payment.completed",
GroupID: "orders",
Handler: m.handlePaymentCompleted,
},
}
}
Generated AsyncAPI spec
For a service with the orders module above plus the built-in auth and users modules, the generated spec looks like:
asyncapi: "2.6.0"
info:
title: My Service
version: "1.0.0"
channels:
order.created:
description: Fired when a new order is placed.
publish:
operationId: order.created
summary: order.created
message:
name: order.created
contentType: application/json
payload:
$ref: '#/components/schemas/OrderCreatedPayload'
order.cancelled:
description: Fired when an order is cancelled by the customer or system.
publish:
operationId: order.cancelled
message:
name: order.cancelled
contentType: application/json
payment.completed:
subscribe:
operationId: consume_payment.completed
message:
name: payment.completed
contentType: application/json
bindings:
kafka:
groupId: orders
auth.login:
publish:
operationId: auth.login
message:
name: auth.login
contentType: application/json
users.user_created:
publish:
operationId: users.user_created
message:
name: users.user_created
contentType: application/json
payload:
$ref: '#/components/schemas/UserCreatedPayload'
components:
schemas:
OrderCreatedPayload:
type: object
required: [order_id, customer_id, total, currency, items_count]
properties:
order_id: { type: string, format: uuid }
customer_id: { type: string, format: uuid }
total: { type: number }
currency: { type: string }
items_count: { type: integer }
UserCreatedPayload:
type: object
required: [user_id, email]
properties:
user_id: { type: string }
email: { type: string }
Channel naming defaults to entity_type.action when no explicit Topic is set on the event.
Payload schema rules
The ir.Packet attached to ir.Event.Payload becomes a components/schemas entry and is
referenced via $ref in the channel's message definition. If Payload is nil, the
channel's message has an open (unconstrained) schema.
Field type → JSON Schema:
ir.FieldType | JSON Schema |
|---|---|
FieldTypeString | string |
FieldTypeInt | integer |
FieldTypeFloat | number |
FieldTypeBool | boolean |
FieldTypeUUID | string + format: uuid |
FieldTypeDateTime | string + format: date-time |
FieldTypeObjectID | string + format: objectid |
FieldTypeObject | object (nested Fields as properties) |
FieldTypeArray | array |
Required: true on a field adds it to the schema's required array.
Reflected Go payload structs (AsyncAPISchemasProvider)
Declaring the payload inline as an ir.Packet duplicates a struct you usually
already have in Go. The recommended path for typed event payloads is to reflect
the Go struct directly instead:
- Implement
AsyncAPISchemasProvideron the module, returning the payload structs (the AsyncAPI counterpart ofSchemasProviderfor OpenAPI):go// modules.AsyncAPISchemasProvider func (m *Module) AsyncAPISchemas() []any { return []any{OrderPaidPayload{}, OrderEventPayload{}} }
Each struct is reflected via the same JSON-Schema engine the OpenAPI side uses (honoursjson+validatetags, nested structs, slices, maps) and lands incomponents/schemaskeyed by its Go type name. Nested named structs are registered too and referenced with$ref. - Link each event to its schema in
module.yamlwithpayload_ref:yamlevents: - topic: order.paid description: An order reached fully paid. payload_ref: OrderPaidPayload
The channel's message payload then$refs#/components/schemas/OrderPaidPayload.payload_refresolves even when the name is not a declaredir.Packet— it points at whatever landed incomponents/schemas, including a reflected struct.
Prefer payload_ref + AsyncAPISchemas() for real Go payloads (single source
of truth, full nesting); use an inline ir.Packet only for schemas with no Go
type. A declared packet of the same name always wins over a reflected struct
(first-wins dedupe).
AsyncAPI Studio UI
The /asyncapi endpoint serves a self-contained HTML page powered by
@asyncapi/react-component 3.x (CDN build).
The Studio renders:
- Sidebar — channel list with
PUB/SUBbadges - Info panel — service title, version, description
- Operations — one per event, expandable with message schema
- Schema browser —
components/schemaswith property types and required fields
┌──────────────────────────────────────────────────────────┐
│ AsyncAPI Studio — My Service 1.0.0 │
├──────────────────────────────────────────────────────────┤
│ Channels │ order.created │
│ │ │
│ PUB order.created │ Fired when a new order is placed. │
│ PUB order.cancelled │ │
│ SUB payment.completed│ Message │
│ PUB auth.login │ OrderCreatedPayload │
│ PUB users.user_created│ order_id string (uuid) ✓ │
│ │ customer_id string (uuid) ✓ │
│ │ total number ✓ │
│ │ currency string ✓ │
│ │ items_count integer ✓ │
└──────────────────────────────────────────────────────────┘
Cross-service schema validation
Both Go services can load any service's AsyncAPI spec at startup and validate every incoming event payload against the declared schema — catching contract drift at runtime:
import "github.com/redelay/go-events/schema"
// Auto-loaded when ASYNCAPI_URL is set (validates all incoming events)
// ASYNCAPI_URL=http://payment-service/asyncapi.json
// Manual validation:
reg, err := schema.LoadFromURL("http://payment-service/asyncapi.json")
if err != nil {
log.Fatal(err)
}
// In your consumer handler:
func (m *Module) handlePaymentCompleted(ctx context.Context, msg *modules.EventMessage) error {
if err := reg.Validate("payment.completed", msg.Payload); err != nil {
return fmt.Errorf("schema violation: %w", err)
}
// ...
}
Referencing from FlowDSL
FlowDSL packet schemas can reference your live AsyncAPI spec instead of duplicating schema definitions — keeping event contracts in sync across all consumers:
# my-pipeline.flowdsl.yaml
components:
packets:
OrderCreatedPayload:
$ref: "https://api.myservice.com/asyncapi.json#/components/schemas/OrderCreatedPayload"
PaymentCompletedPayload:
$ref: "https://api.payments.com/asyncapi.json#/components/schemas/PaymentCompletedPayload"
nodes:
ReceiveOrder:
operationId: receive_order
kind: source
ProcessPayment:
operationId: process_payment
kind: action
edges:
- from: ReceiveOrder
to: ProcessPayment
delivery:
mode: durable
packet: OrderCreatedPayload
See FlowDSL Integration for the complete wiring guide.
Modspec integration
The modspec browser at /modules links
directly to /asyncapi in its top navigation bar. Each module card shows an Events tab
and a Consumers tab with the full list of topics and group IDs — a quick cross-reference
before opening the full Studio UI.
Programmatic spec generation
Generate or import AsyncAPI specs outside of HTTP context — useful for contract testing, schema registries, or publishing to a broker:
import (
"github.com/redelay/go-framework/core/asyncapi"
"github.com/redelay/go-framework/core/ir"
)
// Export IR document → AsyncAPI 2.6 YAML
doc := &ir.Document{ /* ... */ }
err := asyncapi.Export(doc, os.Stdout)
// Import AsyncAPI YAML → IR document
doc, err := asyncapi.ImportYAML(reader)
The core/asyncapi package is used internally by the asyncapi module but can be called
directly for pipelines that need to generate or transform specs independently of an HTTP
server.
Configuration
| Variable | Default | Description |
|---|---|---|
ASYNCAPI_UI | studio | UI provider (studio or disabled) |
ASYNCAPI_SPEC_PATH | /asyncapi.json | URL path for the YAML spec |
ASYNCAPI_UI_PATH | /asyncapi | URL path for the Studio UI |
ASYNCAPI_URL | — | Remote AsyncAPI URL for schema validation on consume |
APP_NAME | redelay | Value of info.title in the spec |
APP_VERSION | 1.0.0 | Value of info.version in the spec |