"""Push-driven координатор: кэш свойств + доступность по состоянию сессии.""" from __future__ import annotations import asyncio import logging from dataclasses import dataclass, field, replace from typing import Optional from homeassistant.core import HomeAssistant, callback from homeassistant.exceptions import HomeAssistantError from homeassistant.helpers.update_coordinator import DataUpdateCoordinator from pyfglair import Error, Prop, PropertyEvent, State, Template, Value, ValueKind from pyfglair.templates import template_info from .const import AVAILABLE_STATES, CONF_DSN, DOMAIN from .fglair_client import FglairClient from .overrides import ConversionSet _LOGGER = logging.getLogger(__name__) @dataclass class FglairData: values: dict[Prop, Value] = field(default_factory=dict) state: State = State.IDLE last_error: Error = Error.NONE available: bool = False @dataclass class FglairRuntime: client: FglairClient coordinator: "FglairCoordinator" conversions: ConversionSet def _writable_props(template: Template): for info in template_info(template): if info.kind != ValueKind.STRING: yield info.prop class FglairCoordinator(DataUpdateCoordinator[FglairData]): """Держит актуальный снимок свойств и рассылает его сущностям.""" def __init__( self, hass: HomeAssistant, client: FglairClient, conversions: ConversionSet, ) -> None: super().__init__( hass, _LOGGER, name=f"{DOMAIN}-{client.data[CONF_DSN]}", update_interval=None, ) self.client = client self.conversions = conversions self._synced = False self._sync_task: Optional[asyncio.Task] = None self._unsubs = [ client.add_state_listener(self._on_state), client.add_property_listener(self._on_property), ] self.data = self._snapshot() # -- снимок ----------------------------------------------------------- def _snapshot(self) -> FglairData: values: dict[Prop, Value] = {} for prop in _writable_props(self.client.template): value = self.client.cached(prop) if value is not None: values[prop] = value state = self.client.state return FglairData( values=values, state=state, last_error=self.client.last_error, available=state in AVAILABLE_STATES, ) async def _async_update_data(self) -> FglairData: return await self.hass.async_add_executor_job(self._snapshot) # -- события ---------------------------------------------------------- @callback def _on_state(self, state: State, error: Error) -> None: self.data = replace( self.data, state=state, last_error=error, available=state in AVAILABLE_STATES, ) self.async_set_updated_data(self.data) if state == State.ONLINE: if not self._synced and self._sync_task is None: self._sync_task = self.hass.async_create_task( self._async_initial_sync() ) else: # После разрыва при следующем ONLINE нужна повторная синхронизация. self._synced = False async def async_sync_now(self) -> None: """Гарантирует первичный GET-батч (если уже online).""" if self._sync_task is not None: await self._sync_task elif self.client.state == State.ONLINE and not self._synced: self._sync_task = self.hass.async_create_task( self._async_initial_sync() ) await self._sync_task # Снимок создавался до ONLINE: перечитываем его без await (события # не могут вклиниться между чтением и присваиванием). self.data = self._snapshot() self.async_set_updated_data(self.data) @callback def _on_property(self, event: PropertyEvent) -> None: values = dict(self.data.values) values[event.prop] = event.value self.data = replace(self.data, values=values) self.async_set_updated_data(self.data) async def _async_initial_sync(self) -> None: try: await self.hass.async_add_executor_job(self._initial_sync) self._synced = True except Exception: # pragma: no cover - защитный путь _LOGGER.exception("первичная синхронизация свойств не удалась") finally: self._sync_task = None def _initial_sync(self) -> None: self.client.batch_begin() for prop in _writable_props(self.client.template): self.client.get_prop(prop) self.client.batch_commit() # -- записи ----------------------------------------------------------- async def async_write( self, updates: list[tuple[Prop, int]] ) -> None: """Одно действие пользователя — один batch_commit().""" if not updates: return ok = await self.hass.async_add_executor_job(self._write, updates) if not ok: raise HomeAssistantError( "FGLair: команда не отправлена (сессия недоступна)" ) values = dict(self.data.values) for prop, _ in updates: value = self.client.cached(prop) if value is not None: values[prop] = value self.data = replace(self.data, values=values) self.async_set_updated_data(self.data) def _write(self, updates: list[tuple[Prop, int]]) -> bool: if not self.client.batch_begin(): return False ok = True for prop, value in updates: if not self.client.set_int(prop, value): ok = False if not self.client.batch_commit(): ok = False return ok async def async_shutdown(self) -> None: for unsub in self._unsubs: unsub() self._unsubs.clear() if self._sync_task is not None: self._sync_task.cancel() self._sync_task = None await super().async_shutdown()