From f29887dcd071683a67584e1e3ba2633a865fa76c Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Sat, 16 Sep 2023 02:03:50 +0530 Subject: [PATCH] Use bit map. (#2075) * Use bit map. Also, duplicate packet detection is impoetant for dropping padding only packets at the publisher side itself. In the last PR, mentioned that it is only for stats. * clean up * Update deps --- go.mod | 2 +- go.sum | 4 +- pkg/sfu/buffer/rtpstats_receiver.go | 83 +++++------------------- pkg/sfu/buffer/rtpstats_receiver_test.go | 10 +-- 4 files changed, 25 insertions(+), 74 deletions(-) diff --git a/go.mod b/go.mod index 0df7de2d2..95aec396b 100644 --- a/go.mod +++ b/go.mod @@ -18,7 +18,7 @@ require ( github.com/jxskiss/base62 v1.1.0 github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1 github.com/livekit/mediatransportutil v0.0.0-20230906055425-e81fd5f6fb3f - github.com/livekit/protocol v1.7.3-0.20230911160509-47d330eafb32 + github.com/livekit/protocol v1.7.3-0.20230915202328-cf9f95141e0e github.com/livekit/psrpc v0.3.3 github.com/mackerelio/go-osstat v0.2.4 github.com/magefile/mage v1.15.0 diff --git a/go.sum b/go.sum index 6119003f1..ae9768dd2 100644 --- a/go.sum +++ b/go.sum @@ -127,8 +127,8 @@ github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1 h1:jm09419p0lqTkD github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1/go.mod h1:Rs3MhFwutWhGwmY1VQsygw28z5bWcnEYmS1OG9OxjOQ= github.com/livekit/mediatransportutil v0.0.0-20230906055425-e81fd5f6fb3f h1:b4ri7hQESRSzJWzXXcmANG2hJ4HTj5LM01Ekm8lnQmg= github.com/livekit/mediatransportutil v0.0.0-20230906055425-e81fd5f6fb3f/go.mod h1:+WIOYwiBMive5T81V8B2wdAc2zQNRjNQiJIcPxMTILY= -github.com/livekit/protocol v1.7.3-0.20230911160509-47d330eafb32 h1:5PdmCpGGXA2hz1pKGgKSJYTjmk3Kkm+kNiW5NOFARCI= -github.com/livekit/protocol v1.7.3-0.20230911160509-47d330eafb32/go.mod h1:zbh0QPUcLGOeZeIO/VeigwWWbudz4Lv+Px94FnVfQH0= +github.com/livekit/protocol v1.7.3-0.20230915202328-cf9f95141e0e h1:WEet0iH/JazBFNhhH+YuZHtXpKefb7mnbCC2al3peyA= +github.com/livekit/protocol v1.7.3-0.20230915202328-cf9f95141e0e/go.mod h1:zbh0QPUcLGOeZeIO/VeigwWWbudz4Lv+Px94FnVfQH0= github.com/livekit/psrpc v0.3.3 h1:+lltbuN39IdaynXhLLxRShgYqYsRMWeeXKzv60oqyWo= github.com/livekit/psrpc v0.3.3/go.mod h1:n6JntEg+zT6Ji8InoyTpV7wusPNwGqqtxmHlkNhDN0U= github.com/mackerelio/go-osstat v0.2.4 h1:qxGbdPkFo65PXOb/F/nhDKpF2nGmGaCFDLXoZjJTtUs= diff --git a/pkg/sfu/buffer/rtpstats_receiver.go b/pkg/sfu/buffer/rtpstats_receiver.go index 0d41e73d7..2b409c990 100644 --- a/pkg/sfu/buffer/rtpstats_receiver.go +++ b/pkg/sfu/buffer/rtpstats_receiver.go @@ -22,6 +22,7 @@ import ( "github.com/livekit/livekit-server/pkg/sfu/utils" "github.com/livekit/protocol/livekit" + protoutils "github.com/livekit/protocol/utils" ) const ( @@ -52,7 +53,7 @@ type RTPStatsReceiver struct { timestamp *utils.WrapAround[uint32, uint64] - history [cHistorySize / 64]uint64 + history *protoutils.Bitmap[uint64] } func NewRTPStatsReceiver(params RTPStatsParams) *RTPStatsReceiver { @@ -60,6 +61,7 @@ func NewRTPStatsReceiver(params RTPStatsParams) *RTPStatsReceiver { rtpStatsBase: newRTPStatsBase(params), sequenceNumber: utils.NewWrapAround[uint16, uint64](), timestamp: utils.NewWrapAround[uint32, uint64](), + history: protoutils.NewBitmap[uint64](cHistorySize), } } @@ -173,14 +175,16 @@ func (r *RTPStatsReceiver) Update( ) } - if !r.isLost(resSN.ExtendedVal, resSN.PreExtendedHighest) { - r.bytesDuplicate += pktSize - r.headerBytesDuplicate += uint64(hdrSize) - r.packetsDuplicate++ - flowState.IsDuplicate = true - } else { - r.packetsLost-- - r.setHistory(resSN.ExtendedVal, resSN.PreExtendedHighest) + if r.isInRange(resSN.ExtendedVal, resSN.PreExtendedHighest) { + if r.history.IsSet(resSN.ExtendedVal) { + r.bytesDuplicate += pktSize + r.headerBytesDuplicate += uint64(hdrSize) + r.packetsDuplicate++ + flowState.IsDuplicate = true + } else { + r.packetsLost-- + r.history.Set(resSN.ExtendedVal) + } } flowState.IsOutOfOrder = true @@ -191,10 +195,10 @@ func (r *RTPStatsReceiver) Update( r.updateGapHistogram(int(gapSN)) // update missing sequence numbers - r.clearHistory(resSN.PreExtendedHighest+1, resSN.ExtendedVal, resSN.PreExtendedHighest) + r.history.ClearRange(resSN.PreExtendedHighest+1, resSN.ExtendedVal-1) r.packetsLost += uint64(gapSN - 1) - r.setHistory(resSN.ExtendedVal, resSN.PreExtendedHighest) + r.history.Set(resSN.ExtendedVal) if timestamp != uint32(resTS.PreExtendedHighest) { // update only on first packet as same timestamp could be in multiple packets. @@ -473,62 +477,9 @@ func (r *RTPStatsReceiver) ToProto() *livekit.RTPStats { ) } -func (r *RTPStatsReceiver) getOutOfOrderHistorySlot(esn uint64, ehsn uint64) (int, int) { +func (r *RTPStatsReceiver) isInRange(esn uint64, ehsn uint64) bool { diff := int64(ehsn - esn) - if diff >= cHistorySize || diff < 0 { - // too old OR too new (i. e. ahead of highest) - return -1, -1 - } - - return int(esn) % len(r.history), int(esn & 63) -} - -func (r *RTPStatsReceiver) getHistorySlot(esn uint64, ehsn uint64) (int, int) { - if int64(esn-ehsn) < 0 { - return r.getOutOfOrderHistorySlot(esn, ehsn) - } - - return int(esn) % len(r.history), int(esn & 63) -} - -func (r *RTPStatsReceiver) setHistory(esn uint64, ehsn uint64) { - slot, offset := r.getHistorySlot(esn, ehsn) - if slot < 0 { - return - } - - r.history[slot] |= (1 << offset) -} - -func (r *RTPStatsReceiver) clearHistory(extStartInclusive uint64, extEndExclusive uint64, ehsn uint64) { - if extEndExclusive <= extStartInclusive { - return - } - - slot, offset := r.getHistorySlot(extStartInclusive, ehsn) - if slot < 0 { - return - } - for esn := extStartInclusive; esn != extEndExclusive; esn++ { - r.history[slot] &= ^(1 << offset) - offset++ - if offset > 63 { - offset -= 64 - slot++ - if slot >= len(r.history) { - slot -= len(r.history) - } - } - } -} - -func (r *RTPStatsReceiver) isLost(esn uint64, ehsn uint64) bool { - slot, offset := r.getHistorySlot(esn, ehsn) - if slot < 0 { - return false - } - - return r.history[slot]&(1<= 0 && diff < cHistorySize } // ---------------------------------- diff --git a/pkg/sfu/buffer/rtpstats_receiver_test.go b/pkg/sfu/buffer/rtpstats_receiver_test.go index 3fb648ba0..34fea0f4b 100644 --- a/pkg/sfu/buffer/rtpstats_receiver_test.go +++ b/pkg/sfu/buffer/rtpstats_receiver_test.go @@ -224,7 +224,7 @@ func Test_RTPStatsReceiver_Update(t *testing.T) { require.Equal(t, uint64(sequenceNumber-1), flowState.LossStartInclusive) require.Equal(t, uint64(sequenceNumber), flowState.LossEndExclusive) require.Equal(t, uint64(17), r.packetsLost) - require.True(t, r.isLost(uint64(sequenceNumber)-1, r.sequenceNumber.GetExtendedHighest())) + require.False(t, r.history.IsSet(uint64(sequenceNumber)-1)) // out-of-order sequenceNumber-- @@ -242,7 +242,7 @@ func Test_RTPStatsReceiver_Update(t *testing.T) { require.False(t, flowState.HasLoss) require.Equal(t, uint64(16), r.packetsLost) require.Equal(t, uint64(4), r.packetsOutOfOrder) - require.False(t, r.isLost(uint64(sequenceNumber), r.sequenceNumber.GetExtendedHighest())) + require.True(t, r.history.IsSet(uint64(sequenceNumber))) // padding only sequenceNumber += 2 @@ -259,9 +259,9 @@ func Test_RTPStatsReceiver_Update(t *testing.T) { require.False(t, flowState.HasLoss) require.Equal(t, uint64(16), r.packetsLost) require.Equal(t, uint64(4), r.packetsOutOfOrder) - require.False(t, r.isLost(uint64(sequenceNumber), r.sequenceNumber.GetExtendedHighest())) - require.False(t, r.isLost(uint64(sequenceNumber)-1, r.sequenceNumber.GetExtendedHighest())) - require.False(t, r.isLost(uint64(sequenceNumber)-2, r.sequenceNumber.GetExtendedHighest())) + require.True(t, r.history.IsSet(uint64(sequenceNumber))) + require.True(t, r.history.IsSet(uint64(sequenceNumber)-1)) + require.True(t, r.history.IsSet(uint64(sequenceNumber)-2)) r.Stop() }