/
/
1"""
2Controller to stream audio to players.
3
4The streams controller hosts a basic, unprotected HTTP-only webserver
5purely to stream audio packets to players.
6"""
7
8from __future__ import annotations
9
10import asyncio
11import logging
12import os
13from collections.abc import AsyncGenerator
14from contextlib import aclosing
15from math import ceil
16from typing import TYPE_CHECKING, cast
17from uuid import uuid4
18
19from aiofiles.os import wrap
20from aiohttp import web
21from music_assistant_models.audio_processing import AudioQueueProcessing
22from music_assistant_models.config_entries import ConfigEntry, ConfigValueOption
23from music_assistant_models.enums import (
24 ConfigEntryType,
25 ContentType,
26 CrossfadeMode,
27 MediaType,
28 PlayerFeature,
29 ProviderType,
30 VolumeNormalizationMode,
31)
32from music_assistant_models.errors import (
33 AudioError,
34 InvalidDataError,
35 MediaNotFoundError,
36 ProviderUnavailableError,
37)
38from music_assistant_models.helpers import get_global_cache_value
39from music_assistant_models.media_items import AudioFormat
40
41from music_assistant.constants import (
42 CONF_BACKGROUND_SCAN_CONCURRENCY,
43 CONF_BIND_IP,
44 CONF_BIND_PORT,
45 CONF_CROSSFADE_DURATION,
46 CONF_CROSSFADE_MODE,
47 CONF_ENTRY_ENABLE_ICY_METADATA,
48 CONF_ENTRY_LOG_LEVEL,
49 CONF_ENTRY_VOLUME_NORMALIZATION_TARGET,
50 CONF_HTTP_PROFILE,
51 CONF_OUTPUT_CODEC,
52 CONF_PLAYER_QUEUES,
53 CONF_PREFER_WAV_FOR_LIVE_SOURCES,
54 CONF_PUBLISH_IP,
55 CONF_VALUE_AUTO,
56 CONF_VOLUME_NORMALIZATION_FIXED_GAIN_RADIO,
57 CONF_VOLUME_NORMALIZATION_FIXED_GAIN_TRACKS,
58 CONF_VOLUME_NORMALIZATION_RADIO,
59 CONF_VOLUME_NORMALIZATION_TRACKS,
60 DEFAULT_BACKGROUND_SCAN_CONCURRENCY,
61 DEFAULT_HOST,
62 DEFAULT_STREAM_HEADERS,
63 DLNA_CONTENT_FEATURES,
64 DLNA_CONTENT_FEATURES_REALTIME,
65 ICY_HEADERS,
66 SILENCE_FILE,
67 VERBOSE_LOG_LEVEL,
68 WILDCARD_BIND_IPS,
69)
70from music_assistant.controllers.players.helpers import AnnounceData
71from music_assistant.controllers.streams.announcements import (
72 DEFAULT_RENDER_TIMEOUT,
73 AnnouncementRenderer,
74)
75from music_assistant.controllers.streams.audio import StreamsAudio, overlay_active
76from music_assistant.controllers.streams.audio_analysis import AudioAnalysisController
77from music_assistant.controllers.streams.audio_processing import (
78 AudioProcessingManager,
79)
80from music_assistant.controllers.streams.constants import (
81 CONF_ALLOW_CROSSFADE_SAME_ALBUM,
82 CONF_BUFFER_SIZE,
83 CONF_BUFFER_SIZE_DEFAULT,
84 CONF_SMART_FADES_LOG_LEVEL,
85 DEFAULT_PORT,
86 DEFAULT_VOLUME_NORMALIZATION_MODE,
87 FLOW_STREAM_LEAD_OUT_SECONDS,
88 OUTCOME_ONLY_NORMALIZATION_MODES,
89 SINGLE_ITEM_READRATE,
90 SINGLE_ITEM_READRATE_INITIAL_BURST,
91 BufferSize,
92 get_available_buffer_sizes,
93)
94from music_assistant.controllers.streams.live_announcements import (
95 LIVE_ANNOUNCEMENT_STREAM_PATH,
96 LiveAnnouncementManager,
97)
98from music_assistant.helpers.audio import (
99 calculate_content_length,
100 create_streaming_wave_header,
101 get_content_length,
102 get_mime_type,
103 store_content_length_in_cache,
104)
105from music_assistant.helpers.ffmpeg import (
106 CACHE_ATTR_FFMPEG_VERSION,
107 CACHE_ATTR_LIBSOXR_PRESENT,
108 check_ffmpeg_version,
109 get_ffmpeg_stream,
110)
111from music_assistant.helpers.ffmpeg import LOGGER as FFMPEG_LOGGER
112from music_assistant.helpers.util import (
113 format_ip_for_url,
114 get_ip_addresses,
115 get_publish_ip_candidates,
116 get_source_ip_for_target,
117 sanitize_http_header_value,
118)
119from music_assistant.helpers.webserver import Webserver, redact_sensitive_headers
120from music_assistant.models.core_controller import CoreController
121from music_assistant.models.music_provider import MusicProvider, ProviderStreamLimitError
122from music_assistant.models.plugin import PluginProvider
123from music_assistant.providers.universal_group.constants import UGP_PREFIX
124from music_assistant.providers.universal_group.player import UniversalGroupPlayer
125
126if TYPE_CHECKING:
127 from music_assistant_models.config_entries import CoreConfig
128 from music_assistant_models.player import PlayerMedia
129 from music_assistant_models.player_queue import PlayerQueue
130 from music_assistant_models.queue_item import QueueItem
131 from music_assistant_models.streamdetails import StreamDetails
132
133 from music_assistant.controllers.players.audio_sources import AudioSourceSession
134 from music_assistant.helpers.json import SerializableType
135 from music_assistant.mass import MusicAssistant
136 from music_assistant.models.player import Player
137
138
139isfile = wrap(os.path.isfile)
140
141
142def _volume_normalization_preference_options() -> list[ConfigValueOption]:
143 """Return the normalization modes that can be picked as a preference."""
144 return [
145 ConfigValueOption(mode.value, title=mode.value.replace("_", " ").title())
146 for mode in VolumeNormalizationMode
147 if mode not in OUTCOME_ONLY_NORMALIZATION_MODES
148 ]
149
150
151def _audio_source_headers(session: AudioSourceSession, output_format_str: str) -> dict[str, str]:
152 """
153 Return the response headers for a live audio source stream.
154
155 Live sources are sender-paced, so they always advertise the realtime DLNA
156 flags. ``icy-name`` is sanitized of every control character, not just
157 newlines, because aiohttp rejects the rest as a header injection attempt.
158
159 :param session: The session whose source is being streamed.
160 :param output_format_str: Output format to derive the content type from.
161 """
162 return {
163 **DEFAULT_STREAM_HEADERS,
164 "icy-name": sanitize_http_header_value(session.source.name),
165 "contentFeatures.dlna.org": DLNA_CONTENT_FEATURES_REALTIME,
166 "Content-Type": get_mime_type(output_format_str),
167 }
168
169
170async def _wav_passthrough_stream(
171 audio_input: AsyncGenerator[bytes], output_format: AudioFormat
172) -> AsyncGenerator[bytes]:
173 """
174 Yield a WAV header followed by raw PCM bytes from ``audio_input``.
175
176 Closes ``audio_input`` when this generator is closed, so a provider waiting in
177 its own finally to release a claim is not left until garbage collection - a
178 reconnect would otherwise block on a claim nobody is holding on purpose.
179
180 :param audio_input: The PCM stream to pass through.
181 :param output_format: Format the WAV header should describe.
182 """
183 async with aclosing(audio_input):
184 yield create_streaming_wave_header(output_format)
185 async for chunk in audio_input:
186 yield chunk
187
188
189def _get_publish_addresses(
190 bind_ip: str, configured_publish_ip: str | None, publish_candidates: tuple[str, ...]
191) -> list[str]:
192 """
193 Return the addresses this host publishes on, best candidate first.
194
195 :param bind_ip: The configured bind IP (a wildcard means all interfaces).
196 :param configured_publish_ip: The explicitly configured publish IP, or None when auto.
197 :param publish_candidates: Host addresses reachable from the local network, ranked.
198 """
199 if configured_publish_ip:
200 # an explicitly configured address is the authoritative answer
201 return [configured_publish_ip]
202 if bind_ip and bind_ip not in WILDCARD_BIND_IPS:
203 # only one interface is served, so no other address can be reached
204 return [bind_ip]
205 # auto-detected: keep the whole ranked list - publish_ip takes the best of them and
206 # the network fingerprint watches all of them to spot an interface change
207 return list(publish_candidates)
208
209
210class StreamsController(CoreController):
211 """Controller to stream audio to players."""
212
213 domain: str = "streams"
214
215 def __init__(self, mass: MusicAssistant) -> None:
216 """Initialize instance."""
217 super().__init__(mass)
218 self._server = Webserver(self.logger, enable_dynamic_routes=True)
219 self.register_dynamic_route = self._server.register_dynamic_route
220 self.unregister_dynamic_route = self._server.unregister_dynamic_route
221 self.manifest.name = "Streamserver"
222 self.manifest.description = (
223 "Music Assistant's core controller that is responsible for "
224 "streaming audio to players on the local network."
225 )
226 self.manifest.icon = "cast-audio"
227 self.announcement_renderer = AnnouncementRenderer()
228 self.live_announcements = LiveAnnouncementManager(mass, self.logger)
229 self._bind_ip: str = "0.0.0.0"
230 self._base_url: str = ""
231 self._configured_publish_ip: str | None = None
232 # every address players may reach this host on, best candidate first; publish_ip is
233 # the first of them and the network fingerprint watches the whole list for changes
234 self._publish_addresses: list[str] = []
235 # the network as it was at the previous setup, to spot a runtime change
236 self._network_fingerprint: tuple[str, str, int, tuple[str, ...]] | None = None
237 self.audio = StreamsAudio(mass)
238 self.audio_processing = AudioProcessingManager(mass)
239 self._audio_analysis = AudioAnalysisController(self)
240 # Number of queue streams (single item or flow) actively serving a player right now.
241 # Audio analysis reads this (via audio_analysis.playback_active) to yield CPU while a
242 # queue stream is live. Announcements are a separate path that never runs analysis.
243 self._active_output_streams = 0
244
245 @property
246 def audio_analysis(self) -> AudioAnalysisController:
247 """Return the AudioAnalysisController instance."""
248 return self._audio_analysis
249
250 def output_stream_active(self) -> bool:
251 """Return whether a queue stream (single item or flow) is actively serving a player."""
252 return self._active_output_streams > 0
253
254 async def get_diagnostics(self) -> dict[str, SerializableType]:
255 """Return diagnostics info for this controller to include in diagnostics reports."""
256 return {
257 "ffmpeg_version": get_global_cache_value(CACHE_ATTR_FFMPEG_VERSION),
258 "libsoxr_support": get_global_cache_value(CACHE_ATTR_LIBSOXR_PRESENT),
259 "active_output_streams": self._active_output_streams,
260 "active_announcements": self.announcement_renderer.active_announcements,
261 "active_announcement_renders": self.announcement_renderer.active_renders,
262 "active_live_announcements": self.live_announcements.active_sessions,
263 "publish_ip_configured": self._configured_publish_ip is not None,
264 }
265
266 @property
267 def base_url(self) -> str:
268 """Return the base_url for the streamserver."""
269 return self._base_url
270
271 @property
272 def bind_ip(self) -> str:
273 """Return the IP address this streamserver is bound to."""
274 return self._bind_ip
275
276 async def get_source_ip(self, target_ip: str | None = None) -> str | None:
277 """
278 Return a local, bindable source IP on the player-facing network.
279
280 For callers that bind a socket or hand a local interface address to a helper
281 process, so their traffic leaves on the network the players live on. The result
282 is always an address of this host, never the advertised address, which may not
283 exist here at all.
284
285 Returns None when no single interface should be pinned, which the caller must
286 read as "bind all interfaces and let the routing table decide".
287
288 :param target_ip: IP address of the device the traffic is meant for. Omit it for
289 a shared consumer that serves every player at once; such a caller can only be
290 pinned by an explicitly configured bind IP.
291 """
292 if self._bind_ip and self._bind_ip not in WILDCARD_BIND_IPS:
293 if target_ip and not _same_ip_family(self._bind_ip, target_ip):
294 return None
295 return self._bind_ip
296 if not target_ip:
297 return None
298 return await get_source_ip_for_target(target_ip) or None
299
300 def get_publish_ip(self, target_ip: str) -> str | None:
301 """
302 Return the address to advertise to the device at ``target_ip``, if one is configured.
303
304 Only an explicitly configured publish IP is returned. An auto-detected one is a
305 guess at this host's primary interface, which on a multi-homed host is not
306 necessarily the network the players live on, so callers that can derive the
307 address from the connection itself must prefer that over the guess.
308
309 Returns None when no publish IP was configured, or when the configured one cannot
310 apply to this device.
311
312 :param target_ip: IP address of the device that would receive the address, used to
313 reject an address of the wrong IP family.
314 """
315 if not self._configured_publish_ip:
316 return None
317 if not _same_ip_family(self._configured_publish_ip, target_ip):
318 return None
319 return self._configured_publish_ip
320
321 @property
322 def smart_fades_available(self) -> bool:
323 """
324 Return whether smart crossfade can be used on this server.
325
326 Requires a large-enough audio buffer (at least balanced) and a loaded
327 smart fades audio analysis provider.
328 """
329 buffer_size = BufferSize(
330 self.mass.config.get_raw_core_config_value(
331 self.domain, CONF_BUFFER_SIZE, CONF_BUFFER_SIZE_DEFAULT
332 )
333 )
334 return (
335 buffer_size != BufferSize.MINIMAL and self.audio_analysis.smart_fades_provider_available
336 )
337
338 def get_crossfade_mode(self, queue: PlayerQueue) -> CrossfadeMode:
339 """
340 Return the effective crossfade mode for a queue.
341
342 Combines the per-play on/off toggle with the crossfade_mode setting and smart fades
343 availability: smart when enabled, selected and available; standard when enabled but smart
344 is not selected/available; disabled otherwise.
345 """
346 if not queue.crossfade_enabled:
347 return CrossfadeMode.DISABLED
348 # default to smart when this server can use it, else standard
349 default_mode = (
350 CrossfadeMode.SMART_CROSSFADE
351 if self.smart_fades_available
352 else CrossfadeMode.STANDARD_CROSSFADE
353 )
354 mode = self.mass.config.get_effective_player_queue_config_value(
355 queue.queue_id, CONF_CROSSFADE_MODE, default_mode
356 )
357 if mode == CrossfadeMode.SMART_CROSSFADE and self.smart_fades_available:
358 return CrossfadeMode.SMART_CROSSFADE
359 return CrossfadeMode.STANDARD_CROSSFADE
360
361 def source_normalizes_audio(self, streamdetails: StreamDetails) -> bool:
362 """
363 Return whether the item's own source already levelled this audio.
364
365 Correcting a level the source set would mean normalizing twice, the second
366 time against a measurement of its own output.
367
368 :param streamdetails: Stream details of the item.
369 """
370 # plugin providers serve playable items too, and only a music provider
371 # declares this (a plugin's live audio is handled by the media type)
372 provider = self.mass.get_provider(streamdetails.provider)
373 return isinstance(provider, MusicProvider) and provider.delivers_normalized_audio(
374 streamdetails
375 )
376
377 def get_source_crossfade_mode(self, queue: PlayerQueue, queue_item: QueueItem) -> CrossfadeMode:
378 """
379 Return the crossfade an item's own source applies, or DISABLED when none does.
380
381 A source that crossfades its own playback is handed the queue's crossfade
382 setting instead of Music Assistant mixing the overlap. This only answers for
383 audio we do not mix ourselves, and only for tracks - the same limit as applies
384 to a fade of our own.
385
386 A source already serving this item's queue answers for the item's own
387 boundaries. Before it does, the answer can only be the setting it will be
388 handed together with a boundary it would own both sides of.
389
390 :param queue: Queue the item is played from.
391 :param queue_item: Queue item to report the fade for.
392 """
393 if queue_item.media_type != MediaType.TRACK or queue_item.streamdetails is None:
394 return CrossfadeMode.DISABLED
395 if not queue_item.streamdetails.is_realtime:
396 # we mix this item's overlap ourselves, so the source's is not in play
397 return CrossfadeMode.DISABLED
398 provider = self.mass.get_provider(queue_item.streamdetails.provider)
399 if not isinstance(provider, MusicProvider):
400 return CrossfadeMode.DISABLED
401 source_fades = provider.delivers_crossfaded_audio(queue_item.streamdetails)
402 if source_fades is None:
403 # the source is not serving this queue yet, so the setting it will be
404 # handed and a boundary it would own are all there is to go on
405 queue_fades = self.get_crossfade_mode(queue) != CrossfadeMode.DISABLED
406 source_fades = queue_fades and self._source_fades_an_adjacent_item(queue, queue_item)
407 return CrossfadeMode.SOURCE if source_fades else CrossfadeMode.DISABLED
408
409 def is_smart_fades_active(self, queue: PlayerQueue) -> bool:
410 """Return whether the queue's effective crossfade mode is smart crossfade."""
411 return self.get_crossfade_mode(queue) == CrossfadeMode.SMART_CROSSFADE
412
413 async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
414 """Return all Config Entries for this core module (if any)."""
415 ip_addresses = await get_ip_addresses(include_ipv6=True)
416 return (
417 ConfigEntry(
418 key=CONF_BUFFER_SIZE,
419 type=ConfigEntryType.STRING,
420 default_value=CONF_BUFFER_SIZE_DEFAULT,
421 # Only offer presets the host's RAM can sustain (Balanced >= 4GB,
422 # Maximum >= 7GB); see get_available_buffer_sizes.
423 options=[ConfigValueOption(size.value) for size in get_available_buffer_sizes()],
424 required=False,
425 category="playback",
426 ),
427 ConfigEntry(
428 key=CONF_VOLUME_NORMALIZATION_RADIO,
429 type=ConfigEntryType.STRING,
430 default_value=DEFAULT_VOLUME_NORMALIZATION_MODE,
431 options=_volume_normalization_preference_options(),
432 category="playback",
433 ),
434 ConfigEntry(
435 key=CONF_VOLUME_NORMALIZATION_TRACKS,
436 type=ConfigEntryType.STRING,
437 default_value=DEFAULT_VOLUME_NORMALIZATION_MODE,
438 options=_volume_normalization_preference_options(),
439 category="playback",
440 ),
441 ConfigEntry(
442 key=CONF_VOLUME_NORMALIZATION_FIXED_GAIN_RADIO,
443 type=ConfigEntryType.FLOAT,
444 range=(-20, 10),
445 default_value=-6,
446 category="playback",
447 ),
448 ConfigEntry(
449 key=CONF_VOLUME_NORMALIZATION_FIXED_GAIN_TRACKS,
450 type=ConfigEntryType.FLOAT,
451 range=(-20, 10),
452 default_value=-6,
453 category="playback",
454 ),
455 CONF_ENTRY_VOLUME_NORMALIZATION_TARGET,
456 ConfigEntry(
457 key=CONF_ALLOW_CROSSFADE_SAME_ALBUM,
458 type=ConfigEntryType.BOOLEAN,
459 default_value=False,
460 category="playback",
461 ),
462 ConfigEntry(
463 key=CONF_PUBLISH_IP,
464 type=ConfigEntryType.STRING,
465 default_value=CONF_VALUE_AUTO,
466 required=False,
467 category="generic",
468 advanced=True,
469 requires_reload=True,
470 ),
471 ConfigEntry(
472 key=CONF_BIND_PORT,
473 type=ConfigEntryType.INTEGER,
474 default_value=DEFAULT_PORT,
475 category="generic",
476 advanced=True,
477 requires_reload=True,
478 ),
479 ConfigEntry(
480 key=CONF_BIND_IP,
481 type=ConfigEntryType.STRING,
482 default_value=DEFAULT_HOST,
483 options=[ConfigValueOption(x, title=x) for x in {DEFAULT_HOST, *ip_addresses}],
484 category="generic",
485 advanced=True,
486 required=False,
487 requires_reload=True,
488 ),
489 ConfigEntry(
490 key=CONF_SMART_FADES_LOG_LEVEL,
491 type=ConfigEntryType.STRING,
492 options=CONF_ENTRY_LOG_LEVEL.options,
493 default_value="GLOBAL",
494 category="audio_analysis",
495 advanced=True,
496 ),
497 ConfigEntry(
498 key=CONF_BACKGROUND_SCAN_CONCURRENCY,
499 type=ConfigEntryType.INTEGER,
500 range=(1, 16),
501 default_value=DEFAULT_BACKGROUND_SCAN_CONCURRENCY,
502 category="audio_analysis",
503 ),
504 )
505
506 async def setup(self, config: CoreConfig) -> None:
507 """Async initialize of module."""
508 # initialize the audio sub-controller (needs mass.streams to be set)
509 self.audio.setup()
510 self._audio_analysis.setup()
511 # copy log level to audio/ffmpeg loggers
512 self.audio.logger.setLevel(self.logger.level)
513 FFMPEG_LOGGER.setLevel(self.logger.level)
514 self._setup_smart_fades_logger(config)
515 # perform check for ffmpeg version
516 await check_ffmpeg_version()
517 # start the webserver
518 self.publish_port = config.get_value(CONF_BIND_PORT, DEFAULT_PORT)
519 configured_publish_ip = str(config.get_value(CONF_PUBLISH_IP) or CONF_VALUE_AUTO)
520 self._configured_publish_ip = (
521 None if configured_publish_ip == CONF_VALUE_AUTO else configured_publish_ip
522 )
523 publish_candidates = await get_publish_ip_candidates(include_ipv6=True)
524 bind_ip = str(config.get_value(CONF_BIND_IP))
525 self._resolve_publish_state(bind_ip, publish_candidates)
526 await self._server.setup(
527 bind_ip=bind_ip,
528 bind_port=cast("int", self.publish_port),
529 static_routes=[
530 (
531 "*",
532 "/flow/{session_id}/{queue_id}/{queue_item_id}/{player_id}.{fmt}",
533 self.serve_queue_flow_stream,
534 ),
535 (
536 "*",
537 "/single/{session_id}/{queue_id}/{queue_item_id}/{player_id}.{fmt}",
538 self.serve_queue_item_stream,
539 ),
540 (
541 "*",
542 "/source/{session_id}/{source_player_id}/{player_id}.{fmt}",
543 self.serve_audio_source_stream,
544 ),
545 (
546 "*",
547 "/command/{session_id}/{queue_id}/{command}.mp3",
548 self.serve_command_request,
549 ),
550 ("*", "/announcement/{player_id}.{fmt}", self.serve_announcement_stream),
551 (
552 "GET",
553 LIVE_ANNOUNCEMENT_STREAM_PATH,
554 self.live_announcements.serve_stream,
555 ),
556 ],
557 )
558 # adopt what the server actually bound to: a configured port of 0 is only resolved
559 # by the OS at bind time and an unavailable bind IP falls back to all interfaces
560 self.publish_port = cast("int", self._server.port)
561 self._resolve_publish_state(self._server.bind_ip or DEFAULT_HOST, publish_candidates)
562 # print a big fat message in the log where the streamserver is running
563 # because this is a common source of issues for people with more complex setups
564 self.logger.log(
565 logging.INFO if self.mass.config.onboard_done else logging.WARNING,
566 "\n\n################################################################################\n"
567 "Started streamserver on %s:%s\n"
568 "This is the IP address that is communicated to players.\n"
569 "If this is incorrect, audio will not play!\n"
570 "See the documentation for how to configure the publish IP for the Streamserver\n"
571 "in Settings --> System --> Streams\n"
572 "################################################################################\n",
573 self.publish_ip,
574 self.publish_port,
575 )
576 await self._reload_network_dependent_providers()
577
578 async def post_setup(self) -> None:
579 """Handle logic after all core controllers have been set up."""
580 # the inbound half of a live announcement rides on the webserver: it is the only
581 # one of the two servers that authenticates (and that browsers reach over https)
582 self.live_announcements.setup()
583
584 async def close(self) -> None:
585 """Cleanup on exit."""
586 await self._audio_analysis.close()
587 await self.live_announcements.close()
588 await self._server.close()
589
590 async def resolve_stream_url(self, player_id: str, media: PlayerMedia) -> str:
591 """
592 Resolve the stream URL for the given PlayerMedia.
593
594 :param player_id: The (protocol) player ID requesting the stream.
595 :param media: The PlayerMedia object for which to resolve the stream URL.
596 :return: The resolved stream URL as a string.
597 """
598 if media.media_type in (MediaType.ANNOUNCEMENT, MediaType.FLOW_STREAM):
599 return media.uri
600 protocol_player = self.mass.players.get_player(player_id)
601 conf_output_codec = cast(
602 "str",
603 protocol_player.config.get_value(CONF_OUTPUT_CODEC, default="flac")
604 if protocol_player
605 else "flac",
606 )
607 prefer_wav_for_live_sources = (
608 media.media_type == MediaType.AUDIO_SOURCE
609 and protocol_player is not None
610 and cast(
611 "bool",
612 protocol_player.config.get_value(CONF_PREFER_WAV_FOR_LIVE_SOURCES, default=False),
613 )
614 )
615 output_codec = (
616 ContentType.WAV
617 if prefer_wav_for_live_sources
618 else ContentType.try_parse(conf_output_codec)
619 )
620 fmt = output_codec.value
621 # handle raw pcm without exact format specifiers
622 if output_codec.is_pcm() and ";" not in fmt:
623 fmt += f";codec=pcm;rate={44100};bitrate={16};channels={2}"
624 if media.media_type == MediaType.AUDIO_SOURCE and not media.queue_item_id:
625 # a source playing on a player, rather than an item in a queue
626 if not media.source_id or not media.queue_session_id:
627 raise InvalidDataError("Can not resolve stream URL: Invalid PlayerMedia data")
628 return (
629 f"{self.base_url}/source/{media.queue_session_id}"
630 f"/{media.source_id}/{player_id}.{fmt}"
631 )
632 session_id = media.queue_session_id
633 queue_item_id = media.queue_item_id
634 if not session_id or not queue_item_id:
635 raise InvalidDataError("Can not resolve stream URL: Invalid PlayerMedia data")
636 queue_id = media.source_id
637 queue = self.mass.player_queues.get(queue_id) if queue_id else None
638 crossfade_needs_flow_mode = (
639 # crossfade only applies to tracks; if the queue has it enabled but the player(protocol)
640 # does not support gapless playback, we need to enforce flow mode
641 media.media_type == MediaType.TRACK
642 and queue is not None
643 and queue.crossfade_enabled
644 and protocol_player
645 and not protocol_player.supports_gapless
646 )
647 # the audio overlay is mixed into the queue's continuous (flow) stream;
648 # per-item requests would restart the overlay at every track boundary
649 overlay_needs_flow_mode = queue is not None and overlay_active(queue)
650 # Determine flow_mode based on the actual player's capabilities.
651 # This is done here (just-in-time) because the player's protocol determines this
652 flow_mode = (
653 protocol_player is not None
654 and (protocol_player.flow_mode or crossfade_needs_flow_mode or overlay_needs_flow_mode)
655 and media.media_type not in (MediaType.RADIO, MediaType.AUDIO_SOURCE)
656 )
657 base_path = "flow" if flow_mode else "single"
658 return (
659 f"{self.base_url}/{base_path}/{session_id}/{queue_id}/{queue_item_id}/{player_id}.{fmt}"
660 )
661
662 async def serve_queue_item_stream(self, request: web.Request) -> web.StreamResponse: # noqa: PLR0915
663 """Stream single queueitem audio to a player."""
664 self._log_request(request)
665 queue_id = request.match_info["queue_id"]
666 player_id = request.match_info["player_id"]
667 if not (queue := self.mass.player_queues.get(queue_id)):
668 raise web.HTTPNotFound(reason=f"Unknown Queue: {queue_id}")
669 session_id = request.match_info["session_id"]
670 pq_data = self.mass.player_queues.queue_data(queue.queue_id)
671 if pq_data.session_id is None or session_id != pq_data.session_id:
672 raise web.HTTPNotFound(reason=f"Unknown (or invalid) session: {session_id}")
673 if not (player := self.mass.players.get_player(player_id)):
674 raise web.HTTPNotFound(reason=f"Unknown Player: {player_id}")
675 queue_item_id = request.match_info["queue_item_id"]
676 queue_item = self.mass.player_queues.get_item(queue_id, queue_item_id)
677 if not queue_item:
678 raise web.HTTPNotFound(reason=f"Unknown Queue item: {queue_item_id}")
679
680 is_audio_source = (
681 queue_item.media_item is not None
682 and queue_item.media_item.media_type == MediaType.AUDIO_SOURCE
683 )
684
685 # HEAD probes for AudioSource items return a minimal response without
686 # touching the plugin. on_source_selected is the lifecycle hook that
687 # claims ownership and fires off transfer/handoff side effects (stop
688 # the previous player, redirect on disallowed switch, etc.), and a
689 # renderer probing with HEAD before GET should not trigger any of
690 # that. The actual GET request goes through the full hook chain.
691 if request.method != "GET" and is_audio_source:
692 # Validate the providing plugin still exists before advertising the
693 # source. Many DLNA renderers cache HEAD responses; returning 200
694 # for a URI whose plugin has been unloaded would lie to the
695 # renderer and the follow-up GET would fail unrecoverably.
696 assert queue_item.media_item is not None
697 if not isinstance(
698 self.mass.get_provider(queue_item.media_item.provider), PluginProvider
699 ):
700 raise web.HTTPNotFound(
701 reason=f"AudioSource provider {queue_item.media_item.provider} unavailable"
702 )
703 # For PCM-fmt URLs, advertise audio/wav in HEAD: most DLNA renderers
704 # key off the HEAD Content-Type to pick a decoder and do not handle
705 # raw PCM (application/octet-stream). The actual GET response will
706 # still wrap the bytes into a WAV container if needed via the same
707 # mime-type translation downstream.
708 head_fmt = request.match_info["fmt"]
709 if ContentType.try_parse(head_fmt).is_pcm():
710 head_fmt = ContentType.WAV.value
711 headers = {
712 **DEFAULT_STREAM_HEADERS,
713 "icy-name": sanitize_http_header_value(queue_item.name),
714 "contentFeatures.dlna.org": DLNA_CONTENT_FEATURES_REALTIME,
715 "Content-Type": get_mime_type(head_fmt),
716 }
717 resp = web.StreamResponse(status=200, reason="OK", headers=headers)
718 await resp.prepare(request)
719 return resp
720
721 # Fire on_source_selected hook for every AudioSource GET — this is the
722 # single point where exclusive plugin sources claim ownership. Firing
723 # unconditionally (regardless of whether streamdetails are cached from
724 # a previous request) means a disconnect/reconnect for the same queue
725 # item re-claims the lock with a fresh session id, instead of streaming
726 # against the stale ownership of the prior request.
727 # Source identity comes from queue_item.media_item because streamdetails
728 # may not exist yet on the first request.
729 # stream_session_id is a fresh per-request token threaded through to
730 # on_source_unselected so the provider can distinguish a stale
731 # teardown (e.g. a same-queue reconnect's first request completing
732 # AFTER its replacement has already started streaming) from the
733 # currently active session's real teardown.
734 audio_source_provider: PluginProvider | None = None
735 audio_source_id: str | None = None
736 stream_session_id = uuid4().hex
737 if (
738 is_audio_source
739 and queue_item.media_item is not None
740 and (prov := self.mass.get_provider(queue_item.media_item.provider))
741 and isinstance(prov, PluginProvider)
742 ):
743 audio_source_id = queue_item.media_item.item_id
744 # Wire the provider into the finally block BEFORE awaiting the
745 # hook: if the provider partially mutates state (claims the lock,
746 # records the session id) and then raises a non-RuntimeError
747 # exception (buggy plugin, asyncio.CancelledError, etc.), the
748 # finally must still fire on_source_unselected so the lock gets
749 # released. The provider's session-id guard makes a spurious
750 # release a no-op if the lock was never actually claimed.
751 audio_source_provider = prov
752 try:
753 await prov.on_source_selected(
754 audio_source_id, player_id, queue_id, stream_session_id
755 )
756 except RuntimeError as err:
757 # Provider intentionally aborts the original request (e.g.
758 # allow_player_switch=False has just redirected play_media to
759 # the configured target). Surface as 404 so the disallowed
760 # player drops the connection cleanly instead of treating an
761 # uncaught 500 as transient and retrying. The provider
762 # contract requires raising BEFORE claiming, so this is a
763 # clean abort — but we still let the finally run, where the
764 # session-id guard makes the unselect a no-op.
765 self.logger.info(
766 "AudioSource %s aborted stream for player %s: %s",
767 audio_source_id,
768 player_id,
769 err,
770 )
771 raise web.HTTPNotFound(reason=str(err))
772
773 try:
774 if not queue_item.streamdetails:
775 try:
776 queue_item.streamdetails = await self.audio.get_stream_details(
777 queue_item=queue_item
778 )
779 except Exception as e:
780 self.logger.error(
781 "Failed to get streamdetails for QueueItem %s: %s", queue_item_id, e
782 )
783 # a source capacity miss is transient, the item itself is fine
784 if not isinstance(e, ProviderStreamLimitError):
785 queue_item.available = False
786 raise web.HTTPNotFound(
787 reason=f"No streamdetails for Queue item: {queue_item_id}"
788 )
789
790 standard_crossfade_duration = self.mass.config.get_raw_core_config_value(
791 CONF_PLAYER_QUEUES, CONF_CROSSFADE_DURATION, 8
792 )
793 if queue_item.media_type != MediaType.TRACK:
794 crossfade_mode = CrossfadeMode.DISABLED
795 elif queue_item.streamdetails.is_realtime:
796 # a realtime source delivers at playback pace, so it has no audio to
797 # spare for an overlap in either direction
798 crossfade_mode = CrossfadeMode.DISABLED
799 else:
800 crossfade_mode = self.get_crossfade_mode(queue)
801 if (
802 crossfade_mode != CrossfadeMode.DISABLED
803 and PlayerFeature.GAPLESS_PLAYBACK not in player.state.supported_features
804 ):
805 self.logger.warning(
806 "Crossfade disabled: Player %s does not support gapless playback, "
807 "consider enabling flow mode to enable crossfade on this player.",
808 player.state.name,
809 )
810 crossfade_mode = CrossfadeMode.DISABLED
811
812 # pick output format based on the streamdetails and player capabilities
813 pcm_format = await self.audio.select_pcm_format(
814 player=player,
815 streamdetails=queue_item.streamdetails,
816 crossfade_enabled=crossfade_mode != CrossfadeMode.DISABLED,
817 overlay_active=(queue_item.media_type == MediaType.RADIO and overlay_active(queue)),
818 )
819 output_format = await self.audio.get_output_format(
820 output_format_str=request.match_info["fmt"],
821 player=player,
822 content_sample_rate=pcm_format.sample_rate,
823 content_bit_depth=pcm_format.bit_depth,
824 media_type=queue_item.media_type,
825 )
826
827 # prepare request, add some DLNA/UPNP compatible headers
828 # icy-name is sanitized (all control chars, not just newlines) to avoid a
829 # "Potential header injection attack" ValueError by aiohttp
830 # see https://github.com/music-assistant/support/issues/4913
831 # and https://github.com/music-assistant/support/issues/5791
832 # use realtime DLNA flags for radio (sender-paced) since the source delivers slowly
833 dlna_features = (
834 DLNA_CONTENT_FEATURES_REALTIME
835 if queue_item.media_type != MediaType.TRACK
836 else DLNA_CONTENT_FEATURES
837 )
838 headers = {
839 **DEFAULT_STREAM_HEADERS,
840 "icy-name": sanitize_http_header_value(queue_item.name),
841 "contentFeatures.dlna.org": dlna_features,
842 "Content-Type": get_mime_type(output_format.output_format_str),
843 }
844
845 resp = web.StreamResponse(status=200, reason="OK", headers=headers)
846 resp.content_type = get_mime_type(output_format.output_format_str)
847 http_profile = player.get_config_value(CONF_HTTP_PROFILE, "default")
848 if http_profile == "forced_content_length" and not queue_item.duration:
849 # just set an insane high content length to make sure the player keeps playing
850 resp.content_length = calculate_content_length(output_format, 12 * 3600)
851 elif http_profile == "forced_content_length" and queue_item.duration:
852 # estimate content length based on effective duration
853 # account for seek position (e.g., crossfade from previous track)
854 seek_pos = queue_item.streamdetails.seek_position if queue_item.streamdetails else 0
855 effective_duration = max(queue_item.duration - seek_pos, 1)
856 # use cached actual bytes-per-second if available (from a previous stream)
857 resp.content_length = await get_content_length(
858 self.mass, queue_item.uri, output_format, effective_duration
859 )
860 elif http_profile == "chunked":
861 resp.enable_chunked_encoding()
862
863 await resp.prepare(request)
864
865 # return early if this is not a GET request
866 if request.method != "GET":
867 return resp
868
869 self._update_audio_processing_context(
870 queue=queue,
871 queue_item=queue_item,
872 pcm_format=pcm_format,
873 overlay_enabled=(
874 queue_item.media_type == MediaType.RADIO and overlay_active(queue)
875 ),
876 session_id=session_id,
877 # a crossfade the source performs itself is never mixed here, so it
878 # is reported apart from the mode that drives our own mixer
879 source_crossfade_mode=self.get_source_crossfade_mode(queue, queue_item),
880 )
881
882 if crossfade_mode != CrossfadeMode.DISABLED:
883 # crossfade is enabled, use special crossfaded single item stream
884 # where the crossfade of the next track is present in the stream of
885 # a single track. This only works if the player supports gapless playback!
886 audio_input = self.audio.get_queue_item_stream_with_smartfade(
887 player=player,
888 queue_item=queue_item,
889 pcm_format=pcm_format,
890 crossfade_mode=crossfade_mode,
891 standard_crossfade_duration=standard_crossfade_duration,
892 session_id=session_id,
893 )
894 else:
895 # no crossfade, just a regular single item stream
896 audio_input = self.audio.get_queue_item_stream(
897 queue_item=queue_item,
898 pcm_format=pcm_format,
899 seek_position=int(queue_item.streamdetails.seek_position),
900 playback_speed=cast(
901 "float", queue_item.extra_attributes.get("playback_speed", 1.0)
902 ),
903 session_id=session_id,
904 )
905 if queue_item.media_type == MediaType.RADIO and overlay_active(queue):
906 # radio plays as a single long-lived stream (never in flow mode),
907 # so mix the audio overlay in here
908 audio_input = self.audio.get_overlay_mixed_stream(queue, audio_input, pcm_format)
909 # stream the audio
910 # this final ffmpeg process in the chain converts raw lossless PCM into
911 # the desired output format for the player including any player specific
912 # filter params such as channels mixing, DSP, resampling and, only if
913 # needed, encoding to lossy formats
914 output_plan = self.audio.get_player_output_plan(
915 player_id=player.player_id,
916 input_format=pcm_format,
917 output_format=output_format,
918 shared_player_ids=player.state.group_members,
919 queue_id=queue_id,
920 session_id=session_id,
921 queue_item_id=queue_item.queue_item_id,
922 )
923 filter_params = output_plan.filter_params
924 # Fast path for live AudioSource: when the player accepts WAV at the
925 # source's exact PCM rate/depth/channels and no filters apply, we
926 # skip the encode ffmpeg entirely and just stream a WAV header
927 # followed by the raw PCM bytes — saves an ffmpeg process and the
928 # latency of its internal buffer on every realtime stream.
929 audio_bytes: AsyncGenerator[bytes]
930 if (
931 queue_item.media_type == MediaType.AUDIO_SOURCE
932 and output_format.content_type == ContentType.WAV
933 and not filter_params
934 and output_format.sample_rate == pcm_format.sample_rate
935 and output_format.bit_depth == pcm_format.bit_depth
936 and output_format.channels == pcm_format.channels
937 ):
938 audio_bytes = _wav_passthrough_stream(audio_input, output_format)
939 else:
940 audio_bytes = get_ffmpeg_stream(
941 audio_input=audio_input,
942 input_format=pcm_format,
943 output_format=output_format,
944 filter_params=filter_params,
945 extra_input_args=[
946 "-readrate",
947 SINGLE_ITEM_READRATE,
948 "-readrate_initial_burst",
949 SINGLE_ITEM_READRATE_INITIAL_BURST,
950 ],
951 )
952 first_chunk_received = False
953 bytes_sent = 0
954 # Mark this player as actively streaming so audio analysis yields CPU to playback
955 # for the duration of the transfer (see audio_analysis.playback_active).
956 self._active_output_streams += 1
957 try:
958 # aclosing guarantees the generator (and thus the ffmpeg process chain
959 # behind it) is torn down immediately when the player disconnects
960 # mid-stream, instead of lingering until garbage collection finalizes
961 # the abandoned generator.
962 async with aclosing(audio_bytes):
963 async for chunk in audio_bytes:
964 if pq_data.session_id != session_id:
965 # playback moved on (or stopped) while this response was open;
966 # the flow path checks the same thing per chunk
967 self.logger.debug(
968 "Ending stream for %s: session %s is no longer current",
969 queue_item.name,
970 session_id,
971 )
972 break
973 try:
974 await resp.write(chunk)
975 bytes_sent += len(chunk)
976 if not first_chunk_received:
977 first_chunk_received = True
978 # inform the queue that the track is now loaded in the buffer
979 # so for example the next track can be enqueued
980 self.mass.player_queues.track_loaded_in_buffer(
981 queue_item.queue_id, queue_item.queue_item_id
982 )
983 except (BrokenPipeError, ConnectionResetError, ConnectionError) as err:
984 if (
985 first_chunk_received
986 and not player.stop_called
987 and queue_item.streamdetails.duration # ignore for radio streams
988 ):
989 # Player disconnected (unexpected) after receiving at least
990 # some data. This could indicate buffering issues, network
991 # problems, or player-specific issues.
992 self.logger.warning(
993 "Player %s disconnected prematurely from stream for %s (%s) - "
994 "error: %s, sent %d bytes, content_length=%s",
995 queue.display_name,
996 queue_item.name,
997 queue_item.uri,
998 err.__class__.__name__,
999 bytes_sent,
1000 resp.content_length,
1001 )
1002 break
1003 finally:
1004 self._active_output_streams -= 1
1005 if queue_item.streamdetails.stream_error:
1006 self.logger.error(
1007 "Error streaming QueueItem %s (%s) to %s",
1008 queue_item.name,
1009 queue_item.uri,
1010 queue.display_name,
1011 )
1012 elif (
1013 bytes_sent > 0
1014 and queue_item.streamdetails
1015 and queue_item.streamdetails.seconds_streamed
1016 and queue_item.duration
1017 ):
1018 # cache the actual encoded bytes-per-second for this URI + output format
1019 # so future content_length estimates are near-exact
1020 self.mass.create_task(
1021 store_content_length_in_cache(
1022 self.mass,
1023 queue_item.uri,
1024 output_format,
1025 bytes_sent,
1026 queue_item.streamdetails.seconds_streamed,
1027 )
1028 )
1029 return resp
1030 finally:
1031 # Paired with on_source_selected — fires regardless of how streaming
1032 # ended (normal completion, client disconnect, exception). Lets
1033 # NAMED_PIPE plugins release ownership without depending on an
1034 # external session event. The stream_session_id is the same token
1035 # passed to on_source_selected; the provider must reject the
1036 # callback if it does not match the currently stored active
1037 # session (otherwise a stale teardown from a superseded same-queue
1038 # request would clear the live claim of its replacement).
1039 if audio_source_provider is not None and audio_source_id is not None:
1040 # Provider teardown failures must not break the response cycle
1041 # (we're already in finally for a stream that ended one way or
1042 # another), but they MUST surface in logs — otherwise a buggy
1043 # plugin leaks _in_use_by_queue forever and there is no trail.
1044 try:
1045 await audio_source_provider.on_source_unselected(
1046 audio_source_id, queue_id, stream_session_id
1047 )
1048 except Exception:
1049 self.logger.warning(
1050 "on_source_unselected raised for provider %s source %s queue %s",
1051 audio_source_provider.instance_id,
1052 audio_source_id,
1053 queue_id,
1054 exc_info=True,
1055 )
1056
1057 async def serve_audio_source_stream(self, request: web.Request) -> web.StreamResponse:
1058 """Stream a live AudioSource playing on a player."""
1059 self._log_request(request)
1060 session, player, prov = self._resolve_audio_source_request(request)
1061 playback_session_id = session.playback_session_id
1062 # the session's own player, never the url's: the consuming player differs for
1063 # protocol and group members, and the claim belongs to the owner
1064 source_player_id = session.player_id
1065
1066 # A renderer probing with HEAD must not trigger the selection side effects
1067 # on_source_selected fires (stopping the previous player, redirecting a
1068 # disallowed switch), so answer it without touching the plugin.
1069 if request.method != "GET":
1070 return await self._serve_audio_source_head(request, session)
1071
1072 stream_session_id = uuid4().hex
1073 # wire the provider in before awaiting the hook: a plugin that claims the
1074 # source and then raises must still get its release
1075 claimed = False
1076 serving = False
1077 try:
1078 try:
1079 claimed = True
1080 await prov.on_source_selected(
1081 # deliberately the owner for both: providers store this id to stop
1082 # or re-target the player later, and the url's player can be an
1083 # ephemeral protocol bridge whose id is invalid by then
1084 session.source_id,
1085 source_player_id,
1086 source_player_id,
1087 stream_session_id,
1088 )
1089 if (
1090 self.mass.players.get_audio_source_session(source_player_id) is not session
1091 or session.playback_session_id != playback_session_id
1092 ):
1093 raise web.HTTPNotFound(reason="AudioSource session was superseded")
1094 session.stream_session_id = stream_session_id
1095 except RuntimeError as err:
1096 # the plugin refuses this player (e.g. it just redirected playback
1097 # elsewhere); a 404 makes the renderer drop the connection instead of
1098 # retrying a 500 as transient
1099 self.logger.info(
1100 "AudioSource %s aborted stream for player %s: %s",
1101 session.source_id,
1102 player.player_id,
1103 err,
1104 )
1105 raise web.HTTPNotFound(reason=str(err)) from err
1106
1107 if (streamdetails := session.streamdetails) is None:
1108 try:
1109 streamdetails = await prov.get_stream_details(
1110 session.source_id, MediaType.AUDIO_SOURCE
1111 )
1112 except Exception as err:
1113 self.logger.error(
1114 "Failed to get streamdetails for AudioSource %s: %s",
1115 session.source_id,
1116 err,
1117 )
1118 raise web.HTTPNotFound(reason="Failed to get stream details") from err
1119 session.attach_streamdetails(streamdetails)
1120
1121 resp, audio_bytes = await self._prepare_audio_source_stream(
1122 request=request,
1123 player=player,
1124 session=session,
1125 streamdetails=streamdetails,
1126 )
1127 serving = True
1128 self._active_output_streams += 1
1129 try:
1130 async with aclosing(audio_bytes):
1131 async for chunk in audio_bytes:
1132 if (
1133 self.mass.players.get_audio_source_session(source_player_id)
1134 is not session
1135 or session.playback_session_id != playback_session_id
1136 or session.stream_session_id != stream_session_id
1137 ):
1138 self.logger.debug(
1139 "Ending stream for %s: a newer request took the source over",
1140 session.source.name,
1141 )
1142 break
1143 try:
1144 await resp.write(chunk)
1145 except BrokenPipeError, ConnectionResetError, ConnectionError:
1146 break
1147 finally:
1148 self._active_output_streams -= 1
1149 return resp
1150 finally:
1151 if claimed:
1152 try:
1153 await prov.on_source_unselected(
1154 session.source_id, source_player_id, stream_session_id
1155 )
1156 except Exception:
1157 self.logger.warning(
1158 "on_source_unselected raised for provider %s source %s player %s",
1159 prov.instance_id,
1160 session.source_id,
1161 source_player_id,
1162 exc_info=True,
1163 )
1164 if not serving:
1165 await self._release_unstarted_audio_source(session, playback_session_id)
1166
1167 async def serve_queue_flow_stream(self, request: web.Request) -> web.StreamResponse: # noqa: PLR0915
1168 """Stream Queue Flow audio to player."""
1169 self._log_request(request)
1170 queue_id = request.match_info["queue_id"]
1171 player_id = request.match_info["player_id"]
1172 if not (queue := self.mass.player_queues.get(queue_id)):
1173 raise web.HTTPNotFound(reason=f"Unknown Queue: {queue_id}")
1174 session_id = request.match_info["session_id"]
1175 queue_data = self.mass.player_queues.queue_data(queue_id)
1176 if queue_data.session_id is None or session_id != queue_data.session_id:
1177 raise web.HTTPNotFound(reason=f"Unknown (or invalid) session: {session_id}")
1178 if not (player := self.mass.players.get_player(player_id)):
1179 raise web.HTTPNotFound(reason=f"Unknown Player: {player_id}")
1180 start_queue_item_id = request.match_info["queue_item_id"]
1181 start_queue_item = self.mass.player_queues.get_item(queue_id, start_queue_item_id)
1182 if not start_queue_item:
1183 raise web.HTTPNotFound(reason=f"Unknown Queue item: {start_queue_item_id}")
1184
1185 # select the PCM format for the flow stream, anchored on the first track
1186 crossfade_mode = (
1187 self.get_crossfade_mode(queue)
1188 if start_queue_item.media_type == MediaType.TRACK
1189 else CrossfadeMode.DISABLED
1190 )
1191 flow_pcm_format = await self.audio.select_flow_pcm_format(
1192 player,
1193 start_streamdetails=start_queue_item.streamdetails,
1194 crossfade_enabled=crossfade_mode != CrossfadeMode.DISABLED,
1195 overlay_active=overlay_active(queue),
1196 )
1197
1198 # work out output format/details
1199 output_format = await self.audio.get_output_format(
1200 output_format_str=request.match_info["fmt"],
1201 player=player,
1202 content_sample_rate=flow_pcm_format.sample_rate,
1203 content_bit_depth=flow_pcm_format.bit_depth,
1204 media_type=start_queue_item.media_type,
1205 )
1206 # work out ICY metadata support
1207 icy_preference = self.mass.config.get_raw_player_config_value(
1208 player_id,
1209 CONF_ENTRY_ENABLE_ICY_METADATA.key,
1210 CONF_ENTRY_ENABLE_ICY_METADATA.default_value,
1211 )
1212 enable_icy = request.headers.get("Icy-MetaData", "") == "1" and icy_preference != "disabled"
1213 icy_meta_interval = 256000 if icy_preference == "full" else 16384
1214
1215 # prepare request, add some DLNA/UPNP compatible headers.
1216 # icy-name (in DEFAULT_STREAM_HEADERS) is always present so players have a
1217 # readable stream name; the rest of the ICY/shoutcast metadata headers are
1218 # only advertised when the client actually requested ICY metadata, rather
1219 # than on every flow response.
1220 headers = {
1221 **DEFAULT_STREAM_HEADERS,
1222 **(ICY_HEADERS if enable_icy else {}),
1223 "contentFeatures.dlna.org": DLNA_CONTENT_FEATURES_REALTIME,
1224 "Content-Type": get_mime_type(output_format.output_format_str),
1225 }
1226 if enable_icy:
1227 headers["icy-metaint"] = str(icy_meta_interval)
1228
1229 resp = web.StreamResponse(status=200, reason="OK", headers=headers)
1230 http_profile = player.get_config_value(CONF_HTTP_PROFILE, "default")
1231 if http_profile == "forced_content_length":
1232 # just set an insane high content length to make sure the player keeps playing
1233 resp.content_length = calculate_content_length(output_format, 12 * 3600)
1234 elif http_profile == "chunked":
1235 resp.enable_chunked_encoding()
1236
1237 await resp.prepare(request)
1238
1239 # return early if this is not a GET request
1240 if request.method != "GET":
1241 return resp
1242
1243 self._update_audio_processing_context(
1244 queue=queue,
1245 queue_item=start_queue_item,
1246 pcm_format=flow_pcm_format,
1247 overlay_enabled=overlay_active(queue),
1248 session_id=session_id,
1249 )
1250 output_plan = self.audio.get_player_output_plan(
1251 player.player_id,
1252 flow_pcm_format,
1253 output_format,
1254 shared_player_ids=player.state.group_members,
1255 queue_id=queue_id,
1256 session_id=session_id,
1257 )
1258
1259 # all checks passed, start streaming!
1260 # this final ffmpeg process in the chain will convert the raw, lossless PCM audio into
1261 # the desired output format for the player including any player specific filter params
1262 # such as channels mixing, DSP, resampling and, only if needed, encoding to lossy formats
1263 self.logger.debug("Start serving Queue flow audio stream for %s", queue.display_name)
1264
1265 # Mark this player as actively streaming so audio analysis yields CPU to playback
1266 # for the duration of the flow stream (see audio_analysis.playback_active).
1267 self._active_output_streams += 1
1268 flow_stream = self.audio.get_queue_flow_stream(
1269 queue=queue,
1270 start_queue_item=start_queue_item,
1271 pcm_format=flow_pcm_format,
1272 session_id=session_id,
1273 protocol_player=player,
1274 )
1275 if overlay_active(queue):
1276 flow_stream = self.audio.get_overlay_mixed_stream(queue, flow_stream, flow_pcm_format)
1277 audio_bytes = get_ffmpeg_stream(
1278 audio_input=flow_stream,
1279 input_format=flow_pcm_format,
1280 output_format=output_format,
1281 filter_params=output_plan.filter_params,
1282 # we need to slowly feed the music to avoid the player stopping and later
1283 # restarting (or completely failing) the audio stream by keeping the buffer short.
1284 # this is reported to be an issue especially with Chromecast players.
1285 # see for example: https://github.com/music-assistant/support/issues/3717
1286 # allow buffer ahead of a few seconds and read rest in (near) realtime
1287 extra_input_args=["-readrate", "1.05", "-readrate_initial_burst", "5"],
1288 chunk_size=icy_meta_interval if enable_icy else calculate_content_length(output_format),
1289 )
1290 client_disconnected = False
1291 try:
1292 # aclosing guarantees the flow stream (and thus the ffmpeg process chain
1293 # behind it) is torn down immediately when the player disconnects
1294 # mid-stream, instead of lingering until garbage collection finalizes
1295 # the abandoned generator.
1296 async with aclosing(audio_bytes):
1297 async for chunk in audio_bytes:
1298 try:
1299 await resp.write(chunk)
1300 except BrokenPipeError, ConnectionResetError, ConnectionError:
1301 # race condition
1302 client_disconnected = True
1303 break
1304
1305 if not enable_icy:
1306 continue
1307
1308 # if icy metadata is enabled, send the icy metadata after the chunk
1309 if (
1310 # use current item here and not buffered item, otherwise
1311 # the icy metadata will be too much ahead
1312 (current_item := queue.current_item)
1313 and current_item.streamdetails
1314 and current_item.streamdetails.stream_title
1315 ):
1316 title = current_item.streamdetails.stream_title
1317 elif queue and current_item and current_item.name:
1318 title = current_item.name
1319 else:
1320 title = "Music Assistant"
1321 metadata = f"StreamTitle='{title}';".encode()
1322 if icy_preference == "full" and current_item and current_item.image:
1323 metadata += f"StreamURL='{current_item.image.path}'".encode()
1324 while len(metadata) % 16 != 0:
1325 metadata += b"\x00"
1326 length = len(metadata)
1327 length_b = chr(int(length / 16)).encode()
1328 await resp.write(length_b + metadata)
1329 finally:
1330 self._active_output_streams -= 1
1331
1332 if not client_disconnected and http_profile == "forced_content_length":
1333 await self._finish_flow_stream(resp, queue_id, session_id)
1334
1335 return resp
1336
1337 async def serve_command_request(self, request: web.Request) -> web.FileResponse:
1338 """Handle special 'command' request for a player."""
1339 self._log_request(request)
1340 queue_id = request.match_info["queue_id"]
1341 session_id = request.match_info["session_id"]
1342 queue_data = self.mass.player_queues.queue_data_or_none(queue_id)
1343 if queue_data is None or queue_data.session_id != session_id:
1344 raise web.HTTPNotFound(reason=f"Unknown (or invalid) session: {session_id}")
1345 command = request.match_info["command"]
1346 if command == "next":
1347 self.mass.create_task(self.mass.player_queues.next(queue_id))
1348 return web.FileResponse(SILENCE_FILE, headers={"icy-name": "Music Assistant"})
1349
1350 async def serve_announcement_stream(self, request: web.Request) -> web.StreamResponse:
1351 """Stream announcement audio to a player."""
1352 self._log_request(request)
1353 player_id = request.match_info["player_id"]
1354 if not (player := self.mass.players.get_player(player_id)):
1355 raise web.HTTPNotFound(reason=f"Unknown Player: {player_id}")
1356 if not (announce_data := self.announcement_renderer.get_for_player(player_id)):
1357 raise web.HTTPNotFound(reason=f"No pending announcements for Player: {player_id}")
1358
1359 # work out output format/details
1360 fmt = request.match_info["fmt"]
1361 audio_format = AudioFormat(content_type=ContentType.try_parse(fmt))
1362
1363 http_profile = self._get_announcement_http_profile(player_id, announce_data)
1364
1365 # return early if this is not a GET request:
1366 # players often probe the url with a HEAD request before fetching it and
1367 # rendering the announcement for such a probe would run the entire (costly)
1368 # TTS/ffmpeg chain twice for a single announcement.
1369 if request.method != "GET":
1370 resp = web.StreamResponse(status=200, reason="OK", headers=DEFAULT_STREAM_HEADERS)
1371 resp.content_type = get_mime_type(audio_format.output_format_str)
1372 if http_profile == "chunked":
1373 resp.enable_chunked_encoding()
1374 await resp.prepare(request)
1375 return resp
1376
1377 if http_profile == "forced_content_length":
1378 # given the fact that an announcement is just a short audio clip,
1379 # just send it over completely at once so we have a fixed content length
1380 data = bytearray()
1381 announcement_stream = self.get_announcement_stream(announce_data, audio_format)
1382 # aclosing guarantees the stream (and thus the ffmpeg process chain behind
1383 # it) is torn down immediately when the request is cancelled, instead of
1384 # lingering until garbage collection finalizes the abandoned generator.
1385 async with aclosing(announcement_stream):
1386 async for chunk in announcement_stream:
1387 data += chunk
1388 return web.Response(
1389 body=bytes(data),
1390 content_type=get_mime_type(audio_format.output_format_str),
1391 headers=DEFAULT_STREAM_HEADERS,
1392 )
1393
1394 resp = web.StreamResponse(status=200, reason="OK", headers=DEFAULT_STREAM_HEADERS)
1395 resp.content_type = get_mime_type(audio_format.output_format_str)
1396 if http_profile == "chunked":
1397 resp.enable_chunked_encoding()
1398
1399 await resp.prepare(request)
1400
1401 # all checks passed, start streaming!
1402 self.logger.debug(
1403 "Start serving audio stream for Announcement %s to %s",
1404 announce_data["announcement_url"],
1405 player.display_name,
1406 )
1407 announcement_stream = self.get_announcement_stream(announce_data, audio_format)
1408 # aclosing guarantees the stream (and thus the ffmpeg process chain behind
1409 # it) is torn down immediately when the player disconnects mid-stream,
1410 # instead of lingering until garbage collection finalizes the abandoned
1411 # generator.
1412 async with aclosing(announcement_stream):
1413 async for chunk in announcement_stream:
1414 try:
1415 await resp.write(chunk)
1416 except BrokenPipeError, ConnectionResetError:
1417 break
1418
1419 self.logger.debug(
1420 "Finished serving audio stream for Announcement %s to %s",
1421 announce_data["announcement_url"],
1422 player.display_name,
1423 )
1424
1425 return resp
1426
1427 def get_command_url(self, player_or_queue_id: str, command: str) -> str | None:
1428 """
1429 Get the url for the special command stream, or None if the queue is not playing.
1430
1431 :param player_or_queue_id: Queue (or player) to send the command to.
1432 :param command: Command the url triggers when fetched.
1433 """
1434 # resolve to the active queue: a protocol player (e.g. the cast child of a
1435 # universal player) does not own the active queue, its parent player does
1436 if active_queue := self.mass.player_queues.get_active_queue(player_or_queue_id):
1437 queue_id = active_queue.queue_id
1438 else:
1439 queue_id = player_or_queue_id
1440 queue_data = self.mass.player_queues.queue_data_or_none(queue_id)
1441 if queue_data is None or (session_id := queue_data.session_id) is None:
1442 return None
1443 return f"{self.base_url}/command/{session_id}/{queue_id}/{command}.mp3"
1444
1445 def get_announcement_url(
1446 self,
1447 player_id: str,
1448 content_type: ContentType = ContentType.MP3,
1449 ) -> str:
1450 """
1451 Get the url that serves the announcement registered for the given player.
1452
1453 :param player_id: The player the announcement is played on.
1454 :param content_type: The format to serve the announcement in.
1455 """
1456 # use stream server to host announcement on local network
1457 # this ensures playback on all players, including ones that do not
1458 # like https hosts and it also offers the pre-announce 'bell'
1459 return f"{self.base_url}/announcement/{player_id}.{content_type.value}"
1460
1461 def get_stream(
1462 self,
1463 media: PlayerMedia,
1464 pcm_format: AudioFormat,
1465 player_id: str | None = None,
1466 force_flow_mode: bool = False,
1467 ) -> AsyncGenerator[bytes]:
1468 """
1469 Get a stream of the given media as raw PCM audio.
1470
1471 This is used as helper for player providers that can consume the raw PCM
1472 audio stream directly (e.g. AirPlay) and not rely on HTTP transport.
1473
1474 :param media: The PlayerMedia to stream.
1475 :param pcm_format: The desired output PCM format.
1476 :param player_id: The player ID requesting the stream. Used to determine
1477 if flow mode should be used based on the player's capabilities.
1478 :param force_flow_mode: Force flow mode regardless of player capabilities.
1479 Used for multi-client streaming scenarios that require continuous streams.
1480 """
1481 # select audio source
1482 if media.media_type == MediaType.ANNOUNCEMENT:
1483 # special case: stream announcement
1484 assert media.custom_data
1485 return self.get_announcement_stream(cast("AnnounceData", media.custom_data), pcm_format)
1486 if (
1487 media.source_id
1488 and media.source_id.startswith(UGP_PREFIX)
1489 and media.uri
1490 and "/ugp/" in media.uri
1491 ):
1492 # special case: member player accessing UGP stream
1493 # Check URI to distinguish from the UGP accessing its own stream
1494 ugp_player = cast("UniversalGroupPlayer", self.mass.players.get_player(media.source_id))
1495 ugp_stream = ugp_player.stream
1496 assert ugp_stream is not None # for type checker
1497 if ugp_stream.base_pcm_format == pcm_format:
1498 # no conversion needed
1499 return ugp_stream.subscribe_raw()
1500 return ugp_stream.get_stream(output_format=pcm_format)
1501 if (
1502 media.media_type == MediaType.AUDIO_SOURCE
1503 and not media.queue_item_id
1504 and media.source_id
1505 and (session := self.mass.players.get_audio_source_session(media.source_id))
1506 ):
1507 # a live source playing on a player rather than an item in a queue
1508 if media.queue_session_id != session.playback_session_id:
1509 # a stale request from a superseded session must not attach to the one
1510 # playing now; the http route rejects the same mismatch with a 404
1511 raise AudioError(
1512 f"Unknown (or invalid) audio source session: {media.queue_session_id}"
1513 )
1514 return self._get_audio_source_session_stream(
1515 session, pcm_format, player_id or media.source_id
1516 )
1517 if media.source_id and media.queue_item_id:
1518 # Queue stream request - determine flow_mode based on player capabilities
1519 # or force it if explicitly requested (e.g., for multi-client streaming)
1520 protocol_player = self.mass.players.get_player(player_id) if player_id else None
1521 queue_id = media.source_id
1522 queue = self.mass.player_queues.get(queue_id)
1523 queue_session_id = media.queue_session_id
1524 crossfade_needs_flow_mode = (
1525 # crossfade only applies to tracks; if the queue has it enabled but the
1526 # player(protocol) does not support gapless playback, we need to enforce flow mode
1527 media.media_type == MediaType.TRACK
1528 and queue is not None
1529 and queue.crossfade_enabled
1530 and protocol_player
1531 and not protocol_player.supports_gapless
1532 )
1533 # the audio overlay is mixed into the queue's continuous (flow) stream;
1534 # per-item requests would restart the overlay at every track boundary
1535 overlay_needs_flow_mode = queue is not None and overlay_active(queue)
1536 flow_mode = (
1537 force_flow_mode
1538 or (protocol_player is not None and protocol_player.flow_mode)
1539 or crossfade_needs_flow_mode
1540 or overlay_needs_flow_mode
1541 )
1542 if media.media_type in (MediaType.RADIO, MediaType.AUDIO_SOURCE):
1543 # flow_mode for live/infinite streams is pointless
1544 flow_mode = False
1545 if flow_mode:
1546 # flow stream request
1547 assert queue
1548 start_queue_item = self.mass.player_queues.get_item(
1549 media.source_id, media.queue_item_id
1550 )
1551 assert start_queue_item
1552 self._update_audio_processing_context(
1553 queue=queue,
1554 queue_item=start_queue_item,
1555 pcm_format=pcm_format,
1556 overlay_enabled=overlay_active(queue),
1557 session_id=queue_session_id,
1558 )
1559 flow_stream = self.audio.get_queue_flow_stream(
1560 queue=queue,
1561 start_queue_item=start_queue_item,
1562 pcm_format=pcm_format,
1563 session_id=queue_session_id,
1564 protocol_player=protocol_player,
1565 )
1566 if overlay_active(queue):
1567 flow_stream = self.audio.get_overlay_mixed_stream(
1568 queue, flow_stream, pcm_format
1569 )
1570 return flow_stream
1571 # single item stream (e.g. radio or non-flow mode)
1572 queue_item = self.mass.player_queues.get_item(media.source_id, media.queue_item_id)
1573 assert queue_item
1574 if queue is not None:
1575 self._update_audio_processing_context(
1576 queue=queue,
1577 queue_item=queue_item,
1578 pcm_format=pcm_format,
1579 overlay_enabled=(
1580 queue_item.media_type == MediaType.RADIO and overlay_active(queue)
1581 ),
1582 session_id=queue_session_id,
1583 source_crossfade_mode=self.get_source_crossfade_mode(queue, queue_item),
1584 )
1585 inner_stream = self.audio.get_queue_item_stream(
1586 queue_item=queue_item,
1587 pcm_format=pcm_format,
1588 seek_position=(
1589 int(queue_item.streamdetails.seek_position) if queue_item.streamdetails else 0
1590 ),
1591 playback_speed=cast(
1592 "float", queue_item.extra_attributes.get("playback_speed", 1.0)
1593 ),
1594 session_id=queue_session_id,
1595 )
1596 if (
1597 queue is not None
1598 and queue_item.media_type == MediaType.RADIO
1599 and overlay_active(queue)
1600 ):
1601 # radio plays as a single long-lived stream, so mix the overlay in here
1602 inner_stream = self.audio.get_overlay_mixed_stream(queue, inner_stream, pcm_format)
1603 # mirror the on_source_selected/unselected lifecycle the HTTP route
1604 # fires, so direct-PCM consumers (AirPlay, Snapcast, UGP) honour the
1605 # plugin contract too
1606 if (
1607 queue_item.media_item is not None
1608 and queue_item.media_item.media_type == MediaType.AUDIO_SOURCE
1609 ):
1610 return self._wrap_with_audio_source_lifecycle(
1611 inner=inner_stream,
1612 queue_item=queue_item,
1613 player_id=player_id or media.source_id,
1614 )
1615 return inner_stream
1616 # assume url or some other direct path
1617 # NOTE: this will fail if its an uri not playable by ffmpeg
1618 return get_ffmpeg_stream(
1619 audio_input=media.uri,
1620 input_format=AudioFormat(content_type=ContentType.try_parse(media.uri)),
1621 output_format=pcm_format,
1622 )
1623
1624 async def get_preview_stream(
1625 self,
1626 provider_instance_id_or_domain: str,
1627 item_id: str,
1628 media_type: MediaType = MediaType.TRACK,
1629 ) -> AsyncGenerator[bytes]:
1630 """Create a 30 seconds preview audioclip for the given media item."""
1631 if not (music_prov := self.mass.get_provider(provider_instance_id_or_domain)):
1632 raise ProviderUnavailableError
1633 if music_prov.type != ProviderType.MUSIC:
1634 msg = f"{provider_instance_id_or_domain} is not a music provider"
1635 raise InvalidDataError(msg)
1636 music_prov = cast("MusicProvider", music_prov)
1637
1638 try:
1639 await self.mass.music.get_item(
1640 media_type,
1641 item_id,
1642 provider_instance_id_or_domain,
1643 allow_update_metadata=False,
1644 )
1645 except MediaNotFoundError as err:
1646 msg = f"Item {item_id} not found in provider {provider_instance_id_or_domain}"
1647 raise InvalidDataError(msg) from err
1648
1649 streamdetails = await music_prov.get_stream_details(item_id, media_type)
1650 pcm_format = AudioFormat(
1651 content_type=ContentType.from_bit_depth(streamdetails.audio_format.bit_depth),
1652 sample_rate=streamdetails.audio_format.sample_rate,
1653 bit_depth=streamdetails.audio_format.bit_depth,
1654 channels=streamdetails.audio_format.channels,
1655 )
1656 async for chunk in get_ffmpeg_stream(
1657 audio_input=self.audio.get_media_stream(
1658 streamdetails=streamdetails, pcm_format=pcm_format
1659 ),
1660 input_format=pcm_format,
1661 output_format=AudioFormat(content_type=ContentType.AAC),
1662 extra_input_args=["-t", "30"],
1663 ):
1664 yield chunk
1665
1666 async def get_announcement_stream(
1667 self, announce_data: AnnounceData, output_format: AudioFormat
1668 ) -> AsyncGenerator[bytes]:
1669 """
1670 Get the audio of an announcement (pre-announce chime + announcement).
1671
1672 Any number of consumers may stream the same announcement at once; its source is
1673 fetched and decoded only once. The audio stays available while the stream is
1674 held open.
1675
1676 :param announce_data: The announcement to stream.
1677 :param output_format: The format to deliver the audio in.
1678 """
1679 render = self.announcement_renderer.acquire(announce_data)
1680 try:
1681 # aclosing guarantees this consumer's ffmpeg encoder is torn down
1682 # immediately when it goes away, instead of lingering until garbage
1683 # collection finalizes the abandoned generator.
1684 stream = render.get_stream(output_format)
1685 async with aclosing(stream):
1686 async for chunk in stream:
1687 yield chunk
1688 finally:
1689 await self.announcement_renderer.release(render)
1690
1691 async def get_announcement_duration(
1692 self, announcement: PlayerMedia, timeout: float = DEFAULT_RENDER_TIMEOUT
1693 ) -> int | None:
1694 """
1695 Get the exact duration (in seconds) of an announcement, once it finished rendering.
1696
1697 Waits for the audio to be rendered in full, so call this while the announcement
1698 plays rather than before handing it to a player. Returns None when the length can
1699 not be determined, e.g. the announcement is no longer playing or its source did
1700 not deliver in time.
1701
1702 :param announcement: The announcement to return the duration for.
1703 :param timeout: Maximum time to wait for the audio to finish rendering.
1704 """
1705 if announcement.duration:
1706 return announcement.duration
1707 if not announcement.custom_data:
1708 return None
1709 render = self.announcement_renderer.get(cast("AnnounceData", announcement.custom_data))
1710 if render is None:
1711 return None
1712 duration = await render.wait_finished(timeout)
1713 return ceil(duration) if duration else None
1714
1715 def _resolve_audio_source_request(
1716 self, request: web.Request
1717 ) -> tuple[AudioSourceSession, Player, PluginProvider]:
1718 """
1719 Resolve a source stream request to its session, consuming player and plugin.
1720
1721 :param request: The stream request to resolve.
1722 :raises web.HTTPNotFound: When any of the three is gone, or the url names a
1723 session that is no longer the one playing.
1724 """
1725 source_player_id = request.match_info["source_player_id"]
1726 player_id = request.match_info["player_id"]
1727 session_id = request.match_info["session_id"]
1728 session = self.mass.players.get_audio_source_session(source_player_id)
1729 if session is None:
1730 raise web.HTTPNotFound(reason=f"No audio source playing on {source_player_id}")
1731 if session_id != session.playback_session_id:
1732 raise web.HTTPNotFound(reason=f"Unknown (or invalid) session: {session_id}")
1733 if not (player := self.mass.players.get_player(player_id)):
1734 raise web.HTTPNotFound(reason=f"Unknown Player: {player_id}")
1735 prov = self.mass.get_provider(session.provider_instance_id)
1736 if not isinstance(prov, PluginProvider):
1737 raise web.HTTPNotFound(
1738 reason=f"AudioSource provider {session.provider_instance_id} unavailable"
1739 )
1740 return session, player, prov
1741
1742 async def _release_unstarted_audio_source(
1743 self, session: AudioSourceSession, playback_session_id: str
1744 ) -> None:
1745 """
1746 Take a source that never started off the player holding it.
1747
1748 The command that pointed the renderer here has already returned, so nothing
1749 else will clear the session: without this the player goes on publishing a
1750 source that never played, with its own queue held inactive behind it.
1751
1752 :param session: The session whose stream failed before any audio flowed.
1753 :param playback_session_id: Playback session active when stream setup started.
1754 """
1755 current_session = self.mass.players.get_audio_source_session(session.player_id)
1756 if (
1757 current_session is not session
1758 or current_session.playback_session_id != playback_session_id
1759 ):
1760 # already superseded, so it is not ours to release
1761 return
1762 self.logger.debug(
1763 "AudioSource %s never started on player %s, releasing it",
1764 session.source_id,
1765 session.player_id,
1766 )
1767 try:
1768 await self.mass.players.deselect_source(
1769 session.player_id,
1770 provider_instance_id=session.provider_instance_id,
1771 source_id=session.source_id,
1772 playback_session_id=playback_session_id,
1773 )
1774 except Exception:
1775 # deselect_source already absorbs the expected stop failures, so anything
1776 # arriving here is a defect worth a trail rather than a silent half-cleanup
1777 self.logger.warning(
1778 "Failed to release AudioSource %s on player %s",
1779 session.source_id,
1780 session.player_id,
1781 exc_info=True,
1782 )
1783
1784 async def _serve_audio_source_head(
1785 self, request: web.Request, session: AudioSourceSession
1786 ) -> web.StreamResponse:
1787 """
1788 Answer a HEAD probe for a live audio source without touching the plugin.
1789
1790 :param request: The probe to answer.
1791 :param session: The session whose source is being probed.
1792 """
1793 head_fmt = request.match_info["fmt"]
1794 if ContentType.try_parse(head_fmt).is_pcm():
1795 # most DLNA renderers pick a decoder from the HEAD content type and
1796 # cannot handle raw PCM
1797 head_fmt = ContentType.WAV.value
1798 resp = web.StreamResponse(
1799 status=200, reason="OK", headers=_audio_source_headers(session, head_fmt)
1800 )
1801 await resp.prepare(request)
1802 return resp
1803
1804 async def _prepare_audio_source_stream(
1805 self,
1806 request: web.Request,
1807 player: Player,
1808 session: AudioSourceSession,
1809 streamdetails: StreamDetails,
1810 ) -> tuple[web.StreamResponse, AsyncGenerator[bytes]]:
1811 """
1812 Open the response for a live audio source and build the audio behind it.
1813
1814 :param request: The stream request being answered.
1815 :param player: The player consuming this stream.
1816 :param session: The session whose source is being streamed.
1817 :param streamdetails: The stream details resolved for that source.
1818 :return: The prepared response and the encoded audio to write to it.
1819 """
1820 pcm_format = await self.audio.select_pcm_format(
1821 player=player,
1822 streamdetails=streamdetails,
1823 crossfade_enabled=False,
1824 overlay_active=False,
1825 )
1826 output_format = await self.audio.get_output_format(
1827 output_format_str=request.match_info["fmt"],
1828 player=player,
1829 content_sample_rate=pcm_format.sample_rate,
1830 content_bit_depth=pcm_format.bit_depth,
1831 media_type=MediaType.AUDIO_SOURCE,
1832 )
1833 resp = web.StreamResponse(
1834 status=200,
1835 reason="OK",
1836 headers=_audio_source_headers(session, output_format.output_format_str),
1837 )
1838 resp.content_type = get_mime_type(output_format.output_format_str)
1839 http_profile = player.get_config_value(CONF_HTTP_PROFILE, "default")
1840 if http_profile == "forced_content_length":
1841 # a live source has no length, so advertise one it will never reach
1842 resp.content_length = calculate_content_length(output_format, 12 * 3600)
1843 elif http_profile == "chunked":
1844 resp.enable_chunked_encoding()
1845 await resp.prepare(request)
1846
1847 audio_input = self.audio.get_audio_source_stream(
1848 streamdetails=streamdetails,
1849 pcm_format=pcm_format,
1850 raise_on_error=False,
1851 display_name=session.source.name,
1852 )
1853 filter_params = self.audio.get_player_output_plan(
1854 player_id=player.player_id,
1855 input_format=pcm_format,
1856 output_format=output_format,
1857 shared_player_ids=player.state.group_members,
1858 ).filter_params
1859 if (
1860 output_format.content_type == ContentType.WAV
1861 and not filter_params
1862 and output_format.sample_rate == pcm_format.sample_rate
1863 and output_format.bit_depth == pcm_format.bit_depth
1864 and output_format.channels == pcm_format.channels
1865 ):
1866 # the player takes the source's exact PCM, so skip the encode ffmpeg and
1867 # its buffer latency and send a WAV header with the raw bytes
1868 return resp, _wav_passthrough_stream(audio_input, output_format)
1869 return resp, get_ffmpeg_stream(
1870 audio_input=audio_input,
1871 input_format=pcm_format,
1872 output_format=output_format,
1873 filter_params=filter_params,
1874 # keep the encode stage from reading further ahead than it needs to: a live
1875 # source's latency is whatever is buffered between it and the player
1876 extra_input_args=[
1877 "-readrate",
1878 SINGLE_ITEM_READRATE,
1879 "-readrate_initial_burst",
1880 SINGLE_ITEM_READRATE_INITIAL_BURST,
1881 ],
1882 )
1883
1884 async def _get_audio_source_session_stream(
1885 self,
1886 session: AudioSourceSession,
1887 pcm_format: AudioFormat,
1888 consumer_player_id: str,
1889 ) -> AsyncGenerator[bytes]:
1890 """
1891 Stream a live source to a consumer that takes raw PCM rather than the http url.
1892
1893 AirPlay, Snapcast, squeezelite's multi-client path, universal groups and the
1894 MSX bridge all consume PCM directly, so they never reach the http route and
1895 need the plugin lifecycle fired here instead — those hooks are what claim and
1896 release the source and kick acquisition side effects into life.
1897
1898 :param session: The live source session playing on its owner.
1899 :param pcm_format: The PCM format the consumer wants.
1900 :param consumer_player_id: The player consuming this stream, which is not
1901 necessarily the one that owns the source.
1902 """
1903 prov = self.mass.get_provider(session.provider_instance_id)
1904 if not isinstance(prov, PluginProvider):
1905 raise AudioError(
1906 f"AudioSource provider {session.provider_instance_id} is not available"
1907 )
1908 playback_session_id = session.playback_session_id
1909 stream_session_id = uuid4().hex
1910 serving = False
1911 try:
1912 try:
1913 await prov.on_source_selected(
1914 session.source_id,
1915 consumer_player_id,
1916 session.player_id,
1917 stream_session_id,
1918 )
1919 except RuntimeError as err:
1920 # the plugin refuses this consumer, e.g. it just redirected playback
1921 raise AudioError(str(err)) from err
1922 if (
1923 self.mass.players.get_audio_source_session(session.player_id) is not session
1924 or session.playback_session_id != playback_session_id
1925 ):
1926 raise AudioError("AudioSource session was superseded")
1927 session.stream_session_id = stream_session_id
1928 if (streamdetails := session.streamdetails) is None:
1929 streamdetails = await prov.get_stream_details(
1930 session.source_id, MediaType.AUDIO_SOURCE
1931 )
1932 session.attach_streamdetails(streamdetails)
1933 serving = True
1934 async for chunk in self.audio.get_audio_source_stream(
1935 streamdetails=streamdetails,
1936 pcm_format=pcm_format,
1937 raise_on_error=False,
1938 display_name=session.source.name,
1939 ):
1940 if (
1941 self.mass.players.get_audio_source_session(session.player_id) is not session
1942 or session.playback_session_id != playback_session_id
1943 or session.stream_session_id != stream_session_id
1944 ):
1945 break
1946 yield chunk
1947 finally:
1948 try:
1949 await prov.on_source_unselected(
1950 session.source_id, session.player_id, stream_session_id
1951 )
1952 except Exception:
1953 self.logger.warning(
1954 "on_source_unselected raised for provider %s source %s player %s",
1955 prov.instance_id,
1956 session.source_id,
1957 session.player_id,
1958 exc_info=True,
1959 )
1960 if not serving:
1961 await self._release_unstarted_audio_source(session, playback_session_id)
1962
1963 async def _wrap_with_audio_source_lifecycle(
1964 self,
1965 inner: AsyncGenerator[bytes],
1966 queue_item: QueueItem,
1967 player_id: str,
1968 ) -> AsyncGenerator[bytes]:
1969 """
1970 Wrap an AudioSource queue item stream with on_source_selected/unselected hooks.
1971
1972 Direct-PCM consumers (AirPlay, Snapcast, UGP, ...) call ``get_stream`` instead
1973 of going through the HTTP route, but the plugin contract requires the
1974 lifecycle hooks to fire for every actual stream request — they're what
1975 claim/release the per-queue exclusive ownership and trigger acquisition
1976 side effects like the Spotify Connect Web API play kick. This wrapper
1977 gives those consumers the same lifecycle the HTTP route already provides.
1978
1979 :param inner: The underlying audio stream generator.
1980 :param queue_item: The AudioSource queue item being streamed.
1981 :param player_id: The protocol player consuming this stream.
1982 """
1983 media_item = queue_item.media_item
1984 assert media_item is not None # caller checked media_type == AUDIO_SOURCE
1985 prov = self.mass.get_provider(media_item.provider)
1986 queue_id = queue_item.queue_id
1987 if not isinstance(prov, PluginProvider):
1988 async for chunk in inner:
1989 yield chunk
1990 return
1991 source_id = media_item.item_id
1992 stream_session_id = uuid4().hex
1993 # single try/finally so on_source_unselected fires even when
1994 # on_source_selected raises after partially claiming state; the
1995 # provider's session_id guard makes a no-op claim release safe.
1996 try:
1997 try:
1998 await prov.on_source_selected(source_id, player_id, queue_id, stream_session_id)
1999 except RuntimeError as err:
2000 # provider intentionally aborts the request — surface as AudioError
2001 raise AudioError(str(err)) from err
2002 async for chunk in inner:
2003 yield chunk
2004 finally:
2005 try:
2006 await prov.on_source_unselected(source_id, queue_id, stream_session_id)
2007 except Exception:
2008 self.logger.exception(
2009 "on_source_unselected raised for provider %s source %s queue %s",
2010 prov.instance_id,
2011 source_id,
2012 queue_id,
2013 )
2014
2015 def _source_fades_an_adjacent_item(self, queue: PlayerQueue, queue_item: QueueItem) -> bool:
2016 """
2017 Return whether the same source serves an item next to this one.
2018
2019 Stands in for a source that is not serving this queue yet and so cannot answer
2020 for the item itself: a source can only fade across a boundary it owns both
2021 sides of, which makes this a necessary condition rather than proof of a fade.
2022 Both neighbours count, the same way one of our own fades credits both of its
2023 sides - the last track of a queue is still the side that was faded into.
2024
2025 :param queue: Queue the item is played from.
2026 :param queue_item: Queue item whose neighbours to check.
2027 """
2028 assert queue_item.streamdetails is not None # guaranteed by the caller
2029 controller = self.mass.player_queues
2030 queue_id = queue.queue_id
2031 neighbours = [controller.get_next_item(queue_id, queue_item.queue_item_id)]
2032 index = controller.index_by_id(queue_id, queue_item.queue_item_id)
2033 if index is not None and index > 0:
2034 neighbours.append(controller.get_item(queue_id, index - 1))
2035 return any(
2036 self._served_by(neighbour, queue_item.streamdetails.provider)
2037 for neighbour in neighbours
2038 )
2039
2040 def _served_by(self, queue_item: QueueItem | None, provider_instance: str) -> bool:
2041 """
2042 Return whether a queue item is a track the given provider instance serves.
2043
2044 :param queue_item: Queue item to check, or None when there is none.
2045 :param provider_instance: Instance id of the provider to match.
2046 """
2047 if queue_item is None or queue_item.media_type != MediaType.TRACK:
2048 return False
2049 if (streamdetails := queue_item.streamdetails) is not None:
2050 # already resolved, so this is the provider that will really serve it
2051 return streamdetails.provider == provider_instance
2052 if (media_item := queue_item.media_item) is None:
2053 return False
2054 return media_item.provider == provider_instance or any(
2055 mapping.provider_instance == provider_instance
2056 for mapping in media_item.provider_mappings
2057 )
2058
2059 def _update_audio_processing_context(
2060 self,
2061 queue: PlayerQueue,
2062 queue_item: QueueItem,
2063 pcm_format: AudioFormat,
2064 overlay_enabled: bool,
2065 session_id: str | None = None,
2066 source_crossfade_mode: CrossfadeMode = CrossfadeMode.DISABLED,
2067 ) -> None:
2068 """
2069 Store the shared processing context selected for a queue item.
2070
2071 Our own crossfade is left out on purpose: only the audio layer knows whether
2072 one really happens, and it reports that itself once the boundary has decided.
2073 A crossfade the source performs is the exception - the audio layer never sees
2074 that one, so it is carried here.
2075
2076 :param queue: Active player queue.
2077 :param queue_item: Queue item being prepared.
2078 :param pcm_format: Shared PCM format leaving queue processing.
2079 :param overlay_enabled: Whether an overlay is mixed into this stream.
2080 :param session_id: Queue session that owns processing-detail updates.
2081 :param source_crossfade_mode: Crossfade the item's own source applies, if any.
2082 """
2083 if queue_item.streamdetails is None:
2084 return
2085 queue_data = self.mass.player_queues.queue_data_or_none(queue.queue_id)
2086 if (
2087 queue_data is None
2088 or (processing_session_id := session_id or queue_data.session_id) is None
2089 or queue_data.session_id != processing_session_id
2090 ):
2091 return
2092 self.audio_processing.start_session(queue.queue_id, processing_session_id)
2093 self.audio_processing.update_item_context(
2094 queue_id=queue.queue_id,
2095 session_id=processing_session_id,
2096 queue_item_id=queue_item.queue_item_id,
2097 queue_processing=AudioQueueProcessing(
2098 pcm_format=pcm_format,
2099 playback_speed=cast(
2100 "float",
2101 queue_item.extra_attributes.get("playback_speed", 1.0),
2102 ),
2103 crossfade_mode=source_crossfade_mode,
2104 overlay_active=overlay_enabled,
2105 ),
2106 alters_audio=queue_item.streamdetails.fade_in,
2107 )
2108
2109 def _get_announcement_http_profile(self, player_id: str, announce_data: AnnounceData) -> str:
2110 """
2111 Resolve the http profile for serving an announcement stream.
2112
2113 Announcement urls are registered under the visible player's id, but the
2114 stream may be fetched by a linked protocol player; the profile must come
2115 from the player that actually performs the fetch.
2116 """
2117 announce_player = None
2118 if announce_player_id := announce_data.get("announce_player_id"):
2119 announce_player = self.mass.players.get_player(announce_player_id)
2120 if announce_player is None:
2121 announce_player = self.mass.players.get_player(player_id)
2122 if announce_player is None:
2123 return "default"
2124 return announce_player.get_output_config_value(CONF_HTTP_PROFILE, "default")
2125
2126 async def _finish_flow_stream(
2127 self, resp: web.StreamResponse, queue_id: str, session_id: str
2128 ) -> None:
2129 """
2130 Close a fully served flow stream, giving the player time to drain when it ends the queue.
2131
2132 :param resp: The flow stream response, already fully written.
2133 :param queue_id: Id of the queue the flow stream belongs to.
2134 :param session_id: Stream session this response was opened for.
2135 """
2136 if self.mass.player_queues.flow_queue_exhausted(queue_id, session_id):
2137 # the player is still holding a few seconds of audio it has not rendered yet
2138 # and drops that as soon as the stream ends, so let it play out first.
2139 # a flow that ends to be restarted right away gets no such grace: there the
2140 # player should go idle as soon as possible so the next stream can start.
2141 self.logger.debug(
2142 "Flow stream for queue %s reached the end of the queue - holding the "
2143 "connection open for %ss so the player can play out its buffer",
2144 queue_id,
2145 FLOW_STREAM_LEAD_OUT_SECONDS,
2146 )
2147 await asyncio.sleep(FLOW_STREAM_LEAD_OUT_SECONDS)
2148 # aiohttp derives keep-alive from the request, so the 'Connection: close' we
2149 # advertise is relayed to the player but never applied to the response itself.
2150 # Without this the player is left waiting on a stream that already ended.
2151 resp.force_close()
2152
2153 def _log_request(self, request: web.Request) -> None:
2154 """Log request."""
2155 if self.logger.isEnabledFor(VERBOSE_LOG_LEVEL):
2156 self.logger.log(
2157 VERBOSE_LOG_LEVEL,
2158 "Got %s request to %s from %s\nheaders: %s\n",
2159 request.method,
2160 request.path,
2161 request.remote,
2162 redact_sensitive_headers(request.headers),
2163 )
2164 else:
2165 self.logger.debug(
2166 "Got %s request to %s from %s (HTTP/%s.%s, connection: %s)",
2167 request.method,
2168 request.path,
2169 request.remote,
2170 request.version.major,
2171 request.version.minor,
2172 request.headers.get("Connection", "-"),
2173 )
2174
2175 async def _reload_network_dependent_providers(self) -> None:
2176 """Reload the providers that captured the streamserver network, if it changed."""
2177 previous = self._network_fingerprint
2178 current = (
2179 self._bind_ip,
2180 str(self.publish_ip),
2181 cast("int", self.publish_port),
2182 tuple(self._publish_addresses),
2183 )
2184 if previous is None or previous == current:
2185 self._network_fingerprint = current
2186 return
2187 # these providers bind or advertise the network while they load, so a plain
2188 # reload is what moves them over - they share no lighter rebind path
2189 instance_ids = [
2190 prov.instance_id
2191 for prov in self.mass.providers
2192 if prov.reload_on_streams_network_change
2193 ]
2194 for instance_id in instance_ids:
2195 try:
2196 config = await self.mass.config.get_provider_config(instance_id)
2197 self.logger.info(
2198 "Streamserver network changed, reloading provider %s",
2199 config.name or config.domain,
2200 )
2201 await self.mass.load_provider_config(config)
2202 except Exception as err:
2203 self.logger.warning(
2204 "Error reloading provider %s: %s",
2205 instance_id,
2206 str(err) or err.__class__.__name__,
2207 exc_info=err,
2208 )
2209 # only mark the new network as applied once the loop completed, so a run cut short
2210 # by a second config change runs again on the next reload
2211 self._network_fingerprint = current
2212
2213 def _setup_smart_fades_logger(self, config: CoreConfig) -> None:
2214 """Set up smart fades logger level."""
2215 log_level = str(config.get_value(CONF_SMART_FADES_LOG_LEVEL))
2216 if log_level == "GLOBAL":
2217 self.audio.smart_fades_mixer.logger.setLevel(self.logger.level)
2218 else:
2219 self.audio.smart_fades_mixer.logger.setLevel(log_level)
2220
2221 def _resolve_publish_state(self, bind_ip: str, publish_candidates: tuple[str, ...]) -> None:
2222 """
2223 Resolve the addresses and base URL to advertise for the given bind address.
2224
2225 Reads ``self.publish_port``, so set that first.
2226
2227 :param bind_ip: Address the streamserver binds to (a wildcard means all interfaces).
2228 :param publish_candidates: Host addresses reachable from the local network, ranked.
2229 """
2230 self._bind_ip = bind_ip
2231 self._publish_addresses = _get_publish_addresses(
2232 bind_ip, self._configured_publish_ip, publish_candidates
2233 )
2234 # the single address players are handed, taken from the top of the ranked list
2235 self.publish_ip = self._publish_addresses[0]
2236 self._base_url = f"http://{format_ip_for_url(self.publish_ip)}:{self.publish_port}"
2237
2238
2239def _same_ip_family(ip: str, other_ip: str) -> bool:
2240 """Return whether two addresses belong to the same IP family."""
2241 return (":" in ip) == (":" in other_ip)
2242