From 9e5b085d8bf0b0847d44071b9afef99c5ae905d0 Mon Sep 17 00:00:00 2001 From: Niklas Ye Date: Sat, 3 Oct 2026 11:58:39 +0200 Subject: [PATCH] Resolve the new-incident insert conflict instead of dropping the payload MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit openIncident's INSERT had no ON CONFLICT clause, relying entirely on incidentForGroup's earlier SELECT to avoid a duplicate. On more than one replica, two webhook deliveries for the very first occurrence of a brand-new groupKey can both pass that SELECT before either INSERTs; the loser then hit incidents_open_group_key_idx's unique violation, which rolled back its whole transaction — including that payload's alert upserts, done earlier in the same transaction. ingest's error is only logged and receiveWebhook answers 200 regardless, so nothing retried it: the loser's alerts silently never existed. Add ON CONFLICT (team_id, group_key) WHERE resolved_at IS NULL DO NOTHING to the INSERT, matching the partial unique index. Postgres only resolves that conflict after the winning transaction commits (or rolls back), so by the time RETURNING comes back empty, existingOpenIncident's follow-up SELECT is guaranteed to see the winner's row. The loser attaches to it instead of failing outright, and the rest of its payload commits normally. Covers both callers, since the dead man's switch sweeper shares this same function. New test (package api_test, fires N webhook deliveries for one groupKey from a synchronized start with distinct fingerprints, so they aren't accidentally serialized by upsertAlerts' own per- fingerprint lock) confirmed meaningful: with the ON CONFLICT clause reverted, it fails 10/10 on a missing alert fingerprint; restored, 0/10. Note while building it: "exactly one incident" alone cannot distinguish fixed from broken, since the DB's own unique index already guarantees that either way — the real signal is the loser's payload surviving. Chart comment updated: all three of the chart's original reasons for Recreate are now addressed in code, though replicas stays at 1 and the strategy stays Recreate pending a deliberate decision to raise it. Co-authored-by: Claude --- .../terdut-server/templates/deployment.yaml | 13 +- internal/api/alertmanager.go | 36 +++++- internal/api/alertmanager_test.go | 112 ++++++++++++++++++ 3 files changed, 155 insertions(+), 6 deletions(-) create mode 100644 internal/api/alertmanager_test.go 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) + } + } +}