/
/
1"""
2Spotify Connect plugin for Music Assistant.
3
4We tie a single player to a single Spotify Connect daemon (go-librespot).
5The provider has multi instance support, so multiple players can be linked to
6multiple Spotify Connect daemons.
7
8go-librespot is driven entirely over its local HTTP+WebSocket API: the WebSocket
9``/events`` stream feeds player/session/metadata/volume state into Music Assistant,
10and transport + volume commands are issued via the REST endpoints. Playback control
11works without a configured Spotify *music* provider or the Spotify Web API.
12"""
13
14from __future__ import annotations
15
16import asyncio
17import json
18import os
19import time
20from collections.abc import AsyncGenerator
21from contextlib import suppress
22from typing import TYPE_CHECKING, Any, cast
23
24from music_assistant_models.enums import (
25 ContentType,
26 MediaType,
27 PlaybackState,
28 ProviderFeature,
29 SourceControl,
30 StreamType,
31)
32from music_assistant_models.errors import AudioError, MediaNotFoundError
33from music_assistant_models.media_items import AudioFormat, AudioSource, ProviderMapping
34from music_assistant_models.streamdetails import StreamDetails, StreamMetadata
35
36from music_assistant.constants import CONF_ENTRY_WARN_PREVIEW
37from music_assistant.helpers.process import AsyncProcess
38from music_assistant.helpers.util import (
39 interface_name_for_ip,
40 is_port_in_use,
41 select_free_port,
42)
43from music_assistant.models.plugin import PluginProvider
44
45from .client import GoLibrespotClient
46from .helpers import generate_device_id, get_go_librespot_binary
47
48if TYPE_CHECKING:
49 from music_assistant_models.config_entries import ConfigEntry, ProviderConfig
50 from music_assistant_models.provider import ProviderManifest
51
52 from music_assistant.mass import MusicAssistant
53 from music_assistant.models import ProviderInstanceType
54
55CONF_MASS_PLAYER_ID = "mass_player_id"
56CONF_PUBLISH_NAME = "publish_name"
57DEFAULT_PUBLISH_NAME = "Music Assistant"
58
59# Special value for auto player selection
60PLAYER_ID_AUTO = "__auto__"
61
62SUPPORTED_FEATURES = {ProviderFeature.AUDIO_SOURCE}
63
64# stable id for the single AudioSource this provider exposes;
65# combined with the provider instance_id this forms the persistent uri
66AUDIO_SOURCE_ID = "main"
67
68# go-librespot volume scale; we pin volume_steps to this so the daemon's 0..max
69# volume maps 1:1 to a 0-100 percentage.
70VOLUME_STEPS = 100
71
72# Read size for pulling PCM off the daemon's stdout.
73STREAM_READ_CHUNK = 16384
74
75# When playback is paused the daemon stops writing PCM. If no PCM arrives for
76# this long while we're not in a 'playing' state, end the stream (clean EOF) so
77# the player leaves the playing state; the next 'playing' event re-streams.
78PAUSE_EOF_TIMEOUT_S = 0.5
79
80# Port range the go-librespot API server binds to (loopback only, one per instance).
81API_PORT_RANGE_START = 38800
82API_PORT_RANGE_END = 38900
83
84# Seconds to wait for the daemon to report 'playing' after a resume request.
85PLAYBACK_START_TIMEOUT_S = 3.0
86
87# Debounce before acting on an externally-triggered 'playing' event (see
88# _deferred_play_media_fire for why).
89PLAY_MEDIA_DEBOUNCE_S = 0.5
90
91# Ignore Spotify volume events for this long after a session becomes active, so
92# the player's own volume wins over librespot's initial value on (re)connect.
93INITIAL_VOLUME_GRACE_S = 3.0
94
95# User-facing message for the "not the active Spotify device" failure.
96# {0} is the Spotify Connect device's published name (see _not_active_error).
97NOT_ACTIVE_DEVICE_MESSAGE = (
98 "'{0}' is not the active Spotify playback device. "
99 "Open the Spotify app, select it as the playback device, and try again."
100)
101
102
103async def setup(
104 mass: MusicAssistant, manifest: ProviderManifest, config: ProviderConfig
105) -> ProviderInstanceType:
106 """Initialize provider(instance) with given configuration."""
107 return SpotifyConnectProvider(mass, manifest, config)
108
109
110class SpotifyConnectProvider(PluginProvider):
111 """Implementation of a Spotify Connect Plugin (backed by go-librespot)."""
112
113 reload_on_streams_network_change = True
114
115 def __init__(
116 self, mass: MusicAssistant, manifest: ProviderManifest, config: ProviderConfig
117 ) -> None:
118 """Initialize MusicProvider."""
119 super().__init__(mass, manifest, config, SUPPORTED_FEATURES)
120 # Configured default player (PLAYER_ID_AUTO or a specific player id)
121 self._default_player_id: str = (
122 cast("str", self.get_setup_value(CONF_MASS_PLAYER_ID)) or PLAYER_ID_AUTO
123 )
124 self._publish_name = (
125 cast("str", self.get_setup_value(CONF_PUBLISH_NAME)) or DEFAULT_PUBLISH_NAME
126 )
127 # Currently active player (the one currently playing or selected)
128 self._active_player_id: str | None = None
129 self.cache_dir = os.path.join(self.mass.cache_path, self.instance_id)
130 self._binary: str | None = None
131 self._api_port: int = 0
132 self._client: GoLibrespotClient | None = None
133 self._stop_called: bool = False
134 self._daemon_task: asyncio.Task[None] | None = None
135 self._events_task: asyncio.Task[None] | None = None
136 self._proc: AsyncProcess | None = None
137 self._restart_error_count = 0
138 self.logger.debug(
139 "Init plugin with name '%s' for player '%s' with instance id '%s'",
140 self.name,
141 self._default_player_id,
142 self.instance_id,
143 )
144 # _audio_format is the original Spotify source codec (Ogg Vorbis 320 kbps),
145 # advertised to clients for display. _decoded_audio_format is the raw PCM
146 # go-librespot actually writes to its stdout after decoding â what
147 # get_audio_stream yields and what the streams controller hands ffmpeg as
148 # the input format. We always emit the source's own format here; MA is
149 # responsible for converting it to whatever each player needs.
150 self._audio_format = AudioFormat(
151 content_type=ContentType.OGG,
152 codec_type=ContentType.VORBIS,
153 sample_rate=44100,
154 bit_depth=16,
155 channels=2,
156 bit_rate=320,
157 )
158 self._decoded_audio_format = AudioFormat(
159 content_type=ContentType.PCM_S16LE,
160 codec_type=ContentType.PCM_S16LE,
161 sample_rate=44100,
162 bit_depth=16,
163 channels=2,
164 )
165 self._stream_metadata = StreamMetadata(title=f"Spotify Connect | {self._publish_name}")
166 self._audio_source = self._build_audio_source()
167 # _in_use_by_queue is the queue currently streaming us. Claimed in
168 # on_source_selected (NOT in get_stream_details â that path also runs
169 # from queue preload, where claiming would block a later cross-queue
170 # handoff). Released in on_source_unselected when the session id
171 # matches, or in _clear_active_player on the daemon's 'inactive' event.
172 self._in_use_by_queue: str | None = None
173 # _active_session_id is the controller-provided token for the current
174 # stream request â used to reject stale on_source_unselected callbacks
175 # after a same-queue reconnect supersedes the previous request.
176 self._active_session_id: str | None = None
177 # tracks the daemon's play/pause state from its 'playing' / 'paused' /
178 # 'inactive' events; gates the resume kick in on_source_selected (skip if
179 # already playing) and the play_media trigger in the event handler.
180 self._playing: bool = False
181 # True while MA is the active Spotify Connect device (set on 'active',
182 # cleared on 'inactive'); gates get_stream_details and transport commands.
183 self._spotify_session_active: bool = False
184 # holds the single in-flight deferred play_media task scheduled from a
185 # 'playing' event; cancelled when a 'paused' / 'stopped' / 'active' event
186 # arrives during the debounce so we don't act on stale state from a dying
187 # session.
188 self._pending_play_media_task: asyncio.Task[None] | None = None
189 self._last_session_active_time: float = 0
190 self._last_volume_sent: int | None = None
191 # Last context/track URIs seen on the event stream. Used to take playback
192 # back (make ourselves the active Spotify device) when the user switched
193 # the active device away in the Spotify app and then presses play in MA.
194 self._last_context_uri: str | None = None
195 self._last_track_uri: str | None = None
196
197 @property
198 def instance_name_postfix(self) -> str | None:
199 """Return the advertised device name as the multi-instance postfix."""
200 return self._publish_name if self._publish_name != DEFAULT_PUBLISH_NAME else None
201
202 async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
203 """Return runtime options for this provider."""
204 return (CONF_ENTRY_WARN_PREVIEW,)
205
206 async def handle_async_init(self) -> None:
207 """Handle async initialization of the provider."""
208 self._binary = get_go_librespot_binary()
209 self._api_port = await select_free_port(
210 API_PORT_RANGE_START, API_PORT_RANGE_END, host="127.0.0.1"
211 )
212 self._client = GoLibrespotClient(
213 self.mass, f"http://127.0.0.1:{self._api_port}", self.logger
214 )
215 # Two self-healing supervisors: one keeps the daemon process alive, the
216 # other keeps the events websocket connected (reconnecting across daemon
217 # restarts). The events runner resets the daemon's restart backoff once
218 # the websocket is healthy again.
219 self._daemon_task = self.mass.create_task(self._daemon_runner())
220 self._events_task = self.mass.create_task(self._events_runner())
221
222 async def unload(self, is_removed: bool = False) -> None:
223 """Handle close/cleanup of the provider."""
224 self._stop_called = True
225 self._cancel_pending_play_media()
226 for task in (self._events_task, self._daemon_task):
227 if task and not task.done():
228 task.cancel()
229 with suppress(asyncio.CancelledError):
230 await task
231
232 @property
233 def active_player_id(self) -> str | None:
234 """Return the currently active player ID for this plugin."""
235 return self._active_player_id
236
237 async def get_audio_sources(self) -> list[AudioSource]:
238 """Return the AudioSources this plugin currently exposes."""
239 return [self._audio_source]
240
241 async def get_stream_details(self, item_id: str, media_type: MediaType) -> StreamDetails:
242 """
243 Return StreamDetails for streaming the Spotify Connect audio.
244
245 Side-effect-free: ownership is claimed in on_source_selected (which the
246 streams controller fires before this method on the actual stream
247 request). Keeping this idempotent means preload paths can fetch
248 streamdetails without claiming the source and blocking a cross-queue
249 handoff.
250
251 Raises AudioError when MA is not the active Spotify Connect device, since
252 playback can only be acquired while a Spotify session is connected to us
253 (entry must come from the Spotify app â see can_initiate below).
254 """
255 if item_id != AUDIO_SOURCE_ID:
256 raise MediaNotFoundError(f"Unknown AudioSource: {item_id}")
257 # Only refuse when we can neither resume nor take playback back. If a last
258 # context is known we let the stream proceed; on_source_selected then takes
259 # playback back (makes us the active device) before audio is pulled.
260 if not self._playing and not self._spotify_session_active and not self._last_context_uri:
261 raise self._not_active_error()
262 # CUSTOM: the core pulls PCM from get_audio_stream. Reading the daemon's
263 # stdout means a consumer is always attached, so go-librespot's non-blocking
264 # pipe open succeeds, and it lets us end the stream cleanly when playback
265 # pauses so the player leaves the playing state. decoded_audio_format tells
266 # the core the PCM is s16le while audio_format keeps the Ogg/Vorbis source
267 # codec for display; MA resamples to each player's format as needed.
268 # `-fflags nobuffer` keeps ffmpeg from buffering ahead on the resample path
269 # (we back-pressure the daemon to realtime in get_audio_stream).
270 # expiration=0: never reuse a cached streamdetails so the active-device
271 # check above re-runs on every play attempt.
272 return StreamDetails(
273 provider=self.instance_id,
274 item_id=item_id,
275 audio_format=self._audio_format,
276 decoded_audio_format=self._decoded_audio_format,
277 media_type=MediaType.AUDIO_SOURCE,
278 stream_type=StreamType.CUSTOM,
279 stream_metadata=self._stream_metadata,
280 extra_input_args=["-fflags", "nobuffer"],
281 expiration=0,
282 )
283
284 async def get_audio_stream(
285 self,
286 streamdetails: StreamDetails,
287 seek_position: int = 0,
288 ) -> AsyncGenerator[bytes]:
289 """
290 Yield raw PCM from the go-librespot daemon's stdout for the live AudioSource.
291
292 We rate-pace the read at the source's native rate so go-librespot (whose
293 pipe backend is not realtime-paced) is back-pressured to ~realtime and
294 stays only a fraction of a second ahead. Without this, on the ffmpeg
295 resample path (when the player's PCM format differs from ours) nothing
296 throttles the daemon â it races seconds ahead and pause/skip lag badly.
297
298 When playback pauses the daemon stops writing PCM; we then end the stream
299 (clean EOF) so the consuming player leaves the playing state. The next
300 ``playing`` event re-triggers playback. ``seek_position`` is ignored â
301 seeking is handled upstream by Spotify, not by replaying the bytestream.
302 """
303 if streamdetails.item_id != AUDIO_SOURCE_ID:
304 raise MediaNotFoundError(f"Unknown AudioSource: {streamdetails.item_id}")
305 proc = self._proc
306 if proc is None:
307 raise AudioError("Spotify Connect daemon is not running")
308 fmt = self._decoded_audio_format
309 bytes_per_second = fmt.sample_rate * fmt.channels * (fmt.bit_depth // 8)
310 loop = self.mass.loop
311 start = loop.time()
312 streamed = 0
313 while True:
314 try:
315 chunk = await asyncio.wait_for(
316 proc.read(STREAM_READ_CHUNK), timeout=PAUSE_EOF_TIMEOUT_S
317 )
318 except TimeoutError:
319 # No PCM for a while. If playback is no longer active (paused /
320 # stopped / session gone) end the stream so the player goes idle;
321 # a brief buffering gap while still playing just keeps waiting.
322 if not self._playing:
323 return
324 continue
325 if not chunk:
326 return # daemon stdout closed (process exited / restarting)
327 yield chunk
328 # pace at native rate â back-pressure the daemon to ~realtime
329 streamed += len(chunk)
330 ahead = streamed / bytes_per_second - (loop.time() - start)
331 if ahead > 0:
332 await asyncio.sleep(ahead)
333
334 async def on_source_selected(
335 self,
336 source_id: str,
337 player_id: str,
338 queue_id: str,
339 stream_session_id: str,
340 ) -> None:
341 """Handle callback when this AudioSource has been selected/started on a player."""
342 if source_id != AUDIO_SOURCE_ID or not player_id:
343 return
344
345 # Cache the queue_id (== user-facing MA player) rather than the
346 # protocol-level player_id. Some protocol players are ephemeral bridges
347 # whose ID is invalid for play_media / queue lookups once torn down.
348 active_player_id = queue_id
349
350 # If there's already a different active player, kick it out. The claim
351 # below replaces the previous queue's claim; the prior stream's
352 # on_source_unselected may fire later, but its session-id guard keeps it
353 # from clobbering the new claim.
354 if self._active_player_id and self._active_player_id != active_player_id:
355 prev_player_id = self._active_player_id
356 self.logger.info(
357 "Source selected on player %s, stopping playback on %s",
358 active_player_id,
359 prev_player_id,
360 )
361 try:
362 await self.mass.players.cmd_stop(prev_player_id)
363 except Exception as err:
364 self.logger.debug("Failed to stop previous player %s: %s", prev_player_id, err)
365
366 # Claim ownership for this queue.
367 self._in_use_by_queue = queue_id
368 self._active_session_id = stream_session_id
369 self._active_player_id = active_player_id
370 self.logger.debug("Active player set to: %s", active_player_id)
371
372 # Only persist the selected player as the new default if not in auto mode
373 if self._default_player_id != PLAYER_ID_AUTO:
374 self._save_last_player_id(active_player_id)
375
376 # Externally triggered: the daemon is already playing â nothing to do.
377 # Otherwise acquire playback, then confirm it actually started.
378 if not self._playing:
379 assert self._client is not None
380 try:
381 if self._spotify_session_active:
382 # Still the active Spotify device (just paused) â resume.
383 await self._client.resume()
384 elif self._last_context_uri:
385 # The user moved the active device away in the Spotify app.
386 # Take playback back by (re)starting the last context on us,
387 # which makes go-librespot the active device again. The track
388 # restarts from its beginning (go-librespot has no
389 # resume-at-position play call).
390 self.logger.info("Taking Spotify playback back to Music Assistant")
391 await self._client.play(
392 self._last_context_uri, skip_to_uri=self._last_track_uri
393 )
394 else:
395 raise self._not_active_error()
396 except AudioError:
397 raise
398 except Exception as err:
399 raise AudioError(f"Failed to acquire Spotify Connect: {err}") from err
400 if not await self._wait_for_playing():
401 raise self._not_active_error()
402
403 # go-librespot reports 100% volume until told otherwise (with
404 # external_volume it ignores initial_volume); push the player's volume
405 # so the Spotify app's absolute volume commands start from the real level.
406 await self._sync_player_volume_to_spotify(active_player_id)
407
408 async def on_source_unselected(
409 self, source_id: str, queue_id: str, stream_session_id: str
410 ) -> None:
411 """Release the queue-scoped exclusive claim when MA tears down the stream."""
412 if source_id != AUDIO_SOURCE_ID:
413 return
414 # Reject stale callbacks: only release if this is still the active
415 # session. A queue_id check alone is not sufficient â same-queue
416 # reconnects would otherwise let an old request's late callback clear
417 # the live claim of the new stream.
418 if self._active_session_id != stream_session_id:
419 return
420 self._active_session_id = None
421 if self._in_use_by_queue == queue_id:
422 self._in_use_by_queue = None
423
424 async def on_source_control(
425 self,
426 source_id: str,
427 action: SourceControl,
428 value: int | None = None,
429 ) -> None:
430 """Proxy playback control commands to go-librespot's REST API."""
431 if source_id != AUDIO_SOURCE_ID:
432 return
433 if not self._playing and not self._spotify_session_active:
434 raise self._not_active_error()
435 assert self._client is not None
436 try:
437 if action == SourceControl.PLAY:
438 await self._client.resume()
439 elif action == SourceControl.PAUSE:
440 await self._client.pause()
441 elif action == SourceControl.NEXT:
442 await self._client.next()
443 elif action == SourceControl.PREVIOUS:
444 await self._client.prev()
445 elif action == SourceControl.SEEK and value is not None:
446 await self._client.seek(value * 1000)
447 except Exception as err:
448 self.logger.warning("Failed to send %s command to go-librespot: %s", action, err)
449 raise
450
451 async def on_volume_change(self, source_id: str, volume: int) -> None:
452 """Sync the Spotify app's volume slider with the player's new volume."""
453 if source_id != AUDIO_SOURCE_ID:
454 return
455 if not self._playing and not self._spotify_session_active:
456 raise self._not_active_error()
457 # Prevent ping-pong: only push if the value actually changed from what we
458 # last sent to / received from the daemon.
459 if self._last_volume_sent == volume:
460 return
461 try:
462 await self._push_volume_to_daemon(volume)
463 except Exception as err:
464 self.logger.warning("Failed to send volume command to go-librespot: %s", err)
465 raise
466
467 def _not_active_error(self) -> AudioError:
468 """Build the localized 'not the active Spotify device' error, naming this device."""
469 return AudioError(
470 NOT_ACTIVE_DEVICE_MESSAGE.format(self._publish_name),
471 translation_key="not_active_device",
472 translation_args=[self._publish_name],
473 translation_owner=self.translation_owner,
474 )
475
476 def _build_audio_source(self) -> AudioSource:
477 """
478 Construct the AudioSource MediaItem.
479
480 go-librespot exposes a full REST control API, so play / pause / seek /
481 next / previous are always available while a session is active â the
482 capability flags are static (no dependency on the Spotify Web API).
483 """
484 return AudioSource(
485 item_id=AUDIO_SOURCE_ID,
486 provider=self.instance_id,
487 name=self.name,
488 provider_mappings={
489 ProviderMapping(
490 item_id=AUDIO_SOURCE_ID,
491 provider_domain=self.domain,
492 provider_instance=self.instance_id,
493 audio_format=self._audio_format,
494 )
495 },
496 can_play_pause=True,
497 can_seek=True,
498 can_next_previous=True,
499 exclusive=True,
500 allow_external_trigger=True,
501 # Cold-start from MA is unreliable (Spotify needs an existing
502 # playback context), so only allow external entry via the Spotify app.
503 can_initiate=False,
504 )
505
506 def _get_target_player_id(self) -> str | None:
507 """
508 Determine the target player ID for playback.
509
510 Priority: an explicitly selected player; else (auto) a currently playing
511 player then the first available; else the configured default player.
512
513 :return: The player ID to use for playback, or None if none available.
514 """
515 if self._active_player_id:
516 if self.mass.players.get_player(self._active_player_id):
517 return self._active_player_id
518 self._active_player_id = None
519
520 if self._default_player_id == PLAYER_ID_AUTO:
521 all_players = list(self.mass.players.all_players(False, False))
522 for player in all_players:
523 if player.state.playback_state == PlaybackState.PLAYING:
524 self.logger.debug("Auto-selecting playing player: %s", player.display_name)
525 return player.player_id
526 if all_players:
527 first_player = all_players[0]
528 self.logger.debug(
529 "Auto-selecting first available player: %s", first_player.display_name
530 )
531 return first_player.player_id
532 return None
533
534 if self.mass.players.get_player(self._default_player_id):
535 return self._default_player_id
536 self.logger.warning(
537 "Configured default player '%s' no longer exists", self._default_player_id
538 )
539 return None
540
541 async def _wait_for_playing(self, timeout: float = PLAYBACK_START_TIMEOUT_S) -> bool:
542 """
543 Wait up to ``timeout`` seconds for the daemon to report it is playing.
544
545 :param timeout: Maximum seconds to wait.
546 :return: True once playback is confirmed, False if the timeout elapses.
547 """
548 deadline = self.mass.loop.time() + timeout
549 while True:
550 if self._playing:
551 return True
552 if self.mass.loop.time() >= deadline:
553 return False
554 await asyncio.sleep(0.1)
555
556 def _cancel_pending_play_media(self) -> None:
557 """Cancel any pending deferred play_media trigger."""
558 task = self._pending_play_media_task
559 if task is not None and not task.done():
560 task.cancel()
561 self._pending_play_media_task = None
562
563 async def _deferred_play_media_fire(self) -> None:
564 """
565 Trigger play_media after a short debounce.
566
567 The daemon can emit a stale 'playing' from a dying session just before it
568 reconnects; acting on it immediately would start a stream for a session
569 that is about to be replaced. Debouncing â and cancelling the task on a
570 later 'paused' / 'stopped' / 'active' event â avoids a playâstopâreplay loop.
571 """
572 try:
573 await asyncio.sleep(PLAY_MEDIA_DEBOUNCE_S)
574 except asyncio.CancelledError:
575 return
576 if not self._playing or self._in_use_by_queue:
577 return
578 target_player_id = self._get_target_player_id()
579 if not target_player_id:
580 self.logger.warning(
581 "Spotify Connect playback started but no player available. "
582 "Select this source on a player to start playback."
583 )
584 return
585 self.logger.info(
586 "Starting Spotify Connect playback [%s] on player %s",
587 self.instance_id,
588 target_player_id,
589 )
590 self._active_player_id = target_player_id
591 self.mass.create_task(
592 self.mass.player_queues.play_media(target_player_id, str(self._audio_source.uri))
593 )
594
595 def _clear_active_player(self) -> None:
596 """Clear the active player and reset playback state when a session ends."""
597 prev_player_id = self._active_player_id
598 self._active_player_id = None
599 self._in_use_by_queue = None
600 self._active_session_id = None
601 self._playing = False
602 if prev_player_id:
603 self.logger.debug("Playback ended on player %s, clearing active player", prev_player_id)
604 self.mass.players.trigger_player_update(prev_player_id)
605
606 def _save_last_player_id(self, player_id: str) -> None:
607 """Persist the selected player ID as the new default."""
608 if self._default_player_id == player_id:
609 return
610 try:
611 self._update_setup_data(CONF_MASS_PLAYER_ID, player_id)
612 self._default_player_id = player_id
613 except Exception as err:
614 self.logger.debug("Failed to persist player ID: %s", err)
615
616 def _write_config(self, source_ip: str | None) -> None:
617 """
618 Write the go-librespot ``config.yml`` for this instance.
619
620 go-librespot reads a YAML config; JSON is valid YAML, so we emit JSON to
621 sidestep an extra dependency and any string-quoting pitfalls (the device
622 name is user-provided). The config dir doubles as the credential/device
623 cache so the Spotify Connect device stays paired across restarts.
624
625 :param source_ip: Local address of the player-facing interface, or None to
626 advertise the Spotify Connect device on all interfaces.
627 """
628 os.makedirs(self.cache_dir, exist_ok=True)
629 config: dict[str, Any] = {
630 "device_name": self._publish_name,
631 "device_type": "speaker",
632 "device_id": generate_device_id(self.instance_id),
633 "bitrate": 320,
634 "audio_backend": "pipe",
635 # write decoded PCM to the daemon's stdout, which we capture and
636 # forward (the process pipe is always attached, so the daemon's
637 # non-blocking pipe open never fails for lack of a reader). s16le is
638 # the Spotify source representation; MA converts it per player.
639 "audio_output_pipe": "/dev/stdout",
640 "audio_output_pipe_format": "s16le",
641 # external_volume: don't let go-librespot attenuate the PCM â MA / the
642 # target player owns the actual volume. We still receive 'volume'
643 # events and push volume back so the Spotify app slider stays in sync.
644 # No initial_volume: go-librespot ignores it with external_volume set;
645 # _sync_player_volume_to_spotify pushes the player's volume instead.
646 "external_volume": True,
647 "volume_steps": VOLUME_STEPS,
648 "zeroconf_enabled": True,
649 "credentials": {"type": "zeroconf", "zeroconf": {"persist_credentials": True}},
650 "server": {"enabled": True, "address": "127.0.0.1", "port": self._api_port},
651 }
652 # Advertise the Spotify Connect device only on the interface the streams
653 # server binds to, so it lands on the right network on multi-homed hosts.
654 # go-librespot selects advertise interfaces by name, so map the IP to one.
655 if source_ip:
656 if iface_name := interface_name_for_ip(source_ip):
657 config["zeroconf_interfaces_to_advertise"] = [iface_name]
658 else:
659 self.logger.debug(
660 "No interface found for stream bind IP %s; advertising on all interfaces",
661 source_ip,
662 )
663 config_file = os.path.join(self.cache_dir, "config.yml")
664 with open(config_file, "w", encoding="utf-8") as fileobj:
665 json.dump(config, fileobj, indent=2)
666
667 async def _daemon_runner(self) -> None:
668 """Run and supervise the go-librespot daemon, restarting it if it exits."""
669 assert self._binary
670 assert self._client
671 # Loop forever; unload() cancels this task and the explicit stop-check below
672 # handles a graceful exit without a restart.
673 while True:
674 # If the API port was taken while the daemon was down, move to a
675 # fresh port instead of crash-looping on a bind error.
676 if await is_port_in_use(self._api_port, host="127.0.0.1"):
677 self._api_port = await select_free_port(
678 API_PORT_RANGE_START, API_PORT_RANGE_END, host="127.0.0.1"
679 )
680 self._client.base_url = f"http://127.0.0.1:{self._api_port}"
681 self.logger.warning(
682 "API port in use by another process; switching to port %s", self._api_port
683 )
684 self._write_config(await self.mass.streams.get_source_ip())
685 proc: AsyncProcess | None = None
686 try:
687 # stdout carries the decoded PCM (audio_output_pipe=/dev/stdout) and
688 # is consumed by get_audio_stream; stderr carries the daemon's logs,
689 # read here. Because the process pipe owns stdout from spawn there is
690 # always a reader, so go-librespot's non-blocking pipe open never
691 # fails for lack of a consumer.
692 self._proc = proc = AsyncProcess(
693 [self._binary, "--config_dir", self.cache_dir],
694 stdout=True,
695 stderr=True,
696 name=f"go-librespot[{self.name}]",
697 )
698 await proc.start()
699 self.logger.info("Started Spotify Connect background daemon [%s]", self.name)
700 async for line in proc.iter_stderr():
701 self.logger.debug("[%s] %s", self.name, line)
702 except asyncio.CancelledError:
703 raise
704 except Exception as err:
705 self.logger.warning("go-librespot daemon error [%s]: %s", self.name, err)
706 finally:
707 if proc:
708 await proc.close()
709 # The daemon â and thus the Spotify session â is gone. Reset session
710 # state so a dead/restarting daemon isn't treated as active and
711 # controllable; a fresh 'active' event re-establishes it on reconnect.
712 self._proc = None
713 self._playing = False
714 self._spotify_session_active = False
715 if self._stop_called:
716 break
717 self.logger.info("Spotify Connect background daemon stopped for %s", self.name)
718 self._restart_error_count += 1
719 if self._restart_error_count >= 5:
720 self.unload_with_error("go-librespot daemon failed to start multiple times.")
721 return
722 await asyncio.sleep(2)
723
724 async def _events_runner(self) -> None:
725 """Keep the go-librespot events websocket connected, reconnecting as needed."""
726 assert self._client is not None
727 while not self._stop_called:
728 try:
729 if not await self._client.wait_until_ready():
730 await asyncio.sleep(2)
731 continue
732 # A live websocket means the daemon is healthy: reset the restart
733 # backoff counter the daemon supervisor uses.
734 self._restart_error_count = 0
735 await self._client.listen_events(self._handle_event)
736 except asyncio.CancelledError:
737 raise
738 except Exception as err:
739 self.logger.debug("go-librespot events websocket dropped: %s", err)
740 if not self._stop_called:
741 await asyncio.sleep(2)
742
743 async def _handle_event(self, event_type: str, data: dict[str, Any]) -> None:
744 """Dispatch a single go-librespot websocket event."""
745 self.logger.debug("Received event [%s]: %s %s", self.name, event_type, data)
746 # Remember the latest context/track so we can take playback back if the
747 # user moves the active device away in the Spotify app (see on_source_selected).
748 if context_uri := data.get("context_uri"):
749 self._last_context_uri = context_uri
750 if track_uri := data.get("uri"):
751 self._last_track_uri = track_uri
752
753 if event_type == "active":
754 self._spotify_session_active = True
755 self._last_session_active_time = time.time()
756 # A (re)activation supersedes any deferred play_media scheduled from a
757 # previous session's stale 'playing'; the fresh 'playing' that follows
758 # schedules a new one.
759 self._cancel_pending_play_media()
760 self.logger.info("Spotify Connect session active for %s", self.name)
761 # A new session starts at the daemon's 100% volume default; push the
762 # target player's volume so the Spotify app's slider is correct from
763 # device selection, before any playback starts.
764 if player_id := self._get_target_player_id():
765 await self._sync_player_volume_to_spotify(player_id)
766 elif event_type == "inactive":
767 self.logger.info("Spotify Connect session inactive for %s", self.name)
768 self._spotify_session_active = False
769 prev_player_id = self._active_player_id
770 self._clear_active_player()
771 if prev_player_id:
772 self.mass.create_task(self.mass.players.cmd_stop(prev_player_id))
773 return
774 elif event_type == "playing":
775 self._playing = True
776 # Externally triggered playback: kick a play_media on the target MA
777 # player so the audio reaches a speaker. Deferred so a rapid
778 # playing/active burst from a reconnecting session can cancel it.
779 if not self._in_use_by_queue and (
780 self._pending_play_media_task is None or self._pending_play_media_task.done()
781 ):
782 self._pending_play_media_task = self.mass.create_task(
783 self._deferred_play_media_fire()
784 )
785 elif event_type in ("paused", "stopped"):
786 self._playing = False
787 # A pause/stop is the definitive "don't start": cancel a deferred fire
788 # from a now-stale 'playing'. The active get_audio_stream sees the PCM
789 # stop and ends the stream (clean EOF), so the player leaves the playing
790 # state; the next 'playing' event re-fires play_media to resume.
791 self._cancel_pending_play_media()
792
793 self._apply_metadata(event_type, data)
794
795 if event_type == "volume":
796 await self._handle_volume_event(data)
797
798 # push metadata update to the active queue item's streamdetails
799 if self._in_use_by_queue:
800 self.mass.streams.update_stream_metadata(
801 self._in_use_by_queue,
802 AUDIO_SOURCE_ID,
803 self.instance_id,
804 self._stream_metadata,
805 )
806
807 def _apply_metadata(self, event_type: str, data: dict[str, Any]) -> None:
808 """Update the live StreamMetadata from a 'metadata' or 'seek' event."""
809 if event_type == "metadata":
810 self._stream_metadata.uri = data.get("uri")
811 if name := data.get("name"):
812 self._stream_metadata.title = name
813 artists = data.get("artist_names") or []
814 self._stream_metadata.artist = artists[0] if artists else None
815 self._stream_metadata.album = data.get("album_name")
816 self._stream_metadata.image_url = data.get("album_cover_url")
817 self._stream_metadata.description = None
818 duration_ms = data.get("duration")
819 self._stream_metadata.duration = duration_ms // 1000 if duration_ms else None
820 self._stream_metadata.elapsed_time = int(data.get("position", 0)) // 1000
821 self._stream_metadata.elapsed_time_last_updated = int(time.time())
822 elif event_type == "seek":
823 self._stream_metadata.elapsed_time = int(data.get("position", 0)) // 1000
824 self._stream_metadata.elapsed_time_last_updated = int(time.time())
825
826 async def _handle_volume_event(self, data: dict[str, Any]) -> None:
827 """Apply a Spotify-side volume change to the linked MA player."""
828 value = data.get("value")
829 max_value = data.get("max") or VOLUME_STEPS
830 if value is None:
831 return
832 volume = int(int(value) / int(max_value) * 100)
833 # Ignore our own echo: go-librespot emits a 'volume' event for the value we
834 # just pushed in on_volume_change; re-applying it would ping-pong.
835 if volume == self._last_volume_sent:
836 return
837 # Ignore the volume that librespot reports right after a session becomes
838 # active â the player's own volume should win in that window.
839 if time.time() - self._last_session_active_time < INITIAL_VOLUME_GRACE_S:
840 self.logger.debug("Ignoring initial volume_changed event after session active")
841 return
842 if not self._in_use_by_queue:
843 return
844 previous_volume = self._last_volume_sent
845 self._last_volume_sent = volume
846 try:
847 await self.mass.players.cmd_volume_set(self._in_use_by_queue, volume)
848 except Exception as err:
849 # Volume sync is best-effort: the player may not support volume, or the
850 # command may fail. Restore the cached value so a retry isn't wrongly
851 # deduped, and never let it bubble up and drop the events loop.
852 self._last_volume_sent = previous_volume
853 self.logger.debug("Could not set volume on %s: %s", self._in_use_by_queue, err)
854
855 async def _sync_player_volume_to_spotify(self, player_id: str) -> None:
856 """
857 Push a player's current volume to go-librespot (best-effort).
858
859 :param player_id: The MA player whose volume to push.
860 """
861 player = self.mass.players.get_player(player_id)
862 if player is None or player.state.volume_level is None:
863 return
864 # clamp: the logical volume can be out of range until volume limit
865 # enforcement runs
866 volume = max(0, min(100, player.state.volume_level))
867 # No dedupe against _last_volume_sent here: it holds the last value
868 # exchanged with the daemon, not the daemon's current volume, which
869 # resets to its 100% default on a new session or daemon restart.
870 try:
871 await self._push_volume_to_daemon(volume)
872 except Exception as err:
873 self.logger.debug("Failed to sync player volume to Spotify: %s", err)
874
875 async def _push_volume_to_daemon(self, volume: int) -> None:
876 """
877 Send an absolute 0-100 volume to go-librespot.
878
879 :param volume: Volume percentage to send.
880 :raises Exception: If the request to the daemon fails.
881 """
882 assert self._client is not None
883 previous_volume = self._last_volume_sent
884 # Record BEFORE the call: go-librespot echoes a 'volume' event back, and
885 # that echo can arrive over the WS while we're still awaiting set_volume.
886 # Recording up front lets _handle_volume_event dedupe it instead of
887 # bouncing it back as a player volume change.
888 self._last_volume_sent = volume
889 try:
890 await self._client.set_volume(round(volume / 100 * VOLUME_STEPS))
891 except Exception:
892 # restore on failure so a retry of this value isn't wrongly deduped
893 self._last_volume_sent = previous_volume
894 raise
895