"""设备连接管理器(方案 C:按需连接 + 空闲保持,全操作串行)。 - 一把 asyncio.Lock 串行所有设备操作(BLE 适配器独占) - 首次操作守候连接(connect_timeout 应覆盖设备 30-60s 广播窗口) - 操作后保持连接;空闲 idle_timeout 秒自动断开让设备睡眠 - 操作中连接丢失(TRANSPORT_ERRORS 全家)→ 重连并重试一次,再失败向上抛 """ from __future__ import annotations import asyncio import time from bleak.exc import BleakError from .client import PPClient from .transports.base import TransportError from .transports.local import TransportError as LocalTransportError # 传输层异常全家(掉线/未连上的全部真实形态): # - transports.base.TransportError:SDK 公开语义异常(bridge/fake 路径) # - transports.local.TransportError:真实 BLE 路径 connect/scan/ID 应答的中文异常 # (SDK 两个同名类互不继承,历史包袱;MCP 层必须同时覆盖) # - bleak.exc.BleakError:写/读路径 bleak 裸抛不包装(如 "Not connected") # - OSError / asyncio.TimeoutError:bleak 后端平台差异形态 TRANSPORT_ERRORS = (TransportError, LocalTransportError, BleakError, OSError, asyncio.TimeoutError) class DeviceManager: def __init__(self, mac: str | None = None, connect_timeout: float = 90.0, idle_timeout: float = 300.0, client_factory=None): self._mac = mac self._connect_timeout = connect_timeout self._idle_timeout = idle_timeout self._factory = client_factory or ( lambda mac, timeout: PPClient.via_local(mac=mac, timeout=timeout)) self._client = None self._lock = asyncio.Lock() self._last_activity: float | None = None self._idle_task: asyncio.Task | None = None self._connect_count = 0 @property def connected(self) -> bool: return self._client is not None async def run(self, op): """串行执行 op(client);掉线重连重试一次。""" async with self._lock: try: return await self._run_with_retry(op) finally: self._touch() async def status(self) -> dict: async with self._lock: await self._ensure_connected() # 先算 idle_seconds 再 _touch:status 自身也是活动,但不能把读数归零 idle = (None if self._last_activity is None else round(time.monotonic() - self._last_activity, 1)) self._touch() info = { "connected": self.connected, "mac": self._mac, "idle_seconds": idle, "connect_count": self._connect_count, "connect_timeout": self._connect_timeout, "idle_timeout": self._idle_timeout, } try: info["device_id"] = await self._client.get_device_id() except Exception: # noqa: BLE001 - 状态查询尽力而为 pass return info async def close(self) -> None: if self._idle_task is not None: self._idle_task.cancel() async with self._lock: await self._drop() # ---------- 内部 ---------- async def _run_with_retry(self, op): await self._ensure_connected() try: return await op(self._client) except TRANSPORT_ERRORS: await self._drop() await self._ensure_connected() return await op(self._client) async def _ensure_connected(self): if self._client is not None: return client = self._factory(self._mac, self._connect_timeout) try: await client.__aenter__() except Exception: # __aenter__ 半途失败(如 connect 后 start_notify 炸)会留下半开连接, # 补一次 __aexit__ 清理;清理自身失败不影响原异常上抛 try: await client.__aexit__(None, None, None) except Exception: # noqa: BLE001 pass raise self._client = client self._connect_count += 1 if self._mac is None: # 自动扫描后记住实际地址 self._mac = getattr(client._t, "address", None) async def _drop(self): client, self._client = self._client, None if client is not None: try: await client.__aexit__(None, None, None) except Exception: # noqa: BLE001 - 断开失败不影响状态 pass def _touch(self): self._last_activity = time.monotonic() if self._idle_task is not None: self._idle_task.cancel() self._idle_task = None if self.connected: self._idle_task = asyncio.create_task(self._idle_watch()) async def _idle_watch(self): try: await asyncio.sleep(self._idle_timeout) async with self._lock: if (self._last_activity is not None and time.monotonic() - self._last_activity >= self._idle_timeout): await self._drop() except asyncio.CancelledError: pass