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