Files
minstrel/internal/library/scanner.go
T
bvandeusen f6d1cf24f0
test-go / test (push) Successful in 1m0s
test-go / integration (push) Successful in 5m10s
feat(library): detect missing files and stop offering them — #2523
Nothing in Minstrel ever noticed a deleted file. The walk only visits
paths that exist, so a row whose file was gone was never scanned, never
errored, never counted — permanently invisible. classifyEvent ignores
fsnotify removals by design, and the safety-net scan is the same walk, so
it covers additions only. Rows accumulated forever.

Found on the operator's library: a completed scan reported
skipped=24185 errored=0 while the MBID backfill (which opens files by DB
path rather than walking) logged ~40 "no such file or directory" across
three reorganised albums. Those rows also kept their pre-#2499 welded
genre, which is how this surfaced — the version-stamped tag re-read can
only reach files the walk visits.

The harm is not cosmetic. tracks is the candidate universe for
recommendation.sql / discover.sql / system_mixes.sql and nothing filtered
on file existence, so a mix could spend a slot on a track that cannot
stream.

Marks rather than deletes. A missing file is a claim about the filesystem
and the filesystem lies transiently — an unmounted volume, a network
blip, a container that started before its media mount attached. Every
sweep in internal/gc resolves a truth INSIDE the database and is safe to
run blind; this one is not, so no deletion happens here. Three guards
refuse to act on ambiguous evidence: every scan root must resolve to a
non-empty directory, the walk must have seen at least one file, and one
reconcile may newly mark at most 25% of the library. Clearing a mark is
never the dangerous direction, so it runs unconditionally — otherwise a
library that tripped the cap could never recover once the mount returned.

Only a full Scan reconciles. The walk's set of seen paths is the
evidence, and ScanFiles has no basis for concluding anything about files
it did not look at.

Excludes marked tracks from all 13 track-emitting queries (radio x2,
system mixes x5, discover x4, most-played x2), the 6 play-history seed
picks, and the genre browse axis. Deliberately NOT filtered: the shared
ListPlaylistTracks read path, because it also serves user-curated
playlists where hiding a track the user added would be wrong — system
playlists shed orphans on their next daily rebuild instead. History and
the taste profile also keep them: those record the past, and a track you
played 200 times still says something about your taste.

Reconcile tallies land in scan_runs so a disappearance is visible rather
than discovered when a mix comes up short.
2026-08-06 14:34:53 -04:00

509 lines
17 KiB
Go

package library
import (
"context"
"errors"
"fmt"
"io/fs"
"log/slog"
"math"
"os"
"os/exec"
"path/filepath"
"strconv"
"strings"
"time"
"github.com/dhowden/tag"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgtype"
"github.com/jackc/pgx/v5/pgxpool"
"git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq"
syncpkg "git.fabledsword.com/bvandeusen/minstrel/internal/sync"
)
// audioExtensions is the set the scanner indexes. Keep in sync with
// `internal/api/media.go` MIME detection — the stream handler must be
// able to serve every extension the scanner indexes, and there is no
// point in adding extensions to the stream handler that the scanner
// will silently skip. Drift #571 caught the divergence after .opus,
// .aac, and .wav were added to media.go but not here.
var audioExtensions = map[string]bool{
".mp3": true,
".m4a": true,
".flac": true,
".ogg": true,
".opus": true,
".aac": true,
".wav": true,
}
// tagReadVersion is the version of this package's tag-extraction logic. Rows
// whose tracks.tag_read_version is lower get their tags re-read on the next
// scan even when the file itself hasn't changed, so a fix reaches an existing
// library without the operator rebuilding it (migration 0054).
//
// Bump this whenever a change to tag extraction should reach already-indexed
// files, and say why below.
//
// 1: genre read from the ID3v2 TCON frame directly and stored ";"-delimited.
// dhowden/tag welds null-separated multi-values into one token
// ("Alternative Rock" + "Rock" -> "Alternative RockRock"), which corrupted
// the genre browse axis and polluted the taste profile's tag vocabulary,
// and left bare ID3v1 numeric references unresolved (#2499).
const tagReadVersion int16 = 1
type Stats struct {
Scanned int `json:"scanned"`
Added int `json:"added"`
Updated int `json:"updated"`
Skipped int `json:"skipped"`
Errored int `json:"errored"`
// Missing / Restored come from the reconcile pass, not the walk (#2523):
// rows whose file the walk didn't find, and rows whose file came back.
// Only a full Scan sets these — see reconcileMissing.
Missing int `json:"missing"`
Restored int `json:"restored"`
}
type Scanner struct {
pool *pgxpool.Pool
logger *slog.Logger
paths []string
}
func New(pool *pgxpool.Pool, logger *slog.Logger, paths []string) *Scanner {
return &Scanner{pool: pool, logger: logger, paths: paths}
}
// Scan walks every configured root and upserts any audio file whose mtime is
// newer than the existing row's updated_at. Walk errors and per-file errors
// are logged + counted; the scan keeps going.
//
// It then reconciles: rows whose file the walk never saw get marked missing,
// and rows whose file has come back get un-marked (#2523). Only a FULL scan may
// do this — the walk's set of seen paths is the evidence, and a partial
// (watcher-driven) scan has no basis for concluding anything about files it
// didn't look at. That's why ScanFiles does not reconcile.
//
// progressCb (may be nil) receives the current Stats snapshot after each
// processed file. Used by the orchestrator to drive partial-tally writes.
func (s *Scanner) Scan(ctx context.Context, progressCb func(Stats)) (Stats, error) {
var stats Stats
q := dbq.New(s.pool)
start := time.Now()
// Every audio path the walk visited. Reconcile diffs this against the table,
// so it costs no extra filesystem I/O — the walk already established which
// files exist. ~100 bytes/path, so a 250k-track library is ~25MB, which is
// worth it to avoid a second stat pass over the whole library.
seen := make(map[string]struct{}, 8192)
for _, root := range s.paths {
if err := filepath.WalkDir(root, func(path string, d fs.DirEntry, err error) error {
if ctx.Err() != nil {
return fs.SkipAll
}
if err != nil {
s.logger.Warn("library scan walk error", "path", path, "err", err)
stats.Errored++
if progressCb != nil {
progressCb(stats)
}
return nil
}
if d.IsDir() {
return nil
}
if !audioExtensions[strings.ToLower(filepath.Ext(path))] {
return nil
}
// Recorded before scanFile so a file that exists but fails to parse
// still counts as present. It's a broken file, not a missing one,
// and marking it missing would hide it from the operator behind the
// wrong explanation.
seen[path] = struct{}{}
if _, _, err := s.scanFile(ctx, q, path, &stats); err != nil {
s.logger.Warn("library scan file error", "path", path, "err", err)
stats.Errored++
}
if progressCb != nil {
progressCb(stats)
}
return nil
}); err != nil {
return stats, fmt.Errorf("library: walk %q: %w", root, err)
}
}
// Reconcile only after a COMPLETE walk. A cancelled scan has a partial
// `seen` set, which would mark everything it hadn't reached yet.
if err := ctx.Err(); err != nil {
return stats, err
}
if err := s.reconcileMissing(ctx, q, seen, &stats); err != nil {
// Not fatal: the walk's results are already persisted and useful. The
// guards deliberately refuse to act on ambiguous evidence, and that
// refusal arrives here as an error.
s.logger.Warn("library scan: reconcile skipped", "err", err)
}
s.logger.Info("library scan complete",
"scanned", stats.Scanned,
"added", stats.Added,
"updated", stats.Updated,
"skipped", stats.Skipped,
"errored", stats.Errored,
"missing", stats.Missing,
"restored", stats.Restored,
"duration_ms", time.Since(start).Milliseconds(),
)
if err := ctx.Err(); err != nil {
return stats, err
}
return stats, nil
}
// scanFile upserts a single audio file. Returns the album ID the track
// belongs to and whether the file was added/updated (false = skipped as
// unchanged), so watcher-driven callers can enrich just the changed albums.
func (s *Scanner) scanFile(
ctx context.Context, q *dbq.Queries, path string, stats *Stats,
) (pgtype.UUID, bool, error) {
stats.Scanned++
info, err := os.Stat(path)
if err != nil {
return pgtype.UUID{}, false, fmt.Errorf("stat: %w", err)
}
mtime := info.ModTime()
existing, err := q.GetTrackByPath(ctx, path)
knownTrack := err == nil
if err != nil && !errors.Is(err, pgx.ErrNoRows) {
return pgtype.UUID{}, false, fmt.Errorf("lookup: %w", err)
}
// Incremental skip: only when the file hasn't changed AND we already have a
// real duration AND the row's tag-derived columns were written by the
// current extraction logic. The duration clause lets older scans that
// recorded duration_ms=0 (before ffprobe was wired) get backfilled without
// forcing the operator to wipe the library; the tag-version clause does the
// same job for tag-extraction fixes (#2499). Once both are current,
// subsequent scans short-circuit as before.
unchanged := knownTrack && !existing.UpdatedAt.Time.Before(mtime)
if unchanged && existing.DurationMs > 0 && existing.TagReadVersion >= tagReadVersion {
stats.Skipped++
return pgtype.UUID{}, false, nil
}
f, err := os.Open(path)
if err != nil {
return pgtype.UUID{}, false, fmt.Errorf("open: %w", err)
}
defer func() { _ = f.Close() }()
meta, err := tag.ReadFrom(f)
if err != nil {
return pgtype.UUID{}, false, fmt.Errorf("tag read: %w", err)
}
albumMBID, artistMBID := extractMBIDs(meta)
recordingMBID := extractRecordingMBID(meta)
artistName := meta.Artist()
if artistName == "" {
artistName = "Unknown Artist"
}
albumTitle := meta.Album()
if albumTitle == "" {
albumTitle = "Unknown Album"
}
trackTitle := meta.Title()
if trackTitle == "" {
trackTitle = strings.TrimSuffix(filepath.Base(path), filepath.Ext(path))
}
artist, err := s.resolveArtist(ctx, q, artistName, artistMBID)
if err != nil {
return pgtype.UUID{}, false, fmt.Errorf("artist: %w", err)
}
album, err := s.resolveAlbum(ctx, q, artist.ID, albumTitle, meta.Year(), albumMBID)
if err != nil {
return pgtype.UUID{}, false, fmt.Errorf("album: %w", err)
}
trackNum, _ := meta.Track()
discNum, _ := meta.Disc()
// An unchanged file being re-read only to refresh tag-derived columns
// doesn't need another ffprobe: the stored duration is still accurate, and
// the file's bytes haven't moved. This keeps a library-wide tag-repair pass
// (a tagReadVersion bump) bound by tag reads rather than costing one
// fork+exec per file.
var durationMs int32
if unchanged && existing.DurationMs > 0 {
durationMs = existing.DurationMs
} else {
probed, perr := probeDurationMs(ctx, path)
if perr != nil {
// Missing duration is degraded UX (clients can't scrub) but not a
// blocker for ingestion. Record the file with 0ms; the next scan
// will retry via the backfill clause in the skip check above.
s.logger.Warn("library scan: ffprobe failed", "path", path, "err", perr)
}
durationMs = probed
}
params := dbq.UpsertTrackParams{
Title: trackTitle,
AlbumID: album.ID,
ArtistID: artist.ID,
DurationMs: durationMs,
FilePath: path,
FileSize: info.Size(),
FileFormat: strings.TrimPrefix(strings.ToLower(filepath.Ext(path)), "."),
// Stamped so a future extraction fix can find this row again.
TagReadVersion: tagReadVersion,
}
if trackNum > 0 {
v := int32(trackNum)
params.TrackNumber = &v
}
if discNum > 0 {
v := int32(discNum)
params.DiscNumber = &v
}
if genres, fellBack := extractGenres(meta, f); len(genres) > 0 {
if fellBack {
// dhowden/tag's welded value — see genre.go. Logged because the
// stored genre for this file is the old, corrupt shape.
s.logger.Warn("library scan: genre frame unreadable, using fallback",
"path", path, "genre", meta.Genre())
}
g := strings.Join(genres, genreDelimiter)
params.Genre = &g
}
// Recording MBID feeds the ListenBrainz similarity pipeline.
// UpsertTrack heals mbid on the file_path conflict, so a re-scan
// of a previously-untagged-into-DB track backfills it for free.
if recordingMBID != "" {
m := recordingMBID
params.Mbid = &m
}
track, err := q.UpsertTrack(ctx, params)
if err != nil {
return pgtype.UUID{}, false, fmt.Errorf("upsert track: %w", err)
}
if err := syncpkg.LogChange(ctx, s.pool, syncpkg.EntityTrack,
syncpkg.FormatUUID(track.ID), syncpkg.OpUpsert); err != nil {
// Best-effort: log but don't fail the scan. The next scan that
// touches this track will re-emit the change.
s.logger.Warn("library scan: LogChange track upsert failed", "track_id", track.ID, "err", err)
}
if knownTrack {
stats.Updated++
} else {
stats.Added++
}
return album.ID, true, nil
}
// ScanFiles processes a specific set of audio file paths (watcher-driven),
// applying the same upsert + delta-skip logic as a full Scan. Non-audio or
// unreadable paths are skipped (logged), never fatal. Returns the distinct
// album IDs whose tracks were added or updated, so the caller can enrich
// just those albums inline rather than waiting for a batch pass.
func (s *Scanner) ScanFiles(ctx context.Context, paths []string) ([]pgtype.UUID, error) {
q := dbq.New(s.pool)
seen := make(map[[16]byte]struct{})
changed := make([]pgtype.UUID, 0)
var stats Stats
for _, path := range paths {
if ctx.Err() != nil {
return changed, ctx.Err()
}
if !audioExtensions[strings.ToLower(filepath.Ext(path))] {
continue
}
albumID, didChange, err := s.scanFile(ctx, q, path, &stats)
if err != nil {
s.logger.Warn("library watch: scan file error", "path", path, "err", err)
continue
}
if didChange && albumID.Valid {
if _, ok := seen[albumID.Bytes]; !ok {
seen[albumID.Bytes] = struct{}{}
changed = append(changed, albumID)
}
}
}
if stats.Added > 0 || stats.Updated > 0 {
s.logger.Info("library watch: scan batch",
"added", stats.Added,
"updated", stats.Updated,
"skipped", stats.Skipped,
"errored", stats.Errored,
)
}
return changed, nil
}
func (s *Scanner) resolveArtist(ctx context.Context, q *dbq.Queries, name, mbid string) (dbq.Artist, error) {
existing, err := q.GetArtistByName(ctx, name)
if err == nil {
// Heal: backfill mbid on a previously-imported row if we have one now.
if mbid != "" && (existing.Mbid == nil || *existing.Mbid == "") {
m := mbid
if uerr := q.SetArtistMbidIfNull(ctx, dbq.SetArtistMbidIfNullParams{
ID: existing.ID,
Mbid: &m,
}); uerr != nil {
s.logger.Warn("library scan: heal artist mbid failed",
"artist_id", existing.ID, "err", uerr)
} else {
existing.Mbid = &m
}
}
return existing, nil
}
if !errors.Is(err, pgx.ErrNoRows) {
return dbq.Artist{}, err
}
params := dbq.UpsertArtistParams{
Name: name,
SortName: sortKey(name),
}
if mbid != "" {
m := mbid
params.Mbid = &m
}
artist, err := q.UpsertArtist(ctx, params)
if err != nil {
return dbq.Artist{}, err
}
if err := syncpkg.LogChange(ctx, s.pool, syncpkg.EntityArtist,
syncpkg.FormatUUID(artist.ID), syncpkg.OpUpsert); err != nil {
s.logger.Warn("library scan: LogChange artist upsert failed", "artist_id", artist.ID, "err", err)
}
return artist, nil
}
func (s *Scanner) resolveAlbum(ctx context.Context, q *dbq.Queries, artistID pgtype.UUID, title string, year int, mbid string) (dbq.Album, error) {
existing, err := q.GetAlbumByArtistAndTitle(ctx, dbq.GetAlbumByArtistAndTitleParams{ArtistID: artistID, Title: title})
if err == nil {
// Heal: backfill mbid on a previously-imported row if we have one now.
if mbid != "" && (existing.Mbid == nil || *existing.Mbid == "") {
m := mbid
if uerr := q.SetAlbumMbidIfNull(ctx, dbq.SetAlbumMbidIfNullParams{
ID: existing.ID,
Mbid: &m,
}); uerr != nil {
if isUniqueViolation(uerr) {
// Another album row already owns this MBID — duplicate
// release in the DB. Leave NULL; operator merges later.
s.logger.Info("library scan: duplicate album mbid (canonical row already owns it)",
"album_id", existing.ID, "mbid", mbid)
} else {
s.logger.Warn("library scan: heal album mbid failed",
"album_id", existing.ID, "err", uerr)
}
} else {
existing.Mbid = &m
}
}
return existing, nil
}
if !errors.Is(err, pgx.ErrNoRows) {
return dbq.Album{}, err
}
params := dbq.UpsertAlbumParams{
Title: title,
SortTitle: sortKey(title),
ArtistID: artistID,
}
if d, ok := releaseDateFromYear(year); ok {
params.ReleaseDate = d
} else if year != 0 {
// year=0 is the no-tag case; anything else getting rejected is a tag
// we couldn't trust (typo, OOB number, …). Log it so users can chase
// down the file but don't fail the album insert over a soft field.
s.logger.Warn("library scan: dropping invalid release year",
"year", year, "album", title)
}
if mbid != "" {
m := mbid
params.Mbid = &m
}
album, err := q.UpsertAlbum(ctx, params)
if err != nil {
return dbq.Album{}, err
}
if err := syncpkg.LogChange(ctx, s.pool, syncpkg.EntityAlbum,
syncpkg.FormatUUID(album.ID), syncpkg.OpUpsert); err != nil {
s.logger.Warn("library scan: LogChange album upsert failed", "album_id", album.ID, "err", err)
}
return album, nil
}
// releaseDateFromYear converts a tag-supplied year into a Postgres date,
// returning ok=false if the year is outside what we'll accept. We're
// deliberately strict (1..9999): Postgres' date type goes much wider, but
// 5-digit years from ID3 tags are always garbage in practice and trigger
// SQLSTATE 22008 (datetime field overflow) at insert time.
func releaseDateFromYear(year int) (pgtype.Date, bool) {
if year < 1 || year > 9999 {
return pgtype.Date{}, false
}
return pgtype.Date{
Time: time.Date(year, 1, 1, 0, 0, 0, 0, time.UTC),
Valid: true,
}, true
}
// sortKey drops a leading "The " for sortable ordering. Non-English articles
// (Los, Die, Les) can be added if users ask — keeping the rule obvious for now.
func sortKey(s string) string {
if len(s) >= 4 && strings.EqualFold(s[:4], "the ") {
return s[4:]
}
return s
}
// probeTimeout bounds how long ffprobe is allowed to inspect a single file.
// 10s is generous for local mp3/flac — hits are usually <100ms — but caps
// blast radius if a pathological file or a slow network mount stalls a scan.
const probeTimeout = 10 * time.Second
// probeDurationMs shells out to ffprobe to extract a track's duration. We
// rely on ffmpeg being in the image (see Dockerfile). The CLI is slow per
// call (fork+exec) but scans are batch-mode; this is simpler than pulling
// a Go-side decoder library and handles every format ffmpeg does.
func probeDurationMs(ctx context.Context, path string) (int32, error) {
probeCtx, cancel := context.WithTimeout(ctx, probeTimeout)
defer cancel()
cmd := exec.CommandContext(probeCtx, "ffprobe",
"-v", "error",
"-show_entries", "format=duration",
"-of", "default=noprint_wrappers=1:nokey=1",
path,
)
out, err := cmd.Output()
if err != nil {
return 0, fmt.Errorf("ffprobe: %w", err)
}
seconds, err := strconv.ParseFloat(strings.TrimSpace(string(out)), 64)
if err != nil {
return 0, fmt.Errorf("parse ffprobe output %q: %w", out, err)
}
if seconds <= 0 || math.IsNaN(seconds) || math.IsInf(seconds, 0) {
return 0, fmt.Errorf("invalid duration: %v", seconds)
}
ms := seconds * 1000
if ms > math.MaxInt32 {
ms = math.MaxInt32
}
return int32(ms), nil
}