mirror of
https://github.com/livekit/livekit.git
synced 2026-07-20 20:01:07 +00:00
Fix time stamp adjustment when starting with dummy packets. (#2053)
* Fix time stamp adjustment when starting with dmummy packets. - Populated extended values in ExtPacket on dummy packet. - Have to pass reference time stamp offset to first packet time adjustment. * display participant version info
This commit is contained in:
@@ -16,6 +16,7 @@ package rtc
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"strconv"
|
||||
@@ -69,12 +70,20 @@ type downTrackState struct {
|
||||
downTrack sfu.DownTrackState
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------
|
||||
|
||||
type participantUpdateInfo struct {
|
||||
version uint32
|
||||
state livekit.ParticipantInfo_State
|
||||
updatedAt time.Time
|
||||
}
|
||||
|
||||
func (p participantUpdateInfo) String() string {
|
||||
return fmt.Sprintf("version: %d, state: %s, updatedAt: %s", p.version, p.state.String(), p.updatedAt.String())
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------
|
||||
|
||||
type ParticipantParams struct {
|
||||
Identity livekit.ParticipantIdentity
|
||||
Name livekit.ParticipantName
|
||||
|
||||
@@ -95,7 +95,7 @@ func (p *ParticipantImpl) SendParticipantUpdate(participantsToUpdate []*livekit.
|
||||
// this is a message delivered out of order, a more recent version of the message had already been
|
||||
// sent.
|
||||
if pi.Version < lastVersion.version {
|
||||
p.params.Logger.Debugw("skipping outdated participant update", "version", pi.Version, "lastVersion", lastVersion)
|
||||
p.params.Logger.Debugw("skipping outdated participant update", "otherParticipant", pi.Identity, "otherPID", pi.Sid, "version", pi.Version, "lastVersion", lastVersion)
|
||||
isValid = false
|
||||
}
|
||||
}
|
||||
|
||||
@@ -863,13 +863,11 @@ func (r *RTPStats) GetRtt() uint32 {
|
||||
return r.rtt
|
||||
}
|
||||
|
||||
func (r *RTPStats) MaybeAdjustFirstPacketTime(srData *RTCPSenderReportData) {
|
||||
func (r *RTPStats) MaybeAdjustFirstPacketTime(ets uint64) {
|
||||
r.lock.Lock()
|
||||
defer r.lock.Unlock()
|
||||
|
||||
if srData != nil {
|
||||
r.maybeAdjustFirstPacketTime(srData.RTPTimestampExt)
|
||||
}
|
||||
r.maybeAdjustFirstPacketTime(ets)
|
||||
}
|
||||
|
||||
func (r *RTPStats) maybeAdjustFirstPacketTime(ets uint64) {
|
||||
|
||||
+10
-3
@@ -1741,7 +1741,6 @@ func (d *DownTrack) sendPaddingOnMute() {
|
||||
// let uptrack have chance to send packet before we send padding
|
||||
time.Sleep(waitBeforeSendPaddingOnMute)
|
||||
|
||||
d.params.Logger.Debugw("sending padding on mute")
|
||||
if d.kind == webrtc.RTPCodecTypeVideo {
|
||||
d.sendPaddingOnMuteForVideo()
|
||||
} else if d.mime == "audio/opus" {
|
||||
@@ -1756,6 +1755,9 @@ func (d *DownTrack) sendPaddingOnMuteForVideo() {
|
||||
if d.rtpStats.IsActive() || d.IsClosed() {
|
||||
return
|
||||
}
|
||||
if i == 0 {
|
||||
d.params.Logger.Debugw("sending padding on mute")
|
||||
}
|
||||
d.WritePaddingRTP(20, true, true)
|
||||
time.Sleep(paddingOnMuteInterval)
|
||||
}
|
||||
@@ -1765,10 +1767,15 @@ func (d *DownTrack) sendSilentFrameOnMuteForOpus() {
|
||||
frameRate := uint32(50)
|
||||
frameDuration := time.Duration(1000/frameRate) * time.Millisecond
|
||||
numFrames := frameRate * uint32(maxPaddingOnMuteDuration/time.Second)
|
||||
first := true
|
||||
for {
|
||||
if d.rtpStats.IsActive() || d.IsClosed() || numFrames <= 0 {
|
||||
return
|
||||
}
|
||||
if first {
|
||||
first = false
|
||||
d.params.Logger.Debugw("sending padding on mute")
|
||||
}
|
||||
snts, _, err := d.forwarder.GetSnTsForBlankFrames(frameRate, 1)
|
||||
if err != nil {
|
||||
d.params.Logger.Warnw("could not get SN/TS for blank frame", err)
|
||||
@@ -1812,8 +1819,8 @@ func (d *DownTrack) sendSilentFrameOnMuteForOpus() {
|
||||
}
|
||||
|
||||
func (d *DownTrack) HandleRTCPSenderReportData(_payloadType webrtc.PayloadType, layer int32, srData *buffer.RTCPSenderReportData) error {
|
||||
if layer == d.forwarder.GetReferenceLayerSpatial() {
|
||||
d.rtpStats.MaybeAdjustFirstPacketTime(srData)
|
||||
if layer == d.forwarder.GetReferenceLayerSpatial() && srData != nil {
|
||||
d.rtpStats.MaybeAdjustFirstPacketTime(srData.RTPTimestampExt + d.forwarder.GetReferenceTimestampOffset())
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
+29
-17
@@ -532,6 +532,13 @@ func (f *Forwarder) GetReferenceLayerSpatial() int32 {
|
||||
return f.referenceLayerSpatial
|
||||
}
|
||||
|
||||
func (f *Forwarder) GetReferenceTimestampOffset() uint64 {
|
||||
f.lock.RLock()
|
||||
defer f.lock.RUnlock()
|
||||
|
||||
return f.refTSOffset
|
||||
}
|
||||
|
||||
func (f *Forwarder) isDeficientLocked() bool {
|
||||
return f.lastAllocation.IsDeficient
|
||||
}
|
||||
@@ -1516,22 +1523,23 @@ func (f *Forwarder) processSourceSwitch(extPkt *buffer.ExtPacket, layer int32) e
|
||||
if err == nil {
|
||||
extExpectedTS = tsExt
|
||||
} else {
|
||||
rtpDiff := uint64(0)
|
||||
if !f.preStartTime.IsZero() && f.refTSOffset == 0 {
|
||||
if !f.preStartTime.IsZero() {
|
||||
timeSinceFirst := time.Since(f.preStartTime)
|
||||
rtpDiff = uint64(timeSinceFirst.Nanoseconds() * int64(f.codec.ClockRate) / 1e9)
|
||||
f.refTSOffset = f.extFirstTS + rtpDiff - extRefTS
|
||||
f.logger.Infow(
|
||||
"calculating refTSOffset",
|
||||
"preStartTime", f.preStartTime.String(),
|
||||
"extFirstTS", f.extFirstTS,
|
||||
"timeSinceFirst", timeSinceFirst,
|
||||
"rtpDiff", rtpDiff,
|
||||
"extRefTS", extRefTS,
|
||||
"refTSOffset", f.refTSOffset,
|
||||
)
|
||||
rtpDiff := uint64(timeSinceFirst.Nanoseconds() * int64(f.codec.ClockRate) / 1e9)
|
||||
extExpectedTS = f.extFirstTS + rtpDiff
|
||||
if f.refTSOffset == 0 {
|
||||
f.refTSOffset = extExpectedTS - extRefTS
|
||||
f.logger.Infow(
|
||||
"calculating refTSOffset",
|
||||
"preStartTime", f.preStartTime.String(),
|
||||
"extFirstTS", f.extFirstTS,
|
||||
"timeSinceFirst", timeSinceFirst,
|
||||
"rtpDiff", rtpDiff,
|
||||
"extRefTS", extRefTS,
|
||||
"refTSOffset", f.refTSOffset,
|
||||
)
|
||||
}
|
||||
}
|
||||
extExpectedTS += rtpDiff
|
||||
}
|
||||
}
|
||||
extRefTS += f.refTSOffset
|
||||
@@ -1746,17 +1754,21 @@ func (f *Forwarder) maybeStart() {
|
||||
f.started = true
|
||||
f.preStartTime = time.Now()
|
||||
|
||||
sequenceNumber := uint16(rand.Intn(1<<14)) + uint16(1<<15) // a random number in third quartile of sequence number space
|
||||
timestamp := uint32(rand.Intn(1<<30)) + uint32(1<<31) // a random number in third quartile of timestamp space
|
||||
extPkt := &buffer.ExtPacket{
|
||||
Packet: &rtp.Packet{
|
||||
Header: rtp.Header{
|
||||
SequenceNumber: uint16(rand.Intn(1<<14)) + uint16(1<<15), // a random number in third quartile of sequence number space
|
||||
Timestamp: uint32(rand.Intn(1<<30)) + uint32(1<<31), // a random number in third quartile of timestamp space
|
||||
SequenceNumber: sequenceNumber,
|
||||
Timestamp: timestamp,
|
||||
},
|
||||
},
|
||||
ExtSequenceNumber: uint64(sequenceNumber),
|
||||
ExtTimestamp: uint64(timestamp),
|
||||
}
|
||||
f.rtpMunger.SetLastSnTs(extPkt)
|
||||
|
||||
f.extFirstTS = uint64(extPkt.Packet.Timestamp)
|
||||
f.extFirstTS = uint64(timestamp)
|
||||
f.logger.Debugw(
|
||||
"starting with dummy forwarding",
|
||||
"sequenceNumber", extPkt.Packet.SequenceNumber,
|
||||
|
||||
Reference in New Issue
Block a user