package handlers import ( "context" "encoding/json" "log/slog" "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. Token wird NUR EINMAL zurückgegeben; Server // speichert keinen Klartext, beim Re-Use blockt der nonce-Tracker. func (h *ClusterHandler) GenerateJoinToken(c *gin.Context) { token, exp, err := h.Tokens.Generate() if err != nil { response.Internal(c, err) return } // Wir liefern auch die primary-fqdn + ca-fingerprint mit, damit // das UI den Join-Befehl als kompletten curl/CLI-String anzeigen // kann. caCert, _, err := h.TLSStore.LoadCA() caFP := "" if err == nil && caCert != nil { caFP = jointoken.CAFingerprint16(caCert.Raw) } response.OK(c, gin.H{ "token": token, "expires_at": exp.UTC().Format(time.RFC3339), "ca_fingerprint": caFP, }) } // 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. consumedBy := c.ClientIP() if _, err := h.Tokens.Consume(c.Request.Context(), req.Token, consumedBy); 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 } response.OK(c, issueCertResponse{ CACert: caPEM, PeerCert: peerCert, }) } 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:// 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: "joining", } 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 } // 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) }