feat: integrate musicfs for post-download metadata enrichment
- Add internal/musicfs package: gRPC client, TriggerRescan (stream consumer), EnrichFiles (BatchUpdateMetadata), DeriveSubdir helper - Add MusicFS config section (enabled, endpoint, origin_id, origin_root, timeout) - Wire MusicFSClient into PollDownloadWorker: on download completion, trigger scoped rescan then push enriched metadata - Fix album auto-persist in monitor workflow: persist artist+album from metadata-agregator before saving torrent/download - Add filterByActiveSeeders: skip torrents where magnet resolution found 0 real active seeders - Add musicfs proto to buf.gen.yaml inputs - Enricher only sends non-empty fields to avoid overwriting existing file tag data with empty values
This commit is contained in:
@@ -2,6 +2,7 @@ version: v2
|
||||
inputs:
|
||||
- directory: proto
|
||||
- directory: ../metadata-agregator/proto
|
||||
- directory: ../musicfs/crates/musicfs-grpc/proto
|
||||
plugins:
|
||||
- local: protoc-gen-go
|
||||
out: gen
|
||||
|
||||
@@ -32,6 +32,7 @@ import (
|
||||
"homelab.lan/music-agregator/internal/hello"
|
||||
"homelab.lan/music-agregator/internal/indexer"
|
||||
"homelab.lan/music-agregator/internal/metadata"
|
||||
"homelab.lan/music-agregator/internal/musicfs"
|
||||
"homelab.lan/music-agregator/internal/torrent"
|
||||
"homelab.lan/music-agregator/internal/workers"
|
||||
)
|
||||
@@ -83,7 +84,7 @@ type riverSetup struct {
|
||||
cacheRefreshWorker *indexer.CacheRefreshWorker
|
||||
}
|
||||
|
||||
func setupRiver(ctx context.Context, cfg config.Config, db *database.DB, torrentClient torrent.TorrentClient, pathMapper *torrent.PathMapper) *riverSetup {
|
||||
func setupRiver(ctx context.Context, cfg config.Config, db *database.DB, torrentClient torrent.TorrentClient, pathMapper *torrent.PathMapper, musicfsClient *musicfs.Client, metadataService *metadata.MetadataService) *riverSetup {
|
||||
cacheWorker := &indexer.CacheRefreshWorker{}
|
||||
pollWorker := &workers.PollDownloadWorker{
|
||||
Downloads: database.NewDownloadRepository(db.Pool),
|
||||
@@ -93,6 +94,10 @@ func setupRiver(ctx context.Context, cfg config.Config, db *database.DB, torrent
|
||||
TorrentClient: torrentClient,
|
||||
PathMapper: pathMapper,
|
||||
Analyzer: analysis.NewReleaseAnalyzer(db),
|
||||
MusicFSClient: musicfsClient,
|
||||
MusicFSOriginID: cfg.MusicFS.OriginID,
|
||||
MusicFSOriginRoot: cfg.MusicFS.OriginRoot,
|
||||
MetadataService: metadataService,
|
||||
}
|
||||
|
||||
riverWorkers := river.NewWorkers()
|
||||
@@ -171,7 +176,25 @@ func serveGrpc(config config.Config) {
|
||||
log.Fatal().Err(err).Msg("failed to create path mapper")
|
||||
}
|
||||
|
||||
rs := setupRiver(ctx, config, db, torrentClient, pathMapper)
|
||||
var musicfsClient *musicfs.Client
|
||||
if config.MusicFS.Enabled {
|
||||
client, conn, err := musicfs.NewClient(config.MusicFS.Endpoint)
|
||||
if err != nil {
|
||||
log.Warn().Err(err).Msg("failed to connect to musicfs, continuing without it")
|
||||
} else {
|
||||
musicfsClient = client
|
||||
defer conn.Close()
|
||||
}
|
||||
}
|
||||
|
||||
metadataClient, metadataConn, err := metadata.NewMetadataClient(config.Metadata.Endpoint)
|
||||
if err != nil {
|
||||
log.Fatal().Err(err).Msg("failed to connect to metadata service")
|
||||
}
|
||||
defer metadataConn.Close()
|
||||
metadataService := metadata.NewMetadataService(metadataClient, db)
|
||||
|
||||
rs := setupRiver(ctx, config, db, torrentClient, pathMapper, musicfsClient, metadataService)
|
||||
|
||||
musiscAgregatorSeerver, err := internal.NewMusicAgregatorServer(config, rs.client, torrentClient, pathMapper, db)
|
||||
if err != nil {
|
||||
|
||||
@@ -45,6 +45,14 @@ type Config struct {
|
||||
Metadata struct {
|
||||
Endpoint string `yaml:"endpoint"`
|
||||
} `yaml:"metadata"`
|
||||
|
||||
MusicFS struct {
|
||||
Enabled bool `yaml:"enabled"`
|
||||
Endpoint string `yaml:"endpoint"`
|
||||
OriginID string `yaml:"origin_id"`
|
||||
OriginRoot string `yaml:"origin_root"`
|
||||
TimeoutSeconds int `yaml:"timeout_seconds"`
|
||||
} `yaml:"musicfs"`
|
||||
}
|
||||
|
||||
type CacheConfig struct {
|
||||
|
||||
@@ -312,6 +312,14 @@ func (w *monitorWorkflow) run(ctx context.Context) error {
|
||||
return err
|
||||
}
|
||||
|
||||
filtered = filterByActiveSeeders(filtered)
|
||||
if len(filtered) == 0 {
|
||||
err := fmt.Errorf("no releases with active seeders for %s", album.GetTitle())
|
||||
log.Warn().Str("album", album.GetTitle()).Msg("no releases with active seeders")
|
||||
w.publisher.PublishError(ctx, pb.MonitorStep_MONITOR_STEP_FILTERING_QUALITY, err, false)
|
||||
return err
|
||||
}
|
||||
|
||||
var best parsedItem
|
||||
if w.mode == pb.InteractionMode_INTERACTION_MODE_MANUAL && len(filtered) > 1 {
|
||||
options := make([]*pb.SelectOption, len(filtered))
|
||||
@@ -402,11 +410,16 @@ func (w *monitorWorkflow) run(ctx context.Context) error {
|
||||
w.publisher.PublishStatus(ctx, pb.MonitorStep_MONITOR_STEP_SAVING, "Saving to database...", nil)
|
||||
|
||||
dbAlbum, _ = w.service.metadata.GetAlbumByExternalID(ctx, album.GetId())
|
||||
if dbAlbum == nil {
|
||||
w.service.metadata.PersistArtist(ctx, album, database.Monitored)
|
||||
w.service.metadata.PersistAlbum(ctx, album, database.Monitored)
|
||||
dbAlbum, _ = w.service.metadata.GetAlbumByExternalID(ctx, album.GetId())
|
||||
}
|
||||
if dbAlbum != nil {
|
||||
w.publisher.SetAlbumID(dbAlbum.ID)
|
||||
w.service.saveTorrentAndDownload(ctx, dbAlbum.ID, best)
|
||||
} else {
|
||||
log.Warn().Str("album_id", w.req.AlbumId).Msg("album not in DB, skipping torrent/download persistence")
|
||||
log.Warn().Str("album_id", w.req.AlbumId).Msg("album not in DB after persist attempt, skipping torrent/download persistence")
|
||||
}
|
||||
|
||||
w.addedHash = best.rel.InfoHash
|
||||
|
||||
@@ -0,0 +1,32 @@
|
||||
package musicfs
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"github.com/rs/zerolog/log"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials/insecure"
|
||||
|
||||
pb "homelab.lan/music-agregator/gen/musicfs/v1"
|
||||
)
|
||||
|
||||
type Client struct {
|
||||
MusicFS pb.MusicFSClient
|
||||
Metadata pb.MetadataServiceClient
|
||||
}
|
||||
|
||||
func NewClient(endpoint string) (*Client, *grpc.ClientConn, error) {
|
||||
log.Trace().Str("endpoint", endpoint).Msg("connecting to musicfs")
|
||||
|
||||
conn, err := grpc.NewClient(endpoint, grpc.WithTransportCredentials(insecure.NewCredentials()))
|
||||
if err != nil {
|
||||
return nil, nil, fmt.Errorf("connecting to musicfs: %w", err)
|
||||
}
|
||||
|
||||
log.Info().Str("endpoint", endpoint).Msg("musicfs connected")
|
||||
|
||||
return &Client{
|
||||
MusicFS: pb.NewMusicFSClient(conn),
|
||||
Metadata: pb.NewMetadataServiceClient(conn),
|
||||
}, conn, nil
|
||||
}
|
||||
@@ -0,0 +1,92 @@
|
||||
package musicfs
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
|
||||
"github.com/rs/zerolog/log"
|
||||
|
||||
metadataPb "homelab.lan/music-agregator/gen/metadata/v1"
|
||||
pb "homelab.lan/music-agregator/gen/musicfs/v1"
|
||||
)
|
||||
|
||||
func (c *Client) EnrichFiles(ctx context.Context, files []SyncedFile, album *metadataPb.Album) error {
|
||||
if len(files) == 0 || album == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
genre := firstGenreName(album.GetGenres())
|
||||
label := labelName(album.GetLabel())
|
||||
albumType := album.GetAlbumType()
|
||||
coverURL := album.GetCoverUrl()
|
||||
|
||||
var items []*pb.BatchUpdateItem
|
||||
for _, f := range files {
|
||||
meta := &pb.UpdateMetadataRequest{
|
||||
FileId: f.FileID,
|
||||
}
|
||||
if genre != "" {
|
||||
meta.Genre = &genre
|
||||
}
|
||||
if label != "" {
|
||||
meta.Label = &label
|
||||
}
|
||||
if albumType != "" {
|
||||
meta.AlbumType = &albumType
|
||||
}
|
||||
if coverURL != "" {
|
||||
meta.CoverUrl = &coverURL
|
||||
}
|
||||
items = append(items, &pb.BatchUpdateItem{
|
||||
FileId: f.FileID,
|
||||
Metadata: meta,
|
||||
})
|
||||
}
|
||||
|
||||
stream, err := c.Metadata.BatchUpdateMetadata(ctx, &pb.BatchUpdateRequest{Items: items})
|
||||
if err != nil {
|
||||
return fmt.Errorf("batch update: %w", err)
|
||||
}
|
||||
|
||||
var updated, failed int
|
||||
for {
|
||||
progress, err := stream.Recv()
|
||||
if err == io.EOF {
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
return fmt.Errorf("batch update stream: %w", err)
|
||||
}
|
||||
if progress.ErrorMessage != nil {
|
||||
log.Warn().
|
||||
Int64("file_id", progress.GetCurrentFileId()).
|
||||
Str("error", progress.GetErrorMessage()).
|
||||
Msg("batch update item failed")
|
||||
failed++
|
||||
} else {
|
||||
updated++
|
||||
}
|
||||
}
|
||||
|
||||
log.Info().
|
||||
Int("updated", updated).
|
||||
Int("failed", failed).
|
||||
Msg("enrichment batch complete")
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func firstGenreName(genres []*metadataPb.Genre) string {
|
||||
if len(genres) == 0 {
|
||||
return ""
|
||||
}
|
||||
return genres[0].GetName()
|
||||
}
|
||||
|
||||
func labelName(label *metadataPb.Label) string {
|
||||
if label == nil {
|
||||
return ""
|
||||
}
|
||||
return label.GetName()
|
||||
}
|
||||
@@ -0,0 +1,20 @@
|
||||
package musicfs
|
||||
|
||||
import (
|
||||
"path/filepath"
|
||||
"strings"
|
||||
|
||||
"github.com/rs/zerolog/log"
|
||||
)
|
||||
|
||||
func DeriveSubdir(originRoot, contentPath string) string {
|
||||
rel, err := filepath.Rel(originRoot, contentPath)
|
||||
if err != nil || strings.HasPrefix(rel, "..") {
|
||||
log.Warn().
|
||||
Str("origin_root", originRoot).
|
||||
Str("content_path", contentPath).
|
||||
Msg("content path outside origin root, falling back to full rescan")
|
||||
return ""
|
||||
}
|
||||
return rel
|
||||
}
|
||||
@@ -0,0 +1,60 @@
|
||||
package musicfs
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
func TestDeriveSubdir(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
originRoot string
|
||||
contentPath string
|
||||
expected string
|
||||
}{
|
||||
{
|
||||
name: "normal subdirectory",
|
||||
originRoot: "/downloads/music",
|
||||
contentPath: "/downloads/music/Metallica - Master of Puppets (1986) [FLAC]",
|
||||
expected: "Metallica - Master of Puppets (1986) [FLAC]",
|
||||
},
|
||||
{
|
||||
name: "nested subdirectory",
|
||||
originRoot: "/downloads/music",
|
||||
contentPath: "/downloads/music/Metal/Metallica/Master of Puppets",
|
||||
expected: "Metal/Metallica/Master of Puppets",
|
||||
},
|
||||
{
|
||||
name: "same directory",
|
||||
originRoot: "/downloads/music",
|
||||
contentPath: "/downloads/music",
|
||||
expected: ".",
|
||||
},
|
||||
{
|
||||
name: "outside origin root falls back to empty",
|
||||
originRoot: "/downloads/music",
|
||||
contentPath: "/tmp/other/path",
|
||||
expected: "",
|
||||
},
|
||||
{
|
||||
name: "parent directory falls back to empty",
|
||||
originRoot: "/downloads/music/sub",
|
||||
contentPath: "/downloads/music",
|
||||
expected: "",
|
||||
},
|
||||
{
|
||||
name: "trailing slash on root",
|
||||
originRoot: "/downloads/music/",
|
||||
contentPath: "/downloads/music/Album",
|
||||
expected: "Album",
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
result := DeriveSubdir(tt.originRoot, tt.contentPath)
|
||||
assert.Equal(t, tt.expected, result)
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
package musicfs
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
|
||||
"github.com/rs/zerolog/log"
|
||||
|
||||
pb "homelab.lan/music-agregator/gen/musicfs/v1"
|
||||
)
|
||||
|
||||
type SyncedFile struct {
|
||||
Path string
|
||||
FileID int64
|
||||
VirtualPath string
|
||||
}
|
||||
|
||||
type RescanResult struct {
|
||||
NewFiles []SyncedFile
|
||||
BytesSynced uint64
|
||||
}
|
||||
|
||||
func (c *Client) TriggerRescan(ctx context.Context, originID string, subdir string) (*RescanResult, error) {
|
||||
stream, err := c.MusicFS.RescanOrigin(ctx, &pb.OriginRequest{
|
||||
OriginId: originID,
|
||||
Subdir: &subdir,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("rescan origin: %w", err)
|
||||
}
|
||||
|
||||
var result RescanResult
|
||||
for {
|
||||
progress, err := stream.Recv()
|
||||
if err == io.EOF {
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("rescan stream: %w", err)
|
||||
}
|
||||
|
||||
log.Trace().
|
||||
Str("phase", progress.GetPhase()).
|
||||
Uint32("current", progress.GetCurrent()).
|
||||
Uint32("total", progress.GetTotal()).
|
||||
Msg("rescan progress")
|
||||
|
||||
if progress.GetPhase() == "complete" {
|
||||
for _, sf := range progress.GetNewFiles() {
|
||||
result.NewFiles = append(result.NewFiles, SyncedFile{
|
||||
Path: sf.GetPath(),
|
||||
FileID: sf.GetFileId(),
|
||||
VirtualPath: sf.GetVirtualPath(),
|
||||
})
|
||||
}
|
||||
result.BytesSynced = progress.GetBytesSynced()
|
||||
}
|
||||
}
|
||||
|
||||
return &result, nil
|
||||
}
|
||||
+33
-1
@@ -571,6 +571,11 @@ func (service *MusicAgregatorService) MonitorAlbum(ctx context.Context, req *pb.
|
||||
return nil, status.Errorf(codes.NotFound, "no releases match quality filter for %s", album.GetTitle())
|
||||
}
|
||||
|
||||
filtered = filterByActiveSeeders(filtered)
|
||||
if len(filtered) == 0 {
|
||||
return nil, status.Errorf(codes.NotFound, "no releases with active seeders for %s", album.GetTitle())
|
||||
}
|
||||
|
||||
best := selectBestRelease(filtered)
|
||||
|
||||
if err := service.addToTorrentClient(best); err != nil {
|
||||
@@ -578,10 +583,15 @@ func (service *MusicAgregatorService) MonitorAlbum(ctx context.Context, req *pb.
|
||||
}
|
||||
|
||||
dbAlbum, _ = service.metadata.GetAlbumByExternalID(ctx, album.GetId())
|
||||
if dbAlbum == nil {
|
||||
service.metadata.PersistArtist(ctx, album, database.Monitored)
|
||||
service.metadata.PersistAlbum(ctx, album, database.Monitored)
|
||||
dbAlbum, _ = service.metadata.GetAlbumByExternalID(ctx, album.GetId())
|
||||
}
|
||||
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")
|
||||
log.Warn().Str("album_id", req.GetAlbumId()).Msg("album not in DB after persist attempt, skipping torrent/download persistence")
|
||||
}
|
||||
|
||||
return service.buildMonitorAlbumResponse(ctx, album, dbAlbum, &best), nil
|
||||
@@ -735,6 +745,28 @@ func filterByQuality(items []parsedItem, quality pb.QualityType) []parsedItem {
|
||||
return filtered
|
||||
}
|
||||
|
||||
func filterByActiveSeeders(items []parsedItem) []parsedItem {
|
||||
var withSeeders []parsedItem
|
||||
for _, p := range items {
|
||||
if p.realSeeders > 0 {
|
||||
withSeeders = append(withSeeders, p)
|
||||
} else if p.realSeeders == -1 {
|
||||
withSeeders = append(withSeeders, p)
|
||||
} else {
|
||||
log.Debug().
|
||||
Str("title", p.item.Title).
|
||||
Int("reported_seeders", p.item.Seeders).
|
||||
Int("real_seeders", p.realSeeders).
|
||||
Msg("filtered out: magnet resolved but no active seeders")
|
||||
}
|
||||
}
|
||||
if len(withSeeders) == 0 {
|
||||
log.Warn().Int("total", len(items)).Msg("no torrents with active seeders, keeping all as fallback")
|
||||
return items
|
||||
}
|
||||
return withSeeders
|
||||
}
|
||||
|
||||
func selectBestRelease(items []parsedItem) parsedItem {
|
||||
best := items[0]
|
||||
for _, p := range items[1:] {
|
||||
|
||||
@@ -10,6 +10,8 @@ import (
|
||||
|
||||
"homelab.lan/music-agregator/internal/analysis"
|
||||
"homelab.lan/music-agregator/internal/database"
|
||||
"homelab.lan/music-agregator/internal/metadata"
|
||||
"homelab.lan/music-agregator/internal/musicfs"
|
||||
"homelab.lan/music-agregator/internal/torrent"
|
||||
)
|
||||
|
||||
@@ -31,6 +33,10 @@ type PollDownloadWorker struct {
|
||||
RiverClient *river.Client[pgx.Tx]
|
||||
PathMapper *torrent.PathMapper
|
||||
Analyzer *analysis.ReleaseAnalyzer
|
||||
MusicFSClient *musicfs.Client
|
||||
MusicFSOriginID string
|
||||
MusicFSOriginRoot string
|
||||
MetadataService *metadata.MetadataService
|
||||
}
|
||||
|
||||
func (w *PollDownloadWorker) Work(ctx context.Context, job *river.Job[PollDownloadArgs]) error {
|
||||
@@ -117,9 +123,60 @@ func (w *PollDownloadWorker) onCompleted(ctx context.Context, args PollDownloadA
|
||||
}
|
||||
}
|
||||
|
||||
if w.MusicFSClient != nil {
|
||||
w.syncAndEnrichMusicFS(ctx, args, contentPath)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (w *PollDownloadWorker) syncAndEnrichMusicFS(ctx context.Context, args PollDownloadArgs, contentPath string) {
|
||||
subdir := musicfs.DeriveSubdir(w.MusicFSOriginRoot, contentPath)
|
||||
|
||||
result, err := w.MusicFSClient.TriggerRescan(ctx, w.MusicFSOriginID, subdir)
|
||||
if err != nil {
|
||||
log.Warn().Err(err).Msg("musicfs rescan failed")
|
||||
return
|
||||
}
|
||||
|
||||
log.Info().
|
||||
Int("new_files", len(result.NewFiles)).
|
||||
Uint64("bytes_synced", result.BytesSynced).
|
||||
Msg("musicfs rescan complete")
|
||||
|
||||
if len(result.NewFiles) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
download, err := w.Downloads.GetByID(ctx, args.DownloadID)
|
||||
if err != nil {
|
||||
log.Warn().Err(err).Msg("failed to get download for enrichment")
|
||||
return
|
||||
}
|
||||
|
||||
if w.MetadataService == nil {
|
||||
log.Warn().Msg("no metadata service configured, skipping enrichment")
|
||||
return
|
||||
}
|
||||
|
||||
album, err := w.MetadataService.GetAlbumByID(ctx, download.AlbumID)
|
||||
if err != nil {
|
||||
log.Warn().Err(err).Str("album_id", download.AlbumID).Msg("failed to get album for enrichment")
|
||||
return
|
||||
}
|
||||
|
||||
albumMeta, err := w.MetadataService.GetAlbum(ctx, album.ExternalID)
|
||||
if err != nil {
|
||||
log.Warn().Err(err).Str("external_id", album.ExternalID).Msg("failed to get album metadata for enrichment")
|
||||
return
|
||||
}
|
||||
|
||||
err = w.MusicFSClient.EnrichFiles(ctx, result.NewFiles, albumMeta)
|
||||
if err != nil {
|
||||
log.Warn().Err(err).Msg("musicfs enrichment failed")
|
||||
}
|
||||
}
|
||||
|
||||
func (w *PollDownloadWorker) reschedule(ctx context.Context, args PollDownloadArgs) error {
|
||||
if w.RiverClient == nil {
|
||||
log.Warn().Str("download_id", args.DownloadID).Msg("no river client, cannot reschedule poll_download")
|
||||
|
||||
Reference in New Issue
Block a user