diff --git a/pkg/sfu/connectionquality/connectionstats.go b/pkg/sfu/connectionquality/connectionstats.go index 76841b2da..ac6506bd8 100644 --- a/pkg/sfu/connectionquality/connectionstats.go +++ b/pkg/sfu/connectionquality/connectionstats.go @@ -74,8 +74,12 @@ func (cs *ConnectionStats) UpdateMute(isMuted bool, at time.Time) { cs.scorer.UpdateMute(isMuted, at) } -func (cs *ConnectionStats) AddTransition(bitrate int64, at time.Time) { - cs.scorer.AddTransition(bitrate, at) +func (cs *ConnectionStats) AddBitrateTransition(bitrate int64, at time.Time) { + cs.scorer.AddBitrateTransition(bitrate, at) +} + +func (cs *ConnectionStats) AddLayerTransition(distance float64, at time.Time) { + cs.scorer.AddLayerTransition(distance, at) } func (cs *ConnectionStats) GetScoreAndQuality() (float32, livekit.ConnectionQuality) { diff --git a/pkg/sfu/connectionquality/connectionstats_test.go b/pkg/sfu/connectionquality/connectionstats_test.go index 6fe4ef218..355a4e270 100644 --- a/pkg/sfu/connectionquality/connectionstats_test.go +++ b/pkg/sfu/connectionquality/connectionstats_test.go @@ -273,8 +273,8 @@ func TestConnectionQuality(t *testing.T) { cs.UpdateMute(false, now.Add(2*time.Second)) // bitrate based calculation can drop quality even if there is no loss - cs.AddTransition(1_000_000, now) - cs.AddTransition(2_000_000, now.Add(2*time.Second)) + cs.AddBitrateTransition(1_000_000, now) + cs.AddBitrateTransition(2_000_000, now.Add(2*time.Second)) streams = map[uint32]*buffer.StreamStatsWithLayers{ 1: &buffer.StreamStatsWithLayers{ @@ -293,13 +293,13 @@ func TestConnectionQuality(t *testing.T) { // a transition to 0 (all layers stopped) should flip quality to EXCELLENT now = now.Add(duration) - cs.AddTransition(0, now) + cs.AddBitrateTransition(0, now) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.6), mos) require.Equal(t, livekit.ConnectionQuality_EXCELLENT, quality) }) - t.Run("codecs - packets", func(t *testing.T) { + t.Run("codecs - packet", func(t *testing.T) { type expectedQuality struct { packetLossPercentage float64 expectedMOS float32 @@ -465,7 +465,7 @@ func TestConnectionQuality(t *testing.T) { } }) - t.Run("bytes", func(t *testing.T) { + t.Run("bitrate", func(t *testing.T) { type transition struct { bitrate int64 offset time.Duration @@ -537,7 +537,7 @@ func TestConnectionQuality(t *testing.T) { cs.Start(&livekit.TrackInfo{Type: livekit.TrackType_VIDEO}, now) for _, tr := range tc.transitions { - cs.AddTransition(tr.bitrate, now.Add(tr.offset)) + cs.AddBitrateTransition(tr.bitrate, now.Add(tr.offset)) } streams := map[uint32]*buffer.StreamStatsWithLayers{ @@ -557,4 +557,90 @@ func TestConnectionQuality(t *testing.T) { }) } }) + + t.Run("layer", func(t *testing.T) { + type transition struct { + distance float64 + offset time.Duration + } + testCases := []struct { + name string + transitions []transition + expectedMOS float32 + expectedQuality livekit.ConnectionQuality + }{ + // NOTE: Because of EWMA (Exponentially Weighted Moving Average), these cut off points are not exact + // each spatial layer missed drops o quality level + { + name: "excellent", + transitions: []transition{ + { + distance: 0.5, + }, + { + distance: 0.0, + offset: 3 * time.Second, + }, + }, + expectedMOS: 4.6, + expectedQuality: livekit.ConnectionQuality_EXCELLENT, + }, + { + name: "good", + transitions: []transition{ + { + distance: 1.0, + }, + { + distance: 1.5, + offset: 2 * time.Second, + }, + }, + expectedMOS: 4.1, + expectedQuality: livekit.ConnectionQuality_GOOD, + }, + { + name: "poor", + transitions: []transition{ + { + distance: 2.0, + }, + { + distance: 2.7, + offset: 1 * time.Second, + }, + }, + expectedMOS: 3.2, + expectedQuality: livekit.ConnectionQuality_POOR, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + cs := newConnectionStats("video/vp8", false) + + duration := 5 * time.Second + now := time.Now() + cs.Start(&livekit.TrackInfo{Type: livekit.TrackType_VIDEO}, now) + + for _, tr := range tc.transitions { + cs.AddLayerTransition(tr.distance, now.Add(tr.offset)) + } + + streams := map[uint32]*buffer.StreamStatsWithLayers{ + 123: &buffer.StreamStatsWithLayers{ + RTPStats: &buffer.RTPDeltaInfo{ + StartTime: now, + Duration: duration, + Packets: 200, + }, + }, + } + cs.updateScore(streams, now.Add(duration)) + mos, quality := cs.GetScoreAndQuality() + require.Greater(t, tc.expectedMOS, mos) + require.Equal(t, tc.expectedQuality, quality) + }) + } + }) } diff --git a/pkg/sfu/connectionquality/scorer.go b/pkg/sfu/connectionquality/scorer.go index 8d77f31e6..1cca32299 100644 --- a/pkg/sfu/connectionquality/scorer.go +++ b/pkg/sfu/connectionquality/scorer.go @@ -19,6 +19,8 @@ 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 + unmuteTimeThreshold = float64(0.5) ) @@ -63,7 +65,7 @@ func (w *windowStat) calculatePacketScore(plw float64) float64 { return score } -func (w *windowStat) calculateByteScore(expectedBitrate int64) float64 { +func (w *windowStat) calculateBitrateScore(expectedBitrate int64) float64 { if expectedBitrate == 0 { // unsupported mode OR all layers stopped return maxScore @@ -100,11 +102,16 @@ func (w *windowStat) String() string { // ------------------------------------------ -type layerTransition struct { +type bitrateTransition struct { startedAt time.Time bitrate int64 } +type layerTransition struct { + startedAt time.Time + distance float64 +} + type qualityScorerParams struct { PacketLossWeight float64 Logger logger.Logger @@ -121,12 +128,13 @@ type qualityScorer struct { mutedAt time.Time unmutedAt time.Time - layersMutedAt time.Time - layersUnmutedAt time.Time + bitrateMutedAt time.Time + bitrateUnmutedAt time.Time maxPPS float64 - transitions []layerTransition + bitrateTransitions []bitrateTransition + layerTransitions []layerTransition } func newQualityScorer(params qualityScorerParams) *qualityScorer { @@ -155,31 +163,42 @@ func (q *qualityScorer) UpdateMute(isMuted bool, at time.Time) { } } -func (q *qualityScorer) AddTransition(bitrate int64, at time.Time) { +func (q *qualityScorer) AddBitrateTransition(bitrate int64, at time.Time) { q.lock.Lock() defer q.lock.Unlock() - q.transitions = append(q.transitions, layerTransition{ + q.bitrateTransitions = append(q.bitrateTransitions, bitrateTransition{ startedAt: at, bitrate: bitrate, }) if bitrate == 0 { - q.layersMutedAt = at + q.bitrateMutedAt = at q.score = maxScore } else { - if q.layersUnmutedAt.IsZero() || q.layersMutedAt.After(q.layersUnmutedAt) { - q.layersUnmutedAt = at + if q.bitrateUnmutedAt.IsZero() || q.bitrateMutedAt.After(q.bitrateUnmutedAt) { + q.bitrateUnmutedAt = at } } } +func (q *qualityScorer) AddLayerTransition(distance float64, at time.Time) { + q.lock.Lock() + defer q.lock.Unlock() + + q.layerTransitions = append(q.layerTransitions, layerTransition{ + startedAt: at, + distance: distance, + }) +} + func (q *qualityScorer) Update(stat *windowStat, at time.Time) { q.lock.Lock() defer q.lock.Unlock() // always update transitions expectedBitrate := q.getExpectedBitsAndUpdateTransitions(at) + expectedDistance := q.getExpectedDistanceAndUpdateTransitions(at) // nothing to do when muted or not unmuted for long enough // NOTE: it is possible that unmute -> mute -> unmute transition happens in the @@ -187,7 +206,7 @@ func (q *qualityScorer) Update(stat *windowStat, at time.Time) { // to stable and quality EXCELLENT for responsiveness. On an unmute, the // entire window data is considered (as long as enough time has passed since // unmute) including the data before mute. - if q.isMuted() || !q.isUnmutedEnough(at) || q.areLayersMuted() { + if q.isMuted() || !q.isUnmutedEnough(at) || q.isBitrateMuted() { q.lastUpdateAt = at return } @@ -199,13 +218,24 @@ func (q *qualityScorer) Update(stat *windowStat, at time.Time) { score = poorScore } else { packetScore := stat.calculatePacketScore(q.getPacketLossWeight(stat)) - byteScore := stat.calculateByteScore(expectedBitrate) - if packetScore < byteScore { + bitrateScore := stat.calculateBitrateScore(expectedBitrate) + layerScore := math.Max(math.Min(maxScore, maxScore-(expectedDistance*distanceWeight)), 0.0) + + minScore := math.Min(packetScore, bitrateScore) + minScore = math.Min(minScore, layerScore) + + switch { + case packetScore == minScore: reason = "packet" score = packetScore - } else { + + case bitrateScore == minScore: reason = "bitrate" - score = byteScore + score = bitrateScore + + case layerScore == minScore: + reason = "layer" + score = layerScore } } factor := increaseFactor @@ -224,6 +254,7 @@ func (q *qualityScorer) Update(stat *windowStat, at time.Time) { "quality", scoreToConnectionQuality(score), "stat", stat, "expectedBitrate", expectedBitrate, + "expectedDistance", expectedDistance, ) } @@ -244,10 +275,10 @@ func (q *qualityScorer) isUnmutedEnough(at time.Time) bool { } var sinceLayersUnmute time.Duration - if q.layersUnmutedAt.IsZero() { + if q.bitrateUnmutedAt.IsZero() { sinceLayersUnmute = at.Sub(q.lastUpdateAt) } else { - sinceLayersUnmute = at.Sub(q.layersUnmutedAt) + sinceLayersUnmute = at.Sub(q.bitrateUnmutedAt) } validDuration := sinceUnmute @@ -260,8 +291,8 @@ func (q *qualityScorer) isUnmutedEnough(at time.Time) bool { return validDuration.Seconds()/sinceLastUpdate.Seconds() > unmuteTimeThreshold } -func (q *qualityScorer) areLayersMuted() bool { - return !q.layersMutedAt.IsZero() && (q.layersUnmutedAt.IsZero() || q.layersMutedAt.After(q.layersUnmutedAt)) +func (q *qualityScorer) isBitrateMuted() bool { + return !q.bitrateMutedAt.IsZero() && (q.bitrateUnmutedAt.IsZero() || q.bitrateMutedAt.After(q.bitrateUnmutedAt)) } func (q *qualityScorer) getPacketLossWeight(stat *windowStat) float64 { @@ -283,39 +314,97 @@ func (q *qualityScorer) getPacketLossWeight(stat *windowStat) float64 { } func (q *qualityScorer) getExpectedBitsAndUpdateTransitions(at time.Time) int64 { - if len(q.transitions) == 0 { + if len(q.bitrateTransitions) == 0 { return 0 } var startedAt time.Time totalBits := float64(0.0) - for idx := 0; idx < len(q.transitions)-1; idx++ { - lt := &q.transitions[idx] - ltNext := &q.transitions[idx+1] + for idx := 0; idx < len(q.bitrateTransitions)-1; idx++ { + bt := &q.bitrateTransitions[idx] + btNext := &q.bitrateTransitions[idx+1] + + if bt.startedAt.After(q.lastUpdateAt) { + startedAt = bt.startedAt + } else { + startedAt = q.lastUpdateAt + } + totalBits += btNext.startedAt.Sub(startedAt).Seconds() * float64(bt.bitrate) + } + + // last transition + bt := &q.bitrateTransitions[len(q.bitrateTransitions)-1] + if bt.startedAt.After(q.lastUpdateAt) { + startedAt = bt.startedAt + } else { + startedAt = q.lastUpdateAt + } + totalBits += at.Sub(startedAt).Seconds() * float64(bt.bitrate) + + // set up last bit rate as the startig bit rate for next analysis window + q.bitrateTransitions = []bitrateTransition{bitrateTransition{ + startedAt: at, + bitrate: bt.bitrate, + }} + + return int64(totalBits) +} + +func (q *qualityScorer) getExpectedDistanceAndUpdateTransitions(at time.Time) float64 { + if len(q.layerTransitions) == 0 { + return 0 + } + + var startedAt time.Time + totalDistance := float64(0.0) + totalDuration := time.Duration(0) + for idx := 0; idx < len(q.layerTransitions)-1; idx++ { + lt := &q.layerTransitions[idx] + ltNext := &q.layerTransitions[idx+1] if lt.startedAt.After(q.lastUpdateAt) { startedAt = lt.startedAt } else { startedAt = q.lastUpdateAt } - totalBits += ltNext.startedAt.Sub(startedAt).Seconds() * float64(lt.bitrate) + dur := ltNext.startedAt.Sub(startedAt) + totalDuration += dur + + dist := lt.distance + if dist < 0.0 { + // negative distances are overshoot, that does not compensate for shortfalls, so use optimal, i. e. 0 distance when overshooting + dist = 0.0 + } + totalDistance += dur.Seconds() * float64(dist) } + // last transition - lt := &q.transitions[len(q.transitions)-1] + lt := &q.layerTransitions[len(q.layerTransitions)-1] if lt.startedAt.After(q.lastUpdateAt) { startedAt = lt.startedAt } else { startedAt = q.lastUpdateAt } - totalBits += at.Sub(startedAt).Seconds() * float64(lt.bitrate) + dur := at.Sub(startedAt) + totalDuration += dur - // set up last bit rate as the startig bit rate for next analysis window - q.transitions = []layerTransition{layerTransition{ + dist := lt.distance + if dist < 0.0 { + dist = 0.0 + } + totalDistance += dur.Seconds() * float64(dist) + + // set up last distance as the startig distance for next analysis window + q.layerTransitions = []layerTransition{layerTransition{ startedAt: at, - bitrate: lt.bitrate, + distance: lt.distance, }} - return int64(totalBits) + if totalDuration == 0 { + return 0 + } + + return totalDistance / totalDuration.Seconds() } func (q *qualityScorer) GetScoreAndQuality() (float32, livekit.ConnectionQuality) { diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index 6d2a31df1..2b6ae507e 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -29,6 +29,7 @@ type TrackSender interface { UpTrackLayersChange() UpTrackBitrateAvailabilityChange() UpTrackMaxPublishedLayerChange(maxPublishedLayer int32) + UpTrackMaxTemporalLayerSeenChange(maxTemporalLayerSeen int32) UpTrackBitrateReport(availableLayers []int32, bitrates Bitrates) WriteRTP(p *buffer.ExtPacket, layer int32) error Close() @@ -877,21 +878,29 @@ func (d *DownTrack) UpTrackMaxPublishedLayerChange(maxPublishedLayer int32) { } } -func (d *DownTrack) maybeAddTransition(bitrate int64) { +func (d *DownTrack) UpTrackMaxTemporalLayerSeenChange(maxTemporalLayerSeen int32) { + d.forwarder.SetMaxTemporalLayerSeen(maxTemporalLayerSeen) +} + +func (d *DownTrack) maybeAddTransition(bitrate int64, distance float64) { if d.kind == webrtc.RTPCodecTypeAudio { return } ti := d.receiver.TrackInfo() - if ti == nil || ti.Source == livekit.TrackSource_SCREEN_SHARE { + if ti == nil { return } - d.connectionStats.AddTransition(bitrate, time.Now()) + if ti.Source == livekit.TrackSource_SCREEN_SHARE { + d.connectionStats.AddLayerTransition(distance, time.Now()) + } + + d.connectionStats.AddBitrateTransition(bitrate, time.Now()) } func (d *DownTrack) UpTrackBitrateReport(_availableLayers []int32, bitrates Bitrates) { - d.maybeAddTransition(d.forwarder.GetOptimalBandwidthNeeded(bitrates)) + d.maybeAddTransition(d.forwarder.GetOptimalBandwidthNeeded(bitrates), d.forwarder.DistanceToDesired(bitrates)) } // OnCloseHandler method to be called on remote tracked removed @@ -974,15 +983,16 @@ func (d *DownTrack) BandwidthRequested() int64 { return d.forwarder.BandwidthRequested(brs) } -func (d *DownTrack) DistanceToDesired() int32 { - return d.forwarder.DistanceToDesired() +func (d *DownTrack) DistanceToDesired() float64 { + _, brs := d.receiver.GetLayeredBitrate() + return d.forwarder.DistanceToDesired(brs) } func (d *DownTrack) AllocateOptimal(allowOvershoot bool) VideoAllocation { al, brs := d.receiver.GetLayeredBitrate() allocation := d.forwarder.AllocateOptimal(al, brs, allowOvershoot) d.maybeStartKeyFrameRequester() - d.maybeAddTransition(allocation.bandwidthNeeded) + d.maybeAddTransition(allocation.bandwidthNeeded, allocation.distanceToDesired) return allocation } @@ -1010,7 +1020,7 @@ func (d *DownTrack) ProvisionalAllocateGetBestWeightedTransition() VideoTransiti func (d *DownTrack) ProvisionalAllocateCommit() VideoAllocation { allocation := d.forwarder.ProvisionalAllocateCommit() d.maybeStartKeyFrameRequester() - d.maybeAddTransition(allocation.bandwidthNeeded) + d.maybeAddTransition(allocation.bandwidthNeeded, allocation.distanceToDesired) return allocation } @@ -1018,7 +1028,7 @@ func (d *DownTrack) AllocateNextHigher(availableChannelCapacity int64, allowOver _, brs := d.receiver.GetLayeredBitrate() allocation, available := d.forwarder.AllocateNextHigher(availableChannelCapacity, brs, allowOvershoot) d.maybeStartKeyFrameRequester() - d.maybeAddTransition(allocation.bandwidthNeeded) + d.maybeAddTransition(allocation.bandwidthNeeded, allocation.distanceToDesired) return allocation, available } @@ -1033,7 +1043,7 @@ func (d *DownTrack) Pause() VideoAllocation { _, brs := d.receiver.GetLayeredBitrate() allocation := d.forwarder.Pause(brs) d.maybeStartKeyFrameRequester() - d.maybeAddTransition(allocation.bandwidthNeeded) + d.maybeAddTransition(allocation.bandwidthNeeded, allocation.distanceToDesired) return allocation } diff --git a/pkg/sfu/forwarder.go b/pkg/sfu/forwarder.go index 768cc9353..1f2277374 100644 --- a/pkg/sfu/forwarder.go +++ b/pkg/sfu/forwarder.go @@ -64,11 +64,11 @@ type VideoAllocation struct { targetLayers VideoLayers requestLayerSpatial int32 maxLayers VideoLayers - distanceToDesired int32 + distanceToDesired float64 } func (v VideoAllocation) String() string { - return fmt.Sprintf("VideoAllocation{pause: %s, def: %+v, bwr: %d, del: %d, bwn: %d, rates: %+v, target: %s, req: %d, max: %s, dist: %d}", + return fmt.Sprintf("VideoAllocation{pause: %s, def: %+v, bwr: %d, del: %d, bwn: %d, rates: %+v, target: %s, req: %d, max: %s, dist: %0.2f}", v.pauseReason, v.isDeficient, v.bandwidthRequested, @@ -94,14 +94,15 @@ var ( // ------------------------------------------------------------------- type VideoAllocationProvisional struct { - muted bool - pubMuted bool - maxPublishedLayer int32 - bitrates Bitrates - maxLayers VideoLayers - currentLayers VideoLayers - parkedLayers VideoLayers - allocatedLayers VideoLayers + muted bool + pubMuted bool + maxPublishedLayer int32 + maxTemporalLayerSeen int32 + bitrates Bitrates + maxLayers VideoLayers + currentLayers VideoLayers + parkedLayers VideoLayers + allocatedLayers VideoLayers } // ------------------------------------------------------------------- @@ -178,7 +179,8 @@ type Forwarder struct { muted bool pubMuted bool - maxPublishedLayer int32 + maxPublishedLayer int32 + maxTemporalLayerSeen int32 started bool lastSSRC uint32 @@ -215,7 +217,8 @@ func NewForwarder( logger: logger, getReferenceLayerRTPTimestamp: getReferenceLayerRTPTimestamp, - maxPublishedLayer: InvalidLayerSpatial, + maxPublishedLayer: InvalidLayerSpatial, + maxTemporalLayerSeen: InvalidLayerTemporal, referenceLayerSpatial: InvalidLayerSpatial, @@ -251,6 +254,18 @@ func (f *Forwarder) SetMaxPublishedLayer(maxPublishedLayer int32) { f.logger.Infow("setting max published layer", "maxPublishedLayer", f.maxPublishedLayer) } +func (f *Forwarder) SetMaxTemporalLayerSeen(maxTemporalLayerSeen int32) { + f.lock.Lock() + defer f.lock.Unlock() + + if maxTemporalLayerSeen <= f.maxTemporalLayerSeen { + return + } + + f.maxTemporalLayerSeen = maxTemporalLayerSeen + f.logger.Infow("setting max temporal layer seen", "maxTemporalLayerSeen", f.maxTemporalLayerSeen) +} + func (f *Forwarder) OnParkedLayersExpired(fn func()) { f.lock.Lock() defer f.lock.Unlock() @@ -482,11 +497,11 @@ func (f *Forwarder) BandwidthRequested(brs Bitrates) int64 { return brs[f.targetLayers.Spatial][f.targetLayers.Temporal] } -func (f *Forwarder) DistanceToDesired() int32 { +func (f *Forwarder) DistanceToDesired(brs Bitrates) float64 { f.lock.RLock() defer f.lock.RUnlock() - return f.lastAllocation.distanceToDesired + return getDistanceToDesired(f.muted, f.pubMuted, f.maxPublishedLayer, f.maxTemporalLayerSeen, brs, f.targetLayers, f.maxLayers) } func (f *Forwarder) GetOptimalBandwidthNeeded(brs Bitrates) int64 { @@ -605,7 +620,7 @@ func (f *Forwarder) AllocateOptimal(availableLayers []int32, brs Bitrates, allow alloc.bandwidthRequested = optimalBandwidthNeeded } alloc.bandwidthDelta = alloc.bandwidthRequested - f.lastAllocation.bandwidthRequested - alloc.distanceToDesired = getDistanceToDesired(f.muted, f.pubMuted, f.maxPublishedLayer, brs, alloc.targetLayers, f.maxLayers) + alloc.distanceToDesired = getDistanceToDesired(f.muted, f.pubMuted, f.maxPublishedLayer, f.maxTemporalLayerSeen, brs, alloc.targetLayers, f.maxLayers) return f.updateAllocation(alloc, "optimal") } @@ -615,14 +630,15 @@ func (f *Forwarder) ProvisionalAllocatePrepare(bitrates Bitrates) { defer f.lock.Unlock() f.provisional = &VideoAllocationProvisional{ - allocatedLayers: InvalidLayers, - muted: f.muted, - pubMuted: f.pubMuted, - maxPublishedLayer: f.maxPublishedLayer, - bitrates: bitrates, - maxLayers: f.maxLayers, - currentLayers: f.currentLayers, - parkedLayers: f.parkedLayers, + allocatedLayers: InvalidLayers, + muted: f.muted, + pubMuted: f.pubMuted, + maxPublishedLayer: f.maxPublishedLayer, + maxTemporalLayerSeen: f.maxTemporalLayerSeen, + bitrates: bitrates, + maxLayers: f.maxLayers, + currentLayers: f.currentLayers, + parkedLayers: f.parkedLayers, } } @@ -928,6 +944,7 @@ func (f *Forwarder) ProvisionalAllocateCommit() VideoAllocation { f.provisional.muted, f.provisional.pubMuted, f.provisional.maxPublishedLayer, + f.provisional.maxTemporalLayerSeen, f.provisional.bitrates, f.provisional.allocatedLayers, f.provisional.maxLayers, @@ -1036,7 +1053,7 @@ func (f *Forwarder) AllocateNextHigher(availableChannelCapacity int64, brs Bitra targetLayers: targetLayers, requestLayerSpatial: targetLayers.Spatial, maxLayers: f.maxLayers, - distanceToDesired: getDistanceToDesired(f.muted, f.pubMuted, f.maxPublishedLayer, brs, targetLayers, f.maxLayers), + distanceToDesired: getDistanceToDesired(f.muted, f.pubMuted, f.maxPublishedLayer, f.maxTemporalLayerSeen, brs, targetLayers, f.maxLayers), } if targetLayers.GreaterThan(f.maxLayers) || bandwidthRequested >= optimalBandwidthNeeded { alloc.isDeficient = false @@ -1183,7 +1200,7 @@ func (f *Forwarder) Pause(brs Bitrates) VideoAllocation { targetLayers: InvalidLayers, requestLayerSpatial: InvalidLayerSpatial, maxLayers: f.maxLayers, - distanceToDesired: getDistanceToDesired(f.muted, f.pubMuted, f.maxPublishedLayer, brs, InvalidLayers, f.maxLayers), + distanceToDesired: getDistanceToDesired(f.muted, f.pubMuted, f.maxPublishedLayer, f.maxTemporalLayerSeen, brs, InvalidLayers, f.maxLayers), } switch { @@ -1680,16 +1697,18 @@ func getDistanceToDesired( muted bool, pubMuted bool, maxPublishedLayer int32, + maxTemporalLayerSeen int32, brs Bitrates, targetLayers VideoLayers, maxLayers VideoLayers, -) int32 { +) float64 { if muted || pubMuted || maxPublishedLayer == InvalidLayerSpatial || !maxLayers.IsValid() { - return 0 + return 0.0 } found := false - distance := int32(0) + distance := float64(0.0) +done: for s := maxLayers.Spatial; s >= 0; s-- { for t := maxLayers.Temporal; t >= 0; t-- { if brs[s][t] == 0 { @@ -1697,20 +1716,16 @@ func getDistanceToDesired( } if s == targetLayers.Spatial && t == targetLayers.Temporal { found = true - break + break done } distance++ } - - if found { - break - } } // maybe overshooting if !found && targetLayers.IsValid() { - distance = 0 + 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 { @@ -1721,5 +1736,9 @@ func getDistanceToDesired( } } - return distance + if maxTemporalLayerSeen < 0 { + maxTemporalLayerSeen = 0 + } + + return distance / float64(maxTemporalLayerSeen+1) } diff --git a/pkg/sfu/forwarder_test.go b/pkg/sfu/forwarder_test.go index 0ba5241dc..2765bb0ce 100644 --- a/pkg/sfu/forwarder_test.go +++ b/pkg/sfu/forwarder_test.go @@ -1094,6 +1094,7 @@ func TestForwarderPause(t *testing.T) { f.SetMaxSpatialLayer(DefaultMaxLayerSpatial) f.SetMaxTemporalLayer(DefaultMaxLayerTemporal) f.SetMaxPublishedLayer(DefaultMaxLayerSpatial) + f.SetMaxTemporalLayerSeen(DefaultMaxLayerTemporal) bitrates := Bitrates{ {1, 2, 3, 4}, @@ -1116,7 +1117,7 @@ func TestForwarderPause(t *testing.T) { targetLayers: InvalidLayers, requestLayerSpatial: InvalidLayerSpatial, maxLayers: DefaultMaxLayers, - distanceToDesired: 12, + distanceToDesired: 3, } result := f.Pause(bitrates) require.Equal(t, expectedResult, result) diff --git a/pkg/sfu/receiver.go b/pkg/sfu/receiver.go index e09ebee45..b0b4f1576 100644 --- a/pkg/sfu/receiver.go +++ b/pkg/sfu/receiver.go @@ -195,6 +195,7 @@ func NewWebRTCReceiver( w.streamTrackerManager.OnAvailableLayersChanged(w.downTrackLayerChange) w.streamTrackerManager.OnBitrateAvailabilityChanged(w.downTrackBitrateAvailabilityChange) w.streamTrackerManager.OnMaxPublishedLayerChanged(w.downTrackMaxPublishedLayerChange) + w.streamTrackerManager.OnMaxTemporalLayerSeenChanged(w.downTrackMaxTemporalLayerSeenChange) w.streamTrackerManager.OnBitrateReport(w.downTrackBitrateReport) for _, opt := range opts { @@ -400,7 +401,7 @@ func (w *WebRTCReceiver) notifyMaxExpectedLayer(layer int32) { } } - w.connectionStats.AddTransition(expectedBitrate, time.Now()) + w.connectionStats.AddBitrateTransition(expectedBitrate, time.Now()) } func (w *WebRTCReceiver) SetMaxExpectedSpatialLayer(layer int32) { @@ -428,10 +429,24 @@ func (w *WebRTCReceiver) downTrackMaxPublishedLayerChange(maxPublishedLayer int3 w.notifyMaxExpectedLayer(maxPublishedLayer) } +func (w *WebRTCReceiver) downTrackMaxTemporalLayerSeenChange(maxTemporalLayerSeen int32) { + for _, dt := range w.downTrackSpreader.GetDownTracks() { + dt.UpTrackMaxTemporalLayerSeenChange(maxTemporalLayerSeen) + } + + if w.trackInfo.Source == livekit.TrackSource_SCREEN_SHARE { + w.connectionStats.AddLayerTransition(w.streamTrackerManager.DistanceToDesired(), time.Now()) + } +} + func (w *WebRTCReceiver) downTrackBitrateReport(availableLayers []int32, bitrates Bitrates) { for _, dt := range w.downTrackSpreader.GetDownTracks() { dt.UpTrackBitrateReport(availableLayers, bitrates) } + + if w.trackInfo.Source == livekit.TrackSource_SCREEN_SHARE { + w.connectionStats.AddLayerTransition(w.streamTrackerManager.DistanceToDesired(), time.Now()) + } } func (w *WebRTCReceiver) GetLayeredBitrate() ([]int32, Bitrates) { diff --git a/pkg/sfu/streamallocator.go b/pkg/sfu/streamallocator.go index 54906c964..d8f2c2136 100644 --- a/pkg/sfu/streamallocator.go +++ b/pkg/sfu/streamallocator.go @@ -1461,7 +1461,7 @@ func (t *Track) BandwidthRequested() int64 { return t.downTrack.BandwidthRequested() } -func (t *Track) DistanceToDesired() int32 { +func (t *Track) DistanceToDesired() float64 { return t.downTrack.DistanceToDesired() } diff --git a/pkg/sfu/streamtrackermanager.go b/pkg/sfu/streamtrackermanager.go index 9d0f7dfd4..ead3c2095 100644 --- a/pkg/sfu/streamtrackermanager.go +++ b/pkg/sfu/streamtrackermanager.go @@ -15,15 +15,16 @@ import ( ) type StreamTrackerManager struct { - logger logger.Logger - trackInfo *livekit.TrackInfo - isSVC bool - maxPublishedLayer int32 - clockRate uint32 + logger logger.Logger + trackInfo *livekit.TrackInfo + isSVC bool + clockRate uint32 trackerConfig config.StreamTrackerConfig - lock sync.RWMutex + lock sync.RWMutex + maxPublishedLayer int32 + maxTemporalLayerSeen int32 trackers [DefaultMaxLayerSpatial + 1]*streamtracker.StreamTracker @@ -36,11 +37,12 @@ type StreamTrackerManager struct { closed core.Fuse - onAvailableLayersChanged func() - onBitrateAvailabilityChanged func() - onMaxPublishedLayerChanged func(maxPublishedLayer int32) - onMaxAvailableLayerChanged func(maxAvailableLayer int32) - onBitrateReport func(availableLayers []int32, bitrates Bitrates) + onAvailableLayersChanged func() + onBitrateAvailabilityChanged func() + onMaxPublishedLayerChanged func(maxPublishedLayer int32) + onMaxTemporalLayerSeenChanged func(maxTemporalLayerSeen int32) + onMaxAvailableLayerChanged func(maxAvailableLayer int32) + onBitrateReport func(availableLayers []int32, bitrates Bitrates) } func NewStreamTrackerManager( @@ -51,12 +53,13 @@ func NewStreamTrackerManager( trackersConfig config.StreamTrackersConfig, ) *StreamTrackerManager { s := &StreamTrackerManager{ - logger: logger, - trackInfo: trackInfo, - isSVC: isSVC, - maxPublishedLayer: InvalidLayerSpatial, - clockRate: clockRate, - closed: core.NewFuse(), + logger: logger, + trackInfo: trackInfo, + isSVC: isSVC, + maxPublishedLayer: InvalidLayerSpatial, + maxTemporalLayerSeen: InvalidLayerTemporal, + clockRate: clockRate, + closed: core.NewFuse(), } switch s.trackInfo.Source { @@ -90,6 +93,12 @@ func (s *StreamTrackerManager) OnMaxPublishedLayerChanged(f func(maxPublishedLay s.onMaxPublishedLayerChanged = f } +func (s *StreamTrackerManager) OnMaxTemporalLayerSeenChanged(f func(maxTemporalLayerSeen int32)) { + s.lock.Lock() + s.onMaxTemporalLayerSeenChanged = f + s.lock.Unlock() +} + func (s *StreamTrackerManager) OnMaxLayerChanged(f func(maxAvailableLayer int32)) { s.onMaxAvailableLayerChanged = f } @@ -282,6 +291,53 @@ func (s *StreamTrackerManager) SetMaxExpectedSpatialLayer(layer int32) int32 { return prev } +func (s *StreamTrackerManager) DistanceToDesired() float64 { + s.lock.RLock() + defer s.lock.RUnlock() + + if s.paused { + return 0 + } + + _, brs := s.getLayeredBitrateLocked() + + maxLayers := InvalidLayers +done: + for s := int32(len(brs)) - 1; s >= 0; s-- { + for t := int32(len(brs[0])) - 1; t >= 0; t-- { + if brs[s][t] != 0 { + maxLayers = VideoLayers{ + Spatial: s, + Temporal: t, + } + break done + } + } + } + + distance := float64(0.0) + for sp := maxLayers.Spatial; sp <= s.getMaxExpectedLayerLocked(); sp++ { + for t := maxLayers.Temporal; t <= s.maxTemporalLayerSeen; t++ { + distance++ + } + } + + if s.maxTemporalLayerSeen < 0 { + return distance + } + + return distance / float64(s.maxTemporalLayerSeen+1) +} + +func (s *StreamTrackerManager) getMaxExpectedLayerLocked() int32 { + // find min of layer + maxExpectedLayer := s.maxExpectedLayer + if maxExpectedLayer > s.maxPublishedLayer { + maxExpectedLayer = s.maxPublishedLayer + } + return maxExpectedLayer +} + func (s *StreamTrackerManager) GetMaxPublishedLayer() int32 { s.lock.RLock() defer s.lock.RUnlock() @@ -293,6 +349,10 @@ func (s *StreamTrackerManager) GetLayeredBitrate() ([]int32, Bitrates) { s.lock.RLock() defer s.lock.RUnlock() + return s.getLayeredBitrateLocked() +} + +func (s *StreamTrackerManager) getLayeredBitrateLocked() ([]int32, Bitrates) { var br Bitrates for i, tracker := range s.trackers { @@ -486,6 +546,33 @@ func (s *StreamTrackerManager) GetReferenceLayerRTPTimestamp(ts uint32, layer in return ts + (srRef.SenderReportData.RTPTimestamp - normalizedTS), nil } +func (s *StreamTrackerManager) updateMaxTemporalLayerSeen(brs Bitrates) { + maxTemporalLayerSeen := InvalidLayerTemporal +done: + for t := int32(len(brs[0])) - 1; t >= 0; t-- { + for s := int32(len(brs)) - 1; s >= 0; s-- { + if brs[s][t] != 0 { + maxTemporalLayerSeen = t + break done + } + } + } + + s.lock.Lock() + if maxTemporalLayerSeen <= s.maxTemporalLayerSeen { + s.lock.Unlock() + return + } + + s.maxTemporalLayerSeen = maxTemporalLayerSeen + onMaxTemporalLayerSeenChanged := s.onMaxTemporalLayerSeenChanged + s.lock.Unlock() + + if onMaxTemporalLayerSeenChanged != nil { + onMaxTemporalLayerSeenChanged(maxTemporalLayerSeen) + } +} + func (s *StreamTrackerManager) bitrateReporter() { ticker := time.NewTicker(time.Second) defer ticker.Stop() @@ -496,13 +583,17 @@ func (s *StreamTrackerManager) bitrateReporter() { return case <-ticker.C: + al, brs := s.GetLayeredBitrate() + s.updateMaxTemporalLayerSeen(brs) + s.lock.RLock() onBitrateReport := s.onBitrateReport s.lock.RUnlock() if onBitrateReport != nil { - onBitrateReport(s.GetLayeredBitrate()) + onBitrateReport(al, brs) } + } } }