mirror of
https://github.com/livekit/livekit.git
synced 2026-08-22 05:39:46 +00:00
Keep track of retransmissions in NodeStats (#677)
This commit is contained in:
@@ -157,7 +157,9 @@ func listNodes(c *cli.Context) error {
|
||||
"ID", "IP Address", "Region",
|
||||
"CPUs", "CPU Usage", "Load",
|
||||
"Clients", "Rooms", "Tracks In/Out",
|
||||
"Bytes In/Out", "Packets In/Out", "Nack", "Bps In/Out", "Pps In/Out", "Nack/Sec",
|
||||
"Bytes In/Out", "Packets In/Out",
|
||||
"Nack", "Retransmits",
|
||||
"Bps In/Out", "Pps In/Out", "Nack/Sec",
|
||||
"Started At", "Updated At",
|
||||
})
|
||||
for _, node := range nodes {
|
||||
@@ -177,6 +179,7 @@ func listNodes(c *cli.Context) error {
|
||||
bytes := fmt.Sprintf("%d / %d", stats.BytesIn, stats.BytesOut)
|
||||
packets := fmt.Sprintf("%d / %d", stats.PacketsIn, stats.PacketsOut)
|
||||
nack := strconv.Itoa(int(stats.NackTotal))
|
||||
retransmit := strconv.Itoa(int(stats.RetransmitPacketsOut))
|
||||
bps := fmt.Sprintf("%.2f / %.2f", stats.BytesInPerSec, stats.BytesOutPerSec)
|
||||
packetsPerSec := fmt.Sprintf("%.2f / %.2f", stats.PacketsInPerSec, stats.PacketsOutPerSec)
|
||||
nackPerSec := fmt.Sprintf("%f", stats.NackPerSec)
|
||||
@@ -188,7 +191,8 @@ func listNodes(c *cli.Context) error {
|
||||
node.Id, node.Ip, node.Region,
|
||||
cpus, cpuUsage, loadAvg,
|
||||
clients, rooms, tracks,
|
||||
bytes, packets, nack, bps, packetsPerSec, nackPerSec,
|
||||
bytes, packets, nack, retransmit,
|
||||
bps, packetsPerSec, nackPerSec,
|
||||
startedAt, updatedAt,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -12,7 +12,7 @@ require (
|
||||
github.com/google/wire v0.5.0
|
||||
github.com/gorilla/websocket v1.4.2
|
||||
github.com/hashicorp/golang-lru v0.5.4
|
||||
github.com/livekit/protocol v0.13.2-0.20220502170852-688e4f627bcf
|
||||
github.com/livekit/protocol v0.13.3-0.20220510071353-084233d23a03
|
||||
github.com/mackerelio/go-osstat v0.2.1
|
||||
github.com/magefile/mage v1.11.0
|
||||
github.com/maxbrunsfeld/counterfeiter/v6 v6.3.0
|
||||
@@ -82,7 +82,6 @@ require (
|
||||
github.com/cpuguy83/go-md2man/v2 v2.0.0 // indirect
|
||||
github.com/d5/tengo/v2 v2.10.1
|
||||
github.com/google/subcommands v1.2.0 // indirect
|
||||
github.com/rs/zerolog v1.26.1
|
||||
github.com/russross/blackfriday/v2 v2.1.0 // indirect
|
||||
golang.org/x/crypto v0.0.0-20220411220226-7b82a4e95df4 // indirect
|
||||
golang.org/x/mod v0.5.1 // indirect
|
||||
|
||||
@@ -25,7 +25,6 @@ github.com/cncf/udpa/go v0.0.0-20210930031921-04548b0d99d4/go.mod h1:6pvJx4me5XP
|
||||
github.com/cncf/xds/go v0.0.0-20210805033703-aa0b78936158/go.mod h1:eXthEFrGJvWHgFFCl3hGmgk+/aYT6PnTQLykKQRLhEs=
|
||||
github.com/cncf/xds/go v0.0.0-20210922020428-25de7278fc84/go.mod h1:eXthEFrGJvWHgFFCl3hGmgk+/aYT6PnTQLykKQRLhEs=
|
||||
github.com/cncf/xds/go v0.0.0-20211011173535-cb28da3451f1/go.mod h1:eXthEFrGJvWHgFFCl3hGmgk+/aYT6PnTQLykKQRLhEs=
|
||||
github.com/coreos/go-systemd/v22 v22.3.2/go.mod h1:Y58oyj3AT4RCenI/lSvhwexgC+NSVTIJ3seZv2GcEnc=
|
||||
github.com/cpuguy83/go-md2man/v2 v2.0.0-20190314233015-f79a8a8ca69d/go.mod h1:maD7wRr/U5Z6m/iR4s+kqSMx2CaBsrgA7czyZG/E6dU=
|
||||
github.com/cpuguy83/go-md2man/v2 v2.0.0 h1:EoUDS0afbrsXAZ9YQ9jdu/mZ2sXgT1/2yyNng4PGlyM=
|
||||
github.com/cpuguy83/go-md2man/v2 v2.0.0/go.mod h1:maD7wRr/U5Z6m/iR4s+kqSMx2CaBsrgA7czyZG/E6dU=
|
||||
@@ -70,7 +69,6 @@ github.com/go-redis/redis/v8 v8.11.3 h1:GCjoYp8c+yQTJfc0n69iwSiHjvuAdruxl7elnZCx
|
||||
github.com/go-redis/redis/v8 v8.11.3/go.mod h1:xNJ9xDG09FsIPwh3bWdk+0oDWHbtF9rPN0F/oD9XeKc=
|
||||
github.com/go-stack/stack v1.8.0/go.mod h1:v0f6uXyyMGvRgIKkXu+yp6POWl0qKG85gN/melR3HDY=
|
||||
github.com/go-task/slim-sprig v0.0.0-20210107165309-348f09dbbbc0/go.mod h1:fyg7847qk6SyHyPtNmDHnmrv/HOrqktSC+C9fM+CJOE=
|
||||
github.com/godbus/dbus/v5 v5.0.4/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA=
|
||||
github.com/gogo/protobuf v1.1.1/go.mod h1:r8qH/GZQm5c6nD/R0oafs1akxWv10x8SbQlK7atdtwQ=
|
||||
github.com/golang/glog v0.0.0-20160126235308-23def4e6c14b/go.mod h1:SBH7ygxi8pfUlaOkMMuAQtPIUF8ecWP5IEl/CR7VP2Q=
|
||||
github.com/golang/mock v1.1.1/go.mod h1:oTYuIxOrZwtPieC+H1uAHpcLFnEyAGVDL/k47Jfbm0A=
|
||||
@@ -131,8 +129,8 @@ github.com/kr/text v0.1.0 h1:45sCR5RtlFHMR4UwH9sdQ5TC8v0qDQCHnXt+kaKSTVE=
|
||||
github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI=
|
||||
github.com/lithammer/shortuuid/v3 v3.0.6 h1:pr15YQyvhiSX/qPxncFtqk+v4xLEpOZObbsY/mKrcvA=
|
||||
github.com/lithammer/shortuuid/v3 v3.0.6/go.mod h1:vMk8ke37EmiewwolSO1NLW8vP4ZaKlRuDIi8tWWmAts=
|
||||
github.com/livekit/protocol v0.13.2-0.20220502170852-688e4f627bcf h1:x7l15X3hAFnGAd1hd9rfHkliAbtaRRcx0agB/Tb7rWQ=
|
||||
github.com/livekit/protocol v0.13.2-0.20220502170852-688e4f627bcf/go.mod h1:BLtSeVmn2rLP37xjzw7gHgaAmkWl3L/L9bPvgSbaOfo=
|
||||
github.com/livekit/protocol v0.13.3-0.20220510071353-084233d23a03 h1:UI0AeS8kCu0yEfwYUEmj3WUyMpS/VBSKojGenQ9WeUM=
|
||||
github.com/livekit/protocol v0.13.3-0.20220510071353-084233d23a03/go.mod h1:BLtSeVmn2rLP37xjzw7gHgaAmkWl3L/L9bPvgSbaOfo=
|
||||
github.com/mackerelio/go-osstat v0.2.1 h1:5AeAcBEutEErAOlDz6WCkEvm6AKYgHTUQrfwm5RbeQc=
|
||||
github.com/mackerelio/go-osstat v0.2.1/go.mod h1:UzRL8dMCCTqG5WdRtsxbuljMpZt9PCAGXqxPst5QtaY=
|
||||
github.com/magefile/mage v1.11.0 h1:C/55Ywp9BpgVVclD3lRnSYCwXTYxmSppIgLeDYlNuls=
|
||||
@@ -237,9 +235,6 @@ github.com/prometheus/procfs v0.6.0/go.mod h1:cz+aTbrPOrUb4q7XlbU9ygM+/jj0fzG6c1
|
||||
github.com/rogpeppe/fastuuid v1.2.0/go.mod h1:jVj6XXZzXRy/MSR5jhDC/2q6DgLz+nrA6LYCDYWNEvQ=
|
||||
github.com/rs/cors v1.8.2 h1:KCooALfAYGs415Cwu5ABvv9n9509fSiG5SQJn/AQo4U=
|
||||
github.com/rs/cors v1.8.2/go.mod h1:XyqrcTp5zjWr1wsJ8PIRZssZ8b/WMcMf71DJnit4EMU=
|
||||
github.com/rs/xid v1.3.0/go.mod h1:trrq9SKmegXys3aeAKXMUTdJsYXVwGY3RLcfgqegfbg=
|
||||
github.com/rs/zerolog v1.26.1 h1:/ihwxqH+4z8UxyI70wM1z9yCvkWcfz/a3mj48k/Zngc=
|
||||
github.com/rs/zerolog v1.26.1/go.mod h1:/wSSJWX7lVrsOwlbyTRSOJvqRlc+WjWlfes+CiJ+tmc=
|
||||
github.com/russross/blackfriday/v2 v2.0.1/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM=
|
||||
github.com/russross/blackfriday/v2 v2.1.0 h1:JIOH55/0cWyOuilr9/qlrm0BSXldqnqwMsf35Ld67mk=
|
||||
github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM=
|
||||
@@ -292,7 +287,6 @@ golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACk
|
||||
golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
|
||||
golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto=
|
||||
golang.org/x/crypto v0.0.0-20210314154223-e6e6c4f2bb5b/go.mod h1:T9bdIzuCu7OtxOm1hfPfRQxPLYneinmdGuTeoZ9dtd4=
|
||||
golang.org/x/crypto v0.0.0-20211215165025-cf75a172585e/go.mod h1:P+XmwS30IXTQdn5tA2iutPOUgjI07+tq3H3K9MVA1s8=
|
||||
golang.org/x/crypto v0.0.0-20220131195533-30dcbda58838/go.mod h1:IxCIyHEi3zRg3s0A5j5BB6A9Jmi73HwBIUl50j+osU4=
|
||||
golang.org/x/crypto v0.0.0-20220411220226-7b82a4e95df4 h1:kUhD7nTDoI3fVd9G4ORWrbV5NY0liEs/Jg2pv5f+bBA=
|
||||
golang.org/x/crypto v0.0.0-20220411220226-7b82a4e95df4/go.mod h1:IxCIyHEi3zRg3s0A5j5BB6A9Jmi73HwBIUl50j+osU4=
|
||||
|
||||
@@ -65,6 +65,8 @@ func GetUpdatedNodeStats(prev *livekit.NodeStats, prevAverage *livekit.NodeStats
|
||||
packetsInNow := packetsIn.Load()
|
||||
packetsOutNow := packetsOut.Load()
|
||||
nackTotalNow := nackTotal.Load()
|
||||
retransmitBytesNow := retransmitBytes.Load()
|
||||
retransmitPacketsNow := retransmitPackets.Load()
|
||||
|
||||
updatedAt := time.Now().Unix()
|
||||
elapsed := updatedAt - prevAverage.UpdatedAt
|
||||
@@ -73,32 +75,38 @@ func GetUpdatedNodeStats(prev *livekit.NodeStats, prevAverage *livekit.NodeStats
|
||||
if bytesInNow != prevAverage.BytesIn ||
|
||||
bytesOutNow != prevAverage.BytesOut ||
|
||||
packetsInNow != prevAverage.PacketsIn ||
|
||||
packetsOutNow != prevAverage.PacketsOut {
|
||||
packetsOutNow != prevAverage.PacketsOut ||
|
||||
retransmitBytesNow != prevAverage.RetransmitBytesOut ||
|
||||
retransmitPacketsNow != prevAverage.RetransmitPacketsOut {
|
||||
computeAverage = true
|
||||
}
|
||||
|
||||
stats := &livekit.NodeStats{
|
||||
StartedAt: prev.StartedAt,
|
||||
UpdatedAt: updatedAt,
|
||||
NumRooms: roomTotal.Load(),
|
||||
NumClients: participantTotal.Load(),
|
||||
NumTracksIn: trackPublishedTotal.Load(),
|
||||
NumTracksOut: trackSubscribedTotal.Load(),
|
||||
BytesIn: bytesInNow,
|
||||
BytesOut: bytesOutNow,
|
||||
PacketsIn: packetsInNow,
|
||||
PacketsOut: packetsOutNow,
|
||||
NackTotal: nackTotalNow,
|
||||
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),
|
||||
StartedAt: prev.StartedAt,
|
||||
UpdatedAt: updatedAt,
|
||||
NumRooms: roomTotal.Load(),
|
||||
NumClients: participantTotal.Load(),
|
||||
NumTracksIn: trackPublishedTotal.Load(),
|
||||
NumTracksOut: trackSubscribedTotal.Load(),
|
||||
BytesIn: bytesInNow,
|
||||
BytesOut: bytesOutNow,
|
||||
PacketsIn: packetsInNow,
|
||||
PacketsOut: packetsOutNow,
|
||||
RetransmitBytesOut: retransmitBytesNow,
|
||||
RetransmitPacketsOut: retransmitPacketsNow,
|
||||
NackTotal: nackTotalNow,
|
||||
BytesInPerSec: prevAverage.BytesInPerSec,
|
||||
BytesOutPerSec: prevAverage.BytesOutPerSec,
|
||||
PacketsInPerSec: prevAverage.PacketsInPerSec,
|
||||
PacketsOutPerSec: prevAverage.PacketsOutPerSec,
|
||||
RetransmitBytesOutPerSec: prevAverage.RetransmitBytesOutPerSec,
|
||||
RetransmitPacketsOutPerSec: prevAverage.RetransmitPacketsOutPerSec,
|
||||
NackPerSec: prevAverage.NackPerSec,
|
||||
NumCpus: numCPUs,
|
||||
CpuLoad: cpuLoad,
|
||||
LoadAvgLast1Min: float32(loadAvg.Loadavg1),
|
||||
LoadAvgLast5Min: float32(loadAvg.Loadavg5),
|
||||
LoadAvgLast15Min: float32(loadAvg.Loadavg15),
|
||||
}
|
||||
|
||||
// update stats
|
||||
@@ -107,6 +115,8 @@ func GetUpdatedNodeStats(prev *livekit.NodeStats, prevAverage *livekit.NodeStats
|
||||
stats.BytesOutPerSec = perSec(prevAverage.BytesOut, bytesOutNow, elapsed)
|
||||
stats.PacketsInPerSec = perSec(prevAverage.PacketsIn, packetsInNow, elapsed)
|
||||
stats.PacketsOutPerSec = perSec(prevAverage.PacketsOut, packetsOutNow, elapsed)
|
||||
stats.RetransmitBytesOutPerSec = perSec(prevAverage.RetransmitBytesOut, retransmitBytesNow, elapsed)
|
||||
stats.RetransmitPacketsOutPerSec = perSec(prevAverage.RetransmitPacketsOut, retransmitPacketsNow, elapsed)
|
||||
stats.NackPerSec = perSec(prevAverage.NackTotal, nackTotalNow, elapsed)
|
||||
}
|
||||
|
||||
|
||||
@@ -8,24 +8,28 @@ import (
|
||||
type Direction string
|
||||
|
||||
const (
|
||||
Incoming Direction = "incoming"
|
||||
Outgoing Direction = "outgoing"
|
||||
Incoming Direction = "incoming"
|
||||
Outgoing Direction = "outgoing"
|
||||
transmissionInitial = "initial"
|
||||
transmissionRetransmit = "retransmit"
|
||||
)
|
||||
|
||||
var (
|
||||
bytesIn atomic.Uint64
|
||||
bytesOut atomic.Uint64
|
||||
packetsIn atomic.Uint64
|
||||
packetsOut atomic.Uint64
|
||||
nackTotal atomic.Uint64
|
||||
bytesIn atomic.Uint64
|
||||
bytesOut atomic.Uint64
|
||||
packetsIn atomic.Uint64
|
||||
packetsOut atomic.Uint64
|
||||
nackTotal atomic.Uint64
|
||||
retransmitBytes atomic.Uint64
|
||||
retransmitPackets atomic.Uint64
|
||||
|
||||
promPacketLabels = []string{"direction"}
|
||||
|
||||
promPacketTotal *prometheus.CounterVec
|
||||
promPacketBytes *prometheus.CounterVec
|
||||
promNackTotal *prometheus.CounterVec
|
||||
promPliTotal *prometheus.CounterVec
|
||||
promFirTotal *prometheus.CounterVec
|
||||
promPacketLabels = []string{"direction", "transmission"}
|
||||
promPacketTotal *prometheus.CounterVec
|
||||
promPacketBytes *prometheus.CounterVec
|
||||
promRTCPLabels = []string{"direction"}
|
||||
promNackTotal *prometheus.CounterVec
|
||||
promPliTotal *prometheus.CounterVec
|
||||
promFirTotal *prometheus.CounterVec
|
||||
)
|
||||
|
||||
func initPacketStats(nodeID string) {
|
||||
@@ -46,19 +50,19 @@ func initPacketStats(nodeID string) {
|
||||
Subsystem: "nack",
|
||||
Name: "total",
|
||||
ConstLabels: prometheus.Labels{"node_id": nodeID},
|
||||
}, promPacketLabels)
|
||||
}, promRTCPLabels)
|
||||
promPliTotal = prometheus.NewCounterVec(prometheus.CounterOpts{
|
||||
Namespace: livekitNamespace,
|
||||
Subsystem: "pli",
|
||||
Name: "total",
|
||||
ConstLabels: prometheus.Labels{"node_id": nodeID},
|
||||
}, promPacketLabels)
|
||||
}, promRTCPLabels)
|
||||
promFirTotal = prometheus.NewCounterVec(prometheus.CounterOpts{
|
||||
Namespace: livekitNamespace,
|
||||
Subsystem: "fir",
|
||||
Name: "total",
|
||||
ConstLabels: prometheus.Labels{"node_id": nodeID},
|
||||
}, promPacketLabels)
|
||||
}, promRTCPLabels)
|
||||
|
||||
prometheus.MustRegister(promPacketTotal)
|
||||
prometheus.MustRegister(promPacketBytes)
|
||||
@@ -67,21 +71,33 @@ func initPacketStats(nodeID string) {
|
||||
prometheus.MustRegister(promFirTotal)
|
||||
}
|
||||
|
||||
func IncrementPackets(direction Direction, count uint64) {
|
||||
promPacketTotal.WithLabelValues(string(direction)).Add(float64(count))
|
||||
func IncrementPackets(direction Direction, count uint64, retransmit bool) {
|
||||
promPacketTotal.WithLabelValues(
|
||||
string(direction),
|
||||
transmissionLabel(retransmit),
|
||||
).Add(float64(count))
|
||||
if direction == Incoming {
|
||||
packetsIn.Add(count)
|
||||
} else {
|
||||
packetsOut.Add(count)
|
||||
if retransmit {
|
||||
retransmitPackets.Add(count)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func IncrementBytes(direction Direction, count uint64) {
|
||||
promPacketBytes.WithLabelValues(string(direction)).Add(float64(count))
|
||||
func IncrementBytes(direction Direction, count uint64, retransmit bool) {
|
||||
promPacketBytes.WithLabelValues(
|
||||
string(direction),
|
||||
transmissionLabel(retransmit),
|
||||
).Add(float64(count))
|
||||
if direction == Incoming {
|
||||
bytesIn.Add(count)
|
||||
} else {
|
||||
bytesOut.Add(count)
|
||||
if retransmit {
|
||||
retransmitBytes.Add(count)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -97,3 +113,11 @@ func IncrementRTCP(direction Direction, nack, pli, fir uint32) {
|
||||
promFirTotal.WithLabelValues(string(direction)).Add(float64(fir))
|
||||
}
|
||||
}
|
||||
|
||||
func transmissionLabel(retransmit bool) string {
|
||||
if retransmit {
|
||||
return transmissionInitial
|
||||
} else {
|
||||
return transmissionRetransmit
|
||||
}
|
||||
}
|
||||
|
||||
@@ -55,16 +55,32 @@ func (t *telemetryServiceInternal) TrackStats(streamType livekit.StreamType, par
|
||||
firs := uint32(0)
|
||||
packets := uint32(0)
|
||||
bytes := uint64(0)
|
||||
retransmitBytes := uint64(0)
|
||||
retransmitPackets := uint32(0)
|
||||
for _, stream := range stat.Streams {
|
||||
nacks += stream.Nacks
|
||||
plis += stream.Plis
|
||||
firs += stream.Firs
|
||||
packets += stream.PrimaryPackets + stream.RetransmitPackets + stream.PaddingPackets
|
||||
bytes += stream.PrimaryBytes + stream.RetransmitBytes + stream.PaddingBytes
|
||||
packets += stream.PrimaryPackets + stream.PaddingPackets
|
||||
bytes += stream.PrimaryBytes + stream.PaddingBytes
|
||||
if streamType == livekit.StreamType_DOWNSTREAM {
|
||||
retransmitPackets += stream.RetransmitPackets
|
||||
retransmitBytes += stream.RetransmitBytes
|
||||
} else {
|
||||
// for upstream, we don't account for these separately for now
|
||||
packets += stream.RetransmitPackets
|
||||
bytes += stream.RetransmitBytes
|
||||
}
|
||||
}
|
||||
prometheus.IncrementRTCP(direction, nacks, plis, firs)
|
||||
prometheus.IncrementPackets(direction, uint64(packets))
|
||||
prometheus.IncrementBytes(direction, bytes)
|
||||
prometheus.IncrementPackets(direction, uint64(packets), false)
|
||||
prometheus.IncrementBytes(direction, bytes, false)
|
||||
if retransmitPackets != 0 {
|
||||
prometheus.IncrementPackets(direction, uint64(retransmitPackets), true)
|
||||
}
|
||||
if retransmitBytes != 0 {
|
||||
prometheus.IncrementBytes(direction, retransmitBytes, true)
|
||||
}
|
||||
|
||||
if w := t.getStatsWorker(participantID); w != nil {
|
||||
w.OnTrackStat(trackID, streamType, stat)
|
||||
|
||||
Reference in New Issue
Block a user