mirror of
https://github.com/livekit/livekit.git
synced 2026-08-28 18:08:17 +00:00
Use seq num offsets cache instead of missing seq num map. (#658)
* Use seq num offsets cache instead of missing seq num map. Map operations can be costly. Use a fixed size array to store offsets. 4096 sequence numbers should be more than 16 seconds for 720p video which should be plenty to look up offset of out-of-order packets. Packets out-of-order more than that will probably be useless anyway. * Move offset cache population to only when we are forwarding * add some debug logs * Remove debug
This commit is contained in:
+38
-15
@@ -20,6 +20,9 @@ const (
|
||||
|
||||
const (
|
||||
RtxGateWindow = 2000
|
||||
|
||||
SnOffsetCacheSize = 4096
|
||||
SnOffsetCacheMask = SnOffsetCacheSize - 1
|
||||
)
|
||||
|
||||
type TranslationParamsRTP struct {
|
||||
@@ -41,7 +44,9 @@ type RTPMungerParams struct {
|
||||
tsOffset uint32
|
||||
lastMarker bool
|
||||
|
||||
missingSNs map[uint16]uint16
|
||||
snOffsets [SnOffsetCacheSize]uint16
|
||||
snOffsetsWritePtr int
|
||||
snOffsetsOccupancy int
|
||||
|
||||
rtxGateSn uint16
|
||||
isInRtxGateRegion bool
|
||||
@@ -56,9 +61,6 @@ type RTPMunger struct {
|
||||
func NewRTPMunger(logger logger.Logger) *RTPMunger {
|
||||
return &RTPMunger{
|
||||
logger: logger,
|
||||
RTPMungerParams: RTPMungerParams{
|
||||
missingSNs: make(map[uint16]uint16, 10),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
@@ -84,8 +86,9 @@ func (r *RTPMunger) UpdateSnTsOffsets(extPkt *buffer.ExtPacket, snAdjust uint16,
|
||||
r.snOffset = extPkt.Packet.SequenceNumber - r.lastSN - snAdjust
|
||||
r.tsOffset = extPkt.Packet.Timestamp - r.lastTS - tsAdjust
|
||||
|
||||
// clear incoming missing sequence numbers on layer/source switch
|
||||
r.missingSNs = make(map[uint16]uint16, 10)
|
||||
// clear offsets cache layer/source switch
|
||||
r.snOffsetsWritePtr = 0
|
||||
r.snOffsetsOccupancy = 0
|
||||
}
|
||||
|
||||
func (r *RTPMunger) PacketDropped(extPkt *buffer.ExtPacket) {
|
||||
@@ -96,19 +99,24 @@ func (r *RTPMunger) PacketDropped(extPkt *buffer.ExtPacket) {
|
||||
r.highestIncomingSN = extPkt.Packet.SequenceNumber
|
||||
r.snOffset += 1
|
||||
r.lastSN = extPkt.Packet.SequenceNumber - r.snOffset
|
||||
|
||||
r.snOffsetsWritePtr = (r.snOffsetsWritePtr - 1) & SnOffsetCacheMask
|
||||
r.snOffsetsOccupancy--
|
||||
if r.snOffsetsOccupancy < 0 {
|
||||
r.logger.Warnw("sequence number offset cache is invalid", nil, "occupancy", r.snOffsetsOccupancy)
|
||||
}
|
||||
}
|
||||
|
||||
func (r *RTPMunger) UpdateAndGetSnTs(extPkt *buffer.ExtPacket) (*TranslationParamsRTP, error) {
|
||||
// if out-of-order, look up missing sequence number cache
|
||||
// if out-of-order, look up sequence number offset cache
|
||||
if !extPkt.Head {
|
||||
snOffset, ok := r.missingSNs[extPkt.Packet.SequenceNumber]
|
||||
if !ok {
|
||||
snOffset, isValid := r.getSnOffset(extPkt.Packet.SequenceNumber)
|
||||
if !isValid {
|
||||
return &TranslationParamsRTP{
|
||||
snOrdering: SequenceNumberOrderingOutOfOrder,
|
||||
}, ErrOutOfOrderSequenceNumberCacheMiss
|
||||
}
|
||||
|
||||
delete(r.missingSNs, extPkt.Packet.SequenceNumber)
|
||||
return &TranslationParamsRTP{
|
||||
snOrdering: SequenceNumberOrderingOutOfOrder,
|
||||
sequenceNumber: extPkt.Packet.SequenceNumber - snOffset,
|
||||
@@ -118,14 +126,9 @@ func (r *RTPMunger) UpdateAndGetSnTs(extPkt *buffer.ExtPacket) (*TranslationPara
|
||||
|
||||
ordering := SequenceNumberOrderingContiguous
|
||||
|
||||
// if there are gaps, record it in missing sequence number cache
|
||||
diff := extPkt.Packet.SequenceNumber - r.highestIncomingSN
|
||||
if diff > 1 {
|
||||
ordering = SequenceNumberOrderingGap
|
||||
|
||||
for i := r.highestIncomingSN + 1; i != extPkt.Packet.SequenceNumber; i++ {
|
||||
r.missingSNs[i] = r.snOffset
|
||||
}
|
||||
} else {
|
||||
// can get duplicate packet due to FEC
|
||||
if diff == 0 {
|
||||
@@ -148,6 +151,16 @@ func (r *RTPMunger) UpdateAndGetSnTs(extPkt *buffer.ExtPacket) (*TranslationPara
|
||||
}
|
||||
}
|
||||
|
||||
// record sn offset
|
||||
for i := r.highestIncomingSN + 1; i != extPkt.Packet.SequenceNumber+1; i++ {
|
||||
r.snOffsets[r.snOffsetsWritePtr] = r.snOffset
|
||||
r.snOffsetsWritePtr = (r.snOffsetsWritePtr + 1) & SnOffsetCacheMask
|
||||
r.snOffsetsOccupancy++
|
||||
}
|
||||
if r.snOffsetsOccupancy > SnOffsetCacheSize {
|
||||
r.snOffsetsOccupancy = SnOffsetCacheSize
|
||||
}
|
||||
|
||||
// in-order incoming packet, may or may not be contiguous.
|
||||
// In the case of loss (i.e. incoming sequence number is not contiguous),
|
||||
// forward even if it is a padding only packet. With temporal scalability,
|
||||
@@ -227,3 +240,13 @@ func (r *RTPMunger) UpdateAndGetPaddingSnTs(num int, clockRate uint32, frameRate
|
||||
func (r *RTPMunger) IsOnFrameBoundary() bool {
|
||||
return r.lastMarker
|
||||
}
|
||||
|
||||
func (r *RTPMunger) getSnOffset(sn uint16) (uint16, bool) {
|
||||
diff := r.highestIncomingSN - sn
|
||||
if int(diff) >= r.snOffsetsOccupancy {
|
||||
return 0, false
|
||||
}
|
||||
|
||||
readPtr := (r.snOffsetsWritePtr - int(diff) - 1) & SnOffsetCacheMask
|
||||
return r.snOffsets[readPtr], true
|
||||
}
|
||||
|
||||
+25
-39
@@ -1,7 +1,6 @@
|
||||
package sfu
|
||||
|
||||
import (
|
||||
"reflect"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
@@ -72,8 +71,8 @@ func TestUpdateSnTsOffsets(t *testing.T) {
|
||||
require.True(t, r.highestIncomingSN == 33332)
|
||||
require.True(t, r.lastSN == 23333)
|
||||
require.True(t, r.lastTS == 0xabcdef)
|
||||
require.EqualValues(t, 9999, r.snOffset)
|
||||
require.EqualValues(t, 0xffffffff, r.tsOffset)
|
||||
require.Equal(t, uint16(9999), r.snOffset)
|
||||
require.Equal(t, uint32(0xffffffff), r.tsOffset)
|
||||
}
|
||||
|
||||
func TestPacketDropped(t *testing.T) {
|
||||
@@ -89,8 +88,8 @@ func TestPacketDropped(t *testing.T) {
|
||||
require.Equal(t, r.highestIncomingSN, uint16(23332))
|
||||
require.Equal(t, r.lastSN, uint16(23333))
|
||||
require.Equal(t, r.lastTS, uint32(0xabcdef))
|
||||
require.EqualValues(t, 0, r.snOffset)
|
||||
require.EqualValues(t, 0, r.tsOffset)
|
||||
require.Equal(t, uint16(0), r.snOffset)
|
||||
require.Equal(t, uint32(0), r.tsOffset)
|
||||
|
||||
// drop a non-head packet, should cause no change in internals
|
||||
params = &testutils.TestExtPacketParams{
|
||||
@@ -102,7 +101,7 @@ func TestPacketDropped(t *testing.T) {
|
||||
r.PacketDropped(extPkt)
|
||||
require.Equal(t, r.highestIncomingSN, uint16(23332))
|
||||
require.Equal(t, r.lastSN, uint16(23333))
|
||||
require.EqualValues(t, 0, r.snOffset)
|
||||
require.Equal(t, uint16(0), r.snOffset)
|
||||
|
||||
// drop a head packet and check offset increases
|
||||
params = &testutils.TestExtPacketParams{
|
||||
@@ -115,7 +114,7 @@ func TestPacketDropped(t *testing.T) {
|
||||
r.PacketDropped(extPkt)
|
||||
require.Equal(t, r.highestIncomingSN, uint16(44444))
|
||||
require.Equal(t, r.lastSN, uint16(44443))
|
||||
require.EqualValues(t, 1, r.snOffset)
|
||||
require.Equal(t, uint16(1), r.snOffset)
|
||||
}
|
||||
|
||||
func TestOutOfOrderSequenceNumber(t *testing.T) {
|
||||
@@ -144,10 +143,11 @@ func TestOutOfOrderSequenceNumber(t *testing.T) {
|
||||
tp, err := r.UpdateAndGetSnTs(extPkt)
|
||||
require.Error(t, err)
|
||||
require.ErrorIs(t, err, ErrOutOfOrderSequenceNumberCacheMiss)
|
||||
require.True(t, reflect.DeepEqual(*tp, tpExpected))
|
||||
require.Equal(t, tpExpected, *tp)
|
||||
|
||||
// add missing sequence number to the cache and try again
|
||||
r.missingSNs[23332] = 10
|
||||
r.snOffsets[SnOffsetCacheSize-1] = 10
|
||||
r.snOffsetsOccupancy++
|
||||
|
||||
tpExpected = TranslationParamsRTP{
|
||||
snOrdering: SequenceNumberOrderingOutOfOrder,
|
||||
@@ -157,7 +157,7 @@ func TestOutOfOrderSequenceNumber(t *testing.T) {
|
||||
|
||||
tp, err = r.UpdateAndGetSnTs(extPkt)
|
||||
require.NoError(t, err)
|
||||
require.True(t, reflect.DeepEqual(*tp, tpExpected))
|
||||
require.Equal(t, tpExpected, *tp)
|
||||
}
|
||||
|
||||
func TestDuplicateSequenceNumber(t *testing.T) {
|
||||
@@ -183,7 +183,7 @@ func TestDuplicateSequenceNumber(t *testing.T) {
|
||||
tp, err := r.UpdateAndGetSnTs(extPkt)
|
||||
require.Error(t, err)
|
||||
require.ErrorIs(t, err, ErrDuplicatePacket)
|
||||
require.True(t, reflect.DeepEqual(*tp, tpExpected))
|
||||
require.Equal(t, tpExpected, *tp)
|
||||
}
|
||||
|
||||
func TestPaddingOnlyPacket(t *testing.T) {
|
||||
@@ -206,10 +206,10 @@ func TestPaddingOnlyPacket(t *testing.T) {
|
||||
tp, err := r.UpdateAndGetSnTs(extPkt)
|
||||
require.Error(t, err)
|
||||
require.ErrorIs(t, err, ErrPaddingOnlyPacket)
|
||||
require.True(t, reflect.DeepEqual(*tp, tpExpected))
|
||||
require.Equal(t, tpExpected, *tp)
|
||||
require.True(t, r.highestIncomingSN == 23333)
|
||||
require.True(t, r.lastSN == 23333)
|
||||
require.EqualValues(t, 1, r.snOffset)
|
||||
require.Equal(t, uint16(1), r.snOffset)
|
||||
|
||||
// padding only packet with a gap should not report an error
|
||||
params = &testutils.TestExtPacketParams{
|
||||
@@ -228,10 +228,10 @@ func TestPaddingOnlyPacket(t *testing.T) {
|
||||
|
||||
tp, err = r.UpdateAndGetSnTs(extPkt)
|
||||
require.NoError(t, err)
|
||||
require.True(t, reflect.DeepEqual(*tp, tpExpected))
|
||||
require.Equal(t, tpExpected, *tp)
|
||||
require.True(t, r.highestIncomingSN == 23335)
|
||||
require.True(t, r.lastSN == 23334)
|
||||
require.EqualValues(t, 1, r.snOffset)
|
||||
require.Equal(t, uint16(1), r.snOffset)
|
||||
}
|
||||
|
||||
func TestGapInSequenceNumber(t *testing.T) {
|
||||
@@ -268,33 +268,19 @@ func TestGapInSequenceNumber(t *testing.T) {
|
||||
|
||||
tp, err := r.UpdateAndGetSnTs(extPkt)
|
||||
require.NoError(t, err)
|
||||
require.True(t, reflect.DeepEqual(*tp, tpExpected))
|
||||
require.Equal(t, tpExpected, *tp)
|
||||
require.True(t, r.highestIncomingSN == 1)
|
||||
require.True(t, r.lastSN == 1)
|
||||
require.EqualValues(t, 0, r.snOffset)
|
||||
require.Equal(t, uint16(0), r.snOffset)
|
||||
|
||||
// ensure missing sequence numbers got recorded in cache
|
||||
|
||||
// last received should not be in cache
|
||||
_, ok := r.missingSNs[65533]
|
||||
require.False(t, ok)
|
||||
|
||||
// three after last received missing with wrap-around
|
||||
offset, ok := r.missingSNs[65534]
|
||||
require.True(t, ok)
|
||||
require.Equal(t, uint16(0), offset)
|
||||
|
||||
offset, ok = r.missingSNs[65535]
|
||||
require.True(t, ok)
|
||||
require.Equal(t, uint16(0), offset)
|
||||
|
||||
offset, ok = r.missingSNs[0]
|
||||
require.True(t, ok)
|
||||
require.Equal(t, uint16(0), offset)
|
||||
|
||||
// current received should not be in cache
|
||||
_, ok = r.missingSNs[1]
|
||||
require.False(t, ok)
|
||||
// last received, three missing in between and current received should all be in cache
|
||||
for i := uint16(65533); i != 2; i++ {
|
||||
offset, ok := r.getSnOffset(i)
|
||||
require.True(t, ok)
|
||||
require.Equal(t, uint16(0), offset)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateAndGetPaddingSnTs(t *testing.T) {
|
||||
@@ -329,7 +315,7 @@ func TestUpdateAndGetPaddingSnTs(t *testing.T) {
|
||||
}
|
||||
snts, err := r.UpdateAndGetPaddingSnTs(numPadding, clockRate, frameRate, true)
|
||||
require.NoError(t, err)
|
||||
require.True(t, reflect.DeepEqual(snts, sntsExpected))
|
||||
require.Equal(t, sntsExpected, snts)
|
||||
|
||||
// now that there is a marker, timestamp should jump on first padding when asked again
|
||||
for i := 0; i < numPadding; i++ {
|
||||
@@ -340,7 +326,7 @@ func TestUpdateAndGetPaddingSnTs(t *testing.T) {
|
||||
}
|
||||
snts, err = r.UpdateAndGetPaddingSnTs(numPadding, clockRate, frameRate, false)
|
||||
require.NoError(t, err)
|
||||
require.True(t, reflect.DeepEqual(snts, sntsExpected))
|
||||
require.Equal(t, sntsExpected, snts)
|
||||
}
|
||||
|
||||
func TestIsOnFrameBoundary(t *testing.T) {
|
||||
|
||||
Reference in New Issue
Block a user