Commit 99c5044cba
Verified · cmc
Layout: unified · split
internal/packlimit/packlimit.go added +131
| @@ -0,0 +1,131 @@ | |||
| 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. | ||
| 8 | package packlimit | ||
| 9 | |||
| 10 | import ( | ||
| 11 | "errors" | ||
| 12 | "sync" | ||
| 13 | "time" | ||
| 14 | ) | ||
| 15 | |||
| 16 | var ( | ||
| 17 | ErrBusy = errors.New("the server is busy generating packs for other clients; try again in a minute") | ||
| 18 | ErrGone = errors.New("client went away while queued") | ||
| 19 | ) | ||
| 20 | |||
| 21 | type Limiter struct { | ||
| 22 | max, per, queue int | ||
| 23 | wait time.Duration | ||
| 24 | |||
| 25 | mu sync.Mutex | ||
| 26 | running int | ||
| 27 | queued int | ||
| 28 | held map[string]int // running, per principal | ||
| 29 | waiting map[string]int // queued, per principal | ||
| 30 | changed chan struct{} // closed and replaced on every release | ||
| 31 | } | ||
| 32 | |||
| 33 | // New returns a limiter, or nil — no limit — when max is not positive. | ||
| 34 | func New(max, per, queue int, wait time.Duration) *Limiter { | ||
| 35 | if max <= 0 { | ||
| 36 | return nil | ||
| 37 | } | ||
| 38 | return &Limiter{max: max, per: per, queue: queue, wait: wait, | ||
| 39 | held: map[string]int{}, waiting: map[string]int{}, changed: make(chan struct{})} | ||
| 40 | } | ||
| 41 | |||
| 42 | // Acquire takes a slot for principal, queueing when none is free. | ||
| 43 | // done, when it closes, ends the wait. Once Acquire returns a nil | ||
| 44 | // error, the caller holds the slot and must call release — once git | ||
| 45 | // has exited — regardless of what its own context has done since: | ||
| 46 | // done closing after that point does not release the slot on the | ||
| 47 | // caller's behalf. | ||
| 48 | func (l *Limiter) Acquire(done <-chan struct{}, principal string) (release func(), err error) { | ||
| 49 | if l == nil { | ||
| 50 | return func() {}, nil | ||
| 51 | } | ||
| 52 | l.mu.Lock() | ||
| 53 | if l.fits(principal) { | ||
| 54 | l.take(principal) | ||
| 55 | l.mu.Unlock() | ||
| 56 | return l.releaser(principal), nil | ||
| 57 | } | ||
| 58 | if l.queued >= l.queue || (l.per > 0 && l.waiting[principal] >= l.per) { | ||
| 59 | l.mu.Unlock() | ||
| 60 | return nil, ErrBusy | ||
| 61 | } | ||
| 62 | l.queued++ | ||
| 63 | l.waiting[principal]++ | ||
| 64 | l.mu.Unlock() | ||
| 65 | defer func() { | ||
| 66 | l.mu.Lock() | ||
| 67 | l.queued-- | ||
| 68 | if l.waiting[principal]--; l.waiting[principal] == 0 { | ||
| 69 | delete(l.waiting, principal) | ||
| 70 | } | ||
| 71 | l.mu.Unlock() | ||
| 72 | }() | ||
| 73 | |||
| 74 | timer := time.NewTimer(l.wait) | ||
| 75 | defer timer.Stop() | ||
| 76 | for { | ||
| 77 | l.mu.Lock() | ||
| 78 | // changed and done can both be ready at once — a slot can | ||
| 79 | // free up at the same moment the caller gives up. select | ||
| 80 | // among the wake sources would then pick between them at | ||
| 81 | // random, so re-check done first, under the lock, on every | ||
| 82 | // pass: this makes the limiter prefer ErrGone whenever both | ||
| 83 | // are ready, instead of leaving it to chance which one a | ||
| 84 | // given pass observes. | ||
| 85 | select { | ||
| 86 | case <-done: | ||
| 87 | l.mu.Unlock() | ||
| 88 | return nil, ErrGone | ||
| 89 | default: | ||
| 90 | } | ||
| 91 | if l.fits(principal) { | ||
| 92 | l.take(principal) | ||
| 93 | l.mu.Unlock() | ||
| 94 | return l.releaser(principal), nil | ||
| 95 | } | ||
| 96 | changed := l.changed | ||
| 97 | l.mu.Unlock() | ||
| 98 | select { | ||
| 99 | case <-changed: | ||
| 100 | case <-timer.C: | ||
| 101 | return nil, ErrBusy | ||
| 102 | case <-done: | ||
| 103 | return nil, ErrGone | ||
| 104 | } | ||
| 105 | } | ||
| 106 | } | ||
| 107 | |||
| 108 | func (l *Limiter) fits(principal string) bool { | ||
| 109 | return l.running < l.max && (l.per <= 0 || l.held[principal] < l.per) | ||
| 110 | } | ||
| 111 | |||
| 112 | func (l *Limiter) take(principal string) { | ||
| 113 | l.running++ | ||
| 114 | l.held[principal]++ | ||
| 115 | } | ||
| 116 | |||
| 117 | func (l *Limiter) releaser(principal string) func() { | ||
| 118 | var once sync.Once | ||
| 119 | return func() { | ||
| 120 | once.Do(func() { | ||
| 121 | l.mu.Lock() | ||
| 122 | defer l.mu.Unlock() | ||
| 123 | l.running-- | ||
| 124 | if l.held[principal]--; l.held[principal] == 0 { | ||
| 125 | delete(l.held, principal) | ||
| 126 | } | ||
| 127 | close(l.changed) | ||
| 128 | l.changed = make(chan struct{}) | ||
| 129 | }) | ||
| 130 | } | ||
| 131 | } | ||
internal/packlimit/packlimit_test.go added +209
| @@ -0,0 +1,209 @@ | |||
| 1 | package packlimit | ||
| 2 | |||
| 3 | import ( | ||
| 4 | "errors" | ||
| 5 | "testing" | ||
| 6 | "time" | ||
| 7 | ) | ||
| 8 | |||
| 9 | func TestGlobalCap(t *testing.T) { | ||
| 10 | l := New(2, 0, 0, time.Second) | ||
| 11 | r1, err1 := l.Acquire(nil, "a") | ||
| 12 | r2, err2 := l.Acquire(nil, "b") | ||
| 13 | if err1 != nil || err2 != nil { | ||
| 14 | t.Fatal(err1, err2) | ||
| 15 | } | ||
| 16 | if _, err := l.Acquire(nil, "c"); !errors.Is(err, ErrBusy) { | ||
| 17 | t.Fatalf("third with no queue: %v", err) | ||
| 18 | } | ||
| 19 | r1() | ||
| 20 | r1() // a second release is a no-op | ||
| 21 | r3, err := l.Acquire(nil, "c") | ||
| 22 | if err != nil { | ||
| 23 | t.Fatal(err) | ||
| 24 | } | ||
| 25 | if _, err := l.Acquire(nil, "d"); !errors.Is(err, ErrBusy) { | ||
| 26 | t.Fatalf("double release freed two slots: %v", err) | ||
| 27 | } | ||
| 28 | r2() | ||
| 29 | r3() | ||
| 30 | } | ||
| 31 | |||
| 32 | func TestPerPrincipalCap(t *testing.T) { | ||
| 33 | l := New(4, 1, 4, 50*time.Millisecond) | ||
| 34 | ra, err := l.Acquire(nil, "a") | ||
| 35 | if err != nil { | ||
| 36 | t.Fatal(err) | ||
| 37 | } | ||
| 38 | defer ra() | ||
| 39 | if _, err := l.Acquire(nil, "a"); !errors.Is(err, ErrBusy) { | ||
| 40 | t.Fatalf("second for a: %v", err) | ||
| 41 | } | ||
| 42 | rb, err := l.Acquire(nil, "b") | ||
| 43 | if err != nil { | ||
| 44 | t.Fatalf("b blocked by a: %v", err) | ||
| 45 | } | ||
| 46 | rb() | ||
| 47 | } | ||
| 48 | |||
| 49 | func TestWaiterGetsReleasedSlot(t *testing.T) { | ||
| 50 | l := New(1, 0, 1, 5*time.Second) | ||
| 51 | r1, _ := l.Acquire(nil, "a") | ||
| 52 | got := make(chan error, 1) | ||
| 53 | go func() { | ||
| 54 | r, err := l.Acquire(nil, "b") | ||
| 55 | if err == nil { | ||
| 56 | r() | ||
| 57 | } | ||
| 58 | got <- err | ||
| 59 | }() | ||
| 60 | waitQueued(t, l, 1) | ||
| 61 | r1() | ||
| 62 | select { | ||
| 63 | case err := <-got: | ||
| 64 | if err != nil { | ||
| 65 | t.Fatal(err) | ||
| 66 | } | ||
| 67 | case <-time.After(2 * time.Second): | ||
| 68 | t.Fatal("waiter never got the slot") | ||
| 69 | } | ||
| 70 | } | ||
| 71 | |||
| 72 | func TestQueueIsBounded(t *testing.T) { | ||
| 73 | l := New(1, 0, 1, 5*time.Second) | ||
| 74 | r1, _ := l.Acquire(nil, "a") | ||
| 75 | defer r1() | ||
| 76 | go l.Acquire(nil, "b") | ||
| 77 | waitQueued(t, l, 1) | ||
| 78 | if _, err := l.Acquire(nil, "c"); !errors.Is(err, ErrBusy) { | ||
| 79 | t.Fatalf("queue over its bound: %v", err) | ||
| 80 | } | ||
| 81 | } | ||
| 82 | |||
| 83 | // A principal cannot fill the queue on its own. | ||
| 84 | func TestPrincipalQueueIsBounded(t *testing.T) { | ||
| 85 | l := New(1, 1, 8, 5*time.Second) | ||
| 86 | r1, _ := l.Acquire(nil, "x") | ||
| 87 | defer r1() | ||
| 88 | go l.Acquire(nil, "a") | ||
| 89 | waitQueued(t, l, 1) | ||
| 90 | if _, err := l.Acquire(nil, "a"); !errors.Is(err, ErrBusy) { | ||
| 91 | t.Fatalf("second waiter for a: %v", err) | ||
| 92 | } | ||
| 93 | } | ||
| 94 | |||
| 95 | func TestClientGoneWhileQueued(t *testing.T) { | ||
| 96 | l := New(1, 0, 1, 5*time.Second) | ||
| 97 | r1, _ := l.Acquire(nil, "a") | ||
| 98 | done := make(chan struct{}) | ||
| 99 | close(done) | ||
| 100 | if _, err := l.Acquire(done, "b"); !errors.Is(err, ErrGone) { | ||
| 101 | t.Fatalf("got %v, want ErrGone", err) | ||
| 102 | } | ||
| 103 | assertQueueEmpty(t, l) | ||
| 104 | r1() | ||
| 105 | assertHeldEmpty(t, l) | ||
| 106 | } | ||
| 107 | |||
| 108 | func TestWaitRunsOut(t *testing.T) { | ||
| 109 | l := New(1, 0, 1, 20*time.Millisecond) | ||
| 110 | r1, _ := l.Acquire(nil, "a") | ||
| 111 | if _, err := l.Acquire(nil, "b"); !errors.Is(err, ErrBusy) { | ||
| 112 | t.Fatalf("got %v, want ErrBusy", err) | ||
| 113 | } | ||
| 114 | assertQueueEmpty(t, l) | ||
| 115 | r1() | ||
| 116 | assertHeldEmpty(t, l) | ||
| 117 | } | ||
| 118 | |||
| 119 | func assertQueueEmpty(t *testing.T, l *Limiter) { | ||
| 120 | t.Helper() | ||
| 121 | l.mu.Lock() | ||
| 122 | defer l.mu.Unlock() | ||
| 123 | if l.queued != 0 || len(l.waiting) != 0 { | ||
| 124 | t.Fatalf("queue not cleaned up: queued=%d waiting=%v", l.queued, l.waiting) | ||
| 125 | } | ||
| 126 | } | ||
| 127 | |||
| 128 | func assertHeldEmpty(t *testing.T, l *Limiter) { | ||
| 129 | t.Helper() | ||
| 130 | l.mu.Lock() | ||
| 131 | defer l.mu.Unlock() | ||
| 132 | if len(l.held) != 0 { | ||
| 133 | t.Fatalf("held not cleaned up: %v", l.held) | ||
| 134 | } | ||
| 135 | } | ||
| 136 | |||
| 137 | func TestNilLimiterNeverWaits(t *testing.T) { | ||
| 138 | var l *Limiter | ||
| 139 | if l = New(0, 1, 1, time.Second); l != nil { | ||
| 140 | t.Fatal("max 0 should mean no limit") | ||
| 141 | } | ||
| 142 | r, err := l.Acquire(nil, "a") | ||
| 143 | if err != nil { | ||
| 144 | t.Fatal(err) | ||
| 145 | } | ||
| 146 | r() | ||
| 147 | } | ||
| 148 | |||
| 149 | // A waiter whose done channel closes just as a slot frees up must not be | ||
| 150 | // granted the slot: it has to see ErrGone, and the slot must go to | ||
| 151 | // someone else instead of leaking to an abandoned caller. | ||
| 152 | func TestGivenUpWaiterNeverGetsSlot(t *testing.T) { | ||
| 153 | l := New(1, 0, 1, 10*time.Second) | ||
| 154 | _, _ = l.Acquire(nil, "a") // holds the only slot | ||
| 155 | |||
| 156 | done := make(chan struct{}) | ||
| 157 | got := make(chan error, 1) | ||
| 158 | go func() { | ||
| 159 | r, err := l.Acquire(done, "b") | ||
| 160 | if err == nil { | ||
| 161 | r() | ||
| 162 | } | ||
| 163 | got <- err | ||
| 164 | }() | ||
| 165 | waitQueued(t, l, 1) | ||
| 166 | |||
| 167 | // Close done and free a's slot in the same critical section, so | ||
| 168 | // changed and done become ready to b's select at the same instant | ||
| 169 | // — the exact race the done-check-before-fits ordering in Acquire | ||
| 170 | // has to win, whichever the select picks. | ||
| 171 | l.mu.Lock() | ||
| 172 | close(done) | ||
| 173 | l.running-- | ||
| 174 | delete(l.held, "a") | ||
| 175 | close(l.changed) | ||
| 176 | l.changed = make(chan struct{}) | ||
| 177 | l.mu.Unlock() | ||
| 178 | |||
| 179 | select { | ||
| 180 | case err := <-got: | ||
| 181 | if !errors.Is(err, ErrGone) { | ||
| 182 | t.Fatalf("got %v, want ErrGone", err) | ||
| 183 | } | ||
| 184 | case <-time.After(2 * time.Second): | ||
| 185 | t.Fatal("b never returned") | ||
| 186 | } | ||
| 187 | |||
| 188 | // The slot must still be free for someone else: b must not hold it. | ||
| 189 | r2, err := l.Acquire(nil, "c") | ||
| 190 | if err != nil { | ||
| 191 | t.Fatalf("slot leaked to the abandoned waiter: %v", err) | ||
| 192 | } | ||
| 193 | r2() | ||
| 194 | } | ||
| 195 | |||
| 196 | func waitQueued(t *testing.T, l *Limiter, n int) { | ||
| 197 | t.Helper() | ||
| 198 | deadline := time.Now().Add(2 * time.Second) | ||
| 199 | for time.Now().Before(deadline) { | ||
| 200 | l.mu.Lock() | ||
| 201 | q := l.queued | ||
| 202 | l.mu.Unlock() | ||
| 203 | if q == n { | ||
| 204 | return | ||
| 205 | } | ||
| 206 | time.Sleep(time.Millisecond) | ||
| 207 | } | ||
| 208 | t.Fatalf("queue never reached %d", n) | ||
| 209 | } | ||