diff --git a/pkg/sfu/buffer/buffer.go b/pkg/sfu/buffer/buffer.go index b94cb0265..3401d603a 100644 --- a/pkg/sfu/buffer/buffer.go +++ b/pkg/sfu/buffer/buffer.go @@ -440,9 +440,12 @@ func (b *Buffer) calc(pkt []byte, arrivalTime time.Time) { // 45 - regular packet - offset = 2 (running offset) - passed through with adjusted sequence number as 43 // 44 - padding only - out-of-order + duplicate - dropped as duplicate // + b.logger.Debugw("PDBG dropping padding", "sn", rtpPacket.Header.SequenceNumber, "esn", flowState.ExtSequenceNumber) // REMOVE 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) } + } else { + b.logger.Debugw("PDBG dropping duplicate padding", "sn", rtpPacket.Header.SequenceNumber, "esn", flowState.ExtSequenceNumber) // REMOVE } return } @@ -453,6 +456,7 @@ func (b *Buffer) calc(pkt []byte, arrivalTime time.Time) { b.logger.Errorw("could not get sequence number adjustment", err, "sn", flowState.ExtSequenceNumber, "payloadSize", len(rtpPacket.Payload)) return } + b.logger.Debugw("PDBG incoming packet", "sn", rtpPacket.Header.SequenceNumber, "esn", flowState.ExtSequenceNumber, "adjust", snAdjustment, "size", len(rtpPacket.Payload)) // REMOVE flowState.ExtSequenceNumber -= snAdjustment rtpPacket.Header.SequenceNumber = uint16(flowState.ExtSequenceNumber) _, err = b.bucket.AddPacketWithSequenceNumber(pkt, rtpPacket.Header.SequenceNumber) @@ -494,6 +498,7 @@ func (b *Buffer) patchExtPacket(ep *ExtPacket, buf []byte) *ExtPacket { } pkt.Payload = buf[payloadStart:payloadEnd] ep.Packet = &pkt + b.logger.Debugw("PDBG forwarding packet", "sn", pkt.Header.SequenceNumber, "esn", ep.ExtSequenceNumber, "size", len(ep.Packet.Payload)) // REMOVE return ep } diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index 44414d7c8..5a3328d26 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -684,7 +684,7 @@ func (d *DownTrack) WriteRTP(extPkt *buffer.ExtPacket, layer int32) error { } } if d.sequencer != nil { - d.sequencer.push( + isPushed := d.sequencer.push( extPkt.Packet.SequenceNumber, tp.rtp.sequenceNumber, tp.rtp.timestamp, @@ -693,8 +693,12 @@ func (d *DownTrack) WriteRTP(extPkt *buffer.ExtPacket, layer int32) error { tp.codecBytes, tp.ddBytes, ) + if !isPushed { + d.params.Logger.Debugw("PDBG sequencer push failed", "isn", extPkt.Packet.SequenceNumber, "iesn", extPkt.ExtSequenceNumber, "layer", layer, "osn", tp.rtp.sequenceNumber, "size", len(payload)) // REMOVE + } } + d.params.Logger.Debugw("PDBG forwarding packet", "isn", extPkt.Packet.SequenceNumber, "iesn", extPkt.ExtSequenceNumber, "layer", layer, "osn", tp.rtp.sequenceNumber, "size", len(payload)) // REMOVE d.pacer.Enqueue(pacer.Packet{ Header: hdr, Extensions: extensions, diff --git a/pkg/sfu/sequencer.go b/pkg/sfu/sequencer.go index a2699427b..e5ed8299a 100644 --- a/pkg/sfu/sequencer.go +++ b/pkg/sfu/sequencer.go @@ -116,13 +116,13 @@ func (s *sequencer) push( layer int8, codecBytes []byte, ddBytes []byte, -) { +) bool { s.Lock() defer s.Unlock() slot, isValid := s.getSlot(offSn) if !isValid { - return + return false } s.meta[s.metaWritePtr] = packetMeta{ @@ -142,6 +142,7 @@ func (s *sequencer) push( if s.metaWritePtr >= len(s.meta) { s.metaWritePtr -= len(s.meta) } + return true } func (s *sequencer) pushPadding(offSn uint16) {