Skip to content

Commit bd99739

Browse files
authored
Export receive time of packet. (#22)
* Export receive time of packet. * staticcheck * handler close * writer closer * not the correct interface
1 parent 11d305e commit bd99739

4 files changed

Lines changed: 41 additions & 38 deletions

File tree

jitter/buffer.go

Lines changed: 26 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,11 @@ import (
2525
"github.com/livekit/protocol/logger"
2626
)
2727

28+
type ExtPacket struct {
29+
ReceivedAt time.Time
30+
*rtp.Packet
31+
}
32+
2833
type Buffer struct {
2934
depacketizer rtp.Depacketizer
3035
latency time.Duration
@@ -58,7 +63,7 @@ type BufferStats struct {
5863
SamplesPopped uint64 // samples sent to handler
5964
}
6065

61-
type PacketFunc func(packets []*rtp.Packet)
66+
type PacketFunc func(packets []ExtPacket)
6267

6368
func NewBuffer(
6469
depacketizer rtp.Depacketizer,
@@ -117,7 +122,7 @@ func (b *Buffer) UpdateLatency(latency time.Duration) {
117122

118123
b.latency = latency
119124
if b.head != nil {
120-
b.timer.Reset(time.Until(b.head.received.Add(latency)))
125+
b.timer.Reset(time.Until(b.head.extPacket.ReceivedAt.Add(latency)))
121126
}
122127
}
123128

@@ -199,10 +204,10 @@ func (b *Buffer) push(pkt *rtp.Packet) {
199204
return
200205
}
201206

202-
beforeHead := before(pkt.SequenceNumber, b.head.packet.SequenceNumber)
203-
afterTail := !before(pkt.SequenceNumber, b.tail.packet.SequenceNumber)
204-
withinHeadRange := withinRange(pkt.SequenceNumber, b.head.packet.SequenceNumber)
205-
withinTailRange := withinRange(pkt.SequenceNumber, b.tail.packet.SequenceNumber)
207+
beforeHead := before(pkt.SequenceNumber, b.head.extPacket.SequenceNumber)
208+
afterTail := !before(pkt.SequenceNumber, b.tail.extPacket.SequenceNumber)
209+
withinHeadRange := withinRange(pkt.SequenceNumber, b.head.extPacket.SequenceNumber)
210+
withinTailRange := withinRange(pkt.SequenceNumber, b.tail.extPacket.SequenceNumber)
206211

207212
switch {
208213
case beforeHead && withinHeadRange:
@@ -221,8 +226,8 @@ func (b *Buffer) push(pkt *rtp.Packet) {
221226
case withinTailRange:
222227
// insert, search from tail
223228
for c := b.tail.prev; c != nil; c = c.prev {
224-
discont = !withinRange(pkt.SequenceNumber, c.packet.SequenceNumber)
225-
if !before(pkt.SequenceNumber, c.packet.SequenceNumber) || discont {
229+
discont = !withinRange(pkt.SequenceNumber, c.extPacket.SequenceNumber)
230+
if !before(pkt.SequenceNumber, c.extPacket.SequenceNumber) || discont {
226231
// insert after c
227232
p.discont = discont && p.start
228233
p.prev = c
@@ -236,8 +241,8 @@ func (b *Buffer) push(pkt *rtp.Packet) {
236241
case withinHeadRange:
237242
// insert, search from head
238243
for c := b.head.next; c != nil; c = c.next {
239-
discont = !withinRange(pkt.SequenceNumber, c.packet.SequenceNumber)
240-
if before(pkt.SequenceNumber, c.packet.SequenceNumber) || discont {
244+
discont = !withinRange(pkt.SequenceNumber, c.extPacket.SequenceNumber)
245+
if before(pkt.SequenceNumber, c.extPacket.SequenceNumber) || discont {
241246
// insert before c
242247
p.prev = c.prev
243248
p.next = c
@@ -266,12 +271,12 @@ func (b *Buffer) popReady() {
266271
for b.head != nil &&
267272
b.head.isComplete() {
268273

269-
if b.head.packet.SequenceNumber == b.prevSN+1 || b.head.discont || !b.initialized {
274+
if b.head.extPacket.SequenceNumber == b.prevSN+1 || b.head.discont || !b.initialized {
270275
// normal
271-
} else if b.head.received.Before(expiry) {
276+
} else if b.head.extPacket.ReceivedAt.Before(expiry) {
272277
// max latency reached
273278
loss = true
274-
b.stats.PacketsLost += uint64(b.head.packet.SequenceNumber - b.prevSN - 1)
279+
b.stats.PacketsLost += uint64(b.head.extPacket.SequenceNumber - b.prevSN - 1)
275280
} else {
276281
break
277282
}
@@ -286,17 +291,17 @@ func (b *Buffer) popReady() {
286291
}
287292

288293
if b.head != nil {
289-
b.timer.Reset(time.Until(b.head.received.Add(b.latency)))
294+
b.timer.Reset(time.Until(b.head.extPacket.ReceivedAt.Add(b.latency)))
290295
}
291296
}
292297

293298
// dropIncompleteExpired drops incomplete expired packets
294299
func (b *Buffer) dropIncompleteExpired(expiry time.Time) {
295300
dropped := false
296301

297-
for b.head != nil && !b.head.isComplete() && b.head.received.Before(expiry) {
302+
for b.head != nil && !b.head.isComplete() && b.head.extPacket.ReceivedAt.Before(expiry) {
298303
if b.initialized && !b.head.discont {
299-
b.stats.PacketsLost += uint64(b.head.packet.SequenceNumber - b.prevSN - 1)
304+
b.stats.PacketsLost += uint64(b.head.extPacket.SequenceNumber - b.prevSN - 1)
300305
}
301306

302307
b.free(b.popHead())
@@ -310,15 +315,15 @@ func (b *Buffer) dropIncompleteExpired(expiry time.Time) {
310315
}
311316
}
312317

313-
func (b *Buffer) popSample() []*rtp.Packet {
314-
sample := make([]*rtp.Packet, 0, b.size)
318+
func (b *Buffer) popSample() []ExtPacket {
319+
sample := make([]ExtPacket, 0, b.size)
315320
end := false
316321
for !end {
317322
c := b.popHead()
318323
end = c.end
319324

320-
if !c.packet.Padding {
321-
sample = append(sample, c.packet)
325+
if !c.extPacket.Padding {
326+
sample = append(sample, c.extPacket)
322327
}
323328

324329
b.stats.PacketsPopped++
@@ -333,7 +338,7 @@ func (b *Buffer) popSample() []*rtp.Packet {
333338

334339
func (b *Buffer) popHead() *packet {
335340
c := b.head
336-
b.prevSN = c.packet.SequenceNumber
341+
b.prevSN = c.extPacket.SequenceNumber
337342
b.head = c.next
338343
if b.head == nil {
339344
b.tail = nil

jitter/buffer_test.go

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -26,8 +26,8 @@ import (
2626

2727
const testBufferLatency = 800 * time.Millisecond
2828

29-
func chanFunc(t testing.TB, out chan<- []*rtp.Packet) PacketFunc {
30-
return func(packets []*rtp.Packet) {
29+
func chanFunc(t testing.TB, out chan<- []ExtPacket) PacketFunc {
30+
return func(packets []ExtPacket) {
3131
select {
3232
case out <- packets:
3333
default:
@@ -37,7 +37,7 @@ func chanFunc(t testing.TB, out chan<- []*rtp.Packet) PacketFunc {
3737
}
3838

3939
func TestJitterBuffer(t *testing.T) {
40-
out := make(chan []*rtp.Packet, 100)
40+
out := make(chan []ExtPacket, 100)
4141
b := NewBuffer(&testDepacketizer{}, testBufferLatency, chanFunc(t, out))
4242
s := newTestStream()
4343

@@ -57,7 +57,7 @@ func TestJitterBuffer(t *testing.T) {
5757
}
5858

5959
func TestSamples(t *testing.T) {
60-
out := make(chan []*rtp.Packet, 100)
60+
out := make(chan []ExtPacket, 100)
6161
b := NewBuffer(&testDepacketizer{}, testBufferLatency, chanFunc(t, out))
6262
s := newTestStream()
6363

@@ -80,7 +80,7 @@ func TestSamples(t *testing.T) {
8080
}
8181

8282
func TestJitter(t *testing.T) {
83-
out := make(chan []*rtp.Packet, 100)
83+
out := make(chan []ExtPacket, 100)
8484
b := NewBuffer(&testDepacketizer{}, testBufferLatency, chanFunc(t, out))
8585
s := newTestStream()
8686

@@ -114,7 +114,7 @@ func TestJitter(t *testing.T) {
114114
}
115115

116116
func TestDiscontinuity(t *testing.T) {
117-
out := make(chan []*rtp.Packet, 100)
117+
out := make(chan []ExtPacket, 100)
118118
b := NewBuffer(&testDepacketizer{}, testBufferLatency, chanFunc(t, out))
119119
s := newTestStream()
120120

@@ -139,7 +139,7 @@ func TestDiscontinuity(t *testing.T) {
139139
}
140140

141141
func TestLostPackets(t *testing.T) {
142-
out := make(chan []*rtp.Packet, 100)
142+
out := make(chan []ExtPacket, 100)
143143
b := NewBuffer(&testDepacketizer{}, testBufferLatency, chanFunc(t, out))
144144
s := newTestStream()
145145

@@ -173,7 +173,7 @@ func TestLostPackets(t *testing.T) {
173173
}
174174

175175
func TestDroppedPackets(t *testing.T) {
176-
out := make(chan []*rtp.Packet, 100)
176+
out := make(chan []ExtPacket, 100)
177177
b := NewBuffer(&testDepacketizer{}, testBufferLatency, chanFunc(t, out))
178178
s := newTestStream()
179179

@@ -232,7 +232,7 @@ func TestDroppedPackets(t *testing.T) {
232232
})
233233
}
234234

235-
func checkSample(t *testing.T, out chan []*rtp.Packet, expected int) {
235+
func checkSample(t *testing.T, out chan []ExtPacket, expected int) {
236236
select {
237237
case sample := <-out:
238238
if expected == 0 {

jitter/packet.go

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -21,11 +21,10 @@ import (
2121
)
2222

2323
type packet struct {
24-
received time.Time
24+
extPacket ExtPacket
2525
prev, next *packet
2626
start, end bool
2727
discont bool
28-
packet *rtp.Packet
2928
}
3029

3130
func (b *Buffer) newPacket(pkt *rtp.Packet) *packet {
@@ -37,12 +36,11 @@ func (b *Buffer) newPacket(pkt *rtp.Packet) *packet {
3736
}
3837
b.pool = p.next
3938

40-
p.received = time.Now()
4139
p.prev = nil
4240
p.next = nil
4341
p.start = b.depacketizer.IsPartitionHead(pkt.Payload)
4442
p.end = b.depacketizer.IsPartitionTail(pkt.Marker, pkt.Payload)
45-
p.packet = pkt
43+
p.extPacket = ExtPacket{time.Now(), pkt}
4644

4745
return p
4846
}
@@ -57,7 +55,7 @@ func (p *packet) isComplete() bool {
5755
}
5856

5957
for c := p; c.next != nil; c = c.next {
60-
if c.next.packet.SequenceNumber != c.packet.SequenceNumber+1 {
58+
if c.next.extPacket.SequenceNumber != c.extPacket.SequenceNumber+1 {
6159
return false
6260
}
6361
if c.next.end {
@@ -72,7 +70,7 @@ func (b *Buffer) free(pkt *packet) {
7270
b.size--
7371

7472
pkt.prev = nil
75-
pkt.packet = nil
73+
pkt.extPacket = ExtPacket{}
7674
pkt.next = b.pool
7775

7876
b.pool = pkt

rtp/jitter.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -33,9 +33,9 @@ func HandleJitter(h HandlerCloser) HandlerCloser {
3333
}
3434
// Jitter buffer expects to be closed (to stop the timer), but handler interface doesn't allow it.
3535
// This should be fine, because GC can now collect timers and goroutines blocked on them if they are not referenced.
36-
handler.buf = jitter.NewBuffer(audioDepacketizer{}, jitterMaxLatency, func(packets []*rtp.Packet) {
36+
handler.buf = jitter.NewBuffer(audioDepacketizer{}, jitterMaxLatency, func(packets []jitter.ExtPacket) {
3737
for _, p := range packets {
38-
handler.handleRTP(p)
38+
handler.handleRTP(p.Packet)
3939
}
4040
})
4141
return handler

0 commit comments

Comments
 (0)