/
/
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
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 queue_item.available = False
747 raise web.HTTPNotFound(
748 reason=f"No streamdetails for Queue item: {queue_item_id}"
749 )
750
751 standard_crossfade_duration = self.mass.config.get_raw_core_config_value(
752 CONF_PLAYER_QUEUES, CONF_CROSSFADE_DURATION, 8
753 )
754 if queue_item.media_type != MediaType.TRACK:
755 crossfade_mode = CrossfadeMode.DISABLED
756 else:
757 crossfade_mode = self.get_crossfade_mode(queue)
758 if (
759 crossfade_mode != CrossfadeMode.DISABLED
760 and PlayerFeature.GAPLESS_PLAYBACK not in player.state.supported_features
761 ):
762 self.logger.warning(
763 "Crossfade disabled: Player %s does not support gapless playback, "
764 "consider enabling flow mode to enable crossfade on this player.",
765 player.state.name,
766 )
767 crossfade_mode = CrossfadeMode.DISABLED
768
769 # pick output format based on the streamdetails and player capabilities
770 pcm_format = await self.audio.select_pcm_format(
771 player=player,
772 streamdetails=queue_item.streamdetails,
773 crossfade_enabled=crossfade_mode != CrossfadeMode.DISABLED,
774 overlay_active=(queue_item.media_type == MediaType.RADIO and overlay_active(queue)),
775 )
776 output_format = await self.audio.get_output_format(
777 output_format_str=request.match_info["fmt"],
778 player=player,
779 content_sample_rate=pcm_format.sample_rate,
780 content_bit_depth=pcm_format.bit_depth,
781 media_type=queue_item.media_type,
782 )
783
784 # prepare request, add some DLNA/UPNP compatible headers
785 # icy-name is sanitized (all control chars, not just newlines) to avoid a
786 # "Potential header injection attack" ValueError by aiohttp
787 # see https://github.com/music-assistant/support/issues/4913
788 # and https://github.com/music-assistant/support/issues/5791
789 # use realtime DLNA flags for radio (sender-paced) since the source delivers slowly
790 dlna_features = (
791 DLNA_CONTENT_FEATURES_REALTIME
792 if queue_item.media_type != MediaType.TRACK
793 else DLNA_CONTENT_FEATURES
794 )
795 headers = {
796 **DEFAULT_STREAM_HEADERS,
797 "icy-name": sanitize_http_header_value(queue_item.name),
798 "contentFeatures.dlna.org": dlna_features,
799 "Content-Type": get_mime_type(output_format.output_format_str),
800 }
801
802 resp = web.StreamResponse(status=200, reason="OK", headers=headers)
803 resp.content_type = get_mime_type(output_format.output_format_str)
804 http_profile = player.get_config_value(CONF_HTTP_PROFILE, "default")
805 if http_profile == "forced_content_length" and not queue_item.duration:
806 # just set an insane high content length to make sure the player keeps playing
807 resp.content_length = calculate_content_length(output_format, 12 * 3600)
808 elif http_profile == "forced_content_length" and queue_item.duration:
809 # estimate content length based on effective duration
810 # account for seek position (e.g., crossfade from previous track)
811 seek_pos = queue_item.streamdetails.seek_position if queue_item.streamdetails else 0
812 effective_duration = max(queue_item.duration - seek_pos, 1)
813 # use cached actual bytes-per-second if available (from a previous stream)
814 resp.content_length = await get_content_length(
815 self.mass, queue_item.uri, output_format, effective_duration
816 )
817 elif http_profile == "chunked":
818 resp.enable_chunked_encoding()
819
820 await resp.prepare(request)
821
822 # return early if this is not a GET request
823 if request.method != "GET":
824 return resp
825
826 self._update_audio_processing_context(
827 queue=queue,
828 queue_item=queue_item,
829 pcm_format=pcm_format,
830 crossfade_mode=crossfade_mode,
831 overlay_enabled=(
832 queue_item.media_type == MediaType.RADIO and overlay_active(queue)
833 ),
834 session_id=session_id,
835 )
836
837 if crossfade_mode != CrossfadeMode.DISABLED:
838 # crossfade is enabled, use special crossfaded single item stream
839 # where the crossfade of the next track is present in the stream of
840 # a single track. This only works if the player supports gapless playback!
841 audio_input = self.audio.get_queue_item_stream_with_smartfade(
842 player=player,
843 queue_item=queue_item,
844 pcm_format=pcm_format,
845 crossfade_mode=crossfade_mode,
846 standard_crossfade_duration=standard_crossfade_duration,
847 session_id=session_id,
848 )
849 else:
850 # no crossfade, just a regular single item stream
851 audio_input = self.audio.get_queue_item_stream(
852 queue_item=queue_item,
853 pcm_format=pcm_format,
854 seek_position=int(queue_item.streamdetails.seek_position),
855 playback_speed=cast(
856 "float", queue_item.extra_attributes.get("playback_speed", 1.0)
857 ),
858 session_id=session_id,
859 )
860 if queue_item.media_type == MediaType.RADIO and overlay_active(queue):
861 # radio plays as a single long-lived stream (never in flow mode),
862 # so mix the audio overlay in here
863 audio_input = self.audio.get_overlay_mixed_stream(queue, audio_input, pcm_format)
864 # stream the audio
865 # this final ffmpeg process in the chain converts raw lossless PCM into
866 # the desired output format for the player including any player specific
867 # filter params such as channels mixing, DSP, resampling and, only if
868 # needed, encoding to lossy formats
869 output_plan = self.audio.get_player_output_plan(
870 player_id=player.player_id,
871 input_format=pcm_format,
872 output_format=output_format,
873 shared_player_ids=player.state.group_members,
874 queue_id=queue_id,
875 session_id=session_id,
876 queue_item_id=queue_item.queue_item_id,
877 )
878 filter_params = output_plan.filter_params
879 # Fast path for live AudioSource: when the player accepts WAV at the
880 # source's exact PCM rate/depth/channels and no filters apply, we
881 # skip the encode ffmpeg entirely and just stream a WAV header
882 # followed by the raw PCM bytes â saves an ffmpeg process and the
883 # latency of its internal buffer on every realtime stream.
884 audio_bytes: AsyncGenerator[bytes]
885 if (
886 queue_item.media_type == MediaType.AUDIO_SOURCE
887 and output_format.content_type == ContentType.WAV
888 and not filter_params
889 and output_format.sample_rate == pcm_format.sample_rate
890 and output_format.bit_depth == pcm_format.bit_depth
891 and output_format.channels == pcm_format.channels
892 ):
893 audio_bytes = _wav_passthrough_stream(audio_input, output_format)
894 else:
895 audio_bytes = get_ffmpeg_stream(
896 audio_input=audio_input,
897 input_format=pcm_format,
898 output_format=output_format,
899 filter_params=filter_params,
900 )
901 first_chunk_received = False
902 bytes_sent = 0
903 # Mark this player as actively streaming so audio analysis yields CPU to playback
904 # for the duration of the transfer (see audio_analysis.playback_active).
905 self._active_output_streams += 1
906 try:
907 # aclosing guarantees the generator (and thus the ffmpeg process chain
908 # behind it) is torn down immediately when the player disconnects
909 # mid-stream, instead of lingering until garbage collection finalizes
910 # the abandoned generator.
911 async with aclosing(audio_bytes):
912 async for chunk in audio_bytes:
913 try:
914 await resp.write(chunk)
915 bytes_sent += len(chunk)
916 if not first_chunk_received:
917 first_chunk_received = True
918 # inform the queue that the track is now loaded in the buffer
919 # so for example the next track can be enqueued
920 self.mass.player_queues.track_loaded_in_buffer(
921 queue_item.queue_id, queue_item.queue_item_id
922 )
923 except (BrokenPipeError, ConnectionResetError, ConnectionError) as err:
924 if (
925 first_chunk_received
926 and not player.stop_called
927 and queue_item.streamdetails.duration # ignore for radio streams
928 ):
929 # Player disconnected (unexpected) after receiving at least
930 # some data. This could indicate buffering issues, network
931 # problems, or player-specific issues.
932 self.logger.warning(
933 "Player %s disconnected prematurely from stream for %s (%s) - "
934 "error: %s, sent %d bytes, content_length=%s",
935 queue.display_name,
936 queue_item.name,
937 queue_item.uri,
938 err.__class__.__name__,
939 bytes_sent,
940 resp.content_length,
941 )
942 break
943 finally:
944 self._active_output_streams -= 1
945 if queue_item.streamdetails.stream_error:
946 self.logger.error(
947 "Error streaming QueueItem %s (%s) to %s",
948 queue_item.name,
949 queue_item.uri,
950 queue.display_name,
951 )
952 elif (
953 bytes_sent > 0
954 and queue_item.streamdetails
955 and queue_item.streamdetails.seconds_streamed
956 and queue_item.duration
957 ):
958 # cache the actual encoded bytes-per-second for this URI + output format
959 # so future content_length estimates are near-exact
960 self.mass.create_task(
961 store_content_length_in_cache(
962 self.mass,
963 queue_item.uri,
964 output_format,
965 bytes_sent,
966 queue_item.streamdetails.seconds_streamed,
967 )
968 )
969 return resp
970 finally:
971 # Paired with on_source_selected â fires regardless of how streaming
972 # ended (normal completion, client disconnect, exception). Lets
973 # NAMED_PIPE plugins release ownership without depending on an
974 # external session event. The stream_session_id is the same token
975 # passed to on_source_selected; the provider must reject the
976 # callback if it does not match the currently stored active
977 # session (otherwise a stale teardown from a superseded same-queue
978 # request would clear the live claim of its replacement).
979 if audio_source_provider is not None and audio_source_id is not None:
980 # Provider teardown failures must not break the response cycle
981 # (we're already in finally for a stream that ended one way or
982 # another), but they MUST surface in logs â otherwise a buggy
983 # plugin leaks _in_use_by_queue forever and there is no trail.
984 try:
985 await audio_source_provider.on_source_unselected(
986 audio_source_id, queue_id, stream_session_id
987 )
988 except Exception:
989 self.logger.warning(
990 "on_source_unselected raised for provider %s source %s queue %s",
991 audio_source_provider.instance_id,
992 audio_source_id,
993 queue_id,
994 exc_info=True,
995 )
996
997 async def serve_queue_flow_stream(self, request: web.Request) -> web.StreamResponse: # noqa: PLR0915
998 """Stream Queue Flow audio to player."""
999 self._log_request(request)
1000 queue_id = request.match_info["queue_id"]
1001 player_id = request.match_info["player_id"]
1002 if not (queue := self.mass.player_queues.get(queue_id)):
1003 raise web.HTTPNotFound(reason=f"Unknown Queue: {queue_id}")
1004 session_id = request.match_info["session_id"]
1005 queue_data = self.mass.player_queues.queue_data(queue_id)
1006 if queue_data.session_id is None or session_id != queue_data.session_id:
1007 raise web.HTTPNotFound(reason=f"Unknown (or invalid) session: {session_id}")
1008 if not (player := self.mass.players.get_player(player_id)):
1009 raise web.HTTPNotFound(reason=f"Unknown Player: {player_id}")
1010 start_queue_item_id = request.match_info["queue_item_id"]
1011 start_queue_item = self.mass.player_queues.get_item(queue_id, start_queue_item_id)
1012 if not start_queue_item:
1013 raise web.HTTPNotFound(reason=f"Unknown Queue item: {start_queue_item_id}")
1014
1015 # select the PCM format for the flow stream, anchored on the first track
1016 crossfade_mode = (
1017 self.get_crossfade_mode(queue)
1018 if start_queue_item.media_type == MediaType.TRACK
1019 else CrossfadeMode.DISABLED
1020 )
1021 flow_pcm_format = await self.audio.select_flow_pcm_format(
1022 player,
1023 start_streamdetails=start_queue_item.streamdetails,
1024 crossfade_enabled=crossfade_mode != CrossfadeMode.DISABLED,
1025 overlay_active=overlay_active(queue),
1026 )
1027
1028 # work out output format/details
1029 output_format = await self.audio.get_output_format(
1030 output_format_str=request.match_info["fmt"],
1031 player=player,
1032 content_sample_rate=flow_pcm_format.sample_rate,
1033 content_bit_depth=flow_pcm_format.bit_depth,
1034 media_type=start_queue_item.media_type,
1035 )
1036 # work out ICY metadata support
1037 icy_preference = self.mass.config.get_raw_player_config_value(
1038 player_id,
1039 CONF_ENTRY_ENABLE_ICY_METADATA.key,
1040 CONF_ENTRY_ENABLE_ICY_METADATA.default_value,
1041 )
1042 enable_icy = request.headers.get("Icy-MetaData", "") == "1" and icy_preference != "disabled"
1043 icy_meta_interval = 256000 if icy_preference == "full" else 16384
1044
1045 # prepare request, add some DLNA/UPNP compatible headers.
1046 # icy-name (in DEFAULT_STREAM_HEADERS) is always present so players have a
1047 # readable stream name; the rest of the ICY/shoutcast metadata headers are
1048 # only advertised when the client actually requested ICY metadata, rather
1049 # than on every flow response.
1050 headers = {
1051 **DEFAULT_STREAM_HEADERS,
1052 **(ICY_HEADERS if enable_icy else {}),
1053 "contentFeatures.dlna.org": DLNA_CONTENT_FEATURES_REALTIME,
1054 "Content-Type": get_mime_type(output_format.output_format_str),
1055 }
1056 if enable_icy:
1057 headers["icy-metaint"] = str(icy_meta_interval)
1058
1059 resp = web.StreamResponse(status=200, reason="OK", headers=headers)
1060 http_profile = player.get_config_value(CONF_HTTP_PROFILE, "default")
1061 if http_profile == "forced_content_length":
1062 # just set an insane high content length to make sure the player keeps playing
1063 resp.content_length = calculate_content_length(output_format, 12 * 3600)
1064 elif http_profile == "chunked":
1065 resp.enable_chunked_encoding()
1066
1067 await resp.prepare(request)
1068
1069 # return early if this is not a GET request
1070 if request.method != "GET":
1071 return resp
1072
1073 self._update_audio_processing_context(
1074 queue=queue,
1075 queue_item=start_queue_item,
1076 pcm_format=flow_pcm_format,
1077 crossfade_mode=crossfade_mode,
1078 overlay_enabled=overlay_active(queue),
1079 session_id=session_id,
1080 )
1081 output_plan = self.audio.get_player_output_plan(
1082 player.player_id,
1083 flow_pcm_format,
1084 output_format,
1085 shared_player_ids=player.state.group_members,
1086 queue_id=queue_id,
1087 session_id=session_id,
1088 )
1089
1090 # all checks passed, start streaming!
1091 # this final ffmpeg process in the chain will convert the raw, lossless PCM audio into
1092 # the desired output format for the player including any player specific filter params
1093 # such as channels mixing, DSP, resampling and, only if needed, encoding to lossy formats
1094 self.logger.debug("Start serving Queue flow audio stream for %s", queue.display_name)
1095
1096 # Mark this player as actively streaming so audio analysis yields CPU to playback
1097 # for the duration of the flow stream (see audio_analysis.playback_active).
1098 self._active_output_streams += 1
1099 flow_stream = self.audio.get_queue_flow_stream(
1100 queue=queue,
1101 start_queue_item=start_queue_item,
1102 pcm_format=flow_pcm_format,
1103 session_id=session_id,
1104 protocol_player=player,
1105 )
1106 if overlay_active(queue):
1107 flow_stream = self.audio.get_overlay_mixed_stream(queue, flow_stream, flow_pcm_format)
1108 audio_bytes = get_ffmpeg_stream(
1109 audio_input=flow_stream,
1110 input_format=flow_pcm_format,
1111 output_format=output_format,
1112 filter_params=output_plan.filter_params,
1113 # we need to slowly feed the music to avoid the player stopping and later
1114 # restarting (or completely failing) the audio stream by keeping the buffer short.
1115 # this is reported to be an issue especially with Chromecast players.
1116 # see for example: https://github.com/music-assistant/support/issues/3717
1117 # allow buffer ahead of a few seconds and read rest in (near) realtime
1118 extra_input_args=["-readrate", "1.05", "-readrate_initial_burst", "5"],
1119 chunk_size=icy_meta_interval if enable_icy else calculate_content_length(output_format),
1120 )
1121 client_disconnected = False
1122 try:
1123 # aclosing guarantees the flow stream (and thus the ffmpeg process chain
1124 # behind it) is torn down immediately when the player disconnects
1125 # mid-stream, instead of lingering until garbage collection finalizes
1126 # the abandoned generator.
1127 async with aclosing(audio_bytes):
1128 async for chunk in audio_bytes:
1129 try:
1130 await resp.write(chunk)
1131 except BrokenPipeError, ConnectionResetError, ConnectionError:
1132 # race condition
1133 client_disconnected = True
1134 break
1135
1136 if not enable_icy:
1137 continue
1138
1139 # if icy metadata is enabled, send the icy metadata after the chunk
1140 if (
1141 # use current item here and not buffered item, otherwise
1142 # the icy metadata will be too much ahead
1143 (current_item := queue.current_item)
1144 and current_item.streamdetails
1145 and current_item.streamdetails.stream_title
1146 ):
1147 title = current_item.streamdetails.stream_title
1148 elif queue and current_item and current_item.name:
1149 title = current_item.name
1150 else:
1151 title = "Music Assistant"
1152 metadata = f"StreamTitle='{title}';".encode()
1153 if icy_preference == "full" and current_item and current_item.image:
1154 metadata += f"StreamURL='{current_item.image.path}'".encode()
1155 while len(metadata) % 16 != 0:
1156 metadata += b"\x00"
1157 length = len(metadata)
1158 length_b = chr(int(length / 16)).encode()
1159 await resp.write(length_b + metadata)
1160 finally:
1161 self._active_output_streams -= 1
1162
1163 if not client_disconnected and http_profile == "forced_content_length":
1164 await self._finish_flow_stream(resp, queue_id, session_id)
1165
1166 return resp
1167
1168 async def serve_command_request(self, request: web.Request) -> web.FileResponse:
1169 """Handle special 'command' request for a player."""
1170 self._log_request(request)
1171 queue_id = request.match_info["queue_id"]
1172 command = request.match_info["command"]
1173 if command == "next":
1174 self.mass.create_task(self.mass.player_queues.next(queue_id))
1175 return web.FileResponse(SILENCE_FILE, headers={"icy-name": "Music Assistant"})
1176
1177 async def serve_announcement_stream(self, request: web.Request) -> web.StreamResponse:
1178 """Stream announcement audio to a player."""
1179 self._log_request(request)
1180 player_id = request.match_info["player_id"]
1181 if not (player := self.mass.players.get_player(player_id)):
1182 raise web.HTTPNotFound(reason=f"Unknown Player: {player_id}")
1183 if not (announce_data := self.announcement_renderer.get_for_player(player_id)):
1184 raise web.HTTPNotFound(reason=f"No pending announcements for Player: {player_id}")
1185
1186 # work out output format/details
1187 fmt = request.match_info["fmt"]
1188 audio_format = AudioFormat(content_type=ContentType.try_parse(fmt))
1189
1190 http_profile = self._get_announcement_http_profile(player_id, announce_data)
1191
1192 # return early if this is not a GET request:
1193 # players often probe the url with a HEAD request before fetching it and
1194 # rendering the announcement for such a probe would run the entire (costly)
1195 # TTS/ffmpeg chain twice for a single announcement.
1196 if request.method != "GET":
1197 resp = web.StreamResponse(status=200, reason="OK", headers=DEFAULT_STREAM_HEADERS)
1198 resp.content_type = get_mime_type(audio_format.output_format_str)
1199 if http_profile == "chunked":
1200 resp.enable_chunked_encoding()
1201 await resp.prepare(request)
1202 return resp
1203
1204 if http_profile == "forced_content_length":
1205 # given the fact that an announcement is just a short audio clip,
1206 # just send it over completely at once so we have a fixed content length
1207 data = bytearray()
1208 announcement_stream = self.get_announcement_stream(announce_data, audio_format)
1209 # aclosing guarantees the stream (and thus the ffmpeg process chain behind
1210 # it) is torn down immediately when the request is cancelled, instead of
1211 # lingering until garbage collection finalizes the abandoned generator.
1212 async with aclosing(announcement_stream):
1213 async for chunk in announcement_stream:
1214 data += chunk
1215 return web.Response(
1216 body=bytes(data),
1217 content_type=get_mime_type(audio_format.output_format_str),
1218 headers=DEFAULT_STREAM_HEADERS,
1219 )
1220
1221 resp = web.StreamResponse(status=200, reason="OK", headers=DEFAULT_STREAM_HEADERS)
1222 resp.content_type = get_mime_type(audio_format.output_format_str)
1223 if http_profile == "chunked":
1224 resp.enable_chunked_encoding()
1225
1226 await resp.prepare(request)
1227
1228 # all checks passed, start streaming!
1229 self.logger.debug(
1230 "Start serving audio stream for Announcement %s to %s",
1231 announce_data["announcement_url"],
1232 player.display_name,
1233 )
1234 announcement_stream = self.get_announcement_stream(announce_data, audio_format)
1235 # aclosing guarantees the stream (and thus the ffmpeg process chain behind
1236 # it) is torn down immediately when the player disconnects mid-stream,
1237 # instead of lingering until garbage collection finalizes the abandoned
1238 # generator.
1239 async with aclosing(announcement_stream):
1240 async for chunk in announcement_stream:
1241 try:
1242 await resp.write(chunk)
1243 except BrokenPipeError, ConnectionResetError:
1244 break
1245
1246 self.logger.debug(
1247 "Finished serving audio stream for Announcement %s to %s",
1248 announce_data["announcement_url"],
1249 player.display_name,
1250 )
1251
1252 return resp
1253
1254 def get_command_url(self, player_or_queue_id: str, command: str) -> str:
1255 """Get the url for the special command stream."""
1256 return f"{self.base_url}/command/{player_or_queue_id}/{command}.mp3"
1257
1258 def get_announcement_url(
1259 self,
1260 player_id: str,
1261 content_type: ContentType = ContentType.MP3,
1262 ) -> str:
1263 """
1264 Get the url that serves the announcement registered for the given player.
1265
1266 :param player_id: The player the announcement is played on.
1267 :param content_type: The format to serve the announcement in.
1268 """
1269 # use stream server to host announcement on local network
1270 # this ensures playback on all players, including ones that do not
1271 # like https hosts and it also offers the pre-announce 'bell'
1272 return f"{self.base_url}/announcement/{player_id}.{content_type.value}"
1273
1274 def get_stream(
1275 self,
1276 media: PlayerMedia,
1277 pcm_format: AudioFormat,
1278 player_id: str | None = None,
1279 force_flow_mode: bool = False,
1280 use_flow_stream_buffering: bool = False,
1281 ) -> AsyncGenerator[bytes]:
1282 """
1283 Get a stream of the given media as raw PCM audio.
1284
1285 This is used as helper for player providers that can consume the raw PCM
1286 audio stream directly (e.g. AirPlay) and not rely on HTTP transport.
1287
1288 :param media: The PlayerMedia to stream.
1289 :param pcm_format: The desired output PCM format.
1290 :param player_id: The player ID requesting the stream. Used to determine
1291 if flow mode should be used based on the player's capabilities.
1292 :param force_flow_mode: Force flow mode regardless of player capabilities.
1293 Used for multi-client streaming scenarios that require continuous streams.
1294 :param use_flow_stream_buffering: Buffer the flow stream to provide headroom
1295 during smart fades transitions. Use for consumers that read directly
1296 (e.g. AirPlay, Snapcast) and can't tolerate stalls.
1297 """
1298 # select audio source
1299 if media.media_type == MediaType.ANNOUNCEMENT:
1300 # special case: stream announcement
1301 assert media.custom_data
1302 return self.get_announcement_stream(cast("AnnounceData", media.custom_data), pcm_format)
1303 if (
1304 media.source_id
1305 and media.source_id.startswith(UGP_PREFIX)
1306 and media.uri
1307 and "/ugp/" in media.uri
1308 ):
1309 # special case: member player accessing UGP stream
1310 # Check URI to distinguish from the UGP accessing its own stream
1311 ugp_player = cast("UniversalGroupPlayer", self.mass.players.get_player(media.source_id))
1312 ugp_stream = ugp_player.stream
1313 assert ugp_stream is not None # for type checker
1314 if ugp_stream.base_pcm_format == pcm_format:
1315 # no conversion needed
1316 return ugp_stream.subscribe_raw()
1317 return ugp_stream.get_stream(output_format=pcm_format)
1318 if media.source_id and media.queue_item_id:
1319 # Queue stream request - determine flow_mode based on player capabilities
1320 # or force it if explicitly requested (e.g., for multi-client streaming)
1321 protocol_player = self.mass.players.get_player(player_id) if player_id else None
1322 queue_id = media.source_id
1323 queue = self.mass.player_queues.get(queue_id)
1324 queue_session_id = cast(
1325 "str | None",
1326 (media.custom_data or {}).get("session_id"),
1327 )
1328 crossfade_needs_flow_mode = (
1329 # crossfade only applies to tracks; if the queue has it enabled but the
1330 # player(protocol) does not support gapless playback, we need to enforce flow mode
1331 media.media_type == MediaType.TRACK
1332 and queue is not None
1333 and queue.crossfade_enabled
1334 and protocol_player
1335 and not protocol_player.supports_gapless
1336 )
1337 # the audio overlay is mixed into the queue's continuous (flow) stream;
1338 # per-item requests would restart the overlay at every track boundary
1339 overlay_needs_flow_mode = queue is not None and overlay_active(queue)
1340 flow_mode = (
1341 force_flow_mode
1342 or (protocol_player is not None and protocol_player.flow_mode)
1343 or crossfade_needs_flow_mode
1344 or overlay_needs_flow_mode
1345 )
1346 if media.media_type in (MediaType.RADIO, MediaType.AUDIO_SOURCE):
1347 # flow_mode for live/infinite streams is pointless
1348 flow_mode = False
1349 if flow_mode:
1350 # flow stream request
1351 assert queue
1352 start_queue_item = self.mass.player_queues.get_item(
1353 media.source_id, media.queue_item_id
1354 )
1355 assert start_queue_item
1356 crossfade_mode = (
1357 self.get_crossfade_mode(queue)
1358 if start_queue_item.media_type == MediaType.TRACK
1359 else CrossfadeMode.DISABLED
1360 )
1361 self._update_audio_processing_context(
1362 queue=queue,
1363 queue_item=start_queue_item,
1364 pcm_format=pcm_format,
1365 crossfade_mode=crossfade_mode,
1366 overlay_enabled=overlay_active(queue),
1367 session_id=queue_session_id,
1368 )
1369 flow_stream = self.audio.get_queue_flow_stream(
1370 queue=queue,
1371 start_queue_item=start_queue_item,
1372 pcm_format=pcm_format,
1373 session_id=queue_session_id,
1374 protocol_player=protocol_player,
1375 )
1376 if overlay_active(queue):
1377 flow_stream = self.audio.get_overlay_mixed_stream(
1378 queue, flow_stream, pcm_format
1379 )
1380 if use_flow_stream_buffering:
1381 return buffered(flow_stream, buffer_size=30, min_buffer_before_yield=1)
1382 return flow_stream
1383 # single item stream (e.g. radio or non-flow mode)
1384 queue_item = self.mass.player_queues.get_item(media.source_id, media.queue_item_id)
1385 assert queue_item
1386 if queue is not None:
1387 self._update_audio_processing_context(
1388 queue=queue,
1389 queue_item=queue_item,
1390 pcm_format=pcm_format,
1391 crossfade_mode=(
1392 self.get_crossfade_mode(queue)
1393 if queue_item.media_type == MediaType.TRACK
1394 else CrossfadeMode.DISABLED
1395 ),
1396 overlay_enabled=(
1397 queue_item.media_type == MediaType.RADIO and overlay_active(queue)
1398 ),
1399 session_id=queue_session_id,
1400 )
1401 inner_stream = self.audio.get_queue_item_stream(
1402 queue_item=queue_item,
1403 pcm_format=pcm_format,
1404 seek_position=(
1405 int(queue_item.streamdetails.seek_position) if queue_item.streamdetails else 0
1406 ),
1407 playback_speed=cast(
1408 "float", queue_item.extra_attributes.get("playback_speed", 1.0)
1409 ),
1410 session_id=queue_session_id,
1411 )
1412 if (
1413 queue is not None
1414 and queue_item.media_type == MediaType.RADIO
1415 and overlay_active(queue)
1416 ):
1417 # radio plays as a single long-lived stream, so mix the overlay in here
1418 inner_stream = self.audio.get_overlay_mixed_stream(queue, inner_stream, pcm_format)
1419 # mirror the on_source_selected/unselected lifecycle the HTTP route
1420 # fires, so direct-PCM consumers (AirPlay, Snapcast, UGP) honour the
1421 # plugin contract too
1422 if (
1423 queue_item.media_item is not None
1424 and queue_item.media_item.media_type == MediaType.AUDIO_SOURCE
1425 ):
1426 return self._wrap_with_audio_source_lifecycle(
1427 inner=inner_stream,
1428 queue_item=queue_item,
1429 player_id=player_id or media.source_id,
1430 )
1431 return inner_stream
1432 # assume url or some other direct path
1433 # NOTE: this will fail if its an uri not playable by ffmpeg
1434 return get_ffmpeg_stream(
1435 audio_input=media.uri,
1436 input_format=AudioFormat(content_type=ContentType.try_parse(media.uri)),
1437 output_format=pcm_format,
1438 )
1439
1440 async def get_preview_stream(
1441 self,
1442 provider_instance_id_or_domain: str,
1443 item_id: str,
1444 media_type: MediaType = MediaType.TRACK,
1445 ) -> AsyncGenerator[bytes]:
1446 """Create a 30 seconds preview audioclip for the given media item."""
1447 if not (music_prov := self.mass.get_provider(provider_instance_id_or_domain)):
1448 raise ProviderUnavailableError
1449 if music_prov.type != ProviderType.MUSIC:
1450 msg = f"{provider_instance_id_or_domain} is not a music provider"
1451 raise InvalidDataError(msg)
1452 music_prov = cast("MusicProvider", music_prov)
1453
1454 try:
1455 await self.mass.music.get_item(
1456 media_type,
1457 item_id,
1458 provider_instance_id_or_domain,
1459 allow_update_metadata=False,
1460 )
1461 except MediaNotFoundError as err:
1462 msg = f"Item {item_id} not found in provider {provider_instance_id_or_domain}"
1463 raise InvalidDataError(msg) from err
1464
1465 streamdetails = await music_prov.get_stream_details(item_id, media_type)
1466 pcm_format = AudioFormat(
1467 content_type=ContentType.from_bit_depth(streamdetails.audio_format.bit_depth),
1468 sample_rate=streamdetails.audio_format.sample_rate,
1469 bit_depth=streamdetails.audio_format.bit_depth,
1470 channels=streamdetails.audio_format.channels,
1471 )
1472 async for chunk in get_ffmpeg_stream(
1473 audio_input=self.audio.get_media_stream(
1474 streamdetails=streamdetails, pcm_format=pcm_format
1475 ),
1476 input_format=pcm_format,
1477 output_format=AudioFormat(content_type=ContentType.AAC),
1478 extra_input_args=["-t", "30"],
1479 ):
1480 yield chunk
1481
1482 async def get_announcement_stream(
1483 self, announce_data: AnnounceData, output_format: AudioFormat
1484 ) -> AsyncGenerator[bytes]:
1485 """
1486 Get the audio of an announcement (pre-announce chime + announcement).
1487
1488 Any number of consumers may stream the same announcement at once; its source is
1489 fetched and decoded only once. The audio stays available while the stream is
1490 held open.
1491
1492 :param announce_data: The announcement to stream.
1493 :param output_format: The format to deliver the audio in.
1494 """
1495 render = self.announcement_renderer.acquire(announce_data)
1496 try:
1497 # aclosing guarantees this consumer's ffmpeg encoder is torn down
1498 # immediately when it goes away, instead of lingering until garbage
1499 # collection finalizes the abandoned generator.
1500 stream = render.get_stream(output_format)
1501 async with aclosing(stream):
1502 async for chunk in stream:
1503 yield chunk
1504 finally:
1505 await self.announcement_renderer.release(render)
1506
1507 async def get_announcement_duration(
1508 self, announcement: PlayerMedia, timeout: float = DEFAULT_RENDER_TIMEOUT
1509 ) -> int | None:
1510 """
1511 Get the exact duration (in seconds) of an announcement, once it finished rendering.
1512
1513 Waits for the audio to be rendered in full, so call this while the announcement
1514 plays rather than before handing it to a player. Returns None when the length can
1515 not be determined, e.g. the announcement is no longer playing or its source did
1516 not deliver in time.
1517
1518 :param announcement: The announcement to return the duration for.
1519 :param timeout: Maximum time to wait for the audio to finish rendering.
1520 """
1521 if announcement.duration:
1522 return announcement.duration
1523 if not announcement.custom_data:
1524 return None
1525 render = self.announcement_renderer.get(cast("AnnounceData", announcement.custom_data))
1526 if render is None:
1527 return None
1528 duration = await render.wait_finished(timeout)
1529 return ceil(duration) if duration else None
1530
1531 async def _wrap_with_audio_source_lifecycle(
1532 self,
1533 inner: AsyncGenerator[bytes],
1534 queue_item: QueueItem,
1535 player_id: str,
1536 ) -> AsyncGenerator[bytes]:
1537 """
1538 Wrap an AudioSource queue item stream with on_source_selected/unselected hooks.
1539
1540 Direct-PCM consumers (AirPlay, Snapcast, UGP, ...) call ``get_stream`` instead
1541 of going through the HTTP route, but the plugin contract requires the
1542 lifecycle hooks to fire for every actual stream request â they're what
1543 claim/release the per-queue exclusive ownership and trigger acquisition
1544 side effects like the Spotify Connect Web API play kick. This wrapper
1545 gives those consumers the same lifecycle the HTTP route already provides.
1546
1547 :param inner: The underlying audio stream generator.
1548 :param queue_item: The AudioSource queue item being streamed.
1549 :param player_id: The protocol player consuming this stream.
1550 """
1551 media_item = queue_item.media_item
1552 assert media_item is not None # caller checked media_type == AUDIO_SOURCE
1553 prov = self.mass.get_provider(media_item.provider)
1554 queue_id = queue_item.queue_id
1555 if not isinstance(prov, PluginProvider):
1556 async for chunk in inner:
1557 yield chunk
1558 return
1559 source_id = media_item.item_id
1560 stream_session_id = uuid4().hex
1561 # single try/finally so on_source_unselected fires even when
1562 # on_source_selected raises after partially claiming state; the
1563 # provider's session_id guard makes a no-op claim release safe.
1564 try:
1565 try:
1566 await prov.on_source_selected(source_id, player_id, queue_id, stream_session_id)
1567 except RuntimeError as err:
1568 # provider intentionally aborts the request â surface as AudioError
1569 raise AudioError(str(err)) from err
1570 async for chunk in inner:
1571 yield chunk
1572 finally:
1573 try:
1574 await prov.on_source_unselected(source_id, queue_id, stream_session_id)
1575 except Exception as err:
1576 self.logger.exception(
1577 "on_source_unselected raised for provider %s source %s queue %s: %s",
1578 prov.instance_id,
1579 source_id,
1580 queue_id,
1581 err,
1582 )
1583
1584 def _update_audio_processing_context(
1585 self,
1586 queue: PlayerQueue,
1587 queue_item: QueueItem,
1588 pcm_format: AudioFormat,
1589 crossfade_mode: CrossfadeMode,
1590 overlay_enabled: bool,
1591 session_id: str | None = None,
1592 ) -> None:
1593 """
1594 Store the shared processing context selected for a queue item.
1595
1596 :param queue: Active player queue.
1597 :param queue_item: Queue item being prepared.
1598 :param pcm_format: Shared PCM format leaving queue processing.
1599 :param crossfade_mode: Effective crossfade mode for the item.
1600 :param overlay_enabled: Whether an overlay is mixed into this stream.
1601 :param session_id: Queue session that owns processing-detail updates.
1602 """
1603 if queue_item.streamdetails is None:
1604 return
1605 queue_data = self.mass.player_queues.queue_data_or_none(queue.queue_id)
1606 if (
1607 queue_data is None
1608 or (processing_session_id := session_id or queue_data.session_id) is None
1609 or queue_data.session_id != processing_session_id
1610 ):
1611 return
1612 self.audio_processing.start_session(queue.queue_id, processing_session_id)
1613 self.audio_processing.update_item_context(
1614 queue_id=queue.queue_id,
1615 session_id=processing_session_id,
1616 queue_item_id=queue_item.queue_item_id,
1617 queue_processing=AudioQueueProcessing(
1618 pcm_format=pcm_format,
1619 playback_speed=cast(
1620 "float",
1621 queue_item.extra_attributes.get("playback_speed", 1.0),
1622 ),
1623 crossfade_mode=crossfade_mode,
1624 overlay_active=overlay_enabled,
1625 ),
1626 alters_audio=queue_item.streamdetails.fade_in,
1627 )
1628
1629 def _get_announcement_http_profile(self, player_id: str, announce_data: AnnounceData) -> str:
1630 """
1631 Resolve the http profile for serving an announcement stream.
1632
1633 Announcement urls are registered under the visible player's id, but the
1634 stream may be fetched by a linked protocol player; the profile must come
1635 from the player that actually performs the fetch.
1636 """
1637 announce_player = None
1638 if announce_player_id := announce_data.get("announce_player_id"):
1639 announce_player = self.mass.players.get_player(announce_player_id)
1640 if announce_player is None:
1641 announce_player = self.mass.players.get_player(player_id)
1642 if announce_player is None:
1643 return "default"
1644 return announce_player.get_output_config_value(CONF_HTTP_PROFILE, "default")
1645
1646 async def _finish_flow_stream(
1647 self, resp: web.StreamResponse, queue_id: str, session_id: str
1648 ) -> None:
1649 """
1650 Close a fully served flow stream, giving the player time to drain when it ends the queue.
1651
1652 :param resp: The flow stream response, already fully written.
1653 :param queue_id: Id of the queue the flow stream belongs to.
1654 :param session_id: Stream session this response was opened for.
1655 """
1656 if self.mass.player_queues.flow_queue_exhausted(queue_id, session_id):
1657 # the player is still holding a few seconds of audio it has not rendered yet
1658 # and drops that as soon as the stream ends, so let it play out first.
1659 # a flow that ends to be restarted right away gets no such grace: there the
1660 # player should go idle as soon as possible so the next stream can start.
1661 self.logger.debug(
1662 "Flow stream for queue %s reached the end of the queue - holding the "
1663 "connection open for %ss so the player can play out its buffer",
1664 queue_id,
1665 FLOW_STREAM_LEAD_OUT_SECONDS,
1666 )
1667 await asyncio.sleep(FLOW_STREAM_LEAD_OUT_SECONDS)
1668 # aiohttp derives keep-alive from the request, so the 'Connection: close' we
1669 # advertise is relayed to the player but never applied to the response itself.
1670 # Without this the player is left waiting on a stream that already ended.
1671 resp.force_close()
1672
1673 def _log_request(self, request: web.Request) -> None:
1674 """Log request."""
1675 if self.logger.isEnabledFor(VERBOSE_LOG_LEVEL):
1676 self.logger.log(
1677 VERBOSE_LOG_LEVEL,
1678 "Got %s request to %s from %s\nheaders: %s\n",
1679 request.method,
1680 request.path,
1681 request.remote,
1682 redact_sensitive_headers(request.headers),
1683 )
1684 else:
1685 self.logger.debug(
1686 "Got %s request to %s from %s (HTTP/%s.%s, connection: %s)",
1687 request.method,
1688 request.path,
1689 request.remote,
1690 request.version.major,
1691 request.version.minor,
1692 request.headers.get("Connection", "-"),
1693 )
1694
1695 async def _reload_network_dependent_providers(self) -> None:
1696 """Reload the providers that captured the streamserver network, if it changed."""
1697 previous = self._network_fingerprint
1698 current = (
1699 self._bind_ip,
1700 str(self.publish_ip),
1701 cast("int", self.publish_port),
1702 tuple(self._publish_addresses),
1703 )
1704 if previous is None or previous == current:
1705 self._network_fingerprint = current
1706 return
1707 # these providers bind or advertise the network while they load, so a plain
1708 # reload is what moves them over - they share no lighter rebind path
1709 instance_ids = [
1710 prov.instance_id
1711 for prov in self.mass.providers
1712 if prov.reload_on_streams_network_change
1713 ]
1714 for instance_id in instance_ids:
1715 try:
1716 config = await self.mass.config.get_provider_config(instance_id)
1717 self.logger.info(
1718 "Streamserver network changed, reloading provider %s",
1719 config.name or config.domain,
1720 )
1721 await self.mass.load_provider_config(config)
1722 except Exception as err:
1723 self.logger.warning(
1724 "Error reloading provider %s: %s",
1725 instance_id,
1726 str(err) or err.__class__.__name__,
1727 exc_info=err,
1728 )
1729 # only mark the new network as applied once the loop completed, so a run cut short
1730 # by a second config change runs again on the next reload
1731 self._network_fingerprint = current
1732
1733 def _setup_smart_fades_logger(self, config: CoreConfig) -> None:
1734 """Set up smart fades logger level."""
1735 log_level = str(config.get_value(CONF_SMART_FADES_LOG_LEVEL))
1736 if log_level == "GLOBAL":
1737 self.audio.smart_fades_mixer.logger.setLevel(self.logger.level)
1738 else:
1739 self.audio.smart_fades_mixer.logger.setLevel(log_level)
1740
1741 def _resolve_publish_state(self, bind_ip: str, publish_candidates: tuple[str, ...]) -> None:
1742 """
1743 Resolve the addresses and base URL to advertise for the given bind address.
1744
1745 Reads ``self.publish_port``, so set that first.
1746
1747 :param bind_ip: Address the streamserver binds to (a wildcard means all interfaces).
1748 :param publish_candidates: Host addresses reachable from the local network, ranked.
1749 """
1750 self._bind_ip = bind_ip
1751 self._publish_addresses = _get_publish_addresses(
1752 bind_ip, self._configured_publish_ip, publish_candidates
1753 )
1754 # the single address players are handed, taken from the top of the ranked list
1755 self.publish_ip = self._publish_addresses[0]
1756 self._base_url = f"http://{format_ip_for_url(self.publish_ip)}:{self.publish_port}"
1757
1758
1759def _same_ip_family(ip: str, other_ip: str) -> bool:
1760 """Return whether two addresses belong to the same IP family."""
1761 return (":" in ip) == (":" in other_ip)
1762