Merge remote-tracking branch 'origin/master' into raja_fr

This commit is contained in:
boks1971
2023-10-26 21:26:01 +05:30
16 changed files with 162 additions and 79 deletions
+2 -2
View File
@@ -18,8 +18,8 @@ require (
github.com/jxskiss/base62 v1.1.0
github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1
github.com/livekit/mediatransportutil v0.0.0-20231017082622-43f077b4e60e
github.com/livekit/protocol v1.8.1-0.20231024024326-07ca9d4e47bd
github.com/livekit/psrpc v0.4.0
github.com/livekit/protocol v1.8.2-0.20231026030639-f8b1277b3c7b
github.com/livekit/psrpc v0.5.0
github.com/mackerelio/go-osstat v0.2.4
github.com/magefile/mage v1.15.0
github.com/maxbrunsfeld/counterfeiter/v6 v6.7.0
+4 -4
View File
@@ -125,10 +125,10 @@ 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-20231017082622-43f077b4e60e h1:yNeIo7MSMUWgoLu7LkNKnBYnJBFPFH9Wq4S6h1kS44M=
github.com/livekit/mediatransportutil v0.0.0-20231017082622-43f077b4e60e/go.mod h1:+WIOYwiBMive5T81V8B2wdAc2zQNRjNQiJIcPxMTILY=
github.com/livekit/protocol v1.8.1-0.20231024024326-07ca9d4e47bd h1:5By2nxIMS9MApyJ654STmtwb0V4MeuaH77lqUTxV3nw=
github.com/livekit/protocol v1.8.1-0.20231024024326-07ca9d4e47bd/go.mod h1:oTWtPGfpZSJGKRrbSvDQK0jiuUylYzhiw/bnGB4Cqko=
github.com/livekit/psrpc v0.4.0 h1:oC4l99HSot/aUza4ZR3ZcZW1SRZm34KCHX3wAbmw6Lo=
github.com/livekit/psrpc v0.4.0/go.mod h1:1XYH1LLoD/YbvBvt6xg2KQ/J3InLXSJK6PL/+DKmuAU=
github.com/livekit/protocol v1.8.2-0.20231026030639-f8b1277b3c7b h1:ExuLaXyk6pGe2DVgXef7YQB0BNA7eDxidmthSkfGB2w=
github.com/livekit/protocol v1.8.2-0.20231026030639-f8b1277b3c7b/go.mod h1:l2WjlZWErS6vBlQaQyCGwWLt1aOx10XfQTsmvLjJWFQ=
github.com/livekit/psrpc v0.5.0 h1:g+yYNSs6Y1/vM7UlFkB2s/ARe2y3RKWZhX8ata5j+eo=
github.com/livekit/psrpc v0.5.0/go.mod h1:1XYH1LLoD/YbvBvt6xg2KQ/J3InLXSJK6PL/+DKmuAU=
github.com/mackerelio/go-osstat v0.2.4 h1:qxGbdPkFo65PXOb/F/nhDKpF2nGmGaCFDLXoZjJTtUs=
github.com/mackerelio/go-osstat v0.2.4/go.mod h1:Zy+qzGdZs3A9cuIqmgbJvwbmLQH9dJvtio5ZjJTbdlQ=
github.com/magefile/mage v1.15.0 h1:BvGheCMAsG3bWUDbZ8AyXXpCNwU9u5CB6sM+HNb9HYg=
+1 -3
View File
@@ -3,7 +3,6 @@ package routing
import (
"github.com/livekit/livekit-server/pkg/config"
"github.com/livekit/livekit-server/pkg/telemetry/prometheus"
"github.com/livekit/protocol/livekit"
"github.com/livekit/protocol/logger"
protopsrpc "github.com/livekit/protocol/psrpc"
"github.com/livekit/protocol/rpc"
@@ -11,9 +10,8 @@ import (
"github.com/livekit/psrpc/pkg/middleware"
)
func NewRoomClient(nodeID livekit.NodeID, bus psrpc.MessageBus, config config.PSRPCConfig) (rpc.TypedRoomClient, error) {
func NewRoomClient(bus psrpc.MessageBus, config config.PSRPCConfig) (rpc.TypedRoomClient, error) {
return rpc.NewTypedRoomClient(
nodeID,
bus,
protopsrpc.WithClientLogger(logger.GetLogger()),
middleware.WithClientMetrics(prometheus.PSRPCMetricsObserver{}),
+3
View File
@@ -322,6 +322,9 @@ func newPeerConnection(params TransportParams, onBandwidthEstimator func(estimat
params.Logger.Infow("client doesn't support prflx over relay, use external ip only as host candidate", "ips", nat1to1Ips)
se.SetNAT1To1IPs(nat1to1Ips, webrtc.ICECandidateTypeHost)
se.SetIPFilter(func(ip net.IP) bool {
if ip.To4() == nil {
return true
}
ipstr := ip.String()
for _, inc := range includeIps {
if inc == ipstr {
+1 -1
View File
@@ -52,7 +52,7 @@ func NewIOInfoService(
}
if bus != nil {
ioServer, err := rpc.NewIOInfoServer(string(nodeID), s, bus)
ioServer, err := rpc.NewIOInfoServer(s, bus)
if err != nil {
return nil, err
}
+1 -1
View File
@@ -119,7 +119,7 @@ func NewLocalRoomManager(
},
}
r.roomServer, err = rpc.NewTypedRoomServer(livekit.NodeID(r.currentNode.Id), r, bus)
r.roomServer, err = rpc.NewTypedRoomServer(r, bus)
if err != nil {
return nil, err
}
+3 -3
View File
@@ -53,7 +53,7 @@ func InitializeServer(conf *config.Config, currentNode routing.LocalNode) (*Live
if err != nil {
return nil, err
}
egressClient, err := rpc.NewEgressClient(nodeID, messageBus)
egressClient, err := rpc.NewEgressClient(messageBus)
if err != nil {
return nil, err
}
@@ -75,7 +75,7 @@ func InitializeServer(conf *config.Config, currentNode routing.LocalNode) (*Live
}
rtcEgressLauncher := NewEgressLauncher(egressClient, ioInfoService)
topicFormatter := routing.NewTopicFormatter()
roomClient, err := routing.NewRoomClient(nodeID, messageBus, psrpcConfig)
roomClient, err := routing.NewRoomClient(messageBus, psrpcConfig)
if err != nil {
return nil, err
}
@@ -85,7 +85,7 @@ func InitializeServer(conf *config.Config, currentNode routing.LocalNode) (*Live
}
egressService := NewEgressService(egressClient, objectStore, ioInfoService, roomService, rtcEgressLauncher)
ingressConfig := getIngressConfig(conf)
ingressClient, err := rpc.NewIngressClient(nodeID, messageBus)
ingressClient, err := rpc.NewIngressClient(messageBus)
if err != nil {
return nil, err
}
@@ -142,6 +142,12 @@ type DependencyDescriptorDecodeTarget struct {
Layer VideoLayer
}
func (dt *DependencyDescriptorDecodeTarget) String() string {
return fmt.Sprintf("DecodeTarget{t: %d, l: %+v}", dt.Target, dt.Layer)
}
// ------------------------------------------------------------------------------
func ProcessFrameDependencyStructure(structure *dd.FrameDependencyStructure) []DependencyDescriptorDecodeTarget {
decodeTargets := make([]DependencyDescriptorDecodeTarget, 0, structure.NumDecodeTargets)
for target := 0; target < structure.NumDecodeTargets; target++ {
+12 -5
View File
@@ -541,12 +541,16 @@ func (r *RTPStatsSender) UpdateFromReceiverReport(rr rtcp.ReceptionReport) (rtt
eis := &s.intervalStats
eis.aggregate(&is)
if is.packetsNotFound != 0 {
timeSinceLastRR := time.Since(r.lastRRTime)
if r.lastRRTime.IsZero() {
timeSinceLastRR = time.Since(r.startTime)
}
if r.metadataCacheOverflowCount%10 == 0 {
r.logger.Infow(
"metadata cache overflow",
"lastRRTime", r.lastRRTime.String(),
"lastRR", r.lastRR,
"sinceLastRR", time.Since(r.lastRRTime).String(),
"timeSinceLastRR", timeSinceLastRR.String(),
"receivedRR", rr,
"extStartSN", r.extStartSN,
"extHighestSN", r.extHighestSN,
@@ -636,10 +640,10 @@ func (r *RTPStatsSender) GetRtcpSenderReport(ssrc uint32, calculatedClockRate ui
At: now,
}
if r.srNewest != nil {
timeSinceLastReport := nowNTP.Time().Sub(r.srNewest.NTPTimestamp.Time()).Seconds()
timeSinceLastReport := nowNTP.Time().Sub(r.srNewest.NTPTimestamp.Time())
rtpDiffSinceLastReport := nowRTPExt - r.srNewest.RTPTimestampExt
windowClockRate := float64(rtpDiffSinceLastReport) / timeSinceLastReport
if timeSinceLastReport > 0.2 && math.Abs(float64(r.params.ClockRate)-windowClockRate) > 0.2*float64(r.params.ClockRate) {
windowClockRate := float64(rtpDiffSinceLastReport) / timeSinceLastReport.Seconds()
if timeSinceLastReport.Seconds() > 0.2 && math.Abs(float64(r.params.ClockRate)-windowClockRate) > 0.2*float64(r.params.ClockRate) {
if r.clockSkewCount%10 == 0 {
r.logger.Infow(
"sending sender report, clock skew",
@@ -655,7 +659,7 @@ func (r *RTPStatsSender) GetRtcpSenderReport(ssrc uint32, calculatedClockRate ui
"nowRTPExtUsingTime", nowRTPExtUsingTime,
"calculatedClockRate", calculatedClockRate,
"nowRTPExtUsingRate", nowRTPExtUsingRate,
"timeSinceLastReport", timeSinceLastReport,
"timeSinceLastReport", timeSinceLastReport.String(),
"rtpDiffSinceLastReport", rtpDiffSinceLastReport,
"windowClockRate", windowClockRate,
"count", r.clockSkewCount,
@@ -700,6 +704,9 @@ func (r *RTPStatsSender) GetRtcpSenderReport(ssrc uint32, calculatedClockRate ui
ntpDiffSinceLast := nowNTP.Time().Sub(r.srNewest.NTPTimestamp.Time())
nowRTPExt = r.srNewest.RTPTimestampExt + uint64(ntpDiffSinceLast.Seconds()*float64(r.params.ClockRate))
nowRTP = uint32(nowRTPExt)
srData.RTPTimestamp = nowRTP
srData.RTPTimestampExt = nowRTPExt
}
r.srNewest = srData
+7 -5
View File
@@ -54,7 +54,7 @@ type TrackSender interface {
ID() string
SubscriberID() livekit.ParticipantID
TrackInfoAvailable()
HandleRTCPSenderReportData(payloadType webrtc.PayloadType, layer int32, srData *buffer.RTCPSenderReportData) error
HandleRTCPSenderReportData(payloadType webrtc.PayloadType, isSVC bool, layer int32, srData *buffer.RTCPSenderReportData) error
HandleTrackFrameRateReport(payloadType webrtc.PayloadType, fps [][]float32) error
}
@@ -1464,9 +1464,11 @@ func (d *DownTrack) handleRTCP(bytes []byte) {
pliOnce := true
sendPliOnce := func() {
_, layer := d.forwarder.CheckSync()
isAnyMuted := d.forwarder.IsAnyMuted()
d.params.Logger.Debugw("received PLI/FIR RTCP", "layer", layer, "isAnyMuted", isAnyMuted)
if pliOnce {
_, layer := d.forwarder.CheckSync()
if layer != buffer.InvalidLayerSpatial && !d.forwarder.IsAnyMuted() {
if layer != buffer.InvalidLayerSpatial && !isAnyMuted {
d.params.Logger.Debugw("sending PLI RTCP", "layer", layer)
d.params.Receiver.SendPLI(layer, false)
d.isNACKThrottled.Store(true)
@@ -1892,8 +1894,8 @@ func (d *DownTrack) sendSilentFrameOnMuteForOpus() {
}
}
func (d *DownTrack) HandleRTCPSenderReportData(_payloadType webrtc.PayloadType, layer int32, srData *buffer.RTCPSenderReportData) error {
if layer == d.forwarder.GetReferenceLayerSpatial() && srData != nil {
func (d *DownTrack) HandleRTCPSenderReportData(_payloadType webrtc.PayloadType, isSVC bool, layer int32, srData *buffer.RTCPSenderReportData) error {
if (layer == d.forwarder.GetReferenceLayerSpatial() || (layer == 0 && isSVC)) && srData != nil {
d.rtpStats.MaybeAdjustFirstPacketTime(srData.RTPTimestamp + uint32(d.forwarder.GetReferenceTimestampOffset()))
}
return nil
+2 -2
View File
@@ -1653,7 +1653,7 @@ func (f *Forwarder) getTranslationParamsCommon(extPkt *buffer.ExtPacket, layer i
f.lastSSRC = extPkt.Packet.SSRC
}
tpRTP, err := f.rtpMunger.UpdateAndGetSnTs(extPkt)
tpRTP, err := f.rtpMunger.UpdateAndGetSnTs(extPkt, tp.marker)
if err != nil {
tp.shouldDrop = true
if err == ErrPaddingOnlyPacket || err == ErrDuplicatePacket || err == ErrOutOfOrderSequenceNumberCacheMiss {
@@ -1691,7 +1691,7 @@ func (f *Forwarder) getTranslationParamsVideo(extPkt *buffer.ExtPacket, layer in
tp.shouldDrop = true
if f.started && result.IsRelevant {
// call to update highest incoming sequence number and other internal structures
if tpRTP, err := f.rtpMunger.UpdateAndGetSnTs(extPkt); err == nil && tpRTP.snOrdering == SequenceNumberOrderingContiguous {
if tpRTP, err := f.rtpMunger.UpdateAndGetSnTs(extPkt, result.RTPMarker); err == nil && tpRTP.snOrdering == SequenceNumberOrderingContiguous {
f.rtpMunger.PacketDropped(extPkt)
}
}
+1 -1
View File
@@ -343,7 +343,7 @@ func (w *WebRTCReceiver) AddUpTrack(track *webrtc.TrackRemote, buff *buffer.Buff
w.streamTrackerManager.SetRTCPSenderReportData(layer, srFirst, srNewest)
w.downTrackSpreader.Broadcast(func(dt TrackSender) {
_ = dt.HandleRTCPSenderReportData(w.codec.PayloadType, layer, srNewest)
_ = dt.HandleRTCPSenderReportData(w.codec.PayloadType, w.isSVC, layer, srNewest)
})
})
+27 -12
View File
@@ -51,14 +51,21 @@ type SnTs struct {
// ----------------------------------------------------------------------
type RTPMungerState struct {
ExtLastSN uint64
ExtSecondLastSN uint64
ExtLastTS uint64
ExtSecondLastTS uint64
ExtLastSN uint64
ExtSecondLastSN uint64
ExtLastTS uint64
ExtSecondLastTS uint64
LastMarker bool
SecondLastMarker bool
}
func (r RTPMungerState) String() string {
return fmt.Sprintf("RTPMungerState{extLastSN: %d, extSecondLastSN: %d, extLastTS: %d, extSecondLastTS: %d)", r.ExtLastSN, r.ExtSecondLastSN, r.ExtLastTS, r.ExtSecondLastTS)
return fmt.Sprintf(
"RTPMungerState{extLastSN: %d, extSecondLastSN: %d, extLastTS: %d, extSecondLastTS: %d, lastMarker: %v, secondLastMarker: %v)",
r.ExtLastSN, r.ExtSecondLastSN,
r.ExtLastTS, r.ExtSecondLastTS,
r.LastMarker, r.SecondLastMarker,
)
}
// ----------------------------------------------------------------------
@@ -77,7 +84,8 @@ type RTPMunger struct {
extSecondLastTS uint64
tsOffset uint64
lastMarker bool
lastMarker bool
secondLastMarker bool
extRtxGateSn uint64
isInRtxGateRegion bool
@@ -100,15 +108,18 @@ func (r *RTPMunger) DebugInfo() map[string]interface{} {
"ExtSecondLastTS": r.extSecondLastTS,
"TSOffset": r.tsOffset,
"LastMarker": r.lastMarker,
"SecondLastMarker": r.secondLastMarker,
}
}
func (r *RTPMunger) GetLast() RTPMungerState {
return RTPMungerState{
ExtLastSN: r.extLastSN,
ExtSecondLastSN: r.extSecondLastSN,
ExtLastTS: r.extLastTS,
ExtSecondLastTS: r.extSecondLastTS,
ExtLastSN: r.extLastSN,
ExtSecondLastSN: r.extSecondLastSN,
ExtLastTS: r.extLastTS,
ExtSecondLastTS: r.extSecondLastTS,
LastMarker: r.lastMarker,
SecondLastMarker: r.secondLastMarker,
}
}
@@ -117,6 +128,8 @@ func (r *RTPMunger) SeedLast(state RTPMungerState) {
r.extSecondLastSN = state.ExtSecondLastSN
r.extLastTS = state.ExtLastTS
r.extSecondLastTS = state.ExtSecondLastTS
r.lastMarker = state.LastMarker
r.secondLastMarker = state.SecondLastMarker
}
func (r *RTPMunger) SetLastSnTs(extPkt *buffer.ExtPacket) {
@@ -164,9 +177,10 @@ func (r *RTPMunger) PacketDropped(extPkt *buffer.ExtPacket) {
r.updateSnOffset()
r.extLastTS = r.extSecondLastTS
r.lastMarker = r.secondLastMarker
}
func (r *RTPMunger) UpdateAndGetSnTs(extPkt *buffer.ExtPacket) (*TranslationParamsRTP, error) {
func (r *RTPMunger) UpdateAndGetSnTs(extPkt *buffer.ExtPacket, marker bool) (*TranslationParamsRTP, error) {
diff := int64(extPkt.ExtSequenceNumber - r.extHighestIncomingSN)
if (diff == 1 && len(extPkt.Packet.Payload) != 0) || diff > 1 {
// in-order - either contiguous packet with payload OR packet following a gap, may or may not have payload
@@ -184,7 +198,8 @@ func (r *RTPMunger) UpdateAndGetSnTs(extPkt *buffer.ExtPacket) (*TranslationPara
r.extLastSN = extMungedSN
r.extSecondLastTS = r.extLastTS
r.extLastTS = extMungedTS
r.lastMarker = extPkt.Packet.Marker
r.secondLastMarker = r.lastMarker
r.lastMarker = marker
if extPkt.KeyFrame {
r.extRtxGateSn = extMungedSN
+22 -23
View File
@@ -103,7 +103,7 @@ func TestPacketDropped(t *testing.T) {
require.Equal(t, uint64(0), snOffset)
require.Equal(t, uint64(0), r.tsOffset)
r.UpdateAndGetSnTs(extPkt) // update sequence number offset
r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker) // update sequence number offset
// drop a non-head packet, should cause no change in internals
params = &testutils.TestExtPacketParams{
@@ -128,7 +128,7 @@ func TestPacketDropped(t *testing.T) {
}
extPkt, _ = testutils.GetTestExtPacket(params)
r.UpdateAndGetSnTs(extPkt) // update sequence number offset
r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker) // update sequence number offset
require.Equal(t, uint64(44444), r.extLastSN)
r.PacketDropped(extPkt)
@@ -147,7 +147,7 @@ func TestPacketDropped(t *testing.T) {
}
extPkt, _ = testutils.GetTestExtPacket(params)
r.UpdateAndGetSnTs(extPkt) // update sequence number offset
r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker) // update sequence number offset
require.Equal(t, r.extLastSN, uint64(44444))
snOffset, err = r.snRangeMap.GetValue(r.extHighestIncomingSN)
require.NoError(t, err)
@@ -165,7 +165,7 @@ func TestOutOfOrderSequenceNumber(t *testing.T) {
}
extPkt, _ := testutils.GetTestExtPacket(params)
r.SetLastSnTs(extPkt)
r.UpdateAndGetSnTs(extPkt)
r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
// should not be able to add a missing sequence number to the cache that is before start
err := r.snRangeMap.ExcludeRange(23332, 23333)
@@ -180,7 +180,7 @@ func TestOutOfOrderSequenceNumber(t *testing.T) {
}
extPkt, _ = testutils.GetTestExtPacket(params)
tp, err := r.UpdateAndGetSnTs(extPkt)
tp, err := r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.Error(t, err)
// add a missing sequence number to the cache
@@ -194,7 +194,7 @@ func TestOutOfOrderSequenceNumber(t *testing.T) {
PayloadSize: 10,
}
extPkt, _ = testutils.GetTestExtPacket(params)
r.UpdateAndGetSnTs(extPkt)
r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
// out-of-order sequence number should be munged from cache
params = &testutils.TestExtPacketParams{
@@ -211,7 +211,7 @@ func TestOutOfOrderSequenceNumber(t *testing.T) {
extTimestamp: 0xabcdef,
}
tp, err = r.UpdateAndGetSnTs(extPkt)
tp, err = r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.NoError(t, err)
require.Equal(t, tpExpected, *tp)
@@ -227,7 +227,7 @@ func TestOutOfOrderSequenceNumber(t *testing.T) {
snOrdering: SequenceNumberOrderingOutOfOrder,
}
tp, err = r.UpdateAndGetSnTs(extPkt)
tp, err = r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.Error(t, err, ErrOutOfOrderSequenceNumberCacheMiss)
require.Equal(t, tpExpected, *tp)
}
@@ -244,15 +244,14 @@ func TestDuplicateSequenceNumber(t *testing.T) {
r.SetLastSnTs(extPkt)
// send first packet through
r.UpdateAndGetSnTs(extPkt)
r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
// send it again - duplicate packet
tpExpected := TranslationParamsRTP{
snOrdering: SequenceNumberOrderingDuplicate,
}
tp, err := r.UpdateAndGetSnTs(extPkt)
require.Error(t, err)
tp, err := r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.ErrorIs(t, err, ErrDuplicatePacket)
require.Equal(t, tpExpected, *tp)
}
@@ -273,7 +272,7 @@ func TestPaddingOnlyPacket(t *testing.T) {
snOrdering: SequenceNumberOrderingContiguous,
}
tp, err := r.UpdateAndGetSnTs(extPkt)
tp, err := r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.Error(t, err)
require.ErrorIs(t, err, ErrPaddingOnlyPacket)
require.Equal(t, tpExpected, *tp)
@@ -296,7 +295,7 @@ func TestPaddingOnlyPacket(t *testing.T) {
extTimestamp: 0xabcdef,
}
tp, err = r.UpdateAndGetSnTs(extPkt)
tp, err = r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.NoError(t, err)
require.Equal(t, tpExpected, *tp)
require.Equal(t, uint64(23335), r.extHighestIncomingSN)
@@ -318,7 +317,7 @@ func TestGapInSequenceNumber(t *testing.T) {
extPkt, _ := testutils.GetTestExtPacket(params)
r.SetLastSnTs(extPkt)
_, err := r.UpdateAndGetSnTs(extPkt)
_, err := r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.NoError(t, err)
// three lost packets
@@ -337,7 +336,7 @@ func TestGapInSequenceNumber(t *testing.T) {
extTimestamp: 0xabcdef,
}
tp, err := r.UpdateAndGetSnTs(extPkt)
tp, err := r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.NoError(t, err)
require.Equal(t, tpExpected, *tp)
require.Equal(t, uint64(65536+1), r.extHighestIncomingSN)
@@ -366,7 +365,7 @@ func TestGapInSequenceNumber(t *testing.T) {
snOrdering: SequenceNumberOrderingContiguous,
}
tp, err = r.UpdateAndGetSnTs(extPkt)
tp, err = r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.ErrorIs(t, err, ErrPaddingOnlyPacket)
require.Equal(t, tpExpected, *tp)
require.Equal(t, uint64(65536+2), r.extHighestIncomingSN)
@@ -390,7 +389,7 @@ func TestGapInSequenceNumber(t *testing.T) {
extTimestamp: 0xabcdef,
}
tp, err = r.UpdateAndGetSnTs(extPkt)
tp, err = r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.NoError(t, err)
require.Equal(t, tpExpected, *tp)
require.Equal(t, uint64(65536+4), r.extHighestIncomingSN)
@@ -417,7 +416,7 @@ func TestGapInSequenceNumber(t *testing.T) {
snOrdering: SequenceNumberOrderingContiguous,
}
tp, err = r.UpdateAndGetSnTs(extPkt)
tp, err = r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.ErrorIs(t, err, ErrPaddingOnlyPacket)
require.Equal(t, tpExpected, *tp)
require.Equal(t, uint64(65536+5), r.extHighestIncomingSN)
@@ -441,7 +440,7 @@ func TestGapInSequenceNumber(t *testing.T) {
extTimestamp: 0xabcdef,
}
tp, err = r.UpdateAndGetSnTs(extPkt)
tp, err = r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.NoError(t, err)
require.Equal(t, tpExpected, *tp)
require.Equal(t, uint64(65536+7), r.extHighestIncomingSN)
@@ -474,7 +473,7 @@ func TestGapInSequenceNumber(t *testing.T) {
extTimestamp: 0xabcdef,
}
tp, err = r.UpdateAndGetSnTs(extPkt)
tp, err = r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.NoError(t, err)
require.Equal(t, tpExpected, *tp)
require.Equal(t, uint64(65536+7), r.extHighestIncomingSN)
@@ -497,7 +496,7 @@ func TestGapInSequenceNumber(t *testing.T) {
extTimestamp: 0xabcdef,
}
tp, err = r.UpdateAndGetSnTs(extPkt)
tp, err = r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.NoError(t, err)
require.Equal(t, tpExpected, *tp)
require.Equal(t, uint64(65536+7), r.extHighestIncomingSN)
@@ -565,7 +564,7 @@ func TestIsOnFrameBoundary(t *testing.T) {
r.SetLastSnTs(extPkt)
// send it through
_, err := r.UpdateAndGetSnTs(extPkt)
_, err := r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.NoError(t, err)
require.False(t, r.IsOnFrameBoundary())
@@ -580,7 +579,7 @@ func TestIsOnFrameBoundary(t *testing.T) {
extPkt, _ = testutils.GetTestExtPacket(params)
// send it through
_, err = r.UpdateAndGetSnTs(extPkt)
_, err = r.UpdateAndGetSnTs(extPkt, extPkt.Packet.Marker)
require.NoError(t, err)
require.True(t, r.IsOnFrameBoundary())
}
@@ -63,6 +63,7 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r
ddwdt := extPkt.DependencyDescriptor
if ddwdt == nil {
// packet doesn't have dependency descriptor
d.logger.Debugw(fmt.Sprintf("drop packet, no DD, incoming %v, sn: %d, isKeyFrame: %v", extPkt.VideoLayer, extPkt.Packet.SequenceNumber, extPkt.KeyFrame))
return
}
@@ -80,11 +81,23 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r
sd, err := d.decisions.GetDecision(extFrameNum)
if err != nil {
// do not mark as dropped as only error is an old frame
d.logger.Debugw(fmt.Sprintf("drop packet on decision error, incoming %v, fn: %d/%d, sn: %d",
incomingLayer,
dd.FrameNumber,
extFrameNum,
extPkt.Packet.SequenceNumber,
), "err", err)
return
}
switch sd {
case selectorDecisionDropped:
// a packet of an alreadty dropped frame, maintain decision
d.logger.Debugw(fmt.Sprintf("drop packet already dropped, incoming %v, fn: %d/%d, sm: %d",
incomingLayer,
dd.FrameNumber,
extFrameNum,
extPkt.Packet.SequenceNumber,
))
return
}
@@ -118,9 +131,7 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r
if err != nil {
d.decodeTargetsLock.RUnlock()
// dtis error, dependency descriptor might lost
d.logger.Debugw(fmt.Sprintf("drop packet for frame detection error, incoming: %v",
incomingLayer,
), "err", err)
d.logger.Debugw(fmt.Sprintf("drop packet for frame detection error, incoming: %v", incomingLayer), "err", err)
d.decisions.AddDropped(extFrameNum)
return
}
@@ -135,23 +146,34 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r
if highestDecodeTarget.Target < 0 {
// no active decode target, do not select
// d.logger.Debugw(fmt.Sprintf("drop packet for no target found, decodeTargets %v, tagetLayer %v, incoming %v",
// d.decodeTargets,
// d.targetLayer,
// incomingLayer,
// ))
d.logger.Debugw(
"drop packet for no target found",
"highestDecodeTarget", highestDecodeTarget,
"decodeTargets", d.decodeTargets,
"tagetLayer", d.targetLayer,
"incoming", incomingLayer,
"fn", dd.FrameNumber,
"efn", extFrameNum,
"sn", extPkt.Packet.SequenceNumber,
"isKeyFrame", extPkt.KeyFrame,
)
d.decisions.AddDropped(extFrameNum)
return
}
// DD-TODO : if bandwidth in congest, could drop the 'Discardable' frame
if dti == dede.DecodeTargetNotPresent {
// d.logger.Debugw(fmt.Sprintf("drop packet for decode target not present, highestDecodeTarget %d, incoming %v, fn: %d/%d",
// highestDecodeTarget,
// incomingLayer,
// dd.FrameNumber,
// extFrameNum,
// ))
d.logger.Debugw(
"drop packet for decode target not present",
"highestDecodeTarget", highestDecodeTarget,
"decodeTargets", d.decodeTargets,
"tagetLayer", d.targetLayer,
"incoming", incomingLayer,
"fn", dd.FrameNumber,
"efn", extFrameNum,
"sn", extPkt.Packet.SequenceNumber,
"isKeyFrame", extPkt.KeyFrame,
)
d.decisions.AddDropped(extFrameNum)
return
}
@@ -171,6 +193,17 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r
}
}
if !isDecodable {
d.logger.Debugw(
"drop packet for not decodable",
"highestDecodeTarget", highestDecodeTarget,
"decodeTargets", d.decodeTargets,
"tagetLayer", d.targetLayer,
"incoming", incomingLayer,
"fn", dd.FrameNumber,
"efn", extFrameNum,
"sn", extPkt.Packet.SequenceNumber,
"isKeyFrame", extPkt.KeyFrame,
)
d.decisions.AddDropped(extFrameNum)
return
}
@@ -188,7 +221,10 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r
"req", d.requestSpatial,
"maxSeen", d.maxSeenLayer,
"feed", extPkt.Packet.SSRC,
"frame", extFrameNum,
"fn", dd.FrameNumber,
"efn", extFrameNum,
"sn", extPkt.Packet.SequenceNumber,
"isKeyFrame", extPkt.KeyFrame,
)
}
@@ -197,7 +233,16 @@ func (d *DependencyDescriptor) Select(extPkt *buffer.ExtPacket, _layer int32) (r
d.previousActiveDecodeTargetsBitmask = d.activeDecodeTargetsBitmask
d.activeDecodeTargetsBitmask = buffer.GetActiveDecodeTargetBitmask(d.currentLayer, ddwdt.DecodeTargets)
d.logger.Debugw("switch to target", "highest", highestDecodeTarget.Layer, "current", d.currentLayer, "bitmask", *d.activeDecodeTargetsBitmask, "frame", extFrameNum)
d.logger.Debugw(
"switch to target",
"highestDecodeTarget", highestDecodeTarget,
"current", d.currentLayer,
"bitmask", *d.activeDecodeTargetsBitmask,
"fn", dd.FrameNumber,
"efn", extFrameNum,
"sn", extPkt.Packet.SequenceNumber,
"isKeyFrame", extPkt.KeyFrame,
)
}
ddExtension := &dede.DependencyDescriptorExtension{
@@ -282,12 +327,19 @@ func (d *DependencyDescriptor) updateActiveDecodeTargets(activeDecodeTargetsBitm
func (d *DependencyDescriptor) CheckSync() (locked bool, layer int32) {
layer = d.GetRequestSpatial()
if !d.currentLayer.IsValid() {
// always declare not locked when trying to resume from nothing
return false, layer
}
d.decodeTargetsLock.RLock()
defer d.decodeTargetsLock.RUnlock()
for _, dt := range d.decodeTargets {
if dt.Active() && dt.Layer.Spatial == layer && dt.Valid() {
d.logger.Debugw(fmt.Sprintf("checking sync, matching decode target, layer: %d, dt: %s, dts: %+v", layer, dt, d.decodeTargets))
return true, layer
}
}
return false, layer
}
@@ -137,7 +137,7 @@ func TestDependencyDescriptor(t *testing.T) {
ddSelector.SetRequestSpatial(1)
// no dd ext, dropped
ret := ddSelector.Select(&buffer.ExtPacket{}, 0)
ret := ddSelector.Select(&buffer.ExtPacket{Packet: &rtp.Packet{}}, 0)
require.False(t, ret.IsSelected)
require.True(t, ret.IsRelevant)
@@ -153,6 +153,7 @@ func TestDependencyDescriptor(t *testing.T) {
},
},
},
Packet: &rtp.Packet{},
}, 0)
require.False(t, ret.IsSelected)
require.True(t, ret.IsRelevant)