diff --git a/pkg/config/config.go b/pkg/config/config.go index 2fe3d7fb4..023ae9f54 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -22,6 +22,8 @@ type CongestionControlProbeMode string const ( CongestionControlProbeModePadding CongestionControlProbeMode = "padding" CongestionControlProbeModeMedia CongestionControlProbeMode = "media" + + StatsUpdateFrequency = time.Second * 10 ) type Config struct { diff --git a/pkg/routing/redisrouter.go b/pkg/routing/redisrouter.go index db6f1f5d7..b5a60e6b8 100644 --- a/pkg/routing/redisrouter.go +++ b/pkg/routing/redisrouter.go @@ -39,6 +39,8 @@ type RedisRouter struct { ctx context.Context isStarted atomic.Bool statsMu sync.Mutex + // previous stats for computing averages + prevStats *livekit.NodeStats pubsub *redis.PubSub cancel func() @@ -479,13 +481,18 @@ func (r *RedisRouter) handleRTCMessage(rm *livekit.RTCNodeMessage) error { } r.statsMu.Lock() - updated, err := prometheus.GetUpdatedNodeStats(r.currentNode.Stats) + if r.prevStats == nil { + r.prevStats = r.currentNode.Stats + } + updated, computedAvg, err := prometheus.GetUpdatedNodeStats(r.currentNode.Stats, r.prevStats) if err != nil { logger.Errorw("could not update node stats", err) - } else { - if updated != nil { - r.currentNode.Stats = updated - } + r.statsMu.Unlock() + return err + } + r.currentNode.Stats = updated + if computedAvg { + r.prevStats = updated } r.statsMu.Unlock() diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index 81d1b2b40..384b4813b 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -624,8 +624,13 @@ func (p *ParticipantImpl) Close(sendLeave bool) error { if onClose != nil { onClose(p, disallowedSubscriptions) } - p.publisher.Close() - p.subscriber.Close() + + // Close peer connections without blocking participant close. If peer connections are gathering candidates + // Close will block. + go func() { + p.publisher.Close() + p.subscriber.Close() + }() return nil } diff --git a/pkg/service/roomservice.go b/pkg/service/roomservice.go index 531fb1c0b..71d4016f8 100644 --- a/pkg/service/roomservice.go +++ b/pkg/service/roomservice.go @@ -10,10 +10,9 @@ import ( "github.com/twitchtv/twirp" "google.golang.org/protobuf/proto" - "github.com/livekit/protocol/livekit" - "github.com/livekit/livekit-server/pkg/config" "github.com/livekit/livekit-server/pkg/routing" + "github.com/livekit/protocol/livekit" ) const ( diff --git a/pkg/telemetry/prometheus/node.go b/pkg/telemetry/prometheus/node.go index 011459f02..d3516ad7f 100644 --- a/pkg/telemetry/prometheus/node.go +++ b/pkg/telemetry/prometheus/node.go @@ -6,13 +6,13 @@ import ( "github.com/mackerelio/go-osstat/loadavg" "github.com/prometheus/client_golang/prometheus" + "github.com/livekit/livekit-server/pkg/config" "github.com/livekit/protocol/livekit" "github.com/livekit/protocol/utils" ) const ( - livekitNamespace string = "livekit" - forceUpdateInterval int64 = 15 + livekitNamespace string = "livekit" ) var ( @@ -49,36 +49,36 @@ func init() { initRoomStats(nodeID) } -func GetUpdatedNodeStats(prev *livekit.NodeStats) (*livekit.NodeStats, error) { +func GetUpdatedNodeStats(prev *livekit.NodeStats, prevAverage *livekit.NodeStats) (*livekit.NodeStats, bool, error) { loadAvg, err := loadavg.Get() if err != nil { - return nil, err + return nil, false, err } cpuLoad, numCPUs, err := getCPUStats() if err != nil { - return nil, err + return nil, false, err } - updatedAt := time.Now().Unix() - elapsed := updatedAt - prev.UpdatedAt - bytesInNow := bytesIn.Load() bytesOutNow := bytesOut.Load() packetsInNow := packetsIn.Load() packetsOutNow := packetsOut.Load() nackTotalNow := nackTotal.Load() - if bytesInNow == prev.BytesIn && - bytesOutNow == prev.BytesOut && - packetsInNow == prev.PacketsIn && - packetsOutNow == prev.PacketsOut && - nackTotalNow == prev.NackTotal && - elapsed < forceUpdateInterval { - return nil, nil + updatedAt := time.Now().Unix() + elapsed := updatedAt - prevAverage.UpdatedAt + // include sufficient buffer to be sure a stats update had taken place + computeAverage := elapsed > int64(config.StatsUpdateFrequency.Seconds()+2) + if bytesInNow != prevAverage.BytesIn || + bytesOutNow != prevAverage.BytesOut || + packetsInNow != prevAverage.PacketsIn || + packetsOutNow != prevAverage.PacketsOut || + nackTotalNow != prevAverage.NackTotal { + computeAverage = true } - return &livekit.NodeStats{ + stats := &livekit.NodeStats{ StartedAt: prev.StartedAt, UpdatedAt: updatedAt, NumRooms: roomTotal.Load(), @@ -90,17 +90,28 @@ func GetUpdatedNodeStats(prev *livekit.NodeStats) (*livekit.NodeStats, error) { PacketsIn: packetsInNow, PacketsOut: packetsOutNow, NackTotal: nackTotalNow, - BytesInPerSec: perSec(prev.BytesIn, bytesInNow, elapsed), - BytesOutPerSec: perSec(prev.BytesOut, bytesOutNow, elapsed), - PacketsInPerSec: perSec(prev.PacketsIn, packetsInNow, elapsed), - PacketsOutPerSec: perSec(prev.PacketsOut, packetsOutNow, elapsed), - NackPerSec: perSec(prev.NackTotal, nackTotalNow, elapsed), + BytesInPerSec: prevAverage.BytesInPerSec, + BytesOutPerSec: prevAverage.BytesOutPerSec, + PacketsInPerSec: prevAverage.PacketsInPerSec, + PacketsOutPerSec: prevAverage.PacketsOutPerSec, + NackPerSec: prevAverage.NackPerSec, NumCpus: numCPUs, CpuLoad: cpuLoad, LoadAvgLast1Min: float32(loadAvg.Loadavg1), LoadAvgLast5Min: float32(loadAvg.Loadavg5), LoadAvgLast15Min: float32(loadAvg.Loadavg15), - }, nil + } + + // update stats + if computeAverage { + stats.BytesInPerSec = perSec(prevAverage.BytesIn, bytesInNow, elapsed) + stats.BytesOutPerSec = perSec(prevAverage.BytesOut, bytesOutNow, elapsed) + stats.PacketsInPerSec = perSec(prevAverage.PacketsIn, packetsInNow, elapsed) + stats.PacketsInPerSec = perSec(prevAverage.PacketsOut, packetsOutNow, elapsed) + stats.NackPerSec = perSec(prevAverage.NackTotal, nackTotalNow, elapsed) + } + + return stats, computeAverage, nil } func perSec(prev, curr uint64, secs int64) float32 { diff --git a/pkg/telemetry/telemetryservice.go b/pkg/telemetry/telemetryservice.go index 55e2556e3..403f5ddbc 100644 --- a/pkg/telemetry/telemetryservice.go +++ b/pkg/telemetry/telemetryservice.go @@ -4,14 +4,13 @@ import ( "context" "time" + "github.com/livekit/livekit-server/pkg/config" "github.com/livekit/livekit-server/pkg/utils" "github.com/livekit/protocol/livekit" "github.com/livekit/protocol/logger" "github.com/livekit/protocol/webhook" ) -const updateFrequency = time.Second * 10 - //go:generate go run github.com/maxbrunsfeld/counterfeiter/v6 . TelemetryService type TelemetryService interface { // stats @@ -57,7 +56,7 @@ func NewTelemetryService(notifier webhook.Notifier, analytics AnalyticsService) } func (t *telemetryService) run() { - ticker := time.NewTicker(updateFrequency) + ticker := time.NewTicker(config.StatsUpdateFrequency) defer ticker.Stop() for { <-ticker.C