From 6f7e6c45568a34f116e99895bfed90992f7d7a78 Mon Sep 17 00:00:00 2001 From: shishirng Date: Wed, 9 Feb 2022 20:45:53 -0500 Subject: [PATCH] Compute delta stats to send downstream (#426) * Compute delta stats to send downstream Signed-off-by: shishir gowda * Update tests: total_packets should be diff between 2 packets First packet was 1, second was 4. diff should be 3 Signed-off-by: shishir gowda * If there are no videoLayers, do not sent in Stats For audio and Downstream tracks, we do not get layers Signed-off-by: shishir gowda * Use prev Max layer for current delta and update layer info for next round --- pkg/telemetry/statsworker.go | 244 ++++++++++++++++--- pkg/telemetry/test/telemetry_service_test.go | 2 +- 2 files changed, 207 insertions(+), 39 deletions(-) diff --git a/pkg/telemetry/statsworker.go b/pkg/telemetry/statsworker.go index ac2107bc3..84cbc6020 100644 --- a/pkg/telemetry/statsworker.go +++ b/pkg/telemetry/statsworker.go @@ -3,11 +3,64 @@ package telemetry import ( "context" - "github.com/livekit/protocol/livekit" "google.golang.org/protobuf/proto" + + "github.com/livekit/protocol/livekit" "google.golang.org/protobuf/types/known/timestamppb" ) +type Stat struct { + Score float32 + Rtt uint32 + Jitter uint32 + TotalPrimaryPackets uint32 + TotalPrimaryBytes uint64 + TotalRetransmitPackets uint32 + TotalRetransmitBytes uint64 + TotalPaddingPackets uint32 + TotalPaddingBytes uint64 + TotalPacketsLost uint32 + TotalFrames uint32 + TotalNacks uint32 + TotalPlis uint32 + TotalFirs uint32 + VideoLayers map[int32]*livekit.AnalyticsVideoLayer + TotalBytes uint64 + TotalPackets uint32 + MaxLayer int32 +} + +func (stat *Stat) ToAnalyticsStats(layers *livekit.AnalyticsVideoLayer) *livekit.AnalyticsStat { + stream := &livekit.AnalyticsStream{ + TotalPrimaryPackets: stat.TotalPrimaryPackets, + TotalPrimaryBytes: stat.TotalPrimaryBytes, + TotalRetransmitPackets: stat.TotalRetransmitPackets, + TotalRetransmitBytes: stat.TotalRetransmitBytes, + TotalPaddingPackets: stat.TotalPaddingPackets, + TotalPaddingBytes: stat.TotalPaddingBytes, + TotalPacketsLost: stat.TotalPacketsLost, + TotalFrames: stat.TotalFrames, + Rtt: stat.Rtt, + Jitter: stat.Jitter, + TotalNacks: stat.TotalNacks, + TotalPlis: stat.TotalPlis, + TotalFirs: stat.TotalFirs, + } + if layers != nil { + stream.VideoLayers = []*livekit.AnalyticsVideoLayer{layers} + } + return &livekit.AnalyticsStat{Streams: []*livekit.AnalyticsStream{stream}, Score: stat.Score} +} + +type Stats struct { + // local stats context used for coalesce/delta calculations + curStats *Stat + prevStats *Stat + + // current stats received per stream + queue []*livekit.AnalyticsStat +} + // StatsWorker handles participant stats type StatsWorker struct { ctx context.Context @@ -16,8 +69,8 @@ type StatsWorker struct { roomName livekit.RoomName participantID livekit.ParticipantID - outgoingPerTrack map[livekit.TrackID][]*livekit.AnalyticsStat - incomingPerTrack map[livekit.TrackID][]*livekit.AnalyticsStat + outgoingPerTrack map[livekit.TrackID]Stats + incomingPerTrack map[livekit.TrackID]Stats } func newStatsWorker( @@ -34,18 +87,22 @@ func newStatsWorker( roomName: roomName, participantID: participantID, - outgoingPerTrack: make(map[livekit.TrackID][]*livekit.AnalyticsStat), - incomingPerTrack: make(map[livekit.TrackID][]*livekit.AnalyticsStat), + outgoingPerTrack: make(map[livekit.TrackID]Stats), + incomingPerTrack: make(map[livekit.TrackID]Stats), } return s } func (s *StatsWorker) appendOutgoing(trackID livekit.TrackID, stat *livekit.AnalyticsStat) { - s.outgoingPerTrack[trackID] = append(s.outgoingPerTrack[trackID], stat) + stats := s.outgoingPerTrack[trackID] + stats.queue = append(stats.queue, stat) + s.outgoingPerTrack[trackID] = stats } func (s *StatsWorker) appendIncoming(trackID livekit.TrackID, stat *livekit.AnalyticsStat) { - s.incomingPerTrack[trackID] = append(s.incomingPerTrack[trackID], stat) + stats := s.incomingPerTrack[trackID] + stats.queue = append(stats.queue, stat) + s.incomingPerTrack[trackID] = stats } func (s *StatsWorker) OnTrackStat(trackID livekit.TrackID, direction livekit.StreamType, stat *livekit.AnalyticsStat) { @@ -69,33 +126,43 @@ func (s *StatsWorker) Update() { func (s *StatsWorker) collectDownstreamStats(ts *timestamppb.Timestamp, stats []*livekit.AnalyticsStat) []*livekit.AnalyticsStat { for trackID, analyticsStats := range s.outgoingPerTrack { - analyticsStat := coalesce(analyticsStats) - if analyticsStat == nil { - continue + analyticsStat := s.getDeltaStats(&analyticsStats, ts, trackID, livekit.StreamType_DOWNSTREAM) + if analyticsStat != nil { + stats = append(stats, analyticsStat) } - - s.patch(analyticsStat, ts, trackID, livekit.StreamType_DOWNSTREAM) - stats = append(stats, analyticsStat) + // clear the queue + analyticsStats.queue = nil + s.outgoingPerTrack[trackID] = analyticsStats } - s.outgoingPerTrack = make(map[livekit.TrackID][]*livekit.AnalyticsStat, 0) - return stats } func (s *StatsWorker) collectUpstreamStats(ts *timestamppb.Timestamp, stats []*livekit.AnalyticsStat) []*livekit.AnalyticsStat { for trackID, analyticsStats := range s.incomingPerTrack { - analyticsStat := coalesce(analyticsStats) - if analyticsStat == nil { - continue + analyticsStat := s.getDeltaStats(&analyticsStats, ts, trackID, livekit.StreamType_UPSTREAM) + if analyticsStat != nil { + stats = append(stats, analyticsStat) } - - s.patch(analyticsStat, ts, trackID, livekit.StreamType_UPSTREAM) - stats = append(stats, analyticsStat) + // clear the queue + analyticsStats.queue = nil + s.incomingPerTrack[trackID] = analyticsStats } - s.incomingPerTrack = make(map[livekit.TrackID][]*livekit.AnalyticsStat, 0) - return stats } +func (s *StatsWorker) getDeltaStats(stats *Stats, ts *timestamppb.Timestamp, trackID livekit.TrackID, kind livekit.StreamType) *livekit.AnalyticsStat { + // merge all streams stats of track + stats.coalesce() + // create deltaStats to send + analyticsStat := stats.computeDeltaStats() + // update stats for next interval + if analyticsStat == nil { + return nil + } + stats.update() + s.patch(analyticsStat, ts, trackID, kind) + + return analyticsStat +} func (s *StatsWorker) patch( analyticsStat *livekit.AnalyticsStat, @@ -115,23 +182,30 @@ func (s *StatsWorker) Close() { s.Update() } -func coalesce(stats []*livekit.AnalyticsStat) *livekit.AnalyticsStat { - if len(stats) == 0 { - return nil +// create a single stream and single video layer post aggregation +func (stats *Stats) update() { + stats.prevStats = stats.curStats + stats.curStats = nil +} + +// create a single stream and single video layer post aggregation +func (stats *Stats) coalesce() { + if len(stats.queue) == 0 { + return } // average score of all available stats score := float32(0.0) - for _, stat := range stats { + for _, stat := range stats.queue { score += stat.Score } - score = score / float32(len(stats)) + score = score / float32(len(stats.queue)) // aggregate streams across all stats maxRTT := make(map[uint32]uint32) maxJitter := make(map[uint32]uint32) analyticsStreams := make(map[uint32]*livekit.AnalyticsStream) - for _, stat := range stats { + for _, stat := range stats.queue { // // For each stream (identified by SSRC) consolidate reports. // For cumulative stats, take the latest report. @@ -163,17 +237,111 @@ func coalesce(stats []*livekit.AnalyticsStat) *livekit.AnalyticsStat { } } - streams := make([]*livekit.AnalyticsStream, 0, len(analyticsStreams)) + curStats := Stat{Score: score, VideoLayers: make(map[int32]*livekit.AnalyticsVideoLayer)} + // find aggregates across streams for ssrc, analyticsStream := range analyticsStreams { - stream := proto.Clone(analyticsStream).(*livekit.AnalyticsStream) - stream.Rtt = maxRTT[ssrc] - stream.Jitter = maxJitter[ssrc] + rtt := maxRTT[ssrc] + if rtt > curStats.Rtt { + curStats.Rtt = rtt + } + + jitter := maxJitter[ssrc] + if jitter > curStats.Jitter { + curStats.Jitter = jitter + } + + curStats.TotalFrames += analyticsStream.TotalFrames + curStats.TotalPrimaryPackets += analyticsStream.TotalPrimaryPackets + curStats.TotalPrimaryBytes += analyticsStream.TotalPrimaryBytes + curStats.TotalRetransmitPackets += analyticsStream.TotalRetransmitPackets + curStats.TotalRetransmitBytes += analyticsStream.TotalRetransmitBytes + curStats.TotalPaddingPackets += analyticsStream.TotalPaddingPackets + curStats.TotalPaddingBytes += analyticsStream.TotalPaddingBytes + curStats.TotalPacketsLost += analyticsStream.TotalPacketsLost + curStats.TotalFrames += analyticsStream.TotalFrames + curStats.TotalNacks += analyticsStream.TotalNacks + curStats.TotalPlis += analyticsStream.TotalPlis + curStats.TotalFirs += analyticsStream.TotalFirs + // add/update new video VideoLayers data to current and sum up video layer bytes/packets + for _, videoLayer := range analyticsStream.VideoLayers { + curStats.VideoLayers[videoLayer.Layer] = proto.Clone(videoLayer).(*livekit.AnalyticsVideoLayer) + curStats.TotalPackets += videoLayer.TotalPackets + curStats.TotalBytes += videoLayer.TotalBytes + } - streams = append(streams, stream) } - return &livekit.AnalyticsStat{ - Score: score, - Streams: streams, - } + // update currentStats + stats.curStats = &curStats + + return +} + +// find delta between curStats and prevStats and prepare proto payload +func (stats *Stats) computeDeltaStats() *livekit.AnalyticsStat { + + if stats.curStats == nil { + return nil + } + + // Stats in both queue/prev contain consolidated single deltaStats + cur := stats.curStats + + var maxLayer int32 + var maxTotalBytes uint64 + //create a map of VideoLayers - to pick max/best layer wrt current and prev + curLayers := make(map[int32]*livekit.AnalyticsVideoLayer) + for _, layer := range cur.VideoLayers { + curLayers[layer.Layer] = layer + // identify layer which sent max data - as VideoLayers can change in current interval + if layer.TotalBytes > maxTotalBytes { + maxTotalBytes = layer.TotalBytes + maxLayer = layer.Layer + } + } + + // no previous stats, prepare stat + if stats.prevStats == nil { + return cur.ToAnalyticsStats(curLayers[maxLayer]) + } + + // we have prevStats, find delta between cur and prev + prev := stats.prevStats + deltaStats := Stat{} + deltaStats.Rtt = cur.Rtt + deltaStats.Jitter = cur.Jitter + deltaStats.TotalPlis = cur.TotalPlis - prev.TotalPlis + deltaStats.TotalFrames = cur.TotalFrames - prev.TotalFrames + deltaStats.TotalNacks = cur.TotalNacks - prev.TotalNacks + deltaStats.TotalFirs = cur.TotalFirs - prev.TotalFirs + deltaStats.TotalPacketsLost = cur.TotalPacketsLost - prev.TotalPacketsLost + deltaStats.TotalPrimaryPackets = cur.TotalPrimaryPackets - prev.TotalPrimaryPackets + deltaStats.TotalRetransmitPackets = cur.TotalRetransmitPackets - prev.TotalRetransmitPackets + deltaStats.TotalPaddingPackets = cur.TotalPaddingPackets - prev.TotalPaddingPackets + deltaStats.TotalPaddingPackets = cur.TotalPaddingPackets - prev.TotalPaddingPackets + deltaStats.TotalPrimaryBytes = cur.TotalPrimaryBytes - prev.TotalPrimaryBytes + deltaStats.TotalPaddingBytes = cur.TotalPaddingBytes - prev.TotalPaddingBytes + deltaStats.TotalRetransmitBytes = cur.TotalRetransmitBytes - prev.TotalRetransmitBytes + + var videoLayer *livekit.AnalyticsVideoLayer + if len(cur.VideoLayers) > 0 && len(prev.VideoLayers) > 0 { + videoLayer = new(livekit.AnalyticsVideoLayer) + // find the current layer for the same layer id as previous, compute current round of delta with it + if curLayer, ok := curLayers[prev.MaxLayer]; ok { + videoLayer.Layer = prev.MaxLayer + videoLayer.TotalFrames = curLayer.TotalFrames - prev.VideoLayers[prev.MaxLayer].TotalFrames + } else { + videoLayer = curLayers[maxLayer] + } + // store new max layer for next round + cur.MaxLayer = maxLayer + // we accumulate bytes/packets across layers + videoLayer.TotalBytes = cur.TotalBytes - prev.TotalBytes + videoLayer.TotalPackets = cur.TotalPackets - prev.TotalPackets + } + // if no packets from any layers, return nil to send no stats + if deltaStats.TotalPackets == 0 && deltaStats.TotalPrimaryPackets == 0 && deltaStats.TotalRetransmitPackets == 0 && deltaStats.TotalPaddingPackets == 0 { + return nil + } + return deltaStats.ToAnalyticsStats(videoLayer) } diff --git a/pkg/telemetry/test/telemetry_service_test.go b/pkg/telemetry/test/telemetry_service_test.go index 8e8a2a3ab..aa42e2b66 100644 --- a/pkg/telemetry/test/telemetry_service_test.go +++ b/pkg/telemetry/test/telemetry_service_test.go @@ -226,7 +226,7 @@ func Test_PacketLostDiffShouldBeSentToTelemetry(t *testing.T) { _, stats = fixture.analytics.SendStatsArgsForCall(1) require.Equal(t, 1, len(stats)) require.Equal(t, livekit.StreamType_DOWNSTREAM, stats[0].Kind) - require.Equal(t, 4, int(stats[0].Streams[0].TotalPacketsLost)) // see diff of TotalLost between pkts2 and pkts1 + require.Equal(t, 3, int(stats[0].Streams[0].TotalPacketsLost)) // see diff of TotalLost between pkts2 and pkts1 } func Test_OnDownStreamRTCP_SeveralTracks(t *testing.T) {