internal/packlimit/packlimit.go

f8b976a97290a20d552056a999511f5d27d8e8ec
gitbay/internal/packlimit/packlimit.go history · blame · raw

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}