/
/
/
1"""All logic for metadata retrieval."""
2
3from __future__ import annotations
4
5import asyncio
6import logging
7import os
8import random
9import sqlite3
10import threading
11from collections import OrderedDict
12from time import time
13from typing import TYPE_CHECKING, cast
14from uuid import NAMESPACE_URL, uuid5
15
16import aiohttp
17from music_assistant_models.auth import Scope
18from music_assistant_models.background_task import TaskSchedule
19from music_assistant_models.config_entries import ConfigEntry, ConfigValueOption
20from music_assistant_models.enums import (
21 AlbumType,
22 ConfigEntryType,
23 MediaType,
24 ProviderFeature,
25 ProviderType,
26)
27from music_assistant_models.errors import MediaNotFoundError, MusicAssistantError
28from music_assistant_models.media_items import BrowseFolder
29
30from music_assistant.constants import (
31 CONF_LANGUAGE,
32 DB_TABLE_ALBUM_ARTISTS,
33 DB_TABLE_ALBUMS,
34 DB_TABLE_ARTISTS,
35 DB_TABLE_PLAYLISTS,
36 VERBOSE_LOG_LEVEL,
37)
38from music_assistant.controllers.tasks.context import (
39 report_current_task_failure,
40 update_current_task_progress,
41 update_current_task_progress_from_index,
42 update_current_task_progress_text,
43)
44from music_assistant.helpers.api import api_command
45from music_assistant.helpers.compare import (
46 ALBUM_RETAIL_SUFFIX_KEYS,
47 album_retail_suffix_sql_match,
48)
49from music_assistant.helpers.images import cleanup_thumb_cache
50from music_assistant.helpers.lyrics import extract_lrc_lyrics, normalize_lrc_lyrics
51from music_assistant.helpers.throttle_retry import Throttler
52from music_assistant.helpers.util import try_parse_int
53from music_assistant.models.core_controller import CoreController
54from music_assistant.models.music_provider import MusicProvider
55
56from .constants import (
57 ALBUM_RECONCILIATION_TASK_ID,
58 CONF_ENABLE_ONLINE_METADATA,
59 CONF_ENABLE_RADIO_METADATA_LOOKUP,
60 CONF_PREFER_LOCAL_GENRES,
61 CONF_THUMB_CACHE_MAX_SIZE,
62 DEFAULT_LANGUAGE,
63 DEFAULT_THUMB_CACHE_MAX_SIZE_MB,
64 LOCALES,
65 METADATA_LOOKUP_TASK_ID_PREFIX,
66 METADATA_SCAN_BATCH_SIZE,
67 MISSING_ARTIST_METADATA_SCAN_TASK_ID,
68 PLAYLIST_METADATA_SCAN_TASK_ID,
69 REFRESH_INTERVAL,
70 THUMB_CACHE_CLEANUP_TASK_ID,
71)
72from .enrichment import MetadataEnrichmentMixin
73from .images import ImageProxyMixin
74from .radio import RadioArtworkMixin
75
76if TYPE_CHECKING:
77 from music_assistant_models.config_entries import CoreConfig
78 from music_assistant_models.media_items import (
79 Album,
80 Artist,
81 Audiobook,
82 MediaItemType,
83 Playlist,
84 Podcast,
85 Track,
86 )
87
88 from music_assistant import MusicAssistant
89 from music_assistant.controllers.music.media.base import MediaControllerBase
90 from music_assistant.helpers.json import SerializableType
91 from music_assistant.models.metadata_provider import MetadataProvider
92
93
94class MetaDataController(
95 ImageProxyMixin, RadioArtworkMixin, MetadataEnrichmentMixin, CoreController
96):
97 """Controller that handles metadata retrieval and management for media items."""
98
99 domain: str = "metadata"
100 config: CoreConfig
101
102 def __init__(self, mass: MusicAssistant) -> None:
103 """Initialize class."""
104 super().__init__(mass)
105 self.cache = self.mass.cache
106 self._pref_lang: str | None = None
107 self.manifest.name = "Metadata controller"
108 self.manifest.description = (
109 "Music Assistant's core controller which handles all metadata for music."
110 )
111 self.manifest.icon = "book-information-variant"
112 self._throttler = Throttler(1, 30)
113 # image-id bookkeeping, all bounded by _IMAGE_ID_LRU_MAX and sharing the
114 # same key/id string objects so the combined footprint stays small:
115 # - _image_id_forward: (provider, path) -> image_id memo so serializing a
116 # known image skips the sha256 and the lock entirely. Read lock-free
117 # (single dict lookup is atomic), mutated only while holding the lock.
118 # - _image_id_lru: image_id -> (provider, path). Write-through hot cache
119 # in front of the cache controller so that resolving an image by id
120 # never blocks on sqlite if the URL was generated recently.
121 # - _image_id_persisted: image_id -> epoch of the last persist to the
122 # cache db, so repeat encounters skip the sqlite write.
123 # The lock is needed because compute_image_id() runs from the executor
124 # thread during outbound websocket serialization.
125 self._image_id_forward: dict[tuple[str, str], str] = {}
126 self._image_id_lru: OrderedDict[str, tuple[str, str]] = OrderedDict()
127 self._image_id_persisted: dict[str, float] = {}
128 self._image_id_lock = threading.Lock()
129 # corrupt metadata rows found by the last scan pass, per table, for diagnostics
130 self._corrupt_metadata_rows: dict[str, list[dict[str, str | int]]] = {}
131
132 async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
133 """Return all Config Entries for this core module (if any)."""
134 return (
135 # deliberately without a default_value: only values differing from the entry default
136 # are persisted, so declaring one would make a chosen DEFAULT_LANGUAGE
137 # indistinguishable from "never chosen". The locale property applies it on read.
138 ConfigEntry(
139 key=CONF_LANGUAGE,
140 type=ConfigEntryType.STRING,
141 required=False,
142 options=[ConfigValueOption(key, title=value) for key, value in LOCALES.items()],
143 ),
144 ConfigEntry(
145 key=CONF_ENABLE_ONLINE_METADATA,
146 type=ConfigEntryType.BOOLEAN,
147 required=False,
148 default_value=True,
149 ),
150 ConfigEntry(
151 key=CONF_PREFER_LOCAL_GENRES,
152 type=ConfigEntryType.BOOLEAN,
153 required=False,
154 default_value=False,
155 ),
156 ConfigEntry(
157 key=CONF_ENABLE_RADIO_METADATA_LOOKUP,
158 type=ConfigEntryType.BOOLEAN,
159 required=False,
160 default_value=True,
161 ),
162 ConfigEntry(
163 key=CONF_THUMB_CACHE_MAX_SIZE,
164 type=ConfigEntryType.INTEGER,
165 required=False,
166 default_value=DEFAULT_THUMB_CACHE_MAX_SIZE_MB,
167 range=(50, 5000),
168 ),
169 )
170
171 async def setup(self, config: CoreConfig) -> None:
172 """Async initialize of module."""
173 self.config = config
174 if not self.logger.isEnabledFor(VERBOSE_LOG_LEVEL):
175 # silence PIL logger
176 logging.getLogger("PIL").setLevel(logging.WARNING)
177 # make sure that our directory with collage images exists
178 self._collage_images_dir = os.path.join(self.mass.cache_path, "collage_images")
179 if not await asyncio.to_thread(os.path.exists, self._collage_images_dir):
180 await asyncio.to_thread(os.mkdir, self._collage_images_dir)
181
182 async def post_setup(self) -> None:
183 """Handle logic after all core controllers have been set up."""
184 # canonical opaque-id endpoint, served by both the public webserver
185 # and the streams server (the latter is what player metadata URLs hit)
186 self.mass.streams.register_dynamic_route("/imageproxy/*", self.handle_imageproxy)
187 self.mass.webserver.register_dynamic_route("/imageproxy/*", self.handle_imageproxy)
188 self._register_maintenance_tasks()
189
190 async def close(self) -> None:
191 """Handle logic on server stop."""
192 self.mass.streams.unregister_dynamic_route("/imageproxy/*")
193 self.mass.webserver.unregister_dynamic_route("/imageproxy/*")
194
195 @property
196 def providers(self) -> list[MetadataProvider]:
197 """Return all loaded/running MetadataProviders."""
198 return sorted(
199 cast("list[MetadataProvider]", self.mass.get_providers(ProviderType.METADATA)),
200 key=lambda p: p.priority,
201 )
202
203 @property
204 def preferred_language(self) -> str:
205 """Return preferred language for metadata (as 2 letter language code 'en')."""
206 return self.locale.split("_")[0]
207
208 @property
209 def locale(self) -> str:
210 """Return preferred language for metadata (as full locale code 'en_EN')."""
211 value = self.mass.config.get_raw_core_config_value(
212 self.domain, CONF_LANGUAGE, DEFAULT_LANGUAGE
213 )
214 return str(value)
215
216 @api_command("metadata/set_default_preferred_language", required_scope=Scope.CONFIG_CORE_WRITE)
217 def set_default_preferred_language(self, lang: str) -> None:
218 """
219 Set the default preferred language.
220
221 Reasoning behind this is that the backend can not make a wise choice for the default,
222 so relies on some external source that knows better to set this info, like the frontend
223 or a streaming provider.
224 Can only be set once (by this call or the user).
225 """
226 if self.mass.config.get_raw_core_config_value(self.domain, CONF_LANGUAGE):
227 return # already set
228 self.set_preferred_language(lang)
229
230 @api_command("metadata/set_preferred_language", required_scope=Scope.LIBRARY_MANAGE)
231 def set_preferred_language(self, lang: str) -> None:
232 """
233 Set the preferred language.
234
235 Note that this will not modify any existing metadata,
236 but will be used for future lookups.
237 """
238 # prefer exact match
239 if lang in LOCALES:
240 self.mass.config.set_raw_core_config_value(self.domain, CONF_LANGUAGE, lang)
241 return
242 # try strict matching on either locale code or region
243 lang = lang.lower().replace("-", "_")
244 for locale_code, lang_name in LOCALES.items():
245 if lang in (locale_code.lower(), lang_name.lower()):
246 self.mass.config.set_raw_core_config_value(self.domain, CONF_LANGUAGE, locale_code)
247 return
248 # attempt loose match on language code or region code
249 for lang_part in (lang[:2], lang[:-2]):
250 for locale_code in tuple(LOCALES):
251 language_code, region_code = locale_code.lower().split("_", 1)
252 if lang_part in (language_code, region_code):
253 self.mass.config.set_raw_core_config_value(
254 self.domain, CONF_LANGUAGE, locale_code
255 )
256 return
257 # if we reach this point, we couldn't match the language
258 self.logger.warning("%s is not a valid language", lang)
259
260 @api_command("metadata/update_metadata", required_scope=Scope.LIBRARY_MANAGE)
261 async def update_metadata(
262 self, item: str | MediaItemType, force_refresh: bool = False
263 ) -> MediaItemType:
264 """Get/update extra/enhanced metadata for/on given MediaItem."""
265 async with self.cache.handle_refresh(force_refresh):
266 if isinstance(item, str):
267 retrieved_item = await self.mass.music.get_item_by_uri(item)
268 if isinstance(retrieved_item, BrowseFolder):
269 raise TypeError("Cannot update metadata on a BrowseFolder item.")
270 item = retrieved_item
271
272 if item.provider != "library":
273 # this shouldn't happen but just in case.
274 raise RuntimeError("Metadata can only be updated for library items")
275
276 async with self._throttler:
277 if item.media_type == MediaType.ARTIST:
278 await self._update_artist_metadata(
279 cast("Artist", item), force_refresh=force_refresh
280 )
281 if item.media_type == MediaType.ALBUM:
282 await self._update_album_metadata(cast("Album", item), force_refresh=force_refresh)
283 if item.media_type == MediaType.TRACK:
284 await self._update_track_metadata(cast("Track", item), force_refresh=force_refresh)
285 if item.media_type == MediaType.PLAYLIST:
286 await self._update_playlist_metadata(
287 cast("Playlist", item), force_refresh=force_refresh
288 )
289 if item.media_type == MediaType.AUDIOBOOK:
290 await self._update_audiobook_metadata(
291 cast("Audiobook", item), force_refresh=force_refresh
292 )
293 if item.media_type == MediaType.PODCAST:
294 await self._update_podcast_metadata(
295 cast("Podcast", item), force_refresh=force_refresh
296 )
297 return item
298
299 def schedule_update_metadata(self, item: MediaItemType) -> None:
300 """Schedule metadata update for given MediaItem."""
301 if item.provider != "library":
302 # this shouldn't happen but just in case.
303 return
304 last_refresh = item.metadata.last_refresh or 0
305 needs_update = (time() - last_refresh) > REFRESH_INTERVAL
306 if not needs_update:
307 return
308 assert item.uri is not None
309 task_id = self._get_metadata_lookup_task_id(item.uri)
310 _item = item
311
312 self.mass.tasks.run_background_task(
313 task_id=task_id,
314 name=f"Update metadata for {item.name}",
315 handler=lambda: self.update_metadata(_item),
316 translation_key="update_metadata",
317 translation_args=[item.name],
318 translation_owner=self.translation_owner,
319 metadata={
320 "task_domain": "metadata_lookup",
321 "item_uri": item.uri,
322 },
323 )
324
325 @api_command("metadata/get_track_lyrics", required_scope=Scope.LIBRARY_READ)
326 async def get_track_lyrics(
327 self,
328 track: Track,
329 ) -> tuple[str | None, str | None]:
330 """
331 Get lyrics for given track from metadata providers.
332
333 Returns a tuple of (lyrics, lrc_lyrics) if found.
334 """
335 lyrics, lrc_lyrics = await self._get_track_lyrics(track)
336 # on-demand lookups are not stored in the library db, so normalize on the way out
337 # promoting LRC formatted text stored in the plain lyrics tag
338 return lyrics, normalize_lrc_lyrics(lrc_lyrics or extract_lrc_lyrics(lyrics))
339
340 async def get_diagnostics(self) -> dict[str, SerializableType] | None:
341 """Return diagnostics info for this controller to include in diagnostics reports."""
342 if not self._corrupt_metadata_rows:
343 return None
344 return {"corrupt_metadata_rows": cast("SerializableType", self._corrupt_metadata_rows)}
345
346 async def _get_track_lyrics(
347 self,
348 track: Track,
349 ) -> tuple[str | None, str | None]:
350 """Look up (lyrics, lrc_lyrics) for the given track."""
351 if track.metadata and track.metadata.lyrics:
352 return track.metadata.lyrics, track.metadata.lrc_lyrics
353
354 if track.provider == "library":
355 # try to update metadata first
356 await self._update_track_metadata(track, force_refresh=False)
357 return track.metadata.lyrics, track.metadata.lrc_lyrics
358
359 # prefer lyrics from the track's own provider
360 track_provider = self.mass.get_provider(track.provider, provider_type=MusicProvider)
361 if track_provider and ProviderFeature.LYRICS in track_provider.supported_features:
362 full_track = await self.mass.music.tracks.get_provider_item(
363 track.item_id, track.provider
364 )
365 if full_track.metadata and full_track.metadata.lyrics:
366 return full_track.metadata.lyrics, full_track.metadata.lrc_lyrics
367
368 # fallback to other metadata providers
369 for provider in self.providers:
370 if ProviderFeature.LYRICS not in provider.supported_features:
371 continue
372 try:
373 metadata = await provider.get_track_metadata(track)
374 except Exception as err:
375 # a provider failure must not abort the lookup â skip to the next provider
376 self.logger.warning(
377 "Error fetching lyrics for %s from provider %s: %s",
378 track.name,
379 provider.name,
380 err,
381 exc_info=err if self.logger.isEnabledFor(10) else None,
382 )
383 continue
384 if metadata and (metadata.lyrics or metadata.lrc_lyrics):
385 return metadata.lyrics, metadata.lrc_lyrics
386 return None, None
387
388 def _register_maintenance_tasks(self) -> None:
389 """Register the recurring metadata maintenance background tasks."""
390 # Spread across the full day so instances don't all hit the shared MusicBrainz mirror at once
391 utc_hour, utc_minute = divmod(random.randint(0, 24 * 60 - 1), 60)
392 desired_schedule = TaskSchedule.daily(hour=utc_hour, minute=utc_minute)
393 self.mass.tasks.register_scheduled_task(
394 task_id=MISSING_ARTIST_METADATA_SCAN_TASK_ID,
395 name="Scan missing artist metadata",
396 handler=self._scan_missing_artist_metadata,
397 schedule=desired_schedule,
398 translation_key="scan_missing_artist_metadata",
399 translation_owner=self.translation_owner,
400 metadata={"task_domain": "metadata_missing_artist_metadata_scan"},
401 allow_retry=True,
402 )
403 self.mass.tasks.register_scheduled_task(
404 task_id=PLAYLIST_METADATA_SCAN_TASK_ID,
405 name="Refresh playlist metadata",
406 handler=self._refresh_playlist_metadata_batch,
407 schedule=desired_schedule,
408 translation_key="refresh_playlist_metadata",
409 translation_owner=self.translation_owner,
410 metadata={"task_domain": "metadata_playlist_metadata_scan"},
411 allow_retry=True,
412 )
413 self.mass.tasks.register_scheduled_task(
414 task_id=THUMB_CACHE_CLEANUP_TASK_ID,
415 name="Cleanup thumbnail cache",
416 handler=self._cleanup_thumb_cache,
417 schedule=desired_schedule,
418 translation_key="cleanup_thumbnail_cache",
419 translation_owner=self.translation_owner,
420 metadata={"task_domain": "metadata_thumb_cache_cleanup"},
421 allow_retry=True,
422 )
423 # runs every hour rather than spread across the day: it is bounded to a small
424 # batch of albums per run, so there is no shared-mirror stampede to avoid
425 self.mass.tasks.register_scheduled_task(
426 task_id=ALBUM_RECONCILIATION_TASK_ID,
427 name="Reconcile duplicate albums",
428 handler=self._reconcile_duplicate_albums,
429 schedule=TaskSchedule.hourly(),
430 translation_key="reconcile_duplicate_albums",
431 translation_owner=self.translation_owner,
432 metadata={"task_domain": "metadata_album_reconciliation"},
433 allow_retry=True,
434 )
435
436 @staticmethod
437 def _get_metadata_lookup_task_id(uri: str) -> str:
438 """Return deterministic task id for a metadata lookup."""
439 return f"{METADATA_LOOKUP_TASK_ID_PREFIX}_{uuid5(NAMESPACE_URL, uri).hex}"
440
441 async def _scan_missing_artist_metadata(self) -> None:
442 """Scan for artists with missing metadata."""
443 update_current_task_progress_text("Searching for artists with missing metadata")
444 missing_images = (
445 f"(json_extract({DB_TABLE_ARTISTS}.metadata,'$.images') ISNULL "
446 f"OR json_extract({DB_TABLE_ARTISTS}.metadata,'$.images') = '[]')"
447 )
448 missing_description = f"json_extract({DB_TABLE_ARTISTS}.metadata,'$.description') ISNULL"
449 never_refreshed = f"json_extract({DB_TABLE_ARTISTS}.metadata,'$.last_refresh') ISNULL"
450 query = f"({missing_images} OR {missing_description}) AND {never_refreshed}"
451 artists = await self._get_scan_batch(self.mass.music.artists, DB_TABLE_ARTISTS, query)
452 if not artists:
453 update_current_task_progress_text("No artists with missing metadata found")
454 return
455 for index, artist in enumerate(artists, 1):
456 try:
457 update_current_task_progress_from_index(
458 index,
459 len(artists),
460 f"Refreshing metadata for artist {index}/{len(artists)}: {artist.name}",
461 )
462 await self._update_artist_metadata(artist, force_refresh=False)
463 except Exception as err:
464 report_current_task_failure(f"{artist.name}: {err}")
465 self.logger.warning(
466 "Error while updating artist metadata for %s: %s",
467 artist.name,
468 str(err),
469 exc_info=err if self.logger.isEnabledFor(10) else None,
470 )
471 update_current_task_progress(100, f"Processed {len(artists)} artist(s)")
472
473 async def _refresh_playlist_metadata_batch(self) -> None:
474 """Refresh metadata for a small batch of library playlists."""
475 update_current_task_progress_text("Searching for playlists needing metadata refresh")
476 refresh_before = int(time() - REFRESH_INTERVAL)
477 query = (
478 f"{DB_TABLE_PLAYLISTS}.is_dynamic = 0 AND ("
479 f"json_extract({DB_TABLE_PLAYLISTS}.metadata,'$.last_refresh') ISNULL "
480 f"OR json_extract({DB_TABLE_PLAYLISTS}.metadata,'$.last_refresh') < {refresh_before})"
481 )
482 playlists = await self._get_scan_batch(self.mass.music.playlists, DB_TABLE_PLAYLISTS, query)
483 if not playlists:
484 update_current_task_progress_text("No playlists require metadata refresh")
485 return
486 for index, playlist in enumerate(playlists, 1):
487 try:
488 update_current_task_progress_from_index(
489 index,
490 len(playlists),
491 f"Refreshing playlist metadata {index}/{len(playlists)}: {playlist.name}",
492 )
493 await self._update_playlist_metadata(playlist, force_refresh=False)
494 except Exception as err:
495 report_current_task_failure(f"{playlist.name}: {err}")
496 self.logger.warning(
497 "Error while refreshing playlist metadata for %s: %s",
498 playlist.name,
499 str(err),
500 exc_info=err if self.logger.isEnabledFor(10) else None,
501 )
502 update_current_task_progress(100, f"Processed {len(playlists)} playlist(s)")
503
504 async def _reconcile_duplicate_albums(self) -> None:
505 """Enrich and re-match a small batch of sparse or possibly duplicated albums."""
506 update_current_task_progress_text("Searching for albums needing reconciliation")
507 # candidates keep retrying at the normal REFRESH_INTERVAL cadence (e.g. after a
508 # transient provider outage), rather than only ever once
509 refresh_before = int(time() - REFRESH_INTERVAL)
510 query = (
511 f"({DB_TABLE_ALBUMS}.album_type = '{AlbumType.UNKNOWN.value}' "
512 f"OR {_duplicate_album_sibling_guard()}) AND ("
513 f"json_extract({DB_TABLE_ALBUMS}.metadata,'$.last_refresh') ISNULL "
514 f"OR json_extract({DB_TABLE_ALBUMS}.metadata,'$.last_refresh') < {refresh_before})"
515 )
516 albums = await self._get_scan_batch(self.mass.music.albums, DB_TABLE_ALBUMS, query)
517 if not albums:
518 update_current_task_progress_text("No albums require reconciliation")
519 return
520 for index, album in enumerate(albums, 1):
521 try:
522 update_current_task_progress_from_index(
523 index,
524 len(albums),
525 f"Reconciling album {index}/{len(albums)}: {album.name}",
526 )
527 # enrich sparse provider data (type/year/metadata) first so the follow-up
528 # match has full album details to work with, then re-fetch the now-enriched
529 # library row before re-matching: match_providers merges a confirmed mapping
530 # into an existing duplicate through the safe add_provider_mappings path
531 try:
532 await self._update_album_metadata(album, force_refresh=False)
533 reconciled_album = await self.mass.music.albums.get_library_item(album.item_id)
534 except MediaNotFoundError:
535 # both rows of a duplicate pair can share a batch, so this row may
536 # already have been merged into its duplicate earlier in the run
537 continue
538 await self.mass.music.albums.match_providers(reconciled_album)
539 except (MusicAssistantError, aiohttp.ClientError, TimeoutError) as err:
540 report_current_task_failure(f"{album.name}: {err}")
541 self.logger.warning(
542 "Error while reconciling album %s: %s",
543 album.name,
544 str(err),
545 exc_info=err if self.logger.isEnabledFor(10) else None,
546 )
547 update_current_task_progress(100, f"Processed {len(albums)} album(s)")
548
549 async def _cleanup_thumb_cache(self) -> None:
550 """Remove oldest thumbnails when the cache folder exceeds the configured limit."""
551 max_size_mb = (
552 try_parse_int(
553 self.config.get_value(CONF_THUMB_CACHE_MAX_SIZE), DEFAULT_THUMB_CACHE_MAX_SIZE_MB
554 )
555 or DEFAULT_THUMB_CACHE_MAX_SIZE_MB
556 )
557 removed = await cleanup_thumb_cache(self.mass.cache_path, max_size_mb * 1024 * 1024)
558 if removed:
559 self.logger.debug("Thumbnail cache cleanup: removed %s file(s)", removed)
560
561 async def _get_scan_batch[ItemCls: MediaItemType](
562 self,
563 media_controller: MediaControllerBase[ItemCls],
564 table: str,
565 query: str,
566 ) -> list[ItemCls]:
567 """Fetch a metadata-scan batch, tolerating rows with corrupt metadata JSON."""
568 try:
569 items = await media_controller.get_library_items_by_query(
570 limit=METADATA_SCAN_BATCH_SIZE,
571 order_by="random",
572 extra_query_parts=[query],
573 )
574 except sqlite3.OperationalError as err:
575 if "malformed JSON" not in str(err):
576 raise
577 await self._report_corrupt_metadata_rows(table)
578 return await media_controller.get_library_items_by_query(
579 limit=METADATA_SCAN_BATCH_SIZE,
580 order_by="random",
581 extra_query_parts=[f"{_valid_metadata_guard(table)} AND {query}"],
582 )
583 # a clean scan proves the table currently holds no corrupt rows
584 self._corrupt_metadata_rows.pop(table, None)
585 return items
586
587 async def _report_corrupt_metadata_rows(self, table: str) -> None:
588 """Report library rows whose metadata column holds invalid JSON."""
589 rows = await self.mass.music.database.get_rows_from_query(
590 f"SELECT item_id, name FROM {table} "
591 f"WHERE {table}.metadata IS NOT NULL AND NOT json_valid({table}.metadata)",
592 limit=25,
593 )
594 # keep the findings for the diagnostics report, replacing the previous
595 # pass so repaired rows drop out again
596 if rows:
597 self._corrupt_metadata_rows[table] = [
598 {"item_id": row["item_id"], "name": row["name"]} for row in rows
599 ]
600 else:
601 self._corrupt_metadata_rows.pop(table, None)
602 for row in rows:
603 message = (
604 f"'{row['name']}' has corrupt metadata and was skipped. To repair, remove "
605 f"'{row['name']}' from the library; it will be re-added with fresh metadata "
606 f"on the next library sync ({table} id {row['item_id']})."
607 )
608 report_current_task_failure(message)
609 self.logger.warning(message)
610
611
612def _duplicate_album_sibling_guard() -> str:
613 """Return a query part that selects albums which may be a duplicate of another library row."""
614 shares_artist = (
615 f"EXISTS (SELECT 1 FROM {DB_TABLE_ALBUM_ARTISTS} own "
616 f"JOIN {DB_TABLE_ALBUM_ARTISTS} other ON other.artist_id = own.artist_id "
617 f"WHERE own.album_id = {DB_TABLE_ALBUMS}.item_id AND other.album_id = dup.item_id)"
618 )
619 # a title that normalizes to nothing (e.g. Ed Sheeran's '+', '=' and '÷') matches every
620 # other such title, so those fall back to their raw spelling like the album comparison does
621 same_title = (
622 f"({DB_TABLE_ALBUMS}.search_name != '' OR "
623 f"REPLACE({DB_TABLE_ALBUMS}.name,' ','') = REPLACE(dup.name,' ',''))"
624 )
625 # a provider that spells out the retail suffix stores the album under the plain name
626 # plus that suffix, so the pair is related from either side. The raw title decides
627 # which side spelled it out, so an ordinary title that merely ends in those letters
628 # ("Step") is left alone.
629 # Every alternative stays an equality on dup.search_name, keeping the name index in use.
630 own_name = f"{DB_TABLE_ALBUMS}.search_name"
631 matches_name = [f"dup.search_name = {own_name}"]
632 for suffix in ALBUM_RETAIL_SUFFIX_KEYS:
633 matches_name.append(
634 f"({album_retail_suffix_sql_match('dup.name', suffix)} "
635 f"AND dup.search_name = {own_name} || '{suffix}')"
636 )
637 matches_name.append(
638 f"({album_retail_suffix_sql_match(f'{DB_TABLE_ALBUMS}.name', suffix)} "
639 f"AND dup.search_name = "
640 f"substr({own_name}, 1, length({own_name}) - {len(suffix)}))"
641 )
642 same_name = " OR ".join(matches_name)
643 # deliberately an identity-only pre-filter: which editions may be merged is decided by
644 # the album comparison, which escalates an ambiguous edition to tracklists and
645 # MusicBrainz and rejects a recording-changing one (live, remix, ...) outright
646 return (
647 f"EXISTS (SELECT 1 FROM {DB_TABLE_ALBUMS} dup "
648 f"WHERE dup.item_id != {DB_TABLE_ALBUMS}.item_id "
649 f"AND ({same_name}) "
650 f"AND {same_title} AND {shares_artist})"
651 )
652
653
654def _valid_metadata_guard(table: str) -> str:
655 """Return a query part that excludes rows with invalid JSON in the metadata column."""
656 # sqlite's json functions raise a fatal 'malformed JSON' error on invalid input,
657 # which would fail the entire scan query because of a single corrupt row
658 return f"({table}.metadata IS NULL OR json_valid({table}.metadata))"
659