2026-05-22 20:26:11 -04:00
package metadata
import (
"context"
2026-05-27 21:01:28 +02:00
"errors"
2026-05-22 20:26:11 -04:00
"fmt"
"log/slog"
"path/filepath"
"slices"
"strconv"
"strings"
"sync"
"sync/atomic"
"time"
2026-05-27 21:01:28 +02:00
"github.com/Silo-Server/silo-server/internal/catalog"
2026-07-02 11:42:28 -04:00
"github.com/Silo-Server/silo-server/internal/librarykind"
2026-05-22 20:26:11 -04:00
"github.com/Silo-Server/silo-server/internal/models"
2026-07-24 18:02:52 +02:00
"github.com/Silo-Server/silo-server/internal/naming"
2026-06-05 19:43:20 -07:00
"github.com/Silo-Server/silo-server/internal/notifications"
2026-05-22 20:26:11 -04:00
)
// MatchWorker processes unmatched files in the background.
type MatchWorker struct {
service * MetadataService
fileLister UnmatchedFileLister
itemLister UnmatchedItemLister
movieClaimer MovieFileClaimer
seriesClaimer SeriesRootClaimer
enableTVSeriesRootQueue bool
2026-06-05 19:43:20 -07:00
realtimeHub * notifications . Hub
2026-06-10 19:25:07 -04:00
// 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
2026-05-22 20:26:11 -04:00
}
2026-06-10 19:25:07 -04:00
// 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 ()) }
2026-07-24 18:18:49 +02:00
func ( w * MatchWorker ) queueClaimSize () int {
return min ( w . workerCount (), w . claimBatchSize ())
}
2026-06-10 19:25:07 -04:00
2026-05-22 20:26:11 -04:00
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 )
}
2026-06-05 19:43:20 -07:00
type MatchSuppressionChecker interface {
IsMatchSuppressed ( ctx context . Context , fileID int ) ( bool , error )
}
2026-05-22 20:26:11 -04:00
type MovieFileClaimer interface {
2026-07-24 18:18:49 +02:00
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 )
2026-05-22 20:26:11 -04:00
}
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 )
2026-07-24 18:18:49 +02:00
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 )
2026-05-22 20:26:11 -04:00
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
}
2026-06-10 19:25:07 -04:00
w := & MatchWorker {
2026-05-22 20:26:11 -04:00
service : service ,
fileLister : fileLister ,
itemLister : itemLister ,
interval : interval ,
}
2026-06-10 19:25:07 -04:00
w . SetConcurrency ( workers , batchSize )
return w
2026-05-22 20:26:11 -04:00
}
// 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
}
2026-06-05 19:43:20 -07:00
// SetRealtimeHub publishes catalog item changes as metadata enrichment commits.
func ( w * MatchWorker ) SetRealtimeHub ( hub * notifications . Hub ) {
if w == nil {
return
}
w . realtimeHub = hub
}
2026-05-22 20:26:11 -04:00
// 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 {
2026-07-24 18:18:49 +02:00
w . processBackgroundSeriesQueue ( ctx )
2026-05-22 20:26:11 -04:00
}
if w . movieClaimer != nil {
2026-07-24 18:18:49 +02:00
w . processBackgroundMovieQueue ( ctx )
2026-05-22 20:26:11 -04:00
}
files , err := w . claimBackgroundFiles ( ctx )
if err != nil {
2026-07-09 20:53:52 +08:00
slog . ErrorContext ( ctx , "metadata: failed to list unmatched files" , "component" , "metadata" , "error" , err )
2026-05-22 20:26:11 -04:00
return
}
if len ( files ) == 0 {
return
}
2026-07-09 20:53:52 +08:00
slog . InfoContext ( ctx , "metadata: processing unmatched files" , "component" , "metadata" , "count" , len ( files ))
2026-05-22 20:26:11 -04:00
w . processFiles ( ctx , files )
}
2026-07-24 18:18:49 +02:00
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
}
}
}
2026-05-22 20:26:11 -04:00
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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to record match attempt" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"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 ) {
2026-07-09 20:53:52 +08:00
slog . InfoContext ( ctx , "metadata: skipping file in disabled library" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: skeleton creation failed" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to ensure series episode links" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: enrichment failed" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: fallback episode synthesis failed after enrichment error" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: fallback episode synthesis failed for unmatched series" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"content_id" , skeleton . ContentID , "error" , fbErr )
}
}
2026-06-05 19:43:20 -07:00
return
2026-05-22 20:26:11 -04:00
}
2026-06-05 19:43:20 -07:00
w . publishCatalogItemChanged ( ctx , file . MediaFolderID , resultContentID ( result , skeleton . ContentID ), "metadata_updated" )
2026-05-22 20:26:11 -04:00
}
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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to load observed-root files" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"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" ,
2026-07-24 18:02:52 +02:00
AlternateIdentities : queuedMatchIdentityAlternates ( representative , skeleton ),
2026-05-22 20:26:11 -04:00
}
return ProcessRequest {
ContentID : skeleton . ContentID ,
Hints : hints ,
FolderID : formatFolderID ( representative . MediaFolderID ),
Mode : ModeInitialMatch ,
}
}
2026-07-24 18:02:52 +02:00
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 )
}
2026-05-22 20:26:11 -04:00
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 {
2026-07-24 18:18:49 +02:00
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
}
2026-05-22 20:26:11 -04:00
}
2026-07-24 18:18:49 +02:00
if claimed > 0 {
return processed , nil
2026-05-22 20:26:11 -04:00
}
}
if w . movieClaimer != nil {
2026-07-24 18:18:49 +02:00
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
}
2026-05-22 20:26:11 -04:00
}
2026-07-24 18:18:49 +02:00
if claimed > 0 {
return processed , nil
2026-05-22 20:26:11 -04:00
}
}
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 ) {
2026-05-26 20:07:43 -04:00
useSeriesQueue , useMovieQueue , err := w . queueUsageForFolder ( ctx , folderID )
2026-05-22 20:26:11 -04:00
if err != nil {
return 0 , err
}
if useSeriesQueue {
2026-07-24 18:18:49 +02:00
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
}
2026-05-22 20:26:11 -04:00
}
2026-07-24 18:18:49 +02:00
if processed > 0 {
return processed , nil
2026-05-26 20:07:43 -04:00
}
2026-05-22 20:26:11 -04:00
}
2026-05-26 20:07:43 -04:00
if useSeriesQueue && ! useMovieQueue {
return processed , nil
2026-05-22 20:26:11 -04:00
}
if useMovieQueue {
2026-07-24 18:18:49 +02:00
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
}
2026-05-26 20:07:43 -04:00
}
2026-07-24 18:18:49 +02:00
if processed > 0 || ! useSeriesQueue {
2026-05-26 20:07:43 -04:00
return processed , nil
}
2026-05-22 20:26:11 -04:00
}
2026-05-26 20:07:43 -04:00
files , err := w . claimScopedFiles ( ctx , folderID , pathPrefix , attemptBefore , scopedFallbackMode ( useSeriesQueue , useMovieQueue ))
2026-05-22 20:26:11 -04:00
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 ()
}
2026-05-26 20:07:43 -04:00
useSeriesQueue , useMovieQueue , err := w . queueUsageForFolder ( ctx , folderID )
2026-05-22 20:26:11 -04:00
if err != nil {
return 0 , err
}
2026-07-09 20:53:52 +08:00
slog . InfoContext ( ctx , "metadata: scoped matcher selected" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"folder_id" , folderID ,
"path_prefix" , pathPrefix ,
"matcher_path" , scopedMatcherPath ( useSeriesQueue , useMovieQueue ),
)
for {
if ctx . Err () != nil {
return processed , ctx . Err ()
}
if useSeriesQueue {
2026-07-24 18:18:49 +02:00
jobs , err := w . seriesClaimer . ClaimByFolderAndPathPrefix ( ctx , folderID , pathPrefix , w . queueClaimSize (), attemptBefore )
2026-05-22 20:26:11 -04:00
if err != nil {
return processed , err
}
2026-05-26 20:07:43 -04:00
if len ( jobs ) > 0 {
batchProcessed , err := w . processSeriesRoots ( ctx , jobs )
if err != nil {
return processed , err
}
processed += batchProcessed
continue
2026-05-22 20:26:11 -04:00
}
}
if useMovieQueue {
2026-07-24 18:18:49 +02:00
jobs , err := w . movieClaimer . ClaimByFolderAndPathPrefix ( ctx , folderID , pathPrefix , w . queueClaimSize (), attemptBefore )
2026-05-22 20:26:11 -04:00
if err != nil {
return processed , err
}
2026-07-24 18:18:49 +02:00
if len ( jobs ) > 0 {
batchProcessed := w . processQueuedMovieFiles ( ctx , jobs )
2026-05-26 20:07:43 -04:00
processed += batchProcessed
continue
2026-05-22 20:26:11 -04:00
}
2026-05-26 20:07:43 -04:00
}
if useSeriesQueue != useMovieQueue {
return processed , nil
2026-05-22 20:26:11 -04:00
}
2026-05-26 20:07:43 -04:00
files , err := w . claimScopedFiles ( ctx , folderID , pathPrefix , attemptBefore , scopedFallbackMode ( useSeriesQueue , useMovieQueue ))
2026-05-22 20:26:11 -04:00
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
)
2026-06-10 19:25:07 -04:00
for i := 0 ; i < w . workerCount (); i ++ {
2026-05-22 20:26:11 -04:00
wg . Go ( func () {
for file := range fileChan {
if ctx . Err () != nil {
return
}
2026-06-05 19:43:20 -07:00
if w . isMatchSuppressed ( ctx , file ) {
continue
}
2026-05-22 20:26:11 -04:00
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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: deferred series episode link failed" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"content_id" , contentID , "error" , err )
}
return true
})
}
return claimedCount
}
2026-06-05 19:43:20 -07:00
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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to check match suppression" , "component" , "metadata" ,
2026-06-05 19:43:20 -07:00
"file_id" , file . ID ,
"path" , file . FilePath ,
"error" , err ,
)
return true
}
if suppressed {
2026-07-09 20:53:52 +08:00
slog . InfoContext ( ctx , "metadata: skipping suppressed unmatched file" , "component" , "metadata" ,
2026-06-05 19:43:20 -07:00
"file_id" , file . ID ,
"path" , file . FilePath ,
)
}
return suppressed
}
2026-07-24 18:18:49 +02:00
func ( w * MatchWorker ) processQueuedMovieFiles ( ctx context . Context , jobs [] models . MovieMatchJob ) int {
if len ( jobs ) == 0 {
2026-05-22 20:26:11 -04:00
return 0
}
2026-07-24 18:18:49 +02:00
defer w . releaseMovieLeases ( jobs )
2026-05-22 20:26:11 -04:00
2026-07-24 18:18:49 +02:00
jobChan := make ( chan models . MovieMatchJob , len ( jobs ))
for _ , job := range jobs {
jobChan <- job
2026-05-22 20:26:11 -04:00
}
2026-07-24 18:18:49 +02:00
close ( jobChan )
2026-05-22 20:26:11 -04:00
var (
wg sync . WaitGroup
processed atomic . Int64
folders sync . Map
)
2026-06-10 19:25:07 -04:00
for i := 0 ; i < w . workerCount (); i ++ {
2026-05-22 20:26:11 -04:00
wg . Go ( func () {
2026-07-24 18:18:49 +02:00
for job := range jobChan {
2026-05-22 20:26:11 -04:00
if ctx . Err () != nil {
return
}
2026-07-24 18:18:49 +02:00
if w . processQueuedMovieFile ( ctx , job , & folders ) {
2026-05-22 20:26:11 -04:00
processed . Add ( 1 )
}
}
})
}
wg . Wait ()
return int ( processed . Load ())
}
2026-07-24 18:18:49 +02:00
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
2026-05-22 20:26:11 -04:00
if file == nil || w == nil || w . service == nil || w . movieClaimer == nil {
return false
}
if ! w . folderEnabled ( ctx , file . MediaFolderID , folderEnabledCache ) {
return false
}
2026-07-24 18:18:49 +02:00
skeleton , reusedLinkedItem , err := w . queuedMovieSkeleton ( ctx , file , job . RerunRequested )
2026-05-22 20:26:11 -04:00
if err != nil {
queueErr := truncateSeriesQueueError ( err . Error ())
2026-07-24 18:18:49 +02:00
if updateErr := w . movieClaimer . UpdateError ( ctx , file . ID , job . LeaseToken , queueErr ); updateErr != nil {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to update movie queue error" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"file_id" , file . ID ,
"path" , file . FilePath ,
"error" , updateErr ,
)
}
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: movie queue skeleton creation failed" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"file_id" , file . ID ,
"path" , file . FilePath ,
"error" , err ,
)
return false
}
2026-06-10 23:18:41 +02:00
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.
2026-07-24 18:18:49 +02:00
if err := w . movieClaimer . Delete ( ctx , file . ID , job . LeaseToken ); err != nil {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to delete skipped movie queue row" , "component" , "metadata" ,
2026-06-10 23:18:41 +02:00
"file_id" , file . ID ,
"path" , file . FilePath ,
"error" , err ,
)
return false
}
return true
}
2026-05-22 20:26:11 -04:00
if skeleton == nil || strings . TrimSpace ( skeleton . ContentID ) == "" {
2026-07-24 18:18:49 +02:00
if updateErr := w . movieClaimer . UpdateError ( ctx , file . ID , job . LeaseToken , truncateSeriesQueueError ( "movie queue claimed without a content id" )); updateErr != nil {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to update movie queue error" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"file_id" , file . ID ,
"path" , file . FilePath ,
"error" , updateErr ,
)
}
return false
}
if skeleton . ItemStatus == "ambiguous" {
2026-07-24 18:18:49 +02:00
if err := w . movieClaimer . Delete ( ctx , file . ID , job . LeaseToken ); err != nil {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to delete ambiguous movie queue row" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"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 ())
2026-07-24 18:18:49 +02:00
if updateErr := w . movieClaimer . UpdateError ( ctx , file . ID , job . LeaseToken , queueErr ); updateErr != nil {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to update movie queue error" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"file_id" , file . ID ,
"path" , file . FilePath ,
"error" , updateErr ,
)
}
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: enrichment failed" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"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 {
2026-07-24 18:18:49 +02:00
if updateErr := w . updateMovieFailure ( ctx , file . ID , job . LeaseToken , result . Decision ); updateErr != nil {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to update movie queue error" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"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
}
2026-06-05 19:43:20 -07:00
w . publishCatalogItemChanged ( ctx , file . MediaFolderID , resultContentID ( result , skeleton . ContentID ), "metadata_updated" )
2026-05-22 20:26:11 -04:00
}
2026-07-24 18:18:49 +02:00
if err := w . movieClaimer . Delete ( ctx , file . ID , job . LeaseToken ); err != nil {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to delete movie queue row" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"file_id" , file . ID ,
"path" , file . FilePath ,
"error" , err ,
)
return false
}
return true
}
2026-07-24 18:18:49 +02:00
func ( w * MatchWorker ) queuedMovieSkeleton ( ctx context . Context , file * models . MediaFile , allowMatched bool ) ( * skeletonResult , bool , error ) {
if skeleton , ok := w . reusableQueuedMovieSkeleton ( ctx , file , allowMatched ); ok {
2026-05-22 20:26:11 -04:00
return skeleton , true , nil
}
2026-07-24 18:02:52 +02:00
currentFile := reparseQueuedFileIdentity ( file , "movie" )
skeleton , err := w . service . createOrFindSkeleton ( ctx , currentFile , currentFile . MediaFolderID )
2026-05-22 20:26:11 -04:00
if err != nil {
return nil , false , err
}
return skeleton , false , nil
}
2026-07-24 18:18:49 +02:00
func ( w * MatchWorker ) reusableQueuedMovieSkeleton ( ctx context . Context , file * models . MediaFile , allowMatched bool ) ( * skeletonResult , bool ) {
2026-05-22 20:26:11 -04:00
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 ))
2026-07-24 18:18:49 +02:00
reusableStatus := isSkeletonLikeStatus ( status ) ||
status == "ambiguous" ||
( allowMatched && status == string ( MatchOutcomeMatched ))
if ! reusableStatus {
2026-05-22 20:26:11 -04:00
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
}
2026-07-24 18:02:52 +02:00
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
}
}
2026-05-22 20:26:11 -04:00
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 ,
2026-07-24 18:02:52 +02:00
TmdbID : tmdbID ,
ImdbID : imdbID ,
TvdbID : tvdbID ,
2026-05-22 20:26:11 -04:00
}, true
}
2026-07-24 18:02:52 +02:00
// 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
}
2026-05-22 20:26:11 -04:00
func scopedMatcherPath ( useSeriesQueue bool , useMovieQueue bool ) string {
switch {
2026-05-26 20:07:43 -04:00
case useSeriesQueue && useMovieQueue :
return "series_root_queue+movie_file_queue"
2026-05-22 20:26:11 -04:00
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
}
2026-07-24 18:18:49 +02:00
defer w . releaseSeriesLeases ( jobs )
2026-05-22 20:26:11 -04:00
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
)
2026-06-10 19:25:07 -04:00
for i := 0 ; i < w . workerCount (); i ++ {
2026-05-22 20:26:11 -04:00
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
}
2026-07-24 18:18:49 +02:00
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 )
}
}
}
2026-05-22 20:26:11 -04:00
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 {
2026-07-24 18:18:49 +02:00
if err := w . seriesClaimer . Delete ( ctx , job . MediaFolderID , job . ObservedRootPath , job . LeaseToken ); err != nil {
2026-05-22 20:26:11 -04:00
return 0 , err
}
2026-07-09 20:53:52 +08:00
slog . InfoContext ( ctx , "metadata: series root dropped because files disappeared" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"folder_id" , job . MediaFolderID ,
"observed_root_path" , job . ObservedRootPath ,
)
return 0 , nil
}
representative := selectRepresentativeGroupFile ( groupFiles )
if representative == nil {
2026-07-24 18:18:49 +02:00
if err := w . seriesClaimer . Delete ( ctx , job . MediaFolderID , job . ObservedRootPath , job . LeaseToken ); err != nil {
2026-05-22 20:26:11 -04:00
return 0 , err
}
return 0 , nil
}
if ! hasUnlinkedGroupFile ( groupFiles ) {
if strings . TrimSpace ( representative . ContentID ) != "" {
2026-07-24 18:18:49 +02:00
if skeleton , ok := w . reusableQueuedMovieSkeleton ( ctx , representative , job . RerunRequested ); ok && skeleton . ItemStatus != "ambiguous" {
2026-05-26 20:07:43 -04:00
req := w . buildProcessRequestForGroup ( ctx , representative , skeleton , groupFiles )
result , processErr := w . service . Process ( ctx , req )
if processErr != nil {
queueErr := truncateSeriesQueueError ( processErr . Error ())
2026-07-24 18:18:49 +02:00
if updateErr := w . seriesClaimer . UpdateError ( ctx , job . MediaFolderID , job . ObservedRootPath , job . LeaseToken , queueErr ); updateErr != nil {
2026-05-26 20:07:43 -04:00
return 0 , updateErr
}
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: enrichment failed" , "component" , "metadata" ,
2026-05-26 20:07:43 -04:00
"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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: fallback episode synthesis failed after series root enrichment error" , "component" , "metadata" ,
2026-05-26 20:07:43 -04:00
"content_id" , skeleton . ContentID , "error" , fbErr )
}
return 0 , nil
}
if result != nil && ! result . Updated {
2026-07-24 18:18:49 +02:00
if updateErr := w . updateSeriesFailure ( ctx , job . MediaFolderID , job . ObservedRootPath , job . LeaseToken , result . Decision ); updateErr != nil {
2026-05-26 20:07:43 -04:00
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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: fallback episode synthesis failed for unmatched series root" , "component" , "metadata" ,
2026-05-26 20:07:43 -04:00
"content_id" , skeleton . ContentID , "error" , fbErr )
}
return 0 , nil
}
2026-06-05 19:43:20 -07:00
w . publishCatalogItemChanged ( ctx , job . MediaFolderID , resultContentID ( result , skeleton . ContentID ), "metadata_updated" )
2026-05-26 20:07:43 -04:00
}
2026-05-22 20:26:11 -04:00
if err := w . service . ensureSeriesEpisodeLinks ( ctx , representative . ContentID ); err != nil {
2026-05-27 21:01:28 +02:00
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.
2026-07-09 20:53:52 +08:00
slog . InfoContext ( ctx , "metadata: series item gone during episode-link ensure (likely concurrent merge); skipping" , "component" , "metadata" ,
2026-05-27 21:01:28 +02:00
"content_id" , representative . ContentID ,
"folder_id" , job . MediaFolderID ,
"observed_root_path" , job . ObservedRootPath )
} else {
2026-07-24 18:18:49 +02:00
if updateErr := w . seriesClaimer . UpdateError ( ctx , job . MediaFolderID , job . ObservedRootPath , job . LeaseToken , truncateSeriesQueueError ( err . Error ())); updateErr != nil {
2026-05-27 21:01:28 +02:00
return 0 , updateErr
}
return 0 , fmt . Errorf ( "ensuring series episode links for %s: %w" , representative . ContentID , err )
2026-05-22 20:26:11 -04:00
}
}
}
2026-07-24 18:18:49 +02:00
if err := w . seriesClaimer . Delete ( ctx , job . MediaFolderID , job . ObservedRootPath , job . LeaseToken ); err != nil {
2026-05-22 20:26:11 -04:00
return 0 , err
}
return len ( groupFiles ), nil
}
2026-07-24 18:02:52 +02:00
currentRepresentative := reparseQueuedFileIdentity ( representative , "series" )
skeleton , err := w . service . createOrFindSkeleton ( ctx , currentRepresentative , job . MediaFolderID )
2026-05-22 20:26:11 -04:00
if err != nil {
queueErr := truncateSeriesQueueError ( err . Error ())
2026-07-24 18:18:49 +02:00
if updateErr := w . seriesClaimer . UpdateError ( ctx , job . MediaFolderID , job . ObservedRootPath , job . LeaseToken , queueErr ); updateErr != nil {
2026-05-22 20:26:11 -04:00
return 0 , updateErr
}
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: series root skeleton creation failed" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"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 {
2026-07-24 18:18:49 +02:00
if updateErr := w . seriesClaimer . UpdateError ( ctx , job . MediaFolderID , job . ObservedRootPath , job . LeaseToken , truncateSeriesQueueError ( err . Error ())); updateErr != nil {
2026-05-22 20:26:11 -04:00
return 0 , updateErr
}
return 0 , fmt . Errorf ( "relinking series root %d/%s: %w" , job . MediaFolderID , job . ObservedRootPath , err )
}
2026-07-24 18:18:49 +02:00
needsInitialMatch := skeleton . IsNew || job . RerunRequested
2026-07-24 18:02:52 +02:00
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 ())
2026-07-24 18:18:49 +02:00
if updateErr := w . seriesClaimer . UpdateError ( ctx , job . MediaFolderID , job . ObservedRootPath , job . LeaseToken , queueErr ); updateErr != nil {
2026-07-24 18:02:52 +02:00
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" {
2026-05-22 20:26:11 -04:00
req := w . buildProcessRequestForGroup ( ctx , representative , skeleton , groupFiles )
result , processErr := w . service . Process ( ctx , req )
if processErr != nil {
queueErr := truncateSeriesQueueError ( processErr . Error ())
2026-07-24 18:18:49 +02:00
if updateErr := w . seriesClaimer . UpdateError ( ctx , job . MediaFolderID , job . ObservedRootPath , job . LeaseToken , queueErr ); updateErr != nil {
2026-05-22 20:26:11 -04:00
return 0 , updateErr
}
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: enrichment failed" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: fallback episode synthesis failed after enrichment error" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"content_id" , skeleton . ContentID ,
"error" , fbErr ,
)
}
return 0 , nil
} else if result != nil && ! result . Updated {
2026-07-24 18:18:49 +02:00
if updateErr := w . updateSeriesFailure ( ctx , job . MediaFolderID , job . ObservedRootPath , job . LeaseToken , result . Decision ); updateErr != nil {
2026-05-22 20:26:11 -04:00
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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: fallback episode synthesis failed for unmatched series" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"content_id" , skeleton . ContentID ,
"error" , fbErr ,
)
}
return 0 , nil
}
2026-06-05 19:43:20 -07:00
w . publishCatalogItemChanged ( ctx , job . MediaFolderID , resultContentID ( result , skeleton . ContentID ), "metadata_updated" )
2026-05-22 20:26:11 -04:00
}
finalContentID , err := w . service . fileRepo . FindContentIDByObservedRootPath ( ctx , job . MediaFolderID , job . ObservedRootPath , "series" )
if err != nil {
2026-07-24 18:18:49 +02:00
if updateErr := w . seriesClaimer . UpdateError ( ctx , job . MediaFolderID , job . ObservedRootPath , job . LeaseToken , truncateSeriesQueueError ( err . Error ())); updateErr != nil {
2026-05-22 20:26:11 -04:00
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 {
2026-05-27 21:01:28 +02:00
if errors . Is ( err , catalog . ErrItemNotFound ) {
2026-07-09 20:53:52 +08:00
slog . InfoContext ( ctx , "metadata: series item gone during episode-link ensure (likely concurrent merge); skipping" , "component" , "metadata" ,
2026-05-27 21:01:28 +02:00
"content_id" , finalContentID ,
"folder_id" , job . MediaFolderID ,
"observed_root_path" , job . ObservedRootPath )
} else {
2026-07-24 18:18:49 +02:00
if updateErr := w . seriesClaimer . UpdateError ( ctx , job . MediaFolderID , job . ObservedRootPath , job . LeaseToken , truncateSeriesQueueError ( err . Error ())); updateErr != nil {
2026-05-27 21:01:28 +02:00
return 0 , updateErr
}
return 0 , fmt . Errorf ( "ensuring series episode links for %s: %w" , finalContentID , err )
2026-05-22 20:26:11 -04:00
}
}
if _ , ok := w . service . confirmedOwnershipItem ( ctx , finalContentID ); ok {
w . service . claimConfirmedSeriesRootOwnership ( ctx , job . MediaFolderID , job . ObservedRootPath , finalContentID , groupFiles )
}
}
2026-07-24 18:18:49 +02:00
if err := w . seriesClaimer . Delete ( ctx , job . MediaFolderID , job . ObservedRootPath , job . LeaseToken ); err != nil {
2026-05-22 20:26:11 -04:00
return 0 , err
}
2026-07-09 20:53:52 +08:00
slog . InfoContext ( ctx , "metadata: series root completed" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"folder_id" , job . MediaFolderID ,
"observed_root_path" , job . ObservedRootPath ,
"content_id" , finalContentID ,
"file_count" , len ( groupFiles ),
)
return len ( groupFiles ), nil
}
2026-07-24 18:18:49 +02:00
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 ))
}
2026-05-22 20:26:11 -04:00
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 ... )
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to update item status" , append ([] any { "component" , "metadata" }, args ... ) ... )
2026-05-22 20:26:11 -04:00
}
}
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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to load folder type for batch compaction" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to load folder state during match" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"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
}
2026-05-26 20:07:43 -04:00
func ( w * MatchWorker ) queueUsageForFolder ( ctx context . Context , folderID int ) ( useSeriesQueue bool , useMovieQueue bool , err error ) {
2026-05-22 20:26:11 -04:00
folderType , err := w . folderType ( ctx , folderID )
if err != nil {
2026-05-26 20:07:43 -04:00
return false , false , err
2026-05-22 20:26:11 -04:00
}
2026-05-26 20:07:43 -04:00
useSeriesQueue = w . enableTVSeriesRootQueue &&
w . seriesClaimer != nil &&
2026-07-02 11:42:28 -04:00
( 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 ))
2026-05-26 20:07:43 -04:00
return useSeriesQueue , useMovieQueue , nil
2026-05-22 20:26:11 -04:00
}
func ( w * MatchWorker ) claimBackgroundFiles ( ctx context . Context ) ([] * models . MediaFile , error ) {
if w . enableTVSeriesRootQueue && w . movieClaimer != nil {
if claimer , ok := w . fileLister .( MixedFileClaimer ); ok {
2026-06-10 19:25:07 -04:00
return claimer . ClaimUnmatchedMixed ( ctx , w . claimBatchSize ())
2026-05-22 20:26:11 -04:00
}
return nil , fmt . Errorf ( "mixed-library file claimer is not configured" )
}
if w . enableTVSeriesRootQueue {
if claimer , ok := w . fileLister .( NonSeriesFileClaimer ); ok {
2026-06-10 19:25:07 -04:00
return claimer . ClaimUnmatchedNonSeries ( ctx , w . claimBatchSize ())
2026-05-22 20:26:11 -04:00
}
return nil , fmt . Errorf ( "non-series file claimer is not configured" )
}
2026-06-10 19:25:07 -04:00
return w . fileLister . ClaimUnmatched ( ctx , w . claimBatchSize ())
2026-05-22 20:26:11 -04:00
}
2026-05-26 20:07:43 -04:00
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 {
2026-06-10 19:25:07 -04:00
return claimer . ClaimUnmatchedMixedByFolderAndPathPrefix ( ctx , folderID , pathPrefix , w . claimBatchSize (), attemptBefore )
2026-05-26 20:07:43 -04:00
}
return nil , fmt . Errorf ( "mixed-library file claimer is not configured" )
case scopedFallbackNonSeries :
2026-05-22 20:26:11 -04:00
if claimer , ok := w . fileLister .( NonSeriesFileClaimer ); ok {
2026-06-10 19:25:07 -04:00
return claimer . ClaimUnmatchedNonSeriesByFolderAndPathPrefix ( ctx , folderID , pathPrefix , w . claimBatchSize (), attemptBefore )
2026-05-22 20:26:11 -04:00
}
return nil , fmt . Errorf ( "non-series file claimer is not configured" )
2026-05-26 20:07:43 -04:00
default :
2026-06-10 19:25:07 -04:00
return w . fileLister . ClaimUnmatchedByFolderAndPathPrefix ( ctx , folderID , pathPrefix , w . claimBatchSize (), attemptBefore )
2026-05-22 20:26:11 -04:00
}
}
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 ++
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: scoped retry failed" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"content_id" , contentID ,
"folder_id" , folderID ,
"path_prefix" , pathPrefix ,
"error" , processErr )
continue
}
if result == nil || ! result . Updated {
stillUnmatched ++
2026-06-05 19:43:20 -07:00
continue
2026-05-22 20:26:11 -04:00
}
2026-06-05 19:43:20 -07:00
w . publishCatalogItemChanged ( ctx , folderID , resultContentID ( result , contentID ), "metadata_updated" )
2026-05-22 20:26:11 -04:00
}
if repaired , repairErr := w . service . repairMatchedDuplicateProviderOwnersByFolderAndPathPrefix ( ctx , folderID , pathPrefix ); repairErr != nil {
return retried , stillUnmatched , repairErr
} else if repaired > 0 {
retried += repaired
2026-07-09 20:53:52 +08:00
slog . InfoContext ( ctx , "metadata: repaired matched duplicate items in scope" , "component" , "metadata" ,
2026-05-22 20:26:11 -04:00
"folder_id" , folderID ,
"path_prefix" , pathPrefix ,
"repaired_items" , repaired ,
)
}
2026-07-23 15:06:32 -04:00
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 ,
)
}
2026-05-22 20:26:11 -04:00
return retried , stillUnmatched , nil
}
func formatFolderID ( id int ) string {
return strconv . Itoa ( id )
}
2026-06-05 19:43:20 -07:00
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 {
2026-07-09 20:53:52 +08:00
slog . WarnContext ( ctx , "metadata: failed to publish catalog item change" , "component" , "metadata" ,
2026-06-05 19:43:20 -07:00
"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
}