feat(watchdog): 主机侧巡检定时器按异常邮件通知平台所有者,operator 增加 phase 与 build_info 指标、卡死存活探针与告警规则

This commit is contained in:
Lemon-miaow committed 2026-09-24 18:06:19 +08:00
1 parent d50492b86f
commit d17524cd67
26 files changed
+2686 -34

No files matched your search

+98
View File
@@ -0,0 +1,98 @@
package metrics
import (
"os"
"reflect"
"regexp"
"strings"
"testing"
"github.com/prometheus/client_golang/prometheus"
"sigs.k8s.io/yaml"
)
// The shipped alert rules live in deploy/alerts. promtool tests what they do
// (felis-alerts_test.yml); these tests pin what promtool cannot see from there.
type ruleGroups struct {
Groups []struct {
Name string `json:"name"`
Rules []struct {
Alert string `json:"alert"`
Expr string `json:"expr"`
} `json:"rules"`
} `json:"groups"`
}
func readYAML(t *testing.T, path string, into any) {
t.Helper()
raw, err := os.ReadFile(path)
if err != nil {
t.Fatal(err)
}
if err := yaml.Unmarshal(raw, into); err != nil {
t.Fatalf("parse %s: %v", path, err)
}
}
// TestPrometheusRuleMatchesPlainRules: the PrometheusRule twin carries exactly
// the groups of the promtool-tested plain file.
func TestPrometheusRuleMatchesPlainRules(t *testing.T) {
var plain, twin map[string]any
readYAML(t, "../../deploy/alerts/felis-alerts.yaml", &plain)
readYAML(t, "../../deploy/alerts/felis-prometheusrule.yaml", &twin)
spec, _ := twin["spec"].(map[string]any)
if !reflect.DeepEqual(plain["groups"], spec["groups"]) {
t.Fatal("deploy/alerts/felis-prometheusrule.yaml spec.groups differs from felis-alerts.yaml groups; copy the plain file's groups over")
}
}
// TestAlertRulesUseRealMetrics: every felis_* series a rule reads is one this
// package exports (or the backup timer's textfile series), so a renamed metric
// cannot leave an alert silently watching nothing.
func TestAlertRulesUseRealMetrics(t *testing.T) {
reg := prometheus.NewRegistry()
if err := Register(reg); err != nil {
t.Fatal(err)
}
families, err := reg.Gather()
if err != nil {
t.Fatal(err)
}
known := map[string]bool{}
for _, mf := range families {
known[mf.GetName()] = true
}
// Collectors with no children yet gather nothing; name them from their
// descriptors instead.
for _, c := range Collectors() {
ch := make(chan *prometheus.Desc, 4)
c.Describe(ch)
close(ch)
for d := range ch {
if m := regexp.MustCompile(`fqName: "([^"]+)"`).FindStringSubmatch(d.String()); m != nil {
known[m[1]] = true
}
}
}
var rules ruleGroups
readYAML(t, "../../deploy/alerts/felis-alerts.yaml", &rules)
series := regexp.MustCompile(`felis_[a-z_]+`)
for _, g := range rules.Groups {
for _, r := range g.Rules {
for _, name := range series.FindAllString(r.Expr, -1) {
if strings.HasPrefix(name, "felis_db_backup_") {
continue // written by felis-db-backup.timer for node-exporter
}
base := name
for _, suffix := range []string{"_bucket", "_count", "_sum"} {
base = strings.TrimSuffix(base, suffix)
}
if !known[name] && !known[base] {
t.Errorf("alert %s reads %s, which no felis collector exports", r.Alert, name)
}
}
}
}
}
+45
View File
@@ -32,6 +32,18 @@ var (
Help: "Current number of Minecraft servers known to the operator, by desired state.",
}, []string{"state"})
// ServerPhase is 1 for each server's current phase and absent for every other
// phase. role is the server's system role (login, lobby), empty for a user
// server, so an alert can single out the login gate; desired is its
// desiredState, so a server that is down on purpose can be told from one that
// failed to come up. SyncServerPhases republishes it from a full List, so a
// deleted server's series goes away instead of freezing at its last phase.
ServerPhase = prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: namespace,
Name: "server_phase",
Help: "1 for each Minecraft server's current phase, by server, system role, phase and desired state.",
}, []string{"server", "role", "phase", "desired"})
// StartDurationSeconds observes the wall-clock time from desiredState=Running
// to a server reporting ready. Buckets are tuned for Minecraft cold starts
// (seconds to a few minutes), not the default sub-second web-latency buckets.
@@ -107,8 +119,25 @@ var (
Name: "audit_write_failures_total",
Help: "Audit rows the API failed to write.",
})
// BuildInfo is 1 for the process serving it, labelled by component
// ("operator", "api") and version. Both processes register every collector,
// so this is the one series that says which of them a scrape reached: an
// alert on absent(felis_build_info{component="api"}) fires when felis-api is
// down or no longer scraped, where every other felis_* series would still be
// present from the operator.
BuildInfo = prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: namespace,
Name: "build_info",
Help: "1 for the Felis component serving these metrics, by component and version.",
}, []string{"component", "version"})
)
// SetBuildInfo marks this process as component at version on felis_build_info.
func SetBuildInfo(component, version string) {
BuildInfo.WithLabelValues(component, version).Set(1)
}
// OTPPurposes are the email-code doors OTPLockoutsTotal is labelled by.
var OTPPurposes = []string{"onboard_email", "login_email", "op_login", "migrate_confirm"}
@@ -149,12 +178,27 @@ func SyncServerGauge(states []string) {
}
}
// ServerPhaseSample is one server's felis_server_phase series.
type ServerPhaseSample struct {
Server, Role, Phase, Desired string
}
// SyncServerPhases republishes felis_server_phase from a full snapshot of the
// fleet, Resetting first for the same reason SyncServerGauge does.
func SyncServerPhases(samples []ServerPhaseSample) {
ServerPhase.Reset()
for _, s := range samples {
ServerPhase.WithLabelValues(s.Server, s.Role, s.Phase, s.Desired).Set(1)
}
}
// Collectors returns every felis_* collector in a stable order. Production and
// tests register the same slice, so the test asserting the full set is exposed
// also pins the production surface.
func Collectors() []prometheus.Collector {
return []prometheus.Collector{
ServersTotal,
ServerPhase,
StartDurationSeconds,
ImageBuildFailuresTotal,
ReaperWorldsDeletedTotal,
@@ -164,6 +208,7 @@ func Collectors() []prometheus.Collector {
AuthFailuresTotal,
SessionsRevokedTotal,
AuditWriteFailuresTotal,
BuildInfo,
}
}
+24
View File
@@ -13,9 +13,11 @@ import (
// silent dashboard/alert break.
var wantNames = []string{
"felis_servers_total",
"felis_server_phase",
"felis_start_duration_seconds",
"felis_image_build_failures_total",
"felis_reaper_worlds_deleted_total",
"felis_build_info",
}
func TestRegisterExposesNamedFelisMetrics(t *testing.T) {
@@ -28,6 +30,8 @@ func TestRegisterExposesNamedFelisMetrics(t *testing.T) {
// family — this proves the exported vars are the ones actually registered,
// not shadow copies.
ServersTotal.WithLabelValues("Running").Set(3)
ServerPhase.WithLabelValues("survival", "", "Running", "Running").Set(1)
SetBuildInfo("operator", "v1.2.3")
StartDurationSeconds.Observe(12.5)
ImageBuildFailuresTotal.Inc()
ReaperWorldsDeletedTotal.Add(2)
@@ -54,6 +58,26 @@ func TestRegisterExposesNamedFelisMetrics(t *testing.T) {
}
}
// TestSyncServerPhasesDropsDeletedServers: each sync publishes exactly the
// servers in the snapshot, so a server that left the fleet or changed phase keeps
// no stale series behind.
func TestSyncServerPhasesDropsDeletedServers(t *testing.T) {
SyncServerPhases([]ServerPhaseSample{
{Server: "login", Role: "login", Phase: "Starting", Desired: "Running"},
{Server: "survival", Phase: "Running", Desired: "Running"},
})
SyncServerPhases([]ServerPhaseSample{
{Server: "login", Role: "login", Phase: "Running", Desired: "Running"},
})
if n := testutil.CollectAndCount(ServerPhase); n != 1 {
t.Fatalf("series after the second sync = %d, want 1", n)
}
if got := testutil.ToFloat64(ServerPhase.WithLabelValues("login", "login", "Running", "Running")); got != 1 {
t.Errorf("login Running = %v, want 1", got)
}
ServerPhase.Reset()
}
func TestSyncServerGaugeResetsStaleStates(t *testing.T) {
// Self-contained (Reset->sync->assert within this one test) so it never
// clobbers another test's ServersTotal children. The correctness point is the
+30 -3
View File
@@ -13,8 +13,8 @@ import (
// defaultGaugeInterval is the republish cadence used when none is configured.
const defaultGaugeInterval = 30 * time.Second
// GaugeSyncer periodically republishes felis_servers_total (spec §23) from a
// full List of MinecraftServers.
// GaugeSyncer periodically republishes felis_servers_total (spec §23) and
// felis_server_phase from a full List of MinecraftServers.
//
// felis_servers_total is a fleet-wide gauge partitioned by desiredState, which a
// per-object Reconcile fundamentally cannot maintain: one reconcile observes a
@@ -62,7 +62,8 @@ func (g *GaugeSyncer) Start(ctx context.Context) error {
}
}
// SyncOnce Lists the fleet once and republishes felis_servers_total from it.
// SyncOnce Lists the fleet once and republishes felis_servers_total and
// felis_server_phase from it.
// Separated from Start so the List->translate->gauge path is unit-testable
// against a fake client, leaving only the ticker loop untested.
func (g *GaugeSyncer) SyncOnce(ctx context.Context) error {
@@ -71,9 +72,35 @@ func (g *GaugeSyncer) SyncOnce(ctx context.Context) error {
return err
}
metrics.SyncServerGauge(serverStates(list.Items))
metrics.SyncServerPhases(serverPhases(list.Items))
return nil
}
// serverPhases maps each server to its felis_server_phase labels. A server the
// operator has not reported on yet is Unknown; desiredState defaults as in
// serverStates.
func serverPhases(items []v1alpha1.MinecraftServer) []metrics.ServerPhaseSample {
out := make([]metrics.ServerPhaseSample, 0, len(items))
for i := range items {
ms := &items[i]
phase := string(ms.Status.Phase)
if phase == "" {
phase = string(v1alpha1.PhaseUnknown)
}
desired := string(ms.Spec.DesiredState)
if desired == "" {
desired = string(v1alpha1.DesiredStopped)
}
out = append(out, metrics.ServerPhaseSample{
Server: ms.Name,
Role: ms.Labels[v1alpha1.LabelSystemRole],
Phase: phase,
Desired: desired,
})
}
return out
}
// serverStates maps each server to its desiredState, defaulting an unset state
// to Stopped (the CRD default, spec §4). Pure, so the labelling rule is covered
// by SyncOnce's fake-client test rather than only at runtime.
+27
View File
@@ -48,3 +48,30 @@ func TestGaugeSyncerPublishesFleetStateFromCluster(t *testing.T) {
t.Errorf("servers_total{state=Stopped} = %v, want 2", got)
}
}
// TestGaugeSyncerPublishesServerPhases: every server gets one felis_server_phase
// series carrying its system role, observed phase (Unknown before the first
// report) and desired state.
func TestGaugeSyncerPublishesServerPhases(t *testing.T) {
login := gaugeServer("login", v1alpha1.DesiredRunning)
login.Labels = map[string]string{v1alpha1.LabelSystemRole: "login"}
login.Status.Phase = v1alpha1.PhaseFailed
fresh := gaugeServer("fresh", "")
c := fake.NewClientBuilder().WithScheme(newScheme(t)).WithObjects(login, fresh).Build()
if err := (&operator.GaugeSyncer{Client: c}).SyncOnce(context.Background()); err != nil {
t.Fatalf("SyncOnce: %v", err)
}
defer metrics.ServerPhase.Reset()
if n := testutil.CollectAndCount(metrics.ServerPhase); n != 2 {
t.Fatalf("felis_server_phase series = %d, want 2", n)
}
for _, labels := range [][]string{
{"login", "login", "Failed", "Running"},
{"fresh", "", "Unknown", "Stopped"},
} {
if got := testutil.ToFloat64(metrics.ServerPhase.WithLabelValues(labels...)); got != 1 {
t.Errorf("felis_server_phase%v = %v, want 1", labels, got)
}
}
}
+13 -1
View File
@@ -78,6 +78,8 @@ type Reconciler struct {
// (internal/maintenance). It is the manager's uncached API reader, so the
// operator needs jobs:list and no Job informer. Nil skips the check.
Jobs client.Reader
// Watch records the passes in flight for the liveness probe. Nil skips it.
Watch *ReconcileWatch
}
func (r *Reconciler) now() metav1.Time {
@@ -109,8 +111,18 @@ func (r *Reconciler) SetupWithManager(mgr ctrl.Manager) error {
// The same server is never reconciled twice at once regardless.
const maxConcurrentReconciles = 4
// Reconcile drives a single MinecraftServer toward spec.desiredState.
// Reconcile drives a single MinecraftServer toward spec.desiredState. One pass
// is bounded by reconcileTimeout and reported to r.Watch while it runs.
func (r *Reconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
ctx, cancel := context.WithTimeout(ctx, reconcileTimeout)
defer cancel()
if r.Watch != nil {
defer r.Watch.begin(req.Name)()
}
return r.reconcile(ctx, req)
}
func (r *Reconciler) reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
var server v1alpha1.MinecraftServer
if err := r.Get(ctx, req.NamespacedName, &server); err != nil {
// Deletion is handled by owner references on the children.
+91
View File
@@ -0,0 +1,91 @@
package operator
import (
"fmt"
"net/http"
"sync"
"time"
)
// reconcileTimeout bounds one reconcile pass. The slowest legitimate pass waits
// on RCON: 5s for a probe, defaultSaveTimeout for the world save ahead of a stop.
// A pass still running past this is waiting on something that will not answer;
// cancelling its context makes the calls that honour it return, and the server
// is retried with backoff.
const reconcileTimeout = 3 * time.Minute
// stuckReconcileAfter is how long one pass may run before the operator reports
// itself unhealthy. It is well past reconcileTimeout, so only a pass blocked in a
// call that ignores its context gets here. Each such pass holds one of the
// maxConcurrentReconciles workers for good; a few of them stall every server's
// start and stop while the Pod still looks healthy. Failing the liveness probe
// gets the operator restarted, the one fix for a goroutine that never returns.
const stuckReconcileAfter = 10 * time.Minute
// ReconcileWatch tracks the reconcile passes in flight so the liveness probe can
// tell a busy operator from a wedged one. An idle operator has nothing in flight
// and is healthy: a quiet fleet reconciles rarely, so the time since the last
// pass says nothing about whether the next one would run.
type ReconcileWatch struct {
// StuckAfter defaults to stuckReconcileAfter.
StuckAfter time.Duration
// Now defaults to time.Now.
Now func() time.Time
mu sync.Mutex
next uint64
inflight map[uint64]inflightPass
}
type inflightPass struct {
server string
started time.Time
}
func (w *ReconcileWatch) now() time.Time {
if w.Now != nil {
return w.Now()
}
return time.Now()
}
// begin records a pass for server and returns the func that ends it.
func (w *ReconcileWatch) begin(server string) func() {
w.mu.Lock()
defer w.mu.Unlock()
if w.inflight == nil {
w.inflight = map[uint64]inflightPass{}
}
w.next++
id := w.next
w.inflight[id] = inflightPass{server: server, started: w.now()}
return func() {
w.mu.Lock()
delete(w.inflight, id)
w.mu.Unlock()
}
}
// Check is a healthz.Checker: it fails while any pass has been running longer
// than StuckAfter, naming the server and how long.
func (w *ReconcileWatch) Check(_ *http.Request) error {
limit := w.StuckAfter
if limit <= 0 {
limit = stuckReconcileAfter
}
now := w.now()
w.mu.Lock()
defer w.mu.Unlock()
var oldest *inflightPass
for id := range w.inflight {
p := w.inflight[id]
if oldest == nil || p.started.Before(oldest.started) {
oldest = &p
}
}
if oldest != nil && now.Sub(oldest.started) > limit {
return fmt.Errorf("reconcile of %s has been running for %s (limit %s)",
oldest.server, now.Sub(oldest.started).Truncate(time.Second), limit)
}
return nil
}
+83
View File
@@ -0,0 +1,83 @@
package operator
import (
"context"
"strings"
"testing"
"time"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
"sigs.k8s.io/controller-runtime/pkg/client/interceptor"
)
// TestReconcileWatchCheck: an idle or busy operator is healthy; one pass running
// past the limit fails the check and names the server; the check recovers once
// that pass returns.
func TestReconcileWatchCheck(t *testing.T) {
now := time.Unix(1_700_000_000, 0)
w := &ReconcileWatch{StuckAfter: 10 * time.Minute, Now: func() time.Time { return now }}
if err := w.Check(nil); err != nil {
t.Fatalf("idle: %v", err)
}
endOld := w.begin("survival")
now = now.Add(9 * time.Minute)
endNew := w.begin("creative")
if err := w.Check(nil); err != nil {
t.Fatalf("9m in: %v", err)
}
now = now.Add(2 * time.Minute)
err := w.Check(nil)
if err == nil || !strings.Contains(err.Error(), "survival") || !strings.Contains(err.Error(), "11m0s") {
t.Fatalf("11m in: err = %v, want survival stuck for 11m0s", err)
}
endOld()
if err := w.Check(nil); err != nil {
t.Fatalf("after the stuck pass returned: %v", err)
}
endNew()
}
// TestReconcileBoundedAndWatched: Reconcile hands the pass a deadline and holds
// an in-flight entry exactly while it runs.
func TestReconcileBoundedAndWatched(t *testing.T) {
scheme := runtime.NewScheme()
if err := v1alpha1.AddToScheme(scheme); err != nil {
t.Fatal(err)
}
watch := &ReconcileWatch{}
var deadline time.Duration
var inflight int
cl := fake.NewClientBuilder().WithScheme(scheme).WithInterceptorFuncs(interceptor.Funcs{
Get: func(ctx context.Context, c client.WithWatch, key client.ObjectKey, obj client.Object, opts ...client.GetOption) error {
if d, ok := ctx.Deadline(); ok && deadline == 0 {
deadline = time.Until(d)
}
watch.mu.Lock()
inflight = len(watch.inflight)
watch.mu.Unlock()
return c.Get(ctx, key, obj, opts...)
},
}).Build()
r := &Reconciler{Client: cl, Scheme: scheme, Watch: watch}
if _, err := r.Reconcile(context.Background(), ctrl.Request{NamespacedName: types.NamespacedName{Namespace: "minecraft", Name: "missing"}}); err != nil {
t.Fatalf("Reconcile: %v", err)
}
if deadline <= 0 || deadline > reconcileTimeout {
t.Errorf("pass deadline = %s, want within %s", deadline, reconcileTimeout)
}
if inflight != 1 {
t.Errorf("in-flight passes during the pass = %d, want 1", inflight)
}
if n := len(watch.inflight); n != 0 {
t.Errorf("in-flight passes after the pass = %d, want 0", n)
}
}
+45 -4
View File
@@ -173,6 +173,21 @@ const (
// DNS; the on-node break-glass console resolves its ClusterIP and dials it directly.
const APIInternalServiceName = SAAPI + "-internal"
// OperatorMetricsServiceName is the ClusterIP Service in front of the operator's
// /metrics, the scrape target for felis_servers_total, felis_server_phase,
// felis_start_duration_seconds and controller-runtime's reconcile series.
const OperatorMetricsServiceName = SAOperator + "-metrics"
// scrapeAnnotations mark a Service for a Prometheus that discovers targets with
// the common prometheus.io/* convention (kubernetes_sd role: endpoints).
func scrapeAnnotations(port int32) map[string]string {
return map[string]string{
"prometheus.io/scrape": "true",
"prometheus.io/port": fmt.Sprint(port),
"prometheus.io/path": "/metrics",
}
}
// APIInternalPort is the felis-api internal-face port, exported for the on-node
// console which builds http://<clusterIP>:APIInternalPort after a Service lookup.
const APIInternalPort = apiInternalPort
@@ -206,6 +221,7 @@ func Workloads(p Params) []Object {
apiService(p),
apiInternalService(p),
OperatorDeployment(p),
operatorMetricsService(p),
registryDeployment(p),
registryService(p),
registryPVC(p),
@@ -430,8 +446,10 @@ func apiInternalService(p Params) *corev1.Service {
p = p.withDefaults()
labels := controlPlanePodLabels(ComponentAPI)
return &corev1.Service{
TypeMeta: metav1.TypeMeta{APIVersion: "v1", Kind: "Service"},
ObjectMeta: metav1.ObjectMeta{Name: APIInternalServiceName, Namespace: p.ControlNamespace, Labels: labels},
TypeMeta: metav1.TypeMeta{APIVersion: "v1", Kind: "Service"},
// The internal face also serves /metrics (felis-api's felis_* series).
ObjectMeta: metav1.ObjectMeta{Name: APIInternalServiceName, Namespace: p.ControlNamespace, Labels: labels,
Annotations: scrapeAnnotations(apiInternalPort)},
Spec: corev1.ServiceSpec{
Type: corev1.ServiceTypeClusterIP,
Selector: labels,
@@ -445,6 +463,28 @@ func apiInternalService(p Params) *corev1.Service {
}
}
// operatorMetricsService gives the operator's metrics listener a stable scrape
// target. ClusterIP only: /metrics is unauthenticated, like felis-api's.
func operatorMetricsService(p Params) *corev1.Service {
p = p.withDefaults()
labels := controlPlanePodLabels(ComponentOperator)
return &corev1.Service{
TypeMeta: metav1.TypeMeta{APIVersion: "v1", Kind: "Service"},
ObjectMeta: metav1.ObjectMeta{Name: OperatorMetricsServiceName, Namespace: p.ControlNamespace, Labels: labels,
Annotations: scrapeAnnotations(operatorMetricsPort)},
Spec: corev1.ServiceSpec{
Type: corev1.ServiceTypeClusterIP,
Selector: labels,
Ports: []corev1.ServicePort{{
Name: "metrics",
Port: operatorMetricsPort,
TargetPort: intstr.FromString("metrics"),
Protocol: corev1.ProtocolTCP,
}},
},
}
}
// OperatorDeployment renders the felis-operator Deployment (spec §5). It runs as
// the felis-operator SA and carries controlPlanePodLabels(operator), the second
// pod the allow-rcon peer admits (the readiness prober dials RCON). It takes NO
@@ -478,8 +518,9 @@ func OperatorDeployment(p Params) *appsv1.Deployment {
{Name: tmpVolume, MountPath: "/tmp"},
},
// controller-runtime serves /healthz and /readyz on the health listener
// (both registered as always-pass pings in cmd/felis/operator.go): the
// probe's contract is "the manager process is up", not a dependency check.
// (cmd/felis/operator.go): /healthz fails while a reconcile pass is stuck,
// so liveness restarts a wedged operator; /readyz waits for the informer
// caches. Neither checks a dependency, so an API blip restarts nothing.
ReadinessProbe: &corev1.Probe{
ProbeHandler: corev1.ProbeHandler{HTTPGet: &corev1.HTTPGetAction{
Path: "/readyz", Port: intstr.FromInt32(operatorHealthPort),
+36 -3
View File
@@ -608,13 +608,14 @@ func TestWorkloads_DeploymentsCarryProbes(t *testing.T) {
}
// TestWorkloads_BundleContents sanity-checks the slice Workloads returns: the two
// control-plane Deployments, the api external+internal Services, and the registry
// control-plane Deployments, the api external+internal Services, the operator
// metrics Service, and the registry
// Deployment/Service/PVC, every one with TypeMeta (so its YAML header renders). The
// internal Service must be present or the login pod's felis-api:8081 path is dead.
func TestWorkloads_BundleContents(t *testing.T) {
objs := Workloads(testParams())
if len(objs) != 8 {
t.Fatalf("Workloads returned %d objects, want 8", len(objs))
if len(objs) != 9 {
t.Fatalf("Workloads returned %d objects, want 9", len(objs))
}
var haveInternalSvc bool
for _, o := range objs {
@@ -1014,3 +1015,35 @@ func TestWorldsRootStaticPV(t *testing.T) {
}
}
}
// TestOperatorMetricsService: the operator's /metrics has a ClusterIP scrape
// target selecting the operator pods on their "metrics" port, marked for
// prometheus.io discovery like the api's internal face.
func TestOperatorMetricsService(t *testing.T) {
p := testParams()
svc := operatorMetricsService(p)
dep := OperatorDeployment(p)
if svc.Name != OperatorMetricsServiceName || svc.Namespace != p.ControlNamespace || svc.Spec.Type != corev1.ServiceTypeClusterIP {
t.Fatalf("Service = %s/%s %s, want %s/%s ClusterIP", svc.Namespace, svc.Name, svc.Spec.Type, p.ControlNamespace, OperatorMetricsServiceName)
}
if !mapSelectorMatches(svc.Spec.Selector, dep.Spec.Template.Labels) {
t.Errorf("selector %v does not select operator pod labels %v", svc.Spec.Selector, dep.Spec.Template.Labels)
}
if len(svc.Spec.Ports) != 1 || svc.Spec.Ports[0].TargetPort.StrVal != "metrics" || svc.Spec.Ports[0].Port != operatorMetricsPort {
t.Errorf("ports = %+v, want %d -> metrics", svc.Spec.Ports, operatorMetricsPort)
}
for _, s := range []*corev1.Service{svc, apiInternalService(p)} {
if s.Annotations["prometheus.io/scrape"] != "true" || s.Annotations["prometheus.io/port"] != fmt.Sprint(s.Spec.Ports[0].Port) {
t.Errorf("%s annotations = %v, want prometheus.io scrape on its port", s.Name, s.Annotations)
}
}
found := false
for _, o := range Workloads(p) {
if s, ok := o.(*corev1.Service); ok && s.Name == OperatorMetricsServiceName {
found = true
}
}
if !found {
t.Error("Workloads does not render the operator metrics Service")
}
}
+425
View File
@@ -0,0 +1,425 @@
package watchdog
import (
"bufio"
"context"
"fmt"
"os"
"strconv"
"strings"
"syscall"
"time"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
"felis.lolicon.best/internal/dbbackup"
"felis.lolicon.best/internal/naming"
"felis.lolicon.best/internal/platform"
appsv1 "k8s.io/api/apps/v1"
batchv1 "k8s.io/api/batch/v1"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"sigs.k8s.io/controller-runtime/pkg/client"
)
// Thresholds. The durations are how long a condition must hold before it is
// mailed (Finding.For): long enough for a rollout, a restart or an installer
// run to heal it, short enough that players are not the first to notice.
const (
controlPlaneFor = 5 * time.Minute
systemServerFor = 10 * time.Minute
failedServerFor = 5 * time.Minute
nodeNotReadyFor = 5 * time.Minute
nodePressureFor = 2 * time.Minute
postgresFor = 3 * time.Minute
kubeAPIFor = 5 * time.Minute
backupFor = 10 * time.Minute
diskLowFor = 15 * time.Minute
diskCriticalFor = 5 * time.Minute
memoryLowFor = 15 * time.Minute
// jobFailureWindow is how far back a failed Job is still news.
jobFailureWindow = 24 * time.Hour
// maxBackupAge is the daily database backup's deadline: a day plus the
// timer's randomized delay and a margin.
maxBackupAge = 26 * time.Hour
// maxReaperAge is the same for the daily reaper CronJob.
maxReaperAge = 26 * time.Hour
diskLowRatio = 0.15
diskCriticalRatio = 0.05
memoryLowRatio = 0.10
)
// controlDeployments are the control-plane Deployments, with what their outage
// takes down.
var controlDeployments = []struct {
name, impact, impactEN string
}{
{platform.SAAPI, "面板、登录验证和内部接口都不可用", "the panel, sign-in and the internal API are down"},
{platform.SAOperator, "服务器无法启动、停止或更新状态", "no server can start, stop or report status"},
{"registry", "游戏镜像拉取和构建都会失败", "game image pulls and builds fail"},
}
// ClusterPrefixes are the keys Cluster.Check produces. When the API server
// cannot be reached, alerts under them keep their state (Report.Unknown).
var ClusterPrefixes = []string{"deployment/", "system-server/", "server-failed/", "job-failed/", "reaper-stale", "node/"}
// Cluster checks the Kubernetes side: the control plane, the system servers,
// the fleet, recent Job failures, the reaper's schedule and the nodes.
type Cluster struct {
Client client.Client
ControlNamespace string
MinecraftNamespace string
}
// Check returns the cluster's findings. An error means the API server could
// not answer; the caller reports that and treats every cluster check as unknown.
func (c Cluster) Check(ctx context.Context, now time.Time) ([]Finding, error) {
var out []Finding
for _, d := range controlDeployments {
f, err := c.deployment(ctx, d.name, d.impact, d.impactEN)
if err != nil {
return nil, err
}
if f != nil {
out = append(out, *f)
}
}
var servers v1alpha1.MinecraftServerList
if err := c.Client.List(ctx, &servers, client.InNamespace(c.MinecraftNamespace)); err != nil {
return nil, fmt.Errorf("list servers: %w", err)
}
for i := range servers.Items {
if f := serverFinding(&servers.Items[i], c.MinecraftNamespace); f != nil {
out = append(out, *f)
}
}
var jobs batchv1.JobList
if err := c.Client.List(ctx, &jobs, client.InNamespace(c.MinecraftNamespace)); err != nil {
return nil, fmt.Errorf("list jobs: %w", err)
}
for i := range jobs.Items {
if f := jobFinding(&jobs.Items[i], now); f != nil {
out = append(out, *f)
}
}
var cj batchv1.CronJob
switch err := c.Client.Get(ctx, client.ObjectKey{Namespace: c.MinecraftNamespace, Name: platform.SAReaper}, &cj); {
case apierrors.IsNotFound(err):
// Retention is off; there is no reaper to be late.
case err != nil:
return nil, fmt.Errorf("get reaper cronjob: %w", err)
default:
if f := reaperFinding(&cj, now); f != nil {
out = append(out, *f)
}
}
var nodes corev1.NodeList
if err := c.Client.List(ctx, &nodes); err != nil {
return nil, fmt.Errorf("list nodes: %w", err)
}
for i := range nodes.Items {
out = append(out, nodeFindings(&nodes.Items[i])...)
}
return out, nil
}
func (c Cluster) deployment(ctx context.Context, name, impact, impactEN string) (*Finding, error) {
var d appsv1.Deployment
err := c.Client.Get(ctx, client.ObjectKey{Namespace: c.ControlNamespace, Name: name}, &d)
if apierrors.IsNotFound(err) {
return &Finding{
Key: "deployment/" + name, Severity: Critical, For: controlPlaneFor,
Summary: fmt.Sprintf("控制面组件 %s 不存在:%s", name, impact),
SummaryEN: fmt.Sprintf("control-plane Deployment %s is missing: %s", name, impactEN),
Hint: "re-run the installer, or `sudo felis apply`, to recreate it",
}, nil
}
if err != nil {
return nil, fmt.Errorf("get deployment %s: %w", name, err)
}
if d.Status.AvailableReplicas > 0 {
return nil, nil
}
want := int32(1)
if d.Spec.Replicas != nil {
want = *d.Spec.Replicas
}
return &Finding{
Key: "deployment/" + name, Severity: Critical, For: controlPlaneFor,
Summary: fmt.Sprintf("控制面组件 %s 没有可用副本(%d/%d 就绪):%s", name, d.Status.ReadyReplicas, want, impact),
SummaryEN: fmt.Sprintf("control-plane Deployment %s has no available replica (%d/%d ready): %s", name, d.Status.ReadyReplicas, want, impactEN),
Hint: fmt.Sprintf("kubectl -n %s get pods; kubectl -n %s logs deploy/%s --all-containers", c.ControlNamespace, c.ControlNamespace, name),
}, nil
}
// serverFinding reports a system server that should run and does not, and a
// user server that failed. A user server that is merely stopped or starting is
// its owner's business.
func serverFinding(ms *v1alpha1.MinecraftServer, ns string) *Finding {
phase := ms.Status.Phase
if phase == "" {
phase = v1alpha1.PhaseUnknown
}
why := failureMessage(ms)
if role := ms.Labels[v1alpha1.LabelSystemRole]; role != "" {
if ms.Spec.DesiredState != v1alpha1.DesiredRunning || phase == v1alpha1.PhaseRunning {
return nil
}
f := &Finding{
Key: "system-server/" + ms.Name, Severity: Warning, For: systemServerFor,
Summary: fmt.Sprintf("系统服务器 %s 处于 %s(应为 Running)%s", ms.Name, phase, why),
SummaryEN: fmt.Sprintf("system server %s is %s, not Running%s", ms.Name, phase, why),
Hint: fmt.Sprintf("kubectl -n %s describe minecraftserver %s; docs/troubleshooting.md §1-§2", ns, ms.Name),
}
if role == naming.SystemLoginServer {
f.Severity = Critical
f.Summary = fmt.Sprintf("登录门 %s 处于 %s(应为 Running),玩家无法进服%s", ms.Name, phase, why)
f.SummaryEN = fmt.Sprintf("login gate %s is %s, not Running: players cannot join%s", ms.Name, phase, why)
}
return f
}
if phase != v1alpha1.PhaseFailed {
return nil
}
return &Finding{
Key: "server-failed/" + ms.Name, Severity: Warning, For: failedServerFor,
Summary: fmt.Sprintf("服务器 %s 启动失败(Failed)%s", ms.Name, why),
SummaryEN: fmt.Sprintf("server %s is Failed%s", ms.Name, why),
Hint: fmt.Sprintf("kubectl -n %s describe minecraftserver %s; docs/troubleshooting.md §2", ns, ms.Name),
}
}
// failureMessage is the first false condition's message, as ": <message>".
func failureMessage(ms *v1alpha1.MinecraftServer) string {
for _, c := range ms.Status.Conditions {
if c.Status == "False" && c.Message != "" {
return ": " + c.Message
}
}
return ""
}
// jobFinding reports a Job in the minecraft namespace (a world backup, a
// restore, a reaper run) that failed within jobFailureWindow, once.
func jobFinding(j *batchv1.Job, now time.Time) *Finding {
for _, c := range j.Status.Conditions {
if c.Type != batchv1.JobFailed || c.Status != corev1.ConditionTrue {
continue
}
if now.Sub(c.LastTransitionTime.Time) > jobFailureWindow {
return nil
}
kind, kindEN := "任务", "Job"
switch {
case strings.HasPrefix(j.Name, "backup-"):
kind, kindEN = "世界备份任务", "world backup Job"
case strings.HasPrefix(j.Name, "restore-"):
kind, kindEN = "世界恢复任务", "world restore Job"
case strings.HasPrefix(j.Name, platform.SAReaper+"-"):
kind, kindEN = "世界回收任务", "world reaper Job"
}
reason := ""
if c.Reason != "" {
reason = " (" + c.Reason + ")"
}
return &Finding{
Key: "job-failed/" + j.Name, Severity: Warning, Event: true,
Summary: fmt.Sprintf("%s %s 失败%s", kind, j.Name, reason),
SummaryEN: fmt.Sprintf("%s %s failed%s", kindEN, j.Name, reason),
Hint: fmt.Sprintf("kubectl -n %s logs job/%s", j.Namespace, j.Name),
}
}
return nil
}
// reaperFinding reports a reaper CronJob that has not succeeded for over a day:
// idle worlds are neither backed up nor released, and expired backups pile up.
func reaperFinding(cj *batchv1.CronJob, now time.Time) *Finding {
if cj.Spec.Suspend != nil && *cj.Spec.Suspend {
return nil
}
last := cj.CreationTimestamp.Time
if cj.Status.LastSuccessfulTime != nil {
last = cj.Status.LastSuccessfulTime.Time
}
if now.Sub(last) <= maxReaperAge {
return nil
}
return &Finding{
Key: "reaper-stale", Severity: Warning,
Summary: fmt.Sprintf("世界回收任务已超过 %s 没有成功运行", roundHours(now.Sub(last))),
SummaryEN: fmt.Sprintf("the world reaper has not succeeded for %s", roundHours(now.Sub(last))),
Hint: fmt.Sprintf("kubectl -n %s get jobs; docs/troubleshooting.md §10", cj.Namespace),
}
}
// nodeFindings reports a node that is not Ready or is under pressure: the
// kubelet is about to evict pods, or already is.
func nodeFindings(n *corev1.Node) []Finding {
var out []Finding
for _, c := range n.Status.Conditions {
switch {
case c.Type == corev1.NodeReady && c.Status != corev1.ConditionTrue:
out = append(out, Finding{
Key: "node/" + n.Name + "/NotReady", Severity: Critical, For: nodeNotReadyFor,
Summary: fmt.Sprintf("节点 %s 未就绪:%s", n.Name, c.Message),
SummaryEN: fmt.Sprintf("node %s is not Ready: %s", n.Name, c.Message),
Hint: "systemctl status k3s; journalctl -u k3s -n 200",
})
case c.Type != corev1.NodeReady && c.Status == corev1.ConditionTrue &&
(c.Type == corev1.NodeDiskPressure || c.Type == corev1.NodeMemoryPressure || c.Type == corev1.NodePIDPressure):
out = append(out, Finding{
Key: "node/" + n.Name + "/" + string(c.Type), Severity: Critical, For: nodePressureFor,
Summary: fmt.Sprintf("节点 %s 报告 %s,kubelet 正在驱逐 Pod", n.Name, c.Type),
SummaryEN: fmt.Sprintf("node %s reports %s: the kubelet is evicting pods", n.Name, c.Type),
Hint: "docs/troubleshooting.md §13b",
})
}
}
return out
}
// KubeAPIDown is the finding for an API server that did not answer.
func KubeAPIDown(err error) Finding {
return Finding{
Key: "kube-api", Severity: Critical, For: kubeAPIFor,
Summary: "Kubernetes API 无法访问:面板、登录和所有服务器的状态都无法确认",
SummaryEN: "the Kubernetes API is unreachable: the panel, sign-in and every server are in doubt",
Hint: fmt.Sprintf("systemctl status k3s; journalctl -u k3s -n 200 (%v)", err),
}
}
// PostgresDown is the finding for a database that did not answer.
func PostgresDown(err error) Finding {
return Finding{
Key: "postgres", Severity: Critical, For: postgresFor,
Summary: "PostgreSQL 无法连接:登录、面板和服务器管理都会失败",
SummaryEN: "PostgreSQL is unreachable: sign-in, the panel and server management fail",
Hint: fmt.Sprintf("systemctl status postgresql; journalctl -u postgresql -n 100 (%v)", err),
}
}
// BackupFinding reports a control-plane database backup older than a day, or
// none at all, in dir.
func BackupFinding(dir string, now time.Time) *Finding {
bundles, err := dbbackup.List(dir)
if err != nil {
return &Finding{
Key: "db-backup", Severity: Critical, For: backupFor,
Summary: fmt.Sprintf("无法读取数据库备份目录 %s", dir),
SummaryEN: fmt.Sprintf("cannot read the database backup directory %s: %v", dir, err),
Hint: "docs/troubleshooting.md §16",
}
}
if len(bundles) > 0 && now.Sub(bundles[0].Created) <= maxBackupAge {
return nil
}
f := &Finding{
Key: "db-backup", Severity: Critical, For: backupFor,
Summary: fmt.Sprintf("%s 里没有任何控制面数据库备份", dir),
SummaryEN: fmt.Sprintf("no control-plane database backup in %s", dir),
Hint: "journalctl -u felis-db-backup -n 50; take one now with `sudo felis db backup` (docs/troubleshooting.md §16)",
}
if len(bundles) > 0 {
age := roundHours(now.Sub(bundles[0].Created))
f.Summary = fmt.Sprintf("最新的控制面数据库备份已是 %s 前(%s)", age, bundles[0].Name)
f.SummaryEN = fmt.Sprintf("the newest control-plane database backup is %s old (%s)", age, bundles[0].Name)
}
return f
}
// DiskFindings reports each filesystem under paths that is running out of
// space. Paths on one filesystem are reported once, under the first of them; a
// path that does not exist is skipped (a feature that is not in use).
func DiskFindings(paths []string) []Finding {
var out []Finding
seen := map[uint64]bool{}
for _, p := range paths {
var st syscall.Stat_t
if err := syscall.Stat(p, &st); err != nil {
continue
}
dev := uint64(st.Dev) // int32 on darwin
if seen[dev] {
continue
}
seen[dev] = true
var fs syscall.Statfs_t
if err := syscall.Statfs(p, &fs); err != nil || fs.Blocks == 0 {
continue
}
free := float64(fs.Bavail) / float64(fs.Blocks)
bsize := uint64(fs.Bsize) // uint32 on darwin
avail := humanBytes(uint64(fs.Bavail) * bsize)
switch {
case free < diskCriticalRatio:
out = append(out, Finding{
Key: "disk/" + p, Severity: Critical, For: diskCriticalFor,
Summary: fmt.Sprintf("%s 所在磁盘只剩 %.1f%%(%s):kubelet 即将驱逐游戏服务器", p, free*100, avail),
SummaryEN: fmt.Sprintf("the filesystem holding %s is %.1f%% free (%s): the kubelet is about to evict game servers", p, free*100, avail),
Hint: "docs/troubleshooting.md §13b",
})
case free < diskLowRatio:
out = append(out, Finding{
Key: "disk/" + p, Severity: Warning, For: diskLowFor,
Summary: fmt.Sprintf("%s 所在磁盘只剩 %.1f%%(%s)", p, free*100, avail),
SummaryEN: fmt.Sprintf("the filesystem holding %s is %.1f%% free (%s)", p, free*100, avail),
Hint: "docs/troubleshooting.md §13b",
})
}
}
return out
}
// MemoryFinding reports sustained low available memory from a /proc/meminfo
// style file. A file that cannot be read (not Linux) reports nothing.
func MemoryFinding(meminfo string) *Finding {
f, err := os.Open(meminfo)
if err != nil {
return nil
}
defer f.Close()
vals := map[string]uint64{}
sc := bufio.NewScanner(f)
for sc.Scan() {
fields := strings.Fields(sc.Text())
if len(fields) < 2 {
continue
}
if v, err := strconv.ParseUint(fields[1], 10, 64); err == nil {
vals[strings.TrimSuffix(fields[0], ":")] = v * 1024
}
}
total, avail := vals["MemTotal"], vals["MemAvailable"]
if total == 0 || float64(avail)/float64(total) >= memoryLowRatio {
return nil
}
return &Finding{
Key: "memory", Severity: Warning, For: memoryLowFor,
Summary: fmt.Sprintf("主机可用内存只剩 %s / %s:有 OOM 风险", humanBytes(avail), humanBytes(total)),
SummaryEN: fmt.Sprintf("host memory available is %s of %s: processes risk being OOM-killed", humanBytes(avail), humanBytes(total)),
Hint: "PostgreSQL, the control plane, the registry and game servers share this node; stop idle servers or lower their memory",
}
}
func roundHours(d time.Duration) string {
return d.Round(time.Hour).String()
}
func humanBytes(b uint64) string {
const unit = 1024
if b < unit {
return fmt.Sprintf("%d B", b)
}
div, exp := uint64(unit), 0
for n := b / unit; n >= unit; n /= unit {
div *= unit
exp++
}
return fmt.Sprintf("%.1f %ciB", float64(b)/float64(div), "KMGTPE"[exp])
}
+187
View File
@@ -0,0 +1,187 @@
package watchdog
import (
"context"
"os"
"path/filepath"
"sort"
"strconv"
"strings"
"testing"
"time"
"felis.lolicon.best/internal/apis/felis/v1alpha1"
"felis.lolicon.best/internal/dbbackup"
appsv1 "k8s.io/api/apps/v1"
batchv1 "k8s.io/api/batch/v1"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
clientgoscheme "k8s.io/client-go/kubernetes/scheme"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
)
func scheme(t *testing.T) *runtime.Scheme {
t.Helper()
s := runtime.NewScheme()
if err := clientgoscheme.AddToScheme(s); err != nil {
t.Fatal(err)
}
if err := v1alpha1.AddToScheme(s); err != nil {
t.Fatal(err)
}
return s
}
func deploy(name string, available int32) *appsv1.Deployment {
d := &appsv1.Deployment{ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "felis"}}
d.Status.AvailableReplicas = available
return d
}
func server(name, role string, desired v1alpha1.DesiredState, phase v1alpha1.Phase) *v1alpha1.MinecraftServer {
ms := &v1alpha1.MinecraftServer{ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "minecraft"}}
if role != "" {
ms.Labels = map[string]string{v1alpha1.LabelSystemRole: role}
}
ms.Spec.DesiredState = desired
ms.Status.Phase = phase
return ms
}
func failedJob(name string, at time.Time) *batchv1.Job {
j := &batchv1.Job{ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "minecraft"}}
j.Status.Conditions = []batchv1.JobCondition{{
Type: batchv1.JobFailed, Status: corev1.ConditionTrue, Reason: "BackoffLimitExceeded",
LastTransitionTime: metav1.NewTime(at),
}}
return j
}
func findingKeys(fs []Finding) []string {
var out []string
for _, f := range fs {
out = append(out, f.Key)
}
sort.Strings(out)
return out
}
// TestClusterCheck: exactly the broken things are reported — a control-plane
// Deployment down or missing, the login gate not running, a failed user server,
// a recent failed backup Job, a late reaper and a node under disk pressure —
// and the healthy or deliberate states are not.
func TestClusterCheck(t *testing.T) {
now := t0
login := server("login", "login", v1alpha1.DesiredRunning, v1alpha1.PhaseStarting)
login.Status.Conditions = []metav1.Condition{{Type: "Ready", Status: "False", Message: "rcon unreachable"}}
reaper := &batchv1.CronJob{ObjectMeta: metav1.ObjectMeta{Name: "felis-reaper", Namespace: "minecraft",
CreationTimestamp: metav1.NewTime(now.Add(-72 * time.Hour))}}
reaper.Status.LastSuccessfulTime = &metav1.Time{Time: now.Add(-50 * time.Hour)}
node := &corev1.Node{ObjectMeta: metav1.ObjectMeta{Name: "n1"}}
node.Status.Conditions = []corev1.NodeCondition{
{Type: corev1.NodeReady, Status: corev1.ConditionTrue},
{Type: corev1.NodeDiskPressure, Status: corev1.ConditionTrue},
{Type: corev1.NodeMemoryPressure, Status: corev1.ConditionFalse},
}
cl := fake.NewClientBuilder().WithScheme(scheme(t)).WithObjects(
deploy("felis-api", 1), deploy("felis-operator", 0),
login,
server("lobby", "lobby", v1alpha1.DesiredRunning, v1alpha1.PhaseRunning),
server("halted", "lobby", v1alpha1.DesiredStopped, v1alpha1.PhaseStopped),
server("broken", "", v1alpha1.DesiredRunning, v1alpha1.PhaseFailed),
server("asleep", "", v1alpha1.DesiredStopped, v1alpha1.PhaseStopped),
failedJob("backup-survival-abc", now.Add(-time.Hour)),
failedJob("backup-old-abc", now.Add(-48*time.Hour)),
reaper, node,
).Build()
got, err := Cluster{Client: cl, ControlNamespace: "felis", MinecraftNamespace: "minecraft"}.Check(context.Background(), now)
if err != nil {
t.Fatalf("Check: %v", err)
}
want := []string{
"deployment/felis-operator", "deployment/registry",
"job-failed/backup-survival-abc", "node/n1/DiskPressure", "reaper-stale",
"server-failed/broken", "system-server/login",
}
if strings.Join(findingKeys(got), ",") != strings.Join(want, ",") {
t.Fatalf("findings = %v, want %v", findingKeys(got), want)
}
for _, f := range got {
switch f.Key {
case "system-server/login":
if f.Severity != Critical || !strings.Contains(f.Summary, "rcon unreachable") || !strings.Contains(f.SummaryEN, "players cannot join") {
t.Errorf("login finding = %+v", f)
}
case "deployment/registry":
if !strings.Contains(f.SummaryEN, "missing") {
t.Errorf("registry finding = %+v, want missing", f)
}
case "job-failed/backup-survival-abc":
if !f.Event || !strings.Contains(f.SummaryEN, "world backup Job") {
t.Errorf("job finding = %+v", f)
}
}
}
}
// TestClusterCheckAPIDown: a List the API server refuses is an error, not an
// empty (healthy-looking) cluster.
func TestClusterCheckAPIDown(t *testing.T) {
cl := fake.NewClientBuilder().WithScheme(runtime.NewScheme()).Build()
if _, err := (Cluster{Client: cl, ControlNamespace: "felis", MinecraftNamespace: "minecraft"}).Check(context.Background(), t0); err == nil {
t.Fatal("Check succeeded against a client that knows no kinds")
}
}
func TestBackupFinding(t *testing.T) {
dir := t.TempDir()
if f := BackupFinding(dir, t0); f == nil || !strings.Contains(f.SummaryEN, "no control-plane database backup") {
t.Fatalf("empty dir: %+v", f)
}
touch := func(at time.Time) {
if err := os.WriteFile(filepath.Join(dir, dbbackup.BundleName(at, "daily")), []byte("x"), 0o600); err != nil {
t.Fatal(err)
}
}
touch(t0.Add(-30 * time.Hour))
if f := BackupFinding(dir, t0); f == nil || !strings.Contains(f.SummaryEN, "30h0m0s old") {
t.Fatalf("stale: %+v", f)
}
touch(t0.Add(-2 * time.Hour))
if f := BackupFinding(dir, t0); f != nil {
t.Fatalf("fresh backup reported: %+v", f)
}
}
func TestMemoryFinding(t *testing.T) {
path := filepath.Join(t.TempDir(), "meminfo")
write := func(avail int) {
body := "MemTotal: 24000000 kB\nMemFree: 100000 kB\nMemAvailable: " + strconv.Itoa(avail) + " kB\n"
if err := os.WriteFile(path, []byte(body), 0o600); err != nil {
t.Fatal(err)
}
}
write(2000000)
if f := MemoryFinding(path); f == nil || f.Key != "memory" {
t.Fatalf("8%% available: %+v", f)
}
write(6000000)
if f := MemoryFinding(path); f != nil {
t.Fatalf("25%% available reported: %+v", f)
}
if f := MemoryFinding(filepath.Join(t.TempDir(), "none")); f != nil {
t.Fatalf("missing meminfo reported: %+v", f)
}
}
// TestDiskFindingsDedupAndSkip: two paths on one filesystem yield at most one
// finding, and a path that does not exist is skipped.
func TestDiskFindingsDedupAndSkip(t *testing.T) {
dir := t.TempDir()
got := DiskFindings([]string{dir, filepath.Join(dir, "."), filepath.Join(dir, "missing")})
if len(got) > 1 {
t.Fatalf("findings = %v, want at most one for one filesystem", findingKeys(got))
}
}
+340
View File
@@ -0,0 +1,340 @@
// Package watchdog is the consumer the platform's health signals otherwise lack:
// `felis watchdog` runs on the host from a systemd timer, checks the control
// plane, the login gate, the fleet, PostgreSQL, the database backups and the
// node, and mails the platform owners when something has stayed wrong long
// enough to matter, again every day while it stays wrong, and once more when it
// clears.
//
// It runs on the host rather than in the cluster so the failures that take the
// cluster down (k3s stopped, the API server wedged, the node out of disk) are
// still reported: the one thing it needs to send mail, the relay, it reaches
// directly, with the credentials it cached on its last good run.
//
// This file is the pure part: turning one run's findings into state changes and
// a message. The probes that produce findings live in probes.go.
package watchdog
import (
"encoding/json"
"errors"
"fmt"
"os"
"path/filepath"
"sort"
"strings"
"time"
)
// Severity orders how urgent a finding is.
type Severity string
const (
Critical Severity = "critical"
Warning Severity = "warning"
)
// Finding is one condition a probe saw on this run.
type Finding struct {
// Key identifies the condition across runs, e.g. "deployment/felis-api". A
// key is also how a condition is found again once it clears.
Key string `json:"key"`
Severity Severity `json:"severity"`
// Summary and SummaryEN name what is wrong, in Chinese and English.
Summary string `json:"summary"`
SummaryEN string `json:"summary_en"`
// Hint says where to look next (a command, a troubleshooting section).
Hint string `json:"hint,omitempty"`
// For is how long the condition must hold before it is mailed, so a rollout
// or a restart that heals itself never pages anyone.
For time.Duration `json:"for"`
// Event marks a one-off (a failed Job): mailed once, with no reminder and no
// resolved notice when it goes away.
Event bool `json:"event,omitempty"`
}
// Report is one run's observations.
type Report struct {
Findings []Finding
// Unknown lists key prefixes no probe could evaluate on this run (the API
// server being down hides every cluster check). Alerts under them keep their
// state rather than reading as cleared.
Unknown []string
}
// Alert is a finding the state file remembers across runs.
type Alert struct {
Finding
FirstSeen time.Time `json:"first_seen"`
// Notified is when the alert was last mailed, at NotifiedSeverity; zero
// while it is pending.
Notified time.Time `json:"notified,omitempty"`
NotifiedSeverity Severity `json:"notified_severity,omitempty"`
// ClearedAt is when a mailed alert was first seen gone. It is mailed as
// resolved only after staying gone for resolveAfter, so a value hovering at
// its threshold does not mail on every crossing.
ClearedAt time.Time `json:"cleared_at,omitempty"`
}
// State is what the watchdog keeps between runs.
type State struct {
Alerts map[string]*Alert `json:"alerts"`
// Recipients and SMTPPassword are cached from the last run that could read
// them (the database and the felis-smtp Secret), so an outage of either can
// still be mailed.
Recipients []string `json:"recipients,omitempty"`
SMTPPassword string `json:"smtp_password,omitempty"`
}
const (
// remindEvery is how often a firing alert is mailed again.
remindEvery = 24 * time.Hour
// resolveAfter is how long a mailed alert must stay gone before it is
// mailed as resolved.
resolveAfter = 10 * time.Minute
)
// Plan is what one run has to tell the owners.
type Plan struct {
Firing []Alert // newly past their For, or worse than when last mailed
Reminders []Alert // still firing a day after the last mail
Resolved []Alert // gone for resolveAfter
// Active is every mailed condition still firing (one-off events aside), for
// the message's summary.
Active []Alert
}
// Empty reports whether the plan has nothing to mail.
func (p Plan) Empty() bool {
return len(p.Firing) == 0 && len(p.Reminders) == 0 && len(p.Resolved) == 0
}
// Observe folds a report into s and returns what is due to be mailed. It only
// tracks when conditions appeared and cleared; Commit records the plan as
// delivered. A run whose mail failed, or that falls in a quiet period, saves
// the state without committing, and the same alerts come due again next run.
func (s *State) Observe(r Report, now time.Time) Plan {
if s.Alerts == nil {
s.Alerts = map[string]*Alert{}
}
seen := map[string]bool{}
var plan Plan
for _, f := range r.Findings {
seen[f.Key] = true
a, ok := s.Alerts[f.Key]
if !ok {
a = &Alert{FirstSeen: now}
s.Alerts[f.Key] = a
}
a.Finding = f
a.ClearedAt = time.Time{}
notified := !a.Notified.IsZero()
switch {
case !notified && now.Sub(a.FirstSeen) >= f.For:
plan.Firing = append(plan.Firing, *a)
case notified && f.Severity == Critical && a.NotifiedSeverity != Critical:
// Worse than when it was mailed (a disk from low to nearly full):
// new news, mailed at once.
plan.Firing = append(plan.Firing, *a)
case notified && !f.Event && now.Sub(a.Notified) >= remindEvery:
plan.Reminders = append(plan.Reminders, *a)
}
}
for key, a := range s.Alerts {
if seen[key] || underAny(key, r.Unknown) {
continue
}
switch {
case a.Notified.IsZero() || a.Event:
// Never mailed, or a one-off: nothing to take back.
delete(s.Alerts, key)
case a.ClearedAt.IsZero():
a.ClearedAt = now
case now.Sub(a.ClearedAt) >= resolveAfter:
plan.Resolved = append(plan.Resolved, *a)
}
}
for _, a := range s.Alerts {
if !a.Notified.IsZero() && a.ClearedAt.IsZero() && !a.Event {
plan.Active = append(plan.Active, *a)
}
}
for _, l := range [][]Alert{plan.Firing, plan.Reminders, plan.Resolved, plan.Active} {
sortAlerts(l)
}
return plan
}
// Commit records plan as mailed at now.
func (s *State) Commit(plan Plan, now time.Time) {
for _, l := range [][]Alert{plan.Firing, plan.Reminders} {
for _, a := range l {
if cur := s.Alerts[a.Key]; cur != nil {
cur.Notified, cur.NotifiedSeverity = now, a.Severity
}
}
}
for _, a := range plan.Resolved {
delete(s.Alerts, a.Key)
}
}
func underAny(key string, prefixes []string) bool {
for _, p := range prefixes {
if strings.HasPrefix(key, p) {
return true
}
}
return false
}
// sortAlerts puts critical alerts first, then orders by key.
func sortAlerts(l []Alert) {
sort.Slice(l, func(i, j int) bool {
if (l[i].Severity == Critical) != (l[j].Severity == Critical) {
return l[i].Severity == Critical
}
return l[i].Key < l[j].Key
})
}
// Message renders the plan as one mail: a subject naming the counts and host,
// and a body listing what fired, what is still wrong and what cleared.
func (p Plan) Message(host string, now time.Time) (subject, body string) {
var parts, partsEN []string
if n := len(p.Firing) + len(p.Reminders); n > 0 {
parts = append(parts, fmt.Sprintf("%d 项异常", n))
partsEN = append(partsEN, fmt.Sprintf("%d firing", n))
}
if n := len(p.Resolved); n > 0 {
parts = append(parts, fmt.Sprintf("%d 项恢复", n))
partsEN = append(partsEN, fmt.Sprintf("%d resolved", n))
}
tag := "告警"
if len(p.Firing)+len(p.Reminders) == 0 {
tag = "恢复"
} else if hasCritical(p.Firing) || hasCritical(p.Reminders) {
tag = "严重告警"
}
subject = fmt.Sprintf("Felis %s(%s):%s · %s", tag, host, strings.Join(parts, ","), strings.Join(partsEN, ", "))
var b strings.Builder
fmt.Fprintf(&b, "Felis 平台巡检 / platform watchdog — %s, %s\r\n", host, now.UTC().Format("2006-01-02 15:04 UTC"))
section := func(title string, l []Alert, since bool) {
if len(l) == 0 {
return
}
fmt.Fprintf(&b, "\r\n%s\r\n", title)
for _, a := range l {
fmt.Fprintf(&b, "\r\n[%s] %s\r\n %s\r\n", a.Severity, a.Summary, a.SummaryEN)
if since {
fmt.Fprintf(&b, " 自 / since %s\r\n", a.FirstSeen.UTC().Format("2006-01-02 15:04 UTC"))
}
if a.Hint != "" {
fmt.Fprintf(&b, " → %s\r\n", a.Hint)
}
}
}
section("== 新出现的异常 / new ==", p.Firing, true)
section("== 仍未恢复(每日提醒)/ still firing (daily reminder) ==", p.Reminders, true)
section("== 已恢复 / resolved ==", p.Resolved, false)
if others := without(p.Active, p.Firing, p.Reminders); len(others) > 0 {
fmt.Fprintf(&b, "\r\n== 其他仍在进行的告警 / also still firing ==\r\n")
for _, a := range others {
fmt.Fprintf(&b, " [%s] %s / %s\r\n", a.Severity, a.Summary, a.SummaryEN)
}
}
// The journal rather than a -dry-run: the unit carries the flags (proxy port,
// disk paths) a bare command line would not.
fmt.Fprintf(&b, "\r\n在主机上运行 `journalctl -u felis-watchdog -n 40` 查看最近一次巡检的全部检查结果。\r\n"+
"Run `journalctl -u felis-watchdog -n 40` on the host for every check of the latest run.\r\n")
return subject, b.String()
}
func hasCritical(l []Alert) bool {
for _, a := range l {
if a.Severity == Critical {
return true
}
}
return false
}
// without returns the alerts of all not keyed in any of skip.
func without(all []Alert, skip ...[]Alert) []Alert {
keys := map[string]bool{}
for _, l := range skip {
for _, a := range l {
keys[a.Key] = true
}
}
var out []Alert
for _, a := range all {
if !keys[a.Key] {
out = append(out, a)
}
}
return out
}
// LoadState reads the state file; a missing file is a fresh state.
func LoadState(path string) (*State, error) {
raw, err := os.ReadFile(path)
if errors.Is(err, os.ErrNotExist) {
return &State{Alerts: map[string]*Alert{}}, nil
}
if err != nil {
return nil, err
}
var s State
if err := json.Unmarshal(raw, &s); err != nil {
return nil, fmt.Errorf("parse %s: %w", path, err)
}
if s.Alerts == nil {
s.Alerts = map[string]*Alert{}
}
return &s, nil
}
// SaveState writes s atomically, readable by root only: it caches the relay
// password.
func SaveState(path string, s *State) error {
raw, err := json.MarshalIndent(s, "", " ")
if err != nil {
return err
}
if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil {
return err
}
tmp, err := os.CreateTemp(filepath.Dir(path), ".state-*")
if err != nil {
return err
}
defer os.Remove(tmp.Name())
if err := tmp.Chmod(0o600); err != nil {
tmp.Close()
return err
}
if _, err := tmp.Write(raw); err != nil {
tmp.Close()
return err
}
if err := tmp.Close(); err != nil {
return err
}
return os.Rename(tmp.Name(), path)
}
// QuietUntil reads the maintenance marker the installer writes while it
// restarts things on purpose: a Unix timestamp, before which nothing is mailed.
// A missing or unreadable marker means no quiet period.
func QuietUntil(path string) time.Time {
raw, err := os.ReadFile(path)
if err != nil {
return time.Time{}
}
var sec int64
if _, err := fmt.Sscan(strings.TrimSpace(string(raw)), &sec); err != nil {
return time.Time{}
}
return time.Unix(sec, 0)
}
+207
View File
@@ -0,0 +1,207 @@
package watchdog
import (
"os"
"path/filepath"
"strings"
"testing"
"time"
)
var t0 = time.Date(2026, 9, 24, 12, 0, 0, 0, time.UTC)
func finding(key string, sev Severity, forDur time.Duration) Finding {
return Finding{Key: key, Severity: sev, For: forDur, Summary: key + " 中文", SummaryEN: key + " en"}
}
func keys(l []Alert) []string {
var out []string
for _, a := range l {
out = append(out, a.Key)
}
return out
}
// run observes r at now and commits whatever came due, like a run whose mail
// went out.
func run(s *State, r Report, now time.Time) Plan {
p := s.Observe(r, now)
s.Commit(p, now)
return p
}
// TestObserveLifecycle walks one condition through pending, firing, the daily
// reminder, clearing, and the resolved notice after it stayed clear.
func TestObserveLifecycle(t *testing.T) {
s := &State{}
f := finding("deployment/felis-api", Critical, 5*time.Minute)
down := Report{Findings: []Finding{f}}
if p := run(s, down, t0); !p.Empty() {
t.Fatalf("first sight mailed %+v, want it pending for 5m", p)
}
if p := run(s, down, t0.Add(4*time.Minute)); !p.Empty() {
t.Fatalf("4m in mailed %v", keys(p.Firing))
}
p := run(s, down, t0.Add(6*time.Minute))
if len(p.Firing) != 1 || p.Firing[0].Key != f.Key || !p.Firing[0].FirstSeen.Equal(t0) {
t.Fatalf("6m in: firing = %+v, want %s since t0", p.Firing, f.Key)
}
if p := run(s, down, t0.Add(8*time.Minute)); !p.Empty() {
t.Fatalf("already mailed, mailed again: %+v", p)
}
p = run(s, down, t0.Add(6*time.Minute+remindEvery))
if len(p.Reminders) != 1 || len(p.Firing) != 0 {
t.Fatalf("a day later: %+v, want one reminder", p)
}
up := Report{}
later := t0.Add(7*time.Minute + remindEvery)
if p := run(s, up, later); !p.Empty() {
t.Fatalf("just cleared, mailed %+v; want resolveAfter to pass first", p)
}
p = run(s, up, later.Add(resolveAfter))
if len(p.Resolved) != 1 || p.Resolved[0].Key != f.Key {
t.Fatalf("resolved = %v, want %s", keys(p.Resolved), f.Key)
}
if len(s.Alerts) != 0 {
t.Errorf("state after resolve = %v, want empty", s.Alerts)
}
}
// TestObserveFlapWithinResolveWindow: a condition that returns before
// resolveAfter is neither resolved nor mailed as new.
func TestObserveFlapWithinResolveWindow(t *testing.T) {
s := &State{}
f := finding("disk//", Warning, 0)
run(s, Report{Findings: []Finding{f}}, t0)
run(s, Report{}, t0.Add(time.Minute))
if p := run(s, Report{Findings: []Finding{f}}, t0.Add(3*time.Minute)); !p.Empty() {
t.Fatalf("flap mailed %+v", p)
}
if a := s.Alerts[f.Key]; a == nil || !a.ClearedAt.IsZero() {
t.Fatalf("alert after flap = %+v, want firing again", a)
}
}
// TestObservePendingNeverMailed: a condition that heals inside its For is
// dropped without a word.
func TestObservePendingNeverMailed(t *testing.T) {
s := &State{}
run(s, Report{Findings: []Finding{finding("deployment/registry", Critical, 5*time.Minute)}}, t0)
if p := run(s, Report{}, t0.Add(2*time.Minute)); !p.Empty() || len(s.Alerts) != 0 {
t.Fatalf("healed pending alert: plan %+v state %v", p, s.Alerts)
}
}
// TestObserveEscalation: a warning that turns critical is mailed again at once.
func TestObserveEscalation(t *testing.T) {
s := &State{}
run(s, Report{Findings: []Finding{finding("disk//", Warning, 0)}}, t0)
p := run(s, Report{Findings: []Finding{finding("disk//", Critical, 5*time.Minute)}}, t0.Add(time.Minute))
if len(p.Firing) != 1 || p.Firing[0].Severity != Critical {
t.Fatalf("escalation = %+v, want the critical finding mailed", p.Firing)
}
if p := run(s, Report{Findings: []Finding{finding("disk//", Critical, 5*time.Minute)}}, t0.Add(2*time.Minute)); !p.Empty() {
t.Fatalf("critical mailed twice: %+v", p)
}
}
// TestObserveEventOnce: a failed Job is mailed once, never reminded, and leaves
// without a resolved notice.
func TestObserveEventOnce(t *testing.T) {
s := &State{}
ev := finding("job-failed/backup-survival-x", Warning, 0)
ev.Event = true
if p := run(s, Report{Findings: []Finding{ev}}, t0); len(p.Firing) != 1 {
t.Fatalf("event not mailed: %+v", p)
}
if p := run(s, Report{Findings: []Finding{ev}}, t0.Add(remindEvery+time.Hour)); !p.Empty() {
t.Fatalf("event reminded: %+v", p)
}
if p := run(s, Report{}, t0.Add(remindEvery+2*time.Hour)); !p.Empty() || len(s.Alerts) != 0 {
t.Fatalf("event gone: plan %+v state %v, want silent removal", p, s.Alerts)
}
}
// TestObserveUnknownCarriesOver: while the API server is down, cluster alerts
// neither resolve nor restart their clocks.
func TestObserveUnknownCarriesOver(t *testing.T) {
s := &State{}
f := finding("system-server/login", Critical, 0)
run(s, Report{Findings: []Finding{f}}, t0)
api := finding("kube-api", Critical, 5*time.Minute)
for i := 1; i <= 3; i++ {
p := run(s, Report{Findings: []Finding{api}, Unknown: ClusterPrefixes}, t0.Add(time.Duration(i)*resolveAfter))
if len(p.Resolved) != 0 {
t.Fatalf("run %d resolved %v while the API was down", i, keys(p.Resolved))
}
}
if a := s.Alerts[f.Key]; a == nil || !a.ClearedAt.IsZero() {
t.Fatalf("login alert = %+v, want still firing", a)
}
}
// TestObserveUncommittedRetries: a plan whose mail failed comes due again.
func TestObserveUncommittedRetries(t *testing.T) {
s := &State{}
f := finding("postgres", Critical, 0)
if p := s.Observe(Report{Findings: []Finding{f}}, t0); len(p.Firing) != 1 {
t.Fatalf("not due: %+v", p)
}
if p := s.Observe(Report{Findings: []Finding{f}}, t0.Add(2*time.Minute)); len(p.Firing) != 1 || !p.Firing[0].FirstSeen.Equal(t0) {
t.Fatalf("after a failed send: %+v, want the same alert due again since t0", p.Firing)
}
}
func TestMessage(t *testing.T) {
s := &State{}
run(s, Report{Findings: []Finding{finding("memory", Warning, 0)}}, t0)
p := run(s, Report{Findings: []Finding{finding("memory", Warning, 0), finding("postgres", Critical, 0)}}, t0.Add(time.Minute))
subject, body := p.Message("node-1", t0.Add(time.Minute))
for _, want := range []string{"严重告警", "node-1", "1 项异常", "1 firing"} {
if !strings.Contains(subject, want) {
t.Errorf("subject %q lacks %q", subject, want)
}
}
for _, want := range []string{"postgres 中文", "postgres en", "== 其他仍在进行的告警 / also still firing ==", "memory 中文 / memory en", "journalctl -u felis-watchdog"} {
if !strings.Contains(body, want) {
t.Errorf("body lacks %q:\n%s", want, body)
}
}
}
func TestStateRoundTripPrivate(t *testing.T) {
path := filepath.Join(t.TempDir(), "watchdog", "state.json")
s := &State{SMTPPassword: "secret", Recipients: []string{"[email protected]"}}
run(s, Report{Findings: []Finding{finding("memory", Warning, time.Hour)}}, t0)
if err := SaveState(path, s); err != nil {
t.Fatalf("SaveState: %v", err)
}
info, err := os.Stat(path)
if err != nil || info.Mode().Perm() != 0o600 {
t.Fatalf("state file mode = %v (%v), want 0600", info.Mode(), err)
}
got, err := LoadState(path)
if err != nil || got.SMTPPassword != "secret" || !got.Alerts["memory"].FirstSeen.Equal(t0) {
t.Fatalf("LoadState = %+v, %v", got, err)
}
fresh, err := LoadState(filepath.Join(t.TempDir(), "missing.json"))
if err != nil || fresh.Alerts == nil {
t.Fatalf("missing state = %+v, %v", fresh, err)
}
}
func TestQuietUntil(t *testing.T) {
dir := t.TempDir()
marker := filepath.Join(dir, "quiet")
if !QuietUntil(marker).IsZero() {
t.Error("missing marker should mean no quiet period")
}
if err := os.WriteFile(marker, []byte("1790000000\n"), 0o644); err != nil {
t.Fatal(err)
}
if got := QuietUntil(marker); got.Unix() != 1790000000 {
t.Errorf("QuietUntil = %v", got)
}
}