diff --git a/README.md b/README.md index ed5b3b0..05442ed 100644 --- a/README.md +++ b/README.md @@ -107,6 +107,13 @@ or proxy cache keeps the image after pixa would refuse the URL. A URL with no expiry gets one year. `immutable` only stops a client revalidating while its copy is fresh. +When several requests for the same image, size, format, quality and fit miss +the cache at once, they share one upstream fetch (or one read of the cached +source) and one transcode: the first request does the work, and the others wait +for its image or its error, holding no upstream connection or processing slot +of their own. A waiting request stops waiting when its own client goes away; +the work goes on for the others. + The login form (`POST /`) is limited to 5 attempts per minute per client address, counting an IPv6 client by its /64; an attempt over the limit is refused with 429 and a `Retry-After` header. Behind a reverse proxy the client diff --git a/TODO.md b/TODO.md index 637e3a5..d0c75a0 100644 --- a/TODO.md +++ b/TODO.md @@ -29,6 +29,17 @@ P2: security: referer blacklist # Completed Steps +- 2026-09-29 share concurrent misses (closes #65): requests that miss the same + variant at once (the same cache key, so quality and fit included) share one + upstream fetch or cached source read and one transcode through + `golang.org/x/sync/singleflight`; the first request's processing runs with a + context that does not end with its own, and the others wait for its image or + error holding no upstream connection or processing slot, and stop waiting + when their own context ends; the request doing the processing waits for it + even then, as before; a request whose context has already ended starts + nothing; each request counts one miss, and the processing counts its fetch + and transcode once; a panic while processing becomes an error for every + waiting request instead of stopping pixad; documented in `README.md`. - 2026-09-29 the container makes `/var/lib/pixa` usable by itself (closes #159): `deploy/docker-entrypoint.sh` creates the directory if it is missing, gives the directory and everything in it to `pixad` when the directory or one diff --git a/go.mod b/go.mod index 4a13f35..ae8006a 100644 --- a/go.mod +++ b/go.mod @@ -20,6 +20,7 @@ require ( github.com/spf13/cobra v1.10.2 go.uber.org/fx v1.24.0 golang.org/x/crypto v0.41.0 + golang.org/x/sync v0.19.0 modernc.org/sqlite v1.42.2 ) @@ -134,7 +135,6 @@ require ( golang.org/x/image v0.34.0 // indirect golang.org/x/net v0.43.0 // indirect golang.org/x/oauth2 v0.30.0 // indirect - golang.org/x/sync v0.19.0 // indirect golang.org/x/sys v0.36.0 // indirect golang.org/x/term v0.34.0 // indirect golang.org/x/text v0.32.0 // indirect diff --git a/internal/imgcache/cache.go b/internal/imgcache/cache.go index a4faab4..a2d2d61 100644 --- a/internal/imgcache/cache.go +++ b/internal/imgcache/cache.go @@ -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) } } diff --git a/internal/imgcache/service.go b/internal/imgcache/service.go index e1152dc..0634aaa 100644 --- a/internal/imgcache/service.go +++ b/internal/imgcache/service.go @@ -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 }