Unverified Commit 82ffa55a authored by Lemon-miaow's avatar Lemon-miaow
Browse files

fix(watchdog): 外部心跳、巡检失败兜底告警、状态文件损坏自动移开

parent 149ab0be
Loading
Loading
Loading
Loading
+1 −1
Changes for cmd/felis/hostcreds_test.go: 1 added line, 1 removed line.
Original line number Diff line number Diff line
@@ -183,7 +183,7 @@ password_ref = "FELIS_TEST_UNSET_RELAY_PW"
	var stdout, stderr bytes.Buffer
	cmdWatchdog([]string{
		"-config", cfgPath, "-state", statePath, "-quiet-file", filepath.Join(dir, "quiet"),
		"-backup-dir", "", "-disk-paths", dir, "-smtp-password-file", pwPath,
		"-backup-dir", "", "-disk-paths", dir, "-smtp-password-file", pwPath, "-heartbeat-file", filepath.Join(dir, "no-heartbeat"),
	}, &stdout, &stderr)
	if !strings.Contains(stdout.String(), "kube-api") {
		t.Fatalf("the run found the API server up; the test needs it down (stdout %s)", stdout.String())
+211 −28
Changes for cmd/felis/watchdog.go: 211 added lines, 28 removed lines.
Original line number Diff line number Diff line
@@ -8,11 +8,13 @@ import (
	"fmt"
	"io"
	"net"
	"net/http"
	"os"
	"strings"
	"time"

	"felis.lolicon.best/internal/config"
	"felis.lolicon.best/internal/mail"
	"felis.lolicon.best/internal/offsite"
	"felis.lolicon.best/internal/platform"
	"felis.lolicon.best/internal/store"
@@ -46,25 +48,35 @@ func cmdWatchdog(args []string, stdout, stderr io.Writer) int {
	offsiteStatus := fs.String("offsite-status", offsite.DefaultStatusFile, "the record `felis offsite sync` leaves, checked when [offsite] is configured")
	toolsStatus := fs.String("build-tools-status", defaultBuildToolsStatus, "the record `felis mirror-build-tools` leaves, checked when builds scan against the registry's DB copy")
	dryRun := fs.Bool("dry-run", false, "print every finding and the mail that is due; send nothing and keep the state as it was")
	heartbeatFile := fs.String("heartbeat-file", defaultHeartbeatFile, "file holding the heartbeat URL each run pings, a dead man's switch at a monitoring service that alerts when the pings stop (no file pings nothing)")
	unitFailed := fs.Bool("unit-failed", false, "report a failed run of felis-watchdog.service instead of checking; felis-watchdog-failed.service runs this through OnFailure=")
	if err := fs.Parse(args); err != nil {
		if errors.Is(err, flag.ErrHelp) {
			return 0
		}
		return 2
	}
	now := time.Now()
	if *unitFailed {
		return watchdogUnitFailed(unitFailedRun{
			cfgPath: *cfgPath, statePath: *statePath, quietPath: *quietPath,
			offsiteStatus: *offsiteStatus, heartbeatFile: *heartbeatFile,
			result: os.Getenv("MONITOR_SERVICE_RESULT"), exitStatus: os.Getenv("MONITOR_EXIT_STATUS"),
			send: watchdogSender, client: http.DefaultClient, now: now,
		}, stdout, stderr)
	}
	cfg, err := config.Load(*cfgPath)
	if err != nil {
		fmt.Fprintf(stderr, "felis watchdog: %v\n", err)
		return 1
	}
	state, err := watchdog.LoadState(*statePath)
	state, aside, err := watchdog.RecoverState(*statePath, now)
	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) {
@@ -72,6 +84,11 @@ func cmdWatchdog(args []string, stdout, stderr io.Writer) int {
			report.Findings = append(report.Findings, *f)
		}
	}
	if aside != "" {
		fmt.Fprintf(stderr, "felis watchdog: %s was unreadable; moved it to %s and started over\n", *statePath, aside)
		f := watchdog.StateSetAside(aside)
		add(&f)
	}

	// The cluster: one unreachable API server stands in for every check behind it.
	minecraftNS := cfg.K8s.Namespace
@@ -97,6 +114,7 @@ func cmdWatchdog(args []string, stdout, stderr io.Writer) int {
		}
		refreshSMTPPassword(ctx, *smtpPasswordFile, secrets, *controlNS, state, stderr)
	}
	state.Relay = cachedRelay(cfg.SMTP)

	if recipients, err := ownerEmails(ctx, cfg.Database.URL); err != nil {
		f := watchdog.PostgresDown(err)
@@ -136,45 +154,207 @@ func cmdWatchdog(args []string, stdout, stderr io.Writer) int {
	plan := state.Observe(report, now)
	host, _ := os.Hostname()
	subject, body := plan.Message(host, now)
	beat := heartbeat{
		standby: standsBy(cfg.Offsite.Enabled(), *offsiteStatus),
		quiet:   now.Before(watchdog.QuietUntil(*quietPath)),
	}
	if beat.url, err = readHeartbeatURL(*heartbeatFile); err != nil {
		fmt.Fprintf(stderr, "felis watchdog: %v; pinging no heartbeat\n", err)
	}
	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"))
		}
		if beat.url != "" {
			fmt.Fprintf(stdout, "felis watchdog: a run pings the heartbeat at %s\n", redactURL(beat.url))
		}
		return 0
	}

	save := func() int {
		if err := watchdog.SaveState(*statePath, state); err != nil {
			fmt.Fprintf(stderr, "felis watchdog: save state: %v\n", err)
	m := configMailer(cfg.SMTP, state.SMTPPassword, watchdogSender)
	unheard, mailFailed := m.deliver(ctx, state, plan, subject, body, mailHold(*quietPath, cfg.Offsite.Enabled(), *offsiteStatus, now), now, stdout, stderr)
	saveErr := watchdog.SaveState(*statePath, state)
	if saveErr != nil {
		fmt.Fprintf(stderr, "felis watchdog: save state: %v\n", saveErr)
	}
	beat.report = failureReport(unheard, mailFailed, saveErr, state.Open(), report)
	beat.fail = beat.report != ""
	beat.send(http.DefaultClient, stdout, stderr)
	if mailFailed || saveErr != nil {
		return 1
	}
	return 0
}
	if plan.Empty() {
		return save()

// failureReport is what the heartbeat's failure ping carries, "" when the run
// pings success: the alerts this run knows of reach no one (a mail that
// failed, or no relay or recipient while something is open), or the state did
// not save and the next run mails the same alerts again.
func failureReport(unheard string, mailFailed bool, saveErr error, open bool, r watchdog.Report) string {
	var why []string
	if unheard != "" && (mailFailed || open) {
		why = append(why, "the alerts reach no one: "+unheard)
	}
	if saveErr != nil {
		why = append(why, "the watchdog state did not save: "+saveErr.Error())
	}
	if len(why) == 0 {
		return ""
	}
	if hold := mailHold(*quietPath, cfg.Offsite.Enabled(), *offsiteStatus, now); hold != "" {
		fmt.Fprintf(stdout, "felis watchdog: %s; holding this mail: %s\n", hold, subject)
		return save()
	return strings.Join(why, "\n") + "\n\n" + findingLines(r)
}

// findingLines is the report as the journal shows it.
func findingLines(r watchdog.Report) string {
	var b strings.Builder
	for _, f := range r.Findings {
		fmt.Fprintf(&b, "[%s] %s: %s\n", f.Severity, f.Key, f.SummaryEN)
	}
	return b.String()
}

// mailer is how a run reaches the owners.
type mailer struct {
	relay    *watchdog.Relay // nil: no [smtp] relay
	password string
	send     func(*mail.SMTP) alertSender
}

// deliver mails plan to the owners unless hold says why it waits, and commits
// it once it reached them, or once it is logged because nothing can reach them.
// unheard is why the owners hear nothing of this run's alerts, "" when they do;
// mailFailed is a mail that did not go out, left uncommitted so the same
// alerts come due again next run.
func (m mailer) deliver(ctx context.Context, state *watchdog.State, plan watchdog.Plan, subject, body, hold string, now time.Time, stdout, stderr io.Writer) (unheard string, mailFailed bool) {
	switch {
	case cfg.SMTP.Host == "":
		fmt.Fprintf(stdout, "felis watchdog: no [smtp] relay configured, so this is logged only: %s\n", subject)
	case m.relay == nil:
		unheard = "no [smtp] relay is configured"
	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)
		unheard = "no owner account has a verified email"
	}
	switch {
	case plan.Empty():
	case hold != "":
		fmt.Fprintf(stdout, "felis watchdog: %s; holding this mail: %s\n", hold, subject)
	case unheard != "":
		fmt.Fprintf(stdout, "felis watchdog: %s, so this is logged only: %s\n", unheard, subject)
		state.Commit(plan, now)
	default:
		if err := sendAlert(ctx, cfg, state, subject, body); err != nil {
			// Not committed: the same alerts come due again next run.
		relay := &mail.SMTP{Host: m.relay.Host, Port: m.relay.Port, From: m.relay.From, Username: m.relay.Username, Password: m.password, RequireTLS: m.relay.RequireTLS}
		if err := sendAlert(ctx, m.send(relay), state.Recipients, subject, body); err != nil {
			fmt.Fprintf(stderr, "felis watchdog: mail %q: %v\n", subject, err)
			save()
			return 1
			return "the alert mail failed: " + err.Error(), true
		}
		fmt.Fprintf(stdout, "felis watchdog: mailed %s: %s\n", strings.Join(state.Recipients, ", "), subject)
	}
		state.Commit(plan, now)
	return save()
	}
	return unheard, false
}

// configMailer reaches the owners through the relay felis.toml configures.
func configMailer(c config.SMTPConfig, cachedPassword string, send func(*mail.SMTP) alertSender) mailer {
	return mailer{relay: cachedRelay(c), password: smtpPassword(c, cachedPassword), send: send}
}

// cachedRelay is what State.Relay keeps of c, nil when no relay is configured.
func cachedRelay(c config.SMTPConfig) *watchdog.Relay {
	if c.Host == "" {
		return nil
	}
	return &watchdog.Relay{Host: c.Host, Port: c.Port, From: c.From, Username: c.Username, RequireTLS: c.TLSRequired()}
}

// smtpPassword is the relay password a mail signs in with: the env var [smtp]
// password_ref names when it is set, else the one the state caches.
func smtpPassword(c config.SMTPConfig, cached string) string {
	if ref := c.PasswordRef; ref != "" && os.Getenv(ref) != "" {
		return os.Getenv(ref)
	}
	return cached
}

// unitFailedRun is one run of felis-watchdog-failed.service.
type unitFailedRun struct {
	cfgPath, statePath, quietPath, offsiteStatus, heartbeatFile string
	// result and exitStatus are what systemd hands an OnFailure= unit
	// (MONITOR_SERVICE_RESULT, MONITOR_EXIT_STATUS; systemd 251 and later).
	result, exitStatus string
	send               func(*mail.SMTP) alertSender
	client             *http.Client
	now                time.Time
}

// watchdogUnitFailed is felis-watchdog-failed.service, which systemd starts
// through OnFailure= when a run of felis-watchdog.service fails: a crash, a
// felis.toml that no longer loads, a hang past the unit's timeout, a mail
// that did not go out. Such a run checks and mails nothing, so this records
// the failure as the alert watchdog/run, due after five failed runs in a row
// and cleared by the next run that succeeds; mails it through the relay the
// last good run cached when felis.toml does not load; and pings the
// heartbeat's failure endpoint. Every other alert keeps its state.
func watchdogUnitFailed(r unitFailedRun, stdout, stderr io.Writer) int {
	detail := failureDetail(r.result, r.exitStatus)
	cfg, cfgErr := config.Load(r.cfgPath)
	if cfgErr != nil {
		detail += "; " + cfgErr.Error()
	}
	fmt.Fprintf(stdout, "felis watchdog: felis-watchdog.service failed: %s\n", detail)
	offsiteOn := cfgErr != nil || cfg.Offsite.Enabled()
	beat := heartbeat{
		fail:    true,
		report:  "felis-watchdog.service failed: " + detail,
		standby: standsBy(offsiteOn, r.offsiteStatus),
		quiet:   r.now.Before(watchdog.QuietUntil(r.quietPath)),
	}
	var err error
	if beat.url, err = readHeartbeatURL(r.heartbeatFile); err != nil {
		fmt.Fprintf(stderr, "felis watchdog: %v; pinging no heartbeat\n", err)
	}
	state, err := watchdog.LoadState(r.statePath)
	if err != nil {
		// The next run that gets that far moves a state that does not parse
		// aside (watchdog.RecoverState).
		fmt.Fprintf(stderr, "felis watchdog: %v; mailing nothing\n", err)
		beat.report += "\nThe watchdog state does not load either, so nothing was mailed: " + err.Error()
		beat.send(r.client, stdout, stderr)
		return 1
	}
	plan := state.Observe(watchdog.Report{Findings: []watchdog.Finding{watchdog.WatchdogFailed(detail)}, Unknown: []string{""}}, r.now)
	host, _ := os.Hostname()
	subject, body := plan.Message(host, r.now)
	m := mailer{relay: state.Relay, password: state.SMTPPassword, send: r.send}
	if cfgErr == nil {
		m = configMailer(cfg.SMTP, state.SMTPPassword, r.send)
	}
	ctx, cancel := context.WithTimeout(context.Background(), time.Minute)
	defer cancel()
	unheard, mailFailed := m.deliver(ctx, state, plan, subject, body, mailHold(r.quietPath, offsiteOn, r.offsiteStatus, r.now), r.now, stdout, stderr)
	code := 0
	if mailFailed {
		code = 1
	}
	if err := watchdog.SaveState(r.statePath, state); err != nil {
		fmt.Fprintf(stderr, "felis watchdog: save state: %v\n", err)
		code = 1
	}
	if unheard != "" {
		beat.report += "\nThe alerts reach no one: " + unheard
	}
	beat.send(r.client, stdout, stderr)
	return code
}

// failureDetail names how felis-watchdog.service failed.
func failureDetail(result, exitStatus string) string {
	switch {
	case result == "":
		return "systemd named no cause (journalctl -u felis-watchdog -n 50)"
	case exitStatus == "":
		return "result " + result
	}
	return "result " + result + ", exit status " + exitStatus
}

// mailHold is why this run's mail waits, "" when it goes out: the installer's
@@ -288,21 +468,24 @@ func proxyFinding(ctx context.Context, addr string) *watchdog.Finding {
	}
}

// alertSender mails one alert; smtpSender is the real one.
type alertSender func(ctx context.Context, to, subject, body string) error

func smtpSender(relay *mail.SMTP) alertSender { return relay.SendNotice }

// watchdogSender is what a run mails through; tests stand a recorder in.
var watchdogSender = smtpSender

// 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 := smtpRelay(cfg.SMTP, password)
func sendAlert(ctx context.Context, send alertSender, recipients []string, subject, body string) error {
	var errs []error
	for _, to := range state.Recipients {
		if err := relay.SendNotice(ctx, to, subject, body); err != nil {
	for _, to := range recipients {
		if err := send(ctx, to, subject, body); err != nil {
			errs = append(errs, fmt.Errorf("%s: %w", to, err))
		}
	}
	if len(errs) == len(state.Recipients) {
	if len(errs) == len(recipients) {
		return errors.Join(errs...)
	}
	return nil
+164 −0
Changes for cmd/felis/watchdog_heartbeat.go: 164 added lines, 0 removed lines.
Original line number Diff line number Diff line
package main

import (
	"context"
	"errors"
	"fmt"
	"io"
	"net/http"
	"net/url"
	"os"
	"strings"
	"time"

	"felis.lolicon.best/internal/offsite"
)

// defaultHeartbeatFile holds the heartbeat URL deploy/bootstrap.sh writes from
// FELIS_WATCHDOG_HEARTBEAT_URL. It lives in /etc/felis, so a host rebuilt from
// a bundle pings the same check once it takes the off-site bucket over.
const defaultHeartbeatFile = "/etc/felis/watchdog-heartbeat-url"

const (
	// heartbeatTimeout bounds one ping; a monitoring service slower than that
	// is as good as down for the run.
	heartbeatTimeout = 10 * time.Second
	// heartbeatReportMax caps the report a failure ping carries: monitoring
	// services keep about the first 10 KB of a ping's body.
	heartbeatReportMax = 10000
)

// heartbeat is the ping one run sends to a dead man's switch at a monitoring
// service (Healthchecks.io, and the services that copy its API), which alerts
// its own users when the pings stop: the host down, the timer gone, the
// watchdog failing before it can mail. No check on the host can report those.
// A run whose alerts reach the owners GETs url; one whose alerts reach no one
// POSTs report to url/fail, or withholds the ping when url has a query, where
// no /fail can be added and the missing ping trips the check instead.
type heartbeat struct {
	url    string // "" pings nothing
	fail   bool
	report string
	// standby is a host standing by for another host's off-site bucket: a
	// rehearsal, or a rebuild not taken over. It pings nothing, or it would
	// keep the check of the host that writes the bucket green after that host
	// died.
	standby bool
	// quiet is the installer's quiet window, which restarts things on
	// purpose: no failure is pinged.
	quiet bool
}

// send pings the heartbeat and logs the outcome.
func (b heartbeat) send(cl *http.Client, stdout, stderr io.Writer) {
	switch {
	case b.url == "":
		return
	case b.standby:
		fmt.Fprintln(stdout, "felis watchdog: this host stands by for the off-site bucket; pinging no heartbeat (the host that writes the bucket pings it)")
		return
	case b.fail && b.quiet:
		fmt.Fprintln(stdout, "felis watchdog: quiet while the installer runs; withholding the failure ping")
		return
	}
	sent, err := b.ping(cl)
	switch {
	case err != nil:
		fmt.Fprintf(stderr, "felis watchdog: heartbeat: %v\n", err)
	case !sent:
		fmt.Fprintf(stdout, "felis watchdog: withholding the heartbeat ping (the URL has a query, so it has no /fail endpoint)\n")
	case b.fail:
		fmt.Fprintf(stdout, "felis watchdog: pinged the heartbeat's failure endpoint at %s\n", redactURL(b.url))
	}
}

// ping sends the request; sent is false when a failure withholds it.
func (b heartbeat) ping(cl *http.Client) (sent bool, err error) {
	ctx, cancel := context.WithTimeout(context.Background(), heartbeatTimeout)
	defer cancel()
	method, target, body := http.MethodGet, b.url, io.Reader(nil)
	if b.fail {
		if strings.Contains(b.url, "?") {
			return false, nil
		}
		method, target = http.MethodPost, strings.TrimSuffix(b.url, "/")+"/fail"
		body = strings.NewReader(clipUTF8(b.report, heartbeatReportMax))
	}
	req, err := http.NewRequestWithContext(ctx, method, target, body)
	if err != nil {
		return false, fmt.Errorf("%s %s: bad URL", method, redactURL(target))
	}
	resp, err := cl.Do(req)
	if err != nil {
		// A url.Error quotes the whole URL, the check's key among it.
		var ue *url.Error
		if errors.As(err, &ue) {
			err = ue.Err
		}
		return false, fmt.Errorf("%s %s: %w", method, redactURL(target), err)
	}
	io.Copy(io.Discard, io.LimitReader(resp.Body, 4096))
	resp.Body.Close()
	if resp.StatusCode/100 != 2 {
		return false, fmt.Errorf("%s %s: %s", method, redactURL(target), resp.Status)
	}
	return true, nil
}

// clipUTF8 cuts s to at most n bytes without splitting a character.
func clipUTF8(s string, n int) string {
	if len(s) <= n {
		return s
	}
	return strings.ToValidUTF8(s[:n], "")
}

// readHeartbeatURL reads the heartbeat URL from path; no file is no heartbeat.
func readHeartbeatURL(path string) (string, error) {
	if path == "" {
		return "", nil
	}
	raw, err := os.ReadFile(path)
	if errors.Is(err, os.ErrNotExist) {
		return "", nil
	}
	if err != nil {
		return "", err
	}
	s := strings.TrimSpace(string(raw))
	if err := checkHeartbeatURL(s); err != nil {
		return "", fmt.Errorf("%s: %w", path, err)
	}
	return s, nil
}

// checkHeartbeatURL accepts an http(s) URL with a host. Its error leaves the
// URL out: the path is the check's key, which anyone who reads it can ping in
// the host's name.
func checkHeartbeatURL(s string) error {
	u, err := url.Parse(s)
	if err != nil || (u.Scheme != "http" && u.Scheme != "https") || u.Host == "" || strings.ContainsAny(s, " \t\r\n") {
		return errors.New("the heartbeat URL is not an http:// or https:// URL")
	}
	return nil
}

// redactURL is a heartbeat URL as logs show it: its scheme and host.
func redactURL(s string) string {
	u, err := url.Parse(s)
	if err != nil || u.Host == "" {
		return "(the heartbeat URL)"
	}
	return u.Scheme + "://" + u.Host + "/..."
}

// standsBy reports whether this host stands by for another host's off-site
// bucket (offsite.Status.Standby), however old that record is: the
// installer's take-over or the host's first write ends it.
func standsBy(offsiteOn bool, statusPath string) bool {
	if !offsiteOn {
		return false
	}
	st, err := offsite.ReadStatus(statusPath)
	return err == nil && st != nil && st.Standby
}
+589 −0

File added.

Preview size limit exceeded, changes collapsed.

+96 −6

File changed.

Preview size limit exceeded, changes collapsed.

Loading