Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
96 changes: 86 additions & 10 deletions src/main_stream_service/main_stream_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
from injector import inject

from core.config.settings import settings
from main_stream_service.yandex_music_api import YandexMusicAPI
from main_stream_service.yandex_music_api import TrackSource, YandexMusicAPI
from ruark_audio_system.exceptions import RuarkDeviceNotFoundError
from ruark_audio_system.ruark_r5_controller import RuarkR5Controller
from ruark_audio_system.volume_store import RuarkVolumeStore
Expand All @@ -20,6 +20,8 @@
SILENCE_RESEND_GRACE,
STREAM_POLL_INTERVAL,
STREAMING_RESTART_DELAY,
TRACK_SOURCE_CACHE_SIZE,
TRACK_SOURCE_TTL,
)
from yandex_station.models import Track
from yandex_station.station_controls import YandexStationControls
Expand Down Expand Up @@ -81,6 +83,8 @@ def __init__(
self._tasks = [] # Хранение фоновых задач
# Ссылка на задачу авто-остановки (Ruark выключен пользователем)
self._shutdown_task: asyncio.Task[None] | None = None
self._source_cache: dict[str, tuple[TrackSource, float]] = {}
self._prefetch_tasks: dict[str, asyncio.Task[None]] = {}

async def start(self):
"""Запуск всех стриминговых процессов."""
Expand Down Expand Up @@ -132,9 +136,17 @@ async def stop(self):
# Отмена всех активных задач
for task in self._tasks:
task.cancel()
for prefetch_task in self._prefetch_tasks.values():
prefetch_task.cancel()

await asyncio.gather(*self._tasks, return_exceptions=True)
await asyncio.gather(
*self._tasks,
*self._prefetch_tasks.values(),
return_exceptions=True,
)
self._tasks.clear()
self._prefetch_tasks.clear()
self._source_cache.clear()
logger.info("✅ Стриминг остановлен")

async def _safe_stop_step(
Expand Down Expand Up @@ -238,6 +250,7 @@ async def _handle_idle_cycle(
track = await self._refresh_track_if_unchanged(track, ctx)
await self._resync_after_progress_jump(track, ctx)
await self._switch_to_new_track(track, ctx)
self._schedule_prefetch(track)
await self._restore_ruark_after_speech(track, ctx)
await self._resume_track_if_silent(track, ctx)
await self._fade_alice_if_playing(track)
Expand Down Expand Up @@ -382,6 +395,75 @@ async def _is_ruark_powered_off(self) -> bool:
logger.debug(f"Не удалось узнать питание Ruark: {e}")
return False

def _get_cached_source(self, track_id: str) -> TrackSource | None:
"""Возвращает свежий источник из кеша или None, если протух."""
cached = self._source_cache.get(track_id)
if not cached:
return None
source, cached_at = cached
if time.monotonic() - cached_at > TRACK_SOURCE_TTL:
del self._source_cache[track_id]
return None
return source

def _store_source(self, track_id: str, source: TrackSource) -> None:
"""Кладёт источник в кеш, вытесняя самую старую запись."""
if len(self._source_cache) >= TRACK_SOURCE_CACHE_SIZE:
oldest = min(
self._source_cache,
key=lambda key: self._source_cache[key][1],
)
del self._source_cache[oldest]
self._source_cache[track_id] = (source, time.monotonic())

async def _resolve_track_source(self, track_id: str) -> TrackSource | None:
"""Источник трека: из кеша или свежим запросом с записью в кеш."""
cached = self._get_cached_source(track_id)
if cached:
logger.info(f"⚡ Ссылка для {track_id} взята из кеша")
return cached
source = await self._yandex_music_api.get_track_source(
track_id=track_id,
quality=settings.stream_quality,
)
if source:
self._store_source(track_id, source)
return source

def _schedule_prefetch(self, track: Track) -> None:
"""Фоново предзагружает ссылки соседних треков очереди.

Следующий трек нужен для мгновенного переключения вперёд,
предыдущий — для возврата назад. Уже закешированные и уже
загружаемые треки пропускаются.
"""
if track.type == "FmRadio":
return
for neighbor_id in (track.next_id, track.prev_id):
if not neighbor_id or neighbor_id == track.id:
continue
if self._get_cached_source(neighbor_id):
continue
if neighbor_id in self._prefetch_tasks:
continue
self._prefetch_tasks[neighbor_id] = asyncio.create_task(
self._prefetch_source(neighbor_id)
)

async def _prefetch_source(self, track_id: str) -> None:
"""Резолвит и кеширует ссылку трека в фоновой задаче."""
try:
source = await self._resolve_track_source(track_id)
if source:
logger.info(
f"📦 Предзагружена ссылка трека {track_id} "
f"({source.codec})"
)
except Exception as e:
logger.debug(f"Предзагрузка трека {track_id} не удалась: {e}")
finally:
self._prefetch_tasks.pop(track_id, None)

async def _resend_current_track(
self, track: Track, ctx: "_CycleContext", reason: str
) -> None:
Expand All @@ -390,10 +472,7 @@ async def _resend_current_track(
f"🔁 Ресинк стрима ({reason}): позиция {track.progress:.0f}s"
)
resync_started = time.monotonic()
source = await self._yandex_music_api.get_track_source(
track_id=track.id,
quality=settings.stream_quality,
)
source = await self._resolve_track_source(track.id)
if not source:
logger.warning("⚠️ Не удалось получить URL трека для ресинка")
return
Expand Down Expand Up @@ -431,10 +510,7 @@ async def _switch_to_new_track(
ctx.track_url = await self._station_controls.get_radio_url()
logger.info(f"🎵 URL радиостанции: {ctx.track_url}")
else:
source = await self._yandex_music_api.get_track_source(
track_id=track.id,
quality=settings.stream_quality,
)
source = await self._resolve_track_source(track.id)
ctx.track_url = source.url if source else None
codec = source.codec if source else "mp3"
resolve_seconds = time.monotonic() - switch_started
Expand Down
5 changes: 5 additions & 0 deletions src/yandex_station/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,3 +20,8 @@
# после отправки потока, прежде чем считать тишину проблемой, секунды
SILENCE_CHECK_INTERVAL = 3.0
SILENCE_RESEND_GRACE = 15.0

# Кеш прямых ссылок на треки: время жизни записи (CDN-ссылки
# протухают) и максимум записей, секунды / штуки
TRACK_SOURCE_TTL = 600.0
TRACK_SOURCE_CACHE_SIZE = 8
3 changes: 3 additions & 0 deletions src/yandex_station/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,3 +12,6 @@ class Track:
duration: float
progress: float
playing: bool
# Соседние треки очереди станции (entityInfo) — для предзагрузки
next_id: str | None = None
prev_id: str | None = None
21 changes: 21 additions & 0 deletions src/yandex_station/station_controls.py
Original file line number Diff line number Diff line change
Expand Up @@ -128,6 +128,7 @@ async def get_current_track(self) -> Track | None:
player_state = await self.get_player_status()
if not player_state:
return None
entity = player_state.get("entityInfo") or {}
return Track(
id=str(player_state.get("id", "0")),
title=str(player_state.get("title", "")),
Expand All @@ -136,11 +137,31 @@ async def get_current_track(self) -> Track | None:
duration=float(player_state.get("duration") or 0),
progress=float(player_state.get("progress") or 0),
playing=bool(player_state.get("playing", False)),
next_id=self._neighbor_track_id(entity, "next"),
prev_id=self._neighbor_track_id(entity, "prev"),
)
except Exception as e:
logger.error(f"❌ Ошибка при получении текущего трека: {e}")
return None

@staticmethod
def _neighbor_track_id(entity: Any, key: str) -> str | None:
"""Извлекает id соседнего трека из entityInfo станции.

Args:
entity (Any): Словарь entityInfo из playerState.
key (str): Ключ соседа: next или prev.
Returns:
str | None: Идентификатор, если сосед — трек.
"""
if not isinstance(entity, dict):
return None
neighbor = entity.get(key)
if not isinstance(neighbor, dict) or neighbor.get("type") != "Track":
return None
neighbor_id = neighbor.get("id")
return str(neighbor_id) if neighbor_id else None

async def get_volume(self) -> float | None:
"""Получение текущего уровня громкости (0.0–1.0)."""
try:
Expand Down
72 changes: 72 additions & 0 deletions tests/test_station_controls.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
from unittest.mock import AsyncMock, MagicMock

from yandex_station.protobuf_parser import Protobuf
from yandex_station.station_controls import YandexStationControls
from yandex_station.station_ws_control import YandexStationClient

BASE_PLAYER_STATE = {
"id": "42",
"title": "title",
"type": "Track",
"subtitle": "artist",
"duration": 200,
"progress": 10,
}


def make_controls(player_state) -> YandexStationControls:
"""Контролы станции с замоканным состоянием плеера."""
ws_client = AsyncMock(spec=YandexStationClient)
ws_client.get_latest_message.return_value = {
"state": {"playerState": player_state, "playing": True}
}
return YandexStationControls(
ws_client=ws_client,
protobuf=MagicMock(spec=Protobuf),
)


async def test_current_track_carries_neighbor_ids():
controls = make_controls(
{
**BASE_PLAYER_STATE,
"entityInfo": {
"next": {"id": "43", "type": "Track"},
"prev": {"id": "41", "type": "Track"},
},
}
)

track = await controls.get_current_track()

assert track is not None
assert track.next_id == "43"
assert track.prev_id == "41"


async def test_non_track_neighbors_are_ignored():
controls = make_controls(
{
**BASE_PLAYER_STATE,
"entityInfo": {
"next": {"id": "gen", "type": "Generative"},
"prev": "мусор",
},
}
)

track = await controls.get_current_track()

assert track is not None
assert track.next_id is None
assert track.prev_id is None


async def test_missing_entity_info_is_safe():
controls = make_controls(dict(BASE_PLAYER_STATE))

track = await controls.get_current_track()

assert track is not None
assert track.next_id is None
assert track.prev_id is None
Loading
Loading