feat(bridge): BLE 桥接能力——bridge_server.py(Windows原生) + BridgeTransport + --via/PPCLOCK_VIA,149 测试全绿
This commit is contained in:
1 parent
17e2512f6a
commit
49b2010088
7 files changed
+658
-2
No files matched your search
@@ -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)
|
||||
+17
-1
@@ -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":
|
||||
|
||||
Reference in new issue
Block a user