/
/
1"""
2Sendspin Source provider implementation.
3
4Exposes every connected Sendspin client with an active `source` role as a Music
5Assistant AudioSource. Decoded PCM from the client is pulled through an
6occupancy-controlled clock bridge so the capture clock never has to match MA's
7consumption clock. See README.md for the design rationale.
8"""
9
10from __future__ import annotations
11
12import asyncio
13import time
14from dataclasses import dataclass, field
15from typing import TYPE_CHECKING, cast
16
17from aiosendspin.audio import AsrcSourceBridge
18from aiosendspin.audio import AudioFormat as SendspinAudioFormat
19from aiosendspin.models.types import role_family
20from aiosendspin.server import (
21 ClientConnectedEvent,
22 ClientDisconnectedEvent,
23 ClientRemovedEvent,
24 SignalState,
25 SourceSignalChangedEvent,
26 SourceStreamEndedEvent,
27 SourceStreamStartedEvent,
28)
29from music_assistant_models.enums import (
30 ContentType,
31 MediaType,
32 QueueOption,
33 StreamType,
34)
35from music_assistant_models.errors import (
36 AudioError,
37 MediaNotFoundError,
38 PlayerCommandFailed,
39 PlayerUnavailableError,
40)
41from music_assistant_models.helpers import create_uri
42from music_assistant_models.media_items import AudioFormat, AudioSource, ProviderMapping
43from music_assistant_models.streamdetails import StreamDetails, StreamMetadata
44
45from music_assistant.models.plugin import PluginProvider
46from music_assistant.providers.sendspin.constants import (
47 CONF_SOURCE_AUTOSTART_TARGET,
48 SOURCE_AUTOSTART_OFF,
49)
50
51from .constants import (
52 AUTOSTART_SIGNAL_ABSENT_HOLD_S,
53 AUTOSTART_SIGNAL_DEBOUNCE_S,
54 CHUNK_DURATION_MS,
55 COLD_START_TIMEOUT_S,
56 CONF_TARGET_LATENCY,
57 DEFAULT_TARGET_LATENCY_MS,
58 OUTPUT_BIT_DEPTH,
59 OUTPUT_CHANNELS,
60 OUTPUT_SAMPLE_RATE,
61 SOURCE_TIMEOUT_S,
62)
63
64if TYPE_CHECKING:
65 from collections.abc import AsyncGenerator, Callable
66
67 from aiosendspin.audio import SourceBridge
68 from aiosendspin.server import (
69 ClientEvent,
70 SendspinClient,
71 SendspinEvent,
72 SendspinServer,
73 SourceStream,
74 )
75 from aiosendspin.server.roles import SourceV1Role
76 from music_assistant_models.config_entries import ProviderConfig
77 from music_assistant_models.enums import ProviderFeature
78 from music_assistant_models.provider import ProviderManifest
79
80 from music_assistant.mass import MusicAssistant
81 from music_assistant.providers.sendspin.provider import SendspinProvider
82
83OUTPUT_FORMAT = AudioFormat(
84 content_type=ContentType.PCM_S16LE,
85 sample_rate=OUTPUT_SAMPLE_RATE,
86 bit_depth=OUTPUT_BIT_DEPTH,
87 channels=OUTPUT_CHANNELS,
88)
89BRIDGE_OUTPUT_FORMAT = SendspinAudioFormat(
90 sample_rate=OUTPUT_SAMPLE_RATE,
91 bit_depth=OUTPUT_BIT_DEPTH,
92 channels=OUTPUT_CHANNELS,
93)
94
95
96@dataclass
97class _SourceSession:
98 """State for one source's active exclusive stream."""
99
100 client_id: str
101 player_id: str
102 owner_player_id: str
103 stream_session_id: str
104 # Retires the prior generator when the same player reclaims the session.
105 generation: int = 0
106 bridge: SourceBridge | None = None
107 ingest_task: asyncio.Task[None] | None = None
108 # Selection time counts toward the source timeout before PCM arrives.
109 last_pcm_monotonic: float = field(default_factory=time.monotonic)
110 pcm_received: asyncio.Event = field(default_factory=asyncio.Event)
111
112
113@dataclass
114class _SourceClientState:
115 """State retained for one source client."""
116
117 session: _SourceSession | None = None
118 unwatch: Callable[[], None] | None = None
119 signal: SignalState | None = None
120 autostart_queue_id: str | None = None
121 autostart_queue_session_id: str | None = None
122 suppressed_autostart_claim: tuple[str, str | None] | None = None
123 selection_lock: asyncio.Lock = field(default_factory=asyncio.Lock)
124
125
126class SendspinSourceProvider(PluginProvider):
127 """Expose Sendspin source-role clients as MA AudioSources."""
128
129 def __init__(
130 self,
131 mass: MusicAssistant,
132 manifest: ProviderManifest,
133 config: ProviderConfig,
134 supported_features: set[ProviderFeature] | None = None,
135 ) -> None:
136 """Initialize the provider."""
137 super().__init__(mass, manifest, config, supported_features)
138 self._clients: dict[str, _SourceClientState] = {}
139 self._server_unsubscribe: Callable[[], None] | None = None
140 self._unloading = False
141
142 async def loaded_in_mass(self) -> None:
143 """Start watching every source client, including ones that are already idle."""
144 await super().loaded_in_mass()
145 self._unloading = False
146 if (sendspin := self._sendspin_provider) is None:
147 return
148 self._server_unsubscribe = sendspin.server_api.add_event_listener(self._on_server_event)
149 for client in sendspin.server_api.connected_clients:
150 self._watch_client(client)
151
152 async def unload(self, is_removed: bool = False) -> None:
153 """Handle unload/close of the provider."""
154 self._unloading = True
155 if self._server_unsubscribe is not None:
156 self._server_unsubscribe()
157 self._server_unsubscribe = None
158 for client_id, state in list(self._clients.items()):
159 self._cancel_pending_autostart(client_id)
160 self._cancel_pending_autostop(client_id)
161 if state.unwatch is not None:
162 state.unwatch()
163 state.unwatch = None
164 await self._teardown_session(client_id)
165 self._clients.clear()
166
167 async def get_audio_sources(self) -> list[AudioSource]:
168 """Return one AudioSource per connected client with an active source role."""
169 if (sendspin := self._sendspin_provider) is None:
170 return []
171 sources: list[AudioSource] = []
172 for client in sendspin.server_api.connected_clients:
173 if self._get_source_role(client) is None:
174 continue
175 info = client.info_or_none
176 name = info.name if info else client.client_id
177 sources.append(
178 AudioSource(
179 item_id=client.client_id,
180 provider=self.instance_id,
181 name=name,
182 provider_mappings={
183 ProviderMapping(
184 item_id=client.client_id,
185 provider_domain=self.domain,
186 provider_instance=self.instance_id,
187 audio_format=OUTPUT_FORMAT,
188 )
189 },
190 can_play_pause=False,
191 can_seek=False,
192 can_next_previous=False,
193 exclusive=True,
194 allow_external_trigger=True,
195 can_initiate=True,
196 )
197 )
198 return sources
199
200 async def get_stream_details(self, item_id: str, media_type: MediaType) -> StreamDetails:
201 """
202 Return StreamDetails for streaming the given source to a queue.
203
204 Side-effect-free: streaming is requested from the client in
205 on_source_selected, which the streams controller fires before this
206 method on the actual stream request (never on queue preload).
207 """
208 client = self._get_client(item_id)
209 if client is None or self._get_source_role(client) is None:
210 raise MediaNotFoundError(f"Unknown or unavailable Sendspin source: {item_id}")
211 info = client.info_or_none
212 return StreamDetails(
213 provider=self.instance_id,
214 item_id=item_id,
215 audio_format=OUTPUT_FORMAT,
216 media_type=media_type,
217 stream_type=StreamType.CUSTOM,
218 stream_metadata=StreamMetadata(title=info.name if info else item_id),
219 )
220
221 async def on_source_selected(
222 self, source_id: str, player_id: str, owner_player_id: str, stream_session_id: str
223 ) -> None:
224 """Claim the source and ask the client to start streaming."""
225 client = self._get_client(source_id)
226 role = self._get_source_role(client) if client else None
227 if client is None or role is None:
228 raise MediaNotFoundError(f"Sendspin source is not connected: {source_id}")
229 state = self._clients.setdefault(source_id, _SourceClientState())
230 # Capture the source's session before waiting so delayed requests keep their
231 # identity: a live source plays on the player, so the queue has no session for it.
232 queue_session_id = self._player_session_id(owner_player_id)
233 try:
234 # Serialize handoffs because stopping the previous player can suspend.
235 async with state.selection_lock:
236 # The client or its role may have changed while waiting for the lock.
237 client = self._get_client(source_id)
238 role = self._get_source_role(client) if client else None
239 if client is None or role is None:
240 raise MediaNotFoundError(f"Sendspin source is not connected: {source_id}")
241 # Reject only the delayed request from an autostart superseded by the user.
242 if state.suppressed_autostart_claim == (owner_player_id, queue_session_id):
243 raise RuntimeError("Superseded autostart stream request")
244 autostart_queue_id = state.autostart_queue_id
245 if autostart_queue_id == owner_player_id:
246 state.autostart_queue_id = None
247 state.autostart_queue_session_id = None
248 self._cancel_pending_autostart(source_id, cancel_running=False)
249 else:
250 if autostart_queue_id is not None:
251 autostart_session_id = (
252 state.autostart_queue_session_id
253 or self._player_session_id(autostart_queue_id)
254 )
255 state.autostart_queue_id = None
256 state.autostart_queue_session_id = None
257 state.suppressed_autostart_claim = (
258 autostart_queue_id,
259 autostart_session_id,
260 )
261 self._cancel_pending_autostart(source_id)
262 if autostart_queue_id is not None:
263 try:
264 # deselect rather than stopping the queue: the source plays on
265 # the player, so stopping the queue would leave the source
266 # published on it while resetting the queue we are preserving
267 await self.mass.players.deselect_source(autostart_queue_id)
268 except (KeyError, PlayerCommandFailed, PlayerUnavailableError) as err:
269 self.logger.debug(
270 "Failed to release autostart player %s: %s",
271 autostart_queue_id,
272 err,
273 )
274 # Keep the running bridge when renderers reopen the URL for the same queue.
275 if (live := state.session) is not None and (
276 live.player_id,
277 live.owner_player_id,
278 ) == (player_id, owner_player_id):
279 live.stream_session_id = stream_session_id
280 live.generation += 1
281 return
282 await self._teardown_session(source_id, superseded_by_player_id=player_id)
283 if self._unloading or self._clients.get(source_id) is not state:
284 raise PlayerUnavailableError(f"Sendspin source is unloading: {source_id}")
285 # Teardown can suspend long enough for the connection role to be recreated.
286 client = self._get_client(source_id)
287 role = self._get_source_role(client) if client else None
288 if client is None or role is None:
289 raise MediaNotFoundError(f"Sendspin source is not connected: {source_id}")
290 state.session = _SourceSession(
291 client_id=source_id,
292 player_id=player_id,
293 owner_player_id=owner_player_id,
294 stream_session_id=stream_session_id,
295 )
296 role.request_start()
297 finally:
298 self._drop_empty_client_state(source_id, state)
299
300 async def on_source_unselected(
301 self, source_id: str, owner_player_id: str, stream_session_id: str
302 ) -> None:
303 """Release the source when MA tears down its stream."""
304 session = self._get_session(source_id)
305 # Reject stale callbacks from superseded same-queue requests.
306 if session is None or session.stream_session_id != stream_session_id:
307 return
308 await self._teardown_session(source_id)
309 if (state := self._clients.get(source_id)) is not None:
310 state.suppressed_autostart_claim = None
311 self._drop_empty_client_state(source_id, state)
312
313 async def get_audio_stream(
314 self, streamdetails: StreamDetails, seek_position: int = 0
315 ) -> AsyncGenerator[bytes]:
316 """
317 Yield fixed-format PCM pulled from the source's clock bridge.
318
319 The pull cadence of this loop is the master clock: the bridge converts
320 the client's drifting capture stream to it and pads silence on underrun,
321 so the stream keeps playing through gaps (an unplugged line-in is silent,
322 not stopped) until SOURCE_TIMEOUT_S passes without source audio. Pulling at
323 the output rate here makes the controller's own realtime pacer a no-op.
324 """
325 session = self._get_session(streamdetails.item_id)
326 if session is None:
327 raise AudioError(f"Sendspin source is not selected: {streamdetails.item_id}")
328 generation = session.generation
329 await self._await_first_audio(session)
330 if self._get_session(session.client_id) is not session or session.generation != generation:
331 return
332 frames_per_chunk = OUTPUT_SAMPLE_RATE * CHUNK_DURATION_MS // 1000
333 period = CHUNK_DURATION_MS / 1000
334 loop = self.mass.loop
335 next_deadline = loop.time()
336 while True:
337 if (
338 self._get_session(session.client_id) is not session
339 or session.generation != generation
340 ):
341 break
342 if time.monotonic() - session.last_pcm_monotonic > SOURCE_TIMEOUT_S:
343 self.logger.info(
344 "No audio from Sendspin source %s for %.0fs, ending stream",
345 session.client_id,
346 SOURCE_TIMEOUT_S,
347 )
348 break
349 if (bridge := session.bridge) is None:
350 break
351 yield bridge.read(frames_per_chunk)
352 next_deadline += period
353 delay = next_deadline - loop.time()
354 if delay > 0:
355 await asyncio.sleep(delay)
356 elif -delay * 1_000_000 > bridge.occupancy_us:
357 # Consumer stalled beyond the buffered audio. Catch-up reads past this
358 # point would fabricate silence into the timeline, so re-anchor instead.
359 next_deadline = loop.time()
360
361 @property
362 def _sendspin_provider(self) -> SendspinProvider | None:
363 return cast("SendspinProvider | None", self.mass.get_provider("sendspin"))
364
365 def _get_client(self, client_id: str) -> SendspinClient | None:
366 if (sendspin := self._sendspin_provider) is None:
367 return None
368 return sendspin.server_api.get_client(client_id)
369
370 @staticmethod
371 def _get_source_role(client: SendspinClient) -> SourceV1Role | None:
372 roles = client.roles_by_family("source")
373 return cast("SourceV1Role", roles[0]) if roles else None
374
375 async def _await_first_audio(self, session: _SourceSession) -> None:
376 """
377 Block until the source actually streams, so a failed acquisition raises.
378
379 The silence-hold only makes sense once audio has flowed: a client that never
380 answers the start command is a broken source, not a quiet one.
381 """
382 try:
383 async with asyncio.timeout(COLD_START_TIMEOUT_S):
384 await session.pcm_received.wait()
385 except TimeoutError:
386 if self._get_session(session.client_id) is not session:
387 return
388 client = self._get_client(session.client_id)
389 info = client.info_or_none if client else None
390 raise AudioError(
391 f"Sendspin source {session.client_id} did not start streaming",
392 translation_key="no_audio",
393 translation_owner=self.translation_owner,
394 translation_args=[info.name if info else session.client_id],
395 ) from None
396
397 def _on_server_event(self, server: SendspinServer, event: SendspinEvent) -> None:
398 if self._unloading:
399 return
400 match event:
401 case ClientConnectedEvent(client_id):
402 # Roles attach after this event, so defer until the next loop turn.
403 self.mass.create_task(self._on_client_connected(client_id), eager_start=False)
404 case ClientRemovedEvent(client_id) | ClientDisconnectedEvent(client_id):
405 self._cancel_pending_autostart(client_id)
406 self._cancel_pending_autostop(client_id)
407 if (state := self._clients.get(client_id)) is not None:
408 state.signal = None
409 state.autostart_queue_id = None
410 state.autostart_queue_session_id = None
411 if state.session is None:
412 state.suppressed_autostart_claim = None
413 if state.unwatch is not None:
414 state.unwatch()
415 state.unwatch = None
416 self._drop_empty_client_state(client_id, state)
417
418 async def _on_client_connected(self, client_id: str) -> None:
419 """Re-arm watching and streaming for a client that just (re)connected."""
420 if self._unloading:
421 return
422 client = self._get_client(client_id)
423 if client is None:
424 return
425 self._watch_client(client)
426 # A reconnect clears the client's start request, so ask again.
427 if self._get_session(client_id) is None:
428 return
429 role = self._get_source_role(client)
430 if role is None or role.stream_active:
431 return
432 self.logger.debug("Re-requesting stream start from %s", client_id)
433 role.request_start()
434
435 def _watch_client(self, client: SendspinClient) -> None:
436 """Subscribe to a source client's events, for signal presence while idle."""
437 # Watch negotiated source roles because pairing can activate the role later
438 # without reconnecting or emitting another event.
439 if "source" not in {role_family(role_id) for role_id in client.negotiated_role_ids}:
440 return
441 state = self._clients.setdefault(client.client_id, _SourceClientState())
442 if state.unwatch is not None:
443 return
444 state.unwatch = client.add_event_listener(self._on_client_event)
445
446 def _on_client_event(self, client: SendspinClient, event: ClientEvent) -> None:
447 client_id = client.client_id
448 if isinstance(event, SourceSignalChangedEvent):
449 self._on_signal_reported(client_id, event.signal)
450 return
451 session = self._get_session(client_id)
452 if session is None:
453 return
454 if isinstance(event, SourceStreamStartedEvent):
455 self.mass.create_task(self._attach_stream(session, event))
456 elif isinstance(event, SourceStreamEndedEvent):
457 self.logger.debug("Sendspin source %s ended its stream", client_id)
458
459 def _on_signal_reported(self, client_id: str, signal: SignalState) -> None:
460 """
461 Drive line-in autostart/autostop from a reported signal presence.
462
463 Only transitions act. The first report for a client is recorded silently so a
464 server restart or a reconnect with the needle already down starts nothing.
465 """
466 state = self._clients.setdefault(client_id, _SourceClientState())
467 previous = state.signal
468 if previous == signal:
469 return
470 state.signal = signal
471 if previous is None:
472 return
473 self.logger.debug("Sendspin source %s signal %s", client_id, signal.value)
474 self._cancel_pending_autostart(client_id)
475 self._cancel_pending_autostop(client_id)
476 if signal == SignalState.PRESENT:
477 self.mass.call_later(
478 AUTOSTART_SIGNAL_DEBOUNCE_S,
479 self._autostart,
480 client_id,
481 task_id=self._autostart_timer_id(client_id),
482 )
483 else:
484 self.mass.call_later(
485 AUTOSTART_SIGNAL_ABSENT_HOLD_S,
486 self._autostop,
487 client_id,
488 task_id=self._autostop_timer_id(client_id),
489 )
490
491 def _cancel_pending_autostart(self, client_id: str, *, cancel_running: bool = True) -> None:
492 self._cancel_scheduled_action(
493 self._autostart_timer_id(client_id), cancel_running=cancel_running
494 )
495
496 def _cancel_pending_autostop(self, client_id: str) -> None:
497 self._cancel_scheduled_action(self._autostop_timer_id(client_id))
498
499 async def _autostart(self, client_id: str) -> None:
500 """Start playing a source whose signal has stayed present."""
501 state = self._clients.get(client_id)
502 if state is None or state.signal != SignalState.PRESENT or state.session is not None:
503 return
504 if (target := await self._resolve_autostart_target(client_id)) is None:
505 return
506 state = self._clients.get(client_id)
507 if state is None or state.signal != SignalState.PRESENT or state.session is not None:
508 return
509 queue_id, uri = target
510 self.logger.info("Line-in signal on %s, starting playback on %s", client_id, queue_id)
511 state.suppressed_autostart_claim = None
512 state.autostart_queue_id = queue_id
513 state.autostart_queue_session_id = None
514 completed = False
515 try:
516 await self.mass.player_queues.play_media(queue_id, uri, option=QueueOption.PLAY)
517 completed = True
518 if self._clients.get(client_id) is state and state.autostart_queue_id == queue_id:
519 state.autostart_queue_session_id = self._player_session_id(queue_id)
520 finally:
521 if (
522 not completed
523 and self._clients.get(client_id) is state
524 and state.autostart_queue_id == queue_id
525 ):
526 state.autostart_queue_id = None
527 state.autostart_queue_session_id = None
528 self._drop_empty_client_state(client_id, state)
529
530 async def _autostop(self, client_id: str) -> None:
531 """Stop a source whose signal has stayed absent, e.g. a record that ended."""
532 state = self._clients.get(client_id)
533 if state is None or state.signal != SignalState.ABSENT or state.session is None:
534 return
535 session = state.session
536 self.logger.info("Line-in signal gone on %s, stopping playback", client_id)
537 try:
538 # Queue stop also cancels pending preload and enqueue timers.
539 await self.mass.players.deselect_source(session.owner_player_id)
540 except (KeyError, PlayerCommandFailed, PlayerUnavailableError) as err:
541 self.logger.debug("Failed to stop player %s: %s", session.owner_player_id, err)
542
543 def _player_session_id(self, player_id: str) -> str | None:
544 """
545 Return the id of whatever playback session the given player is running.
546
547 A live source's session belongs to the player; anything else is the queue's.
548 Used to tell a superseded autostart request apart from the one that replaced it.
549
550 :param player_id: The player to read the session of.
551 """
552 if (session := self.mass.players.get_audio_source_session(player_id)) is not None:
553 return session.playback_session_id
554 return self.mass.player_queues.queue_data(player_id).session_id
555
556 async def _resolve_autostart_target(self, client_id: str) -> tuple[str, str] | None:
557 """Return the queue id and source uri to autostart, if the source is configured."""
558 # Config entry defaults are not persisted until the user saves the page.
559 target_player_id = await self.mass.config.get_player_config_value(
560 client_id, CONF_SOURCE_AUTOSTART_TARGET, default=SOURCE_AUTOSTART_OFF
561 )
562 if not target_player_id or target_player_id == SOURCE_AUTOSTART_OFF:
563 return None
564 player = self.mass.players.get_player(target_player_id)
565 if player is None:
566 self.logger.warning(
567 "Autostart target %s for Sendspin source %s no longer exists",
568 target_player_id,
569 client_id,
570 )
571 return None
572 # Keep a grouped target in its active group. A target already playing a live
573 # source resolves to no queue, which is not a reason to refuse - it is the
574 # player we start on either way.
575 queue = self.mass.players.get_active_queue(player)
576 # mirror the controller's owner resolution: a sync child starts on its leader and
577 # a group member on its group, never on itself - selecting on the child ungroups it
578 target_id = (
579 queue.queue_id
580 if queue
581 else (player.state.synced_to or player.state.active_group or player.player_id)
582 )
583 return target_id, create_uri(MediaType.AUDIO_SOURCE, self.instance_id, client_id)
584
585 async def _attach_stream(
586 self, session: _SourceSession, event: SourceStreamStartedEvent
587 ) -> None:
588 """Route a (re)started source stream into a fresh bridge."""
589 if self._get_session(session.client_id) is not session:
590 return
591 if session.ingest_task is not None:
592 session.ingest_task.cancel()
593 # Provider options are unresolved when the instance first loads.
594 target_latency = (
595 cast("int | None", self.config.get_value(CONF_TARGET_LATENCY))
596 or DEFAULT_TARGET_LATENCY_MS
597 )
598 session.bridge = self._create_bridge(event.audio_format, target_latency)
599 session.last_pcm_monotonic = time.monotonic()
600 session.ingest_task = self.mass.create_task(self._ingest(session, event.handle))
601
602 def _create_bridge(
603 self, input_format: SendspinAudioFormat, target_latency_ms: int
604 ) -> SourceBridge:
605 return AsrcSourceBridge(
606 input_format=input_format,
607 output_format=BRIDGE_OUTPUT_FORMAT,
608 target_latency_ms=target_latency_ms,
609 )
610
611 async def _ingest(self, session: _SourceSession, handle: SourceStream) -> None:
612 """Feed decoded source chunks into the session's bridge until the stream ends."""
613 bridge = session.bridge
614 if bridge is None:
615 return
616 async for pcm, timestamp_us in handle:
617 if self._get_session(session.client_id) is not session or session.bridge is not bridge:
618 break
619 try:
620 bridge.feed(pcm, timestamp_us)
621 except ValueError as err:
622 self.logger.warning("Dropping malformed chunk from %s: %s", session.client_id, err)
623 continue
624 session.last_pcm_monotonic = time.monotonic()
625 session.pcm_received.set()
626 else:
627 if self._get_session(session.client_id) is session and session.bridge is bridge:
628 bridge.flush()
629
630 async def _teardown_session(
631 self, source_id: str, superseded_by_player_id: str | None = None
632 ) -> None:
633 state = self._clients.get(source_id)
634 if state is None or state.session is None:
635 return
636 session, state.session = state.session, None
637 if session.ingest_task is not None:
638 session.ingest_task.cancel()
639 # Stop even when superseding: the replacement session only gets a bridge from a
640 # fresh client_stream/start, which the client sends after a stop/start cycle.
641 if (client := self._get_client(session.client_id)) is not None and (
642 role := self._get_source_role(client)
643 ) is not None:
644 role.request_stop()
645 if superseded_by_player_id is not None and superseded_by_player_id != session.player_id:
646 # Ending the generator leaves the handed-off player draining its buffer over
647 # the new one, so stop it. A same-player re-claim keeps playing.
648 try:
649 await self.mass.players.cmd_stop(session.player_id)
650 except (PlayerCommandFailed, PlayerUnavailableError) as err:
651 self.logger.debug("Failed to stop player %s: %s", session.player_id, err)
652 self._drop_empty_client_state(source_id, state)
653
654 def _get_session(self, client_id: str) -> _SourceSession | None:
655 state = self._clients.get(client_id)
656 return state.session if state is not None else None
657
658 def _drop_empty_client_state(self, client_id: str, state: _SourceClientState) -> None:
659 if (
660 self._clients.get(client_id) is state
661 and state.session is None
662 and state.unwatch is None
663 and state.signal is None
664 and state.autostart_queue_id is None
665 and state.autostart_queue_session_id is None
666 and state.suppressed_autostart_claim is None
667 and not state.selection_lock.locked()
668 ):
669 self._clients.pop(client_id)
670
671 def _cancel_scheduled_action(self, task_id: str, *, cancel_running: bool = True) -> None:
672 self.mass.cancel_timer(task_id)
673 task = self.mass.get_task(task_id)
674 if cancel_running and task is not None and task is not asyncio.current_task():
675 self.mass.cancel_task(task_id)
676
677 def _autostart_timer_id(self, client_id: str) -> str:
678 return f"{self.instance_id}_autostart_{client_id}"
679
680 def _autostop_timer_id(self, client_id: str) -> str:
681 return f"{self.instance_id}_autostop_{client_id}"
682