diff --git a/internal/events/hub.go b/internal/events/hub.go new file mode 100644 index 0000000..51aeaaa --- /dev/null +++ b/internal/events/hub.go @@ -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)}) + } +} \ No newline at end of file diff --git a/internal/events/hub_test.go b/internal/events/hub_test.go new file mode 100644 index 0000000..04abb1c --- /dev/null +++ b/internal/events/hub_test.go @@ -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}) +} \ No newline at end of file