package syncengine import ( "log/slog" "sync" ) type EventBus struct { subscribers map[int64]map[chan Event]struct{} mu sync.RWMutex global chan Event bufferSize int } func NewEventBus(bufferSize int) *EventBus { return &EventBus{ subscribers: make(map[int64]map[chan Event]struct{}), global: make(chan Event, bufferSize), bufferSize: bufferSize, } } func (eb *EventBus) Subscribe(jobID int64) (chan Event, func()) { eb.mu.Lock() defer eb.mu.Unlock() if eb.subscribers[jobID] == nil { eb.subscribers[jobID] = make(map[chan Event]struct{}) } ch := make(chan Event, eb.bufferSize) eb.subscribers[jobID][ch] = struct{}{} unsubscribe := func() { eb.mu.Lock() defer eb.mu.Unlock() if subs, ok := eb.subscribers[jobID]; ok { delete(subs, ch) if len(subs) == 0 { delete(eb.subscribers, jobID) } } close(ch) } return ch, unsubscribe } func (eb *EventBus) SubscribeGlobal() (chan Event, func()) { eb.mu.RLock() ch := make(chan Event, eb.bufferSize) eb.mu.RUnlock() go func() { for evt := range eb.global { select { case ch <- evt: default: slog.Warn("global event subscriber buffer full, dropping event", "type", evt.Type) } } close(ch) }() return ch, func() { close(ch) } } func (eb *EventBus) Publish(evt Event) { eb.mu.RLock() defer eb.mu.RUnlock() if subs, ok := eb.subscribers[evt.JobID]; ok { for ch := range subs { select { case ch <- evt: default: slog.Warn("job event subscriber buffer full, dropping event", "job_id", evt.JobID) } } } select { case eb.global <- evt: default: slog.Warn("global event bus full, dropping event", "type", evt.Type) } } func (eb *EventBus) CloseJobChannels(jobID int64) { eb.mu.Lock() defer eb.mu.Unlock() if subs, ok := eb.subscribers[jobID]; ok { for ch := range subs { close(ch) } delete(eb.subscribers, jobID) } }