/
/
1"""Yandex Ynison plugin provider for Music Assistant."""
2
3from __future__ import annotations
4
5import asyncio
6import hashlib
7import random
8import time
9from collections.abc import AsyncGenerator, Callable
10from contextlib import suppress
11from dataclasses import dataclass
12from typing import TYPE_CHECKING, Any, ClassVar, Literal, cast
13
14from music_assistant_models.config_entries import ConfigEntry, ConfigValueOption
15from music_assistant_models.enums import (
16 ConfigEntryType,
17 ContentType,
18 EventType,
19 MediaType,
20 PlaybackState,
21 ProviderFeature,
22 ProviderType,
23 SourceControl,
24 StreamType,
25)
26from music_assistant_models.errors import (
27 LoginFailed,
28 MediaNotFoundError,
29 PlayerCommandFailed,
30 UnsupportedFeaturedException,
31)
32from music_assistant_models.media_items import AudioSource, ProviderMapping
33from music_assistant_models.streamdetails import StreamDetails, StreamMetadata
34from ya_passport_auth import SecretStr
35from ya_passport_auth.ma import BorrowedCredentialSource
36
37from music_assistant.helpers.ffmpeg import get_ffmpeg_stream
38from music_assistant.helpers.throttle_retry import BYPASS_THROTTLER, ThrottlerManager
39from music_assistant.models.plugin import PluginProvider
40
41from .auth import refresh_music_token
42from .constants import (
43 CONF_ALLOW_PLAYER_SWITCH,
44 CONF_DEVICE_ID,
45 CONF_MASS_PLAYER_ID,
46 CONF_OUTPUT_BIT_DEPTH,
47 CONF_OUTPUT_SAMPLE_RATE,
48 CONF_PUBLISH_NAME,
49 CONF_TOKEN,
50 CONF_X_TOKEN,
51 CONF_YM_INSTANCE,
52 DEFAULT_DISPLAY_NAME,
53 OUTPUT_AUTO,
54 PLAYER_ID_AUTO,
55 YANDEX_MUSIC_CONF_QUALITY,
56 YANDEX_MUSIC_LOSSLESS_QUALITIES,
57 YM_INSTANCE_OWN,
58)
59from .protocols import YandexMusicProviderLike
60from .streaming import (
61 PCM_LOSSLESS_PARAMS,
62 PCM_LOSSY_PARAMS,
63 PROBE_ARGS,
64 make_pcm_format,
65)
66from .ynison_client import (
67 YnisonClient,
68 YnisonDeviceInfo,
69 YnisonSendError,
70 YnisonState,
71 generate_device_id,
72 make_version_block,
73)
74
75if TYPE_CHECKING:
76 from music_assistant_models.config_entries import ProviderConfig
77 from music_assistant_models.event import MassEvent
78 from music_assistant_models.media_items import AudioFormat
79 from music_assistant_models.provider import ProviderManifest
80
81 from music_assistant.mass import MusicAssistant
82
83# How often (seconds) to sync progress to MA UI and Ynison.
84_PROGRESS_SYNC_INTERVAL = 5.0
85
86# Grace window after our own REPLACE/seek during which incoming Ynison
87# progress updates are treated as our own echo (not a user seek).
88_ECHO_GRACE_PERIOD = 3.0
89
90# Bound on the synchronous pre-fetch in _prefetch_format_for_track. A slow
91# pre-fetch is treated like a failed one â fall back to the current format
92# and let the in-stream `_get_stream_details_with_retry` handle retries.
93_PREFETCH_FORMAT_TIMEOUT = 2.5
94
95# Idempotency cache TTL for outbound peer-commands.
96_COMMAND_IDEMPOTENCY_TTL = 1.0
97
98# stable id for the single AudioSource this provider exposes;
99# combined with the provider instance_id this forms the persistent uri
100AUDIO_SOURCE_ID = "main"
101
102# Retry settings for transient Yandex API failures
103_API_MAX_RETRIES = 3
104_API_INITIAL_BACKOFF = 2.0
105_API_MAX_BACKOFF = 30.0
106
107# Cache TTL for stream details (seconds)
108_STREAM_DETAILS_CACHE_TTL = 300 # 5 minutes
109
110# In-memory music-token cache TTL (seconds). Yandex music tokens live ~60 min;
111# 50 min leaves 10 min headroom before the server would reject them. Tied to
112# the borrow-mode-with-only-x_token + 401-storm path described in spec 0004.
113_MUSIC_TOKEN_TTL_S = 50 * 60
114
115# Maximum number of distinct x_token entries kept in the own-mode music-token
116# cache (borrow mode caches inside BorrowedCredentialSource). 4 keeps headroom
117# for an x_token rotation with one refresh in flight.
118_MUSIC_TOKEN_CACHE_MAX = 4
119
120# Accepted non-auto values for output format overrides; mirrors the options
121# offered in CONF_OUTPUT_SAMPLE_RATE / CONF_OUTPUT_BIT_DEPTH config entries.
122# Used defensively to reject stale/tampered values without raising.
123_VALID_SAMPLE_RATES: frozenset[str] = frozenset({"44100", "48000", "96000"})
124_VALID_BIT_DEPTHS: frozenset[str] = frozenset({"16", "24"})
125
126
127@dataclass(frozen=True)
128class _CachedToken:
129 """
130 Music token entry in the in-memory cache.
131
132 `expires_monotonic` is compared against the provider's `_now()` seam.
133 """
134
135 token: SecretStr
136 expires_monotonic: float
137
138
139def _hash_x_token(x_token: str) -> str:
140 """
141 Return the SHA-256 hex digest of an x_token, used as cache key.
142
143 The raw x_token is never stored in dict keys (defence-in-depth against
144 accidental log / dump leakage of the cache structure).
145 """
146 return hashlib.sha256(x_token.encode("utf-8")).hexdigest()
147
148
149class YandexYnisonProvider(PluginProvider):
150 """Implementation of the Yandex Music Connect (Ynison) Plugin."""
151
152 # PluginProvider base does not declare `is_streaming_provider`; MA's
153 # audio-analysis path raises AttributeError for live sources without
154 # an explicit opt-out. Analysing transient external-source tracks
155 # buys nothing.
156 is_streaming_provider: bool = False
157
158 @property
159 def instance_name_postfix(self) -> str | None:
160 """Return display name as instance postfix for multi-instance setups."""
161 name = self._display_name
162 return name if name != DEFAULT_DISPLAY_NAME else None
163
164 def __init__(
165 self,
166 mass: MusicAssistant,
167 manifest: ProviderManifest,
168 config: ProviderConfig,
169 supported_features: set[ProviderFeature],
170 ) -> None:
171 """Initialize the Ynison plugin provider."""
172 super().__init__(mass, manifest, config, supported_features)
173
174 # Setup identity and playback options
175 self._default_player_id: str = (
176 cast("str", self.get_setup_value(CONF_MASS_PLAYER_ID)) or PLAYER_ID_AUTO
177 )
178 allow_switch_value = self.config.get_value(CONF_ALLOW_PLAYER_SWITCH)
179 self._allow_player_switch: bool = (
180 cast("bool", allow_switch_value) if allow_switch_value is not None else True
181 )
182 self._cfg_sample_rate: str = (
183 cast("str", self.config.get_value(CONF_OUTPUT_SAMPLE_RATE)) or OUTPUT_AUTO
184 )
185 self._cfg_bit_depth: str = (
186 cast("str", self.config.get_value(CONF_OUTPUT_BIT_DEPTH)) or OUTPUT_AUTO
187 )
188 self._display_name: str = (
189 cast("str", self.get_setup_value(CONF_PUBLISH_NAME)) or DEFAULT_DISPLAY_NAME
190 )
191
192 # Token source â None = own (manually entered CONF_TOKEN);
193 # otherwise the instance_id of a linked yandex_music provider to borrow from.
194 ym_instance_value = cast("str | None", self.get_setup_value(CONF_YM_INSTANCE))
195 self._ym_instance_id: str | None = (
196 ym_instance_value
197 if ym_instance_value and ym_instance_value != YM_INSTANCE_OWN
198 else None
199 )
200 # Borrow mode: read-only credential source over the linked
201 # yandex_music instance (shared auth layer). The owner stays the
202 # single writer/rotator of persisted credentials; minted music
203 # tokens are cached in-memory inside the source (TTL + LRU +
204 # coalesced refreshes per its spec).
205 self._borrow_source: BorrowedCredentialSource | None = (
206 BorrowedCredentialSource(self.mass, self._ym_instance_id)
207 if self._ym_instance_id is not None
208 else None
209 )
210
211 # Device ID â persist in config so re-registration uses the same ID
212 device_id = cast("str | None", self.config.get_value(CONF_DEVICE_ID))
213 if not device_id:
214 device_id = generate_device_id()
215 self._update_config_value(CONF_DEVICE_ID, device_id)
216 self._device_id: str = device_id
217
218 # Runtime state
219 self._active_player_id: str | None = None
220 self._ynison: YnisonClient | None = None
221 self._runner_task: asyncio.Task[None] | None = None
222 self._on_unload_callbacks: list[Callable[..., None]] = []
223 self._yandex_provider: YandexMusicProviderLike | None = None
224 self._current_streaming_track_id: str | None = None
225 self._track_changed_event = asyncio.Event()
226 self._stream_stop_event = asyncio.Event()
227 self._seek_position_ms: int = 0
228 self._seek_grace_until: float = 0.0
229 self._last_player_update_time: float = 0.0
230 self._actual_duration_ms: int = 0
231 self._prefetched_list: list[dict[str, Any]] | None = None
232 self._prefetch_task: asyncio.Task[Any] | None = None
233 self._normalized_params: dict[str, Any] = PCM_LOSSY_PARAMS
234 self._normalized_format: AudioFormat = make_pcm_format(PCM_LOSSY_PARAMS)
235
236 # Rate limiter for Yandex API calls (max 2 req/s)
237 self._api_throttler = ThrottlerManager(rate_limit=2, period=1.0)
238
239 # Progress tracking â byte counter is the single source of truth
240 # during active streaming; Ynison echoes are detected via
241 # YnisonState.last_update_is_echo and ignored.
242 self._streaming_progress_ms: int = 0
243
244 # AudioSource MediaItem + per-stream state
245 self._stream_metadata = StreamMetadata(
246 title=f"Yandex Music Connect | {self._display_name}",
247 )
248 self._audio_source = AudioSource(
249 item_id=AUDIO_SOURCE_ID,
250 provider=self.instance_id,
251 name=self.name,
252 provider_mappings={
253 ProviderMapping(
254 item_id=AUDIO_SOURCE_ID,
255 provider_domain=self.domain,
256 provider_instance=self.instance_id,
257 # Fresh AudioFormat copy: AudioFormat is mutable and MA's
258 # FFMpeg._log_reader_task sets `input_format.codec_type`
259 # in-place. Sharing `self._normalized_format` here would
260 # let that mutation leak into the ProviderMapping and into
261 # later StreamDetails snapshots.
262 audio_format=make_pcm_format(self._normalized_params),
263 )
264 },
265 can_play_pause=False,
266 can_seek=False,
267 can_next_previous=False,
268 exclusive=True,
269 allow_external_trigger=True,
270 )
271 # _in_use_by_queue tracks the queue currently consuming our stream
272 self._in_use_by_queue: str | None = None
273 # _active_session_id is the controller-provided token for the current
274 # stream request â used to reject stale on_source_unselected callbacks
275 # after a same-queue reconnect supersedes the previous request.
276 self._active_session_id: str | None = None
277
278 # Idempotency cache for outbound peer-commands. Suppresses duplicate
279 # (action, key) pairs inside `_COMMAND_IDEMPOTENCY_TTL` â protects
280 # against echo-storms where the same Ynison broadcast lands on our
281 # state-handler twice in quick succession.
282 self._command_idempotency: dict[tuple[str, str | None], float] = {}
283
284 # "Ynison paused us externally â expect a resume that needs
285 # `play_media` re-issuance." Set in `_pause_playback`, read in
286 # `_activate_playback`. Survives a stray `_stream_stop_event` clear
287 # independent of the stop signal (which covers non-pause stop reasons).
288 self._externally_paused: bool = False
289
290 # In-memory music-token cache keyed by SHA-256(x_token). 50-min TTL,
291 # 4-entry LRU. Coalesces concurrent refresh attempts via a single
292 # asyncio.Lock so a reconnect storm makes at most one Passport call.
293 # `_now` is a seam for tests to advance the clock.
294 self._token_cache: dict[str, _CachedToken] = {}
295 self._token_refresh_lock = asyncio.Lock()
296 self._now: Callable[[], float] = time.monotonic
297
298 async def get_config_entries(self) -> tuple[ConfigEntry, ...]:
299 """
300 Return Config entries to configure this provider.
301
302 Account, player and device identity are collected by the interactive setup flow;
303 only runtime playback options live here.
304 """
305 return (
306 ConfigEntry(
307 key=CONF_ALLOW_PLAYER_SWITCH,
308 type=ConfigEntryType.BOOLEAN,
309 default_value=True,
310 ),
311 ConfigEntry(
312 key=CONF_OUTPUT_SAMPLE_RATE,
313 type=ConfigEntryType.STRING,
314 default_value=OUTPUT_AUTO,
315 options=[
316 ConfigValueOption(OUTPUT_AUTO),
317 ConfigValueOption("44100"),
318 ConfigValueOption("48000"),
319 ConfigValueOption("96000"),
320 ],
321 advanced=True,
322 ),
323 ConfigEntry(
324 key=CONF_OUTPUT_BIT_DEPTH,
325 type=ConfigEntryType.STRING,
326 default_value=OUTPUT_AUTO,
327 options=[
328 ConfigValueOption(OUTPUT_AUTO),
329 ConfigValueOption("16"),
330 ConfigValueOption("24"),
331 ],
332 advanced=True,
333 ),
334 ConfigEntry(
335 key=CONF_DEVICE_ID,
336 type=ConfigEntryType.STRING,
337 hidden=True,
338 required=False,
339 ),
340 )
341
342 # ------------------------------------------------------------------
343 # Provider lifecycle
344 # ------------------------------------------------------------------
345
346 async def handle_async_init(self) -> None:
347 """Handle async initialization of the provider."""
348 if self._ym_instance_id is not None:
349 self.logger.info(
350 "Borrowing credentials from yandex_music instance '%s'",
351 self._ym_instance_id,
352 )
353 else:
354 self.logger.info("Using manually configured Yandex Music token (no auto-refresh)")
355 token = await self._resolve_token()
356
357 device_info = YnisonDeviceInfo(
358 device_id=self._device_id,
359 title=self._display_name,
360 )
361
362 self._ynison = YnisonClient(
363 token=token,
364 device_info=device_info,
365 on_state_update=self._handle_ynison_state,
366 logger=self.logger,
367 on_auth_failure=self._refresh_ynison_token,
368 )
369
370 self._runner_task = self.mass.create_task(self._ynison.connect())
371
372 # Subscribe to provider events to detect linked yandex_music provider
373 self._on_unload_callbacks.append(
374 self.mass.subscribe(
375 self._on_provider_event,
376 EventType.PROVIDERS_UPDATED,
377 )
378 )
379 # Initial check for matching provider
380 self.mass.create_task(self._check_yandex_provider_match())
381
382 async def unload(self, is_removed: bool = False) -> None:
383 """Handle close/cleanup of the provider."""
384 if self._prefetch_task and not self._prefetch_task.done():
385 self._prefetch_task.cancel()
386 with suppress(asyncio.CancelledError):
387 await self._prefetch_task
388
389 if self._ynison:
390 await self._ynison.disconnect()
391
392 if self._runner_task and not self._runner_task.done():
393 self._runner_task.cancel()
394 with suppress(asyncio.CancelledError):
395 await self._runner_task
396
397 for callback in self._on_unload_callbacks:
398 with suppress(KeyError):
399 callback()
400
401 async def get_audio_sources(self) -> list[AudioSource]:
402 """Return the AudioSources this plugin currently exposes."""
403 return [self._audio_source]
404
405 async def get_stream_details(self, item_id: str, media_type: MediaType) -> StreamDetails:
406 """
407 Return StreamDetails for streaming the Yandex Music Connect audio.
408
409 Side-effect-free: ownership is claimed in on_source_selected (which the
410 streams controller fires before this method on the actual stream
411 request). Keeping this idempotent means preload paths like
412 player_queues._load_item can fetch streamdetails without claiming the
413 source and blocking a subsequent cross-queue handoff.
414 """
415 if item_id != AUDIO_SOURCE_ID:
416 raise MediaNotFoundError(f"Unknown AudioSource: {item_id}")
417 return StreamDetails(
418 provider=self.instance_id,
419 item_id=item_id,
420 # Fresh AudioFormat copy per call: MA's ffmpeg mutates
421 # input_format.codec_type in place, so a shared instance would
422 # propagate that mutation into future stream-details snapshots.
423 audio_format=make_pcm_format(self._normalized_params),
424 media_type=MediaType.AUDIO_SOURCE,
425 stream_type=StreamType.CUSTOM,
426 stream_metadata=self._stream_metadata,
427 )
428
429 async def on_source_control(
430 self,
431 source_id: str,
432 action: SourceControl,
433 value: int | None = None,
434 ) -> None:
435 """Proxy playback control commands to Yandex via the linked Yandex Music provider."""
436 if source_id != AUDIO_SOURCE_ID:
437 return
438 if action == SourceControl.PLAY:
439 await self._on_play()
440 elif action == SourceControl.PAUSE:
441 await self._on_pause()
442 elif action == SourceControl.NEXT:
443 await self._on_next()
444 elif action == SourceControl.PREVIOUS:
445 await self._on_previous()
446 elif action == SourceControl.SEEK and value is not None:
447 await self._on_seek(value)
448
449 async def get_audio_stream( # noqa: PLR0915
450 self, streamdetails: StreamDetails, seek_position: int = 0
451 ) -> AsyncGenerator[bytes]:
452 """
453 Return continuous audio stream following Ynison track changes.
454
455 Streams the current track, then waits for track changes and streams
456 the next track automatically. Runs until the source is deselected.
457
458 The PCM format is frozen at session start to match what the outer
459 ffmpeg captured from ``self._normalized_format``. If
460 ``_update_normalized_format()`` fires mid-session (e.g. a provider
461 reload), the new format takes effect only on the *next* session â
462 preventing bit-depth/sample-rate mismatches that cause noise.
463 """
464 self._stream_stop_event.clear()
465 # snapshot the consumer at session start; the rest of this generator
466 # treats the queue_id as the player_id (they are the same by convention).
467 # The lock may legitimately be empty here â MA's `_load_item` preload
468 # path drives the generator to fill an initial audio buffer BEFORE
469 # `on_source_selected` has been dispatched, so `_in_use_by_queue` is
470 # still None on that call. `had_claim` records whether a lock was
471 # already in force at entry; only in that case do we enforce
472 # cross-session invariants on the loop and the `finally` cleanup.
473 player_id = self._in_use_by_queue or ""
474 had_claim = self._in_use_by_queue is not None
475 # Snapshot the active session id too so a same-queue reconnect (which
476 # updates _active_session_id but not _in_use_by_queue) is treated as a
477 # superseding session: the loop exits early, and the finally clear
478 # below skips the release so it doesn't clobber the new claim.
479 captured_session_id = self._active_session_id
480
481 # MA's streams controller may pass a non-zero seek_position (e.g. a
482 # resume initiated through a path that does NOT go through Ynison and
483 # therefore did not set `_seek_position_ms`). Honor it as the seed for
484 # the upcoming track. The Ynison-driven seek path (`_activate_playback`
485 # / `_on_seek`) keeps writing `_seek_position_ms` directly, which
486 # subsequent track iterations consume â only the seed differs.
487 if seek_position > 0 and self._seek_position_ms == 0:
488 self._seek_position_ms = seek_position * 1000
489
490 # Freeze format for this streaming session so every inner ffmpeg
491 # produces data matching the outer ffmpeg's captured input_format.
492 session_params: dict[str, Any] = dict(self._normalized_params)
493 session_fmt: AudioFormat = make_pcm_format(session_params)
494
495 try:
496 while not self._stream_stop_event.is_set() and (
497 # Preload path: no claim was active at entry â drive the
498 # loop purely off Ynison state and the stop event.
499 not had_claim or not self._session_lost(player_id, captured_session_id)
500 ):
501 if not self._ynison or not self._ynison.state.current_track_id:
502 # Wait for a track to appear
503 self._track_changed_event.clear()
504 try:
505 await asyncio.wait_for(self._track_changed_event.wait(), timeout=30.0)
506 except TimeoutError:
507 continue
508 continue
509
510 # Clear event before reading state so any subsequent update
511 # re-sets the event instead of being silently cleared.
512 self._track_changed_event.clear()
513 track_id = self._ynison.state.current_track_id
514 self._current_streaming_track_id = track_id
515
516 # `_pause_playback` set the stop event; finalize.
517 if self._ynison.state.is_paused:
518 return
519
520 if not self._yandex_provider:
521 self.logger.warning(
522 "No linked Yandex Music provider â cannot stream track %s", track_id
523 )
524 self._stream_stop_event.set()
525 if self._in_use_by_queue == player_id:
526 await self.mass.players.cmd_stop(player_id)
527 return
528
529 # Stream the current track
530 seek_ms = self._seek_position_ms
531 self._seek_position_ms = 0
532 bytes_yielded = 0
533 self._streaming_progress_ms = seek_ms
534 last_progress_sync = time.monotonic()
535
536 track_fmt = make_pcm_format(session_params)
537 async for chunk in self._stream_track(
538 track_id, seek_ms=seek_ms, session_params=session_params
539 ):
540 yield chunk
541 bytes_yielded += len(chunk)
542 now_mono = time.monotonic()
543 if now_mono - last_progress_sync >= _PROGRESS_SYNC_INTERVAL:
544 last_progress_sync = now_mono
545 await self._sync_progress(seek_ms, bytes_yielded, player_id, session_fmt)
546 if (
547 self._track_changed_event.is_set()
548 or self._stream_stop_event.is_set()
549 or (had_claim and self._session_lost(player_id, captured_session_id))
550 ):
551 break
552
553 # Align to PCM frame boundary â prevents misalignment in MA's
554 # downstream ffmpeg when a track stream is interrupted mid-chunk.
555 # We pad with zeros (can't un-yield bytes already sent downstream).
556 frame_size = (track_fmt.bit_depth // 8) * track_fmt.channels
557 if frame_size > 0:
558 excess = bytes_yielded % frame_size
559 if excess:
560 yield b"\x00" * (frame_size - excess)
561
562 # Don't clear _current_streaming_track_id yet â keep it set
563 # during advance/wait so Ynison echo of the same track doesn't
564 # trigger a false track-change detection in _activate_playback.
565
566 if self._stream_stop_event.is_set():
567 break
568
569 # Differentiate "track finished naturally" from "inner loop
570 # broke out early". Signalling completion on an
571 # interrupted track makes Yandex auto-advance the queue â
572 # surfaces as an unwanted skip on pause / handoff.
573 broke_for_pause = self._ynison is not None and self._ynison.state.is_paused
574 broke_for_session_change = had_claim and self._session_lost(
575 player_id, captured_session_id
576 )
577 natural_end = (
578 not self._track_changed_event.is_set()
579 and not broke_for_pause
580 and not broke_for_session_change
581 and self._ynison is not None
582 )
583 if natural_end:
584 self.logger.info("Track %s finished, advancing to next", track_id)
585 await self._signal_track_completion()
586 if not await self._wait_for_track_change(track_id):
587 self._stream_stop_event.set()
588 break
589
590 # Clear before next iteration â the new track ID will be set at
591 # the top of the loop from the latest Ynison state.
592 self._current_streaming_track_id = None
593 finally:
594 # Release ownership only if THIS generator owned the claim at
595 # entry AND no one else has superseded it since. The double-guard
596 # protects against a same-queue reconnect refreshing the session
597 # id without changing the queue id; clearing the lock on the old
598 # generator's teardown would otherwise clobber the new session's
599 # claim. `had_claim` keeps the preload path from touching the lock
600 # at all (no claim ever existed to release).
601 if had_claim and not self._session_lost(player_id, captured_session_id):
602 self._in_use_by_queue = None
603 self._current_streaming_track_id = None
604
605 async def on_source_selected(
606 self,
607 source_id: str,
608 player_id: str,
609 queue_id: str,
610 stream_session_id: str,
611 ) -> None:
612 """Handle callback when this AudioSource has been selected/started on a player."""
613 if source_id != AUDIO_SOURCE_ID or not player_id:
614 return
615
616 # Check if manual player switching is allowed
617 if not self._allow_player_switch:
618 current_target = self._get_target_player_id()
619 if player_id != current_target and current_target:
620 # Redirect to the configured target, but only once per
621 # idempotency window. The target may be a sendspin bridge /
622 # sync-group whose stream is consumed under a player id that
623 # never equals `current_target`, so each redirect re-triggers
624 # selection here. Re-issuing `play_media` on every rejection
625 # turns that into an unbounded AudioError storm; the raise
626 # below still aborts every wrong-player stream regardless.
627 if self._idempotent("source_redirect", current_target):
628 self.logger.debug(
629 "Player switching disabled, redirecting selection from %s to %s",
630 player_id,
631 current_target,
632 )
633 await self.mass.player_queues.play_media(
634 current_target, str(self._audio_source.uri)
635 )
636 msg = f"Player switching is disabled; source must remain on {current_target}"
637 raise RuntimeError(msg)
638
639 # Stop previous player if switching. The lock claim a few lines below
640 # replaces the previous queue's claim; the previous stream loop notices
641 # the queue change and exits cleanly.
642 if self._active_player_id and self._active_player_id != player_id:
643 prev_player_id = self._active_player_id
644 self.logger.info(
645 "Source selected on %s, stopping %s",
646 player_id,
647 prev_player_id,
648 )
649 try:
650 await self.mass.players.cmd_stop(prev_player_id)
651 except Exception as err:
652 self.logger.debug(
653 "Failed to stop previous player %s: %s",
654 prev_player_id,
655 err,
656 )
657
658 # Claim ownership for this queue. The lock lives here (not in
659 # get_stream_details) so preload paths can fetch streamdetails without
660 # accidentally blocking a subsequent cross-queue handoff at the actual
661 # stream request.
662 self._in_use_by_queue = queue_id
663 # Record this request's session id so a later on_source_unselected can
664 # tell whether it is the live teardown or a stale callback from a
665 # superseded same-queue request.
666 self._active_session_id = stream_session_id
667 self._active_player_id = player_id
668 self.logger.debug("Active player set to: %s", player_id)
669
670 async def on_source_unselected(
671 self, source_id: str, queue_id: str, stream_session_id: str
672 ) -> None:
673 """Release the queue-scoped exclusive claim when MA tears down the stream."""
674 if source_id != AUDIO_SOURCE_ID:
675 return
676 # Reject stale callbacks: only release if this is still the active
677 # session. A queue_id check alone is not sufficient â same-queue
678 # reconnects (player drops + reopens the same stream URL before the
679 # original request's finally fires) would otherwise let the old
680 # request's late callback clear the live claim of the new stream.
681 if self._active_session_id != stream_session_id:
682 return
683 self._active_session_id = None
684 if self._in_use_by_queue == queue_id:
685 self._in_use_by_queue = None
686
687 async def _wait_for_track_change(self, old_track_id: str, timeout: float = 30.0) -> bool:
688 """
689 Wait for Ynison to report a different track, ignoring echoes.
690
691 After _signal_track_completion sends update_playing_status, Ynison
692 echoes back the same track with updated progress. Only return True
693 once current_track_id actually differs from old_track_id.
694 """
695 deadline = time.monotonic() + timeout
696 while not self._stream_stop_event.is_set():
697 # Check state BEFORE clearing the event. Ynison may have already
698 # advanced between _signal_track_completion() returning and this
699 # method running; clearing first would drop the set() that went
700 # with the state update, leaving us to wait until timeout.
701 # Check is race-free: no await between the read and clear() below.
702 # None means empty/unreadable queue â treat as "not advanced."
703 if self._ynison:
704 current = self._ynison.state.current_track_id
705 if current is not None and current != old_track_id:
706 return True
707 self._track_changed_event.clear()
708 remaining = deadline - time.monotonic()
709 if remaining <= 0:
710 break
711 try:
712 await asyncio.wait_for(self._track_changed_event.wait(), timeout=remaining)
713 except TimeoutError:
714 break
715 self.logger.info("No new track from Ynison after completion, stopping stream")
716 return False
717
718 async def _stream_track(
719 self,
720 track_id: str,
721 seek_ms: int = 0,
722 session_params: dict[str, Any] | None = None,
723 ) -> AsyncGenerator[bytes]:
724 """
725 Stream a single track, normalizing to fixed PCM via per-track ffmpeg.
726
727 Every track is decoded through its own ffmpeg process to produce a
728 fixed PCM output (s16le or s24le based on YM quality setting). This
729 ensures MA's single ffmpeg process never encounters mid-stream format
730 changes (codec, bit depth, sample rate).
731
732 *session_params* â frozen format dict from the enclosing
733 ``get_audio_stream()`` session. Falls back to the current
734 ``_normalized_params`` when called outside a session.
735 """
736 # In-flight stream fetch outranks unrelated 429 cooldowns:
737 # dropping a stream the user is actively trying to play is
738 # worse than risking another captcha. Prefetch deliberately
739 # stays throttled (see `_prefetch_format_for_track`).
740 bypass_token = BYPASS_THROTTLER.set(True)
741 try:
742 stream_details = await self._get_stream_details_with_retry(track_id)
743 except Exception:
744 self.logger.exception("Failed to get stream details for track %s", track_id)
745 self._stream_stop_event.set()
746 return
747 finally:
748 BYPASS_THROTTLER.reset(bypass_token)
749
750 # Re-capture the provider after the above await: _yandex_provider may
751 # have flipped to None while we were fetching stream details. Using
752 # the attribute directly below would race with
753 # _check_yandex_provider_match.
754 provider = self._yandex_provider
755 if provider is None:
756 self.logger.warning(
757 "Linked Yandex Music provider unloaded mid-stream â stopping track %s",
758 track_id,
759 )
760 self._stream_stop_event.set()
761 return
762
763 await self._update_metadata_from_stream(stream_details, seek_ms)
764
765 # No -re here: MA's realtime pacer is the single pacing authority for
766 # AudioSources. Pacing the decode a second time would pin it to realtime
767 # and forbid the small read-ahead that absorbs CDN jitter; back-pressure
768 # through the generator chain still bounds memory.
769 extra_input_args = list(PROBE_ARGS)
770 if seek_ms > 0:
771 extra_input_args += ["-ss", f"{seek_ms / 1000.0:.3f}"]
772
773 # Use session format when available, otherwise current normalized params
774 params = session_params if session_params is not None else self._normalized_params
775 out_fmt = make_pcm_format(params)
776 # Log the output rate + bit depth alongside the source format: with the
777 # passthrough fast path this PCM IS the delivered audio, so the line must
778 # let an operator read rate passthrough vs a resample, not just codec.
779 self.logger.info(
780 "Streaming track %s â %s/%dHz/%dbit: input=%s seek=%dms",
781 track_id,
782 out_fmt.content_type.value,
783 out_fmt.sample_rate,
784 out_fmt.bit_depth,
785 stream_details.audio_format,
786 seek_ms,
787 )
788 async for chunk in get_ffmpeg_stream(
789 audio_input=provider.get_audio_stream(stream_details),
790 input_format=stream_details.audio_format,
791 output_format=out_fmt,
792 extra_input_args=extra_input_args,
793 ):
794 yield chunk
795
796 async def _get_stream_details_with_retry(
797 self,
798 track_id: str,
799 media_type: MediaType = MediaType.TRACK,
800 ) -> StreamDetails:
801 """Fetch stream details with caching, throttling, and retry."""
802 # Capture the linked yandex_music provider into a local ref at entry.
803 # self._yandex_provider can flip to None mid-await when the linked
804 # MusicProvider is unloaded (see _check_yandex_provider_match, which
805 # runs as a background task on provider-loaded/unloaded events).
806 # Dereferencing the attribute after an await would raise
807 # AttributeError and hard-stop the audio generator.
808 provider = self._yandex_provider
809 if provider is None:
810 raise LoginFailed(
811 "Linked Yandex Music provider is not loaded â cannot fetch stream details"
812 )
813
814 cache_key = f"ynison_sd_{track_id}"
815 cached = await self.mass.cache.get(
816 cache_key,
817 provider=self.instance_id,
818 base_class=StreamDetails,
819 )
820 if cached is not None:
821 self.logger.debug("Stream details cache hit for %s", track_id)
822 return cast("StreamDetails", cached)
823
824 backoff = _API_INITIAL_BACKOFF
825 last_err: Exception | None = None
826 for attempt in range(_API_MAX_RETRIES):
827 async with self._api_throttler.acquire() as delay:
828 if delay > 0:
829 self.logger.debug("get_stream_details throttled %.1fs", delay)
830 try:
831 sd = await provider.get_stream_details(track_id, media_type)
832 # StreamDetails.data has serialize="omit", so to_dict()
833 # strips it. Manually include it so cached entries keep
834 # the URL / decryption key needed by get_audio_stream().
835 cache_value = sd.to_dict()
836 cache_value["data"] = sd.data
837 # Respect the provider's expiration (e.g. yandex_music sets
838 # 50 s because CDN URLs expire after ~60 s). Fall back to
839 # our default TTL when the provider does not override.
840 cache_ttl = min(_STREAM_DETAILS_CACHE_TTL, sd.expiration)
841 if cache_ttl > 0:
842 await self.mass.cache.set(
843 cache_key,
844 cache_value,
845 expiration=cache_ttl,
846 provider=self.instance_id,
847 )
848 return sd
849 except asyncio.CancelledError:
850 raise
851 except Exception as err:
852 last_err = err
853 if attempt < _API_MAX_RETRIES - 1:
854 jitter = backoff * random.uniform(0.75, 1.25)
855 self.logger.warning(
856 "get_stream_details attempt %d/%d failed: %s, retrying in %.1fs",
857 attempt + 1,
858 _API_MAX_RETRIES,
859 err,
860 jitter,
861 )
862 await asyncio.sleep(jitter)
863 backoff = min(backoff * 2, _API_MAX_BACKOFF)
864 msg = f"get_stream_details failed after {_API_MAX_RETRIES} attempts for {track_id}"
865 raise RuntimeError(msg) from last_err
866
867 async def _invalidate_stream_cache(self, track_id: str) -> None:
868 """Evict cached stream details for a track so the next fetch is fresh."""
869 cache_key = f"ynison_sd_{track_id}"
870 await self.mass.cache.delete(cache_key, provider=self.instance_id)
871 self.logger.debug("Invalidated stream cache for %s", track_id)
872
873 # ------------------------------------------------------------------
874 # Token handling
875 # ------------------------------------------------------------------
876
877 async def _refresh_via_x_token(self, x_token: str) -> SecretStr:
878 """
879 Refresh the music token from an x_token, caching the result.
880
881 Within :data:`_MUSIC_TOKEN_TTL_S` of a successful refresh, subsequent
882 calls for the same x_token return the cached :class:`SecretStr`
883 without hitting Yandex Passport. Concurrent callers coalesce via
884 :attr:`_token_refresh_lock`.
885
886 :param x_token: Long-lived session token to exchange for a music
887 token. Hashed before use as a cache key; the raw value is
888 never stored in dict keys or logs.
889 :returns: Fresh or cached music-scoped :class:`SecretStr`.
890 :raises LoginFailed: When Yandex explicitly rejects the x_token
891 (propagated from :func:`provider.auth.refresh_music_token`).
892 :raises ResourceTemporarilyUnavailable: On transient Passport
893 failures (network, rate limit) â retry later, credentials
894 are still good.
895 """
896 cache_key = _hash_x_token(x_token)
897 cached = self._token_cache.get(cache_key)
898 now = self._now()
899 if cached is not None and cached.expires_monotonic > now:
900 return cached.token
901
902 async with self._token_refresh_lock:
903 # Double-check inside the lock â a peer caller may have refreshed
904 # while we were waiting for the lock, in which case we reuse
905 # their fresh entry instead of issuing a duplicate Passport call.
906 cached = self._token_cache.get(cache_key)
907 now = self._now()
908 if cached is not None and cached.expires_monotonic > now:
909 return cached.token
910
911 token = await refresh_music_token(SecretStr(x_token))
912 self._store_cached_token(cache_key, token)
913 return token
914
915 def _store_cached_token(self, cache_key: str, token: SecretStr) -> None:
916 """
917 Insert a cache entry, enforcing the LRU bound.
918
919 Refreshing an existing key bumps its position to most-recent. When
920 a new key would push the cache over :data:`_MUSIC_TOKEN_CACHE_MAX`,
921 the oldest entry is evicted first.
922 """
923 # Reordering: pop-then-set positions the (possibly-new) key as
924 # most-recent in Python's insertion-ordered dict.
925 self._token_cache.pop(cache_key, None)
926 while len(self._token_cache) >= _MUSIC_TOKEN_CACHE_MAX:
927 oldest = next(iter(self._token_cache))
928 self._token_cache.pop(oldest)
929 self._token_cache[cache_key] = _CachedToken(
930 token=token,
931 expires_monotonic=self._now() + _MUSIC_TOKEN_TTL_S,
932 )
933
934 def _invalidate_cached_token(self, x_token: str) -> None:
935 """Drop the cache entry for an x_token (e.g. after a 401)."""
936 self._token_cache.pop(_hash_x_token(x_token), None)
937
938 async def _resolve_token(self) -> SecretStr:
939 """
940 Resolve the Yandex Music OAuth token for the Ynison connection.
941
942 In borrow mode: read from the linked yandex_music provider's config.
943 If only x_token is present (YM hasn't refreshed yet), do a cached
944 in-memory refresh without writing back â YM owns token persistence.
945
946 In own mode: return CONF_TOKEN if set; otherwise, when CONF_X_TOKEN
947 is present (QR-with-Remember-session path), cached in-memory refresh.
948 """
949 if self._borrow_source is not None:
950 return await self._borrow_source.resolve_music_token()
951
952 token = cast("str | None", self.get_setup_value(CONF_TOKEN))
953 if token:
954 return SecretStr(token)
955 x_token = cast("str | None", self.get_setup_value(CONF_X_TOKEN))
956 if x_token:
957 self.logger.debug("Own-mode token not present â refreshing from stored x_token")
958 return await self._refresh_via_x_token(x_token)
959 raise LoginFailed("No Yandex Music token configured")
960
961 async def _refresh_ynison_token(self) -> SecretStr:
962 """
963 Refresh the OAuth token for Ynison reconnection.
964
965 Called by YnisonClient on auth failure (401/403) during reconnect.
966
967 In borrow mode: re-read the linked YM instance's x_token and refresh
968 in-memory only (no config writes â YM owns token persistence).
969
970 In own mode: refresh from stored CONF_X_TOKEN when present (QR with
971 "Remember session" enabled). When absent (manual token paste only),
972 surface LoginFailed so the user knows to paste a new token.
973
974 The cached token entry for the current x_token is invalidated up
975 front â this method is reached only on a server-rejected token, so
976 the cached value is provably stale.
977 """
978 if self._borrow_source is not None:
979 ym_music_token, ym_x_token = self._borrow_source.read_tokens()
980 if ym_x_token is None:
981 raise LoginFailed("Cannot refresh: linked Yandex Music instance has no x_token")
982 # Both the minted entry AND the owner's persisted token may be the
983 # value the server just rejected â invalidate both so the source
984 # can't re-serve either; it will mint fresh from x_token.
985 if ym_music_token is not None:
986 self._borrow_source.invalidate(ym_music_token)
987 self._borrow_source.invalidate(ym_x_token)
988 self.logger.info("Refreshing Yandex Music token for Ynison reconnect (borrow mode)")
989 return await self._borrow_source.resolve_music_token()
990
991 x_token = cast("str | None", self.get_setup_value(CONF_X_TOKEN))
992 if x_token:
993 self._invalidate_cached_token(x_token)
994 self.logger.info("Refreshing Yandex Music token for Ynison reconnect (own mode)")
995 return await self._refresh_via_x_token(x_token)
996
997 raise LoginFailed(
998 "Token expired and no stored x_token to refresh from. Re-authenticate "
999 "via QR or paste a fresh Yandex Music token."
1000 )
1001
1002 # ------------------------------------------------------------------
1003 # Ynison state handling
1004 # ------------------------------------------------------------------
1005
1006 async def _handle_ynison_state(self, state: YnisonState) -> None:
1007 """Handle state update from Ynison."""
1008 is_our_device = state.active_device_id == self._device_id
1009
1010 # Detailed queue logging for diagnostics
1011 queue = state.player_state.get("player_queue", {})
1012 playable_list = queue.get("playable_list", [])
1013 current_index = queue.get("current_playable_index", -1)
1014 entity_type = queue.get("entity_type", "")
1015 entity_id = queue.get("entity_id", "")
1016 track_id = state.current_track_id
1017 self.logger.debug(
1018 "Ynison state: active_device=%s (ours=%s) track=%s "
1019 "index=%d/%d entity=%s type=%s paused=%s progress=%dms",
1020 state.active_device_id,
1021 is_our_device,
1022 track_id,
1023 current_index,
1024 len(playable_list),
1025 entity_id[:40] if entity_id else "<none>",
1026 entity_type,
1027 state.is_paused,
1028 state.progress_ms,
1029 )
1030
1031 # Post-reconnect settle window: the first inbound state after a WS
1032 # reconnect may reflect pre-reconnect peer state (active device etc).
1033 # Acting on it would re-issue play_media, mirror a stale paused flag
1034 # to MA, or worst case clobber a fresh local claim. The 2 s window in
1035 # YnisonClient._connect_state gives the server time to emit a state
1036 # broadcast that reflects our re-registered presence; until then we
1037 # only log.
1038 if self._ynison and self._ynison.in_post_reconnect_settle:
1039 self.logger.debug(
1040 "Skipping state inside post-reconnect settle window (track=%s paused=%s)",
1041 track_id,
1042 state.is_paused,
1043 )
1044 return
1045
1046 if is_our_device and not state.is_paused:
1047 self.logger.info(
1048 "Ynison â playing (track=%s progress=%dms)", track_id, state.progress_ms
1049 )
1050 # Pre-fetch next batch when playing second-to-last track
1051 self._maybe_prefetch(current_index, playable_list, entity_id, entity_type)
1052 await self._activate_playback(state)
1053 elif is_our_device and state.is_paused:
1054 self.logger.info(
1055 "Ynison â paused (track=%s progress=%dms)", track_id, state.progress_ms
1056 )
1057 await self._pause_playback()
1058 elif self._in_use_by_queue:
1059 self.logger.info(
1060 "Ynison â other device active (was=%s), clearing",
1061 state.active_device_id,
1062 )
1063 self._clear_active_player()
1064
1065 async def _activate_playback(self, state: YnisonState) -> None: # noqa: PLR0915
1066 """Activate playback on the target MA player."""
1067 target_player_id = self._get_target_player_id()
1068 if not target_player_id:
1069 self.logger.warning("Ynison active on our device but no MA player available")
1070 return
1071
1072 # Resume after pause / fresh start: either signal triggers
1073 # play_media below. `_externally_paused` survives a stray stop-event
1074 # clear; the stop event covers non-pause stop reasons
1075 # (`_stream_track` warning branch, `_clear_active_player`).
1076 needs_reselect = self._stream_stop_event.is_set() or self._externally_paused
1077 self._stream_stop_event.clear()
1078 self._externally_paused = False
1079
1080 # Start playback via the standard play_media flow if not already active.
1081 # Guard on _active_player_id (set immediately) rather than in_use_by_queue
1082 # (set by get_stream_details when the streams controller picks up the request)
1083 # to prevent queuing redundant play_media calls during the ~5s gap.
1084 if self._active_player_id != target_player_id or needs_reselect:
1085 # Pre-fetch the upcoming track's real format BEFORE submitting
1086 # play_media so the AudioSource's provider_mapping carries the
1087 # right audio_format when the streams controller calls
1088 # get_stream_details(). Skip on same-track same-player resume â
1089 # the cached format is still correct for that case.
1090 upcoming = state.current_track_id
1091 switching_player = self._active_player_id != target_player_id
1092 self._active_player_id = target_player_id
1093 if upcoming and (switching_player or upcoming != self._current_streaming_track_id):
1094 await self._prefetch_format_for_track(upcoming)
1095 self.mass.create_task(
1096 self.mass.player_queues.play_media(target_player_id, str(self._audio_source.uri))
1097 )
1098
1099 # Signal track change if track_id changed
1100 significant_change = False
1101 new_track = state.current_track_id
1102 if new_track and new_track != self._current_streaming_track_id:
1103 self.logger.info("Track changed: %s -> %s", self._current_streaming_track_id, new_track)
1104 self._current_streaming_track_id = new_track
1105 self._seek_position_ms = state.progress_ms
1106 self._track_changed_event.set()
1107 significant_change = True
1108 # Grace period: ignore seek detection for a few seconds after
1109 # track change â Ynison echoes can report stale progress that
1110 # looks like a large drift.
1111 self._seek_grace_until = time.monotonic() + _ECHO_GRACE_PERIOD
1112 elif new_track and new_track == self._current_streaming_track_id:
1113 # Same-track resume after pause: explicitly seek to the Ynison position
1114 # so the new stream starts at the right offset.
1115 if needs_reselect:
1116 self._seek_position_ms = state.progress_ms
1117 self._track_changed_event.set()
1118 self._seek_grace_until = time.monotonic() + _ECHO_GRACE_PERIOD
1119 significant_change = True
1120 else:
1121 # Detect seek: compare Ynison progress against our stream position.
1122 # Ignore Ynison echoes (updates authored by our own device_id) to
1123 # prevent feedback loops where our own progress triggers false seeks.
1124 now = time.monotonic()
1125 if now < self._seek_grace_until:
1126 pass # Skip during grace period after track change or seek
1127 elif state.last_update_is_echo:
1128 pass # Echo of our own update â ignore
1129 else:
1130 our_ms = self._streaming_progress_ms
1131 if our_ms >= 0:
1132 verdict = self._classify_drift(state.progress_ms, our_ms)
1133 if verdict == "seek":
1134 drift_ms = abs(state.progress_ms - our_ms)
1135 self.logger.info(
1136 "Seek detected on track %s: "
1137 "expected ~%dms, Ynison at %dms (drift %dms)",
1138 new_track,
1139 our_ms,
1140 state.progress_ms,
1141 int(drift_ms),
1142 )
1143 self._seek_position_ms = state.progress_ms
1144 self._track_changed_event.set()
1145 self._seek_grace_until = now + _ECHO_GRACE_PERIOD
1146 significant_change = True
1147 elif verdict == "queue_rebuild":
1148 self.logger.debug(
1149 "Drift on track %s classified as queue-rebuild "
1150 "echo (Ynison=%dms, ours=%dms) â not seeking",
1151 new_track,
1152 state.progress_ms,
1153 our_ms,
1154 )
1155
1156 # Update metadata from state
1157 self._update_metadata(state)
1158
1159 # Always trigger player update on significant changes;
1160 # throttle regular updates to avoid UI churn (every 5 seconds).
1161 # Use force_update on seek/track change so the server broadcasts a full
1162 # PLAYER_UPDATED event instead of a lightweight elapsed-time-only one
1163 # that the frontend may not handle for AudioSource players.
1164 now_mono = time.monotonic()
1165 if significant_change or needs_reselect or now_mono - self._last_player_update_time >= 5.0:
1166 self.mass.players.trigger_player_update(
1167 target_player_id, force_update=significant_change
1168 )
1169 self._last_player_update_time = now_mono
1170
1171 def _update_metadata(self, state: YnisonState) -> None:
1172 """Update AudioSource metadata from Ynison state."""
1173 meta = self._stream_metadata
1174
1175 # Update duration (prefer actual from stream_details) and elapsed time
1176 best_duration = self._best_duration_ms()
1177 if best_duration:
1178 meta.duration = best_duration // 1000
1179 # Only update elapsed from Ynison when NOT actively streaming â
1180 # during streaming, _sync_progress provides byte-accurate progress.
1181 if state.progress_ms is not None and not self._in_use_by_queue:
1182 meta.elapsed_time = state.progress_ms // 1000
1183 meta.elapsed_time_last_updated = time.time()
1184
1185 # Extract track info from player state if available
1186 queue = state.player_state.get("player_queue", {})
1187 playable_list = queue.get("playable_list", [])
1188 index = queue.get("current_playable_index", 0)
1189 if playable_list and 0 <= index < len(playable_list):
1190 playable = playable_list[index]
1191 title = playable.get("title")
1192 if title:
1193 meta.title = title
1194 cover = playable.get("cover_url_optional")
1195 if cover and not cover.startswith("http"):
1196 cover = f"https://{cover}"
1197 if cover:
1198 # Replace %% placeholder with size
1199 cover = cover.replace("%%", "400x400")
1200 meta.image_url = cover
1201
1202 async def _update_metadata_from_stream(
1203 self, stream_details: StreamDetails, seek_ms: int = 0
1204 ) -> None:
1205 """Update AudioSource metadata from stream details (authoritative for duration)."""
1206 meta = self._stream_metadata
1207 if stream_details.duration:
1208 meta.duration = stream_details.duration
1209 self._actual_duration_ms = stream_details.duration * 1000
1210 # Push the real duration to Ynison so the YM app shows
1211 # the correct value (we send duration_ms=0 on advance to
1212 # prevent stale propagation, so this corrects it).
1213 if self._ynison:
1214 await self._send_progress_to_ynison(
1215 progress_ms=seek_ms,
1216 duration_ms=self._actual_duration_ms,
1217 paused=self._ynison.state.is_paused,
1218 )
1219 meta.elapsed_time = seek_ms // 1000 if seek_ms else 0
1220 meta.elapsed_time_last_updated = time.time()
1221 # `trigger_player_update` expects a player_id; `_in_use_by_queue` is
1222 # a queue identifier which only happens to coincide with player_id
1223 # when there is no protocol bridge. Use `_active_player_id` â the
1224 # real player wrapping our stream (bridge if any).
1225 if self._active_player_id:
1226 self.mass.players.trigger_player_update(self._active_player_id, force_update=True)
1227
1228 async def _send_progress_to_ynison(
1229 self,
1230 progress_ms: int,
1231 duration_ms: int,
1232 paused: bool,
1233 *,
1234 strict: bool = False,
1235 ) -> None:
1236 """
1237 Send progress to Ynison.
1238
1239 Progress is clamped to duration because Ynison rejects updates where
1240 progress > duration (error 400030001) and disconnects the WebSocket.
1241 The byte counter can slightly overshoot duration at end-of-stream.
1242
1243 Echo detection is done upstream via YnisonState.last_update_is_echo,
1244 which is set when Ynison rebroadcasts an update we authored.
1245
1246 :param progress_ms: Current playback position in milliseconds.
1247 :param duration_ms: Current track duration in milliseconds.
1248 :param paused: Whether playback is paused.
1249 :param strict: When ``True``, propagate transport failures as
1250 :class:`provider.ynison_client.YnisonSendError`. Used by user-command
1251 and end-of-track callers. Heartbeat callers leave the default.
1252 """
1253 if duration_ms <= 0:
1254 # Ynison rejects progress > duration; skip until duration is known.
1255 return
1256 if not self._ynison or not self._ynison.connected:
1257 if strict:
1258 raise YnisonSendError("Ynison not connected")
1259 return
1260 progress_ms = min(progress_ms, duration_ms)
1261 await self._ynison.update_playing_status(
1262 progress_ms=progress_ms,
1263 duration_ms=duration_ms,
1264 paused=paused,
1265 strict=strict,
1266 )
1267
1268 def _bytes_to_ms(self, byte_count: int, fmt: AudioFormat | None = None) -> int:
1269 """Convert PCM byte count to milliseconds using the given format."""
1270 bps = (fmt or self._normalized_format).pcm_sample_size
1271 if bps == 0:
1272 return 0
1273 return (byte_count * 1000) // bps
1274
1275 async def _sync_progress(
1276 self,
1277 seek_ms: int,
1278 bytes_yielded: int,
1279 player_id: str | None,
1280 fmt: AudioFormat | None = None,
1281 ) -> None:
1282 """Push real playback progress to MA metadata and Ynison."""
1283 elapsed_ms = seek_ms + self._bytes_to_ms(bytes_yielded, fmt)
1284 self._streaming_progress_ms = elapsed_ms
1285 # Update MA metadata
1286 meta = self._stream_metadata
1287 if meta:
1288 meta.elapsed_time = elapsed_ms // 1000
1289 meta.elapsed_time_last_updated = time.time()
1290 if player_id:
1291 self.mass.players.trigger_player_update(player_id)
1292 # Update Ynison so the Yandex app shows correct position
1293 await self._send_progress_to_ynison(
1294 progress_ms=elapsed_ms,
1295 duration_ms=self._best_duration_ms(),
1296 paused=False,
1297 )
1298
1299 async def _pause_playback(self) -> None:
1300 """
1301 Release the active player on external pause.
1302
1303 ``cmd_stop`` is the only mechanism that flips ``PlaybackState``
1304 to IDLE for an AudioSource queue item; ``cmd_pause`` and
1305 ``queue.pause`` both short-circuit back to ``on_source_control``
1306 and leave MA's state untouched. Pattern matches upstream
1307 ``AriaCastReceiver._handle_playback_state_update``. Resume
1308 re-runs ``play_media`` (preload + ffmpeg startup) so it costs
1309 a few seconds â the alternative kept resume instant but left
1310 MA's UI stuck on PLAYING.
1311 """
1312 target = self._in_use_by_queue
1313 if not target:
1314 self.logger.info("Pause requested but no active queue (_in_use_by_queue is None)")
1315 return
1316 self.logger.info("Pause: cmd_stop(%s)", target)
1317 # stop event ends the audio generator; finally clears the lock.
1318 self._stream_stop_event.set()
1319 try:
1320 await self.mass.players.cmd_stop(target)
1321 except Exception:
1322 # cmd_stop is the only mechanism that flips MA's PlaybackState
1323 # to IDLE for an AudioSource. A silent failure here resurrects
1324 # the very UX bug this code path exists to fix.
1325 self.logger.warning(
1326 "cmd_stop(%s) failed during external pause â MA UI may stay PLAYING",
1327 target,
1328 exc_info=True,
1329 )
1330 return
1331 # Demote `_active_player_id` from the bridge MA streams to
1332 # (e.g. `spb_*`) back to the queue id; queues live on the bare
1333 # UUID. Without this, resume's `play_media(_active_player_id,
1334 # â¦)` would target the bridge and raise
1335 # `PlayerUnavailableError`. Post-success only so a failure
1336 # path keeps the bridge id intact for the next attempt.
1337 self._active_player_id = target
1338 self._externally_paused = True
1339
1340 # ------------------------------------------------------------------
1341 # Player selection
1342 # ------------------------------------------------------------------
1343
1344 def _get_target_player_id(self) -> str | None:
1345 """Determine the target player ID for playback."""
1346 # If there's an active player, validate it still exists
1347 if self._active_player_id:
1348 if self.mass.players.get_player(self._active_player_id):
1349 return self._active_player_id
1350 self._active_player_id = None
1351
1352 # Auto selection
1353 if self._default_player_id == PLAYER_ID_AUTO:
1354 all_players = list(self.mass.players.all_players(False, False))
1355 # Prefer currently playing player
1356 for player in all_players:
1357 if player.state.playback_state == PlaybackState.PLAYING:
1358 self.logger.debug("Auto-selecting playing player: %s", player.display_name)
1359 return str(player.player_id)
1360 # Fallback to first available
1361 if all_players:
1362 return str(all_players[0].player_id)
1363 return None
1364
1365 # Specific configured player
1366 if self.mass.players.get_player(self._default_player_id):
1367 return self._default_player_id
1368
1369 self.logger.warning(
1370 "Configured default player '%s' no longer exists",
1371 self._default_player_id,
1372 )
1373 return None
1374
1375 def _session_lost(self, player_id: str, session_id: str | None) -> bool:
1376 """
1377 Return ``True`` when our claim no longer matches the live session.
1378
1379 :param player_id: Queue id captured at generator entry.
1380 :param session_id: ``_active_session_id`` captured at generator entry.
1381 """
1382 return self._in_use_by_queue != player_id or self._active_session_id != session_id
1383
1384 def _idempotent(self, action: str, key: str | None) -> bool:
1385 """
1386 Return ``True`` if ``(action, key)`` was not seen within the TTL window.
1387
1388 :param action: A short string identifying the command kind.
1389 :param key: Sub-key inside the action namespace, or ``None``.
1390 """
1391 now = time.monotonic()
1392 for stale_key in [
1393 k for k, ts in self._command_idempotency.items() if now - ts > _COMMAND_IDEMPOTENCY_TTL
1394 ]:
1395 self._command_idempotency.pop(stale_key, None)
1396 composite = (action, key)
1397 last = self._command_idempotency.get(composite)
1398 if last is not None and now - last < _COMMAND_IDEMPOTENCY_TTL:
1399 return False
1400 self._command_idempotency[composite] = now
1401 return True
1402
1403 @staticmethod
1404 def _classify_drift(
1405 ynison_ms: int,
1406 our_ms: int,
1407 threshold_ms: int = 3000,
1408 ) -> Literal["ignore", "queue_rebuild", "seek"]:
1409 """
1410 Classify drift between Ynison-reported and our local position.
1411
1412 Returns one of:
1413
1414 - ``"ignore"`` â drift at or below ``threshold_ms``; no seek needed.
1415 - ``"queue_rebuild"`` â Ynison reports near-zero progress while we
1416 are past 5s into the track; treat as a RADIO queue-rebuild echo,
1417 not a user seek (otherwise we'd yank playback to the start every
1418 time the rotor station refills the queue).
1419 - ``"seek"`` â genuine drift; honor it.
1420
1421 :param ynison_ms: Position reported by Ynison in milliseconds.
1422 :param our_ms: Position tracked locally in milliseconds.
1423 :param threshold_ms: Minimum drift to consider non-ignorable.
1424 """
1425 drift = abs(ynison_ms - our_ms)
1426 if drift <= threshold_ms:
1427 return "ignore"
1428 if ynison_ms < 1000 and our_ms > 5000:
1429 return "queue_rebuild"
1430 return "seek"
1431
1432 async def _prefetch_format_for_track(self, track_id: str) -> None:
1433 """
1434 Pre-fetch stream details for *track_id* and adapt PCM format.
1435
1436 Best-effort: bounded by ``_PREFETCH_FORMAT_TIMEOUT`` so a slow Yandex
1437 API does not stall ``_activate_playback``. On timeout / error the
1438 current format stays in place and the in-stream
1439 ``_get_stream_details_with_retry`` handles retries.
1440
1441 :param track_id: Yandex Music track id to query.
1442 """
1443 if not self._yandex_provider:
1444 return
1445 try:
1446 stream_details = await asyncio.wait_for(
1447 self._get_stream_details_with_retry(track_id),
1448 timeout=_PREFETCH_FORMAT_TIMEOUT,
1449 )
1450 except TimeoutError:
1451 self.logger.info(
1452 "Pre-fetch of stream details for %s exceeded %.1fs â "
1453 "keeping current format; in-stream fetch will retry",
1454 track_id,
1455 _PREFETCH_FORMAT_TIMEOUT,
1456 )
1457 return
1458 except Exception:
1459 self.logger.warning(
1460 "Pre-fetch of stream details failed for %s â keeping current format",
1461 track_id,
1462 exc_info=True,
1463 )
1464 return
1465 old_sr = self._normalized_params.get("sample_rate")
1466 old_bd = self._normalized_params.get("bit_depth")
1467 self._update_normalized_format(hint=stream_details.audio_format)
1468 new_sr = self._normalized_params.get("sample_rate")
1469 new_bd = self._normalized_params.get("bit_depth")
1470 if (old_sr, old_bd) != (new_sr, new_bd):
1471 self.logger.info(
1472 "Pre-fetch adapted format for %s: %dHz/%dbit -> %dHz/%dbit (source=%s)",
1473 track_id,
1474 old_sr or 0,
1475 old_bd or 0,
1476 new_sr or 0,
1477 new_bd or 0,
1478 stream_details.audio_format,
1479 )
1480
1481 def _clear_active_player(self) -> None:
1482 """Clear the active player and reset plugin state."""
1483 prev_player_id = self._active_player_id
1484 was_in_use = self._in_use_by_queue == prev_player_id
1485 self._active_player_id = None
1486 self._in_use_by_queue = None
1487 self._active_session_id = None
1488 self._stream_stop_event.set()
1489 self._streaming_progress_ms = 0
1490 self._prefetched_list = None
1491 self._command_idempotency.clear()
1492 self._externally_paused = False
1493 if self._prefetch_task and not self._prefetch_task.done():
1494 self._prefetch_task.cancel()
1495
1496 if prev_player_id:
1497 self.logger.debug(
1498 "Playback ended on player %s, clearing active player",
1499 prev_player_id,
1500 )
1501 if was_in_use:
1502 self.mass.create_task(self.mass.players.cmd_stop(prev_player_id))
1503 self.mass.players.trigger_player_update(prev_player_id)
1504
1505 # ------------------------------------------------------------------
1506 # Yandex Music provider matching
1507 # ------------------------------------------------------------------
1508
1509 def _on_provider_event(self, event: MassEvent) -> None:
1510 """Handle provider added/removed events."""
1511 self.mass.create_task(self._check_yandex_provider_match())
1512
1513 async def _check_yandex_provider_match(self) -> None:
1514 """
1515 Check if a Yandex Music provider is available for audio streaming.
1516
1517 In borrow mode (self._ym_instance_id set), match strictly by instance_id
1518 so that audio and credentials come from the same account. In own mode,
1519 accept any yandex_music music-provider (prior behavior).
1520 """
1521 for provider in self.mass.get_providers():
1522 if provider.domain != "yandex_music" or provider.type != ProviderType.MUSIC:
1523 continue
1524 if self._ym_instance_id is not None and provider.instance_id != self._ym_instance_id:
1525 continue
1526 self.logger.debug("Found Yandex Music provider â enabling playback control")
1527 self._yandex_provider = cast("YandexMusicProviderLike", provider)
1528 self._update_normalized_format()
1529 self._update_source_capabilities()
1530 return
1531
1532 if self._yandex_provider is not None:
1533 self.logger.debug(
1534 "Yandex Music provider no longer available â disabling playback control"
1535 )
1536 self._yandex_provider = None
1537 self._update_source_capabilities()
1538
1539 def _snap_rate_to_player(self, rate: int) -> int:
1540 """
1541 Snap *rate* down to the nearest sample rate the target player accepts.
1542
1543 Best-effort: returns *rate* unchanged when no target player or
1544 supported-rate set can be resolved, and never raises.
1545
1546 :param rate: The sample rate the hint / floor logic chose.
1547 :return: A rate the target player can play (``rate`` itself when it is
1548 already supported or no player is resolvable).
1549 """
1550 # Mirror MA's _select_audio_source_pcm_format so the declared format
1551 # equals what the AudioSource passthrough picks â keeping MA off its
1552 # second resampling ffmpeg.
1553 try:
1554 player_id = self._get_target_player_id()
1555 if not player_id:
1556 return rate
1557 player = self.mass.players.get_player(player_id)
1558 if player is None:
1559 return rate
1560 supported = [sr for sr, _ in player.get_supported_sample_rates()]
1561 if not supported or rate in supported:
1562 return rate
1563 return max((r for r in supported if r <= rate), default=min(supported))
1564 except Exception:
1565 self.logger.debug(
1566 "Could not snap sample rate to player capabilities; keeping %d Hz",
1567 rate,
1568 exc_info=True,
1569 )
1570 return rate
1571
1572 def _update_normalized_format(self, hint: AudioFormat | None = None) -> None:
1573 """
1574 Set PCM normalization profile based on config and YM quality.
1575
1576 Priority: explicit config values > hint from real stream_details >
1577 auto-detection from YM quality. The hint is fed by
1578 ``_prefetch_format_for_track`` when ``CONF_OUTPUT_SAMPLE_RATE`` is
1579 ``auto`` so the AudioSource ``provider_mapping.audio_format`` matches
1580 the actual source rate of the upcoming track. Without a hint, falls
1581 back to YM-quality-based detection (superb/lossless â 24bit/44.1kHz,
1582 else â 16bit/44.1kHz). The resulting auto/hint rate is then snapped
1583 down to the nearest rate the target player supports; a valid explicit
1584 override is delivered verbatim and never snapped.
1585
1586 Creates fresh AudioFormat instances each time to prevent mutation by
1587 MA's FFMpeg._log_reader_task (which sets input_format.codec_type
1588 in-place on the object passed as input_format to the outer ffmpeg).
1589
1590 :param hint: Optional real source AudioFormat (from a stream-details
1591 pre-fetch). Lifts auto mode from the quality-based default to the
1592 track's actual sample rate and bit depth.
1593 """
1594 # Start with auto-detected base from YM quality config
1595 # (yandex_music does not expose get_quality(); read from its ProviderConfig instead)
1596 quality = ""
1597 if self._yandex_provider is not None:
1598 provider_config = getattr(self._yandex_provider, "config", None)
1599 if provider_config is not None and hasattr(provider_config, "get_value"):
1600 config_quality = provider_config.get_value(YANDEX_MUSIC_CONF_QUALITY)
1601 if isinstance(config_quality, str):
1602 quality = config_quality
1603 is_lossless = quality in YANDEX_MUSIC_LOSSLESS_QUALITIES
1604 base = dict(PCM_LOSSLESS_PARAMS if is_lossless else PCM_LOSSY_PARAMS)
1605 # Promote auto-base from the real stream details when available.
1606 # Validate the hint against the same allow-lists we use for explicit
1607 # config overrides â a Yandex API hiccup that returns an unsupported
1608 # rate (or 0) must not poison the AudioSource provider_mapping or the
1609 # outer ffmpeg input_format.
1610 if hint is not None:
1611 if hint.sample_rate and str(hint.sample_rate) in _VALID_SAMPLE_RATES:
1612 base["sample_rate"] = hint.sample_rate
1613 if hint.bit_depth and str(hint.bit_depth) in _VALID_BIT_DEPTHS:
1614 base["bit_depth"] = hint.bit_depth
1615
1616 # Apply config overrides. MA's ConfigEntry options constrain the UI to
1617 # known-good strings, but a stale persisted value or hand-edited config
1618 # could still surface something unparsable or off-list â fall back to
1619 # the auto-detected base with a warning instead of crashing the load.
1620 sample_rate = base["sample_rate"]
1621 bit_depth = base["bit_depth"]
1622 explicit_rate = False
1623 if self._cfg_sample_rate != OUTPUT_AUTO:
1624 if self._cfg_sample_rate in _VALID_SAMPLE_RATES:
1625 sample_rate = int(self._cfg_sample_rate)
1626 explicit_rate = True
1627 else:
1628 self.logger.warning(
1629 "Invalid %s=%r; falling back to auto-detected %d Hz",
1630 CONF_OUTPUT_SAMPLE_RATE,
1631 self._cfg_sample_rate,
1632 sample_rate,
1633 )
1634 # Snap the auto / hint / floor rate to a value the target player accepts
1635 # so the declared format matches what MA's AudioSource passthrough picks
1636 # and no second resampling ffmpeg is spawned. A valid explicit override
1637 # is delivered verbatim and is never snapped.
1638 if not explicit_rate:
1639 sample_rate = self._snap_rate_to_player(sample_rate)
1640 if self._cfg_bit_depth != OUTPUT_AUTO:
1641 if self._cfg_bit_depth in _VALID_BIT_DEPTHS:
1642 bit_depth = int(self._cfg_bit_depth)
1643 else:
1644 self.logger.warning(
1645 "Invalid %s=%r; falling back to auto-detected %d-bit",
1646 CONF_OUTPUT_BIT_DEPTH,
1647 self._cfg_bit_depth,
1648 bit_depth,
1649 )
1650
1651 content_type = ContentType.PCM_S24LE if bit_depth == 24 else ContentType.PCM_S16LE
1652 new_params: dict[str, Any] = {
1653 "content_type": content_type,
1654 "sample_rate": sample_rate,
1655 "bit_depth": bit_depth,
1656 "channels": 2,
1657 }
1658
1659 # Warn if format changes while a player is actively streaming â the
1660 # active session keeps using its frozen snapshot; the new format takes
1661 # effect on the next session.
1662 old = self._normalized_params
1663 if self._in_use_by_queue and (
1664 old.get("content_type") != content_type
1665 or old.get("sample_rate") != sample_rate
1666 or old.get("bit_depth") != bit_depth
1667 ):
1668 self.logger.warning(
1669 "Normalization format changed while streaming â new format "
1670 "(%s/%dHz/%dbit) will apply on next session",
1671 content_type.value,
1672 sample_rate,
1673 bit_depth,
1674 )
1675
1676 self._normalized_params = new_params
1677 # Fresh copy for each caller so no shared mutable state
1678 self._normalized_format = make_pcm_format(self._normalized_params)
1679 # rebuild the AudioSource so its ProviderMapping carries the new audio_format
1680 self._audio_source = self._build_audio_source()
1681 self.logger.debug(
1682 "Normalization format: %s/%dHz/%dbit",
1683 self._normalized_format.content_type.value,
1684 self._normalized_format.sample_rate,
1685 self._normalized_format.bit_depth,
1686 )
1687
1688 def _update_source_capabilities(self) -> None:
1689 """Rebuild AudioSource so capability flags reflect linked provider availability."""
1690 self._audio_source = self._build_audio_source()
1691 # The currently playing queue item carries a SNAPSHOT of the old
1692 # AudioSource â overwrite it so the new capability flags reach the UI
1693 # without waiting for the next play_media. Snapshot current_item and
1694 # re-check identity before the write so a queue advance racing this
1695 # callback can't stamp the new AudioSource onto an item that has
1696 # already moved on. Signal the queue update so the frontend re-renders
1697 # the controls (play/pause, next/prev) live.
1698 if not self._in_use_by_queue:
1699 return
1700 queue_id = self._in_use_by_queue
1701 queue = self.mass.player_queues.get(queue_id)
1702 if queue is None:
1703 return
1704 current_item = queue.current_item
1705 if (
1706 current_item is not None
1707 and current_item.media_item is not None
1708 and current_item.media_item.media_type == MediaType.AUDIO_SOURCE
1709 and current_item.media_item.item_id == AUDIO_SOURCE_ID
1710 and current_item.media_item.provider == self.instance_id
1711 and queue.current_item is current_item
1712 ):
1713 current_item.media_item = self._audio_source
1714 self.mass.player_queues.signal_update(queue_id, items_changed=True)
1715 self.mass.players.trigger_player_update(queue_id)
1716
1717 def _build_audio_source(self) -> AudioSource:
1718 """Construct the AudioSource MediaItem with current capability flags."""
1719 has_provider = self._yandex_provider is not None
1720 return AudioSource(
1721 item_id=AUDIO_SOURCE_ID,
1722 provider=self.instance_id,
1723 name=self.name,
1724 provider_mappings={
1725 ProviderMapping(
1726 item_id=AUDIO_SOURCE_ID,
1727 provider_domain=self.domain,
1728 provider_instance=self.instance_id,
1729 # Fresh AudioFormat copy â `self._normalized_format` is a
1730 # shared mutable that MA's ffmpeg sets `codec_type` on
1731 # in-place. Sharing it would let that mutation leak into
1732 # the rebuilt AudioSource and any future stream-details.
1733 audio_format=make_pcm_format(self._normalized_params),
1734 )
1735 },
1736 can_play_pause=has_provider,
1737 can_seek=has_provider,
1738 can_next_previous=has_provider,
1739 exclusive=True,
1740 allow_external_trigger=True,
1741 )
1742
1743 # ------------------------------------------------------------------
1744 # Playback control callbacks
1745 # ------------------------------------------------------------------
1746
1747 def _best_duration_ms(self) -> int:
1748 """Return the best known duration: actual from stream, or Ynison state as fallback."""
1749 if self._actual_duration_ms > 0:
1750 return self._actual_duration_ms
1751 if self._ynison:
1752 return self._ynison.state.duration_ms
1753 return 0
1754
1755 def _require_connected_ynison(self) -> YnisonClient:
1756 """
1757 Return the live Ynison client or raise an MA player-control error.
1758
1759 :raises UnsupportedFeaturedException: When the provider's Ynison
1760 client has not been initialised yet (pre-`handle_async_init`
1761 or post-`unload`).
1762 :raises PlayerCommandFailed: When the Ynison WebSocket is currently
1763 disconnected (e.g. mid-reconnect after a transient network
1764 error). Surface to MA so the UI shows a clear failure toast
1765 instead of accepting the command and stalling.
1766 """
1767 if not self._ynison:
1768 raise UnsupportedFeaturedException("Ynison client not initialized")
1769 if not self._ynison.connected:
1770 raise PlayerCommandFailed("Ynison WebSocket disconnected")
1771 return self._ynison
1772
1773 async def _on_play(self) -> None:
1774 """Handle play command â send resume to Ynison."""
1775 client = self._require_connected_ynison()
1776 if not self._idempotent("on_play", None):
1777 return
1778 state = client.state
1779 try:
1780 await self._send_progress_to_ynison(
1781 progress_ms=state.progress_ms,
1782 duration_ms=self._best_duration_ms(),
1783 paused=False,
1784 strict=True,
1785 )
1786 except YnisonSendError as exc:
1787 raise PlayerCommandFailed("Ynison send failed") from exc
1788
1789 async def _on_pause(self) -> None:
1790 """Handle pause command â send pause to Ynison."""
1791 client = self._require_connected_ynison()
1792 if not self._idempotent("on_pause", None):
1793 return
1794 state = client.state
1795 try:
1796 await self._send_progress_to_ynison(
1797 progress_ms=state.progress_ms,
1798 duration_ms=self._best_duration_ms(),
1799 paused=True,
1800 strict=True,
1801 )
1802 except YnisonSendError as exc:
1803 raise PlayerCommandFailed("Ynison send failed") from exc
1804
1805 # Entity types that use server-side "radio" queue replenishment.
1806 # Currently only RADIO (personal wave, genre stations).
1807 # Add "WAVE" here if/when Yandex supports it via the same
1808 # rotor_station_tracks API.
1809 _RADIO_ENTITY_TYPES: ClassVar[set[str]] = {"RADIO"}
1810
1811 def _maybe_prefetch(
1812 self,
1813 current_index: int,
1814 playable_list: list[dict[str, Any]],
1815 entity_id: str,
1816 entity_type: str,
1817 ) -> None:
1818 """Kick off background prefetch when nearing the end of the queue."""
1819 if entity_type not in self._RADIO_ENTITY_TYPES:
1820 return
1821 if not self._yandex_provider or not playable_list:
1822 return
1823 # second-to-last or last â trigger prefetch near end of queue
1824 if current_index < len(playable_list) - 2:
1825 return
1826 # Already prefetched or prefetch in progress
1827 if self._prefetched_list is not None:
1828 return
1829 if self._prefetch_task and not self._prefetch_task.done():
1830 return
1831
1832 self.logger.info(
1833 "Pre-fetching tracks (at index %d/%d, entity=%s)",
1834 current_index,
1835 len(playable_list),
1836 entity_id[:40] if entity_id else "<none>",
1837 )
1838
1839 async def _do_prefetch() -> None:
1840 result = await self._replenish_radio_queue(entity_id, entity_type, playable_list)
1841 if result:
1842 self._prefetched_list = result
1843 # Push expanded queue to Ynison immediately so the YM app
1844 # sees upcoming tracks and enables the "next" button.
1845 await self._update_queue_list(result)
1846
1847 self._prefetch_task = self.mass.create_task(_do_prefetch())
1848
1849 async def _signal_track_completion(self) -> None:
1850 """
1851 Signal that the current track finished playing.
1852
1853 Ynison is a state-sync protocol â the active device must advance
1854 current_playable_index itself.
1855
1856 If the next index is within the playable list, we advance immediately.
1857 If we're at the end (typical for RADIO/wave with short queues),
1858 we fetch more tracks via the Yandex Music API, append them to the
1859 playable_list, and then advance.
1860 """
1861 if not self._ynison:
1862 return
1863 state = self._ynison.state
1864 duration = self._best_duration_ms()
1865 queue = state.player_state.get("player_queue", {})
1866 current_index = queue.get("current_playable_index", 0)
1867 playable_list = queue.get("playable_list", [])
1868 entity_type = queue.get("entity_type", "")
1869 entity_id = queue.get("entity_id", "")
1870 next_index = current_index + 1
1871
1872 self.logger.info(
1873 "Track finished at index %d/%d (entity=%s type=%s), "
1874 "advancing to index %d (duration=%dms)",
1875 current_index,
1876 len(playable_list),
1877 entity_id[:40] if entity_id else "<none>",
1878 entity_type,
1879 next_index,
1880 duration,
1881 )
1882 self._actual_duration_ms = 0
1883
1884 # 1. Report that playback reached the end.
1885 # Echo tracking is handled by _send_progress_to_ynison.
1886 # `strict=True`: a dropped end-of-track signal stalls the YM app on
1887 # the just-finished track. We log and continue â the reconnect is
1888 # already scheduled and the queue-advance below sees the same WS state
1889 # â but we don't reraise (this is end-of-stream, there's no command to
1890 # fail back to the user).
1891 try:
1892 await self._send_progress_to_ynison(
1893 progress_ms=duration, duration_ms=duration, paused=False, strict=True
1894 )
1895 except YnisonSendError:
1896 self.logger.warning(
1897 "Track-completion signal dropped (Ynison transport failure); "
1898 "queue advance will retry once the WS reconnects",
1899 exc_info=True,
1900 )
1901
1902 if next_index < len(playable_list):
1903 # 2a. Queue has room â advance immediately.
1904 # Clear stale prefetch data so _maybe_prefetch can trigger for
1905 # the new queue tail on subsequent state updates.
1906 self._prefetched_list = None
1907 await self._advance_queue_index(next_index)
1908 elif entity_type in self._RADIO_ENTITY_TYPES:
1909 # 2b. At end of RADIO queue â use prefetched data or fetch now
1910 expanded: list[dict[str, Any]] | None = None
1911 if self._prefetched_list:
1912 self.logger.info("Using pre-fetched queue (%d items)", len(self._prefetched_list))
1913 expanded = self._prefetched_list
1914 self._prefetched_list = None
1915 elif self._prefetch_task and not self._prefetch_task.done():
1916 self.logger.info("Waiting for in-flight prefetch...")
1917 await self._prefetch_task
1918 expanded = self._prefetched_list
1919 self._prefetched_list = None
1920 else:
1921 expanded = await self._replenish_radio_queue(entity_id, entity_type, playable_list)
1922 if expanded and next_index < len(expanded):
1923 await self._advance_queue_index(next_index, expanded_list=expanded)
1924 elif expanded:
1925 self.logger.warning(
1926 "Expanded queue has %d items but next_index=%d â re-fetching",
1927 len(expanded),
1928 next_index,
1929 )
1930 fresh = await self._replenish_radio_queue(entity_id, entity_type, expanded)
1931 if fresh and next_index < len(fresh):
1932 await self._advance_queue_index(next_index, expanded_list=fresh)
1933 else:
1934 self.logger.warning("Still cannot advance after re-fetch")
1935 else:
1936 self.logger.warning(
1937 "Could not replenish queue (entity=%s type=%s), cannot advance",
1938 entity_id,
1939 entity_type,
1940 )
1941 else:
1942 self.logger.info(
1943 "End of non-radio queue (entity=%s type=%s), playback complete",
1944 entity_id[:40] if entity_id else "<none>",
1945 entity_type,
1946 )
1947
1948 async def _replenish_radio_queue(
1949 self,
1950 entity_id: str,
1951 entity_type: str,
1952 playable_list: list[dict[str, Any]],
1953 ) -> list[dict[str, Any]] | None:
1954 """
1955 Fetch more tracks from Yandex Music API and return expanded playable_list.
1956
1957 The active device is responsible for replenishing RADIO/wave queues.
1958 Ynison only syncs state â it does NOT generate new tracks.
1959 """
1960 if not self._yandex_provider:
1961 self.logger.warning("No yandex_music provider available for radio replenishment")
1962 return None
1963
1964 # Determine the last track ID for pagination
1965 last_track_id: str | None = None
1966 if playable_list:
1967 last_track_id = playable_list[-1].get("playable_id")
1968
1969 self.logger.info(
1970 "Fetching more tracks for %s station %s (queue=%s)",
1971 entity_type,
1972 entity_id,
1973 last_track_id,
1974 )
1975
1976 try:
1977 tracks, batch_id = await self._yandex_provider.get_rotor_station_tracks(
1978 entity_id, queue=last_track_id
1979 )
1980 except Exception:
1981 self.logger.exception("Failed to fetch radio tracks for %s", entity_id)
1982 return None
1983
1984 if not tracks:
1985 self.logger.warning("No tracks returned for station %s", entity_id)
1986 return None
1987
1988 # Determine the 'from' field from existing items
1989 from_field = ""
1990 if playable_list:
1991 from_field = playable_list[0].get("from", "")
1992
1993 # Convert tracks to Ynison playable_list format
1994 new_items: list[dict[str, Any]] = []
1995 for track in tracks:
1996 album_id = ""
1997 if hasattr(track, "albums") and track.albums:
1998 album_id = str(track.albums[0].id) if track.albums[0].id else ""
1999 cover = ""
2000 if hasattr(track, "cover_uri") and track.cover_uri:
2001 cover = track.cover_uri
2002 new_items.append(
2003 {
2004 "playable_id": str(track.id),
2005 "album_id_optional": album_id,
2006 "playable_type": "TRACK",
2007 "from": from_field,
2008 "title": track.title or "",
2009 "cover_url_optional": cover,
2010 }
2011 )
2012
2013 self.logger.info(
2014 "Fetched %d new tracks for station %s (batch=%s)",
2015 len(new_items),
2016 entity_id,
2017 batch_id,
2018 )
2019
2020 return list(playable_list) + new_items
2021
2022 async def _advance_queue_index(
2023 self,
2024 next_index: int,
2025 *,
2026 expanded_list: list[dict[str, Any]] | None = None,
2027 ) -> None:
2028 """
2029 Send update_player_state to advance the queue to next_index.
2030
2031 If expanded_list is provided, it replaces the playable_list
2032 (used after radio queue replenishment).
2033
2034 Waits up to 10 s for reconnection if Ynison is temporarily
2035 disconnected (e.g. after a transient error).
2036 """
2037 if not self._ynison:
2038 return
2039 if not self._ynison.connected:
2040 self.logger.info("Waiting for Ynison reconnection before advancing queueâ¦")
2041 for _ in range(10):
2042 await asyncio.sleep(1)
2043 if not self._ynison or self._ynison.connected:
2044 break
2045 if not self._ynison or not self._ynison.connected:
2046 self.logger.warning("Cannot advance queue â Ynison still disconnected")
2047 return
2048 state = self._ynison.state
2049 queue = state.player_state.get("player_queue", {})
2050 device_id = self._ynison.device_id
2051 new_state = dict(state.player_state)
2052 new_state["player_queue"] = dict(queue)
2053 new_state["player_queue"]["current_playable_index"] = next_index
2054 new_state["player_queue"]["version"] = make_version_block(device_id)
2055 if expanded_list is not None:
2056 new_state["player_queue"]["playable_list"] = expanded_list
2057 new_state["status"] = dict(new_state.get("status", {}))
2058 new_state["status"]["progress_ms"] = "0"
2059 new_state["status"]["duration_ms"] = "0"
2060 new_state["status"]["paused"] = False
2061 new_state["status"]["version"] = make_version_block(device_id)
2062 # `strict=True`: a dropped queue-advance leaves `_wait_for_track_change`
2063 # spinning for its full 30 s timeout. Log and return â the next
2064 # reconnect-broadcast picks up our authored version block and resyncs.
2065 try:
2066 await self._ynison.update_player_state(player_state=new_state, strict=True)
2067 except YnisonSendError:
2068 self.logger.warning(
2069 "Queue-advance dropped (Ynison transport failure); "
2070 "stream will stall until reconnect-broadcast resyncs",
2071 exc_info=True,
2072 )
2073
2074 async def _update_queue_list(self, expanded_list: list[dict[str, Any]]) -> None:
2075 """
2076 Push an expanded playable_list to Ynison without changing index or progress.
2077
2078 Called right after prefetch completes so the YM app sees upcoming
2079 tracks and enables the "next" button.
2080 """
2081 if not self._ynison or not self._ynison.connected:
2082 return
2083 state = self._ynison.state
2084 queue = state.player_state.get("player_queue", {})
2085 device_id = self._ynison.device_id
2086 new_state = dict(state.player_state)
2087 new_state["player_queue"] = dict(queue)
2088 new_state["player_queue"]["playable_list"] = expanded_list
2089 new_state["player_queue"]["version"] = make_version_block(device_id)
2090 await self._ynison.update_player_state(player_state=new_state)
2091
2092 async def _on_next(self) -> None:
2093 """Handle next track command â signal track end so Yandex advances."""
2094 self._require_connected_ynison()
2095 await self._signal_track_completion()
2096
2097 async def _on_previous(self) -> None:
2098 """Handle previous track command â update queue index in Ynison."""
2099 client = self._require_connected_ynison()
2100 queue = client.state.player_state.get("player_queue", {})
2101 current_index = queue.get("current_playable_index", 0)
2102 if current_index > 0:
2103 self._actual_duration_ms = 0
2104 await self._advance_queue_index(current_index - 1)
2105
2106 async def _on_seek(self, position: int) -> None:
2107 """
2108 Handle seek command â send position update to Ynison.
2109
2110 :param position: Position in seconds from Music Assistant.
2111 """
2112 client = self._require_connected_ynison()
2113 seek_ms = position * 1000
2114 state = client.state
2115 try:
2116 await self._send_progress_to_ynison(
2117 progress_ms=seek_ms,
2118 duration_ms=self._best_duration_ms(),
2119 paused=state.is_paused,
2120 strict=True,
2121 )
2122 except YnisonSendError as exc:
2123 # Do not mutate `_seek_position_ms` / `_seek_grace_until` on failure
2124 # â local stream state must not drift past a send that never landed.
2125 raise PlayerCommandFailed("Ynison send failed") from exc
2126 # Also trigger local stream restart so seek takes effect
2127 # immediately without waiting for the Ynison echo.
2128 self._seek_position_ms = seek_ms
2129 self._seek_grace_until = time.monotonic() + _ECHO_GRACE_PERIOD
2130 self._track_changed_event.set()
2131