mirror of
https://github.com/navidrome/navidrome.git
synced 2026-10-10 03:17:27 +02:00
* refactor(persistence): adopt generic deluan/rest repository API Pin deluan/rest to the refactor branch. REST-facing repository methods take a context and return typed values. Drop DataStore.Resource and ResourceRepository; the native API names typed repositories directly through a per-request adapter that later commits remove. * refactor(persistence): base repository helpers take a context * refactor(persistence): LibraryRepository takes a context per call * refactor(persistence): PropertyRepository takes a context per call * refactor(persistence): UserPropsRepository takes a context per call * refactor(persistence): TranscodingRepository takes a context per call * refactor(persistence): ShareRepository takes a context per call * refactor(persistence): PlayerRepository takes a context per call * refactor(persistence): RadioRepository takes a context per call * refactor(persistence): PlayQueueRepository takes a context per call * refactor(persistence): Tag and Genre repositories take a context per call * refactor(persistence): PluginRepository takes a context per call * refactor(persistence): Scrobble repositories take a context per call * refactor(persistence): FolderRepository takes a context per call * refactor(persistence): Artwork repositories take a context per call * refactor(persistence): UserRepository takes a context per call * refactor(persistence): ArtistRepository takes a context per call ReadAll no longer rewrites the shared sort mappings for the role filter; it works on a per-call copy. * test(persistence): assert artist role sort sanitization in ReadAll * refactor(persistence): AlbumRepository takes a context per call * test(persistence): pass the test context to album repository helpers * refactor(persistence): MediaFileRepository takes a context per call * refactor(persistence): Playlist repositories take a context per call * refactor(persistence): build all repositories once per store * refactor(core): REST repository wrappers are built once * refactor(persistence): repositories are stateless Remove the context field from the base repository and the per-request REST adapter. Enable the containedctx linter so no repository can hold a request context again. * chore(lint): skip containedctx in test files * refactor: share simplifications from the stateless repositories sweep Add deleteOwnedAll on sqlRepository and use it in player/share Delete to remove the duplicated bulk-delete loop; have Share.Repository() return model.ShareRepository so subsonic sharing.go drops its repeated type assertions. * chore(core): assert REST wrappers implement Persistable * chore: reformat imports * perf(persistence): build repositories on first use Each transaction store used to construct all 21 repositories up front, paying for filter and sort mapping setup the block never touched. Fields are now sync.OnceValue thunks, so a store only builds what it uses. * fix(persistence): clean plugin references per deleted user A bulk user delete that fails on a later id had already removed the earlier rows but skipped their plugin cleanup. Cleanup now runs right after each successful delete. * fix(core): unload disabled plugins even when a user delete fails A bulk delete can fail on a later id after earlier users were removed and their plugins auto-disabled. The wrapper returned before unloading, leaving those plugins running until the next successful delete or a restart. * chore(deps): pin deluan/rest to v1.0.1 Replaces the pseudo-version of the refactor branch with the tagged release. REST error messages now name the bare type (Artist, not model.Artist). * test: use the spec context instead of context.Background() Replace the context.Background()/context.TODO() calls this branch added to tests with the spec's ctx, GinkgoT().Context(), or t/b.Context(), so repository calls are bound to the running spec's lifetime. * test: declare the spec context once per Describe Set ctx from GinkgoT().Context() first in each top-level BeforeEach and reuse it, building user contexts on top of it instead of repeating inline calls.
412 lines
14 KiB
Go
412 lines
14 KiB
Go
package scanner
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"maps"
|
|
"path/filepath"
|
|
"slices"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
ppl "github.com/google/go-pipeline/pkg/pipeline"
|
|
"github.com/navidrome/navidrome/conf"
|
|
"github.com/navidrome/navidrome/consts"
|
|
"github.com/navidrome/navidrome/core/playlists"
|
|
"github.com/navidrome/navidrome/log"
|
|
"github.com/navidrome/navidrome/model"
|
|
"github.com/navidrome/navidrome/utils/run"
|
|
"github.com/navidrome/navidrome/utils/slice"
|
|
)
|
|
|
|
type scannerImpl struct {
|
|
ds model.DataStore
|
|
pls playlists.Playlists
|
|
}
|
|
|
|
// scanState holds the state of an in-progress scan, to be passed to the various phases
|
|
type scanState struct {
|
|
progress chan<- *ProgressInfo
|
|
fullScan bool
|
|
changesDetected atomic.Bool
|
|
libraries model.Libraries // Store libraries list for consistency across phases
|
|
targets map[int][]string // Optional: map[libraryID][]folderPaths for selective scans
|
|
totalLibraryCount int // Total number of libraries (unfiltered), for cross-library move detection
|
|
}
|
|
|
|
func (s *scanState) sendProgress(info *ProgressInfo) {
|
|
if s.progress != nil {
|
|
s.progress <- info
|
|
}
|
|
}
|
|
|
|
func (s *scanState) isSelectiveScan() bool {
|
|
return len(s.targets) > 0
|
|
}
|
|
|
|
func (s *scanState) sendWarning(msg string) {
|
|
s.sendProgress(&ProgressInfo{Warning: msg})
|
|
}
|
|
|
|
func (s *scanState) sendError(err error) {
|
|
s.sendProgress(&ProgressInfo{Error: err.Error()})
|
|
}
|
|
|
|
// libraryRelativePath rebases an absolute scan target path onto the library root, since the
|
|
// scanner's fs.FS only accepts paths relative to it. Relative paths, and absolute paths outside
|
|
// the library root, are returned unchanged.
|
|
func libraryRelativePath(libPath, folderPath string) string {
|
|
if !filepath.IsAbs(folderPath) {
|
|
return folderPath
|
|
}
|
|
// The library root may be relative (e.g. the default "./music"); it must be made absolute
|
|
// to match against an absolute target, and it resolves against the same cwd as the scanner's fs.
|
|
absLib, err := filepath.Abs(libPath)
|
|
if err != nil {
|
|
return folderPath
|
|
}
|
|
rel, err := filepath.Rel(absLib, folderPath)
|
|
if err != nil || !filepath.IsLocal(rel) {
|
|
return folderPath
|
|
}
|
|
// The scanner's fs.FS is an io/fs, which always uses forward slashes.
|
|
return filepath.ToSlash(rel)
|
|
}
|
|
|
|
func (s *scannerImpl) scanFolders(ctx context.Context, fullScan bool, targets []model.ScanTarget, progress chan<- *ProgressInfo) {
|
|
startTime := time.Now()
|
|
|
|
state := scanState{
|
|
progress: progress,
|
|
fullScan: fullScan,
|
|
changesDetected: atomic.Bool{},
|
|
}
|
|
|
|
// Set changesDetected to true for full scans to ensure all maintenance operations run
|
|
if fullScan {
|
|
state.changesDetected.Store(true)
|
|
}
|
|
|
|
// Get libraries and optionally filter by targets
|
|
allLibs, err := s.ds.Library().GetAll(ctx)
|
|
if err != nil {
|
|
state.sendWarning(fmt.Sprintf("getting libraries: %s", err))
|
|
return
|
|
}
|
|
state.totalLibraryCount = len(allLibs)
|
|
|
|
if len(targets) > 0 {
|
|
// Selective scan: filter libraries and build targets map
|
|
state.targets = make(map[int][]string)
|
|
|
|
libPaths := slice.ToMap(allLibs, func(lib model.Library) (int, string) {
|
|
return lib.ID, lib.Path
|
|
})
|
|
|
|
for _, target := range targets {
|
|
folderPath := libraryRelativePath(libPaths[target.LibraryID], target.FolderPath)
|
|
if folderPath == "" {
|
|
folderPath = "."
|
|
}
|
|
state.targets[target.LibraryID] = append(state.targets[target.LibraryID], folderPath)
|
|
}
|
|
|
|
// Filter libraries to only those in targets
|
|
state.libraries = slice.Filter(allLibs, func(lib model.Library) bool {
|
|
return len(state.targets[lib.ID]) > 0
|
|
})
|
|
|
|
log.Info(ctx, "Scanner: Starting selective scan", "fullScan", state.fullScan, "numLibraries", len(state.libraries), "numTargets", len(targets))
|
|
} else {
|
|
// Full library scan
|
|
state.libraries = allLibs
|
|
log.Info(ctx, "Scanner: Starting scan", "fullScan", state.fullScan, "numLibraries", len(state.libraries))
|
|
}
|
|
|
|
// Store scan type and start time
|
|
scanType := "quick"
|
|
if state.fullScan {
|
|
scanType = "full"
|
|
}
|
|
if state.isSelectiveScan() {
|
|
scanType += "-selective"
|
|
}
|
|
_ = s.ds.Property().Put(ctx, consts.LastScanTypeKey, scanType)
|
|
_ = s.ds.Property().Put(ctx, consts.LastScanStartTimeKey, startTime.Format(time.RFC3339))
|
|
|
|
// if there was a full scan in progress, force a full scan
|
|
if !state.fullScan {
|
|
for _, lib := range state.libraries {
|
|
if lib.FullScanInProgress {
|
|
log.Info(ctx, "Scanner: Interrupted full scan detected", "lib", lib.Name)
|
|
state.fullScan = true
|
|
if state.isSelectiveScan() {
|
|
_ = s.ds.Property().Put(ctx, consts.LastScanTypeKey, "full-selective")
|
|
} else {
|
|
_ = s.ds.Property().Put(ctx, consts.LastScanTypeKey, "full")
|
|
}
|
|
break
|
|
}
|
|
}
|
|
}
|
|
|
|
// Prepare libraries for scanning (initialize LastScanStartedAt if needed)
|
|
err = s.prepareLibrariesForScan(ctx, &state)
|
|
if err != nil {
|
|
log.Error(ctx, "Scanner: Error preparing libraries for scan", err)
|
|
state.sendError(err)
|
|
return
|
|
}
|
|
|
|
err = run.Sequentially(
|
|
// Phase 1: Scan all libraries and import new/updated files
|
|
runPhase[*folderEntry](ctx, 1, createPhaseFolders(ctx, &state, s.ds)),
|
|
|
|
// Phase 2: Process missing files, checking for moves
|
|
runPhase[*missingTracks](ctx, 2, createPhaseMissingTracks(ctx, &state, s.ds)),
|
|
|
|
// Phases 3 and 4 can be run in parallel
|
|
run.Parallel(
|
|
// Phase 3: Refresh all new/changed albums and update artists
|
|
runPhase[*model.Album](ctx, 3, createPhaseRefreshAlbums(ctx, &state, s.ds)),
|
|
|
|
// Phase 4: Import/update playlists
|
|
runPhase[*model.Folder](ctx, 4, createPhasePlaylists(ctx, &state, s.ds, s.pls)),
|
|
),
|
|
|
|
// Final Steps (cannot be parallelized):
|
|
|
|
// Run GC if there were any changes (Remove dangling tracks, empty albums and artists, and orphan annotations)
|
|
s.runGC(ctx, &state),
|
|
|
|
// Queue artwork for entities that never resolved (after GC, so nothing dangling is queued)
|
|
s.runEnqueueMissingArtwork(ctx, &state),
|
|
|
|
// Refresh artist and tags stats
|
|
s.runRefreshStats(ctx, &state),
|
|
|
|
// Update last_scan_completed_at for all libraries
|
|
s.runUpdateLibraries(ctx, &state),
|
|
)
|
|
if err != nil {
|
|
log.Error(ctx, "Scanner: Finished with error", "duration", time.Since(startTime), err)
|
|
_ = s.ds.Property().Put(ctx, consts.LastScanErrorKey, err.Error())
|
|
state.sendError(err)
|
|
return
|
|
}
|
|
|
|
_ = s.ds.Property().Put(ctx, consts.LastScanErrorKey, "")
|
|
|
|
if state.changesDetected.Load() {
|
|
state.sendProgress(&ProgressInfo{ChangesDetected: true})
|
|
}
|
|
|
|
if state.isSelectiveScan() {
|
|
log.Info(ctx, "Scanner: Finished scanning selected folders", "duration", time.Since(startTime), "numTargets", len(targets))
|
|
} else {
|
|
log.Info(ctx, "Scanner: Finished scanning all libraries", "duration", time.Since(startTime))
|
|
}
|
|
}
|
|
|
|
// prepareLibrariesForScan initializes the scan for all libraries in the state.
|
|
// It calls ScanBegin for libraries that haven't started scanning yet (LastScanStartedAt is zero),
|
|
// reloads them to get the updated state, and filters out any libraries that fail to initialize.
|
|
func (s *scannerImpl) prepareLibrariesForScan(ctx context.Context, state *scanState) error {
|
|
var successfulLibs []model.Library
|
|
|
|
for _, lib := range state.libraries {
|
|
if lib.LastScanStartedAt.IsZero() {
|
|
// This is a new scan - mark it as started
|
|
err := s.ds.WithTxRetry(ctx, func(ctx context.Context, tx model.DataStore) error {
|
|
return tx.Library().ScanBegin(ctx, lib.ID, state.fullScan)
|
|
}, "scanner: begin library scan")
|
|
if err != nil {
|
|
log.Error(ctx, "Scanner: Error marking scan start", "lib", lib.Name, err)
|
|
state.sendWarning(err.Error())
|
|
continue
|
|
}
|
|
|
|
// Reload library to get updated state (timestamps, etc.)
|
|
reloadedLib, err := s.ds.Library().Get(ctx, lib.ID)
|
|
if err != nil {
|
|
log.Error(ctx, "Scanner: Error reloading library", "lib", lib.Name, err)
|
|
state.sendWarning(err.Error())
|
|
continue
|
|
}
|
|
lib = *reloadedLib
|
|
} else {
|
|
// This is a resumed scan
|
|
log.Debug(ctx, "Scanner: Resuming previous scan", "lib", lib.Name,
|
|
"lastScanStartedAt", lib.LastScanStartedAt, "fullScan", lib.FullScanInProgress)
|
|
}
|
|
|
|
successfulLibs = append(successfulLibs, lib)
|
|
}
|
|
|
|
if len(successfulLibs) == 0 {
|
|
return fmt.Errorf("no libraries available for scanning")
|
|
}
|
|
|
|
// Update state with only successfully initialized libraries
|
|
state.libraries = successfulLibs
|
|
return nil
|
|
}
|
|
|
|
func (s *scannerImpl) runGC(ctx context.Context, state *scanState) func() error {
|
|
return func() error {
|
|
state.sendProgress(&ProgressInfo{ForceUpdate: true})
|
|
return s.ds.WithTxRetry(ctx, func(ctx context.Context, tx model.DataStore) error {
|
|
if state.changesDetected.Load() {
|
|
start := time.Now()
|
|
|
|
// For selective scans, extract library IDs to scope GC operations
|
|
var libraryIDs []int
|
|
if state.isSelectiveScan() {
|
|
libraryIDs = slices.Collect(maps.Keys(state.targets))
|
|
log.Debug(ctx, "Scanner: Running selective GC", "libraryIDs", libraryIDs)
|
|
}
|
|
|
|
if err := tx.GC(ctx, libraryIDs...); err != nil {
|
|
return fmt.Errorf("running GC: %w", err)
|
|
}
|
|
log.Debug(ctx, "Scanner: GC completed", "elapsed", time.Since(start))
|
|
} else {
|
|
log.Debug(ctx, "Scanner: No changes detected, skipping GC")
|
|
}
|
|
return nil
|
|
}, "scanner: GC")
|
|
}
|
|
}
|
|
|
|
// runEnqueueMissingArtwork is the safety net for entities phase 1 never enqueued.
|
|
func (s *scannerImpl) runEnqueueMissingArtwork(ctx context.Context, state *scanState) func() error {
|
|
return func() error {
|
|
if !state.changesDetected.Load() {
|
|
log.Debug(ctx, "Scanner: No changes detected, skipping artwork enqueue")
|
|
return nil
|
|
}
|
|
start := time.Now()
|
|
var total int64
|
|
for _, kind := range []model.Kind{model.KindAlbumArtwork, model.KindArtistArtwork} {
|
|
var n int64
|
|
err := s.ds.WithTxRetry(ctx, func(ctx context.Context, tx model.DataStore) error {
|
|
var err error
|
|
n, err = tx.ArtworkQueue().EnqueueAllMissing(ctx, kind, model.ArtworkPriorityScan)
|
|
return err
|
|
}, "scanner: enqueue missing artwork")
|
|
if err != nil {
|
|
log.Error(ctx, "Scanner: Error enqueueing missing artwork", "kind", kind, err)
|
|
return fmt.Errorf("enqueueing missing artwork: %w", err)
|
|
}
|
|
total += n
|
|
}
|
|
log.Debug(ctx, "Scanner: Enqueued missing artwork", "items", total, "elapsed", time.Since(start))
|
|
return nil
|
|
}
|
|
}
|
|
|
|
func (s *scannerImpl) runRefreshStats(ctx context.Context, state *scanState) func() error {
|
|
return func() error {
|
|
if !state.changesDetected.Load() {
|
|
log.Debug(ctx, "Scanner: No changes detected, skipping refreshing stats")
|
|
return nil
|
|
}
|
|
start := time.Now()
|
|
stats, err := s.ds.Artist().RefreshStats(ctx, state.fullScan)
|
|
if err != nil {
|
|
log.Error(ctx, "Scanner: Error refreshing artists stats", err)
|
|
return fmt.Errorf("refreshing artists stats: %w", err)
|
|
}
|
|
log.Debug(ctx, "Scanner: Refreshed artist stats", "stats", stats, "elapsed", time.Since(start))
|
|
|
|
start = time.Now()
|
|
err = s.ds.WithTxRetry(ctx, func(ctx context.Context, tx model.DataStore) error {
|
|
return tx.Tag().UpdateCounts(ctx)
|
|
}, "scanner: update tag counts")
|
|
if err != nil {
|
|
log.Error(ctx, "Scanner: Error updating tag counts", err)
|
|
return fmt.Errorf("updating tag counts: %w", err)
|
|
}
|
|
log.Debug(ctx, "Scanner: Updated tag counts", "elapsed", time.Since(start))
|
|
return nil
|
|
}
|
|
}
|
|
|
|
func (s *scannerImpl) runUpdateLibraries(ctx context.Context, state *scanState) func() error {
|
|
return func() error {
|
|
start := time.Now()
|
|
return s.ds.WithTxRetry(ctx, func(ctx context.Context, tx model.DataStore) error {
|
|
for _, lib := range state.libraries {
|
|
if err := tx.Library().ScanEnd(ctx, lib.ID); err != nil {
|
|
return fmt.Errorf("updating last scan completed for %s: %w", lib.Name, err)
|
|
}
|
|
if err := tx.Property().Put(ctx, consts.PIDTrackKey, conf.Server.PID.Track); err != nil {
|
|
return fmt.Errorf("updating track PID conf: %w", err)
|
|
}
|
|
if err := tx.Property().Put(ctx, consts.PIDAlbumKey, conf.Server.PID.Album); err != nil {
|
|
return fmt.Errorf("updating album PID conf: %w", err)
|
|
}
|
|
if state.changesDetected.Load() {
|
|
log.Debug(ctx, "Scanner: Refreshing library stats", "lib", lib.Name)
|
|
if err := tx.Library().RefreshStats(ctx, lib.ID); err != nil {
|
|
return fmt.Errorf("refreshing library stats for %s: %w", lib.Name, err)
|
|
}
|
|
} else {
|
|
log.Debug(ctx, "Scanner: No changes detected, skipping library stats refresh", "lib", lib.Name)
|
|
}
|
|
}
|
|
log.Debug(ctx, "Scanner: Updated libraries after scan", "elapsed", time.Since(start), "numLibraries", len(state.libraries))
|
|
return nil
|
|
}, "scanner: update libraries")
|
|
}
|
|
}
|
|
|
|
type phase[T any] interface {
|
|
producer() ppl.Producer[T]
|
|
stages() []ppl.Stage[T]
|
|
finalize(error) error
|
|
description() string
|
|
}
|
|
|
|
func runPhase[T any](ctx context.Context, phaseNum int, phase phase[T]) func() error {
|
|
return func() error {
|
|
log.Debug(ctx, fmt.Sprintf("Scanner: Starting phase %d: %s", phaseNum, phase.description()))
|
|
start := time.Now()
|
|
|
|
producer := phase.producer()
|
|
stages := phase.stages()
|
|
|
|
// Prepend a counter stage to the phase's pipeline
|
|
counter, countStageFn := countTasks[T]()
|
|
stages = append([]ppl.Stage[T]{ppl.NewStage(countStageFn, ppl.Name("count tasks"))}, stages...)
|
|
|
|
var err error
|
|
if log.IsGreaterOrEqualTo(log.LevelDebug) {
|
|
var m *ppl.Metrics
|
|
m, err = ppl.Measure(producer, stages...)
|
|
log.Info(ctx, "Scanner: "+m.String(), err)
|
|
} else {
|
|
err = ppl.Do(producer, stages...)
|
|
}
|
|
|
|
err = phase.finalize(err)
|
|
|
|
if err != nil {
|
|
log.Error(ctx, fmt.Sprintf("Scanner: Error processing libraries in phase %d", phaseNum), "elapsed", time.Since(start), err)
|
|
} else {
|
|
log.Debug(ctx, fmt.Sprintf("Scanner: Finished phase %d", phaseNum), "elapsed", time.Since(start), "totalTasks", counter.Load())
|
|
}
|
|
|
|
return err
|
|
}
|
|
}
|
|
|
|
func countTasks[T any]() (*atomic.Int64, func(T) (T, error)) {
|
|
counter := atomic.Int64{}
|
|
return &counter, func(in T) (T, error) {
|
|
counter.Add(1)
|
|
return in, nil
|
|
}
|
|
}
|
|
|
|
var _ scanner = (*scannerImpl)(nil)
|