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 <noreply@anthropic.com>

* 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 <noreply@anthropic.com>

* 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 <noreply@anthropic.com>

* 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 <noreply@anthropic.com>

* 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 <noreply@anthropic.com>

* 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 <noreply@anthropic.com>

* 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 <noreply@anthropic.com>

* 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 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
Raja Subramanian
2026-09-29 01:39:05 +05:30
committed by GitHub
co-authored by Claude Opus 4.8
parent 9a6c187afd
commit 6ed445bad1
4 changed files with 214 additions and 123 deletions
+122 -61
View File
@@ -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()))
}
+89 -44
View File
@@ -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) {
-1
View File
@@ -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,
+3 -17
View File
@@ -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))
}