unistore: close event stream on context cancelation (#101293)
* add tests for broacaster * fix sql notifier not closing the stream * fix sql notifier not closing the stream * close sub * fix broadcaster test * fix broadcaster test * suggestion
This commit is contained in:
@@ -104,3 +104,41 @@ func TestCache(t *testing.T) {
|
||||
// slice should return all values
|
||||
require.Equal(t, []int{4, 5, 6, 7, 8, 9, 10, 11, 12, 13}, c.Slice())
|
||||
}
|
||||
|
||||
func TestBroadcaster(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
ch := make(chan int)
|
||||
input := []int{1, 2, 3}
|
||||
go func() {
|
||||
for _, v := range input {
|
||||
ch <- v
|
||||
}
|
||||
}()
|
||||
t.Cleanup(func() {
|
||||
close(ch)
|
||||
})
|
||||
|
||||
b, err := NewBroadcaster(ctx, func(out chan<- int) error {
|
||||
go func() {
|
||||
for v := range ch {
|
||||
out <- v
|
||||
}
|
||||
}()
|
||||
return nil
|
||||
})
|
||||
require.NoError(t, err)
|
||||
|
||||
sub, err := b.Subscribe(ctx)
|
||||
require.NoError(t, err)
|
||||
|
||||
for _, expected := range input {
|
||||
v, ok := <-sub
|
||||
require.True(t, ok)
|
||||
require.Equal(t, expected, v)
|
||||
}
|
||||
|
||||
// cancel the context should close the stream
|
||||
cancel()
|
||||
_, ok := <-sub
|
||||
require.False(t, ok)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user