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