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