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
3 changes: 3 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,9 @@ APP_YA_MUSIC_TOKEN=your_token_here
# PIN-код Ruark (по умолчанию 1234)
APP_RUARK_PIN=1234

# Необязательно: известный IP Ruark — быстрый старт без SSDP-скана сети
# APP_RUARK_HOST=192.168.1.20

# Адрес и порты сервисов (адрес — IP машины в локальной сети)
APP_LOCAL_SERVER_HOST=192.168.1.10
APP_LOCAL_SERVER_PORT_DLNA=8080
Expand Down
2 changes: 2 additions & 0 deletions src/core/config/settings.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@ class Settings(BaseSettings):

# Ruark R5 settings
ruark_pin: str
# Известный IP Ruark: быстрый старт без SSDP-скана всей сети
ruark_host: str | None = None

# Mode settings
debug: bool = False
Expand Down
22 changes: 17 additions & 5 deletions src/main_stream_service/main_stream_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -93,18 +93,28 @@ async def start(self):
return

logger.info("🎵 Запуск стриминга")
# Поиск Ruark и подключение к станции идут параллельно —
# это заметно ускоряет старт стрима
results: tuple[Any, Any] = await asyncio.gather(
self._ruark_controls.connect(),
self._station_controls.start_ws_client(),
return_exceptions=True,
)
ruark_connected, ws_result = results
try:
# Поиск устройств выполняется здесь, а не при старте процесса
if not await self._ruark_controls.connect():
if isinstance(ws_result, BaseException):
raise ws_result
if isinstance(ruark_connected, BaseException):
raise ruark_connected
if not ruark_connected:
raise RuarkDeviceNotFoundError(
f"Устройство "
f"'{self._ruark_controls.device_name}' "
f"не найдено в сети"
)
logger.info("🔄 Запуск WebSocket клиента")
await self._station_controls.start_ws_client()
except Exception as e:
logger.error(f"❌ Не удалось запустить стриминг: {e}")
await self._station_controls.stop_ws_client()
return

self._stream_state_running = True
Expand Down Expand Up @@ -652,7 +662,9 @@ async def _remember_ruark_volume(self):

async def _prepare_devices(self):
logger.info("🔧 Подготовка устройств к стримингу...")
await asyncio.sleep(1)
# Ждём первое сообщение станции вместо фиксированной паузы:
# обычно оно приходит сразу после подключения WebSocket
await self._ws_client.wait_for_state_update(1.0)
await self._station_controls.set_default_volume()
await self._ruark_controls.get_session_id()
if await self._ruark_controls.get_power_status() == "0":
Expand Down
35 changes: 35 additions & 0 deletions src/ruark_audio_system/location_store.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
from logging import getLogger
from pathlib import Path

logger = getLogger(__name__)

DEFAULT_STORE_PATH = Path("cache") / "ruark_location.txt"


class RuarkLocationStore:
"""Хранит location-URL Ruark между запусками для быстрого старта."""

def __init__(self, path: Path | None = None) -> None:
self._path = path or DEFAULT_STORE_PATH

def load(self) -> str | None:
"""Возвращает сохранённый location-URL или None, если его нет."""
try:
location = self._path.read_text().strip()
return location or None
except FileNotFoundError:
return None
except OSError as e:
logger.warning(
f"⚠️ Не удалось прочитать сохранённый адрес Ruark: {e}"
)
return None

def save(self, location: str) -> None:
"""Сохраняет location-URL для следующего запуска."""
try:
self._path.parent.mkdir(parents=True, exist_ok=True)
self._path.write_text(location)
logger.debug(f"💾 Адрес Ruark сохранён: {location}")
except OSError as e:
logger.warning(f"⚠️ Не удалось сохранить адрес Ruark: {e}")
73 changes: 71 additions & 2 deletions src/ruark_audio_system/ruark_r5_controller.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import asyncio
import socket
import urllib.parse
from logging import getLogger
from typing import Any, Dict, List, Literal, Optional
Expand All @@ -11,6 +12,8 @@
from ruark_audio_system.constants import DEFAULT_STREAM_TITLE, META_INFO
from ruark_audio_system.exceptions import RuarkDeviceNotFoundError
from ruark_audio_system.fsapi_client import RuarkFsApiClient
from ruark_audio_system.location_store import RuarkLocationStore
from ruark_audio_system.ssdp import ssdp_locate

logger = getLogger(__name__)

Expand Down Expand Up @@ -40,6 +43,7 @@ def __init__(self, device_name: str = "Ruark R5") -> None:
self._connection_manager: Any = None
self._rendering_control: Any = None
self._fsapi = RuarkFsApiClient(pin=settings.ruark_pin)
self._location_store = RuarkLocationStore()

@property
def av_transport(self) -> Any:
Expand Down Expand Up @@ -102,9 +106,15 @@ async def connect(self, attempts: int = 3, delay: float = 2.0) -> bool:
return False

def refresh_device(self) -> None:
"""Обновление устройства."""
"""Обновление устройства.

Сначала быстрый путь по известному адресу (настройка или кеш
прошлого запуска), при неудаче — полный SSDP-скан сети.
"""
logger.info("🔄 Обновление устройства")
self.device = self.find_device(device_name=self.device_name)
self.device = self._device_from_known_address()
if not self.device:
self.device = self.find_device(device_name=self.device_name)
if not self.device:
logger.warning(
f"⚠ Устройство '{self.device_name}' не найдено в сети!"
Expand All @@ -126,11 +136,70 @@ def refresh_device(self) -> None:
self._rendering_control = self.services.get(
"urn:schemas-upnp-org:service:RenderingControl:1"
)
self._location_store.save(self.device.location)
logger.info(
f"Устройство обновлено: {self.device.friendly_name} "
f"({self.device.location})"
)

def _device_from_known_address(self) -> Optional[upnpclient.Device]:
"""Подключение по известному адресу без сканирования сети.

Приоритет: IP из настроек (APP_RUARK_HOST), затем location
прошлого успешного поиска. Любая неудача — фолбэк на скан.
"""
candidates: list[str] = []
if settings.ruark_host:
location = ssdp_locate(settings.ruark_host)
if location:
candidates.append(location)
cached = self._location_store.load()
if cached and cached not in candidates:
candidates.append(cached)

for location in candidates:
device = self._device_at_location(location)
if device:
return device
return None

def _device_at_location(
self, location: str
) -> Optional[upnpclient.Device]:
"""Проверяет, что по location отвечает именно нужное устройство."""
try:
if not self._is_location_reachable(location):
return None
device = upnpclient.Device(location)
if self.device_name in device.friendly_name:
logger.info(
f"⚡ Ruark найден по известному адресу: {location}"
)
return device
logger.info(
f"ℹ️ По адресу {location} другое устройство: "
f"{device.friendly_name}"
)
except Exception as e:
logger.info(
f"ℹ️ Быстрое подключение по {location} не удалось: {e}"
)
return None

@staticmethod
def _is_location_reachable(location: str, timeout: float = 2.0) -> bool:
"""Быстрая TCP-проверка адреса перед HTTP-запросом описания."""
parsed = urllib.parse.urlparse(location)
if not parsed.hostname or not parsed.port:
return False
try:
with socket.create_connection(
(parsed.hostname, parsed.port), timeout=timeout
):
return True
except OSError:
return False

def find_device(self, device_name: str) -> Optional[upnpclient.Device]:
"""Находит устройство по имени."""
logger.info(f"Начинаем поиск устройства: {device_name}")
Expand Down
44 changes: 44 additions & 0 deletions src/ruark_audio_system/ssdp.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
import socket
from logging import getLogger

logger = getLogger(__name__)

SSDP_PORT = 1900

M_SEARCH_REQUEST = (
"M-SEARCH * HTTP/1.1\r\n"
"HOST: {host}:{port}\r\n"
'MAN: "ssdp:discover"\r\n'
"MX: 1\r\n"
"ST: upnp:rootdevice\r\n"
"\r\n"
)


def ssdp_locate(host: str, timeout: float = 2.0) -> str | None:
"""Запрашивает location описания UPnP-устройства напрямую у IP.

Юникастовый M-SEARCH вместо сканирования всей сети: устройство
с известным адресом отвечает за десятки миллисекунд.

Args:
host (str): IP-адрес устройства.
timeout (float): Ожидание ответа в секундах.
Returns:
str | None: URL описания устройства или None, если нет ответа.
"""
request = M_SEARCH_REQUEST.format(host=host, port=SSDP_PORT).encode()
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
sock.settimeout(timeout)
try:
sock.sendto(request, (host, SSDP_PORT))
data, _ = sock.recvfrom(4096)
for line in data.decode(errors="ignore").splitlines():
if line.lower().startswith("location:"):
return line.split(":", 1)[1].strip()
return None
except OSError as e:
logger.info(f"ℹ️ {host} не ответил на M-SEARCH: {e}")
return None
finally:
sock.close()
47 changes: 47 additions & 0 deletions src/yandex_station/station_store.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
import json
from logging import getLogger
from pathlib import Path

from yandex_station.mdns_device_finder import StationDevice

logger = getLogger(__name__)

DEFAULT_STORE_PATH = Path("cache") / "station_device.json"


class StationDeviceStore:
"""Хранит параметры найденной станции между запусками.

device_id и platform у станции постоянные, host почти всегда
закреплён DHCP — кеш позволяет пропустить mDNS-поиск при старте.
"""

def __init__(self, path: Path | None = None) -> None:
self._path = path or DEFAULT_STORE_PATH

def load(self) -> StationDevice | None:
"""Возвращает сохранённые параметры станции или None."""
try:
data = json.loads(self._path.read_text())
return StationDevice(
device_id=str(data["device_id"]),
platform=str(data["platform"]),
host=str(data["host"]),
port=int(data["port"]),
)
except FileNotFoundError:
return None
except (OSError, ValueError, KeyError, TypeError) as e:
logger.warning(
f"⚠️ Не удалось прочитать сохранённые параметры станции: {e}"
)
return None

def save(self, device: StationDevice) -> None:
"""Сохраняет параметры станции для следующего запуска."""
try:
self._path.parent.mkdir(parents=True, exist_ok=True)
self._path.write_text(json.dumps(dict(device)))
logger.debug(f"💾 Параметры станции сохранены: {device['host']}")
except OSError as e:
logger.warning(f"⚠️ Не удалось сохранить параметры станции: {e}")
36 changes: 34 additions & 2 deletions src/yandex_station/station_ws_control.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,8 @@
ClientNotRunningError,
StationNotFoundError,
)
from yandex_station.mdns_device_finder import DeviceFinder
from yandex_station.mdns_device_finder import DeviceFinder, StationDevice
from yandex_station.station_store import StationDeviceStore

logger = logging.getLogger(__name__)

Expand Down Expand Up @@ -52,6 +53,7 @@ def __init__(
self.state_updated = asyncio.Event()
self._connect_task: asyncio.Task[None] | None = None
self._connected_at: float | None = None
self._device_store = StationDeviceStore()
# Хранение фоновых задач
self.tasks: list[asyncio.Task[None]] = []

Expand All @@ -75,12 +77,23 @@ def _require_station_params(self) -> tuple[str, str, str]:
async def _ensure_device(self) -> None:
"""Находит станцию в сети, если она ещё не найдена.

Сначала проверяется кеш прошлого запуска (быстрая TCP-проверка
адреса), при неудаче — mDNS-поиск с сохранением результата.

Raises:
StationNotFoundError: Если станция не найдена за отведённое время.
"""
if self.device_id:
return

cached = self._device_store.load()
if cached and await self._is_station_reachable(
cached["host"], cached["port"]
):
self._apply_device(cached)
logger.info(f"⚡ Станция взята из кеша: {self.uri}")
return

logger.info("🔍 Поиск Яндекс Станции в сети...")
await asyncio.to_thread(self.device_finder.find_devices)
device = self.device_finder.device
Expand All @@ -90,10 +103,29 @@ async def _ensure_device(self) -> None:
"(mDNS-сервис _yandexio._tcp.local.)"
)

self._device_store.save(device)
self._apply_device(device)
logger.info(f"✅ Станция найдена: {self.uri}")

def _apply_device(self, device: StationDevice) -> None:
"""Заполняет параметры подключения из найденной станции."""
self.device_id = device["device_id"]
self.platform = device["platform"]
self.uri = f"wss://{device['host']}:{device['port']}"
logger.info(f"✅ Станция найдена: {self.uri}")

@staticmethod
async def _is_station_reachable(
host: str, port: int, timeout: float = 1.5
) -> bool:
"""Быстрая TCP-проверка, что станция отвечает по адресу из кеша."""
try:
_, writer = await asyncio.wait_for(
asyncio.open_connection(host, port), timeout
)
writer.close()
return True
except (OSError, asyncio.TimeoutError):
return False

async def run_once(self):
"""Гарантированный однократный запуск WebSocket."""
Expand Down
7 changes: 5 additions & 2 deletions tests/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,5 +21,8 @@ def mock_finder():

@pytest.fixture
def mock_station_client(mock_finder):
"""Клиент станции с замоканным DeviceFinder."""
return YandexStationClient(device_finder=mock_finder)
"""Клиент станции с замоканным DeviceFinder и пустым кешем."""
client = YandexStationClient(device_finder=mock_finder)
client._device_store = MagicMock()
client._device_store.load.return_value = None
return client
Loading
Loading