// Package events — шина событий для обмена UI ↔ Core. // // Две независимые шины: доменная (статусы задач, история, трейсы) и логовая // (сырые строки лога для панели «Логи»). Обе построены на одном типе *Bus, // DOMAIN шина блокирующая (гарантия доставки и порядка), логовая — та же, // но с большим буфером, чтобы не тормозить логирование. package events import ( "sync" ) // Event — доменное событие. Закрытый интерфейс: новые типы добавляются // только внутри пакета. type Event interface { _event() } // Bus — широковещательная шина событий (fan-out). // // Publish блокирует вызывающую горутину до тех пор, пока все подписчики не // получат событие (в копию их буфера). Порядок событий для каждого // подписчика сохраняется. Удаление подписчика происходит горутиной-монтируется // close(done), что снимает блокировку Publish. type Bus struct { bufSize int mu sync.RWMutex subs map[*subscriber]struct{} } type subscriber struct { ch chan Event done chan struct{} } // New создаёт шину с буфером bufSize на каждого подписчика. func New(bufSize int) *Bus { if bufSize < 1 { bufSize = 1 } return &Bus{ bufSize: bufSize, subs: make(map[*subscriber]struct{}), } } // Subscribe регистрирует нового подписчика и возвращает канал событий вместе // с функцией отписки. Рекомендуемый паттерн потребления: // // ch, unsub := bus.Subscribe() // defer unsub() // for { // select { // case e := <-ch: // switch ev := e.(type) { ... } // case <-closeCh: // return // } // } // // Канал не закрывается шиной: выход из горутины подписчика организуется через // закрытие канала приложения либо другого сигнала в select. func (b *Bus) Subscribe() (<-chan Event, func()) { s := &subscriber{ ch: make(chan Event, b.bufSize), done: make(chan struct{}), } b.mu.Lock() b.subs[s] = struct{}{} b.mu.Unlock() var once sync.Once unsubscribe := func() { once.Do(func() { close(s.done) b.mu.Lock() delete(b.subs, s) b.mu.Unlock() }) } return s.ch, unsubscribe } // Publish рассылает событие всем подписчикам и блокируется, пока каждый // подписчик либо примет событие (в свой буфер), либо отпишется. Порядок // рассылки стабилен (по списку подписок). Безопасен для Concurrent вызовов. func (b *Bus) Publish(e Event) { b.mu.RLock() subs := make([]*subscriber, 0, len(b.subs)) for s := range b.subs { subs = append(subs, s) } b.mu.RUnlock() for _, s := range subs { select { case s.ch <- e: case <-s.done: } } } // SubscribersCount — число активных подписчиков (для юнит-тестов и отладки). func (b *Bus) SubscribersCount() int { b.mu.RLock() defer b.mu.RUnlock() return len(b.subs) }