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