Skip to content

Commit 4f7b8f9

Browse files
gm-e2be2b-bot[bot]
authored andcommitted
fix(db): back off serializable transaction retries
GitOrigin-RevId: 76591f66b30fb229ca506daacc33a820cadd82b1
1 parent fd9602b commit 4f7b8f9

2 files changed

Lines changed: 67 additions & 4 deletions

File tree

‎packages/db/pkg/pool/session.go‎

Lines changed: 23 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import (
44
"context"
55
"errors"
66
"fmt"
7+
"math/rand/v2"
78
"time"
89

910
"github.com/jackc/pgerrcode"
@@ -13,9 +14,11 @@ import (
1314
)
1415

1516
const (
16-
sessionReleaseTimeout = 5 * time.Second
17-
serializableAttempts = 10
18-
serializableRetryStep = 5 * time.Millisecond
17+
sessionReleaseTimeout = 5 * time.Second
18+
serializableAttempts = 10
19+
serializableRetryBase = 10 * time.Millisecond
20+
serializableRetryMax = 500 * time.Millisecond
21+
serializableRetryTimeout = 5 * time.Second
1922
)
2023

2124
var (
@@ -92,6 +95,7 @@ func (l *AdvisoryLock) InSerializableTx(ctx context.Context, fn func(context.Con
9295
// InSerializableTxReturn1 runs fn in a SERIALIZABLE transaction on the session
9396
// that holds the lock. Serialization failures and deadlocks replay the whole
9497
// callback; only the value produced by the attempt that commits is returned.
98+
// Attempts and backoff share a five-second budget, bounded by the caller's deadline.
9599
func InSerializableTxReturn1[T any](
96100
ctx context.Context,
97101
lock *AdvisoryLock,
@@ -102,8 +106,15 @@ func InSerializableTxReturn1[T any](
102106
return zero, errLockReleased
103107
}
104108

109+
ctx, cancel := context.WithTimeout(ctx, serializableRetryTimeout)
110+
defer cancel()
111+
105112
var conflict error
106113
for attempt := 1; attempt <= serializableAttempts; attempt++ {
114+
if err := ctx.Err(); err != nil {
115+
return zero, errors.Join(err, conflict)
116+
}
117+
107118
value, err := runInTxReturn1(ctx, lock.conn, fn)
108119
switch {
109120
case err == nil:
@@ -165,7 +176,7 @@ func isSerializationConflict(err error) bool {
165176
}
166177

167178
func waitToReplay(ctx context.Context, attempt int) error {
168-
timer := time.NewTimer(serializableRetryStep * time.Duration(attempt))
179+
timer := time.NewTimer(serializableRetryDelay(attempt))
169180
defer timer.Stop()
170181

171182
select {
@@ -176,6 +187,14 @@ func waitToReplay(ctx context.Context, attempt int) error {
176187
}
177188
}
178189

190+
func serializableRetryDelay(attempt int) time.Duration {
191+
// Equal jitter keeps retries apart without allowing an immediate replay.
192+
// attempt is bounded by serializableAttempts, so the shift cannot overflow.
193+
window := min(2*serializableRetryBase<<(attempt-1), serializableRetryMax)
194+
195+
return window/2 + rand.N(window/2)
196+
}
197+
179198
// Release unlocks the session and returns its connection to the pool. It is
180199
// idempotent. If the unlock result is unknown, Release destroys the session so
181200
// a connection carrying an unknown lock cannot re-enter the pool.

‎packages/db/pkg/pool/session_integration_test.go‎

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,36 @@ import (
1818

1919
const testPostgresImage = "postgres:18-alpine"
2020

21+
func TestSerializableRetryDelay(t *testing.T) {
22+
t.Parallel()
23+
24+
for _, tc := range []struct {
25+
attempt int
26+
minimum time.Duration
27+
maximum time.Duration
28+
}{
29+
{1, 10 * time.Millisecond, 20 * time.Millisecond},
30+
{2, 20 * time.Millisecond, 40 * time.Millisecond},
31+
{3, 40 * time.Millisecond, 80 * time.Millisecond},
32+
{5, 160 * time.Millisecond, 320 * time.Millisecond},
33+
{6, 250 * time.Millisecond, 500 * time.Millisecond},
34+
{9, 250 * time.Millisecond, 500 * time.Millisecond},
35+
} {
36+
t.Run(fmt.Sprintf("attempt-%d", tc.attempt), func(t *testing.T) {
37+
t.Parallel()
38+
39+
delays := make(map[time.Duration]struct{})
40+
for range 32 {
41+
delay := serializableRetryDelay(tc.attempt)
42+
require.GreaterOrEqual(t, delay, tc.minimum)
43+
require.Less(t, delay, tc.maximum)
44+
delays[delay] = struct{}{}
45+
}
46+
assert.Greater(t, len(delays), 1, "retries must not use a deterministic delay")
47+
})
48+
}
49+
}
50+
2151
func TestConnectChecksAndConfiguresPool(t *testing.T) {
2252
t.Parallel()
2353

@@ -96,6 +126,9 @@ func TestAdvisoryLockRetriesSerializableTransactionOnItsSession(t *testing.T) {
96126
attempts := 0
97127
result, err := InSerializableTxReturn1(t.Context(), lock, func(ctx context.Context, tx pgx.Tx) (transactionResult, error) {
98128
attempts++
129+
deadline, ok := ctx.Deadline()
130+
require.True(t, ok, "transaction attempts must have a deadline")
131+
assert.LessOrEqual(t, time.Until(deadline), 5*time.Second)
99132

100133
var value int
101134
var backend int
@@ -194,6 +227,17 @@ func TestSerializableTransactionCleansUpAndBoundsReplays(t *testing.T) {
194227
).Scan(&value))
195228
assert.Equal(t, 3, value)
196229

230+
deadline := time.Now().Add(time.Second)
231+
shorter, cancelShorter := context.WithDeadline(t.Context(), deadline)
232+
defer cancelShorter()
233+
require.NoError(t, lock.InSerializableTx(shorter, func(ctx context.Context, _ pgx.Tx) error {
234+
actual, ok := ctx.Deadline()
235+
require.True(t, ok)
236+
assert.Equal(t, deadline, actual, "transaction budget must not extend the caller's deadline")
237+
238+
return nil
239+
}))
240+
197241
attempts := 0
198242
err = lock.InSerializableTx(t.Context(), func(context.Context, pgx.Tx) error {
199243
attempts++

0 commit comments

Comments
 (0)