From 9e1f5d4cea1037054212a0b7d57ff0431457727a Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 04:08:33 +0000 Subject: [PATCH] Test that a cached source is read only with a processing slot (closes #64) Failing test: with the only processing slot held, a request for a new width of a cached image waits for the slot while the cached file is rewritten; the answer must come from the rewritten file, so the request read none of the source before it had a slot. Also a test that a fetch whose request context ends while it waits for a connection shared by all hosts gives its host's slot back; that one passes already. Model: opus-5-5 --- .../max_connections_internal_test.go | 37 ++++++ ...max_concurrent_processing_internal_test.go | 125 ++++++++++++++++++ 2 files changed, 162 insertions(+) create mode 100644 internal/imgcache/max_concurrent_processing_internal_test.go diff --git a/internal/httpfetcher/max_connections_internal_test.go b/internal/httpfetcher/max_connections_internal_test.go index 235024b..d23816c 100644 --- a/internal/httpfetcher/max_connections_internal_test.go +++ b/internal/httpfetcher/max_connections_internal_test.go @@ -1,6 +1,7 @@ package httpfetcher import ( + "context" "errors" "net" "strconv" @@ -82,6 +83,42 @@ func TestFetchLimitsConnectionsToAllHostsTogether(t *testing.T) { _ = third.Content.Close() } +// TestFetchFreesHostSlotWhenContextEndsWaitingForConnection checks that a +// fetch whose request context ends while it waits for a connection shared +// by all hosts gives its host's slot back. With MaxConnections at 1 and one +// response open, a fetch from another host takes that host's slot and waits; +// its context ends long before the 10 second wait timeout. +func TestFetchFreesHostSlotWhenContextEndsWaitingForConnection(t *testing.T) { + t.Parallel() + + srv := startUpstream(t) + + cfg := DefaultConfig() + cfg.MaxConnections = 1 + + f, _ := newServerFetcher(t, srv, cfg) + + first, err := f.Fetch(testContext(t), imageURLOnPort(81)) + if err != nil { + t.Fatalf("first Fetch() error = %v", err) + } + + defer func() { _ = first.Content.Close() }() + + ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond) + defer cancel() + + _, err = f.Fetch(ctx, imageURLOnPort(82)) + if !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("second Fetch() error = %v, want context.DeadlineExceeded", err) + } + + if held := semLen(f, testPublicHost+":82"); held != 0 { + t.Errorf("the fetch kept its host's slot after its context ended: "+ + "%d held", held) + } +} + // TestFetchReleasesConnectionOnError checks that a fetch that fails after // taking its connection gives it back: with MaxConnections at 1, the slot // must be free after the failure and the next fetch must succeed. diff --git a/internal/imgcache/max_concurrent_processing_internal_test.go b/internal/imgcache/max_concurrent_processing_internal_test.go new file mode 100644 index 0000000..e4c310e --- /dev/null +++ b/internal/imgcache/max_concurrent_processing_internal_test.go @@ -0,0 +1,125 @@ +package imgcache + +import ( + "image/color" + "image/jpeg" + "io" + "os" + "testing" + "time" + + "sneak.berlin/go/pixa/internal/imageprocessor" +) + +// widthOnlyRequest asks for the test photo at width, its height scaled to +// keep the photo's aspect ratio. +func widthOnlyRequest(fixtures *TestFixtures, width int) *ImageRequest { + return &ImageRequest{ + SourceHost: fixtures.GoodHost, + SourcePath: testPathPhoto, + Size: Size{Width: width}, + Format: FormatJPEG, + Quality: 85, + FitMode: FitCover, + } +} + +// holdProcessingSlot takes one of proc's processing slots and returns the +// func that gives it back. Process takes its slot before it reads its input, +// so once it has read a byte from the pipe it holds the slot, until the pipe +// is closed. +func holdProcessingSlot( + t *testing.T, proc *imageprocessor.ImageProcessor, +) func() { + t.Helper() + + input, feed := io.Pipe() + + go func() { + _, _ = proc.Process(t.Context(), input, &imageprocessor.Request{}) + }() + + _, err := feed.Write([]byte{0}) + if err != nil { + t.Fatalf("Process call to hold the slot did not start: %v", err) + } + + release := func() { _ = feed.Close() } + t.Cleanup(release) + + return release +} + +// TestService_Get_WaitsForSlotBeforeReadingCachedSource checks that a +// request whose source is cached holds none of it while it waits for a +// processing slot: it reads the cached file only once it has a slot. With +// the only slot held, a request for a new width of the cached 100x100 photo +// waits; the cached file is then rewritten as a 100x50 image before the slot +// is freed, so the request must answer with that image scaled to 40x20. +func TestService_Get_WaitsForSlotBeforeReadingCachedSource(t *testing.T) { + t.Parallel() + + svc, fixtures := SetupTestService(t) + svc.processor = imageprocessor.New( + imageprocessor.Params{MaxConcurrentProcessing: 1}, + ) + + // A first request caches the photo as a source. + resp, err := svc.Get(t.Context(), widthOnlyRequest(fixtures, 50)) + if err != nil { + t.Fatalf("first Get() error = %v", err) + } + + _ = resp.Content.Close() + + contentHash, _, err := svc.cache.LookupSource(t.Context(), + widthOnlyRequest(fixtures, 50)) + if err != nil || contentHash == "" { + t.Fatalf("LookupSource() = %q, %v; want the cached source", + contentHash, err) + } + + release := holdProcessingSlot(t, svc.processor) + + var ( + waited *ImageResponse + waitedErr error + ) + + done := make(chan struct{}) + + go func() { + defer close(done) + + waited, waitedErr = svc.Get(t.Context(), widthOnlyRequest(fixtures, 40)) + }() + + // Give the request time to reach the slot: had it read the cached source + // before waiting, it would have read it by now. + time.Sleep(100 * time.Millisecond) + + err = os.WriteFile(svc.cache.srcContent.hashToPath(contentHash), + generateTestJPEG(t, 100, 50, color.RGBA{0, 0, 255, 255}), 0o600) + if err != nil { + t.Fatalf("failed to rewrite the cached source: %v", err) + } + + release() + <-done + + if waitedErr != nil { + t.Fatalf("Get() error = %v", waitedErr) + } + + defer func() { _ = waited.Content.Close() }() + + output, err := jpeg.DecodeConfig(waited.Content) + if err != nil { + t.Fatalf("failed to decode the response: %v", err) + } + + if output.Width != 40 || output.Height != 20 { + t.Errorf("response is %dx%d, want 40x20: the request read the cached "+ + "source before it had a processing slot", output.Width, output.Height) + } +}