5 Commits
Author SHA1 Message Date
clawbot ec86b964d5 Read a cached source only once a processing slot is taken (closes #64)
check / check (push) Successful in 3m27s
A request whose source was in the disk cache read the whole file into
memory, then waited for a processing slot, so a burst of new sizes for
one large cached image held one copy per waiting request, with no
ceiling. The service now opens the cached file and hands it to the image
processor, which reads it only after taking its slot. The file's size,
now returned by GetSourceContent, still sends an empty or oversized
cached source to upstream instead. A cached file that fails while being
read now fails the request instead of being fetched again.

Model: opus-5-5
2026-09-29 04:10:55 +00:00
clawbot d57379414b Test that a cached source is read only with a processing slot (closes #64)
Failing test: with the only processing slot held, a request for a new
width of a cached image waits for the slot while the cached file is
rewritten; the answer must come from the rewritten file, so the request
read none of the source before it had a slot. Also a test that a fetch
whose request context ends while it waits for a connection shared by all
hosts gives its host's slot back; that one passes already.

Model: opus-5-5
2026-09-29 04:08:58 +00:00
clawbot 5a20a5f48d Bound concurrent image processing and upstream fetches (closes #64)
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 04:08:58 +00:00
clawbot 7472249394 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 04:08:38 +00:00
clawbot 1b920fe000 Correct trusted_proxies advice and state signature padding (closes #150)
check / check (push) Successful in 14s
The README told operators to set trusted_proxies to the proxy's own
address. A proxy on the Docker host that connects over 127.0.0.1 reaches
pixa from the Docker network's gateway, so that advice made pixa count
every user as one client for the login limit. The login-limit paragraph,
the trusted_proxies entry and config.example.yml now say to use the
address pixa sees for requests through the proxy, that a proxy connecting
through another host address is seen with that address, and how to read
it from the request log.

The signature section now says sig is base64url with the = padding
kept, since pixa compares it exactly, and shows the example's sig for a
stated key, computed with pixa's signer.

Model: opus-5-5
2026-09-29 06:05:44 +02:00
7 changed files with 267 additions and 55 deletions
+37 -16
View File
@@ -112,13 +112,19 @@ copy is fresh.
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
refused with 429 and a `Retry-After` header. Behind a reverse proxy the client
address comes from `X-Forwarded-For` only when the proxy's address is in
`trusted_proxies`; otherwise all users behind the proxy are counted as one
client. 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 the proxy's own
address closes this.
address comes from `X-Forwarded-For` only when the address pixa sees for
requests that come through the proxy is in `trusted_proxies`; otherwise all
users behind the proxy are counted as one client. That address is not always
the proxy's own: a proxy on the Docker host that connects to pixa over
`127.0.0.1` is seen as the gateway of the container's Docker network, such as
`172.17.0.1` on the default bridge, and one that connects through another of the
host's addresses is seen with that address. To be sure, read it as `remoteIP` in
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
@@ -172,19 +178,27 @@ Where:
outside), or `cover` when the URL has no `fit`; a request whose `fit` is
anything else, an empty `fit=` included, is refused with 400
**Example:** resize `https://cdn.example.com/photos/cat.jpg` to 800x600
WebP with expiration 1704067200, default quality and fit:
The URL's `sig` is the HMAC-SHA256 result in base64url (the URL-safe alphabet
of RFC 4648) with the trailing `=` padding kept, 44 characters in all. pixa
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:
`cdn.example.com:/photos/cat.jpg::800:600:webp:1704067200:85:cover`
2. Compute HMAC-SHA256 with your secret key
3. Base64URL-encode the result
2. Compute HMAC-SHA256 of it with the signing key
3. Base64URL-encode the result, keeping the `=` padding:
`-ay7KHpfqmtIGbibDGbUuBDkymi-Ymdn0NkC6j5EJag=`
4. URL:
`/v1/image/cdn.example.com/photos/cat.jpg/800x600.webp?sig=<base64url>&exp=1704067200`
`/v1/image/cdn.example.com/photos/cat.jpg/800x600.webp?sig=-ay7KHpfqmtIGbibDGbUuBDkymi-Ymdn0NkC6j5EJag=&exp=1704067200`
For the same image at quality 40 with fit `contain`, the input ends in
`:40:contain` and the URL is
`/v1/image/cdn.example.com/photos/cat.jpg/800x600.webp?sig=<base64url>&exp=1704067200&q=40&fit=contain`.
`:40:contain`, the signature is `5IwXUx6vf7yefhaUvFzgXZvG2o0Df4RJxPTK3pKq5VU=`,
and the URL is
`/v1/image/cdn.example.com/photos/cat.jpg/800x600.webp?sig=5IwXUx6vf7yefhaUvFzgXZvG2o0Df4RJxPTK3pKq5VU=&exp=1704067200&q=40&fit=contain`.
**Allowlist patterns:**
@@ -248,8 +262,15 @@ Key settings in more detail:
`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, and an explicit list replaces the default. An invalid CIDR aborts
startup. Set this to your proxy's address range if it is not already
covered by the defaults
startup. Set this to the address pixa sees for requests that come through
your proxy, such as `172.17.0.1/32`, when the defaults do not cover it, or
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_max_response_size` — max origin response size
- `downstream_timeout` — client response timeout
+13 -4
View File
@@ -34,10 +34,19 @@ P2: security: referer blacklist
limits the images decoded and encoded at once, and `upstream_connections`
(default 64) the connections to all upstream hosts together, on top of
`upstream_connections_per_host`; a fetch holds its connection until its image
has been processed; a request that finds either limit reached waits up to 10
seconds for a free one, then gets 503 `server busy, try again later`; libvips
runs one worker thread per image with its operation cache off; documented in
`README.md` and `config.example.yml`.
has been processed, and a request whose source is cached reads it only once it
has a processing slot; a request that finds either limit reached waits up to
10 seconds for a free one, then gets 503 `server busy, try again later`;
libvips runs one worker thread per image with its operation cache off;
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
`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`
+7 -1
View File
@@ -50,7 +50,13 @@ allowlist_hosts:
# 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
# explicit list replaces the default. An invalid CIDR aborts startup.
# Uncomment to override the defaults with your proxy's address range.
# Uncomment to override the defaults with the address pixa sees for
# 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:
# - 10.0.0.0/8
# - 2001:db8::/32
@@ -1,6 +1,7 @@
package httpfetcher
import (
"context"
"errors"
"net"
"strconv"
@@ -82,6 +83,42 @@ func TestFetchLimitsConnectionsToAllHostsTogether(t *testing.T) {
_ = 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
// taking its connection gives it back: with MaxConnections at 1, the slot
// must be free after the failure and the next fetch must succeed.
+7 -4
View File
@@ -383,13 +383,16 @@ func (c *Cache) GetSourceMetadataID(
return id, nil
}
// GetSourceContent returns a reader for cached source content by its hash.
func (c *Cache) GetSourceContent(contentHash ContentHash) (io.ReadCloser, error) {
// 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, int64, error) {
if c.disabled {
return nil, ErrNotFound
return nil, 0, ErrNotFound
}
return c.srcContent.Load(contentHash)
return c.srcContent.LoadWithSize(contentHash)
}
// CleanExpired removes expired entries from the cache.
@@ -0,0 +1,125 @@
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)
}
}
+40 -29
View File
@@ -241,38 +241,37 @@ func (s *Service) GenerateSignedURL(
baseURL, path, sig, exp, req.Quality, req.FitMode), nil
}
// loadCachedSource attempts to load source content from cache, returning nil
// if the cached data is unavailable or exceeds maxResponseSize.
func (s *Service) loadCachedSource(contentHash ContentHash) []byte {
reader, err := s.cache.GetSourceContent(contentHash)
// loadCachedSource opens source content from cache, without reading it, and
// returns it with its size; nil if the cached data is unavailable, empty or
// exceeds maxResponseSize.
func (s *Service) loadCachedSource(
contentHash ContentHash,
) (io.ReadCloser, int64) {
reader, size, err := s.cache.GetSourceContent(contentHash)
if err != nil {
s.log.Warn("failed to load cached source, fetching", "error", err)
return nil
return nil, 0
}
// Bound the read to maxResponseSize to prevent unbounded memory use
// from unexpectedly large cached files.
limited := io.LimitReader(reader, s.maxResponseSize+1)
data, err := io.ReadAll(limited)
if size > s.maxResponseSize {
_ = 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",
"hash", contentHash,
"max_bytes", s.maxResponseSize,
)
return nil
return nil, 0
}
return data
if size == 0 {
_ = reader.Close()
return nil, 0
}
return reader, size
}
// processFromSourceOrFetch processes an image, using cached source content
@@ -289,22 +288,27 @@ func (s *Service) processFromSourceOrFetch(
s.log.Warn("source lookup failed", "error", err)
}
var sourceData []byte
var (
source io.ReadCloser
sourceSize int64
)
if contentHash != "" {
s.log.Debug("using cached source", "hash", contentHash)
sourceData = s.loadCachedSource(contentHash)
source, sourceSize = s.loadCachedSource(contentHash)
}
// Fetch from upstream if we don't have source data or it's empty
if len(sourceData) == 0 {
if source == nil {
return s.fetchAndProcess(ctx, req, cacheKey)
}
// Process using cached source; nothing was fetched from upstream
resp, err := s.processAndStore(
ctx, req, cacheKey, sourceData, int64(len(sourceData)),
)
defer func() { _ = source.Close() }()
// Process using cached source; nothing was fetched from upstream. The
// 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
}
@@ -338,6 +342,10 @@ func (s *Service) fetchAndProcess(
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() }()
// Read and validate the source content
@@ -383,17 +391,20 @@ func (s *Service) fetchAndProcess(
// Continue even if caching fails
}
resp, err := s.processAndStore(ctx, req, cacheKey, sourceData, fetchBytes)
resp, err := s.processAndStore(
ctx, req, cacheKey, bytes.NewReader(sourceData), fetchBytes,
)
return resp, fetchBytes, err
}
// processAndStore processes an image and stores the result.
// processAndStore processes the image read from source and stores the
// result.
func (s *Service) processAndStore(
ctx context.Context,
req *ImageRequest,
cacheKey VariantKey,
sourceData []byte,
source io.Reader,
fetchBytes int64,
) (*ImageResponse, error) {
// Process the image
@@ -406,7 +417,7 @@ func (s *Service) processAndStore(
FitMode: imageprocessor.FitMode(req.FitMode),
}
processResult, err := s.processor.Process(ctx, bytes.NewReader(sourceData), processReq)
processResult, err := s.processor.Process(ctx, source, processReq)
if err != nil {
return nil, fmt.Errorf("image processing failed: %w", err)
}