mirror of
https://github.com/livekit/livekit.git
synced 2026-08-28 23:01:19 +00:00
worked for merge master
This commit is contained in:
@@ -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())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+5
-14
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
+10
-6
@@ -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 {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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:]
|
||||
|
||||
@@ -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
|
||||
}
|
||||
Reference in New Issue
Block a user