app/httpclient.go

v1.5.5
skunky-art/app/httpclient.go history · blame · raw

181 lines · 6509 bytes

  1package app
  2
  3import (
  4	"errors"
  5	"net/http"
  6	"net/url"
  7	"strings"
  8	"sync"
  9	"sync/atomic"
 10	"time"
 11)
 12
 13// DeviantArt fronts its API with AWS CloudFront + WAF, which bans egress IPs that
 14// hit it too hard. Under a bot flood, unbounded concurrent handlers each fetch
 15// ~150-200 KB of DA JSON, which both hammers that IP (risking a ban) and can OOM
 16// the process. devianter makes its requests with a bare &http.Client{}, so they go
 17// through http.DefaultTransport — we wrap it here to bound the rate and concurrency
 18// of calls to deviantart.com and to add timeouts. Requests to other hosts (e.g.
 19// wixmp image CDN) are passed straight through, so media stays fast.
 20//
 21// http.ProxyFromEnvironment is preserved, so HTTPS_PROXY (VPN egress) still applies.
 22
 23// Tunables, set from the upstream config block by ExecuteConfig; these are
 24// the defaults for a config that omits it. Slower is gentler on the egress
 25// address, which DeviantArt bans when it asks too often.
 26var (
 27	daMinInterval   = 400 * time.Millisecond // minimum gap between DA request starts
 28	daMaxConcurrent = 2                      // max simultaneous in-flight DA requests
 29)
 30
 31// downloadTimeout bounds a single outbound fetch end to end, so that a stalled
 32// CDN connection cannot pin a request handler open indefinitely.
 33const downloadTimeout = 60 * time.Second
 34
 35type daThrottle struct {
 36	base http.RoundTripper
 37	sem  chan struct{}
 38	mu   sync.Mutex
 39	last time.Time
 40
 41	// waiting counts requests queued for a slot, for load shedding.
 42	waiting atomic.Int64
 43}
 44
 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
 55// RoundTrip applies the rate and concurrency limits to DeviantArt requests and
 56// 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.
 62func (t *daThrottle) RoundTrip(req *http.Request) (*http.Response, error) {
 63	// Only throttle DeviantArt's WAF-protected API host; let everything else fly.
 64	if !strings.Contains(req.URL.Hostname(), "deviantart.com") {
 65		return t.base.RoundTrip(req)
 66	}
 67	ctx := req.Context()
 68
 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	}
 83	defer func() { <-t.sem }()
 84
 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.
 88	t.mu.Lock()
 89	if wait := daMinInterval - time.Since(t.last); wait > 0 {
 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		}
 98	}
 99	t.last = time.Now()
100	t.mu.Unlock()
101
102	return t.base.RoundTrip(req)
103}
104
105// baseTransport is the tuned transport installed by InstallDAThrottle, kept so
106// that per-client transports (see ProxiedTransport) inherit the same timeouts
107// instead of silently bypassing them.
108var baseTransport *http.Transport
109
110// tunedTransport clones the current default transport, preserving its Proxy
111// (ProxyFromEnvironment) and connection-pool defaults, and tightens timeouts to
112// bound hung connections.
113func tunedTransport() *http.Transport {
114	base, ok := http.DefaultTransport.(*http.Transport)
115	if !ok {
116		// Already wrapped, or a non-standard transport is installed. Start from a
117		// fresh one rather than panicking on a type assertion.
118		base = &http.Transport{Proxy: http.ProxyFromEnvironment}
119	}
120
121	t := base.Clone()
122	t.TLSHandshakeTimeout = 10 * time.Second
123	t.ResponseHeaderTimeout = 20 * time.Second
124	t.ExpectContinueTimeout = 2 * time.Second
125	return t
126}
127
128// daCache is the API response cache shared by every transport, or nil when
129// api-cache.enabled is false.
130var daCache *apiCache
131
132// chain wraps base with the throttle and, when enabled, the cache in front
133// of it, so a hit never spends a throttle slot.
134func chain(base http.RoundTripper) http.RoundTripper {
135	rt := throttled(base)
136	if daCache != nil {
137		return daCache.transport(rt)
138	}
139	return rt
140}
141
142// logCacheStatsForever prints one line an hour so an operator can see the
143// cache working without an endpoint. Run it in its own goroutine.
144func logCacheStatsForever(c *apiCache) {
145	for {
146		time.Sleep(time.Hour)
147		hits, misses, stale, entries, held := c.stats()
148		println("api cache:", hits, "hits,", misses, "misses,", stale, "served stale,", entries, "entries,", held>>20, "MB held")
149	}
150}
151
152// InstallDAThrottle wraps http.DefaultTransport with the rate/concurrency limits
153// and timeouts above, and with the API response cache when it is enabled. Call
154// once at startup, after ExecuteConfig and before any DeviantArt request.
155func InstallDAThrottle() {
156	baseTransport = tunedTransport()
157	if CFG.APICache.Enabled {
158		daCache = newAPICache(CFG.APICache.MaxSize<<20, apiCacheTTL, apiCacheStale)
159		go logCacheStatsForever(daCache)
160	}
161	http.DefaultTransport = chain(baseTransport)
162}
163
164// throttled wraps base with the DeviantArt rate and concurrency limits.
165func throttled(base http.RoundTripper) http.RoundTripper {
166	return &daThrottle{base: base, sem: make(chan struct{}, daMaxConcurrent)}
167}
168
169// ProxiedTransport returns a throttled transport routing through proxy. Downloads
170// configured with download-proxy go through here so they keep the timeouts and
171// limits that InstallDAThrottle installs on the default transport.
172func ProxiedTransport(proxy *url.URL) http.RoundTripper {
173	var base *http.Transport
174	if baseTransport != nil {
175		base = baseTransport.Clone()
176	} else {
177		base = tunedTransport()
178	}
179	base.Proxy = http.ProxyURL(proxy)
180	return chain(base)
181}