Commit 32f6c3dc37

32f6c3dc376001abeb8305c5abd90dfa5756989d

parent: ff8e09b0c8

Verified · cmc ci/build: success ci/lint: success ci/test: success

cmc <hello@cleberg.net> · 2026-09-12 01:24 UTC

Drop abandoned requests from the upstream queue, and shed a deep one

A request whose context ends while it waits for a slot or for the
interval now returns at once without taking a turn; before, every
request a client had already given up on still burned an interval,
which under a crawl turned the queue into a wall of timeouts. A
request that would queue longer than a client waits (20s against
devianter's 30s) is refused immediately with errUpstreamBusy, which
the error page maps to a 503 with Retry-After.

Layout: unified · split

app/httpclient.go +46 −4
@@ -1,10 +1,12 @@
11package app
22
33import (
4 "errors"
45 "net/http"
56 "net/url"
67 "strings"
78 "sync"
9 "sync/atomic"
810 "time"
911)
1012
@@ -35,24 +37,64 @@ type daThrottle struct {
3537 sem chan struct{}
3638 mu sync.Mutex
3739 last time.Time
40
41 // waiting counts requests queued for a slot, for load shedding.
42 waiting atomic.Int64
3843}
3944
45// maxQueueWait is the longest a request may expect to queue before it is shed
46// with errUpstreamBusy. devianter gives up after 30 seconds, so a request
47// that would wait longer than this would only time out anyway, while holding
48// a place in the queue that a live request could have used.
49const maxQueueWait = 20 * time.Second
50
51// errUpstreamBusy is returned without an upstream call when the queue is
52// already deeper than a client will wait for. Error maps it to a 503.
53var errUpstreamBusy = errors.New("upstream queue full")
54
4055// RoundTrip applies the rate and concurrency limits to DeviantArt requests and
4156// passes everything else straight through to the base transport.
57//
58// A request whose context ends while it waits is dropped without taking a
59// turn: under a crawl, most queued requests have already been abandoned by
60// their client, and letting each one still burn an interval slot is what
61// turned the queue into a wall of timeouts.
4262func (t *daThrottle) RoundTrip(req *http.Request) (*http.Response, error) {
4363 // Only throttle DeviantArt's WAF-protected API host; let everything else fly.
4464 if !strings.Contains(req.URL.Hostname(), "deviantart.com") {
4565 return t.base.RoundTrip(req)
4666 }
67 ctx := req.Context()
4768
48 // Concurrency cap: block until a slot frees up (backpressure under floods).
49 t.sem <- struct{}{}
69 // Shed before queueing when the queue already implies a wait no client
70 // will sit through.
71 if queued := t.waiting.Load(); time.Duration(queued)*daMinInterval > maxQueueWait {
72 return nil, errUpstreamBusy
73 }
74 t.waiting.Add(1)
75 defer t.waiting.Add(-1)
76
77 // Concurrency cap: wait for a slot, or give up with the caller.
78 select {
79 case t.sem <- struct{}{}:
80 case <-ctx.Done():
81 return nil, ctx.Err()
82 }
5083 defer func() { <-t.sem }()
5184
52 // Rate cap: enforce a minimum interval between request starts.
85 // Rate cap: enforce a minimum interval between request starts. A request
86 // cancelled while waiting leaves last untouched, so the interval it did
87 // not use goes to the next request.
5388 t.mu.Lock()
5489 if wait := daMinInterval - time.Since(t.last); wait > 0 {
55 time.Sleep(wait)
90 timer := time.NewTimer(wait)
91 select {
92 case <-timer.C:
93 case <-ctx.Done():
94 timer.Stop()
95 t.mu.Unlock()
96 return nil, ctx.Err()
97 }
5698 }
5799 t.last = time.Now()
58100 t.mu.Unlock()
app/httpclient_test.go +103
@@ -1,11 +1,15 @@
11package app
22
33import (
4 "context"
5 "errors"
46 "net/http"
57 "net/http/httptest"
68 "sync"
79 "testing"
810 "time"
11
12 "github.com/krazywarez/devianter"
913)
1014
1115// stubTransport records how many requests reached it and returns an empty 200.
@@ -143,3 +147,102 @@ func TestInstallDAThrottlePreservesProxy(t *testing.T) {
143147 t.Error("base transport lost its Proxy func: HTTPS_PROXY / VPN egress would break")
144148 }
145149}
150
151// countingRT counts base round trips for the throttle tests.
152type countingRT struct {
153 mu sync.Mutex
154 calls int
155}
156
157func (c *countingRT) RoundTrip(r *http.Request) (*http.Response, error) {
158 c.mu.Lock()
159 c.calls++
160 c.mu.Unlock()
161 return &http.Response{StatusCode: 200, Body: http.NoBody, Request: r}, nil
162}
163
164func daRequest(ctx context.Context) *http.Request {
165 req, _ := http.NewRequestWithContext(ctx, http.MethodGet, "https://www.deviantart.com/_puppy/x", nil)
166 return req
167}
168
169// TestThrottleShedsADeepQueue pins load shedding: with more requests queued
170// than the client timeout can absorb, a new one is refused at once with
171// errUpstreamBusy and never reaches upstream.
172func TestThrottleShedsADeepQueue(t *testing.T) {
173 interval := daMinInterval
174 daMinInterval = time.Second
175 defer func() { daMinInterval = interval }()
176 base := &countingRT{}
177 th := &daThrottle{base: base, sem: make(chan struct{}, 1)}
178 th.waiting.Store(int64(maxQueueWait/time.Second) + 1)
179
180 start := time.Now()
181 _, err := th.RoundTrip(daRequest(context.Background()))
182
183 if !errors.Is(err, errUpstreamBusy) {
184 t.Fatalf("err = %v, want errUpstreamBusy", err)
185 }
186 if time.Since(start) > 100*time.Millisecond || base.calls != 0 {
187 t.Errorf("shed request took %v and made %d upstream calls, want immediate and none", time.Since(start), base.calls)
188 }
189}
190
191// TestThrottleDropsACancelledRequestWaitingForASlot pins that a request whose
192// client has gone does not sit in the queue: with the only slot held, a
193// cancelled context returns at once.
194func TestThrottleDropsACancelledRequestWaitingForASlot(t *testing.T) {
195 base := &countingRT{}
196 th := &daThrottle{base: base, sem: make(chan struct{}, 1)}
197 th.sem <- struct{}{} // hold the only slot
198 ctx, cancel := context.WithCancel(context.Background())
199 cancel()
200
201 start := time.Now()
202 _, err := th.RoundTrip(daRequest(ctx))
203
204 if !errors.Is(err, context.Canceled) {
205 t.Fatalf("err = %v, want context.Canceled", err)
206 }
207 if time.Since(start) > 100*time.Millisecond || base.calls != 0 {
208 t.Errorf("took %v with %d upstream calls, want immediate and none", time.Since(start), base.calls)
209 }
210}
211
212// TestThrottleCancelledDuringIntervalKeepsTheSlot pins that a request
213// cancelled while waiting out the interval does not consume it: the next
214// live request starts as soon as the original interval allows.
215func TestThrottleCancelledDuringIntervalKeepsTheSlot(t *testing.T) {
216 interval := daMinInterval
217 daMinInterval = 300 * time.Millisecond
218 defer func() { daMinInterval = interval }()
219 base := &countingRT{}
220 th := &daThrottle{base: base, sem: make(chan struct{}, 1)}
221
222 if _, err := th.RoundTrip(daRequest(context.Background())); err != nil {
223 t.Fatal(err)
224 }
225 first := th.last
226
227 ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
228 defer cancel()
229 if _, err := th.RoundTrip(daRequest(ctx)); !errors.Is(err, context.DeadlineExceeded) {
230 t.Fatalf("err = %v, want context.DeadlineExceeded", err)
231 }
232 if th.last != first {
233 t.Error("a cancelled request advanced the interval clock")
234 }
235 if base.calls != 1 {
236 t.Errorf("%d upstream calls, want 1: the cancelled request must not go upstream", base.calls)
237 }
238}
239
240// TestErrorPageMapsAShedRequestTo503 pins the user-facing side: a shed
241// request is a 503 with Retry-After, not a DeviantArt error.
242func TestErrorPageMapsAShedRequestTo503(t *testing.T) {
243 rec := httptest.NewRecorder()
244 skunkyart{Writer: rec, Host: "http://localhost"}.Error(devianter.Error{Error: "devianter: Get ...: " + errUpstreamBusy.Error()})
245 if rec.Code != 503 || rec.Header().Get("Retry-After") != "5" {
246 t.Errorf("status %d Retry-After %q, want 503 and 5", rec.Code, rec.Header().Get("Retry-After"))
247 }
248}
app/util.go +8
@@ -226,6 +226,14 @@ func URLBuilder(host string, strs ...string) string {
226226// neither readable nor safe to echo.
227227func (s skunkyart) Error(dAerr devianter.Error) {
228228 s.Writer.Header().Del("Cache-Control")
229
230 // A shed request is not an upstream failure: the instance is busy and the
231 // client should come back shortly rather than treat the page as broken.
232 if strings.Contains(dAerr.Error, errUpstreamBusy.Error()) {
233 s.Writer.Header().Set("Retry-After", "5")
234 s.ReturnHTTPError(http.StatusServiceUnavailable)
235 return
236 }
229237 s.Writer.WriteHeader(502)
230238
231239 reason, _, _ := strings.Cut(dAerr.Error, "\n")