Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions server/e2e/health_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,17 @@ func TestHealthStandalone(t *testing.T) {
if got := grpcHealth(t, s); got != healthpb.HealthCheckResponse_SERVING {
t.Errorf("grpc health: %v", got)
}

// The metrics, with the connection just made.
code, metrics := monitorGet(t, s, "/_metrics")
if code != 200 {
t.Fatalf("/_metrics: %d %s", code, metrics)
}
for _, want := range []string{"unitdb_build_info{", "unitdb_up_seconds ", `unitdb_dependency_up{dep="store"} 1`, "unitdb_connections "} {
if !strings.Contains(metrics, want) {
t.Errorf("/_metrics has no %q:\n%s", want, metrics)
}
}
}

func TestHealthCluster(t *testing.T) {
Expand Down
41 changes: 41 additions & 0 deletions server/internal/cluster_health.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package internal

import (
"fmt"
"sort"
"sync/atomic"
"time"
)
Expand Down Expand Up @@ -57,3 +58,43 @@ func (c *Cluster) readiness() (bool, string) {
return true, fmt.Sprintf("in the ring: %d of %d nodes; leader %s, %s ago",
len(ring), configured, leader, age.Round(time.Millisecond))
}

// writeMetrics writes this node's view of the cluster: only what is cheap
// and safe to read per scrape. Peers are labelled by their node names, a
// fixed set.
func (c *Cluster) writeMetrics(m *metricsWriter) {
b := func(v bool) float64 {
if v {
return 1
}
return 0
}
m.one("unitdb_cluster_nodes", "gauge", "Nodes in the cluster's configuration, this one included.", float64(len(c.nodes)+1))
m.one("unitdb_cluster_members", "gauge", "Nodes in the ring this node routes by.", float64(len(c.getRingNodes())))
m.one("unitdb_cluster_ring_version", "gauge", "The ring version last seen from the leader; 0 before any.", float64(c.clusterRing.Load()))
m.one("unitdb_cluster_rebuilding", "gauge", "Whether this node is copying its topics from the others.", b(c.rebuilding.Load()))
m.one("unitdb_cluster_leaving", "gauge", "Whether this node is leaving the cluster.", b(c.leaving.Load()))
if last := c.health.lastLeader.Load(); last != 0 {
m.one("unitdb_cluster_leader_age_seconds", "gauge", "Seconds since this node last heard from a leader, or was one.", time.Since(time.Unix(0, last)).Seconds())
}
c.pendingMu.Lock()
pending := len(c.pending)
c.pendingMu.Unlock()
m.one("unitdb_cluster_pending_hints", "gauge", "Replicas waiting in memory for a peer to take them.", float64(pending))

names := make([]string, 0, len(c.nodes))
for name := range c.nodes {
names = append(names, name)
}
sort.Strings(names)
m.head("unitdb_replication_queue", "gauge", "Replicas queued for a peer.")
withoutTLS := 0
for _, name := range names {
n := c.nodes[name]
m.value("unitdb_replication_queue", "peer="+quote(name), float64(len(n.repl)))
if caps, ok := n.capabilities(); !ok || !containsNode(caps.Capabilities, capTLS) {
withoutTLS++
}
}
m.one("unitdb_cluster_peers_without_tls", "gauge", "Peers that don't advertise TLS for cluster traffic.", float64(withoutTLS))
}
14 changes: 13 additions & 1 deletion server/internal/health.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"encoding/json"
"errors"
"fmt"
"io"
"net"
"net/http"
"sort"
Expand All @@ -27,7 +28,8 @@ import (
// reads the results the checks left, and never waits on them.
// - /_status: the same, as JSON, with each check's detail.
//
// The gRPC server also answers grpc.health.v1, SERVING while ready.
// The gRPC server also answers grpc.health.v1, SERVING while ready, and
// /_metrics has the server's metrics (metrics.go).
//
// Each check runs every healthInterval in the background, with a
// healthTimeout, and keeps its last result.
Expand Down Expand Up @@ -57,6 +59,8 @@ type healthMonitor struct {
draining atomic.Bool
grpc *health.Server
srv *http.Server
// metrics writes /_metrics, when set.
metrics func(io.Writer)
}

func newHealthMonitor(started time.Time) *healthMonitor {
Expand Down Expand Up @@ -170,6 +174,14 @@ func (h *healthMonitor) handler() http.Handler {
}
w.Write([]byte("ready"))
})
mux.HandleFunc("/_metrics", func(w http.ResponseWriter, r *http.Request) {
if h.metrics == nil {
http.NotFound(w, r)
return
}
w.Header().Set("content-type", "text/plain; version=0.0.4")
h.metrics(w)
})
mux.HandleFunc("/_status", func(w http.ResponseWriter, r *http.Request) {
why := h.notReady()
status := "ready"
Expand Down
136 changes: 136 additions & 0 deletions server/internal/metrics.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,136 @@
package internal

import (
"fmt"
"io"
"runtime/debug"
"sort"
"strings"
"time"
)

// Metrics, in the Prometheus text format, at /_metrics on the monitor port:
// what /varz computes (connections, subscriptions, messages and bytes, and
// the event durations), the health checks, and in a cluster this node's
// place in it. Labels stay bounded: never a client or a topic.

// buildInfo is the version and commit the binary was built from, when go
// recorded them.
var buildInfo = func() (version, commit string) {
version, commit = "unknown", "unknown"
bi, ok := debug.ReadBuildInfo()
if !ok {
return
}
if bi.Main.Version != "" {
version = bi.Main.Version
}
for _, s := range bi.Settings {
if s.Key == "vcs.revision" && s.Value != "" {
commit = s.Value
if len(commit) > 12 {
commit = commit[:12]
}
}
}
return
}

type metricsWriter struct {
w io.Writer
seen map[string]bool
}

// head writes a metric's HELP and TYPE once.
func (m *metricsWriter) head(name, typ, help string) {
if m.seen[name] {
return
}
m.seen[name] = true
fmt.Fprintf(m.w, "# HELP %s %s\n# TYPE %s %s\n", name, help, name, typ)
}

func (m *metricsWriter) value(name, labels string, v float64) {
if labels != "" {
labels = "{" + labels + "}"
}
fmt.Fprintf(m.w, "%s%s %s\n", name, labels, formatFloat(v))
}

func (m *metricsWriter) one(name, typ, help string, v float64) {
m.head(name, typ, help)
m.value(name, "", v)
}

func formatFloat(v float64) string {
s := fmt.Sprintf("%g", v)
if strings.Contains(s, "e+") {
s = fmt.Sprintf("%.0f", v)
}
return s
}

func quote(s string) string {
r := strings.NewReplacer(`\`, `\\`, `"`, `\"`, "\n", `\n`)
return `"` + r.Replace(s) + `"`
}

// writeMetrics writes the server's metrics.
func (s *_Service) writeMetrics(w io.Writer) {
m := &metricsWriter{w: w, seen: map[string]bool{}}

version, commit := buildInfo()
m.head("unitdb_build_info", "gauge", "The build: its version and commit.")
m.value("unitdb_build_info", "version="+quote(version)+",commit="+quote(commit), 1)
m.one("unitdb_up_seconds", "gauge", "Seconds since the server started.", time.Since(s.start).Seconds())

if s.meter != nil {
m.one("unitdb_connections", "gauge", "Client connections open.", float64(s.meter.Connections.Count()))
m.one("unitdb_subscriptions", "gauge", "Subscriptions held.", float64(s.meter.Subscriptions.Count()))
m.one("unitdb_messages_in_total", "counter", "Messages published to this server.", float64(s.meter.InMsgs.Count()))
m.one("unitdb_messages_out_total", "counter", "Messages delivered to subscribers.", float64(s.meter.OutMsgs.Count()))
m.one("unitdb_bytes_in_total", "counter", "Payload bytes published to this server.", float64(s.meter.InBytes.Count()))
m.one("unitdb_bytes_out_total", "counter", "Payload bytes delivered to subscribers.", float64(s.meter.OutBytes.Count()))

// How long handling each client packet took (hdl_conn.go), in nanoseconds.
ts := s.meter.ConnTimeSeries.Snapshot()
const name = "unitdb_event_duration_seconds"
m.head(name, "summary", "Time to handle a client's packet, from the meter's sample.")
for _, q := range []struct {
q string
v int64
}{{"0.5", int64(ts.P50())}, {"0.75", int64(ts.P75())}, {"0.95", int64(ts.P95())}, {"0.99", int64(ts.P99())}, {"0.999", int64(ts.P999())}} {
m.value(name, "quantile="+quote(q.q), time.Duration(q.v).Seconds())
}
}

if s.health != nil {
m.head("unitdb_dependency_up", "gauge", "Whether a health check passes (1) or fails (0).")
m.head("unitdb_dependency_check_seconds", "gauge", "How long a health check took, last time.")
s.health.mu.RLock()
names := make([]string, 0, len(s.health.results))
for n := range s.health.results {
names = append(names, n)
}
sort.Strings(names)
for _, n := range names {
r := s.health.results[n]
up := 0.0
if r.OK {
up = 1
}
m.value("unitdb_dependency_up", "dep="+quote(n), up)
m.value("unitdb_dependency_check_seconds", "dep="+quote(n), float64(r.LatencyMs)/1000)
}
s.health.mu.RUnlock()
draining := 0.0
if s.health.draining.Load() {
draining = 1
}
m.one("unitdb_draining", "gauge", "Whether the server is shutting down.", draining)
}

if c := Globals.Cluster; c != nil {
c.writeMetrics(m)
}
}
71 changes: 71 additions & 0 deletions server/internal/metrics_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
package internal

import (
"bytes"
"strings"
"testing"
"time"
)

func TestMetrics(t *testing.T) {
s := &_Service{start: time.Now().Add(-time.Minute), meter: NewMeter()}
defer s.meter.UnregisterAll()
s.meter.Connections.Inc(3)
s.meter.InMsgs.Inc(7)
s.meter.ConnTimeSeries.AddTime(2 * time.Millisecond)
s.health = newHealthMonitor(s.start)
s.health.add("store", func() (string, error) { return "open", nil })
s.health.runAll()

var b bytes.Buffer
s.writeMetrics(&b)
out := b.String()
for _, want := range []string{
"# TYPE unitdb_build_info gauge",
"unitdb_connections 3",
"unitdb_messages_in_total 7",
`unitdb_event_duration_seconds{quantile="0.5"}`,
`unitdb_dependency_up{dep="store"} 1`,
"unitdb_draining 0",
} {
if !strings.Contains(out, want) {
t.Errorf("no %q in:\n%s", want, out)
}
}
// Each metric's HELP and TYPE once, and every sample line well formed.
for _, line := range strings.Split(strings.TrimSpace(out), "\n") {
if strings.HasPrefix(line, "#") {
continue
}
if f := strings.Fields(line); len(f) != 2 || !strings.HasPrefix(f[0], "unitdb_") {
t.Errorf("bad sample line %q", line)
}
}
if n := strings.Count(out, "# TYPE unitdb_dependency_up "); n != 1 {
t.Errorf("TYPE of unitdb_dependency_up %d times", n)
}
}

func TestClusterMetrics(t *testing.T) {
c := &Cluster{thisNodeName: "a", nodes: map[string]*ClusterNode{
"b": {name: "b", repl: make(chan replicaItem, 4)},
"c": {name: "c", repl: make(chan replicaItem, 4)},
}}
c.ringNodes = []string{"a", "b"}
c.health.leaderSeen("b")
var b bytes.Buffer
c.writeMetrics(&metricsWriter{w: &b, seen: map[string]bool{}})
out := b.String()
for _, want := range []string{
"unitdb_cluster_nodes 3",
"unitdb_cluster_members 2",
`unitdb_replication_queue{peer="b"} 0`,
`unitdb_replication_queue{peer="c"} 0`,
"unitdb_cluster_peers_without_tls 2",
"unitdb_cluster_leader_age_seconds",
} {
if !strings.Contains(out, want) {
t.Errorf("no %q in:\n%s", want, out)
}
}
}
1 change: 1 addition & 0 deletions server/internal/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -213,6 +213,7 @@ func (s *_Service) listen(addr string) {
// The health checks, and grpc.health.v1 on the gRPC server.
s.health = newHealthMonitor(s.start)
s.health.addServiceChecks()
s.health.metrics = s.writeMetrics
s.grpc.Register = func(g *grpc.Server) {
healthpb.RegisterHealthServer(g, s.health.grpc)
}
Expand Down
Loading