diff --git a/pkg/sfu/connectionquality/connectionstats_test.go b/pkg/sfu/connectionquality/connectionstats_test.go index 13b61a645..5adb322c8 100644 --- a/pkg/sfu/connectionquality/connectionstats_test.go +++ b/pkg/sfu/connectionquality/connectionstats_test.go @@ -901,3 +901,66 @@ func TestConnectionQuality(t *testing.T) { } }) } + +func TestPacketRateMode(t *testing.T) { + type phase struct { + minutes float64 + pps float64 + } + for name, tc := range map[string]struct { + weight float64 + histogram map[int]int // bin -> count + phases []phase + last *windowStat + mode int // pps + quality livekit.ConnectionQuality + }{ + // 120 windows on the 250 pps layer, 600 windows on the 40 pps layer + "most common bin": { + weight: 10.0, + histogram: map[int]int{125: 120, 20: 600}, + last: &windowStat{packets: 200, packetsLost: 8, duration: 5 * time.Second}, + mode: 40, + quality: livekit.ConnectionQuality_GOOD, + }, + // 4% loss on the layer now watched is not EXCELLENT + "follows a layer switch": { + weight: 10.0, + phases: []phase{{10, 250}, {5, 40}}, + last: &windowStat{packets: 200, packetsLost: 8, duration: 5 * time.Second}, + mode: 40, + quality: livekit.ConnectionQuality_GOOD, + }, + // idle and DTX rates do not become the reference, so a lost DTX packet counts less + "keeps the speech rate through a long silence": { + weight: 8.0, + phases: []phase{{5, 50}, {10, 0}, {30, 2.5}}, + last: &windowStat{packets: 12, packetsLost: 1, duration: 5 * time.Second}, + mode: 50, + quality: livekit.ConnectionQuality_EXCELLENT, + }, + } { + t.Run(name, func(t *testing.T) { + q := newQualityScorer(qualityScorerParams{Logger: logger.GetLogger()}) + // a timeline in the past, the mode must follow update times, not the wall clock + at := time.Now().Add(-time.Hour) + q.StartAt(tc.weight, at) + for bin, count := range tc.histogram { + q.ppsHistogram[bin] = count + q.numPPSReadings += count + } + for _, p := range tc.phases { + for range int(p.minutes * 12) { + at = at.Add(5 * time.Second) + q.UpdateAt(&windowStat{packets: uint32(p.pps * 5), duration: 5 * time.Second}, at) + } + } + at = at.Add(tc.last.duration) + q.UpdateAt(tc.last, at) + + require.Equal(t, tc.mode, int(float64(q.ppsMode)*cPPSQuantization)) + _, quality := q.GetScoreAndQuality() + require.Equal(t, tc.quality, quality) + }) + } +} diff --git a/pkg/sfu/connectionquality/scorer.go b/pkg/sfu/connectionquality/scorer.go index 471b291b6..57aa3e2e9 100644 --- a/pkg/sfu/connectionquality/scorer.go +++ b/pkg/sfu/connectionquality/scorer.go @@ -42,6 +42,7 @@ const ( cPPSQuantization = float64(2) cPPSMinReadings = 10 + cPPSModeMinBin = 5 // 10 pps cModeCalculationInterval = 2 * time.Minute ) @@ -235,13 +236,13 @@ func newQualityScorer(params qualityScorerParams) *qualityScorer { layerDistance: utils.NewTimedAggregator[float64](utils.TimedAggregatorParams{ CapNegativeValues: true, }), - modeCalculatedAt: time.Now().Add(-cModeCalculationInterval), } } func (q *qualityScorer) startAtLocked(packetLossWeight float64, at time.Time) { q.packetLossWeight = packetLossWeight q.lastUpdateAt = at + q.modeCalculatedAt = at.Add(-cModeCalculationInterval) } func (q *qualityScorer) StartAt(packetLossWeight float64, at time.Time) { @@ -410,7 +411,7 @@ func (q *qualityScorer) updateAtLocked(stat *windowStat, at time.Time) { return } - aplw := q.getAdjustedPacketLossWeight(stat) + aplw := q.getAdjustedPacketLossWeight(stat, at) reason := "none" var score, packetScore, bitrateScore, layerScore float64 if stat.packets+stat.packetsPadding == 0 { @@ -544,7 +545,7 @@ func (q *qualityScorer) isPaused() bool { return !q.pausedAt.IsZero() && (q.resumedAt.IsZero() || q.pausedAt.After(q.resumedAt)) } -func (q *qualityScorer) getAdjustedPacketLossWeight(stat *windowStat) float64 { +func (q *qualityScorer) getAdjustedPacketLossWeight(stat *windowStat, at time.Time) float64 { if stat == nil || stat.duration <= 0 { return q.packetLossWeight } @@ -566,14 +567,24 @@ func (q *qualityScorer) getAdjustedPacketLossWeight(stat *windowStat) float64 { // calculate mode sparingly, do it under the following conditions // 1. minimum number of readings available (AND) // 2. enough time has elapsed since last calculation - if q.numPPSReadings > cPPSMinReadings && time.Since(q.modeCalculatedAt) > cModeCalculationInterval { - q.ppsMode = 0 - for i := range len(q.ppsHistogram) { - if q.ppsHistogram[i] > q.ppsMode { - q.ppsMode = i + if q.numPPSReadings > cPPSMinReadings && at.Sub(q.modeCalculatedAt) > cModeCalculationInterval { + // the most common rate, ignoring rates as low as DTX or static content so they do not become the reference, + // keep the previous mode through a long stretch of only low rates, for example a long silence + mode, maxCount := 0, 0 + for i := cPPSModeMinBin; i < len(q.ppsHistogram); i++ { + if q.ppsHistogram[i] > maxCount { + mode, maxCount = i, q.ppsHistogram[i] } } - q.modeCalculatedAt = time.Now() + if maxCount > 0 { + q.ppsMode = mode + } + + // halve the counts so that the mode follows a lasting change of rate, for example a layer switch + for i := range q.ppsHistogram { + q.ppsHistogram[i] >>= 1 + } + q.modeCalculatedAt = at q.params.Logger.Debugw("updating pps mode", "expected", stat.packets, "duration", stat.duration.Seconds(), "pps", pps, "ppsMode", q.ppsMode) }