Share one fetch and transcode among concurrent misses (closes #65)
check / check (push) Waiting to run
check / check (push) Waiting to run
Requests that missed the same variant at once each fetched and transcoded it. They now share one call through golang.org/x/sync/singleflight, keyed on the variant cache key. The first request processes the variant with a context that does not end with its own; the others wait for its result, holding no connection or processing slot, and return as soon as their own context ends. The processing request waits even then, as before. Each request counts one miss; the processing counts its fetch and transcode once. A panic while processing becomes an error instead of stopping pixad. Model: opus-5-5
This commit is contained in:
+22
-13
@@ -471,7 +471,8 @@ func (c *Cache) Stats(ctx context.Context) (*CacheStats, error) {
|
||||
return &stats, nil
|
||||
}
|
||||
|
||||
// IncrementStats increments cache statistics.
|
||||
// IncrementStats counts a cache hit or miss, and an upstream fetch that read
|
||||
// fetchBytes bytes, as IncrementUpstreamFetch does.
|
||||
func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64) {
|
||||
var err error
|
||||
|
||||
@@ -495,18 +496,26 @@ func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64)
|
||||
c.log.Warn("failed to count cache hit or miss", "hit", hit, "error", err)
|
||||
}
|
||||
|
||||
if fetchBytes > 0 {
|
||||
_, err = c.db.ExecContext(ctx, `
|
||||
UPDATE cache_stats
|
||||
SET upstream_fetch_count = upstream_fetch_count + 1,
|
||||
upstream_fetch_bytes = upstream_fetch_bytes + ?,
|
||||
last_updated_at = CURRENT_TIMESTAMP
|
||||
WHERE id = 1
|
||||
`, fetchBytes)
|
||||
if err != nil {
|
||||
c.log.Warn("failed to count upstream fetch",
|
||||
"fetch_bytes", fetchBytes, "error", err)
|
||||
}
|
||||
c.IncrementUpstreamFetch(ctx, fetchBytes)
|
||||
}
|
||||
|
||||
// IncrementUpstreamFetch counts one upstream fetch that read fetchBytes bytes.
|
||||
// A fetch that read no bytes is not counted.
|
||||
func (c *Cache) IncrementUpstreamFetch(ctx context.Context, fetchBytes int64) {
|
||||
if fetchBytes <= 0 {
|
||||
return
|
||||
}
|
||||
|
||||
_, err := c.db.ExecContext(ctx, `
|
||||
UPDATE cache_stats
|
||||
SET upstream_fetch_count = upstream_fetch_count + 1,
|
||||
upstream_fetch_bytes = upstream_fetch_bytes + ?,
|
||||
last_updated_at = CURRENT_TIMESTAMP
|
||||
WHERE id = 1
|
||||
`, fetchBytes)
|
||||
if err != nil {
|
||||
c.log.Warn("failed to count upstream fetch",
|
||||
"fetch_bytes", fetchBytes, "error", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+111
-27
@@ -8,9 +8,11 @@ import (
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/url"
|
||||
"runtime/debug"
|
||||
"time"
|
||||
|
||||
"github.com/dustin/go-humanize"
|
||||
"golang.org/x/sync/singleflight"
|
||||
"sneak.berlin/go/pixa/internal/allowlist"
|
||||
"sneak.berlin/go/pixa/internal/httpfetcher"
|
||||
"sneak.berlin/go/pixa/internal/imageprocessor"
|
||||
@@ -29,6 +31,9 @@ type Service struct {
|
||||
log *slog.Logger
|
||||
allowHTTP bool
|
||||
maxResponseSize int64
|
||||
// variantsInProgress lets the requests that miss the same variant at the
|
||||
// same time share one fetch and one transcode.
|
||||
variantsInProgress singleflight.Group
|
||||
}
|
||||
|
||||
// ServiceConfig holds configuration for the image service.
|
||||
@@ -160,14 +165,12 @@ func (s *Service) Get(ctx context.Context, req *ImageRequest) (*ImageResponse, e
|
||||
}
|
||||
}
|
||||
|
||||
// Cache miss - process the cached source or fetch it, then count the
|
||||
// miss with the bytes it fetched from upstream, also when it failed or
|
||||
// the request context has ended meanwhile
|
||||
cacheKey := CacheKey(req)
|
||||
// Cache miss - get the variant, processed once for all the requests that
|
||||
// miss it at the same time, then count this request's miss, also when it
|
||||
// failed or the request context has ended meanwhile
|
||||
response, err := s.processOrWait(ctx, req)
|
||||
|
||||
response, fetchedBytes, err := s.processFromSourceOrFetch(ctx, req, cacheKey)
|
||||
|
||||
s.cache.IncrementStats(context.WithoutCancel(ctx), false, fetchedBytes)
|
||||
s.cache.IncrementStats(context.WithoutCancel(ctx), false, 0)
|
||||
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -241,6 +244,83 @@ func (s *Service) GenerateSignedURL(
|
||||
baseURL, path, sig, exp, req.Quality, req.FitMode), nil
|
||||
}
|
||||
|
||||
// errPanicked is returned when processing a variant panicked.
|
||||
var errPanicked = errors.New("panic while processing image")
|
||||
|
||||
// processOrWait returns the variant req asks for. The first of the requests
|
||||
// that miss a variant at the same time processes it, and singleflight hands
|
||||
// its result, or its error, to the others: they fetch nothing, read no source
|
||||
// and take no upstream connection or processing slot. The processing runs with
|
||||
// a context that does not end with the first request's, so the others are
|
||||
// still served if that client goes away; the fetch timeout and the waits for a
|
||||
// connection and a processing slot still bound it.
|
||||
func (s *Service) processOrWait(
|
||||
ctx context.Context, req *ImageRequest,
|
||||
) (*ImageResponse, error) {
|
||||
// A request that has already ended starts no processing
|
||||
if ctx.Err() != nil {
|
||||
return nil, ctx.Err()
|
||||
}
|
||||
|
||||
cacheKey := CacheKey(req)
|
||||
|
||||
// Closed when this request's own function runs, which singleflight does
|
||||
// only when no other request is processing the variant
|
||||
processing := make(chan struct{})
|
||||
|
||||
results := s.variantsInProgress.DoChan(string(cacheKey),
|
||||
func() (_ any, err error) {
|
||||
close(processing)
|
||||
|
||||
// singleflight would raise a panic again in a goroutine of its
|
||||
// own, where no handler recovers it, and stop pixad
|
||||
defer func() {
|
||||
recovered := recover()
|
||||
if recovered != nil {
|
||||
s.log.Error("panic while processing image",
|
||||
"host", req.SourceHost, "path", req.SourcePath,
|
||||
"panic", recovered, "stack", string(debug.Stack()))
|
||||
|
||||
err = fmt.Errorf("%w: %v", errPanicked, recovered)
|
||||
}
|
||||
}()
|
||||
|
||||
return s.processFromSourceOrFetch(
|
||||
context.WithoutCancel(ctx), req, cacheKey)
|
||||
})
|
||||
|
||||
var result singleflight.Result
|
||||
|
||||
select {
|
||||
case result = <-results:
|
||||
case <-ctx.Done():
|
||||
select {
|
||||
case <-processing:
|
||||
// This request is processing the variant: it waits for the
|
||||
// result, as every request did before misses were shared
|
||||
result = <-results
|
||||
default:
|
||||
// Another request is processing the variant, or this request's
|
||||
// function has not started yet; the processing goes on without it
|
||||
return nil, ctx.Err()
|
||||
}
|
||||
}
|
||||
|
||||
if result.Err != nil {
|
||||
return nil, result.Err
|
||||
}
|
||||
|
||||
variant, _ := result.Val.(*processedVariant)
|
||||
|
||||
return &ImageResponse{
|
||||
Content: io.NopCloser(bytes.NewReader(variant.data)),
|
||||
ContentLength: int64(len(variant.data)),
|
||||
ContentType: variant.contentType,
|
||||
FetchedBytes: variant.fetchedBytes,
|
||||
ETag: formatETag(cacheKey),
|
||||
}, nil
|
||||
}
|
||||
|
||||
// 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.
|
||||
@@ -275,13 +355,12 @@ func (s *Service) loadCachedSource(
|
||||
}
|
||||
|
||||
// processFromSourceOrFetch processes an image, using cached source content
|
||||
// if available. It also returns the number of bytes fetched from upstream,
|
||||
// as fetchAndProcess does, or 0 when the cached source was used.
|
||||
// if available.
|
||||
func (s *Service) processFromSourceOrFetch(
|
||||
ctx context.Context,
|
||||
req *ImageRequest,
|
||||
cacheKey VariantKey,
|
||||
) (*ImageResponse, int64, error) {
|
||||
) (*processedVariant, error) {
|
||||
// Check if we have cached source content
|
||||
contentHash, _, err := s.cache.LookupSource(ctx, req)
|
||||
if err != nil {
|
||||
@@ -308,19 +387,17 @@ func (s *Service) processFromSourceOrFetch(
|
||||
// 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 s.processAndStore(ctx, req, cacheKey, source, sourceSize)
|
||||
}
|
||||
|
||||
// fetchAndProcess fetches from upstream, processes, and caches the result.
|
||||
// It also returns the number of bytes read from upstream, including when
|
||||
// It counts the fetch with the bytes read from upstream, including when
|
||||
// reading the response or a later step fails.
|
||||
func (s *Service) fetchAndProcess(
|
||||
ctx context.Context,
|
||||
req *ImageRequest,
|
||||
cacheKey VariantKey,
|
||||
) (*ImageResponse, int64, error) {
|
||||
) (*processedVariant, error) {
|
||||
// Fetch from upstream
|
||||
sourceURL := req.SourceURL()
|
||||
|
||||
@@ -339,7 +416,7 @@ func (s *Service) fetchAndProcess(
|
||||
}
|
||||
}
|
||||
|
||||
return nil, 0, fmt.Errorf("upstream fetch failed: %w", err)
|
||||
return nil, fmt.Errorf("upstream fetch failed: %w", err)
|
||||
}
|
||||
|
||||
// Closing the body frees the upstream connection. It is closed only
|
||||
@@ -352,8 +429,11 @@ func (s *Service) fetchAndProcess(
|
||||
sourceData, err := io.ReadAll(fetchResult.Content)
|
||||
fetchBytes := int64(len(sourceData))
|
||||
|
||||
// Counted also when the request context has ended meanwhile
|
||||
s.cache.IncrementUpstreamFetch(context.WithoutCancel(ctx), fetchBytes)
|
||||
|
||||
if err != nil {
|
||||
return nil, fetchBytes, fmt.Errorf("failed to read upstream response: %w", err)
|
||||
return nil, fmt.Errorf("failed to read upstream response: %w", err)
|
||||
}
|
||||
|
||||
// Calculate download bitrate
|
||||
@@ -381,7 +461,7 @@ func (s *Service) fetchAndProcess(
|
||||
// Validate magic bytes match content type
|
||||
err = magic.ValidateMagicBytes(sourceData, fetchResult.ContentType)
|
||||
if err != nil {
|
||||
return nil, fetchBytes, fmt.Errorf("content validation failed: %w", err)
|
||||
return nil, fmt.Errorf("content validation failed: %w", err)
|
||||
}
|
||||
|
||||
// Store source content
|
||||
@@ -391,11 +471,17 @@ func (s *Service) fetchAndProcess(
|
||||
// Continue even if caching fails
|
||||
}
|
||||
|
||||
resp, err := s.processAndStore(
|
||||
return s.processAndStore(
|
||||
ctx, req, cacheKey, bytes.NewReader(sourceData), fetchBytes,
|
||||
)
|
||||
}
|
||||
|
||||
return resp, fetchBytes, err
|
||||
// processedVariant is a variant as processAndStore made it. Each request that
|
||||
// shared its processing serves it through a reader of its own.
|
||||
type processedVariant struct {
|
||||
data []byte
|
||||
contentType string
|
||||
fetchedBytes int64
|
||||
}
|
||||
|
||||
// processAndStore processes the image read from source and stores the
|
||||
@@ -406,7 +492,7 @@ func (s *Service) processAndStore(
|
||||
cacheKey VariantKey,
|
||||
source io.Reader,
|
||||
fetchBytes int64,
|
||||
) (*ImageResponse, error) {
|
||||
) (*processedVariant, error) {
|
||||
// Process the image
|
||||
processStart := time.Now()
|
||||
|
||||
@@ -470,12 +556,10 @@ func (s *Service) processAndStore(
|
||||
// Continue even if caching fails
|
||||
}
|
||||
|
||||
return &ImageResponse{
|
||||
Content: io.NopCloser(bytes.NewReader(processedData)),
|
||||
ContentLength: outputSize,
|
||||
ContentType: processResult.ContentType,
|
||||
FetchedBytes: fetchBytes,
|
||||
ETag: formatETag(cacheKey),
|
||||
return &processedVariant{
|
||||
data: processedData,
|
||||
contentType: processResult.ContentType,
|
||||
fetchedBytes: fetchBytes,
|
||||
}, nil
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user