diff --git a/pkg/rtc/mediatrack_test.go b/pkg/rtc/mediatrack_test.go index eda544909..fe5d1a961 100644 --- a/pkg/rtc/mediatrack_test.go +++ b/pkg/rtc/mediatrack_test.go @@ -142,6 +142,7 @@ func TestSubscribedMaxQuality(t *testing.T) { }, }, }) + mt.Start() var lock sync.Mutex actualTrackID := livekit.TrackID("") actualSubscribedQualities := make([]*livekit.SubscribedCodec, 0) @@ -216,6 +217,7 @@ func TestSubscribedMaxQuality(t *testing.T) { DynacastPauseDelay: 100 * time.Millisecond, }, }) + mt.Start() mt.AddCodec(webrtc.MimeTypeVP8) mt.AddCodec(webrtc.MimeTypeAV1) @@ -259,6 +261,7 @@ func TestSubscribedMaxQuality(t *testing.T) { }, }, } + time.Sleep(10 * time.Millisecond) lock.RLock() require.Equal(t, livekit.TrackID("v1"), actualTrackID) require.EqualValues(t, subscribedCodecsAsString(expectedSubscribedQualities), subscribedCodecsAsString(actualSubscribedQualities)) @@ -327,6 +330,7 @@ func TestSubscribedMaxQuality(t *testing.T) { // muting "s2" only should not disable all qualities of vp8, no change of expected qualities mt.notifySubscriberMaxQuality("s2", webrtc.RTPCodecCapability{MimeType: webrtc.MimeTypeVP8}, livekit.VideoQuality_OFF) + time.Sleep(10 * time.Millisecond) lock.RLock() require.Equal(t, livekit.TrackID("v1"), actualTrackID) require.EqualValues(t, subscribedCodecsAsString(expectedSubscribedQualities), subscribedCodecsAsString(actualSubscribedQualities)) @@ -411,6 +415,7 @@ func TestSubscribedMaxQuality(t *testing.T) { }, }, } + time.Sleep(10 * time.Millisecond) lock.RLock() require.Equal(t, livekit.TrackID("v1"), actualTrackID) require.EqualValues(t, subscribedCodecsAsString(expectedSubscribedQualities), subscribedCodecsAsString(actualSubscribedQualities)) diff --git a/pkg/rtc/mediatrackreceiver.go b/pkg/rtc/mediatrackreceiver.go index f6ebb2f51..59736995f 100644 --- a/pkg/rtc/mediatrackreceiver.go +++ b/pkg/rtc/mediatrackreceiver.go @@ -228,11 +228,11 @@ func (t *MediaTrackReceiver) ClearReceiver(mime string) { t.receiversShadow = make([]*simulcastReceiver, len(t.receivers)) copy(t.receiversShadow, t.receivers) - closeSubscription := len(t.receiversShadow) == 0 + stopSubscription := len(t.receiversShadow) == 0 t.lock.Unlock() - if closeSubscription { - t.MediaTrackSubscriptions.Close() + if stopSubscription { + t.MediaTrackSubscriptions.Stop() } } @@ -241,7 +241,7 @@ func (t *MediaTrackReceiver) ClearAllReceivers() { t.receivers = t.receivers[:0] t.receiversShadow = nil t.lock.Unlock() - t.MediaTrackSubscriptions.Close() + t.MediaTrackSubscriptions.Stop() } func (t *MediaTrackReceiver) OnMediaLossUpdate(f func(fractionalLoss uint8)) { @@ -253,20 +253,26 @@ func (t *MediaTrackReceiver) OnVideoLayerUpdate(f func(layers []*livekit.VideoLa } func (t *MediaTrackReceiver) TryClose() bool { - t.lock.Lock() + t.lock.RLock() if len(t.receiversShadow) > 0 { t.lock.Unlock() return false } + t.lock.RUnlock() + t.Close() + + return true +} + +func (t *MediaTrackReceiver) Close() { + t.lock.RLock() onclose := t.onClose - t.lock.Unlock() + t.lock.RUnlock() t.MediaTrackSubscriptions.Close() - for _, f := range onclose { f() } - return true } func (t *MediaTrackReceiver) ID() livekit.TrackID { diff --git a/pkg/rtc/mediatracksubscriptions.go b/pkg/rtc/mediatracksubscriptions.go index 4bb7c9fb0..4b9bff75a 100644 --- a/pkg/rtc/mediatracksubscriptions.go +++ b/pkg/rtc/mediatracksubscriptions.go @@ -45,6 +45,8 @@ type MediaTrackSubscriptions struct { maxSubscribedQualityDebounce func(func()) onSubscribedMaxQualityChange func(subscribedQualities []*livekit.SubscribedCodec, maxSubscribedQualities []types.SubscribedCodecQuality) maxQualityTimer *time.Timer + + qualityNotifyOpQueue *utils.OpsQueue } type MediaTrackSubscriptionsParams struct { @@ -69,12 +71,14 @@ func NewMediaTrackSubscriptions(params MediaTrackSubscriptionsParams) *MediaTrac maxSubscriberNodeQuality: make(map[livekit.NodeID][]types.SubscribedCodecQuality), maxSubscribedQuality: make(map[string]livekit.VideoQuality), maxSubscribedQualityDebounce: debounce.New(params.VideoConfig.DynacastPauseDelay), + qualityNotifyOpQueue: utils.NewOpsQueue(params.Logger, "quality-notify", 100), } return t } func (t *MediaTrackSubscriptions) Start() { + t.qualityNotifyOpQueue.Start() t.startMaxQualityTimer(false) } @@ -82,10 +86,14 @@ func (t *MediaTrackSubscriptions) Restart() { t.startMaxQualityTimer(true) } -func (t *MediaTrackSubscriptions) Close() { +func (t *MediaTrackSubscriptions) Stop() { t.stopMaxQualityTimer() } +func (t *MediaTrackSubscriptions) Close() { + t.qualityNotifyOpQueue.Stop() +} + func (t *MediaTrackSubscriptions) OnNoSubscribers(f func()) { t.onNoSubscribers = f } @@ -618,14 +626,15 @@ func (t *MediaTrackSubscriptions) UpdateQualityChange(force bool) { }) } } - t.maxQualityLock.Unlock() - if t.onSubscribedMaxQualityChange != nil { t.params.Logger.Debugw("subscribedMaxQualityChange", "subscribedCodec", subscribedCodec, "maxSubscribedQualities", maxSubscribedQualities) - t.onSubscribedMaxQualityChange(subscribedCodec, maxSubscribedQualities) + t.qualityNotifyOpQueue.Enqueue(func() { + t.onSubscribedMaxQualityChange(subscribedCodec, maxSubscribedQualities) + }) } + t.maxQualityLock.Unlock() } func (t *MediaTrackSubscriptions) startMaxQualityTimer(force bool) {