mirror of
https://github.com/livekit/livekit.git
synced 2026-08-29 03:18:47 +00:00
Adjust first packet time on down track resume. (#2566)
Allows subscriber sender report to line up better quicker.
This commit is contained in:
@@ -26,6 +26,7 @@ import (
|
||||
"github.com/livekit/protocol/logger"
|
||||
|
||||
"github.com/livekit/livekit-server/pkg/sfu"
|
||||
"github.com/livekit/livekit-server/pkg/sfu/buffer"
|
||||
)
|
||||
|
||||
// wrapper around WebRTC receiver, overriding its ID
|
||||
@@ -330,6 +331,13 @@ func (d *DummyReceiver) GetReferenceLayerRTPTimestamp(ts uint32, layer int32, re
|
||||
return 0, errors.New("receiver not available")
|
||||
}
|
||||
|
||||
func (d *DummyReceiver) GetRTCPSenderReportData(layer int32) (*buffer.RTCPSenderReportData, *buffer.RTCPSenderReportData) {
|
||||
if r, ok := d.receiver.Load().(sfu.TrackReceiver); ok {
|
||||
return r.GetRTCPSenderReportData(layer)
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (d *DummyReceiver) GetTrackStats() *livekit.RTPStats {
|
||||
if r, ok := d.receiver.Load().(sfu.TrackReceiver); ok {
|
||||
return r.GetTrackStats()
|
||||
|
||||
@@ -615,11 +615,19 @@ func (r *RTPStatsSender) MaybeAdjustFirstPacketTime(srFirst *RTCPSenderReportDat
|
||||
r.lock.Lock()
|
||||
defer r.lock.Unlock()
|
||||
|
||||
srFirstCopy := *srFirst
|
||||
r.srFeedFirst = &srFirstCopy
|
||||
if !r.initialized {
|
||||
return
|
||||
}
|
||||
|
||||
srNewestCopy := *srNewest
|
||||
r.srFeedNewest = &srNewestCopy
|
||||
if srFirst != nil {
|
||||
srFirstCopy := *srFirst
|
||||
r.srFeedFirst = &srFirstCopy
|
||||
}
|
||||
|
||||
if srNewest != nil {
|
||||
srNewestCopy := *srNewest
|
||||
r.srFeedNewest = &srNewestCopy
|
||||
}
|
||||
|
||||
r.maybeAdjustFirstPacketTime(ts, uint32(r.extStartTS))
|
||||
}
|
||||
|
||||
+17
-6
@@ -1951,16 +1951,24 @@ func (d *DownTrack) HandleRTCPSenderReportData(
|
||||
srFirst *buffer.RTCPSenderReportData,
|
||||
srNewest *buffer.RTCPSenderReportData,
|
||||
) error {
|
||||
if (layer == d.forwarder.GetReferenceLayerSpatial() || (layer == 0 && isSVC)) && srNewest != nil {
|
||||
d.rtpStats.MaybeAdjustFirstPacketTime(
|
||||
srFirst,
|
||||
srNewest,
|
||||
srNewest.RTPTimestamp+uint32(d.forwarder.GetReferenceTimestampOffset()),
|
||||
)
|
||||
if layer == d.forwarder.GetReferenceLayerSpatial() || (layer == 0 && isSVC) {
|
||||
d.handleRTCPSenderReportData(srFirst, srNewest)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (d *DownTrack) handleRTCPSenderReportData(srFirst *buffer.RTCPSenderReportData, srNewest *buffer.RTCPSenderReportData) {
|
||||
if srNewest == nil {
|
||||
return
|
||||
}
|
||||
|
||||
d.rtpStats.MaybeAdjustFirstPacketTime(
|
||||
srFirst,
|
||||
srNewest,
|
||||
srNewest.RTPTimestamp+uint32(d.forwarder.GetReferenceTimestampOffset()),
|
||||
)
|
||||
}
|
||||
|
||||
type sendPacketMetadata struct {
|
||||
layer int32
|
||||
packetTime time.Time
|
||||
@@ -2010,6 +2018,9 @@ func (d *DownTrack) sendingPacket(hdr *rtp.Header, payloadSize int, spmd *sendPa
|
||||
}
|
||||
|
||||
if spmd.tp.isResuming {
|
||||
// adjust first packet time on a resumption so that sender reports can lock in quicker
|
||||
d.handleRTCPSenderReportData(d.params.Receiver.GetRTCPSenderReportData(d.forwarder.GetReferenceLayerSpatial()))
|
||||
|
||||
if sal := d.getStreamAllocatorListener(); sal != nil {
|
||||
sal.OnResume(d)
|
||||
}
|
||||
|
||||
@@ -85,6 +85,7 @@ type TrackReceiver interface {
|
||||
|
||||
GetCalculatedClockRate(layer int32) uint32
|
||||
GetReferenceLayerRTPTimestamp(ts uint32, layer int32, referenceLayer int32) (uint32, error)
|
||||
GetRTCPSenderReportData(layer int32) (*buffer.RTCPSenderReportData, *buffer.RTCPSenderReportData)
|
||||
|
||||
GetTrackStats() *livekit.RTPStats
|
||||
}
|
||||
@@ -790,6 +791,10 @@ func (w *WebRTCReceiver) GetReferenceLayerRTPTimestamp(ts uint32, layer int32, r
|
||||
return w.streamTrackerManager.GetReferenceLayerRTPTimestamp(ts, layer, referenceLayer)
|
||||
}
|
||||
|
||||
func (w *WebRTCReceiver) GetRTCPSenderReportData(layer int32) (*buffer.RTCPSenderReportData, *buffer.RTCPSenderReportData) {
|
||||
return w.streamTrackerManager.GetRTCPSenderReportData(layer)
|
||||
}
|
||||
|
||||
// closes all track senders in parallel, returns when all are closed
|
||||
func closeTrackSenders(senders []TrackSender) {
|
||||
wg := sync.WaitGroup{}
|
||||
|
||||
@@ -626,6 +626,23 @@ func (s *StreamTrackerManager) SetRTCPSenderReportData(layer int32, srFirst *buf
|
||||
}
|
||||
}
|
||||
|
||||
func (s *StreamTrackerManager) GetRTCPSenderReportData(layer int32) (*buffer.RTCPSenderReportData, *buffer.RTCPSenderReportData) {
|
||||
s.senderReportMu.Lock()
|
||||
defer s.senderReportMu.Unlock()
|
||||
|
||||
if layer < 0 || int(layer) >= len(s.senderReports) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// SVC-TODO: better SVC detection
|
||||
if s.isSVC {
|
||||
// there is only one stream in SVC
|
||||
layer = 0
|
||||
}
|
||||
|
||||
return s.senderReports[layer].first, s.senderReports[layer].newest
|
||||
}
|
||||
|
||||
func (s *StreamTrackerManager) GetCalculatedClockRate(layer int32) uint32 {
|
||||
s.senderReportMu.RLock()
|
||||
defer s.senderReportMu.RUnlock()
|
||||
|
||||
Reference in New Issue
Block a user