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