From aab0a804b97d77f6d32fef12f0aef3e1a52e0c1e Mon Sep 17 00:00:00 2001 From: Silo Server Developer Date: Wed, 27 May 2026 21:01:28 +0200 Subject: [PATCH] fix(matcher): tolerate concurrent-merge ErrItemNotFound in series episode-link ensure A scan drainer and the background MatchWorker can process the same folder concurrently. When a provider-ID merge moves a series' episodes to the survivor and deletes the source, an in-flight ensureSeriesEpisodeLinks(sourceID) hits catalog.ErrItemNotFound and was failing the whole scan. The episodes are already reattached, so this is benign: log and continue (matching the lenient call sites) instead of failing. Genuine errors still abort. --- internal/metadata/worker.go | 32 +++- internal/metadata/worker_test.go | 274 +++++++++++++++++++++++++++++++ 2 files changed, 300 insertions(+), 6 deletions(-) diff --git a/internal/metadata/worker.go b/internal/metadata/worker.go index 72785f9b..bec12554 100644 --- a/internal/metadata/worker.go +++ b/internal/metadata/worker.go @@ -2,6 +2,7 @@ package metadata import ( "context" + "errors" "fmt" "log/slog" "path/filepath" @@ -12,6 +13,7 @@ import ( "sync/atomic" "time" + "github.com/Silo-Server/silo-server/internal/catalog" "github.com/Silo-Server/silo-server/internal/models" ) @@ -837,10 +839,21 @@ func (w *MatchWorker) processSeriesRoot(ctx context.Context, job models.SeriesRo } } if err := w.service.ensureSeriesEpisodeLinks(ctx, representative.ContentID); err != nil { - if updateErr := w.seriesClaimer.UpdateError(ctx, job.MediaFolderID, job.ObservedRootPath, truncateSeriesQueueError(err.Error())); updateErr != nil { - return 0, updateErr + if errors.Is(err, catalog.ErrItemNotFound) { + // The series item was concurrently merged into another (provider-ID + // dedup moves its seasons+episodes to the survivor, then deletes the + // source row). The episodes are already reattached, so there is nothing + // to link here — benign; finish normally instead of failing the batch. + slog.Info("metadata: series item gone during episode-link ensure (likely concurrent merge); skipping", + "content_id", representative.ContentID, + "folder_id", job.MediaFolderID, + "observed_root_path", job.ObservedRootPath) + } else { + if updateErr := w.seriesClaimer.UpdateError(ctx, job.MediaFolderID, job.ObservedRootPath, truncateSeriesQueueError(err.Error())); updateErr != nil { + return 0, updateErr + } + return 0, fmt.Errorf("ensuring series episode links for %s: %w", representative.ContentID, err) } - return 0, fmt.Errorf("ensuring series episode links for %s: %w", representative.ContentID, err) } } if err := w.seriesClaimer.Delete(ctx, job.MediaFolderID, job.ObservedRootPath); err != nil { @@ -933,10 +946,17 @@ func (w *MatchWorker) processSeriesRoot(ctx context.Context, job models.SeriesRo } if strings.TrimSpace(finalContentID) != "" { if err := w.service.ensureSeriesEpisodeLinks(ctx, finalContentID); err != nil { - if updateErr := w.seriesClaimer.UpdateError(ctx, job.MediaFolderID, job.ObservedRootPath, truncateSeriesQueueError(err.Error())); updateErr != nil { - return 0, updateErr + if errors.Is(err, catalog.ErrItemNotFound) { + slog.Info("metadata: series item gone during episode-link ensure (likely concurrent merge); skipping", + "content_id", finalContentID, + "folder_id", job.MediaFolderID, + "observed_root_path", job.ObservedRootPath) + } else { + if updateErr := w.seriesClaimer.UpdateError(ctx, job.MediaFolderID, job.ObservedRootPath, truncateSeriesQueueError(err.Error())); updateErr != nil { + return 0, updateErr + } + return 0, fmt.Errorf("ensuring series episode links for %s: %w", finalContentID, err) } - return 0, fmt.Errorf("ensuring series episode links for %s: %w", finalContentID, err) } if _, ok := w.service.confirmedOwnershipItem(ctx, finalContentID); ok { w.service.claimConfirmedSeriesRootOwnership(ctx, job.MediaFolderID, job.ObservedRootPath, finalContentID, groupFiles) diff --git a/internal/metadata/worker_test.go b/internal/metadata/worker_test.go index 867d3f44..0ecb4d15 100644 --- a/internal/metadata/worker_test.go +++ b/internal/metadata/worker_test.go @@ -8,6 +8,7 @@ import ( "testing" "time" + "github.com/Silo-Server/silo-server/internal/catalog" "github.com/Silo-Server/silo-server/internal/models" ) @@ -1200,3 +1201,276 @@ func TestWorkerProcessBatchByFolderAndPathPrefix_MovieQueueClaimsOnlyOncePerScan t.Fatal("expected queue error to be recorded") } } + +// --------------------------------------------------------------------------- +// Tests for concurrent-merge ErrItemNotFound tolerance (hotfix 2026-05-27) +// --------------------------------------------------------------------------- + +// TestProcessSeriesRoot_Site1_ErrItemNotFoundIsToleratedNotFailed verifies that +// when ensureSeriesEpisodeLinks returns catalog.ErrItemNotFound at Site 1 +// (the "all files already linked" fast path), the batch succeeds instead of +// failing. This covers the concurrent-merge case where the source series row +// was deleted after its episodes were moved to the survivor. +func TestProcessSeriesRoot_Site1_ErrItemNotFoundIsToleratedNotFailed(t *testing.T) { + h := newTestHarness() + ctx := context.Background() + + // Need a series-type folder so queueUsageForFolder enables the series queue. + h.service.folderRepo = &fakeWorkerFolderRepo{ + folders: map[int]*models.MediaFolder{ + 10: {ID: 10, Type: "series", Enabled: true}, + }, + } + + // Pre-populate the item as "matched" so reusableQueuedMovieSkeleton + // returns false (matched is not skeleton-like), bypassing the re-process + // block and falling through directly to ensureSeriesEpisodeLinks. + const contentID = "series-gone-after-merge" + if err := h.itemRepo.Upsert(ctx, &models.MediaItem{ + ContentID: contentID, + Status: "matched", + Title: "Example Show", + Type: "series", + Studios: []string{}, + Networks: []string{}, + Countries: []string{}, + Genres: []string{}, + }); err != nil { + t.Fatalf("upsert item: %v", err) + } + + // All files are already linked — hasUnlinkedGroupFile returns false. + file := &models.MediaFile{ + ID: 1, + MediaFolderID: 10, + FilePath: "/media/shows/Example Show/Season 01/Example.Show.S01E01.mkv", + ObservedRootPath: "/media/shows/Example Show", + GroupKeyVersion: 1, + ContentGroupKey: "v1|series|example_show|2024", + ContentID: contentID, + } + h.fileRepo.setGroupFiles(10, 1, "v1|series|example_show|2024", file) + h.fileRepo.contentIDs[file.ID] = contentID + + // Hook ensureSeriesEpisodeLinks to simulate the source row being gone. + h.service.hooks.ensureSeriesEpisodeLinks = func(_ context.Context, _ string) error { + return fmt.Errorf("loading series item: %w", catalog.ErrItemNotFound) + } + + queueRepo := newFakeSeriesQueueRepo(models.SeriesRootMatchJob{ + MediaFolderID: 10, + ObservedRootPath: "/media/shows/Example Show", + SampleFilePath: file.FilePath, + ObservedFileCount: 1, + }) + + worker := NewMatchWorker(h.service, h.fileRepo, 1, 10, 0) + worker.SetSeriesRootClaimer(queueRepo, true) + processed, err := worker.ProcessAllByFolderAndPathPrefix(ctx, 10, "/media/shows/Example Show", time.Time{}) + if err != nil { + t.Fatalf("expected no error for ErrItemNotFound (benign concurrent merge), got: %v", err) + } + if processed != 1 { + t.Fatalf("processed = %d, want 1", processed) + } + // Queue row must be deleted on normal completion. + if _, ok := queueRepo.deleted["10:/media/shows/Example Show"]; !ok { + t.Fatal("expected queue row to be deleted on normal completion") + } + // No error must have been recorded in the queue. + if queueRepo.errors["10:/media/shows/Example Show"] != "" { + t.Fatalf("expected no queue error, got %q", queueRepo.errors["10:/media/shows/Example Show"]) + } +} + +// TestProcessSeriesRoot_Site1_NonNotFoundErrorStillFails verifies that a +// genuine (non-ErrItemNotFound) error from ensureSeriesEpisodeLinks at Site 1 +// still fails the batch — the fix must not swallow real errors. +func TestProcessSeriesRoot_Site1_NonNotFoundErrorStillFails(t *testing.T) { + h := newTestHarness() + ctx := context.Background() + + h.service.folderRepo = &fakeWorkerFolderRepo{ + folders: map[int]*models.MediaFolder{ + 10: {ID: 10, Type: "series", Enabled: true}, + }, + } + + const contentID = "series-link-failure" + if err := h.itemRepo.Upsert(ctx, &models.MediaItem{ + ContentID: contentID, + Status: "matched", + Title: "Example Show", + Type: "series", + Studios: []string{}, + Networks: []string{}, + Countries: []string{}, + Genres: []string{}, + }); err != nil { + t.Fatalf("upsert item: %v", err) + } + + file := &models.MediaFile{ + ID: 1, + MediaFolderID: 10, + FilePath: "/media/shows/Example Show/Season 01/Example.Show.S01E01.mkv", + ObservedRootPath: "/media/shows/Example Show", + GroupKeyVersion: 1, + ContentGroupKey: "v1|series|example_show|2024", + ContentID: contentID, + } + h.fileRepo.setGroupFiles(10, 1, "v1|series|example_show|2024", file) + h.fileRepo.contentIDs[file.ID] = contentID + + // Return a genuine (non-not-found) error. + linkErr := errors.New("database connection reset") + h.service.hooks.ensureSeriesEpisodeLinks = func(_ context.Context, _ string) error { + return linkErr + } + + queueRepo := newFakeSeriesQueueRepo(models.SeriesRootMatchJob{ + MediaFolderID: 10, + ObservedRootPath: "/media/shows/Example Show", + SampleFilePath: file.FilePath, + ObservedFileCount: 1, + }) + + worker := NewMatchWorker(h.service, h.fileRepo, 1, 10, 0) + worker.SetSeriesRootClaimer(queueRepo, true) + _, err := worker.ProcessAllByFolderAndPathPrefix(ctx, 10, "/media/shows/Example Show", time.Time{}) + if err == nil { + t.Fatal("expected error for genuine link failure, got nil") + } + if !strings.Contains(err.Error(), "ensuring series episode links") { + t.Fatalf("unexpected error message: %v", err) + } +} + +// TestProcessSeriesRoot_Site2_ErrItemNotFoundIsToleratedNotFailed verifies that +// when ensureSeriesEpisodeLinks returns catalog.ErrItemNotFound at Site 2 +// (the new-skeleton / enrichment path), the batch succeeds. +func TestProcessSeriesRoot_Site2_ErrItemNotFoundIsToleratedNotFailed(t *testing.T) { + h := newTestHarness() + ctx := context.Background() + + h.service.folderRepo = &fakeWorkerFolderRepo{ + folders: map[int]*models.MediaFolder{ + 10: {ID: 10, Type: "series", Enabled: true}, + }, + } + + // File has no ContentID yet → hasUnlinkedGroupFile returns true → + // proceeds through createOrFindSkeleton + enrichment path (Site 2). + file := &models.MediaFile{ + ID: 1, + MediaFolderID: 10, + FilePath: "/media/shows/Example Show/Season 01/Example.Show.S01E01.mkv", + ObservedRootPath: "/media/shows/Example Show", + GroupKeyVersion: 1, + ContentGroupKey: "v1|series|example_show|2024", + BaseTitle: "Example Show", + BaseType: "series", + } + h.fileRepo.setGroupFiles(10, 1, "v1|series|example_show|2024", file) + + const skeletonID = "skeleton-series-id" + + h.service.hooks.createOrFindSkeleton = func(_ context.Context, _ *models.MediaFile, _ int) (*skeletonResult, error) { + return &skeletonResult{ + ContentID: skeletonID, + IsNew: true, + ItemStatus: "pending", + Type: "series", + }, nil + } + h.service.hooks.process = func(_ context.Context, _ ProcessRequest) (*ProcessResult, error) { + return &ProcessResult{Updated: true}, nil + } + // Simulate the concurrent-merge: ensureSeriesEpisodeLinks finds the source gone. + h.service.hooks.ensureSeriesEpisodeLinks = func(_ context.Context, _ string) error { + return fmt.Errorf("loading series item: %w", catalog.ErrItemNotFound) + } + + queueRepo := newFakeSeriesQueueRepo(models.SeriesRootMatchJob{ + MediaFolderID: 10, + ObservedRootPath: "/media/shows/Example Show", + SampleFilePath: file.FilePath, + ObservedFileCount: 1, + }) + + worker := NewMatchWorker(h.service, h.fileRepo, 1, 10, 0) + worker.SetSeriesRootClaimer(queueRepo, true) + processed, err := worker.ProcessAllByFolderAndPathPrefix(ctx, 10, "/media/shows/Example Show", time.Time{}) + if err != nil { + t.Fatalf("expected no error for ErrItemNotFound (benign concurrent merge), got: %v", err) + } + if processed != 1 { + t.Fatalf("processed = %d, want 1", processed) + } + if _, ok := queueRepo.deleted["10:/media/shows/Example Show"]; !ok { + t.Fatal("expected queue row to be deleted on normal completion") + } + if queueRepo.errors["10:/media/shows/Example Show"] != "" { + t.Fatalf("expected no queue error, got %q", queueRepo.errors["10:/media/shows/Example Show"]) + } +} + +// TestProcessSeriesRoot_Site2_NonNotFoundErrorStillFails verifies that a genuine +// error from ensureSeriesEpisodeLinks at Site 2 still fails the batch. +func TestProcessSeriesRoot_Site2_NonNotFoundErrorStillFails(t *testing.T) { + h := newTestHarness() + ctx := context.Background() + + h.service.folderRepo = &fakeWorkerFolderRepo{ + folders: map[int]*models.MediaFolder{ + 10: {ID: 10, Type: "series", Enabled: true}, + }, + } + + file := &models.MediaFile{ + ID: 1, + MediaFolderID: 10, + FilePath: "/media/shows/Example Show/Season 01/Example.Show.S01E01.mkv", + ObservedRootPath: "/media/shows/Example Show", + GroupKeyVersion: 1, + ContentGroupKey: "v1|series|example_show|2024", + BaseTitle: "Example Show", + BaseType: "series", + } + h.fileRepo.setGroupFiles(10, 1, "v1|series|example_show|2024", file) + + const skeletonID = "skeleton-series-id-2" + + h.service.hooks.createOrFindSkeleton = func(_ context.Context, _ *models.MediaFile, _ int) (*skeletonResult, error) { + return &skeletonResult{ + ContentID: skeletonID, + IsNew: true, + ItemStatus: "pending", + Type: "series", + }, nil + } + h.service.hooks.process = func(_ context.Context, _ ProcessRequest) (*ProcessResult, error) { + return &ProcessResult{Updated: true}, nil + } + linkErr := errors.New("db timeout") + h.service.hooks.ensureSeriesEpisodeLinks = func(_ context.Context, _ string) error { + return linkErr + } + + queueRepo := newFakeSeriesQueueRepo(models.SeriesRootMatchJob{ + MediaFolderID: 10, + ObservedRootPath: "/media/shows/Example Show", + SampleFilePath: file.FilePath, + ObservedFileCount: 1, + }) + + worker := NewMatchWorker(h.service, h.fileRepo, 1, 10, 0) + worker.SetSeriesRootClaimer(queueRepo, true) + _, err := worker.ProcessAllByFolderAndPathPrefix(ctx, 10, "/media/shows/Example Show", time.Time{}) + if err == nil { + t.Fatal("expected error for genuine link failure at site 2, got nil") + } + if !strings.Contains(err.Error(), "ensuring series episode links") { + t.Fatalf("unexpected error message: %v", err) + } +}