diff --git a/pkg/sfu/buffer/buffer.go b/pkg/sfu/buffer/buffer.go index 5ff0addfb..ddaad261c 100644 --- a/pkg/sfu/buffer/buffer.go +++ b/pkg/sfu/buffer/buffer.go @@ -583,6 +583,20 @@ func (b *Buffer) feedFECLocked( ) (fecRecoveryDelta, func(received int, recovered int, discarded int, bytesReceived int)) { statsBefore := b.fecDecoder.Stats() recovered := b.fecDecoder.DecodeFEC(pkt) + statsAfter := b.fecDecoder.Stats() + + fecPacketsReceived := statsAfter.FECPacketsReceived - statsBefore.FECPacketsReceived + fecBytesReceived := statsAfter.FECBytesReceived - statsBefore.FECBytesReceived + fecPacketsDiscarded := statsAfter.FECPacketsDiscarded - statsBefore.FECPacketsDiscarded + packetsRecovered := statsAfter.PacketsRecovered - statsBefore.PacketsRecovered + if b.rtpStats != nil { + b.rtpStats.UpdateFEC( + fecPacketsReceived, + fecBytesReceived, + fecPacketsDiscarded, + packetsRecovered, + ) + } if len(recovered) > 0 && b.fecPktBuf == nil { b.fecPktBuf = make([]byte, bucket.RTPMaxPktSize) @@ -602,12 +616,11 @@ func (b *Buffer) feedFECLocked( } if cb := b.onFECRecovery; cb != nil { - statsAfter := b.fecDecoder.Stats() return fecRecoveryDelta{ - received: int(statsAfter.FECPacketsReceived - statsBefore.FECPacketsReceived), - recovered: len(recovered), - discarded: int(statsAfter.FECPacketsDiscarded - statsBefore.FECPacketsDiscarded), - bytesReceived: int(statsAfter.FECBytesReceived - statsBefore.FECBytesReceived), + received: int(fecPacketsReceived), + recovered: int(packetsRecovered), + discarded: int(fecPacketsDiscarded), + bytesReceived: int(fecBytesReceived), }, cb } diff --git a/pkg/sfu/buffer/buffer_fec_test.go b/pkg/sfu/buffer/buffer_fec_test.go index 876e58dbb..8d859b4dd 100644 --- a/pkg/sfu/buffer/buffer_fec_test.go +++ b/pkg/sfu/buffer/buffer_fec_test.go @@ -216,6 +216,17 @@ func TestBufferFECRecoversDroppedPacket(t *testing.T) { assert.Equal(t, 1, recoveredDelta) assert.Equal(t, len(fecPackets), receivedDelta) + fecPayloadBytes := 0 + for i := range fecPackets { + fecPayloadBytes += len(fecPackets[i].Payload) + } + rtpStats := primary.GetStats() + require.NotNil(t, rtpStats) + assert.EqualValues(t, len(fecPackets), rtpStats.FecPacketsReceived) + assert.EqualValues(t, fecPayloadBytes, rtpStats.FecBytesReceived) + assert.EqualValues(t, 0, rtpStats.FecPacketsDiscarded) + assert.EqualValues(t, 1, rtpStats.PacketsRecovered) + // the 9 received packets flow through the ext packet pipeline, the // recovered one fills the bucket like an RTX repair extSNBySN := readExtSequenceNumbers(t, primary, len(media)-1) diff --git a/pkg/sfu/rtpstats/rtpstats_base.go b/pkg/sfu/rtpstats/rtpstats_base.go index b14a92b8b..05253d135 100644 --- a/pkg/sfu/rtpstats/rtpstats_base.go +++ b/pkg/sfu/rtpstats/rtpstats_base.go @@ -47,6 +47,10 @@ type RTPDeltaInfo struct { PacketsPadding uint32 BytesPadding uint64 HeaderBytesPadding uint64 + FecPacketsReceived uint32 + FecBytesReceived uint64 + FecPacketsDiscarded uint32 + PacketsRecovered uint32 PacketsLost uint32 PacketsMissing uint32 PacketsOutOfOrder uint32 @@ -75,6 +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("FecPacketsDiscarded", r.FecPacketsDiscarded) + e.AddUint32("PacketsRecovered", r.PacketsRecovered) e.AddUint32("PacketsLost", r.PacketsLost) e.AddUint32("PacketsMissing", r.PacketsMissing) e.AddUint32("PacketsOutOfOrder", r.PacketsOutOfOrder) @@ -103,6 +111,11 @@ type snapshot struct { bytesPadding uint64 headerBytesPadding uint64 + fecPacketsReceived uint64 + fecBytesReceived uint64 + fecPacketsDiscarded uint64 + packetsRecovered uint64 + frames uint32 plis uint32 @@ -125,6 +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("fecPacketsDiscarded", s.fecPacketsDiscarded) + e.AddUint64("packetsRecovered", s.packetsRecovered) e.AddUint32("frames", s.frames) e.AddUint32("plis", s.plis) e.AddUint32("firs", s.firs) @@ -225,6 +242,11 @@ type rtpStatsBase struct { bytesPadding uint64 headerBytesPadding uint64 + fecPacketsReceived uint64 + fecBytesReceived uint64 + fecPacketsDiscarded uint64 + packetsRecovered uint64 + frames uint32 jitter float64 @@ -288,6 +310,11 @@ func (r *rtpStatsBase) seed(from *rtpStatsBase) bool { r.bytesPadding = from.bytesPadding r.headerBytesPadding = from.headerBytesPadding + r.fecPacketsReceived = from.fecPacketsReceived + r.fecBytesReceived = from.fecBytesReceived + r.fecPacketsDiscarded = from.fecPacketsDiscarded + r.packetsRecovered = from.packetsRecovered + r.frames = from.frames r.jitter = from.jitter @@ -327,6 +354,25 @@ func (r *rtpStatsBase) newSnapshotID(extStartSN uint64) uint32 { return id } +func (r *rtpStatsBase) UpdateFEC( + fecPacketsReceived uint64, + fecBytesReceived uint64, + fecPacketsDiscarded uint64, + packetsRecovered uint64, +) { + r.lock.Lock() + defer r.lock.Unlock() + + if r.endTime != 0 { + return + } + + r.fecPacketsReceived += fecPacketsReceived + r.fecBytesReceived += fecBytesReceived + r.fecPacketsDiscarded += fecPacketsDiscarded + r.packetsRecovered += packetsRecovered +} + func (r *rtpStatsBase) UpdateFir(firCount uint32) { r.lock.Lock() defer r.lock.Unlock() @@ -506,8 +552,12 @@ func (r *rtpStatsBase) deltaInfo( } if packetsExpected == 0 { deltaInfo = &RTPDeltaInfo{ - StartTime: time.Unix(0, startTime), - EndTime: time.Unix(0, endTime), + StartTime: time.Unix(0, startTime), + EndTime: time.Unix(0, endTime), + FecPacketsReceived: uint32(now.fecPacketsReceived - then.fecPacketsReceived), + FecBytesReceived: now.fecBytesReceived - then.fecBytesReceived, + FecPacketsDiscarded: uint32(now.fecPacketsDiscarded - then.fecPacketsDiscarded), + PacketsRecovered: uint32(now.packetsRecovered - then.packetsRecovered), } return } @@ -547,6 +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, + FecPacketsDiscarded: uint32(now.fecPacketsDiscarded - then.fecPacketsDiscarded), + PacketsRecovered: uint32(now.packetsRecovered - then.packetsRecovered), PacketsLost: packetsLost, PacketsOutOfOrder: uint32(now.packetsOutOfOrder - then.packetsOutOfOrder), Frames: now.frames - then.frames, @@ -591,6 +645,11 @@ 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("fecPacketsDiscarded", r.fecPacketsDiscarded) + e.AddUint64("packetsRecovered", r.packetsRecovered) + e.AddUint32("frames", r.frames) e.AddFloat64("frameRate", float64(r.frames)/elapsedSeconds) @@ -641,6 +700,11 @@ 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.FecPacketsDiscarded = uint32(r.fecPacketsDiscarded) + p.PacketsRecovered = uint32(r.packetsRecovered) + p.Frames = r.frames p.FrameRate = float64(r.frames) / p.Duration @@ -816,6 +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, + fecPacketsDiscarded: r.fecPacketsDiscarded, + packetsRecovered: r.packetsRecovered, frames: r.frames, plis: r.plis, firs: r.firs, @@ -856,6 +924,11 @@ func AggregateRTPDeltaInfo(deltaInfoList []*RTPDeltaInfo) *RTPDeltaInfo { bytesPadding := uint64(0) headerBytesPadding := uint64(0) + fecPacketsReceived := uint32(0) + fecBytesReceived := uint64(0) + fecPacketsDiscarded := uint32(0) + packetsRecovered := uint32(0) + packetsLost := uint32(0) packetsMissing := uint32(0) packetsOutOfOrder := uint32(0) @@ -894,6 +967,11 @@ func AggregateRTPDeltaInfo(deltaInfoList []*RTPDeltaInfo) *RTPDeltaInfo { bytesPadding += deltaInfo.BytesPadding headerBytesPadding += deltaInfo.HeaderBytesPadding + fecPacketsReceived += deltaInfo.FecPacketsReceived + fecBytesReceived += deltaInfo.FecBytesReceived + fecPacketsDiscarded += deltaInfo.FecPacketsDiscarded + packetsRecovered += deltaInfo.PacketsRecovered + packetsLost += deltaInfo.PacketsLost packetsMissing += deltaInfo.PacketsMissing packetsOutOfOrder += deltaInfo.PacketsOutOfOrder @@ -928,6 +1006,10 @@ func AggregateRTPDeltaInfo(deltaInfoList []*RTPDeltaInfo) *RTPDeltaInfo { PacketsPadding: packetsPadding, BytesPadding: bytesPadding, HeaderBytesPadding: headerBytesPadding, + FecPacketsReceived: fecPacketsReceived, + FecBytesReceived: fecBytesReceived, + FecPacketsDiscarded: fecPacketsDiscarded, + PacketsRecovered: packetsRecovered, PacketsLost: packetsLost, PacketsMissing: packetsMissing, PacketsOutOfOrder: packetsOutOfOrder, diff --git a/pkg/sfu/rtpstats/rtpstats_test.go b/pkg/sfu/rtpstats/rtpstats_test.go index 9039ad773..c55c1b287 100644 --- a/pkg/sfu/rtpstats/rtpstats_test.go +++ b/pkg/sfu/rtpstats/rtpstats_test.go @@ -223,6 +223,52 @@ func Test_RTPStatsReceiver_Update(t *testing.T) { r.Stop() } +func Test_RTPStatsReceiver_UpdateFEC(t *testing.T) { + r := NewRTPStatsReceiver(RTPStatsParams{}) + r.SetClockRate(90000) + snapshotID := r.NewSnapshotId() + + packet := getPacket(100, 3000, 1000) + r.Update( + time.Now().UnixNano(), + packet.SequenceNumber, + packet.Timestamp, + packet.Marker, + packet.Header.MarshalSize(), + len(packet.Payload), + 0, + ) + r.UpdateFEC(3, 900, 1, 2) + + stats := r.ToProto() + require.EqualValues(t, 3, stats.FecPacketsReceived) + require.EqualValues(t, 900, stats.FecBytesReceived) + require.EqualValues(t, 1, stats.FecPacketsDiscarded) + require.EqualValues(t, 2, stats.PacketsRecovered) + + delta := r.DeltaInfo(snapshotID) + require.EqualValues(t, 3, delta.FecPacketsReceived) + require.EqualValues(t, 900, delta.FecBytesReceived) + require.EqualValues(t, 1, delta.FecPacketsDiscarded) + require.EqualValues(t, 2, delta.PacketsRecovered) + + // 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, 1, delta.FecPacketsDiscarded) + require.EqualValues(t, 1, delta.PacketsRecovered) + + 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, 2, stats.FecPacketsDiscarded) + require.EqualValues(t, 3, stats.PacketsRecovered) +} + func Test_RTPStatsReceiver_Restart(t *testing.T) { clockRate := uint32(90000) r := NewRTPStatsReceiver(RTPStatsParams{})