diff --git a/server/e2e/health_test.go b/server/e2e/health_test.go index fbe1328..3366bf4 100644 --- a/server/e2e/health_test.go +++ b/server/e2e/health_test.go @@ -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) { diff --git a/server/internal/cluster_health.go b/server/internal/cluster_health.go index bc2148f..134817d 100644 --- a/server/internal/cluster_health.go +++ b/server/internal/cluster_health.go @@ -2,6 +2,7 @@ package internal import ( "fmt" + "sort" "sync/atomic" "time" ) @@ -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)) +} diff --git a/server/internal/health.go b/server/internal/health.go index b3e5b67..b20ec88 100644 --- a/server/internal/health.go +++ b/server/internal/health.go @@ -5,6 +5,7 @@ import ( "encoding/json" "errors" "fmt" + "io" "net" "net/http" "sort" @@ -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. @@ -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 { @@ -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" diff --git a/server/internal/metrics.go b/server/internal/metrics.go new file mode 100644 index 0000000..7692f44 --- /dev/null +++ b/server/internal/metrics.go @@ -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) + } +} diff --git a/server/internal/metrics_test.go b/server/internal/metrics_test.go new file mode 100644 index 0000000..b63f632 --- /dev/null +++ b/server/internal/metrics_test.go @@ -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) + } + } +} diff --git a/server/internal/service.go b/server/internal/service.go index a9aff5e..c76c48d 100644 --- a/server/internal/service.go +++ b/server/internal/service.go @@ -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) }