Add statsd

This commit is contained in:
9seconds
2021-03-17 11:07:16 +03:00
parent c04f9e392a
commit 8e7207d975
5 changed files with 256 additions and 0 deletions
+1
View File
@@ -16,6 +16,7 @@ require (
github.com/panjf2000/ants v1.3.0 // indirect github.com/panjf2000/ants v1.3.0 // indirect
github.com/pelletier/go-toml v1.8.1 github.com/pelletier/go-toml v1.8.1
github.com/rs/zerolog v1.20.0 // indirect github.com/rs/zerolog v1.20.0 // indirect
github.com/smira/go-statsd v1.3.2 // indirect
github.com/stretchr/objx v0.3.0 // indirect github.com/stretchr/objx v0.3.0 // indirect
github.com/stretchr/testify v1.7.0 github.com/stretchr/testify v1.7.0
github.com/tylertreat/BoomFilters v0.0.0-20200520150052-42a7b4300c0c // indirect github.com/tylertreat/BoomFilters v0.0.0-20200520150052-42a7b4300c0c // indirect
+2
View File
@@ -37,6 +37,8 @@ github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZN
github.com/rs/xid v1.2.1/go.mod h1:+uKXf+4Djp6Md1KODXJxgGQPKngRmWyn10oCKFzNHOQ= github.com/rs/xid v1.2.1/go.mod h1:+uKXf+4Djp6Md1KODXJxgGQPKngRmWyn10oCKFzNHOQ=
github.com/rs/zerolog v1.20.0 h1:38k9hgtUBdxFwE34yS8rTHmHBa4eN16E4DJlv177LNs= github.com/rs/zerolog v1.20.0 h1:38k9hgtUBdxFwE34yS8rTHmHBa4eN16E4DJlv177LNs=
github.com/rs/zerolog v1.20.0/go.mod h1:IzD0RJ65iWH0w97OQQebJEvTZYvsCUm9WVLWBQrJRjo= github.com/rs/zerolog v1.20.0/go.mod h1:IzD0RJ65iWH0w97OQQebJEvTZYvsCUm9WVLWBQrJRjo=
github.com/smira/go-statsd v1.3.2 h1:1EeuzxNZ/TD9apbTOFSM9nulqfcsQFmT4u1A2DREabI=
github.com/smira/go-statsd v1.3.2/go.mod h1:1srXJ9/pbnN04G8f4F1jUzsGOnwkPKXciyqpewGlkC4=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/objx v0.3.0 h1:NGXK3lHquSN08v5vWalVI/L8XU9hdzE/G6xsrze47As= github.com/stretchr/objx v0.3.0 h1:NGXK3lHquSN08v5vWalVI/L8XU9hdzE/G6xsrze47As=
github.com/stretchr/objx v0.3.0/go.mod h1:qt09Ya8vawLte6SNmTgCsAVtYtaKzEcn8ATUoHMkEqE= github.com/stretchr/objx v0.3.0/go.mod h1:qt09Ya8vawLte6SNmTgCsAVtYtaKzEcn8ATUoHMkEqE=
+12
View File
@@ -0,0 +1,12 @@
package stats
const (
MetricActiveConnection = "active_connections"
MetricSessionDuration = "session_duration"
MetricConcurrencyLimited = "concurrency_limited"
TagIPType = "ip_type"
TagIPTypeIPv4 = "ipv4"
TagIPTypeIPv6 = "ipv4"
)
+115
View File
@@ -0,0 +1,115 @@
package stats
import (
"fmt"
"net"
"strings"
"time"
"github.com/9seconds/mtg/v2/events"
"github.com/9seconds/mtg/v2/mtglib"
statsd "github.com/smira/go-statsd"
)
type statsdFakeLogger struct{}
func (s statsdFakeLogger) Printf(msg string, args ...interface{}) {}
type statsdStreamInfo struct {
createdAt time.Time
clientIP net.IP
}
func (s *statsdStreamInfo) ClientIPTag() statsd.Tag {
if s.clientIP.To4() == nil {
return statsd.StringTag(TagIPType, TagIPTypeIPv6)
} else {
return statsd.StringTag(TagIPType, TagIPTypeIPv4)
}
}
type statsdProcessor struct {
streams map[string]*statsdStreamInfo
client *statsd.Client
}
func (s statsdProcessor) EventStart(evt mtglib.EventStart) {
clientInfo := &statsdStreamInfo{
createdAt: evt.CreatedAt,
clientIP: evt.RemoteIP,
}
s.streams[evt.StreamID()] = clientInfo
s.client.GaugeDelta(MetricActiveConnection, 1, clientInfo.ClientIPTag())
}
func (s statsdProcessor) EventFinish(evt mtglib.EventFinish) {
clientInfo, ok := s.streams[evt.StreamID()]
if !ok {
return
}
defer delete(s.streams, evt.StreamID())
duration := evt.CreatedAt.Sub(clientInfo.createdAt)
s.client.GaugeDelta(MetricActiveConnection, -1, clientInfo.ClientIPTag())
s.client.PrecisionTiming(MetricSessionDuration, duration)
}
func (s statsdProcessor) EventConcurrencyLimited(_ mtglib.EventConcurrencyLimited) {
s.client.Incr(MetricConcurrencyLimited, 1)
}
func (s statsdProcessor) Shutdown() {
now := time.Now()
events := make([]mtglib.EventFinish, 0, len(s.streams))
for k := range s.streams {
events = append(events, mtglib.EventFinish{
CreatedAt: now,
ConnID: k,
})
}
for i := range events {
s.EventFinish(events[i])
}
}
type StatsdFactory struct {
client *statsd.Client
}
func (s StatsdFactory) Close() error {
return s.client.Close()
}
func (s StatsdFactory) Make() events.Observer {
return statsdProcessor{
client: s.client,
streams: make(map[string]*statsdStreamInfo),
}
}
func NewStatsd(address, metricPrefix, tagFormat string) (StatsdFactory, error) {
options := []statsd.Option{
statsd.MetricPrefix(metricPrefix),
statsd.Logger(statsdFakeLogger{}),
}
switch strings.ToLower(tagFormat) {
case "datadog":
options = append(options, statsd.TagStyle(statsd.TagFormatDatadog))
case "influxdb":
options = append(options, statsd.TagStyle(statsd.TagFormatInfluxDB))
case "graphite":
options = append(options, statsd.TagStyle(statsd.TagFormatGraphite))
default:
return StatsdFactory{}, fmt.Errorf("unknown tag format %s", tagFormat)
}
return StatsdFactory{
client: statsd.NewClient(address, options...),
}, nil
}
+126
View File
@@ -0,0 +1,126 @@
package stats_test
import (
"bytes"
"net"
"strings"
"testing"
"time"
"github.com/9seconds/mtg/v2/events"
"github.com/9seconds/mtg/v2/mtglib"
"github.com/9seconds/mtg/v2/stats"
statsd "github.com/smira/go-statsd"
"github.com/stretchr/testify/suite"
)
type statsdFakeServer struct {
conn *net.UDPConn
buf *bytes.Buffer
}
func (s statsdFakeServer) Addr() string {
return s.conn.LocalAddr().String()
}
func (s statsdFakeServer) Close() error {
if s.conn != nil {
return s.conn.Close()
}
return nil
}
func (s statsdFakeServer) String() string {
return strings.TrimSpace(s.buf.String())
}
func statsdNewFakeServer() statsdFakeServer {
conn, err := net.ListenUDP("udp", &net.UDPAddr{
IP: net.ParseIP("127.0.0.1"),
Port: 0,
})
if err != nil {
panic(err)
}
buf := &bytes.Buffer{}
go func() {
currentBuffer := make([]byte, 4096)
for {
n, _, err := conn.ReadFromUDP(currentBuffer)
if n > 0 {
buf.Write(currentBuffer[:n])
}
if err != nil {
return
}
}
}()
return statsdFakeServer{
conn: conn,
buf: buf,
}
}
type StatsdTestSuite struct {
suite.Suite
statsdServer statsdFakeServer
factory stats.StatsdFactory
statsd events.Observer
}
func (suite *StatsdTestSuite) SetupTest() {
suite.statsdServer = statsdNewFakeServer()
factory, err := stats.NewStatsd(suite.statsdServer.Addr(), "mtg.", "datadog")
if err != nil {
panic(err)
}
suite.factory = factory
suite.statsd = suite.factory.Make()
}
func (suite *StatsdTestSuite) TearDownTest() {
suite.statsd.Shutdown()
suite.factory.Close()
suite.statsdServer.Close()
}
func (suite *StatsdTestSuite) TestEventStartFinish() {
suite.statsd.EventStart(mtglib.EventStart{
CreatedAt: time.Now(),
ConnID: "connID",
})
time.Sleep(2 * statsd.DefaultFlushInterval)
suite.Equal("mtg.active_connections:+1|g|#ip_type:ipv4", suite.statsdServer.String())
suite.statsd.EventFinish(mtglib.EventFinish{
CreatedAt: time.Now(),
ConnID: "connID",
})
time.Sleep(2 * statsd.DefaultFlushInterval)
suite.Contains(suite.statsdServer.String(), "mtg.session_duration")
}
func (suite *StatsdTestSuite) TestEventConcurrencyLimited() {
suite.statsd.EventConcurrencyLimited(mtglib.EventConcurrencyLimited{
CreatedAt: time.Now(),
})
time.Sleep(2 * statsd.DefaultFlushInterval)
suite.Equal("mtg.concurrency_limited:1|c", suite.statsdServer.String())
}
func TestStatsd(t *testing.T) {
t.Parallel()
suite.Run(t, &StatsdTestSuite{})
}