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