Files
silo-server/internal/metadata/worker.go
22fec4ed2d feat(metadata): add resilient match queue diagnostics (#463)
* feat(metadata): add resilient match queue diagnostics

* fix(metadata): harden match queue lifecycle

---------

Co-authored-by: Quick104 <31828688+Quick104@users.noreply.github.com>
2026-07-24 12:18:49 -04:00

1669 lines
55 KiB
Go

package metadata
import (
"context"
"errors"
"fmt"
"log/slog"
"path/filepath"
"slices"
"strconv"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/Silo-Server/silo-server/internal/catalog"
"github.com/Silo-Server/silo-server/internal/librarykind"
"github.com/Silo-Server/silo-server/internal/models"
"github.com/Silo-Server/silo-server/internal/naming"
"github.com/Silo-Server/silo-server/internal/notifications"
)
// MatchWorker processes unmatched files in the background.
type MatchWorker struct {
service *MetadataService
fileLister UnmatchedFileLister
itemLister UnmatchedItemLister
movieClaimer MovieFileClaimer
seriesClaimer SeriesRootClaimer
enableTVSeriesRootQueue bool
realtimeHub *notifications.Hub
// workers/batchSize are atomics so admin settings changes can resize the
// pool while the long-running match loop reads them each cycle.
workers atomic.Int32
batchSize atomic.Int32
interval time.Duration
}
// SetConcurrency updates the worker pool size and claim batch size. Safe for
// concurrent use; the next match cycle picks the new values up. Out-of-range
// values are ignored.
func (w *MatchWorker) SetConcurrency(workers, batchSize int) {
if w == nil {
return
}
if workers >= 1 {
w.workers.Store(int32(workers))
}
if batchSize > 0 {
w.batchSize.Store(int32(batchSize))
}
}
func (w *MatchWorker) workerCount() int { return int(w.workers.Load()) }
func (w *MatchWorker) claimBatchSize() int { return int(w.batchSize.Load()) }
func (w *MatchWorker) queueClaimSize() int {
return min(w.workerCount(), w.claimBatchSize())
}
type NonSeriesFileClaimer interface {
ClaimUnmatchedNonSeries(ctx context.Context, limit int) ([]*models.MediaFile, error)
ClaimUnmatchedNonSeriesByFolderAndPathPrefix(ctx context.Context, folderID int, pathPrefix string, limit int, attemptBefore time.Time) ([]*models.MediaFile, error)
}
type MixedFileClaimer interface {
ClaimUnmatchedMixed(ctx context.Context, limit int) ([]*models.MediaFile, error)
ClaimUnmatchedMixedByFolderAndPathPrefix(ctx context.Context, folderID int, pathPrefix string, limit int, attemptBefore time.Time) ([]*models.MediaFile, error)
}
type MatchSuppressionChecker interface {
IsMatchSuppressed(ctx context.Context, fileID int) (bool, error)
}
type MovieFileClaimer interface {
Claim(ctx context.Context, limit int) ([]models.MovieMatchJob, error)
ClaimByFolderAndPathPrefix(ctx context.Context, folderID int, pathPrefix string, limit int, attemptBefore time.Time) ([]models.MovieMatchJob, error)
Delete(ctx context.Context, mediaFileID int, leaseToken string) error
UpdateError(ctx context.Context, mediaFileID int, leaseToken, errText string) error
UpdateFailure(ctx context.Context, mediaFileID int, leaseToken string, failure MatchFailure) error
ReleaseLease(ctx context.Context, leaseToken string) (int, error)
}
type SeriesRootClaimer interface {
Claim(ctx context.Context, limit int) ([]models.SeriesRootMatchJob, error)
ClaimByFolderAndPathPrefix(ctx context.Context, folderID int, pathPrefix string, limit int, attemptBefore time.Time) ([]models.SeriesRootMatchJob, error)
Delete(ctx context.Context, folderID int, observedRootPath, leaseToken string) error
UpdateError(ctx context.Context, folderID int, observedRootPath, leaseToken, errText string) error
UpdateFailure(ctx context.Context, folderID int, observedRootPath, leaseToken string, failure MatchFailure) error
ReleaseLease(ctx context.Context, leaseToken string) (int, error)
ListByFolder(ctx context.Context, folderID int, limit int, offset int) ([]models.SeriesRootMatchQueueEntry, int, error)
CountByFolder(ctx context.Context, folderID int) (int, error)
}
// NewMatchWorker creates a new background match worker.
func NewMatchWorker(service *MetadataService, fileLister UnmatchedFileLister, workers, batchSize int, interval time.Duration) *MatchWorker {
if workers < 1 {
workers = 8
}
if batchSize <= 0 {
batchSize = 500
}
if interval == 0 {
interval = 30 * time.Second
}
var itemLister UnmatchedItemLister
if service != nil {
itemLister = service.itemRepo
}
w := &MatchWorker{
service: service,
fileLister: fileLister,
itemLister: itemLister,
interval: interval,
}
w.SetConcurrency(workers, batchSize)
return w
}
// SetSeriesRootClaimer enables native TV root-backed matching when enabled is true.
func (w *MatchWorker) SetSeriesRootClaimer(claimer SeriesRootClaimer, enabled bool) {
if w == nil {
return
}
w.seriesClaimer = claimer
w.enableTVSeriesRootQueue = enabled
}
// SetMovieFileClaimer enables queue-backed movie matching when claimer is non-nil.
func (w *MatchWorker) SetMovieFileClaimer(claimer MovieFileClaimer) {
if w == nil {
return
}
w.movieClaimer = claimer
}
// SetRealtimeHub publishes catalog item changes as metadata enrichment commits.
func (w *MatchWorker) SetRealtimeHub(hub *notifications.Hub) {
if w == nil {
return
}
w.realtimeHub = hub
}
// Run starts the match worker loop. It blocks until ctx is cancelled.
func (w *MatchWorker) Run(ctx context.Context) {
ticker := time.NewTicker(w.interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
w.processUnmatched(ctx)
}
}
}
// processUnmatched fetches a batch of unmatched files and processes them.
func (w *MatchWorker) processUnmatched(ctx context.Context) {
if w.enableTVSeriesRootQueue && w.seriesClaimer != nil {
w.processBackgroundSeriesQueue(ctx)
}
if w.movieClaimer != nil {
w.processBackgroundMovieQueue(ctx)
}
files, err := w.claimBackgroundFiles(ctx)
if err != nil {
slog.ErrorContext(ctx, "metadata: failed to list unmatched files", "component", "metadata", "error", err)
return
}
if len(files) == 0 {
return
}
slog.InfoContext(ctx, "metadata: processing unmatched files", "component", "metadata", "count", len(files))
w.processFiles(ctx, files)
}
func (w *MatchWorker) processBackgroundSeriesQueue(ctx context.Context) {
remaining := w.claimBatchSize()
for remaining > 0 && ctx.Err() == nil {
claimLimit := min(w.queueClaimSize(), remaining)
jobs, err := w.seriesClaimer.Claim(ctx, claimLimit)
if err != nil {
slog.ErrorContext(ctx, "metadata: failed to claim unmatched series roots", "component", "metadata", "error", err)
return
}
if len(jobs) == 0 {
return
}
slog.InfoContext(ctx, "metadata: processing unmatched series roots", "component", "metadata", "count", len(jobs))
if _, err := w.processSeriesRoots(ctx, jobs); err != nil {
slog.ErrorContext(ctx, "metadata: failed to process unmatched series roots", "component", "metadata", "error", err)
return
}
remaining -= len(jobs)
if len(jobs) < claimLimit {
return
}
}
}
func (w *MatchWorker) processBackgroundMovieQueue(ctx context.Context) {
remaining := w.claimBatchSize()
for remaining > 0 && ctx.Err() == nil {
claimLimit := min(w.queueClaimSize(), remaining)
jobs, err := w.movieClaimer.Claim(ctx, claimLimit)
if err != nil {
slog.ErrorContext(ctx, "metadata: failed to claim queued movie files", "component", "metadata", "error", err)
return
}
if len(jobs) == 0 {
return
}
slog.InfoContext(ctx, "metadata: processing queued movie files", "component", "metadata", "count", len(jobs))
w.processQueuedMovieFiles(ctx, jobs)
remaining -= len(jobs)
if len(jobs) < claimLimit {
return
}
}
}
func (w *MatchWorker) processFile(ctx context.Context, file *models.MediaFile) {
if w.fileLister != nil {
defer func() {
stampCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second)
defer cancel()
if err := w.fileLister.MarkMatchAttempted(stampCtx, file.ID); err != nil {
slog.WarnContext(ctx, "metadata: failed to record match attempt", "component", "metadata",
"file_id", file.ID,
"path", file.FilePath,
"error", err)
}
}()
}
w.processFileWithFolderCache(ctx, file, nil, nil)
}
func (w *MatchWorker) processFileWithFolderCache(ctx context.Context, file *models.MediaFile, folderEnabledCache *sync.Map, deferredSeriesLinks *sync.Map) {
if !w.folderEnabled(ctx, file.MediaFolderID, folderEnabledCache) {
slog.InfoContext(ctx, "metadata: skipping file in disabled library", "component", "metadata",
"file_id", file.ID,
"path", file.FilePath,
"folder_id", file.MediaFolderID,
)
return
}
// Phase 0: Create skeleton item or find existing.
skeleton, err := w.service.createOrFindSkeleton(ctx, file, file.MediaFolderID)
if err != nil {
slog.WarnContext(ctx, "metadata: skeleton creation failed", "component", "metadata",
"file_id", file.ID, "path", file.FilePath, "error", err)
return
}
// If the file was linked to an existing item via dedup, skip enrichment.
// The first file for each series creates the skeleton (IsNew=true) and runs
// the full provider pipeline. Subsequent files only need linking.
if !skeleton.IsNew {
// For series files, defer episode linking to a single call per series
// after the batch completes, rather than calling per-file.
if skeleton.Type == "series" && deferredSeriesLinks != nil {
deferredSeriesLinks.Store(skeleton.ContentID, struct{}{})
} else if skeleton.Type == "series" {
// Fallback for callers that don't provide a deferred map.
if err := w.service.ensureSeriesEpisodeLinks(ctx, skeleton.ContentID); err != nil {
slog.WarnContext(ctx, "metadata: failed to ensure series episode links", "component", "metadata",
"content_id", skeleton.ContentID,
"file_id", file.ID,
"path", file.FilePath,
"error", err)
}
}
return
}
if skeleton.ItemStatus == "ambiguous" {
return
}
req := w.buildProcessRequestForGroup(ctx, file, skeleton, nil)
result, err := w.service.Process(ctx, req)
if err != nil {
slog.WarnContext(ctx, "metadata: enrichment failed", "component", "metadata",
"file_id", file.ID, "path", file.FilePath, "error", err)
w.logStatusUpdateFailure(ctx, skeleton.ContentID, "unmatched", "content_id", skeleton.ContentID, "file_id", file.ID, "path", file.FilePath)
// For series items, synthesize fallback episode structure so episodes
// are visible even when no provider match was found.
if skeleton.Type == "series" {
if fbErr := w.service.SynthesizeFallbackEpisodes(ctx, skeleton.ContentID); fbErr != nil {
slog.WarnContext(ctx, "metadata: fallback episode synthesis failed after enrichment error", "component", "metadata",
"content_id", skeleton.ContentID, "error", fbErr)
}
}
return
}
if result != nil && !result.Updated {
w.logStatusUpdateFailure(ctx, skeleton.ContentID, "unmatched", "content_id", skeleton.ContentID, "file_id", file.ID, "path", file.FilePath)
// Same fallback synthesis for series when no provider data was returned.
if skeleton.Type == "series" {
if fbErr := w.service.SynthesizeFallbackEpisodes(ctx, skeleton.ContentID); fbErr != nil {
slog.WarnContext(ctx, "metadata: fallback episode synthesis failed for unmatched series", "component", "metadata",
"content_id", skeleton.ContentID, "error", fbErr)
}
}
return
}
w.publishCatalogItemChanged(ctx, file.MediaFolderID, resultContentID(result, skeleton.ContentID), "metadata_updated")
}
func (w *MatchWorker) buildProcessRequestForGroup(ctx context.Context, representative *models.MediaFile, skeleton *skeletonResult, preloadedGroupFiles []*models.MediaFile) ProcessRequest {
groupFiles := preloadedGroupFiles
if len(groupFiles) == 0 {
groupFiles = []*models.MediaFile{representative}
if skeleton != nil && skeleton.Type == "series" && w.service != nil && w.service.fileRepo != nil {
loadedFiles, err := w.service.fileRepo.ListByObservedRootPath(ctx, representative.MediaFolderID, skeleton.ObservedRootPath)
if err != nil {
slog.WarnContext(ctx, "metadata: failed to load observed-root files", "component", "metadata",
"folder_id", representative.MediaFolderID,
"observed_root_path", skeleton.ObservedRootPath,
"error", err,
)
} else if len(loadedFiles) > 0 {
groupFiles = loadedFiles
}
}
}
groupFilePaths := make([]string, 0, len(groupFiles))
groupFilesForSidecars := make([]*models.MediaFile, 0, len(groupFiles))
for _, groupFile := range groupFiles {
if groupFile == nil {
continue
}
groupFilePaths = append(groupFilePaths, groupFile.FilePath)
groupFilesForSidecars = append(groupFilesForSidecars, groupFile)
}
if len(groupFilesForSidecars) == 0 {
groupFilePaths = []string{representative.FilePath}
groupFilesForSidecars = []*models.MediaFile{representative}
}
slices.Sort(groupFilePaths)
sidecarSearchPaths := w.service.directorySidecarSearchPathsForFiles(ctx, groupFilesForSidecars)
safeObservedRootPath := ""
if w.service.canUseObservedRootForDirectorySidecars(
ctx,
representative.MediaFolderID,
skeleton.ObservedRootPath,
skeleton.GroupKeyVersion,
skeleton.ContentGroupKey,
) {
safeObservedRootPath = skeleton.ObservedRootPath
}
hints := &MatchHints{
FileHash: representative.FileHash,
FilePath: representative.FilePath,
RepresentativeFilePath: representative.FilePath,
ObservedRootPath: safeObservedRootPath,
AllGroupFilePaths: groupFilePaths,
PrimarySidecarSearchPaths: sidecarSearchPaths,
Title: skeleton.Title,
Year: skeleton.Year,
Type: skeleton.Type,
TmdbID: skeleton.TmdbID,
ImdbID: skeleton.ImdbID,
TvdbID: skeleton.TvdbID,
HintSource: "scanner",
AlternateIdentities: queuedMatchIdentityAlternates(representative, skeleton),
}
return ProcessRequest{
ContentID: skeleton.ContentID,
Hints: hints,
FolderID: formatFolderID(representative.MediaFolderID),
Mode: ModeInitialMatch,
}
}
const maxAlternateMatchIdentities = 3
func queuedMatchIdentityAlternates(file *models.MediaFile, skeleton *skeletonResult) []MatchIdentityHint {
if file == nil || skeleton == nil {
return nil
}
seen := map[string]struct{}{
matchIdentityKey(skeleton.Title, skeleton.Year): {},
}
out := make([]MatchIdentityHint, 0, maxAlternateMatchIdentities)
add := func(title string, year int, source string) {
title = strings.TrimSpace(naming.StripComparisonSafeEditionSuffix(title))
key := matchIdentityKey(title, year)
if title == "" || key == "" {
return
}
if _, exists := seen[key]; exists {
return
}
seen[key] = struct{}{}
out = append(out, MatchIdentityHint{Title: title, Year: year, Source: source})
}
itemType := strings.ToLower(strings.TrimSpace(skeleton.Type))
if parsed := naming.ParseFilename(file.FilePath, itemType); parsed != nil {
add(parsed.Title, parsed.Year, "current_path")
}
if itemType == matchContentTypeMovie {
base := strings.TrimSuffix(filepath.Base(file.FilePath), filepath.Ext(file.FilePath))
stem := naming.ParseInferMovieStem(base, skeleton.Title, skeleton.Year)
add(stem.Title, stem.Year, "filename")
observedRoot := strings.TrimSpace(skeleton.ObservedRootPath)
if observedRoot == "" {
observedRoot = filepath.Dir(file.FilePath)
}
rootStem := naming.ParseInferMovieStem(filepath.Base(observedRoot), "", 0)
add(rootStem.Title, rootStem.Year, "release_folder")
}
if len(out) > maxAlternateMatchIdentities {
out = out[:maxAlternateMatchIdentities]
}
return out
}
func matchIdentityKey(title string, year int) string {
title = strings.ToLower(strings.Join(strings.Fields(strings.TrimSpace(title)), " "))
if title == "" {
return ""
}
return title + "\x00" + strconv.Itoa(year)
}
func compactUniquePaths(paths []string) []string {
if len(paths) == 0 {
return nil
}
seen := make(map[string]struct{}, len(paths))
out := make([]string, 0, len(paths))
for _, path := range paths {
clean := filepath.Clean(path)
if clean == "." || clean == "" {
continue
}
if _, ok := seen[clean]; ok {
continue
}
seen[clean] = struct{}{}
out = append(out, clean)
}
return out
}
// ProcessFile applies the normal unmatched-file pipeline to a single file.
func (w *MatchWorker) ProcessFile(ctx context.Context, file *models.MediaFile) {
w.processFile(ctx, file)
}
// ProcessBatch fetches and processes one batch of unmatched files. Returns the
// number of files processed.
func (w *MatchWorker) ProcessBatch(ctx context.Context) (processed int, err error) {
if w.enableTVSeriesRootQueue && w.seriesClaimer != nil {
claimed := 0
for claimed < w.claimBatchSize() {
claimLimit := min(w.queueClaimSize(), w.claimBatchSize()-claimed)
jobs, err := w.seriesClaimer.Claim(ctx, claimLimit)
if err != nil {
return processed, err
}
if len(jobs) == 0 {
break
}
claimed += len(jobs)
batchProcessed, err := w.processSeriesRoots(ctx, jobs)
processed += batchProcessed
if err != nil {
return processed, err
}
if len(jobs) < claimLimit {
break
}
}
if claimed > 0 {
return processed, nil
}
}
if w.movieClaimer != nil {
claimed := 0
for claimed < w.claimBatchSize() {
claimLimit := min(w.queueClaimSize(), w.claimBatchSize()-claimed)
jobs, err := w.movieClaimer.Claim(ctx, claimLimit)
if err != nil {
return processed, err
}
if len(jobs) == 0 {
break
}
claimed += len(jobs)
processed += w.processQueuedMovieFiles(ctx, jobs)
if len(jobs) < claimLimit {
break
}
}
if claimed > 0 {
return processed, nil
}
}
files, err := w.claimBackgroundFiles(ctx)
if err != nil {
return 0, err
}
if len(files) == 0 {
return 0, nil
}
return w.processFiles(ctx, files), nil
}
// ProcessBatchByFolderAndPathPrefix processes unmatched files within a single
// library subtree immediately instead of waiting for the periodic worker loop.
func (w *MatchWorker) ProcessBatchByFolderAndPathPrefix(ctx context.Context, folderID int, pathPrefix string, attemptBefore time.Time) (processed int, err error) {
useSeriesQueue, useMovieQueue, err := w.queueUsageForFolder(ctx, folderID)
if err != nil {
return 0, err
}
if useSeriesQueue {
claimed := 0
for claimed < w.claimBatchSize() {
claimLimit := min(w.queueClaimSize(), w.claimBatchSize()-claimed)
jobs, err := w.seriesClaimer.ClaimByFolderAndPathPrefix(ctx, folderID, pathPrefix, claimLimit, attemptBefore)
if err != nil {
return processed, err
}
if len(jobs) == 0 {
break
}
claimed += len(jobs)
batchProcessed, err := w.processSeriesRoots(ctx, jobs)
processed += batchProcessed
if err != nil {
return processed, err
}
if len(jobs) < claimLimit {
break
}
}
if processed > 0 {
return processed, nil
}
}
if useSeriesQueue && !useMovieQueue {
return processed, nil
}
if useMovieQueue {
claimed := 0
for claimed < w.claimBatchSize() {
claimLimit := min(w.queueClaimSize(), w.claimBatchSize()-claimed)
jobs, err := w.movieClaimer.ClaimByFolderAndPathPrefix(ctx, folderID, pathPrefix, claimLimit, attemptBefore)
if err != nil {
return processed, err
}
if len(jobs) == 0 {
break
}
claimed += len(jobs)
processed += w.processQueuedMovieFiles(ctx, jobs)
if len(jobs) < claimLimit {
break
}
}
if processed > 0 || !useSeriesQueue {
return processed, nil
}
}
files, err := w.claimScopedFiles(ctx, folderID, pathPrefix, attemptBefore, scopedFallbackMode(useSeriesQueue, useMovieQueue))
if err != nil {
return 0, err
}
return w.processFiles(ctx, files), nil
}
// ProcessAllByFolderAndPathPrefix keeps draining scoped unmatched files until
// the subtree is empty.
func (w *MatchWorker) ProcessAllByFolderAndPathPrefix(ctx context.Context, folderID int, pathPrefix string, attemptBefore time.Time) (processed int, err error) {
if attemptBefore.IsZero() {
attemptBefore = time.Now().UTC()
}
useSeriesQueue, useMovieQueue, err := w.queueUsageForFolder(ctx, folderID)
if err != nil {
return 0, err
}
slog.InfoContext(ctx, "metadata: scoped matcher selected", "component", "metadata",
"folder_id", folderID,
"path_prefix", pathPrefix,
"matcher_path", scopedMatcherPath(useSeriesQueue, useMovieQueue),
)
for {
if ctx.Err() != nil {
return processed, ctx.Err()
}
if useSeriesQueue {
jobs, err := w.seriesClaimer.ClaimByFolderAndPathPrefix(ctx, folderID, pathPrefix, w.queueClaimSize(), attemptBefore)
if err != nil {
return processed, err
}
if len(jobs) > 0 {
batchProcessed, err := w.processSeriesRoots(ctx, jobs)
if err != nil {
return processed, err
}
processed += batchProcessed
continue
}
}
if useMovieQueue {
jobs, err := w.movieClaimer.ClaimByFolderAndPathPrefix(ctx, folderID, pathPrefix, w.queueClaimSize(), attemptBefore)
if err != nil {
return processed, err
}
if len(jobs) > 0 {
batchProcessed := w.processQueuedMovieFiles(ctx, jobs)
processed += batchProcessed
continue
}
}
if useSeriesQueue != useMovieQueue {
return processed, nil
}
files, err := w.claimScopedFiles(ctx, folderID, pathPrefix, attemptBefore, scopedFallbackMode(useSeriesQueue, useMovieQueue))
if err != nil {
return processed, err
}
batchProcessed := w.processFiles(ctx, files)
processed += batchProcessed
if batchProcessed == 0 {
return processed, nil
}
}
}
func (w *MatchWorker) processFiles(ctx context.Context, files []*models.MediaFile) int {
if len(files) == 0 {
return 0
}
claimedCount := len(files)
folders := sync.Map{}
selectedFiles := w.collapseClaimedSeriesBatch(ctx, files)
fileChan := make(chan *models.MediaFile, len(selectedFiles))
for _, f := range selectedFiles {
fileChan <- f
}
close(fileChan)
var (
wg sync.WaitGroup
processed atomic.Int64
deferredSeriesLinks sync.Map
)
for i := 0; i < w.workerCount(); i++ {
wg.Go(func() {
for file := range fileChan {
if ctx.Err() != nil {
return
}
if w.isMatchSuppressed(ctx, file) {
continue
}
w.processFileWithFolderCache(ctx, file, &folders, &deferredSeriesLinks)
processed.Add(1)
}
})
}
wg.Wait()
// Run deferred series episode linking once per unique series.
if ctx.Err() == nil {
deferredSeriesLinks.Range(func(key, _ any) bool {
if ctx.Err() != nil {
return false
}
contentID := key.(string)
if err := w.service.ensureSeriesEpisodeLinks(ctx, contentID); err != nil {
slog.WarnContext(ctx, "metadata: deferred series episode link failed", "component", "metadata",
"content_id", contentID, "error", err)
}
return true
})
}
return claimedCount
}
func (w *MatchWorker) isMatchSuppressed(ctx context.Context, file *models.MediaFile) bool {
if file == nil {
return true
}
checker, ok := w.fileLister.(MatchSuppressionChecker)
if !ok {
return false
}
suppressed, err := checker.IsMatchSuppressed(ctx, file.ID)
if err != nil {
slog.WarnContext(ctx, "metadata: failed to check match suppression", "component", "metadata",
"file_id", file.ID,
"path", file.FilePath,
"error", err,
)
return true
}
if suppressed {
slog.InfoContext(ctx, "metadata: skipping suppressed unmatched file", "component", "metadata",
"file_id", file.ID,
"path", file.FilePath,
)
}
return suppressed
}
func (w *MatchWorker) processQueuedMovieFiles(ctx context.Context, jobs []models.MovieMatchJob) int {
if len(jobs) == 0 {
return 0
}
defer w.releaseMovieLeases(jobs)
jobChan := make(chan models.MovieMatchJob, len(jobs))
for _, job := range jobs {
jobChan <- job
}
close(jobChan)
var (
wg sync.WaitGroup
processed atomic.Int64
folders sync.Map
)
for i := 0; i < w.workerCount(); i++ {
wg.Go(func() {
for job := range jobChan {
if ctx.Err() != nil {
return
}
if w.processQueuedMovieFile(ctx, job, &folders) {
processed.Add(1)
}
}
})
}
wg.Wait()
return int(processed.Load())
}
func (w *MatchWorker) releaseMovieLeases(jobs []models.MovieMatchJob) {
if w == nil || w.movieClaimer == nil {
return
}
tokens := make(map[string]struct{})
for _, job := range jobs {
if token := strings.TrimSpace(job.LeaseToken); token != "" {
tokens[token] = struct{}{}
}
}
for token := range tokens {
releaseCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
_, err := w.movieClaimer.ReleaseLease(releaseCtx, token)
cancel()
if err != nil {
slog.Warn("metadata: failed to release unfinished movie match lease", "component", "metadata", "error", err)
}
}
}
func (w *MatchWorker) processQueuedMovieFile(ctx context.Context, job models.MovieMatchJob, folderEnabledCache *sync.Map) bool {
file := job.File
if file == nil || w == nil || w.service == nil || w.movieClaimer == nil {
return false
}
if !w.folderEnabled(ctx, file.MediaFolderID, folderEnabledCache) {
return false
}
skeleton, reusedLinkedItem, err := w.queuedMovieSkeleton(ctx, file, job.RerunRequested)
if err != nil {
queueErr := truncateSeriesQueueError(err.Error())
if updateErr := w.movieClaimer.UpdateError(ctx, file.ID, job.LeaseToken, queueErr); updateErr != nil {
slog.WarnContext(ctx, "metadata: failed to update movie queue error", "component", "metadata",
"file_id", file.ID,
"path", file.FilePath,
"error", updateErr,
)
}
slog.WarnContext(ctx, "metadata: movie queue skeleton creation failed", "component", "metadata",
"file_id", file.ID,
"path", file.FilePath,
"error", err,
)
return false
}
if skeleton != nil && skeleton.ItemStatus == "skipped" {
// Deliberately skipped during skeleton creation (e.g. a misplaced TV
// episode inside a movie library): no item is created on purpose.
// Dequeue immediately; the recorded skipped root keeps the file out of
// future enqueues (see movieQueueFileEligibleCond), so this drains rows
// claimed before the skipped root was recorded.
if err := w.movieClaimer.Delete(ctx, file.ID, job.LeaseToken); err != nil {
slog.WarnContext(ctx, "metadata: failed to delete skipped movie queue row", "component", "metadata",
"file_id", file.ID,
"path", file.FilePath,
"error", err,
)
return false
}
return true
}
if skeleton == nil || strings.TrimSpace(skeleton.ContentID) == "" {
if updateErr := w.movieClaimer.UpdateError(ctx, file.ID, job.LeaseToken, truncateSeriesQueueError("movie queue claimed without a content id")); updateErr != nil {
slog.WarnContext(ctx, "metadata: failed to update movie queue error", "component", "metadata",
"file_id", file.ID,
"path", file.FilePath,
"error", updateErr,
)
}
return false
}
if skeleton.ItemStatus == "ambiguous" {
if err := w.movieClaimer.Delete(ctx, file.ID, job.LeaseToken); err != nil {
slog.WarnContext(ctx, "metadata: failed to delete ambiguous movie queue row", "component", "metadata",
"file_id", file.ID,
"path", file.FilePath,
"error", err,
)
return false
}
return true
}
if skeleton.IsNew || reusedLinkedItem {
req := w.buildProcessRequestForGroup(ctx, file, skeleton, nil)
result, processErr := w.service.Process(ctx, req)
if processErr != nil {
queueErr := truncateSeriesQueueError(processErr.Error())
if updateErr := w.movieClaimer.UpdateError(ctx, file.ID, job.LeaseToken, queueErr); updateErr != nil {
slog.WarnContext(ctx, "metadata: failed to update movie queue error", "component", "metadata",
"file_id", file.ID,
"path", file.FilePath,
"error", updateErr,
)
}
slog.WarnContext(ctx, "metadata: enrichment failed", "component", "metadata",
"file_id", file.ID,
"path", file.FilePath,
"error", processErr,
)
w.logStatusUpdateFailure(ctx, skeleton.ContentID, "unmatched", "content_id", skeleton.ContentID, "file_id", file.ID, "path", file.FilePath)
return false
} else if result != nil && !result.Updated {
if updateErr := w.updateMovieFailure(ctx, file.ID, job.LeaseToken, result.Decision); updateErr != nil {
slog.WarnContext(ctx, "metadata: failed to update movie queue error", "component", "metadata",
"file_id", file.ID,
"path", file.FilePath,
"error", updateErr,
)
}
w.logStatusUpdateFailure(ctx, skeleton.ContentID, "unmatched", "content_id", skeleton.ContentID, "file_id", file.ID, "path", file.FilePath)
return false
}
w.publishCatalogItemChanged(ctx, file.MediaFolderID, resultContentID(result, skeleton.ContentID), "metadata_updated")
}
if err := w.movieClaimer.Delete(ctx, file.ID, job.LeaseToken); err != nil {
slog.WarnContext(ctx, "metadata: failed to delete movie queue row", "component", "metadata",
"file_id", file.ID,
"path", file.FilePath,
"error", err,
)
return false
}
return true
}
func (w *MatchWorker) queuedMovieSkeleton(ctx context.Context, file *models.MediaFile, allowMatched bool) (*skeletonResult, bool, error) {
if skeleton, ok := w.reusableQueuedMovieSkeleton(ctx, file, allowMatched); ok {
return skeleton, true, nil
}
currentFile := reparseQueuedFileIdentity(file, "movie")
skeleton, err := w.service.createOrFindSkeleton(ctx, currentFile, currentFile.MediaFolderID)
if err != nil {
return nil, false, err
}
return skeleton, false, nil
}
func (w *MatchWorker) reusableQueuedMovieSkeleton(ctx context.Context, file *models.MediaFile, allowMatched bool) (*skeletonResult, bool) {
if w == nil || w.service == nil || w.service.itemRepo == nil || file == nil {
return nil, false
}
contentID := strings.TrimSpace(file.ContentID)
if contentID == "" {
return nil, false
}
item, err := w.service.itemRepo.GetByID(ctx, contentID)
if err != nil || item == nil {
return nil, false
}
status := strings.ToLower(strings.TrimSpace(item.Status))
reusableStatus := isSkeletonLikeStatus(status) ||
status == "ambiguous" ||
(allowMatched && status == string(MatchOutcomeMatched))
if !reusableStatus {
return nil, false
}
rootPath := filepath.Dir(file.FilePath)
if file.CanonicalRootPath != "" {
rootPath = filepath.Clean(file.CanonicalRootPath)
}
observedRootPath := filepath.Dir(file.FilePath)
if file.ObservedRootPath != "" {
observedRootPath = filepath.Clean(file.ObservedRootPath)
}
title := strings.TrimSpace(item.Title)
if title == "" {
title = file.BaseTitle
}
itemType := strings.TrimSpace(item.Type)
if itemType == "" {
itemType = file.BaseType
}
if itemType == "" {
itemType = "movie"
}
year := item.Year
if year == 0 {
year = file.BaseYear
}
currentFile := reparseQueuedFileIdentity(file, itemType)
if currentFile.BaseTitle != "" {
title = currentFile.BaseTitle
}
if currentFile.BaseYear != 0 {
year = currentFile.BaseYear
}
if currentFile.BaseType != "" {
itemType = currentFile.BaseType
}
refreshedIDs := trustedStructuredIDsForSkeleton(file.FilePath, observedRootPath, rootPath)
// This is an unmatched-style skeleton being re-evaluated. Structured path
// IDs are current input; IDs that disappeared from the path must not remain
// trusted forever merely because an older parser copied them onto the
// provisional item.
tmdbID, imdbID, tvdbID := "", "", ""
if refreshedIDs != nil {
if refreshedIDs.TmdbID != "" {
tmdbID = refreshedIDs.TmdbID
}
if refreshedIDs.ImdbID != "" {
imdbID = refreshedIDs.ImdbID
}
if refreshedIDs.TvdbID != "" {
tvdbID = refreshedIDs.TvdbID
}
}
return &skeletonResult{
ContentID: item.ContentID,
ItemStatus: item.Status,
RootPath: rootPath,
ObservedRootPath: observedRootPath,
GroupKeyVersion: file.GroupKeyVersion,
ContentGroupKey: strings.TrimSpace(file.ContentGroupKey),
Title: title,
Year: year,
Type: itemType,
TmdbID: tmdbID,
ImdbID: imdbID,
TvdbID: tvdbID,
}, true
}
// reparseQueuedFileIdentity refreshes the in-memory scanner identity from the
// current path without mutating the persisted scan row. Queue entries can live
// across parser releases, including entries that have not created a skeleton
// yet, so both fresh and reusable skeleton paths must use this view.
func reparseQueuedFileIdentity(file *models.MediaFile, fallbackType string) *models.MediaFile {
if file == nil {
return nil
}
current := *file
parseType := strings.TrimSpace(fallbackType)
if parseType == "" {
parseType = strings.TrimSpace(current.BaseType)
}
if parsed := naming.ParseFilename(current.FilePath, parseType); parsed != nil {
if parsed.Title != "" {
current.BaseTitle = parsed.Title
}
current.BaseYear = parsed.Year
if parsed.Type != "" {
current.BaseType = parsed.Type
}
}
return &current
}
func scopedMatcherPath(useSeriesQueue bool, useMovieQueue bool) string {
switch {
case useSeriesQueue && useMovieQueue:
return "series_root_queue+movie_file_queue"
case useSeriesQueue:
return "series_root_queue"
case useMovieQueue:
return "movie_file_queue"
default:
return "file"
}
}
func (w *MatchWorker) processSeriesRoots(ctx context.Context, jobs []models.SeriesRootMatchJob) (int, error) {
if len(jobs) == 0 {
return 0, nil
}
defer w.releaseSeriesLeases(jobs)
runCtx, cancel := context.WithCancel(ctx)
defer cancel()
jobChan := make(chan models.SeriesRootMatchJob, len(jobs))
for _, job := range jobs {
jobChan <- job
}
close(jobChan)
var (
wg sync.WaitGroup
processed atomic.Int64
firstErr error
firstErrMu sync.Mutex
folders sync.Map
)
for i := 0; i < w.workerCount(); i++ {
wg.Go(func() {
for job := range jobChan {
if runCtx.Err() != nil {
return
}
count, err := w.processSeriesRoot(runCtx, job, &folders)
if err != nil {
firstErrMu.Lock()
if firstErr == nil {
firstErr = err
cancel()
}
firstErrMu.Unlock()
return
}
processed.Add(int64(count))
}
})
}
wg.Wait()
return int(processed.Load()), firstErr
}
func (w *MatchWorker) releaseSeriesLeases(jobs []models.SeriesRootMatchJob) {
if w == nil || w.seriesClaimer == nil {
return
}
tokens := make(map[string]struct{})
for _, job := range jobs {
if token := strings.TrimSpace(job.LeaseToken); token != "" {
tokens[token] = struct{}{}
}
}
for token := range tokens {
releaseCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
_, err := w.seriesClaimer.ReleaseLease(releaseCtx, token)
cancel()
if err != nil {
slog.Warn("metadata: failed to release unfinished series match lease", "component", "metadata", "error", err)
}
}
}
func (w *MatchWorker) processSeriesRoot(ctx context.Context, job models.SeriesRootMatchJob, folderEnabledCache *sync.Map) (int, error) {
if !w.folderEnabled(ctx, job.MediaFolderID, folderEnabledCache) {
return 0, nil
}
if w.service == nil || w.service.fileRepo == nil || w.seriesClaimer == nil {
return 0, fmt.Errorf("series root matching requires file repo and queue claimer")
}
groupFiles, err := w.service.fileRepo.ListByObservedRootPath(ctx, job.MediaFolderID, job.ObservedRootPath)
if err != nil {
return 0, fmt.Errorf("loading files for series root %d/%s: %w", job.MediaFolderID, job.ObservedRootPath, err)
}
if len(groupFiles) == 0 {
if err := w.seriesClaimer.Delete(ctx, job.MediaFolderID, job.ObservedRootPath, job.LeaseToken); err != nil {
return 0, err
}
slog.InfoContext(ctx, "metadata: series root dropped because files disappeared", "component", "metadata",
"folder_id", job.MediaFolderID,
"observed_root_path", job.ObservedRootPath,
)
return 0, nil
}
representative := selectRepresentativeGroupFile(groupFiles)
if representative == nil {
if err := w.seriesClaimer.Delete(ctx, job.MediaFolderID, job.ObservedRootPath, job.LeaseToken); err != nil {
return 0, err
}
return 0, nil
}
if !hasUnlinkedGroupFile(groupFiles) {
if strings.TrimSpace(representative.ContentID) != "" {
if skeleton, ok := w.reusableQueuedMovieSkeleton(ctx, representative, job.RerunRequested); ok && skeleton.ItemStatus != "ambiguous" {
req := w.buildProcessRequestForGroup(ctx, representative, skeleton, groupFiles)
result, processErr := w.service.Process(ctx, req)
if processErr != nil {
queueErr := truncateSeriesQueueError(processErr.Error())
if updateErr := w.seriesClaimer.UpdateError(ctx, job.MediaFolderID, job.ObservedRootPath, job.LeaseToken, queueErr); updateErr != nil {
return 0, updateErr
}
slog.WarnContext(ctx, "metadata: enrichment failed", "component", "metadata",
"file_id", representative.ID,
"path", representative.FilePath,
"error", processErr)
w.logStatusUpdateFailure(ctx, skeleton.ContentID, "unmatched",
"content_id", skeleton.ContentID,
"file_id", representative.ID,
"path", representative.FilePath)
if fbErr := w.service.SynthesizeFallbackEpisodes(ctx, skeleton.ContentID); fbErr != nil {
slog.WarnContext(ctx, "metadata: fallback episode synthesis failed after series root enrichment error", "component", "metadata",
"content_id", skeleton.ContentID, "error", fbErr)
}
return 0, nil
}
if result != nil && !result.Updated {
if updateErr := w.updateSeriesFailure(ctx, job.MediaFolderID, job.ObservedRootPath, job.LeaseToken, result.Decision); updateErr != nil {
return 0, updateErr
}
w.logStatusUpdateFailure(ctx, skeleton.ContentID, "unmatched",
"content_id", skeleton.ContentID,
"file_id", representative.ID,
"path", representative.FilePath)
if fbErr := w.service.SynthesizeFallbackEpisodes(ctx, skeleton.ContentID); fbErr != nil {
slog.WarnContext(ctx, "metadata: fallback episode synthesis failed for unmatched series root", "component", "metadata",
"content_id", skeleton.ContentID, "error", fbErr)
}
return 0, nil
}
w.publishCatalogItemChanged(ctx, job.MediaFolderID, resultContentID(result, skeleton.ContentID), "metadata_updated")
}
if err := w.service.ensureSeriesEpisodeLinks(ctx, representative.ContentID); err != nil {
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.InfoContext(ctx, "metadata: series item gone during episode-link ensure (likely concurrent merge); skipping", "component", "metadata",
"content_id", representative.ContentID,
"folder_id", job.MediaFolderID,
"observed_root_path", job.ObservedRootPath)
} else {
if updateErr := w.seriesClaimer.UpdateError(ctx, job.MediaFolderID, job.ObservedRootPath, job.LeaseToken, truncateSeriesQueueError(err.Error())); updateErr != nil {
return 0, updateErr
}
return 0, fmt.Errorf("ensuring series episode links for %s: %w", representative.ContentID, err)
}
}
}
if err := w.seriesClaimer.Delete(ctx, job.MediaFolderID, job.ObservedRootPath, job.LeaseToken); err != nil {
return 0, err
}
return len(groupFiles), nil
}
currentRepresentative := reparseQueuedFileIdentity(representative, "series")
skeleton, err := w.service.createOrFindSkeleton(ctx, currentRepresentative, job.MediaFolderID)
if err != nil {
queueErr := truncateSeriesQueueError(err.Error())
if updateErr := w.seriesClaimer.UpdateError(ctx, job.MediaFolderID, job.ObservedRootPath, job.LeaseToken, queueErr); updateErr != nil {
return 0, updateErr
}
slog.WarnContext(ctx, "metadata: series root skeleton creation failed", "component", "metadata",
"folder_id", job.MediaFolderID,
"observed_root_path", job.ObservedRootPath,
"sample_file_path", job.SampleFilePath,
"error", err,
)
return 0, nil
}
if skeleton == nil || strings.TrimSpace(skeleton.ContentID) == "" {
return 0, nil
}
if _, err := w.service.fileRepo.UpdateContentIDByObservedRootPath(ctx, job.MediaFolderID, job.ObservedRootPath, skeleton.ContentID); err != nil {
if updateErr := w.seriesClaimer.UpdateError(ctx, job.MediaFolderID, job.ObservedRootPath, job.LeaseToken, truncateSeriesQueueError(err.Error())); updateErr != nil {
return 0, updateErr
}
return 0, fmt.Errorf("relinking series root %d/%s: %w", job.MediaFolderID, job.ObservedRootPath, err)
}
needsInitialMatch := skeleton.IsNew || job.RerunRequested
if !needsInitialMatch && strings.TrimSpace(skeleton.ContentID) != "" && w.service.itemRepo != nil {
item, loadErr := w.service.itemRepo.GetByID(ctx, skeleton.ContentID)
if loadErr != nil {
queueErr := truncateSeriesQueueError(loadErr.Error())
if updateErr := w.seriesClaimer.UpdateError(ctx, job.MediaFolderID, job.ObservedRootPath, job.LeaseToken, queueErr); updateErr != nil {
return 0, updateErr
}
return 0, fmt.Errorf("loading linked series skeleton %s: %w", skeleton.ContentID, loadErr)
}
if item != nil {
skeleton.ItemStatus = item.Status
needsInitialMatch = isSkeletonLikeStatus(item.Status)
}
}
if needsInitialMatch && skeleton.ItemStatus != "ambiguous" {
req := w.buildProcessRequestForGroup(ctx, representative, skeleton, groupFiles)
result, processErr := w.service.Process(ctx, req)
if processErr != nil {
queueErr := truncateSeriesQueueError(processErr.Error())
if updateErr := w.seriesClaimer.UpdateError(ctx, job.MediaFolderID, job.ObservedRootPath, job.LeaseToken, queueErr); updateErr != nil {
return 0, updateErr
}
slog.WarnContext(ctx, "metadata: enrichment failed", "component", "metadata",
"file_id", representative.ID,
"path", representative.FilePath,
"error", processErr)
w.logStatusUpdateFailure(ctx, skeleton.ContentID, "unmatched",
"content_id", skeleton.ContentID,
"file_id", representative.ID,
"path", representative.FilePath,
"folder_id", job.MediaFolderID,
"observed_root_path", job.ObservedRootPath,
)
if fbErr := w.service.SynthesizeFallbackEpisodes(ctx, skeleton.ContentID); fbErr != nil {
slog.WarnContext(ctx, "metadata: fallback episode synthesis failed after enrichment error", "component", "metadata",
"content_id", skeleton.ContentID,
"error", fbErr,
)
}
return 0, nil
} else if result != nil && !result.Updated {
if updateErr := w.updateSeriesFailure(ctx, job.MediaFolderID, job.ObservedRootPath, job.LeaseToken, result.Decision); updateErr != nil {
return 0, updateErr
}
w.logStatusUpdateFailure(ctx, skeleton.ContentID, "unmatched",
"content_id", skeleton.ContentID,
"file_id", representative.ID,
"path", representative.FilePath,
"folder_id", job.MediaFolderID,
"observed_root_path", job.ObservedRootPath,
)
if fbErr := w.service.SynthesizeFallbackEpisodes(ctx, skeleton.ContentID); fbErr != nil {
slog.WarnContext(ctx, "metadata: fallback episode synthesis failed for unmatched series", "component", "metadata",
"content_id", skeleton.ContentID,
"error", fbErr,
)
}
return 0, nil
}
w.publishCatalogItemChanged(ctx, job.MediaFolderID, resultContentID(result, skeleton.ContentID), "metadata_updated")
}
finalContentID, err := w.service.fileRepo.FindContentIDByObservedRootPath(ctx, job.MediaFolderID, job.ObservedRootPath, "series")
if err != nil {
if updateErr := w.seriesClaimer.UpdateError(ctx, job.MediaFolderID, job.ObservedRootPath, job.LeaseToken, truncateSeriesQueueError(err.Error())); updateErr != nil {
return 0, updateErr
}
return 0, fmt.Errorf("resolving final content for series root %d/%s: %w", job.MediaFolderID, job.ObservedRootPath, err)
}
if finalContentID == "" {
finalContentID = skeleton.ContentID
}
if strings.TrimSpace(finalContentID) != "" {
if err := w.service.ensureSeriesEpisodeLinks(ctx, finalContentID); err != nil {
if errors.Is(err, catalog.ErrItemNotFound) {
slog.InfoContext(ctx, "metadata: series item gone during episode-link ensure (likely concurrent merge); skipping", "component", "metadata",
"content_id", finalContentID,
"folder_id", job.MediaFolderID,
"observed_root_path", job.ObservedRootPath)
} else {
if updateErr := w.seriesClaimer.UpdateError(ctx, job.MediaFolderID, job.ObservedRootPath, job.LeaseToken, truncateSeriesQueueError(err.Error())); updateErr != nil {
return 0, updateErr
}
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)
}
}
if err := w.seriesClaimer.Delete(ctx, job.MediaFolderID, job.ObservedRootPath, job.LeaseToken); err != nil {
return 0, err
}
slog.InfoContext(ctx, "metadata: series root completed", "component", "metadata",
"folder_id", job.MediaFolderID,
"observed_root_path", job.ObservedRootPath,
"content_id", finalContentID,
"file_count", len(groupFiles),
)
return len(groupFiles), nil
}
func matchFailureFromDecision(decision *MatchDecision) MatchFailure {
kind := MatchOutcomeMetadataEmpty
if decision != nil && strings.TrimSpace(string(decision.Outcome)) != "" {
kind = decision.Outcome
}
message := ErrMetadataNotFound.Error()
if decision != nil {
switch decision.Outcome {
case MatchOutcomeNoCandidates:
message = "no provider candidates returned"
case MatchOutcomeCandidateRejected:
message = "provider candidates did not meet the automatic match threshold"
case MatchOutcomeTrustedIDConflict:
message = "provider candidates conflicted with a trusted external ID"
case MatchOutcomeTrustedIDTypeMismatch:
message = "trusted external ID resolves to the opposite library type"
case MatchOutcomeMetadataEmpty:
message = "selected providers returned no usable metadata"
case MatchOutcomeProviderTransient:
message = "metadata provider is temporarily unavailable"
case MatchOutcomeProviderPermanent:
message = "metadata provider rejected the request permanently"
}
}
return MatchFailure{Kind: kind, Message: message, Decision: decision}
}
func (w *MatchWorker) updateMovieFailure(ctx context.Context, mediaFileID int, leaseToken string, decision *MatchDecision) error {
return w.movieClaimer.UpdateFailure(ctx, mediaFileID, leaseToken, matchFailureFromDecision(decision))
}
func (w *MatchWorker) updateSeriesFailure(ctx context.Context, folderID int, observedRootPath, leaseToken string, decision *MatchDecision) error {
return w.seriesClaimer.UpdateFailure(ctx, folderID, observedRootPath, leaseToken, matchFailureFromDecision(decision))
}
func (w *MatchWorker) logStatusUpdateFailure(ctx context.Context, contentID, status string, attrs ...any) {
if w == nil || w.service == nil {
return
}
if err := w.service.updateItemStatus(ctx, contentID, status); err != nil {
args := append([]any{"content_id", contentID, "status", status, "error", err}, attrs...)
slog.WarnContext(ctx, "metadata: failed to update item status", append([]any{"component", "metadata"}, args...)...)
}
}
func (w *MatchWorker) collapseClaimedSeriesBatch(ctx context.Context, files []*models.MediaFile) []*models.MediaFile {
if len(files) == 0 || w == nil || w.service == nil || w.service.folderRepo == nil {
return files
}
out := make([]*models.MediaFile, 0, len(files))
folderTypes := make(map[int]string)
seenGroups := make(map[string]struct{})
for _, file := range files {
if file == nil {
continue
}
folderType, ok := folderTypes[file.MediaFolderID]
if !ok {
folder, err := w.service.folderRepo.GetByID(ctx, file.MediaFolderID)
if err != nil {
slog.WarnContext(ctx, "metadata: failed to load folder type for batch compaction", "component", "metadata",
"folder_id", file.MediaFolderID,
"file_id", file.ID,
"error", err,
)
}
if folder != nil {
folderType = strings.ToLower(strings.TrimSpace(folder.Type))
}
folderTypes[file.MediaFolderID] = folderType
}
switch folderType {
case "series", "tv", "show", "tvshows":
default:
out = append(out, file)
continue
}
if file.ContentGroupKey == "" {
out = append(out, file)
continue
}
groupKey := fmt.Sprintf("%d:%d:%s", file.MediaFolderID, file.GroupKeyVersion, file.ContentGroupKey)
if _, ok := seenGroups[groupKey]; ok {
continue
}
seenGroups[groupKey] = struct{}{}
out = append(out, file)
}
return out
}
func (w *MatchWorker) folderEnabled(ctx context.Context, folderID int, cache *sync.Map) bool {
if folderID <= 0 || w == nil || w.service == nil || w.service.folderRepo == nil {
return true
}
if cache != nil {
if cached, ok := cache.Load(folderID); ok {
return cached.(bool)
}
}
enabled := true
folder, err := w.service.folderRepo.GetByID(ctx, folderID)
if err != nil {
slog.WarnContext(ctx, "metadata: failed to load folder state during match", "component", "metadata",
"folder_id", folderID,
"error", err,
)
} else if folder != nil {
enabled = folder.Enabled
}
if cache != nil {
cache.Store(folderID, enabled)
}
return enabled
}
func (w *MatchWorker) folderType(ctx context.Context, folderID int) (string, error) {
if folderID <= 0 || w == nil || w.service == nil || w.service.folderRepo == nil {
return "", nil
}
folder, err := w.service.folderRepo.GetByID(ctx, folderID)
if err != nil {
return "", err
}
if folder == nil {
return "", nil
}
return strings.ToLower(strings.TrimSpace(folder.Type)), nil
}
func (w *MatchWorker) queueUsageForFolder(ctx context.Context, folderID int) (useSeriesQueue bool, useMovieQueue bool, err error) {
folderType, err := w.folderType(ctx, folderID)
if err != nil {
return false, false, err
}
useSeriesQueue = w.enableTVSeriesRootQueue &&
w.seriesClaimer != nil &&
(librarykind.IsTV(folderType) || librarykind.IsMixed(folderType))
// Mixed libraries feed the movie queue too: their unmatched files may be
// either kind, and the movie queue is the fallback lane.
useMovieQueue = w.movieClaimer != nil &&
(librarykind.IsMovie(folderType) || librarykind.IsMixed(folderType))
return useSeriesQueue, useMovieQueue, nil
}
func (w *MatchWorker) claimBackgroundFiles(ctx context.Context) ([]*models.MediaFile, error) {
if w.enableTVSeriesRootQueue && w.movieClaimer != nil {
if claimer, ok := w.fileLister.(MixedFileClaimer); ok {
return claimer.ClaimUnmatchedMixed(ctx, w.claimBatchSize())
}
return nil, fmt.Errorf("mixed-library file claimer is not configured")
}
if w.enableTVSeriesRootQueue {
if claimer, ok := w.fileLister.(NonSeriesFileClaimer); ok {
return claimer.ClaimUnmatchedNonSeries(ctx, w.claimBatchSize())
}
return nil, fmt.Errorf("non-series file claimer is not configured")
}
return w.fileLister.ClaimUnmatched(ctx, w.claimBatchSize())
}
type scopedFallbackClaimMode int
const (
scopedFallbackGeneric scopedFallbackClaimMode = iota
scopedFallbackNonSeries
scopedFallbackMixed
)
func scopedFallbackMode(useSeriesQueue bool, useMovieQueue bool) scopedFallbackClaimMode {
switch {
case useSeriesQueue && useMovieQueue:
return scopedFallbackMixed
case useSeriesQueue:
return scopedFallbackNonSeries
default:
return scopedFallbackGeneric
}
}
func (w *MatchWorker) claimScopedFiles(ctx context.Context, folderID int, pathPrefix string, attemptBefore time.Time, mode scopedFallbackClaimMode) ([]*models.MediaFile, error) {
switch mode {
case scopedFallbackMixed:
if claimer, ok := w.fileLister.(MixedFileClaimer); ok {
return claimer.ClaimUnmatchedMixedByFolderAndPathPrefix(ctx, folderID, pathPrefix, w.claimBatchSize(), attemptBefore)
}
return nil, fmt.Errorf("mixed-library file claimer is not configured")
case scopedFallbackNonSeries:
if claimer, ok := w.fileLister.(NonSeriesFileClaimer); ok {
return claimer.ClaimUnmatchedNonSeriesByFolderAndPathPrefix(ctx, folderID, pathPrefix, w.claimBatchSize(), attemptBefore)
}
return nil, fmt.Errorf("non-series file claimer is not configured")
default:
return w.fileLister.ClaimUnmatchedByFolderAndPathPrefix(ctx, folderID, pathPrefix, w.claimBatchSize(), attemptBefore)
}
}
func selectRepresentativeGroupFile(groupFiles []*models.MediaFile) *models.MediaFile {
var (
firstUnlinked *models.MediaFile
firstAny *models.MediaFile
)
for _, file := range groupFiles {
if file == nil {
continue
}
if firstAny == nil || file.ID < firstAny.ID {
firstAny = file
}
if strings.TrimSpace(file.ContentID) == "" && (firstUnlinked == nil || file.ID < firstUnlinked.ID) {
firstUnlinked = file
}
}
if firstUnlinked != nil {
return firstUnlinked
}
return firstAny
}
func hasUnlinkedGroupFile(groupFiles []*models.MediaFile) bool {
for _, file := range groupFiles {
if file != nil && strings.TrimSpace(file.ContentID) == "" {
return true
}
}
return false
}
func truncateSeriesQueueError(errText string) string {
errText = strings.TrimSpace(errText)
if len(errText) <= 1024 {
return errText
}
return errText[:1024]
}
// RetryUnmatchedItemsByFolderAndPathPrefix revisits linked unmatched items in
// scope once. Per-item retry failures are counted as warnings, not fatal.
func (w *MatchWorker) RetryUnmatchedItemsByFolderAndPathPrefix(ctx context.Context, folderID int, pathPrefix string) (retried int, stillUnmatched int, err error) {
if w.service == nil {
return 0, 0, fmt.Errorf("metadata match worker requires a service")
}
if w.itemLister == nil {
return 0, 0, fmt.Errorf("metadata match worker requires an item lister")
}
contentIDs, err := w.itemLister.ListUnmatchedByFolderAndPathPrefix(ctx, folderID, pathPrefix, 0)
if err != nil {
return 0, 0, err
}
for _, contentID := range contentIDs {
if ctx.Err() != nil {
return retried, stillUnmatched, ctx.Err()
}
retried++
result, processErr := w.service.Process(ctx, ProcessRequest{
ContentID: contentID,
FolderID: formatFolderID(folderID),
Mode: ModeScheduledRefresh,
})
if processErr != nil {
stillUnmatched++
slog.WarnContext(ctx, "metadata: scoped retry failed", "component", "metadata",
"content_id", contentID,
"folder_id", folderID,
"path_prefix", pathPrefix,
"error", processErr)
continue
}
if result == nil || !result.Updated {
stillUnmatched++
continue
}
w.publishCatalogItemChanged(ctx, folderID, resultContentID(result, contentID), "metadata_updated")
}
if repaired, repairErr := w.service.repairMatchedDuplicateProviderOwnersByFolderAndPathPrefix(ctx, folderID, pathPrefix); repairErr != nil {
return retried, stillUnmatched, repairErr
} else if repaired > 0 {
retried += repaired
slog.InfoContext(ctx, "metadata: repaired matched duplicate items in scope", "component", "metadata",
"folder_id", folderID,
"path_prefix", pathPrefix,
"repaired_items", repaired,
)
}
repaired, repairErr := w.service.repairAnchoredIdentityMismatchesByFolderAndPathPrefix(ctx, folderID, pathPrefix)
retried += repaired
if repairErr != nil {
return retried, stillUnmatched, repairErr
}
if repaired > 0 {
slog.InfoContext(ctx, "metadata: repaired provider-anchored roots in scope", "component", "metadata",
"folder_id", folderID,
"path_prefix", pathPrefix,
"repaired_roots", repaired,
)
}
return retried, stillUnmatched, nil
}
func formatFolderID(id int) string {
return strconv.Itoa(id)
}
func (w *MatchWorker) publishCatalogItemChanged(ctx context.Context, libraryID int, contentID string, change string) {
if w == nil || w.realtimeHub == nil || libraryID <= 0 || strings.TrimSpace(contentID) == "" {
return
}
if err := w.realtimeHub.PublishCatalogItemChanged(ctx, notifications.MetadataUpdateEvent{
LibraryID: libraryID,
ContentID: contentID,
Change: change,
}); err != nil {
slog.WarnContext(ctx, "metadata: failed to publish catalog item change", "component", "metadata",
"content_id", contentID,
"library_id", libraryID,
"change", change,
"error", err,
)
}
}
func resultContentID(result *ProcessResult, fallback string) string {
if result != nil && strings.TrimSpace(result.ContentID) != "" {
return result.ContentID
}
return fallback
}