51 lines
1.3 KiB
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()
|
|
}
|