feat/new_ui #8
96
internal/events/hub.go
Normal file
96
internal/events/hub.go
Normal file
@@ -0,0 +1,96 @@
|
||||
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)})
|
||||
}
|
||||
}
|
||||
119
internal/events/hub_test.go
Normal file
119
internal/events/hub_test.go
Normal file
@@ -0,0 +1,119 @@
|
||||
package events
|
||||
|
||||
import (
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/kamelion/ratatoskr-go/internal/model"
|
||||
)
|
||||
|
||||
func TestHubDeliversTypedEvent(t *testing.T) {
|
||||
bus := New(16)
|
||||
h := NewHub(bus)
|
||||
|
||||
var got atomic.Value
|
||||
On[TaskStatusChanged](h, func(e TaskStatusChanged) {
|
||||
got.CompareAndSwap(nil, e)
|
||||
})
|
||||
|
||||
h.Start()
|
||||
defer h.Close()
|
||||
|
||||
want := TaskStatusChanged{ID: 7, From: model.StatusReady, To: model.StatusRunning}
|
||||
bus.Publish(want)
|
||||
|
||||
deadline := time.Now().Add(2 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
if raw := got.Load(); raw != nil {
|
||||
ev := raw.(TaskStatusChanged)
|
||||
if ev.ID != 7 || ev.From != model.StatusReady || ev.To != model.StatusRunning {
|
||||
t.Fatalf("unexpected event: %+v", ev)
|
||||
}
|
||||
return
|
||||
}
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
}
|
||||
t.Fatal("handler was not called")
|
||||
}
|
||||
|
||||
func TestHubIgnoresUnrelatedTypes(t *testing.T) {
|
||||
bus := New(16)
|
||||
h := NewHub(bus)
|
||||
|
||||
var calls atomic.Int32
|
||||
On[TaskStatusChanged](h, func(TaskStatusChanged) { calls.Add(1) })
|
||||
|
||||
h.Start()
|
||||
defer h.Close()
|
||||
|
||||
// другие типы событий не должны дойти до этого обработчика
|
||||
bus.Publish(AgentActivity{TaskID: 1, Agent: "dev", Stage: "run"})
|
||||
bus.Publish(LogLine{Level: "log", Text: "x"})
|
||||
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
if n := calls.Load(); n != 0 {
|
||||
t.Fatalf("handler called %d times for unrelated events", n)
|
||||
}
|
||||
}
|
||||
|
||||
func TestHubMultipleHandlersSameType(t *testing.T) {
|
||||
bus := New(16)
|
||||
h := NewHub(bus)
|
||||
|
||||
var a, b atomic.Int32
|
||||
On[HistoryAppended](h, func(HistoryAppended) { a.Add(1) })
|
||||
On[HistoryAppended](h, func(HistoryAppended) { b.Add(1) })
|
||||
|
||||
h.Start()
|
||||
defer h.Close()
|
||||
|
||||
bus.Publish(HistoryAppended{TaskID: 1, Role: "user", Content: "hi"})
|
||||
|
||||
deadline := time.Now().Add(2 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
if a.Load() == 1 && b.Load() == 1 {
|
||||
return
|
||||
}
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
}
|
||||
t.Fatal("not all handlers called")
|
||||
}
|
||||
|
||||
func TestHubCloseStopsDelivery(t *testing.T) {
|
||||
bus := New(16)
|
||||
h := NewHub(bus)
|
||||
|
||||
var calls atomic.Int32
|
||||
On[TaskCreated](h, func(TaskCreated) { calls.Add(1) })
|
||||
|
||||
h.Start()
|
||||
bus.Publish(TaskCreated{ID: 1})
|
||||
|
||||
deadline := time.Now().Add(2 * time.Second)
|
||||
for time.Now().Before(deadline) && calls.Load() == 0 {
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
}
|
||||
if calls.Load() == 0 {
|
||||
t.Fatal("initial delivery failed")
|
||||
}
|
||||
|
||||
h.Close()
|
||||
// после Close Hub отписан — события не доходят
|
||||
bus.Publish(TaskCreated{ID: 2})
|
||||
time.Sleep(150 * time.Millisecond)
|
||||
if n := calls.Load(); n > 1 {
|
||||
t.Fatalf("handler called %d times after Close", n)
|
||||
}
|
||||
}
|
||||
|
||||
func TestHubStartIdempotent(t *testing.T) {
|
||||
bus := New(16)
|
||||
h := NewHub(bus)
|
||||
h.Start()
|
||||
h.Start() // повторный Start не должен создавать вторую горутину
|
||||
h.Close()
|
||||
// если бы было две горутины — Publish блокировался бы на буфере, но закрытие сняло бы блок;
|
||||
// просто проверяем, что Close и повторный Start не падают
|
||||
bus.Publish(TaskCreated{ID: 1})
|
||||
}
|
||||
Reference in New Issue
Block a user