package library import ( "context" "errors" "fmt" "log/slog" "strings" "sync" "time" "unicode" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgtype" "github.com/jackc/pgx/v5/pgxpool" "git.fabledsword.com/bvandeusen/minstrel/internal/acoustid" "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" ) // AcoustID lookup worker (M401 #3921). // // A track whose tags carry no recording MBID is invisible to the ListenBrainz // similarity arm: it is never a seed and never comes back as a result. This // worker fingerprints such tracks and asks AcoustID which recording they are, // writing an id only when the answer is unambiguous (D5). The file's own tag // always outranks it (D4). // // It never gates anything else (rule 164). With no key, switched off or // AcoustID unreachable, it idles and says why through Status; scans and // playback do not notice. // acoustIDLookupTick is how often the worker looks for work. A settings save // starts a pass at once, so this only paces new tracks and retries after an // outage. const acoustIDLookupTick = 15 * time.Minute // acoustIDLookupBatch is how many tracks one query hands the worker. const acoustIDLookupBatch = 50 // acoustIDLookupConcurrency is how many files are fingerprinted at once. The // client serialises the requests themselves to AcoustID's rate, so this only // overlaps fpcalc's decodes; low for the reason the loudness backfill's is. const acoustIDLookupConcurrency = 2 // acoustIDDurationToleranceSec is how far a recording's MusicBrainz length may // be from the file's and still count as the same take when telling candidates // apart. A release's length and a rip's differ by a second or two of // silence; a radio edit or a live take differs by far more. const acoustIDDurationToleranceSec = 3 // Lookup states, as migration 0069's CHECK has them. const ( lookupMatched = "matched" lookupAmbiguous = "ambiguous" lookupNoMatch = "no_match" lookupFailed = "failed" ) // recordingLookup is the AcoustID client as the worker uses it. type recordingLookup interface { Lookup(ctx context.Context, apiKey, fingerprint string, durationSec int) ([]acoustid.Candidate, error) } // AcoustIDLookupStatus is what the admin card shows about the worker. type AcoustIDLookupStatus struct { Running bool LastPassAt time.Time // Problem says why the last pass stopped short, in words for the card: // the key was refused, or AcoustID could not be reached. Empty when the // last pass ran to the end. Problem string } // AcoustIDLookupResult tallies one pass. type AcoustIDLookupResult struct { Processed int Matched int Ambiguous int NoMatch int Failed int Inconclusive int // nothing stored; tried again on a later pass } func (r *AcoustIDLookupResult) add(state string) { r.Processed++ switch state { case lookupMatched: r.Matched++ case lookupAmbiguous: r.Ambiguous++ case lookupNoMatch: r.NoMatch++ case lookupFailed: r.Failed++ default: r.Inconclusive++ } } // errPassStopped ends a pass early: every further lookup would fail the same // way, so asking again would only spend the rate limit. type errPassStopped struct{ reason string } func (e errPassStopped) Error() string { return e.reason } // AcoustIDLookupWorker fills recording MBIDs through AcoustID. type AcoustIDLookupWorker struct { pool *pgxpool.Pool logger *slog.Logger settings *AcoustIDSettingsService client recordingLookup tick time.Duration batch int32 // fingerprint is a field so an integration test pins which tracks a pass // touches, not what fpcalc prints. fingerprint func(ctx context.Context, path string) (lookupFingerprint, error) kick chan struct{} mu sync.Mutex status AcoustIDLookupStatus } // NewAcoustIDLookupWorker builds a worker with the production cadence. func NewAcoustIDLookupWorker( pool *pgxpool.Pool, logger *slog.Logger, settings *AcoustIDSettingsService, client recordingLookup, ) *AcoustIDLookupWorker { return &AcoustIDLookupWorker{ pool: pool, logger: logger, settings: settings, client: client, tick: acoustIDLookupTick, batch: acoustIDLookupBatch, fingerprint: computeLookupFingerprint, kick: make(chan struct{}, 1), } } // Kick asks for a pass now, for the admin card's "look up now". A pass // already running absorbs it. func (w *AcoustIDLookupWorker) Kick() { if w == nil { return } select { case w.kick <- struct{}{}: default: } } // Settings is the settings service the worker reads, for the admin API. A // nil worker has none, which serves the defaults. func (w *AcoustIDLookupWorker) Settings() *AcoustIDSettingsService { if w == nil { return nil } return w.settings } // Status reports the worker's state for the admin card. func (w *AcoustIDLookupWorker) Status() AcoustIDLookupStatus { if w == nil { return AcoustIDLookupStatus{} } w.mu.Lock() defer w.mu.Unlock() return w.status } // Run blocks until ctx is cancelled: one pass at start, then one per tick, on // a kick, and after every settings save. func (w *AcoustIDLookupWorker) Run(ctx context.Context) { w.runOnce(ctx) t := time.NewTicker(w.tick) defer t.Stop() for { select { case <-ctx.Done(): return case <-t.C: case <-w.kick: case <-w.settings.Changed(): } w.runOnce(ctx) } } // runOnce contains a pass so that nothing it does (an error, a panic) can stop // the next one from starting (rule 157). func (w *AcoustIDLookupWorker) runOnce(ctx context.Context) { w.setStatus(func(s *AcoustIDLookupStatus) { s.Running = true }) problem := "" defer func() { if r := recover(); r != nil { w.logger.Error("acoustid lookup: pass panicked", "panic", r) problem = "The last lookup pass stopped on an internal error; see the server log." } w.setStatus(func(s *AcoustIDLookupStatus) { s.Running, s.LastPassAt, s.Problem = false, time.Now(), problem }) }() res, err := w.pass(ctx) var stopped errPassStopped switch { case errors.As(err, &stopped): problem = stopped.reason w.logger.Warn("acoustid lookup: pass stopped", "reason", stopped.reason, "processed", res.Processed) case err != nil && ctx.Err() == nil: problem = "The last lookup pass failed; see the server log." w.logger.Warn("acoustid lookup: pass failed", "err", err, "processed", res.Processed) } if res.Processed > 0 { w.logger.Info("acoustid lookup: pass complete", "processed", res.Processed, "matched", res.Matched, "ambiguous", res.Ambiguous, "no_match", res.NoMatch, "failed", res.Failed, "inconclusive", res.Inconclusive) } } func (w *AcoustIDLookupWorker) setStatus(f func(*AcoustIDLookupStatus)) { w.mu.Lock() f(&w.status) w.mu.Unlock() } // pass walks the queue once, keyset-paged on id so it ends even when every // attempt is inconclusive. Settings are read before every batch, so switching // the lookup off or changing the threshold applies at the next batch. func (w *AcoustIDLookupWorker) pass(ctx context.Context) (AcoustIDLookupResult, error) { q := dbq.New(w.pool) var res AcoustIDLookupResult after := pgtype.UUID{Valid: true} for { if err := ctx.Err(); err != nil { return res, err } cfg := w.settings.Get() if !cfg.Ready() { return res, nil } rows, err := q.ListTracksNeedingAcoustIDLookup(ctx, dbq.ListTracksNeedingAcoustIDLookupParams{ AfterID: after, BatchLimit: w.batch, }) if err != nil { return res, fmt.Errorf("list tracks needing a lookup: %w", err) } if len(rows) == 0 { return res, nil } if err := w.lookupBatch(ctx, cfg, rows, &res); err != nil { return res, err } after = rows[len(rows)-1].ID } } // lookupBatch looks up one batch, acoustIDLookupConcurrency at a time. The // first track that stops the pass stops the batch: the tracks already started // finish, no new ones start. func (w *AcoustIDLookupWorker) lookupBatch( ctx context.Context, cfg AcoustIDSettings, rows []dbq.ListTracksNeedingAcoustIDLookupRow, res *AcoustIDLookupResult, ) error { ctx, cancel := context.WithCancelCause(ctx) defer cancel(nil) var ( mu sync.Mutex wg sync.WaitGroup sem = make(chan struct{}, acoustIDLookupConcurrency) ) for _, row := range rows { if ctx.Err() != nil { break } sem <- struct{}{} // The slot may have come free because a lookup just stopped the pass. if ctx.Err() != nil { <-sem break } wg.Add(1) go func(row dbq.ListTracksNeedingAcoustIDLookupRow) { defer wg.Done() defer func() { <-sem }() defer func() { if r := recover(); r != nil { w.logger.Error("acoustid lookup: track panicked", "path", row.FilePath, "panic", r) } }() state, err := w.lookupTrack(ctx, cfg, row) if err != nil { cancel(err) } mu.Lock() res.add(state) mu.Unlock() }(row) } wg.Wait() if cause := context.Cause(ctx); cause != nil && !errors.Is(cause, context.Canceled) { return cause } return nil } // lookupTrack fingerprints one file, asks AcoustID and stores the answer. It // returns the state stored ("" when nothing was), and an errPassStopped when // no further lookup can succeed this pass. func (w *AcoustIDLookupWorker) lookupTrack( ctx context.Context, cfg AcoustIDSettings, row dbq.ListTracksNeedingAcoustIDLookupRow, ) (string, error) { fp, err := w.fingerprint(ctx, row.FilePath) if err != nil { if isInconclusive(err) { return "", nil } // fpcalc ran and rejected the file: a verdict, settled until it changes. return w.store(ctx, row, lookupDecision{state: lookupFailed, detail: err.Error()}), nil } cands, err := w.client.Lookup(ctx, cfg.APIKey, fp.fingerprint, fp.durationSec) switch { case err == nil: case errors.Is(err, acoustid.ErrInvalidKey): return "", errPassStopped{"AcoustID refused the API key. Check it in Settings."} case errors.Is(err, acoustid.ErrUnavailable): return "", errPassStopped{"AcoustID could not be reached; lookups resume on the next pass."} case errors.Is(err, acoustid.ErrInvalidFingerprint): return w.store(ctx, row, lookupDecision{state: lookupFailed, detail: err.Error()}), nil default: // The caller's cancellation, or a request this client built wrongly. // Neither is a verdict on the track. if ctx.Err() == nil { w.logger.Warn("acoustid lookup: lookup failed", "path", row.FilePath, "err", err) } return "", nil } return w.store(ctx, row, chooseRecording(cands, cfg.MinScore, row.Title, row.DurationMs)), nil } // store records a decision and applies it to tracks.mbid in one transaction, // so the gauge never shows a match whose id was not written. A failed write // stores nothing and the track is tried again. func (w *AcoustIDLookupWorker) store(ctx context.Context, row dbq.ListTracksNeedingAcoustIDLookupRow, d lookupDecision) string { err := pgx.BeginFunc(ctx, w.pool, func(tx pgx.Tx) error { q := dbq.New(tx) params := dbq.RecordAcoustIDLookupParams{ TrackID: row.ID, State: d.state, Candidates: int32(d.candidates), } if d.candidates > 0 { params.BestScore = &d.bestScore } if d.recording != "" { params.RecordingMbid = &d.recording } if d.detail != "" { params.Detail = &d.detail } if err := q.RecordAcoustIDLookup(ctx, params); err != nil { return err } if d.state == lookupMatched { _, err := q.SetTrackMbidFromAcoustID(ctx, dbq.SetTrackMbidFromAcoustIDParams{ID: row.ID, Mbid: &d.recording}) return err } // A re-lookup of changed bytes that no longer matches takes back the // id the earlier lookup wrote. return q.ClearTrackAcoustIDMbid(ctx, row.ID) }) if err != nil { if ctx.Err() == nil { w.logger.Warn("acoustid lookup: storing the result failed", "path", row.FilePath, "err", err) } return "" } return d.state } // lookupDecision is what one lookup settles on. type lookupDecision struct { state string recording string // set only when matched bestScore float32 candidates int detail string } // chooseRecording applies D5 to AcoustID's candidates (best score first): // only recordings at or above minScore count, one of them is a match, and // several are told apart by the track's title and length or not at all. A // wrong MBID would feed similarity the wrong neighbours, which is worse than // none, so anything left unresolved is ambiguous and writes nothing. func chooseRecording(cands []acoustid.Candidate, minScore float32, title string, durationMs int32) lookupDecision { d := lookupDecision{candidates: len(cands)} if len(cands) > 0 { d.bestScore = float32(cands[0].Score) } var above []acoustid.Candidate for _, c := range cands { if float32(c.Score) >= minScore { above = append(above, c) } } switch len(above) { case 0: d.state = lookupNoMatch return d case 1: d.state, d.recording = lookupMatched, above[0].ID return d } want := normalizeTitle(title) var same []acoustid.Candidate for _, c := range above { if normalizeTitle(c.Title) == want && sameLength(c.DurationSec, durationMs) { same = append(same, c) } } if len(same) == 1 { d.state, d.recording = lookupMatched, same[0].ID return d } d.state = lookupAmbiguous return d } // sameLength reports whether a recording's length (0 when MusicBrainz has // none) fits the file's. An unknown length on either side tells nothing, so // it does not rule a candidate out. func sameLength(recordingSec int, fileMs int32) bool { if recordingSec <= 0 || fileMs <= 0 { return true } diff := recordingSec*1000 - int(fileMs) return diff <= acoustIDDurationToleranceSec*1000 && diff >= -acoustIDDurationToleranceSec*1000 } // normalizeTitle compares titles by their letters and digits alone, case // folded, so "Don't Stop" and "Dont stop" match while "Song (Radio Edit)" // and "Song" do not. func normalizeTitle(s string) string { var b strings.Builder for _, r := range strings.ToLower(s) { if unicode.IsLetter(r) || unicode.IsDigit(r) { b.WriteRune(r) } } return b.String() } // AcoustIDCoverage reports where the library's recording MBIDs came from and // what the lookups found, for the admin gauge. func AcoustIDCoverage(ctx context.Context, pool *pgxpool.Pool) (dbq.GetMbidCoverageRow, error) { return dbq.New(pool).GetMbidCoverage(ctx) }