cluster-init-replication: - listen_addresses = '*' damit Cluster-Peers PG auf :5432 erreichen können - max_replication_slots = 20 / max_wal_senders = 10 (verhindert Slot-Erschöpfung bei initaler Tabellen-Synchronisation mit vielen gleichzeitigen Sync-Workern) - pg-replication-secret: Ownership an edgeguard-User (API-Lesezugriff) - detectPGConfig() statt hardcoded PG 16 (System läuft PG 17) cluster-setup-standby: - syncMasterKey(): holt /var/lib/edgeguard/.master_key via mTLS vom Primary — ohne identischen Master-Key können replizierte WireGuard-Keys nicht entschlüsselt werden - render-config: sudo -u edgeguard statt als root (DB-Zugriff) nftables Template: - Port 5432 (PG) + 6379 (KeyDB) für Cluster-Peers (@peer_ipv4/@peer_ipv6) freigegeben handlers/cluster.go: - GET /agent/cluster/master-key: gibt .master_key via mTLS zurück (hex-kodiert) v1.2.15 Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
923 lines
31 KiB
Go
923 lines
31 KiB
Go
package handlers
|
|
|
|
import (
|
|
"context"
|
|
"crypto/x509"
|
|
"encoding/json"
|
|
"encoding/pem"
|
|
"fmt"
|
|
"log/slog"
|
|
"net/http"
|
|
"os"
|
|
"os/exec"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/gin-gonic/gin"
|
|
|
|
"git.netcell-it.de/projekte/edgeguard-native/internal/aggregator"
|
|
"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/handlers/response"
|
|
"git.netcell-it.de/projekte/edgeguard-native/internal/models"
|
|
)
|
|
|
|
// ClusterHandler exposes cluster-state endpoints. /status ist die
|
|
// strukturierte UI-Sicht (local + peers + health), /nodes ist die
|
|
// simple list, /system/load fan-outet via mTLS-Aggregator zu allen
|
|
// Peers und liefert pro Node die /proc-Metriken.
|
|
//
|
|
// Aggregator kann nil sein (clustertls nicht initialisiert) — dann
|
|
// liefert /cluster/system/load nur den lokalen Wert.
|
|
type ClusterHandler struct {
|
|
Store *cluster.Store
|
|
LocalID string
|
|
Aggregator *aggregator.Aggregator
|
|
Version string // laufende Binary-Version, für Rolling-Update-Koordination
|
|
|
|
// TLSStore + Tokens: optional, gesetzt bei Phase 3.4. Erlauben das
|
|
// Generieren von Join-Tokens und das Issue-Cert für joining Peers.
|
|
TLSStore *clustertls.Store
|
|
Tokens *jointoken.Service
|
|
|
|
// PeerReloader: optional, gesetzt bei Phase 3.5. Nach Auto-Register
|
|
// triggert das den firewall-Render damit peer_ipv4 frisch ist.
|
|
PeerReloader PeerReloader
|
|
}
|
|
|
|
func NewClusterHandler(store *cluster.Store, localID string) *ClusterHandler {
|
|
return &ClusterHandler{Store: store, LocalID: localID}
|
|
}
|
|
|
|
// WithAggregator: optionale Aggregator-Konfiguration. Nur wenn vorhanden
|
|
// wird /cluster/system/load die Peers via mTLS abklappern.
|
|
func (h *ClusterHandler) WithAggregator(a *aggregator.Aggregator) *ClusterHandler {
|
|
h.Aggregator = a
|
|
return h
|
|
}
|
|
|
|
// WithJoinFlow: Cluster-CA + Join-Token-Service. Nur wenn beide gesetzt
|
|
// sind exposen wir /cluster/join-tokens (admin) + /cluster/issue-cert (public).
|
|
func (h *ClusterHandler) WithJoinFlow(store *clustertls.Store, tokens *jointoken.Service) *ClusterHandler {
|
|
h.TLSStore = store
|
|
h.Tokens = tokens
|
|
return h
|
|
}
|
|
|
|
func (h *ClusterHandler) Register(rg *gin.RouterGroup) {
|
|
g := rg.Group("/cluster")
|
|
g.GET("/nodes", h.ListNodes)
|
|
g.GET("/status", h.Status)
|
|
g.GET("/system/load", h.SystemLoad)
|
|
g.DELETE("/nodes/:id", h.DeleteNode)
|
|
g.GET("/vip-settings", h.GetVIPSettings)
|
|
g.PUT("/vip-settings", h.UpdateVIPSettings)
|
|
g.POST("/rolling-update", h.RollingUpdate)
|
|
g.GET("/rolling-update/status", h.RollingUpdateStatus)
|
|
if h.TLSStore != nil {
|
|
g.GET("/cert-status", h.CertStatus)
|
|
g.POST("/renew-self", h.RenewSelf)
|
|
}
|
|
if h.TLSStore != nil && h.Tokens != nil {
|
|
g.POST("/join-tokens", h.GenerateJoinToken)
|
|
}
|
|
}
|
|
|
|
// DeleteNode entfernt einen Peer aus ha_nodes. Verweigert für die
|
|
// lokale Node (LocalID) — die kannst du nicht via UI löschen, sonst
|
|
// kommt der nächste Heartbeat-Tick die Row wieder anlegen oder
|
|
// die Cluster-Page wird inkonsistent.
|
|
//
|
|
// Nach erfolgreichem Delete triggert der PeerReloader (falls gesetzt)
|
|
// einen Firewall-Render — peer_ipv4-Set verliert die IP, der entfernte
|
|
// Peer kann nicht mehr auf :8443/:16379 connecten.
|
|
func (h *ClusterHandler) DeleteNode(c *gin.Context) {
|
|
id := c.Param("id")
|
|
if id == "" {
|
|
response.BadRequest(c, simpleError("missing id"))
|
|
return
|
|
}
|
|
if id == h.LocalID {
|
|
response.BadRequest(c, simpleError("cannot remove the local node — would auto-recreate on next heartbeat"))
|
|
return
|
|
}
|
|
if err := h.Store.Delete(c.Request.Context(), id); err != nil {
|
|
if err == cluster.ErrNotFound {
|
|
response.NotFound(c, err)
|
|
return
|
|
}
|
|
response.Internal(c, err)
|
|
return
|
|
}
|
|
if h.PeerReloader != nil {
|
|
go func() {
|
|
rctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer cancel()
|
|
if err := h.PeerReloader(rctx); err != nil {
|
|
slog.Warn("cluster: firewall render after peer delete failed", "error", err)
|
|
}
|
|
}()
|
|
}
|
|
slog.Info("cluster: peer removed", "id", id, "actor", actorOf(c))
|
|
response.NoContent(c)
|
|
}
|
|
|
|
// GetVIPSettings liest die cluster_settings-Singleton-Row (VIP/VRRP-Config).
|
|
func (h *ClusterHandler) GetVIPSettings(c *gin.Context) {
|
|
if h.Store == nil {
|
|
response.NotFound(c, simpleError("cluster store not available"))
|
|
return
|
|
}
|
|
var cs vipSettingsRow
|
|
row := h.Store.Pool.QueryRow(c.Request.Context(),
|
|
`SELECT vip_address, vip_interface, vip_auth_pass, vrrp_router_id FROM cluster_settings WHERE id = 1`)
|
|
if err := row.Scan(&cs.VIPAddress, &cs.VIPInterface, &cs.VIPAuthPass, &cs.VRRPRouterID); err != nil {
|
|
response.Internal(c, err)
|
|
return
|
|
}
|
|
response.OK(c, cs)
|
|
}
|
|
|
|
// UpdateVIPSettings speichert die VIP/VRRP-Konfiguration und triggert
|
|
// einen Keepalived-Config-Render. Viewer-Schutz via RequireAdminForMutations-
|
|
// Middleware auf der authed-Group — kein Extra-Check nötig.
|
|
func (h *ClusterHandler) UpdateVIPSettings(c *gin.Context) {
|
|
var req vipSettingsRow
|
|
if err := c.ShouldBindJSON(&req); err != nil {
|
|
response.BadRequest(c, err)
|
|
return
|
|
}
|
|
if h.Store == nil {
|
|
response.NotFound(c, simpleError("cluster store not available"))
|
|
return
|
|
}
|
|
_, err := h.Store.Pool.Exec(c.Request.Context(), `
|
|
UPDATE cluster_settings
|
|
SET vip_address=$1, vip_interface=$2, vip_auth_pass=$3, vrrp_router_id=$4, updated_at=NOW()
|
|
WHERE id=1`,
|
|
nullIfEmpty(req.VIPAddress), nullIfEmpty(req.VIPInterface),
|
|
nullIfEmpty(req.VIPAuthPass), req.VRRPRouterID)
|
|
if err != nil {
|
|
response.Internal(c, err)
|
|
return
|
|
}
|
|
slog.Info("cluster: VIP settings updated", "vip", req.VIPAddress, "actor", actorOf(c))
|
|
// Keepalived-Config asynchron neu rendern
|
|
if h.PeerReloader != nil {
|
|
go func() {
|
|
rctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
|
|
defer cancel()
|
|
if err := h.PeerReloader(rctx); err != nil {
|
|
slog.Warn("cluster: keepalived render after VIP update failed", "error", err)
|
|
}
|
|
}()
|
|
}
|
|
response.NoContent(c)
|
|
}
|
|
|
|
type vipSettingsRow struct {
|
|
VIPAddress *string `json:"vip_address"`
|
|
VIPInterface *string `json:"vip_interface"`
|
|
VIPAuthPass *string `json:"vip_auth_pass"`
|
|
VRRPRouterID int `json:"vrrp_router_id"`
|
|
}
|
|
|
|
func nullIfEmpty(s *string) *string {
|
|
if s == nil || *s == "" {
|
|
return nil
|
|
}
|
|
return s
|
|
}
|
|
|
|
// RegisterPublic mountet die public (unauth) Endpoints — joining Peers
|
|
// haben noch keine Session/Cert, deshalb läuft /issue-cert vor der
|
|
// requireAuth-Middleware. Aufrufer muss diesen Group auf /api/v1 setzen
|
|
// (NICHT auf authed).
|
|
func (h *ClusterHandler) RegisterPublic(rg *gin.RouterGroup) {
|
|
if h.TLSStore == nil || h.Tokens == nil {
|
|
return
|
|
}
|
|
g := rg.Group("/cluster")
|
|
g.POST("/issue-cert", h.IssueCert)
|
|
}
|
|
|
|
// RegisterAgent mountet die Peer-only-Endpoints auf dem mTLS-Agent-
|
|
// Listener. Auth läuft über das Peer-Cert (RequireAndVerifyClientCert
|
|
// im ServerTLSConfig); der CN des Cert ist die FQDN des Peers.
|
|
//
|
|
// /agent/cluster/peers — joining Peer trägt sich nach erfolgreichem
|
|
// issue-cert hier ein, damit der Primary ihn in ha_nodes hat (mit
|
|
// status='joining') und der nächste Firewall-Render-Lauf seine IP
|
|
// ins peer_ipv4-Set aufnimmt.
|
|
func (h *ClusterHandler) RegisterAgent(rg *gin.RouterGroup) {
|
|
g := rg.Group("/agent/cluster")
|
|
g.POST("/peers", h.AgentRegisterPeer)
|
|
g.GET("/identity", h.AgentIdentity)
|
|
g.GET("/pg-replication-info", h.AgentPGReplicationInfo)
|
|
g.GET("/master-key", h.AgentMasterKey)
|
|
g.GET("/version", h.AgentVersion)
|
|
g.POST("/trigger-update", h.AgentTriggerUpdate)
|
|
}
|
|
|
|
// AgentIdentity gibt die eigene ha_nodes-Row zurück. Wird vom Primary
|
|
// genutzt um joining-Peers aktiv zu reconcilen wenn autoRegister (Push)
|
|
// fehlgeschlagen ist — Pull-Fallback.
|
|
func (h *ClusterHandler) AgentIdentity(c *gin.Context) {
|
|
if h.Store == nil || h.LocalID == "" {
|
|
response.NotFound(c, simpleError("node not registered"))
|
|
return
|
|
}
|
|
node, err := h.Store.Get(c.Request.Context(), h.LocalID)
|
|
if err != nil {
|
|
if err == cluster.ErrNotFound {
|
|
response.NotFound(c, simpleError("local node not in ha_nodes yet"))
|
|
return
|
|
}
|
|
response.Internal(c, err)
|
|
return
|
|
}
|
|
response.OK(c, node)
|
|
}
|
|
|
|
// AgentPGReplicationInfo gibt die Replication-Credentials für pg_basebackup
|
|
// zurück. Nur über den mTLS-Agent-Listener erreichbar. Liest das Passwort
|
|
// aus /var/lib/edgeguard/pg-replication-secret. Gibt 404 zurück wenn die
|
|
// Datei fehlt (cluster-init-replication noch nicht ausgeführt).
|
|
func (h *ClusterHandler) AgentPGReplicationInfo(c *gin.Context) {
|
|
const secretPath = "/var/lib/edgeguard/pg-replication-secret"
|
|
pass, err := readFileString(secretPath)
|
|
if err != nil {
|
|
response.NotFound(c, simpleError("pg-replication-secret nicht gefunden — cluster-init-replication auf dem Primary ausführen"))
|
|
return
|
|
}
|
|
// Host = eigene Public-IP aus ha_nodes (oder Fallback: FQDN)
|
|
host := ""
|
|
if h.Store != nil && h.LocalID != "" {
|
|
if node, err := h.Store.Get(c.Request.Context(), h.LocalID); err == nil {
|
|
if node.PublicIP != nil && *node.PublicIP != "" {
|
|
host = *node.PublicIP
|
|
}
|
|
if host == "" {
|
|
host = node.FQDN
|
|
}
|
|
}
|
|
}
|
|
response.OK(c, gin.H{
|
|
"host": host,
|
|
"port": 5432,
|
|
"user": "edgeguard_replicator",
|
|
"password": strings.TrimSpace(pass),
|
|
})
|
|
}
|
|
|
|
// AgentMasterKey gibt den Secrets-Master-Key zurück, damit cluster-setup-standby
|
|
// ihn auf dem Secondary synchronisieren kann. Nur über den mTLS-Agent-Listener
|
|
// erreichbar. Ohne gemeinsamen Master-Key können replizierte verschlüsselte
|
|
// Felder (WireGuard private keys, PSKs) auf dem Secondary nicht entschlüsselt werden.
|
|
func (h *ClusterHandler) AgentMasterKey(c *gin.Context) {
|
|
const keyPath = "/var/lib/edgeguard/.master_key"
|
|
data, err := os.ReadFile(keyPath)
|
|
if err != nil {
|
|
response.NotFound(c, simpleError("master key nicht gefunden"))
|
|
return
|
|
}
|
|
response.OK(c, gin.H{"key_hex": fmt.Sprintf("%x", data)})
|
|
}
|
|
|
|
func readFileString(path string) (string, error) {
|
|
b, err := os.ReadFile(path)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
return string(b), nil
|
|
}
|
|
|
|
// PeerReloader: optionale Funktion die nach einem Auto-Register
|
|
// Firewall + ggfs. andere Configs regeneriert (damit peer_ipv4-Set
|
|
// frisch ist). Wird vom main.go gesetzt.
|
|
type PeerReloader func(ctx context.Context) error
|
|
|
|
// WithPeerReloader: nach jedem AgentRegisterPeer-Aufruf gefeuert.
|
|
func (h *ClusterHandler) WithPeerReloader(r PeerReloader) *ClusterHandler {
|
|
h.PeerReloader = r
|
|
return h
|
|
}
|
|
|
|
// WithVersion: setzt die laufende Binary-Version für Rolling-Update-Koordination.
|
|
func (h *ClusterHandler) WithVersion(v string) *ClusterHandler {
|
|
h.Version = v
|
|
return h
|
|
}
|
|
|
|
func (h *ClusterHandler) ListNodes(c *gin.Context) {
|
|
nodes, err := h.Store.List(c.Request.Context())
|
|
if err != nil {
|
|
response.Internal(c, err)
|
|
return
|
|
}
|
|
response.OK(c, gin.H{"nodes": nodes, "local_id": h.LocalID})
|
|
}
|
|
|
|
// ClusterStatus ist die UI-zentrierte Sicht: local-Node hervorgehoben,
|
|
// peers separat, mode + health-flag.
|
|
type ClusterStatus struct {
|
|
LocalID string `json:"local_id"`
|
|
LocalNode *models.HANode `json:"local_node,omitempty"`
|
|
Peers []models.HANode `json:"peers"`
|
|
Mode string `json:"mode"` // "single-node" | "cluster"
|
|
Health string `json:"health"` // "ok" | "degraded" | "split-brain"
|
|
DriftFound bool `json:"drift_found"`
|
|
UpdatedAt time.Time `json:"updated_at"`
|
|
}
|
|
|
|
// Status splittet alle Nodes in local + peers, berechnet mode + health.
|
|
// On-demand: bevor wir die Rows lesen refreshen wir den eigenen
|
|
// config_hash, sodass das UI immer aktuelle Werte sieht — auch wenn
|
|
// der 5min-Scheduler-Tick gerade vorher nicht gelaufen ist.
|
|
func (h *ClusterHandler) Status(c *gin.Context) {
|
|
if h.Store != nil && h.Store.Pool != nil && h.LocalID != "" {
|
|
// 2s Timeout — der Hash-Compute braucht im normal-case <50ms.
|
|
// Bei timeout fallen wir auf den (eventuell stale) DB-Wert zurück.
|
|
ctx, cancel := context.WithTimeout(c.Request.Context(), 2*time.Second)
|
|
if _, err := cluster.RefreshLocalHash(ctx, h.Store.Pool, h.LocalID); err != nil {
|
|
slog.Warn("cluster: config_hash refresh failed", "error", err)
|
|
}
|
|
cancel()
|
|
}
|
|
all, err := h.Store.List(c.Request.Context())
|
|
if err != nil {
|
|
response.Internal(c, err)
|
|
return
|
|
}
|
|
out := ClusterStatus{
|
|
LocalID: h.LocalID,
|
|
Peers: []models.HANode{},
|
|
Mode: "single-node",
|
|
Health: "ok",
|
|
UpdatedAt: time.Now().UTC(),
|
|
}
|
|
var localHash *string
|
|
for i := range all {
|
|
n := all[i]
|
|
if n.ID == h.LocalID {
|
|
ln := n
|
|
out.LocalNode = &ln
|
|
localHash = ln.ConfigHash
|
|
continue
|
|
}
|
|
out.Peers = append(out.Peers, n)
|
|
}
|
|
if len(out.Peers) > 0 {
|
|
out.Mode = "cluster"
|
|
}
|
|
|
|
// Pull-Reconcile für joining-Peers: wenn ein Peer via Aggregator
|
|
// erreichbar ist aber noch mit Placeholder-ID in ha_nodes steht,
|
|
// holen wir seine echte Identity aktiv ab. Best-effort goroutine —
|
|
// blockiert die Status-Response nicht.
|
|
if h.Aggregator != nil && h.Store != nil {
|
|
var joining []models.HANode
|
|
for _, p := range out.Peers {
|
|
if p.Status == "joining" || p.Status == "pending" {
|
|
joining = append(joining, p)
|
|
}
|
|
}
|
|
if len(joining) > 0 {
|
|
go h.reconcileJoiningPeers(joining)
|
|
}
|
|
}
|
|
|
|
// Drift-Detection: jeder peer mit anderem config_hash als unser
|
|
// lokaler → Banner-Trigger im UI.
|
|
if localHash != nil && *localHash != "" {
|
|
for _, p := range out.Peers {
|
|
if p.ConfigHash == nil || *p.ConfigHash == "" {
|
|
continue
|
|
}
|
|
if *p.ConfigHash != *localHash {
|
|
out.DriftFound = true
|
|
out.Health = "degraded"
|
|
break
|
|
}
|
|
}
|
|
}
|
|
// Offline-Peers → degraded.
|
|
if !out.DriftFound {
|
|
for _, p := range out.Peers {
|
|
if p.Status != "online" {
|
|
out.Health = "degraded"
|
|
break
|
|
}
|
|
}
|
|
}
|
|
response.OK(c, out)
|
|
}
|
|
|
|
// SystemLoad aggregiert /proc-Metriken aller Peers via mTLS-Aggregator.
|
|
// Liefert ein Array { node_id, fqdn, ok, data, error, duration_ms }.
|
|
// Lokaler Node wird IMMER eingefügt (direkter Call statt mTLS-Roundtrip).
|
|
func (h *ClusterHandler) SystemLoad(c *gin.Context) {
|
|
all, err := h.Store.List(c.Request.Context())
|
|
if err != nil {
|
|
response.Internal(c, err)
|
|
return
|
|
}
|
|
results := make([]aggregator.PeerResult, 0, len(all))
|
|
// Lokalen Node selbst befragen: wir rufen die Resources-Funktion
|
|
// inline statt einen mTLS-Loopback aufzubauen — auch wenn der
|
|
// Agent-Listener läuft, ist ein direkter Call billiger.
|
|
for _, n := range all {
|
|
if n.ID == h.LocalID {
|
|
local := localSystemLoad()
|
|
raw, _ := json.Marshal(local)
|
|
results = append(results, aggregator.PeerResult{
|
|
NodeID: n.ID,
|
|
FQDN: n.FQDN,
|
|
OK: true,
|
|
Data: raw,
|
|
Duration: 0,
|
|
})
|
|
}
|
|
}
|
|
if h.Aggregator != nil {
|
|
peers := make([]models.HANode, 0, len(all))
|
|
for _, n := range all {
|
|
if n.ID == h.LocalID {
|
|
continue
|
|
}
|
|
peers = append(peers, n)
|
|
}
|
|
// /agent/system/resources auf dem Peer-Agent-Listener (:8443 mTLS).
|
|
fan := h.Aggregator.FanOut(c.Request.Context(), peers, "/agent/system/resources", h.LocalID)
|
|
results = append(results, fan...)
|
|
}
|
|
response.OK(c, gin.H{"nodes": results})
|
|
}
|
|
|
|
// localSystemLoad ruft die selben Werte wie /system/resources, aber
|
|
// als bare struct (kein gin.Context). Damit liefert SystemLoad pro Node
|
|
// dasselbe Format wie der Agent-Endpoint.
|
|
func localSystemLoad() any {
|
|
// SystemHandler.Resources nutzt ein internes `resources` struct.
|
|
// Wir duplizieren die /proc-Reads nicht — der Agent-Listener mountet
|
|
// denselben Handler und liefert die JSON-Struktur. Für den lokalen
|
|
// Path liefern wir das Snapshot über computeLocalSystemResources.
|
|
return computeLocalSystemResources()
|
|
}
|
|
|
|
// ── Phase 3.4: Cluster-Join Token Flow ────────────────────────────────
|
|
|
|
// GenerateJoinToken — Admin generiert einen one-shot Bootstrap-Token
|
|
// für einen neuen Peer. Der FQDN des neuen Nodes wird vorab übergeben
|
|
// so dass er sofort in ha_nodes vorregistriert und im UI angezeigt werden
|
|
// kann. Token wird NUR EINMAL zurückgegeben.
|
|
func (h *ClusterHandler) GenerateJoinToken(c *gin.Context) {
|
|
var req struct {
|
|
NodeFQDN string `json:"node_fqdn"`
|
|
}
|
|
// Body optional — wenn leer, läuft der Flow ohne Pre-Register.
|
|
_ = c.ShouldBindJSON(&req)
|
|
|
|
token, exp, err := h.Tokens.Generate()
|
|
if err != nil {
|
|
response.Internal(c, err)
|
|
return
|
|
}
|
|
caCert, _, err := h.TLSStore.LoadCA()
|
|
caFP := ""
|
|
if err == nil && caCert != nil {
|
|
caFP = jointoken.CAFingerprint16(caCert.Raw)
|
|
}
|
|
|
|
// Pre-register: wenn ein FQDN übergeben wurde, schon jetzt in
|
|
// ha_nodes anlegen (status='pending') + Firewall-Reload. So ist
|
|
// der Node bekannt bevor er überhaupt antwortet, und @peer_ipv4
|
|
// wird beim issue-cert (wo wir die IP haben) nur noch updaten.
|
|
if req.NodeFQDN != "" && h.Store != nil {
|
|
go h.preRegisterByFQDN(req.NodeFQDN)
|
|
}
|
|
|
|
response.OK(c, gin.H{
|
|
"token": token,
|
|
"expires_at": exp.UTC().Format(time.RFC3339),
|
|
"ca_fingerprint": caFP,
|
|
"node_fqdn": req.NodeFQDN,
|
|
})
|
|
}
|
|
|
|
// preRegisterByFQDN legt einen ha_nodes-Eintrag mit status='pending' an
|
|
// bevor der Joiner überhaupt die Verbindung aufbaut. Idempotent dank
|
|
// ON CONFLICT. Kein Firewall-Reload hier — die IP ist noch unbekannt;
|
|
// das erledigt preRegisterJoiner wenn die issue-cert-Anfrage eintrifft.
|
|
func (h *ClusterHandler) preRegisterByFQDN(fqdn string) {
|
|
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer cancel()
|
|
nodeID := fmt.Sprintf("prenode-%s", strings.ReplaceAll(fqdn, ".", "-"))
|
|
n := models.HANode{
|
|
ID: nodeID,
|
|
Name: fqdn,
|
|
FQDN: fqdn,
|
|
APIURL: "https://" + fqdn + ":3443",
|
|
Role: "peer",
|
|
Status: "pending",
|
|
}
|
|
if _, err := h.Store.UpsertSelf(ctx, n); err != nil {
|
|
slog.Warn("cluster: pre-register by FQDN failed", "fqdn", fqdn, "error", err)
|
|
}
|
|
}
|
|
|
|
// IssueCert — Joining Peer POSTet seinen CSR + den Token. Wir verifizieren
|
|
// + konsumieren den Token, signieren den CSR mit unserer Cluster-CA und
|
|
// liefern {ca_cert, peer_cert} zurück. PUBLIC Endpoint — keine Session-
|
|
// Auth nötig (der Joiner hat noch keine).
|
|
type issueCertRequest struct {
|
|
Token string `json:"token"`
|
|
CSR string `json:"csr"`
|
|
}
|
|
type issueCertResponse struct {
|
|
CACert string `json:"ca_cert"`
|
|
PeerCert string `json:"peer_cert"`
|
|
}
|
|
|
|
func (h *ClusterHandler) IssueCert(c *gin.Context) {
|
|
var req issueCertRequest
|
|
if err := c.ShouldBindJSON(&req); err != nil {
|
|
response.BadRequest(c, err)
|
|
return
|
|
}
|
|
if req.Token == "" || req.CSR == "" {
|
|
response.BadRequest(c, errInvalidJoinRequest)
|
|
return
|
|
}
|
|
// consumedBy → Remote-IP. Audit-Trail wenn jemand Tokens stiehlt
|
|
// und vom falschen Host einlöst.
|
|
clientIP := c.ClientIP()
|
|
if _, err := h.Tokens.Consume(c.Request.Context(), req.Token, clientIP); err != nil {
|
|
response.BadRequest(c, err)
|
|
return
|
|
}
|
|
// CSR signieren.
|
|
peerCert, err := h.TLSStore.SignCSR(req.CSR, nil)
|
|
if err != nil {
|
|
response.BadRequest(c, err)
|
|
return
|
|
}
|
|
caPEM, err := h.TLSStore.CACertPEM()
|
|
if err != nil {
|
|
response.Internal(c, err)
|
|
return
|
|
}
|
|
|
|
// Pre-register the joining node SYNCHRONOUSLY before returning the
|
|
// cert so that nftables @peer_ipv4 already contains the joiner's IP
|
|
// by the time they call autoRegister on port 8443. A goroutine here
|
|
// caused a race: cert returned → joiner calls autoRegister → nftables
|
|
// not updated yet → connection refused → status stays "joining".
|
|
if h.Store != nil && h.PeerReloader != nil {
|
|
h.preRegisterJoiner(c.Request.Context(), clientIP, req.CSR)
|
|
}
|
|
|
|
response.OK(c, issueCertResponse{
|
|
CACert: caPEM,
|
|
PeerCert: peerCert,
|
|
})
|
|
}
|
|
|
|
// preRegisterJoiner inserts a minimal ha_nodes row for the joining peer
|
|
// (using the CSR CN as FQDN and the HTTP client IP as public_ip), then
|
|
// triggers a firewall reload so @peer_ipv4 contains the new IP before
|
|
// the peer tries to call /agent/cluster/peers on port 8443.
|
|
// Uses a stable deterministic ID so re-joins are idempotent.
|
|
func (h *ClusterHandler) preRegisterJoiner(parent context.Context, clientIP, csrPEM string) {
|
|
ctx, cancel := context.WithTimeout(parent, 10*time.Second)
|
|
defer cancel()
|
|
|
|
fqdn := cnFromCSR(csrPEM)
|
|
if fqdn == "" {
|
|
fqdn = "joining-" + clientIP
|
|
}
|
|
nodeID := fmt.Sprintf("prenode-%s", strings.ReplaceAll(fqdn, ".", "-"))
|
|
|
|
n := models.HANode{
|
|
ID: nodeID,
|
|
Name: fqdn,
|
|
FQDN: fqdn,
|
|
APIURL: "https://" + fqdn + ":3443",
|
|
Role: "peer",
|
|
Status: "joining",
|
|
}
|
|
n.PublicIP = &clientIP
|
|
|
|
if _, err := h.Store.UpsertSelf(ctx, n); err != nil {
|
|
slog.Warn("cluster: pre-register joiner failed", "fqdn", fqdn, "ip", clientIP, "error", err)
|
|
return
|
|
}
|
|
if err := h.PeerReloader(ctx); err != nil {
|
|
slog.Warn("cluster: PeerReloader failed after pre-register", "error", err)
|
|
return
|
|
}
|
|
slog.Info("cluster: joiner pre-registered, firewall updated", "fqdn", fqdn, "ip", clientIP)
|
|
}
|
|
|
|
// cnFromCSR extracts the Subject Common Name from a PEM-encoded CSR.
|
|
// Returns empty string on any parse error.
|
|
func cnFromCSR(csrPEM string) string {
|
|
block, _ := pem.Decode([]byte(csrPEM))
|
|
if block == nil {
|
|
return ""
|
|
}
|
|
csr, err := x509.ParseCertificateRequest(block.Bytes)
|
|
if err != nil {
|
|
return ""
|
|
}
|
|
return csr.Subject.CommonName
|
|
}
|
|
|
|
// reconcileJoiningPeers versucht für jeden Peer im Status "joining" oder
|
|
// "pending" die echte Node-ID via /agent/cluster/identity zu holen und
|
|
// ihn in ha_nodes mit der richtigen ID einzutragen. Self-Healing-
|
|
// Fallback wenn autoRegister (Push von joining-Peer zu Primary) wegen
|
|
// eines temporären Netzwerkproblems fehlgeschlagen ist.
|
|
//
|
|
// Die public_ip des Placeholder-Rows wird in der neuen Row übernommen
|
|
// damit @peer_ipv4 (nftables) korrekt bleibt.
|
|
func (h *ClusterHandler) reconcileJoiningPeers(placeholders []models.HANode) {
|
|
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
|
|
defer cancel()
|
|
|
|
results := h.Aggregator.FanOut(ctx, placeholders, "/agent/cluster/identity", h.LocalID)
|
|
changed := false
|
|
for i, res := range results {
|
|
if !res.OK || len(res.Data) == 0 {
|
|
continue
|
|
}
|
|
var identity models.HANode
|
|
if err := json.Unmarshal(res.Data, &identity); err != nil || identity.ID == "" {
|
|
continue
|
|
}
|
|
placeholder := placeholders[i]
|
|
if identity.ID == placeholder.ID {
|
|
continue // ID bereits korrekt
|
|
}
|
|
// Echte ID gefunden — row mit realer ID anlegen, public_ip aus
|
|
// dem Placeholder-Row übernehmen damit nftables korrekt bleibt.
|
|
n := identity
|
|
n.Status = "online"
|
|
if n.PublicIP == nil {
|
|
n.PublicIP = placeholder.PublicIP
|
|
}
|
|
if n.InternalIP == nil {
|
|
n.InternalIP = placeholder.InternalIP
|
|
}
|
|
// Placeholder zuerst löschen: ha_nodes hat UNIQUE(fqdn). Ohne
|
|
// dieses Delete würde UpsertSelf (ON CONFLICT(id)) mit fqdn-
|
|
// unique-Violation scheitern.
|
|
_ = h.Store.DeletePlaceholdersByFQDN(ctx, n.FQDN, n.ID)
|
|
out, err := h.Store.UpsertSelf(ctx, n)
|
|
if err != nil {
|
|
slog.Warn("cluster: reconcile joining peer: upsert failed",
|
|
"fqdn", n.FQDN, "real_id", n.ID, "error", err)
|
|
continue
|
|
}
|
|
changed = true
|
|
slog.Info("cluster: joining peer reconciled via identity pull",
|
|
"id", out.ID, "fqdn", out.FQDN, "placeholder_id", placeholder.ID)
|
|
}
|
|
if changed && h.PeerReloader != nil {
|
|
rctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer cancel()
|
|
if err := h.PeerReloader(rctx); err != nil {
|
|
slog.Warn("cluster: PeerReloader failed after reconcile", "error", err)
|
|
}
|
|
}
|
|
}
|
|
|
|
// AgentVersion gibt die laufende Binary-Version zurück. Wird vom Rolling-
|
|
// Update-Orchestrator gepollt um zu erkennen wann der Secondary die neue
|
|
// Version hat.
|
|
func (h *ClusterHandler) AgentVersion(c *gin.Context) {
|
|
response.OK(c, gin.H{"version": h.Version})
|
|
}
|
|
|
|
// AgentTriggerUpdate startet den Upgrade-Prozess auf diesem Node via
|
|
// systemd-run (detached). Wird vom Primary via mTLS aufgerufen um den
|
|
// Secondary zuerst zu aktualisieren (Rolling-Update). Pattern identisch
|
|
// zu /system/upgrade — nutzt dieselbe Sudoers-Whitelist aus dem postinst.
|
|
func (h *ClusterHandler) AgentTriggerUpdate(c *gin.Context) {
|
|
const scriptPath = "/var/lib/edgeguard/upgrade.sh"
|
|
const script = `#!/bin/bash
|
|
set -e
|
|
sleep 2
|
|
export DEBIAN_FRONTEND=noninteractive
|
|
dpkg --configure -a || true
|
|
retry_apt() {
|
|
local attempt=0 max=3 wait_for=15
|
|
while [ $attempt -lt $max ]; do
|
|
attempt=$((attempt + 1))
|
|
apt-get update -qq || true
|
|
if apt-get install -y -qq -o Dpkg::Options::=--force-confold \
|
|
edgeguard-api edgeguard-ui edgeguard; then return 0; fi
|
|
[ $attempt -lt $max ] && sleep $wait_for && wait_for=$((wait_for * 2))
|
|
done
|
|
return 1
|
|
}
|
|
retry_apt
|
|
echo "[upgrade] complete"
|
|
rm -f /var/lib/edgeguard/upgrade.sh
|
|
`
|
|
if err := os.WriteFile(scriptPath, []byte(script), 0o755); err != nil {
|
|
response.Internal(c, err)
|
|
return
|
|
}
|
|
const unitName = "edgeguard-upgrade.service"
|
|
_ = exec.Command("sudo", "-n", "/usr/bin/systemctl", "reset-failed", unitName).Run()
|
|
cmd := exec.Command("sudo", "-n", "/usr/bin/systemd-run",
|
|
"--unit="+unitName,
|
|
"--description=EdgeGuard self-upgrade",
|
|
"--collect",
|
|
"bash", scriptPath)
|
|
if err := cmd.Run(); err != nil {
|
|
slog.Warn("cluster: AgentTriggerUpdate: systemd-run failed", "error", err)
|
|
response.Internal(c, err)
|
|
return
|
|
}
|
|
slog.Info("cluster: rolling update triggered on this node by primary mTLS call",
|
|
"client", c.ClientIP())
|
|
c.JSON(http.StatusAccepted, gin.H{"status": "upgrading"})
|
|
}
|
|
|
|
var errInvalidJoinRequest = simpleError("missing token or csr")
|
|
|
|
type simpleError string
|
|
|
|
func (e simpleError) Error() string { return string(e) }
|
|
|
|
// ── Cluster-Cert-Status + Renewal ─────────────────────────────────────
|
|
|
|
// CertStatus liefert Metadata zu CA + Peer-Cert (Common Name, Expiry,
|
|
// days_remaining). UI nutzt das für Expiry-Warnungen.
|
|
func (h *ClusterHandler) CertStatus(c *gin.Context) {
|
|
out := gin.H{
|
|
"has_ca": h.TLSStore.HasCA(),
|
|
"has_peer": h.TLSStore.HasPeer(),
|
|
}
|
|
if h.TLSStore.HasCA() {
|
|
if info, err := h.TLSStore.CACertInfo(); err == nil {
|
|
out["ca"] = info
|
|
}
|
|
}
|
|
if h.TLSStore.HasPeer() {
|
|
if info, err := h.TLSStore.PeerCertInfo(); err == nil {
|
|
out["peer"] = info
|
|
}
|
|
}
|
|
response.OK(c, out)
|
|
}
|
|
|
|
// RenewSelf re-signed das eigene peer.crt mit der eigenen CA. Nur
|
|
// sinnvoll auf einer Founder-Box; Joiner haben keine eigene CA.
|
|
//
|
|
// Nach Renew muss edgeguard-api restartet werden damit der Agent-
|
|
// Listener das neue Cert in seinen TLS-Config-Snapshot lädt — wir
|
|
// triggern das NICHT automatisch (würde die HTTP-Response abreißen);
|
|
// stattdessen liefern wir einen Hinweis im Response.
|
|
func (h *ClusterHandler) RenewSelf(c *gin.Context) {
|
|
if !h.TLSStore.HasCA() {
|
|
response.BadRequest(c, simpleError("no local CA — joiners cannot self-renew"))
|
|
return
|
|
}
|
|
// Common-Name aus dem existierenden Peer-Cert übernehmen damit der
|
|
// Cert weiterhin auf die aktuelle FQDN passt.
|
|
cn := "edgeguard-node"
|
|
if info, err := h.TLSStore.PeerCertInfo(); err == nil && info.CommonName != "" {
|
|
cn = info.CommonName
|
|
}
|
|
if err := h.TLSStore.RenewSelfSigned(cn, []string{cn}, nil, nil); err != nil {
|
|
response.Internal(c, err)
|
|
return
|
|
}
|
|
info, err := h.TLSStore.PeerCertInfo()
|
|
if err != nil {
|
|
response.Internal(c, err)
|
|
return
|
|
}
|
|
response.OK(c, gin.H{
|
|
"peer": info,
|
|
"restart_hint": "systemctl restart edgeguard-api",
|
|
})
|
|
}
|
|
|
|
// ── Phase 3.5: Auto-Register beim Cluster-Join ────────────────────────
|
|
|
|
// registerPeerRequest: vom Joiner via mTLS-POST an /agent/cluster/peers.
|
|
// CN des Client-Cert authentifiziert den Peer. Wir nehmen nur die Felder
|
|
// die wir wirklich brauchen — sonst kann ein joining Peer beliebige
|
|
// ha_nodes-Felder überschreiben.
|
|
type registerPeerRequest struct {
|
|
ID string `json:"id"` // Joiner's eigene node-id
|
|
Name string `json:"name"` // hostname
|
|
FQDN string `json:"fqdn"` // sollte mit Client-Cert-CN matchen
|
|
APIURL string `json:"api_url"` // https://<fqdn>
|
|
PublicIP string `json:"public_ip"` // optional
|
|
InternalIP string `json:"internal_ip"` // mTLS-Listener-IP (für peer_ipv4-Set)
|
|
MgmtIP string `json:"mgmt_ip"` // optional
|
|
Version string `json:"version"`
|
|
ConfigHash *string `json:"config_hash"` // nil=absent (don't change), ""=no user config
|
|
}
|
|
|
|
// AgentRegisterPeer: vom Joiner nach issue-cert via mTLS aufgerufen.
|
|
// Validiert dass der Client-Cert-CN zur fqdn passt (verhindert Cross-
|
|
// Peer-Hijack) und upsertet die Row in ha_nodes mit status='joining'.
|
|
// Phase 3.2 Heartbeat wird die Status auf 'online' ändern sobald der
|
|
// Joiner seinen eigenen Heartbeat-Tick startet.
|
|
func (h *ClusterHandler) AgentRegisterPeer(c *gin.Context) {
|
|
var req registerPeerRequest
|
|
if err := c.ShouldBindJSON(&req); err != nil {
|
|
response.BadRequest(c, err)
|
|
return
|
|
}
|
|
if req.ID == "" || req.FQDN == "" {
|
|
response.BadRequest(c, simpleError("id + fqdn required"))
|
|
return
|
|
}
|
|
// Cert-CN-Check: TLS-Layer hat den Cert bereits gegen unsere CA
|
|
// verifiziert; jetzt prüfen wir dass der CN zur claimed FQDN passt.
|
|
// Sonst könnte ein peer1.example.com Cert nutzen um peer2.example.com
|
|
// in ha_nodes zu schreiben.
|
|
if c.Request.TLS == nil || len(c.Request.TLS.PeerCertificates) == 0 {
|
|
response.Forbidden(c, simpleError("no client cert presented"))
|
|
return
|
|
}
|
|
cn := c.Request.TLS.PeerCertificates[0].Subject.CommonName
|
|
if cn != req.FQDN {
|
|
slog.Warn("cluster: peer cert CN does not match registration FQDN",
|
|
"cn", cn, "fqdn", req.FQDN)
|
|
response.Forbidden(c, simpleError("cert CN does not match fqdn"))
|
|
return
|
|
}
|
|
|
|
// Wir bauen ein models.HANode zusammen + nutzen den existierenden
|
|
// UpsertSelf. (UpsertSelf nimmt eine HANode für einen registrierenden
|
|
// Node, hier ist der „Self" der joining-Peer auf dieser Primary-Seite.
|
|
// Der Name passt nicht 100% semantisch, aber das SQL ist exakt das was
|
|
// wir brauchen.)
|
|
n := models.HANode{
|
|
ID: req.ID,
|
|
Name: req.Name,
|
|
FQDN: req.FQDN,
|
|
APIURL: req.APIURL,
|
|
Role: "peer",
|
|
Status: "online", // peer IS online — it just connected via mTLS
|
|
}
|
|
if req.PublicIP != "" {
|
|
v := req.PublicIP
|
|
n.PublicIP = &v
|
|
}
|
|
if req.InternalIP != "" {
|
|
v := req.InternalIP
|
|
n.InternalIP = &v
|
|
}
|
|
if req.MgmtIP != "" {
|
|
v := req.MgmtIP
|
|
n.MgmtIP = &v
|
|
}
|
|
if req.Version != "" {
|
|
v := req.Version
|
|
n.Version = &v
|
|
}
|
|
if req.ConfigHash != nil {
|
|
n.ConfigHash = req.ConfigHash
|
|
}
|
|
// Placeholder zuerst löschen: ha_nodes hat UNIQUE(fqdn). Der INSERT
|
|
// in UpsertSelf verwendet ON CONFLICT(id) — greift NICHT bei fqdn-
|
|
// Konflikten. Ohne das Delete würde der INSERT mit "duplicate key on
|
|
// ha_nodes_fqdn_unique" scheitern und der Peer bliebe ewig "joining".
|
|
_ = h.Store.DeletePlaceholdersByFQDN(c.Request.Context(), req.FQDN, req.ID)
|
|
|
|
out, err := h.Store.UpsertSelf(c.Request.Context(), n)
|
|
if err != nil {
|
|
response.Internal(c, err)
|
|
return
|
|
}
|
|
|
|
// Firewall-Reload damit peer_ipv4-Set die neue IP aufnimmt. Best-
|
|
// effort: Fehler loggen, Response weiter durchreichen — der Peer
|
|
// hat seine Identity erfolgreich registriert, Operator kann manuell
|
|
// nachrendern.
|
|
if h.PeerReloader != nil {
|
|
go func() {
|
|
rctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer cancel()
|
|
if err := h.PeerReloader(rctx); err != nil {
|
|
slog.Warn("cluster: PeerReloader failed after AgentRegisterPeer", "error", err)
|
|
}
|
|
}()
|
|
}
|
|
|
|
slog.Info("cluster: peer registered via mTLS",
|
|
"id", out.ID, "fqdn", out.FQDN, "role", out.Role, "status", out.Status,
|
|
"client_cn", cn, "remote", c.ClientIP())
|
|
response.OK(c, out)
|
|
}
|