diff --git a/pkg/rtc/wrappedreceiver.go b/pkg/rtc/wrappedreceiver.go index 854d4f792..19cb3a694 100644 --- a/pkg/rtc/wrappedreceiver.go +++ b/pkg/rtc/wrappedreceiver.go @@ -18,6 +18,7 @@ import ( "errors" "strings" "sync" + "time" "github.com/pion/webrtc/v3" "go.uber.org/atomic" @@ -139,6 +140,8 @@ type DummyReceiver struct { maxExpectedLayer int32 pausedValid bool paused bool + + baseTime time.Time } func NewDummyReceiver(trackID livekit.TrackID, streamId string, codec webrtc.RTPCodecParameters, headerExtensions []webrtc.RTPHeaderExtensionParameter) *DummyReceiver { @@ -148,6 +151,7 @@ func NewDummyReceiver(trackID livekit.TrackID, streamId string, codec webrtc.RTP codec: codec, headerExtensions: headerExtensions, downTracks: make(map[livekit.ParticipantID]sfu.TrackSender), + baseTime: time.Now(), } } @@ -324,3 +328,10 @@ func (d *DummyReceiver) GetTrackStats() *livekit.RTPStats { } return nil } + +func (d *DummyReceiver) GetMonotonicNowUnixNano() int64 { + if r, ok := d.receiver.Load().(sfu.TrackReceiver); ok { + return r.GetMonotonicNowUnixNano() + } + return d.baseTime.Add(time.Since(d.baseTime)).UnixNano() +} diff --git a/pkg/sfu/buffer/buffer.go b/pkg/sfu/buffer/buffer.go index ccaf9cd4c..706478ba4 100644 --- a/pkg/sfu/buffer/buffer.go +++ b/pkg/sfu/buffer/buffer.go @@ -91,6 +91,7 @@ type Buffer struct { closed atomic.Bool mime string + baseTime time.Time snRangeMap *utils.RangeMap[uint64, uint64] latestTSForAudioLevelInitialized bool @@ -148,6 +149,7 @@ func NewBuffer(ssrc uint32, maxVideoPkts, maxAudioPkts int) *Buffer { mediaSSRC: ssrc, maxVideoPkts: maxVideoPkts, maxAudioPkts: maxAudioPkts, + baseTime: time.Now(), snRangeMap: utils.NewRangeMap[uint64, uint64](100), pliThrottle: int64(500 * time.Millisecond), logger: l.WithComponent(sutils.ComponentPub).WithComponent(sutils.ComponentSFU), @@ -167,6 +169,13 @@ func (b *Buffer) SetLogger(logger logger.Logger) { } } +func (b *Buffer) SetBaseTime(baseTime time.Time) { + b.Lock() + defer b.Unlock() + + b.baseTime = baseTime +} + func (b *Buffer) SetPaused(paused bool) { b.Lock() defer b.Unlock() @@ -212,7 +221,7 @@ func (b *Buffer) Bind(params webrtc.RTPParameters, codec webrtc.RTPCodecCapabili b.ppsSnapshotId = b.rtpStats.NewSnapshotId() b.clockRate = codec.ClockRate - b.lastReport = time.Now().UnixNano() + b.lastReport = b.getMonotonicNowUnixNano() b.mime = strings.ToLower(codec.MimeType) for _, codecParameter := range params.Codecs { if strings.EqualFold(codecParameter.MimeType, codec.MimeType) { @@ -321,7 +330,7 @@ func (b *Buffer) Write(pkt []byte) (n int, err error) { return } - now := time.Now().UnixNano() + now := b.getMonotonicNowUnixNano() if b.twcc != nil && b.twccExtID != 0 && !b.closed.Load() { if ext := rtpPacket.GetExtension(b.twccExtID); ext != nil { b.twcc.Push(rtpPacket.SSRC, binary.BigEndian.Uint16(ext[0:2]), now, rtpPacket.Marker) @@ -929,7 +938,7 @@ func (b *Buffer) SetSenderReportData(rtpTime uint32, ntpTime uint64, packets uin srData := &RTCPSenderReportData{ RTPTimestamp: rtpTime, NTPTimestamp: mediatransportutil.NtpTime(ntpTime), - At: time.Now(), + At: b.getMonotonicNow(), Packets: packets, Octets: octets, } @@ -1093,7 +1102,7 @@ func (b *Buffer) GetAudioLevel() (float64, bool) { return 0, false } - return b.audioLevel.GetLevel(time.Now().UnixNano()) + return b.audioLevel.GetLevel(b.getMonotonicNowUnixNano()) } func (b *Buffer) OnFpsChanged(f func()) { @@ -1113,6 +1122,16 @@ func (b *Buffer) GetTemporalLayerFpsForSpatial(layer int32) []float32 { return nil } +func (b *Buffer) getMonotonicNowUnixNano() int64 { + return b.baseTime.Add(time.Since(b.baseTime)).UnixNano() +} + +func (b *Buffer) getMonotonicNow() time.Time { + return b.baseTime.Add(time.Since(b.baseTime)) +} + +// --------------------------------------------------------------- + // SVC-TODO: Have to use more conditions to differentiate between // SVC-TODO: SVC and non-SVC (could be single layer or simulcast). // SVC-TODO: May only need to differentiate between simulcast and non-simulcast diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index b94cd308d..9b393818c 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -948,7 +948,7 @@ func (d *DownTrack) WritePaddingRTP(bytesToSend int, paddingOnMute bool, forceMa &hdr, len(payload), &sendPacketMetadata{ - packetTime: time.Now().UnixNano(), + packetTime: d.params.Receiver.GetMonotonicNowUnixNano(), extSequenceNumber: snts[i].extSequenceNumber, extTimestamp: snts[i].extTimestamp, isPadding: true, @@ -1491,7 +1491,7 @@ func (d *DownTrack) writeBlankFrameRTP(duration float32, generation uint32) chan } d.sendingPacket(&hdr, len(payload), &sendPacketMetadata{ - packetTime: time.Now().UnixNano(), + packetTime: d.params.Receiver.GetMonotonicNowUnixNano(), extSequenceNumber: snts[i].extSequenceNumber, extTimestamp: snts[i].extTimestamp, }) @@ -1854,7 +1854,7 @@ func (d *DownTrack) retransmitPackets(nacks []uint16) { len(payload), &sendPacketMetadata{ layer: int32(epm.layer), - packetTime: time.Now().UnixNano(), + packetTime: d.params.Receiver.GetMonotonicNowUnixNano(), extSequenceNumber: epm.extSequenceNumber, extTimestamp: epm.extTimestamp, isRTX: true, @@ -2076,7 +2076,7 @@ func (d *DownTrack) sendSilentFrameOnMuteForOpus() { &hdr, len(payload), &sendPacketMetadata{ - packetTime: time.Now().UnixNano(), + packetTime: d.params.Receiver.GetMonotonicNowUnixNano(), extSequenceNumber: snts[i].extSequenceNumber, extTimestamp: snts[i].extTimestamp, // although this is using empty frames, mark as padding as these are used to trigger Pion OnTrack only diff --git a/pkg/sfu/receiver.go b/pkg/sfu/receiver.go index beb2fb06b..e87471d99 100644 --- a/pkg/sfu/receiver.go +++ b/pkg/sfu/receiver.go @@ -83,6 +83,8 @@ type TrackReceiver interface { GetTemporalLayerFpsForSpatial(layer int32) []float32 GetTrackStats() *livekit.RTPStats + + GetMonotonicNowUnixNano() int64 } // WebRTCReceiver receives a media track @@ -130,6 +132,8 @@ type WebRTCReceiver struct { redPktWriter func(pkt *buffer.ExtPacket, spatialLayer int32) int forwardStats *ForwardStats + + baseTime time.Time } type ReceiverOpts func(w *WebRTCReceiver) *WebRTCReceiver @@ -204,6 +208,7 @@ func NewWebRTCReceiver( onRTCP: onRTCP, isSVC: buffer.IsSvcCodec(track.Codec().MimeType), isRED: buffer.IsRedCodec(track.Codec().MimeType), + baseTime: time.Now(), } for _, opt := range opts { @@ -335,6 +340,7 @@ func (w *WebRTCReceiver) AddUpTrack(track *webrtc.TrackRemote, buff *buffer.Buff layer = buffer.RidToSpatialLayer(track.RID(), w.trackInfo.Load()) } buff.SetLogger(w.logger.WithValues("layer", layer)) + buff.SetBaseTime(w.baseTime) buff.SetAudioLevelParams(audio.AudioLevelParams{ ActiveLevel: w.audioConfig.ActiveLevel, MinPercentile: w.audioConfig.MinPercentile, @@ -831,6 +837,12 @@ func (w *WebRTCReceiver) GetTemporalLayerFpsForSpatial(layer int32) []float32 { return b.GetTemporalLayerFpsForSpatial(layer) } +func (w *WebRTCReceiver) GetMonotonicNowUnixNano() int64 { + return w.baseTime.Add(time.Since(w.baseTime)).UnixNano() +} + +// ----------------------------------------------------------- + // closes all track senders in parallel, returns when all are closed func closeTrackSenders(senders []TrackSender) { wg := sync.WaitGroup{}