3 Commits

Author SHA1 Message Date
Debian
b3dda81b49 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>
2026-06-06 13:24:13 +02:00
Debian
b20ace8763 fix(cluster): periodischer Peer-Heartbeat (30s) + Rolling-Update candidate-aware — v1.2.96
Fix 1 — Peer zeigt fälschlich 'offline': runPrimaryPush (Secondary→Primary, einziger periodischer Cross-Node-ha_nodes-Refresh) tickte mit 5 min, SweepStaleNodes-Threshold ist aber 2 min → Secondary war 2 min online, dann 3 min offline, im 5-min-Takt. Tick auf 30s (4× Marge unter Threshold). Receiver lädt nftables nur bei IP-Änderung → kein Reload-Sturm.
Fix 2 — Rolling-Update konnte nie fertig werden wenn der Secondary die Zielversion schon hatte (baseline==target → Warten auf unmöglichen Flip → 10-min-Timeout). runRollingUpdate ist jetzt candidate-aware: ermittelt apt-Candidate, überspringt den Secondary-Schritt wenn dieser schon aktuell ist, erkennt den Flip via 'erreicht candidate ODER bewegt sich von baseline', und schließt direkt mit 'done' wenn auch der Primary schon aktuell ist. FinishRollingUpdateIfPending setzt hängende updating/waiting-secondary-Phasen beim Boot auf idle zurück (tote Orchestrierungs-Goroutine nach Restart).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-06 13:10:40 +02:00
Debian
053b38e46c fix: AlertWriter graceful flush (#15) + Rolling-Update Robustheit (#19) — v1.2.95
#15 waf/alerts.go: AlertWriter.Close() flusht gepufferte Alerts + stoppt die Goroutine (stop/done-Channels, sync.Once, atomic closed; Kanal wird NIE geschlossen → Send racet ohne Panic). Wiring in cmd/edgeguard-waf nach ListenAndServe (graceful shutdown). -race-Test alerts_test.go.
#19 handlers/cluster_rollingupdate.go: (a) RollingUpdateStatus mutiert State nicht mehr beim GET — terminale Zustände altern in readRollingUpdateState nach 10 min aus (kein verlorenes 'done' bei parallelen Pollern). (b) State-File via sync.Mutex + configgen.AtomicWrite (kein partieller Read / Race zwischen Handler & Goroutine). (c) Version-Flip wird gegen die VORHER erfasste Secondary-Baseline geprüft statt gegen die Primary-Version (verhindert sofort-/nie-Flip).
Bewusst belassen: geteilter upgrade.sh-Pfad ist deterministischer Inhalt + an exakte sudoers-Zeile gebunden → Überschreib-Race benign; MST-Timestamp-Parse locale (Server laufen C-Locale).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-06 11:21:50 +02:00
8 changed files with 336 additions and 68 deletions

View File

@@ -1 +1 @@
1.2.94
1.2.97

View File

@@ -194,6 +194,12 @@ func main() {
}
// runSecondaryConfigRender wird weiter unten gestartet sobald
// 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
@@ -826,12 +832,20 @@ func runSecondaryConfigRender(ctx context.Context, pool *pgxpoolPool, box *secre
}
// runPrimaryPush periodically pushes this secondary node's config_hash to the
// primary via mTLS. The primary's ha_nodes view only gets config_hash written
// during join-time autoRegister — after that the primary never hears about
// hash changes unless we push. Without this, the drift banner shows stale
// hashes from join-time forever.
// primary via mTLS. The primary's ha_nodes view only gets config_hash + last_seen
// written during join-time autoRegister — after that the primary never hears about
// the secondary unless we push. Without this, the drift banner shows stale hashes
// from join-time forever AND the secondary's last_seen freezes → SweepStaleNodes
// marks it offline.
//
// WICHTIG: tick MUSS deutlich unter dem Stale-Threshold (4× 30s = 2 min, siehe
// scheduler.staleThreshold / cluster.SweepStaleNodes) liegen. Sonst flippt der
// Secondary zwischen den Pushes zwangsläufig auf "offline" (bei 5-min-Tick:
// 2 min online, 3 min offline). 30s = 4 Pushes pro Stale-Fenster → ein
// verpasster Push (Netz-Glitch) ist unkritisch. Der Receiver (AgentRegisterPeer)
// lädt nftables nur bei IP-Änderung neu → kein Reload-Sturm durch häufige Pushes.
func runPrimaryPush(ctx context.Context, pool *pgxpoolPool, nodeID, fqdn, version, primaryURL string) {
const tick = 5 * time.Minute
const tick = 30 * time.Second
t := time.NewTicker(tick)
defer t.Stop()
push := func() {
@@ -855,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 {
b := make([]byte, 32)
if _, err := rand.Read(b); err != nil {

View File

@@ -83,6 +83,8 @@ func main() {
slog.Error("waf: SPOE agent stopped", "error", err)
os.Exit(1)
}
// Graceful shutdown (ctx cancelled): gepufferte Alerts flushen.
alertWriter.Close()
}
// reload fetches all domain+waf_config pairs from DB and rebuilds engines.

View File

@@ -866,6 +866,7 @@ type registerPeerRequest struct {
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.
@@ -904,12 +905,21 @@ func (h *ClusterHandler) AgentRegisterPeer(c *gin.Context) {
// 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: "peer",
Role: role,
Status: "online", // peer IS online — it just connected via mTLS
}
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,
"client_cn", cn, "remote", c.ClientIP())
response.OK(c, out)

View File

@@ -7,14 +7,21 @@ import (
"net/http"
"os"
"os/exec"
"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"
"git.netcell-it.de/projekte/edgeguard-native/internal/models"
aptsvc "git.netcell-it.de/projekte/edgeguard-native/internal/services/apt"
)
// ruStateMu serialisiert Lesen/Schreiben der Rolling-Update-State-Datei
// (HTTP-Handler + Hintergrund-Goroutine greifen gleichzeitig zu).
var ruStateMu sync.Mutex
const rollingUpdateStateFile = "/var/lib/edgeguard/rolling-update-state.json"
const (
@@ -26,17 +33,24 @@ const (
phaseFailed = "failed"
)
// FinishRollingUpdateIfPending wird beim API-Start aufgerufen. Wenn die
// State-Datei "updating-primary" enthält, bedeutet das dass der Primary
// gerade erfolgreich neugestartet ist → Update abgeschlossen → "done" schreiben.
// FinishRollingUpdateIfPending wird beim API-Start aufgerufen.
// - "updating-primary": der Primary ist gerade erfolgreich neugestartet →
// Update abgeschlossen → "done".
// - "updating-secondary"/"waiting-secondary": die orchestrierende Goroutine
// lief in DIESEM (jetzt neu gestarteten) Prozess und ist mit ihm gestorben.
// Die Phase kann nicht weiterlaufen → auf "idle" zurücksetzen, sonst zeigt
// die UI ewig "Rolling Update läuft". (Vorher blieb so ein Stand hängen.)
func FinishRollingUpdateIfPending() {
st := readRollingUpdateState()
if st.Phase == phaseUpdatingPrimary {
switch st.Phase {
case phaseUpdatingPrimary:
writeRollingUpdateState(RollingUpdateState{
Phase: phaseDone,
SecondaryID: st.SecondaryID,
SecondaryFQDN: st.SecondaryFQDN,
})
case phaseUpdatingSecondary, phaseWaitingSecondary:
writeRollingUpdateState(RollingUpdateState{Phase: phaseIdle})
}
}
@@ -53,6 +67,8 @@ type RollingUpdateState struct {
}
func readRollingUpdateState() RollingUpdateState {
ruStateMu.Lock()
defer ruStateMu.Unlock()
data, err := os.ReadFile(rollingUpdateStateFile)
if err != nil {
return RollingUpdateState{Phase: phaseIdle, UpdatedAt: time.Now()}
@@ -61,6 +77,13 @@ func readRollingUpdateState() RollingUpdateState {
if err := json.Unmarshal(data, &s); err != nil {
return RollingUpdateState{Phase: phaseIdle, UpdatedAt: time.Now()}
}
// Terminale Zustände altern aus (statt Mutation-on-GET): nach 10 min
// gilt done/failed als idle — so verliert kein paralleler Poller das
// Ergebnis und ein alter Stand bleibt nicht hängen.
if (s.Phase == phaseDone || s.Phase == phaseFailed) && !s.UpdatedAt.IsZero() &&
time.Since(s.UpdatedAt) > 10*time.Minute {
return RollingUpdateState{Phase: phaseIdle, UpdatedAt: time.Now()}
}
return s
}
@@ -71,7 +94,10 @@ func writeRollingUpdateState(s RollingUpdateState) {
slog.Warn("rolling-update: failed to marshal state", "error", err)
return
}
if err := os.WriteFile(rollingUpdateStateFile, data, 0o600); err != nil {
ruStateMu.Lock()
defer ruStateMu.Unlock()
// AtomicWrite (temp+rename) → Leser sehen nie einen partiellen Stand.
if err := configgen.AtomicWrite(rollingUpdateStateFile, data, 0o600); err != nil {
slog.Warn("rolling-update: failed to write state file", "error", err)
}
}
@@ -127,76 +153,108 @@ func (h *ClusterHandler) RollingUpdate(c *gin.Context) {
}
// RollingUpdateStatus gibt den aktuellen Rolling-Update-State zurück.
// Bei phase == "done" wird nach Auslieferung sofort auf idle zurückgesetzt
// damit der nächste Pageload keinen Stale-done vorfindet.
// Read-only — terminale Zustände altern in readRollingUpdateState aus
// (kein Reset-on-GET mehr, das parallelen Pollern das "done" wegnahm).
func (h *ClusterHandler) RollingUpdateStatus(c *gin.Context) {
st := readRollingUpdateState()
response.OK(c, st)
if st.Phase == phaseDone {
writeRollingUpdateState(RollingUpdateState{Phase: phaseIdle})
}
response.OK(c, readRollingUpdateState())
}
func (h *ClusterHandler) runRollingUpdate(secondary *models.HANode) {
ctx := context.Background()
// 1. Secondary triggern
slog.Info("rolling-update: posting trigger-update to secondary", "fqdn", secondary.FQDN)
result := h.Aggregator.PostPeer(ctx, *secondary, "/agent/cluster/trigger-update")
if !result.OK {
// Zielversion = das verfügbare apt-Candidate (worauf wir hochziehen) und
// die aktuelle Secondary-Version als Baseline. Beides steuert, ob der
// Secondary überhaupt etwas zu tun hat.
candidate := rollingCandidateVersion(ctx)
baseline := secondaryVersion(ctx, h, secondary)
// Ist der Secondary bereits auf der Zielversion, gibt es nichts
// hochzuziehen — KEIN Trigger, KEIN Warten. Sonst würde auf einen
// Version-Flip gewartet, der nie kommt → 10-min-Timeout (der frühere Bug,
// wenn beide Nodes schon aktuell waren).
secondaryUpToDate := candidate != "" && baseline != "" && baseline == candidate
if secondaryUpToDate {
slog.Info("rolling-update: secondary already at target — skipping secondary step",
"version", candidate)
} else {
// 1. Secondary triggern
slog.Info("rolling-update: posting trigger-update to secondary", "fqdn", secondary.FQDN)
result := h.Aggregator.PostPeer(ctx, *secondary, "/agent/cluster/trigger-update")
if !result.OK {
writeRollingUpdateState(RollingUpdateState{
Phase: phaseFailed,
SecondaryID: secondary.ID,
SecondaryFQDN: secondary.FQDN,
Error: "trigger-update failed: " + result.Err,
})
slog.Warn("rolling-update: secondary trigger failed", "error", result.Err)
return
}
// 2. Secondary-Version pollen — der Secondary restartet nach dem
// Upgrade, danach zeigt /agent/cluster/version eine neue Version.
writeRollingUpdateState(RollingUpdateState{
Phase: phaseFailed,
Phase: phaseWaitingSecondary,
SecondaryID: secondary.ID,
SecondaryFQDN: secondary.FQDN,
Error: "trigger-update failed: " + result.Err,
})
slog.Warn("rolling-update: secondary trigger failed", "error", result.Err)
return
}
slog.Info("rolling-update: waiting for secondary version flip",
"baseline", baseline, "candidate", candidate)
// 2. Secondary-Version pollen — der Secondary restartet nach dem
// Upgrade, danach zeigt /agent/cluster/version eine neue Version.
writeRollingUpdateState(RollingUpdateState{
Phase: phaseWaitingSecondary,
SecondaryID: secondary.ID,
SecondaryFQDN: secondary.FQDN,
})
slog.Info("rolling-update: waiting for secondary version flip")
// Kurze Wartezeit damit apt auf dem Secondary erst losläuft
time.Sleep(20 * time.Second)
// Kurze Wartezeit damit apt auf dem Secondary erst losläuft
time.Sleep(20 * time.Second)
deadline := time.Now().Add(10 * time.Minute)
versionFlipped := false
for time.Now().Before(deadline) {
results := h.Aggregator.FanOut(ctx, []models.HANode{*secondary}, "/agent/cluster/version", h.LocalID)
if len(results) > 0 && results[0].OK {
var ver struct {
Version string `json:"version"`
}
if err := json.Unmarshal(results[0].Data, &ver); err == nil {
slog.Info("rolling-update: secondary version", "version", ver.Version, "primary", h.Version)
if ver.Version != h.Version {
versionFlipped = true
break
deadline := time.Now().Add(10 * time.Minute)
versionFlipped := false
for time.Now().Before(deadline) {
results := h.Aggregator.FanOut(ctx, []models.HANode{*secondary}, "/agent/cluster/version", h.LocalID)
if len(results) > 0 && results[0].OK {
var ver struct {
Version string `json:"version"`
}
if err := json.Unmarshal(results[0].Data, &ver); err == nil {
slog.Info("rolling-update: secondary version", "version", ver.Version,
"baseline", baseline, "candidate", candidate)
// Erfolg = Secondary hat die Zielversion erreicht (candidate)
// ODER hat sich gegenüber der Baseline überhaupt bewegt
// (Fallback, wenn candidate nicht ermittelbar war).
if ver.Version != "" &&
((candidate != "" && ver.Version == candidate) || ver.Version != baseline) {
versionFlipped = true
break
}
}
}
time.Sleep(10 * time.Second)
}
if !versionFlipped {
writeRollingUpdateState(RollingUpdateState{
Phase: phaseFailed,
SecondaryID: secondary.ID,
SecondaryFQDN: secondary.FQDN,
Error: "timeout (10 min) waiting for secondary version flip",
})
slog.Warn("rolling-update: secondary version flip timeout")
return
}
time.Sleep(10 * time.Second)
}
if !versionFlipped {
// 3. Primary (uns selbst) aktualisieren — identisch zu /system/upgrade.
// Ist der Primary bereits auf der Zielversion (z. B. beide Nodes schon
// aktuell), gibt es nichts zu tun → direkt "done". Sonst liefe ein
// apt-Lauf ohne Paket-Wechsel → kein Restart → Phase hinge ewig in
// "updating-primary".
if candidate != "" && h.Version == candidate {
slog.Info("rolling-update: primary already at target — nothing to upgrade", "version", candidate)
writeRollingUpdateState(RollingUpdateState{
Phase: phaseFailed,
Phase: phaseDone,
SecondaryID: secondary.ID,
SecondaryFQDN: secondary.FQDN,
Error: "timeout (10 min) waiting for secondary version flip",
})
slog.Warn("rolling-update: secondary version flip timeout")
return
}
// 3. Primary (uns selbst) aktualisieren — identisch zu /system/upgrade
writeRollingUpdateState(RollingUpdateState{
Phase: phaseUpdatingPrimary,
SecondaryID: secondary.ID,
@@ -256,3 +314,26 @@ rm -f /var/lib/edgeguard/upgrade.sh
// UI erkennt Version-Flip via /system/health und schließt den Flow.
slog.Info("rolling-update: primary upgrade dispatched, process will restart")
}
// rollingCandidateVersion liefert best-effort die verfügbare apt-Candidate-
// Version des Meta-Pakets "edgeguard" — also die Version, auf die das Rolling-
// Update hochzieht. Leerer String, wenn apt sie nicht ermitteln kann (dann
// fällt runRollingUpdate auf reine Baseline-Flip-Erkennung zurück).
func rollingCandidateVersion(ctx context.Context) string {
vers := aptsvc.PackageVersions(ctx, false)
return vers["edgeguard_available"]
}
// secondaryVersion holt best-effort die laufende Version des Peers via mTLS.
func secondaryVersion(ctx context.Context, h *ClusterHandler, secondary *models.HANode) string {
results := h.Aggregator.FanOut(ctx, []models.HANode{*secondary}, "/agent/cluster/version", h.LocalID)
if len(results) > 0 && results[0].OK {
var ver struct {
Version string `json:"version"`
}
if json.Unmarshal(results[0].Data, &ver) == nil {
return ver.Version
}
}
return ""
}

View File

@@ -128,7 +128,7 @@ func Join(req Request) error {
// synchronous on the primary side.
var autoRegErr error
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
break
} 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
// config_hash (not the stale join-time value).
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 == "" {
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)
if err != nil {
return err
@@ -241,6 +249,9 @@ func autoRegister(primary, tlsDir, commonName, version, nodeID, configHash strin
nodeID = strings.TrimSpace(string(raw))
}
hostname, _ := os.Hostname()
if role == "" {
role = "peer"
}
body, _ := json.Marshal(map[string]string{
"id": nodeID,
"name": hostname,
@@ -248,6 +259,7 @@ func autoRegister(primary, tlsDir, commonName, version, nodeID, configHash strin
"api_url": "https://" + commonName + ":3443",
"version": version,
"config_hash": configHash,
"role": role,
})
pair, err := tls.LoadX509KeyPair(tlsDir+"/peer.crt", tlsDir+"/peer.key")

View File

@@ -3,6 +3,8 @@ package waf
import (
"context"
"log/slog"
"sync"
"sync/atomic"
"time"
"github.com/jackc/pgx/v5/pgxpool"
@@ -26,8 +28,12 @@ type Alert struct {
// AlertWriter accepts Alert values via a buffered channel and writes
// them to PostgreSQL asynchronously so SPOE handling stays low-latency.
type AlertWriter struct {
pool *pgxpool.Pool
ch chan Alert
pool *pgxpool.Pool
ch chan Alert
stop chan struct{}
done chan struct{}
closeOnce sync.Once
closed atomic.Bool
}
// NewAlertWriter creates an AlertWriter and starts its background goroutine.
@@ -36,14 +42,19 @@ func NewAlertWriter(pool *pgxpool.Pool, bufSize int) *AlertWriter {
aw := &AlertWriter{
pool: pool,
ch: make(chan Alert, bufSize),
stop: make(chan struct{}),
done: make(chan struct{}),
}
go aw.run()
return aw
}
// Send enqueues an alert. Drops silently if the channel is full to
// avoid slowing down SPOE request handling.
// Send enqueues an alert. Drops silently if the channel is full (or the
// writer is closing) to avoid slowing down / panicking SPOE handling.
func (aw *AlertWriter) Send(a Alert) {
if aw.closed.Load() {
return
}
select {
case aw.ch <- a:
default:
@@ -51,9 +62,34 @@ func (aw *AlertWriter) Send(a Alert) {
}
}
// Close stops the writer and flushes buffered alerts (best-effort).
// Safe to call multiple times. The channel is never closed → Send never
// panics even if it races with Close.
func (aw *AlertWriter) Close() {
aw.closeOnce.Do(func() {
aw.closed.Store(true)
close(aw.stop)
})
<-aw.done
}
func (aw *AlertWriter) run() {
for a := range aw.ch {
aw.write(a)
defer close(aw.done)
for {
select {
case a := <-aw.ch:
aw.write(a)
case <-aw.stop:
// Restliche gepufferte Alerts noch wegschreiben, dann Ende.
for {
select {
case a := <-aw.ch:
aw.write(a)
default:
return
}
}
}
}
}

View File

@@ -0,0 +1,61 @@
package waf
import (
"context"
"os"
"sync"
"testing"
"time"
"git.netcell-it.de/projekte/edgeguard-native/internal/database"
)
// Beweist Fix #15: AlertWriter.Close() flusht, ist idempotent, und Send/Close
// racen ohne Panic (Kanal wird nie geschlossen). Guarded per EG_FWTEST_DSN.
func TestAlertWriter_CloseFlush(t *testing.T) {
dsn := os.Getenv("EG_FWTEST_DSN")
if dsn == "" {
t.Skip("set EG_FWTEST_DSN to run the alert-writer test")
}
ctx := context.Background()
var mErr error
for i := 0; i < 3; i++ {
if mErr = database.Migrate(ctx, dsn); mErr == nil {
break
}
time.Sleep(700 * time.Millisecond)
}
if mErr != nil {
t.Fatalf("migrate: %v", mErr)
}
pool, err := database.Open(ctx, dsn)
if err != nil {
t.Fatalf("open: %v", err)
}
defer pool.Close()
aw := NewAlertWriter(pool, 64)
for i := 0; i < 20; i++ {
aw.Send(Alert{Hostname: "t.local", ClientIP: "203.0.113.1", Method: "GET", URI: "/", Action: "detected"})
}
// Send parallel zu Close → darf nicht paniken.
var wg sync.WaitGroup
for i := 0; i < 10; i++ {
wg.Add(1)
go func() { defer wg.Done(); aw.Send(Alert{Hostname: "t.local", Action: "detected"}) }()
}
done := make(chan struct{})
go func() { aw.Close(); close(done) }()
select {
case <-done:
case <-time.After(10 * time.Second):
t.Fatal("Close() did not return (flush hung)")
}
wg.Wait()
// Idempotent + Send nach Close ist No-op (kein Panic).
aw.Close()
aw.Send(Alert{Hostname: "after.local", Action: "detected"})
}