// edgeguard-scheduler runs background jobs that don't belong on the // API request path: // // - ACME cert renewal (every 6h, re-issues anything < 30d to expiry) // // Future jobs (cluster heartbeat, backup, audit-log retention) // hang off the same Tick loop. Stays single-process — no leader // election yet (Phase 3). package main import ( "bufio" "context" "encoding/json" "fmt" "log/slog" "net" "os" "os/exec" "regexp" "strconv" "strings" "syscall" "time" "github.com/jackc/pgx/v5/pgxpool" "git.netcell-it.de/projekte/edgeguard-native/internal/cluster" "git.netcell-it.de/projekte/edgeguard-native/internal/cluster/clustertls" "git.netcell-it.de/projekte/edgeguard-native/internal/database" "git.netcell-it.de/projekte/edgeguard-native/internal/license" "git.netcell-it.de/projekte/edgeguard-native/internal/services/acme" "git.netcell-it.de/projekte/edgeguard-native/internal/services/alerts" "git.netcell-it.de/projekte/edgeguard-native/internal/services/audit" "git.netcell-it.de/projekte/edgeguard-native/internal/services/backends" "git.netcell-it.de/projekte/edgeguard-native/internal/services/backup" backupremote "git.netcell-it.de/projekte/edgeguard-native/internal/services/backup/remote" "git.netcell-it.de/projekte/edgeguard-native/internal/services/certrenewer" licsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/license" "git.netcell-it.de/projekte/edgeguard-native/internal/services/setup" "git.netcell-it.de/projekte/edgeguard-native/internal/services/tlscerts" ) var version = "1.2.35" const ( // renewTickInterval — how often we re-evaluate expiring certs. // 6h is enough: LE renewal window is 30 days; missing one tick // makes no difference. Hourly would log too much. renewTickInterval = 6 * time.Hour // certDir matches handlers.NewTLSCertsHandler default — HAProxy // reads from this directory. certDir = "/etc/edgeguard/tls" // licenseTickInterval — daily re-verify against // license.netcell-it.com. Result lands in the licenses table. licenseTickInterval = 24 * time.Hour // backupTickInterval — daily scheduled backup at ~03:00 (Tick // alignment ist approximativ, weil time.Ticker bei Boot startet). // Retention: 14 erfolgreiche Backups (default in backup.Service). backupTickInterval = 24 * time.Hour // staleSweepTickInterval — Phase 3.2: alle 30s prüfen ob Peers // last_seen länger als staleThreshold nicht gemeldet haben → // status='offline'. Symmetrisch zum 30s-API-Heartbeat. staleSweepTickInterval = 30 * time.Second // staleThreshold — Peer gilt als offline wenn last_seen älter als // das ist. 4× Heartbeat-Intervall lässt einen verpassten Tick // (Restart, GC-Pause, kurzer Network-Glitch) zu ohne false positive. staleThreshold = 2 * time.Minute // clusterCertCheckInterval — täglicher Check ob CA + peer-cert // in den nächsten clusterCertWarnDays ablaufen. Bei Hit feuert // ein Alert (dedupe 12h damit der Operator nicht alle 24h dieselbe // Warnung sieht). clusterCertCheckInterval = 24 * time.Hour // clusterCertWarnDays — Schwelle für die Cert-Expiry-Warnung. // Operator hat damit min. 30 Tage Vorlauf für `edgeguard-ctl // cluster-renew-self` oder einen manuellen Re-Join. clusterCertWarnDays = 30 // clusterCertAutoRenewDays — Schwelle ab der wir automatisch // neu signieren (nur Founder mit lokaler CA). Wir liegen 2× vor // der Warn-Schwelle damit ein verpasster Tick + ein verpasster // Restart-Window noch passen. clusterCertAutoRenewDays = 60 // diskCheckInterval — stündliche Disk-Usage-Prüfung. Fire-Schwellen // in runDiskCheck (warning 80%, error 90%). Stündlich ist schnell // genug damit der Operator vor /var = 100% noch Zeit zum Aufräumen // hat, ohne Log-Spam zu produzieren (Dedupe 12h pro Severity). diskCheckInterval = 1 * time.Hour diskWarnPct = 80.0 diskCriticalPct = 90.0 // auditCleanupInterval — täglicher Cleanup. Audit-Rows älter als // auditRetentionDays werden gelöscht. Idempotent — wenn nichts da // ist passiert nichts. auditCleanupInterval = 24 * time.Hour auditRetentionDays = 90 // alertRetentionDays — alert_events wächst sonst unbegrenzt (node-lokale // Health-Events: backend.down, mem.high, cert.expiring …). Läuft im // selben täglichen Tick wie der Audit-Cleanup. Fester Default, kein // Setup-Override (Events sind reine Diagnose-History). alertRetentionDays = 90 // backendDownCheckInterval — alle 2 Minuten HAProxy-Stats lesen und // prüfen ob ein Backend komplett ausgefallen ist (alle Server DOWN). // Dedupe 12h pro Backend → kein Alert-Spam. Frischer Alert wenn das // Backend nach 12h immer noch unten ist. backendDownCheckInterval = 2 * time.Minute // memCheckInterval — alle 5 Minuten /proc/meminfo lesen. Schwellen // warning 85%, critical 95%. Dedupe 1h pro Severity damit bei einem // kurzfristigen Spike nicht jeder Tick feuert. memCheckInterval = 5 * time.Minute memWarnPct = 85.0 memCriticalPct = 95.0 // conntrackCheckInterval — alle 2 Minuten /proc/sys/net/netfilter/ // nf_conntrack_count+max lesen. Eine volle conntrack-Tabelle verwirft // alle neuen Verbindungen ohne jegliche Rückmeldung. 2-Minuten-Takt // erlaubt früh zu warnen bevor die Tabelle überläuft. // Schwellen analog Disk: 80% Warning, 90% Critical. Dedupe 1h. conntrackCheckInterval = 2 * time.Minute conntrackWarnPct = 80.0 conntrackCriticalPct = 90.0 // ntpSyncCheckInterval — alle 10 Minuten chronyc tracking aufrufen. // Keine Sync bedeutet: Uhr driftet → TLS-Cert-Prüfung schlägt fehl // wenn die Abweichung > Toleranz des Gegenstücks (i.d.R. ±1 min), // JWT-Ablauf inconsistent, Cluster-Split-Brain möglich. Dedupe 1h // damit ein kurzer Upstream-Ausfall (Reboot, DHCP-Pause) keinen // Alert-Regen produziert. ntpSyncCheckInterval = 10 * time.Minute // wgTunnelCheckInterval — alle 5 Minuten WireGuard-Client-Tunnels // auf Aktualität prüfen. Client-Tunnels (mode='client') haben genau // einen Peer; wenn dessen letzter Handshake älter als wgStaleSec ist, // ist der Tunnel effektiv tot — Traffic droht lautlos. Dedupe 30min // pro Tunnel damit schnell wiederhergestellte Tunnels nur einmal feuern. wgTunnelCheckInterval = 5 * time.Minute wgStaleSec = int64(5 * 60) // 5 Minuten ohne Handshake = tot ) func main() { slog.SetDefault(slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelInfo}))) slog.Info("edgeguard-scheduler starting", "version", version) ctx := context.Background() pool, err := database.Open(ctx, database.ConnStringFromEnv()) if err != nil { slog.Error("scheduler: DB open failed — sleeping forever", "error", err) select {} } defer pool.Close() tlsRepo := tlscerts.New(pool) setupStore := setup.NewStore(setup.DefaultDir) st, _ := setupStore.Load() var renewer *certrenewer.Service if st != nil && st.ACMEEmail != "" { issuer := acme.New(st.ACMEEmail) renewer = certrenewer.New(tlsRepo, issuer, certDir, 30*24*time.Hour) slog.Info("scheduler: ACME renewer enabled", "email", st.ACMEEmail, "tick", renewTickInterval, "threshold", "30d") } else { slog.Warn("scheduler: setup.acme_email empty — ACME renewal disabled until setup wizard ran") } licRepo := licsvc.New(pool) licClient := license.NewClient() licKeyStore := license.NewKeyStore() nodeID := os.Getenv("EDGEGUARD_NODE_ID") slog.Info("scheduler: license re-verify enabled", "tick", licenseTickInterval) backupSvc := backup.New(pool) backupSvc.RemoteUploader = newSchedRemoteAdapter(backupremote.New(pool)) slog.Info("scheduler: daily backup enabled", "tick", backupTickInterval, "dir", backupSvc.BackupDir, "keep_n", backup.DefaultKeepN) alertSvc := alerts.New(pool) auditRepo := audit.New(pool) alertDedupe := newDedupe(12 * time.Hour) if renewer != nil { runRenewer(ctx, renewer, alertSvc, alertDedupe) } runLicenseVerify(ctx, licClient, licKeyStore, licRepo, nodeID, alertSvc, alertDedupe) // Lokale Node-ID für Heartbeat. EnsureNodeID liefert dieselbe ID // die die API hat (gleiches /var/lib/edgeguard/node-id). localID, _ := cluster.EnsureNodeID("") slog.Info("scheduler: stale-sweeper enabled", "tick", staleSweepTickInterval, "threshold", staleThreshold, "node_id", localID) // Initial-Sweep + initial-Heartbeat damit /cluster/status nach // Scheduler-Boot direkt einen frischen Zustand sieht. runHeartbeat(ctx, pool, localID, version) runStaleSweep(ctx, pool) renewTick := time.NewTicker(renewTickInterval) defer renewTick.Stop() licTick := time.NewTicker(licenseTickInterval) defer licTick.Stop() backupTick := time.NewTicker(backupTickInterval) defer backupTick.Stop() sweepTick := time.NewTicker(staleSweepTickInterval) defer sweepTick.Stop() clusterCertTick := time.NewTicker(clusterCertCheckInterval) defer clusterCertTick.Stop() // Initial-Cert-Check direkt beim Start, sonst sieht der Operator // einen Warning erst nach 24h ab Boot. runClusterCertExpiryCheck(ctx, alertSvc, alertDedupe) diskTick := time.NewTicker(diskCheckInterval) defer diskTick.Stop() // Initial-Disk-Check: wenn die Box schon bei 95% steht beim // Scheduler-Boot, wollen wir keine Stunde auf den ersten Alert // warten. runDiskCheck(ctx, alertSvc, alertDedupe) auditTick := time.NewTicker(auditCleanupInterval) defer auditTick.Stop() backendDownTick := time.NewTicker(backendDownCheckInterval) defer backendDownTick.Stop() // Initial-Check direkt beim Start — wenn ein Backend seit dem letzten // Scheduler-Restart down ist, brauchen wir nicht 2 Minuten zu warten. runBackendDownCheck(ctx, pool, alertSvc, alertDedupe) memTick := time.NewTicker(memCheckInterval) defer memTick.Stop() runMemoryCheck(ctx, alertSvc, alertDedupe) conntrackTick := time.NewTicker(conntrackCheckInterval) defer conntrackTick.Stop() runConntrackCheck(ctx, alertSvc, alertDedupe) ntpSyncTick := time.NewTicker(ntpSyncCheckInterval) defer ntpSyncTick.Stop() // Kein Initial-Check bei Boot: chrony braucht nach dem Start // einige Sekunden bis zur ersten Synchronization — ein // sofortiger Check würde immer feuern. wgTunnelTick := time.NewTicker(wgTunnelCheckInterval) defer wgTunnelTick.Stop() // Kein Initial-Check bei Boot: Tunnels brauchen nach dem Start // des wg-quick-Dienstes einen Moment für den ersten Handshake. for { select { case <-renewTick.C: if renewer != nil { runRenewer(ctx, renewer, alertSvc, alertDedupe) } runCertExpiryCheck(ctx, tlsRepo, alertSvc, alertDedupe) case <-licTick.C: runLicenseVerify(ctx, licClient, licKeyStore, licRepo, nodeID, alertSvc, alertDedupe) case <-backupTick.C: runBackup(ctx, backupSvc, version, alertSvc, setupStore) case <-sweepTick.C: // Symmetrisches Heartbeat aus dem Scheduler — falls die API // pausiert/hängt, hält der Scheduler die eigene Row warm. // Idempotent zur API-Heartbeat-Goroutine. runHeartbeat(ctx, pool, localID, version) runStaleSweep(ctx, pool) case <-clusterCertTick.C: runClusterCertExpiryCheck(ctx, alertSvc, alertDedupe) case <-diskTick.C: runDiskCheck(ctx, alertSvc, alertDedupe) case <-auditTick.C: runAuditCleanup(ctx, auditRepo, setupStore) runAlertCleanup(ctx, alertSvc) case <-backendDownTick.C: runBackendDownCheck(ctx, pool, alertSvc, alertDedupe) case <-memTick.C: runMemoryCheck(ctx, alertSvc, alertDedupe) case <-conntrackTick.C: runConntrackCheck(ctx, alertSvc, alertDedupe) case <-ntpSyncTick.C: runNTPSyncCheck(ctx, alertSvc, alertDedupe) case <-wgTunnelTick.C: runWGClientTunnelCheck(ctx, pool, alertSvc, alertDedupe) } } } // runAuditCleanup löscht audit_log-Rows älter als die konfigurierte // Retention. Operator kann den Wert in den Settings übersteuern; ohne // Setup-Custom fällt's auf auditRetentionDays-Default zurück. // Schutz vor unbounded growth — auf einer aktiven Box wird das Log // nach 1-2 Jahren mehrere GB groß und macht die /logs-Page langsam. // Best-effort: Fehler werden nur geloggt, der Tick läuft beim nächsten // Zyklus wieder. func runAuditCleanup(ctx context.Context, r *audit.Repo, setupStore *setup.Store) { if r == nil { return } keepDays := auditRetentionDays if setupStore != nil { if st, err := setupStore.Load(); err == nil && st != nil && st.AuditRetentionDays > 0 { keepDays = st.AuditRetentionDays } } cctx, cancel := context.WithTimeout(ctx, 30*time.Second) defer cancel() n, err := r.Cleanup(cctx, keepDays) if err != nil { slog.Warn("scheduler: audit cleanup failed", "keep_days", keepDays, "error", err) return } if n > 0 { slog.Info("scheduler: audit cleanup", "deleted", n, "keep_days", keepDays) } } // runAlertCleanup löscht alert_events älter als alertRetentionDays. // Schutz vor unbounded growth — auf einer aktiven Box feuern backend.down/ // mem.high/cert.expiring über Monate tausende Rows (die Tabelle ist // node-lokal, wird also nirgends sonst abgeräumt). Best-effort: Fehler // werden nur geloggt. func runAlertCleanup(ctx context.Context, a *alerts.Service) { if a == nil { return } cctx, cancel := context.WithTimeout(ctx, 30*time.Second) defer cancel() n, err := a.Cleanup(cctx, alertRetentionDays) if err != nil { slog.Warn("scheduler: alert cleanup failed", "keep_days", alertRetentionDays, "error", err) return } if n > 0 { slog.Info("scheduler: alert cleanup", "deleted", n, "keep_days", alertRetentionDays) } } // runDiskCheck prüft die Belegung von / via statfs. Fire-Schwellen: // - >= 90% → Critical (error). Box ist akut gefährdet — beim // nächsten Backup-Run oder größeren apt-Update droht "no space // left" und damit failed services. // - >= 80% → Warning. Operator hat noch Luft aber sollte aufräumen. // - < 80% → kein Alert. // // Dedupe-Keys pro Severity, damit ein lang-belegtes Filesystem nicht // jede Stunde feuert (12h pro Stufe). Bei Übergang warning→critical // gibt's einen frischen Alert weil die Keys verschieden sind. // // Fix-Hint im Body: was der Operator als Erstes prüfen soll // (/var/backups/edgeguard, /var/log/edgeguard, /var/cache/apt). func runDiskCheck(ctx context.Context, a *alerts.Service, d *dedupe) { if a == nil || d == nil { return } var fs syscall.Statfs_t if err := syscall.Statfs("/", &fs); err != nil { slog.Warn("scheduler: disk-check statfs failed", "error", err) return } total := float64(fs.Blocks) * float64(fs.Bsize) free := float64(fs.Bavail) * float64(fs.Bsize) if total <= 0 { return } usedPct := (total - free) * 100 / total freeGB := free / (1024 * 1024 * 1024) totalGB := total / (1024 * 1024 * 1024) var key, title string var sev alerts.Severity switch { case usedPct >= diskCriticalPct: key = "disk.full.critical" sev = alerts.SeverityError title = fmt.Sprintf("Disk kritisch voll: %.0f%%", usedPct) case usedPct >= diskWarnPct: key = "disk.full.warning" sev = alerts.SeverityWarning title = fmt.Sprintf("Disk-Belegung hoch: %.0f%%", usedPct) default: return } if !d.shouldFire(key) { return } desc := fmt.Sprintf( "Wurzel-Filesystem (/) ist zu %.1f%% belegt — noch %.2f GB von %.2f GB frei.\n\n"+ "Häufige Verursacher checken:\n"+ " sudo du -hs /var/backups/edgeguard /var/log/edgeguard /var/cache/apt /var/lib/postgresql\n\n"+ "Backup-Retention ist 14 (default). Manuell aufräumen:\n"+ " ls -lhS /var/backups/edgeguard | head\n"+ " sudo apt-get clean # /var/cache/apt leeren", usedPct, freeGB, totalGB) if _, err := a.Fire(ctx, "disk.full", sev, title, desc); err != nil { slog.Warn("scheduler: disk-check alert fire failed", "error", err) } } // runClusterCertExpiryCheck warnt wenn CA oder peer.crt in < // clusterCertWarnDays Tagen ablaufen (oder schon abgelaufen sind). // Dedupe pro Cert-Typ + 12h. // // Zusätzlich (Phase 1.1.1): wenn das peer.crt < clusterCertAutoRenewDays // remaining hat UND eine lokale CA existiert, wird automatisch neu // signiert. Restart-Hinweis als Info-Alert — wir starten edgeguard-api // nicht selbst neu, das passiert beim nächsten geplanten Update/Reboot. // runMemoryCheck liest /proc/meminfo und feuert bei hoher RAM-Belegung. // Schwellen: warning >= 85%, critical >= 95%. Dedupe 1h pro Severity // damit kurze Spikes (Backup, apt-Upgrade) keine Alert-Flut erzeugen. func runMemoryCheck(ctx context.Context, a *alerts.Service, d *dedupe) { if a == nil || d == nil { return } data, err := os.ReadFile("/proc/meminfo") if err != nil { return } var memTotal, memAvail int64 for _, line := range strings.Split(string(data), "\n") { var key string var val int64 if _, err := fmt.Sscanf(line, "%s %d", &key, &val); err != nil { continue } switch key { case "MemTotal:": memTotal = val case "MemAvailable:": memAvail = val } } if memTotal <= 0 { return } usedPct := float64(memTotal-memAvail) * 100 / float64(memTotal) usedGB := float64(memTotal-memAvail) / 1024 / 1024 totalGB := float64(memTotal) / 1024 / 1024 var key, title string var sev alerts.Severity switch { case usedPct >= memCriticalPct: key = "mem.high.critical" sev = alerts.SeverityError title = fmt.Sprintf("RAM kritisch hoch: %.0f%%", usedPct) case usedPct >= memWarnPct: key = "mem.high.warning" sev = alerts.SeverityWarning title = fmt.Sprintf("RAM-Belegung hoch: %.0f%%", usedPct) default: return } if !d.shouldFire(key) { return } desc := fmt.Sprintf( "RAM-Auslastung: %.1f%% — %.1f von %.1f GB belegt.\n\n"+ "Häufige Ursachen:\n"+ " • Unbound-Cache zu groß (rrset-cache-size in /etc/edgeguard/unbound/unbound.conf)\n"+ " • Squid cache_mem zu groß (64 MB default)\n"+ " • PostgreSQL shared_buffers (default ~128 MB)\n"+ " • Prozesse prüfen: ps aux --sort=-%%mem | head -10", usedPct, usedGB, totalGB) if _, err := a.Fire(ctx, "mem.high", sev, title, desc); err != nil { slog.Warn("scheduler: memory-check alert fire failed", "error", err) } } // runConntrackCheck liest die conntrack-Tabellen-Belegung aus /proc und // feuert bei hoher Auslastung. Eine volle conntrack-Tabelle (100%) // verwirft alle neuen TCP/UDP-Verbindungen ohne ICMP-Rückmeldung — // der Operator sieht auf der Gegenstelle nur Timeouts. // // Schwellen: 80% Warning, 90% Critical (wie Disk, niedriger als RAM weil // der Impact sofortig ist). Dedupe 1h pro Severity. func runConntrackCheck(ctx context.Context, a *alerts.Service, d *dedupe) { if a == nil || d == nil { return } readInt := func(path string) int64 { b, err := os.ReadFile(path) if err != nil { return 0 } v, _ := strconv.ParseInt(strings.TrimSpace(string(b)), 10, 64) return v } count := readInt("/proc/sys/net/netfilter/nf_conntrack_count") max := readInt("/proc/sys/net/netfilter/nf_conntrack_max") if max <= 0 { return } usedPct := float64(count) * 100 / float64(max) var key, title string var sev alerts.Severity switch { case usedPct >= conntrackCriticalPct: key = "conntrack.high.critical" sev = alerts.SeverityError title = fmt.Sprintf("Conntrack-Tabelle kritisch voll: %.0f%%", usedPct) case usedPct >= conntrackWarnPct: key = "conntrack.high.warning" sev = alerts.SeverityWarning title = fmt.Sprintf("Conntrack-Tabelle fast voll: %.0f%%", usedPct) default: return } if !d.shouldFire(key) { return } desc := fmt.Sprintf( "Conntrack-Auslastung: %.1f%% — %d von %d Einträgen belegt.\n\n"+ "Wenn die Tabelle auf 100%% steigt, werden alle neuen Verbindungen\n"+ "ohne Fehlermeldung verworfen (Silent Drop).\n\n"+ "Maßnahmen:\n"+ " • Zeitweilige Spikes: nf_conntrack_max erhöhen\n"+ " (sysctl net.netfilter.nf_conntrack_max)\n"+ " • Leaks: conntrack -L | sort | head zeigt häufige Quellen\n"+ " • Timeouts reduzieren (z.B. nf_conntrack_tcp_timeout_established)", usedPct, count, max) if _, err := a.Fire(ctx, "conntrack.high", sev, title, desc); err != nil { slog.Warn("scheduler: conntrack-check alert fire failed", "error", err) } } // runNTPSyncCheck ruft chronyc tracking auf und feuert einen Alert wenn // chrony keine synchronisierte Zeitquelle hat (Stratum 0 oder ≥ 16). // Zeitdrift > ~1 Minute führt zu TLS-Handshake-Fehlern, JWT-Ablauf- // Inkonsistenzen und möglichen Cluster-Problemen. Dedupe 1h. func runNTPSyncCheck(ctx context.Context, a *alerts.Service, d *dedupe) { if a == nil || d == nil { return } out, err := exec.Command("chronyc", "tracking").Output() if err != nil { // chrony nicht installiert oder nicht gestartet — kein Alert, // weil wir nicht wissen ob chrony hier überhaupt erwartet wird. return } synced, stratum, ref := parseChronyTrackingForAlert(string(out)) if synced { return } const key = "ntp.unsync" if !d.shouldFire(key) { return } refStr := ref if refStr == "" { refStr = "(keine Referenz)" } title := fmt.Sprintf("NTP nicht synchronisiert (Stratum %d)", stratum) desc := fmt.Sprintf( "chrony hat keine synchronisierte Zeitquelle.\n"+ "Referenz: %s Stratum: %d\n\n"+ "Mögliche Ursachen:\n"+ " • Upstream-NTP-Server nicht erreichbar (UDP/123 blockiert?)\n"+ " • Pool-DNS-Einträge lösen nicht auf\n"+ " • chrony läuft, braucht aber noch Zeit nach Boot (warten)\n\n"+ "Prüfen: chronyc sources -v — chronyc tracking", refStr, stratum) if _, err := a.Fire(ctx, "ntp.unsync", alerts.SeverityWarning, title, desc); err != nil { slog.Warn("scheduler: ntp-sync-check alert fire failed", "error", err) } } // parseChronyTrackingForAlert ist eine schlanke Variante des NTP-Handler- // Parsers: liefert nur synced/stratum/reference ohne die vollen Felder. func parseChronyTrackingForAlert(out string) (synced bool, stratum int, reference string) { for _, line := range strings.Split(out, "\n") { line = strings.TrimSpace(line) key, val, ok := strings.Cut(line, ":") if !ok { continue } key = strings.TrimSpace(key) val = strings.TrimSpace(val) switch key { case "Reference ID": if i := strings.Index(val, "("); i >= 0 { reference = strings.Trim(val[i:], "()") } if val != "00000000 ()" { synced = true } case "Stratum": _, _ = fmt.Sscanf(val, "%d", &stratum) if stratum > 0 && stratum < 16 { synced = true } else if stratum == 0 || stratum >= 16 { synced = false } } } return } // runWGClientTunnelCheck prüft alle aktiven WireGuard-Client-Tunnels // (mode='client') auf Handshake-Aktualität. Ein Client-Tunnel hat genau // einen Peer; wenn dessen letzter Handshake älter als wgStaleSec oder // noch nie stattgefunden hat, ist der Tunnel tot — Traffic wird lautlos // verworfen (kein ICMP Unreachable). Dedupe 30min pro Tunnel damit // nach einer Selbstheilung nicht alle paar Minuten neu gefeuert wird. func runWGClientTunnelCheck(ctx context.Context, pool *pgxpool.Pool, a *alerts.Service, d *dedupe) { if a == nil || d == nil || pool == nil { return } // Alle aktiven Client-Interfaces aus DB laden. type wgIface struct{ name string } rows, err := pool.Query(ctx, `SELECT name FROM wg_interfaces WHERE mode = 'client' AND active = true ORDER BY name`) if err != nil { return } defer rows.Close() var ifaces []wgIface for rows.Next() { var n string if err := rows.Scan(&n); err == nil { ifaces = append(ifaces, wgIface{n}) } } rows.Close() if len(ifaces) == 0 { return } now := time.Now().Unix() for _, ifc := range ifaces { out, err := exec.Command("wg", "show", ifc.name, "dump").Output() if err != nil { // Interface existiert nicht mehr im Kernel (wg-quick down) — // das ist selbst schon ein Problem; kein separater Alert hier, // da systemd-Restart-Policy das abdeckt. continue } lines := strings.Split(strings.TrimSpace(string(out)), "\n") // Zeile 0 ist die Interface-Zeile (own key / pubkey / port / fwmark). // Zeile 1 ist die Peer-Zeile: pubkey psk endpoint allowed-ips last-hs rx tx keepalive if len(lines) < 2 { continue } fields := strings.Fields(lines[1]) if len(fields) < 5 { continue } lastHS, _ := strconv.ParseInt(fields[4], 10, 64) stale := lastHS == 0 || (now-lastHS) > wgStaleSec if !stale { continue } key := "wg.tunnel.down." + ifc.name if !d.shouldFire(key) { continue } var detail string if lastHS == 0 { detail = "Noch kein Handshake — Tunnel wurde nie erfolgreich aufgebaut." } else { ageMin := (now - lastHS) / 60 detail = fmt.Sprintf("Letzter Handshake: vor %d Minuten.", ageMin) } title := fmt.Sprintf("WireGuard-Tunnel %s ausgefallen", ifc.name) desc := fmt.Sprintf( "Client-Tunnel %s hat seit >5 Minuten keinen Handshake.\n%s\n\n"+ "Traffic zu den RemoteAllowed-Netzen wird lautlos verworfen.\n\n"+ "Mögliche Ursachen:\n"+ " • Remote-Peer nicht erreichbar (Firewall, Routing)\n"+ " • Remote-Server-Keypair geändert (Public-Key stimmt nicht mehr)\n"+ " • UDP-Port des Peers geblockt\n"+ " • wg-quick-Dienst auf dieser Box gestoppt: systemctl status wg-quick@%s", ifc.name, detail, ifc.name) if _, err := a.Fire(ctx, "wg.tunnel.down", alerts.SeverityError, title, desc); err != nil { slog.Warn("scheduler: wg-tunnel-check alert fire failed", "iface", ifc.name, "error", err) } } } var egBackendRE = regexp.MustCompile(`^eg_backend_(\d+)$`) // nodeHoldsVIP meldet true, wenn dieser Node aktuell mindestens eine // is_vip-Adresse lokal trägt — also der keepalived-MASTER ist. Nur der // Master hält die VLAN-Gateway-VIPs und erreicht damit die Backend- // Subnetze; ein BACKUP-Node hat KEINE VLAN-IP und sieht deshalb JEDES // Backend als L4-down. Spiegelt SystemHandler.VIPStatus (net.Interfaces, // kein Shell-out). func nodeHoldsVIP(ctx context.Context, pool *pgxpool.Pool) bool { if pool == nil { return false } rows, err := pool.Query(ctx, `SELECT address FROM ip_addresses WHERE is_vip = true AND active = true`) if err != nil { return false } defer rows.Close() var vips []string for rows.Next() { var addr string if err := rows.Scan(&addr); err == nil { vips = append(vips, addr) } } if len(vips) == 0 { return false } local := make(map[string]bool) ifaces, err := net.Interfaces() if err != nil { return false } for _, ifc := range ifaces { addrs, err := ifc.Addrs() if err != nil { continue } for _, a := range addrs { if ipnet, ok := a.(*net.IPNet); ok { local[ipnet.IP.String()] = true } } } for _, v := range vips { if local[v] { return true } } return false } // runBackendDownCheck liest HAProxy-Stats via Admin-Socket und feuert // einen Error-Alert für jedes Backend bei dem alle Server DOWN sind // (und mind. einer einen echten Health-Check hat). Dedupe 12h pro Backend. func runBackendDownCheck(ctx context.Context, pool *pgxpool.Pool, a *alerts.Service, d *dedupe) { if a == nil || d == nil { return } // Nur auf dem VIP-Master prüfen. Ein BACKUP-Node hält die VLAN- // Gateway-VIPs nicht und kann die Backend-Subnetze gar nicht erreichen // → jeder Health-Check läuft L4TOUT → Dauer-"backend.down"-Fehlalarm // (Hauptquelle des alert_events-Spams). Der Master bedient den Traffic // und sieht die echten Backend-States. if !nodeHoldsVIP(ctx, pool) { return } dialer := net.Dialer{Timeout: 2 * time.Second} conn, err := dialer.DialContext(ctx, "unix", "/run/haproxy/admin.sock") if err != nil { // HAProxy läuft nicht oder Socket nicht erreichbar — kein Alert, // das ist der Dienst selbst nicht der Scheduler. return } defer func() { _ = conn.Close() }() _ = conn.SetDeadline(time.Now().Add(3 * time.Second)) if _, err := conn.Write([]byte("show stat\n")); err != nil { return } type srvEntry struct { status string hasCheck bool } byBackend := map[string][]srvEntry{} colIdx := map[string]int{} scanner := bufio.NewScanner(conn) scanner.Buffer(make([]byte, 64*1024), 1024*1024) for scanner.Scan() { line := scanner.Text() if line == "" { continue } fields := strings.Split(line, ",") if strings.HasPrefix(line, "# ") { fields[0] = strings.TrimPrefix(fields[0], "# ") for i, name := range fields { colIdx[name] = i } continue } px := fieldAt(fields, colIdx["pxname"]) sv := fieldAt(fields, colIdx["svname"]) if !strings.HasPrefix(px, "eg_backend_") || sv == "BACKEND" || sv == "FRONTEND" || sv == "" { continue } status := fieldAt(fields, colIdx["status"]) byBackend[px] = append(byBackend[px], srvEntry{ status: status, hasCheck: status != "no check", }) } if len(byBackend) == 0 { return } // Friendly Backend-Namen aus DB — best-effort, Fehler = anonyme ID. bkRepo := backends.New(pool) bklist, _ := bkRepo.List(ctx) nameOf := func(id int64) string { for _, b := range bklist { if b.ID == id { return b.Name } } return fmt.Sprintf("#%d", id) } for haName, servers := range byBackend { hasRealCheck, allDown := false, true for _, s := range servers { if s.hasCheck { hasRealCheck = true } if s.status == "UP" { allDown = false break } } if !hasRealCheck || !allDown { continue } m := egBackendRE.FindStringSubmatch(haName) if m == nil { continue } id, _ := strconv.ParseInt(m[1], 10, 64) name := nameOf(id) key := "backend.down." + haName if !d.shouldFire(key) { continue } msg := fmt.Sprintf( "Alle Server in Backend \"%s\" sind DOWN — HAProxy liefert 503 für alle Requests zu diesem Backend.\n\n"+ "HAProxy-Backend-Name: %s\n\n"+ "Nächste Schritte:\n"+ " • Dienst auf Backend-Host prüfen (systemctl status / docker ps)\n"+ " • Health-Check-Pfad erreichbar? (curl http://:)\n"+ " • Firewall-Regeln zwischen EdgeGuard und Backend-Host prüfen", name, haName) if _, err := a.Fire(ctx, "backend.down", alerts.SeverityError, fmt.Sprintf("Backend DOWN: %s", name), msg); err != nil { slog.Warn("scheduler: backend-down alert fire failed", "backend", name, "error", err) } } } func fieldAt(fields []string, i int) string { if i < 0 || i >= len(fields) { return "" } return fields[i] } func runClusterCertExpiryCheck(ctx context.Context, a *alerts.Service, d *dedupe) { if a == nil || d == nil { return } store := clustertls.New("") // Auto-Renew zuerst — danach lesen wir die (eventuell frischen) // Cert-Infos für den Alert-Check ab. tryAutoRenew(ctx, store, a, d) check := func(kind, key string, info *clustertls.CertInfo, err error) { if err != nil { return // Cert nicht vorhanden / unleserlich — keine Warnung. } if info.DaysRemaining > clusterCertWarnDays { return } if !d.shouldFire(key) { return } sev := alerts.SeverityWarning title := "Cluster-" + kind + " läuft bald ab" desc := fmt.Sprintf("%s (CN=%s) läuft in %d Tagen ab (NotAfter=%s).", kind, info.CommonName, info.DaysRemaining, info.NotAfter.Format(time.RFC3339)) if info.DaysRemaining < 0 { sev = alerts.SeverityError title = "Cluster-" + kind + " ist abgelaufen" desc = fmt.Sprintf("%s (CN=%s) ist seit %d Tagen abgelaufen (NotAfter=%s).", kind, info.CommonName, -info.DaysRemaining, info.NotAfter.Format(time.RFC3339)) } desc += "\n\nFix: sudo edgeguard-ctl cluster-renew-self (founder/single-node)\n sudo systemctl restart edgeguard-api" if _, err := a.Fire(ctx, "cluster.cert.expiring", sev, title, desc); err != nil { slog.Warn("scheduler: cluster cert alert fire failed", "kind", kind, "error", err) } } if store.HasCA() { info, err := store.CACertInfo() check("CA", "cluster.cert.expiring:ca", info, err) } if store.HasPeer() { info, err := store.PeerCertInfo() check("peer-Cert", "cluster.cert.expiring:peer", info, err) } } // tryAutoRenew: wenn das peer.crt < clusterCertAutoRenewDays Tage // remaining hat UND wir eine lokale CA haben (= Founder-Node), wird // automatisch ein frisches peer.{crt,key} signiert. Edgeguard-api // muss anschließend manuell restartet werden damit der Listener das // neue Material lädt — wir alarmieren das, restarten aber nicht // selbst (würde laufende Requests + die Heartbeat-Goroutine killen). // // Joiner-Nodes (keine eigene CA) ignorieren wir hier; sie laufen über // einen anderen Renewal-Pfad (Phase 3.6, Renewal-Token via mTLS). func tryAutoRenew(ctx context.Context, store *clustertls.Store, a *alerts.Service, d *dedupe) { if !store.HasPeer() || !store.HasCA() { return } info, err := store.PeerCertInfo() if err != nil { return } if info.DaysRemaining > clusterCertAutoRenewDays { return } // CN aus dem alten Cert übernehmen — sonst würde ein Hostname- // Wechsel mitten in der Renewal unbemerkt durchgehen. cn := info.CommonName if cn == "" { cn = "edgeguard-node" } if err := store.RenewSelfSigned(cn, []string{cn}, nil, nil); err != nil { slog.Warn("scheduler: cluster cert auto-renew failed", "error", err) // Failure-Alert dedupe 12h — Operator soll daran erinnert werden. if d.shouldFire("cluster.cert.auto_renew.failed") { _, _ = a.Fire(ctx, "cluster.cert.auto_renew.failed", alerts.SeverityError, "Cluster-Peer-Cert Auto-Renew fehlgeschlagen", "clustertls.RenewSelfSigned: "+err.Error()+ "\n\nFix: sudo edgeguard-ctl cluster-renew-self") } return } // Success — neuer Cert auf Disk, alter Cert noch im API-Speicher. // Info-Alert mit Restart-Hinweis. Dedupe 24h damit nicht // gefloodet wird wenn der Operator nicht restartet. if d.shouldFire("cluster.cert.auto_renew.ok") { fresh, _ := store.PeerCertInfo() until := info.NotAfter.Format(time.RFC3339) if fresh != nil { until = fresh.NotAfter.Format(time.RFC3339) } _, _ = a.Fire(ctx, "cluster.cert.auto_renewed", alerts.SeverityInfo, "Cluster-Peer-Cert automatisch erneuert", fmt.Sprintf("Neues Peer-Cert auf Disk (CN=%s, gültig bis %s). "+ "Damit edgeguard-api das neue Cert in den mTLS-Listener lädt:\n\n"+ " sudo systemctl restart edgeguard-api\n\n"+ "Bis dahin nutzt der laufende Prozess das alte Cert.", cn, until)) } slog.Info("scheduler: cluster peer cert auto-renewed", "cn", cn, "old_days_remaining", info.DaysRemaining) } // dedupe verhindert dass derselbe Alert-Key (z.B. "cert.expiring:utm-1.netcell-it.de") // öfter als alle 12h gefeuert wird. In-memory — Scheduler-Restart // resettet, was OK ist (Operator soll bei restart wieder einen kennen- // lernen-Event sehen können). type dedupe struct { ttl time.Duration last map[string]time.Time } func newDedupe(ttl time.Duration) *dedupe { return &dedupe{ttl: ttl, last: map[string]time.Time{}} } func (d *dedupe) shouldFire(key string) bool { now := time.Now() if last, ok := d.last[key]; ok && now.Sub(last) < d.ttl { return false } d.last[key] = now return true } // runCertExpiryCheck prüft tls_certs auf bevorstehende Expiry. Warning // bei <14d Restzeit. Dedupe pro Cert-Name 12h damit der scheduler // nicht alle 6h dieselbe Warnung feuert. func runCertExpiryCheck(ctx context.Context, repo *tlscerts.Repo, a *alerts.Service, d *dedupe) { if repo == nil || a == nil { return } certs, err := repo.List(ctx) if err != nil { slog.Warn("scheduler: cert-expiry list failed", "error", err) return } threshold := 14 * 24 * time.Hour now := time.Now() for _, c := range certs { if c.NotAfter == nil { continue } remain := c.NotAfter.Sub(now) if remain > threshold || remain < -90*24*time.Hour { continue } key := "cert.expiring:" + c.Domain if !d.shouldFire(key) { continue } days := int(remain.Hours() / 24) sev := alerts.SeverityWarning if days < 3 { sev = alerts.SeverityError } _, err := a.Fire(ctx, "cert.expiring", sev, "TLS-Zertifikat läuft ab: "+c.Domain, "Cert für "+c.Domain+" läuft in "+strconv.Itoa(days)+" Tagen ab ("+c.NotAfter.Format(time.RFC3339)+"). Renewer-Status: "+c.Status) if err != nil { slog.Warn("scheduler: alert fire failed", "error", err) } } } // runHeartbeat schreibt last_seen + status=online + version + config_hash // auf die eigene ha_nodes-Row. Pool kann nil sein (scheduler-pool-fail // beim Boot) — dann no-op. Errors landen im WARN, kein Abort der Schleife. func runHeartbeat(ctx context.Context, pool *pgxpoolPool, localID, version string) { if pool == nil || localID == "" { return } hbCtx, cancel := context.WithTimeout(ctx, 5*time.Second) defer cancel() if err := cluster.Heartbeat(hbCtx, pool, localID, version); err != nil { slog.Warn("scheduler: heartbeat failed", "error", err) } } // runStaleSweep markiert Peers mit last_seen < NOW()-staleThreshold als // offline. Logged nur wenn Rows betroffen sind (sonst floodet das Log // mit "0 rows" alle 30s). func runStaleSweep(ctx context.Context, pool *pgxpoolPool) { if pool == nil { return } swCtx, cancel := context.WithTimeout(ctx, 5*time.Second) defer cancel() flipped, err := cluster.SweepStaleNodes(swCtx, pool, staleThreshold) if err != nil { slog.Warn("scheduler: stale-sweep failed", "error", err) return } if flipped > 0 { slog.Info("scheduler: marked stale peers offline", "count", flipped, "threshold", staleThreshold) } } // pgxpoolPool ist ein lokaler Alias damit die Signatur stabil bleibt // wenn wir später den pool austauschen wollen (z.B. read-only-replica). type pgxpoolPool = pgxpool.Pool // runBackup führt einen scheduled Backup aus + prunet alte. Failures // loggen wir + alarmieren — verlorene Backups sind kritisch. func runBackup(ctx context.Context, svc *backup.Service, version string, a *alerts.Service, setupStore *setup.Store) { res, err := svc.Run(ctx, backup.KindScheduled, version) if err != nil { slog.Warn("scheduler: backup failed", "error", err, "file", res.File) if a != nil { _, _ = a.Fire(ctx, "backup.failed", alerts.SeverityError, "Backup fehlgeschlagen", "Scheduled Backup konnte nicht erstellt werden: "+err.Error()) } return } slog.Info("scheduler: backup done", "file", res.File, "size", res.SizeBytes, "db_bytes", res.DBDumpBytes, "files_bytes", res.FilesBytes, "sha256", res.SHA256) keepN := backup.DefaultKeepN if setupStore != nil { if st, err := setupStore.Load(); err == nil && st != nil && st.BackupRetentionKeep > 0 { keepN = st.BackupRetentionKeep } } if err := svc.Prune(ctx, keepN); err != nil { slog.Warn("scheduler: backup prune failed", "error", err) } } // runLicenseVerify performs a single re-verify pass. Empty key = no-op // (box stays in trial), so this is safe to call on every tick. // Bei valid:false-Antwort + Stand >7d alt → Warnung an Alerts. func runLicenseVerify(ctx context.Context, c *license.Client, ks *license.KeyStore, repo *licsvc.Repo, nodeID string, a *alerts.Service, d *dedupe) { key := ks.Get() if key == "" { slog.Debug("scheduler: license verify skipped — no key") return } res, err := c.Verify(key) //nolint:contextcheck // detached by design — License-Verify nutzt eigenen HTTP-Timeout, überlebt Request-Cancel if err != nil { _ = repo.MarkError(ctx, key, err.Error()) slog.Warn("scheduler: license verify failed", "error", err) return } payload, _ := json.Marshal(res) status := "active" if !res.Valid { status = "expired" if res.Status == "revoked" { status = "invalid" } } if err := repo.Upsert(ctx, key, status, res.ExpiresAt, nodeID, 0, payload, ""); err != nil { slog.Warn("scheduler: license db upsert failed", "error", err) return } slog.Info("scheduler: license verified", "status", status, "valid", res.Valid, "expires_at", res.ExpiresAt) // Alarm bei ungültiger Lizenz (revoked, expired) — dedupe 12h damit // der Operator nicht alle 24h denselben Alert bekommt. if a != nil && d != nil && !res.Valid { if d.shouldFire("license.invalid") { _, _ = a.Fire(ctx, "license.invalid", alerts.SeverityError, "License "+status, "License-Server liefert valid=false. Reason: "+res.Reason) } } } func runRenewer(ctx context.Context, r *certrenewer.Service, a *alerts.Service, d *dedupe) { res, err := r.Run(ctx) if err != nil { slog.Error("scheduler: renewer run failed", "error", err) if a != nil && d != nil && d.shouldFire("cert.renewer.run_failed") { _, _ = a.Fire(ctx, "cert.renewer.run_failed", alerts.SeverityError, "ACME-Renewer-Lauf fehlgeschlagen", "Certrenewer-Cycle abgebrochen: "+err.Error()) } return } slog.Info("scheduler: renewer pass complete", "checked", res.Checked, "renewed", res.Renewed, "failed", res.Failed, "skipped", res.Skipped) if a != nil && d != nil { for _, domain := range res.FailedDomains { key := "cert.renew_failed:" + domain if !d.shouldFire(key) { continue } _, _ = a.Fire(ctx, "cert.renew_failed", alerts.SeverityError, "Cert-Renewal fehlgeschlagen: "+domain, "Let's Encrypt Erneuerung für "+domain+" ist fehlgeschlagen. "+ "Prüfe ACME-Configuration und DNS-Erreichbarkeit. "+ "Nächster Versuch beim nächsten Renewer-Tick (alle 6h).") } } } // schedRemoteAdapter ist die scheduler-seitige Kopie des Adapters // aus edgeguard-api — gleicher Code, separater Type damit kein // Cross-Binary-Import nötig wird. type schedRemoteAdapter struct{ s *backupremote.Service } func newSchedRemoteAdapter(s *backupremote.Service) backup.RemoteUploader { return schedRemoteAdapter{s: s} } func (a schedRemoteAdapter) UploadAll(ctx context.Context, localPath string) ([]backup.RemoteUploadInfo, error) { res, err := a.s.UploadAll(ctx, localPath) out := make([]backup.RemoteUploadInfo, len(res)) for i, r := range res { out[i] = backup.RemoteUploadInfo{ RemoteID: r.RemoteID, RemoteName: r.RemoteName, OK: r.OK, SizeBytes: r.SizeBytes, DurationMs: r.DurationMs, Error: r.Error, } } return out, err }