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