Files
metadata-agregator/internal/service/metadata.go
T
Alexander 2756cc1043 feat: fuzzy search with parallel DB/MusicBrainz execution
- Add pg_trgm fuzzy search to PostgreSQL artist repository
- Add Lucene fuzzy query (~0.7) for MusicBrainz artist search
- Run DB and MusicBrainz searches in parallel goroutines
- Cache MusicBrainz results to DB for accumulation
- Deduplicate results by external ID
- Filter out special purpose artists (type=Other)
- Add SaveAll method for batch artist inserts
- Move database schema to containers repo
2026-05-10 21:13:17 +02:00

372 lines
10 KiB
Go

package service
import (
"context"
"errors"
"sync"
"github.com/rs/zerolog"
"github.com/metadata-agregator/internal/domain"
"github.com/metadata-agregator/internal/metrics"
"github.com/metadata-agregator/internal/provider"
"github.com/metadata-agregator/internal/repository"
)
type MetadataService struct {
artists repository.ArtistRepository
albums repository.AlbumRepository
tracks repository.TrackRepository
provider provider.Provider
}
func NewMetadataService(
artists repository.ArtistRepository,
albums repository.AlbumRepository,
tracks repository.TrackRepository,
prov provider.Provider,
) *MetadataService {
return &MetadataService{
artists: artists,
albums: albums,
tracks: tracks,
provider: prov,
}
}
func (s *MetadataService) GetArtist(ctx context.Context, id string) (*domain.Artist, error) {
log := zerolog.Ctx(ctx)
artist, err := s.artists.GetByExternalID(ctx, s.provider.Name(), id)
if err == nil {
log.Debug().Str("artist_id", id).Msg("artist cache hit")
metrics.CacheHits.WithLabelValues("artist").Inc()
return artist, nil
}
if !errors.Is(err, repository.ErrNotFound) {
log.Error().Err(err).Str("artist_id", id).Msg("artist cache lookup failed")
return nil, err
}
metrics.CacheMisses.WithLabelValues("artist").Inc()
log.Debug().Str("artist_id", id).Str("provider", s.provider.Name()).Msg("artist cache miss, fetching from provider")
artist, err = s.provider.GetArtist(ctx, id)
if err != nil {
log.Debug().Err(err).Str("artist_id", id).Msg("provider fetch failed")
return nil, err
}
if saveErr := s.artists.Save(ctx, artist); saveErr != nil {
log.Warn().Err(saveErr).Str("artist_id", id).Msg("failed to cache artist")
return artist, nil
}
log.Trace().Str("artist_id", id).Msg("artist cached successfully")
return artist, nil
}
func (s *MetadataService) SearchArtists(ctx context.Context, query string, limit, offset int) (*domain.SearchResult[domain.Artist], error) {
log := zerolog.Ctx(ctx)
var (
wg sync.WaitGroup
dbResult *domain.SearchResult[domain.Artist]
mbResult *domain.SearchResult[domain.Artist]
dbErr error
mbErr error
)
wg.Add(2)
go func() {
defer wg.Done()
dbResult, dbErr = s.artists.Search(ctx, query, limit, offset)
if dbErr == nil && len(dbResult.Items) > 0 {
metrics.CacheHits.WithLabelValues("artist_search").Inc()
}
}()
go func() {
defer wg.Done()
mbResult, mbErr = s.provider.SearchArtists(ctx, query, limit, offset)
if mbErr == nil && len(mbResult.Items) > 0 {
if saveErr := s.artists.SaveAll(ctx, mbResult.Items); saveErr != nil {
log.Warn().Err(saveErr).Msg("failed to cache provider search results")
} else {
log.Debug().Int("count", len(mbResult.Items)).Msg("cached provider search results")
}
}
}()
wg.Wait()
if dbErr != nil {
log.Warn().Err(dbErr).Msg("database search failed")
}
if mbErr != nil {
log.Warn().Err(mbErr).Msg("provider search failed")
}
if dbErr != nil && mbErr != nil {
return nil, mbErr
}
combined := s.deduplicateArtists(dbResult, mbResult)
log.Debug().
Int("db_count", countItems(dbResult)).
Int("provider_count", countItems(mbResult)).
Int("combined_count", len(combined.Items)).
Msg("artist search completed")
return combined, nil
}
func (s *MetadataService) deduplicateArtists(dbResult, mbResult *domain.SearchResult[domain.Artist]) *domain.SearchResult[domain.Artist] {
seen := make(map[string]bool)
var items []domain.Artist
total := 0
addArtist := func(artist domain.Artist) {
if artist.Type == "Other" {
return
}
for _, ext := range artist.ExternalIDs {
key := ext.Source + ":" + ext.SourceID
if seen[key] {
return
}
}
for _, ext := range artist.ExternalIDs {
seen[ext.Source+":"+ext.SourceID] = true
}
items = append(items, artist)
}
if dbResult != nil {
for _, artist := range dbResult.Items {
addArtist(artist)
}
total = dbResult.Total
}
if mbResult != nil {
for _, artist := range mbResult.Items {
addArtist(artist)
}
if mbResult.Total > total {
total = mbResult.Total
}
}
limit := 25
offset := 0
if dbResult != nil {
limit = dbResult.Limit
offset = dbResult.Offset
} else if mbResult != nil {
limit = mbResult.Limit
offset = mbResult.Offset
}
return &domain.SearchResult[domain.Artist]{
Items: items,
Total: total,
Limit: limit,
Offset: offset,
}
}
func countItems[T any](result *domain.SearchResult[T]) int {
if result == nil {
return 0
}
return len(result.Items)
}
func (s *MetadataService) SearchAlbums(ctx context.Context, query string, artist string, limit, offset int) (*domain.SearchResult[domain.Album], error) {
zerolog.Ctx(ctx).Debug().Str("query", query).Str("artist", artist).Str("provider", s.provider.Name()).Msg("searching albums via provider")
return s.provider.SearchAlbums(ctx, query, artist, limit, offset)
}
func (s *MetadataService) GetAlbum(ctx context.Context, id string) (*domain.Album, error) {
log := zerolog.Ctx(ctx)
album, err := s.albums.GetByExternalID(ctx, s.provider.Name(), id)
if err == nil {
log.Debug().Str("album_id", id).Msg("album cache hit")
metrics.CacheHits.WithLabelValues("album").Inc()
return album, nil
}
if !errors.Is(err, repository.ErrNotFound) {
log.Error().Err(err).Str("album_id", id).Msg("album cache lookup failed")
return nil, err
}
metrics.CacheMisses.WithLabelValues("album").Inc()
log.Debug().Str("album_id", id).Str("provider", s.provider.Name()).Msg("album cache miss, fetching from provider")
album, err = s.provider.GetAlbum(ctx, id)
if err != nil {
log.Debug().Err(err).Str("album_id", id).Msg("provider fetch failed")
return nil, err
}
s.ensureArtistsCached(ctx, []domain.Album{*album})
if saveErr := s.albums.Save(ctx, album); saveErr != nil {
log.Warn().Err(saveErr).Str("album_id", id).Msg("failed to cache album")
return album, nil
}
log.Trace().Str("album_id", id).Msg("album cached successfully")
return album, nil
}
func (s *MetadataService) GetArtistAlbums(ctx context.Context, artistID string) ([]domain.Album, error) {
log := zerolog.Ctx(ctx)
cached, err := s.albums.GetAllByArtistID(ctx, artistID)
if err == nil && len(cached) > 0 {
log.Debug().Str("artist_id", artistID).Int("count", len(cached)).Msg("artist albums cache hit")
metrics.CacheHits.WithLabelValues("artist_albums").Inc()
return cached, nil
}
metrics.CacheMisses.WithLabelValues("artist_albums").Inc()
log.Debug().Str("artist_id", artistID).Str("provider", s.provider.Name()).Msg("artist albums cache miss, fetching from provider")
albums, err := s.provider.GetArtistAlbums(ctx, artistID)
if err != nil {
return nil, err
}
s.ensureArtistsCached(ctx, albums)
if saveErr := s.albums.SaveAll(ctx, albums); saveErr != nil {
log.Warn().Err(saveErr).Str("artist_id", artistID).Msg("failed to cache artist albums")
} else {
log.Debug().Str("artist_id", artistID).Int("count", len(albums)).Msg("artist albums cached")
}
return albums, nil
}
func (s *MetadataService) ensureArtistsCached(ctx context.Context, albums []domain.Album) {
log := zerolog.Ctx(ctx)
seen := make(map[string]bool)
for _, album := range albums {
for _, ac := range album.Artists {
id := ac.Artist.ID
if id == "" || seen[id] {
continue
}
seen[id] = true
_, err := s.artists.GetByID(ctx, id)
if err == nil {
continue
}
if _, err := s.GetArtist(ctx, id); err != nil {
log.Warn().Err(err).Str("artist_id", id).Msg("failed to fetch and cache artist from album credit")
}
}
}
}
func (s *MetadataService) GetTrack(ctx context.Context, id string) (*domain.Track, error) {
log := zerolog.Ctx(ctx)
track, err := s.tracks.GetByExternalID(ctx, s.provider.Name(), id)
if err == nil {
log.Debug().Str("track_id", id).Msg("track cache hit")
metrics.CacheHits.WithLabelValues("track").Inc()
return track, nil
}
if !errors.Is(err, repository.ErrNotFound) {
log.Error().Err(err).Str("track_id", id).Msg("track cache lookup failed")
return nil, err
}
metrics.CacheMisses.WithLabelValues("track").Inc()
log.Debug().Str("track_id", id).Str("provider", s.provider.Name()).Msg("track cache miss, fetching from provider")
track, err = s.provider.GetTrack(ctx, id)
if err != nil {
log.Debug().Err(err).Str("track_id", id).Msg("provider fetch failed")
return nil, err
}
if saveErr := s.tracks.Save(ctx, track); saveErr != nil {
log.Warn().Err(saveErr).Str("track_id", id).Msg("failed to cache track")
return track, nil
}
log.Trace().Str("track_id", id).Msg("track cached successfully")
return track, nil
}
func (s *MetadataService) GetTrackByISRC(ctx context.Context, isrc string) (*domain.Track, error) {
log := zerolog.Ctx(ctx)
track, err := s.tracks.GetByISRC(ctx, isrc)
if err == nil {
log.Debug().Str("isrc", isrc).Msg("track ISRC cache hit")
metrics.CacheHits.WithLabelValues("track_isrc").Inc()
return track, nil
}
if !errors.Is(err, repository.ErrNotFound) {
log.Error().Err(err).Str("isrc", isrc).Msg("track ISRC cache lookup failed")
return nil, err
}
metrics.CacheMisses.WithLabelValues("track_isrc").Inc()
log.Debug().Str("isrc", isrc).Str("provider", s.provider.Name()).Msg("track ISRC cache miss, fetching from provider")
track, err = s.provider.GetTrackByISRC(ctx, isrc)
if err != nil {
log.Debug().Err(err).Str("isrc", isrc).Msg("provider ISRC fetch failed")
return nil, err
}
if saveErr := s.tracks.Save(ctx, track); saveErr != nil {
log.Warn().Err(saveErr).Str("isrc", isrc).Msg("failed to cache track")
return track, nil
}
log.Trace().Str("isrc", isrc).Msg("track cached successfully")
return track, nil
}
func (s *MetadataService) GetAlbumTracks(ctx context.Context, albumID string) ([]domain.Track, error) {
log := zerolog.Ctx(ctx)
tracks, err := s.tracks.GetByAlbumID(ctx, albumID)
if err == nil && len(tracks) > 0 {
log.Debug().Str("album_id", albumID).Int("count", len(tracks)).Msg("album tracks served from cache")
metrics.CacheHits.WithLabelValues("album_tracks").Inc()
return tracks, nil
}
metrics.CacheMisses.WithLabelValues("album_tracks").Inc()
log.Debug().Str("album_id", albumID).Str("provider", s.provider.Name()).Msg("album tracks cache miss, querying provider")
tracks, err = s.provider.GetAlbumTracks(ctx, albumID)
if err != nil {
return nil, err
}
album, err := s.albums.GetByExternalID(ctx, s.provider.Name(), albumID)
if err == nil {
if saveErr := s.tracks.SaveAlbumTracks(ctx, album.ID, tracks); saveErr != nil {
log.Warn().Err(saveErr).Str("album_id", albumID).Msg("failed to cache album tracks")
} else {
log.Debug().Str("album_id", albumID).Int("count", len(tracks)).Msg("album tracks cached")
}
}
return tracks, nil
}