Read a cached source only once a processing slot is taken (closes #64)
check / check (push) Successful in 3m10s
check / check (push) Successful in 3m10s
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
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user