forked from livekit/livekit
/
stats.go
84 lines (77 loc) · 2.67 KB
/
stats.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
package telemetry
import (
"github.com/whoyao/livekit/pkg/telemetry/prometheus"
"github.com/whoyao/protocol/livekit"
)
type StatsKey struct {
streamType livekit.StreamType
participantID livekit.ParticipantID
trackID livekit.TrackID
trackSource livekit.TrackSource
trackType livekit.TrackType
track bool
}
func StatsKeyForTrack(streamType livekit.StreamType, participantID livekit.ParticipantID, trackID livekit.TrackID, trackSource livekit.TrackSource, trackType livekit.TrackType) StatsKey {
return StatsKey{
streamType: streamType,
participantID: participantID,
trackID: trackID,
trackSource: trackSource,
trackType: trackType,
track: true,
}
}
func StatsKeyForData(streamType livekit.StreamType, participantID livekit.ParticipantID, trackID livekit.TrackID) StatsKey {
return StatsKey{
streamType: streamType,
participantID: participantID,
trackID: trackID,
}
}
func (t *telemetryService) TrackStats(key StatsKey, stat *livekit.AnalyticsStat) {
t.enqueue(func() {
direction := prometheus.Incoming
if key.streamType == livekit.StreamType_DOWNSTREAM {
direction = prometheus.Outgoing
}
nacks := uint32(0)
plis := uint32(0)
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.PaddingPackets
bytes += stream.PrimaryBytes + stream.PaddingBytes
if key.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
}
if key.track {
prometheus.RecordPacketLoss(direction, key.trackSource, key.trackType, stream.PacketsLost, stream.PrimaryPackets+stream.PaddingPackets)
prometheus.RecordRTT(direction, key.trackSource, key.trackType, stream.Rtt)
prometheus.RecordJitter(direction, key.trackSource, key.trackType, stream.Jitter)
}
}
prometheus.IncrementRTCP(direction, nacks, plis, firs)
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 worker, ok := t.getWorker(key.participantID); ok {
worker.OnTrackStat(key.trackID, key.streamType, stat)
}
})
}