"""Обёртка сессии fgl-aircon: cffi-ядро + доставка событий в asyncio. События забирает выделенный python-поток (``fgl_session_wait_events``) и перекладывает их в event loop через ``loop.call_soon_threadsafe``. Прямые колбэки cffi из потоков ядра не используются: у session-потока ядра стек 8 КБ, вызов интерпретатора оттуда небезопасен. ``stop()/close()`` блокирующие (≤ ~2 c: DELETE-команда и join потоков ядра) — в async-коде используйте ``async_stop``. """ from __future__ import annotations import asyncio import logging import threading from dataclasses import dataclass from typing import Callable, Optional from ._cffi import ffi, lib from .const import Error, Prop, State, Template, ValueKind _LOGGER = logging.getLogger(__name__) StateHandler = Callable[[State, Error], None] PropertyHandler = Callable[["PropertyEvent"], None] LogHandler = Callable[[int, str], None] @dataclass(frozen=True) class Value: kind: ValueKind int_value: int = 0 str_value: str = "" @dataclass(frozen=True) class PropertyEvent: prop: Prop value: Value cmd_id: int = -1 status: int = 0 @dataclass class Config: host: str dsn: str lanip_key: str lanip_key_id: int template: Template = Template.A device_port: int = 80 listen_port: int = 10275 keepalive_ms: int = 15000 max_queue: int = 40 _GLOBAL_LOG: Optional[LogHandler] = None _LOG_PUMP: Optional["_LogPump"] = None _LOG_PUMP_LOCK = threading.Lock() class _LogPump(threading.Thread): """Единственный на процесс читатель кольца логов ядра.""" def __init__(self) -> None: super().__init__(name="pyfglair-log", daemon=True) def run(self) -> None: out = ffi.new("fgl_log_event_t[]", 8) dropped = 0 while True: total = int(lib.fgl_log_events_dropped()) if total > dropped: _LOGGER.warning( "переполнение очереди логов: потеряно %d", total - dropped ) dropped = total count = lib.fgl_log_wait_events(out, 8, 500) handler = _GLOBAL_LOG if handler is None: continue for i in range(count): try: handler( int(out[i].level), ffi.string(out[i].msg).decode("utf-8", "replace"), ) except Exception: _LOGGER.exception("log handler failed") def set_log_handler(handler: Optional[LogHandler]) -> None: """Глобальный приёмник логов ядра (0=debug..3=error).""" global _GLOBAL_LOG, _LOG_PUMP with _LOG_PUMP_LOCK: _GLOBAL_LOG = handler if handler is not None and _LOG_PUMP is None: _LOG_PUMP = _LogPump() _LOG_PUMP.start() def set_log_level(level: int) -> None: lib.fgl_log_set_level(int(level)) def _running_loop() -> Optional[asyncio.AbstractEventLoop]: try: return asyncio.get_running_loop() except RuntimeError: return None def _read_value(raw) -> Value: kind = ValueKind(int(raw.kind)) if kind == ValueKind.STRING: text = ffi.string(raw.s).decode("utf-8", "replace") return Value(kind, 0, text) return Value(kind, int(raw.i), "") class _SessionReader(threading.Thread): """Читает очередь событий сессии и передаёт их в ``Session``.""" def __init__(self, session: "Session") -> None: super().__init__(name="pyfglair-events", daemon=True) self._session = session def run(self) -> None: out = ffi.new("fgl_event_t[]", 16) finalize = ffi.NULL dropped = 0 try: while True: with self._session._lock: if self._session._closed: # Единственный владелец stop+destroy при закрытии — # reader (защищает от гонки с close() и позволяет # close() из колбэка в этом же потоке). finalize = self._session._session self._session._session = ffi.NULL return raw = self._session._session if raw == ffi.NULL: return total = int(lib.fgl_session_events_dropped(raw)) if total > dropped: _LOGGER.warning( "переполнение очереди событий: потеряно %d", total - dropped, ) dropped = total count = lib.fgl_session_wait_events(raw, out, 16, 250) for i in range(count): try: self._session._dispatch_event(out[i]) except Exception: _LOGGER.exception("event dispatch failed") finally: if finalize != ffi.NULL: lib.fgl_session_destroy(finalize) class Session: """Сессия к одному модулю кондиционера. Создание/использование — из event loop потока; методы ядра потокобезопасны. Колбэки вызываются в ``loop``. """ def __init__( self, config: Config, *, loop: Optional[asyncio.AbstractEventLoop] = None, on_state: Optional[StateHandler] = None, on_property: Optional[PropertyHandler] = None, ) -> None: self._config = config self._on_state = on_state self._on_property = on_property self._loop = loop if loop is not None else _running_loop() self._lock = threading.Lock() self._closed = False self._cfg_strings: list = [] def cstr(text: str): buf = ffi.new("char[]", text.encode("utf-8")) self._cfg_strings.append(buf) return buf cfg = ffi.new("fgl_config_t*") cfg.host = cstr(config.host) cfg.device_port = config.device_port cfg.dsn = cstr(config.dsn) cfg.lanip_key = cstr(config.lanip_key) cfg.lanip_key_id = config.lanip_key_id cfg.tmpl = int(config.template) cfg.listen_port = config.listen_port cfg.keepalive_ms = config.keepalive_ms cfg.max_queue = config.max_queue self._session = lib.fgl_session_create(cfg, ffi.NULL) if self._session == ffi.NULL: raise ValueError( "fgl_session_create вернул NULL: проверьте host/dsn/lanip_key " "и параметры сессии" ) self._reader = _SessionReader(self) self._reader.start() def _dispatch_event(self, ev) -> None: if int(ev.type) == 0: self._emit( self._on_state, State(int(ev.data.state.state)), Error(int(ev.data.state.error)), ) return kind = ValueKind(int(ev.data.property.value.kind)) if kind == ValueKind.STRING: value = Value( kind, 0, ffi.string(ev.data.property.value.s).decode("utf-8", "replace"), ) else: value = Value(kind, int(ev.data.property.value.i), "") self._emit( self._on_property, PropertyEvent( Prop(int(ev.data.property.prop)), value, int(ev.data.property.cmd_id), int(ev.data.property.status), ), ) def _emit(self, handler, *args) -> None: if handler is None: return loop = self._loop if loop is not None and loop.is_running(): try: loop.call_soon_threadsafe(handler, *args) return except RuntimeError: _LOGGER.debug("event loop закрыт, колбэк отброшен") return try: handler(*args) except Exception: _LOGGER.exception("callback failed") @property def state(self) -> State: with self._lock: if self._closed: return State.IDLE return State(int(lib.fgl_session_state(self._session))) @property def last_error(self) -> Error: with self._lock: if self._closed: return Error.NONE return Error(int(lib.fgl_session_last_error(self._session))) @property def closed(self) -> bool: return self._closed @property def events_dropped(self) -> int: """Сколько событий потеряно из-за переполнения очереди (диагностика).""" with self._lock: if self._closed or self._session == ffi.NULL: return 0 return int(lib.fgl_session_events_dropped(self._session)) @property def config(self) -> Config: return self._config def start(self) -> bool: with self._lock: if self._closed: return False return bool(lib.fgl_session_start(self._session)) def stop(self) -> None: """Полное завершение сессии (синоним close: stop + destroy).""" self.close() def close(self) -> None: """Останавливает сессию, джойнит reader и освобождает ядро. Из колбэка в reader-потоке (Session без loop) возвращается сразу: освобождение завершит сам reader. Повторный вызов — no-op. """ with self._lock: if self._closed: return self._closed = True if self._reader is threading.current_thread(): return self._reader.join(timeout=10.0) with self._lock: if self._reader.is_alive(): _LOGGER.error( "читатель событий не завершился за 10 с; сессия не " "освобождена" ) return raw = self._session self._session = ffi.NULL if raw != ffi.NULL: # Reader завершился, не освободив сессию (неожиданное исключение). lib.fgl_session_destroy(raw) def __del__(self) -> None: try: if not self._closed: _LOGGER.warning( "Session не закрыта: вызовите close()/async_stop()" ) except Exception: pass async def async_stop(self) -> None: loop = asyncio.get_running_loop() await loop.run_in_executor(None, self.close) def set_int(self, prop: Prop, value: int) -> bool: with self._lock: if self._closed: return False return bool(lib.fgl_session_set_int(self._session, int(prop), value)) def set_bool(self, prop: Prop, value: bool) -> bool: with self._lock: if self._closed: return False return bool( lib.fgl_session_set_bool( self._session, int(prop), 1 if value else 0 ) ) def set_string(self, prop: Prop, value: str) -> bool: with self._lock: if self._closed: return False return bool( lib.fgl_session_set_string( self._session, int(prop), value.encode("utf-8") ) ) def get_prop(self, prop: Prop) -> bool: with self._lock: if self._closed: return False return bool(lib.fgl_session_get_prop(self._session, int(prop))) def batch_begin(self) -> bool: with self._lock: if self._closed: return False return bool(lib.fgl_session_batch_begin(self._session)) def batch_commit(self) -> bool: with self._lock: if self._closed: return False return bool(lib.fgl_session_batch_commit(self._session)) def batch_abort(self) -> bool: with self._lock: if self._closed: return False return bool(lib.fgl_session_batch_abort(self._session)) def cached(self, prop: Prop) -> Optional[Value]: with self._lock: if self._closed: return None out = ffi.new("fgl_value_t*") if not lib.fgl_session_cached(self._session, int(prop), out): return None return _read_value(out) def __enter__(self) -> "Session": return self def __exit__(self, *exc_info) -> None: self.close()