Skip to content

Commit 9ef04d4

Browse files
cptpcrdadriancable
authored andcommitted
Don't use out-of-order packets in sender reports
If a sender report is generated immediately after sending an out-of-order packet (which can happen e.g. in an SFU when forwarding media), the timestamp from the last in-order packet should be extrapolated, since the departure time of an out-of-order packet is is not properly correlated with its RTP timestamp. Leave an option to re-enable the old behavior.
1 parent 38fe7f5 commit 9ef04d4

4 files changed

Lines changed: 149 additions & 7 deletions

File tree

pkg/report/sender_interceptor.go

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -59,6 +59,8 @@ type SenderInterceptor struct {
5959
wg sync.WaitGroup
6060
close chan struct{}
6161
started chan struct{}
62+
63+
useLatestPacket bool
6264
}
6365

6466
func (s *SenderInterceptor) isClosed() bool {
@@ -133,7 +135,7 @@ func (s *SenderInterceptor) loop(rtcpWriter interceptor.RTCPWriter) {
133135
// BindLocalStream lets you modify any outgoing RTP packets. It is called once for per LocalStream. The returned method
134136
// will be called once per rtp packet.
135137
func (s *SenderInterceptor) BindLocalStream(info *interceptor.StreamInfo, writer interceptor.RTPWriter) interceptor.RTPWriter {
136-
stream := newSenderStream(info.SSRC, info.ClockRate)
138+
stream := newSenderStream(info.SSRC, info.ClockRate, s.useLatestPacket)
137139
s.streams.Store(info.SSRC, stream)
138140

139141
return interceptor.RTPWriterFunc(func(header *rtp.Header, payload []byte, a interceptor.Attributes) (int, error) {

pkg/report/sender_interceptor_test.go

Lines changed: 124 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -92,6 +92,130 @@ func TestSenderInterceptor(t *testing.T) {
9292
}, sr)
9393
})
9494

95+
t.Run("out of order RTP packets", func(t *testing.T) {
96+
mt := &test.MockTime{}
97+
f, err := NewSenderInterceptor(
98+
SenderInterval(time.Millisecond*50),
99+
SenderLog(logging.NewDefaultLoggerFactory().NewLogger("test")),
100+
SenderNow(mt.Now),
101+
)
102+
assert.NoError(t, err)
103+
104+
i, err := f.NewInterceptor("")
105+
assert.NoError(t, err)
106+
107+
stream := test.NewMockStream(&interceptor.StreamInfo{
108+
SSRC: 123456,
109+
ClockRate: 90000,
110+
}, i)
111+
defer func() {
112+
assert.NoError(t, stream.Close())
113+
}()
114+
115+
// Write several packets
116+
for i := 0; i < 10; i++ {
117+
assert.NoError(t, stream.WriteRTP(&rtp.Packet{
118+
Header: rtp.Header{
119+
SequenceNumber: uint16(i),
120+
Timestamp: uint32(i),
121+
},
122+
Payload: []byte("\x00\x00"),
123+
}))
124+
}
125+
126+
// Skip a packet, then redeliver it out-of-order
127+
assert.NoError(t, stream.WriteRTP(&rtp.Packet{
128+
Header: rtp.Header{
129+
SequenceNumber: 12,
130+
Timestamp: 12,
131+
},
132+
Payload: []byte("\x00\x00"),
133+
}))
134+
assert.NoError(t, stream.WriteRTP(&rtp.Packet{
135+
Header: rtp.Header{
136+
SequenceNumber: 11,
137+
Timestamp: 11,
138+
},
139+
Payload: []byte("\x00\x00"),
140+
}))
141+
142+
pkts := <-stream.WrittenRTCP()
143+
assert.Equal(t, len(pkts), 1)
144+
sr, ok := pkts[0].(*rtcp.SenderReport)
145+
assert.True(t, ok)
146+
// The out-of-order packet is included in PacketCount and OctetCount, but the RTP
147+
// timestamp of the last in-order packet is used for RTPTime
148+
assert.Equal(t, &rtcp.SenderReport{
149+
SSRC: 123456,
150+
NTPTime: ntp.ToNTP(mt.Now()),
151+
RTPTime: 12,
152+
PacketCount: 12,
153+
OctetCount: 24,
154+
}, sr)
155+
})
156+
157+
t.Run("out of order RTP packets with SenderUseLatestPacket", func(t *testing.T) {
158+
mt := &test.MockTime{}
159+
f, err := NewSenderInterceptor(
160+
SenderInterval(time.Millisecond*50),
161+
SenderLog(logging.NewDefaultLoggerFactory().NewLogger("test")),
162+
SenderNow(mt.Now),
163+
SenderUseLatestPacket(),
164+
)
165+
assert.NoError(t, err)
166+
167+
i, err := f.NewInterceptor("")
168+
assert.NoError(t, err)
169+
170+
stream := test.NewMockStream(&interceptor.StreamInfo{
171+
SSRC: 123456,
172+
ClockRate: 90000,
173+
}, i)
174+
defer func() {
175+
assert.NoError(t, stream.Close())
176+
}()
177+
178+
// Write several packets
179+
for i := 0; i < 10; i++ {
180+
assert.NoError(t, stream.WriteRTP(&rtp.Packet{
181+
Header: rtp.Header{
182+
SequenceNumber: uint16(i),
183+
Timestamp: uint32(i),
184+
},
185+
Payload: []byte("\x00\x00"),
186+
}))
187+
}
188+
189+
// Skip a packet, then redeliver it out-of-order
190+
assert.NoError(t, stream.WriteRTP(&rtp.Packet{
191+
Header: rtp.Header{
192+
SequenceNumber: 12,
193+
Timestamp: 12,
194+
},
195+
Payload: []byte("\x00\x00"),
196+
}))
197+
assert.NoError(t, stream.WriteRTP(&rtp.Packet{
198+
Header: rtp.Header{
199+
SequenceNumber: 11,
200+
Timestamp: 11,
201+
},
202+
Payload: []byte("\x00\x00"),
203+
}))
204+
205+
pkts := <-stream.WrittenRTCP()
206+
assert.Equal(t, len(pkts), 1)
207+
sr, ok := pkts[0].(*rtcp.SenderReport)
208+
assert.True(t, ok)
209+
// The out-of-order packet *is* used for RTPTime
210+
assert.Equal(t, &rtcp.SenderReport{
211+
SSRC: 123456,
212+
NTPTime: ntp.ToNTP(mt.Now()),
213+
RTPTime: 11,
214+
PacketCount: 12,
215+
OctetCount: 24,
216+
}, sr)
217+
})
218+
95219
t.Run("inject ticker", func(t *testing.T) {
96220
mNow := &test.MockTime{}
97221
mTick := &test.MockTicker{

pkg/report/sender_option.go

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,15 @@ func SenderTicker(f TickerFactory) SenderOption {
4444
}
4545
}
4646

47+
// SenderUseLatestPacket sets the interceptor to always use the latest packet, even
48+
// if it appears to be out-of-order.
49+
func SenderUseLatestPacket() SenderOption {
50+
return func(r *SenderInterceptor) error {
51+
r.useLatestPacket = true
52+
return nil
53+
}
54+
}
55+
4756
// enableStartTracking is used by tests to synchronize whether the loop() has begun
4857
// and it's safe to start sending ticks to the ticker.
4958
func enableStartTracking(startedCh chan struct{}) SenderOption {

pkg/report/sender_stream.go

Lines changed: 13 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -17,27 +17,34 @@ type senderStream struct {
1717
clockRate float64
1818
m sync.Mutex
1919

20+
useLatestPacket bool
21+
2022
// data from rtp packets
2123
lastRTPTimeRTP uint32
2224
lastRTPTimeTime time.Time
25+
lastRTPSN uint16
2326
packetCount uint32
2427
octetCount uint32
2528
}
2629

27-
func newSenderStream(ssrc uint32, clockRate uint32) *senderStream {
30+
func newSenderStream(ssrc uint32, clockRate uint32, useLatestPacket bool) *senderStream {
2831
return &senderStream{
29-
ssrc: ssrc,
30-
clockRate: float64(clockRate),
32+
ssrc: ssrc,
33+
clockRate: float64(clockRate),
34+
useLatestPacket: useLatestPacket,
3135
}
3236
}
3337

3438
func (stream *senderStream) processRTP(now time.Time, header *rtp.Header, payload []byte) {
3539
stream.m.Lock()
3640
defer stream.m.Unlock()
3741

38-
// always update time to minimize errors
39-
stream.lastRTPTimeRTP = header.Timestamp
40-
stream.lastRTPTimeTime = now
42+
if stream.useLatestPacket || stream.packetCount == 0 || int16(header.SequenceNumber-stream.lastRTPSN) > 0 {
43+
// Told to consider every packet, or this was the first packet, or it's in-order
44+
stream.lastRTPSN = header.SequenceNumber
45+
stream.lastRTPTimeRTP = header.Timestamp
46+
stream.lastRTPTimeTime = now
47+
}
4148

4249
stream.packetCount++
4350
stream.octetCount += uint32(len(payload))

0 commit comments

Comments
 (0)