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