package eventbus import ( "context" "sync" ) type Event interface { Name() string } type Handler func(ctx context.Context, event Event) error type Bus struct { mu sync.RWMutex handlers map[string][]Handler async bool queue chan asyncEvent } type asyncEvent struct { ctx context.Context event Event } func New() *Bus { return &Bus{ handlers: make(map[string][]Handler), queue: make(chan asyncEvent, 1024), } } func (b *Bus) Subscribe(eventName string, handler Handler) { b.mu.Lock() defer b.mu.Unlock() b.handlers[eventName] = append(b.handlers[eventName], handler) } func (b *Bus) Publish(ctx context.Context, event Event) error { b.mu.RLock() handlers := b.handlers[event.Name()] b.mu.RUnlock() for _, h := range handlers { if err := h(ctx, event); err != nil { return err } } return nil } func (b *Bus) PublishAsync(ctx context.Context, event Event) { b.queue <- asyncEvent{ctx: ctx, event: event} } func (b *Bus) StartWorkers(n int) { for i := 0; i < n; i++ { go func() { for ae := range b.queue { b.mu.RLock() handlers := b.handlers[ae.event.Name()] b.mu.RUnlock() for _, h := range handlers { _ = h(ae.ctx, ae.event) } } }() } } func (b *Bus) Stop() { close(b.queue) }