Loading broker/amp.go +4 −3 Original line number Diff line number Diff line Loading @@ -36,8 +36,9 @@ func ampClientOffers(i *IPC, w http.ResponseWriter, r *http.Request) { arg := messages.Arg{ Body: encPollReq, RemoteAddr: "", RendezvousMethod: messages.RendezvousAmpCache, } err = i.ClientOffers(arg, &response, RendezvousAmpCache) err = i.ClientOffers(arg, &response) } else { response, err = (&messages.ClientPollResponse{ Error: "cannot decode URL path", Loading broker/http.go +4 −3 Original line number Diff line number Diff line Loading @@ -169,10 +169,11 @@ func clientOffers(i *IPC, w http.ResponseWriter, r *http.Request) { arg := messages.Arg{ Body: body, RemoteAddr: "", RendezvousMethod: messages.RendezvousHttp, } var response []byte err = i.ClientOffers(arg, &response, RendezvousHttp) err = i.ClientOffers(arg, &response) if err != nil { log.Println(err) w.WriteHeader(http.StatusInternalServerError) Loading broker/ipc.go +8 −7 Original line number Diff line number Diff line Loading @@ -161,7 +161,8 @@ func sendClientResponse(resp *messages.ClientPollResponse, response *[]byte) err } } func (i *IPC) ClientOffers(arg messages.Arg, response *[]byte, rendezvousMethod RendezvousMethod) error { func (i *IPC) ClientOffers(arg messages.Arg, response *[]byte) error { startTime := time.Now() req, err := messages.DecodeClientPollRequest(arg.Body) Loading Loading @@ -195,12 +196,12 @@ func (i *IPC) ClientOffers(arg messages.Arg, response *[]byte, rendezvousMethod snowflake.offerChannel <- offer } else { i.ctx.metrics.lock.Lock() i.ctx.metrics.clientDeniedCount[rendezvousMethod]++ i.ctx.metrics.promMetrics.ClientPollTotal.With(prometheus.Labels{"nat": offer.natType, "status": "denied", "rendezvous_method": string(rendezvousMethod)}).Inc() i.ctx.metrics.clientDeniedCount[arg.RendezvousMethod]++ i.ctx.metrics.promMetrics.ClientPollTotal.With(prometheus.Labels{"nat": offer.natType, "status": "denied", "rendezvous_method": string(arg.RendezvousMethod)}).Inc() if offer.natType == NATUnrestricted { i.ctx.metrics.clientUnrestrictedDeniedCount[rendezvousMethod]++ i.ctx.metrics.clientUnrestrictedDeniedCount[arg.RendezvousMethod]++ } else { i.ctx.metrics.clientRestrictedDeniedCount[rendezvousMethod]++ i.ctx.metrics.clientRestrictedDeniedCount[arg.RendezvousMethod]++ } i.ctx.metrics.lock.Unlock() resp := &messages.ClientPollResponse{Error: messages.StrNoProxies} Loading @@ -211,8 +212,8 @@ func (i *IPC) ClientOffers(arg messages.Arg, response *[]byte, rendezvousMethod select { case answer := <-snowflake.answerChannel: i.ctx.metrics.lock.Lock() i.ctx.metrics.clientProxyMatchCount[rendezvousMethod]++ i.ctx.metrics.promMetrics.ClientPollTotal.With(prometheus.Labels{"nat": offer.natType, "status": "matched", "rendezvous_method": string(rendezvousMethod)}).Inc() i.ctx.metrics.clientProxyMatchCount[arg.RendezvousMethod]++ i.ctx.metrics.promMetrics.ClientPollTotal.With(prometheus.Labels{"nat": offer.natType, "status": "matched", "rendezvous_method": string(arg.RendezvousMethod)}).Inc() i.ctx.metrics.lock.Unlock() resp := &messages.ClientPollResponse{Answer: answer} err = sendClientResponse(resp, response) Loading broker/metrics.go +14 −22 Original line number Diff line number Diff line Loading @@ -36,14 +36,6 @@ type CountryStats struct { counts map[string]int } type RendezvousMethod string const ( RendezvousHttp RendezvousMethod = "http" RendezvousAmpCache RendezvousMethod = "ampcache" RendezvousSqs RendezvousMethod = "sqs" ) // Implements Observable type Metrics struct { logger *log.Logger Loading @@ -52,10 +44,10 @@ type Metrics struct { countryStats CountryStats clientRoundtripEstimate time.Duration proxyIdleCount uint clientDeniedCount map[RendezvousMethod]uint clientRestrictedDeniedCount map[RendezvousMethod]uint clientUnrestrictedDeniedCount map[RendezvousMethod]uint clientProxyMatchCount map[RendezvousMethod]uint clientDeniedCount map[messages.RendezvousMethod]uint clientRestrictedDeniedCount map[messages.RendezvousMethod]uint clientUnrestrictedDeniedCount map[messages.RendezvousMethod]uint clientProxyMatchCount map[messages.RendezvousMethod]uint proxyPollWithRelayURLExtension uint proxyPollWithoutRelayURLExtension uint Loading Loading @@ -160,10 +152,10 @@ func (m *Metrics) LoadGeoipDatabases(geoipDB string, geoip6DB string) error { func NewMetrics(metricsLogger *log.Logger) (*Metrics, error) { m := new(Metrics) m.clientDeniedCount = make(map[RendezvousMethod]uint) m.clientRestrictedDeniedCount = make(map[RendezvousMethod]uint) m.clientUnrestrictedDeniedCount = make(map[RendezvousMethod]uint) m.clientProxyMatchCount = make(map[RendezvousMethod]uint) m.clientDeniedCount = make(map[messages.RendezvousMethod]uint) m.clientRestrictedDeniedCount = make(map[messages.RendezvousMethod]uint) m.clientUnrestrictedDeniedCount = make(map[messages.RendezvousMethod]uint) m.clientProxyMatchCount = make(map[messages.RendezvousMethod]uint) m.countryStats = CountryStats{ counts: make(map[string]int), Loading Loading @@ -215,7 +207,7 @@ func (m *Metrics) printMetrics() { m.logger.Println("client-unrestricted-denied-count", binCount(sumMapValues(&m.clientUnrestrictedDeniedCount))) m.logger.Println("client-snowflake-match-count", binCount(sumMapValues(&m.clientProxyMatchCount))) for _, rendezvousMethod := range [3]RendezvousMethod{RendezvousHttp, RendezvousAmpCache, RendezvousSqs} { for _, rendezvousMethod := range [3]messages.RendezvousMethod{messages.RendezvousHttp, messages.RendezvousAmpCache, messages.RendezvousSqs} { m.logger.Printf("client-%s-denied-count %d\n", rendezvousMethod, binCount(m.clientDeniedCount[rendezvousMethod])) m.logger.Printf("client-%s-restricted-denied-count %d\n", rendezvousMethod, binCount(m.clientRestrictedDeniedCount[rendezvousMethod])) m.logger.Printf("client-%s-unrestricted-denied-count %d\n", rendezvousMethod, binCount(m.clientUnrestrictedDeniedCount[rendezvousMethod])) Loading @@ -231,13 +223,13 @@ func (m *Metrics) printMetrics() { // Restores all metrics to original values func (m *Metrics) zeroMetrics() { m.proxyIdleCount = 0 m.clientDeniedCount = make(map[RendezvousMethod]uint) m.clientRestrictedDeniedCount = make(map[RendezvousMethod]uint) m.clientUnrestrictedDeniedCount = make(map[RendezvousMethod]uint) m.clientDeniedCount = make(map[messages.RendezvousMethod]uint) m.clientRestrictedDeniedCount = make(map[messages.RendezvousMethod]uint) m.clientUnrestrictedDeniedCount = make(map[messages.RendezvousMethod]uint) m.proxyPollRejectedWithRelayURLExtension = 0 m.proxyPollWithRelayURLExtension = 0 m.proxyPollWithoutRelayURLExtension = 0 m.clientProxyMatchCount = make(map[RendezvousMethod]uint) m.clientProxyMatchCount = make(map[messages.RendezvousMethod]uint) m.countryStats.counts = make(map[string]int) for pType := range m.countryStats.proxies { m.countryStats.proxies[pType] = make(map[string]bool) Loading @@ -253,7 +245,7 @@ func binCount(count uint) uint { return uint((math.Ceil(float64(count) / 8)) * 8) } func sumMapValues(m *map[RendezvousMethod]uint) uint { func sumMapValues(m *map[messages.RendezvousMethod]uint) uint { var s uint = 0 for _, v := range *m { s += v Loading broker/sqs.go +4 −3 Original line number Diff line number Diff line Loading @@ -147,8 +147,9 @@ func (r *sqsHandler) handleMessage(context context.Context, message *types.Messa arg := messages.Arg{ Body: encPollReq, RemoteAddr: "", RendezvousMethod: messages.RendezvousSqs, } err = r.IPC.ClientOffers(arg, &response, RendezvousSqs) err = r.IPC.ClientOffers(arg, &response) if err != nil { log.Printf("SQSHandler: error encountered when handling message: %v\n", err) Loading Loading
broker/amp.go +4 −3 Original line number Diff line number Diff line Loading @@ -36,8 +36,9 @@ func ampClientOffers(i *IPC, w http.ResponseWriter, r *http.Request) { arg := messages.Arg{ Body: encPollReq, RemoteAddr: "", RendezvousMethod: messages.RendezvousAmpCache, } err = i.ClientOffers(arg, &response, RendezvousAmpCache) err = i.ClientOffers(arg, &response) } else { response, err = (&messages.ClientPollResponse{ Error: "cannot decode URL path", Loading
broker/http.go +4 −3 Original line number Diff line number Diff line Loading @@ -169,10 +169,11 @@ func clientOffers(i *IPC, w http.ResponseWriter, r *http.Request) { arg := messages.Arg{ Body: body, RemoteAddr: "", RendezvousMethod: messages.RendezvousHttp, } var response []byte err = i.ClientOffers(arg, &response, RendezvousHttp) err = i.ClientOffers(arg, &response) if err != nil { log.Println(err) w.WriteHeader(http.StatusInternalServerError) Loading
broker/ipc.go +8 −7 Original line number Diff line number Diff line Loading @@ -161,7 +161,8 @@ func sendClientResponse(resp *messages.ClientPollResponse, response *[]byte) err } } func (i *IPC) ClientOffers(arg messages.Arg, response *[]byte, rendezvousMethod RendezvousMethod) error { func (i *IPC) ClientOffers(arg messages.Arg, response *[]byte) error { startTime := time.Now() req, err := messages.DecodeClientPollRequest(arg.Body) Loading Loading @@ -195,12 +196,12 @@ func (i *IPC) ClientOffers(arg messages.Arg, response *[]byte, rendezvousMethod snowflake.offerChannel <- offer } else { i.ctx.metrics.lock.Lock() i.ctx.metrics.clientDeniedCount[rendezvousMethod]++ i.ctx.metrics.promMetrics.ClientPollTotal.With(prometheus.Labels{"nat": offer.natType, "status": "denied", "rendezvous_method": string(rendezvousMethod)}).Inc() i.ctx.metrics.clientDeniedCount[arg.RendezvousMethod]++ i.ctx.metrics.promMetrics.ClientPollTotal.With(prometheus.Labels{"nat": offer.natType, "status": "denied", "rendezvous_method": string(arg.RendezvousMethod)}).Inc() if offer.natType == NATUnrestricted { i.ctx.metrics.clientUnrestrictedDeniedCount[rendezvousMethod]++ i.ctx.metrics.clientUnrestrictedDeniedCount[arg.RendezvousMethod]++ } else { i.ctx.metrics.clientRestrictedDeniedCount[rendezvousMethod]++ i.ctx.metrics.clientRestrictedDeniedCount[arg.RendezvousMethod]++ } i.ctx.metrics.lock.Unlock() resp := &messages.ClientPollResponse{Error: messages.StrNoProxies} Loading @@ -211,8 +212,8 @@ func (i *IPC) ClientOffers(arg messages.Arg, response *[]byte, rendezvousMethod select { case answer := <-snowflake.answerChannel: i.ctx.metrics.lock.Lock() i.ctx.metrics.clientProxyMatchCount[rendezvousMethod]++ i.ctx.metrics.promMetrics.ClientPollTotal.With(prometheus.Labels{"nat": offer.natType, "status": "matched", "rendezvous_method": string(rendezvousMethod)}).Inc() i.ctx.metrics.clientProxyMatchCount[arg.RendezvousMethod]++ i.ctx.metrics.promMetrics.ClientPollTotal.With(prometheus.Labels{"nat": offer.natType, "status": "matched", "rendezvous_method": string(arg.RendezvousMethod)}).Inc() i.ctx.metrics.lock.Unlock() resp := &messages.ClientPollResponse{Answer: answer} err = sendClientResponse(resp, response) Loading
broker/metrics.go +14 −22 Original line number Diff line number Diff line Loading @@ -36,14 +36,6 @@ type CountryStats struct { counts map[string]int } type RendezvousMethod string const ( RendezvousHttp RendezvousMethod = "http" RendezvousAmpCache RendezvousMethod = "ampcache" RendezvousSqs RendezvousMethod = "sqs" ) // Implements Observable type Metrics struct { logger *log.Logger Loading @@ -52,10 +44,10 @@ type Metrics struct { countryStats CountryStats clientRoundtripEstimate time.Duration proxyIdleCount uint clientDeniedCount map[RendezvousMethod]uint clientRestrictedDeniedCount map[RendezvousMethod]uint clientUnrestrictedDeniedCount map[RendezvousMethod]uint clientProxyMatchCount map[RendezvousMethod]uint clientDeniedCount map[messages.RendezvousMethod]uint clientRestrictedDeniedCount map[messages.RendezvousMethod]uint clientUnrestrictedDeniedCount map[messages.RendezvousMethod]uint clientProxyMatchCount map[messages.RendezvousMethod]uint proxyPollWithRelayURLExtension uint proxyPollWithoutRelayURLExtension uint Loading Loading @@ -160,10 +152,10 @@ func (m *Metrics) LoadGeoipDatabases(geoipDB string, geoip6DB string) error { func NewMetrics(metricsLogger *log.Logger) (*Metrics, error) { m := new(Metrics) m.clientDeniedCount = make(map[RendezvousMethod]uint) m.clientRestrictedDeniedCount = make(map[RendezvousMethod]uint) m.clientUnrestrictedDeniedCount = make(map[RendezvousMethod]uint) m.clientProxyMatchCount = make(map[RendezvousMethod]uint) m.clientDeniedCount = make(map[messages.RendezvousMethod]uint) m.clientRestrictedDeniedCount = make(map[messages.RendezvousMethod]uint) m.clientUnrestrictedDeniedCount = make(map[messages.RendezvousMethod]uint) m.clientProxyMatchCount = make(map[messages.RendezvousMethod]uint) m.countryStats = CountryStats{ counts: make(map[string]int), Loading Loading @@ -215,7 +207,7 @@ func (m *Metrics) printMetrics() { m.logger.Println("client-unrestricted-denied-count", binCount(sumMapValues(&m.clientUnrestrictedDeniedCount))) m.logger.Println("client-snowflake-match-count", binCount(sumMapValues(&m.clientProxyMatchCount))) for _, rendezvousMethod := range [3]RendezvousMethod{RendezvousHttp, RendezvousAmpCache, RendezvousSqs} { for _, rendezvousMethod := range [3]messages.RendezvousMethod{messages.RendezvousHttp, messages.RendezvousAmpCache, messages.RendezvousSqs} { m.logger.Printf("client-%s-denied-count %d\n", rendezvousMethod, binCount(m.clientDeniedCount[rendezvousMethod])) m.logger.Printf("client-%s-restricted-denied-count %d\n", rendezvousMethod, binCount(m.clientRestrictedDeniedCount[rendezvousMethod])) m.logger.Printf("client-%s-unrestricted-denied-count %d\n", rendezvousMethod, binCount(m.clientUnrestrictedDeniedCount[rendezvousMethod])) Loading @@ -231,13 +223,13 @@ func (m *Metrics) printMetrics() { // Restores all metrics to original values func (m *Metrics) zeroMetrics() { m.proxyIdleCount = 0 m.clientDeniedCount = make(map[RendezvousMethod]uint) m.clientRestrictedDeniedCount = make(map[RendezvousMethod]uint) m.clientUnrestrictedDeniedCount = make(map[RendezvousMethod]uint) m.clientDeniedCount = make(map[messages.RendezvousMethod]uint) m.clientRestrictedDeniedCount = make(map[messages.RendezvousMethod]uint) m.clientUnrestrictedDeniedCount = make(map[messages.RendezvousMethod]uint) m.proxyPollRejectedWithRelayURLExtension = 0 m.proxyPollWithRelayURLExtension = 0 m.proxyPollWithoutRelayURLExtension = 0 m.clientProxyMatchCount = make(map[RendezvousMethod]uint) m.clientProxyMatchCount = make(map[messages.RendezvousMethod]uint) m.countryStats.counts = make(map[string]int) for pType := range m.countryStats.proxies { m.countryStats.proxies[pType] = make(map[string]bool) Loading @@ -253,7 +245,7 @@ func binCount(count uint) uint { return uint((math.Ceil(float64(count) / 8)) * 8) } func sumMapValues(m *map[RendezvousMethod]uint) uint { func sumMapValues(m *map[messages.RendezvousMethod]uint) uint { var s uint = 0 for _, v := range *m { s += v Loading
broker/sqs.go +4 −3 Original line number Diff line number Diff line Loading @@ -147,8 +147,9 @@ func (r *sqsHandler) handleMessage(context context.Context, message *types.Messa arg := messages.Arg{ Body: encPollReq, RemoteAddr: "", RendezvousMethod: messages.RendezvousSqs, } err = r.IPC.ClientOffers(arg, &response, RendezvousSqs) err = r.IPC.ClientOffers(arg, &response) if err != nil { log.Printf("SQSHandler: error encountered when handling message: %v\n", err) Loading