Files
edgeguard-native/internal/handlers/cluster.go
noroot 808f6fc055 feat(cluster): Logical Replication wird beim Join automatisch eingerichtet
Bisher war der zweite Node nach dem Join zwar im Cluster registriert und
in der UI sichtbar, replizierte aber keine einzige geteilte Tabelle —
dafuer musste jemand manuell `edgeguard-ctl cluster-setup-standby`
ausfuehren. Wer das uebersah, merkte es erst beim Failover: der neue
Primary stand ohne Domains, Backends, Firewall-Regeln und WireGuard-Keys
da. Das ist jetzt Teil des Join-Vorgangs.

Beide Seiten muessen dafuer vorbereitet sein:

1) Primary, beim Erzeugen des Join-Tokens: ein frisch installierter
   Single-Node hat weder Replikations-Rolle noch PUBLICATION noch
   wal_level=logical. Ohne das liefe das spaetere CREATE SUBSCRIPTION in
   ein 404. Der Token wird deshalb erst ausgegeben, nachdem die
   Publisher-Seite steht — inklusive des einmaligen PG-Restarts
   (wal_level ist ein postmaster-Parameter), der bewusst hier passiert,
   solange der Admin danebensteht und noch kein Peer Traffic erwartet.

   WICHTIG dabei: setupReplicationPrimary rotiert bei jedem Lauf das
   Replikations-Passwort (ALTER ROLE … PASSWORD). Auf einem Cluster mit
   bereits angebundenem Subscriber wuerde ein zweiter Token-Klick dessen
   Connection-String ungueltig machen und die Replikation still
   anhalten. Deshalb laeuft die Initialisierung nur, wenn PUBLICATION
   und Secret nicht bereits existieren.

2) Neuer Node, nach erfolgreichem Join: cluster-setup-standby laeuft
   detached (die Initialkopie dauert je nach Datenmenge Minuten), der
   Wizard pollt GET /setup/replication-status und zeigt running/done/
   failed an. Schlaegt es fehl, steht das manuelle Kommando inkl.
   Primary-Host direkt daneben statt nur einer Fehlermeldung.

Beides braucht root (psql als postgres, pg_hba, PG-Restart), die API
laeuft als unprivilegierter edgeguard → Aufruf via sudo mit gepinnten
Regeln. Das einzige variable Argument (Primary-Host) wird vorher gegen
Hostname/IP-Syntax geprueft; der Aufruf laeuft ohne Shell. Test dafuer
liegt bei.

Ausserdem zwei Doku-Korrekturen: architecture.md behauptete,
cluster-join richte die Replikation gleich mit ein (tut es nicht,
clusterjoin.Join macht nur Cert + Registrierung), und der Hinweistext
von cluster-join verwies noch auf "PG-Basebackup + KeyDB, Phase 3.5" —
beides laut Doku laengst verworfen.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-11 11:40:44 +02:00

1192 lines
42 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"
aptsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/apt"
"git.netcell-it.de/projekte/edgeguard-native/internal/services/audit"
)
// 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-Coordination
// 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
// Audit + NodeID: optional, gesetzt via WithAudit. Nötig für
// protokollierte, mutierende Aktionen wie den Replication-Repair.
Audit *audit.Repo
NodeID string
}
const (
// pgPublicationName + pgReplicationSecretPath spiegeln die Werte aus
// cmd/edgeguard-ctl (egPubName / egReplSecret) — beide Seiten muessen
// dasselbe meinen.
pgPublicationName = "edgeguard_shared"
pgReplicationSecretPath = "/var/lib/edgeguard/pg-replication-secret"
)
func NewClusterHandler(store *cluster.Store, localID string) *ClusterHandler {
return &ClusterHandler{Store: store, LocalID: localID}
}
// WithAggregator: optionale Aggregator-Configuration. 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
}
// WithAudit setzt den Audit-Repo + NodeID für protokollierte Aktionen.
func (h *ClusterHandler) WithAudit(a *audit.Repo, nodeID string) *ClusterHandler {
h.Audit = a
h.NodeID = nodeID
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)
g.POST("/repair-replication", h.RepairReplication)
g.GET("/repair-replication/status", h.RepairReplicationStatus)
g.GET("/vip-status", h.VIPStatus)
g.POST("/vip-test", h.VIPTest)
g.GET("/update-channel", h.UpdateChannel)
g.POST("/update-channel", h.SetUpdateChannel)
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 inconsistent.
//
// 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,
hb_interface, hb_src_ip, hb_peer_ip, hb_router_id, gw_check_ip
FROM cluster_settings WHERE id = 1`)
if err := row.Scan(&cs.VIPAddress, &cs.VIPInterface, &cs.VIPAuthPass, &cs.VRRPRouterID,
&cs.HBInterface, &cs.HBSrcIP, &cs.HBPeerIP, &cs.HBRouterID, &cs.GWCheckIP); err != nil {
response.Internal(c, err)
return
}
response.OK(c, cs)
}
// UpdateVIPSettings speichert die VIP/VRRP-Configuration 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,
hb_interface=$5, hb_src_ip=$6, hb_peer_ip=$7, hb_router_id=$8, gw_check_ip=$9,
updated_at=NOW()
WHERE id=1`,
nullIfEmpty(req.VIPAddress), nullIfEmpty(req.VIPInterface),
nullIfEmpty(req.VIPAuthPass), req.VRRPRouterID,
nullIfEmpty(req.HBInterface), nullIfEmpty(req.HBSrcIP),
nullIfEmpty(req.HBPeerIP), req.HBRouterID, nullIfEmpty(req.GWCheckIP))
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"`
HBInterface *string `json:"hb_interface"`
HBSrcIP *string `json:"hb_src_ip"`
HBPeerIP *string `json:"hb_peer_ip"`
HBRouterID int `json:"hb_router_id"`
GWCheckIP *string `json:"gw_check_ip"`
}
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)
g.POST("/set-channel", h.AgentSetChannel)
g.GET("/channel", h.AgentChannel)
g.GET("/active-ips", h.AgentActiveIPs)
g.POST("/vip-cmd", h.AgentVIPCmd)
g.GET("/tls-certs", h.AgentTLSCerts)
g.POST("/repair-replication", h.AgentRepairReplication)
g.GET("/repair-replication/status", h.AgentRepairReplicationStatus)
}
// 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) {
pass, err := readFileString(pgReplicationSecretPath)
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-Coordination.
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)
// Publisher-Seite sicherstellen, BEVOR ein Token rausgeht. Ein frisch
// installierter Single-Node hat weder Replikations-Rolle noch
// PUBLICATION noch wal_level=logical — der beitretende Node bekaeme
// beim CREATE SUBSCRIPTION nur ein 404 ("pg-replication-secret nicht
// gefunden") und stuende ohne replizierte Config da. Idempotent; der
// PG-Restart (nur beim allerersten Mal noetig, wal_level ist ein
// postmaster-Parameter) passiert hier bewusst, solange der Admin
// danebensteht und noch kein zweiter Node Traffic erwartet.
if err := h.ensureReplicationPublisher(c.Request.Context()); err != nil {
slog.Error("cluster: publisher setup before join-token failed", "error", err)
response.Internal(c, err)
return
}
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)
}
// ptrStr dereferences a *string safely for comparison; nil → "".
func ptrStr(s *string) string {
if s == nil {
return ""
}
return *s
}
// 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
# --allow-downgrades: nur relevant nach testing→stable-Kanalwechsel
# (Testing-Versionen sortieren datumsbasiert höher als Stable-Semver).
# No-Op im Normalfall, da die Candidate sonst immer >= installed ist.
if apt-get install -y -qq --allow-downgrades -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"})
}
// ── Update-Kanal (stable/testing) ──────────────────────────────────────
//
// Kanal-Modell wie enconf (Suite=Codename, Komponente=Kanal, siehe
// internal/services/apt.Channel/SetChannel) — an EdgeGuards fixes
// Primary/Standby-Paar angepasst statt generischer Server-Flotte: der
// Kanal wird auf beiden Nodes synchron gehalten (wie config_hash),
// kein Node-Override. Reines Umschreiben der sources.list + `apt-get
// update` ist risikofrei (kein Service-Restart, keine VIP-Auswirkung)
// — das eigentliche Downgrade/Upgrade auf die neue Kanal-Version läuft
// danach ganz normal über den bestehenden (sicheren, Standby-zuerst)
// Rolling-Update-Flow, der --allow-downgrades jetzt mit unterstützt.
type updateChannelResponse struct {
Channel string `json:"channel"`
PeerChannel string `json:"peer_channel,omitempty"`
PeerReached bool `json:"peer_reached"`
PeerDrifted bool `json:"peer_drifted"`
}
// UpdateChannel liefert den lokalen Kanal + (falls Cluster) den Kanal
// des Peers zur Drift-Erkennung — analog zum config_hash-Vergleich.
func (h *ClusterHandler) UpdateChannel(c *gin.Context) {
resp := updateChannelResponse{Channel: aptsvc.Channel()}
peer := h.peerNode(c.Request.Context())
if peer != nil && h.Aggregator != nil {
results := h.Aggregator.FanOut(c.Request.Context(), []models.HANode{*peer}, "/agent/cluster/channel", h.LocalID)
if len(results) > 0 && results[0].OK {
var body struct {
Channel string `json:"channel"`
}
if json.Unmarshal(results[0].Data, &body) == nil {
resp.PeerReached = true
resp.PeerChannel = body.Channel
resp.PeerDrifted = body.Channel != resp.Channel
}
}
}
response.OK(c, resp)
}
// SetUpdateChannel setzt den Kanal lokal und — falls ein Peer existiert
// — synchron auch auf dem Peer via mTLS. Löst KEIN Paket-Update aus;
// das übernimmt der Admin danach ganz normal über den Update-Banner /
// Rolling-Update, der die neue Candidate-Version dann bereits sieht.
func (h *ClusterHandler) SetUpdateChannel(c *gin.Context) {
var req struct {
Channel string `json:"channel"`
}
if err := c.ShouldBindJSON(&req); err != nil {
response.BadRequest(c, err)
return
}
if req.Channel != "stable" && req.Channel != "testing" {
response.BadRequest(c, fmt.Errorf("channel must be 'stable' or 'testing'"))
return
}
if err := aptsvc.SetChannel(c.Request.Context(), req.Channel); err != nil {
response.Internal(c, err)
return
}
resp := updateChannelResponse{Channel: req.Channel}
if peer := h.peerNode(c.Request.Context()); peer != nil && h.Aggregator != nil {
body, _ := json.Marshal(req)
result := h.Aggregator.PostPeerWithBody(c.Request.Context(), *peer, "/agent/cluster/set-channel", body)
resp.PeerReached = result.OK
if !result.OK {
slog.Warn("cluster: set-channel on peer failed", "peer", peer.FQDN, "error", result.Err)
}
}
if h.Audit != nil {
_ = h.Audit.Log(c.Request.Context(), actorOf(c), "system.update_channel.set",
"", gin.H{"channel": req.Channel}, h.NodeID)
}
response.OK(c, resp)
}
// AgentChannel: mTLS-Peer-Read des lokalen Kanals (für Drift-Anzeige).
func (h *ClusterHandler) AgentChannel(c *gin.Context) {
response.OK(c, gin.H{"channel": aptsvc.Channel()})
}
// AgentSetChannel: mTLS-Peer-Write — wird vom Primary aufgerufen um den
// Kanal auf diesem (Standby-)Node synchron zu setzen.
func (h *ClusterHandler) AgentSetChannel(c *gin.Context) {
var req struct {
Channel string `json:"channel"`
}
if err := c.ShouldBindJSON(&req); err != nil {
response.BadRequest(c, err)
return
}
if req.Channel != "stable" && req.Channel != "testing" {
response.BadRequest(c, fmt.Errorf("channel must be 'stable' or 'testing'"))
return
}
if err := aptsvc.SetChannel(c.Request.Context(), req.Channel); err != nil {
response.Internal(c, err)
return
}
slog.Info("cluster: update channel set on this node by primary mTLS call",
"channel", req.Channel, "client", c.ClientIP())
response.OK(c, gin.H{"channel": req.Channel})
}
// peerNode liefert die einzige andere ha_nodes-Row (best-effort, nil
// wenn Standalone oder Store fehlt) — gleiches Muster wie in
// RollingUpdate für die Secondary-Ermittlung.
func (h *ClusterHandler) peerNode(ctx context.Context) *models.HANode {
if h.Store == nil {
return nil
}
nodes, err := h.Store.List(ctx)
if err != nil {
return nil
}
for i := range nodes {
if nodes[i].ID != h.LocalID {
return &nodes[i]
}
}
return nil
}
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
// triggering 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
Role string `json:"role"` // "" → "peer" (joining peer); "primary" beim Push des Primary
}
// 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.)
// Rolle aus dem Request (default "peer"). Ein joining-Peer sendet keine
// Rolle → "peer". Der Primary-Push sendet "primary", damit die vom
// Secondary ausgelieferte UI den Primary korrekt als primary zeigt.
// Cert-CN authentifiziert die FQDN; role ist node-lokal/Anzeige (echte
// Rollenerkennung läuft über pg_publication).
role := strings.TrimSpace(req.Role)
if role == "" {
role = "peer"
}
n := models.HANode{
ID: req.ID,
Name: req.Name,
FQDN: req.FQDN,
APIURL: req.APIURL,
Role: role,
Status: "online", // peer IS online — it just connected via mTLS
}
if req.PublicIP != "" {
v := req.PublicIP
n.PublicIP = &v
} else if ip := c.ClientIP(); ip != "" {
// Der Push-Payload (autoRegister, Primary→Secondary) trägt KEINE
// public_ip → sonst bliebe sie NULL und der Peer fehlt im nft-
// peer_ipv4-Set → VRRP-Adverts nur via conntrack → Flapping. Der
// pushende Peer verbindet sich über mTLS von seiner EIGEN-IP (nicht
// der VIP — der Kernel nimmt die primäre Interface-IP als Source),
// genau wie preRegisterJoiner die Joiner-IP übernimmt. Selbstheilend.
n.PublicIP = &ip
}
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)
// Snapshot der aktuellen IPs VOR dem Upsert — zum Vergleich danach.
// Nur wenn sich public_ip oder internal_ip ändert, müssen wir nftables
// neu laden (@peer_ipv4-Set). Periodische Pushes vom Secondary (alle
// 5 min) ändern nur version/config_hash, nicht die IPs → kein Reset.
existing, _ := h.Store.Get(c.Request.Context(), req.ID)
out, err := h.Store.UpsertSelf(c.Request.Context(), n)
if err != nil {
response.Internal(c, err)
return
}
// Firewall-Reload nur wenn sich die Peer-IP geändert hat oder der
// Peer neu eingetragen wurde. Verhindert Counter-Reset alle 5 min
// durch den periodischen Secondary-Push (runPrimaryPush).
ipChanged := existing == nil ||
ptrStr(existing.PublicIP) != ptrStr(out.PublicIP) ||
ptrStr(existing.InternalIP) != ptrStr(out.InternalIP)
if ipChanged && 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)
}
}()
}
// Bei neuem Peer / IP-Wechsel als Info loggen (relevantes Ereignis),
// sonst Debug — die periodischen 30s-Pushes (runPrimaryPush/runPeerPush)
// würden sonst das Log fluten.
logFn := slog.Debug
if ipChanged {
logFn = slog.Info
}
logFn("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)
}
// ensureReplicationPublisher richtet die lokale PG-Instanz als Logical-
// Replication-Publisher ein (Rolle + Secret, wal_level=logical, pg_hba,
// Grants, PUBLICATION). Idempotent — auf einem bereits eingerichteten
// Primary ist es ein No-Op.
//
// Braucht root (psql als postgres, pg_hba schreiben, ggf. PG-Restart), die
// API laeuft als unprivilegierter `edgeguard` → Aufruf via sudo mit
// gepinnter Regel, wie bei den uebrigen privilegierten Operationen.
func (h *ClusterHandler) ensureReplicationPublisher(ctx context.Context) error {
// WICHTIG: nur ausfuehren wenn die Publisher-Seite noch NICHT steht.
// setupReplicationPrimary generiert bei JEDEM Lauf ein neues
// Replikations-Passwort (ALTER ROLE … PASSWORD). Auf einem Cluster mit
// bereits angebundenem Subscriber wuerde dessen gespeicherter
// Connection-String damit ungueltig und die Replikation bliebe still
// stehen — ein zweiter Token-Klick duerfte das niemals ausloesen.
// Das Passwort laesst sich nicht wiederverwenden (in PG nur gehasht),
// deshalb ist "schon eingerichtet" hier ein hartes Abbruchkriterium.
if h.Store != nil {
var hasPub bool
if err := h.Store.Pool.QueryRow(ctx,
`SELECT EXISTS(SELECT 1 FROM pg_publication WHERE pubname = $1)`,
pgPublicationName).Scan(&hasPub); err == nil && hasPub {
if _, err := os.Stat(pgReplicationSecretPath); err == nil {
slog.Info("cluster: replication publisher already set up — skipping init")
return nil
}
}
}
cmd := exec.Command("sudo", "-n", "/usr/bin/edgeguard-ctl", //nolint:noctx // System-Setup, darf nicht am Request-Context haengen
"cluster-init-replication")
out, err := cmd.CombinedOutput()
if err != nil {
return fmt.Errorf("cluster-init-replication: %w: %s",
err, strings.TrimSpace(string(out)))
}
slog.Info("cluster: replication publisher ensured")
return nil
}