Files
Lyra/docs/superpowers/plans/2026-07-10-lyra-pipeline-core.md
T
Jonathan b62b4145da docs: add pipeline-core implementation plan (slice 1, plan 2)
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-10 18:19:06 +02:00

42 KiB

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)

# 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:

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:

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:

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:

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:

cd worker
worker/.venv/bin/python -m pytest tests/test_quality.py -v

Expected: PASS — all three tests green.

  • Step 6: Commit
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:

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:

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:

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:

cd worker
worker/.venv/bin/python -m pytest tests/test_confidence.py -v

Expected: PASS — all five tests green.

  • Step 5: Commit
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:

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:

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:

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:

cd worker
worker/.venv/bin/python -m pytest tests/test_ranker.py -v

Expected: PASS — all four tests green.

  • Step 5: Commit
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:

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:

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:

worker/lyra_worker/adapters/base.py:

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:

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:

cd worker
worker/.venv/bin/python -m pytest tests/test_adapters.py -v

Expected: PASS — all five tests green.

  • Step 6: Commit
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:

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:

  candidates Candidate[]

In model Request { ... } add:

  libraryItem LibraryItem?
  • Step 2: Create and apply the migration

Run (Postgres up via docker compose up -d db):

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:

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
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:

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:

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:

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:

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:

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
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:

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:

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:

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:

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:

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:

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):

[ -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
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.