release / web (push) Successful in 1m44s
release / go (push) Successful in 2m1s
release / govulncheck (push) Successful in 17s
release / integration (push) Successful in 5m22s
release / android (push) Successful in 5m48s
release / Build signed APK (releases and dev) (push) Successful in 5m53s
release / Attach APK to the Release (tag releases only) (push) Skipped
release / Build + push container image (push) Successful in 2m6s
release / Verify release artifacts (tag releases only) (push) Skipped
Migration 0069 adds tracks.mbid_source (tag | acoustid), a lookup state per track (matched | ambiguous | no_match | failed) and the acoustid_settings row (off, no key, min score 0.85). The file's tag outranks a lookup (D4). UpsertTrack keeps a looked-up id through a re-read that finds no tag id and replaces it as soon as one appears. SetTrackMbidFromAcoustID refuses to write over a tag id. The worker fingerprints each untagged track with fpcalc's compressed print, looks it up and writes an id only when D5 settles it: one recording at or above the threshold, or one left after matching title and length. Ambiguous and no-match results write nothing. A key AcoustID refuses, or the service being unreachable, stops the pass and is reported in the worker's status. It never counts as a verdict on a track. A changed file drops its lookup in the scan. The re-lookup takes back an id that no longer matches. Admin API: GET /api/admin/library/acoustid (settings, status, coverage by source), PUT …/acoustid-settings (write-only key), POST …/acoustid/run, GET …/acoustid/unsettled. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
455 lines
14 KiB
Go
455 lines
14 KiB
Go
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)
|
|
}
|