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