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