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