Fix settle-window drainer cancellation in library ingest

This commit is contained in:
Quick
2026-05-27 13:37:33 -04:00
parent 886c492c6a
commit abdccd5904
2 changed files with 95 additions and 36 deletions
+13 -12
View File
@@ -165,7 +165,8 @@ func (e *Executor) ingest(ctx context.Context, folder *models.MediaFolder, mode
CurrentScope: claim.path,
})
runStartedAt := time.Now().UTC()
drainerCtx, stopDrainers := context.WithCancel(scanCtx)
drainerStopCtx, stopDrainers := context.WithCancel(context.Background())
defer stopDrainers()
drainerErrCh := make(chan error, 1)
var drainerWG sync.WaitGroup
if len(matchScopes) > 0 {
@@ -180,24 +181,21 @@ func (e *Executor) ingest(ctx context.Context, folder *models.MediaFolder, mode
ticker := time.NewTicker(scopedDrainInterval)
defer ticker.Stop()
for {
if drainerCtx.Err() != nil {
select {
case <-scanCtx.Done():
return
case <-drainerStopCtx.Done():
return
default:
}
started := time.Now()
processed, err := e.matcher.ProcessBatchByFolderAndPathPrefix(drainerCtx, folder.ID, scopePath, runStartedAt)
processed, err := e.matcher.ProcessBatchByFolderAndPathPrefix(scanCtx, folder.ID, scopePath, runStartedAt)
concurrentMatchDuration.Add(time.Since(started).Nanoseconds())
if processed > 0 {
concurrentMatched.Add(int64(processed))
}
if err != nil {
// A cancelled drainer context means this is a deliberate
// shutdown (the settle window ended and stopDrainers ran, or
// the whole scan is being torn down), not a real failure. The
// in-flight ProcessBatch call returns context.Canceled here, so
// exit cleanly without escalating. Genuine external
// cancellation still reaches the run via the scanCtx checks in
// the main goroutine, so this does not swallow real cancels.
if drainerCtx.Err() != nil {
if scanCtx.Err() != nil {
return
}
select {
@@ -208,7 +206,9 @@ func (e *Executor) ingest(ctx context.Context, folder *models.MediaFolder, mode
return
}
select {
case <-drainerCtx.Done():
case <-scanCtx.Done():
return
case <-drainerStopCtx.Done():
return
case <-ticker.C:
}
@@ -234,6 +234,7 @@ func (e *Executor) ingest(ctx context.Context, folder *models.MediaFolder, mode
scanStarted := time.Now()
scanMatchScopes, scanResult, err := e.scan(scanProgressCtx, folder, mode, claim.path)
if err != nil {
cancel()
stopDrainers()
drainerWG.Wait()
if len(matchScopes) > 0 {
+82 -24
View File
@@ -2,6 +2,7 @@ package libraryingest
import (
"context"
"sync"
"sync/atomic"
"testing"
"time"
@@ -10,25 +11,50 @@ import (
"github.com/Silo-Server/silo-server/internal/scanner"
)
// settleBlockingMatcher simulates a TV match drainer whose provider lookup is
// still in flight when the settle window expires: the batch call blocks until
// its context is cancelled, then returns context.Canceled — exactly what
// stopDrainers triggers when it cancels the drainer context.
type settleBlockingMatcher struct {
batchCalls atomic.Int64
// settleControlledMatcher simulates a TV match drainer whose provider lookup is
// still in flight when the settle window expires. The batch only returns once
// the test releases it, so the test can catch stopDrainers cancelling the active
// batch context instead of only stopping the drainer loop.
type settleControlledMatcher struct {
batchCalls atomic.Int64
batchCompleted atomic.Bool
batchCtxCanceled atomic.Bool
processAllBeforeBatch atomic.Bool
batchStarted chan struct{}
releaseBatch chan struct{}
closeBatchStartedOnce sync.Once
}
func (m *settleBlockingMatcher) ProcessBatchByFolderAndPathPrefix(ctx context.Context, _ int, _ string, _ time.Time) (int, error) {
func newSettleControlledMatcher() *settleControlledMatcher {
return &settleControlledMatcher{
batchStarted: make(chan struct{}),
releaseBatch: make(chan struct{}),
}
}
func (m *settleControlledMatcher) ProcessBatchByFolderAndPathPrefix(ctx context.Context, _ int, _ string, _ time.Time) (int, error) {
m.batchCalls.Add(1)
<-ctx.Done()
return 0, ctx.Err()
m.closeBatchStartedOnce.Do(func() {
close(m.batchStarted)
})
select {
case <-ctx.Done():
m.batchCtxCanceled.Store(true)
return 0, ctx.Err()
case <-m.releaseBatch:
m.batchCompleted.Store(true)
return 1, nil
}
}
func (m *settleBlockingMatcher) ProcessAllByFolderAndPathPrefix(context.Context, int, string, time.Time) (int, error) {
func (m *settleControlledMatcher) ProcessAllByFolderAndPathPrefix(context.Context, int, string, time.Time) (int, error) {
if !m.batchCompleted.Load() {
m.processAllBeforeBatch.Store(true)
}
return 0, nil
}
func (m *settleBlockingMatcher) RetryUnmatchedItemsByFolderAndPathPrefix(context.Context, int, string) (int, int, error) {
func (m *settleControlledMatcher) RetryUnmatchedItemsByFolderAndPathPrefix(context.Context, int, string) (int, int, error) {
return 0, 0, nil
}
@@ -54,29 +80,61 @@ func (s *settleStubScanner) FinalizeVariantsByPathPrefix(context.Context, *model
return nil
}
// TestIngestFolderCompletesWhenDrainerCanceledMidBatch is a regression test for
// the settle-window cancellation bug: a TV library full scan was recorded as
// "cancelled" because a drainer batch in flight when stopDrainers() ran returned
// context.Canceled, which the drainer escalated as a fatal scan error. The
// ingest must instead complete normally.
func TestIngestFolderCompletesWhenDrainerCanceledMidBatch(t *testing.T) {
matcher := &settleBlockingMatcher{}
// TestIngestFolderLetsActiveDrainerBatchFinishAfterSettleWindow is a regression
// test for the settle-window cancellation bug: a TV library full scan was
// recorded as "cancelled" because stopDrainers cancelled a drainer batch that
// was already in flight. The active batch must be allowed to finish; otherwise
// rows it already claimed can be excluded from the final scoped matcher by the
// runStartedAt attempt window.
func TestIngestFolderLetsActiveDrainerBatchFinishAfterSettleWindow(t *testing.T) {
const settleWindow = 25 * time.Millisecond
matcher := newSettleControlledMatcher()
exec := &Executor{
scanner: &settleStubScanner{result: &scanner.ScanResult{New: 1}},
matcher: matcher,
now: time.Now,
tvDrainSettleWindow: 50 * time.Millisecond,
tvDrainSettleWindow: settleWindow,
}
folder := &models.MediaFolder{ID: 5, Type: "series", Paths: []string{"/tv"}}
result, err := exec.IngestFolder(context.Background(), folder)
if err != nil {
t.Fatalf("expected ingest to complete, got error: %v", err)
type ingestResult struct {
result *Result
err error
}
if result == nil || result.Skipped {
t.Fatalf("expected a non-skipped result, got %+v", result)
done := make(chan ingestResult, 1)
go func() {
result, err := exec.IngestFolder(context.Background(), folder)
done <- ingestResult{result: result, err: err}
}()
select {
case <-matcher.batchStarted:
case <-time.After(time.Second):
t.Fatal("drainer never started a batch")
}
time.Sleep(2 * settleWindow)
close(matcher.releaseBatch)
var got ingestResult
select {
case got = <-done:
case <-time.After(time.Second):
t.Fatal("ingest did not complete after releasing the active drainer batch")
}
if got.err != nil {
t.Fatalf("expected ingest to complete, got error: %v", got.err)
}
if got.result == nil || got.result.Skipped {
t.Fatalf("expected a non-skipped result, got %+v", got.result)
}
if matcher.batchCalls.Load() == 0 {
t.Fatal("drainer never ran a batch; test did not exercise the settle-window shutdown path")
}
if matcher.batchCtxCanceled.Load() {
t.Fatal("settle-window shutdown cancelled the active drainer batch context")
}
if matcher.processAllBeforeBatch.Load() {
t.Fatal("final scoped matcher ran before the active drainer batch completed")
}
}