From bc2bd4510a954557a53d4310538619adc115bb81 Mon Sep 17 00:00:00 2001 From: 9seconds Date: Sat, 27 Mar 2021 22:03:43 +0300 Subject: [PATCH] Rework stats --- events/event_stream.go | 4 +- events/event_stream_test.go | 8 +- events/init.go | 2 +- events/init_test.go | 2 +- events/multi_observer.go | 4 +- events/noop.go | 2 +- events/noop_test.go | 2 +- mtglib/conns.go | 4 +- mtglib/events.go | 28 ++++- mtglib/events_test.go | 23 ++++- mtglib/init.go | 1 + mtglib/proxy.go | 6 ++ stats/init.go | 36 ++++--- stats/pools.go | 20 ++++ stats/prometheus.go | 197 +++++++++++++++--------------------- stats/prometheus_test.go | 20 ++-- stats/statsd.go | 72 +++++-------- stats/statsd_test.go | 19 ++-- stats/stream_info.go | 55 +++++++--- 19 files changed, 273 insertions(+), 232 deletions(-) create mode 100644 stats/pools.go diff --git a/events/event_stream.go b/events/event_stream.go index 1726c81..c8fe758 100644 --- a/events/event_stream.go +++ b/events/event_stream.go @@ -73,8 +73,8 @@ func eventStreamProcessor(ctx context.Context, eventChan <-chan mtglib.Event, ob observer.EventStart(typedEvt) case mtglib.EventConnectedToDC: observer.EventConnectedToDC(typedEvt) - case mtglib.EventTraffic: - observer.EventTraffic(typedEvt) + case mtglib.EventTelegramTraffic: + observer.EventTelegramTraffic(typedEvt) case mtglib.EventFinish: observer.EventFinish(typedEvt) case mtglib.EventIPBlocklisted: diff --git a/events/event_stream_test.go b/events/event_stream_test.go index 90dee03..8daa939 100644 --- a/events/event_stream_test.go +++ b/events/event_stream_test.go @@ -90,8 +90,8 @@ func (suite *EventStreamTestSuite) TestEventConnectedToDC() { time.Sleep(100 * time.Millisecond) } -func (suite *EventStreamTestSuite) TestEventTraffic() { - evt := mtglib.EventTraffic{ +func (suite *EventStreamTestSuite) TestEventTelegramTraffic() { + evt := mtglib.EventTelegramTraffic{ CreatedAt: time.Now(), ConnID: "connID", Traffic: 1024, @@ -100,10 +100,10 @@ func (suite *EventStreamTestSuite) TestEventTraffic() { for _, v := range []*ObserverMock{suite.observerMock1, suite.observerMock2} { v. - On("EventTraffic", mock.Anything). + On("EventTelegramTraffic", mock.Anything). Once(). Run(func(args mock.Arguments) { - caught := args.Get(0).(mtglib.EventTraffic) + caught := args.Get(0).(mtglib.EventTelegramTraffic) suite.Equal(evt.CreatedAt, caught.CreatedAt) suite.Equal(evt.ConnID, caught.ConnID) diff --git a/events/init.go b/events/init.go index 4720a3b..fbfe330 100644 --- a/events/init.go +++ b/events/init.go @@ -6,7 +6,7 @@ type Observer interface { EventStart(mtglib.EventStart) EventFinish(mtglib.EventFinish) EventConnectedToDC(mtglib.EventConnectedToDC) - EventTraffic(mtglib.EventTraffic) + EventTelegramTraffic(mtglib.EventTelegramTraffic) EventConcurrencyLimited(mtglib.EventConcurrencyLimited) EventIPBlocklisted(mtglib.EventIPBlocklisted) diff --git a/events/init_test.go b/events/init_test.go index 58e3635..96089c0 100644 --- a/events/init_test.go +++ b/events/init_test.go @@ -17,7 +17,7 @@ func (o *ObserverMock) EventConnectedToDC(evt mtglib.EventConnectedToDC) { o.Called(evt) } -func (o *ObserverMock) EventTraffic(evt mtglib.EventTraffic) { +func (o *ObserverMock) EventTelegramTraffic(evt mtglib.EventTelegramTraffic) { o.Called(evt) } diff --git a/events/multi_observer.go b/events/multi_observer.go index 1dec489..292ab2c 100644 --- a/events/multi_observer.go +++ b/events/multi_observer.go @@ -40,7 +40,7 @@ func (m multiObserver) EventConnectedToDC(evt mtglib.EventConnectedToDC) { wg.Wait() } -func (m multiObserver) EventTraffic(evt mtglib.EventTraffic) { +func (m multiObserver) EventTelegramTraffic(evt mtglib.EventTelegramTraffic) { wg := &sync.WaitGroup{} wg.Add(len(m.observers)) @@ -48,7 +48,7 @@ func (m multiObserver) EventTraffic(evt mtglib.EventTraffic) { go func(obs Observer) { defer wg.Done() - obs.EventTraffic(evt) + obs.EventTelegramTraffic(evt) }(v) } diff --git a/events/noop.go b/events/noop.go index 8dcbf57..9f37062 100644 --- a/events/noop.go +++ b/events/noop.go @@ -19,7 +19,7 @@ type noopObserver struct{} func (n noopObserver) EventStart(_ mtglib.EventStart) {} func (n noopObserver) EventConnectedToDC(_ mtglib.EventConnectedToDC) {} -func (n noopObserver) EventTraffic(_ mtglib.EventTraffic) {} +func (n noopObserver) EventTelegramTraffic(_ mtglib.EventTelegramTraffic) {} func (n noopObserver) EventFinish(_ mtglib.EventFinish) {} func (n noopObserver) EventConcurrencyLimited(_ mtglib.EventConcurrencyLimited) {} func (n noopObserver) EventIPBlocklisted(_ mtglib.EventIPBlocklisted) {} diff --git a/events/noop_test.go b/events/noop_test.go index 2621e38..6c9b768 100644 --- a/events/noop_test.go +++ b/events/noop_test.go @@ -31,7 +31,7 @@ func (suite *NoopTestSuite) SetupSuite() { RemoteIP: net.ParseIP("127.1.0.1"), DC: 2, }, - "traffic": mtglib.EventTraffic{ + "telegram-traffic": mtglib.EventTelegramTraffic{ CreatedAt: time.Now(), ConnID: "connID", Traffic: 1000, diff --git a/mtglib/conns.go b/mtglib/conns.go index 07b7086..3575957 100644 --- a/mtglib/conns.go +++ b/mtglib/conns.go @@ -21,7 +21,7 @@ func (c connTelegramTraffic) Read(b []byte) (int, error) { n, err := c.Conn.Read(b) if n > 0 { - c.stream.Send(c.ctx, EventTraffic{ + c.stream.Send(c.ctx, EventTelegramTraffic{ CreatedAt: time.Now(), ConnID: c.connID, Traffic: uint(n), @@ -36,7 +36,7 @@ func (c connTelegramTraffic) Write(b []byte) (int, error) { n, err := c.Conn.Write(b) if n > 0 { - c.stream.Send(c.ctx, EventTraffic{ + c.stream.Send(c.ctx, EventTelegramTraffic{ CreatedAt: time.Now(), ConnID: c.connID, Traffic: uint(n), diff --git a/mtglib/events.go b/mtglib/events.go index 0e49e48..959ffea 100644 --- a/mtglib/events.go +++ b/mtglib/events.go @@ -15,6 +15,10 @@ func (e EventStart) StreamID() string { return e.ConnID } +func (e EventStart) Timestamp() time.Time { + return e.CreatedAt +} + type EventConnectedToDC struct { CreatedAt time.Time ConnID string @@ -26,17 +30,25 @@ func (e EventConnectedToDC) StreamID() string { return e.ConnID } -type EventTraffic struct { +func (e EventConnectedToDC) Timestamp() time.Time { + return e.CreatedAt +} + +type EventTelegramTraffic struct { CreatedAt time.Time ConnID string Traffic uint IsRead bool } -func (e EventTraffic) StreamID() string { +func (e EventTelegramTraffic) StreamID() string { return e.ConnID } +func (e EventTelegramTraffic) Timestamp() time.Time { + return e.CreatedAt +} + type EventFinish struct { CreatedAt time.Time ConnID string @@ -46,6 +58,10 @@ func (e EventFinish) StreamID() string { return e.ConnID } +func (e EventFinish) Timestamp() time.Time { + return e.CreatedAt +} + type EventConcurrencyLimited struct { CreatedAt time.Time } @@ -54,6 +70,10 @@ func (e EventConcurrencyLimited) StreamID() string { return "" } +func (e EventConcurrencyLimited) Timestamp() time.Time { + return e.CreatedAt +} + type EventIPBlocklisted struct { CreatedAt time.Time RemoteIP net.IP @@ -62,3 +82,7 @@ type EventIPBlocklisted struct { func (e EventIPBlocklisted) StreamID() string { return "" } + +func (e EventIPBlocklisted) Timestamp() time.Time { + return e.CreatedAt +} diff --git a/mtglib/events_test.go b/mtglib/events_test.go index 133fce3..d0745aa 100644 --- a/mtglib/events_test.go +++ b/mtglib/events_test.go @@ -21,6 +21,7 @@ func (suite *EventsTestSuite) TestEventStart() { } suite.Equal("CONNID", evt.StreamID()) + suite.WithinDuration(time.Now(), evt.Timestamp(), 10*time.Millisecond) } func (suite *EventsTestSuite) TestEventFinish() { @@ -30,6 +31,7 @@ func (suite *EventsTestSuite) TestEventFinish() { } suite.Equal("CONNID", evt.StreamID()) + suite.WithinDuration(time.Now(), evt.Timestamp(), 10*time.Millisecond) } func (suite *EventsTestSuite) TestEventConnectedToDC() { @@ -41,10 +43,11 @@ func (suite *EventsTestSuite) TestEventConnectedToDC() { } suite.Equal("CONNID", evt.StreamID()) + suite.WithinDuration(time.Now(), evt.Timestamp(), 10*time.Millisecond) } -func (suite *EventsTestSuite) TestEventTraffic() { - evt := mtglib.EventTraffic{ +func (suite *EventsTestSuite) TestEventTelegramTraffic() { + evt := mtglib.EventTelegramTraffic{ CreatedAt: time.Now(), ConnID: "CONNID", Traffic: 3, @@ -52,14 +55,26 @@ func (suite *EventsTestSuite) TestEventTraffic() { } suite.Equal("CONNID", evt.StreamID()) + suite.WithinDuration(time.Now(), evt.Timestamp(), 10*time.Millisecond) } func (suite *EventsTestSuite) TestEventConcurrencyLimited() { - suite.Empty(mtglib.EventConcurrencyLimited{}.StreamID()) + evt := mtglib.EventConcurrencyLimited{ + CreatedAt: time.Now(), + } + + suite.Empty(evt.StreamID()) + suite.WithinDuration(time.Now(), evt.Timestamp(), 10*time.Millisecond) } func (suite *EventsTestSuite) TestEventIPBlocklisted() { - suite.Empty(mtglib.EventIPBlocklisted{}.StreamID()) + evt := mtglib.EventIPBlocklisted{ + CreatedAt: time.Now(), + RemoteIP: net.ParseIP("10.0.0.10"), + } + + suite.Empty(evt.StreamID()) + suite.WithinDuration(time.Now(), evt.Timestamp(), 10*time.Millisecond) } func TestEvents(t *testing.T) { diff --git a/mtglib/init.go b/mtglib/init.go index 44da825..0660e9e 100644 --- a/mtglib/init.go +++ b/mtglib/init.go @@ -45,6 +45,7 @@ type IPBlocklist interface { type Event interface { StreamID() string + Timestamp() time.Time } type EventStream interface { diff --git a/mtglib/proxy.go b/mtglib/proxy.go index b6a8f12..50b038c 100644 --- a/mtglib/proxy.go +++ b/mtglib/proxy.go @@ -95,6 +95,9 @@ func (p *Proxy) Serve(listener net.Listener) error { if addr := conn.RemoteAddr().(*net.TCPAddr).IP; p.ipBlocklist.Contains(addr) { conn.Close() + p.logger. + BindStr("ip", conn.RemoteAddr().(*net.TCPAddr).IP.String()). + Info("ip was blacklisted") p.eventStream.Send(p.ctx, EventIPBlocklisted{ CreatedAt: time.Now(), RemoteIP: addr, @@ -110,6 +113,9 @@ func (p *Proxy) Serve(listener net.Listener) error { case errors.Is(err, ants.ErrPoolClosed): return nil case errors.Is(err, ants.ErrPoolOverload): + p.logger. + BindStr("ip", conn.RemoteAddr().(*net.TCPAddr).IP.String()). + Info("connection was concurrency limited") p.eventStream.Send(p.ctx, EventConcurrencyLimited{ CreatedAt: time.Now(), }) diff --git a/stats/init.go b/stats/init.go index 4bf1aa8..ddeea9d 100644 --- a/stats/init.go +++ b/stats/init.go @@ -6,21 +6,27 @@ const ( DefaultStatsdMetricPrefix = DefaultMetricPrefix + "." DefaultStatsdTagFormat = "datadog" - MetricClientConnections = "client_connections" - MetricTelegramConnections = "telegram_connections" - MetricTraffic = "traffic" - MetricSessionDuration = "session_duration" - MetricSessionTraffic = "session_traffic" - MetricConcurrencyLimited = "concurrency_limited" - MetricIPBlocklisted = "ip_blocklisted" + MetricClientConnections = "client_connections" + MetricTelegramConnections = "telegram_connections" + MetricDomainDisguisingConnections = "domain_disguising_connections" - TagIPType = "ip_type" - TagTelegramIP = "ip" - TagDC = "dc" - TagDirection = "direction" + MetricTelegramTraffic = "telegram_traffic" + MetricDomainDisguisingTraffic = "domain_disguising_traffic" - TagIPTypeIPv4 = "ipv4" - TagIPTypeIPv6 = "ipv6" - TagDirectionTelegram = "telegram" - TagDirectionClient = "client" + MetricDomainDisguising = "domain_disguising" + MetricConcurrencyLimited = "concurrency_limited" + MetricIPBlocklisted = "ip_blocklisted" + MetricReplayAttacks = "replay_attacks" + + TagIPFamily = "ip_family" + TagIPFamilyIPv4 = "ipv4" + TagIPFamilyIPv6 = "ipv6" + + TagTelegramIP = "telegram_ip" + + TagDC = "dc" + + TagDirection = "direction" + TagDirectionToClient = "to_client" + TagDirectionFromClient = "from_client" ) diff --git a/stats/pools.go b/stats/pools.go new file mode 100644 index 0000000..aaa6c95 --- /dev/null +++ b/stats/pools.go @@ -0,0 +1,20 @@ +package stats + +import "sync" + +var streamInfoPool = sync.Pool{ + New: func() interface{} { + return &streamInfo{ + tags: map[string]string{}, + } + }, +} + +func acquireStreamInfo() *streamInfo { + return streamInfoPool.Get().(*streamInfo) +} + +func releaseStreamInfo(info *streamInfo) { + info.Reset() + streamInfoPool.Put(info) +} diff --git a/stats/prometheus.go b/stats/prometheus.go index 9069bbf..e002934 100644 --- a/stats/prometheus.go +++ b/stats/prometheus.go @@ -4,8 +4,6 @@ import ( "context" "net" "net/http" - "strconv" - "time" "github.com/9seconds/mtg/v2/events" "github.com/9seconds/mtg/v2/mtglib" @@ -19,89 +17,61 @@ type prometheusProcessor struct { } func (p prometheusProcessor) EventStart(evt mtglib.EventStart) { - sInfo := &streamInfo{ - createdAt: evt.CreatedAt, - clientIP: evt.RemoteIP, - } - p.streams[evt.StreamID()] = sInfo + info := acquireStreamInfo() + info.SetStartTime(evt.CreatedAt) + info.SetClientIP(evt.RemoteIP) + p.streams[evt.StreamID()] = info - p.factory.metricClientConnections.WithLabelValues(sInfo.GetClientIPType()).Inc() + p.factory.metricClientConnections. + WithLabelValues(info.V(TagIPFamily)). + Inc() } func (p prometheusProcessor) EventConnectedToDC(evt mtglib.EventConnectedToDC) { - sInfo, ok := p.streams[evt.StreamID()] + info, ok := p.streams[evt.StreamID()] if !ok { return } - sInfo.remoteIP = evt.RemoteIP - sInfo.dc = evt.DC + info.SetTelegramIP(evt.RemoteIP) + info.SetDC(evt.DC) - p.factory.metricTelegramConnections.WithLabelValues( - sInfo.GetRemoteIPType(), - sInfo.remoteIP.String(), - strconv.Itoa(sInfo.dc)).Inc() + p.factory.metricTelegramConnections. + WithLabelValues(info.V(TagTelegramIP), info.V(TagDC)). + Inc() } -func (p prometheusProcessor) EventTraffic(evt mtglib.EventTraffic) { - sInfo, ok := p.streams[evt.StreamID()] +func (p prometheusProcessor) EventTelegramTraffic(evt mtglib.EventTelegramTraffic) { + info, ok := p.streams[evt.StreamID()] if !ok { return } - labels := []string{ - sInfo.GetRemoteIPType(), - sInfo.remoteIP.String(), - strconv.Itoa(sInfo.dc), - } - - if evt.IsRead { - sInfo.bytesRecvFromTelegram += evt.Traffic - - labels = append(labels, TagDirectionClient) - } else { - sInfo.bytesSentToTelegram += evt.Traffic - - labels = append(labels, TagDirectionTelegram) - } - - p.factory.metricTraffic.WithLabelValues(labels...).Add(float64(evt.Traffic)) + p.factory.metricTelegramTraffic. + WithLabelValues(info.V(TagTelegramIP), info.V(TagDC), getDirection(evt.IsRead)). + Add(float64(evt.Traffic)) } func (p prometheusProcessor) EventFinish(evt mtglib.EventFinish) { - sInfo, ok := p.streams[evt.StreamID()] + info, ok := p.streams[evt.StreamID()] if !ok { return } - defer delete(p.streams, evt.StreamID()) + defer func() { + delete(p.streams, evt.StreamID()) + releaseStreamInfo(info) + }() - duration := evt.CreatedAt.Sub(sInfo.createdAt) + p.factory.metricClientConnections. + WithLabelValues(info.V(TagIPFamily)). + Dec() - p.factory.metricClientConnections.WithLabelValues(sInfo.GetClientIPType()).Dec() - p.factory.metricSessionDuration.Observe(float64(duration) / float64(time.Second)) - - if sInfo.remoteIP == nil { - return + if info.V(TagTelegramIP) != "" { + p.factory.metricTelegramConnections. + WithLabelValues(info.V(TagTelegramIP), info.V(TagDC)). + Dec() } - - labels := []string{ - sInfo.GetRemoteIPType(), - sInfo.remoteIP.String(), - strconv.Itoa(sInfo.dc), - } - - p.factory.metricTelegramConnections.WithLabelValues(labels...).Dec() - - labels = append(labels, TagDirectionClient) - p.factory.metricSessionTraffic. - WithLabelValues(labels...). - Observe(float64(sInfo.bytesRecvFromTelegram)) - - labels[3] = TagDirectionTelegram - p.factory.metricSessionTraffic. - WithLabelValues(labels...). - Observe(float64(sInfo.bytesSentToTelegram)) } func (p prometheusProcessor) EventConcurrencyLimited(_ mtglib.EventConcurrencyLimited) { @@ -109,11 +79,7 @@ func (p prometheusProcessor) EventConcurrencyLimited(_ mtglib.EventConcurrencyLi } func (p prometheusProcessor) EventIPBlocklisted(evt mtglib.EventIPBlocklisted) { - if evt.RemoteIP.To4() == nil { - p.factory.metricIPBlocklisted.WithLabelValues(TagIPTypeIPv6).Inc() - } else { - p.factory.metricIPBlocklisted.WithLabelValues(TagIPTypeIPv4).Inc() - } + p.factory.metricIPBlocklisted.Inc() } func (p prometheusProcessor) Shutdown() { @@ -123,13 +89,17 @@ func (p prometheusProcessor) Shutdown() { type PrometheusFactory struct { httpServer *http.Server - metricClientConnections *prometheus.GaugeVec - metricTelegramConnections *prometheus.GaugeVec - metricTraffic *prometheus.CounterVec - metricIPBlocklisted *prometheus.CounterVec - metricSessionTraffic *prometheus.HistogramVec - metricConcurrencyLimited prometheus.Counter - metricSessionDuration prometheus.Histogram + metricClientConnections *prometheus.GaugeVec + metricTelegramConnections *prometheus.GaugeVec + metricDomainDisguisingConnections *prometheus.GaugeVec + + metricTelegramTraffic *prometheus.CounterVec + metricDomainDisguisingTraffic *prometheus.CounterVec + + metricDomainDisguising prometheus.Counter + metricConcurrencyLimited prometheus.Counter + metricIPBlocklisted prometheus.Counter + metricReplayAttacks prometheus.Counter } func (p *PrometheusFactory) Make() events.Observer { @@ -164,70 +134,63 @@ func NewPrometheus(metricPrefix, httpPath string) *PrometheusFactory { // nolint metricClientConnections: prometheus.NewGaugeVec(prometheus.GaugeOpts{ Namespace: metricPrefix, Name: MetricClientConnections, - Help: "A number of connections under active processing.", - }, []string{TagIPType}), + Help: "A number of actively processing client connections.", + }, []string{TagIPFamily}), metricTelegramConnections: prometheus.NewGaugeVec(prometheus.GaugeOpts{ Namespace: metricPrefix, Name: MetricTelegramConnections, Help: "A number of connections to Telegram servers.", - }, []string{TagIPType, TagTelegramIP, TagDC}), - metricSessionDuration: prometheus.NewHistogram(prometheus.HistogramOpts{ + }, []string{TagTelegramIP, TagDC}), + metricDomainDisguisingConnections: prometheus.NewGaugeVec(prometheus.GaugeOpts{ Namespace: metricPrefix, - Name: MetricSessionDuration, - Help: "Session duration.", - Buckets: []float64{ // per 30 seconds - 30, - 60, - 90, - 120, - 150, - 180, - 210, - 240, - 270, - 300, - }, + Name: MetricDomainDisguisingConnections, + Help: "A number of connections which talk with disguising domain.", + }, []string{TagIPFamily}), + + metricTelegramTraffic: prometheus.NewCounterVec(prometheus.CounterOpts{ + Namespace: metricPrefix, + Name: MetricTelegramTraffic, + Help: "Traffic which is generated talking with Telegram servers.", + }, []string{TagTelegramIP, TagDC, TagDirection}), + metricDomainDisguisingTraffic: prometheus.NewCounterVec(prometheus.CounterOpts{ + Namespace: metricPrefix, + Name: MetricDomainDisguisingTraffic, + Help: "Traffic which is generated talking with disguising domain.", + }, []string{TagDirection}), + + metricDomainDisguising: prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: metricPrefix, + Name: MetricDomainDisguising, + Help: "A number of routings to disguising domain.", }), - metricSessionTraffic: prometheus.NewHistogramVec(prometheus.HistogramOpts{ - Namespace: metricPrefix, - Name: MetricSessionTraffic, - Help: "A traffic size which flew via proxy within a single session.", - Buckets: []float64{ // per 1mb - 1 * 1024 * 1024, - 2 * 1024 * 1024, - 3 * 1024 * 1024, - 4 * 1024 * 1024, - 5 * 1024 * 1024, - 6 * 1024 * 1024, - 7 * 1024 * 1024, - 8 * 1024 * 1024, - 9 * 1024 * 1024, - }, - }, []string{TagIPType, TagTelegramIP, TagDC, TagDirection}), - metricTraffic: prometheus.NewCounterVec(prometheus.CounterOpts{ - Namespace: metricPrefix, - Name: MetricTraffic, - Help: "Traffic which is sent through this proxy.", - }, []string{TagIPType, TagTelegramIP, TagDC, TagDirection}), metricConcurrencyLimited: prometheus.NewCounter(prometheus.CounterOpts{ Namespace: metricPrefix, Name: MetricConcurrencyLimited, Help: "A number of sessions that were rejected by concurrency limiter.", }), - metricIPBlocklisted: prometheus.NewCounterVec(prometheus.CounterOpts{ + metricIPBlocklisted: prometheus.NewCounter(prometheus.CounterOpts{ Namespace: metricPrefix, Name: MetricIPBlocklisted, - Help: "A number of rejected sessions due to ip blocklisting", - }, []string{TagIPType}), + Help: "A number of rejected sessions due to ip blocklisting.", + }), + metricReplayAttacks: prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: metricPrefix, + Name: MetricReplayAttacks, + Help: "A number of detected replay attacks.", + }), } registry.MustRegister(factory.metricClientConnections) registry.MustRegister(factory.metricTelegramConnections) - registry.MustRegister(factory.metricTraffic) - registry.MustRegister(factory.metricSessionTraffic) - registry.MustRegister(factory.metricSessionDuration) + registry.MustRegister(factory.metricDomainDisguisingConnections) + + registry.MustRegister(factory.metricTelegramTraffic) + registry.MustRegister(factory.metricDomainDisguisingTraffic) + + registry.MustRegister(factory.metricDomainDisguising) registry.MustRegister(factory.metricConcurrencyLimited) registry.MustRegister(factory.metricIPBlocklisted) + registry.MustRegister(factory.metricReplayAttacks) return factory } diff --git a/stats/prometheus_test.go b/stats/prometheus_test.go index c1425b5..d3db5a0 100644 --- a/stats/prometheus_test.go +++ b/stats/prometheus_test.go @@ -64,7 +64,7 @@ func (suite *PrometheusTestSuite) TestEventStartFinish() { data, err := suite.Get() suite.NoError(err) - suite.Contains(data, `mtg_client_connections{ip_type="ipv4"} 1`) + suite.Contains(data, `mtg_client_connections{ip_family="ipv4"} 1`) suite.prometheus.EventConnectedToDC(mtglib.EventConnectedToDC{ CreatedAt: time.Now(), @@ -76,9 +76,9 @@ func (suite *PrometheusTestSuite) TestEventStartFinish() { data, err = suite.Get() suite.NoError(err) - suite.Contains(data, `mtg_telegram_connections{dc="4",ip="10.0.0.1",ip_type="ipv4"} 1`) + suite.Contains(data, `mtg_telegram_connections{dc="4",telegram_ip="10.0.0.1"} 1`) - suite.prometheus.EventTraffic(mtglib.EventTraffic{ + suite.prometheus.EventTelegramTraffic(mtglib.EventTelegramTraffic{ CreatedAt: time.Now(), ConnID: "connID", Traffic: 200, @@ -88,9 +88,9 @@ func (suite *PrometheusTestSuite) TestEventStartFinish() { data, err = suite.Get() suite.NoError(err) - suite.Contains(data, `mtg_traffic{dc="4",direction="client",ip="10.0.0.1",ip_type="ipv4"} 200`) + suite.Contains(data, `mtg_telegram_traffic{dc="4",direction="to_client",telegram_ip="10.0.0.1"} 200`) - suite.prometheus.EventTraffic(mtglib.EventTraffic{ + suite.prometheus.EventTelegramTraffic(mtglib.EventTelegramTraffic{ CreatedAt: time.Now(), ConnID: "connID", Traffic: 100, @@ -100,7 +100,7 @@ func (suite *PrometheusTestSuite) TestEventStartFinish() { data, err = suite.Get() suite.NoError(err) - suite.Contains(data, `mtg_traffic{dc="4",direction="telegram",ip="10.0.0.1",ip_type="ipv4"} 100`) + suite.Contains(data, `mtg_telegram_traffic{dc="4",direction="from_client",telegram_ip="10.0.0.1"} 100`) suite.prometheus.EventFinish(mtglib.EventFinish{ CreatedAt: time.Now(), @@ -110,10 +110,8 @@ func (suite *PrometheusTestSuite) TestEventStartFinish() { data, err = suite.Get() suite.NoError(err) - suite.Contains(data, `mtg_client_connections{ip_type="ipv4"} 0`) - suite.Contains(data, `mtg_telegram_connections{dc="4",ip="10.0.0.1",ip_type="ipv4"} 0`) - suite.Contains(data, `mtg_traffic{dc="4",direction="client",ip="10.0.0.1",ip_type="ipv4"} 200`) - suite.Contains(data, `mtg_traffic{dc="4",direction="telegram",ip="10.0.0.1",ip_type="ipv4"} 100`) + suite.Contains(data, `mtg_client_connections{ip_family="ipv4"} 0`) + suite.Contains(data, `mtg_telegram_connections{dc="4",telegram_ip="10.0.0.1"} 0`) } func (suite *PrometheusTestSuite) TestEventConcurrencyLimited() { @@ -138,7 +136,7 @@ func (suite *PrometheusTestSuite) TestEventIPBlocklisted() { data, err := suite.Get() suite.NoError(err) - suite.Contains(data, `mtg_ip_blocklisted{ip_type="ipv6"} 1`) + suite.Contains(data, `mtg_ip_blocklisted 1`) } func TestPrometheus(t *testing.T) { diff --git a/stats/statsd.go b/stats/statsd.go index beea0f8..61e3ced 100644 --- a/stats/statsd.go +++ b/stats/statsd.go @@ -17,74 +17,64 @@ type statsdProcessor struct { } func (s statsdProcessor) EventStart(evt mtglib.EventStart) { - sInfo := &streamInfo{ - createdAt: evt.CreatedAt, - clientIP: evt.RemoteIP, - } - s.streams[evt.StreamID()] = sInfo + info := acquireStreamInfo() + info.SetStartTime(evt.CreatedAt) + info.SetClientIP(evt.RemoteIP) + s.streams[evt.StreamID()] = info s.client.GaugeDelta(MetricClientConnections, 1, - statsd.StringTag(TagIPType, sInfo.GetClientIPType())) + info.TV(TagIPFamily)) } func (s statsdProcessor) EventConnectedToDC(evt mtglib.EventConnectedToDC) { - sInfo, ok := s.streams[evt.StreamID()] + info, ok := s.streams[evt.StreamID()] if !ok { return } - sInfo.remoteIP = evt.RemoteIP - sInfo.dc = evt.DC + info.SetTelegramIP(evt.RemoteIP) + info.SetDC(evt.DC) s.client.GaugeDelta(MetricTelegramConnections, 1, - statsd.StringTag(TagIPType, sInfo.GetRemoteIPType()), - statsd.StringTag(TagTelegramIP, sInfo.remoteIP.String()), - statsd.IntTag(TagDC, sInfo.dc)) + info.TV(TagTelegramIP), + info.TV(TagDC)) } -func (s statsdProcessor) EventTraffic(evt mtglib.EventTraffic) { - sInfo, ok := s.streams[evt.StreamID()] +func (s statsdProcessor) EventTelegramTraffic(evt mtglib.EventTelegramTraffic) { + info, ok := s.streams[evt.StreamID()] if !ok { return } - tags := []statsd.Tag{ - statsd.StringTag(TagIPType, sInfo.GetRemoteIPType()), - statsd.StringTag(TagTelegramIP, sInfo.remoteIP.String()), - statsd.IntTag(TagDC, sInfo.dc), - } - - if evt.IsRead { - tags = append(tags, statsd.StringTag(TagDirection, TagDirectionClient)) - s.client.Incr(MetricTraffic, int64(evt.Traffic), tags...) - } else { - tags = append(tags, statsd.StringTag(TagDirection, TagDirectionTelegram)) - s.client.Incr(MetricTraffic, int64(evt.Traffic), tags...) - } + s.client.Incr(MetricTelegramTraffic, + int64(evt.Traffic), + info.TV(TagTelegramIP), + info.TV(TagDC), + statsd.StringTag(TagDirection, getDirection(evt.IsRead))) } func (s statsdProcessor) EventFinish(evt mtglib.EventFinish) { - sInfo, ok := s.streams[evt.StreamID()] + info, ok := s.streams[evt.StreamID()] if !ok { return } - defer delete(s.streams, evt.StreamID()) + defer func() { + delete(s.streams, evt.StreamID()) + releaseStreamInfo(info) + }() s.client.GaugeDelta(MetricClientConnections, -1, - statsd.StringTag(TagIPType, sInfo.GetClientIPType())) - s.client.PrecisionTiming(MetricSessionDuration, - evt.CreatedAt.Sub(sInfo.createdAt)) + info.TV(TagIPFamily)) - if sInfo.remoteIP != nil { + if info.V(TagTelegramIP) != "" { s.client.GaugeDelta(MetricTelegramConnections, -1, - statsd.StringTag(TagIPType, sInfo.GetRemoteIPType()), - statsd.StringTag(TagTelegramIP, sInfo.remoteIP.String()), - statsd.IntTag(TagDC, sInfo.dc)) + info.TV(TagTelegramIP), + info.TV(TagDC)) } } @@ -93,15 +83,7 @@ func (s statsdProcessor) EventConcurrencyLimited(_ mtglib.EventConcurrencyLimite } func (s statsdProcessor) EventIPBlocklisted(evt mtglib.EventIPBlocklisted) { - var tag statsd.Tag - - if evt.RemoteIP.To4() == nil { - tag = statsd.StringTag(TagIPType, TagIPTypeIPv6) - } else { - tag = statsd.StringTag(TagIPType, TagIPTypeIPv4) - } - - s.client.Incr(MetricIPBlocklisted, 1, tag) + s.client.Incr(MetricIPBlocklisted, 1) } func (s statsdProcessor) Shutdown() { diff --git a/stats/statsd_test.go b/stats/statsd_test.go index 9129755..217ca5d 100644 --- a/stats/statsd_test.go +++ b/stats/statsd_test.go @@ -111,7 +111,7 @@ func (suite *StatsdTestSuite) TestEventStartFinish() { RemoteIP: net.ParseIP("10.0.0.10"), }) time.Sleep(statsdSleepTime) - suite.Equal("mtg.client_connections:+1|g|#ip_type:ipv4", suite.statsdServer.String()) + suite.Equal("mtg.client_connections:+1|g|#ip_family:ipv4", suite.statsdServer.String()) suite.statsd.EventConnectedToDC(mtglib.EventConnectedToDC{ CreatedAt: time.Now(), @@ -121,9 +121,9 @@ func (suite *StatsdTestSuite) TestEventStartFinish() { }) time.Sleep(statsdSleepTime) suite.Contains(suite.statsdServer.String(), - "mtg.telegram_connections:+1|g|#ip_type:ipv4,ip:10.1.0.10,dc:2") + "mtg.telegram_connections:+1|g|#telegram_ip:10.1.0.10,dc:2") - suite.statsd.EventTraffic(mtglib.EventTraffic{ + suite.statsd.EventTelegramTraffic(mtglib.EventTelegramTraffic{ CreatedAt: time.Now(), ConnID: "connID", Traffic: 30, @@ -131,9 +131,9 @@ func (suite *StatsdTestSuite) TestEventStartFinish() { }) time.Sleep(statsdSleepTime) suite.Contains(suite.statsdServer.String(), - "mtg.traffic:30|c|#ip_type:ipv4,ip:10.1.0.10,dc:2,direction:client") + "mtg.telegram_traffic:30|c|#telegram_ip:10.1.0.10,dc:2,direction:to_client") - suite.statsd.EventTraffic(mtglib.EventTraffic{ + suite.statsd.EventTelegramTraffic(mtglib.EventTelegramTraffic{ CreatedAt: time.Now(), ConnID: "connID", Traffic: 90, @@ -141,18 +141,17 @@ func (suite *StatsdTestSuite) TestEventStartFinish() { }) time.Sleep(statsdSleepTime) suite.Contains(suite.statsdServer.String(), - "mtg.traffic:90|c|#ip_type:ipv4,ip:10.1.0.10,dc:2,direction:telegram") + "mtg.telegram_traffic:90|c|#telegram_ip:10.1.0.10,dc:2,direction:from_client") suite.statsd.EventFinish(mtglib.EventFinish{ CreatedAt: time.Now(), ConnID: "connID", }) time.Sleep(statsdSleepTime) - suite.Contains(suite.statsdServer.String(), "mtg.session_duration") suite.Contains(suite.statsdServer.String(), - "mtg.telegram_connections:-1|g|#ip_type:ipv4,ip:10.1.0.10,dc:2") + "mtg.telegram_connections:-1|g|#telegram_ip:10.1.0.10,dc:2") suite.Contains(suite.statsdServer.String(), - "mtg.client_connections:-1|g|#ip_type:ipv4") + "mtg.client_connections:-1|g|#ip_family:ipv4") } func (suite *StatsdTestSuite) TestEventConcurrencyLimited() { @@ -171,7 +170,7 @@ func (suite *StatsdTestSuite) TestEventIPBlocklisted() { }) time.Sleep(statsdSleepTime) - suite.Equal("mtg.ip_blocklisted:1|c|#ip_type:ipv4", suite.statsdServer.String()) + suite.Equal("mtg.ip_blocklisted:1|c", suite.statsdServer.String()) } func TestStatsd(t *testing.T) { diff --git a/stats/stream_info.go b/stats/stream_info.go index 2827792..0aeb904 100644 --- a/stats/stream_info.go +++ b/stats/stream_info.go @@ -2,30 +2,57 @@ package stats import ( "net" + "strconv" "time" + + statsd "github.com/smira/go-statsd" ) type streamInfo struct { - createdAt time.Time - clientIP net.IP - remoteIP net.IP - dc int - bytesSentToTelegram uint - bytesRecvFromTelegram uint + startTime time.Time + tags map[string]string } -func (s *streamInfo) GetClientIPType() string { - return s.getIPType(s.clientIP) +func (s *streamInfo) SetStartTime(tme time.Time) { + s.startTime = tme } -func (s *streamInfo) GetRemoteIPType() string { - return s.getIPType(s.remoteIP) +func (s *streamInfo) SetClientIP(ip net.IP) { + if ip.To4() != nil { + s.tags[TagIPFamily] = TagIPFamilyIPv4 + } else { + s.tags[TagIPFamily] = TagIPFamilyIPv6 + } } -func (s *streamInfo) getIPType(ip net.IP) string { - if ip.To4() == nil { - return TagIPTypeIPv6 +func (s *streamInfo) SetTelegramIP(ip net.IP) { + s.tags[TagTelegramIP] = ip.String() +} + +func (s *streamInfo) SetDC(dc int) { + s.tags[TagDC] = strconv.Itoa(dc) +} + +func (s *streamInfo) V(key string) string { + return s.tags[key] +} + +func (s *streamInfo) TV(key string) statsd.Tag { + return statsd.StringTag(key, s.tags[key]) +} + +func (s *streamInfo) Reset() { + s.startTime = time.Time{} + + for k := range s.tags { + delete(s.tags, k) + } +} + +func getDirection(isRead bool) string { + if isRead { // for telegram + return TagDirectionToClient } - return TagIPTypeIPv4 + return TagDirectionFromClient }