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() }