Files
ratatoskr-go/internal/events/bus_test.go
ki.sagidullin 1f7ab9a67e feat(events): EventBus + domain-модель статусов
- internal/model: единый источник статусов задачи (Status, TraceStatus,
  машина переходов), без зависимости от storage.
- internal/storage: совместимый мост (type Status = model.Status,
  re-export констант) — внешний код не меняется.
- internal/events: шина событий (fan-out, блокирующий Publish с гарантией
  порядка), события задач/трейсов, отдельная логовая шина + LogWriter,
  Publisher/NilPublisher для внедрения в Core/Worker.
- docs/ui-spec.md: спецификация десктопного UI (Fyne).
2026-08-20 01:02:31 +05:00

119 lines
2.7 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package events
import (
"testing"
"time"
"github.com/kamelion/ratatoskr-go/internal/model"
)
// receiveOne помогает получить одно событие с таймаутом.
func receiveOne(t *testing.T, ch <-chan Event) Event {
t.Helper()
select {
case e := <-ch:
return e
case <-time.After(2 * time.Second):
t.Fatal("timeout waiting for event")
return nil
}
}
func TestBusPublishToSubscriber(t *testing.T) {
bus := New(10)
ch, unsub := bus.Subscribe()
defer unsub()
want := TaskStatusChanged{ID: 42, From: model.StatusReady, To: model.StatusRunning}
bus.Publish(want)
got := receiveOne(t, ch)
ev, ok := got.(TaskStatusChanged)
if !ok {
t.Fatalf("got %T, want TaskStatusChanged", got)
}
if ev.ID != 42 || ev.From != model.StatusReady || ev.To != model.StatusRunning {
t.Fatalf("unexpected event: %+v", ev)
}
}
func TestBusFanout(t *testing.T) {
bus := New(10)
ch1, unsub1 := bus.Subscribe()
defer unsub1()
ch2, unsub2 := bus.Subscribe()
defer unsub2()
e := LogLine{Level: "log", Text: "hello"}
bus.Publish(e)
if got := receiveOne(t, ch1); got != e {
t.Fatalf("subscriber 1 got %#v, want %#v", got, e)
}
if got := receiveOne(t, ch2); got != e {
t.Fatalf("subscriber 2 got %#v, want %#v", got, e)
}
}
func TestBusPreservesOrder(t *testing.T) {
bus := New(64)
ch, unsub := bus.Subscribe()
defer unsub()
const n = 25
for i := 0; i < n; i++ {
bus.Publish(TraceAppended{TaskID: int64(i)})
}
for i := 0; i < n; i++ {
ev := receiveOne(t, ch)
ta, ok := ev.(TraceAppended)
if !ok {
t.Fatalf("got %T, want TraceAppended", ev)
}
if ta.TaskID != int64(i) {
t.Fatalf("out of order: got %d, want %d", ta.TaskID, i)
}
}
}
func TestUnsubscribeStopsDelivery(t *testing.T) {
bus := New(10)
ch, unsub := bus.Subscribe()
bus.Publish(TaskCreated{ID: 1})
receiveOne(t, ch)
unsub()
if got := bus.SubscribersCount(); got != 0 {
t.Fatalf("SubscribersCount = %d, want 0", got)
}
// Убеждаемся, что Publish не блокируется навечно отписанным подписчиком.
bus.Publish(TaskCreated{ID: 2})
select {
case got := <-ch:
t.Fatalf("received %#v after unsubscribe", got)
case <-time.After(200 * time.Millisecond):
}
}
func TestBusSubscribersCount(t *testing.T) {
bus := New(10)
if got := bus.SubscribersCount(); got != 0 {
t.Fatalf("initial count = %d, want 0", got)
}
_, unsub1 := bus.Subscribe()
_, unsub2 := bus.Subscribe()
if got := bus.SubscribersCount(); got != 2 {
t.Fatalf("count = %d, want 2", got)
}
unsub1()
unsub2()
if got := bus.SubscribersCount(); got != 0 {
t.Fatalf("after unsub count = %d, want 0", got)
}
}
func TestNilPublisherNoOp(t *testing.T) {
NilPublisher{}.Publish(TaskCreated{ID: 1}) // must not panic
}