mirror of
https://github.com/livekit/livekit.git
synced 2026-08-28 02:54:12 +00:00
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
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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=
|
||||
|
||||
@@ -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<<offset) == 0
|
||||
return diff >= 0 && diff < cHistorySize
|
||||
}
|
||||
|
||||
// ----------------------------------
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user