feat(ble): add BLE UART daemon RPC API

This commit is contained in:
Zhou Xiao
2026-04-28 13:31:22 +08:00
parent 75f5de65b2
commit 3103b1c442
5 changed files with 330 additions and 0 deletions
@@ -0,0 +1,7 @@
# SPDX-FileCopyrightText: 2026 Espressif Systems (Shanghai) CO LTD
# SPDX-License-Identifier: Apache-2.0
from .api import run_daemon
from .api import run_daemon_send
from .api import run_daemon_status
__all__ = ['run_daemon', 'run_daemon_status', 'run_daemon_send']
@@ -0,0 +1,92 @@
# SPDX-FileCopyrightText: 2026 Espressif Systems (Shanghai) CO LTD
# SPDX-License-Identifier: Apache-2.0
import json
from typing import Any
from urllib.error import HTTPError
from urllib.error import URLError
from urllib.request import Request
from urllib.request import urlopen
import uvicorn
from .server import app as daemon_app
def _daemon_url(host: str, port: int, path: str) -> str:
return f'http://{host}:{port}{path}'
def _request_json(
method: str,
url: str,
payload: dict[str, Any] | None = None,
timeout: float = 10.0,
) -> dict[str, Any]:
data = json.dumps(payload).encode() if payload is not None else None
headers = {'Content-Type': 'application/json'} if payload is not None else {}
request = Request(url, data=data, headers=headers, method=method)
try:
with urlopen(request, timeout=timeout) as response:
body = response.read().decode()
except HTTPError as e:
detail = e.read().decode(errors='replace')
raise RuntimeError(f'Daemon request failed with HTTP {e.code}: {detail}') from e
except TimeoutError as e:
raise RuntimeError(f'Timed out waiting for BLE UART Daemon: {url}') from e
except URLError as e:
raise RuntimeError(f'Failed to connect to BLE UART Daemon: {e.reason}') from e
if not body:
return {}
try:
result = json.loads(body)
except json.JSONDecodeError as e:
raise RuntimeError(f'Invalid daemon response: {body!r}') from e
if not isinstance(result, dict):
raise RuntimeError(f'Invalid daemon response: {result!r}')
return result
def run_daemon(device_id: str, host: str, port: int) -> None:
daemon_app.state.device_id = device_id
uvicorn.run(daemon_app, host=host, port=port)
def run_daemon_status(host: str = '127.0.0.1', port: int = 8888) -> None:
try:
status = _request_json('GET', _daemon_url(host, port, '/status'))
except RuntimeError as e:
print(e)
raise SystemExit(1) from e
print(json.dumps(status, indent=2))
def run_daemon_send(
data: str,
op: str = 'raw',
json_payload: bool = False,
timeout: float = 10.0,
host: str = '127.0.0.1',
port: int = 8888,
) -> None:
try:
payload_data: Any = json.loads(data) if json_payload else data
except json.JSONDecodeError as e:
print(f'Invalid JSON payload: {e}')
raise SystemExit(1) from e
try:
response = _request_json(
'POST',
_daemon_url(host, port, '/request'),
payload={'op': op, 'data': payload_data, 'timeout': timeout},
timeout=timeout + 1.0,
)
except RuntimeError as e:
print(e)
raise SystemExit(1) from e
result = response.get('data', response.get('response', ''))
print(result if isinstance(result, str) else json.dumps(result))
@@ -0,0 +1,103 @@
# SPDX-FileCopyrightText: 2026 Espressif Systems (Shanghai) CO LTD
# SPDX-License-Identifier: Apache-2.0
import asyncio
import json
from json import JSONDecodeError
from typing import Any
from loguru import logger
from ..core.constants import DEFAULT_READ_BUFFER_LIMIT
PROTOCOL_VERSION = 1
def encode_jsonl_request(request_id: str, op: str, data: Any) -> str:
return json.dumps({'v': PROTOCOL_VERSION, 'id': request_id, 'op': op, 'data': data}) + '\n'
def drain_jsonl_messages(
buffer: bytearray, data: bytes, buffer_limit: int = DEFAULT_READ_BUFFER_LIMIT
) -> list[dict[str, Any]]:
buffer.extend(data)
messages: list[dict[str, Any]] = []
while True:
try:
newline_index = buffer.index(b'\n')
except ValueError:
break
line = bytes(buffer[:newline_index])
del buffer[: newline_index + 1]
if not line:
continue
try:
message = json.loads(line.decode())
except (JSONDecodeError, UnicodeDecodeError):
logger.warning(f'Invalid JSONL message from device, dropping data: {line!r}')
continue
if not isinstance(message, dict):
logger.warning(f'JSONL message is not an object, dropping data: {message!r}')
continue
messages.append(message)
if len(buffer) > buffer_limit:
logger.warning(f'JSONL receive buffer exceeded {buffer_limit} bytes, dropping buffered data')
buffer.clear()
return messages
def resolve_pending_response(pending_requests: dict[str, asyncio.Future[Any]], message: dict[str, Any]) -> bool:
request_id = message.get('id')
if not isinstance(request_id, str):
return False
future = pending_requests.get(request_id)
if future is None:
return False
pending_requests.pop(request_id)
if future.done():
return True
version = message.get('v')
if version is not None and version != PROTOCOL_VERSION:
future.set_exception(RuntimeError(f'Invalid protocol response: unsupported version {version!r}'))
return True
if 'ok' in message:
ok = message['ok']
if not isinstance(ok, bool):
future.set_exception(RuntimeError('Invalid protocol response: ok must be boolean'))
return True
if ok:
if 'data' not in message:
future.set_exception(RuntimeError('Invalid protocol response: missing data for successful response'))
return True
future.set_result(message['data'])
return True
error = message.get('error')
if not isinstance(error, str) or not error:
future.set_exception(RuntimeError('Invalid protocol response: error must be non-empty string'))
return True
future.set_exception(RuntimeError(error))
return True
if 'error' in message:
future.set_exception(RuntimeError(str(message['error'])))
return True
if 'response' in message:
future.set_result(message['response'])
return True
future.set_exception(RuntimeError('Invalid protocol response: missing ok/data, response, or error'))
return True
@@ -0,0 +1,16 @@
# SPDX-FileCopyrightText: 2026 Espressif Systems (Shanghai) CO LTD
# SPDX-License-Identifier: Apache-2.0
from typing import Any
from pydantic import BaseModel
from pydantic import Field
MAX_OP_LENGTH = 64
MAX_REQUEST_DATA_BYTES = 4096
class BLEUARTRequestPayload(BaseModel):
op: str = Field('raw', min_length=1, max_length=MAX_OP_LENGTH, description='Operation name to send to BLE device')
data: Any = Field(..., description='Request payload to send to BLE device')
timeout: float = Field(10.0, gt=0, description='Response timeout in seconds')
@@ -0,0 +1,112 @@
# SPDX-FileCopyrightText: 2026 Espressif Systems (Shanghai) CO LTD
# SPDX-License-Identifier: Apache-2.0
import asyncio
import json
from collections.abc import AsyncIterator
from contextlib import asynccontextmanager
from uuid import uuid4
from fastapi import FastAPI
from fastapi import HTTPException
from loguru import logger
from ..core import BLEUARTBridge
from .jsonl import PROTOCOL_VERSION
from .jsonl import drain_jsonl_messages
from .jsonl import encode_jsonl_request
from .jsonl import resolve_pending_response
from .models import MAX_REQUEST_DATA_BYTES
from .models import BLEUARTRequestPayload
@asynccontextmanager
async def lifespan(app: FastAPI) -> AsyncIterator[None]:
# Server initialization
app.state.bridge = BLEUARTBridge(app.state.device_id)
app.state.request_lock = asyncio.Lock()
# Set BLE UART Bridge RX callback
loop = asyncio.get_running_loop()
app.state.rx_buffer = bytearray()
app.state.pending_requests = {}
def _handle_rx_data(data: bytes) -> None:
for message in drain_jsonl_messages(app.state.rx_buffer, data):
if resolve_pending_response(app.state.pending_requests, message):
continue
logger.debug(f'Received unsolicited BLE UART message: {message!r}')
def _rx_handler(data: bytearray) -> None:
try:
loop.call_soon_threadsafe(_handle_rx_data, bytes(data))
except RuntimeError:
logger.warning(f'Event loop is unavailable, dropping data: {data.decode(errors="replace")}')
app.state.bridge.add_rx_handler(_rx_handler)
# Try to connect to the device
if not await app.state.bridge.connect():
logger.error('Failed to start BLE UART Daemon!')
raise RuntimeError('Failed to start BLE UART Daemon!')
yield
# Disconnect from the device
await app.state.bridge.disconnect()
app = FastAPI(title='BLE UART Daemon', lifespan=lifespan)
def _request_data_size(data: object) -> int:
return len(json.dumps(data).encode())
@app.get('/status')
async def status() -> dict:
bridge: BLEUARTBridge | None = getattr(app.state, 'bridge', None)
pending_requests: dict | None = getattr(app.state, 'pending_requests', None)
return {
'device_id': getattr(app.state, 'device_id', None),
'connection_state': bridge.connection_state.value if bridge else 'DISCONNECTED',
'is_connected': bridge.is_connected if bridge else False,
'pending_requests': len(pending_requests or {}),
'single_flight': True,
'max_request_data_bytes': MAX_REQUEST_DATA_BYTES,
'protocol': f'esp-jsonl-rpc-lite-v{PROTOCOL_VERSION}',
}
@app.post('/request')
async def request(payload: BLEUARTRequestPayload) -> dict:
if _request_data_size(payload.data) > MAX_REQUEST_DATA_BYTES:
raise HTTPException(status_code=413, detail=f'Request data exceeds {MAX_REQUEST_DATA_BYTES} bytes')
# Request with coroutine lock
async with app.state.request_lock:
request_id = uuid4().hex
response_future = asyncio.get_running_loop().create_future()
app.state.pending_requests[request_id] = response_future
try:
# Send command to BLE device
success = await app.state.bridge.send(
encode_jsonl_request(request_id, payload.op, payload.data),
with_response=True,
)
if not success:
raise HTTPException(status_code=500, detail='Failed to send data to device')
# Wait for BLE device to respond
try:
response = await asyncio.wait_for(response_future, timeout=payload.timeout)
except asyncio.TimeoutError:
raise HTTPException(status_code=504, detail='Response timeout')
except RuntimeError as e:
raise HTTPException(status_code=502, detail=str(e))
finally:
app.state.pending_requests.pop(request_id, None)
return {'ok': True, 'data': response}