Files

307 lines
8.8 KiB
Python

import asyncio
from collections.abc import Awaitable, Callable
from contextlib import suppress
from datetime import datetime, timezone
from .controller import TionController
from .models import TionState
TionOperation = Callable[
[TionController],
Awaitable[TionState]
]
class TionService:
"""
Долгоживущий сервис работы с Tion.
Отвечает за:
- подключение;
- периодический опрос состояния;
- online/offline;
- last_seen;
- автоматическое переподключение;
- синхронизацию команд с polling.
"""
def __init__(
self,
controller: TionController,
poll_interval: float = 5.0,
):
if poll_interval <= 0:
raise ValueError("poll_interval must be greater than 0")
self._controller = controller
self._poll_interval = poll_interval
self._state: TionState | None = None
self._online = False
self._last_seen: datetime | None = None
self._last_error: str | None = None
self._running = False
self._poll_task: asyncio.Task | None = None
# Защищает последовательность:
#
# reconnect -> command -> update state
#
# от вмешательства polling или другой команды.
self._operation_lock = asyncio.Lock()
# Защищает start / stop.
self._lifecycle_lock = asyncio.Lock()
# ------------------------------------------------------------------
# Properties
# ------------------------------------------------------------------
@property
def state(self) -> TionState | None:
"""
Последнее успешно полученное состояние Tion.
Bluetooth-запрос не выполняется.
"""
return self._state
@property
def online(self) -> bool:
return self._online
@property
def running(self) -> bool:
return self._running
@property
def last_seen(self) -> datetime | None:
"""
Время последнего успешного обмена с Tion.
"""
return self._last_seen
@property
def last_error(self) -> str | None:
"""
Последняя ошибка связи.
После успешного обмена сбрасывается в None.
"""
return self._last_error
# ------------------------------------------------------------------
# Lifecycle
# ------------------------------------------------------------------
async def start(self) -> None:
"""
Запустить сервис.
Первая попытка подключения и чтения состояния выполняется сразу.
После этого запускается фоновый polling.
"""
async with self._lifecycle_lock:
if self._running:
return
self._running = True
# Сразу пытаемся получить состояние.
# Если Tion недоступен, сервис всё равно продолжит работу.
await self.refresh_state()
self._poll_task = asyncio.create_task(
self._poll_loop(),
name="tion-poll",
)
async def stop(self) -> None:
"""
Остановить polling и корректно закрыть BLE-соединение.
"""
async with self._lifecycle_lock:
if not self._running:
return
self._running = False
poll_task = self._poll_task
self._poll_task = None
if poll_task is not None:
poll_task.cancel()
with suppress(asyncio.CancelledError):
await poll_task
async with self._operation_lock:
await self._safe_disconnect()
self._online = False
# ------------------------------------------------------------------
# State
# ------------------------------------------------------------------
async def refresh_state(self) -> TionState | None:
"""
Принудительно обновить состояние Tion.
При ошибке:
- online становится False;
- last_error обновляется;
- старый state сохраняется;
- исключение наружу не выбрасывается.
Возвращает None при ошибке.
"""
async with self._operation_lock:
return await self._execute_locked(
lambda controller: controller.get_state(),
raise_on_error=False,
)
# ------------------------------------------------------------------
# Commands
# ------------------------------------------------------------------
async def execute(
self,
operation: TionOperation,
) -> TionState:
"""
Выполнить любую команду TionController.
Пример:
await service.execute(
lambda tion: tion.set_speed(3)
)
После команды состояние Service автоматически обновляется.
"""
async with self._operation_lock:
state = await self._execute_locked(
operation,
raise_on_error=True,
)
# Здесь None невозможен, потому что raise_on_error=True.
assert state is not None
return state
# ------------------------------------------------------------------
# Internal
# ------------------------------------------------------------------
async def _execute_locked(
self,
operation: TionOperation,
*,
raise_on_error: bool,
) -> TionState | None:
"""
Выполнение BLE-операции.
Вызывается только при занятом _operation_lock.
"""
try:
await self._ensure_connected()
state = await operation(self._controller)
except asyncio.CancelledError:
raise
except Exception as exc:
self._mark_offline(exc)
# После BLE-ошибки считаем соединение повреждённым.
# На следующей попытке будет создано новое.
await self._safe_disconnect()
if raise_on_error:
raise
return None
self._mark_online(state)
return state
async def _ensure_connected(self) -> None:
"""
Убедиться, что имеется рабочее BLE-соединение.
Если физического соединения нет, старое состояние подключения
сбрасывается и выполняется новое connect().
"""
if self._controller.connected:
return
await self._safe_disconnect()
await self._controller.connect()
async def _safe_disconnect(self) -> None:
"""
Закрыть соединение, не распространяя ошибку disconnect наружу.
"""
with suppress(Exception):
await self._controller.disconnect()
def _mark_online(self, state: TionState) -> None:
self._state = state
self._online = True
self._last_seen = datetime.now(timezone.utc)
self._last_error = None
def _mark_offline(self, exc: Exception) -> None:
self._online = False
self._last_error = (
f"{type(exc).__name__}: {exc}"
)
# ------------------------------------------------------------------
# Background polling
# ------------------------------------------------------------------
async def _poll_loop(self) -> None:
"""
Фоновый цикл обновления состояния.
"""
while self._running:
await asyncio.sleep(self._poll_interval)
if not self._running:
break
await self.refresh_state()
# ------------------------------------------------------------------
# Context manager
# ------------------------------------------------------------------
async def __aenter__(self) -> "TionService":
await self.start()
return self
async def __aexit__(
self,
exc_type,
exc_value,
traceback,
) -> None:
await self.stop()