From 8e7207d975829b6a8f7aa1f7cdb7415b0ebbfe0e Mon Sep 17 00:00:00 2001 From: 9seconds Date: Wed, 17 Mar 2021 11:07:16 +0300 Subject: [PATCH] Add statsd --- go.mod | 1 + go.sum | 2 + stats/init.go | 12 +++++ stats/statsd.go | 115 +++++++++++++++++++++++++++++++++++++++ stats/statsd_test.go | 126 +++++++++++++++++++++++++++++++++++++++++++ 5 files changed, 256 insertions(+) create mode 100644 stats/init.go create mode 100644 stats/statsd.go create mode 100644 stats/statsd_test.go diff --git a/go.mod b/go.mod index f9caefe..281af91 100644 --- a/go.mod +++ b/go.mod @@ -16,6 +16,7 @@ require ( github.com/panjf2000/ants v1.3.0 // indirect github.com/pelletier/go-toml v1.8.1 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/testify v1.7.0 github.com/tylertreat/BoomFilters v0.0.0-20200520150052-42a7b4300c0c // indirect diff --git a/go.sum b/go.sum index bcb02ae..3aee538 100644 --- a/go.sum +++ b/go.sum @@ -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/zerolog v1.20.0 h1:38k9hgtUBdxFwE34yS8rTHmHBa4eN16E4DJlv177LNs= 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.3.0 h1:NGXK3lHquSN08v5vWalVI/L8XU9hdzE/G6xsrze47As= github.com/stretchr/objx v0.3.0/go.mod h1:qt09Ya8vawLte6SNmTgCsAVtYtaKzEcn8ATUoHMkEqE= diff --git a/stats/init.go b/stats/init.go new file mode 100644 index 0000000..7210888 --- /dev/null +++ b/stats/init.go @@ -0,0 +1,12 @@ +package stats + +const ( + MetricActiveConnection = "active_connections" + MetricSessionDuration = "session_duration" + MetricConcurrencyLimited = "concurrency_limited" + + TagIPType = "ip_type" + + TagIPTypeIPv4 = "ipv4" + TagIPTypeIPv6 = "ipv4" +) diff --git a/stats/statsd.go b/stats/statsd.go new file mode 100644 index 0000000..92e0979 --- /dev/null +++ b/stats/statsd.go @@ -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 +} diff --git a/stats/statsd_test.go b/stats/statsd_test.go new file mode 100644 index 0000000..13b6f02 --- /dev/null +++ b/stats/statsd_test.go @@ -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{}) +}