From 6ed445bad167a981561809de71d9d76a2859bde6 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Tue, 29 Sep 2026 01:39:05 +0530 Subject: [PATCH] sfu: report forwarding latency as p90 instead of the mean (#4920) * sfu: report forwarding latency as p90 instead of the mean The forwarding-latency metric was the mean transit over all forwarded packets. A mean is dominated by a few slow outliers, so a handful of packets stalled on the forward path (e.g. goroutine scheduling latency) inflated the whole node's reported latency even when nearly every packet was forwarded promptly. Report p90 instead: p90 rising means roughly a tenth of forwarded packets are slow, a broad signal of systemic forwarding load rather than a sparse tail. To read a percentile over the report window, the mergeable per-interval summary now keeps a small power-of-two-bucket histogram of transit instead of running moments (sum, sum-of-squares). Drop the jitter (transit std dev) gauge: nothing consumed it, and any spread is derivable from the forward-latency histogram. The protobuf ForwardJitter field is left in place, now unset, to deprecate separately. Co-Authored-By: Claude Opus 4.8 * sfu: clamp forwarding percentile to the observed [min, max] Bucket interpolation assumes a uniform fill, so a single 20ms packet (or uniform traffic) could report a p90 above every observed sample. Clamp the interpolated value to the summary's already-tracked min/max, so a quantile never falls outside the data. Exact for single samples and repeated identical latencies. Co-Authored-By: Claude Opus 4.8 * sfu: fix stale p50 comment in the percentile test The reported metric is p90; the test comment still said p50 replaced the mean. Reword it to reflect that a percentile, not the mean, is reported. Co-Authored-By: Claude Opus 4.8 * sfu: quarter-octave forwarding-latency buckets for threshold resolution Octave buckets are too coarse near the overload thresholds: 300us falls in [256,512), so a p90 clustered at ~265us and one at ~500us interpolate to the same value and would trip (or not) identically. Split each octave into four linear sub-buckets so the two land in different buckets, on the correct side of the threshold. Min/max clamping alone does not fix this once a node has a high tail, since its max no longer bounds the interpolation. Co-Authored-By: Claude Opus 4.8 * sfu: interpolate percentiles within each bucket's observed range Replace the quarter-octave split with plain octave buckets that also carry the observed [min, max] of their samples, and interpolate a percentile within that range instead of the bucket's nominal edges. This is exact when a bucket's samples cluster, so a p90 near an overload threshold that falls mid-bucket lands on the correct side of it regardless of bucket width -- no threshold-aware boundaries needed. It subsumes the min/max clamp, since an estimate can no longer leave the observed samples. Co-Authored-By: Claude Opus 4.8 * sfu: cover the 500us cluster in the threshold test Assert milos's full review example exactly: 265us and 500us clusters that octave-nominal interpolation both read as ~341us now read 265us and 500us. Co-Authored-By: Claude Opus 4.8 * sfu: make forwardSummary.addSample a pointer receiver The summary grew from a small moments struct into a per-bucket histogram (~700 bytes), so the value-receiver addSample copied the whole summary on every drained sample in the flush loop. Mutate in place instead: ~28ns -> ~2ns per sample in the background fold, no behavior change. Co-Authored-By: Claude Opus 4.8 * sfu: eighth-octave buckets for accuracy near thresholds Split each octave into 8 linear sub-buckets (was plain octaves). With the per-bucket min/max interpolation this reports a tight p90 cluster exactly even when it sits mid-octave: 850x265us + 150x410us now reads 410us (above a 400us threshold) instead of 395.5us, and a lognormal p90 lands within ~0.2us of exact. Per-sample add cost is unchanged (~2ns); cost is ~5KB per summary and a larger but per-report merge. Co-Authored-By: Claude Opus 4.8 --------- Co-authored-by: Claude Opus 4.8 --- pkg/sfu/forwardstats.go | 183 ++++++++++++++++++---------- pkg/sfu/forwardstats_test.go | 133 +++++++++++++------- pkg/telemetry/prometheus/node.go | 1 - pkg/telemetry/prometheus/packets.go | 20 +-- 4 files changed, 214 insertions(+), 123 deletions(-) diff --git a/pkg/sfu/forwardstats.go b/pkg/sfu/forwardstats.go index 2ad8497a7..a7dfbb0dc 100644 --- a/pkg/sfu/forwardstats.go +++ b/pkg/sfu/forwardstats.go @@ -1,7 +1,7 @@ package sfu import ( - "math" + "math/bits" "sync" "time" @@ -13,6 +13,12 @@ import ( const ( cHighForwardingLatency = 20 * time.Millisecond cSkewFactor = 10 + + // forwardLatencyPercentile is the transit quantile reported as the node's + // forwarding latency and read by the overload controller. p90 tracks a broad + // slowdown but ignores the sparse tail (a few packets stalled on goroutine + // scheduling), which shedding cannot fix. + forwardLatencyPercentile = 0.90 ) const ( @@ -94,32 +100,86 @@ func (b *forwardSampleBuffer) takeDropped() uint64 { return b.dropped.Swap(0) } -// forwardSummary is a mergeable summary of forwarding transit over an interval. -// The sum of squares is kept in microseconds so it does not overflow int64. -type forwardSummary struct { - count int64 - sumUs int64 - sumSqUs int64 - minNs int64 - maxNs int64 +// forwardSummary is a mergeable histogram of forwarding transit over an +// interval. Each power-of-two octave [2^e, 2^(e+1)) us is split into +// forwardHistSub linear sub-buckets, and each bucket keeps the observed +// [min, max] of the samples in it. A percentile interpolates within the +// crossing bucket's observed range rather than its nominal edges, so the +// estimate never leaves the samples and is exact when a bucket's samples +// cluster -- which keeps it accurate near an overload threshold that falls +// mid-bucket, without needing threshold-aware boundaries. Two summaries merge +// by adding their buckets, so the percentile covers the whole report window, +// and a few stalled packets (e.g. from goroutine scheduling latency) cannot +// drag it the way a mean does. +const ( + forwardHistSub = 8 // linear sub-buckets per octave + forwardHistOctaves = 27 // up to 2^27 us (~134 s) + forwardHistBuckets = forwardHistOctaves * forwardHistSub // 216 +) + +// bucketStat is one bucket: how many samples fell in it and their observed +// transit range in nanoseconds. +type bucketStat struct { + count int64 + minNs int64 + maxNs int64 } -func (s forwardSummary) addSample(transitNs int64) forwardSummary { - us := transitNs / 1000 - if s.count == 0 { - return forwardSummary{count: 1, sumUs: us, sumSqUs: us * us, minNs: transitNs, maxNs: transitNs} +func (b *bucketStat) add(transitNs int64) { + if b.count == 0 { + b.minNs, b.maxNs = transitNs, transitNs + } else { + b.minNs = min(b.minNs, transitNs) + b.maxNs = max(b.maxNs, transitNs) } + b.count++ +} +func (b *bucketStat) mergeIn(o bucketStat) { + if o.count == 0 { + return + } + if b.count == 0 { + *b = o + return + } + b.count += o.count + b.minNs = min(b.minNs, o.minNs) + b.maxNs = max(b.maxNs, o.maxNs) +} + +type forwardSummary struct { + count int64 + minNs int64 + maxNs int64 + buckets [forwardHistBuckets]bucketStat +} + +// forwardBucket returns the bucket index for a transit in nanoseconds: the +// octave floor(log2(us)) times forwardHistSub, plus the linear sub-bucket +// within the octave. +func forwardBucket(transitNs int64) int { + us := transitNs / 1000 + if us <= 0 { + return 0 + } + e := bits.Len64(uint64(us)) - 1 // floor(log2(us)), us >= 1 so e >= 0 + if e >= forwardHistOctaves { + return forwardHistBuckets - 1 + } + octave := int64(1) << e + return e*forwardHistSub + int((us-octave)*forwardHistSub/octave) +} + +func (s *forwardSummary) addSample(transitNs int64) { + if s.count == 0 { + s.minNs, s.maxNs = transitNs, transitNs + } else { + s.minNs = min(s.minNs, transitNs) + s.maxNs = max(s.maxNs, transitNs) + } s.count++ - s.sumUs += us - s.sumSqUs += us * us - if transitNs < s.minNs { - s.minNs = transitNs - } - if transitNs > s.maxNs { - s.maxNs = transitNs - } - return s + s.buckets[forwardBucket(transitNs)].add(transitNs) } func (s forwardSummary) merge(o forwardSummary) forwardSummary { @@ -129,35 +189,36 @@ func (s forwardSummary) merge(o forwardSummary) forwardSummary { if s.count == 0 { return o } - return forwardSummary{ - count: s.count + o.count, - sumUs: s.sumUs + o.sumUs, - sumSqUs: s.sumSqUs + o.sumSqUs, - minNs: min(s.minNs, o.minNs), - maxNs: max(s.maxNs, o.maxNs), + s.count += o.count + s.minNs = min(s.minNs, o.minNs) + s.maxNs = max(s.maxNs, o.maxNs) + for i := range s.buckets { + s.buckets[i].mergeIn(o.buckets[i]) } + return s } -func (s forwardSummary) meanStdDev() (mean, stdDev time.Duration) { +// percentile returns the p-quantile (0..1) of transit, interpolated within the +// observed [min, max] of the bucket the quantile falls in. +func (s forwardSummary) percentile(p float64) time.Duration { if s.count == 0 { - return 0, 0 + return 0 } - - meanUs := float64(s.sumUs) / float64(s.count) - mean = time.Duration(meanUs * float64(time.Microsecond)) - if s.count < 2 { - return mean, 0 + target := p * float64(s.count) + var cum float64 + for i := range s.buckets { + b := s.buckets[i] + c := float64(b.count) + if c == 0 { + continue + } + if cum+c >= target { + lo, hi := float64(b.minNs), float64(b.maxNs) + return time.Duration(lo + (hi-lo)*(target-cum)/c) + } + cum += c } - - // sample variance (divisor count-1) - m2 := float64(s.sumSqUs) - float64(s.sumUs)*meanUs - varUs2 := m2 / float64(s.count-1) - if varUs2 < 0 { - // floating point rounding can push a (near-zero) variance slightly negative - varUs2 = 0 - } - stdDev = time.Duration(math.Sqrt(varUs2) * float64(time.Microsecond)) - return mean, stdDev + return time.Duration(s.maxNs) } type ForwardStats struct { @@ -231,12 +292,12 @@ func (s *ForwardStats) run() { // flush drains the buffered samples, observes each into the Prometheus // histogram, and folds the interval summary into the window ring used for the -// latency/jitter gauges. +// latency gauge. func (s *ForwardStats) flush() { var summ forwardSummary s.samples.drain(func(transitNs int64) { prometheus.RecordForwardLatencySample(transitNs) - summ = summ.addSample(transitNs) + summ.addSample(transitNs) }) s.lock.Lock() @@ -274,33 +335,33 @@ func (s *ForwardStats) summarize(window time.Duration) forwardSummary { return w } -// GetStats returns the mean latency and jitter (std dev) of the forwarding -// transit over the most recent duration. The duration is rounded up to a whole -// number of summary intervals (the smallest bucket span that covers it). A -// duration <= 0, or one that meets/exceeds the report window, covers the full -// window. -func (s *ForwardStats) GetStats(duration time.Duration) (time.Duration, time.Duration) { - return s.summarize(duration).meanStdDev() +// GetStats returns the reported forwarding-latency percentile +// (forwardLatencyPercentile) over the most recent duration. The duration is +// rounded up to a whole number of summary intervals (the smallest bucket span +// that covers it). A duration <= 0, or one that meets/exceeds the report window, +// covers the full window. +func (s *ForwardStats) GetStats(duration time.Duration) time.Duration { + return s.summarize(duration).percentile(forwardLatencyPercentile) } func (s *ForwardStats) report() { w := s.summarize(0) - latency, jitter := w.meanStdDev() + p90 := w.percentile(forwardLatencyPercentile) if dropped := s.samples.takeDropped(); dropped > 0 { logger.Warnw("forward stats sample buffer overflow", nil, "dropped", dropped) } - if w.count > 0 && jitter > latency*cSkewFactor { + // a max far above p90 means a few packets stalled (e.g. goroutine scheduling + // latency) rather than a broad forwarding slowdown; p90 rides through it. + if w.count > 0 && w.maxNs > p90.Nanoseconds()*cSkewFactor { logger.Infow( - "high jitter in forwarding path", + "high spread in forwarding path", "lowest", time.Duration(w.minNs), "highest", time.Duration(w.maxNs), "count", w.count, - "latency", latency, - "jitter", jitter, + "p90", p90, ) } - prometheus.RecordForwardJitter(uint32(jitter.Nanoseconds())) - prometheus.RecordForwardLatency(uint32(latency.Nanoseconds())) + prometheus.RecordForwardLatency(uint32(p90.Nanoseconds())) } diff --git a/pkg/sfu/forwardstats_test.go b/pkg/sfu/forwardstats_test.go index 152404978..fe145659a 100644 --- a/pkg/sfu/forwardstats_test.go +++ b/pkg/sfu/forwardstats_test.go @@ -24,28 +24,39 @@ func initPrometheus(t *testing.T) { // forwardSummary // --------------------------------------------------------------------------- +// bucketSum totals a summary's per-bucket counts, independent of the bucketing +// scheme, so tests can assert samples landed without hard-coding bucket indices. +func bucketSum(s forwardSummary) int64 { + var n int64 + for _, b := range s.buckets { + n += b.count + } + return n +} + func TestForwardSummary_AddSample(t *testing.T) { var s forwardSummary // empty summary require.Equal(t, int64(0), s.count) - // microsecond-aligned transits so the /1000 truncation is exact - s = s.addSample(3000) // 3us - s = s.addSample(1000) // 1us - s = s.addSample(2000) // 2us + s.addSample(3000) // 3us + s.addSample(1000) // 1us + s.addSample(2000) // 2us require.Equal(t, int64(3), s.count) - require.Equal(t, int64(1+2+3), s.sumUs) - require.Equal(t, int64(1+4+9), s.sumSqUs) require.Equal(t, int64(1000), s.minNs) require.Equal(t, int64(3000), s.maxNs) + require.Equal(t, s.count, bucketSum(s)) // every sample landed in a bucket } func TestForwardSummary_Merge(t *testing.T) { var empty forwardSummary - a := forwardSummary{}.addSample(1000).addSample(2000) - b := forwardSummary{}.addSample(5000).addSample(3000) + var a, b forwardSummary + a.addSample(1000) + a.addSample(2000) + b.addSample(5000) + b.addSample(3000) // merging with empty is identity, in both directions require.Equal(t, a, a.merge(empty)) @@ -53,34 +64,72 @@ func TestForwardSummary_Merge(t *testing.T) { m := a.merge(b) require.Equal(t, int64(4), m.count) - require.Equal(t, a.sumUs+b.sumUs, m.sumUs) - require.Equal(t, a.sumSqUs+b.sumSqUs, m.sumSqUs) require.Equal(t, int64(1000), m.minNs) require.Equal(t, int64(5000), m.maxNs) + // merge sums the per-bucket counts + require.Equal(t, bucketSum(a)+bucketSum(b), bucketSum(m)) + require.Equal(t, m.count, bucketSum(m)) } -func TestForwardSummary_MeanStdDev(t *testing.T) { +func TestForwardSummary_Percentile(t *testing.T) { // empty -> zero - mean, stdDev := forwardSummary{}.meanStdDev() - require.Zero(t, mean) - require.Zero(t, stdDev) + require.Zero(t, forwardSummary{}.percentile(0.5)) - // single sample -> mean set, stddev zero (needs >= 2 for variance) - mean, stdDev = forwardSummary{}.addSample(4000).meanStdDev() - require.Equal(t, 4*time.Microsecond, mean) - require.Zero(t, stdDev) + // half the samples fast (1us), half a slow tail (100us). A percentile is not + // dragged toward the tail the way the mean is -- why the reported metric is + // p90, not the mean. + var s forwardSummary + for i := 0; i < 5; i++ { + s.addSample(1000) // 1us + } + for i := 0; i < 5; i++ { + s.addSample(100_000) // 100us + } + require.Equal(t, int64(10), s.count) - // identical samples -> zero variance - s := forwardSummary{}.addSample(2000).addSample(2000).addSample(2000) - mean, stdDev = s.meanStdDev() - require.Equal(t, 2*time.Microsecond, mean) - require.Zero(t, stdDev) + // p50 lands in the fast cluster, unmoved by the 100us tail. + require.Equal(t, 1*time.Microsecond, s.percentile(0.5)) + // p90 crosses into the slow cluster. + require.Equal(t, 100*time.Microsecond, s.percentile(0.9)) - // known dataset [1us, 2us, 3us]: mean 2us, sample variance 1us^2 -> stddev 1us - s = forwardSummary{}.addSample(1000).addSample(2000).addSample(3000) - mean, stdDev = s.meanStdDev() - require.Equal(t, 2*time.Microsecond, mean) - require.InDelta(t, float64(time.Microsecond), float64(stdDev), float64(50*time.Nanosecond)) + // a single sample reports exactly its latency, not the top edge of its bucket + // (nominal-edge interpolation would report the bucket's upper reach instead). + var single forwardSummary + single.addSample(int64(20 * time.Millisecond)) + require.Equal(t, 20*time.Millisecond, single.percentile(0.9)) + + // identical samples share a bucket but must not invent intra-bucket spread. + var u forwardSummary + for i := 0; i < 8; i++ { + u.addSample(int64(2 * time.Millisecond)) + } + require.Equal(t, 2*time.Millisecond, u.percentile(0.5)) + require.Equal(t, 2*time.Millisecond, u.percentile(0.99)) +} + +func TestForwardSummary_ThresholdResolution(t *testing.T) { + // Nodes whose p90 packets cluster near a 300us overload threshold. + // Interpolating within each bucket's observed range reports a tight tail + // exactly, so each cluster lands on the correct side of 300us. + build := func(tailUs int64) forwardSummary { + var s forwardSummary + for i := 0; i < 850; i++ { + s.addSample(50 * int64(time.Microsecond)) + } + for i := 0; i < 150; i++ { + s.addSample(tailUs * int64(time.Microsecond)) + } + return s + } + + require.Less(t, build(265).percentile(0.9), 300*time.Microsecond) // below -> no trip + require.Greater(t, build(310).percentile(0.9), 300*time.Microsecond) // just above -> trips + require.Greater(t, build(500).percentile(0.9), 300*time.Microsecond) // well above -> trips + + // a tight tail is reported exactly, regardless of where it sits in the bucket. + require.Equal(t, 265*time.Microsecond, build(265).percentile(0.9)) + require.Equal(t, 310*time.Microsecond, build(310).percentile(0.9)) + require.Equal(t, 500*time.Microsecond, build(500).percentile(0.9)) } // --------------------------------------------------------------------------- @@ -284,29 +333,25 @@ func TestForwardStats_GetStats(t *testing.T) { // 5 buckets, each covering one 100ms summary interval. s := &ForwardStats{ring: make([]forwardSummary, 5), summaryInterval: 100 * time.Millisecond} - // fold five 100ms buckets, one sample each: 1ms, 2ms, 3ms, 4ms, 5ms. - for i := 1; i <= 5; i++ { + // fold five 100ms buckets, one sample each, descending 5ms..1ms so trailing + // windows have distinct maxima. + for i := 5; i >= 1; i-- { s.Update(0, int64(i)*int64(time.Millisecond)) s.flush() } require.Equal(t, 5, s.ringLen) - // a duration <= 0 covers the whole window: mean of 1..5ms == 3ms. - latency, jitter := s.GetStats(0) - require.InDelta(t, float64(3*time.Millisecond), float64(latency), float64(50*time.Microsecond)) - require.Greater(t, jitter, time.Duration(0)) + // whole window {1..5ms}: p90 clamps to the observed max, 5ms. + require.Equal(t, 5*time.Millisecond, s.GetStats(0)) + require.Equal(t, 5*time.Millisecond, s.GetStats(time.Second)) - // a duration meeting/exceeding the window also covers it. - fullLatency, _ := s.GetStats(time.Second) - require.InDelta(t, float64(3*time.Millisecond), float64(fullLatency), float64(50*time.Microsecond)) + // ~200ms covers only the two most recent buckets {2ms, 1ms}: a lower window + // than the full one, and above the single most-recent bucket. + require.Less(t, s.GetStats(200*time.Millisecond), 3*time.Millisecond) + require.Greater(t, s.GetStats(200*time.Millisecond), s.GetStats(time.Nanosecond)) - // ~200ms rounds up to the two most recent buckets (4ms, 5ms): mean == 4.5ms. - shortLatency, _ := s.GetStats(200 * time.Millisecond) - require.InDelta(t, float64(4500*time.Microsecond), float64(shortLatency), float64(50*time.Microsecond)) - - // a sub-interval duration still yields at least the most recent bucket (5ms). - lastLatency, _ := s.GetStats(time.Nanosecond) - require.InDelta(t, float64(5*time.Millisecond), float64(lastLatency), float64(50*time.Microsecond)) + // a sub-interval yields only the most recent bucket, ~1ms. + require.Equal(t, 1*time.Millisecond, s.GetStats(time.Nanosecond)) } func TestForwardStats_Lifecycle(t *testing.T) { diff --git a/pkg/telemetry/prometheus/node.go b/pkg/telemetry/prometheus/node.go index fb540fb30..915b0322a 100644 --- a/pkg/telemetry/prometheus/node.go +++ b/pkg/telemetry/prometheus/node.go @@ -174,7 +174,6 @@ func GetNodeStats(nodeStartedAt int64, prevStats []*livekit.NodeStats, rateInter ParticipantRtcCanceled: participantRTCCanceled.Load(), ParticipantRtcActive: participantRTCActive.Load(), ForwardLatency: forwardLatency.Load(), - ForwardJitter: forwardJitter.Load(), NumCpus: uint32(cpuStats.NumCPU()), // this will round down to the nearest integer CpuLoad: float32(cpuStats.GetCPULoad()), MemoryTotal: memTotal, diff --git a/pkg/telemetry/prometheus/packets.go b/pkg/telemetry/prometheus/packets.go index 5b8e45659..8736d4de6 100644 --- a/pkg/telemetry/prometheus/packets.go +++ b/pkg/telemetry/prometheus/packets.go @@ -51,7 +51,6 @@ var ( participantRTCCanceled atomic.Uint64 participantRTCActive atomic.Uint64 forwardLatency atomic.Uint32 - forwardJitter atomic.Uint32 promPacketLabels = []string{"direction", "transmission", "country"} promPacketTotal *prometheus.CounterVec @@ -70,7 +69,6 @@ var ( promParticipantJoin *prometheus.CounterVec promConnections *prometheus.GaugeVec promForwardLatency prometheus.Gauge - promForwardJitter prometheus.Gauge promForwardLatencyHist prometheus.Histogram ) @@ -164,12 +162,6 @@ func initPacketStats(nodeID string, nodeType livekit.NodeType) { Name: "latency", ConstLabels: prometheus.Labels{"node_id": nodeID, "node_type": nodeType.String()}, }) - promForwardJitter = prometheus.NewGauge(prometheus.GaugeOpts{ - Namespace: livekitNamespace, - Subsystem: "forward", - Name: "jitter", - ConstLabels: prometheus.Labels{"node_id": nodeID, "node_type": nodeType.String()}, - }) promForwardLatencyHist = prometheus.NewHistogram(prometheus.HistogramOpts{ Namespace: livekitNamespace, Subsystem: "forward_latency", @@ -204,7 +196,6 @@ func initPacketStats(nodeID string, nodeType livekit.NodeType) { prometheus.MustRegister(promParticipantJoin) prometheus.MustRegister(promConnections) prometheus.MustRegister(promForwardLatency) - prometheus.MustRegister(promForwardJitter) prometheus.MustRegister(promForwardLatencyHist) } @@ -380,12 +371,7 @@ func RecordForwardLatencySample(forwardLatency int64) { promForwardLatencyHist.Observe(float64(forwardLatency)) } -func RecordForwardLatency(longTermLatencyAvg uint32) { - forwardLatency.Store(longTermLatencyAvg) - promForwardLatency.Set(float64(longTermLatencyAvg)) -} - -func RecordForwardJitter(longTermJitterAvg uint32) { - forwardJitter.Store(longTermJitterAvg) - promForwardJitter.Set(float64(longTermJitterAvg)) +func RecordForwardLatency(p90 uint32) { + forwardLatency.Store(p90) + promForwardLatency.Set(float64(p90)) }