Files
Lyra/worker/lyra_worker/pipeline.py
T
Jonathan 675d964bbe fix(worker): smarter Soulseek peer selection + stall/progress-aware transfers
Draining the queue via Soulseek was stalling: slow peers (~100KB/s) never finished a
full FLAC album within the fixed 10-min timeout, got cancelled at ~89%, and the
fall-through restarted from scratch on the next peer — while single-song folders named
like the album were also being tried.

- Stall detection: abandon a peer with no byte progress for ~90s (queued/dead) instead
  of waiting the full timeout; a peer that IS progressing keeps its slot up to a 60-min
  backstop, so a slow-but-working transfer can finish.
- Rank search candidates by peer: free upload slot, then speed, then shortest queue.
- Drop candidates with far fewer tracks than the release (single-song folders) up front,
  so no download attempt is wasted on them.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-15 19:54:14 +02:00

504 lines
24 KiB
Python

import os
import random
import shutil
import threading
import time
from concurrent.futures import ThreadPoolExecutor
from dataclasses import replace
from typing import Callable, Sequence
import psycopg
from lyra_worker.adapters.base import SourceAdapter
from lyra_worker.confidence import score_confidence
from lyra_worker.library import album_dir, import_album, list_audio_files, staging_dir
from lyra_worker.qobuz_gate import gate_open, record_download
from lyra_worker.quality import quality_class
from lyra_worker.ranker import rank_candidates
from lyra_worker.types import Candidate, MBTarget
_AUDIO_EXT = {".flac", ".mp3", ".m4a", ".opus", ".ogg", ".aac", ".wav"}
# Reject a download whose total playtime falls below this fraction of MusicBrainz's expected
# total — catches truncated files / preview clips substituted for real tracks that keep the
# track COUNT right. Generous, so edition/encoding differences don't false-positive; bonus
# tracks (over-long) are always fine.
_DURATION_MIN_RATIO = 0.85
# Cap how many ranked candidates a single job will actually download+verify before giving up.
# Without this, the incomplete-download fall-through can grind through dozens of sources (a popular
# album returns ~200 Soulseek candidates), each a slow attempt that also leaves an abandoned slskd
# transfer. The best candidates rank first, so 6 attempts is plenty.
_MAX_DOWNLOAD_ATTEMPTS = 6
# The best quality any non-Qobuz source offers (Soulseek FLAC = CD-lossless). A library copy at or
# above this can only be improved by hi-res Qobuz, so when Qobuz is gated by pacing an upgrade of
# such a copy is deferred without even searching.
_CD_LOSSLESS_CLASS = 2
def _measure_staged_duration_s(staging: str) -> float | None:
"""Sum the playtime (seconds) of the non-empty staged audio via mutagen. Returns None if
nothing could be measured (no audio, or an unreadable container) so the caller skips the
check rather than failing a good download on a probe hiccup. Lazy mutagen import; not
exercised by the offline fakes (which stage no real audio)."""
try:
import mutagen
except Exception:
return None
total = 0.0
measured = False
for root, _dirs, files in os.walk(staging):
for f in files:
if os.path.splitext(f)[1].lower() not in _AUDIO_EXT:
continue
path = os.path.join(root, f)
try:
if os.path.getsize(path) == 0:
continue
audio = mutagen.File(path)
length = getattr(getattr(audio, "info", None), "length", None)
except Exception:
continue
if length:
total += float(length)
measured = True
return total if measured else None
def _count_staged_audio(staging: str) -> int:
"""Count NON-EMPTY audio files anywhere under the staging dir (streamrip nests them in a
subfolder). 0-byte files are skipped: streamrip leaves an empty placeholder when a track
download fails (e.g. Qobuz's >100-header responses), and a dead track must not count toward
completeness — otherwise a broken album imports as 'complete'."""
n = 0
for root, _dirs, files in os.walk(staging):
for f in files:
if os.path.splitext(f)[1].lower() in _AUDIO_EXT:
try:
if os.path.getsize(os.path.join(root, f)) > 0:
n += 1
except OSError:
pass
return n
def _download_problem(result, staging: str, target, winner, measure_duration) -> str | None:
"""Why a completed download is unusable — ``"incomplete download"`` / ``"download too short
(...)"`` — or None if it's good. Completeness is judged against the CHOSEN SOURCE's own track
count (``winner.track_count``), NOT MusicBrainz's canonical release: MB picks an arbitrary
releases[0] that is often a deluxe/expanded edition with more tracks (and runtime) than the
standard album a source legitimately delivers, so gating on ``target.track_count`` falsely
rejects complete standard-edition downloads (The Script "No Sound Without Silence" 11 vs MB 12,
Skillet "Unleashed" 12 vs MB's 20-track deluxe). Pure — no DB I/O — so the download loop can
call it per candidate and fall through to the next source on a problem."""
expected = winner.track_count
# (a) Truncated: fewer files than the source promised (a track failed mid-download). Count the
# NON-EMPTY audio on disk too: a silently-skipped track leaves a 0-byte placeholder the adapter
# still counts. (`staged and ...` keeps fake adapters that stage nothing on the count path.)
staged = _count_staged_audio(staging)
truncated = result.track_count < expected or bool(staged and staged < expected)
# (b) Implausibly small vs MB's album — a lone-track "full album" video / wrong match, as
# distinct from a legitimately smaller edition. A real edition keeps more than half MB's tracks.
too_small = target.track_count is not None and expected * 2 <= target.track_count
if truncated or too_small:
return "incomplete download"
# Duration: a truncated file or short preview substituted for a real track shows up as playtime
# well under the expected total. Scale MB's total to the DELIVERED edition's size so a smaller
# edition isn't judged against a larger one's runtime. Only when the total is known (measured is
# None for the offline fakes).
if target.total_duration_s:
expected_total = target.total_duration_s
if target.track_count:
expected_total *= winner.track_count / target.track_count
measured = measure_duration(staging)
if measured is not None and measured < expected_total * _DURATION_MIN_RATIO:
return f"download too short ({measured / 60:.0f} of ~{expected_total / 60:.0f} min)"
return None
def _poll_download_progress(job_id: str, staging: str, expected: int, stop: "threading.Event",
reports: "threading.Event") -> None:
"""Until `stop`, periodically write Job.downloadProgress = files-in-staging / expected
(capped 0.99) — the FALLBACK for adapters (streamrip/Qobuz) that report no byte-level
progress. Skips writing while `reports` is set, i.e. the in-flight adapter is driving
progress itself via on_progress (yt-dlp/slskd). Uses its own short-lived connection
(psycopg conns aren't shareable across threads). A missing DSN or expected<=0 just no-ops."""
dsn = os.environ.get("DATABASE_URL")
if not dsn or expected <= 0:
return
conn = psycopg.connect(dsn)
try:
while not stop.is_set():
if not reports.is_set(): # the adapter isn't reporting real progress → estimate
frac = min(_count_staged_audio(staging) / expected, 0.99)
with conn.cursor() as cur:
cur.execute('UPDATE "Job" SET "downloadProgress" = %s WHERE id = %s', (frac, job_id))
conn.commit()
stop.wait(1.5)
except Exception as e: # a progress poller must never affect the job
print(f"pipeline: download progress poller error: {e}", flush=True)
finally:
conn.close()
def _search_adapter(adapter: SourceAdapter, target: MBTarget) -> list[Candidate]:
"""Search one adapter, swallowing failures — a down source contributes no candidates rather
than crashing the job. Runs on a worker thread (searches are parallelized across sources)."""
try:
return list(adapter.search(target))
except Exception as e:
print(f"pipeline: adapter {adapter.name} search failed: {e}", flush=True)
return []
def _set_state(conn: psycopg.Connection, job_id: str, state: str, stage: str) -> None:
with conn.cursor() as cur:
cur.execute(
'UPDATE "Job" SET state = %s, "currentStage" = %s, "updatedAt" = now() WHERE id = %s',
(state, stage, job_id),
)
conn.commit()
def _set_download_progress(conn: psycopg.Connection, job_id: str, frac: float) -> None:
with conn.cursor() as cur:
cur.execute('UPDATE "Job" SET "downloadProgress" = %s WHERE id = %s', (frac, job_id))
conn.commit()
def _make_on_progress(conn: psycopg.Connection, job_id: str, reports: "threading.Event"):
"""Build the on_progress(pct) callback threaded into adapter.download. It marks the
adapter as reporting real byte-level progress (so the file-count poller yields) and
throttle-writes Job.downloadProgress. Called synchronously from adapter.download on the
pipeline's own thread, so it safely reuses `conn`."""
state = {"frac": -1.0, "t": 0.0}
def on_progress(pct: float) -> None:
reports.set()
frac = 0.0 if pct < 0 else 1.0 if pct > 1 else float(pct)
now = time.monotonic()
# throttle DB writes: on a >=1% move or every 0.5s, not once per received byte
if frac - state["frac"] >= 0.01 or (now - state["t"]) >= 0.5:
state["frac"] = frac
state["t"] = now
try:
_set_download_progress(conn, job_id, frac)
except Exception as e: # a progress write must never fail the download
print(f"pipeline: on_progress write failed: {e}", flush=True)
return on_progress
def _request_id(conn: psycopg.Connection, job_id: str) -> str:
with conn.cursor() as cur:
cur.execute('SELECT "requestId" FROM "Job" WHERE id = %s', (job_id,))
return cur.fetchone()[0]
def _fail(conn: psycopg.Connection, job_id: str, reason: str) -> None:
with conn.cursor() as cur:
cur.execute(
'UPDATE "Job" SET state = \'needs_attention\', error = %s, "updatedAt" = now() WHERE id = %s',
(reason, job_id),
)
cur.execute(
'UPDATE "Request" SET status = \'needs_attention\' WHERE id = (SELECT "requestId" FROM "Job" WHERE id = %s)',
(job_id,),
)
conn.commit()
def _load_target(conn: psycopg.Connection, job_id: str) -> MBTarget:
with conn.cursor() as cur:
cur.execute(
'SELECT artist, album FROM "Request" WHERE id = (SELECT "requestId" FROM "Job" WHERE id = %s)',
(job_id,),
)
artist, album = cur.fetchone()
return MBTarget(artist=artist, album=album)
def _already_in_library(conn: psycopg.Connection, target: MBTarget) -> bool:
with conn.cursor() as cur:
cur.execute(
'SELECT 1 FROM "LibraryItem" WHERE artist = %s AND album = %s LIMIT 1',
(target.artist, target.album),
)
return cur.fetchone() is not None
def _is_upgrade_job(conn: psycopg.Connection, job_id: str) -> bool:
with conn.cursor() as cur:
cur.execute(
'SELECT "monitoredReleaseId" IS NOT NULL FROM "Request" '
'WHERE id = (SELECT "requestId" FROM "Job" WHERE id = %s)',
(job_id,),
)
row = cur.fetchone()
return bool(row and row[0])
def _is_force_job(conn: psycopg.Connection, job_id: str) -> bool:
"""A user-forced re-acquire (Library 'Replace / upgrade'): skip the already-in-library
dedupe so an owned album is re-downloaded. The import step still keeps the new copy only
if it's higher quality, so a forced upgrade can never downgrade what's on disk."""
with conn.cursor() as cur:
cur.execute(
'SELECT force FROM "Request" WHERE id = (SELECT "requestId" FROM "Job" WHERE id = %s)',
(job_id,),
)
row = cur.fetchone()
return bool(row and row[0])
def _already_in_library_at_cutoff(conn: psycopg.Connection, target: MBTarget, cutoff: int) -> bool:
with conn.cursor() as cur:
cur.execute(
'SELECT 1 FROM "LibraryItem" WHERE artist = %s AND album = %s '
'AND "qualityClass" >= %s LIMIT 1',
(target.artist, target.album, cutoff),
)
return cur.fetchone() is not None
def _library_quality(conn: psycopg.Connection, target: MBTarget) -> int | None:
with conn.cursor() as cur:
cur.execute(
'SELECT "qualityClass" FROM "LibraryItem" WHERE artist = %s AND album = %s',
(target.artist, target.album),
)
row = cur.fetchone()
return row[0] if row else None
def _keep_existing_and_complete(conn: psycopg.Connection, job_id: str) -> None:
"""Finish a job without importing — keep the existing library copy — and mark its request
completed. Used when a (deferred) upgrade has nothing better to offer right now; the release
stays below the cutoff and is re-attempted later."""
_set_state(conn, job_id, "imported", "import")
with conn.cursor() as cur:
cur.execute("UPDATE \"Request\" SET status = 'completed' WHERE id = %s",
(_request_id(conn, job_id),))
conn.commit()
def _persist_candidates(conn: psycopg.Connection, job_id: str, candidates: list[Candidate]) -> None:
with conn.cursor() as cur:
for c in candidates:
cur.execute(
'INSERT INTO "Candidate" (id, "jobId", source, format, "qualityClass", '
'"trackCount", confidence, "sourceRef", chosen, "createdAt") '
"VALUES (gen_random_uuid()::text, %s, %s, %s, %s, %s, %s, %s, false, now())",
(job_id, c.source, c.quality.fmt, quality_class(c.quality),
c.track_count, c.confidence, c.source_ref),
)
conn.commit()
def _mark_chosen(conn: psycopg.Connection, job_id: str, source_ref: str) -> None:
"""Mark exactly one candidate (the one now being downloaded) as chosen, clearing any
prior choice. Called at each download attempt so `chosen` reflects the source in flight."""
with conn.cursor() as cur:
cur.execute(
'UPDATE "Candidate" SET chosen = ("sourceRef" = %s) WHERE "jobId" = %s',
(source_ref, job_id),
)
conn.commit()
def _import(conn: psycopg.Connection, job_id: str, target: MBTarget,
winner: Candidate, path: str) -> None:
request_id = _request_id(conn, job_id)
track_names = list_audio_files(path) # the real files just imported into the album folder
with conn.cursor() as cur:
cur.execute(
'INSERT INTO "LibraryItem" (id, "requestId", artist, album, path, source, '
'format, "qualityClass", "trackNames", "rgMbid", "artistMbid", "importedAt") '
"VALUES (gen_random_uuid()::text, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, now()) "
'ON CONFLICT (artist, album) DO UPDATE SET '
' "requestId" = EXCLUDED."requestId", path = EXCLUDED.path, source = EXCLUDED.source, '
' format = EXCLUDED.format, "qualityClass" = EXCLUDED."qualityClass", '
' "trackNames" = EXCLUDED."trackNames", "rgMbid" = COALESCE(EXCLUDED."rgMbid", "LibraryItem"."rgMbid"), '
' "artistMbid" = COALESCE(EXCLUDED."artistMbid", "LibraryItem"."artistMbid"), "importedAt" = now() '
'WHERE EXCLUDED."qualityClass" > "LibraryItem"."qualityClass"',
(request_id, target.artist, target.album, path, winner.source,
winner.quality.fmt, quality_class(winner.quality), track_names,
target.rg_mbid or None, target.artist_mbid or None),
)
cur.execute(
"UPDATE \"Request\" SET status = 'completed' WHERE id = %s", (request_id,)
)
conn.commit()
def run_pipeline(
conn: psycopg.Connection,
job_id: str,
adapters: Sequence[SourceAdapter],
resolver=None,
tagger=None,
min_confidence: float = 0.7,
dest_root: str = "/music",
staging_root: str | None = None,
upgrade_cutoff: int | None = None,
measure_duration: Callable[[str], float | None] = _measure_staged_duration_s,
) -> None:
"""Real staged acquisition using source-agnostic adapters. Fakes in this plan."""
names = [a.name for a in adapters]
if len(names) != len(set(names)):
raise ValueError(f"adapter names must be unique, got {names}")
# 1. intake
_set_state(conn, job_id, "matching", "intake")
target = _load_target(conn, job_id)
if resolver is not None:
resolved = resolver.resolve(target.artist, target.album)
if resolved is not None:
target = resolved
# dedupe: a monitor-driven upgrade job proceeds unless the existing copy already meets
# the cutoff; a plain request short-circuits on any existing copy (slice-1 behavior). A
# user-forced re-acquire skips the dedupe entirely (keep-if-better still guards the import).
if _is_force_job(conn, job_id):
_dedupe_hit = False
elif upgrade_cutoff is not None and _is_upgrade_job(conn, job_id):
_dedupe_hit = _already_in_library_at_cutoff(conn, target, upgrade_cutoff)
else:
_dedupe_hit = _already_in_library(conn, target)
if _dedupe_hit:
_set_state(conn, job_id, "imported", "import")
with conn.cursor() as cur:
cur.execute(
"UPDATE \"Request\" SET status = 'completed' WHERE id = %s",
(_request_id(conn, job_id),),
)
conn.commit()
return
# Pre-search upgrade skip: a >= CD-lossless copy can only be improved by hi-res Qobuz. If Qobuz
# is gated by pacing, don't even search (no other source can beat it) — defer to when Qobuz
# frees up. This avoids re-searching every already-owned album (~2 min Soulseek search each) on
# every monitor cycle. A user-forced re-acquire always proceeds.
_existing_q = _library_quality(conn, target)
if _existing_q is not None and _existing_q >= _CD_LOSSLESS_CLASS \
and not _is_force_job(conn, job_id) and not gate_open(conn):
_keep_existing_and_complete(conn, job_id)
return
# 2. match — search all adapters CONCURRENTLY (Soulseek's search polls up to ~2 min; run it
# alongside Qobuz's instant search so total match time = the slowest source, not their sum).
# Only Qobuz touches streamrip's shared asyncio loop and there's a single Qobuz adapter, so no
# two threads drive that loop at once. Results are collected in adapter order (map preserves
# input order) so candidate ordering stays deterministic.
_set_state(conn, job_id, "matching", "match")
found: list[Candidate] = []
with ThreadPoolExecutor(max_workers=max(1, len(adapters))) as pool:
per_adapter = list(pool.map(lambda a: (a, _search_adapter(a, target)), adapters))
for adapter, results in per_adapter:
for c in results:
found.append(replace(c, source_tier=adapter.tier, confidence=score_confidence(target, c)))
# Drop candidates that can't be the full album — a result with far fewer tracks than the release
# has (e.g. a single-song Soulseek folder named like the album). Same threshold as the
# completeness 'too small' guard, applied up front so we never waste a download attempt on it.
# Unknown counts (0) and legitimately smaller editions (e.g. 11 vs 12) are kept.
if target.track_count:
found = [c for c in found if c.track_count == 0 or c.track_count * 2 > target.track_count]
_persist_candidates(conn, job_id, found)
# Qobuz pacing: when the budget gate is closed (off-hours / daily cap / spacing), drop Qobuz
# candidates so the job falls through to another source and Qobuz stays under the radar — UNLESS
# Qobuz is the only option (nothing to fall through to, so don't strand the album).
if any(c.source == "qobuz" for c in found) and any(c.source != "qobuz" for c in found):
if not gate_open(conn):
found = [c for c in found if c.source != "qobuz"]
# 3. rank
_set_state(conn, job_id, "matched", "rank")
ranked = rank_candidates(target, found, min_confidence)
if not ranked:
_fail(conn, job_id, "no candidate above confidence threshold")
return
# Upgrade guard: for an album already in the library, skip the download when no currently-
# available source beats the copy we have (ranked[0] is the highest-quality candidate). This
# happens when Qobuz is gated by pacing and only a same-quality Soulseek copy is left — the
# download would be discarded by keep-if-better, so defer it. The release stays below the cutoff
# and is re-attempted when a better source frees up. A user-forced re-acquire always proceeds.
existing_q = _library_quality(conn, target)
if existing_q is not None and not _is_force_job(conn, job_id) \
and quality_class(ranked[0].quality) <= existing_q:
_keep_existing_and_complete(conn, job_id) # keep the existing copy; nothing better available
return
# 4. download (fall-through) into an isolated per-job staging dir
_set_state(conn, job_id, "downloading", "download")
by_source = {a.name: a for a in adapters}
staging = staging_dir(staging_root or f"{dest_root}/.staging", job_id)
winner = None
last_problem = None # why the most recent downloaded-but-rejected candidate was unusable
_set_download_progress(conn, job_id, 0.0) # reset for this run
expected = target.track_count or (ranked[0].track_count if ranked else 0)
_stop = threading.Event()
_reports = threading.Event() # set while the in-flight adapter reports real byte progress
_poller = threading.Thread(
target=_poll_download_progress, args=(job_id, staging, expected, _stop, _reports), daemon=True
)
_poller.start()
try:
try:
attempts = 0
for candidate in ranked:
adapter = by_source.get(candidate.source)
if adapter is None:
continue
if attempts >= _MAX_DOWNLOAD_ATTEMPTS:
break # don't grind through every peer/source of a popular album
attempts += 1
_mark_chosen(conn, job_id, candidate.source_ref) # reflect the source now in flight
shutil.rmtree(staging, ignore_errors=True) # clean slate per attempt
_reports.clear() # this attempt hasn't reported yet → poller estimates until it does
result = adapter.download(candidate, staging, _make_on_progress(conn, job_id, _reports))
if not result.ok:
continue
# Verify completeness on the staging copy. A download that lands incomplete or too
# short falls through to the NEXT-ranked source (e.g. Qobuz repeatedly fails one
# track → try the deluxe edition or Soulseek) rather than failing the whole job on
# the first source's shortfall.
last_problem = _download_problem(result, staging, target, candidate, measure_duration)
if last_problem is None:
winner = candidate
break
finally:
_stop.set()
_poller.join(timeout=3)
if winner is None:
# No source produced a complete album: surface why the last one was rejected (incomplete
# / too short), or "all downloads failed" if none even downloaded.
_fail(conn, job_id, last_problem or "all downloads failed")
return
_set_download_progress(conn, job_id, 1.0) # download done — UI shows 100% into Finishing
if winner.source == "qobuz":
# Count this Qobuz download against today's budget and set the next randomized spacing.
try:
record_download(conn, random.random())
except Exception as e: # pacing bookkeeping must never fail a good download
print(f"pipeline: qobuz pacing record failed for job {job_id}: {e}", flush=True)
# 5. promote + tag the verified winner
_set_state(conn, job_id, "tagging", "tag")
final = album_dir(dest_root, target)
existing_q = _library_quality(conn, target)
new_q = quality_class(winner.quality)
_set_state(conn, job_id, "imported", "import")
if existing_q is None or new_q > existing_q:
import_album(staging, final) # clean audio(+cover) move, atomically replacing any prior copy
if tagger is not None:
try:
tagger.tag_album(final, target)
except Exception as e: # a tagging failure must not discard a good download
print(f"pipeline: tagging failed for job {job_id}: {e}", flush=True)
_import(conn, job_id, target, winner, final)
else:
# an existing copy is already at >= this quality — keep it untouched, just complete
_keep_existing_and_complete(conn, job_id)
finally:
shutil.rmtree(staging, ignore_errors=True) # a partial/failed download never persists