Skip to content

Commit 7b969fd

Browse files
gustavobertoiclaude
andcommitted
feat(lock): distributed lock seam + pg_advisory_lock for a remote backend (spec 21)
The concurrency spine's local gofrs/flock cannot serialize two developers on two machines mutating the SAME remote backend's ledger rows / provisioning / port allocation — the central unsolved problem for a team/cloud shared backend (spec 21). This lands the correctness gate: a Locker seam whose remote implementation is a session-scoped pg_advisory_lock on the shared cluster Postgres (Q-REMOTE-LOCK RESOLVED — the DB we already run is the coordinator, so zero new infra and no daemon; session-scoped means it auto-releases on a crash so another machine reconciles cleanly). - internal/lock/distlock.go: the Locker interface, FileLocker (today's flock, verbatim), and LockerFor(remote, subject, path, connect) — remote+connector → PGLocker, else the local flock (a remote backend with no reachable cluster yet degrades safely to single-machine correctness). internal/lock stays a dependency-free leaf (plain bool, no docker import → no cycle). - internal/lock/pglock.go: PGLocker takes a session-scoped pg_advisory_lock keyed by AdvisoryKey (FNV-64a subject hash → deterministic bigint, so every client of one cluster hashes the same subject to the same key). It polls pg_try_advisory_lock (not the blocking form) to honor ctx cancel/timeout, and opens a fresh session per WithLock so the lock never outlives its section. AdvisoryConn is the injectable session seam (transaction-pooling pgbouncer breaks session advisory locks — documented). - internal/provision/advisory.go: the pgx/v5-backed AdvisoryConn + PGLockConnector (the production session), kept out of internal/lock to preserve the leaf. Tests (offline, no live PG): AdvisoryKey determinism, acquire/run/release order, fn-error release, connect/try errors, ctx-timeout on a held key, LockerFor selection, FileLocker, and a 50-goroutine mutual-exclusion test over an in-memory advisory server (compare-and-set) run under -race. This is the foundation; wiring the orchestrator's ~25 lock call sites through the selected Locker + reaching the remote cluster over an SSH forward are the follow-ups. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1 parent 2d9cec1 commit 7b969fd

5 files changed

Lines changed: 437 additions & 0 deletions

File tree

‎internal/lock/distlock.go‎

Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,52 @@
1+
package lock
2+
3+
import "context"
4+
5+
// Locker serializes every mutation of the machine-global ledger or the shared
6+
// Docker stack. It is the seam that lets the concurrency spine span a REMOTE
7+
// team backend (spec 21): the local implementation is the gofrs/flock advisory
8+
// lock (a single machine); the remote implementation is a session-scoped
9+
// pg_advisory_lock on the shared cluster Postgres, which serializes across
10+
// MACHINES against one backend.
11+
//
12+
// This exists because a local flock cannot serialize two developers on two
13+
// laptops mutating the SAME remote ledger rows / provisioning / port allocation
14+
// (spec 21 §"the central unsolved problem"). Q-REMOTE-LOCK is RESOLVED: use
15+
// pg_advisory_lock on the cluster DB as the coordinator — zero new infra,
16+
// crash-safe (session-scoped → auto-released on disconnect), reusing pgx.
17+
type Locker interface {
18+
// WithLock runs fn while holding the lock, releasing it before returning
19+
// (on success or error). It blocks until the lock is held or ctx is done.
20+
WithLock(ctx context.Context, fn func() error) error
21+
}
22+
23+
// FileLocker is the LOCAL Locker: the coarse gofrs/flock at a lockfile path
24+
// (today's behavior, verbatim). An empty Path runs fn unlocked (tests).
25+
type FileLocker struct{ Path string }
26+
27+
// NewFileLocker returns a FileLocker bound to a lockfile path (under XDG_RUNTIME_DIR).
28+
func NewFileLocker(path string) FileLocker { return FileLocker{Path: path} }
29+
30+
// WithLock takes the flock at Path, runs fn, and releases it.
31+
func (f FileLocker) WithLock(ctx context.Context, fn func() error) error {
32+
if f.Path == "" {
33+
return fn()
34+
}
35+
return WithLock(ctx, f.Path, fn)
36+
}
37+
38+
// LockerFor selects the concurrency primitive for a backend. A remote backend
39+
// with a reachable cluster Postgres (connect != nil) serializes across machines
40+
// on a pg_advisory_lock keyed by subject; everything else uses the local flock
41+
// at path. connect is nil until remote host-reachability is wired, so a remote
42+
// backend degrades safely to the local flock (single-machine correctness holds;
43+
// cross-machine serialization is the follow-up that needs the cluster reachable).
44+
//
45+
// The parameter is a plain bool, not a docker.Backend, so internal/lock stays a
46+
// dependency-free leaf (no import cycle).
47+
func LockerFor(remote bool, subject, path string, connect PGConnector) Locker {
48+
if remote && connect != nil {
49+
return NewPGLocker(subject, connect)
50+
}
51+
return NewFileLocker(path)
52+
}

‎internal/lock/distlock_test.go‎

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,45 @@
1+
package lock
2+
3+
import (
4+
"context"
5+
"path/filepath"
6+
"testing"
7+
)
8+
9+
func TestFileLocker_RunsFn(t *testing.T) {
10+
l := NewFileLocker(filepath.Join(t.TempDir(), "x.lock"))
11+
ran := false
12+
if err := l.WithLock(context.Background(), func() error { ran = true; return nil }); err != nil {
13+
t.Fatal(err)
14+
}
15+
if !ran {
16+
t.Fatal("fn did not run under the file lock")
17+
}
18+
}
19+
20+
func TestFileLocker_EmptyPathRunsUnlocked(t *testing.T) {
21+
ran := false
22+
if err := (FileLocker{}).WithLock(context.Background(), func() error { ran = true; return nil }); err != nil {
23+
t.Fatal(err)
24+
}
25+
if !ran {
26+
t.Fatal("empty-path FileLocker should run fn unlocked")
27+
}
28+
}
29+
30+
func TestLockerFor_Selects(t *testing.T) {
31+
connect := func(context.Context) (AdvisoryConn, error) { return nil, nil }
32+
33+
// Remote + a connector → distributed pg lock.
34+
if _, ok := LockerFor(true, "subj", "/tmp/x.lock", connect).(*PGLocker); !ok {
35+
t.Error("remote + connector should select *PGLocker")
36+
}
37+
// Remote but no connector (reachability not wired) → safe local fallback.
38+
if _, ok := LockerFor(true, "subj", "/tmp/x.lock", nil).(FileLocker); !ok {
39+
t.Error("remote + nil connector should fall back to FileLocker")
40+
}
41+
// Local → file lock.
42+
if _, ok := LockerFor(false, "subj", "/tmp/x.lock", connect).(FileLocker); !ok {
43+
t.Error("local should select FileLocker")
44+
}
45+
}

‎internal/lock/pglock.go‎

Lines changed: 108 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,108 @@
1+
package lock
2+
3+
import (
4+
"context"
5+
"errors"
6+
"fmt"
7+
"hash/fnv"
8+
"time"
9+
)
10+
11+
// This file is the DISTRIBUTED lock (spec 21, Q-REMOTE-LOCK RESOLVED): a
12+
// session-scoped pg_advisory_lock on the shared cluster Postgres. It serializes
13+
// two developers on two machines against ONE remote backend — the guarantee a
14+
// local flock structurally cannot provide.
15+
//
16+
// Why pg_advisory_lock and not a coordinator daemon: the team Postgres we already
17+
// run IS the coordinator (zero new infra), the lock is session-scoped so it is
18+
// auto-released when the holder disconnects or crashes (a kill -9 on one machine
19+
// lets another machine's next command reconcile cleanly), and it reuses the pgx
20+
// path provisioning already uses. Gotcha (spec 21): advisory locks live on a
21+
// SESSION — a transaction-pooling pgbouncer silently breaks them, so the
22+
// connector must hand out a direct/session connection, not a pooled one.
23+
24+
// advisoryRetry is how often PGLocker re-polls pg_try_advisory_lock while waiting,
25+
// mirroring the flock's TryLockContext cadence.
26+
const advisoryRetry = 100 * time.Millisecond
27+
28+
// AdvisoryConn is the minimal session-connection surface PGLocker needs. It MUST
29+
// be a single dedicated Postgres session (never a pool / transaction-mode
30+
// pgbouncer). Injectable so the locker is unit-testable without a live server.
31+
type AdvisoryConn interface {
32+
// TryLock runs `SELECT pg_try_advisory_lock($1)` and reports whether the lock
33+
// was granted on this session (non-blocking).
34+
TryLock(ctx context.Context, key int64) (bool, error)
35+
// Unlock runs `SELECT pg_advisory_unlock($1)` for this session.
36+
Unlock(ctx context.Context, key int64) error
37+
// Close ends the session (a crash-safe backstop that releases any held lock).
38+
Close(ctx context.Context) error
39+
}
40+
41+
// PGConnector opens a fresh AdvisoryConn (a dedicated session). PGLocker opens
42+
// one per WithLock and closes it after, so the advisory lock never outlives the
43+
// critical section even if a release is missed.
44+
type PGConnector func(ctx context.Context) (AdvisoryConn, error)
45+
46+
// PGLocker is the distributed Locker. It hashes Subject to a deterministic
47+
// bigint key and takes a session-scoped pg_advisory_lock on it.
48+
type PGLocker struct {
49+
Subject string
50+
Connect PGConnector
51+
// retry is injectable for tests; 0 → advisoryRetry.
52+
retry time.Duration
53+
}
54+
55+
// NewPGLocker returns a PGLocker for a lock subject and a session connector.
56+
func NewPGLocker(subject string, connect PGConnector) *PGLocker {
57+
return &PGLocker{Subject: subject, Connect: connect}
58+
}
59+
60+
// AdvisoryKey maps a lock subject to the deterministic int64 pg_advisory_lock
61+
// key: FNV-64a over the subject, reinterpreted as a signed bigint. Deterministic
62+
// across machines and runs, so every client of one cluster hashes the same
63+
// subject to the same key. A hash collision only causes extra (safe)
64+
// serialization between two subjects — never a missed exclusion.
65+
func AdvisoryKey(subject string) int64 {
66+
h := fnv.New64a()
67+
_, _ = h.Write([]byte(subject))
68+
return int64(h.Sum64())
69+
}
70+
71+
// WithLock opens a dedicated session, polls pg_try_advisory_lock until the key is
72+
// held or ctx is done, runs fn, then releases the lock and closes the session.
73+
// Polling (rather than the blocking pg_advisory_lock) honors ctx cancellation and
74+
// timeouts instead of parking a backend indefinitely.
75+
func (p *PGLocker) WithLock(ctx context.Context, fn func() error) error {
76+
if p.Connect == nil {
77+
return errors.New("pg advisory lock: no session connector configured")
78+
}
79+
retry := p.retry
80+
if retry == 0 {
81+
retry = advisoryRetry
82+
}
83+
key := AdvisoryKey(p.Subject)
84+
85+
conn, err := p.Connect(ctx)
86+
if err != nil {
87+
return fmt.Errorf("open advisory-lock session: %w", err)
88+
}
89+
defer func() { _ = conn.Close(ctx) }()
90+
91+
for {
92+
ok, err := conn.TryLock(ctx, key)
93+
if err != nil {
94+
return fmt.Errorf("acquire advisory lock: %w", err)
95+
}
96+
if ok {
97+
break
98+
}
99+
select {
100+
case <-ctx.Done():
101+
return fmt.Errorf("timed out waiting for the shared-backend lock (subject %q); another machine may be holding it: %w", p.Subject, ctx.Err())
102+
case <-time.After(retry):
103+
}
104+
}
105+
// Held. Release (Unlock) runs before Close because deferred calls are LIFO.
106+
defer func() { _ = conn.Unlock(ctx, key) }()
107+
return fn()
108+
}

‎internal/lock/pglock_test.go‎

Lines changed: 173 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,173 @@
1+
package lock
2+
3+
import (
4+
"context"
5+
"errors"
6+
"sync"
7+
"testing"
8+
"time"
9+
)
10+
11+
func TestAdvisoryKeyDeterministic(t *testing.T) {
12+
subject := "devstack-shared-ledger"
13+
first, second := AdvisoryKey(subject), AdvisoryKey(subject)
14+
if first != second {
15+
t.Fatalf("same subject must hash to the same key (%d != %d)", first, second)
16+
}
17+
if AdvisoryKey("a") == AdvisoryKey("b") {
18+
t.Fatal("different subjects should (almost surely) differ")
19+
}
20+
}
21+
22+
// fakeAdvisory is an in-memory model of ONE cluster's advisory-lock table: a
23+
// single map[key]bool guarded by a real mutex, so concurrent PGLockers exercise
24+
// genuine mutual exclusion (compare-and-set), like the real server session.
25+
type fakeAdvisory struct {
26+
mu sync.Mutex
27+
held map[int64]bool
28+
}
29+
30+
func newFakeAdvisory() *fakeAdvisory { return &fakeAdvisory{held: map[int64]bool{}} }
31+
32+
// session is one connection to the fake server.
33+
type fakeSession struct {
34+
srv *fakeAdvisory
35+
closed bool
36+
// hooks to simulate failures
37+
tryErr error
38+
closeErr error
39+
}
40+
41+
func (s *fakeSession) TryLock(_ context.Context, key int64) (bool, error) {
42+
if s.tryErr != nil {
43+
return false, s.tryErr
44+
}
45+
s.srv.mu.Lock()
46+
defer s.srv.mu.Unlock()
47+
if s.srv.held[key] {
48+
return false, nil
49+
}
50+
s.srv.held[key] = true
51+
return true, nil
52+
}
53+
54+
func (s *fakeSession) Unlock(_ context.Context, key int64) error {
55+
s.srv.mu.Lock()
56+
defer s.srv.mu.Unlock()
57+
delete(s.srv.held, key)
58+
return nil
59+
}
60+
61+
func (s *fakeSession) Close(context.Context) error {
62+
s.closed = true
63+
return s.closeErr
64+
}
65+
66+
func (f *fakeAdvisory) connector() PGConnector {
67+
return func(context.Context) (AdvisoryConn, error) {
68+
return &fakeSession{srv: f}, nil
69+
}
70+
}
71+
72+
func TestPGLocker_AcquireRunRelease(t *testing.T) {
73+
srv := newFakeAdvisory()
74+
l := &PGLocker{Subject: "s", Connect: srv.connector()}
75+
ran := false
76+
if err := l.WithLock(context.Background(), func() error {
77+
ran = true
78+
// While fn runs, the key must be held.
79+
srv.mu.Lock()
80+
defer srv.mu.Unlock()
81+
if !srv.held[AdvisoryKey("s")] {
82+
t.Error("key should be held while fn runs")
83+
}
84+
return nil
85+
}); err != nil {
86+
t.Fatal(err)
87+
}
88+
if !ran {
89+
t.Fatal("fn did not run")
90+
}
91+
// Released after WithLock returns.
92+
if srv.held[AdvisoryKey("s")] {
93+
t.Error("key should be released after WithLock")
94+
}
95+
}
96+
97+
func TestPGLocker_PropagatesFnError(t *testing.T) {
98+
srv := newFakeAdvisory()
99+
l := &PGLocker{Subject: "s", Connect: srv.connector()}
100+
sentinel := errors.New("boom")
101+
if err := l.WithLock(context.Background(), func() error { return sentinel }); !errors.Is(err, sentinel) {
102+
t.Fatalf("want fn error propagated, got %v", err)
103+
}
104+
if srv.held[AdvisoryKey("s")] {
105+
t.Error("key must be released even when fn errors")
106+
}
107+
}
108+
109+
func TestPGLocker_NoConnector(t *testing.T) {
110+
l := &PGLocker{Subject: "s"}
111+
if err := l.WithLock(context.Background(), func() error { return nil }); err == nil {
112+
t.Fatal("want error when no connector is configured")
113+
}
114+
}
115+
116+
func TestPGLocker_ConnectError(t *testing.T) {
117+
l := &PGLocker{Subject: "s", Connect: func(context.Context) (AdvisoryConn, error) {
118+
return nil, errors.New("dial refused")
119+
}}
120+
ran := false
121+
err := l.WithLock(context.Background(), func() error { ran = true; return nil })
122+
if err == nil || ran {
123+
t.Fatalf("connect error must abort before fn (err=%v ran=%v)", err, ran)
124+
}
125+
}
126+
127+
func TestPGLocker_CtxTimeoutWhenHeld(t *testing.T) {
128+
srv := newFakeAdvisory()
129+
// Pre-hold the key so TryLock always returns false.
130+
srv.held[AdvisoryKey("s")] = true
131+
l := &PGLocker{Subject: "s", Connect: srv.connector(), retry: time.Millisecond}
132+
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Millisecond)
133+
defer cancel()
134+
ran := false
135+
err := l.WithLock(ctx, func() error { ran = true; return nil })
136+
if err == nil || ran {
137+
t.Fatalf("a permanently-held key must time out before fn (err=%v ran=%v)", err, ran)
138+
}
139+
}
140+
141+
func TestPGLocker_TryLockErrorSurfaces(t *testing.T) {
142+
l := &PGLocker{Subject: "s", Connect: func(context.Context) (AdvisoryConn, error) {
143+
return &fakeSession{srv: newFakeAdvisory(), tryErr: errors.New("conn reset")}, nil
144+
}}
145+
if err := l.WithLock(context.Background(), func() error { return nil }); err == nil {
146+
t.Fatal("a TryLock error must surface")
147+
}
148+
}
149+
150+
// TestPGLocker_Serializes is the concurrency guarantee: N goroutines each take
151+
// the lock and do a deliberately non-atomic read-modify-write; mutual exclusion
152+
// must make the final count exact. Run under -race.
153+
func TestPGLocker_Serializes(t *testing.T) {
154+
srv := newFakeAdvisory()
155+
l := &PGLocker{Subject: "ledger", Connect: srv.connector(), retry: time.Millisecond}
156+
const n = 50
157+
counter := 0
158+
var wg sync.WaitGroup
159+
for range n {
160+
wg.Go(func() {
161+
_ = l.WithLock(context.Background(), func() error {
162+
c := counter
163+
time.Sleep(time.Microsecond) // widen the race window
164+
counter = c + 1
165+
return nil
166+
})
167+
})
168+
}
169+
wg.Wait()
170+
if counter != n {
171+
t.Fatalf("mutual exclusion failed: counter = %d, want %d", counter, n)
172+
}
173+
}

0 commit comments

Comments
 (0)