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
This commit is contained in:
2026-09-29 03:31:44 +00:00
parent 23ff05e293
commit 6d6c76937b
10 changed files with 262 additions and 49 deletions
+78 -27
View File
@@ -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
}