From 591888f712acfb91ca8ecfce97431a4c56cec0ec Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Sun, 2 Mar 2025 11:52:17 +0530 Subject: [PATCH] Fix missing RTCP sender report when forwarding RED as Opus. (#3480) With publish RED and subscribe Opus, the RTCP sender reports were not sent to down track as publisher sender reports were not forwarded to the down track. --- pkg/sfu/receiver.go | 19 ++++++++++++++++--- pkg/sfu/redprimaryreceiver.go | 12 ++++++++++++ pkg/sfu/redreceiver.go | 12 ++++++++++++ 3 files changed, 40 insertions(+), 3 deletions(-) diff --git a/pkg/sfu/receiver.go b/pkg/sfu/receiver.go index ee307a863..4793820c8 100644 --- a/pkg/sfu/receiver.go +++ b/pkg/sfu/receiver.go @@ -143,6 +143,12 @@ type TrackReceiver interface { } type redPktWriteFunc func(pkt *buffer.ExtPacket, spatialLayer int32) int +type redSenderReportWriteFunc func( + payloadType webrtc.PayloadType, + isSVC bool, + layer int32, + publisherSRData *livekit.RTCPSenderReportState, +) // WebRTCReceiver receives a media track type WebRTCReceiver struct { @@ -185,9 +191,10 @@ type WebRTCReceiver struct { onStatsUpdate func(w *WebRTCReceiver, stat *livekit.AnalyticsStat) onMaxLayerChange func(maxLayer int32) - primaryReceiver atomic.Pointer[RedPrimaryReceiver] - redReceiver atomic.Pointer[RedReceiver] - redPktWriter atomic.Value // redPktWriteFunc + primaryReceiver atomic.Pointer[RedPrimaryReceiver] + redReceiver atomic.Pointer[RedReceiver] + redPktWriter atomic.Value // redPktWriteFunc + redSenderReportWriter atomic.Value // redPktWriteFunc forwardStats *ForwardStats } @@ -405,6 +412,10 @@ func (w *WebRTCReceiver) AddUpTrack(track TrackRemote, buff *buffer.Buffer) erro w.downTrackSpreader.Broadcast(func(dt TrackSender) { _ = dt.HandleRTCPSenderReportData(w.codec.PayloadType, w.isSVC, layer, srData) }) + + if f := w.redSenderReportWriter.Load(); f != nil { + f.(redSenderReportWriteFunc)(w.codec.PayloadType, w.isSVC, layer, srData) + } }) if w.Kind() == webrtc.RTPCodecTypeVideo && layer == 0 { @@ -866,6 +877,7 @@ func (w *WebRTCReceiver) GetPrimaryReceiverForRed() TrackReceiver { }) if w.primaryReceiver.CompareAndSwap(nil, pr) { w.redPktWriter.Store(redPktWriteFunc(pr.ForwardRTP)) + w.redSenderReportWriter.Store(redSenderReportWriteFunc(pr.ForwardRTCPSenderReport)) } } return w.primaryReceiver.Load() @@ -883,6 +895,7 @@ func (w *WebRTCReceiver) GetRedReceiver() TrackReceiver { }) if w.redReceiver.CompareAndSwap(nil, pr) { w.redPktWriter.Store(redPktWriteFunc(pr.ForwardRTP)) + w.redSenderReportWriter.Store(redSenderReportWriteFunc(pr.ForwardRTCPSenderReport)) } } return w.redReceiver.Load() diff --git a/pkg/sfu/redprimaryreceiver.go b/pkg/sfu/redprimaryreceiver.go index f066205b3..ffa9028cc 100644 --- a/pkg/sfu/redprimaryreceiver.go +++ b/pkg/sfu/redprimaryreceiver.go @@ -21,6 +21,7 @@ import ( "go.uber.org/atomic" "github.com/pion/rtp" + "github.com/pion/webrtc/v4" "github.com/livekit/livekit-server/pkg/sfu/buffer" "github.com/livekit/protocol/livekit" @@ -94,6 +95,17 @@ func (r *RedPrimaryReceiver) ForwardRTP(pkt *buffer.ExtPacket, spatialLayer int3 return count } +func (r *RedPrimaryReceiver) ForwardRTCPSenderReport( + payloadType webrtc.PayloadType, + isSVC bool, + layer int32, + publisherSRData *livekit.RTCPSenderReportState, +) { + r.downTrackSpreader.Broadcast(func(dt TrackSender) { + _ = dt.HandleRTCPSenderReportData(payloadType, isSVC, layer, publisherSRData) + }) +} + func (r *RedPrimaryReceiver) AddDownTrack(track TrackSender) error { if r.closed.Load() { return ErrReceiverClosed diff --git a/pkg/sfu/redreceiver.go b/pkg/sfu/redreceiver.go index 7fbd9a6eb..d095f85dc 100644 --- a/pkg/sfu/redreceiver.go +++ b/pkg/sfu/redreceiver.go @@ -21,6 +21,7 @@ import ( "go.uber.org/atomic" "github.com/pion/rtp" + "github.com/pion/webrtc/v4" "github.com/livekit/livekit-server/pkg/sfu/buffer" "github.com/livekit/mediatransportutil/pkg/bucket" @@ -89,6 +90,17 @@ func (r *RedReceiver) ForwardRTP(pkt *buffer.ExtPacket, spatialLayer int32) int }) } +func (r *RedReceiver) ForwardRTCPSenderReport( + payloadType webrtc.PayloadType, + isSVC bool, + layer int32, + publisherSRData *livekit.RTCPSenderReportState, +) { + r.downTrackSpreader.Broadcast(func(dt TrackSender) { + _ = dt.HandleRTCPSenderReportData(payloadType, isSVC, layer, publisherSRData) + }) +} + func (r *RedReceiver) AddDownTrack(track TrackSender) error { if r.closed.Load() { return ErrReceiverClosed