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>
This commit is contained in:
@@ -53,6 +53,14 @@ type ClusterHandler struct {
|
||||
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}
|
||||
}
|
||||
@@ -284,8 +292,7 @@ func (h *ClusterHandler) AgentIdentity(c *gin.Context) {
|
||||
// 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)
|
||||
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
|
||||
@@ -518,6 +525,20 @@ func (h *ClusterHandler) GenerateJoinToken(c *gin.Context) {
|
||||
// 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)
|
||||
@@ -1128,3 +1149,43 @@ func (h *ClusterHandler) AgentRegisterPeer(c *gin.Context) {
|
||||
"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
|
||||
}
|
||||
|
||||
@@ -64,6 +64,7 @@ func (h *SetupHandler) Register(rg *gin.RouterGroup) {
|
||||
g.POST("/complete", h.Complete)
|
||||
g.POST("/complete-node", h.CompleteAsNode)
|
||||
g.POST("/join-cluster", h.JoinCluster)
|
||||
g.GET("/replication-status", h.ReplicationStatus)
|
||||
}
|
||||
|
||||
// RegisterAuthed mountet die Endpoints die nach abgeschlossenem Setup
|
||||
@@ -206,6 +207,13 @@ func (h *SetupHandler) JoinCluster(c *gin.Context) {
|
||||
go h.preRegisterPrimary(body.PrimaryFQDN)
|
||||
}
|
||||
|
||||
// Logical Replication automatisch einrichten. Ohne diesen Schritt waere
|
||||
// der Node zwar im Cluster registriert, wuerde aber keinerlei geteilte
|
||||
// Config (Domains, Backends, Firewall-Rules, WireGuard, …) bekommen —
|
||||
// was frueher erst beim Failover auffiel. Laeuft detached, der Wizard
|
||||
// pollt /setup/replication-status.
|
||||
h.startReplicationSetup(body.PrimaryFQDN)
|
||||
|
||||
response.OK(c, gin.H{
|
||||
"completed": st.Completed,
|
||||
"is_cluster_node": st.IsClusterNode,
|
||||
|
||||
195
internal/handlers/setup_replication.go
Normal file
195
internal/handlers/setup_replication.go
Normal file
@@ -0,0 +1,195 @@
|
||||
package handlers
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"log/slog"
|
||||
"net"
|
||||
"os"
|
||||
"os/exec"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
|
||||
"git.netcell-it.de/projekte/edgeguard-native/internal/configgen"
|
||||
"git.netcell-it.de/projekte/edgeguard-native/internal/handlers/response"
|
||||
)
|
||||
|
||||
// Automatische Logical-Replication-Einrichtung beim Cluster-Join.
|
||||
//
|
||||
// Früher war das ein manueller Schritt: nach dem Join musste der Operator
|
||||
// auf dem neuen Node `edgeguard-ctl cluster-setup-standby <primary>`
|
||||
// ausführen. Wer das übersah, hatte einen Node, der im Cluster sichtbar
|
||||
// war, aber KEINE geteilte Config replizierte — und merkte es erst beim
|
||||
// Failover. Deshalb läuft es jetzt direkt aus dem Join heraus.
|
||||
//
|
||||
// Der eigentliche Ablauf bleibt im CLI (`cluster-setup-standby`): er
|
||||
// braucht root (psql als postgres-User, pg_hba, render-config), die API
|
||||
// läuft als unprivilegierter `edgeguard`. Aufruf daher via sudo mit
|
||||
// gepinnter Regel — gleiches Muster wie bei apt-get/systemctl/tee.
|
||||
//
|
||||
// Weil die Initialkopie der geteilten Tabellen Minuten dauern kann, läuft
|
||||
// das detached; der Setup-Wizard pollt GET /setup/replication-status.
|
||||
|
||||
const replicationStateFile = "/var/lib/edgeguard/replication-setup-state.json"
|
||||
|
||||
const (
|
||||
replPhaseIdle = "idle"
|
||||
replPhaseRunning = "running"
|
||||
replPhaseDone = "done"
|
||||
replPhaseFailed = "failed"
|
||||
)
|
||||
|
||||
// replStateMu serialisiert Lesen/Schreiben der State-Datei (HTTP-Handler
|
||||
// + Hintergrund-Goroutine greifen gleichzeitig zu).
|
||||
var replStateMu sync.Mutex
|
||||
|
||||
// ReplicationSetupState hält den Fortschritt der Standby-Einrichtung.
|
||||
// Persistiert, damit der Status einen API-Neustart übersteht — der ist
|
||||
// der letzte Schritt des Setups und würde den Zustand sonst verlieren.
|
||||
type ReplicationSetupState struct {
|
||||
Phase string `json:"phase"`
|
||||
Primary string `json:"primary,omitempty"`
|
||||
Error string `json:"error,omitempty"`
|
||||
Log string `json:"log,omitempty"`
|
||||
StartedAt time.Time `json:"started_at,omitempty"`
|
||||
UpdatedAt time.Time `json:"updated_at,omitempty"`
|
||||
}
|
||||
|
||||
func readReplicationState() ReplicationSetupState {
|
||||
replStateMu.Lock()
|
||||
defer replStateMu.Unlock()
|
||||
raw, err := os.ReadFile(replicationStateFile)
|
||||
if err != nil {
|
||||
return ReplicationSetupState{Phase: replPhaseIdle}
|
||||
}
|
||||
var st ReplicationSetupState
|
||||
if err := json.Unmarshal(raw, &st); err != nil {
|
||||
return ReplicationSetupState{Phase: replPhaseIdle}
|
||||
}
|
||||
if st.Phase == "" {
|
||||
st.Phase = replPhaseIdle
|
||||
}
|
||||
// Ein "running", das älter als das CLI-Timeout ist, kann nur von einem
|
||||
// gestorbenen Prozess stammen (z. B. OOM-Kill). Sonst haengt der Wizard
|
||||
// ewig im Spinner.
|
||||
if st.Phase == replPhaseRunning && !st.StartedAt.IsZero() &&
|
||||
time.Since(st.StartedAt) > 15*time.Minute {
|
||||
st.Phase = replPhaseFailed
|
||||
st.Error = "Zeitüberschreitung — Einrichtung lief länger als 15 Minuten. " +
|
||||
"Manuell nachholen: sudo edgeguard-ctl cluster-setup-standby " + st.Primary
|
||||
}
|
||||
return st
|
||||
}
|
||||
|
||||
func writeReplicationState(st ReplicationSetupState) {
|
||||
st.UpdatedAt = time.Now()
|
||||
raw, err := json.Marshal(st)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
replStateMu.Lock()
|
||||
defer replStateMu.Unlock()
|
||||
if err := configgen.AtomicWrite(replicationStateFile, raw, 0o640); err != nil {
|
||||
slog.Warn("setup: replication state write failed", "error", err)
|
||||
}
|
||||
}
|
||||
|
||||
// validPrimaryHost laesst nur das durch, was ein Hostname oder eine IP
|
||||
// sein kann. exec.Command startet keine Shell, Metazeichen koennen also
|
||||
// ohnehin nichts ausloesen — die Pruefung haelt aber Unsinn von der
|
||||
// sudo-Regel fern und liefert dem Operator einen klaren Fehler statt
|
||||
// eines kryptischen CLI-Abbruchs.
|
||||
func validPrimaryHost(h string) bool {
|
||||
h = strings.TrimSpace(h)
|
||||
if h == "" || len(h) > 253 {
|
||||
return false
|
||||
}
|
||||
if net.ParseIP(h) != nil {
|
||||
return true
|
||||
}
|
||||
for _, label := range strings.Split(h, ".") {
|
||||
if label == "" {
|
||||
return false
|
||||
}
|
||||
for _, r := range label {
|
||||
isAlnum := (r >= 'a' && r <= 'z') || (r >= 'A' && r <= 'Z') || (r >= '0' && r <= '9')
|
||||
if !isAlnum && r != '-' {
|
||||
return false
|
||||
}
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// startReplicationSetup richtet diesen Node im Hintergrund als Logical-
|
||||
// Replication-Subscriber ein. Nicht-blockierend: der Join-Request
|
||||
// antwortet sofort, der Wizard pollt den Status.
|
||||
func (h *SetupHandler) startReplicationSetup(primary string) {
|
||||
primary = strings.ToLower(strings.TrimSpace(primary))
|
||||
if !validPrimaryHost(primary) {
|
||||
writeReplicationState(ReplicationSetupState{
|
||||
Phase: replPhaseFailed,
|
||||
Error: "ungültiger Primary-Host: " + primary,
|
||||
})
|
||||
return
|
||||
}
|
||||
|
||||
writeReplicationState(ReplicationSetupState{
|
||||
Phase: replPhaseRunning,
|
||||
Primary: primary,
|
||||
StartedAt: time.Now(),
|
||||
})
|
||||
|
||||
go func() {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
slog.Error("setup: replication setup panic", "panic", r)
|
||||
writeReplicationState(ReplicationSetupState{
|
||||
Phase: replPhaseFailed, Primary: primary,
|
||||
Error: "interner Fehler bei der Replikations-Einrichtung",
|
||||
})
|
||||
}
|
||||
}()
|
||||
|
||||
slog.Info("setup: starting logical replication setup", "primary", primary)
|
||||
// Kein Request-Context: der Join-Request ist längst beantwortet,
|
||||
// und ein Abbruch mitten im CREATE SUBSCRIPTION wäre schlimmer
|
||||
// als ein Weiterlaufen.
|
||||
cmd := exec.Command("sudo", "-n", "/usr/bin/edgeguard-ctl", //nolint:noctx // detached by design — darf nicht am Request haengen
|
||||
"cluster-setup-standby", primary)
|
||||
out, err := cmd.CombinedOutput()
|
||||
logTail := tailString(string(out), 4000)
|
||||
|
||||
if err != nil {
|
||||
slog.Warn("setup: logical replication setup failed",
|
||||
"primary", primary, "error", err, "output", logTail)
|
||||
writeReplicationState(ReplicationSetupState{
|
||||
Phase: replPhaseFailed, Primary: primary,
|
||||
Error: err.Error(), Log: logTail,
|
||||
})
|
||||
return
|
||||
}
|
||||
slog.Info("setup: logical replication setup finished", "primary", primary)
|
||||
writeReplicationState(ReplicationSetupState{
|
||||
Phase: replPhaseDone, Primary: primary, Log: logTail,
|
||||
})
|
||||
}()
|
||||
}
|
||||
|
||||
// tailString kuerzt lange CLI-Ausgaben auf die letzten n Bytes — der
|
||||
// interessante Teil (Fehler, Abschlussmeldung) steht am Ende.
|
||||
func tailString(s string, n int) string {
|
||||
if len(s) <= n {
|
||||
return s
|
||||
}
|
||||
return "…" + s[len(s)-n:]
|
||||
}
|
||||
|
||||
// ReplicationStatus liefert den Fortschritt der automatischen Standby-
|
||||
// Einrichtung. Liegt bewusst auf der Setup-Gruppe (pre-auth): der Wizard
|
||||
// pollt es, bevor auf dem neuen Node ueberhaupt ein Login moeglich ist.
|
||||
func (h *SetupHandler) ReplicationStatus(c *gin.Context) {
|
||||
response.OK(c, readReplicationState())
|
||||
}
|
||||
53
internal/handlers/setup_replication_test.go
Normal file
53
internal/handlers/setup_replication_test.go
Normal file
@@ -0,0 +1,53 @@
|
||||
package handlers
|
||||
|
||||
import "testing"
|
||||
|
||||
// validPrimaryHost bewacht das einzige variable Argument einer sudo-Regel
|
||||
// (`edgeguard-ctl cluster-setup-standby *`). Der Aufruf laeuft zwar ohne
|
||||
// Shell, aber die Pruefung soll trotzdem halten was sie verspricht.
|
||||
func TestValidPrimaryHost(t *testing.T) {
|
||||
valid := []string{
|
||||
"utm-1.netcell-it.de",
|
||||
"primary",
|
||||
"10.0.5.1",
|
||||
"89.163.205.6",
|
||||
"2001:db8::1",
|
||||
"a-b-c.example.com",
|
||||
}
|
||||
for _, h := range valid {
|
||||
if !validPrimaryHost(h) {
|
||||
t.Errorf("validPrimaryHost(%q) = false, erwartet true", h)
|
||||
}
|
||||
}
|
||||
|
||||
invalid := []string{
|
||||
"",
|
||||
" ",
|
||||
"host; rm -rf /",
|
||||
"host && reboot",
|
||||
"host|tee",
|
||||
"host$(id)",
|
||||
"host`id`",
|
||||
"--tls-dir=/tmp/evil",
|
||||
"host with space",
|
||||
"host\nsecond-line",
|
||||
"..",
|
||||
"host..example.com",
|
||||
"/etc/passwd",
|
||||
}
|
||||
for _, h := range invalid {
|
||||
if validPrimaryHost(h) {
|
||||
t.Errorf("validPrimaryHost(%q) = true, erwartet false", h)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidPrimaryHostRejectsOverlongName(t *testing.T) {
|
||||
long := make([]byte, 254)
|
||||
for i := range long {
|
||||
long[i] = 'a'
|
||||
}
|
||||
if validPrimaryHost(string(long)) {
|
||||
t.Error("Hostname > 253 Zeichen muss abgelehnt werden")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user