feat: monitor sweep, real MB browser, and worker-loop wiring
This commit is contained in:
@@ -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
|
||||
@@ -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()
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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)
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user