Author SHA1 Message Date
clawbot 4009490242 Share one fetch and transcode among concurrent misses (closes #65)
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
2026-09-29 11:02:03 +00:00
clawbot bcd5363d99 Test that concurrent misses for one variant share one fetch (closes #65)
Requests that miss the same variant while its fetch is held must make
one fetch and one transcode between them, all get the same image and
each count one miss. Variants differing only in quality or fit stay
apart. A waiting request whose context ends returns at once; the first
request's client leaving does not stop the work. A shared failure
reaches every request, the negative cache answers the next one, and an
uncached failure is fetched again. A request that has already ended
fetches nothing, and a panic comes back as an error. All but the
variants test fail on next.

Model: opus-5-5
2026-09-29 11:01:33 +00:00
9 changed files with 596 additions and 283 deletions
+7
View File
@@ -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
+11
View File
@@ -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
-80
View File
@@ -1,80 +0,0 @@
package main
import (
"errors"
"testing"
"go.uber.org/fx"
)
// errTestHook is the error returned by the test hooks that fail.
var errTestHook = errors.New("test hook failed")
// TestRunAppExitCode checks the exit code runApp returns: the one a
// shutdown request carries, 0 for a request without one (as for SIGINT or
// SIGTERM), and 1 when the app fails to start or to stop.
func TestRunAppExitCode(t *testing.T) {
t.Parallel()
cases := []struct {
name string
hook func(shutdowner fx.Shutdowner) fx.Hook
want int
}{
{
name: "shutdown requested with exit code 1",
hook: func(shutdowner fx.Shutdowner) fx.Hook {
return fx.StartHook(func() error {
return shutdowner.Shutdown(fx.ExitCode(1))
})
},
want: 1,
},
{
name: "shutdown requested without an exit code",
hook: func(shutdowner fx.Shutdowner) fx.Hook {
return fx.StartHook(func() error {
return shutdowner.Shutdown()
})
},
want: 0,
},
{
name: "start fails",
hook: func(fx.Shutdowner) fx.Hook {
return fx.StartHook(func() error { return errTestHook })
},
want: 1,
},
{
name: "stop fails",
hook: func(shutdowner fx.Shutdowner) fx.Hook {
return fx.StartStopHook(
func() error { return shutdowner.Shutdown() },
func() error { return errTestHook },
)
},
want: 1,
},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
app := fx.New(
fx.NopLogger,
fx.Invoke(func(lc fx.Lifecycle, shutdowner fx.Shutdowner) {
lc.Append(tc.hook(shutdowner))
}),
)
got := runApp(app)
t.Logf("runApp() = %d", got)
if got != tc.want {
t.Errorf("runApp() = %d, want %d", got, tc.want)
}
})
}
}
+1 -1
View File
@@ -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
@@ -300,69 +300,3 @@ func TestProcessReleasesSlotOnError(t *testing.T) {
})
}
}
// TestWaitForProcessing holds a processing slot with a Process call that
// cannot finish reading its input. WaitForProcessing must report that image
// when its context ends first, wait for it otherwise, return 0 once it has
// finished, and give back the slots it took while waiting.
func TestWaitForProcessing(t *testing.T) {
t.Parallel()
proc := New(Params{MaxConcurrentProcessing: 2})
gate := make(chan struct{})
entered := make(chan struct{}, 1)
results := make(chan error, 1)
openGate := sync.OnceFunc(func() { close(gate) })
t.Cleanup(openGate)
processInBackground(proc, &gatedReader{
data: bytes.NewReader(createTestJPEG(t, 10, 10)), gate: gate,
entered: entered, counter: &readingCounter{},
}, results)
waitForEntries(t, entered, 1)
ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond)
defer cancel()
stillProcessing := proc.WaitForProcessing(ctx)
t.Logf("WaitForProcessing() after its context ended: %d", stillProcessing)
if stillProcessing != 1 {
t.Errorf("WaitForProcessing() after its context ended = %d, want 1",
stillProcessing)
}
waited := make(chan int, 1)
go func() { waited <- proc.WaitForProcessing(t.Context()) }()
select {
case got := <-waited:
t.Fatalf("WaitForProcessing() = %d while an image was being processed",
got)
case <-time.After(100 * time.Millisecond):
}
openGate()
err := <-results
if err != nil {
t.Errorf("Process() error = %v, want nil", err)
}
select {
case got := <-waited:
if got != 0 {
t.Errorf("WaitForProcessing() once processing finished = %d, want 0",
got)
}
case <-time.After(5 * time.Second):
t.Fatal("WaitForProcessing() did not return once processing finished")
}
if held := len(proc.processingSemaphore); held != 0 {
t.Errorf("%d slots still held after WaitForProcessing() returned", held)
}
}
+13 -4
View File
@@ -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,8 +496,17 @@ 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, `
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 + ?,
@@ -507,7 +517,6 @@ func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64)
c.log.Warn("failed to count upstream fetch",
"fetch_bytes", fetchBytes, "error", err)
}
}
}
// IncrementTransformCount counts one image transcoded by the image processor.
@@ -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)
}
}
+111 -27
View File
@@ -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
}
-96
View File
@@ -1,96 +0,0 @@
package server
import (
"log/slog"
"net"
"testing"
"time"
"go.uber.org/fx"
"go.uber.org/fx/fxtest"
"sneak.berlin/go/pixa/internal/config"
"sneak.berlin/go/pixa/internal/globals"
"sneak.berlin/go/pixa/internal/logger"
)
// shutdownRecorder is an fx.Shutdowner that sends the options of each
// shutdown request on requests.
type shutdownRecorder struct {
requests chan []fx.ShutdownOption
}
func (r shutdownRecorder) Shutdown(opts ...fx.ShutdownOption) error {
r.requests <- opts
return nil
}
// TestSentryInitFailureFailsStartup checks that a Sentry DSN that cannot be
// used makes the server's start hook fail, so fx stops what has already
// started, instead of the process exiting from a goroutine.
func TestSentryInitFailureFailsStartup(t *testing.T) {
t.Parallel()
lc := fxtest.NewLifecycle(t)
log, err := logger.New(lc, logger.Params{Globals: &globals.Globals{}})
if err != nil {
t.Fatalf("logger.New() error = %v", err)
}
_, err = New(lc, Params{
Logger: log,
Globals: &globals.Globals{Appname: "pixad"},
Config: &config.Config{SentryDSN: "not-a-dsn"},
})
if err != nil {
t.Fatalf("New() error = %v", err)
}
err = lc.Start(t.Context())
t.Logf("Start() error = %v", err)
if err == nil {
t.Fatal("Start() error = nil, want the Sentry initialization error")
}
}
// TestListenErrorRequestsShutdownWithExitCode1 occupies the server's port
// and checks that the listen error asks fx to shut down with exit code 1.
func TestListenErrorRequestsShutdownWithExitCode1(t *testing.T) {
t.Parallel()
busy, err := (&net.ListenConfig{}).Listen(t.Context(), "tcp", ":0")
if err != nil {
t.Fatalf("Listen() error = %v", err)
}
t.Cleanup(func() { _ = busy.Close() })
addr, ok := busy.Addr().(*net.TCPAddr)
if !ok {
t.Fatalf("listener address %v is not a TCP address", busy.Addr())
}
requests := make(chan []fx.ShutdownOption, 1)
s := &Server{
log: slog.New(slog.DiscardHandler),
config: &config.Config{Port: addr.Port},
shutdowner: shutdownRecorder{requests: requests},
}
s.httpServer = s.newHTTPServer()
go s.serveUntilShutdown()
select {
case opts := <-requests:
t.Logf("shutdown options = %v", opts)
if len(opts) != 1 || opts[0] != fx.ExitCode(1) {
t.Errorf("shutdown options = %v, want [fx.ExitCode(1)]", opts)
}
case <-time.After(5 * time.Second):
t.Fatal("no shutdown was requested after the listen error")
}
}