diff --git a/docs/superpowers/plans/2026-07-10-lyra-pipeline-core.md b/docs/superpowers/plans/2026-07-10-lyra-pipeline-core.md new file mode 100644 index 0000000..129d65b --- /dev/null +++ b/docs/superpowers/plans/2026-07-10-lyra-pipeline-core.md @@ -0,0 +1,1187 @@ +# Lyra Pipeline Core Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Replace the walking skeleton's fake stage-walk with a real, source-agnostic six-stage acquisition pipeline — a uniform `SourceAdapter` contract, quality-first ranking gated by a confidence score, candidate persistence, fall-through on failure, and dedupe-on-import — driven entirely by **fake** adapters so it runs with zero credentials. + +**Architecture:** The Python worker keeps its claim loop and state machine (Plan 1). The fake `run_pipeline` is replaced by staged functions (`intake → match → rank → download → tag → import`) that talk only to the `SourceAdapter` interface. Source complexity is quarantined behind that interface; this plan ships three fake adapters plus a failing one. New Postgres tables `Candidate` and `LibraryItem` (Prisma-owned) record what Match found and what Import committed. + +**Tech Stack:** Python 3.12 (worker), psycopg v3, pytest; Prisma 6 / PostgreSQL 17 (schema). No new third-party Python dependencies — confidence scoring uses stdlib `difflib`. + +## Global Constraints + +- Python 3.12 in the worker container; local test runs may use newer 3.x — no version-specific syntax beyond 3.12. Database access via `psycopg` v3; raw SQL only in the worker (no ORM/DDL — **Prisma owns the schema**). +- No new third-party Python packages. Confidence/text similarity uses `difflib.SequenceMatcher` from the stdlib. +- The claim loop, job state machine, and existing job states are unchanged from Plan 1. Job states remain exactly: `requested, matching, matched, downloading, tagging, imported, needs_attention`. Request statuses remain exactly: `pending, completed, needs_attention`. +- **Quality-first ranking order (best→worst), verbatim from the design spec:** Qobuz hi-res → Qobuz lossless → Soulseek FLAC → Soulseek MP3 → YouTube. Encode as `(quality_class desc, source_tier asc)` where `quality_class`: 3=hi-res lossless, 2=CD lossless, 1=lossy; `source_tier`: qobuz=0, soulseek=1, youtube=2. +- Adapters are queried source-agnostically; adding a source is a new adapter and nothing else. Ranking/confidence never contain source-specific branches. +- Fall-through is mandatory: within ranked candidates, a failed download tries the next; if none succeed (or none pass the confidence gate), the job → `needs_attention`, never silently dropped. +- All new SQL uses quoted Prisma identifiers (`"Job"`, `"Request"`, `"Candidate"`, `"LibraryItem"`, camelCase columns) and `%s` parameters (never string interpolation). +- Every task ends with a commit. + +## Shared interfaces (defined by the tasks below; listed here so tasks stay consistent) + +```python +# lyra_worker/types.py (Task 1) +@dataclass(frozen=True) +class MBTarget: + artist: str + album: str + track_count: int | None = None # expected tracks (None until Plan 4's MusicBrainz) + total_duration_s: int | None = None + +@dataclass(frozen=True) +class Quality: + fmt: str # "FLAC" | "MP3" | "OPUS" | "AAC" + lossless: bool + bit_depth: int | None = None # e.g. 24 (lossless only) + sample_rate: int | None = None # e.g. 96000 (lossless only) + bitrate_kbps: int | None = None # e.g. 320 (lossy only) + +@dataclass(frozen=True) +class Candidate: + source: str # "qobuz" | "soulseek" | "youtube" + source_ref: str # opaque handle passed back to download() + matched_artist: str + matched_album: str + quality: Quality + track_count: int + confidence: float = 0.0 # 0..1, filled by score_confidence() + +@dataclass(frozen=True) +class DownloadResult: + ok: bool + path: str | None = None + track_count: int = 0 + error: str | None = None + +# lyra_worker/quality.py (Task 1): quality_class(Quality)->int ; source_tier(str)->int ; rank_key(Candidate)->tuple +# lyra_worker/confidence.py (Task 2): score_confidence(target: MBTarget, c: Candidate) -> float +# lyra_worker/ranker.py (Task 3): rank_candidates(target, candidates, min_confidence=0.7) -> list[Candidate] +# lyra_worker/adapters/base.py (Task 4): class SourceAdapter(Protocol): name; tier; health()->bool; +# search(target)->list[Candidate]; download(candidate, dest, on_progress)->DownloadResult +# lyra_worker/adapters/fakes.py (Task 4): FakeQobuz, FakeSoulseek, FakeYouTube, FailingAdapter +# lyra_worker/pipeline.py (Task 6, rewritten): run_pipeline(conn, job_id, adapters, min_confidence=0.7, dest_root="/music") -> None +``` + +--- + +### Task 1: Domain types + quality ranking + +**Files:** +- Create: `worker/lyra_worker/types.py` +- Create: `worker/lyra_worker/quality.py` +- Test: `worker/tests/test_quality.py` + +**Interfaces:** +- Consumes: nothing. +- Produces: the dataclasses `MBTarget`, `Quality`, `Candidate`, `DownloadResult` (see Shared interfaces), and in `quality.py`: `quality_class(q: Quality) -> int`, `source_tier(source: str) -> int`, `rank_key(c: Candidate) -> tuple[int, int]` where a **larger** tuple is better (`(quality_class, -source_tier)`). + +- [ ] **Step 1: Write the domain types** + +`worker/lyra_worker/types.py`: +```python +from dataclasses import dataclass + + +@dataclass(frozen=True) +class MBTarget: + artist: str + album: str + track_count: int | None = None + total_duration_s: int | None = None + + +@dataclass(frozen=True) +class Quality: + fmt: str + lossless: bool + bit_depth: int | None = None + sample_rate: int | None = None + bitrate_kbps: int | None = None + + +@dataclass(frozen=True) +class Candidate: + source: str + source_ref: str + matched_artist: str + matched_album: str + quality: Quality + track_count: int + confidence: float = 0.0 + + +@dataclass(frozen=True) +class DownloadResult: + ok: bool + path: str | None = None + track_count: int = 0 + error: str | None = None +``` + +- [ ] **Step 2: Write the failing test** + +`worker/tests/test_quality.py`: +```python +from lyra_worker.quality import quality_class, source_tier, rank_key +from lyra_worker.types import Candidate, Quality + +HIRES = Quality(fmt="FLAC", lossless=True, bit_depth=24, sample_rate=96000) +CD = Quality(fmt="FLAC", lossless=True, bit_depth=16, sample_rate=44100) +MP3 = Quality(fmt="MP3", lossless=False, bitrate_kbps=320) + + +def _cand(source, quality): + return Candidate(source=source, source_ref="x", matched_artist="A", + matched_album="B", quality=quality, track_count=10) + + +def test_quality_class_tiers(): + assert quality_class(HIRES) == 3 + assert quality_class(CD) == 2 + assert quality_class(MP3) == 1 + + +def test_source_tier_order(): + assert source_tier("qobuz") == 0 + assert source_tier("soulseek") == 1 + assert source_tier("youtube") == 2 + + +def test_design_ranking_order(): + # Qobuz hi-res > Qobuz lossless > Soulseek FLAC > Soulseek MP3 > YouTube + ordered = [ + _cand("qobuz", HIRES), + _cand("qobuz", CD), + _cand("soulseek", CD), + _cand("soulseek", MP3), + _cand("youtube", MP3), + ] + ranks = [rank_key(c) for c in ordered] + assert ranks == sorted(ranks, reverse=True) + assert len(set(ranks)) == 5 # strict total order, no ties +``` + +- [ ] **Step 3: Run the test to verify it fails** + +Run: +```bash +cd worker +worker/.venv/bin/python -m pytest tests/test_quality.py -v +``` +Expected: FAIL — `ModuleNotFoundError: No module named 'lyra_worker.quality'`. + +- [ ] **Step 4: Implement quality ranking** + +`worker/lyra_worker/quality.py`: +```python +from lyra_worker.types import Candidate, Quality + +_SOURCE_TIERS = {"qobuz": 0, "soulseek": 1, "youtube": 2} + + +def quality_class(q: Quality) -> int: + """3 = hi-res lossless, 2 = CD lossless, 1 = lossy.""" + if q.lossless: + hires = (q.bit_depth is not None and q.bit_depth > 16) or ( + q.sample_rate is not None and q.sample_rate > 44100 + ) + return 3 if hires else 2 + return 1 + + +def source_tier(source: str) -> int: + """Lower is better. Unknown sources rank worst.""" + return _SOURCE_TIERS.get(source, 99) + + +def rank_key(c: Candidate) -> tuple[int, int]: + """Sort key; larger is better. Quality class first, then source tier.""" + return (quality_class(c.quality), -source_tier(c.source)) +``` + +- [ ] **Step 5: Run the test to verify it passes** + +Run: +```bash +cd worker +worker/.venv/bin/python -m pytest tests/test_quality.py -v +``` +Expected: PASS — all three tests green. + +- [ ] **Step 6: Commit** + +```bash +git add worker/lyra_worker/types.py worker/lyra_worker/quality.py worker/tests/test_quality.py +git commit -m "feat: add domain types and quality ranking" +``` + +--- + +### Task 2: Confidence scoring + +**Files:** +- Create: `worker/lyra_worker/confidence.py` +- Test: `worker/tests/test_confidence.py` + +**Interfaces:** +- Consumes: `MBTarget`, `Candidate` (Task 1). +- Produces: `score_confidence(target: MBTarget, c: Candidate) -> float` returning 0..1. Weighted blend: 0.5·artist_similarity + 0.4·album_similarity + 0.1·track_factor. Similarity via `difflib.SequenceMatcher` on casefolded/stripped strings. `track_factor` = 1.0 when `target.track_count is None` (unknown) or equal; otherwise `max(0, 1 - abs(diff)/max(target.track_count, 1))`. + +- [ ] **Step 1: Write the failing test** + +`worker/tests/test_confidence.py`: +```python +from lyra_worker.confidence import score_confidence +from lyra_worker.types import Candidate, MBTarget, Quality + +Q = Quality(fmt="FLAC", lossless=True, bit_depth=16, sample_rate=44100) + + +def _cand(artist, album, tracks=10): + return Candidate(source="qobuz", source_ref="x", matched_artist=artist, + matched_album=album, quality=Q, track_count=tracks) + + +def test_exact_match_scores_high(): + target = MBTarget(artist="Radiohead", album="In Rainbows", track_count=10) + assert score_confidence(target, _cand("Radiohead", "In Rainbows", 10)) > 0.95 + + +def test_case_and_whitespace_insensitive(): + target = MBTarget(artist="Radiohead", album="In Rainbows") + assert score_confidence(target, _cand(" radiohead ", "IN RAINBOWS")) > 0.95 + + +def test_wrong_album_scores_low(): + target = MBTarget(artist="Radiohead", album="In Rainbows", track_count=10) + assert score_confidence(target, _cand("Radiohead", "OK Computer", 12)) < 0.7 + + +def test_unknown_track_count_does_not_penalize(): + target = MBTarget(artist="Radiohead", album="In Rainbows", track_count=None) + # track_count unknown -> track_factor neutral; name match still drives a high score + assert score_confidence(target, _cand("Radiohead", "In Rainbows", 99)) > 0.95 + + +def test_track_count_mismatch_lowers_score(): + target = MBTarget(artist="Radiohead", album="In Rainbows", track_count=10) + good = score_confidence(target, _cand("Radiohead", "In Rainbows", 10)) + bad = score_confidence(target, _cand("Radiohead", "In Rainbows", 4)) + assert bad < good +``` + +- [ ] **Step 2: Run the test to verify it fails** + +Run: +```bash +cd worker +worker/.venv/bin/python -m pytest tests/test_confidence.py -v +``` +Expected: FAIL — `ModuleNotFoundError: No module named 'lyra_worker.confidence'`. + +- [ ] **Step 3: Implement confidence scoring** + +`worker/lyra_worker/confidence.py`: +```python +from difflib import SequenceMatcher + +from lyra_worker.types import Candidate, MBTarget + + +def _sim(a: str, b: str) -> float: + return SequenceMatcher(None, a.strip().casefold(), b.strip().casefold()).ratio() + + +def score_confidence(target: MBTarget, c: Candidate) -> float: + """0..1 confidence that candidate c is the requested target.""" + artist_sim = _sim(target.artist, c.matched_artist) + album_sim = _sim(target.album, c.matched_album) + + if target.track_count is None or target.track_count == c.track_count: + track_factor = 1.0 + else: + diff = abs(target.track_count - c.track_count) + track_factor = max(0.0, 1.0 - diff / max(target.track_count, 1)) + + return 0.5 * artist_sim + 0.4 * album_sim + 0.1 * track_factor +``` + +- [ ] **Step 4: Run the test to verify it passes** + +Run: +```bash +cd worker +worker/.venv/bin/python -m pytest tests/test_confidence.py -v +``` +Expected: PASS — all five tests green. + +- [ ] **Step 5: Commit** + +```bash +git add worker/lyra_worker/confidence.py worker/tests/test_confidence.py +git commit -m "feat: add confidence scoring" +``` + +--- + +### Task 3: Ranker + +**Files:** +- Create: `worker/lyra_worker/ranker.py` +- Test: `worker/tests/test_ranker.py` + +**Interfaces:** +- Consumes: `MBTarget`, `Candidate` (Task 1), `score_confidence` (Task 2), `rank_key` (Task 1). +- Produces: `rank_candidates(target: MBTarget, candidates: list[Candidate], min_confidence: float = 0.7) -> list[Candidate]`. It (1) scores each candidate's confidence (returning **new** Candidate instances with `confidence` set), (2) drops any below `min_confidence`, (3) returns the survivors sorted best-first by `(rank_key desc, confidence desc)`. Empty input or all-below-threshold → `[]`. + +- [ ] **Step 1: Write the failing test** + +`worker/tests/test_ranker.py`: +```python +from lyra_worker.ranker import rank_candidates +from lyra_worker.types import Candidate, MBTarget, Quality + +HIRES = Quality(fmt="FLAC", lossless=True, bit_depth=24, sample_rate=96000) +CD = Quality(fmt="FLAC", lossless=True, bit_depth=16, sample_rate=44100) +MP3 = Quality(fmt="MP3", lossless=False, bitrate_kbps=320) +TARGET = MBTarget(artist="Radiohead", album="In Rainbows", track_count=10) + + +def _cand(source, quality, artist="Radiohead", album="In Rainbows", tracks=10): + return Candidate(source=source, source_ref=f"{source}:{quality.fmt}", + matched_artist=artist, matched_album=album, + quality=quality, track_count=tracks) + + +def test_returns_empty_for_no_candidates(): + assert rank_candidates(TARGET, []) == [] + + +def test_best_quality_wins_among_correct_matches(): + ranked = rank_candidates(TARGET, [ + _cand("youtube", MP3), + _cand("qobuz", HIRES), + _cand("soulseek", CD), + ]) + assert [c.source for c in ranked] == ["qobuz", "soulseek", "youtube"] + assert all(c.confidence > 0.9 for c in ranked) + + +def test_confidence_gate_drops_wrong_album_even_if_higher_quality(): + ranked = rank_candidates(TARGET, [ + _cand("qobuz", HIRES, album="OK Computer", tracks=12), # wrong release, hi-res + _cand("soulseek", CD), # correct, lossless + ]) + assert [c.source for c in ranked] == ["soulseek"] + + +def test_all_below_threshold_returns_empty(): + ranked = rank_candidates(TARGET, [ + _cand("qobuz", HIRES, artist="Nirvana", album="Nevermind"), + ]) + assert ranked == [] +``` + +- [ ] **Step 2: Run the test to verify it fails** + +Run: +```bash +cd worker +worker/.venv/bin/python -m pytest tests/test_ranker.py -v +``` +Expected: FAIL — `ModuleNotFoundError: No module named 'lyra_worker.ranker'`. + +- [ ] **Step 3: Implement the ranker** + +`worker/lyra_worker/ranker.py`: +```python +from dataclasses import replace + +from lyra_worker.confidence import score_confidence +from lyra_worker.quality import rank_key +from lyra_worker.types import Candidate, MBTarget + + +def rank_candidates( + target: MBTarget, candidates: list[Candidate], min_confidence: float = 0.7 +) -> list[Candidate]: + """Score, gate by confidence, and sort best-first (quality, then confidence).""" + scored = [replace(c, confidence=score_confidence(target, c)) for c in candidates] + kept = [c for c in scored if c.confidence >= min_confidence] + kept.sort(key=lambda c: (rank_key(c), c.confidence), reverse=True) + return kept +``` + +- [ ] **Step 4: Run the test to verify it passes** + +Run: +```bash +cd worker +worker/.venv/bin/python -m pytest tests/test_ranker.py -v +``` +Expected: PASS — all four tests green. + +- [ ] **Step 5: Commit** + +```bash +git add worker/lyra_worker/ranker.py worker/tests/test_ranker.py +git commit -m "feat: add candidate ranker with confidence gate" +``` + +--- + +### Task 4: SourceAdapter contract + fake adapters + +**Files:** +- Create: `worker/lyra_worker/adapters/__init__.py` +- Create: `worker/lyra_worker/adapters/base.py` +- Create: `worker/lyra_worker/adapters/fakes.py` +- Test: `worker/tests/test_adapters.py` + +**Interfaces:** +- Consumes: `MBTarget`, `Candidate`, `DownloadResult`, `Quality` (Task 1). +- Produces: + - `base.py`: a `SourceAdapter` `Protocol` with attributes `name: str`, `tier: int` and methods `health() -> bool`, `search(target: MBTarget) -> list[Candidate]`, `download(candidate: Candidate, dest: str, on_progress: Callable[[float], None]) -> DownloadResult`. + - `fakes.py`: `FakeQobuz` (returns one hi-res candidate, download succeeds), `FakeSoulseek` (returns one CD-lossless candidate, download succeeds), `FakeYouTube` (returns one lossy candidate, download succeeds), `FailingAdapter` (search returns one candidate but download always fails). Each fake's `search` echoes the target's artist/album into the candidate so confidence is high; each accepts constructor kwargs to override behavior for tests (e.g. `FakeQobuz(matches=False)` returns `[]`). + +- [ ] **Step 1: Write the contract + fakes test** + +`worker/tests/test_adapters.py`: +```python +from lyra_worker.adapters.fakes import FakeQobuz, FakeSoulseek, FakeYouTube, FailingAdapter +from lyra_worker.types import MBTarget + +TARGET = MBTarget(artist="Radiohead", album="In Rainbows", track_count=10) +ALL_FAKES = [FakeQobuz(), FakeSoulseek(), FakeYouTube(), FailingAdapter()] + + +def _noop(_pct: float) -> None: + pass + + +def test_adapters_expose_name_and_tier(): + assert FakeQobuz().name == "qobuz" and FakeQobuz().tier == 0 + assert FakeSoulseek().name == "soulseek" and FakeSoulseek().tier == 1 + assert FakeYouTube().name == "youtube" and FakeYouTube().tier == 2 + + +def test_search_returns_wellformed_candidates(): + for adapter in ALL_FAKES: + for c in adapter.search(TARGET): + assert c.source == adapter.name + assert c.matched_artist and c.matched_album + assert c.track_count > 0 + assert c.source_ref + + +def test_matches_false_yields_no_candidates(): + assert FakeQobuz(matches=False).search(TARGET) == [] + + +def test_successful_download_reports_progress_and_ok(): + seen = [] + result = FakeQobuz().download(FakeQobuz().search(TARGET)[0], "/tmp/x", seen.append) + assert result.ok is True + assert result.track_count > 0 + assert seen and seen[-1] == 1.0 # progress reaches 100% + + +def test_failing_adapter_download_fails(): + fa = FailingAdapter() + result = fa.download(fa.search(TARGET)[0], "/tmp/x", _noop) + assert result.ok is False + assert result.error +``` + +- [ ] **Step 2: Run the test to verify it fails** + +Run: +```bash +cd worker +worker/.venv/bin/python -m pytest tests/test_adapters.py -v +``` +Expected: FAIL — `ModuleNotFoundError: No module named 'lyra_worker.adapters'`. + +- [ ] **Step 3: Implement the contract** + +`worker/lyra_worker/adapters/__init__.py`: +```python +``` + +`worker/lyra_worker/adapters/base.py`: +```python +from typing import Callable, Protocol, runtime_checkable + +from lyra_worker.types import Candidate, DownloadResult, MBTarget + + +@runtime_checkable +class SourceAdapter(Protocol): + name: str + tier: int + + def health(self) -> bool: + """Is this source configured and reachable?""" + ... + + def search(self, target: MBTarget) -> list[Candidate]: + """Return zero or more candidates matching target. No downloading.""" + ... + + def download( + self, candidate: Candidate, dest: str, on_progress: Callable[[float], None] + ) -> DownloadResult: + """Fetch candidate into dest, reporting progress 0.0..1.0.""" + ... +``` + +- [ ] **Step 4: Implement the fakes** + +`worker/lyra_worker/adapters/fakes.py`: +```python +from typing import Callable + +from lyra_worker.types import Candidate, DownloadResult, MBTarget, Quality + +_QUALITIES = { + "qobuz": Quality(fmt="FLAC", lossless=True, bit_depth=24, sample_rate=96000), + "soulseek": Quality(fmt="FLAC", lossless=True, bit_depth=16, sample_rate=44100), + "youtube": Quality(fmt="OPUS", lossless=False, bitrate_kbps=160), +} + + +class _BaseFake: + name = "fake" + tier = 99 + + def __init__(self, matches: bool = True): + self._matches = matches + + def health(self) -> bool: + return True + + def search(self, target: MBTarget) -> list[Candidate]: + if not self._matches: + return [] + tracks = target.track_count or 10 + return [ + Candidate( + source=self.name, + source_ref=f"{self.name}:{target.artist}:{target.album}", + matched_artist=target.artist, + matched_album=target.album, + quality=_QUALITIES.get(self.name, _QUALITIES["youtube"]), + track_count=tracks, + ) + ] + + def download( + self, candidate: Candidate, dest: str, on_progress: Callable[[float], None] + ) -> DownloadResult: + for pct in (0.25, 0.5, 0.75, 1.0): + on_progress(pct) + return DownloadResult(ok=True, path=dest, track_count=candidate.track_count) + + +class FakeQobuz(_BaseFake): + name = "qobuz" + tier = 0 + + +class FakeSoulseek(_BaseFake): + name = "soulseek" + tier = 1 + + +class FakeYouTube(_BaseFake): + name = "youtube" + tier = 2 + + +class FailingAdapter(_BaseFake): + name = "qobuz" # pretends to be a real source that then fails to download + tier = 0 + + def download( + self, candidate: Candidate, dest: str, on_progress: Callable[[float], None] + ) -> DownloadResult: + return DownloadResult(ok=False, error="simulated download failure") +``` + +- [ ] **Step 5: Run the test to verify it passes** + +Run: +```bash +cd worker +worker/.venv/bin/python -m pytest tests/test_adapters.py -v +``` +Expected: PASS — all five tests green. + +- [ ] **Step 6: Commit** + +```bash +git add worker/lyra_worker/adapters worker/tests/test_adapters.py +git commit -m "feat: add SourceAdapter contract and fake adapters" +``` + +--- + +### Task 5: Candidate + LibraryItem schema + +**Files:** +- Modify: `web/prisma/schema.prisma` +- Test: (migration verification — see steps) + +**Interfaces:** +- Consumes: the existing `Job`/`Request` models (Plan 1). +- Produces: + - `Candidate` table: `id` (cuid), `jobId` (FK → Job, onDelete Cascade), `source`, `format`, `qualityClass` (Int), `trackCount` (Int), `confidence` (Float), `sourceRef`, `chosen` (Boolean default false), `createdAt`. + - `LibraryItem` table: `id` (cuid), `requestId` (FK → Request, onDelete Cascade), `artist`, `album`, `path`, `source`, `format`, `qualityClass` (Int), `importedAt` (default now). Unique on (`artist`, `album`) so Import can dedupe. + +- [ ] **Step 1: Add the models to the Prisma schema** + +Append to `web/prisma/schema.prisma`: +```prisma +model Candidate { + id String @id @default(cuid()) + job Job @relation(fields: [jobId], references: [id], onDelete: Cascade) + jobId String + source String + format String + qualityClass Int + trackCount Int + confidence Float + sourceRef String + chosen Boolean @default(false) + createdAt DateTime @default(now()) +} + +model LibraryItem { + id String @id @default(cuid()) + request Request @relation(fields: [requestId], references: [id], onDelete: Cascade) + requestId String + artist String + album String + path String + source String + format String + qualityClass Int + importedAt DateTime @default(now()) + + @@unique([artist, album]) +} +``` + +Add the back-relations to the existing models (Prisma requires both sides). In `model Job { ... }` add: +```prisma + candidates Candidate[] +``` +In `model Request { ... }` add: +```prisma + libraryItem LibraryItem? +``` + +- [ ] **Step 2: Create and apply the migration** + +Run (Postgres up via `docker compose up -d db`): +```bash +cd web +export DATABASE_URL=postgresql://lyra:lyra@localhost:5432/lyra +npx prisma migrate dev --name add_candidate_library_item +``` +Expected: a new migration under `web/prisma/migrations/` is created and applied; "Your database is now in sync with your schema." + +- [ ] **Step 3: Verify the tables exist** + +Run: +```bash +cd web +export DATABASE_URL=postgresql://lyra:lyra@localhost:5432/lyra +npx prisma db execute --stdin <<'SQL' +SELECT 1 FROM "Candidate" LIMIT 0; +SELECT 1 FROM "LibraryItem" LIMIT 0; +SQL +echo "tables OK" +``` +Expected: no error; prints `tables OK`. + +- [ ] **Step 4: Commit** + +```bash +git add web/prisma/schema.prisma web/prisma/migrations +git commit -m "feat: add Candidate and LibraryItem schema" +``` + +--- + +### Task 6: Real staged pipeline + +**Files:** +- Modify (rewrite): `worker/lyra_worker/pipeline.py` +- Test: `worker/tests/test_pipeline.py` (replace the Plan 1 fake-pipeline test) + +**Interfaces:** +- Consumes: `claim_next` behavior (job already in `matching`), `MBTarget`/`Candidate`/`DownloadResult` (Task 1), `rank_candidates` (Task 3), the `SourceAdapter` fakes (Task 4), the `Candidate`/`LibraryItem` tables (Task 5). +- Produces: `run_pipeline(conn, job_id, adapters: list[SourceAdapter], min_confidence: float = 0.7, dest_root: str = "/music") -> None`. Stages: + 1. **intake** — read the job's Request row; build `MBTarget(artist, album)` (track_count None for now). Set stage `intake`. If the album already exists in `LibraryItem` (same artist+album), skip to import-complete without re-downloading (dedupe): set job `imported`, Request `completed`, return. + 2. **match** — set state `matching`/stage `match`; call every adapter's `search(target)`; collect candidates. Persist each as a `Candidate` row (chosen=false), storing `qualityClass` via `quality_class`. + 3. **rank** — set state `matched`/stage `rank`; `rank_candidates(target, all_candidates, min_confidence)`. If empty → `_fail(job, "no candidate above confidence threshold")` (job `needs_attention`, Request `needs_attention`), return. + 4. **download** — set state `downloading`/stage `download`; iterate ranked candidates; for each, find the adapter whose `name == candidate.source` and call `download(...)`; on `ok` mark that `Candidate` row `chosen=true` and proceed; on failure try the next (fall-through). If all fail → `_fail(job, "all downloads failed")`, return. + 5. **tag** — set state `tagging`/stage `tag`; verify `result.track_count == winning_candidate.track_count` (integrity); mismatch → `_fail`, return. (Real tagging is Plan 4.) + 6. **import** — set state `imported`/stage `import`; insert a `LibraryItem` (artist, album, path=result.path, source, format, qualityClass); set Request `completed`. + +- [ ] **Step 1: Write the failing tests** + +Replace `worker/tests/test_pipeline.py` with: +```python +from lyra_worker.adapters.fakes import FakeQobuz, FakeSoulseek, FakeYouTube, FailingAdapter +from lyra_worker.claim import claim_next +from lyra_worker.pipeline import run_pipeline +from tests.conftest import insert_request + + +def _job_state(conn, job_id): + with conn.cursor() as cur: + cur.execute('SELECT state, "currentStage" FROM "Job" WHERE id = %s', (job_id,)) + return cur.fetchone() + + +def _request_status(conn, job_id): + with conn.cursor() as cur: + cur.execute( + 'SELECT status FROM "Request" WHERE id = (SELECT "requestId" FROM "Job" WHERE id = %s)', + (job_id,), + ) + return cur.fetchone()[0] + + +def test_success_picks_best_source_and_imports(conn): + job_id = insert_request(conn, artist="Radiohead", album="In Rainbows") + claim_next(conn) + run_pipeline(conn, job_id, [FakeYouTube(), FakeSoulseek(), FakeQobuz()], dest_root="/tmp/lib") + + assert _job_state(conn, job_id) == ("imported", "import") + assert _request_status(conn, job_id) == "completed" + with conn.cursor() as cur: + cur.execute('SELECT source FROM "Candidate" WHERE "jobId" = %s AND chosen = true', (job_id,)) + assert cur.fetchone()[0] == "qobuz" # highest quality won + cur.execute('SELECT source, format FROM "LibraryItem" WHERE artist = %s', ("Radiohead",)) + row = cur.fetchone() + assert row[0] == "qobuz" + + +def test_falls_through_when_best_download_fails(conn): + job_id = insert_request(conn, artist="Radiohead", album="In Rainbows") + claim_next(conn) + # FailingAdapter is source 'qobuz' (tier 0) whose download fails; Soulseek should win. + run_pipeline(conn, job_id, [FailingAdapter(), FakeSoulseek()], dest_root="/tmp/lib") + + assert _job_state(conn, job_id) == ("imported", "import") + with conn.cursor() as cur: + cur.execute('SELECT source FROM "Candidate" WHERE "jobId" = %s AND chosen = true', (job_id,)) + assert cur.fetchone()[0] == "soulseek" + + +def test_no_match_goes_to_needs_attention(conn): + job_id = insert_request(conn, artist="Radiohead", album="In Rainbows") + claim_next(conn) + run_pipeline(conn, job_id, [FakeQobuz(matches=False)], dest_root="/tmp/lib") + + assert _job_state(conn, job_id)[0] == "needs_attention" + assert _request_status(conn, job_id) == "needs_attention" + + +def test_all_downloads_fail_goes_to_needs_attention(conn): + job_id = insert_request(conn, artist="Radiohead", album="In Rainbows") + claim_next(conn) + run_pipeline(conn, job_id, [FailingAdapter()], dest_root="/tmp/lib") + + assert _job_state(conn, job_id)[0] == "needs_attention" + + +def test_incomplete_download_goes_to_needs_attention(conn): + from lyra_worker.adapters.fakes import FakeQobuz + from lyra_worker.types import DownloadResult + + class TruncatingQobuz(FakeQobuz): + def download(self, candidate, dest, on_progress): + on_progress(1.0) + # reports fewer tracks than the candidate promised -> integrity check fails + return DownloadResult(ok=True, path=dest, track_count=candidate.track_count - 1) + + job_id = insert_request(conn, artist="Radiohead", album="In Rainbows") + claim_next(conn) + run_pipeline(conn, job_id, [TruncatingQobuz()], dest_root="/tmp/lib") + + assert _job_state(conn, job_id)[0] == "needs_attention" + assert _request_status(conn, job_id) == "needs_attention" + + +def test_dedupe_skips_already_in_library(conn): + # First acquisition + job1 = insert_request(conn, artist="Radiohead", album="In Rainbows") + claim_next(conn) + run_pipeline(conn, job1, [FakeQobuz()], dest_root="/tmp/lib") + # Second request for the same album + job2 = insert_request(conn, artist="Radiohead", album="In Rainbows") + claim_next(conn) + run_pipeline(conn, job2, [FakeQobuz()], dest_root="/tmp/lib") + + assert _job_state(conn, job2) == ("imported", "import") + with conn.cursor() as cur: + # dedupe path creates no candidates for the second job + cur.execute('SELECT count(*) FROM "Candidate" WHERE "jobId" = %s', (job2,)) + assert cur.fetchone()[0] == 0 + cur.execute('SELECT count(*) FROM "LibraryItem" WHERE artist = %s AND album = %s', + ("Radiohead", "In Rainbows")) + assert cur.fetchone()[0] == 1 # not duplicated +``` + +- [ ] **Step 2: Run the tests to verify they fail** + +Run: +```bash +cd worker +worker/.venv/bin/python -m pytest tests/test_pipeline.py -v +``` +Expected: FAIL — `run_pipeline` now takes an `adapters` argument the old implementation lacks (TypeError), or import/logic errors. This confirms the tests exercise the new signature. + +- [ ] **Step 3: Rewrite the pipeline** + +Replace `worker/lyra_worker/pipeline.py` with: +```python +from typing import Sequence + +import psycopg + +from lyra_worker.adapters.base import SourceAdapter +from lyra_worker.quality import quality_class +from lyra_worker.ranker import rank_candidates +from lyra_worker.types import Candidate, MBTarget + + +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 _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 _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: + with conn.cursor() as cur: + cur.execute( + 'UPDATE "Candidate" SET chosen = true WHERE "jobId" = %s AND "sourceRef" = %s', + (job_id, source_ref), + ) + conn.commit() + + +def _import(conn: psycopg.Connection, job_id: str, target: MBTarget, + winner: Candidate, path: str) -> None: + request_id = _request_id(conn, job_id) + with conn.cursor() as cur: + cur.execute( + 'INSERT INTO "LibraryItem" (id, "requestId", artist, album, path, source, ' + 'format, "qualityClass", "importedAt") ' + "VALUES (gen_random_uuid()::text, %s, %s, %s, %s, %s, %s, %s, now()) " + 'ON CONFLICT (artist, album) DO NOTHING', + (request_id, target.artist, target.album, path, winner.source, + winner.quality.fmt, quality_class(winner.quality)), + ) + 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], + min_confidence: float = 0.7, + dest_root: str = "/music", +) -> None: + """Real staged acquisition using source-agnostic adapters. Fakes in this plan.""" + # 1. intake + _set_state(conn, job_id, "matching", "intake") + target = _load_target(conn, job_id) + if _already_in_library(conn, target): + _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 + + # 2. match + _set_state(conn, job_id, "matching", "match") + candidates: list[Candidate] = [] + for adapter in adapters: + candidates.extend(adapter.search(target)) + + # 3. rank + _set_state(conn, job_id, "matched", "rank") + ranked = rank_candidates(target, candidates, min_confidence) + _persist_candidates(conn, job_id, ranked) + if not ranked: + _fail(conn, job_id, "no candidate above confidence threshold") + return + + # 4. download (fall-through) + _set_state(conn, job_id, "downloading", "download") + by_source = {a.name: a for a in adapters} + dest = f"{dest_root}/{target.artist}/{target.album}" + winner = None + result = None + for candidate in ranked: + adapter = by_source.get(candidate.source) + if adapter is None: + continue + result = adapter.download(candidate, dest, lambda _pct: None) + if result.ok: + winner = candidate + _mark_chosen(conn, job_id, candidate.source_ref) + break + if winner is None or result is None or not result.ok: + _fail(conn, job_id, "all downloads failed") + return + + # 5. tag (integrity check; real tagging is Plan 4) + _set_state(conn, job_id, "tagging", "tag") + if result.track_count != winner.track_count: + _fail(conn, job_id, "incomplete download") + return + + # 6. import + _set_state(conn, job_id, "imported", "import") + _import(conn, job_id, target, winner, result.path or dest) +``` + +- [ ] **Step 4: Run the tests to verify they pass** + +Run: +```bash +cd worker +worker/.venv/bin/python -m pytest tests/test_pipeline.py -v +``` +Expected: PASS — all five tests green. + +- [ ] **Step 5: Run the full worker suite** + +Run: +```bash +cd worker +export DATABASE_URL=postgresql://lyra:lyra@localhost:5432/lyra +worker/.venv/bin/python -m pytest -v +``` +Expected: PASS — quality, confidence, ranker, adapters, claim, db, and pipeline suites all green, output pristine. + +- [ ] **Step 6: Commit** + +```bash +git add worker/lyra_worker/pipeline.py worker/tests/test_pipeline.py +git commit -m "feat: replace fake pipeline with staged adapter-driven pipeline" +``` + +--- + +### Task 7: Wire adapters into the worker loop + end-to-end + +**Files:** +- Create: `worker/lyra_worker/registry.py` +- Modify: `worker/lyra_worker/main.py` +- Test: `worker/tests/test_registry.py` + +**Interfaces:** +- Consumes: the fakes (Task 4), `run_pipeline` new signature (Task 6). +- Produces: + - `registry.py`: `build_adapters() -> list[SourceAdapter]` returning the currently-enabled adapters. For this plan it returns `[FakeQobuz(), FakeSoulseek(), FakeYouTube()]` (real adapters replace these in Plan 3). A module-level comment marks this as the swap point. + - `main.py`: `run_forever()` builds adapters once via `build_adapters()` and passes them to `run_pipeline(conn, job_id, adapters)`. The `STAGE_DELAY` constant and its use are removed (the real pipeline has real stages; no artificial delay). + +- [ ] **Step 1: Write the failing test** + +`worker/tests/test_registry.py`: +```python +from lyra_worker.registry import build_adapters + + +def test_registry_returns_the_three_fake_sources(): + adapters = build_adapters() + names = sorted(a.name for a in adapters) + assert names == ["qobuz", "soulseek", "youtube"] + for a in adapters: + assert a.health() is True +``` + +- [ ] **Step 2: Run the test to verify it fails** + +Run: +```bash +cd worker +worker/.venv/bin/python -m pytest tests/test_registry.py -v +``` +Expected: FAIL — `ModuleNotFoundError: No module named 'lyra_worker.registry'`. + +- [ ] **Step 3: Implement the registry** + +`worker/lyra_worker/registry.py`: +```python +from lyra_worker.adapters.base import SourceAdapter +from lyra_worker.adapters.fakes import FakeQobuz, FakeSoulseek, FakeYouTube + + +def build_adapters() -> list[SourceAdapter]: + """The enabled source adapters, best-tier first. + + SWAP POINT: Plan 3 replaces these fakes with real Qobuz/Soulseek/YouTube + adapters (reading credentials from config). Nothing else in the pipeline changes. + """ + return [FakeQobuz(), FakeSoulseek(), FakeYouTube()] +``` + +- [ ] **Step 4: Run the test to verify it passes** + +Run: +```bash +cd worker +worker/.venv/bin/python -m pytest tests/test_registry.py -v +``` +Expected: PASS. + +- [ ] **Step 5: Update the worker loop** + +Replace `worker/lyra_worker/main.py` with: +```python +import time + +from lyra_worker.claim import claim_next +from lyra_worker.db import wait_for_db +from lyra_worker.pipeline import run_pipeline +from lyra_worker.registry import build_adapters + +IDLE_SLEEP = 2.0 + + +def run_forever() -> None: + conn = wait_for_db() + adapters = build_adapters() + print("worker: waiting for jobs", flush=True) + try: + while True: + 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) + print(f"worker: finished job {job_id}", flush=True) + finally: + conn.close() + + +if __name__ == "__main__": + run_forever() +``` + +- [ ] **Step 6: Run the full worker suite** + +Run: +```bash +cd worker +export DATABASE_URL=postgresql://lyra:lyra@localhost:5432/lyra +worker/.venv/bin/python -m pytest -v +``` +Expected: PASS — every suite green, output pristine. + +- [ ] **Step 7: End-to-end through the real pipeline in Docker** + +Run (from repo root; if host port 8770 is busy, this still uses it — the stack maps 8770→3000): +```bash +[ -f .env ] || cp .env.example .env +docker compose build +docker compose up -d +sleep 20 +curl -s -X POST http://localhost:8770/api/requests \ + -H "content-type: application/json" \ + -d '{"artist":"Radiohead","album":"In Rainbows"}' >/dev/null +for i in $(seq 1 15); do + sleep 2 + state=$(curl -s http://localhost:8770/api/requests | python3 -c 'import sys,json;print(json.load(sys.stdin)["requests"][0]["job"]["state"])') + echo "state=$state" + [ "$state" = "imported" ] && break +done +docker compose exec -T db psql -U lyra -d lyra -c 'SELECT source, chosen FROM "Candidate" ORDER BY chosen DESC;' +docker compose down +``` +Expected: the request reaches `imported`; the `Candidate` query shows `qobuz` as the chosen source (highest quality among the three fakes). + +- [ ] **Step 8: Commit** + +```bash +git add worker/lyra_worker/registry.py worker/lyra_worker/main.py worker/tests/test_registry.py +git commit -m "feat: wire adapter registry into worker loop" +``` + +--- + +## Notes for the next plan (Real adapters — Plan 3) + +`registry.build_adapters()` is the single swap point: Plan 3 replaces the three fakes with real adapters (`QobuzAdapter` via streamrip, `SoulseekAdapter` via the slskd REST API, `YouTubeAdapter` via yt-dlp), each implementing the same `SourceAdapter` protocol and reading credentials from a `config` table (added then). The pipeline, ranker, confidence, and schema in this plan stay unchanged. Plan 4 adds real MusicBrainz resolution in `intake` (filling `MBTarget.track_count`/`total_duration_s`) and real tagging/file-move in `tag`/`import`. +```