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