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