#!/usr/bin/env python3 """ppclock BLE bridge 服务端 —— 在蓝牙物理所在机(host-win)原生运行。 依赖:pip install bleak(仅此一个,Windows 用 WinRT 后端)。 运行:python bridge_server.py [--host 0.0.0.0] [--port 8971] 协议:JSONL over TCP。每行一个请求 {"id":N,"op":...}, 响应 {"id":N,"ok":bool,"data":...,"error":...};OTA 期间穿插 {"event":"ota_progress",...} 事件。单 BLE 设备会话按连接串行化。 RPC: ping → {pong, version, platform} scan {timeout} → {devices:[{name,address,rssi}]} connect {address} → {}(连接+服务发现+订阅 RXTX notify) write {char,data,response}→ {}(char: "epd"|"rxtx",data: hex) read {char} → {data: hex} request_device_id {timeout}→ {device_id}(EFEF + 收集 14B notify) ota {image} → {blocks, bytes}(image: hex;完整 SUOTA) close → {} """ from __future__ import annotations import argparse import asyncio import json import platform import sys PROTOCOL_VERSION = 1 SERVER_BUILD = 10 # 2026-07-29: 状态码 8/16 按会话信息码处理(固件反汇编实证)+ 阶段状态取证 + SPOTA 写特征入 CHARS EPD_SERVICE_UUID = "13187b10-eba9-a3ba-044e-83d3217d9a38" EPD_CHAR_UUID = "4b646063-6264-f3a7-8941-e65356ea82fe" RXTX_SERVICE_UUID = "00001f10-0000-1000-8000-00805f9b34fb" RXTX_CHAR_UUID = "00001f1f-0000-1000-8000-00805f9b34fb" # 固件分析(firmware-analysis §2/§3):0x221f 服务的 0x331f 特征(handle 0x0B) # 与 0x1f1f(handle 7)同入命令分发器;部分固件 notify 能力挂在 331f 上 ALT_NOTIFY_UUID = "0000331f-0000-1000-8000-00805f9b34fb" DEVICE_NAME_PREFIX = "NRF-" SPOTA_SERVICE_UUID = "0000fef5-0000-1000-8000-00805f9b34fb" SPOTA_SERV_STATUS_UUID = "5f78df94-798c-46f5-990a-b3eb6a065c88" # 实机 h=38 notify(64b4e8b5 仅 read) SPOTA_MEM_DEV_UUID = "8082caa8-41a6-4021-91c6-56f9b954cc34" SPOTA_GPIO_MAP_UUID = "724249f0-5ec3-4b5f-8804-42345af08651" SPOTA_PATCH_DATA_UUID = "457871e8-d516-4ca1-9116-57d0b17b9cb2" SPOTA_PATCH_LEN_UUID = "9d84b9a3-000c-49d8-9183-855b673fda31" MEM_DEV_SPI_FLASH = bytes([0x00, 0x00, 0x00, 0x13]) GPIO_MAP_DEFAULT = bytes([0x00, 0x00, 0x06, 0x05]) END_SIGNAL = bytes([0x00, 0x00, 0x00, 0xFE]) REBOOT_SIGNAL = bytes([0x00, 0x00, 0x00, 0xFD]) BLOCK_SIZE = 240 WRITE_CHUNK = 20 STATUS_OK = 2 CHARS = {"epd": EPD_CHAR_UUID, "rxtx": RXTX_CHAR_UUID, "alt": ALT_NOTIFY_UUID, "dis_model": "00002a24-0000-1000-8000-00805f9b34fb", "dis_fw": "00002a26-0000-1000-8000-00805f9b34fb", "spota_mem_info": "6c53db25-47a1-45fe-a022-7c92fb334fd4", "spota_status_ro": "64b4e8b5-0de5-401b-a21d-acc8db3b913a", "spota_mem_dev": "8082caa8-41a6-4021-91c6-56f9b954cc34", "spota_gpio": "724249f0-5ec3-4b5f-8804-42345af08651", "spota_patch_len": "9d84b9a3-000c-49d8-9183-855b673fda31", "spota_patch_data": "457871e8-d516-4ca1-9116-57d0b17b9cb2"} # 实机 GATT(2026-07-28 bridge 诊断):1f1f(h=53)=仅write 无 notify; # 331f(h=57)=notify+write+read,为真交互通道;EPD(h=48)=notify+write def parse_device_id(data: bytes) -> str: if all(0x20 <= b < 0x7F for b in data): return data.decode("ascii") return data.hex().upper() class BleSession: """单设备 BLE 会话(bleak),与 ppclock.transport.BLETransport 同语义。""" def __init__(self): self.client = None self.notify_char = None self._id_buf = bytearray() self._id_event = asyncio.Event() async def scan(self, timeout: float): from bleak import BleakScanner found = await BleakScanner.discover(timeout=timeout, return_adv=True) out = [] for addr, (dev, adv) in found.items(): name = dev.name or adv.local_name or "" if name.startswith(DEVICE_NAME_PREFIX): out.append({"name": name, "address": dev.address, "rssi": adv.rssi}) return {"devices": out} async def connect(self, address: str, timeout: float = 15.0, subscribe: bool = True): from bleak import BleakClient, BleakScanner await self.close() # Windows/WinRT 按地址直连常解析失败——先扫描拿到 BLEDevice(含地址类型/广播数据)再连 deadline = asyncio.get_event_loop().time() + timeout dev = None while True: remain = deadline - asyncio.get_event_loop().time() if remain <= 0: break dev = await BleakScanner.find_device_by_address( address, timeout=min(5.0, max(1.0, remain))) if dev is not None: break if dev is None: raise RuntimeError(f"设备解析失败(扫描未见 {address})") self.client = BleakClient(dev, timeout=timeout) await self.client.connect() # notify 全订阅:331f(主)+ 1f1f(部分固件)+ EPD(设备 ID 也可能经此应答) # subscribe=False 用于 OTA 净连接(避免订阅抖动影响烧录链路) self.notify_char = [] if subscribe: for uuid in (ALT_NOTIFY_UUID, RXTX_CHAR_UUID, EPD_CHAR_UUID): try: await asyncio.wait_for( self.client.start_notify(uuid, self._on_notify), 10) self.notify_char.append(uuid) except Exception: continue table = self._gatt_table() return {"address": address, "name": dev.name, "notify_char": self.notify_char, "gatt": table} def _gatt_table(self): table = [] for s in self.client.services: for c in s.characteristics: table.append({"service": str(s.uuid), "char": str(c.uuid), "handle": c.handle, "props": list(c.properties)}) return table def _on_notify(self, _sender, data: bytearray): self._id_buf.extend(data) if len(self._id_buf) >= 14: self._id_event.set() async def write(self, char: str, data: bytes, response: bool): await self.client.write_gatt_char(CHARS[char], data, response=response) return {"sent": len(data)} async def read(self, char: str): return {"data": bytes(await self.client.read_gatt_char(CHARS[char])).hex()} async def request_device_id(self, timeout: float = 8.0): self._id_buf.clear() self._id_event.clear() # EFEF 走 331f(实机证据:1f1f 仅 ATT-ACK 不执行命令,331f 为活通道) await self.client.write_gatt_char(ALT_NOTIFY_UUID, bytes([0xEF, 0xEF]), response=True) await asyncio.wait_for(self._id_event.wait(), timeout) return {"device_id": parse_device_id(bytes(self._id_buf[:14]))} async def ota(self, image: bytes, status_timeout: float = 60.0): if SPOTA_SERVICE_UUID not in [str(s.uuid) for s in self.client.services]: raise RuntimeError("设备不存在 fef5 SUOTA 服务") status_q: asyncio.Queue[int] = asyncio.Queue() def _on_status(_s, data: bytearray): status_q.put_nowait(int.from_bytes(bytes(data), "little")) acc = 0 for b in image: acc ^= b payload = image + bytes([acc]) blocks = [payload[i:i + BLOCK_SIZE] for i in range(0, len(payload), BLOCK_SIZE)] async def w(char, data, to=15.0): await asyncio.wait_for( self.client.write_gatt_char(char, data, response=True), to) async def drain(label: str, settle: float = 0.5): """排干状态队列并如实上报(阶段状态码取证)。""" await asyncio.sleep(settle) got = [] while not status_q.empty(): got.append(status_q.get_nowait()) if got: yield {"event": "ota_stage", "stage": f"{label}_status", "codes": got} await asyncio.wait_for( self.client.start_notify(SPOTA_SERV_STATUS_UUID, _on_status), 10) try: await w(SPOTA_MEM_DEV_UUID, MEM_DEV_SPI_FLASH) yield {"event": "ota_stage", "stage": "mem_dev_ok"} async for ev in drain("mem_dev", 1.0): yield ev await w(SPOTA_GPIO_MAP_UUID, GPIO_MAP_DEFAULT) yield {"event": "ota_stage", "stage": "gpio_map_ok"} async for ev in drain("gpio_map", 1.0): yield ev for idx, block in enumerate(blocks, 1): await w(SPOTA_PATCH_LEN_UUID, len(block).to_bytes(2, "little")) for i in range(0, len(block), WRITE_CHUNK): await self.client.write_gatt_char(SPOTA_PATCH_DATA_UUID, block[i:i + WRITE_CHUNK], response=False) yield {"event": "ota_block_sent", "block": idx, "total": len(blocks)} # 块确认:只认 2;8/16 为会话信息码(固件 0x07FCB770/0x07FCB780 证据),记录并跳过 while True: status = await asyncio.wait_for(status_q.get(), status_timeout) if status == STATUS_OK: break if status in (8, 16): yield {"event": "ota_stage", "stage": f"block{idx}_info_status", "codes": [status]} continue raise RuntimeError(f"块 {idx}/{len(blocks)} 失败: 状态码 {status}") yield {"event": "ota_progress", "block": idx, "total": len(blocks)} await w(SPOTA_MEM_DEV_UUID, END_SIGNAL) yield {"event": "ota_stage", "stage": "end_signal_ok"} while True: status = await asyncio.wait_for(status_q.get(), status_timeout) if status == STATUS_OK: break if status in (8, 16): yield {"event": "ota_stage", "stage": "end_info_status", "codes": [status]} continue raise RuntimeError(f"镜像校验 失败: 状态码 {status}") await w(SPOTA_MEM_DEV_UUID, REBOOT_SIGNAL) yield {"event": "ota_stage", "stage": "reboot_sent"} finally: try: await asyncio.wait_for( self.client.stop_notify(SPOTA_SERV_STATUS_UUID), 5) except Exception: pass async def close(self): if self.client is not None: try: if self.client.is_connected: await self.client.disconnect() finally: self.client = None class BridgeServer: def __init__(self): self.ble = BleSession() self._ble_lock = asyncio.Lock() async def handle(self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter): peer = writer.get_extra_info("peername") print(f"[bridge] client {peer} connected", flush=True) async def send(obj): writer.write((json.dumps(obj, ensure_ascii=False) + "\n").encode()) await writer.drain() try: async for line in reader: line = line.strip() if not line: continue try: req = json.loads(line) rid = req.get("id") op = req["op"] except (KeyError, ValueError) as e: await send({"id": None, "ok": False, "error": f"bad request: {e}"}) continue try: async with self._ble_lock: if op == "ota": gen = self.ble.ota(bytes.fromhex(req["image"])) async for ev in gen: await send(ev) data = {"done": True} else: data = await self._dispatch(op, req) await send({"id": rid, "ok": True, "data": data}) except Exception as e: await send({"id": rid, "ok": False, "error": f"{type(e).__name__}: {e}"}) finally: await self.ble.close() writer.close() print(f"[bridge] client {peer} disconnected", flush=True) async def _dispatch(self, op: str, req: dict): if op == "ping": return {"pong": True, "version": PROTOCOL_VERSION, "build": SERVER_BUILD, "platform": platform.system()} if op == "scan": return await self.ble.scan(float(req.get("timeout", 5.0))) if op == "connect": return await self.ble.connect(req["address"], float(req.get("timeout", 15.0)), bool(req.get("subscribe", True))) if op == "write": return await self.ble.write(req["char"], bytes.fromhex(req["data"]), bool(req.get("response", True))) if op == "read": return await self.ble.read(req["char"]) if op == "request_device_id": return await self.ble.request_device_id(float(req.get("timeout", 8.0))) if op == "close": return await self.ble.close() or {} if op == "services": return {"gatt": self.ble._gatt_table(), "notify_char": self.ble.notify_char} raise ValueError(f"unknown op {op!r}") async def main(): ap = argparse.ArgumentParser(description="ppclock BLE bridge server") ap.add_argument("--host", default="0.0.0.0") ap.add_argument("--port", type=int, default=8971) args = ap.parse_args() server = BridgeServer() # limit: OTA 镜像请求 ~145KB(hex),远超默认 64KiB —— 必须显式放大(2026-07-29 根因) srv = await asyncio.start_server(server.handle, args.host, args.port, limit=4 * 1024 * 1024) print(f"[bridge] listening on {args.host}:{args.port} " f"(protocol v{PROTOCOL_VERSION}, platform {platform.system()})", flush=True) async with srv: await srv.serve_forever() if __name__ == "__main__": try: asyncio.run(main()) except KeyboardInterrupt: sys.exit(0)