diff --git a/TODO.md b/TODO.md index c30341d..8a84f56 100644 --- a/TODO.md +++ b/TODO.md @@ -34,10 +34,11 @@ 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 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 diff --git a/internal/imgcache/cache.go b/internal/imgcache/cache.go index e15b90a..ddac1b6 100644 --- a/internal/imgcache/cache.go +++ b/internal/imgcache/cache.go @@ -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. diff --git a/internal/imgcache/service.go b/internal/imgcache/service.go index 0792a28..e1152dc 100644 --- a/internal/imgcache/service.go +++ b/internal/imgcache/service.go @@ -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) - _ = reader.Close() + 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) }