Leftovers
This commit is contained in:
@@ -201,7 +201,7 @@ func serveGrpc(config config.Config) {
|
||||
|
||||
rs := setupRiver(ctx, config, db, torrentClient, pathMapper, musicfsClient, metadataService, bus)
|
||||
|
||||
musiscAgregatorSeerver, err := internal.NewMusicAgregatorServer(config, rs.client, torrentClient, pathMapper, db, bus)
|
||||
musiscAgregatorSeerver, err := internal.NewMusicAgregatorServer(config, rs.client, torrentClient, pathMapper, db, bus, musicfsClient)
|
||||
if err != nil {
|
||||
log.Fatal().Err(err).Msg("failed to create MusicAgregatorServer")
|
||||
}
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -27,6 +27,7 @@ const (
|
||||
MusicAgregatorService_SearchArtists_FullMethodName = "/music_agregator.v1.MusicAgregatorService/SearchArtists"
|
||||
MusicAgregatorService_GetArtistAlbums_FullMethodName = "/music_agregator.v1.MusicAgregatorService/GetArtistAlbums"
|
||||
MusicAgregatorService_SubscribeEvents_FullMethodName = "/music_agregator.v1.MusicAgregatorService/SubscribeEvents"
|
||||
MusicAgregatorService_ResyncAlbum_FullMethodName = "/music_agregator.v1.MusicAgregatorService/ResyncAlbum"
|
||||
)
|
||||
|
||||
// MusicAgregatorServiceClient is the client API for MusicAgregatorService service.
|
||||
@@ -42,6 +43,7 @@ type MusicAgregatorServiceClient interface {
|
||||
SearchArtists(ctx context.Context, in *SearchArtistsRequest, opts ...grpc.CallOption) (*SearchArtistsResponse, error)
|
||||
GetArtistAlbums(ctx context.Context, in *GetArtistAlbumsRequest, opts ...grpc.CallOption) (*GetArtistAlbumsResponse, error)
|
||||
SubscribeEvents(ctx context.Context, in *SubscribeEventsRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[AlbumEvent], error)
|
||||
ResyncAlbum(ctx context.Context, in *ResyncAlbumRequest, opts ...grpc.CallOption) (*ResyncAlbumResponse, error)
|
||||
}
|
||||
|
||||
type musicAgregatorServiceClient struct {
|
||||
@@ -145,6 +147,16 @@ func (c *musicAgregatorServiceClient) SubscribeEvents(ctx context.Context, in *S
|
||||
// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name.
|
||||
type MusicAgregatorService_SubscribeEventsClient = grpc.ServerStreamingClient[AlbumEvent]
|
||||
|
||||
func (c *musicAgregatorServiceClient) ResyncAlbum(ctx context.Context, in *ResyncAlbumRequest, opts ...grpc.CallOption) (*ResyncAlbumResponse, error) {
|
||||
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
|
||||
out := new(ResyncAlbumResponse)
|
||||
err := c.cc.Invoke(ctx, MusicAgregatorService_ResyncAlbum_FullMethodName, in, out, cOpts...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// MusicAgregatorServiceServer is the server API for MusicAgregatorService service.
|
||||
// All implementations must embed UnimplementedMusicAgregatorServiceServer
|
||||
// for forward compatibility.
|
||||
@@ -158,6 +170,7 @@ type MusicAgregatorServiceServer interface {
|
||||
SearchArtists(context.Context, *SearchArtistsRequest) (*SearchArtistsResponse, error)
|
||||
GetArtistAlbums(context.Context, *GetArtistAlbumsRequest) (*GetArtistAlbumsResponse, error)
|
||||
SubscribeEvents(*SubscribeEventsRequest, grpc.ServerStreamingServer[AlbumEvent]) error
|
||||
ResyncAlbum(context.Context, *ResyncAlbumRequest) (*ResyncAlbumResponse, error)
|
||||
mustEmbedUnimplementedMusicAgregatorServiceServer()
|
||||
}
|
||||
|
||||
@@ -192,6 +205,9 @@ func (UnimplementedMusicAgregatorServiceServer) GetArtistAlbums(context.Context,
|
||||
func (UnimplementedMusicAgregatorServiceServer) SubscribeEvents(*SubscribeEventsRequest, grpc.ServerStreamingServer[AlbumEvent]) error {
|
||||
return status.Error(codes.Unimplemented, "method SubscribeEvents not implemented")
|
||||
}
|
||||
func (UnimplementedMusicAgregatorServiceServer) ResyncAlbum(context.Context, *ResyncAlbumRequest) (*ResyncAlbumResponse, error) {
|
||||
return nil, status.Error(codes.Unimplemented, "method ResyncAlbum not implemented")
|
||||
}
|
||||
func (UnimplementedMusicAgregatorServiceServer) mustEmbedUnimplementedMusicAgregatorServiceServer() {}
|
||||
func (UnimplementedMusicAgregatorServiceServer) testEmbeddedByValue() {}
|
||||
|
||||
@@ -339,6 +355,24 @@ func _MusicAgregatorService_SubscribeEvents_Handler(srv interface{}, stream grpc
|
||||
// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name.
|
||||
type MusicAgregatorService_SubscribeEventsServer = grpc.ServerStreamingServer[AlbumEvent]
|
||||
|
||||
func _MusicAgregatorService_ResyncAlbum_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
|
||||
in := new(ResyncAlbumRequest)
|
||||
if err := dec(in); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if interceptor == nil {
|
||||
return srv.(MusicAgregatorServiceServer).ResyncAlbum(ctx, in)
|
||||
}
|
||||
info := &grpc.UnaryServerInfo{
|
||||
Server: srv,
|
||||
FullMethod: MusicAgregatorService_ResyncAlbum_FullMethodName,
|
||||
}
|
||||
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
|
||||
return srv.(MusicAgregatorServiceServer).ResyncAlbum(ctx, req.(*ResyncAlbumRequest))
|
||||
}
|
||||
return interceptor(ctx, in, info, handler)
|
||||
}
|
||||
|
||||
// MusicAgregatorService_ServiceDesc is the grpc.ServiceDesc for MusicAgregatorService service.
|
||||
// It's only intended for direct use with grpc.RegisterService,
|
||||
// and not to be introspected or modified (even as a copy)
|
||||
@@ -370,6 +404,10 @@ var MusicAgregatorService_ServiceDesc = grpc.ServiceDesc{
|
||||
MethodName: "GetArtistAlbums",
|
||||
Handler: _MusicAgregatorService_GetArtistAlbums_Handler,
|
||||
},
|
||||
{
|
||||
MethodName: "ResyncAlbum",
|
||||
Handler: _MusicAgregatorService_ResyncAlbum_Handler,
|
||||
},
|
||||
},
|
||||
Streams: []grpc.StreamDesc{
|
||||
{
|
||||
|
||||
@@ -173,6 +173,19 @@ func (r *DownloadRepository) GetByID(ctx context.Context, id string) (*Download,
|
||||
return d, nil
|
||||
}
|
||||
|
||||
func (r *DownloadRepository) GetLatestCompleted(ctx context.Context, albumID string) (*Download, error) {
|
||||
d := &Download{}
|
||||
err := r.pool.QueryRow(ctx,
|
||||
`SELECT id, torrent_id, album_id, format, quality, state, qbit_hash, save_path, error_message, queued_at, started_at, completed_at, created_at, updated_at
|
||||
FROM downloads WHERE album_id = $1 AND state = 'completed'
|
||||
ORDER BY completed_at DESC LIMIT 1`, albumID,
|
||||
).Scan(&d.ID, &d.TorrentID, &d.AlbumID, &d.Format, &d.Quality, &d.State, &d.QbitHash, &d.SavePath, &d.ErrorMessage, &d.QueuedAt, &d.StartedAt, &d.CompletedAt, &d.CreatedAt, &d.UpdatedAt)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("getting latest completed download: %w", err)
|
||||
}
|
||||
return d, nil
|
||||
}
|
||||
|
||||
func (r *DownloadRepository) GetByAlbumIDs(ctx context.Context, albumIDs []string) (map[string][]*Download, error) {
|
||||
if len(albumIDs) == 0 {
|
||||
return nil, nil
|
||||
|
||||
@@ -222,7 +222,10 @@ func (w *monitorWorkflow) run(ctx context.Context) error {
|
||||
return err
|
||||
}
|
||||
|
||||
parsed := w.service.parseSearchResults(searchResult, album)
|
||||
parsed := w.service.parseSearchResults(ctx, searchResult, album, func(phase string, current, total int) {
|
||||
w.publisher.PublishStatus(ctx, pb.MonitorStep_MONITOR_STEP_PARSING_RESULTS,
|
||||
fmt.Sprintf("%s: %d/%d", phase, current, total), nil)
|
||||
})
|
||||
|
||||
if len(parsed) > 0 {
|
||||
summaries := make([]*pb.TorrentSummary, len(parsed))
|
||||
@@ -250,6 +253,13 @@ func (w *monitorWorkflow) run(ctx context.Context) error {
|
||||
return err
|
||||
}
|
||||
|
||||
parsed = filterByAlbumMatch(parsed, album.GetTitle(), artistName)
|
||||
if len(parsed) == 0 {
|
||||
err := fmt.Errorf("no torrents match album title %q", album.GetTitle())
|
||||
w.publisher.PublishError(ctx, pb.MonitorStep_MONITOR_STEP_PARSING_RESULTS, err, false)
|
||||
return err
|
||||
}
|
||||
|
||||
if w.mode == pb.InteractionMode_INTERACTION_MODE_MANUAL && len(parsed) > 1 {
|
||||
options := make([]*pb.SelectOption, len(parsed))
|
||||
defaultIDs := make([]string, len(parsed))
|
||||
@@ -366,7 +376,7 @@ func (w *monitorWorkflow) run(ctx context.Context) error {
|
||||
}
|
||||
best = filtered[selectedIdx]
|
||||
} else {
|
||||
best = selectBestRelease(filtered)
|
||||
best = selectBestRelease(filtered, album.GetTitle())
|
||||
}
|
||||
|
||||
w.publisher.PublishStatus(ctx, pb.MonitorStep_MONITOR_STEP_SELECTING_RELEASE,
|
||||
|
||||
+7
-2
@@ -16,6 +16,7 @@ import (
|
||||
"homelab.lan/music-agregator/internal/config"
|
||||
"homelab.lan/music-agregator/internal/database"
|
||||
"homelab.lan/music-agregator/internal/eventbus"
|
||||
"homelab.lan/music-agregator/internal/musicfs"
|
||||
"homelab.lan/music-agregator/internal/torrent"
|
||||
)
|
||||
|
||||
@@ -26,8 +27,8 @@ type MusicAgregatorServer struct {
|
||||
pb.UnimplementedMusicAgregatorServiceServer
|
||||
}
|
||||
|
||||
func NewMusicAgregatorServer(cfg config.Config, riverClient *river.Client[pgx.Tx], torrentClient torrent.TorrentClient, pathMapper *torrent.PathMapper, db *database.DB, bus *eventbus.EventBus) (*MusicAgregatorServer, error) {
|
||||
service, err := NewMusicAgregatorService(cfg, riverClient, torrentClient, pathMapper, db)
|
||||
func NewMusicAgregatorServer(cfg config.Config, riverClient *river.Client[pgx.Tx], torrentClient torrent.TorrentClient, pathMapper *torrent.PathMapper, db *database.DB, bus *eventbus.EventBus, musicfsClient *musicfs.Client) (*MusicAgregatorServer, error) {
|
||||
service, err := NewMusicAgregatorService(cfg, riverClient, torrentClient, pathMapper, db, musicfsClient)
|
||||
if err != nil {
|
||||
log.Err(err).Msg("failed to create MusicAgregatorService")
|
||||
return nil, err
|
||||
@@ -245,6 +246,10 @@ func (s *MusicAgregatorServer) GetArtistAlbums(ctx context.Context, req *pb.GetA
|
||||
return s.service.GetArtistAlbums(ctx, req)
|
||||
}
|
||||
|
||||
func (s *MusicAgregatorServer) ResyncAlbum(ctx context.Context, req *pb.ResyncAlbumRequest) (*pb.ResyncAlbumResponse, error) {
|
||||
return s.service.ResyncAlbum(ctx, req.GetAlbumId())
|
||||
}
|
||||
|
||||
func (s *MusicAgregatorServer) Register(server *grpc.Server) {
|
||||
pb.RegisterMusicAgregatorServiceServer(server, s)
|
||||
}
|
||||
|
||||
+338
-42
@@ -5,8 +5,11 @@ import (
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"regexp"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
@@ -23,7 +26,9 @@ import (
|
||||
"homelab.lan/music-agregator/internal/config"
|
||||
"homelab.lan/music-agregator/internal/database"
|
||||
"homelab.lan/music-agregator/internal/indexer"
|
||||
albumMatcher "homelab.lan/music-agregator/internal/matcher"
|
||||
"homelab.lan/music-agregator/internal/metadata"
|
||||
"homelab.lan/music-agregator/internal/musicfs"
|
||||
"homelab.lan/music-agregator/internal/release"
|
||||
"homelab.lan/music-agregator/internal/torrent"
|
||||
torrentParser "homelab.lan/music-agregator/internal/tracker"
|
||||
@@ -54,11 +59,12 @@ type MusicAgregatorService struct {
|
||||
analyzer *analysis.ReleaseAnalyzer
|
||||
workflowRuns *database.WorkflowRunRepository
|
||||
albumEvents *database.AlbumEventRepository
|
||||
musicfsClient *musicfs.Client
|
||||
shutdownCtx context.Context
|
||||
shutdownCancel context.CancelFunc
|
||||
}
|
||||
|
||||
func NewMusicAgregatorService(cfg config.Config, riverClient *river.Client[pgx.Tx], torrentClient torrent.TorrentClient, pathMapper *torrent.PathMapper, db *database.DB) (*MusicAgregatorService, error) {
|
||||
func NewMusicAgregatorService(cfg config.Config, riverClient *river.Client[pgx.Tx], torrentClient torrent.TorrentClient, pathMapper *torrent.PathMapper, db *database.DB, musicfsClient *musicfs.Client) (*MusicAgregatorService, error) {
|
||||
idx, err := indexer.NewIndexerService(cfg, riverClient, nil)
|
||||
if err != nil {
|
||||
log.Err(err).Msg("failed to create IndexerService")
|
||||
@@ -96,6 +102,7 @@ func NewMusicAgregatorService(cfg config.Config, riverClient *river.Client[pgx.T
|
||||
analyzer: analysis.NewReleaseAnalyzer(db),
|
||||
workflowRuns: database.NewWorkflowRunRepository(db.Pool),
|
||||
albumEvents: database.NewAlbumEventRepository(db.Pool),
|
||||
musicfsClient: musicfsClient,
|
||||
shutdownCtx: ctx,
|
||||
shutdownCancel: cancel,
|
||||
}, nil
|
||||
@@ -560,12 +567,21 @@ func (service *MusicAgregatorService) MonitorAlbum(ctx context.Context, req *pb.
|
||||
return nil, err
|
||||
}
|
||||
|
||||
parsed := service.parseSearchResults(searchResult, album)
|
||||
parsed := service.parseSearchResults(ctx, searchResult, album, nil)
|
||||
|
||||
if len(parsed) == 0 {
|
||||
return nil, status.Errorf(codes.NotFound, "no torrents found for %s", album.GetTitle())
|
||||
}
|
||||
|
||||
albumArtist := ""
|
||||
if len(album.GetArtists()) > 0 {
|
||||
albumArtist = album.GetArtists()[0].GetArtist().GetName()
|
||||
}
|
||||
parsed = filterByAlbumMatch(parsed, album.GetTitle(), albumArtist)
|
||||
if len(parsed) == 0 {
|
||||
return nil, status.Errorf(codes.NotFound, "no torrents match album title %q", album.GetTitle())
|
||||
}
|
||||
|
||||
filtered := filterByQuality(parsed, req.GetQuality())
|
||||
if len(filtered) == 0 {
|
||||
return nil, status.Errorf(codes.NotFound, "no releases match quality filter for %s", album.GetTitle())
|
||||
@@ -576,7 +592,7 @@ func (service *MusicAgregatorService) MonitorAlbum(ctx context.Context, req *pb.
|
||||
return nil, status.Errorf(codes.NotFound, "no releases with active seeders for %s", album.GetTitle())
|
||||
}
|
||||
|
||||
best := selectBestRelease(filtered)
|
||||
best := selectBestRelease(filtered, album.GetTitle())
|
||||
|
||||
if err := service.addToTorrentClient(best); err != nil {
|
||||
return nil, err
|
||||
@@ -630,64 +646,237 @@ func (service *MusicAgregatorService) searchIndexer(album *metadataPb.Album, tra
|
||||
}
|
||||
|
||||
func buildSearchQueries(artist, title string) []string {
|
||||
base := title
|
||||
if artist != "" {
|
||||
base = artist + " " + title
|
||||
}
|
||||
artist = normalizeSearchQuery(artist)
|
||||
title = normalizeSearchQuery(title)
|
||||
|
||||
queries := []string{base}
|
||||
clean := cleanAlbumTitle(title)
|
||||
|
||||
stripped := stripParenthetical(title)
|
||||
if stripped != title {
|
||||
q := stripped
|
||||
if artist != "" {
|
||||
q = artist + " " + stripped
|
||||
}
|
||||
var queries []string
|
||||
seen := make(map[string]bool)
|
||||
add := func(q string) {
|
||||
q = strings.TrimSpace(q)
|
||||
if q != "" && !seen[strings.ToLower(q)] {
|
||||
seen[strings.ToLower(q)] = true
|
||||
queries = append(queries, q)
|
||||
}
|
||||
}
|
||||
|
||||
withArtist := func(t string) string {
|
||||
if artist != "" {
|
||||
return artist + " " + t
|
||||
}
|
||||
return t
|
||||
}
|
||||
|
||||
add(withArtist(clean))
|
||||
|
||||
noPunct := stripPunctuation(clean)
|
||||
add(withArtist(noPunct))
|
||||
|
||||
if artist != "" {
|
||||
add(artist)
|
||||
}
|
||||
|
||||
return queries
|
||||
}
|
||||
|
||||
func stripParenthetical(s string) string {
|
||||
idx := strings.Index(s, "(")
|
||||
if idx > 0 {
|
||||
return strings.TrimSpace(s[:idx])
|
||||
var unicodeReplacer = strings.NewReplacer(
|
||||
"\u2019", "'", // right single quote
|
||||
"\u2018", "'", // left single quote
|
||||
"\u201C", "\"", // left double quote
|
||||
"\u201D", "\"", // right double quote
|
||||
"\u2013", "-", // en dash
|
||||
"\u2014", "-", // em dash
|
||||
"\u2026", "...", // ellipsis
|
||||
)
|
||||
|
||||
func normalizeSearchQuery(s string) string {
|
||||
s = unicodeReplacer.Replace(s)
|
||||
for strings.Contains(s, " ") {
|
||||
s = strings.ReplaceAll(s, " ", " ")
|
||||
}
|
||||
return s
|
||||
return strings.TrimSpace(s)
|
||||
}
|
||||
|
||||
func (service *MusicAgregatorService) parseSearchResults(searchResult *indexer.SearchResponse, album *metadataPb.Album) []parsedItem {
|
||||
parser := torrentParser.NewGenericParser()
|
||||
var parsed []parsedItem
|
||||
var editionTagPattern = regexp.MustCompile(`(?i)\b(remaster(?:ed)?|deluxe(?:\s+edition)?|expanded(?:\s+edition)?|anniversary(?:\s+edition)?|special\s+edition|bonus\s+track|limited\s+edition|collector'?s?\s+edition|super\s+deluxe|standard\s+edition)\b`)
|
||||
var bracketPattern = regexp.MustCompile(`[\(\[\{][^)\]\}]*[\)\]\}]`)
|
||||
|
||||
for _, item := range searchResult.Items {
|
||||
func cleanAlbumTitle(title string) string {
|
||||
clean := bracketPattern.ReplaceAllString(title, "")
|
||||
clean = editionTagPattern.ReplaceAllString(clean, "")
|
||||
clean = strings.TrimRight(clean, " -–—:")
|
||||
for strings.Contains(clean, " ") {
|
||||
clean = strings.ReplaceAll(clean, " ", " ")
|
||||
}
|
||||
clean = strings.TrimSpace(clean)
|
||||
if clean == "" {
|
||||
return title
|
||||
}
|
||||
return clean
|
||||
}
|
||||
|
||||
func extractEditionTags(title string) []string {
|
||||
matches := editionTagPattern.FindAllString(title, -1)
|
||||
tags := make([]string, len(matches))
|
||||
for i, m := range matches {
|
||||
tags[i] = strings.ToLower(strings.TrimSpace(m))
|
||||
}
|
||||
return tags
|
||||
}
|
||||
|
||||
func stripPunctuation(s string) string {
|
||||
var b strings.Builder
|
||||
for _, r := range s {
|
||||
switch {
|
||||
case r >= 'a' && r <= 'z', r >= 'A' && r <= 'Z', r >= '0' && r <= '9', r == ' ':
|
||||
b.WriteRune(r)
|
||||
default:
|
||||
b.WriteRune(' ')
|
||||
}
|
||||
}
|
||||
result := b.String()
|
||||
for strings.Contains(result, " ") {
|
||||
result = strings.ReplaceAll(result, " ", " ")
|
||||
}
|
||||
return strings.TrimSpace(result)
|
||||
}
|
||||
|
||||
type resolveProgressFunc func(phase string, current, total int)
|
||||
|
||||
const maxMagnetResolves = 5
|
||||
|
||||
func (service *MusicAgregatorService) parseSearchResults(ctx context.Context, searchResult *indexer.SearchResponse, album *metadataPb.Album, onProgress resolveProgressFunc) []parsedItem {
|
||||
var candidates []*indexer.SearchItemResult
|
||||
for i := range searchResult.Items {
|
||||
item := searchResult.Items[i]
|
||||
if item.DownloadLink == "" {
|
||||
log.Trace().Str("title", item.Title).Msg("skipping item without download link")
|
||||
continue
|
||||
}
|
||||
|
||||
if item.Seeders == 0 {
|
||||
log.Warn().Str("title", item.Title).Str("tracker", item.Tracker).Msg("skipping torrent with no seeders")
|
||||
continue
|
||||
}
|
||||
|
||||
out := service.resolveRelease(parser, item, album)
|
||||
|
||||
log.Debug().
|
||||
Str("title", item.Title).
|
||||
Str("format", out.rel.Format.String()).
|
||||
Int("tracks", out.rel.TrackCount).
|
||||
Bool("lossless", out.rel.Format.IsLossless()).
|
||||
Int("seeders", item.Seeders).
|
||||
Int("real_seeders", out.realSeeders).
|
||||
Str("tracker", item.Tracker).
|
||||
Msg("release parsed")
|
||||
|
||||
parsed = append(parsed, parsedItem{item: item, rel: out.rel, torrentData: out.torrentData, realSeeders: out.realSeeders})
|
||||
candidates = append(candidates, item)
|
||||
}
|
||||
|
||||
log.Debug().Int("total", len(searchResult.Items)).Int("parsed", len(parsed)).Msg("parsing complete")
|
||||
total := len(candidates)
|
||||
if total == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
progress := func(phase string, current, t int) {
|
||||
if onProgress != nil {
|
||||
onProgress(phase, current, t)
|
||||
}
|
||||
}
|
||||
|
||||
progress("Pre-parsing titles", 0, total)
|
||||
parser := torrentParser.NewGenericParser()
|
||||
parsed := make([]parsedItem, total)
|
||||
for i, item := range candidates {
|
||||
rel := parser.Parse(item.Title)
|
||||
parsed[i] = parsedItem{
|
||||
item: item,
|
||||
rel: rel,
|
||||
realSeeders: -1,
|
||||
}
|
||||
}
|
||||
log.Info().Int("candidates", total).Msg("phase 1: title pre-parse complete")
|
||||
progress("Pre-parsed titles", total, total)
|
||||
|
||||
progress("Checking cache", 0, total)
|
||||
var cacheHits int
|
||||
for i := range parsed {
|
||||
p := &parsed[i]
|
||||
infoHash := torrentParser.ExtractInfoHash(p.item.DownloadLink)
|
||||
if infoHash == "" {
|
||||
if h, ok := p.item.Attrs["infohash"]; ok {
|
||||
infoHash = strings.ToLower(h)
|
||||
}
|
||||
}
|
||||
if infoHash == "" {
|
||||
continue
|
||||
}
|
||||
|
||||
cached, err := service.torrents.GetByInfoHash(ctx, infoHash)
|
||||
if err != nil || len(cached.TorrentFile) == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
fullParser := torrentParser.NewGenericParser()
|
||||
p.rel = fullParser.ParseTorrent(cached.TorrentFile, album)
|
||||
p.torrentData = cached.TorrentFile
|
||||
cacheHits++
|
||||
|
||||
log.Debug().
|
||||
Str("title", p.item.Title).
|
||||
Str("hash", infoHash).
|
||||
Str("format", p.rel.Format.String()).
|
||||
Msg("cache hit")
|
||||
}
|
||||
log.Info().Int("cache_hits", cacheHits).Int("candidates", total).Msg("phase 2: cache lookup complete")
|
||||
progress("Cache hits", cacheHits, total)
|
||||
|
||||
var needsResolve []int
|
||||
for i := range parsed {
|
||||
if parsed[i].torrentData == nil {
|
||||
needsResolve = append(needsResolve, i)
|
||||
}
|
||||
}
|
||||
|
||||
sort.Slice(needsResolve, func(a, b int) bool {
|
||||
return parsed[needsResolve[a]].item.Seeders > parsed[needsResolve[b]].item.Seeders
|
||||
})
|
||||
if len(needsResolve) > maxMagnetResolves {
|
||||
needsResolve = needsResolve[:maxMagnetResolves]
|
||||
}
|
||||
|
||||
resolveCount := len(needsResolve)
|
||||
if resolveCount > 0 {
|
||||
log.Info().Int("to_resolve", resolveCount).Int("skipped", total-cacheHits-resolveCount).Msg("phase 3: resolving top candidates")
|
||||
progress("Resolving torrents", 0, resolveCount)
|
||||
|
||||
const maxWorkers = 10
|
||||
sem := make(chan struct{}, maxWorkers)
|
||||
var wg sync.WaitGroup
|
||||
var resolved atomic.Int32
|
||||
|
||||
for _, idx := range needsResolve {
|
||||
wg.Add(1)
|
||||
go func(i int) {
|
||||
defer wg.Done()
|
||||
sem <- struct{}{}
|
||||
defer func() { <-sem }()
|
||||
|
||||
p := &parsed[i]
|
||||
resolveParser := torrentParser.NewGenericParser()
|
||||
out := service.resolveRelease(resolveParser, p.item, album)
|
||||
p.rel = out.rel
|
||||
p.torrentData = out.torrentData
|
||||
p.realSeeders = out.realSeeders
|
||||
|
||||
n := int(resolved.Add(1))
|
||||
progress("Resolving torrents", n, resolveCount)
|
||||
|
||||
log.Debug().
|
||||
Str("title", p.item.Title).
|
||||
Str("format", out.rel.Format.String()).
|
||||
Int("real_seeders", out.realSeeders).
|
||||
Msg("resolved")
|
||||
}(idx)
|
||||
}
|
||||
|
||||
wg.Wait()
|
||||
}
|
||||
|
||||
log.Debug().
|
||||
Int("total", total).
|
||||
Int("cache_hits", cacheHits).
|
||||
Int("resolved", resolveCount).
|
||||
Int("title_only", total-cacheHits-resolveCount).
|
||||
Msg("parsing complete")
|
||||
|
||||
return parsed
|
||||
}
|
||||
|
||||
@@ -729,6 +918,26 @@ func (service *MusicAgregatorService) resolveRelease(parser *torrentParser.Gener
|
||||
}
|
||||
}
|
||||
|
||||
func filterByAlbumMatch(items []parsedItem, albumTitle, artistName string) []parsedItem {
|
||||
var matched []parsedItem
|
||||
for _, p := range items {
|
||||
score := albumMatcher.AlbumMatchScore(albumTitle, p.item.Title, artistName)
|
||||
if score >= albumMatcher.MatchThreshold {
|
||||
matched = append(matched, p)
|
||||
} else {
|
||||
log.Debug().
|
||||
Str("title", p.item.Title).
|
||||
Float64("score", score).
|
||||
Str("wanted", albumTitle).
|
||||
Msg("filtered out by album match")
|
||||
}
|
||||
}
|
||||
if len(matched) == 0 {
|
||||
log.Warn().Str("album", albumTitle).Int("total", len(items)).Msg("no torrents match album title")
|
||||
}
|
||||
return matched
|
||||
}
|
||||
|
||||
func filterByQuality(items []parsedItem, quality pb.QualityType) []parsedItem {
|
||||
var filtered []parsedItem
|
||||
for _, p := range items {
|
||||
@@ -767,10 +976,12 @@ func filterByActiveSeeders(items []parsedItem) []parsedItem {
|
||||
return withSeeders
|
||||
}
|
||||
|
||||
func selectBestRelease(items []parsedItem) parsedItem {
|
||||
func selectBestRelease(items []parsedItem, albumTitle string) parsedItem {
|
||||
editionTags := extractEditionTags(albumTitle)
|
||||
|
||||
best := items[0]
|
||||
for _, p := range items[1:] {
|
||||
if betterRelease(p, best) {
|
||||
if betterRelease(p, best, editionTags) {
|
||||
best = p
|
||||
}
|
||||
}
|
||||
@@ -782,12 +993,19 @@ func selectBestRelease(items []parsedItem) parsedItem {
|
||||
Int("real_seeders", best.realSeeders).
|
||||
Str("tracker", best.item.Tracker).
|
||||
Str("hash", best.rel.InfoHash).
|
||||
Strs("edition_tags", editionTags).
|
||||
Msg("best release selected")
|
||||
|
||||
return best
|
||||
}
|
||||
|
||||
func betterRelease(candidate, current parsedItem) bool {
|
||||
func betterRelease(candidate, current parsedItem, editionTags []string) bool {
|
||||
cEdition := editionMatchScore(candidate.item.Title, editionTags)
|
||||
bEdition := editionMatchScore(current.item.Title, editionTags)
|
||||
if cEdition != bEdition {
|
||||
return cEdition > bEdition
|
||||
}
|
||||
|
||||
cHasReal := candidate.realSeeders > 0
|
||||
bHasReal := current.realSeeders > 0
|
||||
|
||||
@@ -803,6 +1021,20 @@ func betterRelease(candidate, current parsedItem) bool {
|
||||
return candidate.item.Seeders > current.item.Seeders
|
||||
}
|
||||
|
||||
func editionMatchScore(torrentTitle string, editionTags []string) int {
|
||||
if len(editionTags) == 0 {
|
||||
return 0
|
||||
}
|
||||
lower := strings.ToLower(torrentTitle)
|
||||
score := 0
|
||||
for _, tag := range editionTags {
|
||||
if strings.Contains(lower, tag) {
|
||||
score++
|
||||
}
|
||||
}
|
||||
return score
|
||||
}
|
||||
|
||||
func (service *MusicAgregatorService) addToTorrentClient(best parsedItem) error {
|
||||
if best.rel.InfoHash != "" {
|
||||
existing, err := service.torrentClient.Find(torrent.FindOptions{Hash: best.rel.InfoHash})
|
||||
@@ -1054,6 +1286,70 @@ func downloadTorrentData(url string) ([]byte, error) {
|
||||
return data, nil
|
||||
}
|
||||
|
||||
func (service *MusicAgregatorService) ResyncAlbum(ctx context.Context, albumID string) (*pb.ResyncAlbumResponse, error) {
|
||||
if service.musicfsClient == nil {
|
||||
return nil, status.Error(codes.FailedPrecondition, "musicfs not configured")
|
||||
}
|
||||
|
||||
download, err := service.downloads.GetLatestCompleted(ctx, albumID)
|
||||
if err != nil {
|
||||
return nil, status.Errorf(codes.NotFound, "no completed download for album %s: %v", albumID, err)
|
||||
}
|
||||
|
||||
subdir := ""
|
||||
if download.QbitHash != "" {
|
||||
results, err := service.torrentClient.Find(torrent.FindOptions{Hash: download.QbitHash})
|
||||
if err == nil && len(results) > 0 {
|
||||
contentPath := results[0].ContentPath
|
||||
if service.pathMapper != nil {
|
||||
contentPath = service.pathMapper.ToHost(contentPath)
|
||||
}
|
||||
subdir = musicfs.DeriveSubdir(service.config.MusicFS.OriginRoot, contentPath)
|
||||
}
|
||||
}
|
||||
|
||||
log.Info().
|
||||
Str("album_id", albumID).
|
||||
Str("subdir", subdir).
|
||||
Msg("resyncing album with musicfs")
|
||||
|
||||
result, err := service.musicfsClient.TriggerRescan(ctx, service.config.MusicFS.OriginID, subdir)
|
||||
if err != nil {
|
||||
return nil, status.Errorf(codes.Internal, "musicfs rescan failed: %v", err)
|
||||
}
|
||||
|
||||
log.Info().
|
||||
Int("new_files", len(result.NewFiles)).
|
||||
Uint64("bytes_synced", result.BytesSynced).
|
||||
Msg("musicfs rescan complete")
|
||||
|
||||
var filesEnriched int
|
||||
if len(result.NewFiles) > 0 {
|
||||
dbAlbum, err := service.metadata.GetAlbumByID(ctx, albumID)
|
||||
if err != nil {
|
||||
log.Warn().Err(err).Str("album_id", albumID).Msg("failed to get album for enrichment")
|
||||
} else {
|
||||
albumMeta, err := service.metadata.GetAlbum(ctx, dbAlbum.ExternalID)
|
||||
if err != nil {
|
||||
log.Warn().Err(err).Str("external_id", dbAlbum.ExternalID).Msg("failed to get album metadata for enrichment")
|
||||
} else {
|
||||
err = service.musicfsClient.EnrichFiles(ctx, result.NewFiles, albumMeta)
|
||||
if err != nil {
|
||||
log.Warn().Err(err).Msg("musicfs enrichment failed")
|
||||
} else {
|
||||
filesEnriched = len(result.NewFiles)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return &pb.ResyncAlbumResponse{
|
||||
FilesSynced: int32(len(result.NewFiles)),
|
||||
FilesEnriched: int32(filesEnriched),
|
||||
BytesSynced: result.BytesSynced,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (service *MusicAgregatorService) SearchArtists(ctx context.Context, req *pb.SearchArtistsRequest) (*pb.SearchArtistsResponse, error) {
|
||||
resp, err := service.metadata.SearchArtists(ctx, req.GetQuery(), req.GetLimit(), req.GetOffset())
|
||||
if err != nil {
|
||||
|
||||
@@ -3,6 +3,8 @@ package tracker
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/anacrolix/torrent"
|
||||
@@ -11,6 +13,18 @@ import (
|
||||
"github.com/rs/zerolog/log"
|
||||
)
|
||||
|
||||
func ExtractInfoHash(magnetURI string) string {
|
||||
u, err := url.Parse(magnetURI)
|
||||
if err != nil {
|
||||
return ""
|
||||
}
|
||||
xt := u.Query().Get("xt")
|
||||
if strings.HasPrefix(xt, "urn:btih:") {
|
||||
return strings.ToLower(xt[len("urn:btih:"):])
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
type ResolveResult struct {
|
||||
Data []byte
|
||||
ConnectedSeeders int
|
||||
|
||||
@@ -13,6 +13,17 @@ service MusicAgregatorService {
|
||||
rpc SearchArtists(SearchArtistsRequest) returns (SearchArtistsResponse) {}
|
||||
rpc GetArtistAlbums(GetArtistAlbumsRequest) returns (GetArtistAlbumsResponse) {}
|
||||
rpc SubscribeEvents(SubscribeEventsRequest) returns (stream AlbumEvent) {}
|
||||
rpc ResyncAlbum(ResyncAlbumRequest) returns (ResyncAlbumResponse) {}
|
||||
}
|
||||
|
||||
message ResyncAlbumRequest {
|
||||
string album_id = 1;
|
||||
}
|
||||
|
||||
message ResyncAlbumResponse {
|
||||
int32 files_synced = 1;
|
||||
int32 files_enriched = 2;
|
||||
uint64 bytes_synced = 3;
|
||||
}
|
||||
|
||||
message MonitorAlbumRequest {
|
||||
|
||||
Reference in New Issue
Block a user