From 7fef374b193d3f44f238d70188ce90d67476a964 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Tue, 11 Feb 2025 16:05:00 +0530 Subject: [PATCH] Split down stream snapshot into sender view and receiver view. (#3422) Receiver view is used for connection quality. Sender view is used for analytics. One thing that this introduces is that sender view uses the packet loss information from receiver view as true loss is available only in the RTCP Receiver Reports received from the remote side. So, the time alignment is off, i. e. receiver report happens periodically and it includes information till the time at which it was sent from remote side, but sender could have sent more packets after that time. The split should ensure that analytics does not rely on remote side sending proper receiver repoerts albeit at slight misalignment of loss statistic for remotes that send RTCP RR (which should be majority of the cases) --- pkg/sfu/buffer/streamstats.go | 2 + pkg/sfu/connectionquality/connectionstats.go | 49 +- pkg/sfu/downtrack.go | 15 +- pkg/sfu/playoutdelay.go | 2 +- pkg/sfu/rtpstats/rtpstats_base.go | 17 +- pkg/sfu/rtpstats/rtpstats_sender.go | 511 +++++++++++++------ 6 files changed, 413 insertions(+), 183 deletions(-) diff --git a/pkg/sfu/buffer/streamstats.go b/pkg/sfu/buffer/streamstats.go index e38cd1980..e45cb2eec 100644 --- a/pkg/sfu/buffer/streamstats.go +++ b/pkg/sfu/buffer/streamstats.go @@ -19,4 +19,6 @@ import "github.com/livekit/livekit-server/pkg/sfu/rtpstats" type StreamStatsWithLayers struct { RTPStats *rtpstats.RTPDeltaInfo Layers map[int32]*rtpstats.RTPDeltaInfo + + RTPStatsRemoteView *rtpstats.RTPDeltaInfo } diff --git a/pkg/sfu/connectionquality/connectionstats.go b/pkg/sfu/connectionquality/connectionstats.go index f4f09b7c5..ea6ff8bcc 100644 --- a/pkg/sfu/connectionquality/connectionstats.go +++ b/pkg/sfu/connectionquality/connectionstats.go @@ -264,7 +264,12 @@ func (cs *ConnectionStats) updateScoreFromReceiverReport(at time.Time) (float32, return mos, streams } - agg := toAggregateDeltaInfo(streams) + agg := toAggregateDeltaInfo(streams, true) + if agg == nil { + // no receiver report in the window + mos, _ := cs.scorer.GetMOSAndQuality() + return mos, streams + } if streamingStartedAt.After(agg.StartTime) { agg.StartTime = streamingStartedAt } @@ -287,11 +292,12 @@ func (cs *ConnectionStats) updateScoreAt(at time.Time) (float32, map[uint32]*buf return mos, nil } - deltaInfoList := make([]*rtpstats.RTPDeltaInfo, 0, len(streams)) - for _, s := range streams { - deltaInfoList = append(deltaInfoList, s.RTPStats) + agg := toAggregateDeltaInfo(streams, false) + if agg == nil { + // no receiver report in the window + mos, _ := cs.scorer.GetMOSAndQuality() + return mos, streams } - agg := rtpstats.AggregateRTPDeltaInfo(deltaInfoList) return cs.updateScoreWithAggregate(agg, cs.params.ReceiverProvider.GetLastSenderReportTime(), at), streams } @@ -323,7 +329,7 @@ func (cs *ConnectionStats) getStat() { if cs.onStatsUpdate != nil && len(streams) != 0 { analyticsStreams := make([]*livekit.AnalyticsStream, 0, len(streams)) for ssrc, stream := range streams { - as := toAnalyticsStream(ssrc, stream.RTPStats) + as := toAnalyticsStream(ssrc, stream.RTPStats, stream.RTPStatsRemoteView) // // add video layer if either @@ -413,21 +419,36 @@ func getPacketLossWeight(mimeType mime.MimeType, isFecEnabled bool) float64 { return plw } -func toAggregateDeltaInfo(streams map[uint32]*buffer.StreamStatsWithLayers) *rtpstats.RTPDeltaInfo { +func toAggregateDeltaInfo(streams map[uint32]*buffer.StreamStatsWithLayers, useRemoteView bool) *rtpstats.RTPDeltaInfo { deltaInfoList := make([]*rtpstats.RTPDeltaInfo, 0, len(streams)) for _, s := range streams { - deltaInfoList = append(deltaInfoList, s.RTPStats) + if useRemoteView { + if s.RTPStatsRemoteView != nil { + deltaInfoList = append(deltaInfoList, s.RTPStatsRemoteView) + } + } else { + if s.RTPStats != nil { + deltaInfoList = append(deltaInfoList, s.RTPStats) + } + } } return rtpstats.AggregateRTPDeltaInfo(deltaInfoList) } -func toAnalyticsStream(ssrc uint32, deltaStats *rtpstats.RTPDeltaInfo) *livekit.AnalyticsStream { - // discount the feed side loss when reporting forwarded track stats +func toAnalyticsStream( + ssrc uint32, + deltaStats *rtpstats.RTPDeltaInfo, + deltaStatsRemoteView *rtpstats.RTPDeltaInfo, +) *livekit.AnalyticsStream { + // discount the feed side loss when reporting forwarded track stats, packetsLost := deltaStats.PacketsLost - if deltaStats.PacketsMissing > packetsLost { - packetsLost = 0 - } else { - packetsLost -= deltaStats.PacketsMissing + if deltaStatsRemoteView != nil { + packetsLost = deltaStatsRemoteView.PacketsLost + if deltaStatsRemoteView.PacketsMissing > packetsLost { + packetsLost = 0 + } else { + packetsLost -= deltaStatsRemoteView.PacketsMissing + } } return &livekit.AnalyticsStream{ StartTime: timestamppb.New(deltaStats.StartTime), diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index 79f28a805..4a751ac36 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -2370,14 +2370,15 @@ func (d *DownTrack) GetTrackStats() *livekit.RTPStats { return rtpstats.ReconcileRTPStatsWithRTX(d.rtpStats.ToProto(), d.rtpStatsRTX.ToProto()) } -func (d *DownTrack) deltaStats(ds *rtpstats.RTPDeltaInfo) map[uint32]*buffer.StreamStatsWithLayers { - if ds == nil { +func (d *DownTrack) deltaStats(ds *rtpstats.RTPDeltaInfo, dsrv *rtpstats.RTPDeltaInfo) map[uint32]*buffer.StreamStatsWithLayers { + if ds == nil && dsrv == nil { return nil } streamStats := make(map[uint32]*buffer.StreamStatsWithLayers, 1) streamStats[d.ssrc] = &buffer.StreamStatsWithLayers{ - RTPStats: ds, + RTPStats: ds, + RTPStatsRemoteView: dsrv, Layers: map[int32]*rtpstats.RTPDeltaInfo{ 0: ds, }, @@ -2387,11 +2388,11 @@ func (d *DownTrack) deltaStats(ds *rtpstats.RTPDeltaInfo) map[uint32]*buffer.Str } func (d *DownTrack) GetDeltaStatsSender() map[uint32]*buffer.StreamStatsWithLayers { + ds, dsrv := d.rtpStats.DeltaInfoSender(d.deltaStatsSenderSnapshotId) + dsRTX, dsrvRTX := d.rtpStatsRTX.DeltaInfoSender(d.deltaStatsRTXSenderSnapshotId) return d.deltaStats( - rtpstats.ReconcileRTPDeltaInfoWithRTX( - d.rtpStats.DeltaInfoSender(d.deltaStatsSenderSnapshotId), - d.rtpStatsRTX.DeltaInfoSender(d.deltaStatsRTXSenderSnapshotId), - ), + rtpstats.ReconcileRTPDeltaInfoWithRTX(ds, dsRTX), + rtpstats.ReconcileRTPDeltaInfoWithRTX(dsrv, dsrvRTX), ) } diff --git a/pkg/sfu/playoutdelay.go b/pkg/sfu/playoutdelay.go index 5fb70b001..42217f747 100644 --- a/pkg/sfu/playoutdelay.go +++ b/pkg/sfu/playoutdelay.go @@ -119,7 +119,7 @@ func (c *PlayoutDelayController) SeedState(pdcs PlayoutDelayControllerState) { func (c *PlayoutDelayController) SetJitter(jitter uint32) { c.lock.Lock() - deltaInfoSender := c.rtpStats.DeltaInfoSender(c.senderSnapshotID) + deltaInfoSender, _ := c.rtpStats.DeltaInfoSender(c.senderSnapshotID) var nackPercent uint32 if deltaInfoSender != nil && deltaInfoSender.Packets > 0 { nackPercent = deltaInfoSender.Nacks * 100 / deltaInfoSender.Packets diff --git a/pkg/sfu/rtpstats/rtpstats_base.go b/pkg/sfu/rtpstats/rtpstats_base.go index bbf9aa262..15fbf7720 100644 --- a/pkg/sfu/rtpstats/rtpstats_base.go +++ b/pkg/sfu/rtpstats/rtpstats_base.go @@ -134,6 +134,18 @@ func (s *snapshot) MarshalLogObject(e zapcore.ObjectEncoder) error { return nil } +func (s *snapshot) maybeUpdateMaxRTT(rtt uint32) { + if rtt > s.maxRtt { + s.maxRtt = rtt + } +} + +func (s *snapshot) maybeUpdateMaxJitter(jitter float64) { + if jitter > s.maxJitter { + s.maxJitter = jitter + } +} + // ------------------------------------------------------------------ type wrappedRTPDriftLogger struct { @@ -646,10 +658,7 @@ func (r *rtpStatsBase) updateJitter(ets uint64, packetTime int64) float64 { } for i := uint32(0); i < r.nextSnapshotID-cFirstSnapshotID; i++ { - s := &r.snapshots[i] - if r.jitter > s.maxJitter { - s.maxJitter = r.jitter - } + r.snapshots[i].maybeUpdateMaxJitter(r.jitter) } } diff --git a/pkg/sfu/rtpstats/rtpstats_sender.go b/pkg/sfu/rtpstats/rtpstats_sender.go index 648e80049..6c20d84f7 100644 --- a/pkg/sfu/rtpstats/rtpstats_sender.go +++ b/pkg/sfu/rtpstats/rtpstats_sender.go @@ -99,11 +99,11 @@ func (is *intervalStats) MarshalLogObject(e zapcore.ObjectEncoder) error { // ------------------------------------------------------------------- type wrappedReceptionReportsLogger struct { - *senderSnapshot + *senderSnapshotReceiverView } func (w wrappedReceptionReportsLogger) MarshalLogObject(e zapcore.ObjectEncoder) error { - for i, rr := range w.senderSnapshot.processedReceptionReports { + for i, rr := range w.senderSnapshotReceiverView.processedReceptionReports { e.AddReflected(fmt.Sprintf("%d", i), rr) } @@ -112,7 +112,7 @@ func (w wrappedReceptionReportsLogger) MarshalLogObject(e zapcore.ObjectEncoder) // ------------------------------------------------------------------- -type senderSnapshot struct { +type senderSnapshotWindow struct { isValid bool startTime int64 @@ -131,8 +131,7 @@ type senderSnapshot struct { packetsOutOfOrderFeed uint64 - packetsLostFeed uint64 - packetsLostFromRR uint64 + packetsLostFeed uint64 frames uint32 @@ -141,17 +140,10 @@ type senderSnapshot struct { plis uint32 firs uint32 - maxRtt uint32 maxJitterFeed float64 - maxJitter float64 - - extLastRRSN uint64 - intervalStats intervalStats - processedReceptionReports []rtcp.ReceptionReport - metadataCacheOverflowCount int } -func (s *senderSnapshot) MarshalLogObject(e zapcore.ObjectEncoder) error { +func (s *senderSnapshotWindow) MarshalLogObject(e zapcore.ObjectEncoder) error { if s == nil { return nil } @@ -169,13 +161,51 @@ func (s *senderSnapshot) MarshalLogObject(e zapcore.ObjectEncoder) error { e.AddUint64("headerBytesDuplicate", s.headerBytesDuplicate) e.AddUint64("packetsOutOfOrderFeed", s.packetsOutOfOrderFeed) e.AddUint64("packetsLostFeed", s.packetsLostFeed) - e.AddUint64("packetsLostFromRR", s.packetsLostFromRR) e.AddUint32("frames", s.frames) e.AddUint32("nacks", s.nacks) + e.AddUint32("nackRepeated", s.nackRepeated) e.AddUint32("plis", s.plis) e.AddUint32("firs", s.firs) - e.AddUint32("maxRtt", s.maxRtt) e.AddFloat64("maxJitterFeed", s.maxJitterFeed) + return nil +} + +func (s *senderSnapshotWindow) maybeReinit(oldESN uint64, newESN uint64) { + if s.extStartSN == oldESN { + s.extStartSN = newESN + } +} + +func (s *senderSnapshotWindow) maybeUpdateMaxJitterFeed(jitter float64) { + if jitter > s.maxJitterFeed { + s.maxJitterFeed = jitter + } +} + +// --------- + +type senderSnapshotReceiverView struct { + senderSnapshotWindow + + packetsLost uint64 + + maxRtt uint32 + maxJitter float64 + + extLastRRSN uint64 + intervalStats intervalStats + processedReceptionReports []rtcp.ReceptionReport + metadataCacheOverflowCount int +} + +func (s *senderSnapshotReceiverView) MarshalLogObject(e zapcore.ObjectEncoder) error { + if s == nil { + return nil + } + + s.senderSnapshotWindow.MarshalLogObject(e) + e.AddUint64("packetsLost", s.packetsLost) + e.AddUint32("maxRtt", s.maxRtt) e.AddFloat64("maxJitter", s.maxJitter) e.AddUint64("extLastRRSN", s.extLastRRSN) e.AddObject("intervalStats", &s.intervalStats) @@ -184,6 +214,62 @@ func (s *senderSnapshot) MarshalLogObject(e zapcore.ObjectEncoder) error { return nil } +func (s *senderSnapshotReceiverView) maybeReinit(oldESN uint64, newESN uint64) { + if s.extStartSN == oldESN { + s.extStartSN = newESN + if s.extLastRRSN == (oldESN - 1) { + s.extLastRRSN = newESN - 1 + } + } +} + +func (s *senderSnapshotReceiverView) maybeUpdateMaxRTT(rtt uint32) { + if rtt > s.maxRtt { + s.maxRtt = rtt + } +} + +func (s *senderSnapshotReceiverView) maybeUpdateMaxJitter(jitter float64) { + if jitter > s.maxJitter { + s.maxJitter = jitter + } +} + +// --------- + +type senderSnapshot struct { + senderView senderSnapshotWindow + receiverView senderSnapshotReceiverView +} + +func (s *senderSnapshot) MarshalLogObject(e zapcore.ObjectEncoder) error { + if s == nil { + return nil + } + + e.AddObject("senderView", &s.senderView) + e.AddObject("receiverView", &s.receiverView) + return nil +} + +func (s *senderSnapshot) maybeReinit(oldESN uint64, newESN uint64) { + s.senderView.maybeReinit(oldESN, newESN) + s.receiverView.maybeReinit(oldESN, newESN) +} + +func (s *senderSnapshot) maybeUpdateMaxJitterFeed(jitter float64) { + s.senderView.maybeUpdateMaxJitterFeed(jitter) + s.receiverView.maybeUpdateMaxJitterFeed(jitter) +} + +func (s *senderSnapshot) maybeUpdateMaxRTT(rtt uint32) { + s.receiverView.maybeUpdateMaxRTT(rtt) +} + +func (s *senderSnapshot) maybeUpdateMaxJitter(jitter float64) { + s.receiverView.maybeUpdateMaxJitter(jitter) +} + // ------------------------------------------------------------------- type rttMarker struct { @@ -390,13 +476,7 @@ func (r *RTPStatsSender) Update( } } 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 - } - } + r.senderSnapshots[i].maybeReinit(r.extStartSN, extSequenceNumber) } ulgr().Infow( @@ -497,10 +577,7 @@ func (r *RTPStatsSender) Update( jitter := r.updateJitter(extTimestamp, packetTime) for i := uint32(0); i < r.nextSenderSnapshotID-cFirstSnapshotID; i++ { - s := &r.senderSnapshots[i] - if jitter > s.maxJitterFeed { - s.maxJitterFeed = jitter - } + r.senderSnapshots[i].maybeUpdateMaxJitterFeed(jitter) } } } @@ -559,7 +636,7 @@ func (r *RTPStatsSender) UpdateFromReceiverReport(rr rtcp.ReceptionReport) (rtt } extReceivedRRSN := extHighestSNFromRR + (r.extStartSN & 0xFFFF_FFFF_FFFF_0000) - if int64(r.extHighestSN-extReceivedRRSN) > (1 << 16) { + if r.extHighestSNFromRR != extHighestSNFromRR && int64(r.extHighestSN-extReceivedRRSN) > (1<<16) { // there are cases where remote does not send RTCP Receiver Report for extended periods of time, // some times several minutes, in that interval the sequence number rolls over, // @@ -570,6 +647,11 @@ func (r *RTPStatsSender) UpdateFromReceiverReport(rr rtcp.ReceptionReport) (rtt // // catch up till diffrence between highest sent and highest received via receiver report is // less than full 16-bit range. + // + // in a different flavor, there are clients that do not report properly, + // i. e. never update the last received sequence number, + // so skip any catch up if the last receeved sequence number reported in + // RTCP RR does not change. r.logger.Infow( "receiver report missed rollover, adjusting", "timeSinceLastRR", timeSinceLastRR(), @@ -662,28 +744,26 @@ func (r *RTPStatsSender) UpdateFromReceiverReport(rr rtcp.ReceptionReport) (rtt // update snapshots for i := uint32(0); i < r.nextSnapshotID-cFirstSnapshotID; i++ { s := &r.snapshots[i] - if isRttChanged && rtt > s.maxRtt { - s.maxRtt = rtt + if isRttChanged { + s.maybeUpdateMaxRTT(rtt) } } for i := uint32(0); i < r.nextSenderSnapshotID-cFirstSnapshotID; i++ { s := &r.senderSnapshots[i] - if isRttChanged && rtt > s.maxRtt { - s.maxRtt = rtt + if isRttChanged { + s.maybeUpdateMaxRTT(rtt) } - if r.jitterFromRR > s.maxJitter { - s.maxJitter = r.jitterFromRR - } + s.maybeUpdateMaxJitter(r.jitterFromRR) // on every RR, calculate delta since last RR using packet metadata cache - is := r.getIntervalStats(s.extLastRRSN+1, extReceivedRRSN+1, r.extHighestSN) - eis := &s.intervalStats + is := r.getIntervalStats(s.receiverView.extLastRRSN+1, extReceivedRRSN+1, r.extHighestSN) + eis := &s.receiverView.intervalStats eis.aggregate(&is) if is.packetsNotFoundMetadata != 0 { - s.metadataCacheOverflowCount++ - if (s.metadataCacheOverflowCount-1)%10 == 0 { + s.receiverView.metadataCacheOverflowCount++ + if (s.receiverView.metadataCacheOverflowCount-1)%10 == 0 { r.logger.Infow( "metadata cache overflow", "senderSnapshotID", i+cFirstSnapshotID, @@ -691,16 +771,16 @@ func (r *RTPStatsSender) UpdateFromReceiverReport(rr rtcp.ReceptionReport) (rtt "timeSinceLastRR", timeSinceLastRR(), "receivedRR", rr, "extReceivedRRSN", extReceivedRRSN, - "packetsInInterval", extReceivedRRSN-s.extLastRRSN, + "packetsInInterval", extReceivedRRSN-s.receiverView.extLastRRSN, "intervalStats", &is, "aggregateIntervalStats", eis, - "count", s.metadataCacheOverflowCount, + "count", s.receiverView.metadataCacheOverflowCount, "rtpStats", lockedRTPStatsSenderLogEncoder{r}, ) } } - s.extLastRRSN = extReceivedRRSN - s.processedReceptionReports = append(s.processedReceptionReports, rr) + s.receiverView.extLastRRSN = extReceivedRRSN + s.receiverView.processedReceptionReports = append(s.receiverView.processedReceptionReports, rr) } return @@ -860,92 +940,150 @@ func (r *RTPStatsSender) DeltaInfo(snapshotID uint32) *RTPDeltaInfo { return deltaInfo } -func (r *RTPStatsSender) DeltaInfoSender(senderSnapshotID uint32) *RTPDeltaInfo { +func (r *RTPStatsSender) DeltaInfoSender(senderSnapshotID uint32) (*RTPDeltaInfo, *RTPDeltaInfo) { r.lock.Lock() defer r.lock.Unlock() - if r.lastRRTime == 0 { - return nil + + var deltaStatsSenderView *RTPDeltaInfo + thenSenderView, nowSenderView := r.getAndResetSenderSnapshotWindow(senderSnapshotID) + if thenSenderView != nil && nowSenderView != nil { + startTime := thenSenderView.startTime + endTime := nowSenderView.startTime + + packetsExpected := uint32(nowSenderView.extStartSN - thenSenderView.extStartSN) + if packetsExpected > cNumSequenceNumbers { + r.logger.Warnw( + "too many packets expected in delta (sender)", nil, + "senderSnapshotID", senderSnapshotID, + "senderSnapshotNow", nowSenderView, + "senderSnapshotThen", thenSenderView, + "packetsExpected", packetsExpected, + "duration", time.Duration(endTime-startTime), + "rtpStats", lockedRTPStatsSenderLogEncoder{r}, + ) + } else if packetsExpected != 0 { + packetsLostFeed := uint32(nowSenderView.packetsLostFeed - thenSenderView.packetsLostFeed) + if int32(packetsLostFeed) < 0 { + packetsLostFeed = 0 + } + if packetsLostFeed > packetsExpected { + r.logger.Warnw( + "unexpected number of packets lost", nil, + "senderSnapshotID", senderSnapshotID, + "senderSnapshotNow", nowSenderView, + "senderSnapshotThen", thenSenderView, + "packetsExpected", packetsExpected, + "packetsLostFeed", packetsLostFeed, + "duration", time.Duration(endTime-startTime), + "rtpStats", lockedRTPStatsSenderLogEncoder{r}, + ) + packetsLostFeed = packetsExpected + } + + maxJitterTime := thenSenderView.maxJitterFeed / float64(r.params.ClockRate) * 1e6 + + deltaStatsSenderView = &RTPDeltaInfo{ + StartTime: time.Unix(0, startTime), + EndTime: time.Unix(0, endTime), + Packets: packetsExpected - uint32(nowSenderView.packetsPadding-thenSenderView.packetsPadding), + Bytes: nowSenderView.bytes - thenSenderView.bytes, + HeaderBytes: nowSenderView.headerBytes - thenSenderView.headerBytes, + PacketsDuplicate: uint32(nowSenderView.packetsDuplicate - thenSenderView.packetsDuplicate), + BytesDuplicate: nowSenderView.bytesDuplicate - thenSenderView.bytesDuplicate, + HeaderBytesDuplicate: nowSenderView.headerBytesDuplicate - thenSenderView.headerBytesDuplicate, + PacketsPadding: uint32(nowSenderView.packetsPadding - thenSenderView.packetsPadding), + BytesPadding: nowSenderView.bytesPadding - thenSenderView.bytesPadding, + HeaderBytesPadding: nowSenderView.headerBytesPadding - thenSenderView.headerBytesPadding, + PacketsMissing: packetsLostFeed, + PacketsOutOfOrder: uint32(nowSenderView.packetsOutOfOrderFeed - thenSenderView.packetsOutOfOrderFeed), + Frames: nowSenderView.frames - thenSenderView.frames, + JitterMax: maxJitterTime, + Nacks: nowSenderView.nacks - thenSenderView.nacks, + NackRepeated: nowSenderView.nackRepeated - thenSenderView.nackRepeated, + Plis: nowSenderView.plis - thenSenderView.plis, + Firs: nowSenderView.firs - thenSenderView.firs, + } + } } - then, now := r.getAndResetSenderSnapshot(senderSnapshotID) - if now == nil || then == nil { - return nil + var deltaStatsReceiverView *RTPDeltaInfo + if r.lastRRTime != 0 { + thenReceiverView, nowReceiverView := r.getAndResetSenderSnapshotReceiverView(senderSnapshotID) + if thenReceiverView != nil && nowReceiverView != nil { + startTime := thenReceiverView.startTime + endTime := nowReceiverView.startTime + + packetsExpected := uint32(nowReceiverView.extStartSN - thenReceiverView.extStartSN) + if packetsExpected > cNumSequenceNumbers { + r.logger.Warnw( + "too many packets expected in delta (sender - receiver view)", nil, + "senderSnapshotID", senderSnapshotID, + "senderSnapshotNow", nowReceiverView, + "senderSnapshotThen", thenReceiverView, + "packetsExpected", packetsExpected, + "duration", time.Duration(endTime-startTime), + "rtpStats", lockedRTPStatsSenderLogEncoder{r}, + ) + } else if packetsExpected != 0 { + // do not process if no RTCP RR (OR) publisher is not producing any data + packetsLost := uint32(nowReceiverView.packetsLost - thenReceiverView.packetsLost) + if int32(packetsLost) < 0 { + packetsLost = 0 + } + packetsLostFeed := uint32(nowReceiverView.packetsLostFeed - thenReceiverView.packetsLostFeed) + if int32(packetsLostFeed) < 0 { + packetsLostFeed = 0 + } + if packetsLost > packetsExpected { + r.logger.Warnw( + "unexpected number of packets lost (receiver view)", nil, + "senderSnapshotID", senderSnapshotID, + "senderSnapshotNow", nowReceiverView, + "senderSnapshotThen", thenReceiverView, + "packetsExpected", packetsExpected, + "packetsLost", packetsLost, + "packetsLostFeed", packetsLostFeed, + "duration", time.Duration(endTime-startTime), + "rtpStats", lockedRTPStatsSenderLogEncoder{r}, + ) + packetsLost = packetsExpected + } + + // discount jitter from publisher side + internal processing + maxJitter := thenReceiverView.maxJitter - thenReceiverView.maxJitterFeed + if maxJitter < 0.0 { + maxJitter = 0.0 + } + maxJitterTime := maxJitter / float64(r.params.ClockRate) * 1e6 + + deltaStatsReceiverView = &RTPDeltaInfo{ + StartTime: time.Unix(0, startTime), + EndTime: time.Unix(0, endTime), + Packets: packetsExpected - uint32(nowReceiverView.packetsPadding-thenReceiverView.packetsPadding), + Bytes: nowReceiverView.bytes - thenReceiverView.bytes, + HeaderBytes: nowReceiverView.headerBytes - thenReceiverView.headerBytes, + PacketsDuplicate: uint32(nowReceiverView.packetsDuplicate - thenReceiverView.packetsDuplicate), + BytesDuplicate: nowReceiverView.bytesDuplicate - thenReceiverView.bytesDuplicate, + HeaderBytesDuplicate: nowReceiverView.headerBytesDuplicate - thenReceiverView.headerBytesDuplicate, + PacketsPadding: uint32(nowReceiverView.packetsPadding - thenReceiverView.packetsPadding), + BytesPadding: nowReceiverView.bytesPadding - thenReceiverView.bytesPadding, + HeaderBytesPadding: nowReceiverView.headerBytesPadding - thenReceiverView.headerBytesPadding, + PacketsLost: packetsLost, + PacketsMissing: packetsLostFeed, + PacketsOutOfOrder: uint32(nowReceiverView.packetsOutOfOrderFeed - thenReceiverView.packetsOutOfOrderFeed), + Frames: nowReceiverView.frames - thenReceiverView.frames, + RttMax: thenReceiverView.maxRtt, + JitterMax: maxJitterTime, + Nacks: nowReceiverView.nacks - thenReceiverView.nacks, + NackRepeated: nowReceiverView.nackRepeated - thenReceiverView.nackRepeated, + Plis: nowReceiverView.plis - thenReceiverView.plis, + Firs: nowReceiverView.firs - thenReceiverView.firs, + } + } + } } - startTime := then.startTime - endTime := now.startTime - - packetsExpected := uint32(now.extStartSN - then.extStartSN) - if packetsExpected > cNumSequenceNumbers { - r.logger.Warnw( - "too many packets expected in delta (sender)", nil, - "senderSnapshotID", senderSnapshotID, - "senderSnapshotNow", now, - "senderSnapshotThen", then, - "packetsExpected", packetsExpected, - "duration", time.Duration(endTime-startTime), - "rtpStats", lockedRTPStatsSenderLogEncoder{r}, - ) - return nil - } - if packetsExpected == 0 { - // not received RTCP RR (OR) publisher is not producing any data - return nil - } - - packetsLost := uint32(now.packetsLostFromRR - then.packetsLostFromRR) - 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", nil, - "senderSnapshotID", senderSnapshotID, - "senderSnapshotNow", now, - "senderSnapshotThen", then, - "packetsExpected", packetsExpected, - "packetsLost", packetsLost, - "packetsLostFeed", packetsLostFeed, - "duration", time.Duration(endTime-startTime), - "rtpStats", lockedRTPStatsSenderLogEncoder{r}, - ) - packetsLost = packetsExpected - } - - // discount jitter from publisher side + internal processing - maxJitter := then.maxJitter - then.maxJitterFeed - if maxJitter < 0.0 { - maxJitter = 0.0 - } - maxJitterTime := maxJitter / float64(r.params.ClockRate) * 1e6 - - return &RTPDeltaInfo{ - StartTime: time.Unix(0, startTime), - EndTime: time.Unix(0, endTime), - 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(now.packetsPadding - then.packetsPadding), - BytesPadding: now.bytesPadding - then.bytesPadding, - HeaderBytesPadding: now.headerBytesPadding - then.headerBytesPadding, - PacketsLost: packetsLost, - PacketsMissing: packetsLostFeed, - PacketsOutOfOrder: uint32(now.packetsOutOfOrderFeed - then.packetsOutOfOrderFeed), - Frames: now.frames - then.frames, - RttMax: then.maxRtt, - JitterMax: maxJitterTime, - Nacks: now.nacks - then.nacks, - NackRepeated: now.nackRepeated - then.nackRepeated, - Plis: now.plis - then.plis, - Firs: now.firs - then.firs, - } + return deltaStatsSenderView, deltaStatsReceiverView } func (r *RTPStatsSender) MarshalLogObject(e zapcore.ObjectEncoder) error { @@ -980,53 +1118,95 @@ func (r *RTPStatsSender) ToProto() *livekit.RTPStats { return p } -func (r *RTPStatsSender) getAndResetSenderSnapshot(senderSnapshotID uint32) (*senderSnapshot, *senderSnapshot) { +func (r *RTPStatsSender) getAndResetSenderSnapshotWindow(senderSnapshotID uint32) (*senderSnapshotWindow, *senderSnapshotWindow) { + if !r.initialized { + return nil, nil + } + + idx := senderSnapshotID - cFirstSnapshotID + then := r.senderSnapshots[idx] + if !then.senderView.isValid { + then.senderView = initSenderSnapshotWindow(r.startTime, r.extStartSN) + r.senderSnapshots[idx] = then + } + + // snapshot now + r.senderSnapshots[idx].senderView = r.getSenderSnapshotWindow(mono.UnixNano()) + return &then.senderView, &r.senderSnapshots[idx].senderView +} + +func (r *RTPStatsSender) getSenderSnapshotWindow(startTime int64) senderSnapshotWindow { + return senderSnapshotWindow{ + isValid: true, + startTime: startTime, + extStartSN: r.extHighestSN + 1, + bytes: r.bytes, + headerBytes: r.headerBytes, + packetsPadding: r.packetsPadding, + bytesPadding: r.bytesPadding, + headerBytesPadding: r.headerBytesPadding, + packetsDuplicate: r.packetsDuplicate, + bytesDuplicate: r.bytesDuplicate, + headerBytesDuplicate: r.headerBytesDuplicate, + packetsOutOfOrderFeed: r.packetsOutOfOrder, + packetsLostFeed: r.packetsLost, + frames: r.frames, + nacks: r.nacks, + nackRepeated: r.nackRepeated, + plis: r.plis, + firs: r.firs, + maxJitterFeed: r.jitter, + } +} + +func (r *RTPStatsSender) getAndResetSenderSnapshotReceiverView(senderSnapshotID uint32) (*senderSnapshotReceiverView, *senderSnapshotReceiverView) { if !r.initialized || r.lastRRTime == 0 { return nil, nil } idx := senderSnapshotID - cFirstSnapshotID then := r.senderSnapshots[idx] - if !then.isValid { - then = initSenderSnapshot(r.startTime, r.extStartSN) + if !then.receiverView.isValid { + then.receiverView = initSenderSnapshotReceiverView(r.startTime, r.extStartSN) r.senderSnapshots[idx] = then } // snapshot now - now := r.getSenderSnapshot(r.lastRRTime, &then) - r.senderSnapshots[idx] = now - return &then, &now + r.senderSnapshots[idx].receiverView = r.getSenderSnapshotReceiverView(r.lastRRTime, &then.receiverView) + return &then.receiverView, &r.senderSnapshots[idx].receiverView } -func (r *RTPStatsSender) getSenderSnapshot(startTime int64, s *senderSnapshot) senderSnapshot { +func (r *RTPStatsSender) getSenderSnapshotReceiverView(startTime int64, s *senderSnapshotReceiverView) senderSnapshotReceiverView { if s == nil { - return senderSnapshot{} + return senderSnapshotReceiverView{} } - 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, - packetsOutOfOrderFeed: s.packetsOutOfOrderFeed + s.intervalStats.packetsOutOfOrderFeed, - packetsLostFeed: s.packetsLostFeed + s.intervalStats.packetsLostFeed, - packetsLostFromRR: r.packetsLostFromRR, - frames: s.frames + s.intervalStats.frames, - nacks: r.nacks, - nackRepeated: r.nackRepeated, - plis: r.plis, - firs: r.firs, - maxRtt: r.rtt, - maxJitterFeed: r.jitter, - maxJitter: r.jitterFromRR, - extLastRRSN: s.extLastRRSN, + return senderSnapshotReceiverView{ + senderSnapshotWindow: senderSnapshotWindow{ + 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, + packetsOutOfOrderFeed: s.packetsOutOfOrderFeed + s.intervalStats.packetsOutOfOrderFeed, + packetsLostFeed: s.packetsLostFeed + s.intervalStats.packetsLostFeed, + frames: s.frames + s.intervalStats.frames, + nacks: r.nacks, + nackRepeated: r.nackRepeated, + plis: r.plis, + firs: r.firs, + maxJitterFeed: r.jitter, + }, + packetsLost: r.packetsLostFromRR, + maxRtt: r.rtt, + maxJitter: r.jitterFromRR, + extLastRRSN: s.extLastRRSN, } } @@ -1187,9 +1367,26 @@ func (r lockedRTPStatsSenderLogEncoder) MarshalLogObject(e zapcore.ObjectEncoder func initSenderSnapshot(startTime int64, extStartSN uint64) senderSnapshot { return senderSnapshot{ - isValid: true, - startTime: startTime, - extStartSN: extStartSN, + senderView: initSenderSnapshotWindow(startTime, extStartSN), + receiverView: initSenderSnapshotReceiverView(startTime, extStartSN), + } +} + +func initSenderSnapshotWindow(startTime int64, extStartSN uint64) senderSnapshotWindow { + return senderSnapshotWindow{ + isValid: true, + startTime: startTime, + extStartSN: extStartSN, + } +} + +func initSenderSnapshotReceiverView(startTime int64, extStartSN uint64) senderSnapshotReceiverView { + return senderSnapshotReceiverView{ + senderSnapshotWindow: senderSnapshotWindow{ + isValid: true, + startTime: startTime, + extStartSN: extStartSN, + }, extLastRRSN: extStartSN - 1, } }