From a67c552f2944cd81335f44a97343966d0b19dd38 Mon Sep 17 00:00:00 2001 From: Jonathan Date: Mon, 13 Jul 2026 00:57:00 +0200 Subject: [PATCH] feat(worker): chunk the discovery sweep (one chunk per loop iteration) Co-Authored-By: Claude Opus 4.8 (1M context) --- worker/lyra_worker/main.py | 44 ++++++++++++------- worker/tests/test_discovery_trigger.py | 58 +++++++++++++++++++------- 2 files changed, 72 insertions(+), 30 deletions(-) diff --git a/worker/lyra_worker/main.py b/worker/lyra_worker/main.py index acdd8e3..af3fd62 100644 --- a/worker/lyra_worker/main.py +++ b/worker/lyra_worker/main.py @@ -6,7 +6,7 @@ import psycopg 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.discovery import DiscoveryConfig, run_discovery +from lyra_worker.discovery import DiscoveryConfig, DiscoveryResult, run_discovery from lyra_worker.library import clear_staging_root from lyra_worker.monitor import MonitorConfig, fulfill_owned_releases, reconcile, sweep from lyra_worker.pipeline import run_pipeline @@ -96,12 +96,9 @@ def maybe_run_scan(conn, resolver, probe, browser, config, dest_root: str = DEST _set_config(conn, "scan.progress", f"{imported_total}/{skipped_total}") -def _run_discovery(conn, sources, browser, config=None) -> None: +def _run_discovery(conn, sources, browser, config=None) -> DiscoveryResult: cfg = DiscoveryConfig.from_config(config if config is not None else get_config(conn)) - result = run_discovery(conn, sources, browser, cfg) - _set_config(conn, "discover.result", f"artists {result.artists}, albums {result.albums}") - _set_config(conn, "discover.requested", "false") - print(f"worker: discovery done — {result.artists} artists, {result.albums} albums", flush=True) + return run_discovery(conn, sources, browser, cfg) def finish_job(conn, job_id, mcfg) -> None: @@ -117,21 +114,40 @@ def finish_job(conn, job_id, mcfg) -> None: def maybe_run_discovery(conn, sources, browser, config, dcfg, now, last_discover_tick, tick_seconds: float = DISCOVER_TICK_SECONDS) -> float: - """Run discovery at most once per loop iteration, whether triggered by the schedule - (enabled + interval elapsed) or a one-shot ``discover.requested`` flag. A single run - satisfies both, so the two triggers must not fire it twice off the same stale config - snapshot. Returns the (possibly advanced) last_discover_tick.""" + """Advance the discovery sweep by at most one chunk (``discover.chunkSize`` seeds) per call, + so the worker loop keeps claiming jobs between chunks. A sweep starts on ``discover.requested`` + or a due schedule tick, then continues via ``discover.inProgress`` until a chunk drains 0 seeds + (every eligible seed refreshed), at which point ``discover.result`` is written. Returns the + (possibly advanced) last_discover_tick.""" requested = str(config.get("discover.requested", "")).strip().lower() in _TRUE + in_progress = str(config.get("discover.inProgress", "")).strip().lower() in _TRUE due = dcfg.enabled and now - last_discover_tick >= tick_seconds - if not (requested or due): + if not (requested or due or in_progress): return last_discover_tick + + starting = not in_progress + if starting: # begin a fresh sweep + _set_config(conn, "discover.inProgress", "true") + _set_config(conn, "discover.requested", "false") + try: - _run_discovery(conn, sources, browser, config) # also clears discover.requested + result = _run_discovery(conn, sources, browser, config) # one chunk except Exception as e: # a discovery error must never kill the worker print(f"worker: discovery run failed: {e}", flush=True) conn.rollback() - if requested: - _set_config(conn, "discover.requested", "false") + _set_config(conn, "discover.inProgress", "false") # abandon; a new request restarts it + _set_config(conn, "discover.requested", "false") + return now if due else last_discover_tick + + a_prev, b_prev = (0, 0) if starting else _parse_progress(config.get("discover.progress", "")) + a_total, b_total = a_prev + result.artists, b_prev + result.albums + if result.seeds == 0: # every eligible seed refreshed -> sweep drained + _set_config(conn, "discover.result", f"artists {a_total}, albums {b_total}") + _set_config(conn, "discover.progress", "") + _set_config(conn, "discover.inProgress", "false") + print(f"worker: discovery done — {a_total} artists, {b_total} albums", flush=True) + else: + _set_config(conn, "discover.progress", f"{a_total}/{b_total}") return now if due else last_discover_tick diff --git a/worker/tests/test_discovery_trigger.py b/worker/tests/test_discovery_trigger.py index 5b15d6a..0c27833 100644 --- a/worker/tests/test_discovery_trigger.py +++ b/worker/tests/test_discovery_trigger.py @@ -1,7 +1,7 @@ from lyra_worker import main from lyra_worker.adapters.fakes import FakeMbBrowser, FakeSimilaritySource -from lyra_worker.discovery import DiscoveryConfig -from lyra_worker.main import _run_discovery, maybe_run_discovery +from lyra_worker.discovery import DiscoveryConfig, DiscoveryResult +from lyra_worker.main import maybe_run_discovery from lyra_worker.similarity.base import SimilarArtist from tests.conftest import insert_watched_artist @@ -23,25 +23,49 @@ def _get(conn, key): return row[0] if row else None -def test_run_discovery_writes_result_and_clears_flag(conn): - insert_watched_artist(conn, mbid="s1", name="Seed") +def test_maybe_run_discovery_drains_across_chunks_and_finishes(conn): + # The sweep must advance one chunk (discover.chunkSize seeds) per call, clear + # discover.requested exactly once at start, and finally drain on the 0-seed chunk: + # writing discover.result and clearing discover.inProgress. + for i in range(3): + insert_watched_artist(conn, mbid=f"s{i}", name=f"Seed{i}") _set(conn, "discover.requested", "true") - src = FakeSimilaritySource(similar={"s1": [SimilarArtist("c1", "Cand", 0.9)]}) - - _run_discovery(conn, [src], FakeMbBrowser()) - - assert _get(conn, "discover.requested") == "false" # flag cleared - assert "artists" in (_get(conn, "discover.result") or "") # summary written - with conn.cursor() as cur: - cur.execute('SELECT count(*) FROM "DiscoverySuggestion" WHERE kind = \'artist\'') - assert cur.fetchone()[0] == 1 + _set(conn, "discover.chunkSize", "2") + src = FakeSimilaritySource( + similar={f"s{i}": [SimilarArtist(f"c{i}", f"C{i}", 0.9)] for i in range(3)} + ) + cfg = {"discover.enabled": "false", "discover.requested": "true", "discover.chunkSize": "2"} + tick = 0.0 + # iteration 1: starts sweep, processes 2 seeds, still in progress + tick = maybe_run_discovery( + conn, [src], FakeMbBrowser(), cfg, DiscoveryConfig.from_config(cfg), 1.0, tick + ) + assert _get(conn, "discover.inProgress") == "true" + assert _get(conn, "discover.requested") == "false" # cleared at sweep start + # refresh the snapshot the way run_forever's loop does (re-reads config each iteration) + cfg2 = {**cfg, "discover.requested": "false", "discover.inProgress": "true", + "discover.progress": _get(conn, "discover.progress")} + # iteration 2: processes the last seed, still in progress + maybe_run_discovery( + conn, [src], FakeMbBrowser(), cfg2, DiscoveryConfig.from_config(cfg2), 2.0, tick + ) + assert _get(conn, "discover.inProgress") == "true" + cfg3 = {**cfg2, "discover.progress": _get(conn, "discover.progress")} + # iteration 3: 0 seeds left -> drains, writes result, clears inProgress + maybe_run_discovery( + conn, [src], FakeMbBrowser(), cfg3, DiscoveryConfig.from_config(cfg3), 3.0, tick + ) + assert _get(conn, "discover.inProgress") == "false" + assert _get(conn, "discover.progress") == "" + assert "artists" in (_get(conn, "discover.result") or "") def test_maybe_run_discovery_runs_once_when_due_and_requested(conn, monkeypatch): # A scheduled tick that is due AND a pending one-shot request in the same loop # iteration must trigger discovery exactly once, not twice. calls = [] - monkeypatch.setattr(main, "_run_discovery", lambda *a, **k: calls.append(1)) + monkeypatch.setattr(main, "_run_discovery", + lambda *a, **k: (calls.append(1), DiscoveryResult(0, 0, seeds=0))[1]) config = {"discover.enabled": "true", "discover.requested": "true"} dcfg = DiscoveryConfig.from_config(config) new_tick = maybe_run_discovery( @@ -55,7 +79,8 @@ def test_maybe_run_discovery_requested_when_disabled_does_not_advance_tick(conn, # A one-shot request runs even when discovery is disabled, but must not advance the # schedule tick (there is no active schedule). calls = [] - monkeypatch.setattr(main, "_run_discovery", lambda *a, **k: calls.append(1)) + monkeypatch.setattr(main, "_run_discovery", + lambda *a, **k: (calls.append(1), DiscoveryResult(0, 0, seeds=0))[1]) config = {"discover.enabled": "false", "discover.requested": "true"} dcfg = DiscoveryConfig.from_config(config) new_tick = maybe_run_discovery( @@ -67,7 +92,8 @@ def test_maybe_run_discovery_requested_when_disabled_does_not_advance_tick(conn, def test_maybe_run_discovery_idle_does_nothing(conn, monkeypatch): calls = [] - monkeypatch.setattr(main, "_run_discovery", lambda *a, **k: calls.append(1)) + monkeypatch.setattr(main, "_run_discovery", + lambda *a, **k: (calls.append(1), DiscoveryResult(0, 0, seeds=0))[1]) config = {"discover.enabled": "true", "discover.requested": "false"} dcfg = DiscoveryConfig.from_config(config) # enabled but not yet due (tick was recent relative to now)