Stop cache eviction in progress at shutdown (closes #102) #173

Merged
clawbot merged 4 commits from issue-102-cancellable-eviction into next 2026-10-04 09:24:42 +02:00
5 changed files with 324 additions and 44 deletions
+9
View File
@@ -29,6 +29,15 @@ P2: security: referer blacklist
# Completed Steps
- 2026-10-04 shutdown stops cache eviction in progress (closes #102):
`StartEviction` runs the eviction goroutine with its own context, which
`StopEviction` cancels, so a pass in progress stops at its next database
call, file, row or eviction candidate instead of running to completion, and
no pass starts after it, so a stop logs at most one warning;
`StopEviction` takes a context and, when that context ends before the
goroutine exits, stops waiting and returns its error; the handlers' stop hook
passes fx's stop context, so an eviction still running when fx's stop
deadline ends fails the stop and makes the exit code 1.
- 2026-10-04 dead code in `internal/imgcache` is gone (closes #73): `Purge`,
which only returned an error and which nothing called, is no longer part of
the `ImageCache` interface or `Service`; the `SignatureValidator`,
+5 -12
View File
@@ -58,23 +58,16 @@ func New(lc fx.Lifecycle, params Params) (*Handlers, error) {
}
lc.Append(fx.Hook{
// The eviction goroutine must outlive OnStart, so it cannot
// inherit this hook's context. It makes its own instead, which
// leaves it uncancellable: an in-flight pass runs to completion
// during OnStop regardless of the shutdown deadline. Making the
// loop cancellable changes shutdown semantics and is tracked
// separately in issue #102, rather than being folded into the
// lint-conformance change that surfaced it.
//nolint:contextcheck // see issue #102
//nolint:contextcheck // the eviction loop outlives OnStart; OnStop cancels it
OnStart: func(_ context.Context) error {
return s.initImageService()
},
OnStop: func(_ context.Context) error {
if s.imgCache != nil {
s.imgCache.StopEviction()
OnStop: func(ctx context.Context) error {
if s.imgCache == nil {
return nil
}
return nil
return s.imgCache.StopEviction(ctx)
},
})
+3 -5
View File
@@ -11,7 +11,6 @@ import (
"io"
"log/slog"
"path/filepath"
"sync"
"time"
lru "github.com/hashicorp/golang-lru/v2"
@@ -69,11 +68,11 @@ type Cache struct {
// Eviction machinery. The channels are created in NewCache so
// stores can signal write pressure without racing StartEviction.
// evictionCancel, set by StartEviction, cancels the eviction
// goroutine's context.
evictionPressure chan struct{}
evictionStop chan struct{}
evictionDone chan struct{}
evictionStarted bool
evictionStopOnce sync.Once
evictionCancel context.CancelFunc
// metaCache holds the content types of the variants most recently
// stored or served, so a hit does not read the variant's .meta file.
@@ -112,7 +111,6 @@ func NewCache(db *sql.DB, config CacheConfig) (*Cache, error) {
log: log,
disabled: config.DisableDiskCache,
evictionPressure: make(chan struct{}, 1),
evictionStop: make(chan struct{}),
evictionDone: make(chan struct{}),
metaCache: metaCache,
contentLocks: newContentLock(),
+64 -23
View File
@@ -117,7 +117,10 @@ func (c *Cache) EvictToLimit(ctx context.Context) error {
// evictBatch fetches one batch of LRU candidates across variants and
// source blobs and evicts them oldest-first until excessBytes are
// freed or the batch is exhausted. It returns the bytes freed.
// freed or the batch is exhausted. It returns the bytes freed. A
// candidate that fails once ctx is cancelled (every one started after
// that fails at its first database call) ends the batch with ctx's
// error, without a warning.
func (c *Cache) evictBatch(ctx context.Context, excessBytes int64) (int64, error) {
candidates, err := c.evictionCandidates(ctx)
if err != nil {
@@ -133,6 +136,10 @@ func (c *Cache) evictBatch(ctx context.Context, excessBytes int64) (int64, error
err := c.evictCandidate(ctx, candidate)
if err != nil {
if ctx.Err() != nil {
return freed, ctx.Err()
}
c.log.Warn("failed to evict cache entry",
"cache_key", candidate.cacheKey,
"content_hash", candidate.contentHash,
@@ -424,37 +431,44 @@ func (c *Cache) notifyWritePressure() {
// startup and again on every periodic tick thereafter, and evicts to
// the configured limit on the given periodic interval and on
// write-pressure notifications. It is a no-op on a disabled cache or
// when already started.
// when already started. The goroutine outlives the caller, so it runs
// with its own context, which StopEviction cancels.
func (c *Cache) StartEviction(interval time.Duration) {
if c.disabled || c.evictionStarted {
if c.disabled || c.evictionCancel != nil {
return
}
c.evictionStarted = true
ctx, cancel := context.WithCancel(context.Background())
c.evictionCancel = cancel
go c.evictionLoop(interval)
go c.evictionLoop(ctx, interval)
}
// StopEviction stops the background eviction goroutine and waits for
// it to exit. It is safe to call when eviction was never started, and
// safe to call more than once.
func (c *Cache) StopEviction() {
if !c.evictionStarted {
return
// StopEviction cancels the background eviction goroutine, which
// interrupts a pass in progress, and waits for it to exit or for ctx to
// end, whichever comes first. In the second case it returns an error
// wrapping ctx's error. It is safe to call when eviction was never
// started, and safe to call more than once.
func (c *Cache) StopEviction(ctx context.Context) error {
if c.evictionCancel == nil {
return nil
}
c.evictionStopOnce.Do(func() {
close(c.evictionStop)
<-c.evictionDone
})
c.evictionCancel()
select {
case <-c.evictionDone:
return nil
case <-ctx.Done():
return fmt.Errorf("cache eviction still running: %w", ctx.Err())
}
}
// evictionLoop is the body of the background eviction goroutine.
func (c *Cache) evictionLoop(interval time.Duration) {
// evictionLoop is the body of the background eviction goroutine. It
// returns when ctx is cancelled, and starts no pass after that.
func (c *Cache) evictionLoop(ctx context.Context, interval time.Duration) {
defer close(c.evictionDone)
ctx := context.Background()
c.runReconciliationPass(ctx)
c.runEvictionPass(ctx)
@@ -463,7 +477,7 @@ func (c *Cache) evictionLoop(interval time.Duration) {
for {
select {
case <-c.evictionStop:
case <-ctx.Done():
return
case <-ticker.C:
// Reconciliation walks the cache directories, so it only
@@ -484,8 +498,13 @@ func (c *Cache) evictionLoop(interval time.Duration) {
}
// runEvictionPass runs one eviction pass, logging failures instead of
// propagating them (the loop must keep running).
// propagating them (the loop must keep running). It does nothing once
// ctx is cancelled.
func (c *Cache) runEvictionPass(ctx context.Context) {
if ctx.Err() != nil {
return
}
err := c.EvictToLimit(ctx)
if err != nil {
c.log.Warn("cache eviction pass failed", "error", err)
@@ -493,8 +512,13 @@ func (c *Cache) runEvictionPass(ctx context.Context) {
}
// runReconciliationPass runs one reconciliation pass, logging failures
// instead of propagating them (the loop must keep running).
// instead of propagating them (the loop must keep running). It does
// nothing once ctx is cancelled.
func (c *Cache) runReconciliationPass(ctx context.Context) {
if ctx.Err() != nil {
return
}
err := c.reconcileAccounting(ctx)
if err != nil {
c.log.Warn("cache accounting reconciliation failed", "error", err)
@@ -511,7 +535,8 @@ func (c *Cache) runReconciliationPass(ctx context.Context) {
// know (and rows whose files are gone), and sweeps stale temp files
// left behind by crashed writes. Running it periodically, not just
// once, bounds how long such drift can accumulate unaccounted for on a
// long-running process to one eviction interval.
// long-running process to one eviction interval. Once ctx is cancelled,
// it stops at the next file or row and returns ctx's error.
func (c *Cache) reconcileAccounting(ctx context.Context) error {
if c.disabled {
return nil
@@ -546,6 +571,10 @@ func (c *Cache) reconcileVariantFiles(ctx context.Context) error {
return filepath.WalkDir(
c.variants.baseDir,
func(path string, entry fs.DirEntry, err error) error {
if ctx.Err() != nil {
return ctx.Err()
}
if err != nil || entry.IsDir() {
return err
}
@@ -634,6 +663,10 @@ func (c *Cache) reconcileVariantRows(ctx context.Context) error {
}
for _, key := range keys {
if ctx.Err() != nil {
return ctx.Err()
}
if c.variants.Exists(key) {
continue
}
@@ -699,6 +732,10 @@ func (c *Cache) reconcileSourceFiles(ctx context.Context) error {
return filepath.WalkDir(
c.srcContent.baseDir,
func(path string, entry fs.DirEntry, err error) error {
if ctx.Err() != nil {
return ctx.Err()
}
if err != nil || entry.IsDir() {
return err
}
@@ -760,6 +797,10 @@ func (c *Cache) reconcileSourceRows(ctx context.Context) error {
}
for _, hash := range hashes {
if ctx.Err() != nil {
return ctx.Err()
}
if c.srcContent.Exists(hash) {
continue
}
+243 -4
View File
@@ -4,9 +4,12 @@ import (
"bytes"
"context"
"database/sql"
"errors"
"io/fs"
"log/slog"
"os"
"path/filepath"
"strings"
"testing"
"time"
@@ -657,7 +660,7 @@ func TestEvictionRunsUnderWritePressure(t *testing.T) {
// An interval far longer than the test ensures only write
// pressure can trigger eviction here.
cache.StartEviction(time.Hour)
defer cache.StopEviction()
defer func() { _ = cache.StopEviction(t.Context()) }()
keys := []VariantKey{
testVariantKeyOne, testVariantKeyTwo, testVariantKeyThree,
@@ -690,7 +693,7 @@ func TestEvictionRunsOnPeriodicSchedule(t *testing.T) {
// write-pressure notification fires and only the periodic ticker
// can trigger eviction.
cache.StartEviction(100 * time.Millisecond)
defer cache.StopEviction()
defer func() { _ = cache.StopEviction(t.Context()) }()
keys := []VariantKey{
testVariantKeyOne, testVariantKeyTwo, testVariantKeyThree,
@@ -751,7 +754,7 @@ func TestStartEvictionReconcilesAccountingWithDisk(t *testing.T) {
}
cache.StartEviction(time.Hour)
defer cache.StopEviction()
defer func() { _ = cache.StopEviction(t.Context()) }()
deadline := time.Now().Add(5 * time.Second)
@@ -808,7 +811,7 @@ func TestPeriodicReconciliationAdoptsFileThatAppearsAfterStartup(t *testing.T) {
const interval = 100 * time.Millisecond
cache.StartEviction(interval)
defer cache.StopEviction()
defer func() { _ = cache.StopEviction(t.Context()) }()
// Let startup reconciliation run and settle on an empty cache
// before introducing the untracked file, so the adoption we assert
@@ -862,6 +865,242 @@ func TestPeriodicReconciliationAdoptsFileThatAppearsAfterStartup(t *testing.T) {
}
}
// TestStopEvictionInterruptsPassInProgress holds the test database's
// only connection, so the startup reconciliation pass waits for it, and
// checks that StopEviction stops that pass instead of waiting for the
// connection to come free, and that the stop logs one warning: the
// interrupted reconciliation's, with no eviction pass started after it.
func TestStopEvictionInterruptsPassInProgress(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<30)
var logBuf bytes.Buffer
cache.log = slog.New(slog.NewJSONHandler(&logBuf, nil))
conn, err := cache.db.Conn(t.Context())
if err != nil {
t.Fatalf("failed to take the database connection: %v", err)
}
defer func() { _ = conn.Close() }()
cache.StartEviction(time.Hour)
// The pass is in progress once it waits for the connection.
deadline := time.Now().Add(5 * time.Second)
for cache.db.Stats().WaitCount == 0 {
if time.Now().After(deadline) {
t.Fatal("the reconciliation pass never waited for the database")
}
time.Sleep(10 * time.Millisecond)
}
ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second)
defer cancel()
err = cache.StopEviction(ctx)
t.Logf("StopEviction() error = %v", err)
if err != nil {
t.Fatalf("StopEviction() error = %v, want nil: the pass waiting for "+
"the database did not stop", err)
}
t.Logf("log output: %s", logBuf.String())
warnings := strings.Count(logBuf.String(), `"level":"WARN"`)
if warnings != 1 {
t.Errorf("the stop logged %d warnings, want 1", warnings)
}
}
// TestEvictToLimitStopsAtNextCandidateOnceCancelled cancels the context
// while the oldest of three source blobs is being evicted, and checks that
// EvictToLimit then returns context.Canceled without evicting the other
// two or logging a warning for either of them.
func TestEvictToLimitStopsAtNextCandidateOnceCancelled(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1)
var logBuf bytes.Buffer
cache.log = slog.New(slog.NewJSONHandler(&logBuf, nil))
hashes := []ContentHash{
storeEvictionTestSource(t, cache, "cancel.example.com", "/a.jpg",
bytes.Repeat([]byte{0x61}, 1000)),
storeEvictionTestSource(t, cache, "cancel.example.com", "/b.jpg",
bytes.Repeat([]byte{0x62}, 1000)),
storeEvictionTestSource(t, cache, "cancel.example.com", "/c.jpg",
bytes.Repeat([]byte{0x63}, 1000)),
}
base := time.Now().Add(-time.Hour)
for i, hash := range hashes {
setSourceLastAccessed(t, cache, hash, base.Add(time.Duration(i)*time.Minute))
}
ctx, cancel := context.WithCancel(t.Context())
defer cancel()
cache.evictSourceBlobTestHook = func(ContentHash) { cancel() }
err := cache.EvictToLimit(ctx)
t.Logf("EvictToLimit() error = %v", err)
t.Logf("log output: %s", logBuf.String())
if !errors.Is(err, context.Canceled) {
t.Errorf("EvictToLimit() error = %v, want context.Canceled", err)
}
if cache.srcContent.Exists(hashes[0]) {
t.Errorf("source blob %s, evicted when the context was cancelled, "+
"is still on disk", hashes[0])
}
for _, hash := range hashes[1:] {
if !cache.srcContent.Exists(hash) {
t.Errorf("source blob %s was evicted after the context was cancelled", hash)
}
}
if strings.Contains(logBuf.String(), `"level":"WARN"`) {
t.Errorf("EvictToLimit logged a warning after the context was cancelled")
}
assertNoDanglingReferences(t, cache)
}
// TestStopEvictionReturnsWhenItsContextEnds pauses an eviction pass where
// cancellation cannot reach it, after a source blob's rows are deleted and
// before its file is removed, and checks that StopEviction returns its
// context's error when that context ends instead of waiting for the pass.
// Once the pass goes on, the goroutine exits and no row points at a
// missing file.
func TestStopEvictionReturnsWhenItsContextEnds(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1)
paused := make(chan struct{})
resume := make(chan struct{})
cache.evictSourceBlobTestHook = func(ContentHash) {
close(paused)
<-resume
}
hash := storeEvictionTestSource(t, cache, "stop.example.com", "/a.jpg",
bytes.Repeat([]byte{0x61}, 1000))
cache.StartEviction(time.Hour)
select {
case <-paused:
case <-time.After(5 * time.Second):
t.Fatal("the eviction pass never reached the source blob")
}
ctx, cancel := context.WithTimeout(t.Context(), 50*time.Millisecond)
defer cancel()
err := cache.StopEviction(ctx)
t.Logf("StopEviction() error = %v", err)
if !errors.Is(err, context.DeadlineExceeded) {
t.Errorf("StopEviction() error = %v, want context.DeadlineExceeded", err)
}
close(resume)
err = cache.StopEviction(t.Context())
if err != nil {
t.Fatalf("second StopEviction() error = %v, want nil", err)
}
assertNoDanglingReferences(t, cache)
if cache.srcContent.Exists(hash) {
t.Errorf("source blob %s is still on disk after its rows were deleted", hash)
}
}
// TestReconciliationWalksStopOnceCancelled checks that both directory
// walks of a reconciliation pass return the context's error once it is
// cancelled, leaving in place a stale temp file they would otherwise
// remove.
func TestReconciliationWalksStopOnceCancelled(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<30)
ctx, cancel := context.WithCancel(t.Context())
cancel()
staleTime := time.Now().Add(-2 * staleTempFileAge)
walks := map[string]func(context.Context) error{
cache.variants.baseDir: cache.reconcileVariantFiles,
cache.srcContent.baseDir: cache.reconcileSourceFiles,
}
for dir, walk := range walks {
tempFile := filepath.Join(dir, tempFilePrefix+"stale")
err := os.WriteFile(tempFile, []byte("partial"), 0o600)
if err != nil {
t.Fatalf("failed to write temp file: %v", err)
}
err = os.Chtimes(tempFile, staleTime, staleTime)
if err != nil {
t.Fatalf("failed to backdate temp file: %v", err)
}
err = walk(ctx)
t.Logf("walk of %s: error = %v", dir, err)
if !errors.Is(err, context.Canceled) {
t.Errorf("walk of %s: error = %v, want context.Canceled", dir, err)
}
_, err = os.Stat(tempFile)
if err != nil {
t.Errorf("walk of %s went on after cancellation: %v", dir, err)
}
}
}
// TestReconciliationPassLogsNoWarningOnceCancelled checks that a
// reconciliation pass run with an already cancelled context logs no
// warning, so a periodic tick the loop takes after a stop adds no
// warning to the one from the pass the stop interrupted.
func TestReconciliationPassLogsNoWarningOnceCancelled(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<30)
var logBuf bytes.Buffer
cache.log = slog.New(slog.NewJSONHandler(&logBuf, nil))
ctx, cancel := context.WithCancel(t.Context())
cancel()
cache.runReconciliationPass(ctx)
t.Logf("log output: %s", logBuf.String())
if strings.Contains(logBuf.String(), `"level":"WARN"`) {
t.Errorf("runReconciliationPass logged a warning with a cancelled context")
}
}
// TestEvictSourceBlobExcludesConcurrentStoreOfIdenticalContent exercises
// the exact TOCTOU window between evictSourceBlob's row-deletion
// transaction commit and its content file unlink: a concurrent