From c66e30042550e06ab43e87088da16abcff9a1f85 Mon Sep 17 00:00:00 2001 From: 9seconds Date: Mon, 9 Jul 2018 11:29:47 +0300 Subject: [PATCH] Stats management utilities --- proxy/proxy.go | 2 +- stats/channels.go | 147 ++++++++++++++++++++++++++++++++++++++++++++++ stats/server.go | 11 ++-- stats/stats.go | 31 ++++++++-- 4 files changed, 182 insertions(+), 9 deletions(-) create mode 100644 stats/channels.go diff --git a/proxy/proxy.go b/proxy/proxy.go index a8399eb..e0f3a4e 100644 --- a/proxy/proxy.go +++ b/proxy/proxy.go @@ -39,7 +39,7 @@ func (p *Proxy) Serve() error { func (p *Proxy) accept(conn net.Conn) { connID := uuid.NewV4().String() - log := zap.S().With("connection_id", connID) + log := zap.S().With("connection_id", connID).Named("main") defer func() { conn.Close() diff --git a/stats/channels.go b/stats/channels.go new file mode 100644 index 0000000..775bd2e --- /dev/null +++ b/stats/channels.go @@ -0,0 +1,147 @@ +package stats + +import ( + "net" + "sync/atomic" + "time" + + "github.com/9seconds/mtg/mtproto" +) + +const ( + crashesChanLength = 1 + connectionsChanLength = 20 + trafficChanLength = 5000 +) + +var ( + CrashesChan = make(chan struct{}, crashesChanLength) + ConnectionsChan = make(chan *connectionData, connectionsChanLength) + TrafficChan = make(chan *trafficData, trafficChanLength) +) + +type connectionData struct { + connectionType mtproto.ConnectionType + addr *net.TCPAddr + connected bool +} + +type trafficData struct { + traffic int + ingress bool +} + +func crashManager() { + for range CrashesChan { + instance.mutex.RLock() + + instance.Crashes++ + + instance.mutex.RUnlock() + } +} + +func connectionManager() { + for event := range ConnectionsChan { + instance.mutex.RLock() + + isIPv4 := event.addr.IP.To4() == nil + var inc uint32 = 1 + if !event.connected { + inc = ^uint32(0) + } + + switch event.connectionType { + case mtproto.ConnectionTypeAbridged: + if isIPv4 { + atomic.AddUint32(&instance.ActiveConnections.Abridged.IPv4, inc) + if event.connected { + atomic.AddUint32(&instance.AllConnections.Abridged.IPv4, inc) + } + } else { + atomic.AddUint32(&instance.ActiveConnections.Abridged.IPv6, inc) + if event.connected { + atomic.AddUint32(&instance.AllConnections.Abridged.IPv6, inc) + } + } + default: + if isIPv4 { + atomic.AddUint32(&instance.ActiveConnections.Intermediate.IPv4, inc) + if event.connected { + atomic.AddUint32(&instance.AllConnections.Intermediate.IPv4, inc) + } + } else { + atomic.AddUint32(&instance.ActiveConnections.Intermediate.IPv6, inc) + if event.connected { + atomic.AddUint32(&instance.AllConnections.Intermediate.IPv6, inc) + } + } + } + + instance.mutex.RUnlock() + } +} + +func trafficManager() { + speedChan := time.Tick(time.Second) + + for { + select { + case event := <-TrafficChan: + instance.mutex.RLock() + + if event.ingress { + instance.Traffic.Ingress += trafficValue(event.traffic) + instance.speedCurrent.Ingress += trafficSpeedValue(event.traffic) + } else { + instance.Traffic.Egress += trafficValue(event.traffic) + instance.speedCurrent.Egress += trafficSpeedValue(event.traffic) + } + + instance.mutex.RUnlock() + case <-speedChan: + instance.mutex.RLock() + + instance.Speed.Ingress = instance.speedCurrent.Ingress + instance.Speed.Egress = instance.speedCurrent.Egress + instance.speedCurrent.Ingress = trafficSpeedValue(0) + instance.speedCurrent.Egress = trafficSpeedValue(0) + + instance.mutex.RUnlock() + } + } +} + +func NewCrash() { + CrashesChan <- struct{}{} +} + +func ClientConnected(connectionType mtproto.ConnectionType, addr *net.TCPAddr) { + ConnectionsChan <- &connectionData{ + connectionType: connectionType, + addr: addr, + connected: true, + } +} + +func ClientDisconnected(connectionType mtproto.ConnectionType, addr *net.TCPAddr) { + ConnectionsChan <- &connectionData{ + connectionType: connectionType, + addr: addr, + connected: false, + } +} + +func IngressTraffic(traffic int) { + TrafficChan <- &trafficData{ + traffic: traffic, + ingress: true, + } +} + +func EgressTraffic(traffic int) { + TrafficChan <- &trafficData{ + traffic: traffic, + ingress: false, + } +} diff --git a/stats/server.go b/stats/server.go index c3935d2..3fb6f99 100644 --- a/stats/server.go +++ b/stats/server.go @@ -13,12 +13,15 @@ var instance *stats func Start(conf *config.Config) { instance = &stats{ - URLs: conf.GetURLs(), - Uptime: uptime(time.Now()), - speedCurrent: &speed{}, - mutex: &sync.RWMutex{}, + URLs: conf.GetURLs(), + Uptime: uptime(time.Now()), + mutex: &sync.RWMutex{}, } + go crashManager() + go connectionManager() + go trafficManager() + http.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "application/json") diff --git a/stats/stats.go b/stats/stats.go index 4754674..5c6099c 100644 --- a/stats/stats.go +++ b/stats/stats.go @@ -49,9 +49,31 @@ func (t trafficSpeedValue) MarshalJSON() ([]byte, error) { } type connections struct { - All uint32 `json:"all"` - Abridged uint32 `json:"abridged"` - Intermediate uint32 `json:"intermediate"` + All connectionType `json:"all"` + Abridged connectionType `json:"abridged"` + Intermediate connectionType `json:"intermediate"` +} + +func (c connections) MarshalJSON() ([]byte, error) { + c.All.IPv4 = c.Abridged.IPv4 + c.Intermediate.IPv4 + c.All.IPv6 = c.Abridged.IPv6 + c.Intermediate.IPv6 + + value := struct { + All connectionType `json:"all"` + Abridged connectionType `json:"abridged"` + Intermediate connectionType `json:"intermediate"` + }{ + All: c.All, + Abridged: c.Abridged, + Intermediate: c.Intermediate, + } + + return json.Marshal(value) +} + +type connectionType struct { + IPv6 uint32 `json:"ipv6"` + IPv4 uint32 `json:"ipv4"` } type traffic struct { @@ -71,7 +93,8 @@ type stats struct { Traffic traffic `json:"traffic"` Speed speed `json:"speed"` Uptime uptime `json:"uptime"` + Crashes uint32 `json:"crashes"` - speedCurrent *speed + speedCurrent speed mutex *sync.RWMutex }