Files
edgeguard-native/internal/handlers/cluster.go
Debian 357113e7be fix(cluster): Alle Placeholder-Rows bei autoRegister bereinigen
AgentRegisterPeer löschte bisher nur den neuen prenode-{fqdn}-Placeholder.
Alte pre-{timestamp}-Rows (aus Versionen vor 1.1.158) blieben stehen und
zeigten dauerhaft status=joining.

Fix: DeletePlaceholdersByFQDN löscht ALLE ha_nodes-Rows mit gleicher FQDN
außer der echten Node-ID — unabhängig vom ID-Format.

Auch preRegisterByFQDN nutzt jetzt das stabile prenode-{fqdn}-Format.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-29 18:33:50 +02:00

628 lines
21 KiB
Go

package handlers
import (
"context"
"crypto/x509"
"encoding/json"
"encoding/pem"
"fmt"
"log/slog"
"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
// 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)
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)
}
// 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)
}
// 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
}
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"
}
// 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
}
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"`
}
// 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
}
out, err := h.Store.UpsertSelf(c.Request.Context(), n)
if err != nil {
response.Internal(c, err)
return
}
// Remove ALL placeholder rows for this FQDN (both legacy "pre-{timestamp}"
// and current "prenode-{fqdn}" style) — the real row just took their place.
_ = h.Store.DeletePlaceholdersByFQDN(c.Request.Context(), req.FQDN, req.ID)
// 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)
}