From 39a0875d58044db5dae4910ea978be5d1cd378da Mon Sep 17 00:00:00 2001 From: Steve Cliff Date: Sat, 22 Aug 2026 11:31:56 +0100 Subject: [PATCH] test(alert): drain notifications before teardown --- internal/alert/engine.go | 11 +++++++++++ internal/alert/rules.go | 8 ++++---- internal/alert/rules_test.go | 1 + 3 files changed, 16 insertions(+), 4 deletions(-) diff --git a/internal/alert/engine.go b/internal/alert/engine.go index 8a536c9..807019d 100644 --- a/internal/alert/engine.go +++ b/internal/alert/engine.go @@ -61,9 +61,20 @@ type Engine struct { stuckThresholds map[string]time.Duration closeOnce sync.Once + notifyWG sync.WaitGroup done chan struct{} } +func (e *Engine) dispatchNotification(ctx context.Context, payload notification.Payload) { + e.notifyWG.Add(1) + go func() { + defer e.notifyWG.Done() + e.hub.Dispatch(ctx, payload) + }() +} + +func (e *Engine) waitNotifications() { e.notifyWG.Wait() } + // NewEngine builds the engine. agentOfflineFloor + tickPeriod default // to 15min and 60s respectively when zero. func NewEngine(st *store.Store, hub *notification.Hub) *Engine { diff --git a/internal/alert/rules.go b/internal/alert/rules.go index d00c91e..6c30067 100644 --- a/internal/alert/rules.go +++ b/internal/alert/rules.go @@ -60,7 +60,7 @@ func (e *Engine) raiseAndNotify(ctx context.Context, hostID, kind, dedupKey, sev if err == nil { hostName = host.Name } - go e.hub.Dispatch(ctx, notification.Payload{ + e.dispatchNotification(ctx, notification.Payload{ Event: notification.EventRaised, AlertID: id, Severity: severity, @@ -85,7 +85,7 @@ func (e *Engine) Acknowledge(ctx context.Context, alertID, userID string, when t return nil //nolint:nilerr } p := alertPayload(ctx, e.store, notification.EventAcknowledged, a) - go e.hub.Dispatch(context.WithoutCancel(ctx), p) + e.dispatchNotification(context.WithoutCancel(ctx), p) return nil } @@ -99,7 +99,7 @@ func (e *Engine) Resolve(ctx context.Context, alertID string, when time.Time) er return nil } p := alertPayload(ctx, e.store, notification.EventResolved, a) - go e.hub.Dispatch(context.WithoutCancel(ctx), p) + e.dispatchNotification(context.WithoutCancel(ctx), p) return nil } @@ -164,7 +164,7 @@ func (e *Engine) resolveAndNotify(ctx context.Context, hostID, kind, dedupKey st if a.Kind != kind || a.DedupKey != dedupKey { continue } - go e.hub.Dispatch(ctx, notification.Payload{ + e.dispatchNotification(ctx, notification.Payload{ Event: notification.EventResolved, AlertID: a.ID, Severity: a.Severity, diff --git a/internal/alert/rules_test.go b/internal/alert/rules_test.go index ec6a2d3..8a0fa0e 100644 --- a/internal/alert/rules_test.go +++ b/internal/alert/rules_test.go @@ -25,6 +25,7 @@ func setupEngine(t *testing.T) (*Engine, *store.Store, string) { aead, _ := crypto.NewAEAD(key) hub := notification.NewHub(st, aead, "https://rm.example") eng := NewEngine(st, hub) + t.Cleanup(eng.waitNotifications) hostID := ulid.Make().String() if err := st.CreateHost(context.Background(), store.Host{ ID: hostID, Name: "alfa-01", OS: "linux", Arch: "amd64",