diff --git a/pkg/rfc8888/interceptor.go b/pkg/rfc8888/interceptor.go index 9a47e09b..83149ac3 100644 --- a/pkg/rfc8888/interceptor.go +++ b/pkg/rfc8888/interceptor.go @@ -142,6 +142,8 @@ func (s *SenderInterceptor) BindRemoteStream( func (s *SenderInterceptor) Close() error { s.log.Trace("close") defer s.wg.Wait() + s.lock.Lock() + defer s.lock.Unlock() if !s.isClosed() { close(s.close) diff --git a/pkg/rfc8888/interceptor_test.go b/pkg/rfc8888/interceptor_test.go index f535af64..891061ba 100644 --- a/pkg/rfc8888/interceptor_test.go +++ b/pkg/rfc8888/interceptor_test.go @@ -4,6 +4,7 @@ package rfc8888 import ( + "sync" "testing" "time" @@ -323,3 +324,27 @@ func TestReadAfterClose(t *testing.T) { assert.Fail(t, "read after close blocked") } } + +func TestConcurrentClose(t *testing.T) { + f, err := NewSenderInterceptor() + assert.NoError(t, err) + + intcp, err := f.NewInterceptor("") + assert.NoError(t, err) + + intcp.BindRTCPWriter(interceptor.RTCPWriterFunc( + func(_ []rtcp.Packet, _ interceptor.Attributes) (int, error) { + return 0, nil + }, + )) + + var wg sync.WaitGroup + for range 10 { + wg.Add(1) + go func() { + defer wg.Done() + assert.NoError(t, intcp.Close()) + }() + } + wg.Wait() +}