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 }