mirror of
https://github.com/firestar5683/StarPilot.git
synced 2026-09-08 17:13:45 +08:00
499 lines
20 KiB
Python
499 lines
20 KiB
Python
import os
|
|
import socket
|
|
import threading
|
|
import time
|
|
import uuid
|
|
|
|
from typing import Any
|
|
|
|
from jeepney import DBusAddress, MatchRule, new_error, new_method_call, new_method_return
|
|
from jeepney.fds import FileDescriptor
|
|
from jeepney.io.threading import DBusRouter, open_dbus_connection
|
|
from jeepney.low_level import HeaderFields, MessageType
|
|
from jeepney.wrappers import Properties
|
|
|
|
from openpilot.common.swaglog import cloudlog
|
|
from openpilot.starpilot.system.bluetooth.protocol import SPP_UUID, device_capabilities, show_pairing_device
|
|
|
|
|
|
BLUEZ = "org.bluez"
|
|
OBJECT_MANAGER = "org.freedesktop.DBus.ObjectManager"
|
|
ADAPTER_IFACE = "org.bluez.Adapter1"
|
|
DEVICE_IFACE = "org.bluez.Device1"
|
|
AGENT_MANAGER_IFACE = "org.bluez.AgentManager1"
|
|
AGENT_IFACE = "org.bluez.Agent1"
|
|
AGENT_PATH = "/link/firestar/starpilot/agent"
|
|
PROFILE_MANAGER_IFACE = "org.bluez.ProfileManager1"
|
|
PROFILE_IFACE = "org.bluez.Profile1"
|
|
OBDYSSEY_PROFILE_PATH = "/link/firestar/starpilot/obdyssey"
|
|
|
|
|
|
class BlueZError(RuntimeError):
|
|
"""A BlueZ D-Bus error with both its stable name and human detail."""
|
|
|
|
def __init__(self, error_name: str, detail: str, *, path: str = "", interface: str = "", member: str = ""):
|
|
self.error_name = error_name
|
|
self.detail = detail
|
|
self.path = path
|
|
self.interface = interface
|
|
self.member = member
|
|
message = error_name if not detail or detail == error_name else f"{error_name}: {detail}"
|
|
super().__init__(message)
|
|
|
|
@property
|
|
def method(self) -> str:
|
|
return f"{self.interface}.{self.member}" if self.interface and self.member else self.member
|
|
|
|
|
|
def _socket_from_dbus_fd(fd: Any) -> socket.socket:
|
|
"""Take ownership of a Profile1 NewConnection file descriptor."""
|
|
if isinstance(fd, FileDescriptor):
|
|
try:
|
|
return fd.to_socket()
|
|
except Exception:
|
|
try:
|
|
fd.close()
|
|
except Exception:
|
|
pass
|
|
raise
|
|
if isinstance(fd, int):
|
|
# Raw integer FDs are not owned by the message wrapper, so duplicate them
|
|
# before handing ownership to the socket object.
|
|
return socket.socket(fileno=os.dup(fd))
|
|
raise TypeError(f"Unsupported D-Bus file descriptor: {type(fd).__name__}")
|
|
|
|
|
|
def _dbus_text(value: Any) -> str:
|
|
if isinstance(value, bytes):
|
|
return value.decode("utf-8", errors="replace")
|
|
return str(value)
|
|
|
|
|
|
def unwrap_variant(value: Any) -> Any:
|
|
if isinstance(value, tuple) and len(value) == 2 and isinstance(value[0], str):
|
|
return unwrap_variant(value[1])
|
|
if isinstance(value, dict):
|
|
return {key: unwrap_variant(item) for key, item in value.items()}
|
|
if isinstance(value, list):
|
|
return [unwrap_variant(item) for item in value]
|
|
return value
|
|
|
|
|
|
class PairingAgent:
|
|
def __init__(self):
|
|
self._condition = threading.Condition()
|
|
self._prompt: dict[str, Any] | None = None
|
|
self._response: tuple[bool, str] | None = None
|
|
|
|
@property
|
|
def prompt(self) -> dict[str, Any] | None:
|
|
with self._condition:
|
|
return dict(self._prompt) if self._prompt is not None else None
|
|
|
|
def clear(self) -> None:
|
|
with self._condition:
|
|
self._prompt = None
|
|
self._response = None
|
|
self._condition.notify_all()
|
|
|
|
def display(self, kind: str, device_path: str, value: str) -> None:
|
|
with self._condition:
|
|
self._prompt = {"id": uuid.uuid4().hex, "kind": kind, "device_path": device_path, "value": value, "display_only": True}
|
|
|
|
def request(self, kind: str, device_path: str, value: str = "", timeout: float = 60.0) -> tuple[bool, str]:
|
|
prompt_id = uuid.uuid4().hex
|
|
with self._condition:
|
|
self._response = None
|
|
self._prompt = {"id": prompt_id, "kind": kind, "device_path": device_path, "value": value, "display_only": False}
|
|
deadline = time.monotonic() + timeout
|
|
while self._response is None:
|
|
remaining = deadline - time.monotonic()
|
|
if remaining <= 0:
|
|
self._prompt = None
|
|
return False, ""
|
|
self._condition.wait(remaining)
|
|
response = self._response
|
|
self._response = None
|
|
self._prompt = None
|
|
return response
|
|
|
|
def respond(self, prompt_id: str, accepted: bool, value: str = "") -> bool:
|
|
with self._condition:
|
|
if self._prompt is None or self._prompt.get("id") != prompt_id or self._prompt.get("display_only"):
|
|
return False
|
|
self._response = accepted, value
|
|
self._condition.notify_all()
|
|
return True
|
|
|
|
|
|
class _BlueZConnection:
|
|
def __init__(self, enable_fds: bool = False):
|
|
self.router = DBusRouter(open_dbus_connection(bus="SYSTEM", enable_fds=enable_fds))
|
|
|
|
def close(self) -> None:
|
|
self.router.close()
|
|
|
|
def _call(self, path: str, interface: str, member: str, signature: str | None = None, body: tuple = (), timeout: float = 15.0):
|
|
address = DBusAddress(path, bus_name=BLUEZ, interface=interface)
|
|
message = new_method_call(address, member, signature, body) if signature is not None else new_method_call(address, member)
|
|
reply = self.router.send_and_get_reply(message, timeout=timeout)
|
|
if reply.header.message_type == MessageType.error:
|
|
error_name = _dbus_text(reply.header.fields.get(HeaderFields.error_name, "org.bluez.Error.Failed"))
|
|
detail = _dbus_text(reply.body[0]) if reply.body else error_name
|
|
raise BlueZError(error_name, detail, path=path, interface=interface, member=member)
|
|
return reply.body
|
|
|
|
def managed_objects(self) -> dict[str, dict[str, dict[str, Any]]]:
|
|
body = self._call("/", OBJECT_MANAGER, "GetManagedObjects")
|
|
return unwrap_variant(body[0]) if body else {}
|
|
|
|
def adapter(self, objects: dict[str, Any] | None = None) -> tuple[str, dict[str, Any]]:
|
|
objects = self.managed_objects() if objects is None else objects
|
|
for path, interfaces in objects.items():
|
|
if ADAPTER_IFACE in interfaces:
|
|
return path, interfaces[ADAPTER_IFACE]
|
|
raise RuntimeError("Bluetooth adapter is not available")
|
|
|
|
def devices(self, objects: dict[str, Any] | None = None) -> list[dict[str, Any]]:
|
|
objects = self.managed_objects() if objects is None else objects
|
|
devices = []
|
|
for path, interfaces in objects.items():
|
|
if DEVICE_IFACE not in interfaces:
|
|
continue
|
|
props = interfaces[DEVICE_IFACE]
|
|
uuids = [str(value).lower() for value in props.get("UUIDs", [])]
|
|
audio, controller, serial = device_capabilities(uuids, int(props.get("Class", 0)), str(props.get("Icon", "")))
|
|
device = {
|
|
"path": path,
|
|
"address": str(props.get("Address", "")),
|
|
"name": str(props.get("Alias") or props.get("Name") or props.get("Address") or "Unknown device"),
|
|
"paired": bool(props.get("Paired", False)),
|
|
"trusted": bool(props.get("Trusted", False)),
|
|
"connected": bool(props.get("Connected", False)),
|
|
"blocked": bool(props.get("Blocked", False)),
|
|
"rssi": int(props["RSSI"]) if "RSSI" in props else None,
|
|
"uuids": uuids,
|
|
"audio": audio,
|
|
"controller": controller,
|
|
"serial": serial,
|
|
}
|
|
if show_pairing_device(device["address"], device["name"], device["paired"], device["trusted"], device["connected"],
|
|
device["blocked"], audio, controller, serial):
|
|
devices.append(device)
|
|
return sorted(devices, key=lambda device: (not device["connected"], not device["paired"], -(device["rssi"] or -127), device["name"].lower()))
|
|
|
|
def device_for_address(self, address: str) -> dict[str, Any]:
|
|
normalized = address.upper()
|
|
for device in self.devices():
|
|
if device["address"].upper() == normalized:
|
|
return device
|
|
raise RuntimeError(f"Bluetooth device {address} was not found")
|
|
|
|
|
|
class BlueZClient(_BlueZConnection):
|
|
def __init__(self):
|
|
super().__init__(enable_fds=False)
|
|
self.agent = PairingAgent()
|
|
self._agent_filter = self.router.filter(MatchRule(type="method_call", interface=AGENT_IFACE, path=AGENT_PATH), bufsize=20)
|
|
self._agent_queue = self._agent_filter.__enter__()
|
|
self._agent_thread = threading.Thread(target=self._agent_loop, daemon=True)
|
|
self._agent_thread.start()
|
|
self._register_agent()
|
|
|
|
def close(self) -> None:
|
|
try:
|
|
self._call("/org/bluez", AGENT_MANAGER_IFACE, "UnregisterAgent", "o", (AGENT_PATH,))
|
|
except Exception:
|
|
pass
|
|
self._agent_filter.__exit__(None, None, None)
|
|
super().close()
|
|
|
|
def _register_agent(self) -> None:
|
|
self._call("/org/bluez", AGENT_MANAGER_IFACE, "RegisterAgent", "os", (AGENT_PATH, "KeyboardDisplay"))
|
|
self._call("/org/bluez", AGENT_MANAGER_IFACE, "RequestDefaultAgent", "o", (AGENT_PATH,))
|
|
|
|
def _agent_loop(self) -> None:
|
|
while True:
|
|
message = self._agent_queue.get()
|
|
member = message.header.fields.get(HeaderFields.member, "")
|
|
try:
|
|
response_signature = None
|
|
response_body: tuple = ()
|
|
device_path = str(message.body[0]) if message.body else ""
|
|
if member == "Release":
|
|
self.agent.clear()
|
|
elif member == "RequestPinCode":
|
|
accepted, value = self.agent.request("pin", device_path)
|
|
if not accepted:
|
|
raise PermissionError
|
|
response_signature, response_body = "s", (value,)
|
|
elif member == "DisplayPinCode":
|
|
self.agent.display("display_pin", device_path, str(message.body[1]))
|
|
elif member == "RequestPasskey":
|
|
accepted, value = self.agent.request("passkey", device_path)
|
|
if not accepted:
|
|
raise PermissionError
|
|
response_signature, response_body = "u", (int(value),)
|
|
elif member == "DisplayPasskey":
|
|
self.agent.display("display_passkey", device_path, f"{int(message.body[1]):06d}")
|
|
elif member == "RequestConfirmation":
|
|
accepted, _ = self.agent.request("confirmation", device_path, f"{int(message.body[1]):06d}")
|
|
if not accepted:
|
|
raise PermissionError
|
|
elif member in ("RequestAuthorization", "AuthorizeService"):
|
|
accepted, _ = self.agent.request("authorization", device_path)
|
|
if not accepted:
|
|
raise PermissionError
|
|
elif member == "Cancel":
|
|
self.agent.clear()
|
|
else:
|
|
raise RuntimeError(f"Unsupported pairing request: {member}")
|
|
self.router.send(new_method_return(message, response_signature, response_body))
|
|
except PermissionError:
|
|
self.router.send(new_error(message, "org.bluez.Error.Rejected", "s", ("Pairing rejected",)))
|
|
except Exception as error:
|
|
self.router.send(new_error(message, "org.bluez.Error.Canceled", "s", (str(error),)))
|
|
|
|
def status(self) -> dict[str, Any]:
|
|
objects = self.managed_objects()
|
|
_, adapter = self.adapter(objects)
|
|
return {
|
|
"powered": bool(adapter.get("Powered", False)),
|
|
"discovering": bool(adapter.get("Discovering", False)),
|
|
"devices": self.devices(objects),
|
|
"prompt": self.agent.prompt,
|
|
}
|
|
|
|
def set_powered(self, powered: bool) -> None:
|
|
path, _ = self.adapter()
|
|
address = DBusAddress(path, bus_name=BLUEZ, interface=ADAPTER_IFACE)
|
|
reply = self.router.send_and_get_reply(Properties(address).set("Powered", "b", powered), timeout=10.0)
|
|
if reply.header.message_type == MessageType.error:
|
|
raise RuntimeError(str(reply.body[0] if reply.body else "Unable to change Bluetooth power"))
|
|
|
|
def start_discovery(self) -> None:
|
|
path, _ = self.adapter()
|
|
self._call(path, ADAPTER_IFACE, "StartDiscovery")
|
|
|
|
def stop_discovery(self) -> None:
|
|
path, props = self.adapter()
|
|
if props.get("Discovering", False):
|
|
self._call(path, ADAPTER_IFACE, "StopDiscovery")
|
|
|
|
def set_device_property(self, address: str, name: str, signature: str, value: Any) -> None:
|
|
device = self.device_for_address(address)
|
|
dbus_address = DBusAddress(device["path"], bus_name=BLUEZ, interface=DEVICE_IFACE)
|
|
reply = self.router.send_and_get_reply(Properties(dbus_address).set(name, signature, value), timeout=10.0)
|
|
if reply.header.message_type == MessageType.error:
|
|
raise RuntimeError(str(reply.body[0] if reply.body else f"Unable to set {name}"))
|
|
|
|
def pair(self, address: str) -> None:
|
|
device = self.device_for_address(address)
|
|
self._call(device["path"], DEVICE_IFACE, "Pair", timeout=90.0)
|
|
self.set_device_property(address, "Trusted", "b", True)
|
|
self.agent.clear()
|
|
|
|
def connect(self, address: str) -> None:
|
|
device = self.device_for_address(address)
|
|
self._call(device["path"], DEVICE_IFACE, "Connect", timeout=30.0)
|
|
|
|
def disconnect(self, address: str) -> None:
|
|
device = self.device_for_address(address)
|
|
self._call(device["path"], DEVICE_IFACE, "Disconnect", timeout=15.0)
|
|
|
|
def remove(self, address: str) -> None:
|
|
adapter_path, _ = self.adapter()
|
|
device = self.device_for_address(address)
|
|
self._call(adapter_path, ADAPTER_IFACE, "RemoveDevice", "o", (device["path"],))
|
|
|
|
|
|
class BlueZProfileClient(_BlueZConnection):
|
|
def __init__(self, profile_path: str = OBDYSSEY_PROFILE_PATH, profile_uuid: str = SPP_UUID):
|
|
super().__init__(enable_fds=True)
|
|
self.profile_path = profile_path
|
|
self.profile_uuid = profile_uuid
|
|
self._lock = threading.Lock()
|
|
self._new_conn_event = threading.Event()
|
|
self._closed = threading.Event()
|
|
self._last_conn_path = ""
|
|
self._last_conn_sock: socket.socket | None = None
|
|
self._active_conn_path = ""
|
|
self._active_conn_sock: socket.socket | None = None
|
|
self._profile_filter = None
|
|
self._profile_queue = None
|
|
self._profile_registered = False
|
|
self._close_lock = threading.Lock()
|
|
self._profile_thread: threading.Thread | None = None
|
|
|
|
try:
|
|
self._profile_filter = self.router.filter(MatchRule(type="method_call", interface=PROFILE_IFACE, path=self.profile_path), bufsize=20)
|
|
self._profile_queue = self._profile_filter.__enter__()
|
|
self._profile_thread = threading.Thread(target=self._profile_loop, daemon=True)
|
|
self._profile_thread.start()
|
|
self._register_profile()
|
|
except Exception:
|
|
self.close()
|
|
raise
|
|
|
|
def close(self) -> None:
|
|
with self._close_lock:
|
|
if self._closed.is_set():
|
|
return
|
|
if self._profile_registered:
|
|
self._unregister_profile()
|
|
self._profile_registered = False
|
|
self._closed.set()
|
|
profile_queue = self._profile_queue
|
|
if profile_queue is not None:
|
|
try:
|
|
profile_queue.put(None, timeout=1.0)
|
|
except Exception:
|
|
pass
|
|
try:
|
|
if self._profile_filter is not None:
|
|
self._profile_filter.__exit__(None, None, None)
|
|
self._profile_filter = None
|
|
self._profile_queue = None
|
|
finally:
|
|
if self._profile_thread is not None and self._profile_thread is not threading.current_thread():
|
|
self._profile_thread.join(timeout=1.0)
|
|
self._close_connection_sockets()
|
|
super().close()
|
|
|
|
def _register_profile(self) -> None:
|
|
options = {"Role": ("s", "client"), "Name": ("s", "OBDyssey")}
|
|
# Keep cleanup armed before the remote call. A timeout can happen after
|
|
# BlueZ has registered the profile but before its reply reaches us.
|
|
self._profile_registered = True
|
|
self._call("/org/bluez", PROFILE_MANAGER_IFACE, "RegisterProfile", "osa{sv}", (self.profile_path, self.profile_uuid, options))
|
|
|
|
def _unregister_profile(self) -> None:
|
|
try:
|
|
self._call("/org/bluez", PROFILE_MANAGER_IFACE, "UnregisterProfile", "o", (self.profile_path,))
|
|
except Exception:
|
|
pass
|
|
|
|
def _profile_loop(self) -> None:
|
|
profile_queue = self._profile_queue
|
|
while profile_queue is not None:
|
|
message = profile_queue.get()
|
|
if message is None:
|
|
break
|
|
member = message.header.fields.get(HeaderFields.member, "")
|
|
if self._closed.is_set():
|
|
if member == "NewConnection" and len(message.body) > 1:
|
|
try:
|
|
_socket_from_dbus_fd(message.body[1]).close()
|
|
except Exception:
|
|
pass
|
|
try:
|
|
self.router.send(new_error(message, "org.bluez.Error.Canceled", "s", ("Profile is closed",)))
|
|
except Exception:
|
|
pass
|
|
continue
|
|
sock: socket.socket | None = None
|
|
try:
|
|
if member == "NewConnection":
|
|
device_path = str(message.body[0]) if len(message.body) > 0 else ""
|
|
fd = message.body[1] if len(message.body) > 1 else None
|
|
if fd is None:
|
|
raise RuntimeError("Profile1 NewConnection did not include a file descriptor")
|
|
sock = _socket_from_dbus_fd(fd)
|
|
self._store_connection_socket(device_path, sock)
|
|
sock = None
|
|
self.router.send(new_method_return(message))
|
|
elif member == "RequestDisconnection":
|
|
device_path = str(message.body[0]) if message.body else ""
|
|
self._close_connection_sockets(device_path)
|
|
self.router.send(new_method_return(message))
|
|
elif member == "Release":
|
|
self._close_connection_sockets()
|
|
self.router.send(new_method_return(message))
|
|
else:
|
|
self.router.send(new_method_return(message))
|
|
except Exception as error:
|
|
if sock is not None:
|
|
try:
|
|
sock.close()
|
|
except Exception:
|
|
pass
|
|
elif member == "NewConnection":
|
|
self._close_connection_sockets(device_path)
|
|
self.router.send(new_error(message, "org.bluez.Error.Failed", "s", (str(error),)))
|
|
|
|
def _store_connection_socket(self, device_path: str, sock: socket.socket) -> None:
|
|
old_socks: list[socket.socket] = []
|
|
with self._lock:
|
|
if self._last_conn_sock is not None:
|
|
old_socks.append(self._last_conn_sock)
|
|
if self._active_conn_sock is not None:
|
|
old_socks.append(self._active_conn_sock)
|
|
self._last_conn_path = device_path
|
|
self._last_conn_sock = sock
|
|
self._active_conn_path = ""
|
|
self._active_conn_sock = None
|
|
self._new_conn_event.set()
|
|
for old_sock in old_socks:
|
|
try:
|
|
old_sock.close()
|
|
except Exception:
|
|
pass
|
|
|
|
def _close_connection_sockets(self, device_path: str | None = None) -> None:
|
|
sockets: list[socket.socket] = []
|
|
with self._lock:
|
|
if self._last_conn_sock is not None and (device_path is None or self._last_conn_path == device_path):
|
|
sockets.append(self._last_conn_sock)
|
|
self._last_conn_sock = None
|
|
self._last_conn_path = ""
|
|
if self._active_conn_sock is not None and (device_path is None or self._active_conn_path == device_path):
|
|
if self._active_conn_sock not in sockets:
|
|
sockets.append(self._active_conn_sock)
|
|
self._active_conn_sock = None
|
|
self._active_conn_path = ""
|
|
self._new_conn_event.set()
|
|
for sock in sockets:
|
|
try:
|
|
sock.close()
|
|
except Exception:
|
|
pass
|
|
|
|
def connect_profile(self, address: str, timeout: float = 30.0) -> socket.socket:
|
|
device = self.device_for_address(address)
|
|
device_path = device["path"]
|
|
self._close_connection_sockets(device_path)
|
|
with self._lock:
|
|
self._new_conn_event.clear()
|
|
|
|
try:
|
|
self._call(device_path, DEVICE_IFACE, "ConnectProfile", "s", (self.profile_uuid,), timeout=timeout)
|
|
|
|
deadline = time.monotonic() + timeout
|
|
while time.monotonic() < deadline and not self._closed.is_set():
|
|
remaining = deadline - time.monotonic()
|
|
if self._new_conn_event.wait(timeout=min(0.2, max(0.01, remaining))):
|
|
with self._lock:
|
|
if self._last_conn_path == device_path and self._last_conn_sock is not None:
|
|
sock = self._last_conn_sock
|
|
self._last_conn_sock = None
|
|
self._last_conn_path = ""
|
|
self._active_conn_path = device_path
|
|
self._active_conn_sock = sock
|
|
return sock
|
|
self._new_conn_event.clear()
|
|
raise TimeoutError(f"Timed out waiting for SPP connection to {address}")
|
|
except Exception as err:
|
|
method = err.method if isinstance(err, BlueZError) else f"{DEVICE_IFACE}.ConnectProfile"
|
|
error_name = err.error_name if isinstance(err, BlueZError) else type(err).__name__
|
|
detail = err.detail if isinstance(err, BlueZError) else str(err)
|
|
cloudlog.error(
|
|
f"OBDyssey SPP connection failed method={method} address={address} device_path={device_path} "
|
|
+ f"device_uuids={device.get('uuids', [])} error_name={error_name} detail={detail}"
|
|
)
|
|
self._close_connection_sockets()
|
|
raise
|
|
|
|
def disconnect_profile(self, address: str, timeout: float = 15.0) -> None:
|
|
device = self.device_for_address(address)
|
|
self._call(device["path"], DEVICE_IFACE, "DisconnectProfile", "s", (self.profile_uuid,), timeout=timeout)
|