From 7649e4ffab39143ba7ab41919e955f2cf4064a89 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Tue, 27 Feb 2024 15:45:32 +0530 Subject: [PATCH] Post data and signal stats once in 5 minutes (#2518) --- pkg/config/config.go | 5 +++-- pkg/telemetry/signalanddatastats.go | 25 ++++++++++++++++--------- 2 files changed, 19 insertions(+), 11 deletions(-) diff --git a/pkg/config/config.go b/pkg/config/config.go index c7d86c0bd..0530fb045 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -47,8 +47,9 @@ const ( StreamTrackerTypePacket StreamTrackerType = "packet" StreamTrackerTypeFrame StreamTrackerType = "frame" - StatsUpdateInterval = time.Second * 10 - TelemetryStatsUpdateInterval = time.Second * 30 + StatsUpdateInterval = time.Second * 10 + TelemetryStatsUpdateInterval = time.Second * 30 + TelemetryNonMediaStatsUpdateInterval = time.Minute * 5 ) var ( diff --git a/pkg/telemetry/signalanddatastats.go b/pkg/telemetry/signalanddatastats.go index 5ab5855ab..ffb223d55 100644 --- a/pkg/telemetry/signalanddatastats.go +++ b/pkg/telemetry/signalanddatastats.go @@ -20,6 +20,7 @@ import ( "go.uber.org/atomic" + "github.com/frostbyte73/core" "github.com/livekit/livekit-server/pkg/config" "github.com/livekit/protocol/livekit" "github.com/livekit/protocol/utils" @@ -53,7 +54,7 @@ type BytesTrackStats struct { totalSendBytes, totalRecvBytes atomic.Uint64 totalSendMessages, totalRecvMessages atomic.Uint32 telemetry TelemetryService - isStopped atomic.Bool + done core.Fuse } func NewBytesTrackStats(trackID livekit.TrackID, pID livekit.ParticipantID, telemetry TelemetryService) *BytesTrackStats { @@ -61,6 +62,7 @@ func NewBytesTrackStats(trackID livekit.TrackID, pID livekit.ParticipantID, tele trackID: trackID, pID: pID, telemetry: telemetry, + done: core.NewFuse(), } go s.reporter() return s @@ -91,7 +93,7 @@ func (s *BytesTrackStats) GetTrafficTotals() *TrafficTotals { } func (s *BytesTrackStats) Stop() { - s.isStopped.Store(true) + s.done.Break() } func (s *BytesTrackStats) report() { @@ -119,15 +121,20 @@ func (s *BytesTrackStats) report() { } func (s *BytesTrackStats) reporter() { - ticker := time.NewTicker(config.TelemetryStatsUpdateInterval) - defer ticker.Stop() - - for !s.isStopped.Load() { - <-ticker.C + ticker := time.NewTicker(config.TelemetryNonMediaStatsUpdateInterval) + defer func() { + ticker.Stop() s.report() - } + }() - s.report() + for { + select { + case <-s.done.Watch(): + return + case <-ticker.C: + s.report() + } + } } // -----------------------------------------------------------------------