diff --git a/docs/source/reference/package-apis/drivers/bt-peer.md b/docs/source/reference/package-apis/drivers/bt-peer.md new file mode 100644 index 000000000..843bc16a2 --- /dev/null +++ b/docs/source/reference/package-apis/drivers/bt-peer.md @@ -0,0 +1,59 @@ +# BT Peer Driver + +`jumpstarter-driver-bt-peer` provides a Bluetooth peer device powered by [bumble](https://github.com/google/bumble). +It can pair, connect, and stream A2DP audio to a DUT over BR/EDR. +Transport-agnostic: works over TCP (rootcanal/netsim), USB dongle, serial +UART, or any bumble transport string. + +## Installation + +```{code-block} console +:substitutions: +$ pip3 install --extra-index-url {{index_url}} jumpstarter-driver-bt-peer +``` + +## Configuration + +```yaml +export: + bt_peer: + type: jumpstarter_driver_bt_peer.driver.BtPeer + config: + transport: "tcp-client:127.0.0.1:7300" +``` + +### Config parameters + +| Parameter | Description | Type | Required | Default | +| --------- | ----------- | ---- | -------- | ------- | +| transport | Bumble transport string (e.g. `tcp-client:host:port`, `usb:0`, `serial:/dev/ttyUSB0`) | str | no | `tcp-client:127.0.0.1:7300` | + +## API Reference + +```{eval-rst} +.. autoclass:: jumpstarter_driver_bt_peer.client.BtPeerClient() + :members: +``` + +### CLI + +```console +jumpstarter ⚡ local ➤ j bt_peer +Usage: j bt_peer [OPTIONS] COMMAND [ARGS]... + + Bluetooth peer device (bumble). + +Options: + --help Show this message and exit. + +Commands: + address Show the peer's Bluetooth address. + connect Connect to a remote device. + connections Show active connections. + events Show events since timestamp. + pair Authenticate and encrypt a connection. + start Start the BT peer device. + stop Stop the BT peer device. + wait-connection Wait for an incoming connection. + wait-disconnection Wait for a disconnection. +``` diff --git a/docs/source/reference/package-apis/drivers/index.md b/docs/source/reference/package-apis/drivers/index.md index 43a6f75b8..0eadf8441 100644 --- a/docs/source/reference/package-apis/drivers/index.md +++ b/docs/source/reference/package-apis/drivers/index.md @@ -29,6 +29,7 @@ Drivers that provide various communication interfaces: - {doc}`ADB ` (`jumpstarter-driver-adb`) - Android Debug Bridge tunneling - {doc}`BLE ` (`jumpstarter-driver-ble`) - Bluetooth Low Energy communication +- {doc}`BT Peer ` (`jumpstarter-driver-bt-peer`) - Bluetooth peer device powered by bumble - {doc}`CAN ` (`jumpstarter-driver-can`) - Controller Area Network communication - {doc}`HTTP ` (`jumpstarter-driver-http`) - HTTP communication - {doc}`mitmproxy ` (`jumpstarter-driver-mitmproxy`) - HTTP/HTTPS interception, mocking, and traffic recording @@ -101,6 +102,7 @@ General-purpose utility drivers: adb.md androidemulator.md ble.md +bt-peer.md can.md corellium.md doip.md diff --git a/python/packages/jumpstarter-driver-bt-peer/.gitignore b/python/packages/jumpstarter-driver-bt-peer/.gitignore new file mode 100644 index 000000000..dbb6e9b82 --- /dev/null +++ b/python/packages/jumpstarter-driver-bt-peer/.gitignore @@ -0,0 +1,4 @@ +__pycache__/ +.coverage +coverage.xml +htmlcov/ diff --git a/python/packages/jumpstarter-driver-bt-peer/README.md b/python/packages/jumpstarter-driver-bt-peer/README.md new file mode 100644 index 000000000..52801d9a6 --- /dev/null +++ b/python/packages/jumpstarter-driver-bt-peer/README.md @@ -0,0 +1,47 @@ +# BtPeer Driver + +`jumpstarter-driver-bt-peer` provides a Bluetooth peer device powered by +[bumble](https://github.com/google/bumble). It can pair, connect, and stream +A2DP audio to a DUT over BR/EDR. + +## Installation + +```shell +pip3 install --extra-index-url https://pkg.jumpstarter.dev/simple/ jumpstarter-driver-bt-peer +``` + +## Configuration + +Example configuration: + +```yaml +export: + bt_peer: + type: jumpstarter_driver_bt_peer.driver.BtPeer + config: + transport: "tcp-client:127.0.0.1:7300" # bumble transport string +``` + +## Usage + +Start the peer, pair with a DUT, and verify the connection: + +```bash +j bt_peer start '{"name": "Bumble-Phone"}' +j bt_peer address +j bt_peer wait-connection --timeout 60 +j bt_peer connections +j bt_peer pair --handle 0 +j bt_peer stop +``` + +The `transport` config accepts any bumble transport string: +- `tcp-client:host:port` - rootcanal / netsim +- `usb:0` — USB HCI dongle +- `serial:/dev/ttyUSB0` - serial UART + +## API Reference + +```{eval-rst} +.. autoclass:: jumpstarter_driver_bt_peer.driver.BtPeer() +``` diff --git a/python/packages/jumpstarter-driver-bt-peer/examples/exporter.yaml b/python/packages/jumpstarter-driver-bt-peer/examples/exporter.yaml new file mode 100644 index 000000000..92666c91a --- /dev/null +++ b/python/packages/jumpstarter-driver-bt-peer/examples/exporter.yaml @@ -0,0 +1,12 @@ +apiVersion: jumpstarter.dev/v1alpha1 +kind: ExporterConfig +metadata: + namespace: default + name: bt-peer +endpoint: grpc.jumpstarter.192.168.0.203.nip.io:8082 +token: "" +export: + bt_peer: + type: jumpstarter_driver_bt_peer.driver.BtPeer + config: + transport: "tcp-client:127.0.0.1:7300" diff --git a/python/packages/jumpstarter-driver-bt-peer/jumpstarter_driver_bt_peer/__init__.py b/python/packages/jumpstarter-driver-bt-peer/jumpstarter_driver_bt_peer/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/python/packages/jumpstarter-driver-bt-peer/jumpstarter_driver_bt_peer/client.py b/python/packages/jumpstarter-driver-bt-peer/jumpstarter_driver_bt_peer/client.py new file mode 100644 index 000000000..54ba33f9b --- /dev/null +++ b/python/packages/jumpstarter-driver-bt-peer/jumpstarter_driver_bt_peer/client.py @@ -0,0 +1,139 @@ +import json +from typing import Any + +import click + +from jumpstarter.client import DriverClient + + +def _parse(raw: str) -> dict[str, Any] | list[Any] | str: + try: + return json.loads(raw) + except (json.JSONDecodeError, TypeError): + return raw + + +def _parse_dict(raw: str) -> dict[str, Any]: + result = _parse(raw) + if not isinstance(result, dict): + raise ValueError(f"expected dict, got {type(result).__name__}: {raw!r}") + return result + + +def _parse_list(raw: str) -> list[Any]: + result = _parse(raw) + if not isinstance(result, list): + raise ValueError(f"expected list, got {type(result).__name__}: {raw!r}") + return result + + +def _echo(obj: object) -> None: + if isinstance(obj, (dict, list)): + click.echo(json.dumps(obj, indent=2)) + else: + click.echo(obj) + + +class BtPeerClient(DriverClient): + """Client for the Bluetooth peer driver.""" + + def start_peer(self, config_json: str = "{}") -> dict[str, Any]: + return _parse_dict(self.call("start_peer", config_json)) + + def stop_peer(self) -> dict[str, Any]: + return _parse_dict(self.call("stop_peer")) + + def wait_connection(self, timeout: int = 30) -> dict[str, Any]: + return _parse_dict(self.call("wait_connection", timeout)) + + def wait_disconnection(self, timeout: int = 30) -> dict[str, Any]: + return _parse_dict(self.call("wait_disconnection", timeout)) + + def get_events(self, since: str = "0") -> list[Any]: + return _parse_list(self.call("get_events", since)) + + def get_address(self) -> str: + return self.call("get_address") + + def pair(self, handle: int = 0) -> dict[str, Any]: + return _parse_dict(self.call("pair", handle)) + + def connect_to(self, address: str, timeout: int = 30) -> dict[str, Any]: + return _parse_dict(self.call("connect_to", address, timeout)) + + def get_connections(self) -> list[Any]: + return _parse_list(self.call("get_connections")) + + def cli(self): # noqa: C901 + @click.group() + def bt_peer(): + """Bluetooth peer device (bumble).""" + + @bt_peer.command("start") + @click.argument("config", default="{}") + def start_cmd(config: str): + """Start the BT peer device. + + CONFIG is JSON: {"name": "...", "classic_enabled": true, "class_of_device": 123} + """ + try: + parsed = json.loads(config) + except json.JSONDecodeError as exc: + raise click.BadParameter( + f"CONFIG must be valid JSON: {exc.msg}", + param_hint="CONFIG", + ) from exc + if not isinstance(parsed, dict): + raise click.BadParameter( + "CONFIG must be a JSON object", + param_hint="CONFIG", + ) + _echo(self.start_peer(config)) + + @bt_peer.command("stop") + def stop_cmd(): + """Stop the BT peer device.""" + _echo(self.stop_peer()) + + @bt_peer.command("wait-connection") + @click.option("--timeout", "-t", default=30, help="Timeout in seconds") + def wait_connection_cmd(timeout: int): + """Wait for an incoming connection.""" + _echo(self.wait_connection(timeout)) + + @bt_peer.command("wait-disconnection") + @click.option("--timeout", "-t", default=30, help="Timeout in seconds") + def wait_disconnection_cmd(timeout: int): + """Wait for a disconnection.""" + _echo(self.wait_disconnection(timeout)) + + @bt_peer.command("events") + @click.option("--since", default="0", help="Timestamp filter") + def events_cmd(since: str): + """Show events since timestamp.""" + _echo(self.get_events(since)) + + @bt_peer.command("address") + def address_cmd(): + """Show the peer's Bluetooth address.""" + click.echo(self.get_address()) + + @bt_peer.command("pair") + @click.option("--handle", "-h", default=0, help="Connection handle") + def pair_cmd(handle: int): + """Authenticate and encrypt a connection.""" + _echo(self.pair(handle)) + + @bt_peer.command("connect") + @click.argument("address") + @click.option("--timeout", "-t", default=30, help="Timeout in seconds") + def connect_cmd(address: str, timeout: int): + """Connect to a remote device (e.g. the CVD).""" + _echo(self.connect_to(address, timeout)) + + @bt_peer.command("connections") + def connections_cmd(): + """Show active connections.""" + _echo(self.get_connections()) + + return bt_peer diff --git a/python/packages/jumpstarter-driver-bt-peer/jumpstarter_driver_bt_peer/driver.py b/python/packages/jumpstarter-driver-bt-peer/jumpstarter_driver_bt_peer/driver.py new file mode 100644 index 000000000..3e839ab68 --- /dev/null +++ b/python/packages/jumpstarter-driver-bt-peer/jumpstarter_driver_bt_peer/driver.py @@ -0,0 +1,454 @@ +import json +import time +from collections import deque +from dataclasses import dataclass, field +from typing import Any + +import anyio +from bumble.a2dp import ( + A2DP_SBC_CODEC_TYPE, + SbcMediaCodecInformation, + make_audio_source_service_sdp_records, +) +from bumble.avdtp import ( + AVDTP_AUDIO_MEDIA_TYPE, + Listener, + MediaCodecCapabilities, + MediaPacketPump, +) +from bumble.core import BT_BR_EDR_TRANSPORT, ClassOfDevice, DeviceClass +from bumble.device import Connection, Device, DeviceConfiguration +from bumble.hci import HCI_PERIPHERAL_ROLE +from bumble.host import Host +from bumble.pairing import PairingConfig, PairingDelegate +from bumble.rtp import MediaPacket +from bumble.transport import open_transport + +from jumpstarter.driver import Driver, export + + +class BtPeerError(Exception): + pass + + +class AutoAcceptDelegate(PairingDelegate): + """Auto-accepts all pairing requests""" + + def __init__(self, io_capability=None): + super().__init__( + io_capability=io_capability or PairingDelegate.IoCapability.NO_OUTPUT_NO_INPUT, + local_initiator_key_distribution=PairingDelegate.DEFAULT_KEY_DISTRIBUTION, + local_responder_key_distribution=PairingDelegate.DEFAULT_KEY_DISTRIBUTION, + ) + + async def confirm(self, auto=False) -> bool: + return True + + async def compare_numbers(self, number: int, digits: int) -> bool: + return True + + async def accept(self) -> bool: + return True + + async def get_number(self) -> int | None: + return 0 + + +@dataclass(kw_only=True) +class BtPeer(Driver): + """Bluetooth peer device powered by bumble. + + Exposes a configurable Bluetooth device that can pair, connect, and + stream A2DP audio to a DUT over BR/EDR. Transport-agnostic: works over + TCP (rootcanal/netsim), USB dongle, serial UART, or any bumble transport. + """ + + transport: str = "tcp-client:127.0.0.1:7300" + + _device: Device | None = field(default=None, init=False, repr=False) + _transport: Any = field(default=None, init=False, repr=False) + _avdtp_listener: Listener | None = field(default=None, init=False, repr=False) + _events: deque = field(default_factory=lambda: deque(maxlen=1000), init=False, repr=False) + _connections: dict[int, Connection] = field(default_factory=dict, init=False, repr=False) + + @classmethod + def client(cls) -> str: + return "jumpstarter_driver_bt_peer.client.BtPeerClient" + + def _emit(self, event_type: str, data: dict | None = None) -> None: + entry = { + "ts": time.time(), + "type": event_type, + **(data or {}), + } + self._events.append(entry) + self.logger.info("event: %s %s", event_type, json.dumps(data or {})) + + def _on_connection(self, connection: Connection) -> None: + self._connections[connection.handle] = connection + self._emit("connection", { + "handle": connection.handle, + "address": str(connection.peer_address), + "transport": "BR/EDR" if connection.transport == BT_BR_EDR_TRANSPORT else "LE", + "role": "peripheral" if connection.role == HCI_PERIPHERAL_ROLE else "central", + "encrypted": connection.is_encrypted, + }) + + def on_disconnect(reason): + self._connections.pop(connection.handle, None) + self._emit("disconnection", { + "handle": connection.handle, + "address": str(connection.peer_address), + "reason": reason, + }) + + connection.on("disconnection", on_disconnect) + + def _on_avdtp_connection(self, protocol) -> None: + S = SbcMediaCodecInformation + sbc_capabilities = MediaCodecCapabilities( + media_type=AVDTP_AUDIO_MEDIA_TYPE, + media_codec_type=A2DP_SBC_CODEC_TYPE, + media_codec_information=S( + sampling_frequency=S.SamplingFrequency.SF_48000 | S.SamplingFrequency.SF_44100, + channel_mode=S.ChannelMode.JOINT_STEREO | S.ChannelMode.STEREO, + block_length=S.BlockLength.BL_16 | S.BlockLength.BL_12 | S.BlockLength.BL_8, + subbands=S.Subbands.S_8, + allocation_method=S.AllocationMethod.LOUDNESS, + minimum_bitpool_value=2, + maximum_bitpool_value=53, + ), + ) + + async def silence(): + """Yield timestamped RTP packets so MediaPacketPump can stream.""" + sequence_number = 0 + timestamp_seconds = 0.0 + while True: + packet = MediaPacket( + 2, + 0, + 0, + 0, + sequence_number, + int(timestamp_seconds * 8000), + 0, + [], + 96, + b"\x00", + ) + packet.timestamp_seconds = timestamp_seconds + yield packet + sequence_number = (sequence_number + 1) & 0xFFFF + timestamp_seconds += 0.02 + await anyio.sleep(0.02) + + pump = MediaPacketPump(silence()) + protocol.add_source(sbc_capabilities, pump) + self._emit("avdtp_connected", {}) + + @export + async def start_peer(self, config_json: str = "{}") -> str: + """Start the Bluetooth peer device. + + config_json fields: + name: device name (default: "Bumble-Phone") + classic_enabled: enable BR/EDR (default: true) + class_of_device: CoD integer (default: smartphone with audio/telephony) + + Returns JSON with assigned address. + """ + if self._device is not None: + raise BtPeerError("peer already running — call stop_peer first") + + try: + config = json.loads(config_json) if config_json else {} + except json.JSONDecodeError as exc: + raise BtPeerError(f"invalid config JSON: {exc.msg}") from exc + if not isinstance(config, dict): + raise BtPeerError( + f"config must be a JSON object, got {type(config).__name__}" + ) + + name = config.get("name", "Bumble-Phone") + classic_enabled = config.get("classic_enabled", True) + cod = config.get("class_of_device", None) + + self.logger.info("opening HCI transport: %s", self.transport) + + self._transport = await open_transport(self.transport) + + try: + device_config = DeviceConfiguration() + device_config.name = name + device_config.classic_enabled = classic_enabled + + if cod is not None: + device_config.class_of_device = cod + else: + device_config.class_of_device = int(ClassOfDevice( + ClassOfDevice.MajorServiceClasses.AUDIO | ClassOfDevice.MajorServiceClasses.TELEPHONY, + ClassOfDevice.MajorDeviceClass.PHONE, + DeviceClass.PHONE_SMARTPHONE_MINOR_DEVICE_CLASS, + )) + + device_config.classic_sc_enabled = True + device_config.classic_ssp_enabled = True + + host = Host( + controller_source=self._transport.source, + controller_sink=self._transport.sink, + ) + self._device = Device(config=device_config, host=host) + + self._device.pairing_config_factory = lambda connection: PairingConfig( + sc=True, + mitm=False, + bonding=True, + delegate=AutoAcceptDelegate(), + ) + + service_record_handle = 0x00010001 + self._device.sdp_service_records = { + service_record_handle: make_audio_source_service_sdp_records( + service_record_handle + ) + } + + self._avdtp_listener = Listener.for_device(self._device) + self._avdtp_listener.on("connection", self._on_avdtp_connection) + + self._device.on("connection", self._on_connection) + + await self._device.power_on() + + if classic_enabled: + await self._device.set_discoverable(True) + await self._device.set_connectable(True) + except Exception: + if self._device is not None: + try: + await self._device.power_off() + except Exception: + pass + self._device = None + self._avdtp_listener = None + if self._transport is not None: + await self._transport.close() + self._transport = None + raise + + address = str(self._device.public_address) + self._emit("peer_started", {"address": address, "name": name}) + + return json.dumps({"address": address, "name": name}) + + @export + async def stop_peer(self) -> str: + """Stop the Bluetooth peer device and clean up.""" + if self._device is None and self._transport is None: + return json.dumps({"status": "not_running"}) + + self._avdtp_listener = None + device = self._device + transport = self._transport + self._device = None + self._transport = None + self._connections.clear() + self._events.clear() + + first_error: Exception | None = None + try: + if device is not None: + await device.power_off() + except Exception as exc: + first_error = exc + finally: + try: + if transport is not None: + await transport.close() + except Exception as exc: + if first_error is None: + first_error = exc + + self._emit("peer_stopped", {}) + + if first_error is not None: + raise BtPeerError(f"failed to stop peer: {first_error}") from first_error + + return json.dumps({"status": "stopped"}) + + @export + async def wait_connection(self, timeout: int = 30) -> str: + """Wait for an incoming BR/EDR or LE connection. + + Returns JSON with connection details (handle, address, transport, role, encrypted). + Raises on timeout. + """ + if self._device is None: + raise BtPeerError("peer not running — call start_peer first") + + if self._connections: + conn = next(iter(self._connections.values())) + return json.dumps({ + "handle": conn.handle, + "address": str(conn.peer_address), + "transport": "BR/EDR" if conn.transport == BT_BR_EDR_TRANSPORT else "LE", + "role": "peripheral" if conn.role == HCI_PERIPHERAL_ROLE else "central", + "encrypted": conn.is_encrypted, + }) + + event = anyio.Event() + result_holder: list[Connection] = [] + + def on_connect(connection: Connection): + result_holder.append(connection) + event.set() + + self._device.on("connection", on_connect) + try: + with anyio.fail_after(timeout): + await event.wait() + except TimeoutError: + raise BtPeerError(f"no connection within {timeout}s") from None + finally: + self._device.remove_listener("connection", on_connect) + + conn = result_holder[0] + return json.dumps({ + "handle": conn.handle, + "address": str(conn.peer_address), + "transport": "BR/EDR" if conn.transport == BT_BR_EDR_TRANSPORT else "LE", + "role": "peripheral" if conn.role == HCI_PERIPHERAL_ROLE else "central", + "encrypted": conn.is_encrypted, + }) + + @export + async def wait_disconnection(self, timeout: int = 30) -> str: + """Wait for any active connection to disconnect. + + Returns JSON with disconnection details (handle, address, reason). + Raises on timeout if no disconnection occurs. + """ + if self._device is None: + raise BtPeerError("peer not running") + if not self._connections: + return json.dumps({"status": "no_connections"}) + + event = anyio.Event() + result_holder: list[dict] = [] + registered: list[tuple[Connection, Any]] = [] + + for conn in list(self._connections.values()): + def on_disconnect(reason, c=conn): + result_holder.append({ + "handle": c.handle, + "address": str(c.peer_address), + "reason": reason, + }) + event.set() + + conn.on("disconnection", on_disconnect) + registered.append((conn, on_disconnect)) + + try: + with anyio.fail_after(timeout): + await event.wait() + except TimeoutError: + raise BtPeerError(f"no disconnection within {timeout}s") from None + finally: + for registered_conn, registered_handler in registered: + registered_conn.remove_listener("disconnection", registered_handler) + + return json.dumps(result_holder[0]) + + @export + async def pair(self, handle: int = 0) -> str: + """Initiate authentication + encryption on a connection. + + AutoAcceptDelegate handles numeric comparison automatically. + """ + if self._device is None: + raise BtPeerError("peer not running") + + handle = int(handle) + conn = self._connections.get(handle) + if conn is None: + raise BtPeerError(f"no connection with handle {handle}") + + await conn.authenticate() + await conn.encrypt() + + self._emit("paired", { + "handle": handle, + "address": str(conn.peer_address), + "encrypted": conn.is_encrypted, + }) + + return json.dumps({ + "handle": handle, + "address": str(conn.peer_address), + "encrypted": conn.is_encrypted, + }) + + @export + async def connect_to(self, address: str, timeout: int = 30) -> str: + """Initiate outgoing BR/EDR connection to a remote device. + + Returns JSON with connection details. + """ + if self._device is None: + raise BtPeerError("peer not running - call start_peer first") + + from bumble.hci import Address as HciAddress + + target = HciAddress( + address, + address_type=HciAddress.PUBLIC_DEVICE_ADDRESS, + ) + self.logger.info("connecting to %s", target) + try: + with anyio.fail_after(timeout): + connection = await self._device.connect( + target, transport=BT_BR_EDR_TRANSPORT + ) + except TimeoutError: + raise BtPeerError(f"connect to {address} timed out after {timeout}s") from None + + return json.dumps({ + "handle": connection.handle, + "address": str(connection.peer_address), + "transport": "BR/EDR" if connection.transport == BT_BR_EDR_TRANSPORT else "LE", + "role": "central", + "encrypted": connection.is_encrypted, + }) + + @export + def get_events(self, since: str = "0") -> str: + """Get events since a timestamp. Pass "0" for all events.""" + try: + since_ts = float(since) + except (TypeError, ValueError) as exc: + raise BtPeerError(f"invalid since value: {since!r}") from exc + events = [e for e in self._events if e["ts"] > since_ts] + return json.dumps(events) + + @export + def get_address(self) -> str: + """Return the peer's assigned Bluetooth address.""" + if self._device is None: + raise BtPeerError("peer not running") + return str(self._device.public_address) + + @export + def get_connections(self) -> str: + """Return currently active connections as JSON array.""" + conns = [] + for conn in self._connections.values(): + conns.append({ + "handle": conn.handle, + "address": str(conn.peer_address), + "transport": "BR/EDR" if conn.transport == BT_BR_EDR_TRANSPORT else "LE", + "encrypted": conn.is_encrypted, + }) + return json.dumps(conns) diff --git a/python/packages/jumpstarter-driver-bt-peer/jumpstarter_driver_bt_peer/driver_test.py b/python/packages/jumpstarter-driver-bt-peer/jumpstarter_driver_bt_peer/driver_test.py new file mode 100644 index 000000000..1a08a74bd --- /dev/null +++ b/python/packages/jumpstarter-driver-bt-peer/jumpstarter_driver_bt_peer/driver_test.py @@ -0,0 +1,615 @@ +import json +from unittest.mock import AsyncMock, MagicMock, patch + +import pytest + +from .driver import AutoAcceptDelegate, BtPeer, BtPeerError + + +@pytest.fixture +def anyio_backend(): + return "asyncio" + + +def test_bt_peer_default_transport(): + peer = BtPeer() + assert peer.transport == "tcp-client:127.0.0.1:7300" + + +def test_bt_peer_custom_transport(): + peer = BtPeer(transport="usb:0") + assert peer.transport == "usb:0" + + +def test_auto_accept_delegate_init(): + delegate = AutoAcceptDelegate() + assert delegate is not None + + +@pytest.mark.anyio +async def test_auto_accept_delegate_confirms(): + delegate = AutoAcceptDelegate() + assert await delegate.confirm() is True + assert await delegate.compare_numbers(123456, 6) is True + assert await delegate.accept() is True + assert await delegate.get_number() == 0 + + +def _make_mock_transport(): + transport = MagicMock() + transport.source = MagicMock() + transport.sink = MagicMock() + transport.close = AsyncMock() + return transport + + +def _make_mock_device(address="DA:4C:10:DE:00:01"): + device = MagicMock() + device.public_address = address + device.power_on = AsyncMock() + device.power_off = AsyncMock() + device.set_discoverable = AsyncMock() + device.set_connectable = AsyncMock() + device.on = MagicMock() + device.remove_listener = MagicMock() + device.connect = AsyncMock() + device.sdp_service_records = {} + device.pairing_config_factory = None + return device + + +@pytest.mark.anyio +async def test_start_peer_and_stop(): + peer = BtPeer() + mock_transport = _make_mock_transport() + mock_device = _make_mock_device() + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + return_value=mock_device, + ), + patch("jumpstarter_driver_bt_peer.driver.Listener") as mock_listener_cls, + ): + mock_avdtp = MagicMock() + mock_listener_cls.for_device.return_value = mock_avdtp + + result = await peer.start_peer('{"name": "Test-Phone"}') + parsed = json.loads(result) + assert parsed["name"] == "Test-Phone" + assert parsed["address"] == "DA:4C:10:DE:00:01" + mock_device.power_on.assert_awaited_once() + mock_device.set_discoverable.assert_awaited_once_with(True) + mock_device.set_connectable.assert_awaited_once_with(True) + + with pytest.raises(BtPeerError, match="already running"): + await peer.start_peer() + + result = await peer.stop_peer() + parsed = json.loads(result) + assert parsed["status"] == "stopped" + mock_device.power_off.assert_awaited_once() + mock_transport.close.assert_awaited_once() + + +@pytest.mark.anyio +async def test_stop_peer_not_running(): + peer = BtPeer() + result = await peer.stop_peer() + parsed = json.loads(result) + assert parsed["status"] == "not_running" + + +@pytest.mark.anyio +async def test_start_peer_rejects_non_object_config(): + peer = BtPeer() + with pytest.raises(BtPeerError, match="JSON object"): + await peer.start_peer("[]") + with pytest.raises(BtPeerError, match="invalid config JSON"): + await peer.start_peer("{") + + +@pytest.mark.anyio +async def test_stop_peer_cleans_up_when_power_off_fails(): + peer = BtPeer() + mock_transport = _make_mock_transport() + mock_device = _make_mock_device() + mock_device.power_off = AsyncMock(side_effect=RuntimeError("power off failed")) + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + return_value=mock_device, + ), + patch("jumpstarter_driver_bt_peer.driver.Listener") as mock_listener_cls, + ): + mock_listener_cls.for_device.return_value = MagicMock() + await peer.start_peer() + with pytest.raises(BtPeerError, match="failed to stop peer"): + await peer.stop_peer() + + assert peer._device is None + assert peer._transport is None + mock_transport.close.assert_awaited_once() + + +@pytest.mark.anyio +async def test_start_peer_default_config(): + peer = BtPeer() + mock_transport = _make_mock_transport() + mock_device = _make_mock_device() + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + return_value=mock_device, + ), + patch("jumpstarter_driver_bt_peer.driver.Listener") as mock_listener_cls, + ): + mock_listener_cls.for_device.return_value = MagicMock() + + result = await peer.start_peer() + parsed = json.loads(result) + assert parsed["name"] == "Bumble-Phone" + + await peer.stop_peer() + + +@pytest.mark.anyio +async def test_start_peer_custom_cod(): + peer = BtPeer() + mock_transport = _make_mock_transport() + mock_device = _make_mock_device() + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + return_value=mock_device, + ) as device_cls, + patch("jumpstarter_driver_bt_peer.driver.Listener") as mock_listener_cls, + ): + mock_listener_cls.for_device.return_value = MagicMock() + + result = await peer.start_peer('{"name": "Test", "class_of_device": 42}') + parsed = json.loads(result) + assert parsed["name"] == "Test" + config = device_cls.call_args[1]["config"] + assert config.class_of_device == 42 + + await peer.stop_peer() + + +@pytest.mark.anyio +async def test_start_peer_classic_disabled(): + peer = BtPeer() + mock_transport = _make_mock_transport() + mock_device = _make_mock_device() + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + return_value=mock_device, + ), + patch("jumpstarter_driver_bt_peer.driver.Listener") as mock_listener_cls, + ): + mock_listener_cls.for_device.return_value = MagicMock() + + await peer.start_peer('{"classic_enabled": false}') + mock_device.set_discoverable.assert_not_awaited() + mock_device.set_connectable.assert_not_awaited() + + await peer.stop_peer() + + +@pytest.mark.anyio +async def test_get_address(): + peer = BtPeer() + with pytest.raises(BtPeerError, match="not running"): + peer.get_address() + + mock_transport = _make_mock_transport() + mock_device = _make_mock_device("AA:BB:CC:DD:EE:FF") + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + return_value=mock_device, + ), + patch("jumpstarter_driver_bt_peer.driver.Listener") as mock_listener_cls, + ): + mock_listener_cls.for_device.return_value = MagicMock() + await peer.start_peer() + assert peer.get_address() == "AA:BB:CC:DD:EE:FF" + await peer.stop_peer() + + +@pytest.mark.anyio +async def test_get_events(): + peer = BtPeer() + mock_transport = _make_mock_transport() + mock_device = _make_mock_device() + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + return_value=mock_device, + ), + patch("jumpstarter_driver_bt_peer.driver.Listener") as mock_listener_cls, + ): + mock_listener_cls.for_device.return_value = MagicMock() + await peer.start_peer() + events_json = peer.get_events("0") + events = json.loads(events_json) + assert isinstance(events, list) + assert len(events) == 1 + assert events[0]["type"] == "peer_started" + with pytest.raises(BtPeerError, match="invalid since value"): + peer.get_events("abc") + await peer.stop_peer() + + +@pytest.mark.anyio +async def test_get_connections_empty(): + peer = BtPeer() + mock_transport = _make_mock_transport() + mock_device = _make_mock_device() + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + return_value=mock_device, + ), + patch("jumpstarter_driver_bt_peer.driver.Listener") as mock_listener_cls, + ): + mock_listener_cls.for_device.return_value = MagicMock() + await peer.start_peer() + conns = json.loads(peer.get_connections()) + assert conns == [] + await peer.stop_peer() + + +@pytest.mark.anyio +async def test_on_connection_tracking(): + peer = BtPeer() + + mock_conn = MagicMock() + mock_conn.handle = 42 + mock_conn.peer_address = "11:22:33:44:55:66" + mock_conn.transport = 0 + mock_conn.role = 1 + mock_conn.is_encrypted = False + mock_conn.on = MagicMock() + + peer._on_connection(mock_conn) + + assert 42 in peer._connections + conns = json.loads(peer.get_connections()) + assert len(conns) == 1 + assert conns[0]["handle"] == 42 + assert conns[0]["address"] == "11:22:33:44:55:66" + + +@pytest.mark.anyio +async def test_on_connection_disconnect(): + peer = BtPeer() + + mock_conn = MagicMock() + mock_conn.handle = 7 + mock_conn.peer_address = "AA:BB:CC:DD:EE:FF" + mock_conn.transport = 0 + mock_conn.role = 0 + mock_conn.is_encrypted = True + mock_conn.on = MagicMock() + + peer._on_connection(mock_conn) + assert 7 in peer._connections + + disconnect_handler = mock_conn.on.call_args[0][1] + disconnect_handler(0x13) + assert 7 not in peer._connections + + +@pytest.mark.anyio +async def test_wait_connection_not_running(): + peer = BtPeer() + with pytest.raises(BtPeerError, match="not running"): + await peer.wait_connection() + + +@pytest.mark.anyio +async def test_wait_connection_timeout(): + peer = BtPeer() + mock_transport = _make_mock_transport() + mock_device = _make_mock_device() + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + return_value=mock_device, + ), + patch("jumpstarter_driver_bt_peer.driver.Listener") as mock_listener_cls, + ): + mock_listener_cls.for_device.return_value = MagicMock() + await peer.start_peer() + with pytest.raises(BtPeerError, match="no connection within"): + await peer.wait_connection(timeout=0) + await peer.stop_peer() + + +@pytest.mark.anyio +async def test_wait_disconnection_not_running(): + peer = BtPeer() + with pytest.raises(BtPeerError, match="not running"): + await peer.wait_disconnection() + + +@pytest.mark.anyio +async def test_wait_disconnection_no_connections(): + peer = BtPeer() + mock_transport = _make_mock_transport() + mock_device = _make_mock_device() + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + return_value=mock_device, + ), + patch("jumpstarter_driver_bt_peer.driver.Listener") as mock_listener_cls, + ): + mock_listener_cls.for_device.return_value = MagicMock() + await peer.start_peer() + result = json.loads(await peer.wait_disconnection()) + assert result["status"] == "no_connections" + await peer.stop_peer() + + +@pytest.mark.anyio +async def test_wait_disconnection_removes_handlers(): + peer = BtPeer() + mock_conn = MagicMock() + mock_conn.handle = 1 + mock_conn.peer_address = "11:22:33:44:55:66" + mock_conn.transport = 0 + mock_conn.role = 1 + mock_conn.is_encrypted = False + mock_conn.on = MagicMock() + mock_conn.remove_listener = MagicMock() + peer._connections[1] = mock_conn + peer._device = MagicMock() + + with pytest.raises(BtPeerError, match="no disconnection within"): + await peer.wait_disconnection(timeout=0) + + mock_conn.on.assert_called_once() + assert mock_conn.on.call_args.args[0] == "disconnection" + mock_conn.remove_listener.assert_called_once_with( + "disconnection", mock_conn.on.call_args.args[1] + ) + + +@pytest.mark.anyio +async def test_pair_not_running(): + peer = BtPeer() + with pytest.raises(BtPeerError, match="not running"): + await peer.pair() + + +@pytest.mark.anyio +async def test_pair_no_connection(): + peer = BtPeer() + mock_transport = _make_mock_transport() + mock_device = _make_mock_device() + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + return_value=mock_device, + ), + patch("jumpstarter_driver_bt_peer.driver.Listener") as mock_listener_cls, + ): + mock_listener_cls.for_device.return_value = MagicMock() + await peer.start_peer() + with pytest.raises(BtPeerError, match="no connection with handle"): + await peer.pair(99) + await peer.stop_peer() + + +@pytest.mark.anyio +async def test_pair_success(): + peer = BtPeer() + + mock_conn = MagicMock() + mock_conn.handle = 1 + mock_conn.peer_address = "11:22:33:44:55:66" + mock_conn.transport = 0 + mock_conn.role = 1 + mock_conn.is_encrypted = True + mock_conn.on = MagicMock() + mock_conn.authenticate = AsyncMock() + mock_conn.encrypt = AsyncMock() + + peer._on_connection(mock_conn) + + mock_transport = _make_mock_transport() + mock_device = _make_mock_device() + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + return_value=mock_device, + ), + patch("jumpstarter_driver_bt_peer.driver.Listener") as mock_listener_cls, + ): + mock_listener_cls.for_device.return_value = MagicMock() + await peer.start_peer() + result = json.loads(await peer.pair(1)) + assert result["encrypted"] is True + assert result["handle"] == 1 + mock_conn.authenticate.assert_awaited_once() + mock_conn.encrypt.assert_awaited_once() + await peer.stop_peer() + + +@pytest.mark.anyio +async def test_connect_to_not_running(): + peer = BtPeer() + with pytest.raises(BtPeerError, match="not running"): + await peer.connect_to("AA:BB:CC:DD:EE:FF") + + +@pytest.mark.anyio +async def test_connect_to(): + peer = BtPeer() + mock_transport = _make_mock_transport() + mock_device = _make_mock_device() + + mock_outgoing = MagicMock() + mock_outgoing.handle = 5 + mock_outgoing.peer_address = "AA:BB:CC:DD:EE:FF" + mock_outgoing.transport = 0 + mock_outgoing.is_encrypted = False + mock_device.connect = AsyncMock(return_value=mock_outgoing) + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + return_value=mock_device, + ), + patch("jumpstarter_driver_bt_peer.driver.Listener") as mock_listener_cls, + ): + mock_listener_cls.for_device.return_value = MagicMock() + await peer.start_peer() + result = json.loads(await peer.connect_to("AA:BB:CC:DD:EE:FF")) + assert result["handle"] == 5 + assert result["role"] == "central" + await peer.stop_peer() + + +@pytest.mark.anyio +async def test_emit_events(): + peer = BtPeer() + peer._emit("test_event", {"key": "value"}) + assert len(peer._events) == 1 + assert peer._events[0]["type"] == "test_event" + assert peer._events[0]["key"] == "value" + assert "ts" in peer._events[0] + + +@pytest.mark.anyio +async def test_on_avdtp_connection(): + peer = BtPeer() + mock_protocol = MagicMock() + + peer._on_avdtp_connection(mock_protocol) + + mock_protocol.add_source.assert_called_once() + pump = mock_protocol.add_source.call_args.args[1] + packets = pump.packets + first = await packets.__anext__() + assert hasattr(first, "timestamp_seconds") + assert bytes(first) + events = list(peer._events) + assert any(e["type"] == "avdtp_connected" for e in events) + + +@pytest.mark.anyio +async def test_start_peer_exception_cleanup(): + peer = BtPeer() + mock_transport = _make_mock_transport() + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + side_effect=RuntimeError("device init failed"), + ), + ): + with pytest.raises(RuntimeError, match="device init failed"): + await peer.start_peer() + + assert peer._device is None + assert peer._transport is None + mock_transport.close.assert_awaited_once() + + +@pytest.mark.anyio +async def test_start_peer_exception_cleanup_after_device_init(): + peer = BtPeer() + mock_transport = _make_mock_transport() + mock_device = _make_mock_device() + mock_device.power_on = AsyncMock(side_effect=RuntimeError("power on failed")) + + with ( + patch( + "jumpstarter_driver_bt_peer.driver.open_transport", + new=AsyncMock(return_value=mock_transport), + ), + patch( + "jumpstarter_driver_bt_peer.driver.Device", + return_value=mock_device, + ), + patch("jumpstarter_driver_bt_peer.driver.Listener") as mock_listener_cls, + ): + mock_listener_cls.for_device.return_value = MagicMock() + with pytest.raises(RuntimeError, match="power on failed"): + await peer.start_peer() + + assert peer._device is None + assert peer._transport is None + mock_device.power_off.assert_awaited_once() + mock_transport.close.assert_awaited_once() diff --git a/python/packages/jumpstarter-driver-bt-peer/pyproject.toml b/python/packages/jumpstarter-driver-bt-peer/pyproject.toml new file mode 100644 index 000000000..ecf0d4785 --- /dev/null +++ b/python/packages/jumpstarter-driver-bt-peer/pyproject.toml @@ -0,0 +1,46 @@ +[project] +name = "jumpstarter-driver-bt-peer" +dynamic = ["version", "urls"] +description = "Bluetooth peer driver powered by bumble" +readme = "README.md" +license = "Apache-2.0" +authors = [ + { name = "Benny Zlotnik", email = "bzlotnik@redhat.com" } +] +requires-python = ">=3.11" +dependencies = [ + "anyio>=4.10.0", + "bumble>=0.0.233", + "jumpstarter", +] + +[project.entry-points."jumpstarter.drivers"] +BtPeer = "jumpstarter_driver_bt_peer.driver:BtPeer" + +[tool.hatch.version] +source = "vcs" +raw-options = { 'root' = '../../../'} + +[tool.hatch.metadata.hooks.vcs.urls] +Homepage = "https://jumpstarter.dev" +source_archive = "https://github.com/jumpstarter-dev/jumpstarter/archive/{commit_hash}.zip" + +[tool.pytest.ini_options] +addopts = "--cov --cov-report=html --cov-report=xml" +log_cli = true +log_cli_level = "INFO" +testpaths = ["jumpstarter_driver_bt_peer"] + +[build-system] +requires = ["hatchling", "hatch-vcs", "hatch-pin-jumpstarter"] +build-backend = "hatchling.build" + +[tool.hatch.build.hooks.pin_jumpstarter] +name = "pin_jumpstarter" + +[dependency-groups] +dev = [ + "pytest-anyio>=0.0.0", + "pytest-cov>=6.0.0", + "pytest>=8.3.3", +] diff --git a/python/pyproject.toml b/python/pyproject.toml index f0f1dcef4..fd98bc597 100644 --- a/python/pyproject.toml +++ b/python/pyproject.toml @@ -10,6 +10,7 @@ jumpstarter-cli-driver = { workspace = true } jumpstarter-driver-adb = { workspace = true } jumpstarter-driver-androidemulator = { workspace = true } jumpstarter-driver-ble = { workspace = true } +jumpstarter-driver-bt-peer = { workspace = true } jumpstarter-driver-can = { workspace = true } jumpstarter-driver-composite = { workspace = true } jumpstarter-driver-doip = { workspace = true }