diff --git a/server/e2e/harness_test.go b/server/e2e/harness_test.go index d65470d..4f5fdad 100644 --- a/server/e2e/harness_test.go +++ b/server/e2e/harness_test.go @@ -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. @@ -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) @@ -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 { @@ -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() { diff --git a/server/e2e/health_test.go b/server/e2e/health_test.go new file mode 100644 index 0000000..fbe1328 --- /dev/null +++ b/server/e2e/health_test.go @@ -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"]) + } + } +} diff --git a/server/internal/cluster.go b/server/internal/cluster.go index edcdf93..e32ce72 100644 --- a/server/internal/cluster.go +++ b/server/internal/cluster.go @@ -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 diff --git a/server/internal/cluster_health.go b/server/internal/cluster_health.go new file mode 100644 index 0000000..bc2148f --- /dev/null +++ b/server/internal/cluster_health.go @@ -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)) +} diff --git a/server/internal/cluster_leader.go b/server/internal/cluster_leader.go index fb39270..d87bc12 100644 --- a/server/internal/cluster_leader.go +++ b/server/internal/cluster_leader.go @@ -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 { @@ -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) } diff --git a/server/internal/config/config.go b/server/internal/config/config.go index 268a698..88d8ff6 100644 --- a/server/internal/config/config.go +++ b/server/internal/config/config.go @@ -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"` diff --git a/server/internal/health.go b/server/internal/health.go new file mode 100644 index 0000000..b3e5b67 --- /dev/null +++ b/server/internal/health.go @@ -0,0 +1,255 @@ +package internal + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "net" + "net/http" + "sort" + "sync" + "sync/atomic" + "time" + + "google.golang.org/grpc/health" + healthpb "google.golang.org/grpc/health/grpc_health_v1" + + "github.com/unit-io/unitdb/server/internal/pkg/log" + "github.com/unit-io/unitdb/server/internal/store" +) + +// Health checks, on a port of their own (monitor_listen), apart from the +// clients': +// +// - /_healthz: the process serves requests. It checks nothing else. +// - /_readyz: 200 "ready", or 503 and why: draining, or a check failed. It +// 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. +// +// Each check runs every healthInterval in the background, with a +// healthTimeout, and keeps its last result. + +const ( + healthInterval = 5 * time.Second + healthTimeout = 3 * time.Second +) + +type healthResult struct { + OK bool `json:"ok"` + Detail string `json:"detail"` + CheckedAt time.Time `json:"checked_at"` + LatencyMs int64 `json:"latency_ms"` +} + +type healthCheck struct { + name string + run func() (string, error) +} + +type healthMonitor struct { + started time.Time + checks []healthCheck + mu sync.RWMutex + results map[string]healthResult + draining atomic.Bool + grpc *health.Server + srv *http.Server +} + +func newHealthMonitor(started time.Time) *healthMonitor { + return &healthMonitor{ + started: started, + results: map[string]healthResult{}, + grpc: health.NewServer(), + } +} + +// add adds a check, before run. +func (h *healthMonitor) add(name string, run func() (string, error)) { + h.checks = append(h.checks, healthCheck{name, run}) +} + +// runAll runs every check once, each within healthTimeout. +func (h *healthMonitor) runAll() { + var wg sync.WaitGroup + for _, c := range h.checks { + wg.Add(1) + go func(c healthCheck) { + defer wg.Done() + start := time.Now() + type answer struct { + detail string + err error + } + done := make(chan answer, 1) + go func() { + d, err := c.run() + done <- answer{d, err} + }() + r := healthResult{CheckedAt: time.Now().UTC()} + select { + case a := <-done: + r.OK, r.Detail = a.err == nil, a.detail + if a.err != nil { + r.Detail = a.err.Error() + } + case <-time.After(healthTimeout): + r.Detail = fmt.Sprintf("no answer in %s", healthTimeout) + } + r.LatencyMs = time.Since(start).Milliseconds() + h.mu.Lock() + h.results[c.name] = r + h.mu.Unlock() + }(c) + } + wg.Wait() + h.updateGRPC() +} + +// run checks now, then every healthInterval, until ctx is done. +func (h *healthMonitor) run(ctx context.Context) { + h.runAll() + t := time.NewTicker(healthInterval) + defer t.Stop() + for { + select { + case <-ctx.Done(): + return + case <-t.C: + h.runAll() + } + } +} + +// notReady returns why the server isn't ready, or "". +func (h *healthMonitor) notReady() string { + if h.draining.Load() { + return "draining" + } + h.mu.RLock() + defer h.mu.RUnlock() + for _, c := range h.checks { + r, ok := h.results[c.name] + if !ok { + return c.name + ": not checked yet" + } + if !r.OK { + return c.name + ": " + r.Detail + } + } + return "" +} + +func (h *healthMonitor) updateGRPC() { + status := healthpb.HealthCheckResponse_SERVING + if h.notReady() != "" { + status = healthpb.HealthCheckResponse_NOT_SERVING + } + h.grpc.SetServingStatus("", status) + h.grpc.SetServingStatus("unitdb.Unitdb", status) +} + +// drain stops the server being ready: SIGTERM, or Close. +func (h *healthMonitor) drain() { + h.draining.Store(true) + h.updateGRPC() +} + +func (h *healthMonitor) handler() http.Handler { + mux := http.NewServeMux() + mux.HandleFunc("/_healthz", func(w http.ResponseWriter, r *http.Request) { + w.Write([]byte("ok")) + }) + mux.HandleFunc("/_readyz", func(w http.ResponseWriter, r *http.Request) { + if why := h.notReady(); why != "" { + http.Error(w, why, http.StatusServiceUnavailable) + return + } + w.Write([]byte("ready")) + }) + mux.HandleFunc("/_status", func(w http.ResponseWriter, r *http.Request) { + why := h.notReady() + status := "ready" + switch { + case h.draining.Load(): + status = "draining" + case why != "": + status = "not_ready" + } + h.mu.RLock() + checks := make(map[string]healthResult, len(h.results)) + for k, v := range h.results { + checks[k] = v + } + h.mu.RUnlock() + names := make([]string, 0, len(checks)) + for k := range checks { + names = append(names, k) + } + sort.Strings(names) + body := map[string]interface{}{ + "status": status, + "service": "unitdb", + "started": h.started.UTC().Format(time.RFC3339), + "checks": checks, + } + if why != "" { + body["reason"] = why + } + w.Header().Set("content-type", "application/json") + enc := json.NewEncoder(w) + enc.SetIndent("", " ") + enc.Encode(body) + }) + return mux +} + +// serve serves the checks on addr, until close. +func (h *healthMonitor) serve(addr string) error { + l, err := net.Listen("tcp", addr) + if err != nil { + return err + } + h.srv = &http.Server{Handler: h.handler(), ReadHeaderTimeout: 5 * time.Second} + go func() { + if err := h.srv.Serve(l); err != nil && !errors.Is(err, http.ErrServerClosed) { + log.Error("health", "serve: "+err.Error()) + } + }() + log.Info("health", "checks served at "+addr) + return nil +} + +func (h *healthMonitor) close() { + if h.srv != nil { + h.srv.Close() + } +} + +// addServiceChecks adds the server's checks: the store, and in a cluster, +// this node's place in it. +func (h *healthMonitor) addServiceChecks() { + h.add("store", func() (string, error) { + if !store.IsOpen() { + return "", errors.New("the store is not open") + } + // A write and a read of the store's own key, which nothing + // replicates. + if err := store.Probe(); err != nil { + return "", err + } + return store.GetAdapterName() + " open, takes a write", nil + }) + if Globals.Cluster != nil { + h.add("cluster", func() (string, error) { + ok, detail := Globals.Cluster.readiness() + if !ok { + return "", errors.New(detail) + } + return detail, nil + }) + } +} diff --git a/server/internal/health_test.go b/server/internal/health_test.go new file mode 100644 index 0000000..6badb53 --- /dev/null +++ b/server/internal/health_test.go @@ -0,0 +1,148 @@ +package internal + +import ( + "context" + "encoding/json" + "errors" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + healthpb "google.golang.org/grpc/health/grpc_health_v1" +) + +func get(t *testing.T, h http.Handler, path string) (int, string) { + t.Helper() + w := httptest.NewRecorder() + h.ServeHTTP(w, httptest.NewRequest("GET", path, nil)) + return w.Code, strings.TrimSpace(w.Body.String()) +} + +func grpcStatus(t *testing.T, h *healthMonitor) healthpb.HealthCheckResponse_ServingStatus { + t.Helper() + r, err := h.grpc.Check(context.Background(), &healthpb.HealthCheckRequest{Service: "unitdb.Unitdb"}) + if err != nil { + t.Fatalf("grpc health: %v", err) + } + return r.Status +} + +func TestHealthReadiness(t *testing.T) { + h := newHealthMonitor(time.Now()) + storeOK := true + h.add("store", func() (string, error) { + if !storeOK { + return "", errors.New("a read failed") + } + return "open", nil + }) + web := h.handler() + + // Liveness never depends on the checks. + if code, body := get(t, web, "/_healthz"); code != 200 || body != "ok" { + t.Errorf("/_healthz: %d %q", code, body) + } + if code, body := get(t, web, "/_readyz"); code != 503 || body != "store: not checked yet" { + t.Errorf("/_readyz before a check: %d %q", code, body) + } + + h.runAll() + if code, body := get(t, web, "/_readyz"); code != 200 || body != "ready" { + t.Errorf("/_readyz: %d %q", code, body) + } + if s := grpcStatus(t, h); s != healthpb.HealthCheckResponse_SERVING { + t.Errorf("grpc: %v", s) + } + + storeOK = false + h.runAll() + if code, body := get(t, web, "/_readyz"); code != 503 || body != "store: a read failed" { + t.Errorf("/_readyz with the store failing: %d %q", code, body) + } + if s := grpcStatus(t, h); s != healthpb.HealthCheckResponse_NOT_SERVING { + t.Errorf("grpc with the store failing: %v", s) + } + if code, _ := get(t, web, "/_healthz"); code != 200 { + t.Errorf("/_healthz with the store failing: %d", code) + } + + // The detail, as JSON. + _, body := get(t, web, "/_status") + var status struct { + Status string `json:"status"` + Reason string `json:"reason"` + Checks map[string]healthResult `json:"checks"` + } + if err := json.Unmarshal([]byte(body), &status); err != nil { + t.Fatalf("/_status: %v: %s", err, body) + } + if status.Status != "not_ready" || status.Reason != "store: a read failed" || + status.Checks["store"].OK || status.Checks["store"].CheckedAt.IsZero() { + t.Errorf("/_status: %s", body) + } + + storeOK = true + h.runAll() + h.drain() + if code, body := get(t, web, "/_readyz"); code != 503 || body != "draining" { + t.Errorf("/_readyz draining: %d %q", code, body) + } + if s := grpcStatus(t, h); s != healthpb.HealthCheckResponse_NOT_SERVING { + t.Errorf("grpc draining: %v", s) + } +} + +func TestHealthCheckTimeout(t *testing.T) { + h := newHealthMonitor(time.Now()) + block := make(chan struct{}) + defer close(block) + h.add("slow", func() (string, error) { + <-block + return "late", nil + }) + start := time.Now() + h.runAll() + if d := time.Since(start); d > healthTimeout+time.Second { + t.Errorf("runAll took %s", d) + } + if why := h.notReady(); !strings.HasPrefix(why, "slow: no answer in") { + t.Errorf("notReady: %q", why) + } +} + +func TestClusterReadiness(t *testing.T) { + var none *Cluster + if ok, detail := none.readiness(); !ok || detail != "standalone" { + t.Errorf("no cluster: %v %q", ok, detail) + } + + c := &Cluster{thisNodeName: "a", nodes: map[string]*ClusterNode{"b": {}, "c": {}}} + c.fo = &clusterFailover{heartBeat: 100 * time.Millisecond, voteTimeout: 3} + if ok, detail := c.readiness(); ok || !strings.Contains(detail, "not in the ring") { + t.Errorf("before joining: %v %q", ok, detail) + } + c.ringNodes = []string{"a", "b", "c"} + if ok, detail := c.readiness(); ok || detail != "no leader heard from yet" { + t.Errorf("no leader: %v %q", ok, detail) + } + c.health.leaderSeen("b") + if ok, detail := c.readiness(); !ok || !strings.Contains(detail, "leader b") { + t.Errorf("ready: %v %q", ok, detail) + } + c.health.lastLeader.Store(time.Now().Add(-time.Second).UnixNano()) + if ok, detail := c.readiness(); ok || !strings.HasPrefix(detail, "no leader for") { + t.Errorf("a stale leader: %v %q", ok, detail) + } + c.health.leaderSeen("b") + c.rebuilding.Store(true) + if ok, detail := c.readiness(); ok || !strings.Contains(detail, "catching up") { + t.Errorf("rebuilding: %v %q", ok, detail) + } + c.rebuilding.Store(false) + c.leaving.Store(true) + if ok, detail := c.readiness(); ok || detail != "leaving the cluster" { + t.Errorf("leaving: %v %q", ok, detail) + } +} diff --git a/server/internal/net/hdl_grpc.go b/server/internal/net/hdl_grpc.go index cea0b76..6274a9c 100644 --- a/server/internal/net/hdl_grpc.go +++ b/server/internal/net/hdl_grpc.go @@ -92,6 +92,9 @@ func (s *GrpcServer) Serve(list net.Listener) error { srv := grpc.NewServer(opts...) pbx.RegisterUnitdbServer(srv, unitdbService{s: s}) + if s.Register != nil { + s.Register(srv) + } s.Lock() s.stop = srv.Stop s.Unlock() diff --git a/server/internal/net/server.go b/server/internal/net/server.go index ae3a06d..333f67d 100644 --- a/server/internal/net/server.go +++ b/server/internal/net/server.go @@ -25,6 +25,8 @@ import ( "os/signal" "sync" "syscall" + + "google.golang.org/grpc" ) const ( @@ -106,6 +108,9 @@ type server struct { opts *options stop func() // stops serving, if the server supports it Handler Handler //The handler to invoke when a connection is accepted + // Register, for a gRPC server, registers more services on it before it + // serves, such as grpc.health.v1. + Register func(*grpc.Server) } func signalHandler() <-chan bool { diff --git a/server/internal/service.go b/server/internal/service.go index 6ed3603..a9aff5e 100644 --- a/server/internal/service.go +++ b/server/internal/service.go @@ -29,6 +29,9 @@ import ( "syscall" "time" + "google.golang.org/grpc" + healthpb "google.golang.org/grpc/health/grpc_health_v1" + "github.com/unit-io/unitdb/server/internal/config" "github.com/unit-io/unitdb/server/internal/keys" lp "github.com/unit-io/unitdb/server/internal/net" @@ -63,6 +66,7 @@ type _Service struct { tcp *lp.TcpServer // The underlying TCP server. grpc *lp.GrpcServer // The underlying GRPC server. meter *Meter // The metircs to measure timeseries on message events + health *healthMonitor // The health checks, on their own port stats *stats.Stats // Shutdown. @@ -206,6 +210,19 @@ func (s *_Service) listen(addr string) { l.SetReadTimeout(120 * time.Second) + // The health checks, and grpc.health.v1 on the gRPC server. + s.health = newHealthMonitor(s.start) + s.health.addServiceChecks() + s.grpc.Register = func(g *grpc.Server) { + healthpb.RegisterHealthServer(g, s.health.grpc) + } + if s.config.MonitorListen != "" { + if err := s.health.serve(s.config.MonitorListen); err != nil { + panic(err) + } + } + go s.health.run(s.context) + // Configure the protos if s.config.GrpcListen != "" { grpcList, err := netListener(s.config.GrpcListen) @@ -274,6 +291,12 @@ func (s *_Service) Close() { } func (s *_Service) close() { + // Not ready from now on: the load balancer moves away. + if s.health != nil { + s.health.drain() + defer s.health.close() + } + // Leave the cluster first, while this node's clients are still served: // the others take over what it holds, and its clients' subscriptions. Globals.Cluster.drain() diff --git a/server/internal/store/probe_test.go b/server/internal/store/probe_test.go new file mode 100644 index 0000000..5296c81 --- /dev/null +++ b/server/internal/store/probe_test.go @@ -0,0 +1,21 @@ +package store_test + +import ( + "testing" + + "github.com/unit-io/unitdb/server/internal/store" +) + +// TestProbe checks that the health probe writes and reads back, repeatedly, +// sealed and not. +func TestProbe(t *testing.T) { + for _, seal := range []bool{false, true} { + newStore(t, 0, ring(0), seal) + for i := 0; i < 3; i++ { + if err := store.Probe(); err != nil { + t.Fatalf("sealed %v: %v", seal, err) + } + } + store.Close() + } +} diff --git a/server/internal/store/security.go b/server/internal/store/security.go index b246e16..c32f95b 100644 --- a/server/internal/store/security.go +++ b/server/internal/store/security.go @@ -16,6 +16,14 @@ package store +import ( + "bytes" + "crypto/rand" + "errors" + "fmt" + "hash/fnv" +) + // securityTopic is the topic the cluster's security state is kept under, in // the node's own namespace (namespaces.go). var securityTopic = sysTopic(sysSecurity, "state") @@ -68,3 +76,34 @@ func (SecurityStore) Legacy() (ids, records [][]byte, err error) { func (SecurityStore) DeleteLegacy(id []byte) error { return adp.Delete(legacySecurityStoreId, id, legacySecurityTopic) } + +// probeKey is the memdb key the health probe writes: a hash of a string no +// other record's key comes from. +var probeKey = func() uint64 { + h := fnv.New64a() + h.Write([]byte("\x00unitdb-health-probe")) + return h.Sum64() +}() + +// Probe writes a record and reads it back, for health checks. It writes this +// node's store only, through the adapter, so nothing replicates it, and under +// a key of its own, so it can't clash with a real one. Sealing applies, as +// to every record. +func Probe() error { + var b [16]byte + if _, err := rand.Read(b[:]); err != nil { + return err + } + adp.DeleteMessage(probeKey) + if err := adp.PutMessage(probeKey, b[:]); err != nil { + return fmt.Errorf("store probe: write: %v", err) + } + got, err := adp.GetMessage(probeKey) + if err != nil { + return fmt.Errorf("store probe: read: %v", err) + } + if !bytes.Equal(got, b[:]) { + return errors.New("store probe: read back something else than it wrote") + } + return nil +} diff --git a/server/main.go b/server/main.go index bef80f6..7ad9923 100644 --- a/server/main.go +++ b/server/main.go @@ -40,6 +40,7 @@ func main() { var configfile = flag.String("config", "unitdb.conf", "Path to config file.") var listenOn = flag.String("listen", "", "Override address and port to listen on for HTTP(S) clients.") var listenGrpcOn = flag.String("grpc_listen", "", "Override address and port to listen on for GRPC clients.") + var monitorOn = flag.String("monitor_listen", "", "Override address and port to serve the health checks on (/_healthz, /_readyz, /_status).") var clusterSelf = flag.String("cluster_self", "", "Override the name of the current cluster node") var dbPath = flag.String("db_path", "/tmp/unitdb", "Override the db path.") var varzPath = flag.String("varz", "/varz", "Expose runtime stats at the given endpoint, e.g. /varz. Disabled if not set") @@ -73,6 +74,10 @@ func main() { cfg.GrpcListen = *listenGrpcOn } + if *monitorOn != "" { + cfg.MonitorListen = *monitorOn + } + if *dbPath != "" { cfg.DBPath = *dbPath }