package syncengine import ( "log/slog" "sync" ) type EventBus struct { subscribers map[int64]map[chan Event]struct{} mu sync.RWMutex bufferSize int globalSubs []globalSub } type globalSub struct { ch chan Event } func NewEventBus(bufferSize int) *EventBus { return &EventBus{ subscribers: make(map[int64]map[chan Event]struct{}), bufferSize: bufferSize, globalSubs: nil, } } 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()) { ch := make(chan Event, eb.bufferSize) eb.mu.Lock() eb.globalSubs = append(eb.globalSubs, globalSub{ch: ch}) eb.mu.Unlock() return ch, func() { eb.mu.Lock() for i, s := range eb.globalSubs { if s.ch == ch { eb.globalSubs = append(eb.globalSubs[:i], eb.globalSubs[i+1:]...) break } } eb.mu.Unlock() 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) } } } for _, sub := range eb.globalSubs { select { case sub.ch <- evt: default: slog.Warn("global event subscriber buffer 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) } }