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)
This commit is contained in:
Raja Subramanian
2025-02-11 16:05:00 +05:30
committed by GitHub
parent 5e1431f433
commit 7fef374b19
6 changed files with 413 additions and 183 deletions
+2
View File
@@ -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
}
+35 -14
View File
@@ -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),
+8 -7
View File
@@ -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),
)
}
+1 -1
View File
@@ -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
+13 -4
View File
@@ -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)
}
}
+354 -157
View File
@@ -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,
}
}