Bisher pushte nur der Secondary seine Liveness an den Primary (runPrimaryPush). Der Primary pushte nichts → in der lokalen ha_nodes des Secondary fror die Primary-Row nach dem Boot ein → die vom Secondary ausgelieferte UI zeigte den Primary als offline. Neu: runPeerPush auf dem Primary/Founder pusht alle 30s self (role=primary) an jeden Peer via mTLS (/agent/cluster/peers). PushSelfToPeer(role) generalisiert PushSelfToPrimary; registerPeerRequest+AgentRegisterPeer akzeptieren ein role-Feld (default 'peer' → joining-Peer-Verhalten unverändert). Peer-Register-Log bei Routine-Pushes auf Debug (Info nur bei neuem Peer/IP-Wechsel) gegen 30s-Spam. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
926 lines
38 KiB
Go
926 lines
38 KiB
Go
// Command edgeguard-api serves the management REST API on
|
||
// 127.0.0.1:9443. HAProxy (or a dev curl) terminates TLS in front of
|
||
// it; this process is plain HTTP behind that.
|
||
package main
|
||
|
||
import (
|
||
"context"
|
||
"crypto/rand"
|
||
"log"
|
||
"log/slog"
|
||
"net/http"
|
||
"os"
|
||
"path/filepath"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/gin-gonic/gin"
|
||
"github.com/jackc/pgx/v5/pgxpool"
|
||
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/cluster"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/database"
|
||
firewallrender "git.netcell-it.de/projekte/edgeguard-native/internal/firewall"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/haproxy"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/handlers"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/license"
|
||
licsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/license"
|
||
chronyrender "git.netcell-it.de/projekte/edgeguard-native/internal/chrony"
|
||
squidrender "git.netcell-it.de/projekte/edgeguard-native/internal/squid"
|
||
unboundrender "git.netcell-it.de/projekte/edgeguard-native/internal/unbound"
|
||
wgrender "git.netcell-it.de/projekte/edgeguard-native/internal/wireguard"
|
||
kearender "git.netcell-it.de/projekte/edgeguard-native/internal/kea"
|
||
radiusrender "git.netcell-it.de/projekte/edgeguard-native/internal/freeradius"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/handlers/response"
|
||
"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/backendservers"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/backup"
|
||
backupremote "git.netcell-it.de/projekte/edgeguard-native/internal/services/backup/remote"
|
||
dnssvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/dns"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/aggregator"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/cluster/clustertls"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/cluster/jointoken"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/clusterjoin"
|
||
aptsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/apt"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/domainheaders"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/domains"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/firewall"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/firewalllog"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/syslogs"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/forwardproxy"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/ipaddresses"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/networkifs"
|
||
ntpsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/ntp"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/routingrules"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/secrets"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/staticroutes"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/session"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/setup"
|
||
"git.netcell-it.de/projekte/edgeguard-native/internal/services/tlscerts"
|
||
wgsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/wireguard"
|
||
dhcpsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/dhcp"
|
||
oidcsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/oidc"
|
||
radiussvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/radius"
|
||
usersvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/users"
|
||
wafsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/waf"
|
||
)
|
||
|
||
var version = "1.2.35"
|
||
|
||
func main() {
|
||
addr := os.Getenv("EDGEGUARD_API_ADDR")
|
||
if addr == "" {
|
||
addr = "127.0.0.1:9443"
|
||
}
|
||
|
||
dataDir := os.Getenv("EDGEGUARD_DATA_DIR")
|
||
if dataDir == "" {
|
||
dataDir = setup.DefaultDir
|
||
}
|
||
setupStore := setup.NewStore(dataDir)
|
||
|
||
signer, err := session.NewSignerFromPath("")
|
||
if err != nil {
|
||
// /var/lib/edgeguard not writable in dev → fall back to a
|
||
// process-local secret so `go run` works without sudo. Tokens
|
||
// won't survive a restart, which is fine for an unprivileged
|
||
// developer machine.
|
||
slog.Warn("session signer: persisted secret unavailable, using ephemeral",
|
||
"error", err)
|
||
signer = session.NewSigner(randomEphemeralSecret(), nil, 0)
|
||
}
|
||
|
||
gin.SetMode(gin.ReleaseMode)
|
||
r := gin.New()
|
||
r.Use(handlers.Recover())
|
||
|
||
// Health endpoints are mounted *before* SetupGate so they answer
|
||
// 200 even on a virgin box. UI uses /api/v1/system/health for the
|
||
// post-upgrade version-flip poll.
|
||
r.GET("/healthz", func(c *gin.Context) {
|
||
response.OK(c, gin.H{"status": "ok", "version": version})
|
||
})
|
||
r.GET("/api/health", func(c *gin.Context) {
|
||
response.OK(c, gin.H{"status": "ok", "version": version})
|
||
})
|
||
|
||
// ACME HTTP-01 webroot — HAProxy proxies these through pre-setup
|
||
// so certbot can issue the first cert. Webroot location matches
|
||
// certbot's default; override via EDGEGUARD_ACME_WEBROOT for
|
||
// dev/tests.
|
||
acmeWebroot := os.Getenv("EDGEGUARD_ACME_WEBROOT")
|
||
handlers.NewACMEHandler(acmeWebroot).Register(r)
|
||
|
||
v1 := r.Group("/api/v1")
|
||
v1.Use(handlers.SetupGate(setupStore))
|
||
|
||
requireAuth := handlers.RequireAuth(signer)
|
||
|
||
setupHdl := handlers.NewSetupHandler(setupStore).WithVersion(version)
|
||
setupHdl.Register(v1)
|
||
|
||
// systemHdl exists früh damit sowohl der frühe (DB-pool nicht
|
||
// nötige) Pfad als auch der späte WithMaintenance-Hookup gehen.
|
||
systemHdl := handlers.NewSystemHandler(version)
|
||
systemHdl.Register(v1)
|
||
authHdl := handlers.NewAuthHandler(setupStore, signer)
|
||
authHdl.Register(v1, requireAuth)
|
||
|
||
// Background-Refresh für apt-cache: hält die Apt-Lists alle 5 min
|
||
// frisch, damit der UI-Update-Banner kurz nach `make publish` ein
|
||
// verfügbares Update sieht — ohne den Background-Timer wäre der
|
||
// Cache nur nach UI-Polls aktuell und der Throttle würde
|
||
// Aktualisierungen zwischen den Polls schlucken.
|
||
aptsvc.StartBackgroundRefresh(context.Background())
|
||
|
||
// agentHdl wird vom Agent-Listener mit-gemountet (Phase 3.5).
|
||
// Nil-safe — wenn DB nicht offen ist, läuft der Agent-Listener
|
||
// nur mit den read-only System-Endpoints.
|
||
var agentHdl *handlers.ClusterHandler
|
||
|
||
// Open the DB pool best-effort. Without a reachable PG, CRUD
|
||
// handlers stay unregistered and only Auth/Setup/System answer —
|
||
// good enough for `go run` on a developer machine that has no
|
||
// postgres-16 yet.
|
||
pool, err := openDBBestEffort()
|
||
if err != nil {
|
||
slog.Warn("DB pool unavailable, CRUD endpoints disabled",
|
||
"error", err)
|
||
} else {
|
||
slog.Info("DB pool open, registering CRUD handlers")
|
||
|
||
nodeID, nodeErr := cluster.EnsureNodeID("")
|
||
if nodeErr != nil {
|
||
slog.Warn("node-id not persisted, using ephemeral",
|
||
"id", nodeID, "error", nodeErr)
|
||
}
|
||
clusterStore := cluster.NewStore(pool)
|
||
|
||
// Self-register in ha_nodes — only if setup is complete
|
||
// (we want the operator-defined FQDN, not the OS hostname,
|
||
// to land in api_url). Failures are logged but non-fatal.
|
||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||
st, _ := setupStore.Load()
|
||
if st != nil && st.Completed {
|
||
// Auto-create /etc/edgeguard/node.conf falls fehlt.
|
||
_, _ = cluster.EnsureLocalConfig("")
|
||
if _, err := cluster.EnsureSelfRegistered(ctx, clusterStore, st.FQDN, "primary", version); err != nil {
|
||
slog.Warn("self-register in ha_nodes failed", "error", err)
|
||
}
|
||
}
|
||
cancel()
|
||
|
||
// Phase 3.2: alle 30s Heartbeat (last_seen, status, version,
|
||
// config_hash) für die eigene Row. Goroutine läuft so lange wie
|
||
// die API — beim graceful Shutdown stoppt sie via ctx.Done().
|
||
// Hält den eigenen Node-Status auch dann frisch wenn der
|
||
// Scheduler gerade down ist; ein crashender API stoppt den
|
||
// Heartbeat → Peer-Sweeper markiert binnen 2 min "offline".
|
||
if nodeID != "" {
|
||
go runClusterHeartbeat(context.Background(), pool, nodeID, version)
|
||
}
|
||
|
||
// Secondary: push config_hash to primary every 5 min so the primary's
|
||
// ha_nodes reflects actual state. Without this, the primary retains the
|
||
// stale hash written at join-time and the drift banner never clears.
|
||
// st.IsClusterNode + PrimaryFQDN are only set on joined secondary nodes.
|
||
if nodeID != "" && st != nil && st.IsClusterNode && st.PrimaryFQDN != "" {
|
||
if primaryURL, normErr := clusterjoin.NormalizePrimaryURL(st.PrimaryFQDN); normErr == nil {
|
||
go runPrimaryPush(context.Background(), pool, nodeID, st.FQDN, version, primaryURL)
|
||
} else {
|
||
slog.Warn("cluster: cannot normalize primary URL for push", "primary", st.PrimaryFQDN, "error", normErr)
|
||
}
|
||
// runSecondaryConfigRender wird weiter unten gestartet sobald
|
||
// clusterAggregator verfügbar ist (braucht mTLS-Client für Cert-Sync).
|
||
} else if nodeID != "" && st != nil && st.Completed && st.FQDN != "" {
|
||
// Primary/Founder (kein joined Secondary): self (role=primary) an
|
||
// alle Peers pushen, damit deren lokale ha_nodes den Primary frisch
|
||
// hält — sonst zeigt die vom Secondary ausgelieferte UI den Primary
|
||
// als offline. No-op solange keine Peers existieren (Single-Node).
|
||
go runPeerPush(context.Background(), pool, clusterStore, nodeID, st.FQDN, version)
|
||
}
|
||
|
||
// Phase 3.3: Cluster-CA + Peer-Cert. Founder-Pfad — auf einem
|
||
// frisch installierten Single-Node generieren wir die CA und
|
||
// signieren uns selbst, damit der Agent-Listener auf :8443
|
||
// gleich hochfahren kann. Joining-Nodes (Phase 3.4) werden den
|
||
// Pfad nicht durchlaufen: sie kriegen CA+Peer-Cert vom Primary
|
||
// gepusht und finden die Files bereits vor.
|
||
clusterTLSStore := clustertls.New("")
|
||
if !clusterTLSStore.HasCA() {
|
||
org := "edgeguard.local"
|
||
if st != nil && st.FQDN != "" {
|
||
org = st.FQDN
|
||
}
|
||
if err := clusterTLSStore.InitCA(org, nil); err != nil {
|
||
slog.Warn("cluster-tls: InitCA failed", "error", err)
|
||
} else {
|
||
slog.Info("cluster-tls: CA generated", "dir", clusterTLSStore.Dir, "org", org)
|
||
}
|
||
}
|
||
if clusterTLSStore.HasCA() && !clusterTLSStore.HasPeer() {
|
||
cn := "edgeguard-node"
|
||
var dnsNames []string
|
||
if st != nil && st.FQDN != "" {
|
||
cn = st.FQDN
|
||
dnsNames = []string{st.FQDN}
|
||
}
|
||
if err := clusterTLSStore.EnsureSelfSigned(cn, dnsNames, nil, nil); err != nil {
|
||
slog.Warn("cluster-tls: EnsureSelfSigned failed", "error", err)
|
||
} else {
|
||
slog.Info("cluster-tls: peer cert generated", "cn", cn)
|
||
}
|
||
}
|
||
|
||
// Aggregator nur aufsetzen wenn Cert-Material da ist — sonst
|
||
// fan-out scheitert sowieso am Handshake.
|
||
var clusterAggregator *aggregator.Aggregator
|
||
if clusterTLSStore.HasPeer() {
|
||
if clientTLS, err := clusterTLSStore.ClientTLSConfig(); err == nil {
|
||
clusterAggregator = aggregator.New(clientTLS)
|
||
slog.Info("cluster: aggregator ready", "agent_port", aggregator.DefaultAgentPort)
|
||
} else {
|
||
slog.Warn("cluster: ClientTLSConfig failed", "error", err)
|
||
}
|
||
}
|
||
|
||
// Secondary-Config-Render: jetzt wo der Aggregator bereit ist starten.
|
||
// Aggregator wird für Cert-Sync (mTLS GET /agent/cluster/tls-certs) benötigt.
|
||
if nodeID != "" && st != nil && st.IsClusterNode && st.PrimaryFQDN != "" {
|
||
go runSecondaryConfigRender(context.Background(), pool, secrets.New(""), clusterAggregator, nodeID)
|
||
}
|
||
|
||
auditRepo := audit.New(pool)
|
||
domainsRepo := domains.New(pool)
|
||
domainHeadersRepo := domainheaders.New(pool)
|
||
backendsRepo := backends.New(pool)
|
||
backendServersRepo := backendservers.New(pool)
|
||
routingRepo := routingrules.New(pool)
|
||
ifsRepo := networkifs.New(pool)
|
||
ipsRepo := ipaddresses.New(pool)
|
||
tlsRepo := tlscerts.New(pool)
|
||
fwZones := firewall.NewZonesRepo(pool)
|
||
fwAddrObj := firewall.NewAddressObjectsRepo(pool)
|
||
fwAddrGrp := firewall.NewAddressGroupsRepo(pool)
|
||
fwSvc := firewall.NewServicesRepo(pool)
|
||
fwSvcGrp := firewall.NewServiceGroupsRepo(pool)
|
||
fwRules := firewall.NewRulesRepo(pool)
|
||
fwNAT := firewall.NewNATRulesRepo(pool)
|
||
secretsBox := secrets.New("")
|
||
wgIfaces := wgsvc.NewInterfacesRepo(pool)
|
||
wgPeers := wgsvc.NewPeersRepo(pool)
|
||
fwdProxyRepo := forwardproxy.New(pool)
|
||
dnsRepo := dnssvc.New(pool)
|
||
ntpRepo := ntpsvc.New(pool)
|
||
|
||
// ACME (Let's Encrypt). Email comes from setup.json — the
|
||
// wizard collects acme_email and the issuer registers an
|
||
// account on first /tls-certs/issue call.
|
||
var acmeService handlers.LetsEncryptIssuer
|
||
if st != nil && st.ACMEEmail != "" {
|
||
acmeService = acme.New(st.ACMEEmail)
|
||
}
|
||
|
||
// HAProxy reload — re-rendert haproxy.cfg + sudo systemctl
|
||
// reload haproxy. Wird in Domains/Backends/RoutingRules-Handler
|
||
// injiziert, damit jede Änderung ohne expliziten render-config-
|
||
// Aufruf live geht. Errors werden geloggt, nicht failed
|
||
// (Row schon committed, Operator kann manuell re-triggern).
|
||
// Maintenance-Endpoints brauchen den Reloader — späte Wiring
|
||
// nachdem haproxyReloader-closure existiert.
|
||
haproxyReloaderForLater := func(ctx context.Context) error {
|
||
return haproxy.New(pool).Render(ctx)
|
||
}
|
||
systemHdl.WithMaintenance(setupStore, haproxyReloaderForLater)
|
||
|
||
// Audit-Wiring (Phase Polish): Settings + Auth-Mutationen
|
||
// landen jetzt im audit_log. Nodes-id ist die persistente
|
||
// /var/lib/edgeguard/node-id.
|
||
systemHdl.WithAudit(auditRepo, nodeID)
|
||
systemHdl.WithDB(pool)
|
||
systemHdl.WithConfigPreviewers(map[string]func(context.Context) (string, error){
|
||
"haproxy": haproxy.New(pool).RenderToString,
|
||
"nftables": firewallrender.New(pool).RenderToString,
|
||
"squid": squidrender.New(pool).RenderToString,
|
||
"unbound": unboundrender.New(pool).RenderToString,
|
||
"chrony": chronyrender.New(pool).RenderToString,
|
||
"wireguard": wgrender.New(pool, secretsBox).RenderToString,
|
||
})
|
||
setupHdl.WithAudit(auditRepo, nodeID)
|
||
setupHdl.WithClusterSupport(clusterStore, func(ctx context.Context) error {
|
||
return firewallrender.New(pool).Render(ctx)
|
||
})
|
||
// Cluster-Node-Startup: Primary in lokalen ha_nodes eintragen damit
|
||
// nftables @peer_ipv4 korrekt ist — auch ohne erneuten Join.
|
||
go setupHdl.StartupPeerSync()
|
||
usersRepo := usersvc.New(pool)
|
||
authHdl.WithAudit(auditRepo, nodeID).WithUsers(usersRepo).WithClusterTLS(clusterTLSStore)
|
||
systemHdl.WithUsers(usersRepo)
|
||
|
||
haproxyReloader := func(ctx context.Context) error {
|
||
return haproxy.New(pool).Render(ctx)
|
||
}
|
||
|
||
authed := v1.Group("")
|
||
authed.Use(requireAuth, handlers.RequireAdminForMutations())
|
||
setupHdl.RegisterAuthed(authed)
|
||
handlers.NewUsersHandler(usersRepo, auditRepo, nodeID).Register(authed)
|
||
|
||
// OIDC/Keycloak SSO — public Flow-Endpoints auf v1 (hinter SetupGate),
|
||
// Admin-Settings auf authed (PUT nur admin via RequireAdminForMutations).
|
||
oidcRepo := oidcsvc.New(pool, secretsBox)
|
||
oidcHdl := handlers.NewOIDCHandler(oidcRepo, oidcsvc.NewClient(oidcRepo), usersRepo, signer, setupStore).
|
||
WithAudit(auditRepo, nodeID)
|
||
oidcHdl.RegisterPublic(v1)
|
||
oidcHdl.RegisterAdmin(authed)
|
||
handlers.NewDomainsHandler(domainsRepo, routingRepo, domainHeadersRepo, auditRepo, nodeID, haproxyReloader).Register(authed)
|
||
handlers.NewBackendsHandler(backendsRepo, auditRepo, nodeID, haproxyReloader).Register(authed)
|
||
handlers.NewBackendServersHandler(backendServersRepo, auditRepo, nodeID, haproxyReloader).Register(authed)
|
||
handlers.NewRoutingRulesHandler(routingRepo, auditRepo, nodeID, haproxyReloader).Register(authed)
|
||
handlers.NewNetworksHandler(ifsRepo, ipsRepo, fwZones, auditRepo, nodeID).Register(authed)
|
||
handlers.NewIPAddressesHandler(ipsRepo, auditRepo, nodeID).Register(authed)
|
||
handlers.NewRoutesHandler(staticroutes.New(pool), staticroutes.NewGenerator(pool),
|
||
auditRepo, nodeID).Register(authed)
|
||
// Phase 3.4 — Join-Token-Service. CA-Fingerprint kommt aus
|
||
// dem clustertls.Store; ohne CA = nil Tokens, GenerateToken
|
||
// scheitert, IssueCert wird gar nicht erst gemountet.
|
||
var joinTokens *jointoken.Service
|
||
if clusterTLSStore.HasCA() {
|
||
joinTokens = jointoken.New(pool, func() (string, error) {
|
||
caCert, _, err := clusterTLSStore.LoadCA()
|
||
if err != nil {
|
||
return "", err
|
||
}
|
||
return jointoken.CAFingerprint16(caCert.Raw), nil
|
||
})
|
||
}
|
||
// PeerReloader: nach Auto-Register triggert das den firewall-
|
||
// Render damit peer_ipv4 frisch ist und der mTLS-Listener für
|
||
// den neuen Peer erreichbar wird. Best-effort.
|
||
peerReloader := func(ctx context.Context) error {
|
||
return firewallrender.New(pool).Render(ctx)
|
||
}
|
||
clusterHdl := handlers.NewClusterHandler(clusterStore, nodeID).
|
||
WithAggregator(clusterAggregator).
|
||
WithJoinFlow(clusterTLSStore, joinTokens).
|
||
WithPeerReloader(peerReloader).
|
||
WithAudit(auditRepo, nodeID).
|
||
WithVersion(version)
|
||
clusterHdl.Register(authed)
|
||
// /cluster/issue-cert läuft PUBLIC — joining Peer hat noch
|
||
// keine Session/Cert. Token + Nonce-Tracking ist die einzige
|
||
// Auth-Stufe.
|
||
clusterHdl.RegisterPublic(v1)
|
||
// Agent-Listener (mTLS) bekommt clusterHdl mit, damit Joiner
|
||
// sich via /agent/cluster/peers eintragen können.
|
||
agentHdl = clusterHdl
|
||
handlers.NewAuditHandler(auditRepo).Register(authed)
|
||
handlers.NewHAProxyStatsHandler().Register(authed)
|
||
|
||
// Firewall-Log (Phase 2): Tailer für /var/log/edgeguard/
|
||
// firewall.jsonl + HTTP-Tail + WebSocket-Live-Stream.
|
||
fwLogTailer := firewalllog.NewTailer(firewalllog.DefaultLogPath, 1000)
|
||
handlers.StartFirewallLogTailer(context.Background(), fwLogTailer)
|
||
handlers.NewFirewallLogHandler(fwLogTailer, firewalllog.DefaultLogPath).Register(authed)
|
||
|
||
// /logs (Phase 4): aggregierter Reader für journalctl + audit_log
|
||
handlers.NewLogsHandler(syslogs.New(auditRepo)).Register(authed)
|
||
|
||
// /backups — manueller Trigger + Liste + Download. Scheduled-
|
||
// Jobs laufen im edgeguard-scheduler.
|
||
backupSvc := backup.New(pool)
|
||
backupSvc.RemoteUploader = newBackupRemoteAdapter(backupremote.New(pool))
|
||
handlers.NewBackupHandler(backupSvc, auditRepo, nodeID, version).Register(authed)
|
||
handlers.NewBackupRemotesHandler(pool, auditRepo, nodeID).Register(authed)
|
||
handlers.NewDiagnosticsHandler().Register(authed)
|
||
handlers.NewAlertsHandler(alerts.New(pool), auditRepo, nodeID).Register(authed)
|
||
handlers.NewTLSCertsHandler(tlsRepo, auditRepo, nodeID, acmeService).Register(authed)
|
||
// Firewall reload: nach jeder Mutation den Renderer neu fahren
|
||
// (writes ruleset.nft + sudo nft -f). Errors loggen, nicht failen.
|
||
fwReloader := func(ctx context.Context) error {
|
||
return firewallrender.New(pool).Render(ctx)
|
||
}
|
||
handlers.NewFirewallHandler(fwZones, fwAddrObj, fwAddrGrp, fwSvc, fwSvcGrp, fwRules, fwNAT, auditRepo, nodeID, fwReloader, pool).Register(authed)
|
||
handlers.NewCrowdSecHandler(auditRepo, nodeID).Register(authed)
|
||
handlers.NewWafHandler(wafsvc.New(pool), auditRepo, nodeID, haproxyReloader).Register(authed)
|
||
|
||
// withFW wraps a service-reloader so that AFTER the service is
|
||
// reloaded, the firewall is also re-rendered. Necessary for
|
||
// services whose state feeds the auto-FW-rule generator (DNS
|
||
// listen-IPs, Squid ACL count, WG listen-port, NTP serve-clients).
|
||
// Service-Reload-Errors propagieren; FW-Errors werden nur
|
||
// geloggt (DB-Row ist commited, FW kann nachgezogen werden).
|
||
withFW := func(svc func(context.Context) error) func(context.Context) error {
|
||
return func(ctx context.Context) error {
|
||
if err := svc(ctx); err != nil {
|
||
return err
|
||
}
|
||
if err := fwReloader(ctx); err != nil {
|
||
slog.Warn("firewall: re-render after service mutation failed", "error", err)
|
||
}
|
||
return nil
|
||
}
|
||
}
|
||
|
||
// WireGuard reload: re-render /etc/edgeguard/wireguard/*.conf
|
||
// + restart wg-quick@<iface>. Same pattern as the haproxy +
|
||
// firewall reloaders. WG braucht FW-Trigger (server-mode
|
||
// listen-port wird Auto-Rule).
|
||
wgReloader := func(ctx context.Context) error {
|
||
return wgrender.New(pool, secretsBox).Render(ctx)
|
||
}
|
||
handlers.NewWireguardHandler(wgIfaces, wgPeers, secretsBox, auditRepo, nodeID, withFW(wgReloader)).Register(authed)
|
||
|
||
// Squid forward-proxy reload — re-render squid.conf + reload
|
||
// squid.service. sudoers im postinst whitelistet das. ACL-Count
|
||
// triggert Auto-FW-Rule für tcp/3128.
|
||
squidReloader := func(ctx context.Context) error {
|
||
return squidrender.New(pool).Render(ctx)
|
||
}
|
||
handlers.NewForwardProxyHandler(fwdProxyRepo, auditRepo, nodeID, withFW(squidReloader)).Register(authed)
|
||
|
||
// Unbound DNS reload — re-render edgeguard.conf + restart
|
||
// unbound. Listen-IPs triggern Auto-FW-Rule für udp/tcp 53.
|
||
unboundReloader := func(ctx context.Context) error {
|
||
return unboundrender.New(pool).Render(ctx)
|
||
}
|
||
handlers.NewDNSHandler(dnsRepo, auditRepo, nodeID, withFW(unboundReloader)).Register(authed)
|
||
|
||
// Chrony NTP reload — re-render edgeguard.conf + restart chrony.
|
||
// Listen-IPs + serve_clients triggern Auto-FW-Rule für udp/123.
|
||
chronyReloader := func(ctx context.Context) error {
|
||
return chronyrender.New(pool).Render(ctx)
|
||
}
|
||
handlers.NewNTPHandler(ntpRepo, auditRepo, nodeID, withFW(chronyReloader)).Register(authed)
|
||
|
||
// DHCP (Kea) — re-render kea-dhcp4.conf + manage service lifecycle.
|
||
keaReloader := func(ctx context.Context) error {
|
||
return kearender.New(pool).Render(ctx)
|
||
}
|
||
handlers.NewDHCPHandler(dhcpsvc.New(pool), auditRepo, nodeID, withFW(keaReloader)).Register(authed)
|
||
|
||
// RADIUS (FreeRADIUS) — re-render clients.conf + authorize + service lifecycle.
|
||
radiusReloader := func(ctx context.Context) error {
|
||
return radiusrender.New(pool, secretsBox).Render(ctx)
|
||
}
|
||
handlers.NewRADIUSHandler(radiussvc.New(pool, secretsBox), auditRepo, nodeID, withFW(radiusReloader)).Register(authed)
|
||
|
||
// Wire all service reloaders into systemHdl so RenderConfigs
|
||
// re-renders every service from DB state in one shot.
|
||
systemHdl.WithAllReloaders(map[string]func(context.Context) error{
|
||
"nftables": fwReloader,
|
||
"wireguard": wgReloader,
|
||
"squid": squidReloader,
|
||
"unbound": unboundReloader,
|
||
"chrony": chronyReloader,
|
||
"kea": keaReloader,
|
||
"freeradius": radiusReloader,
|
||
})
|
||
|
||
// License — node-local key store + DB-mirror of last verify
|
||
// result. Real verify runs against license.netcell-it.com via
|
||
// internal/license; the scheduler triggers daily re-verify.
|
||
licRepo := licsvc.New(pool)
|
||
licClient := license.NewClient()
|
||
licKeyStore := license.NewKeyStore()
|
||
handlers.NewLicenseHandler(licRepo, licKeyStore, licClient, auditRepo, nodeID).Register(authed)
|
||
// Kick off periodic re-verify in this process so a long-running
|
||
// api answers /license/status with fresh data even without the
|
||
// scheduler. StartPeriodicVerification is a no-op when the key
|
||
// is empty.
|
||
licClient.StartPeriodicVerification(licKeyStore.Get())
|
||
|
||
// Startup-Render nftables: stellt sicher dass Template-Änderungen
|
||
// aus einem Update (z.B. neue WireGuard forward-Chain-Auto-Regel)
|
||
// sofort nach dem API-Restart aktiv werden — ohne dass der
|
||
// Operator manuell eine Mutation triggern müsste. nft -f ist
|
||
// idempotent und atomar; kein Dienst wird neu gestartet.
|
||
go func() {
|
||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||
defer cancel()
|
||
if err := firewallrender.New(pool).Render(ctx); err != nil {
|
||
slog.Warn("startup: nftables render failed", "error", err)
|
||
}
|
||
}()
|
||
}
|
||
|
||
mountUI(r)
|
||
|
||
// Phase 3.3: zweiter Listener auf :8443 mit mTLS für Cluster-Peer-
|
||
// Reads. RequireAndVerifyClientCert gegen unsere Cluster-CA — wer
|
||
// keinen CA-signierten Cert hat, kommt nicht durch den Handshake.
|
||
// Listener wird nur gestartet wenn Cert-Material vorhanden ist;
|
||
// auf einer frisch installierten Box hat die Init-Phase oben das
|
||
// schon erledigt.
|
||
startAgentListener(version, agentHdl, systemHdl)
|
||
|
||
// Nach einem Upgrade-Neustart: wenn die State-Datei "updating-primary"
|
||
// enthält, sind wir gerade neu gestartet → Update abgeschlossen → "done".
|
||
handlers.FinishRollingUpdateIfPending()
|
||
|
||
log.Printf("edgeguard-api %s listening on %s", version, addr)
|
||
srv := &http.Server{Addr: addr, Handler: r}
|
||
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
||
log.Fatalf("edgeguard-api: %v", err)
|
||
}
|
||
}
|
||
|
||
// startAgentListener startet den mTLS-Agent-Listener auf :8443 als
|
||
// Goroutine. Mountet nur read-only Endpoints (siehe SystemHandler.
|
||
// RegisterAgent — health + resources). Fehler im Cert-Load = no-op
|
||
// + log; Fehler beim Listen.Serve loggen wir aber lassen die API
|
||
// weiterlaufen.
|
||
func startAgentListener(version string, clusterHdl *handlers.ClusterHandler, sysHdl *handlers.SystemHandler) {
|
||
store := clustertls.New("")
|
||
serverTLS, err := store.ServerTLSConfig()
|
||
if err != nil {
|
||
slog.Info("cluster: agent listener disabled (no cert material)", "error", err)
|
||
return
|
||
}
|
||
addr := os.Getenv("EDGEGUARD_AGENT_LISTEN")
|
||
if addr == "" {
|
||
// 0.0.0.0:8443 — Auth via mTLS, also unbedenklich auf Public-IP.
|
||
// nft anti-lockout-Regel + Peer-IP-Set bestimmen wer überhaupt
|
||
// connecten darf. Loopback-Tests gehen direkt.
|
||
addr = "0.0.0.0:8443"
|
||
}
|
||
r := gin.New()
|
||
r.Use(gin.Recovery())
|
||
// Kein /api/v1-Prefix auf dem Agent-Listener: das Versioning kommt
|
||
// hier implizit aus dem Binary (Peer-Roundtrip ist immer same-major).
|
||
// Aggregator-Aufrufer sehen /agent/... direkt.
|
||
root := r.Group("")
|
||
// Nutze den gewiredeten systemHdl (mit Users + Setup) damit
|
||
// AgentAuthCheck Credentials gegen die echte DB prüfen kann.
|
||
if sysHdl != nil {
|
||
sysHdl.RegisterAgent(root)
|
||
} else {
|
||
handlers.NewSystemHandler(version).RegisterAgent(root)
|
||
}
|
||
if clusterHdl != nil {
|
||
// Phase 3.5: /agent/cluster/peers (Auto-Register).
|
||
clusterHdl.RegisterAgent(root)
|
||
}
|
||
|
||
srv := &http.Server{
|
||
Addr: addr,
|
||
Handler: r,
|
||
TLSConfig: serverTLS,
|
||
ReadHeaderTimeout: 10 * time.Second,
|
||
}
|
||
go func() {
|
||
slog.Info("cluster: agent (mTLS) listener starting", "addr", addr)
|
||
if err := srv.ListenAndServeTLS("", ""); err != nil && err != http.ErrServerClosed {
|
||
slog.Error("cluster: agent listener", "error", err)
|
||
}
|
||
}()
|
||
}
|
||
|
||
// mountUI serves the management UI — Vite-built static assets under
|
||
// /usr/share/edgeguard/ui/ — with SPA fallback (any path that isn't
|
||
// /api/* or /healthz and isn't a real file → index.html). When the
|
||
// dist directory is missing (dev box without `bun run build`), a
|
||
// placeholder HTML page is served at /.
|
||
func mountUI(r *gin.Engine) {
|
||
uiDir := os.Getenv("EDGEGUARD_UI_DIR")
|
||
if uiDir == "" {
|
||
uiDir = "/usr/share/edgeguard/ui"
|
||
}
|
||
indexPath := filepath.Join(uiDir, "index.html")
|
||
|
||
if _, err := os.Stat(indexPath); err != nil {
|
||
slog.Warn("UI dist not found, serving placeholder",
|
||
"ui_dir", uiDir, "error", err)
|
||
r.NoRoute(func(c *gin.Context) {
|
||
path := c.Request.URL.Path
|
||
if isAPIPath(path) {
|
||
c.Status(http.StatusNotFound)
|
||
return
|
||
}
|
||
c.Data(http.StatusOK, "text/html; charset=utf-8", []byte(uiPlaceholder))
|
||
})
|
||
return
|
||
}
|
||
|
||
r.NoRoute(func(c *gin.Context) {
|
||
path := c.Request.URL.Path
|
||
if isAPIPath(path) {
|
||
c.Status(http.StatusNotFound)
|
||
return
|
||
}
|
||
// Serve real file when one exists for the requested path.
|
||
// filepath.Clean blocks `..` traversal; the join still pins
|
||
// the result inside uiDir even with shenanigans.
|
||
clean := filepath.Clean(path)
|
||
if !strings.HasPrefix(clean, "/") {
|
||
clean = "/" + clean
|
||
}
|
||
full := filepath.Join(uiDir, clean)
|
||
if !strings.HasPrefix(full, uiDir) {
|
||
c.Status(http.StatusForbidden)
|
||
return
|
||
}
|
||
if info, err := os.Stat(full); err == nil && !info.IsDir() {
|
||
// Vite hashed assets are immutable — cache them forever.
|
||
// index.html must never be cached so updates take effect.
|
||
if strings.HasPrefix(clean, "/assets/") {
|
||
c.Header("Cache-Control", "public, max-age=31536000, immutable")
|
||
} else {
|
||
c.Header("Cache-Control", "no-cache, no-store, must-revalidate")
|
||
}
|
||
c.File(full)
|
||
return
|
||
}
|
||
// SPA fallback — React Router renders the right page.
|
||
c.Header("Cache-Control", "no-cache, no-store, must-revalidate")
|
||
c.File(indexPath)
|
||
})
|
||
}
|
||
|
||
// isAPIPath returns true for paths the API owns; UI serves
|
||
// everything else. /healthz and /api/health are technically API
|
||
// surfaces but don't need to fall through to index.html either.
|
||
func isAPIPath(p string) bool {
|
||
return strings.HasPrefix(p, "/api/") || p == "/healthz" || p == "/api/health"
|
||
}
|
||
|
||
const uiPlaceholder = `<!doctype html>
|
||
<html lang="en"><head><meta charset="utf-8"><title>EdgeGuard</title></head>
|
||
<body style="font-family: -apple-system, sans-serif; max-width: 640px; margin: 4em auto; line-height: 1.5;">
|
||
<h1>EdgeGuard</h1>
|
||
<p>The management UI has not been built yet. From the project root, run:</p>
|
||
<pre style="background:#f4f4f4;padding:12px;border-radius:4px;">cd management-ui && bun install && bun run build</pre>
|
||
<p>Then the same URL will serve the React SPA. The REST API is fully functional at
|
||
<code>/api/v1/*</code> regardless.</p>
|
||
</body></html>`
|
||
|
||
// openDBBestEffort opens the pool with a 3s timeout. Returns the
|
||
// non-nil error so callers can decide whether to register CRUD or
|
||
// degrade gracefully.
|
||
func openDBBestEffort() (*pgxpoolPool, error) {
|
||
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
|
||
defer cancel()
|
||
dsn := database.ConnStringFromEnv()
|
||
return database.Open(ctx, dsn)
|
||
}
|
||
|
||
// pgxpoolPool aliases the concrete pool type so we don't import it in
|
||
// main.go on every platform — keeps the import block lean.
|
||
type pgxpoolPool = pgxpool.Pool
|
||
|
||
// nodeIDOrHostname returns the node identifier audit_log entries are
|
||
// stamped with. v1 just uses /etc/machine-id (or the hostname on dev
|
||
// machines without one). Phase 3's cluster store will replace this.
|
||
func nodeIDOrHostname() string {
|
||
if b, err := os.ReadFile("/etc/machine-id"); err == nil {
|
||
s := string(b)
|
||
s = stripTrailingNewline(s)
|
||
if s != "" {
|
||
return s
|
||
}
|
||
}
|
||
if h, err := os.Hostname(); err == nil {
|
||
return h
|
||
}
|
||
return "unknown"
|
||
}
|
||
|
||
func stripTrailingNewline(s string) string {
|
||
for len(s) > 0 && (s[len(s)-1] == '\n' || s[len(s)-1] == '\r') {
|
||
s = s[:len(s)-1]
|
||
}
|
||
return s
|
||
}
|
||
|
||
// randomEphemeralSecret is the fallback for dev environments where
|
||
// /var/lib/edgeguard isn't writable. Tokens issued with this secret
|
||
// die on restart — production reads/writes the persistent file via
|
||
// session.NewSignerFromPath.
|
||
// backupRemoteAdapter überbrückt backup.RemoteUploader (Interface)
|
||
// und remote.Service. Die Field-Names sind gleich; nur der Type ist
|
||
// verschieden weil sonst Import-Cycle backup→remote→backup entstehen
|
||
// würde.
|
||
type backupRemoteAdapter struct{ s *backupremote.Service }
|
||
|
||
func newBackupRemoteAdapter(s *backupremote.Service) backup.RemoteUploader {
|
||
return backupRemoteAdapter{s: s}
|
||
}
|
||
func (a backupRemoteAdapter) 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
|
||
}
|
||
|
||
// runClusterHeartbeat tickt alle 30s und bumpt die eigene ha_nodes-Row
|
||
// (last_seen, status, version, config_hash) via cluster.Heartbeat.
|
||
// Fehler werden geloggt aber nicht zurückgegeben — der nächste Tick
|
||
// versucht es erneut. Beendet beim ctx.Done() (graceful API shutdown).
|
||
func runClusterHeartbeat(ctx context.Context, pool *pgxpoolPool, localID, version string) {
|
||
const tick = 30 * time.Second
|
||
t := time.NewTicker(tick)
|
||
defer t.Stop()
|
||
// Erster Schlag direkt nach Start damit die UI nicht 30s wartet.
|
||
if err := cluster.Heartbeat(ctx, pool, localID, version); err != nil {
|
||
slog.Warn("cluster: initial heartbeat failed", "error", err)
|
||
}
|
||
slog.Info("cluster: heartbeat goroutine started", "tick", tick.String(), "node_id", localID)
|
||
for {
|
||
select {
|
||
case <-ctx.Done():
|
||
slog.Info("cluster: heartbeat goroutine stopping")
|
||
return
|
||
case <-t.C:
|
||
hbCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
|
||
if err := cluster.Heartbeat(hbCtx, pool, localID, version); err != nil {
|
||
slog.Warn("cluster: heartbeat failed", "error", err)
|
||
}
|
||
cancel()
|
||
}
|
||
}
|
||
}
|
||
|
||
// runSecondaryConfigRender läuft auf Secondary-Nodes und re-rendert alle
|
||
// Service-Configs wenn die Logical Replication Änderungen vom Primary
|
||
// geliefert hat. Erkennt das an einem geänderten config_hash.
|
||
// Tick: 5 min — balanciert Reaktionszeit gegen Reload-Overhead.
|
||
//
|
||
// Cert-Sync läuft auf jedem Tick unabhängig vom config_hash, da certbot-
|
||
// Renewals auf dem Primary den Hash nicht ändern.
|
||
func runSecondaryConfigRender(ctx context.Context, pool *pgxpoolPool, box *secrets.Box, agg *aggregator.Aggregator, localID string) {
|
||
const tick = 5 * time.Minute
|
||
t := time.NewTicker(tick)
|
||
defer t.Stop()
|
||
var lastHash string
|
||
render := func() {
|
||
rCtx, cancel := context.WithTimeout(ctx, 90*time.Second)
|
||
defer cancel()
|
||
|
||
// TLS-Zertifikate bei jedem Tick synchronisieren — unabhängig vom
|
||
// config_hash, da certbot-Renewals den Hash nicht berühren.
|
||
if err := handlers.SyncTLSCertsFromPrimary(rCtx, pool, agg, localID); err != nil {
|
||
slog.Warn("cluster: cert sync failed", "error", err)
|
||
}
|
||
|
||
hash, err := cluster.ComputeConfigHash(rCtx, pool)
|
||
if err != nil || hash == lastHash {
|
||
return
|
||
}
|
||
lastHash = hash
|
||
slog.Info("cluster: secondary config changed via replication, re-rendering", "hash", hash)
|
||
// HAProxy
|
||
if err := haproxy.New(pool).Render(rCtx); err != nil {
|
||
slog.Warn("cluster: secondary haproxy render failed", "error", err)
|
||
}
|
||
// nftables
|
||
if err := firewallrender.New(pool).Render(rCtx); err != nil {
|
||
slog.Warn("cluster: secondary nftables render failed", "error", err)
|
||
}
|
||
// WireGuard — Interface-Configs + wg-quick@<iface> reload
|
||
if err := wgrender.New(pool, box).Render(rCtx); err != nil {
|
||
slog.Warn("cluster: secondary wireguard render failed", "error", err)
|
||
}
|
||
// Squid forward proxy
|
||
if err := squidrender.New(pool).Render(rCtx); err != nil {
|
||
slog.Warn("cluster: secondary squid render failed", "error", err)
|
||
}
|
||
// Unbound DNS
|
||
if err := unboundrender.New(pool).Render(rCtx); err != nil {
|
||
slog.Warn("cluster: secondary unbound render failed", "error", err)
|
||
}
|
||
// Chrony NTP
|
||
if err := chronyrender.New(pool).Render(rCtx); err != nil {
|
||
slog.Warn("cluster: secondary chrony render failed", "error", err)
|
||
}
|
||
// Netzwerk-Interfaces (VLAN/Bridge/Bond) — erstellt Interface-Objekte,
|
||
// weist aber KEINE IPs zu (das ist node-spezifisch und darf nicht aus
|
||
// der Replikation kommen — sonst IP-Konflikt mit dem Primary).
|
||
if err := networkifs.NewGenerator(networkifs.New(pool)).Render(rCtx); err != nil {
|
||
slog.Warn("cluster: secondary interfaces render failed", "error", err)
|
||
}
|
||
// IP-Adressen werden auf dem Secondary NICHT aus der Replikation
|
||
// angewendet. Jeder Node konfiguriert seine eigenen IPs statisch
|
||
// (z.B. /etc/network/interfaces). Floating-Service-IPs werden von
|
||
// Keepalived verwaltet — nicht vom Renderer.
|
||
}
|
||
// Initialer Check nach kurzem Delay (Replication braucht einen Moment)
|
||
select {
|
||
case <-ctx.Done():
|
||
return
|
||
case <-time.After(30 * time.Second):
|
||
render()
|
||
}
|
||
for {
|
||
select {
|
||
case <-ctx.Done():
|
||
return
|
||
case <-t.C:
|
||
render()
|
||
}
|
||
}
|
||
}
|
||
|
||
// runPrimaryPush periodically pushes this secondary node's config_hash to the
|
||
// primary via mTLS. The primary's ha_nodes view only gets config_hash + last_seen
|
||
// written during join-time autoRegister — after that the primary never hears about
|
||
// the secondary unless we push. Without this, the drift banner shows stale hashes
|
||
// from join-time forever AND the secondary's last_seen freezes → SweepStaleNodes
|
||
// marks it offline.
|
||
//
|
||
// WICHTIG: tick MUSS deutlich unter dem Stale-Threshold (4× 30s = 2 min, siehe
|
||
// scheduler.staleThreshold / cluster.SweepStaleNodes) liegen. Sonst flippt der
|
||
// Secondary zwischen den Pushes zwangsläufig auf "offline" (bei 5-min-Tick:
|
||
// 2 min online, 3 min offline). 30s = 4 Pushes pro Stale-Fenster → ein
|
||
// verpasster Push (Netz-Glitch) ist unkritisch. Der Receiver (AgentRegisterPeer)
|
||
// lädt nftables nur bei IP-Änderung neu → kein Reload-Sturm durch häufige Pushes.
|
||
func runPrimaryPush(ctx context.Context, pool *pgxpoolPool, nodeID, fqdn, version, primaryURL string) {
|
||
const tick = 30 * time.Second
|
||
t := time.NewTicker(tick)
|
||
defer t.Stop()
|
||
push := func() {
|
||
pCtx, cancel := context.WithTimeout(ctx, 15*time.Second)
|
||
defer cancel()
|
||
hash, _ := cluster.ComputeConfigHash(pCtx, pool)
|
||
if err := clusterjoin.PushSelfToPrimary(primaryURL, "", nodeID, fqdn, version, hash); err != nil {
|
||
slog.Warn("cluster: push-to-primary failed", "error", err)
|
||
} else {
|
||
slog.Debug("cluster: config_hash pushed to primary", "hash", hash)
|
||
}
|
||
}
|
||
push() // immediate push on API startup
|
||
for {
|
||
select {
|
||
case <-ctx.Done():
|
||
return
|
||
case <-t.C:
|
||
push()
|
||
}
|
||
}
|
||
}
|
||
|
||
// runPeerPush läuft auf dem Primary/Founder und pusht alle 30s die eigene
|
||
// Identität (role=primary) an jeden Peer via mTLS — das Gegenstück zu
|
||
// runPrimaryPush (Secondary→Primary). Zusammen ergibt das einen
|
||
// bidirektionalen Cross-Node-Heartbeat: beide Nodes sehen sich gegenseitig
|
||
// als online, egal von welchem Node die UI ausgeliefert wird. Tick wie
|
||
// runPrimaryPush deutlich unter dem 2-min-Stale-Threshold. No-op solange
|
||
// keine Peers existieren (Single-Node) bzw. wenn ein Peer down ist (Debug-Log).
|
||
func runPeerPush(ctx context.Context, pool *pgxpoolPool, store *cluster.Store, nodeID, fqdn, version string) {
|
||
const tick = 30 * time.Second
|
||
t := time.NewTicker(tick)
|
||
defer t.Stop()
|
||
push := func() {
|
||
pCtx, cancel := context.WithTimeout(ctx, 25*time.Second)
|
||
defer cancel()
|
||
peers, err := store.List(pCtx)
|
||
if err != nil {
|
||
slog.Warn("cluster: peer-push list failed", "error", err)
|
||
return
|
||
}
|
||
hash, _ := cluster.ComputeConfigHash(pCtx, pool)
|
||
for i := range peers {
|
||
p := peers[i]
|
||
if p.ID == nodeID {
|
||
continue // nicht an sich selbst pushen
|
||
}
|
||
target := p.APIURL
|
||
if target == "" {
|
||
target = "https://" + p.FQDN
|
||
}
|
||
if err := clusterjoin.PushSelfToPeer(target, "", nodeID, fqdn, version, hash, "primary"); err != nil {
|
||
slog.Debug("cluster: push-to-peer failed", "peer", p.FQDN, "error", err)
|
||
}
|
||
}
|
||
}
|
||
push() // immediate push on API startup
|
||
for {
|
||
select {
|
||
case <-ctx.Done():
|
||
return
|
||
case <-t.C:
|
||
push()
|
||
}
|
||
}
|
||
}
|
||
|
||
func randomEphemeralSecret() []byte {
|
||
b := make([]byte, 32)
|
||
if _, err := rand.Read(b); err != nil {
|
||
// Should never happen on a sane Linux box; fall back to a
|
||
// time-based filler so the process can at least start.
|
||
log.Printf("WARN: crypto/rand read failed: %v", err)
|
||
}
|
||
return b
|
||
}
|