diff --git a/Dockerfile b/Dockerfile index 8afc3dd..e82d2d0 100644 --- a/Dockerfile +++ b/Dockerfile @@ -32,10 +32,13 @@ RUN script/bootstrap --cgo COPY . . # Without -v first; on a failure, again with -v for the details, and -# the step fails even if the second run passes. -RUN go test -count=1 -timeout 90s -race -cover ./... || \ +# the step fails even if the second run passes. -parallel 4: by default +# a package runs as many of its tests at once as the host has CPUs, and +# on a busy host they then wait so long to be scheduled that a test +# that times a request can see it return a second late. +RUN go test -count=1 -timeout 90s -race -parallel 4 -cover ./... || \ { echo "--- Rerunning with -v for details ---"; \ - go test -count=1 -timeout 90s -race -v ./...; exit 1; } + go test -count=1 -timeout 90s -race -parallel 4 -v ./...; exit 1; } # Build stage. Nothing is wanted from the two phases above: these copies # make BuildKit build them first, so this stage runs only when lint and diff --git a/README.md b/README.md index 8b6b032..04768c7 100644 --- a/README.md +++ b/README.md @@ -83,10 +83,12 @@ which answers 200 whenever pixa is running, in maintenance mode too (see `maintenance_mode`). On SIGTERM or SIGINT pixa stops accepting connections, gives the requests in -progress and the images being processed 5 seconds to finish, and exits: with 0, -or with 1 when images were still being processed after those 5 seconds or -another part of pixa failed to stop. A request not finished by then is cut off. -`docker stop` waits 10 seconds before it kills the container. +progress, the images being processed and the counts of cache hits, misses, +fetches and conversions being written to the database 5 seconds to finish, and +exits: with 0, or with 1 when images were still being processed or counts still +being written after those 5 seconds, or another part of pixa failed to stop. A +request not finished by then is cut off. `docker stop` waits 10 seconds before +it kills the container. Outside Docker, pixa needs libvips (the image has 8.16) and libheif to run, as it uses libvips through CGO. pixad does not start unless libvips has its JPEG XL diff --git a/TODO.md b/TODO.md index 7bee655..b2971a4 100644 --- a/TODO.md +++ b/TODO.md @@ -30,6 +30,21 @@ P2: security: per-IP rate limiting on the image routes # Completed Steps +- 2026-10-08 a request past its deadline no longer waits for a database write + (closes #224). On a busy host, `TestService_Get_ReturnsByItsDeadline` failed + because its request returned over a second after its deadline: it was still + waiting for its miss count to be written to the database, a write with no + deadline, which in pixad waits for the one database connection that every + request shares. The hit, miss, upstream fetch and transform counts are now + each written in a goroutine of their own, which the request waits for only + until its deadline; a count not written by then is written after the request + returns. On shutdown pixad waits for those writes within the same 5 seconds as + for the images being processed, and exits with 1 when some are unfinished. The + test phase of the `Dockerfile` also runs at most 4 tests of a package at once + (`go test -parallel 4`), where it ran as many as the host has CPUs: on a busy + host they then wait so long to be scheduled that a test that times a request + can fail. `-p`, how many packages are tested at once, is unchanged, as capping + it also slows compiling. - 2026-10-08 every output format is saved with settings pixa sets on purpose (closes #232): each format has its own govips export, as JPEG XL does, in place of govips' generic `Export`, which sent libvips a zero for some settings diff --git a/internal/handlers/handlers.go b/internal/handlers/handlers.go index 0414026..6575f45 100644 --- a/internal/handlers/handlers.go +++ b/internal/handlers/handlers.go @@ -93,6 +93,13 @@ func (s *Handlers) WaitForProcessing(ctx context.Context) int { return s.imgSvc.WaitForProcessing(ctx) } +// WaitForCountWrites waits until every count a request has started writing +// to the database is written, or until ctx ends, and reports whether they +// all were. +func (s *Handlers) WaitForCountWrites(ctx context.Context) bool { + return s.imgSvc.WaitForCountWrites(ctx) +} + // newCacheConfig builds the image cache's configuration from cfg. // cache_max_bytes: 0 disables the disk cache entirely; any other value // is the eviction limit in bytes; when it is omitted, the cache works diff --git a/internal/imgcache/count_writes_internal_test.go b/internal/imgcache/count_writes_internal_test.go new file mode 100644 index 0000000..19642b9 --- /dev/null +++ b/internal/imgcache/count_writes_internal_test.go @@ -0,0 +1,116 @@ +package imgcache + +import ( + "context" + "errors" + "testing" + "time" +) + +// holdDatabase takes the one connection of the test service's database, so +// that every other query waits for it, and returns the func that frees it. +func holdDatabase(t *testing.T, svc *Service) func() { + t.Helper() + + conn, err := svc.cache.db.Conn(t.Context()) + if err != nil { + t.Fatalf("failed to take the database connection: %v", err) + } + + release := func() { _ = conn.Close() } + t.Cleanup(release) + + return release +} + +// TestService_Get_ReturnsByItsDeadlineWhileTheDatabaseIsBusy holds the +// database's one connection while a request whose fetch is held reaches its +// deadline. The request must still return by its deadline with the +// deadline's error, and its miss must be counted once the connection is free. +func TestService_Get_ReturnsByItsDeadlineWhileTheDatabaseIsBusy(t *testing.T) { + t.Parallel() + + svc, fixtures, fetcher := setupHeldFetchService(t) + + const timeout = 200 * time.Millisecond + + ctx, cancel := context.WithTimeout(t.Context(), timeout) + defer cancel() + + results := startGet(ctx, svc, photoVariant(fixtures, 85, FitCover)) + + // The request has made its database reads by the time it fetches. + <-fetcher.started + + releaseDatabase := holdDatabase(t, svc) + + select { + case got := <-results: + deadline, _ := ctx.Deadline() + t.Logf("Get() returned %v after its deadline, error = %v", + time.Since(deadline), got.err) + + if !errors.Is(got.err, context.DeadlineExceeded) { + t.Errorf("Get() error = %v, want %v", got.err, context.DeadlineExceeded) + } + case <-time.After(timeout + time.Second): + t.Fatal("request did not return by its deadline while the database was busy") + } + + releaseDatabase() + + waitCtx, cancelWait := context.WithTimeout(t.Context(), 5*time.Second) + defer cancelWait() + + if !svc.WaitForCountWrites(waitCtx) { + t.Fatal("the miss was not counted once the database was free") + } + + want := cacheStatsCounters{missCount: 1} + if got := readCacheStatsCounters(t, svc.cache); got != want { + t.Errorf("counters = %+v, want %+v", got, want) + } +} + +// TestService_WaitForCountWrites holds the database's one connection while a +// request past its deadline counts a miss. WaitForCountWrites must report the +// count unwritten when its context ends, and written once the connection is +// free. +func TestService_WaitForCountWrites(t *testing.T) { + t.Parallel() + + svc, _, _ := setupHeldFetchService(t) + + releaseDatabase := holdDatabase(t, svc) + + ended, cancel := context.WithDeadline(t.Context(), time.Now()) + defer cancel() + + svc.writeCount(ended, func(ctx context.Context) { + svc.cache.IncrementStats(ctx, false, 0) + }) + + shortCtx, cancelShort := context.WithTimeout(t.Context(), 100*time.Millisecond) + defer cancelShort() + + written := svc.WaitForCountWrites(shortCtx) + t.Logf("WaitForCountWrites() while the database was busy = %t", written) + + if written { + t.Fatal("WaitForCountWrites() = true while the database was busy") + } + + releaseDatabase() + + waitCtx, cancelWait := context.WithTimeout(t.Context(), 5*time.Second) + defer cancelWait() + + if !svc.WaitForCountWrites(waitCtx) { + t.Fatal("WaitForCountWrites() = false once the database was free") + } + + want := cacheStatsCounters{missCount: 1} + if got := readCacheStatsCounters(t, svc.cache); got != want { + t.Errorf("counters = %+v, want %+v", got, want) + } +} diff --git a/internal/imgcache/service.go b/internal/imgcache/service.go index f196b26..90e3e0c 100644 --- a/internal/imgcache/service.go +++ b/internal/imgcache/service.go @@ -9,6 +9,7 @@ import ( "log/slog" "net/url" "runtime/debug" + "sync" "time" "github.com/dustin/go-humanize" @@ -36,6 +37,9 @@ type Service struct { // variantsInProgress lets the requests that miss the same variant at the // same time share one fetch and one transcode. variantsInProgress singleflight.Group + // countWrites holds the count writes writeCount has started, for + // WaitForCountWrites. + countWrites sync.WaitGroup } // ServiceConfig holds configuration for the image service. @@ -161,8 +165,9 @@ func (s *Service) Get(ctx context.Context, req *ImageRequest) (*ImageResponse, e s.log.Error("failed to get cached variant", "key", result.CacheKey, "error", err) // Fall through to re-process } else { - // Counted also when the request context has ended meanwhile - s.cache.IncrementStats(context.WithoutCancel(ctx), true, 0) + s.writeCount(ctx, func(ctx context.Context) { + s.cache.IncrementStats(ctx, true, 0) + }) return &ImageResponse{ Content: reader, @@ -179,7 +184,9 @@ func (s *Service) Get(ctx context.Context, req *ImageRequest) (*ImageResponse, e // failed or the request context has ended meanwhile response, err := s.processOrWait(ctx, req) - s.cache.IncrementStats(context.WithoutCancel(ctx), false, 0) + s.writeCount(ctx, func(ctx context.Context) { + s.cache.IncrementStats(ctx, false, 0) + }) if err != nil { return nil, err @@ -208,6 +215,25 @@ func (s *Service) WaitForProcessing(ctx context.Context) int { return s.processor.WaitForProcessing(ctx) } +// WaitForCountWrites waits until every count a request has started writing +// to the database is written, or until ctx ends, and reports whether they +// all were. +func (s *Service) WaitForCountWrites(ctx context.Context) bool { + written := make(chan struct{}) + + go func() { + s.countWrites.Wait() + close(written) + }() + + select { + case <-written: + return true + case <-ctx.Done(): + return false + } +} + // ValidateRequest validates the request signature if required. func (s *Service) ValidateRequest(req *ImageRequest) error { // Check if host is allowed (no signature required) @@ -254,6 +280,38 @@ func (s *Service) GenerateSignedURL( baseURL, path, sig, exp, req.Quality, req.FitMode), nil } +// writeCount runs write, which adds to a count in the database, in a +// goroutine of its own. The write gets ctx without its cancellation or +// deadline, so an ended request is still counted. writeCount waits for it +// until ctx's deadline, even if ctx is cancelled first, so a request +// returns with its count written unless its deadline has passed; past +// that, it does not wait for the database connection, which every request +// shares. +func (s *Service) writeCount(ctx context.Context, write func(context.Context)) { + written := make(chan struct{}) + + s.countWrites.Go(func() { + defer close(written) + + write(context.WithoutCancel(ctx)) + }) + + deadline, hasDeadline := ctx.Deadline() + if !hasDeadline { + <-written + + return + } + + deadlineTimer := time.NewTimer(time.Until(deadline)) + defer deadlineTimer.Stop() + + select { + case <-written: + case <-deadlineTimer.C: + } +} + // errPanicked is returned when processing a variant panicked. var errPanicked = errors.New("panic while processing image") @@ -452,8 +510,9 @@ 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) + s.writeCount(ctx, func(ctx context.Context) { + s.cache.IncrementUpstreamFetch(ctx, fetchBytes) + }) if err != nil { return nil, fmt.Errorf("failed to read upstream response: %w", err) @@ -534,8 +593,7 @@ func (s *Service) processAndStore( processDuration := time.Since(processStart) - // Counted also when the request context has ended meanwhile - s.cache.IncrementTransformCount(context.WithoutCancel(ctx)) + s.writeCount(ctx, s.cache.IncrementTransformCount) // Read processed content processedData, err := io.ReadAll(processResult.Content) diff --git a/internal/server/server.go b/internal/server/server.go index 7be7d67..a253fd6 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -29,6 +29,10 @@ const ( // still being processed once ShutdownTimeout has passed. var errStillProcessing = errors.New("images still being processed at shutdown") +// errStillWritingCounts is returned by the server's stop hook when counts +// are still being written to the database once ShutdownTimeout has passed. +var errStillWritingCounts = errors.New("counts still being written at shutdown") + // Params defines dependencies for Server. type Params struct { fx.In @@ -117,9 +121,11 @@ func (s *Server) enableSentry() error { } // cleanShutdown stops the HTTP server, waits for the images still being -// processed, then flushes Sentry. The first two share ShutdownTimeout. It -// returns errStillProcessing when images are still being processed after -// that, as their work is abandoned. +// processed and then for the counts still being written to the database, +// which closes after this hook, then flushes Sentry. The first three share +// ShutdownTimeout. It returns errStillProcessing and errStillWritingCounts +// for the images and counts still unfinished after that, as their work is +// abandoned. func (s *Server) cleanShutdown(ctx context.Context) error { s.log.Info("shutting down") @@ -132,17 +138,26 @@ func (s *Server) cleanShutdown(ctx context.Context) error { } stillProcessing := s.h.WaitForProcessing(ctxShutdown) + countsWritten := s.h.WaitForCountWrites(ctxShutdown) if s.sentryEnabled { sentry.Flush(SentryFlushTimeout) } + var unfinished []error + if stillProcessing > 0 { s.log.Error("images still being processed at shutdown", "count", stillProcessing) - return errStillProcessing + unfinished = append(unfinished, errStillProcessing) } - return nil + if !countsWritten { + s.log.Error("counts still being written at shutdown") + + unfinished = append(unfinished, errStillWritingCounts) + } + + return errors.Join(unfinished...) }