kra-oa/app/system/internal/worker/task_scheduler_test.go

51 lines
1.3 KiB
Go

package worker
import "testing"
func TestCloseSubscribersClosesEveryAlertStream(t *testing.T) {
scheduler := &TaskScheduler{subscribers: map[uint]map[chan []byte]struct{}{}}
first := scheduler.Subscribe(1)
second := scheduler.Subscribe(1)
third := scheduler.Subscribe(2)
scheduler.closeSubscribers()
for _, ch := range []chan []byte{first, second, third} {
if _, open := <-ch; open {
t.Fatal("subscriber channel remains open after scheduler shutdown")
}
}
if len(scheduler.subscribers) != 0 {
t.Fatalf("subscriber registry was not cleared: %#v", scheduler.subscribers)
}
// Handlers defer Unsubscribe. It must remain safe after shutdown has already
// closed and removed every channel.
scheduler.Unsubscribe(1, first)
}
func TestSubscribeCapsConnectionsPerUser(t *testing.T) {
scheduler := &TaskScheduler{subscribers: map[uint]map[chan []byte]struct{}{}}
channels := make([]chan []byte, 0, 11)
for i := 0; i < 11; i++ {
channels = append(channels, scheduler.Subscribe(7))
}
if got := len(scheduler.subscribers[7]); got != 10 {
t.Fatalf("subscriber count = %d, want 10", got)
}
closed := 0
for _, ch := range channels {
select {
case _, open := <-ch:
if !open {
closed++
}
default:
}
}
if closed != 1 {
t.Fatalf("evicted subscriber count = %d, want 1", closed)
}
scheduler.closeSubscribers()
}