import os import re import uuid from dataclasses import dataclass import psycopg from lyra_worker.library import list_audio_files from lyra_worker.probe import AudioProbe _AUDIO_EXT = {".flac", ".mp3", ".m4a", ".opus", ".ogg", ".aac", ".wav"} _YEAR = re.compile(r"^(.*) \((\d{4})\)$") @dataclass(frozen=True) class ScanResult: imported: int skipped: int def _parse_album_dir(name: str) -> str: m = _YEAR.match(name) return m.group(1) if m else name def _first_audio(album_path: str) -> str | None: for name in sorted(os.listdir(album_path)): if os.path.splitext(name)[1].lower() in _AUDIO_EXT: return os.path.join(album_path, name) return None def _ensure_artist(conn, target) -> str | None: """Upsert the WatchedArtist (by mbid) and return its id, or None if the target has no artist mbid.""" if not target.artist_mbid: return None with conn.cursor() as cur: cur.execute( 'INSERT INTO "WatchedArtist" (id, mbid, name, "autoMonitorFuture", "monitorFrom", "createdAt") ' "VALUES (gen_random_uuid()::text, %s, %s, false, now(), now()) " 'ON CONFLICT (mbid) DO NOTHING', (target.artist_mbid, target.artist), ) cur.execute('SELECT id FROM "WatchedArtist" WHERE mbid = %s', (target.artist_mbid,)) row = cur.fetchone() conn.commit() return row[0] if row else None def _populate_discography(conn, browser, watched_id, artist_mbid: str, artist_name: str) -> None: """Store the artist's full release-group list as unmonitored rows (like a web follow).""" for rg in browser.browse_release_groups(artist_mbid): with conn.cursor() as cur: cur.execute( 'INSERT INTO "MonitoredRelease" (id, "watchedArtistId", "artistMbid", "artistName", ' '"rgMbid", album, "primaryType", "secondaryTypes", "firstReleaseDate", monitored, state, "createdAt") ' "VALUES (gen_random_uuid()::text, %s, %s, %s, %s, %s, %s, %s, %s, false, 'wanted', now()) " 'ON CONFLICT ("rgMbid") DO UPDATE SET ' ' "primaryType" = EXCLUDED."primaryType", ' ' "secondaryTypes" = EXCLUDED."secondaryTypes", ' ' "firstReleaseDate" = EXCLUDED."firstReleaseDate", ' ' "watchedArtistId" = COALESCE("MonitoredRelease"."watchedArtistId", EXCLUDED."watchedArtistId")', (watched_id, artist_mbid, artist_name, rg.rg_mbid, rg.title, rg.primary_type or None, list(rg.secondary_types), rg.first_release_date or None), ) conn.commit() def _record_owned(conn, target, quality_class: int, fmt: str, path: str, watched_id) -> bool: """Mark the owned release monitored+fulfilled and record the LibraryItem. Returns True if newly imported.""" with conn.cursor() as cur: if target.rg_mbid: cur.execute( 'INSERT INTO "MonitoredRelease" (id, "watchedArtistId", "artistMbid", "artistName", ' '"rgMbid", album, "secondaryTypes", monitored, state, "currentQualityClass", ' '"firstGrabbedAt", "createdAt") ' "VALUES (gen_random_uuid()::text, %s, %s, %s, %s, %s, '{}', true, 'fulfilled', %s, now(), now()) " 'ON CONFLICT ("rgMbid") DO UPDATE SET monitored = true, state = \'fulfilled\', ' ' "currentQualityClass" = EXCLUDED."currentQualityClass", ' ' "firstGrabbedAt" = COALESCE("MonitoredRelease"."firstGrabbedAt", now()), ' ' "watchedArtistId" = COALESCE("MonitoredRelease"."watchedArtistId", EXCLUDED."watchedArtistId")', (watched_id, target.artist_mbid, target.artist, target.rg_mbid, target.album, quality_class), ) track_names = list_audio_files(path) # capture the real on-disk tracks cur.execute( 'INSERT INTO "LibraryItem" (id, "requestId", artist, album, path, source, format, ' '"qualityClass", "trackNames", "importedAt") ' "VALUES (gen_random_uuid()::text, NULL, %s, %s, %s, 'scan', %s, %s, %s, now()) " # refresh trackNames on re-scan (populates items imported before this column); # xmax=0 ⇒ this was an INSERT (a genuinely new library item) vs an UPDATE. 'ON CONFLICT (artist, album) DO UPDATE SET "trackNames" = EXCLUDED."trackNames" ' "RETURNING (xmax = 0) AS inserted", (target.artist, target.album, path, fmt, quality_class, track_names), ) newly = cur.fetchone()[0] # a new LibraryItem means this album wasn't recorded before conn.commit() return newly def _iter_album_entries(dest_root: str): """Yield (artist_name, album_folder, album_path) for each candidate album directory, in the stable nested-sorted order the scan cursor relies on.""" if not os.path.isdir(dest_root): return for artist_name in sorted(os.listdir(dest_root)): artist_path = os.path.join(dest_root, artist_name) if artist_name.startswith(".") or not os.path.isdir(artist_path): continue for album_folder in sorted(os.listdir(artist_path)): album_path = os.path.join(artist_path, album_folder) if album_folder.startswith(".") or not os.path.isdir(album_path): continue yield artist_name, album_folder, album_path def _process_album(conn, resolver, probe, browser, artist_name, album_folder, album_path, browsed) -> str: """Identify and record one album folder. Returns 'imported' (newly recorded), 'recorded' (matched but already present), 'skipped' (no MB match) or 'none' (no audio).""" audio = _first_audio(album_path) if audio is None: return "none" # no audio → not an album pr = probe.probe(audio) album_name = _parse_album_dir(album_folder) target = resolver.resolve(artist_name, album_name) if artist_name and album_name else None if target is None: # fall back to the file's embedded tags if pr.artist and pr.album and (pr.artist, pr.album) != (artist_name, album_name): target = resolver.resolve(pr.artist, pr.album) if target is None: return "skipped" watched_id = _ensure_artist(conn, target) if watched_id is not None and target.artist_mbid not in browsed: _populate_discography(conn, browser, watched_id, target.artist_mbid, target.artist) browsed.add(target.artist_mbid) return "imported" if _record_owned(conn, target, pr.quality_class, pr.fmt, album_path, watched_id) else "recorded" def build_worklist(conn: psycopg.Connection, scan_id: str, dest_root: str = "/music") -> int: """Walk the tree once and persist this scan's album worklist. Idempotent (ON CONFLICT DO NOTHING), so a re-run after a crash adds only missing rows. Returns the number of rows inserted.""" inserted = 0 with conn.cursor() as cur: for artist_name, album_folder, album_path in _iter_album_entries(dest_root): cur.execute( 'INSERT INTO "ScanWorkItem" (id, "scanId", artist, album, path, done) ' "VALUES (gen_random_uuid()::text, %s, %s, %s, %s, false) " 'ON CONFLICT ("scanId", path) DO NOTHING', (scan_id, artist_name, album_folder, album_path), ) inserted += cur.rowcount conn.commit() return inserted def scan_chunk(conn: psycopg.Connection, resolver, probe: AudioProbe, browser, scan_id: str, limit: int) -> tuple[int, int, bool]: """Process up to `limit` not-yet-done worklist items for `scan_id`, in stable (artist, album) order, marking each done. Returns (imported, skipped, done) where `done` is True when no undone rows remain. Idempotent: all record writes are ON CONFLICT, and an item is marked done only after it is processed.""" with conn.cursor() as cur: cur.execute( 'SELECT id, artist, album, path FROM "ScanWorkItem" ' 'WHERE "scanId" = %s AND NOT done ORDER BY artist, album LIMIT %s', (scan_id, limit), ) items = cur.fetchall() imported = skipped = 0 browsed: set[str] = set() # artists whose discography we've populated this chunk for item_id, artist_name, album_folder, album_path in items: outcome = _process_album(conn, resolver, probe, browser, artist_name, album_folder, album_path, browsed) if outcome == "imported": imported += 1 elif outcome == "skipped": skipped += 1 with conn.cursor() as cur: cur.execute('UPDATE "ScanWorkItem" SET done = true WHERE id = %s', (item_id,)) conn.commit() with conn.cursor() as cur: cur.execute('SELECT count(*) FROM "ScanWorkItem" WHERE "scanId" = %s AND NOT done', (scan_id,)) remaining = cur.fetchone()[0] return imported, skipped, remaining == 0 def clear_worklist(conn: psycopg.Connection, scan_id: str) -> None: with conn.cursor() as cur: cur.execute('DELETE FROM "ScanWorkItem" WHERE "scanId" = %s', (scan_id,)) conn.commit() def scan_library(conn: psycopg.Connection, resolver, probe: AudioProbe, browser, dest_root: str = "/music") -> ScanResult: """Walk `{dest_root}/{Artist}/{Album}/` in one unbounded pass (build the worklist, drain it, clean up), recording matched albums as have + monitored + followed plus each followed artist's full discography. Test/one-shot convenience; the worker loop uses the chunked build_worklist + scan_chunk so a large library doesn't block job-claim.""" scan_id = "oneshot-" + uuid.uuid4().hex build_worklist(conn, scan_id, dest_root) imported = skipped = 0 done = False while not done: imp, skp, done = scan_chunk(conn, resolver, probe, browser, scan_id, limit=1000) imported += imp skipped += skp clear_worklist(conn, scan_id) return ScanResult(imported, skipped)