""" Index-time enrichment for the Music group app: embedded tag/cover extraction (mutagen) and filename-parse fallback for a newly-added audio IndexEntry (docs/musicbay.md §2.1, §6). Runs through its own small bounded worker pool, the same discipline as the Videos app's `enrich.py` — separate from any other pool, never blocking a scan or the watchdog. Unlike `enrich.py`, this one shells out to nothing: `mutagen` is pure Python, synchronous I/O only, so there is no subprocess to spawn, no pipe to drain, and no ffmpeg-shaped deadlock risk here at all — the bounded pool exists to keep a large library's indexing burst bounded, not to contain a process. Reads run via `asyncio.to_thread` so they never block the event loop. MusicBrainz lookups are **not** done here. Tag/cover extraction is free and local, so it runs for every audio file the Music app is enabled for, regardless of whether MusicBrainz itself is turned on for the group — the flat view (docs/musicbay.md §5.2) needs nothing more than this. MusicBrainz is a separate, lazy, per-request enrichment (`music_meta_req`, handled in webrtc_server.py), the same "fetched on demand, cached once" shape TMDB already uses. """ import asyncio import logging import re from collections.abc import Awaitable, Callable from pathlib import Path import blake3 from meshbay_common.protocol import IndexEntry from mutagen import File as MutagenFile from meshbay_node.indexer import title_parse from meshbay_node.media_cache import MediaCache log = logging.getLogger(__name__) # Higher than the video pool's default (2): mutagen reads a few KB of tag # data synchronously, no subprocess, no decode — cheap enough that a wider # pool doesn't cost much and finishes a large library's initial scan sooner. DEFAULT_MAX_CONCURRENT = 4 READ_TIMEOUT_SECS = 10 _TRACK_NO_RE = re.compile(r"\d+") def _extract_cover(mf) -> bytes | None: """ Best-effort embedded cover art across the tag formats mutagen exposes differently: ID3 (MP3) keeps pictures as APIC frames on `.tags`, FLAC exposes `.pictures` on the file object itself, MP4/M4A keeps a `covr` atom on `.tags`. Returns the first picture found, or None — most of a real library has no embedded art at all, which is not an error. """ tags = mf.tags if tags is not None and hasattr(tags, "getall"): pics = tags.getall("APIC") if pics: return bytes(pics[0].data) pictures = getattr(mf, "pictures", None) if pictures: return bytes(pictures[0].data) if tags is not None and hasattr(tags, "get"): covr = tags.get("covr") if covr: return bytes(covr[0]) return None def _read_tags_and_cover(path: Path) -> tuple[dict, float | None, bytes | None]: """ Synchronous — always called via asyncio.to_thread. Returns a partial `tags` dict (only keys actually found: title/artist/album/track_no), duration in seconds (None if unreadable), and raw cover bytes (None if absent). Never raises for an unreadable/corrupt file — the caller falls back to filename parsing entirely in that case. """ tags: dict = {} duration: float | None = None try: easy = MutagenFile(str(path), easy=True) except Exception: easy = None if easy is not None: if easy.info is not None: duration = getattr(easy.info, "length", None) for field in ("title", "artist", "album"): values = easy.get(field) if values and str(values[0]).strip(): tags[field] = str(values[0]).strip() track_raw = easy.get("tracknumber") if track_raw: m = _TRACK_NO_RE.match(str(track_raw[0])) if m: tags["track_no"] = int(m.group()) cover: bytes | None = None try: raw = MutagenFile(str(path)) except Exception: raw = None if raw is not None: try: cover = _extract_cover(raw) except Exception: cover = None return tags, duration, cover def _artist_album_from_ancestors(file_path: Path) -> tuple[str | None, str | None]: """ `Artist/Album/track.mp3` is the common shape (docs/musicbay.md §2.1) — used only to fill whatever the tags left empty. No attempt to validate against the group's actual roots (enrich.py's ancestor walks don't either): a flat `Artist/track.mp3` layout, or a various-artists compilation folder, just yields a plausible-but-not-guaranteed album name from the immediate parent and nothing further up — good enough for a fallback, not asserted as accurate. """ album_folder = file_path.parent if album_folder == album_folder.parent: return None, None artist_folder = album_folder.parent album = album_folder.name or None artist = artist_folder.name if artist_folder != artist_folder.parent else None return artist, album class AudioEnricher: """Owns the node's bounded audio index-time enrichment pool.""" def __init__(self, media_cache: MediaCache, max_concurrent: int = DEFAULT_MAX_CONCURRENT): self._media_cache = media_cache self._sem = asyncio.Semaphore(max_concurrent) self._tasks: set[asyncio.Task] = set() def spawn( self, entry: IndexEntry, file_path: Path, on_done: Callable[[str, dict], Awaitable[None]], ) -> asyncio.Task: """ Fire-and-forget one file's enrichment — same contract as `enrich.Enricher.spawn`: `on_done(file_id, fields)` is awaited with the index fields to merge once ready, never blocks the caller, and the returned task must be held by the caller for the same reason `WebRTCPeerSession._spawn` holds streaming tasks (a bare `ensure_future` can be garbage-collected mid-flight). """ task = asyncio.ensure_future(self._run(entry, file_path, on_done)) self._tasks.add(task) def _cleanup(t: asyncio.Task) -> None: self._tasks.discard(t) if not t.cancelled() and t.exception(): log.error("Audio enrichment failed for %s: %s", entry.id[:12], t.exception(), exc_info=t.exception()) task.add_done_callback(_cleanup) return task async def _run( self, entry: IndexEntry, file_path: Path, on_done: Callable[[str, dict], Awaitable[None]], ) -> None: async with self._sem: fields: dict = {} try: tags, duration, cover = await asyncio.wait_for( asyncio.to_thread(_read_tags_and_cover, file_path), timeout=READ_TIMEOUT_SECS) except Exception as e: log.warning("Tag read failed for %s: %s", file_path, e) tags, duration, cover = {}, None, None if duration: fields["duration"] = int(duration) parsed = title_parse.parse_track_filename(entry.name) fields["display_title"] = tags.get("title") or parsed.title or parsed.naive_title fields["track_no"] = tags.get("track_no") if "track_no" in tags else parsed.track_no artist = tags.get("artist") album = tags.get("album") if not artist or not album: fallback_artist, fallback_album = await asyncio.to_thread( _artist_album_from_ancestors, file_path) artist = artist or fallback_artist album = album or fallback_album fields["artist"] = artist fields["album"] = album if cover: thumb_hash = blake3.blake3(cover).hexdigest() await self._media_cache.put_thumb(thumb_hash, entry.id, cover) fields["thumb_hash"] = thumb_hash await on_done(entry.id, fields)