diff --git a/pkg/sfu/buffer/buffer.go b/pkg/sfu/buffer/buffer.go index f8dc63e06..89f776608 100644 --- a/pkg/sfu/buffer/buffer.go +++ b/pkg/sfu/buffer/buffer.go @@ -132,6 +132,7 @@ type Buffer struct { packetNotFoundCount atomic.Uint32 packetTooOldCount atomic.Uint32 extPacketTooMuchCount atomic.Uint32 + invalidPacketCount atomic.Uint32 primaryBufferForRTX *Buffer rtxPktBuf []byte @@ -317,24 +318,25 @@ func (b *Buffer) Write(pkt []byte) (n int, err error) { } if err = utils.ValidateRTPPacket(&rtpPacket, b.payloadType, b.mediaSSRC); err != nil { - b.logger.Warnw( - "validating RTP packet failed", err, - "version", rtpPacket.Version, - "padding", rtpPacket.Padding, - "marker", rtpPacket.Marker, - "expectedPayloadType", b.payloadType, - "payloadType", rtpPacket.PayloadType, - "sequenceNumber", rtpPacket.SequenceNumber, - "timestamp", rtpPacket.Timestamp, - "expectedSSRC", b.mediaSSRC, - "ssrc", rtpPacket.SSRC, - "numExtensions", len(rtpPacket.Extensions), - "payloadSize", len(rtpPacket.Payload), - "rtpStats", b.rtpStats, - "snRangeMap", b.snRangeMap, - ) - b.Unlock() - return + invalidPacketCount := b.invalidPacketCount.Inc() + if (invalidPacketCount-1)%100 == 0 { + b.logger.Warnw( + "validating RTP packet failed", err, + "version", rtpPacket.Version, + "padding", rtpPacket.Padding, + "marker", rtpPacket.Marker, + "expectedPayloadType", b.payloadType, + "payloadType", rtpPacket.PayloadType, + "sequenceNumber", rtpPacket.SequenceNumber, + "timestamp", rtpPacket.Timestamp, + "expectedSSRC", b.mediaSSRC, + "ssrc", rtpPacket.SSRC, + "numExtensions", len(rtpPacket.Extensions), + "payloadSize", len(rtpPacket.Payload), + "rtpStats", b.rtpStats, + "snRangeMap", b.snRangeMap, + ) + } } now := time.Now() diff --git a/pkg/sfu/utils/helpers.go b/pkg/sfu/utils/helpers.go index f3f12161e..074842777 100644 --- a/pkg/sfu/utils/helpers.go +++ b/pkg/sfu/utils/helpers.go @@ -16,6 +16,7 @@ package utils import ( "errors" + "fmt" "strings" "github.com/pion/interceptor" @@ -54,18 +55,24 @@ func GetHeaderExtensionID(extensions []interceptor.RTPHeaderExtension, extension return 0 } +var ( + ErrInvalidRTPVersion = errors.New("invalid RTP version") + ErrRTPPayloadTypeMismatch = errors.New("RTP payload type mismatch") + ErrRTPSSRCMismatch = errors.New("RTP SSRC mismatch") +) + // ValidateRTPPacket checks for a valid RTP packet and returns an error if fields are incorrect func ValidateRTPPacket(pkt *rtp.Packet, expectedPayloadType uint8, expectedSSRC uint32) error { if pkt.Version != 2 { - return errors.New("invalid RTP version") + return fmt.Errorf("%w, expected: 2, actual: %d", ErrInvalidRTPVersion, pkt.Version) } if expectedPayloadType != 0 && pkt.PayloadType != expectedPayloadType { - return errors.New("invalid RTP payload type") + return fmt.Errorf("%w, expected: %d, actual: %d", ErrRTPPayloadTypeMismatch, expectedPayloadType, pkt.PayloadType) } if expectedSSRC != 0 && pkt.SSRC != expectedSSRC { - return errors.New("invalid RTP SSRC") + return fmt.Errorf("%w, expected: %d, actual: %d", ErrRTPSSRCMismatch, expectedSSRC, pkt.SSRC) } return nil