ha(H1): pyfglair — cffi-ядро (очередь событий), wheel, CLI discover/monitor
- CMake: shared-таргет fgl-aircon-shared (FGL_BUILD_SHARED, PIC, SOVERSION) - C-API: fgl_prop_raw_range + кольцевые очереди событий/логов (fgl_session_poll/wait_events, fgl_log_poll/wait_events): прямые колбэки из потоков ядра со стеком 8 КБ небезопасны для интерпретатора Python - pyfglair: _corebuild (cmake->wheel), _cffi (ABI), Session (reader-поток -> loop.call_soon_threadsafe), templates (интроспекция C-таблиц), provision (aiohttp), CLI discover/monitor + esphome-secrets - setup.py/pyproject: платформенный wheel py3-none-linux_x86_64 с .so - tests/pyfglair: templates, provision (мок-облако), session + CLI против mock_ac.py; scripts/py-ci.sh; ci.sh запускает python-часть
This commit is contained in:
@@ -0,0 +1,354 @@
|
||||
"""Обёртка сессии 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
|
||||
|
||||
|
||||
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)
|
||||
while True:
|
||||
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
|
||||
_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
|
||||
try:
|
||||
while True:
|
||||
with self._session._lock:
|
||||
if self._session._closed:
|
||||
# Освобождает сессию reader: close() мог быть вызван
|
||||
# из колбэка в этом же потоке (loop отсутствует).
|
||||
finalize = self._session._session
|
||||
self._session._session = ffi.NULL
|
||||
return
|
||||
raw = self._session._session
|
||||
if raw == ffi.NULL:
|
||||
return
|
||||
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 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:
|
||||
with self._lock:
|
||||
if self._closed:
|
||||
return
|
||||
raw = self._session
|
||||
lib.fgl_session_stop(raw)
|
||||
|
||||
def close(self) -> None:
|
||||
with self._lock:
|
||||
if self._closed:
|
||||
return
|
||||
self._closed = True
|
||||
raw = self._session
|
||||
lib.fgl_session_stop(raw)
|
||||
if self._reader is threading.current_thread():
|
||||
return # сессию освободит reader (см. _SessionReader.run)
|
||||
self._reader.join(timeout=3.0)
|
||||
if self._reader.is_alive():
|
||||
_LOGGER.error(
|
||||
"читатель событий не завершился; сессия не освобождена"
|
||||
)
|
||||
|
||||
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()
|
||||
Reference in New Issue
Block a user