From 1596bab23e7146ed8ad158e1a8c616766a03d281 Mon Sep 17 00:00:00 2001 From: Jonathan Date: Fri, 10 Jul 2026 23:59:28 +0200 Subject: [PATCH] feat: wire the mutagen tagger into the worker loop --- worker/lyra_worker/main.py | 5 ++-- worker/lyra_worker/registry.py | 7 +++++ worker/tests/test_tag_pipeline.py | 49 +++++++++++++++++++++++++++++++ 3 files changed, 59 insertions(+), 2 deletions(-) create mode 100644 worker/tests/test_tag_pipeline.py diff --git a/worker/lyra_worker/main.py b/worker/lyra_worker/main.py index f1dd1bf..d9b0b42 100644 --- a/worker/lyra_worker/main.py +++ b/worker/lyra_worker/main.py @@ -4,7 +4,7 @@ from lyra_worker.claim import claim_next from lyra_worker.config import get_config from lyra_worker.db import wait_for_db from lyra_worker.pipeline import run_pipeline -from lyra_worker.registry import build_adapters, build_resolver +from lyra_worker.registry import build_adapters, build_resolver, build_tagger IDLE_SLEEP = 2.0 @@ -13,6 +13,7 @@ def run_forever() -> None: conn = wait_for_db() adapters = build_adapters(get_config(conn)) resolver = build_resolver() + tagger = build_tagger() print(f"worker: {len(adapters)} adapter(s) enabled: {[a.name for a in adapters]}", flush=True) print("worker: waiting for jobs", flush=True) try: @@ -22,7 +23,7 @@ def run_forever() -> None: time.sleep(IDLE_SLEEP) continue print(f"worker: claimed job {job_id}", flush=True) - run_pipeline(conn, job_id, adapters, resolver=resolver) + run_pipeline(conn, job_id, adapters, resolver=resolver, tagger=tagger) print(f"worker: finished job {job_id}", flush=True) finally: conn.close() diff --git a/worker/lyra_worker/registry.py b/worker/lyra_worker/registry.py index 768523f..a123117 100644 --- a/worker/lyra_worker/registry.py +++ b/worker/lyra_worker/registry.py @@ -1,4 +1,5 @@ from lyra_worker._musicbrainz import MusicBrainzResolver +from lyra_worker._mutagen import MutagenTagger from lyra_worker.adapters._slskd import SlskdClient from lyra_worker.adapters._streamrip import StreamripClient from lyra_worker.adapters._ytdlp import YtDlpClient @@ -7,6 +8,7 @@ from lyra_worker.adapters.qobuz import QobuzAdapter from lyra_worker.adapters.soulseek import SoulseekAdapter from lyra_worker.adapters.youtube import YouTubeAdapter from lyra_worker.resolver import MbResolver +from lyra_worker.tagger import Tagger def build_adapters(config: dict) -> list[SourceAdapter]: @@ -26,3 +28,8 @@ def build_adapters(config: dict) -> list[SourceAdapter]: def build_resolver() -> MbResolver: """The MusicBrainz resolver used to canonicalize requests in intake.""" return MusicBrainzResolver() + + +def build_tagger() -> Tagger: + """The tagger that writes canonical metadata onto downloaded files.""" + return MutagenTagger() diff --git a/worker/tests/test_tag_pipeline.py b/worker/tests/test_tag_pipeline.py new file mode 100644 index 0000000..14ecad0 --- /dev/null +++ b/worker/tests/test_tag_pipeline.py @@ -0,0 +1,49 @@ +import os + +from lyra_worker.adapters.youtube import YouTubeAdapter +from lyra_worker.claim import claim_next +from lyra_worker.pipeline import run_pipeline +from lyra_worker.registry import build_tagger +from lyra_worker.types import MBTarget +from tests.conftest import insert_request +from tests.test_youtube_adapter import FakeYtClient + + +def test_build_tagger_returns_a_tagger(): + t = build_tagger() + assert hasattr(t, "tag_album") + + +def test_download_dir_receives_files_and_tagger_runs(conn, tmp_path): + # A fake YouTube client that actually writes a file into dest, so a real MutagenTagger-shaped + # tagger has something to operate on. Here we use a recording tagger to assert it's invoked + # on the album directory the download wrote to. + written = {} + + def fake_download(source_ref, dest, on_progress): + os.makedirs(dest, exist_ok=True) + open(os.path.join(dest, "track.opus"), "wb").close() + written["dir"] = dest + on_progress(1.0) + return {"track_count": 1, "path": dest} + + client = FakeYtClient(results=[{"source_ref": "yt1", "title": "Continuum", "artist": "John Mayer", "track_count": 1}]) + client.download = fake_download + + class RecordingTagger: + def __init__(self): + self.seen = None + + def tag_album(self, album_dir, target): + self.seen = album_dir + + tagger = RecordingTagger() + job_id = insert_request(conn, artist="John Mayer", album="Continuum") + claim_next(conn) + run_pipeline(conn, job_id, [YouTubeAdapter(client)], tagger=tagger, dest_root=str(tmp_path)) + + with conn.cursor() as cur: + cur.execute('SELECT state FROM "Job" WHERE id = %s', (job_id,)) + assert cur.fetchone()[0] == "imported" + assert tagger.seen == written["dir"] # tagger ran on the exact directory the download wrote to + assert os.path.isdir(tagger.seen)