/
/
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 ("*", "/command/{queue_id}/{command}.mp3", self.serve_command_request),
457 ("*", "/announcement/{player_id}.{fmt}", self.serve_announcement_stream),
458 (
459 "GET",
460 LIVE_ANNOUNCEMENT_STREAM_PATH,
461 self.live_announcements.serve_stream,
462 ),
463 ],
464 )
465 # adopt what the server actually bound to: a configured port of 0 is only resolved
466 # by the OS at bind time and an unavailable bind IP falls back to all interfaces
467 self.publish_port = cast("int", self._server.port)
468 self._resolve_publish_state(self._server.bind_ip or DEFAULT_HOST, publish_candidates)
469 # print a big fat message in the log where the streamserver is running
470 # because this is a common source of issues for people with more complex setups
471 self.logger.log(
472 logging.INFO if self.mass.config.onboard_done else logging.WARNING,
473 "\n\n################################################################################\n"
474 "Started streamserver on %s:%s\n"
475 "This is the IP address that is communicated to players.\n"
476 "If this is incorrect, audio will not play!\n"
477 "See the documentation for how to configure the publish IP for the Streamserver\n"
478 "in Settings --> System --> Streams\n"
479 "################################################################################\n",
480 self.publish_ip,
481 self.publish_port,
482 )
483 await self._reload_network_dependent_providers()
484
485 async def post_setup(self) -> None:
486 """Handle logic after all core controllers have been set up."""
487 # the inbound half of a live announcement rides on the webserver: it is the only
488 # one of the two servers that authenticates (and that browsers reach over https)
489 self.live_announcements.setup()
490
491 async def close(self) -> None:
492 """Cleanup on exit."""
493 await self._audio_analysis.close()
494 await self.live_announcements.close()
495 await self._server.close()
496
497 async def resolve_stream_url(self, player_id: str, media: PlayerMedia) -> str:
498 """
499 Resolve the stream URL for the given PlayerMedia.
500
501 :param player_id: The (protocol) player ID requesting the stream.
502 :param media: The PlayerMedia object for which to resolve the stream URL.
503 :return: The resolved stream URL as a string.
504 """
505 if media.media_type in (MediaType.ANNOUNCEMENT, MediaType.FLOW_STREAM):
506 return media.uri
507 protocol_player = self.mass.players.get_player(player_id)
508 conf_output_codec = cast(
509 "str",
510 protocol_player.config.get_value(CONF_OUTPUT_CODEC, default="flac")
511 if protocol_player
512 else "flac",
513 )
514 prefer_wav_for_live_sources = (
515 media.media_type == MediaType.AUDIO_SOURCE
516 and protocol_player is not None
517 and cast(
518 "bool",
519 protocol_player.config.get_value(CONF_PREFER_WAV_FOR_LIVE_SOURCES, default=False),
520 )
521 )
522 output_codec = (
523 ContentType.WAV
524 if prefer_wav_for_live_sources
525 else ContentType.try_parse(conf_output_codec)
526 )
527 fmt = output_codec.value
528 # handle raw pcm without exact format specifiers
529 if output_codec.is_pcm() and ";" not in fmt:
530 fmt += f";codec=pcm;rate={44100};bitrate={16};channels={2}"
531 extra_data = media.custom_data or {}
532 session_id = extra_data.get("session_id")
533 queue_item_id = media.queue_item_id
534 if not session_id or not queue_item_id:
535 raise InvalidDataError("Can not resolve stream URL: Invalid PlayerMedia data")
536 queue_id = media.source_id
537 queue = self.mass.player_queues.get(queue_id) if queue_id else None
538 crossfade_needs_flow_mode = (
539 # crossfade only applies to tracks; if the queue has it enabled but the player(protocol)
540 # does not support gapless playback, we need to enforce flow mode
541 media.media_type == MediaType.TRACK
542 and queue is not None
543 and queue.crossfade_enabled
544 and protocol_player
545 and not protocol_player.supports_gapless
546 )
547 # the audio overlay is mixed into the queue's continuous (flow) stream;
548 # per-item requests would restart the overlay at every track boundary
549 overlay_needs_flow_mode = queue is not None and overlay_active(queue)
550 # Determine flow_mode based on the actual player's capabilities.
551 # This is done here (just-in-time) because the player's protocol determines this
552 flow_mode = (
553 protocol_player is not None
554 and (protocol_player.flow_mode or crossfade_needs_flow_mode or overlay_needs_flow_mode)
555 and media.media_type not in (MediaType.RADIO, MediaType.AUDIO_SOURCE)
556 )
557 base_path = "flow" if flow_mode else "single"
558 return (
559 f"{self.base_url}/{base_path}/{session_id}/{queue_id}/{queue_item_id}/{player_id}.{fmt}"
560 )
561
562 def update_stream_metadata(
563 self,
564 queue_id: str,
565 source_id: str,
566 provider: str,
567 stream_metadata: StreamMetadata,
568 ) -> None:
569 """
570 Push a live stream metadata update for an active AudioSource queue item.
571
572 Used by plugin providers exposing an AudioSource (e.g. AirPlay receiver,
573 Spotify Connect) to surface live track-change info without restarting
574 the stream. The radio ICY metadata path uses the same underlying field
575 via ``_update_radio_stream_metadata`` but goes through a separate path.
576
577 The update is rejected silently unless the queue's current item is an
578 AudioSource owned by ``provider`` with ``item_id == source_id``. This
579 guard prevents a late callback (e.g. the provider firing one more
580 metadata update after MA has already moved the queue on to a track or
581 a different AudioSource) from stamping unrelated metadata over the
582 wrong item.
583
584 :param queue_id: The queue whose active item should receive the update.
585 :param source_id: The AudioSource.item_id emitting this metadata.
586 :param provider: The provider instance id emitting this metadata.
587 :param stream_metadata: The new stream metadata to attach.
588 """
589 queue = self.mass.player_queues.get(queue_id)
590 if queue is None:
591 return
592 current_item = queue.current_item
593 if current_item is None or current_item.streamdetails is None:
594 return
595 sd = current_item.streamdetails
596 if (
597 sd.media_type != MediaType.AUDIO_SOURCE
598 or sd.provider != provider
599 or sd.item_id != source_id
600 ):
601 # Log at debug so misbehaving providers firing constantly are
602 # diagnosable (count alone is the signal) without spamming higher
603 # log levels for the legitimate transition cases.
604 self.logger.debug(
605 "Rejected metadata update for queue %s from provider %s source %s "
606 "(current item: %s)",
607 queue_id,
608 provider,
609 source_id,
610 sd.uri if sd else "none",
611 )
612 return
613 # Re-check identity *after* preparing the write so a queue advance that
614 # races with this callback can't slip in between the guard above and
615 # the mutation below. Plugins fire these from arbitrary executor /
616 # event-loop threads (AirPlay metadata reader, Spotify webservice
617 # handler, AriaCast WebSocket reader); the GIL keeps each attribute
618 # write atomic but not the read-then-write sequence.
619 if queue.current_item is not current_item or current_item.streamdetails is not sd:
620 return
621 sd.stream_metadata = stream_metadata
622 sd.stream_metadata_last_updated = time.time()
623 self.mass.player_queues.signal_update(queue_id)
624
625 async def serve_queue_item_stream(self, request: web.Request) -> web.StreamResponse: # noqa: PLR0915
626 """Stream single queueitem audio to a player."""
627 self._log_request(request)
628 queue_id = request.match_info["queue_id"]
629 player_id = request.match_info["player_id"]
630 if not (queue := self.mass.player_queues.get(queue_id)):
631 raise web.HTTPNotFound(reason=f"Unknown Queue: {queue_id}")
632 session_id = request.match_info["session_id"]
633 pq_data = self.mass.player_queues.queue_data(queue.queue_id)
634 if pq_data.session_id is None or session_id != pq_data.session_id:
635 raise web.HTTPNotFound(reason=f"Unknown (or invalid) session: {session_id}")
636 if not (player := self.mass.players.get_player(player_id)):
637 raise web.HTTPNotFound(reason=f"Unknown Player: {player_id}")
638 queue_item_id = request.match_info["queue_item_id"]
639 queue_item = self.mass.player_queues.get_item(queue_id, queue_item_id)
640 if not queue_item:
641 raise web.HTTPNotFound(reason=f"Unknown Queue item: {queue_item_id}")
642
643 is_audio_source = (
644 queue_item.media_item is not None
645 and queue_item.media_item.media_type == MediaType.AUDIO_SOURCE
646 )
647
648 # HEAD probes for AudioSource items return a minimal response without
649 # touching the plugin. on_source_selected is the lifecycle hook that
650 # claims ownership and fires off transfer/handoff side effects (stop
651 # the previous player, redirect on disallowed switch, etc.), and a
652 # renderer probing with HEAD before GET should not trigger any of
653 # that. The actual GET request goes through the full hook chain.
654 if request.method != "GET" and is_audio_source:
655 # Validate the providing plugin still exists before advertising the
656 # source. Many DLNA renderers cache HEAD responses; returning 200
657 # for a URI whose plugin has been unloaded would lie to the
658 # renderer and the follow-up GET would fail unrecoverably.
659 assert queue_item.media_item is not None
660 if not isinstance(
661 self.mass.get_provider(queue_item.media_item.provider), PluginProvider
662 ):
663 raise web.HTTPNotFound(
664 reason=f"AudioSource provider {queue_item.media_item.provider} unavailable"
665 )
666 # For PCM-fmt URLs, advertise audio/wav in HEAD: most DLNA renderers
667 # key off the HEAD Content-Type to pick a decoder and do not handle
668 # raw PCM (application/octet-stream). The actual GET response will
669 # still wrap the bytes into a WAV container if needed via the same
670 # mime-type translation downstream.
671 head_fmt = request.match_info["fmt"]
672 if ContentType.try_parse(head_fmt).is_pcm():
673 head_fmt = ContentType.WAV.value
674 headers = {
675 **DEFAULT_STREAM_HEADERS,
676 "icy-name": sanitize_http_header_value(queue_item.name),
677 "contentFeatures.dlna.org": DLNA_CONTENT_FEATURES_REALTIME,
678 "Content-Type": get_mime_type(head_fmt),
679 }
680 resp = web.StreamResponse(status=200, reason="OK", headers=headers)
681 await resp.prepare(request)
682 return resp
683
684 # Fire on_source_selected hook for every AudioSource GET â this is the
685 # single point where exclusive plugin sources claim ownership. Firing
686 # unconditionally (regardless of whether streamdetails are cached from
687 # a previous request) means a disconnect/reconnect for the same queue
688 # item re-claims the lock with a fresh session id, instead of streaming
689 # against the stale ownership of the prior request.
690 # Source identity comes from queue_item.media_item because streamdetails
691 # may not exist yet on the first request.
692 # stream_session_id is a fresh per-request token threaded through to
693 # on_source_unselected so the provider can distinguish a stale
694 # teardown (e.g. a same-queue reconnect's first request completing
695 # AFTER its replacement has already started streaming) from the
696 # currently active session's real teardown.
697 audio_source_provider: PluginProvider | None = None
698 audio_source_id: str | None = None
699 stream_session_id = uuid4().hex
700 if (
701 is_audio_source
702 and queue_item.media_item is not None
703 and (prov := self.mass.get_provider(queue_item.media_item.provider))
704 and isinstance(prov, PluginProvider)
705 ):
706 audio_source_id = queue_item.media_item.item_id
707 # Wire the provider into the finally block BEFORE awaiting the
708 # hook: if the provider partially mutates state (claims the lock,
709 # records the session id) and then raises a non-RuntimeError
710 # exception (buggy plugin, asyncio.CancelledError, etc.), the
711 # finally must still fire on_source_unselected so the lock gets
712 # released. The provider's session-id guard makes a spurious
713 # release a no-op if the lock was never actually claimed.
714 audio_source_provider = prov
715 try:
716 await prov.on_source_selected(
717 audio_source_id, player_id, queue_id, stream_session_id
718 )
719 except RuntimeError as err:
720 # Provider intentionally aborts the original request (e.g.
721 # allow_player_switch=False has just redirected play_media to
722 # the configured target). Surface as 404 so the disallowed
723 # player drops the connection cleanly instead of treating an
724 # uncaught 500 as transient and retrying. The provider
725 # contract requires raising BEFORE claiming, so this is a
726 # clean abort â but we still let the finally run, where the
727 # session-id guard makes the unselect a no-op.
728 self.logger.info(
729 "AudioSource %s aborted stream for player %s: %s",
730 audio_source_id,
731 player_id,
732 err,
733 )
734 raise web.HTTPNotFound(reason=str(err))
735
736 try:
737 if not queue_item.streamdetails:
738 try:
739 queue_item.streamdetails = await self.audio.get_stream_details(
740 queue_item=queue_item
741 )
742 except Exception as e:
743 self.logger.error(
744 "Failed to get streamdetails for QueueItem %s: %s", queue_item_id, e
745 )
746 # a source capacity miss is transient, the item itself is fine
747 if not isinstance(e, ProviderStreamLimitError):
748 queue_item.available = False
749 raise web.HTTPNotFound(
750 reason=f"No streamdetails for Queue item: {queue_item_id}"
751 )
752
753 standard_crossfade_duration = self.mass.config.get_raw_core_config_value(
754 CONF_PLAYER_QUEUES, CONF_CROSSFADE_DURATION, 8
755 )
756 if queue_item.media_type != MediaType.TRACK:
757 crossfade_mode = CrossfadeMode.DISABLED
758 else:
759 crossfade_mode = self.get_crossfade_mode(queue)
760 if (
761 crossfade_mode != CrossfadeMode.DISABLED
762 and PlayerFeature.GAPLESS_PLAYBACK not in player.state.supported_features
763 ):
764 self.logger.warning(
765 "Crossfade disabled: Player %s does not support gapless playback, "
766 "consider enabling flow mode to enable crossfade on this player.",
767 player.state.name,
768 )
769 crossfade_mode = CrossfadeMode.DISABLED
770
771 # pick output format based on the streamdetails and player capabilities
772 pcm_format = await self.audio.select_pcm_format(
773 player=player,
774 streamdetails=queue_item.streamdetails,
775 crossfade_enabled=crossfade_mode != CrossfadeMode.DISABLED,
776 overlay_active=(queue_item.media_type == MediaType.RADIO and overlay_active(queue)),
777 )
778 output_format = await self.audio.get_output_format(
779 output_format_str=request.match_info["fmt"],
780 player=player,
781 content_sample_rate=pcm_format.sample_rate,
782 content_bit_depth=pcm_format.bit_depth,
783 media_type=queue_item.media_type,
784 )
785
786 # prepare request, add some DLNA/UPNP compatible headers
787 # icy-name is sanitized (all control chars, not just newlines) to avoid a
788 # "Potential header injection attack" ValueError by aiohttp
789 # see https://github.com/music-assistant/support/issues/4913
790 # and https://github.com/music-assistant/support/issues/5791
791 # use realtime DLNA flags for radio (sender-paced) since the source delivers slowly
792 dlna_features = (
793 DLNA_CONTENT_FEATURES_REALTIME
794 if queue_item.media_type != MediaType.TRACK
795 else DLNA_CONTENT_FEATURES
796 )
797 headers = {
798 **DEFAULT_STREAM_HEADERS,
799 "icy-name": sanitize_http_header_value(queue_item.name),
800 "contentFeatures.dlna.org": dlna_features,
801 "Content-Type": get_mime_type(output_format.output_format_str),
802 }
803
804 resp = web.StreamResponse(status=200, reason="OK", headers=headers)
805 resp.content_type = get_mime_type(output_format.output_format_str)
806 http_profile = player.get_config_value(CONF_HTTP_PROFILE, "default")
807 if http_profile == "forced_content_length" and not queue_item.duration:
808 # just set an insane high content length to make sure the player keeps playing
809 resp.content_length = calculate_content_length(output_format, 12 * 3600)
810 elif http_profile == "forced_content_length" and queue_item.duration:
811 # estimate content length based on effective duration
812 # account for seek position (e.g., crossfade from previous track)
813 seek_pos = queue_item.streamdetails.seek_position if queue_item.streamdetails else 0
814 effective_duration = max(queue_item.duration - seek_pos, 1)
815 # use cached actual bytes-per-second if available (from a previous stream)
816 resp.content_length = await get_content_length(
817 self.mass, queue_item.uri, output_format, effective_duration
818 )
819 elif http_profile == "chunked":
820 resp.enable_chunked_encoding()
821
822 await resp.prepare(request)
823
824 # return early if this is not a GET request
825 if request.method != "GET":
826 return resp
827
828 self._update_audio_processing_context(
829 queue=queue,
830 queue_item=queue_item,
831 pcm_format=pcm_format,
832 crossfade_mode=crossfade_mode,
833 overlay_enabled=(
834 queue_item.media_type == MediaType.RADIO and overlay_active(queue)
835 ),
836 session_id=session_id,
837 )
838
839 if crossfade_mode != CrossfadeMode.DISABLED:
840 # crossfade is enabled, use special crossfaded single item stream
841 # where the crossfade of the next track is present in the stream of
842 # a single track. This only works if the player supports gapless playback!
843 audio_input = self.audio.get_queue_item_stream_with_smartfade(
844 player=player,
845 queue_item=queue_item,
846 pcm_format=pcm_format,
847 crossfade_mode=crossfade_mode,
848 standard_crossfade_duration=standard_crossfade_duration,
849 session_id=session_id,
850 )
851 else:
852 # no crossfade, just a regular single item stream
853 audio_input = self.audio.get_queue_item_stream(
854 queue_item=queue_item,
855 pcm_format=pcm_format,
856 seek_position=int(queue_item.streamdetails.seek_position),
857 playback_speed=cast(
858 "float", queue_item.extra_attributes.get("playback_speed", 1.0)
859 ),
860 session_id=session_id,
861 )
862 if queue_item.media_type == MediaType.RADIO and overlay_active(queue):
863 # radio plays as a single long-lived stream (never in flow mode),
864 # so mix the audio overlay in here
865 audio_input = self.audio.get_overlay_mixed_stream(queue, audio_input, pcm_format)
866 # stream the audio
867 # this final ffmpeg process in the chain converts raw lossless PCM into
868 # the desired output format for the player including any player specific
869 # filter params such as channels mixing, DSP, resampling and, only if
870 # needed, encoding to lossy formats
871 output_plan = self.audio.get_player_output_plan(
872 player_id=player.player_id,
873 input_format=pcm_format,
874 output_format=output_format,
875 shared_player_ids=player.state.group_members,
876 queue_id=queue_id,
877 session_id=session_id,
878 queue_item_id=queue_item.queue_item_id,
879 )
880 filter_params = output_plan.filter_params
881 # Fast path for live AudioSource: when the player accepts WAV at the
882 # source's exact PCM rate/depth/channels and no filters apply, we
883 # skip the encode ffmpeg entirely and just stream a WAV header
884 # followed by the raw PCM bytes â saves an ffmpeg process and the
885 # latency of its internal buffer on every realtime stream.
886 audio_bytes: AsyncGenerator[bytes]
887 if (
888 queue_item.media_type == MediaType.AUDIO_SOURCE
889 and output_format.content_type == ContentType.WAV
890 and not filter_params
891 and output_format.sample_rate == pcm_format.sample_rate
892 and output_format.bit_depth == pcm_format.bit_depth
893 and output_format.channels == pcm_format.channels
894 ):
895 audio_bytes = _wav_passthrough_stream(audio_input, output_format)
896 else:
897 audio_bytes = get_ffmpeg_stream(
898 audio_input=audio_input,
899 input_format=pcm_format,
900 output_format=output_format,
901 filter_params=filter_params,
902 )
903 first_chunk_received = False
904 bytes_sent = 0
905 # Mark this player as actively streaming so audio analysis yields CPU to playback
906 # for the duration of the transfer (see audio_analysis.playback_active).
907 self._active_output_streams += 1
908 try:
909 # aclosing guarantees the generator (and thus the ffmpeg process chain
910 # behind it) is torn down immediately when the player disconnects
911 # mid-stream, instead of lingering until garbage collection finalizes
912 # the abandoned generator.
913 async with aclosing(audio_bytes):
914 async for chunk in audio_bytes:
915 try:
916 await resp.write(chunk)
917 bytes_sent += len(chunk)
918 if not first_chunk_received:
919 first_chunk_received = True
920 # inform the queue that the track is now loaded in the buffer
921 # so for example the next track can be enqueued
922 self.mass.player_queues.track_loaded_in_buffer(
923 queue_item.queue_id, queue_item.queue_item_id
924 )
925 except (BrokenPipeError, ConnectionResetError, ConnectionError) as err:
926 if (
927 first_chunk_received
928 and not player.stop_called
929 and queue_item.streamdetails.duration # ignore for radio streams
930 ):
931 # Player disconnected (unexpected) after receiving at least
932 # some data. This could indicate buffering issues, network
933 # problems, or player-specific issues.
934 self.logger.warning(
935 "Player %s disconnected prematurely from stream for %s (%s) - "
936 "error: %s, sent %d bytes, content_length=%s",
937 queue.display_name,
938 queue_item.name,
939 queue_item.uri,
940 err.__class__.__name__,
941 bytes_sent,
942 resp.content_length,
943 )
944 break
945 finally:
946 self._active_output_streams -= 1
947 if queue_item.streamdetails.stream_error:
948 self.logger.error(
949 "Error streaming QueueItem %s (%s) to %s",
950 queue_item.name,
951 queue_item.uri,
952 queue.display_name,
953 )
954 elif (
955 bytes_sent > 0
956 and queue_item.streamdetails
957 and queue_item.streamdetails.seconds_streamed
958 and queue_item.duration
959 ):
960 # cache the actual encoded bytes-per-second for this URI + output format
961 # so future content_length estimates are near-exact
962 self.mass.create_task(
963 store_content_length_in_cache(
964 self.mass,
965 queue_item.uri,
966 output_format,
967 bytes_sent,
968 queue_item.streamdetails.seconds_streamed,
969 )
970 )
971 return resp
972 finally:
973 # Paired with on_source_selected â fires regardless of how streaming
974 # ended (normal completion, client disconnect, exception). Lets
975 # NAMED_PIPE plugins release ownership without depending on an
976 # external session event. The stream_session_id is the same token
977 # passed to on_source_selected; the provider must reject the
978 # callback if it does not match the currently stored active
979 # session (otherwise a stale teardown from a superseded same-queue
980 # request would clear the live claim of its replacement).
981 if audio_source_provider is not None and audio_source_id is not None:
982 # Provider teardown failures must not break the response cycle
983 # (we're already in finally for a stream that ended one way or
984 # another), but they MUST surface in logs â otherwise a buggy
985 # plugin leaks _in_use_by_queue forever and there is no trail.
986 try:
987 await audio_source_provider.on_source_unselected(
988 audio_source_id, queue_id, stream_session_id
989 )
990 except Exception:
991 self.logger.warning(
992 "on_source_unselected raised for provider %s source %s queue %s",
993 audio_source_provider.instance_id,
994 audio_source_id,
995 queue_id,
996 exc_info=True,
997 )
998
999 async def serve_queue_flow_stream(self, request: web.Request) -> web.StreamResponse: # noqa: PLR0915
1000 """Stream Queue Flow audio to player."""
1001 self._log_request(request)
1002 queue_id = request.match_info["queue_id"]
1003 player_id = request.match_info["player_id"]
1004 if not (queue := self.mass.player_queues.get(queue_id)):
1005 raise web.HTTPNotFound(reason=f"Unknown Queue: {queue_id}")
1006 session_id = request.match_info["session_id"]
1007 queue_data = self.mass.player_queues.queue_data(queue_id)
1008 if queue_data.session_id is None or session_id != queue_data.session_id:
1009 raise web.HTTPNotFound(reason=f"Unknown (or invalid) session: {session_id}")
1010 if not (player := self.mass.players.get_player(player_id)):
1011 raise web.HTTPNotFound(reason=f"Unknown Player: {player_id}")
1012 start_queue_item_id = request.match_info["queue_item_id"]
1013 start_queue_item = self.mass.player_queues.get_item(queue_id, start_queue_item_id)
1014 if not start_queue_item:
1015 raise web.HTTPNotFound(reason=f"Unknown Queue item: {start_queue_item_id}")
1016
1017 # select the PCM format for the flow stream, anchored on the first track
1018 crossfade_mode = (
1019 self.get_crossfade_mode(queue)
1020 if start_queue_item.media_type == MediaType.TRACK
1021 else CrossfadeMode.DISABLED
1022 )
1023 flow_pcm_format = await self.audio.select_flow_pcm_format(
1024 player,
1025 start_streamdetails=start_queue_item.streamdetails,
1026 crossfade_enabled=crossfade_mode != CrossfadeMode.DISABLED,
1027 overlay_active=overlay_active(queue),
1028 )
1029
1030 # work out output format/details
1031 output_format = await self.audio.get_output_format(
1032 output_format_str=request.match_info["fmt"],
1033 player=player,
1034 content_sample_rate=flow_pcm_format.sample_rate,
1035 content_bit_depth=flow_pcm_format.bit_depth,
1036 media_type=start_queue_item.media_type,
1037 )
1038 # work out ICY metadata support
1039 icy_preference = self.mass.config.get_raw_player_config_value(
1040 player_id,
1041 CONF_ENTRY_ENABLE_ICY_METADATA.key,
1042 CONF_ENTRY_ENABLE_ICY_METADATA.default_value,
1043 )
1044 enable_icy = request.headers.get("Icy-MetaData", "") == "1" and icy_preference != "disabled"
1045 icy_meta_interval = 256000 if icy_preference == "full" else 16384
1046
1047 # prepare request, add some DLNA/UPNP compatible headers.
1048 # icy-name (in DEFAULT_STREAM_HEADERS) is always present so players have a
1049 # readable stream name; the rest of the ICY/shoutcast metadata headers are
1050 # only advertised when the client actually requested ICY metadata, rather
1051 # than on every flow response.
1052 headers = {
1053 **DEFAULT_STREAM_HEADERS,
1054 **(ICY_HEADERS if enable_icy else {}),
1055 "contentFeatures.dlna.org": DLNA_CONTENT_FEATURES_REALTIME,
1056 "Content-Type": get_mime_type(output_format.output_format_str),
1057 }
1058 if enable_icy:
1059 headers["icy-metaint"] = str(icy_meta_interval)
1060
1061 resp = web.StreamResponse(status=200, reason="OK", headers=headers)
1062 http_profile = player.get_config_value(CONF_HTTP_PROFILE, "default")
1063 if http_profile == "forced_content_length":
1064 # just set an insane high content length to make sure the player keeps playing
1065 resp.content_length = calculate_content_length(output_format, 12 * 3600)
1066 elif http_profile == "chunked":
1067 resp.enable_chunked_encoding()
1068
1069 await resp.prepare(request)
1070
1071 # return early if this is not a GET request
1072 if request.method != "GET":
1073 return resp
1074
1075 self._update_audio_processing_context(
1076 queue=queue,
1077 queue_item=start_queue_item,
1078 pcm_format=flow_pcm_format,
1079 crossfade_mode=crossfade_mode,
1080 overlay_enabled=overlay_active(queue),
1081 session_id=session_id,
1082 )
1083 output_plan = self.audio.get_player_output_plan(
1084 player.player_id,
1085 flow_pcm_format,
1086 output_format,
1087 shared_player_ids=player.state.group_members,
1088 queue_id=queue_id,
1089 session_id=session_id,
1090 )
1091
1092 # all checks passed, start streaming!
1093 # this final ffmpeg process in the chain will convert the raw, lossless PCM audio into
1094 # the desired output format for the player including any player specific filter params
1095 # such as channels mixing, DSP, resampling and, only if needed, encoding to lossy formats
1096 self.logger.debug("Start serving Queue flow audio stream for %s", queue.display_name)
1097
1098 # Mark this player as actively streaming so audio analysis yields CPU to playback
1099 # for the duration of the flow stream (see audio_analysis.playback_active).
1100 self._active_output_streams += 1
1101 flow_stream = self.audio.get_queue_flow_stream(
1102 queue=queue,
1103 start_queue_item=start_queue_item,
1104 pcm_format=flow_pcm_format,
1105 session_id=session_id,
1106 protocol_player=player,
1107 )
1108 if overlay_active(queue):
1109 flow_stream = self.audio.get_overlay_mixed_stream(queue, flow_stream, flow_pcm_format)
1110 audio_bytes = get_ffmpeg_stream(
1111 audio_input=flow_stream,
1112 input_format=flow_pcm_format,
1113 output_format=output_format,
1114 filter_params=output_plan.filter_params,
1115 # we need to slowly feed the music to avoid the player stopping and later
1116 # restarting (or completely failing) the audio stream by keeping the buffer short.
1117 # this is reported to be an issue especially with Chromecast players.
1118 # see for example: https://github.com/music-assistant/support/issues/3717
1119 # allow buffer ahead of a few seconds and read rest in (near) realtime
1120 extra_input_args=["-readrate", "1.05", "-readrate_initial_burst", "5"],
1121 chunk_size=icy_meta_interval if enable_icy else calculate_content_length(output_format),
1122 )
1123 client_disconnected = False
1124 try:
1125 # aclosing guarantees the flow stream (and thus the ffmpeg process chain
1126 # behind it) is torn down immediately when the player disconnects
1127 # mid-stream, instead of lingering until garbage collection finalizes
1128 # the abandoned generator.
1129 async with aclosing(audio_bytes):
1130 async for chunk in audio_bytes:
1131 try:
1132 await resp.write(chunk)
1133 except BrokenPipeError, ConnectionResetError, ConnectionError:
1134 # race condition
1135 client_disconnected = True
1136 break
1137
1138 if not enable_icy:
1139 continue
1140
1141 # if icy metadata is enabled, send the icy metadata after the chunk
1142 if (
1143 # use current item here and not buffered item, otherwise
1144 # the icy metadata will be too much ahead
1145 (current_item := queue.current_item)
1146 and current_item.streamdetails
1147 and current_item.streamdetails.stream_title
1148 ):
1149 title = current_item.streamdetails.stream_title
1150 elif queue and current_item and current_item.name:
1151 title = current_item.name
1152 else:
1153 title = "Music Assistant"
1154 metadata = f"StreamTitle='{title}';".encode()
1155 if icy_preference == "full" and current_item and current_item.image:
1156 metadata += f"StreamURL='{current_item.image.path}'".encode()
1157 while len(metadata) % 16 != 0:
1158 metadata += b"\x00"
1159 length = len(metadata)
1160 length_b = chr(int(length / 16)).encode()
1161 await resp.write(length_b + metadata)
1162 finally:
1163 self._active_output_streams -= 1
1164
1165 if not client_disconnected and http_profile == "forced_content_length":
1166 await self._finish_flow_stream(resp, queue_id, session_id)
1167
1168 return resp
1169
1170 async def serve_command_request(self, request: web.Request) -> web.FileResponse:
1171 """Handle special 'command' request for a player."""
1172 self._log_request(request)
1173 queue_id = request.match_info["queue_id"]
1174 command = request.match_info["command"]
1175 if command == "next":
1176 self.mass.create_task(self.mass.player_queues.next(queue_id))
1177 return web.FileResponse(SILENCE_FILE, headers={"icy-name": "Music Assistant"})
1178
1179 async def serve_announcement_stream(self, request: web.Request) -> web.StreamResponse:
1180 """Stream announcement audio to a player."""
1181 self._log_request(request)
1182 player_id = request.match_info["player_id"]
1183 if not (player := self.mass.players.get_player(player_id)):
1184 raise web.HTTPNotFound(reason=f"Unknown Player: {player_id}")
1185 if not (announce_data := self.announcement_renderer.get_for_player(player_id)):
1186 raise web.HTTPNotFound(reason=f"No pending announcements for Player: {player_id}")
1187
1188 # work out output format/details
1189 fmt = request.match_info["fmt"]
1190 audio_format = AudioFormat(content_type=ContentType.try_parse(fmt))
1191
1192 http_profile = self._get_announcement_http_profile(player_id, announce_data)
1193
1194 # return early if this is not a GET request:
1195 # players often probe the url with a HEAD request before fetching it and
1196 # rendering the announcement for such a probe would run the entire (costly)
1197 # TTS/ffmpeg chain twice for a single announcement.
1198 if request.method != "GET":
1199 resp = web.StreamResponse(status=200, reason="OK", headers=DEFAULT_STREAM_HEADERS)
1200 resp.content_type = get_mime_type(audio_format.output_format_str)
1201 if http_profile == "chunked":
1202 resp.enable_chunked_encoding()
1203 await resp.prepare(request)
1204 return resp
1205
1206 if http_profile == "forced_content_length":
1207 # given the fact that an announcement is just a short audio clip,
1208 # just send it over completely at once so we have a fixed content length
1209 data = bytearray()
1210 announcement_stream = self.get_announcement_stream(announce_data, audio_format)
1211 # aclosing guarantees the stream (and thus the ffmpeg process chain behind
1212 # it) is torn down immediately when the request is cancelled, instead of
1213 # lingering until garbage collection finalizes the abandoned generator.
1214 async with aclosing(announcement_stream):
1215 async for chunk in announcement_stream:
1216 data += chunk
1217 return web.Response(
1218 body=bytes(data),
1219 content_type=get_mime_type(audio_format.output_format_str),
1220 headers=DEFAULT_STREAM_HEADERS,
1221 )
1222
1223 resp = web.StreamResponse(status=200, reason="OK", headers=DEFAULT_STREAM_HEADERS)
1224 resp.content_type = get_mime_type(audio_format.output_format_str)
1225 if http_profile == "chunked":
1226 resp.enable_chunked_encoding()
1227
1228 await resp.prepare(request)
1229
1230 # all checks passed, start streaming!
1231 self.logger.debug(
1232 "Start serving audio stream for Announcement %s to %s",
1233 announce_data["announcement_url"],
1234 player.display_name,
1235 )
1236 announcement_stream = self.get_announcement_stream(announce_data, audio_format)
1237 # aclosing guarantees the stream (and thus the ffmpeg process chain behind
1238 # it) is torn down immediately when the player disconnects mid-stream,
1239 # instead of lingering until garbage collection finalizes the abandoned
1240 # generator.
1241 async with aclosing(announcement_stream):
1242 async for chunk in announcement_stream:
1243 try:
1244 await resp.write(chunk)
1245 except BrokenPipeError, ConnectionResetError:
1246 break
1247
1248 self.logger.debug(
1249 "Finished serving audio stream for Announcement %s to %s",
1250 announce_data["announcement_url"],
1251 player.display_name,
1252 )
1253
1254 return resp
1255
1256 def get_command_url(self, player_or_queue_id: str, command: str) -> str:
1257 """Get the url for the special command stream."""
1258 return f"{self.base_url}/command/{player_or_queue_id}/{command}.mp3"
1259
1260 def get_announcement_url(
1261 self,
1262 player_id: str,
1263 content_type: ContentType = ContentType.MP3,
1264 ) -> str:
1265 """
1266 Get the url that serves the announcement registered for the given player.
1267
1268 :param player_id: The player the announcement is played on.
1269 :param content_type: The format to serve the announcement in.
1270 """
1271 # use stream server to host announcement on local network
1272 # this ensures playback on all players, including ones that do not
1273 # like https hosts and it also offers the pre-announce 'bell'
1274 return f"{self.base_url}/announcement/{player_id}.{content_type.value}"
1275
1276 def get_stream(
1277 self,
1278 media: PlayerMedia,
1279 pcm_format: AudioFormat,
1280 player_id: str | None = None,
1281 force_flow_mode: bool = False,
1282 use_flow_stream_buffering: bool = False,
1283 ) -> AsyncGenerator[bytes]:
1284 """
1285 Get a stream of the given media as raw PCM audio.
1286
1287 This is used as helper for player providers that can consume the raw PCM
1288 audio stream directly (e.g. AirPlay) and not rely on HTTP transport.
1289
1290 :param media: The PlayerMedia to stream.
1291 :param pcm_format: The desired output PCM format.
1292 :param player_id: The player ID requesting the stream. Used to determine
1293 if flow mode should be used based on the player's capabilities.
1294 :param force_flow_mode: Force flow mode regardless of player capabilities.
1295 Used for multi-client streaming scenarios that require continuous streams.
1296 :param use_flow_stream_buffering: Buffer the flow stream to provide headroom
1297 during smart fades transitions. Use for consumers that read directly
1298 (e.g. AirPlay, Snapcast) and can't tolerate stalls.
1299 """
1300 # select audio source
1301 if media.media_type == MediaType.ANNOUNCEMENT:
1302 # special case: stream announcement
1303 assert media.custom_data
1304 return self.get_announcement_stream(cast("AnnounceData", media.custom_data), pcm_format)
1305 if (
1306 media.source_id
1307 and media.source_id.startswith(UGP_PREFIX)
1308 and media.uri
1309 and "/ugp/" in media.uri
1310 ):
1311 # special case: member player accessing UGP stream
1312 # Check URI to distinguish from the UGP accessing its own stream
1313 ugp_player = cast("UniversalGroupPlayer", self.mass.players.get_player(media.source_id))
1314 ugp_stream = ugp_player.stream
1315 assert ugp_stream is not None # for type checker
1316 if ugp_stream.base_pcm_format == pcm_format:
1317 # no conversion needed
1318 return ugp_stream.subscribe_raw()
1319 return ugp_stream.get_stream(output_format=pcm_format)
1320 if media.source_id and media.queue_item_id:
1321 # Queue stream request - determine flow_mode based on player capabilities
1322 # or force it if explicitly requested (e.g., for multi-client streaming)
1323 protocol_player = self.mass.players.get_player(player_id) if player_id else None
1324 queue_id = media.source_id
1325 queue = self.mass.player_queues.get(queue_id)
1326 queue_session_id = cast(
1327 "str | None",
1328 (media.custom_data or {}).get("session_id"),
1329 )
1330 crossfade_needs_flow_mode = (
1331 # crossfade only applies to tracks; if the queue has it enabled but the
1332 # player(protocol) does not support gapless playback, we need to enforce flow mode
1333 media.media_type == MediaType.TRACK
1334 and queue is not None
1335 and queue.crossfade_enabled
1336 and protocol_player
1337 and not protocol_player.supports_gapless
1338 )
1339 # the audio overlay is mixed into the queue's continuous (flow) stream;
1340 # per-item requests would restart the overlay at every track boundary
1341 overlay_needs_flow_mode = queue is not None and overlay_active(queue)
1342 flow_mode = (
1343 force_flow_mode
1344 or (protocol_player is not None and protocol_player.flow_mode)
1345 or crossfade_needs_flow_mode
1346 or overlay_needs_flow_mode
1347 )
1348 if media.media_type in (MediaType.RADIO, MediaType.AUDIO_SOURCE):
1349 # flow_mode for live/infinite streams is pointless
1350 flow_mode = False
1351 if flow_mode:
1352 # flow stream request
1353 assert queue
1354 start_queue_item = self.mass.player_queues.get_item(
1355 media.source_id, media.queue_item_id
1356 )
1357 assert start_queue_item
1358 crossfade_mode = (
1359 self.get_crossfade_mode(queue)
1360 if start_queue_item.media_type == MediaType.TRACK
1361 else CrossfadeMode.DISABLED
1362 )
1363 self._update_audio_processing_context(
1364 queue=queue,
1365 queue_item=start_queue_item,
1366 pcm_format=pcm_format,
1367 crossfade_mode=crossfade_mode,
1368 overlay_enabled=overlay_active(queue),
1369 session_id=queue_session_id,
1370 )
1371 flow_stream = self.audio.get_queue_flow_stream(
1372 queue=queue,
1373 start_queue_item=start_queue_item,
1374 pcm_format=pcm_format,
1375 session_id=queue_session_id,
1376 protocol_player=protocol_player,
1377 )
1378 if overlay_active(queue):
1379 flow_stream = self.audio.get_overlay_mixed_stream(
1380 queue, flow_stream, pcm_format
1381 )
1382 if use_flow_stream_buffering:
1383 return buffered(flow_stream, buffer_size=30, min_buffer_before_yield=1)
1384 return flow_stream
1385 # single item stream (e.g. radio or non-flow mode)
1386 queue_item = self.mass.player_queues.get_item(media.source_id, media.queue_item_id)
1387 assert queue_item
1388 if queue is not None:
1389 self._update_audio_processing_context(
1390 queue=queue,
1391 queue_item=queue_item,
1392 pcm_format=pcm_format,
1393 crossfade_mode=(
1394 self.get_crossfade_mode(queue)
1395 if queue_item.media_type == MediaType.TRACK
1396 else CrossfadeMode.DISABLED
1397 ),
1398 overlay_enabled=(
1399 queue_item.media_type == MediaType.RADIO and overlay_active(queue)
1400 ),
1401 session_id=queue_session_id,
1402 )
1403 inner_stream = self.audio.get_queue_item_stream(
1404 queue_item=queue_item,
1405 pcm_format=pcm_format,
1406 seek_position=(
1407 int(queue_item.streamdetails.seek_position) if queue_item.streamdetails else 0
1408 ),
1409 playback_speed=cast(
1410 "float", queue_item.extra_attributes.get("playback_speed", 1.0)
1411 ),
1412 session_id=queue_session_id,
1413 )
1414 if (
1415 queue is not None
1416 and queue_item.media_type == MediaType.RADIO
1417 and overlay_active(queue)
1418 ):
1419 # radio plays as a single long-lived stream, so mix the overlay in here
1420 inner_stream = self.audio.get_overlay_mixed_stream(queue, inner_stream, pcm_format)
1421 # mirror the on_source_selected/unselected lifecycle the HTTP route
1422 # fires, so direct-PCM consumers (AirPlay, Snapcast, UGP) honour the
1423 # plugin contract too
1424 if (
1425 queue_item.media_item is not None
1426 and queue_item.media_item.media_type == MediaType.AUDIO_SOURCE
1427 ):
1428 return self._wrap_with_audio_source_lifecycle(
1429 inner=inner_stream,
1430 queue_item=queue_item,
1431 player_id=player_id or media.source_id,
1432 )
1433 return inner_stream
1434 # assume url or some other direct path
1435 # NOTE: this will fail if its an uri not playable by ffmpeg
1436 return get_ffmpeg_stream(
1437 audio_input=media.uri,
1438 input_format=AudioFormat(content_type=ContentType.try_parse(media.uri)),
1439 output_format=pcm_format,
1440 )
1441
1442 async def get_preview_stream(
1443 self,
1444 provider_instance_id_or_domain: str,
1445 item_id: str,
1446 media_type: MediaType = MediaType.TRACK,
1447 ) -> AsyncGenerator[bytes]:
1448 """Create a 30 seconds preview audioclip for the given media item."""
1449 if not (music_prov := self.mass.get_provider(provider_instance_id_or_domain)):
1450 raise ProviderUnavailableError
1451 if music_prov.type != ProviderType.MUSIC:
1452 msg = f"{provider_instance_id_or_domain} is not a music provider"
1453 raise InvalidDataError(msg)
1454 music_prov = cast("MusicProvider", music_prov)
1455
1456 try:
1457 await self.mass.music.get_item(
1458 media_type,
1459 item_id,
1460 provider_instance_id_or_domain,
1461 allow_update_metadata=False,
1462 )
1463 except MediaNotFoundError as err:
1464 msg = f"Item {item_id} not found in provider {provider_instance_id_or_domain}"
1465 raise InvalidDataError(msg) from err
1466
1467 streamdetails = await music_prov.get_stream_details(item_id, media_type)
1468 pcm_format = AudioFormat(
1469 content_type=ContentType.from_bit_depth(streamdetails.audio_format.bit_depth),
1470 sample_rate=streamdetails.audio_format.sample_rate,
1471 bit_depth=streamdetails.audio_format.bit_depth,
1472 channels=streamdetails.audio_format.channels,
1473 )
1474 async for chunk in get_ffmpeg_stream(
1475 audio_input=self.audio.get_media_stream(
1476 streamdetails=streamdetails, pcm_format=pcm_format
1477 ),
1478 input_format=pcm_format,
1479 output_format=AudioFormat(content_type=ContentType.AAC),
1480 extra_input_args=["-t", "30"],
1481 ):
1482 yield chunk
1483
1484 async def get_announcement_stream(
1485 self, announce_data: AnnounceData, output_format: AudioFormat
1486 ) -> AsyncGenerator[bytes]:
1487 """
1488 Get the audio of an announcement (pre-announce chime + announcement).
1489
1490 Any number of consumers may stream the same announcement at once; its source is
1491 fetched and decoded only once. The audio stays available while the stream is
1492 held open.
1493
1494 :param announce_data: The announcement to stream.
1495 :param output_format: The format to deliver the audio in.
1496 """
1497 render = self.announcement_renderer.acquire(announce_data)
1498 try:
1499 # aclosing guarantees this consumer's ffmpeg encoder is torn down
1500 # immediately when it goes away, instead of lingering until garbage
1501 # collection finalizes the abandoned generator.
1502 stream = render.get_stream(output_format)
1503 async with aclosing(stream):
1504 async for chunk in stream:
1505 yield chunk
1506 finally:
1507 await self.announcement_renderer.release(render)
1508
1509 async def get_announcement_duration(
1510 self, announcement: PlayerMedia, timeout: float = DEFAULT_RENDER_TIMEOUT
1511 ) -> int | None:
1512 """
1513 Get the exact duration (in seconds) of an announcement, once it finished rendering.
1514
1515 Waits for the audio to be rendered in full, so call this while the announcement
1516 plays rather than before handing it to a player. Returns None when the length can
1517 not be determined, e.g. the announcement is no longer playing or its source did
1518 not deliver in time.
1519
1520 :param announcement: The announcement to return the duration for.
1521 :param timeout: Maximum time to wait for the audio to finish rendering.
1522 """
1523 if announcement.duration:
1524 return announcement.duration
1525 if not announcement.custom_data:
1526 return None
1527 render = self.announcement_renderer.get(cast("AnnounceData", announcement.custom_data))
1528 if render is None:
1529 return None
1530 duration = await render.wait_finished(timeout)
1531 return ceil(duration) if duration else None
1532
1533 async def _wrap_with_audio_source_lifecycle(
1534 self,
1535 inner: AsyncGenerator[bytes],
1536 queue_item: QueueItem,
1537 player_id: str,
1538 ) -> AsyncGenerator[bytes]:
1539 """
1540 Wrap an AudioSource queue item stream with on_source_selected/unselected hooks.
1541
1542 Direct-PCM consumers (AirPlay, Snapcast, UGP, ...) call ``get_stream`` instead
1543 of going through the HTTP route, but the plugin contract requires the
1544 lifecycle hooks to fire for every actual stream request â they're what
1545 claim/release the per-queue exclusive ownership and trigger acquisition
1546 side effects like the Spotify Connect Web API play kick. This wrapper
1547 gives those consumers the same lifecycle the HTTP route already provides.
1548
1549 :param inner: The underlying audio stream generator.
1550 :param queue_item: The AudioSource queue item being streamed.
1551 :param player_id: The protocol player consuming this stream.
1552 """
1553 media_item = queue_item.media_item
1554 assert media_item is not None # caller checked media_type == AUDIO_SOURCE
1555 prov = self.mass.get_provider(media_item.provider)
1556 queue_id = queue_item.queue_id
1557 if not isinstance(prov, PluginProvider):
1558 async for chunk in inner:
1559 yield chunk
1560 return
1561 source_id = media_item.item_id
1562 stream_session_id = uuid4().hex
1563 # single try/finally so on_source_unselected fires even when
1564 # on_source_selected raises after partially claiming state; the
1565 # provider's session_id guard makes a no-op claim release safe.
1566 try:
1567 try:
1568 await prov.on_source_selected(source_id, player_id, queue_id, stream_session_id)
1569 except RuntimeError as err:
1570 # provider intentionally aborts the request â surface as AudioError
1571 raise AudioError(str(err)) from err
1572 async for chunk in inner:
1573 yield chunk
1574 finally:
1575 try:
1576 await prov.on_source_unselected(source_id, queue_id, stream_session_id)
1577 except Exception as err:
1578 self.logger.exception(
1579 "on_source_unselected raised for provider %s source %s queue %s: %s",
1580 prov.instance_id,
1581 source_id,
1582 queue_id,
1583 err,
1584 )
1585
1586 def _update_audio_processing_context(
1587 self,
1588 queue: PlayerQueue,
1589 queue_item: QueueItem,
1590 pcm_format: AudioFormat,
1591 crossfade_mode: CrossfadeMode,
1592 overlay_enabled: bool,
1593 session_id: str | None = None,
1594 ) -> None:
1595 """
1596 Store the shared processing context selected for a queue item.
1597
1598 :param queue: Active player queue.
1599 :param queue_item: Queue item being prepared.
1600 :param pcm_format: Shared PCM format leaving queue processing.
1601 :param crossfade_mode: Effective crossfade mode for the item.
1602 :param overlay_enabled: Whether an overlay is mixed into this stream.
1603 :param session_id: Queue session that owns processing-detail updates.
1604 """
1605 if queue_item.streamdetails is None:
1606 return
1607 queue_data = self.mass.player_queues.queue_data_or_none(queue.queue_id)
1608 if (
1609 queue_data is None
1610 or (processing_session_id := session_id or queue_data.session_id) is None
1611 or queue_data.session_id != processing_session_id
1612 ):
1613 return
1614 self.audio_processing.start_session(queue.queue_id, processing_session_id)
1615 self.audio_processing.update_item_context(
1616 queue_id=queue.queue_id,
1617 session_id=processing_session_id,
1618 queue_item_id=queue_item.queue_item_id,
1619 queue_processing=AudioQueueProcessing(
1620 pcm_format=pcm_format,
1621 playback_speed=cast(
1622 "float",
1623 queue_item.extra_attributes.get("playback_speed", 1.0),
1624 ),
1625 crossfade_mode=crossfade_mode,
1626 overlay_active=overlay_enabled,
1627 ),
1628 alters_audio=queue_item.streamdetails.fade_in,
1629 )
1630
1631 def _get_announcement_http_profile(self, player_id: str, announce_data: AnnounceData) -> str:
1632 """
1633 Resolve the http profile for serving an announcement stream.
1634
1635 Announcement urls are registered under the visible player's id, but the
1636 stream may be fetched by a linked protocol player; the profile must come
1637 from the player that actually performs the fetch.
1638 """
1639 announce_player = None
1640 if announce_player_id := announce_data.get("announce_player_id"):
1641 announce_player = self.mass.players.get_player(announce_player_id)
1642 if announce_player is None:
1643 announce_player = self.mass.players.get_player(player_id)
1644 if announce_player is None:
1645 return "default"
1646 return announce_player.get_output_config_value(CONF_HTTP_PROFILE, "default")
1647
1648 async def _finish_flow_stream(
1649 self, resp: web.StreamResponse, queue_id: str, session_id: str
1650 ) -> None:
1651 """
1652 Close a fully served flow stream, giving the player time to drain when it ends the queue.
1653
1654 :param resp: The flow stream response, already fully written.
1655 :param queue_id: Id of the queue the flow stream belongs to.
1656 :param session_id: Stream session this response was opened for.
1657 """
1658 if self.mass.player_queues.flow_queue_exhausted(queue_id, session_id):
1659 # the player is still holding a few seconds of audio it has not rendered yet
1660 # and drops that as soon as the stream ends, so let it play out first.
1661 # a flow that ends to be restarted right away gets no such grace: there the
1662 # player should go idle as soon as possible so the next stream can start.
1663 self.logger.debug(
1664 "Flow stream for queue %s reached the end of the queue - holding the "
1665 "connection open for %ss so the player can play out its buffer",
1666 queue_id,
1667 FLOW_STREAM_LEAD_OUT_SECONDS,
1668 )
1669 await asyncio.sleep(FLOW_STREAM_LEAD_OUT_SECONDS)
1670 # aiohttp derives keep-alive from the request, so the 'Connection: close' we
1671 # advertise is relayed to the player but never applied to the response itself.
1672 # Without this the player is left waiting on a stream that already ended.
1673 resp.force_close()
1674
1675 def _log_request(self, request: web.Request) -> None:
1676 """Log request."""
1677 if self.logger.isEnabledFor(VERBOSE_LOG_LEVEL):
1678 self.logger.log(
1679 VERBOSE_LOG_LEVEL,
1680 "Got %s request to %s from %s\nheaders: %s\n",
1681 request.method,
1682 request.path,
1683 request.remote,
1684 redact_sensitive_headers(request.headers),
1685 )
1686 else:
1687 self.logger.debug(
1688 "Got %s request to %s from %s (HTTP/%s.%s, connection: %s)",
1689 request.method,
1690 request.path,
1691 request.remote,
1692 request.version.major,
1693 request.version.minor,
1694 request.headers.get("Connection", "-"),
1695 )
1696
1697 async def _reload_network_dependent_providers(self) -> None:
1698 """Reload the providers that captured the streamserver network, if it changed."""
1699 previous = self._network_fingerprint
1700 current = (
1701 self._bind_ip,
1702 str(self.publish_ip),
1703 cast("int", self.publish_port),
1704 tuple(self._publish_addresses),
1705 )
1706 if previous is None or previous == current:
1707 self._network_fingerprint = current
1708 return
1709 # these providers bind or advertise the network while they load, so a plain
1710 # reload is what moves them over - they share no lighter rebind path
1711 instance_ids = [
1712 prov.instance_id
1713 for prov in self.mass.providers
1714 if prov.reload_on_streams_network_change
1715 ]
1716 for instance_id in instance_ids:
1717 try:
1718 config = await self.mass.config.get_provider_config(instance_id)
1719 self.logger.info(
1720 "Streamserver network changed, reloading provider %s",
1721 config.name or config.domain,
1722 )
1723 await self.mass.load_provider_config(config)
1724 except Exception as err:
1725 self.logger.warning(
1726 "Error reloading provider %s: %s",
1727 instance_id,
1728 str(err) or err.__class__.__name__,
1729 exc_info=err,
1730 )
1731 # only mark the new network as applied once the loop completed, so a run cut short
1732 # by a second config change runs again on the next reload
1733 self._network_fingerprint = current
1734
1735 def _setup_smart_fades_logger(self, config: CoreConfig) -> None:
1736 """Set up smart fades logger level."""
1737 log_level = str(config.get_value(CONF_SMART_FADES_LOG_LEVEL))
1738 if log_level == "GLOBAL":
1739 self.audio.smart_fades_mixer.logger.setLevel(self.logger.level)
1740 else:
1741 self.audio.smart_fades_mixer.logger.setLevel(log_level)
1742
1743 def _resolve_publish_state(self, bind_ip: str, publish_candidates: tuple[str, ...]) -> None:
1744 """
1745 Resolve the addresses and base URL to advertise for the given bind address.
1746
1747 Reads ``self.publish_port``, so set that first.
1748
1749 :param bind_ip: Address the streamserver binds to (a wildcard means all interfaces).
1750 :param publish_candidates: Host addresses reachable from the local network, ranked.
1751 """
1752 self._bind_ip = bind_ip
1753 self._publish_addresses = _get_publish_addresses(
1754 bind_ip, self._configured_publish_ip, publish_candidates
1755 )
1756 # the single address players are handed, taken from the top of the ranked list
1757 self.publish_ip = self._publish_addresses[0]
1758 self._base_url = f"http://{format_ip_for_url(self.publish_ip)}:{self.publish_port}"
1759
1760
1761def _same_ip_family(ip: str, other_ip: str) -> bool:
1762 """Return whether two addresses belong to the same IP family."""
1763 return (":" in ip) == (":" in other_ip)
1764