"""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.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 and not self._synced and self._sync_task is None: self._sync_task = self.hass.async_create_task( self._async_initial_sync() ) @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 await self.hass.async_add_executor_job(self._write, updates) 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]]) -> None: self.client.batch_begin() for prop, value in updates: self.client.set_int(prop, value) self.client.batch_commit() 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()