mirror of
https://github.com/livekit/livekit.git
synced 2026-08-29 07:39:09 +00:00
Compute delta stats to send downstream (#426)
* Compute delta stats to send downstream Signed-off-by: shishir gowda <shishir@livekit.io> * 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 <shishir@livekit.io> * 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 <shishir@livekit.io> * Use prev Max layer for current delta and update layer info for next round
This commit is contained in:
+206
-38
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user