// 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" "errors" "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/aggregator" chronyrender "git.netcell-it.de/projekte/edgeguard-native/internal/chrony" "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/cluster/jointoken" "git.netcell-it.de/projekte/edgeguard-native/internal/crowdsec" "git.netcell-it.de/projekte/edgeguard-native/internal/database" firewallrender "git.netcell-it.de/projekte/edgeguard-native/internal/firewall" radiusrender "git.netcell-it.de/projekte/edgeguard-native/internal/freeradius" "git.netcell-it.de/projekte/edgeguard-native/internal/handlers" "git.netcell-it.de/projekte/edgeguard-native/internal/handlers/response" "git.netcell-it.de/projekte/edgeguard-native/internal/haproxy" kearender "git.netcell-it.de/projekte/edgeguard-native/internal/kea" "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" aptsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/apt" "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" "git.netcell-it.de/projekte/edgeguard-native/internal/services/clusterjoin" dhcpsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/dhcp" dnssvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/dns" "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/forwardproxy" "git.netcell-it.de/projekte/edgeguard-native/internal/services/ipaddresses" licsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/license" "git.netcell-it.de/projekte/edgeguard-native/internal/services/networkifs" ntpsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/ntp" oidcsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/oidc" radiussvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/radius" "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/session" "git.netcell-it.de/projekte/edgeguard-native/internal/services/setup" "git.netcell-it.de/projekte/edgeguard-native/internal/services/staticroutes" "git.netcell-it.de/projekte/edgeguard-native/internal/services/syslogs" "git.netcell-it.de/projekte/edgeguard-native/internal/services/tlscerts" usersvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/users" wafsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/waf" wgsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/wireguard" 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" ) 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-triggering). // 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 persistence // /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, "crowdsec-whitelist": crowdsec.NewWhitelistGenerator(pool).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) } // Domain-Mutationen rendern zusätzlich die CrowdSec-Admin-Whitelist neu // (Flag crowdsec_trusted → host-genaue Ausnahme). No-op ohne CrowdSec. // Beide laufen unabhängig; Fehler werden zusammengefasst (nur geloggt). crowdsecWL := crowdsec.NewWhitelistGenerator(pool) domainsReloader := func(ctx context.Context) error { return errors.Join(haproxy.New(pool).Render(ctx), crowdsecWL.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, domainsReloader).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 committed, 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@. 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) } // Öffentlicher WG-Endpoint-Host für Peer-Configs = FQDN dieser Node // (aus setup.json). Verhindert den REPLACE_WITH_PUBLIC_HOST-Platzhalter, // an dem Clients sonst keinen Tunnel aufbauen können. wgPublicHost := "" if sst, serr := setupStore.Load(); serr == nil && sst != nil { wgPublicHost = sst.FQDN } handlers.NewWireguardHandler(wgIfaces, wgPeers, secretsBox, auditRepo, nodeID, withFW(wgReloader)).WithPublicHost(wgPublicHost).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 triggering 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 triggering 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 triggering 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) // ReadHeaderTimeout kappt Slowloris-artige Header-Stalls (gosec G112). // ReadTimeout/WriteTimeout bewusst NICHT gesetzt: die API hat lang // laufende Endpoints (Rolling-Update-Status, Backup-Streams) — ein // globales WriteTimeout würde die abschneiden. IdleTimeout hält // Keep-Alive-Verbindungen in Grenzen. srv := &http.Server{ Addr: addr, Handler: r, ReadHeaderTimeout: 15 * time.Second, IdleTimeout: 120 * time.Second, } 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 = ` EdgeGuard

EdgeGuard

The management UI has not been built yet. From the project root, run:

cd management-ui && bun install && bun run build

Then the same URL will serve the React SPA. The REST API is fully functional at /api/v1/* regardless.

` // 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 // 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@ 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 { //nolint:contextcheck // detached by design — Heartbeat-Push nutzt eigenen Timeout, überlebt Request-Cancel 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 { //nolint:contextcheck // detached by design — Heartbeat-Push nutzt eigenen Timeout, überlebt Request-Cancel 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 }