From c7f2c3873350ec15220776d9c72d353f83dd200f Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Sun, 4 Oct 2026 03:53:52 +0000 Subject: [PATCH] Stop cache eviction in progress at shutdown (closes #102) StartEviction runs the eviction goroutine with its own context, which StopEviction cancels in place of the old stop channel, so a pass in progress stops at its next database call, file, row or eviction candidate instead of running to completion. StopEviction takes a context: when it ends before the goroutine exits, StopEviction stops waiting and returns an error wrapping it. The handlers' stop hook passes fx's stop context, so an eviction still running at fx's stop deadline fails the stop and the exit code is 1. The contextcheck suppression on the start hook stays, with a one-line reason: the loop outlives OnStart. Model: opus-5-5 --- TODO.md | 8 ++ internal/handlers/handlers.go | 17 +-- internal/imgcache/cache.go | 8 +- internal/imgcache/eviction.go | 71 ++++++--- internal/imgcache/eviction_internal_test.go | 150 +++++++++++++++++++- 5 files changed, 212 insertions(+), 42 deletions(-) diff --git a/TODO.md b/TODO.md index 404abd1..9636d71 100644 --- a/TODO.md +++ b/TODO.md @@ -29,6 +29,14 @@ 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; + `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 upstream host semaphores and variant `.meta` files no longer outlive their use (closes #87): the fetcher counts the fetches holding or waiting for a slot of each upstream host's semaphore and removes the host's diff --git a/internal/handlers/handlers.go b/internal/handlers/handlers.go index fafb27b..cecb5d7 100644 --- a/internal/handlers/handlers.go +++ b/internal/handlers/handlers.go @@ -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) }, }) diff --git a/internal/imgcache/cache.go b/internal/imgcache/cache.go index a2d2d61..5cd47ae 100644 --- a/internal/imgcache/cache.go +++ b/internal/imgcache/cache.go @@ -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(), diff --git a/internal/imgcache/eviction.go b/internal/imgcache/eviction.go index 9f90b8f..7ec9273 100644 --- a/internal/imgcache/eviction.go +++ b/internal/imgcache/eviction.go @@ -117,7 +117,8 @@ 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. Once ctx +// is cancelled, it stops at the next candidate and returns ctx's error. func (c *Cache) evictBatch(ctx context.Context, excessBytes int64) (int64, error) { candidates, err := c.evictionCandidates(ctx) if err != nil { @@ -131,6 +132,10 @@ func (c *Cache) evictBatch(ctx context.Context, excessBytes int64) (int64, error break } + if ctx.Err() != nil { + return freed, ctx.Err() + } + err := c.evictCandidate(ctx, candidate) if err != nil { c.log.Warn("failed to evict cache entry", @@ -424,37 +429,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. +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 +475,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 @@ -511,7 +523,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 +559,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 +651,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 +720,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 +785,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 } diff --git a/internal/imgcache/eviction_internal_test.go b/internal/imgcache/eviction_internal_test.go index f7084d3..68f0b10 100644 --- a/internal/imgcache/eviction_internal_test.go +++ b/internal/imgcache/eviction_internal_test.go @@ -4,6 +4,7 @@ import ( "bytes" "context" "database/sql" + "errors" "io/fs" "os" "path/filepath" @@ -657,7 +658,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 +691,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 +752,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 +809,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 +863,147 @@ 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. +func TestStopEvictionInterruptsPassInProgress(t *testing.T) { + t.Parallel() + + cache, _ := newEvictionTestCache(t, 1<<30) + + 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) + } +} + +// 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) + } + } +} + // TestEvictSourceBlobExcludesConcurrentStoreOfIdenticalContent exercises // the exact TOCTOU window between evictSourceBlob's row-deletion // transaction commit and its content file unlink: a concurrent