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

  1. Stop() atomicamente define flag closed e enfileira ação Stop
  2. Dispatcher limpa mapa de subscribers
  3. 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
  4. WaitGroup completa

Consulte também