Fix 1 — Peer zeigt fälschlich 'offline': runPrimaryPush (Secondary→Primary, einziger periodischer Cross-Node-ha_nodes-Refresh) tickte mit 5 min, SweepStaleNodes-Threshold ist aber 2 min → Secondary war 2 min online, dann 3 min offline, im 5-min-Takt. Tick auf 30s (4× Marge unter Threshold). Receiver lädt nftables nur bei IP-Änderung → kein Reload-Sturm. Fix 2 — Rolling-Update konnte nie fertig werden wenn der Secondary die Zielversion schon hatte (baseline==target → Warten auf unmöglichen Flip → 10-min-Timeout). runRollingUpdate ist jetzt candidate-aware: ermittelt apt-Candidate, überspringt den Secondary-Schritt wenn dieser schon aktuell ist, erkennt den Flip via 'erreicht candidate ODER bewegt sich von baseline', und schließt direkt mit 'done' wenn auch der Primary schon aktuell ist. FinishRollingUpdateIfPending setzt hängende updating/waiting-secondary-Phasen beim Boot auf idle zurück (tote Orchestrierungs-Goroutine nach Restart). Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
875 lines
36 KiB
Go
875 lines
36 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).
|
||
}
|
||
|
||
// 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()
|
||
}
|
||
}
|
||
}
|
||
|
||
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
|
||
}
|