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
This commit is contained in:
2026-10-04 05:42:52 +00:00
parent 5b17d1f555
commit c7f2c38733
5 changed files with 212 additions and 42 deletions
+50 -21
View File
@@ -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
}