diff --git a/internal/db/dbq/models.go b/internal/db/dbq/models.go index c14ff1a3..3fb3d24f 100644 --- a/internal/db/dbq/models.go +++ b/internal/db/dbq/models.go @@ -241,6 +241,12 @@ type ArtistSimilarity struct { FetchedAt pgtype.Timestamptz } +type ArtistSimilarityFetch struct { + ArtistID pgtype.UUID + FetchedAt pgtype.Timestamptz + Returned int32 +} + type ArtistSimilarityUnmatched struct { SeedArtistID pgtype.UUID CandidateMbid string @@ -436,6 +442,12 @@ type LidarrRequest struct { LidarrAddConfirmedAt pgtype.Timestamptz } +type ListenbrainzSimilarRecording struct { + SeedTrackID pgtype.UUID + RecordingMbid string + Score float64 +} + type LoudnessSetting struct { ID bool Enabled bool @@ -771,6 +783,12 @@ type TrackSimilarity struct { FetchedAt pgtype.Timestamptz } +type TrackSimilarityFetch struct { + TrackID pgtype.UUID + FetchedAt pgtype.Timestamptz + Returned int32 +} + type TrackTag struct { TrackID pgtype.UUID Tag string diff --git a/internal/db/dbq/similarity.sql.go b/internal/db/dbq/similarity.sql.go index 153bcb26..480c64cd 100644 --- a/internal/db/dbq/similarity.sql.go +++ b/internal/db/dbq/similarity.sql.go @@ -11,6 +11,16 @@ import ( "github.com/jackc/pgx/v5/pgtype" ) +const deleteSimilarRecordingsForSeed = `-- name: DeleteSimilarRecordingsForSeed :exec +DELETE FROM listenbrainz_similar_recordings WHERE seed_track_id = $1 +` + +// A fetch replaces the seed's cached answer whole. +func (q *Queries) DeleteSimilarRecordingsForSeed(ctx context.Context, seedTrackID pgtype.UUID) error { + _, err := q.db.Exec(ctx, deleteSimilarRecordingsForSeed, seedTrackID) + return err +} + const getArtistsByMBIDs = `-- name: GetArtistsByMBIDs :many SELECT id, mbid FROM artists WHERE mbid = ANY($1::text[]) ` @@ -40,49 +50,40 @@ func (q *Queries) GetArtistsByMBIDs(ctx context.Context, dollar_1 []string) ([]G return items, nil } -const getTracksByMBIDs = `-- name: GetTracksByMBIDs :many -SELECT id, mbid FROM tracks WHERE mbid = ANY($1::text[]) +const insertSimilarRecordings = `-- name: InsertSimilarRecordings :exec +INSERT INTO listenbrainz_similar_recordings (seed_track_id, recording_mbid, score) +SELECT $1::uuid, r.mbid, max(r.score) +FROM (SELECT unnest($2::text[]) AS mbid, + unnest($3::float8[]) AS score) r +WHERE r.mbid <> '' +GROUP BY r.mbid ` -type GetTracksByMBIDsRow struct { - ID pgtype.UUID - Mbid *string +type InsertSimilarRecordingsParams struct { + SeedTrackID pgtype.UUID + Mbids []string + Scores []float64 } -// Bulk in-library lookup: maps a slice of MBIDs back to local track IDs. -func (q *Queries) GetTracksByMBIDs(ctx context.Context, dollar_1 []string) ([]GetTracksByMBIDsRow, error) { - rows, err := q.db.Query(ctx, getTracksByMBIDs, dollar_1) - if err != nil { - return nil, err - } - defer rows.Close() - var items []GetTracksByMBIDsRow - for rows.Next() { - var i GetTracksByMBIDsRow - if err := rows.Scan(&i.ID, &i.Mbid); err != nil { - return nil, err - } - items = append(items, i) - } - if err := rows.Err(); err != nil { - return nil, err - } - return items, nil +// The seed's whole answer, in or out of the library. An MBID LB lists twice +// keeps its best score. +func (q *Queries) InsertSimilarRecordings(ctx context.Context, arg InsertSimilarRecordingsParams) error { + _, err := q.db.Exec(ctx, insertSimilarRecordings, arg.SeedTrackID, arg.Mbids, arg.Scores) + return err } const listPlayedArtistsNeedingSimilarity = `-- name: ListPlayedArtistsNeedingSimilarity :many -SELECT DISTINCT ar.id, ar.mbid +SELECT ar.id, ar.mbid FROM artists ar -JOIN tracks t ON t.artist_id = ar.id -JOIN play_events pe ON pe.track_id = t.id +LEFT JOIN artist_similarity_fetches f ON f.artist_id = ar.id WHERE ar.mbid IS NOT NULL - AND NOT EXISTS ( - SELECT 1 FROM artist_similarity asim - WHERE asim.artist_a_id = ar.id - AND asim.source = 'listenbrainz' - AND asim.fetched_at > now() - interval '7 days' + AND EXISTS ( + SELECT 1 FROM tracks t + JOIN play_events pe ON pe.track_id = t.id + WHERE t.artist_id = ar.id ) -ORDER BY ar.id + AND (f.artist_id IS NULL OR f.fetched_at < now() - interval '30 days') +ORDER BY f.fetched_at NULLS FIRST, ar.id LIMIT $1 ` @@ -91,6 +92,7 @@ type ListPlayedArtistsNeedingSimilarityRow struct { Mbid *string } +// The artist queue, with the same shape and the same reason. func (q *Queries) ListPlayedArtistsNeedingSimilarity(ctx context.Context, limit int32) ([]ListPlayedArtistsNeedingSimilarityRow, error) { rows, err := q.db.Query(ctx, listPlayedArtistsNeedingSimilarity, limit) if err != nil { @@ -112,17 +114,14 @@ func (q *Queries) ListPlayedArtistsNeedingSimilarity(ctx context.Context, limit } const listPlayedTracksNeedingSimilarity = `-- name: ListPlayedTracksNeedingSimilarity :many -SELECT DISTINCT t.id, t.mbid +SELECT t.id, t.mbid FROM tracks t -JOIN play_events pe ON pe.track_id = t.id +LEFT JOIN track_similarity_fetches f ON f.track_id = t.id WHERE t.mbid IS NOT NULL - AND NOT EXISTS ( - SELECT 1 FROM track_similarity ts - WHERE ts.track_a_id = t.id - AND ts.source = 'listenbrainz' - AND ts.fetched_at > now() - interval '7 days' - ) -ORDER BY t.id + AND t.missing_since IS NULL + AND EXISTS (SELECT 1 FROM play_events pe WHERE pe.track_id = t.id) + AND (f.track_id IS NULL OR f.fetched_at < now() - interval '30 days') +ORDER BY f.fetched_at NULLS FIRST, t.id LIMIT $1 ` @@ -131,9 +130,12 @@ type ListPlayedTracksNeedingSimilarityRow struct { Mbid *string } -// Tracks with at least one play, an MBID, and no fresh listenbrainz row -// (no row at all, OR fetched_at older than 7 days). Used by the worker -// to find work each tick. Bounded by $1. +// The worker's track queue: played, present tracks with a recording MBID that +// ListenBrainz has never answered for, or not in 30 days. Never-fetched first, +// then the oldest answer. Freshness is read from track_similarity_fetches, not +// from the edges, so a seed whose answer matched nothing in the library still +// counts as fetched (#5296); reading the edges is what let 25 such seeds hold +// the head of the queue forever. func (q *Queries) ListPlayedTracksNeedingSimilarity(ctx context.Context, limit int32) ([]ListPlayedTracksNeedingSimilarityRow, error) { rows, err := q.db.Query(ctx, listPlayedTracksNeedingSimilarity, limit) if err != nil { @@ -225,6 +227,86 @@ func (q *Queries) ListSimilarArtistsForArtist(ctx context.Context, arg ListSimil return items, nil } +const recordArtistSimilarityFetch = `-- name: RecordArtistSimilarityFetch :exec +INSERT INTO artist_similarity_fetches (artist_id, fetched_at, returned) +VALUES ($1, now(), $2) +ON CONFLICT (artist_id) DO UPDATE SET fetched_at = now(), returned = EXCLUDED.returned +` + +type RecordArtistSimilarityFetchParams struct { + ArtistID pgtype.UUID + Returned int32 +} + +func (q *Queries) RecordArtistSimilarityFetch(ctx context.Context, arg RecordArtistSimilarityFetchParams) error { + _, err := q.db.Exec(ctx, recordArtistSimilarityFetch, arg.ArtistID, arg.Returned) + return err +} + +const recordTrackSimilarityFetch = `-- name: RecordTrackSimilarityFetch :exec +INSERT INTO track_similarity_fetches (track_id, fetched_at, returned) +VALUES ($1, now(), $2) +ON CONFLICT (track_id) DO UPDATE SET fetched_at = now(), returned = EXCLUDED.returned +` + +type RecordTrackSimilarityFetchParams struct { + TrackID pgtype.UUID + Returned int32 +} + +func (q *Queries) RecordTrackSimilarityFetch(ctx context.Context, arg RecordTrackSimilarityFetchParams) error { + _, err := q.db.Exec(ctx, recordTrackSimilarityFetch, arg.TrackID, arg.Returned) + return err +} + +const resolveListenBrainzTrackEdges = `-- name: ResolveListenBrainzTrackEdges :exec +WITH want AS ( + SELECT DISTINCT ON (c.seed_track_id, c.recording_mbid) + c.seed_track_id AS track_a_id, t.id AS track_b_id, c.score + FROM listenbrainz_similar_recordings c + JOIN tracks t ON t.mbid = c.recording_mbid + AND t.missing_since IS NULL + AND t.id <> c.seed_track_id + WHERE $1::uuid[] IS NULL + OR c.seed_track_id = ANY($1::uuid[]) + ORDER BY c.seed_track_id, c.recording_mbid, t.id +), +upserted AS ( + INSERT INTO track_similarity (track_a_id, track_b_id, score, source, fetched_at) + SELECT track_a_id, track_b_id, score, 'listenbrainz', now() FROM want + ON CONFLICT (track_a_id, track_b_id, source) DO UPDATE + SET score = EXCLUDED.score, fetched_at = EXCLUDED.fetched_at + WHERE track_similarity.score IS DISTINCT FROM EXCLUDED.score + RETURNING 1 +) +DELETE FROM track_similarity ts +WHERE ts.source = 'listenbrainz' + AND ($1::uuid[] IS NULL OR ts.track_a_id = ANY($1::uuid[])) + AND EXISTS (SELECT 1 FROM track_similarity_fetches f WHERE f.track_id = ts.track_a_id) + AND NOT EXISTS ( + SELECT 1 FROM want w + WHERE w.track_a_id = ts.track_a_id AND w.track_b_id = ts.track_b_id + ) +` + +// Derives the listenbrainz edges in track_similarity from the cached answers: +// every cached recording that is a present track in the library becomes an +// edge, whenever it got there. Scoped to the given seeds, or every seed when +// seed_ids is NULL. +// +// One track per recording MBID: two copies of a song (an mp3 beside its flac, +// or one recording on two releases) would otherwise be two candidates, and the +// pool would offer the same song twice. No cap per seed: LB's answer is already +// at most 50. +// +// Only edges of seeds the cache owns are removed, those with a fetch row. +// Edges written before the cache existed, or copied onto a survivor by a fold +// (merge.sql), stay until their seed is fetched. +func (q *Queries) ResolveListenBrainzTrackEdges(ctx context.Context, seedIds []pgtype.UUID) error { + _, err := q.db.Exec(ctx, resolveListenBrainzTrackEdges, seedIds) + return err +} + const upsertArtistSimilarity = `-- name: UpsertArtistSimilarity :exec INSERT INTO artist_similarity (artist_a_id, artist_b_id, score, source, fetched_at) VALUES ($1, $2, $3, 'listenbrainz', now()) @@ -274,21 +356,3 @@ func (q *Queries) UpsertArtistSimilarityUnmatched(ctx context.Context, arg Upser ) return err } - -const upsertTrackSimilarity = `-- name: UpsertTrackSimilarity :exec -INSERT INTO track_similarity (track_a_id, track_b_id, score, source, fetched_at) -VALUES ($1, $2, $3, 'listenbrainz', now()) -ON CONFLICT (track_a_id, track_b_id, source) -DO UPDATE SET score = EXCLUDED.score, fetched_at = EXCLUDED.fetched_at -` - -type UpsertTrackSimilarityParams struct { - TrackAID pgtype.UUID - TrackBID pgtype.UUID - Score float64 -} - -func (q *Queries) UpsertTrackSimilarity(ctx context.Context, arg UpsertTrackSimilarityParams) error { - _, err := q.db.Exec(ctx, upsertTrackSimilarity, arg.TrackAID, arg.TrackBID, arg.Score) - return err -} diff --git a/internal/db/migrations/0072_similarity_cache.down.sql b/internal/db/migrations/0072_similarity_cache.down.sql new file mode 100644 index 00000000..fb3b01a0 --- /dev/null +++ b/internal/db/migrations/0072_similarity_cache.down.sql @@ -0,0 +1,3 @@ +DROP TABLE IF EXISTS artist_similarity_fetches; +DROP TABLE IF EXISTS track_similarity_fetches; +DROP TABLE IF EXISTS listenbrainz_similar_recordings; diff --git a/internal/db/migrations/0072_similarity_cache.up.sql b/internal/db/migrations/0072_similarity_cache.up.sql new file mode 100644 index 00000000..c23ca242 --- /dev/null +++ b/internal/db/migrations/0072_similarity_cache.up.sql @@ -0,0 +1,43 @@ +-- #5296: keep ListenBrainz's whole similar-recordings answer, and record every +-- fetch, so the library match is made locally and can be made again. +-- +-- Until now the worker kept only the recordings already in the library (at most +-- 20 of LB's 50) and threw the rest away. Two things followed: +-- +-- * A recording that reached the library later — a Lidarr import, or an MBID +-- filled in by the AcoustID lookup — was not linked until its seed was +-- fetched again. +-- * A seed with no in-library match wrote nothing at all, and freshness was +-- read from track_similarity, so that seed was never fresh. With the queue +-- ordered by id, 25 such seeds held its head and were re-asked every hour +-- while 2,449 others waited (#3879). +-- +-- The listenbrainz rows in track_similarity are now derived from this cache by +-- ResolveListenBrainzTrackEdges. + +CREATE TABLE listenbrainz_similar_recordings ( + seed_track_id uuid NOT NULL REFERENCES tracks(id) ON DELETE CASCADE, + recording_mbid text NOT NULL, + score DOUBLE PRECISION NOT NULL, + PRIMARY KEY (seed_track_id, recording_mbid) +); + +-- The resolve joins the cache to tracks by recording MBID. +CREATE INDEX listenbrainz_similar_recordings_mbid_idx + ON listenbrainz_similar_recordings (recording_mbid); + +-- One row per seed that ListenBrainz has answered for, an empty answer +-- included. The worker's queue reads this, not the edges. +CREATE TABLE track_similarity_fetches ( + track_id uuid PRIMARY KEY REFERENCES tracks(id) ON DELETE CASCADE, + fetched_at timestamptz NOT NULL DEFAULT now(), + returned integer NOT NULL +); + +-- The same for artists. Their answer is already kept, in artist_similarity and +-- artist_similarity_unmatched; this only stops an empty answer being re-asked. +CREATE TABLE artist_similarity_fetches ( + artist_id uuid PRIMARY KEY REFERENCES artists(id) ON DELETE CASCADE, + fetched_at timestamptz NOT NULL DEFAULT now(), + returned integer NOT NULL +); diff --git a/internal/db/queries/similarity.sql b/internal/db/queries/similarity.sql index 94989f94..70949a35 100644 --- a/internal/db/queries/similarity.sql +++ b/internal/db/queries/similarity.sql @@ -1,47 +1,103 @@ -- name: ListPlayedTracksNeedingSimilarity :many --- Tracks with at least one play, an MBID, and no fresh listenbrainz row --- (no row at all, OR fetched_at older than 7 days). Used by the worker --- to find work each tick. Bounded by $1. -SELECT DISTINCT t.id, t.mbid +-- The worker's track queue: played, present tracks with a recording MBID that +-- ListenBrainz has never answered for, or not in 30 days. Never-fetched first, +-- then the oldest answer. Freshness is read from track_similarity_fetches, not +-- from the edges, so a seed whose answer matched nothing in the library still +-- counts as fetched (#5296); reading the edges is what let 25 such seeds hold +-- the head of the queue forever. +SELECT t.id, t.mbid FROM tracks t -JOIN play_events pe ON pe.track_id = t.id +LEFT JOIN track_similarity_fetches f ON f.track_id = t.id WHERE t.mbid IS NOT NULL - AND NOT EXISTS ( - SELECT 1 FROM track_similarity ts - WHERE ts.track_a_id = t.id - AND ts.source = 'listenbrainz' - AND ts.fetched_at > now() - interval '7 days' - ) -ORDER BY t.id + AND t.missing_since IS NULL + AND EXISTS (SELECT 1 FROM play_events pe WHERE pe.track_id = t.id) + AND (f.track_id IS NULL OR f.fetched_at < now() - interval '30 days') +ORDER BY f.fetched_at NULLS FIRST, t.id LIMIT $1; -- name: ListPlayedArtistsNeedingSimilarity :many -SELECT DISTINCT ar.id, ar.mbid +-- The artist queue, with the same shape and the same reason. +SELECT ar.id, ar.mbid FROM artists ar -JOIN tracks t ON t.artist_id = ar.id -JOIN play_events pe ON pe.track_id = t.id +LEFT JOIN artist_similarity_fetches f ON f.artist_id = ar.id WHERE ar.mbid IS NOT NULL - AND NOT EXISTS ( - SELECT 1 FROM artist_similarity asim - WHERE asim.artist_a_id = ar.id - AND asim.source = 'listenbrainz' - AND asim.fetched_at > now() - interval '7 days' + AND EXISTS ( + SELECT 1 FROM tracks t + JOIN play_events pe ON pe.track_id = t.id + WHERE t.artist_id = ar.id ) -ORDER BY ar.id + AND (f.artist_id IS NULL OR f.fetched_at < now() - interval '30 days') +ORDER BY f.fetched_at NULLS FIRST, ar.id LIMIT $1; --- name: GetTracksByMBIDs :many --- Bulk in-library lookup: maps a slice of MBIDs back to local track IDs. -SELECT id, mbid FROM tracks WHERE mbid = ANY($1::text[]); - -- name: GetArtistsByMBIDs :many SELECT id, mbid FROM artists WHERE mbid = ANY($1::text[]); --- name: UpsertTrackSimilarity :exec -INSERT INTO track_similarity (track_a_id, track_b_id, score, source, fetched_at) -VALUES ($1, $2, $3, 'listenbrainz', now()) -ON CONFLICT (track_a_id, track_b_id, source) -DO UPDATE SET score = EXCLUDED.score, fetched_at = EXCLUDED.fetched_at; +-- name: DeleteSimilarRecordingsForSeed :exec +-- A fetch replaces the seed's cached answer whole. +DELETE FROM listenbrainz_similar_recordings WHERE seed_track_id = $1; + +-- name: InsertSimilarRecordings :exec +-- The seed's whole answer, in or out of the library. An MBID LB lists twice +-- keeps its best score. +INSERT INTO listenbrainz_similar_recordings (seed_track_id, recording_mbid, score) +SELECT sqlc.arg(seed_track_id)::uuid, r.mbid, max(r.score) +FROM (SELECT unnest(sqlc.arg(mbids)::text[]) AS mbid, + unnest(sqlc.arg(scores)::float8[]) AS score) r +WHERE r.mbid <> '' +GROUP BY r.mbid; + +-- name: RecordTrackSimilarityFetch :exec +INSERT INTO track_similarity_fetches (track_id, fetched_at, returned) +VALUES ($1, now(), $2) +ON CONFLICT (track_id) DO UPDATE SET fetched_at = now(), returned = EXCLUDED.returned; + +-- name: RecordArtistSimilarityFetch :exec +INSERT INTO artist_similarity_fetches (artist_id, fetched_at, returned) +VALUES ($1, now(), $2) +ON CONFLICT (artist_id) DO UPDATE SET fetched_at = now(), returned = EXCLUDED.returned; + +-- name: ResolveListenBrainzTrackEdges :exec +-- Derives the listenbrainz edges in track_similarity from the cached answers: +-- every cached recording that is a present track in the library becomes an +-- edge, whenever it got there. Scoped to the given seeds, or every seed when +-- seed_ids is NULL. +-- +-- One track per recording MBID: two copies of a song (an mp3 beside its flac, +-- or one recording on two releases) would otherwise be two candidates, and the +-- pool would offer the same song twice. No cap per seed: LB's answer is already +-- at most 50. +-- +-- Only edges of seeds the cache owns are removed, those with a fetch row. +-- Edges written before the cache existed, or copied onto a survivor by a fold +-- (merge.sql), stay until their seed is fetched. +WITH want AS ( + SELECT DISTINCT ON (c.seed_track_id, c.recording_mbid) + c.seed_track_id AS track_a_id, t.id AS track_b_id, c.score + FROM listenbrainz_similar_recordings c + JOIN tracks t ON t.mbid = c.recording_mbid + AND t.missing_since IS NULL + AND t.id <> c.seed_track_id + WHERE sqlc.narg(seed_ids)::uuid[] IS NULL + OR c.seed_track_id = ANY(sqlc.narg(seed_ids)::uuid[]) + ORDER BY c.seed_track_id, c.recording_mbid, t.id +), +upserted AS ( + INSERT INTO track_similarity (track_a_id, track_b_id, score, source, fetched_at) + SELECT track_a_id, track_b_id, score, 'listenbrainz', now() FROM want + ON CONFLICT (track_a_id, track_b_id, source) DO UPDATE + SET score = EXCLUDED.score, fetched_at = EXCLUDED.fetched_at + WHERE track_similarity.score IS DISTINCT FROM EXCLUDED.score + RETURNING 1 +) +DELETE FROM track_similarity ts +WHERE ts.source = 'listenbrainz' + AND (sqlc.narg(seed_ids)::uuid[] IS NULL OR ts.track_a_id = ANY(sqlc.narg(seed_ids)::uuid[])) + AND EXISTS (SELECT 1 FROM track_similarity_fetches f WHERE f.track_id = ts.track_a_id) + AND NOT EXISTS ( + SELECT 1 FROM want w + WHERE w.track_a_id = ts.track_a_id AND w.track_b_id = ts.track_b_id + ); -- name: UpsertArtistSimilarity :exec INSERT INTO artist_similarity (artist_a_id, artist_b_id, score, source, fetched_at) diff --git a/internal/dbtest/reset.go b/internal/dbtest/reset.go index abcb16ae..4dc79322 100644 --- a/internal/dbtest/reset.go +++ b/internal/dbtest/reset.go @@ -40,7 +40,10 @@ const TestUserPrefix = "test-" var dataTables = []string{ "artist_similarity", "track_similarity", - "artist_similarity_unmatched", // M5c + "artist_similarity_unmatched", // M5c + "listenbrainz_similar_recordings", // #5296 + "track_similarity_fetches", + "artist_similarity_fetches", "scrobble_queue", "contextual_likes", "general_likes_albums", diff --git a/internal/library/acoustid_lookup_test.go b/internal/library/acoustid_lookup_test.go index b9fe4cd7..be51fd9d 100644 --- a/internal/library/acoustid_lookup_test.go +++ b/internal/library/acoustid_lookup_test.go @@ -349,9 +349,25 @@ func TestAcoustIDLookup_FilledTrackReachesSimilarity_Integration(t *testing.T) { if n := seeds(); n != 1 { t.Errorf("after the lookup, %d similarity seeds, want the played track", n) } - got, err := q.GetTracksByMBIDs(ctx, []string{"rec-played"}) - if err != nil || len(got) != 1 || got[0].ID != tr.ID { - t.Errorf("GetTracksByMBIDs = %+v (err %v), want the played track", got, err) + + // And it is reachable as a target: an answer cached for another seed + // before the lookup filled the MBID now resolves to it, with no new + // fetch (#5296). + other, _, _ := seedTrack(t, pool, filepath.Join(dir, "other.mp3")) + if _, err := pool.Exec(ctx, `INSERT INTO listenbrainz_similar_recordings (seed_track_id, recording_mbid, score) + VALUES ($1, 'rec-played', 0.8)`, other.ID); err != nil { + t.Fatal(err) + } + if err := q.ResolveListenBrainzTrackEdges(ctx, nil); err != nil { + t.Fatal(err) + } + var edges int + if err := pool.QueryRow(ctx, `SELECT count(*) FROM track_similarity + WHERE track_a_id = $1 AND track_b_id = $2 AND source = 'listenbrainz'`, other.ID, tr.ID).Scan(&edges); err != nil { + t.Fatal(err) + } + if edges != 1 { + t.Errorf("cached answer naming the filled MBID resolved to %d edges, want 1", edges) } } diff --git a/internal/similarity/worker.go b/internal/similarity/worker.go index ba6ff365..c2d5078a 100644 --- a/internal/similarity/worker.go +++ b/internal/similarity/worker.go @@ -1,18 +1,29 @@ // 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, filters returned MBIDs to the local -// library, and stores the top-K edges in track_similarity / -// artist_similarity for M4c's radio candidate-pool builder. +// 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" @@ -20,23 +31,21 @@ import ( "git.fabledsword.com/bvandeusen/minstrel/internal/scrobble/listenbrainz" ) -// Worker drains played-tracks-and-artists needing similarity and POSTs -// the results into track_similarity / artist_similarity. Failures are -// passively retried via the timer (no durable queue table — losing one -// tick's worth of refresh attempts is "1 hour of staleness," fine). +// 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 + 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; -// at 25/h a freshly-played library converges in hours, not days, while -// staying well under ListenBrainz rate limits (429s abort the 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, @@ -64,15 +73,23 @@ func (w *Worker) Run(ctx context.Context) { } } -// tickOnce drains one batch of tracks and one batch of artists. Per-row -// errors are logged and skipped (passive retry via timer). 429 aborts the -// entire tick. +// 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 } - return w.tickArtists(ctx, q) + 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 { @@ -85,16 +102,23 @@ func (w *Worker) tickTracks(ctx context.Context, q *dbq.Queries) error { continue // defensive — query already filters NULL } results, err := w.client.SimilarRecordings(ctx, *r.Mbid, 100) - if err != nil { - var ra *listenbrainz.RetryAfterError - if errors.As(err, &ra) { - w.logger.Warn("similarity: 429 — aborting tick", "retry_after", ra.Wait) - return nil - } + 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 } - w.upsertTrackSimilar(ctx, q, r.ID, results) + 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 } @@ -109,64 +133,65 @@ func (w *Worker) tickArtists(ctx context.Context, q *dbq.Queries) error { continue } results, err := w.client.SimilarArtists(ctx, *r.Mbid, 100) - if err != nil { - var ra *listenbrainz.RetryAfterError - if errors.As(err, &ra) { - w.logger.Warn("similarity: 429 — aborting tick", "retry_after", ra.Wait) - return nil - } + 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 } -// upsertTrackSimilar filters returned MBIDs to those in our library, takes -// top-K by score, and upserts rows. -func (w *Worker) upsertTrackSimilar(ctx context.Context, q *dbq.Queries, trackAID pgtype.UUID, results []listenbrainz.SimilarRecording) { - if len(results) == 0 { - return - } - sort.Slice(results, func(i, j int) bool { return results[i].Score > results[j].Score }) +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) } - rows, err := q.GetTracksByMBIDs(ctx, mbids) - if err != nil { - w.logger.Warn("similarity: GetTracksByMBIDs", "err", err) - return - } - idByMBID := make(map[string]pgtype.UUID, len(rows)) - for _, r := range rows { - if r.Mbid != nil { - idByMBID[*r.Mbid] = r.ID + 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) } - } - - taken := 0 - for _, r := range results { - if taken >= w.topK { - break + 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) + } } - localID, ok := idByMBID[r.MBID] - if !ok { - continue + if err := q.RecordTrackSimilarityFetch(ctx, dbq.RecordTrackSimilarityFetchParams{ + TrackID: seedID, Returned: int32(len(results)), + }); err != nil { + return fmt.Errorf("record fetch: %w", err) } - if localID == trackAID { - continue // defensive — DB CHECK constraint also catches self-edges + if err := q.ResolveListenBrainzTrackEdges(ctx, []pgtype.UUID{seedID}); err != nil { + return fmt.Errorf("resolve edges: %w", err) } - if uerr := q.UpsertTrackSimilarity(ctx, dbq.UpsertTrackSimilarityParams{ - TrackAID: trackAID, TrackBID: localID, Score: r.Score, - }); uerr != nil { - w.logger.Warn("similarity: UpsertTrackSimilarity", "err", uerr) - continue - } - taken++ - } + return nil + }) } func (w *Worker) upsertArtistSimilar(ctx context.Context, q *dbq.Queries, artistAID pgtype.UUID, results []listenbrainz.SimilarArtist) { diff --git a/internal/similarity/worker_integration_test.go b/internal/similarity/worker_integration_test.go index d323d978..940bff5c 100644 --- a/internal/similarity/worker_integration_test.go +++ b/internal/similarity/worker_integration_test.go @@ -10,6 +10,7 @@ import ( "net/http/httptest" "os" "strings" + "sync" "sync/atomic" "testing" "time" @@ -208,7 +209,9 @@ func TestTickOnce_MapsLBResponseToLocalLibrary(t *testing.T) { } } -func TestTickOnce_TopKEnforced(t *testing.T) { +// The 20-per-seed cap is gone with the cache (#5296): LB's answer is at most +// 50, and every recording of it that is in the library becomes an edge. +func TestTickOnce_KeepsEveryInLibraryMatch(t *testing.T) { f := newFixture(t) mbidSeed := "11111111-1111-1111-1111-111111111111" seed := seedTrack(t, f, "Seed", &mbidSeed) @@ -230,30 +233,25 @@ func TestTickOnce_TopKEnforced(t *testing.T) { if err := w.tickOnce(context.Background()); err != nil { t.Fatalf("tickOnce: %v", err) } - if got := countTrackSim(t, f, seed.ID); got != 20 { - t.Errorf("got %d rows, want exactly 20 (all 25 in-library; top 20 by score)", got) + if got := countTrackSim(t, f, seed.ID); got != 25 { + t.Errorf("got %d rows, want all 25 in-library matches", got) } } -func TestTickOnce_RespectsSevenDayCap(t *testing.T) { +func TestTickOnce_RecentFetchIsNotRepeated(t *testing.T) { f := newFixture(t) mbid := "11111111-1111-1111-1111-111111111111" trackA := seedTrack(t, f, "A", &mbid) markPlayed(t, f, trackA.ID) - otherMbid := "55555555-5555-5555-5555-555555555555" - other := seedTrack(t, f, "Other", &otherMbid) - if _, err := f.pool.Exec(context.Background(), - `INSERT INTO track_similarity (track_a_id, track_b_id, score, source, fetched_at) - VALUES ($1, $2, 0.5, 'listenbrainz', now())`, trackA.ID, other.ID); err != nil { - t.Fatalf("seed sim: %v", err) - } - beforeCount := countTrackSim(t, f, trackA.ID) - srv := stubLB(`[{"recording_mbid":"`+otherMbid+`","score":0.99}]`, `[]`, http.StatusOK) + recordFetch(t, f, trackA.ID, "now() - interval '29 days'") + srv, asked := recordingLB(func(string) string { return `[]` }) defer srv.Close() w := newTestWorker(f, srv.URL) - _ = w.tickOnce(context.Background()) - if got := countTrackSim(t, f, trackA.ID); got != beforeCount { - t.Errorf("worker re-queried fresh track: before=%d after=%d", beforeCount, got) + if err := w.tickOnce(context.Background()); err != nil { + t.Fatalf("tickOnce: %v", err) + } + if n := len(asked()); n != 0 { + t.Errorf("a seed fetched 29 days ago was asked again (%d requests)", n) } } @@ -266,9 +264,10 @@ func TestTickOnce_RefreshesStaleRow(t *testing.T) { other := seedTrack(t, f, "Other", &otherMbid) if _, err := f.pool.Exec(context.Background(), `INSERT INTO track_similarity (track_a_id, track_b_id, score, source, fetched_at) - VALUES ($1, $2, 0.5, 'listenbrainz', now() - interval '8 days')`, trackA.ID, other.ID); err != nil { + VALUES ($1, $2, 0.5, 'listenbrainz', now() - interval '31 days')`, trackA.ID, other.ID); err != nil { t.Fatalf("seed sim: %v", err) } + recordFetch(t, f, trackA.ID, "now() - interval '31 days'") srv := stubLB(`[{"recording_mbid":"`+otherMbid+`","score":0.95}]`, `[]`, http.StatusOK) defer srv.Close() w := newTestWorker(f, srv.URL) @@ -467,3 +466,147 @@ func TestUpsertArtistSimilar_PersistsUnmatchedToTable(t *testing.T) { _ = inLibArtist } + +// recordFetch stamps a seed as answered by ListenBrainz at the given SQL time. +func recordFetch(t *testing.T, f fixture, trackID pgtype.UUID, at string) { + t.Helper() + if _, err := f.pool.Exec(context.Background(), + `INSERT INTO track_similarity_fetches (track_id, fetched_at, returned) VALUES ($1, `+at+`, 0)`, + trackID); err != nil { + t.Fatalf("record fetch: %v", err) + } +} + +// recordingLB answers similar-recordings with answer(mbid) and similar-artists +// with an empty list, and reports which recording MBIDs it was asked about. +func recordingLB(answer func(mbid string) string) (*httptest.Server, func() []string) { + var mu sync.Mutex + var asked []string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if !strings.Contains(r.URL.Path, "/similar-recordings/") { + _, _ = w.Write([]byte(`[]`)) + return + } + mbid := r.URL.Query().Get("recording_mbids") + mu.Lock() + asked = append(asked, mbid) + mu.Unlock() + _, _ = w.Write([]byte(answer(mbid))) + })) + return srv, func() []string { + mu.Lock() + defer mu.Unlock() + return append([]string(nil), asked...) + } +} + +// The starvation #3879 measured: a seed whose answer matched nothing in the +// library wrote nothing, stayed due, and was asked again every tick ahead of +// everything behind it. Recording the fetch moves the queue on. +func TestTickOnce_UnmatchedAnswerDoesNotHoldTheQueue(t *testing.T) { + f := newFixture(t) + mbids := []string{ + "11111111-1111-1111-1111-111111111111", + "22222222-2222-2222-2222-222222222222", + "33333333-3333-3333-3333-333333333333", + } + for i := range mbids { + tr := seedTrack(t, f, fmt.Sprintf("S%d", i), &mbids[i]) + markPlayed(t, f, tr.ID) + } + srv, asked := recordingLB(func(string) string { + return `[{"recording_mbid":"99999999-9999-9999-9999-999999999999","score":0.9}]` + }) + defer srv.Close() + w := newTestWorker(f, srv.URL) + w.batch = 1 + for i := 0; i < len(mbids); i++ { + if err := w.tickOnce(context.Background()); err != nil { + t.Fatalf("tick %d: %v", i, err) + } + } + seen := map[string]bool{} + for _, m := range asked() { + seen[m] = true + } + if len(seen) != len(mbids) { + t.Errorf("three ticks of one asked about %d distinct seeds (%v), want all %d", + len(seen), asked(), len(mbids)) + } + var fetched int + _ = f.pool.QueryRow(context.Background(), `SELECT count(*) FROM track_similarity_fetches`).Scan(&fetched) + if fetched != len(mbids) { + t.Errorf("track_similarity_fetches = %d, want %d", fetched, len(mbids)) + } +} + +// The point of the cache (#5296): a recording LB named while it was not in the +// library becomes an edge once it arrives, on the next tick's resolve, without +// asking LB again. One edge per recording when the library holds two copies, +// and none to a copy that is missing. +func TestTickOnce_CachedAnswerResolvesWhenTheRecordingArrives(t *testing.T) { + f := newFixture(t) + ctx := context.Background() + seedMbid := "11111111-1111-1111-1111-111111111111" + laterMbid := "77777777-7777-7777-7777-777777777777" + goneMbid := "88888888-8888-8888-8888-888888888888" + seed := seedTrack(t, f, "Seed", &seedMbid) + markPlayed(t, f, seed.ID) + srv, asked := recordingLB(func(string) string { + return `[{"recording_mbid":"` + laterMbid + `","score":0.9}, + {"recording_mbid":"` + goneMbid + `","score":0.8}]` + }) + defer srv.Close() + w := newTestWorker(f, srv.URL) + if err := w.tickOnce(ctx); err != nil { + t.Fatalf("tick 1: %v", err) + } + if got := countTrackSim(t, f, seed.ID); got != 0 { + t.Fatalf("edges before the recording is in the library = %d, want 0", got) + } + var cached int + _ = f.pool.QueryRow(ctx, `SELECT count(*) FROM listenbrainz_similar_recordings WHERE seed_track_id = $1`, + seed.ID).Scan(&cached) + if cached != 2 { + t.Fatalf("cached answer rows = %d, want 2 (both out of library)", cached) + } + + // The recording arrives twice (an mp3 beside its flac); the other one + // arrives only as a missing row. + _ = seedTrack(t, f, "Later mp3", &laterMbid) + _ = seedTrack(t, f, "Later flac", &laterMbid) + gone := seedTrack(t, f, "Gone", &goneMbid) + if _, err := f.pool.Exec(ctx, `UPDATE tracks SET missing_since = now() WHERE id = $1`, gone.ID); err != nil { + t.Fatal(err) + } + if err := w.tickOnce(ctx); err != nil { + t.Fatalf("tick 2: %v", err) + } + if got := countTrackSim(t, f, seed.ID); got != 1 { + t.Errorf("edges after the recording arrived = %d, want 1 (one per recording, none to a missing copy)", got) + } + if n := len(asked()); n != 1 { + t.Errorf("LB asked %d times, want 1: the resolve needs no new fetch", n) + } +} + +// A 4xx answer is permanent for that MBID; it is recorded like an empty one +// so the seed waits its turn instead of being re-asked every hour. +func TestTickOnce_PermanentErrorCountsAsFetched(t *testing.T) { + f := newFixture(t) + mbid := "11111111-1111-1111-1111-111111111111" + trackA := seedTrack(t, f, "A", &mbid) + markPlayed(t, f, trackA.ID) + srv := stubLB(``, ``, http.StatusBadRequest) + defer srv.Close() + w := newTestWorker(f, srv.URL) + if err := w.tickOnce(context.Background()); err != nil { + t.Fatalf("tickOnce: %v", err) + } + var tracks, artists int + _ = f.pool.QueryRow(context.Background(), `SELECT count(*) FROM track_similarity_fetches`).Scan(&tracks) + _ = f.pool.QueryRow(context.Background(), `SELECT count(*) FROM artist_similarity_fetches`).Scan(&artists) + if tracks != 1 || artists != 1 { + t.Errorf("fetches recorded: tracks=%d artists=%d, want 1 and 1", tracks, artists) + } +}