From bd616d607450a185657e4b358924c5470a4c5276 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Sat, 20 Jul 2024 10:25:06 +0530 Subject: [PATCH] Split ICE candidate queue. (#2885) Shared ICE candidate queue meant only one of PUBLISHER/SUBSCRIBER pc got final candidate notification. Split the queue. --- pkg/rtc/participant.go | 2 +- pkg/rtc/participant_signal.go | 9 ++++++++- pkg/sfu/utils/wraparound.go | 8 ++++---- test/client/client.go | 10 ++++++++-- 4 files changed, 21 insertions(+), 8 deletions(-) diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index bb2fc66bd..dc0828c66 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -191,7 +191,7 @@ type ParticipantImpl struct { *UpTrackManager *SubscriptionManager - icQueue atomic.Pointer[webrtc.ICECandidate] + icQueue [2]atomic.Pointer[webrtc.ICECandidate] // keeps track of unpublished tracks in order to reuse trackID unpublishedTracks []*livekit.TrackInfo diff --git a/pkg/rtc/participant_signal.go b/pkg/rtc/participant_signal.go index 64f5ab393..8762c0270 100644 --- a/pkg/rtc/participant_signal.go +++ b/pkg/rtc/participant_signal.go @@ -20,6 +20,7 @@ import ( "time" "github.com/pion/webrtc/v3" + "go.uber.org/atomic" "github.com/livekit/protocol/livekit" "github.com/livekit/protocol/logger" @@ -265,7 +266,13 @@ func (p *ParticipantImpl) sendDisconnectUpdatesForReconnect() error { } func (p *ParticipantImpl) sendICECandidate(ic *webrtc.ICECandidate, target livekit.SignalTarget) error { - prevIC := p.icQueue.Swap(ic) + var icQueue *atomic.Pointer[webrtc.ICECandidate] + if target == livekit.SignalTarget_PUBLISHER { + icQueue = &p.icQueue[0] + } else { + icQueue = &p.icQueue[1] + } + prevIC := icQueue.Swap(ic) if prevIC == nil { return nil } diff --git a/pkg/sfu/utils/wraparound.go b/pkg/sfu/utils/wraparound.go index a1a0441ca..d18d5a4ab 100644 --- a/pkg/sfu/utils/wraparound.go +++ b/pkg/sfu/utils/wraparound.go @@ -157,7 +157,7 @@ func (w *WrapAround[T, ET]) GetExtendedHighest() ET { } func (w *WrapAround[T, ET]) updateExtendedHighest() { - w.extendedHighest = getExtendedHighest(w.cycles, w.highest) + w.extendedHighest = getExtended(w.cycles, w.highest) } func (w *WrapAround[T, ET]) maybeAdjustStart(val T) (result WrapAroundUpdateResult[ET]) { @@ -172,7 +172,7 @@ func (w *WrapAround[T, ET]) maybeAdjustStart(val T) (result WrapAroundUpdateResu cycles -= w.fullRange } result.PreExtendedHighest = w.extendedHighest - result.ExtendedVal = getExtendedHighest(cycles, val) + result.ExtendedVal = getExtended(cycles, val) return } @@ -201,7 +201,7 @@ func (w *WrapAround[T, ET]) maybeAdjustStart(val T) (result WrapAroundUpdateResu } } result.PreExtendedHighest = w.extendedHighest - result.ExtendedVal = getExtendedHighest(cycles, val) + result.ExtendedVal = getExtended(cycles, val) return } @@ -211,6 +211,6 @@ func (w *WrapAround[T, ET]) isWrapBack(earlier T, later T) bool { // ------------------------------------ -func getExtendedHighest[T number, ET extendedNumber](cycles ET, val T) ET { +func getExtended[T number, ET extendedNumber](cycles ET, val T) ET { return cycles + ET(val) } diff --git a/test/client/client.go b/test/client/client.go index 49dbcad2d..117932737 100644 --- a/test/client/client.go +++ b/test/client/client.go @@ -69,7 +69,7 @@ type RTCClient struct { signalRequestInterceptor SignalRequestInterceptor signalResponseInterceptor SignalResponseInterceptor - icQueue atomic.Pointer[webrtc.ICECandidate] + icQueue [2]atomic.Pointer[webrtc.ICECandidate] subscriberAsPrimary atomic.Bool publisherFullyEstablished atomic.Bool @@ -555,7 +555,13 @@ func (c *RTCClient) sendRequest(msg *livekit.SignalRequest) error { } func (c *RTCClient) SendIceCandidate(ic *webrtc.ICECandidate, target livekit.SignalTarget) error { - prevIC := c.icQueue.Swap(ic) + var icQueue *atomic.Pointer[webrtc.ICECandidate] + if target == livekit.SignalTarget_PUBLISHER { + icQueue = &c.icQueue[0] + } else { + icQueue = &c.icQueue[1] + } + prevIC := icQueue.Swap(ic) if prevIC == nil { return nil }