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