Author SHA1 Message Date
sneak df7fff0ce1 Stop waiting for count writes after a request's deadline (closes #224)
check / check (push) Canceled after 0s
TestService_Get_ReturnsByItsDeadline failed on a busy host because its
request, past its deadline, still waited for its miss count to be
written, a database write with no deadline; in pixad that write waits
for the one database connection 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 it returns. Shutdown waits
for those writes within ShutdownTimeout and reports any left unfinished.
The test phase also runs go test with -parallel 4: on a busy host, as
many tests at once as there are CPUs wait so long to be scheduled that
a timed request can fail.

Model: opus-5-5
2026-10-08 13:29:06 +00:00
9 changed files with 243 additions and 586 deletions
+6 -6
View File
@@ -83,12 +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, then writes to the
database the counts of cache hits, misses, fetches and conversions that requests
could not write by their deadline, and exits: with 0, or with 1 when images were
still being processed after those 5 seconds, some of those counts could not be
written, 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
+15 -13
View File
@@ -30,19 +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 the database to
count it (closes #224). On a busy host, `TestService_Get_ReturnsByItsDeadline`
failed because its request, past its deadline, still waited for its miss count
to be written, a write with no deadline, which in pixad waits for the one
database connection every request shares. Each count write now keeps the
request's deadline but not its cancellation. A count not written by then is
kept in memory, where `Stats` sees it, and one goroutine of the cache writes
those counts, one write at a time, and once more at shutdown before the
database closes. 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 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
+9 -11
View File
@@ -4,7 +4,6 @@ package handlers
import (
"context"
"encoding/json"
"errors"
"log/slog"
"net/http"
"time"
@@ -72,21 +71,16 @@ func New(lc fx.Lifecycle, params Params) (*Handlers, error) {
}
lc.Append(fx.Hook{
//nolint:contextcheck // the cache's goroutines outlive OnStart; OnStop stops them
//nolint:contextcheck // the eviction loop outlives OnStart; OnStop cancels it
OnStart: func(_ context.Context) error {
return s.initImageService()
},
// The pending counts are written here, before the database's stop
// hook closes it.
OnStop: func(ctx context.Context) error {
if s.imgCache == nil {
return nil
}
return errors.Join(
s.imgCache.StopEviction(ctx),
s.imgCache.StopPendingCountWrites(ctx),
)
return s.imgCache.StopEviction(ctx)
},
})
@@ -99,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
@@ -128,9 +129,6 @@ func (s *Handlers) initImageService() error {
// write-pressure passes. No-op when the disk cache is disabled.
cache.StartEviction(imgcache.DefaultEvictionInterval)
// Writes the counts requests could not write by their deadline
cache.StartPendingCountWrites()
// Create the fetcher config
fetcherCfg := httpfetcher.DefaultConfig()
fetcherCfg.AllowHTTP = s.config.AllowHTTP
+8 -71
View File
@@ -11,8 +11,6 @@ import (
"io"
"log/slog"
"path/filepath"
"sync"
"sync/atomic"
"time"
lru "github.com/hashicorp/golang-lru/v2"
@@ -81,26 +79,6 @@ type Cache struct {
evictionDone chan struct{}
evictionCancel context.CancelFunc
// The pending counts: hits, misses, upstream fetches and transforms
// whose write to the database missed its request's deadline, kept
// here until they are written by the goroutine that
// StartPendingCountWrites starts. As for eviction, the channels are
// created in NewCache: pendingCountsAdded wakes that goroutine, and
// pendingCountsDone is closed when it returns. pendingCountsCancel,
// set by StartPendingCountWrites, cancels its context, which tells it
// to finish.
pendingHits atomic.Int64
pendingMisses atomic.Int64
pendingUpstreamFetches atomic.Int64
pendingUpstreamFetchBytes atomic.Int64
pendingTransforms atomic.Int64
pendingCountsAdded chan struct{}
pendingCountsDone chan struct{}
pendingCountsCancel context.CancelFunc
// Held through each write of the pending counts, and by Stats to read them
pendingCountsWriteMutex sync.Mutex
// metaCache holds the content types of the variants most recently
// stored or served, so a hit does not read the variant's .meta file.
// It never stands in for the variant file, which is always opened.
@@ -161,9 +139,6 @@ func newCache(
metaCache: metaCache,
contentLocks: newContentLock(),
pendingCountsAdded: make(chan struct{}, 1),
pendingCountsDone: make(chan struct{}),
reconciliationPageSize: defaultReconciliationPageSize,
}
@@ -499,21 +474,12 @@ func (c *Cache) CleanExpired(ctx context.Context) error {
func (c *Cache) Stats(ctx context.Context) (*CacheStats, error) {
var stats CacheStats
// So that no write of the pending counts falls between the two reads
c.pendingCountsWriteMutex.Lock()
// Fetch hit/miss counts from the stats table
err := c.db.QueryRowContext(ctx, `
SELECT hit_count, miss_count
FROM cache_stats WHERE id = 1
`).Scan(&stats.HitCount, &stats.MissCount)
// Hits and misses not yet written count too (see StartPendingCountWrites)
stats.HitCount += c.pendingHits.Load()
stats.MissCount += c.pendingMisses.Load()
c.pendingCountsWriteMutex.Unlock()
if err != nil && !errors.Is(err, sql.ErrNoRows) {
return nil, fmt.Errorf("failed to get cache stats: %w", err)
}
@@ -545,30 +511,19 @@ func (c *Cache) Stats(ctx context.Context) (*CacheStats, error) {
}
// IncrementStats counts a cache hit or miss, and an upstream fetch that read
// fetchBytes bytes, as IncrementUpstreamFetch does. Like the other Increment
// methods, it writes to the database with ctx's deadline but not its
// cancellation, so an ended request is still counted but does not wait for
// the database past its deadline; a count not written by then becomes a
// pending count (see StartPendingCountWrites).
// fetchBytes bytes, as IncrementUpstreamFetch does.
func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64) {
countCtx, cancel := withoutCancelKeepingDeadline(ctx)
defer cancel()
var err error
pendingCount := &c.pendingMisses
if hit {
pendingCount = &c.pendingHits
_, err = c.db.ExecContext(countCtx, `
_, err = c.db.ExecContext(ctx, `
UPDATE cache_stats
SET hit_count = hit_count + 1,
last_updated_at = CURRENT_TIMESTAMP
WHERE id = 1
`)
} else {
_, err = c.db.ExecContext(countCtx, `
_, err = c.db.ExecContext(ctx, `
UPDATE cache_stats
SET miss_count = miss_count + 1,
last_updated_at = CURRENT_TIMESTAMP
@@ -576,10 +531,7 @@ func (c *Cache) IncrementStats(ctx context.Context, hit bool, fetchBytes int64)
`)
}
switch {
case errors.Is(err, context.DeadlineExceeded):
c.addPendingCount(pendingCount, 1)
case err != nil:
if err != nil {
c.log.Warn("failed to count cache hit or miss", "hit", hit, "error", err)
}
@@ -593,22 +545,14 @@ func (c *Cache) IncrementUpstreamFetch(ctx context.Context, fetchBytes int64) {
return
}
countCtx, cancel := withoutCancelKeepingDeadline(ctx)
defer cancel()
_, err := c.db.ExecContext(countCtx, `
_, 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)
switch {
case errors.Is(err, context.DeadlineExceeded):
c.addPendingCount(&c.pendingUpstreamFetches, 1)
c.addPendingCount(&c.pendingUpstreamFetchBytes, fetchBytes)
case err != nil:
if err != nil {
c.log.Warn("failed to count upstream fetch",
"fetch_bytes", fetchBytes, "error", err)
}
@@ -616,20 +560,13 @@ func (c *Cache) IncrementUpstreamFetch(ctx context.Context, fetchBytes int64) {
// IncrementTransformCount counts one image transcoded by the image processor.
func (c *Cache) IncrementTransformCount(ctx context.Context) {
countCtx, cancel := withoutCancelKeepingDeadline(ctx)
defer cancel()
_, err := c.db.ExecContext(countCtx, `
_, err := c.db.ExecContext(ctx, `
UPDATE cache_stats
SET transform_count = transform_count + 1,
last_updated_at = CURRENT_TIMESTAMP
WHERE id = 1
`)
switch {
case errors.Is(err, context.DeadlineExceeded):
c.addPendingCount(&c.pendingTransforms, 1)
case err != nil:
if err != nil {
c.log.Warn("failed to count transform", "error", err)
}
}
@@ -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(writeCtx context.Context) {
svc.cache.IncrementStats(writeCtx, 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)
}
}
-141
View File
@@ -1,141 +0,0 @@
package imgcache
import (
"context"
"fmt"
"sync/atomic"
)
// StartPendingCountWrites starts the goroutine that writes the pending
// counts to the database whenever there are some, one UPDATE at a time:
// however many requests pass their deadline before their counts are
// written, at most this one write waits for the database. It is a no-op
// when already started. The goroutine outlives the caller, so it runs with
// its own context, which StopPendingCountWrites cancels to tell it to
// finish.
func (c *Cache) StartPendingCountWrites() {
if c.pendingCountsCancel != nil {
return
}
ctx, cancel := context.WithCancel(context.Background())
c.pendingCountsCancel = cancel
go func() {
c.pendingCountWriteLoop(ctx)
}()
}
// StopPendingCountWrites tells the goroutine StartPendingCountWrites
// started, if it did, to finish, and waits for its write under way to end
// and for it to return. It then writes the pending counts left, so that
// they are written before the database closes. It waits for all of this at
// most until ctx ends, then logs the counts not written and returns an
// error.
func (c *Cache) StopPendingCountWrites(ctx context.Context) error {
var err error
if c.pendingCountsCancel != nil {
c.pendingCountsCancel()
select {
case <-c.pendingCountsDone:
case <-ctx.Done():
// The write under way is left to finish: the database's
// close waits for a query under way.
err = fmt.Errorf("pending counts still being written: %w", ctx.Err())
}
}
if err == nil {
err = c.writePendingCounts(ctx)
}
if err != nil {
c.log.Error("counts not written at shutdown",
"hits", c.pendingHits.Load(),
"misses", c.pendingMisses.Load(),
"upstream_fetches", c.pendingUpstreamFetches.Load(),
"upstream_fetch_bytes", c.pendingUpstreamFetchBytes.Load(),
"transforms", c.pendingTransforms.Load(),
"error", err,
)
return err
}
return nil
}
// pendingCountWriteLoop is the body of the goroutine
// StartPendingCountWrites starts. It returns when ctx is cancelled, but
// does not cut off a write under way then. A write that fails leaves the
// counts pending, for the next write.
func (c *Cache) pendingCountWriteLoop(ctx context.Context) {
defer close(c.pendingCountsDone)
for {
select {
case <-ctx.Done():
return
case <-c.pendingCountsAdded:
}
err := c.writePendingCounts(context.WithoutCancel(ctx))
if err != nil {
c.log.Warn("failed to write pending counts", "error", err)
}
}
}
// addPendingCount adds n to pendingCount, one of the pending counts, and
// wakes the goroutine that writes them. pendingCountsAdded has capacity one,
// so a wakeup already waiting is enough.
func (c *Cache) addPendingCount(pendingCount *atomic.Int64, n int64) {
pendingCount.Add(n)
select {
case c.pendingCountsAdded <- struct{}{}:
default:
}
}
// writePendingCounts adds the pending counts to the cache_stats row in one
// UPDATE, then takes what it wrote out of them, leaving any added
// meanwhile. When the UPDATE fails, they all stay pending.
func (c *Cache) writePendingCounts(ctx context.Context) error {
c.pendingCountsWriteMutex.Lock()
defer c.pendingCountsWriteMutex.Unlock()
hits := c.pendingHits.Load()
misses := c.pendingMisses.Load()
upstreamFetches := c.pendingUpstreamFetches.Load()
upstreamFetchBytes := c.pendingUpstreamFetchBytes.Load()
transforms := c.pendingTransforms.Load()
if hits+misses+upstreamFetches+upstreamFetchBytes+transforms == 0 {
return nil
}
_, err := c.db.ExecContext(ctx, `
UPDATE cache_stats
SET hit_count = hit_count + ?,
miss_count = miss_count + ?,
upstream_fetch_count = upstream_fetch_count + ?,
upstream_fetch_bytes = upstream_fetch_bytes + ?,
transform_count = transform_count + ?,
last_updated_at = CURRENT_TIMESTAMP
WHERE id = 1
`, hits, misses, upstreamFetches, upstreamFetchBytes, transforms)
if err != nil {
return fmt.Errorf("failed to write pending counts: %w", err)
}
c.pendingHits.Add(-hits)
c.pendingMisses.Add(-misses)
c.pendingUpstreamFetches.Add(-upstreamFetches)
c.pendingUpstreamFetchBytes.Add(-upstreamFetchBytes)
c.pendingTransforms.Add(-transforms)
return nil
}
@@ -1,316 +0,0 @@
package imgcache
import (
"context"
"database/sql"
"errors"
"testing"
"time"
)
// pendingCountsWait is how long a test waits for the database's connection,
// for a count to reach the database, for a write to start waiting for the
// connection, or for StopPendingCountWrites to return.
const pendingCountsWait = 5 * time.Second
// holdDatabase takes the one connection of cache's database, so that every
// other query waits for it, and returns the func that frees it.
func holdDatabase(t *testing.T, cache *Cache) func() {
t.Helper()
ctx, cancel := context.WithTimeout(t.Context(), pendingCountsWait)
defer cancel()
conn, err := cache.db.Conn(ctx)
if err != nil {
t.Fatalf("failed to take the database connection: %v", err)
}
return func() { _ = conn.Close() }
}
// waitForConnectionWaits waits until db.Stats().WaitCount, the number of
// times a caller has waited for a connection, reaches waits.
func waitForConnectionWaits(t *testing.T, db *sql.DB, waits int64) {
t.Helper()
deadline := time.Now().Add(pendingCountsWait)
for db.Stats().WaitCount < waits {
if time.Now().After(deadline) {
t.Fatalf("callers waited for the database connection %d times, want %d",
db.Stats().WaitCount, waits)
}
time.Sleep(10 * time.Millisecond)
}
}
// waitForCounters waits until the cache_stats row holds want.
func waitForCounters(t *testing.T, cache *Cache, want cacheStatsCounters) {
t.Helper()
deadline := time.Now().Add(pendingCountsWait)
for {
got := readCacheStatsCounters(t, cache)
if got == want {
return
}
if time.Now().After(deadline) {
t.Fatalf("counters = %+v, want %+v", got, want)
}
time.Sleep(10 * time.Millisecond)
}
}
// 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 reach the database once the
// connection is free.
func TestService_Get_ReturnsByItsDeadlineWhileTheDatabaseIsBusy(t *testing.T) {
t.Parallel()
svc, fixtures, fetcher := setupHeldFetchService(t)
svc.cache.StartPendingCountWrites()
defer func() { _ = svc.cache.StopPendingCountWrites(t.Context()) }()
// Room for the request to reach its fetch on a busy host
const timeout = 2 * time.Second
ctx, cancel := context.WithTimeout(t.Context(), timeout)
defer cancel()
deadline, _ := ctx.Deadline()
results := startGet(ctx, svc, photoVariant(fixtures, 85, FitCover))
// The request has made its database reads by the time it fetches.
select {
case <-fetcher.started:
case <-time.After(timeout):
t.Fatal("request did not reach its fetch by its deadline")
}
releaseDatabase := holdDatabase(t, svc.cache)
defer releaseDatabase()
select {
case got := <-results:
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(time.Until(deadline) + time.Second):
t.Fatal("request did not return by its deadline while the database was busy")
}
releaseDatabase()
waitForCounters(t, svc.cache, cacheStatsCounters{missCount: 1})
}
// TestIncrementStats_ManyCountsPastTheirDeadlineLeaveOneWriteWaiting holds the
// database's one connection while many misses are counted past their
// deadline. At most one write may then wait for the connection, and every
// miss must reach the database once it is free.
func TestIncrementStats_ManyCountsPastTheirDeadlineLeaveOneWriteWaiting(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<30)
cache.StartPendingCountWrites()
defer func() { _ = cache.StopPendingCountWrites(t.Context()) }()
releaseDatabase := holdDatabase(t, cache)
defer releaseDatabase()
waitsBefore := cache.db.Stats().WaitCount
// Writes given a context past its deadline fail at once.
ended, cancel := context.WithDeadline(t.Context(), time.Now())
defer cancel()
const misses = 20
for range misses {
cache.IncrementStats(ended, false, 0)
}
// Once one write waits for the connection, any others start waiting too.
waitForConnectionWaits(t, cache.db, waitsBefore+1)
time.Sleep(arrivalWait)
waits := cache.db.Stats().WaitCount - waitsBefore
t.Logf("%d writes waited for the database", waits)
if waits > 1 {
t.Errorf("%d writes waited for the database, want at most 1", waits)
}
releaseDatabase()
waitForCounters(t, cache, cacheStatsCounters{missCount: misses})
}
// TestStats_IncludesPendingCounts counts hits and misses whose writes miss
// their deadline, and checks that Stats adds them to the database's counts
// before they are written.
func TestStats_IncludesPendingCounts(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<30)
_, err := cache.db.ExecContext(t.Context(),
`UPDATE cache_stats SET hit_count = 75, miss_count = 25 WHERE id = 1`)
if err != nil {
t.Fatal(err)
}
// Writes given a context past its deadline fail at once.
ended, cancel := context.WithDeadline(t.Context(), time.Now())
defer cancel()
cache.IncrementStats(ended, true, 0)
cache.IncrementStats(ended, true, 0)
cache.IncrementStats(ended, false, 0)
stats, err := cache.Stats(t.Context())
if err != nil {
t.Fatalf("Stats() error = %v", err)
}
if stats.HitCount != 77 || stats.MissCount != 26 {
t.Errorf("HitCount = %d, MissCount = %d, want 77 and 26",
stats.HitCount, stats.MissCount)
}
want := cacheStatsCounters{hitCount: 75, missCount: 25}
if got := readCacheStatsCounters(t, cache); got != want {
t.Errorf("counters in the database = %+v, want %+v", got, want)
}
}
// TestStopPendingCountWrites_WritesPendingCounts holds the database's one
// connection while the goroutine that writes the pending counts waits for it,
// counts more, then stops that goroutine. The stop must leave the write under
// way to finish rather than cut it off, and every count must reach the
// database once the connection is free.
func TestStopPendingCountWrites_WritesPendingCounts(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<30)
cache.StartPendingCountWrites()
defer func() { _ = cache.StopPendingCountWrites(t.Context()) }()
releaseDatabase := holdDatabase(t, cache)
defer releaseDatabase()
waitsBefore := cache.db.Stats().WaitCount
// Writes given a context past its deadline fail at once.
ended, cancel := context.WithDeadline(t.Context(), time.Now())
defer cancel()
cache.IncrementStats(ended, true, 0)
// The goroutine's write waits for the connection.
waitForConnectionWaits(t, cache.db, waitsBefore+1)
// Counted after that write read the pending counts
cache.IncrementStats(ended, false, 0)
cache.IncrementUpstreamFetch(ended, 1024)
cache.IncrementTransformCount(ended)
stopped := make(chan error, 1)
go func() {
stopped <- cache.StopPendingCountWrites(t.Context())
}()
// Had the stop cut off the write under way, its own write would wait for
// the connection too.
time.Sleep(arrivalWait)
waits := cache.db.Stats().WaitCount - waitsBefore
t.Logf("%d writes waited for the database", waits)
if waits > 1 {
t.Errorf("%d writes waited for the database, want 1: "+
"the stop cut off the write under way", waits)
}
releaseDatabase()
select {
case err := <-stopped:
if err != nil {
t.Fatalf("StopPendingCountWrites() error = %v", err)
}
case <-time.After(pendingCountsWait):
t.Fatal("StopPendingCountWrites() did not return once the database was free")
}
want := cacheStatsCounters{1, 1, 1, 1024, 1}
if got := readCacheStatsCounters(t, cache); got != want {
t.Errorf("counters = %+v, want %+v", got, want)
}
}
// TestStopPendingCountWrites_ReturnsWhenItsContextEnds holds the database's
// one connection while the goroutine that writes the pending counts waits for
// it, then stops that goroutine with a context that ends first. The stop must
// return the context's error, and the write under way must still reach the
// database once the connection is free.
func TestStopPendingCountWrites_ReturnsWhenItsContextEnds(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<30)
cache.StartPendingCountWrites()
defer func() { _ = cache.StopPendingCountWrites(t.Context()) }()
releaseDatabase := holdDatabase(t, cache)
defer releaseDatabase()
waitsBefore := cache.db.Stats().WaitCount
// Writes given a context past its deadline fail at once.
ended, cancel := context.WithDeadline(t.Context(), time.Now())
defer cancel()
cache.IncrementStats(ended, true, 0)
// The goroutine's write waits for the connection.
waitForConnectionWaits(t, cache.db, waitsBefore+1)
stopCtx, cancelStop := context.WithCancel(t.Context())
cancelStop()
stopped := make(chan error, 1)
go func() {
stopped <- cache.StopPendingCountWrites(stopCtx)
}()
select {
case err := <-stopped:
if !errors.Is(err, context.Canceled) {
t.Errorf("StopPendingCountWrites() error = %v, want %v",
err, context.Canceled)
}
case <-time.After(pendingCountsWait):
t.Fatal("StopPendingCountWrites() did not return once its context ended")
}
releaseDatabase()
waitForCounters(t, cache, cacheStatsCounters{hitCount: 1})
}
+69 -23
View File
@@ -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(ctx, true, 0)
s.writeCount(ctx, func(writeCtx context.Context) {
s.cache.IncrementStats(writeCtx, 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(ctx, false, 0)
s.writeCount(ctx, func(writeCtx context.Context) {
s.cache.IncrementStats(writeCtx, false, 0)
})
if err != nil {
return nil, err
@@ -208,6 +215,53 @@ 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
}
}
// writeCount runs write, which adds to a count in the database, in a
// goroutine of its own, and waits for it, though not past ctx's deadline.
// The write gets ctx without its cancellation or deadline, so an ended
// request is still counted. A request past its deadline thus does not wait
// for the database connection, which every request shares; its count is
// written after it returns.
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
}
select {
case <-written:
case <-time.After(time.Until(deadline)):
}
}
// ValidateRequest validates the request signature if required.
func (s *Service) ValidateRequest(req *ImageRequest) error {
// Check if host is allowed (no signature required)
@@ -300,8 +354,14 @@ func (s *Service) processOrWait(
}
}()
processingCtx, cancel := withoutCancelKeepingDeadline(ctx)
defer cancel()
processingCtx := context.WithoutCancel(ctx)
if deadline, ok := ctx.Deadline(); ok {
var cancel context.CancelFunc
processingCtx, cancel = context.WithDeadline(processingCtx, deadline)
defer cancel()
}
return s.processFromSourceOrFetch(processingCtx, req, cacheKey)
})
@@ -338,20 +398,6 @@ func (s *Service) processOrWait(
}, nil
}
// withoutCancelKeepingDeadline returns ctx without its cancellation but with
// its deadline, if it has one, and the func that releases the returned
// context.
func withoutCancelKeepingDeadline(
ctx context.Context,
) (context.Context, context.CancelFunc) {
deadline, hasDeadline := ctx.Deadline()
if !hasDeadline {
return context.WithoutCancel(ctx), func() {}
}
return context.WithDeadline(context.WithoutCancel(ctx), deadline)
}
// 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.
@@ -460,8 +506,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(ctx, fetchBytes)
s.writeCount(ctx, func(writeCtx context.Context) {
s.cache.IncrementUpstreamFetch(writeCtx, fetchBytes)
})
if err != nil {
return nil, fmt.Errorf("failed to read upstream response: %w", err)
@@ -542,8 +589,7 @@ func (s *Service) processAndStore(
processDuration := time.Since(processStart)
// Counted also when the request context has ended meanwhile
s.cache.IncrementTransformCount(ctx)
s.writeCount(ctx, s.cache.IncrementTransformCount)
// Read processed content
processedData, err := io.ReadAll(processResult.Content)
+20 -5
View File
@@ -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...)
}