diff --git a/tools/ble/ble_uart_bridge/src/daemon/__init__.py b/tools/ble/ble_uart_bridge/src/daemon/__init__.py new file mode 100644 index 00000000000..1aefe89058c --- /dev/null +++ b/tools/ble/ble_uart_bridge/src/daemon/__init__.py @@ -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'] diff --git a/tools/ble/ble_uart_bridge/src/daemon/api.py b/tools/ble/ble_uart_bridge/src/daemon/api.py new file mode 100644 index 00000000000..212f76a60d8 --- /dev/null +++ b/tools/ble/ble_uart_bridge/src/daemon/api.py @@ -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)) diff --git a/tools/ble/ble_uart_bridge/src/daemon/jsonl.py b/tools/ble/ble_uart_bridge/src/daemon/jsonl.py new file mode 100644 index 00000000000..5e0702426dc --- /dev/null +++ b/tools/ble/ble_uart_bridge/src/daemon/jsonl.py @@ -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 diff --git a/tools/ble/ble_uart_bridge/src/daemon/models.py b/tools/ble/ble_uart_bridge/src/daemon/models.py new file mode 100644 index 00000000000..1eb59ea3aca --- /dev/null +++ b/tools/ble/ble_uart_bridge/src/daemon/models.py @@ -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') diff --git a/tools/ble/ble_uart_bridge/src/daemon/server.py b/tools/ble/ble_uart_bridge/src/daemon/server.py new file mode 100644 index 00000000000..2e5afbc2f02 --- /dev/null +++ b/tools/ble/ble_uart_bridge/src/daemon/server.py @@ -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}