Files
weave-scope/vendor/github.com/uber-go/tally/m3/reporter.go
Marc Carré 652cc90f98 Update weaveworks/common to latest version
```
$ gvt delete github.com/weaveworks/common
$ gvt fetch --revision 4d96fd8dcf2c7b417912c6219b310112cb4a4626 github.com/weaveworks/common
2018/07/23 15:31:11 Fetching: github.com/weaveworks/common
2018/07/23 15:31:14 · Skipping (existing): github.com/golang/protobuf/ptypes/any
2018/07/23 15:31:14 · Fetching recursive dependency: github.com/pkg/errors
2018/07/23 15:31:16 · Skipping (existing): github.com/aws/aws-sdk-go/aws
2018/07/23 15:31:16 · Fetching recursive dependency: github.com/sirupsen/logrus
2018/07/23 15:31:18 ·· Skipping (existing): golang.org/x/sys/unix
2018/07/23 15:31:18 ·· Skipping (existing): golang.org/x/crypto/ssh/terminal
2018/07/23 15:31:18 · Skipping (existing): google.golang.org/grpc/status
2018/07/23 15:31:18 · Skipping (existing): github.com/gorilla/mux
2018/07/23 15:31:18 · Fetching recursive dependency: github.com/opentracing-contrib/go-stdlib/nethttp
2018/07/23 15:31:20 ·· Skipping (existing): github.com/opentracing/opentracing-go/ext
2018/07/23 15:31:20 ·· Skipping (existing): github.com/opentracing/opentracing-go/log
2018/07/23 15:31:20 ·· Skipping (existing): github.com/opentracing/opentracing-go
2018/07/23 15:31:20 · Skipping (existing): github.com/prometheus/client_golang/prometheus
2018/07/23 15:31:20 · Skipping (existing): google.golang.org/grpc
2018/07/23 15:31:20 · Skipping (existing): github.com/pmezard/go-difflib/difflib
2018/07/23 15:31:20 · Fetching recursive dependency: github.com/go-kit/kit/log
2018/07/23 15:31:23 ·· Fetching recursive dependency: github.com/go-logfmt/logfmt
2018/07/23 15:31:25 ··· Fetching recursive dependency: github.com/kr/logfmt
2018/07/23 15:31:27 ·· Fetching recursive dependency: github.com/go-stack/stack
2018/07/23 15:31:29 · Fetching recursive dependency: google.golang.org/genproto/googleapis/rpc/status
2018/07/23 15:31:37 ·· Skipping (existing): github.com/golang/protobuf/proto
2018/07/23 15:31:37 ·· Skipping (existing): github.com/golang/protobuf/ptypes/any
2018/07/23 15:31:37 · Skipping (existing): github.com/opentracing/opentracing-go/log
2018/07/23 15:31:37 · Fetching recursive dependency: github.com/sercand/kuberesolver
2018/07/23 15:31:39 ·· Skipping (existing): google.golang.org/grpc/grpclog
2018/07/23 15:31:39 ·· Skipping (existing): google.golang.org/grpc/resolver
2018/07/23 15:31:39 ·· Skipping (existing): golang.org/x/net/context
2018/07/23 15:31:39 · Skipping (existing): google.golang.org/grpc/metadata
2018/07/23 15:31:39 · Skipping (existing): github.com/opentracing/opentracing-go/ext
2018/07/23 15:31:39 · Skipping (existing): github.com/armon/go-socks5
2018/07/23 15:31:39 · Skipping (existing): github.com/opentracing/opentracing-go
2018/07/23 15:31:39 · Skipping (existing): github.com/davecgh/go-spew/spew
2018/07/23 15:31:39 · Skipping (existing): github.com/golang/protobuf/ptypes
2018/07/23 15:31:39 · Skipping (existing): github.com/golang/protobuf/proto
2018/07/23 15:31:39 · Fetching recursive dependency: github.com/grpc-ecosystem/grpc-opentracing/go/otgrpc
2018/07/23 15:31:41 ·· Skipping (existing): github.com/opentracing/opentracing-go/log
2018/07/23 15:31:41 ·· Skipping (existing): golang.org/x/net/context
2018/07/23 15:31:41 ·· Skipping (existing): google.golang.org/grpc/codes
2018/07/23 15:31:41 ·· Skipping (existing): github.com/golang/protobuf/proto
2018/07/23 15:31:41 ·· Skipping (existing): github.com/opentracing/opentracing-go
2018/07/23 15:31:41 ·· Skipping (existing): github.com/opentracing/opentracing-go/ext
2018/07/23 15:31:41 ·· Skipping (existing): google.golang.org/grpc
2018/07/23 15:31:41 ·· Skipping (existing): google.golang.org/grpc/metadata
2018/07/23 15:31:41 ·· Skipping (existing): google.golang.org/grpc/status
2018/07/23 15:31:41 · Fetching recursive dependency: github.com/uber/jaeger-client-go/config
2018/07/23 15:31:44 ·· Fetching recursive dependency: github.com/uber/jaeger-client-go/internal/throttler/remote
2018/07/23 15:31:44 ··· Fetching recursive dependency: github.com/uber/jaeger-client-go/utils
2018/07/23 15:31:44 ···· Fetching recursive dependency: github.com/uber/jaeger-client-go/thrift
2018/07/23 15:31:44 ···· Fetching recursive dependency: github.com/uber/jaeger-client-go/thrift-gen/agent
2018/07/23 15:31:44 ····· Fetching recursive dependency: github.com/uber/jaeger-client-go/thrift-gen/jaeger
2018/07/23 15:31:44 ····· Fetching recursive dependency: github.com/uber/jaeger-client-go/thrift-gen/zipkincore
2018/07/23 15:31:44 ··· Fetching recursive dependency: github.com/uber/jaeger-client-go
2018/07/23 15:31:44 ···· Fetching recursive dependency: github.com/crossdock/crossdock-go
2018/07/23 15:31:46 ····· Skipping (existing): github.com/davecgh/go-spew/spew
2018/07/23 15:31:46 ····· Skipping (existing): golang.org/x/net/context/ctxhttp
2018/07/23 15:31:46 ····· Skipping (existing): golang.org/x/net/context
2018/07/23 15:31:46 ····· Skipping (existing): github.com/pmezard/go-difflib/difflib
2018/07/23 15:31:46 ···· Skipping (existing): github.com/opentracing/opentracing-go/log
2018/07/23 15:31:46 ···· Fetching recursive dependency: go.uber.org/zap/zapcore
2018/07/23 15:31:49 ····· Fetching recursive dependency: go.uber.org/atomic
2018/07/23 15:31:51 ····· Fetching recursive dependency: go.uber.org/zap/internal/bufferpool
2018/07/23 15:31:51 ······ Fetching recursive dependency: go.uber.org/zap/buffer
2018/07/23 15:31:51 ····· Fetching recursive dependency: go.uber.org/multierr
2018/07/23 15:31:54 ····· Fetching recursive dependency: go.uber.org/zap/internal/exit
2018/07/23 15:31:54 ····· Fetching recursive dependency: go.uber.org/zap/internal/color
2018/07/23 15:31:54 ···· Fetching recursive dependency: go.uber.org/zap
2018/07/23 15:31:54 ···· Skipping (existing): github.com/opentracing/opentracing-go
2018/07/23 15:31:54 ···· Skipping (existing): github.com/opentracing/opentracing-go/ext
2018/07/23 15:31:54 ···· Fetching recursive dependency: github.com/uber/jaeger-lib/metrics
2018/07/23 15:31:56 ····· Fetching recursive dependency: github.com/uber-go/tally
2018/07/23 15:31:58 ······ Fetching recursive dependency: github.com/m3db/prometheus_client_golang/prometheus/promhttp
2018/07/23 15:32:00 ······· Skipping (existing): github.com/prometheus/client_golang/prometheus
2018/07/23 15:32:00 ······· Skipping (existing): github.com/prometheus/common/expfmt
2018/07/23 15:32:00 ······· Skipping (existing): github.com/prometheus/client_model/go
2018/07/23 15:32:00 ······ Fetching recursive dependency: gopkg.in/validator.v2
2018/07/23 15:32:06 ······ Fetching recursive dependency: github.com/cactus/go-statsd-client/statsd
2018/07/23 15:32:08 ······ Skipping (existing): gopkg.in/yaml.v2
2018/07/23 15:32:08 ······ Fetching recursive dependency: github.com/m3db/prometheus_client_golang/prometheus
2018/07/23 15:32:08 ······· Skipping (existing): github.com/prometheus/procfs
2018/07/23 15:32:08 ······· Skipping (existing): github.com/prometheus/client_model/go
2018/07/23 15:32:08 ······· Skipping (existing): github.com/prometheus/common/expfmt
2018/07/23 15:32:08 ······· Skipping (existing): golang.org/x/net/context
2018/07/23 15:32:08 ······· Skipping (existing): github.com/beorn7/perks/quantile
2018/07/23 15:32:08 ······· Skipping (existing): github.com/golang/protobuf/proto
2018/07/23 15:32:08 ······· Skipping (existing): github.com/prometheus/common/model
2018/07/23 15:32:08 ······· Skipping (existing): github.com/prometheus/client_golang/prometheus
2018/07/23 15:32:08 ······ Fetching recursive dependency: github.com/apache/thrift/lib/go/thrift
2018/07/23 15:32:13 ····· Skipping (existing): github.com/stretchr/testify/assert
2018/07/23 15:32:13 ····· Fetching recursive dependency: github.com/go-kit/kit/metrics/influx
2018/07/23 15:32:13 ······ Fetching recursive dependency: github.com/influxdata/influxdb/client/v2
2018/07/23 15:32:17 ······· Fetching recursive dependency: github.com/influxdata/influxdb/models
2018/07/23 15:32:17 ········ Fetching recursive dependency: github.com/influxdata/influxdb/pkg/escape
2018/07/23 15:32:17 ······ Fetching recursive dependency: github.com/go-kit/kit/metrics
2018/07/23 15:32:17 ······· Fetching recursive dependency: github.com/performancecopilot/speed
2018/07/23 15:32:19 ······· Fetching recursive dependency: github.com/aws/aws-sdk-go-v2/aws
2018/07/23 15:32:28 ········ Fetching recursive dependency: github.com/aws/aws-sdk-go-v2/internal/sdk
2018/07/23 15:32:28 ········ Skipping (existing): github.com/go-ini/ini
2018/07/23 15:32:28 ········ Fetching recursive dependency: github.com/aws/aws-sdk-go-v2/service/sts
2018/07/23 15:32:28 ········· Fetching recursive dependency: github.com/aws/aws-sdk-go-v2/private/protocol/query
2018/07/23 15:32:28 ·········· Fetching recursive dependency: github.com/aws/aws-sdk-go-v2/private/protocol
2018/07/23 15:32:28 ········· Fetching recursive dependency: github.com/aws/aws-sdk-go-v2/internal/awsutil
2018/07/23 15:32:28 ·········· Skipping (existing): github.com/jmespath/go-jmespath
2018/07/23 15:32:28 ······· Fetching recursive dependency: github.com/aws/aws-sdk-go-v2/service/cloudwatch
2018/07/23 15:32:29 ······· Skipping (existing): github.com/aws/aws-sdk-go/aws
2018/07/23 15:32:29 ······· Skipping (existing): github.com/prometheus/client_golang/prometheus
2018/07/23 15:32:29 ······· Skipping (existing): github.com/aws/aws-sdk-go/service/cloudwatch
2018/07/23 15:32:29 ······· Skipping (existing): github.com/aws/aws-sdk-go/service/cloudwatch/cloudwatchiface
2018/07/23 15:32:29 ······· Fetching recursive dependency: golang.org/x/sync/errgroup
2018/07/23 15:32:31 ········ Skipping (existing): golang.org/x/net/context
2018/07/23 15:32:31 ······· Fetching recursive dependency: github.com/go-kit/kit/util/conn
2018/07/23 15:32:31 ······· Fetching recursive dependency: github.com/VividCortex/gohistogram
2018/07/23 15:32:33 ····· Skipping (existing): github.com/prometheus/client_golang/prometheus
2018/07/23 15:32:33 ····· Fetching recursive dependency: github.com/codahale/hdrhistogram
2018/07/23 15:32:35 ·· Skipping (existing): github.com/opentracing/opentracing-go
2018/07/23 15:32:35 · Fetching recursive dependency: github.com/mwitkow/go-grpc-middleware
2018/07/23 15:32:37 ·· Fetching recursive dependency: github.com/grpc-ecosystem/go-grpc-middleware/logging
2018/07/23 15:32:39 ··· Fetching recursive dependency: github.com/grpc-ecosystem/go-grpc-middleware
2018/07/23 15:32:39 ···· Fetching recursive dependency: github.com/golang/protobuf/jsonpb
2018/07/23 15:32:42 ····· Skipping (existing): github.com/golang/protobuf/ptypes/timestamp
2018/07/23 15:32:42 ····· Skipping (existing): github.com/golang/protobuf/proto
2018/07/23 15:32:42 ····· Skipping (existing): github.com/golang/protobuf/ptypes/duration
2018/07/23 15:32:42 ····· Skipping (existing): github.com/golang/protobuf/ptypes/any
2018/07/23 15:32:42 ····· Skipping (existing): github.com/golang/protobuf/ptypes/struct
2018/07/23 15:32:42 ····· Skipping (existing): github.com/golang/protobuf/ptypes/wrappers
2018/07/23 15:32:42 ···· Skipping (existing): google.golang.org/grpc/metadata
2018/07/23 15:32:42 ···· Fetching recursive dependency: github.com/stretchr/testify/suite
2018/07/23 15:32:45 ····· Skipping (existing): github.com/stretchr/testify/assert
2018/07/23 15:32:45 ····· Fetching recursive dependency: github.com/stretchr/testify/require
2018/07/23 15:32:45 ······ Skipping (existing): github.com/stretchr/testify/assert
2018/07/23 15:32:45 ···· Skipping (existing): google.golang.org/grpc/peer
2018/07/23 15:32:45 ···· Skipping (existing): golang.org/x/net/context
2018/07/23 15:32:45 ···· Skipping (existing): golang.org/x/net/trace
2018/07/23 15:32:45 ···· Fetching recursive dependency: github.com/gogo/protobuf/gogoproto
2018/07/23 15:32:48 ····· Fetching recursive dependency: github.com/gogo/protobuf/protoc-gen-gogo/descriptor
2018/07/23 15:32:48 ······ Skipping (existing): github.com/gogo/protobuf/proto
2018/07/23 15:32:48 ····· Skipping (existing): github.com/gogo/protobuf/proto
2018/07/23 15:32:48 ···· Skipping (existing): google.golang.org/grpc/credentials
2018/07/23 15:32:48 ···· Skipping (existing): google.golang.org/grpc
2018/07/23 15:32:48 ···· Skipping (existing): github.com/opentracing/opentracing-go
2018/07/23 15:32:48 ···· Skipping (existing): google.golang.org/grpc/codes
2018/07/23 15:32:48 ···· Skipping (existing): github.com/golang/protobuf/proto
2018/07/23 15:32:48 ···· Skipping (existing): google.golang.org/grpc/grpclog
2018/07/23 15:32:48 ···· Skipping (existing): github.com/opentracing/opentracing-go/ext
2018/07/23 15:32:48 ···· Skipping (existing): github.com/opentracing/opentracing-go/log
2018/07/23 15:32:48 ··· Skipping (existing): golang.org/x/net/context
2018/07/23 15:32:48 ··· Skipping (existing): google.golang.org/grpc
2018/07/23 15:32:48 ··· Skipping (existing): google.golang.org/grpc/grpclog
2018/07/23 15:32:48 ··· Skipping (existing): google.golang.org/grpc/codes
2018/07/23 15:32:48 ··· Skipping (existing): github.com/golang/protobuf/proto
2018/07/23 15:32:48 ·· Skipping (existing): github.com/opentracing/opentracing-go
2018/07/23 15:32:48 ·· Skipping (existing): google.golang.org/grpc
2018/07/23 15:32:48 ·· Skipping (existing): golang.org/x/net/context
2018/07/23 15:32:48 ·· Skipping (existing): google.golang.org/grpc/codes
2018/07/23 15:32:48 ·· Skipping (existing): google.golang.org/grpc/grpclog
2018/07/23 15:32:48 ·· Skipping (existing): github.com/opentracing/opentracing-go/log
2018/07/23 15:32:48 ·· Skipping (existing): google.golang.org/grpc/metadata
2018/07/23 15:32:48 ·· Skipping (existing): google.golang.org/grpc/peer
2018/07/23 15:32:48 ·· Skipping (existing): google.golang.org/grpc/credentials
2018/07/23 15:32:48 ·· Skipping (existing): github.com/golang/protobuf/proto
2018/07/23 15:32:48 ·· Skipping (existing): golang.org/x/net/trace
2018/07/23 15:32:48 ·· Skipping (existing): github.com/opentracing/opentracing-go/ext
2018/07/23 15:32:48 · Fetching recursive dependency: github.com/weaveworks/promrus
2018/07/23 15:32:53 ·· Skipping (existing): gopkg.in/yaml.v2
2018/07/23 15:32:53 ·· Skipping (existing): golang.org/x/net/context/ctxhttp
2018/07/23 15:32:53 ·· Fetching recursive dependency: github.com/stretchr/objx
2018/07/23 15:32:55 ·· Fetching recursive dependency: gopkg.in/alecthomas/kingpin.v2
2018/07/23 15:32:58 ··· Fetching recursive dependency: github.com/alecthomas/units
2018/07/23 15:33:00 ··· Fetching recursive dependency: github.com/alecthomas/template
2018/07/23 15:33:02 ·· Fetching recursive dependency: github.com/julienschmidt/httprouter
2018/07/23 15:33:05 ·· Skipping (existing): golang.org/x/net/context
2018/07/23 15:33:05 · Skipping (existing): github.com/aws/aws-sdk-go/aws/credentials
2018/07/23 15:33:05 · Skipping (existing): github.com/golang/protobuf/ptypes/empty
2018/07/23 15:33:05 · Skipping (existing): golang.org/x/net/context
2018/07/23 15:33:05 · Skipping (existing): golang.org/x/tools/cover
2018/07/23 15:33:05 · Skipping (existing): github.com/mgutz/ansi
```
2018-07-23 20:10:13 +02:00

621 lines
15 KiB
Go

// Copyright (c) 2018 Uber Technologies, Inc.
//
// Permission is hereby granted, free of charge, to any person obtaining a copy
// of this software and associated documentation files (the "Software"), to deal
// in the Software without restriction, including without limitation the rights
// to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
// copies of the Software, and to permit persons to whom the Software is
// furnished to do so, subject to the following conditions:
//
// The above copyright notice and this permission notice shall be included in
// all copies or substantial portions of the Software.
//
// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
// IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
// FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
// AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
// LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
// THE SOFTWARE.
package m3
import (
"errors"
"fmt"
"io"
"math"
"os"
"strconv"
"sync"
"time"
"github.com/uber-go/tally"
"github.com/uber-go/tally/m3/customtransports"
m3thrift "github.com/uber-go/tally/m3/thrift"
"github.com/uber-go/tally/m3/thriftudp"
"github.com/apache/thrift/lib/go/thrift"
)
// Protocol describes a M3 thrift transport protocol.
type Protocol int
// Compact and Binary represent the compact and
// binary thrift protocols respectively.
const (
Compact Protocol = iota
Binary
)
const (
// ServiceTag is the name of the M3 service tag.
ServiceTag = "service"
// EnvTag is the name of the M3 env tag.
EnvTag = "env"
// HostTag is the name of the M3 host tag.
HostTag = "host"
// DefaultMaxQueueSize is the default M3 reporter queue size.
DefaultMaxQueueSize = 4096
// DefaultMaxPacketSize is the default M3 reporter max packet size.
DefaultMaxPacketSize = int32(1440)
// DefaultHistogramBucketIDName is the default histogram bucket ID tag name
DefaultHistogramBucketIDName = "bucketid"
// DefaultHistogramBucketName is the default histogram bucket name tag name
DefaultHistogramBucketName = "bucket"
// DefaultHistogramBucketTagPrecision is the default
// precision to use when formatting the metric tag
// with the histogram bucket bound values.
DefaultHistogramBucketTagPrecision = uint(6)
emitMetricBatchOverhead = 19
minMetricBucketIDTagLength = 4
)
// Initialize max vars in init function to avoid lint error.
var (
maxInt64 int64
maxFloat64 float64
)
func init() {
maxInt64 = math.MaxInt64
maxFloat64 = math.MaxFloat64
}
type metricType int
const (
counterType metricType = iota + 1
timerType
gaugeType
)
var (
errNoHostPorts = errors.New("at least one entry for HostPorts is required")
errCommonTagSize = errors.New("common tags serialized size exceeds packet size")
)
// Reporter is an M3 reporter.
type Reporter interface {
tally.CachedStatsReporter
io.Closer
}
// reporter is a metrics backend that reports metrics to a local or
// remote M3 collector, metrics are batched together and emitted
// via either thrift compact or binary protocol in batch UDP packets.
type reporter struct {
client *m3thrift.M3Client
curBatch *m3thrift.MetricBatch
curBatchLock sync.Mutex
calc *customtransport.TCalcTransport
calcProto thrift.TProtocol
calcLock sync.Mutex
commonTags map[*m3thrift.MetricTag]bool
freeBytes int32
processors sync.WaitGroup
resourcePool *resourcePool
bucketIDTagName string
bucketTagName string
bucketValFmt string
closeChan chan struct{}
metCh chan sizedMetric
}
// Options is a set of options for the M3 reporter.
type Options struct {
HostPorts []string
Service string
Env string
CommonTags map[string]string
IncludeHost bool
Protocol Protocol
MaxQueueSize int
MaxPacketSizeBytes int32
HistogramBucketIDName string
HistogramBucketName string
HistogramBucketTagPrecision uint
}
// NewReporter creates a new M3 reporter.
func NewReporter(opts Options) (Reporter, error) {
if opts.MaxQueueSize <= 0 {
opts.MaxQueueSize = DefaultMaxQueueSize
}
if opts.MaxPacketSizeBytes <= 0 {
opts.MaxPacketSizeBytes = DefaultMaxPacketSize
}
if opts.HistogramBucketIDName == "" {
opts.HistogramBucketIDName = DefaultHistogramBucketIDName
}
if opts.HistogramBucketName == "" {
opts.HistogramBucketName = DefaultHistogramBucketName
}
if opts.HistogramBucketTagPrecision == 0 {
opts.HistogramBucketTagPrecision = DefaultHistogramBucketTagPrecision
}
// Create M3 thrift client
var trans thrift.TTransport
var err error
if len(opts.HostPorts) == 0 {
err = errNoHostPorts
} else if len(opts.HostPorts) == 1 {
trans, err = thriftudp.NewTUDPClientTransport(opts.HostPorts[0], "")
} else {
trans, err = thriftudp.NewTMultiUDPClientTransport(opts.HostPorts, "")
}
if err != nil {
return nil, err
}
var protocolFactory thrift.TProtocolFactory
if opts.Protocol == Compact {
protocolFactory = thrift.NewTCompactProtocolFactory()
} else {
protocolFactory = thrift.NewTBinaryProtocolFactoryDefault()
}
client := m3thrift.NewM3ClientFactory(trans, protocolFactory)
resourcePool := newResourcePool(protocolFactory)
// Create common tags
tags := resourcePool.getTagList()
for k, v := range opts.CommonTags {
tags[createTag(resourcePool, k, v)] = true
}
if opts.CommonTags[ServiceTag] == "" {
if opts.Service == "" {
return nil, fmt.Errorf("%s common tag is required", ServiceTag)
}
tags[createTag(resourcePool, ServiceTag, opts.Service)] = true
}
if opts.CommonTags[EnvTag] == "" {
if opts.Env == "" {
return nil, fmt.Errorf("%s common tag is required", EnvTag)
}
tags[createTag(resourcePool, EnvTag, opts.Env)] = true
}
if opts.IncludeHost {
if opts.CommonTags[HostTag] == "" {
hostname, err := os.Hostname()
if err != nil {
return nil, fmt.Errorf("error resolving host tag: %v", err)
}
tags[createTag(resourcePool, HostTag, hostname)] = true
}
}
// Calculate size of common tags
batch := resourcePool.getBatch()
batch.CommonTags = tags
batch.Metrics = []*m3thrift.Metric{}
proto := resourcePool.getProto()
batch.Write(proto)
calc := proto.Transport().(*customtransport.TCalcTransport)
numOverheadBytes := emitMetricBatchOverhead + calc.GetCount()
calc.ResetCount()
freeBytes := opts.MaxPacketSizeBytes - numOverheadBytes
if freeBytes <= 0 {
return nil, errCommonTagSize
}
r := &reporter{
client: client,
curBatch: batch,
calc: calc,
calcProto: proto,
commonTags: tags,
freeBytes: freeBytes,
resourcePool: resourcePool,
bucketIDTagName: opts.HistogramBucketIDName,
bucketTagName: opts.HistogramBucketName,
bucketValFmt: "%." + strconv.Itoa(int(opts.HistogramBucketTagPrecision)) + "f",
metCh: make(chan sizedMetric, opts.MaxQueueSize),
}
r.processors.Add(1)
go r.process()
return r, nil
}
// AllocateCounter implements tally.CachedStatsReporter.
func (r *reporter) AllocateCounter(
name string, tags map[string]string,
) tally.CachedCount {
return r.allocateCounter(name, tags)
}
func (r *reporter) allocateCounter(
name string, tags map[string]string,
) cachedMetric {
counter := r.newMetric(name, tags, counterType)
size := r.calculateSize(counter)
return cachedMetric{counter, r, size}
}
// AllocateGauge implements tally.CachedStatsReporter.
func (r *reporter) AllocateGauge(
name string, tags map[string]string,
) tally.CachedGauge {
gauge := r.newMetric(name, tags, gaugeType)
size := r.calculateSize(gauge)
return cachedMetric{gauge, r, size}
}
// AllocateTimer implements tally.CachedStatsReporter.
func (r *reporter) AllocateTimer(
name string, tags map[string]string,
) tally.CachedTimer {
timer := r.newMetric(name, tags, timerType)
size := r.calculateSize(timer)
return cachedMetric{timer, r, size}
}
// AllocateHistogram implements tally.CachedStatsReporter.
func (r *reporter) AllocateHistogram(
name string,
tags map[string]string,
buckets tally.Buckets,
) tally.CachedHistogram {
var (
cachedValueBuckets []cachedHistogramBucket
cachedDurationBuckets []cachedHistogramBucket
)
bucketIDLen := len(strconv.Itoa(buckets.Len()))
bucketIDLen = int(math.Max(float64(bucketIDLen),
float64(minMetricBucketIDTagLength)))
bucketIDLenStr := strconv.Itoa(bucketIDLen)
bucketIDFmt := "%0" + bucketIDLenStr + "d"
for i, pair := range tally.BucketPairs(buckets) {
valueTags, durationTags :=
make(map[string]string), make(map[string]string)
for k, v := range tags {
valueTags[k], durationTags[k] = v, v
}
idTagValue := fmt.Sprintf(bucketIDFmt, i)
valueTags[r.bucketIDTagName] = idTagValue
valueTags[r.bucketTagName] = fmt.Sprintf("%s-%s",
r.valueBucketString(pair.LowerBoundValue()),
r.valueBucketString(pair.UpperBoundValue()))
cachedValueBuckets = append(cachedValueBuckets,
cachedHistogramBucket{pair.UpperBoundValue(),
pair.UpperBoundDuration(),
r.allocateCounter(name, valueTags)})
durationTags[r.bucketIDTagName] = idTagValue
durationTags[r.bucketTagName] = fmt.Sprintf("%s-%s",
r.durationBucketString(pair.LowerBoundDuration()),
r.durationBucketString(pair.UpperBoundDuration()))
cachedDurationBuckets = append(cachedDurationBuckets,
cachedHistogramBucket{pair.UpperBoundValue(),
pair.UpperBoundDuration(),
r.allocateCounter(name, durationTags)})
}
return cachedHistogram{r, name, tags, buckets,
cachedValueBuckets, cachedDurationBuckets}
}
func (r *reporter) valueBucketString(v float64) string {
if v == math.MaxFloat64 {
return "infinity"
}
if v == -math.MaxFloat64 {
return "-infinity"
}
return fmt.Sprintf(r.bucketValFmt, v)
}
func (r *reporter) durationBucketString(d time.Duration) string {
if d == 0 {
return "0"
}
if d == time.Duration(math.MaxInt64) {
return "infinity"
}
if d == time.Duration(math.MinInt64) {
return "-infinity"
}
return d.String()
}
func (r *reporter) newMetric(
name string,
tags map[string]string,
t metricType,
) *m3thrift.Metric {
var (
m = r.resourcePool.getMetric()
metVal = r.resourcePool.getValue()
)
m.Name = name
if tags != nil {
metTags := r.resourcePool.getTagList()
for k, v := range tags {
val := v
metTag := r.resourcePool.getTag()
metTag.TagName = k
metTag.TagValue = &val
metTags[metTag] = true
}
m.Tags = metTags
} else {
m.Tags = nil
}
m.Timestamp = &maxInt64
switch t {
case counterType:
c := r.resourcePool.getCount()
c.I64Value = &maxInt64
metVal.Count = c
case gaugeType:
g := r.resourcePool.getGauge()
g.DValue = &maxFloat64
metVal.Gauge = g
case timerType:
t := r.resourcePool.getTimer()
t.I64Value = &maxInt64
metVal.Timer = t
}
m.MetricValue = metVal
return m
}
func (r *reporter) calculateSize(m *m3thrift.Metric) int32 {
r.calcLock.Lock()
m.Write(r.calcProto)
size := r.calc.GetCount()
r.calc.ResetCount()
r.calcLock.Unlock()
return size
}
func (r *reporter) reportCopyMetric(
m *m3thrift.Metric,
size int32,
t metricType,
iValue int64,
dValue float64,
) {
copy := r.resourcePool.getMetric()
copy.Name = m.Name
copy.Tags = m.Tags
timestampNano := time.Now().UnixNano()
copy.Timestamp = &timestampNano
copy.MetricValue = r.resourcePool.getValue()
switch t {
case counterType:
c := r.resourcePool.getCount()
c.I64Value = &iValue
copy.MetricValue.Count = c
case gaugeType:
g := r.resourcePool.getGauge()
g.DValue = &dValue
copy.MetricValue.Gauge = g
case timerType:
t := r.resourcePool.getTimer()
t.I64Value = &iValue
copy.MetricValue.Timer = t
}
// NB(r): This is to avoid sending on a closed channel,
// it's faster to actually defer/recover than acquire
// a read lock here to ensure we aren't closed: benchmarked
// this with BenchmarkTimer in reporter_benchmark_test.go
defer func() {
recover()
}()
select {
case r.metCh <- sizedMetric{copy, size}:
default:
}
}
// Flush implements tally.CachedStatsReporter.
func (r *reporter) Flush() {
// Avoid send on a closed channel
defer func() {
recover()
}()
r.metCh <- sizedMetric{}
}
// Close waits for metrics to be flushed before closing the backend.
func (r *reporter) Close() (err error) {
defer func() {
if r := recover(); r != nil {
err = fmt.Errorf("close error occurred: %v", r)
}
}()
close(r.metCh)
r.processors.Wait()
return
}
func (r *reporter) Capabilities() tally.Capabilities {
return r
}
func (r *reporter) Reporting() bool {
return true
}
func (r *reporter) Tagging() bool {
return true
}
func (r *reporter) process() {
mets := make([]*m3thrift.Metric, 0, (r.freeBytes / 10))
bytes := int32(0)
for smet := range r.metCh {
if smet.m == nil {
// Explicit flush requested
if len(mets) > 0 {
mets = r.flush(mets)
bytes = 0
}
continue
}
if bytes+smet.size > r.freeBytes {
mets = r.flush(mets)
bytes = 0
}
mets = append(mets, smet.m)
bytes += smet.size
}
if len(mets) > 0 {
// Final flush
r.flush(mets)
}
r.processors.Done()
}
func (r *reporter) flush(
mets []*m3thrift.Metric,
) []*m3thrift.Metric {
r.curBatchLock.Lock()
r.curBatch.Metrics = mets
r.client.EmitMetricBatch(r.curBatch)
r.curBatch.Metrics = nil
r.curBatchLock.Unlock()
r.resourcePool.releaseShallowMetrics(mets)
for i := range mets {
mets[i] = nil
}
return mets[:0]
}
func createTag(
pool *resourcePool,
tagName, tagValue string,
) *m3thrift.MetricTag {
tag := pool.getTag()
tag.TagName = tagName
if tagValue != "" {
tag.TagValue = &tagValue
}
return tag
}
type cachedMetric struct {
metric *m3thrift.Metric
reporter *reporter
size int32
}
func (c cachedMetric) ReportCount(value int64) {
c.reporter.reportCopyMetric(c.metric, c.size, counterType, value, 0)
}
func (c cachedMetric) ReportGauge(value float64) {
c.reporter.reportCopyMetric(c.metric, c.size, gaugeType, 0, value)
}
func (c cachedMetric) ReportTimer(interval time.Duration) {
val := int64(interval)
c.reporter.reportCopyMetric(c.metric, c.size, timerType, val, 0)
}
func (c cachedMetric) ReportSamples(value int64) {
c.reporter.reportCopyMetric(c.metric, c.size, counterType, value, 0)
}
type noopMetric struct {
}
func (c noopMetric) ReportCount(value int64) {
}
func (c noopMetric) ReportGauge(value float64) {
}
func (c noopMetric) ReportTimer(interval time.Duration) {
}
func (c noopMetric) ReportSamples(value int64) {
}
type cachedHistogram struct {
r *reporter
name string
tags map[string]string
buckets tally.Buckets
cachedValueBuckets []cachedHistogramBucket
cachedDurationBuckets []cachedHistogramBucket
}
type cachedHistogramBucket struct {
valueUpperBound float64
durationUpperBound time.Duration
metric cachedMetric
}
func (h cachedHistogram) ValueBucket(
bucketLowerBound, bucketUpperBound float64,
) tally.CachedHistogramBucket {
for _, b := range h.cachedValueBuckets {
if b.valueUpperBound >= bucketUpperBound {
return b.metric
}
}
return noopMetric{}
}
func (h cachedHistogram) DurationBucket(
bucketLowerBound, bucketUpperBound time.Duration,
) tally.CachedHistogramBucket {
for _, b := range h.cachedDurationBuckets {
if b.durationUpperBound >= bucketUpperBound {
return b.metric
}
}
return noopMetric{}
}
type sizedMetric struct {
m *m3thrift.Metric
size int32
}