mirror of
https://github.com/livekit/livekit.git
synced 2026-08-27 22:34:25 +00:00
Increase frequency of status updates and longer availability threshold (#628)
* Increase frequency of status updates and longer avail. threshold. * better fix. * fix room close test failure due to slow peer connection Close * Perform avg computation more frequently if data has changed
This commit is contained in:
@@ -22,6 +22,8 @@ type CongestionControlProbeMode string
|
||||
const (
|
||||
CongestionControlProbeModePadding CongestionControlProbeMode = "padding"
|
||||
CongestionControlProbeModeMedia CongestionControlProbeMode = "media"
|
||||
|
||||
StatsUpdateFrequency = time.Second * 10
|
||||
)
|
||||
|
||||
type Config struct {
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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 (
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user