diff --git a/README.md b/README.md index 7a3a013..c05372c 100644 --- a/README.md +++ b/README.md @@ -235,6 +235,8 @@ variables set by the file's `env:` section are checked the same way. | `PIXA_TRUSTED_PROXIES` | `trusted_proxies` | CIDR ranges of proxies whose `X-Forwarded-For` is believed; default RFC 1918 | | `PIXA_ALLOW_HTTP` | `allow_http` | Allow plain-HTTP upstreams, for testing only; default `false` | | `PIXA_UPSTREAM_CONNECTIONS_PER_HOST` | `upstream_connections_per_host` | Concurrent connections per upstream host; default `20` | +| `PIXA_UPSTREAM_CONNECTIONS` | `upstream_connections` | Concurrent connections to all upstream hosts together; default `64` | +| `PIXA_MAX_CONCURRENT_PROCESSING` | `max_concurrent_processing` | Images processed at once; default the number of CPUs | | `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_DOWNSTREAM_TIMEOUT` | `downstream_timeout` | Time allowed for answering one client request; default `60s` | @@ -292,6 +294,16 @@ Key settings in more detail: - `cache_max_bytes` — disk cache size limit in bytes; `0` disables the disk cache entirely; omitted defaults to 75% of the free space on the filesystem containing `/cache/` (minimum 500 MiB) +- `upstream_connections` — the most connections to upstream hosts at once, all + hosts together, on top of `upstream_connections_per_host`; default `64`. A + fetch holds its connection until its image has been processed. A fetch that + finds all of them in use waits up to 10 seconds for one to free up; if none + does, the request is answered 503 with the error + `server busy, try again later` +- `max_concurrent_processing` — the most images decoded and encoded at once; + default the number of CPUs pixa can use (`GOMAXPROCS`), which follows a + container's CPU limit. A request that finds all of them in use waits up to 10 + seconds for one to free up; if none does, it is answered 503 the same way See `config.example.yml` for all options with defaults. diff --git a/TODO.md b/TODO.md index e6ad439..cb9f1f2 100644 --- a/TODO.md +++ b/TODO.md @@ -25,11 +25,19 @@ The disk cache is now size-bounded with LRU eviction # Next Step -P1: rate limit global concurrent upstream fetches to prevent resource -exhaustion +P2: security: referer blacklist # Completed Steps +- 2026-09-29 bound concurrent image processing and upstream fetches (closes + #64): `max_concurrent_processing` (default the number of CPUs pixa can use) + 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`. - 2026-09-29 migrations at the path `REPO_POLICIES.md` sets (closes #96): the migration files moved, contents unchanged, from `internal/database/schema/` to `internal/db/migrations/` as `000_migration.sql` and `001_schema.sql`; the @@ -308,7 +316,6 @@ exhaustion # Future Steps - P2: security - - referer blacklist - per-IP rate limiting on the image routes - per-origin rate limiting - P2: HTTP response handling diff --git a/config.example.yml b/config.example.yml index 27c2422..36cf388 100644 --- a/config.example.yml +++ b/config.example.yml @@ -71,6 +71,18 @@ allow_http: false # Maximum concurrent connections per upstream host (default: 20) upstream_connections_per_host: 20 +# Maximum concurrent connections to all upstream hosts together, on top of +# the per-host limit (default: 64). A fetch holds its connection until its +# image has been processed. A fetch that finds none free waits up to 10 +# seconds for one, and if none frees up the request is answered 503. +upstream_connections: 64 + +# Maximum number of images decoded and encoded at once (default: the +# number of CPUs pixa can use, which follows a container's CPU limit). A +# request that finds none free waits up to 10 seconds for one, and if none +# frees up it is answered 503. +# max_concurrent_processing: 4 + # Time allowed for one fetch from an upstream host (default: 30s) upstream_fetch_timeout: 30s diff --git a/internal/config/config.go b/internal/config/config.go index 5c9504f..d815802 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -10,6 +10,7 @@ import ( "net/url" "os" "path/filepath" + "runtime" "sort" "strconv" "strings" @@ -25,6 +26,7 @@ const ( DefaultPort = 8080 DefaultStateDir = "/var/lib/pixa" DefaultUpstreamConnectionsPerHost = 20 + DefaultUpstreamConnections = 64 DefaultAccessControlAllowOrigin = "*" DefaultUpstreamFetchTimeout = 30 * time.Second DefaultUpstreamMaxResponseSize = 50 << 20 // 50 MiB @@ -46,6 +48,8 @@ const ( keyAllowlistHosts = "allowlist_hosts" keyAllowHTTP = "allow_http" keyUpstreamConnectionsPerHost = "upstream_connections_per_host" + keyUpstreamConnections = "upstream_connections" + keyMaxConcurrentProcessing = "max_concurrent_processing" keyCacheMaxBytes = "cache_max_bytes" keyBlockedNetworks = "blocked_networks" keyTrustedProxies = "trusted_proxies" @@ -79,7 +83,7 @@ var ( errNotAValidURL = errors.New("not a valid URL") errPortOutOfRange = errors.New("outside the valid port range") errSizeOutOfRange = errors.New("outside the accepted range") - errTooFewConnections = errors.New("must be at least 1") + errMustBeAtLeastOne = errors.New("must be at least 1") errValueTooShort = errors.New("value too short") errPlaceholderKey = errors.New( "is the placeholder from config.example.yml; " + @@ -126,6 +130,12 @@ type Config struct { AllowHTTP bool // Allow non-TLS upstream (testing only) UpstreamConnectionsPerHost int // Max concurrent connections per upstream host + // UpstreamConnections is the most concurrent connections to all + // upstream hosts together, on top of the per-host limit. + // MaxConcurrentProcessing is the most images processed at once. + UpstreamConnections int + MaxConcurrentProcessing int + // UpstreamFetchTimeout is the time allowed for one fetch from an // upstream host. UpstreamMaxResponseSize is the largest upstream // response accepted, in bytes, and also the image processor's input @@ -271,6 +281,12 @@ func newFromSmartConfig(sc *smartconfig.Config) (*Config, error) { AllowHTTP: loader.boolVal(keyAllowHTTP, false), UpstreamConnectionsPerHost: loader.intVal( keyUpstreamConnectionsPerHost, DefaultUpstreamConnectionsPerHost), + UpstreamConnections: loader.intVal( + keyUpstreamConnections, DefaultUpstreamConnections), + // Decoding and encoding are CPU-bound, so the default is one image + // per CPU Go uses, which follows a container's CPU limit. + MaxConcurrentProcessing: loader.intVal( + keyMaxConcurrentProcessing, runtime.GOMAXPROCS(0)), UpstreamFetchTimeout: loader.durationVal( keyUpstreamFetchTimeout, DefaultUpstreamFetchTimeout), UpstreamMaxResponseSize: loader.int64Val( @@ -392,7 +408,8 @@ func isKnownConfigKey(key string) bool { switch key { case keyDebug, keyMaintenanceMode, keyPort, keyStateDir, keySentryDSN, keyDBURL, keyMetrics, keySigningKey, keyAllowlistHosts, keyAllowHTTP, - keyUpstreamConnectionsPerHost, keyCacheMaxBytes, keyBlockedNetworks, + keyUpstreamConnectionsPerHost, keyUpstreamConnections, + keyMaxConcurrentProcessing, keyCacheMaxBytes, keyBlockedNetworks, keyTrustedProxies, keyAccessControlAllowOrigin, keyUpstreamFetchTimeout, keyUpstreamMaxResponseSize, keyDownstreamTimeout, "env": return true @@ -419,6 +436,8 @@ func envVarNames() map[string]string { keyAllowlistHosts: "PIXA_ALLOWLIST_HOSTS", keyAllowHTTP: "PIXA_ALLOW_HTTP", keyUpstreamConnectionsPerHost: "PIXA_UPSTREAM_CONNECTIONS_PER_HOST", + keyUpstreamConnections: "PIXA_UPSTREAM_CONNECTIONS", + keyMaxConcurrentProcessing: "PIXA_MAX_CONCURRENT_PROCESSING", keyCacheMaxBytes: "PIXA_CACHE_MAX_BYTES", keyBlockedNetworks: "PIXA_BLOCKED_NETWORKS", keyTrustedProxies: "PIXA_TRUSTED_PROXIES", @@ -562,10 +581,9 @@ func (c *Config) validate() error { settingName(keyPort), c.Port, errPortOutOfRange, maxPort) } - if c.UpstreamConnectionsPerHost < 1 { - return fmt.Errorf("%s: value %d %w", - settingName(keyUpstreamConnectionsPerHost), - c.UpstreamConnectionsPerHost, errTooFewConnections) + err = c.validateConcurrencyLimits() + if err != nil { + return err } if c.StateDir == "" { @@ -684,6 +702,30 @@ func (c *Config) validateAccessControlAllowOrigin() error { return nil } +// validateConcurrencyLimits checks that the two upstream connection limits +// and the image processing limit are at least 1. +func (c *Config) validateConcurrencyLimits() error { + if c.UpstreamConnectionsPerHost < 1 { + return fmt.Errorf("%s: value %d %w", + settingName(keyUpstreamConnectionsPerHost), + c.UpstreamConnectionsPerHost, errMustBeAtLeastOne) + } + + if c.UpstreamConnections < 1 { + return fmt.Errorf("%s: value %d %w", + settingName(keyUpstreamConnections), + c.UpstreamConnections, errMustBeAtLeastOne) + } + + if c.MaxConcurrentProcessing < 1 { + return fmt.Errorf("%s: value %d %w", + settingName(keyMaxConcurrentProcessing), + c.MaxConcurrentProcessing, errMustBeAtLeastOne) + } + + return nil +} + // validateAllowlistHost checks that an allowlist_hosts entry is a bare // hostname, optionally with a leading dot for suffix matching. URLs, // paths, and whitespace indicate a misconfigured entry. An entry with diff --git a/internal/handlers/handlers.go b/internal/handlers/handlers.go index 3da580d..b350595 100644 --- a/internal/handlers/handlers.go +++ b/internal/handlers/handlers.go @@ -113,15 +113,17 @@ func (s *Handlers) initImageService() error { fetcherCfg.MaxConnectionsPerHost = s.config.UpstreamConnectionsPerHost } + fetcherCfg.MaxConnections = s.config.UpstreamConnections fetcherCfg.BlockedNetworks = s.config.BlockedNetworks // Create the service svc, err := imgcache.NewService(&imgcache.ServiceConfig{ - Cache: cache, - FetcherConfig: fetcherCfg, - SigningKey: s.config.SigningKey, - Allowlist: s.config.AllowlistHosts, - Logger: s.log, + Cache: cache, + FetcherConfig: fetcherCfg, + SigningKey: s.config.SigningKey, + Allowlist: s.config.AllowlistHosts, + MaxConcurrentProcessing: s.config.MaxConcurrentProcessing, + Logger: s.log, }) if err != nil { return err diff --git a/internal/handlers/image.go b/internal/handlers/image.go index 561ca1d..3c94bc8 100644 --- a/internal/handlers/image.go +++ b/internal/handlers/image.go @@ -12,6 +12,7 @@ import ( "github.com/go-chi/chi/v5" "sneak.berlin/go/pixa/internal/encurl" "sneak.berlin/go/pixa/internal/httpfetcher" + "sneak.berlin/go/pixa/internal/imageprocessor" "sneak.berlin/go/pixa/internal/imgcache" ) @@ -217,6 +218,14 @@ func (s *Handlers) respondImageError( return } + if errors.Is(err, httpfetcher.ErrTooManyConnections) || + errors.Is(err, imageprocessor.ErrTooManyImages) { + s.respondError(w, "server busy, try again later", + http.StatusServiceUnavailable) + + return + } + s.respondError(w, "internal error", http.StatusInternalServerError) } diff --git a/internal/handlers/imageenc.go b/internal/handlers/imageenc.go index 5effd8b..56af31b 100644 --- a/internal/handlers/imageenc.go +++ b/internal/handlers/imageenc.go @@ -12,6 +12,7 @@ import ( "sneak.berlin/go/pixa/internal/encurl" "sneak.berlin/go/pixa/internal/httpfetcher" + "sneak.berlin/go/pixa/internal/imageprocessor" "sneak.berlin/go/pixa/internal/imgcache" ) @@ -124,6 +125,10 @@ func (s *Handlers) handleImageError(w http.ResponseWriter, err error) { s.respondError(w, "upstream error", http.StatusBadGateway) case errors.Is(err, httpfetcher.ErrUpstreamTimeout): s.respondError(w, "upstream timeout", http.StatusGatewayTimeout) + case errors.Is(err, httpfetcher.ErrTooManyConnections), + errors.Is(err, imageprocessor.ErrTooManyImages): + s.respondError(w, "server busy, try again later", + http.StatusServiceUnavailable) default: s.log.Error("image request failed", "error", err) s.respondError(w, "internal error", http.StatusInternalServerError) diff --git a/internal/httpfetcher/httpfetcher.go b/internal/httpfetcher/httpfetcher.go index 64e27eb..d65f2da 100644 --- a/internal/httpfetcher/httpfetcher.go +++ b/internal/httpfetcher/httpfetcher.go @@ -1,5 +1,6 @@ // Package httpfetcher fetches content from upstream HTTP origins with SSRF -// protection, per-host connection limits, and content-type validation. +// protection, connection limits per host and for all hosts together, and +// content-type validation. package httpfetcher import ( @@ -28,8 +29,13 @@ const ( DefaultIdleConnTimeout = 90 * time.Second DefaultMaxRedirects = 10 DefaultMaxConnectionsPerHost = 20 + DefaultMaxConnections = 64 ) +// ConnectionWaitTimeout is how long Fetch waits for a free connection when +// MaxConnections fetches are already in progress. +const ConnectionWaitTimeout = 10 * time.Second + // MIME content types. const ( contentTypeJPEG = "image/jpeg" @@ -70,6 +76,7 @@ var ( ErrInvalidContentType = errors.New("invalid or unsupported content type") ErrUpstreamError = errors.New("upstream server error") ErrUpstreamTimeout = errors.New("upstream request timeout") + ErrTooManyConnections = errors.New("too many concurrent upstream connections") ) // Internal fetcher errors. @@ -122,6 +129,9 @@ type Config struct { AllowHTTP bool // MaxConnectionsPerHost limits concurrent connections to each upstream host. MaxConnectionsPerHost int + // MaxConnections limits concurrent connections to all upstream hosts + // together. + MaxConnections int // BlockedNetworks are operator-supplied CIDR ranges refused by the // dialer, in addition to the always-enforced built-in ranges. BlockedNetworks []netip.Prefix @@ -143,15 +153,22 @@ func DefaultConfig() *Config { }, AllowHTTP: false, MaxConnectionsPerHost: DefaultMaxConnectionsPerHost, + MaxConnections: DefaultMaxConnections, } } -// HTTPFetcher implements Fetcher with SSRF protection and per-host connection limits. +// HTTPFetcher implements Fetcher with SSRF protection and connection limits +// per host and for all hosts together. type HTTPFetcher struct { client *http.Client config *Config hostSems map[string]chan struct{} // per-host semaphores hostSemMu sync.Mutex // protects hostSems map + // allHostsSemaphore has one slot per connection allowed to all hosts + // together (config.MaxConnections). + allHostsSemaphore chan struct{} + // connectionWaitTimeout is ConnectionWaitTimeout; tests shorten it. + connectionWaitTimeout time.Duration } // New creates a new HTTPFetcher with SSRF protection. @@ -192,13 +209,18 @@ func New(config *Config) *HTTPFetcher { } return &HTTPFetcher{ - client: client, - config: config, - hostSems: make(map[string]chan struct{}), + client: client, + config: config, + hostSems: make(map[string]chan struct{}), + allHostsSemaphore: make(chan struct{}, config.MaxConnections), + connectionWaitTimeout: ConnectionWaitTimeout, } } -// Fetch retrieves content from the given URL with SSRF protection. +// Fetch retrieves content from the given URL with SSRF protection. When +// MaxConnections fetches are already in progress, it waits up to +// ConnectionWaitTimeout for one to finish, then fails with +// ErrTooManyConnections. func (f *HTTPFetcher) Fetch(ctx context.Context, url string) (*FetchResult, error) { // Validate URL before making request err := validateURL(ctx, url, f.config.AllowHTTP) @@ -206,24 +228,17 @@ func (f *HTTPFetcher) Fetch(ctx context.Context, url string) (*FetchResult, erro return nil, err } - // Extract host for rate limiting - host := extractHost(url) - - // Acquire semaphore slot for this host - sem := f.getHostSemaphore(host) - select { - case sem <- struct{}{}: - // Acquired slot - case <-ctx.Done(): - return nil, ctx.Err() + release, err := f.acquireConnection(ctx, extractHost(url)) + if err != nil { + return nil, err } - // If we fail before returning a result, release the slot + // If we fail before returning a result, release the connection success := false defer func() { if !success { - <-sem + release() } }() @@ -267,17 +282,52 @@ func (f *HTTPFetcher) Fetch(ctx context.Context, url string) (*FetchResult, erro return nil, fmt.Errorf("upstream request failed: %w", err) } - result, err := f.buildResult(resp, remoteAddr, fetchDuration, sem) + result, err := f.buildResult(resp, remoteAddr, fetchDuration, release) if err != nil { return nil, err } - // Mark success so defer doesn't release the semaphore + // Mark success so defer doesn't release the connection; closing the + // result's Content does success = true return result, nil } +// acquireConnection takes a slot for host, then one of the slots shared by +// all hosts, and returns the func that gives both back. The host's slot +// comes first, so fetches queued for one busy host hold no shared slot. +// Only the wait for a shared slot is bounded: after connectionWaitTimeout +// it fails with ErrTooManyConnections. +func (f *HTTPFetcher) acquireConnection( + ctx context.Context, host string, +) (func(), error) { + hostSem := f.getHostSemaphore(host) + + select { + case hostSem <- struct{}{}: + case <-ctx.Done(): + return nil, ctx.Err() + } + + select { + case f.allHostsSemaphore <- struct{}{}: + case <-time.After(f.connectionWaitTimeout): + <-hostSem + + return nil, ErrTooManyConnections + case <-ctx.Done(): + <-hostSem + + return nil, ctx.Err() + } + + return func() { + <-hostSem + <-f.allHostsSemaphore + }, nil +} + // getHostSemaphore returns the semaphore for a host, creating it if necessary. func (f *HTTPFetcher) getHostSemaphore(host string) chan struct{} { f.hostSemMu.Lock() @@ -293,12 +343,12 @@ func (f *HTTPFetcher) getHostSemaphore(host string) chan struct{} { } // buildResult validates the upstream response and assembles a FetchResult -// whose Content releases the host semaphore slot when closed. +// whose Content calls release when closed. func (f *HTTPFetcher) buildResult( resp *http.Response, remoteAddr string, fetchDuration time.Duration, - sem chan struct{}, + release func(), ) (*FetchResult, error) { // Extract HTTP version (strip "HTTP/" prefix) httpVersion := strings.TrimPrefix(resp.Proto, "HTTP/") @@ -333,7 +383,7 @@ func (f *HTTPFetcher) buildResult( } return &FetchResult{ - Content: &semaphoreReleasingReadCloser{limitedBody, resp.Body, sem}, + Content: &semaphoreReleasingReadCloser{limitedBody, resp.Body, release}, ContentLength: resp.ContentLength, ContentType: contentType, Headers: resp.Header, @@ -574,17 +624,18 @@ func (r *limitedReader) Read(p []byte) (int, error) { return n, err } -// semaphoreReleasingReadCloser releases a semaphore slot when closed. +// semaphoreReleasingReadCloser releases the fetch's connection slots when +// closed. type semaphoreReleasingReadCloser struct { *limitedReader - closer io.Closer - sem chan struct{} + closer io.Closer + release func() } func (r *semaphoreReleasingReadCloser) Close() error { err := r.closer.Close() - <-r.sem // Release semaphore slot + r.release() return err } diff --git a/internal/imageprocessor/imageprocessor.go b/internal/imageprocessor/imageprocessor.go index dd77ac1..6a9a92d 100644 --- a/internal/imageprocessor/imageprocessor.go +++ b/internal/imageprocessor/imageprocessor.go @@ -7,7 +7,9 @@ import ( "errors" "fmt" "io" + "runtime" "sync" + "time" "github.com/davidbyttow/govips/v2/vips" ) @@ -17,11 +19,21 @@ import ( //nolint:gochecknoglobals // package-level sync.Once for one-time vips init var vipsOnce sync.Once -// initVips initializes libvips with quiet logging. +// initVips initializes libvips with quiet logging, one worker thread per +// image and no operation cache. Process already works on one image per CPU +// by default, so more threads per image would only compete for the CPUs. +// Each request decodes different source bytes, so the operation cache +// would rarely be hit and would hold memory outside MaxConcurrentProcessing; +// repeated requests are served from pixa's disk cache instead. func initVips() { vipsOnce.Do(func() { vips.LoggingSettings(nil, vips.LogLevelError) - vips.Startup(nil) + vips.Startup(&vips.Config{ + ConcurrencyLevel: 1, + MaxCacheSize: 0, + MaxCacheMem: 0, + MaxCacheFiles: 0, + }) }) } @@ -106,9 +118,23 @@ var ErrInputDataTooLarge = errors.New("input data exceeds maximum allowed size") // not supported. var ErrUnsupportedOutputFormat = errors.New("unsupported output format") +// ErrTooManyImages is returned when MaxConcurrentProcessing images are being +// processed and none finishes within ProcessingWaitTimeout. +var ErrTooManyImages = errors.New("too many images being processed at once") + +// ProcessingWaitTimeout is how long Process waits for a free slot when +// MaxConcurrentProcessing images are already being processed. +const ProcessingWaitTimeout = 10 * time.Second + // ImageProcessor implements image transformation using libvips via govips. type ImageProcessor struct { maxInputBytes int64 + // processingSemaphore has one slot per image that may be processed at + // once. Process holds a slot from before it reads its input until it + // returns, so the input, the decoded image and the output all count. + processingSemaphore chan struct{} + // processingWaitTimeout is ProcessingWaitTimeout; tests shorten it. + processingWaitTimeout time.Duration } // Params holds configuration for creating an ImageProcessor. @@ -117,6 +143,9 @@ type Params struct { // MaxInputBytes is the maximum allowed input size in bytes. // If <= 0, DefaultMaxInputBytes is used. MaxInputBytes int64 + // MaxConcurrentProcessing is the most images processed at once. + // If <= 0, the number of CPUs Go uses (runtime.GOMAXPROCS(0)) is used. + MaxConcurrentProcessing int } // New creates a new image processor with the given parameters. @@ -129,17 +158,34 @@ func New(params Params) *ImageProcessor { maxInputBytes = DefaultMaxInputBytes } + maxConcurrentProcessing := params.MaxConcurrentProcessing + if maxConcurrentProcessing <= 0 { + maxConcurrentProcessing = runtime.GOMAXPROCS(0) + } + return &ImageProcessor{ - maxInputBytes: maxInputBytes, + maxInputBytes: maxInputBytes, + processingSemaphore: make(chan struct{}, maxConcurrentProcessing), + processingWaitTimeout: ProcessingWaitTimeout, } } -// Process transforms an image according to the request. +// Process transforms an image according to the request. When +// MaxConcurrentProcessing images are already being processed, it waits up +// to ProcessingWaitTimeout for one to finish, then fails with +// ErrTooManyImages. func (p *ImageProcessor) Process( - _ context.Context, + ctx context.Context, input io.Reader, req *Request, ) (*Result, error) { + release, err := p.acquireSlot(ctx) + if err != nil { + return nil, err + } + + defer release() + // Read input with a size limit to prevent unbounded memory consumption. // We read at most maxInputBytes+1 so we can detect if the input exceeds // the limit without consuming additional memory. @@ -285,6 +331,29 @@ func FormatToMIME(format Format) string { } } +// acquireSlot takes a slot in processingSemaphore, waiting at most +// processingWaitTimeout for one to free up, and returns the func that gives +// it back. A free slot is taken even when ctx has ended; only the wait for +// one stops when ctx ends, as the rest of Process does not check ctx. +func (p *ImageProcessor) acquireSlot(ctx context.Context) (func(), error) { + release := func() { <-p.processingSemaphore } + + select { + case p.processingSemaphore <- struct{}{}: + return release, nil + default: + } + + select { + case p.processingSemaphore <- struct{}{}: + return release, nil + case <-time.After(p.processingWaitTimeout): + return nil, ErrTooManyImages + case <-ctx.Done(): + return nil, ctx.Err() + } +} + // detectFormat returns the format string from a vips image. func (p *ImageProcessor) detectFormat(img *vips.ImageRef) string { format := img.Format() diff --git a/internal/imgcache/service.go b/internal/imgcache/service.go index d85dee8..0792a28 100644 --- a/internal/imgcache/service.go +++ b/internal/imgcache/service.go @@ -43,6 +43,9 @@ type ServiceConfig struct { SigningKey string // Allowlist is the list of hosts that don't require signatures Allowlist []string + // MaxConcurrentProcessing is the most images processed at once; zero + // uses the image processor's default, one per CPU + MaxConcurrentProcessing int // Logger for logging Logger *slog.Logger } @@ -91,9 +94,10 @@ func NewService(cfg *ServiceConfig) (*Service, error) { } maxResponseSize := fetcherCfg.MaxResponseSize - processor := imageprocessor.New( - imageprocessor.Params{MaxInputBytes: maxResponseSize}, - ) + processor := imageprocessor.New(imageprocessor.Params{ + MaxInputBytes: maxResponseSize, + MaxConcurrentProcessing: cfg.MaxConcurrentProcessing, + }) return &Service{ cache: cfg.Cache,