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