diff --git a/pkg/sfu/buffer/rtpstats.go b/pkg/sfu/buffer/rtpstats.go index 9a0ba637a..c6ea61498 100644 --- a/pkg/sfu/buffer/rtpstats.go +++ b/pkg/sfu/buffer/rtpstats.go @@ -79,7 +79,6 @@ type Snapshot struct { } type SnInfo struct { - pktTime int64 hdrSize uint16 pktSize uint16 isPaddingOnly bool diff --git a/pkg/sfu/connectionquality/connectionstats_test.go b/pkg/sfu/connectionquality/connectionstats_test.go index 355a4e270..f29d4f3e8 100644 --- a/pkg/sfu/connectionquality/connectionstats_test.go +++ b/pkg/sfu/connectionquality/connectionstats_test.go @@ -606,7 +606,7 @@ func TestConnectionQuality(t *testing.T) { distance: 2.0, }, { - distance: 2.7, + distance: 2.0, offset: 1 * time.Second, }, }, diff --git a/pkg/sfu/connectionquality/scorer.go b/pkg/sfu/connectionquality/scorer.go index 1cca32299..f094c8ecc 100644 --- a/pkg/sfu/connectionquality/scorer.go +++ b/pkg/sfu/connectionquality/scorer.go @@ -19,7 +19,7 @@ const ( increaseFactor = float64(0.4) // slow increase decreaseFactor = float64(0.8) // fast decrease - distanceWeight = float64(20.0) // each spatial layer missed drops a quality level + distanceWeight = float64(25.0) // each spatial layer missed drops a quality level unmuteTimeThreshold = float64(0.5) ) diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index 57811c13a..6bd4c0ea0 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -946,22 +946,12 @@ func (d *DownTrack) UpTrackMaxTemporalLayerSeenChange(maxTemporalLayerSeen int32 d.forwarder.SetMaxTemporalLayerSeen(maxTemporalLayerSeen) } -func (d *DownTrack) maybeAddTransition(bitrate int64, distance float64) { +func (d *DownTrack) maybeAddTransition(_bitrate int64, distance float64) { if d.kind == webrtc.RTPCodecTypeAudio { return } - ti := d.receiver.TrackInfo() - if ti == nil { - return - } - - if ti.Source == livekit.TrackSource_SCREEN_SHARE { - d.connectionStats.AddLayerTransition(distance, time.Now()) - return - } - - d.connectionStats.AddBitrateTransition(bitrate, time.Now()) + d.connectionStats.AddLayerTransition(distance, time.Now()) } func (d *DownTrack) UpTrackBitrateReport(_availableLayers []int32, bitrates Bitrates) { diff --git a/pkg/sfu/forwarder.go b/pkg/sfu/forwarder.go index 1f2277374..1b0dea4a6 100644 --- a/pkg/sfu/forwarder.go +++ b/pkg/sfu/forwarder.go @@ -1702,43 +1702,71 @@ func getDistanceToDesired( targetLayers VideoLayers, maxLayers VideoLayers, ) float64 { - if muted || pubMuted || maxPublishedLayer == InvalidLayerSpatial || !maxLayers.IsValid() { + if muted || pubMuted || maxPublishedLayer == InvalidLayerSpatial || maxTemporalLayerSeen == InvalidLayerTemporal || !maxLayers.IsValid() { return 0.0 } - found := false - distance := float64(0.0) + adjustedMaxLayers := maxLayers + + // max available spatial is min(subscribedMax, publishedMax, availableMax) + // subscribedMax = subscriber requested max spatial layer + // publishedMax = max spatial layer ever published + // availableMax = based on bit rate measurement, available max spatial layer + maxAvailableSpatial := InvalidLayerSpatial done: - for s := maxLayers.Spatial; s >= 0; s-- { - for t := maxLayers.Temporal; t >= 0; t-- { - if brs[s][t] == 0 { - continue - } - if s == targetLayers.Spatial && t == targetLayers.Temporal { - found = true + for s := int32(len(brs)) - 1; s >= 0; s-- { + for t := int32(len(brs[0])) - 1; t >= 0; t-- { + if brs[s][t] != 0 { + maxAvailableSpatial = s break done } - - distance++ } } + if maxAvailableSpatial < adjustedMaxLayers.Spatial { + adjustedMaxLayers.Spatial = maxAvailableSpatial + } - // maybe overshooting - if !found && targetLayers.IsValid() { - distance = 0.0 - for s := targetLayers.Spatial; s > maxLayers.Spatial; s-- { - for t := maxLayers.Temporal; t >= 0; t-- { - if targetLayers.Temporal < t || brs[s][t] == 0 { - continue - } - distance-- + if maxPublishedLayer < adjustedMaxLayers.Spatial { + adjustedMaxLayers.Spatial = maxPublishedLayer + } + + // max available temporal is min(subscribedMax, temporalLayerSeenMax, availableMax) + // subscribedMax = subscriber requested max temporal layer + // temporalLayerSeenMax = max temporal layer ever published/seen + // availableMax = based on bit rate measurement, available max temporal in the adjusted max spatial layer + maxAvailableTemporal := InvalidLayerTemporal + if adjustedMaxLayers.Spatial != InvalidLayerSpatial { + for t := int32(len(brs[0])) - 1; t >= 0; t-- { + if brs[adjustedMaxLayers.Spatial][t] != 0 { + maxAvailableTemporal = t + break } } } - - if maxTemporalLayerSeen < 0 { - maxTemporalLayerSeen = 0 + if maxAvailableTemporal < adjustedMaxLayers.Temporal { + adjustedMaxLayers.Temporal = maxAvailableTemporal } - return distance / float64(maxTemporalLayerSeen+1) + if maxTemporalLayerSeen < adjustedMaxLayers.Temporal { + adjustedMaxLayers.Temporal = maxTemporalLayerSeen + } + + if !adjustedMaxLayers.IsValid() { + adjustedMaxLayers = VideoLayers{Spatial: 0, Temporal: 0} + } + + // adjust target layers if they are invalid, i. e. not streaming + adjustedTargetLayers := targetLayers + if !targetLayers.IsValid() { + adjustedTargetLayers = VideoLayers{Spatial: 0, Temporal: 0} + } + + distance := + ((adjustedMaxLayers.Spatial - adjustedTargetLayers.Spatial) * (maxTemporalLayerSeen + 1)) + + (adjustedMaxLayers.Temporal - adjustedTargetLayers.Temporal) + if !targetLayers.IsValid() { + distance++ + } + + return float64(distance) / float64(maxTemporalLayerSeen+1) } diff --git a/pkg/sfu/forwarder_test.go b/pkg/sfu/forwarder_test.go index 2765bb0ce..5398a5a61 100644 --- a/pkg/sfu/forwarder_test.go +++ b/pkg/sfu/forwarder_test.go @@ -216,6 +216,7 @@ func TestForwarderAllocateOptimal(t *testing.T) { f.parkedLayers = InvalidLayers // when max layers changes, target is opportunistic, but requested spatial layer should be at max + f.SetMaxTemporalLayerSeen(3) f.maxLayers = VideoLayers{Spatial: 1, Temporal: 3} expectedResult = VideoAllocation{ pauseReason: VideoPauseReasonNone, @@ -252,7 +253,7 @@ func TestForwarderAllocateOptimal(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: expectedTargetLayers.Spatial, maxLayers: DefaultMaxLayers, - distanceToDesired: 0, + distanceToDesired: -0.5, } result = f.AllocateOptimal(nil, bitrates, true) require.Equal(t, expectedResult, result) @@ -275,7 +276,7 @@ func TestForwarderAllocateOptimal(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: expectedTargetLayers.Spatial, maxLayers: DefaultMaxLayers, - distanceToDesired: 0, + distanceToDesired: -0.75, } result = f.AllocateOptimal(nil, emptyBitrates, true) require.Equal(t, expectedResult, result) @@ -294,7 +295,7 @@ func TestForwarderAllocateOptimal(t *testing.T) { targetLayers: DefaultMaxLayers, requestLayerSpatial: 2, maxLayers: DefaultMaxLayers, - distanceToDesired: 0, + distanceToDesired: -0.5, } result = f.AllocateOptimal([]int32{0, 1}, bitrates, true) require.Equal(t, expectedResult, result) @@ -311,7 +312,7 @@ func TestForwarderAllocateOptimal(t *testing.T) { targetLayers: DefaultMaxLayers, requestLayerSpatial: 2, maxLayers: DefaultMaxLayers, - distanceToDesired: 0, + distanceToDesired: -2.75, } result = f.AllocateOptimal([]int32{0, 1}, emptyBitrates, false) require.Equal(t, expectedResult, result) @@ -331,7 +332,7 @@ func TestForwarderAllocateOptimal(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: 2, maxLayers: DefaultMaxLayers, - distanceToDesired: 0, + distanceToDesired: -0.5, } result = f.AllocateOptimal([]int32{0, 1}, bitrates, true) require.Equal(t, expectedResult, result) @@ -352,7 +353,7 @@ func TestForwarderAllocateOptimal(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: 1, maxLayers: DefaultMaxLayers, - distanceToDesired: 1, + distanceToDesired: 0.5, } result = f.AllocateOptimal([]int32{1}, bitrates, true) require.Equal(t, expectedResult, result) @@ -374,7 +375,7 @@ func TestForwarderAllocateOptimal(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: 0, maxLayers: f.maxLayers, - distanceToDesired: 0, + distanceToDesired: -0.25, } result = f.AllocateOptimal([]int32{0, 1}, emptyBitrates, true) require.Equal(t, expectedResult, result) @@ -396,7 +397,7 @@ func TestForwarderAllocateOptimal(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: 2, maxLayers: f.maxLayers, - distanceToDesired: 0, + distanceToDesired: -2.75, } result = f.AllocateOptimal([]int32{0, 1}, emptyBitrates, true) require.Equal(t, expectedResult, result) @@ -408,6 +409,7 @@ func TestForwarderProvisionalAllocate(t *testing.T) { f.SetMaxSpatialLayer(DefaultMaxLayerSpatial) f.SetMaxTemporalLayer(DefaultMaxLayerTemporal) f.SetMaxPublishedLayer(DefaultMaxLayerSpatial) + f.SetMaxTemporalLayerSeen(DefaultMaxLayerTemporal) bitrates := Bitrates{ {1, 2, 3, 4}, @@ -447,7 +449,7 @@ func TestForwarderProvisionalAllocate(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: expectedTargetLayers.Spatial, maxLayers: DefaultMaxLayers, - distanceToDesired: 5, + distanceToDesired: 1.25, } result := f.ProvisionalAllocateCommit() require.Equal(t, expectedResult, result) @@ -474,7 +476,7 @@ func TestForwarderProvisionalAllocate(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: expectedTargetLayers.Spatial, maxLayers: DefaultMaxLayers, - distanceToDesired: 11, + distanceToDesired: 2.75, } result = f.ProvisionalAllocateCommit() require.Equal(t, expectedResult, result) @@ -521,7 +523,7 @@ func TestForwarderProvisionalAllocate(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: expectedTargetLayers.Spatial, maxLayers: expectedMaxLayers, - distanceToDesired: -4, + distanceToDesired: -1.75, } result = f.ProvisionalAllocateCommit() require.Equal(t, expectedResult, result) @@ -565,7 +567,7 @@ func TestForwarderProvisionalAllocate(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: expectedTargetLayers.Spatial, maxLayers: expectedMaxLayers, - distanceToDesired: 0, + distanceToDesired: 0.25, } result = f.ProvisionalAllocateCommit() require.Equal(t, expectedResult, result) @@ -598,7 +600,7 @@ func TestForwarderProvisionalAllocate(t *testing.T) { targetLayers: InvalidLayers, requestLayerSpatial: InvalidLayerSpatial, maxLayers: expectedMaxLayers, - distanceToDesired: 0, + distanceToDesired: 0.25, } result = f.ProvisionalAllocateCommit() require.Equal(t, expectedResult, result) @@ -649,6 +651,7 @@ func TestForwarderProvisionalAllocateGetCooperativeTransition(t *testing.T) { f.SetMaxSpatialLayer(DefaultMaxLayerSpatial) f.SetMaxTemporalLayer(DefaultMaxLayerTemporal) f.SetMaxPublishedLayer(DefaultMaxLayerSpatial) + f.SetMaxTemporalLayerSeen(DefaultMaxLayerTemporal) bitrates := Bitrates{ {1, 2, 3, 4}, @@ -678,7 +681,7 @@ func TestForwarderProvisionalAllocateGetCooperativeTransition(t *testing.T) { targetLayers: expectedLayers, requestLayerSpatial: expectedLayers.Spatial, maxLayers: DefaultMaxLayers, - distanceToDesired: 9, + distanceToDesired: 2.25, } result := f.ProvisionalAllocateCommit() require.Equal(t, expectedResult, result) @@ -707,7 +710,7 @@ func TestForwarderProvisionalAllocateGetCooperativeTransition(t *testing.T) { targetLayers: expectedLayers, requestLayerSpatial: expectedLayers.Spatial, maxLayers: DefaultMaxLayers, - distanceToDesired: 0, + distanceToDesired: 0.0, } result = f.ProvisionalAllocateCommit() require.Equal(t, expectedResult, result) @@ -776,7 +779,7 @@ func TestForwarderProvisionalAllocateGetCooperativeTransition(t *testing.T) { targetLayers: expectedLayers, requestLayerSpatial: expectedLayers.Spatial, maxLayers: expectedMaxLayers, - distanceToDesired: -1, + distanceToDesired: -1.0, } result = f.ProvisionalAllocateCommit() require.Equal(t, expectedResult, result) @@ -815,7 +818,7 @@ func TestForwarderProvisionalAllocateGetCooperativeTransition(t *testing.T) { targetLayers: expectedLayers, requestLayerSpatial: expectedLayers.Spatial, maxLayers: expectedMaxLayers, - distanceToDesired: 0, + distanceToDesired: -0.5, } result = f.ProvisionalAllocateCommit() require.Equal(t, expectedResult, result) @@ -830,7 +833,7 @@ func TestForwarderProvisionalAllocateGetCooperativeTransition(t *testing.T) { targetLayers: expectedLayers, requestLayerSpatial: expectedLayers.Spatial, maxLayers: expectedMaxLayers, - distanceToDesired: 0, + distanceToDesired: -0.5, } result = f.ProvisionalAllocateCommit() require.Equal(t, expectedResult, result) @@ -883,6 +886,7 @@ func TestForwarderAllocateNextHigher(t *testing.T) { f.SetMaxSpatialLayer(DefaultMaxLayerSpatial) f.SetMaxTemporalLayer(DefaultMaxLayerTemporal) f.SetMaxPublishedLayer(DefaultMaxLayerSpatial) + f.SetMaxTemporalLayerSeen(DefaultMaxLayerTemporal) // when not in deficient state, does not boost result, boosted = f.AllocateNextHigher(ChannelCapacityInfinity, bitrates, false) @@ -918,7 +922,7 @@ func TestForwarderAllocateNextHigher(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: expectedTargetLayers.Spatial, maxLayers: DefaultMaxLayers, - distanceToDesired: 3, + distanceToDesired: 2.0, } result, boosted = f.AllocateNextHigher(ChannelCapacityInfinity, bitrates, false) require.Equal(t, expectedResult, result) @@ -946,7 +950,7 @@ func TestForwarderAllocateNextHigher(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: expectedTargetLayers.Spatial, maxLayers: DefaultMaxLayers, - distanceToDesired: 2, + distanceToDesired: 1.25, } result, boosted = f.AllocateNextHigher(ChannelCapacityInfinity, bitrates, false) require.Equal(t, expectedResult, result) @@ -970,7 +974,7 @@ func TestForwarderAllocateNextHigher(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: expectedTargetLayers.Spatial, maxLayers: DefaultMaxLayers, - distanceToDesired: 1, + distanceToDesired: 0.5, } result, boosted = f.AllocateNextHigher(ChannelCapacityInfinity, bitrates, false) require.Equal(t, expectedResult, result) @@ -992,7 +996,7 @@ func TestForwarderAllocateNextHigher(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: expectedTargetLayers.Spatial, maxLayers: DefaultMaxLayers, - distanceToDesired: 0, + distanceToDesired: 0.0, } result, boosted = f.AllocateNextHigher(ChannelCapacityInfinity, bitrates, false) require.Equal(t, expectedResult, result) @@ -1027,7 +1031,7 @@ func TestForwarderAllocateNextHigher(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: expectedTargetLayers.Spatial, maxLayers: DefaultMaxLayers, - distanceToDesired: 4, + distanceToDesired: 2.25, } result, boosted = f.AllocateNextHigher(ChannelCapacityInfinity, bitrates, false) require.Equal(t, expectedResult, result) @@ -1045,7 +1049,7 @@ func TestForwarderAllocateNextHigher(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: expectedTargetLayers.Spatial, maxLayers: DefaultMaxLayers, - distanceToDesired: 4, + distanceToDesired: 2.25, } result, boosted = f.AllocateNextHigher(0, bitrates, false) require.Equal(t, expectedResult, result) @@ -1079,7 +1083,7 @@ func TestForwarderAllocateNextHigher(t *testing.T) { targetLayers: expectedTargetLayers, requestLayerSpatial: expectedTargetLayers.Spatial, maxLayers: expectedMaxLayers, - distanceToDesired: -1, + distanceToDesired: -1.0, } // overshoot should return (1, 0) even if there is not enough capacity result, boosted = f.AllocateNextHigher(bitrates[1][0]-1, bitrates, true) diff --git a/pkg/sfu/receiver.go b/pkg/sfu/receiver.go index 3eab251e4..58a3221f8 100644 --- a/pkg/sfu/receiver.go +++ b/pkg/sfu/receiver.go @@ -385,26 +385,10 @@ func (w *WebRTCReceiver) AddDownTrack(track TrackSender) error { return nil } -func (w *WebRTCReceiver) notifyMaxExpectedLayer(layer int32) { - if w.Kind() == webrtc.RTPCodecTypeAudio || w.trackInfo.Source == livekit.TrackSource_SCREEN_SHARE { - // screen share tracks have highly variable bitrate, do not use bit rate based quality for those - return - } - - expectedBitrate := int64(0) - for _, vl := range w.trackInfo.Layers { - l := buffer.VideoQualityToSpatialLayer(vl.Quality, w.trackInfo) - if l <= layer { - expectedBitrate += int64(vl.Bitrate) - } - } - - w.connectionStats.AddBitrateTransition(expectedBitrate, time.Now()) -} - func (w *WebRTCReceiver) SetMaxExpectedSpatialLayer(layer int32) { w.streamTrackerManager.SetMaxExpectedSpatialLayer(layer) - w.notifyMaxExpectedLayer(layer) + + w.connectionStats.AddLayerTransition(w.streamTrackerManager.DistanceToDesired(), time.Now()) } // StreamTrackerManagerListener.OnAvailableLayersChanged @@ -427,7 +411,7 @@ func (w *WebRTCReceiver) OnMaxPublishedLayerChanged(maxPublishedLayer int32) { dt.UpTrackMaxPublishedLayerChange(maxPublishedLayer) } - w.notifyMaxExpectedLayer(maxPublishedLayer) + w.connectionStats.AddLayerTransition(w.streamTrackerManager.DistanceToDesired(), time.Now()) } // StreamTrackerManagerListener.OnMaxTemporalLayerSeenChanged @@ -436,9 +420,7 @@ func (w *WebRTCReceiver) OnMaxTemporalLayerSeenChanged(maxTemporalLayerSeen int3 dt.UpTrackMaxTemporalLayerSeenChange(maxTemporalLayerSeen) } - if w.trackInfo.Source == livekit.TrackSource_SCREEN_SHARE { - w.connectionStats.AddLayerTransition(w.streamTrackerManager.DistanceToDesired(), time.Now()) - } + w.connectionStats.AddLayerTransition(w.streamTrackerManager.DistanceToDesired(), time.Now()) } // StreamTrackerManagerListener.OnMaxAvailableLayerChanged diff --git a/pkg/sfu/streamtrackermanager.go b/pkg/sfu/streamtrackermanager.go index 7b2257d79..d6c6ba63b 100644 --- a/pkg/sfu/streamtrackermanager.go +++ b/pkg/sfu/streamtrackermanager.go @@ -310,14 +310,10 @@ done: return 0.0 } - distance := float64(0.0) - for sp := maxLayers.Spatial; sp <= s.getMaxExpectedLayerLocked(); sp++ { - for t := maxLayers.Temporal; t <= s.maxTemporalLayerSeen; t++ { - distance++ - } - } - - return distance / float64(s.maxTemporalLayerSeen+1) + distance := + ((s.getMaxExpectedLayerLocked() - maxLayers.Spatial) * (s.maxTemporalLayerSeen + 1)) + + (s.maxTemporalLayerSeen - maxLayers.Temporal) + return float64(distance) / float64(s.maxTemporalLayerSeen+1) } func (s *StreamTrackerManager) getMaxExpectedLayerLocked() int32 {