mirror of
https://github.com/livekit/livekit.git
synced 2026-08-27 22:34:25 +00:00
Make a copy of available layers in forwarder. (#709)
Should be okay as that is usually only three elements max. Stream tracker + manager send available layers. As stream trackers run in different go routines, need to think some more about chances of out-of-order operations. But, making a copy in forwarder so that it is not referncing moving data. Also, when setting up provisional allocation, make a copy (that was the intention, i. e. a provisional allocation should have stable data to work with).
This commit is contained in:
+39
-11
@@ -322,7 +322,12 @@ func (f *Forwarder) UpTrackLayersChange(availableLayers []int32) {
|
||||
f.lock.Lock()
|
||||
defer f.lock.Unlock()
|
||||
|
||||
f.availableLayers = availableLayers
|
||||
if len(availableLayers) > 0 {
|
||||
f.availableLayers = make([]int32, len(availableLayers))
|
||||
copy(f.availableLayers, availableLayers)
|
||||
} else {
|
||||
f.availableLayers = nil
|
||||
}
|
||||
}
|
||||
|
||||
func (f *Forwarder) getOptimalBandwidthNeeded(brs Bitrates) int64 {
|
||||
@@ -481,11 +486,15 @@ func (f *Forwarder) AllocateOptimal(brs Bitrates) VideoAllocation {
|
||||
change: change,
|
||||
bandwidthRequested: bandwidthRequested,
|
||||
bandwidthDelta: bandwidthRequested - f.lastAllocation.bandwidthRequested,
|
||||
availableLayers: f.availableLayers,
|
||||
bitrates: brs,
|
||||
targetLayers: targetLayers,
|
||||
distanceToDesired: f.getDistanceToDesired(brs, targetLayers),
|
||||
}
|
||||
if len(f.availableLayers) > 0 {
|
||||
f.lastAllocation.availableLayers = make([]int32, len(f.availableLayers))
|
||||
copy(f.lastAllocation.availableLayers, f.availableLayers)
|
||||
}
|
||||
|
||||
f.setTargetLayers(f.lastAllocation.targetLayers)
|
||||
if f.targetLayers == InvalidLayers {
|
||||
f.resyncLocked()
|
||||
@@ -499,11 +508,14 @@ func (f *Forwarder) ProvisionalAllocatePrepare(bitrates Bitrates) {
|
||||
defer f.lock.Unlock()
|
||||
|
||||
f.provisional = &VideoAllocationProvisional{
|
||||
layers: InvalidLayers,
|
||||
muted: f.muted,
|
||||
bitrates: bitrates,
|
||||
availableLayers: f.availableLayers,
|
||||
maxLayers: f.maxLayers,
|
||||
layers: InvalidLayers,
|
||||
muted: f.muted,
|
||||
bitrates: bitrates,
|
||||
maxLayers: f.maxLayers,
|
||||
}
|
||||
if len(f.availableLayers) > 0 {
|
||||
f.provisional.availableLayers = make([]int32, len(f.availableLayers))
|
||||
copy(f.provisional.availableLayers, f.availableLayers)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -776,11 +788,15 @@ func (f *Forwarder) ProvisionalAllocateCommit() VideoAllocation {
|
||||
change: change,
|
||||
bandwidthRequested: bandwidthRequested,
|
||||
bandwidthDelta: bandwidthRequested - f.lastAllocation.bandwidthRequested,
|
||||
availableLayers: f.provisional.availableLayers,
|
||||
bitrates: f.provisional.bitrates,
|
||||
targetLayers: f.provisional.layers,
|
||||
distanceToDesired: f.getDistanceToDesired(f.provisional.bitrates, f.provisional.layers),
|
||||
}
|
||||
if len(f.availableLayers) > 0 {
|
||||
f.lastAllocation.availableLayers = make([]int32, len(f.availableLayers))
|
||||
copy(f.lastAllocation.availableLayers, f.provisional.availableLayers)
|
||||
}
|
||||
|
||||
f.setTargetLayers(f.lastAllocation.targetLayers)
|
||||
if f.targetLayers == InvalidLayers {
|
||||
f.resyncLocked()
|
||||
@@ -842,11 +858,15 @@ func (f *Forwarder) AllocateNextHigher(availableChannelCapacity int64, brs Bitra
|
||||
change: VideoStreamingChangeNone,
|
||||
bandwidthRequested: bandwidthRequested,
|
||||
bandwidthDelta: bandwidthRequested - alreadyAllocated,
|
||||
availableLayers: f.availableLayers,
|
||||
bitrates: brs,
|
||||
targetLayers: targetLayers,
|
||||
distanceToDesired: f.getDistanceToDesired(brs, targetLayers),
|
||||
}
|
||||
if len(f.availableLayers) > 0 {
|
||||
f.lastAllocation.availableLayers = make([]int32, len(f.availableLayers))
|
||||
copy(f.lastAllocation.availableLayers, f.availableLayers)
|
||||
}
|
||||
|
||||
f.setTargetLayers(f.lastAllocation.targetLayers)
|
||||
return f.lastAllocation, true
|
||||
}
|
||||
@@ -882,11 +902,15 @@ func (f *Forwarder) AllocateNextHigher(availableChannelCapacity int64, brs Bitra
|
||||
change: change,
|
||||
bandwidthRequested: bandwidthRequested,
|
||||
bandwidthDelta: bandwidthRequested - alreadyAllocated,
|
||||
availableLayers: f.availableLayers,
|
||||
bitrates: brs,
|
||||
targetLayers: targetLayers,
|
||||
distanceToDesired: f.getDistanceToDesired(brs, targetLayers),
|
||||
}
|
||||
if len(f.availableLayers) > 0 {
|
||||
f.lastAllocation.availableLayers = make([]int32, len(f.availableLayers))
|
||||
copy(f.lastAllocation.availableLayers, f.availableLayers)
|
||||
}
|
||||
|
||||
f.setTargetLayers(f.lastAllocation.targetLayers)
|
||||
return f.lastAllocation, true
|
||||
}
|
||||
@@ -985,11 +1009,15 @@ func (f *Forwarder) Pause(brs Bitrates) VideoAllocation {
|
||||
change: change,
|
||||
bandwidthRequested: 0,
|
||||
bandwidthDelta: 0 - f.lastAllocation.bandwidthRequested,
|
||||
availableLayers: f.availableLayers,
|
||||
bitrates: brs,
|
||||
targetLayers: InvalidLayers,
|
||||
distanceToDesired: f.getDistanceToDesired(brs, InvalidLayers),
|
||||
}
|
||||
if len(f.availableLayers) > 0 {
|
||||
f.lastAllocation.availableLayers = make([]int32, len(f.availableLayers))
|
||||
copy(f.lastAllocation.availableLayers, f.availableLayers)
|
||||
}
|
||||
|
||||
f.setTargetLayers(f.lastAllocation.targetLayers)
|
||||
if f.targetLayers == InvalidLayers {
|
||||
f.resyncLocked()
|
||||
|
||||
@@ -157,7 +157,7 @@ func TestForwarderUpTrackLayersChange(t *testing.T) {
|
||||
|
||||
availableLayers = []int32{}
|
||||
f.UpTrackLayersChange(availableLayers)
|
||||
require.Equal(t, availableLayers, f.availableLayers)
|
||||
require.Nil(t, f.availableLayers)
|
||||
}
|
||||
|
||||
func TestForwarderAllocate(t *testing.T) {
|
||||
|
||||
Reference in New Issue
Block a user