From 1937631cccac769afb25eef7613b6f8ca441701c Mon Sep 17 00:00:00 2001 From: cnderrauber Date: Mon, 27 Feb 2023 22:15:11 +0800 Subject: [PATCH] worked for merge master --- pkg/rtc/mediatracksubscriptions.go | 30 ++++--- pkg/rtc/participant.go | 19 ++-- pkg/rtc/subscribedtrack.go | 10 ++- pkg/rtc/subscriptionmanager.go | 3 +- pkg/rtc/transport.go | 16 ++-- pkg/rtc/types/interfaces.go | 21 ++++- .../typesfakes/fake_local_participant.go | 82 +++++++++++++++++ .../types/typesfakes/fake_subscribed_track.go | 86 +++++++++++++++--- pkg/sfu/audioselection/forwarder.go | 6 +- pkg/sfu/audioselection/nullaudiodowntrack.go | 87 +++++++++++++++++++ 10 files changed, 309 insertions(+), 51 deletions(-) create mode 100644 pkg/sfu/audioselection/nullaudiodowntrack.go diff --git a/pkg/rtc/mediatracksubscriptions.go b/pkg/rtc/mediatracksubscriptions.go index 145d8cc59..a7d62e210 100644 --- a/pkg/rtc/mediatracksubscriptions.go +++ b/pkg/rtc/mediatracksubscriptions.go @@ -13,6 +13,7 @@ import ( "github.com/livekit/livekit-server/pkg/rtc/types" "github.com/livekit/livekit-server/pkg/sfu" + "github.com/livekit/livekit-server/pkg/sfu/audioselection" "github.com/livekit/livekit-server/pkg/telemetry" ) @@ -89,20 +90,25 @@ func (t *MediaTrackSubscriptions) AddSubscriber(sub types.LocalParticipant, wr * if t.params.MediaTrack.Kind() == livekit.TrackType_AUDIO /*&& audioselection.AudioCodecCanbeMux(*t.params.MediaTrack.ToProto(), wr.codecs) */ { wr.DetermineReceiver(opusCodecCapability) + ndt := audioselection.NewNullAudioDowntrack() subTrack := NewSubscribedTrack(SubscribedTrackParams{ PublisherID: t.params.MediaTrack.PublisherID(), PublisherIdentity: t.params.MediaTrack.PublisherIdentity(), PublisherVersion: t.params.MediaTrack.PublisherVersion(), Subscriber: sub, MediaTrack: t.params.MediaTrack, - DownTrack: nil, + DownTrack: ndt, AdaptiveStream: sub.GetAdaptiveStream(), + IsMuxedTrack: true, }) t.subscribedTracksMu.Lock() t.subscribedTracks[subscriberID] = subTrack t.subscribedTracksMu.Unlock() sub.VerifySubscribeParticipantInfo(subTrack.PublisherID(), subTrack.PublisherVersion()) sub.AddMuxAudioTrack(subTrack.PublisherID(), trackID, wr) + + go subTrack.Bound() + subTrack.SetPublisherMuted(t.params.MediaTrack.IsMuted()) return subTrack, nil } @@ -281,20 +287,22 @@ func (t *MediaTrackSubscriptions) RemoveSubscriber(subscriberID livekit.Particip } func (t *MediaTrackSubscriptions) closeSubscribedTrack(subTrack types.SubscribedTrack, willBeResumed bool) { - dt := subTrack.DownTrack() sub := subTrack.Subscriber() - if dt == nil { + if subTrack.IsMuxedTrack() { sub.RemoveMuxAudioTrack(t.params.MediaTrack.ID()) - return - } + go t.downTrackClosed(sub, willBeResumed) + } else { + if dt, ok := subTrack.DownTrack().(*sfu.DownTrack); ok { + dt.CloseWithFlush(!willBeResumed) - dt.CloseWithFlush(!willBeResumed) - - if willBeResumed { - tr := dt.GetTransceiver() - if tr != nil { - sub.CacheDownTrack(subTrack.ID(), tr, dt.GetState()) + if willBeResumed { + tr := dt.GetTransceiver() + if tr != nil { + sub.CacheDownTrack(subTrack.ID(), tr, dt.GetState()) + } + } } + } } diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index 51bca7c9a..0367a98f3 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -843,14 +843,13 @@ func (p *ParticipantImpl) GetConnectionQuality() *livekit.ConnectionQualityInfo } subscribedTracks := p.SubscriptionManager.GetSubscribedTracks() - subscriberScores := make(map[livekit.TrackID]float32, len(subscribedTracks)) - // TODO-mux: calculate score for muxed tracks + // subscriberScores := make(map[livekit.TrackID]float32, len(subscribedTracks)) for _, subTrack := range subscribedTracks { - if subTrack.IsMuted() || subTrack.MediaTrack().IsMuted() || subTrack.DownTrack() == nil { + if subTrack.IsMuted() || subTrack.MediaTrack().IsMuted() { continue } score := subTrack.DownTrack().GetConnectionScore() - subscriberScores[subTrack.ID()] = score + // subscriberScores[subTrack.ID()] = score totalScore += score numTracks++ } @@ -1287,18 +1286,10 @@ func (p *ParticipantImpl) subscriberRTCPWorker() { var srs []rtcp.Packet var sd []rtcp.SourceDescriptionChunk subscribedTracks := p.SubscriptionManager.GetSubscribedTracks() - downtracks := p.audioForwarder.GetDowntracks() p.lock.RLock() for _, subTrack := range subscribedTracks { - if subTrack.DownTrack() == nil { - continue - } - downtracks = append(downtracks, subTrack.DownTrack()) - } - - for _, dt := range downtracks { - sr := dt.CreateSenderReport() - chunks := dt.CreateSourceDescriptionChunks() + sr := subTrack.DownTrack().CreateSenderReport() + chunks := subTrack.DownTrack().CreateSourceDescriptionChunks() if sr == nil || chunks == nil { continue } diff --git a/pkg/rtc/subscribedtrack.go b/pkg/rtc/subscribedtrack.go index c3e6cee92..667705667 100644 --- a/pkg/rtc/subscribedtrack.go +++ b/pkg/rtc/subscribedtrack.go @@ -12,7 +12,6 @@ import ( "github.com/livekit/protocol/logger" "github.com/livekit/livekit-server/pkg/rtc/types" - "github.com/livekit/livekit-server/pkg/sfu" "github.com/livekit/livekit-server/pkg/sfu/buffer" ) @@ -26,8 +25,9 @@ type SubscribedTrackParams struct { PublisherVersion uint32 Subscriber types.LocalParticipant MediaTrack types.MediaTrack - DownTrack *sfu.DownTrack + DownTrack types.DownTrack AdaptiveStream bool + IsMuxedTrack bool } type SubscribedTrack struct { @@ -126,6 +126,10 @@ func (t *SubscribedTrack) IsBound() bool { return t.bound.Load() } +func (t *SubscribedTrack) IsMuxedTrack() bool { + return t.params.IsMuxedTrack +} + func (t *SubscribedTrack) ID() livekit.TrackID { return livekit.TrackID(t.params.DownTrack.ID()) } @@ -154,7 +158,7 @@ func (t *SubscribedTrack) Subscriber() types.LocalParticipant { return t.params.Subscriber } -func (t *SubscribedTrack) DownTrack() *sfu.DownTrack { +func (t *SubscribedTrack) DownTrack() types.DownTrack { return t.params.DownTrack } diff --git a/pkg/rtc/subscriptionmanager.go b/pkg/rtc/subscriptionmanager.go index d60a3766b..2e2c32458 100644 --- a/pkg/rtc/subscriptionmanager.go +++ b/pkg/rtc/subscriptionmanager.go @@ -25,7 +25,6 @@ import ( "github.com/pion/webrtc/v3/pkg/rtcerr" "github.com/livekit/livekit-server/pkg/rtc/types" - "github.com/livekit/livekit-server/pkg/sfu" "github.com/livekit/livekit-server/pkg/telemetry" "github.com/livekit/protocol/livekit" "github.com/livekit/protocol/logger" @@ -91,7 +90,7 @@ func (m *SubscriptionManager) Close(willBeResumed bool) { <-m.doneCh subTracks := m.GetSubscribedTracks() - downTracksToClose := make([]*sfu.DownTrack, 0, len(subTracks)) + downTracksToClose := make([]types.DownTrack, 0, len(subTracks)) for _, st := range subTracks { dt := st.DownTrack() // nil check exists primarily for tests diff --git a/pkg/rtc/transport.go b/pkg/rtc/transport.go index a02e8418d..8cc1ae0cb 100644 --- a/pkg/rtc/transport.go +++ b/pkg/rtc/transport.go @@ -1051,11 +1051,13 @@ func (t *PCTransport) AddTrackToStreamAllocator(subTrack types.SubscribedTrack) return } - t.streamAllocator.AddTrack(subTrack.DownTrack(), sfu.AddTrackParams{ - Source: subTrack.MediaTrack().Source(), - IsSimulcast: subTrack.MediaTrack().IsSimulcast(), - PublisherID: subTrack.MediaTrack().PublisherID(), - }) + if dt, ok := subTrack.DownTrack().(*sfu.DownTrack); ok { + t.streamAllocator.AddTrack(dt, sfu.AddTrackParams{ + Source: subTrack.MediaTrack().Source(), + IsSimulcast: subTrack.MediaTrack().IsSimulcast(), + PublisherID: subTrack.MediaTrack().PublisherID(), + }) + } } func (t *PCTransport) RemoveTrackFromStreamAllocator(subTrack types.SubscribedTrack) { @@ -1063,7 +1065,9 @@ func (t *PCTransport) RemoveTrackFromStreamAllocator(subTrack types.SubscribedTr return } - t.streamAllocator.RemoveTrack(subTrack.DownTrack()) + if dt, ok := subTrack.DownTrack().(*sfu.DownTrack); ok { + t.streamAllocator.RemoveTrack(dt) + } } func (t *PCTransport) GetICEConnectionType() types.ICEConnectionType { diff --git a/pkg/rtc/types/interfaces.go b/pkg/rtc/types/interfaces.go index 969bf054d..21f02f9c7 100644 --- a/pkg/rtc/types/interfaces.go +++ b/pkg/rtc/types/interfaces.go @@ -427,7 +427,7 @@ type SubscribedTrack interface { SubscriberID() livekit.ParticipantID SubscriberIdentity() livekit.ParticipantIdentity Subscriber() LocalParticipant - DownTrack() *sfu.DownTrack + DownTrack() DownTrack MediaTrack() MediaTrack RTPSender() *webrtc.RTPSender IsMuted() bool @@ -436,6 +436,25 @@ type SubscribedTrack interface { // selects appropriate video layer according to subscriber preferences UpdateVideoLayer() NeedsNegotiation() bool + IsMuxedTrack() bool +} + +type DownTrack interface { + ID() string + CloseWithFlush(flush bool) + Resync() + Codec() webrtc.RTPCodecCapability + DebugInfo() map[string]interface{} + GetConnectionScore() float32 + SetActivePaddingOnMuteUpTrack() + SetConnected() + SetMaxSpatialLayer(spatialLayer int32) + SetMaxTemporalLayer(temporalLayer int32) + Kind() webrtc.RTPCodecType + Mute(muted bool) + PubMute(pubMuted bool) + CreateSenderReport() *rtcp.SenderReport + CreateSourceDescriptionChunks() []rtcp.SourceDescriptionChunk } type ChangeNotifier interface { diff --git a/pkg/rtc/types/typesfakes/fake_local_participant.go b/pkg/rtc/types/typesfakes/fake_local_participant.go index 3baa5c697..52a1937d0 100644 --- a/pkg/rtc/types/typesfakes/fake_local_participant.go +++ b/pkg/rtc/types/typesfakes/fake_local_participant.go @@ -23,6 +23,13 @@ type FakeLocalParticipant struct { arg1 webrtc.ICECandidateInit arg2 livekit.SignalTarget } + AddMuxAudioTrackStub func(livekit.ParticipantID, livekit.TrackID, sfu.TrackReceiver) + addMuxAudioTrackMutex sync.RWMutex + addMuxAudioTrackArgsForCall []struct { + arg1 livekit.ParticipantID + arg2 livekit.TrackID + arg3 sfu.TrackReceiver + } AddTrackStub func(*livekit.AddTrackRequest) addTrackMutex sync.RWMutex addTrackArgsForCall []struct { @@ -501,6 +508,11 @@ type FakeLocalParticipant struct { protocolVersionReturnsOnCall map[int]struct { result1 types.ProtocolVersion } + RemoveMuxAudioTrackStub func(livekit.TrackID) + removeMuxAudioTrackMutex sync.RWMutex + removeMuxAudioTrackArgsForCall []struct { + arg1 livekit.TrackID + } RemovePublishedTrackStub func(types.MediaTrack, bool, bool) removePublishedTrackMutex sync.RWMutex removePublishedTrackArgsForCall []struct { @@ -852,6 +864,40 @@ func (fake *FakeLocalParticipant) AddICECandidateArgsForCall(i int) (webrtc.ICEC return argsForCall.arg1, argsForCall.arg2 } +func (fake *FakeLocalParticipant) AddMuxAudioTrack(arg1 livekit.ParticipantID, arg2 livekit.TrackID, arg3 sfu.TrackReceiver) { + fake.addMuxAudioTrackMutex.Lock() + fake.addMuxAudioTrackArgsForCall = append(fake.addMuxAudioTrackArgsForCall, struct { + arg1 livekit.ParticipantID + arg2 livekit.TrackID + arg3 sfu.TrackReceiver + }{arg1, arg2, arg3}) + stub := fake.AddMuxAudioTrackStub + fake.recordInvocation("AddMuxAudioTrack", []interface{}{arg1, arg2, arg3}) + fake.addMuxAudioTrackMutex.Unlock() + if stub != nil { + fake.AddMuxAudioTrackStub(arg1, arg2, arg3) + } +} + +func (fake *FakeLocalParticipant) AddMuxAudioTrackCallCount() int { + fake.addMuxAudioTrackMutex.RLock() + defer fake.addMuxAudioTrackMutex.RUnlock() + return len(fake.addMuxAudioTrackArgsForCall) +} + +func (fake *FakeLocalParticipant) AddMuxAudioTrackCalls(stub func(livekit.ParticipantID, livekit.TrackID, sfu.TrackReceiver)) { + fake.addMuxAudioTrackMutex.Lock() + defer fake.addMuxAudioTrackMutex.Unlock() + fake.AddMuxAudioTrackStub = stub +} + +func (fake *FakeLocalParticipant) AddMuxAudioTrackArgsForCall(i int) (livekit.ParticipantID, livekit.TrackID, sfu.TrackReceiver) { + fake.addMuxAudioTrackMutex.RLock() + defer fake.addMuxAudioTrackMutex.RUnlock() + argsForCall := fake.addMuxAudioTrackArgsForCall[i] + return argsForCall.arg1, argsForCall.arg2, argsForCall.arg3 +} + func (fake *FakeLocalParticipant) AddTrack(arg1 *livekit.AddTrackRequest) { fake.addTrackMutex.Lock() fake.addTrackArgsForCall = append(fake.addTrackArgsForCall, struct { @@ -3430,6 +3476,38 @@ func (fake *FakeLocalParticipant) ProtocolVersionReturnsOnCall(i int, result1 ty }{result1} } +func (fake *FakeLocalParticipant) RemoveMuxAudioTrack(arg1 livekit.TrackID) { + fake.removeMuxAudioTrackMutex.Lock() + fake.removeMuxAudioTrackArgsForCall = append(fake.removeMuxAudioTrackArgsForCall, struct { + arg1 livekit.TrackID + }{arg1}) + stub := fake.RemoveMuxAudioTrackStub + fake.recordInvocation("RemoveMuxAudioTrack", []interface{}{arg1}) + fake.removeMuxAudioTrackMutex.Unlock() + if stub != nil { + fake.RemoveMuxAudioTrackStub(arg1) + } +} + +func (fake *FakeLocalParticipant) RemoveMuxAudioTrackCallCount() int { + fake.removeMuxAudioTrackMutex.RLock() + defer fake.removeMuxAudioTrackMutex.RUnlock() + return len(fake.removeMuxAudioTrackArgsForCall) +} + +func (fake *FakeLocalParticipant) RemoveMuxAudioTrackCalls(stub func(livekit.TrackID)) { + fake.removeMuxAudioTrackMutex.Lock() + defer fake.removeMuxAudioTrackMutex.Unlock() + fake.RemoveMuxAudioTrackStub = stub +} + +func (fake *FakeLocalParticipant) RemoveMuxAudioTrackArgsForCall(i int) livekit.TrackID { + fake.removeMuxAudioTrackMutex.RLock() + defer fake.removeMuxAudioTrackMutex.RUnlock() + argsForCall := fake.removeMuxAudioTrackArgsForCall[i] + return argsForCall.arg1 +} + func (fake *FakeLocalParticipant) RemovePublishedTrack(arg1 types.MediaTrack, arg2 bool, arg3 bool) { fake.removePublishedTrackMutex.Lock() fake.removePublishedTrackArgsForCall = append(fake.removePublishedTrackArgsForCall, struct { @@ -5174,6 +5252,8 @@ func (fake *FakeLocalParticipant) Invocations() map[string][][]interface{} { defer fake.invocationsMutex.RUnlock() fake.addICECandidateMutex.RLock() defer fake.addICECandidateMutex.RUnlock() + fake.addMuxAudioTrackMutex.RLock() + defer fake.addMuxAudioTrackMutex.RUnlock() fake.addTrackMutex.RLock() defer fake.addTrackMutex.RUnlock() fake.addTrackToSubscriberMutex.RLock() @@ -5284,6 +5364,8 @@ func (fake *FakeLocalParticipant) Invocations() map[string][][]interface{} { defer fake.onTrackUpdatedMutex.RUnlock() fake.protocolVersionMutex.RLock() defer fake.protocolVersionMutex.RUnlock() + fake.removeMuxAudioTrackMutex.RLock() + defer fake.removeMuxAudioTrackMutex.RUnlock() fake.removePublishedTrackMutex.RLock() defer fake.removePublishedTrackMutex.RUnlock() fake.removeTrackFromSubscriberMutex.RLock() diff --git a/pkg/rtc/types/typesfakes/fake_subscribed_track.go b/pkg/rtc/types/typesfakes/fake_subscribed_track.go index b14ceb2c4..3b8b8313b 100644 --- a/pkg/rtc/types/typesfakes/fake_subscribed_track.go +++ b/pkg/rtc/types/typesfakes/fake_subscribed_track.go @@ -5,7 +5,6 @@ import ( "sync" "github.com/livekit/livekit-server/pkg/rtc/types" - "github.com/livekit/livekit-server/pkg/sfu" "github.com/livekit/protocol/livekit" webrtc "github.com/pion/webrtc/v3" ) @@ -21,15 +20,15 @@ type FakeSubscribedTrack struct { closeArgsForCall []struct { arg1 bool } - DownTrackStub func() *sfu.DownTrack + DownTrackStub func() types.DownTrack downTrackMutex sync.RWMutex downTrackArgsForCall []struct { } downTrackReturns struct { - result1 *sfu.DownTrack + result1 types.DownTrack } downTrackReturnsOnCall map[int]struct { - result1 *sfu.DownTrack + result1 types.DownTrack } IDStub func() livekit.TrackID iDMutex sync.RWMutex @@ -61,6 +60,16 @@ type FakeSubscribedTrack struct { isMutedReturnsOnCall map[int]struct { result1 bool } + IsMuxedTrackStub func() bool + isMuxedTrackMutex sync.RWMutex + isMuxedTrackArgsForCall []struct { + } + isMuxedTrackReturns struct { + result1 bool + } + isMuxedTrackReturnsOnCall map[int]struct { + result1 bool + } MediaTrackStub func() types.MediaTrack mediaTrackMutex sync.RWMutex mediaTrackArgsForCall []struct { @@ -238,7 +247,7 @@ func (fake *FakeSubscribedTrack) CloseArgsForCall(i int) bool { return argsForCall.arg1 } -func (fake *FakeSubscribedTrack) DownTrack() *sfu.DownTrack { +func (fake *FakeSubscribedTrack) DownTrack() types.DownTrack { fake.downTrackMutex.Lock() ret, specificReturn := fake.downTrackReturnsOnCall[len(fake.downTrackArgsForCall)] fake.downTrackArgsForCall = append(fake.downTrackArgsForCall, struct { @@ -262,32 +271,32 @@ func (fake *FakeSubscribedTrack) DownTrackCallCount() int { return len(fake.downTrackArgsForCall) } -func (fake *FakeSubscribedTrack) DownTrackCalls(stub func() *sfu.DownTrack) { +func (fake *FakeSubscribedTrack) DownTrackCalls(stub func() types.DownTrack) { fake.downTrackMutex.Lock() defer fake.downTrackMutex.Unlock() fake.DownTrackStub = stub } -func (fake *FakeSubscribedTrack) DownTrackReturns(result1 *sfu.DownTrack) { +func (fake *FakeSubscribedTrack) DownTrackReturns(result1 types.DownTrack) { fake.downTrackMutex.Lock() defer fake.downTrackMutex.Unlock() fake.DownTrackStub = nil fake.downTrackReturns = struct { - result1 *sfu.DownTrack + result1 types.DownTrack }{result1} } -func (fake *FakeSubscribedTrack) DownTrackReturnsOnCall(i int, result1 *sfu.DownTrack) { +func (fake *FakeSubscribedTrack) DownTrackReturnsOnCall(i int, result1 types.DownTrack) { fake.downTrackMutex.Lock() defer fake.downTrackMutex.Unlock() fake.DownTrackStub = nil if fake.downTrackReturnsOnCall == nil { fake.downTrackReturnsOnCall = make(map[int]struct { - result1 *sfu.DownTrack + result1 types.DownTrack }) } fake.downTrackReturnsOnCall[i] = struct { - result1 *sfu.DownTrack + result1 types.DownTrack }{result1} } @@ -450,6 +459,59 @@ func (fake *FakeSubscribedTrack) IsMutedReturnsOnCall(i int, result1 bool) { }{result1} } +func (fake *FakeSubscribedTrack) IsMuxedTrack() bool { + fake.isMuxedTrackMutex.Lock() + ret, specificReturn := fake.isMuxedTrackReturnsOnCall[len(fake.isMuxedTrackArgsForCall)] + fake.isMuxedTrackArgsForCall = append(fake.isMuxedTrackArgsForCall, struct { + }{}) + stub := fake.IsMuxedTrackStub + fakeReturns := fake.isMuxedTrackReturns + fake.recordInvocation("IsMuxedTrack", []interface{}{}) + fake.isMuxedTrackMutex.Unlock() + if stub != nil { + return stub() + } + if specificReturn { + return ret.result1 + } + return fakeReturns.result1 +} + +func (fake *FakeSubscribedTrack) IsMuxedTrackCallCount() int { + fake.isMuxedTrackMutex.RLock() + defer fake.isMuxedTrackMutex.RUnlock() + return len(fake.isMuxedTrackArgsForCall) +} + +func (fake *FakeSubscribedTrack) IsMuxedTrackCalls(stub func() bool) { + fake.isMuxedTrackMutex.Lock() + defer fake.isMuxedTrackMutex.Unlock() + fake.IsMuxedTrackStub = stub +} + +func (fake *FakeSubscribedTrack) IsMuxedTrackReturns(result1 bool) { + fake.isMuxedTrackMutex.Lock() + defer fake.isMuxedTrackMutex.Unlock() + fake.IsMuxedTrackStub = nil + fake.isMuxedTrackReturns = struct { + result1 bool + }{result1} +} + +func (fake *FakeSubscribedTrack) IsMuxedTrackReturnsOnCall(i int, result1 bool) { + fake.isMuxedTrackMutex.Lock() + defer fake.isMuxedTrackMutex.Unlock() + fake.IsMuxedTrackStub = nil + if fake.isMuxedTrackReturnsOnCall == nil { + fake.isMuxedTrackReturnsOnCall = make(map[int]struct { + result1 bool + }) + } + fake.isMuxedTrackReturnsOnCall[i] = struct { + result1 bool + }{result1} +} + func (fake *FakeSubscribedTrack) MediaTrack() types.MediaTrack { fake.mediaTrackMutex.Lock() ret, specificReturn := fake.mediaTrackReturnsOnCall[len(fake.mediaTrackArgsForCall)] @@ -1062,6 +1124,8 @@ func (fake *FakeSubscribedTrack) Invocations() map[string][][]interface{} { defer fake.isBoundMutex.RUnlock() fake.isMutedMutex.RLock() defer fake.isMutedMutex.RUnlock() + fake.isMuxedTrackMutex.RLock() + defer fake.isMuxedTrackMutex.RUnlock() fake.mediaTrackMutex.RLock() defer fake.mediaTrackMutex.RUnlock() fake.needsNegotiationMutex.RLock() diff --git a/pkg/sfu/audioselection/forwarder.go b/pkg/sfu/audioselection/forwarder.go index d3a7f5212..d567e94cb 100644 --- a/pkg/sfu/audioselection/forwarder.go +++ b/pkg/sfu/audioselection/forwarder.go @@ -162,13 +162,13 @@ func (f *SelectionForwarder) updateForward() { return f.sources[i].audioLevel > f.sources[j].audioLevel }) - var activeteSources, idleSources []*sourceInfo + var activateSources, idleSources []*sourceInfo for i, source := range f.sources { if i >= f.params.ActiveDowntracks && source.active { idleSources = append(idleSources, source) } else { if source.audioLevel > f.params.ActiveLevelThreshold && !source.active { - activeteSources = append(activeteSources, source) + activateSources = append(activateSources, source) } else if source.audioLevel <= f.params.ActiveLevelThreshold && source.active { idleSources = append(idleSources, source) } @@ -176,7 +176,7 @@ func (f *SelectionForwarder) updateForward() { } var forwardChanged bool - for _, source := range activeteSources { + for _, source := range activateSources { if len(f.idleDowntracks) == 0 && len(idleSources) > 0 { f.deactiveSource(idleSources[0]) idleSources = idleSources[1:] diff --git a/pkg/sfu/audioselection/nullaudiodowntrack.go b/pkg/sfu/audioselection/nullaudiodowntrack.go new file mode 100644 index 000000000..0ef177fba --- /dev/null +++ b/pkg/sfu/audioselection/nullaudiodowntrack.go @@ -0,0 +1,87 @@ +package audioselection + +import ( + "github.com/pion/rtcp" + "github.com/pion/webrtc/v3" + + "github.com/livekit/protocol/utils" +) + +// implements types.DownTrack + +/* +ID() string +CloseWithFlush(flush bool) +Resync() +Codec() webrtc.RTPCodecCapability +DebugInfo() map[string]interface{} +GetConnectionScore() float32 +SetActivePaddingOnMuteUpTrack() +SetConnected() +SetMaxSpatialLayer(spatialLayer int32) +SetMaxTemporalLayer(temporalLayer int32) +Kind() webrtc.RTPCodecType +Mute(muted bool) +PubMute(pubMuted bool) +CreateSenderReport() *rtcp.SenderReport +CreateSourceDescriptionChunks() []rtcp.SourceDescriptionChunk +*/ +type NullAudioDowntrack struct { + id string +} + +func NewNullAudioDowntrack() *NullAudioDowntrack { + return &NullAudioDowntrack{id: utils.NewGuid("TRAN_")} +} + +func (d *NullAudioDowntrack) ID() string { + return d.id +} + +func (d *NullAudioDowntrack) CloseWithFlush(flush bool) { +} + +func (d *NullAudioDowntrack) Resync() { +} + +func (d *NullAudioDowntrack) Codec() webrtc.RTPCodecCapability { + return webrtc.RTPCodecCapability{MimeType: webrtc.MimeTypeOpus, ClockRate: 48000, Channels: 2, SDPFmtpLine: "minptime=10;useinbandfec=1"} +} + +func (d *NullAudioDowntrack) DebugInfo() map[string]interface{} { + return nil +} + +func (d *NullAudioDowntrack) GetConnectionScore() float32 { + return 5.0 +} + +func (d *NullAudioDowntrack) SetActivePaddingOnMuteUpTrack() { +} + +func (d *NullAudioDowntrack) SetConnected() { +} + +func (d *NullAudioDowntrack) SetMaxSpatialLayer(spatialLayer int32) { +} + +func (d *NullAudioDowntrack) SetMaxTemporalLayer(temporalLayer int32) { +} + +func (d *NullAudioDowntrack) Kind() webrtc.RTPCodecType { + return webrtc.RTPCodecTypeAudio +} + +func (d *NullAudioDowntrack) Mute(muted bool) { +} + +func (d *NullAudioDowntrack) PubMute(pubMuted bool) { +} + +func (d *NullAudioDowntrack) CreateSenderReport() *rtcp.SenderReport { + return nil +} + +func (d *NullAudioDowntrack) CreateSourceDescriptionChunks() []rtcp.SourceDescriptionChunk { + return nil +}