Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
df7fff0ce1 |
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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})
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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...)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user