Simplifying forwarding logic a bit (#1349)

* Notes on wht to do

- Should targetLayers be altered while doing opportunistic locking
- Should targetLayers be altered in any other path than stream allocator path?
- Lock to layer as long as it is <= opportunistic layer
- When not congested, opportunistic can be highest
- When congested, opportunistic could be nil or lowest if paused is not allowed
- When muting, can we hold on to current layers (or keep it as previous) and
  restore on unmute.
- Store current/target in forwarder state and restore on seeding
- Watch for looking for targetLayers, etc. when looking to insert padding
  packets. There may be an assumption about restarting on key frame and hence
  okay to insert padding when target layers are invalid. This may not be true
  any more when doing opportunistic forwarding.
- Can we distinguish between publisher mute or dynacast (i. e. publisher side
  stopping) vs subscriber mute and do something useful? Publisher side mute
  could mean continuity in sequence numbers on a restart (might be able to
  catch it with opportunistic forwarding). But, there is the challenge of
  unmute from publisher via signalling channel vs media. If media is arriving,
  should subscribers do opportunistic forwarding before publisher mute state
  update happens?
- Maybe introduce a mode where forwarding continues to a frame end (of course
  with a time limit just in case the end of frame packet is lost) and then
  insert silence/padding packets?
- Ensure that audio blank frame insertion does not suffer from frame boundary
  issues.

* pub/sub mute separate + more notes on things to check

* WIP commit, more notes

* WIP commit

* WIP commit

* WIP commit

* WIP commit

* WIP commit

* WIP commit

* clean up

* slightly better comments

* Do not stop on unmute

* do not inject blank frames when pub muted

* do not forward on audio publisher mute
This commit is contained in:
Raja Subramanian
2023-02-01 21:57:53 +05:30
committed by GitHub
parent 7e5ba6a3b0
commit ffadb94e3a
8 changed files with 485 additions and 202 deletions
+2 -2
View File
@@ -227,8 +227,8 @@ func (t *SubscribedTrack) SetRTPSender(sender *webrtc.RTPSender) {
}
func (t *SubscribedTrack) updateDownTrackMute() {
muted := t.subMuted.Load() || t.pubMuted.Load()
t.DownTrack().Mute(muted)
t.DownTrack().Mute(t.subMuted.Load())
t.DownTrack().PubMute(t.pubMuted.Load())
}
func (t *SubscribedTrack) spatialLayerFromSettings(settings *livekit.UpdateTrackSettings) int32 {
+54 -4
View File
@@ -5,6 +5,7 @@ import (
"errors"
"fmt"
"io"
"math"
"strings"
"sync"
"time"
@@ -249,6 +250,11 @@ func NewDownTrack(
codec: codecs[0].RTPCodecCapability,
}
d.forwarder = NewForwarder(d.kind, d.logger, d.receiver.GetReferenceLayerRTPTimestamp)
d.forwarder.OnParkedLayersExpired(func() {
if d.onSubscriptionChanged != nil {
d.onSubscriptionChanged(d)
}
})
d.connectionStats = connectionquality.NewConnectionStats(connectionquality.ConnectionStatsParams{
MimeType: codecs[0].MimeType, // LK-TODO have to notify on codec change
@@ -348,7 +354,12 @@ func (d *DownTrack) Unbind(_ webrtc.TrackLocalContext) error {
}
func (d *DownTrack) TrackInfoAvailable() {
d.connectionStats.Start(d.receiver.TrackInfo())
ti := d.receiver.TrackInfo()
if ti == nil {
return
}
d.forwarder.SetNumAdvertisedLayers(int32(math.Max(1, float64(len(ti.Layers)))))
d.connectionStats.Start(ti)
}
// ID is the unique identifier for this Track. This should be unique for the
@@ -427,7 +438,7 @@ func (d *DownTrack) stopKeyFrameRequester() {
}
func (d *DownTrack) keyFrameRequester(generation uint32, layer int32) {
if d.IsClosed() {
if d.IsClosed() || layer == InvalidLayerSpatial {
return
}
interval := 2 * d.rtpStats.GetRtt()
@@ -649,14 +660,44 @@ func (d *DownTrack) WritePaddingRTP(bytesToSend int, paddingOnMute bool) int {
return bytesSent
}
// Mute enables or disables media forwarding
// Mute enables or disables media forwarding - subscriber triggered
func (d *DownTrack) Mute(muted bool) {
changed, maxLayers := d.forwarder.Mute(muted)
d.handleMute(muted, false, changed, maxLayers)
}
// PubMute enables or disables media forwarding - publisher side
func (d *DownTrack) PubMute(pubMuted bool) {
changed, maxLayers := d.forwarder.PubMute(pubMuted)
d.handleMute(pubMuted, true, changed, maxLayers)
}
func (d *DownTrack) handleMute(muted bool, isPub bool, changed bool, maxLayers VideoLayers) {
if !changed {
return
}
if d.onMaxLayerChanged != nil && d.kind == webrtc.RTPCodecTypeVideo {
//
// Subscriber mute changes trigger a max layer notification.
// That could result in encoding layers getting turned on/off on publisher side
// (depending on aggregate layer requirements of all subscribers of the track).
//
// Publisher mute changes should not trigger notification.
// If publisher turns off all layers because of subscribers indicating
// no layers required due to publisher mute (bit of circular dependency),
// there will be a delay in layers turning back on when unmute happens.
// Unmute path will require
// 1. unmute signalling out-of-band from publisher received by down track(s)
// 2. down track(s) notifying max layer
// 3. out-of-band notification about max layer sent back to the publisher
// 4. publisher starts layer(s)
// Ideally, on publisher mute, whatever layers were active reamin active and
// can be restarted by publisher immediately on unmute.
//
// Note that while publisher mute is active, subscriber changes can also happen
// and that could turn on/off layers on publisher side.
//
if !isPub && d.onMaxLayerChanged != nil && d.kind == webrtc.RTPCodecTypeVideo {
notifyLayer := InvalidLayerSpatial
if !muted {
//
@@ -676,6 +717,14 @@ func (d *DownTrack) Mute(muted bool) {
// when muting, send a few silence frames to ensure residual noise does not
// put the comfort noise generator on decoder side in a bad state where it
// generates noise that is not so comfortable.
//
// One possibility is not to inject blank frames when publisher is muted
// and let forwarding continue. When publisher is muted, unless the media
// stream is stopped, publisher will send silence frames which should have
// comfort noise information. But, in case the publisher stops at an
// inopportune frame (due to media stream stop or injecting audio from a file),
// the decoder could be in a noisy state. So, inject blank frames on publisher
// mute too.
d.blankFramesGeneration.Inc()
if d.kind == webrtc.RTPCodecTypeAudio && muted {
d.writeBlankFrameRTP(RTPBlankFramesMuteSeconds, d.blankFramesGeneration.Load())
@@ -1452,6 +1501,7 @@ func (d *DownTrack) DebugInfo() map[string]interface{} {
"MimeType": d.codec.MimeType,
"Bound": d.bound.Load(),
"Muted": d.forwarder.IsMuted(),
"PubMuted": d.forwarder.IsPubMuted(),
"CurrentSpatialLayer": d.forwarder.CurrentLayers().Spatial,
"Stats": stats,
}
+305 -147
View File
@@ -5,6 +5,7 @@ import (
"math"
"strings"
"sync"
"time"
"github.com/pion/webrtc/v3"
@@ -16,9 +17,10 @@ import (
// Forwarder
const (
FlagPauseOnDowngrade = true
FlagFilterRTX = true
TransitionCostSpatial = 10
FlagPauseOnDowngrade = true
FlagFilterRTX = true
TransitionCostSpatial = 10
ParkedLayersWaitDuration = 2 * time.Second
)
// -------------------------------------------------------------------
@@ -61,6 +63,7 @@ type VideoAllocationState int
const (
VideoAllocationStateNone VideoAllocationState = iota
VideoAllocationStateMuted
VideoAllocationStatePubMuted
VideoAllocationStateFeedDry
VideoAllocationStateAwaitingMeasurement
VideoAllocationStateOptimal
@@ -73,6 +76,8 @@ func (v VideoAllocationState) String() string {
return "NONE"
case VideoAllocationStateMuted:
return "MUTED"
case VideoAllocationStatePubMuted:
return "PUB_MUTED"
case VideoAllocationStateFeedDry:
return "FEED_DRY"
case VideoAllocationStateAwaitingMeasurement:
@@ -113,10 +118,13 @@ var (
type VideoAllocationProvisional struct {
muted bool
pubMuted bool
bitrates Bitrates
availableLayers []int32
exemptedLayers []int32
maxLayers VideoLayers
currentLayers VideoLayers
parkedLayers VideoLayers
allocatedLayers VideoLayers
}
@@ -151,6 +159,7 @@ type TranslationParams struct {
}
// -------------------------------------------------------------------
type VideoLayers = buffer.VideoLayer
const (
@@ -188,15 +197,20 @@ type Forwarder struct {
logger logger.Logger
getReferenceLayerRTPTimestamp func(ts uint32, layer int32, referenceLayer int32) (uint32, error)
muted bool
numAdvertisedLayers int32
muted bool
pubMuted bool
started bool
lastSSRC uint32
referenceLayerSpatial int32
maxLayers VideoLayers
currentLayers VideoLayers
targetLayers VideoLayers
maxLayers VideoLayers
currentLayers VideoLayers
targetLayers VideoLayers
parkedLayers VideoLayers // layers that can resume without key frame
parkedLayersTimer *time.Timer
provisional *VideoAllocationProvisional
@@ -211,6 +225,8 @@ type Forwarder struct {
isTemporalSupported bool
ddLayerSelector *DDVideoLayerSelector
onParkedLayersExpired func()
}
func NewForwarder(
@@ -225,9 +241,10 @@ func NewForwarder(
referenceLayerSpatial: InvalidLayerSpatial,
// start off with nothing, let streamallocator set things
// start off with nothing, let streamallocator/opportunistic forwarder set the target
currentLayers: InvalidLayers,
targetLayers: InvalidLayers,
parkedLayers: InvalidLayers,
lastAllocation: VideoAllocationDefault,
@@ -243,6 +260,27 @@ func NewForwarder(
return f
}
func (f *Forwarder) SetNumAdvertisedLayers(numAdvertisedLayers int32) {
f.lock.Lock()
defer f.lock.Unlock()
f.numAdvertisedLayers = numAdvertisedLayers
}
func (f *Forwarder) OnParkedLayersExpired(fn func()) {
f.lock.Lock()
defer f.lock.Unlock()
f.onParkedLayersExpired = fn
}
func (f *Forwarder) getOnParkedLayersExpired() func() {
f.lock.RLock()
defer f.lock.RUnlock()
return f.onParkedLayersExpired
}
func (f *Forwarder) DetermineCodec(codec webrtc.RTPCodecCapability) {
f.lock.Lock()
defer f.lock.Unlock()
@@ -258,7 +296,7 @@ func (f *Forwarder) DetermineCodec(codec webrtc.RTPCodecCapability) {
f.vp8Munger = NewVP8Munger(f.logger)
case "video/av1":
// TODO : we only enable dd layer selector for av1 now, at future we can
// enable it for vp9 too
// enable it for vp8 too
f.ddLayerSelector = NewDDVideoLayerSelector(f.logger)
}
}
@@ -292,13 +330,13 @@ func (f *Forwarder) SeedState(state ForwarderState) {
f.lock.Lock()
defer f.lock.Unlock()
f.referenceLayerSpatial = state.ReferenceLayerSpatial
f.rtpMunger.SeedLast(state.RTP)
if f.vp8Munger != nil {
f.vp8Munger.SeedLast(state.VP8)
}
f.started = true
f.referenceLayerSpatial = state.ReferenceLayerSpatial
}
func (f *Forwarder) Mute(muted bool) (bool, VideoLayers) {
@@ -327,6 +365,34 @@ func (f *Forwarder) IsMuted() bool {
return f.muted
}
func (f *Forwarder) PubMute(pubMuted bool) (bool, VideoLayers) {
f.lock.Lock()
defer f.lock.Unlock()
if f.pubMuted == pubMuted {
return false, f.maxLayers
}
f.logger.Debugw("setting forwarder pub mute", "pubMuted", pubMuted)
f.pubMuted = pubMuted
// Do not resync on publisher mute as forwarding can continue on unmute using same layers.
// On unmute, park current layers as streaming can continue without a key frame when publisher starts the stream.
if !pubMuted && f.targetLayers.IsValid() && f.currentLayers.Spatial == f.targetLayers.Spatial {
f.setupParkedLayers(f.targetLayers)
f.currentLayers = InvalidLayers
}
return true, f.maxLayers
}
func (f *Forwarder) IsPubMuted() bool {
f.lock.RLock()
defer f.lock.RUnlock()
return f.pubMuted
}
func (f *Forwarder) SetMaxSpatialLayer(spatialLayer int32) (bool, VideoLayers, VideoLayers) {
f.lock.Lock()
defer f.lock.Unlock()
@@ -338,6 +404,8 @@ func (f *Forwarder) SetMaxSpatialLayer(spatialLayer int32) (bool, VideoLayers, V
f.logger.Infow("setting max spatial layer", "layer", spatialLayer)
f.maxLayers.Spatial = spatialLayer
f.clearParkedLayers()
return true, f.maxLayers, f.currentLayers
}
@@ -352,6 +420,8 @@ func (f *Forwarder) SetMaxTemporalLayer(temporalLayer int32) (bool, VideoLayers,
f.logger.Infow("setting max temporal layer", "layer", temporalLayer)
f.maxLayers.Temporal = temporalLayer
f.clearParkedLayers()
return true, f.maxLayers, f.currentLayers
}
@@ -387,7 +457,7 @@ func (f *Forwarder) GetForwardingStatus() ForwardingStatus {
f.lock.RLock()
defer f.lock.RUnlock()
if f.muted || len(f.availableLayers) == 0 {
if f.muted || f.pubMuted || len(f.availableLayers) == 0 {
return ForwardingStatusOptimal
}
@@ -406,7 +476,7 @@ func (f *Forwarder) IsReducedQuality() (int32, bool) {
f.lock.RLock()
defer f.lock.RUnlock()
if f.muted || len(f.availableLayers) == 0 || f.targetLayers.Spatial == InvalidLayerSpatial {
if f.muted || f.pubMuted || len(f.availableLayers) == 0 || f.targetLayers.Spatial == InvalidLayerSpatial {
return 0, false
}
@@ -423,7 +493,7 @@ func (f *Forwarder) IsReducedQuality() (int32, bool) {
distance = 0
}
return distance, f.lastAllocation.state == VideoAllocationStateDeficient
return distance, f.isDeficientLocked()
}
func (f *Forwarder) UpTrackLayersChange(availableLayers []int32, exemptedLayers []int32) {
@@ -446,7 +516,7 @@ func (f *Forwarder) UpTrackLayersChange(availableLayers []int32, exemptedLayers
}
func (f *Forwarder) getOptimalBandwidthNeeded(brs Bitrates, maxLayers VideoLayers) int64 {
if f.muted {
if f.muted || f.pubMuted {
return 0
}
@@ -502,7 +572,7 @@ func (f *Forwarder) bitrateAvailable(brs Bitrates) bool {
}
func (f *Forwarder) getDistanceToDesired(brs Bitrates, targetLayers VideoLayers, maxLayers VideoLayers) int32 {
if f.muted {
if f.muted || f.pubMuted {
return 0
}
@@ -542,11 +612,15 @@ func (f *Forwarder) getDistanceToDesired(brs Bitrates, targetLayers VideoLayers,
return distance
}
func (f *Forwarder) isDeficientLocked() bool {
return f.lastAllocation.state == VideoAllocationStateDeficient
}
func (f *Forwarder) IsDeficient() bool {
f.lock.RLock()
defer f.lock.RUnlock()
return f.lastAllocation.state == VideoAllocationStateDeficient
return f.isDeficientLocked()
}
func (f *Forwarder) BandwidthRequested(brs Bitrates) int64 {
@@ -586,21 +660,33 @@ func (f *Forwarder) AllocateOptimal(brs Bitrates, allowOvershoot bool) VideoAllo
case f.muted:
alloc.state = VideoAllocationStateMuted
case len(f.availableLayers) == 0:
// feed is dry
alloc.state = VideoAllocationStateFeedDry
case f.pubMuted:
alloc.state = VideoAllocationStatePubMuted
// leave it at current layers for opportunistic resume
alloc.targetLayers = f.currentLayers
case f.parkedLayers.IsValid():
// if parked on a layer, let it continue
alloc.state = VideoAllocationStateOptimal
alloc.targetLayers = f.parkedLayers
case !f.bitrateAvailable(brs):
// feed bitrate not yet calculated for all available layers
alloc.state = VideoAllocationStateAwaitingMeasurement
//
// Resume with the highest layer available <= max subscribed layer
// If already resumed, move allocation to the highest available layer <= max subscribed layer
// Target the highest layer available <= max subscribed layer.
//
alloc.targetLayers = VideoLayers{
Spatial: int32(math.Min(float64(f.maxLayers.Spatial), float64(f.availableLayers[len(f.availableLayers)-1]))),
Temporal: int32(math.Max(0, float64(f.maxLayers.Temporal))),
// If target already set, move allocation to the highest available layer <= max subscribed layer only if that is higher than existing target.
// It is possible that in situations like coming out of publisher mute, target is already set higher due to parked layers.
//
if f.availableLayers[len(f.availableLayers)-1] > f.targetLayers.Spatial {
alloc.targetLayers = VideoLayers{
Spatial: int32(math.Min(float64(f.maxLayers.Spatial), float64(f.availableLayers[len(f.availableLayers)-1]))),
Temporal: int32(math.Max(0, float64(f.maxLayers.Temporal))),
}
} else {
alloc.targetLayers = f.targetLayers
}
default:
@@ -656,28 +742,19 @@ func (f *Forwarder) AllocateOptimal(brs Bitrates, allowOvershoot bool) VideoAllo
}
}
// feed may be dry, leave target at current if already started for opportunistic resume
if alloc.bandwidthRequested == 0 && f.maxLayers.IsValid() {
// if overshoot was allowed and it did not also find a layer,
// keep target at exempted layer (if available) and the current layer is at that level.
// i. e. exempted layer may really have stopped, so a layer switch to an exempted layer should
// not happen as layer switch will send PLI requests. Just letting it continue at the current
// layer if the current is exempted will protect against any stream tracker misdetects
// OR latch on to the layer quicker when it restarts
if f.currentLayers.IsValid() {
for _, s := range f.exemptedLayers {
if s <= f.maxLayers.Spatial && f.currentLayers.Spatial == s {
alloc.targetLayers = f.currentLayers
alloc.bandwidthRequested = brs[alloc.targetLayers.Spatial][alloc.targetLayers.Temporal]
alloc.state = VideoAllocationStateDeficient
break
}
alloc.state = VideoAllocationStateFeedDry
if f.started {
alloc.targetLayers = f.currentLayers
} else {
// opportunisitically latch on to anything
alloc.targetLayers = VideoLayers{
Spatial: DefaultMaxLayerSpatial,
Temporal: DefaultMaxLayerTemporal,
}
}
}
if !alloc.targetLayers.IsValid() && f.maxLayers.IsValid() {
alloc.state = VideoAllocationStateDeficient
}
}
if !alloc.targetLayers.IsValid() {
@@ -696,8 +773,11 @@ func (f *Forwarder) ProvisionalAllocatePrepare(bitrates Bitrates) {
f.provisional = &VideoAllocationProvisional{
allocatedLayers: InvalidLayers,
muted: f.muted,
pubMuted: f.pubMuted,
bitrates: bitrates,
maxLayers: f.maxLayers,
currentLayers: f.currentLayers,
parkedLayers: f.parkedLayers,
}
if len(f.availableLayers) > 0 {
f.provisional.availableLayers = make([]int32, len(f.availableLayers))
@@ -713,31 +793,17 @@ func (f *Forwarder) ProvisionalAllocate(availableChannelCapacity int64, layers V
f.lock.Lock()
defer f.lock.Unlock()
if f.provisional.muted || !f.provisional.maxLayers.IsValid() || (!allowOvershoot && layers.GreaterThan(f.provisional.maxLayers)) {
if f.provisional.muted || f.provisional.pubMuted || !f.provisional.maxLayers.IsValid() || (!allowOvershoot && layers.GreaterThan(f.provisional.maxLayers)) {
return 0
}
maybeAdoptExempted := func() int64 {
br := int64(0)
if f.currentLayers.IsValid() {
for _, s := range f.provisional.exemptedLayers {
if s <= f.provisional.maxLayers.Spatial && f.currentLayers.Spatial == s {
f.provisional.allocatedLayers = f.currentLayers
br = f.provisional.bitrates[f.provisional.allocatedLayers.Spatial][f.provisional.allocatedLayers.Temporal]
break
}
}
}
return br
}
requiredBitrate := f.provisional.bitrates[layers.Spatial][layers.Temporal]
if requiredBitrate == 0 {
return maybeAdoptExempted()
return 0
}
alreadyAllocatedBitrate := int64(0)
if f.provisional.allocatedLayers != InvalidLayers {
if f.provisional.allocatedLayers.IsValid() {
alreadyAllocatedBitrate = f.provisional.bitrates[f.provisional.allocatedLayers.Spatial][f.provisional.allocatedLayers.Temporal]
}
@@ -759,7 +825,7 @@ func (f *Forwarder) ProvisionalAllocate(availableChannelCapacity int64, layers V
return requiredBitrate - alreadyAllocatedBitrate
}
return maybeAdoptExempted()
return 0
}
func (f *Forwarder) ProvisionalAllocateGetCooperativeTransition(allowOvershoot bool) VideoTransition {
@@ -785,18 +851,21 @@ func (f *Forwarder) ProvisionalAllocateGetCooperativeTransition(allowOvershoot b
f.lock.Lock()
defer f.lock.Unlock()
if f.provisional.muted {
if f.provisional.muted || f.provisional.pubMuted {
f.provisional.allocatedLayers = InvalidLayers
if f.provisional.pubMuted {
// leave it at current for opportunistic forwarding, there is still bandwidth saving with publisher mute
f.provisional.allocatedLayers = f.provisional.currentLayers
}
return VideoTransition{
from: f.targetLayers,
to: InvalidLayers,
to: f.provisional.allocatedLayers,
bandwidthDelta: 0 - f.lastAllocation.bandwidthRequested,
// LK-TODO should this take current bitrate of current target layers?
}
}
// check if we should preserve current target
if f.targetLayers != InvalidLayers {
if f.targetLayers.IsValid() {
// what is the highest that is available
maximalLayers := InvalidLayers
maximalBandwidthRequired := int64(0)
@@ -814,9 +883,10 @@ func (f *Forwarder) ProvisionalAllocateGetCooperativeTransition(allowOvershoot b
}
}
if maximalLayers != InvalidLayers {
if maximalLayers.IsValid() {
if !f.targetLayers.GreaterThan(maximalLayers) && f.provisional.bitrates[f.targetLayers.Spatial][f.targetLayers.Temporal] != 0 {
// currently streaming and wanting an upgrade, just preserve current target in the cooperative scheme of things
// currently streaming and maybe wanting an upgrade (f.targetLayers <= maximalLayers),
// just preserve current target in the cooperative scheme of things
f.provisional.allocatedLayers = f.targetLayers
return VideoTransition{
from: f.targetLayers,
@@ -826,7 +896,7 @@ func (f *Forwarder) ProvisionalAllocateGetCooperativeTransition(allowOvershoot b
}
if f.targetLayers.GreaterThan(maximalLayers) {
// maximalLayers <= f.targetLayers, make the down move
// maximalLayers < f.targetLayers, make the down move
f.provisional.allocatedLayers = maximalLayers
return VideoTransition{
from: f.targetLayers,
@@ -870,30 +940,23 @@ func (f *Forwarder) ProvisionalAllocateGetCooperativeTransition(allowOvershoot b
0, f.provisional.maxLayers.Spatial,
0, f.provisional.maxLayers.Temporal,
)
}
// could not find a minimal layer, overshoot if allowed
if bandwidthRequired == 0 && f.provisional.maxLayers.IsValid() && allowOvershoot {
targetLayers, bandwidthRequired = findNextLayer(
f.provisional.maxLayers.Spatial+1, DefaultMaxLayerSpatial,
0, DefaultMaxLayerTemporal,
)
}
// adopt exempted layer if current is at one of the exempted layers below maximum
if bandwidthRequired == 0 && f.provisional.maxLayers.IsValid() && f.currentLayers.IsValid() {
for _, s := range f.provisional.exemptedLayers {
if s <= f.provisional.maxLayers.Spatial && f.currentLayers.Spatial == s {
targetLayers = f.currentLayers
bandwidthRequired = f.provisional.bitrates[targetLayers.Spatial][targetLayers.Temporal]
break
}
// could not find a minimal layer, overshoot if allowed
if bandwidthRequired == 0 && f.provisional.maxLayers.IsValid() && allowOvershoot {
targetLayers, bandwidthRequired = findNextLayer(
f.provisional.maxLayers.Spatial+1, DefaultMaxLayerSpatial,
0, DefaultMaxLayerTemporal,
)
}
}
// turn off if nothing found, not even an exempted layer to continue with
if bandwidthRequired == 0 && (!f.currentLayers.IsValid() || f.currentLayers != targetLayers) {
targetLayers = InvalidLayers
// if nothing available, just leave target at current to enable opportunistic forwarding in case current resumes
if targetLayers == InvalidLayers {
if f.provisional.parkedLayers.IsValid() {
targetLayers = f.provisional.parkedLayers
} else {
targetLayers = f.provisional.currentLayers
}
}
f.provisional.allocatedLayers = targetLayers
@@ -907,7 +970,8 @@ func (f *Forwarder) ProvisionalAllocateGetCooperativeTransition(allowOvershoot b
func (f *Forwarder) ProvisionalAllocateGetBestWeightedTransition() VideoTransition {
//
// This is called when a track needs a change (could be mute/unmute, subscribed layers changed, published layers changed)
// when channel is congested.
// when channel is congested. This is called on tracks other than the one needing the change. When the track
// needing the change requires bits, this is called to check if this track can contribute some bits to the pool.
//
// The goal is to keep all tracks streaming as much as possible. So, the track that needs a change needs bandwidth to be unpaused.
//
@@ -922,13 +986,16 @@ func (f *Forwarder) ProvisionalAllocateGetBestWeightedTransition() VideoTransiti
f.lock.Lock()
defer f.lock.Unlock()
if f.provisional.muted {
if f.provisional.muted || f.provisional.pubMuted {
f.provisional.allocatedLayers = InvalidLayers
if f.provisional.pubMuted {
// leave it at current for opportunistic forwarding, there is still bandwidth saving with publisher mute
f.provisional.allocatedLayers = f.provisional.currentLayers
}
return VideoTransition{
from: f.targetLayers,
to: InvalidLayers,
to: f.provisional.allocatedLayers,
bandwidthDelta: 0 - f.lastAllocation.bandwidthRequested,
// LK-TODO should this take current bitrate of current target layers?
}
}
@@ -946,28 +1013,18 @@ func (f *Forwarder) ProvisionalAllocateGetBestWeightedTransition() VideoTransiti
}
if maxReachableLayerTemporal == InvalidLayerTemporal {
// stick to an exempted layer if available
if f.currentLayers.IsValid() {
for _, s := range f.provisional.exemptedLayers {
if s <= f.provisional.maxLayers.Spatial && f.currentLayers.Spatial == s {
f.provisional.allocatedLayers = f.currentLayers
return VideoTransition{
from: f.targetLayers,
to: f.provisional.allocatedLayers,
bandwidthDelta: 0 - f.lastAllocation.bandwidthRequested,
// LK-TODO should this take current bitrate of current target layers?
}
}
}
// feed has gone dry, just leave target at current to enable opportunistic forwarding in case current resumes.
// Note that this is giving back bits and opportunistic forwarding resuming might trigger congestion again,
// but that should be handled by stream allocator.
if f.provisional.parkedLayers.IsValid() {
f.provisional.allocatedLayers = f.provisional.parkedLayers
} else {
f.provisional.allocatedLayers = f.provisional.currentLayers
}
// feed has gone dry,
f.provisional.allocatedLayers = InvalidLayers
return VideoTransition{
from: f.targetLayers,
to: InvalidLayers,
to: f.provisional.allocatedLayers,
bandwidthDelta: 0 - f.lastAllocation.bandwidthRequested,
// LK-TODO should this take current bitrate of current target layers?
}
}
@@ -1029,13 +1086,26 @@ func (f *Forwarder) ProvisionalAllocateCommit() VideoAllocation {
case f.provisional.muted:
alloc.state = VideoAllocationStateMuted
case f.provisional.pubMuted:
alloc.state = VideoAllocationStatePubMuted
case len(f.provisional.availableLayers) == 0:
// feed is dry
alloc.state = VideoAllocationStateFeedDry
case f.provisional.allocatedLayers == InvalidLayers:
case !f.provisional.allocatedLayers.IsValid():
alloc.state = VideoAllocationStateDeficient
// leave target at current if current layer is exempted for opportunistic forwarding
if f.provisional.currentLayers.IsValid() {
for _, s := range f.provisional.exemptedLayers {
if s <= f.provisional.maxLayers.Spatial && f.provisional.currentLayers.Spatial == s {
f.provisional.allocatedLayers = f.provisional.currentLayers
alloc.targetLayers = f.provisional.allocatedLayers
break
}
}
}
default:
optimalBandwidthNeeded := f.getOptimalBandwidthNeeded(f.provisional.bitrates, f.provisional.maxLayers)
bandwidthRequested := f.provisional.bitrates[f.provisional.allocatedLayers.Spatial][f.provisional.allocatedLayers.Temporal]
@@ -1055,6 +1125,7 @@ func (f *Forwarder) ProvisionalAllocateCommit() VideoAllocation {
}
}
f.clearParkedLayers()
return f.updateAllocation(alloc, "cooperative")
}
@@ -1074,7 +1145,7 @@ func (f *Forwarder) AllocateNextHigher(availableChannelCapacity int64, brs Bitra
}
// if targets are still pending, don't increase
if f.targetLayers != InvalidLayers && f.targetLayers != f.currentLayers {
if f.targetLayers.IsValid() && f.targetLayers != f.currentLayers {
f.lastAllocation.change = VideoStreamingChangeNone
return f.lastAllocation, false
}
@@ -1082,7 +1153,7 @@ func (f *Forwarder) AllocateNextHigher(availableChannelCapacity int64, brs Bitra
optimalBandwidthNeeded := f.getOptimalBandwidthNeeded(brs, f.maxLayers)
alreadyAllocated := int64(0)
if f.targetLayers != InvalidLayers {
if f.targetLayers.IsValid() {
alreadyAllocated = brs[f.targetLayers.Spatial][f.targetLayers.Temporal]
}
@@ -1130,7 +1201,7 @@ func (f *Forwarder) AllocateNextHigher(availableChannelCapacity int64, brs Bitra
boosted := false
// try moving temporal layer up in currently streaming spatial layer
if f.targetLayers != InvalidLayers {
if f.targetLayers.IsValid() {
done, allocation, boosted = doAllocation(
f.targetLayers.Spatial, f.targetLayers.Spatial,
f.targetLayers.Temporal+1, f.maxLayers.Temporal,
@@ -1177,12 +1248,12 @@ func (f *Forwarder) GetNextHigherTransition(brs Bitrates, allowOvershoot bool) (
}
// if targets are still pending, don't increase
if f.targetLayers != InvalidLayers && f.targetLayers != f.currentLayers {
if f.targetLayers.IsValid() && f.targetLayers != f.currentLayers {
return VideoTransition{}, false
}
alreadyAllocated := int64(0)
if f.targetLayers != InvalidLayers {
if f.targetLayers.IsValid() {
alreadyAllocated = brs[f.targetLayers.Spatial][f.targetLayers.Temporal]
}
@@ -1215,7 +1286,7 @@ func (f *Forwarder) GetNextHigherTransition(brs Bitrates, allowOvershoot bool) (
isAvailable := false
// try moving temporal layer up in currently streaming spatial layer
if f.targetLayers != InvalidLayers {
if f.targetLayers.IsValid() {
done, transition, isAvailable = findNextHigher(
f.targetLayers.Spatial, f.targetLayers.Spatial,
f.targetLayers.Temporal+1, f.maxLayers.Temporal,
@@ -1265,6 +1336,9 @@ func (f *Forwarder) Pause(brs Bitrates) VideoAllocation {
case f.muted:
alloc.state = VideoAllocationStateMuted
case f.pubMuted:
alloc.state = VideoAllocationStatePubMuted
case len(f.availableLayers) == 0:
// feed is dry
alloc.state = VideoAllocationStateFeedDry
@@ -1274,13 +1348,14 @@ func (f *Forwarder) Pause(brs Bitrates) VideoAllocation {
alloc.state = VideoAllocationStateDeficient
}
f.clearParkedLayers()
return f.updateAllocation(alloc, "pause")
}
func (f *Forwarder) updateAllocation(alloc VideoAllocation, reason string) VideoAllocation {
if f.targetLayers == InvalidLayers && alloc.targetLayers.IsValid() {
if !f.targetLayers.IsValid() && alloc.targetLayers.IsValid() {
alloc.change = VideoStreamingChangeResuming
} else if f.targetLayers != InvalidLayers && !alloc.targetLayers.IsValid() {
} else if f.targetLayers.IsValid() && !alloc.targetLayers.IsValid() {
alloc.change = VideoStreamingChangePausing
}
@@ -1323,6 +1398,30 @@ func (f *Forwarder) Resync() {
func (f *Forwarder) resyncLocked() {
f.currentLayers = InvalidLayers
f.lastSSRC = 0
f.clearParkedLayers()
}
func (f *Forwarder) clearParkedLayers() {
f.parkedLayers = InvalidLayers
if f.parkedLayersTimer != nil {
f.parkedLayersTimer.Stop()
f.parkedLayersTimer = nil
}
}
func (f *Forwarder) setupParkedLayers(parkedLayers VideoLayers) {
f.clearParkedLayers()
f.parkedLayers = parkedLayers
f.parkedLayersTimer = time.AfterFunc(ParkedLayersWaitDuration, func() {
f.lock.Lock()
f.clearParkedLayers()
f.lock.Unlock()
if onParkedLayersExpired := f.getOnParkedLayersExpired(); onParkedLayersExpired != nil {
onParkedLayersExpired()
}
})
}
func (f *Forwarder) CheckSync() (locked bool, layer int32) {
@@ -1330,8 +1429,7 @@ func (f *Forwarder) CheckSync() (locked bool, layer int32) {
defer f.lock.RUnlock()
layer = f.targetLayers.Spatial
locked = f.targetLayers.Spatial == f.currentLayers.Spatial
locked = f.targetLayers.Spatial == f.currentLayers.Spatial || f.parkedLayers.IsValid()
return
}
@@ -1355,8 +1453,7 @@ func (f *Forwarder) FilterRTX(nacks []uint16) (filtered []uint16, disallowedLaye
// Without the curb, when congestion hits, RTX rate could be so high that it further congests the channel.
//
for layer := int32(0); layer < DefaultMaxLayerSpatial+1; layer++ {
if f.lastAllocation.state == VideoAllocationStateDeficient &&
(f.targetLayers.Spatial < f.currentLayers.Spatial || layer > f.currentLayers.Spatial) {
if f.isDeficientLocked() && (f.targetLayers.Spatial < f.currentLayers.Spatial || layer > f.currentLayers.Spatial) {
disallowedLayers[layer] = true
}
}
@@ -1368,6 +1465,7 @@ func (f *Forwarder) GetTranslationParams(extPkt *buffer.ExtPacket, layer int32)
f.lock.Lock()
defer f.lock.Unlock()
// Video: Do not drop on publisher mute to enable resume on publisher unmute without a key frame.
if f.muted {
return &TranslationParams{
shouldDrop: true,
@@ -1376,6 +1474,13 @@ func (f *Forwarder) GetTranslationParams(extPkt *buffer.ExtPacket, layer int32)
switch f.kind {
case webrtc.RTPCodecTypeAudio:
// Audio: Blank frames are injected on publisher mute to ensure decoder does not get stuck at a noise frame. So, do not forward.
if f.pubMuted {
return &TranslationParams{
shouldDrop: true,
}, nil
}
return f.getTranslationParamsAudio(extPkt, layer)
case webrtc.RTPCodecTypeVideo:
return f.getTranslationParamsVideo(extPkt, layer)
@@ -1451,7 +1556,7 @@ func (f *Forwarder) getTranslationParamsAudio(extPkt *buffer.ExtPacket, layer in
func (f *Forwarder) getTranslationParamsVideo(extPkt *buffer.ExtPacket, layer int32) (*TranslationParams, error) {
tp := &TranslationParams{}
if f.targetLayers == InvalidLayers {
if !f.targetLayers.IsValid() {
// stream is paused by streamallocator
tp.shouldDrop = true
return tp, nil
@@ -1462,37 +1567,90 @@ func (f *Forwarder) getTranslationParamsVideo(extPkt *buffer.ExtPacket, layer in
tp.shouldDrop = true
f.rtpMunger.PacketDropped(extPkt)
return tp, nil
} else if tp.switchingToTargetLayer {
// lock to target layer
f.logger.Infow("locking to target layer", "current", f.currentLayers, "target", f.targetLayers)
f.currentLayers.Spatial = f.targetLayers.Spatial
if !f.isTemporalSupported {
f.currentLayers.Temporal = f.targetLayers.Temporal
}
// TODO : we switch to target layer immediately now since we assume all frame chain is integrity
// if we have frame chain check, should switch only if target chain is not broken and decodable
// if f.ddLayerSelector != nil {
// f.ddLayerSelector.SelectLayer(f.currentLayers)
// }
if f.currentLayers.Spatial >= f.maxLayers.Spatial || f.currentLayers.Spatial == (f.numAdvertisedLayers-1) {
tp.isSwitchingToMaxLayer = true
}
}
}
if f.targetLayers.Spatial != f.currentLayers.Spatial {
if f.targetLayers.Spatial == layer {
if extPkt.KeyFrame || tp.switchingToTargetLayer {
// lock to target layer
f.logger.Infow("locking to target layer", "current", f.currentLayers, "target", f.targetLayers)
f.currentLayers.Spatial = f.targetLayers.Spatial
if !f.isTemporalSupported {
f.currentLayers.Temporal = f.targetLayers.Temporal
} else {
if f.currentLayers.Spatial != f.targetLayers.Spatial {
// Three things to check when not locked to target
// 1. Resumable layer - don't need a key frame
// 2. Opportunistic layer upgrade - needs a key frame
// 3. Need to downgrade - needs a key frame
found := false
if f.parkedLayers.IsValid() {
if f.parkedLayers.Spatial == layer {
f.logger.Infow("resuming at parked layer", "current", f.currentLayers, "target", f.targetLayers, "parked", f.parkedLayers)
f.currentLayers = f.parkedLayers
found = true
}
// TODO : we switch to target layer immediately now since we assume all frame chain is integrity
// if we have frame chain check, should switch only if target chain is not broken and decodable
// if f.ddLayerSelector != nil {
// f.ddLayerSelector.SelectLayer(f.currentLayers)
// }
if f.currentLayers.Spatial >= f.maxLayers.Spatial {
} else {
if extPkt.KeyFrame {
if layer > f.currentLayers.Spatial && layer <= f.targetLayers.Spatial {
f.logger.Infow("upgrading layer", "current", f.currentLayers, "target", f.targetLayers, "upgrade", layer)
found = true
}
if layer < f.currentLayers.Spatial && layer >= f.targetLayers.Spatial {
f.logger.Infow("downgrading layer", "current", f.currentLayers, "target", f.targetLayers, "downgrade", layer)
found = true
}
if found {
f.currentLayers.Spatial = layer
if !f.isTemporalSupported {
f.currentLayers.Temporal = extPkt.Temporal
}
}
}
}
if found {
f.clearParkedLayers()
if f.currentLayers.Spatial >= f.maxLayers.Spatial || f.currentLayers.Spatial == (f.numAdvertisedLayers-1) {
tp.isSwitchingToMaxLayer = true
// if maximum is attained, adjust target to enable fast path layer check in per-packet path
f.logger.Infow(
"reached max layer",
"current", f.currentLayers,
"target", f.targetLayers,
"max", f.maxLayers,
"advertised", f.numAdvertisedLayers,
)
f.targetLayers.Spatial = f.currentLayers.Spatial
}
}
} else {
// if locked to higher than nax layer due to overshoot, check if it can be dialed back
if f.targetLayers.Spatial > f.maxLayers.Spatial {
if layer <= f.maxLayers.Spatial && extPkt.KeyFrame {
f.logger.Infow("adjusting overshoot", "current", f.currentLayers, "target", f.targetLayers, "adjuted", layer)
f.currentLayers.Spatial = layer
f.targetLayers.Spatial = layer
}
}
}
}
// if we have layer selector, let it decide whether to drop or not
if f.ddLayerSelector == nil && f.currentLayers.Spatial != layer {
if f.currentLayers.Spatial != layer {
tp.shouldDrop = true
return tp, nil
}
if FlagPauseOnDowngrade && f.targetLayers.Spatial < f.currentLayers.Spatial && f.lastAllocation.state == VideoAllocationStateDeficient {
if FlagPauseOnDowngrade && f.targetLayers.Spatial < f.currentLayers.Spatial && f.isDeficientLocked() {
//
// If target layer is lower than both the current and
// maximum subscribed layer, it is due to bandwidth
@@ -1504,12 +1662,12 @@ func (f *Forwarder) getTranslationParamsVideo(extPkt *buffer.ExtPacket, layer in
// switch point to get a smoother stream till the higher
// layer key frame arrives.
//
// Note that in the case of client subscription layer restriction
// coinciding with server restriction due to bandwidth limitation,
// Note that it is possible for client subscription layer restriction
// to coincide with server restriction due to bandwidth limitation,
// In the case of subscription change, higher should continue streaming
// to ensure smooth transition.
//
// To differentiate, drop only when in DEFICIENT state.
// To differentiate between the two cases, drop only when in DEFICIENT state.
//
tp.shouldDrop = true
tp.isDroppingRelevant = true
@@ -1559,12 +1717,12 @@ func (f *Forwarder) GetSnTsForPadding(num int) ([]SnTs, error) {
defer f.lock.Unlock()
// padding is used for probing. Padding packets should be
// at frame boundaries only to ensure decoder sequencer does
// at only the frame boundaries to ensure decoder sequencer does
// not get out-of-sync. But, when a stream is paused,
// force a frame marker as a restart of the stream will
// start with a key frame which will reset the decoder.
forceMarker := false
if f.targetLayers == InvalidLayers {
if !f.targetLayers.IsValid() {
forceMarker = true
}
return f.rtpMunger.UpdateAndGetPaddingSnTs(num, 0, 0, forceMarker)
+83 -25
View File
@@ -194,29 +194,55 @@ func TestForwarderAllocateOptimal(t *testing.T) {
require.Equal(t, expectedResult, result)
require.Equal(t, expectedResult, f.lastAllocation)
// feed dry state
f.Mute(false)
// feed dry state, target set to maximum for opportunistic start
f.lastAllocation.state = VideoAllocationStateNone
disable(f)
expectedTargetLayers := VideoLayers{
Spatial: DefaultMaxLayerSpatial,
Temporal: DefaultMaxLayerTemporal,
}
expectedResult = VideoAllocation{
state: VideoAllocationStateFeedDry,
change: VideoStreamingChangeNone,
change: VideoStreamingChangeResuming,
bandwidthRequested: 0,
bandwidthDelta: 0,
availableLayers: nil,
bitrates: emptyBitrates,
targetLayers: InvalidLayers,
targetLayers: expectedTargetLayers,
distanceToDesired: 0,
}
result = f.AllocateOptimal(emptyBitrates, true)
require.Equal(t, expectedResult, result)
require.Equal(t, expectedResult, f.lastAllocation)
require.Equal(t, expectedTargetLayers, f.TargetLayers())
f.parkedLayers = VideoLayers{
Spatial: 0,
Temporal: 1,
}
expectedResult = VideoAllocation{
state: VideoAllocationStateOptimal,
change: VideoStreamingChangeNone,
bandwidthRequested: 0,
bandwidthDelta: 0,
availableLayers: nil,
bitrates: emptyBitrates,
targetLayers: f.parkedLayers,
distanceToDesired: 0,
}
result = f.AllocateOptimal(emptyBitrates, true)
require.Equal(t, expectedResult, result)
require.Equal(t, expectedResult, f.lastAllocation)
require.Equal(t, f.parkedLayers, f.TargetLayers())
f.parkedLayers = InvalidLayers
// awaiting measurement, i.e. bitrates are not available, but layers available
f.lastAllocation.state = VideoAllocationStateNone
disable(f)
f.UpTrackLayersChange([]int32{0}, []int32{})
expectedTargetLayers := VideoLayers{
expectedTargetLayers = VideoLayers{
Spatial: 0,
Temporal: DefaultMaxLayerTemporal,
}
@@ -236,16 +262,17 @@ func TestForwarderAllocateOptimal(t *testing.T) {
require.Equal(t, expectedTargetLayers, f.TargetLayers())
require.Equal(t, InvalidLayers, f.CurrentLayers())
// layers are available, but all layers under max are exempted
// layers are available, but all layers <= max are exempted.
// forwarder not started yet, so should set target to max.
f.SetMaxSpatialLayer(0)
f.currentLayers = VideoLayers{Spatial: 0, Temporal: 0}
f.UpTrackLayersChange([]int32{0, 1, 2}, []int32{0})
expectedTargetLayers = VideoLayers{
Spatial: 0,
Temporal: 0,
Spatial: DefaultMaxLayerSpatial,
Temporal: DefaultMaxLayerTemporal,
}
expectedResult = VideoAllocation{
state: VideoAllocationStateDeficient,
state: VideoAllocationStateFeedDry,
change: VideoStreamingChangeNone,
bandwidthRequested: 0,
bandwidthDelta: 0,
@@ -260,23 +287,49 @@ func TestForwarderAllocateOptimal(t *testing.T) {
require.Equal(t, expectedResult, f.lastAllocation)
require.Equal(t, expectedTargetLayers, f.TargetLayers())
// if current is not exempt, should not choose that
f.currentLayers = VideoLayers{Spatial: 1, Temporal: 0}
// when started, should just leave it at current for opportunistic resume
f.started = true
expectedTargetLayers = VideoLayers{
Spatial: 0,
Temporal: 0,
}
expectedResult = VideoAllocation{
state: VideoAllocationStateDeficient,
change: VideoStreamingChangePausing,
state: VideoAllocationStateFeedDry,
change: VideoStreamingChangeNone,
bandwidthRequested: 0,
bandwidthDelta: 0,
availableLayers: []int32{0, 1, 2},
exemptedLayers: []int32{0},
bitrates: emptyBitrates,
targetLayers: InvalidLayers,
targetLayers: expectedTargetLayers,
distanceToDesired: 0,
}
result = f.AllocateOptimal(emptyBitrates, true)
require.Equal(t, expectedResult, result)
require.Equal(t, expectedResult, f.lastAllocation)
require.Equal(t, InvalidLayers, f.TargetLayers())
require.Equal(t, expectedTargetLayers, f.TargetLayers())
// if current is not exempt, should still leave at current for opportunistic resume
f.currentLayers = VideoLayers{Spatial: 1, Temporal: 0}
expectedTargetLayers = VideoLayers{
Spatial: 1,
Temporal: 0,
}
expectedResult = VideoAllocation{
state: VideoAllocationStateFeedDry,
change: VideoStreamingChangeNone,
bandwidthRequested: 0,
bandwidthDelta: 0,
availableLayers: []int32{0, 1, 2},
exemptedLayers: []int32{0},
bitrates: emptyBitrates,
targetLayers: expectedTargetLayers,
distanceToDesired: 0,
}
result = f.AllocateOptimal(emptyBitrates, true)
require.Equal(t, expectedResult, result)
require.Equal(t, expectedResult, f.lastAllocation)
require.Equal(t, expectedTargetLayers, f.TargetLayers())
// allocate using bitrates, allocation should choose optimal
f.currentLayers = InvalidLayers
@@ -289,7 +342,7 @@ func TestForwarderAllocateOptimal(t *testing.T) {
}
expectedResult = VideoAllocation{
state: VideoAllocationStateOptimal,
change: VideoStreamingChangeResuming,
change: VideoStreamingChangeNone,
bandwidthRequested: bitrates[1][3],
bandwidthDelta: bitrates[1][3],
availableLayers: []int32{0, 1, 2},
@@ -332,7 +385,7 @@ func TestForwarderAllocateOptimal(t *testing.T) {
// when not allowing overshoot, should not be able to find a layer
expectedResult = VideoAllocation{
state: VideoAllocationStateDeficient,
state: VideoAllocationStateFeedDry,
change: VideoStreamingChangePausing,
bandwidthRequested: 0,
bandwidthDelta: -sparseBitrates[2][1],
@@ -477,7 +530,7 @@ func TestForwarderProvisionalAllocate(t *testing.T) {
//
// Test exemptedLayers
// Even if overshoot is allowed, but higher layers do not have bit rates and the current layer is empted,,
// Even if overshoot is allowed, but if higher layers do not have bit rates and the current layer is exempted,
// should continue with current layer.
//
f.SetMaxSpatialLayer(0)
@@ -797,37 +850,39 @@ func TestForwarderProvisionalAllocateGetCooperativeTransition(t *testing.T) {
require.Equal(t, expectedResult, f.lastAllocation)
require.Equal(t, expectedLayers, f.TargetLayers())
// Test exemptedLayers not matching current layers
// Test exemptedLayers not matching current layers,
// leave target at current if no transition available to facilitate opportunistic forwarding
f.currentLayers = VideoLayers{Spatial: 2, Temporal: 2}
f.targetLayers = f.currentLayers
f.availableLayers = availableLayers
f.exemptedLayers = exemptedLayers
f.ProvisionalAllocatePrepare(bitrates)
expectedLayers = VideoLayers{Spatial: 2, Temporal: 2}
expectedTransition = VideoTransition{
from: VideoLayers{Spatial: 2, Temporal: 2},
to: InvalidLayers,
to: expectedLayers,
bandwidthDelta: 0,
}
transition = f.ProvisionalAllocateGetCooperativeTransition(true)
require.Equal(t, expectedTransition, transition)
// committing should set target to InvalidLayers
// committing should set target to current layers to enable opportunistic forwarding
expectedResult = VideoAllocation{
state: VideoAllocationStateDeficient,
change: VideoStreamingChangePausing,
state: VideoAllocationStateOptimal,
change: VideoStreamingChangeNone,
bandwidthRequested: 0,
bandwidthDelta: 0,
availableLayers: availableLayers,
exemptedLayers: exemptedLayers,
bitrates: bitrates,
targetLayers: InvalidLayers,
targetLayers: expectedLayers,
distanceToDesired: 0,
}
result = f.ProvisionalAllocateCommit()
require.Equal(t, expectedResult, result)
require.Equal(t, expectedResult, f.lastAllocation)
require.Equal(t, InvalidLayers, f.TargetLayers())
require.Equal(t, expectedLayers, f.TargetLayers())
}
func TestForwarderProvisionalAllocateGetBestWeightedTransition(t *testing.T) {
@@ -882,7 +937,10 @@ func TestForwarderAllocateNextHigher(t *testing.T) {
// if layers have not caught up, should not allocate next layer
f.lastAllocation.state = VideoAllocationStateDeficient
f.targetLayers.Spatial = 0
f.targetLayers = VideoLayers{
Spatial: 0,
Temporal: 0,
}
expectedResult := VideoAllocation{
state: VideoAllocationStateDeficient,
change: VideoStreamingChangeNone,
+3
View File
@@ -197,6 +197,7 @@ func NewWebRTCReceiver(
w.streamTrackerManager = NewStreamTrackerManager(logger, trackInfo, w.isSVC, w.codec.ClockRate, trackersConfig)
w.streamTrackerManager.OnAvailableLayersChanged(w.downTrackLayerChange)
w.streamTrackerManager.OnBitrateAvailabilityChanged(w.downTrackBitrateAvailabilityChange)
w.streamTrackerManager.Start()
for _, opt := range opts {
w = opt(w)
@@ -562,6 +563,8 @@ func (w *WebRTCReceiver) forwardRTP(layer int32) {
if pr := w.redReceiver.Load(); pr != nil {
pr.(*RedReceiver).Close()
}
w.streamTrackerManager.Stop()
})
w.streamTrackerManager.RemoveTracker(layer)
+4 -4
View File
@@ -204,7 +204,7 @@ func (r *RTPMunger) UpdateAndGetSnTs(extPkt *buffer.ExtPacket) (*TranslationPara
r.isInRtxGateRegion = true
}
if r.isInRtxGateRegion && (mungedSN-r.rtxGateSn) > RtxGateWindow {
if r.isInRtxGateRegion && (mungedSN-r.rtxGateSn) < (1<<15) && (mungedSN-r.rtxGateSn) > RtxGateWindow {
r.isInRtxGateRegion = false
}
@@ -235,10 +235,10 @@ func (r *RTPMunger) UpdateAndGetPaddingSnTs(num int, clockRate uint32, frameRate
if !r.lastMarker {
if !forceMarker {
return nil, ErrPaddingNotOnFrameBoundary
} else {
// if forcing frame end, use timestamp of latest received frame for the first one
tsOffset = 1
}
// if forcing frame end, use timestamp of latest received frame for the first one
tsOffset = 1
}
vals := make([]SnTs, num)
+27 -13
View File
@@ -7,6 +7,7 @@ import (
"github.com/livekit/livekit-server/pkg/config"
"github.com/livekit/livekit-server/pkg/sfu/buffer"
"github.com/livekit/livekit-server/pkg/sfu/streamtracker"
"github.com/livekit/livekit-server/pkg/utils"
"github.com/livekit/protocol/livekit"
"github.com/livekit/protocol/logger"
)
@@ -20,6 +21,8 @@ type StreamTrackerManager struct {
trackerConfig config.StreamTrackerConfig
layersNotifyOpQueue *utils.OpsQueue
lock sync.RWMutex
trackers [DefaultMaxLayerSpatial + 1]*streamtracker.StreamTracker
@@ -42,11 +45,12 @@ func NewStreamTrackerManager(
trackersConfig config.StreamTrackersConfig,
) *StreamTrackerManager {
s := &StreamTrackerManager{
logger: logger,
trackInfo: trackInfo,
isSVC: isSVC,
maxPublishedLayer: 0,
clockRate: clockRate,
logger: logger,
trackInfo: trackInfo,
isSVC: isSVC,
maxPublishedLayer: 0,
clockRate: clockRate,
layersNotifyOpQueue: utils.NewOpsQueue(logger, "layer-notify", 10),
}
switch s.trackInfo.Source {
@@ -69,6 +73,14 @@ func NewStreamTrackerManager(
return s
}
func (s *StreamTrackerManager) Start() {
s.layersNotifyOpQueue.Start()
}
func (s *StreamTrackerManager) Stop() {
s.layersNotifyOpQueue.Stop()
}
func (s *StreamTrackerManager) OnAvailableLayersChanged(f func(availableLayers []int32, exemptedLayers []int32)) {
s.onAvailableLayersChanged = f
}
@@ -133,12 +145,14 @@ func (s *StreamTrackerManager) AddTracker(layer int32) *streamtracker.StreamTrac
s.logger.Debugw("StreamTrackerManager add track", "layer", layer)
tracker.OnStatusChanged(func(status streamtracker.StreamStatus) {
s.logger.Debugw("StreamTrackerManager OnStatusChanged", "layer", layer, "status", status)
if status == streamtracker.StreamStatusStopped {
s.removeAvailableLayer(layer)
} else {
s.addAvailableLayer(layer)
}
s.layersNotifyOpQueue.Enqueue(func() {
s.logger.Debugw("StreamTrackerManager OnStatusChanged", "layer", layer, "status", status)
if status == streamtracker.StreamStatusStopped {
s.removeAvailableLayer(layer)
} else {
s.addAvailableLayer(layer)
}
})
})
tracker.OnBitrateAvailable(func() {
if s.onBitrateAvailabilityChanged != nil {
@@ -225,9 +239,9 @@ func (s *StreamTrackerManager) SetMaxExpectedSpatialLayer(layer int32) int32 {
//
// Some higher layer is expected to start.
// If the layer was not stopped (i.e. it will still be in available layers),
// If the layer was not detected as stopped (i.e. it is still in available layers),
// don't need to do anything. If not, reset the stream tracker so that
// the layer is declared available on the first packet
// the layer is declared available on the first packet.
//
// NOTE: There may be a race between checking if a layer is available and
// resetting the tracker, i.e. the track may stop just after checking.
+7 -7
View File
@@ -17,7 +17,7 @@ type targetLayer struct {
type DDVideoLayerSelector struct {
logger logger.Logger
// TODO : fields for frame chain detect
// DD-TODO : fields for frame chain detect
// frameNumberWrapper Uint16Wrapper
// expectKeyFrame bool
@@ -43,7 +43,7 @@ func (s *DDVideoLayerSelector) Select(expPkt *buffer.ExtPacket, tp *TranslationP
if expPkt.DependencyDescriptor.AttachedStructure != nil {
// update decode target layer and active decode targets
// TODO : these targets info can be shared by all the downtracks, no need calculate in every selector
// DD-TODO : these targets info can be shared by all the downtracks, no need calculate in every selector
s.updateDependencyStructure(expPkt.DependencyDescriptor.AttachedStructure)
}
@@ -52,7 +52,7 @@ func (s *DDVideoLayerSelector) Select(expPkt *buffer.ExtPacket, tp *TranslationP
return true
}
// TODO : we don't have a rtp queue to ensure the order of packets now,
// DD-TODO : we don't have a rtp queue to ensure the order of packets now,
// so we don't know packet is lost/out of order, that cause us can't detect
// frame integrity, entire frame is forwareded, whether frame chain is broken.
// So use a simple check here, assume all the reference frame is forwarded and
@@ -69,7 +69,7 @@ func (s *DDVideoLayerSelector) Select(expPkt *buffer.ExtPacket, tp *TranslationP
// find target match with selected layer
if dt.Layer.Spatial <= s.layer.Spatial && dt.Layer.Temporal <= s.layer.Temporal {
if activeDecodeTargets == nil || ((*activeDecodeTargets)&(1<<dt.Target) != 0) {
// TODO : check frame chain integrity
// DD-TODO : check frame chain integrity
currentTarget = dt.Target
// s.logger.Debugw("select target", "target", currentTarget, "layer", dt.layer, "dtis", expPkt.DependencyDescriptor.FrameDependencies.DecodeTargetIndications)
break
@@ -92,7 +92,7 @@ func (s *DDVideoLayerSelector) Select(expPkt *buffer.ExtPacket, tp *TranslationP
return false
}
// TODO : if bandwidth in congest, could drop the 'Discardable' packet
// DD-TODO : if bandwidth in congest, could drop the 'Discardable' packet
if dti := dtis[currentTarget]; dti == dd.DecodeTargetNotPresent {
// s.logger.Debugw(fmt.Sprintf("drop packet for decode target not present, dtis %v, currentTarget %d, s:%d, t:%d", dtis, currentTarget,
// expPkt.DependencyDescriptor.FrameDependencies.SpatialId, expPkt.DependencyDescriptor.FrameDependencies.TemporalId))
@@ -101,7 +101,7 @@ func (s *DDVideoLayerSelector) Select(expPkt *buffer.ExtPacket, tp *TranslationP
tp.switchingToTargetLayer = true
}
// TODO : add frame to forwarded queue if entire frame is forwarded
// DD-TODO : add frame to forwarded queue if entire frame is forwarded
// s.logger.Debugw("select packet", "target", currentTarget, "layer", s.layer)
tp.ddExtension = &dd.DependencyDescriptorExtension{
@@ -176,7 +176,7 @@ func (s *DDVideoLayerSelector) updateDependencyStructure(structure *dd.FrameDepe
s.logger.Debugw(fmt.Sprintf("update decode targets: %v", s.decodeTargetLayer))
}
// TODO : use generic wrapper when updated to go 1.18
// DD-TODO : use generic wrapper when updated to go 1.18
type Uint16Wrapper struct {
last_value *uint16
lastUnwrapped int32