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
37 changes: 36 additions & 1 deletion src/main_stream_service/main_stream_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@
SILENCE_CHECK_INTERVAL,
SILENCE_RESEND_GRACE,
STREAM_POLL_INTERVAL,
STREAMING_FAILURE_ANNOUNCE_AFTER,
STREAMING_FAILURE_STREAK_RESET,
STREAMING_RESTART_DELAY,
TRACK_SOURCE_CACHE_SIZE,
TRACK_SOURCE_TTL,
Expand All @@ -29,6 +31,12 @@

logger = getLogger(__name__)

# Фразы голосовых уведомлений об ошибках (озвучивает станция)
RUARK_MISSING_PHRASE = "Стрим не запустился: не вижу Руарк в сети"
STREAM_BROKEN_PHRASE = (
"Стрим сломался и сам не восстанавливается, загляни в логи"
)


@dataclass
class _CycleContext:
Expand Down Expand Up @@ -114,6 +122,9 @@ async def start(self):
)
except Exception as e:
logger.error(f"❌ Не удалось запустить стриминг: {e}")
# Станция доступна — озвучиваем причину перед остановкой
if not isinstance(ws_result, BaseException):
await self._announce_error(RUARK_MISSING_PHRASE)
await self._station_controls.stop_ws_client()
return

Expand Down Expand Up @@ -628,7 +639,14 @@ async def _unmute_before_track_end(
await self._station_controls.unmute()

async def _wrap_streaming(self):
"""Следит за потоком стриминга и перезапускает его при падении."""
"""Следит за потоком стриминга и перезапускает его при падении.

Серия быстрых падений подряд означает, что сам перезапуск
не помогает (например, умер DLNA-сервер) — один раз
озвучивается ошибка, чтобы тишина не осталась без объяснения.
"""
failure_streak = 0
last_failure_at = 0.0
while self._stream_state_running:
try:
logger.info("🚀 Запуск потока стриминга")
Expand All @@ -638,13 +656,30 @@ async def _wrap_streaming(self):
break
except Exception as e:
logger.error(f"❌ Поток стриминга упал с ошибкой: {e}")
now = time.monotonic()
if now - last_failure_at > STREAMING_FAILURE_STREAK_RESET:
failure_streak = 0
last_failure_at = now
failure_streak += 1
if failure_streak == STREAMING_FAILURE_ANNOUNCE_AFTER:
# Станция замьючена стримом — вернуть звук,
# иначе объявление не услышать
await self._safe_stop_step(self._station_controls.unmute)
await self._announce_error(STREAM_BROKEN_PHRASE)
logger.info(
f"🔁 Перезапуск стриминга через "
f"{STREAMING_RESTART_DELAY} секунд..."
)
await asyncio.sleep(STREAMING_RESTART_DELAY)
logger.debug("🔄 Перезапуск потока после падения")

async def _announce_error(self, phrase: str) -> None:
"""Озвучивает ошибку голосом станции, не роняя основной поток."""
try:
await self._station_controls.say(phrase)
except Exception as e:
logger.warning(f"⚠️ Не удалось озвучить ошибку: {e}")

async def _remember_ruark_volume(self):
"""Запоминает пользовательскую громкость Ruark для нового сеанса."""
try:
Expand Down
9 changes: 9 additions & 0 deletions src/yandex_station/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,3 +25,12 @@
# протухают) и максимум записей, секунды / штуки
TRACK_SOURCE_TTL = 600.0
TRACK_SOURCE_CACHE_SIZE = 8

# Локальный TTS: sendText с этим префиксом заставляет Алису
# произнести фразу вслух
TTS_REPEAT_PREFIX = "Повтори за мной"

# Голосовые уведомления об ошибках: сколько падений цикла подряд
# терпеть до объявления и пауза, сбрасывающая счётчик, секунды
STREAMING_FAILURE_ANNOUNCE_AFTER = 3
STREAMING_FAILURE_STREAK_RESET = 60.0
10 changes: 9 additions & 1 deletion src/yandex_station/station_controls.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,11 @@

from injector import inject

from yandex_station.constants import ALICE_ACTIVE_STATES, FADE_TIME
from yandex_station.constants import (
ALICE_ACTIVE_STATES,
FADE_TIME,
TTS_REPEAT_PREFIX,
)
from yandex_station.models import Track
from yandex_station.protobuf_parser import Protobuf
from yandex_station.station_ws_control import YandexStationClient
Expand Down Expand Up @@ -62,6 +66,10 @@ async def send_text(self, text: str):
except Exception as e:
logger.error(f"❌ Ошибка при отправке текстового сообщения: {e}")

async def say(self, phrase: str) -> None:
"""Произносит фразу голосом Алисы через локальный TTS."""
await self.send_text(f"{TTS_REPEAT_PREFIX} {phrase}")

async def get_current_state(self) -> dict[str, Any] | None:
"""Получение текущего состояния станции."""
try:
Expand Down
11 changes: 11 additions & 0 deletions tests/test_station_controls.py
Original file line number Diff line number Diff line change
Expand Up @@ -70,3 +70,14 @@ async def test_missing_entity_info_is_safe():
assert track is not None
assert track.next_id is None
assert track.prev_id is None


async def test_say_sends_repeat_command_to_station():
"""say() превращает фразу в команду локального TTS."""
controls = make_controls(dict(BASE_PLAYER_STATE))

await controls.say("тестовая фраза")

controls._ws_client.send_command.assert_awaited_once_with(
{"command": "sendText", "text": "Повтори за мной тестовая фраза"}
)
78 changes: 78 additions & 0 deletions tests/test_streaming_loop.py
Original file line number Diff line number Diff line change
Expand Up @@ -538,6 +538,84 @@ async def test_start_stops_ws_client_when_ruark_not_found():
station.stop_ws_client.assert_called_once()


async def test_start_announces_missing_ruark_by_voice():
"""Ruark не найден — станция голосом объясняет причину."""
station = make_station_controls(make_track())
ruark = make_ruark()
ruark.connect.return_value = False
ruark.device_name = "Ruark R5"
manager = make_manager(station, ruark)

await manager.start()

station.say.assert_awaited_once()
phrase = station.say.call_args.args[0]
assert "Руарк" in phrase


async def test_start_stays_silent_when_station_unreachable():
"""Недоступна сама станция — озвучивать ошибку нечем."""
station = make_station_controls(make_track())
station.start_ws_client.side_effect = RuntimeError("станция не найдена")
manager = make_manager(station, make_ruark())

await manager.start()

assert manager._stream_state_running is False
station.say.assert_not_called()


async def test_repeated_crashes_announce_error_once(fast_sleep, monkeypatch):
"""Серия падений цикла — одно голосовое уведомление, без спама."""
monkeypatch.setattr(
"main_stream_service.main_stream_manager.STREAMING_RESTART_DELAY",
0.0,
)
station = make_station_controls(make_track())
manager = make_manager(station, make_ruark())
manager.streaming = AsyncMock(side_effect=RuntimeError("бум"))
manager._stream_state_running = True

task = asyncio.create_task(manager._wrap_streaming())
deadline = time.monotonic() + 1.0
while time.monotonic() < deadline and not station.say.called:
await REAL_SLEEP(0.005)
# Даём циклу упасть ещё несколько раз после объявления
await REAL_SLEEP(0.05)
manager._stream_state_running = False
task.cancel()
await asyncio.gather(task, return_exceptions=True)

station.say.assert_awaited_once_with(
"Стрим сломался и сам не восстанавливается, загляни в логи"
)
station.unmute.assert_called()


async def test_short_failure_streak_stays_silent(fast_sleep, monkeypatch):
"""Пара падений с восстановлением — голосовое уведомление не нужно."""
monkeypatch.setattr(
"main_stream_service.main_stream_manager.STREAMING_RESTART_DELAY",
0.0,
)
station = make_station_controls(make_track())
manager = make_manager(station, make_ruark())
calls = {"count": 0}

async def flaky_streaming():
calls["count"] += 1
if calls["count"] <= 2:
raise RuntimeError("бум")
manager._stream_state_running = False

manager.streaming = flaky_streaming
manager._stream_state_running = True

await manager._wrap_streaming()

station.say.assert_not_called()


async def test_stop_survives_unreachable_ruark():
"""Недоступный Ruark не прерывает остановку стриминга."""
station = make_station_controls(make_track())
Expand Down
Loading