From f3e13ef27fe7bcce7691ed32c38b476b4e9077b7 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Mon, 5 Oct 2026 19:21:30 +0530 Subject: [PATCH] fix: use the most common recent packet rate in connection quality (#4940) * fix: use the most common recent packet rate in connection quality The loop that finds the mode packet rate compared a bin's count with the index of the bin picked so far, so it did not find the most common rate. Take the bin with the highest count, skipping rates below 10 pps so that idle, DTX and static content rates do not become the reference. Halve the counts at each mode update so that the mode follows a lasting change of rate, and keep the previous mode through a long stretch of low rates. Co-Authored-By: Claude Opus 5.5 * fix: start the mode timer from the start time Co-Authored-By: Claude Opus 5.5 --------- Co-authored-by: Claude Opus 5.5 --- .../connectionquality/connectionstats_test.go | 63 +++++++++++++++++++ pkg/sfu/connectionquality/scorer.go | 29 ++++++--- 2 files changed, 83 insertions(+), 9 deletions(-) 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) }