Persist metadata to DB, poll download worker, metadata service layer

This commit is contained in:
Alexander
2026-05-08 11:00:04 +02:00
parent 66264e1314
commit 60c94935b2
10 changed files with 489 additions and 58 deletions
+111 -33
View File
@@ -11,17 +11,18 @@ import (
"github.com/jackc/pgx/v5"
"github.com/riverqueue/river"
"github.com/rs/zerolog/log"
"google.golang.org/grpc"
metadataPb "homelab.lan/music-agregator/gen/metadata/v1"
pb "homelab.lan/music-agregator/gen/music_agregator/v1"
"homelab.lan/music-agregator/internal/config"
"homelab.lan/music-agregator/internal/database"
"homelab.lan/music-agregator/internal/indexer"
"homelab.lan/music-agregator/internal/metadata"
"homelab.lan/music-agregator/internal/release"
"homelab.lan/music-agregator/internal/torrent"
torrentParser "homelab.lan/music-agregator/internal/tracker"
"homelab.lan/music-agregator/internal/workers"
)
type parsedItem struct {
@@ -32,21 +33,23 @@ type parsedItem struct {
type MusicAgregatorService struct {
config config.Config
metadataClient metadataPb.MetadataServiceClient
metadataConn *grpc.ClientConn
metadata *metadata.MetadataService
indexer *indexer.IndexerService
torrentClient torrent.TorrentClient
magnetResolver *torrentParser.MagnetResolver
riverClient *river.Client[pgx.Tx]
torrents *database.TorrentRepository
downloads *database.DownloadRepository
}
func NewMusicAgregatorService(cfg config.Config, riverClient *river.Client[pgx.Tx]) (*MusicAgregatorService, error) {
indexer, err := indexer.NewIndexerService(cfg, riverClient)
func NewMusicAgregatorService(cfg config.Config, riverClient *river.Client[pgx.Tx], db *database.DB) (*MusicAgregatorService, error) {
idx, err := indexer.NewIndexerService(cfg, riverClient, nil)
if err != nil {
log.Err(err).Msg("failed to create IndexerService")
return nil, err
}
metadataClient, conn, err := metadata.NewMetadataClient(cfg.Metadata.Endpoint)
metadataClient, _, err := metadata.NewMetadataClient(cfg.Metadata.Endpoint)
if err != nil {
log.Err(err).Msg("failed to create metadata client")
return nil, err
@@ -66,29 +69,39 @@ func NewMusicAgregatorService(cfg config.Config, riverClient *river.Client[pgx.T
return &MusicAgregatorService{
config: cfg,
metadataClient: metadataClient,
metadataConn: conn,
indexer: indexer,
metadata: metadata.NewMetadataService(metadataClient, db),
indexer: idx,
torrentClient: torrentClient,
magnetResolver: magnetResolver,
riverClient: riverClient,
torrents: database.NewTorrentRepository(db.Pool),
downloads: database.NewDownloadRepository(db.Pool),
}, nil
}
func (s *MusicAgregatorService) Close() {
if s.metadataConn != nil {
s.metadataConn.Close()
}
if s.magnetResolver != nil {
s.magnetResolver.Close()
}
}
func (service *MusicAgregatorService) MonitorAlbum(ctx context.Context, req *pb.MonitorAlbumRequest) (*pb.MonitorAlbumResponse, error) {
album, err := service.fetchAlbumMetadata(ctx, req.GetAlbumId())
album, err := service.metadata.GetAlbum(ctx, req.GetAlbumId())
if err != nil {
log.Error().Err(err).Str("album_id", req.GetAlbumId()).Msg("failed to get album")
return nil, err
}
dbAlbum, _ := service.metadata.GetAlbumByExternalID(ctx, req.GetAlbumId())
if dbAlbum != nil {
qualityStr := normalizeQuality(req.GetQuality(), 0, 0)
owned, err := service.downloads.HasAlbumInQuality(ctx, dbAlbum.ID, req.GetQuality().String(), qualityStr)
if err == nil && owned {
log.Info().Str("album", dbAlbum.Title).Str("quality", qualityStr).Msg("album already owned in requested quality")
return &pb.MonitorAlbumResponse{}, nil
}
}
searchResult, err := service.searchIndexer(album, req.GetIndexerOptions().GetTracker())
if err != nil {
return nil, err
@@ -108,31 +121,17 @@ func (service *MusicAgregatorService) MonitorAlbum(ctx context.Context, req *pb.
return nil, err
}
if dbAlbum != nil {
service.saveTorrentAndDownload(ctx, dbAlbum.ID, best)
} else {
log.Warn().Str("album_id", req.GetAlbumId()).Msg("album not in DB, skipping torrent/download persistence")
}
return &pb.MonitorAlbumResponse{
Release: buildMonitoredRelease(best),
}, nil
}
func (service *MusicAgregatorService) fetchAlbumMetadata(ctx context.Context, albumID string) (*metadataPb.Album, error) {
resp, err := service.metadataClient.GetAlbum(ctx, &metadataPb.GetAlbumRequest{
Identifier: &metadataPb.GetAlbumRequest_Id{Id: albumID},
})
if err != nil {
log.Error().Err(err).Str("album_id", albumID).Msg("metadata GetAlbum failed")
return nil, err
}
album := resp.GetAlbum()
artistName := ""
if len(album.GetArtists()) > 0 {
artistName = album.GetArtists()[0].GetArtist().GetName()
}
log.Debug().Str("album_id", albumID).Str("title", album.GetTitle()).Str("artist", artistName).Msg("album metadata fetched")
return album, nil
}
func (service *MusicAgregatorService) searchIndexer(album *metadataPb.Album, tracker string) (*indexer.SearchResponse, error) {
artistName := ""
if len(album.GetArtists()) > 0 {
@@ -277,6 +276,85 @@ func (service *MusicAgregatorService) addToTorrentClient(best parsedItem) error
return nil
}
func (service *MusicAgregatorService) saveTorrentAndDownload(ctx context.Context, dbAlbumID string, best parsedItem) {
quality := normalizeQuality(pb.QualityType_QUALITY_UNSPECIFIED, best.rel.BitDepth, best.rel.SampleRate)
dbTorrent := &database.Torrent{
AlbumID: dbAlbumID,
InfoHash: best.rel.InfoHash,
Tracker: best.item.Tracker,
Title: best.item.Title,
Format: best.rel.Format.String(),
Quality: quality,
Source: best.rel.Source.String(),
BitDepth: best.rel.BitDepth,
SampleRate: best.rel.SampleRate,
Seeders: best.item.Seeders,
Peers: best.item.Peers,
Size: best.rel.TotalAudioSize,
TrackCount: best.rel.TrackCount,
HasCoverArt: best.rel.HasCoverArt,
HasCueSheet: best.rel.HasCueSheet,
HasRipLog: best.rel.HasRipLog,
DownloadLink: best.item.DownloadLink,
TorrentFile: best.torrentData,
}
if err := service.torrents.Create(ctx, dbTorrent); err != nil {
log.Error().Err(err).Str("hash", best.rel.InfoHash).Msg("failed to save torrent to DB")
return
}
savedTorrent, err := service.torrents.GetByInfoHash(ctx, best.rel.InfoHash)
if err != nil {
log.Error().Err(err).Msg("failed to retrieve saved torrent")
return
}
download := &database.Download{
TorrentID: savedTorrent.ID,
AlbumID: dbAlbumID,
Format: best.rel.Format.String(),
Quality: quality,
State: "downloading",
QbitHash: best.rel.InfoHash,
}
if err := service.downloads.Create(ctx, download); err != nil {
log.Error().Err(err).Msg("failed to save download to DB")
return
}
if service.riverClient != nil {
_, err := service.riverClient.Insert(ctx, workers.PollDownloadArgs{
DownloadID: download.ID,
TorrentHash: best.rel.InfoHash,
CheckInterval: 30 * time.Second,
}, &river.InsertOpts{
ScheduledAt: time.Now().Add(30 * time.Second),
})
if err != nil {
log.Error().Err(err).Msg("failed to schedule download poll job")
} else {
log.Debug().Str("download_id", download.ID).Str("hash", best.rel.InfoHash).Msg("download poll job scheduled")
}
}
log.Info().Str("hash", best.rel.InfoHash).Str("download_id", download.ID).Msg("torrent and download saved to DB")
}
func normalizeQuality(quality pb.QualityType, bitDepth int, sampleRate int) string {
if bitDepth > 0 && sampleRate > 0 {
return fmt.Sprintf("%d-%d", bitDepth, sampleRate/1000)
}
switch quality {
case pb.QualityType_QUALITY_LOSSLESS:
return "16-44"
case pb.QualityType_QUALITY_LOSSY:
return "320"
default:
return ""
}
}
func buildMonitoredRelease(p parsedItem) *pb.MonitoredRelease {
return &pb.MonitoredRelease{
InfoHash: p.rel.InfoHash,