Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions internal/rtpbuffer/rtpbuffer.go
Original file line number Diff line number Diff line change
Expand Up @@ -116,3 +116,5 @@ func (r *RTPBuffer) Get(seq uint16) *RetainablePacket {

return pkt
}

func (r *RTPBuffer) Started() bool { return r.started }
175 changes: 133 additions & 42 deletions pkg/jitterbuffer/jitter_buffer.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (
"errors"
"sync"

"github.com/pion/interceptor/internal/rtpbuffer"
"github.com/pion/rtp"
)

Expand Down Expand Up @@ -66,16 +67,19 @@ type (
// order, and allows removing in either sequence number order or via a
// provided timestamp.
type JitterBuffer struct {
packets *PriorityQueue
minStartCount uint16
overflowLen uint16
lastSequence uint16
playoutHead uint16
playoutReady bool
state State
stats Stats
listeners map[Event][]EventListener
mutex sync.Mutex
packetFactory rtpbuffer.PacketFactory
reorderBuffer *rtpbuffer.RTPBuffer
playbackBuffer *RingBuffer
minStartCount uint16
overflowLen uint16
lastSequence uint16
expectedSequence uint16
playoutHead uint16
playoutReady bool
state State
stats Stats
listeners map[Event][]EventListener
mutex sync.Mutex
}

// Stats Track interesting statistics for the life of this JitterBuffer
Expand All @@ -91,21 +95,33 @@ type Stats struct {
overflowCount uint32
}

var (
// ErrInvalidOperation may be returned if a Pop or Find operation is performed on an playback buffer.
ErrInvalidOperation = errors.New("attempt to find or pop on an empty list")
// ErrNotFound will be returned if the packet cannot be found in the playblack buffer.
ErrNotFound = errors.New("packet not found")
)

// New will initialize a jitter buffer and its associated statistics.
func New(opts ...Option) *JitterBuffer {
jb := &JitterBuffer{
state: Buffering,
stats: Stats{0, 0, 0},
minStartCount: 50,
overflowLen: 100,
packets: NewQueue(),
overflowLen: 1024,
listeners: make(map[Event][]EventListener),
}

for _, o := range opts {
o(jb)
}

if jb.packetFactory == nil {
jb.packetFactory = rtpbuffer.NewPacketFactoryCopy()
}
jb.reorderBuffer, _ = rtpbuffer.NewRTPBuffer(jb.overflowLen)
jb.playbackBuffer = NewRingBuffer(jb.overflowLen)

return jb
}

Expand All @@ -117,6 +133,14 @@ func WithMinimumPacketCount(count uint16) Option {
}
}

// DisableCopy bypasses copy of underlying packets. It should be used when
// you are not re-using underlying buffers of packets that have been written.
func DisableCopy() Option {
return func(jb *JitterBuffer) {
jb.packetFactory = &rtpbuffer.PacketFactoryNoOp{}
}
}

// Listen will register an event listener
// The jitter buffer may emit events correspnding, interested listerns should
// look at Event for available events.
Expand All @@ -139,12 +163,14 @@ func (jb *JitterBuffer) SetPlayoutHead(playoutHead uint16) {
defer jb.mutex.Unlock()

jb.playoutHead = playoutHead
jb.expectedSequence = playoutHead
jb.drain()
}

func (jb *JitterBuffer) updateStats(lastPktSeqNo uint16) {
// If we have at least one packet, and the next packet being pushed in is not
// at the expected sequence number increment the out of order count
if jb.packets.Length() > 0 && lastPktSeqNo != (jb.lastSequence+1) {
if jb.reorderBuffer.Started() && lastPktSeqNo != (jb.lastSequence+1) {
jb.stats.outOfOrderCount++
}
jb.lastSequence = lastPktSeqNo
Expand All @@ -157,21 +183,24 @@ func (jb *JitterBuffer) Push(packet *rtp.Packet) {
jb.mutex.Lock()
defer jb.mutex.Unlock()

if jb.packets.Length() == 0 {
jb.emit(StartBuffering)
rPacket, err := jb.packetFactory.NewPacket(&packet.Header, packet.Payload, 0, 0)
if err != nil {
return
}

if jb.packets.Length() > jb.overflowLen {
jb.stats.overflowCount++
jb.emit(BufferOverflow)
if jb.playbackBuffer.Length() == 0 {
jb.emit(StartBuffering)
}

if !jb.playoutReady && jb.packets.Length() == 0 {
if !jb.reorderBuffer.Started() {
jb.playoutHead = packet.SequenceNumber
jb.expectedSequence = packet.SequenceNumber
}

jb.updateStats(packet.SequenceNumber)
jb.packets.Push(packet, packet.SequenceNumber)

jb.reorderBuffer.Add(rPacket)
jb.drain()

jb.updateState()
}

Expand All @@ -183,7 +212,7 @@ func (jb *JitterBuffer) emit(event Event) {

func (jb *JitterBuffer) updateState() {
// For now, we only look at the number of packets captured in the play buffer
if jb.packets.Length() >= jb.minStartCount && jb.state == Buffering {
if jb.playbackBuffer.Length() >= jb.minStartCount && jb.state == Buffering {
jb.state = Emitting
jb.playoutReady = true
jb.emit(BeginPlayback)
Expand All @@ -200,29 +229,52 @@ func (jb *JitterBuffer) updateState() {
func (jb *JitterBuffer) Peek(playoutHead bool) (*rtp.Packet, error) {
jb.mutex.Lock()
defer jb.mutex.Unlock()
if jb.packets.Length() < 1 {
if !jb.reorderBuffer.Started() {
return nil, ErrBufferUnderrun
}

var packet *rtpbuffer.RetainablePacket
if playoutHead && jb.state == Emitting {
return jb.packets.Find(jb.playoutHead)
packet = jb.playbackBuffer.Peek()
if packet != nil {
if err := packet.Retain(); err != nil {
return nil, ErrNotFound
}
}
} else {
packet = jb.reorderBuffer.Get(jb.lastSequence)
}

return jb.packets.Find(jb.lastSequence)
if packet == nil {
return nil, ErrNotFound
}

return jb.takePacket(packet), nil
}

// Pop an RTP packet from the jitter buffer at the current playout head.
func (jb *JitterBuffer) Pop() (*rtp.Packet, error) {
packet, err := jb.popRetainable()
if err != nil {
return nil, err
}

return jb.takePacket(packet), nil
}

// Same as Pop, except it returns a RetainablePacket.
func (jb *JitterBuffer) popRetainable() (*rtpbuffer.RetainablePacket, error) {
jb.mutex.Lock()
defer jb.mutex.Unlock()
if jb.state != Emitting {
return nil, ErrPopWhileBuffering
}
packet, err := jb.packets.PopAt(jb.playoutHead)
if err != nil {
packet := jb.playbackBuffer.Pop()
if packet == nil {
jb.stats.underflowCount++
jb.emit(BufferUnderflow)

return nil, err
return nil, ErrNotFound
}
jb.playoutHead = (jb.playoutHead + 1)
jb.updateState()
Expand All @@ -237,30 +289,30 @@ func (jb *JitterBuffer) PopAtSequence(sq uint16) (*rtp.Packet, error) {
if jb.state != Emitting {
return nil, ErrPopWhileBuffering
}
packet, err := jb.packets.PopAt(sq)
if err != nil {
packet := jb.playbackBuffer.PopAt(sq)
if packet == nil {
jb.stats.underflowCount++
jb.emit(BufferUnderflow)

return nil, err
return nil, ErrNotFound
}
jb.playoutHead = (jb.playoutHead + 1)
jb.playoutHead = sq + 1
jb.updateState()

return packet, nil
return jb.takePacket(packet), nil
}

// PeekAtSequence will return an RTP packet from the jitter buffer at the specified Sequence
// without removing it from the buffer.
func (jb *JitterBuffer) PeekAtSequence(sq uint16) (*rtp.Packet, error) {
jb.mutex.Lock()
defer jb.mutex.Unlock()
packet, err := jb.packets.Find(sq)
if err != nil {
return nil, err
packet := jb.reorderBuffer.Get(sq)
if packet == nil {
return nil, ErrNotFound
}

return packet, nil
return jb.takePacket(packet), nil
}

// PopAtTimestamp pops an RTP packet from the jitter buffer with the provided timestamp
Expand All @@ -271,26 +323,65 @@ func (jb *JitterBuffer) PopAtTimestamp(ts uint32) (*rtp.Packet, error) {
if jb.state != Emitting {
return nil, ErrPopWhileBuffering
}
packet, err := jb.packets.PopAtTimestamp(ts)
if err != nil {
packet := jb.playbackBuffer.PopAtTimestamp(ts)
if packet == nil {
jb.stats.underflowCount++
jb.emit(BufferUnderflow)

return nil, err
return nil, ErrNotFound
}
jb.playoutHead = packet.Header().SequenceNumber + 1
jb.updateState()

return packet, nil
return jb.takePacket(packet), nil
}

// Unwrap the packet.
func (jb *JitterBuffer) takePacket(rPacket *rtpbuffer.RetainablePacket) *rtp.Packet {
header := *rPacket.Header()
payload := rPacket.Payload()
if _, ok := jb.packetFactory.(*rtpbuffer.PacketFactoryNoOp); !ok {
out := make([]byte, len(payload))
copy(out, payload)
payload = out
}
rPacket.Release()

return &rtp.Packet{
Header: header,
Payload: payload,
}
}

// Move the packets into the playback buffer as long as they are in order.
func (jb *JitterBuffer) drain() {
for range jb.overflowLen {
rPacket := jb.reorderBuffer.Get(jb.expectedSequence)
if rPacket == nil {
break
}
if !jb.playbackBuffer.Push(rPacket) {
rPacket.Release()
jb.stats.overflowCount++
jb.emit(BufferOverflow)
break
}
jb.expectedSequence++
}
}

// Clear will empty the buffer and optionally reset the state.
func (jb *JitterBuffer) Clear(resetState bool) {
jb.mutex.Lock()
defer jb.mutex.Unlock()
jb.packets.Clear()
jb.reorderBuffer.Clear()
jb.playbackBuffer.Clear()

if resetState {
jb.lastSequence = 0
jb.expectedSequence = 0
jb.playoutHead = 0
jb.playoutReady = false
jb.state = Buffering
jb.stats = Stats{0, 0, 0}
jb.minStartCount = 50
Expand Down
Loading