diff --git a/pkg/sfu/buffer/rtpstats.go b/pkg/sfu/buffer/rtpstats.go index 9bacf08f7..1f0f2694e 100644 --- a/pkg/sfu/buffer/rtpstats.go +++ b/pkg/sfu/buffer/rtpstats.go @@ -705,7 +705,7 @@ func (r *RTPStats) UpdateFromReceiverReport(rr rtcp.ReceptionReport) (rtt uint32 return } -func (r *RTPStats) LastReceiverReport() time.Time { +func (r *RTPStats) LastReceiverReportTime() time.Time { r.lock.RLock() defer r.lock.RUnlock() diff --git a/pkg/sfu/connectionquality/connectionstats.go b/pkg/sfu/connectionquality/connectionstats.go index 565eccabb..641290fa8 100644 --- a/pkg/sfu/connectionquality/connectionstats.go +++ b/pkg/sfu/connectionquality/connectionstats.go @@ -42,6 +42,7 @@ type ConnectionStatsParams struct { GetDeltaStats func() map[uint32]*buffer.StreamStatsWithLayers GetDeltaStatsOverridden func() map[uint32]*buffer.StreamStatsWithLayers GetLastReceiverReportTime func() time.Time + GetTotalPacketsSent func() uint64 Logger logger.Logger } @@ -54,6 +55,7 @@ type ConnectionStats struct { onStatsUpdate func(cs *ConnectionStats, stat *livekit.AnalyticsStat) lock sync.RWMutex + packetsSent uint64 streamingStartedAt time.Time scorer *qualityScorer @@ -213,13 +215,11 @@ func (cs *ConnectionStats) updateScoreWithAggregate(agg *buffer.RTPDeltaInfo, at } func (cs *ConnectionStats) updateScoreFromReceiverReport(at time.Time) (float32, map[uint32]*buffer.StreamStatsWithLayers) { - if cs.params.GetDeltaStatsOverridden == nil || cs.params.GetLastReceiverReportTime == nil { + if cs.params.GetDeltaStatsOverridden == nil || cs.params.GetLastReceiverReportTime == nil || cs.params.GetTotalPacketsSent == nil { return MinMOS, nil } - cs.lock.RLock() - streamingStartedAt := cs.streamingStartedAt - cs.lock.RUnlock() + streamingStartedAt := cs.updateStreamingStart(at) if streamingStartedAt.IsZero() { // not streaming, just return current score mos, _ := cs.scorer.GetMOSAndQuality() @@ -260,6 +260,11 @@ func (cs *ConnectionStats) updateScoreFromReceiverReport(at time.Time) (float32, } func (cs *ConnectionStats) updateScoreAt(at time.Time) (float32, map[uint32]*buffer.StreamStatsWithLayers) { + if cs.params.GetDeltaStatsOverridden != nil { + // receiver report based quality scoring, use stats from receiver report for scoring + return cs.updateScoreFromReceiverReport(at) + } + if cs.params.GetDeltaStats == nil { return MinMOS, nil } @@ -275,33 +280,25 @@ func (cs *ConnectionStats) updateScoreAt(at time.Time) (float32, map[uint32]*buf deltaInfoList = append(deltaInfoList, s.RTPStats) } agg := buffer.AggregateRTPDeltaInfo(deltaInfoList) - if agg != nil && agg.Packets > 0 { - // not very accurate as streaming could have started part way in the window, but don't need accurate time - cs.maybeSetStreamingStart(agg.StartTime) - } else { - cs.clearStreamingStart() - } - - if cs.params.GetDeltaStatsOverridden != nil { - // receiver report based quality scoring, use stats from receiver report for scoring - return cs.updateScoreFromReceiverReport(at) - } - return cs.updateScoreWithAggregate(agg, at), streams } -func (cs *ConnectionStats) maybeSetStreamingStart(at time.Time) { +func (cs *ConnectionStats) updateStreamingStart(at time.Time) time.Time { cs.lock.Lock() - if cs.streamingStartedAt.IsZero() { - cs.streamingStartedAt = at - } - cs.lock.Unlock() -} + defer cs.lock.Unlock() -func (cs *ConnectionStats) clearStreamingStart() { - cs.lock.Lock() - cs.streamingStartedAt = time.Time{} - cs.lock.Unlock() + packetsSent := cs.params.GetTotalPacketsSent() + if packetsSent > cs.packetsSent { + if cs.streamingStartedAt.IsZero() { + // the start could be anywhere after last update, but using `at` as this is not required to be accurate + cs.streamingStartedAt = at + } + } else { + cs.streamingStartedAt = time.Time{} + } + cs.packetsSent = packetsSent + + return cs.streamingStartedAt } func (cs *ConnectionStats) getStat() { diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index d48b98460..6216fa0d8 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -129,14 +129,13 @@ var ( type DownTrackState struct { RTPStats *buffer.RTPStats - DeltaStatsSnapshotId uint32 DeltaStatsOverriddenSnapshotId uint32 ForwarderState ForwarderState } func (d DownTrackState) String() string { - return fmt.Sprintf("DownTrackState{rtpStats: %s, delta: %d, deltaOverridden: %d, forwarder: %s}", - d.RTPStats.ToString(), d.DeltaStatsSnapshotId, d.DeltaStatsOverriddenSnapshotId, d.ForwarderState.String()) + return fmt.Sprintf("DownTrackState{rtpStats: %s, deltaOverridden: %d, forwarder: %s}", + d.RTPStats.ToString(), d.DeltaStatsOverriddenSnapshotId, d.ForwarderState.String()) } // ------------------------------------------------------------------- @@ -248,7 +247,6 @@ type DownTrack struct { blankFramesGeneration atomic.Uint32 connectionStats *connectionquality.ConnectionStats - deltaStatsSnapshotId uint32 deltaStatsOverriddenSnapshotId uint32 isNACKThrottled atomic.Bool @@ -310,15 +308,14 @@ func NewDownTrack(params DowntrackParams) (*DownTrack, error) { IsReceiverReportDriven: true, Logger: params.Logger, }) - d.deltaStatsSnapshotId = d.rtpStats.NewSnapshotId() d.deltaStatsOverriddenSnapshotId = d.rtpStats.NewSnapshotId() d.connectionStats = connectionquality.NewConnectionStats(connectionquality.ConnectionStatsParams{ MimeType: codecs[0].MimeType, // LK-TODO have to notify on codec change IsFECEnabled: strings.EqualFold(codecs[0].MimeType, webrtc.MimeTypeOpus) && strings.Contains(strings.ToLower(codecs[0].SDPFmtpLine), "fec"), - GetDeltaStats: d.getDeltaStats, GetDeltaStatsOverridden: d.getDeltaStatsOverridden, - GetLastReceiverReportTime: func() time.Time { return d.rtpStats.LastReceiverReport() }, + GetLastReceiverReportTime: func() time.Time { return d.rtpStats.LastReceiverReportTime() }, + GetTotalPacketsSent: func() uint64 { return d.rtpStats.GetTotalPacketsPrimary() }, Logger: params.Logger.WithValues("direction", "down"), }) d.connectionStats.OnStatsUpdate(func(_cs *connectionquality.ConnectionStats, stat *livekit.AnalyticsStat) { @@ -328,7 +325,6 @@ func NewDownTrack(params DowntrackParams) (*DownTrack, error) { }) // set initial playout delay to minimum value - if d.params.PlayoutDelayLimit.GetEnabled() && d.params.PlayoutDelayLimit.GetMin() > 0 { delay := rtpextension.PlayoutDelayFromValue( uint16(d.params.PlayoutDelayLimit.GetMin()), @@ -730,13 +726,11 @@ func (d *DownTrack) WritePaddingRTP(bytesToSend int, paddingOnMute bool, forceMa return 0 } - // LK-TODO-START // Ideally should look at header extensions negotiated for // track and decide if padding can be sent. But, browsers behave // in unexpected ways when using audio for bandwidth estimation and // padding is mainly used to probe for excess available bandwidth. // So, to be safe, limit to video tracks - // LK-TODO-END if d.kind == webrtc.RTPCodecTypeAudio { return 0 } @@ -750,6 +744,12 @@ func (d *DownTrack) WritePaddingRTP(bytesToSend int, paddingOnMute bool, forceMa return 0 } + // Hold sending padding packets till first RTCP-RR is received for this RTP stream. + // That is definitive proof that the remote side knows about this RTP stream. + if d.rtpStats.LastReceiverReportTime().IsZero() && !paddingOnMute { + return 0 + } + // RTP padding maximum is 255 bytes. Break it up. // Use 20 byte as estimate of RTP header size (12 byte header + 8 byte extension) num := (bytesToSend + RTPPaddingMaxPayloadSize + RTPPaddingEstimatedHeaderSize - 1) / (RTPPaddingMaxPayloadSize + RTPPaddingEstimatedHeaderSize) @@ -762,16 +762,8 @@ func (d *DownTrack) WritePaddingRTP(bytesToSend int, paddingOnMute bool, forceMa return 0 } - // LK-TODO Look at load balancing a la sfu.Receiver to spread across available CPUs bytesSent := 0 for i := 0; i < len(snts); i++ { - // LK-TODO-START - // Hold sending padding packets till first RTCP-RR is received for this RTP stream. - // That is definitive proof that the remote side knows about this RTP stream. - // The packet count check at the beginning of this function gates sending padding - // on as yet unstarted streams which is a reasonable check. - // LK-TODO-END - hdr := rtp.Header{ Version: 2, Padding: true, @@ -812,7 +804,7 @@ func (d *DownTrack) WritePaddingRTP(bytesToSend int, paddingOnMute bool, forceMa bytesSent += hdr.MarshalSize() + len(payload) } - // STREAM_ALLOCATOR-TODO: change this to pull this counter from stream allocator so that counter can be update in pacer callback + // STREAM_ALLOCATOR-TODO: change this to pull this counter from stream allocator so that counter can be updated in pacer callback return bytesSent } @@ -979,7 +971,6 @@ func (d *DownTrack) MaxLayer() buffer.VideoLayer { func (d *DownTrack) GetState() DownTrackState { dts := DownTrackState{ RTPStats: d.rtpStats, - DeltaStatsSnapshotId: d.deltaStatsSnapshotId, DeltaStatsOverriddenSnapshotId: d.deltaStatsOverriddenSnapshotId, ForwarderState: d.forwarder.GetState(), } @@ -988,7 +979,6 @@ func (d *DownTrack) GetState() DownTrackState { func (d *DownTrack) SeedState(state DownTrackState) { d.rtpStats.Seed(state.RTPStats) - d.deltaStatsSnapshotId = state.DeltaStatsSnapshotId d.deltaStatsOverriddenSnapshotId = state.DeltaStatsOverriddenSnapshotId d.forwarder.SeedState(state.ForwarderState) } @@ -1698,10 +1688,6 @@ func (d *DownTrack) deltaStats(ds *buffer.RTPDeltaInfo) map[uint32]*buffer.Strea return streamStats } -func (d *DownTrack) getDeltaStats() map[uint32]*buffer.StreamStatsWithLayers { - return d.deltaStats(d.rtpStats.DeltaInfo(d.deltaStatsSnapshotId)) -} - func (d *DownTrack) getDeltaStatsOverridden() map[uint32]*buffer.StreamStatsWithLayers { return d.deltaStats(d.rtpStats.DeltaInfoOverridden(d.deltaStatsOverriddenSnapshotId)) }