Bus de eventos :id=event-bus
El bus de eventos procesa acciones pub/sub encoladas en una goroutine despachadora y entrega los eventos coincidentes a los canales de los suscriptores.
Los fragmentos Go son partes de implementación y extensión. Suponen un contexto de componentes, logger, handlers y tipos de eventos de la aplicación ya existentes.
Estructura 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
}
Arquitectura del 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
El bus almacena estado en una estructura simple:
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 las mutaciones pasan por la goroutine despachadora, lo que elimina las condiciones de carrera sin bloqueos complejos.
Acciones
Cuatro tipos de acciones fluyen a través de la cola:
| Acción | Comportamiento |
|---|---|
| Subscribe | Agrega subscriber al mapa, responde en canal done |
| Unsubscribe | Remueve subscriber, responde en canal done |
Send |
Entrega el evento a los suscriptores coincidentes |
Stop |
Limpia los suscriptores, drena la cola y sale del bucle |
Subscribe y Unsubscribe bloquean hasta que el despachador confirma. Send envía sin esperar respuesta. El bus acepta como máximo DefaultMaxSubscribers suscripciones (4096 de forma predeterminada); las que superan el límite fallan con ErrSubscribersCapReached.
Subscribe se rechaza con ErrSubscribersCapReached una vez que el bus mantiene DefaultMaxSubscribers (4096) suscripciones activas.
Subscribe falla inmediatamente cuando el contexto de la suscripcion ya esta cancelado, y de nuevo en el dispatcher si se cancela antes de que se tome la decision de propiedad — el bus nunca toma un canal que no instalo.
Unsubscribe es una barrera de propiedad, no una sugerencia de mejor esfuerzo. Retorna solo despues de que el dispatcher confirma, de modo que el llamador puede liberar el canal sabiendo que el bus no mantiene ninguna referencia de envio en vuelo. Cuando llega despues de Stop, la confirmacion espera a que el dispatcher termine de entregar el lote que ya dreno.
Stop es igualmente terminal: un segundo Stop concurrente no retorna temprano por la bandera de ya-cerrado, sino que espera a que el dispatcher drene y salga.
Intercambio de Cola
El despachador intercambia segmentos para evitar asignaciones en estado estable:
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
}
Se alternan dos segmentos: uno para el procesamiento y otro para las nuevas llegadas. El canal actionReady tiene un búfer de 1, por lo que la señalización nunca bloquea y múltiples operaciones de encolado se agrupan en una sola activación.
Coincidencia de patrones :id=pattern-matching
Las suscripciones compilan patrones una vez en tiempo de subscribe:
type sub struct {
subID SubscriberID
ctx context.Context
system *wildcard.Wildcard
kind *wildcard.Wildcard
eventCh chan<- Event
}
El paquete de comodines admite cuatro tipos de patrón:
| Patrón | Matchea |
|---|---|
registry |
Solo match exacto |
* |
Cualquier segmento único |
** |
Cero o más segmentos |
(a|b) |
Alternación dentro de segmento |
Los patrones se dividen en . así que registry.* matchea registry.create pero no registry.entry.create. El patrón registry.** matchea los tres: registry, registry.create, y registry.entry.create.
Entrega de Eventos
Durante el procesamiento de Send, el despachador recorre los suscriptores:
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:
}
}
Si el contexto de un subscriber es cancelado, se marca para remoción durante ese pase de entrega. El contexto del evento también puede cancelar entrega a mitad de iteración.
Puente de procesos Lua :id=bridge-de-proceso-lua
El despachador de eventos conecta eventos de Go con procesos de Lua. Se suscribe una vez a todos los eventos ("**") y enruta internamente según las suscripciones de los procesos:
type Dispatcher struct {
bus event.Bus
node relay.Node
subID SubscriberID
eventC chan event.Event
mu sync.RWMutex
subs map[string]*subscription // topic -> subscription
}
Cuando un proceso Lua se suscribe mediante events.subscribe(), el despachador almacena el patrón y el PID de destino. Los eventos coincidentes se empaquetan y envían mediante el relé:
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 :id=tipos-helper
Subscriber
Envuelve suscripción de canal con callback:
handler, err := eventbus.NewSubscriber(ctx, bus, "registry", "entry.*",
func(evt Event) {
// handle
})
if err != nil {
return err
}
defer handler.Close()
Genera dos goroutines: una lee eventos y llama al handler, otra espera cancelación de contexto para desuscribir.
EventRouter
Gestiona múltiples handlers con 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() y Handle(). El router crea un Subscriber para cada uno y cierra todos en Stop.
AwaitService
Solicitud-respuesta sobre pub/sub. Mantiene una única suscripción por cada par (system, kind) y enruta los eventos a los 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() // devuelve AwaitResult{Event, Accepted, Error}
Prepare registra el waiter antes de que se envíe el evento disparador, evitando la carrera en la que la respuesta llega antes de que la espera esté registrada. Wait bloquea hasta que llega un evento con Path coincidente o expira el timeout (por defecto DefaultAwaitTimeout, 30s, cuando no es positivo). Accepted es true cuando el kind del evento es accept, *.accept o *.accepted; en caso contrario el kind se trata como rechazo y cualquier error en Data se expone como Error. La función de conveniencia Await(ctx, system, kind, path, timeout) combina Prepare y Wait. La infraestructura de arranque registra un AwaitService en el contexto (event.GetAwaitService).
Apagado :id=shutdown
Stop()atómicamente establece flag closed y encola acción Stop- El despachador limpia el mapa de suscriptores
- Acciones restantes en cola son drenadas:
- Solicitudes Subscribe obtienen error "bus is closed"
- Solicitudes Unsubscribe completan inmediatamente
- Los eventos
Sendse descartan
- WaitGroup completa
Ver También
- Registry - Productor principal de eventos
- Command Dispatch - Routing proceso-a-handler