diff --git a/pkg/rtc/subscribedtrack.go b/pkg/rtc/subscribedtrack.go index b01896cc1..3a08c731f 100644 --- a/pkg/rtc/subscribedtrack.go +++ b/pkg/rtc/subscribedtrack.go @@ -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 { diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index 1a0b00003..751b48a10 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -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, } diff --git a/pkg/sfu/forwarder.go b/pkg/sfu/forwarder.go index 4554b4161..255465ee2 100644 --- a/pkg/sfu/forwarder.go +++ b/pkg/sfu/forwarder.go @@ -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) diff --git a/pkg/sfu/forwarder_test.go b/pkg/sfu/forwarder_test.go index ae4acfe53..5248bd587 100644 --- a/pkg/sfu/forwarder_test.go +++ b/pkg/sfu/forwarder_test.go @@ -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, diff --git a/pkg/sfu/receiver.go b/pkg/sfu/receiver.go index 3f3f85b27..95c606d98 100644 --- a/pkg/sfu/receiver.go +++ b/pkg/sfu/receiver.go @@ -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) diff --git a/pkg/sfu/rtpmunger.go b/pkg/sfu/rtpmunger.go index 5ee1045b7..d42a4fe5f 100644 --- a/pkg/sfu/rtpmunger.go +++ b/pkg/sfu/rtpmunger.go @@ -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) diff --git a/pkg/sfu/streamtrackermanager.go b/pkg/sfu/streamtrackermanager.go index 498fdd870..1fdd3fadc 100644 --- a/pkg/sfu/streamtrackermanager.go +++ b/pkg/sfu/streamtrackermanager.go @@ -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. diff --git a/pkg/sfu/videolayerselector.go b/pkg/sfu/videolayerselector.go index 109aa99e3..baad634ed 100644 --- a/pkg/sfu/videolayerselector.go +++ b/pkg/sfu/videolayerselector.go @@ -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<