fix(bridge): _rpc 双向 timeout 拆分(payload 服务端预算 vs 客户端等待预算),149 全绿
This commit is contained in:
1 parent
92a8661088
commit
c4f577905f
1 file changed
+9
-6
@@ -76,7 +76,8 @@ class BridgeTransport:
|
|||||||
if not devices:
|
if not devices:
|
||||||
raise BridgeError(f"bridge 未发现 {self.name_prefix!r} 前缀设备")
|
raise BridgeError(f"bridge 未发现 {self.name_prefix!r} 前缀设备")
|
||||||
self.address = devices[0]["address"]
|
self.address = devices[0]["address"]
|
||||||
await self._rpc("connect", address=self.address, timeout=self.timeout)
|
await self._rpc("connect", 90.0 + self.timeout,
|
||||||
|
address=self.address, timeout=self.timeout)
|
||||||
except Exception:
|
except Exception:
|
||||||
await self._teardown()
|
await self._teardown()
|
||||||
raise
|
raise
|
||||||
@@ -122,15 +123,16 @@ class BridgeTransport:
|
|||||||
fut.set_exception(BridgeError("bridge 连接中断"))
|
fut.set_exception(BridgeError("bridge 连接中断"))
|
||||||
self._pending.clear()
|
self._pending.clear()
|
||||||
|
|
||||||
async def _rpc(self, op: str, timeout: float = 30.0, **kw) -> dict:
|
async def _rpc(self, op: str, _timeout: float = 30.0, **kw) -> dict:
|
||||||
self._next_id += 1
|
self._next_id += 1
|
||||||
rid = self._next_id
|
rid = self._next_id
|
||||||
fut = asyncio.get_running_loop().create_future()
|
fut = asyncio.get_running_loop().create_future()
|
||||||
self._pending[rid] = fut
|
self._pending[rid] = fut
|
||||||
self._writer.write(json.dumps({"id": rid, "op": op, **kw}).encode() + b"\n")
|
payload = {"id": rid, "op": op, **kw}
|
||||||
|
self._writer.write(json.dumps(payload).encode() + b"\n")
|
||||||
await self._writer.drain()
|
await self._writer.drain()
|
||||||
try:
|
try:
|
||||||
resp = await asyncio.wait_for(fut, timeout)
|
resp = await asyncio.wait_for(fut, _timeout)
|
||||||
finally:
|
finally:
|
||||||
self._pending.pop(rid, None)
|
self._pending.pop(rid, None)
|
||||||
if not resp.get("ok"):
|
if not resp.get("ok"):
|
||||||
@@ -147,12 +149,13 @@ class BridgeTransport:
|
|||||||
return bytes.fromhex((await self._rpc("read", char="rxtx"))["data"])
|
return bytes.fromhex((await self._rpc("read", char="rxtx"))["data"])
|
||||||
|
|
||||||
async def request_device_id(self, timeout: float = 8.0) -> str:
|
async def request_device_id(self, timeout: float = 8.0) -> str:
|
||||||
return (await self._rpc("request_device_id", timeout=timeout + 5))["device_id"]
|
return (await self._rpc("request_device_id", timeout + 15,
|
||||||
|
timeout=timeout))["device_id"]
|
||||||
|
|
||||||
async def run_ota(self, image: bytes, on_progress=None):
|
async def run_ota(self, image: bytes, on_progress=None):
|
||||||
self._on_progress = on_progress
|
self._on_progress = on_progress
|
||||||
try:
|
try:
|
||||||
return await self._rpc("ota", timeout=600.0, image=image.hex())
|
return await self._rpc("ota", 600.0, image=image.hex())
|
||||||
finally:
|
finally:
|
||||||
self._on_progress = None
|
self._on_progress = None
|
||||||
|
|
||||||
|
|||||||
Reference in new issue
Block a user