Files
minstrel/internal/similarity/worker.go
T
bvandeusenandClaude Opus 5.5 e5dac9ddf0
release / govulncheck (push) Successful in 45s
release / web (push) Successful in 1m27s
release / go (push) Successful in 1m51s
release / integration (push) Successful in 5m36s
release / android (push) Successful in 7m47s
release / Build signed APK (releases and dev) (push) Successful in 8m29s
release / Attach APK to the Release (tag releases only) (push) Skipped
release / Build + push container image (push) Successful in 1m47s
release / Verify release artifacts (tag releases only) (push) Skipped
fix(similarity): keep ListenBrainz's whole answer and resolve it locally (#5296)
The worker kept only the similar recordings already in the library, at most
20 of ListenBrainz's 50, and judged freshness by the edges it had written.
Two failures followed, both measured on the operator's library (#3879):

- A seed whose answer matched nothing wrote nothing, so it was never fresh.
  With the queue ordered by id, 25 such seeds held its head and were
  re-asked every hour; 17 of 2,466 played seeds had any edges.
- A recording that reached the library after its seed was fetched (a
  Lidarr import, an MBID from the AcoustID lookup) was never linked until
  a refetch, which for the stuck seeds never came.

Now every answer is cached whole in listenbrainz_similar_recordings and
every answer, an empty one or a permanent 4xx included, is recorded in
track_similarity_fetches. The queue reads the fetch record: never-fetched
first, then the oldest, refreshed after 30 days. The listenbrainz edges are
derived in SQL from the cache, one present track per recording and no cap,
for the seed just fetched and for every seed once per tick, so new arrivals
link within the hour without asking ListenBrainz again.

Artists get the same queue fix via artist_similarity_fetches; their answer
was already kept in artist_similarity and artist_similarity_unmatched.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-07 19:56:22 -04:00

269 lines
8.4 KiB
Go

// Package similarity owns the inbound ListenBrainz similarity ingest
// pipeline. A periodic worker queries the LB Labs API
// (labs.api.listenbrainz.org similar-recordings / similar-artists) for
// tracks the user has played.
//
// For recordings it keeps LB's whole answer in listenbrainz_similar_recordings
// and derives the track_similarity edges from it in SQL, against whatever is
// in the library at the time (#5296). Matching once at fetch time missed every
// recording that reached the library afterwards, through Lidarr or an MBID the
// AcoustID lookup filled in. For artists it stores the in-library matches in
// artist_similarity and the rest in artist_similarity_unmatched, as before.
//
// Every answer is recorded as a fetch, an empty one included, and the queue
// reads those records: freshness used to be read from the edges, so a seed
// that matched nothing was never fresh and held the head of the queue.
package similarity
import (
"context"
"errors"
"fmt"
"log/slog"
"sort"
"time"
"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"
"git.fabledsword.com/bvandeusen/minstrel/internal/scrobble/listenbrainz"
)
// Worker drains played tracks and artists needing similarity. A transient
// failure is not recorded as a fetch, so the timer retries it next tick.
type Worker struct {
pool *pgxpool.Pool
client *listenbrainz.Client
logger *slog.Logger
tick time.Duration
batch int32
topK int // per artist, for each of the matched and unmatched stores
}
// NewWorker constructs a worker with production defaults: 1h tick,
// batch=25, topK=20. batch is tracks AND artists processed per tick, well
// under ListenBrainz rate limits (429s abort the tick). With never-fetched
// seeds first, a freshly played library is covered in days.
func NewWorker(pool *pgxpool.Pool, client *listenbrainz.Client, logger *slog.Logger) *Worker {
return &Worker{
pool: pool,
client: client,
logger: logger,
tick: 1 * time.Hour,
batch: 25,
topK: 20,
}
}
// Run blocks until ctx is cancelled, ticking every w.tick.
func (w *Worker) Run(ctx context.Context) {
t := time.NewTicker(w.tick)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
if err := w.tickOnce(ctx); err != nil {
w.logger.Error("similarity: tick failed", "err", err)
}
}
}
}
// tickOnce drains one batch of tracks and one batch of artists, then
// re-resolves every cached answer against the library. Per-row errors are
// logged and skipped. 429 aborts the fetching, not the resolve.
func (w *Worker) tickOnce(ctx context.Context) error {
q := dbq.New(w.pool)
if err := w.tickTracks(ctx, q); err != nil {
return err
}
if err := w.tickArtists(ctx, q); err != nil {
return err
}
// Links recordings that reached the library since their seed was fetched:
// a scan, a Lidarr import, an MBID from the AcoustID lookup.
if err := q.ResolveListenBrainzTrackEdges(ctx, nil); err != nil {
return fmt.Errorf("resolve similarity edges: %w", err)
}
return nil
}
func (w *Worker) tickTracks(ctx context.Context, q *dbq.Queries) error {
rows, err := q.ListPlayedTracksNeedingSimilarity(ctx, w.batch)
if err != nil {
return err
}
for _, r := range rows {
if r.Mbid == nil {
continue // defensive — query already filters NULL
}
results, err := w.client.SimilarRecordings(ctx, *r.Mbid, 100)
switch {
case err == nil:
case isRateLimited(err):
w.logger.Warn("similarity: 429 — aborting tick")
return nil
case errors.Is(err, listenbrainz.ErrPermanent):
// LB refused this MBID outright; asking again next hour won't
// change that. Recorded as an empty answer so it waits its turn.
w.logger.Warn("similarity: similar-recordings refused", "track_id", r.ID, "err", err)
results = nil
default:
w.logger.Warn("similarity: similar-recordings failed", "track_id", r.ID, "err", err)
continue
}
if err := w.storeTrackSimilar(ctx, r.ID, results); err != nil {
w.logger.Warn("similarity: store similar recordings", "track_id", r.ID, "err", err)
}
}
return nil
}
func (w *Worker) tickArtists(ctx context.Context, q *dbq.Queries) error {
rows, err := q.ListPlayedArtistsNeedingSimilarity(ctx, w.batch)
if err != nil {
return err
}
for _, r := range rows {
if r.Mbid == nil {
continue
}
results, err := w.client.SimilarArtists(ctx, *r.Mbid, 100)
switch {
case err == nil:
case isRateLimited(err):
w.logger.Warn("similarity: 429 — aborting tick")
return nil
case errors.Is(err, listenbrainz.ErrPermanent):
w.logger.Warn("similarity: similar-artists refused", "artist_id", r.ID, "err", err)
results = nil
default:
w.logger.Warn("similarity: similar-artists failed", "artist_id", r.ID, "err", err)
continue
}
w.upsertArtistSimilar(ctx, q, r.ID, results)
if err := q.RecordArtistSimilarityFetch(ctx, dbq.RecordArtistSimilarityFetchParams{
ArtistID: r.ID, Returned: int32(len(results)),
}); err != nil {
w.logger.Warn("similarity: record artist fetch", "artist_id", r.ID, "err", err)
}
}
return nil
}
func isRateLimited(err error) bool {
var ra *listenbrainz.RetryAfterError
return errors.As(err, &ra)
}
// storeTrackSimilar replaces the seed's cached answer with this one, records
// the fetch, and resolves the seed's edges, in one transaction so the queue
// never sees a fetch whose answer is not stored.
func (w *Worker) storeTrackSimilar(ctx context.Context, seedID pgtype.UUID, results []listenbrainz.SimilarRecording) error {
mbids := make([]string, 0, len(results))
scores := make([]float64, 0, len(results))
for _, r := range results {
mbids = append(mbids, r.MBID)
scores = append(scores, r.Score)
}
return pgx.BeginFunc(ctx, w.pool, func(tx pgx.Tx) error {
q := dbq.New(tx)
if err := q.DeleteSimilarRecordingsForSeed(ctx, seedID); err != nil {
return fmt.Errorf("clear cached answer: %w", err)
}
if len(mbids) > 0 {
if err := q.InsertSimilarRecordings(ctx, dbq.InsertSimilarRecordingsParams{
SeedTrackID: seedID, Mbids: mbids, Scores: scores,
}); err != nil {
return fmt.Errorf("cache answer: %w", err)
}
}
if err := q.RecordTrackSimilarityFetch(ctx, dbq.RecordTrackSimilarityFetchParams{
TrackID: seedID, Returned: int32(len(results)),
}); err != nil {
return fmt.Errorf("record fetch: %w", err)
}
if err := q.ResolveListenBrainzTrackEdges(ctx, []pgtype.UUID{seedID}); err != nil {
return fmt.Errorf("resolve edges: %w", err)
}
return nil
})
}
func (w *Worker) upsertArtistSimilar(ctx context.Context, q *dbq.Queries, artistAID pgtype.UUID, results []listenbrainz.SimilarArtist) {
if len(results) == 0 {
return
}
sort.Slice(results, func(i, j int) bool { return results[i].Score > results[j].Score })
mbids := make([]string, 0, len(results))
for _, r := range results {
mbids = append(mbids, r.MBID)
}
rows, err := q.GetArtistsByMBIDs(ctx, mbids)
if err != nil {
w.logger.Warn("similarity: GetArtistsByMBIDs", "err", err)
return
}
idByMBID := make(map[string]pgtype.UUID, len(rows))
for _, r := range rows {
if r.Mbid != nil {
idByMBID[*r.Mbid] = r.ID
}
}
// Matched: in-library similars → artist_similarity (existing path).
takenMatched := 0
for _, r := range results {
if takenMatched >= w.topK {
break
}
localID, ok := idByMBID[r.MBID]
if !ok {
continue
}
if localID == artistAID {
continue
}
if uerr := q.UpsertArtistSimilarity(ctx, dbq.UpsertArtistSimilarityParams{
ArtistAID: artistAID, ArtistBID: localID, Score: r.Score,
}); uerr != nil {
w.logger.Warn("similarity: UpsertArtistSimilarity", "err", uerr)
continue
}
takenMatched++
}
// Unmatched: out-of-library similars → artist_similarity_unmatched (M5c).
// Same top-K cap as the matched path. Skip rows missing a name — we can't
// render a suggestion card without one.
takenUnmatched := 0
for _, r := range results {
if takenUnmatched >= w.topK {
break
}
if _, inLib := idByMBID[r.MBID]; inLib {
continue
}
if r.Name == "" {
w.logger.Debug("similarity: skipping unmatched similar with empty name", "mbid", r.MBID)
continue
}
if uerr := q.UpsertArtistSimilarityUnmatched(ctx, dbq.UpsertArtistSimilarityUnmatchedParams{
SeedArtistID: artistAID,
CandidateMbid: r.MBID,
CandidateName: r.Name,
Score: r.Score,
Source: "listenbrainz",
}); uerr != nil {
w.logger.Warn("similarity: UpsertArtistSimilarityUnmatched", "err", uerr)
continue
}
takenUnmatched++
}
}