diff --git a/CMakeLists.txt b/CMakeLists.txt index 47fd8b5..f1dac7b 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -72,7 +72,7 @@ else() set(USE_SHARED_MBEDTLS_LIBRARY OFF CACHE BOOL "" FORCE) set(USE_STATIC_MBEDTLS_LIBRARY ON CACHE BOOL "" FORCE) FetchContent_MakeAvailable(mbedtls) - set(FGL_MBEDTLS_TARGET mbedcrypto STATIC) + set(FGL_MBEDTLS_TARGET mbedcrypto) set(FGL_MBEDTLS_INCLUDE_DIR "") # Статический mbedcrypto нужен и для shared-библиотеки (pyfglair). set_target_properties(mbedcrypto PROPERTIES POSITION_INDEPENDENT_CODE ON) diff --git a/MANIFEST.in b/MANIFEST.in new file mode 100644 index 0000000..a01a29f --- /dev/null +++ b/MANIFEST.in @@ -0,0 +1,5 @@ +include CMakeLists.txt README.md pyproject.toml setup.py +recursive-include include *.h *.hpp +recursive-include src *.cpp *.hpp +recursive-include third_party * +recursive-include pyfglair *.py py.typed diff --git a/docs/PLAN_HOME_ASSISTANT.md b/docs/PLAN_HOME_ASSISTANT.md index aa22421..5c4a8c7 100644 --- a/docs/PLAN_HOME_ASSISTANT.md +++ b/docs/PLAN_HOME_ASSISTANT.md @@ -41,7 +41,7 @@ Fallback, если сборка wheel станет блокером: сборк * `pyfglair/provision.py` — облачный discovery (перенос логики `docs/legacy/aircon/discovery.py` на современный aiohttp): e-mail/ пароль/регион → dsn, host/ip, oem_model, lanip_key, lanip_key_id; - * CLI: `python -m pyfglair discover` / `monitor` / `--output esphome-secrets` + * CLI: `python -m pyfglair discover` / `monitor` / `--format esphome-secrets` (печать блока для secrets.yaml ESPHome). 2. **`custom_components/fglair/`**: ``` diff --git a/include/fgl-aircon/c_api.h b/include/fgl-aircon/c_api.h index 89c9f86..c9d7d55 100644 --- a/include/fgl-aircon/c_api.h +++ b/include/fgl-aircon/c_api.h @@ -160,7 +160,6 @@ typedef struct { fgl_error_t error; } state; fgl_property_event_t property; - fgl_log_event_t log; } data; } fgl_event_t; @@ -171,6 +170,9 @@ int fgl_session_poll_events(fgl_session_t* s, fgl_event_t* out, int max_events); // Возвращает число заполненных out. int fgl_session_wait_events(fgl_session_t* s, fgl_event_t* out, int max_events, int timeout_ms); +// Счётчик потерянных (переполнение кольца) событий сессии/логов. +uint64_t fgl_session_events_dropped(const fgl_session_t* s); +uint64_t fgl_log_events_dropped(void); // Логи ядра (общий кольцевой буфер; включается при создании C-API-сессии). int fgl_log_poll_events(fgl_log_event_t* out, int max_events); diff --git a/pyfglair/__main__.py b/pyfglair/__main__.py index 96380dd..2235cba 100644 --- a/pyfglair/__main__.py +++ b/pyfglair/__main__.py @@ -15,7 +15,7 @@ import sys from typing import Any, Optional from .const import State, Template, ValueKind -from .provision import REGIONS +from .provision import REGIONS, ProvisionError def _build_parser() -> argparse.ArgumentParser: @@ -200,13 +200,16 @@ async def _run_monitor(args: argparse.Namespace) -> int: def main(argv: Optional[list[str]] = None) -> int: args = _build_parser().parse_args(argv) - if args.command == "discover": - return asyncio.run(_run_discover(args)) - if args.command == "monitor": - try: + try: + if args.command == "discover": + return asyncio.run(_run_discover(args)) + if args.command == "monitor": return asyncio.run(_run_monitor(args)) - except KeyboardInterrupt: - return 130 + except ProvisionError as err: + print(f"pyfglair: {err}", file=sys.stderr) + return 2 + except KeyboardInterrupt: + return 130 return 2 diff --git a/pyfglair/_cffi.py b/pyfglair/_cffi.py index 3d6d31a..3a13c6e 100644 --- a/pyfglair/_cffi.py +++ b/pyfglair/_cffi.py @@ -7,6 +7,7 @@ from __future__ import annotations import ctypes.util import os +import re from pathlib import Path from cffi import FFI @@ -149,7 +150,6 @@ typedef struct { fgl_error_t error; } state; fgl_property_event_t property; - fgl_log_event_t log; } data; } fgl_event_t; @@ -177,6 +177,8 @@ int fgl_session_cached(const fgl_session_t* s, fgl_prop_t prop, int fgl_session_poll_events(fgl_session_t* s, fgl_event_t* out, int max_events); int fgl_session_wait_events(fgl_session_t* s, fgl_event_t* out, int max_events, int timeout_ms); +uint64_t fgl_session_events_dropped(const fgl_session_t* s); +uint64_t fgl_log_events_dropped(void); int fgl_log_poll_events(fgl_log_event_t* out, int max_events); int fgl_log_wait_events(fgl_log_event_t* out, int max_events, int timeout_ms); @@ -203,11 +205,13 @@ def find_library() -> str: raise FileNotFoundError(f"FGL_AIRCON_LIB={env} не существует") return env pkg = Path(__file__).resolve().parent + pattern = re.compile(r"^libfgl-aircon\.(so(\.[0-9.]+)?|dylib|\d+\.dylib)$") candidates = sorted( - p for p in pkg.glob("libfgl-aircon.so*") if p.is_file() + (p for p in pkg.iterdir() if p.is_file() and pattern.match(p.name)), + key=lambda p: (p.name.count("."), len(p.name)), ) if candidates: - return str(max(candidates, key=lambda p: p.stat().st_mtime)) + return str(candidates[0]) system = ctypes.util.find_library("fgl-aircon") if system: return system diff --git a/pyfglair/_corebuild.py b/pyfglair/_corebuild.py index 016cba2..b09f51d 100644 --- a/pyfglair/_corebuild.py +++ b/pyfglair/_corebuild.py @@ -61,10 +61,13 @@ def build_core( build_dir: Path | None = None, jobs: int | None = None, clean: bool = False, + bundled_mbedtls: bool | None = None, ) -> Path: """Собирает shared-ядро и (при dest) копирует его в каталог пакета. - Возвращает путь к скопированной/собранной библиотеке. + Для wheel mbedtls по умолчанию встраивается (bundled), чтобы .so не + зависел от системной soname; FGL_BUNDLED_MBEDTLS=0 — использовать + системный. """ src = Path(source) if source else source_root() if src is None: @@ -72,6 +75,9 @@ def build_core( "не найден корень fgl-aircon (CMakeLists.txt); задайте " "FGLAIR_SOURCE_ROOT" ) + if bundled_mbedtls is None: + env = os.environ.get("FGL_BUNDLED_MBEDTLS") + bundled_mbedtls = env != "0" src = Path(src).resolve() build = Path(build_dir) if build_dir else src / "build-pyfglair" if clean and build.exists(): @@ -84,6 +90,7 @@ def build_core( "-DFGL_BUILD_TESTS=OFF", "-DFGL_BUILD_EXAMPLES=OFF", "-DFGL_BUILD_SHARED=ON", + "-DFGL_BUNDLED_MBEDTLS=" + ("ON" if bundled_mbedtls else "OFF"), ]) _run([ "cmake", "--build", str(build), "--target", "fgl-aircon-shared", diff --git a/pyfglair/const.py b/pyfglair/const.py index 04fc043..8d4cd0e 100644 --- a/pyfglair/const.py +++ b/pyfglair/const.py @@ -73,9 +73,3 @@ class ValueKind(IntEnum): INT = 0 BOOL = 1 STRING = 2 - - -TEMPLATE_NAMES = {Template.A: "A", Template.B: "B", Template.F: "F"} - -STATE_NAMES = {state: state.name.lower() for state in State} -ERROR_NAMES = {error: error.name.lower() for error in Error} diff --git a/pyfglair/provision.py b/pyfglair/provision.py index 5cc442b..713a754 100644 --- a/pyfglair/provision.py +++ b/pyfglair/provision.py @@ -140,13 +140,12 @@ async def _get_devices(session, token, region, base_url, insecure): ) if status != 200 or not isinstance(data, list): raise ProvisionError(f"Ошибка списка устройств ({status}): {data}") - devices = [ - item.get("device", {}) - for item in data - if isinstance(item, dict) - ] - if any(not dev.get("dsn") for dev in devices): - raise ProvisionError("Облако вернуло устройство без dsn") + devices = [] + for item in data: + dev = item.get("device") if isinstance(item, dict) else None + if not isinstance(dev, dict) or not isinstance(dev.get("dsn"), str): + raise ProvisionError("Облако вернуло неожиданный формат устройства") + devices.append(dev) return devices @@ -158,7 +157,12 @@ async def _get_lanip(session, token, region, base_url, insecure, dsn): ) if status != 200 or not isinstance(data, dict): raise ProvisionError(f"Ошибка lan.json для {dsn} ({status}): {data}") - return data.get("lanip") or {} + lanip = data.get("lanip") + if lanip is not None and not isinstance(lanip, dict): + raise ProvisionError( + f"Облако вернуло неожиданный lanip для {dsn}: {lanip!r}" + ) + return lanip or {} async def discover( diff --git a/pyfglair/session.py b/pyfglair/session.py index d95f27b..a0c4526 100644 --- a/pyfglair/session.py +++ b/pyfglair/session.py @@ -56,6 +56,7 @@ class Config: _GLOBAL_LOG: Optional[LogHandler] = None _LOG_PUMP: Optional["_LogPump"] = None +_LOG_PUMP_LOCK = threading.Lock() class _LogPump(threading.Thread): @@ -66,7 +67,14 @@ class _LogPump(threading.Thread): 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: @@ -84,10 +92,11 @@ class _LogPump(threading.Thread): 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() + 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: @@ -119,18 +128,27 @@ class _SessionReader(threading.Thread): 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: - # Освобождает сессию reader: close() мог быть вызван - # из колбэка в этом же потоке (loop отсутствует). + # Единственный владелец 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: @@ -252,6 +270,14 @@ class Session: 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 @@ -263,26 +289,43 @@ class Session: 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) + """Полное завершение сессии (синоним 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 - 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( - "читатель событий не завершился; сессия не освобождена" - ) + 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() diff --git a/pyproject.toml b/pyproject.toml index 18b61d9..573b3f9 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -6,6 +6,7 @@ build-backend = "setuptools.build_meta" name = "pyfglair" version = "1.0.0" description = "Local control of Fujitsu General (FGLair/Ayla) air conditioners: cffi bindings to the fgl-aircon C++ core" +readme = "README.md" requires-python = ">=3.10" dependencies = [ "cffi>=1.15", @@ -17,13 +18,16 @@ classifiers = [ "Topic :: Home Automation", ] +[project.scripts] +pyfglair = "pyfglair.__main__:main" + [tool.setuptools] packages = ["pyfglair"] include-package-data = true zip-safe = false [tool.setuptools.package-data] -pyfglair = ["libfgl-aircon.so*", "py.typed"] +pyfglair = ["py.typed"] [tool.pytest.ini_options] testpaths = ["tests/pyfglair"] diff --git a/setup.py b/setup.py index 521f623..ba4f108 100644 --- a/setup.py +++ b/setup.py @@ -30,7 +30,9 @@ class build_py(_build_py): def run(self): dest = Path(self.build_lib) / "pyfglair" dest.mkdir(parents=True, exist_ok=True) - _load_corebuild().build_core(dest=dest) + for stale in dest.glob("libfgl-aircon.*"): + stale.unlink() + _load_corebuild().build_core(dest=dest, bundled_mbedtls=True) super().run() diff --git a/src/aircon/c_api.cpp b/src/aircon/c_api.cpp index 0af1197..ec9406a 100644 --- a/src/aircon/c_api.cpp +++ b/src/aircon/c_api.cpp @@ -2,6 +2,7 @@ // обязаны совпадать с C++ enum'ами — проверено static_assert'ами. #include "fgl-aircon/c_api.h" +#include #include #include #include @@ -110,6 +111,11 @@ struct Ring { } return n; } + + uint64_t dropped_count() { + std::lock_guard lock(mtx); + return dropped; + } }; constexpr size_t kSessionEventCap = 64; @@ -118,8 +124,8 @@ constexpr size_t kLogEventCap = 64; Ring g_log_ring; struct LogWrap { - fgl_log_sink_fn fn; - void* ctx; + std::atomic fn{nullptr}; + std::atomic ctx{nullptr}; }; LogWrap g_log_wrap; std::once_flag g_log_installed; @@ -131,8 +137,9 @@ void log_wrap_trampoline(int level, const char* msg, size_t len, void*) { memcpy(ev.msg, msg, n); ev.msg[n] = '\0'; g_log_ring.push(ev); - if (g_log_wrap.fn != nullptr) { - g_log_wrap.fn(level, msg, len, g_log_wrap.ctx); + auto fn = g_log_wrap.fn.load(std::memory_order_acquire); + if (fn != nullptr) { + fn(level, msg, len, g_log_wrap.ctx.load(std::memory_order_relaxed)); } } @@ -144,8 +151,8 @@ void ensure_log_pipe() { } // namespace void fgl_log_set_sink(fgl_log_sink_fn sink, void* ctx) { - g_log_wrap.fn = sink; - g_log_wrap.ctx = ctx; + g_log_wrap.ctx.store(ctx, std::memory_order_relaxed); + g_log_wrap.fn.store(sink, std::memory_order_release); // Ядро всегда пишет в кольцевой буфер (его читает pyfglair); sink — // дополнительный прямой колбэк для C/C++ (только малые стеки!). ensure_log_pipe(); @@ -292,6 +299,14 @@ int fgl_session_wait_events(fgl_session_t* s, fgl_event_t* out, return s->events->wait(out, max_events, timeout_ms); } +uint64_t fgl_session_events_dropped(const fgl_session_t* s) { + return (s != nullptr && s->events != nullptr) + ? s->events->dropped_count() + : 0; +} + +uint64_t fgl_log_events_dropped(void) { return g_log_ring.dropped_count(); } + fgl_state_t fgl_session_state(const fgl_session_t* s) { return s != nullptr ? static_cast(s->s->state()) : FGL_STATE_IDLE; diff --git a/tests/pyfglair/conftest.py b/tests/pyfglair/conftest.py index 29f588c..c3569e3 100644 --- a/tests/pyfglair/conftest.py +++ b/tests/pyfglair/conftest.py @@ -90,7 +90,13 @@ def ac_env() -> dict: @pytest.fixture async def cloud(): """Мок облака Ayla: sign_in / devices / lan.json.""" - state = {"url": "", "fail_login": False, "no_key": False} + state = { + "url": "", + "fail_login": False, + "no_key": False, + "bad_devices": False, + "bad_lanip": False, + } async def sign_in(request): if state["fail_login"]: @@ -99,6 +105,8 @@ async def cloud(): async def devices(request): assert request.headers["Authorization"] == "auth_token tok-123" + if state["bad_devices"]: + return web.json_response([{"device": "not-an-object"}]) return web.json_response([ { "device": { @@ -114,6 +122,8 @@ async def cloud(): async def lan(request): if state["no_key"]: return web.json_response({"lanip": {}}) + if state["bad_lanip"]: + return web.json_response({"lanip": ["unexpected"]}) return web.json_response({ "lanip": {"lanip_key": "TW9ja0tleQ==", "lanip_key_id": KEY_ID} }) diff --git a/tests/pyfglair/test_cli.py b/tests/pyfglair/test_cli.py index 85f52c2..018985a 100644 --- a/tests/pyfglair/test_cli.py +++ b/tests/pyfglair/test_cli.py @@ -40,6 +40,38 @@ async def test_cli_discover_esphome_secrets(cloud, capfd): assert "bedroom_lanip_key_id: 64201" in out +async def test_cli_discover_error(cloud, capfd): + cloud["fail_login"] = True + rc = await asyncio.to_thread( + main, + [ + "discover", "--api-base", cloud["url"], + "--email", "user@example.com", "--password", "bad", + ], + ) + assert rc == 2 + err = capfd.readouterr().err + assert "Ошибка входа" in err + assert "Traceback" not in err + + +async def test_cli_discover_out_file(cloud, capfd, tmp_path): + out_file = tmp_path / "config_bedroom.json" + rc = await asyncio.to_thread( + main, + [ + "discover", "--api-base", cloud["url"], + "--email", "user@example.com", "--password", "secret", + "--out", str(out_file), + ], + ) + assert rc == 0 + capfd.readouterr() + assert out_file.stat().st_mode & 0o777 == 0o600 + config = json.loads(out_file.read_text()) + assert config["dsn"] == "AC000W00MOCK0001" + + async def test_cli_monitor(mock_ac, ac_env, capfd): proc = mock_ac() rc = await asyncio.to_thread( diff --git a/tests/pyfglair/test_core_api.py b/tests/pyfglair/test_core_api.py new file mode 100644 index 0000000..ff51aaf --- /dev/null +++ b/tests/pyfglair/test_core_api.py @@ -0,0 +1,55 @@ +"""Низкоуровневые проверки C-API: очередь событий, счётчики, ошибки входа.""" +from __future__ import annotations + +from pyfglair._cffi import ffi, lib + + +def _make_session(): + cfg = ffi.new("fgl_config_t*") + cfg.host = ffi.new("char[]", b"127.0.0.1") + cfg.device_port = 9 + cfg.dsn = ffi.new("char[]", b"AC000W00MOCK0001") + cfg.lanip_key = ffi.new("char[]", b"key") + cfg.lanip_key_id = 1 + cfg.tmpl = 0 + cfg.listen_port = 0 + cfg.keepalive_ms = 1000 + cfg.max_queue = 40 + session = lib.fgl_session_create(cfg, ffi.NULL) + assert session != ffi.NULL + return session + + +def test_poll_wait_empty(): + session = _make_session() + try: + out = ffi.new("fgl_event_t[]", 4) + assert lib.fgl_session_poll_events(session, out, 4) == 0 + assert lib.fgl_session_wait_events(session, out, 4, 0) == 0 + assert lib.fgl_session_events_dropped(session) == 0 + finally: + lib.fgl_session_destroy(session) + + +def test_poll_invalid_args(): + session = _make_session() + try: + out = ffi.new("fgl_event_t[]", 4) + assert lib.fgl_session_poll_events(session, ffi.NULL, 4) == 0 + assert lib.fgl_session_poll_events(session, out, 0) == 0 + assert lib.fgl_session_wait_events(session, out, 0, 0) == 0 + assert lib.fgl_session_poll_events(ffi.NULL, out, 4) == 0 + assert lib.fgl_session_events_dropped(ffi.NULL) == 0 + finally: + lib.fgl_session_destroy(session) + + +def test_create_null_config(): + assert lib.fgl_session_create(ffi.NULL, ffi.NULL) == ffi.NULL + + +def test_log_poll_empty(): + out = ffi.new("fgl_log_event_t[]", 4) + assert lib.fgl_log_poll_events(out, 4) >= 0 + assert lib.fgl_log_wait_events(out, 4, 0) >= 0 + assert lib.fgl_log_events_dropped() >= 0 diff --git a/tests/pyfglair/test_provision.py b/tests/pyfglair/test_provision.py index c0470c6..2a27186 100644 --- a/tests/pyfglair/test_provision.py +++ b/tests/pyfglair/test_provision.py @@ -47,3 +47,15 @@ async def test_discover_no_key(cloud): async def test_discover_unknown_region(): with pytest.raises(ProvisionError, match="регион"): await discover("user@example.com", "secret", "xx") + + +async def test_discover_malformed_devices(cloud): + cloud["bad_devices"] = True + with pytest.raises(ProvisionError, match="формат устройства"): + await discover("user@example.com", "secret", "eu", base_url=cloud["url"]) + + +async def test_discover_malformed_lanip(cloud): + cloud["bad_lanip"] = True + with pytest.raises(ProvisionError, match="lanip"): + await discover("user@example.com", "secret", "eu", base_url=cloud["url"]) diff --git a/tests/pyfglair/test_session.py b/tests/pyfglair/test_session.py index 7e2ebfc..3c1a765 100644 --- a/tests/pyfglair/test_session.py +++ b/tests/pyfglair/test_session.py @@ -74,6 +74,114 @@ async def test_session_against_mock(mock_ac, ac_env): assert proc.wait_line("DELETE", 10) is not None +async def test_close_idempotent(mock_ac, ac_env): + proc = mock_ac() + cfg = Config( + host="127.0.0.1", + device_port=proc.port, + dsn=ac_env["dsn"], + lanip_key=ac_env["key"], + lanip_key_id=ac_env["key_id"], + ) + session = Session(cfg) + assert session.start() + await session.async_stop() + await session.async_stop() + assert session.close() is None + assert session.closed + assert session.set_int(Prop.FAN_SPEED, 1) is False + assert session.cached(Prop.FAN_SPEED) is None + + +async def test_close_offline_module(ac_env): + import socket + + sock = socket.socket() + sock.bind(("127.0.0.1", 0)) + dead_port = sock.getsockname()[1] + sock.close() + + cfg = Config( + host="127.0.0.1", + device_port=dead_port, + dsn=ac_env["dsn"], + lanip_key=ac_env["key"], + lanip_key_id=ac_env["key_id"], + keepalive_ms=500, + ) + session = Session(cfg) + session.start() + await asyncio.sleep(0.5) + await session.async_stop() + assert session.closed + assert session.events_dropped == 0 + + +async def test_close_from_callback(mock_ac, ac_env): + """Без loop колбэк идёт в reader-потоке; close() оттуда не дедлочится.""" + import threading + + proc = mock_ac() + holder: list[Session] = [] + called = threading.Event() + + def on_state(state, error): + if state == State.REGISTERING and holder: + holder[0].close() + called.set() + + cfg = Config( + host="127.0.0.1", + device_port=proc.port, + dsn=ac_env["dsn"], + lanip_key=ac_env["key"], + lanip_key_id=ac_env["key_id"], + ) + + def build(): + session = Session(cfg, loop=None, on_state=on_state) + holder.append(session) + session.start() + return session + + session = await asyncio.to_thread(build) + assert await asyncio.to_thread(called.wait, 10), "колбэк не пришёл" + deadline = time.monotonic() + 10 + while not session.closed and time.monotonic() < deadline: + await asyncio.sleep(0.05) + assert session.closed + assert session.close() is None + + +async def test_log_delivery(mock_ac, ac_env): + from pyfglair import set_log_handler, set_log_level + + logs: list[tuple[int, str]] = [] + set_log_level(0) + set_log_handler(lambda level, msg: logs.append((level, msg))) + try: + proc = mock_ac() + cfg = Config( + host="127.0.0.1", + device_port=proc.port, + dsn=ac_env["dsn"], + lanip_key=ac_env["key"], + lanip_key_id=ac_env["key_id"], + keepalive_ms=1000, + ) + session = Session(cfg) + session.start() + deadline = time.monotonic() + 10 + while not logs and time.monotonic() < deadline: + await asyncio.sleep(0.1) + await session.async_stop() + finally: + set_log_handler(None) + set_log_level(1) + assert logs, "логи ядра не доставлены" + assert all(isinstance(level, int) and message for level, message in logs) + + async def test_callbacks_run_in_loop_thread(mock_ac, ac_env): proc = mock_ac() loop = asyncio.get_running_loop()