update RTPstats

This commit is contained in:
David Chen
2026-09-24 17:02:04 -07:00
parent c2f88c3aaf
commit 874f8e49f6
9 changed files with 87 additions and 88 deletions
+3 -6
View File
@@ -21,8 +21,8 @@ require (
github.com/jxskiss/base62 v1.1.0
github.com/livekit/mageutil v0.0.0-20250511045019-0f1ff63f7731
github.com/livekit/mediatransportutil v0.0.0-20260821083140-f234b534b095
github.com/livekit/protocol v1.52.0
github.com/livekit/psrpc v0.7.7
github.com/livekit/protocol v1.52.1
github.com/livekit/psrpc v0.8.0
github.com/mackerelio/go-osstat v0.2.8
github.com/magefile/mage v1.17.2
github.com/mitchellh/go-homedir v1.1.0
@@ -128,9 +128,6 @@ require (
github.com/mdlayher/socket v0.6.1 // indirect
github.com/moby/docker-image-spec v1.3.1 // indirect
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
github.com/nats-io/nats.go v1.53.1 // indirect
github.com/nats-io/nkeys v0.4.16 // indirect
github.com/nats-io/nuid v1.0.1 // indirect
github.com/opencontainers/go-digest v1.0.0 // indirect
github.com/opencontainers/image-spec v1.1.1 // indirect
github.com/pion/logging v0.2.4
@@ -139,7 +136,7 @@ require (
github.com/pion/srtp/v3 v3.0.13 // indirect
github.com/pion/stun/v3 v3.1.7
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect
github.com/prometheus/client_model v0.6.2 // indirect
github.com/prometheus/client_model v0.6.3 // indirect
github.com/prometheus/common v0.70.1 // indirect
github.com/prometheus/procfs v0.21.1 // indirect
github.com/urfave/cli/v3 v3.10.1
+6 -6
View File
@@ -164,10 +164,10 @@ github.com/livekit/mageutil v0.0.0-20250511045019-0f1ff63f7731 h1:9x+U2HGLrSw5AT
github.com/livekit/mageutil v0.0.0-20250511045019-0f1ff63f7731/go.mod h1:Rs3MhFwutWhGwmY1VQsygw28z5bWcnEYmS1OG9OxjOQ=
github.com/livekit/mediatransportutil v0.0.0-20260821083140-f234b534b095 h1:BcliKAXoMhl/nWmzQweQ5kmh4Qqagxl4s3Z5pvM/7AY=
github.com/livekit/mediatransportutil v0.0.0-20260821083140-f234b534b095/go.mod h1:o8CFmAdrVwzJNOCsQCLUzXRjokkufNshnQHOe4fRaqU=
github.com/livekit/protocol v1.52.0 h1:K7XWd1cJOekMCigjkyBnpmM841J8TX/YWXYWBMK9wYc=
github.com/livekit/protocol v1.52.0/go.mod h1:YUI8xd1bvdEWtUWwgvsIkUFb4Ll98ngsi7ifLdT5HV0=
github.com/livekit/psrpc v0.7.7 h1:eZ/jYlayQ3Y+C/3+NPtwTJBXZo/aN4vJf2MmBK6AfT4=
github.com/livekit/psrpc v0.7.7/go.mod h1:Twno03W8gpTNRpmc2cQwSoxYEo/nJxf+65iMq+9IuHo=
github.com/livekit/protocol v1.52.1 h1:w8wRjRlj/57AJBxwomEzpEMzf+lT0qOoLQLO7awdTMU=
github.com/livekit/protocol v1.52.1/go.mod h1:t5CV/LCsWbI9zsSTHTJKDyLtByqvvVkioSP7jgqpIGA=
github.com/livekit/psrpc v0.8.0 h1:uo3x+DL8pK2OIaCxg/s9sLUO2cIocS5HwdPhhI+m+fQ=
github.com/livekit/psrpc v0.8.0/go.mod h1:6KnU8vULoizX0Er9NNfeKJDV8ug87j3fuo19M0bhj6Q=
github.com/livekit/webrtc-pion/v4 v4.2.18-warp.1 h1:fH+v4W+NFp9FfPzON6FaUFNmazGcctaAhb2P+Ksf+1s=
github.com/livekit/webrtc-pion/v4 v4.2.18-warp.1/go.mod h1:rbKGHo2OpNUImWTvRIV776/3xjjq/t47H3IZiTtwluc=
github.com/mackerelio/go-osstat v0.2.8 h1:I2duicTaCGWoM53XwAwA9OIe1inu0xnVs8/pqOWWVr4=
@@ -272,8 +272,8 @@ github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRI
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/prometheus/client_golang v1.24.1 h1:JnJkREXzWxUdCuPFpIWZiPispT9xVV59uiuyR2bPlnU=
github.com/prometheus/client_golang v1.24.1/go.mod h1:F+oSRECHg4sse5ucfYpYDeIv/hu68Zo0uoHKetWnzcE=
github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk=
github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE=
github.com/prometheus/client_model v0.6.3 h1:O0jaTVAYNxTHYInEPFJt5I3+sN8zqBtVMPTB1qyxiEo=
github.com/prometheus/client_model v0.6.3/go.mod h1:gpN5P9S7Rr6Yr92PiQ+Ixvhf6JZEkF1dnxsYL2aPBEM=
github.com/prometheus/common v0.70.1 h1:1HvjP4D5oL3t8RsPlwxA9onvvStjtIHYE5XuuwOi/PY=
github.com/prometheus/common v0.70.1/go.mod h1:VdFUQDMZK3VLkurFUVhia6uys/0suUp86TJz5qbJRhc=
github.com/prometheus/procfs v0.21.1 h1:GljZCt+zSTS+NZq88cyQ1LjZ+RCHp3uVuabBWA5+OJI=
+2 -2
View File
@@ -18,7 +18,7 @@ import (
"context"
"github.com/dennwc/iters"
"github.com/livekit/livekit-server/pkg/service"
"github.com/livekit/psrpc"
"github.com/livekit/psrpc/pkg/bus/redisbus"
"slices"
"testing"
@@ -28,7 +28,7 @@ import (
func ioStoreDocker(t testing.TB) (*service.IOInfoService, *service.RedisStore) {
r := redisClientDocker(t)
bus := psrpc.NewRedisMessageBus(r)
bus := redisbus.New(r)
rs := service.NewRedisStore(r)
io, err := service.NewIOInfoService(bus, rs, rs, rs, nil)
require.NoError(t, err)
+2 -1
View File
@@ -35,6 +35,7 @@ import (
"github.com/livekit/protocol/utils"
"github.com/livekit/protocol/webhook"
"github.com/livekit/psrpc"
"github.com/livekit/psrpc/pkg/bus/redisbus"
"github.com/livekit/psrpc/pkg/middleware/otelpsrpc"
"github.com/livekit/livekit-server/pkg/agent"
@@ -197,7 +198,7 @@ func getMessageBus(rc redis.UniversalClient) psrpc.MessageBus {
if rc == nil {
return psrpc.NewLocalMessageBus()
}
return psrpc.NewRedisMessageBus(rc)
return redisbus.New(rc)
}
func getEgressStore(s ObjectStore) EgressStore {
+2 -1
View File
@@ -21,6 +21,7 @@ import (
"github.com/livekit/protocol/utils"
"github.com/livekit/protocol/webhook"
"github.com/livekit/psrpc"
"github.com/livekit/psrpc/pkg/bus/redisbus"
"github.com/livekit/psrpc/pkg/middleware/otelpsrpc"
"github.com/pion/turn/v5"
"github.com/pkg/errors"
@@ -258,7 +259,7 @@ func getMessageBus(rc redis.UniversalClient) psrpc.MessageBus {
if rc == nil {
return psrpc.NewLocalMessageBus()
}
return psrpc.NewRedisMessageBus(rc)
return redisbus.New(rc)
}
func getEgressStore(s ObjectStore) EgressStore {
+9 -9
View File
@@ -581,16 +581,16 @@ func (b *Buffer) feedFECLocked(
recovered := b.fecDecoder.DecodeFEC(pkt)
statsAfter := b.fecDecoder.Stats()
fecPacketsReceived := statsAfter.FECPacketsReceived - statsBefore.FECPacketsReceived
fecBytesReceived := statsAfter.FECBytesReceived - statsBefore.FECBytesReceived
fecPackets := statsAfter.FECPacketsReceived - statsBefore.FECPacketsReceived
fecBytes := statsAfter.FECBytesReceived - statsBefore.FECBytesReceived
fecPacketsDiscarded := statsAfter.FECPacketsDiscarded - statsBefore.FECPacketsDiscarded
packetsRecovered := statsAfter.PacketsRecovered - statsBefore.PacketsRecovered
fecPacketsRecovered := statsAfter.PacketsRecovered - statsBefore.PacketsRecovered
if b.rtpStats != nil {
b.rtpStats.UpdateFEC(
fecPacketsReceived,
fecBytesReceived,
fecPackets,
fecBytes,
fecPacketsDiscarded,
packetsRecovered,
fecPacketsRecovered,
)
}
@@ -613,10 +613,10 @@ func (b *Buffer) feedFECLocked(
if cb := b.onFECRecovery; cb != nil {
return fecRecoveryDelta{
received: int(fecPacketsReceived),
recovered: int(packetsRecovered),
received: int(fecPackets),
recovered: int(fecPacketsRecovered),
discarded: int(fecPacketsDiscarded),
bytesReceived: int(fecBytesReceived),
bytesReceived: int(fecBytes),
}, cb
}
+3 -3
View File
@@ -194,10 +194,10 @@ func TestBufferFECRecoversDroppedPacket(t *testing.T) {
}
rtpStats := primary.GetStats()
require.NotNil(t, rtpStats)
assert.EqualValues(t, len(fecPackets), rtpStats.FecPacketsReceived)
assert.EqualValues(t, fecPayloadBytes, rtpStats.FecBytesReceived)
assert.EqualValues(t, len(fecPackets), rtpStats.FecPackets)
assert.EqualValues(t, fecPayloadBytes, rtpStats.FecBytes)
assert.EqualValues(t, 0, rtpStats.FecPacketsDiscarded)
assert.EqualValues(t, 1, rtpStats.PacketsRecovered)
assert.EqualValues(t, 1, rtpStats.FecPacketsRecovered)
// the 9 received packets flow through the ext packet pipeline, the
// recovered one fills the bucket like an RTX repair
+48 -48
View File
@@ -47,10 +47,10 @@ type RTPDeltaInfo struct {
PacketsPadding uint32
BytesPadding uint64
HeaderBytesPadding uint64
FecPacketsReceived uint32
FecBytesReceived uint64
FecPackets uint32
FecBytes uint64
FecPacketsDiscarded uint32
PacketsRecovered uint32
FecPacketsRecovered uint32
PacketsLost uint32
PacketsMissing uint32
PacketsOutOfOrder uint32
@@ -79,10 +79,10 @@ func (r *RTPDeltaInfo) MarshalLogObject(e zapcore.ObjectEncoder) error {
e.AddUint32("PacketsPadding", r.PacketsPadding)
e.AddUint64("BytesPadding", r.BytesPadding)
e.AddUint64("HeaderBytesPadding", r.HeaderBytesPadding)
e.AddUint32("FecPacketsReceived", r.FecPacketsReceived)
e.AddUint64("FecBytesReceived", r.FecBytesReceived)
e.AddUint32("FecPackets", r.FecPackets)
e.AddUint64("FecBytes", r.FecBytes)
e.AddUint32("FecPacketsDiscarded", r.FecPacketsDiscarded)
e.AddUint32("PacketsRecovered", r.PacketsRecovered)
e.AddUint32("FecPacketsRecovered", r.FecPacketsRecovered)
e.AddUint32("PacketsLost", r.PacketsLost)
e.AddUint32("PacketsMissing", r.PacketsMissing)
e.AddUint32("PacketsOutOfOrder", r.PacketsOutOfOrder)
@@ -111,10 +111,10 @@ type snapshot struct {
bytesPadding uint64
headerBytesPadding uint64
fecPacketsReceived uint64
fecBytesReceived uint64
fecPackets uint64
fecBytes uint64
fecPacketsDiscarded uint64
packetsRecovered uint64
fecPacketsRecovered uint64
frames uint32
@@ -138,10 +138,10 @@ func (s *snapshot) MarshalLogObject(e zapcore.ObjectEncoder) error {
e.AddUint64("packetsPadding", s.packetsPadding)
e.AddUint64("bytesPadding", s.bytesPadding)
e.AddUint64("headerBytesPadding", s.headerBytesPadding)
e.AddUint64("fecPacketsReceived", s.fecPacketsReceived)
e.AddUint64("fecBytesReceived", s.fecBytesReceived)
e.AddUint64("fecPackets", s.fecPackets)
e.AddUint64("fecBytes", s.fecBytes)
e.AddUint64("fecPacketsDiscarded", s.fecPacketsDiscarded)
e.AddUint64("packetsRecovered", s.packetsRecovered)
e.AddUint64("fecPacketsRecovered", s.fecPacketsRecovered)
e.AddUint32("frames", s.frames)
e.AddUint32("plis", s.plis)
e.AddUint32("firs", s.firs)
@@ -242,10 +242,10 @@ type rtpStatsBase struct {
bytesPadding uint64
headerBytesPadding uint64
fecPacketsReceived uint64
fecBytesReceived uint64
fecPackets uint64
fecBytes uint64
fecPacketsDiscarded uint64
packetsRecovered uint64
fecPacketsRecovered uint64
frames uint32
@@ -310,10 +310,10 @@ func (r *rtpStatsBase) seed(from *rtpStatsBase) bool {
r.bytesPadding = from.bytesPadding
r.headerBytesPadding = from.headerBytesPadding
r.fecPacketsReceived = from.fecPacketsReceived
r.fecBytesReceived = from.fecBytesReceived
r.fecPackets = from.fecPackets
r.fecBytes = from.fecBytes
r.fecPacketsDiscarded = from.fecPacketsDiscarded
r.packetsRecovered = from.packetsRecovered
r.fecPacketsRecovered = from.fecPacketsRecovered
r.frames = from.frames
@@ -355,10 +355,10 @@ func (r *rtpStatsBase) newSnapshotID(extStartSN uint64) uint32 {
}
func (r *rtpStatsBase) UpdateFEC(
fecPacketsReceived uint64,
fecBytesReceived uint64,
fecPackets uint64,
fecBytes uint64,
fecPacketsDiscarded uint64,
packetsRecovered uint64,
fecPacketsRecovered uint64,
) {
r.lock.Lock()
defer r.lock.Unlock()
@@ -367,10 +367,10 @@ func (r *rtpStatsBase) UpdateFEC(
return
}
r.fecPacketsReceived += fecPacketsReceived
r.fecBytesReceived += fecBytesReceived
r.fecPackets += fecPackets
r.fecBytes += fecBytes
r.fecPacketsDiscarded += fecPacketsDiscarded
r.packetsRecovered += packetsRecovered
r.fecPacketsRecovered += fecPacketsRecovered
}
func (r *rtpStatsBase) UpdateFir(firCount uint32) {
@@ -554,10 +554,10 @@ func (r *rtpStatsBase) deltaInfo(
deltaInfo = &RTPDeltaInfo{
StartTime: time.Unix(0, startTime),
EndTime: time.Unix(0, endTime),
FecPacketsReceived: uint32(now.fecPacketsReceived - then.fecPacketsReceived),
FecBytesReceived: now.fecBytesReceived - then.fecBytesReceived,
FecPackets: uint32(now.fecPackets - then.fecPackets),
FecBytes: now.fecBytes - then.fecBytes,
FecPacketsDiscarded: uint32(now.fecPacketsDiscarded - then.fecPacketsDiscarded),
PacketsRecovered: uint32(now.packetsRecovered - then.packetsRecovered),
FecPacketsRecovered: uint32(now.fecPacketsRecovered - then.fecPacketsRecovered),
}
return
}
@@ -597,10 +597,10 @@ func (r *rtpStatsBase) deltaInfo(
PacketsPadding: uint32(packetsPadding),
BytesPadding: now.bytesPadding - then.bytesPadding,
HeaderBytesPadding: now.headerBytesPadding - then.headerBytesPadding,
FecPacketsReceived: uint32(now.fecPacketsReceived - then.fecPacketsReceived),
FecBytesReceived: now.fecBytesReceived - then.fecBytesReceived,
FecPackets: uint32(now.fecPackets - then.fecPackets),
FecBytes: now.fecBytes - then.fecBytes,
FecPacketsDiscarded: uint32(now.fecPacketsDiscarded - then.fecPacketsDiscarded),
PacketsRecovered: uint32(now.packetsRecovered - then.packetsRecovered),
FecPacketsRecovered: uint32(now.fecPacketsRecovered - then.fecPacketsRecovered),
PacketsLost: packetsLost,
PacketsOutOfOrder: uint32(now.packetsOutOfOrder - then.packetsOutOfOrder),
Frames: now.frames - then.frames,
@@ -645,10 +645,10 @@ func (r *rtpStatsBase) marshalLogObject(
e.AddFloat64("bitratePadding", float64(r.bytesPadding)*8.0/elapsedSeconds)
e.AddUint64("headerBytesPadding", r.headerBytesPadding)
e.AddUint64("fecPacketsReceived", r.fecPacketsReceived)
e.AddUint64("fecBytesReceived", r.fecBytesReceived)
e.AddUint64("fecPackets", r.fecPackets)
e.AddUint64("fecBytes", r.fecBytes)
e.AddUint64("fecPacketsDiscarded", r.fecPacketsDiscarded)
e.AddUint64("packetsRecovered", r.packetsRecovered)
e.AddUint64("fecPacketsRecovered", r.fecPacketsRecovered)
e.AddUint32("frames", r.frames)
e.AddFloat64("frameRate", float64(r.frames)/elapsedSeconds)
@@ -700,10 +700,10 @@ func (r *rtpStatsBase) toProto(
p.BitratePadding = float64(r.bytesPadding) * 8.0 / p.Duration
p.HeaderBytesPadding = r.headerBytesPadding
p.FecPacketsReceived = uint32(r.fecPacketsReceived)
p.FecBytesReceived = r.fecBytesReceived
p.FecPackets = uint32(r.fecPackets)
p.FecBytes = r.fecBytes
p.FecPacketsDiscarded = uint32(r.fecPacketsDiscarded)
p.PacketsRecovered = uint32(r.packetsRecovered)
p.FecPacketsRecovered = uint32(r.fecPacketsRecovered)
p.Frames = r.frames
p.FrameRate = float64(r.frames) / p.Duration
@@ -880,10 +880,10 @@ func (r *rtpStatsBase) getSnapshot(startTime int64, extStartSN uint64) snapshot
packetsPadding: r.packetsPadding,
bytesPadding: r.bytesPadding,
headerBytesPadding: r.headerBytesPadding,
fecPacketsReceived: r.fecPacketsReceived,
fecBytesReceived: r.fecBytesReceived,
fecPackets: r.fecPackets,
fecBytes: r.fecBytes,
fecPacketsDiscarded: r.fecPacketsDiscarded,
packetsRecovered: r.packetsRecovered,
fecPacketsRecovered: r.fecPacketsRecovered,
frames: r.frames,
plis: r.plis,
firs: r.firs,
@@ -924,10 +924,10 @@ func AggregateRTPDeltaInfo(deltaInfoList []*RTPDeltaInfo) *RTPDeltaInfo {
bytesPadding := uint64(0)
headerBytesPadding := uint64(0)
fecPacketsReceived := uint32(0)
fecBytesReceived := uint64(0)
fecPackets := uint32(0)
fecBytes := uint64(0)
fecPacketsDiscarded := uint32(0)
packetsRecovered := uint32(0)
fecPacketsRecovered := uint32(0)
packetsLost := uint32(0)
packetsMissing := uint32(0)
@@ -967,10 +967,10 @@ func AggregateRTPDeltaInfo(deltaInfoList []*RTPDeltaInfo) *RTPDeltaInfo {
bytesPadding += deltaInfo.BytesPadding
headerBytesPadding += deltaInfo.HeaderBytesPadding
fecPacketsReceived += deltaInfo.FecPacketsReceived
fecBytesReceived += deltaInfo.FecBytesReceived
fecPackets += deltaInfo.FecPackets
fecBytes += deltaInfo.FecBytes
fecPacketsDiscarded += deltaInfo.FecPacketsDiscarded
packetsRecovered += deltaInfo.PacketsRecovered
fecPacketsRecovered += deltaInfo.FecPacketsRecovered
packetsLost += deltaInfo.PacketsLost
packetsMissing += deltaInfo.PacketsMissing
@@ -1006,10 +1006,10 @@ func AggregateRTPDeltaInfo(deltaInfoList []*RTPDeltaInfo) *RTPDeltaInfo {
PacketsPadding: packetsPadding,
BytesPadding: bytesPadding,
HeaderBytesPadding: headerBytesPadding,
FecPacketsReceived: fecPacketsReceived,
FecBytesReceived: fecBytesReceived,
FecPackets: fecPackets,
FecBytes: fecBytes,
FecPacketsDiscarded: fecPacketsDiscarded,
PacketsRecovered: packetsRecovered,
FecPacketsRecovered: fecPacketsRecovered,
PacketsLost: packetsLost,
PacketsMissing: packetsMissing,
PacketsOutOfOrder: packetsOutOfOrder,
+12 -12
View File
@@ -241,32 +241,32 @@ func Test_RTPStatsReceiver_UpdateFEC(t *testing.T) {
r.UpdateFEC(3, 900, 1, 2)
stats := r.ToProto()
require.EqualValues(t, 3, stats.FecPacketsReceived)
require.EqualValues(t, 900, stats.FecBytesReceived)
require.EqualValues(t, 3, stats.FecPackets)
require.EqualValues(t, 900, stats.FecBytes)
require.EqualValues(t, 1, stats.FecPacketsDiscarded)
require.EqualValues(t, 2, stats.PacketsRecovered)
require.EqualValues(t, 2, stats.FecPacketsRecovered)
delta := r.DeltaInfo(snapshotID)
require.EqualValues(t, 3, delta.FecPacketsReceived)
require.EqualValues(t, 900, delta.FecBytesReceived)
require.EqualValues(t, 3, delta.FecPackets)
require.EqualValues(t, 900, delta.FecBytes)
require.EqualValues(t, 1, delta.FecPacketsDiscarded)
require.EqualValues(t, 2, delta.PacketsRecovered)
require.EqualValues(t, 2, delta.FecPacketsRecovered)
// FEC-only intervals are reported even if no primary packet arrived.
r.UpdateFEC(2, 500, 1, 1)
delta = r.DeltaInfo(snapshotID)
require.EqualValues(t, 2, delta.FecPacketsReceived)
require.EqualValues(t, 500, delta.FecBytesReceived)
require.EqualValues(t, 2, delta.FecPackets)
require.EqualValues(t, 500, delta.FecBytes)
require.EqualValues(t, 1, delta.FecPacketsDiscarded)
require.EqualValues(t, 1, delta.PacketsRecovered)
require.EqualValues(t, 1, delta.FecPacketsRecovered)
r.Stop()
r.UpdateFEC(1, 100, 1, 1)
stats = r.ToProto()
require.EqualValues(t, 5, stats.FecPacketsReceived)
require.EqualValues(t, 1400, stats.FecBytesReceived)
require.EqualValues(t, 5, stats.FecPackets)
require.EqualValues(t, 1400, stats.FecBytes)
require.EqualValues(t, 2, stats.FecPacketsDiscarded)
require.EqualValues(t, 3, stats.PacketsRecovered)
require.EqualValues(t, 3, stats.FecPacketsRecovered)
}
func Test_RTPStatsReceiver_Restart(t *testing.T) {