mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-07 19:37:51 +00:00
commit
17da0443d3
13 changed files with 545 additions and 17 deletions
8
Gopkg.lock
generated
8
Gopkg.lock
generated
|
|
@ -161,6 +161,12 @@
|
|||
packages = ["."]
|
||||
revision = "e2103e2c35297fb7e17febb81e49b312087a2372"
|
||||
|
||||
[[projects]]
|
||||
name = "github.com/sony/gobreaker"
|
||||
packages = ["."]
|
||||
revision = "e9556a45379ef1da12e54847edb2fb3d7d566f36"
|
||||
version = "0.3.0"
|
||||
|
||||
[[projects]]
|
||||
branch = "master"
|
||||
name = "github.com/spf13/afero"
|
||||
|
|
@ -230,6 +236,6 @@
|
|||
[solve-meta]
|
||||
analyzer-name = "dep"
|
||||
analyzer-version = 1
|
||||
inputs-digest = "e8e78a7c61547d8f4d967c8deac9334b3151e7c10c4f7909e80758956e9c8204"
|
||||
inputs-digest = "0a7eaef315413bc55acf42c5566982fb71be8c11f33162127953572738967b35"
|
||||
solver-name = "gps-cdcl"
|
||||
solver-version = 1
|
||||
|
|
|
|||
|
|
@ -99,6 +99,7 @@ type Config struct {
|
|||
Service string `toml:"service"`
|
||||
Host string `toml:"host"`
|
||||
PollInterval Duration `toml:"poll-interval"`
|
||||
Diagnostics bool `toml:"diagnostics"`
|
||||
} `toml:"metric"`
|
||||
|
||||
TLS TLSConfig
|
||||
|
|
@ -116,6 +117,7 @@ func NewConfig() *Config {
|
|||
c.Cluster.Hosts = []string{}
|
||||
c.AntiEntropy.Interval = Duration(DefaultAntiEntropyInterval)
|
||||
c.Metric.Service = DefaultMetrics
|
||||
c.Metric.Diagnostics = true
|
||||
c.TLS = TLSConfig{}
|
||||
return c
|
||||
}
|
||||
|
|
|
|||
|
|
@ -44,6 +44,7 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) {
|
|||
flags.StringVarP(&srv.Config.Cluster.Type, "cluster.type", "", "gossip", "Determine how the cluster handles membership and state sharing. Choose from [static, gossip]")
|
||||
flags.StringVarP(&srv.Config.Metric.Service, "metric.service", "", "nop", "Default URI on which pilosa should listen.")
|
||||
flags.StringVarP(&srv.Config.Metric.Host, "metric.host", "", "", "Default URI to send metrics.")
|
||||
flags.BoolVarP((&srv.Config.Metric.Diagnostics), "metric.diagnostics", "", true, "Enabled diagnostics reporting.")
|
||||
flags.DurationVarP((*time.Duration)(&srv.Config.Metric.PollInterval), "metric.poll-interval", "", time.Minute*0, "Polling interval metrics.")
|
||||
SetTLSConfig(flags, &srv.Config.TLS.CertificatePath, &srv.Config.TLS.CertificateKeyPath, &srv.Config.TLS.SkipVerify)
|
||||
}
|
||||
|
|
|
|||
216
diagnostics/diagnostics.go
Normal file
216
diagnostics/diagnostics.go
Normal file
|
|
@ -0,0 +1,216 @@
|
|||
package diagnostics
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"log"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/sony/gobreaker"
|
||||
)
|
||||
|
||||
// TODO: unique Cluster ID
|
||||
|
||||
// Default version check URL.
|
||||
const (
|
||||
DefaultVersionCheckURL = "https://diagnostics.pilosa.com/v0/version"
|
||||
)
|
||||
|
||||
type versionResponse struct {
|
||||
Version string `json:"version"`
|
||||
Message string `json:"message"`
|
||||
}
|
||||
|
||||
// Diagnostics represents a client to the Pilosa cluster.
|
||||
type Diagnostics struct {
|
||||
mu sync.Mutex
|
||||
wg sync.WaitGroup
|
||||
closing chan struct{}
|
||||
host string
|
||||
VersionURL string
|
||||
version string
|
||||
lastVersion string
|
||||
startTime int64
|
||||
start time.Time
|
||||
|
||||
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 {
|
||||
|
||||
return &Diagnostics{
|
||||
closing: make(chan struct{}),
|
||||
host: host,
|
||||
VersionURL: DefaultVersionCheckURL,
|
||||
startTime: time.Now().Unix(),
|
||||
start: time.Now(),
|
||||
client: http.DefaultClient,
|
||||
metrics: make(map[string]interface{}),
|
||||
logOutput: ioutil.Discard,
|
||||
}
|
||||
}
|
||||
|
||||
// SetVersion of locally running Pilosa Cluster to check against master.
|
||||
func (d *Diagnostics) SetVersion(v string) {
|
||||
d.version = v
|
||||
d.Set("Version", v)
|
||||
}
|
||||
|
||||
// SetInterval of the diagnostic go routine and match with the circuit breaker timeout.
|
||||
func (d *Diagnostics) SetInterval(i time.Duration) {
|
||||
d.interval = i
|
||||
}
|
||||
|
||||
// schedule start the diagnostics service ticker.
|
||||
func (d *Diagnostics) schedule() {
|
||||
ticker := time.NewTicker(d.interval)
|
||||
defer ticker.Stop()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-d.closing:
|
||||
return
|
||||
case <-ticker.C:
|
||||
d.CheckVersion()
|
||||
d.Flush()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Flush sends the current metrics.
|
||||
func (d *Diagnostics) Flush() error {
|
||||
d.mu.Lock()
|
||||
d.metrics["uptime"] = (time.Now().Unix() - d.startTime)
|
||||
buf, _ := d.Encode()
|
||||
d.mu.Unlock()
|
||||
|
||||
_, 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
|
||||
body, err := ioutil.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return body, nil
|
||||
})
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
// Open configures the circuit breaker used by the HTTP client.
|
||||
func (d *Diagnostics) Open() {
|
||||
var st gobreaker.Settings
|
||||
if d.interval > 0 {
|
||||
st.Timeout = d.interval * 2
|
||||
}
|
||||
d.cb = gobreaker.NewCircuitBreaker(st)
|
||||
|
||||
d.logger().Printf("Pilosa is currently configured to send small diagnostics reports to our team every hour. More information here: https://www.pilosa.com/docs/latest/administration/#diagnostics")
|
||||
}
|
||||
|
||||
// Close notify goroutine to stop.
|
||||
func (d *Diagnostics) Close() error {
|
||||
close(d.closing)
|
||||
d.wg.Wait()
|
||||
return nil
|
||||
}
|
||||
|
||||
// CheckVersion of the local build against Pilosa master.
|
||||
func (d *Diagnostics) CheckVersion() error {
|
||||
var rsp versionResponse
|
||||
req, err := http.NewRequest("GET", d.VersionURL, nil)
|
||||
resp, err := d.client.Do(req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return fmt.Errorf("http: status=%d", resp.StatusCode)
|
||||
} else if err := json.NewDecoder(resp.Body).Decode(&rsp); err != nil {
|
||||
return fmt.Errorf("json decode: %s", err)
|
||||
}
|
||||
|
||||
// Same a version as last test
|
||||
if rsp.Version == d.lastVersion {
|
||||
return nil
|
||||
}
|
||||
|
||||
d.lastVersion = rsp.Version
|
||||
if err := d.CompareVersion(rsp.Version); err != nil {
|
||||
d.logger().Printf("%s\n", err.Error())
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// CompareVersion check version strings.
|
||||
func (d *Diagnostics) CompareVersion(value string) error {
|
||||
currentVersion := VersionSegments(value)
|
||||
localVersion := VersionSegments(d.version)
|
||||
|
||||
if localVersion[0] < currentVersion[0] { //Major
|
||||
return fmt.Errorf("Warning: You are running Pilosa %s. A newer version (%s) is available: https://github.com/pilosa/pilosa/releases", d.version, value)
|
||||
} else if localVersion[1] < currentVersion[1] { // Minor
|
||||
return fmt.Errorf("Warning: You are running Pilosa %s. The latest Minor release is %s: https://github.com/pilosa/pilosa/releases", d.version, value)
|
||||
} else if localVersion[2] < currentVersion[2] { // Patch
|
||||
return fmt.Errorf("There is a new patch release of Pilosa availbale: %s: https://github.com/pilosa/pilosa/releases", value)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// 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 interface{}) {
|
||||
d.mu.Lock()
|
||||
defer d.mu.Unlock()
|
||||
d.metrics[name] = value
|
||||
}
|
||||
|
||||
// SetLogger Set the logger output type.
|
||||
func (d *Diagnostics) SetLogger(logger io.Writer) {
|
||||
d.logOutput = logger
|
||||
}
|
||||
|
||||
// logger returns a logger that writes to LogOutput.
|
||||
func (d *Diagnostics) logger() *log.Logger {
|
||||
return log.New(d.logOutput, "", log.LstdFlags)
|
||||
}
|
||||
|
||||
// VersionSegments returns the numeric segments of the version as a slice of ints.
|
||||
func VersionSegments(segments string) []int {
|
||||
segments = strings.Trim(segments, "v")
|
||||
segments = strings.Split(segments, "-")[0]
|
||||
s := strings.Split(segments, ".")
|
||||
segmentSlice := make([]int, len(s))
|
||||
for i, v := range s {
|
||||
segmentSlice[i], _ = strconv.Atoi(v)
|
||||
}
|
||||
return segmentSlice
|
||||
}
|
||||
158
diagnostics/diagnostics_test.go
Normal file
158
diagnostics/diagnostics_test.go
Normal file
|
|
@ -0,0 +1,158 @@
|
|||
package diagnostics_test
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"io/ioutil"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"reflect"
|
||||
"runtime"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/pilosa/pilosa/diagnostics"
|
||||
)
|
||||
|
||||
func TestDiagnosticsClient(t *testing.T) {
|
||||
// Mock server.
|
||||
server := httptest.NewServer(nil)
|
||||
defer server.Close()
|
||||
|
||||
// Create a new client.
|
||||
d := diagnostics.New(server.URL)
|
||||
d.SetLogger(ioutil.Discard)
|
||||
d.Open()
|
||||
defer d.Close()
|
||||
|
||||
d.Set("gg", 10)
|
||||
d.Set("ss", "ss")
|
||||
|
||||
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(`{"gg":10,"ss":"ss"}`)
|
||||
if eq, err = compareJSON(data, output1); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !eq {
|
||||
t.Fatalf("unexpected diagnostics: %+v", string(data))
|
||||
}
|
||||
|
||||
// Test the metrics after a flush.
|
||||
d.Flush()
|
||||
data, err = d.Encode()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
output2 := []byte(`{"gg":10,"ss":"ss","uptime":0}`)
|
||||
if eq, err = compareJSON(data, output2); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !eq {
|
||||
t.Fatalf("unexpected diagnostics after flush: %+v", string(data))
|
||||
}
|
||||
}
|
||||
|
||||
func TestDiagnosticsVersion_Parse(t *testing.T) {
|
||||
version := "0.1.1"
|
||||
vs := diagnostics.VersionSegments(version)
|
||||
|
||||
output := []int{0, 1, 1}
|
||||
if !reflect.DeepEqual(vs, output) {
|
||||
t.Fatalf("unexpected version: %+v", vs)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDiagnosticsVersion_Compare(t *testing.T) {
|
||||
d := diagnostics.New("localhost:10101")
|
||||
d.Open()
|
||||
defer d.Close()
|
||||
|
||||
version := "v0.1.1"
|
||||
d.SetVersion(version)
|
||||
|
||||
err := d.CompareVersion("v1.7.0")
|
||||
if !strings.Contains(err.Error(), "A newer version") {
|
||||
t.Fatalf("Expected a newer version is available, actual error: %s", err)
|
||||
}
|
||||
err = d.CompareVersion("1.7.0")
|
||||
if !strings.Contains(err.Error(), "A newer version") {
|
||||
t.Fatalf("Expected a newer version is available, actual error: %s", err)
|
||||
}
|
||||
err = d.CompareVersion("0.7.0")
|
||||
if !strings.Contains(err.Error(), "The latest Minor release is") {
|
||||
t.Fatalf("Expected Minor Version Missmatch, actual error: %s", err)
|
||||
}
|
||||
err = d.CompareVersion("0.1.2")
|
||||
if !strings.Contains(err.Error(), "There is a new patch release of Pilosa") {
|
||||
t.Fatalf("Expected Patch Version Missmatch, actual error: %s", err)
|
||||
}
|
||||
err = d.CompareVersion("0.1.1")
|
||||
if err != nil {
|
||||
t.Fatalf("Versions should match")
|
||||
}
|
||||
}
|
||||
|
||||
func TestDiagnosticsVersion_Check(t *testing.T) {
|
||||
// Mock server.
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.WriteHeader(http.StatusOK)
|
||||
json.NewEncoder(w).Encode(versionResponse{
|
||||
Version: "1.1.1",
|
||||
})
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
// Create a new client.
|
||||
d := diagnostics.New("localhost:10101")
|
||||
defer d.Close()
|
||||
|
||||
version := "0.1.1"
|
||||
d.SetVersion(version)
|
||||
d.VersionURL = server.URL
|
||||
|
||||
d.CheckVersion()
|
||||
}
|
||||
|
||||
type versionResponse struct {
|
||||
Version string `json:"version"`
|
||||
}
|
||||
|
||||
func compareJSON(a, b []byte) (bool, error) {
|
||||
var j1, j2 interface{}
|
||||
if err := json.Unmarshal(a, &j1); err != nil {
|
||||
return false, err
|
||||
}
|
||||
if err := json.Unmarshal(b, &j2); err != nil {
|
||||
return false, err
|
||||
}
|
||||
return reflect.DeepEqual(j1, j2), nil
|
||||
}
|
||||
|
||||
func BenchmarkDiagnostics(b *testing.B) {
|
||||
// Mock server.
|
||||
server := httptest.NewServer(nil)
|
||||
defer server.Close()
|
||||
|
||||
// Create a new client.
|
||||
d := diagnostics.New(server.URL)
|
||||
d.SetLogger(ioutil.Discard)
|
||||
defer d.Close()
|
||||
|
||||
prev := runtime.GOMAXPROCS(4)
|
||||
defer runtime.GOMAXPROCS(prev)
|
||||
|
||||
b.ResetTimer()
|
||||
|
||||
b.RunParallel(func(pb *testing.PB) {
|
||||
for pb.Next() {
|
||||
d.Set("cc", 1)
|
||||
d.Set("gg", "test")
|
||||
}
|
||||
})
|
||||
}
|
||||
|
|
@ -112,6 +112,27 @@ Note: This will only work when the replication factor is >= 2
|
|||
- Restart the cluster
|
||||
- Wait for the 1st sync (10 minutes) to validate Index connections
|
||||
|
||||
#### Diagnostics
|
||||
|
||||
Each Pilosa cluster is configured by default to share anonymous usage details with Pilosa Corp. These metrics allow us to understand how Pilosa is used by the community and improve the technology to suit your needs. Diagnostics are sent to Pilosa every hour. Each of the metrics are detailed below as well as opt-out instructions.
|
||||
|
||||
<strong id="version">Version:</strong> Version string of the build.
|
||||
<strong id="host">Host:</strong> Host URI.
|
||||
<strong id="cluster">Cluster:</strong> List of nodes in the Cluster.
|
||||
<strong id="num_nodes">NumNodes:</strong> Number of nodes in the Cluster.
|
||||
<strong id="num_cpu">NumCPU:</strong> Number of Cores per Node
|
||||
<strong id="bsa_enabled">BSIEnabled:</strong> Bit Slice Index Frames in use.
|
||||
<strong id="time_quantum_enabled">TimeQuantumEnabled:</strong> Time Quantum Frames in use.
|
||||
<strong id="inverse_enabled">InverseEnabled:</strong> Inverse Frames in use.
|
||||
<strong id="num_indexes">NumIndexes:</strong> Number of Indexes in the Cluster.
|
||||
<strong id="num_frames">NumFrames:</strong> Number of Frames in the Cluster.
|
||||
<strong id="num_slices">NumSlices:</strong> Number of Slices in the Cluster.
|
||||
<strong id="num_views">NumViews:</strong> Number of Views in the Cluster.
|
||||
<strong id="open_files">OpenFiles:</strong> Open file handle count.
|
||||
<strong id="go_routines">GoRoutines:</strong> Go routine count.
|
||||
|
||||
You can opt-out of the Pilosa diagnostics reporting by setting either the command line configuration option `--metric.diagnostics=false`, use the `PILOSA_METRIC_DIAGNOSTICS` environment variable, or the TOML configuration file `[metric]` `diagnostics` option.
|
||||
|
||||
#### Metrics
|
||||
|
||||
Pilosa can be configured to emit metrics pertaining to its internal processes in one of two formats: Expvar or StatsD. Metric recording is disabled by default.
|
||||
|
|
|
|||
|
|
@ -219,6 +219,19 @@ Any flag that has a value that is a comma separated list on the command line bec
|
|||
poll-interval = "0m15s"
|
||||
```
|
||||
|
||||
##### Metric Diagnostics
|
||||
|
||||
* Description: Enable diagnostic reporting. To disable diagnostics set to false.
|
||||
* Flag: `metric.diagnostics`
|
||||
* Env: `PILOSA_METRIC_DIAGNOSTICS`
|
||||
* Config:
|
||||
|
||||
```toml
|
||||
[metric]
|
||||
diagnostics = true
|
||||
```
|
||||
|
||||
|
||||
##### TLS Certificate
|
||||
|
||||
* Description: Path to the TLS certificate to use for serving HTTPS. Usually has one of`.crt` or `.pem` extensions.
|
||||
|
|
|
|||
|
|
@ -123,11 +123,14 @@ func (h *Holder) Open() error {
|
|||
h.wg.Add(1)
|
||||
go func() { defer h.wg.Done(); h.monitorCacheFlush() }()
|
||||
|
||||
h.Stats.Open()
|
||||
return nil
|
||||
}
|
||||
|
||||
// Close closes all open fragments.
|
||||
func (h *Holder) Close() error {
|
||||
h.Stats.Close()
|
||||
|
||||
// Notify goroutines of closing and wait for completion.
|
||||
close(h.closing)
|
||||
h.wg.Wait()
|
||||
|
|
|
|||
92
server.go
92
server.go
|
|
@ -33,6 +33,7 @@ import (
|
|||
|
||||
"github.com/CAFxX/gcnotifier"
|
||||
"github.com/gogo/protobuf/proto"
|
||||
"github.com/pilosa/pilosa/diagnostics"
|
||||
"github.com/pilosa/pilosa/internal"
|
||||
"golang.org/x/net/context"
|
||||
)
|
||||
|
|
@ -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: 0,
|
||||
|
||||
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
|
||||
}
|
||||
|
|
@ -223,10 +230,10 @@ func (s *Server) Addr() net.Addr {
|
|||
return s.ln.Addr()
|
||||
}
|
||||
|
||||
// Logger returns a logger that writes to LogOutput
|
||||
func (s *Server) Logger() *log.Logger { return log.New(s.LogOutput, "", log.LstdFlags) }
|
||||
|
||||
func (s *Server) monitorAntiEntropy() {
|
||||
t := time.Now()
|
||||
ticker := time.NewTicker(s.AntiEntropyInterval)
|
||||
defer ticker.Stop()
|
||||
|
||||
|
|
@ -240,7 +247,7 @@ func (s *Server) monitorAntiEntropy() {
|
|||
case <-ticker.C:
|
||||
s.Holder.Stats.Count("AntiEntropy", 1, 1.0)
|
||||
}
|
||||
|
||||
t := time.Now()
|
||||
s.Logger().Printf("holder sync beginning")
|
||||
|
||||
// Initialize syncer with local holder and remote client.
|
||||
|
|
@ -259,9 +266,9 @@ func (s *Server) monitorAntiEntropy() {
|
|||
|
||||
// Record successful sync in log.
|
||||
s.Logger().Printf("holder sync complete")
|
||||
dif := time.Since(t)
|
||||
s.Holder.Stats.Histogram("AntiEntropyDuration", float64(dif), 1.0)
|
||||
}
|
||||
dif := time.Since(t)
|
||||
s.Holder.Stats.Histogram("AntiEntropyDuration", float64(dif), 1.0)
|
||||
}
|
||||
|
||||
// monitorMaxSlices periodically pulls the highest slice from each node in the cluster.
|
||||
|
|
@ -490,9 +497,66 @@ func (s *Server) checkMaxSlices(scheme string, hostPort string) (map[string]uint
|
|||
return s.defaultClient.MaxSliceByIndex(ctx)
|
||||
}
|
||||
|
||||
// monitorDiagnostics periodically polls the the Pilosa Indexes for cluster info.
|
||||
func (s *Server) monitorDiagnostics() {
|
||||
if s.DiagnosticInterval <= 0 {
|
||||
s.Logger().Printf("diagnostics disabled")
|
||||
return
|
||||
}
|
||||
|
||||
s.diagnostics.SetLogger(s.LogOutput)
|
||||
s.diagnostics.SetVersion(Version)
|
||||
s.diagnostics.SetInterval(s.DiagnosticInterval)
|
||||
s.diagnostics.Open()
|
||||
s.diagnostics.Set("Host", s.URI.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
|
||||
|
||||
// Flush the diagnostics metrics at startup, then on each tick interval
|
||||
flush := func() {
|
||||
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()
|
||||
}
|
||||
|
||||
ticker := time.NewTicker(s.DiagnosticInterval)
|
||||
defer ticker.Stop()
|
||||
flush()
|
||||
for {
|
||||
// Wait for tick or a close.
|
||||
select {
|
||||
case <-s.closing:
|
||||
return
|
||||
case <-ticker.C:
|
||||
flush()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// monitorRuntime periodically polls the Go runtime metrics.
|
||||
func (s *Server) monitorRuntime() {
|
||||
// Disable metrics when poll interval is zero
|
||||
// Disable metrics when poll interval is zero.
|
||||
if s.MetricInterval <= 0 {
|
||||
return
|
||||
}
|
||||
|
|
@ -512,18 +576,18 @@ func (s *Server) monitorRuntime() {
|
|||
case <-s.closing:
|
||||
return
|
||||
case <-gcn.AfterGC():
|
||||
// GC just ran
|
||||
// GC just ran.
|
||||
s.Holder.Stats.Count("garbage_collection", 1, 1.0)
|
||||
case <-ticker.C:
|
||||
}
|
||||
|
||||
// Record the number of go routines
|
||||
// Record the number of go routines.
|
||||
s.Holder.Stats.Gauge("goroutines", float64(runtime.NumGoroutine()), 1.0)
|
||||
|
||||
// Open File handles
|
||||
// Open File handles.
|
||||
s.Holder.Stats.Gauge("OpenFiles", float64(CountOpenFiles()), 1.0)
|
||||
|
||||
// Runtime memory metrics
|
||||
// Runtime memory metrics.
|
||||
runtime.ReadMemStats(&m)
|
||||
s.Holder.Stats.Gauge("HeapAlloc", float64(m.HeapAlloc), 1.0)
|
||||
s.Holder.Stats.Gauge("HeapInuse", float64(m.HeapInuse), 1.0)
|
||||
|
|
@ -541,7 +605,7 @@ func (s *Server) createDefaultClient() {
|
|||
s.defaultClient = NewInternalHTTPClientFromURI(nil, &ClientOptions{TLS: s.TLS})
|
||||
}
|
||||
|
||||
// CountOpenFiles on opperating systems that support lsof
|
||||
// CountOpenFiles on operating systems that support lsof.
|
||||
func CountOpenFiles() int {
|
||||
count := 0
|
||||
|
||||
|
|
|
|||
|
|
@ -30,6 +30,7 @@ import (
|
|||
"time"
|
||||
|
||||
"crypto/tls"
|
||||
|
||||
"github.com/pilosa/pilosa"
|
||||
"github.com/pilosa/pilosa/gossip"
|
||||
"github.com/pilosa/pilosa/statsd"
|
||||
|
|
@ -43,6 +44,9 @@ func init() {
|
|||
const (
|
||||
// DefaultDataDir is the default data directory.
|
||||
DefaultDataDir = "~/.pilosa"
|
||||
|
||||
// DefaultDiagnosticsInterval is the default sync frequency diagnostic metrics.
|
||||
DefaultDiagnosticsInterval = 1 * time.Hour
|
||||
)
|
||||
|
||||
// Command represents the state of the pilosa server command.
|
||||
|
|
@ -143,6 +147,9 @@ func (m *Command) SetupServer() error {
|
|||
m.Server.Logger().Printf("Using data from: %s\n", m.Config.DataDir)
|
||||
m.Server.Holder.Path = m.Config.DataDir
|
||||
m.Server.MetricInterval = time.Duration(m.Config.Metric.PollInterval)
|
||||
if m.Config.Metric.Diagnostics {
|
||||
m.Server.DiagnosticInterval = time.Duration(DefaultDiagnosticsInterval)
|
||||
}
|
||||
m.Server.Holder.Stats, err = NewStatsClient(m.Config.Metric.Service, m.Config.Metric.Host)
|
||||
if err != nil {
|
||||
return err
|
||||
|
|
|
|||
36
stats.go
36
stats.go
|
|
@ -58,6 +58,12 @@ type StatsClient interface {
|
|||
|
||||
// SetLogger Set the logger output type
|
||||
SetLogger(logger io.Writer)
|
||||
|
||||
// Starts the service
|
||||
Open()
|
||||
|
||||
// Closes the client
|
||||
Close() error
|
||||
}
|
||||
|
||||
// NopStatsClient represents a client that doesn't do anything.
|
||||
|
|
@ -74,6 +80,8 @@ func (c *nopStatsClient) Histogram(name string, value float64, rate float64)
|
|||
func (c *nopStatsClient) Set(name string, value string, rate float64) {}
|
||||
func (c *nopStatsClient) Timing(name string, value time.Duration, rate float64) {}
|
||||
func (c *nopStatsClient) SetLogger(logger io.Writer) {}
|
||||
func (c *nopStatsClient) Open() {}
|
||||
func (c *nopStatsClient) Close() error { return nil }
|
||||
|
||||
// ExpvarStatsClient writes stats out to expvars.
|
||||
type ExpvarStatsClient struct {
|
||||
|
|
@ -145,10 +153,16 @@ func (c *ExpvarStatsClient) Timing(name string, value time.Duration, rate float6
|
|||
c.mu.Unlock()
|
||||
}
|
||||
|
||||
// SetLogger has no logger
|
||||
// SetLogger has no logger.
|
||||
func (c *ExpvarStatsClient) SetLogger(logger io.Writer) {
|
||||
}
|
||||
|
||||
// Open no-op.
|
||||
func (c *ExpvarStatsClient) Open() {}
|
||||
|
||||
// Close no-op.
|
||||
func (c *ExpvarStatsClient) Close() error { return nil }
|
||||
|
||||
// MultiStatsClient joins multiple stats clients together.
|
||||
type MultiStatsClient []StatsClient
|
||||
|
||||
|
|
@ -211,13 +225,31 @@ func (a MultiStatsClient) Timing(name string, value time.Duration, rate float64)
|
|||
}
|
||||
}
|
||||
|
||||
// SetLogger Sets the StatsD logger output type
|
||||
// SetLogger Sets the StatsD logger output type.
|
||||
func (a MultiStatsClient) SetLogger(logger io.Writer) {
|
||||
for _, c := range a {
|
||||
c.SetLogger(logger)
|
||||
}
|
||||
}
|
||||
|
||||
// Open starts the stat service.
|
||||
func (a MultiStatsClient) Open() {
|
||||
for _, c := range a {
|
||||
c.Open()
|
||||
}
|
||||
}
|
||||
|
||||
// Close shuts down the stats clients.
|
||||
func (a MultiStatsClient) Close() error {
|
||||
for _, c := range a {
|
||||
err := c.Close()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// UnionStringSlice returns a sorted set of tags which combine a & b.
|
||||
func UnionStringSlice(a, b []string) []string {
|
||||
// Sort both sets first.
|
||||
|
|
|
|||
|
|
@ -344,3 +344,5 @@ func (c *MockStats) Histogram(name string, value float64, rate float64) {}
|
|||
func (c *MockStats) Set(name string, value string, rate float64) {}
|
||||
func (c *MockStats) Timing(name string, value time.Duration, rate float64) {}
|
||||
func (c *MockStats) SetLogger(logger io.Writer) {}
|
||||
func (c *MockStats) Open() {}
|
||||
func (c *MockStats) Close() error { return nil }
|
||||
|
|
|
|||
|
|
@ -58,6 +58,9 @@ func NewStatsClient(host string) (*StatsClient, error) {
|
|||
}, nil
|
||||
}
|
||||
|
||||
// Open no-op
|
||||
func (c *StatsClient) Open() {}
|
||||
|
||||
// Close closes the connection to the agent.
|
||||
func (c *StatsClient) Close() error {
|
||||
return c.client.Close()
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue