diff --git a/go.mod b/go.mod index 8ec9cd47e..4c01e2787 100644 --- a/go.mod +++ b/go.mod @@ -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 diff --git a/go.sum b/go.sum index 53e811338..630b60db0 100644 --- a/go.sum +++ b/go.sum @@ -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= diff --git a/pkg/routing/roomclient.go b/pkg/routing/roomclient.go index d511fc282..f5b6d2bf4 100644 --- a/pkg/routing/roomclient.go +++ b/pkg/routing/roomclient.go @@ -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{}), diff --git a/pkg/rtc/transport.go b/pkg/rtc/transport.go index 7d285eb27..1a20d6540 100644 --- a/pkg/rtc/transport.go +++ b/pkg/rtc/transport.go @@ -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 { diff --git a/pkg/service/ioservice.go b/pkg/service/ioservice.go index a993daa83..4a653fbe5 100644 --- a/pkg/service/ioservice.go +++ b/pkg/service/ioservice.go @@ -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 } diff --git a/pkg/service/roommanager.go b/pkg/service/roommanager.go index 14caa049e..f37bfeee8 100644 --- a/pkg/service/roommanager.go +++ b/pkg/service/roommanager.go @@ -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 } diff --git a/pkg/service/wire_gen.go b/pkg/service/wire_gen.go index 69412c1fe..cbc3d041a 100644 --- a/pkg/service/wire_gen.go +++ b/pkg/service/wire_gen.go @@ -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 } diff --git a/pkg/sfu/buffer/dependencydescriptorparser.go b/pkg/sfu/buffer/dependencydescriptorparser.go index a3a2be794..ca69bd763 100644 --- a/pkg/sfu/buffer/dependencydescriptorparser.go +++ b/pkg/sfu/buffer/dependencydescriptorparser.go @@ -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++ { diff --git a/pkg/sfu/buffer/rtpstats_sender.go b/pkg/sfu/buffer/rtpstats_sender.go index ff262c11c..567a5efad 100644 --- a/pkg/sfu/buffer/rtpstats_sender.go +++ b/pkg/sfu/buffer/rtpstats_sender.go @@ -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 diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index 22d745eb5..9becfabee 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -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 diff --git a/pkg/sfu/forwarder.go b/pkg/sfu/forwarder.go index 887c48566..49c955291 100644 --- a/pkg/sfu/forwarder.go +++ b/pkg/sfu/forwarder.go @@ -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) } } diff --git a/pkg/sfu/receiver.go b/pkg/sfu/receiver.go index d0956dcf3..0f29156e1 100644 --- a/pkg/sfu/receiver.go +++ b/pkg/sfu/receiver.go @@ -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) }) }) diff --git a/pkg/sfu/rtpmunger.go b/pkg/sfu/rtpmunger.go index 1fcda126e..467ec2b3f 100644 --- a/pkg/sfu/rtpmunger.go +++ b/pkg/sfu/rtpmunger.go @@ -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 diff --git a/pkg/sfu/rtpmunger_test.go b/pkg/sfu/rtpmunger_test.go index 7c202ce1e..128d37b54 100644 --- a/pkg/sfu/rtpmunger_test.go +++ b/pkg/sfu/rtpmunger_test.go @@ -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()) } diff --git a/pkg/sfu/videolayerselector/dependencydescriptor.go b/pkg/sfu/videolayerselector/dependencydescriptor.go index afff336c3..1f8a9390b 100644 --- a/pkg/sfu/videolayerselector/dependencydescriptor.go +++ b/pkg/sfu/videolayerselector/dependencydescriptor.go @@ -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 } diff --git a/pkg/sfu/videolayerselector/dependencydescriptor_test.go b/pkg/sfu/videolayerselector/dependencydescriptor_test.go index e9e158fda..a33f3b15f 100644 --- a/pkg/sfu/videolayerselector/dependencydescriptor_test.go +++ b/pkg/sfu/videolayerselector/dependencydescriptor_test.go @@ -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)