Files
edgeguard-native/cmd/edgeguard-api/main.go
Debian bc6db1fc2b feat(cluster): Rolling Update — Secondary-first upgrade orchestration (v1.2.3)
POST /cluster/rolling-update startet den gestaffelten Upgrade-Prozess:
1. Secondary via mTLS /agent/cluster/trigger-update anstoßen
2. /agent/cluster/version pollen bis Secondary Version-Flip zeigt (max 10 min)
3. Primary self-upgrade via systemd-run (identisch zu /system/upgrade)

State wird in /var/lib/edgeguard/rolling-update-state.json persistiert:
Phasen: updating-secondary → waiting-secondary → updating-primary.
"done" wird nicht geschrieben — Prozess stirbt beim Upgrade. UI erkennt
Abschluss via /system/health version-flip (analog Single-Node-Upgrade).

UI: UpdateBanner erkennt Cluster-Modus (/cluster/status mode="cluster")
und tauscht den "Install now"-Button gegen "Rolling Update (Cluster)" aus.
Multi-Step-Modal zeigt die drei Phasen; ab updating-primary wechselt der
Client auf /system/health polling.

Aggregator.PostPeer: neuer einzel-POST-Helper für mTLS-trigger-update.
WithVersion(): ClusterHandler bekommt Binary-Version für /agent/cluster/version.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-29 23:40:49 +02:00

788 lines
32 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"
"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"
usersvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/users"
)
var version = "1.2.3"
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)
}
// Logical Replication liefert Änderungen automatisch — aber Service-
// Configs (haproxy.cfg, nftables …) müssen nach jeder Änderung neu
// gerendert werden. Diese Goroutine erkennt hash-Änderungen und rendert.
go runSecondaryConfigRender(context.Background(), pool)
}
// 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)
}
}
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)
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).
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)
// 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)
// 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,
})
// 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)
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() {
c.File(full)
return
}
// SPA fallback — React Router renders the right page.
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 &amp;&amp; bun install &amp;&amp; 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.
func runSecondaryConfigRender(ctx context.Context, pool *pgxpoolPool) {
const tick = 5 * time.Minute
t := time.NewTicker(tick)
defer t.Stop()
var lastHash string
render := func() {
rCtx, cancel := context.WithTimeout(ctx, 60*time.Second)
defer cancel()
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)
}
// Weitere Dienste (Squid, Unbound, Chrony, WireGuard) werden bei
// Änderungen an ihren spezifischen Tabellen ebenfalls neu gerendert.
// render-config ohne Reload: die Dienste merken Änderungen selbst
// (HAProxy/nftables über systemctl reload, der oben bereits läuft).
}
// 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 written
// during join-time autoRegister — after that the primary never hears about
// hash changes unless we push. Without this, the drift banner shows stale
// hashes from join-time forever.
func runPrimaryPush(ctx context.Context, pool *pgxpoolPool, nodeID, fqdn, version, primaryURL string) {
const tick = 5 * time.Minute
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()
}
}
}
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
}