simplifying the diagnostics client. Using circuit breaker to manage the diagnostics http connection.

This commit is contained in:
Michael Baird 2017-10-16 12:08:48 -05:00
parent d91865465c
commit 64fd7d3429
4 changed files with 105 additions and 158 deletions

View file

@ -13,16 +13,14 @@ import (
"sync"
"time"
"github.com/pilosa/pilosa"
"github.com/sony/gobreaker"
)
// TODO: white list of statsd metrics to use
// TODO: unique Cluster ID
// TODO: how should this be disabled, config
// Default interval to sync diagnostics metrics.
const (
DefaultDiagnosticsInterval = 10 * time.Second
DefaultDiagnosticsInterval = 1 * time.Hour
DefaultVersionCheckURL = "https://diagnostics.pilosa.com/v0/version"
)
@ -42,17 +40,20 @@ type Diagnostics struct {
startTime int64
start time.Time
counts map[string]int64
metrics map[string]string
metrics map[string]interface{}
client *http.Client
interval time.Duration
cb *gobreaker.CircuitBreaker
logOutput io.Writer
}
// New returns a pointer to a new Diagnostics Client given an addr in the format "hostname:port".
func New(host string) *Diagnostics {
var st gobreaker.Settings
st.Timeout = DefaultDiagnosticsInterval * 2
return &Diagnostics{
closing: make(chan struct{}),
host: host,
@ -60,17 +61,17 @@ func New(host string) *Diagnostics {
startTime: time.Now().Unix(),
start: time.Now(),
client: http.DefaultClient,
counts: make(map[string]int64),
metrics: make(map[string]string),
metrics: make(map[string]interface{}),
interval: DefaultDiagnosticsInterval,
logOutput: ioutil.Discard,
cb: gobreaker.NewCircuitBreaker(st),
}
}
// SetVersion of locally running Pilosa Cluster to check against master.
func (d *Diagnostics) SetVersion(v string) {
d.version = v
d.Set("Version", v, 1.0)
d.Set("Version", v)
}
// schedule start the diagnostics service ticker.
@ -92,35 +93,28 @@ func (d *Diagnostics) schedule() {
// Flush sends the current metrics.
func (d *Diagnostics) Flush() error {
d.mu.Lock()
d.metrics["uptime"] = strconv.FormatInt((time.Now().Unix() - d.startTime), 10)
buf, _ := d.MarshalJSON()
d.Reset()
d.metrics["uptime"] = (time.Now().Unix() - d.startTime)
buf, _ := d.Encode()
d.mu.Unlock()
// d.logger().Println(string(buf))
req, err := http.NewRequest("POST", d.host, bytes.NewReader(buf))
req.Header.Set("Content-Type", "application/json")
resp, err := d.client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
_, err := d.cb.Execute(func() (interface{}, error) {
req, err := http.NewRequest("POST", d.host, bytes.NewReader(buf))
req.Header.Set("Content-Type", "application/json")
resp, err := d.client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
// TODO verify response
// Read response into buffer.
// body, err := ioutil.ReadAll(resp.Body)
// if err != nil {
// return err
// }
// TODO verify response
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return nil, err
}
return body, nil
})
// TODO circuit breaker
return nil
}
// Reset clears the incremented metrics.
func (d *Diagnostics) Reset() {
d.counts = make(map[string]int64)
return err
}
// Open starts the diagnostics metric go routine.
@ -175,90 +169,18 @@ func (d *Diagnostics) CompareVersion(value string) error {
return nil
}
// MarshalJSON custom marshall string and int maps together.
func (d *Diagnostics) MarshalJSON() ([]byte, error) {
buffer := bytes.NewBufferString("{")
length := len(d.counts)
count := 0
for key, value := range d.counts {
jsonValue, err := json.Marshal(value)
if err != nil {
return nil, err
}
buffer.WriteString(fmt.Sprintf("\"%s\":%s", key, string(jsonValue)))
count++
if count < length {
buffer.WriteString(",")
}
}
if length > 0 {
buffer.WriteString(",")
}
length = len(d.metrics)
count = 0
for key, value := range d.metrics {
jsonValue, err := json.Marshal(value)
if err != nil {
return nil, err
}
buffer.WriteString(fmt.Sprintf("\"%s\":%s", key, string(jsonValue)))
count++
if count < length {
buffer.WriteString(",")
}
}
buffer.WriteString("}")
return buffer.Bytes(), nil
}
// Stats interface implementation.
// Tags no-op.
func (d *Diagnostics) Tags() []string {
return nil
}
// WithTags no-op.
func (d *Diagnostics) WithTags(tags ...string) pilosa.StatsClient {
return d
}
// Count tracks the number of times something occurs per diagnostic period.
func (d *Diagnostics) Count(name string, value int64, rate float64) {
d.mu.Lock()
defer d.mu.Unlock()
d.counts[name] += value
}
// CountWithCustomTags Tracks the number of times something occurs per diagnostic period.
func (d *Diagnostics) CountWithCustomTags(name string, value int64, rate float64, tags []string) {
d.mu.Lock()
defer d.mu.Unlock()
d.counts[name] += value
}
// Gauge records the value of a metric.
func (d *Diagnostics) Gauge(name string, value float64, rate float64) {
d.Set(name, strconv.FormatFloat(value, 'f', -1, 64), rate)
}
// Histogram is a no-op.
func (d *Diagnostics) Histogram(name string, value float64, rate float64) {
// Encode metrics maps into the json message format
func (d *Diagnostics) Encode() ([]byte, error) {
return json.Marshal(d.metrics)
}
// Set adds a key value metric.
func (d *Diagnostics) Set(name string, value string, rate float64) {
func (d *Diagnostics) Set(name string, value interface{}) {
d.mu.Lock()
defer d.mu.Unlock()
d.metrics[name] = value
}
// Timing no-op.
func (d *Diagnostics) Timing(name string, value time.Duration, rate float64) {
}
// SetLogger Set the logger output type.
func (d *Diagnostics) SetLogger(logger io.Writer) {
d.logOutput = logger

View file

@ -9,7 +9,6 @@ import (
"runtime"
"strings"
"testing"
"time"
"github.com/pilosa/pilosa/diagnostics"
)
@ -24,31 +23,17 @@ func TestDiagnosticsClient(t *testing.T) {
d.SetLogger(ioutil.Discard)
defer d.Close()
dur, _ := time.ParseDuration("123us")
d.CountWithCustomTags("ct", 1, 1.0, []string{"foo:bar"})
d.Count("cc", 1, 1.0)
d.Gauge("gg", 10, 1.0)
d.Histogram("hh", 1, 1.0)
d.Timing("tt", dur, 1.0)
d.Set("ss", "ss", 1.0)
d.Set("gg", 10)
d.Set("ss", "ss")
d1 := d.WithTags("test")
if !reflect.DeepEqual(d, d1) {
t.Fatalf("Diagnostics is a singleton")
}
if s := d.Tags(); s != nil {
t.Fatalf("No Diagnostics Tags")
}
data, err := d.MarshalJSON()
data, err := d.Encode()
if err != nil {
t.Fatal(err)
}
// Test the recorded metrics, note that some types are skipped.
var eq bool
output1 := []byte(`{"ct":1,"cc":1,"gg":"10","ss":"ss"}`)
output1 := []byte(`{"gg":10,"ss":"ss"}`)
if eq, err = compareJSON(data, output1); err != nil {
t.Fatal(err)
}
@ -58,12 +43,12 @@ func TestDiagnosticsClient(t *testing.T) {
// Test the metrics after a flush.
d.Flush()
data, err = d.MarshalJSON()
data, err = d.Encode()
if err != nil {
t.Fatal(err)
}
output2 := []byte(`{"gg":"10","ss":"ss","uptime":"0"}`)
output2 := []byte(`{"gg":10,"ss":"ss","uptime":0}`)
if eq, err = compareJSON(data, output2); err != nil {
t.Fatal(err)
}
@ -160,8 +145,8 @@ func BenchmarkDiagnostics(b *testing.B) {
b.RunParallel(func(pb *testing.PB) {
for pb.Next() {
d.Count("cc", 1, 1.0)
d.Gauge("gg", 10, 1.0)
d.Set("cc", 1)
d.Set("gg", "test")
}
})
}

View file

@ -34,6 +34,7 @@ import (
"github.com/CAFxX/gcnotifier"
"github.com/gogo/protobuf/proto"
"github.com/pilosa/pilosa/diagnostics"
"github.com/pilosa/pilosa/internal"
)
@ -41,6 +42,7 @@ import (
const (
DefaultAntiEntropyInterval = 10 * time.Minute
DefaultPollingInterval = 60 * time.Second
DefaultDiagnosticServer = "https://diagnostics.pilosa.com/v0/diagnostics"
)
// Server represents a holder wrapped by a running HTTP server.
@ -59,14 +61,16 @@ type Server struct {
// Cluster configuration.
// Host is replaced with actual host after opening if port is ":0".
Network string
URI *URI
Cluster *Cluster
Network string
URI *URI
Cluster *Cluster
diagnostics *diagnostics.Diagnostics
// Background monitoring intervals.
AntiEntropyInterval time.Duration
PollingInterval time.Duration
MetricInterval time.Duration
DiagnosticInterval time.Duration
// TLS configuration
TLS *tls.Config
@ -88,12 +92,14 @@ func NewServer() *Server {
Handler: NewHandler(),
Broadcaster: NopBroadcaster,
BroadcastReceiver: NopBroadcastReceiver,
diagnostics: diagnostics.New(DefaultDiagnosticServer),
Network: "tcp",
AntiEntropyInterval: DefaultAntiEntropyInterval,
PollingInterval: DefaultPollingInterval,
MetricInterval: 0,
DiagnosticInterval: diagnostics.DefaultDiagnosticsInterval
LogOutput: os.Stderr,
}
@ -191,10 +197,11 @@ func (s *Server) Open() error {
}()
// Start background monitoring.
s.wg.Add(3)
s.wg.Add(4)
go func() { defer s.wg.Done(); s.monitorAntiEntropy() }()
go func() { defer s.wg.Done(); s.monitorMaxSlices() }()
go func() { defer s.wg.Done(); s.monitorRuntime() }()
go func() { defer s.wg.Done(); s.monitorDiagnostics() }()
return nil
}
@ -507,14 +514,57 @@ func (s *Server) checkMaxSlices(scheme string, hostPort string) (map[string]uint
return pb.MaxSlices, nil
}
// monitorDiagnostics periodically polls the the Pilosa Indexes for cluster info.
func (s *Server) monitorDiagnostics() {
if s.DiagnosticInterval <= 0 {
return
}
s.diagnostics.SetLogger(s.LogOutput)
s.diagnostics.SetVersion(Version)
s.diagnostics.Set("Host", s.Host)
s.diagnostics.Set("Cluster", strings.Join(s.Cluster.NodeSetHosts(), ","))
s.diagnostics.Set("NumNodes", len(s.Cluster.Nodes))
s.diagnostics.Set("NumCPU", runtime.NumCPU())
// TODO: unique cluster ID
ticker := time.NewTicker(s.DiagnosticInterval)
defer ticker.Stop()
for {
// Wait for tick or a close.
select {
case <-s.closing:
return
case <-ticker.C:
numFrames := 0
numSlices := uint64(0)
for _, index := range s.Holder.Indexes() {
numSlices += index.MaxSlice() + 1
for _, f := range index.Frames() {
numFrames++
if f.rangeEnabled {
s.diagnostics.Set("BSIEnabled", true)
}
if f.timeQuantum != "" {
s.diagnostics.Set("TimeQuantumEnabled", true)
}
}
}
s.diagnostics.Set("NumIndexes", len(s.Holder.Indexes()))
s.diagnostics.Set("NumFrames", numFrames)
s.diagnostics.Set("NumSlices", numSlices)
s.diagnostics.Set("OpenFiles", CountOpenFiles())
s.diagnostics.Set("GoRoutines", runtime.NumGoroutine())
s.diagnostics.CheckVersion()
s.diagnostics.Flush()
}
}
}
// monitorRuntime periodically polls the Go runtime metrics.
func (s *Server) monitorRuntime() {
s.Holder.Stats.Set("Host", s.Host, 1.0)
s.Holder.Stats.Set("Cluster", strings.Join(s.Cluster.NodeSetHosts(), ","), 1.0)
s.Holder.Stats.Set("NumNodes", strconv.Itoa(len(s.Cluster.Nodes)), 1.0)
s.Holder.Stats.Set("NumCPU", strconv.Itoa(runtime.NumCPU()), 1.0)
// TODO should we force this to run for diagnostics?
// Disable metrics when poll interval is zero.
if s.MetricInterval <= 0 {
return

View file

@ -31,7 +31,6 @@ import (
"crypto/tls"
"github.com/pilosa/pilosa"
"github.com/pilosa/pilosa/diagnostics"
"github.com/pilosa/pilosa/gossip"
"github.com/pilosa/pilosa/statsd"
)
@ -42,8 +41,7 @@ func init() {
const (
// DefaultDataDir is the default data directory.
DefaultDataDir = "~/.pilosa"
DefaultDiagnosticServer = "https://requestb.in/w3uukzw3"
DefaultDataDir = "~/.pilosa"
)
// Command represents the state of the pilosa server command.
@ -247,20 +245,12 @@ func (m *Command) Close() error {
// NewStatsClient creates a stats client from the config
func NewStatsClient(name string, host string) (pilosa.StatsClient, error) {
ms := make(pilosa.MultiStatsClient, 1)
d := diagnostics.New(DefaultDiagnosticServer)
d.SetVersion(pilosa.Version)
ms[0] = d
switch name {
case "expvar":
ms = append(ms, pilosa.NewExpvarStatsClient())
return pilosa.NewExpvarStatsClient(), nil
case "statsd":
r, err := statsd.NewStatsClient(host)
if err != nil {
return nil, err
}
ms = append(ms, r)
return statsd.NewStatsClient(host)
default:
return pilosa.NopStatsClient, nil
}
return ms, nil
}