Compare commits
8 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 5c4e0bdd0e | |||
| 9e5b085d8b | |||
| 0050738ca0 | |||
| 42180948d1 | |||
| c83c7c2a8b | |||
| 43beda9a30 | |||
| 710521a73c | |||
| d9492913ed |
@@ -15,5 +15,5 @@ type: application
|
||||
# appVersion and image.tag in values.yaml no longer agree, and that is not an oversight:
|
||||
# image.tag stays "latest", which is what a local install actually pulls. appVersion is
|
||||
# metadata and drives nothing.
|
||||
version: 0.35.0
|
||||
appVersion: "v0.35.0"
|
||||
version: 0.36.0
|
||||
appVersion: "v0.36.0"
|
||||
|
||||
@@ -10,9 +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 and the notifier are unsynchronised singletons, and two replicas
|
||||
# overlapping during a rollout would both page for the same incident.
|
||||
# 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:
|
||||
|
||||
@@ -0,0 +1,122 @@
|
||||
package api
|
||||
|
||||
// This file is internal (package api, not api_test) because withAdvisoryLock is
|
||||
// unexported and these tests exercise its locking semantics directly rather than
|
||||
// through the full StartArchiver/StartNotifier loop, which would make the "does
|
||||
// not run while held" case timing-dependent instead of deterministic. It opens a
|
||||
// plain connection to TERDUT_TEST_DSN rather than reusing testdb_test.go's
|
||||
// newTestDB, since that helper lives in the separate, already-compiled
|
||||
// api_test package and a Postgres advisory lock needs no schema or migration
|
||||
// to exercise.
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
_ "github.com/jackc/pgx/v5/stdlib"
|
||||
)
|
||||
|
||||
// advisoryTestDB opens a plain, unmigrated connection to the test database. An
|
||||
// unset DSN fails rather than skips, matching testdb_test.go's rationale: a
|
||||
// suite that quietly tests nothing is worse than one that does not run.
|
||||
func advisoryTestDB(t *testing.T) *sql.DB {
|
||||
t.Helper()
|
||||
|
||||
dsn := os.Getenv("TERDUT_TEST_DSN")
|
||||
if dsn == "" {
|
||||
t.Fatalf("TERDUT_TEST_DSN is not set: these tests need Postgres.\n" +
|
||||
"Run `make test-db` for a local one, then\n" +
|
||||
" export TERDUT_TEST_DSN=postgres://terdut:terdut@localhost:5432/terdut_test?sslmode=disable")
|
||||
}
|
||||
|
||||
db, err := sql.Open("pgx", dsn)
|
||||
if err != nil {
|
||||
t.Fatalf("connect to TERDUT_TEST_DSN: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { db.Close() })
|
||||
return db
|
||||
}
|
||||
|
||||
func TestWithAdvisoryLock_RunsWhenFree(t *testing.T) {
|
||||
db := advisoryTestDB(t)
|
||||
ctx := context.Background()
|
||||
|
||||
ran := false
|
||||
withAdvisoryLock(ctx, db, archiverLockKey, "test", func() { ran = true })
|
||||
|
||||
if !ran {
|
||||
t.Fatal("fn did not run although the lock was free")
|
||||
}
|
||||
}
|
||||
|
||||
func TestWithAdvisoryLock_SkipsWhileHeldElsewhere(t *testing.T) {
|
||||
db := advisoryTestDB(t)
|
||||
ctx := context.Background()
|
||||
|
||||
// Hold the lock on a connection of our own, standing in for another
|
||||
// replica mid-pass.
|
||||
holder, err := db.Conn(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("acquire holder connection: %v", err)
|
||||
}
|
||||
defer holder.Close()
|
||||
if _, err := holder.ExecContext(ctx, "SELECT pg_advisory_lock($1)", archiverLockKey); err != nil {
|
||||
t.Fatalf("pre-acquire lock: %v", err)
|
||||
}
|
||||
|
||||
ran := false
|
||||
withAdvisoryLock(ctx, db, archiverLockKey, "test", func() { ran = true })
|
||||
if ran {
|
||||
t.Fatal("fn ran although another connection already held the lock")
|
||||
}
|
||||
|
||||
if _, err := holder.ExecContext(ctx, "SELECT pg_advisory_unlock($1)", archiverLockKey); err != nil {
|
||||
t.Fatalf("release held lock: %v", err)
|
||||
}
|
||||
|
||||
// Now that the holder released it, the next caller should get it.
|
||||
ran = false
|
||||
withAdvisoryLock(ctx, db, archiverLockKey, "test", func() { ran = true })
|
||||
if !ran {
|
||||
t.Fatal("fn did not run after the other connection released the lock")
|
||||
}
|
||||
}
|
||||
|
||||
func TestWithAdvisoryLock_ReleasesAfterFnReturns(t *testing.T) {
|
||||
db := advisoryTestDB(t)
|
||||
ctx := context.Background()
|
||||
|
||||
withAdvisoryLock(ctx, db, notifierLockKey, "test", func() {})
|
||||
|
||||
// If the first call had leaked the lock, this one would see it held and
|
||||
// skip, leaving ran false.
|
||||
ran := false
|
||||
withAdvisoryLock(ctx, db, notifierLockKey, "test", func() { ran = true })
|
||||
if !ran {
|
||||
t.Fatal("fn did not run on a later call: the earlier call leaked its lock")
|
||||
}
|
||||
}
|
||||
|
||||
func TestWithAdvisoryLock_KeysAreIndependent(t *testing.T) {
|
||||
db := advisoryTestDB(t)
|
||||
ctx := context.Background()
|
||||
|
||||
holder, err := db.Conn(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("acquire holder connection: %v", err)
|
||||
}
|
||||
defer holder.Close()
|
||||
if _, err := holder.ExecContext(ctx, "SELECT pg_advisory_lock($1)", archiverLockKey); err != nil {
|
||||
t.Fatalf("pre-acquire archiver lock: %v", err)
|
||||
}
|
||||
defer holder.ExecContext(ctx, "SELECT pg_advisory_unlock($1)", archiverLockKey)
|
||||
|
||||
// Holding archiverLockKey must not block notifierLockKey.
|
||||
ran := false
|
||||
withAdvisoryLock(ctx, db, notifierLockKey, "test", func() { ran = true })
|
||||
if !ran {
|
||||
t.Fatal("fn did not run under a different key although only archiverLockKey was held")
|
||||
}
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -15,6 +15,12 @@ const (
|
||||
// expiryGrace absorbs clock skew and notification latency before an alert
|
||||
// whose ends_at watermark has passed is treated as stale.
|
||||
expiryGrace = 5 * time.Minute
|
||||
|
||||
// archiverLockKey is the Postgres advisory lock the sweeper takes for the
|
||||
// duration of each pass, so that running more than one replica does not run
|
||||
// the sweep concurrently on all of them. Its value has no meaning beyond
|
||||
// being distinct from notifierLockKey.
|
||||
archiverLockKey int64 = 7265_0001
|
||||
)
|
||||
|
||||
// StartArchiver runs the alert sweeper until ctx is cancelled, starting with an
|
||||
@@ -23,15 +29,25 @@ const (
|
||||
// the fallback, not the setting: each pass reads the current value from the
|
||||
// settings table, so an administrator's change takes effect on the next tick
|
||||
// instead of at the next restart.
|
||||
//
|
||||
// Each pass runs under archiverLockKey (see withAdvisoryLock), so that on more
|
||||
// than one replica only whichever instance's tick takes the lock first actually
|
||||
// sweeps; the rest skip that tick rather than racing the same pass.
|
||||
func StartArchiver(ctx context.Context, db *sql.DB, archiveAfter, staleAfter time.Duration, notify NotifyConfig) {
|
||||
ticker := time.NewTicker(sweepInterval)
|
||||
defer ticker.Stop()
|
||||
|
||||
Sweep(ctx, db, archiveAfter, staleAfter, notify)
|
||||
sweep := func() {
|
||||
withAdvisoryLock(ctx, db, archiverLockKey, "sweeper", func() {
|
||||
Sweep(ctx, db, archiveAfter, staleAfter, notify)
|
||||
})
|
||||
}
|
||||
|
||||
sweep()
|
||||
for {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
Sweep(ctx, db, archiveAfter, staleAfter, notify)
|
||||
sweep()
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
|
||||
@@ -1,8 +1,11 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"log"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
@@ -77,3 +80,39 @@ func decodeJSON(r *http.Request, v any) error {
|
||||
func errResp(msg string) map[string]string {
|
||||
return map[string]string{"error": msg}
|
||||
}
|
||||
|
||||
// withAdvisoryLock runs fn only if it can take the named Postgres advisory lock on a
|
||||
// dedicated connection, and skips fn otherwise. This is what keeps the archiver and
|
||||
// notifier safe to run on more than one replica: whichever instance's tick gets there
|
||||
// first does the work; the rest see the lock held and simply wait for their next tick
|
||||
// instead of running the same pass concurrently.
|
||||
//
|
||||
// pg_try_advisory_lock is session-scoped, so taking and releasing it must happen on the
|
||||
// same connection, reserved via db.Conn rather than borrowed from the pool's shared
|
||||
// connections fn itself may use — and released (unlocked, then closed) before returning,
|
||||
// since a session lock otherwise outlives this call and leaks onto whatever reuses the
|
||||
// pooled connection next.
|
||||
func withAdvisoryLock(ctx context.Context, db *sql.DB, key int64, name string, fn func()) {
|
||||
conn, err := db.Conn(ctx)
|
||||
if err != nil {
|
||||
log.Printf("%s: advisory lock: acquire connection: %v", name, err)
|
||||
return
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
var locked bool
|
||||
if err := conn.QueryRowContext(ctx, "SELECT pg_try_advisory_lock($1)", key).Scan(&locked); err != nil {
|
||||
log.Printf("%s: advisory lock: %v", name, err)
|
||||
return
|
||||
}
|
||||
if !locked {
|
||||
return // another replica is already running this pass
|
||||
}
|
||||
defer func() {
|
||||
if _, err := conn.ExecContext(ctx, "SELECT pg_advisory_unlock($1)", key); err != nil {
|
||||
log.Printf("%s: advisory unlock: %v", name, err)
|
||||
}
|
||||
}()
|
||||
|
||||
fn()
|
||||
}
|
||||
|
||||
@@ -38,6 +38,13 @@ const (
|
||||
ackTokenTTL = 24 * time.Hour
|
||||
)
|
||||
|
||||
// notifierLockKey is the Postgres advisory lock the notifier takes for the
|
||||
// duration of each pass, so that running more than one replica does not
|
||||
// deliver (or double-deliver) the same notification from more than one of
|
||||
// them at once. Its value has no meaning beyond being distinct from
|
||||
// archiverLockKey.
|
||||
const notifierLockKey int64 = 7265_0002
|
||||
|
||||
// Notification kinds, recording why a push was sent.
|
||||
const (
|
||||
notifyTriggered = "triggered"
|
||||
@@ -97,6 +104,10 @@ var notifyClient = &http.Client{Timeout: 10 * time.Second}
|
||||
|
||||
// StartNotifier delivers queued notifications until ctx is cancelled, starting
|
||||
// with an immediate pass so a restart flushes whatever the last one left behind.
|
||||
//
|
||||
// Each pass runs under notifierLockKey (see withAdvisoryLock), so that on more
|
||||
// than one replica only whichever instance's tick takes the lock first actually
|
||||
// delivers; the rest skip that tick rather than racing the same pass.
|
||||
func StartNotifier(ctx context.Context, db *sql.DB, cfg NotifyConfig) {
|
||||
if !cfg.enabled() {
|
||||
log.Print("notifier: disabled (no ntfy URL configured)")
|
||||
@@ -107,11 +118,17 @@ func StartNotifier(ctx context.Context, db *sql.DB, cfg NotifyConfig) {
|
||||
ticker := time.NewTicker(notifyInterval)
|
||||
defer ticker.Stop()
|
||||
|
||||
NotifySweep(ctx, db, cfg)
|
||||
sweep := func() {
|
||||
withAdvisoryLock(ctx, db, notifierLockKey, "notifier", func() {
|
||||
NotifySweep(ctx, db, cfg)
|
||||
})
|
||||
}
|
||||
|
||||
sweep()
|
||||
for {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
NotifySweep(ctx, db, cfg)
|
||||
sweep()
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"embed"
|
||||
"fmt"
|
||||
@@ -62,6 +63,16 @@ func Open(dsn string) (*sql.DB, error) {
|
||||
}
|
||||
}
|
||||
|
||||
// migrationLockKey is the Postgres advisory lock Migrate holds for its whole
|
||||
// run. Two replicas starting at once would otherwise race the check-then-apply
|
||||
// loop below against schema_migrations: the loser could crash on a
|
||||
// duplicate-key insert, or contend with the winner's uncommitted DDL. Blocking
|
||||
// (pg_advisory_lock, not pg_try_advisory_lock as the archiver and notifier
|
||||
// use): on boot there is no later tick to defer to, so the right behaviour is
|
||||
// to wait for the other replica to finish migrating, not to skip ahead and
|
||||
// start serving against an unmigrated schema.
|
||||
const migrationLockKey int64 = 7265_0003
|
||||
|
||||
// Migrate applies every embedded migration that has not been applied yet, in
|
||||
// filename order, recording each in schema_migrations.
|
||||
//
|
||||
@@ -69,6 +80,22 @@ func Open(dsn string) (*sql.DB, error) {
|
||||
// migration that failed half way used to leave the schema in whatever state it
|
||||
// had reached. Postgres has transactional DDL, so the rollback is real.
|
||||
func Migrate(db *sql.DB) error {
|
||||
ctx := context.Background()
|
||||
conn, err := db.Conn(ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("migrate: acquire connection: %w", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
if _, err := conn.ExecContext(ctx, "SELECT pg_advisory_lock($1)", migrationLockKey); err != nil {
|
||||
return fmt.Errorf("migrate: acquire advisory lock: %w", err)
|
||||
}
|
||||
defer func() {
|
||||
if _, err := conn.ExecContext(ctx, "SELECT pg_advisory_unlock($1)", migrationLockKey); err != nil {
|
||||
log.Printf("migrate: release advisory lock: %v", err)
|
||||
}
|
||||
}()
|
||||
|
||||
if _, err := db.Exec(`CREATE TABLE IF NOT EXISTS schema_migrations (
|
||||
version TEXT PRIMARY KEY,
|
||||
applied_at BIGINT NOT NULL DEFAULT FLOOR(EXTRACT(EPOCH FROM now()))::bigint
|
||||
|
||||
@@ -0,0 +1,125 @@
|
||||
package db_test
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
"git.ryuvia.com/niklas/terdut-server/internal/db"
|
||||
|
||||
_ "github.com/jackc/pgx/v5/stdlib"
|
||||
)
|
||||
|
||||
// TERDUT_TEST_DSN must point at a database the test role may create schemas
|
||||
// in; see internal/api/testdb_test.go for the fuller rationale this mirrors.
|
||||
// An unset DSN fails rather than skips, deliberately.
|
||||
const testDSNEnv = "TERDUT_TEST_DSN"
|
||||
|
||||
// TestMigrate_ConcurrentCallersDoNotRace reproduces two replicas starting at
|
||||
// once against a brand-new, unmigrated schema: both call db.Migrate at the
|
||||
// same time. Before migrationLockKey, the loser could crash on a
|
||||
// duplicate-key insert into schema_migrations, or contend with the winner's
|
||||
// uncommitted DDL; with the advisory lock, one blocks until the other
|
||||
// finishes and both return cleanly.
|
||||
func TestMigrate_ConcurrentCallersDoNotRace(t *testing.T) {
|
||||
dsn := os.Getenv(testDSNEnv)
|
||||
if dsn == "" {
|
||||
t.Fatalf("%s is not set: these tests need Postgres.\n"+
|
||||
"Run `make test-db` for a local one, then\n"+
|
||||
" export %s=postgres://terdut:terdut@localhost:5432/terdut_test?sslmode=disable",
|
||||
testDSNEnv, testDSNEnv)
|
||||
}
|
||||
|
||||
schema := fmt.Sprintf("migrate_race_%d", os.Getpid())
|
||||
admin, err := sql.Open("pgx", dsn)
|
||||
if err != nil {
|
||||
t.Fatalf("connect to %s: %v", testDSNEnv, err)
|
||||
}
|
||||
defer admin.Close()
|
||||
if _, err := admin.Exec("CREATE SCHEMA " + schema); err != nil {
|
||||
t.Fatalf("create schema %s: %v", schema, err)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
cleanup, err := sql.Open("pgx", dsn)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
defer cleanup.Close()
|
||||
if _, err := cleanup.Exec("DROP SCHEMA " + schema + " CASCADE"); err != nil {
|
||||
t.Logf("drop schema %s: %v", schema, err)
|
||||
}
|
||||
})
|
||||
|
||||
scoped := withSearchPath(dsn, schema)
|
||||
|
||||
const callers = 2
|
||||
errs := make([]error, callers)
|
||||
var wg sync.WaitGroup
|
||||
for i := range callers {
|
||||
wg.Add(1)
|
||||
go func(i int) {
|
||||
defer wg.Done()
|
||||
database, err := db.Open(scoped)
|
||||
if err != nil {
|
||||
errs[i] = fmt.Errorf("open: %w", err)
|
||||
return
|
||||
}
|
||||
defer database.Close()
|
||||
errs[i] = db.Migrate(database)
|
||||
}(i)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
for i, err := range errs {
|
||||
if err != nil {
|
||||
t.Fatalf("Migrate #%d: %v", i, err)
|
||||
}
|
||||
}
|
||||
|
||||
entries, err := os.ReadDir("migrations")
|
||||
if err != nil {
|
||||
t.Fatalf("read migrations dir: %v", err)
|
||||
}
|
||||
var want int
|
||||
for _, e := range entries {
|
||||
if !e.IsDir() && strings.HasSuffix(e.Name(), ".sql") {
|
||||
want++
|
||||
}
|
||||
}
|
||||
|
||||
check, err := sql.Open("pgx", scoped)
|
||||
if err != nil {
|
||||
t.Fatalf("connect for verification: %v", err)
|
||||
}
|
||||
defer check.Close()
|
||||
|
||||
var got int
|
||||
if err := check.QueryRow("SELECT COUNT(*) FROM schema_migrations").Scan(&got); err != nil {
|
||||
t.Fatalf("count schema_migrations: %v", err)
|
||||
}
|
||||
if got != want {
|
||||
t.Fatalf("schema_migrations has %d row(s) after two concurrent Migrate calls, want %d (one per migration file, no duplicates)", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
// withSearchPath pins a DSN to one schema. Copied from
|
||||
// internal/api/testdb_test.go rather than shared: that helper lives in the
|
||||
// api_test package, a separate compiled package this one cannot import.
|
||||
func withSearchPath(dsn, schema string) string {
|
||||
opt := "-csearch_path=" + schema
|
||||
|
||||
if strings.HasPrefix(dsn, "postgres://") || strings.HasPrefix(dsn, "postgresql://") {
|
||||
u, err := url.Parse(dsn)
|
||||
if err == nil {
|
||||
q := u.Query()
|
||||
q.Set("options", opt)
|
||||
u.RawQuery = q.Encode()
|
||||
return u.String()
|
||||
}
|
||||
}
|
||||
return dsn + " options='" + opt + "'"
|
||||
}
|
||||
@@ -326,6 +326,7 @@ input:focus, textarea:focus { outline: none; border-color: var(--accent); box-sh
|
||||
/* ---------- chips ---------- */
|
||||
|
||||
.chips {
|
||||
position: relative;
|
||||
display: flex; gap: 6px;
|
||||
padding: 12px 16px 8px;
|
||||
overflow-x: auto; scrollbar-width: none;
|
||||
@@ -341,11 +342,14 @@ input:focus, textarea:focus { outline: none; border-color: var(--accent); box-sh
|
||||
}
|
||||
.chip[aria-selected="true"] { background: var(--text); border-color: var(--text); color: var(--bg); }
|
||||
.chip .count { margin-left: 4px; opacity: 0.7; }
|
||||
/* Pinned to the visible right edge of the scrolling row (sticky, not
|
||||
absolute, so it tracks the scroll position rather than the content). */
|
||||
/* An overlay, not a flex item: absolute against .chips' own (non-scrolling)
|
||||
box stays flush with its real right edge regardless of scroll position,
|
||||
which turned out not to be true of position:sticky here — as a flex
|
||||
item, its sticky offset interacted with the row's gap and its own
|
||||
negative margin, landing short of the edge by about one gap's width. */
|
||||
.chips-fade {
|
||||
position: sticky; right: -1px; flex: none;
|
||||
width: 24px; margin-left: -24px;
|
||||
position: absolute; top: 0; right: 0; bottom: 0;
|
||||
width: 24px;
|
||||
background: linear-gradient(to right, transparent, var(--bg));
|
||||
pointer-events: none;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user