feat(cluster): bidirektionaler Peer-Heartbeat (Primary→Secondary Push) — v1.2.97
Bisher pushte nur der Secondary seine Liveness an den Primary (runPrimaryPush). Der Primary pushte nichts → in der lokalen ha_nodes des Secondary fror die Primary-Row nach dem Boot ein → die vom Secondary ausgelieferte UI zeigte den Primary als offline. Neu: runPeerPush auf dem Primary/Founder pusht alle 30s self (role=primary) an jeden Peer via mTLS (/agent/cluster/peers). PushSelfToPeer(role) generalisiert PushSelfToPrimary; registerPeerRequest+AgentRegisterPeer akzeptieren ein role-Feld (default 'peer' → joining-Peer-Verhalten unverändert). Peer-Register-Log bei Routine-Pushes auf Debug (Info nur bei neuem Peer/IP-Wechsel) gegen 30s-Spam. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -194,6 +194,12 @@ func main() {
|
|||||||
}
|
}
|
||||||
// runSecondaryConfigRender wird weiter unten gestartet sobald
|
// runSecondaryConfigRender wird weiter unten gestartet sobald
|
||||||
// clusterAggregator verfügbar ist (braucht mTLS-Client für Cert-Sync).
|
// clusterAggregator verfügbar ist (braucht mTLS-Client für Cert-Sync).
|
||||||
|
} else if nodeID != "" && st != nil && st.Completed && st.FQDN != "" {
|
||||||
|
// Primary/Founder (kein joined Secondary): self (role=primary) an
|
||||||
|
// alle Peers pushen, damit deren lokale ha_nodes den Primary frisch
|
||||||
|
// hält — sonst zeigt die vom Secondary ausgelieferte UI den Primary
|
||||||
|
// als offline. No-op solange keine Peers existieren (Single-Node).
|
||||||
|
go runPeerPush(context.Background(), pool, clusterStore, nodeID, st.FQDN, version)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Phase 3.3: Cluster-CA + Peer-Cert. Founder-Pfad — auf einem
|
// Phase 3.3: Cluster-CA + Peer-Cert. Founder-Pfad — auf einem
|
||||||
@@ -863,6 +869,51 @@ func runPrimaryPush(ctx context.Context, pool *pgxpoolPool, nodeID, fqdn, versio
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// runPeerPush läuft auf dem Primary/Founder und pusht alle 30s die eigene
|
||||||
|
// Identität (role=primary) an jeden Peer via mTLS — das Gegenstück zu
|
||||||
|
// runPrimaryPush (Secondary→Primary). Zusammen ergibt das einen
|
||||||
|
// bidirektionalen Cross-Node-Heartbeat: beide Nodes sehen sich gegenseitig
|
||||||
|
// als online, egal von welchem Node die UI ausgeliefert wird. Tick wie
|
||||||
|
// runPrimaryPush deutlich unter dem 2-min-Stale-Threshold. No-op solange
|
||||||
|
// keine Peers existieren (Single-Node) bzw. wenn ein Peer down ist (Debug-Log).
|
||||||
|
func runPeerPush(ctx context.Context, pool *pgxpoolPool, store *cluster.Store, nodeID, fqdn, version string) {
|
||||||
|
const tick = 30 * time.Second
|
||||||
|
t := time.NewTicker(tick)
|
||||||
|
defer t.Stop()
|
||||||
|
push := func() {
|
||||||
|
pCtx, cancel := context.WithTimeout(ctx, 25*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
peers, err := store.List(pCtx)
|
||||||
|
if err != nil {
|
||||||
|
slog.Warn("cluster: peer-push list failed", "error", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
hash, _ := cluster.ComputeConfigHash(pCtx, pool)
|
||||||
|
for i := range peers {
|
||||||
|
p := peers[i]
|
||||||
|
if p.ID == nodeID {
|
||||||
|
continue // nicht an sich selbst pushen
|
||||||
|
}
|
||||||
|
target := p.APIURL
|
||||||
|
if target == "" {
|
||||||
|
target = "https://" + p.FQDN
|
||||||
|
}
|
||||||
|
if err := clusterjoin.PushSelfToPeer(target, "", nodeID, fqdn, version, hash, "primary"); err != nil {
|
||||||
|
slog.Debug("cluster: push-to-peer failed", "peer", p.FQDN, "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
push() // immediate push on API startup
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
return
|
||||||
|
case <-t.C:
|
||||||
|
push()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func randomEphemeralSecret() []byte {
|
func randomEphemeralSecret() []byte {
|
||||||
b := make([]byte, 32)
|
b := make([]byte, 32)
|
||||||
if _, err := rand.Read(b); err != nil {
|
if _, err := rand.Read(b); err != nil {
|
||||||
|
|||||||
@@ -866,6 +866,7 @@ type registerPeerRequest struct {
|
|||||||
MgmtIP string `json:"mgmt_ip"` // optional
|
MgmtIP string `json:"mgmt_ip"` // optional
|
||||||
Version string `json:"version"`
|
Version string `json:"version"`
|
||||||
ConfigHash *string `json:"config_hash"` // nil=absent (don't change), ""=no user config
|
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.
|
// AgentRegisterPeer: vom Joiner nach issue-cert via mTLS aufgerufen.
|
||||||
@@ -904,12 +905,21 @@ func (h *ClusterHandler) AgentRegisterPeer(c *gin.Context) {
|
|||||||
// Node, hier ist der „Self" der joining-Peer auf dieser Primary-Seite.
|
// 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
|
// Der Name passt nicht 100% semantisch, aber das SQL ist exakt das was
|
||||||
// wir brauchen.)
|
// 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{
|
n := models.HANode{
|
||||||
ID: req.ID,
|
ID: req.ID,
|
||||||
Name: req.Name,
|
Name: req.Name,
|
||||||
FQDN: req.FQDN,
|
FQDN: req.FQDN,
|
||||||
APIURL: req.APIURL,
|
APIURL: req.APIURL,
|
||||||
Role: "peer",
|
Role: role,
|
||||||
Status: "online", // peer IS online — it just connected via mTLS
|
Status: "online", // peer IS online — it just connected via mTLS
|
||||||
}
|
}
|
||||||
if req.PublicIP != "" {
|
if req.PublicIP != "" {
|
||||||
@@ -965,7 +975,14 @@ func (h *ClusterHandler) AgentRegisterPeer(c *gin.Context) {
|
|||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
|
|
||||||
slog.Info("cluster: peer registered via mTLS",
|
// 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,
|
"id", out.ID, "fqdn", out.FQDN, "role", out.Role, "status", out.Status,
|
||||||
"client_cn", cn, "remote", c.ClientIP())
|
"client_cn", cn, "remote", c.ClientIP())
|
||||||
response.OK(c, out)
|
response.OK(c, out)
|
||||||
|
|||||||
@@ -128,7 +128,7 @@ func Join(req Request) error {
|
|||||||
// synchronous on the primary side.
|
// synchronous on the primary side.
|
||||||
var autoRegErr error
|
var autoRegErr error
|
||||||
for i := 0; i < 3; i++ {
|
for i := 0; i < 3; i++ {
|
||||||
if err := autoRegister(primary, tlsDir, req.CommonName, req.Version, req.NodeID, ""); err == nil {
|
if err := autoRegister(primary, tlsDir, req.CommonName, req.Version, req.NodeID, "", "peer"); err == nil {
|
||||||
autoRegErr = nil
|
autoRegErr = nil
|
||||||
break
|
break
|
||||||
} else {
|
} else {
|
||||||
@@ -222,13 +222,21 @@ func issueCert(primary, token, csr string, insecure bool) (caCert, peerCert stri
|
|||||||
// goroutine so the primary's ha_nodes always reflects the secondary's actual
|
// goroutine so the primary's ha_nodes always reflects the secondary's actual
|
||||||
// config_hash (not the stale join-time value).
|
// config_hash (not the stale join-time value).
|
||||||
func PushSelfToPrimary(primaryURL, tlsDir, nodeID, fqdn, version, configHash string) error {
|
func PushSelfToPrimary(primaryURL, tlsDir, nodeID, fqdn, version, configHash string) error {
|
||||||
|
return PushSelfToPeer(primaryURL, tlsDir, nodeID, fqdn, version, configHash, "peer")
|
||||||
|
}
|
||||||
|
|
||||||
|
// PushSelfToPeer sendet die eigene Identität an einen beliebigen Peer (mTLS,
|
||||||
|
// /agent/cluster/peers). role bestimmt, mit welcher Rolle sich dieser Node
|
||||||
|
// beim Empfänger einträgt: ein Secondary pusht "peer" an den Primary, der
|
||||||
|
// Primary pusht "primary" an jeden Secondary (bidirektionaler Heartbeat).
|
||||||
|
func PushSelfToPeer(peerURL, tlsDir, nodeID, fqdn, version, configHash, role string) error {
|
||||||
if tlsDir == "" {
|
if tlsDir == "" {
|
||||||
tlsDir = clustertls.DefaultDir
|
tlsDir = clustertls.DefaultDir
|
||||||
}
|
}
|
||||||
return autoRegister(primaryURL, tlsDir, fqdn, version, nodeID, configHash)
|
return autoRegister(peerURL, tlsDir, fqdn, version, nodeID, configHash, role)
|
||||||
}
|
}
|
||||||
|
|
||||||
func autoRegister(primary, tlsDir, commonName, version, nodeID, configHash string) error {
|
func autoRegister(primary, tlsDir, commonName, version, nodeID, configHash, role string) error {
|
||||||
u, err := url.Parse(primary)
|
u, err := url.Parse(primary)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -241,6 +249,9 @@ func autoRegister(primary, tlsDir, commonName, version, nodeID, configHash strin
|
|||||||
nodeID = strings.TrimSpace(string(raw))
|
nodeID = strings.TrimSpace(string(raw))
|
||||||
}
|
}
|
||||||
hostname, _ := os.Hostname()
|
hostname, _ := os.Hostname()
|
||||||
|
if role == "" {
|
||||||
|
role = "peer"
|
||||||
|
}
|
||||||
body, _ := json.Marshal(map[string]string{
|
body, _ := json.Marshal(map[string]string{
|
||||||
"id": nodeID,
|
"id": nodeID,
|
||||||
"name": hostname,
|
"name": hostname,
|
||||||
@@ -248,6 +259,7 @@ func autoRegister(primary, tlsDir, commonName, version, nodeID, configHash strin
|
|||||||
"api_url": "https://" + commonName + ":3443",
|
"api_url": "https://" + commonName + ":3443",
|
||||||
"version": version,
|
"version": version,
|
||||||
"config_hash": configHash,
|
"config_hash": configHash,
|
||||||
|
"role": role,
|
||||||
})
|
})
|
||||||
|
|
||||||
pair, err := tls.LoadX509KeyPair(tlsDir+"/peer.crt", tlsDir+"/peer.key")
|
pair, err := tls.LoadX509KeyPair(tlsDir+"/peer.crt", tlsDir+"/peer.key")
|
||||||
|
|||||||
Reference in New Issue
Block a user