From 133f6219ae97ee03ed5f207dd14440266443fb03 Mon Sep 17 00:00:00 2001 From: Jonathan Date: Sat, 11 Jul 2026 12:35:07 +0200 Subject: [PATCH] feat: monitor sweep, real MB browser, and worker-loop wiring --- worker/lyra_worker/_mbbrowser.py | 44 +++++++++++++++++++++++++++++ worker/lyra_worker/main.py | 22 ++++++++++++++- worker/lyra_worker/monitor.py | 6 ++++ worker/lyra_worker/registry.py | 7 +++++ worker/tests/test_mbbrowser_live.py | 18 ++++++++++++ worker/tests/test_monitor_sweep.py | 22 +++++++++++++++ 6 files changed, 118 insertions(+), 1 deletion(-) create mode 100644 worker/lyra_worker/_mbbrowser.py create mode 100644 worker/tests/test_mbbrowser_live.py create mode 100644 worker/tests/test_monitor_sweep.py diff --git a/worker/lyra_worker/_mbbrowser.py b/worker/lyra_worker/_mbbrowser.py new file mode 100644 index 0000000..4253e89 --- /dev/null +++ b/worker/lyra_worker/_mbbrowser.py @@ -0,0 +1,44 @@ +from lyra_worker.browser import ArtistHit, ReleaseGroupInfo + + +class MusicBrainzBrowser: + """Real MbBrowser via musicbrainzngs (imported lazily; not unit-tested offline).""" + + def __init__(self, app_name: str = "Lyra", version: str = "0.1", contact: str = "lyra@localhost"): + self._app, self._version, self._contact = app_name, version, contact + + def _ua(self): + import musicbrainzngs + + musicbrainzngs.set_useragent(self._app, self._version, self._contact) + return musicbrainzngs + + def search_artist(self, name: str) -> list[ArtistHit]: + mb = self._ua() + res = mb.search_artists(query=name, limit=8) + return [ + ArtistHit(mbid=a["id"], name=a.get("name", ""), disambiguation=a.get("disambiguation", "")) + for a in res.get("artist-list", []) + ] + + def browse_release_groups(self, artist_mbid: str) -> list[ReleaseGroupInfo]: + mb = self._ua() + out: list[ReleaseGroupInfo] = [] + offset = 0 + while True: + res = mb.browse_release_groups(artist=artist_mbid, limit=100, offset=offset) + groups = res.get("release-group-list", []) + for g in groups: + out.append( + ReleaseGroupInfo( + rg_mbid=g["id"], + title=g.get("title", ""), + primary_type=g.get("primary-type", "") or "", + secondary_types=tuple(g.get("secondary-type-list", []) or ()), + first_release_date=g.get("first-release-date", "") or "", + ) + ) + offset += len(groups) + if len(groups) < 100 or offset >= int(res.get("release-group-count", offset)): + break + return out diff --git a/worker/lyra_worker/main.py b/worker/lyra_worker/main.py index d9b0b42..4d8af2a 100644 --- a/worker/lyra_worker/main.py +++ b/worker/lyra_worker/main.py @@ -3,10 +3,12 @@ import time 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.monitor import MonitorConfig, reconcile, sweep from lyra_worker.pipeline import run_pipeline -from lyra_worker.registry import build_adapters, build_resolver, build_tagger +from lyra_worker.registry import build_adapters, build_browser, build_resolver, build_tagger IDLE_SLEEP = 2.0 +MONITOR_TICK_SECONDS = 60.0 def run_forever() -> None: @@ -14,16 +16,34 @@ def run_forever() -> None: adapters = build_adapters(get_config(conn)) resolver = build_resolver() tagger = build_tagger() + browser = build_browser() print(f"worker: {len(adapters)} adapter(s) enabled: {[a.name for a in adapters]}", flush=True) print("worker: waiting for jobs", flush=True) + last_tick = 0.0 try: while True: + mcfg = MonitorConfig.from_config(get_config(conn)) + now = time.monotonic() + if mcfg.enabled and now - last_tick >= MONITOR_TICK_SECONDS: + try: + sweep(conn, browser, mcfg) + except Exception as e: # a monitor error must never kill the worker + print(f"worker: monitor sweep failed: {e}", flush=True) + conn.rollback() + last_tick = now + job_id = claim_next(conn) if job_id is None: time.sleep(IDLE_SLEEP) continue print(f"worker: claimed job {job_id}", flush=True) run_pipeline(conn, job_id, adapters, resolver=resolver, tagger=tagger) + if mcfg.enabled: + try: + reconcile(conn, job_id, mcfg) + except Exception as e: + print(f"worker: reconcile failed for job {job_id}: {e}", flush=True) + conn.rollback() print(f"worker: finished job {job_id}", flush=True) finally: conn.close() diff --git a/worker/lyra_worker/monitor.py b/worker/lyra_worker/monitor.py index a884056..b7fcbf8 100644 --- a/worker/lyra_worker/monitor.py +++ b/worker/lyra_worker/monitor.py @@ -172,3 +172,9 @@ def reconcile(conn: psycopg.Connection, job_id: str, cfg: MonitorConfig) -> None (quality, new_state, rel_id), ) conn.commit() + + +def sweep(conn: psycopg.Connection, browser: MbBrowser, cfg: MonitorConfig) -> None: + """One monitor pass: discover new releases, then enqueue everything due.""" + discover(conn, browser, cfg) + enqueue_due(conn, cfg) diff --git a/worker/lyra_worker/registry.py b/worker/lyra_worker/registry.py index a123117..adead8b 100644 --- a/worker/lyra_worker/registry.py +++ b/worker/lyra_worker/registry.py @@ -1,3 +1,4 @@ +from lyra_worker._mbbrowser import MusicBrainzBrowser from lyra_worker._musicbrainz import MusicBrainzResolver from lyra_worker._mutagen import MutagenTagger from lyra_worker.adapters._slskd import SlskdClient @@ -7,6 +8,7 @@ from lyra_worker.adapters.base import SourceAdapter from lyra_worker.adapters.qobuz import QobuzAdapter from lyra_worker.adapters.soulseek import SoulseekAdapter from lyra_worker.adapters.youtube import YouTubeAdapter +from lyra_worker.browser import MbBrowser from lyra_worker.resolver import MbResolver from lyra_worker.tagger import Tagger @@ -33,3 +35,8 @@ def build_resolver() -> MbResolver: def build_tagger() -> Tagger: """The tagger that writes canonical metadata onto downloaded files.""" return MutagenTagger() + + +def build_browser() -> MbBrowser: + """The MusicBrainz browser used by the background monitor.""" + return MusicBrainzBrowser() diff --git a/worker/tests/test_mbbrowser_live.py b/worker/tests/test_mbbrowser_live.py new file mode 100644 index 0000000..ee4b430 --- /dev/null +++ b/worker/tests/test_mbbrowser_live.py @@ -0,0 +1,18 @@ +import os + +import pytest + +pytestmark = pytest.mark.skipif( + not os.environ.get("LYRA_LIVE_TESTS"), + reason="hits the real MusicBrainz API; set LYRA_LIVE_TESTS=1 to run", +) + + +def test_search_and_browse_real_artist(): + from lyra_worker._mbbrowser import MusicBrainzBrowser + + browser = MusicBrainzBrowser() + hits = browser.search_artist("John Mayer") + assert hits and any("John Mayer" in h.name for h in hits) + releases = browser.browse_release_groups(hits[0].mbid) + assert any(r.title for r in releases) diff --git a/worker/tests/test_monitor_sweep.py b/worker/tests/test_monitor_sweep.py new file mode 100644 index 0000000..45a7836 --- /dev/null +++ b/worker/tests/test_monitor_sweep.py @@ -0,0 +1,22 @@ +from lyra_worker.adapters.fakes import FakeMbBrowser +from lyra_worker.browser import ReleaseGroupInfo +from lyra_worker.monitor import MonitorConfig, sweep +from tests.conftest import insert_watched_artist + +CFG = MonitorConfig(enabled=True) + + +def test_sweep_discovers_then_enqueues(conn): + insert_watched_artist(conn, mbid="a1", name="John Mayer", auto_monitor_future=True) + browser = FakeMbBrowser( + releases={"a1": [ReleaseGroupInfo("rg-new", "New", "Album", (), "2999-01-01")]} + ) + sweep(conn, browser, CFG) + with conn.cursor() as cur: + cur.execute( + 'SELECT count(*) FROM "Request" r ' + 'JOIN "MonitoredRelease" mr ON mr.id = r."monitoredReleaseId" ' + 'WHERE mr."rgMbid" = %s', + ("rg-new",), + ) + assert cur.fetchone()[0] == 1 # the newly-discovered monitored release got a job