From 97048a923c2278dcd9ca46fad605ed852ac4c6f8 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Sat, 16 Sep 2023 18:54:18 +0530 Subject: [PATCH] Reducing rtp stats memory consumption - part 2 (#2078) * WIP commit * move a struct to sender only * Snapshot intervals * make receiver history 4K too --- pkg/sfu/buffer/rtpstats_base.go | 151 +--------- pkg/sfu/buffer/rtpstats_receiver.go | 9 +- pkg/sfu/buffer/rtpstats_sender.go | 441 +++++++++++++++++++++------- 3 files changed, 352 insertions(+), 249 deletions(-) diff --git a/pkg/sfu/buffer/rtpstats_base.go b/pkg/sfu/buffer/rtpstats_base.go index 3fad1907c..d25341c8b 100644 --- a/pkg/sfu/buffer/rtpstats_base.go +++ b/pkg/sfu/buffer/rtpstats_base.go @@ -30,8 +30,6 @@ const ( cGapHistogramNumBins = 101 cNumSequenceNumbers = 65536 cFirstSnapshotID = 1 - cSnInfoSize = 8192 - cSnInfoMask = cSnInfoSize - 1 cFirstPacketTimeAdjustWindow = 2 * time.Minute cFirstPacketTimeAdjustThreshold = 5 * time.Second @@ -53,18 +51,6 @@ func RTPDriftToString(r *livekit.RTPDrift) string { // ------------------------------------------------------- -type intervalStats struct { - packets uint64 - bytes uint64 - headerBytes uint64 - packetsPadding uint64 - bytesPadding uint64 - headerBytesPadding uint64 - packetsLost uint64 - packetsOutOfOrder uint64 - frames uint32 -} - type RTPDeltaInfo struct { StartTime time.Time Duration time.Duration @@ -119,14 +105,6 @@ type snapshot struct { maxJitter float64 } -type snInfo struct { - hdrSize uint16 - pktSize uint16 - isPaddingOnly bool - marker bool - isOutOfOrder bool -} - type RTCPSenderReportData struct { RTPTimestamp uint32 RTPTimestampExt uint64 @@ -177,8 +155,6 @@ type rtpStatsBase struct { jitter float64 maxJitter float64 - snInfos [cSnInfoSize]snInfo - gapHistogram [cGapHistogramNumBins]uint32 nacks uint32 @@ -251,8 +227,6 @@ func (r *rtpStatsBase) seed(from *rtpStatsBase) bool { r.jitter = from.jitter r.maxJitter = from.maxJitter - r.snInfos = from.snInfos - r.gapHistogram = from.gapHistogram r.nacks = from.nacks @@ -309,14 +283,14 @@ func (r *rtpStatsBase) newSnapshotID(extStartSN uint64) uint32 { id := r.nextSnapshotID r.nextSnapshotID++ - if cap(r.snapshots) < int(r.nextSnapshotID) { - snapshots := make([]snapshot, r.nextSnapshotID) + if cap(r.snapshots) < int(r.nextSnapshotID-cFirstSnapshotID) { + snapshots := make([]snapshot, r.nextSnapshotID-cFirstSnapshotID) copy(snapshots, r.snapshots) r.snapshots = snapshots } if r.initialized { - r.snapshots[id] = r.initSnapshot(time.Now(), extStartSN) + r.snapshots[id-cFirstSnapshotID] = r.initSnapshot(time.Now(), extStartSN) } return id } @@ -467,7 +441,8 @@ func (r *rtpStatsBase) UpdateRtt(rtt uint32) { r.maxRtt = rtt } - for _, s := range r.snapshots { + for i := uint32(0); i < r.nextSnapshotID-cFirstSnapshotID; i++ { + s := &r.snapshots[i] if rtt > s.maxRtt { s.maxRtt = rtt } @@ -545,7 +520,6 @@ func (r *rtpStatsBase) getTotalPacketsPrimary(extStartSN, extHighestSN uint64) u func (r *rtpStatsBase) deltaInfo(snapshotID uint32, extStartSN uint64, extHighestSN uint64) *RTPDeltaInfo { then, now := r.getAndResetSnapshot(snapshotID, extStartSN, extHighestSN) - if now == nil || then == nil { return nil } @@ -772,107 +746,6 @@ func (r *rtpStatsBase) toProto( return p } -func (r *rtpStatsBase) getSnInfoOutOfOrderSlot(esn uint64, ehsn uint64) int { - offset := int64(ehsn - esn) - if offset >= cSnInfoSize || offset < 0 { - // too old OR too new (i. e. ahead of highest) - return -1 - } - - return int(esn & cSnInfoMask) -} - -func (r *rtpStatsBase) setSnInfo(esn uint64, ehsn uint64, pktSize uint16, hdrSize uint16, payloadSize uint16, marker bool, isOutOfOrder bool) { - var slot int - if int64(esn-ehsn) < 0 { - slot = r.getSnInfoOutOfOrderSlot(esn, ehsn) - if slot < 0 { - return - } - } else { - slot = int(esn & cSnInfoMask) - } - - snInfo := &r.snInfos[slot] - snInfo.pktSize = pktSize - snInfo.hdrSize = hdrSize - snInfo.isPaddingOnly = payloadSize == 0 - snInfo.marker = marker - snInfo.isOutOfOrder = isOutOfOrder -} - -func (r *rtpStatsBase) clearSnInfos(extStartInclusive uint64, extEndExclusive uint64) { - if extEndExclusive <= extStartInclusive { - return - } - - for esn := extStartInclusive; esn != extEndExclusive; esn++ { - snInfo := &r.snInfos[esn&cSnInfoMask] - snInfo.pktSize = 0 - snInfo.hdrSize = 0 - snInfo.isPaddingOnly = false - snInfo.marker = false - } -} - -func (r *rtpStatsBase) isSnInfoLost(esn uint64, ehsn uint64) bool { - slot := r.getSnInfoOutOfOrderSlot(esn, ehsn) - if slot < 0 { - return false - } - - return r.snInfos[slot].pktSize == 0 -} - -func (r *rtpStatsBase) getIntervalStats(extStartInclusive uint64, extEndExclusive uint64, ehsn uint64) (intervalStats intervalStats) { - packetsNotFound := uint32(0) - processESN := func(esn uint64, ehsn uint64) { - slot := r.getSnInfoOutOfOrderSlot(esn, ehsn) - if slot < 0 { - packetsNotFound++ - return - } - - snInfo := &r.snInfos[slot] - switch { - case snInfo.pktSize == 0: - intervalStats.packetsLost++ - - case snInfo.isPaddingOnly: - intervalStats.packetsPadding++ - intervalStats.bytesPadding += uint64(snInfo.pktSize) - intervalStats.headerBytesPadding += uint64(snInfo.hdrSize) - - default: - intervalStats.packets++ - intervalStats.bytes += uint64(snInfo.pktSize) - intervalStats.headerBytes += uint64(snInfo.hdrSize) - if snInfo.isOutOfOrder { - intervalStats.packetsOutOfOrder++ - } - } - - if snInfo.marker { - intervalStats.frames++ - } - } - - for esn := extStartInclusive; esn != extEndExclusive; esn++ { - processESN(esn, ehsn) - } - - if packetsNotFound != 0 { - r.logger.Errorw( - "could not find some packets", nil, - "start", extStartInclusive, - "end", extEndExclusive, - "count", packetsNotFound, - "highestSN", ehsn, - ) - } - return -} - func (r *rtpStatsBase) updateJitter(ets uint64, packetTime time.Time) float64 { // Do not update jitter on multiple packets of same frame. // All packets of a frame have the same time stamp. @@ -896,7 +769,8 @@ func (r *rtpStatsBase) updateJitter(ets uint64, packetTime time.Time) float64 { r.maxJitter = r.jitter } - for _, s := range r.snapshots { + for i := uint32(0); i < r.nextSnapshotID-cFirstSnapshotID; i++ { + s := &r.snapshots[i] if r.jitter > s.maxJitter { s.maxJitter = r.jitter } @@ -914,15 +788,16 @@ func (r *rtpStatsBase) getAndResetSnapshot(snapshotID uint32, extStartSN uint64, return nil, nil } - then := r.snapshots[snapshotID] + idx := snapshotID - cFirstSnapshotID + then := r.snapshots[idx] if !then.isValid { then = r.initSnapshot(r.startTime, extStartSN) - r.snapshots[snapshotID] = then + r.snapshots[idx] = then } // snapshot now now := r.getSnapshot(time.Now(), extHighestSN+1) - r.snapshots[snapshotID] = now + r.snapshots[idx] = now return &then, &now } @@ -983,7 +858,7 @@ func (r *rtpStatsBase) updateGapHistogram(gap int) { func (r *rtpStatsBase) initSnapshot(startTime time.Time, extStartSN uint64) snapshot { return snapshot{ isValid: true, - startTime: time.Now(), + startTime: startTime, extStartSN: extStartSN, } } @@ -991,7 +866,7 @@ func (r *rtpStatsBase) initSnapshot(startTime time.Time, extStartSN uint64) snap func (r *rtpStatsBase) getSnapshot(startTime time.Time, extStartSN uint64) snapshot { return snapshot{ isValid: true, - startTime: time.Now(), + startTime: startTime, extStartSN: extStartSN, bytes: r.bytes, headerBytes: r.headerBytes, diff --git a/pkg/sfu/buffer/rtpstats_receiver.go b/pkg/sfu/buffer/rtpstats_receiver.go index 2b409c990..fff35b2f5 100644 --- a/pkg/sfu/buffer/rtpstats_receiver.go +++ b/pkg/sfu/buffer/rtpstats_receiver.go @@ -26,7 +26,7 @@ import ( ) const ( - cHistorySize = 2048 + cHistorySize = 4096 ) type RTPFlowState struct { @@ -69,7 +69,7 @@ func (r *RTPStatsReceiver) NewSnapshotId() uint32 { r.lock.Lock() defer r.lock.Unlock() - return r.newSnapshotID(r.sequenceNumber.GetExtendedStart()) + return r.newSnapshotID(r.sequenceNumber.GetExtendedHighest()) } func (r *RTPStatsReceiver) Update( @@ -114,7 +114,7 @@ func (r *RTPStatsReceiver) Update( resTS = r.timestamp.Update(timestamp) // initialize snapshots if any - for i := uint32(cFirstSnapshotID); i < r.nextSnapshotID; i++ { + for i := uint32(0); i < r.nextSnapshotID-cFirstSnapshotID; i++ { r.snapshots[i] = r.initSnapshot(r.startTime, r.sequenceNumber.GetExtendedStart()) } @@ -154,7 +154,8 @@ func (r *RTPStatsReceiver) Update( r.packetsLost += resSN.PreExtendedStart - resSN.ExtendedVal extStartSN := r.sequenceNumber.GetExtendedStart() - for _, s := range r.snapshots { + for i := uint32(0); i < r.nextSnapshotID-cFirstSnapshotID; i++ { + s := &r.snapshots[i] if s.extStartSN == resSN.PreExtendedStart { s.extStartSN = extStartSN } diff --git a/pkg/sfu/buffer/rtpstats_sender.go b/pkg/sfu/buffer/rtpstats_sender.go index 670032e83..5a66c65d4 100644 --- a/pkg/sfu/buffer/rtpstats_sender.go +++ b/pkg/sfu/buffer/rtpstats_sender.go @@ -25,11 +25,91 @@ import ( "github.com/livekit/protocol/livekit" ) +const ( + cSnInfoSize = 4096 + cSnInfoMask = cSnInfoSize - 1 +) + +type snInfoFlag byte + +const ( + snInfoFlagMarker snInfoFlag = 1 << iota + snInfoFlagPadding + snInfoFlagOutOfOrder +) + +type snInfo struct { + pktSize uint16 + hdrSize uint8 + flags snInfoFlag +} + +// ------------------------------------------------------------------- + +type intervalStats struct { + packets uint64 + bytes uint64 + headerBytes uint64 + packetsPadding uint64 + bytesPadding uint64 + headerBytesPadding uint64 + packetsLost uint64 + packetsOutOfOrder uint64 + frames uint32 +} + +func (is *intervalStats) aggregate(other *intervalStats) { + if is == nil || other == nil { + return + } + + is.packets += other.packets + is.bytes += other.bytes + is.headerBytes += other.headerBytes + is.packetsPadding += other.packetsPadding + is.bytesPadding += other.bytesPadding + is.headerBytesPadding += other.headerBytesPadding + is.packetsLost += other.packetsLost + is.packetsOutOfOrder += other.packetsOutOfOrder + is.frames += other.frames +} + +// ------------------------------------------------------------------- + type senderSnapshot struct { - snapshot - extStartSNFromRR uint64 - packetsLostFromRR uint64 - maxJitterFromRR float64 + isValid bool + + startTime time.Time + + extStartSN uint64 + bytes uint64 + headerBytes uint64 + + packetsPadding uint64 + bytesPadding uint64 + headerBytesPadding uint64 + + packetsDuplicate uint64 + bytesDuplicate uint64 + headerBytesDuplicate uint64 + + packetsOutOfOrder uint64 + + packetsLostFeed uint64 + packetsLost uint64 + + frames uint32 + + nacks uint32 + plis uint32 + firs uint32 + + maxRtt uint32 + maxJitterFeed float64 + maxJitter float64 + + extLastRRSN uint64 + intervalStats intervalStats } type RTPStatsSender struct { @@ -50,6 +130,8 @@ type RTPStatsSender struct { jitterFromRR float64 maxJitterFromRR float64 + snInfos [cSnInfoSize]snInfo + nextSenderSnapshotID uint32 senderSnapshots []senderSnapshot } @@ -85,6 +167,8 @@ func (r *RTPStatsSender) Seed(from *RTPStatsSender) { r.jitterFromRR = from.jitterFromRR r.maxJitterFromRR = from.maxJitterFromRR + r.snInfos = from.snInfos + r.nextSenderSnapshotID = from.nextSenderSnapshotID r.senderSnapshots = make([]senderSnapshot, cap(from.senderSnapshots)) copy(r.senderSnapshots, from.senderSnapshots) @@ -94,7 +178,7 @@ func (r *RTPStatsSender) NewSnapshotId() uint32 { r.lock.Lock() defer r.lock.Unlock() - return r.newSnapshotID(r.extStartSN) + return r.newSnapshotID(r.extHighestSN) } func (r *RTPStatsSender) NewSenderSnapshotId() uint32 { @@ -104,17 +188,14 @@ func (r *RTPStatsSender) NewSenderSnapshotId() uint32 { id := r.nextSenderSnapshotID r.nextSenderSnapshotID++ - if cap(r.senderSnapshots) < int(r.nextSenderSnapshotID) { - senderSnapshots := make([]senderSnapshot, r.nextSenderSnapshotID) + if cap(r.senderSnapshots) < int(r.nextSenderSnapshotID-cFirstSnapshotID) { + senderSnapshots := make([]senderSnapshot, r.nextSenderSnapshotID-cFirstSnapshotID) copy(senderSnapshots, r.senderSnapshots) r.senderSnapshots = senderSnapshots } if r.initialized { - r.senderSnapshots[id] = senderSnapshot{ - snapshot: r.initSnapshot(time.Now(), r.extStartSN), - extStartSNFromRR: r.extStartSN, - } + r.senderSnapshots[id-cFirstSnapshotID] = r.initSenderSnapshot(time.Now(), r.extHighestSN) } return id } @@ -155,22 +236,11 @@ func (r *RTPStatsSender) Update( r.extHighestTS = extTimestamp // initialize snapshots if any - for i := uint32(cFirstSnapshotID); i < r.nextSnapshotID; i++ { - r.snapshots[i] = snapshot{ - isValid: true, - startTime: r.startTime, - extStartSN: r.extStartSN, - } + for i := uint32(0); i < r.nextSnapshotID-cFirstSnapshotID; i++ { + r.snapshots[i] = r.initSnapshot(r.startTime, r.extStartSN) } - for i := uint32(cFirstSnapshotID); i < r.nextSenderSnapshotID; i++ { - r.senderSnapshots[i] = senderSnapshot{ - snapshot: snapshot{ - isValid: true, - startTime: r.startTime, - extStartSN: r.extStartSN, - }, - extStartSNFromRR: r.extStartSN, - } + for i := uint32(0); i < r.nextSenderSnapshotID-cFirstSnapshotID; i++ { + r.senderSnapshots[i] = r.initSenderSnapshot(r.startTime, r.extStartSN) } r.logger.Debugw( @@ -195,14 +265,19 @@ func (r *RTPStatsSender) Update( r.packetsLost += r.extStartSN - extSequenceNumber // adjust start of snapshots - for _, s := range r.snapshots { + for i := uint32(0); i < r.nextSnapshotID-cFirstSnapshotID; i++ { + s := &r.snapshots[i] if s.extStartSN == r.extStartSN { s.extStartSN = extSequenceNumber } } - for _, s := range r.senderSnapshots { + for i := uint32(0); i < r.nextSenderSnapshotID-cFirstSnapshotID; i++ { + s := &r.senderSnapshots[i] if s.extStartSN == r.extStartSN { s.extStartSN = extSequenceNumber + if s.extLastRRSN == (r.extStartSN - 1) { + s.extLastRRSN = extSequenceNumber - 1 + } } } @@ -224,7 +299,7 @@ func (r *RTPStatsSender) Update( isDuplicate = true } else { r.packetsLost-- - r.setSnInfo(extSequenceNumber, r.extHighestSN, uint16(pktSize), uint16(hdrSize), uint16(payloadSize), marker, true) + r.setSnInfo(extSequenceNumber, r.extHighestSN, uint16(pktSize), uint8(hdrSize), uint16(payloadSize), marker, true) } } else { // in-order // update gap histogram @@ -234,7 +309,7 @@ func (r *RTPStatsSender) Update( r.clearSnInfos(r.extHighestSN+1, extSequenceNumber) r.packetsLost += uint64(gapSN - 1) - r.setSnInfo(extSequenceNumber, r.extHighestSN, uint16(pktSize), uint16(hdrSize), uint16(payloadSize), marker, false) + r.setSnInfo(extSequenceNumber, r.extHighestSN, uint16(pktSize), uint8(hdrSize), uint16(payloadSize), marker, false) if extTimestamp != r.extHighestTS { // update only on first packet as same timestamp could be in multiple packets. @@ -259,9 +334,10 @@ func (r *RTPStatsSender) Update( } jitter := r.updateJitter(extTimestamp, packetTime) - for _, s := range r.senderSnapshots { - if jitter > s.maxJitter { - s.maxJitter = jitter + for i := uint32(0); i < r.nextSenderSnapshotID-cFirstSnapshotID; i++ { + s := &r.senderSnapshots[i] + if jitter > s.maxJitterFeed { + s.maxJitterFeed = jitter } } } @@ -307,46 +383,7 @@ func (r *RTPStatsSender) UpdateFromReceiverReport(rr rtcp.ReceptionReport) (rtt } } - if r.lastRRTime.IsZero() || r.extHighestSNFromRR <= extHighestSNFromRR { - r.extHighestSNFromRR = extHighestSNFromRR - - packetsLostFromRR := r.packetsLostFromRR&0xFFFF_FFFF_0000_0000 + uint64(rr.TotalLost) - if (rr.TotalLost-r.lastRR.TotalLost) < (1<<31) && rr.TotalLost < r.lastRR.TotalLost { - packetsLostFromRR += (1 << 32) - } - r.packetsLostFromRR = packetsLostFromRR - - if isRttChanged { - r.rtt = rtt - if rtt > r.maxRtt { - r.maxRtt = rtt - } - } - - r.jitterFromRR = float64(rr.Jitter) - if r.jitterFromRR > r.maxJitterFromRR { - r.maxJitterFromRR = r.jitterFromRR - } - - // update snapshots - for _, s := range r.snapshots { - if isRttChanged && rtt > s.maxRtt { - s.maxRtt = rtt - } - } - for _, s := range r.senderSnapshots { - if isRttChanged && rtt > s.maxRtt { - s.maxRtt = rtt - } - - if r.jitterFromRR > s.maxJitterFromRR { - s.maxJitterFromRR = r.jitterFromRR - } - } - - r.lastRRTime = time.Now() - r.lastRR = rr - } else { + if !r.lastRRTime.IsZero() && r.extHighestSNFromRR > extHighestSNFromRR { r.logger.Debugw( fmt.Sprintf("receiver report potentially out of order, highestSN: existing: %d, received: %d", r.extHighestSNFromRR, extHighestSNFromRR), "lastRRTime", r.lastRRTime, @@ -354,7 +391,57 @@ func (r *RTPStatsSender) UpdateFromReceiverReport(rr rtcp.ReceptionReport) (rtt "sinceLastRR", time.Since(r.lastRRTime), "receivedRR", rr, ) + return } + + r.extHighestSNFromRR = extHighestSNFromRR + + packetsLostFromRR := r.packetsLostFromRR&0xFFFF_FFFF_0000_0000 + uint64(rr.TotalLost) + if (rr.TotalLost-r.lastRR.TotalLost) < (1<<31) && rr.TotalLost < r.lastRR.TotalLost { + packetsLostFromRR += (1 << 32) + } + r.packetsLostFromRR = packetsLostFromRR + + if isRttChanged { + r.rtt = rtt + if rtt > r.maxRtt { + r.maxRtt = rtt + } + } + + r.jitterFromRR = float64(rr.Jitter) + if r.jitterFromRR > r.maxJitterFromRR { + r.maxJitterFromRR = r.jitterFromRR + } + + // update snapshots + for i := uint32(0); i < r.nextSnapshotID-cFirstSnapshotID; i++ { + s := &r.snapshots[i] + if isRttChanged && rtt > s.maxRtt { + s.maxRtt = rtt + } + } + + extLastRRSN := r.extHighestSNFromRR + (r.extStartSN & 0xFFFF_FFFF_FFFF_0000) + for i := uint32(0); i < r.nextSenderSnapshotID-cFirstSnapshotID; i++ { + s := &r.senderSnapshots[i] + if isRttChanged && rtt > s.maxRtt { + s.maxRtt = rtt + } + + if r.jitterFromRR > s.maxJitter { + s.maxJitter = r.jitterFromRR + } + + // on every RR, calculate delta since last RR using packet metadata cache + is := r.getIntervalStats(s.extLastRRSN+1, extLastRRSN+1, r.extHighestSN) + eis := &s.intervalStats + eis.aggregate(&is) + s.extLastRRSN = extLastRRSN + } + + r.lastRRTime = time.Now() + r.lastRR = rr return } @@ -492,11 +579,11 @@ func (r *RTPStatsSender) DeltaInfoSender(senderSnapshotID uint32) *RTPDeltaInfo startTime := then.startTime endTime := now.startTime - packetsExpected := now.extStartSNFromRR - then.extStartSNFromRR + packetsExpected := uint32(now.extStartSN - then.extStartSN) if packetsExpected > cNumSequenceNumbers { r.logger.Warnw( "too many packets expected in delta (sender)", - fmt.Errorf("start: %d, end: %d, expected: %d", then.extStartSNFromRR, now.extStartSNFromRR, packetsExpected), + fmt.Errorf("start: %d, end: %d, expected: %d", then.extStartSN, now.extStartSN, packetsExpected), ) return nil } @@ -505,29 +592,31 @@ func (r *RTPStatsSender) DeltaInfoSender(senderSnapshotID uint32) *RTPDeltaInfo return nil } - intervalStats := r.getIntervalStats(then.extStartSNFromRR, now.extStartSNFromRR, r.extHighestSN) - packetsLost := now.packetsLostFromRR - then.packetsLostFromRR + packetsLost := uint32(now.packetsLost - then.packetsLost) if int32(packetsLost) < 0 { packetsLost = 0 } - + packetsLostFeed := uint32(now.packetsLostFeed - then.packetsLostFeed) + if int32(packetsLostFeed) < 0 { + packetsLostFeed = 0 + } if packetsLost > packetsExpected { r.logger.Warnw( "unexpected number of packets lost", fmt.Errorf( - "start: %d, end: %d, expected: %d, lost: report: %d, interval: %d", - then.extStartSNFromRR, - now.extStartSNFromRR, + "start: %d, end: %d, expected: %d, lost: report: %d, feed: %d", + then.extStartSN, + now.extStartSN, packetsExpected, - now.packetsLostFromRR-then.packetsLostFromRR, - intervalStats.packetsLost, + packetsLost, + packetsLostFeed, ), ) packetsLost = packetsExpected } // discount jitter from publisher side + internal processing - maxJitter := then.maxJitterFromRR - then.maxJitter + maxJitter := then.maxJitter - then.maxJitterFeed if maxJitter < 0.0 { maxJitter = 0.0 } @@ -536,19 +625,19 @@ func (r *RTPStatsSender) DeltaInfoSender(senderSnapshotID uint32) *RTPDeltaInfo return &RTPDeltaInfo{ StartTime: startTime, Duration: endTime.Sub(startTime), - Packets: uint32(packetsExpected - intervalStats.packetsPadding), - Bytes: intervalStats.bytes, - HeaderBytes: intervalStats.headerBytes, + Packets: packetsExpected - uint32(now.packetsPadding-then.packetsPadding), + Bytes: now.bytes - then.bytes, + HeaderBytes: now.headerBytes - then.headerBytes, PacketsDuplicate: uint32(now.packetsDuplicate - then.packetsDuplicate), BytesDuplicate: now.bytesDuplicate - then.bytesDuplicate, HeaderBytesDuplicate: now.headerBytesDuplicate - then.headerBytesDuplicate, - PacketsPadding: uint32(intervalStats.packetsPadding), - BytesPadding: intervalStats.bytesPadding, - HeaderBytesPadding: intervalStats.headerBytesPadding, - PacketsLost: uint32(packetsLost), - PacketsMissing: uint32(intervalStats.packetsLost), - PacketsOutOfOrder: uint32(intervalStats.packetsOutOfOrder), - Frames: intervalStats.frames, + PacketsPadding: uint32(now.packetsPadding - then.packetsPadding), + BytesPadding: now.bytesPadding - then.bytesPadding, + HeaderBytesPadding: now.headerBytesPadding - then.headerBytesPadding, + PacketsLost: packetsLost, + PacketsMissing: packetsLostFeed, + PacketsOutOfOrder: uint32(now.packetsOutOfOrder - then.packetsOutOfOrder), + Frames: now.frames - then.frames, RttMax: then.maxRtt, JitterMax: maxJitterTime, Nacks: now.nacks - then.nacks, @@ -584,24 +673,162 @@ func (r *RTPStatsSender) getAndResetSenderSnapshot(senderSnapshotID uint32) (*se return nil, nil } - then := r.senderSnapshots[senderSnapshotID] + idx := senderSnapshotID - cFirstSnapshotID + then := r.senderSnapshots[idx] if !then.isValid { - then = senderSnapshot{ - snapshot: r.initSnapshot(r.startTime, r.extStartSN), - extStartSNFromRR: r.extStartSN, - } - r.senderSnapshots[senderSnapshotID] = then + then = r.initSenderSnapshot(r.startTime, r.extStartSN) + r.senderSnapshots[idx] = then } // snapshot now - now := senderSnapshot{ - snapshot: r.getSnapshot(r.lastRRTime, r.extHighestSN+1), - extStartSNFromRR: r.extHighestSNFromRR + (r.extStartSN & 0xFFFF_FFFF_FFFF_0000) + 1, - packetsLostFromRR: r.packetsLostFromRR, - maxJitterFromRR: r.jitterFromRR, - } - r.senderSnapshots[senderSnapshotID] = now + now := r.getSenderSnapshot(r.lastRRTime, &then) + r.senderSnapshots[idx] = now return &then, &now } +func (r *RTPStatsSender) initSenderSnapshot(startTime time.Time, extStartSN uint64) senderSnapshot { + return senderSnapshot{ + isValid: true, + startTime: startTime, + extStartSN: extStartSN, + extLastRRSN: extStartSN - 1, + } +} + +func (r *RTPStatsSender) getSenderSnapshot(startTime time.Time, s *senderSnapshot) senderSnapshot { + if s == nil { + return senderSnapshot{} + } + + return senderSnapshot{ + isValid: true, + startTime: startTime, + extStartSN: s.extLastRRSN + 1, + bytes: s.bytes + s.intervalStats.bytes, + headerBytes: s.headerBytes + s.intervalStats.headerBytes, + packetsPadding: s.packetsPadding + s.intervalStats.packetsPadding, + bytesPadding: s.bytesPadding + s.intervalStats.bytesPadding, + headerBytesPadding: s.headerBytesPadding + s.intervalStats.headerBytesPadding, + packetsDuplicate: r.packetsDuplicate, + bytesDuplicate: r.bytesDuplicate, + headerBytesDuplicate: r.headerBytesDuplicate, + packetsLostFeed: r.packetsLost, + packetsOutOfOrder: s.packetsOutOfOrder + s.intervalStats.packetsOutOfOrder, + frames: s.frames + s.intervalStats.frames, + nacks: r.nacks, + plis: r.plis, + firs: r.firs, + maxRtt: r.rtt, + maxJitterFeed: r.jitter, + maxJitter: r.jitterFromRR, + extLastRRSN: s.extLastRRSN, + } +} + +func (r *RTPStatsSender) getSnInfoOutOfOrderSlot(esn uint64, ehsn uint64) int { + offset := int64(ehsn - esn) + if offset >= cSnInfoSize || offset < 0 { + // too old OR too new (i. e. ahead of highest) + return -1 + } + + return int(esn & cSnInfoMask) +} + +func (r *RTPStatsSender) setSnInfo(esn uint64, ehsn uint64, pktSize uint16, hdrSize uint8, payloadSize uint16, marker bool, isOutOfOrder bool) { + var slot int + if int64(esn-ehsn) < 0 { + slot = r.getSnInfoOutOfOrderSlot(esn, ehsn) + if slot < 0 { + return + } + } else { + slot = int(esn & cSnInfoMask) + } + + snInfo := &r.snInfos[slot] + snInfo.pktSize = pktSize + snInfo.hdrSize = hdrSize + if marker { + snInfo.flags |= snInfoFlagMarker + } + if payloadSize == 0 { + snInfo.flags |= snInfoFlagPadding + } + if isOutOfOrder { + snInfo.flags |= snInfoFlagOutOfOrder + } +} + +func (r *RTPStatsSender) clearSnInfos(extStartInclusive uint64, extEndExclusive uint64) { + if extEndExclusive <= extStartInclusive { + return + } + + for esn := extStartInclusive; esn != extEndExclusive; esn++ { + snInfo := &r.snInfos[esn&cSnInfoMask] + snInfo.pktSize = 0 + snInfo.hdrSize = 0 + snInfo.flags = 0 + } +} + +func (r *RTPStatsSender) isSnInfoLost(esn uint64, ehsn uint64) bool { + slot := r.getSnInfoOutOfOrderSlot(esn, ehsn) + if slot < 0 { + return false + } + + return r.snInfos[slot].pktSize == 0 +} + +func (r *RTPStatsSender) getIntervalStats(extStartInclusive uint64, extEndExclusive uint64, ehsn uint64) (intervalStats intervalStats) { + packetsNotFound := uint32(0) + processESN := func(esn uint64, ehsn uint64) { + slot := r.getSnInfoOutOfOrderSlot(esn, ehsn) + if slot < 0 { + packetsNotFound++ + return + } + + snInfo := &r.snInfos[slot] + switch { + case snInfo.pktSize == 0: + intervalStats.packetsLost++ + + case snInfo.flags&snInfoFlagPadding != 0: + intervalStats.packetsPadding++ + intervalStats.bytesPadding += uint64(snInfo.pktSize) + intervalStats.headerBytesPadding += uint64(snInfo.hdrSize) + + default: + intervalStats.packets++ + intervalStats.bytes += uint64(snInfo.pktSize) + intervalStats.headerBytes += uint64(snInfo.hdrSize) + if (snInfo.flags & snInfoFlagOutOfOrder) != 0 { + intervalStats.packetsOutOfOrder++ + } + } + + if (snInfo.flags & snInfoFlagMarker) != 0 { + intervalStats.frames++ + } + } + + for esn := extStartInclusive; esn != extEndExclusive; esn++ { + processESN(esn, ehsn) + } + + if packetsNotFound != 0 { + r.logger.Errorw( + "could not find some packets", nil, + "start", extStartInclusive, + "end", extEndExclusive, + "count", packetsNotFound, + "highestSN", ehsn, + ) + } + return +} + // -------------------------------------------------------------------