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
29 changes: 17 additions & 12 deletions server/e2e/harness_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -135,7 +135,9 @@ type server struct {
dbPath string
tcpAddr string
grpcAddr string
logs *syncBuffer
// monitorAddr serves the health checks: /_healthz, /_readyz, /_status.
monitorAddr string
logs *syncBuffer
// env is added to the server's environment.
env []string
// noWait starts the server without waiting for it to be ready.
Expand Down Expand Up @@ -242,6 +244,7 @@ func startServerWith(t *testing.T, opts serverOpts) *server {
confName := fmt.Sprintf("e2e-%d.conf", freePort(t))
tcpPort := freePort(t)
grpcPort := freePort(t)
monitorPort := freePort(t)
dbPath, err := os.MkdirTemp("", "unitdb-e2e-db")
if err != nil {
t.Fatal(err)
Expand All @@ -250,13 +253,14 @@ func startServerWith(t *testing.T, opts serverOpts) *server {
return fmt.Sprintf(`{
"listen": "127.0.0.1:%d",
"grpc_listen": "127.0.0.1:%d",
"monitor_listen": "127.0.0.1:%d",
"logging_level": %q,
"allow_insecure": %t,
%s
"encryption_config": {"key": %q, "identifier": "local", "sealed": false, "timestamp": 1522325758},
"cluster_config": %s,
"store_config": {"reset": false, "adapters": {"unitdb": {"database": "unitdb", "mem_size": 500000000}}}
}`, tcpPort, grpcPort, opts.logLevel, opts.allowInsecure, opts.extra, opts.key, cluster)
}`, tcpPort, grpcPort, monitorPort, opts.logLevel, opts.allowInsecure, opts.extra, opts.key, cluster)
}
confPath := filepath.Join(binDir, confName)
if err := os.WriteFile(confPath, []byte(confWith(opts.cluster)), 0644); err != nil {
Expand All @@ -273,16 +277,17 @@ func startServerWith(t *testing.T, opts serverOpts) *server {
}

s := &server{
t: t,
cmd: cmd,
dbPath: dbPath,
tcpAddr: fmt.Sprintf("127.0.0.1:%d", tcpPort),
grpcAddr: fmt.Sprintf("127.0.0.1:%d", grpcPort),
logs: logs,
env: opts.env,
noWait: opts.expectExit,
confPath: confPath,
confWith: confWith,
t: t,
cmd: cmd,
dbPath: dbPath,
tcpAddr: fmt.Sprintf("127.0.0.1:%d", tcpPort),
grpcAddr: fmt.Sprintf("127.0.0.1:%d", grpcPort),
monitorAddr: fmt.Sprintf("127.0.0.1:%d", monitorPort),
logs: logs,
env: opts.env,
noWait: opts.expectExit,
confPath: confPath,
confWith: confWith,
}
s.watch()
t.Cleanup(func() {
Expand Down
122 changes: 122 additions & 0 deletions server/e2e/health_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,122 @@
package e2e

// The health checks: /_healthz, /_readyz and /_status on the monitor port,
// and grpc.health.v1 on the gRPC port, for a standalone server and a
// cluster's nodes.

import (
"context"
"encoding/json"
"io"
"net/http"
"strings"
"testing"
"time"

"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
healthpb "google.golang.org/grpc/health/grpc_health_v1"
)

type healthStatus struct {
Status string `json:"status"`
Reason string `json:"reason"`
Checks map[string]struct {
OK bool `json:"ok"`
Detail string `json:"detail"`
} `json:"checks"`
}

func monitorGet(t *testing.T, s *server, path string) (int, string) {
t.Helper()
c := &http.Client{Timeout: 5 * time.Second}
resp, err := c.Get("http://" + s.monitorAddr + path)
if err != nil {
return 0, err.Error()
}
defer resp.Body.Close()
b, _ := io.ReadAll(resp.Body)
return resp.StatusCode, strings.TrimSpace(string(b))
}

// untilReady waits for s's /_readyz to answer 200, and returns its status.
func untilReady(t *testing.T, s *server, within time.Duration) healthStatus {
t.Helper()
deadline := time.Now().Add(within)
for {
code, body := monitorGet(t, s, "/_readyz")
if code == http.StatusOK {
break
}
if time.Now().After(deadline) {
_, status := monitorGet(t, s, "/_status")
t.Fatalf("not ready after %s: %d %q\n%s", within, code, body, status)
}
time.Sleep(200 * time.Millisecond)
}
_, body := monitorGet(t, s, "/_status")
var st healthStatus
if err := json.Unmarshal([]byte(body), &st); err != nil {
t.Fatalf("/_status: %v: %s", err, body)
}
return st
}

func grpcHealth(t *testing.T, s *server) healthpb.HealthCheckResponse_ServingStatus {
t.Helper()
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
// NewClient connects on the first call: Check waits for it, within ctx.
conn, err := grpc.NewClient(s.grpcAddr, grpc.WithTransportCredentials(insecure.NewCredentials()))
if err != nil {
t.Fatalf("grpc client: %v", err)
}
defer conn.Close()
r, err := healthpb.NewHealthClient(conn).Check(ctx, &healthpb.HealthCheckRequest{}, grpc.WaitForReady(true))
if err != nil {
t.Fatalf("grpc health: %v", err)
}
return r.Status
}

func TestHealthStandalone(t *testing.T) {
s := startServer(t)

if code, body := monitorGet(t, s, "/_healthz"); code != 200 || body != "ok" {
t.Errorf("/_healthz: %d %q", code, body)
}
st := untilReady(t, s, 15*time.Second)
if st.Status != "ready" || !st.Checks["store"].OK {
t.Errorf("/_status: %+v", st)
}
if _, ok := st.Checks["cluster"]; ok {
t.Errorf("a standalone server has a cluster check: %+v", st)
}
if got := grpcHealth(t, s); got != healthpb.HealthCheckResponse_SERVING {
t.Errorf("grpc health: %v", got)
}
}

func TestHealthCluster(t *testing.T) {
c := startCluster(t, "one", "two", "three")
for _, n := range c.nodes {
st := untilReady(t, n.server, 30*time.Second)
if !strings.Contains(st.Checks["cluster"].Detail, "in the ring: 3 of 3 nodes") {
t.Errorf("%s: cluster check: %+v", n.name, st.Checks["cluster"])
}
}

// A node shut down answers nothing; the others stay ready.
gone := c.nodes[2]
gone.shutdown()
if code, _ := monitorGet(t, gone.server, "/_readyz"); code == http.StatusOK {
t.Errorf("%s is ready after its shutdown", gone.name)
}
for _, n := range c.nodes[:2] {
// The leader may have been the one that left: an election first.
st := untilReady(t, n.server, 30*time.Second)
if !st.Checks["cluster"].OK {
t.Errorf("%s: %+v", n.name, st.Checks["cluster"])
}
}
}
2 changes: 2 additions & 0 deletions server/internal/cluster.go
Original file line number Diff line number Diff line change
Expand Up @@ -650,6 +650,8 @@ type Cluster struct {
leaving atomic.Bool
// Set once the cluster is shut down.
stopped atomic.Bool
// When this node last heard from a leader, for its readiness.
health clusterHealth
// Time it waits at most for that.
drainTimeout time.Duration

Expand Down
59 changes: 59 additions & 0 deletions server/internal/cluster_health.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
package internal

import (
"fmt"
"sync/atomic"
"time"
)

// clusterHealth is what the health checks read of the cluster: when this
// node last heard from a leader, or was one, and which. The failover runner sets
// it; the fields it keeps itself aren't safe to read from elsewhere.
type clusterHealth struct {
lastLeader atomic.Int64 // unix nanoseconds; 0 before any
leader atomic.Value // string
}

// leaderSeen records a ping from leader, or this node pinging as the leader.
func (h *clusterHealth) leaderSeen(leader string) {
h.leader.Store(leader)
h.lastLeader.Store(time.Now().UnixNano())
}

// readiness says whether this node can serve its part of the cluster: it
// isn't leaving or stopped, isn't copying its topics from the others, is in
// the ring, and has heard from a leader lately. A nil cluster (a single
// server) is ready.
func (c *Cluster) readiness() (bool, string) {
if c == nil {
return true, "standalone"
}
switch {
case c.stopped.Load():
return false, "stopped"
case c.leaving.Load():
return false, "leaving the cluster"
case c.rebuilding.Load():
return false, "catching up: copying its topics from the other nodes"
}
ring := c.getRingNodes()
configured := len(c.nodes) + 1
if !containsNode(ring, c.thisNodeName) {
return false, fmt.Sprintf("not in the ring yet (the ring has %d of %d nodes)", len(ring), configured)
}
if c.fo == nil {
return true, fmt.Sprintf("in the ring: %d of %d nodes", len(ring), configured)
}
last := c.health.lastLeader.Load()
if last == 0 {
return false, "no leader heard from yet"
}
age := time.Since(time.Unix(0, last))
// When a follower would start an election.
if stale := c.fo.heartBeat * time.Duration(c.fo.voteTimeout+1); age > stale {
return false, fmt.Sprintf("no leader for %s", age.Round(time.Second))
}
leader, _ := c.health.leader.Load().(string)
return true, fmt.Sprintf("in the ring: %d of %d nodes; leader %s, %s ago",
len(ring), configured, leader, age.Round(time.Millisecond))
}
2 changes: 2 additions & 0 deletions server/internal/cluster_leader.go
Original file line number Diff line number Diff line change
Expand Up @@ -364,6 +364,7 @@ func (c *Cluster) run() {
if c.fo.leader == c.thisNodeName {
// I'm the leader, send pings
c.sendPings()
c.health.leaderSeen(c.thisNodeName)
} else {
missed++
if missed >= c.fo.voteTimeout {
Expand Down Expand Up @@ -399,6 +400,7 @@ func (c *Cluster) run() {
}

missed = 0
c.health.leaderSeen(ping.Leader)
if c.leaving.Load() && !containsNode(ping.Nodes, c.thisNodeName) {
c.fo.pingsWithoutSelf.Add(1)
}
Expand Down
4 changes: 4 additions & 0 deletions server/internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,10 @@ type Config struct {
// Can be overridden from the command line, see option --listen.
GrpcListen string `json:"grpc_listen"`

// MonitorListen is where the health checks are served: /_healthz,
// /_readyz and /_status. Off when empty.
MonitorListen string `json:"monitor_listen"`

// Default logging level is "InfoLevel" so to enable the debug log set the "LogLevel" to "DebugLevel".
LoggingLevel string `json:"logging_level"`

Expand Down
Loading
Loading