diff --git a/charts/terdut-server/templates/deployment.yaml b/charts/terdut-server/templates/deployment.yaml index e12ec2f..f2eb2b0 100644 --- a/charts/terdut-server/templates/deployment.yaml +++ b/charts/terdut-server/templates/deployment.yaml @@ -10,11 +10,14 @@ spec: selector: matchLabels: {{- include "terdut-server.selectorLabels" . | nindent 6 }} - # Recreate, not RollingUpdate, even though the PVC that forced it is gone: the - # sweeper, notifier and migration runner now take a Postgres advisory lock - # each, so two replicas overlapping during a rollout no longer both page for - # the same incident or race applying a migration, but new-incident creation - # on the first webhook for a brand-new groupKey is still unguarded. + # Recreate, not RollingUpdate, even though the PVC that forced it is gone: + # the sweeper, notifier and migration runner take a Postgres advisory lock + # each, and new-incident creation on the first webhook for a brand-new + # groupKey resolves its own insert conflict — so two replicas overlapping + # during a rollout no longer double-page, race a migration, or drop a + # webhook payload. Nothing left here actually requires Recreate anymore; + # it stays the default pending a deliberate decision to raise replicas + # above 1 and move to RollingUpdate. strategy: type: Recreate template: diff --git a/internal/api/alertmanager.go b/internal/api/alertmanager.go index 1badf72..012513f 100644 --- a/internal/api/alertmanager.go +++ b/internal/api/alertmanager.go @@ -376,6 +376,17 @@ func incidentForGroup(ctx context.Context, tx *sql.Tx, notify NotifyConfig, team // its own. Hence the querier rather than a *sql.Tx. A nil severity leaves the // column for refreshSeverity to fill from the member alerts; the sweeper passes // one because its incidents have no members to derive it from. +// +// Both callers get here only after their own SELECT found no open incident for +// this group_key — but on more than one replica, two webhook deliveries for the +// very first occurrence of a brand-new group_key can both pass that SELECT +// before either INSERTs. ON CONFLICT DO NOTHING against +// incidents_open_group_key_idx is what makes the loser's INSERT a no-op instead +// of a unique-violation error that would otherwise roll back its entire +// payload; existingOpenIncident then hands it the winner's row. Postgres +// resolves that conflict only once the winner's transaction has committed (or +// rolled back), so by the time this RETURNING comes back empty, the winner's +// row is guaranteed visible to that follow-up SELECT. func openIncident(ctx context.Context, q querier, notify NotifyConfig, teamID int64, groupKey, title string, groupLabels map[string]string, severity *string) (int64, error) { onCall, err := currentOnCall(ctx, q, teamID) if err != nil { @@ -391,10 +402,18 @@ func openIncident(ctx context.Context, q querier, notify NotifyConfig, teamID in err = q.QueryRowContext(ctx, ` INSERT INTO incidents (team_id, group_key, title, group_labels, signature, status, severity, triggered_at, assigned_to) VALUES ($1, $2, $3, $4::jsonb, $5, 'triggered', $6, $7, $8) + ON CONFLICT (team_id, group_key) WHERE resolved_at IS NULL DO NOTHING RETURNING id`, teamID, groupKey, title, string(labelsJSON), incidentSignature(groupLabels, title), severity, time.Now().Unix(), onCall).Scan(&id) - if err != nil { + switch { + case err == sql.ErrNoRows: + // Lost the race: someone else's incident for this group_key exists now. + // Everything below — the trigger event, assignment, page, escalation + // clock — already happened for that row when it was created; attach to + // it rather than fail this call (and the whole payload) outright. + return existingOpenIncident(ctx, q, teamID, groupKey) + case err != nil: return 0, err } @@ -423,6 +442,21 @@ func openIncident(ctx context.Context, q querier, notify NotifyConfig, teamID in return id, nil } +// existingOpenIncident looks up the open incident openIncident's own INSERT just +// lost a conflict against — the same lookup incidentForGroup does before ever +// calling openIncident, repeated here for the caller that arrived second. +func existingOpenIncident(ctx context.Context, q querier, teamID int64, groupKey string) (int64, error) { + var id int64 + err := q.QueryRowContext(ctx, + "SELECT id FROM incidents WHERE team_id = $1 AND group_key = $2 AND resolved_at IS NULL", + teamID, groupKey, + ).Scan(&id) + if err != nil { + return 0, err + } + return id, nil +} + // linkAlert adds an alert to an incident, emitting a timeline entry only the // first time. Re-sends of an already-linked alert are silent. func linkAlert(ctx context.Context, tx *sql.Tx, incidentID, alertID int64) error { diff --git a/internal/api/alertmanager_test.go b/internal/api/alertmanager_test.go new file mode 100644 index 0000000..eceb4ab --- /dev/null +++ b/internal/api/alertmanager_test.go @@ -0,0 +1,112 @@ +package api_test + +import ( + "bytes" + "encoding/json" + "fmt" + "net/http" + "sync" + "testing" +) + +// TestWebhook_ConcurrentFirstOccurrenceOpensOneIncident reproduces two +// replicas racing the very first webhook delivery for a brand-new group_key: +// both see no open incident yet (incidentForGroup's own SELECT finds +// nothing) and race openIncident's INSERT. +// +// The DB's own unique index already guarantees at most one incident either +// way, with or without this fix — so "exactly one incident" alone cannot +// tell a fixed run from a broken one. What ON CONFLICT handling actually +// changes is what happens to the *loser*: before it, the loser's INSERT hit +// incidents_open_group_key_idx's unique violation, which — since +// upsertAlerts ran earlier in that same transaction — rolled back its whole +// payload, alert insert included. ingest's error is only logged and +// receiveWebhook answers 200 regardless, so nothing ever retried it: the +// loser's alert silently never existed. That is the regression signal this +// test checks — every caller's fingerprint must show up in /api/alerts, not +// just the winner's. +func TestWebhook_ConcurrentFirstOccurrenceOpensOneIncident(t *testing.T) { + s := newTS(t) + + const callers = 8 + const groupKey = "race-group" + + // Every caller needs its own fingerprint. A shared one would serialize all + // of them at upsertAlerts' own ON CONFLICT (team_id, fingerprint) row lock, + // long before any of them reached incidentForGroup — which would hide the + // very race this test exists to force. + bodies := make([][]byte, callers) + for i := range callers { + payload := map[string]any{ + "version": "4", + "status": "firing", + "groupKey": groupKey, + "groupLabels": map[string]string{"alertname": "RaceAlert"}, + "alerts": []map[string]any{amAlert(fmt.Sprintf("fp-race-%d", i), "RaceAlert", "firing", + "2026-05-20T10:00:00Z", "0001-01-01T00:00:00Z", nil)}, + } + bodies[i], _ = json.Marshal(payload) + } + + // A start line, so every request is fired as close to simultaneously as + // goroutine scheduling allows, rather than trickling out one dial at a + // time — the race window is the gap between incidentForGroup's SELECT and + // openIncident's INSERT, which a staggered start could easily miss. + var ready sync.WaitGroup + start := make(chan struct{}) + statuses := make([]int, callers) + var wg sync.WaitGroup + for i := range callers { + ready.Add(1) + wg.Add(1) + go func(i int) { + defer wg.Done() + ready.Done() + <-start + resp, err := http.Post(s.URL+"/api/integrations/"+s.ingestKey+"/alertmanager", + "application/json", bytes.NewReader(bodies[i])) + if err != nil { + t.Errorf("post webhook #%d: %v", i, err) + return + } + defer resp.Body.Close() + statuses[i] = resp.StatusCode + }(i) + } + ready.Wait() + close(start) + wg.Wait() + + for i, code := range statuses { + if code != http.StatusOK { + t.Errorf("webhook #%d returned %d, want 200", i, code) + } + } + + var matched []any + for _, inc := range listIncidents(t, s, "") { + if inc["group_key"] == groupKey { + matched = append(matched, inc["id"]) + } + } + if len(matched) != 1 { + t.Fatalf("expected exactly 1 incident for group_key %q after %d concurrent deliveries, got %d: %v", + groupKey, callers, len(matched), matched) + } + + var alerts []map[string]any + decode(t, s.req(t, http.MethodGet, "/api/alerts", nil), &alerts) + seen := map[string]bool{} + for _, a := range alerts { + if fp, ok := a["fingerprint"].(string); ok { + seen[fp] = true + } + } + for i := range callers { + fp := fmt.Sprintf("fp-race-%d", i) + if !seen[fp] { + t.Errorf("alert %q is missing: its delivery's whole payload was silently rolled back "+ + "when it lost the race for the incident", fp) + } + } +}