- Hub слушает *Bus в собственной горутине и вызывает зарегистрированные обработчики по типу события (On[T]); порядок сохраняется. - Мост к UI: колбэки выполняются в горутине Hub → внутри можно переложить работу на поток Fyne (fyne.Do) или thread-safe binding. - Close останавливает горутину и отписывается; Start идемпотентен. - Юнит-тесты: доставка по типу, игнор посторонних типов, несколько обработчиков одного типа, остановка после Close.
96 lines
2.8 KiB
Go
96 lines
2.8 KiB
Go
package events
|
|
|
|
import (
|
|
"reflect"
|
|
"sync"
|
|
)
|
|
|
|
// Handler — колбэк-обработчик события конкретного типа.
|
|
type Handler[T Event] func(e T)
|
|
|
|
// Hub — типизированный подписчик шины.
|
|
//
|
|
// Слушает *Bus в собственной горутине и вызывает зарегистрированные обработчики
|
|
// для событий соответствующих типов. Порядок обработки сохраняется (порядок
|
|
// шины). Обработчики одного типа вызываются в порядке регистрации.
|
|
//
|
|
// Это мост между шиной и UI-потоком: колбэки выполняются в горутине Hub, поэтому
|
|
// внутри них нужно либо перекладывать работу на поток Fyne (fyne.Do), либо
|
|
// пользоваться только thread-safe структурами (binding).
|
|
type Hub struct {
|
|
bus *Bus
|
|
done chan struct{}
|
|
once sync.Once
|
|
ch <-chan Event
|
|
unsub func()
|
|
|
|
mu sync.Mutex
|
|
handlers map[reflect.Type][]any
|
|
}
|
|
|
|
// NewHub создаёт подписчик на указанную шину (пока не запущен).
|
|
func NewHub(bus *Bus) *Hub {
|
|
return &Hub{
|
|
bus: bus,
|
|
done: make(chan struct{}),
|
|
handlers: make(map[reflect.Type][]any),
|
|
}
|
|
}
|
|
|
|
// On регистрирует обработчик для типа события T. Безопасно вызывать до Start
|
|
// и из других горутин.
|
|
func On[T Event](h *Hub, fn Handler[T]) {
|
|
h.mu.Lock()
|
|
defer h.mu.Unlock()
|
|
var zero T
|
|
typ := reflect.TypeOf(zero)
|
|
h.handlers[typ] = append(h.handlers[typ], fn)
|
|
}
|
|
|
|
// Start подписывается на шину и запускает горутину чтения событий.
|
|
func (h *Hub) Start() {
|
|
if h.ch != nil {
|
|
return
|
|
}
|
|
h.ch, h.unsub = h.bus.Subscribe()
|
|
go h.run()
|
|
}
|
|
|
|
// Close останавливает горутину и отписывается от шины. Идемпотентен.
|
|
func (h *Hub) Close() {
|
|
h.once.Do(func() {
|
|
close(h.done)
|
|
if h.unsub != nil {
|
|
h.unsub()
|
|
}
|
|
})
|
|
}
|
|
|
|
// run — цикл чтения событий и диспетчеризации.
|
|
func (h *Hub) run() {
|
|
defer h.Close()
|
|
for {
|
|
select {
|
|
case e, ok := <-h.ch:
|
|
if !ok {
|
|
return
|
|
}
|
|
h.dispatch(e)
|
|
case <-h.done:
|
|
return
|
|
}
|
|
}
|
|
}
|
|
|
|
// dispatch вызывает все обработчики, зарегистрированные для типа события e.
|
|
func (h *Hub) dispatch(e Event) {
|
|
typ := reflect.TypeOf(e)
|
|
|
|
h.mu.Lock()
|
|
fns := append([]any(nil), h.handlers[typ]...)
|
|
h.mu.Unlock()
|
|
|
|
for _, fn := range fns {
|
|
reflect.ValueOf(fn).Call([]reflect.Value{reflect.ValueOf(e)})
|
|
}
|
|
} |