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