-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathbudget.go
More file actions
275 lines (255 loc) · 6.32 KB
/
Copy pathbudget.go
File metadata and controls
275 lines (255 loc) · 6.32 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
package tpool
import (
"container/list"
"context"
"errors"
"sync"
"sync/atomic"
"time"
)
// waiter is a blocked Acquire waiting for tokens.
type waiter struct {
n int64
ready chan struct{}
ctx context.Context
}
const (
// waiterFlag is the high bit of state: set while waiters are queued.
waiterFlag int64 = -1 << 63
// usedMask isolates the token-count bits of state.
usedMask int64 = 1<<63 - 1
)
// Budget is a resizable semaphore. Acquire blocks (in FIFO order) while used
// tokens exceed the limit; Release and SetLimit wake waiters.
//
// Locking model: the uncontended path is lock-free. used and a "waiters
// present" flag share one atomic word (state), so a tenant that never hits
// the ceiling acquires and releases a token without touching mu. mu is only
// acquired to enqueue/wake waiters or hand tokens to them (notify), and by
// the reconciler/SetLimit.
//
// FIFO: the fast path yields whenever the waiters flag is set, and the flag
// is set in the same atomic op that proves exhaustion -- so a token freed by
// a concurrent fast Release is never lost to a queued waiter: any Release
// whose Add lands after that proof reads the flag set and notifies under mu.
type Budget struct {
mu sync.Mutex
limit atomic.Int64
state atomic.Int64 // waiterFlag bit | used token count
waiters list.List
}
var ErrBudgetExhausted = errors.New("tpool: budget limit exceeded")
// NewBudget returns a Budget with the given token limit.
func NewBudget(n int64) *Budget {
b := &Budget{}
b.limit.Store(n)
return b
}
// Acquire reserves n tokens, blocking until the limit allows it or ctx ends.
func (b *Budget) Acquire(ctx context.Context, n int64) error {
if n > b.limit.Load() {
return ErrBudgetExhausted
}
// Lock-free fast path: no waiters, no contention on mu.
if b.reserve(n) {
return nil
}
b.mu.Lock()
for {
s := b.state.Load()
used := s & usedMask
if s >= 0 && b.limit.Load()-used >= n {
// No waiters queued and tokens are free: grant directly.
if b.state.CompareAndSwap(s, s+n) {
b.mu.Unlock()
return nil
}
continue
}
if ctx.Err() != nil {
b.mu.Unlock()
return ctx.Err()
}
if s >= 0 {
// Exhausted and flag not yet set: set it in the same op that
// proved exhaustion, so concurrent fast Releases can't miss
// us (see FIFO comment on Budget).
if !b.state.CompareAndSwap(s, s|waiterFlag) {
continue
}
}
break
}
w := &waiter{n: n, ready: make(chan struct{}, 1), ctx: ctx}
e := b.waiters.PushBack(w)
b.mu.Unlock()
for {
select {
case <-ctx.Done():
b.mu.Lock()
if w.n > 0 {
b.state.Add(-w.n)
w.n = 0
b.notify()
}
b.waiters.Remove(e)
b.clearFlagIfEmpty()
b.mu.Unlock()
return ctx.Err()
case <-w.ready:
if w.n > 0 {
return nil
}
}
}
}
// reserve is the lock-free fast path: it claims n tokens only when no waiter
// is queued, preserving FIFO fairness without taking mu.
func (b *Budget) reserve(n int64) bool {
for {
s := b.state.Load()
if s < 0 {
return false
}
if (s&usedMask)+n > b.limit.Load() {
return false
}
if b.state.CompareAndSwap(s, s+n) {
return true
}
}
}
// Release returns n tokens.
func (b *Budget) Release(n int64) {
for {
s := b.state.Load()
if s < 0 {
// Waiters queued: free the token and hand off under mu so the
// oldest waiter gets it before any newcomer.
b.mu.Lock()
b.state.Add(-n)
b.notify()
b.mu.Unlock()
return
}
if b.state.CompareAndSwap(s, s-n) {
return
}
}
}
// TryAcquire reserves n tokens if available without blocking. Honors the
// waiters flag so it never jumps the queue.
func (b *Budget) TryAcquire(n int64) bool {
for {
s := b.state.Load()
if s < 0 || (s&usedMask)+n > b.limit.Load() {
return false
}
if b.state.CompareAndSwap(s, s+n) {
return true
}
}
}
// SetLimit changes the token limit, waking cashable waiters.
func (b *Budget) SetLimit(n int64) {
b.mu.Lock()
b.limit.Store(n)
b.notify()
b.mu.Unlock()
}
// notify hands available tokens to queued waiters in FIFO order. Must hold mu.
func (b *Budget) notify() {
for e := b.waiters.Front(); e != nil; {
w := e.Value.(*waiter)
s := b.state.Load()
used := s & usedMask
limit := b.limit.Load()
if used > limit || limit-used < w.n {
break
}
if !b.state.CompareAndSwap(s, s+w.n) {
continue
}
next := e.Next()
b.waiters.Remove(e)
if w.ctx.Err() != nil {
b.state.Add(-w.n)
w.n = 0
}
w.ready <- struct{}{}
e = next
}
b.clearFlagIfEmpty()
}
// clearFlagIfEmpty drops the waiters flag once no waiters remain. Must hold mu.
func (b *Budget) clearFlagIfEmpty() {
if b.waiters.Len() != 0 {
return
}
for {
s := b.state.Load()
if s >= 0 {
return
}
if b.state.CompareAndSwap(s, s&usedMask) {
return
}
}
}
func (b *Budget) snapshot() (limit, used int64) {
return b.limit.Load(), b.state.Load() & usedMask
}
// waitersLen returns how many Acquire calls are blocked on the budget queue.
func (b *Budget) waitersLen() int {
b.mu.Lock()
defer b.mu.Unlock()
return b.waiters.Len()
}
// reconcile repairs leaked tokens: a socket that died without releasing its
// token (a failed dial runs BeforeConnect but no BeforeClose). used is exact
// bookkeeping otherwise, so reconcile only tightens it — never loosens.
// Loosening would re-book tokens freed by async destroys (puddle still counts
// the dying socket) and overcommit the ceiling.
func (b *Budget) reconcile(n int64) {
b.mu.Lock()
for {
s := b.state.Load()
used := s & usedMask
if n >= used {
break
}
if b.state.CompareAndSwap(s, s-(used-n)) {
b.notify()
break
}
}
b.mu.Unlock()
}
// reconciler periodically re-anchors budget usage to physical pool stats.
func (p *Pool) reconciler(ctx context.Context) {
t := time.NewTicker(p.cfg.ReconcileInterval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
p.budget.reconcile(p.liveConns())
}
}
}
// liveConns counts connections that hold a budget token: acquired, idle and
// constructing. TotalConns already includes constructing ones, so adding
// ConstructingConns again would double-count zombies blocked on acquireToken.
func (p *Pool) liveConns() int64 {
p.mu.RLock()
defer p.mu.RUnlock()
var total int64
for _, tp := range p.tenants {
total += int64(tp.base.Stat().TotalConns())
if tp.burst != nil {
total += int64(tp.burst.Stat().TotalConns())
}
}
return total
}