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
20 changes: 14 additions & 6 deletions pkg/rtpfb/interceptor.go
Original file line number Diff line number Diff line change
Expand Up @@ -220,18 +220,21 @@ func (i *Interceptor) BindRTCPReader(reader interceptor.RTCPReader) interceptor.
//nolint:cyclop
func (i *Interceptor) processFeedback(ts time.Time, pkts []rtcp.Packet) (time.Duration, []PacketReport) {
shortestRTT := time.Duration(math.MaxInt64)
var ackDelay time.Duration
measured := false

for _, pkt := range pkts {
switch fb := pkt.(type) {
case *rtcp.CCFeedbackReport:
var acksPerSSRC map[uint32][]acknowledgement
ackDelay, acksPerSSRC = convertCCFB(ts, fb)
ackDelay, acksPerSSRC := convertCCFB(ts, fb)
for ssrc, acks := range acksPerSSRC {
for _, ack := range acks {
rtt, ok := i.history.onCCFBFeedback(ts, ssrc, ack)
if ok && rtt < shortestRTT {
shortestRTT = rtt
if !ok {
continue
}
if corrected := max(rtt-ackDelay, 0); corrected < shortestRTT {
shortestRTT = corrected
measured = true
}
}
}
Expand All @@ -240,10 +243,15 @@ func (i *Interceptor) processFeedback(ts time.Time, pkts []rtcp.Packet) (time.Du
rtt, ok := i.history.onTWCCFeedback(ts, ack)
if ok && rtt < shortestRTT {
shortestRTT = rtt
measured = true
}
}
}
}

return shortestRTT - ackDelay, i.history.buildReport()
if !measured {
return 0, i.history.buildReport()
}

return shortestRTT, i.history.buildReport()
}
120 changes: 118 additions & 2 deletions pkg/rtpfb/interceptor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,18 @@ type ackListEntry struct {
type mockHistory struct {
log []PacketReport
acks []ackListEntry

// rtts returns the RTT and match result for an acknowledgement. If nil,
// the mock reports a match with a zero RTT.
rtts func(ack acknowledgement) (time.Duration, bool)
}

func (m *mockHistory) rtt(ack acknowledgement) (time.Duration, bool) {
if m.rtts == nil {
return 0, true
}

return m.rtts(ack)
}

// addOutgoing implements packetLog.
Expand Down Expand Up @@ -59,7 +71,7 @@ func (m *mockHistory) onCCFBFeedback(ts time.Time, ssrc uint32, ack acknowledgem
ack: ack,
})

return 0, true
return m.rtt(ack)
}

// onTWCCFeedback implements packetLog.
Expand All @@ -70,7 +82,7 @@ func (m *mockHistory) onTWCCFeedback(ts time.Time, ack acknowledgement) (time.Du
ack: ack,
})

return 0, true
return m.rtt(ack)
}

func TestInterceptor(t *testing.T) {
Expand Down Expand Up @@ -216,3 +228,107 @@ func TestInterceptor(t *testing.T) {
}
})
}

func TestProcessFeedbackRTT(t *testing.T) {
mockTimeStamp := time.Time{}.Add(120 * time.Second)

// ccfbReport builds a report acknowledging a single packet. The arrival
// time offset determines the ack delay: ato/1024 seconds.
ccfbReport := func(seqNr uint16, ato uint16) *rtcp.CCFeedbackReport {
return &rtcp.CCFeedbackReport{
ReportBlocks: []rtcp.CCFeedbackReportBlock{
{
MediaSSRC: 0,
BeginSequence: seqNr,
MetricBlocks: []rtcp.CCFeedbackMetricBlock{
{Received: true, ArrivalTimeOffset: ato},
},
},
},
ReportTimestamp: 0,
}
}
twccReport := &rtcp.TransportLayerCC{
BaseSequenceNumber: 1,
PacketStatusCount: 1,
PacketChunks: []rtcp.PacketStatusChunk{
&rtcp.RunLengthChunk{
PacketStatusSymbol: rtcp.TypeTCCPacketReceivedSmallDelta,
RunLength: 1,
},
},
RecvDeltas: []*rtcp.RecvDelta{
{Type: rtcp.TypeTCCPacketReceivedSmallDelta, Delta: 0},
},
}

cases := []struct {
name string
packets []rtcp.Packet
rtts func(ack acknowledgement) (time.Duration, bool)
expected time.Duration
}{
{
name: "no_match_returns_zero",
packets: []rtcp.Packet{ccfbReport(1, 512)},
rtts: func(acknowledgement) (time.Duration, bool) {
return 0, false
},
expected: 0,
},
{
name: "subtracts_ack_delay",
packets: []rtcp.Packet{ccfbReport(1, 512)},
rtts: func(acknowledgement) (time.Duration, bool) {
return 800 * time.Millisecond, true
},
expected: 300 * time.Millisecond,
},
{
name: "ack_delay_is_per_report",
packets: []rtcp.Packet{ccfbReport(1, 512), ccfbReport(2, 1024)},
rtts: func(ack acknowledgement) (time.Duration, bool) {
if ack.sequenceNumber == 1 {
return 800 * time.Millisecond, true
}

return 1200 * time.Millisecond, true
},
expected: 200 * time.Millisecond,
},
{
name: "ack_delay_larger_than_rtt_clamps_to_zero",
packets: []rtcp.Packet{ccfbReport(1, 1024)},
rtts: func(acknowledgement) (time.Duration, bool) {
return 500 * time.Millisecond, true
},
expected: 0,
},
{
name: "twcc_rtt_is_not_corrected",
packets: []rtcp.Packet{twccReport},
rtts: func(acknowledgement) (time.Duration, bool) {
return 300 * time.Millisecond, true
},
expected: 300 * time.Millisecond,
},
}

for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
mh := &mockHistory{rtts: tc.rtts}
mt := func() time.Time {
return mockTimeStamp
}
f, err := NewInterceptor(timeFactory(mt), setHistory(mh))
assert.NoError(t, err)
i, err := f.NewInterceptor("")
assert.NoError(t, err)

ic, ok := i.(*Interceptor)
assert.True(t, ok)
rtt, _ := ic.processFeedback(mockTimeStamp, tc.packets)
assert.Equal(t, tc.expected, rtt)
})
}
}
9 changes: 7 additions & 2 deletions pkg/rtpfb/report.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,12 @@ type PacketReport struct {
// acknowledged packets that were still in the history and not yet included in
// an earlier Report.
type Report struct {
Arrival time.Time
RTT time.Duration
Arrival time.Time

// RTT is the estimated round trip time. It is zero if the feedback packet
// did not acknowledge any packet that is still known to the interceptor,
// i.e. no RTT could be measured.
RTT time.Duration

PacketReports []PacketReport
}
Loading