diff --git a/pkg/sfu/buffer/buffer.go b/pkg/sfu/buffer/buffer.go index 2b82f3760..e48692ba2 100644 --- a/pkg/sfu/buffer/buffer.go +++ b/pkg/sfu/buffer/buffer.go @@ -551,7 +551,13 @@ func (b *Buffer) calc(rawPkt []byte, rtpPacket *rtp.Packet, arrivalTime time.Tim // 44 - padding only - out-of-order + duplicate - dropped as duplicate // if err := b.snRangeMap.ExcludeRange(flowState.ExtSequenceNumber, flowState.ExtSequenceNumber+1); err != nil { - b.logger.Errorw("could not exclude range", err, "sn", rtpPacket.SequenceNumber, "esn", flowState.ExtSequenceNumber) + b.logger.Errorw( + "could not exclude range", err, + "sn", rtpPacket.SequenceNumber, + "esn", flowState.ExtSequenceNumber, + "rtpStats", b.rtpStats, + "snRangeMap", b.snRangeMap, + ) } } return @@ -560,7 +566,14 @@ func (b *Buffer) calc(rawPkt []byte, rtpPacket *rtp.Packet, arrivalTime time.Tim // add to RTX buffer using sequence number after accounting for dropped padding only packets snAdjustment, err := b.snRangeMap.GetValue(flowState.ExtSequenceNumber) if err != nil { - b.logger.Errorw("could not get sequence number adjustment", err, "sn", flowState.ExtSequenceNumber, "payloadSize", len(rtpPacket.Payload)) + b.logger.Errorw( + "could not get sequence number adjustment", err, + "sn", rtpPacket.SequenceNumber, + "esn", flowState.ExtSequenceNumber, + "payloadSize", len(rtpPacket.Payload), + "rtpStats", b.rtpStats, + "snRangeMap", b.snRangeMap, + ) return } flowState.ExtSequenceNumber -= snAdjustment @@ -577,10 +590,18 @@ func (b *Buffer) calc(rawPkt []byte, rtpPacket *rtp.Packet, arrivalTime time.Tim "snAdjustment", snAdjustment, "incomingSequenceNumber", flowState.ExtSequenceNumber+snAdjustment, "rtpStats", b.rtpStats, + "snRangeMap", b.snRangeMap, ) } } else if err != bucket.ErrRTXPacket { - b.logger.Warnw("could not add packet to bucket", err) + b.logger.Warnw( + "could not add packet to bucket", err, + "flowState", &flowState, + "snAdjustment", snAdjustment, + "incomingSequenceNumber", flowState.ExtSequenceNumber+snAdjustment, + "rtpStats", b.rtpStats, + "snRangeMap", b.snRangeMap, + ) } return } @@ -605,7 +626,14 @@ func (b *Buffer) patchExtPacket(ep *ExtPacket, buf []byte) *ExtPacket { if err != nil { packetNotFoundCount := b.packetNotFoundCount.Inc() if (packetNotFoundCount-1)%20 == 0 { - b.logger.Warnw("could not get packet from bucket", err, "sn", ep.Packet.SequenceNumber, "headSN", b.bucket.HeadSequenceNumber(), "count", packetNotFoundCount) + b.logger.Warnw( + "could not get packet from bucket", err, + "sn", ep.Packet.SequenceNumber, + "headSN", b.bucket.HeadSequenceNumber(), + "count", packetNotFoundCount, + "rtpStats", b.rtpStats, + "snRangeMap", b.snRangeMap, + ) } return nil } diff --git a/pkg/sfu/utils/rangemap.go b/pkg/sfu/utils/rangemap.go index fe5480832..6ca2332eb 100644 --- a/pkg/sfu/utils/rangemap.go +++ b/pkg/sfu/utils/rangemap.go @@ -19,6 +19,8 @@ import ( "fmt" "math" "unsafe" + + "go.uber.org/zap/zapcore" ) const ( @@ -40,12 +42,23 @@ type valueType interface { uint32 | uint64 } +// --------------------------------------------------- + type rangeVal[RT rangeType, VT valueType] struct { start RT end RT value VT } +func (r rangeVal[RT, VT]) MarshalLogObject(e zapcore.ObjectEncoder) error { + e.AddUint64("start", uint64(r.start)) + e.AddUint64("end", uint64(r.end)) + e.AddUint64("value", uint64(r.value)) + return nil +} + +// --------------------------------------------------- + type RangeMap[RT rangeType, VT valueType] struct { halfRange RT @@ -63,6 +76,21 @@ func NewRangeMap[RT rangeType, VT valueType](size int) *RangeMap[RT, VT] { return r } +func (r *RangeMap[RT, VT]) MarshalLogObject(e zapcore.ObjectEncoder) error { + e.AddInt("numRanges", len(r.ranges)) + + // just the last 10 ranges max + startIdx := len(r.ranges) - 10 + if startIdx < 0 { + startIdx = 0 + } + for i := startIdx; i < len(r.ranges); i++ { + e.AddObject(fmt.Sprintf("range[%d]", i), r.ranges[i]) + } + + return nil +} + func (r *RangeMap[RT, VT]) ClearAndResetValue(start RT, val VT) { r.initRanges(start, val) }