internal/packlimit/packlimit.go
205 lines · 5658 bytes
1// Package packlimit bounds concurrent git pack generation. upload-pack
2// and upload-archive over SSH, smart HTTP and git:// draw on one
3// budget: a global cap, a cap per principal (an account, or a client
4// address on the anonymous transports), and a bounded queue whose
5// waiters give up after a fixed wait or when the client goes away.
6// Waiters are not served in order; a new arrival can take a freed slot
7// ahead of them, and the wait bounds how long any one of them waits.
8package packlimit
9
10import (
11 "errors"
12 "log/slog"
13 "net/netip"
14 "strings"
15 "sync"
16 "time"
17)
18
19var (
20 ErrBusy = errors.New("the server is busy generating packs for other clients; try again in a minute")
21 ErrGone = errors.New("client went away while queued")
22)
23
24type Limiter struct {
25 max, per, queue int
26 wait time.Duration
27
28 // Principals starting with class may hold at most classCap slots
29 // between them; classCap 0 is no class cap.
30 class string
31 classCap int
32
33 mu sync.Mutex
34 running int
35 classHeld int
36 queued int
37 held map[string]int // running, per principal
38 waiting map[string]int // queued, per principal
39 changed chan struct{} // closed and replaced on every release
40 warned map[string]time.Time // last refusal logged, per transport
41}
42
43// New returns a limiter, or nil — no limit — when max is not positive.
44func New(max, per, queue int, wait time.Duration) *Limiter {
45 if max <= 0 {
46 return nil
47 }
48 return &Limiter{max: max, per: per, queue: queue, wait: wait,
49 held: map[string]int{}, waiting: map[string]int{}, changed: make(chan struct{}),
50 warned: map[string]time.Time{}}
51}
52
53// Refused logs that a request on transport was turned away with err, at
54// most once a minute per transport. It names the principal's class
55// (user or ip), never the principal: an address is personal data.
56func (l *Limiter) Refused(transport, principal string, err error) {
57 if l == nil {
58 return
59 }
60 now := time.Now()
61 l.mu.Lock()
62 last, seen := l.warned[transport]
63 if seen && now.Sub(last) < time.Minute {
64 l.mu.Unlock()
65 return
66 }
67 l.warned[transport] = now
68 l.mu.Unlock()
69 class, _, _ := strings.Cut(principal, ":")
70 reason := "busy"
71 if errors.Is(err, ErrGone) {
72 reason = "gone"
73 }
74 slog.Warn("pack limit: request turned away (logged at most once a minute per transport)",
75 "transport", transport, "class", class, "reason", reason)
76}
77
78// CapClass caps the slots that principals starting with prefix may hold
79// between them. Call it before the limiter is in use.
80func (l *Limiter) CapClass(prefix string, n int) {
81 if l == nil {
82 return
83 }
84 l.class, l.classCap = prefix, n
85}
86
87// AddrPrincipal is the principal for an unauthenticated client at addr:
88// an IPv4 address as is, an IPv6 address by its /64, since one host
89// commonly holds a whole /64. An address that does not parse is used
90// as given.
91func AddrPrincipal(addr string) string {
92 a, err := netip.ParseAddr(addr)
93 if err != nil {
94 return "ip:" + addr
95 }
96 a = a.WithZone("").Unmap()
97 if a.Is4() {
98 return "ip:" + a.String()
99 }
100 return "ip:" + netip.PrefixFrom(a, 64).Masked().String()
101}
102
103// Acquire takes a slot for principal, queueing when none is free.
104// done, when it closes, ends the wait. Once Acquire returns a nil
105// error, the caller holds the slot and must call release — once git
106// has exited — regardless of what its own context has done since:
107// done closing after that point does not release the slot on the
108// caller's behalf.
109func (l *Limiter) Acquire(done <-chan struct{}, principal string) (release func(), err error) {
110 if l == nil {
111 return func() {}, nil
112 }
113 l.mu.Lock()
114 if l.fits(principal) {
115 l.take(principal)
116 l.mu.Unlock()
117 return l.releaser(principal), nil
118 }
119 if l.queued >= l.queue || (l.per > 0 && l.waiting[principal] >= l.per) {
120 l.mu.Unlock()
121 return nil, ErrBusy
122 }
123 l.queued++
124 l.waiting[principal]++
125 l.mu.Unlock()
126 defer func() {
127 l.mu.Lock()
128 l.queued--
129 if l.waiting[principal]--; l.waiting[principal] == 0 {
130 delete(l.waiting, principal)
131 }
132 l.mu.Unlock()
133 }()
134
135 timer := time.NewTimer(l.wait)
136 defer timer.Stop()
137 for {
138 l.mu.Lock()
139 // changed and done can both be ready at once — a slot can
140 // free up at the same moment the caller gives up. select
141 // among the wake sources would then pick between them at
142 // random, so re-check done first, under the lock, on every
143 // pass: this makes the limiter prefer ErrGone whenever both
144 // are ready, instead of leaving it to chance which one a
145 // given pass observes.
146 select {
147 case <-done:
148 l.mu.Unlock()
149 return nil, ErrGone
150 default:
151 }
152 if l.fits(principal) {
153 l.take(principal)
154 l.mu.Unlock()
155 return l.releaser(principal), nil
156 }
157 changed := l.changed
158 l.mu.Unlock()
159 select {
160 case <-changed:
161 case <-timer.C:
162 return nil, ErrBusy
163 case <-done:
164 return nil, ErrGone
165 }
166 }
167}
168
169func (l *Limiter) fits(principal string) bool {
170 if l.inClass(principal) && l.classHeld >= l.classCap {
171 return false
172 }
173 return l.running < l.max && (l.per <= 0 || l.held[principal] < l.per)
174}
175
176func (l *Limiter) inClass(principal string) bool {
177 return l.classCap > 0 && strings.HasPrefix(principal, l.class)
178}
179
180func (l *Limiter) take(principal string) {
181 l.running++
182 l.held[principal]++
183 if l.inClass(principal) {
184 l.classHeld++
185 }
186}
187
188func (l *Limiter) releaser(principal string) func() {
189 var once sync.Once
190 return func() {
191 once.Do(func() {
192 l.mu.Lock()
193 defer l.mu.Unlock()
194 l.running--
195 if l.inClass(principal) {
196 l.classHeld--
197 }
198 if l.held[principal]--; l.held[principal] == 0 {
199 delete(l.held, principal)
200 }
201 close(l.changed)
202 l.changed = make(chan struct{})
203 })
204 }
205}