/
/
1"""
2Music Assistant Snapcast source stream.
3
4This module implements a Music Assistant-managed Snapcast stream that is exposed to the
5Snapcast server as a TCP source. The stream is produced by running an FFmpeg pipeline
6which pulls audio from Music Assistant and pushes it to the Snapcast source URI.
7
8Optionally, a Unix socket server can be started to provide a control channel for a
9Snapcast control script (used by the built-in Snapcast server integration).
10"""
11
12from __future__ import annotations
13
14import asyncio
15import random
16import time
17import urllib.parse
18from contextlib import suppress
19from typing import TYPE_CHECKING, cast
20
21from music_assistant_models.enums import ContentType
22from music_assistant_models.media_items import AudioFormat
23
24from music_assistant.controllers.streams.audio_processing import (
25 AudioOutputPlan,
26 get_media_session_id,
27)
28from music_assistant.helpers.ffmpeg import FFMpeg
29from music_assistant.providers.snapcast.socket_server import SnapcastSocketServer
30
31from .constants import (
32 CONTROL_SOCKET_PATH_TEMPLATE,
33 snapcast_sampleformat_query,
34)
35
36if TYPE_CHECKING:
37 from music_assistant_models.player import PlayerMedia
38
39 from music_assistant.helpers.dsp import ComplexFilter
40
41 from .provider import SnapCastProvider
42 from .snap_cntrl_proto import SnapstreamProto
43
44
45class SnapcastMAStream:
46 """
47 A Music Assistant-managed Snapcast stream.
48
49 The stream lifecycle is:
50 - setup: ensure required server resources exist (Snapcast source, optional socket server)
51 - start_stream: start the FFmpeg streaming task
52 - request_stop_stream / wait_for_stopped: stop streaming and await termination
53 - destroy: stop streaming, remove Snapcast source, and stop ancillary services
54
55 If `cntrl_queue_id` is provided, a Unix socket server is started to allow a Snapcast
56 control script to communicate with Music Assistant.
57 """
58
59 def __init__(
60 self,
61 provider: SnapCastProvider,
62 media: PlayerMedia,
63 stream_name: str,
64 source_id: str | None = None,
65 filter_settings_owner: str | None = None,
66 use_cntrl_script: bool = False,
67 destroy_on_stop: bool = False,
68 ) -> None:
69 """
70 Initialize the stream.
71
72 Args:
73 provider: The Snapcast provider instance.
74 media: The media item to stream.
75 stream_name: Name used to register the stream on the Snapcast server.
76 cntrl_queue_id: If set, enables the control socket server used by the control script.
77 filter_settings_owner: Player/entity id used to fetch DSP/filter parameters.
78 destroy_on_stop: If true, delete this MA stream once streaming stops.
79 """
80 self.media = media
81 self.stream_name = stream_name
82 self.snap_stream: SnapstreamProto | None = None
83
84 self._provider = provider
85 self._logger = provider.logger
86 self._mass = provider.mass
87 self._source_id = source_id
88 self._use_cntrl_script = use_cntrl_script
89 self._cntrl_queue_id = source_id if use_cntrl_script else None
90 self._filter_settings_owner = filter_settings_owner
91 self._destroy_on_stop = destroy_on_stop
92
93 self._lifecycle_lock = asyncio.Lock()
94 self._destroyed = False
95 self._setup_done = False
96 self._is_streaming = False
97 self._restart_requested: bool = False
98 self._stop_requested: bool = False
99 self._streaming_started_at: float | None = None
100 self._output_plan: AudioOutputPlan | None = None
101
102 self._socket_server: SnapcastSocketServer | None = None
103 self._socket_path: str | None = None
104 self._streamer_task: asyncio.Task[None] | None = None
105 self._stop_streamer_evt = asyncio.Event()
106 self._streamer_started_evt = asyncio.Event()
107 self._stop_timer: asyncio.Handle | None = None
108 self._stop_timer_started_at: float | None = None
109 self._filter_settings: list[str | ComplexFilter] | None = None
110
111 @property
112 def source_id(self) -> str | None:
113 """Return the source id this stream was created for."""
114 return self._source_id
115
116 @property
117 def stream_id(self) -> str | None:
118 """Return the Snapcast stream identifier, if registered."""
119 if self.snap_stream:
120 return self.snap_stream.identifier
121 return None
122
123 @property
124 def is_streaming(self) -> bool:
125 """Return True if the FFmpeg streaming task is currently running."""
126 return self._is_streaming
127
128 @property
129 def playback_started_at(self) -> float | None:
130 """
131 Return when the playback started at the clients.
132
133 return The (UTC) timestamp when the playback was started on the client
134 or None if not started yet or not streaming.
135 """
136 if self._streaming_started_at is None:
137 return None
138 if self._provider._use_builtin_server:
139 buffer_ms = self._provider._snapcast_server_buffer_size
140 if time.time() - self._streaming_started_at < buffer_ms / 1000.0:
141 return None
142 return self._streaming_started_at + buffer_ms / 1000.0
143 return self._streaming_started_at
144
145 async def setup(self) -> None:
146 """
147 Prepare the Snapcast stream resources.
148
149 Ensures a Snapcast source exists on the server. If `cntrl_queue_id` is set,
150 also starts the Unix socket server used by the control script.
151 """
152 async with self._lifecycle_lock:
153 if self._destroyed:
154 raise RuntimeError("Session is destroyed")
155 if self._setup_done:
156 return
157 if self._provider._snapserver is None:
158 raise RuntimeError("Snapserver needs to be setup first")
159
160 if self._cntrl_queue_id:
161 await self._start_socket_server()
162
163 await self._register_tcp_server_source()
164 self._setup_done = True
165
166 async def destroy(self) -> None:
167 """
168 Stop streaming and tear down all resources.
169
170 This stops the streamer task (if running), removes the Snapcast source,
171 and stops the optional control socket server.
172 """
173 async with self._lifecycle_lock:
174 if self._destroyed:
175 return
176 self._destroyed = True
177
178 self.request_stop_stream()
179 await self.wait_for_stopped()
180 await self._remove_snap_source()
181 await self._stop_socket_server()
182
183 async def start_stream(self, allow_restart: bool = False) -> None:
184 """
185 Start streaming the configured media to the Snapcast source.
186
187 Raises:
188 RuntimeError: If the streamer task is already running.
189 """
190 await self.setup()
191 async with self._lifecycle_lock:
192 if self._streamer_task and not self._streamer_task.done():
193 if not allow_restart:
194 raise RuntimeError("streamer already running")
195 if self._stop_requested or self._stop_streamer_evt.is_set():
196 # stop in flight; _on_streamer_done will start the fresh run
197 self._restart_requested = True
198 else:
199 self._restart_if_running()
200 return
201
202 self._stop_requested = False
203 self._restart_requested = False
204 self._stop_streamer_evt.clear()
205 self._streamer_started_evt.clear()
206 self._streamer_task = self._mass.create_task(self._streamer_task_impl())
207 self._streamer_task.add_done_callback(self._on_streamer_done)
208
209 async def wait_for_started(self, timeout_sec: float | None = None) -> None:
210 """
211 Wait until the streamer task signals it has started.
212
213 Args:
214 timeout_sec: Optional timeout in seconds.
215 """
216 try:
217 await asyncio.wait_for(self._streamer_started_evt.wait(), timeout_sec)
218 except TimeoutError:
219 self._logger.warning(
220 "Timeout waiting for stream %s to start; Canceling...",
221 self.stream_name,
222 )
223
224 def update_media(self, media: PlayerMedia) -> None:
225 """Update the media to play and restart the stream if required."""
226 if media != self.media:
227 self.media = media
228 self._restart_if_running()
229
230 def update_filter_settings(self, from_player: str | None = None) -> None:
231 """Update the filter setting."""
232 take_from = from_player or self._filter_settings_owner
233 if not take_from:
234 raise RuntimeError("No player provided to read filter settings from.")
235 stream_format = self._provider.stream_audio_format
236 output_format = self._get_transport_format()
237 self._output_plan = self._mass.streams.audio.get_player_output_plan(
238 take_from,
239 stream_format,
240 output_format,
241 handoff_format=stream_format,
242 )
243 self._register_output_plan()
244 new_settings = self._output_plan.filter_params
245 if from_player:
246 self._filter_settings_owner = from_player
247 if new_settings != self._filter_settings:
248 self._restart_if_running()
249
250 def request_stop_stream(self) -> None:
251 """
252 Request the streamer task to stop.
253
254 This is cooperative: the streamer task will stop when it observes the stop event.
255 Any pending inactivity stop timer is canceled.
256 """
257 self._stop_requested = True
258 self._restart_requested = False # explicit stop cancels any pending restart
259 self._stop_streamer_evt.set()
260
261 self._stop_timer_started_at = None
262 if self._stop_timer:
263 self._stop_timer.cancel()
264
265 def set_in_use(self, in_use: bool) -> None:
266 """
267 Mark the stream as in-use or idle.
268
269 When marked idle, a delayed stop is scheduled. When marked in-use, any pending
270 delayed stop is canceled.
271 """
272 if in_use:
273 self._stop_timer_started_at = None
274 if self._stop_timer:
275 self._stop_timer.cancel()
276 elif self._stop_timer_started_at is None and not self._stop_requested:
277 self._stop_timer_started_at = self._mass.loop.time()
278 self._stop_timer = self._mass.loop.call_later(3.0, self.request_stop_stream)
279
280 async def wait_for_stopped(self, timeout_sec: float | None = None) -> None:
281 """
282 Wait for the streamer task to finish.
283
284 If the task does not finish within the timeout, it is canceled and awaited.
285
286 Args:
287 timeout_sec: Optional timeout in seconds.
288 """
289 curr_task = self._streamer_task
290 if not curr_task:
291 return
292 try:
293 await asyncio.wait_for(curr_task, timeout_sec)
294 except asyncio.CancelledError:
295 self._logger.warning("Streamer task got canceled")
296 except TimeoutError:
297 self._logger.warning(
298 "Timeout waiting for stream %s to finish; Canceling...",
299 self.stream_name,
300 )
301 curr_task.cancel()
302 await asyncio.gather(curr_task, return_exceptions=True)
303
304 async def _streamer_task_impl(self) -> None:
305 """
306 Streamer task implementation.
307
308 Runs FFmpeg to push audio to the Snapcast TCP source until FFmpeg exits or a stop
309 request is received. After exit, waits briefly for the Snapcast stream to report
310 an idle state.
311 """
312 stream_path = self._snap_get_stream_path()
313 if stream_path is None:
314 raise RuntimeError("The path to stream to is not set")
315
316 self._logger.debug("Start streaming to %s", stream_path)
317 self._stop_streamer_evt.clear()
318 self._streamer_started_evt.clear()
319 if self._filter_settings_owner:
320 stream_format = self._provider.stream_audio_format
321 output_format = self._get_transport_format()
322 self._output_plan = self._mass.streams.audio.get_player_output_plan(
323 self._filter_settings_owner,
324 stream_format,
325 output_format,
326 handoff_format=stream_format,
327 )
328 self._register_output_plan()
329 self._filter_settings = self._output_plan.filter_params
330 stream_format = self._provider.stream_audio_format
331 audio_source = self._mass.streams.get_stream(
332 self.media,
333 stream_format,
334 self._filter_settings_owner,
335 use_flow_stream_buffering=True,
336 )
337 try:
338 async with FFMpeg(
339 audio_input=audio_source,
340 input_format=stream_format,
341 output_format=stream_format,
342 filter_params=self._filter_settings or [],
343 audio_output=stream_path,
344 extra_input_args=["-y", "-re"],
345 ) as ffmpeg_proc:
346 wait_ffmpeg = self._mass.create_task(ffmpeg_proc.wait())
347 wait_stop = self._mass.create_task(self._stop_streamer_evt.wait())
348 self._streaming_started_at = time.time()
349 self._streamer_started_evt.set()
350 self._is_streaming = True
351
352 done, pending = await asyncio.wait(
353 {wait_ffmpeg, wait_stop},
354 return_when=asyncio.FIRST_COMPLETED,
355 )
356
357 if wait_stop in done and wait_ffmpeg not in done:
358 self._logger.debug("Stopping stream %s requested.", self.stream_name)
359 wait_ffmpeg.cancel()
360 await asyncio.gather(wait_ffmpeg, return_exceptions=True)
361 return
362
363 await wait_ffmpeg
364 for t in pending:
365 t.cancel()
366 await asyncio.gather(*pending, return_exceptions=True)
367 except asyncio.CancelledError:
368 self._logger.debug("Snapcast stream %s cancelled", self.stream_name)
369 raise
370 except Exception as err:
371 self._logger.error("Snapcast stream %s error: %s", self.stream_name, err, exc_info=err)
372 raise
373 finally:
374 self._is_streaming = False
375 self._logger.debug("Finished streaming to %s", stream_path)
376 await self._wait_stream_idle()
377
378 async def _wait_stream_idle(self) -> None:
379 """Wait for the Snapcast stream to become idle after streaming ends."""
380 try:
381
382 async def wait_until_idle() -> None:
383 while True:
384 stream_is_idle = False
385 with suppress(KeyError):
386 snap_stream = self._provider._snapserver.stream(self.stream_name)
387 stream_is_idle = snap_stream.status == "idle"
388 if self._mass.closing or stream_is_idle:
389 break
390 await asyncio.sleep(0.25)
391
392 await asyncio.wait_for(wait_until_idle(), timeout=10.0)
393 except TimeoutError:
394 self._logger.warning(
395 "Timeout waiting for stream %s to become idle",
396 self.stream_name,
397 )
398 finally:
399 self._streaming_started_at = None
400
401 def _on_streamer_done(self, t: asyncio.Task[None]) -> None:
402 """Handle streamer task completion and optional cleanup."""
403 restart = False
404 try:
405 t.result()
406 except asyncio.CancelledError:
407 self._logger.debug("Streamer task cancelled: %s", self.stream_name)
408 except Exception:
409 self._logger.exception("Streamer task failed")
410 finally:
411 restart = self._restart_requested and not self._destroyed
412
413 if self._streamer_task is t:
414 self._streamer_task = None
415
416 # reset per-run state
417 self._restart_requested = False
418 self._stop_requested = False
419 self._stop_streamer_evt.clear()
420 self._streamer_started_evt.clear()
421
422 if restart:
423 self._mass.create_task(self._restart_stream_locked())
424 elif self._destroy_on_stop:
425 self._mass.create_task(self._provider.delete_ma_stream(self.stream_name))
426
427 def _restart_if_running(self) -> None:
428 """Request a running stream to restart."""
429 t = self._streamer_task
430 if not t or t.done():
431 return
432
433 if self._stop_requested or self._stop_streamer_evt.is_set():
434 return
435
436 self._restart_requested = True
437 self._stop_requested = True
438 self._stop_streamer_evt.set()
439
440 self._stop_timer_started_at = None
441 if self._stop_timer:
442 self._stop_timer.cancel()
443
444 async def _restart_stream_locked(self) -> None:
445 """Restart the streamer under the lifecycle lock."""
446 async with self._lifecycle_lock:
447 if self._destroyed:
448 return
449 if self._streamer_task and not self._streamer_task.done():
450 return
451
452 # reset state and start a fresh run
453 self._stop_requested = False
454 self._restart_requested = False
455 self._stop_streamer_evt.clear()
456 self._streamer_started_evt.clear()
457
458 self._streamer_task = self._mass.create_task(self._streamer_task_impl())
459 self._streamer_task.add_done_callback(self._on_streamer_done)
460
461 def _get_transport_format(self) -> AudioFormat:
462 """Return the format Snapserver sends to its clients."""
463 stream_format = self._provider.stream_audio_format
464 if self._provider._use_builtin_server:
465 codec_name = str(self._provider._snapcast_server_transport_codec)
466 else:
467 stream_data = self.snap_stream._stream if self.snap_stream else {}
468 uri_data = stream_data.get("uri", {})
469 query_data = uri_data.get("query", {}) if isinstance(uri_data, dict) else {}
470 codec_name = str(query_data.get("codec", "") if isinstance(query_data, dict) else "")
471 codec_name = codec_name.partition(":")[0].lower()
472 pcm_type = ContentType.PCM_S24LE if stream_format.bit_depth == 24 else ContentType.PCM_S16LE
473 content_type, codec_type = {
474 "flac": (ContentType.FLAC, ContentType.FLAC),
475 "ogg": (ContentType.OGG, ContentType.VORBIS),
476 "opus": (ContentType.OPUS, ContentType.OPUS),
477 "pcm": (pcm_type, pcm_type),
478 }.get(codec_name, (ContentType.UNKNOWN, ContentType.UNKNOWN))
479 return AudioFormat(
480 content_type=content_type,
481 codec_type=codec_type,
482 sample_rate=stream_format.sample_rate,
483 bit_depth=stream_format.bit_depth,
484 channels=stream_format.channels,
485 )
486
487 def _register_output_plan(self) -> None:
488 """Register the shared Snapcast path for every connected group member."""
489 queue_id = self.media.source_id
490 session_id = get_media_session_id(self.media)
491 if self._output_plan is None or queue_id is None or session_id is None:
492 return
493 player_ids = set(self._output_plan.output_details.player_ids)
494 if self.snap_stream:
495 for group in self._provider._snapserver.groups:
496 if group.stream != self.snap_stream.identifier:
497 continue
498 for client_id in group.clients:
499 if player_id := self._provider._get_ma_id(client_id):
500 player_ids.add(player_id)
501 for player_id in player_ids:
502 self._mass.streams.audio_processing.update_output(
503 player_id,
504 self._output_plan,
505 queue_id=queue_id,
506 session_id=session_id,
507 )
508
509 def _find_local_stream_by_name(self, name: str) -> SnapstreamProto | None:
510 """
511 Look up a snapserver stream by its name (not id) in the local cache.
512
513 :param name: Stream name to look up.
514 :return: The matching SnapstreamProto, or None if not found.
515 """
516 for s in self._provider._snapserver.streams:
517 if getattr(s, "name", None) == name:
518 return s
519 return None
520
521 def _stream_matches_configured_format(self, stream: SnapstreamProto) -> bool:
522 """
523 Return whether a snapserver stream uses the configured PCM sample format.
524
525 :param stream: Existing snapserver stream that may be adopted.
526 """
527 expected = self._provider.stream_audio_format
528 expected_sampleformat = f"{expected.sample_rate}:{expected.bit_depth}:{expected.channels}"
529 stream_data = getattr(stream, "_stream", None)
530 uri_data = stream_data.get("uri", {}) if isinstance(stream_data, dict) else {}
531 query_data = uri_data.get("query", {}) if isinstance(uri_data, dict) else {}
532 if not isinstance(query_data, dict):
533 return False
534 if str(query_data.get("sampleformat", "")) != expected_sampleformat:
535 return False
536 if expected.bit_depth == 24:
537 packed = str(query_data.get("packed_s24le", "")).lower()
538 return packed in {"1", "true", "yes"}
539 return True
540
541 @staticmethod
542 def _is_name_collision_error(result: object) -> bool:
543 """
544 Detect snapserver's 'Stream with name X already exists' error.
545
546 :param result: Result returned by snapserver.stream_add_stream.
547 :return: True if the result indicates a stream-name collision.
548 """
549 if not isinstance(result, dict):
550 return False
551 data = result.get("data", "")
552 return (
553 isinstance(data, str)
554 and data.startswith("Stream with name")
555 and "already exists" in data
556 )
557
558 def _pick_port_avoiding(self, used: set[int]) -> int | None:
559 """
560 Pick a random TCP source port within the unchanged upstream range.
561
562 Avoids ports already tried in the current retry loop.
563
564 :param used: Set of ports that should not be returned again.
565 :return: A port in [4953, 5153] not in `used`, or None if none could be
566 found within 20 attempts (extremely rare; the range has 201 values).
567 """
568 for _ in range(20):
569 port = random.randint(4953, 4953 + 200)
570 if port not in used:
571 return port
572 return None
573
574 async def _register_tcp_server_source(self) -> None:
575 """Create a Snapcast TCP source for this stream (or reuse an existing one)."""
576 # This runs under the per-stream `_lifecycle_lock` (acquired in setup()), not the
577 # provider-wide `_snapcast_ma_streams_lock`. Adoption is safe because each MA-managed
578 # stream has at most one outstanding setup() call at a time.
579
580 # prefer to reuse existing stream if possible
581 if self.snap_stream:
582 return
583
584 # The control script is used only for music streams in the builtin server
585 extra_args = ""
586 if (cntrl_queue_id := self._cntrl_queue_id) is not None:
587 # Create socket server for control script communication
588 socket_path = self._socket_path
589 if socket_path is None:
590 raise RuntimeError("socket_path needs to be set if cntrl_queue_id is set")
591 extra_args = (
592 f"&controlscript={urllib.parse.quote_plus('control.py')}"
593 f"&controlscriptparams=--queueid={urllib.parse.quote_plus(cntrl_queue_id)}%20"
594 f"--socket={urllib.parse.quote_plus(socket_path)}%20"
595 f"--streamserver-ip={self._mass.streams.publish_ip}%20"
596 f"--streamserver-port={self._mass.streams.publish_port}"
597 )
598
599 tried_ports: set[int] = set()
600 attempts = 50
601 loop_succeeded = False
602 try:
603 while attempts:
604 attempts -= 1
605 port = self._pick_port_avoiding(tried_ports)
606 if port is None:
607 break
608 tried_ports.add(port)
609 result = await self._provider._snapserver.stream_add_stream(
610 # 24-bit requires Snapserver packed_s24le support (snapcast/snapcast#1532)
611 f"tcp://0.0.0.0:{port}?{snapcast_sampleformat_query(self._provider.stream_audio_format)}"
612 f"&idle_threshold={self._provider._snapcast_stream_idle_threshold}"
613 f"{extra_args}&name={self.stream_name}"
614 )
615 if isinstance(result, dict) and "id" in result:
616 self.snap_stream = self._provider._snapserver.stream(result["id"])
617 self.snap_stream.set_callback(self._snap_on_stream_update)
618 loop_succeeded = True
619 return
620
621 if self._is_name_collision_error(result):
622 adopted = self._find_local_stream_by_name(self.stream_name)
623 if adopted is None:
624 # Local cache may be stale (e.g. after an MA restart). Resync.
625 status, _ = await self._provider._snapserver.status()
626 if isinstance(status, dict):
627 self._provider._snapserver.synchronize(status)
628 adopted = self._find_local_stream_by_name(self.stream_name)
629 if adopted is not None:
630 if not self._stream_matches_configured_format(adopted):
631 self._logger.info(
632 "Orphaned snapserver stream %s (name=%s) has a different "
633 "sample format; removing it so it can be recreated",
634 adopted.identifier,
635 self.stream_name,
636 )
637 await self._provider._snapserver.stream_remove_stream(
638 adopted.identifier
639 )
640 continue
641 self._logger.info(
642 "Adopted orphaned snapserver stream %s (name=%s)",
643 adopted.identifier,
644 self.stream_name,
645 )
646 self.snap_stream = adopted
647 self.snap_stream.set_callback(self._snap_on_stream_update)
648 loop_succeeded = True
649 return
650 elif isinstance(result, dict) and result.get("code") == -32603:
651 # Forward-compat hint: snapserver may rephrase the name-collision
652 # error in a future release; surface mismatches in debug logs.
653 self._logger.debug(
654 "snapserver returned internal error not matching "
655 "name-collision pattern: %s",
656 result,
657 )
658
659 # if the port is already taken, the result will be an error
660 self._logger.warning("stream_add_stream failed: %s", result)
661 continue
662 finally:
663 # Invariant: if we leave _register_tcp_server_source without setting
664 # self.snap_stream (retries exhausted OR exception during a retry),
665 # any socket_server we started must be stopped.
666 if not loop_succeeded and self._socket_server:
667 await self._stop_socket_server()
668
669 msg = f"Unable to register snapserver stream {self.stream_name!r} after 50 attempts"
670 raise RuntimeError(msg)
671
672 async def _remove_snap_source(self) -> None:
673 """Remove the Snapcast source created for this stream and detach groups."""
674 if self._mass.closing or self.snap_stream is None:
675 return
676
677 for snap_group in self._provider._snapserver.groups:
678 if snap_group.stream != self.snap_stream.identifier:
679 continue
680 self._logger.debug(f"Set stream of group {snap_group.name} to default.")
681 await snap_group.set_stream("default")
682
683 with suppress(KeyError, AttributeError):
684 snap_stream = self._provider._snapserver.stream(self.stream_name)
685 await self._provider._snapserver.stream_remove_stream(snap_stream.identifier)
686
687 if self._socket_server:
688 await self._stop_socket_server()
689 self._snap_on_stream_update()
690
691 return
692
693 def _snap_get_stream_path(self) -> str | None:
694 """Return the Snapcast TCP URI to stream to."""
695 if self.snap_stream is None:
696 return None
697
698 uri = self.snap_stream._stream.get("uri", {})
699 uri_host = uri.get("host", "")
700 stream_path = self.snap_stream.path or f"tcp://{uri_host}"
701 return stream_path.replace("0.0.0.0", self._provider._snapcast_server_host)
702
703 def _snap_on_stream_update(self, stream: SnapstreamProto | None = None) -> None:
704 """Handle Snapcast stream updates and trigger group member refresh."""
705 if self.snap_stream is None:
706 return
707
708 for snap_group in self._provider._snapserver.groups:
709 if snap_group.stream != self.snap_stream.identifier:
710 continue
711 self._provider.poke_group_members(snap_group)
712 self._register_output_plan()
713
714 async def _start_socket_server(self) -> str:
715 """
716 Get or create a socket server for the given queue.
717
718 :return: The path to the Unix socket.
719 """
720 if self._socket_server:
721 return self._socket_server.socket_path
722
723 if self._cntrl_queue_id is None:
724 raise RuntimeError("Socket server require _cntrl_queue_id to be set")
725
726 socket_path = CONTROL_SOCKET_PATH_TEMPLATE.format(queue_id=self._cntrl_queue_id)
727 socket_server = SnapcastSocketServer(
728 mass=self._mass,
729 queue_id=self._cntrl_queue_id,
730 socket_path=socket_path,
731 streamserver_ip=str(self._mass.streams.publish_ip),
732 streamserver_port=cast("int", self._mass.streams.publish_port),
733 )
734 await socket_server.start()
735 self._socket_server = socket_server
736 self._socket_path = socket_path
737 self._logger.debug(
738 "Created socket server for queue %s at %s", self._cntrl_queue_id, socket_path
739 )
740 return socket_path
741
742 async def _stop_socket_server(self) -> None:
743 """Stop and remove the socket server for the given queue."""
744 if not self._socket_server:
745 return
746
747 await self._socket_server.stop()
748 self._socket_server = None
749 self._logger.debug("Stopped socket server for queue %s", self._cntrl_queue_id)
750