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()