From 49b2010088b4aacb12c63af0364bc4d6597e41e8 Mon Sep 17 00:00:00 2001 From: agent Date: Tue, 28 Jul 2026 14:23:30 +0000 Subject: [PATCH] =?UTF-8?q?feat(bridge):=20BLE=20=E6=A1=A5=E6=8E=A5?= =?UTF-8?q?=E8=83=BD=E5=8A=9B=E2=80=94=E2=80=94bridge=5Fserver.py(Windows?= =?UTF-8?q?=E5=8E=9F=E7=94=9F)=20+=20BridgeTransport=20+=20--via/PPCLOCK?= =?UTF-8?q?=5FVIA=EF=BC=8C149=20=E6=B5=8B=E8=AF=95=E5=85=A8=E7=BB=BF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- PROGRESS.md | 7 ++ README.md | 36 ++++++- src/ppclock/bridge.py | 160 ++++++++++++++++++++++++++++ src/ppclock/cli.py | 18 +++- tests/test_bridge.py | 151 ++++++++++++++++++++++++++ tests/test_cli.py | 52 +++++++++ tools/bridge_server.py | 236 +++++++++++++++++++++++++++++++++++++++++ 7 files changed, 658 insertions(+), 2 deletions(-) create mode 100644 src/ppclock/bridge.py create mode 100644 tests/test_bridge.py create mode 100644 tools/bridge_server.py diff --git a/PROGRESS.md b/PROGRESS.md index 86c3fd3..dbf3d9c 100644 --- a/PROGRESS.md +++ b/PROGRESS.md @@ -91,6 +91,13 @@ - 调用方式:`ppclock --via bridge://192.168.61.35:8971 ...` 或 `PPCLOCK_VIA=bridge://...` - field_test.py 经 `PPCLOCK_VIA` 环境变量继承,无需改动 +### 2026-07-28:bridge 实现完成(149 测试全绿) + +- `tools/bridge_server.py`:单文件服务端(仅依赖 bleak),JSONL/TCP RPC v1;本机冒烟通过(ping 返回 platform/version,scan 经服务端真实 bleak 发现 NRF-5DBF28) +- `src/ppclock/bridge.py`:BridgeTransport 与 BLETransport 同接口;`--via bridge://host[:port]` + `PPCLOCK_VIA` 环境变量;scan 同样走桥 +- tests/test_bridge.py(12 测试:假服务端全协议)+ test_cli.py 新增 5 个 --via 测试 +- **下一步**:用户在 host-win 执行 `pip install bleak && python tools/bridge_server.py`,然后本机 `PPCLOCK_VIA=bridge://192.168.61.35 .venv/bin/python tools/field_test.py` 即可完成全链路真机测试 + ### 会话时间线 - 08:14 目标设定(治理+规划+执行) diff --git a/README.md b/README.md index 5eba709..3c57b00 100644 --- a/README.md +++ b/README.md @@ -86,10 +86,44 @@ CLI 全部为单次执行语义,调度交给系统 cron: 结果自动写入 `docs/field-test-result.md`。 +## BLE bridge(蓝牙在他机时) + +当本机无蓝牙(或蓝牙经 USB/IP 透传不稳定)时,把 BLE 操作桥接到蓝牙物理所在机: + +**服务端**(蓝牙所在机,如 Windows host-win;单文件,仅依赖 bleak): + +```bash +pip install bleak +python tools/bridge_server.py --host 0.0.0.0 --port 8971 +``` + +**客户端**(本机,全部 CLI 命令透明可用): + +```bash +ppclock --via bridge://192.168.61.35:8971 --json scan +export PPCLOCK_VIA=bridge://192.168.61.35:8971 # 或环境变量一次设定 +ppclock --json image photo.jpg --slot 0 +PPCLOCK_VIA=bridge://192.168.61.35 .venv/bin/python tools/field_test.py +``` + +协议:JSONL over TCP(`ping/scan/connect/write/read/request_device_id/ota/close`), +OTA 进度以事件流回。仅限可信 LAN 使用(无鉴权)。 + +## 本机蓝牙经 USB/IP 透传的已知问题 + +大 MTU GATT 流量会触发断流(详见 DECISIONS.md 2026-07-28)。保留缓解: +`/sys/bus/usb/devices/1-1/power/control=on`(禁用 runtime PM)。重启后失效,固化: + +```bash +echo 'ACTION=="add", SUBSYSTEM=="usb", ATTR{idVendor}=="8087", ATTR{idProduct}=="0029", ATTR{power/control}="on"' | sudo tee /etc/udev/rules.d/99-ax200-pm.rules +``` + +根治请用上面的 bridge 模式。 + ## 开发 ```bash -.venv/bin/python -m pytest tests/ -q # 132 测试 +.venv/bin/python -m pytest tests/ -q # 149 测试 ``` 治理规则见 `GOVERNANCE.md`;计划见 `PLAN.md`;协议规范见 `docs/protocol.md`。 diff --git a/src/ppclock/bridge.py b/src/ppclock/bridge.py new file mode 100644 index 0000000..292d7b8 --- /dev/null +++ b/src/ppclock/bridge.py @@ -0,0 +1,160 @@ +"""BridgeTransport —— 经 JSONL/TCP 调用远端 bridge 服务端(tools/bridge_server.py)。 +与 BLETransport 同接口,CLI 全命令透明可用。协议见服务端 docstring。 +""" +from __future__ import annotations + +import asyncio +import json + +DEFAULT_PORT = 8971 + + +class BridgeError(Exception): + pass + + +def parse_via(via: str): + """'bridge://host[:port]' → (host, port);'local'/None → None。""" + if via in (None, "", "local"): + return None + if not via.startswith("bridge://"): + raise ValueError(f"--via 仅支持 bridge://host[:port],得到 {via!r}") + rest = via[len("bridge://"):] + host, _, port_s = rest.partition(":") + if not host: + raise ValueError(f"bridge URI 缺 host: {via!r}") + return host, int(port_s) if port_s else DEFAULT_PORT + + +async def scan(host: str, port: int = DEFAULT_PORT, timeout: float = 5.0): + """一次性 RPC:扫描设备。""" + reader, writer = await asyncio.open_connection(host, port) + try: + writer.write(json.dumps({"id": 1, "op": "scan", "timeout": timeout}).encode() + b"\n") + await writer.drain() + resp = json.loads(await reader.readline()) + if not resp.get("ok"): + raise BridgeError(resp.get("error", "scan failed")) + return resp["data"]["devices"] + finally: + writer.close() + + +class BridgeTransport: + """与 BLETransport 同接口:__aenter__/close/write_epd/write_rxtx/read_rxtx/ + request_device_id/run_ota/delay。""" + + def __init__(self, address: str | None = None, timeout: float = 10.0, + *, host: str, port: int = DEFAULT_PORT, + name_prefix: str = "NRF-"): + self.address = address + self.timeout = timeout + self.host = host + self.port = port + self.name_prefix = name_prefix + self._reader: asyncio.StreamReader | None = None + self._writer: asyncio.StreamWriter | None = None + self._reader_task: asyncio.Task | None = None + self._pending: dict[int, asyncio.Future] = {} + self._next_id = 0 + self._on_progress = None + + async def __aenter__(self): + await self._open() + return self + + async def __aexit__(self, *exc): + await self.close() + + async def _open(self): + self._reader, self._writer = await asyncio.open_connection(self.host, self.port) + self._reader_task = asyncio.create_task(self._read_loop()) + try: + if not self.address: + devices = (await self._rpc("scan", timeout=self.timeout))["devices"] + devices = [d for d in devices if d["name"].startswith(self.name_prefix)] + if not devices: + raise BridgeError(f"bridge 未发现 {self.name_prefix!r} 前缀设备") + self.address = devices[0]["address"] + await self._rpc("connect", address=self.address, timeout=self.timeout) + except Exception: + await self._teardown() + raise + + async def close(self): + try: + if self._writer: + await self._rpc("close") + except Exception: + pass + await self._teardown() + + async def _teardown(self): + if self._reader_task: + self._reader_task.cancel() + try: + await self._reader_task + except (asyncio.CancelledError, Exception): + pass + self._reader_task = None + if self._writer: + self._writer.close() + try: + await self._writer.wait_closed() + except Exception: + pass + self._writer = None + + async def _read_loop(self): + try: + async for line in self._reader: + msg = json.loads(line) + if "event" in msg: + if self._on_progress: + self._on_progress(msg) + else: + fut = self._pending.pop(msg.get("id"), None) + if fut and not fut.done(): + fut.set_result(msg) + except Exception: + for fut in self._pending.values(): + if not fut.done(): + fut.set_exception(BridgeError("bridge 连接中断")) + self._pending.clear() + + async def _rpc(self, op: str, timeout: float = 30.0, **kw) -> dict: + self._next_id += 1 + rid = self._next_id + fut = asyncio.get_running_loop().create_future() + self._pending[rid] = fut + self._writer.write(json.dumps({"id": rid, "op": op, **kw}).encode() + b"\n") + await self._writer.drain() + try: + resp = await asyncio.wait_for(fut, timeout) + finally: + self._pending.pop(rid, None) + if not resp.get("ok"): + raise BridgeError(resp.get("error", op)) + return resp.get("data") or {} + + async def write_epd(self, data: bytes, response: bool = True): + await self._rpc("write", char="epd", data=data.hex(), response=response) + + async def write_rxtx(self, data: bytes, response: bool = True): + await self._rpc("write", char="rxtx", data=data.hex(), response=response) + + async def read_rxtx(self) -> bytes: + return bytes.fromhex((await self._rpc("read", char="rxtx"))["data"]) + + async def request_device_id(self, timeout: float = 8.0) -> str: + return (await self._rpc("request_device_id", timeout=timeout + 5))["device_id"] + + async def run_ota(self, image: bytes, on_progress=None): + self._on_progress = on_progress + try: + return await self._rpc("ota", timeout=600.0, image=image.hex()) + finally: + self._on_progress = None + + async def delay(self, seconds: float): + await asyncio.sleep(seconds) diff --git a/src/ppclock/cli.py b/src/ppclock/cli.py index 548b0c7..1f8450e 100644 --- a/src/ppclock/cli.py +++ b/src/ppclock/cli.py @@ -1,10 +1,13 @@ -"""ppclock CLI — 离线、agent 友好(--json 机器可读输出)。""" +"""ppclock CLI — 离线、agent 友好(--json 机器可读输出)。 +传输选择:默认本机 BLE;`--via bridge://host[:port]`(或环境变量 PPCLOCK_VIA) +经 bridge 服务端(tools/bridge_server.py)远程使用他机蓝牙。""" from __future__ import annotations import argparse import asyncio import datetime import json +import os import sys from . import commands as C @@ -39,6 +42,12 @@ def _emit(args, cmd, data=None, error=None): def _make_transport(args): """工厂函数——测试用 monkeypatch 替换为 FakeTransport。""" + from .bridge import BridgeTransport, parse_via + via = parse_via(getattr(args, "via", None)) + if via: + host, port = via + return BridgeTransport(address=args.mac, timeout=args.timeout, + host=host, port=port) from .transport import BLETransport return BLETransport(address=args.mac, timeout=args.timeout) @@ -86,6 +95,8 @@ def build_parser() -> argparse.ArgumentParser: p.add_argument("--mac", help="设备 MAC(省略则自动扫描 NRF- 前缀设备)") p.add_argument("--timeout", type=float, default=10.0) p.add_argument("--json", action="store_true", help="机器可读 JSON 输出") + p.add_argument("--via", default=os.environ.get("PPCLOCK_VIA"), + help="传输:local(默认) 或 bridge://host[:port](桥接他机蓝牙)") sub = p.add_subparsers(dest="cmd", required=True) sub.add_parser("scan", help="扫描设备") @@ -299,6 +310,11 @@ async def _run_batch(args) -> list: async def _dispatch(args) -> dict | list | None: if args.cmd == "scan": + from .bridge import parse_via + via = parse_via(getattr(args, "via", None)) + if via: + from .bridge import scan as bridge_scan + return {"devices": await bridge_scan(via[0], via[1], timeout=args.timeout)} from .transport import scan return {"devices": await scan(timeout=args.timeout)} if args.cmd == "batch": diff --git a/tests/test_bridge.py b/tests/test_bridge.py new file mode 100644 index 0000000..db8bf9e --- /dev/null +++ b/tests/test_bridge.py @@ -0,0 +1,151 @@ +"""BridgeTransport 测试 — 用内存假 bridge 服务端验证 JSONL 协议与同构接口。""" +import asyncio +import json + +import pytest +import pytest_asyncio + +from ppclock import bridge + + +class FakeBridgeServer: + """模拟 tools/bridge_server.py 的 JSONL 协议。""" + + def __init__(self): + self.received = [] + self.device_id = "89794900980101" + + async def start(self): + self.server = await asyncio.start_server(self._handle, "127.0.0.1", 0) + return self.server.sockets[0].getsockname()[1] + + async def _handle(self, reader, writer): + async for line in reader: + req = json.loads(line) + self.received.append(req) + op, rid = req["op"], req["id"] + if op == "ping": + data = {"pong": True, "version": 1, "platform": "FakeOS"} + elif op == "scan": + data = {"devices": [{"name": "NRF-FAKE", "address": "AA:BB:CC:DD:EE:FF", + "rssi": -50}]} + elif op == "connect": + data = {"address": req["address"]} + elif op == "write": + data = {"sent": len(req["data"]) // 2} + elif op == "read": + data = {"data": "beef"} + elif op == "request_device_id": + data = {"device_id": self.device_id} + elif op == "ota": + for i in (1, 2): + writer.write((json.dumps({"event": "ota_progress", "block": i, + "total": 2}) + "\n").encode()) + await writer.drain() + data = {"done": True} + elif op == "close": + data = {} + elif op == "fail_me": + writer.write((json.dumps({"id": rid, "ok": False, + "error": "RuntimeError: boom"}) + "\n").encode()) + await writer.drain() + continue + else: + continue + writer.write((json.dumps({"id": rid, "ok": True, "data": data}) + "\n").encode()) + await writer.drain() + + +@pytest_asyncio.fixture +async def server(): + s = FakeBridgeServer() + port = await s.start() + yield s, port + s.server.close() + + +@pytest.mark.asyncio +async def test_scan_via_rpc(server): + s, port = server + devices = await bridge.scan("127.0.0.1", port) + assert devices[0]["name"] == "NRF-FAKE" + + +@pytest.mark.asyncio +async def test_connect_autoscan(server): + """无 address 时自动扫描选首台 NRF-(与 BLETransport 语义一致)""" + s, port = server + async with bridge.BridgeTransport(host="127.0.0.1", port=port) as t: + assert t.address == "AA:BB:CC:DD:EE:FF" + ops = [r["op"] for r in s.received] + assert ops[:2] == ["scan", "connect"] + + +@pytest.mark.asyncio +async def test_connect_explicit_address(server): + s, port = server + async with bridge.BridgeTransport(address="18:BC:5A:5D:BF:28", + host="127.0.0.1", port=port) as t: + assert t.address == "18:BC:5A:5D:BF:28" + conn = next(r for r in s.received if r["op"] == "connect") + assert conn["address"] == "18:BC:5A:5D:BF:28" + + +@pytest.mark.asyncio +async def test_write_channels_and_hex(server): + s, port = server + async with bridge.BridgeTransport(address="X", host="127.0.0.1", port=port) as t: + await t.write_epd(bytes([0x03, 0xFF])) + await t.write_rxtx(bytes([0xE2]), response=False) + w = [r for r in s.received if r["op"] == "write"] + assert w[0]["char"] == "epd" and w[0]["data"] == "03ff" + assert w[1]["char"] == "rxtx" and w[1]["data"] == "e2" and w[1]["response"] is False + + +@pytest.mark.asyncio +async def test_read_rxtx(server): + s, port = server + async with bridge.BridgeTransport(address="X", host="127.0.0.1", port=port) as t: + assert await t.read_rxtx() == bytes([0xBE, 0xEF]) + + +@pytest.mark.asyncio +async def test_request_device_id(server): + s, port = server + async with bridge.BridgeTransport(address="X", host="127.0.0.1", port=port) as t: + assert await t.request_device_id() == "89794900980101" + + +@pytest.mark.asyncio +async def test_run_ota_streams_progress(server): + s, port = server + async with bridge.BridgeTransport(address="X", host="127.0.0.1", port=port) as t: + events = [] + result = await t.run_ota(b"\x70\x51" + b"\x00" * 100, on_progress=events.append) + assert result == {"done": True} + assert [e["block"] for e in events] == [1, 2] + + +@pytest.mark.asyncio +async def test_error_raises_bridge_error(server): + s, port = server + async with bridge.BridgeTransport(address="X", host="127.0.0.1", port=port) as t: + with pytest.raises(bridge.BridgeError, match="boom"): + await t._rpc("fail_me") + + +class TestParseVia: + def test_bridge_uri(self): + host, port = bridge.parse_via("bridge://192.168.61.35:8971") + assert (host, port) == ("192.168.61.35", 8971) + + def test_default_port(self): + host, port = bridge.parse_via("bridge://192.168.61.35") + assert (host, port) == ("192.168.61.35", 8971) + + def test_bad_scheme(self): + with pytest.raises(ValueError): + bridge.parse_via("tcp://x") + + def test_local_is_none(self): + assert bridge.parse_via("local") is None diff --git a/tests/test_cli.py b/tests/test_cli.py index 411d9c6..58badde 100644 --- a/tests/test_cli.py +++ b/tests/test_cli.py @@ -231,3 +231,55 @@ class TestActivateAuto: expected = keygen_code(bytes([0x9A, 0x01, 0x99, 0x50, 0x10, 0x20])) assert t.rxtx_writes[1] == bytes([0xEF]) + expected assert out["data"]["code"] == expected.hex() + + +_ORIG_MAKE = cli._make_transport + + +class TestViaBridge: + def test_via_bridge_creates_bridge_transport(self): + import argparse + from ppclock.bridge import BridgeTransport + args = argparse.Namespace(mac="AA:BB", timeout=5.0, via="bridge://192.168.61.35:8971") + t = _ORIG_MAKE(args) + assert isinstance(t, BridgeTransport) + assert (t.host, t.port) == ("192.168.61.35", 8971) + assert t.address == "AA:BB" + + def test_via_default_port(self): + import argparse + args = argparse.Namespace(mac=None, timeout=5.0, via="bridge://10.0.0.2") + t = _ORIG_MAKE(args) + assert t.port == 8971 + + def test_via_local_creates_ble(self): + import argparse + from ppclock.transport import BLETransport + args = argparse.Namespace(mac=None, timeout=5.0, via=None) + assert isinstance(_ORIG_MAKE(args), BLETransport) + + def test_scan_uses_bridge(self, capsys, monkeypatch): + import ppclock.bridge as br + called = {} + + async def fake_scan(host, port, timeout=5.0): + called["args"] = (host, port) + return [{"name": "NRF-BR", "address": "CC:DD", "rssi": -60}] + monkeypatch.setattr(br, "scan", fake_scan) + code, out = run_cli(capsys, ["--json", "--via", "bridge://10.0.0.9:9999", "scan"]) + assert code == 0 + assert called["args"] == ("10.0.0.9", 9999) + assert out["data"]["devices"][0]["name"] == "NRF-BR" + + def test_env_var_fallback(self, capsys, monkeypatch): + import ppclock.bridge as br + called = {} + + async def fake_scan(host, port, timeout=5.0): + called["args"] = (host, port) + return [] + monkeypatch.setattr(br, "scan", fake_scan) + monkeypatch.setenv("PPCLOCK_VIA", "bridge://10.0.0.7") + code, _ = run_cli(capsys, ["--json", "scan"]) + assert code == 0 + assert called["args"] == ("10.0.0.7", 8971) diff --git a/tools/bridge_server.py b/tools/bridge_server.py new file mode 100644 index 0000000..bae15ee --- /dev/null +++ b/tools/bridge_server.py @@ -0,0 +1,236 @@ +#!/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 + +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" +DEVICE_NAME_PREFIX = "NRF-" + +SPOTA_SERVICE_UUID = "0000fef5-0000-1000-8000-00805f9b34fb" +SPOTA_SERV_STATUS_UUID = "64b4e8b5-0de5-401b-a21d-acc8db3b913a" +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} + + +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._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): + from bleak import BleakClient + await self.close() + self.client = BleakClient(address, timeout=timeout) + await self.client.connect() + await self.client.start_notify(RXTX_CHAR_UUID, self._on_notify) + return {"address": address} + + 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() + await self.client.write_gatt_char(RXTX_CHAR_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")) + + async def wait_ok(what: str): + status = await asyncio.wait_for(status_q.get(), status_timeout) + if status != STATUS_OK: + raise RuntimeError(f"{what} 失败: 状态码 {status}") + + 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)] + + await self.client.start_notify(SPOTA_SERV_STATUS_UUID, _on_status) + events = [] + try: + await self.client.write_gatt_char(SPOTA_MEM_DEV_UUID, MEM_DEV_SPI_FLASH, response=True) + await self.client.write_gatt_char(SPOTA_GPIO_MAP_UUID, GPIO_MAP_DEFAULT, response=True) + for idx, block in enumerate(blocks, 1): + await self.client.write_gatt_char(SPOTA_PATCH_LEN_UUID, + len(block).to_bytes(2, "little"), response=True) + 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) + await wait_ok(f"块 {idx}/{len(blocks)}") + events.append({"block": idx, "total": len(blocks)}) + yield {"event": "ota_progress", "block": idx, "total": len(blocks)} + await self.client.write_gatt_char(SPOTA_MEM_DEV_UUID, END_SIGNAL, response=True) + await wait_ok("镜像校验") + await self.client.write_gatt_char(SPOTA_MEM_DEV_UUID, REBOOT_SIGNAL, response=True) + finally: + try: + await self.client.stop_notify(SPOTA_SERV_STATUS_UUID) + 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, + "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))) + 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 {} + 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() + srv = await asyncio.start_server(server.handle, args.host, args.port) + 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)