From d17524cd67c256049abc4c79cfe376134461f66d Mon Sep 17 00:00:00 2001 From: Lemon-miaow Date: Thu, 24 Sep 2026 18:06:19 +0800 Subject: [PATCH] =?UTF-8?q?feat(watchdog):=20=E4=B8=BB=E6=9C=BA=E4=BE=A7?= =?UTF-8?q?=E5=B7=A1=E6=A3=80=E5=AE=9A=E6=97=B6=E5=99=A8=E6=8C=89=E5=BC=82?= =?UTF-8?q?=E5=B8=B8=E9=82=AE=E4=BB=B6=E9=80=9A=E7=9F=A5=E5=B9=B3=E5=8F=B0?= =?UTF-8?q?=E6=89=80=E6=9C=89=E8=80=85=EF=BC=8Coperator=20=E5=A2=9E?= =?UTF-8?q?=E5=8A=A0=20phase=20=E4=B8=8E=20build=5Finfo=20=E6=8C=87?= =?UTF-8?q?=E6=A0=87=E3=80=81=E5=8D=A1=E6=AD=BB=E5=AD=98=E6=B4=BB=E6=8E=A2?= =?UTF-8?q?=E9=92=88=E4=B8=8E=E5=91=8A=E8=AD=A6=E8=A7=84=E5=88=99?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- cmd/felis/api.go | 3 + cmd/felis/operator.go | 40 ++- cmd/felis/operator_test.go | 28 ++ cmd/felis/run.go | 2 + cmd/felis/watchdog.go | 249 ++++++++++++++ cmd/felis/watchdog_test.go | 42 +++ deploy/alerts/felis-alerts.yaml | 126 ++++++- deploy/alerts/felis-alerts_test.yml | 216 ++++++++++++ deploy/alerts/felis-prometheusrule.yaml | 110 ++++++ deploy/bootstrap.sh | 67 ++++ deploy/bootstrap_test.sh | 69 ++++ docs/troubleshooting.md | 106 +++++- internal/metrics/alerts_test.go | 98 ++++++ internal/metrics/metrics.go | 45 +++ internal/metrics/metrics_test.go | 24 ++ internal/operator/gauge.go | 33 +- internal/operator/gauge_test.go | 27 ++ internal/operator/reconciler.go | 14 +- internal/operator/watch.go | 91 +++++ internal/operator/watch_test.go | 83 +++++ internal/platform/workloads.go | 49 ++- internal/platform/workloads_test.go | 39 ++- internal/watchdog/probes.go | 425 ++++++++++++++++++++++++ internal/watchdog/probes_test.go | 187 +++++++++++ internal/watchdog/watchdog.go | 340 +++++++++++++++++++ internal/watchdog/watchdog_test.go | 207 ++++++++++++ 26 files changed, 2686 insertions(+), 34 deletions(-) create mode 100644 cmd/felis/operator_test.go create mode 100644 cmd/felis/watchdog.go create mode 100644 cmd/felis/watchdog_test.go create mode 100644 internal/metrics/alerts_test.go create mode 100644 internal/operator/watch.go create mode 100644 internal/operator/watch_test.go create mode 100644 internal/watchdog/probes.go create mode 100644 internal/watchdog/probes_test.go create mode 100644 internal/watchdog/watchdog.go create mode 100644 internal/watchdog/watchdog_test.go diff --git a/cmd/felis/api.go b/cmd/felis/api.go index d7358a7..706e703 100644 --- a/cmd/felis/api.go +++ b/cmd/felis/api.go @@ -19,6 +19,7 @@ import ( "felis.lolicon.best/internal/fileedit" "felis.lolicon.best/internal/imagepin" "felis.lolicon.best/internal/mail" + "felis.lolicon.best/internal/metrics" "felis.lolicon.best/internal/naming" "felis.lolicon.best/internal/panel" "felis.lolicon.best/internal/passkey" @@ -110,6 +111,8 @@ func cmdAPI(args []string, stdout, stderr io.Writer) int { return 1 } + metrics.SetBuildInfo("api", resolvedVersion()) + token := os.Getenv("FELIS_SERVICE_TOKEN") if token == "" { fmt.Fprintln(stderr, "felis api: warning: FELIS_SERVICE_TOKEN unset — internal face will reject all callers") diff --git a/cmd/felis/operator.go b/cmd/felis/operator.go index c1e57f9..3a3534e 100644 --- a/cmd/felis/operator.go +++ b/cmd/felis/operator.go @@ -1,11 +1,15 @@ package main import ( + "context" + "errors" "flag" "fmt" "io" "log/slog" + "net/http" "os" + "time" "felis.lolicon.best/internal/apis/felis/v1alpha1" felismetrics "felis.lolicon.best/internal/metrics" @@ -75,16 +79,19 @@ func cmdOperator(args []string, _, stderr io.Writer) int { } fmt.Fprintf(stderr, "felis operator: watching namespace %q\n", *namespace) - // Register the two probe endpoints. controller-runtime only mounts /healthz and - // /readyz once at least one check is registered, so a bare listener would 404. - // The checks are the canonical always-pass ping: the probes' contract is "the - // manager process is up and serving", and a dependency hiccup (e.g. an API blip) - // must not restart the operator. - if err := mgr.AddHealthzCheck("ping", healthz.Ping); err != nil { + // /healthz fails while a reconcile pass has been stuck past its limit, so the + // liveness probe restarts an operator whose workers are wedged (a Pod whose + // process answers but no server starts or stops). /readyz waits for the + // informer caches: until they sync the operator acts on nothing, and one that + // never syncs (lost RBAC, an unreachable API) never reports Available. + // A dependency hiccup fails neither: the caches ride through API blips, and + // each pass is bounded well inside the stuck limit. + watch := &operator.ReconcileWatch{} + if err := mgr.AddHealthzCheck("reconcile", watch.Check); err != nil { fmt.Fprintf(stderr, "felis operator: register healthz check: %v\n", err) return 1 } - if err := mgr.AddReadyzCheck("ping", healthz.Ping); err != nil { + if err := mgr.AddReadyzCheck("informers", cacheSynced(mgr.GetCache())); err != nil { fmt.Fprintf(stderr, "felis operator: register readyz check: %v\n", err) return 1 } @@ -98,6 +105,7 @@ func cmdOperator(args []string, _, stderr io.Writer) int { fmt.Fprintf(stderr, "felis operator: register metrics: %v\n", err) return 1 } + felismetrics.SetBuildInfo("operator", resolvedVersion()) r := &operator.Reconciler{ Client: mgr.GetClient(), @@ -109,7 +117,8 @@ func cmdOperator(args []string, _, stderr io.Writer) int { FelisImage: os.Getenv("FELIS_IMAGE"), // Uncached: the maintenance-lock check lists Jobs only when a server is // about to start, which does not justify a namespace-wide Job informer. - Jobs: mgr.GetAPIReader(), + Jobs: mgr.GetAPIReader(), + Watch: watch, } if err := r.SetupWithManager(mgr); err != nil { fmt.Fprintf(stderr, "felis operator: setup controller: %v\n", err) @@ -131,3 +140,18 @@ func cmdOperator(args []string, _, stderr io.Writer) int { } return 0 } + +// cacheSynced is a readyz check that passes once every informer the manager +// started has synced. It waits at most a second, well inside the probe timeout. +func cacheSynced(c interface { + WaitForCacheSync(ctx context.Context) bool +}) healthz.Checker { + return func(req *http.Request) error { + ctx, cancel := context.WithTimeout(req.Context(), time.Second) + defer cancel() + if !c.WaitForCacheSync(ctx) { + return errors.New("informer caches not synced") + } + return nil + } +} diff --git a/cmd/felis/operator_test.go b/cmd/felis/operator_test.go new file mode 100644 index 0000000..1953619 --- /dev/null +++ b/cmd/felis/operator_test.go @@ -0,0 +1,28 @@ +package main + +import ( + "context" + "net/http/httptest" + "testing" +) + +type fakeCache bool + +func (f fakeCache) WaitForCacheSync(ctx context.Context) bool { + if !f { + <-ctx.Done() + } + return bool(f) +} + +// TestCacheSynced: the operator reports ready only once its informers synced, +// and a check against caches that never sync returns within its own deadline. +func TestCacheSynced(t *testing.T) { + req := httptest.NewRequest("GET", "/readyz", nil) + if err := cacheSynced(fakeCache(true))(req); err != nil { + t.Errorf("synced: %v", err) + } + if err := cacheSynced(fakeCache(false))(req); err == nil { + t.Error("unsynced caches reported ready") + } +} diff --git a/cmd/felis/run.go b/cmd/felis/run.go index 4a8bcc6..79821ad 100644 --- a/cmd/felis/run.go +++ b/cmd/felis/run.go @@ -27,6 +27,7 @@ Commands: apply Create a MinecraftServer CRD (direct K8s write; use -f server.json) setup Run host bootstrap + first-run setup console (TUI; requires root/sudo) converge Fill in fields a newer desired spec added to already-installed system servers + watchdog Check the platform once and mail the owners what has gone wrong (run by felis-watchdog.timer) version Print the build stamp of this binary update Report which platform components have updates available breakGlass Open the local break-glass emergency console (TUI; requires root/sudo) @@ -67,6 +68,7 @@ var commands = map[string]func(args []string, stdout, stderr io.Writer) int{ "pin-images": cmdPinImages, "version": cmdVersion, "update": cmdUpdate, + "watchdog": cmdWatchdog, } // run dispatches a subcommand. It is separate from main so the router is diff --git a/cmd/felis/watchdog.go b/cmd/felis/watchdog.go new file mode 100644 index 0000000..1d83cda --- /dev/null +++ b/cmd/felis/watchdog.go @@ -0,0 +1,249 @@ +package main + +import ( + "context" + "database/sql" + "errors" + "flag" + "fmt" + "io" + "net" + "os" + "strings" + "time" + + "felis.lolicon.best/internal/config" + "felis.lolicon.best/internal/mail" + "felis.lolicon.best/internal/platform" + "felis.lolicon.best/internal/store" + "felis.lolicon.best/internal/watchdog" + corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// proxyFor is how long the game proxy may refuse connections before it is +// mailed: a restart takes seconds. +const proxyFor = 3 * time.Minute + +// cmdWatchdog runs one pass of the platform watchdog (internal/watchdog): it +// checks the cluster, PostgreSQL, the game proxy, the database backups and the +// host, prints every finding, and mails the platform owners what came due. +// deploy/bootstrap.sh runs it every two minutes from felis-watchdog.timer. +func cmdWatchdog(args []string, stdout, stderr io.Writer) int { + fs := flag.NewFlagSet("watchdog", flag.ContinueOnError) + fs.SetOutput(stderr) + cfgPath := fs.String("config", "/etc/felis/felis.toml", "path to felis.toml (the host copy, which reaches PostgreSQL on 127.0.0.1)") + statePath := fs.String("state", "/var/lib/felis/watchdog/state.json", "state kept between runs (root only: it caches the relay password)") + quietPath := fs.String("quiet-file", "/run/felis/watchdog-quiet-until", "Unix time before which nothing is mailed; the installer writes it while it restarts things on purpose") + backupDir := fs.String("backup-dir", "/var/lib/felis/db-backups", `control-plane database backups to check for freshness ("" skips the check)`) + diskPaths := fs.String("disk-paths", "/,/var/lib/rancher/k3s,/var/lib/postgresql,/var/lib/felis", "comma-separated paths whose filesystems must keep free space") + proxyAddr := fs.String("proxy-addr", "", `game proxy address to dial, e.g. 127.0.0.1:25565 ("" skips the check)`) + controlNS := fs.String("control-namespace", platform.DefaultControlNamespace, "namespace of the control plane") + dryRun := fs.Bool("dry-run", false, "print every finding and the mail that is due; send nothing and keep the state as it was") + if err := fs.Parse(args); err != nil { + if errors.Is(err, flag.ErrHelp) { + return 0 + } + return 2 + } + cfg, err := config.Load(*cfgPath) + if err != nil { + fmt.Fprintf(stderr, "felis watchdog: %v\n", err) + return 1 + } + state, err := watchdog.LoadState(*statePath) + if err != nil { + fmt.Fprintf(stderr, "felis watchdog: %v\n", err) + return 1 + } + ctx, cancel := context.WithTimeout(context.Background(), 90*time.Second) + defer cancel() + now := time.Now() + + var report watchdog.Report + add := func(f *watchdog.Finding) { + if f != nil { + report.Findings = append(report.Findings, *f) + } + } + + // The cluster: one unreachable API server stands in for every check behind it. + minecraftNS := cfg.K8s.Namespace + if minecraftNS == "" { + minecraftNS = platform.DefaultMinecraftNamespace + } + cl, err := buildSystemServerClient() + var found []watchdog.Finding + if err == nil { + found, err = watchdog.Cluster{Client: cl, ControlNamespace: *controlNS, MinecraftNamespace: minecraftNS}.Check(ctx, now) + } + if err != nil { + f := watchdog.KubeAPIDown(err) + add(&f) + report.Unknown = append(report.Unknown, watchdog.ClusterPrefixes...) + } else { + report.Findings = append(report.Findings, found...) + if cfg.SMTP.Host != "" { + refreshSMTPPassword(ctx, cl, *controlNS, state, stderr) + } + } + + if recipients, err := ownerEmails(ctx, cfg.Database.URL); err != nil { + f := watchdog.PostgresDown(err) + add(&f) + } else { + state.Recipients = recipients + } + + if *proxyAddr != "" { + add(proxyFinding(ctx, *proxyAddr)) + } + if *backupDir != "" { + add(watchdog.BackupFinding(*backupDir, now)) + } + report.Findings = append(report.Findings, watchdog.DiskFindings(splitList(*diskPaths))...) + add(watchdog.MemoryFinding("/proc/meminfo")) + + if len(report.Findings) == 0 { + fmt.Fprintln(stdout, "felis watchdog: every check passed") + } + for _, f := range report.Findings { + fmt.Fprintf(stdout, "felis watchdog: [%s] %s: %s\n", f.Severity, f.Key, f.SummaryEN) + } + + plan := state.Observe(report, now) + host, _ := os.Hostname() + subject, body := plan.Message(host, now) + if *dryRun { + if plan.Empty() { + fmt.Fprintln(stdout, "felis watchdog: nothing is due to be mailed") + } else { + fmt.Fprintf(stdout, "felis watchdog: due to be mailed to %s:\nSubject: %s\n\n%s", strings.Join(state.Recipients, ", "), subject, strings.ReplaceAll(body, "\r\n", "\n")) + } + return 0 + } + + save := func() int { + if err := watchdog.SaveState(*statePath, state); err != nil { + fmt.Fprintf(stderr, "felis watchdog: save state: %v\n", err) + return 1 + } + return 0 + } + if plan.Empty() { + return save() + } + if until := watchdog.QuietUntil(*quietPath); now.Before(until) { + fmt.Fprintf(stdout, "felis watchdog: quiet until %s (installer running); holding this mail: %s\n", until.UTC().Format(time.RFC3339), subject) + return save() + } + switch { + case cfg.SMTP.Host == "": + fmt.Fprintf(stdout, "felis watchdog: no [smtp] relay configured, so this is logged only: %s\n", subject) + case len(state.Recipients) == 0: + fmt.Fprintf(stdout, "felis watchdog: no owner account has a verified email, so this is logged only: %s\n", subject) + default: + if err := sendAlert(ctx, cfg, state, subject, body); err != nil { + // Not committed: the same alerts come due again next run. + fmt.Fprintf(stderr, "felis watchdog: mail %q: %v\n", subject, err) + save() + return 1 + } + fmt.Fprintf(stdout, "felis watchdog: mailed %s: %s\n", strings.Join(state.Recipients, ", "), subject) + } + state.Commit(plan, now) + return save() +} + +// refreshSMTPPassword caches the relay password from the felis-smtp Secret, or +// forgets it when the Secret is gone (a relay without AUTH). An env var named by +// [smtp] password_ref, when set, wins at send time instead. +func refreshSMTPPassword(ctx context.Context, cl client.Client, ns string, state *watchdog.State, stderr io.Writer) { + var sec corev1.Secret + err := cl.Get(ctx, client.ObjectKey{Namespace: ns, Name: platform.SMTPSecretName}, &sec) + switch { + case apierrors.IsNotFound(err): + state.SMTPPassword = "" + case err != nil: + fmt.Fprintf(stderr, "felis watchdog: read %s/%s (keeping the cached relay password): %v\n", ns, platform.SMTPSecretName, err) + default: + state.SMTPPassword = string(sec.Data[platform.SMTPSecretPasswordKey]) + } +} + +// ownerEmails pings PostgreSQL and returns the verified addresses of the +// enabled owner accounts, the people who can act on an alert. +func ownerEmails(ctx context.Context, url string) ([]string, error) { + ctx, cancel := context.WithTimeout(ctx, 15*time.Second) + defer cancel() + drv, err := store.Open(ctx, url) + if err != nil { + return nil, err + } + defer drv.Close() + rows, err := drv.DB().QueryContext(ctx, + `SELECT email FROM users + WHERE role = 'owner' AND email_verified AND COALESCE(email, '') <> '' + AND NOT disabled AND deleted_at IS NULL + ORDER BY email`) + if err != nil { + return nil, err + } + defer rows.Close() + var out []string + for rows.Next() { + var email sql.NullString + if err := rows.Scan(&email); err != nil { + return nil, err + } + out = append(out, email.String) + } + return out, rows.Err() +} + +// proxyFinding dials the game proxy; players reach every server through it. +func proxyFinding(ctx context.Context, addr string) *watchdog.Finding { + d := net.Dialer{Timeout: 5 * time.Second} + conn, err := d.DialContext(ctx, "tcp", addr) + if err == nil { + conn.Close() + return nil + } + return &watchdog.Finding{ + Key: "proxy", Severity: watchdog.Critical, For: proxyFor, + Summary: fmt.Sprintf("游戏代理 %s 无法连接:玩家进不了任何服务器", addr), + SummaryEN: fmt.Sprintf("the game proxy at %s refuses connections: players cannot reach any server", addr), + Hint: fmt.Sprintf("systemctl status felis-velocity; journalctl -u felis-velocity -n 200 (%v)", err), + } +} + +// sendAlert mails subject/body to every recipient; it fails only when no +// recipient got it. +func sendAlert(ctx context.Context, cfg *config.Config, state *watchdog.State, subject, body string) error { + password := state.SMTPPassword + if ref := cfg.SMTP.PasswordRef; ref != "" && os.Getenv(ref) != "" { + password = os.Getenv(ref) + } + relay := &mail.SMTP{Host: cfg.SMTP.Host, Port: cfg.SMTP.Port, From: cfg.SMTP.From, Username: cfg.SMTP.Username, Password: password} + var errs []error + for _, to := range state.Recipients { + if err := relay.SendNotice(ctx, to, subject, body); err != nil { + errs = append(errs, fmt.Errorf("%s: %w", to, err)) + } + } + if len(errs) == len(state.Recipients) { + return errors.Join(errs...) + } + return nil +} + +func splitList(s string) []string { + var out []string + for _, p := range strings.Split(s, ",") { + if p = strings.TrimSpace(p); p != "" { + out = append(out, p) + } + } + return out +} diff --git a/cmd/felis/watchdog_test.go b/cmd/felis/watchdog_test.go new file mode 100644 index 0000000..8ff042f --- /dev/null +++ b/cmd/felis/watchdog_test.go @@ -0,0 +1,42 @@ +package main + +import ( + "context" + "net" + "strings" + "testing" +) + +// TestProxyFinding: a listening proxy is healthy; a closed port is the critical +// "players cannot reach any server" finding. +func TestProxyFinding(t *testing.T) { + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + addr := ln.Addr().String() + go func() { + for { + c, err := ln.Accept() + if err != nil { + return + } + c.Close() + } + }() + if f := proxyFinding(context.Background(), addr); f != nil { + t.Fatalf("listening proxy reported: %+v", f) + } + ln.Close() + f := proxyFinding(context.Background(), addr) + if f == nil || f.Key != "proxy" || !strings.Contains(f.SummaryEN, addr) { + t.Fatalf("closed proxy = %+v, want the proxy finding", f) + } +} + +func TestSplitList(t *testing.T) { + got := splitList(" /, /var/lib/felis ,,") + if strings.Join(got, "|") != "/|/var/lib/felis" { + t.Fatalf("splitList = %q", got) + } +} diff --git a/deploy/alerts/felis-alerts.yaml b/deploy/alerts/felis-alerts.yaml index 11c472b..cf235b4 100644 --- a/deploy/alerts/felis-alerts.yaml +++ b/deploy/alerts/felis-alerts.yaml @@ -1,14 +1,20 @@ # Felis alert rules — plain Prometheus format (also the promtool-tested source # for felis-prometheusrule.yaml). See docs/troubleshooting.md §14 for scraping -# and loading instructions. +# and loading instructions. These are for a deployment that brings its own +# Prometheus; every install already runs `felis watchdog` on the host, which +# checks the same conditions without one and mails the owners (§14). # -# felis_* series come from two processes: -# - felis-operator pod :8080/metrics → felis_servers_total, felis_start_duration_seconds -# - felis-api internal :8081/metrics → felis_image_build_failures_total, +# felis_* series come from these processes: +# - felis-operator-metrics Service :8080 → felis_servers_total, felis_server_phase, +# felis_start_duration_seconds, +# felis_build_info{component="operator"}, +# controller_runtime_*, workqueue_* +# - felis-api-internal Service :8081 → felis_image_build_failures_total, # felis_mail_total, felis_rate_limited_total, # felis_auth_otp_lockouts_total, # felis_auth_failures_total, -# felis_audit_write_failures_total +# felis_audit_write_failures_total, +# felis_build_info{component="api"} # - node-exporter textfile collector → felis_db_backup_* (felis-db-backup.timer) # node_* / kube_* series come from node-exporter / kube-state-metrics. groups: @@ -37,6 +43,116 @@ groups: observed when readiness is first reached). A start that never completes records nothing — cross-check desiredState=Running servers with no ready phase (troubleshooting §1). + - name: felis.platform.rules + rules: + - alert: FelisOperatorDown + expr: absent(felis_build_info{component="operator"}) + for: 10m + labels: + severity: critical + annotations: + summary: "felis-operator is down or not scraped" + description: >- + No felis_build_info{component="operator"} series for 10 minutes. Without + the operator no server starts, stops or recovers. Check + `kubectl -n felis get deploy felis-operator` and its log; if the pod is + healthy, the felis-operator-metrics Service is not being scraped + (troubleshooting §14). + - alert: FelisAPIDown + expr: absent(felis_build_info{component="api"}) + for: 10m + labels: + severity: critical + annotations: + summary: "felis-api is down or not scraped" + description: >- + No felis_build_info{component="api"} series for 10 minutes. The panel, + sign-in and the proxy's player lookups all go through felis-api. Check + `kubectl -n felis get deploy felis-api` and its log; if the pod is + healthy, the felis-api-internal Service is not being scraped + (troubleshooting §14). + - alert: FelisLoginGateDown + expr: felis_server_phase{role="login",desired="Running",phase!="Running"} == 1 + for: 10m + labels: + severity: critical + annotations: + summary: "the login gate {{ $labels.server }} is {{ $labels.phase }}" + description: >- + Every player connection passes through the login server first, so no one + can join. The MinecraftServer's conditions carry the reason: + `kubectl -n minecraft describe minecraftserver {{ $labels.server }}` + (troubleshooting §1, §2). + - alert: FelisSystemServerDown + expr: felis_server_phase{role!="",role!="login",desired="Running",phase!="Running"} == 1 + for: 10m + labels: + severity: warning + annotations: + summary: "system server {{ $labels.server }} ({{ $labels.role }}) is {{ $labels.phase }}" + description: >- + Players who sign in are sent to the lobby; while it is down they stay at + the gate. `kubectl -n minecraft describe minecraftserver {{ $labels.server }}` + shows the reason (troubleshooting §1, §2). + - alert: FelisServerFailed + expr: felis_server_phase{role="",phase="Failed"} == 1 + for: 5m + labels: + severity: warning + annotations: + summary: "server {{ $labels.server }} is Failed" + description: >- + The operator gave up on this server (a crash loop, an image that will not + pull, a world volume that will not mount). Its conditions carry the + reason: `kubectl -n minecraft describe minecraftserver {{ $labels.server }}` + (troubleshooting §2). + - alert: FelisReconcileErrors + expr: sum(increase(controller_runtime_reconcile_errors_total{controller="minecraftserver"}[15m])) > 10 + for: 5m + labels: + severity: warning + annotations: + summary: "the operator failed over 10 reconciles in 15 minutes" + description: >- + Server changes are being retried instead of applied. The felis-operator + log names each failing server and its error. + - alert: FelisReconcileStuck + expr: max(workqueue_longest_running_processor_seconds{name="minecraftserver"}) > 300 + for: 5m + labels: + severity: critical + annotations: + summary: "an operator reconcile has been running for over 5 minutes" + description: >- + Each reconcile is bounded at 3 minutes, so this one is ignoring its + deadline and holding a worker. The liveness probe restarts the operator + once a pass passes 10 minutes; the log from before the restart shows + where it hung. + - name: felis.jobs.rules + rules: + - alert: FelisWorldJobFailed + expr: kube_job_failed{namespace="minecraft",condition="true"} == 1 + labels: + severity: warning + annotations: + summary: "Job {{ $labels.job_name }} failed" + description: >- + A world backup, restore or reaper run failed; after a failed backup that + world's newest archive is older than planned. + `kubectl -n minecraft logs job/{{ $labels.job_name }}` has the error + (troubleshooting §10). + - alert: FelisReaperStale + expr: time() - kube_cronjob_status_last_successful_time{namespace="minecraft",cronjob="felis-reaper"} > 26 * 3600 + for: 10m + labels: + severity: warning + annotations: + summary: "the world reaper has not succeeded in over 26h" + description: >- + felis-reaper runs daily; idle worlds are neither backed up nor reclaimed + while it fails. `kubectl -n minecraft get jobs --sort-by=.metadata.creationTimestamp` + lists its runs, and the newest one's log + shows why (troubleshooting §10). - name: felis.node.rules rules: - alert: FelisNodeDiskSpaceLow diff --git a/deploy/alerts/felis-alerts_test.yml b/deploy/alerts/felis-alerts_test.yml index 4a47e25..ce43628 100644 --- a/deploy/alerts/felis-alerts_test.yml +++ b/deploy/alerts/felis-alerts_test.yml @@ -291,3 +291,219 @@ tests: The actions went through but their audit rows were lost. The felis-api log names each lost row (`audit: lost ...`); the usual cause is PostgreSQL being unreachable or out of disk. + - name: operator and api presence + interval: 1m + input_series: + - series: 'felis_build_info{component="operator",version="v1",job="felis-operator",instance="op-0"}' + values: '1x30' + # felis-api stops being scraped after 5m; the series goes stale 5m later. + - series: 'felis_build_info{component="api",version="v1",job="felis-api",instance="api-0"}' + values: '1x5' + alert_rule_test: + - eval_time: 25m + alertname: FelisOperatorDown + exp_alerts: [] + - eval_time: 15m + alertname: FelisAPIDown + exp_alerts: [] + - eval_time: 25m + alertname: FelisAPIDown + exp_alerts: + - exp_labels: + severity: critical + component: api + exp_annotations: + summary: "felis-api is down or not scraped" + description: >- + No felis_build_info{component="api"} series for 10 minutes. The panel, + sign-in and the proxy's player lookups all go through felis-api. Check + `kubectl -n felis get deploy felis-api` and its log; if the pod is + healthy, the felis-api-internal Service is not being scraped + (troubleshooting §14). + - name: operator never scraped + interval: 1m + input_series: + - series: 'felis_build_info{component="api",version="v1",job="felis-api",instance="api-0"}' + values: '1x30' + alert_rule_test: + - eval_time: 5m + alertname: FelisOperatorDown + exp_alerts: [] + - eval_time: 15m + alertname: FelisOperatorDown + exp_alerts: + - exp_labels: + severity: critical + component: operator + exp_annotations: + summary: "felis-operator is down or not scraped" + description: >- + No felis_build_info{component="operator"} series for 10 minutes. Without + the operator no server starts, stops or recovers. Check + `kubectl -n felis get deploy felis-operator` and its log; if the pod is + healthy, the felis-operator-metrics Service is not being scraped + (troubleshooting §14). + - name: system and user servers down + interval: 1m + input_series: + - series: 'felis_server_phase{server="login",role="login",phase="Starting",desired="Running"}' + values: '1x20' + - series: 'felis_server_phase{server="lobby",role="lobby",phase="Failed",desired="Running"}' + values: '1x20' + # Stopped on purpose: not an outage. + - series: 'felis_server_phase{server="lobby2",role="lobby",phase="Stopped",desired="Stopped"}' + values: '1x20' + # A user server carries no role label (the operator publishes role=""). + - series: 'felis_server_phase{server="survival",phase="Failed",desired="Running"}' + values: '1x20' + - series: 'felis_server_phase{server="creative",phase="Running",desired="Running"}' + values: '1x20' + alert_rule_test: + - eval_time: 9m + alertname: FelisLoginGateDown + exp_alerts: [] + - eval_time: 11m + alertname: FelisLoginGateDown + exp_alerts: + - exp_labels: + severity: critical + server: login + role: login + phase: Starting + desired: Running + exp_annotations: + summary: "the login gate login is Starting" + description: >- + Every player connection passes through the login server first, so no one + can join. The MinecraftServer's conditions carry the reason: + `kubectl -n minecraft describe minecraftserver login` + (troubleshooting §1, §2). + - eval_time: 11m + alertname: FelisSystemServerDown + exp_alerts: + - exp_labels: + severity: warning + server: lobby + role: lobby + phase: Failed + desired: Running + exp_annotations: + summary: "system server lobby (lobby) is Failed" + description: >- + Players who sign in are sent to the lobby; while it is down they stay at + the gate. `kubectl -n minecraft describe minecraftserver lobby` + shows the reason (troubleshooting §1, §2). + - eval_time: 3m + alertname: FelisServerFailed + exp_alerts: [] + - eval_time: 6m + alertname: FelisServerFailed + exp_alerts: + - exp_labels: + severity: warning + server: survival + phase: Failed + desired: Running + exp_annotations: + summary: "server survival is Failed" + description: >- + The operator gave up on this server (a crash loop, an image that will not + pull, a world volume that will not mount). Its conditions carry the + reason: `kubectl -n minecraft describe minecraftserver survival` + (troubleshooting §2). + - name: operator reconcile errors and a stuck pass + interval: 1m + input_series: + # Two failed reconciles a minute from 6m on. + - series: 'controller_runtime_reconcile_errors_total{controller="minecraftserver",job="felis-operator"}' + values: '0x5 0+2x20' + # One pass that started at 5m and never returns. + - series: 'workqueue_longest_running_processor_seconds{name="minecraftserver",controller="minecraftserver",job="felis-operator"}' + values: '0x5 60+60x20' + alert_rule_test: + - eval_time: 5m + alertname: FelisReconcileErrors + exp_alerts: [] + - eval_time: 25m + alertname: FelisReconcileErrors + exp_alerts: + - exp_labels: + severity: warning + exp_annotations: + summary: "the operator failed over 10 reconciles in 15 minutes" + description: >- + Server changes are being retried instead of applied. The felis-operator + log names each failing server and its error. + - eval_time: 12m + alertname: FelisReconcileStuck + exp_alerts: [] + - eval_time: 20m + alertname: FelisReconcileStuck + exp_alerts: + - exp_labels: + severity: critical + exp_annotations: + summary: "an operator reconcile has been running for over 5 minutes" + description: >- + Each reconcile is bounded at 3 minutes, so this one is ignoring its + deadline and holding a worker. The liveness probe restarts the operator + once a pass passes 10 minutes; the log from before the restart shows + where it hung. + - name: world job failures and a late reaper + interval: 1m + input_series: + - series: 'kube_job_failed{namespace="minecraft",job_name="backup-survival-abc",condition="true"}' + values: '0x2 1x10' + - series: 'kube_job_failed{namespace="minecraft",job_name="backup-survival-abc",condition="false"}' + values: '1x2 0x10' + # A build Job in another namespace is FelisImageBuildFailures' business. + - series: 'kube_job_failed{namespace="felis-build",job_name="build-x",condition="true"}' + values: '1x12' + # Evaluation starts at the epoch, so "over a day ago" is a negative timestamp. + - series: 'kube_cronjob_status_last_successful_time{namespace="minecraft",cronjob="felis-reaper"}' + values: '-100000x30' + alert_rule_test: + - eval_time: 1m + alertname: FelisWorldJobFailed + exp_alerts: [] + - eval_time: 5m + alertname: FelisWorldJobFailed + exp_alerts: + - exp_labels: + severity: warning + namespace: minecraft + job_name: backup-survival-abc + condition: "true" + exp_annotations: + summary: "Job backup-survival-abc failed" + description: >- + A world backup, restore or reaper run failed; after a failed backup that + world's newest archive is older than planned. + `kubectl -n minecraft logs job/backup-survival-abc` has the error + (troubleshooting §10). + - eval_time: 5m + alertname: FelisReaperStale + exp_alerts: [] + - eval_time: 15m + alertname: FelisReaperStale + exp_alerts: + - exp_labels: + severity: warning + namespace: minecraft + cronjob: felis-reaper + exp_annotations: + summary: "the world reaper has not succeeded in over 26h" + description: >- + felis-reaper runs daily; idle worlds are neither backed up nor reclaimed + while it fails. `kubectl -n minecraft get jobs --sort-by=.metadata.creationTimestamp` + lists its runs, and the newest one's log + shows why (troubleshooting §10). + - name: a reaper that ran yesterday stays quiet + interval: 1m + input_series: + - series: 'kube_cronjob_status_last_successful_time{namespace="minecraft",cronjob="felis-reaper"}' + values: '-50000x30' + alert_rule_test: + - eval_time: 25m + alertname: FelisReaperStale + exp_alerts: [] diff --git a/deploy/alerts/felis-prometheusrule.yaml b/deploy/alerts/felis-prometheusrule.yaml index 556ef4e..3e65c6a 100644 --- a/deploy/alerts/felis-prometheusrule.yaml +++ b/deploy/alerts/felis-prometheusrule.yaml @@ -37,6 +37,116 @@ spec: observed when readiness is first reached). A start that never completes records nothing — cross-check desiredState=Running servers with no ready phase (troubleshooting §1). + - name: felis.platform.rules + rules: + - alert: FelisOperatorDown + expr: absent(felis_build_info{component="operator"}) + for: 10m + labels: + severity: critical + annotations: + summary: "felis-operator is down or not scraped" + description: >- + No felis_build_info{component="operator"} series for 10 minutes. Without + the operator no server starts, stops or recovers. Check + `kubectl -n felis get deploy felis-operator` and its log; if the pod is + healthy, the felis-operator-metrics Service is not being scraped + (troubleshooting §14). + - alert: FelisAPIDown + expr: absent(felis_build_info{component="api"}) + for: 10m + labels: + severity: critical + annotations: + summary: "felis-api is down or not scraped" + description: >- + No felis_build_info{component="api"} series for 10 minutes. The panel, + sign-in and the proxy's player lookups all go through felis-api. Check + `kubectl -n felis get deploy felis-api` and its log; if the pod is + healthy, the felis-api-internal Service is not being scraped + (troubleshooting §14). + - alert: FelisLoginGateDown + expr: felis_server_phase{role="login",desired="Running",phase!="Running"} == 1 + for: 10m + labels: + severity: critical + annotations: + summary: "the login gate {{ $labels.server }} is {{ $labels.phase }}" + description: >- + Every player connection passes through the login server first, so no one + can join. The MinecraftServer's conditions carry the reason: + `kubectl -n minecraft describe minecraftserver {{ $labels.server }}` + (troubleshooting §1, §2). + - alert: FelisSystemServerDown + expr: felis_server_phase{role!="",role!="login",desired="Running",phase!="Running"} == 1 + for: 10m + labels: + severity: warning + annotations: + summary: "system server {{ $labels.server }} ({{ $labels.role }}) is {{ $labels.phase }}" + description: >- + Players who sign in are sent to the lobby; while it is down they stay at + the gate. `kubectl -n minecraft describe minecraftserver {{ $labels.server }}` + shows the reason (troubleshooting §1, §2). + - alert: FelisServerFailed + expr: felis_server_phase{role="",phase="Failed"} == 1 + for: 5m + labels: + severity: warning + annotations: + summary: "server {{ $labels.server }} is Failed" + description: >- + The operator gave up on this server (a crash loop, an image that will not + pull, a world volume that will not mount). Its conditions carry the + reason: `kubectl -n minecraft describe minecraftserver {{ $labels.server }}` + (troubleshooting §2). + - alert: FelisReconcileErrors + expr: sum(increase(controller_runtime_reconcile_errors_total{controller="minecraftserver"}[15m])) > 10 + for: 5m + labels: + severity: warning + annotations: + summary: "the operator failed over 10 reconciles in 15 minutes" + description: >- + Server changes are being retried instead of applied. The felis-operator + log names each failing server and its error. + - alert: FelisReconcileStuck + expr: max(workqueue_longest_running_processor_seconds{name="minecraftserver"}) > 300 + for: 5m + labels: + severity: critical + annotations: + summary: "an operator reconcile has been running for over 5 minutes" + description: >- + Each reconcile is bounded at 3 minutes, so this one is ignoring its + deadline and holding a worker. The liveness probe restarts the operator + once a pass passes 10 minutes; the log from before the restart shows + where it hung. + - name: felis.jobs.rules + rules: + - alert: FelisWorldJobFailed + expr: kube_job_failed{namespace="minecraft",condition="true"} == 1 + labels: + severity: warning + annotations: + summary: "Job {{ $labels.job_name }} failed" + description: >- + A world backup, restore or reaper run failed; after a failed backup that + world's newest archive is older than planned. + `kubectl -n minecraft logs job/{{ $labels.job_name }}` has the error + (troubleshooting §10). + - alert: FelisReaperStale + expr: time() - kube_cronjob_status_last_successful_time{namespace="minecraft",cronjob="felis-reaper"} > 26 * 3600 + for: 10m + labels: + severity: warning + annotations: + summary: "the world reaper has not succeeded in over 26h" + description: >- + felis-reaper runs daily; idle worlds are neither backed up nor reclaimed + while it fails. `kubectl -n minecraft get jobs --sort-by=.metadata.creationTimestamp` + lists its runs, and the newest one's log + shows why (troubleshooting §10). - name: felis.node.rules rules: - alert: FelisNodeDiskSpaceLow diff --git a/deploy/bootstrap.sh b/deploy/bootstrap.sh index 62f5eb2..a09a291 100644 --- a/deploy/bootstrap.sh +++ b/deploy/bootstrap.sh @@ -252,6 +252,13 @@ GOROOT_DIR="/opt/felis/go" NANO_SERVICE="/etc/systemd/system/felis-nano.service" DB_BACKUP_SERVICE="/etc/systemd/system/felis-db-backup.service" DB_BACKUP_TIMER="/etc/systemd/system/felis-db-backup.timer" +WATCHDOG_SERVICE="/etc/systemd/system/felis-watchdog.service" +WATCHDOG_TIMER="/etc/systemd/system/felis-watchdog.timer" +WATCHDOG_STATE="/var/lib/felis/watchdog/state.json" +# While this marker holds a future Unix time, felis watchdog mails nothing: an install +# restarts the control plane and the system servers on purpose. cleanup removes it; the +# time in it is the backstop for an installer killed before its EXIT trap runs. +WATCHDOG_QUIET_FILE="/run/felis/watchdog-quiet-until" VELOCITY_DIR="/opt/felis/velocity" VELOCITY_USER="felis-velocity" VELOCITY_SERVICE="/etc/systemd/system/felis-velocity.service" @@ -336,6 +343,9 @@ cleanup() { for path in "${TEMP_PATHS[@]-}"; do [ -n "$path" ] && rm -rf -- "$path" || true done + # A failed install leaves something broken the owners should hear about, so the + # watchdog speaks again the moment the installer exits, however it exits. + rm -f -- "$WATCHDOG_QUIET_FILE" 2>/dev/null || true } remember_temp() { TEMP_PATHS+=("$1"); } @@ -2415,6 +2425,60 @@ run_migrations() { ok "migrations applied" } +# Silence felis watchdog's mail for the rest of this install (see WATCHDOG_QUIET_FILE). +# Two hours covers a slow source build; cleanup lifts it as soon as the installer exits. +quiet_watchdog() { + install -d -m 0755 "$(dirname "$WATCHDOG_QUIET_FILE")" + printf '%s\n' "$(( $(date +%s) + 7200 ))" > "$WATCHDOG_QUIET_FILE" +} + +# The platform watchdog: every two minutes it checks the control plane, the login gate, +# the fleet, PostgreSQL, the game proxy, the database backups and the host's disks and +# memory, and mails the owners (their verified addresses, over the [smtp] relay) what +# has stayed wrong long enough to matter. It runs on the host so a k3s that is down is +# still reported. The first run happens now, so a broken unit shows up in this install. +install_watchdog_timer() { + local disks="/,/var/lib/rancher/k3s,/var/lib/postgresql,/var/lib/felis" path + for path in "$FELIS_WORLDS_HOST_PATH" "$FELIS_ARCHIVE_LOCAL_PATH" "$FELIS_DB_BACKUP_DIR"; do + if [ -n "$path" ]; then disks="${disks},${path}"; fi + done + install -d -m 0700 "$(dirname "$WATCHDOG_STATE")" + cat > "$WATCHDOG_SERVICE" < "$WATCHDOG_TIMER" <&2 || true + warn "the first watchdog run failed (log above); nothing will be mailed until it runs: sudo systemctl start felis-watchdog.service" + fi +} + # The daily database backup. The first run happens now, so a broken pipeline (pg_dump # missing, directory unwritable) shows up in this install rather than in the first # restore someone needs. @@ -3019,6 +3083,7 @@ main() { main_nano return fi + quiet_watchdog pause_package_background_timers detect_node_ip ensure_swap @@ -3073,6 +3138,8 @@ main() { install_velocity # After deploy_bundle: the bundle's MinecraftServer export reads the cluster. install_db_backup_timer + # Last: its first run should see the platform as this install leaves it. + install_watchdog_timer mark_bootstrap_done summary } diff --git a/deploy/bootstrap_test.sh b/deploy/bootstrap_test.sh index f42a930..ada0ade 100644 --- a/deploy/bootstrap_test.sh +++ b/deploy/bootstrap_test.sh @@ -1152,6 +1152,75 @@ expect "a failed first backup shows its log" "JOURNAL: pg_dump: connection refus expect "a failed first backup is a loud warning" "WARN: the first database backup failed" "$out" rm -rf "$tdir" +wblock="$(awk '/^install_watchdog_timer\(\) \{/,/^}/' "$BS")" +[ -n "$wblock" ] || { echo "FAIL: no install_watchdog_timer found in $BS"; exit 1; } +[ "$(printf '%s\n' "$wblock" | wc -l)" -lt 60 ] \ + || { echo "FAIL: the extracted block is not install_watchdog_timer -- did its closing brace move?"; exit 1; } +qblock="$(awk '/^quiet_watchdog\(\) \{/,/^}/' "$BS")" +[ -n "$qblock" ] || { echo "FAIL: no quiet_watchdog found in $BS"; exit 1; } + +tdir="$(mktemp -d)" +run_watchdog_timer() { # $1: exit status of the first run, $2: FELIS_WORLDS_HOST_PATH + FIRST="$1" FELIS_WORLDS_HOST_PATH="$2" WATCHDOG_SERVICE="$tdir/felis-watchdog.service" WATCHDOG_TIMER="$tdir/felis-watchdog.timer" \ + WATCHDOG_STATE="$tdir/watchdog/state.json" WATCHDOG_QUIET_FILE=/run/felis/watchdog-quiet-until \ + FELIS_DB_BACKUP_DIR=/var/lib/felis/db-backups FELIS_ARCHIVE_LOCAL_PATH=/var/lib/felis/archives FELIS_GAME_PORT=25577 \ + HOST_BIN=/usr/local/bin/felis STATE_DIR=/etc/felis bash -c ' + set -Eeuo pipefail + ok() { printf "OK: %s\n" "$*"; }; warn() { printf "WARN: %s\n" "$*"; } + systemctl() { printf "SYSTEMCTL: %s\n" "$*"; [ "$1" != start ] || return "$FIRST"; } + journalctl() { printf "JOURNAL: parse /etc/felis/felis.host.toml\n"; } + '"$wblock"' + install_watchdog_timer' 2>&1 +} + +out="$(run_watchdog_timer 0 "")" +unit="$(cat "$tdir/felis-watchdog.service")" +timer="$(cat "$tdir/felis-watchdog.timer")" +expect "the watchdog runs the host binary against the host config, dialing the proxy's port" \ + "ExecStart=/usr/local/bin/felis watchdog -config /etc/felis/felis.host.toml -state $tdir/watchdog/state.json -quiet-file /run/felis/watchdog-quiet-until -backup-dir /var/lib/felis/db-backups -proxy-addr 127.0.0.1:25577 -disk-paths /,/var/lib/rancher/k3s,/var/lib/postgresql,/var/lib/felis,/var/lib/felis/archives,/var/lib/felis/db-backups" "$unit" +expect "a wedged run is killed before the next one is due twice over" "TimeoutStartSec=3min" "$unit" +expect "the watchdog runs every two minutes" "OnUnitActiveSec=2min" "$timer" +expect "the watchdog starts soon after boot" "OnBootSec=3min" "$timer" +expect "the watchdog timer is enabled" "SYSTEMCTL: enable --now felis-watchdog.timer" "$out" +expect "the first watchdog run happens during the install" "SYSTEMCTL: start felis-watchdog.service" "$out" +expect "a working first run is reported" "OK: watchdog: checks every 2 minutes" "$out" +if [ "$(stat -c %a "$tdir/watchdog" 2>/dev/null || stat -f %Lp "$tdir/watchdog")" = 700 ]; then + echo "PASS the watchdog state directory is private (it caches the relay password)" +else + echo "FAIL the watchdog state directory must be 0700"; fails=$((fails + 1)) +fi + +out="$(run_watchdog_timer 0 /srv/worlds)" +expect "a custom worlds root is watched for free space" "-disk-paths /,/var/lib/rancher/k3s,/var/lib/postgresql,/var/lib/felis,/srv/worlds," "$(cat "$tdir/felis-watchdog.service")" + +out="$(run_watchdog_timer 1 "")" +expect "a failed first watchdog run shows its log" "JOURNAL: parse /etc/felis/felis.host.toml" "$out" +expect "a failed first watchdog run is a loud warning" "WARN: the first watchdog run failed" "$out" + +out="$(WATCHDOG_QUIET_FILE="$tdir/run/quiet" bash -c ' + set -Eeuo pipefail + '"$qblock"' + quiet_watchdog; cat "$WATCHDOG_QUIET_FILE"; date +%s' 2>&1)" +until_ts="$(printf '%s\n' "$out" | sed -n 1p)"; now_ts="$(printf '%s\n' "$out" | sed -n 2p)" +if [ -n "$until_ts" ] && [ "$((until_ts - now_ts))" -ge 3600 ] && [ "$((until_ts - now_ts))" -le 7200 ]; then + echo "PASS the install quiets the watchdog for a bounded while" +else + echo "FAIL quiet_watchdog wrote '$until_ts' at $now_ts, want now+1h..2h"; fails=$((fails + 1)) +fi +rm -rf "$tdir" + +cblock="$(awk '/^cleanup\(\) \{/,/^}/' "$BS")" +case "$cblock" in + *'rm -f -- "$WATCHDOG_QUIET_FILE"'*) echo "PASS the installer lifts the watchdog's quiet period on exit" ;; + *) echo "FAIL cleanup must remove WATCHDOG_QUIET_FILE, or a failed install stays silent"; fails=$((fails + 1)) ;; +esac +main_block="$(awk '/^main\(\) \{/,/^}/' "$BS")" +case "$main_block" in + *quiet_watchdog*pause_package_background_timers*install_db_backup_timer*install_watchdog_timer*mark_bootstrap_done*) + echo "PASS main quiets the watchdog first and installs it last" ;; + *) echo "FAIL main must call quiet_watchdog before any restart and install_watchdog_timer after the backup timer"; fails=$((fails + 1)) ;; +esac + # --------------------------------------------------------------------------------------- if [ "$fails" -eq 0 ]; then echo "ALL PASS" diff --git a/docs/troubleshooting.md b/docs/troubleshooting.md index b08558a..2a67578 100644 --- a/docs/troubleshooting.md +++ b/docs/troubleshooting.md @@ -964,11 +964,75 @@ space. --- -## 14. Metrics for diagnosis (spec §23) +## 14. Health alerts, and metrics for diagnosis (spec §23) + +Every full install runs `felis watchdog` from `felis-watchdog.timer`, which +mails the owners when something breaks, with no monitoring stack needed. The +metrics and Prometheus rules below are for a deployment that also runs its own +Prometheus. + +### The watchdog: what mails the owners + +Every two minutes the host checks: + +| Check | Mailed after | Severity | +|---|---|---| +| `felis-api`, `felis-operator` or `registry` Deployment has no available pod (or is missing) | 5 min | critical | +| Kubernetes API unreachable (k3s down): the cluster checks below are then unknown and keep their state | 5 min | critical | +| Login gate not `Running` while it should be | 10 min | critical | +| Other system servers (lobby) not `Running` while they should be | 10 min | warning | +| A user server in `Failed` (§2) | 5 min | warning | +| A backup, restore or reaper Job failed in the last 24h | at once, once | warning | +| The reaper CronJob last succeeded over 26h ago (§10) | at once | warning | +| Node `NotReady`, or kubelet reports Disk/Memory/PID pressure (§13b) | 2–5 min | critical | +| PostgreSQL unreachable | 3 min | critical | +| The game proxy (`felis-velocity`) refuses connections on the game port | 3 min | critical | +| Newest control-plane database backup over 26h old, or none (§16) | 10 min | critical | +| A watched filesystem below 15% free (below 5%: critical) | 15 min (5 min) | warning | +| Host memory available below 10% | 15 min | warning | + +How it mails: + +- **Recipients** are the verified email addresses of enabled owner accounts. + The relay is the `[smtp]` one sign-in codes use. +- **One mail per run**, holding everything that came due: new problems, a + daily reminder for each problem still open, and a resolved notice once a + problem has stayed gone for 10 minutes. A condition that heals before its + delay is never mailed; one that turns critical is mailed again at once. +- **During an install** nothing is mailed. `bootstrap.sh` writes + `/run/felis/watchdog-quiet-until` and removes it when it exits. +- **Caching:** the relay password and the recipient list are cached in + `/var/lib/felis/watchdog/state.json` (root-only). An outage of PostgreSQL or + of the API server can therefore still be mailed. +- **No relay or no verified owner address:** each alert is written to the + journal only. + +Commands: + +```bash +journalctl -u felis-watchdog -n 40 # every check of the latest runs +sudo felis watchdog -dry-run # run the checks now; mail nothing, change nothing +systemctl list-timers felis-watchdog.timer # when it last and next runs +``` + +A healthy run logs `every check passed`. Otherwise it logs one line per +finding, and the mail's subject once one is sent. + +A `-dry-run` from a shell uses the command's defaults, and those do not include +the game-proxy check. The unit carries `-proxy-addr 127.0.0.1:` and +the disk list the install chose. `systemctl cat felis-watchdog` shows both. + +### Metrics All four mandated metrics have real producers; scrape them when triaging: - `felis_servers_total` — managed server count. +- `felis_server_phase{server,role,phase,desired}` — 1 for each server's + current phase. `role` is `login`/`lobby` for system servers and empty for + user servers. `desired` tells a server that is down on purpose from one that + failed to come up. +- `felis_build_info{component,version}` — 1 on the process serving it + (`operator` or `api`). Its absence is how an alert tells which process is down. - `felis_start_duration_seconds` — histogram, observed once per start when readiness is first reached (`ReadySignalAt − StartRequestedAt`). A start that never completes (§1) contributes **nothing** here — absence of observations is @@ -980,12 +1044,16 @@ All four mandated metrics have real producers; scrape them when triaging: ### Scraping -The series come from two processes: +The series come from two processes. Both Services carry the +`prometheus.io/scrape|port|path` annotations, so a Prometheus that discovers +annotated Service endpoints picks them up as is. -- `felis-operator` pod `:8080/metrics` — `felis_servers_total`, - `felis_start_duration_seconds` (no Service; scrape pod-scoped, e.g. a - PodMonitor targeting port `metrics`). +- `felis-operator` `:8080/metrics` (Service `felis-operator-metrics`) — + `felis_servers_total`, `felis_server_phase`, `felis_start_duration_seconds`, + `felis_build_info{component="operator"}`, and controller-runtime's + `controller_runtime_reconcile_*` / `workqueue_*` series. - `felis-api` internal face `:8081/metrics` (Service `felis-api-internal`) — + `felis_build_info{component="api"}`, `felis_image_build_failures_total`, and the sign-in series of §17 (`felis_mail_total`, `felis_rate_limited_total`, `felis_auth_otp_lockouts_total`, `felis_auth_failures_total`, @@ -999,11 +1067,27 @@ The series come from two processes: ### Alert rules -`deploy/alerts/` ships ready-made rules: build failures, slow starts, node -disk/memory thresholds, the kubelet `DiskPressure` condition, control-plane -database backup freshness (§16; needs node-exporter's textfile collector), and -sign-in abuse: the mail budget, relay failures, throttled floods, account -code locks and the refused sign-in rate, plus lost audit rows (§17). +`deploy/alerts/` ships ready-made rules for these groups: + +- **Control plane:** `felis-operator` or `felis-api` down or unscraped + (`FelisOperatorDown`, `FelisAPIDown`). +- **Servers:** the login gate or another system server not running + (`FelisLoginGateDown`, `FelisSystemServerDown`), and a user server in + `Failed` (`FelisServerFailed`). +- **Operator reconcile:** errors piling up (`FelisReconcileErrors`), and a + pass stuck past its 3-minute bound (`FelisReconcileStuck`). The operator's + liveness probe restarts a pod whose pass passes 10 minutes; its readiness + waits for the informer caches to sync. +- **World Jobs (kube-state-metrics):** a failed backup/restore/reaper Job + (`FelisWorldJobFailed`) and a reaper that has not succeeded in 26h + (`FelisReaperStale`). +- **Builds and starts:** build failures and slow starts. +- **Node:** disk and memory thresholds, and the kubelet `DiskPressure` + condition. +- **Database backups:** control-plane backup freshness (§16; needs + node-exporter's textfile collector). +- **Sign-in abuse (§17):** the mail budget, relay failures, throttled floods, + account code locks and the refused sign-in rate, plus lost audit rows. - Plain Prometheus: add `felis-alerts.yaml` to `rule_files`. Check and unit-test it standalone with `promtool check rules felis-alerts.yaml` and @@ -1380,6 +1464,8 @@ for 10 seconds (the Free plan's limits). | PVC left behind after delete | §13 | | Node out of disk; pods evicted / ImagePullBackOff | §13b | | Which metric to scrape | §14 | +| An alert mail from the watchdog; nothing is mailed when something breaks | §14 | +| `FelisOperatorDown` / `FelisAPIDown` / `FelisLoginGateDown` / `FelisReconcileStuck` | §14, §1, §2 | | Upgrade / roll back a bad control-plane image | §15 | | `image_change_unconfirmed` / `image_not_in_registry` / `registry_unavailable`; move a world to a newer Minecraft | §15b | | Database backup overdue / `FelisDBBackupStale` / panel shows 从未备份 | §16 | diff --git a/internal/metrics/alerts_test.go b/internal/metrics/alerts_test.go new file mode 100644 index 0000000..98b76e6 --- /dev/null +++ b/internal/metrics/alerts_test.go @@ -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) + } + } + } + } +} diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index d8bb69a..96208dc 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -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, } } diff --git a/internal/metrics/metrics_test.go b/internal/metrics/metrics_test.go index 9a96cbc..96206bb 100644 --- a/internal/metrics/metrics_test.go +++ b/internal/metrics/metrics_test.go @@ -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 diff --git a/internal/operator/gauge.go b/internal/operator/gauge.go index ca51b6a..145cd2d 100644 --- a/internal/operator/gauge.go +++ b/internal/operator/gauge.go @@ -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. diff --git a/internal/operator/gauge_test.go b/internal/operator/gauge_test.go index 75f392a..9f7cd0b 100644 --- a/internal/operator/gauge_test.go +++ b/internal/operator/gauge_test.go @@ -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) + } + } +} diff --git a/internal/operator/reconciler.go b/internal/operator/reconciler.go index 4bb1387..a85f856 100644 --- a/internal/operator/reconciler.go +++ b/internal/operator/reconciler.go @@ -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. diff --git a/internal/operator/watch.go b/internal/operator/watch.go new file mode 100644 index 0000000..91692c0 --- /dev/null +++ b/internal/operator/watch.go @@ -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 +} diff --git a/internal/operator/watch_test.go b/internal/operator/watch_test.go new file mode 100644 index 0000000..3c0fe9b --- /dev/null +++ b/internal/operator/watch_test.go @@ -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) + } +} diff --git a/internal/platform/workloads.go b/internal/platform/workloads.go index 3fb7d18..b9102a1 100644 --- a/internal/platform/workloads.go +++ b/internal/platform/workloads.go @@ -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://: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), diff --git a/internal/platform/workloads_test.go b/internal/platform/workloads_test.go index 1ff8f02..ad3c7c0 100644 --- a/internal/platform/workloads_test.go +++ b/internal/platform/workloads_test.go @@ -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") + } +} diff --git a/internal/watchdog/probes.go b/internal/watchdog/probes.go new file mode 100644 index 0000000..f5f9388 --- /dev/null +++ b/internal/watchdog/probes.go @@ -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 ": ". +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]) +} diff --git a/internal/watchdog/probes_test.go b/internal/watchdog/probes_test.go new file mode 100644 index 0000000..e04c027 --- /dev/null +++ b/internal/watchdog/probes_test.go @@ -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)) + } +} diff --git a/internal/watchdog/watchdog.go b/internal/watchdog/watchdog.go new file mode 100644 index 0000000..62227d7 --- /dev/null +++ b/internal/watchdog/watchdog.go @@ -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) +} diff --git a/internal/watchdog/watchdog_test.go b/internal/watchdog/watchdog_test.go new file mode 100644 index 0000000..2914c29 --- /dev/null +++ b/internal/watchdog/watchdog_test.go @@ -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{"owner@example.com"}} + 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) + } +}