Skip to content

Commit 3490efb

Browse files
authored
Merge pull request #601 from chiemezie1/feat/worker-leader-election
feat: add advisory-lock leader election for workers
2 parents 3c00dbb + 668d51e commit 3490efb

5 files changed

Lines changed: 370 additions & 31 deletions

File tree

WORKER_IMPLEMENTATION.md

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -121,6 +121,11 @@ Expected coverage: 95%+
121121
- Workers coordinate via shared store
122122
- Horizontal scaling supported
123123

124+
### Worker leader election
125+
- Singleton jobs now use PostgreSQL advisory locks so only one instance can run a given job at a time.
126+
- The partition rollover and statement archive jobs wrap each run in a leader guard that releases the lock on shutdown and after the run completes.
127+
- A Prometheus gauge named `worker_leader_status{job="..."}` reports whether the current process holds the leader lock for each singleton job.
128+
124129
## Production Deployment
125130

126131
### Database Integration

internal/worker/leader.go

Lines changed: 166 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,166 @@
1+
package worker
2+
3+
import (
4+
"context"
5+
"database/sql"
6+
"errors"
7+
"fmt"
8+
"sync"
9+
10+
"github.com/prometheus/client_golang/prometheus"
11+
"github.com/prometheus/client_golang/prometheus/promauto"
12+
)
13+
14+
// LeaderLocker exposes PostgreSQL advisory-lock behavior for singleton worker jobs.
15+
type LeaderLocker interface {
16+
AcquireLock(ctx context.Context, key int64) (bool, error)
17+
ReleaseLock(ctx context.Context, key int64) error
18+
Close() error
19+
}
20+
21+
// postgresLeaderLocker holds a dedicated SQL connection so advisory locks can be
22+
// released by the same session that acquired them.
23+
type postgresLeaderLocker struct {
24+
conn *sql.Conn
25+
}
26+
27+
func NewPostgresLeaderLocker(db *sql.DB) (LeaderLocker, error) {
28+
if db == nil {
29+
return nil, errors.New("database connection required for leader election")
30+
}
31+
32+
conn, err := db.Conn(context.Background())
33+
if err != nil {
34+
return nil, fmt.Errorf("acquire leader-election connection: %w", err)
35+
}
36+
37+
return &postgresLeaderLocker{conn: conn}, nil
38+
}
39+
40+
func (l *postgresLeaderLocker) AcquireLock(ctx context.Context, key int64) (bool, error) {
41+
if l == nil || l.conn == nil {
42+
return false, errors.New("leader-election connection unavailable")
43+
}
44+
45+
var acquired bool
46+
err := l.conn.QueryRowContext(ctx, "SELECT pg_try_advisory_lock($1)", key).Scan(&acquired)
47+
if err != nil {
48+
return false, fmt.Errorf("acquire advisory lock: %w", err)
49+
}
50+
return acquired, nil
51+
}
52+
53+
func (l *postgresLeaderLocker) ReleaseLock(ctx context.Context, key int64) error {
54+
if l == nil || l.conn == nil {
55+
return nil
56+
}
57+
58+
var released bool
59+
err := l.conn.QueryRowContext(ctx, "SELECT pg_advisory_unlock($1)", key).Scan(&released)
60+
if err != nil {
61+
return fmt.Errorf("release advisory lock: %w", err)
62+
}
63+
_ = released
64+
return nil
65+
}
66+
67+
func (l *postgresLeaderLocker) Close() error {
68+
if l == nil || l.conn == nil {
69+
return nil
70+
}
71+
conn := l.conn
72+
l.conn = nil
73+
return conn.Close()
74+
}
75+
76+
var (
77+
leaderStatusMetricsOnce sync.Once
78+
leaderStatusMetrics *prometheus.GaugeVec
79+
)
80+
81+
func initLeaderStatusMetrics() {
82+
leaderStatusMetricsOnce.Do(func() {
83+
leaderStatusMetrics = promauto.NewGaugeVec(prometheus.GaugeOpts{
84+
Name: "leader_status",
85+
Help: "Whether a singleton worker job currently holds the advisory lock.",
86+
}, []string{"job"})
87+
})
88+
}
89+
90+
func setLeaderStatus(job string, value float64) {
91+
initLeaderStatusMetrics()
92+
leaderStatusMetrics.WithLabelValues(job).Set(value)
93+
}
94+
95+
var errNotLeader = errors.New("not leader")
96+
97+
type leaderGuard struct {
98+
locker LeaderLocker
99+
job string
100+
key int64
101+
102+
mu sync.Mutex
103+
hold bool
104+
}
105+
106+
func newLeaderGuard(locker LeaderLocker, job string, key int64) *leaderGuard {
107+
return &leaderGuard{locker: locker, job: job, key: key}
108+
}
109+
110+
func (g *leaderGuard) Run(ctx context.Context, fn func(context.Context) error) error {
111+
if g == nil || g.locker == nil {
112+
return fn(ctx)
113+
}
114+
115+
acquired, err := g.locker.AcquireLock(ctx, g.key)
116+
if err != nil {
117+
setLeaderStatus(g.job, 0)
118+
return fmt.Errorf("acquire leader lock for %s: %w", g.job, err)
119+
}
120+
if !acquired {
121+
setLeaderStatus(g.job, 0)
122+
return errNotLeader
123+
}
124+
125+
g.mu.Lock()
126+
g.hold = true
127+
g.mu.Unlock()
128+
setLeaderStatus(g.job, 1)
129+
130+
defer func() {
131+
g.mu.Lock()
132+
hold := g.hold
133+
g.hold = false
134+
g.mu.Unlock()
135+
if hold {
136+
setLeaderStatus(g.job, 0)
137+
_ = g.locker.ReleaseLock(context.Background(), g.key)
138+
}
139+
}()
140+
141+
return fn(ctx)
142+
}
143+
144+
func (g *leaderGuard) Release(ctx context.Context) {
145+
if g == nil || g.locker == nil {
146+
return
147+
}
148+
149+
g.mu.Lock()
150+
if !g.hold {
151+
g.mu.Unlock()
152+
return
153+
}
154+
g.hold = false
155+
g.mu.Unlock()
156+
157+
setLeaderStatus(g.job, 0)
158+
_ = g.locker.ReleaseLock(ctx, g.key)
159+
}
160+
161+
func (g *leaderGuard) Close() error {
162+
if g == nil || g.locker == nil {
163+
return nil
164+
}
165+
return g.locker.Close()
166+
}

internal/worker/leader_test.go

Lines changed: 83 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,83 @@
1+
package worker
2+
3+
import (
4+
"context"
5+
"errors"
6+
"testing"
7+
)
8+
9+
type fakeLeaderLocker struct {
10+
acquired bool
11+
acquireResult bool
12+
acquireErr error
13+
releaseErr error
14+
releases int
15+
}
16+
17+
func (f *fakeLeaderLocker) AcquireLock(ctx context.Context, key int64) (bool, error) {
18+
if f.acquireErr != nil {
19+
return false, f.acquireErr
20+
}
21+
f.acquired = f.acquireResult
22+
return f.acquireResult, nil
23+
}
24+
25+
func (f *fakeLeaderLocker) ReleaseLock(ctx context.Context, key int64) error {
26+
if f.releaseErr != nil {
27+
return f.releaseErr
28+
}
29+
f.releases++
30+
f.acquired = false
31+
return nil
32+
}
33+
34+
func (f *fakeLeaderLocker) Close() error {
35+
return nil
36+
}
37+
38+
func TestLeaderLockerAcquireAndRelease(t *testing.T) {
39+
locker := &fakeLeaderLocker{acquireResult: true}
40+
41+
acquired, err := locker.AcquireLock(context.Background(), 123)
42+
if err != nil {
43+
t.Fatalf("AcquireLock returned error: %v", err)
44+
}
45+
if !acquired {
46+
t.Fatal("expected lock acquisition to succeed")
47+
}
48+
49+
if err := locker.ReleaseLock(context.Background(), 123); err != nil {
50+
t.Fatalf("ReleaseLock returned error: %v", err)
51+
}
52+
if locker.releases != 1 {
53+
t.Fatalf("expected one release, got %d", locker.releases)
54+
}
55+
}
56+
57+
func TestLeaderGuardRunWhenLockIsNotAvailable(t *testing.T) {
58+
locker := &fakeLeaderLocker{acquireResult: false}
59+
guard := newLeaderGuard(locker, "test-job", 42)
60+
61+
err := guard.Run(context.Background(), func(ctx context.Context) error {
62+
t.Fatal("callback should not run when leadership is not acquired")
63+
return nil
64+
})
65+
if !errors.Is(err, errNotLeader) {
66+
t.Fatalf("expected errNotLeader, got %v", err)
67+
}
68+
}
69+
70+
func TestLeaderGuardRunReleasesLockAfterCallback(t *testing.T) {
71+
locker := &fakeLeaderLocker{acquireResult: true}
72+
guard := newLeaderGuard(locker, "test-job", 42)
73+
74+
err := guard.Run(context.Background(), func(ctx context.Context) error {
75+
return nil
76+
})
77+
if err != nil {
78+
t.Fatalf("Run returned unexpected error: %v", err)
79+
}
80+
if locker.releases != 1 {
81+
t.Fatalf("expected one release after callback, got %d", locker.releases)
82+
}
83+
}

0 commit comments

Comments
 (0)