Similarity cache and round-robin artist arms (#5296, #5297) #142

Merged
bvandeusen merged 2 commits from dev into main 2026-10-07 20:03:00 -04:00
12 changed files with 676 additions and 207 deletions
+18
View File
@@ -241,6 +241,12 @@ type ArtistSimilarity struct {
FetchedAt pgtype.Timestamptz FetchedAt pgtype.Timestamptz
} }
type ArtistSimilarityFetch struct {
ArtistID pgtype.UUID
FetchedAt pgtype.Timestamptz
Returned int32
}
type ArtistSimilarityUnmatched struct { type ArtistSimilarityUnmatched struct {
SeedArtistID pgtype.UUID SeedArtistID pgtype.UUID
CandidateMbid string CandidateMbid string
@@ -436,6 +442,12 @@ type LidarrRequest struct {
LidarrAddConfirmedAt pgtype.Timestamptz LidarrAddConfirmedAt pgtype.Timestamptz
} }
type ListenbrainzSimilarRecording struct {
SeedTrackID pgtype.UUID
RecordingMbid string
Score float64
}
type LoudnessSetting struct { type LoudnessSetting struct {
ID bool ID bool
Enabled bool Enabled bool
@@ -771,6 +783,12 @@ type TrackSimilarity struct {
FetchedAt pgtype.Timestamptz FetchedAt pgtype.Timestamptz
} }
type TrackSimilarityFetch struct {
TrackID pgtype.UUID
FetchedAt pgtype.Timestamptz
Returned int32
}
type TrackTag struct { type TrackTag struct {
TrackID pgtype.UUID TrackID pgtype.UUID
Tag string Tag string
+22 -6
View File
@@ -828,13 +828,23 @@ lb_similar AS (
LIMIT $5 LIMIT $5
), ),
similar_artists AS ( similar_artists AS (
SELECT t.id AS track_id, asim.score * 0.5 AS sim_score -- Round-robin across the similar artists (#5297): each artist's first
-- track, best-scoring artist first, then each one's second, and so on.
-- Ordering by artist score alone let the closest similar artist's
-- catalogue fill the whole LIMIT; on the operator's library this arm held
-- exactly one artist for all 17 seeds measured (#3879).
SELECT track_id, sim_score
FROM (
SELECT t.id AS track_id, asim.score * 0.5 AS sim_score, asim.score AS artist_score,
row_number() OVER (PARTITION BY asim.artist_b_id
ORDER BY md5(t.id::text || $12::text)) AS artist_turn
FROM artist_similarity asim FROM artist_similarity asim
JOIN tracks t ON t.artist_id = asim.artist_b_id JOIN tracks t ON t.artist_id = asim.artist_b_id AND t.missing_since IS NULL
JOIN seed_artist sa ON asim.artist_a_id = sa.artist_id JOIN seed_artist sa ON asim.artist_a_id = sa.artist_id
WHERE asim.source = 'listenbrainz' WHERE asim.source = 'listenbrainz'
AND t.id NOT IN (SELECT id FROM excluded_ids) AND t.id NOT IN (SELECT id FROM excluded_ids)
ORDER BY asim.score DESC, md5(t.id::text || $12::text) ) ranked
ORDER BY artist_turn, artist_score DESC, md5(track_id::text || $12::text)
LIMIT $6 LIMIT $6
), ),
tag_overlap AS ( tag_overlap AS (
@@ -882,14 +892,20 @@ coplay_artists AS (
-- the coplay worker). Mirrors similar_artists but from local co-occurrence -- the coplay worker). Mirrors similar_artists but from local co-occurrence
-- instead of ListenBrainz; empty on single-user servers. Same 0.5 damp as -- instead of ListenBrainz; empty on single-user servers. Same 0.5 damp as
-- similar_artists since it's artist-level. -- similar_artists since it's artist-level.
SELECT t.id AS track_id, asim.score * 0.5 AS sim_score -- Round-robin across artists for the same reason as similar_artists.
SELECT track_id, sim_score
FROM (
SELECT t.id AS track_id, asim.score * 0.5 AS sim_score, asim.score AS artist_score,
row_number() OVER (PARTITION BY asim.artist_b_id
ORDER BY md5(t.id::text || $12::text)) AS artist_turn
FROM artist_similarity asim FROM artist_similarity asim
JOIN tracks t ON t.artist_id = asim.artist_b_id JOIN tracks t ON t.artist_id = asim.artist_b_id AND t.missing_since IS NULL
JOIN seed_artist sa ON asim.artist_a_id = sa.artist_id JOIN seed_artist sa ON asim.artist_a_id = sa.artist_id
WHERE asim.source = 'user_cooccurrence' WHERE asim.source = 'user_cooccurrence'
AND t.id NOT IN (SELECT id FROM excluded_ids) AND t.id NOT IN (SELECT id FROM excluded_ids)
AND t.id <> $2 AND t.id <> $2
ORDER BY asim.score DESC, md5(t.id::text || $12::text) ) ranked
ORDER BY artist_turn, artist_score DESC, md5(track_id::text || $12::text)
LIMIT $11 LIMIT $11
), ),
random_fill AS ( random_fill AS (
+127 -63
View File
@@ -11,6 +11,16 @@ import (
"github.com/jackc/pgx/v5/pgtype" "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 const getArtistsByMBIDs = `-- name: GetArtistsByMBIDs :many
SELECT id, mbid FROM artists WHERE mbid = ANY($1::text[]) 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 return items, nil
} }
const getTracksByMBIDs = `-- name: GetTracksByMBIDs :many const insertSimilarRecordings = `-- name: InsertSimilarRecordings :exec
SELECT id, mbid FROM tracks WHERE mbid = ANY($1::text[]) 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 { type InsertSimilarRecordingsParams struct {
ID pgtype.UUID SeedTrackID pgtype.UUID
Mbid *string Mbids []string
Scores []float64
} }
// Bulk in-library lookup: maps a slice of MBIDs back to local track IDs. // The seed's whole answer, in or out of the library. An MBID LB lists twice
func (q *Queries) GetTracksByMBIDs(ctx context.Context, dollar_1 []string) ([]GetTracksByMBIDsRow, error) { // keeps its best score.
rows, err := q.db.Query(ctx, getTracksByMBIDs, dollar_1) func (q *Queries) InsertSimilarRecordings(ctx context.Context, arg InsertSimilarRecordingsParams) error {
if err != nil { _, err := q.db.Exec(ctx, insertSimilarRecordings, arg.SeedTrackID, arg.Mbids, arg.Scores)
return nil, err return 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
} }
const listPlayedArtistsNeedingSimilarity = `-- name: ListPlayedArtistsNeedingSimilarity :many const listPlayedArtistsNeedingSimilarity = `-- name: ListPlayedArtistsNeedingSimilarity :many
SELECT DISTINCT ar.id, ar.mbid SELECT ar.id, ar.mbid
FROM artists ar FROM artists ar
JOIN tracks t ON t.artist_id = ar.id LEFT JOIN artist_similarity_fetches f ON f.artist_id = ar.id
JOIN play_events pe ON pe.track_id = t.id
WHERE ar.mbid IS NOT NULL WHERE ar.mbid IS NOT NULL
AND NOT EXISTS ( AND EXISTS (
SELECT 1 FROM artist_similarity asim SELECT 1 FROM tracks t
WHERE asim.artist_a_id = ar.id JOIN play_events pe ON pe.track_id = t.id
AND asim.source = 'listenbrainz' WHERE t.artist_id = ar.id
AND asim.fetched_at > now() - interval '7 days'
) )
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 LIMIT $1
` `
@@ -91,6 +92,7 @@ type ListPlayedArtistsNeedingSimilarityRow struct {
Mbid *string Mbid *string
} }
// The artist queue, with the same shape and the same reason.
func (q *Queries) ListPlayedArtistsNeedingSimilarity(ctx context.Context, limit int32) ([]ListPlayedArtistsNeedingSimilarityRow, error) { func (q *Queries) ListPlayedArtistsNeedingSimilarity(ctx context.Context, limit int32) ([]ListPlayedArtistsNeedingSimilarityRow, error) {
rows, err := q.db.Query(ctx, listPlayedArtistsNeedingSimilarity, limit) rows, err := q.db.Query(ctx, listPlayedArtistsNeedingSimilarity, limit)
if err != nil { if err != nil {
@@ -112,17 +114,14 @@ func (q *Queries) ListPlayedArtistsNeedingSimilarity(ctx context.Context, limit
} }
const listPlayedTracksNeedingSimilarity = `-- name: ListPlayedTracksNeedingSimilarity :many const listPlayedTracksNeedingSimilarity = `-- name: ListPlayedTracksNeedingSimilarity :many
SELECT DISTINCT t.id, t.mbid SELECT t.id, t.mbid
FROM tracks t 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 WHERE t.mbid IS NOT NULL
AND NOT EXISTS ( AND t.missing_since IS NULL
SELECT 1 FROM track_similarity ts AND EXISTS (SELECT 1 FROM play_events pe WHERE pe.track_id = t.id)
WHERE ts.track_a_id = t.id AND (f.track_id IS NULL OR f.fetched_at < now() - interval '30 days')
AND ts.source = 'listenbrainz' ORDER BY f.fetched_at NULLS FIRST, t.id
AND ts.fetched_at > now() - interval '7 days'
)
ORDER BY t.id
LIMIT $1 LIMIT $1
` `
@@ -131,9 +130,12 @@ type ListPlayedTracksNeedingSimilarityRow struct {
Mbid *string Mbid *string
} }
// Tracks with at least one play, an MBID, and no fresh listenbrainz row // The worker's track queue: played, present tracks with a recording MBID that
// (no row at all, OR fetched_at older than 7 days). Used by the worker // ListenBrainz has never answered for, or not in 30 days. Never-fetched first,
// to find work each tick. Bounded by $1. // 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) { func (q *Queries) ListPlayedTracksNeedingSimilarity(ctx context.Context, limit int32) ([]ListPlayedTracksNeedingSimilarityRow, error) {
rows, err := q.db.Query(ctx, listPlayedTracksNeedingSimilarity, limit) rows, err := q.db.Query(ctx, listPlayedTracksNeedingSimilarity, limit)
if err != nil { if err != nil {
@@ -225,6 +227,86 @@ func (q *Queries) ListSimilarArtistsForArtist(ctx context.Context, arg ListSimil
return items, nil 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 const upsertArtistSimilarity = `-- name: UpsertArtistSimilarity :exec
INSERT INTO artist_similarity (artist_a_id, artist_b_id, score, source, fetched_at) INSERT INTO artist_similarity (artist_a_id, artist_b_id, score, source, fetched_at)
VALUES ($1, $2, $3, 'listenbrainz', now()) VALUES ($1, $2, $3, 'listenbrainz', now())
@@ -274,21 +356,3 @@ func (q *Queries) UpsertArtistSimilarityUnmatched(ctx context.Context, arg Upser
) )
return err 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
}
@@ -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;
@@ -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
);
+22 -6
View File
@@ -96,13 +96,23 @@ lb_similar AS (
LIMIT $5 LIMIT $5
), ),
similar_artists AS ( similar_artists AS (
SELECT t.id AS track_id, asim.score * 0.5 AS sim_score -- Round-robin across the similar artists (#5297): each artist's first
-- track, best-scoring artist first, then each one's second, and so on.
-- Ordering by artist score alone let the closest similar artist's
-- catalogue fill the whole LIMIT; on the operator's library this arm held
-- exactly one artist for all 17 seeds measured (#3879).
SELECT track_id, sim_score
FROM (
SELECT t.id AS track_id, asim.score * 0.5 AS sim_score, asim.score AS artist_score,
row_number() OVER (PARTITION BY asim.artist_b_id
ORDER BY md5(t.id::text || $12::text)) AS artist_turn
FROM artist_similarity asim FROM artist_similarity asim
JOIN tracks t ON t.artist_id = asim.artist_b_id JOIN tracks t ON t.artist_id = asim.artist_b_id AND t.missing_since IS NULL
JOIN seed_artist sa ON asim.artist_a_id = sa.artist_id JOIN seed_artist sa ON asim.artist_a_id = sa.artist_id
WHERE asim.source = 'listenbrainz' WHERE asim.source = 'listenbrainz'
AND t.id NOT IN (SELECT id FROM excluded_ids) AND t.id NOT IN (SELECT id FROM excluded_ids)
ORDER BY asim.score DESC, md5(t.id::text || $12::text) ) ranked
ORDER BY artist_turn, artist_score DESC, md5(track_id::text || $12::text)
LIMIT $6 LIMIT $6
), ),
tag_overlap AS ( tag_overlap AS (
@@ -150,14 +160,20 @@ coplay_artists AS (
-- the coplay worker). Mirrors similar_artists but from local co-occurrence -- the coplay worker). Mirrors similar_artists but from local co-occurrence
-- instead of ListenBrainz; empty on single-user servers. Same 0.5 damp as -- instead of ListenBrainz; empty on single-user servers. Same 0.5 damp as
-- similar_artists since it's artist-level. -- similar_artists since it's artist-level.
SELECT t.id AS track_id, asim.score * 0.5 AS sim_score -- Round-robin across artists for the same reason as similar_artists.
SELECT track_id, sim_score
FROM (
SELECT t.id AS track_id, asim.score * 0.5 AS sim_score, asim.score AS artist_score,
row_number() OVER (PARTITION BY asim.artist_b_id
ORDER BY md5(t.id::text || $12::text)) AS artist_turn
FROM artist_similarity asim FROM artist_similarity asim
JOIN tracks t ON t.artist_id = asim.artist_b_id JOIN tracks t ON t.artist_id = asim.artist_b_id AND t.missing_since IS NULL
JOIN seed_artist sa ON asim.artist_a_id = sa.artist_id JOIN seed_artist sa ON asim.artist_a_id = sa.artist_id
WHERE asim.source = 'user_cooccurrence' WHERE asim.source = 'user_cooccurrence'
AND t.id NOT IN (SELECT id FROM excluded_ids) AND t.id NOT IN (SELECT id FROM excluded_ids)
AND t.id <> $2 AND t.id <> $2
ORDER BY asim.score DESC, md5(t.id::text || $12::text) ) ranked
ORDER BY artist_turn, artist_score DESC, md5(track_id::text || $12::text)
LIMIT $11 LIMIT $11
), ),
random_fill AS ( random_fill AS (
+86 -30
View File
@@ -1,47 +1,103 @@
-- name: ListPlayedTracksNeedingSimilarity :many -- name: ListPlayedTracksNeedingSimilarity :many
-- Tracks with at least one play, an MBID, and no fresh listenbrainz row -- The worker's track queue: played, present tracks with a recording MBID that
-- (no row at all, OR fetched_at older than 7 days). Used by the worker -- ListenBrainz has never answered for, or not in 30 days. Never-fetched first,
-- to find work each tick. Bounded by $1. -- then the oldest answer. Freshness is read from track_similarity_fetches, not
SELECT DISTINCT t.id, t.mbid -- 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 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 WHERE t.mbid IS NOT NULL
AND NOT EXISTS ( AND t.missing_since IS NULL
SELECT 1 FROM track_similarity ts AND EXISTS (SELECT 1 FROM play_events pe WHERE pe.track_id = t.id)
WHERE ts.track_a_id = t.id AND (f.track_id IS NULL OR f.fetched_at < now() - interval '30 days')
AND ts.source = 'listenbrainz' ORDER BY f.fetched_at NULLS FIRST, t.id
AND ts.fetched_at > now() - interval '7 days'
)
ORDER BY t.id
LIMIT $1; LIMIT $1;
-- name: ListPlayedArtistsNeedingSimilarity :many -- 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 FROM artists ar
JOIN tracks t ON t.artist_id = ar.id LEFT JOIN artist_similarity_fetches f ON f.artist_id = ar.id
JOIN play_events pe ON pe.track_id = t.id
WHERE ar.mbid IS NOT NULL WHERE ar.mbid IS NOT NULL
AND NOT EXISTS ( AND EXISTS (
SELECT 1 FROM artist_similarity asim SELECT 1 FROM tracks t
WHERE asim.artist_a_id = ar.id JOIN play_events pe ON pe.track_id = t.id
AND asim.source = 'listenbrainz' WHERE t.artist_id = ar.id
AND asim.fetched_at > now() - interval '7 days'
) )
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; 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 -- name: GetArtistsByMBIDs :many
SELECT id, mbid FROM artists WHERE mbid = ANY($1::text[]); SELECT id, mbid FROM artists WHERE mbid = ANY($1::text[]);
-- name: UpsertTrackSimilarity :exec -- name: DeleteSimilarRecordingsForSeed :exec
INSERT INTO track_similarity (track_a_id, track_b_id, score, source, fetched_at) -- A fetch replaces the seed's cached answer whole.
VALUES ($1, $2, $3, 'listenbrainz', now()) DELETE FROM listenbrainz_similar_recordings WHERE seed_track_id = $1;
ON CONFLICT (track_a_id, track_b_id, source)
DO UPDATE SET score = EXCLUDED.score, fetched_at = EXCLUDED.fetched_at; -- 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 -- name: UpsertArtistSimilarity :exec
INSERT INTO artist_similarity (artist_a_id, artist_b_id, score, source, fetched_at) INSERT INTO artist_similarity (artist_a_id, artist_b_id, score, source, fetched_at)
+3
View File
@@ -41,6 +41,9 @@ var dataTables = []string{
"artist_similarity", "artist_similarity",
"track_similarity", "track_similarity",
"artist_similarity_unmatched", // M5c "artist_similarity_unmatched", // M5c
"listenbrainz_similar_recordings", // #5296
"track_similarity_fetches",
"artist_similarity_fetches",
"scrobble_queue", "scrobble_queue",
"contextual_likes", "contextual_likes",
"general_likes_albums", "general_likes_albums",
+19 -3
View File
@@ -349,9 +349,25 @@ func TestAcoustIDLookup_FilledTrackReachesSimilarity_Integration(t *testing.T) {
if n := seeds(); n != 1 { if n := seeds(); n != 1 {
t.Errorf("after the lookup, %d similarity seeds, want the played track", n) 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 { // And it is reachable as a target: an answer cached for another seed
t.Errorf("GetTracksByMBIDs = %+v (err %v), want the played track", got, err) // 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)
} }
} }
@@ -429,3 +429,69 @@ func TestLoadCandidatesFromSimilarity_DifferentSeedsCanDrawDifferently(t *testin
"varying with the seed at all", len(seen)) "varying with the seed at all", len(seen))
} }
} }
// #5297: an artist-level arm takes a track from each related artist before a
// second from any. Ordered by artist score alone, the closest artist's
// catalogue filled the whole LIMIT; on the operator's library the
// similar_artists arm held one artist for every seed measured (#3879).
func TestLoadCandidatesFromSimilarity_ArtistArmsRoundRobin(t *testing.T) {
for _, source := range []string{"listenbrainz", "user_cooccurrence"} {
t.Run(source, func(t *testing.T) {
f := newFixture(t, 1)
ctx := context.Background()
seed := f.tracks[0]
artistOf := map[[16]byte]string{}
for i, score := range []float64{0.9, 0.8, 0.7} {
name := fmt.Sprintf("Related %d", i)
ar, err := f.q.UpsertArtist(ctx, dbq.UpsertArtistParams{Name: name, SortName: name})
if err != nil {
t.Fatal(err)
}
al, err := f.q.UpsertAlbum(ctx, dbq.UpsertAlbumParams{Title: name, SortTitle: name, ArtistID: ar.ID})
if err != nil {
t.Fatal(err)
}
// More tracks per artist than the arm's limit, so the closest
// artist alone could fill it.
for j := 0; j < 5; j++ {
tr, err := f.q.UpsertTrack(ctx, dbq.UpsertTrackParams{
Title: fmt.Sprintf("%s #%d", name, j), AlbumID: al.ID, ArtistID: ar.ID,
FilePath: fmt.Sprintf("/tmp/related-%d-%d.flac", i, j), DurationMs: 180_000,
})
if err != nil {
t.Fatal(err)
}
artistOf[tr.ID.Bytes] = name
}
if _, err := f.pool.Exec(ctx,
`INSERT INTO artist_similarity (artist_a_id, artist_b_id, score, source) VALUES ($1, $2, $3, $4)`,
seed.ArtistID, ar.ID, score, source); err != nil {
t.Fatal(err)
}
}
limits := CandidateSourceLimits{}
if source == "listenbrainz" {
limits.SimilarArtist = 3
} else {
limits.UserCoplay = 3
}
got, err := LoadCandidatesFromSimilarity(
ctx, f.q, f.user, seed.ID, 1, SessionVector{Seed: true}, nil, limits, "test-seed",
)
if err != nil {
t.Fatalf("load: %v", err)
}
artists := map[string]int{}
for _, c := range got {
if name, ok := artistOf[c.Track.ID.Bytes]; ok {
artists[name]++
}
}
if len(got) != 3 || len(artists) != 3 {
t.Errorf("arm of 3 drew %d candidates from %d related artists (%v), want one from each of 3",
len(got), len(artists), artists)
}
})
}
}
+85 -60
View File
@@ -1,18 +1,29 @@
// Package similarity owns the inbound ListenBrainz similarity ingest // Package similarity owns the inbound ListenBrainz similarity ingest
// pipeline. A periodic worker queries the LB Labs API // pipeline. A periodic worker queries the LB Labs API
// (labs.api.listenbrainz.org similar-recordings / similar-artists) for // (labs.api.listenbrainz.org similar-recordings / similar-artists) for
// tracks the user has played, filters returned MBIDs to the local // tracks the user has played.
// library, and stores the top-K edges in track_similarity / //
// artist_similarity for M4c's radio candidate-pool builder. // 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 package similarity
import ( import (
"context" "context"
"errors" "errors"
"fmt"
"log/slog" "log/slog"
"sort" "sort"
"time" "time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgtype" "github.com/jackc/pgx/v5/pgtype"
"github.com/jackc/pgx/v5/pgxpool" "github.com/jackc/pgx/v5/pgxpool"
@@ -20,23 +31,21 @@ import (
"git.fabledsword.com/bvandeusen/minstrel/internal/scrobble/listenbrainz" "git.fabledsword.com/bvandeusen/minstrel/internal/scrobble/listenbrainz"
) )
// Worker drains played-tracks-and-artists needing similarity and POSTs // Worker drains played tracks and artists needing similarity. A transient
// the results into track_similarity / artist_similarity. Failures are // failure is not recorded as a fetch, so the timer retries it next tick.
// passively retried via the timer (no durable queue table — losing one
// tick's worth of refresh attempts is "1 hour of staleness," fine).
type Worker struct { type Worker struct {
pool *pgxpool.Pool pool *pgxpool.Pool
client *listenbrainz.Client client *listenbrainz.Client
logger *slog.Logger logger *slog.Logger
tick time.Duration tick time.Duration
batch int32 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, // NewWorker constructs a worker with production defaults: 1h tick,
// batch=25, topK=20. batch is tracks AND artists processed per tick; // batch=25, topK=20. batch is tracks AND artists processed per tick, well
// at 25/h a freshly-played library converges in hours, not days, while // under ListenBrainz rate limits (429s abort the tick). With never-fetched
// staying well under ListenBrainz rate limits (429s abort the tick). // seeds first, a freshly played library is covered in days.
func NewWorker(pool *pgxpool.Pool, client *listenbrainz.Client, logger *slog.Logger) *Worker { func NewWorker(pool *pgxpool.Pool, client *listenbrainz.Client, logger *slog.Logger) *Worker {
return &Worker{ return &Worker{
pool: pool, 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 // tickOnce drains one batch of tracks and one batch of artists, then
// errors are logged and skipped (passive retry via timer). 429 aborts the // re-resolves every cached answer against the library. Per-row errors are
// entire tick. // logged and skipped. 429 aborts the fetching, not the resolve.
func (w *Worker) tickOnce(ctx context.Context) error { func (w *Worker) tickOnce(ctx context.Context) error {
q := dbq.New(w.pool) q := dbq.New(w.pool)
if err := w.tickTracks(ctx, q); err != nil { if err := w.tickTracks(ctx, q); err != nil {
return err 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 { 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 continue // defensive — query already filters NULL
} }
results, err := w.client.SimilarRecordings(ctx, *r.Mbid, 100) results, err := w.client.SimilarRecordings(ctx, *r.Mbid, 100)
if err != nil { switch {
var ra *listenbrainz.RetryAfterError case err == nil:
if errors.As(err, &ra) { case isRateLimited(err):
w.logger.Warn("similarity: 429 — aborting tick", "retry_after", ra.Wait) w.logger.Warn("similarity: 429 — aborting tick")
return nil 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) w.logger.Warn("similarity: similar-recordings failed", "track_id", r.ID, "err", err)
continue 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 return nil
} }
@@ -109,64 +133,65 @@ func (w *Worker) tickArtists(ctx context.Context, q *dbq.Queries) error {
continue continue
} }
results, err := w.client.SimilarArtists(ctx, *r.Mbid, 100) results, err := w.client.SimilarArtists(ctx, *r.Mbid, 100)
if err != nil { switch {
var ra *listenbrainz.RetryAfterError case err == nil:
if errors.As(err, &ra) { case isRateLimited(err):
w.logger.Warn("similarity: 429 — aborting tick", "retry_after", ra.Wait) w.logger.Warn("similarity: 429 — aborting tick")
return nil 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) w.logger.Warn("similarity: similar-artists failed", "artist_id", r.ID, "err", err)
continue continue
} }
w.upsertArtistSimilar(ctx, q, r.ID, results) 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 return nil
} }
// upsertTrackSimilar filters returned MBIDs to those in our library, takes func isRateLimited(err error) bool {
// top-K by score, and upserts rows. var ra *listenbrainz.RetryAfterError
func (w *Worker) upsertTrackSimilar(ctx context.Context, q *dbq.Queries, trackAID pgtype.UUID, results []listenbrainz.SimilarRecording) { return errors.As(err, &ra)
if len(results) == 0 { }
return
}
sort.Slice(results, func(i, j int) bool { return results[i].Score > results[j].Score })
// 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)) mbids := make([]string, 0, len(results))
scores := make([]float64, 0, len(results))
for _, r := range results { for _, r := range results {
mbids = append(mbids, r.MBID) mbids = append(mbids, r.MBID)
scores = append(scores, r.Score)
} }
rows, err := q.GetTracksByMBIDs(ctx, mbids) return pgx.BeginFunc(ctx, w.pool, func(tx pgx.Tx) error {
if err != nil { q := dbq.New(tx)
w.logger.Warn("similarity: GetTracksByMBIDs", "err", err) if err := q.DeleteSimilarRecordingsForSeed(ctx, seedID); err != nil {
return return fmt.Errorf("clear cached answer: %w", err)
} }
idByMBID := make(map[string]pgtype.UUID, len(rows)) if len(mbids) > 0 {
for _, r := range rows { if err := q.InsertSimilarRecordings(ctx, dbq.InsertSimilarRecordingsParams{
if r.Mbid != nil { SeedTrackID: seedID, Mbids: mbids, Scores: scores,
idByMBID[*r.Mbid] = r.ID }); err != nil {
return fmt.Errorf("cache answer: %w", err)
} }
} }
if err := q.RecordTrackSimilarityFetch(ctx, dbq.RecordTrackSimilarityFetchParams{
taken := 0 TrackID: seedID, Returned: int32(len(results)),
for _, r := range results { }); err != nil {
if taken >= w.topK { return fmt.Errorf("record fetch: %w", err)
break
} }
localID, ok := idByMBID[r.MBID] if err := q.ResolveListenBrainzTrackEdges(ctx, []pgtype.UUID{seedID}); err != nil {
if !ok { return fmt.Errorf("resolve edges: %w", err)
continue
}
if localID == trackAID {
continue // defensive — DB CHECK constraint also catches self-edges
}
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) { func (w *Worker) upsertArtistSimilar(ctx context.Context, q *dbq.Queries, artistAID pgtype.UUID, results []listenbrainz.SimilarArtist) {
+160 -17
View File
@@ -10,6 +10,7 @@ import (
"net/http/httptest" "net/http/httptest"
"os" "os"
"strings" "strings"
"sync"
"sync/atomic" "sync/atomic"
"testing" "testing"
"time" "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) f := newFixture(t)
mbidSeed := "11111111-1111-1111-1111-111111111111" mbidSeed := "11111111-1111-1111-1111-111111111111"
seed := seedTrack(t, f, "Seed", &mbidSeed) seed := seedTrack(t, f, "Seed", &mbidSeed)
@@ -230,30 +233,25 @@ func TestTickOnce_TopKEnforced(t *testing.T) {
if err := w.tickOnce(context.Background()); err != nil { if err := w.tickOnce(context.Background()); err != nil {
t.Fatalf("tickOnce: %v", err) t.Fatalf("tickOnce: %v", err)
} }
if got := countTrackSim(t, f, seed.ID); got != 20 { if got := countTrackSim(t, f, seed.ID); got != 25 {
t.Errorf("got %d rows, want exactly 20 (all 25 in-library; top 20 by score)", got) 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) f := newFixture(t)
mbid := "11111111-1111-1111-1111-111111111111" mbid := "11111111-1111-1111-1111-111111111111"
trackA := seedTrack(t, f, "A", &mbid) trackA := seedTrack(t, f, "A", &mbid)
markPlayed(t, f, trackA.ID) markPlayed(t, f, trackA.ID)
otherMbid := "55555555-5555-5555-5555-555555555555" recordFetch(t, f, trackA.ID, "now() - interval '29 days'")
other := seedTrack(t, f, "Other", &otherMbid) srv, asked := recordingLB(func(string) string { return `[]` })
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)
defer srv.Close() defer srv.Close()
w := newTestWorker(f, srv.URL) w := newTestWorker(f, srv.URL)
_ = w.tickOnce(context.Background()) if err := w.tickOnce(context.Background()); err != nil {
if got := countTrackSim(t, f, trackA.ID); got != beforeCount { t.Fatalf("tickOnce: %v", err)
t.Errorf("worker re-queried fresh track: before=%d after=%d", beforeCount, got) }
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) other := seedTrack(t, f, "Other", &otherMbid)
if _, err := f.pool.Exec(context.Background(), if _, err := f.pool.Exec(context.Background(),
`INSERT INTO track_similarity (track_a_id, track_b_id, score, source, fetched_at) `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) 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) srv := stubLB(`[{"recording_mbid":"`+otherMbid+`","score":0.95}]`, `[]`, http.StatusOK)
defer srv.Close() defer srv.Close()
w := newTestWorker(f, srv.URL) w := newTestWorker(f, srv.URL)
@@ -467,3 +466,147 @@ func TestUpsertArtistSimilar_PersistsUnmatchedToTable(t *testing.T) {
_ = inLibArtist _ = 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)
}
}