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