/
/
1"""
2Audio streaming helpers that interact with core controllers and providers.
3
4This module contains all audio stream acquisition and processing functions
5that need access to the MusicAssistant instance. Generic audio utilities
6that do not need controller interaction live in helpers/audio.py.
7"""
8
9from __future__ import annotations
10
11import asyncio
12import logging
13import os
14import re
15import time
16from collections.abc import AsyncGenerator, Iterable
17from contextlib import aclosing, asynccontextmanager, nullcontext
18from dataclasses import dataclass
19from functools import partial
20from typing import TYPE_CHECKING, Any, cast
21from urllib.parse import urlparse
22from weakref import WeakValueDictionary
23
24import aiofiles
25import aiofiles.os
26import aiohttp
27import shortuuid
28from aiohttp import ClientConnectorSSLError, ClientResponseError, ClientTimeout
29from music_assistant_models.audio_processing import (
30 AudioDSPDetails,
31 AudioOutputDetails,
32 AudioQueueProcessing,
33)
34from music_assistant_models.dsp import (
35 AudioChannel,
36 ConvolutionFilter,
37 DSPConfig,
38 DSPFilter,
39 DSPState,
40)
41from music_assistant_models.enums import (
42 ContentType,
43 CrossfadeMode,
44 MediaType,
45 PlayerFeature,
46 ProviderFeature,
47 ProviderType,
48 StreamType,
49 VolumeNormalizationMode,
50)
51from music_assistant_models.errors import (
52 AudioError,
53 InvalidDataError,
54 MediaNotFoundError,
55 MusicAssistantError,
56 ProviderPermissionDenied,
57 ProviderUnavailableError,
58 QueueEmpty,
59 RetriesExhausted,
60)
61from music_assistant_models.media_items import Album, AudioFormat, Track
62from music_assistant_models.player_queue import PlayLogEntry
63from music_assistant_models.streamdetails import MultiPartPath, StreamMetadata
64
65from music_assistant.constants import (
66 CONF_CROSSFADE_DURATION,
67 CONF_ENTRY_CROSSFADE_DIFFERENT_SAMPLE_RATES,
68 CONF_ENTRY_VOLUME_NORMALIZATION_TARGET,
69 CONF_FLOW_MODE_SAMPLE_RATE,
70 CONF_OUTPUT_CHANNELS,
71 CONF_PLAYER_QUEUES,
72 CONF_VALUE_DISABLED,
73 CONF_VALUE_ENABLED,
74 CONF_VOLUME_NORMALIZATION,
75 CONF_VOLUME_NORMALIZATION_FIXED_GAIN_RADIO,
76 CONF_VOLUME_NORMALIZATION_FIXED_GAIN_TRACKS,
77 CONF_VOLUME_NORMALIZATION_RADIO,
78 CONF_VOLUME_NORMALIZATION_TARGET,
79 CONF_VOLUME_NORMALIZATION_TRACKS,
80 DSP_IRS_DIRNAME,
81 FLOW_MODE_SAMPLE_RATE_48000,
82 FLOW_MODE_SAMPLE_RATE_96000,
83 FLOW_MODE_SAMPLE_RATE_BIT_PERFECT,
84 FLOW_MODE_SAMPLE_RATE_HIGHEST,
85 FLOW_MODE_SAMPLE_RATE_SMART,
86 INTERNAL_PCM_FORMAT,
87 MASS_LOGGER_NAME,
88 STREAM_STALL_TIMEOUT,
89 STREAM_START_TIMEOUT,
90 VERBOSE_LOG_LEVEL,
91)
92from music_assistant.controllers.streams.audio_analysis import (
93 LOUDNESS_ANALYSIS_DOMAIN,
94)
95from music_assistant.controllers.streams.audio_buffer import AudioBuffer
96from music_assistant.controllers.streams.audio_processing import (
97 AudioOutputPlan,
98 get_normalization_details,
99)
100from music_assistant.controllers.streams.constants import (
101 CACHE_CATEGORY_RESOLVED_RADIO_URL,
102 CACHE_PROVIDER,
103 CONF_ALLOW_CROSSFADE_SAME_ALBUM,
104 STREAM_SLOT_MATCH_TIMEOUT,
105 STREAM_SLOT_PLAYBACK_WAIT_TIMEOUT,
106 STREAM_SLOT_WAIT_TIMEOUT,
107 STREAMDETAILS_INBAND_TITLE_HANDOFF_KEY,
108 STREAMDETAILS_INBAND_TITLE_KEY,
109)
110from music_assistant.controllers.streams.ogg_handler import get_chained_ogg_stream
111from music_assistant.controllers.streams.smart_fades import SmartFadesMixer
112from music_assistant.controllers.streams.smart_fades.fades import SmartFade, StandardCrossFade
113from music_assistant.controllers.streams.smart_fades.helpers import (
114 MIN_EFFECTIVE_FADE_BUFFER,
115 SMART_CROSSFADE_DURATION,
116)
117from music_assistant.helpers import ssl as ssl_util
118from music_assistant.helpers.aiohttp_client import encoded_request_url
119from music_assistant.helpers.audio import (
120 HTTP_HEADERS,
121 HTTP_HEADERS_ICY,
122 build_concat_filelist,
123 calculate_content_length,
124 get_bit_rate,
125 get_normalization_mode,
126 get_parts_from_position,
127 is_grouping_preventing_dsp,
128 iter_pcm_slices,
129 parse_extinf_metadata,
130 realtime_pcm_pacer,
131 resample_pcm_audio,
132 resolve_output_player_ids,
133)
134from music_assistant.helpers.dsp import ComplexFilter, filter_to_ffmpeg_params
135from music_assistant.helpers.ffmpeg import (
136 FFMpeg,
137 get_ffmpeg_overlay_stream,
138 get_ffmpeg_stream,
139)
140from music_assistant.helpers.named_pipe import read_named_pipe
141from music_assistant.helpers.playlists import (
142 HLS_CONTENT_TYPES,
143 PLAYLIST_CONTENT_TYPES,
144 PLAYLIST_READ_TIMEOUT,
145 IsHLSPlaylist,
146 PlaylistItem,
147 parse_m3u,
148 parse_playlist_data,
149 read_playlist_body,
150)
151from music_assistant.helpers.throttle_retry import BYPASS_THROTTLER
152from music_assistant.helpers.util import (
153 clean_stream_title,
154 detect_charset,
155 parse_quoted_stream_title,
156 parse_title_and_version,
157 remove_file,
158)
159from music_assistant.models.music_provider import MusicProvider, ProviderStreamLimitError
160
161if TYPE_CHECKING:
162 from music_assistant_models.media_items import ProviderMapping
163 from music_assistant_models.player_queue import PlayerQueue
164 from music_assistant_models.queue_item import QueueItem
165 from music_assistant_models.streamdetails import StreamDetails
166
167 from music_assistant.mass import MusicAssistant
168 from music_assistant.models.player import Player
169 from music_assistant.models.plugin import PluginProvider
170 from music_assistant.models.provider import Provider
171
172# ruff: noqa: PLR0915
173
174# Seconds of PCM at the start of a track that are yielded straight to the player,
175# never held back for a crossfade.
176WARMUP_DURATION = 8
177# Minimum overlap used when the full incoming crossfade buffer was not prepared.
178MIN_CROSSFADE_FALLBACK_DURATION = 5
179
180# Chunk size for the realtime AudioSource path; small enough to keep ffmpegâconsumer
181# latency below ~50 ms while still amortising per-chunk overhead.
182AUDIO_SOURCE_CHUNK_SECONDS = 0.02
183
184# Terminal errors get_icy_radio_stream raises once a single mirror is exhausted; the
185# multi-mirror reader treats these as the signal to fail over to the next URL.
186RADIO_MIRROR_FAILOVER_ERRORS = (
187 MediaNotFoundError,
188 ProviderPermissionDenied,
189 ProviderUnavailableError,
190 RetriesExhausted,
191 InvalidDataError,
192)
193
194
195@dataclass
196class CrossfadeData:
197 """Data class to hold crossfade data."""
198
199 data: bytes
200 fade_in_media_duration: float
201 pcm_format: AudioFormat # Format of the 'data' bytes (current/previous track's format)
202 queue_item_id: str
203 # Offset for the fade_in track's elapsed time calculation, to account for crossfade duration and trim
204 elapsed_time_offset: float = 0.0
205 # Normalization mode the intro PCM was baked with, used to pin the next track's body to the same mode
206 normalization_mode: VolumeNormalizationMode | None = None
207
208
209def _snap_supported_rate_up(target: int, supported_sample_rates: list[int]) -> int:
210 """Snap target up, falling back to its highest supported divisor or the maximum."""
211 if target in supported_sample_rates:
212 return target
213 higher = [r for r in supported_sample_rates if r > target]
214 if higher:
215 return min(higher)
216 same_family = [r for r in supported_sample_rates if target % r == 0]
217 return max(same_family) if same_family else max(supported_sample_rates)
218
219
220def _snap_supported_rate_down(target: int, supported_sample_rates: list[int]) -> int:
221 """Snap target down to the highest supported rate <= target, falling back to min."""
222 if target in supported_sample_rates:
223 return target
224 lower = [r for r in supported_sample_rates if r < target]
225 return max(lower) if lower else min(supported_sample_rates)
226
227
228def _pcm_formats_match(a: AudioFormat, b: AudioFormat) -> bool:
229 """Return True if two PCM formats describe identical raw bytes."""
230 return (
231 a.content_type == b.content_type
232 and a.sample_rate == b.sample_rate
233 and a.bit_depth == b.bit_depth
234 and a.channels == b.channels
235 )
236
237
238def overlay_active(queue: PlayerQueue) -> bool:
239 """Return True if the given queue has an audio overlay enabled and a source selected."""
240 return queue.overlay_enabled and queue.overlay_source is not None
241
242
243class StreamsAudio:
244 """Audio stream acquisition and processing for the streams controller."""
245
246 def __init__(self, mass: MusicAssistant) -> None:
247 """
248 Initialize StreamsAudio.
249
250 :param mass: The MusicAssistant instance.
251 """
252 self.mass = mass
253 self.logger = logging.getLogger(f"{MASS_LOGGER_NAME}.streams.audio")
254 self._crossfade_data: dict[str, CrossfadeData] = {}
255 self._smart_fades_mixer: SmartFadesMixer | None = None
256 # serializes buffer preparation per queue item, so concurrent callers share
257 # the single source (and the single capacity reselection) instead of racing
258 self._audio_buffer_locks: WeakValueDictionary[tuple[str, str], asyncio.Lock] = (
259 WeakValueDictionary()
260 )
261
262 def setup(self) -> None:
263 """Set up the audio sub-controller (called after all core controllers are created)."""
264 self._smart_fades_mixer = SmartFadesMixer(self.mass.streams)
265
266 @property
267 def smart_fades_mixer(self) -> SmartFadesMixer:
268 """Return the smart fades mixer."""
269 assert self._smart_fades_mixer is not None, "StreamsAudio.setup() not called"
270 return self._smart_fades_mixer
271
272 # --- Public methods ---
273
274 async def get_stream_details(
275 self,
276 queue_item: QueueItem,
277 seek_position: int = 0,
278 fade_in: bool = False,
279 prefer_album_loudness: bool = False,
280 excluded_provider_instances: set[str] | None = None,
281 ) -> StreamDetails:
282 """
283 Get streamdetails for the given QueueItem.
284
285 This is called just-in-time when a PlayerQueue wants a MediaItem to be played.
286 Do not try to request streamdetails too much in advance as this is expiring data.
287
288 :param queue_item: Queue item to resolve.
289 :param seek_position: Requested playback position in seconds.
290 :param fade_in: Whether playback should fade in.
291 :param prefer_album_loudness: Whether album loudness should be preferred.
292 :param excluded_provider_instances: Provider instances to skip during this selection.
293 """
294 mass = self.mass
295 streamdetails: StreamDetails | None = None
296 excluded_provider_instances = excluded_provider_instances or set()
297 time_start = time.time()
298 self.logger.debug("Getting streamdetails for %s", queue_item.uri)
299
300 if not queue_item.media_item and not queue_item.streamdetails:
301 # in case of a non-media item queue item, the streamdetails should already be provided
302 # this should not happen, but guard it just in case
303 raise MediaNotFoundError(
304 f"Unable to retrieve streamdetails for {queue_item.name} ({queue_item.uri})"
305 )
306
307 if (
308 queue_item.streamdetails
309 # cached details of an excluded instance are exactly what we select away from
310 and queue_item.streamdetails.provider not in excluded_provider_instances
311 and (
312 # reuse if the buffer can serve this seek position (fast seek path)
313 (
314 queue_item.streamdetails.buffer
315 and queue_item.streamdetails.buffer.is_valid(int(seek_position * 1000))
316 )
317 # or reuse if streamdetails hasn't expired yet (new buffer will be created)
318 or (queue_item.streamdetails.created_at + queue_item.streamdetails.expiration)
319 > time.time()
320 )
321 ):
322 streamdetails = queue_item.streamdetails
323 else:
324 # need to (re)create streamdetails
325 # retrieve streamdetails from provider
326
327 media_item = queue_item.media_item
328 assert media_item is not None # for type checking
329 preferred_providers: list[str] = []
330 if (
331 (pq_data := mass.player_queues.queue_data_or_none(queue_item.queue_id))
332 and pq_data.userid
333 and (playback_user := await mass.webserver.auth.get_user(pq_data.userid))
334 and playback_user.provider_filter
335 ):
336 # handle steering into user preferred providerinstance
337 preferred_providers = playback_user.provider_filter
338 candidates = self._get_streamdetail_candidates(
339 media_item.provider_mappings,
340 preferred_providers,
341 excluded_provider_instances,
342 )
343 streamdetails = await self._request_streamdetails(candidates, media_item.media_type)
344
345 if not streamdetails:
346 msg = f"Unable to retrieve streamdetails for {queue_item.name} ({queue_item.uri})"
347 raise MediaNotFoundError(msg)
348
349 # work out how to handle radio stream
350 if (
351 streamdetails.stream_type in (StreamType.ICY, StreamType.HLS, StreamType.HTTP)
352 and streamdetails.media_type == MediaType.RADIO
353 and isinstance(streamdetails.path, str)
354 ):
355 resolved_url, stream_type = await self.resolve_radio_stream(streamdetails.path)
356 streamdetails.path = resolved_url
357 streamdetails.stream_type = stream_type
358 # Set up metadata monitoring callback for HLS radio streams, if not already set
359 if (
360 stream_type == StreamType.HLS
361 and not streamdetails.stream_metadata_update_callback
362 ):
363 streamdetails.stream_metadata_update_callback = partial(
364 self._update_hls_radio_metadata
365 )
366 streamdetails.stream_metadata_update_interval = 5
367
368 # providers report an unknown duration as either None or 0
369 if not streamdetails.duration:
370 if queue_item.media_item and queue_item.media_item.duration:
371 streamdetails.duration = queue_item.media_item.duration
372 elif queue_item.duration:
373 streamdetails.duration = queue_item.duration
374 if seek_position and not streamdetails.allow_seek:
375 self.logger.warning("seeking is not possible on this stream!")
376 seek_position = 0
377 elif seek_position and not streamdetails.duration:
378 self.logger.warning("seeking is not possible on duration-less streams!")
379 seek_position = 0
380
381 if streamdetails.media_type in (MediaType.RADIO, MediaType.AUDIO_SOURCE):
382 # radio stations and live audio sources hand over their audio at playback pace
383 streamdetails.is_realtime = True
384
385 # set queue_id on the streamdetails so we know what is being streamed
386 streamdetails.queue_id = queue_item.queue_id
387 # handle skip/fade_in details
388 streamdetails.seek_position = seek_position
389 streamdetails.fade_in = fade_in
390
391 streamdetails.prefer_album_loudness = prefer_album_loudness
392 conf_volume_normalization_target = float(
393 mass.streams.get_config_value(CONF_VOLUME_NORMALIZATION_TARGET, return_type=int)
394 )
395 # guard against invalid volume normalization values
396 # range and default_value are guaranteed to be set for this constant
397 volume_range = CONF_ENTRY_VOLUME_NORMALIZATION_TARGET.range
398 assert volume_range is not None
399 if (
400 conf_volume_normalization_target < volume_range[0]
401 or conf_volume_normalization_target >= volume_range[1]
402 ):
403 default_val = CONF_ENTRY_VOLUME_NORMALIZATION_TARGET.default_value
404 assert isinstance(default_val, (int, float))
405 conf_volume_normalization_target = float(default_val)
406 self.logger.warning(
407 "Invalid volume normalization target configured, resetting to default of %s LUFS",
408 CONF_ENTRY_VOLUME_NORMALIZATION_TARGET.default_value,
409 )
410 streamdetails.target_loudness = conf_volume_normalization_target
411 volume_normalization_enabled = (
412 mass.config.get_effective_player_queue_config_value(
413 streamdetails.queue_id, CONF_VOLUME_NORMALIZATION, CONF_VALUE_ENABLED
414 )
415 != CONF_VALUE_DISABLED
416 )
417 streamdetails.volume_normalization_mode = get_normalization_mode(
418 self._get_volume_normalization_preference(streamdetails),
419 volume_normalization_enabled,
420 streamdetails,
421 )
422
423 self.logger.debug(
424 "Retrieved streamdetails for %s in %s milliseconds",
425 queue_item.uri,
426 int((time.time() - time_start) * 1000),
427 )
428 return streamdetails
429
430 async def get_audio_buffer(
431 self,
432 queue_item: QueueItem,
433 seek_position_ms: int = 0,
434 reason: str = "",
435 capacity_wait_timeout: float = STREAM_SLOT_PLAYBACK_WAIT_TIMEOUT,
436 allow_provider_match: bool = True,
437 ) -> AudioBuffer:
438 """
439 Return a ready AudioBuffer for the given queue item.
440
441 Compatible provider mappings are reselected while the owning provider has no free
442 source-stream slot. Other AudioErrors propagate as on a direct buffer request.
443
444 :param queue_item: Queue item whose source should be buffered.
445 :param seek_position_ms: Position in milliseconds to start from.
446 :param reason: Caller context for logging (e.g. 'prepare_next', 'streaming').
447 :param capacity_wait_timeout: Total seconds to spend waiting for source capacity.
448 :param allow_provider_match: Whether an on-demand cross-provider match may widen
449 the candidates when all are saturated.
450 :raises ProviderStreamLimitError: If no source slot becomes available within the budget.
451 """
452 lock_key = (queue_item.queue_id, queue_item.queue_item_id)
453 if (buffer_lock := self._audio_buffer_locks.get(lock_key)) is None:
454 buffer_lock = asyncio.Lock()
455 self._audio_buffer_locks[lock_key] = buffer_lock
456 async with buffer_lock:
457 return await self._get_audio_buffer(
458 queue_item, seek_position_ms, reason, capacity_wait_timeout, allow_provider_match
459 )
460
461 async def get_media_stream(
462 self,
463 streamdetails: StreamDetails,
464 pcm_format: AudioFormat,
465 seek_position: int = 0,
466 filter_params: list[str] | None = None,
467 chunk_seconds: float = 1.0,
468 source_wait_timeout: float | None = STREAM_SLOT_WAIT_TIMEOUT,
469 ) -> AsyncGenerator[bytes]:
470 """
471 Get audio stream for given media details as raw PCM.
472
473 :param streamdetails: Details of the stream to fetch.
474 :param pcm_format: Target PCM format the consumer expects.
475 :param seek_position: Seek offset in seconds (only honoured when the
476 source allows seeking; ignored for live AudioSources).
477 :param filter_params: Optional ffmpeg filter expressions.
478 :param chunk_seconds: Size of each yielded chunk in seconds of audio.
479 Defaults to 1 s for track-like sources; callers streaming live
480 AudioSources should pass a much smaller value (e.g. 0.02) to keep
481 end-to-end latency low.
482 :param source_wait_timeout: Maximum seconds to wait for a free source-stream slot
483 on the providing music provider, or None to wait without a timeout.
484 :raises ProviderStreamLimitError: If the provider has no free slot within the timeout.
485 """
486 media_stream = self._get_media_stream(
487 streamdetails,
488 pcm_format,
489 seek_position,
490 filter_params,
491 chunk_seconds,
492 )
493 # resolve the exact owning instance (even when flagged unavailable) so the
494 # slot is charged to the account that issued the streamdetails
495 provider = self.mass.get_provider(streamdetails.provider, return_unavailable=True)
496 stream_slot = (
497 provider.acquire_stream_slot(source_wait_timeout)
498 if isinstance(provider, MusicProvider)
499 else nullcontext()
500 )
501 async with stream_slot, aclosing(media_stream):
502 async for chunk in media_stream:
503 yield chunk
504
505 async def resolve_radio_stream(self, url: str) -> tuple[str, StreamType]:
506 """
507 Resolve a streaming radio URL.
508
509 Unwraps playlists and determines stream type (ICY, HLS, SHOUTCAST, IN_BAND, HTTP).
510
511 :param url: Radio stream URL to resolve
512 """
513 mass = self.mass
514 if cache := await mass.cache.get(
515 key=url, provider=CACHE_PROVIDER, category=CACHE_CATEGORY_RESOLVED_RADIO_URL
516 ):
517 if TYPE_CHECKING:
518 cache = cast("tuple[str, str]", cache)
519 return (cache[0], StreamType(cache[1]))
520
521 stream_type = StreamType.HTTP
522 timeout = ClientTimeout(total=None, connect=10, sock_read=5)
523 playlist_data: bytes | None = None
524 playlist_charset: str | None = None
525
526 try:
527 async with self._connect_radio_stream(
528 url, headers=HTTP_HEADERS_ICY, allow_redirects=True, timeout=timeout
529 ) as resp:
530 headers = resp.headers
531 resp.raise_for_status()
532 if not resp.headers:
533 raise InvalidDataError("no headers found")
534 # media types are case insensitive, the comparisons below are all lower case
535 content_type = headers.get("content-type", "").lower()
536 # a server declaring HLS settles it: a media playlist is free to carry none
537 # of the tags the parser recognises an HLS playlist by
538 is_hls = any(hls_type in content_type for hls_type in HLS_CONTENT_TYPES)
539 if not is_hls and (
540 url.endswith((".m3u", ".m3u8", ".pls"))
541 or ".m3u?" in url
542 or ".m3u8?" in url
543 or ".pls?" in url
544 or any(
545 playlist_type in content_type for playlist_type in PLAYLIST_CONTENT_TYPES
546 )
547 ):
548 # take the playlist from this very response: a separate request would
549 # go out with another user agent and stricter TLS than the rest of the
550 # radio paths, so a host could answer it differently
551 try:
552 # the probe has no total timeout, so bound the body on its own:
553 # a server trickling bytes would otherwise stall resolving for hours
554 async with asyncio.timeout(PLAYLIST_READ_TIMEOUT):
555 playlist_data = await read_playlist_body(resp.content)
556 except aiohttp.ClientError as err:
557 # the endpoint answered as a playlist, so a truncated body is a bad
558 # playlist - not a reason to fall back to streaming the URL directly
559 raise InvalidDataError(f"Error while fetching playlist {url}") from err
560 playlist_charset = resp.charset
561
562 if headers.get("icy-metaint") is not None:
563 stream_type = StreamType.ICY
564 elif is_hls:
565 stream_type = StreamType.HLS
566 elif content_type in ("application/ogg", "audio/ogg"):
567 # Ogg streams (Opus/Vorbis) have in-band metadata via Vorbis comments
568 stream_type = StreamType.IN_BAND
569
570 if playlist_data is not None:
571 try:
572 substreams = await parse_playlist_data(url, playlist_data, playlist_charset)
573 if not any(x for x in substreams if x.length):
574 for line in substreams:
575 if not line.is_url:
576 continue
577 return await self.resolve_radio_stream(line.path)
578 raise InvalidDataError("No content found in playlist")
579 except IsHLSPlaylist:
580 stream_type = StreamType.HLS
581
582 except TimeoutError as err:
583 self.logger.warning("Timeout while parsing radio URL %s", url)
584 raise InvalidDataError(f"Timeout connecting to {url}") from err
585
586 except aiohttp.ClientResponseError as err:
587 if err.status == 404:
588 raise MediaNotFoundError(f"Radio stream not found: {url}") from err
589 if err.status == 403:
590 raise InvalidDataError(f"Access denied to radio stream: {url}") from err
591 if err.status >= 500:
592 raise InvalidDataError(
593 f"Radio stream server error (HTTP {err.status}): {url}"
594 ) from err
595 if err.status == 400:
596 # 400 errors might be from legacy Shoutcast servers
597 return await self._handle_client_error_for_radio_stream(url, err, stream_type)
598 raise InvalidDataError(f"HTTP error {err.status} from {url}") from err
599
600 except aiohttp.ClientError as err:
601 return await self._handle_client_error_for_radio_stream(url, err, stream_type)
602
603 return await self._cache_radio_result(url, stream_type)
604
605 async def get_icy_radio_stream(
606 self, url: str, streamdetails: StreamDetails
607 ) -> AsyncGenerator[bytes]:
608 """
609 Stream radio audio with ICY metadata support, reconnecting on disconnect.
610
611 Requires icy-metaint header support. Stream type should be validated
612 by resolve_radio_stream() before calling this function.
613
614 :param url: Radio stream URL
615 :param streamdetails: StreamDetails to update with metadata
616 """
617 self.logger.debug("Start streaming radio with ICY metadata from url %s", url)
618 timeout = ClientTimeout(total=0, connect=30, sock_read=5 * 60)
619 # Budget for *consecutive* reconnects that delivered no audio. A connection
620 # that actually streamed data resets it, so a healthy long-running stream can
621 # reconnect indefinitely while a dead/looping one bails out instead of spinning.
622 failed_reconnects = 0
623 max_failed_reconnects = 25
624
625 while True:
626 streamed_data = False
627 try:
628 async with self._connect_radio_stream(
629 url, allow_redirects=True, headers=HTTP_HEADERS_ICY, timeout=timeout
630 ) as resp:
631 # surface a non-200 (e.g. on reconnect) as a ClientResponseError so the
632 # terminal/HTTP handling below applies instead of failing on the header
633 resp.raise_for_status()
634 meta_int_str = resp.headers.get("icy-metaint")
635 if not meta_int_str:
636 raise InvalidDataError(f"No icy-metaint header for radio stream: {url}")
637 try:
638 meta_int = int(meta_int_str)
639 except ValueError as err:
640 raise InvalidDataError(
641 f"Invalid icy-metaint value for radio stream: {url}"
642 ) from err
643 if meta_int <= 0:
644 raise InvalidDataError(f"Invalid icy-metaint value for radio stream: {url}")
645 # readexactly raises IncompleteReadError when the server closes the
646 # connection mid-frame; that (and the network errors below) drops us
647 # out to the reconnect handler so a live stream survives the blip.
648 while True:
649 chunk = await resp.content.readexactly(meta_int)
650 streamed_data = True
651 yield chunk
652 meta_byte = await resp.content.readexactly(1)
653 if meta_byte == b"\x00":
654 continue
655 meta_length = ord(meta_byte) * 16
656 meta_data = await resp.content.readexactly(meta_length)
657 self._parse_icy_metadata(meta_data, streamdetails)
658 except asyncio.CancelledError:
659 self.logger.debug("ICY radio stream cancelled for %s", url)
660 raise
661 except aiohttp.ClientResponseError as err:
662 if err.status == 404:
663 raise MediaNotFoundError(f"Radio stream not found: {url}") from err
664 if err.status == 403:
665 raise ProviderPermissionDenied(f"Radio stream access denied: {url}") from err
666 raise ProviderUnavailableError(
667 f"Radio stream returned HTTP {err.status}: {err}"
668 ) from err
669 except (
670 asyncio.IncompleteReadError,
671 aiohttp.ClientConnectionError,
672 aiohttp.ClientPayloadError,
673 aiohttp.ServerDisconnectedError,
674 ) as err:
675 if streamed_data:
676 # a healthy session that dropped - reconnect without spending budget
677 failed_reconnects = 0
678 self.logger.debug("ICY radio stream dropped, reconnecting: %s", err)
679 else:
680 failed_reconnects += 1
681 if failed_reconnects > max_failed_reconnects:
682 raise RetriesExhausted(
683 f"ICY radio stream failed after {max_failed_reconnects} "
684 f"reconnects without data: {err}"
685 ) from err
686 self.logger.warning(
687 "ICY radio stream reconnect produced no data (%d/%d): %s",
688 failed_reconnects,
689 max_failed_reconnects,
690 err,
691 )
692 await asyncio.sleep(0.5)
693
694 async def get_reconnecting_icy_radio_stream(
695 self, url: str | list[MultiPartPath], streamdetails: StreamDetails
696 ) -> AsyncGenerator[bytes]:
697 """
698 Yield ICY radio audio with metadata, failing over across mirror URLs.
699
700 A single URL is delegated to :meth:`get_icy_radio_stream`, which already reconnects
701 on disconnect. Multiple URLs are treated as interchangeable mirrors and tried in turn;
702 a mirror that delivers audio resets the failover budget, so a healthy mirror keeps
703 streaming while a set of unreachable mirrors raises the last error instead of spinning.
704
705 :param url: One stream URL, or a list of mirror URLs to fail over between.
706 :param streamdetails: StreamDetails to update with metadata.
707 """
708 urls = self._normalize_reconnecting_urls(url)
709 if len(urls) == 1:
710 async for chunk in self.get_icy_radio_stream(urls[0], streamdetails):
711 yield chunk
712 return
713
714 url_index = 0
715 failed_rotations = 0
716 max_failed_rotations = len(urls) * 2
717 last_err: MusicAssistantError | None = None
718 while failed_rotations <= max_failed_rotations:
719 current_url = urls[url_index % len(urls)]
720 url_index += 1
721 delivered_audio = False
722 try:
723 async for chunk in self.get_icy_radio_stream(current_url, streamdetails):
724 delivered_audio = True
725 failed_rotations = 0
726 # release the previous failure while healthy: it pins the full
727 # exception traceback (with frames) for the lifetime of the stream
728 last_err = None
729 yield chunk
730 return
731 except RADIO_MIRROR_FAILOVER_ERRORS as err:
732 last_err = err
733 if not delivered_audio:
734 failed_rotations += 1
735 self.logger.warning(
736 "ICY radio mirror %s failed, trying next url (%d/%d): %s",
737 current_url,
738 failed_rotations,
739 max_failed_rotations,
740 err,
741 )
742 if last_err is not None:
743 raise last_err
744
745 async def get_reconnecting_radio_stream(self, url: str) -> AsyncGenerator[bytes]:
746 """
747 Yield continuous radio stream data, automatically reconnecting on disconnect.
748
749 :param url: URL of the radio stream.
750 """
751 timeout = ClientTimeout(total=None, connect=30, sock_read=5 * 60)
752 reconnect_count = 0
753 max_reconnects = 1000 # Allow many reconnects for long-running radio
754
755 while reconnect_count <= max_reconnects:
756 try:
757 async with self._connect_radio_stream(
758 url, allow_redirects=True, headers=HTTP_HEADERS, timeout=timeout
759 ) as resp:
760 chunk_count = 0
761 async for chunk in resp.content.iter_any():
762 chunk_count += 1
763 yield chunk
764
765 # Connection closed normally - reconnect
766 self.logger.debug(
767 "Radio stream connection closed after %d chunks, reconnecting... "
768 "(reconnect #%d)",
769 chunk_count,
770 reconnect_count,
771 )
772 reconnect_count += 1
773 await asyncio.sleep(0.1) # Brief delay before reconnect
774
775 except asyncio.CancelledError:
776 self.logger.debug("Radio stream cancelled for %s", url)
777 raise
778 except (
779 aiohttp.ClientConnectionError,
780 aiohttp.ClientPayloadError,
781 aiohttp.ServerDisconnectedError,
782 ) as err:
783 # Transient network errors - retry
784 self.logger.warning("Radio stream error (reconnect #%d): %s", reconnect_count, err)
785 reconnect_count += 1
786 if reconnect_count > max_reconnects:
787 raise RetriesExhausted(
788 f"Radio stream failed after {max_reconnects} reconnects: {err}"
789 ) from err
790 await asyncio.sleep(0.5)
791 except aiohttp.ClientResponseError as err:
792 if err.status == 404:
793 raise MediaNotFoundError(f"Radio stream not found: {url}") from err
794 if err.status == 403:
795 raise ProviderPermissionDenied(f"Radio stream access denied: {url}") from err
796 # Other HTTP errors (5xx etc) - could be temporary
797 raise ProviderUnavailableError(
798 f"Radio stream returned HTTP {err.status}: {err}"
799 ) from err
800
801 self.logger.warning("Radio stream reached max reconnects (%d) for %s", max_reconnects, url)
802
803 async def get_hls_substream(self, url: str) -> PlaylistItem:
804 """Select the (highest quality) HLS substream for given HLS playlist/URL."""
805 mass = self.mass
806 timeout = ClientTimeout(total=None, connect=30, sock_read=5 * 60)
807 # fetch master playlist and select (best) child playlist
808 # https://datatracker.ietf.org/doc/html/draft-pantos-http-live-streaming-19#section-10
809 async with mass.http_session_no_ssl.get(
810 encoded_request_url(url), allow_redirects=True, headers=HTTP_HEADERS, timeout=timeout
811 ) as resp:
812 resp.raise_for_status()
813 raw_data = await resp.read()
814 encoding = await detect_charset(raw_data, preferred=resp.charset)
815 master_m3u_data = raw_data.decode(encoding, errors="replace")
816 substreams = parse_m3u(master_m3u_data)
817 # There is a chance that we did not get a master playlist with subplaylists
818 # but just a single master/sub playlist with the actual audio stream(s)
819 # so we need to detect if the playlist child's contain audio streams or
820 # sub-playlists.
821 if any(
822 x
823 for x in substreams
824 if (x.length or x.path.endswith((".mp4", ".aac")))
825 and not x.path.endswith((".m3u", ".m3u8"))
826 ):
827 return PlaylistItem(path=url, key=substreams[0].key)
828 # sort substreams on best quality (highest bandwidth) when available
829 if any(x for x in substreams if x.stream_info):
830 substreams.sort(
831 key=lambda x: int(
832 x.stream_info.get("BANDWIDTH", "0") if x.stream_info is not None else 0
833 ),
834 reverse=True,
835 )
836 substream = substreams[0]
837 if not substream.path.startswith("http"):
838 # path is relative, stitch it together
839 base_path = url.rsplit("/", 1)[0]
840 substream.path = base_path + "/" + substream.path
841 return substream
842
843 async def get_multi_file_stream(
844 self,
845 streamdetails: StreamDetails,
846 seek_position: int = 0,
847 ) -> AsyncGenerator[bytes]:
848 """
849 Return audio stream for a concatenation of multiple files.
850
851 Arguments:
852 seek_position: The position to seek to in seconds
853 """
854 if not isinstance(streamdetails.path, list):
855 raise InvalidDataError("Multi-file streamdetails requires a list of MultiPartPath")
856 parts, seek_position = get_parts_from_position(streamdetails.path, seek_position)
857 files_list = [part.path for part in parts]
858
859 # concat input files
860 temp_file = f"/tmp/{shortuuid.random(20)}.txt" # noqa: S108
861 async with aiofiles.open(temp_file, "w") as f:
862 await f.write(build_concat_filelist(files_list))
863
864 try:
865 async for chunk in get_ffmpeg_stream(
866 audio_input=temp_file,
867 input_format=streamdetails.audio_format,
868 output_format=AudioFormat(
869 content_type=ContentType.NUT,
870 sample_rate=streamdetails.audio_format.sample_rate,
871 bit_depth=streamdetails.audio_format.bit_depth,
872 channels=streamdetails.audio_format.channels,
873 ),
874 extra_input_args=[
875 "-safe",
876 "0",
877 "-f",
878 "concat",
879 "-i",
880 temp_file,
881 "-ss",
882 str(seek_position),
883 ],
884 ):
885 yield chunk
886 finally:
887 await remove_file(temp_file)
888
889 def get_player_output_plan(
890 self,
891 player_id: str,
892 input_format: AudioFormat,
893 output_format: AudioFormat,
894 *,
895 shared_player_ids: Iterable[str] | None = None,
896 handoff_format: AudioFormat | None = None,
897 queue_id: str | None = None,
898 session_id: str | None = None,
899 queue_item_id: str | None = None,
900 ) -> AudioOutputPlan:
901 """
902 Return executable filters and matching output details for a player.
903
904 :param player_id: Destination player identifier.
905 :param input_format: PCM format entering player-specific processing.
906 :param output_format: Furthest downstream output format known to the server.
907 :param shared_player_ids: Additional players receiving this identical output path.
908 An empty iterable marks a path that can gain shared destinations later.
909 :param handoff_format: Earlier provider handoff format when it differs.
910 :param queue_id: Explicit queue identifier for the processing snapshot.
911 :param session_id: Explicit queue session identifier for the processing snapshot.
912 :param queue_item_id: Queue item for a single-item output path.
913 """
914 filter_params: list[str | ComplexFilter] = []
915 player = self.mass.players.get_player(player_id)
916 destination_player_id = (
917 player.protocol_parent_id if player and player.protocol_parent_id else player_id
918 )
919 resolved_shared_player_ids = (
920 resolve_output_player_ids(self.mass, shared_player_ids) - {destination_player_id}
921 if shared_player_ids is not None
922 else None
923 )
924 destination_player_ids = {destination_player_id, *(resolved_shared_player_ids or ())}
925 if player:
926 dsp_config_id = self._resolve_player_dsp_config_id(player)
927 dsp = self._resolve_player_dsp_config(player)
928 configured_dsp = self.mass.config.get_player_dsp_config(dsp_config_id)
929 if configured_dsp.enabled and not dsp.enabled and is_grouping_preventing_dsp(player):
930 dsp_state = DSPState.DISABLED_BY_UNSUPPORTED_GROUP
931 else:
932 dsp_state = DSPState.ENABLED if dsp.enabled else DSPState.DISABLED
933 else:
934 dsp_config_id = player_id
935 dsp = self.mass.config.get_player_dsp_config(player_id)
936 dsp_state = DSPState.ENABLED if dsp.enabled else DSPState.DISABLED
937
938 enabled_filters = [dsp_filter for dsp_filter in dsp.filters if dsp_filter.enabled]
939 # a neutral filter (0 dB gain, centered balance) emits no params; exclude
940 # it so it is not reported as an active, non-bit-perfect stage
941 effective_filters: list[DSPFilter] = []
942 if dsp.enabled:
943 if dsp.input_gain != 0:
944 filter_params.append(f"volume={dsp.input_gain}dB")
945 ir_dir = os.path.join(self.mass.storage_path, DSP_IRS_DIRNAME)
946 known_ir_ids = {record["ir_id"] for record in self.mass.config.get_dsp_irs()}
947 for dsp_filter in enabled_filters:
948 if isinstance(dsp_filter, ConvolutionFilter) and dsp_filter.ir_id:
949 # ffmpeg fails to open the graph if the impulse response file is gone,
950 # which costs the player all audio, so drop the filter instead
951 if dsp_filter.ir_id not in known_ir_ids:
952 self.logger.warning(
953 "Skipping the convolution filter of player %s: "
954 "impulse response %s is not stored",
955 player_id,
956 dsp_filter.ir_id,
957 )
958 continue
959 params = filter_to_ffmpeg_params(dsp_filter, input_format, ir_dir=ir_dir)
960 if not params:
961 continue
962 filter_params.extend(params)
963 effective_filters.append(dsp_filter)
964 if dsp.output_gain != 0:
965 filter_params.append(f"volume={dsp.output_gain}dB")
966
967 channel_value = self._get_output_channels(player, player_id)
968 source_channel = None
969 channel_mix = ""
970 # a single channel source is already the downmix and holds no FL/FR to select
971 # from, where a pan would silently resolve every gain to zero
972 if input_format.channels > 1:
973 if channel_value == "left":
974 source_channel = AudioChannel.FL
975 channel_mix = "FL"
976 elif channel_value == "right":
977 source_channel = AudioChannel.FR
978 channel_mix = "FR"
979 elif channel_value == "mono":
980 # both source channels feed the downmix, so report ALL to keep the
981 # output from ever being presented as bit perfect
982 source_channel = AudioChannel.ALL
983 channel_mix = "0.5*FL+0.5*FR"
984 if channel_mix:
985 # the pan runs in the command that emits the handoff format, and it feeds
986 # every channel of it explicitly: leaving ffmpeg to upmix from a single
987 # channel costs 3 dB through its rematrix
988 if (handoff_format or output_format).channels == 1:
989 filter_params.append(f"pan=mono|c0={channel_mix}")
990 else:
991 filter_params.append(f"pan=stereo|c0={channel_mix}|c1={channel_mix}")
992
993 output_details = AudioOutputDetails(
994 player_ids=sorted(destination_player_ids),
995 dsp=AudioDSPDetails(
996 state=dsp_state,
997 input_gain=dsp.input_gain if dsp.enabled else 0.0,
998 filters=effective_filters,
999 output_gain=dsp.output_gain if dsp.enabled else 0.0,
1000 preset_id=dsp.preset_id,
1001 ),
1002 source_channel=source_channel,
1003 output_format=output_format,
1004 )
1005 output_plan = AudioOutputPlan(
1006 filter_params=filter_params,
1007 output_details=output_details,
1008 input_format=input_format,
1009 handoff_format=handoff_format,
1010 dsp_config_id=dsp_config_id,
1011 )
1012 if queue_id is not None and session_id is not None:
1013 self.mass.streams.audio_processing.update_output(
1014 destination_player_id,
1015 output_plan,
1016 shared_player_ids=resolved_shared_player_ids,
1017 queue_id=queue_id,
1018 session_id=session_id,
1019 queue_item_id=queue_item_id,
1020 )
1021 self.logger.log(
1022 VERBOSE_LOG_LEVEL,
1023 "Generated ffmpeg params for player %s: %s",
1024 player_id,
1025 filter_params,
1026 )
1027 return output_plan
1028
1029 async def get_output_format(
1030 self,
1031 output_format_str: str,
1032 player: Player,
1033 content_sample_rate: int,
1034 content_bit_depth: int,
1035 media_type: MediaType = MediaType.UNKNOWN,
1036 ) -> AudioFormat:
1037 """Parse (player specific) output format details for given format string."""
1038 content_type: ContentType = ContentType.try_parse(output_format_str)
1039 player_supported_rates = player.get_supported_sample_rates()
1040 supported_sample_rates = [sr for sr, _ in player_supported_rates]
1041 if content_sample_rate in supported_sample_rates:
1042 output_sample_rate = content_sample_rate
1043 else:
1044 output_sample_rate = max(supported_sample_rates)
1045 # only consider bit depths that are actually paired with the chosen sample rate
1046 bit_depths_for_rate = [
1047 bd for (sr, bd) in player_supported_rates if sr == output_sample_rate
1048 ]
1049 output_bit_depth = min(content_bit_depth, max(bit_depths_for_rate, default=16))
1050
1051 if not content_type.is_lossless():
1052 # no point in having a higher bit depth for lossy formats
1053 output_bit_depth = 16
1054 output_sample_rate = min(48000, output_sample_rate)
1055 if media_type not in (MediaType.TRACK, MediaType.AUDIO_SOURCE, MediaType.FLOW_STREAM):
1056 # no point in having a higher bit depth for non-track media types (e.g. TTS, radio)
1057 output_bit_depth = min(output_bit_depth, 16)
1058 if output_format_str == "pcm":
1059 content_type = ContentType.from_bit_depth(output_bit_depth)
1060
1061 output_channels_str = self._get_output_channels(player, player.player_id)
1062 fmt = AudioFormat(
1063 content_type=content_type,
1064 sample_rate=output_sample_rate,
1065 bit_depth=output_bit_depth,
1066 channels=1 if output_channels_str != "stereo" else 2,
1067 )
1068 fmt.bit_rate = get_bit_rate(fmt)
1069 return fmt
1070
1071 async def select_pcm_format(
1072 self,
1073 player: Player,
1074 streamdetails: StreamDetails,
1075 crossfade_enabled: bool,
1076 overlay_active: bool = False,
1077 ) -> AudioFormat:
1078 """
1079 Select the internal PCM format for streaming a single queue item.
1080
1081 Used by the per-item (non-flow) stream path. The sample rate is the highest
1082 rate the player supports that is <= the source rate, so the source is never
1083 upsampled. The bit depth follows the source unless audio processing
1084 (crossfade, volume normalization, DSP) is active â those need F32 headroom
1085 to avoid clipping/precision loss. Surround sources are folded down to stereo.
1086 Realtime AudioSource items skip all processing and get a pure passthrough
1087 format (source rate/bit depth when the player supports them).
1088
1089 :param player: The player requesting the stream.
1090 :param streamdetails: Stream details for the current item.
1091 :param crossfade_enabled: Whether crossfade is enabled for this stream.
1092 :param overlay_active: Whether an audio overlay will be mixed into this stream.
1093 """
1094 if streamdetails.media_type == MediaType.AUDIO_SOURCE:
1095 return self._select_audio_source_pcm_format(player, streamdetails)
1096 supported_sample_rates = [sr for sr, _ in player.get_supported_sample_rates()]
1097 # snap-down: pick the highest supported rate <= source. when the source rate
1098 # is below every supported rate (e.g. 22 kHz content on a 44.1k-only player),
1099 # fall back to the lowest supported rate instead of a hardcoded 48 kHz that
1100 # the player may not actually support.
1101 output_sample_rate = max(
1102 (r for r in supported_sample_rates if r <= streamdetails.audio_format.sample_rate),
1103 default=min(supported_sample_rates),
1104 )
1105 content_type, bit_depth = self._pick_pcm_bit_depth(
1106 (player,),
1107 streamdetails,
1108 crossfade_enabled,
1109 overlay_active,
1110 )
1111 pcm_format = AudioFormat(
1112 sample_rate=output_sample_rate,
1113 content_type=content_type,
1114 bit_depth=bit_depth,
1115 # fold surround sources down to stereo right at the decode step: no
1116 # output format carries more than two channels, so a wider PCM format
1117 # only makes every bytes-to-seconds sum on the stream come out short
1118 channels=min(streamdetails.audio_format.channels, 2),
1119 )
1120 if crossfade_enabled or overlay_active:
1121 pcm_format.channels = 2
1122 return pcm_format
1123
1124 async def select_flow_pcm_format(
1125 self,
1126 player: Player,
1127 start_streamdetails: StreamDetails | None = None,
1128 crossfade_enabled: bool = False,
1129 overlay_active: bool = False,
1130 fallback_sample_rate: int | None = None,
1131 output_players: Iterable[Player] | None = None,
1132 ) -> AudioFormat:
1133 """
1134 Select the internal PCM format for a Queue Flow Mode stream.
1135
1136 Used by the gapless flow path that stitches multiple queue items into one
1137 continuous PCM stream. The sample rate is driven by the player's
1138 ``CONF_FLOW_MODE_SAMPLE_RATE`` setting (smart/bit_perfect/48k/96k/highest)
1139 â for the anchored modes it follows the first track's rate, for the fixed
1140 modes it snaps to the configured rate. The bit depth follows the first
1141 track's source unless audio processing is active (then F32 for headroom),
1142 avoiding an unnecessary up-convert to 32-bit when none of the consumers
1143 will benefit from it. When the first item is a realtime AudioSource, the
1144 flow mode config is ignored and a pure passthrough format is used so the
1145 source audio is delivered with minimum overhead and latency.
1146
1147 :param player: The player the flow stream is being prepared for.
1148 :param start_streamdetails: Stream details of the first track in the flow.
1149 Required for the anchored modes ('smart' / 'bit_perfect') and for the
1150 bit-depth optimization. May be omitted for the fixed-rate modes â when
1151 omitted the bit depth defaults to F32.
1152 :param crossfade_enabled: Whether the queue will use crossfade transitions.
1153 :param overlay_active: Whether an audio overlay will be mixed into the stream.
1154 :param fallback_sample_rate: Preferred rate when the first item format is unknown.
1155 :param output_players: All players consuming the shared PCM stream. Their common
1156 sample rates and processing requirements determine the session format.
1157 """
1158 players = tuple(output_players) if output_players is not None else (player,)
1159 if not players:
1160 raise AudioError("At least one output player is required")
1161 supported_sample_rates = sorted(
1162 set.intersection(
1163 *(
1164 {sample_rate for sample_rate, _ in item.get_supported_sample_rates()}
1165 for item in players
1166 )
1167 )
1168 )
1169 if not supported_sample_rates:
1170 raise AudioError("Output players do not share a supported sample rate")
1171 if start_streamdetails is not None and (
1172 start_streamdetails.media_type == MediaType.AUDIO_SOURCE
1173 ):
1174 return self._select_audio_source_pcm_format(
1175 player,
1176 start_streamdetails,
1177 supported_sample_rates=supported_sample_rates,
1178 )
1179 flow_mode_conf = cast(
1180 "str",
1181 player.config.get_value(CONF_FLOW_MODE_SAMPLE_RATE, FLOW_MODE_SAMPLE_RATE_SMART),
1182 )
1183
1184 if flow_mode_conf == FLOW_MODE_SAMPLE_RATE_HIGHEST:
1185 output_sample_rate = max(supported_sample_rates)
1186 elif flow_mode_conf == FLOW_MODE_SAMPLE_RATE_48000:
1187 # for the fixed-rate modes, the user picked a specific bandwidth/quality
1188 # ceiling; prefer the highest supported rate <= target
1189 output_sample_rate = _snap_supported_rate_down(48000, supported_sample_rates)
1190 elif flow_mode_conf == FLOW_MODE_SAMPLE_RATE_96000:
1191 output_sample_rate = _snap_supported_rate_down(96000, supported_sample_rates)
1192 else:
1193 # smart or bit_perfect (default): anchor the flow at the starting track's
1194 # sample rate; if the player doesn't natively support it, upsample to the
1195 # closest higher supported rate
1196 target_rate = (
1197 start_streamdetails.audio_format.sample_rate
1198 if start_streamdetails
1199 else (
1200 fallback_sample_rate
1201 if fallback_sample_rate is not None
1202 else max(supported_sample_rates)
1203 )
1204 )
1205 output_sample_rate = _snap_supported_rate_up(target_rate, supported_sample_rates)
1206
1207 content_type, bit_depth = self._pick_pcm_bit_depth(
1208 players, start_streamdetails, crossfade_enabled, overlay_active
1209 )
1210 return AudioFormat(
1211 content_type=content_type,
1212 sample_rate=output_sample_rate,
1213 bit_depth=bit_depth,
1214 channels=2,
1215 )
1216
1217 async def get_audio_source_stream(
1218 self,
1219 queue_item: QueueItem,
1220 pcm_format: AudioFormat,
1221 raise_on_error: bool = True,
1222 ) -> AsyncGenerator[bytes]:
1223 """
1224 Get the realtime PCM stream for an AudioSource queue item.
1225
1226 AudioSources are live/realtime: bytes flow at the producer's pace, with
1227 no pre-buffering, no loudness hydration, no volume normalization, no
1228 crossfade/fade-in, no playback-speed shift, no next-track preload. The
1229 path stays as small as possible to keep end-to-end latency low.
1230
1231 Fast path: when the source PCM format already matches the consumer's
1232 ``pcm_format``, the provider's bytes are paced in Python and forwarded
1233 directly â no ffmpeg in the data path.
1234
1235 Slow path: when formats differ, ffmpeg resamples/recodes the stream
1236 (with ``-readrate`` pacing) via ``get_media_stream``.
1237
1238 :param queue_item: The AudioSource queue item to stream.
1239 :param pcm_format: Output PCM format the consumer wants.
1240 :param raise_on_error: Re-raise stream errors instead of swallowing them.
1241 """
1242 streamdetails = queue_item.streamdetails
1243 assert streamdetails
1244 logger = self.logger.getChild("audio_source_stream")
1245 bytes_received = 0
1246 try:
1247 async for chunk in self._iter_audio_source_pcm(streamdetails, pcm_format):
1248 bytes_received += len(chunk)
1249 yield chunk
1250 except AudioError as err:
1251 streamdetails.stream_error = True
1252 # revoke availability when the stream never produced any audio
1253 if bytes_received == 0 and not isinstance(err, ProviderStreamLimitError):
1254 queue_item.available = False
1255 if raise_on_error:
1256 raise
1257 logger.error(
1258 "AudioError while streaming AudioSource %s (%s): %s",
1259 queue_item.name,
1260 streamdetails.uri,
1261 err,
1262 )
1263 except asyncio.CancelledError:
1264 raise
1265 except Exception:
1266 streamdetails.stream_error = True
1267 if raise_on_error:
1268 raise
1269 logger.exception(
1270 "Unexpected error while streaming AudioSource %s (%s)",
1271 queue_item.name,
1272 streamdetails.uri,
1273 )
1274 finally:
1275 streamdetails.seconds_streamed = bytes_received / pcm_format.pcm_sample_size
1276
1277 async def get_queue_item_stream(
1278 self,
1279 queue_item: QueueItem,
1280 pcm_format: AudioFormat,
1281 seek_position: float = 0,
1282 playback_speed: float = 1.0,
1283 raise_on_error: bool = True,
1284 normalization_override: VolumeNormalizationMode | None = None,
1285 session_id: str | None = None,
1286 prepared_buffer: AudioBuffer | None = None,
1287 exact_seek: bool = False,
1288 ) -> AsyncGenerator[bytes]:
1289 """
1290 Get the (PCM) audio stream for a single queue item.
1291
1292 Audio is always served from the AudioBuffer which stores raw decoded PCM.
1293 Volume normalization and other filters are applied on-the-fly when reading
1294 from the buffer.
1295
1296 AudioSource items dispatch to ``get_audio_source_stream`` instead: they
1297 are realtime and bypass the buffering/normalization/filter machinery.
1298
1299 :param normalization_override: Force this volume normalization mode instead of
1300 re-evaluating it from the (possibly just-updated) loudness measurement. Used by
1301 the crossfade path to keep a track's replayed intro and its body on the same mode.
1302 :param session_id: Queue session that owns processing-detail updates.
1303 :param prepared_buffer: Existing buffer that must be used without opening a new source.
1304 :param exact_seek: Preserve millisecond precision instead of user-seek quantization.
1305 """
1306 streamdetails = queue_item.streamdetails
1307 assert streamdetails
1308
1309 # streamdetails are cached and reused for retries; reset this before any
1310 # media-type-specific dispatch so AudioSource failures do not stick.
1311 streamdetails.stream_error = False
1312
1313 if queue_item.media_type == MediaType.AUDIO_SOURCE:
1314 async for chunk in self.get_audio_source_stream(
1315 queue_item=queue_item,
1316 pcm_format=pcm_format,
1317 raise_on_error=raise_on_error,
1318 ):
1319 yield chunk
1320 return
1321 filter_params: list[str] = []
1322
1323 logger = self.logger.getChild("queue_item_stream")
1324
1325 if normalization_override is not None:
1326 # crossfade path pins the body to the intro's mode; skip hydration/re-eval that could flip it
1327 streamdetails.volume_normalization_mode = normalization_override
1328 else:
1329 # hydrate loudness from audio analysis (just-in-time, so that a measurement
1330 # completed during a previous play is picked up here). A live analyzer run
1331 # may have already populated streamdetails.loudness in memory â don't clobber
1332 # that, and don't clobber a value set upstream by the music provider.
1333 if streamdetails.loudness is None:
1334 if analysis := await self.mass.streams.audio_analysis.get_audio_analysis(
1335 streamdetails.item_id,
1336 streamdetails.provider,
1337 media_type=streamdetails.media_type,
1338 # use the authoritative EBU R128 value, not another provider's loudness proxy
1339 priority=(LOUDNESS_ANALYSIS_DOMAIN,),
1340 ):
1341 if analysis.loudness_integrated is not None:
1342 streamdetails.loudness = round(analysis.loudness_integrated, 2)
1343 if analysis.loudness_album is not None and streamdetails.loudness_album is None:
1344 streamdetails.loudness_album = round(analysis.loudness_album, 2)
1345
1346 # re-evaluate normalization mode: the background loudness analyzer may have
1347 # updated streamdetails.loudness since get_stream_details was called
1348 if streamdetails.queue_id:
1349 volume_normalization_enabled = (
1350 self.mass.config.get_effective_player_queue_config_value(
1351 streamdetails.queue_id, CONF_VOLUME_NORMALIZATION, CONF_VALUE_ENABLED
1352 )
1353 != CONF_VALUE_DISABLED
1354 )
1355 streamdetails.volume_normalization_mode = get_normalization_mode(
1356 self._get_volume_normalization_preference(streamdetails),
1357 volume_normalization_enabled,
1358 streamdetails,
1359 )
1360
1361 # get or create the AudioBuffer (stores raw decoded PCM). This runs before the
1362 # filters are built because a source-capacity reselection can hand back another
1363 # provider's streamdetails, which everything below must then work with.
1364 seek_position_ms = int(seek_position * 1000)
1365 try:
1366 if prepared_buffer is not None:
1367 if streamdetails.buffer is not prepared_buffer or not prepared_buffer.is_valid(
1368 seek_position_ms
1369 ):
1370 raise AudioError("Prepared crossfade buffer is no longer available")
1371 audio_buffer = prepared_buffer
1372 else:
1373 audio_buffer = await self.get_audio_buffer(
1374 queue_item, seek_position_ms=seek_position_ms, reason="streaming"
1375 )
1376 except AudioError as err:
1377 streamdetails.stream_error = True
1378 if raise_on_error:
1379 raise
1380 logger.error(
1381 "AudioError while preparing queue item %s (%s): %s",
1382 queue_item.name,
1383 streamdetails.uri,
1384 err,
1385 )
1386 return
1387 streamdetails = queue_item.streamdetails
1388 assert streamdetails # for type checking
1389 if normalization_override is not None:
1390 # a capacity reselection hands back freshly resolved details, so the
1391 # crossfade's intro/body normalization pin must be re-applied to them
1392 streamdetails.volume_normalization_mode = normalization_override
1393
1394 # handle volume normalization
1395 gain_correct: float | None = None
1396 if streamdetails.volume_normalization_mode == VolumeNormalizationMode.DYNAMIC:
1397 filter_rule = (
1398 f"loudnorm=I={streamdetails.target_loudness}"
1399 ":TP=-2.0:LRA=10.0:offset=0.0:print_format=json"
1400 )
1401 filter_params.append(filter_rule)
1402 elif streamdetails.volume_normalization_mode == VolumeNormalizationMode.FIXED_GAIN:
1403 config_key = (
1404 CONF_VOLUME_NORMALIZATION_FIXED_GAIN_TRACKS
1405 if streamdetails.media_type == MediaType.TRACK
1406 else CONF_VOLUME_NORMALIZATION_FIXED_GAIN_RADIO
1407 )
1408 gain_value = self.mass.streams.get_config_value(config_key, return_type=float)
1409 gain_correct = round(gain_value, 2)
1410 filter_params.append(f"volume={gain_correct}dB")
1411 elif streamdetails.volume_normalization_mode == VolumeNormalizationMode.MEASUREMENT_ONLY:
1412 target_loudness = (
1413 float(streamdetails.target_loudness)
1414 if streamdetails.target_loudness is not None
1415 else 0.0
1416 )
1417 if streamdetails.prefer_album_loudness and streamdetails.loudness_album is not None:
1418 gain_correct = target_loudness - float(streamdetails.loudness_album)
1419 elif streamdetails.loudness is not None:
1420 gain_correct = target_loudness - float(streamdetails.loudness)
1421 else:
1422 gain_correct = 0.0
1423 gain_correct = round(gain_correct, 2)
1424 filter_params.append(f"volume={gain_correct}dB")
1425 streamdetails.volume_normalization_gain_correct = gain_correct
1426
1427 # handle playback speed
1428 if playback_speed != 1.0:
1429 filter_params.append(f"atempo={playback_speed}")
1430
1431 # handle optional fade-in
1432 if streamdetails.fade_in:
1433 filter_params.insert(0, "afade=type=in:start_time=0:duration=3")
1434
1435 logger.log(
1436 VERBOSE_LOG_LEVEL,
1437 "Starting queue item stream for %s (%s)"
1438 " - using fade-in: %s"
1439 " - using volume normalization: %s"
1440 " - using playback speed: %s",
1441 queue_item.name,
1442 streamdetails.uri,
1443 streamdetails.fade_in,
1444 streamdetails.volume_normalization_mode,
1445 playback_speed,
1446 )
1447
1448 if (
1449 streamdetails.queue_id
1450 and (queue_data := self.mass.player_queues.queue_data_or_none(streamdetails.queue_id))
1451 and (processing_session_id := session_id or queue_data.session_id)
1452 ):
1453 self.mass.streams.audio_processing.update_item_runtime(
1454 queue_id=streamdetails.queue_id,
1455 session_id=processing_session_id,
1456 queue_item_id=queue_item.queue_item_id,
1457 input_format=audio_buffer.pcm_format,
1458 pcm_format=pcm_format,
1459 normalization=get_normalization_details(streamdetails, gain_correct),
1460 playback_speed=playback_speed,
1461 alters_audio=streamdetails.fade_in,
1462 )
1463 # read from buffer with filters applied (volume normalization, speed, fade-in, etc.)
1464 # if no processing needed, this yields directly from the buffer
1465 media_stream_gen = audio_buffer.get_stream(
1466 output_format=pcm_format,
1467 seek_position_ms=seek_position_ms,
1468 filter_params=filter_params or None,
1469 exact_seek=exact_seek,
1470 )
1471
1472 first_chunk_received = False
1473 bytes_received = 0
1474 finished = False
1475 next_buffer_triggered = False
1476 stream_started_at = asyncio.get_event_loop().time()
1477 try:
1478 async for chunk in media_stream_gen:
1479 bytes_received += len(chunk)
1480 if not first_chunk_received:
1481 first_chunk_received = True
1482 logger.log(
1483 VERBOSE_LOG_LEVEL,
1484 "First audio chunk received for %s (%s) after %.2f seconds",
1485 queue_item.name,
1486 streamdetails.uri,
1487 asyncio.get_event_loop().time() - stream_started_at,
1488 )
1489 # trigger pre-buffering of the next item well before end
1490 # to ensure the raw PCM is ready when the next item needs to be streamed.
1491 # tracks and sound effects are finite files that fill and close immediately;
1492 # live sources (radio, audio_source) open an upstream connection that would
1493 # sit idle and likely time out before the player actually consumes it.
1494 if (
1495 not next_buffer_triggered
1496 and streamdetails.duration
1497 and (queue := self.mass.player_queues.get_active_queue(queue_item.queue_id))
1498 and queue.next_item
1499 and queue.next_item.queue_item_id != queue_item.queue_item_id
1500 and queue.next_item.media_type in (MediaType.TRACK, MediaType.SOUND_EFFECT)
1501 and (bytes_received / pcm_format.pcm_sample_size + seek_position)
1502 >= streamdetails.duration - 60
1503 ):
1504 next_buffer_triggered = True
1505 self.mass.player_queues.prepare_next_audio_buffer(queue_item.queue_id)
1506 yield chunk
1507 del chunk
1508 finished = True
1509 except AudioError as err:
1510 streamdetails.stream_error = True
1511 # revoke availability when the stream never produced any audio
1512 if bytes_received == 0 and not isinstance(err, ProviderStreamLimitError):
1513 queue_item.available = False
1514 if raise_on_error:
1515 raise
1516 logger.error(
1517 "AudioError while streaming queue item %s (%s): %s",
1518 queue_item.name,
1519 streamdetails.uri,
1520 err,
1521 )
1522 except asyncio.CancelledError:
1523 raise
1524 except Exception:
1525 streamdetails.stream_error = True
1526 if raise_on_error:
1527 raise
1528 logger.exception(
1529 "Unexpected error while streaming queue item %s (%s)",
1530 queue_item.name,
1531 streamdetails.uri,
1532 )
1533 finally:
1534 seconds_streamed = bytes_received / pcm_format.pcm_sample_size
1535 streamdetails.seconds_streamed = seconds_streamed
1536 logger.log(
1537 VERBOSE_LOG_LEVEL,
1538 "stream %s for %s in %.2f seconds - seconds streamed/buffered: %.2f",
1539 "aborted" if not finished else "finished",
1540 streamdetails.uri,
1541 asyncio.get_event_loop().time() - stream_started_at,
1542 seconds_streamed,
1543 )
1544 self._notify_provider_streamed(streamdetails, finished, seconds_streamed)
1545
1546 async def get_queue_item_stream_with_smartfade(
1547 self,
1548 player: Player,
1549 queue_item: QueueItem,
1550 pcm_format: AudioFormat,
1551 crossfade_mode: CrossfadeMode = CrossfadeMode.SMART_CROSSFADE,
1552 standard_crossfade_duration: int = 10,
1553 session_id: str | None = None,
1554 ) -> AsyncGenerator[bytes]:
1555 """
1556 Return one queue item with a crossfade into the next item.
1557
1558 :param player: Player consuming the stream.
1559 :param queue_item: Queue item to stream.
1560 :param pcm_format: Shared PCM format.
1561 :param crossfade_mode: Effective crossfade mode.
1562 :param standard_crossfade_duration: Configured standard crossfade duration.
1563 :param session_id: Queue session that owns processing-detail updates.
1564 """
1565 queue = self.mass.player_queues.get(queue_item.queue_id)
1566 if not queue:
1567 raise RuntimeError(f"Queue {queue_item.queue_id} not found")
1568
1569 streamdetails = queue_item.streamdetails
1570 assert streamdetails
1571 crossfade_data = self._crossfade_data.get(queue.queue_id)
1572
1573 if crossfade_data and streamdetails.seek_position > 0:
1574 # don't do crossfade when seeking into track
1575 self.logger.debug(
1576 "Discarding crossfade data for queue %s - seeking into track (pos=%s)",
1577 queue.display_name,
1578 streamdetails.seek_position,
1579 )
1580 crossfade_data = None
1581 if crossfade_data and (crossfade_data.queue_item_id != queue_item.queue_item_id):
1582 # edge case alert: the next item changed just while we were preloading/crossfading
1583 self.logger.warning(
1584 "Skipping crossfade data for queue %s - next item changed!"
1585 " (expected queue_item_id=%s, got=%s)",
1586 queue.display_name,
1587 crossfade_data.queue_item_id,
1588 queue_item.queue_item_id,
1589 )
1590 crossfade_data = None
1591 self._crossfade_data.pop(queue.queue_id, None)
1592 elif not crossfade_data:
1593 self.logger.debug(
1594 "No crossfade data available for queue %s (queue_item_id=%s)",
1595 queue.display_name,
1596 queue_item.queue_item_id,
1597 )
1598
1599 self.logger.debug(
1600 "Start Streaming queue track: %s (%s) for queue %s on player %s"
1601 "- crossfade mode: %s "
1602 "- crossfading from previous track: %s ",
1603 queue_item.streamdetails.uri if queue_item.streamdetails else "Unknown URI",
1604 queue_item.name,
1605 queue.display_name,
1606 player.name,
1607 crossfade_mode,
1608 "true" if crossfade_data else "false",
1609 )
1610
1611 buffer = bytearray()
1612 bytes_written = 0
1613 # calculate crossfade buffer size
1614 crossfade_buffer_duration = (
1615 SMART_CROSSFADE_DURATION
1616 if crossfade_mode == CrossfadeMode.SMART_CROSSFADE
1617 else standard_crossfade_duration
1618 )
1619 crossfade_buffer_duration = min(
1620 crossfade_buffer_duration,
1621 int(streamdetails.duration / 2)
1622 if streamdetails.duration
1623 else crossfade_buffer_duration,
1624 )
1625 # skip crossfade if buffer would be too small to be meaningful
1626 if crossfade_buffer_duration < MIN_CROSSFADE_FALLBACK_DURATION:
1627 crossfade_buffer_duration = 0
1628 # Ensure crossfade buffer size is aligned to frame boundaries
1629 # Frame size = bytes_per_sample * channels
1630 bytes_per_sample = pcm_format.bit_depth // 8
1631 frame_size = bytes_per_sample * pcm_format.channels
1632 crossfade_buffer_size = int(pcm_format.pcm_sample_size * crossfade_buffer_duration)
1633 # Round down to nearest frame boundary
1634 crossfade_buffer_size = (crossfade_buffer_size // frame_size) * frame_size
1635 fade_out_data: bytes | None = None
1636 uncredited_tail_bytes = 0
1637
1638 # pin the body to DYNAMIC when the intro was baked DYNAMIC,
1639 # else a late measurement flips it and causes a volume jump
1640 norm_override: VolumeNormalizationMode | None = None
1641 if crossfade_data and crossfade_data.normalization_mode == VolumeNormalizationMode.DYNAMIC:
1642 norm_override = VolumeNormalizationMode.DYNAMIC
1643
1644 exact_buffer_seek = crossfade_data is not None
1645 if crossfade_data:
1646 # reported media-time (TRIM + CF) is decoupled from the raw buffer seek below (X)
1647 streamdetails.seek_position = crossfade_data.elapsed_time_offset
1648 # yield the POST portion (resample if previous track's format differs)
1649 if crossfade_data.pcm_format != pcm_format:
1650 async for _chunk in resample_pcm_audio(
1651 crossfade_data.data, crossfade_data.pcm_format, pcm_format
1652 ):
1653 yield _chunk
1654 bytes_written += len(_chunk)
1655 else:
1656 for pcm_slice in iter_pcm_slices(crossfade_data.data, pcm_format, 1000):
1657 yield pcm_slice
1658 await asyncio.sleep(0)
1659 bytes_written += len(crossfade_data.data)
1660 # skip past the source media already consumed by the crossfade
1661 discard_position = crossfade_data.fade_in_media_duration
1662 crossfade_data = None
1663 self._crossfade_data.pop(queue.queue_id, None)
1664 else:
1665 discard_position = float(streamdetails.seek_position)
1666
1667 # Yield the first WARMUP_DURATION worth of audio immediately so playback starts
1668 # right away. After that, start accumulating the crossfade holdback buffer.
1669 warmup_size = int(pcm_format.pcm_sample_size * WARMUP_DURATION)
1670 warmup_bytes = 0
1671 total_chunks_received = 0
1672 holdback_armed = False
1673 playback_speed = cast("float", queue_item.extra_attributes.get("playback_speed", 1.0))
1674 async for chunk in self.get_queue_item_stream(
1675 queue_item,
1676 pcm_format,
1677 seek_position=discard_position,
1678 playback_speed=playback_speed,
1679 normalization_override=norm_override,
1680 session_id=session_id,
1681 exact_seek=exact_buffer_seek,
1682 ):
1683 total_chunks_received += 1
1684
1685 if warmup_bytes < warmup_size:
1686 # warmup: yield directly, don't buffer
1687 yield chunk
1688 warmup_bytes += len(chunk)
1689 bytes_written += len(chunk)
1690 del chunk
1691 continue
1692
1693 if not holdback_armed:
1694 holdback_armed = self._crossfade_holdback_allowed(
1695 queue_item.streamdetails or streamdetails,
1696 crossfade_buffer_duration,
1697 playback_speed,
1698 )
1699 if not holdback_armed:
1700 # holding audio back now would only shrink the player's lead
1701 yield chunk
1702 bytes_written += len(chunk)
1703 del chunk
1704 continue
1705
1706 buffer.extend(chunk)
1707 del chunk
1708 if len(buffer) < crossfade_buffer_size:
1709 await asyncio.sleep(0)
1710 continue
1711 # yield everything above the crossfade buffer
1712 while len(buffer) > crossfade_buffer_size:
1713 yield bytes(buffer[: pcm_format.pcm_sample_size])
1714 bytes_written += pcm_format.pcm_sample_size
1715 del buffer[: pcm_format.pcm_sample_size]
1716 await asyncio.sleep(0)
1717
1718 #### HANDLE END OF TRACK
1719
1720 # get next track for crossfade
1721 crossfade_start_time = asyncio.get_event_loop().time()
1722 next_queue_item: QueueItem | None
1723 try:
1724 self.logger.debug(
1725 "Preloading NEXT track for crossfade for queue %s", queue.display_name
1726 )
1727 next_queue_item = await self.mass.player_queues.load_next_queue_item(
1728 queue.queue_id, queue_item.queue_item_id
1729 )
1730 # set index_in_buffer to prevent our next track is overwritten while preloading
1731 if next_queue_item.streamdetails is None:
1732 raise InvalidDataError(
1733 f"No streamdetails for next queue item {next_queue_item.queue_item_id}"
1734 )
1735 queue.index_in_buffer = self.mass.player_queues.index_by_id(
1736 queue.queue_id, next_queue_item.queue_item_id
1737 )
1738 except QueueEmpty:
1739 # end of queue reached, no next item
1740 next_queue_item = None
1741
1742 crossfade_allowed = False
1743 transition_mode = CrossfadeMode.DISABLED
1744 fade_in_buffer_duration = 0.0
1745 fade_in_playback_speed = 1.0
1746 # a fade needs enough of the outgoing track to overlap with; a holdback that
1747 # armed late (or not at all) leaves less than that
1748 min_fade_out_size = int(pcm_format.pcm_sample_size * MIN_CROSSFADE_FALLBACK_DURATION)
1749 if len(buffer) >= min_fade_out_size and next_queue_item and next_queue_item.streamdetails:
1750 fade_in_playback_speed = cast(
1751 "float", next_queue_item.extra_attributes.get("playback_speed", 1.0)
1752 )
1753 next_pcm = await self.select_pcm_format(
1754 player=player,
1755 streamdetails=next_queue_item.streamdetails,
1756 crossfade_enabled=True,
1757 )
1758 crossfade_allowed = self.crossfade_allowed(
1759 queue_item,
1760 crossfade_mode=crossfade_mode,
1761 player_id=player.player_id,
1762 flow_mode=False,
1763 next_queue_item=next_queue_item,
1764 sample_rate=pcm_format.sample_rate,
1765 next_sample_rate=next_pcm.sample_rate,
1766 )
1767 if crossfade_allowed:
1768 transition_mode, fade_in_buffer_duration = self._select_buffered_crossfade(
1769 next_queue_item.streamdetails,
1770 crossfade_mode,
1771 standard_crossfade_duration,
1772 fade_in_playback_speed,
1773 )
1774 crossfade_allowed = transition_mode != CrossfadeMode.DISABLED
1775 if not crossfade_allowed:
1776 # no crossfade enabled/allowed, just yield the buffer last part
1777 bytes_written += len(buffer)
1778 for pcm_slice in iter_pcm_slices(bytes(buffer), pcm_format, 1000):
1779 yield pcm_slice
1780 await asyncio.sleep(0)
1781 else:
1782 assert next_queue_item is not None
1783 assert next_queue_item.streamdetails is not None
1784 assert next_queue_item.streamdetails.buffer is not None
1785 fade_in_audio_buffer = cast("AudioBuffer", next_queue_item.streamdetails.buffer)
1786 # the remaining buffer is the fade-out tail of the current track
1787 fade_out_data = bytes(buffer)
1788 buffer = bytearray()
1789 fade_in_buffer_size = int(pcm_format.pcm_sample_size * fade_in_buffer_duration)
1790 fade_in_buffer_size = (fade_in_buffer_size // frame_size) * frame_size
1791 # initialized before the try block â the except handler reads these
1792 first_part_written = 0
1793 second_part_buf = bytearray()
1794 try:
1795 # wrap the next track's stream in a counting generator that caps
1796 # at the resident fade-in size and tracks how many bytes were consumed
1797 fade_in_bytes_consumed = 0
1798
1799 _next_item = next_queue_item
1800
1801 async def _limited_fade_in() -> AsyncGenerator[bytes]:
1802 nonlocal fade_in_bytes_consumed
1803 fade_in_stream = self.get_queue_item_stream(
1804 _next_item,
1805 pcm_format,
1806 playback_speed=fade_in_playback_speed,
1807 session_id=session_id,
1808 prepared_buffer=fade_in_audio_buffer,
1809 )
1810 async with aclosing(fade_in_stream):
1811 async for chunk in fade_in_stream:
1812 remaining = fade_in_buffer_size - fade_in_bytes_consumed
1813 if remaining <= 0:
1814 break
1815 if len(chunk) >= remaining:
1816 fade_in_bytes_consumed += remaining
1817 yield chunk[:remaining]
1818 break
1819 fade_in_bytes_consumed += len(chunk)
1820 yield chunk
1821
1822 smart_fade = await self.smart_fades_mixer.build(
1823 fade_in_streamdetails=next_queue_item.streamdetails,
1824 fade_out_streamdetails=streamdetails,
1825 pcm_format=pcm_format,
1826 standard_crossfade_duration=standard_crossfade_duration,
1827 mode=transition_mode,
1828 fade_out_data=fade_out_data,
1829 fade_in_bytes_len=fade_in_buffer_size,
1830 )
1831 crossfade_timing = smart_fade.timing_info
1832 # Split mix output at end-of-overlap: PRE+CF to A, POST to B's intro.
1833 fadeout_share_bytes = int(
1834 (crossfade_timing.pre_crossfade_duration + crossfade_timing.crossfade_duration)
1835 * pcm_format.pcm_sample_size
1836 )
1837 fadeout_share_bytes = (fadeout_share_bytes // frame_size) * frame_size
1838 async for mix_chunk in self.smart_fades_mixer.mix(
1839 smart_fade,
1840 fade_in_part=_limited_fade_in(),
1841 fade_out_part=fade_out_data,
1842 pcm_format=pcm_format,
1843 ):
1844 if first_part_written < fadeout_share_bytes:
1845 # split this chunk so A gets exactly fadeout_share_bytes
1846 remaining = fadeout_share_bytes - first_part_written
1847 if len(mix_chunk) > remaining:
1848 yield mix_chunk[:remaining]
1849 first_part_written += remaining
1850 bytes_written += remaining
1851 second_part_buf.extend(mix_chunk[remaining:])
1852 else:
1853 yield mix_chunk
1854 first_part_written += len(mix_chunk)
1855 bytes_written += len(mix_chunk)
1856 else:
1857 second_part_buf.extend(mix_chunk)
1858 # tail consumed by the mix but not credited to bytes_written
1859 uncredited_tail_bytes = len(fade_out_data) - first_part_written
1860 self._crossfade_data[queue_item.queue_id] = CrossfadeData(
1861 data=bytes(second_part_buf),
1862 fade_in_media_duration=(fade_in_bytes_consumed / pcm_format.pcm_sample_size)
1863 * fade_in_playback_speed,
1864 pcm_format=pcm_format,
1865 queue_item_id=next_queue_item.queue_item_id,
1866 elapsed_time_offset=(
1867 crossfade_timing.fadein_trimmed_duration
1868 + crossfade_timing.crossfade_duration
1869 )
1870 * fade_in_playback_speed,
1871 normalization_mode=next_queue_item.streamdetails.volume_normalization_mode,
1872 )
1873 crossfade_elapsed = asyncio.get_event_loop().time() - crossfade_start_time
1874 self.logger.debug(
1875 "Stored crossfade data for queue %s"
1876 " - next queue_item_id: %s (preparation took %.1fs)",
1877 queue.display_name,
1878 next_queue_item.queue_item_id,
1879 crossfade_elapsed,
1880 )
1881 except Exception as err:
1882 if first_part_written or second_part_buf:
1883 # partial mix already played â concat'd fade_out_data would duplicate audio
1884 raise
1885 # crossfade failed, fall back to just yielding the fade_out_data
1886 self.logger.warning(
1887 "Crossfade failed for queue %s: %s",
1888 queue.display_name,
1889 err,
1890 )
1891 next_queue_item = None
1892 for pcm_slice in iter_pcm_slices(fade_out_data, pcm_format, 1000):
1893 yield pcm_slice
1894 await asyncio.sleep(0)
1895 bytes_written += len(fade_out_data)
1896 del fade_out_data
1897 # make sure the buffer gets cleaned up
1898 del buffer
1899 # a capacity reselection inside the stream replaces the queue item's details,
1900 # so rebind before the writebacks land on an orphaned object
1901 streamdetails = queue_item.streamdetails or streamdetails
1902 # update duration details based on the actual pcm data we sent
1903 # this also accounts for crossfade and silence stripping
1904 seconds_streamed = bytes_written / pcm_format.pcm_sample_size
1905 streamdetails.seconds_streamed = seconds_streamed
1906 # an externally aborted source ends in a clean EOF mid-track, so the
1907 # streamed length must not be written back as the item's duration
1908 source_buffer = streamdetails.buffer
1909 if source_buffer is None or not source_buffer.cancelled:
1910 uncredited_tail_seconds = uncredited_tail_bytes / pcm_format.pcm_sample_size
1911 # streamdetails.duration is in media-time; seconds_streamed is stream-time
1912 # (post-atempo), so we scale by playback_speed to recover media-time.
1913 streamdetails.duration = int(
1914 streamdetails.seek_position
1915 + (seconds_streamed + uncredited_tail_seconds) * playback_speed
1916 )
1917 # propagate accurate duration to queue_item so UI displays it
1918 queue_item.duration = streamdetails.duration
1919 self.logger.debug(
1920 "Finished Streaming queue track: %s (%s) on queue %s "
1921 "- crossfade data prepared for next track: %s",
1922 streamdetails.uri,
1923 queue_item.name,
1924 queue.display_name,
1925 (
1926 next_queue_item.name
1927 if next_queue_item and queue_item.queue_id in self._crossfade_data
1928 else "N/A"
1929 ),
1930 )
1931
1932 async def get_queue_flow_stream(
1933 self,
1934 queue: PlayerQueue,
1935 start_queue_item: QueueItem,
1936 pcm_format: AudioFormat,
1937 session_id: str | None = None,
1938 protocol_player: Player | None = None,
1939 ) -> AsyncGenerator[bytes]:
1940 """
1941 Get a flow stream of all tracks in the queue as raw PCM audio.
1942
1943 yields chunks of exactly 1 second of audio in the given pcm_format.
1944
1945 :param queue: Queue being streamed.
1946 :param start_queue_item: First queue item in the flow stream.
1947 :param pcm_format: Shared PCM format for the complete flow stream.
1948 :param session_id: Queue session that owns processing-detail updates.
1949 :param protocol_player: The protocol player actually consuming the flow stream.
1950 Must be the same player that was used to select ``pcm_format`` so
1951 restart decisions are made against the correct supported sample rates
1952 and flow mode configuration. Falls back to the queue's player when omitted.
1953 """
1954 # ruff: noqa: PLR0915
1955 assert pcm_format.content_type.is_pcm()
1956 queue_track = None
1957 last_fadeout_part: bytes = b""
1958 last_streamdetails: StreamDetails | None = None
1959 last_play_log_entry: PlayLogEntry | None = None
1960 # Snapshot the queue's current session_id. PlayerQueues rotates this on
1961 # every new stream session, so if a newer producer takes over the queue
1962 # (rapid track switch, sync-group reform, dynamic leader handoff) the
1963 # snapshot will no longer match and we exit cleanly on the next yield or
1964 # playlog append â preventing two producers from writing to the same
1965 # pq_data.flow_mode_stream_log.
1966 pq_data = self.mass.player_queues.queue_data(queue.queue_id)
1967 flow_session_id = session_id or pq_data.session_id
1968 if flow_session_id is None or pq_data.session_id != flow_session_id:
1969 self.logger.debug(
1970 "Ignoring stale flow stream for queue %s (session %s, active %s)",
1971 queue.display_name,
1972 flow_session_id,
1973 pq_data.session_id,
1974 )
1975 return
1976 queue.flow_mode = True
1977 # A session can also be handed a second producer, which the session check does not
1978 # catch: players such as DLNA renderers sometimes open the same flow url twice to
1979 # probe the audio. Append to the list published here rather than to whatever the
1980 # queue currently holds, so the entries of a producer that has since been replaced
1981 # end up in a list nobody reads instead of interleaving with the live one's.
1982 flow_log: list[PlayLogEntry] = []
1983 pq_data.flow_mode_stream_log = flow_log
1984 if not start_queue_item:
1985 # this can happen in some (edge case) race conditions
1986 return
1987 pcm_sample_size = pcm_format.pcm_sample_size
1988 if start_queue_item.media_type != MediaType.TRACK:
1989 # no crossfade on non-tracks
1990 crossfade_mode = CrossfadeMode.DISABLED
1991 standard_crossfade_duration = 0
1992 else:
1993 crossfade_mode = self.mass.streams.get_crossfade_mode(queue)
1994 # crossfade duration is a global (queue controller) setting; fallback matches
1995 # CONF_ENTRY_CROSSFADE_DURATION's default
1996 standard_crossfade_duration = self.mass.config.get_raw_core_config_value(
1997 CONF_PLAYER_QUEUES, CONF_CROSSFADE_DURATION, 8
1998 )
1999 flow_mode_sample_rate_conf, flow_supported_sample_rates = self._flow_restart_context(
2000 queue.queue_id, protocol_player
2001 )
2002 # note: get_crossfade_mode() already falls back to standard when smart fades aren't
2003 # available (no analysis provider / minimal buffer), so crossfade_mode is safe to use.
2004 self.logger.info(
2005 "Start Queue Flow stream for Queue %s - crossfade: %s %s",
2006 queue.display_name,
2007 crossfade_mode,
2008 f"({standard_crossfade_duration}s)"
2009 if crossfade_mode == CrossfadeMode.STANDARD_CROSSFADE
2010 else "",
2011 )
2012 total_chunks_received = 0
2013
2014 def _superseded() -> bool:
2015 """Return True if a newer stream session has taken over this queue."""
2016 return pq_data.session_id != flow_session_id
2017
2018 queue_exhausted = False
2019 while True:
2020 # bail out early if a newer producer has taken over this queue,
2021 # so we don't append another entry to a stream log we no longer own
2022 if _superseded():
2023 self.logger.debug(
2024 "Flow stream for queue %s superseded (session %s -> %s) "
2025 "- exiting before next track",
2026 queue.display_name,
2027 flow_session_id,
2028 pq_data.session_id,
2029 )
2030 return
2031 # get (next) queue item to stream
2032 if queue_track is None:
2033 queue_track = start_queue_item
2034 else:
2035 try:
2036 queue_track = await self.mass.player_queues.load_next_queue_item(
2037 queue.queue_id, queue_track.queue_item_id
2038 )
2039 except QueueEmpty:
2040 queue_exhausted = True
2041 break
2042
2043 if self._flow_stream_needs_restart(
2044 queue_track,
2045 pcm_format,
2046 flow_supported_sample_rates,
2047 flow_mode_sample_rate_conf,
2048 is_first_track=queue_track is start_queue_item,
2049 ):
2050 break
2051
2052 if queue_track.streamdetails is None:
2053 self.logger.error(
2054 "No StreamDetails for queue item %s (%s) on queue %s - skipping track",
2055 queue_track.queue_item_id,
2056 queue_track.name,
2057 queue.display_name,
2058 )
2059 continue
2060 # a realtime source delivers at playback pace, so it has no audio to spare
2061 # for an overlap in either direction
2062 item_crossfade_mode = (
2063 CrossfadeMode.DISABLED if queue_track.streamdetails.is_realtime else crossfade_mode
2064 )
2065 if flow_session_id is not None:
2066 self.mass.streams.audio_processing.update_item_context(
2067 queue_id=queue.queue_id,
2068 session_id=flow_session_id,
2069 queue_item_id=queue_track.queue_item_id,
2070 queue_processing=AudioQueueProcessing(
2071 pcm_format=pcm_format,
2072 playback_speed=cast(
2073 "float",
2074 queue_track.extra_attributes.get("playback_speed", 1.0),
2075 ),
2076 crossfade_mode=item_crossfade_mode,
2077 overlay_active=overlay_active(queue),
2078 ),
2079 alters_audio=queue_track.streamdetails.fade_in,
2080 )
2081
2082 self.logger.debug(
2083 "Start Streaming queue track: %s (%s) for queue %s",
2084 queue_track.streamdetails.uri,
2085 queue_track.name,
2086 queue.display_name,
2087 )
2088 # last chance to bail before mutating the stream log: a newer producer
2089 # may have taken over while we were awaiting load_next_queue_item
2090 if _superseded():
2091 self.logger.debug(
2092 "Flow stream for queue %s superseded - exiting before playlog append",
2093 queue.display_name,
2094 )
2095 return
2096 track_playback_speed = cast(
2097 "float", queue_track.extra_attributes.get("playback_speed", 1.0)
2098 )
2099 # calculate crossfade buffer size
2100 crossfade_buffer_duration = (
2101 SMART_CROSSFADE_DURATION
2102 if item_crossfade_mode == CrossfadeMode.SMART_CROSSFADE
2103 else standard_crossfade_duration
2104 )
2105 crossfade_buffer_duration = min(
2106 crossfade_buffer_duration,
2107 int(queue_track.streamdetails.duration / 2)
2108 if queue_track.streamdetails.duration
2109 else crossfade_buffer_duration,
2110 )
2111 # skip crossfade if buffer would be too small to be meaningful
2112 if crossfade_buffer_duration < MIN_CROSSFADE_FALLBACK_DURATION:
2113 crossfade_buffer_duration = 0
2114 # Ensure crossfade buffer size is aligned to frame boundaries
2115 # Frame size = bytes_per_sample * channels
2116 bytes_per_sample = pcm_format.bit_depth // 8
2117 frame_size = bytes_per_sample * pcm_format.channels
2118 crossfade_buffer_size = int(pcm_format.pcm_sample_size * crossfade_buffer_duration)
2119 # Round down to nearest frame boundary
2120 crossfade_buffer_size = (crossfade_buffer_size // frame_size) * frame_size
2121 warmup_size = int(pcm_format.pcm_sample_size * WARMUP_DURATION)
2122
2123 # raw_seek_position feeds the PCM buffer; streamdetails.seek_position
2124 # (overwritten below) only drives reported elapsed time.
2125 raw_seek_position = queue_track.streamdetails.seek_position
2126 # Build eagerly so seek_position is set before PlayLogEntry is appended â
2127 # consumer-paced mix() would otherwise let the queue briefly report 0.
2128 crossfade_smart_fade: SmartFade | None = None
2129 collect_started = 0.0
2130 collect_resident = 0.0
2131 incoming_crossfade_size = crossfade_buffer_size
2132 incoming_audio_buffer: AudioBuffer | None = None
2133 if last_fadeout_part and last_streamdetails:
2134 transition_mode = CrossfadeMode.DISABLED
2135 incoming_duration = 0.0
2136 if crossfade_buffer_size > 0 and item_crossfade_mode != CrossfadeMode.DISABLED:
2137 transition_mode, incoming_duration = self._select_buffered_crossfade(
2138 queue_track.streamdetails,
2139 item_crossfade_mode,
2140 standard_crossfade_duration,
2141 track_playback_speed,
2142 )
2143 if transition_mode == CrossfadeMode.DISABLED:
2144 # nothing to fade into: flush the held-back tail of the previous track
2145 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2146 yield pcm_slice
2147 await asyncio.sleep(0)
2148 last_fadeout_part = b""
2149 last_streamdetails = None
2150 last_play_log_entry = None
2151 else:
2152 assert queue_track.streamdetails.buffer is not None
2153 incoming_audio_buffer = cast("AudioBuffer", queue_track.streamdetails.buffer)
2154 incoming_crossfade_size = int(pcm_format.pcm_sample_size * incoming_duration)
2155 incoming_crossfade_size = (incoming_crossfade_size // frame_size) * frame_size
2156 # no audio is emitted while the incoming overlap is collected, so a
2157 # slow collection here is heard as a gap at the transition
2158 collect_started = asyncio.get_event_loop().time()
2159 collect_resident = incoming_audio_buffer.duration_available
2160 crossfade_smart_fade = await self.smart_fades_mixer.build(
2161 fade_in_streamdetails=queue_track.streamdetails,
2162 fade_out_streamdetails=last_streamdetails,
2163 pcm_format=pcm_format,
2164 standard_crossfade_duration=standard_crossfade_duration,
2165 mode=transition_mode,
2166 fade_out_data=last_fadeout_part,
2167 fade_in_bytes_len=incoming_crossfade_size,
2168 )
2169 timing_info = crossfade_smart_fade.timing_info
2170 if isinstance(crossfade_smart_fade, StandardCrossFade):
2171 # A standard fade blends its overlap and passes everything after it
2172 # through untouched, so only the overlap has to be in hand before
2173 # the transition can start. Holding back the rest buys nothing and
2174 # keeps the player waiting - a smart fade does need its full window,
2175 # which is only chosen when the analysis it needs is already there.
2176 blended_seconds = (
2177 timing_info.fadein_trimmed_duration + timing_info.crossfade_duration
2178 )
2179 blended_size = int(pcm_format.pcm_sample_size * blended_seconds)
2180 incoming_crossfade_size = min(
2181 incoming_crossfade_size,
2182 (blended_size // frame_size) * frame_size,
2183 )
2184 queue_track.streamdetails.seek_position = (
2185 raw_seek_position
2186 + (timing_info.fadein_trimmed_duration + timing_info.crossfade_duration)
2187 * track_playback_speed
2188 )
2189 # append to play log so the queue controller can work out which track is playing
2190 play_log_entry = PlayLogEntry(queue_track.queue_item_id)
2191 flow_log.append(play_log_entry)
2192
2193 bytes_written = 0
2194 crossfade_buffer = bytearray()
2195 warmup_bytes = 0
2196 first_chunk_received = False
2197 holdback_armed = False
2198
2199 async for chunk in self.get_queue_item_stream(
2200 queue_track,
2201 pcm_format=pcm_format,
2202 seek_position=int(raw_seek_position),
2203 playback_speed=cast(
2204 "float", queue_track.extra_attributes.get("playback_speed", 1.0)
2205 ),
2206 raise_on_error=False,
2207 session_id=flow_session_id,
2208 prepared_buffer=incoming_audio_buffer,
2209 ):
2210 # if a newer producer has taken over this queue, stop sending
2211 # audio and exit cleanly before the outer-loop end-of-track
2212 # bookkeeping mutates seconds_streamed / duration on the log
2213 if _superseded():
2214 self.logger.debug(
2215 "Flow stream for queue %s superseded - stopping chunk yield",
2216 queue.display_name,
2217 )
2218 return
2219 total_chunks_received += 1
2220 if not first_chunk_received:
2221 first_chunk_received = True
2222 # inform the queue that the track is now loaded in the buffer
2223 # so the next track can be preloaded
2224 self.mass.player_queues.track_loaded_in_buffer(
2225 queue.queue_id, queue_track.queue_item_id
2226 )
2227
2228 if item_crossfade_mode == CrossfadeMode.DISABLED:
2229 # no cross/smart fade: yield chunks directly without intermediate buffer
2230 yield chunk
2231 bytes_written += len(chunk)
2232 del chunk
2233 continue
2234
2235 # Warmup: yield chunks directly until we have streamed WARMUP_DURATION
2236 # worth of audio, so playback starts immediately. Skip warmup when
2237 # crossfade data from the previous track is pending â we need a full
2238 # buffer for the mix.
2239 if warmup_bytes < warmup_size and not last_fadeout_part:
2240 yield chunk
2241 warmup_bytes += len(chunk)
2242 bytes_written += len(chunk)
2243 del chunk
2244 continue
2245
2246 if not last_fadeout_part and not holdback_armed:
2247 holdback_armed = self._crossfade_holdback_allowed(
2248 queue_track.streamdetails,
2249 crossfade_buffer_duration,
2250 track_playback_speed,
2251 )
2252 if not holdback_armed:
2253 # holding audio back now would only shrink the player's lead
2254 yield chunk
2255 bytes_written += len(chunk)
2256 del chunk
2257 continue
2258
2259 # smart fades enabled: accumulate chunks in crossfade buffer
2260 crossfade_buffer.extend(chunk)
2261 del chunk
2262 required_buffer_size = (
2263 incoming_crossfade_size if last_fadeout_part else crossfade_buffer_size
2264 )
2265 if len(crossfade_buffer) < required_buffer_size:
2266 await asyncio.sleep(0)
2267 continue
2268 # handle crossfade of previous track and new track
2269 if (
2270 last_fadeout_part
2271 and last_streamdetails
2272 and crossfade_smart_fade is not None
2273 and last_play_log_entry is not None
2274 ):
2275 self.logger.debug(
2276 "Collected %.1fs of incoming audio for the transition into %s in %.1fs"
2277 " (%.1fs was resident when it started)",
2278 len(crossfade_buffer) / pcm_sample_size,
2279 queue_track.name,
2280 asyncio.get_event_loop().time() - collect_started,
2281 collect_resident,
2282 )
2283 fadein_part = bytes(crossfade_buffer[:incoming_crossfade_size])
2284 remaining_bytes = bytes(crossfade_buffer[incoming_crossfade_size:])
2285 try:
2286 crossfade_bytes_written = 0
2287 async for mix_chunk in self.smart_fades_mixer.mix(
2288 crossfade_smart_fade,
2289 fade_in_part=fadein_part,
2290 fade_out_part=last_fadeout_part,
2291 pcm_format=pcm_format,
2292 ):
2293 yield mix_chunk
2294 crossfade_bytes_written += len(mix_chunk)
2295 except Exception as mix_err:
2296 if crossfade_bytes_written:
2297 # partial mix already played â concat'd tail would duplicate audio
2298 raise
2299 self.logger.warning(
2300 "Crossfade mixer failed for %s, falling back to simple concat: %s",
2301 queue_track.name,
2302 mix_err,
2303 )
2304 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2305 yield pcm_slice
2306 await asyncio.sleep(0)
2307 # full tail was pre-counted and is now yielded as-is
2308 crossfade_bytes_written = 0
2309 remaining_bytes = bytes(crossfade_buffer)
2310 # mix failed â undo the eager seek_position
2311 queue_track.streamdetails.seek_position = raw_seek_position
2312 if crossfade_bytes_written:
2313 # Split mix output at end-of-overlap: PRE+CF to A, POST to B.
2314 fadeout_share_seconds = (
2315 timing_info.pre_crossfade_duration + timing_info.crossfade_duration
2316 )
2317 fadeout_share = int(fadeout_share_seconds * pcm_sample_size)
2318 fadeout_share = (fadeout_share // frame_size) * frame_size
2319 fadeout_share = min(fadeout_share, crossfade_bytes_written)
2320 fadein_share = crossfade_bytes_written - fadeout_share
2321 bytes_written += fadein_share
2322 if last_play_log_entry:
2323 assert last_play_log_entry.seconds_streamed is not None
2324 # correct pre-counted full tail to the timing-based share
2325 last_play_log_entry.seconds_streamed += (
2326 fadeout_share - len(last_fadeout_part)
2327 ) / pcm_sample_size
2328 if remaining_bytes:
2329 for pcm_slice in iter_pcm_slices(remaining_bytes, pcm_format, 1000):
2330 yield pcm_slice
2331 await asyncio.sleep(0)
2332 bytes_written += len(remaining_bytes)
2333 del remaining_bytes
2334 last_fadeout_part = b""
2335 last_streamdetails = None
2336 crossfade_buffer = bytearray()
2337 warmup_bytes = 0
2338
2339 # yield everything above the crossfade buffer size
2340 while len(crossfade_buffer) > crossfade_buffer_size:
2341 yield bytes(crossfade_buffer[:pcm_sample_size])
2342 bytes_written += pcm_sample_size
2343 del crossfade_buffer[:pcm_sample_size]
2344 await asyncio.sleep(0)
2345
2346 # A source error after partial audio must not look like a completed item.
2347 # Progress reporting skips items with stream_error, so the item is not
2348 # marked played; move on to the next queue item like the zero-audio path.
2349 if first_chunk_received and queue_track.streamdetails.stream_error:
2350 if _superseded():
2351 return
2352 self.logger.warning(
2353 "Track %s (%s) on queue %s aborted by a stream error - skipping",
2354 queue_track.name,
2355 queue_track.streamdetails.uri,
2356 queue.display_name,
2357 )
2358 # the audio sent so far will still play out; keep the play log entry
2359 # honest about how much of this item was actually streamed
2360 play_log_entry.seconds_streamed = bytes_written / pcm_sample_size
2361 if last_fadeout_part:
2362 # crossfade into this item never happened â undo the eager seek_position
2363 queue_track.streamdetails.seek_position = raw_seek_position
2364 continue
2365
2366 #### HANDLE END OF TRACK
2367 if not first_chunk_received:
2368 self.logger.warning(
2369 "Track %s (%s) on queue %s produced no audio data - skipping",
2370 queue_track.name,
2371 queue_track.streamdetails.uri if queue_track.streamdetails else "unknown",
2372 queue.display_name,
2373 )
2374 queue_track.streamdetails.stream_error = True
2375 play_log_entry.seconds_streamed = 0
2376 if last_fadeout_part:
2377 queue_track.streamdetails.seek_position = raw_seek_position
2378 continue
2379 if last_fadeout_part:
2380 # edge case: we did not get enough data to make the crossfade
2381 # attribute these bytes to the previous track (they are its tail)
2382 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2383 yield pcm_slice
2384 await asyncio.sleep(0)
2385 # no crossfade happened â undo the eager seek_position
2386 queue_track.streamdetails.seek_position = raw_seek_position
2387 # full tail was pre-counted and is now yielded as-is
2388 last_fadeout_part = b""
2389 # a fade needs enough of the outgoing track to overlap with; a holdback that
2390 # armed late (or not at all) leaves less than that
2391 min_fade_out_size = int(pcm_sample_size * MIN_CROSSFADE_FALLBACK_DURATION)
2392 if len(crossfade_buffer) >= min_fade_out_size and self.crossfade_allowed(
2393 queue_track,
2394 crossfade_mode=item_crossfade_mode,
2395 player_id=queue.queue_id,
2396 flow_mode=True,
2397 ):
2398 last_fadeout_part = bytes(crossfade_buffer[-crossfade_buffer_size:])
2399 last_streamdetails = queue_track.streamdetails
2400 last_play_log_entry = play_log_entry
2401 remaining_bytes = bytes(crossfade_buffer[:-crossfade_buffer_size])
2402 if remaining_bytes:
2403 for pcm_slice in iter_pcm_slices(remaining_bytes, pcm_format, 1000):
2404 yield pcm_slice
2405 await asyncio.sleep(0)
2406 bytes_written += len(remaining_bytes)
2407 del remaining_bytes
2408 elif item_crossfade_mode != CrossfadeMode.DISABLED and crossfade_buffer:
2409 bytes_written += len(crossfade_buffer)
2410 for pcm_slice in iter_pcm_slices(bytes(crossfade_buffer), pcm_format, 1000):
2411 yield pcm_slice
2412 await asyncio.sleep(0)
2413 crossfade_buffer = bytearray()
2414
2415 # update duration details based on the actual pcm data we sent
2416 # this also accounts for crossfade and silence stripping
2417 seconds_streamed = bytes_written / pcm_sample_size
2418 queue_track.streamdetails.seconds_streamed = seconds_streamed
2419 play_log_entry.seconds_streamed = seconds_streamed
2420 # an externally aborted source ends in a clean EOF mid-track, so the
2421 # streamed length must not be written back as the item's duration
2422 source_buffer = queue_track.streamdetails.buffer
2423 source_aborted = source_buffer is not None and source_buffer.cancelled
2424 if not source_aborted:
2425 # the held-back crossfade tail still counts as this track's media-time
2426 tail_seconds = len(last_fadeout_part) / pcm_sample_size
2427 # streamdetails.duration is in media-time; seconds_streamed is stream-time
2428 # (post-atempo), so we scale by the track's playback_speed to recover media-time.
2429 queue_track.streamdetails.duration = int(
2430 queue_track.streamdetails.seek_position
2431 + (seconds_streamed + tail_seconds) * track_playback_speed
2432 )
2433 # propagate accurate duration to queue_item so UI displays it
2434 queue_track.duration = queue_track.streamdetails.duration
2435 play_log_entry.duration = queue_track.streamdetails.duration
2436 if last_play_log_entry is play_log_entry and last_fadeout_part:
2437 # Pre-count the full crossfade tail so the queue index calculation
2438 # doesn't undercount while waiting for the next track's crossfade mix.
2439 # This will be corrected to crossfade_total/2 once the mix completes.
2440 assert play_log_entry.seconds_streamed is not None
2441 play_log_entry.seconds_streamed += len(last_fadeout_part) / pcm_sample_size
2442 self.logger.debug(
2443 "Finished Streaming queue track: %s (%s) on queue %s",
2444 queue_track.streamdetails.uri,
2445 queue_track.name,
2446 queue.display_name,
2447 )
2448 #### HANDLE END OF QUEUE FLOW STREAM
2449 # skip end-of-queue bookkeeping if a newer producer has superseded us;
2450 # the new producer owns queue_buffer_completed and the play log now
2451 if _superseded():
2452 self.logger.debug(
2453 "Flow stream for queue %s superseded - skipping end-of-queue handling",
2454 queue.display_name,
2455 )
2456 return
2457 # end of queue flow: make sure we yield the last_fadeout_part
2458 if last_fadeout_part:
2459 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2460 yield pcm_slice
2461 await asyncio.sleep(0)
2462 # correct seconds streamed - the duration already includes the tail
2463 last_part_seconds = len(last_fadeout_part) / pcm_sample_size
2464 streamdetails = queue_track.streamdetails
2465 assert streamdetails is not None
2466 streamdetails.seconds_streamed = (
2467 streamdetails.seconds_streamed or 0
2468 ) + last_part_seconds
2469 # also update the play log entry so elapsed time tracking stays in sync
2470 if last_play_log_entry:
2471 assert last_play_log_entry.seconds_streamed is not None
2472 # full tail was pre-counted and is now yielded as-is
2473 last_play_log_entry.duration = streamdetails.duration
2474 last_fadeout_part = b""
2475 self.logger.info("Finished Queue Flow stream for Queue %s", queue.display_name)
2476 # only signal completion if we are still the active producer â a later
2477 # producer would (incorrectly) see this as its own completion otherwise
2478 if not _superseded():
2479 # inform the queue controller that all audio data has been generated
2480 # so it can handle the case where new items were added after the flow stream ended
2481 self.mass.player_queues.queue_buffer_completed(queue.queue_id, queue_exhausted)
2482
2483 async def get_overlay_mixed_stream(
2484 self,
2485 queue: PlayerQueue,
2486 audio_input: AsyncGenerator[bytes],
2487 pcm_format: AudioFormat,
2488 ) -> AsyncGenerator[bytes]:
2489 """
2490 Mix the queue's audio overlay (looping sound effect) into the given PCM stream.
2491
2492 The mixed output has the exact same PCM format, duration and chunking as the
2493 input stream. If the overlay source can not be resolved, the original stream
2494 is passed through unchanged so playback is never interrupted.
2495
2496 :param queue: The PlayerQueue holding the overlay source and volume.
2497 :param audio_input: The audio stream (raw PCM in ``pcm_format``) to mix into.
2498 :param pcm_format: PCM format of both the input and the mixed output.
2499 """
2500 overlay_input = await self._resolve_overlay_input(queue)
2501 if overlay_input is None:
2502 # overlay source unavailable: degrade gracefully to music-only
2503 async for chunk in audio_input:
2504 yield chunk
2505 return
2506 async for chunk in get_ffmpeg_overlay_stream(
2507 audio_input=audio_input,
2508 overlay_input=overlay_input,
2509 pcm_format=pcm_format,
2510 overlay_volume=queue.overlay_volume,
2511 chunk_size=pcm_format.pcm_sample_size,
2512 ):
2513 yield chunk
2514
2515 def crossfade_allowed(
2516 self,
2517 queue_item: QueueItem,
2518 crossfade_mode: CrossfadeMode,
2519 player_id: str,
2520 flow_mode: bool = False,
2521 next_queue_item: QueueItem | None = None,
2522 sample_rate: int | None = None,
2523 next_sample_rate: int | None = None,
2524 ) -> bool:
2525 """Get the crossfade config for a queue item."""
2526 if crossfade_mode == CrossfadeMode.DISABLED:
2527 return False
2528 if not (self.mass.player_queues.get(queue_item.queue_id)):
2529 return False # just a guard
2530 if not (self.mass.players.get_player(player_id)):
2531 return False # just a guard
2532 if queue_item.media_type != MediaType.TRACK:
2533 self.logger.debug("Skipping crossfade: current item is not a track")
2534 return False
2535 # check if the next item is part of the same album
2536 next_item = next_queue_item or self.mass.player_queues.get_next_item(
2537 queue_item.queue_id, queue_item.queue_item_id
2538 )
2539 if not next_item:
2540 # there is no next item!
2541 return False
2542 # check if next item is a track
2543 if next_item.media_type != MediaType.TRACK:
2544 self.logger.debug("Skipping crossfade: next item is not a track")
2545 return False
2546 if (
2547 isinstance(queue_item.media_item, Track)
2548 and isinstance(next_item.media_item, Track)
2549 and queue_item.media_item.album
2550 and next_item.media_item.album
2551 and queue_item.media_item.album == next_item.media_item.album
2552 and not self.mass.config.get_raw_core_config_value(
2553 "streams", CONF_ALLOW_CROSSFADE_SAME_ALBUM, False
2554 )
2555 ):
2556 # in general, crossfade is not desired for tracks of the same (gapless) album
2557 # because we have no accurate way to determine if the album is gapless or not,
2558 # for now we just never crossfade between tracks of the same album
2559 self.logger.debug("Skipping crossfade: next item is part of the same album")
2560 return False
2561
2562 # check if we're allowed to crossfade on different sample rates
2563 if (
2564 not flow_mode
2565 and sample_rate
2566 and next_sample_rate
2567 and sample_rate != next_sample_rate
2568 and not self.mass.config.get_raw_player_config_value(
2569 player_id,
2570 CONF_ENTRY_CROSSFADE_DIFFERENT_SAMPLE_RATES.key,
2571 CONF_ENTRY_CROSSFADE_DIFFERENT_SAMPLE_RATES.default_value,
2572 )
2573 ):
2574 self.logger.debug(
2575 "Skipping crossfade: player(protocol) does not support gapless playback "
2576 "with different sample rates (%s vs %s)",
2577 sample_rate,
2578 next_sample_rate,
2579 )
2580 return False
2581
2582 return True
2583
2584 def clear_crossfade_data(self, queue_id: str) -> None:
2585 """
2586 Clear any pending crossfade data for a queue.
2587
2588 :param queue_id: The queue ID to clear crossfade data for.
2589 """
2590 if queue_id in self._crossfade_data:
2591 self.logger.debug("Clearing crossfade data for queue %s", queue_id)
2592 del self._crossfade_data[queue_id]
2593
2594 async def get_shoutcast_stream(
2595 self, url: str, streamdetails: StreamDetails
2596 ) -> AsyncGenerator[bytes]:
2597 """
2598 Yield audio from a legacy Shoutcast server, with ICY metadata parsed inline.
2599
2600 :param url: Shoutcast stream URL.
2601 :param streamdetails: StreamDetails to update with ICY metadata as it arrives.
2602 """
2603 self.logger.debug("Start streaming from legacy Shoutcast server: %s", url)
2604
2605 parsed = urlparse(url)
2606 host = parsed.hostname
2607 port = parsed.port or 80
2608 path = parsed.path or "/"
2609 if parsed.query:
2610 path = f"{path}?{parsed.query}"
2611
2612 try:
2613 # Open raw socket connection
2614 reader, writer = await asyncio.wait_for(asyncio.open_connection(host, port), timeout=30)
2615 except TimeoutError as err:
2616 raise AudioError(f"Timeout connecting to Shoutcast stream {url}") from err
2617 except (OSError, ConnectionError) as err:
2618 raise AudioError(f"Failed to connect to Shoutcast stream {url}") from err
2619
2620 try:
2621 # Send HTTP request with ICY metadata header
2622 request = (
2623 f"GET {path} HTTP/1.1\r\n"
2624 f"Host: {host}\r\n"
2625 f"User-Agent: {HTTP_HEADERS['User-Agent']}\r\n"
2626 f"Icy-MetaData: 1\r\n\r\n"
2627 )
2628 writer.write(request.encode())
2629 await writer.drain()
2630
2631 # Read and parse response line
2632 try:
2633 response_line = await asyncio.wait_for(reader.readline(), timeout=10)
2634 except TimeoutError as err:
2635 raise AudioError("Timeout reading Shoutcast response") from err
2636
2637 if not response_line.startswith(b"ICY"):
2638 raise InvalidDataError("Invalid Shoutcast response")
2639
2640 # Read headers until empty line
2641 headers: dict[str, str] = {}
2642 while True:
2643 try:
2644 line = await asyncio.wait_for(reader.readline(), timeout=5)
2645 except TimeoutError as err:
2646 raise AudioError("Timeout reading Shoutcast headers") from err
2647
2648 if line in (b"\r\n", b"\n", b""):
2649 break
2650
2651 if b":" in line:
2652 try:
2653 key, value = line.decode("latin-1", errors="ignore").split(":", 1)
2654 headers[key.strip().lower()] = value.strip()
2655 except UnicodeDecodeError, ValueError:
2656 continue
2657
2658 # Get metadata interval
2659 meta_int_str = headers.get("icy-metaint")
2660 if not meta_int_str:
2661 raise InvalidDataError("No icy-metaint header in Shoutcast response")
2662
2663 try:
2664 meta_int = int(meta_int_str)
2665 except ValueError as err:
2666 raise InvalidDataError("Invalid icy-metaint value") from err
2667
2668 self.logger.debug("Connected to Shoutcast stream %s (icy-metaint: %s)", url, meta_int)
2669
2670 # Stream audio data with metadata parsing
2671 while True:
2672 try:
2673 # Read audio chunk
2674 audio_chunk = await reader.readexactly(meta_int)
2675 yield audio_chunk
2676
2677 # Read metadata length
2678 meta_byte = await reader.readexactly(1)
2679 if meta_byte == b"\x00":
2680 continue
2681
2682 meta_length = ord(meta_byte) * 16
2683 meta_data = await reader.readexactly(meta_length)
2684 self._parse_icy_metadata(meta_data, streamdetails)
2685
2686 except asyncio.exceptions.IncompleteReadError:
2687 # End of stream
2688 break
2689
2690 finally:
2691 writer.close()
2692 await writer.wait_closed()
2693
2694 # --- Private methods ---
2695
2696 def _notify_provider_streamed(
2697 self, streamdetails: StreamDetails, finished: bool, seconds_streamed: float
2698 ) -> None:
2699 """Report a (mostly) streamed item back to the provider that owns it."""
2700 if not finished and seconds_streamed < 90:
2701 return
2702 provider = self.mass.get_provider(streamdetails.provider)
2703 # plugin providers serve playable items too, but on_streamed is MusicProvider-only
2704 if provider is None or provider.type != ProviderType.MUSIC:
2705 return
2706 music_prov = cast("MusicProvider", provider)
2707 self.mass.create_task(music_prov.on_streamed(streamdetails))
2708
2709 def _get_volume_normalization_preference(
2710 self, streamdetails: StreamDetails
2711 ) -> VolumeNormalizationMode:
2712 """Return the configured normalization preference for the stream's media type."""
2713 conf_key = (
2714 CONF_VOLUME_NORMALIZATION_RADIO
2715 if streamdetails.media_type == MediaType.RADIO
2716 else CONF_VOLUME_NORMALIZATION_TRACKS
2717 )
2718 return VolumeNormalizationMode(
2719 self.mass.streams.get_config_value(conf_key, return_type=str)
2720 )
2721
2722 def _update_radio_stream_metadata(
2723 self,
2724 streamdetails: StreamDetails,
2725 artist: str | None,
2726 title: str,
2727 image_url: str | None = None,
2728 album: str | None = None,
2729 ) -> None:
2730 """
2731 Update radio stream metadata and trigger artwork lookup.
2732
2733 :param streamdetails: The stream details to update.
2734 :param artist: Artist name (will be normalized).
2735 :param title: Track title (will be cleaned for display).
2736 :param image_url: Optional image URL from stream metadata.
2737 :param album: Optional album name.
2738 """
2739 station_image_url = image_url or self.mass.metadata.get_radio_stream_station_image(
2740 streamdetails
2741 )
2742 artist_normalized = (
2743 self.mass.metadata.normalize_radio_artist_name(artist) if artist else None
2744 )
2745 display_title, _ = parse_title_and_version(title, strip_for_display=True)
2746
2747 streamdetails.stream_metadata = StreamMetadata(
2748 title=display_title,
2749 artist=artist_normalized,
2750 album=album,
2751 image_url=station_image_url,
2752 )
2753 streamdetails.stream_metadata_last_updated = time.time()
2754 if streamdetails.queue_id:
2755 self.mass.player_queues.signal_update(streamdetails.queue_id)
2756
2757 # Fetch artwork in background (track, album then artist)
2758 if artist and title and not image_url:
2759 self.mass.call_later(
2760 0.2,
2761 self.mass.metadata.update_radio_stream_artwork,
2762 streamdetails,
2763 task_id=f"update_radio_artwork_{streamdetails.queue_id}",
2764 )
2765
2766 async def _cache_radio_result(
2767 self,
2768 url: str,
2769 stream_type: StreamType,
2770 resolved_url: str | None = None,
2771 ) -> tuple[str, StreamType]:
2772 """Cache and return a radio stream resolution result."""
2773 result = (resolved_url or url, stream_type)
2774 await self.mass.cache.set(
2775 url,
2776 result,
2777 expiration=3600 * 3,
2778 provider=CACHE_PROVIDER,
2779 category=CACHE_CATEGORY_RESOLVED_RADIO_URL,
2780 )
2781 return result
2782
2783 async def _handle_client_error_for_radio_stream(
2784 self, url: str, err: aiohttp.ClientError, fallback_stream_type: StreamType
2785 ) -> tuple[str, StreamType]:
2786 """Handle aiohttp client errors during radio stream resolution."""
2787 # Prefer the final post-redirect URL: aiohttp follows redirects before raising,
2788 # but the original url may just point at a redirector rather than the ICY endpoint.
2789 request_info = getattr(err, "request_info", None)
2790 validate_url = str(request_info.url) if request_info is not None else url
2791
2792 # Check if this is a Shoutcast/ICY response that aiohttp can't parse
2793 if isinstance(err, aiohttp.ClientResponseError) and "ICY" in str(err).upper():
2794 self.logger.debug(
2795 "ICY response detected for %s, validating Shoutcast stream", validate_url
2796 )
2797 if await self._validate_shoutcast_stream(validate_url):
2798 return await self._cache_radio_result(
2799 url, StreamType.SHOUTCAST, resolved_url=validate_url
2800 )
2801 self.logger.warning(
2802 "ICY response detected but Shoutcast validation failed for %s", validate_url
2803 )
2804 return await self._cache_radio_result(
2805 url, fallback_stream_type, resolved_url=validate_url
2806 )
2807
2808 # Other aiohttp errors - might still be Shoutcast, check it
2809 self.logger.debug("aiohttp error for %s, checking if legacy Shoutcast stream", validate_url)
2810 if await self._validate_shoutcast_stream(validate_url):
2811 return await self._cache_radio_result(
2812 url, StreamType.SHOUTCAST, resolved_url=validate_url
2813 )
2814
2815 # Unknown error - still try to stream
2816 self.logger.warning(
2817 "Failed to parse radio URL %s: %s - attempting direct stream", validate_url, str(err)
2818 )
2819 return await self._cache_radio_result(url, fallback_stream_type, resolved_url=validate_url)
2820
2821 async def _get_audio_buffer(
2822 self,
2823 queue_item: QueueItem,
2824 seek_position_ms: int,
2825 reason: str,
2826 capacity_wait_timeout: float,
2827 allow_provider_match: bool,
2828 ) -> AudioBuffer:
2829 """
2830 Create or reuse a ready AudioBuffer within one queue-item preparation lock.
2831
2832 :param queue_item: Queue item whose source should be buffered.
2833 :param seek_position_ms: Position in milliseconds to start from.
2834 :param reason: Caller context for logging.
2835 :param capacity_wait_timeout: Total seconds to spend waiting for source capacity.
2836 :param allow_provider_match: Whether an on-demand cross-provider match may widen
2837 the candidates when all are saturated.
2838 """
2839 loop = asyncio.get_running_loop()
2840 # the playback intent lives on the details we start from; keep it across a reselection
2841 initial_streamdetails = queue_item.streamdetails
2842 seek_position = (
2843 int(initial_streamdetails.seek_position)
2844 if initial_streamdetails
2845 else seek_position_ms // 1000
2846 )
2847 fade_in = bool(initial_streamdetails and initial_streamdetails.fade_in)
2848 prefer_album_loudness = bool(
2849 initial_streamdetails and initial_streamdetails.prefer_album_loudness
2850 )
2851 all_candidate_instances = {
2852 provider.instance_id
2853 for mapping in (
2854 queue_item.media_item.provider_mappings if queue_item.media_item else ()
2855 )
2856 if mapping.available
2857 for provider in self._get_mapping_providers(mapping)
2858 }
2859 if initial_streamdetails is not None:
2860 all_candidate_instances.add(initial_streamdetails.provider)
2861 # a track may also exist on streaming providers it has no mapping for yet; such a
2862 # match is only searched once, and only when every known candidate is saturated
2863 match_pending = (
2864 allow_provider_match
2865 and isinstance(queue_item.media_item, Track)
2866 and self._has_alternative_match_providers(queue_item.media_item)
2867 )
2868
2869 deadline = loop.time() + capacity_wait_timeout
2870 busy_instances: set[str] = set()
2871 final_pass = False
2872 last_capacity_error: ProviderStreamLimitError | None = None
2873 last_failed_streamdetails: StreamDetails | None = None
2874 while True:
2875 if queue_item.streamdetails is None or (
2876 queue_item.streamdetails.provider in busy_instances and not final_pass
2877 ):
2878 try:
2879 queue_item.streamdetails = await self.get_stream_details(
2880 queue_item,
2881 seek_position=seek_position,
2882 fade_in=fade_in,
2883 prefer_album_loudness=prefer_album_loudness,
2884 excluded_provider_instances=busy_instances,
2885 )
2886 except (AudioError, MediaNotFoundError) as err:
2887 if last_capacity_error is None:
2888 raise
2889 if final_pass:
2890 # capacity was the root cause, surface the typed (actionable) error
2891 raise last_capacity_error from err
2892 # no usable alternative mapping: restore the capacity-blocked details
2893 # and spend the remaining budget blocking on that provider's slot
2894 final_pass = True
2895 continue
2896 finally:
2897 if queue_item.streamdetails is None:
2898 # never leave the queue item without streamdetails on any exit,
2899 # including a cancellation or a non-audio provider failure
2900 queue_item.streamdetails = last_failed_streamdetails
2901 streamdetails = queue_item.streamdetails
2902 assert streamdetails is not None # for type checking
2903 remaining = max(deadline - loop.time(), 0)
2904 alternatives_left = bool(
2905 all_candidate_instances - busy_instances - {streamdetails.provider}
2906 )
2907 # probe (0s) whenever a reselection can still follow: a free slot is still
2908 # acquired instantly, while a busy one fails fast instead of spending the
2909 # whole budget on this candidate. Block only on the last resort.
2910 source_wait = (
2911 0.0
2912 if (not final_pass and (alternatives_left or busy_instances or match_pending))
2913 else remaining
2914 )
2915 try:
2916 return await AudioBuffer.get_buffer(
2917 mass=self.mass,
2918 streamdetails=streamdetails,
2919 seek_position_ms=seek_position_ms,
2920 wait_ready=True,
2921 reason=reason,
2922 source_wait_timeout=source_wait,
2923 )
2924 except ProviderStreamLimitError as err:
2925 last_capacity_error = err
2926 last_failed_streamdetails = streamdetails
2927 busy_instances.add(err.provider_instance)
2928 if final_pass or loop.time() >= deadline:
2929 raise
2930 if all_candidate_instances.issubset(busy_instances):
2931 discovered: set[str] = set()
2932 if match_pending:
2933 match_pending = False
2934 try:
2935 discovered = await self._discover_alternative_provider_mappings(
2936 queue_item, busy_instances, max(deadline - loop.time(), 0)
2937 )
2938 except Exception as err:
2939 # discovery is best-effort: any failure falls back to the
2940 # final blocking wait instead of replacing the typed error
2941 self.logger.warning(
2942 "Alternative provider search for %s failed: %s",
2943 queue_item.name,
2944 err,
2945 )
2946 if discovered:
2947 all_candidate_instances.update(discovered)
2948 else:
2949 # every candidate is saturated: one last blocking wait on the best one
2950 busy_instances.clear()
2951 final_pass = True
2952 queue_item.streamdetails = None
2953 except AudioError:
2954 if last_capacity_error is None or final_pass:
2955 raise
2956 # a broken alternate must not turn a transient capacity miss into a hard
2957 # failure: restore the blocked details and spend the rest of the budget there
2958 queue_item.streamdetails = last_failed_streamdetails
2959 final_pass = True
2960
2961 def _get_streamdetail_candidates(
2962 self,
2963 provider_mappings: Iterable[ProviderMapping],
2964 preferred_providers: list[str],
2965 excluded_provider_instances: set[str],
2966 ) -> list[tuple[ProviderMapping, Provider]]:
2967 """
2968 Return mapping candidates in steering, quality, and instance-fallback order.
2969
2970 :param provider_mappings: Mappings attached to the media item.
2971 :param preferred_providers: Provider instances tried before widening to the rest.
2972 :param excluded_provider_instances: Provider instances unavailable to this attempt.
2973 :return: Ordered provider mapping candidates.
2974 """
2975 ordered_mappings = sorted(
2976 provider_mappings, key=lambda mapping: mapping.quality or 0, reverse=True
2977 )
2978 preferred_candidates: list[tuple[ProviderMapping, Provider]] = []
2979 fallback_candidates: list[tuple[ProviderMapping, Provider]] = []
2980 seen_candidates: set[tuple[str, str]] = set()
2981 for mapping in ordered_mappings:
2982 if not mapping.available:
2983 self.logger.debug("Skipping unavailable %s", mapping)
2984 continue
2985 for provider in self._get_mapping_providers(mapping):
2986 candidate_id = (provider.instance_id, mapping.item_id)
2987 if (
2988 candidate_id in seen_candidates
2989 or provider.instance_id in excluded_provider_instances
2990 ):
2991 continue
2992 seen_candidates.add(candidate_id)
2993 candidate = (mapping, provider)
2994 if provider.instance_id in preferred_providers:
2995 preferred_candidates.append(candidate)
2996 else:
2997 fallback_candidates.append(candidate)
2998 return [*preferred_candidates, *fallback_candidates]
2999
3000 def _get_mapping_providers(self, mapping: ProviderMapping) -> list[Provider]:
3001 """
3002 Return the mapped provider followed by compatible instances of its streaming catalog.
3003
3004 :param mapping: Provider mapping whose item ID will be requested.
3005 :return: Loaded provider instances that can resolve the mapping.
3006 """
3007 providers: list[Provider] = []
3008 if (
3009 primary_provider := self.mass.get_provider(
3010 mapping.provider_instance, return_unavailable=True
3011 )
3012 ) and primary_provider.available:
3013 providers.append(primary_provider)
3014 # another account of the same streaming catalog serves the same item ID,
3015 # so it can stand in when the mapped instance can not
3016 for provider in self.mass.providers:
3017 if (
3018 not isinstance(provider, MusicProvider)
3019 or not provider.available
3020 or not provider.is_streaming_provider
3021 or provider.domain != mapping.provider_domain
3022 or provider in providers
3023 ):
3024 continue
3025 providers.append(provider)
3026 if not providers:
3027 self.logger.debug("Skipping %s - provider not available", mapping)
3028 return providers
3029
3030 def _is_match_candidate_provider(
3031 self, provider: MusicProvider, known_domains: set[str]
3032 ) -> bool:
3033 """
3034 Return whether a provider is eligible to search a track match on.
3035
3036 :param provider: Music provider to check.
3037 :param known_domains: Provider domains the track already has mappings for.
3038 """
3039 return (
3040 provider.available
3041 and provider.is_streaming_provider
3042 and ProviderFeature.SEARCH in provider.supported_features
3043 and provider.domain not in known_domains
3044 and MediaType.TRACK in provider.supported_media_types
3045 )
3046
3047 def _has_alternative_match_providers(self, media_item: Track) -> bool:
3048 """
3049 Return whether any configured streaming provider could carry an unmapped match.
3050
3051 :param media_item: Track whose existing mappings define the known provider domains.
3052 """
3053 known_domains = {mapping.provider_domain for mapping in media_item.provider_mappings}
3054 return any(
3055 self._is_match_candidate_provider(provider, known_domains)
3056 for provider in self.mass.music.providers
3057 )
3058
3059 async def _discover_alternative_provider_mappings(
3060 self, queue_item: QueueItem, busy_instances: set[str], remaining: float
3061 ) -> set[str]:
3062 """
3063 Search other streaming providers for the queue item's track and widen its mappings.
3064
3065 A found mapping is added to the media item (and persisted for library items) so the
3066 capacity reselection can continue on the discovered provider.
3067
3068 :param queue_item: Queue item whose track should be matched on another provider.
3069 :param busy_instances: Provider instances already known to be saturated.
3070 :param remaining: Seconds left of the caller's capacity budget.
3071 :return: Provider instances able to serve the discovered mappings.
3072 """
3073 media_item = queue_item.media_item
3074 if not isinstance(media_item, Track):
3075 return set()
3076 known_domains = {mapping.provider_domain for mapping in media_item.provider_mappings}
3077 eligible = [
3078 provider
3079 for provider in self.mass.music.providers
3080 if self._is_match_candidate_provider(provider, known_domains)
3081 and provider.instance_id not in busy_instances
3082 and provider.has_available_stream_slot
3083 ]
3084 if not eligible:
3085 return set()
3086 # mirror the playback user's provider steering for the search order
3087 if (
3088 (pq_data := self.mass.player_queues.queue_data_or_none(queue_item.queue_id))
3089 and pq_data.userid
3090 and (playback_user := await self.mass.webserver.auth.get_user(pq_data.userid))
3091 and playback_user.provider_filter
3092 ):
3093 preferred = set(playback_user.provider_filter)
3094 eligible.sort(key=lambda provider: provider.instance_id not in preferred)
3095 # one instance per domain: a found mapping widens to sibling instances anyway
3096 candidates: list[MusicProvider] = []
3097 for provider in eligible:
3098 if provider.domain in known_domains:
3099 continue
3100 known_domains.add(provider.domain)
3101 candidates.append(provider)
3102 # the track's own album is free, sufficient evidence for the strict compare and
3103 # avoids match_provider's multi-provider album lookup on every call
3104 ref_albums = [media_item.album] if isinstance(media_item.album, Album) else []
3105 matches: list[ProviderMapping] = []
3106 try:
3107 async with asyncio.timeout(min(STREAM_SLOT_MATCH_TIMEOUT, remaining)):
3108 for provider in candidates:
3109 # one failing provider must not end the search on the others
3110 try:
3111 matches = await self.mass.music.tracks.match_provider(
3112 media_item, provider, strict=True, ref_albums=ref_albums
3113 )
3114 except Exception as err:
3115 self.logger.debug("Searching a match on %s failed: %s", provider.name, err)
3116 continue
3117 if matches:
3118 break
3119 except TimeoutError:
3120 self.logger.debug("Searching an alternative provider for %s timed out", media_item.name)
3121 if not matches:
3122 return set()
3123 media_item.provider_mappings.update(matches)
3124 if media_item.provider == "library":
3125 # persist in the background so future plays have the mapping ahead of time;
3126 # cancellation of this playback must never interrupt the library write
3127 self.mass.create_task(
3128 self.mass.music.tracks.add_provider_mappings(media_item.item_id, matches)
3129 )
3130 self.logger.info(
3131 "All known sources for %s are at their stream limit, "
3132 "using a matching track found on %s",
3133 media_item.name,
3134 matches[0].provider_domain,
3135 )
3136 return {
3137 provider.instance_id
3138 for mapping in matches
3139 for provider in self._get_mapping_providers(mapping)
3140 }
3141
3142 async def _request_streamdetails(
3143 self,
3144 candidates: Iterable[tuple[ProviderMapping, Provider]],
3145 media_type: MediaType,
3146 ) -> StreamDetails | None:
3147 """
3148 Request stream details from ordered provider mapping candidates.
3149
3150 :param candidates: Candidates in mapping and compatible-instance order.
3151 :param media_type: Media type requested from each provider.
3152 :return: The first resolved stream details, or None when every candidate failed.
3153 :raises AudioError: The last (actionable) audio error when no candidate resolved.
3154 """
3155 last_audio_error: AudioError | None = None
3156 for mapping, provider in candidates:
3157 # music and plugin providers share this signature, so either type can own the item
3158 token = BYPASS_THROTTLER.set(True)
3159 try:
3160 stream_prov = cast("MusicProvider | PluginProvider", provider)
3161 return await stream_prov.get_stream_details(mapping.item_id, media_type)
3162 except AudioError as err:
3163 # remember the last one so its (actionable) message can be re-raised
3164 last_audio_error = err
3165 self.logger.warning("%s", err)
3166 except MusicAssistantError as err:
3167 self.logger.warning("%s", err)
3168 finally:
3169 BYPASS_THROTTLER.reset(token)
3170 if last_audio_error is not None:
3171 raise last_audio_error
3172 return None
3173
3174 async def _get_media_stream(
3175 self,
3176 streamdetails: StreamDetails,
3177 pcm_format: AudioFormat,
3178 seek_position: int,
3179 filter_params: list[str] | None,
3180 chunk_seconds: float,
3181 ) -> AsyncGenerator[bytes]:
3182 """
3183 Stream one provider source as raw PCM.
3184
3185 :param streamdetails: Details of the stream to fetch.
3186 :param pcm_format: Target PCM format the consumer expects.
3187 :param seek_position: Requested seek offset in seconds.
3188 :param filter_params: Optional ffmpeg filter expressions.
3189 :param chunk_seconds: Size of each yielded chunk in seconds of audio.
3190 """
3191 mass = self.mass
3192 logger = self.logger.getChild("media_stream")
3193 logger.log(VERBOSE_LOG_LEVEL, "Starting media stream for %s", streamdetails.uri)
3194 # copy: the args below are appended per call, while the StreamDetails is cached on
3195 # the queue item and reused across calls (retry, seek, background analysis)
3196 extra_input_args = list(streamdetails.extra_input_args or [])
3197 # the resolver below zeroes out seek_position where the seek is delegated to the
3198 # source itself, so keep the requested position for the duration writeback
3199 requested_seek_position = seek_position
3200
3201 # work out audio source for these streamdetails
3202 audio_source, seek_position, extra_input_args = await self._resolve_media_stream_source(
3203 streamdetails, seek_position, extra_input_args
3204 )
3205
3206 # pace ffmpeg at native rate for live sources; the producer (e.g.
3207 # librespot's pipe backend) may otherwise write faster than realtime.
3208 # The initial burst grants a small bounded read-ahead so downstream
3209 # jitter does not immediately underrun the player. Providers that need
3210 # different pacing can pass their own -re/-readrate args to override.
3211 if (
3212 streamdetails.media_type == MediaType.AUDIO_SOURCE
3213 and "-re" not in extra_input_args
3214 and "-readrate" not in extra_input_args
3215 ):
3216 extra_input_args += ["-readrate", "1", "-readrate_initial_burst", "0.5"]
3217
3218 # handle seek support
3219 if seek_position and streamdetails.duration and streamdetails.allow_seek:
3220 extra_input_args += ["-ss", str(int(seek_position))]
3221
3222 bytes_sent = 0
3223 finished = False
3224 cancelled = False
3225 first_chunk_received = False
3226 ffmpeg_loglevel = "debug" if self.logger.isEnabledFor(VERBOSE_LOG_LEVEL) else "info"
3227 # When a provider hands us already-decoded audio (e.g. Spotify Connect /
3228 # AirPlay receivers piping PCM after their own decode), audio_format is
3229 # the original source format meant for display while decoded_audio_format
3230 # is what ffmpeg actually needs to read off the wire.
3231 ffmpeg_input_format = streamdetails.decoded_audio_format or streamdetails.audio_format
3232 ffmpeg_proc = FFMpeg(
3233 audio_input=audio_source,
3234 input_format=ffmpeg_input_format,
3235 output_format=pcm_format,
3236 filter_params=filter_params,
3237 extra_input_args=extra_input_args,
3238 collect_log_history=True,
3239 loglevel=ffmpeg_loglevel,
3240 )
3241
3242 try:
3243 await ffmpeg_proc.start()
3244 assert ffmpeg_proc.proc is not None # for type checking
3245 if logger.isEnabledFor(VERBOSE_LOG_LEVEL):
3246 logger.log(
3247 VERBOSE_LOG_LEVEL,
3248 "Started media stream for %s - using streamtype: %s "
3249 "- pcm format: %s - ffmpeg PID: %s",
3250 streamdetails.uri,
3251 streamdetails.stream_type,
3252 pcm_format.content_type.value,
3253 ffmpeg_proc.proc.pid,
3254 )
3255 else:
3256 logger.debug(
3257 "Started media stream for %s - using streamtype: %s",
3258 streamdetails.uri,
3259 streamdetails.stream_type,
3260 )
3261 stream_start = mass.loop.time()
3262 chunk_size = calculate_content_length(pcm_format, chunk_seconds)
3263 chunk_iter = ffmpeg_proc.iter_chunked(chunk_size)
3264 while True:
3265 # Time the read, not the yield: catches a stalled source, ignores backpressure.
3266 read_timeout = (
3267 STREAM_START_TIMEOUT if not first_chunk_received else STREAM_STALL_TIMEOUT
3268 )
3269 try:
3270 async with asyncio.timeout(read_timeout):
3271 chunk = await anext(chunk_iter)
3272 except StopAsyncIteration:
3273 break
3274 except TimeoutError as err:
3275 raise AudioError(f"Source stalled: no audio for {read_timeout}s") from err
3276 if not first_chunk_received:
3277 # At this point ffmpeg has started and should now know the codec used
3278 # for encoding the audio.
3279 # Note: ffmpeg_proc.input_format is the same object as
3280 # ffmpeg_input_format, so sample_rate / bit_depth / bit_rate
3281 # parsed from the ffmpeg log already live on streamdetails too.
3282 first_chunk_received = True
3283 # Skip the codec_type writeback when the provider declared a
3284 # decoded format: audio_format already holds the authoritative
3285 # source codec and the probed value would just be the
3286 # post-decode wire format (e.g. PCM for Spotify Connect).
3287 if streamdetails.decoded_audio_format is None:
3288 streamdetails.audio_format.codec_type = ffmpeg_proc.input_format.codec_type
3289 # Some providers omit (or report 0 for) the item duration; ffmpeg can
3290 # usually probe it from the source. Only apply when missing so we
3291 # don't clobber an accurate provider value with a rounded one.
3292 if ffmpeg_proc.parsed_duration is not None and not streamdetails.duration:
3293 streamdetails.duration = ffmpeg_proc.parsed_duration
3294 logger.debug(
3295 "First chunk received after %.2f seconds (codec detected: %s)",
3296 mass.loop.time() - stream_start,
3297 ffmpeg_proc.input_format.codec_type,
3298 )
3299 yield chunk
3300 bytes_sent += len(chunk)
3301
3302 # end of audio/track reached
3303 logger.debug("End of media stream reached for %s", streamdetails.uri)
3304 # wait until stderr also completed reading
3305 await ffmpeg_proc.wait_with_timeout(5)
3306 logger.log(
3307 VERBOSE_LOG_LEVEL,
3308 "FFmpeg process ended with return code %s for %s",
3309 ffmpeg_proc.returncode,
3310 streamdetails.uri,
3311 )
3312 # a nested source raises through the stdin feeder, where ffmpeg's own exit
3313 # would otherwise flatten it into a generic AudioError
3314 if isinstance(ffmpeg_proc.stdin_feeder_exception, ProviderStreamLimitError):
3315 raise ffmpeg_proc.stdin_feeder_exception
3316 if ffmpeg_proc.returncode not in (0, None):
3317 log_trail = "\n".join(list(ffmpeg_proc.log_history)[-5:])
3318 raise AudioError(f"FFMpeg exited with code {ffmpeg_proc.returncode}: {log_trail}")
3319 if bytes_sent == 0:
3320 # edge case: no audio data was received at all
3321 raise AudioError("No audio was received")
3322 finished = True
3323 except (Exception, GeneratorExit, asyncio.CancelledError) as err:
3324 if isinstance(err, asyncio.CancelledError | GeneratorExit):
3325 # we were cancelled, just raise
3326 cancelled = True
3327 raise
3328 if isinstance(ffmpeg_proc.stdin_feeder_exception, ProviderStreamLimitError):
3329 raise ffmpeg_proc.stdin_feeder_exception
3330 # dump the last 10 lines of the log in case of an unclean exit
3331 logger.warning("\n".join(list(ffmpeg_proc.log_history)[-10:]))
3332 raise AudioError(f"Error while streaming: {err}") from err
3333 finally:
3334 # always ensure close is called which also handles all cleanup
3335 await ffmpeg_proc.close()
3336 # determine how many seconds we've received
3337 # for pcm output we can calculate this easily
3338 seconds_received = bytes_sent / pcm_format.pcm_sample_size if bytes_sent else 0
3339 # store accurate duration, but only for a playthrough from the very start:
3340 # a seeked stream yields the remaining audio, not the item's full length
3341 if finished and not requested_seek_position and seconds_received:
3342 streamdetails.duration = int(seconds_received)
3343
3344 logger.log(
3345 VERBOSE_LOG_LEVEL,
3346 "stream %s (with code %s) for %s",
3347 "cancelled" if cancelled else "finished" if finished else "aborted",
3348 ffmpeg_proc.returncode,
3349 streamdetails.uri,
3350 )
3351
3352 def _crossfade_holdback_allowed(
3353 self, streamdetails: StreamDetails, tail_seconds: float, playback_speed: float = 1.0
3354 ) -> bool:
3355 """
3356 Return whether the outgoing tail may be held back for a crossfade.
3357
3358 :param streamdetails: Stream details of the track being streamed.
3359 :param tail_seconds: Length of the tail to hold back, in seconds of playback.
3360 :param playback_speed: Playback-speed multiplier of the track.
3361 """
3362 if tail_seconds <= 0 or playback_speed <= 0 or streamdetails.is_realtime:
3363 return False
3364 audio_buffer = cast("AudioBuffer | None", streamdetails.buffer)
3365 if audio_buffer is None or audio_buffer.has_error:
3366 # a failed source is skipped without a fade, so its remaining audio is
3367 # better off played out than held back for one
3368 return False
3369 # While the source is still delivering, it is what limits playback: withholding
3370 # a tail on top of that eats into the lead the player needs. Once the source is
3371 # done the remaining audio is resident, so the tail comes for free. A buffer that
3372 # is too small to ever hold the tail is the exception - waiting for the source
3373 # there would only lose the fade.
3374 return audio_buffer.eof or audio_buffer.max_size_seconds / playback_speed < tail_seconds
3375
3376 def _select_buffered_crossfade(
3377 self,
3378 streamdetails: StreamDetails,
3379 crossfade_mode: CrossfadeMode,
3380 standard_crossfade_duration: int,
3381 playback_speed: float = 1.0,
3382 ) -> tuple[CrossfadeMode, float]:
3383 """
3384 Select a crossfade that can be completed from resident incoming PCM.
3385
3386 :param streamdetails: Incoming track stream details.
3387 :param crossfade_mode: Requested crossfade mode.
3388 :param standard_crossfade_duration: Configured standard overlap in seconds.
3389 :param playback_speed: Incoming track playback-speed multiplier.
3390 :return: Effective mode and resident fade-in duration in seconds.
3391 """
3392 audio_buffer = streamdetails.buffer
3393 if (
3394 crossfade_mode == CrossfadeMode.DISABLED
3395 or playback_speed <= 0
3396 # a realtime source has no more than its banked lead: fading in would spend
3397 # that at the start of the track and leave nothing for the rest of it
3398 or streamdetails.is_realtime
3399 or audio_buffer is None
3400 or audio_buffer.has_error
3401 or not audio_buffer.is_valid()
3402 ):
3403 return CrossfadeMode.DISABLED, 0
3404
3405 available_seconds = audio_buffer.duration_available / playback_speed
3406 if (
3407 crossfade_mode == CrossfadeMode.SMART_CROSSFADE
3408 and audio_buffer.ready.is_set()
3409 and available_seconds >= MIN_EFFECTIVE_FADE_BUFFER
3410 ):
3411 return crossfade_mode, min(SMART_CROSSFADE_DURATION, available_seconds)
3412 if (
3413 crossfade_mode == CrossfadeMode.STANDARD_CROSSFADE
3414 and standard_crossfade_duration >= MIN_CROSSFADE_FALLBACK_DURATION
3415 and audio_buffer.ready.is_set()
3416 and available_seconds >= standard_crossfade_duration
3417 ):
3418 return crossfade_mode, standard_crossfade_duration
3419
3420 fallback_duration = min(standard_crossfade_duration, available_seconds)
3421 if fallback_duration < MIN_CROSSFADE_FALLBACK_DURATION:
3422 return CrossfadeMode.DISABLED, 0
3423 self.logger.debug(
3424 "Using %s second standard crossfade for %s from resident audio",
3425 fallback_duration,
3426 streamdetails.uri,
3427 )
3428 return CrossfadeMode.STANDARD_CROSSFADE, fallback_duration
3429
3430 async def _resolve_media_stream_source(
3431 self,
3432 streamdetails: StreamDetails,
3433 seek_position: int,
3434 extra_input_args: list[str],
3435 ) -> tuple[str | AsyncGenerator[bytes], int, list[str]]:
3436 """
3437 Resolve the input consumed by ffmpeg for the given stream details.
3438
3439 :param streamdetails: Details of the stream to fetch.
3440 :param seek_position: Requested seek offset in seconds.
3441 :param extra_input_args: Provider-supplied ffmpeg input arguments.
3442 :return: The ffmpeg input, the remaining seek offset and the ffmpeg input arguments.
3443 """
3444 stream_type = streamdetails.stream_type
3445 if stream_type == StreamType.CUSTOM:
3446 # MusicProvider and PluginProvider both expose get_audio_stream with the same shape.
3447 # Pin the exact instance: a domain fallback would stream from a sibling account
3448 # while the source-stream slot is charged to the instance that issued the details.
3449 provider = self.mass.get_provider(streamdetails.provider, return_unavailable=True)
3450 if provider is None or not provider.available:
3451 raise ProviderUnavailableError(
3452 f"Provider {streamdetails.provider} for stream is no longer available"
3453 )
3454 provider = cast("MusicProvider | PluginProvider", provider)
3455 audio_source = provider.get_audio_stream(
3456 streamdetails, seek_position=seek_position if streamdetails.can_seek else 0
3457 )
3458 return audio_source, 0 if streamdetails.can_seek else seek_position, extra_input_args
3459 if stream_type == StreamType.ICY:
3460 assert streamdetails.path is not None
3461 assert isinstance(streamdetails.path, (str, list))
3462 audio_source = self.get_reconnecting_icy_radio_stream(streamdetails.path, streamdetails)
3463 return audio_source, 0, extra_input_args
3464 if stream_type == StreamType.SHOUTCAST:
3465 assert isinstance(streamdetails.path, str)
3466 return self.get_shoutcast_stream(streamdetails.path, streamdetails), 0, extra_input_args
3467 if stream_type == StreamType.IN_BAND:
3468 assert isinstance(streamdetails.path, str) # for type checking
3469
3470 # For IN_BAND (OGG/Opus) radio streams, use chained OGG handler.
3471 # This handles the chained OGG format by stitching logical bitstreams together
3472 # so FFmpeg sees a single continuous stream. Metadata is extracted in-band.
3473 audio_source = get_chained_ogg_stream(
3474 self.mass,
3475 streamdetails.path,
3476 metadata_callback=partial(self._handle_inband_metadata, streamdetails),
3477 )
3478 # seeking not possible on radio streams
3479 return audio_source, 0, extra_input_args
3480 if stream_type == StreamType.HLS:
3481 assert isinstance(streamdetails.path, str) # for type checking
3482 substream = await self.get_hls_substream(streamdetails.path)
3483 if streamdetails.media_type == MediaType.RADIO:
3484 # HLS streams (especially the BBC) struggle when they're played directly
3485 # with ffmpeg, where they just stop after some minutes,
3486 # so we tell ffmpeg to loop around in this case.
3487 extra_input_args += ["-stream_loop", "-1", "-re"]
3488 return substream.path, seek_position, extra_input_args
3489
3490 # all other stream types (HTTP, FILE, etc)
3491 if stream_type == StreamType.ENCRYPTED_HTTP:
3492 assert streamdetails.decryption_key is not None # for type checking
3493 extra_input_args += ["-decryption_key", streamdetails.decryption_key]
3494 if isinstance(streamdetails.path, list):
3495 # multi part stream, which handles the seek itself
3496 return self.get_multi_file_stream(streamdetails, seek_position), 0, extra_input_args
3497 # regular single file/url stream
3498 assert isinstance(streamdetails.path, str) # for type checking
3499 return streamdetails.path, seek_position, extra_input_args
3500
3501 async def _iter_audio_source_pcm(
3502 self,
3503 streamdetails: StreamDetails,
3504 pcm_format: AudioFormat,
3505 ) -> AsyncGenerator[bytes]:
3506 """Yield PCM for an AudioSource, bypassing ffmpeg when formats match."""
3507 if _pcm_formats_match(streamdetails.audio_format, pcm_format):
3508 source_gen = self._open_audio_source_generator(streamdetails)
3509 async for chunk in realtime_pcm_pacer(source_gen, pcm_format):
3510 yield chunk
3511 return
3512 # format mismatch â fall back to ffmpeg for resampling (still small chunks)
3513 async for chunk in self.get_media_stream(
3514 streamdetails=streamdetails,
3515 pcm_format=pcm_format,
3516 filter_params=None,
3517 chunk_seconds=AUDIO_SOURCE_CHUNK_SECONDS,
3518 ):
3519 yield chunk
3520
3521 def _open_audio_source_generator(self, streamdetails: StreamDetails) -> AsyncGenerator[bytes]:
3522 """Open the raw PCM generator for an AudioSource (CUSTOM or NAMED_PIPE)."""
3523 if streamdetails.stream_type == StreamType.CUSTOM:
3524 # pin the exact instance, see _resolve_media_stream_source
3525 provider = self.mass.get_provider(streamdetails.provider, return_unavailable=True)
3526 if provider is None or not provider.available:
3527 raise ProviderUnavailableError(
3528 f"Provider {streamdetails.provider} for stream is no longer available"
3529 )
3530 provider = cast("MusicProvider | PluginProvider", provider)
3531 return provider.get_audio_stream(streamdetails)
3532 if streamdetails.stream_type == StreamType.NAMED_PIPE:
3533 assert isinstance(streamdetails.path, str) # for type checking
3534 return read_named_pipe(streamdetails.path)
3535 raise AudioError(f"Unsupported stream_type {streamdetails.stream_type} for AudioSource")
3536
3537 def _handle_inband_metadata(
3538 self, streamdetails: StreamDetails, metadata: dict[str, str]
3539 ) -> None:
3540 """Handle metadata extracted from a chained Ogg stream."""
3541 title = metadata.get("title", "")
3542 artist = metadata.get("artist", "")
3543 album = metadata.get("album", "")
3544 if not artist and " - " in title:
3545 artist, title = title.split(" - ", 1)
3546 if not (title or artist):
3547 return
3548
3549 stream_title = f"{artist} - {title}" if artist and title else title or artist
3550 cleaned_title = clean_stream_title(stream_title)
3551 if not cleaned_title:
3552 return
3553 if self._record_inband_stream_title(streamdetails, cleaned_title):
3554 return
3555 if cleaned_title != streamdetails.stream_title:
3556 self.logger.log(VERBOSE_LOG_LEVEL, "In-band metadata: %s", cleaned_title)
3557 streamdetails.stream_title = cleaned_title
3558 self._update_radio_stream_metadata(
3559 streamdetails,
3560 artist=artist or None,
3561 title=title or cleaned_title,
3562 album=album or None,
3563 )
3564
3565 def _record_inband_stream_title(self, streamdetails: StreamDetails, cleaned_title: str) -> bool:
3566 """
3567 Record an in-band stream title for provider-owned metadata, if applicable.
3568
3569 When a provider opts into owning stream_metadata (and stream_title is only
3570 a derived view of it), writing either from the stream reader would fight the
3571 provider. The cleaned in-band title is recorded on StreamDetails.data instead,
3572 as the identity signal for the provider callback.
3573
3574 :param streamdetails: StreamDetails carrying the stream.
3575 :param cleaned_title: Cleaned in-band stream title.
3576 :returns: True when recorded (the caller must not write stream metadata);
3577 False when no provider callback exists and normal handling applies.
3578 """
3579 if (
3580 streamdetails.stream_metadata_update_callback is None
3581 or streamdetails.data is None
3582 or not streamdetails.data.get(STREAMDETAILS_INBAND_TITLE_HANDOFF_KEY)
3583 ):
3584 return False
3585 if streamdetails.data.get(STREAMDETAILS_INBAND_TITLE_KEY) != cleaned_title:
3586 # occupancy approximates how far this detection leads audible playback
3587 buffer = streamdetails.buffer
3588 self.logger.debug(
3589 "In-band stream title: %s (buffer occupancy: %ss)",
3590 cleaned_title,
3591 buffer.size_seconds if buffer is not None else "unknown",
3592 )
3593 streamdetails.data[STREAMDETAILS_INBAND_TITLE_KEY] = cleaned_title
3594 return True
3595
3596 def _parse_icy_metadata(self, meta_data: bytes, streamdetails: StreamDetails) -> None:
3597 """
3598 Parse ICY metadata and update streamdetails.
3599
3600 Sets the cleaned stream title and, when the title parses as "Artist - Track",
3601 triggers a radio-artwork metadata update.
3602
3603 :param meta_data: Raw metadata bytes from an ICY stream chunk.
3604 :param streamdetails: StreamDetails to update with parsed title and metadata.
3605 """
3606 if not meta_data:
3607 return
3608
3609 meta_data = meta_data.rstrip(b"\0")
3610 # Match StreamTitle, handling apostrophes in titles
3611 stream_title_re = re.search(rb"StreamTitle='(.*?)';", meta_data)
3612
3613 if not stream_title_re:
3614 self.logger.log(
3615 VERBOSE_LOG_LEVEL,
3616 "ICY metadata does not contain StreamTitle field. Raw: %s",
3617 meta_data.decode("utf-8", errors="replace")[:200],
3618 )
3619 return
3620
3621 try:
3622 # in 99% of the cases the stream title is utf-8 encoded
3623 stream_title = stream_title_re.group(1).decode("utf-8")
3624 except UnicodeDecodeError:
3625 # fallback to iso-8859-1
3626 stream_title = stream_title_re.group(1).decode("iso-8859-1", errors="replace")
3627
3628 cleaned_stream_title = clean_stream_title(stream_title)
3629
3630 if not cleaned_stream_title:
3631 return
3632
3633 if self._record_inband_stream_title(streamdetails, cleaned_stream_title):
3634 return
3635
3636 if cleaned_stream_title == streamdetails.stream_title:
3637 return
3638
3639 self.logger.log(VERBOSE_LOG_LEVEL, "ICY Radio streamtitle original: %s", stream_title)
3640 self.logger.log(
3641 VERBOSE_LOG_LEVEL, "ICY Radio streamtitle cleaned: %s", cleaned_stream_title
3642 )
3643 streamdetails.stream_title = cleaned_stream_title
3644
3645 # Prefer station-provided cover art from the ICY 'StreamUrl' field (when it is
3646 # an image) over the MusicBrainz artwork lookup in _update_radio_stream_metadata.
3647 image_url = self._parse_icy_image_url(meta_data)
3648
3649 # Parse the original title for structured fields first so stations that announce
3650 # an album can refine the artwork lookup; fall back to the "Artist - Track" split.
3651 album: str | None = None
3652 if parsed := parse_quoted_stream_title(stream_title):
3653 track_name, artist_name_raw, album = parsed
3654 elif " - " in cleaned_stream_title:
3655 artist_name_raw, track_name = (
3656 part.strip() for part in cleaned_stream_title.split(" - ", 1)
3657 )
3658 else:
3659 return
3660
3661 if artist_name_raw and track_name:
3662 self.logger.debug(
3663 "ICY metadata: artist='%s', track='%s', album='%s'",
3664 artist_name_raw,
3665 track_name,
3666 album,
3667 )
3668 self._update_radio_stream_metadata(
3669 streamdetails,
3670 artist=artist_name_raw,
3671 title=track_name,
3672 album=album,
3673 image_url=image_url,
3674 )
3675
3676 def _parse_icy_image_url(self, meta_data: bytes) -> str | None:
3677 """
3678 Return a PNG or JPEG cover-art URL from the ICY 'StreamUrl' field, if present.
3679
3680 :param meta_data: Raw metadata bytes from an ICY stream chunk.
3681 """
3682 # The trailing semicolon is optional to match sources that omit it.
3683 stream_url_re = re.search(rb"StreamUrl='([^']*)'", meta_data)
3684 if not stream_url_re:
3685 return None
3686 try:
3687 image_url = stream_url_re.group(1).decode("utf-8").strip()
3688 except UnicodeDecodeError:
3689 return None
3690 if not image_url:
3691 return None
3692 # StreamUrl is not a standardized artwork field (reference clients such as VLC
3693 # ignore it and it conventionally holds a station website link), so only accept
3694 # values that point at a PNG or JPEG image.
3695 parsed = urlparse(image_url)
3696 if parsed.scheme not in ("http", "https"):
3697 return None
3698 if not parsed.path.lower().endswith((".png", ".jpg", ".jpeg")):
3699 return None
3700 self.logger.debug("ICY metadata: StreamUrl image='%s'", image_url)
3701 return image_url
3702
3703 async def _validate_shoutcast_stream(self, url: str) -> bool:
3704 """
3705 Return True if the URL responds with a legacy Shoutcast "ICY 200 OK" line.
3706
3707 :param url: The URL to validate.
3708 """
3709 try:
3710 parsed = urlparse(url)
3711 host = parsed.hostname
3712 port = parsed.port or 80
3713 path = parsed.path or "/"
3714 if parsed.query:
3715 path = f"{path}?{parsed.query}"
3716
3717 # Open raw socket connection with timeout
3718 reader, writer = await asyncio.wait_for(asyncio.open_connection(host, port), timeout=10)
3719 try:
3720 # Send minimal HTTP request with ICY metadata header
3721 request = f"GET {path} HTTP/1.1\r\nHost: {host}\r\nIcy-MetaData: 1\r\n\r\n"
3722 writer.write(request.encode())
3723 await writer.drain()
3724
3725 # Read just the response line
3726 response_line = await asyncio.wait_for(reader.readline(), timeout=5)
3727 finally:
3728 writer.close()
3729 await writer.wait_closed()
3730
3731 # Check if response starts with "ICY"
3732 decoded_line = response_line.decode("latin-1", errors="ignore").strip()
3733 return decoded_line.startswith("ICY")
3734
3735 except TimeoutError:
3736 self.logger.debug("Timeout during Shoutcast validation for %s", url)
3737 return False
3738 except OSError, ConnectionError:
3739 self.logger.debug("Connection failed during Shoutcast validation for %s", url)
3740 return False
3741 except UnicodeDecodeError:
3742 self.logger.debug("Invalid response encoding during Shoutcast validation for %s", url)
3743 return False
3744
3745 def _resolve_player_dsp_config(self, player: Player) -> DSPConfig:
3746 """
3747 Resolve the effective DSP config for a player.
3748
3749 Single source of truth shared by every code path that needs to know
3750 whether DSP will run for this player. Protocol wrappers defer to their
3751 parent player; single-leg ``player_group`` instances that don't expose
3752 ``MULTI_DEVICE_DSP`` defer to their first member; players whose grouping
3753 context prevents DSP get a disabled config back regardless.
3754
3755 :param player: The player to resolve DSP config for.
3756 """
3757 dsp_player_id = self._resolve_player_dsp_config_id(player)
3758 dsp = self.mass.config.get_player_dsp_config(dsp_player_id)
3759 if is_grouping_preventing_dsp(player):
3760 dsp.enabled = False
3761 elif player.provider.domain == "player_group" and (
3762 PlayerFeature.MULTI_DEVICE_DSP not in player.state.supported_features
3763 ):
3764 if not player.state.group_members:
3765 dsp.enabled = False
3766 return dsp
3767
3768 def _resolve_player_dsp_config_id(self, player: Player) -> str:
3769 """
3770 Return the player identifier that supplies the effective DSP config.
3771
3772 :param player: Player whose DSP config source should be resolved.
3773 """
3774 dsp_player_id = player.protocol_parent_id or player.player_id
3775 if (
3776 not is_grouping_preventing_dsp(player)
3777 and player.provider.domain == "player_group"
3778 and PlayerFeature.MULTI_DEVICE_DSP not in player.state.supported_features
3779 and player.state.group_members
3780 ):
3781 child_player = self.mass.players.get_player(player.state.group_members[0])
3782 assert child_player is not None
3783 dsp_player_id = child_player.player_id
3784 return dsp_player_id
3785
3786 def _get_output_channels(self, player: Player | None, player_id: str) -> str:
3787 """
3788 Return the configured output channels for the rendering player.
3789
3790 The value may be stored on the rendering player(protocol) itself (the
3791 protocol section of the config UI) or on its visible parent player (the
3792 native section); the rendering player's own stored value wins.
3793 """
3794 parent_id = player.protocol_parent_id if player and player.protocol_parent_id else player_id
3795 parent_value = self.mass.config.get_raw_player_config_value(
3796 parent_id, CONF_OUTPUT_CHANNELS, "stereo"
3797 )
3798 return self.mass.config.get_raw_player_config_value(
3799 player.player_id if player else player_id, CONF_OUTPUT_CHANNELS, parent_value
3800 )
3801
3802 def _pick_pcm_bit_depth(
3803 self,
3804 players: Iterable[Player],
3805 streamdetails: StreamDetails | None,
3806 crossfade_enabled: bool,
3807 overlay_active: bool = False,
3808 ) -> tuple[ContentType, int]:
3809 """
3810 Return ``(content_type, bit_depth)`` for an internal PCM stream.
3811
3812 F32 is chosen when audio processing (crossfade, audio overlay, volume
3813 normalization, DSP) will run on the stream â those need the extra
3814 headroom to avoid clipping and precision loss. Otherwise the source's
3815 native bit depth is reused so we don't waste memory upcasting a 16-bit
3816 stream to 32-bit just to pass it through. When the source is unknown
3817 (no streamdetails) we fall back to F32 conservatively.
3818 """
3819 if streamdetails is None:
3820 return INTERNAL_PCM_FORMAT.content_type, INTERNAL_PCM_FORMAT.bit_depth
3821 needs_headroom = (
3822 crossfade_enabled
3823 or overlay_active
3824 or streamdetails.volume_normalization_mode != VolumeNormalizationMode.DISABLED
3825 or any(self._resolve_player_dsp_config(player).enabled for player in players)
3826 )
3827 if needs_headroom:
3828 return INTERNAL_PCM_FORMAT.content_type, INTERNAL_PCM_FORMAT.bit_depth
3829 bit_depth = streamdetails.audio_format.bit_depth
3830 return ContentType.from_bit_depth(bit_depth), bit_depth
3831
3832 def _select_audio_source_pcm_format(
3833 self,
3834 player: Player,
3835 streamdetails: StreamDetails,
3836 supported_sample_rates: Iterable[int] | None = None,
3837 ) -> AudioFormat:
3838 """
3839 Return a passthrough PCM format for a realtime AudioSource item.
3840
3841 The format matches the source's native sample rate, bit depth and
3842 channel count whenever the player can accept them; if the player does
3843 not support the source's sample rate, it is snapped down to the
3844 closest supported rate. No F32 widening â realtime sources skip every
3845 processing stage that would otherwise need it. Surround sources are
3846 still folded down to stereo, which every output path requires anyway.
3847
3848 :param player: The player requesting the stream.
3849 :param streamdetails: Stream details for the AudioSource item.
3850 :param supported_sample_rates: Rates shared by every output player, if applicable.
3851 """
3852 resolved_sample_rates = (
3853 list(supported_sample_rates)
3854 if supported_sample_rates is not None
3855 else [sample_rate for sample_rate, _ in player.get_supported_sample_rates()]
3856 )
3857 source_rate = streamdetails.audio_format.sample_rate
3858 if source_rate in resolved_sample_rates:
3859 output_sample_rate = source_rate
3860 else:
3861 output_sample_rate = max(
3862 (rate for rate in resolved_sample_rates if rate <= source_rate),
3863 default=min(resolved_sample_rates),
3864 )
3865 bit_depth = streamdetails.audio_format.bit_depth
3866 return AudioFormat(
3867 content_type=ContentType.from_bit_depth(bit_depth),
3868 sample_rate=output_sample_rate,
3869 bit_depth=bit_depth,
3870 # a realtime source may announce more channels than anything downstream can
3871 # carry (a VBAN stream can be configured up to 8), and player handoff formats
3872 # copy this count straight through, so fold it here
3873 channels=min(streamdetails.audio_format.channels, 2),
3874 )
3875
3876 def _flow_restart_context(
3877 self, queue_id: str, protocol_player: Player | None
3878 ) -> tuple[str, list[int]]:
3879 """
3880 Resolve the flow mode config and supported sample rates for restart decisions.
3881
3882 Prefers the protocol player actually consuming the flow stream over the
3883 queue's (wrapper) player, whose config may lack the audio specific entries.
3884 """
3885 if protocol_player is None:
3886 protocol_player = self.mass.players.get_player(queue_id)
3887 if protocol_player is None:
3888 flow_mode_sample_rate_conf = self.mass.config.get_raw_player_config_value(
3889 queue_id, CONF_FLOW_MODE_SAMPLE_RATE, FLOW_MODE_SAMPLE_RATE_SMART
3890 )
3891 return flow_mode_sample_rate_conf, []
3892 flow_mode_sample_rate_conf = cast(
3893 "str",
3894 protocol_player.config.get_value(
3895 CONF_FLOW_MODE_SAMPLE_RATE, FLOW_MODE_SAMPLE_RATE_SMART
3896 ),
3897 )
3898 supported_sample_rates = sorted(
3899 {sr for sr, _ in protocol_player.get_supported_sample_rates()}
3900 )
3901 return flow_mode_sample_rate_conf, supported_sample_rates
3902
3903 def _flow_stream_needs_restart(
3904 self,
3905 queue_track: QueueItem,
3906 pcm_format: AudioFormat,
3907 supported_sample_rates: list[int],
3908 flow_mode_sample_rate_conf: str,
3909 is_first_track: bool,
3910 ) -> bool:
3911 """
3912 Return True if the upcoming queue track requires exiting the flow stream.
3913
3914 Covers every case where the flow loop should break and hand control back to
3915 the queue controller for restart:
3916
3917 - Live media (radio, audio sources): cannot be played inside a flow,
3918 the controller will fall back to a single-item stream.
3919 - Sample rate mismatch ('smart' / 'bit_perfect' modes only): the next
3920 track's sample rate (snapped up to the closest supported player rate,
3921 mirroring select_flow_pcm_format's anchoring logic) is incompatible with
3922 the current flow rate, so a new flow must be opened.
3923
3924 The first (anchor) track is always allowed to continue for the sample
3925 rate check; select_flow_pcm_format has already snapped the flow rate to it.
3926
3927 :param queue_track: The upcoming queue item.
3928 :param pcm_format: The current flow stream's PCM format.
3929 :param supported_sample_rates: Sorted list of the player's supported rates.
3930 :param flow_mode_sample_rate_conf: The flow mode sample rate config value.
3931 :param is_first_track: Whether this is the first track of the flow stream.
3932 """
3933 # live audio (radio, plugin or audio source) cannot be flowed; let the
3934 # queue controller fall back to single-item streaming for this item
3935 if queue_track.media_type in (MediaType.RADIO, MediaType.AUDIO_SOURCE):
3936 self.logger.info(
3937 "Live media item %s (%s, %s) encountered in flow stream "
3938 "- breaking out to single item stream",
3939 queue_track.queue_item_id,
3940 queue_track.name,
3941 queue_track.media_type,
3942 )
3943 return True
3944
3945 if is_first_track or queue_track.streamdetails is None:
3946 return False
3947 raw_next_rate = queue_track.streamdetails.audio_format.sample_rate
3948 if not raw_next_rate or not supported_sample_rates:
3949 return False
3950 effective_next_rate = _snap_supported_rate_up(raw_next_rate, supported_sample_rates)
3951
3952 # branch order mirrors select_flow_pcm_format: fixed-rate modes resample
3953 # everything to the chosen rate (no restart); bit_perfect restarts on any
3954 # mismatch; anything else falls through to smart-anchor behavior so
3955 # unknown/legacy config values don't silently pin the flow forever.
3956 if flow_mode_sample_rate_conf in (
3957 FLOW_MODE_SAMPLE_RATE_48000,
3958 FLOW_MODE_SAMPLE_RATE_96000,
3959 FLOW_MODE_SAMPLE_RATE_HIGHEST,
3960 ):
3961 needs_restart = False
3962 elif flow_mode_sample_rate_conf == FLOW_MODE_SAMPLE_RATE_BIT_PERFECT:
3963 needs_restart = effective_next_rate != pcm_format.sample_rate
3964 else:
3965 needs_restart = effective_next_rate > pcm_format.sample_rate
3966
3967 if needs_restart:
3968 self.logger.info(
3969 "Track %s (%s) sample rate %s (snapped to %s) incompatible with flow rate %s "
3970 "(mode: %s) - breaking out to restart flow stream",
3971 queue_track.queue_item_id,
3972 queue_track.name,
3973 raw_next_rate,
3974 effective_next_rate,
3975 pcm_format.sample_rate,
3976 flow_mode_sample_rate_conf,
3977 )
3978 return needs_restart
3979
3980 @asynccontextmanager
3981 async def _connect_radio_stream(self, url: str, **kwargs: Any) -> AsyncGenerator[Any]:
3982 """
3983 Connect to a radio stream URL with fallback for legacy SSL/TLS configurations.
3984
3985 Some radio servers use outdated TLS configurations that reject modern
3986 cipher suites. Since radio streams are public broadcast content,
3987 relaxing cipher requirements is acceptable.
3988
3989 :param url: The radio stream URL to connect to.
3990 :param kwargs: Additional keyword arguments passed to aiohttp get().
3991 """
3992 request_url = encoded_request_url(url)
3993 try:
3994 async with self.mass.http_session_no_ssl.get(request_url, **kwargs) as resp:
3995 yield resp
3996 except ClientConnectorSSLError:
3997 self.logger.info(
3998 "SSL handshake failed for %s, retrying with permissive cipher configuration", url
3999 )
4000 insecure_ssl_context = ssl_util.client_context_no_verify(
4001 ssl_util.SSLCipherList.INSECURE
4002 )
4003 async with self.mass.http_session_no_ssl.get(
4004 request_url, ssl=insecure_ssl_context, **kwargs
4005 ) as resp:
4006 yield resp
4007
4008 async def _update_hls_radio_metadata(
4009 self,
4010 streamdetails: StreamDetails,
4011 elapsed_time: int,
4012 ) -> None:
4013 """
4014 Update HLS radio stream metadata by fetching the playlist.
4015
4016 Fetches the HLS playlist and extracts metadata from EXTINF lines.
4017
4018 :param streamdetails: StreamDetails object to update with metadata
4019 :param elapsed_time: Current playback position in seconds (unused for live radio)
4020 """
4021 mass = self.mass
4022 try:
4023 # Get the actual media playlist URL from cache or resolve it
4024 # We cache the media_playlist_url in streamdetails.data to avoid re-resolving
4025 if streamdetails.data is None:
4026 streamdetails.data = {}
4027 media_playlist_url = streamdetails.data.get("hls_media_playlist_url")
4028 if not media_playlist_url:
4029 try:
4030 assert isinstance(streamdetails.path, str) # for type checking
4031 substream = await self.get_hls_substream(streamdetails.path)
4032 media_playlist_url = substream.path
4033 streamdetails.data["hls_media_playlist_url"] = media_playlist_url
4034 except Exception as err:
4035 self.logger.warning(
4036 "Failed to resolve HLS substream for metadata monitoring: %s", err
4037 )
4038 return
4039
4040 # Fetch the media playlist
4041 timeout = ClientTimeout(total=0, connect=10, sock_read=30)
4042 try:
4043 async with mass.http_session_no_ssl.get(
4044 encoded_request_url(media_playlist_url), timeout=timeout
4045 ) as resp:
4046 resp.raise_for_status()
4047 playlist_content = await resp.text()
4048 except ClientResponseError as err:
4049 # Session token likely expired (410/403) â drop cache so next poll re-resolves
4050 if err.status in (403, 410):
4051 streamdetails.data.pop("hls_media_playlist_url", None)
4052 raise
4053
4054 # Parse the playlist and look for EXTINF metadata
4055 # The most recent segment usually has the current metadata
4056 lines = playlist_content.strip().split("\n")
4057 for line in reversed(lines):
4058 if line.startswith("#EXTINF:"):
4059 # Extract metadata from EXTINF line
4060 metadata = parse_extinf_metadata(line)
4061
4062 # Build stream title from title and artist
4063 title = metadata.get("title", "")
4064 artist = metadata.get("artist", "")
4065 image_url = (
4066 metadata.get("image") or metadata.get("artwork") or metadata.get("cover")
4067 )
4068 if not artist and " - " in title:
4069 artist, title = title.split(" - ", 1)
4070 if title or artist:
4071 # Format as "Artist - Title"
4072 if artist and title:
4073 stream_title = f"{artist} - {title}"
4074 elif title:
4075 stream_title = title
4076 else:
4077 stream_title = artist
4078
4079 # Clean the stream title
4080 cleaned_title = clean_stream_title(stream_title)
4081
4082 # Only update if changed
4083 if cleaned_title != streamdetails.stream_title and cleaned_title:
4084 self.logger.log(
4085 VERBOSE_LOG_LEVEL, "HLS Radio metadata updated: %s", cleaned_title
4086 )
4087 streamdetails.stream_title = cleaned_title
4088 self._update_radio_stream_metadata(
4089 streamdetails,
4090 artist=artist or None,
4091 title=title or cleaned_title,
4092 image_url=image_url,
4093 )
4094
4095 # Only check the most recent EXTINF
4096 break
4097
4098 except Exception as err:
4099 self.logger.debug("Error fetching HLS metadata: %s", err)
4100
4101 @staticmethod
4102 def _normalize_reconnecting_urls(url: str | list[MultiPartPath]) -> list[str]:
4103 """Normalize a single URL or a sequence into a non-empty list."""
4104 if isinstance(url, str):
4105 return [url]
4106 if not url:
4107 msg = "Radio stream requires at least one URL"
4108 raise InvalidDataError(msg)
4109 return [part.path for part in url]
4110
4111 async def _resolve_overlay_input(self, queue: PlayerQueue) -> str | None:
4112 """
4113 Resolve the queue's overlay source to a file path or URL for ffmpeg.
4114
4115 Returns None (with a warning logged) when the source can not be resolved,
4116 so the caller can degrade to music-only playback.
4117 """
4118 if not (mapping := queue.overlay_source):
4119 return None
4120 try:
4121 provider = self.mass.get_provider(mapping.provider)
4122 if provider is None:
4123 raise MediaNotFoundError(f"Provider {mapping.provider} is not available")
4124 stream_prov = cast("MusicProvider | PluginProvider", provider)
4125 streamdetails = await stream_prov.get_stream_details(
4126 mapping.item_id, MediaType.SOUND_EFFECT
4127 )
4128 except Exception as err:
4129 self.logger.warning(
4130 "Audio overlay source %s is unavailable (%s) - continuing without overlay",
4131 mapping.uri,
4132 str(err) or err.__class__.__name__,
4133 )
4134 return None
4135 if streamdetails.stream_type not in (StreamType.LOCAL_FILE, StreamType.HTTP) or not (
4136 isinstance(streamdetails.path, str)
4137 ):
4138 self.logger.warning(
4139 "Audio overlay source %s uses unsupported stream type %s "
4140 "- continuing without overlay",
4141 mapping.uri,
4142 streamdetails.stream_type,
4143 )
4144 return None
4145 if streamdetails.stream_type == StreamType.LOCAL_FILE and not await aiofiles.os.path.isfile(
4146 streamdetails.path
4147 ):
4148 # guard against stale sources: feeding a missing file to the mixer would
4149 # kill the whole (music) stream instead of just the overlay
4150 self.logger.warning(
4151 "Audio overlay source %s does not exist - continuing without overlay",
4152 streamdetails.path,
4153 )
4154 return None
4155 return streamdetails.path
4156