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