// 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++ } }