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/concurrent_misses_internal_test.go b/internal/imgcache/concurrent_misses_internal_test.go new file mode 100644 index 0000000..53b2db5 --- /dev/null +++ b/internal/imgcache/concurrent_misses_internal_test.go @@ -0,0 +1,444 @@ +package imgcache + +import ( + "bytes" + "context" + "errors" + "image/jpeg" + "io" + "io/fs" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + "sneak.berlin/go/pixa/internal/httpfetcher" + "sneak.berlin/go/pixa/internal/magic" +) + +// arrivalWait is how long a test gives requests it has started to reach the +// point where they wait for a held fetch. +const arrivalWait = 100 * time.Millisecond + +// heldFetcher counts the fetches it is asked for and holds each one until +// releaseFetches is called, so that a test can have requests arrive while a +// fetch is in progress. started receives once for every fetch. +type heldFetcher struct { + upstream httpfetcher.Fetcher + fetches atomic.Int32 + started chan struct{} + release chan struct{} + releaseOnce sync.Once +} + +func (f *heldFetcher) Fetch( + ctx context.Context, url string, +) (*httpfetcher.FetchResult, error) { + f.fetches.Add(1) + + f.started <- struct{}{} + + select { + case <-f.release: + case <-ctx.Done(): + return nil, ctx.Err() + } + + return f.upstream.Fetch(ctx, url) +} + +// releaseFetches lets every held fetch, and every later one, go on. +func (f *heldFetcher) releaseFetches() { + f.releaseOnce.Do(func() { close(f.release) }) +} + +// setupHeldFetchService returns a test service whose fetches go through a +// heldFetcher. Its database is limited to one connection: each connection to +// an in-memory SQLite database opens a new, empty one, so requests running at +// once must share the connection that holds the schema. +func setupHeldFetchService(t *testing.T) (*Service, *TestFixtures, *heldFetcher) { + t.Helper() + + svc, fixtures := SetupTestService(t) + svc.cache.db.SetMaxOpenConns(1) + + fetcher := &heldFetcher{ + upstream: svc.fetcher, + started: make(chan struct{}, 100), + release: make(chan struct{}), + } + svc.fetcher = fetcher + + t.Cleanup(fetcher.releaseFetches) + + return svc, fixtures, fetcher +} + +// photoVariant asks for the test photo, 100x100, at 50x25 as a JPEG of the +// given quality and fit mode. Each call returns a new request, as Get writes +// to the request it is given. +func photoVariant(fixtures *TestFixtures, quality int, fit FitMode) *ImageRequest { + return &ImageRequest{ + SourceHost: fixtures.GoodHost, + SourcePath: testPathPhoto, + Size: Size{Width: 50, Height: 25}, + Format: FormatJPEG, + Quality: quality, + FitMode: fit, + } +} + +// getResult is what one Get call returned, with the image read out. +type getResult struct { + image []byte + err error +} + +// startGet calls Get in a goroutine of its own and delivers what it returned +// on the channel. +func startGet( + ctx context.Context, svc *Service, req *ImageRequest, +) <-chan getResult { + results := make(chan getResult, 1) + + go func() { + resp, err := svc.Get(ctx, req) + if err != nil { + results <- getResult{err: err} + + return + } + + defer func() { _ = resp.Content.Close() }() + + image, err := io.ReadAll(resp.Content) + results <- getResult{image: image, err: err} + }() + + return results +} + +// jpegSize returns the width and height of the JPEG image in data, or 0 and 0 +// if data is not one. +func jpegSize(data []byte) (int, int) { + config, err := jpeg.DecodeConfig(bytes.NewReader(data)) + if err != nil { + return 0, 0 + } + + return config.Width, config.Height +} + +// TestService_Get_ConcurrentMissesShareOneFetch starts several requests for +// one uncached variant while the first one's fetch is held. Between them they +// must fetch the source once and transcode it once, every one must be answered +// with the same 50x25 JPEG, and each must count one miss. +func TestService_Get_ConcurrentMissesShareOneFetch(t *testing.T) { + t.Parallel() + + svc, fixtures, fetcher := setupHeldFetchService(t) + + const requests = 8 + + pending := make([]<-chan getResult, 0, requests) + for range requests { + pending = append(pending, + startGet(t.Context(), svc, photoVariant(fixtures, 85, FitCover))) + } + + <-fetcher.started + time.Sleep(arrivalWait) + fetcher.releaseFetches() + + var first []byte + + for i, results := range pending { + got := <-results + if got.err != nil { + t.Fatalf("request %d: Get() error = %v", i, got.err) + } + + if width, height := jpegSize(got.image); width != 50 || height != 25 { + t.Errorf("request %d: image is %dx%d, want a 50x25 JPEG", i, width, height) + } + + if first == nil { + first = got.image + } else if !bytes.Equal(got.image, first) { + t.Errorf("request %d: image differs from request 0's", i) + } + } + + if fetches := fetcher.fetches.Load(); fetches != 1 { + t.Errorf("%d requests made %d upstream fetches, want 1", requests, fetches) + } + + // NewTestFS builds the same files the test service's fetcher serves. + testFS, _ := NewTestFS(t) + + photo, err := fs.ReadFile(testFS, fixtures.GoodHostJPEG) + if err != nil { + t.Fatal(err) + } + + want := cacheStatsCounters{0, requests, 1, int64(len(photo)), 1} + + if got := readCacheStatsCounters(t, svc.cache); got != want { + t.Errorf("counters = %+v, want %+v", got, want) + } +} + +// TestService_Get_ConcurrentVariantsStayApart requests three variants of the +// test photo at once that differ only in quality or fit. Each must be made by +// a fetch and a transcode of its own, and each answer must be its own variant. +func TestService_Get_ConcurrentVariantsStayApart(t *testing.T) { + t.Parallel() + + svc, fixtures, fetcher := setupHeldFetchService(t) + + cover := startGet(t.Context(), svc, photoVariant(fixtures, 85, FitCover)) + lowQuality := startGet(t.Context(), svc, photoVariant(fixtures, 40, FitCover)) + contain := startGet(t.Context(), svc, photoVariant(fixtures, 85, FitContain)) + + for range 3 { + select { + case <-fetcher.started: + case <-time.After(5 * time.Second): + t.Fatal("fewer fetches started than variants requested: " + + "variants differing in quality or fit were merged") + } + } + + fetcher.releaseFetches() + + images := make(map[string][]byte) + + for _, variant := range []struct { + name string + results <-chan getResult + width, height int + }{ + {"q=85 fit=cover", cover, 50, 25}, + {"q=40 fit=cover", lowQuality, 50, 25}, + {"q=85 fit=contain", contain, 25, 25}, + } { + got := <-variant.results + if got.err != nil { + t.Fatalf("%s: Get() error = %v", variant.name, got.err) + } + + width, height := jpegSize(got.image) + if width != variant.width || height != variant.height { + t.Errorf("%s: image is %dx%d, want a %dx%d JPEG", variant.name, + width, height, variant.width, variant.height) + } + + images[variant.name] = got.image + } + + if bytes.Equal(images["q=85 fit=cover"], images["q=40 fit=cover"]) { + t.Error("q=40 was answered with the q=85 image") + } +} + +// TestService_Get_WaiterStopsWhenItsContextEnds has a second request for a +// variant join the first one's held fetch, then ends the second request's +// context. The second request must return at once with the context's error, +// while the fetch is still held, and the first must still be answered. +func TestService_Get_WaiterStopsWhenItsContextEnds(t *testing.T) { + t.Parallel() + + svc, fixtures, fetcher := setupHeldFetchService(t) + + first := startGet(t.Context(), svc, photoVariant(fixtures, 85, FitCover)) + + <-fetcher.started + + waiterCtx, cancelWaiter := context.WithCancel(t.Context()) + waiter := startGet(waiterCtx, svc, photoVariant(fixtures, 85, FitCover)) + + time.Sleep(arrivalWait) + cancelWaiter() + + select { + case got := <-waiter: + if !errors.Is(got.err, context.Canceled) { + t.Errorf("waiting request: Get() error = %v, want %v", + got.err, context.Canceled) + } + case <-time.After(time.Second): + t.Fatal("waiting request did not return when its context ended") + } + + fetcher.releaseFetches() + + got := <-first + if got.err != nil { + t.Fatalf("first request: Get() error = %v", got.err) + } + + if width, height := jpegSize(got.image); width != 50 || height != 25 { + t.Errorf("first request: image is %dx%d, want a 50x25 JPEG", width, height) + } + + if fetches := fetcher.fetches.Load(); fetches != 1 { + t.Errorf("upstream fetches = %d, want 1", fetches) + } +} + +// TestService_Get_FirstRequestLeavingKeepsTheWork ends the context of the +// request whose fetch is held, after a second request has joined it. The fetch +// and transcode must go on and answer the second request. The first request +// waits for its own work, as every request did before misses were shared, and +// is answered too. +func TestService_Get_FirstRequestLeavingKeepsTheWork(t *testing.T) { + t.Parallel() + + svc, fixtures, fetcher := setupHeldFetchService(t) + + firstCtx, cancelFirst := context.WithCancel(t.Context()) + first := startGet(firstCtx, svc, photoVariant(fixtures, 85, FitCover)) + + <-fetcher.started + + second := startGet(t.Context(), svc, photoVariant(fixtures, 85, FitCover)) + + time.Sleep(arrivalWait) + cancelFirst() + fetcher.releaseFetches() + + for name, results := range map[string]<-chan getResult{ + "first request": first, "second request": second, + } { + got := <-results + if got.err != nil { + t.Fatalf("%s: Get() error = %v", name, got.err) + } + + if width, height := jpegSize(got.image); width != 50 || height != 25 { + t.Errorf("%s: image is %dx%d, want a 50x25 JPEG", name, width, height) + } + } + + if fetches := fetcher.fetches.Load(); fetches != 1 { + t.Errorf("upstream fetches = %d, want 1", fetches) + } +} + +// TestService_Get_ConcurrentMissesShareAFailure has several requests for an +// image that cannot be served arrive while its fetch is held: the one fetch +// answers all of them with its error. The request after them is answered from +// the negative cache when the failure is kept there, and fetches again when it +// is not. +func TestService_Get_ConcurrentMissesShareAFailure(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + path string + wantErr error // what every request at once gets + wantNextErr error // what the request after them gets + wantFetches int32 // fetches once the request after them is answered + }{ + {"upstream answers 404, kept in the negative cache", + "/images/missing.jpg", httpfetcher.ErrUpstreamError, + ErrNegativeCached, 1}, + {"source fails the magic byte check, not kept", + "/images/text.png", magic.ErrUnknownFormat, magic.ErrUnknownFormat, 2}, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + svc, fixtures, fetcher := setupHeldFetchService(t) + + request := func() *ImageRequest { + req := photoVariant(fixtures, 85, FitCover) + req.SourcePath = tc.path + + return req + } + + pending := make([]<-chan getResult, 0, 4) + for range 4 { + pending = append(pending, startGet(t.Context(), svc, request())) + } + + <-fetcher.started + time.Sleep(arrivalWait) + fetcher.releaseFetches() + + for i, results := range pending { + if got := <-results; !errors.Is(got.err, tc.wantErr) { + t.Errorf("request %d: Get() error = %v, want %v", i, got.err, tc.wantErr) + } + } + + _, err := svc.Get(t.Context(), request()) + if !errors.Is(err, tc.wantNextErr) { + t.Errorf("next request: Get() error = %v, want %v", err, tc.wantNextErr) + } + + if fetches := fetcher.fetches.Load(); fetches != tc.wantFetches { + t.Errorf("upstream fetches = %d, want %d", fetches, tc.wantFetches) + } + }) + } +} + +// TestService_Get_EndedRequestFetchesNothing checks that a request whose +// context has already ended when it misses the cache starts no fetch. +func TestService_Get_EndedRequestFetchesNothing(t *testing.T) { + t.Parallel() + + svc, fixtures, fetcher := setupHeldFetchService(t) + + ctx, cancel := context.WithCancel(t.Context()) + cancel() + + _, err := svc.Get(ctx, photoVariant(fixtures, 85, FitCover)) + if !errors.Is(err, context.Canceled) { + t.Errorf("Get() error = %v, want %v", err, context.Canceled) + } + + if fetches := fetcher.fetches.Load(); fetches != 0 { + t.Errorf("upstream fetches = %d, want 0", fetches) + } +} + +// panickingFetcher panics on every fetch. +type panickingFetcher struct{} + +func (panickingFetcher) Fetch( + context.Context, string, +) (*httpfetcher.FetchResult, error) { + panic("upstream fetcher panicked") +} + +// TestService_Get_PanicBecomesAnError checks that a panic while a variant is +// being made reaches its request as an error naming the panic, instead of +// being raised again. +func TestService_Get_PanicBecomesAnError(t *testing.T) { + t.Parallel() + + svc, fixtures := SetupTestService(t) + svc.fetcher = panickingFetcher{} + + var err error + + func() { + defer func() { + if recovered := recover(); recovered != nil { + t.Fatalf("Get() panicked: %v", recovered) + } + }() + + _, err = svc.Get(t.Context(), photoVariant(fixtures, 85, FitCover)) + }() + + if err == nil || !strings.Contains(err.Error(), "upstream fetcher panicked") { + t.Errorf("Get() error = %v, want one naming the panic", 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 }