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