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