Stop cache eviction in progress at shutdown (closes #102) #173
@@ -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 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`,
|
||||
|
||||
@@ -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)
|
||||
},
|
||||
})
|
||||
|
||||
|
||||
@@ -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(),
|
||||
|
||||
+50
-21
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
在新工单中引用
屏蔽一个用户