/
/
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 as err:
1266 streamdetails.stream_error = True
1267 if raise_on_error:
1268 raise
1269 logger.exception(
1270 "Unexpected error while streaming AudioSource %s (%s): %s",
1271 queue_item.name,
1272 streamdetails.uri,
1273 err,
1274 )
1275 finally:
1276 streamdetails.seconds_streamed = bytes_received / pcm_format.pcm_sample_size
1277
1278 async def get_queue_item_stream(
1279 self,
1280 queue_item: QueueItem,
1281 pcm_format: AudioFormat,
1282 seek_position: float = 0,
1283 playback_speed: float = 1.0,
1284 raise_on_error: bool = True,
1285 normalization_override: VolumeNormalizationMode | None = None,
1286 session_id: str | None = None,
1287 prepared_buffer: AudioBuffer | None = None,
1288 exact_seek: bool = False,
1289 ) -> AsyncGenerator[bytes]:
1290 """
1291 Get the (PCM) audio stream for a single queue item.
1292
1293 Audio is always served from the AudioBuffer which stores raw decoded PCM.
1294 Volume normalization and other filters are applied on-the-fly when reading
1295 from the buffer.
1296
1297 AudioSource items dispatch to ``get_audio_source_stream`` instead: they
1298 are realtime and bypass the buffering/normalization/filter machinery.
1299
1300 :param normalization_override: Force this volume normalization mode instead of
1301 re-evaluating it from the (possibly just-updated) loudness measurement. Used by
1302 the crossfade path to keep a track's replayed intro and its body on the same mode.
1303 :param session_id: Queue session that owns processing-detail updates.
1304 :param prepared_buffer: Existing buffer that must be used without opening a new source.
1305 :param exact_seek: Preserve millisecond precision instead of user-seek quantization.
1306 """
1307 streamdetails = queue_item.streamdetails
1308 assert streamdetails
1309
1310 # streamdetails are cached and reused for retries; reset this before any
1311 # media-type-specific dispatch so AudioSource failures do not stick.
1312 streamdetails.stream_error = False
1313
1314 if queue_item.media_type == MediaType.AUDIO_SOURCE:
1315 async for chunk in self.get_audio_source_stream(
1316 queue_item=queue_item,
1317 pcm_format=pcm_format,
1318 raise_on_error=raise_on_error,
1319 ):
1320 yield chunk
1321 return
1322 filter_params: list[str] = []
1323
1324 logger = self.logger.getChild("queue_item_stream")
1325
1326 if normalization_override is not None:
1327 # crossfade path pins the body to the intro's mode; skip hydration/re-eval that could flip it
1328 streamdetails.volume_normalization_mode = normalization_override
1329 else:
1330 # hydrate loudness from audio analysis (just-in-time, so that a measurement
1331 # completed during a previous play is picked up here). A live analyzer run
1332 # may have already populated streamdetails.loudness in memory â don't clobber
1333 # that, and don't clobber a value set upstream by the music provider.
1334 if streamdetails.loudness is None:
1335 if analysis := await self.mass.streams.audio_analysis.get_audio_analysis(
1336 streamdetails.item_id,
1337 streamdetails.provider,
1338 media_type=streamdetails.media_type,
1339 # use the authoritative EBU R128 value, not another provider's loudness proxy
1340 priority=(LOUDNESS_ANALYSIS_DOMAIN,),
1341 ):
1342 if analysis.loudness_integrated is not None:
1343 streamdetails.loudness = round(analysis.loudness_integrated, 2)
1344 if analysis.loudness_album is not None and streamdetails.loudness_album is None:
1345 streamdetails.loudness_album = round(analysis.loudness_album, 2)
1346
1347 # re-evaluate normalization mode: the background loudness analyzer may have
1348 # updated streamdetails.loudness since get_stream_details was called
1349 if streamdetails.queue_id:
1350 volume_normalization_enabled = (
1351 self.mass.config.get_effective_player_queue_config_value(
1352 streamdetails.queue_id, CONF_VOLUME_NORMALIZATION, CONF_VALUE_ENABLED
1353 )
1354 != CONF_VALUE_DISABLED
1355 )
1356 streamdetails.volume_normalization_mode = get_normalization_mode(
1357 self._get_volume_normalization_preference(streamdetails),
1358 volume_normalization_enabled,
1359 streamdetails,
1360 )
1361
1362 # get or create the AudioBuffer (stores raw decoded PCM). This runs before the
1363 # filters are built because a source-capacity reselection can hand back another
1364 # provider's streamdetails, which everything below must then work with.
1365 seek_position_ms = int(seek_position * 1000)
1366 try:
1367 if prepared_buffer is not None:
1368 if streamdetails.buffer is not prepared_buffer or not prepared_buffer.is_valid(
1369 seek_position_ms
1370 ):
1371 raise AudioError("Prepared crossfade buffer is no longer available")
1372 audio_buffer = prepared_buffer
1373 else:
1374 audio_buffer = await self.get_audio_buffer(
1375 queue_item, seek_position_ms=seek_position_ms, reason="streaming"
1376 )
1377 except AudioError as err:
1378 streamdetails.stream_error = True
1379 if raise_on_error:
1380 raise
1381 logger.error(
1382 "AudioError while preparing queue item %s (%s): %s",
1383 queue_item.name,
1384 streamdetails.uri,
1385 err,
1386 )
1387 return
1388 streamdetails = queue_item.streamdetails
1389 assert streamdetails # for type checking
1390 if normalization_override is not None:
1391 # a capacity reselection hands back freshly resolved details, so the
1392 # crossfade's intro/body normalization pin must be re-applied to them
1393 streamdetails.volume_normalization_mode = normalization_override
1394
1395 # handle volume normalization
1396 gain_correct: float | None = None
1397 if streamdetails.volume_normalization_mode == VolumeNormalizationMode.DYNAMIC:
1398 filter_rule = (
1399 f"loudnorm=I={streamdetails.target_loudness}"
1400 ":TP=-2.0:LRA=10.0:offset=0.0:print_format=json"
1401 )
1402 filter_params.append(filter_rule)
1403 elif streamdetails.volume_normalization_mode == VolumeNormalizationMode.FIXED_GAIN:
1404 config_key = (
1405 CONF_VOLUME_NORMALIZATION_FIXED_GAIN_TRACKS
1406 if streamdetails.media_type == MediaType.TRACK
1407 else CONF_VOLUME_NORMALIZATION_FIXED_GAIN_RADIO
1408 )
1409 gain_value = self.mass.streams.get_config_value(config_key, return_type=float)
1410 gain_correct = round(gain_value, 2)
1411 filter_params.append(f"volume={gain_correct}dB")
1412 elif streamdetails.volume_normalization_mode == VolumeNormalizationMode.MEASUREMENT_ONLY:
1413 target_loudness = (
1414 float(streamdetails.target_loudness)
1415 if streamdetails.target_loudness is not None
1416 else 0.0
1417 )
1418 if streamdetails.prefer_album_loudness and streamdetails.loudness_album is not None:
1419 gain_correct = target_loudness - float(streamdetails.loudness_album)
1420 elif streamdetails.loudness is not None:
1421 gain_correct = target_loudness - float(streamdetails.loudness)
1422 else:
1423 gain_correct = 0.0
1424 gain_correct = round(gain_correct, 2)
1425 filter_params.append(f"volume={gain_correct}dB")
1426 streamdetails.volume_normalization_gain_correct = gain_correct
1427
1428 # handle playback speed
1429 if playback_speed != 1.0:
1430 filter_params.append(f"atempo={playback_speed}")
1431
1432 # handle optional fade-in
1433 if streamdetails.fade_in:
1434 filter_params.insert(0, "afade=type=in:start_time=0:duration=3")
1435
1436 logger.log(
1437 VERBOSE_LOG_LEVEL,
1438 "Starting queue item stream for %s (%s)"
1439 " - using fade-in: %s"
1440 " - using volume normalization: %s"
1441 " - using playback speed: %s",
1442 queue_item.name,
1443 streamdetails.uri,
1444 streamdetails.fade_in,
1445 streamdetails.volume_normalization_mode,
1446 playback_speed,
1447 )
1448
1449 if (
1450 streamdetails.queue_id
1451 and (queue_data := self.mass.player_queues.queue_data_or_none(streamdetails.queue_id))
1452 and (processing_session_id := session_id or queue_data.session_id)
1453 ):
1454 self.mass.streams.audio_processing.update_item_runtime(
1455 queue_id=streamdetails.queue_id,
1456 session_id=processing_session_id,
1457 queue_item_id=queue_item.queue_item_id,
1458 input_format=audio_buffer.pcm_format,
1459 pcm_format=pcm_format,
1460 normalization=get_normalization_details(streamdetails, gain_correct),
1461 playback_speed=playback_speed,
1462 alters_audio=streamdetails.fade_in,
1463 )
1464 # read from buffer with filters applied (volume normalization, speed, fade-in, etc.)
1465 # if no processing needed, this yields directly from the buffer
1466 media_stream_gen = audio_buffer.get_stream(
1467 output_format=pcm_format,
1468 seek_position_ms=seek_position_ms,
1469 filter_params=filter_params or None,
1470 exact_seek=exact_seek,
1471 )
1472
1473 first_chunk_received = False
1474 bytes_received = 0
1475 finished = False
1476 next_buffer_triggered = False
1477 stream_started_at = asyncio.get_event_loop().time()
1478 try:
1479 async for chunk in media_stream_gen:
1480 bytes_received += len(chunk)
1481 if not first_chunk_received:
1482 first_chunk_received = True
1483 logger.log(
1484 VERBOSE_LOG_LEVEL,
1485 "First audio chunk received for %s (%s) after %.2f seconds",
1486 queue_item.name,
1487 streamdetails.uri,
1488 asyncio.get_event_loop().time() - stream_started_at,
1489 )
1490 # trigger pre-buffering of the next item well before end
1491 # to ensure the raw PCM is ready when the next item needs to be streamed.
1492 # tracks and sound effects are finite files that fill and close immediately;
1493 # live sources (radio, audio_source) open an upstream connection that would
1494 # sit idle and likely time out before the player actually consumes it.
1495 if (
1496 not next_buffer_triggered
1497 and streamdetails.duration
1498 and (queue := self.mass.player_queues.get_active_queue(queue_item.queue_id))
1499 and queue.next_item
1500 and queue.next_item.queue_item_id != queue_item.queue_item_id
1501 and queue.next_item.media_type in (MediaType.TRACK, MediaType.SOUND_EFFECT)
1502 and (bytes_received / pcm_format.pcm_sample_size + seek_position)
1503 >= streamdetails.duration - 60
1504 ):
1505 next_buffer_triggered = True
1506 self.mass.player_queues.prepare_next_audio_buffer(queue_item.queue_id)
1507 yield chunk
1508 del chunk
1509 finished = True
1510 except AudioError as err:
1511 streamdetails.stream_error = True
1512 # revoke availability when the stream never produced any audio
1513 if bytes_received == 0 and not isinstance(err, ProviderStreamLimitError):
1514 queue_item.available = False
1515 if raise_on_error:
1516 raise
1517 logger.error(
1518 "AudioError while streaming queue item %s (%s): %s",
1519 queue_item.name,
1520 streamdetails.uri,
1521 err,
1522 )
1523 except asyncio.CancelledError:
1524 raise
1525 except Exception as err:
1526 streamdetails.stream_error = True
1527 if raise_on_error:
1528 raise
1529 logger.exception(
1530 "Unexpected error while streaming queue item %s (%s): %s",
1531 queue_item.name,
1532 streamdetails.uri,
1533 err,
1534 )
1535 finally:
1536 seconds_streamed = bytes_received / pcm_format.pcm_sample_size
1537 streamdetails.seconds_streamed = seconds_streamed
1538 logger.log(
1539 VERBOSE_LOG_LEVEL,
1540 "stream %s for %s in %.2f seconds - seconds streamed/buffered: %.2f",
1541 "aborted" if not finished else "finished",
1542 streamdetails.uri,
1543 asyncio.get_event_loop().time() - stream_started_at,
1544 seconds_streamed,
1545 )
1546 self._notify_provider_streamed(streamdetails, finished, seconds_streamed)
1547
1548 async def get_queue_item_stream_with_smartfade(
1549 self,
1550 player: Player,
1551 queue_item: QueueItem,
1552 pcm_format: AudioFormat,
1553 crossfade_mode: CrossfadeMode = CrossfadeMode.SMART_CROSSFADE,
1554 standard_crossfade_duration: int = 10,
1555 session_id: str | None = None,
1556 ) -> AsyncGenerator[bytes]:
1557 """
1558 Return one queue item with a crossfade into the next item.
1559
1560 :param player: Player consuming the stream.
1561 :param queue_item: Queue item to stream.
1562 :param pcm_format: Shared PCM format.
1563 :param crossfade_mode: Effective crossfade mode.
1564 :param standard_crossfade_duration: Configured standard crossfade duration.
1565 :param session_id: Queue session that owns processing-detail updates.
1566 """
1567 queue = self.mass.player_queues.get(queue_item.queue_id)
1568 if not queue:
1569 raise RuntimeError(f"Queue {queue_item.queue_id} not found")
1570
1571 streamdetails = queue_item.streamdetails
1572 assert streamdetails
1573 crossfade_data = self._crossfade_data.get(queue.queue_id)
1574
1575 if crossfade_data and streamdetails.seek_position > 0:
1576 # don't do crossfade when seeking into track
1577 self.logger.debug(
1578 "Discarding crossfade data for queue %s - seeking into track (pos=%s)",
1579 queue.display_name,
1580 streamdetails.seek_position,
1581 )
1582 crossfade_data = None
1583 if crossfade_data and (crossfade_data.queue_item_id != queue_item.queue_item_id):
1584 # edge case alert: the next item changed just while we were preloading/crossfading
1585 self.logger.warning(
1586 "Skipping crossfade data for queue %s - next item changed!"
1587 " (expected queue_item_id=%s, got=%s)",
1588 queue.display_name,
1589 crossfade_data.queue_item_id,
1590 queue_item.queue_item_id,
1591 )
1592 crossfade_data = None
1593 self._crossfade_data.pop(queue.queue_id, None)
1594 elif not crossfade_data:
1595 self.logger.debug(
1596 "No crossfade data available for queue %s (queue_item_id=%s)",
1597 queue.display_name,
1598 queue_item.queue_item_id,
1599 )
1600
1601 self.logger.debug(
1602 "Start Streaming queue track: %s (%s) for queue %s on player %s"
1603 "- crossfade mode: %s "
1604 "- crossfading from previous track: %s ",
1605 queue_item.streamdetails.uri if queue_item.streamdetails else "Unknown URI",
1606 queue_item.name,
1607 queue.display_name,
1608 player.name,
1609 crossfade_mode,
1610 "true" if crossfade_data else "false",
1611 )
1612
1613 buffer = bytearray()
1614 bytes_written = 0
1615 # calculate crossfade buffer size
1616 crossfade_buffer_duration = (
1617 SMART_CROSSFADE_DURATION
1618 if crossfade_mode == CrossfadeMode.SMART_CROSSFADE
1619 else standard_crossfade_duration
1620 )
1621 crossfade_buffer_duration = min(
1622 crossfade_buffer_duration,
1623 int(streamdetails.duration / 2)
1624 if streamdetails.duration
1625 else crossfade_buffer_duration,
1626 )
1627 # skip crossfade if buffer would be too small to be meaningful
1628 if crossfade_buffer_duration < MIN_CROSSFADE_FALLBACK_DURATION:
1629 crossfade_buffer_duration = 0
1630 # Ensure crossfade buffer size is aligned to frame boundaries
1631 # Frame size = bytes_per_sample * channels
1632 bytes_per_sample = pcm_format.bit_depth // 8
1633 frame_size = bytes_per_sample * pcm_format.channels
1634 crossfade_buffer_size = int(pcm_format.pcm_sample_size * crossfade_buffer_duration)
1635 # Round down to nearest frame boundary
1636 crossfade_buffer_size = (crossfade_buffer_size // frame_size) * frame_size
1637 fade_out_data: bytes | None = None
1638 uncredited_tail_bytes = 0
1639
1640 # pin the body to DYNAMIC when the intro was baked DYNAMIC,
1641 # else a late measurement flips it and causes a volume jump
1642 norm_override: VolumeNormalizationMode | None = None
1643 if crossfade_data and crossfade_data.normalization_mode == VolumeNormalizationMode.DYNAMIC:
1644 norm_override = VolumeNormalizationMode.DYNAMIC
1645
1646 exact_buffer_seek = crossfade_data is not None
1647 if crossfade_data:
1648 # reported media-time (TRIM + CF) is decoupled from the raw buffer seek below (X)
1649 streamdetails.seek_position = crossfade_data.elapsed_time_offset
1650 # yield the POST portion (resample if previous track's format differs)
1651 if crossfade_data.pcm_format != pcm_format:
1652 async for _chunk in resample_pcm_audio(
1653 crossfade_data.data, crossfade_data.pcm_format, pcm_format
1654 ):
1655 yield _chunk
1656 bytes_written += len(_chunk)
1657 else:
1658 for pcm_slice in iter_pcm_slices(crossfade_data.data, pcm_format, 1000):
1659 yield pcm_slice
1660 await asyncio.sleep(0)
1661 bytes_written += len(crossfade_data.data)
1662 # skip past the source media already consumed by the crossfade
1663 discard_position = crossfade_data.fade_in_media_duration
1664 crossfade_data = None
1665 self._crossfade_data.pop(queue.queue_id, None)
1666 else:
1667 discard_position = float(streamdetails.seek_position)
1668
1669 # Yield the first WARMUP_DURATION worth of audio immediately so playback starts
1670 # right away. After that, start accumulating the crossfade holdback buffer.
1671 warmup_size = int(pcm_format.pcm_sample_size * WARMUP_DURATION)
1672 warmup_bytes = 0
1673 total_chunks_received = 0
1674 holdback_armed = False
1675 playback_speed = cast("float", queue_item.extra_attributes.get("playback_speed", 1.0))
1676 async for chunk in self.get_queue_item_stream(
1677 queue_item,
1678 pcm_format,
1679 seek_position=discard_position,
1680 playback_speed=playback_speed,
1681 normalization_override=norm_override,
1682 session_id=session_id,
1683 exact_seek=exact_buffer_seek,
1684 ):
1685 total_chunks_received += 1
1686
1687 if warmup_bytes < warmup_size:
1688 # warmup: yield directly, don't buffer
1689 yield chunk
1690 warmup_bytes += len(chunk)
1691 bytes_written += len(chunk)
1692 del chunk
1693 continue
1694
1695 if not holdback_armed:
1696 holdback_armed = self._crossfade_holdback_allowed(
1697 queue_item.streamdetails or streamdetails,
1698 crossfade_buffer_duration,
1699 playback_speed,
1700 )
1701 if not holdback_armed:
1702 # holding audio back now would only shrink the player's lead
1703 yield chunk
1704 bytes_written += len(chunk)
1705 del chunk
1706 continue
1707
1708 buffer.extend(chunk)
1709 del chunk
1710 if len(buffer) < crossfade_buffer_size:
1711 await asyncio.sleep(0)
1712 continue
1713 # yield everything above the crossfade buffer
1714 while len(buffer) > crossfade_buffer_size:
1715 yield bytes(buffer[: pcm_format.pcm_sample_size])
1716 bytes_written += pcm_format.pcm_sample_size
1717 del buffer[: pcm_format.pcm_sample_size]
1718 await asyncio.sleep(0)
1719
1720 #### HANDLE END OF TRACK
1721
1722 # get next track for crossfade
1723 crossfade_start_time = asyncio.get_event_loop().time()
1724 next_queue_item: QueueItem | None
1725 try:
1726 self.logger.debug(
1727 "Preloading NEXT track for crossfade for queue %s", queue.display_name
1728 )
1729 next_queue_item = await self.mass.player_queues.load_next_queue_item(
1730 queue.queue_id, queue_item.queue_item_id
1731 )
1732 # set index_in_buffer to prevent our next track is overwritten while preloading
1733 if next_queue_item.streamdetails is None:
1734 raise InvalidDataError(
1735 f"No streamdetails for next queue item {next_queue_item.queue_item_id}"
1736 )
1737 queue.index_in_buffer = self.mass.player_queues.index_by_id(
1738 queue.queue_id, next_queue_item.queue_item_id
1739 )
1740 except QueueEmpty:
1741 # end of queue reached, no next item
1742 next_queue_item = None
1743
1744 crossfade_allowed = False
1745 transition_mode = CrossfadeMode.DISABLED
1746 fade_in_buffer_duration = 0.0
1747 fade_in_playback_speed = 1.0
1748 # a fade needs enough of the outgoing track to overlap with; a holdback that
1749 # armed late (or not at all) leaves less than that
1750 min_fade_out_size = int(pcm_format.pcm_sample_size * MIN_CROSSFADE_FALLBACK_DURATION)
1751 if len(buffer) >= min_fade_out_size and next_queue_item and next_queue_item.streamdetails:
1752 fade_in_playback_speed = cast(
1753 "float", next_queue_item.extra_attributes.get("playback_speed", 1.0)
1754 )
1755 next_pcm = await self.select_pcm_format(
1756 player=player,
1757 streamdetails=next_queue_item.streamdetails,
1758 crossfade_enabled=True,
1759 )
1760 crossfade_allowed = self.crossfade_allowed(
1761 queue_item,
1762 crossfade_mode=crossfade_mode,
1763 player_id=player.player_id,
1764 flow_mode=False,
1765 next_queue_item=next_queue_item,
1766 sample_rate=pcm_format.sample_rate,
1767 next_sample_rate=next_pcm.sample_rate,
1768 )
1769 if crossfade_allowed:
1770 transition_mode, fade_in_buffer_duration = self._select_buffered_crossfade(
1771 next_queue_item.streamdetails,
1772 crossfade_mode,
1773 standard_crossfade_duration,
1774 fade_in_playback_speed,
1775 )
1776 crossfade_allowed = transition_mode != CrossfadeMode.DISABLED
1777 if not crossfade_allowed:
1778 # no crossfade enabled/allowed, just yield the buffer last part
1779 bytes_written += len(buffer)
1780 for pcm_slice in iter_pcm_slices(bytes(buffer), pcm_format, 1000):
1781 yield pcm_slice
1782 await asyncio.sleep(0)
1783 else:
1784 assert next_queue_item is not None
1785 assert next_queue_item.streamdetails is not None
1786 assert next_queue_item.streamdetails.buffer is not None
1787 fade_in_audio_buffer = cast("AudioBuffer", next_queue_item.streamdetails.buffer)
1788 # the remaining buffer is the fade-out tail of the current track
1789 fade_out_data = bytes(buffer)
1790 buffer = bytearray()
1791 fade_in_buffer_size = int(pcm_format.pcm_sample_size * fade_in_buffer_duration)
1792 fade_in_buffer_size = (fade_in_buffer_size // frame_size) * frame_size
1793 # initialized before the try block â the except handler reads these
1794 first_part_written = 0
1795 second_part_buf = bytearray()
1796 try:
1797 # wrap the next track's stream in a counting generator that caps
1798 # at the resident fade-in size and tracks how many bytes were consumed
1799 fade_in_bytes_consumed = 0
1800
1801 _next_item = next_queue_item
1802
1803 async def _limited_fade_in() -> AsyncGenerator[bytes]:
1804 nonlocal fade_in_bytes_consumed
1805 fade_in_stream = self.get_queue_item_stream(
1806 _next_item,
1807 pcm_format,
1808 playback_speed=fade_in_playback_speed,
1809 session_id=session_id,
1810 prepared_buffer=fade_in_audio_buffer,
1811 )
1812 async with aclosing(fade_in_stream):
1813 async for chunk in fade_in_stream:
1814 remaining = fade_in_buffer_size - fade_in_bytes_consumed
1815 if remaining <= 0:
1816 break
1817 if len(chunk) >= remaining:
1818 fade_in_bytes_consumed += remaining
1819 yield chunk[:remaining]
1820 break
1821 fade_in_bytes_consumed += len(chunk)
1822 yield chunk
1823
1824 smart_fade = await self.smart_fades_mixer.build(
1825 fade_in_streamdetails=next_queue_item.streamdetails,
1826 fade_out_streamdetails=streamdetails,
1827 pcm_format=pcm_format,
1828 standard_crossfade_duration=standard_crossfade_duration,
1829 mode=transition_mode,
1830 fade_out_data=fade_out_data,
1831 fade_in_bytes_len=fade_in_buffer_size,
1832 )
1833 crossfade_timing = smart_fade.timing_info
1834 # Split mix output at end-of-overlap: PRE+CF to A, POST to B's intro.
1835 fadeout_share_bytes = int(
1836 (crossfade_timing.pre_crossfade_duration + crossfade_timing.crossfade_duration)
1837 * pcm_format.pcm_sample_size
1838 )
1839 fadeout_share_bytes = (fadeout_share_bytes // frame_size) * frame_size
1840 async for mix_chunk in self.smart_fades_mixer.mix(
1841 smart_fade,
1842 fade_in_part=_limited_fade_in(),
1843 fade_out_part=fade_out_data,
1844 pcm_format=pcm_format,
1845 ):
1846 if first_part_written < fadeout_share_bytes:
1847 # split this chunk so A gets exactly fadeout_share_bytes
1848 remaining = fadeout_share_bytes - first_part_written
1849 if len(mix_chunk) > remaining:
1850 yield mix_chunk[:remaining]
1851 first_part_written += remaining
1852 bytes_written += remaining
1853 second_part_buf.extend(mix_chunk[remaining:])
1854 else:
1855 yield mix_chunk
1856 first_part_written += len(mix_chunk)
1857 bytes_written += len(mix_chunk)
1858 else:
1859 second_part_buf.extend(mix_chunk)
1860 # tail consumed by the mix but not credited to bytes_written
1861 uncredited_tail_bytes = len(fade_out_data) - first_part_written
1862 self._crossfade_data[queue_item.queue_id] = CrossfadeData(
1863 data=bytes(second_part_buf),
1864 fade_in_media_duration=(fade_in_bytes_consumed / pcm_format.pcm_sample_size)
1865 * fade_in_playback_speed,
1866 pcm_format=pcm_format,
1867 queue_item_id=next_queue_item.queue_item_id,
1868 elapsed_time_offset=(
1869 crossfade_timing.fadein_trimmed_duration
1870 + crossfade_timing.crossfade_duration
1871 )
1872 * fade_in_playback_speed,
1873 normalization_mode=next_queue_item.streamdetails.volume_normalization_mode,
1874 )
1875 crossfade_elapsed = asyncio.get_event_loop().time() - crossfade_start_time
1876 self.logger.debug(
1877 "Stored crossfade data for queue %s"
1878 " - next queue_item_id: %s (preparation took %.1fs)",
1879 queue.display_name,
1880 next_queue_item.queue_item_id,
1881 crossfade_elapsed,
1882 )
1883 except Exception as err:
1884 if first_part_written or second_part_buf:
1885 # partial mix already played â concat'd fade_out_data would duplicate audio
1886 raise
1887 # crossfade failed, fall back to just yielding the fade_out_data
1888 self.logger.warning(
1889 "Crossfade failed for queue %s: %s",
1890 queue.display_name,
1891 err,
1892 )
1893 next_queue_item = None
1894 for pcm_slice in iter_pcm_slices(fade_out_data, pcm_format, 1000):
1895 yield pcm_slice
1896 await asyncio.sleep(0)
1897 bytes_written += len(fade_out_data)
1898 del fade_out_data
1899 # make sure the buffer gets cleaned up
1900 del buffer
1901 # a capacity reselection inside the stream replaces the queue item's details,
1902 # so rebind before the writebacks land on an orphaned object
1903 streamdetails = queue_item.streamdetails or streamdetails
1904 # update duration details based on the actual pcm data we sent
1905 # this also accounts for crossfade and silence stripping
1906 seconds_streamed = bytes_written / pcm_format.pcm_sample_size
1907 streamdetails.seconds_streamed = seconds_streamed
1908 # an externally aborted source ends in a clean EOF mid-track, so the
1909 # streamed length must not be written back as the item's duration
1910 source_buffer = streamdetails.buffer
1911 if source_buffer is None or not source_buffer.cancelled:
1912 uncredited_tail_seconds = uncredited_tail_bytes / pcm_format.pcm_sample_size
1913 # streamdetails.duration is in media-time; seconds_streamed is stream-time
1914 # (post-atempo), so we scale by playback_speed to recover media-time.
1915 streamdetails.duration = int(
1916 streamdetails.seek_position
1917 + (seconds_streamed + uncredited_tail_seconds) * playback_speed
1918 )
1919 # propagate accurate duration to queue_item so UI displays it
1920 queue_item.duration = streamdetails.duration
1921 self.logger.debug(
1922 "Finished Streaming queue track: %s (%s) on queue %s "
1923 "- crossfade data prepared for next track: %s",
1924 streamdetails.uri,
1925 queue_item.name,
1926 queue.display_name,
1927 (
1928 next_queue_item.name
1929 if next_queue_item and queue_item.queue_id in self._crossfade_data
1930 else "N/A"
1931 ),
1932 )
1933
1934 async def get_queue_flow_stream(
1935 self,
1936 queue: PlayerQueue,
1937 start_queue_item: QueueItem,
1938 pcm_format: AudioFormat,
1939 session_id: str | None = None,
1940 protocol_player: Player | None = None,
1941 ) -> AsyncGenerator[bytes]:
1942 """
1943 Get a flow stream of all tracks in the queue as raw PCM audio.
1944
1945 yields chunks of exactly 1 second of audio in the given pcm_format.
1946
1947 :param queue: Queue being streamed.
1948 :param start_queue_item: First queue item in the flow stream.
1949 :param pcm_format: Shared PCM format for the complete flow stream.
1950 :param session_id: Queue session that owns processing-detail updates.
1951 :param protocol_player: The protocol player actually consuming the flow stream.
1952 Must be the same player that was used to select ``pcm_format`` so
1953 restart decisions are made against the correct supported sample rates
1954 and flow mode configuration. Falls back to the queue's player when omitted.
1955 """
1956 # ruff: noqa: PLR0915
1957 assert pcm_format.content_type.is_pcm()
1958 queue_track = None
1959 last_fadeout_part: bytes = b""
1960 last_streamdetails: StreamDetails | None = None
1961 last_play_log_entry: PlayLogEntry | None = None
1962 # Snapshot the queue's current session_id. PlayerQueues rotates this on
1963 # every new stream session, so if a newer producer takes over the queue
1964 # (rapid track switch, sync-group reform, dynamic leader handoff) the
1965 # snapshot will no longer match and we exit cleanly on the next yield or
1966 # playlog append â preventing two producers from writing to the same
1967 # pq_data.flow_mode_stream_log.
1968 pq_data = self.mass.player_queues.queue_data(queue.queue_id)
1969 flow_session_id = session_id or pq_data.session_id
1970 if flow_session_id is None or pq_data.session_id != flow_session_id:
1971 self.logger.debug(
1972 "Ignoring stale flow stream for queue %s (session %s, active %s)",
1973 queue.display_name,
1974 flow_session_id,
1975 pq_data.session_id,
1976 )
1977 return
1978 queue.flow_mode = True
1979 # A session can also be handed a second producer, which the session check does not
1980 # catch: players such as DLNA renderers sometimes open the same flow url twice to
1981 # probe the audio. Append to the list published here rather than to whatever the
1982 # queue currently holds, so the entries of a producer that has since been replaced
1983 # end up in a list nobody reads instead of interleaving with the live one's.
1984 flow_log: list[PlayLogEntry] = []
1985 pq_data.flow_mode_stream_log = flow_log
1986 if not start_queue_item:
1987 # this can happen in some (edge case) race conditions
1988 return
1989 pcm_sample_size = pcm_format.pcm_sample_size
1990 if start_queue_item.media_type != MediaType.TRACK:
1991 # no crossfade on non-tracks
1992 crossfade_mode = CrossfadeMode.DISABLED
1993 standard_crossfade_duration = 0
1994 else:
1995 crossfade_mode = self.mass.streams.get_crossfade_mode(queue)
1996 # crossfade duration is a global (queue controller) setting; fallback matches
1997 # CONF_ENTRY_CROSSFADE_DURATION's default
1998 standard_crossfade_duration = self.mass.config.get_raw_core_config_value(
1999 CONF_PLAYER_QUEUES, CONF_CROSSFADE_DURATION, 8
2000 )
2001 flow_mode_sample_rate_conf, flow_supported_sample_rates = self._flow_restart_context(
2002 queue.queue_id, protocol_player
2003 )
2004 # note: get_crossfade_mode() already falls back to standard when smart fades aren't
2005 # available (no analysis provider / minimal buffer), so crossfade_mode is safe to use.
2006 self.logger.info(
2007 "Start Queue Flow stream for Queue %s - crossfade: %s %s",
2008 queue.display_name,
2009 crossfade_mode,
2010 f"({standard_crossfade_duration}s)"
2011 if crossfade_mode == CrossfadeMode.STANDARD_CROSSFADE
2012 else "",
2013 )
2014 total_chunks_received = 0
2015
2016 def _superseded() -> bool:
2017 """Return True if a newer stream session has taken over this queue."""
2018 return pq_data.session_id != flow_session_id
2019
2020 queue_exhausted = False
2021 while True:
2022 # bail out early if a newer producer has taken over this queue,
2023 # so we don't append another entry to a stream log we no longer own
2024 if _superseded():
2025 self.logger.debug(
2026 "Flow stream for queue %s superseded (session %s -> %s) "
2027 "- exiting before next track",
2028 queue.display_name,
2029 flow_session_id,
2030 pq_data.session_id,
2031 )
2032 return
2033 # get (next) queue item to stream
2034 if queue_track is None:
2035 queue_track = start_queue_item
2036 else:
2037 try:
2038 queue_track = await self.mass.player_queues.load_next_queue_item(
2039 queue.queue_id, queue_track.queue_item_id
2040 )
2041 except QueueEmpty:
2042 queue_exhausted = True
2043 break
2044
2045 if self._flow_stream_needs_restart(
2046 queue_track,
2047 pcm_format,
2048 flow_supported_sample_rates,
2049 flow_mode_sample_rate_conf,
2050 is_first_track=queue_track is start_queue_item,
2051 ):
2052 break
2053
2054 if queue_track.streamdetails is None:
2055 self.logger.error(
2056 "No StreamDetails for queue item %s (%s) on queue %s - skipping track",
2057 queue_track.queue_item_id,
2058 queue_track.name,
2059 queue.display_name,
2060 )
2061 continue
2062 # a realtime source delivers at playback pace, so it has no audio to spare
2063 # for an overlap in either direction
2064 item_crossfade_mode = (
2065 CrossfadeMode.DISABLED if queue_track.streamdetails.is_realtime else crossfade_mode
2066 )
2067 if flow_session_id is not None:
2068 self.mass.streams.audio_processing.update_item_context(
2069 queue_id=queue.queue_id,
2070 session_id=flow_session_id,
2071 queue_item_id=queue_track.queue_item_id,
2072 queue_processing=AudioQueueProcessing(
2073 pcm_format=pcm_format,
2074 playback_speed=cast(
2075 "float",
2076 queue_track.extra_attributes.get("playback_speed", 1.0),
2077 ),
2078 crossfade_mode=item_crossfade_mode,
2079 overlay_active=overlay_active(queue),
2080 ),
2081 alters_audio=queue_track.streamdetails.fade_in,
2082 )
2083
2084 self.logger.debug(
2085 "Start Streaming queue track: %s (%s) for queue %s",
2086 queue_track.streamdetails.uri,
2087 queue_track.name,
2088 queue.display_name,
2089 )
2090 # last chance to bail before mutating the stream log: a newer producer
2091 # may have taken over while we were awaiting load_next_queue_item
2092 if _superseded():
2093 self.logger.debug(
2094 "Flow stream for queue %s superseded - exiting before playlog append",
2095 queue.display_name,
2096 )
2097 return
2098 track_playback_speed = cast(
2099 "float", queue_track.extra_attributes.get("playback_speed", 1.0)
2100 )
2101 # calculate crossfade buffer size
2102 crossfade_buffer_duration = (
2103 SMART_CROSSFADE_DURATION
2104 if item_crossfade_mode == CrossfadeMode.SMART_CROSSFADE
2105 else standard_crossfade_duration
2106 )
2107 crossfade_buffer_duration = min(
2108 crossfade_buffer_duration,
2109 int(queue_track.streamdetails.duration / 2)
2110 if queue_track.streamdetails.duration
2111 else crossfade_buffer_duration,
2112 )
2113 # skip crossfade if buffer would be too small to be meaningful
2114 if crossfade_buffer_duration < MIN_CROSSFADE_FALLBACK_DURATION:
2115 crossfade_buffer_duration = 0
2116 # Ensure crossfade buffer size is aligned to frame boundaries
2117 # Frame size = bytes_per_sample * channels
2118 bytes_per_sample = pcm_format.bit_depth // 8
2119 frame_size = bytes_per_sample * pcm_format.channels
2120 crossfade_buffer_size = int(pcm_format.pcm_sample_size * crossfade_buffer_duration)
2121 # Round down to nearest frame boundary
2122 crossfade_buffer_size = (crossfade_buffer_size // frame_size) * frame_size
2123 warmup_size = int(pcm_format.pcm_sample_size * WARMUP_DURATION)
2124
2125 # raw_seek_position feeds the PCM buffer; streamdetails.seek_position
2126 # (overwritten below) only drives reported elapsed time.
2127 raw_seek_position = queue_track.streamdetails.seek_position
2128 # Build eagerly so seek_position is set before PlayLogEntry is appended â
2129 # consumer-paced mix() would otherwise let the queue briefly report 0.
2130 crossfade_smart_fade: SmartFade | None = None
2131 collect_started = 0.0
2132 collect_resident = 0.0
2133 incoming_crossfade_size = crossfade_buffer_size
2134 incoming_audio_buffer: AudioBuffer | None = None
2135 if last_fadeout_part and last_streamdetails:
2136 transition_mode = CrossfadeMode.DISABLED
2137 incoming_duration = 0.0
2138 if crossfade_buffer_size > 0 and item_crossfade_mode != CrossfadeMode.DISABLED:
2139 transition_mode, incoming_duration = self._select_buffered_crossfade(
2140 queue_track.streamdetails,
2141 item_crossfade_mode,
2142 standard_crossfade_duration,
2143 track_playback_speed,
2144 )
2145 if transition_mode == CrossfadeMode.DISABLED:
2146 # nothing to fade into: flush the held-back tail of the previous track
2147 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2148 yield pcm_slice
2149 await asyncio.sleep(0)
2150 last_fadeout_part = b""
2151 last_streamdetails = None
2152 last_play_log_entry = None
2153 else:
2154 assert queue_track.streamdetails.buffer is not None
2155 incoming_audio_buffer = cast("AudioBuffer", queue_track.streamdetails.buffer)
2156 incoming_crossfade_size = int(pcm_format.pcm_sample_size * incoming_duration)
2157 incoming_crossfade_size = (incoming_crossfade_size // frame_size) * frame_size
2158 # no audio is emitted while the incoming overlap is collected, so a
2159 # slow collection here is heard as a gap at the transition
2160 collect_started = asyncio.get_event_loop().time()
2161 collect_resident = incoming_audio_buffer.duration_available
2162 crossfade_smart_fade = await self.smart_fades_mixer.build(
2163 fade_in_streamdetails=queue_track.streamdetails,
2164 fade_out_streamdetails=last_streamdetails,
2165 pcm_format=pcm_format,
2166 standard_crossfade_duration=standard_crossfade_duration,
2167 mode=transition_mode,
2168 fade_out_data=last_fadeout_part,
2169 fade_in_bytes_len=incoming_crossfade_size,
2170 )
2171 timing_info = crossfade_smart_fade.timing_info
2172 if isinstance(crossfade_smart_fade, StandardCrossFade):
2173 # A standard fade blends its overlap and passes everything after it
2174 # through untouched, so only the overlap has to be in hand before
2175 # the transition can start. Holding back the rest buys nothing and
2176 # keeps the player waiting - a smart fade does need its full window,
2177 # which is only chosen when the analysis it needs is already there.
2178 blended_seconds = (
2179 timing_info.fadein_trimmed_duration + timing_info.crossfade_duration
2180 )
2181 blended_size = int(pcm_format.pcm_sample_size * blended_seconds)
2182 incoming_crossfade_size = min(
2183 incoming_crossfade_size,
2184 (blended_size // frame_size) * frame_size,
2185 )
2186 queue_track.streamdetails.seek_position = (
2187 raw_seek_position
2188 + (timing_info.fadein_trimmed_duration + timing_info.crossfade_duration)
2189 * track_playback_speed
2190 )
2191 # append to play log so the queue controller can work out which track is playing
2192 play_log_entry = PlayLogEntry(queue_track.queue_item_id)
2193 flow_log.append(play_log_entry)
2194
2195 bytes_written = 0
2196 crossfade_buffer = bytearray()
2197 warmup_bytes = 0
2198 first_chunk_received = False
2199 holdback_armed = False
2200
2201 async for chunk in self.get_queue_item_stream(
2202 queue_track,
2203 pcm_format=pcm_format,
2204 seek_position=int(raw_seek_position),
2205 playback_speed=cast(
2206 "float", queue_track.extra_attributes.get("playback_speed", 1.0)
2207 ),
2208 raise_on_error=False,
2209 session_id=flow_session_id,
2210 prepared_buffer=incoming_audio_buffer,
2211 ):
2212 # if a newer producer has taken over this queue, stop sending
2213 # audio and exit cleanly before the outer-loop end-of-track
2214 # bookkeeping mutates seconds_streamed / duration on the log
2215 if _superseded():
2216 self.logger.debug(
2217 "Flow stream for queue %s superseded - stopping chunk yield",
2218 queue.display_name,
2219 )
2220 return
2221 total_chunks_received += 1
2222 if not first_chunk_received:
2223 first_chunk_received = True
2224 # inform the queue that the track is now loaded in the buffer
2225 # so the next track can be preloaded
2226 self.mass.player_queues.track_loaded_in_buffer(
2227 queue.queue_id, queue_track.queue_item_id
2228 )
2229
2230 if item_crossfade_mode == CrossfadeMode.DISABLED:
2231 # no cross/smart fade: yield chunks directly without intermediate buffer
2232 yield chunk
2233 bytes_written += len(chunk)
2234 del chunk
2235 continue
2236
2237 # Warmup: yield chunks directly until we have streamed WARMUP_DURATION
2238 # worth of audio, so playback starts immediately. Skip warmup when
2239 # crossfade data from the previous track is pending â we need a full
2240 # buffer for the mix.
2241 if warmup_bytes < warmup_size and not last_fadeout_part:
2242 yield chunk
2243 warmup_bytes += len(chunk)
2244 bytes_written += len(chunk)
2245 del chunk
2246 continue
2247
2248 if not last_fadeout_part and not holdback_armed:
2249 holdback_armed = self._crossfade_holdback_allowed(
2250 queue_track.streamdetails,
2251 crossfade_buffer_duration,
2252 track_playback_speed,
2253 )
2254 if not holdback_armed:
2255 # holding audio back now would only shrink the player's lead
2256 yield chunk
2257 bytes_written += len(chunk)
2258 del chunk
2259 continue
2260
2261 # smart fades enabled: accumulate chunks in crossfade buffer
2262 crossfade_buffer.extend(chunk)
2263 del chunk
2264 required_buffer_size = (
2265 incoming_crossfade_size if last_fadeout_part else crossfade_buffer_size
2266 )
2267 if len(crossfade_buffer) < required_buffer_size:
2268 await asyncio.sleep(0)
2269 continue
2270 # handle crossfade of previous track and new track
2271 if (
2272 last_fadeout_part
2273 and last_streamdetails
2274 and crossfade_smart_fade is not None
2275 and last_play_log_entry is not None
2276 ):
2277 self.logger.debug(
2278 "Collected %.1fs of incoming audio for the transition into %s in %.1fs"
2279 " (%.1fs was resident when it started)",
2280 len(crossfade_buffer) / pcm_sample_size,
2281 queue_track.name,
2282 asyncio.get_event_loop().time() - collect_started,
2283 collect_resident,
2284 )
2285 fadein_part = bytes(crossfade_buffer[:incoming_crossfade_size])
2286 remaining_bytes = bytes(crossfade_buffer[incoming_crossfade_size:])
2287 try:
2288 crossfade_bytes_written = 0
2289 async for mix_chunk in self.smart_fades_mixer.mix(
2290 crossfade_smart_fade,
2291 fade_in_part=fadein_part,
2292 fade_out_part=last_fadeout_part,
2293 pcm_format=pcm_format,
2294 ):
2295 yield mix_chunk
2296 crossfade_bytes_written += len(mix_chunk)
2297 except Exception as mix_err:
2298 if crossfade_bytes_written:
2299 # partial mix already played â concat'd tail would duplicate audio
2300 raise
2301 self.logger.warning(
2302 "Crossfade mixer failed for %s, falling back to simple concat: %s",
2303 queue_track.name,
2304 mix_err,
2305 )
2306 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2307 yield pcm_slice
2308 await asyncio.sleep(0)
2309 # full tail was pre-counted and is now yielded as-is
2310 crossfade_bytes_written = 0
2311 remaining_bytes = bytes(crossfade_buffer)
2312 # mix failed â undo the eager seek_position
2313 queue_track.streamdetails.seek_position = raw_seek_position
2314 if crossfade_bytes_written:
2315 # Split mix output at end-of-overlap: PRE+CF to A, POST to B.
2316 fadeout_share_seconds = (
2317 timing_info.pre_crossfade_duration + timing_info.crossfade_duration
2318 )
2319 fadeout_share = int(fadeout_share_seconds * pcm_sample_size)
2320 fadeout_share = (fadeout_share // frame_size) * frame_size
2321 fadeout_share = min(fadeout_share, crossfade_bytes_written)
2322 fadein_share = crossfade_bytes_written - fadeout_share
2323 bytes_written += fadein_share
2324 if last_play_log_entry:
2325 assert last_play_log_entry.seconds_streamed is not None
2326 # correct pre-counted full tail to the timing-based share
2327 last_play_log_entry.seconds_streamed += (
2328 fadeout_share - len(last_fadeout_part)
2329 ) / pcm_sample_size
2330 if remaining_bytes:
2331 for pcm_slice in iter_pcm_slices(remaining_bytes, pcm_format, 1000):
2332 yield pcm_slice
2333 await asyncio.sleep(0)
2334 bytes_written += len(remaining_bytes)
2335 del remaining_bytes
2336 last_fadeout_part = b""
2337 last_streamdetails = None
2338 crossfade_buffer = bytearray()
2339 warmup_bytes = 0
2340
2341 # yield everything above the crossfade buffer size
2342 while len(crossfade_buffer) > crossfade_buffer_size:
2343 yield bytes(crossfade_buffer[:pcm_sample_size])
2344 bytes_written += pcm_sample_size
2345 del crossfade_buffer[:pcm_sample_size]
2346 await asyncio.sleep(0)
2347
2348 # A source error after partial audio must not look like a completed item.
2349 # Progress reporting skips items with stream_error, so the item is not
2350 # marked played; move on to the next queue item like the zero-audio path.
2351 if first_chunk_received and queue_track.streamdetails.stream_error:
2352 if _superseded():
2353 return
2354 self.logger.warning(
2355 "Track %s (%s) on queue %s aborted by a stream error - skipping",
2356 queue_track.name,
2357 queue_track.streamdetails.uri,
2358 queue.display_name,
2359 )
2360 # the audio sent so far will still play out; keep the play log entry
2361 # honest about how much of this item was actually streamed
2362 play_log_entry.seconds_streamed = bytes_written / pcm_sample_size
2363 if last_fadeout_part:
2364 # crossfade into this item never happened â undo the eager seek_position
2365 queue_track.streamdetails.seek_position = raw_seek_position
2366 continue
2367
2368 #### HANDLE END OF TRACK
2369 if not first_chunk_received:
2370 self.logger.warning(
2371 "Track %s (%s) on queue %s produced no audio data - skipping",
2372 queue_track.name,
2373 queue_track.streamdetails.uri if queue_track.streamdetails else "unknown",
2374 queue.display_name,
2375 )
2376 queue_track.streamdetails.stream_error = True
2377 play_log_entry.seconds_streamed = 0
2378 if last_fadeout_part:
2379 queue_track.streamdetails.seek_position = raw_seek_position
2380 continue
2381 if last_fadeout_part:
2382 # edge case: we did not get enough data to make the crossfade
2383 # attribute these bytes to the previous track (they are its tail)
2384 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2385 yield pcm_slice
2386 await asyncio.sleep(0)
2387 # no crossfade happened â undo the eager seek_position
2388 queue_track.streamdetails.seek_position = raw_seek_position
2389 # full tail was pre-counted and is now yielded as-is
2390 last_fadeout_part = b""
2391 # a fade needs enough of the outgoing track to overlap with; a holdback that
2392 # armed late (or not at all) leaves less than that
2393 min_fade_out_size = int(pcm_sample_size * MIN_CROSSFADE_FALLBACK_DURATION)
2394 if len(crossfade_buffer) >= min_fade_out_size and self.crossfade_allowed(
2395 queue_track,
2396 crossfade_mode=item_crossfade_mode,
2397 player_id=queue.queue_id,
2398 flow_mode=True,
2399 ):
2400 last_fadeout_part = bytes(crossfade_buffer[-crossfade_buffer_size:])
2401 last_streamdetails = queue_track.streamdetails
2402 last_play_log_entry = play_log_entry
2403 remaining_bytes = bytes(crossfade_buffer[:-crossfade_buffer_size])
2404 if remaining_bytes:
2405 for pcm_slice in iter_pcm_slices(remaining_bytes, pcm_format, 1000):
2406 yield pcm_slice
2407 await asyncio.sleep(0)
2408 bytes_written += len(remaining_bytes)
2409 del remaining_bytes
2410 elif item_crossfade_mode != CrossfadeMode.DISABLED and crossfade_buffer:
2411 bytes_written += len(crossfade_buffer)
2412 for pcm_slice in iter_pcm_slices(bytes(crossfade_buffer), pcm_format, 1000):
2413 yield pcm_slice
2414 await asyncio.sleep(0)
2415 crossfade_buffer = bytearray()
2416
2417 # update duration details based on the actual pcm data we sent
2418 # this also accounts for crossfade and silence stripping
2419 seconds_streamed = bytes_written / pcm_sample_size
2420 queue_track.streamdetails.seconds_streamed = seconds_streamed
2421 play_log_entry.seconds_streamed = seconds_streamed
2422 # an externally aborted source ends in a clean EOF mid-track, so the
2423 # streamed length must not be written back as the item's duration
2424 source_buffer = queue_track.streamdetails.buffer
2425 source_aborted = source_buffer is not None and source_buffer.cancelled
2426 if not source_aborted:
2427 # the held-back crossfade tail still counts as this track's media-time
2428 tail_seconds = len(last_fadeout_part) / pcm_sample_size
2429 # streamdetails.duration is in media-time; seconds_streamed is stream-time
2430 # (post-atempo), so we scale by the track's playback_speed to recover media-time.
2431 queue_track.streamdetails.duration = int(
2432 queue_track.streamdetails.seek_position
2433 + (seconds_streamed + tail_seconds) * track_playback_speed
2434 )
2435 # propagate accurate duration to queue_item so UI displays it
2436 queue_track.duration = queue_track.streamdetails.duration
2437 play_log_entry.duration = queue_track.streamdetails.duration
2438 if last_play_log_entry is play_log_entry and last_fadeout_part:
2439 # Pre-count the full crossfade tail so the queue index calculation
2440 # doesn't undercount while waiting for the next track's crossfade mix.
2441 # This will be corrected to crossfade_total/2 once the mix completes.
2442 assert play_log_entry.seconds_streamed is not None
2443 play_log_entry.seconds_streamed += len(last_fadeout_part) / pcm_sample_size
2444 self.logger.debug(
2445 "Finished Streaming queue track: %s (%s) on queue %s",
2446 queue_track.streamdetails.uri,
2447 queue_track.name,
2448 queue.display_name,
2449 )
2450 #### HANDLE END OF QUEUE FLOW STREAM
2451 # skip end-of-queue bookkeeping if a newer producer has superseded us;
2452 # the new producer owns queue_buffer_completed and the play log now
2453 if _superseded():
2454 self.logger.debug(
2455 "Flow stream for queue %s superseded - skipping end-of-queue handling",
2456 queue.display_name,
2457 )
2458 return
2459 # end of queue flow: make sure we yield the last_fadeout_part
2460 if last_fadeout_part:
2461 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2462 yield pcm_slice
2463 await asyncio.sleep(0)
2464 # correct seconds streamed - the duration already includes the tail
2465 last_part_seconds = len(last_fadeout_part) / pcm_sample_size
2466 streamdetails = queue_track.streamdetails
2467 assert streamdetails is not None
2468 streamdetails.seconds_streamed = (
2469 streamdetails.seconds_streamed or 0
2470 ) + last_part_seconds
2471 # also update the play log entry so elapsed time tracking stays in sync
2472 if last_play_log_entry:
2473 assert last_play_log_entry.seconds_streamed is not None
2474 # full tail was pre-counted and is now yielded as-is
2475 last_play_log_entry.duration = streamdetails.duration
2476 last_fadeout_part = b""
2477 self.logger.info("Finished Queue Flow stream for Queue %s", queue.display_name)
2478 # only signal completion if we are still the active producer â a later
2479 # producer would (incorrectly) see this as its own completion otherwise
2480 if not _superseded():
2481 # inform the queue controller that all audio data has been generated
2482 # so it can handle the case where new items were added after the flow stream ended
2483 self.mass.player_queues.queue_buffer_completed(queue.queue_id, queue_exhausted)
2484
2485 async def get_overlay_mixed_stream(
2486 self,
2487 queue: PlayerQueue,
2488 audio_input: AsyncGenerator[bytes],
2489 pcm_format: AudioFormat,
2490 ) -> AsyncGenerator[bytes]:
2491 """
2492 Mix the queue's audio overlay (looping sound effect) into the given PCM stream.
2493
2494 The mixed output has the exact same PCM format, duration and chunking as the
2495 input stream. If the overlay source can not be resolved, the original stream
2496 is passed through unchanged so playback is never interrupted.
2497
2498 :param queue: The PlayerQueue holding the overlay source and volume.
2499 :param audio_input: The audio stream (raw PCM in ``pcm_format``) to mix into.
2500 :param pcm_format: PCM format of both the input and the mixed output.
2501 """
2502 overlay_input = await self._resolve_overlay_input(queue)
2503 if overlay_input is None:
2504 # overlay source unavailable: degrade gracefully to music-only
2505 async for chunk in audio_input:
2506 yield chunk
2507 return
2508 async for chunk in get_ffmpeg_overlay_stream(
2509 audio_input=audio_input,
2510 overlay_input=overlay_input,
2511 pcm_format=pcm_format,
2512 overlay_volume=queue.overlay_volume,
2513 chunk_size=pcm_format.pcm_sample_size,
2514 ):
2515 yield chunk
2516
2517 def crossfade_allowed(
2518 self,
2519 queue_item: QueueItem,
2520 crossfade_mode: CrossfadeMode,
2521 player_id: str,
2522 flow_mode: bool = False,
2523 next_queue_item: QueueItem | None = None,
2524 sample_rate: int | None = None,
2525 next_sample_rate: int | None = None,
2526 ) -> bool:
2527 """Get the crossfade config for a queue item."""
2528 if crossfade_mode == CrossfadeMode.DISABLED:
2529 return False
2530 if not (self.mass.player_queues.get(queue_item.queue_id)):
2531 return False # just a guard
2532 if not (self.mass.players.get_player(player_id)):
2533 return False # just a guard
2534 if queue_item.media_type != MediaType.TRACK:
2535 self.logger.debug("Skipping crossfade: current item is not a track")
2536 return False
2537 # check if the next item is part of the same album
2538 next_item = next_queue_item or self.mass.player_queues.get_next_item(
2539 queue_item.queue_id, queue_item.queue_item_id
2540 )
2541 if not next_item:
2542 # there is no next item!
2543 return False
2544 # check if next item is a track
2545 if next_item.media_type != MediaType.TRACK:
2546 self.logger.debug("Skipping crossfade: next item is not a track")
2547 return False
2548 if (
2549 isinstance(queue_item.media_item, Track)
2550 and isinstance(next_item.media_item, Track)
2551 and queue_item.media_item.album
2552 and next_item.media_item.album
2553 and queue_item.media_item.album == next_item.media_item.album
2554 and not self.mass.config.get_raw_core_config_value(
2555 "streams", CONF_ALLOW_CROSSFADE_SAME_ALBUM, False
2556 )
2557 ):
2558 # in general, crossfade is not desired for tracks of the same (gapless) album
2559 # because we have no accurate way to determine if the album is gapless or not,
2560 # for now we just never crossfade between tracks of the same album
2561 self.logger.debug("Skipping crossfade: next item is part of the same album")
2562 return False
2563
2564 # check if we're allowed to crossfade on different sample rates
2565 if (
2566 not flow_mode
2567 and sample_rate
2568 and next_sample_rate
2569 and sample_rate != next_sample_rate
2570 and not self.mass.config.get_raw_player_config_value(
2571 player_id,
2572 CONF_ENTRY_CROSSFADE_DIFFERENT_SAMPLE_RATES.key,
2573 CONF_ENTRY_CROSSFADE_DIFFERENT_SAMPLE_RATES.default_value,
2574 )
2575 ):
2576 self.logger.debug(
2577 "Skipping crossfade: player(protocol) does not support gapless playback "
2578 "with different sample rates (%s vs %s)",
2579 sample_rate,
2580 next_sample_rate,
2581 )
2582 return False
2583
2584 return True
2585
2586 def clear_crossfade_data(self, queue_id: str) -> None:
2587 """
2588 Clear any pending crossfade data for a queue.
2589
2590 :param queue_id: The queue ID to clear crossfade data for.
2591 """
2592 if queue_id in self._crossfade_data:
2593 self.logger.debug("Clearing crossfade data for queue %s", queue_id)
2594 del self._crossfade_data[queue_id]
2595
2596 async def get_shoutcast_stream(
2597 self, url: str, streamdetails: StreamDetails
2598 ) -> AsyncGenerator[bytes]:
2599 """
2600 Yield audio from a legacy Shoutcast server, with ICY metadata parsed inline.
2601
2602 :param url: Shoutcast stream URL.
2603 :param streamdetails: StreamDetails to update with ICY metadata as it arrives.
2604 """
2605 self.logger.debug("Start streaming from legacy Shoutcast server: %s", url)
2606
2607 parsed = urlparse(url)
2608 host = parsed.hostname
2609 port = parsed.port or 80
2610 path = parsed.path or "/"
2611 if parsed.query:
2612 path = f"{path}?{parsed.query}"
2613
2614 try:
2615 # Open raw socket connection
2616 reader, writer = await asyncio.wait_for(asyncio.open_connection(host, port), timeout=30)
2617 except TimeoutError as err:
2618 raise AudioError(f"Timeout connecting to Shoutcast stream {url}") from err
2619 except (OSError, ConnectionError) as err:
2620 raise AudioError(f"Failed to connect to Shoutcast stream {url}") from err
2621
2622 try:
2623 # Send HTTP request with ICY metadata header
2624 request = (
2625 f"GET {path} HTTP/1.1\r\n"
2626 f"Host: {host}\r\n"
2627 f"User-Agent: {HTTP_HEADERS['User-Agent']}\r\n"
2628 f"Icy-MetaData: 1\r\n\r\n"
2629 )
2630 writer.write(request.encode())
2631 await writer.drain()
2632
2633 # Read and parse response line
2634 try:
2635 response_line = await asyncio.wait_for(reader.readline(), timeout=10)
2636 except TimeoutError as err:
2637 raise AudioError("Timeout reading Shoutcast response") from err
2638
2639 if not response_line.startswith(b"ICY"):
2640 raise InvalidDataError("Invalid Shoutcast response")
2641
2642 # Read headers until empty line
2643 headers: dict[str, str] = {}
2644 while True:
2645 try:
2646 line = await asyncio.wait_for(reader.readline(), timeout=5)
2647 except TimeoutError as err:
2648 raise AudioError("Timeout reading Shoutcast headers") from err
2649
2650 if line in (b"\r\n", b"\n", b""):
2651 break
2652
2653 if b":" in line:
2654 try:
2655 key, value = line.decode("latin-1", errors="ignore").split(":", 1)
2656 headers[key.strip().lower()] = value.strip()
2657 except UnicodeDecodeError, ValueError:
2658 continue
2659
2660 # Get metadata interval
2661 meta_int_str = headers.get("icy-metaint")
2662 if not meta_int_str:
2663 raise InvalidDataError("No icy-metaint header in Shoutcast response")
2664
2665 try:
2666 meta_int = int(meta_int_str)
2667 except ValueError as err:
2668 raise InvalidDataError("Invalid icy-metaint value") from err
2669
2670 self.logger.debug("Connected to Shoutcast stream %s (icy-metaint: %s)", url, meta_int)
2671
2672 # Stream audio data with metadata parsing
2673 while True:
2674 try:
2675 # Read audio chunk
2676 audio_chunk = await reader.readexactly(meta_int)
2677 yield audio_chunk
2678
2679 # Read metadata length
2680 meta_byte = await reader.readexactly(1)
2681 if meta_byte == b"\x00":
2682 continue
2683
2684 meta_length = ord(meta_byte) * 16
2685 meta_data = await reader.readexactly(meta_length)
2686 self._parse_icy_metadata(meta_data, streamdetails)
2687
2688 except asyncio.exceptions.IncompleteReadError:
2689 # End of stream
2690 break
2691
2692 finally:
2693 writer.close()
2694 await writer.wait_closed()
2695
2696 # --- Private methods ---
2697
2698 def _notify_provider_streamed(
2699 self, streamdetails: StreamDetails, finished: bool, seconds_streamed: float
2700 ) -> None:
2701 """Report a (mostly) streamed item back to the provider that owns it."""
2702 if not finished and seconds_streamed < 90:
2703 return
2704 provider = self.mass.get_provider(streamdetails.provider)
2705 # plugin providers serve playable items too, but on_streamed is MusicProvider-only
2706 if provider is None or provider.type != ProviderType.MUSIC:
2707 return
2708 music_prov = cast("MusicProvider", provider)
2709 self.mass.create_task(music_prov.on_streamed(streamdetails))
2710
2711 def _get_volume_normalization_preference(
2712 self, streamdetails: StreamDetails
2713 ) -> VolumeNormalizationMode:
2714 """Return the configured normalization preference for the stream's media type."""
2715 conf_key = (
2716 CONF_VOLUME_NORMALIZATION_RADIO
2717 if streamdetails.media_type == MediaType.RADIO
2718 else CONF_VOLUME_NORMALIZATION_TRACKS
2719 )
2720 return VolumeNormalizationMode(
2721 self.mass.streams.get_config_value(conf_key, return_type=str)
2722 )
2723
2724 def _update_radio_stream_metadata(
2725 self,
2726 streamdetails: StreamDetails,
2727 artist: str | None,
2728 title: str,
2729 image_url: str | None = None,
2730 album: str | None = None,
2731 ) -> None:
2732 """
2733 Update radio stream metadata and trigger artwork lookup.
2734
2735 :param streamdetails: The stream details to update.
2736 :param artist: Artist name (will be normalized).
2737 :param title: Track title (will be cleaned for display).
2738 :param image_url: Optional image URL from stream metadata.
2739 :param album: Optional album name.
2740 """
2741 station_image_url = image_url or self.mass.metadata.get_radio_stream_station_image(
2742 streamdetails
2743 )
2744 artist_normalized = (
2745 self.mass.metadata.normalize_radio_artist_name(artist) if artist else None
2746 )
2747 display_title, _ = parse_title_and_version(title, strip_for_display=True)
2748
2749 streamdetails.stream_metadata = StreamMetadata(
2750 title=display_title,
2751 artist=artist_normalized,
2752 album=album,
2753 image_url=station_image_url,
2754 )
2755 streamdetails.stream_metadata_last_updated = time.time()
2756 if streamdetails.queue_id:
2757 self.mass.player_queues.signal_update(streamdetails.queue_id)
2758
2759 # Fetch artwork in background (track, album then artist)
2760 if artist and title and not image_url:
2761 self.mass.call_later(
2762 0.2,
2763 self.mass.metadata.update_radio_stream_artwork,
2764 streamdetails,
2765 task_id=f"update_radio_artwork_{streamdetails.queue_id}",
2766 )
2767
2768 async def _cache_radio_result(
2769 self,
2770 url: str,
2771 stream_type: StreamType,
2772 resolved_url: str | None = None,
2773 ) -> tuple[str, StreamType]:
2774 """Cache and return a radio stream resolution result."""
2775 result = (resolved_url or url, stream_type)
2776 await self.mass.cache.set(
2777 url,
2778 result,
2779 expiration=3600 * 3,
2780 provider=CACHE_PROVIDER,
2781 category=CACHE_CATEGORY_RESOLVED_RADIO_URL,
2782 )
2783 return result
2784
2785 async def _handle_client_error_for_radio_stream(
2786 self, url: str, err: aiohttp.ClientError, fallback_stream_type: StreamType
2787 ) -> tuple[str, StreamType]:
2788 """Handle aiohttp client errors during radio stream resolution."""
2789 # Prefer the final post-redirect URL: aiohttp follows redirects before raising,
2790 # but the original url may just point at a redirector rather than the ICY endpoint.
2791 request_info = getattr(err, "request_info", None)
2792 validate_url = str(request_info.url) if request_info is not None else url
2793
2794 # Check if this is a Shoutcast/ICY response that aiohttp can't parse
2795 if isinstance(err, aiohttp.ClientResponseError) and "ICY" in str(err).upper():
2796 self.logger.debug(
2797 "ICY response detected for %s, validating Shoutcast stream", validate_url
2798 )
2799 if await self._validate_shoutcast_stream(validate_url):
2800 return await self._cache_radio_result(
2801 url, StreamType.SHOUTCAST, resolved_url=validate_url
2802 )
2803 self.logger.warning(
2804 "ICY response detected but Shoutcast validation failed for %s", validate_url
2805 )
2806 return await self._cache_radio_result(
2807 url, fallback_stream_type, resolved_url=validate_url
2808 )
2809
2810 # Other aiohttp errors - might still be Shoutcast, check it
2811 self.logger.debug("aiohttp error for %s, checking if legacy Shoutcast stream", validate_url)
2812 if await self._validate_shoutcast_stream(validate_url):
2813 return await self._cache_radio_result(
2814 url, StreamType.SHOUTCAST, resolved_url=validate_url
2815 )
2816
2817 # Unknown error - still try to stream
2818 self.logger.warning(
2819 "Failed to parse radio URL %s: %s - attempting direct stream", validate_url, str(err)
2820 )
2821 return await self._cache_radio_result(url, fallback_stream_type, resolved_url=validate_url)
2822
2823 async def _get_audio_buffer(
2824 self,
2825 queue_item: QueueItem,
2826 seek_position_ms: int,
2827 reason: str,
2828 capacity_wait_timeout: float,
2829 allow_provider_match: bool,
2830 ) -> AudioBuffer:
2831 """
2832 Create or reuse a ready AudioBuffer within one queue-item preparation lock.
2833
2834 :param queue_item: Queue item whose source should be buffered.
2835 :param seek_position_ms: Position in milliseconds to start from.
2836 :param reason: Caller context for logging.
2837 :param capacity_wait_timeout: Total seconds to spend waiting for source capacity.
2838 :param allow_provider_match: Whether an on-demand cross-provider match may widen
2839 the candidates when all are saturated.
2840 """
2841 loop = asyncio.get_running_loop()
2842 # the playback intent lives on the details we start from; keep it across a reselection
2843 initial_streamdetails = queue_item.streamdetails
2844 seek_position = (
2845 int(initial_streamdetails.seek_position)
2846 if initial_streamdetails
2847 else seek_position_ms // 1000
2848 )
2849 fade_in = bool(initial_streamdetails and initial_streamdetails.fade_in)
2850 prefer_album_loudness = bool(
2851 initial_streamdetails and initial_streamdetails.prefer_album_loudness
2852 )
2853 all_candidate_instances = {
2854 provider.instance_id
2855 for mapping in (
2856 queue_item.media_item.provider_mappings if queue_item.media_item else ()
2857 )
2858 if mapping.available
2859 for provider in self._get_mapping_providers(mapping)
2860 }
2861 if initial_streamdetails is not None:
2862 all_candidate_instances.add(initial_streamdetails.provider)
2863 # a track may also exist on streaming providers it has no mapping for yet; such a
2864 # match is only searched once, and only when every known candidate is saturated
2865 match_pending = (
2866 allow_provider_match
2867 and isinstance(queue_item.media_item, Track)
2868 and self._has_alternative_match_providers(queue_item.media_item)
2869 )
2870
2871 deadline = loop.time() + capacity_wait_timeout
2872 busy_instances: set[str] = set()
2873 final_pass = False
2874 last_capacity_error: ProviderStreamLimitError | None = None
2875 last_failed_streamdetails: StreamDetails | None = None
2876 while True:
2877 if queue_item.streamdetails is None or (
2878 queue_item.streamdetails.provider in busy_instances and not final_pass
2879 ):
2880 try:
2881 queue_item.streamdetails = await self.get_stream_details(
2882 queue_item,
2883 seek_position=seek_position,
2884 fade_in=fade_in,
2885 prefer_album_loudness=prefer_album_loudness,
2886 excluded_provider_instances=busy_instances,
2887 )
2888 except (AudioError, MediaNotFoundError) as err:
2889 if last_capacity_error is None:
2890 raise
2891 if final_pass:
2892 # capacity was the root cause, surface the typed (actionable) error
2893 raise last_capacity_error from err
2894 # no usable alternative mapping: restore the capacity-blocked details
2895 # and spend the remaining budget blocking on that provider's slot
2896 final_pass = True
2897 continue
2898 finally:
2899 if queue_item.streamdetails is None:
2900 # never leave the queue item without streamdetails on any exit,
2901 # including a cancellation or a non-audio provider failure
2902 queue_item.streamdetails = last_failed_streamdetails
2903 streamdetails = queue_item.streamdetails
2904 assert streamdetails is not None # for type checking
2905 remaining = max(deadline - loop.time(), 0)
2906 alternatives_left = bool(
2907 all_candidate_instances - busy_instances - {streamdetails.provider}
2908 )
2909 # probe (0s) whenever a reselection can still follow: a free slot is still
2910 # acquired instantly, while a busy one fails fast instead of spending the
2911 # whole budget on this candidate. Block only on the last resort.
2912 source_wait = (
2913 0.0
2914 if (not final_pass and (alternatives_left or busy_instances or match_pending))
2915 else remaining
2916 )
2917 try:
2918 return await AudioBuffer.get_buffer(
2919 mass=self.mass,
2920 streamdetails=streamdetails,
2921 seek_position_ms=seek_position_ms,
2922 wait_ready=True,
2923 reason=reason,
2924 source_wait_timeout=source_wait,
2925 )
2926 except ProviderStreamLimitError as err:
2927 last_capacity_error = err
2928 last_failed_streamdetails = streamdetails
2929 busy_instances.add(err.provider_instance)
2930 if final_pass or loop.time() >= deadline:
2931 raise
2932 if all_candidate_instances.issubset(busy_instances):
2933 discovered: set[str] = set()
2934 if match_pending:
2935 match_pending = False
2936 try:
2937 discovered = await self._discover_alternative_provider_mappings(
2938 queue_item, busy_instances, max(deadline - loop.time(), 0)
2939 )
2940 except Exception as err:
2941 # discovery is best-effort: any failure falls back to the
2942 # final blocking wait instead of replacing the typed error
2943 self.logger.warning(
2944 "Alternative provider search for %s failed: %s",
2945 queue_item.name,
2946 err,
2947 )
2948 if discovered:
2949 all_candidate_instances.update(discovered)
2950 else:
2951 # every candidate is saturated: one last blocking wait on the best one
2952 busy_instances.clear()
2953 final_pass = True
2954 queue_item.streamdetails = None
2955 except AudioError:
2956 if last_capacity_error is None or final_pass:
2957 raise
2958 # a broken alternate must not turn a transient capacity miss into a hard
2959 # failure: restore the blocked details and spend the rest of the budget there
2960 queue_item.streamdetails = last_failed_streamdetails
2961 final_pass = True
2962
2963 def _get_streamdetail_candidates(
2964 self,
2965 provider_mappings: Iterable[ProviderMapping],
2966 preferred_providers: list[str],
2967 excluded_provider_instances: set[str],
2968 ) -> list[tuple[ProviderMapping, Provider]]:
2969 """
2970 Return mapping candidates in steering, quality, and instance-fallback order.
2971
2972 :param provider_mappings: Mappings attached to the media item.
2973 :param preferred_providers: Provider instances tried before widening to the rest.
2974 :param excluded_provider_instances: Provider instances unavailable to this attempt.
2975 :return: Ordered provider mapping candidates.
2976 """
2977 ordered_mappings = sorted(
2978 provider_mappings, key=lambda mapping: mapping.quality or 0, reverse=True
2979 )
2980 preferred_candidates: list[tuple[ProviderMapping, Provider]] = []
2981 fallback_candidates: list[tuple[ProviderMapping, Provider]] = []
2982 seen_candidates: set[tuple[str, str]] = set()
2983 for mapping in ordered_mappings:
2984 if not mapping.available:
2985 self.logger.debug("Skipping unavailable %s", mapping)
2986 continue
2987 for provider in self._get_mapping_providers(mapping):
2988 candidate_id = (provider.instance_id, mapping.item_id)
2989 if (
2990 candidate_id in seen_candidates
2991 or provider.instance_id in excluded_provider_instances
2992 ):
2993 continue
2994 seen_candidates.add(candidate_id)
2995 candidate = (mapping, provider)
2996 if provider.instance_id in preferred_providers:
2997 preferred_candidates.append(candidate)
2998 else:
2999 fallback_candidates.append(candidate)
3000 return [*preferred_candidates, *fallback_candidates]
3001
3002 def _get_mapping_providers(self, mapping: ProviderMapping) -> list[Provider]:
3003 """
3004 Return the mapped provider followed by compatible instances of its streaming catalog.
3005
3006 :param mapping: Provider mapping whose item ID will be requested.
3007 :return: Loaded provider instances that can resolve the mapping.
3008 """
3009 providers: list[Provider] = []
3010 if (
3011 primary_provider := self.mass.get_provider(
3012 mapping.provider_instance, return_unavailable=True
3013 )
3014 ) and primary_provider.available:
3015 providers.append(primary_provider)
3016 # another account of the same streaming catalog serves the same item ID,
3017 # so it can stand in when the mapped instance can not
3018 for provider in self.mass.providers:
3019 if (
3020 not isinstance(provider, MusicProvider)
3021 or not provider.available
3022 or not provider.is_streaming_provider
3023 or provider.domain != mapping.provider_domain
3024 or provider in providers
3025 ):
3026 continue
3027 providers.append(provider)
3028 if not providers:
3029 self.logger.debug("Skipping %s - provider not available", mapping)
3030 return providers
3031
3032 def _is_match_candidate_provider(
3033 self, provider: MusicProvider, known_domains: set[str]
3034 ) -> bool:
3035 """
3036 Return whether a provider is eligible to search a track match on.
3037
3038 :param provider: Music provider to check.
3039 :param known_domains: Provider domains the track already has mappings for.
3040 """
3041 return (
3042 provider.available
3043 and provider.is_streaming_provider
3044 and ProviderFeature.SEARCH in provider.supported_features
3045 and provider.domain not in known_domains
3046 and MediaType.TRACK in provider.supported_media_types
3047 )
3048
3049 def _has_alternative_match_providers(self, media_item: Track) -> bool:
3050 """
3051 Return whether any configured streaming provider could carry an unmapped match.
3052
3053 :param media_item: Track whose existing mappings define the known provider domains.
3054 """
3055 known_domains = {mapping.provider_domain for mapping in media_item.provider_mappings}
3056 return any(
3057 self._is_match_candidate_provider(provider, known_domains)
3058 for provider in self.mass.music.providers
3059 )
3060
3061 async def _discover_alternative_provider_mappings(
3062 self, queue_item: QueueItem, busy_instances: set[str], remaining: float
3063 ) -> set[str]:
3064 """
3065 Search other streaming providers for the queue item's track and widen its mappings.
3066
3067 A found mapping is added to the media item (and persisted for library items) so the
3068 capacity reselection can continue on the discovered provider.
3069
3070 :param queue_item: Queue item whose track should be matched on another provider.
3071 :param busy_instances: Provider instances already known to be saturated.
3072 :param remaining: Seconds left of the caller's capacity budget.
3073 :return: Provider instances able to serve the discovered mappings.
3074 """
3075 media_item = queue_item.media_item
3076 if not isinstance(media_item, Track):
3077 return set()
3078 known_domains = {mapping.provider_domain for mapping in media_item.provider_mappings}
3079 eligible = [
3080 provider
3081 for provider in self.mass.music.providers
3082 if self._is_match_candidate_provider(provider, known_domains)
3083 and provider.instance_id not in busy_instances
3084 and provider.has_available_stream_slot
3085 ]
3086 if not eligible:
3087 return set()
3088 # mirror the playback user's provider steering for the search order
3089 if (
3090 (pq_data := self.mass.player_queues.queue_data_or_none(queue_item.queue_id))
3091 and pq_data.userid
3092 and (playback_user := await self.mass.webserver.auth.get_user(pq_data.userid))
3093 and playback_user.provider_filter
3094 ):
3095 preferred = set(playback_user.provider_filter)
3096 eligible.sort(key=lambda provider: provider.instance_id not in preferred)
3097 # one instance per domain: a found mapping widens to sibling instances anyway
3098 candidates: list[MusicProvider] = []
3099 for provider in eligible:
3100 if provider.domain in known_domains:
3101 continue
3102 known_domains.add(provider.domain)
3103 candidates.append(provider)
3104 # the track's own album is free, sufficient evidence for the strict compare and
3105 # avoids match_provider's multi-provider album lookup on every call
3106 ref_albums = [media_item.album] if isinstance(media_item.album, Album) else []
3107 matches: list[ProviderMapping] = []
3108 try:
3109 async with asyncio.timeout(min(STREAM_SLOT_MATCH_TIMEOUT, remaining)):
3110 for provider in candidates:
3111 # one failing provider must not end the search on the others
3112 try:
3113 matches = await self.mass.music.tracks.match_provider(
3114 media_item, provider, strict=True, ref_albums=ref_albums
3115 )
3116 except Exception as err:
3117 self.logger.debug("Searching a match on %s failed: %s", provider.name, err)
3118 continue
3119 if matches:
3120 break
3121 except TimeoutError:
3122 self.logger.debug("Searching an alternative provider for %s timed out", media_item.name)
3123 if not matches:
3124 return set()
3125 media_item.provider_mappings.update(matches)
3126 if media_item.provider == "library":
3127 # persist in the background so future plays have the mapping ahead of time;
3128 # cancellation of this playback must never interrupt the library write
3129 self.mass.create_task(
3130 self.mass.music.tracks.add_provider_mappings(media_item.item_id, matches)
3131 )
3132 self.logger.info(
3133 "All known sources for %s are at their stream limit, "
3134 "using a matching track found on %s",
3135 media_item.name,
3136 matches[0].provider_domain,
3137 )
3138 return {
3139 provider.instance_id
3140 for mapping in matches
3141 for provider in self._get_mapping_providers(mapping)
3142 }
3143
3144 async def _request_streamdetails(
3145 self,
3146 candidates: Iterable[tuple[ProviderMapping, Provider]],
3147 media_type: MediaType,
3148 ) -> StreamDetails | None:
3149 """
3150 Request stream details from ordered provider mapping candidates.
3151
3152 :param candidates: Candidates in mapping and compatible-instance order.
3153 :param media_type: Media type requested from each provider.
3154 :return: The first resolved stream details, or None when every candidate failed.
3155 :raises AudioError: The last (actionable) audio error when no candidate resolved.
3156 """
3157 last_audio_error: AudioError | None = None
3158 for mapping, provider in candidates:
3159 # music and plugin providers share this signature, so either type can own the item
3160 token = BYPASS_THROTTLER.set(True)
3161 try:
3162 stream_prov = cast("MusicProvider | PluginProvider", provider)
3163 return await stream_prov.get_stream_details(mapping.item_id, media_type)
3164 except AudioError as err:
3165 # remember the last one so its (actionable) message can be re-raised
3166 last_audio_error = err
3167 self.logger.warning("%s", err)
3168 except MusicAssistantError as err:
3169 self.logger.warning("%s", err)
3170 finally:
3171 BYPASS_THROTTLER.reset(token)
3172 if last_audio_error is not None:
3173 raise last_audio_error
3174 return None
3175
3176 async def _get_media_stream(
3177 self,
3178 streamdetails: StreamDetails,
3179 pcm_format: AudioFormat,
3180 seek_position: int,
3181 filter_params: list[str] | None,
3182 chunk_seconds: float,
3183 ) -> AsyncGenerator[bytes]:
3184 """
3185 Stream one provider source as raw PCM.
3186
3187 :param streamdetails: Details of the stream to fetch.
3188 :param pcm_format: Target PCM format the consumer expects.
3189 :param seek_position: Requested seek offset in seconds.
3190 :param filter_params: Optional ffmpeg filter expressions.
3191 :param chunk_seconds: Size of each yielded chunk in seconds of audio.
3192 """
3193 mass = self.mass
3194 logger = self.logger.getChild("media_stream")
3195 logger.log(VERBOSE_LOG_LEVEL, "Starting media stream for %s", streamdetails.uri)
3196 # copy: the args below are appended per call, while the StreamDetails is cached on
3197 # the queue item and reused across calls (retry, seek, background analysis)
3198 extra_input_args = list(streamdetails.extra_input_args or [])
3199 # the resolver below zeroes out seek_position where the seek is delegated to the
3200 # source itself, so keep the requested position for the duration writeback
3201 requested_seek_position = seek_position
3202
3203 # work out audio source for these streamdetails
3204 audio_source, seek_position, extra_input_args = await self._resolve_media_stream_source(
3205 streamdetails, seek_position, extra_input_args
3206 )
3207
3208 # pace ffmpeg at native rate for live sources; the producer (e.g.
3209 # librespot's pipe backend) may otherwise write faster than realtime.
3210 # The initial burst grants a small bounded read-ahead so downstream
3211 # jitter does not immediately underrun the player. Providers that need
3212 # different pacing can pass their own -re/-readrate args to override.
3213 if (
3214 streamdetails.media_type == MediaType.AUDIO_SOURCE
3215 and "-re" not in extra_input_args
3216 and "-readrate" not in extra_input_args
3217 ):
3218 extra_input_args += ["-readrate", "1", "-readrate_initial_burst", "0.5"]
3219
3220 # handle seek support
3221 if seek_position and streamdetails.duration and streamdetails.allow_seek:
3222 extra_input_args += ["-ss", str(int(seek_position))]
3223
3224 bytes_sent = 0
3225 finished = False
3226 cancelled = False
3227 first_chunk_received = False
3228 ffmpeg_loglevel = "debug" if self.logger.isEnabledFor(VERBOSE_LOG_LEVEL) else "info"
3229 # When a provider hands us already-decoded audio (e.g. Spotify Connect /
3230 # AirPlay receivers piping PCM after their own decode), audio_format is
3231 # the original source format meant for display while decoded_audio_format
3232 # is what ffmpeg actually needs to read off the wire.
3233 ffmpeg_input_format = streamdetails.decoded_audio_format or streamdetails.audio_format
3234 ffmpeg_proc = FFMpeg(
3235 audio_input=audio_source,
3236 input_format=ffmpeg_input_format,
3237 output_format=pcm_format,
3238 filter_params=filter_params,
3239 extra_input_args=extra_input_args,
3240 collect_log_history=True,
3241 loglevel=ffmpeg_loglevel,
3242 )
3243
3244 try:
3245 await ffmpeg_proc.start()
3246 assert ffmpeg_proc.proc is not None # for type checking
3247 if logger.isEnabledFor(VERBOSE_LOG_LEVEL):
3248 logger.log(
3249 VERBOSE_LOG_LEVEL,
3250 "Started media stream for %s - using streamtype: %s "
3251 "- pcm format: %s - ffmpeg PID: %s",
3252 streamdetails.uri,
3253 streamdetails.stream_type,
3254 pcm_format.content_type.value,
3255 ffmpeg_proc.proc.pid,
3256 )
3257 else:
3258 logger.debug(
3259 "Started media stream for %s - using streamtype: %s",
3260 streamdetails.uri,
3261 streamdetails.stream_type,
3262 )
3263 stream_start = mass.loop.time()
3264 chunk_size = calculate_content_length(pcm_format, chunk_seconds)
3265 chunk_iter = ffmpeg_proc.iter_chunked(chunk_size)
3266 while True:
3267 # Time the read, not the yield: catches a stalled source, ignores backpressure.
3268 read_timeout = (
3269 STREAM_START_TIMEOUT if not first_chunk_received else STREAM_STALL_TIMEOUT
3270 )
3271 try:
3272 async with asyncio.timeout(read_timeout):
3273 chunk = await anext(chunk_iter)
3274 except StopAsyncIteration:
3275 break
3276 except TimeoutError as err:
3277 raise AudioError(f"Source stalled: no audio for {read_timeout}s") from err
3278 if not first_chunk_received:
3279 # At this point ffmpeg has started and should now know the codec used
3280 # for encoding the audio.
3281 # Note: ffmpeg_proc.input_format is the same object as
3282 # ffmpeg_input_format, so sample_rate / bit_depth / bit_rate
3283 # parsed from the ffmpeg log already live on streamdetails too.
3284 first_chunk_received = True
3285 # Skip the codec_type writeback when the provider declared a
3286 # decoded format: audio_format already holds the authoritative
3287 # source codec and the probed value would just be the
3288 # post-decode wire format (e.g. PCM for Spotify Connect).
3289 if streamdetails.decoded_audio_format is None:
3290 streamdetails.audio_format.codec_type = ffmpeg_proc.input_format.codec_type
3291 # Some providers omit (or report 0 for) the item duration; ffmpeg can
3292 # usually probe it from the source. Only apply when missing so we
3293 # don't clobber an accurate provider value with a rounded one.
3294 if ffmpeg_proc.parsed_duration is not None and not streamdetails.duration:
3295 streamdetails.duration = ffmpeg_proc.parsed_duration
3296 logger.debug(
3297 "First chunk received after %.2f seconds (codec detected: %s)",
3298 mass.loop.time() - stream_start,
3299 ffmpeg_proc.input_format.codec_type,
3300 )
3301 yield chunk
3302 bytes_sent += len(chunk)
3303
3304 # end of audio/track reached
3305 logger.debug("End of media stream reached for %s", streamdetails.uri)
3306 # wait until stderr also completed reading
3307 await ffmpeg_proc.wait_with_timeout(5)
3308 logger.log(
3309 VERBOSE_LOG_LEVEL,
3310 "FFmpeg process ended with return code %s for %s",
3311 ffmpeg_proc.returncode,
3312 streamdetails.uri,
3313 )
3314 # a nested source raises through the stdin feeder, where ffmpeg's own exit
3315 # would otherwise flatten it into a generic AudioError
3316 if isinstance(ffmpeg_proc.stdin_feeder_exception, ProviderStreamLimitError):
3317 raise ffmpeg_proc.stdin_feeder_exception
3318 if ffmpeg_proc.returncode not in (0, None):
3319 log_trail = "\n".join(list(ffmpeg_proc.log_history)[-5:])
3320 raise AudioError(f"FFMpeg exited with code {ffmpeg_proc.returncode}: {log_trail}")
3321 if bytes_sent == 0:
3322 # edge case: no audio data was received at all
3323 raise AudioError("No audio was received")
3324 finished = True
3325 except (Exception, GeneratorExit, asyncio.CancelledError) as err:
3326 if isinstance(err, asyncio.CancelledError | GeneratorExit):
3327 # we were cancelled, just raise
3328 cancelled = True
3329 raise
3330 if isinstance(ffmpeg_proc.stdin_feeder_exception, ProviderStreamLimitError):
3331 raise ffmpeg_proc.stdin_feeder_exception
3332 # dump the last 10 lines of the log in case of an unclean exit
3333 logger.warning("\n".join(list(ffmpeg_proc.log_history)[-10:]))
3334 raise AudioError(f"Error while streaming: {err}") from err
3335 finally:
3336 # always ensure close is called which also handles all cleanup
3337 await ffmpeg_proc.close()
3338 # determine how many seconds we've received
3339 # for pcm output we can calculate this easily
3340 seconds_received = bytes_sent / pcm_format.pcm_sample_size if bytes_sent else 0
3341 # store accurate duration, but only for a playthrough from the very start:
3342 # a seeked stream yields the remaining audio, not the item's full length
3343 if finished and not requested_seek_position and seconds_received:
3344 streamdetails.duration = int(seconds_received)
3345
3346 logger.log(
3347 VERBOSE_LOG_LEVEL,
3348 "stream %s (with code %s) for %s",
3349 "cancelled" if cancelled else "finished" if finished else "aborted",
3350 ffmpeg_proc.returncode,
3351 streamdetails.uri,
3352 )
3353
3354 def _crossfade_holdback_allowed(
3355 self, streamdetails: StreamDetails, tail_seconds: float, playback_speed: float = 1.0
3356 ) -> bool:
3357 """
3358 Return whether the outgoing tail may be held back for a crossfade.
3359
3360 :param streamdetails: Stream details of the track being streamed.
3361 :param tail_seconds: Length of the tail to hold back, in seconds of playback.
3362 :param playback_speed: Playback-speed multiplier of the track.
3363 """
3364 if tail_seconds <= 0 or playback_speed <= 0 or streamdetails.is_realtime:
3365 return False
3366 audio_buffer = cast("AudioBuffer | None", streamdetails.buffer)
3367 if audio_buffer is None or audio_buffer.has_error:
3368 # a failed source is skipped without a fade, so its remaining audio is
3369 # better off played out than held back for one
3370 return False
3371 # While the source is still delivering, it is what limits playback: withholding
3372 # a tail on top of that eats into the lead the player needs. Once the source is
3373 # done the remaining audio is resident, so the tail comes for free. A buffer that
3374 # is too small to ever hold the tail is the exception - waiting for the source
3375 # there would only lose the fade.
3376 return audio_buffer.eof or audio_buffer.max_size_seconds / playback_speed < tail_seconds
3377
3378 def _select_buffered_crossfade(
3379 self,
3380 streamdetails: StreamDetails,
3381 crossfade_mode: CrossfadeMode,
3382 standard_crossfade_duration: int,
3383 playback_speed: float = 1.0,
3384 ) -> tuple[CrossfadeMode, float]:
3385 """
3386 Select a crossfade that can be completed from resident incoming PCM.
3387
3388 :param streamdetails: Incoming track stream details.
3389 :param crossfade_mode: Requested crossfade mode.
3390 :param standard_crossfade_duration: Configured standard overlap in seconds.
3391 :param playback_speed: Incoming track playback-speed multiplier.
3392 :return: Effective mode and resident fade-in duration in seconds.
3393 """
3394 audio_buffer = streamdetails.buffer
3395 if (
3396 crossfade_mode == CrossfadeMode.DISABLED
3397 or playback_speed <= 0
3398 # a realtime source has no more than its banked lead: fading in would spend
3399 # that at the start of the track and leave nothing for the rest of it
3400 or streamdetails.is_realtime
3401 or audio_buffer is None
3402 or audio_buffer.has_error
3403 or not audio_buffer.is_valid()
3404 ):
3405 return CrossfadeMode.DISABLED, 0
3406
3407 available_seconds = audio_buffer.duration_available / playback_speed
3408 if (
3409 crossfade_mode == CrossfadeMode.SMART_CROSSFADE
3410 and audio_buffer.ready.is_set()
3411 and available_seconds >= MIN_EFFECTIVE_FADE_BUFFER
3412 ):
3413 return crossfade_mode, min(SMART_CROSSFADE_DURATION, available_seconds)
3414 if (
3415 crossfade_mode == CrossfadeMode.STANDARD_CROSSFADE
3416 and standard_crossfade_duration >= MIN_CROSSFADE_FALLBACK_DURATION
3417 and audio_buffer.ready.is_set()
3418 and available_seconds >= standard_crossfade_duration
3419 ):
3420 return crossfade_mode, standard_crossfade_duration
3421
3422 fallback_duration = min(standard_crossfade_duration, available_seconds)
3423 if fallback_duration < MIN_CROSSFADE_FALLBACK_DURATION:
3424 return CrossfadeMode.DISABLED, 0
3425 self.logger.debug(
3426 "Using %s second standard crossfade for %s from resident audio",
3427 fallback_duration,
3428 streamdetails.uri,
3429 )
3430 return CrossfadeMode.STANDARD_CROSSFADE, fallback_duration
3431
3432 async def _resolve_media_stream_source(
3433 self,
3434 streamdetails: StreamDetails,
3435 seek_position: int,
3436 extra_input_args: list[str],
3437 ) -> tuple[str | AsyncGenerator[bytes], int, list[str]]:
3438 """
3439 Resolve the input consumed by ffmpeg for the given stream details.
3440
3441 :param streamdetails: Details of the stream to fetch.
3442 :param seek_position: Requested seek offset in seconds.
3443 :param extra_input_args: Provider-supplied ffmpeg input arguments.
3444 :return: The ffmpeg input, the remaining seek offset and the ffmpeg input arguments.
3445 """
3446 stream_type = streamdetails.stream_type
3447 if stream_type == StreamType.CUSTOM:
3448 # MusicProvider and PluginProvider both expose get_audio_stream with the same shape.
3449 # Pin the exact instance: a domain fallback would stream from a sibling account
3450 # while the source-stream slot is charged to the instance that issued the details.
3451 provider = self.mass.get_provider(streamdetails.provider, return_unavailable=True)
3452 if provider is None or not provider.available:
3453 raise ProviderUnavailableError(
3454 f"Provider {streamdetails.provider} for stream is no longer available"
3455 )
3456 provider = cast("MusicProvider | PluginProvider", provider)
3457 audio_source = provider.get_audio_stream(
3458 streamdetails, seek_position=seek_position if streamdetails.can_seek else 0
3459 )
3460 return audio_source, 0 if streamdetails.can_seek else seek_position, extra_input_args
3461 if stream_type == StreamType.ICY:
3462 assert streamdetails.path is not None
3463 assert isinstance(streamdetails.path, (str, list))
3464 audio_source = self.get_reconnecting_icy_radio_stream(streamdetails.path, streamdetails)
3465 return audio_source, 0, extra_input_args
3466 if stream_type == StreamType.SHOUTCAST:
3467 assert isinstance(streamdetails.path, str)
3468 return self.get_shoutcast_stream(streamdetails.path, streamdetails), 0, extra_input_args
3469 if stream_type == StreamType.IN_BAND:
3470 assert isinstance(streamdetails.path, str) # for type checking
3471
3472 # For IN_BAND (OGG/Opus) radio streams, use chained OGG handler.
3473 # This handles the chained OGG format by stitching logical bitstreams together
3474 # so FFmpeg sees a single continuous stream. Metadata is extracted in-band.
3475 audio_source = get_chained_ogg_stream(
3476 self.mass,
3477 streamdetails.path,
3478 metadata_callback=partial(self._handle_inband_metadata, streamdetails),
3479 )
3480 # seeking not possible on radio streams
3481 return audio_source, 0, extra_input_args
3482 if stream_type == StreamType.HLS:
3483 assert isinstance(streamdetails.path, str) # for type checking
3484 substream = await self.get_hls_substream(streamdetails.path)
3485 if streamdetails.media_type == MediaType.RADIO:
3486 # HLS streams (especially the BBC) struggle when they're played directly
3487 # with ffmpeg, where they just stop after some minutes,
3488 # so we tell ffmpeg to loop around in this case.
3489 extra_input_args += ["-stream_loop", "-1", "-re"]
3490 return substream.path, seek_position, extra_input_args
3491
3492 # all other stream types (HTTP, FILE, etc)
3493 if stream_type == StreamType.ENCRYPTED_HTTP:
3494 assert streamdetails.decryption_key is not None # for type checking
3495 extra_input_args += ["-decryption_key", streamdetails.decryption_key]
3496 if isinstance(streamdetails.path, list):
3497 # multi part stream, which handles the seek itself
3498 return self.get_multi_file_stream(streamdetails, seek_position), 0, extra_input_args
3499 # regular single file/url stream
3500 assert isinstance(streamdetails.path, str) # for type checking
3501 return streamdetails.path, seek_position, extra_input_args
3502
3503 async def _iter_audio_source_pcm(
3504 self,
3505 streamdetails: StreamDetails,
3506 pcm_format: AudioFormat,
3507 ) -> AsyncGenerator[bytes]:
3508 """Yield PCM for an AudioSource, bypassing ffmpeg when formats match."""
3509 if _pcm_formats_match(streamdetails.audio_format, pcm_format):
3510 source_gen = self._open_audio_source_generator(streamdetails)
3511 async for chunk in realtime_pcm_pacer(source_gen, pcm_format):
3512 yield chunk
3513 return
3514 # format mismatch â fall back to ffmpeg for resampling (still small chunks)
3515 async for chunk in self.get_media_stream(
3516 streamdetails=streamdetails,
3517 pcm_format=pcm_format,
3518 filter_params=None,
3519 chunk_seconds=AUDIO_SOURCE_CHUNK_SECONDS,
3520 ):
3521 yield chunk
3522
3523 def _open_audio_source_generator(self, streamdetails: StreamDetails) -> AsyncGenerator[bytes]:
3524 """Open the raw PCM generator for an AudioSource (CUSTOM or NAMED_PIPE)."""
3525 if streamdetails.stream_type == StreamType.CUSTOM:
3526 # pin the exact instance, see _resolve_media_stream_source
3527 provider = self.mass.get_provider(streamdetails.provider, return_unavailable=True)
3528 if provider is None or not provider.available:
3529 raise ProviderUnavailableError(
3530 f"Provider {streamdetails.provider} for stream is no longer available"
3531 )
3532 provider = cast("MusicProvider | PluginProvider", provider)
3533 return provider.get_audio_stream(streamdetails)
3534 if streamdetails.stream_type == StreamType.NAMED_PIPE:
3535 assert isinstance(streamdetails.path, str) # for type checking
3536 return read_named_pipe(streamdetails.path)
3537 raise AudioError(f"Unsupported stream_type {streamdetails.stream_type} for AudioSource")
3538
3539 def _handle_inband_metadata(
3540 self, streamdetails: StreamDetails, metadata: dict[str, str]
3541 ) -> None:
3542 """Handle metadata extracted from a chained Ogg stream."""
3543 title = metadata.get("title", "")
3544 artist = metadata.get("artist", "")
3545 album = metadata.get("album", "")
3546 if not artist and " - " in title:
3547 artist, title = title.split(" - ", 1)
3548 if not (title or artist):
3549 return
3550
3551 stream_title = f"{artist} - {title}" if artist and title else title or artist
3552 cleaned_title = clean_stream_title(stream_title)
3553 if not cleaned_title:
3554 return
3555 if self._record_inband_stream_title(streamdetails, cleaned_title):
3556 return
3557 if cleaned_title != streamdetails.stream_title:
3558 self.logger.log(VERBOSE_LOG_LEVEL, "In-band metadata: %s", cleaned_title)
3559 streamdetails.stream_title = cleaned_title
3560 self._update_radio_stream_metadata(
3561 streamdetails,
3562 artist=artist or None,
3563 title=title or cleaned_title,
3564 album=album or None,
3565 )
3566
3567 def _record_inband_stream_title(self, streamdetails: StreamDetails, cleaned_title: str) -> bool:
3568 """
3569 Record an in-band stream title for provider-owned metadata, if applicable.
3570
3571 When a provider opts into owning stream_metadata (and stream_title is only
3572 a derived view of it), writing either from the stream reader would fight the
3573 provider. The cleaned in-band title is recorded on StreamDetails.data instead,
3574 as the identity signal for the provider callback.
3575
3576 :param streamdetails: StreamDetails carrying the stream.
3577 :param cleaned_title: Cleaned in-band stream title.
3578 :returns: True when recorded (the caller must not write stream metadata);
3579 False when no provider callback exists and normal handling applies.
3580 """
3581 if (
3582 streamdetails.stream_metadata_update_callback is None
3583 or streamdetails.data is None
3584 or not streamdetails.data.get(STREAMDETAILS_INBAND_TITLE_HANDOFF_KEY)
3585 ):
3586 return False
3587 if streamdetails.data.get(STREAMDETAILS_INBAND_TITLE_KEY) != cleaned_title:
3588 # occupancy approximates how far this detection leads audible playback
3589 buffer = streamdetails.buffer
3590 self.logger.debug(
3591 "In-band stream title: %s (buffer occupancy: %ss)",
3592 cleaned_title,
3593 buffer.size_seconds if buffer is not None else "unknown",
3594 )
3595 streamdetails.data[STREAMDETAILS_INBAND_TITLE_KEY] = cleaned_title
3596 return True
3597
3598 def _parse_icy_metadata(self, meta_data: bytes, streamdetails: StreamDetails) -> None:
3599 """
3600 Parse ICY metadata and update streamdetails.
3601
3602 Sets the cleaned stream title and, when the title parses as "Artist - Track",
3603 triggers a radio-artwork metadata update.
3604
3605 :param meta_data: Raw metadata bytes from an ICY stream chunk.
3606 :param streamdetails: StreamDetails to update with parsed title and metadata.
3607 """
3608 if not meta_data:
3609 return
3610
3611 meta_data = meta_data.rstrip(b"\0")
3612 # Match StreamTitle, handling apostrophes in titles
3613 stream_title_re = re.search(rb"StreamTitle='(.*?)';", meta_data)
3614
3615 if not stream_title_re:
3616 self.logger.log(
3617 VERBOSE_LOG_LEVEL,
3618 "ICY metadata does not contain StreamTitle field. Raw: %s",
3619 meta_data.decode("utf-8", errors="replace")[:200],
3620 )
3621 return
3622
3623 try:
3624 # in 99% of the cases the stream title is utf-8 encoded
3625 stream_title = stream_title_re.group(1).decode("utf-8")
3626 except UnicodeDecodeError:
3627 # fallback to iso-8859-1
3628 stream_title = stream_title_re.group(1).decode("iso-8859-1", errors="replace")
3629
3630 cleaned_stream_title = clean_stream_title(stream_title)
3631
3632 if not cleaned_stream_title:
3633 return
3634
3635 if self._record_inband_stream_title(streamdetails, cleaned_stream_title):
3636 return
3637
3638 if cleaned_stream_title == streamdetails.stream_title:
3639 return
3640
3641 self.logger.log(VERBOSE_LOG_LEVEL, "ICY Radio streamtitle original: %s", stream_title)
3642 self.logger.log(
3643 VERBOSE_LOG_LEVEL, "ICY Radio streamtitle cleaned: %s", cleaned_stream_title
3644 )
3645 streamdetails.stream_title = cleaned_stream_title
3646
3647 # Prefer station-provided cover art from the ICY 'StreamUrl' field (when it is
3648 # an image) over the MusicBrainz artwork lookup in _update_radio_stream_metadata.
3649 image_url = self._parse_icy_image_url(meta_data)
3650
3651 # Parse the original title for structured fields first so stations that announce
3652 # an album can refine the artwork lookup; fall back to the "Artist - Track" split.
3653 album: str | None = None
3654 if parsed := parse_quoted_stream_title(stream_title):
3655 track_name, artist_name_raw, album = parsed
3656 elif " - " in cleaned_stream_title:
3657 artist_name_raw, track_name = (
3658 part.strip() for part in cleaned_stream_title.split(" - ", 1)
3659 )
3660 else:
3661 return
3662
3663 if artist_name_raw and track_name:
3664 self.logger.debug(
3665 "ICY metadata: artist='%s', track='%s', album='%s'",
3666 artist_name_raw,
3667 track_name,
3668 album,
3669 )
3670 self._update_radio_stream_metadata(
3671 streamdetails,
3672 artist=artist_name_raw,
3673 title=track_name,
3674 album=album,
3675 image_url=image_url,
3676 )
3677
3678 def _parse_icy_image_url(self, meta_data: bytes) -> str | None:
3679 """
3680 Return a PNG or JPEG cover-art URL from the ICY 'StreamUrl' field, if present.
3681
3682 :param meta_data: Raw metadata bytes from an ICY stream chunk.
3683 """
3684 # The trailing semicolon is optional to match sources that omit it.
3685 stream_url_re = re.search(rb"StreamUrl='([^']*)'", meta_data)
3686 if not stream_url_re:
3687 return None
3688 try:
3689 image_url = stream_url_re.group(1).decode("utf-8").strip()
3690 except UnicodeDecodeError:
3691 return None
3692 if not image_url:
3693 return None
3694 # StreamUrl is not a standardized artwork field (reference clients such as VLC
3695 # ignore it and it conventionally holds a station website link), so only accept
3696 # values that point at a PNG or JPEG image.
3697 parsed = urlparse(image_url)
3698 if parsed.scheme not in ("http", "https"):
3699 return None
3700 if not parsed.path.lower().endswith((".png", ".jpg", ".jpeg")):
3701 return None
3702 self.logger.debug("ICY metadata: StreamUrl image='%s'", image_url)
3703 return image_url
3704
3705 async def _validate_shoutcast_stream(self, url: str) -> bool:
3706 """
3707 Return True if the URL responds with a legacy Shoutcast "ICY 200 OK" line.
3708
3709 :param url: The URL to validate.
3710 """
3711 try:
3712 parsed = urlparse(url)
3713 host = parsed.hostname
3714 port = parsed.port or 80
3715 path = parsed.path or "/"
3716 if parsed.query:
3717 path = f"{path}?{parsed.query}"
3718
3719 # Open raw socket connection with timeout
3720 reader, writer = await asyncio.wait_for(asyncio.open_connection(host, port), timeout=10)
3721 try:
3722 # Send minimal HTTP request with ICY metadata header
3723 request = f"GET {path} HTTP/1.1\r\nHost: {host}\r\nIcy-MetaData: 1\r\n\r\n"
3724 writer.write(request.encode())
3725 await writer.drain()
3726
3727 # Read just the response line
3728 response_line = await asyncio.wait_for(reader.readline(), timeout=5)
3729 finally:
3730 writer.close()
3731 await writer.wait_closed()
3732
3733 # Check if response starts with "ICY"
3734 decoded_line = response_line.decode("latin-1", errors="ignore").strip()
3735 return decoded_line.startswith("ICY")
3736
3737 except TimeoutError:
3738 self.logger.debug("Timeout during Shoutcast validation for %s", url)
3739 return False
3740 except OSError, ConnectionError:
3741 self.logger.debug("Connection failed during Shoutcast validation for %s", url)
3742 return False
3743 except UnicodeDecodeError:
3744 self.logger.debug("Invalid response encoding during Shoutcast validation for %s", url)
3745 return False
3746
3747 def _resolve_player_dsp_config(self, player: Player) -> DSPConfig:
3748 """
3749 Resolve the effective DSP config for a player.
3750
3751 Single source of truth shared by every code path that needs to know
3752 whether DSP will run for this player. Protocol wrappers defer to their
3753 parent player; single-leg ``player_group`` instances that don't expose
3754 ``MULTI_DEVICE_DSP`` defer to their first member; players whose grouping
3755 context prevents DSP get a disabled config back regardless.
3756
3757 :param player: The player to resolve DSP config for.
3758 """
3759 dsp_player_id = self._resolve_player_dsp_config_id(player)
3760 dsp = self.mass.config.get_player_dsp_config(dsp_player_id)
3761 if is_grouping_preventing_dsp(player):
3762 dsp.enabled = False
3763 elif player.provider.domain == "player_group" and (
3764 PlayerFeature.MULTI_DEVICE_DSP not in player.state.supported_features
3765 ):
3766 if not player.state.group_members:
3767 dsp.enabled = False
3768 return dsp
3769
3770 def _resolve_player_dsp_config_id(self, player: Player) -> str:
3771 """
3772 Return the player identifier that supplies the effective DSP config.
3773
3774 :param player: Player whose DSP config source should be resolved.
3775 """
3776 dsp_player_id = player.protocol_parent_id or player.player_id
3777 if (
3778 not is_grouping_preventing_dsp(player)
3779 and player.provider.domain == "player_group"
3780 and PlayerFeature.MULTI_DEVICE_DSP not in player.state.supported_features
3781 and player.state.group_members
3782 ):
3783 child_player = self.mass.players.get_player(player.state.group_members[0])
3784 assert child_player is not None
3785 dsp_player_id = child_player.player_id
3786 return dsp_player_id
3787
3788 def _get_output_channels(self, player: Player | None, player_id: str) -> str:
3789 """
3790 Return the configured output channels for the rendering player.
3791
3792 The value may be stored on the rendering player(protocol) itself (the
3793 protocol section of the config UI) or on its visible parent player (the
3794 native section); the rendering player's own stored value wins.
3795 """
3796 parent_id = player.protocol_parent_id if player and player.protocol_parent_id else player_id
3797 parent_value = self.mass.config.get_raw_player_config_value(
3798 parent_id, CONF_OUTPUT_CHANNELS, "stereo"
3799 )
3800 return self.mass.config.get_raw_player_config_value(
3801 player.player_id if player else player_id, CONF_OUTPUT_CHANNELS, parent_value
3802 )
3803
3804 def _pick_pcm_bit_depth(
3805 self,
3806 players: Iterable[Player],
3807 streamdetails: StreamDetails | None,
3808 crossfade_enabled: bool,
3809 overlay_active: bool = False,
3810 ) -> tuple[ContentType, int]:
3811 """
3812 Return ``(content_type, bit_depth)`` for an internal PCM stream.
3813
3814 F32 is chosen when audio processing (crossfade, audio overlay, volume
3815 normalization, DSP) will run on the stream â those need the extra
3816 headroom to avoid clipping and precision loss. Otherwise the source's
3817 native bit depth is reused so we don't waste memory upcasting a 16-bit
3818 stream to 32-bit just to pass it through. When the source is unknown
3819 (no streamdetails) we fall back to F32 conservatively.
3820 """
3821 if streamdetails is None:
3822 return INTERNAL_PCM_FORMAT.content_type, INTERNAL_PCM_FORMAT.bit_depth
3823 needs_headroom = (
3824 crossfade_enabled
3825 or overlay_active
3826 or streamdetails.volume_normalization_mode != VolumeNormalizationMode.DISABLED
3827 or any(self._resolve_player_dsp_config(player).enabled for player in players)
3828 )
3829 if needs_headroom:
3830 return INTERNAL_PCM_FORMAT.content_type, INTERNAL_PCM_FORMAT.bit_depth
3831 bit_depth = streamdetails.audio_format.bit_depth
3832 return ContentType.from_bit_depth(bit_depth), bit_depth
3833
3834 def _select_audio_source_pcm_format(
3835 self,
3836 player: Player,
3837 streamdetails: StreamDetails,
3838 supported_sample_rates: Iterable[int] | None = None,
3839 ) -> AudioFormat:
3840 """
3841 Return a passthrough PCM format for a realtime AudioSource item.
3842
3843 The format matches the source's native sample rate, bit depth and
3844 channel count whenever the player can accept them; if the player does
3845 not support the source's sample rate, it is snapped down to the
3846 closest supported rate. No F32 widening â realtime sources skip every
3847 processing stage that would otherwise need it. Surround sources are
3848 still folded down to stereo, which every output path requires anyway.
3849
3850 :param player: The player requesting the stream.
3851 :param streamdetails: Stream details for the AudioSource item.
3852 :param supported_sample_rates: Rates shared by every output player, if applicable.
3853 """
3854 resolved_sample_rates = (
3855 list(supported_sample_rates)
3856 if supported_sample_rates is not None
3857 else [sample_rate for sample_rate, _ in player.get_supported_sample_rates()]
3858 )
3859 source_rate = streamdetails.audio_format.sample_rate
3860 if source_rate in resolved_sample_rates:
3861 output_sample_rate = source_rate
3862 else:
3863 output_sample_rate = max(
3864 (rate for rate in resolved_sample_rates if rate <= source_rate),
3865 default=min(resolved_sample_rates),
3866 )
3867 bit_depth = streamdetails.audio_format.bit_depth
3868 return AudioFormat(
3869 content_type=ContentType.from_bit_depth(bit_depth),
3870 sample_rate=output_sample_rate,
3871 bit_depth=bit_depth,
3872 # a realtime source may announce more channels than anything downstream can
3873 # carry (a VBAN stream can be configured up to 8), and player handoff formats
3874 # copy this count straight through, so fold it here
3875 channels=min(streamdetails.audio_format.channels, 2),
3876 )
3877
3878 def _flow_restart_context(
3879 self, queue_id: str, protocol_player: Player | None
3880 ) -> tuple[str, list[int]]:
3881 """
3882 Resolve the flow mode config and supported sample rates for restart decisions.
3883
3884 Prefers the protocol player actually consuming the flow stream over the
3885 queue's (wrapper) player, whose config may lack the audio specific entries.
3886 """
3887 if protocol_player is None:
3888 protocol_player = self.mass.players.get_player(queue_id)
3889 if protocol_player is None:
3890 flow_mode_sample_rate_conf = self.mass.config.get_raw_player_config_value(
3891 queue_id, CONF_FLOW_MODE_SAMPLE_RATE, FLOW_MODE_SAMPLE_RATE_SMART
3892 )
3893 return flow_mode_sample_rate_conf, []
3894 flow_mode_sample_rate_conf = cast(
3895 "str",
3896 protocol_player.config.get_value(
3897 CONF_FLOW_MODE_SAMPLE_RATE, FLOW_MODE_SAMPLE_RATE_SMART
3898 ),
3899 )
3900 supported_sample_rates = sorted(
3901 {sr for sr, _ in protocol_player.get_supported_sample_rates()}
3902 )
3903 return flow_mode_sample_rate_conf, supported_sample_rates
3904
3905 def _flow_stream_needs_restart(
3906 self,
3907 queue_track: QueueItem,
3908 pcm_format: AudioFormat,
3909 supported_sample_rates: list[int],
3910 flow_mode_sample_rate_conf: str,
3911 is_first_track: bool,
3912 ) -> bool:
3913 """
3914 Return True if the upcoming queue track requires exiting the flow stream.
3915
3916 Covers every case where the flow loop should break and hand control back to
3917 the queue controller for restart:
3918
3919 - Live media (radio, audio sources): cannot be played inside a flow,
3920 the controller will fall back to a single-item stream.
3921 - Sample rate mismatch ('smart' / 'bit_perfect' modes only): the next
3922 track's sample rate (snapped up to the closest supported player rate,
3923 mirroring select_flow_pcm_format's anchoring logic) is incompatible with
3924 the current flow rate, so a new flow must be opened.
3925
3926 The first (anchor) track is always allowed to continue for the sample
3927 rate check; select_flow_pcm_format has already snapped the flow rate to it.
3928
3929 :param queue_track: The upcoming queue item.
3930 :param pcm_format: The current flow stream's PCM format.
3931 :param supported_sample_rates: Sorted list of the player's supported rates.
3932 :param flow_mode_sample_rate_conf: The flow mode sample rate config value.
3933 :param is_first_track: Whether this is the first track of the flow stream.
3934 """
3935 # live audio (radio, plugin or audio source) cannot be flowed; let the
3936 # queue controller fall back to single-item streaming for this item
3937 if queue_track.media_type in (MediaType.RADIO, MediaType.AUDIO_SOURCE):
3938 self.logger.info(
3939 "Live media item %s (%s, %s) encountered in flow stream "
3940 "- breaking out to single item stream",
3941 queue_track.queue_item_id,
3942 queue_track.name,
3943 queue_track.media_type,
3944 )
3945 return True
3946
3947 if is_first_track or queue_track.streamdetails is None:
3948 return False
3949 raw_next_rate = queue_track.streamdetails.audio_format.sample_rate
3950 if not raw_next_rate or not supported_sample_rates:
3951 return False
3952 effective_next_rate = _snap_supported_rate_up(raw_next_rate, supported_sample_rates)
3953
3954 # branch order mirrors select_flow_pcm_format: fixed-rate modes resample
3955 # everything to the chosen rate (no restart); bit_perfect restarts on any
3956 # mismatch; anything else falls through to smart-anchor behavior so
3957 # unknown/legacy config values don't silently pin the flow forever.
3958 if flow_mode_sample_rate_conf in (
3959 FLOW_MODE_SAMPLE_RATE_48000,
3960 FLOW_MODE_SAMPLE_RATE_96000,
3961 FLOW_MODE_SAMPLE_RATE_HIGHEST,
3962 ):
3963 needs_restart = False
3964 elif flow_mode_sample_rate_conf == FLOW_MODE_SAMPLE_RATE_BIT_PERFECT:
3965 needs_restart = effective_next_rate != pcm_format.sample_rate
3966 else:
3967 needs_restart = effective_next_rate > pcm_format.sample_rate
3968
3969 if needs_restart:
3970 self.logger.info(
3971 "Track %s (%s) sample rate %s (snapped to %s) incompatible with flow rate %s "
3972 "(mode: %s) - breaking out to restart flow stream",
3973 queue_track.queue_item_id,
3974 queue_track.name,
3975 raw_next_rate,
3976 effective_next_rate,
3977 pcm_format.sample_rate,
3978 flow_mode_sample_rate_conf,
3979 )
3980 return needs_restart
3981
3982 @asynccontextmanager
3983 async def _connect_radio_stream(self, url: str, **kwargs: Any) -> AsyncGenerator[Any]:
3984 """
3985 Connect to a radio stream URL with fallback for legacy SSL/TLS configurations.
3986
3987 Some radio servers use outdated TLS configurations that reject modern
3988 cipher suites. Since radio streams are public broadcast content,
3989 relaxing cipher requirements is acceptable.
3990
3991 :param url: The radio stream URL to connect to.
3992 :param kwargs: Additional keyword arguments passed to aiohttp get().
3993 """
3994 request_url = encoded_request_url(url)
3995 try:
3996 async with self.mass.http_session_no_ssl.get(request_url, **kwargs) as resp:
3997 yield resp
3998 except ClientConnectorSSLError:
3999 self.logger.info(
4000 "SSL handshake failed for %s, retrying with permissive cipher configuration", url
4001 )
4002 insecure_ssl_context = ssl_util.client_context_no_verify(
4003 ssl_util.SSLCipherList.INSECURE
4004 )
4005 async with self.mass.http_session_no_ssl.get(
4006 request_url, ssl=insecure_ssl_context, **kwargs
4007 ) as resp:
4008 yield resp
4009
4010 async def _update_hls_radio_metadata(
4011 self,
4012 streamdetails: StreamDetails,
4013 elapsed_time: int,
4014 ) -> None:
4015 """
4016 Update HLS radio stream metadata by fetching the playlist.
4017
4018 Fetches the HLS playlist and extracts metadata from EXTINF lines.
4019
4020 :param streamdetails: StreamDetails object to update with metadata
4021 :param elapsed_time: Current playback position in seconds (unused for live radio)
4022 """
4023 mass = self.mass
4024 try:
4025 # Get the actual media playlist URL from cache or resolve it
4026 # We cache the media_playlist_url in streamdetails.data to avoid re-resolving
4027 if streamdetails.data is None:
4028 streamdetails.data = {}
4029 media_playlist_url = streamdetails.data.get("hls_media_playlist_url")
4030 if not media_playlist_url:
4031 try:
4032 assert isinstance(streamdetails.path, str) # for type checking
4033 substream = await self.get_hls_substream(streamdetails.path)
4034 media_playlist_url = substream.path
4035 streamdetails.data["hls_media_playlist_url"] = media_playlist_url
4036 except Exception as err:
4037 self.logger.warning(
4038 "Failed to resolve HLS substream for metadata monitoring: %s", err
4039 )
4040 return
4041
4042 # Fetch the media playlist
4043 timeout = ClientTimeout(total=0, connect=10, sock_read=30)
4044 try:
4045 async with mass.http_session_no_ssl.get(
4046 encoded_request_url(media_playlist_url), timeout=timeout
4047 ) as resp:
4048 resp.raise_for_status()
4049 playlist_content = await resp.text()
4050 except ClientResponseError as err:
4051 # Session token likely expired (410/403) â drop cache so next poll re-resolves
4052 if err.status in (403, 410):
4053 streamdetails.data.pop("hls_media_playlist_url", None)
4054 raise
4055
4056 # Parse the playlist and look for EXTINF metadata
4057 # The most recent segment usually has the current metadata
4058 lines = playlist_content.strip().split("\n")
4059 for line in reversed(lines):
4060 if line.startswith("#EXTINF:"):
4061 # Extract metadata from EXTINF line
4062 metadata = parse_extinf_metadata(line)
4063
4064 # Build stream title from title and artist
4065 title = metadata.get("title", "")
4066 artist = metadata.get("artist", "")
4067 image_url = (
4068 metadata.get("image") or metadata.get("artwork") or metadata.get("cover")
4069 )
4070 if not artist and " - " in title:
4071 artist, title = title.split(" - ", 1)
4072 if title or artist:
4073 # Format as "Artist - Title"
4074 if artist and title:
4075 stream_title = f"{artist} - {title}"
4076 elif title:
4077 stream_title = title
4078 else:
4079 stream_title = artist
4080
4081 # Clean the stream title
4082 cleaned_title = clean_stream_title(stream_title)
4083
4084 # Only update if changed
4085 if cleaned_title != streamdetails.stream_title and cleaned_title:
4086 self.logger.log(
4087 VERBOSE_LOG_LEVEL, "HLS Radio metadata updated: %s", cleaned_title
4088 )
4089 streamdetails.stream_title = cleaned_title
4090 self._update_radio_stream_metadata(
4091 streamdetails,
4092 artist=artist or None,
4093 title=title or cleaned_title,
4094 image_url=image_url,
4095 )
4096
4097 # Only check the most recent EXTINF
4098 break
4099
4100 except Exception as err:
4101 self.logger.debug("Error fetching HLS metadata: %s", err)
4102
4103 @staticmethod
4104 def _normalize_reconnecting_urls(url: str | list[MultiPartPath]) -> list[str]:
4105 """Normalize a single URL or a sequence into a non-empty list."""
4106 if isinstance(url, str):
4107 return [url]
4108 if not url:
4109 msg = "Radio stream requires at least one URL"
4110 raise InvalidDataError(msg)
4111 return [part.path for part in url]
4112
4113 async def _resolve_overlay_input(self, queue: PlayerQueue) -> str | None:
4114 """
4115 Resolve the queue's overlay source to a file path or URL for ffmpeg.
4116
4117 Returns None (with a warning logged) when the source can not be resolved,
4118 so the caller can degrade to music-only playback.
4119 """
4120 if not (mapping := queue.overlay_source):
4121 return None
4122 try:
4123 provider = self.mass.get_provider(mapping.provider)
4124 if provider is None:
4125 raise MediaNotFoundError(f"Provider {mapping.provider} is not available")
4126 stream_prov = cast("MusicProvider | PluginProvider", provider)
4127 streamdetails = await stream_prov.get_stream_details(
4128 mapping.item_id, MediaType.SOUND_EFFECT
4129 )
4130 except Exception as err:
4131 self.logger.warning(
4132 "Audio overlay source %s is unavailable (%s) - continuing without overlay",
4133 mapping.uri,
4134 str(err) or err.__class__.__name__,
4135 )
4136 return None
4137 if streamdetails.stream_type not in (StreamType.LOCAL_FILE, StreamType.HTTP) or not (
4138 isinstance(streamdetails.path, str)
4139 ):
4140 self.logger.warning(
4141 "Audio overlay source %s uses unsupported stream type %s "
4142 "- continuing without overlay",
4143 mapping.uri,
4144 streamdetails.stream_type,
4145 )
4146 return None
4147 if streamdetails.stream_type == StreamType.LOCAL_FILE and not await aiofiles.os.path.isfile(
4148 streamdetails.path
4149 ):
4150 # guard against stale sources: feeding a missing file to the mixer would
4151 # kill the whole (music) stream instead of just the overlay
4152 self.logger.warning(
4153 "Audio overlay source %s does not exist - continuing without overlay",
4154 streamdetails.path,
4155 )
4156 return None
4157 return streamdetails.path
4158