/
/
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
13import time
14from collections.abc import AsyncGenerator
15from contextlib import aclosing
16from math import ceil
17from typing import TYPE_CHECKING, cast
18from uuid import uuid4
19
20from aiofiles.os import wrap
21from aiohttp import web
22from music_assistant_models.audio_processing import AudioQueueProcessing
23from music_assistant_models.config_entries import ConfigEntry, ConfigValueOption
24from music_assistant_models.enums import (
25 ConfigEntryType,
26 ContentType,
27 CrossfadeMode,
28 MediaType,
29 PlayerFeature,
30 ProviderType,
31 VolumeNormalizationMode,
32)
33from music_assistant_models.errors import (
34 AudioError,
35 InvalidDataError,
36 MediaNotFoundError,
37 ProviderUnavailableError,
38)
39from music_assistant_models.helpers import get_global_cache_value
40from music_assistant_models.media_items import AudioFormat
41
42from music_assistant.constants import (
43 CONF_BACKGROUND_SCAN_CONCURRENCY,
44 CONF_BIND_IP,
45 CONF_BIND_PORT,
46 CONF_CROSSFADE_DURATION,
47 CONF_CROSSFADE_MODE,
48 CONF_ENTRY_ENABLE_ICY_METADATA,
49 CONF_ENTRY_LOG_LEVEL,
50 CONF_ENTRY_VOLUME_NORMALIZATION_TARGET,
51 CONF_HTTP_PROFILE,
52 CONF_OUTPUT_CODEC,
53 CONF_PLAYER_QUEUES,
54 CONF_PREFER_WAV_FOR_LIVE_SOURCES,
55 CONF_PUBLISH_IP,
56 CONF_VALUE_AUTO,
57 CONF_VOLUME_NORMALIZATION_FIXED_GAIN_RADIO,
58 CONF_VOLUME_NORMALIZATION_FIXED_GAIN_TRACKS,
59 CONF_VOLUME_NORMALIZATION_RADIO,
60 CONF_VOLUME_NORMALIZATION_TRACKS,
61 DEFAULT_BACKGROUND_SCAN_CONCURRENCY,
62 DEFAULT_HOST,
63 DEFAULT_STREAM_HEADERS,
64 DLNA_CONTENT_FEATURES,
65 DLNA_CONTENT_FEATURES_REALTIME,
66 ICY_HEADERS,
67 SILENCE_FILE,
68 VERBOSE_LOG_LEVEL,
69 WILDCARD_BIND_IPS,
70)
71from music_assistant.controllers.players.helpers import AnnounceData
72from music_assistant.controllers.streams.announcements import (
73 DEFAULT_RENDER_TIMEOUT,
74 AnnouncementRenderer,
75)
76from music_assistant.controllers.streams.audio import StreamsAudio, overlay_active
77from music_assistant.controllers.streams.audio_analysis import AudioAnalysisController
78from music_assistant.controllers.streams.audio_processing import (
79 AudioProcessingManager,
80)
81from music_assistant.controllers.streams.constants import (
82 CONF_ALLOW_CROSSFADE_SAME_ALBUM,
83 CONF_BUFFER_SIZE,
84 CONF_BUFFER_SIZE_DEFAULT,
85 CONF_SMART_FADES_LOG_LEVEL,
86 DEFAULT_PORT,
87 FLOW_STREAM_LEAD_OUT_SECONDS,
88 BufferSize,
89 get_available_buffer_sizes,
90)
91from music_assistant.controllers.streams.live_announcements import (
92 LIVE_ANNOUNCEMENT_STREAM_PATH,
93 LiveAnnouncementManager,
94)
95from music_assistant.helpers.audio import (
96 calculate_content_length,
97 create_streaming_wave_header,
98 get_content_length,
99 get_mime_type,
100 store_content_length_in_cache,
101)
102from music_assistant.helpers.buffered_generator import buffered
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 StreamMetadata
130
131 from music_assistant.helpers.json import SerializableType
132 from music_assistant.mass import MusicAssistant
133
134
135isfile = wrap(os.path.isfile)
136
137
138async def _wav_passthrough_stream(
139 audio_input: AsyncGenerator[bytes], output_format: AudioFormat
140) -> AsyncGenerator[bytes]:
141 """Yield a WAV header followed by raw PCM bytes from ``audio_input``."""
142 yield create_streaming_wave_header(output_format)
143 async for chunk in audio_input:
144 yield chunk
145
146
147def _get_publish_addresses(
148 bind_ip: str, configured_publish_ip: str | None, publish_candidates: tuple[str, ...]
149) -> list[str]:
150 """
151 Return the addresses this host publishes on, best candidate first.
152
153 :param bind_ip: The configured bind IP (a wildcard means all interfaces).
154 :param configured_publish_ip: The explicitly configured publish IP, or None when auto.
155 :param publish_candidates: Host addresses reachable from the local network, ranked.
156 """
157 if configured_publish_ip:
158 # an explicitly configured address is the authoritative answer
159 return [configured_publish_ip]
160 if bind_ip and bind_ip not in WILDCARD_BIND_IPS:
161 # only one interface is served, so no other address can be reached
162 return [bind_ip]
163 # auto-detected: keep the whole ranked list - publish_ip takes the best of them and
164 # the network fingerprint watches all of them to spot an interface change
165 return list(publish_candidates)
166
167
168class StreamsController(CoreController):
169 """Controller to stream audio to players."""
170
171 domain: str = "streams"
172
173 def __init__(self, mass: MusicAssistant) -> None:
174 """Initialize instance."""
175 super().__init__(mass)
176 self._server = Webserver(self.logger, enable_dynamic_routes=True)
177 self.register_dynamic_route = self._server.register_dynamic_route
178 self.unregister_dynamic_route = self._server.unregister_dynamic_route
179 self.manifest.name = "Streamserver"
180 self.manifest.description = (
181 "Music Assistant's core controller that is responsible for "
182 "streaming audio to players on the local network."
183 )
184 self.manifest.icon = "cast-audio"
185 self.announcement_renderer = AnnouncementRenderer()
186 self.live_announcements = LiveAnnouncementManager(mass, self.logger)
187 self._bind_ip: str = "0.0.0.0"
188 self._base_url: str = ""
189 self._configured_publish_ip: str | None = None
190 # every address players may reach this host on, best candidate first; publish_ip is
191 # the first of them and the network fingerprint watches the whole list for changes
192 self._publish_addresses: list[str] = []
193 # the network as it was at the previous setup, to spot a runtime change
194 self._network_fingerprint: tuple[str, str, int, tuple[str, ...]] | None = None
195 self.audio = StreamsAudio(mass)
196 self.audio_processing = AudioProcessingManager(mass)
197 self._audio_analysis = AudioAnalysisController(self)
198 # Number of queue streams (single item or flow) actively serving a player right now.
199 # Audio analysis reads this (via audio_analysis.playback_active) to yield CPU while a
200 # queue stream is live. Announcements are a separate path that never runs analysis.
201 self._active_output_streams = 0
202
203 @property
204 def audio_analysis(self) -> AudioAnalysisController:
205 """Return the AudioAnalysisController instance."""
206 return self._audio_analysis
207
208 def output_stream_active(self) -> bool:
209 """Return whether a queue stream (single item or flow) is actively serving a player."""
210 return self._active_output_streams > 0
211
212 async def get_diagnostics(self) -> dict[str, SerializableType]:
213 """Return diagnostics info for this controller to include in diagnostics reports."""
214 return {
215 "ffmpeg_version": get_global_cache_value(CACHE_ATTR_FFMPEG_VERSION),
216 "libsoxr_support": get_global_cache_value(CACHE_ATTR_LIBSOXR_PRESENT),
217 "active_output_streams": self._active_output_streams,
218 "active_announcements": self.announcement_renderer.active_announcements,
219 "active_announcement_renders": self.announcement_renderer.active_renders,
220 "active_live_announcements": self.live_announcements.active_sessions,
221 "publish_ip_configured": self._configured_publish_ip is not None,
222 }
223
224 @property
225 def base_url(self) -> str:
226 """Return the base_url for the streamserver."""
227 return self._base_url
228
229 @property
230 def bind_ip(self) -> str:
231 """Return the IP address this streamserver is bound to."""
232 return self._bind_ip
233
234 async def get_source_ip(self, target_ip: str | None = None) -> str | None:
235 """
236 Return a local, bindable source IP on the player-facing network.
237
238 For callers that bind a socket or hand a local interface address to a helper
239 process, so their traffic leaves on the network the players live on. The result
240 is always an address of this host, never the advertised address, which may not
241 exist here at all.
242
243 Returns None when no single interface should be pinned, which the caller must
244 read as "bind all interfaces and let the routing table decide".
245
246 :param target_ip: IP address of the device the traffic is meant for. Omit it for
247 a shared consumer that serves every player at once; such a caller can only be
248 pinned by an explicitly configured bind IP.
249 """
250 if self._bind_ip and self._bind_ip not in WILDCARD_BIND_IPS:
251 if target_ip and not _same_ip_family(self._bind_ip, target_ip):
252 return None
253 return self._bind_ip
254 if not target_ip:
255 return None
256 return await get_source_ip_for_target(target_ip) or None
257
258 def get_publish_ip(self, target_ip: str) -> str | None:
259 """
260 Return the address to advertise to the device at ``target_ip``, if one is configured.
261
262 Only an explicitly configured publish IP is returned. An auto-detected one is a
263 guess at this host's primary interface, which on a multi-homed host is not
264 necessarily the network the players live on, so callers that can derive the
265 address from the connection itself must prefer that over the guess.
266
267 Returns None when no publish IP was configured, or when the configured one cannot
268 apply to this device.
269
270 :param target_ip: IP address of the device that would receive the address, used to
271 reject an address of the wrong IP family.
272 """
273 if not self._configured_publish_ip:
274 return None
275 if not _same_ip_family(self._configured_publish_ip, target_ip):
276 return None
277 return self._configured_publish_ip
278
279 @property
280 def smart_fades_available(self) -> bool:
281 """
282 Return whether smart crossfade can be used on this server.
283
284 Requires a large-enough audio buffer (at least balanced) and a loaded
285 smart fades audio analysis provider.
286 """
287 buffer_size = BufferSize(
288 self.mass.config.get_raw_core_config_value(
289 self.domain, CONF_BUFFER_SIZE, CONF_BUFFER_SIZE_DEFAULT
290 )
291 )
292 return (
293 buffer_size != BufferSize.MINIMAL and self.audio_analysis.smart_fades_provider_available
294 )
295
296 def get_crossfade_mode(self, queue: PlayerQueue) -> CrossfadeMode:
297 """
298 Return the effective crossfade mode for a queue.
299
300 Combines the per-play on/off toggle with the crossfade_mode setting and smart fades
301 availability: smart when enabled, selected and available; standard when enabled but smart
302 is not selected/available; disabled otherwise.
303 """
304 if not queue.crossfade_enabled:
305 return CrossfadeMode.DISABLED
306 # default to smart when this server can use it, else standard
307 default_mode = (
308 CrossfadeMode.SMART_CROSSFADE
309 if self.smart_fades_available
310 else CrossfadeMode.STANDARD_CROSSFADE
311 )
312 mode = self.mass.config.get_effective_player_queue_config_value(
313 queue.queue_id, CONF_CROSSFADE_MODE, default_mode
314 )
315 if mode == CrossfadeMode.SMART_CROSSFADE and self.smart_fades_available:
316 return CrossfadeMode.SMART_CROSSFADE
317 return CrossfadeMode.STANDARD_CROSSFADE
318
319 def is_smart_fades_active(self, queue: PlayerQueue) -> bool:
320 """Return whether the queue's effective crossfade mode is smart crossfade."""
321 return self.get_crossfade_mode(queue) == CrossfadeMode.SMART_CROSSFADE
322
323 async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
324 """Return all Config Entries for this core module (if any)."""
325 ip_addresses = await get_ip_addresses(include_ipv6=True)
326 return (
327 ConfigEntry(
328 key=CONF_BUFFER_SIZE,
329 type=ConfigEntryType.STRING,
330 default_value=CONF_BUFFER_SIZE_DEFAULT,
331 # Only offer presets the host's RAM can sustain (Balanced >= 4GB,
332 # Maximum >= 7GB); see get_available_buffer_sizes.
333 options=[ConfigValueOption(size.value) for size in get_available_buffer_sizes()],
334 required=False,
335 category="playback",
336 ),
337 ConfigEntry(
338 key=CONF_VOLUME_NORMALIZATION_RADIO,
339 type=ConfigEntryType.STRING,
340 default_value=VolumeNormalizationMode.FALLBACK_DYNAMIC,
341 options=[
342 ConfigValueOption(x.value, title=x.value.replace("_", " ").title())
343 for x in VolumeNormalizationMode
344 ],
345 category="playback",
346 ),
347 ConfigEntry(
348 key=CONF_VOLUME_NORMALIZATION_TRACKS,
349 type=ConfigEntryType.STRING,
350 default_value=VolumeNormalizationMode.FALLBACK_DYNAMIC,
351 options=[
352 ConfigValueOption(x.value, title=x.value.replace("_", " ").title())
353 for x in VolumeNormalizationMode
354 ],
355 category="playback",
356 ),
357 ConfigEntry(
358 key=CONF_VOLUME_NORMALIZATION_FIXED_GAIN_RADIO,
359 type=ConfigEntryType.FLOAT,
360 range=(-20, 10),
361 default_value=-6,
362 category="playback",
363 ),
364 ConfigEntry(
365 key=CONF_VOLUME_NORMALIZATION_FIXED_GAIN_TRACKS,
366 type=ConfigEntryType.FLOAT,
367 range=(-20, 10),
368 default_value=-6,
369 category="playback",
370 ),
371 CONF_ENTRY_VOLUME_NORMALIZATION_TARGET,
372 ConfigEntry(
373 key=CONF_ALLOW_CROSSFADE_SAME_ALBUM,
374 type=ConfigEntryType.BOOLEAN,
375 default_value=False,
376 category="playback",
377 ),
378 ConfigEntry(
379 key=CONF_PUBLISH_IP,
380 type=ConfigEntryType.STRING,
381 default_value=CONF_VALUE_AUTO,
382 required=False,
383 category="generic",
384 advanced=True,
385 requires_reload=True,
386 ),
387 ConfigEntry(
388 key=CONF_BIND_PORT,
389 type=ConfigEntryType.INTEGER,
390 default_value=DEFAULT_PORT,
391 category="generic",
392 advanced=True,
393 requires_reload=True,
394 ),
395 ConfigEntry(
396 key=CONF_BIND_IP,
397 type=ConfigEntryType.STRING,
398 default_value=DEFAULT_HOST,
399 options=[ConfigValueOption(x, title=x) for x in {DEFAULT_HOST, *ip_addresses}],
400 category="generic",
401 advanced=True,
402 required=False,
403 requires_reload=True,
404 ),
405 ConfigEntry(
406 key=CONF_SMART_FADES_LOG_LEVEL,
407 type=ConfigEntryType.STRING,
408 options=CONF_ENTRY_LOG_LEVEL.options,
409 default_value="GLOBAL",
410 category="audio_analysis",
411 advanced=True,
412 ),
413 ConfigEntry(
414 key=CONF_BACKGROUND_SCAN_CONCURRENCY,
415 type=ConfigEntryType.INTEGER,
416 range=(1, 16),
417 default_value=DEFAULT_BACKGROUND_SCAN_CONCURRENCY,
418 category="audio_analysis",
419 ),
420 )
421
422 async def setup(self, config: CoreConfig) -> None:
423 """Async initialize of module."""
424 # initialize the audio sub-controller (needs mass.streams to be set)
425 self.audio.setup()
426 self._audio_analysis.setup()
427 # copy log level to audio/ffmpeg loggers
428 self.audio.logger.setLevel(self.logger.level)
429 FFMPEG_LOGGER.setLevel(self.logger.level)
430 self._setup_smart_fades_logger(config)
431 # perform check for ffmpeg version
432 await check_ffmpeg_version()
433 # start the webserver
434 self.publish_port = config.get_value(CONF_BIND_PORT, DEFAULT_PORT)
435 configured_publish_ip = str(config.get_value(CONF_PUBLISH_IP) or CONF_VALUE_AUTO)
436 self._configured_publish_ip = (
437 None if configured_publish_ip == CONF_VALUE_AUTO else configured_publish_ip
438 )
439 publish_candidates = await get_publish_ip_candidates(include_ipv6=True)
440 bind_ip = str(config.get_value(CONF_BIND_IP))
441 self._resolve_publish_state(bind_ip, publish_candidates)
442 await self._server.setup(
443 bind_ip=bind_ip,
444 bind_port=cast("int", self.publish_port),
445 static_routes=[
446 (
447 "*",
448 "/flow/{session_id}/{queue_id}/{queue_item_id}/{player_id}.{fmt}",
449 self.serve_queue_flow_stream,
450 ),
451 (
452 "*",
453 "/single/{session_id}/{queue_id}/{queue_item_id}/{player_id}.{fmt}",
454 self.serve_queue_item_stream,
455 ),
456 (
457 "*",
458 "/command/{session_id}/{queue_id}/{command}.mp3",
459 self.serve_command_request,
460 ),
461 ("*", "/announcement/{player_id}.{fmt}", self.serve_announcement_stream),
462 (
463 "GET",
464 LIVE_ANNOUNCEMENT_STREAM_PATH,
465 self.live_announcements.serve_stream,
466 ),
467 ],
468 )
469 # adopt what the server actually bound to: a configured port of 0 is only resolved
470 # by the OS at bind time and an unavailable bind IP falls back to all interfaces
471 self.publish_port = cast("int", self._server.port)
472 self._resolve_publish_state(self._server.bind_ip or DEFAULT_HOST, publish_candidates)
473 # print a big fat message in the log where the streamserver is running
474 # because this is a common source of issues for people with more complex setups
475 self.logger.log(
476 logging.INFO if self.mass.config.onboard_done else logging.WARNING,
477 "\n\n################################################################################\n"
478 "Started streamserver on %s:%s\n"
479 "This is the IP address that is communicated to players.\n"
480 "If this is incorrect, audio will not play!\n"
481 "See the documentation for how to configure the publish IP for the Streamserver\n"
482 "in Settings --> System --> Streams\n"
483 "################################################################################\n",
484 self.publish_ip,
485 self.publish_port,
486 )
487 await self._reload_network_dependent_providers()
488
489 async def post_setup(self) -> None:
490 """Handle logic after all core controllers have been set up."""
491 # the inbound half of a live announcement rides on the webserver: it is the only
492 # one of the two servers that authenticates (and that browsers reach over https)
493 self.live_announcements.setup()
494
495 async def close(self) -> None:
496 """Cleanup on exit."""
497 await self._audio_analysis.close()
498 await self.live_announcements.close()
499 await self._server.close()
500
501 async def resolve_stream_url(self, player_id: str, media: PlayerMedia) -> str:
502 """
503 Resolve the stream URL for the given PlayerMedia.
504
505 :param player_id: The (protocol) player ID requesting the stream.
506 :param media: The PlayerMedia object for which to resolve the stream URL.
507 :return: The resolved stream URL as a string.
508 """
509 if media.media_type in (MediaType.ANNOUNCEMENT, MediaType.FLOW_STREAM):
510 return media.uri
511 protocol_player = self.mass.players.get_player(player_id)
512 conf_output_codec = cast(
513 "str",
514 protocol_player.config.get_value(CONF_OUTPUT_CODEC, default="flac")
515 if protocol_player
516 else "flac",
517 )
518 prefer_wav_for_live_sources = (
519 media.media_type == MediaType.AUDIO_SOURCE
520 and protocol_player is not None
521 and cast(
522 "bool",
523 protocol_player.config.get_value(CONF_PREFER_WAV_FOR_LIVE_SOURCES, default=False),
524 )
525 )
526 output_codec = (
527 ContentType.WAV
528 if prefer_wav_for_live_sources
529 else ContentType.try_parse(conf_output_codec)
530 )
531 fmt = output_codec.value
532 # handle raw pcm without exact format specifiers
533 if output_codec.is_pcm() and ";" not in fmt:
534 fmt += f";codec=pcm;rate={44100};bitrate={16};channels={2}"
535 session_id = media.queue_session_id
536 queue_item_id = media.queue_item_id
537 if not session_id or not queue_item_id:
538 raise InvalidDataError("Can not resolve stream URL: Invalid PlayerMedia data")
539 queue_id = media.source_id
540 queue = self.mass.player_queues.get(queue_id) if queue_id else None
541 crossfade_needs_flow_mode = (
542 # crossfade only applies to tracks; if the queue has it enabled but the player(protocol)
543 # does not support gapless playback, we need to enforce flow mode
544 media.media_type == MediaType.TRACK
545 and queue is not None
546 and queue.crossfade_enabled
547 and protocol_player
548 and not protocol_player.supports_gapless
549 )
550 # the audio overlay is mixed into the queue's continuous (flow) stream;
551 # per-item requests would restart the overlay at every track boundary
552 overlay_needs_flow_mode = queue is not None and overlay_active(queue)
553 # Determine flow_mode based on the actual player's capabilities.
554 # This is done here (just-in-time) because the player's protocol determines this
555 flow_mode = (
556 protocol_player is not None
557 and (protocol_player.flow_mode or crossfade_needs_flow_mode or overlay_needs_flow_mode)
558 and media.media_type not in (MediaType.RADIO, MediaType.AUDIO_SOURCE)
559 )
560 base_path = "flow" if flow_mode else "single"
561 return (
562 f"{self.base_url}/{base_path}/{session_id}/{queue_id}/{queue_item_id}/{player_id}.{fmt}"
563 )
564
565 def update_stream_metadata(
566 self,
567 queue_id: str,
568 source_id: str,
569 provider: str,
570 stream_metadata: StreamMetadata,
571 ) -> None:
572 """
573 Push a live stream metadata update for an active AudioSource queue item.
574
575 Used by plugin providers exposing an AudioSource (e.g. AirPlay receiver,
576 Spotify Connect) to surface live track-change info without restarting
577 the stream. The radio ICY metadata path uses the same underlying field
578 via ``_update_radio_stream_metadata`` but goes through a separate path.
579
580 The update is rejected silently unless the queue's current item is an
581 AudioSource owned by ``provider`` with ``item_id == source_id``. This
582 guard prevents a late callback (e.g. the provider firing one more
583 metadata update after MA has already moved the queue on to a track or
584 a different AudioSource) from stamping unrelated metadata over the
585 wrong item.
586
587 :param queue_id: The queue whose active item should receive the update.
588 :param source_id: The AudioSource.item_id emitting this metadata.
589 :param provider: The provider instance id emitting this metadata.
590 :param stream_metadata: The new stream metadata to attach.
591 """
592 queue = self.mass.player_queues.get(queue_id)
593 if queue is None:
594 return
595 current_item = queue.current_item
596 if current_item is None or current_item.streamdetails is None:
597 return
598 sd = current_item.streamdetails
599 if (
600 sd.media_type != MediaType.AUDIO_SOURCE
601 or sd.provider != provider
602 or sd.item_id != source_id
603 ):
604 # Log at debug so misbehaving providers firing constantly are
605 # diagnosable (count alone is the signal) without spamming higher
606 # log levels for the legitimate transition cases.
607 self.logger.debug(
608 "Rejected metadata update for queue %s from provider %s source %s "
609 "(current item: %s)",
610 queue_id,
611 provider,
612 source_id,
613 sd.uri if sd else "none",
614 )
615 return
616 # Re-check identity *after* preparing the write so a queue advance that
617 # races with this callback can't slip in between the guard above and
618 # the mutation below. Plugins fire these from arbitrary executor /
619 # event-loop threads (AirPlay metadata reader, Spotify webservice
620 # handler, AriaCast WebSocket reader); the GIL keeps each attribute
621 # write atomic but not the read-then-write sequence.
622 if queue.current_item is not current_item or current_item.streamdetails is not sd:
623 return
624 sd.stream_metadata = stream_metadata
625 sd.stream_metadata_last_updated = time.time()
626 self.mass.player_queues.signal_update(queue_id)
627
628 async def serve_queue_item_stream(self, request: web.Request) -> web.StreamResponse: # noqa: PLR0915
629 """Stream single queueitem audio to a player."""
630 self._log_request(request)
631 queue_id = request.match_info["queue_id"]
632 player_id = request.match_info["player_id"]
633 if not (queue := self.mass.player_queues.get(queue_id)):
634 raise web.HTTPNotFound(reason=f"Unknown Queue: {queue_id}")
635 session_id = request.match_info["session_id"]
636 pq_data = self.mass.player_queues.queue_data(queue.queue_id)
637 if pq_data.session_id is None or session_id != pq_data.session_id:
638 raise web.HTTPNotFound(reason=f"Unknown (or invalid) session: {session_id}")
639 if not (player := self.mass.players.get_player(player_id)):
640 raise web.HTTPNotFound(reason=f"Unknown Player: {player_id}")
641 queue_item_id = request.match_info["queue_item_id"]
642 queue_item = self.mass.player_queues.get_item(queue_id, queue_item_id)
643 if not queue_item:
644 raise web.HTTPNotFound(reason=f"Unknown Queue item: {queue_item_id}")
645
646 is_audio_source = (
647 queue_item.media_item is not None
648 and queue_item.media_item.media_type == MediaType.AUDIO_SOURCE
649 )
650
651 # HEAD probes for AudioSource items return a minimal response without
652 # touching the plugin. on_source_selected is the lifecycle hook that
653 # claims ownership and fires off transfer/handoff side effects (stop
654 # the previous player, redirect on disallowed switch, etc.), and a
655 # renderer probing with HEAD before GET should not trigger any of
656 # that. The actual GET request goes through the full hook chain.
657 if request.method != "GET" and is_audio_source:
658 # Validate the providing plugin still exists before advertising the
659 # source. Many DLNA renderers cache HEAD responses; returning 200
660 # for a URI whose plugin has been unloaded would lie to the
661 # renderer and the follow-up GET would fail unrecoverably.
662 assert queue_item.media_item is not None
663 if not isinstance(
664 self.mass.get_provider(queue_item.media_item.provider), PluginProvider
665 ):
666 raise web.HTTPNotFound(
667 reason=f"AudioSource provider {queue_item.media_item.provider} unavailable"
668 )
669 # For PCM-fmt URLs, advertise audio/wav in HEAD: most DLNA renderers
670 # key off the HEAD Content-Type to pick a decoder and do not handle
671 # raw PCM (application/octet-stream). The actual GET response will
672 # still wrap the bytes into a WAV container if needed via the same
673 # mime-type translation downstream.
674 head_fmt = request.match_info["fmt"]
675 if ContentType.try_parse(head_fmt).is_pcm():
676 head_fmt = ContentType.WAV.value
677 headers = {
678 **DEFAULT_STREAM_HEADERS,
679 "icy-name": sanitize_http_header_value(queue_item.name),
680 "contentFeatures.dlna.org": DLNA_CONTENT_FEATURES_REALTIME,
681 "Content-Type": get_mime_type(head_fmt),
682 }
683 resp = web.StreamResponse(status=200, reason="OK", headers=headers)
684 await resp.prepare(request)
685 return resp
686
687 # Fire on_source_selected hook for every AudioSource GET â this is the
688 # single point where exclusive plugin sources claim ownership. Firing
689 # unconditionally (regardless of whether streamdetails are cached from
690 # a previous request) means a disconnect/reconnect for the same queue
691 # item re-claims the lock with a fresh session id, instead of streaming
692 # against the stale ownership of the prior request.
693 # Source identity comes from queue_item.media_item because streamdetails
694 # may not exist yet on the first request.
695 # stream_session_id is a fresh per-request token threaded through to
696 # on_source_unselected so the provider can distinguish a stale
697 # teardown (e.g. a same-queue reconnect's first request completing
698 # AFTER its replacement has already started streaming) from the
699 # currently active session's real teardown.
700 audio_source_provider: PluginProvider | None = None
701 audio_source_id: str | None = None
702 stream_session_id = uuid4().hex
703 if (
704 is_audio_source
705 and queue_item.media_item is not None
706 and (prov := self.mass.get_provider(queue_item.media_item.provider))
707 and isinstance(prov, PluginProvider)
708 ):
709 audio_source_id = queue_item.media_item.item_id
710 # Wire the provider into the finally block BEFORE awaiting the
711 # hook: if the provider partially mutates state (claims the lock,
712 # records the session id) and then raises a non-RuntimeError
713 # exception (buggy plugin, asyncio.CancelledError, etc.), the
714 # finally must still fire on_source_unselected so the lock gets
715 # released. The provider's session-id guard makes a spurious
716 # release a no-op if the lock was never actually claimed.
717 audio_source_provider = prov
718 try:
719 await prov.on_source_selected(
720 audio_source_id, player_id, queue_id, stream_session_id
721 )
722 except RuntimeError as err:
723 # Provider intentionally aborts the original request (e.g.
724 # allow_player_switch=False has just redirected play_media to
725 # the configured target). Surface as 404 so the disallowed
726 # player drops the connection cleanly instead of treating an
727 # uncaught 500 as transient and retrying. The provider
728 # contract requires raising BEFORE claiming, so this is a
729 # clean abort â but we still let the finally run, where the
730 # session-id guard makes the unselect a no-op.
731 self.logger.info(
732 "AudioSource %s aborted stream for player %s: %s",
733 audio_source_id,
734 player_id,
735 err,
736 )
737 raise web.HTTPNotFound(reason=str(err))
738
739 try:
740 if not queue_item.streamdetails:
741 try:
742 queue_item.streamdetails = await self.audio.get_stream_details(
743 queue_item=queue_item
744 )
745 except Exception as e:
746 self.logger.error(
747 "Failed to get streamdetails for QueueItem %s: %s", queue_item_id, e
748 )
749 # a source capacity miss is transient, the item itself is fine
750 if not isinstance(e, ProviderStreamLimitError):
751 queue_item.available = False
752 raise web.HTTPNotFound(
753 reason=f"No streamdetails for Queue item: {queue_item_id}"
754 )
755
756 standard_crossfade_duration = self.mass.config.get_raw_core_config_value(
757 CONF_PLAYER_QUEUES, CONF_CROSSFADE_DURATION, 8
758 )
759 if queue_item.media_type != MediaType.TRACK:
760 crossfade_mode = CrossfadeMode.DISABLED
761 else:
762 crossfade_mode = self.get_crossfade_mode(queue)
763 if (
764 crossfade_mode != CrossfadeMode.DISABLED
765 and PlayerFeature.GAPLESS_PLAYBACK not in player.state.supported_features
766 ):
767 self.logger.warning(
768 "Crossfade disabled: Player %s does not support gapless playback, "
769 "consider enabling flow mode to enable crossfade on this player.",
770 player.state.name,
771 )
772 crossfade_mode = CrossfadeMode.DISABLED
773
774 # pick output format based on the streamdetails and player capabilities
775 pcm_format = await self.audio.select_pcm_format(
776 player=player,
777 streamdetails=queue_item.streamdetails,
778 crossfade_enabled=crossfade_mode != CrossfadeMode.DISABLED,
779 overlay_active=(queue_item.media_type == MediaType.RADIO and overlay_active(queue)),
780 )
781 output_format = await self.audio.get_output_format(
782 output_format_str=request.match_info["fmt"],
783 player=player,
784 content_sample_rate=pcm_format.sample_rate,
785 content_bit_depth=pcm_format.bit_depth,
786 media_type=queue_item.media_type,
787 )
788
789 # prepare request, add some DLNA/UPNP compatible headers
790 # icy-name is sanitized (all control chars, not just newlines) to avoid a
791 # "Potential header injection attack" ValueError by aiohttp
792 # see https://github.com/music-assistant/support/issues/4913
793 # and https://github.com/music-assistant/support/issues/5791
794 # use realtime DLNA flags for radio (sender-paced) since the source delivers slowly
795 dlna_features = (
796 DLNA_CONTENT_FEATURES_REALTIME
797 if queue_item.media_type != MediaType.TRACK
798 else DLNA_CONTENT_FEATURES
799 )
800 headers = {
801 **DEFAULT_STREAM_HEADERS,
802 "icy-name": sanitize_http_header_value(queue_item.name),
803 "contentFeatures.dlna.org": dlna_features,
804 "Content-Type": get_mime_type(output_format.output_format_str),
805 }
806
807 resp = web.StreamResponse(status=200, reason="OK", headers=headers)
808 resp.content_type = get_mime_type(output_format.output_format_str)
809 http_profile = player.get_config_value(CONF_HTTP_PROFILE, "default")
810 if http_profile == "forced_content_length" and not queue_item.duration:
811 # just set an insane high content length to make sure the player keeps playing
812 resp.content_length = calculate_content_length(output_format, 12 * 3600)
813 elif http_profile == "forced_content_length" and queue_item.duration:
814 # estimate content length based on effective duration
815 # account for seek position (e.g., crossfade from previous track)
816 seek_pos = queue_item.streamdetails.seek_position if queue_item.streamdetails else 0
817 effective_duration = max(queue_item.duration - seek_pos, 1)
818 # use cached actual bytes-per-second if available (from a previous stream)
819 resp.content_length = await get_content_length(
820 self.mass, queue_item.uri, output_format, effective_duration
821 )
822 elif http_profile == "chunked":
823 resp.enable_chunked_encoding()
824
825 await resp.prepare(request)
826
827 # return early if this is not a GET request
828 if request.method != "GET":
829 return resp
830
831 self._update_audio_processing_context(
832 queue=queue,
833 queue_item=queue_item,
834 pcm_format=pcm_format,
835 crossfade_mode=crossfade_mode,
836 overlay_enabled=(
837 queue_item.media_type == MediaType.RADIO and overlay_active(queue)
838 ),
839 session_id=session_id,
840 )
841
842 if crossfade_mode != CrossfadeMode.DISABLED:
843 # crossfade is enabled, use special crossfaded single item stream
844 # where the crossfade of the next track is present in the stream of
845 # a single track. This only works if the player supports gapless playback!
846 audio_input = self.audio.get_queue_item_stream_with_smartfade(
847 player=player,
848 queue_item=queue_item,
849 pcm_format=pcm_format,
850 crossfade_mode=crossfade_mode,
851 standard_crossfade_duration=standard_crossfade_duration,
852 session_id=session_id,
853 )
854 else:
855 # no crossfade, just a regular single item stream
856 audio_input = self.audio.get_queue_item_stream(
857 queue_item=queue_item,
858 pcm_format=pcm_format,
859 seek_position=int(queue_item.streamdetails.seek_position),
860 playback_speed=cast(
861 "float", queue_item.extra_attributes.get("playback_speed", 1.0)
862 ),
863 session_id=session_id,
864 )
865 if queue_item.media_type == MediaType.RADIO and overlay_active(queue):
866 # radio plays as a single long-lived stream (never in flow mode),
867 # so mix the audio overlay in here
868 audio_input = self.audio.get_overlay_mixed_stream(queue, audio_input, pcm_format)
869 # stream the audio
870 # this final ffmpeg process in the chain converts raw lossless PCM into
871 # the desired output format for the player including any player specific
872 # filter params such as channels mixing, DSP, resampling and, only if
873 # needed, encoding to lossy formats
874 output_plan = self.audio.get_player_output_plan(
875 player_id=player.player_id,
876 input_format=pcm_format,
877 output_format=output_format,
878 shared_player_ids=player.state.group_members,
879 queue_id=queue_id,
880 session_id=session_id,
881 queue_item_id=queue_item.queue_item_id,
882 )
883 filter_params = output_plan.filter_params
884 # Fast path for live AudioSource: when the player accepts WAV at the
885 # source's exact PCM rate/depth/channels and no filters apply, we
886 # skip the encode ffmpeg entirely and just stream a WAV header
887 # followed by the raw PCM bytes â saves an ffmpeg process and the
888 # latency of its internal buffer on every realtime stream.
889 audio_bytes: AsyncGenerator[bytes]
890 if (
891 queue_item.media_type == MediaType.AUDIO_SOURCE
892 and output_format.content_type == ContentType.WAV
893 and not filter_params
894 and output_format.sample_rate == pcm_format.sample_rate
895 and output_format.bit_depth == pcm_format.bit_depth
896 and output_format.channels == pcm_format.channels
897 ):
898 audio_bytes = _wav_passthrough_stream(audio_input, output_format)
899 else:
900 audio_bytes = get_ffmpeg_stream(
901 audio_input=audio_input,
902 input_format=pcm_format,
903 output_format=output_format,
904 filter_params=filter_params,
905 )
906 first_chunk_received = False
907 bytes_sent = 0
908 # Mark this player as actively streaming so audio analysis yields CPU to playback
909 # for the duration of the transfer (see audio_analysis.playback_active).
910 self._active_output_streams += 1
911 try:
912 # aclosing guarantees the generator (and thus the ffmpeg process chain
913 # behind it) is torn down immediately when the player disconnects
914 # mid-stream, instead of lingering until garbage collection finalizes
915 # the abandoned generator.
916 async with aclosing(audio_bytes):
917 async for chunk in audio_bytes:
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_queue_flow_stream(self, request: web.Request) -> web.StreamResponse: # noqa: PLR0915
1003 """Stream Queue Flow audio to player."""
1004 self._log_request(request)
1005 queue_id = request.match_info["queue_id"]
1006 player_id = request.match_info["player_id"]
1007 if not (queue := self.mass.player_queues.get(queue_id)):
1008 raise web.HTTPNotFound(reason=f"Unknown Queue: {queue_id}")
1009 session_id = request.match_info["session_id"]
1010 queue_data = self.mass.player_queues.queue_data(queue_id)
1011 if queue_data.session_id is None or session_id != queue_data.session_id:
1012 raise web.HTTPNotFound(reason=f"Unknown (or invalid) session: {session_id}")
1013 if not (player := self.mass.players.get_player(player_id)):
1014 raise web.HTTPNotFound(reason=f"Unknown Player: {player_id}")
1015 start_queue_item_id = request.match_info["queue_item_id"]
1016 start_queue_item = self.mass.player_queues.get_item(queue_id, start_queue_item_id)
1017 if not start_queue_item:
1018 raise web.HTTPNotFound(reason=f"Unknown Queue item: {start_queue_item_id}")
1019
1020 # select the PCM format for the flow stream, anchored on the first track
1021 crossfade_mode = (
1022 self.get_crossfade_mode(queue)
1023 if start_queue_item.media_type == MediaType.TRACK
1024 else CrossfadeMode.DISABLED
1025 )
1026 flow_pcm_format = await self.audio.select_flow_pcm_format(
1027 player,
1028 start_streamdetails=start_queue_item.streamdetails,
1029 crossfade_enabled=crossfade_mode != CrossfadeMode.DISABLED,
1030 overlay_active=overlay_active(queue),
1031 )
1032
1033 # work out output format/details
1034 output_format = await self.audio.get_output_format(
1035 output_format_str=request.match_info["fmt"],
1036 player=player,
1037 content_sample_rate=flow_pcm_format.sample_rate,
1038 content_bit_depth=flow_pcm_format.bit_depth,
1039 media_type=start_queue_item.media_type,
1040 )
1041 # work out ICY metadata support
1042 icy_preference = self.mass.config.get_raw_player_config_value(
1043 player_id,
1044 CONF_ENTRY_ENABLE_ICY_METADATA.key,
1045 CONF_ENTRY_ENABLE_ICY_METADATA.default_value,
1046 )
1047 enable_icy = request.headers.get("Icy-MetaData", "") == "1" and icy_preference != "disabled"
1048 icy_meta_interval = 256000 if icy_preference == "full" else 16384
1049
1050 # prepare request, add some DLNA/UPNP compatible headers.
1051 # icy-name (in DEFAULT_STREAM_HEADERS) is always present so players have a
1052 # readable stream name; the rest of the ICY/shoutcast metadata headers are
1053 # only advertised when the client actually requested ICY metadata, rather
1054 # than on every flow response.
1055 headers = {
1056 **DEFAULT_STREAM_HEADERS,
1057 **(ICY_HEADERS if enable_icy else {}),
1058 "contentFeatures.dlna.org": DLNA_CONTENT_FEATURES_REALTIME,
1059 "Content-Type": get_mime_type(output_format.output_format_str),
1060 }
1061 if enable_icy:
1062 headers["icy-metaint"] = str(icy_meta_interval)
1063
1064 resp = web.StreamResponse(status=200, reason="OK", headers=headers)
1065 http_profile = player.get_config_value(CONF_HTTP_PROFILE, "default")
1066 if http_profile == "forced_content_length":
1067 # just set an insane high content length to make sure the player keeps playing
1068 resp.content_length = calculate_content_length(output_format, 12 * 3600)
1069 elif http_profile == "chunked":
1070 resp.enable_chunked_encoding()
1071
1072 await resp.prepare(request)
1073
1074 # return early if this is not a GET request
1075 if request.method != "GET":
1076 return resp
1077
1078 self._update_audio_processing_context(
1079 queue=queue,
1080 queue_item=start_queue_item,
1081 pcm_format=flow_pcm_format,
1082 crossfade_mode=crossfade_mode,
1083 overlay_enabled=overlay_active(queue),
1084 session_id=session_id,
1085 )
1086 output_plan = self.audio.get_player_output_plan(
1087 player.player_id,
1088 flow_pcm_format,
1089 output_format,
1090 shared_player_ids=player.state.group_members,
1091 queue_id=queue_id,
1092 session_id=session_id,
1093 )
1094
1095 # all checks passed, start streaming!
1096 # this final ffmpeg process in the chain will convert the raw, lossless PCM audio into
1097 # the desired output format for the player including any player specific filter params
1098 # such as channels mixing, DSP, resampling and, only if needed, encoding to lossy formats
1099 self.logger.debug("Start serving Queue flow audio stream for %s", queue.display_name)
1100
1101 # Mark this player as actively streaming so audio analysis yields CPU to playback
1102 # for the duration of the flow stream (see audio_analysis.playback_active).
1103 self._active_output_streams += 1
1104 flow_stream = self.audio.get_queue_flow_stream(
1105 queue=queue,
1106 start_queue_item=start_queue_item,
1107 pcm_format=flow_pcm_format,
1108 session_id=session_id,
1109 protocol_player=player,
1110 )
1111 if overlay_active(queue):
1112 flow_stream = self.audio.get_overlay_mixed_stream(queue, flow_stream, flow_pcm_format)
1113 audio_bytes = get_ffmpeg_stream(
1114 audio_input=flow_stream,
1115 input_format=flow_pcm_format,
1116 output_format=output_format,
1117 filter_params=output_plan.filter_params,
1118 # we need to slowly feed the music to avoid the player stopping and later
1119 # restarting (or completely failing) the audio stream by keeping the buffer short.
1120 # this is reported to be an issue especially with Chromecast players.
1121 # see for example: https://github.com/music-assistant/support/issues/3717
1122 # allow buffer ahead of a few seconds and read rest in (near) realtime
1123 extra_input_args=["-readrate", "1.05", "-readrate_initial_burst", "5"],
1124 chunk_size=icy_meta_interval if enable_icy else calculate_content_length(output_format),
1125 )
1126 client_disconnected = False
1127 try:
1128 # aclosing guarantees the flow stream (and thus the ffmpeg process chain
1129 # behind it) is torn down immediately when the player disconnects
1130 # mid-stream, instead of lingering until garbage collection finalizes
1131 # the abandoned generator.
1132 async with aclosing(audio_bytes):
1133 async for chunk in audio_bytes:
1134 try:
1135 await resp.write(chunk)
1136 except BrokenPipeError, ConnectionResetError, ConnectionError:
1137 # race condition
1138 client_disconnected = True
1139 break
1140
1141 if not enable_icy:
1142 continue
1143
1144 # if icy metadata is enabled, send the icy metadata after the chunk
1145 if (
1146 # use current item here and not buffered item, otherwise
1147 # the icy metadata will be too much ahead
1148 (current_item := queue.current_item)
1149 and current_item.streamdetails
1150 and current_item.streamdetails.stream_title
1151 ):
1152 title = current_item.streamdetails.stream_title
1153 elif queue and current_item and current_item.name:
1154 title = current_item.name
1155 else:
1156 title = "Music Assistant"
1157 metadata = f"StreamTitle='{title}';".encode()
1158 if icy_preference == "full" and current_item and current_item.image:
1159 metadata += f"StreamURL='{current_item.image.path}'".encode()
1160 while len(metadata) % 16 != 0:
1161 metadata += b"\x00"
1162 length = len(metadata)
1163 length_b = chr(int(length / 16)).encode()
1164 await resp.write(length_b + metadata)
1165 finally:
1166 self._active_output_streams -= 1
1167
1168 if not client_disconnected and http_profile == "forced_content_length":
1169 await self._finish_flow_stream(resp, queue_id, session_id)
1170
1171 return resp
1172
1173 async def serve_command_request(self, request: web.Request) -> web.FileResponse:
1174 """Handle special 'command' request for a player."""
1175 self._log_request(request)
1176 queue_id = request.match_info["queue_id"]
1177 session_id = request.match_info["session_id"]
1178 queue_data = self.mass.player_queues.queue_data_or_none(queue_id)
1179 if queue_data is None or queue_data.session_id != session_id:
1180 raise web.HTTPNotFound(reason=f"Unknown (or invalid) session: {session_id}")
1181 command = request.match_info["command"]
1182 if command == "next":
1183 self.mass.create_task(self.mass.player_queues.next(queue_id))
1184 return web.FileResponse(SILENCE_FILE, headers={"icy-name": "Music Assistant"})
1185
1186 async def serve_announcement_stream(self, request: web.Request) -> web.StreamResponse:
1187 """Stream announcement audio to a player."""
1188 self._log_request(request)
1189 player_id = request.match_info["player_id"]
1190 if not (player := self.mass.players.get_player(player_id)):
1191 raise web.HTTPNotFound(reason=f"Unknown Player: {player_id}")
1192 if not (announce_data := self.announcement_renderer.get_for_player(player_id)):
1193 raise web.HTTPNotFound(reason=f"No pending announcements for Player: {player_id}")
1194
1195 # work out output format/details
1196 fmt = request.match_info["fmt"]
1197 audio_format = AudioFormat(content_type=ContentType.try_parse(fmt))
1198
1199 http_profile = self._get_announcement_http_profile(player_id, announce_data)
1200
1201 # return early if this is not a GET request:
1202 # players often probe the url with a HEAD request before fetching it and
1203 # rendering the announcement for such a probe would run the entire (costly)
1204 # TTS/ffmpeg chain twice for a single announcement.
1205 if request.method != "GET":
1206 resp = web.StreamResponse(status=200, reason="OK", headers=DEFAULT_STREAM_HEADERS)
1207 resp.content_type = get_mime_type(audio_format.output_format_str)
1208 if http_profile == "chunked":
1209 resp.enable_chunked_encoding()
1210 await resp.prepare(request)
1211 return resp
1212
1213 if http_profile == "forced_content_length":
1214 # given the fact that an announcement is just a short audio clip,
1215 # just send it over completely at once so we have a fixed content length
1216 data = bytearray()
1217 announcement_stream = self.get_announcement_stream(announce_data, audio_format)
1218 # aclosing guarantees the stream (and thus the ffmpeg process chain behind
1219 # it) is torn down immediately when the request is cancelled, instead of
1220 # lingering until garbage collection finalizes the abandoned generator.
1221 async with aclosing(announcement_stream):
1222 async for chunk in announcement_stream:
1223 data += chunk
1224 return web.Response(
1225 body=bytes(data),
1226 content_type=get_mime_type(audio_format.output_format_str),
1227 headers=DEFAULT_STREAM_HEADERS,
1228 )
1229
1230 resp = web.StreamResponse(status=200, reason="OK", headers=DEFAULT_STREAM_HEADERS)
1231 resp.content_type = get_mime_type(audio_format.output_format_str)
1232 if http_profile == "chunked":
1233 resp.enable_chunked_encoding()
1234
1235 await resp.prepare(request)
1236
1237 # all checks passed, start streaming!
1238 self.logger.debug(
1239 "Start serving audio stream for Announcement %s to %s",
1240 announce_data["announcement_url"],
1241 player.display_name,
1242 )
1243 announcement_stream = self.get_announcement_stream(announce_data, audio_format)
1244 # aclosing guarantees the stream (and thus the ffmpeg process chain behind
1245 # it) is torn down immediately when the player disconnects mid-stream,
1246 # instead of lingering until garbage collection finalizes the abandoned
1247 # generator.
1248 async with aclosing(announcement_stream):
1249 async for chunk in announcement_stream:
1250 try:
1251 await resp.write(chunk)
1252 except BrokenPipeError, ConnectionResetError:
1253 break
1254
1255 self.logger.debug(
1256 "Finished serving audio stream for Announcement %s to %s",
1257 announce_data["announcement_url"],
1258 player.display_name,
1259 )
1260
1261 return resp
1262
1263 def get_command_url(self, player_or_queue_id: str, command: str) -> str | None:
1264 """
1265 Get the url for the special command stream, or None if the queue is not playing.
1266
1267 :param player_or_queue_id: Queue to send the command to.
1268 :param command: Command the url triggers when fetched.
1269 """
1270 queue_data = self.mass.player_queues.queue_data_or_none(player_or_queue_id)
1271 if queue_data is None or (session_id := queue_data.session_id) is None:
1272 return None
1273 return f"{self.base_url}/command/{session_id}/{player_or_queue_id}/{command}.mp3"
1274
1275 def get_announcement_url(
1276 self,
1277 player_id: str,
1278 content_type: ContentType = ContentType.MP3,
1279 ) -> str:
1280 """
1281 Get the url that serves the announcement registered for the given player.
1282
1283 :param player_id: The player the announcement is played on.
1284 :param content_type: The format to serve the announcement in.
1285 """
1286 # use stream server to host announcement on local network
1287 # this ensures playback on all players, including ones that do not
1288 # like https hosts and it also offers the pre-announce 'bell'
1289 return f"{self.base_url}/announcement/{player_id}.{content_type.value}"
1290
1291 def get_stream(
1292 self,
1293 media: PlayerMedia,
1294 pcm_format: AudioFormat,
1295 player_id: str | None = None,
1296 force_flow_mode: bool = False,
1297 use_flow_stream_buffering: bool = False,
1298 ) -> AsyncGenerator[bytes]:
1299 """
1300 Get a stream of the given media as raw PCM audio.
1301
1302 This is used as helper for player providers that can consume the raw PCM
1303 audio stream directly (e.g. AirPlay) and not rely on HTTP transport.
1304
1305 :param media: The PlayerMedia to stream.
1306 :param pcm_format: The desired output PCM format.
1307 :param player_id: The player ID requesting the stream. Used to determine
1308 if flow mode should be used based on the player's capabilities.
1309 :param force_flow_mode: Force flow mode regardless of player capabilities.
1310 Used for multi-client streaming scenarios that require continuous streams.
1311 :param use_flow_stream_buffering: Buffer the flow stream to provide headroom
1312 during smart fades transitions. Use for consumers that read directly
1313 (e.g. AirPlay, Snapcast) and can't tolerate stalls.
1314 """
1315 # select audio source
1316 if media.media_type == MediaType.ANNOUNCEMENT:
1317 # special case: stream announcement
1318 assert media.custom_data
1319 return self.get_announcement_stream(cast("AnnounceData", media.custom_data), pcm_format)
1320 if (
1321 media.source_id
1322 and media.source_id.startswith(UGP_PREFIX)
1323 and media.uri
1324 and "/ugp/" in media.uri
1325 ):
1326 # special case: member player accessing UGP stream
1327 # Check URI to distinguish from the UGP accessing its own stream
1328 ugp_player = cast("UniversalGroupPlayer", self.mass.players.get_player(media.source_id))
1329 ugp_stream = ugp_player.stream
1330 assert ugp_stream is not None # for type checker
1331 if ugp_stream.base_pcm_format == pcm_format:
1332 # no conversion needed
1333 return ugp_stream.subscribe_raw()
1334 return ugp_stream.get_stream(output_format=pcm_format)
1335 if media.source_id and media.queue_item_id:
1336 # Queue stream request - determine flow_mode based on player capabilities
1337 # or force it if explicitly requested (e.g., for multi-client streaming)
1338 protocol_player = self.mass.players.get_player(player_id) if player_id else None
1339 queue_id = media.source_id
1340 queue = self.mass.player_queues.get(queue_id)
1341 queue_session_id = media.queue_session_id
1342 crossfade_needs_flow_mode = (
1343 # crossfade only applies to tracks; if the queue has it enabled but the
1344 # player(protocol) does not support gapless playback, we need to enforce flow mode
1345 media.media_type == MediaType.TRACK
1346 and queue is not None
1347 and queue.crossfade_enabled
1348 and protocol_player
1349 and not protocol_player.supports_gapless
1350 )
1351 # the audio overlay is mixed into the queue's continuous (flow) stream;
1352 # per-item requests would restart the overlay at every track boundary
1353 overlay_needs_flow_mode = queue is not None and overlay_active(queue)
1354 flow_mode = (
1355 force_flow_mode
1356 or (protocol_player is not None and protocol_player.flow_mode)
1357 or crossfade_needs_flow_mode
1358 or overlay_needs_flow_mode
1359 )
1360 if media.media_type in (MediaType.RADIO, MediaType.AUDIO_SOURCE):
1361 # flow_mode for live/infinite streams is pointless
1362 flow_mode = False
1363 if flow_mode:
1364 # flow stream request
1365 assert queue
1366 start_queue_item = self.mass.player_queues.get_item(
1367 media.source_id, media.queue_item_id
1368 )
1369 assert start_queue_item
1370 crossfade_mode = (
1371 self.get_crossfade_mode(queue)
1372 if start_queue_item.media_type == MediaType.TRACK
1373 else CrossfadeMode.DISABLED
1374 )
1375 self._update_audio_processing_context(
1376 queue=queue,
1377 queue_item=start_queue_item,
1378 pcm_format=pcm_format,
1379 crossfade_mode=crossfade_mode,
1380 overlay_enabled=overlay_active(queue),
1381 session_id=queue_session_id,
1382 )
1383 flow_stream = self.audio.get_queue_flow_stream(
1384 queue=queue,
1385 start_queue_item=start_queue_item,
1386 pcm_format=pcm_format,
1387 session_id=queue_session_id,
1388 protocol_player=protocol_player,
1389 )
1390 if overlay_active(queue):
1391 flow_stream = self.audio.get_overlay_mixed_stream(
1392 queue, flow_stream, pcm_format
1393 )
1394 if use_flow_stream_buffering:
1395 return buffered(flow_stream, buffer_size=30, min_buffer_before_yield=1)
1396 return flow_stream
1397 # single item stream (e.g. radio or non-flow mode)
1398 queue_item = self.mass.player_queues.get_item(media.source_id, media.queue_item_id)
1399 assert queue_item
1400 if queue is not None:
1401 self._update_audio_processing_context(
1402 queue=queue,
1403 queue_item=queue_item,
1404 pcm_format=pcm_format,
1405 crossfade_mode=(
1406 self.get_crossfade_mode(queue)
1407 if queue_item.media_type == MediaType.TRACK
1408 else CrossfadeMode.DISABLED
1409 ),
1410 overlay_enabled=(
1411 queue_item.media_type == MediaType.RADIO and overlay_active(queue)
1412 ),
1413 session_id=queue_session_id,
1414 )
1415 inner_stream = self.audio.get_queue_item_stream(
1416 queue_item=queue_item,
1417 pcm_format=pcm_format,
1418 seek_position=(
1419 int(queue_item.streamdetails.seek_position) if queue_item.streamdetails else 0
1420 ),
1421 playback_speed=cast(
1422 "float", queue_item.extra_attributes.get("playback_speed", 1.0)
1423 ),
1424 session_id=queue_session_id,
1425 )
1426 if (
1427 queue is not None
1428 and queue_item.media_type == MediaType.RADIO
1429 and overlay_active(queue)
1430 ):
1431 # radio plays as a single long-lived stream, so mix the overlay in here
1432 inner_stream = self.audio.get_overlay_mixed_stream(queue, inner_stream, pcm_format)
1433 # mirror the on_source_selected/unselected lifecycle the HTTP route
1434 # fires, so direct-PCM consumers (AirPlay, Snapcast, UGP) honour the
1435 # plugin contract too
1436 if (
1437 queue_item.media_item is not None
1438 and queue_item.media_item.media_type == MediaType.AUDIO_SOURCE
1439 ):
1440 return self._wrap_with_audio_source_lifecycle(
1441 inner=inner_stream,
1442 queue_item=queue_item,
1443 player_id=player_id or media.source_id,
1444 )
1445 return inner_stream
1446 # assume url or some other direct path
1447 # NOTE: this will fail if its an uri not playable by ffmpeg
1448 return get_ffmpeg_stream(
1449 audio_input=media.uri,
1450 input_format=AudioFormat(content_type=ContentType.try_parse(media.uri)),
1451 output_format=pcm_format,
1452 )
1453
1454 async def get_preview_stream(
1455 self,
1456 provider_instance_id_or_domain: str,
1457 item_id: str,
1458 media_type: MediaType = MediaType.TRACK,
1459 ) -> AsyncGenerator[bytes]:
1460 """Create a 30 seconds preview audioclip for the given media item."""
1461 if not (music_prov := self.mass.get_provider(provider_instance_id_or_domain)):
1462 raise ProviderUnavailableError
1463 if music_prov.type != ProviderType.MUSIC:
1464 msg = f"{provider_instance_id_or_domain} is not a music provider"
1465 raise InvalidDataError(msg)
1466 music_prov = cast("MusicProvider", music_prov)
1467
1468 try:
1469 await self.mass.music.get_item(
1470 media_type,
1471 item_id,
1472 provider_instance_id_or_domain,
1473 allow_update_metadata=False,
1474 )
1475 except MediaNotFoundError as err:
1476 msg = f"Item {item_id} not found in provider {provider_instance_id_or_domain}"
1477 raise InvalidDataError(msg) from err
1478
1479 streamdetails = await music_prov.get_stream_details(item_id, media_type)
1480 pcm_format = AudioFormat(
1481 content_type=ContentType.from_bit_depth(streamdetails.audio_format.bit_depth),
1482 sample_rate=streamdetails.audio_format.sample_rate,
1483 bit_depth=streamdetails.audio_format.bit_depth,
1484 channels=streamdetails.audio_format.channels,
1485 )
1486 async for chunk in get_ffmpeg_stream(
1487 audio_input=self.audio.get_media_stream(
1488 streamdetails=streamdetails, pcm_format=pcm_format
1489 ),
1490 input_format=pcm_format,
1491 output_format=AudioFormat(content_type=ContentType.AAC),
1492 extra_input_args=["-t", "30"],
1493 ):
1494 yield chunk
1495
1496 async def get_announcement_stream(
1497 self, announce_data: AnnounceData, output_format: AudioFormat
1498 ) -> AsyncGenerator[bytes]:
1499 """
1500 Get the audio of an announcement (pre-announce chime + announcement).
1501
1502 Any number of consumers may stream the same announcement at once; its source is
1503 fetched and decoded only once. The audio stays available while the stream is
1504 held open.
1505
1506 :param announce_data: The announcement to stream.
1507 :param output_format: The format to deliver the audio in.
1508 """
1509 render = self.announcement_renderer.acquire(announce_data)
1510 try:
1511 # aclosing guarantees this consumer's ffmpeg encoder is torn down
1512 # immediately when it goes away, instead of lingering until garbage
1513 # collection finalizes the abandoned generator.
1514 stream = render.get_stream(output_format)
1515 async with aclosing(stream):
1516 async for chunk in stream:
1517 yield chunk
1518 finally:
1519 await self.announcement_renderer.release(render)
1520
1521 async def get_announcement_duration(
1522 self, announcement: PlayerMedia, timeout: float = DEFAULT_RENDER_TIMEOUT
1523 ) -> int | None:
1524 """
1525 Get the exact duration (in seconds) of an announcement, once it finished rendering.
1526
1527 Waits for the audio to be rendered in full, so call this while the announcement
1528 plays rather than before handing it to a player. Returns None when the length can
1529 not be determined, e.g. the announcement is no longer playing or its source did
1530 not deliver in time.
1531
1532 :param announcement: The announcement to return the duration for.
1533 :param timeout: Maximum time to wait for the audio to finish rendering.
1534 """
1535 if announcement.duration:
1536 return announcement.duration
1537 if not announcement.custom_data:
1538 return None
1539 render = self.announcement_renderer.get(cast("AnnounceData", announcement.custom_data))
1540 if render is None:
1541 return None
1542 duration = await render.wait_finished(timeout)
1543 return ceil(duration) if duration else None
1544
1545 async def _wrap_with_audio_source_lifecycle(
1546 self,
1547 inner: AsyncGenerator[bytes],
1548 queue_item: QueueItem,
1549 player_id: str,
1550 ) -> AsyncGenerator[bytes]:
1551 """
1552 Wrap an AudioSource queue item stream with on_source_selected/unselected hooks.
1553
1554 Direct-PCM consumers (AirPlay, Snapcast, UGP, ...) call ``get_stream`` instead
1555 of going through the HTTP route, but the plugin contract requires the
1556 lifecycle hooks to fire for every actual stream request â they're what
1557 claim/release the per-queue exclusive ownership and trigger acquisition
1558 side effects like the Spotify Connect Web API play kick. This wrapper
1559 gives those consumers the same lifecycle the HTTP route already provides.
1560
1561 :param inner: The underlying audio stream generator.
1562 :param queue_item: The AudioSource queue item being streamed.
1563 :param player_id: The protocol player consuming this stream.
1564 """
1565 media_item = queue_item.media_item
1566 assert media_item is not None # caller checked media_type == AUDIO_SOURCE
1567 prov = self.mass.get_provider(media_item.provider)
1568 queue_id = queue_item.queue_id
1569 if not isinstance(prov, PluginProvider):
1570 async for chunk in inner:
1571 yield chunk
1572 return
1573 source_id = media_item.item_id
1574 stream_session_id = uuid4().hex
1575 # single try/finally so on_source_unselected fires even when
1576 # on_source_selected raises after partially claiming state; the
1577 # provider's session_id guard makes a no-op claim release safe.
1578 try:
1579 try:
1580 await prov.on_source_selected(source_id, player_id, queue_id, stream_session_id)
1581 except RuntimeError as err:
1582 # provider intentionally aborts the request â surface as AudioError
1583 raise AudioError(str(err)) from err
1584 async for chunk in inner:
1585 yield chunk
1586 finally:
1587 try:
1588 await prov.on_source_unselected(source_id, queue_id, stream_session_id)
1589 except Exception as err:
1590 self.logger.exception(
1591 "on_source_unselected raised for provider %s source %s queue %s: %s",
1592 prov.instance_id,
1593 source_id,
1594 queue_id,
1595 err,
1596 )
1597
1598 def _update_audio_processing_context(
1599 self,
1600 queue: PlayerQueue,
1601 queue_item: QueueItem,
1602 pcm_format: AudioFormat,
1603 crossfade_mode: CrossfadeMode,
1604 overlay_enabled: bool,
1605 session_id: str | None = None,
1606 ) -> None:
1607 """
1608 Store the shared processing context selected for a queue item.
1609
1610 :param queue: Active player queue.
1611 :param queue_item: Queue item being prepared.
1612 :param pcm_format: Shared PCM format leaving queue processing.
1613 :param crossfade_mode: Effective crossfade mode for the item.
1614 :param overlay_enabled: Whether an overlay is mixed into this stream.
1615 :param session_id: Queue session that owns processing-detail updates.
1616 """
1617 if queue_item.streamdetails is None:
1618 return
1619 queue_data = self.mass.player_queues.queue_data_or_none(queue.queue_id)
1620 if (
1621 queue_data is None
1622 or (processing_session_id := session_id or queue_data.session_id) is None
1623 or queue_data.session_id != processing_session_id
1624 ):
1625 return
1626 self.audio_processing.start_session(queue.queue_id, processing_session_id)
1627 self.audio_processing.update_item_context(
1628 queue_id=queue.queue_id,
1629 session_id=processing_session_id,
1630 queue_item_id=queue_item.queue_item_id,
1631 queue_processing=AudioQueueProcessing(
1632 pcm_format=pcm_format,
1633 playback_speed=cast(
1634 "float",
1635 queue_item.extra_attributes.get("playback_speed", 1.0),
1636 ),
1637 crossfade_mode=crossfade_mode,
1638 overlay_active=overlay_enabled,
1639 ),
1640 alters_audio=queue_item.streamdetails.fade_in,
1641 )
1642
1643 def _get_announcement_http_profile(self, player_id: str, announce_data: AnnounceData) -> str:
1644 """
1645 Resolve the http profile for serving an announcement stream.
1646
1647 Announcement urls are registered under the visible player's id, but the
1648 stream may be fetched by a linked protocol player; the profile must come
1649 from the player that actually performs the fetch.
1650 """
1651 announce_player = None
1652 if announce_player_id := announce_data.get("announce_player_id"):
1653 announce_player = self.mass.players.get_player(announce_player_id)
1654 if announce_player is None:
1655 announce_player = self.mass.players.get_player(player_id)
1656 if announce_player is None:
1657 return "default"
1658 return announce_player.get_output_config_value(CONF_HTTP_PROFILE, "default")
1659
1660 async def _finish_flow_stream(
1661 self, resp: web.StreamResponse, queue_id: str, session_id: str
1662 ) -> None:
1663 """
1664 Close a fully served flow stream, giving the player time to drain when it ends the queue.
1665
1666 :param resp: The flow stream response, already fully written.
1667 :param queue_id: Id of the queue the flow stream belongs to.
1668 :param session_id: Stream session this response was opened for.
1669 """
1670 if self.mass.player_queues.flow_queue_exhausted(queue_id, session_id):
1671 # the player is still holding a few seconds of audio it has not rendered yet
1672 # and drops that as soon as the stream ends, so let it play out first.
1673 # a flow that ends to be restarted right away gets no such grace: there the
1674 # player should go idle as soon as possible so the next stream can start.
1675 self.logger.debug(
1676 "Flow stream for queue %s reached the end of the queue - holding the "
1677 "connection open for %ss so the player can play out its buffer",
1678 queue_id,
1679 FLOW_STREAM_LEAD_OUT_SECONDS,
1680 )
1681 await asyncio.sleep(FLOW_STREAM_LEAD_OUT_SECONDS)
1682 # aiohttp derives keep-alive from the request, so the 'Connection: close' we
1683 # advertise is relayed to the player but never applied to the response itself.
1684 # Without this the player is left waiting on a stream that already ended.
1685 resp.force_close()
1686
1687 def _log_request(self, request: web.Request) -> None:
1688 """Log request."""
1689 if self.logger.isEnabledFor(VERBOSE_LOG_LEVEL):
1690 self.logger.log(
1691 VERBOSE_LOG_LEVEL,
1692 "Got %s request to %s from %s\nheaders: %s\n",
1693 request.method,
1694 request.path,
1695 request.remote,
1696 redact_sensitive_headers(request.headers),
1697 )
1698 else:
1699 self.logger.debug(
1700 "Got %s request to %s from %s (HTTP/%s.%s, connection: %s)",
1701 request.method,
1702 request.path,
1703 request.remote,
1704 request.version.major,
1705 request.version.minor,
1706 request.headers.get("Connection", "-"),
1707 )
1708
1709 async def _reload_network_dependent_providers(self) -> None:
1710 """Reload the providers that captured the streamserver network, if it changed."""
1711 previous = self._network_fingerprint
1712 current = (
1713 self._bind_ip,
1714 str(self.publish_ip),
1715 cast("int", self.publish_port),
1716 tuple(self._publish_addresses),
1717 )
1718 if previous is None or previous == current:
1719 self._network_fingerprint = current
1720 return
1721 # these providers bind or advertise the network while they load, so a plain
1722 # reload is what moves them over - they share no lighter rebind path
1723 instance_ids = [
1724 prov.instance_id
1725 for prov in self.mass.providers
1726 if prov.reload_on_streams_network_change
1727 ]
1728 for instance_id in instance_ids:
1729 try:
1730 config = await self.mass.config.get_provider_config(instance_id)
1731 self.logger.info(
1732 "Streamserver network changed, reloading provider %s",
1733 config.name or config.domain,
1734 )
1735 await self.mass.load_provider_config(config)
1736 except Exception as err:
1737 self.logger.warning(
1738 "Error reloading provider %s: %s",
1739 instance_id,
1740 str(err) or err.__class__.__name__,
1741 exc_info=err,
1742 )
1743 # only mark the new network as applied once the loop completed, so a run cut short
1744 # by a second config change runs again on the next reload
1745 self._network_fingerprint = current
1746
1747 def _setup_smart_fades_logger(self, config: CoreConfig) -> None:
1748 """Set up smart fades logger level."""
1749 log_level = str(config.get_value(CONF_SMART_FADES_LOG_LEVEL))
1750 if log_level == "GLOBAL":
1751 self.audio.smart_fades_mixer.logger.setLevel(self.logger.level)
1752 else:
1753 self.audio.smart_fades_mixer.logger.setLevel(log_level)
1754
1755 def _resolve_publish_state(self, bind_ip: str, publish_candidates: tuple[str, ...]) -> None:
1756 """
1757 Resolve the addresses and base URL to advertise for the given bind address.
1758
1759 Reads ``self.publish_port``, so set that first.
1760
1761 :param bind_ip: Address the streamserver binds to (a wildcard means all interfaces).
1762 :param publish_candidates: Host addresses reachable from the local network, ranked.
1763 """
1764 self._bind_ip = bind_ip
1765 self._publish_addresses = _get_publish_addresses(
1766 bind_ip, self._configured_publish_ip, publish_candidates
1767 )
1768 # the single address players are handed, taken from the top of the ranked list
1769 self.publish_ip = self._publish_addresses[0]
1770 self._base_url = f"http://{format_ip_for_url(self.publish_ip)}:{self.publish_port}"
1771
1772
1773def _same_ip_family(ip: str, other_ip: str) -> bool:
1774 """Return whether two addresses belong to the same IP family."""
1775 return (":" in ip) == (":" in other_ip)
1776