diff --git a/go.mod b/go.mod index c09719e8c..60cf8b8b1 100644 --- a/go.mod +++ b/go.mod @@ -19,7 +19,7 @@ require ( github.com/jxskiss/base62 v1.1.0 github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1 github.com/livekit/mediatransportutil v0.0.0-20240730083616-559fa5ece598 - github.com/livekit/protocol v1.23.1-0.20240926100400-c7f09c48faae + github.com/livekit/protocol v1.23.1-0.20241003052959-b216a9275d12 github.com/livekit/psrpc v0.6.1-0.20240924010758-9f0a4268a3b9 github.com/mackerelio/go-osstat v0.2.5 github.com/magefile/mage v1.15.0 diff --git a/go.sum b/go.sum index 9ad7fd595..baff4dde1 100644 --- a/go.sum +++ b/go.sum @@ -165,8 +165,8 @@ github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1 h1:jm09419p0lqTkD github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1/go.mod h1:Rs3MhFwutWhGwmY1VQsygw28z5bWcnEYmS1OG9OxjOQ= github.com/livekit/mediatransportutil v0.0.0-20240730083616-559fa5ece598 h1:yLlkHk2feSLHstD9n4VKg7YEBR4rLODTI4WE8gNBEnQ= github.com/livekit/mediatransportutil v0.0.0-20240730083616-559fa5ece598/go.mod h1:jwKUCmObuiEDH0iiuJHaGMXwRs3RjrB4G6qqgkr/5oE= -github.com/livekit/protocol v1.23.1-0.20240926100400-c7f09c48faae h1:ggrNio5LP5WO/iJMig/bHvxGY/7Li+tEVDK/FOmC5kc= -github.com/livekit/protocol v1.23.1-0.20240926100400-c7f09c48faae/go.mod h1:nxRzmQBKSYK64gqr7ABWwt78hvrgiO2wYuCojRYb7Gs= +github.com/livekit/protocol v1.23.1-0.20241003052959-b216a9275d12 h1:qpPbhqJmNAoNeO39MCmmXqOrd1F8skmtT7xq/pfcVrM= +github.com/livekit/protocol v1.23.1-0.20241003052959-b216a9275d12/go.mod h1:nxRzmQBKSYK64gqr7ABWwt78hvrgiO2wYuCojRYb7Gs= github.com/livekit/psrpc v0.6.1-0.20240924010758-9f0a4268a3b9 h1:33oBjGpVD9tYkDXQU42tnHl8eCX9G6PVUToBVuCUyOs= github.com/livekit/psrpc v0.6.1-0.20240924010758-9f0a4268a3b9/go.mod h1:CQUBSPfYYAaevg1TNCc6/aYsa8DJH4jSRFdCeSZk5u0= github.com/mackerelio/go-osstat v0.2.5 h1:+MqTbZUhoIt4m8qzkVoXUJg1EuifwlAJSk4Yl2GXh+o= diff --git a/pkg/sfu/buffer/rtpstats_receiver.go b/pkg/sfu/buffer/rtpstats_receiver.go index eb03b495a..669f609b6 100644 --- a/pkg/sfu/buffer/rtpstats_receiver.go +++ b/pkg/sfu/buffer/rtpstats_receiver.go @@ -151,26 +151,24 @@ func (r *RTPStatsReceiver) Update( var tsRolloverCount int var snRolloverCount int - getLoggingFields := func() []interface{} { - return []interface{}{ - "resSN", resSN, - "gapSN", gapSN, - "resTS", resTS, - "gapTS", int64(resTS.ExtendedVal - resTS.PreExtendedHighest), - "timeSinceHighest", time.Duration(timeSinceHighest), - "snRolloverCount", snRolloverCount, - "expectedTSJump", expectedTSJump, - "tsRolloverCount", tsRolloverCount, - "packetTime", time.Unix(0, packetTime).String(), - "sequenceNumber", sequenceNumber, - "timestamp", timestamp, - "marker", marker, - "hdrSize", hdrSize, - "payloadSize", payloadSize, - "paddingSize", paddingSize, - "rtpStats", lockedRTPStatsReceiverLogEncoder{r}, - } - } + logger := r.logger.WithUnlikelyValues( + "resSN", resSN, + "gapSN", gapSN, + "resTS", resTS, + "gapTS", int64(resTS.ExtendedVal-resTS.PreExtendedHighest), + "timeSinceHighest", time.Duration(timeSinceHighest), + "snRolloverCount", snRolloverCount, + "expectedTSJump", expectedTSJump, + "tsRolloverCount", tsRolloverCount, + "packetTime", time.Unix(0, packetTime).String(), + "sequenceNumber", sequenceNumber, + "timestamp", timestamp, + "marker", marker, + "hdrSize", hdrSize, + "payloadSize", payloadSize, + "paddingSize", paddingSize, + "rtpStats", lockedRTPStatsReceiverLogEncoder{r}, + ) if !r.initialized { if payloadSize == 0 { @@ -209,10 +207,7 @@ func (r *RTPStatsReceiver) Update( timeSinceHighest = packetTime - r.highestTime tsRolloverCount = r.getTSRolloverCount(timeSinceHighest, timestamp) if tsRolloverCount >= 0 { - r.logger.Warnw( - "potential time stamp roll over", nil, - getLoggingFields()..., - ) + logger.Warnw("potential time stamp roll over", nil) } resTS = r.timestamp.Rollover(timestamp, tsRolloverCount) if resTS.IsUnhandled { @@ -236,10 +231,7 @@ func (r *RTPStatsReceiver) Update( if gapTS > int64(float64(expectedTSJump)*cTSJumpTooHighFactor) { r.sequenceNumber.UndoUpdate(resSN) r.timestamp.UndoUpdate(resTS) - r.logger.Warnw( - "dropping old packet, timestamp", nil, - getLoggingFields()..., - ) + logger.Warnw("dropping old packet, timestamp", nil) flowState.IsNotHandled = true return } @@ -250,10 +242,7 @@ func (r *RTPStatsReceiver) Update( if gapTS < 0 && gapSN > 0 { r.sequenceNumber.UndoUpdate(resSN) r.timestamp.UndoUpdate(resTS) - r.logger.Warnw( - "dropping old packet, sequence number", nil, - getLoggingFields()..., - ) + logger.Warnw("dropping old packet, sequence number", nil) flowState.IsNotHandled = true return } @@ -272,10 +261,7 @@ func (r *RTPStatsReceiver) Update( return } - r.logger.Warnw( - "forcing sequence number rollover", nil, - getLoggingFields()..., - ) + logger.Warnw("forcing sequence number rollover", nil) } } gapSN = int64(resSN.ExtendedVal - resSN.PreExtendedHighest) @@ -303,9 +289,9 @@ func (r *RTPStatsReceiver) Update( if !flowState.IsDuplicate && -gapSN >= cSequenceNumberLargeJumpThreshold { r.largeJumpNegativeCount++ if (r.largeJumpNegativeCount-1)%100 == 0 { - r.logger.Warnw( + logger.Warnw( "large sequence number gap negative", nil, - append(getLoggingFields(), "count", r.largeJumpNegativeCount)..., + "count", r.largeJumpNegativeCount, ) } } @@ -313,9 +299,9 @@ func (r *RTPStatsReceiver) Update( if gapSN >= cSequenceNumberLargeJumpThreshold { r.largeJumpCount++ if (r.largeJumpCount-1)%100 == 0 { - r.logger.Warnw( + logger.Warnw( "large sequence number gap", nil, - append(getLoggingFields(), "count", r.largeJumpCount)..., + "count", r.largeJumpCount, ) } } @@ -323,9 +309,9 @@ func (r *RTPStatsReceiver) Update( if resTS.ExtendedVal < resTS.PreExtendedHighest { r.timeReversedCount++ if (r.timeReversedCount-1)%100 == 0 { - r.logger.Warnw( + logger.Warnw( "time reversed", nil, - append(getLoggingFields(), "count", r.timeReversedCount)..., + "count", r.timeReversedCount, ) } } diff --git a/pkg/sfu/buffer/rtpstats_sender.go b/pkg/sfu/buffer/rtpstats_sender.go index ce00d7bef..1f9715e96 100644 --- a/pkg/sfu/buffer/rtpstats_sender.go +++ b/pkg/sfu/buffer/rtpstats_sender.go @@ -292,20 +292,18 @@ func (r *RTPStatsSender) Update( pktSize := uint64(hdrSize + payloadSize + paddingSize) isDuplicate := false gapSN := int64(extSequenceNumber - r.extHighestSN) - getLoggingFields := func() []interface{} { - return []interface{}{ - "currSN", extSequenceNumber, - "gapSN", gapSN, - "currTS", extTimestamp, - "gapTS", int64(extTimestamp - r.extHighestTS), - "packetTime", packetTime, - "marker", marker, - "hdrSize", hdrSize, - "payloadSize", payloadSize, - "paddingSize", paddingSize, - "rtpStats", lockedRTPStatsSenderLogEncoder{r}, - } - } + logger := r.logger.WithUnlikelyValues( + "currSN", extSequenceNumber, + "gapSN", gapSN, + "currTS", extTimestamp, + "gapTS", int64(extTimestamp-r.extHighestTS), + "packetTime", packetTime, + "marker", marker, + "hdrSize", hdrSize, + "payloadSize", payloadSize, + "paddingSize", paddingSize, + "rtpStats", lockedRTPStatsSenderLogEncoder{r}, + ) if gapSN <= 0 { // duplicate OR out-of-order if payloadSize == 0 && extSequenceNumber < r.extStartSN { // do not start on a padding only packet @@ -332,12 +330,10 @@ func (r *RTPStatsSender) Update( } } - r.logger.Infow( + logger.Infow( "adjusting start sequence number", - append(getLoggingFields(), - "snAfter", extSequenceNumber, - "tsAfter", extTimestamp, - )..., + "snAfter", extSequenceNumber, + "tsAfter", extTimestamp, ) r.extStartSN = extSequenceNumber } @@ -359,9 +355,9 @@ func (r *RTPStatsSender) Update( if !isDuplicate && -gapSN >= cSequenceNumberLargeJumpThreshold { r.largeJumpNegativeCount++ if (r.largeJumpNegativeCount-1)%100 == 0 { - r.logger.Warnw( + logger.Warnw( "large sequence number gap negative", nil, - append(getLoggingFields(), "count", r.largeJumpNegativeCount)..., + "count", r.largeJumpNegativeCount, ) } } @@ -369,9 +365,9 @@ func (r *RTPStatsSender) Update( if gapSN >= cSequenceNumberLargeJumpThreshold { r.largeJumpCount++ if (r.largeJumpCount-1)%100 == 0 { - r.logger.Warnw( + logger.Warnw( "large sequence number gap", nil, - append(getLoggingFields(), "count", r.largeJumpCount)..., + "count", r.largeJumpCount, ) } } @@ -379,9 +375,9 @@ func (r *RTPStatsSender) Update( if extTimestamp < r.extHighestTS { r.timeReversedCount++ if (r.timeReversedCount-1)%100 == 0 { - r.logger.Warnw( + logger.Warnw( "time reversed", nil, - append(getLoggingFields(), "count", r.timeReversedCount)..., + "count", r.timeReversedCount, ) } } @@ -399,12 +395,10 @@ func (r *RTPStatsSender) Update( } if extTimestamp < r.extStartTS { - r.logger.Infow( + logger.Infow( "adjusting start timestamp", - append(getLoggingFields(), - "snAfter", extSequenceNumber, - "tsAfter", extTimestamp, - )..., + "snAfter", extSequenceNumber, + "tsAfter", extTimestamp, ) r.extStartTS = extTimestamp } @@ -647,21 +641,20 @@ func (r *RTPStatsSender) GetRtcpSenderReport(ssrc uint32, publisherSRData *livek Octets: octetCount, } - getFields := func() []interface{} { - return []interface{}{ - "curr", WrappedRTCPSenderReportStateLogger{srData}, - "feed", WrappedRTCPSenderReportStateLogger{publisherSRData}, - "tsOffset", tsOffset, - "timeNow", time.Now().String(), - "now", time.Unix(0, now).String(), - "timeSinceHighest", time.Unix(0, now).Sub(time.Unix(0, r.highestTime)).String(), - "timeSinceFirst", time.Unix(0, now).Sub(time.Unix(0, r.firstTime)).String(), - "timeSincePublisherSRAdjusted", timeSincePublisherSRAdjusted.String(), - "timeSincePublisherSR", time.Since(time.Unix(0, publisherSRData.At)).String(), - "nowRTPExt", nowRTPExt, - "rtpStats", lockedRTPStatsSenderLogEncoder{r}, - } - } + logger := r.logger.WithUnlikelyValues( + "curr", WrappedRTCPSenderReportStateLogger{srData}, + "feed", WrappedRTCPSenderReportStateLogger{publisherSRData}, + "tsOffset", tsOffset, + "timeNow", time.Now().String(), + "now", time.Unix(0, now).String(), + "timeSinceHighest", time.Unix(0, now).Sub(time.Unix(0, r.highestTime)).String(), + "timeSinceFirst", time.Unix(0, now).Sub(time.Unix(0, r.firstTime)).String(), + "timeSincePublisherSRAdjusted", timeSincePublisherSRAdjusted.String(), + "timeSincePublisherSR", time.Since(time.Unix(0, publisherSRData.At)).String(), + "nowRTPExt", nowRTPExt, + "rtpStats", lockedRTPStatsSenderLogEncoder{r}, + ) + if r.srNewest != nil && nowRTPExt >= r.srNewest.RtpTimestampExt { timeSinceLastReport := nowNTP.Time().Sub(mediatransportutil.NtpTime(r.srNewest.NtpTimestamp).Time()) rtpDiffSinceLastReport := nowRTPExt - r.srNewest.RtpTimestampExt @@ -669,14 +662,13 @@ func (r *RTPStatsSender) GetRtcpSenderReport(ssrc uint32, publisherSRData *livek if timeSinceLastReport.Seconds() > 0.2 && math.Abs(float64(r.params.ClockRate)-windowClockRate) > 0.2*float64(r.params.ClockRate) { r.clockSkewCount++ if (r.clockSkewCount-1)%100 == 0 { - fields := append( - getFields(), + logger.Infow( + "sending sender report, clock skew", "timeSinceLastReport", timeSinceLastReport.String(), "rtpDiffSinceLastReport", rtpDiffSinceLastReport, "windowClockRate", windowClockRate, "count", r.clockSkewCount, ) - r.logger.Infow("sending sender report, clock skew", fields...) } } } @@ -684,7 +676,7 @@ func (r *RTPStatsSender) GetRtcpSenderReport(ssrc uint32, publisherSRData *livek if r.srNewest != nil && nowRTPExt < r.srNewest.RtpTimestampExt { // If report being generated is behind the last report, skip it. // Should not happen. - r.logger.Infow("sending sender report, out-of-order, skipping", getFields()...) + logger.Infow("sending sender report, out-of-order, skipping") return nil }