/
/
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.abc import AsyncGenerator, Iterable
17from contextlib import asynccontextmanager
18from dataclasses import dataclass
19from functools import partial
20from typing import TYPE_CHECKING, Any, cast
21from urllib.parse import urlparse
22
23import aiofiles
24import aiofiles.os
25import aiohttp
26import shortuuid
27from aiohttp import ClientConnectorSSLError, ClientResponseError, ClientTimeout
28from music_assistant_models.audio_processing import (
29 AudioDSPDetails,
30 AudioOutputDetails,
31 AudioQueueProcessing,
32)
33from music_assistant_models.dsp import (
34 AudioChannel,
35 ConvolutionFilter,
36 DSPConfig,
37 DSPFilter,
38 DSPState,
39)
40from music_assistant_models.enums import (
41 ContentType,
42 CrossfadeMode,
43 MediaType,
44 PlayerFeature,
45 ProviderType,
46 StreamType,
47 VolumeNormalizationMode,
48)
49from music_assistant_models.errors import (
50 AudioError,
51 InvalidDataError,
52 MediaNotFoundError,
53 MusicAssistantError,
54 ProviderPermissionDenied,
55 ProviderUnavailableError,
56 QueueEmpty,
57 RetriesExhausted,
58)
59from music_assistant_models.media_items import AudioFormat, Track
60from music_assistant_models.player_queue import PlayLogEntry
61from music_assistant_models.streamdetails import MultiPartPath, StreamMetadata
62
63from music_assistant.constants import (
64 CONF_CROSSFADE_DURATION,
65 CONF_ENTRY_CROSSFADE_DIFFERENT_SAMPLE_RATES,
66 CONF_ENTRY_VOLUME_NORMALIZATION_TARGET,
67 CONF_FLOW_MODE_SAMPLE_RATE,
68 CONF_OUTPUT_CHANNELS,
69 CONF_PLAYER_QUEUES,
70 CONF_VALUE_DISABLED,
71 CONF_VALUE_ENABLED,
72 CONF_VOLUME_NORMALIZATION,
73 CONF_VOLUME_NORMALIZATION_FIXED_GAIN_RADIO,
74 CONF_VOLUME_NORMALIZATION_FIXED_GAIN_TRACKS,
75 CONF_VOLUME_NORMALIZATION_RADIO,
76 CONF_VOLUME_NORMALIZATION_TARGET,
77 CONF_VOLUME_NORMALIZATION_TRACKS,
78 DSP_IRS_DIRNAME,
79 FLOW_MODE_SAMPLE_RATE_48000,
80 FLOW_MODE_SAMPLE_RATE_96000,
81 FLOW_MODE_SAMPLE_RATE_BIT_PERFECT,
82 FLOW_MODE_SAMPLE_RATE_HIGHEST,
83 FLOW_MODE_SAMPLE_RATE_SMART,
84 INTERNAL_PCM_FORMAT,
85 MASS_LOGGER_NAME,
86 STREAM_STALL_TIMEOUT,
87 STREAM_START_TIMEOUT,
88 VERBOSE_LOG_LEVEL,
89)
90from music_assistant.controllers.streams.audio_analysis import (
91 LOUDNESS_ANALYSIS_DOMAIN,
92)
93from music_assistant.controllers.streams.audio_buffer import AudioBuffer
94from music_assistant.controllers.streams.audio_processing import (
95 AudioOutputPlan,
96 get_normalization_details,
97)
98from music_assistant.controllers.streams.constants import (
99 CACHE_CATEGORY_RESOLVED_RADIO_URL,
100 CACHE_PROVIDER,
101 CONF_ALLOW_CROSSFADE_SAME_ALBUM,
102)
103from music_assistant.controllers.streams.ogg_handler import get_chained_ogg_stream
104from music_assistant.controllers.streams.smart_fades import SmartFadesMixer
105from music_assistant.controllers.streams.smart_fades.fades import SmartFade
106from music_assistant.controllers.streams.smart_fades.helpers import SMART_CROSSFADE_DURATION
107from music_assistant.helpers import ssl as ssl_util
108from music_assistant.helpers.aiohttp_client import encoded_request_url
109from music_assistant.helpers.audio import (
110 HTTP_HEADERS,
111 HTTP_HEADERS_ICY,
112 build_concat_filelist,
113 calculate_content_length,
114 get_bit_rate,
115 get_normalization_mode,
116 get_parts_from_position,
117 is_grouping_preventing_dsp,
118 iter_pcm_slices,
119 parse_extinf_metadata,
120 realtime_pcm_pacer,
121 resample_pcm_audio,
122 resolve_output_player_ids,
123)
124from music_assistant.helpers.dsp import ComplexFilter, filter_to_ffmpeg_params
125from music_assistant.helpers.ffmpeg import (
126 FFMpeg,
127 get_ffmpeg_overlay_stream,
128 get_ffmpeg_stream,
129)
130from music_assistant.helpers.named_pipe import read_named_pipe
131from music_assistant.helpers.playlists import (
132 HLS_CONTENT_TYPES,
133 PLAYLIST_CONTENT_TYPES,
134 PLAYLIST_READ_TIMEOUT,
135 IsHLSPlaylist,
136 PlaylistItem,
137 parse_m3u,
138 parse_playlist_data,
139 read_playlist_body,
140)
141from music_assistant.helpers.throttle_retry import BYPASS_THROTTLER
142from music_assistant.helpers.util import (
143 clean_stream_title,
144 detect_charset,
145 parse_quoted_stream_title,
146 parse_title_and_version,
147 remove_file,
148)
149
150if TYPE_CHECKING:
151 from music_assistant_models.player_queue import PlayerQueue
152 from music_assistant_models.queue_item import QueueItem
153 from music_assistant_models.streamdetails import StreamDetails
154
155 from music_assistant.mass import MusicAssistant
156 from music_assistant.models.music_provider import MusicProvider
157 from music_assistant.models.player import Player
158 from music_assistant.models.plugin import PluginProvider
159
160# ruff: noqa: PLR0915
161
162# Seconds of PCM yielded directly to the player before the crossfade holdback starts buffering.
163WARMUP_DURATION = 8
164
165# Chunk size for the realtime AudioSource path; small enough to keep ffmpegâconsumer
166# latency below ~50 ms while still amortising per-chunk overhead.
167AUDIO_SOURCE_CHUNK_SECONDS = 0.02
168
169# Terminal errors get_icy_radio_stream raises once a single mirror is exhausted; the
170# multi-mirror reader treats these as the signal to fail over to the next URL.
171RADIO_MIRROR_FAILOVER_ERRORS = (
172 MediaNotFoundError,
173 ProviderPermissionDenied,
174 ProviderUnavailableError,
175 RetriesExhausted,
176 InvalidDataError,
177)
178
179
180@dataclass
181class CrossfadeData:
182 """Data class to hold crossfade data."""
183
184 data: bytes
185 fade_in_size: int
186 pcm_format: AudioFormat # Format of the 'data' bytes (current/previous track's format)
187 fade_in_pcm_format: AudioFormat # Format for 'fade_in_size' (next track's format)
188 queue_item_id: str
189 # Offset for the fade_in track's elapsed time calculation, to account for crossfade duration and trim
190 elapsed_time_offset: float = 0.0
191 # Normalization mode the intro PCM was baked with, used to pin the next track's body to the same mode
192 normalization_mode: VolumeNormalizationMode | None = None
193
194
195def _snap_supported_rate_up(target: int, supported_sample_rates: list[int]) -> int:
196 """Snap target up, falling back to its highest supported divisor or the maximum."""
197 if target in supported_sample_rates:
198 return target
199 higher = [r for r in supported_sample_rates if r > target]
200 if higher:
201 return min(higher)
202 same_family = [r for r in supported_sample_rates if target % r == 0]
203 return max(same_family) if same_family else max(supported_sample_rates)
204
205
206def _snap_supported_rate_down(target: int, supported_sample_rates: list[int]) -> int:
207 """Snap target down to the highest supported rate <= target, falling back to min."""
208 if target in supported_sample_rates:
209 return target
210 lower = [r for r in supported_sample_rates if r < target]
211 return max(lower) if lower else min(supported_sample_rates)
212
213
214def _pcm_formats_match(a: AudioFormat, b: AudioFormat) -> bool:
215 """Return True if two PCM formats describe identical raw bytes."""
216 return (
217 a.content_type == b.content_type
218 and a.sample_rate == b.sample_rate
219 and a.bit_depth == b.bit_depth
220 and a.channels == b.channels
221 )
222
223
224def overlay_active(queue: PlayerQueue) -> bool:
225 """Return True if the given queue has an audio overlay enabled and a source selected."""
226 return queue.overlay_enabled and queue.overlay_source is not None
227
228
229class StreamsAudio:
230 """Audio stream acquisition and processing for the streams controller."""
231
232 def __init__(self, mass: MusicAssistant) -> None:
233 """
234 Initialize StreamsAudio.
235
236 :param mass: The MusicAssistant instance.
237 """
238 self.mass = mass
239 self.logger = logging.getLogger(f"{MASS_LOGGER_NAME}.streams.audio")
240 self._crossfade_data: dict[str, CrossfadeData] = {}
241 self._smart_fades_mixer: SmartFadesMixer | None = None
242
243 def setup(self) -> None:
244 """Set up the audio sub-controller (called after all core controllers are created)."""
245 self._smart_fades_mixer = SmartFadesMixer(self.mass.streams)
246
247 @property
248 def smart_fades_mixer(self) -> SmartFadesMixer:
249 """Return the smart fades mixer."""
250 assert self._smart_fades_mixer is not None, "StreamsAudio.setup() not called"
251 return self._smart_fades_mixer
252
253 # --- Public methods ---
254
255 async def get_stream_details(
256 self,
257 queue_item: QueueItem,
258 seek_position: int = 0,
259 fade_in: bool = False,
260 prefer_album_loudness: bool = False,
261 ) -> StreamDetails:
262 """
263 Get streamdetails for the given QueueItem.
264
265 This is called just-in-time when a PlayerQueue wants a MediaItem to be played.
266 Do not try to request streamdetails too much in advance as this is expiring data.
267 """
268 mass = self.mass
269 streamdetails: StreamDetails | None = None
270 time_start = time.time()
271 self.logger.debug("Getting streamdetails for %s", queue_item.uri)
272
273 if not queue_item.media_item and not queue_item.streamdetails:
274 # in case of a non-media item queue item, the streamdetails should already be provided
275 # this should not happen, but guard it just in case
276 raise MediaNotFoundError(
277 f"Unable to retrieve streamdetails for {queue_item.name} ({queue_item.uri})"
278 )
279
280 if queue_item.streamdetails and (
281 # reuse if the buffer can serve this seek position (fast seek path)
282 (
283 queue_item.streamdetails.buffer
284 and queue_item.streamdetails.buffer.is_valid(int(seek_position * 1000))
285 )
286 # or reuse if streamdetails hasn't expired yet (new buffer will be created)
287 or (queue_item.streamdetails.created_at + queue_item.streamdetails.expiration)
288 > time.time()
289 ):
290 streamdetails = queue_item.streamdetails
291 else:
292 # need to (re)create streamdetails
293 # retrieve streamdetails from provider
294
295 media_item = queue_item.media_item
296 assert media_item is not None # for type checking
297 preferred_providers: list[str] = []
298 if (
299 (pq_data := mass.player_queues.queue_data_or_none(queue_item.queue_id))
300 and pq_data.userid
301 and (playback_user := await mass.webserver.auth.get_user(pq_data.userid))
302 and playback_user.provider_filter
303 ):
304 # handle steering into user preferred providerinstance
305 preferred_providers = playback_user.provider_filter
306 else:
307 preferred_providers = [x.provider_instance for x in media_item.provider_mappings]
308 # Remember the last AudioError so we can re-raise its (actionable)
309 # message instead of the generic MediaNotFoundError below.
310 last_audio_error: AudioError | None = None
311 attempted: set[tuple[str, str]] = set()
312 for allow_other_provider in (False, True):
313 if streamdetails:
314 break
315 # sort by quality and check item's availability
316 for prov_media in sorted(
317 media_item.provider_mappings, key=lambda x: x.quality or 0, reverse=True
318 ):
319 if not prov_media.available:
320 self.logger.debug(f"Skipping unavailable {prov_media}")
321 continue
322 if (
323 not allow_other_provider
324 and prov_media.provider_instance not in preferred_providers
325 ):
326 continue
327 # the second pass is there to widen to providers the steering held back,
328 # not to give a mapping a second chance: without a user provider filter
329 # the first pass already tried them all, so re-attempting one just buys
330 # the same failure at the cost of another provider round-trip
331 attempt = (prov_media.provider_instance, prov_media.item_id)
332 if attempt in attempted:
333 continue
334 # guard that provider is available
335 provider = mass.get_provider(prov_media.provider_instance)
336 if not provider:
337 self.logger.debug(f"Skipping {prov_media} - provider not available")
338 continue # provider not available ?
339 attempted.add(attempt)
340 # get streamdetails from provider; music and plugin providers
341 # share this signature, so either type can own the item.
342 try:
343 BYPASS_THROTTLER.set(True)
344 stream_prov = cast("MusicProvider | PluginProvider", provider)
345 streamdetails = await stream_prov.get_stream_details(
346 prov_media.item_id, media_item.media_type
347 )
348 except AudioError as err:
349 last_audio_error = err
350 self.logger.warning(str(err))
351 except MusicAssistantError as err:
352 self.logger.warning(str(err))
353 else:
354 break
355 finally:
356 BYPASS_THROTTLER.set(False)
357
358 if not streamdetails:
359 if last_audio_error is not None:
360 raise last_audio_error
361 msg = f"Unable to retrieve streamdetails for {queue_item.name} ({queue_item.uri})"
362 raise MediaNotFoundError(msg)
363
364 # work out how to handle radio stream
365 if (
366 streamdetails.stream_type in (StreamType.ICY, StreamType.HLS, StreamType.HTTP)
367 and streamdetails.media_type == MediaType.RADIO
368 and isinstance(streamdetails.path, str)
369 ):
370 resolved_url, stream_type = await self.resolve_radio_stream(streamdetails.path)
371 streamdetails.path = resolved_url
372 streamdetails.stream_type = stream_type
373 # Set up metadata monitoring callback for HLS radio streams, if not already set
374 if (
375 stream_type == StreamType.HLS
376 and not streamdetails.stream_metadata_update_callback
377 ):
378 streamdetails.stream_metadata_update_callback = partial(
379 self._update_hls_radio_metadata
380 )
381 streamdetails.stream_metadata_update_interval = 5
382
383 # providers report an unknown duration as either None or 0
384 if not streamdetails.duration:
385 if queue_item.media_item and queue_item.media_item.duration:
386 streamdetails.duration = queue_item.media_item.duration
387 elif queue_item.duration:
388 streamdetails.duration = queue_item.duration
389 if seek_position and not streamdetails.allow_seek:
390 self.logger.warning("seeking is not possible on this stream!")
391 seek_position = 0
392 elif seek_position and not streamdetails.duration:
393 self.logger.warning("seeking is not possible on duration-less streams!")
394 seek_position = 0
395
396 # set queue_id on the streamdetails so we know what is being streamed
397 streamdetails.queue_id = queue_item.queue_id
398 # handle skip/fade_in details
399 streamdetails.seek_position = seek_position
400 streamdetails.fade_in = fade_in
401
402 streamdetails.prefer_album_loudness = prefer_album_loudness
403 conf_volume_normalization_target = float(
404 mass.streams.get_config_value(CONF_VOLUME_NORMALIZATION_TARGET, return_type=int)
405 )
406 # guard against invalid volume normalization values
407 # range and default_value are guaranteed to be set for this constant
408 volume_range = CONF_ENTRY_VOLUME_NORMALIZATION_TARGET.range
409 assert volume_range is not None
410 if (
411 conf_volume_normalization_target < volume_range[0]
412 or conf_volume_normalization_target >= volume_range[1]
413 ):
414 default_val = CONF_ENTRY_VOLUME_NORMALIZATION_TARGET.default_value
415 assert isinstance(default_val, (int, float))
416 conf_volume_normalization_target = float(default_val)
417 self.logger.warning(
418 "Invalid volume normalization target configured, resetting to default of %s LUFS",
419 CONF_ENTRY_VOLUME_NORMALIZATION_TARGET.default_value,
420 )
421 streamdetails.target_loudness = conf_volume_normalization_target
422 volume_normalization_enabled = (
423 mass.config.get_effective_player_queue_config_value(
424 streamdetails.queue_id, CONF_VOLUME_NORMALIZATION, CONF_VALUE_ENABLED
425 )
426 != CONF_VALUE_DISABLED
427 )
428 streamdetails.volume_normalization_mode = get_normalization_mode(
429 self._get_volume_normalization_preference(streamdetails),
430 volume_normalization_enabled,
431 streamdetails,
432 )
433
434 self.logger.debug(
435 "Retrieved streamdetails for %s in %s milliseconds",
436 queue_item.uri,
437 int((time.time() - time_start) * 1000),
438 )
439 return streamdetails
440
441 async def get_media_stream(
442 self,
443 streamdetails: StreamDetails,
444 pcm_format: AudioFormat,
445 seek_position: int = 0,
446 filter_params: list[str] | None = None,
447 chunk_seconds: float = 1.0,
448 ) -> AsyncGenerator[bytes]:
449 """
450 Get audio stream for given media details as raw PCM.
451
452 :param streamdetails: Details of the stream to fetch.
453 :param pcm_format: Target PCM format the consumer expects.
454 :param seek_position: Seek offset in seconds (only honoured when the
455 source allows seeking; ignored for live AudioSources).
456 :param filter_params: Optional ffmpeg filter expressions.
457 :param chunk_seconds: Size of each yielded chunk in seconds of audio.
458 Defaults to 1 s for track-like sources; callers streaming live
459 AudioSources should pass a much smaller value (e.g. 0.02) to keep
460 end-to-end latency low.
461 """
462 mass = self.mass
463 logger = self.logger.getChild("media_stream")
464 logger.log(VERBOSE_LOG_LEVEL, "Starting media stream for %s", streamdetails.uri)
465 # copy: the args below are appended per call, while the StreamDetails is cached on
466 # the queue item and reused across calls (retry, seek, background analysis)
467 extra_input_args = list(streamdetails.extra_input_args or [])
468 # the branches below zero out seek_position where the seek is delegated to the
469 # source itself, so keep the requested position for the duration writeback
470 requested_seek_position = seek_position
471
472 # work out audio source for these streamdetails
473 audio_source: str | AsyncGenerator[bytes]
474 stream_type = streamdetails.stream_type
475 if stream_type == StreamType.CUSTOM:
476 # MusicProvider and PluginProvider both expose get_audio_stream with the same shape
477 provider = mass.get_provider(streamdetails.provider)
478 if provider is None:
479 raise ProviderUnavailableError(
480 f"Provider {streamdetails.provider} for stream is no longer available"
481 )
482 provider = cast("MusicProvider | PluginProvider", provider)
483 audio_source = provider.get_audio_stream(
484 streamdetails, seek_position=seek_position if streamdetails.can_seek else 0
485 )
486 seek_position = 0 if streamdetails.can_seek else seek_position
487 elif stream_type == StreamType.ICY:
488 assert streamdetails.path is not None
489 assert isinstance(streamdetails.path, (str, list))
490 audio_source = self.get_reconnecting_icy_radio_stream(streamdetails.path, streamdetails)
491 seek_position = 0
492 elif stream_type == StreamType.SHOUTCAST:
493 assert isinstance(streamdetails.path, str)
494 audio_source = self.get_shoutcast_stream(streamdetails.path, streamdetails)
495 seek_position = 0
496 elif stream_type == StreamType.IN_BAND:
497 assert isinstance(streamdetails.path, str) # for type checking
498
499 # For IN_BAND (OGG/Opus) radio streams, use chained OGG handler.
500 # This handles the chained OGG format by stitching logical bitstreams together
501 # so FFmpeg sees a single continuous stream. Metadata is extracted in-band.
502 def _on_inband_metadata(metadata: dict[str, str]) -> None:
503 """Handle metadata extracted from the OGG stream."""
504 title = metadata.get("title", "")
505 artist = metadata.get("artist", "")
506 album = metadata.get("album", "")
507 if not artist and " - " in title:
508 artist, title = title.split(" - ", 1)
509 if title or artist:
510 if artist and title:
511 stream_title = f"{artist} - {title}"
512 elif title:
513 stream_title = title
514 else:
515 stream_title = artist
516 cleaned_title = clean_stream_title(stream_title)
517 if cleaned_title and cleaned_title != streamdetails.stream_title:
518 self.logger.log(VERBOSE_LOG_LEVEL, "In-band metadata: %s", cleaned_title)
519 streamdetails.stream_title = cleaned_title
520 self._update_radio_stream_metadata(
521 streamdetails,
522 artist=artist or None,
523 title=title or cleaned_title,
524 album=album or None,
525 )
526
527 audio_source = get_chained_ogg_stream(
528 mass, streamdetails.path, metadata_callback=_on_inband_metadata
529 )
530 seek_position = 0 # seeking not possible on radio streams
531 elif stream_type == StreamType.HLS:
532 assert isinstance(streamdetails.path, str) # for type checking
533 substream = await self.get_hls_substream(streamdetails.path)
534 audio_source = substream.path
535 if streamdetails.media_type == MediaType.RADIO:
536 # HLS streams (especially the BBC) struggle when they're played directly
537 # with ffmpeg, where they just stop after some minutes,
538 # so we tell ffmpeg to loop around in this case.
539 extra_input_args += ["-stream_loop", "-1", "-re"]
540 else:
541 # all other stream types (HTTP, FILE, etc)
542 if stream_type == StreamType.ENCRYPTED_HTTP:
543 assert streamdetails.decryption_key is not None # for type checking
544 extra_input_args += ["-decryption_key", streamdetails.decryption_key]
545 if isinstance(streamdetails.path, list):
546 # multi part stream
547 audio_source = self.get_multi_file_stream(streamdetails, seek_position)
548 seek_position = 0 # handled by get_multi_file_stream
549 else:
550 # regular single file/url stream
551 assert isinstance(streamdetails.path, str) # for type checking
552 audio_source = streamdetails.path
553
554 # pace ffmpeg at native rate for live sources; the producer (e.g.
555 # librespot's pipe backend) may otherwise write faster than realtime
556 if streamdetails.media_type == MediaType.AUDIO_SOURCE and "-re" not in extra_input_args:
557 extra_input_args += ["-re"]
558
559 # handle seek support
560 if seek_position and streamdetails.duration and streamdetails.allow_seek:
561 extra_input_args += ["-ss", str(int(seek_position))]
562
563 bytes_sent = 0
564 finished = False
565 cancelled = False
566 first_chunk_received = False
567 ffmpeg_loglevel = "debug" if self.logger.isEnabledFor(VERBOSE_LOG_LEVEL) else "info"
568 # When a provider hands us already-decoded audio (e.g. Spotify Connect /
569 # AirPlay receivers piping PCM after their own decode), audio_format is
570 # the original source format meant for display while decoded_audio_format
571 # is what ffmpeg actually needs to read off the wire.
572 ffmpeg_input_format = streamdetails.decoded_audio_format or streamdetails.audio_format
573 ffmpeg_proc = FFMpeg(
574 audio_input=audio_source,
575 input_format=ffmpeg_input_format,
576 output_format=pcm_format,
577 filter_params=filter_params,
578 extra_input_args=extra_input_args,
579 collect_log_history=True,
580 loglevel=ffmpeg_loglevel,
581 )
582
583 try:
584 await ffmpeg_proc.start()
585 assert ffmpeg_proc.proc is not None # for type checking
586 if logger.isEnabledFor(VERBOSE_LOG_LEVEL):
587 logger.log(
588 VERBOSE_LOG_LEVEL,
589 "Started media stream for %s - using streamtype: %s "
590 "- pcm format: %s - ffmpeg PID: %s",
591 streamdetails.uri,
592 streamdetails.stream_type,
593 pcm_format.content_type.value,
594 ffmpeg_proc.proc.pid,
595 )
596 else:
597 logger.debug(
598 "Started media stream for %s - using streamtype: %s",
599 streamdetails.uri,
600 streamdetails.stream_type,
601 )
602 stream_start = mass.loop.time()
603 chunk_size = calculate_content_length(pcm_format, chunk_seconds)
604 chunk_iter = ffmpeg_proc.iter_chunked(chunk_size)
605 while True:
606 # Time the read, not the yield: catches a stalled source, ignores backpressure.
607 read_timeout = (
608 STREAM_START_TIMEOUT if not first_chunk_received else STREAM_STALL_TIMEOUT
609 )
610 try:
611 async with asyncio.timeout(read_timeout):
612 chunk = await anext(chunk_iter)
613 except StopAsyncIteration:
614 break
615 except TimeoutError as err:
616 raise AudioError(f"Source stalled: no audio for {read_timeout}s") from err
617 if not first_chunk_received:
618 # At this point ffmpeg has started and should now know the codec used
619 # for encoding the audio.
620 # Note: ffmpeg_proc.input_format is the same object as
621 # ffmpeg_input_format, so sample_rate / bit_depth / bit_rate
622 # parsed from the ffmpeg log already live on streamdetails too.
623 first_chunk_received = True
624 # Skip the codec_type writeback when the provider declared a
625 # decoded format: audio_format already holds the authoritative
626 # source codec and the probed value would just be the
627 # post-decode wire format (e.g. PCM for Spotify Connect).
628 if streamdetails.decoded_audio_format is None:
629 streamdetails.audio_format.codec_type = ffmpeg_proc.input_format.codec_type
630 # Some providers omit (or report 0 for) the item duration; ffmpeg can
631 # usually probe it from the source. Only apply when missing so we
632 # don't clobber an accurate provider value with a rounded one.
633 if ffmpeg_proc.parsed_duration is not None and not streamdetails.duration:
634 streamdetails.duration = ffmpeg_proc.parsed_duration
635 logger.debug(
636 "First chunk received after %.2f seconds (codec detected: %s)",
637 mass.loop.time() - stream_start,
638 ffmpeg_proc.input_format.codec_type,
639 )
640 yield chunk
641 bytes_sent += len(chunk)
642
643 # end of audio/track reached
644 logger.debug("End of media stream reached for %s", streamdetails.uri)
645 # wait until stderr also completed reading
646 await ffmpeg_proc.wait_with_timeout(5)
647 logger.log(
648 VERBOSE_LOG_LEVEL,
649 "FFmpeg process ended with return code %s for %s",
650 ffmpeg_proc.returncode,
651 streamdetails.uri,
652 )
653 if ffmpeg_proc.returncode not in (0, None):
654 log_trail = "\n".join(list(ffmpeg_proc.log_history)[-5:])
655 raise AudioError(f"FFMpeg exited with code {ffmpeg_proc.returncode}: {log_trail}")
656 if bytes_sent == 0:
657 # edge case: no audio data was received at all
658 raise AudioError("No audio was received")
659 finished = True
660 except (Exception, GeneratorExit, asyncio.CancelledError) as err:
661 if isinstance(err, asyncio.CancelledError | GeneratorExit):
662 # we were cancelled, just raise
663 cancelled = True
664 raise
665 # dump the last 10 lines of the log in case of an unclean exit
666 logger.warning("\n".join(list(ffmpeg_proc.log_history)[-10:]))
667 raise AudioError(f"Error while streaming: {err}") from err
668 finally:
669 # always ensure close is called which also handles all cleanup
670 await ffmpeg_proc.close()
671 # determine how many seconds we've received
672 # for pcm output we can calculate this easily
673 seconds_received = bytes_sent / pcm_format.pcm_sample_size if bytes_sent else 0
674 # store accurate duration, but only for a playthrough from the very start:
675 # a seeked stream yields the remaining audio, not the item's full length
676 if finished and not requested_seek_position and seconds_received:
677 streamdetails.duration = int(seconds_received)
678
679 logger.log(
680 VERBOSE_LOG_LEVEL,
681 "stream %s (with code %s) for %s",
682 "cancelled" if cancelled else "finished" if finished else "aborted",
683 ffmpeg_proc.returncode,
684 streamdetails.uri,
685 )
686
687 async def resolve_radio_stream(self, url: str) -> tuple[str, StreamType]:
688 """
689 Resolve a streaming radio URL.
690
691 Unwraps playlists and determines stream type (ICY, HLS, SHOUTCAST, IN_BAND, HTTP).
692
693 :param url: Radio stream URL to resolve
694 """
695 mass = self.mass
696 if cache := await mass.cache.get(
697 key=url, provider=CACHE_PROVIDER, category=CACHE_CATEGORY_RESOLVED_RADIO_URL
698 ):
699 if TYPE_CHECKING:
700 cache = cast("tuple[str, str]", cache)
701 return (cache[0], StreamType(cache[1]))
702
703 stream_type = StreamType.HTTP
704 timeout = ClientTimeout(total=None, connect=10, sock_read=5)
705 playlist_data: bytes | None = None
706 playlist_charset: str | None = None
707
708 try:
709 async with self._connect_radio_stream(
710 url, headers=HTTP_HEADERS_ICY, allow_redirects=True, timeout=timeout
711 ) as resp:
712 headers = resp.headers
713 resp.raise_for_status()
714 if not resp.headers:
715 raise InvalidDataError("no headers found")
716 # media types are case insensitive, the comparisons below are all lower case
717 content_type = headers.get("content-type", "").lower()
718 # a server declaring HLS settles it: a media playlist is free to carry none
719 # of the tags the parser recognises an HLS playlist by
720 is_hls = any(hls_type in content_type for hls_type in HLS_CONTENT_TYPES)
721 if not is_hls and (
722 url.endswith((".m3u", ".m3u8", ".pls"))
723 or ".m3u?" in url
724 or ".m3u8?" in url
725 or ".pls?" in url
726 or any(
727 playlist_type in content_type for playlist_type in PLAYLIST_CONTENT_TYPES
728 )
729 ):
730 # take the playlist from this very response: a separate request would
731 # go out with another user agent and stricter TLS than the rest of the
732 # radio paths, so a host could answer it differently
733 try:
734 # the probe has no total timeout, so bound the body on its own:
735 # a server trickling bytes would otherwise stall resolving for hours
736 async with asyncio.timeout(PLAYLIST_READ_TIMEOUT):
737 playlist_data = await read_playlist_body(resp.content)
738 except aiohttp.ClientError as err:
739 # the endpoint answered as a playlist, so a truncated body is a bad
740 # playlist - not a reason to fall back to streaming the URL directly
741 raise InvalidDataError(f"Error while fetching playlist {url}") from err
742 playlist_charset = resp.charset
743
744 if headers.get("icy-metaint") is not None:
745 stream_type = StreamType.ICY
746 elif is_hls:
747 stream_type = StreamType.HLS
748 elif content_type in ("application/ogg", "audio/ogg"):
749 # Ogg streams (Opus/Vorbis) have in-band metadata via Vorbis comments
750 stream_type = StreamType.IN_BAND
751
752 if playlist_data is not None:
753 try:
754 substreams = await parse_playlist_data(url, playlist_data, playlist_charset)
755 if not any(x for x in substreams if x.length):
756 for line in substreams:
757 if not line.is_url:
758 continue
759 return await self.resolve_radio_stream(line.path)
760 raise InvalidDataError("No content found in playlist")
761 except IsHLSPlaylist:
762 stream_type = StreamType.HLS
763
764 except TimeoutError as err:
765 self.logger.warning("Timeout while parsing radio URL %s", url)
766 raise InvalidDataError(f"Timeout connecting to {url}") from err
767
768 except aiohttp.ClientResponseError as err:
769 if err.status == 404:
770 raise MediaNotFoundError(f"Radio stream not found: {url}") from err
771 if err.status == 403:
772 raise InvalidDataError(f"Access denied to radio stream: {url}") from err
773 if err.status >= 500:
774 raise InvalidDataError(
775 f"Radio stream server error (HTTP {err.status}): {url}"
776 ) from err
777 if err.status == 400:
778 # 400 errors might be from legacy Shoutcast servers
779 return await self._handle_client_error_for_radio_stream(url, err, stream_type)
780 raise InvalidDataError(f"HTTP error {err.status} from {url}") from err
781
782 except aiohttp.ClientError as err:
783 return await self._handle_client_error_for_radio_stream(url, err, stream_type)
784
785 return await self._cache_radio_result(url, stream_type)
786
787 async def get_icy_radio_stream(
788 self, url: str, streamdetails: StreamDetails
789 ) -> AsyncGenerator[bytes]:
790 """
791 Stream radio audio with ICY metadata support, reconnecting on disconnect.
792
793 Requires icy-metaint header support. Stream type should be validated
794 by resolve_radio_stream() before calling this function.
795
796 :param url: Radio stream URL
797 :param streamdetails: StreamDetails to update with metadata
798 """
799 self.logger.debug("Start streaming radio with ICY metadata from url %s", url)
800 timeout = ClientTimeout(total=0, connect=30, sock_read=5 * 60)
801 # Budget for *consecutive* reconnects that delivered no audio. A connection
802 # that actually streamed data resets it, so a healthy long-running stream can
803 # reconnect indefinitely while a dead/looping one bails out instead of spinning.
804 failed_reconnects = 0
805 max_failed_reconnects = 25
806
807 while True:
808 streamed_data = False
809 try:
810 async with self._connect_radio_stream(
811 url, allow_redirects=True, headers=HTTP_HEADERS_ICY, timeout=timeout
812 ) as resp:
813 # surface a non-200 (e.g. on reconnect) as a ClientResponseError so the
814 # terminal/HTTP handling below applies instead of failing on the header
815 resp.raise_for_status()
816 meta_int_str = resp.headers.get("icy-metaint")
817 if not meta_int_str:
818 raise InvalidDataError(f"No icy-metaint header for radio stream: {url}")
819 try:
820 meta_int = int(meta_int_str)
821 except ValueError as err:
822 raise InvalidDataError(
823 f"Invalid icy-metaint value for radio stream: {url}"
824 ) from err
825 if meta_int <= 0:
826 raise InvalidDataError(f"Invalid icy-metaint value for radio stream: {url}")
827 # readexactly raises IncompleteReadError when the server closes the
828 # connection mid-frame; that (and the network errors below) drops us
829 # out to the reconnect handler so a live stream survives the blip.
830 while True:
831 chunk = await resp.content.readexactly(meta_int)
832 streamed_data = True
833 yield chunk
834 meta_byte = await resp.content.readexactly(1)
835 if meta_byte == b"\x00":
836 continue
837 meta_length = ord(meta_byte) * 16
838 meta_data = await resp.content.readexactly(meta_length)
839 self._parse_icy_metadata(meta_data, streamdetails)
840 except asyncio.CancelledError:
841 self.logger.debug("ICY radio stream cancelled for %s", url)
842 raise
843 except aiohttp.ClientResponseError as err:
844 if err.status == 404:
845 raise MediaNotFoundError(f"Radio stream not found: {url}") from err
846 if err.status == 403:
847 raise ProviderPermissionDenied(f"Radio stream access denied: {url}") from err
848 raise ProviderUnavailableError(
849 f"Radio stream returned HTTP {err.status}: {err}"
850 ) from err
851 except (
852 asyncio.IncompleteReadError,
853 aiohttp.ClientConnectionError,
854 aiohttp.ClientPayloadError,
855 aiohttp.ServerDisconnectedError,
856 ) as err:
857 if streamed_data:
858 # a healthy session that dropped - reconnect without spending budget
859 failed_reconnects = 0
860 self.logger.debug("ICY radio stream dropped, reconnecting: %s", err)
861 else:
862 failed_reconnects += 1
863 if failed_reconnects > max_failed_reconnects:
864 raise RetriesExhausted(
865 f"ICY radio stream failed after {max_failed_reconnects} "
866 f"reconnects without data: {err}"
867 ) from err
868 self.logger.warning(
869 "ICY radio stream reconnect produced no data (%d/%d): %s",
870 failed_reconnects,
871 max_failed_reconnects,
872 err,
873 )
874 await asyncio.sleep(0.5)
875
876 async def get_reconnecting_icy_radio_stream(
877 self, url: str | list[MultiPartPath], streamdetails: StreamDetails
878 ) -> AsyncGenerator[bytes]:
879 """
880 Yield ICY radio audio with metadata, failing over across mirror URLs.
881
882 A single URL is delegated to :meth:`get_icy_radio_stream`, which already reconnects
883 on disconnect. Multiple URLs are treated as interchangeable mirrors and tried in turn;
884 a mirror that delivers audio resets the failover budget, so a healthy mirror keeps
885 streaming while a set of unreachable mirrors raises the last error instead of spinning.
886
887 :param url: One stream URL, or a list of mirror URLs to fail over between.
888 :param streamdetails: StreamDetails to update with metadata.
889 """
890 urls = self._normalize_reconnecting_urls(url)
891 if len(urls) == 1:
892 async for chunk in self.get_icy_radio_stream(urls[0], streamdetails):
893 yield chunk
894 return
895
896 url_index = 0
897 failed_rotations = 0
898 max_failed_rotations = len(urls) * 2
899 last_err: MusicAssistantError | None = None
900 while failed_rotations <= max_failed_rotations:
901 current_url = urls[url_index % len(urls)]
902 url_index += 1
903 delivered_audio = False
904 try:
905 async for chunk in self.get_icy_radio_stream(current_url, streamdetails):
906 delivered_audio = True
907 failed_rotations = 0
908 # release the previous failure while healthy: it pins the full
909 # exception traceback (with frames) for the lifetime of the stream
910 last_err = None
911 yield chunk
912 return
913 except RADIO_MIRROR_FAILOVER_ERRORS as err:
914 last_err = err
915 if not delivered_audio:
916 failed_rotations += 1
917 self.logger.warning(
918 "ICY radio mirror %s failed, trying next url (%d/%d): %s",
919 current_url,
920 failed_rotations,
921 max_failed_rotations,
922 err,
923 )
924 if last_err is not None:
925 raise last_err
926
927 async def get_reconnecting_radio_stream(self, url: str) -> AsyncGenerator[bytes]:
928 """
929 Yield continuous radio stream data, automatically reconnecting on disconnect.
930
931 :param url: URL of the radio stream.
932 """
933 timeout = ClientTimeout(total=None, connect=30, sock_read=5 * 60)
934 reconnect_count = 0
935 max_reconnects = 1000 # Allow many reconnects for long-running radio
936
937 while reconnect_count <= max_reconnects:
938 try:
939 async with self._connect_radio_stream(
940 url, allow_redirects=True, headers=HTTP_HEADERS, timeout=timeout
941 ) as resp:
942 chunk_count = 0
943 async for chunk in resp.content.iter_any():
944 chunk_count += 1
945 yield chunk
946
947 # Connection closed normally - reconnect
948 self.logger.debug(
949 "Radio stream connection closed after %d chunks, reconnecting... "
950 "(reconnect #%d)",
951 chunk_count,
952 reconnect_count,
953 )
954 reconnect_count += 1
955 await asyncio.sleep(0.1) # Brief delay before reconnect
956
957 except asyncio.CancelledError:
958 self.logger.debug("Radio stream cancelled for %s", url)
959 raise
960 except (
961 aiohttp.ClientConnectionError,
962 aiohttp.ClientPayloadError,
963 aiohttp.ServerDisconnectedError,
964 ) as err:
965 # Transient network errors - retry
966 self.logger.warning("Radio stream error (reconnect #%d): %s", reconnect_count, err)
967 reconnect_count += 1
968 if reconnect_count > max_reconnects:
969 raise RetriesExhausted(
970 f"Radio stream failed after {max_reconnects} reconnects: {err}"
971 ) from err
972 await asyncio.sleep(0.5)
973 except aiohttp.ClientResponseError as err:
974 if err.status == 404:
975 raise MediaNotFoundError(f"Radio stream not found: {url}") from err
976 if err.status == 403:
977 raise ProviderPermissionDenied(f"Radio stream access denied: {url}") from err
978 # Other HTTP errors (5xx etc) - could be temporary
979 raise ProviderUnavailableError(
980 f"Radio stream returned HTTP {err.status}: {err}"
981 ) from err
982
983 self.logger.warning("Radio stream reached max reconnects (%d) for %s", max_reconnects, url)
984
985 async def get_hls_substream(self, url: str) -> PlaylistItem:
986 """Select the (highest quality) HLS substream for given HLS playlist/URL."""
987 mass = self.mass
988 timeout = ClientTimeout(total=None, connect=30, sock_read=5 * 60)
989 # fetch master playlist and select (best) child playlist
990 # https://datatracker.ietf.org/doc/html/draft-pantos-http-live-streaming-19#section-10
991 async with mass.http_session_no_ssl.get(
992 encoded_request_url(url), allow_redirects=True, headers=HTTP_HEADERS, timeout=timeout
993 ) as resp:
994 resp.raise_for_status()
995 raw_data = await resp.read()
996 encoding = await detect_charset(raw_data, preferred=resp.charset)
997 master_m3u_data = raw_data.decode(encoding, errors="replace")
998 substreams = parse_m3u(master_m3u_data)
999 # There is a chance that we did not get a master playlist with subplaylists
1000 # but just a single master/sub playlist with the actual audio stream(s)
1001 # so we need to detect if the playlist child's contain audio streams or
1002 # sub-playlists.
1003 if any(
1004 x
1005 for x in substreams
1006 if (x.length or x.path.endswith((".mp4", ".aac")))
1007 and not x.path.endswith((".m3u", ".m3u8"))
1008 ):
1009 return PlaylistItem(path=url, key=substreams[0].key)
1010 # sort substreams on best quality (highest bandwidth) when available
1011 if any(x for x in substreams if x.stream_info):
1012 substreams.sort(
1013 key=lambda x: int(
1014 x.stream_info.get("BANDWIDTH", "0") if x.stream_info is not None else 0
1015 ),
1016 reverse=True,
1017 )
1018 substream = substreams[0]
1019 if not substream.path.startswith("http"):
1020 # path is relative, stitch it together
1021 base_path = url.rsplit("/", 1)[0]
1022 substream.path = base_path + "/" + substream.path
1023 return substream
1024
1025 async def get_multi_file_stream(
1026 self,
1027 streamdetails: StreamDetails,
1028 seek_position: int = 0,
1029 ) -> AsyncGenerator[bytes]:
1030 """
1031 Return audio stream for a concatenation of multiple files.
1032
1033 Arguments:
1034 seek_position: The position to seek to in seconds
1035 """
1036 if not isinstance(streamdetails.path, list):
1037 raise InvalidDataError("Multi-file streamdetails requires a list of MultiPartPath")
1038 parts, seek_position = get_parts_from_position(streamdetails.path, seek_position)
1039 files_list = [part.path for part in parts]
1040
1041 # concat input files
1042 temp_file = f"/tmp/{shortuuid.random(20)}.txt" # noqa: S108
1043 async with aiofiles.open(temp_file, "w") as f:
1044 await f.write(build_concat_filelist(files_list))
1045
1046 try:
1047 async for chunk in get_ffmpeg_stream(
1048 audio_input=temp_file,
1049 input_format=streamdetails.audio_format,
1050 output_format=AudioFormat(
1051 content_type=ContentType.NUT,
1052 sample_rate=streamdetails.audio_format.sample_rate,
1053 bit_depth=streamdetails.audio_format.bit_depth,
1054 channels=streamdetails.audio_format.channels,
1055 ),
1056 extra_input_args=[
1057 "-safe",
1058 "0",
1059 "-f",
1060 "concat",
1061 "-i",
1062 temp_file,
1063 "-ss",
1064 str(seek_position),
1065 ],
1066 ):
1067 yield chunk
1068 finally:
1069 await remove_file(temp_file)
1070
1071 def get_player_output_plan(
1072 self,
1073 player_id: str,
1074 input_format: AudioFormat,
1075 output_format: AudioFormat,
1076 *,
1077 shared_player_ids: Iterable[str] | None = None,
1078 handoff_format: AudioFormat | None = None,
1079 queue_id: str | None = None,
1080 session_id: str | None = None,
1081 queue_item_id: str | None = None,
1082 ) -> AudioOutputPlan:
1083 """
1084 Return executable filters and matching output details for a player.
1085
1086 :param player_id: Destination player identifier.
1087 :param input_format: PCM format entering player-specific processing.
1088 :param output_format: Furthest downstream output format known to the server.
1089 :param shared_player_ids: Additional players receiving this identical output path.
1090 An empty iterable marks a path that can gain shared destinations later.
1091 :param handoff_format: Earlier provider handoff format when it differs.
1092 :param queue_id: Explicit queue identifier for the processing snapshot.
1093 :param session_id: Explicit queue session identifier for the processing snapshot.
1094 :param queue_item_id: Queue item for a single-item output path.
1095 """
1096 filter_params: list[str | ComplexFilter] = []
1097 player = self.mass.players.get_player(player_id)
1098 destination_player_id = (
1099 player.protocol_parent_id if player and player.protocol_parent_id else player_id
1100 )
1101 resolved_shared_player_ids = (
1102 resolve_output_player_ids(self.mass, shared_player_ids) - {destination_player_id}
1103 if shared_player_ids is not None
1104 else None
1105 )
1106 destination_player_ids = {destination_player_id, *(resolved_shared_player_ids or ())}
1107 if player:
1108 dsp_config_id = self._resolve_player_dsp_config_id(player)
1109 dsp = self._resolve_player_dsp_config(player)
1110 configured_dsp = self.mass.config.get_player_dsp_config(dsp_config_id)
1111 if configured_dsp.enabled and not dsp.enabled and is_grouping_preventing_dsp(player):
1112 dsp_state = DSPState.DISABLED_BY_UNSUPPORTED_GROUP
1113 else:
1114 dsp_state = DSPState.ENABLED if dsp.enabled else DSPState.DISABLED
1115 else:
1116 dsp_config_id = player_id
1117 dsp = self.mass.config.get_player_dsp_config(player_id)
1118 dsp_state = DSPState.ENABLED if dsp.enabled else DSPState.DISABLED
1119
1120 enabled_filters = [dsp_filter for dsp_filter in dsp.filters if dsp_filter.enabled]
1121 # a neutral filter (0 dB gain, centered balance) emits no params; exclude
1122 # it so it is not reported as an active, non-bit-perfect stage
1123 effective_filters: list[DSPFilter] = []
1124 if dsp.enabled:
1125 if dsp.input_gain != 0:
1126 filter_params.append(f"volume={dsp.input_gain}dB")
1127 ir_dir = os.path.join(self.mass.storage_path, DSP_IRS_DIRNAME)
1128 known_ir_ids = {record["ir_id"] for record in self.mass.config.get_dsp_irs()}
1129 for dsp_filter in enabled_filters:
1130 if isinstance(dsp_filter, ConvolutionFilter) and dsp_filter.ir_id:
1131 # ffmpeg fails to open the graph if the impulse response file is gone,
1132 # which costs the player all audio, so drop the filter instead
1133 if dsp_filter.ir_id not in known_ir_ids:
1134 self.logger.warning(
1135 "Skipping the convolution filter of player %s: "
1136 "impulse response %s is not stored",
1137 player_id,
1138 dsp_filter.ir_id,
1139 )
1140 continue
1141 params = filter_to_ffmpeg_params(dsp_filter, input_format, ir_dir=ir_dir)
1142 if not params:
1143 continue
1144 filter_params.extend(params)
1145 effective_filters.append(dsp_filter)
1146 if dsp.output_gain != 0:
1147 filter_params.append(f"volume={dsp.output_gain}dB")
1148
1149 channel_value = self._get_output_channels(player, player_id)
1150 source_channel = None
1151 channel_mix = ""
1152 # a single channel source is already the downmix and holds no FL/FR to select
1153 # from, where a pan would silently resolve every gain to zero
1154 if input_format.channels > 1:
1155 if channel_value == "left":
1156 source_channel = AudioChannel.FL
1157 channel_mix = "FL"
1158 elif channel_value == "right":
1159 source_channel = AudioChannel.FR
1160 channel_mix = "FR"
1161 elif channel_value == "mono":
1162 # both source channels feed the downmix, so report ALL to keep the
1163 # output from ever being presented as bit perfect
1164 source_channel = AudioChannel.ALL
1165 channel_mix = "0.5*FL+0.5*FR"
1166 if channel_mix:
1167 # the pan runs in the command that emits the handoff format, and it feeds
1168 # every channel of it explicitly: leaving ffmpeg to upmix from a single
1169 # channel costs 3 dB through its rematrix
1170 if (handoff_format or output_format).channels == 1:
1171 filter_params.append(f"pan=mono|c0={channel_mix}")
1172 else:
1173 filter_params.append(f"pan=stereo|c0={channel_mix}|c1={channel_mix}")
1174
1175 output_details = AudioOutputDetails(
1176 player_ids=sorted(destination_player_ids),
1177 dsp=AudioDSPDetails(
1178 state=dsp_state,
1179 input_gain=dsp.input_gain if dsp.enabled else 0.0,
1180 filters=effective_filters,
1181 output_gain=dsp.output_gain if dsp.enabled else 0.0,
1182 preset_id=dsp.preset_id,
1183 ),
1184 source_channel=source_channel,
1185 output_format=output_format,
1186 )
1187 output_plan = AudioOutputPlan(
1188 filter_params=filter_params,
1189 output_details=output_details,
1190 input_format=input_format,
1191 handoff_format=handoff_format,
1192 dsp_config_id=dsp_config_id,
1193 )
1194 if queue_id is not None and session_id is not None:
1195 self.mass.streams.audio_processing.update_output(
1196 destination_player_id,
1197 output_plan,
1198 shared_player_ids=resolved_shared_player_ids,
1199 queue_id=queue_id,
1200 session_id=session_id,
1201 queue_item_id=queue_item_id,
1202 )
1203 self.logger.log(
1204 VERBOSE_LOG_LEVEL,
1205 "Generated ffmpeg params for player %s: %s",
1206 player_id,
1207 filter_params,
1208 )
1209 return output_plan
1210
1211 async def get_output_format(
1212 self,
1213 output_format_str: str,
1214 player: Player,
1215 content_sample_rate: int,
1216 content_bit_depth: int,
1217 media_type: MediaType = MediaType.UNKNOWN,
1218 ) -> AudioFormat:
1219 """Parse (player specific) output format details for given format string."""
1220 content_type: ContentType = ContentType.try_parse(output_format_str)
1221 player_supported_rates = player.get_supported_sample_rates()
1222 supported_sample_rates = [sr for sr, _ in player_supported_rates]
1223 if content_sample_rate in supported_sample_rates:
1224 output_sample_rate = content_sample_rate
1225 else:
1226 output_sample_rate = max(supported_sample_rates)
1227 # only consider bit depths that are actually paired with the chosen sample rate
1228 bit_depths_for_rate = [
1229 bd for (sr, bd) in player_supported_rates if sr == output_sample_rate
1230 ]
1231 output_bit_depth = min(content_bit_depth, max(bit_depths_for_rate, default=16))
1232
1233 if not content_type.is_lossless():
1234 # no point in having a higher bit depth for lossy formats
1235 output_bit_depth = 16
1236 output_sample_rate = min(48000, output_sample_rate)
1237 if media_type not in (MediaType.TRACK, MediaType.AUDIO_SOURCE, MediaType.FLOW_STREAM):
1238 # no point in having a higher bit depth for non-track media types (e.g. TTS, radio)
1239 output_bit_depth = min(output_bit_depth, 16)
1240 if output_format_str == "pcm":
1241 content_type = ContentType.from_bit_depth(output_bit_depth)
1242
1243 output_channels_str = self._get_output_channels(player, player.player_id)
1244 fmt = AudioFormat(
1245 content_type=content_type,
1246 sample_rate=output_sample_rate,
1247 bit_depth=output_bit_depth,
1248 channels=1 if output_channels_str != "stereo" else 2,
1249 )
1250 fmt.bit_rate = get_bit_rate(fmt)
1251 return fmt
1252
1253 async def select_pcm_format(
1254 self,
1255 player: Player,
1256 streamdetails: StreamDetails,
1257 crossfade_enabled: bool,
1258 overlay_active: bool = False,
1259 ) -> AudioFormat:
1260 """
1261 Select the internal PCM format for streaming a single queue item.
1262
1263 Used by the per-item (non-flow) stream path. The sample rate is the highest
1264 rate the player supports that is <= the source rate, so the source is never
1265 upsampled. The bit depth follows the source unless audio processing
1266 (crossfade, volume normalization, DSP) is active â those need F32 headroom
1267 to avoid clipping/precision loss. Surround sources are folded down to stereo.
1268 Realtime AudioSource items skip all processing and get a pure passthrough
1269 format (source rate/bit depth when the player supports them).
1270
1271 :param player: The player requesting the stream.
1272 :param streamdetails: Stream details for the current item.
1273 :param crossfade_enabled: Whether crossfade is enabled for this stream.
1274 :param overlay_active: Whether an audio overlay will be mixed into this stream.
1275 """
1276 if streamdetails.media_type == MediaType.AUDIO_SOURCE:
1277 return self._select_audio_source_pcm_format(player, streamdetails)
1278 supported_sample_rates = [sr for sr, _ in player.get_supported_sample_rates()]
1279 # snap-down: pick the highest supported rate <= source. when the source rate
1280 # is below every supported rate (e.g. 22 kHz content on a 44.1k-only player),
1281 # fall back to the lowest supported rate instead of a hardcoded 48 kHz that
1282 # the player may not actually support.
1283 output_sample_rate = max(
1284 (r for r in supported_sample_rates if r <= streamdetails.audio_format.sample_rate),
1285 default=min(supported_sample_rates),
1286 )
1287 content_type, bit_depth = self._pick_pcm_bit_depth(
1288 (player,),
1289 streamdetails,
1290 crossfade_enabled,
1291 overlay_active,
1292 )
1293 pcm_format = AudioFormat(
1294 sample_rate=output_sample_rate,
1295 content_type=content_type,
1296 bit_depth=bit_depth,
1297 # fold surround sources down to stereo right at the decode step: no
1298 # output format carries more than two channels, so a wider PCM format
1299 # only makes every bytes-to-seconds sum on the stream come out short
1300 channels=min(streamdetails.audio_format.channels, 2),
1301 )
1302 if crossfade_enabled or overlay_active:
1303 pcm_format.channels = 2
1304 return pcm_format
1305
1306 async def select_flow_pcm_format(
1307 self,
1308 player: Player,
1309 start_streamdetails: StreamDetails | None = None,
1310 crossfade_enabled: bool = False,
1311 overlay_active: bool = False,
1312 fallback_sample_rate: int | None = None,
1313 output_players: Iterable[Player] | None = None,
1314 ) -> AudioFormat:
1315 """
1316 Select the internal PCM format for a Queue Flow Mode stream.
1317
1318 Used by the gapless flow path that stitches multiple queue items into one
1319 continuous PCM stream. The sample rate is driven by the player's
1320 ``CONF_FLOW_MODE_SAMPLE_RATE`` setting (smart/bit_perfect/48k/96k/highest)
1321 â for the anchored modes it follows the first track's rate, for the fixed
1322 modes it snaps to the configured rate. The bit depth follows the first
1323 track's source unless audio processing is active (then F32 for headroom),
1324 avoiding an unnecessary up-convert to 32-bit when none of the consumers
1325 will benefit from it. When the first item is a realtime AudioSource, the
1326 flow mode config is ignored and a pure passthrough format is used so the
1327 source audio is delivered with minimum overhead and latency.
1328
1329 :param player: The player the flow stream is being prepared for.
1330 :param start_streamdetails: Stream details of the first track in the flow.
1331 Required for the anchored modes ('smart' / 'bit_perfect') and for the
1332 bit-depth optimization. May be omitted for the fixed-rate modes â when
1333 omitted the bit depth defaults to F32.
1334 :param crossfade_enabled: Whether the queue will use crossfade transitions.
1335 :param overlay_active: Whether an audio overlay will be mixed into the stream.
1336 :param fallback_sample_rate: Preferred rate when the first item format is unknown.
1337 :param output_players: All players consuming the shared PCM stream. Their common
1338 sample rates and processing requirements determine the session format.
1339 """
1340 players = tuple(output_players) if output_players is not None else (player,)
1341 if not players:
1342 raise AudioError("At least one output player is required")
1343 supported_sample_rates = sorted(
1344 set.intersection(
1345 *(
1346 {sample_rate for sample_rate, _ in item.get_supported_sample_rates()}
1347 for item in players
1348 )
1349 )
1350 )
1351 if not supported_sample_rates:
1352 raise AudioError("Output players do not share a supported sample rate")
1353 if start_streamdetails is not None and (
1354 start_streamdetails.media_type == MediaType.AUDIO_SOURCE
1355 ):
1356 return self._select_audio_source_pcm_format(
1357 player,
1358 start_streamdetails,
1359 supported_sample_rates=supported_sample_rates,
1360 )
1361 flow_mode_conf = cast(
1362 "str",
1363 player.config.get_value(CONF_FLOW_MODE_SAMPLE_RATE, FLOW_MODE_SAMPLE_RATE_SMART),
1364 )
1365
1366 if flow_mode_conf == FLOW_MODE_SAMPLE_RATE_HIGHEST:
1367 output_sample_rate = max(supported_sample_rates)
1368 elif flow_mode_conf == FLOW_MODE_SAMPLE_RATE_48000:
1369 # for the fixed-rate modes, the user picked a specific bandwidth/quality
1370 # ceiling; prefer the highest supported rate <= target
1371 output_sample_rate = _snap_supported_rate_down(48000, supported_sample_rates)
1372 elif flow_mode_conf == FLOW_MODE_SAMPLE_RATE_96000:
1373 output_sample_rate = _snap_supported_rate_down(96000, supported_sample_rates)
1374 else:
1375 # smart or bit_perfect (default): anchor the flow at the starting track's
1376 # sample rate; if the player doesn't natively support it, upsample to the
1377 # closest higher supported rate
1378 target_rate = (
1379 start_streamdetails.audio_format.sample_rate
1380 if start_streamdetails
1381 else (
1382 fallback_sample_rate
1383 if fallback_sample_rate is not None
1384 else max(supported_sample_rates)
1385 )
1386 )
1387 output_sample_rate = _snap_supported_rate_up(target_rate, supported_sample_rates)
1388
1389 content_type, bit_depth = self._pick_pcm_bit_depth(
1390 players, start_streamdetails, crossfade_enabled, overlay_active
1391 )
1392 return AudioFormat(
1393 content_type=content_type,
1394 sample_rate=output_sample_rate,
1395 bit_depth=bit_depth,
1396 channels=2,
1397 )
1398
1399 async def get_audio_source_stream(
1400 self,
1401 queue_item: QueueItem,
1402 pcm_format: AudioFormat,
1403 raise_on_error: bool = True,
1404 ) -> AsyncGenerator[bytes]:
1405 """
1406 Get the realtime PCM stream for an AudioSource queue item.
1407
1408 AudioSources are live/realtime: bytes flow at the producer's pace, with
1409 no pre-buffering, no loudness hydration, no volume normalization, no
1410 crossfade/fade-in, no playback-speed shift, no next-track preload. The
1411 path stays as small as possible to keep end-to-end latency low.
1412
1413 Fast path: when the source PCM format already matches the consumer's
1414 ``pcm_format``, the provider's bytes are paced in Python and forwarded
1415 directly â no ffmpeg in the data path.
1416
1417 Slow path: when formats differ, ffmpeg resamples/recodes the stream
1418 (with ``-re`` for rate pacing) via ``get_media_stream``.
1419
1420 :param queue_item: The AudioSource queue item to stream.
1421 :param pcm_format: Output PCM format the consumer wants.
1422 :param raise_on_error: Re-raise stream errors instead of swallowing them.
1423 """
1424 streamdetails = queue_item.streamdetails
1425 assert streamdetails
1426 logger = self.logger.getChild("audio_source_stream")
1427 bytes_received = 0
1428 try:
1429 async for chunk in self._iter_audio_source_pcm(streamdetails, pcm_format):
1430 bytes_received += len(chunk)
1431 yield chunk
1432 except AudioError as err:
1433 streamdetails.stream_error = True
1434 # revoke availability when the stream never produced any audio
1435 if bytes_received == 0:
1436 queue_item.available = False
1437 if raise_on_error:
1438 raise
1439 logger.error(
1440 "AudioError while streaming AudioSource %s (%s): %s",
1441 queue_item.name,
1442 streamdetails.uri,
1443 err,
1444 )
1445 except asyncio.CancelledError:
1446 raise
1447 except Exception as err:
1448 streamdetails.stream_error = True
1449 if raise_on_error:
1450 raise
1451 logger.exception(
1452 "Unexpected error while streaming AudioSource %s (%s): %s",
1453 queue_item.name,
1454 streamdetails.uri,
1455 err,
1456 )
1457 finally:
1458 streamdetails.seconds_streamed = bytes_received / pcm_format.pcm_sample_size
1459
1460 async def get_queue_item_stream(
1461 self,
1462 queue_item: QueueItem,
1463 pcm_format: AudioFormat,
1464 seek_position: int = 0,
1465 playback_speed: float = 1.0,
1466 raise_on_error: bool = True,
1467 normalization_override: VolumeNormalizationMode | None = None,
1468 session_id: str | None = None,
1469 ) -> AsyncGenerator[bytes]:
1470 """
1471 Get the (PCM) audio stream for a single queue item.
1472
1473 Audio is always served from the AudioBuffer which stores raw decoded PCM.
1474 Volume normalization and other filters are applied on-the-fly when reading
1475 from the buffer.
1476
1477 AudioSource items dispatch to ``get_audio_source_stream`` instead: they
1478 are realtime and bypass the buffering/normalization/filter machinery.
1479
1480 :param normalization_override: Force this volume normalization mode instead of
1481 re-evaluating it from the (possibly just-updated) loudness measurement. Used by
1482 the crossfade path to keep a track's replayed intro and its body on the same mode.
1483 :param session_id: Queue session that owns processing-detail updates.
1484 """
1485 streamdetails = queue_item.streamdetails
1486 assert streamdetails
1487
1488 # streamdetails are cached and reused for retries; reset this before any
1489 # media-type-specific dispatch so AudioSource failures do not stick.
1490 streamdetails.stream_error = False
1491
1492 if queue_item.media_type == MediaType.AUDIO_SOURCE:
1493 async for chunk in self.get_audio_source_stream(
1494 queue_item=queue_item,
1495 pcm_format=pcm_format,
1496 raise_on_error=raise_on_error,
1497 ):
1498 yield chunk
1499 return
1500 filter_params: list[str] = []
1501
1502 logger = self.logger.getChild("queue_item_stream")
1503
1504 if normalization_override is not None:
1505 # crossfade path pins the body to the intro's mode; skip hydration/re-eval that could flip it
1506 streamdetails.volume_normalization_mode = normalization_override
1507 else:
1508 # hydrate loudness from audio analysis (just-in-time, so that a measurement
1509 # completed during a previous play is picked up here). A live analyzer run
1510 # may have already populated streamdetails.loudness in memory â don't clobber
1511 # that, and don't clobber a value set upstream by the music provider.
1512 if streamdetails.loudness is None:
1513 if analysis := await self.mass.streams.audio_analysis.get_audio_analysis(
1514 streamdetails.item_id,
1515 streamdetails.provider,
1516 media_type=streamdetails.media_type,
1517 # use the authoritative EBU R128 value, not another provider's loudness proxy
1518 priority=(LOUDNESS_ANALYSIS_DOMAIN,),
1519 ):
1520 if analysis.loudness_integrated is not None:
1521 streamdetails.loudness = round(analysis.loudness_integrated, 2)
1522 if analysis.loudness_album is not None and streamdetails.loudness_album is None:
1523 streamdetails.loudness_album = round(analysis.loudness_album, 2)
1524
1525 # re-evaluate normalization mode: the background loudness analyzer may have
1526 # updated streamdetails.loudness since get_stream_details was called
1527 if streamdetails.queue_id:
1528 volume_normalization_enabled = (
1529 self.mass.config.get_effective_player_queue_config_value(
1530 streamdetails.queue_id, CONF_VOLUME_NORMALIZATION, CONF_VALUE_ENABLED
1531 )
1532 != CONF_VALUE_DISABLED
1533 )
1534 streamdetails.volume_normalization_mode = get_normalization_mode(
1535 self._get_volume_normalization_preference(streamdetails),
1536 volume_normalization_enabled,
1537 streamdetails,
1538 )
1539
1540 # handle volume normalization
1541 gain_correct: float | None = None
1542 if streamdetails.volume_normalization_mode == VolumeNormalizationMode.DYNAMIC:
1543 filter_rule = (
1544 f"loudnorm=I={streamdetails.target_loudness}"
1545 ":TP=-2.0:LRA=10.0:offset=0.0:print_format=json"
1546 )
1547 filter_params.append(filter_rule)
1548 elif streamdetails.volume_normalization_mode == VolumeNormalizationMode.FIXED_GAIN:
1549 config_key = (
1550 CONF_VOLUME_NORMALIZATION_FIXED_GAIN_TRACKS
1551 if streamdetails.media_type == MediaType.TRACK
1552 else CONF_VOLUME_NORMALIZATION_FIXED_GAIN_RADIO
1553 )
1554 gain_value = self.mass.streams.get_config_value(config_key, return_type=float)
1555 gain_correct = round(gain_value, 2)
1556 filter_params.append(f"volume={gain_correct}dB")
1557 elif streamdetails.volume_normalization_mode == VolumeNormalizationMode.MEASUREMENT_ONLY:
1558 target_loudness = (
1559 float(streamdetails.target_loudness)
1560 if streamdetails.target_loudness is not None
1561 else 0.0
1562 )
1563 if streamdetails.prefer_album_loudness and streamdetails.loudness_album is not None:
1564 gain_correct = target_loudness - float(streamdetails.loudness_album)
1565 elif streamdetails.loudness is not None:
1566 gain_correct = target_loudness - float(streamdetails.loudness)
1567 else:
1568 gain_correct = 0.0
1569 gain_correct = round(gain_correct, 2)
1570 filter_params.append(f"volume={gain_correct}dB")
1571 streamdetails.volume_normalization_gain_correct = gain_correct
1572
1573 # handle playback speed
1574 if playback_speed != 1.0:
1575 filter_params.append(f"atempo={playback_speed}")
1576
1577 # handle optional fade-in
1578 if streamdetails.fade_in:
1579 filter_params.insert(0, "afade=type=in:start_time=0:duration=3")
1580
1581 logger.log(
1582 VERBOSE_LOG_LEVEL,
1583 "Starting queue item stream for %s (%s)"
1584 " - using fade-in: %s"
1585 " - using volume normalization: %s"
1586 " - using playback speed: %s",
1587 queue_item.name,
1588 streamdetails.uri,
1589 streamdetails.fade_in,
1590 streamdetails.volume_normalization_mode,
1591 playback_speed,
1592 )
1593
1594 # get or create the AudioBuffer (stores raw decoded PCM)
1595 seek_position_ms = int(seek_position * 1000)
1596 audio_buffer = await AudioBuffer.get_buffer(
1597 mass=self.mass,
1598 streamdetails=streamdetails,
1599 seek_position_ms=seek_position_ms,
1600 reason="streaming",
1601 )
1602 if (
1603 streamdetails.queue_id
1604 and (queue_data := self.mass.player_queues.queue_data_or_none(streamdetails.queue_id))
1605 and (processing_session_id := session_id or queue_data.session_id)
1606 ):
1607 self.mass.streams.audio_processing.update_item_runtime(
1608 queue_id=streamdetails.queue_id,
1609 session_id=processing_session_id,
1610 queue_item_id=queue_item.queue_item_id,
1611 input_format=audio_buffer.pcm_format,
1612 pcm_format=pcm_format,
1613 normalization=get_normalization_details(streamdetails, gain_correct),
1614 playback_speed=playback_speed,
1615 alters_audio=streamdetails.fade_in,
1616 )
1617 # read from buffer with filters applied (volume normalization, speed, fade-in, etc.)
1618 # if no processing needed, this yields directly from the buffer
1619 media_stream_gen = audio_buffer.get_stream(
1620 output_format=pcm_format,
1621 seek_position_ms=seek_position_ms,
1622 filter_params=filter_params or None,
1623 )
1624
1625 first_chunk_received = False
1626 bytes_received = 0
1627 finished = False
1628 next_buffer_triggered = False
1629 stream_started_at = asyncio.get_event_loop().time()
1630 try:
1631 async for chunk in media_stream_gen:
1632 bytes_received += len(chunk)
1633 if not first_chunk_received:
1634 first_chunk_received = True
1635 logger.log(
1636 VERBOSE_LOG_LEVEL,
1637 "First audio chunk received for %s (%s) after %.2f seconds",
1638 queue_item.name,
1639 streamdetails.uri,
1640 asyncio.get_event_loop().time() - stream_started_at,
1641 )
1642 # trigger pre-buffering of the next item well before end
1643 # to ensure the raw PCM is ready when the next item needs to be streamed.
1644 # tracks and sound effects are finite files that fill and close immediately;
1645 # live sources (radio, audio_source) open an upstream connection that would
1646 # sit idle and likely time out before the player actually consumes it.
1647 if (
1648 not next_buffer_triggered
1649 and streamdetails.duration
1650 and (queue := self.mass.player_queues.get_active_queue(queue_item.queue_id))
1651 and queue.next_item
1652 and queue.next_item.queue_item_id != queue_item.queue_item_id
1653 and queue.next_item.media_type in (MediaType.TRACK, MediaType.SOUND_EFFECT)
1654 and (bytes_received / pcm_format.pcm_sample_size + seek_position)
1655 >= streamdetails.duration - 60
1656 ):
1657 next_buffer_triggered = True
1658 self.mass.player_queues.prepare_next_audio_buffer(queue_item.queue_id)
1659 yield chunk
1660 del chunk
1661 finished = True
1662 except AudioError as err:
1663 streamdetails.stream_error = True
1664 # revoke availability when the stream never produced any audio
1665 if bytes_received == 0:
1666 queue_item.available = False
1667 if raise_on_error:
1668 raise
1669 logger.error(
1670 "AudioError while streaming queue item %s (%s): %s",
1671 queue_item.name,
1672 streamdetails.uri,
1673 err,
1674 )
1675 except asyncio.CancelledError:
1676 raise
1677 except Exception as err:
1678 streamdetails.stream_error = True
1679 if raise_on_error:
1680 raise
1681 logger.exception(
1682 "Unexpected error while streaming queue item %s (%s): %s",
1683 queue_item.name,
1684 streamdetails.uri,
1685 err,
1686 )
1687 finally:
1688 seconds_streamed = bytes_received / pcm_format.pcm_sample_size
1689 streamdetails.seconds_streamed = seconds_streamed
1690 logger.log(
1691 VERBOSE_LOG_LEVEL,
1692 "stream %s for %s in %.2f seconds - seconds streamed/buffered: %.2f",
1693 "aborted" if not finished else "finished",
1694 streamdetails.uri,
1695 asyncio.get_event_loop().time() - stream_started_at,
1696 seconds_streamed,
1697 )
1698 self._notify_provider_streamed(streamdetails, finished, seconds_streamed)
1699
1700 async def get_queue_item_stream_with_smartfade(
1701 self,
1702 player: Player,
1703 queue_item: QueueItem,
1704 pcm_format: AudioFormat,
1705 crossfade_mode: CrossfadeMode = CrossfadeMode.SMART_CROSSFADE,
1706 standard_crossfade_duration: int = 10,
1707 session_id: str | None = None,
1708 ) -> AsyncGenerator[bytes]:
1709 """
1710 Return one queue item with a crossfade into the next item.
1711
1712 :param player: Player consuming the stream.
1713 :param queue_item: Queue item to stream.
1714 :param pcm_format: Shared PCM format.
1715 :param crossfade_mode: Effective crossfade mode.
1716 :param standard_crossfade_duration: Configured standard crossfade duration.
1717 :param session_id: Queue session that owns processing-detail updates.
1718 """
1719 queue = self.mass.player_queues.get(queue_item.queue_id)
1720 if not queue:
1721 raise RuntimeError(f"Queue {queue_item.queue_id} not found")
1722
1723 streamdetails = queue_item.streamdetails
1724 assert streamdetails
1725 crossfade_data = self._crossfade_data.get(queue.queue_id)
1726
1727 if crossfade_data and streamdetails.seek_position > 0:
1728 # don't do crossfade when seeking into track
1729 self.logger.debug(
1730 "Discarding crossfade data for queue %s - seeking into track (pos=%s)",
1731 queue.display_name,
1732 streamdetails.seek_position,
1733 )
1734 crossfade_data = None
1735 if crossfade_data and (crossfade_data.queue_item_id != queue_item.queue_item_id):
1736 # edge case alert: the next item changed just while we were preloading/crossfading
1737 self.logger.warning(
1738 "Skipping crossfade data for queue %s - next item changed!"
1739 " (expected queue_item_id=%s, got=%s)",
1740 queue.display_name,
1741 crossfade_data.queue_item_id,
1742 queue_item.queue_item_id,
1743 )
1744 crossfade_data = None
1745 self._crossfade_data.pop(queue.queue_id, None)
1746 elif not crossfade_data:
1747 self.logger.debug(
1748 "No crossfade data available for queue %s (queue_item_id=%s)",
1749 queue.display_name,
1750 queue_item.queue_item_id,
1751 )
1752
1753 self.logger.debug(
1754 "Start Streaming queue track: %s (%s) for queue %s on player %s"
1755 "- crossfade mode: %s "
1756 "- crossfading from previous track: %s ",
1757 queue_item.streamdetails.uri if queue_item.streamdetails else "Unknown URI",
1758 queue_item.name,
1759 queue.display_name,
1760 player.name,
1761 crossfade_mode,
1762 "true" if crossfade_data else "false",
1763 )
1764
1765 buffer = bytearray()
1766 bytes_written = 0
1767 # calculate crossfade buffer size
1768 crossfade_buffer_duration = (
1769 SMART_CROSSFADE_DURATION
1770 if crossfade_mode == CrossfadeMode.SMART_CROSSFADE
1771 else standard_crossfade_duration
1772 )
1773 crossfade_buffer_duration = min(
1774 crossfade_buffer_duration,
1775 int(streamdetails.duration / 2)
1776 if streamdetails.duration
1777 else crossfade_buffer_duration,
1778 )
1779 # skip crossfade if buffer would be too small to be meaningful
1780 if crossfade_buffer_duration < 5:
1781 crossfade_buffer_duration = 0
1782 # Ensure crossfade buffer size is aligned to frame boundaries
1783 # Frame size = bytes_per_sample * channels
1784 bytes_per_sample = pcm_format.bit_depth // 8
1785 frame_size = bytes_per_sample * pcm_format.channels
1786 crossfade_buffer_size = int(pcm_format.pcm_sample_size * crossfade_buffer_duration)
1787 # Round down to nearest frame boundary
1788 crossfade_buffer_size = (crossfade_buffer_size // frame_size) * frame_size
1789 fade_out_data: bytes | None = None
1790 uncredited_tail_bytes = 0
1791
1792 # pin the body to DYNAMIC when the intro was baked DYNAMIC,
1793 # else a late measurement flips it and causes a volume jump
1794 norm_override: VolumeNormalizationMode | None = None
1795 if crossfade_data and crossfade_data.normalization_mode == VolumeNormalizationMode.DYNAMIC:
1796 norm_override = VolumeNormalizationMode.DYNAMIC
1797
1798 if crossfade_data:
1799 # reported media-time (TRIM + CF) is decoupled from the raw buffer seek below (X)
1800 streamdetails.seek_position = crossfade_data.elapsed_time_offset
1801 # yield the POST portion (resample if previous track's format differs)
1802 if crossfade_data.pcm_format != pcm_format:
1803 async for _chunk in resample_pcm_audio(
1804 crossfade_data.data, crossfade_data.pcm_format, pcm_format
1805 ):
1806 yield _chunk
1807 bytes_written += len(_chunk)
1808 else:
1809 for pcm_slice in iter_pcm_slices(crossfade_data.data, pcm_format, 1000):
1810 yield pcm_slice
1811 await asyncio.sleep(0)
1812 bytes_written += len(crossfade_data.data)
1813 # skip past the audio already consumed by the crossfade
1814 fade_in_duration_seconds = (
1815 crossfade_data.fade_in_size / crossfade_data.fade_in_pcm_format.pcm_sample_size
1816 )
1817 discard_seconds = int(fade_in_duration_seconds)
1818 fractional_seconds = fade_in_duration_seconds - discard_seconds
1819 discard_leftover = int(fractional_seconds * pcm_format.pcm_sample_size)
1820 discard_leftover = (discard_leftover // frame_size) * frame_size
1821 crossfade_data = None
1822 self._crossfade_data.pop(queue.queue_id, None)
1823 else:
1824 discard_seconds = int(streamdetails.seek_position)
1825 discard_leftover = 0
1826
1827 # Yield the first WARMUP_DURATION worth of audio immediately so playback starts
1828 # right away. After that, start accumulating the crossfade holdback buffer.
1829 warmup_size = int(pcm_format.pcm_sample_size * WARMUP_DURATION)
1830 warmup_bytes = 0
1831 total_chunks_received = 0
1832 playback_speed = cast("float", queue_item.extra_attributes.get("playback_speed", 1.0))
1833 async for chunk in self.get_queue_item_stream(
1834 queue_item,
1835 pcm_format,
1836 seek_position=discard_seconds,
1837 playback_speed=playback_speed,
1838 normalization_override=norm_override,
1839 session_id=session_id,
1840 ):
1841 total_chunks_received += 1
1842 if discard_leftover:
1843 chunk = chunk[discard_leftover:] # noqa: PLW2901
1844 discard_leftover = 0
1845
1846 if warmup_bytes < warmup_size:
1847 # warmup: yield directly, don't buffer
1848 yield chunk
1849 warmup_bytes += len(chunk)
1850 bytes_written += len(chunk)
1851 del chunk
1852 continue
1853
1854 buffer.extend(chunk)
1855 del chunk
1856 if len(buffer) < crossfade_buffer_size:
1857 await asyncio.sleep(0)
1858 continue
1859 # yield everything above the crossfade buffer
1860 while len(buffer) > crossfade_buffer_size:
1861 yield bytes(buffer[: pcm_format.pcm_sample_size])
1862 bytes_written += pcm_format.pcm_sample_size
1863 del buffer[: pcm_format.pcm_sample_size]
1864 await asyncio.sleep(0)
1865
1866 #### HANDLE END OF TRACK
1867
1868 # get next track for crossfade
1869 crossfade_start_time = asyncio.get_event_loop().time()
1870 next_queue_item: QueueItem | None
1871 try:
1872 self.logger.debug(
1873 "Preloading NEXT track for crossfade for queue %s", queue.display_name
1874 )
1875 next_queue_item = await self.mass.player_queues.load_next_queue_item(
1876 queue.queue_id, queue_item.queue_item_id
1877 )
1878 # set index_in_buffer to prevent our next track is overwritten while preloading
1879 if next_queue_item.streamdetails is None:
1880 raise InvalidDataError(
1881 f"No streamdetails for next queue item {next_queue_item.queue_item_id}"
1882 )
1883 queue.index_in_buffer = self.mass.player_queues.index_by_id(
1884 queue.queue_id, next_queue_item.queue_item_id
1885 )
1886 except QueueEmpty:
1887 # end of queue reached, no next item
1888 next_queue_item = None
1889
1890 crossfade_allowed = False
1891 if next_queue_item and next_queue_item.streamdetails:
1892 next_pcm = await self.select_pcm_format(
1893 player=player,
1894 streamdetails=next_queue_item.streamdetails,
1895 crossfade_enabled=True,
1896 )
1897 crossfade_allowed = self.crossfade_allowed(
1898 queue_item,
1899 crossfade_mode=crossfade_mode,
1900 player_id=player.player_id,
1901 flow_mode=False,
1902 next_queue_item=next_queue_item,
1903 sample_rate=pcm_format.sample_rate,
1904 next_sample_rate=next_pcm.sample_rate,
1905 )
1906 if not crossfade_allowed:
1907 # no crossfade enabled/allowed, just yield the buffer last part
1908 bytes_written += len(buffer)
1909 for pcm_slice in iter_pcm_slices(bytes(buffer), pcm_format, 1000):
1910 yield pcm_slice
1911 await asyncio.sleep(0)
1912 else:
1913 assert next_queue_item is not None
1914 # the remaining buffer is the fade-out tail of the current track
1915 fade_out_data = bytes(buffer)
1916 buffer = bytearray()
1917 # initialized before the try block â the except handler reads these
1918 first_part_written = 0
1919 second_part_buf = bytearray()
1920 try:
1921 # wrap the next track's stream in a counting generator that caps
1922 # at crossfade_buffer_size and tracks how many bytes were consumed
1923 fade_in_bytes_consumed = 0
1924
1925 _next_item = next_queue_item
1926
1927 async def _limited_fade_in() -> AsyncGenerator[bytes]:
1928 nonlocal fade_in_bytes_consumed
1929 async for chunk in self.get_queue_item_stream(
1930 _next_item,
1931 pcm_format,
1932 playback_speed=cast(
1933 "float", _next_item.extra_attributes.get("playback_speed", 1.0)
1934 ),
1935 session_id=session_id,
1936 ):
1937 remaining = crossfade_buffer_size - fade_in_bytes_consumed
1938 if remaining <= 0:
1939 break
1940 if len(chunk) > remaining:
1941 fade_in_bytes_consumed += remaining
1942 yield chunk[:remaining]
1943 break
1944 fade_in_bytes_consumed += len(chunk)
1945 yield chunk
1946
1947 smart_fade = await self.smart_fades_mixer.build(
1948 fade_in_streamdetails=cast("StreamDetails", next_queue_item.streamdetails),
1949 fade_out_streamdetails=streamdetails,
1950 pcm_format=pcm_format,
1951 standard_crossfade_duration=standard_crossfade_duration,
1952 mode=crossfade_mode,
1953 fade_out_data=fade_out_data,
1954 fade_in_bytes_len=crossfade_buffer_size,
1955 )
1956 crossfade_timing = smart_fade.timing_info
1957 # Split mix output at end-of-overlap: PRE+CF to A, POST to B's intro.
1958 fadeout_share_bytes = int(
1959 (crossfade_timing.pre_crossfade_duration + crossfade_timing.crossfade_duration)
1960 * pcm_format.pcm_sample_size
1961 )
1962 fadeout_share_bytes = (fadeout_share_bytes // frame_size) * frame_size
1963 async for mix_chunk in self.smart_fades_mixer.mix(
1964 smart_fade,
1965 fade_in_part=_limited_fade_in(),
1966 fade_out_part=fade_out_data,
1967 pcm_format=pcm_format,
1968 ):
1969 if first_part_written < fadeout_share_bytes:
1970 # split this chunk so A gets exactly fadeout_share_bytes
1971 remaining = fadeout_share_bytes - first_part_written
1972 if len(mix_chunk) > remaining:
1973 yield mix_chunk[:remaining]
1974 first_part_written += remaining
1975 bytes_written += remaining
1976 second_part_buf.extend(mix_chunk[remaining:])
1977 else:
1978 yield mix_chunk
1979 first_part_written += len(mix_chunk)
1980 bytes_written += len(mix_chunk)
1981 else:
1982 second_part_buf.extend(mix_chunk)
1983 # tail consumed by the mix but not credited to bytes_written
1984 uncredited_tail_bytes = len(fade_out_data) - first_part_written
1985 self._crossfade_data[queue_item.queue_id] = CrossfadeData(
1986 data=bytes(second_part_buf),
1987 fade_in_size=fade_in_bytes_consumed,
1988 pcm_format=pcm_format,
1989 fade_in_pcm_format=pcm_format,
1990 queue_item_id=next_queue_item.queue_item_id,
1991 elapsed_time_offset=(
1992 crossfade_timing.fadein_trimmed_duration
1993 + crossfade_timing.crossfade_duration
1994 ),
1995 normalization_mode=cast(
1996 "StreamDetails", next_queue_item.streamdetails
1997 ).volume_normalization_mode,
1998 )
1999 crossfade_elapsed = asyncio.get_event_loop().time() - crossfade_start_time
2000 self.logger.debug(
2001 "Stored crossfade data for queue %s"
2002 " - next queue_item_id: %s (preparation took %.1fs)",
2003 queue.display_name,
2004 next_queue_item.queue_item_id,
2005 crossfade_elapsed,
2006 )
2007 except Exception as err:
2008 if first_part_written or second_part_buf:
2009 # partial mix already played â concat'd fade_out_data would duplicate audio
2010 raise
2011 # crossfade failed, fall back to just yielding the fade_out_data
2012 self.logger.warning(
2013 "Crossfade failed for queue %s: %s",
2014 queue.display_name,
2015 err,
2016 )
2017 next_queue_item = None
2018 for pcm_slice in iter_pcm_slices(fade_out_data, pcm_format, 1000):
2019 yield pcm_slice
2020 await asyncio.sleep(0)
2021 bytes_written += len(fade_out_data)
2022 del fade_out_data
2023 # make sure the buffer gets cleaned up
2024 del buffer
2025 # update duration details based on the actual pcm data we sent
2026 # this also accounts for crossfade and silence stripping
2027 seconds_streamed = bytes_written / pcm_format.pcm_sample_size
2028 streamdetails.seconds_streamed = seconds_streamed
2029 uncredited_tail_seconds = uncredited_tail_bytes / pcm_format.pcm_sample_size
2030 # streamdetails.duration is in media-time; seconds_streamed is stream-time
2031 # (post-atempo), so we scale by playback_speed to recover media-time.
2032 streamdetails.duration = int(
2033 streamdetails.seek_position
2034 + (seconds_streamed + uncredited_tail_seconds) * playback_speed
2035 )
2036 # propagate accurate duration to queue_item so UI displays it
2037 queue_item.duration = streamdetails.duration
2038 self.logger.debug(
2039 "Finished Streaming queue track: %s (%s) on queue %s "
2040 "- crossfade data prepared for next track: %s",
2041 streamdetails.uri,
2042 queue_item.name,
2043 queue.display_name,
2044 next_queue_item.name if next_queue_item else "N/A",
2045 )
2046
2047 async def get_queue_flow_stream(
2048 self,
2049 queue: PlayerQueue,
2050 start_queue_item: QueueItem,
2051 pcm_format: AudioFormat,
2052 session_id: str | None = None,
2053 protocol_player: Player | None = None,
2054 ) -> AsyncGenerator[bytes]:
2055 """
2056 Get a flow stream of all tracks in the queue as raw PCM audio.
2057
2058 yields chunks of exactly 1 second of audio in the given pcm_format.
2059
2060 :param queue: Queue being streamed.
2061 :param start_queue_item: First queue item in the flow stream.
2062 :param pcm_format: Shared PCM format for the complete flow stream.
2063 :param session_id: Queue session that owns processing-detail updates.
2064 :param protocol_player: The protocol player actually consuming the flow stream.
2065 Must be the same player that was used to select ``pcm_format`` so
2066 restart decisions are made against the correct supported sample rates
2067 and flow mode configuration. Falls back to the queue's player when omitted.
2068 """
2069 # ruff: noqa: PLR0915
2070 assert pcm_format.content_type.is_pcm()
2071 queue_track = None
2072 last_fadeout_part: bytes = b""
2073 last_streamdetails: StreamDetails | None = None
2074 last_play_log_entry: PlayLogEntry | None = None
2075 # Snapshot the queue's current session_id. PlayerQueues rotates this on
2076 # every new stream session, so if a newer producer takes over the queue
2077 # (rapid track switch, sync-group reform, dynamic leader handoff) the
2078 # snapshot will no longer match and we exit cleanly on the next yield or
2079 # playlog append â preventing two producers from writing to the same
2080 # pq_data.flow_mode_stream_log.
2081 pq_data = self.mass.player_queues.queue_data(queue.queue_id)
2082 flow_session_id = session_id or pq_data.session_id
2083 if flow_session_id is None or pq_data.session_id != flow_session_id:
2084 self.logger.debug(
2085 "Ignoring stale flow stream for queue %s (session %s, active %s)",
2086 queue.display_name,
2087 flow_session_id,
2088 pq_data.session_id,
2089 )
2090 return
2091 queue.flow_mode = True
2092 pq_data.flow_mode_stream_log = []
2093 if not start_queue_item:
2094 # this can happen in some (edge case) race conditions
2095 return
2096 pcm_sample_size = pcm_format.pcm_sample_size
2097 if start_queue_item.media_type != MediaType.TRACK:
2098 # no crossfade on non-tracks
2099 crossfade_mode = CrossfadeMode.DISABLED
2100 standard_crossfade_duration = 0
2101 else:
2102 crossfade_mode = self.mass.streams.get_crossfade_mode(queue)
2103 # crossfade duration is a global (queue controller) setting; fallback matches
2104 # CONF_ENTRY_CROSSFADE_DURATION's default
2105 standard_crossfade_duration = self.mass.config.get_raw_core_config_value(
2106 CONF_PLAYER_QUEUES, CONF_CROSSFADE_DURATION, 8
2107 )
2108 flow_mode_sample_rate_conf, flow_supported_sample_rates = self._flow_restart_context(
2109 queue.queue_id, protocol_player
2110 )
2111 # note: get_crossfade_mode() already falls back to standard when smart fades aren't
2112 # available (no analysis provider / minimal buffer), so crossfade_mode is safe to use.
2113 self.logger.info(
2114 "Start Queue Flow stream for Queue %s - crossfade: %s %s",
2115 queue.display_name,
2116 crossfade_mode,
2117 f"({standard_crossfade_duration}s)"
2118 if crossfade_mode == CrossfadeMode.STANDARD_CROSSFADE
2119 else "",
2120 )
2121 total_chunks_received = 0
2122
2123 def _superseded() -> bool:
2124 """Return True if a newer stream session has taken over this queue."""
2125 return pq_data.session_id != flow_session_id
2126
2127 queue_exhausted = False
2128 while True:
2129 # bail out early if a newer producer has taken over this queue,
2130 # so we don't append another entry to a stream log we no longer own
2131 if _superseded():
2132 self.logger.debug(
2133 "Flow stream for queue %s superseded (session %s -> %s) "
2134 "- exiting before next track",
2135 queue.display_name,
2136 flow_session_id,
2137 pq_data.session_id,
2138 )
2139 return
2140 # get (next) queue item to stream
2141 if queue_track is None:
2142 queue_track = start_queue_item
2143 else:
2144 try:
2145 queue_track = await self.mass.player_queues.load_next_queue_item(
2146 queue.queue_id, queue_track.queue_item_id
2147 )
2148 except QueueEmpty:
2149 queue_exhausted = True
2150 break
2151
2152 if self._flow_stream_needs_restart(
2153 queue_track,
2154 pcm_format,
2155 flow_supported_sample_rates,
2156 flow_mode_sample_rate_conf,
2157 is_first_track=queue_track is start_queue_item,
2158 ):
2159 break
2160
2161 if queue_track.streamdetails is None:
2162 self.logger.error(
2163 "No StreamDetails for queue item %s (%s) on queue %s - skipping track",
2164 queue_track.queue_item_id,
2165 queue_track.name,
2166 queue.display_name,
2167 )
2168 continue
2169 if flow_session_id is not None:
2170 self.mass.streams.audio_processing.update_item_context(
2171 queue_id=queue.queue_id,
2172 session_id=flow_session_id,
2173 queue_item_id=queue_track.queue_item_id,
2174 queue_processing=AudioQueueProcessing(
2175 pcm_format=pcm_format,
2176 playback_speed=cast(
2177 "float",
2178 queue_track.extra_attributes.get("playback_speed", 1.0),
2179 ),
2180 crossfade_mode=crossfade_mode,
2181 overlay_active=overlay_active(queue),
2182 ),
2183 alters_audio=queue_track.streamdetails.fade_in,
2184 )
2185
2186 self.logger.debug(
2187 "Start Streaming queue track: %s (%s) for queue %s",
2188 queue_track.streamdetails.uri,
2189 queue_track.name,
2190 queue.display_name,
2191 )
2192 # last chance to bail before mutating the stream log: a newer producer
2193 # may have taken over while we were awaiting load_next_queue_item
2194 if _superseded():
2195 self.logger.debug(
2196 "Flow stream for queue %s superseded - exiting before playlog append",
2197 queue.display_name,
2198 )
2199 return
2200 track_playback_speed = cast(
2201 "float", queue_track.extra_attributes.get("playback_speed", 1.0)
2202 )
2203 # calculate crossfade buffer size
2204 crossfade_buffer_duration = (
2205 SMART_CROSSFADE_DURATION
2206 if crossfade_mode == CrossfadeMode.SMART_CROSSFADE
2207 else standard_crossfade_duration
2208 )
2209 crossfade_buffer_duration = min(
2210 crossfade_buffer_duration,
2211 int(queue_track.streamdetails.duration / 2)
2212 if queue_track.streamdetails.duration
2213 else crossfade_buffer_duration,
2214 )
2215 # skip crossfade if buffer would be too small to be meaningful
2216 if crossfade_buffer_duration < 5:
2217 crossfade_buffer_duration = 0
2218 # Ensure crossfade buffer size is aligned to frame boundaries
2219 # Frame size = bytes_per_sample * channels
2220 bytes_per_sample = pcm_format.bit_depth // 8
2221 frame_size = bytes_per_sample * pcm_format.channels
2222 crossfade_buffer_size = int(pcm_format.pcm_sample_size * crossfade_buffer_duration)
2223 # Round down to nearest frame boundary
2224 crossfade_buffer_size = (crossfade_buffer_size // frame_size) * frame_size
2225 warmup_size = int(pcm_format.pcm_sample_size * WARMUP_DURATION)
2226
2227 # raw_seek_position feeds the PCM buffer; streamdetails.seek_position
2228 # (overwritten below) only drives reported elapsed time.
2229 raw_seek_position = queue_track.streamdetails.seek_position
2230 # Build eagerly so seek_position is set before PlayLogEntry is appended â
2231 # consumer-paced mix() would otherwise let the queue briefly report 0.
2232 crossfade_smart_fade: SmartFade | None = None
2233 if (
2234 last_fadeout_part
2235 and last_streamdetails
2236 and crossfade_buffer_size > 0
2237 and crossfade_mode != CrossfadeMode.DISABLED
2238 ):
2239 crossfade_smart_fade = await self.smart_fades_mixer.build(
2240 fade_in_streamdetails=queue_track.streamdetails,
2241 fade_out_streamdetails=last_streamdetails,
2242 pcm_format=pcm_format,
2243 standard_crossfade_duration=standard_crossfade_duration,
2244 mode=crossfade_mode,
2245 fade_out_data=last_fadeout_part,
2246 fade_in_bytes_len=crossfade_buffer_size,
2247 )
2248 timing_info = crossfade_smart_fade.timing_info
2249 queue_track.streamdetails.seek_position = (
2250 raw_seek_position
2251 + timing_info.fadein_trimmed_duration
2252 + timing_info.crossfade_duration
2253 )
2254 # append to play log so the queue controller can work out which track is playing
2255 play_log_entry = PlayLogEntry(queue_track.queue_item_id)
2256 pq_data.flow_mode_stream_log.append(play_log_entry)
2257
2258 bytes_written = 0
2259 crossfade_buffer = bytearray()
2260 warmup_bytes = 0
2261 first_chunk_received = False
2262
2263 async for chunk in self.get_queue_item_stream(
2264 queue_track,
2265 pcm_format=pcm_format,
2266 seek_position=int(raw_seek_position),
2267 playback_speed=cast(
2268 "float", queue_track.extra_attributes.get("playback_speed", 1.0)
2269 ),
2270 raise_on_error=False,
2271 session_id=flow_session_id,
2272 ):
2273 # if a newer producer has taken over this queue, stop sending
2274 # audio and exit cleanly before the outer-loop end-of-track
2275 # bookkeeping mutates seconds_streamed / duration on the log
2276 if _superseded():
2277 self.logger.debug(
2278 "Flow stream for queue %s superseded - stopping chunk yield",
2279 queue.display_name,
2280 )
2281 return
2282 total_chunks_received += 1
2283 if not first_chunk_received:
2284 first_chunk_received = True
2285 # inform the queue that the track is now loaded in the buffer
2286 # so the next track can be preloaded
2287 self.mass.player_queues.track_loaded_in_buffer(
2288 queue.queue_id, queue_track.queue_item_id
2289 )
2290
2291 if crossfade_mode == CrossfadeMode.DISABLED:
2292 # no cross/smart fade: yield chunks directly without intermediate buffer
2293 yield chunk
2294 bytes_written += len(chunk)
2295 del chunk
2296 continue
2297
2298 # Warmup: yield chunks directly until we have streamed WARMUP_DURATION
2299 # worth of audio, so playback starts immediately. Skip warmup when
2300 # crossfade data from the previous track is pending â we need a full
2301 # buffer for the mix.
2302 if warmup_bytes < warmup_size and not last_fadeout_part:
2303 yield chunk
2304 warmup_bytes += len(chunk)
2305 bytes_written += len(chunk)
2306 del chunk
2307 continue
2308
2309 # smart fades enabled: accumulate chunks in crossfade buffer
2310 crossfade_buffer.extend(chunk)
2311 del chunk
2312 if len(crossfade_buffer) < crossfade_buffer_size:
2313 await asyncio.sleep(0)
2314 continue
2315
2316 # handle crossfade of previous track and new track
2317 if (
2318 last_fadeout_part
2319 and last_streamdetails
2320 and crossfade_smart_fade is not None
2321 and last_play_log_entry is not None
2322 ):
2323 fadein_part = bytes(crossfade_buffer[:crossfade_buffer_size])
2324 remaining_bytes = bytes(crossfade_buffer[crossfade_buffer_size:])
2325 try:
2326 crossfade_bytes_written = 0
2327 async for mix_chunk in self.smart_fades_mixer.mix(
2328 crossfade_smart_fade,
2329 fade_in_part=fadein_part,
2330 fade_out_part=last_fadeout_part,
2331 pcm_format=pcm_format,
2332 ):
2333 yield mix_chunk
2334 crossfade_bytes_written += len(mix_chunk)
2335 except Exception as mix_err:
2336 if crossfade_bytes_written:
2337 # partial mix already played â concat'd tail would duplicate audio
2338 raise
2339 self.logger.warning(
2340 "Crossfade mixer failed for %s, falling back to simple concat: %s",
2341 queue_track.name,
2342 mix_err,
2343 )
2344 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2345 yield pcm_slice
2346 await asyncio.sleep(0)
2347 # full tail was pre-counted and is now yielded as-is
2348 crossfade_bytes_written = 0
2349 remaining_bytes = bytes(crossfade_buffer)
2350 # mix failed â undo the eager seek_position
2351 queue_track.streamdetails.seek_position = raw_seek_position
2352 if crossfade_bytes_written:
2353 # Split mix output at end-of-overlap: PRE+CF to A, POST to B.
2354 fadeout_share_seconds = (
2355 timing_info.pre_crossfade_duration + timing_info.crossfade_duration
2356 )
2357 fadeout_share = int(fadeout_share_seconds * pcm_sample_size)
2358 fadeout_share = (fadeout_share // frame_size) * frame_size
2359 fadeout_share = min(fadeout_share, crossfade_bytes_written)
2360 fadein_share = crossfade_bytes_written - fadeout_share
2361 bytes_written += fadein_share
2362 if last_play_log_entry:
2363 assert last_play_log_entry.seconds_streamed is not None
2364 # correct pre-counted full tail to the timing-based share
2365 last_play_log_entry.seconds_streamed += (
2366 fadeout_share - len(last_fadeout_part)
2367 ) / pcm_sample_size
2368 if remaining_bytes:
2369 for pcm_slice in iter_pcm_slices(remaining_bytes, pcm_format, 1000):
2370 yield pcm_slice
2371 await asyncio.sleep(0)
2372 bytes_written += len(remaining_bytes)
2373 del remaining_bytes
2374 last_fadeout_part = b""
2375 last_streamdetails = None
2376 crossfade_buffer = bytearray()
2377 warmup_bytes = 0
2378
2379 # yield everything above the crossfade buffer size
2380 while len(crossfade_buffer) > crossfade_buffer_size:
2381 yield bytes(crossfade_buffer[:pcm_sample_size])
2382 bytes_written += pcm_sample_size
2383 del crossfade_buffer[:pcm_sample_size]
2384 await asyncio.sleep(0)
2385
2386 # A source error after partial audio must not look like a completed item.
2387 # Progress reporting skips items with stream_error, so the item is not
2388 # marked played; move on to the next queue item like the zero-audio path.
2389 if first_chunk_received and queue_track.streamdetails.stream_error:
2390 if _superseded():
2391 return
2392 self.logger.warning(
2393 "Track %s (%s) on queue %s aborted by a stream error - skipping",
2394 queue_track.name,
2395 queue_track.streamdetails.uri,
2396 queue.display_name,
2397 )
2398 # the audio sent so far will still play out; keep the play log entry
2399 # honest about how much of this item was actually streamed
2400 play_log_entry.seconds_streamed = bytes_written / pcm_sample_size
2401 if last_fadeout_part:
2402 # crossfade into this item never happened â undo the eager seek_position
2403 queue_track.streamdetails.seek_position = raw_seek_position
2404 continue
2405
2406 #### HANDLE END OF TRACK
2407 if not first_chunk_received:
2408 self.logger.warning(
2409 "Track %s (%s) on queue %s produced no audio data - skipping",
2410 queue_track.name,
2411 queue_track.streamdetails.uri if queue_track.streamdetails else "unknown",
2412 queue.display_name,
2413 )
2414 queue_track.streamdetails.stream_error = True
2415 play_log_entry.seconds_streamed = 0
2416 if last_fadeout_part:
2417 queue_track.streamdetails.seek_position = raw_seek_position
2418 continue
2419 if last_fadeout_part:
2420 # edge case: we did not get enough data to make the crossfade
2421 # attribute these bytes to the previous track (they are its tail)
2422 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2423 yield pcm_slice
2424 await asyncio.sleep(0)
2425 # no crossfade happened â undo the eager seek_position
2426 queue_track.streamdetails.seek_position = raw_seek_position
2427 # full tail was pre-counted and is now yielded as-is
2428 last_fadeout_part = b""
2429 if self.crossfade_allowed(
2430 queue_track,
2431 crossfade_mode=crossfade_mode,
2432 player_id=queue.queue_id,
2433 flow_mode=True,
2434 ):
2435 last_fadeout_part = bytes(crossfade_buffer[-crossfade_buffer_size:])
2436 last_streamdetails = queue_track.streamdetails
2437 last_play_log_entry = play_log_entry
2438 remaining_bytes = bytes(crossfade_buffer[:-crossfade_buffer_size])
2439 if remaining_bytes:
2440 for pcm_slice in iter_pcm_slices(remaining_bytes, pcm_format, 1000):
2441 yield pcm_slice
2442 await asyncio.sleep(0)
2443 bytes_written += len(remaining_bytes)
2444 del remaining_bytes
2445 elif crossfade_mode != CrossfadeMode.DISABLED and crossfade_buffer:
2446 bytes_written += len(crossfade_buffer)
2447 for pcm_slice in iter_pcm_slices(bytes(crossfade_buffer), pcm_format, 1000):
2448 yield pcm_slice
2449 await asyncio.sleep(0)
2450 crossfade_buffer = bytearray()
2451
2452 # update duration details based on the actual pcm data we sent
2453 # this also accounts for crossfade and silence stripping
2454 seconds_streamed = bytes_written / pcm_sample_size
2455 queue_track.streamdetails.seconds_streamed = seconds_streamed
2456 # the held-back crossfade tail still counts as this track's media-time
2457 tail_seconds = len(last_fadeout_part) / pcm_sample_size
2458 # streamdetails.duration is in media-time; seconds_streamed is stream-time
2459 # (post-atempo), so we scale by the track's playback_speed to recover media-time.
2460 queue_track.streamdetails.duration = int(
2461 queue_track.streamdetails.seek_position
2462 + (seconds_streamed + tail_seconds) * track_playback_speed
2463 )
2464 # propagate accurate duration to queue_item so UI displays it
2465 queue_track.duration = queue_track.streamdetails.duration
2466 play_log_entry.seconds_streamed = seconds_streamed
2467 play_log_entry.duration = queue_track.streamdetails.duration
2468 if last_play_log_entry is play_log_entry and last_fadeout_part:
2469 # Pre-count the full crossfade tail so the queue index calculation
2470 # doesn't undercount while waiting for the next track's crossfade mix.
2471 # This will be corrected to crossfade_total/2 once the mix completes.
2472 assert play_log_entry.seconds_streamed is not None
2473 play_log_entry.seconds_streamed += len(last_fadeout_part) / pcm_sample_size
2474 self.logger.debug(
2475 "Finished Streaming queue track: %s (%s) on queue %s",
2476 queue_track.streamdetails.uri,
2477 queue_track.name,
2478 queue.display_name,
2479 )
2480 #### HANDLE END OF QUEUE FLOW STREAM
2481 # skip end-of-queue bookkeeping if a newer producer has superseded us;
2482 # the new producer owns queue_buffer_completed and the play log now
2483 if _superseded():
2484 self.logger.debug(
2485 "Flow stream for queue %s superseded - skipping end-of-queue handling",
2486 queue.display_name,
2487 )
2488 return
2489 # end of queue flow: make sure we yield the last_fadeout_part
2490 if last_fadeout_part:
2491 for pcm_slice in iter_pcm_slices(last_fadeout_part, pcm_format, 1000):
2492 yield pcm_slice
2493 await asyncio.sleep(0)
2494 # correct seconds streamed - the duration already includes the tail
2495 last_part_seconds = len(last_fadeout_part) / pcm_sample_size
2496 streamdetails = queue_track.streamdetails
2497 assert streamdetails is not None
2498 streamdetails.seconds_streamed = (
2499 streamdetails.seconds_streamed or 0
2500 ) + last_part_seconds
2501 # also update the play log entry so elapsed time tracking stays in sync
2502 if last_play_log_entry:
2503 assert last_play_log_entry.seconds_streamed is not None
2504 # full tail was pre-counted and is now yielded as-is
2505 last_play_log_entry.duration = streamdetails.duration
2506 last_fadeout_part = b""
2507 self.logger.info("Finished Queue Flow stream for Queue %s", queue.display_name)
2508 # only signal completion if we are still the active producer â a later
2509 # producer would (incorrectly) see this as its own completion otherwise
2510 if not _superseded():
2511 # inform the queue controller that all audio data has been generated
2512 # so it can handle the case where new items were added after the flow stream ended
2513 self.mass.player_queues.queue_buffer_completed(queue.queue_id, queue_exhausted)
2514
2515 async def get_overlay_mixed_stream(
2516 self,
2517 queue: PlayerQueue,
2518 audio_input: AsyncGenerator[bytes],
2519 pcm_format: AudioFormat,
2520 ) -> AsyncGenerator[bytes]:
2521 """
2522 Mix the queue's audio overlay (looping sound effect) into the given PCM stream.
2523
2524 The mixed output has the exact same PCM format, duration and chunking as the
2525 input stream. If the overlay source can not be resolved, the original stream
2526 is passed through unchanged so playback is never interrupted.
2527
2528 :param queue: The PlayerQueue holding the overlay source and volume.
2529 :param audio_input: The audio stream (raw PCM in ``pcm_format``) to mix into.
2530 :param pcm_format: PCM format of both the input and the mixed output.
2531 """
2532 overlay_input = await self._resolve_overlay_input(queue)
2533 if overlay_input is None:
2534 # overlay source unavailable: degrade gracefully to music-only
2535 async for chunk in audio_input:
2536 yield chunk
2537 return
2538 async for chunk in get_ffmpeg_overlay_stream(
2539 audio_input=audio_input,
2540 overlay_input=overlay_input,
2541 pcm_format=pcm_format,
2542 overlay_volume=queue.overlay_volume,
2543 chunk_size=pcm_format.pcm_sample_size,
2544 ):
2545 yield chunk
2546
2547 def crossfade_allowed(
2548 self,
2549 queue_item: QueueItem,
2550 crossfade_mode: CrossfadeMode,
2551 player_id: str,
2552 flow_mode: bool = False,
2553 next_queue_item: QueueItem | None = None,
2554 sample_rate: int | None = None,
2555 next_sample_rate: int | None = None,
2556 ) -> bool:
2557 """Get the crossfade config for a queue item."""
2558 if crossfade_mode == CrossfadeMode.DISABLED:
2559 return False
2560 if not (self.mass.player_queues.get(queue_item.queue_id)):
2561 return False # just a guard
2562 if not (self.mass.players.get_player(player_id)):
2563 return False # just a guard
2564 if queue_item.media_type != MediaType.TRACK:
2565 self.logger.debug("Skipping crossfade: current item is not a track")
2566 return False
2567 # check if the next item is part of the same album
2568 next_item = next_queue_item or self.mass.player_queues.get_next_item(
2569 queue_item.queue_id, queue_item.queue_item_id
2570 )
2571 if not next_item:
2572 # there is no next item!
2573 return False
2574 # check if next item is a track
2575 if next_item.media_type != MediaType.TRACK:
2576 self.logger.debug("Skipping crossfade: next item is not a track")
2577 return False
2578 if (
2579 isinstance(queue_item.media_item, Track)
2580 and isinstance(next_item.media_item, Track)
2581 and queue_item.media_item.album
2582 and next_item.media_item.album
2583 and queue_item.media_item.album == next_item.media_item.album
2584 and not self.mass.config.get_raw_core_config_value(
2585 "streams", CONF_ALLOW_CROSSFADE_SAME_ALBUM, False
2586 )
2587 ):
2588 # in general, crossfade is not desired for tracks of the same (gapless) album
2589 # because we have no accurate way to determine if the album is gapless or not,
2590 # for now we just never crossfade between tracks of the same album
2591 self.logger.debug("Skipping crossfade: next item is part of the same album")
2592 return False
2593
2594 # check if we're allowed to crossfade on different sample rates
2595 if (
2596 not flow_mode
2597 and sample_rate
2598 and next_sample_rate
2599 and sample_rate != next_sample_rate
2600 and not self.mass.config.get_raw_player_config_value(
2601 player_id,
2602 CONF_ENTRY_CROSSFADE_DIFFERENT_SAMPLE_RATES.key,
2603 CONF_ENTRY_CROSSFADE_DIFFERENT_SAMPLE_RATES.default_value,
2604 )
2605 ):
2606 self.logger.debug(
2607 "Skipping crossfade: player(protocol) does not support gapless playback "
2608 "with different sample rates (%s vs %s)",
2609 sample_rate,
2610 next_sample_rate,
2611 )
2612 return False
2613
2614 return True
2615
2616 def clear_crossfade_data(self, queue_id: str) -> None:
2617 """
2618 Clear any pending crossfade data for a queue.
2619
2620 :param queue_id: The queue ID to clear crossfade data for.
2621 """
2622 if queue_id in self._crossfade_data:
2623 self.logger.debug("Clearing crossfade data for queue %s", queue_id)
2624 del self._crossfade_data[queue_id]
2625
2626 async def get_shoutcast_stream(
2627 self, url: str, streamdetails: StreamDetails
2628 ) -> AsyncGenerator[bytes]:
2629 """
2630 Yield audio from a legacy Shoutcast server, with ICY metadata parsed inline.
2631
2632 :param url: Shoutcast stream URL.
2633 :param streamdetails: StreamDetails to update with ICY metadata as it arrives.
2634 """
2635 self.logger.debug("Start streaming from legacy Shoutcast server: %s", url)
2636
2637 parsed = urlparse(url)
2638 host = parsed.hostname
2639 port = parsed.port or 80
2640 path = parsed.path or "/"
2641 if parsed.query:
2642 path = f"{path}?{parsed.query}"
2643
2644 try:
2645 # Open raw socket connection
2646 reader, writer = await asyncio.wait_for(asyncio.open_connection(host, port), timeout=30)
2647 except TimeoutError as err:
2648 raise AudioError(f"Timeout connecting to Shoutcast stream {url}") from err
2649 except (OSError, ConnectionError) as err:
2650 raise AudioError(f"Failed to connect to Shoutcast stream {url}") from err
2651
2652 try:
2653 # Send HTTP request with ICY metadata header
2654 request = (
2655 f"GET {path} HTTP/1.1\r\n"
2656 f"Host: {host}\r\n"
2657 f"User-Agent: {HTTP_HEADERS['User-Agent']}\r\n"
2658 f"Icy-MetaData: 1\r\n\r\n"
2659 )
2660 writer.write(request.encode())
2661 await writer.drain()
2662
2663 # Read and parse response line
2664 try:
2665 response_line = await asyncio.wait_for(reader.readline(), timeout=10)
2666 except TimeoutError as err:
2667 raise AudioError("Timeout reading Shoutcast response") from err
2668
2669 if not response_line.startswith(b"ICY"):
2670 raise InvalidDataError("Invalid Shoutcast response")
2671
2672 # Read headers until empty line
2673 headers: dict[str, str] = {}
2674 while True:
2675 try:
2676 line = await asyncio.wait_for(reader.readline(), timeout=5)
2677 except TimeoutError as err:
2678 raise AudioError("Timeout reading Shoutcast headers") from err
2679
2680 if line in (b"\r\n", b"\n", b""):
2681 break
2682
2683 if b":" in line:
2684 try:
2685 key, value = line.decode("latin-1", errors="ignore").split(":", 1)
2686 headers[key.strip().lower()] = value.strip()
2687 except UnicodeDecodeError, ValueError:
2688 continue
2689
2690 # Get metadata interval
2691 meta_int_str = headers.get("icy-metaint")
2692 if not meta_int_str:
2693 raise InvalidDataError("No icy-metaint header in Shoutcast response")
2694
2695 try:
2696 meta_int = int(meta_int_str)
2697 except ValueError as err:
2698 raise InvalidDataError("Invalid icy-metaint value") from err
2699
2700 self.logger.debug("Connected to Shoutcast stream %s (icy-metaint: %s)", url, meta_int)
2701
2702 # Stream audio data with metadata parsing
2703 while True:
2704 try:
2705 # Read audio chunk
2706 audio_chunk = await reader.readexactly(meta_int)
2707 yield audio_chunk
2708
2709 # Read metadata length
2710 meta_byte = await reader.readexactly(1)
2711 if meta_byte == b"\x00":
2712 continue
2713
2714 meta_length = ord(meta_byte) * 16
2715 meta_data = await reader.readexactly(meta_length)
2716 self._parse_icy_metadata(meta_data, streamdetails)
2717
2718 except asyncio.exceptions.IncompleteReadError:
2719 # End of stream
2720 break
2721
2722 finally:
2723 writer.close()
2724 await writer.wait_closed()
2725
2726 # --- Private methods ---
2727
2728 def _notify_provider_streamed(
2729 self, streamdetails: StreamDetails, finished: bool, seconds_streamed: float
2730 ) -> None:
2731 """Report a (mostly) streamed item back to the provider that owns it."""
2732 if not finished and seconds_streamed < 90:
2733 return
2734 provider = self.mass.get_provider(streamdetails.provider)
2735 # plugin providers serve playable items too, but on_streamed is MusicProvider-only
2736 if provider is None or provider.type != ProviderType.MUSIC:
2737 return
2738 music_prov = cast("MusicProvider", provider)
2739 self.mass.create_task(music_prov.on_streamed(streamdetails))
2740
2741 def _get_volume_normalization_preference(
2742 self, streamdetails: StreamDetails
2743 ) -> VolumeNormalizationMode:
2744 """Return the configured normalization preference for the stream's media type."""
2745 conf_key = (
2746 CONF_VOLUME_NORMALIZATION_RADIO
2747 if streamdetails.media_type == MediaType.RADIO
2748 else CONF_VOLUME_NORMALIZATION_TRACKS
2749 )
2750 return VolumeNormalizationMode(
2751 self.mass.streams.get_config_value(conf_key, return_type=str)
2752 )
2753
2754 def _update_radio_stream_metadata(
2755 self,
2756 streamdetails: StreamDetails,
2757 artist: str | None,
2758 title: str,
2759 image_url: str | None = None,
2760 album: str | None = None,
2761 ) -> None:
2762 """
2763 Update radio stream metadata and trigger artwork lookup.
2764
2765 :param streamdetails: The stream details to update.
2766 :param artist: Artist name (will be normalized).
2767 :param title: Track title (will be cleaned for display).
2768 :param image_url: Optional image URL from stream metadata.
2769 :param album: Optional album name.
2770 """
2771 station_image_url = image_url or self.mass.metadata.get_radio_stream_station_image(
2772 streamdetails
2773 )
2774 artist_normalized = (
2775 self.mass.metadata.normalize_radio_artist_name(artist) if artist else None
2776 )
2777 display_title, _ = parse_title_and_version(title, strip_for_display=True)
2778
2779 streamdetails.stream_metadata = StreamMetadata(
2780 title=display_title,
2781 artist=artist_normalized,
2782 album=album,
2783 image_url=station_image_url,
2784 )
2785 streamdetails.stream_metadata_last_updated = time.time()
2786 if streamdetails.queue_id:
2787 self.mass.player_queues.signal_update(streamdetails.queue_id)
2788
2789 # Fetch artwork in background (track, album then artist)
2790 if artist and title and not image_url:
2791 self.mass.call_later(
2792 0.2,
2793 self.mass.metadata.update_radio_stream_artwork,
2794 streamdetails,
2795 task_id=f"update_radio_artwork_{streamdetails.queue_id}",
2796 )
2797
2798 async def _cache_radio_result(
2799 self,
2800 url: str,
2801 stream_type: StreamType,
2802 resolved_url: str | None = None,
2803 ) -> tuple[str, StreamType]:
2804 """Cache and return a radio stream resolution result."""
2805 result = (resolved_url or url, stream_type)
2806 await self.mass.cache.set(
2807 url,
2808 result,
2809 expiration=3600 * 3,
2810 provider=CACHE_PROVIDER,
2811 category=CACHE_CATEGORY_RESOLVED_RADIO_URL,
2812 )
2813 return result
2814
2815 async def _handle_client_error_for_radio_stream(
2816 self, url: str, err: aiohttp.ClientError, fallback_stream_type: StreamType
2817 ) -> tuple[str, StreamType]:
2818 """Handle aiohttp client errors during radio stream resolution."""
2819 # Prefer the final post-redirect URL: aiohttp follows redirects before raising,
2820 # but the original url may just point at a redirector rather than the ICY endpoint.
2821 request_info = getattr(err, "request_info", None)
2822 validate_url = str(request_info.url) if request_info is not None else url
2823
2824 # Check if this is a Shoutcast/ICY response that aiohttp can't parse
2825 if isinstance(err, aiohttp.ClientResponseError) and "ICY" in str(err).upper():
2826 self.logger.debug(
2827 "ICY response detected for %s, validating Shoutcast stream", validate_url
2828 )
2829 if await self._validate_shoutcast_stream(validate_url):
2830 return await self._cache_radio_result(
2831 url, StreamType.SHOUTCAST, resolved_url=validate_url
2832 )
2833 self.logger.warning(
2834 "ICY response detected but Shoutcast validation failed for %s", validate_url
2835 )
2836 return await self._cache_radio_result(
2837 url, fallback_stream_type, resolved_url=validate_url
2838 )
2839
2840 # Other aiohttp errors - might still be Shoutcast, check it
2841 self.logger.debug("aiohttp error for %s, checking if legacy Shoutcast stream", validate_url)
2842 if await self._validate_shoutcast_stream(validate_url):
2843 return await self._cache_radio_result(
2844 url, StreamType.SHOUTCAST, resolved_url=validate_url
2845 )
2846
2847 # Unknown error - still try to stream
2848 self.logger.warning(
2849 "Failed to parse radio URL %s: %s - attempting direct stream", validate_url, str(err)
2850 )
2851 return await self._cache_radio_result(url, fallback_stream_type, resolved_url=validate_url)
2852
2853 async def _iter_audio_source_pcm(
2854 self,
2855 streamdetails: StreamDetails,
2856 pcm_format: AudioFormat,
2857 ) -> AsyncGenerator[bytes]:
2858 """Yield PCM for an AudioSource, bypassing ffmpeg when formats match."""
2859 if _pcm_formats_match(streamdetails.audio_format, pcm_format):
2860 source_gen = self._open_audio_source_generator(streamdetails)
2861 async for chunk in realtime_pcm_pacer(source_gen, pcm_format):
2862 yield chunk
2863 return
2864 # format mismatch â fall back to ffmpeg for resampling (still small chunks)
2865 async for chunk in self.get_media_stream(
2866 streamdetails=streamdetails,
2867 pcm_format=pcm_format,
2868 filter_params=None,
2869 chunk_seconds=AUDIO_SOURCE_CHUNK_SECONDS,
2870 ):
2871 yield chunk
2872
2873 def _open_audio_source_generator(self, streamdetails: StreamDetails) -> AsyncGenerator[bytes]:
2874 """Open the raw PCM generator for an AudioSource (CUSTOM or NAMED_PIPE)."""
2875 if streamdetails.stream_type == StreamType.CUSTOM:
2876 provider = self.mass.get_provider(streamdetails.provider)
2877 if provider is None:
2878 raise ProviderUnavailableError(
2879 f"Provider {streamdetails.provider} for stream is no longer available"
2880 )
2881 provider = cast("MusicProvider | PluginProvider", provider)
2882 return provider.get_audio_stream(streamdetails)
2883 if streamdetails.stream_type == StreamType.NAMED_PIPE:
2884 assert isinstance(streamdetails.path, str) # for type checking
2885 return read_named_pipe(streamdetails.path)
2886 raise AudioError(f"Unsupported stream_type {streamdetails.stream_type} for AudioSource")
2887
2888 def _parse_icy_metadata(self, meta_data: bytes, streamdetails: StreamDetails) -> None:
2889 """
2890 Parse ICY metadata and update streamdetails.
2891
2892 Sets the cleaned stream title and, when the title parses as "Artist - Track",
2893 triggers a radio-artwork metadata update.
2894
2895 :param meta_data: Raw metadata bytes from an ICY stream chunk.
2896 :param streamdetails: StreamDetails to update with parsed title and metadata.
2897 """
2898 if not meta_data:
2899 return
2900
2901 meta_data = meta_data.rstrip(b"\0")
2902 # Match StreamTitle, handling apostrophes in titles
2903 stream_title_re = re.search(rb"StreamTitle='(.*?)';", meta_data)
2904
2905 if not stream_title_re:
2906 self.logger.log(
2907 VERBOSE_LOG_LEVEL,
2908 "ICY metadata does not contain StreamTitle field. Raw: %s",
2909 meta_data.decode("utf-8", errors="replace")[:200],
2910 )
2911 return
2912
2913 try:
2914 # in 99% of the cases the stream title is utf-8 encoded
2915 stream_title = stream_title_re.group(1).decode("utf-8")
2916 except UnicodeDecodeError:
2917 # fallback to iso-8859-1
2918 stream_title = stream_title_re.group(1).decode("iso-8859-1", errors="replace")
2919
2920 cleaned_stream_title = clean_stream_title(stream_title)
2921
2922 if not cleaned_stream_title:
2923 return
2924
2925 if cleaned_stream_title == streamdetails.stream_title:
2926 return
2927
2928 self.logger.log(VERBOSE_LOG_LEVEL, "ICY Radio streamtitle original: %s", stream_title)
2929 self.logger.log(
2930 VERBOSE_LOG_LEVEL, "ICY Radio streamtitle cleaned: %s", cleaned_stream_title
2931 )
2932 streamdetails.stream_title = cleaned_stream_title
2933
2934 # Prefer station-provided cover art from the ICY 'StreamUrl' field (when it is
2935 # an image) over the MusicBrainz artwork lookup in _update_radio_stream_metadata.
2936 image_url = self._parse_icy_image_url(meta_data)
2937
2938 # Parse the original title for structured fields first so stations that announce
2939 # an album can refine the artwork lookup; fall back to the "Artist - Track" split.
2940 album: str | None = None
2941 if parsed := parse_quoted_stream_title(stream_title):
2942 track_name, artist_name_raw, album = parsed
2943 elif " - " in cleaned_stream_title:
2944 artist_name_raw, track_name = (
2945 part.strip() for part in cleaned_stream_title.split(" - ", 1)
2946 )
2947 else:
2948 return
2949
2950 if artist_name_raw and track_name:
2951 self.logger.debug(
2952 "ICY metadata: artist='%s', track='%s', album='%s'",
2953 artist_name_raw,
2954 track_name,
2955 album,
2956 )
2957 self._update_radio_stream_metadata(
2958 streamdetails,
2959 artist=artist_name_raw,
2960 title=track_name,
2961 album=album,
2962 image_url=image_url,
2963 )
2964
2965 def _parse_icy_image_url(self, meta_data: bytes) -> str | None:
2966 """
2967 Return a PNG or JPEG cover-art URL from the ICY 'StreamUrl' field, if present.
2968
2969 :param meta_data: Raw metadata bytes from an ICY stream chunk.
2970 """
2971 # The trailing semicolon is optional to match sources that omit it.
2972 stream_url_re = re.search(rb"StreamUrl='([^']*)'", meta_data)
2973 if not stream_url_re:
2974 return None
2975 try:
2976 image_url = stream_url_re.group(1).decode("utf-8").strip()
2977 except UnicodeDecodeError:
2978 return None
2979 if not image_url:
2980 return None
2981 # StreamUrl is not a standardized artwork field (reference clients such as VLC
2982 # ignore it and it conventionally holds a station website link), so only accept
2983 # values that point at a PNG or JPEG image.
2984 parsed = urlparse(image_url)
2985 if parsed.scheme not in ("http", "https"):
2986 return None
2987 if not parsed.path.lower().endswith((".png", ".jpg", ".jpeg")):
2988 return None
2989 self.logger.debug("ICY metadata: StreamUrl image='%s'", image_url)
2990 return image_url
2991
2992 async def _validate_shoutcast_stream(self, url: str) -> bool:
2993 """
2994 Return True if the URL responds with a legacy Shoutcast "ICY 200 OK" line.
2995
2996 :param url: The URL to validate.
2997 """
2998 try:
2999 parsed = urlparse(url)
3000 host = parsed.hostname
3001 port = parsed.port or 80
3002 path = parsed.path or "/"
3003 if parsed.query:
3004 path = f"{path}?{parsed.query}"
3005
3006 # Open raw socket connection with timeout
3007 reader, writer = await asyncio.wait_for(asyncio.open_connection(host, port), timeout=10)
3008 try:
3009 # Send minimal HTTP request with ICY metadata header
3010 request = f"GET {path} HTTP/1.1\r\nHost: {host}\r\nIcy-MetaData: 1\r\n\r\n"
3011 writer.write(request.encode())
3012 await writer.drain()
3013
3014 # Read just the response line
3015 response_line = await asyncio.wait_for(reader.readline(), timeout=5)
3016 finally:
3017 writer.close()
3018 await writer.wait_closed()
3019
3020 # Check if response starts with "ICY"
3021 decoded_line = response_line.decode("latin-1", errors="ignore").strip()
3022 return decoded_line.startswith("ICY")
3023
3024 except TimeoutError:
3025 self.logger.debug("Timeout during Shoutcast validation for %s", url)
3026 return False
3027 except OSError, ConnectionError:
3028 self.logger.debug("Connection failed during Shoutcast validation for %s", url)
3029 return False
3030 except UnicodeDecodeError:
3031 self.logger.debug("Invalid response encoding during Shoutcast validation for %s", url)
3032 return False
3033
3034 def _resolve_player_dsp_config(self, player: Player) -> DSPConfig:
3035 """
3036 Resolve the effective DSP config for a player.
3037
3038 Single source of truth shared by every code path that needs to know
3039 whether DSP will run for this player. Protocol wrappers defer to their
3040 parent player; single-leg ``player_group`` instances that don't expose
3041 ``MULTI_DEVICE_DSP`` defer to their first member; players whose grouping
3042 context prevents DSP get a disabled config back regardless.
3043
3044 :param player: The player to resolve DSP config for.
3045 """
3046 dsp_player_id = self._resolve_player_dsp_config_id(player)
3047 dsp = self.mass.config.get_player_dsp_config(dsp_player_id)
3048 if is_grouping_preventing_dsp(player):
3049 dsp.enabled = False
3050 elif player.provider.domain == "player_group" and (
3051 PlayerFeature.MULTI_DEVICE_DSP not in player.state.supported_features
3052 ):
3053 if not player.state.group_members:
3054 dsp.enabled = False
3055 return dsp
3056
3057 def _resolve_player_dsp_config_id(self, player: Player) -> str:
3058 """
3059 Return the player identifier that supplies the effective DSP config.
3060
3061 :param player: Player whose DSP config source should be resolved.
3062 """
3063 dsp_player_id = player.protocol_parent_id or player.player_id
3064 if (
3065 not is_grouping_preventing_dsp(player)
3066 and player.provider.domain == "player_group"
3067 and PlayerFeature.MULTI_DEVICE_DSP not in player.state.supported_features
3068 and player.state.group_members
3069 ):
3070 child_player = self.mass.players.get_player(player.state.group_members[0])
3071 assert child_player is not None
3072 dsp_player_id = child_player.player_id
3073 return dsp_player_id
3074
3075 def _get_output_channels(self, player: Player | None, player_id: str) -> str:
3076 """
3077 Return the configured output channels for the rendering player.
3078
3079 The value may be stored on the rendering player(protocol) itself (the
3080 protocol section of the config UI) or on its visible parent player (the
3081 native section); the rendering player's own stored value wins.
3082 """
3083 parent_id = player.protocol_parent_id if player and player.protocol_parent_id else player_id
3084 parent_value = self.mass.config.get_raw_player_config_value(
3085 parent_id, CONF_OUTPUT_CHANNELS, "stereo"
3086 )
3087 return self.mass.config.get_raw_player_config_value(
3088 player.player_id if player else player_id, CONF_OUTPUT_CHANNELS, parent_value
3089 )
3090
3091 def _pick_pcm_bit_depth(
3092 self,
3093 players: Iterable[Player],
3094 streamdetails: StreamDetails | None,
3095 crossfade_enabled: bool,
3096 overlay_active: bool = False,
3097 ) -> tuple[ContentType, int]:
3098 """
3099 Return ``(content_type, bit_depth)`` for an internal PCM stream.
3100
3101 F32 is chosen when audio processing (crossfade, audio overlay, volume
3102 normalization, DSP) will run on the stream â those need the extra
3103 headroom to avoid clipping and precision loss. Otherwise the source's
3104 native bit depth is reused so we don't waste memory upcasting a 16-bit
3105 stream to 32-bit just to pass it through. When the source is unknown
3106 (no streamdetails) we fall back to F32 conservatively.
3107 """
3108 if streamdetails is None:
3109 return INTERNAL_PCM_FORMAT.content_type, INTERNAL_PCM_FORMAT.bit_depth
3110 needs_headroom = (
3111 crossfade_enabled
3112 or overlay_active
3113 or streamdetails.volume_normalization_mode != VolumeNormalizationMode.DISABLED
3114 or any(self._resolve_player_dsp_config(player).enabled for player in players)
3115 )
3116 if needs_headroom:
3117 return INTERNAL_PCM_FORMAT.content_type, INTERNAL_PCM_FORMAT.bit_depth
3118 bit_depth = streamdetails.audio_format.bit_depth
3119 return ContentType.from_bit_depth(bit_depth), bit_depth
3120
3121 def _select_audio_source_pcm_format(
3122 self,
3123 player: Player,
3124 streamdetails: StreamDetails,
3125 supported_sample_rates: Iterable[int] | None = None,
3126 ) -> AudioFormat:
3127 """
3128 Return a passthrough PCM format for a realtime AudioSource item.
3129
3130 The format matches the source's native sample rate, bit depth and
3131 channel count whenever the player can accept them; if the player does
3132 not support the source's sample rate, it is snapped down to the
3133 closest supported rate. No F32 widening â realtime sources skip every
3134 processing stage that would otherwise need it. Surround sources are
3135 still folded down to stereo, which every output path requires anyway.
3136
3137 :param player: The player requesting the stream.
3138 :param streamdetails: Stream details for the AudioSource item.
3139 :param supported_sample_rates: Rates shared by every output player, if applicable.
3140 """
3141 resolved_sample_rates = (
3142 list(supported_sample_rates)
3143 if supported_sample_rates is not None
3144 else [sample_rate for sample_rate, _ in player.get_supported_sample_rates()]
3145 )
3146 source_rate = streamdetails.audio_format.sample_rate
3147 if source_rate in resolved_sample_rates:
3148 output_sample_rate = source_rate
3149 else:
3150 output_sample_rate = max(
3151 (rate for rate in resolved_sample_rates if rate <= source_rate),
3152 default=min(resolved_sample_rates),
3153 )
3154 bit_depth = streamdetails.audio_format.bit_depth
3155 return AudioFormat(
3156 content_type=ContentType.from_bit_depth(bit_depth),
3157 sample_rate=output_sample_rate,
3158 bit_depth=bit_depth,
3159 # a realtime source may announce more channels than anything downstream can
3160 # carry (a VBAN stream can be configured up to 8), and player handoff formats
3161 # copy this count straight through, so fold it here
3162 channels=min(streamdetails.audio_format.channels, 2),
3163 )
3164
3165 def _flow_restart_context(
3166 self, queue_id: str, protocol_player: Player | None
3167 ) -> tuple[str, list[int]]:
3168 """
3169 Resolve the flow mode config and supported sample rates for restart decisions.
3170
3171 Prefers the protocol player actually consuming the flow stream over the
3172 queue's (wrapper) player, whose config may lack the audio specific entries.
3173 """
3174 if protocol_player is None:
3175 protocol_player = self.mass.players.get_player(queue_id)
3176 if protocol_player is None:
3177 flow_mode_sample_rate_conf = self.mass.config.get_raw_player_config_value(
3178 queue_id, CONF_FLOW_MODE_SAMPLE_RATE, FLOW_MODE_SAMPLE_RATE_SMART
3179 )
3180 return flow_mode_sample_rate_conf, []
3181 flow_mode_sample_rate_conf = cast(
3182 "str",
3183 protocol_player.config.get_value(
3184 CONF_FLOW_MODE_SAMPLE_RATE, FLOW_MODE_SAMPLE_RATE_SMART
3185 ),
3186 )
3187 supported_sample_rates = sorted(
3188 {sr for sr, _ in protocol_player.get_supported_sample_rates()}
3189 )
3190 return flow_mode_sample_rate_conf, supported_sample_rates
3191
3192 def _flow_stream_needs_restart(
3193 self,
3194 queue_track: QueueItem,
3195 pcm_format: AudioFormat,
3196 supported_sample_rates: list[int],
3197 flow_mode_sample_rate_conf: str,
3198 is_first_track: bool,
3199 ) -> bool:
3200 """
3201 Return True if the upcoming queue track requires exiting the flow stream.
3202
3203 Covers every case where the flow loop should break and hand control back to
3204 the queue controller for restart:
3205
3206 - Live media (radio, audio sources): cannot be played inside a flow,
3207 the controller will fall back to a single-item stream.
3208 - Sample rate mismatch ('smart' / 'bit_perfect' modes only): the next
3209 track's sample rate (snapped up to the closest supported player rate,
3210 mirroring select_flow_pcm_format's anchoring logic) is incompatible with
3211 the current flow rate, so a new flow must be opened.
3212
3213 The first (anchor) track is always allowed to continue for the sample
3214 rate check; select_flow_pcm_format has already snapped the flow rate to it.
3215
3216 :param queue_track: The upcoming queue item.
3217 :param pcm_format: The current flow stream's PCM format.
3218 :param supported_sample_rates: Sorted list of the player's supported rates.
3219 :param flow_mode_sample_rate_conf: The flow mode sample rate config value.
3220 :param is_first_track: Whether this is the first track of the flow stream.
3221 """
3222 # live audio (radio, plugin or audio source) cannot be flowed; let the
3223 # queue controller fall back to single-item streaming for this item
3224 if queue_track.media_type in (MediaType.RADIO, MediaType.AUDIO_SOURCE):
3225 self.logger.info(
3226 "Live media item %s (%s, %s) encountered in flow stream "
3227 "- breaking out to single item stream",
3228 queue_track.queue_item_id,
3229 queue_track.name,
3230 queue_track.media_type,
3231 )
3232 return True
3233
3234 if is_first_track or queue_track.streamdetails is None:
3235 return False
3236 raw_next_rate = queue_track.streamdetails.audio_format.sample_rate
3237 if not raw_next_rate or not supported_sample_rates:
3238 return False
3239 effective_next_rate = _snap_supported_rate_up(raw_next_rate, supported_sample_rates)
3240
3241 # branch order mirrors select_flow_pcm_format: fixed-rate modes resample
3242 # everything to the chosen rate (no restart); bit_perfect restarts on any
3243 # mismatch; anything else falls through to smart-anchor behavior so
3244 # unknown/legacy config values don't silently pin the flow forever.
3245 if flow_mode_sample_rate_conf in (
3246 FLOW_MODE_SAMPLE_RATE_48000,
3247 FLOW_MODE_SAMPLE_RATE_96000,
3248 FLOW_MODE_SAMPLE_RATE_HIGHEST,
3249 ):
3250 needs_restart = False
3251 elif flow_mode_sample_rate_conf == FLOW_MODE_SAMPLE_RATE_BIT_PERFECT:
3252 needs_restart = effective_next_rate != pcm_format.sample_rate
3253 else:
3254 needs_restart = effective_next_rate > pcm_format.sample_rate
3255
3256 if needs_restart:
3257 self.logger.info(
3258 "Track %s (%s) sample rate %s (snapped to %s) incompatible with flow rate %s "
3259 "(mode: %s) - breaking out to restart flow stream",
3260 queue_track.queue_item_id,
3261 queue_track.name,
3262 raw_next_rate,
3263 effective_next_rate,
3264 pcm_format.sample_rate,
3265 flow_mode_sample_rate_conf,
3266 )
3267 return needs_restart
3268
3269 @asynccontextmanager
3270 async def _connect_radio_stream(self, url: str, **kwargs: Any) -> AsyncGenerator[Any]:
3271 """
3272 Connect to a radio stream URL with fallback for legacy SSL/TLS configurations.
3273
3274 Some radio servers use outdated TLS configurations that reject modern
3275 cipher suites. Since radio streams are public broadcast content,
3276 relaxing cipher requirements is acceptable.
3277
3278 :param url: The radio stream URL to connect to.
3279 :param kwargs: Additional keyword arguments passed to aiohttp get().
3280 """
3281 request_url = encoded_request_url(url)
3282 try:
3283 async with self.mass.http_session_no_ssl.get(request_url, **kwargs) as resp:
3284 yield resp
3285 except ClientConnectorSSLError:
3286 self.logger.info(
3287 "SSL handshake failed for %s, retrying with permissive cipher configuration", url
3288 )
3289 insecure_ssl_context = ssl_util.client_context_no_verify(
3290 ssl_util.SSLCipherList.INSECURE
3291 )
3292 async with self.mass.http_session_no_ssl.get(
3293 request_url, ssl=insecure_ssl_context, **kwargs
3294 ) as resp:
3295 yield resp
3296
3297 async def _update_hls_radio_metadata(
3298 self,
3299 streamdetails: StreamDetails,
3300 elapsed_time: int,
3301 ) -> None:
3302 """
3303 Update HLS radio stream metadata by fetching the playlist.
3304
3305 Fetches the HLS playlist and extracts metadata from EXTINF lines.
3306
3307 :param streamdetails: StreamDetails object to update with metadata
3308 :param elapsed_time: Current playback position in seconds (unused for live radio)
3309 """
3310 mass = self.mass
3311 try:
3312 # Get the actual media playlist URL from cache or resolve it
3313 # We cache the media_playlist_url in streamdetails.data to avoid re-resolving
3314 if streamdetails.data is None:
3315 streamdetails.data = {}
3316 media_playlist_url = streamdetails.data.get("hls_media_playlist_url")
3317 if not media_playlist_url:
3318 try:
3319 assert isinstance(streamdetails.path, str) # for type checking
3320 substream = await self.get_hls_substream(streamdetails.path)
3321 media_playlist_url = substream.path
3322 streamdetails.data["hls_media_playlist_url"] = media_playlist_url
3323 except Exception as err:
3324 self.logger.warning(
3325 "Failed to resolve HLS substream for metadata monitoring: %s", err
3326 )
3327 return
3328
3329 # Fetch the media playlist
3330 timeout = ClientTimeout(total=0, connect=10, sock_read=30)
3331 try:
3332 async with mass.http_session_no_ssl.get(
3333 encoded_request_url(media_playlist_url), timeout=timeout
3334 ) as resp:
3335 resp.raise_for_status()
3336 playlist_content = await resp.text()
3337 except ClientResponseError as err:
3338 # Session token likely expired (410/403) â drop cache so next poll re-resolves
3339 if err.status in (403, 410):
3340 streamdetails.data.pop("hls_media_playlist_url", None)
3341 raise
3342
3343 # Parse the playlist and look for EXTINF metadata
3344 # The most recent segment usually has the current metadata
3345 lines = playlist_content.strip().split("\n")
3346 for line in reversed(lines):
3347 if line.startswith("#EXTINF:"):
3348 # Extract metadata from EXTINF line
3349 metadata = parse_extinf_metadata(line)
3350
3351 # Build stream title from title and artist
3352 title = metadata.get("title", "")
3353 artist = metadata.get("artist", "")
3354 image_url = (
3355 metadata.get("image") or metadata.get("artwork") or metadata.get("cover")
3356 )
3357 if not artist and " - " in title:
3358 artist, title = title.split(" - ", 1)
3359 if title or artist:
3360 # Format as "Artist - Title"
3361 if artist and title:
3362 stream_title = f"{artist} - {title}"
3363 elif title:
3364 stream_title = title
3365 else:
3366 stream_title = artist
3367
3368 # Clean the stream title
3369 cleaned_title = clean_stream_title(stream_title)
3370
3371 # Only update if changed
3372 if cleaned_title != streamdetails.stream_title and cleaned_title:
3373 self.logger.log(
3374 VERBOSE_LOG_LEVEL, "HLS Radio metadata updated: %s", cleaned_title
3375 )
3376 streamdetails.stream_title = cleaned_title
3377 self._update_radio_stream_metadata(
3378 streamdetails,
3379 artist=artist or None,
3380 title=title or cleaned_title,
3381 image_url=image_url,
3382 )
3383
3384 # Only check the most recent EXTINF
3385 break
3386
3387 except Exception as err:
3388 self.logger.debug("Error fetching HLS metadata: %s", err)
3389
3390 @staticmethod
3391 def _normalize_reconnecting_urls(url: str | list[MultiPartPath]) -> list[str]:
3392 """Normalize a single URL or a sequence into a non-empty list."""
3393 if isinstance(url, str):
3394 return [url]
3395 if not url:
3396 msg = "Radio stream requires at least one URL"
3397 raise InvalidDataError(msg)
3398 return [part.path for part in url]
3399
3400 async def _resolve_overlay_input(self, queue: PlayerQueue) -> str | None:
3401 """
3402 Resolve the queue's overlay source to a file path or URL for ffmpeg.
3403
3404 Returns None (with a warning logged) when the source can not be resolved,
3405 so the caller can degrade to music-only playback.
3406 """
3407 if not (mapping := queue.overlay_source):
3408 return None
3409 try:
3410 provider = self.mass.get_provider(mapping.provider)
3411 if provider is None:
3412 raise MediaNotFoundError(f"Provider {mapping.provider} is not available")
3413 stream_prov = cast("MusicProvider | PluginProvider", provider)
3414 streamdetails = await stream_prov.get_stream_details(
3415 mapping.item_id, MediaType.SOUND_EFFECT
3416 )
3417 except Exception as err:
3418 self.logger.warning(
3419 "Audio overlay source %s is unavailable (%s) - continuing without overlay",
3420 mapping.uri,
3421 str(err) or err.__class__.__name__,
3422 )
3423 return None
3424 if streamdetails.stream_type not in (StreamType.LOCAL_FILE, StreamType.HTTP) or not (
3425 isinstance(streamdetails.path, str)
3426 ):
3427 self.logger.warning(
3428 "Audio overlay source %s uses unsupported stream type %s "
3429 "- continuing without overlay",
3430 mapping.uri,
3431 streamdetails.stream_type,
3432 )
3433 return None
3434 if streamdetails.stream_type == StreamType.LOCAL_FILE and not await aiofiles.os.path.isfile(
3435 streamdetails.path
3436 ):
3437 # guard against stale sources: feeding a missing file to the mixer would
3438 # kill the whole (music) stream instead of just the overlay
3439 self.logger.warning(
3440 "Audio overlay source %s does not exist - continuing without overlay",
3441 streamdetails.path,
3442 )
3443 return None
3444 return streamdetails.path
3445