2 Commits
Author SHA1 Message Date
clawbot 6d6c76937b Bound concurrent image processing and upstream fetches (closes #64)
check / check (push) Successful in 4m3s
max_concurrent_processing (default: the number of CPUs Go uses) bounds
the images processed at once, and upstream_connections (default 64) the
fetches from all upstream hosts together, beside the per-host limit. A
request that finds either full waits up to 10 seconds, then gets 503
"server busy, try again later". The processor holds its slot from before
it reads the input until it returns, and takes a free slot even after the
request context has ended; a fetch holds its connection until the
response body is closed, after its image is processed. libvips now starts
with one worker thread per image and no operation cache. Both settings
have PIXA_ variables and are in README.md and config.example.yml.

Model: opus-5-5
2026-09-29 03:31:44 +00:00
clawbot 23ff05e293 Test the processing and upstream connection limits (closes #64)
Failing tests for two limits that do not exist yet. Config:
max_concurrent_processing and upstream_connections, their defaults, and
valid and invalid values from the file and the environment. Image
processor: never more images at once than its limit, waiting and then
failing with ErrTooManyImages when no slot frees, and freeing its slot on
every error. Fetcher: connections to all hosts counted together, apart
from the per-host limit, and freed on errors. Both image routes answer
503 when either wait gives up. TestEnvironmentSetsEveryKey sets the two
new variables, as it compares the whole config. The tests do not compile
until the limits exist.

Model: opus-5-5
2026-09-29 03:26:28 +00:00
7 changed files with 55 additions and 267 deletions
+16 -37
View File
@@ -112,19 +112,13 @@ copy is fresh.
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
address comes from `X-Forwarded-For` only when the address pixa sees for address comes from `X-Forwarded-For` only when the proxy's address is in
requests that come through the proxy is in `trusted_proxies`; otherwise all `trusted_proxies`; otherwise all users behind the proxy are counted as one
users behind the proxy are counted as one client. That address is not always client. With the default `trusted_proxies` (the RFC 1918 ranges), a client
the proxy's own: a proxy on the Docker host that connects to pixa over with a private address can choose the address it is counted by through its own
`127.0.0.1` is seen as the gateway of the container's Docker network, such as `X-Forwarded-For`, whether it connects directly or through the proxy, because
`172.17.0.1` on the default bridge, and one that connects through another of the its own address is trusted too. Setting `trusted_proxies` to the proxy's own
host's addresses is seen with that address. To be sure, read it as `remoteIP` in address closes this.
pixa's request log while it is not in `trusted_proxies` (see `trusted_proxies`
under Configuration). With the default `trusted_proxies` (the RFC 1918 ranges),
a client with a private address can choose the address it is counted by through
its own `X-Forwarded-For`, whether it connects directly or through the proxy,
because its own address is trusted too. Setting `trusted_proxies` to only the
address pixa sees for requests that come through the proxy closes this.
### Image Metadata ### Image Metadata
@@ -178,27 +172,19 @@ Where:
outside), or `cover` when the URL has no `fit`; a request whose `fit` is outside), or `cover` when the URL has no `fit`; a request whose `fit` is
anything else, an empty `fit=` included, is refused with 400 anything else, an empty `fit=` included, is refused with 400
The URL's `sig` is the HMAC-SHA256 result in base64url (the URL-safe alphabet **Example:** resize `https://cdn.example.com/photos/cat.jpg` to 800x600
of RFC 4648) with the trailing `=` padding kept, 44 characters in all. pixa WebP with expiration 1704067200, default quality and fit:
compares it exactly, so a signature encoded without padding, as Node's
`base64url` and Go's `base64.RawURLEncoding` do, is refused with 401.
**Example:** with the signing key `example-signing-key-for-documentation`,
resize `https://cdn.example.com/photos/cat.jpg` to 800x600 WebP with
expiration 1704067200, default quality and fit:
1. Build input: 1. Build input:
`cdn.example.com:/photos/cat.jpg::800:600:webp:1704067200:85:cover` `cdn.example.com:/photos/cat.jpg::800:600:webp:1704067200:85:cover`
2. Compute HMAC-SHA256 of it with the signing key 2. Compute HMAC-SHA256 with your secret key
3. Base64URL-encode the result, keeping the `=` padding: 3. Base64URL-encode the result
`-ay7KHpfqmtIGbibDGbUuBDkymi-Ymdn0NkC6j5EJag=`
4. URL: 4. URL:
`/v1/image/cdn.example.com/photos/cat.jpg/800x600.webp?sig=-ay7KHpfqmtIGbibDGbUuBDkymi-Ymdn0NkC6j5EJag=&exp=1704067200` `/v1/image/cdn.example.com/photos/cat.jpg/800x600.webp?sig=<base64url>&exp=1704067200`
For the same image at quality 40 with fit `contain`, the input ends in For the same image at quality 40 with fit `contain`, the input ends in
`:40:contain`, the signature is `5IwXUx6vf7yefhaUvFzgXZvG2o0Df4RJxPTK3pKq5VU=`, `:40:contain` and the URL is
and the URL is `/v1/image/cdn.example.com/photos/cat.jpg/800x600.webp?sig=<base64url>&exp=1704067200&q=40&fit=contain`.
`/v1/image/cdn.example.com/photos/cat.jpg/800x600.webp?sig=5IwXUx6vf7yefhaUvFzgXZvG2o0Df4RJxPTK3pKq5VU=&exp=1704067200&q=40&fit=contain`.
**Allowlist patterns:** **Allowlist patterns:**
@@ -262,15 +248,8 @@ Key settings in more detail:
`172.16.0.0/12`, `192.168.0.0/16`), since pixa is deployed behind a `172.16.0.0/12`, `192.168.0.0/16`), since pixa is deployed behind a
proxy on a private network; an explicitly empty list (`[]`) trusts no proxy on a private network; an explicitly empty list (`[]`) trusts no
one, and an explicit list replaces the default. An invalid CIDR aborts one, and an explicit list replaces the default. An invalid CIDR aborts
startup. Set this to the address pixa sees for requests that come through startup. Set this to your proxy's address range if it is not already
your proxy, such as `172.17.0.1/32`, when the defaults do not cover it, or covered by the defaults
to trust nothing else (see the login limit under Routes). For a proxy on
the Docker host that connects to pixa over `127.0.0.1`, that address is the
gateway of the container's Docker network (`172.17.0.1` on the default
bridge), not the proxy's own address; a proxy that connects through another of
the host's addresses is seen with that address. To be sure which address it
is, set this to `[]` (or `PIXA_TRUSTED_PROXIES` to empty), send a request
through the proxy, and read `remoteIP` in pixa's request log line for it
- `upstream_fetch_timeout` — timeout for origin requests - `upstream_fetch_timeout` — timeout for origin requests
- `upstream_max_response_size` — max origin response size - `upstream_max_response_size` — max origin response size
- `downstream_timeout` — client response timeout - `downstream_timeout` — client response timeout
+4 -13
View File
@@ -34,19 +34,10 @@ P2: security: referer blacklist
limits the images decoded and encoded at once, and `upstream_connections` limits the images decoded and encoded at once, and `upstream_connections`
(default 64) the connections to all upstream hosts together, on top of (default 64) the connections to all upstream hosts together, on top of
`upstream_connections_per_host`; a fetch holds its connection until its image `upstream_connections_per_host`; a fetch holds its connection until its image
has been processed, and a request whose source is cached reads it only once it has been processed; a request that finds either limit reached waits up to 10
has a processing slot; a request that finds either limit reached waits up to seconds for a free one, then gets 503 `server busy, try again later`; libvips
10 seconds for a free one, then gets 503 `server busy, try again later`; runs one worker thread per image with its operation cache off; documented in
libvips runs one worker thread per image with its operation cache off; `README.md` and `config.example.yml`.
documented in `README.md` and `config.example.yml`.
- 2026-09-29 `trusted_proxies` advice and signature padding in `README.md`
(closes #150): the login-limit paragraph, the `trusted_proxies` entry and
`config.example.yml` say to set `trusted_proxies` to the address pixa sees for
requests that come through the proxy, which the request log shows as
`remoteIP` while it is not trusted; for a proxy on the Docker host that
connects over `127.0.0.1` that is the Docker network's gateway, not the
proxy's own address; the signature section says `sig` is base64url with the
`=` padding kept, and gives the example's `sig` for a stated signing key.
- 2026-09-29 fixed uid and gid for `pixad` (closes #151): the image creates the - 2026-09-29 fixed uid and gid for `pixad` (closes #151): the image creates the
`pixad` group with gid 65532 and the `pixad` user with uid 65532, instead of `pixad` group with gid 65532 and the `pixad` user with uid 65532, instead of
the first free uid 1000, so a bind-mounted `/var/lib/pixa` given to `pixad` the first free uid 1000, so a bind-mounted `/var/lib/pixa` given to `pixad`
+1 -7
View File
@@ -50,13 +50,7 @@ allowlist_hosts:
# 172.16.0.0/12, 192.168.0.0/16), since pixa is deployed behind a proxy on # 172.16.0.0/12, 192.168.0.0/16), since pixa is deployed behind a proxy on
# a private network. An explicitly empty list ([]) trusts no one; an # a private network. An explicitly empty list ([]) trusts no one; an
# explicit list replaces the default. An invalid CIDR aborts startup. # explicit list replaces the default. An invalid CIDR aborts startup.
# Uncomment to override the defaults with the address pixa sees for # Uncomment to override the defaults with your proxy's address range.
# requests that come through your proxy. That is not always the proxy's own
# address: a proxy on the Docker host that connects over 127.0.0.1 is seen
# as the gateway of the container's Docker network (172.17.0.1 on the
# default bridge), and one that connects through another host address is
# seen with that address. To be sure, look it up in the request log as the
# trusted_proxies entry in README.md describes.
# trusted_proxies: # trusted_proxies:
# - 10.0.0.0/8 # - 10.0.0.0/8
# - 2001:db8::/32 # - 2001:db8::/32
@@ -1,7 +1,6 @@
package httpfetcher package httpfetcher
import ( import (
"context"
"errors" "errors"
"net" "net"
"strconv" "strconv"
@@ -83,42 +82,6 @@ func TestFetchLimitsConnectionsToAllHostsTogether(t *testing.T) {
_ = third.Content.Close() _ = third.Content.Close()
} }
// TestFetchFreesHostSlotWhenContextEndsWaitingForConnection checks that a
// fetch whose request context ends while it waits for a connection shared
// by all hosts gives its host's slot back. With MaxConnections at 1 and one
// response open, a fetch from another host takes that host's slot and waits;
// its context ends long before the 10 second wait timeout.
func TestFetchFreesHostSlotWhenContextEndsWaitingForConnection(t *testing.T) {
t.Parallel()
srv := startUpstream(t)
cfg := DefaultConfig()
cfg.MaxConnections = 1
f, _ := newServerFetcher(t, srv, cfg)
first, err := f.Fetch(testContext(t), imageURLOnPort(81))
if err != nil {
t.Fatalf("first Fetch() error = %v", err)
}
defer func() { _ = first.Content.Close() }()
ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond)
defer cancel()
_, err = f.Fetch(ctx, imageURLOnPort(82))
if !errors.Is(err, context.DeadlineExceeded) {
t.Fatalf("second Fetch() error = %v, want context.DeadlineExceeded", err)
}
if held := semLen(f, testPublicHost+":82"); held != 0 {
t.Errorf("the fetch kept its host's slot after its context ended: "+
"%d held", held)
}
}
// TestFetchReleasesConnectionOnError checks that a fetch that fails after // TestFetchReleasesConnectionOnError checks that a fetch that fails after
// taking its connection gives it back: with MaxConnections at 1, the slot // taking its connection gives it back: with MaxConnections at 1, the slot
// must be free after the failure and the next fetch must succeed. // must be free after the failure and the next fetch must succeed.
+4 -7
View File
@@ -383,16 +383,13 @@ func (c *Cache) GetSourceMetadataID(
return id, nil return id, nil
} }
// GetSourceContent returns a reader for cached source content by its hash, // GetSourceContent returns a reader for cached source content by its hash.
// and the content's size in bytes. func (c *Cache) GetSourceContent(contentHash ContentHash) (io.ReadCloser, error) {
func (c *Cache) GetSourceContent(
contentHash ContentHash,
) (io.ReadCloser, int64, error) {
if c.disabled { if c.disabled {
return nil, 0, ErrNotFound return nil, ErrNotFound
} }
return c.srcContent.LoadWithSize(contentHash) return c.srcContent.Load(contentHash)
} }
// CleanExpired removes expired entries from the cache. // CleanExpired removes expired entries from the cache.
@@ -1,125 +0,0 @@
package imgcache
import (
"image/color"
"image/jpeg"
"io"
"os"
"testing"
"time"
"sneak.berlin/go/pixa/internal/imageprocessor"
)
// widthOnlyRequest asks for the test photo at width, its height scaled to
// keep the photo's aspect ratio.
func widthOnlyRequest(fixtures *TestFixtures, width int) *ImageRequest {
return &ImageRequest{
SourceHost: fixtures.GoodHost,
SourcePath: testPathPhoto,
Size: Size{Width: width},
Format: FormatJPEG,
Quality: 85,
FitMode: FitCover,
}
}
// holdProcessingSlot takes one of proc's processing slots and returns the
// func that gives it back. Process takes its slot before it reads its input,
// so once it has read a byte from the pipe it holds the slot, until the pipe
// is closed.
func holdProcessingSlot(
t *testing.T, proc *imageprocessor.ImageProcessor,
) func() {
t.Helper()
input, feed := io.Pipe()
go func() {
_, _ = proc.Process(t.Context(), input, &imageprocessor.Request{})
}()
_, err := feed.Write([]byte{0})
if err != nil {
t.Fatalf("Process call to hold the slot did not start: %v", err)
}
release := func() { _ = feed.Close() }
t.Cleanup(release)
return release
}
// TestService_Get_WaitsForSlotBeforeReadingCachedSource checks that a
// request whose source is cached holds none of it while it waits for a
// processing slot: it reads the cached file only once it has a slot. With
// the only slot held, a request for a new width of the cached 100x100 photo
// waits; the cached file is then rewritten as a 100x50 image before the slot
// is freed, so the request must answer with that image scaled to 40x20.
func TestService_Get_WaitsForSlotBeforeReadingCachedSource(t *testing.T) {
t.Parallel()
svc, fixtures := SetupTestService(t)
svc.processor = imageprocessor.New(
imageprocessor.Params{MaxConcurrentProcessing: 1},
)
// A first request caches the photo as a source.
resp, err := svc.Get(t.Context(), widthOnlyRequest(fixtures, 50))
if err != nil {
t.Fatalf("first Get() error = %v", err)
}
_ = resp.Content.Close()
contentHash, _, err := svc.cache.LookupSource(t.Context(),
widthOnlyRequest(fixtures, 50))
if err != nil || contentHash == "" {
t.Fatalf("LookupSource() = %q, %v; want the cached source",
contentHash, err)
}
release := holdProcessingSlot(t, svc.processor)
var (
waited *ImageResponse
waitedErr error
)
done := make(chan struct{})
go func() {
defer close(done)
waited, waitedErr = svc.Get(t.Context(), widthOnlyRequest(fixtures, 40))
}()
// Give the request time to reach the slot: had it read the cached source
// before waiting, it would have read it by now.
time.Sleep(100 * time.Millisecond)
err = os.WriteFile(svc.cache.srcContent.hashToPath(contentHash),
generateTestJPEG(t, 100, 50, color.RGBA{0, 0, 255, 255}), 0o600)
if err != nil {
t.Fatalf("failed to rewrite the cached source: %v", err)
}
release()
<-done
if waitedErr != nil {
t.Fatalf("Get() error = %v", waitedErr)
}
defer func() { _ = waited.Content.Close() }()
output, err := jpeg.DecodeConfig(waited.Content)
if err != nil {
t.Fatalf("failed to decode the response: %v", err)
}
if output.Width != 40 || output.Height != 20 {
t.Errorf("response is %dx%d, want 40x20: the request read the cached "+
"source before it had a processing slot", output.Width, output.Height)
}
}
+30 -41
View File
@@ -241,37 +241,38 @@ func (s *Service) GenerateSignedURL(
baseURL, path, sig, exp, req.Quality, req.FitMode), nil baseURL, path, sig, exp, req.Quality, req.FitMode), nil
} }
// loadCachedSource opens source content from cache, without reading it, and // loadCachedSource attempts to load source content from cache, returning nil
// returns it with its size; nil if the cached data is unavailable, empty or // if the cached data is unavailable or exceeds maxResponseSize.
// exceeds maxResponseSize. func (s *Service) loadCachedSource(contentHash ContentHash) []byte {
func (s *Service) loadCachedSource( reader, err := s.cache.GetSourceContent(contentHash)
contentHash ContentHash,
) (io.ReadCloser, int64) {
reader, size, err := s.cache.GetSourceContent(contentHash)
if err != nil { if err != nil {
s.log.Warn("failed to load cached source, fetching", "error", err) s.log.Warn("failed to load cached source, fetching", "error", err)
return nil, 0 return nil
} }
if size > s.maxResponseSize { // Bound the read to maxResponseSize to prevent unbounded memory use
_ = reader.Close() // from unexpectedly large cached files.
limited := io.LimitReader(reader, s.maxResponseSize+1)
data, err := io.ReadAll(limited)
_ = reader.Close()
if err != nil {
s.log.Warn("failed to read cached source, fetching", "error", err)
return nil
}
if int64(len(data)) > s.maxResponseSize {
s.log.Warn("cached source exceeds max response size, discarding", s.log.Warn("cached source exceeds max response size, discarding",
"hash", contentHash, "hash", contentHash,
"max_bytes", s.maxResponseSize, "max_bytes", s.maxResponseSize,
) )
return nil, 0 return nil
} }
if size == 0 { return data
_ = reader.Close()
return nil, 0
}
return reader, size
} }
// processFromSourceOrFetch processes an image, using cached source content // processFromSourceOrFetch processes an image, using cached source content
@@ -288,27 +289,22 @@ func (s *Service) processFromSourceOrFetch(
s.log.Warn("source lookup failed", "error", err) s.log.Warn("source lookup failed", "error", err)
} }
var ( var sourceData []byte
source io.ReadCloser
sourceSize int64
)
if contentHash != "" { if contentHash != "" {
s.log.Debug("using cached source", "hash", contentHash) s.log.Debug("using cached source", "hash", contentHash)
source, sourceSize = s.loadCachedSource(contentHash) sourceData = s.loadCachedSource(contentHash)
} }
// Fetch from upstream if we don't have source data or it's empty // Fetch from upstream if we don't have source data or it's empty
if source == nil { if len(sourceData) == 0 {
return s.fetchAndProcess(ctx, req, cacheKey) return s.fetchAndProcess(ctx, req, cacheKey)
} }
defer func() { _ = source.Close() }() // Process using cached source; nothing was fetched from upstream
resp, err := s.processAndStore(
// Process using cached source; nothing was fetched from upstream. The ctx, req, cacheKey, sourceData, int64(len(sourceData)),
// image processor reads the source only once it has a processing slot, )
// so a request waiting for one holds none of it in memory.
resp, err := s.processAndStore(ctx, req, cacheKey, source, sourceSize)
return resp, 0, err return resp, 0, err
} }
@@ -342,10 +338,6 @@ func (s *Service) fetchAndProcess(
return nil, 0, fmt.Errorf("upstream fetch failed: %w", err) return nil, 0, fmt.Errorf("upstream fetch failed: %w", err)
} }
// Closing the body frees the upstream connection. It is closed only
// after processing, so the fetcher's connection limit also bounds the
// fetched sources held in memory while their requests wait for a
// processing slot.
defer func() { _ = fetchResult.Content.Close() }() defer func() { _ = fetchResult.Content.Close() }()
// Read and validate the source content // Read and validate the source content
@@ -391,20 +383,17 @@ func (s *Service) fetchAndProcess(
// Continue even if caching fails // Continue even if caching fails
} }
resp, err := s.processAndStore( resp, err := s.processAndStore(ctx, req, cacheKey, sourceData, fetchBytes)
ctx, req, cacheKey, bytes.NewReader(sourceData), fetchBytes,
)
return resp, fetchBytes, err return resp, fetchBytes, err
} }
// processAndStore processes the image read from source and stores the // processAndStore processes an image and stores the result.
// result.
func (s *Service) processAndStore( func (s *Service) processAndStore(
ctx context.Context, ctx context.Context,
req *ImageRequest, req *ImageRequest,
cacheKey VariantKey, cacheKey VariantKey,
source io.Reader, sourceData []byte,
fetchBytes int64, fetchBytes int64,
) (*ImageResponse, error) { ) (*ImageResponse, error) {
// Process the image // Process the image
@@ -417,7 +406,7 @@ func (s *Service) processAndStore(
FitMode: imageprocessor.FitMode(req.FitMode), FitMode: imageprocessor.FitMode(req.FitMode),
} }
processResult, err := s.processor.Process(ctx, source, processReq) processResult, err := s.processor.Process(ctx, bytes.NewReader(sourceData), processReq)
if err != nil { if err != nil {
return nil, fmt.Errorf("image processing failed: %w", err) return nil, fmt.Errorf("image processing failed: %w", err)
} }