From f97242c8ba1a440a417e49fbe107ac46a5deb486 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Wed, 18 Oct 2023 21:48:41 +0530 Subject: [PATCH] Use 32-bit time stamp to get reference time stamp on a switch. (#2153) * Use 32-bit time stamp to get reference time stamp on a switch. With relay and dyncast and migration, it is possible that different layers of a simulcast get out of sync in terms of extended type, i. e. layer 0 could keep running and its timestamp could have wrapped around and bumped the extended timestamp. But, another layer could start and stop. One possible solution is sending the extended timestamp across relay. But, that breaks down during migration if publisher has started afresh. Subscriber could still be using extended range. So, use 32-bit timestamp to infer reference timestamp and patch it with expected extended time stamp to derive the extended reference. * use calculated value * make it test friendly --- pkg/rtc/wrappedreceiver.go | 4 ++-- pkg/sfu/buffer/rtpstats_base.go | 12 +++++------ pkg/sfu/buffer/rtpstats_receiver.go | 30 +++++++++++++++++++++++--- pkg/sfu/buffer/rtpstats_sender.go | 33 +++++++++++++++++++++++++---- pkg/sfu/downtrack.go | 2 +- pkg/sfu/forwarder.go | 30 +++++++++++++++++++++----- pkg/sfu/receiver.go | 6 +++--- pkg/sfu/rtpmunger.go | 18 +++++++++++++++- pkg/sfu/streamtrackermanager.go | 10 ++++----- 9 files changed, 115 insertions(+), 30 deletions(-) diff --git a/pkg/rtc/wrappedreceiver.go b/pkg/rtc/wrappedreceiver.go index 853d2c4f4..593a62487 100644 --- a/pkg/rtc/wrappedreceiver.go +++ b/pkg/rtc/wrappedreceiver.go @@ -317,9 +317,9 @@ func (d *DummyReceiver) GetCalculatedClockRate(layer int32) uint32 { return 0 } -func (d *DummyReceiver) GetReferenceLayerRTPTimestamp(ets uint64, layer int32, referenceLayer int32) (uint64, error) { +func (d *DummyReceiver) GetReferenceLayerRTPTimestamp(ts uint32, layer int32, referenceLayer int32) (uint32, error) { if r, ok := d.receiver.Load().(sfu.TrackReceiver); ok { - return r.GetReferenceLayerRTPTimestamp(ets, layer, referenceLayer) + return r.GetReferenceLayerRTPTimestamp(ts, layer, referenceLayer) } return 0, errors.New("receiver not available") } diff --git a/pkg/sfu/buffer/rtpstats_base.go b/pkg/sfu/buffer/rtpstats_base.go index 9f8404e77..0c77299ae 100644 --- a/pkg/sfu/buffer/rtpstats_base.go +++ b/pkg/sfu/buffer/rtpstats_base.go @@ -453,7 +453,7 @@ func (r *rtpStatsBase) GetRtt() uint32 { return r.rtt } -func (r *rtpStatsBase) maybeAdjustFirstPacketTime(ets uint64, extStartTS uint64) { +func (r *rtpStatsBase) maybeAdjustFirstPacketTime(ts uint32, startTS uint32) { if time.Since(r.startTime) > cFirstPacketTimeAdjustWindow { return } @@ -464,7 +464,7 @@ func (r *rtpStatsBase) maybeAdjustFirstPacketTime(ets uint64, extStartTS uint64) // abnormal delay (maybe due to pacing or maybe due to queuing // in some network element along the way), push back first time // to an earlier instance. - samplesDiff := int64(ets - extStartTS) + samplesDiff := int32(ts - startTS) if samplesDiff < 0 { // out-of-order, skip return @@ -482,8 +482,8 @@ func (r *rtpStatsBase) maybeAdjustFirstPacketTime(ets uint64, extStartTS uint64) "before", r.firstTime.String(), "after", firstTime.String(), "adjustment", r.firstTime.Sub(firstTime).String(), - "extNowTS", ets, - "extStartTS", extStartTS, + "nowTS", ts, + "startTS", startTS, ) if r.firstTime.Sub(firstTime) > cFirstPacketTimeAdjustThreshold { r.logger.Infow("first packet time adjustment too big, ignoring", @@ -492,8 +492,8 @@ func (r *rtpStatsBase) maybeAdjustFirstPacketTime(ets uint64, extStartTS uint64) "before", r.firstTime.String(), "after", firstTime.String(), "adjustment", r.firstTime.Sub(firstTime).String(), - "extNowTS", ets, - "extStartTS", extStartTS, + "nowTS", ts, + "startTS", startTS, ) } else { r.firstTime = firstTime diff --git a/pkg/sfu/buffer/rtpstats_receiver.go b/pkg/sfu/buffer/rtpstats_receiver.go index 0a2ccf0a9..5902f57e8 100644 --- a/pkg/sfu/buffer/rtpstats_receiver.go +++ b/pkg/sfu/buffer/rtpstats_receiver.go @@ -151,7 +151,19 @@ func (r *RTPStatsReceiver) Update( } } if -gapSN >= cNumSequenceNumbers { - r.logger.Warnw("large sequence number gap negative", nil, "prev", resSN.PreExtendedHighest, "curr", resSN.ExtendedVal, "gap", gapSN) + r.logger.Warnw( + "large sequence number gap negative", nil, + "prev", resSN.PreExtendedHighest, + "curr", resSN.ExtendedVal, + "gap", gapSN, + "packetTime", packetTime.String(), + "sequenceNumber", sequenceNumber, + "timestamp", timestamp, + "marker", marker, + "hdrSize", hdrSize, + "payloadSize", payloadSize, + "paddingSize", paddingSize, + ) } if gapSN != 0 { @@ -205,7 +217,19 @@ func (r *RTPStatsReceiver) Update( flowState.ExtTimestamp = resTS.ExtendedVal } else { // in-order if gapSN >= cNumSequenceNumbers { - r.logger.Warnw("large sequence number gap", nil, "prev", resSN.PreExtendedHighest, "curr", resSN.ExtendedVal, "gap", gapSN) + r.logger.Warnw( + "large sequence number gap", nil, + "prev", resSN.PreExtendedHighest, + "curr", resSN.ExtendedVal, + "gap", gapSN, + "packetTime", packetTime.String(), + "sequenceNumber", sequenceNumber, + "timestamp", timestamp, + "marker", marker, + "hdrSize", hdrSize, + "payloadSize", payloadSize, + "paddingSize", paddingSize, + ) } // update gap histogram @@ -284,7 +308,7 @@ func (r *RTPStatsReceiver) SetRtcpSenderReportData(srData *RTCPSenderReportData) srDataCopy := *srData srDataCopy.RTPTimestampExt = uint64(srDataCopy.RTPTimestamp) + tsCycles - r.maybeAdjustFirstPacketTime(srDataCopy.RTPTimestampExt, r.timestamp.GetExtendedStart()) + r.maybeAdjustFirstPacketTime(srDataCopy.RTPTimestamp, r.timestamp.GetStart()) if r.srNewest != nil && srDataCopy.RTPTimestampExt < r.srNewest.RTPTimestampExt { // This can happen when a track is replaced with a null and then restored - diff --git a/pkg/sfu/buffer/rtpstats_sender.go b/pkg/sfu/buffer/rtpstats_sender.go index d3ea160a5..2ac196904 100644 --- a/pkg/sfu/buffer/rtpstats_sender.go +++ b/pkg/sfu/buffer/rtpstats_sender.go @@ -282,7 +282,19 @@ func (r *RTPStatsSender) Update( return } if -gapSN >= cNumSequenceNumbers { - r.logger.Warnw("large sequence number gap negative", nil, "prev", r.extHighestSN, "curr", extSequenceNumber, "gap", gapSN) + r.logger.Warnw( + "large sequence number gap negative", nil, + "prev", r.extHighestSN, + "curr", extSequenceNumber, + "gap", gapSN, + "packetTime", packetTime.String(), + "sequenceNumber", extSequenceNumber, + "timestamp", extTimestamp, + "marker", marker, + "hdrSize", hdrSize, + "payloadSize", payloadSize, + "paddingSize", paddingSize, + ) } if extSequenceNumber < r.extStartSN { @@ -341,7 +353,19 @@ func (r *RTPStatsSender) Update( } } else { // in-order if gapSN >= cNumSequenceNumbers { - r.logger.Warnw("large sequence number gap", nil, "prev", r.extHighestSN, "curr", extSequenceNumber, "gap", gapSN) + r.logger.Warnw( + "large sequence number gap", nil, + "prev", r.extHighestSN, + "curr", extSequenceNumber, + "gap", gapSN, + "packetTime", packetTime.String(), + "sequenceNumber", extSequenceNumber, + "timestamp", extTimestamp, + "marker", marker, + "hdrSize", hdrSize, + "payloadSize", payloadSize, + "paddingSize", paddingSize, + ) } // update gap histogram @@ -438,6 +462,7 @@ func (r *RTPStatsSender) UpdateFromReceiverReport(rr rtcp.ReceptionReport) (rtt } } + // This is 24-bit max in the protocol. So, technically doesn't need extended type. But, done for consistency. packetsLostFromRR := r.packetsLostFromRR&0xFFFF_FFFF_0000_0000 + uint64(rr.TotalLost) if (rr.TotalLost-r.lastRR.TotalLost) < (1<<31) && rr.TotalLost < r.lastRR.TotalLost { packetsLostFromRR += (1 << 32) @@ -512,11 +537,11 @@ func (r *RTPStatsSender) LastReceiverReportTime() time.Time { return r.lastRRTime } -func (r *RTPStatsSender) MaybeAdjustFirstPacketTime(ets uint64) { +func (r *RTPStatsSender) MaybeAdjustFirstPacketTime(ts uint32) { r.lock.Lock() defer r.lock.Unlock() - r.maybeAdjustFirstPacketTime(ets, r.extStartTS) + r.maybeAdjustFirstPacketTime(ts, uint32(r.extStartTS)) } func (r *RTPStatsSender) GetExpectedRTPTimestamp(at time.Time) (expectedTSExt uint64, err error) { diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index 929386f02..b26080526 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -1883,7 +1883,7 @@ func (d *DownTrack) sendSilentFrameOnMuteForOpus() { func (d *DownTrack) HandleRTCPSenderReportData(_payloadType webrtc.PayloadType, layer int32, srData *buffer.RTCPSenderReportData) error { if layer == d.forwarder.GetReferenceLayerSpatial() && srData != nil { - d.rtpStats.MaybeAdjustFirstPacketTime(srData.RTPTimestampExt + d.forwarder.GetReferenceTimestampOffset()) + d.rtpStats.MaybeAdjustFirstPacketTime(srData.RTPTimestamp + uint32(d.forwarder.GetReferenceTimestampOffset())) } return nil } diff --git a/pkg/sfu/forwarder.go b/pkg/sfu/forwarder.go index 503d13a7f..cebb5704b 100644 --- a/pkg/sfu/forwarder.go +++ b/pkg/sfu/forwarder.go @@ -42,6 +42,7 @@ const ( TransitionCostSpatial = 10 ResumeBehindThresholdSeconds = float64(0.2) // 200ms + ResumeBehindHighTresholdSeconds = float64(2.0) // 2 seconds LayerSwitchBehindThresholdSeconds = float64(0.05) // 50ms SwitchAheadThresholdSeconds = float64(0.025) // 25ms ) @@ -187,7 +188,7 @@ type Forwarder struct { codec webrtc.RTPCodecCapability kind webrtc.RTPCodecType logger logger.Logger - getReferenceLayerRTPTimestamp func(ets uint64, layer int32, referenceLayer int32) (uint64, error) + getReferenceLayerRTPTimestamp func(ts uint32, layer int32, referenceLayer int32) (uint32, error) getExpectedRTPTimestamp func(at time.Time) (uint64, error) muted bool @@ -215,7 +216,7 @@ type Forwarder struct { func NewForwarder( kind webrtc.RTPCodecType, logger logger.Logger, - getReferenceLayerRTPTimestamp func(ets uint64, layer int32, referenceLayer int32) (uint64, error), + getReferenceLayerRTPTimestamp func(ts uint32, layer int32, referenceLayer int32) (uint32, error), getExpectedRTPTimestamp func(at time.Time) (uint64, error), ) *Forwarder { f := &Forwarder{ @@ -1499,11 +1500,11 @@ func (f *Forwarder) processSourceSwitch(extPkt *buffer.ExtPacket, layer int32) e // But, cases like muting/unmuting, clock vagaries, pacing, etc. make them not satisfy those conditions always. rtpMungerState := f.rtpMunger.GetLast() extLastTS := rtpMungerState.ExtLastTS - extRefTS := extLastTS extExpectedTS := extLastTS + extRefTS := extExpectedTS switchingAt := time.Now() if f.getReferenceLayerRTPTimestamp != nil { - ets, err := f.getReferenceLayerRTPTimestamp(extPkt.ExtTimestamp, layer, f.referenceLayerSpatial) + ts, err := f.getReferenceLayerRTPTimestamp(extPkt.Packet.Timestamp, layer, f.referenceLayerSpatial) if err != nil { // error out if extRefTS is not available. It can happen when there is no sender report // for the layer being switched to. Can especially happen at the start of the track when layer switches are @@ -1513,7 +1514,15 @@ func (f *Forwarder) processSourceSwitch(extPkt *buffer.ExtPacket, layer int32) e return err } - extRefTS = ets + extRefTS = (extRefTS & 0xFFFF_FFFF_0000_0000) + uint64(ts) + + expectedTS32 := uint32(extExpectedTS) + if (ts-expectedTS32) < 1<<31 && ts < expectedTS32 { + extRefTS += (1 << 32) + } + if (expectedTS32-ts) < 1<<31 && expectedTS32 < ts && extRefTS >= 1<<32 { + extRefTS -= (1 << 32) + } } if f.getExpectedRTPTimestamp != nil { @@ -1569,6 +1578,17 @@ func (f *Forwarder) processSourceSwitch(extPkt *buffer.ExtPacket, layer int32) e if f.resumeBehindThreshold > 0 && diffSeconds > f.resumeBehindThreshold { logTransition("resume, reference too far behind", extExpectedTS, extRefTS, extLastTS, diffSeconds) extNextTS = extExpectedTS + } else if diffSeconds > ResumeBehindHighTresholdSeconds { + // could be due to incorrect reference calculation + f.logger.Infow( + "resume, reference very far behind", + "layer", layer, + "extExpectedTS", extExpectedTS, + "extRefTS", extRefTS, + "extLastTS", extLastTS, + "diffSeconds", diffSeconds, + ) + extNextTS = extExpectedTS } else { extNextTS = extRefTS } diff --git a/pkg/sfu/receiver.go b/pkg/sfu/receiver.go index 434222078..e7ba1fb93 100644 --- a/pkg/sfu/receiver.go +++ b/pkg/sfu/receiver.go @@ -81,7 +81,7 @@ type TrackReceiver interface { GetTemporalLayerFpsForSpatial(layer int32) []float32 GetCalculatedClockRate(layer int32) uint32 - GetReferenceLayerRTPTimestamp(ets uint64, layer int32, referenceLayer int32) (uint64, error) + GetReferenceLayerRTPTimestamp(ts uint32, layer int32, referenceLayer int32) (uint32, error) } // WebRTCReceiver receives a media track @@ -777,8 +777,8 @@ func (w *WebRTCReceiver) GetCalculatedClockRate(layer int32) uint32 { return w.streamTrackerManager.GetCalculatedClockRate(layer) } -func (w *WebRTCReceiver) GetReferenceLayerRTPTimestamp(ets uint64, layer int32, referenceLayer int32) (uint64, error) { - return w.streamTrackerManager.GetReferenceLayerRTPTimestamp(ets, layer, referenceLayer) +func (w *WebRTCReceiver) GetReferenceLayerRTPTimestamp(ts uint32, layer int32, referenceLayer int32) (uint32, error) { + return w.streamTrackerManager.GetReferenceLayerRTPTimestamp(ts, layer, referenceLayer) } // closes all track senders in parallel, returns when all are closed diff --git a/pkg/sfu/rtpmunger.go b/pkg/sfu/rtpmunger.go index b162158fe..8cfb4be27 100644 --- a/pkg/sfu/rtpmunger.go +++ b/pkg/sfu/rtpmunger.go @@ -202,9 +202,25 @@ func (r *RTPMunger) UpdateAndGetSnTs(extPkt *buffer.ExtPacket) (*TranslationPara }, ErrOutOfOrderSequenceNumberCacheMiss } + extSequenceNumber := extPkt.ExtSequenceNumber - snOffset + if extSequenceNumber >= r.extLastSN { + // should not happen, just being paranoid + r.logger.Errorw( + "unexpected packet ordering", nil, + "extIncomingSN", extPkt.ExtSequenceNumber, + "extHighestIncominSN", r.extHighestIncomingSN, + "extLastSN", r.extLastSN, + "snOffsetIncoming", snOffset, + "snOffsetHighest", r.snOffset, + ) + return &TranslationParamsRTP{ + snOrdering: SequenceNumberOrderingOutOfOrder, + }, ErrOutOfOrderSequenceNumberCacheMiss + } + return &TranslationParamsRTP{ snOrdering: SequenceNumberOrderingOutOfOrder, - extSequenceNumber: extPkt.ExtSequenceNumber - snOffset, + extSequenceNumber: extSequenceNumber, extTimestamp: extPkt.ExtTimestamp - r.tsOffset, }, nil } diff --git a/pkg/sfu/streamtrackermanager.go b/pkg/sfu/streamtrackermanager.go index 19fd266e9..9aa82f91e 100644 --- a/pkg/sfu/streamtrackermanager.go +++ b/pkg/sfu/streamtrackermanager.go @@ -76,7 +76,7 @@ type StreamTrackerManager struct { senderReportMu sync.RWMutex senderReports [buffer.DefaultMaxLayerSpatial + 1]endsSenderReport - layerOffsets [buffer.DefaultMaxLayerSpatial + 1][buffer.DefaultMaxLayerSpatial + 1]uint64 + layerOffsets [buffer.DefaultMaxLayerSpatial + 1][buffer.DefaultMaxLayerSpatial + 1]uint32 closed core.Fuse @@ -563,10 +563,10 @@ func (s *StreamTrackerManager) updateLayerOffsetLocked(ref, other int32) { rtpDiff := ntpDiff.Nanoseconds() * int64(s.clockRate) / 1e9 // calculate other layer's time stamp at the same time as ref layer's NTP time - normalizedOtherTS := srOther.RTPTimestampExt + uint64(rtpDiff) + normalizedOtherTS := srOther.RTPTimestamp + uint32(rtpDiff) // now both layers' time stamp refer to the same NTP time and the diff is the offset between the layers - offset := srRef.RTPTimestampExt - normalizedOtherTS + offset := srRef.RTPTimestamp - normalizedOtherTS // use minimal offset to indicate value availability in the extremely unlikely case of // both layers using the same timestamp @@ -643,7 +643,7 @@ func (s *StreamTrackerManager) GetCalculatedClockRate(layer int32) uint32 { return uint32(float64(rdsf) / tsf.Seconds()) } -func (s *StreamTrackerManager) GetReferenceLayerRTPTimestamp(ets uint64, layer int32, referenceLayer int32) (uint64, error) { +func (s *StreamTrackerManager) GetReferenceLayerRTPTimestamp(ts uint32, layer int32, referenceLayer int32) (uint32, error) { s.senderReportMu.RLock() defer s.senderReportMu.RUnlock() @@ -655,7 +655,7 @@ func (s *StreamTrackerManager) GetReferenceLayerRTPTimestamp(ets uint64, layer i return 0, fmt.Errorf("offset unavailable, target: %d, reference: %d", layer, referenceLayer) } - return ets + s.layerOffsets[referenceLayer][layer], nil + return ts + s.layerOffsets[referenceLayer][layer], nil } func (s *StreamTrackerManager) GetMaxTemporalLayerSeen() int32 {