diff --git a/pkg/sfu/buffer/rtpstats_base.go b/pkg/sfu/buffer/rtpstats_base.go index 11b89b9cc..3dac1c4af 100644 --- a/pkg/sfu/buffer/rtpstats_base.go +++ b/pkg/sfu/buffer/rtpstats_base.go @@ -539,8 +539,8 @@ func (r *rtpStatsBase) deltaInfo(snapshotID uint32, extStartSN uint64, extHighes packetsExpected := now.extStartSN - then.extStartSN if packetsExpected > cNumSequenceNumbers { - r.logger.Errorw( - "too many packets expected in delta", nil, + r.logger.Infow( + "too many packets expected in delta", "startSN", then.extStartSN, "endSN", now.extStartSN, "packetsExpected", packetsExpected, diff --git a/pkg/sfu/buffer/rtpstats_receiver.go b/pkg/sfu/buffer/rtpstats_receiver.go index 164c0928b..1272a7424 100644 --- a/pkg/sfu/buffer/rtpstats_receiver.go +++ b/pkg/sfu/buffer/rtpstats_receiver.go @@ -52,6 +52,9 @@ type RTPStatsReceiver struct { timestamp *utils.WrapAround[uint32, uint64] history *protoutils.Bitmap[uint64] + + clockSkewCount int + outOfOrderSsenderReportCount int } func NewRTPStatsReceiver(params RTPStatsParams) *RTPStatsReceiver { @@ -111,7 +114,7 @@ func (r *RTPStatsReceiver) Update( r.snapshots[i] = r.initSnapshot(r.startTime, r.sequenceNumber.GetExtendedStart()) } - r.logger.Infow( + r.logger.Debugw( "rtp receiver stream start", "startTime", r.startTime.String(), "firstTime", r.firstTime.String(), @@ -318,14 +321,18 @@ func (r *RTPStatsReceiver) SetRtcpSenderReportData(srData *RTCPSenderReportData) if (timeSinceLast > 0.2 && math.Abs(float64(r.params.ClockRate)-calculatedClockRateFromLast) > 0.2*float64(r.params.ClockRate)) || (timeSinceFirst > 0.2 && math.Abs(float64(r.params.ClockRate)-calculatedClockRateFromFirst) > 0.2*float64(r.params.ClockRate)) { - r.logger.Infow( - "clock rate skew", - "first", r.srFirst.ToString(), - "last", r.srNewest.ToString(), - "current", srDataCopy.ToString(), - "calculatedFirst", calculatedClockRateFromFirst, - "calculatedLast", calculatedClockRateFromLast, - ) + if r.clockSkewCount%10 == 0 { + r.logger.Infow( + "clock rate skew", + "first", r.srFirst.ToString(), + "last", r.srNewest.ToString(), + "current", srDataCopy.ToString(), + "calculatedFirst", calculatedClockRateFromFirst, + "calculatedLast", calculatedClockRateFromLast, + "count", r.clockSkewCount, + ) + } + r.clockSkewCount++ } } @@ -334,11 +341,16 @@ func (r *RTPStatsReceiver) SetRtcpSenderReportData(srData *RTCPSenderReportData) // i. e. muting replacing with null and unmute restoring the original track. // Under such a condition reset the sender reports to start from this point. // Resetting will ensure sample rate calculations do not go haywire due to negative time. - r.logger.Infow( - "received sender report, out-of-order, resetting", - "last", r.srNewest.ToString(), - "current", srDataCopy.ToString(), - ) + if r.outOfOrderSsenderReportCount%10 == 0 { + r.logger.Infow( + "received sender report, out-of-order, resetting", + "last", r.srNewest.ToString(), + "current", srDataCopy.ToString(), + "count", r.outOfOrderSsenderReportCount, + ) + } + r.outOfOrderSsenderReportCount++ + r.srFirst = nil } diff --git a/pkg/sfu/buffer/rtpstats_sender.go b/pkg/sfu/buffer/rtpstats_sender.go index 47bf2411b..269e448e2 100644 --- a/pkg/sfu/buffer/rtpstats_sender.go +++ b/pkg/sfu/buffer/rtpstats_sender.go @@ -156,6 +156,10 @@ type RTPStatsSender struct { nextSenderSnapshotID uint32 senderSnapshots []senderSnapshot + + clockSkewCount int + outOfOrderSenderReportCount int + metadataCacheOverflowCount int } func NewRTPStatsSender(params RTPStatsParams) *RTPStatsSender { @@ -265,7 +269,7 @@ func (r *RTPStatsSender) Update( r.senderSnapshots[i] = r.initSenderSnapshot(r.startTime, r.extStartSN) } - r.logger.Infow( + r.logger.Debugw( "rtp sender stream start", "startTime", r.startTime.String(), "firstTime", r.firstTime.String(), @@ -510,22 +514,26 @@ func (r *RTPStatsSender) UpdateFromReceiverReport(rr rtcp.ReceptionReport) (rtt eis := &s.intervalStats eis.aggregate(&is) if is.packetsNotFound != 0 { - r.logger.Warnw( - "potential sequence number de-sync", nil, - "lastRRTime", r.lastRRTime.String(), - "lastRR", r.lastRR, - "sinceLastRR", time.Since(r.lastRRTime).String(), - "receivedRR", rr, - "extStartSN", r.extStartSN, - "extHighestSN", r.extHighestSN, - "extLastRRSN", s.extLastRRSN, - "extReceivedRRSN", extReceivedRRSN, - "packetsInInterval", extReceivedRRSN-s.extLastRRSN, - "intervalStats", is.ToString(), - "aggregateIntervalStats", eis.ToString(), - "extHighestSNFromRR", r.extHighestSNFromRR, - "packetsLostFromRR", r.packetsLostFromRR, - ) + if r.metadataCacheOverflowCount%10 == 0 { + r.logger.Infow( + "metadata cache overflow", + "lastRRTime", r.lastRRTime.String(), + "lastRR", r.lastRR, + "sinceLastRR", time.Since(r.lastRRTime).String(), + "receivedRR", rr, + "extStartSN", r.extStartSN, + "extHighestSN", r.extHighestSN, + "extLastRRSN", s.extLastRRSN, + "extReceivedRRSN", extReceivedRRSN, + "packetsInInterval", extReceivedRRSN-s.extLastRRSN, + "intervalStats", is.ToString(), + "aggregateIntervalStats", eis.ToString(), + "extHighestSNFromRR", r.extHighestSNFromRR, + "packetsLostFromRR", r.packetsLostFromRR, + "count", r.metadataCacheOverflowCount, + ) + } + r.metadataCacheOverflowCount++ } s.extLastRRSN = extReceivedRRSN } @@ -605,24 +613,28 @@ func (r *RTPStatsSender) GetRtcpSenderReport(ssrc uint32, calculatedClockRate ui rtpDiffSinceLastReport := nowRTPExt - r.srNewest.RTPTimestampExt windowClockRate := float64(rtpDiffSinceLastReport) / timeSinceLastReport if timeSinceLastReport > 0.2 && math.Abs(float64(r.params.ClockRate)-windowClockRate) > 0.2*float64(r.params.ClockRate) { - r.logger.Infow( - "sending sender report, clock skew", - "last", r.srNewest.ToString(), - "curr", srData.ToString(), - "timeNow", time.Now().String(), - "extStartTS", r.extStartTS, - "extHighestTS", r.extHighestTS, - "highestTime", r.highestTime.String(), - "timeSinceHighest", timeSinceHighest.String(), - "firstTime", r.firstTime.String(), - "timeSinceFirst", timeSinceFirst.String(), - "nowRTPExtUsingTime", nowRTPExtUsingTime, - "calculatedClockRate", calculatedClockRate, - "nowRTPExtUsingRate", nowRTPExtUsingRate, - "timeSinceLastReport", timeSinceLastReport, - "rtpDiffSinceLastReport", rtpDiffSinceLastReport, - "windowClockRate", windowClockRate, - ) + if r.clockSkewCount%10 == 0 { + r.logger.Infow( + "sending sender report, clock skew", + "last", r.srNewest.ToString(), + "curr", srData.ToString(), + "timeNow", time.Now().String(), + "extStartTS", r.extStartTS, + "extHighestTS", r.extHighestTS, + "highestTime", r.highestTime.String(), + "timeSinceHighest", timeSinceHighest.String(), + "firstTime", r.firstTime.String(), + "timeSinceFirst", timeSinceFirst.String(), + "nowRTPExtUsingTime", nowRTPExtUsingTime, + "calculatedClockRate", calculatedClockRate, + "nowRTPExtUsingRate", nowRTPExtUsingRate, + "timeSinceLastReport", timeSinceLastReport, + "rtpDiffSinceLastReport", rtpDiffSinceLastReport, + "windowClockRate", windowClockRate, + "count", r.clockSkewCount, + ) + } + r.clockSkewCount++ } } @@ -638,21 +650,26 @@ func (r *RTPStatsSender) GetRtcpSenderReport(ssrc uint32, calculatedClockRate ui // result in this module not having calculated clock rate of publisher side. // - When the above happens, current will be generated using highestTS which could be behind. // That could end up behind the last report's timestamp in extreme cases - r.logger.Infow( - "sending sender report, out-of-order, repairing", - "last", r.srNewest.ToString(), - "curr", srData.ToString(), - "timeNow", time.Now().String(), - "extStartTS", r.extStartTS, - "extHighestTS", r.extHighestTS, - "highestTime", r.highestTime.String(), - "timeSinceHighest", timeSinceHighest.String(), - "firstTime", r.firstTime.String(), - "timeSinceFirst", timeSinceFirst.String(), - "nowRTPExtUsingTime", nowRTPExtUsingTime, - "calculatedClockRate", calculatedClockRate, - "nowRTPExtUsingRate", nowRTPExtUsingRate, - ) + if r.outOfOrderSenderReportCount%10 == 0 { + r.logger.Infow( + "sending sender report, out-of-order, repairing", + "last", r.srNewest.ToString(), + "curr", srData.ToString(), + "timeNow", time.Now().String(), + "extStartTS", r.extStartTS, + "extHighestTS", r.extHighestTS, + "highestTime", r.highestTime.String(), + "timeSinceHighest", timeSinceHighest.String(), + "firstTime", r.firstTime.String(), + "timeSinceFirst", timeSinceFirst.String(), + "nowRTPExtUsingTime", nowRTPExtUsingTime, + "calculatedClockRate", calculatedClockRate, + "nowRTPExtUsingRate", nowRTPExtUsingRate, + "count", r.outOfOrderSenderReportCount, + ) + } + r.outOfOrderSenderReportCount++ + ntpDiffSinceLast := nowNTP.Time().Sub(r.srNewest.NTPTimestamp.Time()) nowRTPExt = r.srNewest.RTPTimestampExt + uint64(ntpDiffSinceLast.Seconds()*float64(r.params.ClockRate)) nowRTP = uint32(nowRTPExt) diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index b26080526..880f047e8 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -672,7 +672,7 @@ func (d *DownTrack) WriteRTP(extPkt *buffer.ExtPacket, layer int32) error { tp, err := d.forwarder.GetTranslationParams(extPkt, layer) if tp.shouldDrop { if err != nil { - d.params.Logger.Errorw("write rtp packet failed", err) + d.params.Logger.Errorw("could not get translation params", err) } return err } @@ -692,7 +692,7 @@ func (d *DownTrack) WriteRTP(extPkt *buffer.ExtPacket, layer int32) error { hdr, err := d.getTranslatedRTPHeader(extPkt, tp) if err != nil { - d.params.Logger.Errorw("write rtp packet failed", err) + d.params.Logger.Errorw("could not get translated RTP header", err) if poolEntity != nil { PacketFactory.Put(poolEntity) } @@ -1447,7 +1447,7 @@ func (d *DownTrack) getH264BlankFrame(_frameEndNeeded bool) ([]byte, error) { func (d *DownTrack) handleRTCP(bytes []byte) { pkts, err := rtcp.Unmarshal(bytes) if err != nil { - d.params.Logger.Errorw("unmarshal rtcp receiver packets err", err) + d.params.Logger.Errorw("could not unmarshal rtcp receiver packets", err) return } @@ -1611,7 +1611,7 @@ func (d *DownTrack) retransmitPackets(nacks []uint16) { var pkt rtp.Packet if err = pkt.Unmarshal(pktBuff[:n]); err != nil { - d.params.Logger.Errorw("unmarshalling rtp packet failed in retransmit", err) + d.params.Logger.Errorw("could not unmarshal rtp packet in retransmit", err) continue } pkt.Header.Marker = epm.marker @@ -1625,7 +1625,7 @@ func (d *DownTrack) retransmitPackets(nacks []uint16) { if d.mime == "video/vp8" && len(pkt.Payload) > 0 && len(epm.codecBytes) != 0 { var incomingVP8 buffer.VP8 if err = incomingVP8.Unmarshal(pkt.Payload); err != nil { - d.params.Logger.Errorw("unmarshalling VP8 packet err", err) + d.params.Logger.Errorw("could not unmarshal VP8 packet", err) PacketFactory.Put(poolEntity) continue } diff --git a/pkg/sfu/forwarder.go b/pkg/sfu/forwarder.go index ea6ea29c1..f37bea03f 100644 --- a/pkg/sfu/forwarder.go +++ b/pkg/sfu/forwarder.go @@ -1581,14 +1581,7 @@ func (f *Forwarder) processSourceSwitch(extPkt *buffer.ExtPacket, layer int32) e extNextTS = extExpectedTS } else if diffSeconds > ResumeBehindHighTresholdSeconds { // could be due to incorrect reference calculation - f.logger.Infow( - "resume, reference very far behind", - "layer", layer, - "extExpectedTS", extExpectedTS, - "extRefTS", extRefTS, - "extLastTS", extLastTS, - "diffSeconds", diffSeconds, - ) + logTransition("resume, reference very far behind", extExpectedTS, extRefTS, extLastTS, diffSeconds) extNextTS = extExpectedTS } else { extNextTS = extRefTS