From 76e426c5d3ae6dcc570c670b35f25b0c541328ab Mon Sep 17 00:00:00 2001 From: Alexander Date: Sun, 17 May 2026 23:32:44 +0200 Subject: [PATCH] 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 --- buf.gen.yaml | 1 + cmd/music-agregator/main.go | 41 +++++++++++--- internal/config/config.go | 8 +++ internal/monitor_workflow.go | 15 ++++- internal/musicfs/client.go | 32 +++++++++++ internal/musicfs/enricher.go | 92 +++++++++++++++++++++++++++++++ internal/musicfs/subdir.go | 20 +++++++ internal/musicfs/subdir_test.go | 60 ++++++++++++++++++++ internal/musicfs/sync.go | 62 +++++++++++++++++++++ internal/service.go | 34 +++++++++++- internal/workers/poll_download.go | 73 +++++++++++++++++++++--- 11 files changed, 419 insertions(+), 19 deletions(-) create mode 100644 internal/musicfs/client.go create mode 100644 internal/musicfs/enricher.go create mode 100644 internal/musicfs/subdir.go create mode 100644 internal/musicfs/subdir_test.go create mode 100644 internal/musicfs/sync.go diff --git a/buf.gen.yaml b/buf.gen.yaml index 588644e..328095c 100644 --- a/buf.gen.yaml +++ b/buf.gen.yaml @@ -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 diff --git a/cmd/music-agregator/main.go b/cmd/music-agregator/main.go index 00642b3..e86826e 100644 --- a/cmd/music-agregator/main.go +++ b/cmd/music-agregator/main.go @@ -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,16 +84,20 @@ 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), - DownloadFiles: database.NewDownloadFileRepository(db.Pool), - AlbumReleases: database.NewAlbumReleaseRepository(db.Pool), - TrackReleases: database.NewTrackReleaseRepository(db.Pool), - TorrentClient: torrentClient, - PathMapper: pathMapper, - Analyzer: analysis.NewReleaseAnalyzer(db), + Downloads: database.NewDownloadRepository(db.Pool), + DownloadFiles: database.NewDownloadFileRepository(db.Pool), + AlbumReleases: database.NewAlbumReleaseRepository(db.Pool), + TrackReleases: database.NewTrackReleaseRepository(db.Pool), + 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 { diff --git a/internal/config/config.go b/internal/config/config.go index c2e5627..942ed7f 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -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 { diff --git a/internal/monitor_workflow.go b/internal/monitor_workflow.go index 462213d..2fdb0bd 100644 --- a/internal/monitor_workflow.go +++ b/internal/monitor_workflow.go @@ -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 diff --git a/internal/musicfs/client.go b/internal/musicfs/client.go new file mode 100644 index 0000000..6c650f0 --- /dev/null +++ b/internal/musicfs/client.go @@ -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 +} diff --git a/internal/musicfs/enricher.go b/internal/musicfs/enricher.go new file mode 100644 index 0000000..d143eaa --- /dev/null +++ b/internal/musicfs/enricher.go @@ -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() +} diff --git a/internal/musicfs/subdir.go b/internal/musicfs/subdir.go new file mode 100644 index 0000000..65db4df --- /dev/null +++ b/internal/musicfs/subdir.go @@ -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 +} diff --git a/internal/musicfs/subdir_test.go b/internal/musicfs/subdir_test.go new file mode 100644 index 0000000..0738021 --- /dev/null +++ b/internal/musicfs/subdir_test.go @@ -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) + }) + } +} diff --git a/internal/musicfs/sync.go b/internal/musicfs/sync.go new file mode 100644 index 0000000..1353a3f --- /dev/null +++ b/internal/musicfs/sync.go @@ -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 +} diff --git a/internal/service.go b/internal/service.go index bb2723e..6871627 100644 --- a/internal/service.go +++ b/internal/service.go @@ -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:] { diff --git a/internal/workers/poll_download.go b/internal/workers/poll_download.go index 91c124c..d574421 100644 --- a/internal/workers/poll_download.go +++ b/internal/workers/poll_download.go @@ -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" ) @@ -23,14 +25,18 @@ func (PollDownloadArgs) Kind() string { return "poll_download" } type PollDownloadWorker struct { river.WorkerDefaults[PollDownloadArgs] - TorrentClient torrent.TorrentClient - Downloads *database.DownloadRepository - DownloadFiles *database.DownloadFileRepository - AlbumReleases *database.AlbumReleaseRepository - TrackReleases *database.TrackReleaseRepository - RiverClient *river.Client[pgx.Tx] - PathMapper *torrent.PathMapper - Analyzer *analysis.ReleaseAnalyzer + TorrentClient torrent.TorrentClient + Downloads *database.DownloadRepository + DownloadFiles *database.DownloadFileRepository + AlbumReleases *database.AlbumReleaseRepository + TrackReleases *database.TrackReleaseRepository + 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")