- close() из другого потока после close из reader-колбэка теперь джойнит reader (гарантия, что destroy завершён к возврату); - предупреждения о переполнении очередей не чаще 1/30 с; - events_dropped сохраняет последнее значение после close; - find_library выбирает .so по версии (а не по длине имени).
408 lines
14 KiB
Python
408 lines
14 KiB
Python
"""Обёртка сессии 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
|
|
import time
|
|
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
|
|
warned_at = 0.0
|
|
while True:
|
|
total = int(lib.fgl_log_events_dropped())
|
|
if total > dropped:
|
|
if time.monotonic() - warned_at > 30.0:
|
|
_LOGGER.warning(
|
|
"переполнение очереди логов: потеряно %d",
|
|
total - dropped,
|
|
)
|
|
warned_at = time.monotonic()
|
|
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
|
|
warned_at = 0.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))
|
|
self._session._dropped = total
|
|
if total > dropped:
|
|
if time.monotonic() - warned_at > 30.0:
|
|
_LOGGER.warning(
|
|
"переполнение очереди событий: потеряно %d",
|
|
total - dropped,
|
|
)
|
|
warned_at = time.monotonic()
|
|
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._dropped = 0
|
|
|
|
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 self._dropped
|
|
self._dropped = int(lib.fgl_session_events_dropped(self._session))
|
|
return self._dropped
|
|
|
|
@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. Повторный вызов из другого потока
|
|
дожидается фактического освобождения (join идемпотентен).
|
|
"""
|
|
with self._lock:
|
|
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()
|