/
/
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 pq_data.flow_mode_stream_log = []
1958 if not start_queue_item:
1959 # this can happen in some (edge case) race conditions
1960 return
1961 pcm_sample_size = pcm_format.pcm_sample_size
1962 if start_queue_item.media_type != MediaType.TRACK:
1963 # no crossfade on non-tracks
1964 crossfade_mode = CrossfadeMode.DISABLED
1965 standard_crossfade_duration = 0
1966 else:
1967 crossfade_mode = self.mass.streams.get_crossfade_mode(queue)
1968 # crossfade duration is a global (queue controller) setting; fallback matches
1969 # CONF_ENTRY_CROSSFADE_DURATION's default
1970 standard_crossfade_duration = self.mass.config.get_raw_core_config_value(
1971 CONF_PLAYER_QUEUES, CONF_CROSSFADE_DURATION, 8
1972 )
1973 flow_mode_sample_rate_conf, flow_supported_sample_rates = self._flow_restart_context(
1974 queue.queue_id, protocol_player
1975 )
1976 # note: get_crossfade_mode() already falls back to standard when smart fades aren't
1977 # available (no analysis provider / minimal buffer), so crossfade_mode is safe to use.
1978 self.logger.info(
1979 "Start Queue Flow stream for Queue %s - crossfade: %s %s",
1980 queue.display_name,
1981 crossfade_mode,
1982 f"({standard_crossfade_duration}s)"
1983 if crossfade_mode == CrossfadeMode.STANDARD_CROSSFADE
1984 else "",
1985 )
1986 total_chunks_received = 0
1987
1988 def _superseded() -> bool:
1989 """Return True if a newer stream session has taken over this queue."""
1990 return pq_data.session_id != flow_session_id
1991
1992 queue_exhausted = False
1993 while True:
1994 # bail out early if a newer producer has taken over this queue,
1995 # so we don't append another entry to a stream log we no longer own
1996 if _superseded():
1997 self.logger.debug(
1998 "Flow stream for queue %s superseded (session %s -> %s) "
1999 "- exiting before next track",
2000 queue.display_name,
2001 flow_session_id,
2002 pq_data.session_id,
2003 )
2004 return
2005 # get (next) queue item to stream
2006 if queue_track is None:
2007 queue_track = start_queue_item
2008 else:
2009 try:
2010 queue_track = await self.mass.player_queues.load_next_queue_item(
2011 queue.queue_id, queue_track.queue_item_id
2012 )
2013 except QueueEmpty:
2014 queue_exhausted = True
2015 break
2016
2017 if self._flow_stream_needs_restart(
2018 queue_track,
2019 pcm_format,
2020 flow_supported_sample_rates,
2021 flow_mode_sample_rate_conf,
2022 is_first_track=queue_track is start_queue_item,
2023 ):
2024 break
2025
2026 if queue_track.streamdetails is None:
2027 self.logger.error(
2028 "No StreamDetails for queue item %s (%s) on queue %s - skipping track",
2029 queue_track.queue_item_id,
2030 queue_track.name,
2031 queue.display_name,
2032 )
2033 continue
2034 if flow_session_id is not None:
2035 self.mass.streams.audio_processing.update_item_context(
2036 queue_id=queue.queue_id,
2037 session_id=flow_session_id,
2038 queue_item_id=queue_track.queue_item_id,
2039 queue_processing=AudioQueueProcessing(
2040 pcm_format=pcm_format,
2041 playback_speed=cast(
2042 "float",
2043 queue_track.extra_attributes.get("playback_speed", 1.0),
2044 ),
2045 crossfade_mode=crossfade_mode,
2046 overlay_active=overlay_active(queue),
2047 ),
2048 alters_audio=queue_track.streamdetails.fade_in,
2049 )
2050
2051 self.logger.debug(
2052 "Start Streaming queue track: %s (%s) for queue %s",
2053 queue_track.streamdetails.uri,
2054 queue_track.name,
2055 queue.display_name,
2056 )
2057 # last chance to bail before mutating the stream log: a newer producer
2058 # may have taken over while we were awaiting load_next_queue_item
2059 if _superseded():
2060 self.logger.debug(
2061 "Flow stream for queue %s superseded - exiting before playlog append",
2062 queue.display_name,
2063 )
2064 return
2065 track_playback_speed = cast(
2066 "float", queue_track.extra_attributes.get("playback_speed", 1.0)
2067 )
2068 # calculate crossfade buffer size
2069 crossfade_buffer_duration = (
2070 SMART_CROSSFADE_DURATION
2071 if crossfade_mode == CrossfadeMode.SMART_CROSSFADE
2072 else standard_crossfade_duration
2073 )
2074 crossfade_buffer_duration = min(
2075 crossfade_buffer_duration,
2076 int(queue_track.streamdetails.duration / 2)
2077 if queue_track.streamdetails.duration
2078 else crossfade_buffer_duration,
2079 )
2080 # skip crossfade if buffer would be too small to be meaningful
2081 if crossfade_buffer_duration < MIN_CROSSFADE_FALLBACK_DURATION:
2082 crossfade_buffer_duration = 0
2083 # Ensure crossfade buffer size is aligned to frame boundaries
2084 # Frame size = bytes_per_sample * channels
2085 bytes_per_sample = pcm_format.bit_depth // 8
2086 frame_size = bytes_per_sample * pcm_format.channels
2087 crossfade_buffer_size = int(pcm_format.pcm_sample_size * crossfade_buffer_duration)
2088 # Round down to nearest frame boundary
2089 crossfade_buffer_size = (crossfade_buffer_size // frame_size) * frame_size
2090 warmup_size = int(pcm_format.pcm_sample_size * WARMUP_DURATION)
2091
2092 # raw_seek_position feeds the PCM buffer; streamdetails.seek_position
2093 # (overwritten below) only drives reported elapsed time.
2094 raw_seek_position = queue_track.streamdetails.seek_position
2095 # Build eagerly so seek_position is set before PlayLogEntry is appended â
2096 # consumer-paced mix() would otherwise let the queue briefly report 0.
2097 crossfade_smart_fade: SmartFade | None = None
2098 incoming_crossfade_size = crossfade_buffer_size
2099 incoming_audio_buffer: AudioBuffer | None = None
2100 if last_fadeout_part and last_streamdetails:
2101 transition_mode = CrossfadeMode.DISABLED
2102 incoming_duration = 0.0
2103 if crossfade_buffer_size > 0 and crossfade_mode != CrossfadeMode.DISABLED:
2104 transition_mode, incoming_duration = self._select_buffered_crossfade(
2105 queue_track.streamdetails,
2106 crossfade_mode,
2107 standard_crossfade_duration,
2108 track_playback_speed,
2109 )
2110 if transition_mode == CrossfadeMode.DISABLED:
2111 # nothing to fade into: flush the held-back tail of the previous track
2112 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2113 yield pcm_slice
2114 await asyncio.sleep(0)
2115 last_fadeout_part = b""
2116 last_streamdetails = None
2117 last_play_log_entry = None
2118 else:
2119 assert queue_track.streamdetails.buffer is not None
2120 incoming_audio_buffer = cast("AudioBuffer", queue_track.streamdetails.buffer)
2121 incoming_crossfade_size = int(pcm_format.pcm_sample_size * incoming_duration)
2122 incoming_crossfade_size = (incoming_crossfade_size // frame_size) * frame_size
2123 crossfade_smart_fade = await self.smart_fades_mixer.build(
2124 fade_in_streamdetails=queue_track.streamdetails,
2125 fade_out_streamdetails=last_streamdetails,
2126 pcm_format=pcm_format,
2127 standard_crossfade_duration=standard_crossfade_duration,
2128 mode=transition_mode,
2129 fade_out_data=last_fadeout_part,
2130 fade_in_bytes_len=incoming_crossfade_size,
2131 )
2132 timing_info = crossfade_smart_fade.timing_info
2133 queue_track.streamdetails.seek_position = (
2134 raw_seek_position
2135 + (timing_info.fadein_trimmed_duration + timing_info.crossfade_duration)
2136 * track_playback_speed
2137 )
2138 # append to play log so the queue controller can work out which track is playing
2139 play_log_entry = PlayLogEntry(queue_track.queue_item_id)
2140 pq_data.flow_mode_stream_log.append(play_log_entry)
2141
2142 bytes_written = 0
2143 crossfade_buffer = bytearray()
2144 warmup_bytes = 0
2145 first_chunk_received = False
2146
2147 async for chunk in self.get_queue_item_stream(
2148 queue_track,
2149 pcm_format=pcm_format,
2150 seek_position=int(raw_seek_position),
2151 playback_speed=cast(
2152 "float", queue_track.extra_attributes.get("playback_speed", 1.0)
2153 ),
2154 raise_on_error=False,
2155 session_id=flow_session_id,
2156 prepared_buffer=incoming_audio_buffer,
2157 ):
2158 # if a newer producer has taken over this queue, stop sending
2159 # audio and exit cleanly before the outer-loop end-of-track
2160 # bookkeeping mutates seconds_streamed / duration on the log
2161 if _superseded():
2162 self.logger.debug(
2163 "Flow stream for queue %s superseded - stopping chunk yield",
2164 queue.display_name,
2165 )
2166 return
2167 total_chunks_received += 1
2168 if not first_chunk_received:
2169 first_chunk_received = True
2170 # inform the queue that the track is now loaded in the buffer
2171 # so the next track can be preloaded
2172 self.mass.player_queues.track_loaded_in_buffer(
2173 queue.queue_id, queue_track.queue_item_id
2174 )
2175
2176 if crossfade_mode == CrossfadeMode.DISABLED:
2177 # no cross/smart fade: yield chunks directly without intermediate buffer
2178 yield chunk
2179 bytes_written += len(chunk)
2180 del chunk
2181 continue
2182
2183 # Warmup: yield chunks directly until we have streamed WARMUP_DURATION
2184 # worth of audio, so playback starts immediately. Skip warmup when
2185 # crossfade data from the previous track is pending â we need a full
2186 # buffer for the mix.
2187 if warmup_bytes < warmup_size and not last_fadeout_part:
2188 yield chunk
2189 warmup_bytes += len(chunk)
2190 bytes_written += len(chunk)
2191 del chunk
2192 continue
2193
2194 # smart fades enabled: accumulate chunks in crossfade buffer
2195 crossfade_buffer.extend(chunk)
2196 del chunk
2197 required_buffer_size = (
2198 incoming_crossfade_size if last_fadeout_part else crossfade_buffer_size
2199 )
2200 if len(crossfade_buffer) < required_buffer_size:
2201 await asyncio.sleep(0)
2202 continue
2203
2204 # handle crossfade of previous track and new track
2205 if (
2206 last_fadeout_part
2207 and last_streamdetails
2208 and crossfade_smart_fade is not None
2209 and last_play_log_entry is not None
2210 ):
2211 fadein_part = bytes(crossfade_buffer[:incoming_crossfade_size])
2212 remaining_bytes = bytes(crossfade_buffer[incoming_crossfade_size:])
2213 try:
2214 crossfade_bytes_written = 0
2215 async for mix_chunk in self.smart_fades_mixer.mix(
2216 crossfade_smart_fade,
2217 fade_in_part=fadein_part,
2218 fade_out_part=last_fadeout_part,
2219 pcm_format=pcm_format,
2220 ):
2221 yield mix_chunk
2222 crossfade_bytes_written += len(mix_chunk)
2223 except Exception as mix_err:
2224 if crossfade_bytes_written:
2225 # partial mix already played â concat'd tail would duplicate audio
2226 raise
2227 self.logger.warning(
2228 "Crossfade mixer failed for %s, falling back to simple concat: %s",
2229 queue_track.name,
2230 mix_err,
2231 )
2232 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2233 yield pcm_slice
2234 await asyncio.sleep(0)
2235 # full tail was pre-counted and is now yielded as-is
2236 crossfade_bytes_written = 0
2237 remaining_bytes = bytes(crossfade_buffer)
2238 # mix failed â undo the eager seek_position
2239 queue_track.streamdetails.seek_position = raw_seek_position
2240 if crossfade_bytes_written:
2241 # Split mix output at end-of-overlap: PRE+CF to A, POST to B.
2242 fadeout_share_seconds = (
2243 timing_info.pre_crossfade_duration + timing_info.crossfade_duration
2244 )
2245 fadeout_share = int(fadeout_share_seconds * pcm_sample_size)
2246 fadeout_share = (fadeout_share // frame_size) * frame_size
2247 fadeout_share = min(fadeout_share, crossfade_bytes_written)
2248 fadein_share = crossfade_bytes_written - fadeout_share
2249 bytes_written += fadein_share
2250 if last_play_log_entry:
2251 assert last_play_log_entry.seconds_streamed is not None
2252 # correct pre-counted full tail to the timing-based share
2253 last_play_log_entry.seconds_streamed += (
2254 fadeout_share - len(last_fadeout_part)
2255 ) / pcm_sample_size
2256 if remaining_bytes:
2257 for pcm_slice in iter_pcm_slices(remaining_bytes, pcm_format, 1000):
2258 yield pcm_slice
2259 await asyncio.sleep(0)
2260 bytes_written += len(remaining_bytes)
2261 del remaining_bytes
2262 last_fadeout_part = b""
2263 last_streamdetails = None
2264 crossfade_buffer = bytearray()
2265 warmup_bytes = 0
2266
2267 # yield everything above the crossfade buffer size
2268 while len(crossfade_buffer) > crossfade_buffer_size:
2269 yield bytes(crossfade_buffer[:pcm_sample_size])
2270 bytes_written += pcm_sample_size
2271 del crossfade_buffer[:pcm_sample_size]
2272 await asyncio.sleep(0)
2273
2274 # A source error after partial audio must not look like a completed item.
2275 # Progress reporting skips items with stream_error, so the item is not
2276 # marked played; move on to the next queue item like the zero-audio path.
2277 if first_chunk_received and queue_track.streamdetails.stream_error:
2278 if _superseded():
2279 return
2280 self.logger.warning(
2281 "Track %s (%s) on queue %s aborted by a stream error - skipping",
2282 queue_track.name,
2283 queue_track.streamdetails.uri,
2284 queue.display_name,
2285 )
2286 # the audio sent so far will still play out; keep the play log entry
2287 # honest about how much of this item was actually streamed
2288 play_log_entry.seconds_streamed = bytes_written / pcm_sample_size
2289 if last_fadeout_part:
2290 # crossfade into this item never happened â undo the eager seek_position
2291 queue_track.streamdetails.seek_position = raw_seek_position
2292 continue
2293
2294 #### HANDLE END OF TRACK
2295 if not first_chunk_received:
2296 self.logger.warning(
2297 "Track %s (%s) on queue %s produced no audio data - skipping",
2298 queue_track.name,
2299 queue_track.streamdetails.uri if queue_track.streamdetails else "unknown",
2300 queue.display_name,
2301 )
2302 queue_track.streamdetails.stream_error = True
2303 play_log_entry.seconds_streamed = 0
2304 if last_fadeout_part:
2305 queue_track.streamdetails.seek_position = raw_seek_position
2306 continue
2307 if last_fadeout_part:
2308 # edge case: we did not get enough data to make the crossfade
2309 # attribute these bytes to the previous track (they are its tail)
2310 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2311 yield pcm_slice
2312 await asyncio.sleep(0)
2313 # no crossfade happened â undo the eager seek_position
2314 queue_track.streamdetails.seek_position = raw_seek_position
2315 # full tail was pre-counted and is now yielded as-is
2316 last_fadeout_part = b""
2317 if self.crossfade_allowed(
2318 queue_track,
2319 crossfade_mode=crossfade_mode,
2320 player_id=queue.queue_id,
2321 flow_mode=True,
2322 ):
2323 last_fadeout_part = bytes(crossfade_buffer[-crossfade_buffer_size:])
2324 last_streamdetails = queue_track.streamdetails
2325 last_play_log_entry = play_log_entry
2326 remaining_bytes = bytes(crossfade_buffer[:-crossfade_buffer_size])
2327 if remaining_bytes:
2328 for pcm_slice in iter_pcm_slices(remaining_bytes, pcm_format, 1000):
2329 yield pcm_slice
2330 await asyncio.sleep(0)
2331 bytes_written += len(remaining_bytes)
2332 del remaining_bytes
2333 elif crossfade_mode != CrossfadeMode.DISABLED and crossfade_buffer:
2334 bytes_written += len(crossfade_buffer)
2335 for pcm_slice in iter_pcm_slices(bytes(crossfade_buffer), pcm_format, 1000):
2336 yield pcm_slice
2337 await asyncio.sleep(0)
2338 crossfade_buffer = bytearray()
2339
2340 # update duration details based on the actual pcm data we sent
2341 # this also accounts for crossfade and silence stripping
2342 seconds_streamed = bytes_written / pcm_sample_size
2343 queue_track.streamdetails.seconds_streamed = seconds_streamed
2344 play_log_entry.seconds_streamed = seconds_streamed
2345 # an externally aborted source ends in a clean EOF mid-track, so the
2346 # streamed length must not be written back as the item's duration
2347 source_buffer = queue_track.streamdetails.buffer
2348 source_aborted = source_buffer is not None and source_buffer.cancelled
2349 if not source_aborted:
2350 # the held-back crossfade tail still counts as this track's media-time
2351 tail_seconds = len(last_fadeout_part) / pcm_sample_size
2352 # streamdetails.duration is in media-time; seconds_streamed is stream-time
2353 # (post-atempo), so we scale by the track's playback_speed to recover media-time.
2354 queue_track.streamdetails.duration = int(
2355 queue_track.streamdetails.seek_position
2356 + (seconds_streamed + tail_seconds) * track_playback_speed
2357 )
2358 # propagate accurate duration to queue_item so UI displays it
2359 queue_track.duration = queue_track.streamdetails.duration
2360 play_log_entry.duration = queue_track.streamdetails.duration
2361 if last_play_log_entry is play_log_entry and last_fadeout_part:
2362 # Pre-count the full crossfade tail so the queue index calculation
2363 # doesn't undercount while waiting for the next track's crossfade mix.
2364 # This will be corrected to crossfade_total/2 once the mix completes.
2365 assert play_log_entry.seconds_streamed is not None
2366 play_log_entry.seconds_streamed += len(last_fadeout_part) / pcm_sample_size
2367 self.logger.debug(
2368 "Finished Streaming queue track: %s (%s) on queue %s",
2369 queue_track.streamdetails.uri,
2370 queue_track.name,
2371 queue.display_name,
2372 )
2373 #### HANDLE END OF QUEUE FLOW STREAM
2374 # skip end-of-queue bookkeeping if a newer producer has superseded us;
2375 # the new producer owns queue_buffer_completed and the play log now
2376 if _superseded():
2377 self.logger.debug(
2378 "Flow stream for queue %s superseded - skipping end-of-queue handling",
2379 queue.display_name,
2380 )
2381 return
2382 # end of queue flow: make sure we yield the last_fadeout_part
2383 if last_fadeout_part:
2384 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2385 yield pcm_slice
2386 await asyncio.sleep(0)
2387 # correct seconds streamed - the duration already includes the tail
2388 last_part_seconds = len(last_fadeout_part) / pcm_sample_size
2389 streamdetails = queue_track.streamdetails
2390 assert streamdetails is not None
2391 streamdetails.seconds_streamed = (
2392 streamdetails.seconds_streamed or 0
2393 ) + last_part_seconds
2394 # also update the play log entry so elapsed time tracking stays in sync
2395 if last_play_log_entry:
2396 assert last_play_log_entry.seconds_streamed is not None
2397 # full tail was pre-counted and is now yielded as-is
2398 last_play_log_entry.duration = streamdetails.duration
2399 last_fadeout_part = b""
2400 self.logger.info("Finished Queue Flow stream for Queue %s", queue.display_name)
2401 # only signal completion if we are still the active producer â a later
2402 # producer would (incorrectly) see this as its own completion otherwise
2403 if not _superseded():
2404 # inform the queue controller that all audio data has been generated
2405 # so it can handle the case where new items were added after the flow stream ended
2406 self.mass.player_queues.queue_buffer_completed(queue.queue_id, queue_exhausted)
2407
2408 async def get_overlay_mixed_stream(
2409 self,
2410 queue: PlayerQueue,
2411 audio_input: AsyncGenerator[bytes],
2412 pcm_format: AudioFormat,
2413 ) -> AsyncGenerator[bytes]:
2414 """
2415 Mix the queue's audio overlay (looping sound effect) into the given PCM stream.
2416
2417 The mixed output has the exact same PCM format, duration and chunking as the
2418 input stream. If the overlay source can not be resolved, the original stream
2419 is passed through unchanged so playback is never interrupted.
2420
2421 :param queue: The PlayerQueue holding the overlay source and volume.
2422 :param audio_input: The audio stream (raw PCM in ``pcm_format``) to mix into.
2423 :param pcm_format: PCM format of both the input and the mixed output.
2424 """
2425 overlay_input = await self._resolve_overlay_input(queue)
2426 if overlay_input is None:
2427 # overlay source unavailable: degrade gracefully to music-only
2428 async for chunk in audio_input:
2429 yield chunk
2430 return
2431 async for chunk in get_ffmpeg_overlay_stream(
2432 audio_input=audio_input,
2433 overlay_input=overlay_input,
2434 pcm_format=pcm_format,
2435 overlay_volume=queue.overlay_volume,
2436 chunk_size=pcm_format.pcm_sample_size,
2437 ):
2438 yield chunk
2439
2440 def crossfade_allowed(
2441 self,
2442 queue_item: QueueItem,
2443 crossfade_mode: CrossfadeMode,
2444 player_id: str,
2445 flow_mode: bool = False,
2446 next_queue_item: QueueItem | None = None,
2447 sample_rate: int | None = None,
2448 next_sample_rate: int | None = None,
2449 ) -> bool:
2450 """Get the crossfade config for a queue item."""
2451 if crossfade_mode == CrossfadeMode.DISABLED:
2452 return False
2453 if not (self.mass.player_queues.get(queue_item.queue_id)):
2454 return False # just a guard
2455 if not (self.mass.players.get_player(player_id)):
2456 return False # just a guard
2457 if queue_item.media_type != MediaType.TRACK:
2458 self.logger.debug("Skipping crossfade: current item is not a track")
2459 return False
2460 # check if the next item is part of the same album
2461 next_item = next_queue_item or self.mass.player_queues.get_next_item(
2462 queue_item.queue_id, queue_item.queue_item_id
2463 )
2464 if not next_item:
2465 # there is no next item!
2466 return False
2467 # check if next item is a track
2468 if next_item.media_type != MediaType.TRACK:
2469 self.logger.debug("Skipping crossfade: next item is not a track")
2470 return False
2471 if (
2472 isinstance(queue_item.media_item, Track)
2473 and isinstance(next_item.media_item, Track)
2474 and queue_item.media_item.album
2475 and next_item.media_item.album
2476 and queue_item.media_item.album == next_item.media_item.album
2477 and not self.mass.config.get_raw_core_config_value(
2478 "streams", CONF_ALLOW_CROSSFADE_SAME_ALBUM, False
2479 )
2480 ):
2481 # in general, crossfade is not desired for tracks of the same (gapless) album
2482 # because we have no accurate way to determine if the album is gapless or not,
2483 # for now we just never crossfade between tracks of the same album
2484 self.logger.debug("Skipping crossfade: next item is part of the same album")
2485 return False
2486
2487 # check if we're allowed to crossfade on different sample rates
2488 if (
2489 not flow_mode
2490 and sample_rate
2491 and next_sample_rate
2492 and sample_rate != next_sample_rate
2493 and not self.mass.config.get_raw_player_config_value(
2494 player_id,
2495 CONF_ENTRY_CROSSFADE_DIFFERENT_SAMPLE_RATES.key,
2496 CONF_ENTRY_CROSSFADE_DIFFERENT_SAMPLE_RATES.default_value,
2497 )
2498 ):
2499 self.logger.debug(
2500 "Skipping crossfade: player(protocol) does not support gapless playback "
2501 "with different sample rates (%s vs %s)",
2502 sample_rate,
2503 next_sample_rate,
2504 )
2505 return False
2506
2507 return True
2508
2509 def clear_crossfade_data(self, queue_id: str) -> None:
2510 """
2511 Clear any pending crossfade data for a queue.
2512
2513 :param queue_id: The queue ID to clear crossfade data for.
2514 """
2515 if queue_id in self._crossfade_data:
2516 self.logger.debug("Clearing crossfade data for queue %s", queue_id)
2517 del self._crossfade_data[queue_id]
2518
2519 async def get_shoutcast_stream(
2520 self, url: str, streamdetails: StreamDetails
2521 ) -> AsyncGenerator[bytes]:
2522 """
2523 Yield audio from a legacy Shoutcast server, with ICY metadata parsed inline.
2524
2525 :param url: Shoutcast stream URL.
2526 :param streamdetails: StreamDetails to update with ICY metadata as it arrives.
2527 """
2528 self.logger.debug("Start streaming from legacy Shoutcast server: %s", url)
2529
2530 parsed = urlparse(url)
2531 host = parsed.hostname
2532 port = parsed.port or 80
2533 path = parsed.path or "/"
2534 if parsed.query:
2535 path = f"{path}?{parsed.query}"
2536
2537 try:
2538 # Open raw socket connection
2539 reader, writer = await asyncio.wait_for(asyncio.open_connection(host, port), timeout=30)
2540 except TimeoutError as err:
2541 raise AudioError(f"Timeout connecting to Shoutcast stream {url}") from err
2542 except (OSError, ConnectionError) as err:
2543 raise AudioError(f"Failed to connect to Shoutcast stream {url}") from err
2544
2545 try:
2546 # Send HTTP request with ICY metadata header
2547 request = (
2548 f"GET {path} HTTP/1.1\r\n"
2549 f"Host: {host}\r\n"
2550 f"User-Agent: {HTTP_HEADERS['User-Agent']}\r\n"
2551 f"Icy-MetaData: 1\r\n\r\n"
2552 )
2553 writer.write(request.encode())
2554 await writer.drain()
2555
2556 # Read and parse response line
2557 try:
2558 response_line = await asyncio.wait_for(reader.readline(), timeout=10)
2559 except TimeoutError as err:
2560 raise AudioError("Timeout reading Shoutcast response") from err
2561
2562 if not response_line.startswith(b"ICY"):
2563 raise InvalidDataError("Invalid Shoutcast response")
2564
2565 # Read headers until empty line
2566 headers: dict[str, str] = {}
2567 while True:
2568 try:
2569 line = await asyncio.wait_for(reader.readline(), timeout=5)
2570 except TimeoutError as err:
2571 raise AudioError("Timeout reading Shoutcast headers") from err
2572
2573 if line in (b"\r\n", b"\n", b""):
2574 break
2575
2576 if b":" in line:
2577 try:
2578 key, value = line.decode("latin-1", errors="ignore").split(":", 1)
2579 headers[key.strip().lower()] = value.strip()
2580 except UnicodeDecodeError, ValueError:
2581 continue
2582
2583 # Get metadata interval
2584 meta_int_str = headers.get("icy-metaint")
2585 if not meta_int_str:
2586 raise InvalidDataError("No icy-metaint header in Shoutcast response")
2587
2588 try:
2589 meta_int = int(meta_int_str)
2590 except ValueError as err:
2591 raise InvalidDataError("Invalid icy-metaint value") from err
2592
2593 self.logger.debug("Connected to Shoutcast stream %s (icy-metaint: %s)", url, meta_int)
2594
2595 # Stream audio data with metadata parsing
2596 while True:
2597 try:
2598 # Read audio chunk
2599 audio_chunk = await reader.readexactly(meta_int)
2600 yield audio_chunk
2601
2602 # Read metadata length
2603 meta_byte = await reader.readexactly(1)
2604 if meta_byte == b"\x00":
2605 continue
2606
2607 meta_length = ord(meta_byte) * 16
2608 meta_data = await reader.readexactly(meta_length)
2609 self._parse_icy_metadata(meta_data, streamdetails)
2610
2611 except asyncio.exceptions.IncompleteReadError:
2612 # End of stream
2613 break
2614
2615 finally:
2616 writer.close()
2617 await writer.wait_closed()
2618
2619 # --- Private methods ---
2620
2621 def _notify_provider_streamed(
2622 self, streamdetails: StreamDetails, finished: bool, seconds_streamed: float
2623 ) -> None:
2624 """Report a (mostly) streamed item back to the provider that owns it."""
2625 if not finished and seconds_streamed < 90:
2626 return
2627 provider = self.mass.get_provider(streamdetails.provider)
2628 # plugin providers serve playable items too, but on_streamed is MusicProvider-only
2629 if provider is None or provider.type != ProviderType.MUSIC:
2630 return
2631 music_prov = cast("MusicProvider", provider)
2632 self.mass.create_task(music_prov.on_streamed(streamdetails))
2633
2634 def _get_volume_normalization_preference(
2635 self, streamdetails: StreamDetails
2636 ) -> VolumeNormalizationMode:
2637 """Return the configured normalization preference for the stream's media type."""
2638 conf_key = (
2639 CONF_VOLUME_NORMALIZATION_RADIO
2640 if streamdetails.media_type == MediaType.RADIO
2641 else CONF_VOLUME_NORMALIZATION_TRACKS
2642 )
2643 return VolumeNormalizationMode(
2644 self.mass.streams.get_config_value(conf_key, return_type=str)
2645 )
2646
2647 def _update_radio_stream_metadata(
2648 self,
2649 streamdetails: StreamDetails,
2650 artist: str | None,
2651 title: str,
2652 image_url: str | None = None,
2653 album: str | None = None,
2654 ) -> None:
2655 """
2656 Update radio stream metadata and trigger artwork lookup.
2657
2658 :param streamdetails: The stream details to update.
2659 :param artist: Artist name (will be normalized).
2660 :param title: Track title (will be cleaned for display).
2661 :param image_url: Optional image URL from stream metadata.
2662 :param album: Optional album name.
2663 """
2664 station_image_url = image_url or self.mass.metadata.get_radio_stream_station_image(
2665 streamdetails
2666 )
2667 artist_normalized = (
2668 self.mass.metadata.normalize_radio_artist_name(artist) if artist else None
2669 )
2670 display_title, _ = parse_title_and_version(title, strip_for_display=True)
2671
2672 streamdetails.stream_metadata = StreamMetadata(
2673 title=display_title,
2674 artist=artist_normalized,
2675 album=album,
2676 image_url=station_image_url,
2677 )
2678 streamdetails.stream_metadata_last_updated = time.time()
2679 if streamdetails.queue_id:
2680 self.mass.player_queues.signal_update(streamdetails.queue_id)
2681
2682 # Fetch artwork in background (track, album then artist)
2683 if artist and title and not image_url:
2684 self.mass.call_later(
2685 0.2,
2686 self.mass.metadata.update_radio_stream_artwork,
2687 streamdetails,
2688 task_id=f"update_radio_artwork_{streamdetails.queue_id}",
2689 )
2690
2691 async def _cache_radio_result(
2692 self,
2693 url: str,
2694 stream_type: StreamType,
2695 resolved_url: str | None = None,
2696 ) -> tuple[str, StreamType]:
2697 """Cache and return a radio stream resolution result."""
2698 result = (resolved_url or url, stream_type)
2699 await self.mass.cache.set(
2700 url,
2701 result,
2702 expiration=3600 * 3,
2703 provider=CACHE_PROVIDER,
2704 category=CACHE_CATEGORY_RESOLVED_RADIO_URL,
2705 )
2706 return result
2707
2708 async def _handle_client_error_for_radio_stream(
2709 self, url: str, err: aiohttp.ClientError, fallback_stream_type: StreamType
2710 ) -> tuple[str, StreamType]:
2711 """Handle aiohttp client errors during radio stream resolution."""
2712 # Prefer the final post-redirect URL: aiohttp follows redirects before raising,
2713 # but the original url may just point at a redirector rather than the ICY endpoint.
2714 request_info = getattr(err, "request_info", None)
2715 validate_url = str(request_info.url) if request_info is not None else url
2716
2717 # Check if this is a Shoutcast/ICY response that aiohttp can't parse
2718 if isinstance(err, aiohttp.ClientResponseError) and "ICY" in str(err).upper():
2719 self.logger.debug(
2720 "ICY response detected for %s, validating Shoutcast stream", validate_url
2721 )
2722 if await self._validate_shoutcast_stream(validate_url):
2723 return await self._cache_radio_result(
2724 url, StreamType.SHOUTCAST, resolved_url=validate_url
2725 )
2726 self.logger.warning(
2727 "ICY response detected but Shoutcast validation failed for %s", validate_url
2728 )
2729 return await self._cache_radio_result(
2730 url, fallback_stream_type, resolved_url=validate_url
2731 )
2732
2733 # Other aiohttp errors - might still be Shoutcast, check it
2734 self.logger.debug("aiohttp error for %s, checking if legacy Shoutcast stream", validate_url)
2735 if await self._validate_shoutcast_stream(validate_url):
2736 return await self._cache_radio_result(
2737 url, StreamType.SHOUTCAST, resolved_url=validate_url
2738 )
2739
2740 # Unknown error - still try to stream
2741 self.logger.warning(
2742 "Failed to parse radio URL %s: %s - attempting direct stream", validate_url, str(err)
2743 )
2744 return await self._cache_radio_result(url, fallback_stream_type, resolved_url=validate_url)
2745
2746 async def _get_audio_buffer(
2747 self,
2748 queue_item: QueueItem,
2749 seek_position_ms: int,
2750 reason: str,
2751 capacity_wait_timeout: float,
2752 allow_provider_match: bool,
2753 ) -> AudioBuffer:
2754 """
2755 Create or reuse a ready AudioBuffer within one queue-item preparation lock.
2756
2757 :param queue_item: Queue item whose source should be buffered.
2758 :param seek_position_ms: Position in milliseconds to start from.
2759 :param reason: Caller context for logging.
2760 :param capacity_wait_timeout: Total seconds to spend waiting for source capacity.
2761 :param allow_provider_match: Whether an on-demand cross-provider match may widen
2762 the candidates when all are saturated.
2763 """
2764 loop = asyncio.get_running_loop()
2765 # the playback intent lives on the details we start from; keep it across a reselection
2766 initial_streamdetails = queue_item.streamdetails
2767 seek_position = (
2768 int(initial_streamdetails.seek_position)
2769 if initial_streamdetails
2770 else seek_position_ms // 1000
2771 )
2772 fade_in = bool(initial_streamdetails and initial_streamdetails.fade_in)
2773 prefer_album_loudness = bool(
2774 initial_streamdetails and initial_streamdetails.prefer_album_loudness
2775 )
2776 all_candidate_instances = {
2777 provider.instance_id
2778 for mapping in (
2779 queue_item.media_item.provider_mappings if queue_item.media_item else ()
2780 )
2781 if mapping.available
2782 for provider in self._get_mapping_providers(mapping)
2783 }
2784 if initial_streamdetails is not None:
2785 all_candidate_instances.add(initial_streamdetails.provider)
2786 # a track may also exist on streaming providers it has no mapping for yet; such a
2787 # match is only searched once, and only when every known candidate is saturated
2788 match_pending = (
2789 allow_provider_match
2790 and isinstance(queue_item.media_item, Track)
2791 and self._has_alternative_match_providers(queue_item.media_item)
2792 )
2793
2794 deadline = loop.time() + capacity_wait_timeout
2795 busy_instances: set[str] = set()
2796 final_pass = False
2797 last_capacity_error: ProviderStreamLimitError | None = None
2798 last_failed_streamdetails: StreamDetails | None = None
2799 while True:
2800 if queue_item.streamdetails is None or (
2801 queue_item.streamdetails.provider in busy_instances and not final_pass
2802 ):
2803 try:
2804 queue_item.streamdetails = await self.get_stream_details(
2805 queue_item,
2806 seek_position=seek_position,
2807 fade_in=fade_in,
2808 prefer_album_loudness=prefer_album_loudness,
2809 excluded_provider_instances=busy_instances,
2810 )
2811 except (AudioError, MediaNotFoundError) as err:
2812 if last_capacity_error is None:
2813 raise
2814 if final_pass:
2815 # capacity was the root cause, surface the typed (actionable) error
2816 raise last_capacity_error from err
2817 # no usable alternative mapping: restore the capacity-blocked details
2818 # and spend the remaining budget blocking on that provider's slot
2819 final_pass = True
2820 continue
2821 finally:
2822 if queue_item.streamdetails is None:
2823 # never leave the queue item without streamdetails on any exit,
2824 # including a cancellation or a non-audio provider failure
2825 queue_item.streamdetails = last_failed_streamdetails
2826 streamdetails = queue_item.streamdetails
2827 assert streamdetails is not None # for type checking
2828 remaining = max(deadline - loop.time(), 0)
2829 alternatives_left = bool(
2830 all_candidate_instances - busy_instances - {streamdetails.provider}
2831 )
2832 # probe (0s) whenever a reselection can still follow: a free slot is still
2833 # acquired instantly, while a busy one fails fast instead of spending the
2834 # whole budget on this candidate. Block only on the last resort.
2835 source_wait = (
2836 0.0
2837 if (not final_pass and (alternatives_left or busy_instances or match_pending))
2838 else remaining
2839 )
2840 try:
2841 return await AudioBuffer.get_buffer(
2842 mass=self.mass,
2843 streamdetails=streamdetails,
2844 seek_position_ms=seek_position_ms,
2845 wait_ready=True,
2846 reason=reason,
2847 source_wait_timeout=source_wait,
2848 )
2849 except ProviderStreamLimitError as err:
2850 last_capacity_error = err
2851 last_failed_streamdetails = streamdetails
2852 busy_instances.add(err.provider_instance)
2853 if final_pass or loop.time() >= deadline:
2854 raise
2855 if all_candidate_instances.issubset(busy_instances):
2856 discovered: set[str] = set()
2857 if match_pending:
2858 match_pending = False
2859 try:
2860 discovered = await self._discover_alternative_provider_mappings(
2861 queue_item, busy_instances, max(deadline - loop.time(), 0)
2862 )
2863 except Exception as err:
2864 # discovery is best-effort: any failure falls back to the
2865 # final blocking wait instead of replacing the typed error
2866 self.logger.warning(
2867 "Alternative provider search for %s failed: %s",
2868 queue_item.name,
2869 err,
2870 )
2871 if discovered:
2872 all_candidate_instances.update(discovered)
2873 else:
2874 # every candidate is saturated: one last blocking wait on the best one
2875 busy_instances.clear()
2876 final_pass = True
2877 queue_item.streamdetails = None
2878 except AudioError:
2879 if last_capacity_error is None or final_pass:
2880 raise
2881 # a broken alternate must not turn a transient capacity miss into a hard
2882 # failure: restore the blocked details and spend the rest of the budget there
2883 queue_item.streamdetails = last_failed_streamdetails
2884 final_pass = True
2885
2886 def _get_streamdetail_candidates(
2887 self,
2888 provider_mappings: Iterable[ProviderMapping],
2889 preferred_providers: list[str],
2890 excluded_provider_instances: set[str],
2891 ) -> list[tuple[ProviderMapping, Provider]]:
2892 """
2893 Return mapping candidates in steering, quality, and instance-fallback order.
2894
2895 :param provider_mappings: Mappings attached to the media item.
2896 :param preferred_providers: Provider instances tried before widening to the rest.
2897 :param excluded_provider_instances: Provider instances unavailable to this attempt.
2898 :return: Ordered provider mapping candidates.
2899 """
2900 ordered_mappings = sorted(
2901 provider_mappings, key=lambda mapping: mapping.quality or 0, reverse=True
2902 )
2903 preferred_candidates: list[tuple[ProviderMapping, Provider]] = []
2904 fallback_candidates: list[tuple[ProviderMapping, Provider]] = []
2905 seen_candidates: set[tuple[str, str]] = set()
2906 for mapping in ordered_mappings:
2907 if not mapping.available:
2908 self.logger.debug("Skipping unavailable %s", mapping)
2909 continue
2910 for provider in self._get_mapping_providers(mapping):
2911 candidate_id = (provider.instance_id, mapping.item_id)
2912 if (
2913 candidate_id in seen_candidates
2914 or provider.instance_id in excluded_provider_instances
2915 ):
2916 continue
2917 seen_candidates.add(candidate_id)
2918 candidate = (mapping, provider)
2919 if provider.instance_id in preferred_providers:
2920 preferred_candidates.append(candidate)
2921 else:
2922 fallback_candidates.append(candidate)
2923 return [*preferred_candidates, *fallback_candidates]
2924
2925 def _get_mapping_providers(self, mapping: ProviderMapping) -> list[Provider]:
2926 """
2927 Return the mapped provider followed by compatible instances of its streaming catalog.
2928
2929 :param mapping: Provider mapping whose item ID will be requested.
2930 :return: Loaded provider instances that can resolve the mapping.
2931 """
2932 providers: list[Provider] = []
2933 if (
2934 primary_provider := self.mass.get_provider(
2935 mapping.provider_instance, return_unavailable=True
2936 )
2937 ) and primary_provider.available:
2938 providers.append(primary_provider)
2939 # another account of the same streaming catalog serves the same item ID,
2940 # so it can stand in when the mapped instance can not
2941 for provider in self.mass.providers:
2942 if (
2943 not isinstance(provider, MusicProvider)
2944 or not provider.available
2945 or not provider.is_streaming_provider
2946 or provider.domain != mapping.provider_domain
2947 or provider in providers
2948 ):
2949 continue
2950 providers.append(provider)
2951 if not providers:
2952 self.logger.debug("Skipping %s - provider not available", mapping)
2953 return providers
2954
2955 def _is_match_candidate_provider(
2956 self, provider: MusicProvider, known_domains: set[str]
2957 ) -> bool:
2958 """
2959 Return whether a provider is eligible to search a track match on.
2960
2961 :param provider: Music provider to check.
2962 :param known_domains: Provider domains the track already has mappings for.
2963 """
2964 return (
2965 provider.available
2966 and provider.is_streaming_provider
2967 and ProviderFeature.SEARCH in provider.supported_features
2968 and provider.domain not in known_domains
2969 and MediaType.TRACK in provider.supported_media_types
2970 )
2971
2972 def _has_alternative_match_providers(self, media_item: Track) -> bool:
2973 """
2974 Return whether any configured streaming provider could carry an unmapped match.
2975
2976 :param media_item: Track whose existing mappings define the known provider domains.
2977 """
2978 known_domains = {mapping.provider_domain for mapping in media_item.provider_mappings}
2979 return any(
2980 self._is_match_candidate_provider(provider, known_domains)
2981 for provider in self.mass.music.providers
2982 )
2983
2984 async def _discover_alternative_provider_mappings(
2985 self, queue_item: QueueItem, busy_instances: set[str], remaining: float
2986 ) -> set[str]:
2987 """
2988 Search other streaming providers for the queue item's track and widen its mappings.
2989
2990 A found mapping is added to the media item (and persisted for library items) so the
2991 capacity reselection can continue on the discovered provider.
2992
2993 :param queue_item: Queue item whose track should be matched on another provider.
2994 :param busy_instances: Provider instances already known to be saturated.
2995 :param remaining: Seconds left of the caller's capacity budget.
2996 :return: Provider instances able to serve the discovered mappings.
2997 """
2998 media_item = queue_item.media_item
2999 if not isinstance(media_item, Track):
3000 return set()
3001 known_domains = {mapping.provider_domain for mapping in media_item.provider_mappings}
3002 eligible = [
3003 provider
3004 for provider in self.mass.music.providers
3005 if self._is_match_candidate_provider(provider, known_domains)
3006 and provider.instance_id not in busy_instances
3007 and provider.has_available_stream_slot
3008 ]
3009 if not eligible:
3010 return set()
3011 # mirror the playback user's provider steering for the search order
3012 if (
3013 (pq_data := self.mass.player_queues.queue_data_or_none(queue_item.queue_id))
3014 and pq_data.userid
3015 and (playback_user := await self.mass.webserver.auth.get_user(pq_data.userid))
3016 and playback_user.provider_filter
3017 ):
3018 preferred = set(playback_user.provider_filter)
3019 eligible.sort(key=lambda provider: provider.instance_id not in preferred)
3020 # one instance per domain: a found mapping widens to sibling instances anyway
3021 candidates: list[MusicProvider] = []
3022 for provider in eligible:
3023 if provider.domain in known_domains:
3024 continue
3025 known_domains.add(provider.domain)
3026 candidates.append(provider)
3027 # the track's own album is free, sufficient evidence for the strict compare and
3028 # avoids match_provider's multi-provider album lookup on every call
3029 ref_albums = [media_item.album] if isinstance(media_item.album, Album) else []
3030 matches: list[ProviderMapping] = []
3031 try:
3032 async with asyncio.timeout(min(STREAM_SLOT_MATCH_TIMEOUT, remaining)):
3033 for provider in candidates:
3034 # one failing provider must not end the search on the others
3035 try:
3036 matches = await self.mass.music.tracks.match_provider(
3037 media_item, provider, strict=True, ref_albums=ref_albums
3038 )
3039 except Exception as err:
3040 self.logger.debug("Searching a match on %s failed: %s", provider.name, err)
3041 continue
3042 if matches:
3043 break
3044 except TimeoutError:
3045 self.logger.debug("Searching an alternative provider for %s timed out", media_item.name)
3046 if not matches:
3047 return set()
3048 media_item.provider_mappings.update(matches)
3049 if media_item.provider == "library":
3050 # persist in the background so future plays have the mapping ahead of time;
3051 # cancellation of this playback must never interrupt the library write
3052 self.mass.create_task(
3053 self.mass.music.tracks.add_provider_mappings(media_item.item_id, matches)
3054 )
3055 self.logger.info(
3056 "All known sources for %s are at their stream limit, "
3057 "using a matching track found on %s",
3058 media_item.name,
3059 matches[0].provider_domain,
3060 )
3061 return {
3062 provider.instance_id
3063 for mapping in matches
3064 for provider in self._get_mapping_providers(mapping)
3065 }
3066
3067 async def _request_streamdetails(
3068 self,
3069 candidates: Iterable[tuple[ProviderMapping, Provider]],
3070 media_type: MediaType,
3071 ) -> StreamDetails | None:
3072 """
3073 Request stream details from ordered provider mapping candidates.
3074
3075 :param candidates: Candidates in mapping and compatible-instance order.
3076 :param media_type: Media type requested from each provider.
3077 :return: The first resolved stream details, or None when every candidate failed.
3078 :raises AudioError: The last (actionable) audio error when no candidate resolved.
3079 """
3080 last_audio_error: AudioError | None = None
3081 for mapping, provider in candidates:
3082 # music and plugin providers share this signature, so either type can own the item
3083 token = BYPASS_THROTTLER.set(True)
3084 try:
3085 stream_prov = cast("MusicProvider | PluginProvider", provider)
3086 return await stream_prov.get_stream_details(mapping.item_id, media_type)
3087 except AudioError as err:
3088 # remember the last one so its (actionable) message can be re-raised
3089 last_audio_error = err
3090 self.logger.warning("%s", err)
3091 except MusicAssistantError as err:
3092 self.logger.warning("%s", err)
3093 finally:
3094 BYPASS_THROTTLER.reset(token)
3095 if last_audio_error is not None:
3096 raise last_audio_error
3097 return None
3098
3099 async def _get_media_stream(
3100 self,
3101 streamdetails: StreamDetails,
3102 pcm_format: AudioFormat,
3103 seek_position: int,
3104 filter_params: list[str] | None,
3105 chunk_seconds: float,
3106 ) -> AsyncGenerator[bytes]:
3107 """
3108 Stream one provider source as raw PCM.
3109
3110 :param streamdetails: Details of the stream to fetch.
3111 :param pcm_format: Target PCM format the consumer expects.
3112 :param seek_position: Requested seek offset in seconds.
3113 :param filter_params: Optional ffmpeg filter expressions.
3114 :param chunk_seconds: Size of each yielded chunk in seconds of audio.
3115 """
3116 mass = self.mass
3117 logger = self.logger.getChild("media_stream")
3118 logger.log(VERBOSE_LOG_LEVEL, "Starting media stream for %s", streamdetails.uri)
3119 # copy: the args below are appended per call, while the StreamDetails is cached on
3120 # the queue item and reused across calls (retry, seek, background analysis)
3121 extra_input_args = list(streamdetails.extra_input_args or [])
3122 # the resolver below zeroes out seek_position where the seek is delegated to the
3123 # source itself, so keep the requested position for the duration writeback
3124 requested_seek_position = seek_position
3125
3126 # work out audio source for these streamdetails
3127 audio_source, seek_position, extra_input_args = await self._resolve_media_stream_source(
3128 streamdetails, seek_position, extra_input_args
3129 )
3130
3131 # pace ffmpeg at native rate for live sources; the producer (e.g.
3132 # librespot's pipe backend) may otherwise write faster than realtime.
3133 # The initial burst grants a small bounded read-ahead so downstream
3134 # jitter does not immediately underrun the player. Providers that need
3135 # different pacing can pass their own -re/-readrate args to override.
3136 if (
3137 streamdetails.media_type == MediaType.AUDIO_SOURCE
3138 and "-re" not in extra_input_args
3139 and "-readrate" not in extra_input_args
3140 ):
3141 extra_input_args += ["-readrate", "1", "-readrate_initial_burst", "0.5"]
3142
3143 # handle seek support
3144 if seek_position and streamdetails.duration and streamdetails.allow_seek:
3145 extra_input_args += ["-ss", str(int(seek_position))]
3146
3147 bytes_sent = 0
3148 finished = False
3149 cancelled = False
3150 first_chunk_received = False
3151 ffmpeg_loglevel = "debug" if self.logger.isEnabledFor(VERBOSE_LOG_LEVEL) else "info"
3152 # When a provider hands us already-decoded audio (e.g. Spotify Connect /
3153 # AirPlay receivers piping PCM after their own decode), audio_format is
3154 # the original source format meant for display while decoded_audio_format
3155 # is what ffmpeg actually needs to read off the wire.
3156 ffmpeg_input_format = streamdetails.decoded_audio_format or streamdetails.audio_format
3157 ffmpeg_proc = FFMpeg(
3158 audio_input=audio_source,
3159 input_format=ffmpeg_input_format,
3160 output_format=pcm_format,
3161 filter_params=filter_params,
3162 extra_input_args=extra_input_args,
3163 collect_log_history=True,
3164 loglevel=ffmpeg_loglevel,
3165 )
3166
3167 try:
3168 await ffmpeg_proc.start()
3169 assert ffmpeg_proc.proc is not None # for type checking
3170 if logger.isEnabledFor(VERBOSE_LOG_LEVEL):
3171 logger.log(
3172 VERBOSE_LOG_LEVEL,
3173 "Started media stream for %s - using streamtype: %s "
3174 "- pcm format: %s - ffmpeg PID: %s",
3175 streamdetails.uri,
3176 streamdetails.stream_type,
3177 pcm_format.content_type.value,
3178 ffmpeg_proc.proc.pid,
3179 )
3180 else:
3181 logger.debug(
3182 "Started media stream for %s - using streamtype: %s",
3183 streamdetails.uri,
3184 streamdetails.stream_type,
3185 )
3186 stream_start = mass.loop.time()
3187 chunk_size = calculate_content_length(pcm_format, chunk_seconds)
3188 chunk_iter = ffmpeg_proc.iter_chunked(chunk_size)
3189 while True:
3190 # Time the read, not the yield: catches a stalled source, ignores backpressure.
3191 read_timeout = (
3192 STREAM_START_TIMEOUT if not first_chunk_received else STREAM_STALL_TIMEOUT
3193 )
3194 try:
3195 async with asyncio.timeout(read_timeout):
3196 chunk = await anext(chunk_iter)
3197 except StopAsyncIteration:
3198 break
3199 except TimeoutError as err:
3200 raise AudioError(f"Source stalled: no audio for {read_timeout}s") from err
3201 if not first_chunk_received:
3202 # At this point ffmpeg has started and should now know the codec used
3203 # for encoding the audio.
3204 # Note: ffmpeg_proc.input_format is the same object as
3205 # ffmpeg_input_format, so sample_rate / bit_depth / bit_rate
3206 # parsed from the ffmpeg log already live on streamdetails too.
3207 first_chunk_received = True
3208 # Skip the codec_type writeback when the provider declared a
3209 # decoded format: audio_format already holds the authoritative
3210 # source codec and the probed value would just be the
3211 # post-decode wire format (e.g. PCM for Spotify Connect).
3212 if streamdetails.decoded_audio_format is None:
3213 streamdetails.audio_format.codec_type = ffmpeg_proc.input_format.codec_type
3214 # Some providers omit (or report 0 for) the item duration; ffmpeg can
3215 # usually probe it from the source. Only apply when missing so we
3216 # don't clobber an accurate provider value with a rounded one.
3217 if ffmpeg_proc.parsed_duration is not None and not streamdetails.duration:
3218 streamdetails.duration = ffmpeg_proc.parsed_duration
3219 logger.debug(
3220 "First chunk received after %.2f seconds (codec detected: %s)",
3221 mass.loop.time() - stream_start,
3222 ffmpeg_proc.input_format.codec_type,
3223 )
3224 yield chunk
3225 bytes_sent += len(chunk)
3226
3227 # end of audio/track reached
3228 logger.debug("End of media stream reached for %s", streamdetails.uri)
3229 # wait until stderr also completed reading
3230 await ffmpeg_proc.wait_with_timeout(5)
3231 logger.log(
3232 VERBOSE_LOG_LEVEL,
3233 "FFmpeg process ended with return code %s for %s",
3234 ffmpeg_proc.returncode,
3235 streamdetails.uri,
3236 )
3237 # a nested source raises through the stdin feeder, where ffmpeg's own exit
3238 # would otherwise flatten it into a generic AudioError
3239 if isinstance(ffmpeg_proc.stdin_feeder_exception, ProviderStreamLimitError):
3240 raise ffmpeg_proc.stdin_feeder_exception
3241 if ffmpeg_proc.returncode not in (0, None):
3242 log_trail = "\n".join(list(ffmpeg_proc.log_history)[-5:])
3243 raise AudioError(f"FFMpeg exited with code {ffmpeg_proc.returncode}: {log_trail}")
3244 if bytes_sent == 0:
3245 # edge case: no audio data was received at all
3246 raise AudioError("No audio was received")
3247 finished = True
3248 except (Exception, GeneratorExit, asyncio.CancelledError) as err:
3249 if isinstance(err, asyncio.CancelledError | GeneratorExit):
3250 # we were cancelled, just raise
3251 cancelled = True
3252 raise
3253 if isinstance(ffmpeg_proc.stdin_feeder_exception, ProviderStreamLimitError):
3254 raise ffmpeg_proc.stdin_feeder_exception
3255 # dump the last 10 lines of the log in case of an unclean exit
3256 logger.warning("\n".join(list(ffmpeg_proc.log_history)[-10:]))
3257 raise AudioError(f"Error while streaming: {err}") from err
3258 finally:
3259 # always ensure close is called which also handles all cleanup
3260 await ffmpeg_proc.close()
3261 # determine how many seconds we've received
3262 # for pcm output we can calculate this easily
3263 seconds_received = bytes_sent / pcm_format.pcm_sample_size if bytes_sent else 0
3264 # store accurate duration, but only for a playthrough from the very start:
3265 # a seeked stream yields the remaining audio, not the item's full length
3266 if finished and not requested_seek_position and seconds_received:
3267 streamdetails.duration = int(seconds_received)
3268
3269 logger.log(
3270 VERBOSE_LOG_LEVEL,
3271 "stream %s (with code %s) for %s",
3272 "cancelled" if cancelled else "finished" if finished else "aborted",
3273 ffmpeg_proc.returncode,
3274 streamdetails.uri,
3275 )
3276
3277 def _select_buffered_crossfade(
3278 self,
3279 streamdetails: StreamDetails,
3280 crossfade_mode: CrossfadeMode,
3281 standard_crossfade_duration: int,
3282 playback_speed: float = 1.0,
3283 ) -> tuple[CrossfadeMode, float]:
3284 """
3285 Select a crossfade that can be completed from resident incoming PCM.
3286
3287 :param streamdetails: Incoming track stream details.
3288 :param crossfade_mode: Requested crossfade mode.
3289 :param standard_crossfade_duration: Configured standard overlap in seconds.
3290 :param playback_speed: Incoming track playback-speed multiplier.
3291 :return: Effective mode and resident fade-in duration in seconds.
3292 """
3293 audio_buffer = streamdetails.buffer
3294 if (
3295 crossfade_mode == CrossfadeMode.DISABLED
3296 or playback_speed <= 0
3297 or audio_buffer is None
3298 or audio_buffer.has_error
3299 or not audio_buffer.is_valid()
3300 ):
3301 return CrossfadeMode.DISABLED, 0
3302
3303 available_seconds = audio_buffer.duration_available / playback_speed
3304 if (
3305 crossfade_mode == CrossfadeMode.SMART_CROSSFADE
3306 and audio_buffer.ready.is_set()
3307 and available_seconds >= MIN_EFFECTIVE_FADE_BUFFER
3308 ):
3309 return crossfade_mode, min(SMART_CROSSFADE_DURATION, available_seconds)
3310 if (
3311 crossfade_mode == CrossfadeMode.STANDARD_CROSSFADE
3312 and standard_crossfade_duration >= MIN_CROSSFADE_FALLBACK_DURATION
3313 and audio_buffer.ready.is_set()
3314 and available_seconds >= standard_crossfade_duration
3315 ):
3316 return crossfade_mode, standard_crossfade_duration
3317
3318 fallback_duration = min(standard_crossfade_duration, available_seconds)
3319 if fallback_duration < MIN_CROSSFADE_FALLBACK_DURATION:
3320 return CrossfadeMode.DISABLED, 0
3321 self.logger.debug(
3322 "Using %s second standard crossfade for %s from resident audio",
3323 fallback_duration,
3324 streamdetails.uri,
3325 )
3326 return CrossfadeMode.STANDARD_CROSSFADE, fallback_duration
3327
3328 async def _resolve_media_stream_source(
3329 self,
3330 streamdetails: StreamDetails,
3331 seek_position: int,
3332 extra_input_args: list[str],
3333 ) -> tuple[str | AsyncGenerator[bytes], int, list[str]]:
3334 """
3335 Resolve the input consumed by ffmpeg for the given stream details.
3336
3337 :param streamdetails: Details of the stream to fetch.
3338 :param seek_position: Requested seek offset in seconds.
3339 :param extra_input_args: Provider-supplied ffmpeg input arguments.
3340 :return: The ffmpeg input, the remaining seek offset and the ffmpeg input arguments.
3341 """
3342 stream_type = streamdetails.stream_type
3343 if stream_type == StreamType.CUSTOM:
3344 # MusicProvider and PluginProvider both expose get_audio_stream with the same shape.
3345 # Pin the exact instance: a domain fallback would stream from a sibling account
3346 # while the source-stream slot is charged to the instance that issued the details.
3347 provider = self.mass.get_provider(streamdetails.provider, return_unavailable=True)
3348 if provider is None or not provider.available:
3349 raise ProviderUnavailableError(
3350 f"Provider {streamdetails.provider} for stream is no longer available"
3351 )
3352 provider = cast("MusicProvider | PluginProvider", provider)
3353 audio_source = provider.get_audio_stream(
3354 streamdetails, seek_position=seek_position if streamdetails.can_seek else 0
3355 )
3356 return audio_source, 0 if streamdetails.can_seek else seek_position, extra_input_args
3357 if stream_type == StreamType.ICY:
3358 assert streamdetails.path is not None
3359 assert isinstance(streamdetails.path, (str, list))
3360 audio_source = self.get_reconnecting_icy_radio_stream(streamdetails.path, streamdetails)
3361 return audio_source, 0, extra_input_args
3362 if stream_type == StreamType.SHOUTCAST:
3363 assert isinstance(streamdetails.path, str)
3364 return self.get_shoutcast_stream(streamdetails.path, streamdetails), 0, extra_input_args
3365 if stream_type == StreamType.IN_BAND:
3366 assert isinstance(streamdetails.path, str) # for type checking
3367
3368 # For IN_BAND (OGG/Opus) radio streams, use chained OGG handler.
3369 # This handles the chained OGG format by stitching logical bitstreams together
3370 # so FFmpeg sees a single continuous stream. Metadata is extracted in-band.
3371 audio_source = get_chained_ogg_stream(
3372 self.mass,
3373 streamdetails.path,
3374 metadata_callback=partial(self._handle_inband_metadata, streamdetails),
3375 )
3376 # seeking not possible on radio streams
3377 return audio_source, 0, extra_input_args
3378 if stream_type == StreamType.HLS:
3379 assert isinstance(streamdetails.path, str) # for type checking
3380 substream = await self.get_hls_substream(streamdetails.path)
3381 if streamdetails.media_type == MediaType.RADIO:
3382 # HLS streams (especially the BBC) struggle when they're played directly
3383 # with ffmpeg, where they just stop after some minutes,
3384 # so we tell ffmpeg to loop around in this case.
3385 extra_input_args += ["-stream_loop", "-1", "-re"]
3386 return substream.path, seek_position, extra_input_args
3387
3388 # all other stream types (HTTP, FILE, etc)
3389 if stream_type == StreamType.ENCRYPTED_HTTP:
3390 assert streamdetails.decryption_key is not None # for type checking
3391 extra_input_args += ["-decryption_key", streamdetails.decryption_key]
3392 if isinstance(streamdetails.path, list):
3393 # multi part stream, which handles the seek itself
3394 return self.get_multi_file_stream(streamdetails, seek_position), 0, extra_input_args
3395 # regular single file/url stream
3396 assert isinstance(streamdetails.path, str) # for type checking
3397 return streamdetails.path, seek_position, extra_input_args
3398
3399 async def _iter_audio_source_pcm(
3400 self,
3401 streamdetails: StreamDetails,
3402 pcm_format: AudioFormat,
3403 ) -> AsyncGenerator[bytes]:
3404 """Yield PCM for an AudioSource, bypassing ffmpeg when formats match."""
3405 if _pcm_formats_match(streamdetails.audio_format, pcm_format):
3406 source_gen = self._open_audio_source_generator(streamdetails)
3407 async for chunk in realtime_pcm_pacer(source_gen, pcm_format):
3408 yield chunk
3409 return
3410 # format mismatch â fall back to ffmpeg for resampling (still small chunks)
3411 async for chunk in self.get_media_stream(
3412 streamdetails=streamdetails,
3413 pcm_format=pcm_format,
3414 filter_params=None,
3415 chunk_seconds=AUDIO_SOURCE_CHUNK_SECONDS,
3416 ):
3417 yield chunk
3418
3419 def _open_audio_source_generator(self, streamdetails: StreamDetails) -> AsyncGenerator[bytes]:
3420 """Open the raw PCM generator for an AudioSource (CUSTOM or NAMED_PIPE)."""
3421 if streamdetails.stream_type == StreamType.CUSTOM:
3422 # pin the exact instance, see _resolve_media_stream_source
3423 provider = self.mass.get_provider(streamdetails.provider, return_unavailable=True)
3424 if provider is None or not provider.available:
3425 raise ProviderUnavailableError(
3426 f"Provider {streamdetails.provider} for stream is no longer available"
3427 )
3428 provider = cast("MusicProvider | PluginProvider", provider)
3429 return provider.get_audio_stream(streamdetails)
3430 if streamdetails.stream_type == StreamType.NAMED_PIPE:
3431 assert isinstance(streamdetails.path, str) # for type checking
3432 return read_named_pipe(streamdetails.path)
3433 raise AudioError(f"Unsupported stream_type {streamdetails.stream_type} for AudioSource")
3434
3435 def _handle_inband_metadata(
3436 self, streamdetails: StreamDetails, metadata: dict[str, str]
3437 ) -> None:
3438 """Handle metadata extracted from a chained Ogg stream."""
3439 title = metadata.get("title", "")
3440 artist = metadata.get("artist", "")
3441 album = metadata.get("album", "")
3442 if not artist and " - " in title:
3443 artist, title = title.split(" - ", 1)
3444 if not (title or artist):
3445 return
3446
3447 stream_title = f"{artist} - {title}" if artist and title else title or artist
3448 cleaned_title = clean_stream_title(stream_title)
3449 if not cleaned_title:
3450 return
3451 if self._record_inband_stream_title(streamdetails, cleaned_title):
3452 return
3453 if cleaned_title != streamdetails.stream_title:
3454 self.logger.log(VERBOSE_LOG_LEVEL, "In-band metadata: %s", cleaned_title)
3455 streamdetails.stream_title = cleaned_title
3456 self._update_radio_stream_metadata(
3457 streamdetails,
3458 artist=artist or None,
3459 title=title or cleaned_title,
3460 album=album or None,
3461 )
3462
3463 def _record_inband_stream_title(self, streamdetails: StreamDetails, cleaned_title: str) -> bool:
3464 """
3465 Record an in-band stream title for provider-owned metadata, if applicable.
3466
3467 When a provider opts into owning stream_metadata (and stream_title is only
3468 a derived view of it), writing either from the stream reader would fight the
3469 provider. The cleaned in-band title is recorded on StreamDetails.data instead,
3470 as the identity signal for the provider callback.
3471
3472 :param streamdetails: StreamDetails carrying the stream.
3473 :param cleaned_title: Cleaned in-band stream title.
3474 :returns: True when recorded (the caller must not write stream metadata);
3475 False when no provider callback exists and normal handling applies.
3476 """
3477 if (
3478 streamdetails.stream_metadata_update_callback is None
3479 or streamdetails.data is None
3480 or not streamdetails.data.get(STREAMDETAILS_INBAND_TITLE_HANDOFF_KEY)
3481 ):
3482 return False
3483 if streamdetails.data.get(STREAMDETAILS_INBAND_TITLE_KEY) != cleaned_title:
3484 # occupancy approximates how far this detection leads audible playback
3485 buffer = streamdetails.buffer
3486 self.logger.debug(
3487 "In-band stream title: %s (buffer occupancy: %ss)",
3488 cleaned_title,
3489 buffer.size_seconds if buffer is not None else "unknown",
3490 )
3491 streamdetails.data[STREAMDETAILS_INBAND_TITLE_KEY] = cleaned_title
3492 return True
3493
3494 def _parse_icy_metadata(self, meta_data: bytes, streamdetails: StreamDetails) -> None:
3495 """
3496 Parse ICY metadata and update streamdetails.
3497
3498 Sets the cleaned stream title and, when the title parses as "Artist - Track",
3499 triggers a radio-artwork metadata update.
3500
3501 :param meta_data: Raw metadata bytes from an ICY stream chunk.
3502 :param streamdetails: StreamDetails to update with parsed title and metadata.
3503 """
3504 if not meta_data:
3505 return
3506
3507 meta_data = meta_data.rstrip(b"\0")
3508 # Match StreamTitle, handling apostrophes in titles
3509 stream_title_re = re.search(rb"StreamTitle='(.*?)';", meta_data)
3510
3511 if not stream_title_re:
3512 self.logger.log(
3513 VERBOSE_LOG_LEVEL,
3514 "ICY metadata does not contain StreamTitle field. Raw: %s",
3515 meta_data.decode("utf-8", errors="replace")[:200],
3516 )
3517 return
3518
3519 try:
3520 # in 99% of the cases the stream title is utf-8 encoded
3521 stream_title = stream_title_re.group(1).decode("utf-8")
3522 except UnicodeDecodeError:
3523 # fallback to iso-8859-1
3524 stream_title = stream_title_re.group(1).decode("iso-8859-1", errors="replace")
3525
3526 cleaned_stream_title = clean_stream_title(stream_title)
3527
3528 if not cleaned_stream_title:
3529 return
3530
3531 if self._record_inband_stream_title(streamdetails, cleaned_stream_title):
3532 return
3533
3534 if cleaned_stream_title == streamdetails.stream_title:
3535 return
3536
3537 self.logger.log(VERBOSE_LOG_LEVEL, "ICY Radio streamtitle original: %s", stream_title)
3538 self.logger.log(
3539 VERBOSE_LOG_LEVEL, "ICY Radio streamtitle cleaned: %s", cleaned_stream_title
3540 )
3541 streamdetails.stream_title = cleaned_stream_title
3542
3543 # Prefer station-provided cover art from the ICY 'StreamUrl' field (when it is
3544 # an image) over the MusicBrainz artwork lookup in _update_radio_stream_metadata.
3545 image_url = self._parse_icy_image_url(meta_data)
3546
3547 # Parse the original title for structured fields first so stations that announce
3548 # an album can refine the artwork lookup; fall back to the "Artist - Track" split.
3549 album: str | None = None
3550 if parsed := parse_quoted_stream_title(stream_title):
3551 track_name, artist_name_raw, album = parsed
3552 elif " - " in cleaned_stream_title:
3553 artist_name_raw, track_name = (
3554 part.strip() for part in cleaned_stream_title.split(" - ", 1)
3555 )
3556 else:
3557 return
3558
3559 if artist_name_raw and track_name:
3560 self.logger.debug(
3561 "ICY metadata: artist='%s', track='%s', album='%s'",
3562 artist_name_raw,
3563 track_name,
3564 album,
3565 )
3566 self._update_radio_stream_metadata(
3567 streamdetails,
3568 artist=artist_name_raw,
3569 title=track_name,
3570 album=album,
3571 image_url=image_url,
3572 )
3573
3574 def _parse_icy_image_url(self, meta_data: bytes) -> str | None:
3575 """
3576 Return a PNG or JPEG cover-art URL from the ICY 'StreamUrl' field, if present.
3577
3578 :param meta_data: Raw metadata bytes from an ICY stream chunk.
3579 """
3580 # The trailing semicolon is optional to match sources that omit it.
3581 stream_url_re = re.search(rb"StreamUrl='([^']*)'", meta_data)
3582 if not stream_url_re:
3583 return None
3584 try:
3585 image_url = stream_url_re.group(1).decode("utf-8").strip()
3586 except UnicodeDecodeError:
3587 return None
3588 if not image_url:
3589 return None
3590 # StreamUrl is not a standardized artwork field (reference clients such as VLC
3591 # ignore it and it conventionally holds a station website link), so only accept
3592 # values that point at a PNG or JPEG image.
3593 parsed = urlparse(image_url)
3594 if parsed.scheme not in ("http", "https"):
3595 return None
3596 if not parsed.path.lower().endswith((".png", ".jpg", ".jpeg")):
3597 return None
3598 self.logger.debug("ICY metadata: StreamUrl image='%s'", image_url)
3599 return image_url
3600
3601 async def _validate_shoutcast_stream(self, url: str) -> bool:
3602 """
3603 Return True if the URL responds with a legacy Shoutcast "ICY 200 OK" line.
3604
3605 :param url: The URL to validate.
3606 """
3607 try:
3608 parsed = urlparse(url)
3609 host = parsed.hostname
3610 port = parsed.port or 80
3611 path = parsed.path or "/"
3612 if parsed.query:
3613 path = f"{path}?{parsed.query}"
3614
3615 # Open raw socket connection with timeout
3616 reader, writer = await asyncio.wait_for(asyncio.open_connection(host, port), timeout=10)
3617 try:
3618 # Send minimal HTTP request with ICY metadata header
3619 request = f"GET {path} HTTP/1.1\r\nHost: {host}\r\nIcy-MetaData: 1\r\n\r\n"
3620 writer.write(request.encode())
3621 await writer.drain()
3622
3623 # Read just the response line
3624 response_line = await asyncio.wait_for(reader.readline(), timeout=5)
3625 finally:
3626 writer.close()
3627 await writer.wait_closed()
3628
3629 # Check if response starts with "ICY"
3630 decoded_line = response_line.decode("latin-1", errors="ignore").strip()
3631 return decoded_line.startswith("ICY")
3632
3633 except TimeoutError:
3634 self.logger.debug("Timeout during Shoutcast validation for %s", url)
3635 return False
3636 except OSError, ConnectionError:
3637 self.logger.debug("Connection failed during Shoutcast validation for %s", url)
3638 return False
3639 except UnicodeDecodeError:
3640 self.logger.debug("Invalid response encoding during Shoutcast validation for %s", url)
3641 return False
3642
3643 def _resolve_player_dsp_config(self, player: Player) -> DSPConfig:
3644 """
3645 Resolve the effective DSP config for a player.
3646
3647 Single source of truth shared by every code path that needs to know
3648 whether DSP will run for this player. Protocol wrappers defer to their
3649 parent player; single-leg ``player_group`` instances that don't expose
3650 ``MULTI_DEVICE_DSP`` defer to their first member; players whose grouping
3651 context prevents DSP get a disabled config back regardless.
3652
3653 :param player: The player to resolve DSP config for.
3654 """
3655 dsp_player_id = self._resolve_player_dsp_config_id(player)
3656 dsp = self.mass.config.get_player_dsp_config(dsp_player_id)
3657 if is_grouping_preventing_dsp(player):
3658 dsp.enabled = False
3659 elif player.provider.domain == "player_group" and (
3660 PlayerFeature.MULTI_DEVICE_DSP not in player.state.supported_features
3661 ):
3662 if not player.state.group_members:
3663 dsp.enabled = False
3664 return dsp
3665
3666 def _resolve_player_dsp_config_id(self, player: Player) -> str:
3667 """
3668 Return the player identifier that supplies the effective DSP config.
3669
3670 :param player: Player whose DSP config source should be resolved.
3671 """
3672 dsp_player_id = player.protocol_parent_id or player.player_id
3673 if (
3674 not is_grouping_preventing_dsp(player)
3675 and player.provider.domain == "player_group"
3676 and PlayerFeature.MULTI_DEVICE_DSP not in player.state.supported_features
3677 and player.state.group_members
3678 ):
3679 child_player = self.mass.players.get_player(player.state.group_members[0])
3680 assert child_player is not None
3681 dsp_player_id = child_player.player_id
3682 return dsp_player_id
3683
3684 def _get_output_channels(self, player: Player | None, player_id: str) -> str:
3685 """
3686 Return the configured output channels for the rendering player.
3687
3688 The value may be stored on the rendering player(protocol) itself (the
3689 protocol section of the config UI) or on its visible parent player (the
3690 native section); the rendering player's own stored value wins.
3691 """
3692 parent_id = player.protocol_parent_id if player and player.protocol_parent_id else player_id
3693 parent_value = self.mass.config.get_raw_player_config_value(
3694 parent_id, CONF_OUTPUT_CHANNELS, "stereo"
3695 )
3696 return self.mass.config.get_raw_player_config_value(
3697 player.player_id if player else player_id, CONF_OUTPUT_CHANNELS, parent_value
3698 )
3699
3700 def _pick_pcm_bit_depth(
3701 self,
3702 players: Iterable[Player],
3703 streamdetails: StreamDetails | None,
3704 crossfade_enabled: bool,
3705 overlay_active: bool = False,
3706 ) -> tuple[ContentType, int]:
3707 """
3708 Return ``(content_type, bit_depth)`` for an internal PCM stream.
3709
3710 F32 is chosen when audio processing (crossfade, audio overlay, volume
3711 normalization, DSP) will run on the stream â those need the extra
3712 headroom to avoid clipping and precision loss. Otherwise the source's
3713 native bit depth is reused so we don't waste memory upcasting a 16-bit
3714 stream to 32-bit just to pass it through. When the source is unknown
3715 (no streamdetails) we fall back to F32 conservatively.
3716 """
3717 if streamdetails is None:
3718 return INTERNAL_PCM_FORMAT.content_type, INTERNAL_PCM_FORMAT.bit_depth
3719 needs_headroom = (
3720 crossfade_enabled
3721 or overlay_active
3722 or streamdetails.volume_normalization_mode != VolumeNormalizationMode.DISABLED
3723 or any(self._resolve_player_dsp_config(player).enabled for player in players)
3724 )
3725 if needs_headroom:
3726 return INTERNAL_PCM_FORMAT.content_type, INTERNAL_PCM_FORMAT.bit_depth
3727 bit_depth = streamdetails.audio_format.bit_depth
3728 return ContentType.from_bit_depth(bit_depth), bit_depth
3729
3730 def _select_audio_source_pcm_format(
3731 self,
3732 player: Player,
3733 streamdetails: StreamDetails,
3734 supported_sample_rates: Iterable[int] | None = None,
3735 ) -> AudioFormat:
3736 """
3737 Return a passthrough PCM format for a realtime AudioSource item.
3738
3739 The format matches the source's native sample rate, bit depth and
3740 channel count whenever the player can accept them; if the player does
3741 not support the source's sample rate, it is snapped down to the
3742 closest supported rate. No F32 widening â realtime sources skip every
3743 processing stage that would otherwise need it. Surround sources are
3744 still folded down to stereo, which every output path requires anyway.
3745
3746 :param player: The player requesting the stream.
3747 :param streamdetails: Stream details for the AudioSource item.
3748 :param supported_sample_rates: Rates shared by every output player, if applicable.
3749 """
3750 resolved_sample_rates = (
3751 list(supported_sample_rates)
3752 if supported_sample_rates is not None
3753 else [sample_rate for sample_rate, _ in player.get_supported_sample_rates()]
3754 )
3755 source_rate = streamdetails.audio_format.sample_rate
3756 if source_rate in resolved_sample_rates:
3757 output_sample_rate = source_rate
3758 else:
3759 output_sample_rate = max(
3760 (rate for rate in resolved_sample_rates if rate <= source_rate),
3761 default=min(resolved_sample_rates),
3762 )
3763 bit_depth = streamdetails.audio_format.bit_depth
3764 return AudioFormat(
3765 content_type=ContentType.from_bit_depth(bit_depth),
3766 sample_rate=output_sample_rate,
3767 bit_depth=bit_depth,
3768 # a realtime source may announce more channels than anything downstream can
3769 # carry (a VBAN stream can be configured up to 8), and player handoff formats
3770 # copy this count straight through, so fold it here
3771 channels=min(streamdetails.audio_format.channels, 2),
3772 )
3773
3774 def _flow_restart_context(
3775 self, queue_id: str, protocol_player: Player | None
3776 ) -> tuple[str, list[int]]:
3777 """
3778 Resolve the flow mode config and supported sample rates for restart decisions.
3779
3780 Prefers the protocol player actually consuming the flow stream over the
3781 queue's (wrapper) player, whose config may lack the audio specific entries.
3782 """
3783 if protocol_player is None:
3784 protocol_player = self.mass.players.get_player(queue_id)
3785 if protocol_player is None:
3786 flow_mode_sample_rate_conf = self.mass.config.get_raw_player_config_value(
3787 queue_id, CONF_FLOW_MODE_SAMPLE_RATE, FLOW_MODE_SAMPLE_RATE_SMART
3788 )
3789 return flow_mode_sample_rate_conf, []
3790 flow_mode_sample_rate_conf = cast(
3791 "str",
3792 protocol_player.config.get_value(
3793 CONF_FLOW_MODE_SAMPLE_RATE, FLOW_MODE_SAMPLE_RATE_SMART
3794 ),
3795 )
3796 supported_sample_rates = sorted(
3797 {sr for sr, _ in protocol_player.get_supported_sample_rates()}
3798 )
3799 return flow_mode_sample_rate_conf, supported_sample_rates
3800
3801 def _flow_stream_needs_restart(
3802 self,
3803 queue_track: QueueItem,
3804 pcm_format: AudioFormat,
3805 supported_sample_rates: list[int],
3806 flow_mode_sample_rate_conf: str,
3807 is_first_track: bool,
3808 ) -> bool:
3809 """
3810 Return True if the upcoming queue track requires exiting the flow stream.
3811
3812 Covers every case where the flow loop should break and hand control back to
3813 the queue controller for restart:
3814
3815 - Live media (radio, audio sources): cannot be played inside a flow,
3816 the controller will fall back to a single-item stream.
3817 - Sample rate mismatch ('smart' / 'bit_perfect' modes only): the next
3818 track's sample rate (snapped up to the closest supported player rate,
3819 mirroring select_flow_pcm_format's anchoring logic) is incompatible with
3820 the current flow rate, so a new flow must be opened.
3821
3822 The first (anchor) track is always allowed to continue for the sample
3823 rate check; select_flow_pcm_format has already snapped the flow rate to it.
3824
3825 :param queue_track: The upcoming queue item.
3826 :param pcm_format: The current flow stream's PCM format.
3827 :param supported_sample_rates: Sorted list of the player's supported rates.
3828 :param flow_mode_sample_rate_conf: The flow mode sample rate config value.
3829 :param is_first_track: Whether this is the first track of the flow stream.
3830 """
3831 # live audio (radio, plugin or audio source) cannot be flowed; let the
3832 # queue controller fall back to single-item streaming for this item
3833 if queue_track.media_type in (MediaType.RADIO, MediaType.AUDIO_SOURCE):
3834 self.logger.info(
3835 "Live media item %s (%s, %s) encountered in flow stream "
3836 "- breaking out to single item stream",
3837 queue_track.queue_item_id,
3838 queue_track.name,
3839 queue_track.media_type,
3840 )
3841 return True
3842
3843 if is_first_track or queue_track.streamdetails is None:
3844 return False
3845 raw_next_rate = queue_track.streamdetails.audio_format.sample_rate
3846 if not raw_next_rate or not supported_sample_rates:
3847 return False
3848 effective_next_rate = _snap_supported_rate_up(raw_next_rate, supported_sample_rates)
3849
3850 # branch order mirrors select_flow_pcm_format: fixed-rate modes resample
3851 # everything to the chosen rate (no restart); bit_perfect restarts on any
3852 # mismatch; anything else falls through to smart-anchor behavior so
3853 # unknown/legacy config values don't silently pin the flow forever.
3854 if flow_mode_sample_rate_conf in (
3855 FLOW_MODE_SAMPLE_RATE_48000,
3856 FLOW_MODE_SAMPLE_RATE_96000,
3857 FLOW_MODE_SAMPLE_RATE_HIGHEST,
3858 ):
3859 needs_restart = False
3860 elif flow_mode_sample_rate_conf == FLOW_MODE_SAMPLE_RATE_BIT_PERFECT:
3861 needs_restart = effective_next_rate != pcm_format.sample_rate
3862 else:
3863 needs_restart = effective_next_rate > pcm_format.sample_rate
3864
3865 if needs_restart:
3866 self.logger.info(
3867 "Track %s (%s) sample rate %s (snapped to %s) incompatible with flow rate %s "
3868 "(mode: %s) - breaking out to restart flow stream",
3869 queue_track.queue_item_id,
3870 queue_track.name,
3871 raw_next_rate,
3872 effective_next_rate,
3873 pcm_format.sample_rate,
3874 flow_mode_sample_rate_conf,
3875 )
3876 return needs_restart
3877
3878 @asynccontextmanager
3879 async def _connect_radio_stream(self, url: str, **kwargs: Any) -> AsyncGenerator[Any]:
3880 """
3881 Connect to a radio stream URL with fallback for legacy SSL/TLS configurations.
3882
3883 Some radio servers use outdated TLS configurations that reject modern
3884 cipher suites. Since radio streams are public broadcast content,
3885 relaxing cipher requirements is acceptable.
3886
3887 :param url: The radio stream URL to connect to.
3888 :param kwargs: Additional keyword arguments passed to aiohttp get().
3889 """
3890 request_url = encoded_request_url(url)
3891 try:
3892 async with self.mass.http_session_no_ssl.get(request_url, **kwargs) as resp:
3893 yield resp
3894 except ClientConnectorSSLError:
3895 self.logger.info(
3896 "SSL handshake failed for %s, retrying with permissive cipher configuration", url
3897 )
3898 insecure_ssl_context = ssl_util.client_context_no_verify(
3899 ssl_util.SSLCipherList.INSECURE
3900 )
3901 async with self.mass.http_session_no_ssl.get(
3902 request_url, ssl=insecure_ssl_context, **kwargs
3903 ) as resp:
3904 yield resp
3905
3906 async def _update_hls_radio_metadata(
3907 self,
3908 streamdetails: StreamDetails,
3909 elapsed_time: int,
3910 ) -> None:
3911 """
3912 Update HLS radio stream metadata by fetching the playlist.
3913
3914 Fetches the HLS playlist and extracts metadata from EXTINF lines.
3915
3916 :param streamdetails: StreamDetails object to update with metadata
3917 :param elapsed_time: Current playback position in seconds (unused for live radio)
3918 """
3919 mass = self.mass
3920 try:
3921 # Get the actual media playlist URL from cache or resolve it
3922 # We cache the media_playlist_url in streamdetails.data to avoid re-resolving
3923 if streamdetails.data is None:
3924 streamdetails.data = {}
3925 media_playlist_url = streamdetails.data.get("hls_media_playlist_url")
3926 if not media_playlist_url:
3927 try:
3928 assert isinstance(streamdetails.path, str) # for type checking
3929 substream = await self.get_hls_substream(streamdetails.path)
3930 media_playlist_url = substream.path
3931 streamdetails.data["hls_media_playlist_url"] = media_playlist_url
3932 except Exception as err:
3933 self.logger.warning(
3934 "Failed to resolve HLS substream for metadata monitoring: %s", err
3935 )
3936 return
3937
3938 # Fetch the media playlist
3939 timeout = ClientTimeout(total=0, connect=10, sock_read=30)
3940 try:
3941 async with mass.http_session_no_ssl.get(
3942 encoded_request_url(media_playlist_url), timeout=timeout
3943 ) as resp:
3944 resp.raise_for_status()
3945 playlist_content = await resp.text()
3946 except ClientResponseError as err:
3947 # Session token likely expired (410/403) â drop cache so next poll re-resolves
3948 if err.status in (403, 410):
3949 streamdetails.data.pop("hls_media_playlist_url", None)
3950 raise
3951
3952 # Parse the playlist and look for EXTINF metadata
3953 # The most recent segment usually has the current metadata
3954 lines = playlist_content.strip().split("\n")
3955 for line in reversed(lines):
3956 if line.startswith("#EXTINF:"):
3957 # Extract metadata from EXTINF line
3958 metadata = parse_extinf_metadata(line)
3959
3960 # Build stream title from title and artist
3961 title = metadata.get("title", "")
3962 artist = metadata.get("artist", "")
3963 image_url = (
3964 metadata.get("image") or metadata.get("artwork") or metadata.get("cover")
3965 )
3966 if not artist and " - " in title:
3967 artist, title = title.split(" - ", 1)
3968 if title or artist:
3969 # Format as "Artist - Title"
3970 if artist and title:
3971 stream_title = f"{artist} - {title}"
3972 elif title:
3973 stream_title = title
3974 else:
3975 stream_title = artist
3976
3977 # Clean the stream title
3978 cleaned_title = clean_stream_title(stream_title)
3979
3980 # Only update if changed
3981 if cleaned_title != streamdetails.stream_title and cleaned_title:
3982 self.logger.log(
3983 VERBOSE_LOG_LEVEL, "HLS Radio metadata updated: %s", cleaned_title
3984 )
3985 streamdetails.stream_title = cleaned_title
3986 self._update_radio_stream_metadata(
3987 streamdetails,
3988 artist=artist or None,
3989 title=title or cleaned_title,
3990 image_url=image_url,
3991 )
3992
3993 # Only check the most recent EXTINF
3994 break
3995
3996 except Exception as err:
3997 self.logger.debug("Error fetching HLS metadata: %s", err)
3998
3999 @staticmethod
4000 def _normalize_reconnecting_urls(url: str | list[MultiPartPath]) -> list[str]:
4001 """Normalize a single URL or a sequence into a non-empty list."""
4002 if isinstance(url, str):
4003 return [url]
4004 if not url:
4005 msg = "Radio stream requires at least one URL"
4006 raise InvalidDataError(msg)
4007 return [part.path for part in url]
4008
4009 async def _resolve_overlay_input(self, queue: PlayerQueue) -> str | None:
4010 """
4011 Resolve the queue's overlay source to a file path or URL for ffmpeg.
4012
4013 Returns None (with a warning logged) when the source can not be resolved,
4014 so the caller can degrade to music-only playback.
4015 """
4016 if not (mapping := queue.overlay_source):
4017 return None
4018 try:
4019 provider = self.mass.get_provider(mapping.provider)
4020 if provider is None:
4021 raise MediaNotFoundError(f"Provider {mapping.provider} is not available")
4022 stream_prov = cast("MusicProvider | PluginProvider", provider)
4023 streamdetails = await stream_prov.get_stream_details(
4024 mapping.item_id, MediaType.SOUND_EFFECT
4025 )
4026 except Exception as err:
4027 self.logger.warning(
4028 "Audio overlay source %s is unavailable (%s) - continuing without overlay",
4029 mapping.uri,
4030 str(err) or err.__class__.__name__,
4031 )
4032 return None
4033 if streamdetails.stream_type not in (StreamType.LOCAL_FILE, StreamType.HTTP) or not (
4034 isinstance(streamdetails.path, str)
4035 ):
4036 self.logger.warning(
4037 "Audio overlay source %s uses unsupported stream type %s "
4038 "- continuing without overlay",
4039 mapping.uri,
4040 streamdetails.stream_type,
4041 )
4042 return None
4043 if streamdetails.stream_type == StreamType.LOCAL_FILE and not await aiofiles.os.path.isfile(
4044 streamdetails.path
4045 ):
4046 # guard against stale sources: feeding a missing file to the mixer would
4047 # kill the whole (music) stream instead of just the overlay
4048 self.logger.warning(
4049 "Audio overlay source %s does not exist - continuing without overlay",
4050 streamdetails.path,
4051 )
4052 return None
4053 return streamdetails.path
4054