mirror of
https://github.com/espressif/esp-idf.git
synced 2026-10-01 18:50:34 +03:00
feat(ble): add BLE UART daemon RPC API
This commit is contained in:
@@ -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}
|
||||
Reference in New Issue
Block a user