253 lines
9.8 KiB
Python
253 lines
9.8 KiB
Python
"""BlueZ GATT adapter. Failure is isolated from the display service."""
|
|
import asyncio
|
|
import logging
|
|
import threading
|
|
|
|
from dbus_next import Variant, DBusError, BusType, PropertyAccess
|
|
from dbus_next.aio import MessageBus
|
|
from dbus_next.service import ServiceInterface, method, dbus_property
|
|
|
|
from .protocol import fragments
|
|
|
|
log = logging.getLogger(__name__)
|
|
ROOT = '/org/qimiaoscreen/mobile'
|
|
SERVICE = ROOT + '/service0'
|
|
AD = ROOT + '/advertisement0'
|
|
SERVICE_UUID = '9f57a001-6c31-4c58-bc22-1f728b641001'
|
|
RX_UUID = '9f57a002-6c31-4c58-bc22-1f728b641001'
|
|
TX_UUID = '9f57a003-6c31-4c58-bc22-1f728b641001'
|
|
|
|
|
|
class ObjectManager(ServiceInterface):
|
|
def __init__(self, objects):
|
|
super().__init__('org.freedesktop.DBus.ObjectManager')
|
|
self.objects = objects
|
|
|
|
@method()
|
|
def GetManagedObjects(self) -> 'a{oa{sa{sv}}}':
|
|
return {path: {obj.name: obj.properties()} for path, obj in self.objects.items()}
|
|
|
|
|
|
class GattService(ServiceInterface):
|
|
def __init__(self):
|
|
super().__init__('org.bluez.GattService1')
|
|
|
|
@dbus_property(access=PropertyAccess.READ)
|
|
def UUID(self) -> 's':
|
|
return SERVICE_UUID
|
|
|
|
@dbus_property(access=PropertyAccess.READ)
|
|
def Primary(self) -> 'b':
|
|
return True
|
|
|
|
def properties(self):
|
|
return dict(UUID=Variant('s', self.UUID), Primary=Variant('b', True))
|
|
|
|
|
|
class Characteristic(ServiceInterface):
|
|
def __init__(self, uuid, runtime):
|
|
super().__init__('org.bluez.GattCharacteristic1')
|
|
self.uuid, self.runtime = uuid, runtime
|
|
self.notifying = False
|
|
|
|
@dbus_property(access=PropertyAccess.READ)
|
|
def UUID(self) -> 's':
|
|
return self.uuid
|
|
|
|
@dbus_property(access=PropertyAccess.READ)
|
|
def Service(self) -> 'o':
|
|
return SERVICE
|
|
|
|
@dbus_property(access=PropertyAccess.READ)
|
|
def Flags(self) -> 'as':
|
|
return ['write'] if self.uuid == RX_UUID else ['notify']
|
|
|
|
@dbus_property(access=PropertyAccess.READ)
|
|
def Value(self) -> 'ay':
|
|
return b''
|
|
|
|
@dbus_property(access=PropertyAccess.READ)
|
|
def Notifying(self) -> 'b':
|
|
return self.notifying
|
|
|
|
@method()
|
|
async def WriteValue(self, value: 'ay', options: 'a{sv}'):
|
|
if self.uuid != RX_UUID:
|
|
raise DBusError('org.bluez.Error.NotPermitted', 'Not permitted')
|
|
await self.runtime.write(value, options)
|
|
|
|
@method()
|
|
def StartNotify(self):
|
|
if self.uuid != TX_UUID:
|
|
raise DBusError('org.bluez.Error.NotSupported', 'Not supported')
|
|
self.notifying = True
|
|
|
|
@method()
|
|
def StopNotify(self):
|
|
self.notifying = False
|
|
|
|
def properties(self):
|
|
return dict(UUID=Variant('s', self.UUID), Service=Variant('o', SERVICE), Flags=Variant('as', self.Flags))
|
|
|
|
|
|
class Advertisement(ServiceInterface):
|
|
def __init__(self, short_id):
|
|
super().__init__('org.bluez.LEAdvertisement1')
|
|
self.local_name = 'QMS-' + short_id
|
|
|
|
@dbus_property(access=PropertyAccess.READ)
|
|
def Type(self) -> 's':
|
|
return 'peripheral'
|
|
|
|
@dbus_property(access=PropertyAccess.READ)
|
|
def ServiceUUIDs(self) -> 'as':
|
|
return [SERVICE_UUID]
|
|
|
|
@dbus_property(access=PropertyAccess.READ)
|
|
def LocalName(self) -> 's':
|
|
return self.local_name
|
|
|
|
@method()
|
|
def Release(self):
|
|
pass
|
|
|
|
|
|
class BluezRuntime:
|
|
def __init__(self, sessions, enabled=False):
|
|
self.sessions = sessions
|
|
self.enabled = enabled
|
|
self.stop_event = threading.Event()
|
|
self.thread = None
|
|
self.state = dict(available=False, reason='disabled')
|
|
self.bus = self.manager = self.ad_manager = None
|
|
self.advertising = False
|
|
self.tx = None
|
|
self.message_id = 0
|
|
self.write_lock = None
|
|
self.counters = dict(rx_fragments=0, tx_fragments=0, rx_bytes=0, tx_bytes=0, last_response_kind=0, mtu=23)
|
|
|
|
def status(self):
|
|
return {**self.state, **self.sessions.status(), 'transport': dict(self.counters)}
|
|
|
|
def start(self):
|
|
if not self.enabled or self.thread:
|
|
return
|
|
self.thread = threading.Thread(target=lambda: asyncio.run(self._run()), name='mobile-ble', daemon=True)
|
|
self.thread.start()
|
|
|
|
def close(self):
|
|
self.stop_event.set()
|
|
if self.thread:
|
|
self.thread.join(timeout=8)
|
|
|
|
async def _proxy(self, path):
|
|
introspection = await asyncio.wait_for(self.bus.introspect('org.bluez', path), 5)
|
|
return self.bus.get_proxy_object('org.bluez', path, introspection)
|
|
|
|
async def _disconnect(self, owner):
|
|
was_owner = self.sessions.owner == owner
|
|
try:
|
|
proxy = await self._proxy(owner)
|
|
await asyncio.wait_for(proxy.get_interface('org.bluez.Device1').call_disconnect(), 3)
|
|
except Exception:
|
|
pass
|
|
self.sessions.disconnect(owner)
|
|
if was_owner:
|
|
self.message_id = 0
|
|
|
|
async def write(self, value, options):
|
|
self.counters['rx_fragments'] += 1
|
|
self.counters['rx_bytes'] += len(value)
|
|
owner = options.get('device')
|
|
if not owner or options.get('offset', Variant('q', 0)).value != 0:
|
|
raise DBusError('org.bluez.Error.NotAuthorized', 'Invalid request')
|
|
owner = owner.value
|
|
if not self.sessions.connect(owner):
|
|
await self._disconnect(owner)
|
|
raise DBusError('org.bluez.Error.NotAuthorized', 'Busy')
|
|
async with self.write_lock:
|
|
try:
|
|
if not self.tx.notifying:
|
|
raise ValueError()
|
|
response = await asyncio.to_thread(self.sessions.accept, owner, bytes(value))
|
|
if response is None:
|
|
return
|
|
mtu = min(self.sessions.peer_mtu, max(23, options.get('mtu', Variant('q', 23)).value))
|
|
self.counters['mtu'] = mtu
|
|
self.counters['last_response_kind'] = response[0]
|
|
for fragment in fragments(response[0], self.message_id, response[1], mtu):
|
|
self.tx.emit_properties_changed({'Value': fragment})
|
|
self.counters['tx_fragments'] += 1
|
|
self.counters['tx_bytes'] += len(fragment)
|
|
await asyncio.sleep(0.002)
|
|
self.message_id += 1
|
|
except Exception:
|
|
await self._disconnect(owner)
|
|
raise DBusError('org.bluez.Error.NotAuthorized', 'Session rejected') from None
|
|
|
|
async def _run(self):
|
|
while not self.stop_event.is_set():
|
|
try:
|
|
await self._serve()
|
|
except Exception:
|
|
self.state = dict(available=False, reason='bluetooth_unavailable')
|
|
log.warning('Mobile BLE unavailable; display service remains active')
|
|
finally:
|
|
if self.sessions.owner is not None:
|
|
self.sessions.disconnect(self.sessions.owner)
|
|
self.message_id = 0
|
|
self.advertising = False
|
|
if self.bus:
|
|
self.bus.disconnect()
|
|
self.bus = None
|
|
for _ in range(5):
|
|
if self.stop_event.is_set():
|
|
return
|
|
await asyncio.sleep(1)
|
|
|
|
async def _serve(self):
|
|
self.bus = await asyncio.wait_for(MessageBus(bus_type=BusType.SYSTEM).connect(), 5)
|
|
root_proxy = await self._proxy('/')
|
|
self.manager = root_proxy.get_interface('org.freedesktop.DBus.ObjectManager')
|
|
objects = await asyncio.wait_for(self.manager.call_get_managed_objects(), 5)
|
|
adapters = [path for path, interfaces in objects.items() if 'org.bluez.GattManager1' in interfaces and 'org.bluez.LEAdvertisingManager1' in interfaces]
|
|
if not adapters:
|
|
raise RuntimeError('No GATT peripheral adapter')
|
|
adapter = await self._proxy(sorted(adapters)[0])
|
|
properties = adapter.get_interface('org.freedesktop.DBus.Properties')
|
|
await properties.call_set('org.bluez.Adapter1', 'Powered', Variant('b', True))
|
|
await properties.call_set('org.bluez.Adapter1', 'Pairable', Variant('b', False))
|
|
self.ad_manager = adapter.get_interface('org.bluez.LEAdvertisingManager1')
|
|
self.tx = Characteristic(TX_UUID, self)
|
|
exports = {SERVICE: GattService(), SERVICE + '/rx': Characteristic(RX_UUID, self), SERVICE + '/tx': self.tx}
|
|
self.bus.export(ROOT, ObjectManager(exports))
|
|
for path, interface in exports.items():
|
|
self.bus.export(path, interface)
|
|
self.bus.export(AD, Advertisement(self.sessions.identity()['short_id']))
|
|
await asyncio.wait_for(adapter.get_interface('org.bluez.GattManager1').call_register_application(ROOT, {}), 5)
|
|
self.write_lock = asyncio.Lock()
|
|
self.state = dict(available=True, reason=None)
|
|
while not self.stop_event.is_set():
|
|
objects = await asyncio.wait_for(self.manager.call_get_managed_objects(), 5)
|
|
connected = [path for path, interfaces in objects.items()
|
|
if interfaces.get('org.bluez.Device1', {}).get('Connected', Variant('b', False)).value]
|
|
owner = self.sessions.owner
|
|
if owner and owner not in connected:
|
|
self.sessions.disconnect(owner)
|
|
self.message_id = 0
|
|
for path in connected:
|
|
if not self.sessions.connect(path):
|
|
await self._disconnect(path)
|
|
if self.sessions.expired():
|
|
await self._disconnect(self.sessions.owner)
|
|
should_advertise = self.sessions.owner is None
|
|
if should_advertise != self.advertising:
|
|
if should_advertise:
|
|
await asyncio.wait_for(self.ad_manager.call_register_advertisement(AD, {}), 5)
|
|
else:
|
|
await asyncio.wait_for(self.ad_manager.call_unregister_advertisement(AD), 5)
|
|
self.advertising = should_advertise
|
|
await asyncio.sleep(0.5)
|
|
if self.sessions.owner:
|
|
await self._disconnect(self.sessions.owner)
|