diff --git a/pkg/sfu/connectionquality/connectionstats.go b/pkg/sfu/connectionquality/connectionstats.go index cfd91ee86..5b03cfe21 100644 --- a/pkg/sfu/connectionquality/connectionstats.go +++ b/pkg/sfu/connectionquality/connectionstats.go @@ -33,17 +33,25 @@ const ( noReceiverReportTooLongThreshold = 30 * time.Second ) +type ConnectionStatsReceiverProvider interface { + GetDeltaStats() map[uint32]*buffer.StreamStatsWithLayers +} + +type ConnectionStatsSenderProvider interface { + GetDeltaStatsSender() map[uint32]*buffer.StreamStatsWithLayers + GetLastReceiverReportTime() time.Time + GetTotalPacketsSent() uint64 +} + type ConnectionStatsParams struct { - UpdateInterval time.Duration - MimeType string - IsFECEnabled bool - IncludeRTT bool - IncludeJitter bool - GetDeltaStats func() map[uint32]*buffer.StreamStatsWithLayers - GetDeltaStatsSender func() map[uint32]*buffer.StreamStatsWithLayers - GetLastReceiverReportTime func() time.Time - GetTotalPacketsSent func() uint64 - Logger logger.Logger + UpdateInterval time.Duration + MimeType string + IsFECEnabled bool + IncludeRTT bool + IncludeJitter bool + ReceiverProvider ConnectionStatsReceiverProvider + SenderProvider ConnectionStatsSenderProvider + Logger logger.Logger } type ConnectionStats struct { @@ -215,7 +223,7 @@ func (cs *ConnectionStats) updateScoreWithAggregate(agg *buffer.RTPDeltaInfo, at } func (cs *ConnectionStats) updateScoreFromReceiverReport(at time.Time) (float32, map[uint32]*buffer.StreamStatsWithLayers) { - if cs.params.GetDeltaStatsSender == nil || cs.params.GetLastReceiverReportTime == nil || cs.params.GetTotalPacketsSent == nil { + if cs.params.SenderProvider == nil { return MinMOS, nil } @@ -226,10 +234,10 @@ func (cs *ConnectionStats) updateScoreFromReceiverReport(at time.Time) (float32, return mos, nil } - streams := cs.params.GetDeltaStatsSender() + streams := cs.params.SenderProvider.GetDeltaStatsSender() if len(streams) == 0 { // check for receiver report not received for a while - marker := cs.params.GetLastReceiverReportTime() + marker := cs.params.SenderProvider.GetLastReceiverReportTime() if marker.IsZero() || streamingStartedAt.After(marker) { marker = streamingStartedAt } @@ -246,7 +254,7 @@ func (cs *ConnectionStats) updateScoreFromReceiverReport(at time.Time) (float32, // delta stat duration could be large due to not receiving receiver report for a long time (for example, due to mute), // adjust to streaming start if necessary agg := toAggregateDeltaInfo(streams) - if streamingStartedAt.After(cs.params.GetLastReceiverReportTime()) { + if streamingStartedAt.After(cs.params.SenderProvider.GetLastReceiverReportTime()) { // last receiver report was before streaming started, wait for next one mos, _ := cs.scorer.GetMOSAndQuality() return mos, streams @@ -260,16 +268,16 @@ func (cs *ConnectionStats) updateScoreFromReceiverReport(at time.Time) (float32, } func (cs *ConnectionStats) updateScoreAt(at time.Time) (float32, map[uint32]*buffer.StreamStatsWithLayers) { - if cs.params.GetDeltaStatsSender != nil { + if cs.params.SenderProvider != nil { // receiver report based quality scoring, use stats from receiver report for scoring return cs.updateScoreFromReceiverReport(at) } - if cs.params.GetDeltaStats == nil { + if cs.params.ReceiverProvider == nil { return MinMOS, nil } - streams := cs.params.GetDeltaStats() + streams := cs.params.ReceiverProvider.GetDeltaStats() if len(streams) == 0 { mos, _ := cs.scorer.GetMOSAndQuality() return mos, nil @@ -287,7 +295,7 @@ func (cs *ConnectionStats) updateStreamingStart(at time.Time) time.Time { cs.lock.Lock() defer cs.lock.Unlock() - packetsSent := cs.params.GetTotalPacketsSent() + packetsSent := cs.params.SenderProvider.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 diff --git a/pkg/sfu/connectionquality/connectionstats_test.go b/pkg/sfu/connectionquality/connectionstats_test.go index a8aec4b81..73c9c649b 100644 --- a/pkg/sfu/connectionquality/connectionstats_test.go +++ b/pkg/sfu/connectionquality/connectionstats_test.go @@ -31,25 +31,42 @@ func newConnectionStats( isFECEnabled bool, includeRTT bool, includeJitter bool, - getDeltaStats func() map[uint32]*buffer.StreamStatsWithLayers, + receiverProvider ConnectionStatsReceiverProvider, ) *ConnectionStats { return NewConnectionStats(ConnectionStatsParams{ - MimeType: mimeType, - IsFECEnabled: isFECEnabled, - IncludeRTT: includeRTT, - IncludeJitter: includeJitter, - GetDeltaStats: getDeltaStats, - Logger: logger.GetLogger(), + MimeType: mimeType, + IsFECEnabled: isFECEnabled, + IncludeRTT: includeRTT, + IncludeJitter: includeJitter, + ReceiverProvider: receiverProvider, + Logger: logger.GetLogger(), }) } +// ----------------------------------------------- + +type testReceiverProvider struct { + streams map[uint32]*buffer.StreamStatsWithLayers +} + +func newTestReceiverProvider() *testReceiverProvider { + return &testReceiverProvider{} +} + +func (trp *testReceiverProvider) setStreams(streams map[uint32]*buffer.StreamStatsWithLayers) { + trp.streams = streams +} + +func (trp *testReceiverProvider) GetDeltaStats() map[uint32]*buffer.StreamStatsWithLayers { + return trp.streams +} + +// ----------------------------------------------- + func TestConnectionQuality(t *testing.T) { + trp := newTestReceiverProvider() t.Run("quality scorer operation", func(t *testing.T) { - var streams map[uint32]*buffer.StreamStatsWithLayers - getDeltaStats := func() map[uint32]*buffer.StreamStatsWithLayers { - return streams - } - cs := newConnectionStats("audio/opus", false, true, true, getDeltaStats) + cs := newConnectionStats("audio/opus", false, true, true, trp) duration := 5 * time.Second now := time.Now() @@ -63,7 +80,7 @@ func TestConnectionQuality(t *testing.T) { require.Equal(t, livekit.ConnectionQuality_EXCELLENT, quality) // best conditions (no loss, jitter/rtt = 0) - quality should stay EXCELLENT - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -71,7 +88,7 @@ func TestConnectionQuality(t *testing.T) { Packets: 250, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.6), mos) @@ -79,7 +96,7 @@ func TestConnectionQuality(t *testing.T) { // introduce loss and the score should drop - 12% loss for Opus -> POOR now = now.Add(duration) - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -96,7 +113,7 @@ func TestConnectionQuality(t *testing.T) { PacketsLost: 0, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(2.1), mos) @@ -106,7 +123,7 @@ func TestConnectionQuality(t *testing.T) { // although significant loss (12%) in the previous window, lowest score is // bound so that climbing back does not take too long even under excellent conditions. now = now.Add(duration) - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -114,7 +131,7 @@ func TestConnectionQuality(t *testing.T) { Packets: 250, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.1), mos) @@ -122,7 +139,7 @@ func TestConnectionQuality(t *testing.T) { // should stay at GOOD if conditions continue to be good now = now.Add(duration) - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -130,7 +147,7 @@ func TestConnectionQuality(t *testing.T) { Packets: 250, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.1), mos) @@ -138,7 +155,7 @@ func TestConnectionQuality(t *testing.T) { // should climb up to EXCELLENT if conditions continue to be good now = now.Add(duration) - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -146,7 +163,7 @@ func TestConnectionQuality(t *testing.T) { Packets: 250, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.6), mos) @@ -154,7 +171,7 @@ func TestConnectionQuality(t *testing.T) { // introduce loss and the score should drop - 5% loss for Opus -> GOOD now = now.Add(duration) - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -163,7 +180,7 @@ func TestConnectionQuality(t *testing.T) { PacketsLost: 13, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.1), mos) @@ -171,7 +188,7 @@ func TestConnectionQuality(t *testing.T) { // should stay at GOOD quality for another iteration even if the conditions improve now = now.Add(duration) - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -179,7 +196,7 @@ func TestConnectionQuality(t *testing.T) { Packets: 250, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.1), mos) @@ -187,7 +204,7 @@ func TestConnectionQuality(t *testing.T) { // should climb up to EXCELLENT if conditions continue to be good now = now.Add(duration) - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -195,7 +212,7 @@ func TestConnectionQuality(t *testing.T) { Packets: 250, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.6), mos) @@ -203,7 +220,7 @@ func TestConnectionQuality(t *testing.T) { // mute when quality is POOR should return quality to EXCELLENT now = now.Add(duration) - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -212,7 +229,7 @@ func TestConnectionQuality(t *testing.T) { PacketsLost: 30, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(2.1), mos) @@ -228,7 +245,7 @@ func TestConnectionQuality(t *testing.T) { // that means even if the next update has 0 packets, it should hold state and stay at EXCELLENT quality cs.UpdateMuteAt(false, now.Add(3*time.Second)) - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -236,7 +253,7 @@ func TestConnectionQuality(t *testing.T) { Packets: 0, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.6), mos) @@ -244,7 +261,7 @@ func TestConnectionQuality(t *testing.T) { // next update with no packets should knock quality down now = now.Add(duration) - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -252,7 +269,7 @@ func TestConnectionQuality(t *testing.T) { Packets: 0, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(2.1), mos) @@ -265,7 +282,7 @@ func TestConnectionQuality(t *testing.T) { // with lesser number of packet (simulating DTX). // even higher loss (like 10%) should not knock down quality due to quadratic weighting of packet loss ratio - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -274,7 +291,7 @@ func TestConnectionQuality(t *testing.T) { PacketsLost: 5, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.6), mos) @@ -287,7 +304,7 @@ func TestConnectionQuality(t *testing.T) { // RTT and jitter can knock quality down. // at 2% loss, quality should stay at EXCELLENT purely based on loss, but with added RTT/jitter, should drop to GOOD - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -298,7 +315,7 @@ func TestConnectionQuality(t *testing.T) { JitterMax: 30000, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.1), mos) @@ -313,7 +330,7 @@ func TestConnectionQuality(t *testing.T) { cs.AddBitrateTransitionAt(1_000_000, now) cs.AddBitrateTransitionAt(2_000_000, now.Add(2*time.Second)) - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -322,7 +339,7 @@ func TestConnectionQuality(t *testing.T) { Bytes: 8_000_000 / 8 / 5, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.1), mos) @@ -332,7 +349,7 @@ func TestConnectionQuality(t *testing.T) { cs.AddBitrateTransitionAt(1_000_000, now) cs.AddBitrateTransitionAt(2_000_000, now.Add(2*time.Second)) - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -341,7 +358,7 @@ func TestConnectionQuality(t *testing.T) { Bytes: 8_000_000 / 8 / 5, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.1), mos) @@ -356,7 +373,7 @@ func TestConnectionQuality(t *testing.T) { // unmute layer cs.UpdateLayerMuteAt(false, now.Add(2*time.Second)) - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -365,7 +382,7 @@ func TestConnectionQuality(t *testing.T) { Bytes: 8_000_000 / 8 / 5, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.6), mos) @@ -383,7 +400,7 @@ func TestConnectionQuality(t *testing.T) { // although conditions are perfect, climbing back from POOR (because of pause above) // will only climb to GOOD. - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -392,7 +409,7 @@ func TestConnectionQuality(t *testing.T) { Bytes: 8_000_000 / 8 / 5, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality = cs.GetScoreAndQuality() require.Greater(t, float32(4.1), mos) @@ -400,11 +417,7 @@ func TestConnectionQuality(t *testing.T) { }) t.Run("quality scorer dependent rtt", func(t *testing.T) { - var streams map[uint32]*buffer.StreamStatsWithLayers - getDeltaStats := func() map[uint32]*buffer.StreamStatsWithLayers { - return streams - } - cs := newConnectionStats("audio/opus", false, false, true, getDeltaStats) + cs := newConnectionStats("audio/opus", false, false, true, trp) duration := 5 * time.Second now := time.Now() @@ -414,7 +427,7 @@ func TestConnectionQuality(t *testing.T) { // RTT does not knock quality down because it is dependent and hence not taken into account // at 2% loss, quality should stay at EXCELLENT purely based on loss. With high RTT (700 ms) // quality should drop to GOOD if RTT were taken into consideration - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -424,7 +437,7 @@ func TestConnectionQuality(t *testing.T) { RttMax: 700, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality := cs.GetScoreAndQuality() require.Greater(t, float32(4.6), mos) @@ -432,11 +445,7 @@ func TestConnectionQuality(t *testing.T) { }) t.Run("quality scorer dependent jitter", func(t *testing.T) { - var streams map[uint32]*buffer.StreamStatsWithLayers - getDeltaStats := func() map[uint32]*buffer.StreamStatsWithLayers { - return streams - } - cs := newConnectionStats("audio/opus", false, true, false, getDeltaStats) + cs := newConnectionStats("audio/opus", false, true, false, trp) duration := 5 * time.Second now := time.Now() @@ -446,7 +455,7 @@ func TestConnectionQuality(t *testing.T) { // Jitter does not knock quality down because it is dependent and hence not taken into account // at 2% loss, quality should stay at EXCELLENT purely based on loss. With high jitter (200 ms) // quality should drop to GOOD if jitter were taken into consideration - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 1: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -456,7 +465,7 @@ func TestConnectionQuality(t *testing.T) { JitterMax: 200, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality := cs.GetScoreAndQuality() require.Greater(t, float32(4.6), mos) @@ -601,18 +610,14 @@ func TestConnectionQuality(t *testing.T) { for _, tc := range testCases { t.Run(tc.name, func(t *testing.T) { - var streams map[uint32]*buffer.StreamStatsWithLayers - getDeltaStats := func() map[uint32]*buffer.StreamStatsWithLayers { - return streams - } - cs := newConnectionStats(tc.mimeType, tc.isFECEnabled, true, true, getDeltaStats) + cs := newConnectionStats(tc.mimeType, tc.isFECEnabled, true, true, trp) duration := 5 * time.Second now := time.Now() cs.StartAt(&livekit.TrackInfo{Type: livekit.TrackType_AUDIO}, now.Add(-duration)) for _, eq := range tc.expectedQualities { - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 123: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -621,7 +626,7 @@ func TestConnectionQuality(t *testing.T) { PacketsLost: uint32(math.Ceil(eq.packetLossPercentage * float64(tc.packetsExpected) / 100.0)), }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality := cs.GetScoreAndQuality() require.Greater(t, eq.expectedMOS, mos) @@ -698,11 +703,7 @@ func TestConnectionQuality(t *testing.T) { for _, tc := range testCases { t.Run(tc.name, func(t *testing.T) { - var streams map[uint32]*buffer.StreamStatsWithLayers - getDeltaStats := func() map[uint32]*buffer.StreamStatsWithLayers { - return streams - } - cs := newConnectionStats("video/vp8", false, true, true, getDeltaStats) + cs := newConnectionStats("video/vp8", false, true, true, trp) duration := 5 * time.Second now := time.Now() @@ -712,7 +713,7 @@ func TestConnectionQuality(t *testing.T) { cs.AddBitrateTransitionAt(tr.bitrate, now.Add(tr.offset)) } - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 123: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -721,7 +722,7 @@ func TestConnectionQuality(t *testing.T) { Bytes: tc.bytes, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality := cs.GetScoreAndQuality() require.Greater(t, tc.expectedMOS, mos) @@ -789,11 +790,7 @@ func TestConnectionQuality(t *testing.T) { for _, tc := range testCases { t.Run(tc.name, func(t *testing.T) { - var streams map[uint32]*buffer.StreamStatsWithLayers - getDeltaStats := func() map[uint32]*buffer.StreamStatsWithLayers { - return streams - } - cs := newConnectionStats("video/vp8", false, true, true, getDeltaStats) + cs := newConnectionStats("video/vp8", false, true, true, trp) duration := 5 * time.Second now := time.Now() @@ -803,7 +800,7 @@ func TestConnectionQuality(t *testing.T) { cs.AddLayerTransitionAt(tr.distance, now.Add(tr.offset)) } - streams = map[uint32]*buffer.StreamStatsWithLayers{ + trp.setStreams(map[uint32]*buffer.StreamStatsWithLayers{ 123: { RTPStats: &buffer.RTPDeltaInfo{ StartTime: now, @@ -811,7 +808,7 @@ func TestConnectionQuality(t *testing.T) { Packets: 200, }, }, - } + }) cs.updateScoreAt(now.Add(duration)) mos, quality := cs.GetScoreAndQuality() require.Greater(t, tc.expectedMOS, mos) diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index 430e3e8e0..fb445f222 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -311,12 +311,10 @@ func NewDownTrack(params DowntrackParams) (*DownTrack, error) { d.deltaStatsSenderSnapshotId = d.rtpStats.NewSenderSnapshotId() 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"), - GetDeltaStatsSender: d.getDeltaStatsSender, - GetLastReceiverReportTime: func() time.Time { return d.rtpStats.LastReceiverReportTime() }, - GetTotalPacketsSent: func() uint64 { return d.rtpStats.GetTotalPacketsPrimary() }, - Logger: params.Logger.WithValues("direction", "down"), + 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"), + SenderProvider: d, + Logger: params.Logger.WithValues("direction", "down"), }) d.connectionStats.OnStatsUpdate(func(_cs *connectionquality.ConnectionStats, stat *livekit.AnalyticsStat) { if onStatsUpdate := d.getOnStatsUpdate(); onStatsUpdate != nil { @@ -1715,9 +1713,16 @@ func (d *DownTrack) deltaStats(ds *buffer.RTPDeltaInfo) map[uint32]*buffer.Strea return streamStats } -func (d *DownTrack) getDeltaStatsSender() map[uint32]*buffer.StreamStatsWithLayers { +func (d *DownTrack) GetDeltaStatsSender() map[uint32]*buffer.StreamStatsWithLayers { return d.deltaStats(d.rtpStats.DeltaInfoSender(d.deltaStatsSenderSnapshotId)) - return nil +} + +func (d *DownTrack) GetLastReceiverReportTime() time.Time { + return d.rtpStats.LastReceiverReportTime() +} + +func (d *DownTrack) GetTotalPacketsSent() uint64 { + return d.rtpStats.GetTotalPacketsPrimary() } func (d *DownTrack) GetNackStats() (totalPackets uint32, totalRepeatedNACKs uint32) { diff --git a/pkg/sfu/receiver.go b/pkg/sfu/receiver.go index a153ce53a..275e653d8 100644 --- a/pkg/sfu/receiver.go +++ b/pkg/sfu/receiver.go @@ -223,10 +223,10 @@ func NewWebRTCReceiver( }) w.connectionStats = connectionquality.NewConnectionStats(connectionquality.ConnectionStatsParams{ - MimeType: w.codec.MimeType, - IsFECEnabled: strings.EqualFold(w.codec.MimeType, webrtc.MimeTypeOpus) && strings.Contains(strings.ToLower(w.codec.SDPFmtpLine), "fec"), - GetDeltaStats: w.getDeltaStats, - Logger: w.logger.WithValues("direction", "up"), + MimeType: w.codec.MimeType, + IsFECEnabled: strings.EqualFold(w.codec.MimeType, webrtc.MimeTypeOpus) && strings.Contains(strings.ToLower(w.codec.SDPFmtpLine), "fec"), + ReceiverProvider: w, + Logger: w.logger.WithValues("direction", "up"), }) w.connectionStats.OnStatsUpdate(func(_cs *connectionquality.ConnectionStats, stat *livekit.AnalyticsStat) { if w.onStatsUpdate != nil { @@ -595,7 +595,7 @@ func (w *WebRTCReceiver) GetAudioLevel() (float64, bool) { return 0, false } -func (w *WebRTCReceiver) getDeltaStats() map[uint32]*buffer.StreamStatsWithLayers { +func (w *WebRTCReceiver) GetDeltaStats() map[uint32]*buffer.StreamStatsWithLayers { w.bufferMu.RLock() defer w.bufferMu.RUnlock()