Bound concurrent image processing and upstream fetches (closes #64) #148

Merged
clawbot merged 6 commits from issue-64-concurrency-limits into next 2026-09-29 10:32:18 +02:00
3 changed files with 53 additions and 38 deletions
Showing only changes of commit d800aa62a7 - Show all commits
+5 -4
View File
@@ -34,10 +34,11 @@ 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; a request that finds either limit reached waits up to 10 has been processed, and a request whose source is cached reads it only once it
seconds for a free one, then gets 503 `server busy, try again later`; libvips has a processing slot; a request that finds either limit reached waits up to
runs one worker thread per image with its operation cache off; documented in 10 seconds for a free one, then gets 503 `server busy, try again later`;
`README.md` and `config.example.yml`. libvips runs one worker thread per image with its operation cache off;
documented in `README.md` and `config.example.yml`.
- 2026-09-29 Dockerfiles install through `script/bootstrap` (closes #95): the - 2026-09-29 Dockerfiles install through `script/bootstrap` (closes #95): the
`Dockerfile` lint and build stages and `Dockerfile.lint` copy `script/`, `Dockerfile` lint and build stages and `Dockerfile.lint` copy `script/`,
`go.mod` and `go.sum`, then run `script/bootstrap` in place of their own `go.mod` and `go.sum`, then run `script/bootstrap` in place of their own
+7 -4
View File
@@ -383,13 +383,16 @@ 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,
func (c *Cache) GetSourceContent(contentHash ContentHash) (io.ReadCloser, error) { // and the content's size in bytes.
func (c *Cache) GetSourceContent(
contentHash ContentHash,
) (io.ReadCloser, int64, error) {
if c.disabled { 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. // CleanExpired removes expired entries from the cache.
+40 -29
View File
@@ -241,38 +241,37 @@ func (s *Service) GenerateSignedURL(
baseURL, path, sig, exp, req.Quality, req.FitMode), nil baseURL, path, sig, exp, req.Quality, req.FitMode), nil
} }
// loadCachedSource attempts to load source content from cache, returning nil // loadCachedSource opens source content from cache, without reading it, and
// if the cached data is unavailable or exceeds maxResponseSize. // returns it with its size; nil if the cached data is unavailable, empty or
func (s *Service) loadCachedSource(contentHash ContentHash) []byte { // exceeds maxResponseSize.
reader, err := s.cache.GetSourceContent(contentHash) func (s *Service) loadCachedSource(
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 return nil, 0
} }
// Bound the read to maxResponseSize to prevent unbounded memory use if size > s.maxResponseSize {
// from unexpectedly large cached files.
limited := io.LimitReader(reader, s.maxResponseSize+1)
data, err := io.ReadAll(limited)
_ = reader.Close() _ = 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 return nil, 0
} }
return data if size == 0 {
_ = reader.Close()
return nil, 0
}
return reader, size
} }
// processFromSourceOrFetch processes an image, using cached source content // processFromSourceOrFetch processes an image, using cached source content
@@ -289,22 +288,27 @@ func (s *Service) processFromSourceOrFetch(
s.log.Warn("source lookup failed", "error", err) s.log.Warn("source lookup failed", "error", err)
} }
var sourceData []byte var (
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)
sourceData = s.loadCachedSource(contentHash) source, sourceSize = 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 len(sourceData) == 0 { if source == nil {
return s.fetchAndProcess(ctx, req, cacheKey) return s.fetchAndProcess(ctx, req, cacheKey)
} }
// Process using cached source; nothing was fetched from upstream defer func() { _ = source.Close() }()
resp, err := s.processAndStore(
ctx, req, cacheKey, sourceData, int64(len(sourceData)), // 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 return resp, 0, err
} }
@@ -338,6 +342,10 @@ 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
@@ -383,17 +391,20 @@ func (s *Service) fetchAndProcess(
// Continue even if caching fails // 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 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( func (s *Service) processAndStore(
ctx context.Context, ctx context.Context,
req *ImageRequest, req *ImageRequest,
cacheKey VariantKey, cacheKey VariantKey,
sourceData []byte, source io.Reader,
fetchBytes int64, fetchBytes int64,
) (*ImageResponse, error) { ) (*ImageResponse, error) {
// Process the image // Process the image
@@ -406,7 +417,7 @@ func (s *Service) processAndStore(
FitMode: imageprocessor.FitMode(req.FitMode), 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 { if err != nil {
return nil, fmt.Errorf("image processing failed: %w", err) return nil, fmt.Errorf("image processing failed: %w", err)
} }