Files
qianbian/tools/bridge_server.py
T

333 lines
15 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/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 = 11 # 2026-07-29: ota 支持 gpio_map 参数与 max_blocks 探针模式(GPIO 引脚实证用)
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([0x05, 0x06, 0x00, 0x40])
# GPIO_MAP 实证(2026-07-29 gpio_probe):BE 序 = APK 整数 0x05060040,唯一电气有效
# ({00,00,06,05}/{40,00,06,05} 均致 0x38000 读失败=状态 22;C 读成功=状态 20 系产品头缺失,非引脚问题)
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,
gpio_map: bytes = GPIO_MAP_DEFAULT, max_blocks: int = 0):
"""max_blocks>0 时只烧前 N 块即停(GPIO 探针用,不写 END)。"""
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)
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):
if max_blocks and idx > max_blocks:
yield {"event": "ota_stage", "stage": f"probe_stop_at_{max_blocks}"}
return
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"]),
gpio_map=bytes.fromhex(req["gpio_map"]) if req.get("gpio_map") else GPIO_MAP_DEFAULT,
max_blocks=int(req.get("max_blocks", 0)))
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)