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