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