/
/
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 queue_id: str
103 stream_session_id: str
104 # Retires the prior generator when the same queue 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, queue_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 queue session before waiting so delayed requests keep their identity.
231 queue_session_id = self.mass.player_queues.queue_data(queue_id).session_id
232 try:
233 # Serialize handoffs because stopping the previous player can suspend.
234 async with state.selection_lock:
235 # The client or its role may have changed while waiting for the lock.
236 client = self._get_client(source_id)
237 role = self._get_source_role(client) if client else None
238 if client is None or role is None:
239 raise MediaNotFoundError(f"Sendspin source is not connected: {source_id}")
240 # Reject only the delayed request from an autostart superseded by the user.
241 if state.suppressed_autostart_claim == (queue_id, queue_session_id):
242 raise RuntimeError("Superseded autostart stream request")
243 autostart_queue_id = state.autostart_queue_id
244 if autostart_queue_id == queue_id:
245 state.autostart_queue_id = None
246 state.autostart_queue_session_id = None
247 self._cancel_pending_autostart(source_id, cancel_running=False)
248 else:
249 if autostart_queue_id is not None:
250 autostart_session_id = (
251 state.autostart_queue_session_id
252 or self.mass.player_queues.queue_data(autostart_queue_id).session_id
253 )
254 state.autostart_queue_id = None
255 state.autostart_queue_session_id = None
256 state.suppressed_autostart_claim = (
257 autostart_queue_id,
258 autostart_session_id,
259 )
260 self._cancel_pending_autostart(source_id)
261 if autostart_queue_id is not None:
262 try:
263 await self.mass.player_queues.stop(autostart_queue_id)
264 except (KeyError, PlayerCommandFailed, PlayerUnavailableError) as err:
265 self.logger.debug(
266 "Failed to stop autostart queue %s: %s",
267 autostart_queue_id,
268 err,
269 )
270 # Keep the running bridge when renderers reopen the URL for the same queue.
271 if (live := state.session) is not None and (
272 live.player_id,
273 live.queue_id,
274 ) == (player_id, queue_id):
275 live.stream_session_id = stream_session_id
276 live.generation += 1
277 return
278 await self._teardown_session(source_id, superseded_by_player_id=player_id)
279 if self._unloading or self._clients.get(source_id) is not state:
280 raise PlayerUnavailableError(f"Sendspin source is unloading: {source_id}")
281 # Teardown can suspend long enough for the connection role to be recreated.
282 client = self._get_client(source_id)
283 role = self._get_source_role(client) if client else None
284 if client is None or role is None:
285 raise MediaNotFoundError(f"Sendspin source is not connected: {source_id}")
286 state.session = _SourceSession(
287 client_id=source_id,
288 player_id=player_id,
289 queue_id=queue_id,
290 stream_session_id=stream_session_id,
291 )
292 role.request_start()
293 finally:
294 self._drop_empty_client_state(source_id, state)
295
296 async def on_source_unselected(
297 self, source_id: str, queue_id: str, stream_session_id: str
298 ) -> None:
299 """Release the source when MA tears down its stream."""
300 session = self._get_session(source_id)
301 # Reject stale callbacks from superseded same-queue requests.
302 if session is None or session.stream_session_id != stream_session_id:
303 return
304 await self._teardown_session(source_id)
305 if (state := self._clients.get(source_id)) is not None:
306 state.suppressed_autostart_claim = None
307 self._drop_empty_client_state(source_id, state)
308
309 async def get_audio_stream(
310 self, streamdetails: StreamDetails, seek_position: int = 0
311 ) -> AsyncGenerator[bytes]:
312 """
313 Yield fixed-format PCM pulled from the source's clock bridge.
314
315 The pull cadence of this loop is the master clock: the bridge converts
316 the client's drifting capture stream to it and pads silence on underrun,
317 so the stream keeps playing through gaps (an unplugged line-in is silent,
318 not stopped) until SOURCE_TIMEOUT_S passes without source audio. Pulling at
319 the output rate here makes the controller's own realtime pacer a no-op.
320 """
321 session = self._get_session(streamdetails.item_id)
322 if session is None:
323 raise AudioError(f"Sendspin source is not selected: {streamdetails.item_id}")
324 generation = session.generation
325 await self._await_first_audio(session)
326 if self._get_session(session.client_id) is not session or session.generation != generation:
327 return
328 frames_per_chunk = OUTPUT_SAMPLE_RATE * CHUNK_DURATION_MS // 1000
329 period = CHUNK_DURATION_MS / 1000
330 loop = self.mass.loop
331 next_deadline = loop.time()
332 while True:
333 if (
334 self._get_session(session.client_id) is not session
335 or session.generation != generation
336 ):
337 break
338 if time.monotonic() - session.last_pcm_monotonic > SOURCE_TIMEOUT_S:
339 self.logger.info(
340 "No audio from Sendspin source %s for %.0fs, ending stream",
341 session.client_id,
342 SOURCE_TIMEOUT_S,
343 )
344 break
345 if (bridge := session.bridge) is None:
346 break
347 yield bridge.read(frames_per_chunk)
348 next_deadline += period
349 delay = next_deadline - loop.time()
350 if delay > 0:
351 await asyncio.sleep(delay)
352 elif -delay * 1_000_000 > bridge.occupancy_us:
353 # Consumer stalled beyond the buffered audio. Catch-up reads past this
354 # point would fabricate silence into the timeline, so re-anchor instead.
355 next_deadline = loop.time()
356
357 @property
358 def _sendspin_provider(self) -> SendspinProvider | None:
359 return cast("SendspinProvider | None", self.mass.get_provider("sendspin"))
360
361 def _get_client(self, client_id: str) -> SendspinClient | None:
362 if (sendspin := self._sendspin_provider) is None:
363 return None
364 return sendspin.server_api.get_client(client_id)
365
366 @staticmethod
367 def _get_source_role(client: SendspinClient) -> SourceV1Role | None:
368 roles = client.roles_by_family("source")
369 return cast("SourceV1Role", roles[0]) if roles else None
370
371 async def _await_first_audio(self, session: _SourceSession) -> None:
372 """
373 Block until the source actually streams, so a failed acquisition raises.
374
375 The silence-hold only makes sense once audio has flowed: a client that never
376 answers the start command is a broken source, not a quiet one.
377 """
378 try:
379 async with asyncio.timeout(COLD_START_TIMEOUT_S):
380 await session.pcm_received.wait()
381 except TimeoutError:
382 if self._get_session(session.client_id) is not session:
383 return
384 client = self._get_client(session.client_id)
385 info = client.info_or_none if client else None
386 raise AudioError(
387 f"Sendspin source {session.client_id} did not start streaming",
388 translation_key="no_audio",
389 translation_owner=self.translation_owner,
390 translation_args=[info.name if info else session.client_id],
391 ) from None
392
393 def _on_server_event(self, server: SendspinServer, event: SendspinEvent) -> None:
394 if self._unloading:
395 return
396 match event:
397 case ClientConnectedEvent(client_id):
398 # Roles attach after this event, so defer until the next loop turn.
399 self.mass.create_task(self._on_client_connected(client_id), eager_start=False)
400 case ClientRemovedEvent(client_id) | ClientDisconnectedEvent(client_id):
401 self._cancel_pending_autostart(client_id)
402 self._cancel_pending_autostop(client_id)
403 if (state := self._clients.get(client_id)) is not None:
404 state.signal = None
405 state.autostart_queue_id = None
406 state.autostart_queue_session_id = None
407 if state.session is None:
408 state.suppressed_autostart_claim = None
409 if state.unwatch is not None:
410 state.unwatch()
411 state.unwatch = None
412 self._drop_empty_client_state(client_id, state)
413
414 async def _on_client_connected(self, client_id: str) -> None:
415 """Re-arm watching and streaming for a client that just (re)connected."""
416 if self._unloading:
417 return
418 client = self._get_client(client_id)
419 if client is None:
420 return
421 self._watch_client(client)
422 # A reconnect clears the client's start request, so ask again.
423 if self._get_session(client_id) is None:
424 return
425 role = self._get_source_role(client)
426 if role is None or role.stream_active:
427 return
428 self.logger.debug("Re-requesting stream start from %s", client_id)
429 role.request_start()
430
431 def _watch_client(self, client: SendspinClient) -> None:
432 """Subscribe to a source client's events, for signal presence while idle."""
433 # Watch negotiated source roles because pairing can activate the role later
434 # without reconnecting or emitting another event.
435 if "source" not in {role_family(role_id) for role_id in client.negotiated_role_ids}:
436 return
437 state = self._clients.setdefault(client.client_id, _SourceClientState())
438 if state.unwatch is not None:
439 return
440 state.unwatch = client.add_event_listener(self._on_client_event)
441
442 def _on_client_event(self, client: SendspinClient, event: ClientEvent) -> None:
443 client_id = client.client_id
444 if isinstance(event, SourceSignalChangedEvent):
445 self._on_signal_reported(client_id, event.signal)
446 return
447 session = self._get_session(client_id)
448 if session is None:
449 return
450 if isinstance(event, SourceStreamStartedEvent):
451 self.mass.create_task(self._attach_stream(session, event))
452 elif isinstance(event, SourceStreamEndedEvent):
453 self.logger.debug("Sendspin source %s ended its stream", client_id)
454
455 def _on_signal_reported(self, client_id: str, signal: SignalState) -> None:
456 """
457 Drive line-in autostart/autostop from a reported signal presence.
458
459 Only transitions act. The first report for a client is recorded silently so a
460 server restart or a reconnect with the needle already down starts nothing.
461 """
462 state = self._clients.setdefault(client_id, _SourceClientState())
463 previous = state.signal
464 if previous == signal:
465 return
466 state.signal = signal
467 if previous is None:
468 return
469 self.logger.debug("Sendspin source %s signal %s", client_id, signal.value)
470 self._cancel_pending_autostart(client_id)
471 self._cancel_pending_autostop(client_id)
472 if signal == SignalState.PRESENT:
473 self.mass.call_later(
474 AUTOSTART_SIGNAL_DEBOUNCE_S,
475 self._autostart,
476 client_id,
477 task_id=self._autostart_timer_id(client_id),
478 )
479 else:
480 self.mass.call_later(
481 AUTOSTART_SIGNAL_ABSENT_HOLD_S,
482 self._autostop,
483 client_id,
484 task_id=self._autostop_timer_id(client_id),
485 )
486
487 def _cancel_pending_autostart(self, client_id: str, *, cancel_running: bool = True) -> None:
488 self._cancel_scheduled_action(
489 self._autostart_timer_id(client_id), cancel_running=cancel_running
490 )
491
492 def _cancel_pending_autostop(self, client_id: str) -> None:
493 self._cancel_scheduled_action(self._autostop_timer_id(client_id))
494
495 async def _autostart(self, client_id: str) -> None:
496 """Start playing a source whose signal has stayed present."""
497 state = self._clients.get(client_id)
498 if state is None or state.signal != SignalState.PRESENT or state.session is not None:
499 return
500 if (target := await self._resolve_autostart_target(client_id)) is None:
501 return
502 state = self._clients.get(client_id)
503 if state is None or state.signal != SignalState.PRESENT or state.session is not None:
504 return
505 queue_id, uri = target
506 self.logger.info("Line-in signal on %s, starting playback on %s", client_id, queue_id)
507 state.suppressed_autostart_claim = None
508 state.autostart_queue_id = queue_id
509 state.autostart_queue_session_id = None
510 completed = False
511 try:
512 await self.mass.player_queues.play_media(queue_id, uri, option=QueueOption.PLAY)
513 completed = True
514 if self._clients.get(client_id) is state and state.autostart_queue_id == queue_id:
515 state.autostart_queue_session_id = self.mass.player_queues.queue_data(
516 queue_id
517 ).session_id
518 finally:
519 if (
520 not completed
521 and self._clients.get(client_id) is state
522 and state.autostart_queue_id == queue_id
523 ):
524 state.autostart_queue_id = None
525 state.autostart_queue_session_id = None
526 self._drop_empty_client_state(client_id, state)
527
528 async def _autostop(self, client_id: str) -> None:
529 """Stop a source whose signal has stayed absent, e.g. a record that ended."""
530 state = self._clients.get(client_id)
531 if state is None or state.signal != SignalState.ABSENT or state.session is None:
532 return
533 session = state.session
534 self.logger.info("Line-in signal gone on %s, stopping playback", client_id)
535 try:
536 # Queue stop also cancels pending preload and enqueue timers.
537 await self.mass.player_queues.stop(session.queue_id)
538 except (KeyError, PlayerCommandFailed, PlayerUnavailableError) as err:
539 self.logger.debug("Failed to stop queue %s: %s", session.queue_id, err)
540
541 async def _resolve_autostart_target(self, client_id: str) -> tuple[str, str] | None:
542 """Return the queue id and source uri to autostart, if the source is configured."""
543 # Config entry defaults are not persisted until the user saves the page.
544 target_player_id = await self.mass.config.get_player_config_value(
545 client_id, CONF_SOURCE_AUTOSTART_TARGET, default=SOURCE_AUTOSTART_OFF
546 )
547 if not target_player_id or target_player_id == SOURCE_AUTOSTART_OFF:
548 return None
549 player = self.mass.players.get_player(target_player_id)
550 if player is None:
551 self.logger.warning(
552 "Autostart target %s for Sendspin source %s no longer exists",
553 target_player_id,
554 client_id,
555 )
556 return None
557 # Keep a grouped target in its active group.
558 queue = self.mass.players.get_active_queue(player)
559 if queue is None:
560 self.logger.debug("Autostart target %s has no queue", target_player_id)
561 return None
562 return queue.queue_id, create_uri(MediaType.AUDIO_SOURCE, self.instance_id, client_id)
563
564 async def _attach_stream(
565 self, session: _SourceSession, event: SourceStreamStartedEvent
566 ) -> None:
567 """Route a (re)started source stream into a fresh bridge."""
568 if self._get_session(session.client_id) is not session:
569 return
570 if session.ingest_task is not None:
571 session.ingest_task.cancel()
572 # Provider options are unresolved when the instance first loads.
573 target_latency = (
574 cast("int | None", self.config.get_value(CONF_TARGET_LATENCY))
575 or DEFAULT_TARGET_LATENCY_MS
576 )
577 session.bridge = self._create_bridge(event.audio_format, target_latency)
578 session.last_pcm_monotonic = time.monotonic()
579 session.ingest_task = self.mass.create_task(self._ingest(session, event.handle))
580
581 def _create_bridge(
582 self, input_format: SendspinAudioFormat, target_latency_ms: int
583 ) -> SourceBridge:
584 return AsrcSourceBridge(
585 input_format=input_format,
586 output_format=BRIDGE_OUTPUT_FORMAT,
587 target_latency_ms=target_latency_ms,
588 )
589
590 async def _ingest(self, session: _SourceSession, handle: SourceStream) -> None:
591 """Feed decoded source chunks into the session's bridge until the stream ends."""
592 bridge = session.bridge
593 if bridge is None:
594 return
595 async for pcm, timestamp_us in handle:
596 if self._get_session(session.client_id) is not session or session.bridge is not bridge:
597 break
598 try:
599 bridge.feed(pcm, timestamp_us)
600 except ValueError as err:
601 self.logger.warning("Dropping malformed chunk from %s: %s", session.client_id, err)
602 continue
603 session.last_pcm_monotonic = time.monotonic()
604 session.pcm_received.set()
605 else:
606 if self._get_session(session.client_id) is session and session.bridge is bridge:
607 bridge.flush()
608
609 async def _teardown_session(
610 self, source_id: str, superseded_by_player_id: str | None = None
611 ) -> None:
612 state = self._clients.get(source_id)
613 if state is None or state.session is None:
614 return
615 session, state.session = state.session, None
616 if session.ingest_task is not None:
617 session.ingest_task.cancel()
618 # Stop even when superseding: the replacement session only gets a bridge from a
619 # fresh client_stream/start, which the client sends after a stop/start cycle.
620 if (client := self._get_client(session.client_id)) is not None and (
621 role := self._get_source_role(client)
622 ) is not None:
623 role.request_stop()
624 if superseded_by_player_id is not None and superseded_by_player_id != session.player_id:
625 # Ending the generator leaves the handed-off player draining its buffer over
626 # the new one, so stop it. A same-player re-claim keeps playing.
627 try:
628 await self.mass.players.cmd_stop(session.player_id)
629 except (PlayerCommandFailed, PlayerUnavailableError) as err:
630 self.logger.debug("Failed to stop player %s: %s", session.player_id, err)
631 self._drop_empty_client_state(source_id, state)
632
633 def _get_session(self, client_id: str) -> _SourceSession | None:
634 state = self._clients.get(client_id)
635 return state.session if state is not None else None
636
637 def _drop_empty_client_state(self, client_id: str, state: _SourceClientState) -> None:
638 if (
639 self._clients.get(client_id) is state
640 and state.session is None
641 and state.unwatch is None
642 and state.signal is None
643 and state.autostart_queue_id is None
644 and state.autostart_queue_session_id is None
645 and state.suppressed_autostart_claim is None
646 and not state.selection_lock.locked()
647 ):
648 self._clients.pop(client_id)
649
650 def _cancel_scheduled_action(self, task_id: str, *, cancel_running: bool = True) -> None:
651 self.mass.cancel_timer(task_id)
652 task = self.mass.get_task(task_id)
653 if cancel_running and task is not None and task is not asyncio.current_task():
654 self.mass.cancel_task(task_id)
655
656 def _autostart_timer_id(self, client_id: str) -> str:
657 return f"{self.instance_id}_autostart_{client_id}"
658
659 def _autostop_timer_id(self, client_id: str) -> str:
660 return f"{self.instance_id}_autostop_{client_id}"
661