diff --git a/pkg/rtc/dynacast/dynacastmanager_test.go b/pkg/rtc/dynacast/dynacastmanager_test.go index 6d48e67ea..43d939546 100644 --- a/pkg/rtc/dynacast/dynacastmanager_test.go +++ b/pkg/rtc/dynacast/dynacastmanager_test.go @@ -498,6 +498,53 @@ func TestCodecRegression(t *testing.T) { }) } +func TestResendCommittedQuality(t *testing.T) { + var lock sync.Mutex + var notified []string + dm := NewDynacastManagerVideo(DynacastManagerVideoParams{ + // keeps the downgrade below pending for the whole test + DynacastPauseDelay: time.Minute, + Listener: &testDynacastManagerListener{ + onSubscribedMaxQualityChange: func(subscribedQualities []*livekit.SubscribedCodec) { + lock.Lock() + notified = append(notified, subscribedCodecsAsString(subscribedQualities)) + lock.Unlock() + }, + }, + }) + defer dm.Close() + notifications := func() []string { + lock.Lock() + defer lock.Unlock() + return slices.Clone(notified) + } + + // nothing committed yet, nothing to resend + dm.ResendCommittedQuality() + time.Sleep(100 * time.Millisecond) + require.Empty(t, notifications()) + + dm.NotifySubscriberMaxQuality("s1", mime.MimeTypeVP8, livekit.VideoQuality_LOW) + require.Eventually(t, func() bool { return len(notifications()) == 1 }, 5*time.Second, 10*time.Millisecond) + + // the downgrade waits for the debounce, LOW is still the committed quality + dm.NotifySubscriberMaxQuality("s1", mime.MimeTypeVP8, livekit.VideoQuality_OFF) + time.Sleep(100 * time.Millisecond) + require.Len(t, notifications(), 1) + + // the resend notifies LOW again, without committing the pending OFF + dm.ResendCommittedQuality() + require.Eventually(t, func() bool { return len(notifications()) == 2 }, 5*time.Second, 10*time.Millisecond) + require.Equal(t, notifications()[0], notifications()[1]) + + // no layer paused, nothing to resend + dm.NotifySubscriberMaxQuality("s1", mime.MimeTypeVP8, livekit.VideoQuality_HIGH) + require.Eventually(t, func() bool { return len(notifications()) == 3 }, 5*time.Second, 10*time.Millisecond) + dm.ResendCommittedQuality() + time.Sleep(100 * time.Millisecond) + require.Len(t, notifications(), 3) +} + func subscribedCodecsAsString(c1 []*livekit.SubscribedCodec) string { slices.SortFunc(c1, func(a, b *livekit.SubscribedCodec) int { return strings.Compare(a.Codec, b.Codec) diff --git a/pkg/rtc/dynacast/dynacastmanagervideo.go b/pkg/rtc/dynacast/dynacastmanagervideo.go index 8178b4627..86c04e818 100644 --- a/pkg/rtc/dynacast/dynacastmanagervideo.go +++ b/pkg/rtc/dynacast/dynacastmanagervideo.go @@ -104,6 +104,21 @@ func (d *dynacastManagerVideo) ForceQuality(quality livekit.VideoQuality) { d.enqueueSubscribedQualityChange() } +// ResendCommittedQuality notifies the listener again of the last committed subscribed qualities, +// for a publisher that may have lost them. Nothing new is committed: a debounced downgrade stays pending. +// Only needed if some layer is paused, i.e. some committed quality is not HIGH. +func (d *dynacastManagerVideo) ResendCommittedQuality() { + d.lock.Lock() + defer d.lock.Unlock() + + for _, quality := range d.committedMaxSubscribedQuality { + if quality != livekit.VideoQuality_HIGH { + d.enqueueSubscribedQualityChange() + return + } + } +} + func (d *dynacastManagerVideo) NotifySubscriberMaxQuality( subscriberID livekit.ParticipantID, mime mime.MimeType, diff --git a/pkg/rtc/dynacast/interfaces.go b/pkg/rtc/dynacast/interfaces.go index c544ff54a..a64cfc962 100644 --- a/pkg/rtc/dynacast/interfaces.go +++ b/pkg/rtc/dynacast/interfaces.go @@ -53,6 +53,7 @@ type DynacastManager interface { Restart() Close() ForceUpdate() + ResendCommittedQuality() ForceQuality(quality livekit.VideoQuality) ForceEnable(enabled bool) @@ -88,6 +89,7 @@ func (d *dynacastManagerNull) HandleCodecRegression(fromMime, toMime mime.MimeTy func (d *dynacastManagerNull) Restart() {} func (d *dynacastManagerNull) Close() {} func (d *dynacastManagerNull) ForceUpdate() {} +func (d *dynacastManagerNull) ResendCommittedQuality() {} func (d *dynacastManagerNull) ForceQuality(quality livekit.VideoQuality) {} func (d *dynacastManagerNull) ForceEnable(enabled bool) {} func (d *dynacastManagerNull) NotifySubscriberMaxQuality( diff --git a/pkg/rtc/mediatrack.go b/pkg/rtc/mediatrack.go index be1436fce..9793e2103 100644 --- a/pkg/rtc/mediatrack.go +++ b/pkg/rtc/mediatrack.go @@ -644,6 +644,13 @@ func (t *MediaTrack) SetMuted(muted bool) { t.MediaTrackReceiver.SetMuted(muted) } +// ResendSubscribedQuality sends the publisher the subscribed qualities of the track again. +func (t *MediaTrack) ResendSubscribedQuality() { + if t.dynacastManager != nil { + t.dynacastManager.ResendCommittedQuality() + } +} + // OnTrackSubscribed is called when the track is subscribed by a non-hidden subscriber // this allows the publisher to know when they should start sending data func (t *MediaTrack) OnTrackSubscribed() { diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index af3f8e2d3..4491c6047 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -74,6 +74,9 @@ const ( sdBatchSize = 30 rttUpdateInterval = 5 * time.Second + publisherAnswerDynacastResendDelay = 2 * time.Second + publisherAnswerDynacastResendMaxDelay = 8 * time.Second + disconnectCleanupDuration = 5 * time.Second migrationWaitDuration = 3 * time.Second migrationWaitContinuousMsgDuration = 2 * time.Second @@ -277,6 +280,9 @@ type ParticipantImpl struct { // timer that's set when disconnect is detected on primary PC disconnectTimer *time.Timer migrationTimer *time.Timer + // re-send of the subscribed qualities pending after publisher answers + dynacastResendTimer *time.Timer + dynacastResendDeadline time.Time migratedInAt atomic.Pointer[time.Time] @@ -1268,7 +1274,38 @@ func (p *ParticipantImpl) onPublisherAnswer(answer webrtc.SessionDescription, an "midToTrackID", midToTrackID, ) - return p.sendSdpAnswer(answer, answerId, midToTrackID) + if err := p.sendSdpAnswer(answer, answerId, midToTrackID); err != nil { + return err + } + + // The answer does not carry the pause state of the simulcast layers (SetIgnoreRidPauseForRecv), + // so applying it re-enables in the publisher the layers paused by dynacast. Send the subscribed + // qualities again with a delay to ensure the answer has been applied, once for a burst of answers: + // the delay restarts with each answer, up to a maximum from the first one + p.lock.Lock() + if p.dynacastResendTimer == nil { + p.dynacastResendDeadline = time.Now().Add(publisherAnswerDynacastResendMaxDelay) + p.dynacastResendTimer = time.AfterFunc(publisherAnswerDynacastResendDelay, p.resendSubscribedQualities) + } else { + p.dynacastResendTimer.Reset(min(publisherAnswerDynacastResendDelay, time.Until(p.dynacastResendDeadline))) + } + p.lock.Unlock() + return nil +} + +func (p *ParticipantImpl) resendSubscribedQualities() { + p.lock.Lock() + p.dynacastResendTimer = nil + p.lock.Unlock() + + if p.IsClosed() || p.IsDisconnected() { + return + } + for _, track := range p.GetPublishedTracks() { + if mt, ok := track.(*MediaTrack); ok { + mt.ResendSubscribedQuality() + } + } } func (p *ParticipantImpl) GetAnswer() (webrtc.SessionDescription, uint32, error) {