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