Barramento de eventos :id=event-bus
O event bus processa ações pub/sub enfileiradas em uma única goroutine de dispatcher e entrega eventos correspondentes aos channels dos subscribers.
Os exemplos em Go são fragmentos de implementação e extensão. Eles pressupõem um contexto de componentes, logger, handlers e tipos de eventos da aplicação já existentes.
Estrutura de Evento
type Event struct {
System string // Component/module (e.g., "registry", "process")
Kind string // Event type (e.g., "create", "update", "exit")
Path string // Entity identifier
Data any // Payload
Aux any // In-process dispatcher context; not propagated to processes
}
Arquitetura do Bus
flowchart LR
subgraph Publishers
P1[Component]
P2[Component]
end
subgraph Bus
Q[actionQueue]
D[dispatcher goroutine]
S[subscribers map]
end
subgraph Subscribers
S1[chan Event]
S2[chan Event]
end
P1 & P2 -->|enqueue| Q
Q -->|signal| D
D -->|match & deliver| S1 & S2
D <-->|manage| S
O bus armazena estado em uma estrutura simples:
type Bus struct {
subscribers map[SubscriberID]sub
subscriberCounter uint64
maxSubscribers int
actionQueue []action
spareQueue []action
actionMu sync.Mutex
actionReady chan struct{} // buffered=1
closed atomic.Bool
}
Todas as mutações passam pela goroutine do dispatcher, eliminando race conditions sem locking complexo.
Ações
Quatro tipos de ação fluem pela fila:
| Ação | Comportamento |
|---|---|
| Subscribe | Adiciona subscriber ao map, responde no done channel |
| Unsubscribe | Remove subscriber, responde no done channel |
| Send | Entrega evento para subscribers correspondentes |
| Stop | Limpa subscribers, drena fila, sai do loop |
Subscribe e Unsubscribe bloqueiam até que o dispatcher confirme. Send é fire-and-forget. O bus aceita no máximo DefaultMaxSubscribers inscrições (4096 por padrão); inscrições além desse limite falham com ErrSubscribersCapReached.
Subscribe é rejeitado com ErrSubscribersCapReached assim que o barramento atinge DefaultMaxSubscribers (4096) assinaturas ativas.
Subscribe falha imediatamente quando o contexto da inscrição já está cancelado, e novamente no dispatcher se for cancelado antes da decisão de posse ser tomada — o bus nunca assume um canal que não instalou.
Unsubscribe é uma barreira de posse, não uma dica de melhor esforço. Ele retorna apenas depois que o dispatcher confirma, então o chamador pode liberar o canal sabendo que o bus não mantém nenhuma referência de envio em andamento. Quando chega depois de Stop, a confirmação aguarda o dispatcher terminar de entregar o lote que já havia drenado.
Stop é igualmente terminal: um segundo Stop concorrente não retorna antecipadamente pela flag de já fechado, mas aguarda o dispatcher drenar e sair.
Troca de Fila
O dispatcher usa troca de slices para evitar alocações em estado estável:
func (b *Bus) processActions() bool {
b.actionMu.Lock()
actions := b.actionQueue
b.actionQueue = b.spareQueue[:0]
b.spareQueue = nil
b.actionMu.Unlock()
for i := range actions {
// process action
}
clear(actions)
b.actionMu.Lock()
b.spareQueue = actions[:0]
b.actionMu.Unlock()
return true
}
Dois slices alternam: um para processamento, um para novas chegadas. O channel actionReady tem buffer de 1, então sinalizar nunca bloqueia e múltiplos enqueues coalescem em um wakeup.
Correspondência de padrões
Inscrições compilam padrões uma vez no momento da inscrição:
type sub struct {
subID SubscriberID
ctx context.Context
system *wildcard.Wildcard
kind *wildcard.Wildcard
eventCh chan<- Event
}
O pacote wildcard oferece quatro tipos de padrão:
| Padrão | Corresponde |
|---|---|
registry |
Apenas match exato |
* |
Qualquer segmento único |
** |
Zero ou mais segmentos |
(a|b) |
Alternação dentro do segmento |
Padrões dividem em . então registry.* corresponde registry.create mas não registry.entry.create. O padrão registry.** corresponde todos os três de registry, registry.create, e registry.entry.create.
Entrega de Eventos
Durante processamento de Send, o dispatcher itera subscribers:
for id, s := range b.subscribers {
if s.system != nil && !s.system.Match(a.event.System) {
continue
}
if s.kind != nil && !s.kind.Match(a.event.Kind) {
continue
}
select {
case <-a.ctx.Done():
goto cleanup
case <-s.ctx.Done():
expiredSubs = append(expiredSubs, id)
case s.eventCh <- a.event:
}
}
Se o contexto de um subscriber for cancelado, ele é marcado para remoção durante aquela passagem de entrega. O contexto do evento também pode cancelar entrega no meio da iteração.
Ponte de Processo Lua
O dispatcher de eventos faz ponte de eventos Go para processos Lua. Ele se inscreve uma vez em todos os eventos ("**") e roteia internamente baseado em inscrições de processos:
type Dispatcher struct {
bus event.Bus
node relay.Node
subID SubscriberID
eventC chan event.Event
mu sync.RWMutex
subs map[string]*subscription // topic -> subscription
}
Quando um processo Lua se inscreve via events.subscribe(), o dispatcher armazena o padrão e PID alvo. Eventos correspondentes são empacotados e enviados via relay:
func (d *Dispatcher) routeEvent(evt event.Event) {
d.mu.RLock()
defer d.mu.RUnlock()
for _, sub := range d.subs {
if !matchPattern(sub.system, evt.System) {
continue
}
if sub.kind != "" && sub.kind != "*" && !matchPattern(sub.kind, evt.Kind) {
continue
}
data := map[string]any{
"system": evt.System,
"kind": evt.Kind,
"path": evt.Path,
}
if evt.Data != nil {
data["data"] = evt.Data
}
pkg := relay.NewPackage(pid.PID{}, sub.pid, sub.topic, payload.New(data))
d.node.Send(pkg)
}
}
Tipos Auxiliares
Subscriber
Encapsula inscrição de channel com callback:
handler, err := eventbus.NewSubscriber(ctx, bus, "registry", "entry.*",
func(evt Event) {
// handle
})
if err != nil {
return err
}
defer handler.Close()
Cria duas goroutines: uma lê eventos e chama o handler, outra aguarda cancelamento de contexto para desinscrição.
EventRouter
Gerencia múltiplos handlers com ciclo de vida centralizado:
router, err := eventbus.StartRouter(ctx, bus,
WithHandlers(handler1, handler2),
WithLogger(log))
if err != nil {
return err
}
defer router.Stop()
Cada handler implementa Pattern() e Handle(). O router cria um Subscriber para cada e fecha todos em Stop.
AwaitService
Requisição-resposta sobre pub/sub. Mantém uma única inscrição por par (system, kind) e roteia eventos para os waiters por Path:
svc := eventbus.NewAwaitService(bus)
svc.Start(ctx)
defer svc.Stop()
waiter, _ := svc.Prepare(ctx, "test", "response.(accept|reject)", "test/path", 5*time.Second)
defer waiter.Close()
bus.Send(ctx, triggeringEvent)
result := waiter.Wait() // retorna AwaitResult{Event, Accepted, Error}
Prepare registra o waiter antes de o evento acionador ser enviado, evitando a race em que a resposta chega antes de a espera ser registrada. Wait bloqueia até que um evento com Path correspondente chegue ou o timeout (padrão DefaultAwaitTimeout, 30s, quando não positivo) expire. Accepted é true quando o kind do evento é accept, *.accept ou *.accepted; caso contrário o kind é tratado como rejeição e qualquer error em Data aparece como Error. A conveniência Await(ctx, system, kind, path, timeout) combina Prepare e Wait. A infraestrutura de boot registra um AwaitService no contexto (event.GetAwaitService).
Encerramento
Stop()atomicamente define flag closed e enfileira ação Stop- Dispatcher limpa mapa de subscribers
- Ações restantes na fila são drenadas:
- Requisições Subscribe recebem erro "bus is closed"
- Requisições Unsubscribe completam imediatamente
- Eventos Send são descartados
- WaitGroup completa
Consulte também
- Registro — Principal produtor de eventos
- Despacho de comandos — Roteamento de processos para handlers