From 7eb3362d0aab410f15cdd08f4a2c9f2e446268db Mon Sep 17 00:00:00 2001 From: David Zhao Date: Tue, 10 May 2022 15:25:24 -0700 Subject: [PATCH] Keep track of retransmissions in NodeStats (#677) --- cmd/server/commands.go | 8 ++- go.mod | 3 +- go.sum | 10 +--- pkg/telemetry/prometheus/node.go | 54 +++++++++++-------- pkg/telemetry/prometheus/packets.go | 66 +++++++++++++++-------- pkg/telemetry/telemetryserviceinternal.go | 24 +++++++-- 6 files changed, 106 insertions(+), 59 deletions(-) diff --git a/cmd/server/commands.go b/cmd/server/commands.go index a8570788a..5e96a1fc4 100644 --- a/cmd/server/commands.go +++ b/cmd/server/commands.go @@ -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, }) } diff --git a/go.mod b/go.mod index 42c50f1e4..de12f2b60 100644 --- a/go.mod +++ b/go.mod @@ -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 diff --git a/go.sum b/go.sum index 9c163700b..af907281e 100644 --- a/go.sum +++ b/go.sum @@ -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= diff --git a/pkg/telemetry/prometheus/node.go b/pkg/telemetry/prometheus/node.go index 25c394cf0..d07b1d451 100644 --- a/pkg/telemetry/prometheus/node.go +++ b/pkg/telemetry/prometheus/node.go @@ -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) } diff --git a/pkg/telemetry/prometheus/packets.go b/pkg/telemetry/prometheus/packets.go index c3c6c5fc5..17f1020fb 100644 --- a/pkg/telemetry/prometheus/packets.go +++ b/pkg/telemetry/prometheus/packets.go @@ -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 + } +} diff --git a/pkg/telemetry/telemetryserviceinternal.go b/pkg/telemetry/telemetryserviceinternal.go index 8a189b3dd..68fca98e0 100644 --- a/pkg/telemetry/telemetryserviceinternal.go +++ b/pkg/telemetry/telemetryserviceinternal.go @@ -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)