Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
4009490242 | ||
|
|
bcd5363d99 |
@@ -107,6 +107,13 @@ or proxy cache keeps the image after pixa would refuse the URL. A URL with no
|
|||||||
expiry gets one year. `immutable` only stops a client revalidating while its
|
expiry gets one year. `immutable` only stops a client revalidating while its
|
||||||
copy is fresh.
|
copy is fresh.
|
||||||
|
|
||||||
|
When several requests for the same image, size, format, quality and fit miss
|
||||||
|
the cache at once, they share one upstream fetch (or one read of the cached
|
||||||
|
source) and one transcode: the first request does the work, and the others wait
|
||||||
|
for its image or its error, holding no upstream connection or processing slot
|
||||||
|
of their own. A waiting request stops waiting when its own client goes away;
|
||||||
|
the work goes on for the others.
|
||||||
|
|
||||||
The login form (`POST /`) is limited to 5 attempts per minute per client
|
The login form (`POST /`) is limited to 5 attempts per minute per client
|
||||||
address, counting an IPv6 client by its /64; an attempt over the limit is
|
address, counting an IPv6 client by its /64; an attempt over the limit is
|
||||||
refused with 429 and a `Retry-After` header. Behind a reverse proxy the client
|
refused with 429 and a `Retry-After` header. Behind a reverse proxy the client
|
||||||
@@ -238,7 +245,7 @@ variables set by the file's `env:` section are checked the same way.
|
|||||||
| `PIXA_UPSTREAM_FETCH_TIMEOUT` | `upstream_fetch_timeout` | Time allowed for one fetch from an upstream host; default `30s` |
|
| `PIXA_UPSTREAM_FETCH_TIMEOUT` | `upstream_fetch_timeout` | Time allowed for one fetch from an upstream host; default `30s` |
|
||||||
| `PIXA_UPSTREAM_MAX_RESPONSE_SIZE` | `upstream_max_response_size` | Largest upstream response accepted, in bytes; default 50 MiB |
|
| `PIXA_UPSTREAM_MAX_RESPONSE_SIZE` | `upstream_max_response_size` | Largest upstream response accepted, in bytes; default 50 MiB |
|
||||||
| `PIXA_DOWNSTREAM_TIMEOUT` | `downstream_timeout` | Time allowed for answering one client request; default `60s` |
|
| `PIXA_DOWNSTREAM_TIMEOUT` | `downstream_timeout` | Time allowed for answering one client request; default `60s` |
|
||||||
| `PIXA_ACCESS_CONTROL_ALLOW_ORIGIN` | `access_control_allow_origin` | CORS origin allowed to read image responses: `*` or one origin; default `*` |
|
| `PIXA_ACCESS_CONTROL_ALLOW_ORIGIN` | `access_control_allow_origin` | CORS origin allowed to read responses: `*` or one origin; default `*` |
|
||||||
| `PIXA_METRICS_USERNAME` | `metrics.username` | Username for `/metrics`, which is served only when both are set |
|
| `PIXA_METRICS_USERNAME` | `metrics.username` | Username for `/metrics`, which is served only when both are set |
|
||||||
| `PIXA_METRICS_PASSWORD` | `metrics.password` | Password for `/metrics`; set together with the username |
|
| `PIXA_METRICS_PASSWORD` | `metrics.password` | Password for `/metrics`; set together with the username |
|
||||||
| `PIXA_SENTRY_DSN` | `sentry_dsn` | Sentry DSN for error reporting; empty disables it |
|
| `PIXA_SENTRY_DSN` | `sentry_dsn` | Sentry DSN for error reporting; empty disables it |
|
||||||
@@ -247,9 +254,8 @@ variables set by the file's `env:` section are checked the same way.
|
|||||||
|
|
||||||
Key settings in more detail:
|
Key settings in more detail:
|
||||||
|
|
||||||
- `access_control_allow_origin` — the origin a browser lets read the responses
|
- `access_control_allow_origin` — the origin a browser lets read pixa's
|
||||||
of the image routes, `/v1/image/` and `/v1/e/`, sent as the CORS
|
responses, sent as the CORS `Access-Control-Allow-Origin` header: `*`, the
|
||||||
`Access-Control-Allow-Origin` header; no other route sends it. `*`, the
|
|
||||||
default, is any site; otherwise one `http` or `https` origin such as
|
default, is any site; otherwise one `http` or `https` origin such as
|
||||||
`https://example.com`, whose host is a lowercase host name (letters,
|
`https://example.com`, whose host is a lowercase host name (letters,
|
||||||
digits, hyphens and dots, with a letter in its last part) or an IP address
|
digits, hyphens and dots, with a letter in its last part) or an IP address
|
||||||
|
|||||||
@@ -29,12 +29,17 @@ P2: security: referer blacklist
|
|||||||
|
|
||||||
# Completed Steps
|
# Completed Steps
|
||||||
|
|
||||||
- 2026-09-29 only the image routes send CORS headers (closes #98): the CORS
|
- 2026-09-29 share concurrent misses (closes #65): requests that miss the same
|
||||||
middleware, with the `access_control_allow_origin` origin, moved from the
|
variant at once (the same cache key, so quality and fit included) share one
|
||||||
router root onto a `/v1` subrouter holding `/v1/image/` and `/v1/e/`, where it
|
upstream fetch or cached source read and one transcode through
|
||||||
still answers a preflight `OPTIONS` request; the login and URL generator
|
`golang.org/x/sync/singleflight`; the first request's processing runs with a
|
||||||
pages, `/metrics` and the other routes send no `Access-Control-Allow-Origin`;
|
context that does not end with its own, and the others wait for its image or
|
||||||
documented in `README.md` and `config.example.yml`.
|
error holding no upstream connection or processing slot, and stop waiting
|
||||||
|
when their own context ends; the request doing the processing waits for it
|
||||||
|
even then, as before; a request whose context has already ended starts
|
||||||
|
nothing; each request counts one miss, and the processing counts its fetch
|
||||||
|
and transcode once; a panic while processing becomes an error for every
|
||||||
|
waiting request instead of stopping pixad; documented in `README.md`.
|
||||||
- 2026-09-29 the container makes `/var/lib/pixa` usable by itself (closes
|
- 2026-09-29 the container makes `/var/lib/pixa` usable by itself (closes
|
||||||
#159): `deploy/docker-entrypoint.sh` creates the directory if it is missing,
|
#159): `deploy/docker-entrypoint.sh` creates the directory if it is missing,
|
||||||
gives the directory and everything in it to `pixad` when the directory or one
|
gives the directory and everything in it to `pixad` when the directory or one
|
||||||
|
|||||||
+2
-3
@@ -103,9 +103,8 @@ upstream_max_response_size: 52428800
|
|||||||
# longer than upstream_fetch_timeout plus 20 seconds.
|
# longer than upstream_fetch_timeout plus 20 seconds.
|
||||||
downstream_timeout: 60s
|
downstream_timeout: 60s
|
||||||
|
|
||||||
# The origin a browser lets read the responses of the image routes,
|
# The origin a browser lets read pixa's responses, sent as the CORS
|
||||||
# /v1/image/ and /v1/e/, sent as the CORS Access-Control-Allow-Origin
|
# Access-Control-Allow-Origin header: "*" (the default) is any site;
|
||||||
# header; no other route sends it. "*" (the default) is any site;
|
|
||||||
# otherwise one http or https origin such as https://example.com, whose
|
# otherwise one http or https origin such as https://example.com, whose
|
||||||
# host is a lowercase host name (letters, digits, hyphens and dots, with a
|
# host is a lowercase host name (letters, digits, hyphens and dots, with a
|
||||||
# letter in its last part) or an IP address (IPv6 in brackets, in its
|
# letter in its last part) or an IP address (IPv6 in brackets, in its
|
||||||
|
|||||||
@@ -20,6 +20,7 @@ require (
|
|||||||
github.com/spf13/cobra v1.10.2
|
github.com/spf13/cobra v1.10.2
|
||||||
go.uber.org/fx v1.24.0
|
go.uber.org/fx v1.24.0
|
||||||
golang.org/x/crypto v0.41.0
|
golang.org/x/crypto v0.41.0
|
||||||
|
golang.org/x/sync v0.19.0
|
||||||
modernc.org/sqlite v1.42.2
|
modernc.org/sqlite v1.42.2
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -134,7 +135,6 @@ require (
|
|||||||
golang.org/x/image v0.34.0 // indirect
|
golang.org/x/image v0.34.0 // indirect
|
||||||
golang.org/x/net v0.43.0 // indirect
|
golang.org/x/net v0.43.0 // indirect
|
||||||
golang.org/x/oauth2 v0.30.0 // indirect
|
golang.org/x/oauth2 v0.30.0 // indirect
|
||||||
golang.org/x/sync v0.19.0 // indirect
|
|
||||||
golang.org/x/sys v0.36.0 // indirect
|
golang.org/x/sys v0.36.0 // indirect
|
||||||
golang.org/x/term v0.34.0 // indirect
|
golang.org/x/term v0.34.0 // indirect
|
||||||
golang.org/x/text v0.32.0 // indirect
|
golang.org/x/text v0.32.0 // indirect
|
||||||
|
|||||||
@@ -471,7 +471,8 @@ func (c *Cache) Stats(ctx context.Context) (*CacheStats, error) {
|
|||||||
return &stats, nil
|
return &stats, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// IncrementStats increments cache statistics.
|
// IncrementStats counts a cache hit or miss, and an upstream fetch that read
|
||||||
|
// fetchBytes bytes, as IncrementUpstreamFetch does.
|
||||||
func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64) {
|
func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64) {
|
||||||
var err error
|
var err error
|
||||||
|
|
||||||
@@ -495,8 +496,17 @@ func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64)
|
|||||||
c.log.Warn("failed to count cache hit or miss", "hit", hit, "error", err)
|
c.log.Warn("failed to count cache hit or miss", "hit", hit, "error", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
if fetchBytes > 0 {
|
c.IncrementUpstreamFetch(ctx, fetchBytes)
|
||||||
_, err = c.db.ExecContext(ctx, `
|
}
|
||||||
|
|
||||||
|
// IncrementUpstreamFetch counts one upstream fetch that read fetchBytes bytes.
|
||||||
|
// A fetch that read no bytes is not counted.
|
||||||
|
func (c *Cache) IncrementUpstreamFetch(ctx context.Context, fetchBytes int64) {
|
||||||
|
if fetchBytes <= 0 {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
_, err := c.db.ExecContext(ctx, `
|
||||||
UPDATE cache_stats
|
UPDATE cache_stats
|
||||||
SET upstream_fetch_count = upstream_fetch_count + 1,
|
SET upstream_fetch_count = upstream_fetch_count + 1,
|
||||||
upstream_fetch_bytes = upstream_fetch_bytes + ?,
|
upstream_fetch_bytes = upstream_fetch_bytes + ?,
|
||||||
@@ -508,7 +518,6 @@ func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64)
|
|||||||
"fetch_bytes", fetchBytes, "error", err)
|
"fetch_bytes", fetchBytes, "error", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
// IncrementTransformCount counts one image transcoded by the image processor.
|
// IncrementTransformCount counts one image transcoded by the image processor.
|
||||||
func (c *Cache) IncrementTransformCount(ctx context.Context) {
|
func (c *Cache) IncrementTransformCount(ctx context.Context) {
|
||||||
|
|||||||
@@ -0,0 +1,444 @@
|
|||||||
|
package imgcache
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"image/jpeg"
|
||||||
|
"io"
|
||||||
|
"io/fs"
|
||||||
|
"strings"
|
||||||
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"sneak.berlin/go/pixa/internal/httpfetcher"
|
||||||
|
"sneak.berlin/go/pixa/internal/magic"
|
||||||
|
)
|
||||||
|
|
||||||
|
// arrivalWait is how long a test gives requests it has started to reach the
|
||||||
|
// point where they wait for a held fetch.
|
||||||
|
const arrivalWait = 100 * time.Millisecond
|
||||||
|
|
||||||
|
// heldFetcher counts the fetches it is asked for and holds each one until
|
||||||
|
// releaseFetches is called, so that a test can have requests arrive while a
|
||||||
|
// fetch is in progress. started receives once for every fetch.
|
||||||
|
type heldFetcher struct {
|
||||||
|
upstream httpfetcher.Fetcher
|
||||||
|
fetches atomic.Int32
|
||||||
|
started chan struct{}
|
||||||
|
release chan struct{}
|
||||||
|
releaseOnce sync.Once
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *heldFetcher) Fetch(
|
||||||
|
ctx context.Context, url string,
|
||||||
|
) (*httpfetcher.FetchResult, error) {
|
||||||
|
f.fetches.Add(1)
|
||||||
|
|
||||||
|
f.started <- struct{}{}
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-f.release:
|
||||||
|
case <-ctx.Done():
|
||||||
|
return nil, ctx.Err()
|
||||||
|
}
|
||||||
|
|
||||||
|
return f.upstream.Fetch(ctx, url)
|
||||||
|
}
|
||||||
|
|
||||||
|
// releaseFetches lets every held fetch, and every later one, go on.
|
||||||
|
func (f *heldFetcher) releaseFetches() {
|
||||||
|
f.releaseOnce.Do(func() { close(f.release) })
|
||||||
|
}
|
||||||
|
|
||||||
|
// setupHeldFetchService returns a test service whose fetches go through a
|
||||||
|
// heldFetcher. Its database is limited to one connection: each connection to
|
||||||
|
// an in-memory SQLite database opens a new, empty one, so requests running at
|
||||||
|
// once must share the connection that holds the schema.
|
||||||
|
func setupHeldFetchService(t *testing.T) (*Service, *TestFixtures, *heldFetcher) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
svc, fixtures := SetupTestService(t)
|
||||||
|
svc.cache.db.SetMaxOpenConns(1)
|
||||||
|
|
||||||
|
fetcher := &heldFetcher{
|
||||||
|
upstream: svc.fetcher,
|
||||||
|
started: make(chan struct{}, 100),
|
||||||
|
release: make(chan struct{}),
|
||||||
|
}
|
||||||
|
svc.fetcher = fetcher
|
||||||
|
|
||||||
|
t.Cleanup(fetcher.releaseFetches)
|
||||||
|
|
||||||
|
return svc, fixtures, fetcher
|
||||||
|
}
|
||||||
|
|
||||||
|
// photoVariant asks for the test photo, 100x100, at 50x25 as a JPEG of the
|
||||||
|
// given quality and fit mode. Each call returns a new request, as Get writes
|
||||||
|
// to the request it is given.
|
||||||
|
func photoVariant(fixtures *TestFixtures, quality int, fit FitMode) *ImageRequest {
|
||||||
|
return &ImageRequest{
|
||||||
|
SourceHost: fixtures.GoodHost,
|
||||||
|
SourcePath: testPathPhoto,
|
||||||
|
Size: Size{Width: 50, Height: 25},
|
||||||
|
Format: FormatJPEG,
|
||||||
|
Quality: quality,
|
||||||
|
FitMode: fit,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// getResult is what one Get call returned, with the image read out.
|
||||||
|
type getResult struct {
|
||||||
|
image []byte
|
||||||
|
err error
|
||||||
|
}
|
||||||
|
|
||||||
|
// startGet calls Get in a goroutine of its own and delivers what it returned
|
||||||
|
// on the channel.
|
||||||
|
func startGet(
|
||||||
|
ctx context.Context, svc *Service, req *ImageRequest,
|
||||||
|
) <-chan getResult {
|
||||||
|
results := make(chan getResult, 1)
|
||||||
|
|
||||||
|
go func() {
|
||||||
|
resp, err := svc.Get(ctx, req)
|
||||||
|
if err != nil {
|
||||||
|
results <- getResult{err: err}
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
defer func() { _ = resp.Content.Close() }()
|
||||||
|
|
||||||
|
image, err := io.ReadAll(resp.Content)
|
||||||
|
results <- getResult{image: image, err: err}
|
||||||
|
}()
|
||||||
|
|
||||||
|
return results
|
||||||
|
}
|
||||||
|
|
||||||
|
// jpegSize returns the width and height of the JPEG image in data, or 0 and 0
|
||||||
|
// if data is not one.
|
||||||
|
func jpegSize(data []byte) (int, int) {
|
||||||
|
config, err := jpeg.DecodeConfig(bytes.NewReader(data))
|
||||||
|
if err != nil {
|
||||||
|
return 0, 0
|
||||||
|
}
|
||||||
|
|
||||||
|
return config.Width, config.Height
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestService_Get_ConcurrentMissesShareOneFetch starts several requests for
|
||||||
|
// one uncached variant while the first one's fetch is held. Between them they
|
||||||
|
// must fetch the source once and transcode it once, every one must be answered
|
||||||
|
// with the same 50x25 JPEG, and each must count one miss.
|
||||||
|
func TestService_Get_ConcurrentMissesShareOneFetch(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
svc, fixtures, fetcher := setupHeldFetchService(t)
|
||||||
|
|
||||||
|
const requests = 8
|
||||||
|
|
||||||
|
pending := make([]<-chan getResult, 0, requests)
|
||||||
|
for range requests {
|
||||||
|
pending = append(pending,
|
||||||
|
startGet(t.Context(), svc, photoVariant(fixtures, 85, FitCover)))
|
||||||
|
}
|
||||||
|
|
||||||
|
<-fetcher.started
|
||||||
|
time.Sleep(arrivalWait)
|
||||||
|
fetcher.releaseFetches()
|
||||||
|
|
||||||
|
var first []byte
|
||||||
|
|
||||||
|
for i, results := range pending {
|
||||||
|
got := <-results
|
||||||
|
if got.err != nil {
|
||||||
|
t.Fatalf("request %d: Get() error = %v", i, got.err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if width, height := jpegSize(got.image); width != 50 || height != 25 {
|
||||||
|
t.Errorf("request %d: image is %dx%d, want a 50x25 JPEG", i, width, height)
|
||||||
|
}
|
||||||
|
|
||||||
|
if first == nil {
|
||||||
|
first = got.image
|
||||||
|
} else if !bytes.Equal(got.image, first) {
|
||||||
|
t.Errorf("request %d: image differs from request 0's", i)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if fetches := fetcher.fetches.Load(); fetches != 1 {
|
||||||
|
t.Errorf("%d requests made %d upstream fetches, want 1", requests, fetches)
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewTestFS builds the same files the test service's fetcher serves.
|
||||||
|
testFS, _ := NewTestFS(t)
|
||||||
|
|
||||||
|
photo, err := fs.ReadFile(testFS, fixtures.GoodHostJPEG)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
want := cacheStatsCounters{0, requests, 1, int64(len(photo)), 1}
|
||||||
|
|
||||||
|
if got := readCacheStatsCounters(t, svc.cache); got != want {
|
||||||
|
t.Errorf("counters = %+v, want %+v", got, want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestService_Get_ConcurrentVariantsStayApart requests three variants of the
|
||||||
|
// test photo at once that differ only in quality or fit. Each must be made by
|
||||||
|
// a fetch and a transcode of its own, and each answer must be its own variant.
|
||||||
|
func TestService_Get_ConcurrentVariantsStayApart(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
svc, fixtures, fetcher := setupHeldFetchService(t)
|
||||||
|
|
||||||
|
cover := startGet(t.Context(), svc, photoVariant(fixtures, 85, FitCover))
|
||||||
|
lowQuality := startGet(t.Context(), svc, photoVariant(fixtures, 40, FitCover))
|
||||||
|
contain := startGet(t.Context(), svc, photoVariant(fixtures, 85, FitContain))
|
||||||
|
|
||||||
|
for range 3 {
|
||||||
|
select {
|
||||||
|
case <-fetcher.started:
|
||||||
|
case <-time.After(5 * time.Second):
|
||||||
|
t.Fatal("fewer fetches started than variants requested: " +
|
||||||
|
"variants differing in quality or fit were merged")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fetcher.releaseFetches()
|
||||||
|
|
||||||
|
images := make(map[string][]byte)
|
||||||
|
|
||||||
|
for _, variant := range []struct {
|
||||||
|
name string
|
||||||
|
results <-chan getResult
|
||||||
|
width, height int
|
||||||
|
}{
|
||||||
|
{"q=85 fit=cover", cover, 50, 25},
|
||||||
|
{"q=40 fit=cover", lowQuality, 50, 25},
|
||||||
|
{"q=85 fit=contain", contain, 25, 25},
|
||||||
|
} {
|
||||||
|
got := <-variant.results
|
||||||
|
if got.err != nil {
|
||||||
|
t.Fatalf("%s: Get() error = %v", variant.name, got.err)
|
||||||
|
}
|
||||||
|
|
||||||
|
width, height := jpegSize(got.image)
|
||||||
|
if width != variant.width || height != variant.height {
|
||||||
|
t.Errorf("%s: image is %dx%d, want a %dx%d JPEG", variant.name,
|
||||||
|
width, height, variant.width, variant.height)
|
||||||
|
}
|
||||||
|
|
||||||
|
images[variant.name] = got.image
|
||||||
|
}
|
||||||
|
|
||||||
|
if bytes.Equal(images["q=85 fit=cover"], images["q=40 fit=cover"]) {
|
||||||
|
t.Error("q=40 was answered with the q=85 image")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestService_Get_WaiterStopsWhenItsContextEnds has a second request for a
|
||||||
|
// variant join the first one's held fetch, then ends the second request's
|
||||||
|
// context. The second request must return at once with the context's error,
|
||||||
|
// while the fetch is still held, and the first must still be answered.
|
||||||
|
func TestService_Get_WaiterStopsWhenItsContextEnds(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
svc, fixtures, fetcher := setupHeldFetchService(t)
|
||||||
|
|
||||||
|
first := startGet(t.Context(), svc, photoVariant(fixtures, 85, FitCover))
|
||||||
|
|
||||||
|
<-fetcher.started
|
||||||
|
|
||||||
|
waiterCtx, cancelWaiter := context.WithCancel(t.Context())
|
||||||
|
waiter := startGet(waiterCtx, svc, photoVariant(fixtures, 85, FitCover))
|
||||||
|
|
||||||
|
time.Sleep(arrivalWait)
|
||||||
|
cancelWaiter()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case got := <-waiter:
|
||||||
|
if !errors.Is(got.err, context.Canceled) {
|
||||||
|
t.Errorf("waiting request: Get() error = %v, want %v",
|
||||||
|
got.err, context.Canceled)
|
||||||
|
}
|
||||||
|
case <-time.After(time.Second):
|
||||||
|
t.Fatal("waiting request did not return when its context ended")
|
||||||
|
}
|
||||||
|
|
||||||
|
fetcher.releaseFetches()
|
||||||
|
|
||||||
|
got := <-first
|
||||||
|
if got.err != nil {
|
||||||
|
t.Fatalf("first request: Get() error = %v", got.err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if width, height := jpegSize(got.image); width != 50 || height != 25 {
|
||||||
|
t.Errorf("first request: image is %dx%d, want a 50x25 JPEG", width, height)
|
||||||
|
}
|
||||||
|
|
||||||
|
if fetches := fetcher.fetches.Load(); fetches != 1 {
|
||||||
|
t.Errorf("upstream fetches = %d, want 1", fetches)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestService_Get_FirstRequestLeavingKeepsTheWork ends the context of the
|
||||||
|
// request whose fetch is held, after a second request has joined it. The fetch
|
||||||
|
// and transcode must go on and answer the second request. The first request
|
||||||
|
// waits for its own work, as every request did before misses were shared, and
|
||||||
|
// is answered too.
|
||||||
|
func TestService_Get_FirstRequestLeavingKeepsTheWork(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
svc, fixtures, fetcher := setupHeldFetchService(t)
|
||||||
|
|
||||||
|
firstCtx, cancelFirst := context.WithCancel(t.Context())
|
||||||
|
first := startGet(firstCtx, svc, photoVariant(fixtures, 85, FitCover))
|
||||||
|
|
||||||
|
<-fetcher.started
|
||||||
|
|
||||||
|
second := startGet(t.Context(), svc, photoVariant(fixtures, 85, FitCover))
|
||||||
|
|
||||||
|
time.Sleep(arrivalWait)
|
||||||
|
cancelFirst()
|
||||||
|
fetcher.releaseFetches()
|
||||||
|
|
||||||
|
for name, results := range map[string]<-chan getResult{
|
||||||
|
"first request": first, "second request": second,
|
||||||
|
} {
|
||||||
|
got := <-results
|
||||||
|
if got.err != nil {
|
||||||
|
t.Fatalf("%s: Get() error = %v", name, got.err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if width, height := jpegSize(got.image); width != 50 || height != 25 {
|
||||||
|
t.Errorf("%s: image is %dx%d, want a 50x25 JPEG", name, width, height)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if fetches := fetcher.fetches.Load(); fetches != 1 {
|
||||||
|
t.Errorf("upstream fetches = %d, want 1", fetches)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestService_Get_ConcurrentMissesShareAFailure has several requests for an
|
||||||
|
// image that cannot be served arrive while its fetch is held: the one fetch
|
||||||
|
// answers all of them with its error. The request after them is answered from
|
||||||
|
// the negative cache when the failure is kept there, and fetches again when it
|
||||||
|
// is not.
|
||||||
|
func TestService_Get_ConcurrentMissesShareAFailure(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
tests := []struct {
|
||||||
|
name string
|
||||||
|
path string
|
||||||
|
wantErr error // what every request at once gets
|
||||||
|
wantNextErr error // what the request after them gets
|
||||||
|
wantFetches int32 // fetches once the request after them is answered
|
||||||
|
}{
|
||||||
|
{"upstream answers 404, kept in the negative cache",
|
||||||
|
"/images/missing.jpg", httpfetcher.ErrUpstreamError,
|
||||||
|
ErrNegativeCached, 1},
|
||||||
|
{"source fails the magic byte check, not kept",
|
||||||
|
"/images/text.png", magic.ErrUnknownFormat, magic.ErrUnknownFormat, 2},
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, tc := range tests {
|
||||||
|
t.Run(tc.name, func(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
svc, fixtures, fetcher := setupHeldFetchService(t)
|
||||||
|
|
||||||
|
request := func() *ImageRequest {
|
||||||
|
req := photoVariant(fixtures, 85, FitCover)
|
||||||
|
req.SourcePath = tc.path
|
||||||
|
|
||||||
|
return req
|
||||||
|
}
|
||||||
|
|
||||||
|
pending := make([]<-chan getResult, 0, 4)
|
||||||
|
for range 4 {
|
||||||
|
pending = append(pending, startGet(t.Context(), svc, request()))
|
||||||
|
}
|
||||||
|
|
||||||
|
<-fetcher.started
|
||||||
|
time.Sleep(arrivalWait)
|
||||||
|
fetcher.releaseFetches()
|
||||||
|
|
||||||
|
for i, results := range pending {
|
||||||
|
if got := <-results; !errors.Is(got.err, tc.wantErr) {
|
||||||
|
t.Errorf("request %d: Get() error = %v, want %v", i, got.err, tc.wantErr)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
_, err := svc.Get(t.Context(), request())
|
||||||
|
if !errors.Is(err, tc.wantNextErr) {
|
||||||
|
t.Errorf("next request: Get() error = %v, want %v", err, tc.wantNextErr)
|
||||||
|
}
|
||||||
|
|
||||||
|
if fetches := fetcher.fetches.Load(); fetches != tc.wantFetches {
|
||||||
|
t.Errorf("upstream fetches = %d, want %d", fetches, tc.wantFetches)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestService_Get_EndedRequestFetchesNothing checks that a request whose
|
||||||
|
// context has already ended when it misses the cache starts no fetch.
|
||||||
|
func TestService_Get_EndedRequestFetchesNothing(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
svc, fixtures, fetcher := setupHeldFetchService(t)
|
||||||
|
|
||||||
|
ctx, cancel := context.WithCancel(t.Context())
|
||||||
|
cancel()
|
||||||
|
|
||||||
|
_, err := svc.Get(ctx, photoVariant(fixtures, 85, FitCover))
|
||||||
|
if !errors.Is(err, context.Canceled) {
|
||||||
|
t.Errorf("Get() error = %v, want %v", err, context.Canceled)
|
||||||
|
}
|
||||||
|
|
||||||
|
if fetches := fetcher.fetches.Load(); fetches != 0 {
|
||||||
|
t.Errorf("upstream fetches = %d, want 0", fetches)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// panickingFetcher panics on every fetch.
|
||||||
|
type panickingFetcher struct{}
|
||||||
|
|
||||||
|
func (panickingFetcher) Fetch(
|
||||||
|
context.Context, string,
|
||||||
|
) (*httpfetcher.FetchResult, error) {
|
||||||
|
panic("upstream fetcher panicked")
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestService_Get_PanicBecomesAnError checks that a panic while a variant is
|
||||||
|
// being made reaches its request as an error naming the panic, instead of
|
||||||
|
// being raised again.
|
||||||
|
func TestService_Get_PanicBecomesAnError(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
svc, fixtures := SetupTestService(t)
|
||||||
|
svc.fetcher = panickingFetcher{}
|
||||||
|
|
||||||
|
var err error
|
||||||
|
|
||||||
|
func() {
|
||||||
|
defer func() {
|
||||||
|
if recovered := recover(); recovered != nil {
|
||||||
|
t.Fatalf("Get() panicked: %v", recovered)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
_, err = svc.Get(t.Context(), photoVariant(fixtures, 85, FitCover))
|
||||||
|
}()
|
||||||
|
|
||||||
|
if err == nil || !strings.Contains(err.Error(), "upstream fetcher panicked") {
|
||||||
|
t.Errorf("Get() error = %v, want one naming the panic", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
+111
-27
@@ -8,9 +8,11 @@ import (
|
|||||||
"io"
|
"io"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"net/url"
|
"net/url"
|
||||||
|
"runtime/debug"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/dustin/go-humanize"
|
"github.com/dustin/go-humanize"
|
||||||
|
"golang.org/x/sync/singleflight"
|
||||||
"sneak.berlin/go/pixa/internal/allowlist"
|
"sneak.berlin/go/pixa/internal/allowlist"
|
||||||
"sneak.berlin/go/pixa/internal/httpfetcher"
|
"sneak.berlin/go/pixa/internal/httpfetcher"
|
||||||
"sneak.berlin/go/pixa/internal/imageprocessor"
|
"sneak.berlin/go/pixa/internal/imageprocessor"
|
||||||
@@ -29,6 +31,9 @@ type Service struct {
|
|||||||
log *slog.Logger
|
log *slog.Logger
|
||||||
allowHTTP bool
|
allowHTTP bool
|
||||||
maxResponseSize int64
|
maxResponseSize int64
|
||||||
|
// variantsInProgress lets the requests that miss the same variant at the
|
||||||
|
// same time share one fetch and one transcode.
|
||||||
|
variantsInProgress singleflight.Group
|
||||||
}
|
}
|
||||||
|
|
||||||
// ServiceConfig holds configuration for the image service.
|
// ServiceConfig holds configuration for the image service.
|
||||||
@@ -160,14 +165,12 @@ func (s *Service) Get(ctx context.Context, req *ImageRequest) (*ImageResponse, e
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Cache miss - process the cached source or fetch it, then count the
|
// Cache miss - get the variant, processed once for all the requests that
|
||||||
// miss with the bytes it fetched from upstream, also when it failed or
|
// miss it at the same time, then count this request's miss, also when it
|
||||||
// the request context has ended meanwhile
|
// failed or the request context has ended meanwhile
|
||||||
cacheKey := CacheKey(req)
|
response, err := s.processOrWait(ctx, req)
|
||||||
|
|
||||||
response, fetchedBytes, err := s.processFromSourceOrFetch(ctx, req, cacheKey)
|
s.cache.IncrementStats(context.WithoutCancel(ctx), false, 0)
|
||||||
|
|
||||||
s.cache.IncrementStats(context.WithoutCancel(ctx), false, fetchedBytes)
|
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -241,6 +244,83 @@ func (s *Service) GenerateSignedURL(
|
|||||||
baseURL, path, sig, exp, req.Quality, req.FitMode), nil
|
baseURL, path, sig, exp, req.Quality, req.FitMode), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// errPanicked is returned when processing a variant panicked.
|
||||||
|
var errPanicked = errors.New("panic while processing image")
|
||||||
|
|
||||||
|
// processOrWait returns the variant req asks for. The first of the requests
|
||||||
|
// that miss a variant at the same time processes it, and singleflight hands
|
||||||
|
// its result, or its error, to the others: they fetch nothing, read no source
|
||||||
|
// and take no upstream connection or processing slot. The processing runs with
|
||||||
|
// a context that does not end with the first request's, so the others are
|
||||||
|
// still served if that client goes away; the fetch timeout and the waits for a
|
||||||
|
// connection and a processing slot still bound it.
|
||||||
|
func (s *Service) processOrWait(
|
||||||
|
ctx context.Context, req *ImageRequest,
|
||||||
|
) (*ImageResponse, error) {
|
||||||
|
// A request that has already ended starts no processing
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
return nil, ctx.Err()
|
||||||
|
}
|
||||||
|
|
||||||
|
cacheKey := CacheKey(req)
|
||||||
|
|
||||||
|
// Closed when this request's own function runs, which singleflight does
|
||||||
|
// only when no other request is processing the variant
|
||||||
|
processing := make(chan struct{})
|
||||||
|
|
||||||
|
results := s.variantsInProgress.DoChan(string(cacheKey),
|
||||||
|
func() (_ any, err error) {
|
||||||
|
close(processing)
|
||||||
|
|
||||||
|
// singleflight would raise a panic again in a goroutine of its
|
||||||
|
// own, where no handler recovers it, and stop pixad
|
||||||
|
defer func() {
|
||||||
|
recovered := recover()
|
||||||
|
if recovered != nil {
|
||||||
|
s.log.Error("panic while processing image",
|
||||||
|
"host", req.SourceHost, "path", req.SourcePath,
|
||||||
|
"panic", recovered, "stack", string(debug.Stack()))
|
||||||
|
|
||||||
|
err = fmt.Errorf("%w: %v", errPanicked, recovered)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
return s.processFromSourceOrFetch(
|
||||||
|
context.WithoutCancel(ctx), req, cacheKey)
|
||||||
|
})
|
||||||
|
|
||||||
|
var result singleflight.Result
|
||||||
|
|
||||||
|
select {
|
||||||
|
case result = <-results:
|
||||||
|
case <-ctx.Done():
|
||||||
|
select {
|
||||||
|
case <-processing:
|
||||||
|
// This request is processing the variant: it waits for the
|
||||||
|
// result, as every request did before misses were shared
|
||||||
|
result = <-results
|
||||||
|
default:
|
||||||
|
// Another request is processing the variant, or this request's
|
||||||
|
// function has not started yet; the processing goes on without it
|
||||||
|
return nil, ctx.Err()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if result.Err != nil {
|
||||||
|
return nil, result.Err
|
||||||
|
}
|
||||||
|
|
||||||
|
variant, _ := result.Val.(*processedVariant)
|
||||||
|
|
||||||
|
return &ImageResponse{
|
||||||
|
Content: io.NopCloser(bytes.NewReader(variant.data)),
|
||||||
|
ContentLength: int64(len(variant.data)),
|
||||||
|
ContentType: variant.contentType,
|
||||||
|
FetchedBytes: variant.fetchedBytes,
|
||||||
|
ETag: formatETag(cacheKey),
|
||||||
|
}, nil
|
||||||
|
}
|
||||||
|
|
||||||
// loadCachedSource opens source content from cache, without reading it, and
|
// loadCachedSource opens source content from cache, without reading it, and
|
||||||
// returns it with its size; nil if the cached data is unavailable, empty or
|
// returns it with its size; nil if the cached data is unavailable, empty or
|
||||||
// exceeds maxResponseSize.
|
// exceeds maxResponseSize.
|
||||||
@@ -275,13 +355,12 @@ func (s *Service) loadCachedSource(
|
|||||||
}
|
}
|
||||||
|
|
||||||
// processFromSourceOrFetch processes an image, using cached source content
|
// processFromSourceOrFetch processes an image, using cached source content
|
||||||
// if available. It also returns the number of bytes fetched from upstream,
|
// if available.
|
||||||
// as fetchAndProcess does, or 0 when the cached source was used.
|
|
||||||
func (s *Service) processFromSourceOrFetch(
|
func (s *Service) processFromSourceOrFetch(
|
||||||
ctx context.Context,
|
ctx context.Context,
|
||||||
req *ImageRequest,
|
req *ImageRequest,
|
||||||
cacheKey VariantKey,
|
cacheKey VariantKey,
|
||||||
) (*ImageResponse, int64, error) {
|
) (*processedVariant, error) {
|
||||||
// Check if we have cached source content
|
// Check if we have cached source content
|
||||||
contentHash, _, err := s.cache.LookupSource(ctx, req)
|
contentHash, _, err := s.cache.LookupSource(ctx, req)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -308,19 +387,17 @@ func (s *Service) processFromSourceOrFetch(
|
|||||||
// Process using cached source; nothing was fetched from upstream. The
|
// Process using cached source; nothing was fetched from upstream. The
|
||||||
// image processor reads the source only once it has a processing slot,
|
// image processor reads the source only once it has a processing slot,
|
||||||
// so a request waiting for one holds none of it in memory.
|
// so a request waiting for one holds none of it in memory.
|
||||||
resp, err := s.processAndStore(ctx, req, cacheKey, source, sourceSize)
|
return s.processAndStore(ctx, req, cacheKey, source, sourceSize)
|
||||||
|
|
||||||
return resp, 0, err
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// fetchAndProcess fetches from upstream, processes, and caches the result.
|
// fetchAndProcess fetches from upstream, processes, and caches the result.
|
||||||
// It also returns the number of bytes read from upstream, including when
|
// It counts the fetch with the bytes read from upstream, including when
|
||||||
// reading the response or a later step fails.
|
// reading the response or a later step fails.
|
||||||
func (s *Service) fetchAndProcess(
|
func (s *Service) fetchAndProcess(
|
||||||
ctx context.Context,
|
ctx context.Context,
|
||||||
req *ImageRequest,
|
req *ImageRequest,
|
||||||
cacheKey VariantKey,
|
cacheKey VariantKey,
|
||||||
) (*ImageResponse, int64, error) {
|
) (*processedVariant, error) {
|
||||||
// Fetch from upstream
|
// Fetch from upstream
|
||||||
sourceURL := req.SourceURL()
|
sourceURL := req.SourceURL()
|
||||||
|
|
||||||
@@ -339,7 +416,7 @@ func (s *Service) fetchAndProcess(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil, 0, fmt.Errorf("upstream fetch failed: %w", err)
|
return nil, fmt.Errorf("upstream fetch failed: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Closing the body frees the upstream connection. It is closed only
|
// Closing the body frees the upstream connection. It is closed only
|
||||||
@@ -352,8 +429,11 @@ func (s *Service) fetchAndProcess(
|
|||||||
sourceData, err := io.ReadAll(fetchResult.Content)
|
sourceData, err := io.ReadAll(fetchResult.Content)
|
||||||
fetchBytes := int64(len(sourceData))
|
fetchBytes := int64(len(sourceData))
|
||||||
|
|
||||||
|
// Counted also when the request context has ended meanwhile
|
||||||
|
s.cache.IncrementUpstreamFetch(context.WithoutCancel(ctx), fetchBytes)
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fetchBytes, fmt.Errorf("failed to read upstream response: %w", err)
|
return nil, fmt.Errorf("failed to read upstream response: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Calculate download bitrate
|
// Calculate download bitrate
|
||||||
@@ -381,7 +461,7 @@ func (s *Service) fetchAndProcess(
|
|||||||
// Validate magic bytes match content type
|
// Validate magic bytes match content type
|
||||||
err = magic.ValidateMagicBytes(sourceData, fetchResult.ContentType)
|
err = magic.ValidateMagicBytes(sourceData, fetchResult.ContentType)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fetchBytes, fmt.Errorf("content validation failed: %w", err)
|
return nil, fmt.Errorf("content validation failed: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Store source content
|
// Store source content
|
||||||
@@ -391,11 +471,17 @@ func (s *Service) fetchAndProcess(
|
|||||||
// Continue even if caching fails
|
// Continue even if caching fails
|
||||||
}
|
}
|
||||||
|
|
||||||
resp, err := s.processAndStore(
|
return s.processAndStore(
|
||||||
ctx, req, cacheKey, bytes.NewReader(sourceData), fetchBytes,
|
ctx, req, cacheKey, bytes.NewReader(sourceData), fetchBytes,
|
||||||
)
|
)
|
||||||
|
}
|
||||||
|
|
||||||
return resp, fetchBytes, err
|
// processedVariant is a variant as processAndStore made it. Each request that
|
||||||
|
// shared its processing serves it through a reader of its own.
|
||||||
|
type processedVariant struct {
|
||||||
|
data []byte
|
||||||
|
contentType string
|
||||||
|
fetchedBytes int64
|
||||||
}
|
}
|
||||||
|
|
||||||
// processAndStore processes the image read from source and stores the
|
// processAndStore processes the image read from source and stores the
|
||||||
@@ -406,7 +492,7 @@ func (s *Service) processAndStore(
|
|||||||
cacheKey VariantKey,
|
cacheKey VariantKey,
|
||||||
source io.Reader,
|
source io.Reader,
|
||||||
fetchBytes int64,
|
fetchBytes int64,
|
||||||
) (*ImageResponse, error) {
|
) (*processedVariant, error) {
|
||||||
// Process the image
|
// Process the image
|
||||||
processStart := time.Now()
|
processStart := time.Now()
|
||||||
|
|
||||||
@@ -470,12 +556,10 @@ func (s *Service) processAndStore(
|
|||||||
// Continue even if caching fails
|
// Continue even if caching fails
|
||||||
}
|
}
|
||||||
|
|
||||||
return &ImageResponse{
|
return &processedVariant{
|
||||||
Content: io.NopCloser(bytes.NewReader(processedData)),
|
data: processedData,
|
||||||
ContentLength: outputSize,
|
contentType: processResult.ContentType,
|
||||||
ContentType: processResult.ContentType,
|
fetchedBytes: fetchBytes,
|
||||||
FetchedBytes: fetchBytes,
|
|
||||||
ETag: formatETag(cacheKey),
|
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,64 +0,0 @@
|
|||||||
package server
|
|
||||||
|
|
||||||
import (
|
|
||||||
"net/http"
|
|
||||||
"net/http/httptest"
|
|
||||||
"testing"
|
|
||||||
)
|
|
||||||
|
|
||||||
// TestCORSOnlyOnImageRoutes verifies that the image routes answer with the
|
|
||||||
// configured access_control_allow_origin, a preflight request included, and
|
|
||||||
// that the login and URL generator pages send no Access-Control-Allow-Origin,
|
|
||||||
// so no other site can read them. /metrics is left out: its middleware
|
|
||||||
// registers with the process-wide Prometheus registry, which only one test
|
|
||||||
// in this package can do.
|
|
||||||
func TestCORSOnlyOnImageRoutes(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
const appOrigin = "https://app.example.com"
|
|
||||||
|
|
||||||
s := newTestServer(t)
|
|
||||||
s.config.AccessControlAllowOrigin = appOrigin
|
|
||||||
s.SetupRoutes()
|
|
||||||
|
|
||||||
requests := []struct {
|
|
||||||
method string
|
|
||||||
path string
|
|
||||||
want string
|
|
||||||
}{
|
|
||||||
{http.MethodGet, unsignedImagePath, appOrigin},
|
|
||||||
{http.MethodHead, unsignedImagePath, appOrigin},
|
|
||||||
{http.MethodOptions, unsignedImagePath, appOrigin},
|
|
||||||
{http.MethodGet, encryptedImagePath, appOrigin},
|
|
||||||
{http.MethodGet, "/", ""},
|
|
||||||
{http.MethodOptions, "/", ""},
|
|
||||||
{http.MethodPost, "/generate", ""},
|
|
||||||
{http.MethodGet, "/logout", ""},
|
|
||||||
}
|
|
||||||
|
|
||||||
for _, tc := range requests {
|
|
||||||
t.Run(tc.method+" "+tc.path, func(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
req := httptest.NewRequestWithContext(
|
|
||||||
t.Context(), tc.method, tc.path, nil)
|
|
||||||
req.Header.Set("Origin", appOrigin)
|
|
||||||
|
|
||||||
// An OPTIONS request naming the method it asks about is the
|
|
||||||
// preflight a browser sends before some cross-origin requests.
|
|
||||||
if tc.method == http.MethodOptions {
|
|
||||||
req.Header.Set("Access-Control-Request-Method", http.MethodGet)
|
|
||||||
}
|
|
||||||
|
|
||||||
rec := httptest.NewRecorder()
|
|
||||||
s.ServeHTTP(rec, req)
|
|
||||||
t.Logf("status %d", rec.Code)
|
|
||||||
|
|
||||||
got := rec.Header().Get("Access-Control-Allow-Origin")
|
|
||||||
if got != tc.want {
|
|
||||||
t.Errorf("Access-Control-Allow-Origin = %q, want %q",
|
|
||||||
got, tc.want)
|
|
||||||
}
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -13,10 +13,6 @@ import (
|
|||||||
// unsignedImagePath is an image URL that carries no signature.
|
// unsignedImagePath is an image URL that carries no signature.
|
||||||
const unsignedImagePath = "/v1/image/cdn.example.com/cat.jpg/100x100.jpeg"
|
const unsignedImagePath = "/v1/image/cdn.example.com/cat.jpg/100x100.jpeg"
|
||||||
|
|
||||||
// encryptedImagePath is an encrypted image URL whose token cannot be
|
|
||||||
// decrypted.
|
|
||||||
const encryptedImagePath = "/v1/e/token/cat.jpg"
|
|
||||||
|
|
||||||
// TestMaintenanceModeRefusesImageRequests verifies that while maintenance
|
// TestMaintenanceModeRefusesImageRequests verifies that while maintenance
|
||||||
// mode is on, both image routes answer 503 Service Unavailable with a
|
// mode is on, both image routes answer 503 Service Unavailable with a
|
||||||
// Retry-After header and the JSON error body the image handlers send.
|
// Retry-After header and the JSON error body the image handlers send.
|
||||||
@@ -32,7 +28,7 @@ func TestMaintenanceModeRefusesImageRequests(t *testing.T) {
|
|||||||
}{
|
}{
|
||||||
{http.MethodGet, unsignedImagePath},
|
{http.MethodGet, unsignedImagePath},
|
||||||
{http.MethodHead, unsignedImagePath},
|
{http.MethodHead, unsignedImagePath},
|
||||||
{http.MethodGet, encryptedImagePath},
|
{http.MethodGet, "/v1/e/token/cat.jpg"},
|
||||||
}
|
}
|
||||||
|
|
||||||
for _, tc := range requests {
|
for _, tc := range requests {
|
||||||
@@ -99,7 +95,7 @@ func TestImageRequestsServedWithoutMaintenanceMode(t *testing.T) {
|
|||||||
}{
|
}{
|
||||||
{http.MethodGet, unsignedImagePath, http.StatusUnauthorized},
|
{http.MethodGet, unsignedImagePath, http.StatusUnauthorized},
|
||||||
{http.MethodHead, unsignedImagePath, http.StatusUnauthorized},
|
{http.MethodHead, unsignedImagePath, http.StatusUnauthorized},
|
||||||
{http.MethodGet, encryptedImagePath, http.StatusBadRequest},
|
{http.MethodGet, "/v1/e/token/cat.jpg", http.StatusBadRequest},
|
||||||
}
|
}
|
||||||
|
|
||||||
for _, tc := range requests {
|
for _, tc := range requests {
|
||||||
|
|||||||
@@ -38,6 +38,7 @@ func (s *Server) SetupRoutes() {
|
|||||||
s.router.Use(s.mw.Metrics())
|
s.router.Use(s.mw.Metrics())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
s.router.Use(s.mw.CORS())
|
||||||
s.router.Use(middleware.Timeout(s.config.DownstreamTimeout))
|
s.router.Use(middleware.Timeout(s.config.DownstreamTimeout))
|
||||||
|
|
||||||
if s.sentryEnabled {
|
if s.sentryEnabled {
|
||||||
@@ -73,31 +74,22 @@ func (s *Server) SetupRoutes() {
|
|||||||
|
|
||||||
s.router.Get("/logout", s.h.HandleLogout())
|
s.router.Get("/logout", s.h.HandleLogout())
|
||||||
|
|
||||||
// Image routes, the only ones that send CORS headers, as pages on other
|
// Image routes, refused while maintenance mode is on. Only these: the
|
||||||
// sites read them. They are a subrouter rather than a group: a group's
|
// image's Docker HEALTHCHECK requests the health check, a 503 there
|
||||||
// middleware runs only for a request that matches one of its routes,
|
// would make the container unhealthy, and upaas marks a deploy failed
|
||||||
// and a browser's preflight OPTIONS request matches none, so the CORS
|
|
||||||
// middleware could not answer it.
|
|
||||||
s.router.Route("/v1", func(r chi.Router) {
|
|
||||||
r.Use(s.mw.CORS())
|
|
||||||
|
|
||||||
// Refused while maintenance mode is on. Only these: the image's
|
|
||||||
// Docker HEALTHCHECK requests the health check, a 503 there would
|
|
||||||
// make the container unhealthy, and upaas marks a deploy failed
|
|
||||||
// when its container is unhealthy.
|
// when its container is unhealthy.
|
||||||
r.Group(func(r chi.Router) {
|
s.router.Group(func(r chi.Router) {
|
||||||
r.Use(s.refuseDuringMaintenance)
|
r.Use(s.refuseDuringMaintenance)
|
||||||
|
|
||||||
// Main image proxy route
|
// Main image proxy route
|
||||||
// /v1/image/<host>/<path>/<width>x<height>.<format>
|
// /v1/image/<host>/<path>/<width>x<height>.<format>
|
||||||
r.Get("/image/*", s.h.HandleImage())
|
r.Get("/v1/image/*", s.h.HandleImage())
|
||||||
r.Head("/image/*", s.h.HandleImage())
|
r.Head("/v1/image/*", s.h.HandleImage())
|
||||||
|
|
||||||
// Encrypted image URL route
|
// Encrypted image URL route
|
||||||
// The trailing filename (e.g., /img.jpg) is ignored but helps
|
// The trailing filename (e.g., /img.jpg) is ignored but helps
|
||||||
// browsers with content type
|
// browsers with content type
|
||||||
r.Get("/e/{token}/*", s.h.HandleImageEnc())
|
r.Get("/v1/e/{token}/*", s.h.HandleImageEnc())
|
||||||
})
|
|
||||||
})
|
})
|
||||||
|
|
||||||
// Metrics endpoint with auth
|
// Metrics endpoint with auth
|
||||||
|
|||||||
Reference in New Issue
Block a user