diff options
| -rw-r--r-- | player-server/internal/scanner/file_discoverer.go | 75 | ||||
| -rw-r--r-- | player-server/internal/scanner/probe_worker.go | 185 | ||||
| -rw-r--r-- | player-server/internal/scanner/scan_writer.go | 56 | ||||
| -rw-r--r-- | player-server/internal/scanner/scanner.go | 242 |
4 files changed, 339 insertions, 219 deletions
diff --git a/player-server/internal/scanner/file_discoverer.go b/player-server/internal/scanner/file_discoverer.go new file mode 100644 index 0000000..bb77d61 --- /dev/null +++ b/player-server/internal/scanner/file_discoverer.go @@ -0,0 +1,75 @@ +// Package scanner implements media library scanning logic. +package scanner + +import ( + "fmt" + "io/fs" + "path/filepath" + "strings" + + "codeberg.org/snonux/player/internal/mediatype" +) + +// fileDiscoverer walks a set directory and collects media file paths and +// cover image locations. It is responsible solely for filesystem traversal — +// it does not probe, store, or generate thumbnails. +type fileDiscoverer struct { + fs FS +} + +// newFileDiscoverer creates a fileDiscoverer backed by the given filesystem abstraction. +func newFileDiscoverer(fs FS) *fileDiscoverer { + return &fileDiscoverer{fs: fs} +} + +// Discover walks setPath and returns the absolute paths of all supported media +// files, skipping hidden directories (names starting with "."). +func (d *fileDiscoverer) Discover(setPath string) ([]string, error) { + var files []string + walkErr := d.fs.WalkDir(setPath, func(path string, entry fs.DirEntry, err error) error { + if err != nil { + return fmt.Errorf("walk %q: %w", path, err) + } + // Skip hidden directories (e.g. .git, .thumbnails) to avoid + // accidentally ingesting dot-prefixed paths. + if entry.IsDir() { + if strings.HasPrefix(entry.Name(), ".") && path != setPath { + return filepath.SkipDir + } + return nil + } + if !mediatype.IsSupportedExt(path) { + return nil + } + files = append(files, path) + return nil + }) + if walkErr != nil { + return nil, walkErr + } + return files, nil +} + +// gatherCoverImages walks the set directory and records the first cover image +// per directory. The returned map is keyed by absolute directory path; values +// are the image path relative to setPath so callers can join them as needed. +func (d *fileDiscoverer) gatherCoverImages(setPath string) map[string]string { + coverImages := make(map[string]string) + _ = d.fs.WalkDir(setPath, func(path string, entry fs.DirEntry, err error) error { + if err != nil || entry.IsDir() || !mediatype.IsCoverImageExt(path) { + // Skip hidden subdirectories, but tolerate walk errors silently + // since missing cover images are non-fatal. + if entry != nil && entry.IsDir() && strings.HasPrefix(entry.Name(), ".") && path != setPath { + return filepath.SkipDir + } + return nil + } + relPath, _ := filepath.Rel(setPath, path) + dir := filepath.Dir(path) + if _, ok := coverImages[dir]; !ok { + coverImages[dir] = relPath + } + return nil + }) + return coverImages +} diff --git a/player-server/internal/scanner/probe_worker.go b/player-server/internal/scanner/probe_worker.go new file mode 100644 index 0000000..10766b6 --- /dev/null +++ b/player-server/internal/scanner/probe_worker.go @@ -0,0 +1,185 @@ +// Package scanner implements media library scanning logic. +package scanner + +import ( + "context" + "fmt" + "log/slog" + "path/filepath" + "strings" + + "codeberg.org/snonux/player/internal/clock" + "codeberg.org/snonux/player/internal/mediatype" + "codeberg.org/snonux/player/internal/model" + "codeberg.org/snonux/player/internal/probe" + "codeberg.org/snonux/player/internal/thumb" +) + +// fileResult carries a successfully probed media record back to the scanWriter. +type fileResult struct { + media *model.Media + path string // absolute path used for logging +} + +// probeWorker probes individual media files via ffprobe and resolves thumbnail +// paths. It handles concurrency: multiple goroutines call run() in parallel, +// each reading from pathChan and writing probed fileResults to resultChan. +// All filesystem probing and thumbnail resolution happens here; no DB writes. +type probeWorker struct { + prober probe.Prober + thumbMkr thumb.Maker + fs FS + clock clock.Clock + logger *slog.Logger +} + +// newProbeWorker creates a probeWorker with the required dependencies. +func newProbeWorker(prober probe.Prober, maker thumb.Maker, fs FS, clk clock.Clock, logger *slog.Logger) *probeWorker { + return &probeWorker{ + prober: prober, + thumbMkr: maker, + fs: fs, + clock: clk, + logger: logger, + } +} + +// run consumes file paths from pathChan, probes each one, and sends the result +// to resultChan. It exits early when scanCtx is cancelled or sendErr is called. +// The caller is responsible for closing pathChan; run returns when pathChan is drained. +func (pw *probeWorker) run( + ctx context.Context, + scanCtx context.Context, + pathChan <-chan string, + resultChan chan<- fileResult, + setPath string, + setID int64, + setName string, + existing map[string]model.Media, + coverImages map[string]string, + progress *model.ScanProgress, + sendErr func(error), +) { + for path := range pathChan { + if scanCtx.Err() != nil { + continue + } + // Use scanCtx (not ctx) so ffprobe/ffmpeg subprocesses cancel promptly + // when another worker fails or TriggerRescan restarts the scan. + result, err := pw.probeFile(scanCtx, path, setPath, setID, setName, existing, coverImages, progress) + if err != nil { + sendErr(err) + return + } + if result == nil { + continue + } + select { + case resultChan <- *result: + case <-scanCtx.Done(): + return + } + } +} + +// probeFile probes a single file and builds a media record ready for persistence. +// Returns nil when the file already exists in existing or cannot be probed +// (unrecognised format — a warning is logged and the file is skipped). +func (pw *probeWorker) probeFile( + ctx context.Context, + path, setPath string, + setID int64, + setName string, + existing map[string]model.Media, + coverImages map[string]string, + progress *model.ScanProgress, +) (*fileResult, error) { + relPath, err := filepath.Rel(setPath, path) + if err != nil { + return nil, fmt.Errorf("rel path for %q: %w", path, err) + } + relPath = filepath.ToSlash(relPath) + + if progress != nil { + progress.IncrementFile() + } + + _, alreadyExists := existing[relPath] + pw.logger.Debug("scanner file checked", "set", setName, "path", relPath, "existing", alreadyExists) + if alreadyExists { + return nil, nil + } + + info, err := pw.fs.Stat(path) + if err != nil { + return nil, fmt.Errorf("stat %q: %w", path, err) + } + + meta, err := pw.prober.Probe(ctx, path) + if err != nil { + pw.logger.Warn("scanner skipping unprobeable file", "path", path, "err", err) + return nil, nil + } + meta.FileSizeBytes = info.Size() + + mediaType := mediatype.TypeForExt(path) + thumbnailPath, err := pw.buildThumbnailPath(ctx, path, setPath, mediaType, coverImages, meta) + if err != nil { + return nil, err + } + + media := &model.Media{ + SetID: setID, + RelPath: relPath, + FileName: filepath.Base(path), + AbsPath: path, + Type: mediaType, + Duration: meta.Duration, + Codec: meta.Codec, + Resolution: meta.Resolution, + Bitrate: meta.Bitrate, + FileSizeBytes: meta.FileSizeBytes, + Width: meta.Width, + Height: meta.Height, + EXIFCamera: meta.EXIFCamera, + EXIFLens: meta.EXIFLens, + EXIFDate: meta.EXIFDate, + EXIFISO: meta.EXIFISO, + EXIFFNumber: meta.EXIFFNumber, + EXIFExposure: meta.EXIFExposure, + EXIFFocalLength: meta.EXIFFocalLength, + ThumbnailPath: thumbnailPath, + CreatedAt: pw.clock.Now(), + } + + return &fileResult{media: media, path: path}, nil +} + +// buildThumbnailPath resolves the thumbnail path for a new media file. +// Video and image thumbnails are produced via thumb.Maker; audio uses a +// nearby cover image; SVG images are served as-is (no raster thumbnail needed). +func (pw *probeWorker) buildThumbnailPath(ctx context.Context, path, setPath string, mediaType model.MediaType, coverImages map[string]string, meta *model.Metadata) (string, error) { + switch mediaType { + case model.MediaTypeVideo: + return pw.thumbMkr.MakeVideo(ctx, path, setPath, meta.Duration) + case model.MediaTypeAudio: + return findCoverImage(path, coverImages, setPath), nil + case model.MediaTypeImage: + ext := strings.ToLower(filepath.Ext(path)) + if ext == ".svg" { + // SVG is a vector format; serve the original file directly. + return path, nil + } + thumbPath, err := pw.thumbMkr.MakeImage(ctx, path, setPath) + if err != nil { + return "", err + } + if thumbPath != "" { + if _, statErr := pw.fs.Stat(thumbPath); statErr == nil { + return thumbPath, nil + } + } + return path, nil + } + return "", nil +} diff --git a/player-server/internal/scanner/scan_writer.go b/player-server/internal/scanner/scan_writer.go new file mode 100644 index 0000000..aa4ad0c --- /dev/null +++ b/player-server/internal/scanner/scan_writer.go @@ -0,0 +1,56 @@ +// Package scanner implements media library scanning logic. +package scanner + +import ( + "context" + "fmt" + "log/slog" + "path/filepath" + "sync/atomic" + + "codeberg.org/snonux/player/internal/repository" +) + +// scanWriter persists probed media results to the database. It runs in a +// single goroutine to serialise SQLite writes and avoid concurrent-write +// errors. All probing and thumbnail generation happens in probeWorker before +// results arrive here. +type scanWriter struct { + store repository.ScannerStore + logger *slog.Logger +} + +// newScanWriter creates a scanWriter backed by the given store. +func newScanWriter(store repository.ScannerStore, logger *slog.Logger) *scanWriter { + return &scanWriter{store: store, logger: logger} +} + +// run reads fileResults from resultChan and inserts each into the store. +// It logs progress every 25 files and accumulates the count in newFiles. +// Exits when resultChan is closed or scanCtx is cancelled. +func (sw *scanWriter) run( + ctx context.Context, + scanCtx context.Context, + resultChan <-chan fileResult, + setName string, + setPath string, + newFiles *int32, + sendErr func(error), +) { + for result := range resultChan { + if scanCtx.Err() != nil { + continue + } + if _, err := sw.store.CreateMedia(ctx, result.media); err != nil { + sendErr(fmt.Errorf("create media %q: %w", result.path, err)) + continue + } + nf := atomic.AddInt32(newFiles, 1) + // Log progress on the first insert and every 25 thereafter to give + // operators visibility into long-running scans without flooding logs. + if nf == 1 || nf%25 == 0 { + relPath, _ := filepath.Rel(setPath, result.path) + sw.logger.Info("scanner set progress", "name", setName, "new_media", nf, "latest", filepath.ToSlash(relPath)) + } + } +} diff --git a/player-server/internal/scanner/scanner.go b/player-server/internal/scanner/scanner.go index a669a6b..a6f8fd7 100644 --- a/player-server/internal/scanner/scanner.go +++ b/player-server/internal/scanner/scanner.go @@ -4,16 +4,13 @@ package scanner import ( "context" "fmt" - "io/fs" "log/slog" "path/filepath" "runtime" "strings" "sync" - "sync/atomic" "codeberg.org/snonux/player/internal/clock" - "codeberg.org/snonux/player/internal/mediatype" "codeberg.org/snonux/player/internal/model" "codeberg.org/snonux/player/internal/probe" "codeberg.org/snonux/player/internal/repository" @@ -31,6 +28,11 @@ type Scanner interface { // directory creation, generator invocation, and failure-tolerant warning // to a thumb.Maker. This keeps SRP intact — FSScanner orchestrates the // scan, thumb.Maker decides how thumbnails get produced on disk. +// +// Scan itself is orchestrated by delegating to three focused collaborators: +// - fileDiscoverer: walks the filesystem to find media files and cover images +// - probeWorker: runs ffprobe in parallel workers to build media records +// - scanWriter: persists probed results to the database type FSScanner struct { store repository.ScannerStore prober probe.Prober @@ -88,13 +90,14 @@ func (s *FSScanner) log() *slog.Logger { } // Scan walks immediate subdirectories of root, treating each as a set. +// It orchestrates fileDiscoverer, probeWorker, and scanWriter collaborators. func (s *FSScanner) Scan(ctx context.Context, root string, progress *model.ScanProgress) error { entries, err := s.fs.ReadDir(root) if err != nil { return fmt.Errorf("read media root %q: %w", root, err) } - // Count total sets for progress. + // Count total sets for progress reporting. var setCount int for _, entry := range entries { if entry.IsDir() { @@ -205,126 +208,6 @@ func (s *FSScanner) reconcileOrphans(ctx context.Context, existing map[string]mo } } -// gatherCoverImages walks the set and records the first cover image per directory. -func (s *FSScanner) gatherCoverImages(setPath string) map[string]string { - coverImages := make(map[string]string) - _ = s.fs.WalkDir(setPath, func(path string, d fs.DirEntry, err error) error { - if err != nil || d.IsDir() || !mediatype.IsCoverImageExt(path) { - if d != nil && d.IsDir() && strings.HasPrefix(d.Name(), ".") && path != setPath { - return filepath.SkipDir - } - return nil - } - relPath, _ := filepath.Rel(setPath, path) - dir := filepath.Dir(path) - if _, ok := coverImages[dir]; !ok { - coverImages[dir] = relPath - } - return nil - }) - return coverImages -} - -// buildThumbnailPath resolves the thumbnail path for a new media file. -// Video and image thumbnails are produced via thumb.Maker so the scanner -// stays out of mkdir / path-derivation / generator policy. Audio uses a -// nearby cover image when one was discovered during gatherCoverImages. -// SVG images are served as-is (vector — no raster thumbnail makes sense). -func (s *FSScanner) buildThumbnailPath(ctx context.Context, path, setPath string, mediaType model.MediaType, coverImages map[string]string, meta *model.Metadata) (string, error) { - switch mediaType { - case model.MediaTypeVideo: - return s.thumbMkr.MakeVideo(ctx, path, setPath, meta.Duration) - case model.MediaTypeAudio: - return findCoverImage(path, coverImages, setPath), nil - case model.MediaTypeImage: - ext := strings.ToLower(filepath.Ext(path)) - if ext == ".svg" { - return path, nil - } - thumbPath, err := s.thumbMkr.MakeImage(ctx, path, setPath) - if err != nil { - return "", err - } - if thumbPath != "" { - if _, statErr := s.fs.Stat(thumbPath); statErr == nil { - return thumbPath, nil - } - } - return path, nil - } - return "", nil -} - -// fileResult carries a successfully probed media record back to the writer. -type fileResult struct { - media *model.Media - path string // absolute path for logging -} - -// probeFile probes a single file and builds a media record. -// It returns nil when the file already exists or is unprobeable. -func (s *FSScanner) probeFile(ctx context.Context, path, setPath string, setID int64, setName string, existing map[string]model.Media, coverImages map[string]string, progress *model.ScanProgress) (*fileResult, error) { - relPath, err := filepath.Rel(setPath, path) - if err != nil { - return nil, fmt.Errorf("rel path for %q: %w", path, err) - } - relPath = filepath.ToSlash(relPath) - - if progress != nil { - progress.IncrementFile() - } - - _, alreadyExists := existing[relPath] - s.log().Debug("scanner file checked", "set", setName, "path", relPath, "existing", alreadyExists) - if alreadyExists { - return nil, nil - } - - info, err := s.fs.Stat(path) - if err != nil { - return nil, fmt.Errorf("stat %q: %w", path, err) - } - - meta, err := s.prober.Probe(ctx, path) - if err != nil { - s.log().Warn("scanner skipping unprobeable file", "path", path, "err", err) - return nil, nil - } - meta.FileSizeBytes = info.Size() - - mediaType := mediatype.TypeForExt(path) - thumbnailPath, err := s.buildThumbnailPath(ctx, path, setPath, mediaType, coverImages, meta) - if err != nil { - return nil, err - } - - media := &model.Media{ - SetID: setID, - RelPath: relPath, - FileName: filepath.Base(path), - AbsPath: path, - Type: mediaType, - Duration: meta.Duration, - Codec: meta.Codec, - Resolution: meta.Resolution, - Bitrate: meta.Bitrate, - FileSizeBytes: meta.FileSizeBytes, - Width: meta.Width, - Height: meta.Height, - EXIFCamera: meta.EXIFCamera, - EXIFLens: meta.EXIFLens, - EXIFDate: meta.EXIFDate, - EXIFISO: meta.EXIFISO, - EXIFFNumber: meta.EXIFFNumber, - EXIFExposure: meta.EXIFExposure, - EXIFFocalLength: meta.EXIFFocalLength, - ThumbnailPath: thumbnailPath, - CreatedAt: s.clock.Now(), - } - - return &fileResult{media: media, path: path}, nil -} - // updateAudioThumbnails patches existing audio tracks when a new cover image appears. func (s *FSScanner) updateAudioThumbnails(ctx context.Context, mediaList []model.Media, coverImages map[string]string, setPath string) { for _, m := range mediaList { @@ -340,8 +223,9 @@ func (s *FSScanner) updateAudioThumbnails(ctx context.Context, mediaList []model } } -// scanSet scans a single set using a pool of workers for probing and a single -// writer goroutine for SQLite inserts. +// scanSet scans a single set using fileDiscoverer, probeWorker, and scanWriter. +// fileDiscoverer collects the file list; probeWorker probes files in parallel; +// scanWriter persists results to SQLite via a single writer goroutine. func (s *FSScanner) scanSet(ctx context.Context, root, setPath string, progress *model.ScanProgress) error { workers := s.workers if workers <= 0 { @@ -362,9 +246,10 @@ func (s *FSScanner) scanSet(ctx context.Context, root, setPath string, progress return err } - coverImages := s.gatherCoverImages(setPath) - - files, err := s.collectFiles(setPath) + // Use fileDiscoverer to collect files and cover images. + disc := newFileDiscoverer(s.fs) + coverImages := disc.gatherCoverImages(setPath) + files, err := disc.Discover(setPath) if err != nil { return fmt.Errorf("scan set %q: %w", setName, err) } @@ -384,7 +269,7 @@ func (s *FSScanner) scanSet(ctx context.Context, root, setPath string, progress } pathChan := make(chan string, len(files)) - resultChan := make(chan fileResult, s.workers) + resultChan := make(chan fileResult, workers) scanCtx, cancel := context.WithCancel(ctx) defer cancel() @@ -395,21 +280,25 @@ func (s *FSScanner) scanSet(ctx context.Context, root, setPath string, progress errOnce.Do(func() { errChan <- err; cancel() }) } + // probeWorker probes files concurrently and sends results to resultChan. + pw := newProbeWorker(s.prober, s.thumbMkr, s.fs, s.clock, s.log()) var workerWg sync.WaitGroup for i := 0; i < workers; i++ { workerWg.Add(1) go func() { defer workerWg.Done() - s.probeWorkerLoop(ctx, scanCtx, pathChan, resultChan, setPath, setID, setName, existing, coverImages, progress, sendErr) + pw.run(ctx, scanCtx, pathChan, resultChan, setPath, setID, setName, existing, coverImages, progress, sendErr) }() } + // scanWriter persists results sequentially to avoid SQLite write conflicts. + sw := newScanWriter(s.store, s.log()) var newFiles int32 var writerWg sync.WaitGroup writerWg.Add(1) go func() { defer writerWg.Done() - s.writerLoop(ctx, scanCtx, resultChan, setName, setPath, &newFiles, sendErr) + sw.run(ctx, scanCtx, resultChan, setName, setPath, &newFiles, sendErr) }() // Feed the worker pool. @@ -450,94 +339,9 @@ func (s *FSScanner) scanSet(ctx context.Context, root, setPath string, progress } // collectFiles walks the set and returns the absolute paths of all supported media files. +// Kept for backward compatibility with existing tests that call it directly. func (s *FSScanner) collectFiles(setPath string) ([]string, error) { - var files []string - walkErr := s.fs.WalkDir(setPath, func(path string, d fs.DirEntry, err error) error { - if err != nil { - return fmt.Errorf("walk %q: %w", path, err) - } - if d.IsDir() { - if strings.HasPrefix(d.Name(), ".") && path != setPath { - return filepath.SkipDir - } - return nil - } - if !mediatype.IsSupportedExt(path) { - return nil - } - files = append(files, path) - return nil - }) - if walkErr != nil { - return nil, walkErr - } - return files, nil -} - -// probeWorkerLoop consumes file paths, probes each one, and sends the result to resultChan. -// It stops early if scanCtx is cancelled or if sendErr reports a fatal error. -func (s *FSScanner) probeWorkerLoop( - ctx context.Context, - scanCtx context.Context, - pathChan <-chan string, - resultChan chan<- fileResult, - setPath string, - setID int64, - setName string, - existing map[string]model.Media, - coverImages map[string]string, - progress *model.ScanProgress, - sendErr func(error), -) { - for path := range pathChan { - if scanCtx.Err() != nil { - continue - } - // Use scanCtx (not ctx) so ffprobe/ffmpeg subprocesses spawned by - // probeFile cancel promptly when scanCtx is cancelled — e.g. another - // worker failed or TriggerRescan restarted the scan. Passing the - // parent ctx here would leave ffprobe running after cancel. - result, err := s.probeFile(scanCtx, path, setPath, setID, setName, existing, coverImages, progress) - if err != nil { - sendErr(err) - return - } - if result == nil { - continue - } - select { - case resultChan <- *result: - case <-scanCtx.Done(): - return - } - } -} - -// writerLoop reads probed results and inserts them into the store. -// It logs progress every 25 files and tracks the total newFiles count. -func (s *FSScanner) writerLoop( - ctx context.Context, - scanCtx context.Context, - resultChan <-chan fileResult, - setName string, - setPath string, - newFiles *int32, - sendErr func(error), -) { - for result := range resultChan { - if scanCtx.Err() != nil { - continue - } - if _, err := s.store.CreateMedia(ctx, result.media); err != nil { - sendErr(fmt.Errorf("create media %q: %w", result.path, err)) - continue - } - nf := atomic.AddInt32(newFiles, 1) - if nf == 1 || nf%25 == 0 { - relPath, _ := filepath.Rel(setPath, result.path) - s.log().Info("scanner set progress", "name", setName, "new_media", nf, "latest", filepath.ToSlash(relPath)) - } - } + return newFileDiscoverer(s.fs).Discover(setPath) } func findCoverImage(filePath string, coverImages map[string]string, setPath string) string { |
