From 3b0077f2fe95aa068ac8470688ec9545898d7c31 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Tue, 7 Jan 2025 10:58:31 +0530 Subject: [PATCH] Log connection quality changes. (#3311) Also remove the connection quality drop prom as it is unused and also adds state/complexity. --- pkg/rtc/participant.go | 59 ++++++----------------------- pkg/telemetry/prometheus/quality.go | 12 +----- 2 files changed, 13 insertions(+), 58 deletions(-) diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index fbcbdd228..11a1c79ad 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -253,7 +253,7 @@ type ParticipantImpl struct { supervisor *supervisor.ParticipantSupervisor - tracksQuality map[livekit.TrackID]livekit.ConnectionQuality + connectionQuality livekit.ConnectionQuality metricTimestamper *metric.MetricTimestamper metricsCollector *metric.MetricsCollector @@ -292,9 +292,9 @@ func NewParticipant(params ParticipantParams) (*ParticipantImpl, error) { params.SID, params.Telemetry, ), - tracksQuality: make(map[livekit.TrackID]livekit.ConnectionQuality), - pubLogger: params.Logger.WithComponent(sutils.ComponentPub), - subLogger: params.Logger.WithComponent(sutils.ComponentSub), + connectionQuality: livekit.ConnectionQuality_EXCELLENT, + pubLogger: params.Logger.WithComponent(sutils.ComponentPub), + subLogger: params.Logger.WithComponent(sutils.ComponentSub), } if !params.DisableSupervisor { p.supervisor = supervisor.NewParticipantSupervisor(supervisor.ParticipantSupervisorParams{Logger: params.Logger}) @@ -1249,17 +1249,10 @@ func (p *ParticipantImpl) OnICEConfigChanged(f func(participant types.LocalParti } func (p *ParticipantImpl) GetConnectionQuality() *livekit.ConnectionQualityInfo { - numTracks := 0 minQuality := livekit.ConnectionQuality_EXCELLENT minScore := connectionquality.MaxMOS - numUpDrops := 0 - numDownDrops := 0 - - availableTracks := make(map[livekit.TrackID]bool) for _, pt := range p.GetPublishedTracks() { - numTracks++ - score, quality := pt.(types.LocalMediaTrack).GetConnectionScoreAndQuality() if utils.IsConnectionQualityLower(minQuality, quality) { minQuality = quality @@ -1267,24 +1260,10 @@ func (p *ParticipantImpl) GetConnectionQuality() *livekit.ConnectionQualityInfo } else if quality == minQuality && score < minScore { minScore = score } - - p.lock.Lock() - trackID := pt.ID() - if prevQuality, ok := p.tracksQuality[trackID]; ok { - if utils.IsConnectionQualityLower(prevQuality, quality) { - numUpDrops++ - } - } - p.tracksQuality[trackID] = quality - p.lock.Unlock() - - availableTracks[trackID] = true } subscribedTracks := p.SubscriptionManager.GetSubscribedTracks() for _, subTrack := range subscribedTracks { - numTracks++ - score, quality := subTrack.DownTrack().GetConnectionScoreAndQuality() if utils.IsConnectionQualityLower(minQuality, quality) { minQuality = quality @@ -1292,35 +1271,21 @@ func (p *ParticipantImpl) GetConnectionQuality() *livekit.ConnectionQualityInfo } else if quality == minQuality && score < minScore { minScore = score } - - p.lock.Lock() - trackID := subTrack.ID() - if prevQuality, ok := p.tracksQuality[trackID]; ok { - if utils.IsConnectionQualityLower(prevQuality, quality) { - numDownDrops++ - } - } - p.tracksQuality[trackID] = quality - p.lock.Unlock() - - availableTracks[trackID] = true } - prometheus.RecordQuality(minQuality, minScore, numUpDrops, numDownDrops) - - // remove unavailable tracks from track quality cache - p.lock.Lock() - for trackID := range p.tracksQuality { - if !availableTracks[trackID] { - delete(p.tracksQuality, trackID) - } - } - p.lock.Unlock() + prometheus.RecordQuality(minQuality, minScore) if minQuality == livekit.ConnectionQuality_LOST && !p.ProtocolVersion().SupportsConnectionQualityLost() { minQuality = livekit.ConnectionQuality_POOR } + p.lock.Lock() + if minQuality != p.connectionQuality { + p.params.Logger.Debugw("connection quality changed", "from", p.connectionQuality, "to", minQuality) + } + p.connectionQuality = minQuality + p.lock.Unlock() + return &livekit.ConnectionQualityInfo{ ParticipantSid: string(p.ID()), Quality: minQuality, diff --git a/pkg/telemetry/prometheus/quality.go b/pkg/telemetry/prometheus/quality.go index c9eaf6b65..b55a3f8b6 100644 --- a/pkg/telemetry/prometheus/quality.go +++ b/pkg/telemetry/prometheus/quality.go @@ -23,7 +23,6 @@ import ( var ( qualityRating prometheus.Histogram qualityScore prometheus.Histogram - qualityDrop *prometheus.CounterVec ) func initQualityStats(nodeID string, nodeType livekit.NodeType) { @@ -41,21 +40,12 @@ func initQualityStats(nodeID string, nodeType livekit.NodeType) { ConstLabels: prometheus.Labels{"node_id": nodeID, "node_type": nodeType.String()}, Buckets: []float64{1.0, 2.0, 2.5, 3.0, 3.25, 3.5, 3.75, 4.0, 4.25, 4.5}, }) - qualityDrop = prometheus.NewCounterVec(prometheus.CounterOpts{ - Namespace: livekitNamespace, - Subsystem: "quality", - Name: "drop", - ConstLabels: prometheus.Labels{"node_id": nodeID, "node_type": nodeType.String()}, - }, []string{"direction"}) prometheus.MustRegister(qualityRating) prometheus.MustRegister(qualityScore) - prometheus.MustRegister(qualityDrop) } -func RecordQuality(rating livekit.ConnectionQuality, score float32, numUpDrops int, numDownDrops int) { +func RecordQuality(rating livekit.ConnectionQuality, score float32) { qualityRating.Observe(float64(rating)) qualityScore.Observe(float64(score)) - qualityDrop.WithLabelValues("up").Add(float64(numUpDrops)) - qualityDrop.WithLabelValues("down").Add(float64(numDownDrops)) }