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, } }