Commit a41fc9d6a6
Verified · cmc ci/build: success ci/lint: success ci/test: success
Layout: unified · split
SETUP.md +5
| @@ -35,6 +35,11 @@ below apply. A file named with `-c` must exist. | |||
| 35 | recently used entries are dropped past this. | 35 | recently used entries are dropped past this. |
| 36 | * `ttl` — How long a response is reused, in the time units below. Default | 36 | * `ttl` — How long a response is reused, in the time units below. Default |
| 37 | `5i`. | 37 | `5i`. |
| 38 | * `stale` — How long past `ttl` a response is kept to be served when | ||
| 39 | DeviantArt fails or blocks the instance, so a short ban does not take | ||
| 40 | the daily deviations, popular searches and feeds down. Default `1h`; | ||
| 41 | `0i` keeps nothing past `ttl`. After a block the instance also leaves | ||
| 42 | DeviantArt alone for a minute instead of retrying every request. | ||
| 38 | * `rate-limit` — Per-client budget for page, feed and API requests, so one | 43 | * `rate-limit` — Per-client budget for page, feed and API requests, so one |
| 39 | crawler cannot spend the whole upstream budget. Media, avatars and static | 44 | crawler cannot spend the whole upstream budget. Media, avatars and static |
| 40 | files are not counted. Over budget answers 429 with `Retry-After`. | 45 | files are not counted. Over budget answers 429 with `Retry-After`. |
app/apicache.go +98 −28
| @@ -21,18 +21,31 @@ import ( | |||
| 21 | type apiCache struct { | 21 | type apiCache struct { |
| 22 | maxBytes int64 | 22 | maxBytes int64 |
| 23 | ttl time.Duration | 23 | ttl time.Duration |
| 24 | now func() time.Time | 24 | // stale is how long past ttl an entry is kept to be served when upstream |
| 25 | // is blocked or unreachable. Zero keeps nothing past ttl. | ||
| 26 | stale time.Duration | ||
| 27 | now func() time.Time | ||
| 25 | 28 | ||
| 26 | mu sync.Mutex | 29 | mu sync.Mutex |
| 27 | entries map[string]*cacheEntry | 30 | entries map[string]*cacheEntry |
| 28 | lru *list.List // front is most recently used | 31 | lru *list.List // front is most recently used |
| 29 | held int64 | 32 | held int64 |
| 30 | hits int64 | 33 | hits int64 |
| 31 | misses int64 | 34 | misses int64 |
| 35 | staleHits int64 | ||
| 36 | |||
| 37 | // blockedUntil is set when DeviantArt answers with a block. Until then a | ||
| 38 | // miss with nothing stale to serve gets blockResp back without an | ||
| 39 | // upstream call, so a banned instance stops hammering the WAF. | ||
| 40 | blockedUntil time.Time | ||
| 41 | blockResp *cacheEntry | ||
| 32 | 42 | ||
| 33 | flight singleflight.Group | 43 | flight singleflight.Group |
| 34 | } | 44 | } |
| 35 | 45 | ||
| 46 | // blockBackoff is how long upstream is left alone after a block response. | ||
| 47 | const blockBackoff = time.Minute | ||
| 48 | |||
| 36 | // cacheEntry is one buffered response. header is a clone of the upstream | 49 | // cacheEntry is one buffered response. header is a clone of the upstream |
| 37 | // header; body is the whole body, read once. | 50 | // header; body is the whole body, read once. |
| 38 | type cacheEntry struct { | 51 | type cacheEntry struct { |
| @@ -44,10 +57,11 @@ type cacheEntry struct { | |||
| 44 | elem *list.Element | 57 | elem *list.Element |
| 45 | } | 58 | } |
| 46 | 59 | ||
| 47 | func newAPICache(maxBytes int64, ttl time.Duration) *apiCache { | 60 | func newAPICache(maxBytes int64, ttl, stale time.Duration) *apiCache { |
| 48 | return &apiCache{ | 61 | return &apiCache{ |
| 49 | maxBytes: maxBytes, | 62 | maxBytes: maxBytes, |
| 50 | ttl: ttl, | 63 | ttl: ttl, |
| 64 | stale: stale, | ||
| 51 | now: time.Now, | 65 | now: time.Now, |
| 52 | entries: map[string]*cacheEntry{}, | 66 | entries: map[string]*cacheEntry{}, |
| 53 | lru: list.New(), | 67 | lru: list.New(), |
| @@ -87,16 +101,29 @@ type cachedTransport struct { | |||
| 87 | base http.RoundTripper | 101 | base http.RoundTripper |
| 88 | } | 102 | } |
| 89 | 103 | ||
| 90 | // RoundTrip serves a hit from memory. A miss is fetched once per key however | 104 | // RoundTrip serves a fresh hit from memory. A miss is fetched once per key |
| 91 | // many callers are waiting, buffered, stored if it is a 200, and handed to | 105 | // however many callers are waiting, buffered, stored if it is a 200, and |
| 92 | // every waiter as its own response. | 106 | // handed to every waiter as its own response. |
| 107 | // | ||
| 108 | // When upstream fails or answers with a block, a stale entry is served | ||
| 109 | // instead if one is still held, so a short ban does not take the popular | ||
| 110 | // pages down. A block also starts a backoff during which misses with nothing | ||
| 111 | // stale get the block response back without an upstream call. | ||
| 93 | func (t *cachedTransport) RoundTrip(req *http.Request) (*http.Response, error) { | 112 | func (t *cachedTransport) RoundTrip(req *http.Request) (*http.Response, error) { |
| 94 | if !cacheable(req) { | 113 | if !cacheable(req) { |
| 95 | return t.base.RoundTrip(req) | 114 | return t.base.RoundTrip(req) |
| 96 | } | 115 | } |
| 97 | key := cacheKey(req) | 116 | key := cacheKey(req) |
| 98 | if e := t.cache.get(key); e != nil { | 117 | old, fresh := t.cache.get(key) |
| 99 | return e.response(req), nil | 118 | if fresh { |
| 119 | return old.response(req), nil | ||
| 120 | } | ||
| 121 | if blocked := t.cache.blockedResponse(); blocked != nil { | ||
| 122 | if old != nil { | ||
| 123 | t.cache.countStale() | ||
| 124 | return old.response(req), nil | ||
| 125 | } | ||
| 126 | return blocked.response(req), nil | ||
| 100 | } | 127 | } |
| 101 | 128 | ||
| 102 | v, err, _ := t.cache.flight.Do(key, func() (any, error) { | 129 | v, err, _ := t.cache.flight.Do(key, func() (any, error) { |
| @@ -110,21 +137,57 @@ func (t *cachedTransport) RoundTrip(req *http.Request) (*http.Response, error) { | |||
| 110 | return nil, err | 137 | return nil, err |
| 111 | } | 138 | } |
| 112 | e := &cacheEntry{key: key, status: resp.StatusCode, header: resp.Header.Clone(), body: body} | 139 | e := &cacheEntry{key: key, status: resp.StatusCode, header: resp.Header.Clone(), body: body} |
| 113 | if e.status == http.StatusOK { | 140 | switch e.status { |
| 141 | case http.StatusOK: | ||
| 114 | t.cache.put(e) | 142 | t.cache.put(e) |
| 143 | case http.StatusForbidden, http.StatusTooManyRequests: | ||
| 144 | t.cache.block(e) | ||
| 115 | } | 145 | } |
| 116 | return e, nil | 146 | return e, nil |
| 117 | }) | 147 | }) |
| 118 | if err != nil { | 148 | if err != nil { |
| 149 | if old != nil { | ||
| 150 | t.cache.countStale() | ||
| 151 | return old.response(req), nil | ||
| 152 | } | ||
| 119 | return nil, err | 153 | return nil, err |
| 120 | } | 154 | } |
| 121 | e, ok := v.(*cacheEntry) | 155 | e, ok := v.(*cacheEntry) |
| 122 | if !ok { | 156 | if !ok { |
| 123 | return nil, io.ErrUnexpectedEOF | 157 | return nil, io.ErrUnexpectedEOF |
| 124 | } | 158 | } |
| 159 | if e.status != http.StatusOK && old != nil { | ||
| 160 | t.cache.countStale() | ||
| 161 | return old.response(req), nil | ||
| 162 | } | ||
| 125 | return e.response(req), nil | 163 | return e.response(req), nil |
| 126 | } | 164 | } |
| 127 | 165 | ||
| 166 | // block records a block response and starts the backoff. | ||
| 167 | func (c *apiCache) block(e *cacheEntry) { | ||
| 168 | c.mu.Lock() | ||
| 169 | defer c.mu.Unlock() | ||
| 170 | c.blockedUntil = c.now().Add(blockBackoff) | ||
| 171 | c.blockResp = e | ||
| 172 | } | ||
| 173 | |||
| 174 | // blockedResponse returns the last block response while the backoff runs, | ||
| 175 | // or nil once it is over. | ||
| 176 | func (c *apiCache) blockedResponse() *cacheEntry { | ||
| 177 | c.mu.Lock() | ||
| 178 | defer c.mu.Unlock() | ||
| 179 | if c.blockResp != nil && c.now().Before(c.blockedUntil) { | ||
| 180 | return c.blockResp | ||
| 181 | } | ||
| 182 | return nil | ||
| 183 | } | ||
| 184 | |||
| 185 | func (c *apiCache) countStale() { | ||
| 186 | c.mu.Lock() | ||
| 187 | c.staleHits++ | ||
| 188 | c.mu.Unlock() | ||
| 189 | } | ||
| 190 | |||
| 128 | // response builds a fresh http.Response over the buffered body, so each | 191 | // response builds a fresh http.Response over the buffered body, so each |
| 129 | // caller can read and close its own. | 192 | // caller can read and close its own. |
| 130 | func (e *cacheEntry) response(req *http.Request) *http.Response { | 193 | func (e *cacheEntry) response(req *http.Request) *http.Response { |
| @@ -141,25 +204,31 @@ func (e *cacheEntry) response(req *http.Request) *http.Response { | |||
| 141 | } | 204 | } |
| 142 | } | 205 | } |
| 143 | 206 | ||
| 144 | // get returns the live entry for key, marking it most recently used, or nil. | 207 | // get returns the entry for key and whether it is still fresh. A fresh hit is |
| 145 | // An expired entry is dropped on the way out. | 208 | // marked most recently used. An entry past ttl but within the stale window is |
| 146 | func (c *apiCache) get(key string) *cacheEntry { | 209 | // returned as not fresh, for the caller to fall back on; one past the stale |
| 210 | // window is dropped. | ||
| 211 | func (c *apiCache) get(key string) (*cacheEntry, bool) { | ||
| 147 | c.mu.Lock() | 212 | c.mu.Lock() |
| 148 | defer c.mu.Unlock() | 213 | defer c.mu.Unlock() |
| 149 | 214 | ||
| 150 | e := c.entries[key] | 215 | e := c.entries[key] |
| 151 | if e == nil { | 216 | if e == nil { |
| 152 | c.misses++ | 217 | c.misses++ |
| 153 | return nil | 218 | return nil, false |
| 154 | } | 219 | } |
| 155 | if !c.now().Before(e.expires) { | 220 | now := c.now() |
| 156 | c.remove(e) | 221 | if now.Before(e.expires) { |
| 157 | c.misses++ | 222 | c.lru.MoveToFront(e.elem) |
| 158 | return nil | 223 | c.hits++ |
| 224 | return e, true | ||
| 225 | } | ||
| 226 | c.misses++ | ||
| 227 | if now.Before(e.expires.Add(c.stale)) { | ||
| 228 | return e, false | ||
| 159 | } | 229 | } |
| 160 | c.lru.MoveToFront(e.elem) | 230 | c.remove(e) |
| 161 | c.hits++ | 231 | return nil, false |
| 162 | return e | ||
| 163 | } | 232 | } |
| 164 | 233 | ||
| 165 | // put stores e, evicting from the least recently used end until it fits. A | 234 | // put stores e, evicting from the least recently used end until it fits. A |
| @@ -196,9 +265,10 @@ func (c *apiCache) remove(e *cacheEntry) { | |||
| 196 | c.held -= int64(len(e.body)) | 265 | c.held -= int64(len(e.body)) |
| 197 | } | 266 | } |
| 198 | 267 | ||
| 199 | // stats reports the counters for the hourly log line. | 268 | // stats reports the counters for the hourly log line. stale counts the |
| 200 | func (c *apiCache) stats() (hits, misses int64, entries int, held int64) { | 269 | // misses that were answered from an expired entry because upstream failed. |
| 270 | func (c *apiCache) stats() (hits, misses, stale int64, entries int, held int64) { | ||
| 201 | c.mu.Lock() | 271 | c.mu.Lock() |
| 202 | defer c.mu.Unlock() | 272 | defer c.mu.Unlock() |
| 203 | return c.hits, c.misses, len(c.entries), c.held | 273 | return c.hits, c.misses, c.staleHits, len(c.entries), c.held |
| 204 | } | 274 | } |
app/apicache_stale_test.go added +132
| @@ -0,0 +1,132 @@ | |||
| 1 | package app | ||
| 2 | |||
| 3 | import ( | ||
| 4 | "errors" | ||
| 5 | "net/http" | ||
| 6 | "testing" | ||
| 7 | "time" | ||
| 8 | ) | ||
| 9 | |||
| 10 | // failingRT is an upstream that can be switched between a scripted status and | ||
| 11 | // a transport error mid-test. | ||
| 12 | type failingRT struct { | ||
| 13 | fakeRT | ||
| 14 | err error | ||
| 15 | } | ||
| 16 | |||
| 17 | func (f *failingRT) RoundTrip(r *http.Request) (*http.Response, error) { | ||
| 18 | f.mu.Lock() | ||
| 19 | err := f.err | ||
| 20 | f.mu.Unlock() | ||
| 21 | if err != nil { | ||
| 22 | f.mu.Lock() | ||
| 23 | f.calls++ | ||
| 24 | f.mu.Unlock() | ||
| 25 | return nil, err | ||
| 26 | } | ||
| 27 | return f.fakeRT.RoundTrip(r) | ||
| 28 | } | ||
| 29 | |||
| 30 | func (f *failingRT) set(status int, body string, err error) { | ||
| 31 | f.mu.Lock() | ||
| 32 | defer f.mu.Unlock() | ||
| 33 | f.status, f.body, f.err = status, body, err | ||
| 34 | } | ||
| 35 | |||
| 36 | // staleCache returns a cache with a one minute TTL and one hour stale window | ||
| 37 | // over a controllable clock, warmed with one good response for puppyURL. | ||
| 38 | func staleCache(t *testing.T) (*apiCache, *failingRT, http.RoundTripper, *time.Time) { | ||
| 39 | t.Helper() | ||
| 40 | up := &failingRT{} | ||
| 41 | up.set(200, `{"good":1}`, nil) | ||
| 42 | c := newAPICache(1<<20, time.Minute, time.Hour) | ||
| 43 | now := time.Now() | ||
| 44 | c.now = func() time.Time { return now } | ||
| 45 | rt := c.transport(up) | ||
| 46 | get(t, rt, puppyURL) | ||
| 47 | return c, up, rt, &now | ||
| 48 | } | ||
| 49 | |||
| 50 | func TestStaleEntryServedWhenUpstreamBlocks(t *testing.T) { | ||
| 51 | c, up, rt, now := staleCache(t) | ||
| 52 | *now = now.Add(2 * time.Minute) // past ttl, inside the stale window | ||
| 53 | up.set(403, "<html>blocked</html>", nil) | ||
| 54 | |||
| 55 | status, body := get(t, rt, puppyURL) | ||
| 56 | |||
| 57 | if status != 200 || body != `{"good":1}` { | ||
| 58 | t.Errorf("got %d %q, want the stale 200 body", status, body) | ||
| 59 | } | ||
| 60 | if up.count() != 2 { | ||
| 61 | t.Errorf("upstream called %d times, want 2: one warm-up and one attempt that hit the block", up.count()) | ||
| 62 | } | ||
| 63 | if _, _, stale, _, _ := c.stats(); stale != 1 { | ||
| 64 | t.Errorf("stale counter is %d, want 1", stale) | ||
| 65 | } | ||
| 66 | } | ||
| 67 | |||
| 68 | func TestStaleEntryServedOnTransportError(t *testing.T) { | ||
| 69 | _, up, rt, now := staleCache(t) | ||
| 70 | *now = now.Add(2 * time.Minute) | ||
| 71 | up.set(0, "", errors.New("dial tcp: connection refused")) | ||
| 72 | |||
| 73 | status, body := get(t, rt, puppyURL) | ||
| 74 | |||
| 75 | if status != 200 || body != `{"good":1}` { | ||
| 76 | t.Errorf("got %d %q, want the stale 200 body", status, body) | ||
| 77 | } | ||
| 78 | } | ||
| 79 | |||
| 80 | func TestBlockBackoffSkipsUpstream(t *testing.T) { | ||
| 81 | _, up, rt, now := staleCache(t) | ||
| 82 | *now = now.Add(2 * time.Minute) | ||
| 83 | up.set(403, "<html>blocked</html>", nil) | ||
| 84 | get(t, rt, puppyURL) // triggers the block | ||
| 85 | calls := up.count() | ||
| 86 | |||
| 87 | other := "https://www.deviantart.com/_puppy/dabrowse/search/all?q=other" | ||
| 88 | status, body := get(t, rt, other) | ||
| 89 | if status != 403 || body != "<html>blocked</html>" { | ||
| 90 | t.Errorf("uncached key during backoff got %d %q, want the block response", status, body) | ||
| 91 | } | ||
| 92 | if up.count() != calls { | ||
| 93 | t.Errorf("upstream called during backoff (%d -> %d), want none", calls, up.count()) | ||
| 94 | } | ||
| 95 | |||
| 96 | *now = now.Add(blockBackoff + time.Second) | ||
| 97 | up.set(200, `{"back":1}`, nil) | ||
| 98 | if status, _ := get(t, rt, other); status != 200 || up.count() != calls+1 { | ||
| 99 | t.Errorf("after backoff: status %d, calls %d, want 200 and one more upstream call", status, up.count()) | ||
| 100 | } | ||
| 101 | } | ||
| 102 | |||
| 103 | func TestStaleEntryDroppedAfterTheWindow(t *testing.T) { | ||
| 104 | c, up, rt, now := staleCache(t) | ||
| 105 | *now = now.Add(time.Minute + time.Hour + time.Second) // past ttl and stale | ||
| 106 | up.set(403, "<html>blocked</html>", nil) | ||
| 107 | |||
| 108 | status, _ := get(t, rt, puppyURL) | ||
| 109 | |||
| 110 | if status != 403 { | ||
| 111 | t.Errorf("got %d, want the 403 passed through once nothing stale is held", status) | ||
| 112 | } | ||
| 113 | if _, _, _, entries, _ := c.stats(); entries != 0 { | ||
| 114 | t.Errorf("%d entries held, want 0", entries) | ||
| 115 | } | ||
| 116 | } | ||
| 117 | |||
| 118 | func TestZeroStaleKeepsNothingPastTTL(t *testing.T) { | ||
| 119 | up := &failingRT{} | ||
| 120 | up.set(200, "x", nil) | ||
| 121 | c := newAPICache(1<<20, time.Minute, 0) | ||
| 122 | now := time.Now() | ||
| 123 | c.now = func() time.Time { return now } | ||
| 124 | rt := c.transport(up) | ||
| 125 | get(t, rt, puppyURL) | ||
| 126 | now = now.Add(2 * time.Minute) | ||
| 127 | up.set(403, "blocked", nil) | ||
| 128 | |||
| 129 | if status, _ := get(t, rt, puppyURL); status != 403 { | ||
| 130 | t.Errorf("got %d with stale 0, want 403: nothing may be served past ttl", status) | ||
| 131 | } | ||
| 132 | } | ||
app/apicache_test.go +18 −16
| @@ -56,7 +56,7 @@ func get(t *testing.T, rt http.RoundTripper, url string) (int, string) { | |||
| 56 | 56 | ||
| 57 | func TestSecondRequestIsServedFromCache(t *testing.T) { | 57 | func TestSecondRequestIsServedFromCache(t *testing.T) { |
| 58 | up := &fakeRT{status: 200, body: `{"a":1}`} | 58 | up := &fakeRT{status: 200, body: `{"a":1}`} |
| 59 | rt := newAPICache(1<<20, time.Minute).transport(up) | 59 | rt := newAPICache(1<<20, time.Minute, 0).transport(up) |
| 60 | 60 | ||
| 61 | get(t, rt, puppyURL) | 61 | get(t, rt, puppyURL) |
| 62 | status, body := get(t, rt, puppyURL) | 62 | status, body := get(t, rt, puppyURL) |
| @@ -71,7 +71,7 @@ func TestSecondRequestIsServedFromCache(t *testing.T) { | |||
| 71 | 71 | ||
| 72 | func TestExpiredEntryIsRefetched(t *testing.T) { | 72 | func TestExpiredEntryIsRefetched(t *testing.T) { |
| 73 | up := &fakeRT{status: 200, body: `{}`} | 73 | up := &fakeRT{status: 200, body: `{}`} |
| 74 | c := newAPICache(1<<20, time.Minute) | 74 | c := newAPICache(1<<20, time.Minute, 0) |
| 75 | now := time.Now() | 75 | now := time.Now() |
| 76 | c.now = func() time.Time { return now } | 76 | c.now = func() time.Time { return now } |
| 77 | rt := c.transport(up) | 77 | rt := c.transport(up) |
| @@ -85,24 +85,26 @@ func TestExpiredEntryIsRefetched(t *testing.T) { | |||
| 85 | } | 85 | } |
| 86 | } | 86 | } |
| 87 | 87 | ||
| 88 | // A 500 here rather than a 403: a 403 is a block and starts the backoff, | ||
| 89 | // which is covered in apicache_stale_test.go. | ||
| 88 | func TestNon200IsNotStored(t *testing.T) { | 90 | func TestNon200IsNotStored(t *testing.T) { |
| 89 | up := &fakeRT{status: 403, body: "blocked"} | 91 | up := &fakeRT{status: 500, body: "upstream broke"} |
| 90 | rt := newAPICache(1<<20, time.Minute).transport(up) | 92 | rt := newAPICache(1<<20, time.Minute, 0).transport(up) |
| 91 | 93 | ||
| 92 | status, body := get(t, rt, puppyURL) | 94 | status, body := get(t, rt, puppyURL) |
| 93 | get(t, rt, puppyURL) | 95 | get(t, rt, puppyURL) |
| 94 | 96 | ||
| 95 | if status != 403 || body != "blocked" { | 97 | if status != 500 || body != "upstream broke" { |
| 96 | t.Errorf("first response is %d %q, want the upstream 403 passed through", status, body) | 98 | t.Errorf("first response is %d %q, want the upstream 500 passed through", status, body) |
| 97 | } | 99 | } |
| 98 | if up.count() != 2 { | 100 | if up.count() != 2 { |
| 99 | t.Errorf("upstream called %d times, want 2: a 403 must not be cached", up.count()) | 101 | t.Errorf("upstream called %d times, want 2: a 500 must not be cached", up.count()) |
| 100 | } | 102 | } |
| 101 | } | 103 | } |
| 102 | 104 | ||
| 103 | func TestBypassesSessionAndOtherHosts(t *testing.T) { | 105 | func TestBypassesSessionAndOtherHosts(t *testing.T) { |
| 104 | up := &fakeRT{status: 200, body: "x"} | 106 | up := &fakeRT{status: 200, body: "x"} |
| 105 | rt := newAPICache(1<<20, time.Minute).transport(up) | 107 | rt := newAPICache(1<<20, time.Minute, 0).transport(up) |
| 106 | 108 | ||
| 107 | for _, url := range []string{ | 109 | for _, url := range []string{ |
| 108 | "https://www.deviantart.com/_puppy", | 110 | "https://www.deviantart.com/_puppy", |
| @@ -119,7 +121,7 @@ func TestBypassesSessionAndOtherHosts(t *testing.T) { | |||
| 119 | 121 | ||
| 120 | func TestKeyIgnoresCSRFToken(t *testing.T) { | 122 | func TestKeyIgnoresCSRFToken(t *testing.T) { |
| 121 | up := &fakeRT{status: 200, body: "x"} | 123 | up := &fakeRT{status: 200, body: "x"} |
| 122 | rt := newAPICache(1<<20, time.Minute).transport(up) | 124 | rt := newAPICache(1<<20, time.Minute, 0).transport(up) |
| 123 | 125 | ||
| 124 | get(t, rt, puppyURL) | 126 | get(t, rt, puppyURL) |
| 125 | get(t, rt, strings.Replace(puppyURL, "csrf_token=abc", "csrf_token=def", 1)) | 127 | get(t, rt, strings.Replace(puppyURL, "csrf_token=abc", "csrf_token=def", 1)) |
| @@ -131,7 +133,7 @@ func TestKeyIgnoresCSRFToken(t *testing.T) { | |||
| 131 | 133 | ||
| 132 | func TestByteBoundEvictsLeastRecentlyUsed(t *testing.T) { | 134 | func TestByteBoundEvictsLeastRecentlyUsed(t *testing.T) { |
| 133 | up := &fakeRT{status: 200, body: strings.Repeat("x", 100)} | 135 | up := &fakeRT{status: 200, body: strings.Repeat("x", 100)} |
| 134 | rt := newAPICache(250, time.Minute).transport(up) | 136 | rt := newAPICache(250, time.Minute, 0).transport(up) |
| 135 | a := "https://www.deviantart.com/_puppy/a?p=1" | 137 | a := "https://www.deviantart.com/_puppy/a?p=1" |
| 136 | b := "https://www.deviantart.com/_puppy/b?p=1" | 138 | b := "https://www.deviantart.com/_puppy/b?p=1" |
| 137 | c := "https://www.deviantart.com/_puppy/c?p=1" | 139 | c := "https://www.deviantart.com/_puppy/c?p=1" |
| @@ -151,7 +153,7 @@ func TestByteBoundEvictsLeastRecentlyUsed(t *testing.T) { | |||
| 151 | 153 | ||
| 152 | func TestConcurrentMissesMakeOneUpstreamCall(t *testing.T) { | 154 | func TestConcurrentMissesMakeOneUpstreamCall(t *testing.T) { |
| 153 | up := &fakeRT{status: 200, body: "x", delay: 50 * time.Millisecond} | 155 | up := &fakeRT{status: 200, body: "x", delay: 50 * time.Millisecond} |
| 154 | rt := newAPICache(1<<20, time.Minute).transport(up) | 156 | rt := newAPICache(1<<20, time.Minute, 0).transport(up) |
| 155 | 157 | ||
| 156 | var wg sync.WaitGroup | 158 | var wg sync.WaitGroup |
| 157 | for range 20 { | 159 | for range 20 { |
| @@ -166,16 +168,16 @@ func TestConcurrentMissesMakeOneUpstreamCall(t *testing.T) { | |||
| 166 | 168 | ||
| 167 | func TestStatsCountHitsAndMisses(t *testing.T) { | 169 | func TestStatsCountHitsAndMisses(t *testing.T) { |
| 168 | up := &fakeRT{status: 200, body: "abc"} | 170 | up := &fakeRT{status: 200, body: "abc"} |
| 169 | c := newAPICache(1<<20, time.Minute) | 171 | c := newAPICache(1<<20, time.Minute, 0) |
| 170 | rt := c.transport(up) | 172 | rt := c.transport(up) |
| 171 | 173 | ||
| 172 | get(t, rt, puppyURL) | 174 | get(t, rt, puppyURL) |
| 173 | get(t, rt, puppyURL) | 175 | get(t, rt, puppyURL) |
| 174 | get(t, rt, puppyURL) | 176 | get(t, rt, puppyURL) |
| 175 | 177 | ||
| 176 | hits, misses, entries, held := c.stats() | 178 | hits, misses, stale, entries, held := c.stats() |
| 177 | if hits != 2 || misses != 1 || entries != 1 || held != 3 { | 179 | if hits != 2 || misses != 1 || stale != 0 || entries != 1 || held != 3 { |
| 178 | t.Errorf("stats = %d hits, %d misses, %d entries, %d bytes; want 2, 1, 1, 3", hits, misses, entries, held) | 180 | t.Errorf("stats = %d hits, %d misses, %d stale, %d entries, %d bytes; want 2, 1, 0, 1, 3", hits, misses, stale, entries, held) |
| 179 | } | 181 | } |
| 180 | } | 182 | } |
| 181 | 183 | ||
| @@ -185,7 +187,7 @@ func TestStatsCountHitsAndMisses(t *testing.T) { | |||
| 185 | func TestHitDoesNotConsumeAThrottleSlot(t *testing.T) { | 187 | func TestHitDoesNotConsumeAThrottleSlot(t *testing.T) { |
| 186 | up := &fakeRT{status: 200, body: "x"} | 188 | up := &fakeRT{status: 200, body: "x"} |
| 187 | th := &daThrottle{base: up, sem: make(chan struct{}, 1)} | 189 | th := &daThrottle{base: up, sem: make(chan struct{}, 1)} |
| 188 | rt := newAPICache(1<<20, time.Minute).transport(th) | 190 | rt := newAPICache(1<<20, time.Minute, 0).transport(th) |
| 189 | 191 | ||
| 190 | get(t, rt, puppyURL) // populate through the throttle | 192 | get(t, rt, puppyURL) // populate through the throttle |
| 191 | 193 | ||
app/config.go +12 −2
| @@ -32,6 +32,7 @@ type apiCacheConfig struct { | |||
| 32 | Enabled bool `json:"enabled"` | 32 | Enabled bool `json:"enabled"` |
| 33 | MaxSize int64 `json:"max-size"` | 33 | MaxSize int64 `json:"max-size"` |
| 34 | TTL string `json:"ttl"` | 34 | TTL string `json:"ttl"` |
| 35 | Stale string `json:"stale"` | ||
| 35 | } | 36 | } |
| 36 | 37 | ||
| 37 | type rateLimitConfig struct { | 38 | type rateLimitConfig struct { |
| @@ -75,6 +76,7 @@ var CFG = config{ | |||
| 75 | Enabled: true, | 76 | Enabled: true, |
| 76 | MaxSize: 64, | 77 | MaxSize: 64, |
| 77 | TTL: "5i", | 78 | TTL: "5i", |
| 79 | Stale: "1h", | ||
| 78 | }, | 80 | }, |
| 79 | RateLimit: rateLimitConfig{ | 81 | RateLimit: rateLimitConfig{ |
| 80 | PerMinute: 60, | 82 | PerMinute: 60, |
| @@ -88,8 +90,9 @@ var CFG = config{ | |||
| 88 | 90 | ||
| 89 | var lifetimeParsed int64 | 91 | var lifetimeParsed int64 |
| 90 | 92 | ||
| 91 | // apiCacheTTL is api-cache.ttl parsed, set by ExecuteConfig. | 93 | // apiCacheTTL and apiCacheStale are api-cache.ttl and api-cache.stale parsed, |
| 92 | var apiCacheTTL time.Duration | 94 | // set by ExecuteConfig. |
| 95 | var apiCacheTTL, apiCacheStale time.Duration | ||
| 93 | 96 | ||
| 94 | // parseLifetime reads a duration in the config's unit syntax: a number | 97 | // parseLifetime reads a duration in the config's unit syntax: a number |
| 95 | // followed by i (minutes), h (hours), d (days), w (weeks), m (30-day | 98 | // followed by i (minutes), h (hours), d (days), w (weeks), m (30-day |
| @@ -216,6 +219,13 @@ func ExecuteConfig() { | |||
| 216 | exit("config: api-cache.ttl: "+err.Error(), 1) | 219 | exit("config: api-cache.ttl: "+err.Error(), 1) |
| 217 | } | 220 | } |
| 218 | apiCacheTTL = d | 221 | apiCacheTTL = d |
| 222 | if CFG.APICache.Stale != "" { | ||
| 223 | d, err := parseLifetime(CFG.APICache.Stale) | ||
| 224 | if err != nil { | ||
| 225 | exit("config: api-cache.stale: "+err.Error(), 1) | ||
| 226 | } | ||
| 227 | apiCacheStale = d | ||
| 228 | } | ||
| 219 | } | 229 | } |
| 220 | 230 | ||
| 221 | // per-minute 0 turns the limit off; a burst below one token would | 231 | // per-minute 0 turns the limit off; a burst below one token would |
app/config_test.go +2 −2
| @@ -32,8 +32,8 @@ func TestParseLifetimeRejectsBadInput(t *testing.T) { | |||
| 32 | } | 32 | } |
| 33 | 33 | ||
| 34 | func TestAPICacheDefaults(t *testing.T) { | 34 | func TestAPICacheDefaults(t *testing.T) { |
| 35 | if !CFG.APICache.Enabled || CFG.APICache.MaxSize != 64 || CFG.APICache.TTL != "5i" { | 35 | if !CFG.APICache.Enabled || CFG.APICache.MaxSize != 64 || CFG.APICache.TTL != "5i" || CFG.APICache.Stale != "1h" { |
| 36 | t.Errorf("defaults are %+v, want enabled, 64 MB, 5i", CFG.APICache) | 36 | t.Errorf("defaults are %+v, want enabled, 64 MB, 5i, stale 1h", CFG.APICache) |
| 37 | } | 37 | } |
| 38 | } | 38 | } |
| 39 | 39 | ||
app/httpclient.go +3 −3
| @@ -100,8 +100,8 @@ func chain(base http.RoundTripper) http.RoundTripper { | |||
| 100 | func logCacheStatsForever(c *apiCache) { | 100 | func logCacheStatsForever(c *apiCache) { |
| 101 | for { | 101 | for { |
| 102 | time.Sleep(time.Hour) | 102 | time.Sleep(time.Hour) |
| 103 | hits, misses, entries, held := c.stats() | 103 | hits, misses, stale, entries, held := c.stats() |
| 104 | println("api cache:", hits, "hits,", misses, "misses,", entries, "entries,", held>>20, "MB held") | 104 | println("api cache:", hits, "hits,", misses, "misses,", stale, "served stale,", entries, "entries,", held>>20, "MB held") |
| 105 | } | 105 | } |
| 106 | } | 106 | } |
| 107 | 107 | ||
| @@ -111,7 +111,7 @@ func logCacheStatsForever(c *apiCache) { | |||
| 111 | func InstallDAThrottle() { | 111 | func InstallDAThrottle() { |
| 112 | baseTransport = tunedTransport() | 112 | baseTransport = tunedTransport() |
| 113 | if CFG.APICache.Enabled { | 113 | if CFG.APICache.Enabled { |
| 114 | daCache = newAPICache(CFG.APICache.MaxSize<<20, apiCacheTTL) | 114 | daCache = newAPICache(CFG.APICache.MaxSize<<20, apiCacheTTL, apiCacheStale) |
| 115 | go logCacheStatsForever(daCache) | 115 | go logCacheStatsForever(daCache) |
| 116 | } | 116 | } |
| 117 | http.DefaultTransport = chain(baseTransport) | 117 | http.DefaultTransport = chain(baseTransport) |