From 93addc3494c8eca11ba7fae1423c89461f56e9d8 Mon Sep 17 00:00:00 2001 From: Dennis Moschina <45356478+DennisMoschina@users.noreply.github.com> Date: Tue, 21 Jul 2026 15:37:05 +0200 Subject: [PATCH 01/11] feat!: await BLE notification subscriptions BREAKING CHANGE: BleGattManager.subscribe and SensorHandler.subscribeToSensorData now return Future> so callers can await notification readiness before issuing dependent writes. --- lib/src/managers/ble_gatt_manager.dart | 5 +- lib/src/managers/ble_manager.dart | 31 +- lib/src/managers/esense_sensor_handler.dart | 24 +- .../managers/open_earable_sensor_manager.dart | 17 +- .../managers/open_ring_sensor_handler.dart | 98 +++-- lib/src/managers/sensor_handler.dart | 5 +- lib/src/managers/v2_sensor_handler.dart | 16 +- .../open_ring/open_ring_sensor.dart | 48 ++- lib/src/models/devices/cosinuss_one.dart | 340 ++++++++++-------- lib/src/models/devices/esense_factory.dart | 69 ++-- .../models/devices/open_earable_factory.dart | 56 +-- lib/src/models/devices/open_earable_v1.dart | 78 ++-- lib/src/models/devices/open_earable_v2.dart | 115 +++--- lib/src/models/devices/open_ring.dart | 61 ++-- lib/src/models/devices/polar_factory.dart | 128 ++++--- .../v2_sensor_scheme_reader.dart | 2 +- 16 files changed, 647 insertions(+), 446 deletions(-) diff --git a/lib/src/managers/ble_gatt_manager.dart b/lib/src/managers/ble_gatt_manager.dart index d4246195..e93ef47d 100644 --- a/lib/src/managers/ble_gatt_manager.dart +++ b/lib/src/managers/ble_gatt_manager.dart @@ -24,7 +24,10 @@ abstract class BleGattManager { }); /// Subscribes to a specific characteristic of the connected device. - Stream> subscribe({ + /// + /// The returned future completes only after the underlying GATT + /// notification subscription has been enabled. + Future>> subscribe({ required String deviceId, required String serviceId, required String characteristicId, diff --git a/lib/src/managers/ble_manager.dart b/lib/src/managers/ble_manager.dart index ef71e7ec..e787e558 100644 --- a/lib/src/managers/ble_manager.dart +++ b/lib/src/managers/ble_manager.dart @@ -15,6 +15,7 @@ class BleManager extends BleGattManager { int get mtu => _mtu; final Map>> _streamControllers = {}; + final Map> _subscriptionSetups = {}; /// A stream of discovered devices during scanning. StreamController? _scanStreamController; @@ -49,6 +50,7 @@ class BleManager extends BleGattManager { for (final key in keys) { logger.d("Closing stream for $key due to device disconnection"); _streamControllers.remove(key)?.close(); + _subscriptionSetups.remove(key); } } @@ -312,11 +314,11 @@ class BleManager extends BleGattManager { /// Subscribes to a specific characteristic of the connected Earable device. @override - Stream> subscribe({ + Future>> subscribe({ required String deviceId, required String serviceId, required String characteristicId, - }) { + }) async { logger.d( "Subscribing to $deviceId, service $serviceId, characteristic $characteristicId", ); @@ -327,24 +329,33 @@ class BleManager extends BleGattManager { StreamController>? streamController = _streamControllers[streamIdentifier]; streamController ??= StreamController>.broadcast(); - if (!_streamControllers.containsKey(streamIdentifier)) { - UniversalBle.subscribeNotifications( - deviceId, - serviceId, - characteristicId, - ); - _streamControllers[streamIdentifier] = streamController; + _streamControllers[streamIdentifier] = streamController; + + _subscriptionSetups[streamIdentifier] ??= + UniversalBle.subscribeNotifications( + deviceId, + serviceId, + characteristicId, + ); + + try { + await _subscriptionSetups[streamIdentifier]; + } catch (_) { + _subscriptionSetups.remove(streamIdentifier); + _streamControllers.remove(streamIdentifier); + await streamController.close(); + rethrow; } streamController.onCancel = () { if (_streamControllers.containsKey(streamIdentifier)) { _streamControllers.remove(streamIdentifier)?.close(); + _subscriptionSetups.remove(streamIdentifier); UniversalBle.unsubscribe( deviceId, serviceId, characteristicId, ); - _streamControllers.remove(streamIdentifier); } }; diff --git a/lib/src/managers/esense_sensor_handler.dart b/lib/src/managers/esense_sensor_handler.dart index f378ba00..6600edaf 100644 --- a/lib/src/managers/esense_sensor_handler.dart +++ b/lib/src/managers/esense_sensor_handler.dart @@ -56,7 +56,9 @@ class EsenseSensorHandler extends SensorHandler { } @override - Stream> subscribeToSensorData(int sensorId) { + Future>> subscribeToSensorData( + int sensorId, + ) async { if (!_bleGattManager.isConnected(_discoveredDevice.id)) { throw Exception("Can't subscribe to sensor data. Earable not connected"); } @@ -68,13 +70,13 @@ class EsenseSensorHandler extends SensorHandler { final streamController = StreamController>(); - final subscription = _bleGattManager - .subscribe( - deviceId: _discoveredDevice.id, - serviceId: esenseServiceUuid, - characteristicId: esenseSensorDataCharacteristicUuid, - ) - .listen( + final dataStream = await _bleGattManager.subscribe( + deviceId: _discoveredDevice.id, + serviceId: esenseServiceUuid, + characteristicId: esenseSensorDataCharacteristicUuid, + ); + + final subscription = dataStream.listen( (data) async { if (data.isEmpty) return; @@ -310,7 +312,9 @@ class EsenseSensorHandler extends SensorHandler { ); } - logger.t("Loaded IMU ranges: Accel=$_cachedAccelRange, Gyro=$_cachedGyroRange"); + logger.t( + "Loaded IMU ranges: Accel=$_cachedAccelRange, Gyro=$_cachedGyroRange", + ); return (_cachedAccelRange!, _cachedGyroRange!); } @@ -347,7 +351,7 @@ class EsenseSensorHandler extends SensorHandler { // Make *new* mutable, dynamic-typed inner maps final accel = Map.from(result['Accelerometer'] as Map); - final gyro = Map.from(result['Gyroscope'] as Map); + final gyro = Map.from(result['Gyroscope'] as Map); // Accelerometer to g for (final key in const ['x', 'y', 'z']) { diff --git a/lib/src/managers/open_earable_sensor_manager.dart b/lib/src/managers/open_earable_sensor_manager.dart index ef859d4e..438d0703 100644 --- a/lib/src/managers/open_earable_sensor_manager.dart +++ b/lib/src/managers/open_earable_sensor_manager.dart @@ -31,7 +31,8 @@ class OpenEarableSensorHandler extends SensorHandler { SensorSchemeReader? sensorSchemeParser, SensorValueParser? sensorValueParser, }) : _bleManager = bleManager, - _sensorSchemeParser = sensorSchemeParser ?? EdgeMlSensorSchemeReader(bleManager, deviceId), + _sensorSchemeParser = sensorSchemeParser ?? + EdgeMlSensorSchemeReader(bleManager, deviceId), _sensorValueParser = sensorValueParser ?? EdgeMlSensorValueParser() { _readSensorScheme(); } @@ -63,20 +64,22 @@ class OpenEarableSensorHandler extends SensorHandler { /// - 1: Barometer data /// Returns a [Stream] of sensor data as a [Map] of sensor values. @override - Stream> subscribeToSensorData(int sensorId) { + Future>> subscribeToSensorData( + int sensorId, + ) async { if (!_bleManager.isConnected(deviceId)) { Exception("Can't subscribe to sensor data. Earable not connected"); } StreamController> streamController = StreamController(); int lastTimestamp = 0; - final subscription = _bleManager - .subscribe( + final dataStream = await _bleManager.subscribe( deviceId: deviceId, serviceId: sensorServiceUuid, characteristicId: sensorDataCharacteristicUuid, - ) - .listen( + ); + + final subscription = dataStream.listen( (data) async { if (data.isNotEmpty && data[0] == sensorId) { List> parsedDataList = await _parseData(data); @@ -144,7 +147,7 @@ class OpenEarableSensorHandler extends SensorHandler { /// Parses raw sensor data bytes into a [Map] of sensor values. Future>> _parseData(List data) async { ByteData byteData = ByteData.sublistView(Uint8List.fromList(data)); - + return _sensorValueParser.parse(byteData, _sensorSchemes!); } diff --git a/lib/src/managers/open_ring_sensor_handler.dart b/lib/src/managers/open_ring_sensor_handler.dart index 6081ec92..6a461064 100644 --- a/lib/src/managers/open_ring_sensor_handler.dart +++ b/lib/src/managers/open_ring_sensor_handler.dart @@ -32,7 +32,7 @@ class OpenRingSensorHandler extends SensorHandler { OpenRingGatt.cmdPPGQ2, }; - Stream>? _sensorDataStream; + Future>>? _sensorDataStream; Future _commandQueue = Future.value(); final _OpenRingDesiredState _desiredState = _OpenRingDesiredState(); int _applyVersion = 0; @@ -53,14 +53,17 @@ class OpenRingSensorHandler extends SensorHandler { _sensorValueParser = sensorValueParser; @override - Stream> subscribeToSensorData(int sensorId) { + Future>> subscribeToSensorData( + int sensorId, + ) async { if (!_bleManager.isConnected(_discoveredDevice.id)) { throw Exception("Can't subscribe to sensor data. Earable not connected"); } _sensorDataStream ??= _createSensorDataStream(); - return _sensorDataStream!.where((data) { + final sensorDataStream = await _sensorDataStream!; + return sensorDataStream.where((data) { final dynamic cmd = data['cmd']; return cmd is int && cmd == sensorId; }); @@ -395,9 +398,8 @@ class OpenRingSensorHandler extends SensorHandler { _onInitialStreamingDetected?.call(); } - Stream> _createSensorDataStream() { + Future>> _createSensorDataStream() async { late final StreamController> streamController; - // ignore: cancel_subscriptions StreamSubscription>? bleSubscription; final scheduler = _OpenRingPacedScheduler( @@ -461,56 +463,52 @@ class OpenRingSensorHandler extends SensorHandler { } streamController = StreamController>.broadcast( - onListen: () { - bleSubscription ??= _bleManager - .subscribe( - deviceId: _discoveredDevice.id, - serviceId: OpenRingGatt.service, - characteristicId: OpenRingGatt.rxChar, - ) - .listen( - (data) { - scheduler.resetIfRequested( - _transportTimingResetCounter, - const [OpenRingGatt.cmdPPGQ2], - ); - - final int? rawCmd = data.length > 2 ? data[2] : null; - if (rawCmd != null) { - scheduler.notePacketQueued(rawCmd); - } - - final int arrivalMs = scheduler.nowMonotonicMs; - final int queueKey = rawCmd ?? -1; - final Future previousQueue = - processingQueueByCmd[queueKey] ?? Future.value(); - - processingQueueByCmd[queueKey] = previousQueue - .then((_) => processPacket(data, arrivalMs, rawCmd)) - .catchError((error) { - logger.e( - 'Error while parsing OpenRing sensor packet: $error', - ); - }); - }, - onError: (error) { - logger.e('Error while subscribing to sensor data: $error'); - if (!streamController.isClosed) { - streamController.addError(error); - } - }, - ); - }, onCancel: () { if (!streamController.hasListener) { - final subscription = bleSubscription; - bleSubscription = null; + _sensorDataStream = null; processingQueueByCmd.clear(); scheduler.clear(); + unawaited(bleSubscription?.cancel()); + bleSubscription = null; + } + }, + ); - if (subscription != null) { - unawaited(subscription.cancel()); - } + final dataStream = await _bleManager.subscribe( + deviceId: _discoveredDevice.id, + serviceId: OpenRingGatt.service, + characteristicId: OpenRingGatt.rxChar, + ); + + bleSubscription = dataStream.listen( + (data) { + scheduler.resetIfRequested( + _transportTimingResetCounter, + const [OpenRingGatt.cmdPPGQ2], + ); + + final int? rawCmd = data.length > 2 ? data[2] : null; + if (rawCmd != null) { + scheduler.notePacketQueued(rawCmd); + } + + final int arrivalMs = scheduler.nowMonotonicMs; + final int queueKey = rawCmd ?? -1; + final Future previousQueue = + processingQueueByCmd[queueKey] ?? Future.value(); + + processingQueueByCmd[queueKey] = previousQueue + .then((_) => processPacket(data, arrivalMs, rawCmd)) + .catchError((error) { + logger.e( + 'Error while parsing OpenRing sensor packet: $error', + ); + }); + }, + onError: (error) { + logger.e('Error while subscribing to sensor data: $error'); + if (!streamController.isClosed) { + streamController.addError(error); } }, ); diff --git a/lib/src/managers/sensor_handler.dart b/lib/src/managers/sensor_handler.dart index b0d0943b..123b8d91 100644 --- a/lib/src/managers/sensor_handler.dart +++ b/lib/src/managers/sensor_handler.dart @@ -3,7 +3,10 @@ abstract class SensorHandler { /// /// The [sensorId] parameter specifies the ID of the sensor to subscribe to. /// Returns a [Stream] of sensor data as a [Map] of sensor values. - Stream> subscribeToSensorData(int sensorId); + /// + /// The returned future completes only after the underlying notification + /// subscription has been enabled. + Future>> subscribeToSensorData(int sensorId); /// Writes the sensor configuration to the OpenEarable device. /// diff --git a/lib/src/managers/v2_sensor_handler.dart b/lib/src/managers/v2_sensor_handler.dart index 2f5a52d1..74bdefde 100644 --- a/lib/src/managers/v2_sensor_handler.dart +++ b/lib/src/managers/v2_sensor_handler.dart @@ -26,7 +26,9 @@ class V2SensorHandler extends SensorHandler { _sensorValueParser = sensorValueParser; @override - Stream> subscribeToSensorData(int sensorId) { + Future>> subscribeToSensorData( + int sensorId, + ) async { if (!_bleManager.isConnected(_discoveredDevice.id)) { throw Exception("Can't subscribe to sensor data. Earable not connected"); } @@ -37,13 +39,13 @@ class V2SensorHandler extends SensorHandler { StreamController> streamController = StreamController(); - final subscription = _bleManager - .subscribe( + final dataStream = await _bleManager.subscribe( deviceId: _discoveredDevice.id, serviceId: sensorServiceUuid, characteristicId: sensorDataCharacteristicUuid, - ) - .listen( + ); + + final subscription = dataStream.listen( (data) async { if (data.isNotEmpty && data[0] == sensorId) { List> parsedData = await _parseData(data); @@ -82,14 +84,14 @@ class V2SensorHandler extends SensorHandler { ); } - /// Parses raw sensor data bytes into a [Map] of sensor values. + /// Parses raw sensor data bytes into a [Map] of sensor values. Future>> _parseData(List data) async { ByteData byteData = ByteData.sublistView(Uint8List.fromList(data)); if (_sensorSchemes == null) { await _readSensorScheme(); } - + return _sensorValueParser.parse(byteData, _sensorSchemes!); } diff --git a/lib/src/models/capabilities/sensor_specializations/open_ring/open_ring_sensor.dart b/lib/src/models/capabilities/sensor_specializations/open_ring/open_ring_sensor.dart index a363db8e..dda35337 100644 --- a/lib/src/models/capabilities/sensor_specializations/open_ring/open_ring_sensor.dart +++ b/lib/src/models/capabilities/sensor_specializations/open_ring/open_ring_sensor.dart @@ -44,20 +44,40 @@ class OpenRingSensor extends Sensor { Stream get sensorStream => _sensorStreamController.stream; void _handleListen() { - _sensorSubscription ??= - sensorHandler.subscribeToSensorData(sensorId).listen( - (data) { - final SensorDoubleValue? sensorValue = _toSensorValue(data); - if (sensorValue != null && !_sensorStreamController.isClosed) { - _sensorStreamController.add(sensorValue); - } - }, - onError: (error, stack) { - if (!_sensorStreamController.isClosed) { - _sensorStreamController.addError(error, stack); - } - }, - ); + if (_sensorSubscription != null) { + return; + } + + unawaited(_subscribeToSensorData()); + } + + Future _subscribeToSensorData() async { + try { + final sensorDataStream = + await sensorHandler.subscribeToSensorData(sensorId); + if (_sensorStreamController.isClosed || + _sensorSubscription != null || + !_sensorStreamController.hasListener) { + return; + } + _sensorSubscription = sensorDataStream.listen( + (data) { + final SensorDoubleValue? sensorValue = _toSensorValue(data); + if (sensorValue != null && !_sensorStreamController.isClosed) { + _sensorStreamController.add(sensorValue); + } + }, + onError: (error, stack) { + if (!_sensorStreamController.isClosed) { + _sensorStreamController.addError(error, stack); + } + }, + ); + } catch (error, stack) { + if (!_sensorStreamController.isClosed) { + _sensorStreamController.addError(error, stack); + } + } } Future _handleCancel() async { diff --git a/lib/src/models/devices/cosinuss_one.dart b/lib/src/models/devices/cosinuss_one.dart index 7c511585..9609c13d 100644 --- a/lib/src/models/devices/cosinuss_one.dart +++ b/lib/src/models/devices/cosinuss_one.dart @@ -118,28 +118,36 @@ class CosinussOne extends Wearable @override Stream get batteryPercentageStream { StreamController streamController = StreamController(); - - StreamSubscription subscription = _bleManager - .subscribe( - deviceId: _discoveredDevice.id, - serviceId: batteryServiceUuid, - characteristicId: _batteryLevelCharacteristicUuid, - ) - .listen((data) { - streamController.add(data[0]); - }); - - readBatteryPercentage().then((percentage) { - streamController.add(percentage); - streamController.close(); - }).catchError((error) { - streamController.addError(error); - streamController.close(); - }); + StreamSubscription>? subscription; + + streamController.onListen = () async { + try { + final batteryStream = await _bleManager.subscribe( + deviceId: _discoveredDevice.id, + serviceId: batteryServiceUuid, + characteristicId: _batteryLevelCharacteristicUuid, + ); + if (streamController.isClosed) { + return; + } + subscription = batteryStream.listen((data) { + streamController.add(data[0]); + }); + + final percentage = await readBatteryPercentage(); + streamController.add(percentage); + await streamController.close(); + } catch (error, stack) { + if (!streamController.isClosed) { + streamController.addError(error, stack); + await streamController.close(); + } + } + }; // Cancel BLE subscription when canceling stream streamController.onCancel = () { - subscription.cancel(); + subscription?.cancel(); }; return streamController.stream; @@ -201,41 +209,52 @@ class _CosinussOneSensor extends Sensor { Stream _createAccStream() { StreamController streamController = StreamController(); + StreamSubscription>? subscription; int startTime = DateTime.now().millisecondsSinceEpoch; - _bleManager.write( - deviceId: _discoveredDevice.id, - serviceId: CosinussOne.ppgAndAccServiceUuid, - characteristicId: "0000a001-1212-efde-1523-785feabcd123", - byteData: _sensorBluetoothCharacteristics, - ); - - StreamSubscription subscription = _bleManager - .subscribe( - deviceId: _discoveredDevice.id, - serviceId: CosinussOne.ppgAndAccServiceUuid, - characteristicId: "0000a001-1212-efde-1523-785feabcd123", - ) - .listen((data) { - Int8List bytes = Int8List.fromList(data); - - // description based on placing the earable into your right ear canal - int accX = bytes[14]; - int accY = bytes[16]; - int accZ = bytes[18]; - - streamController.add( - SensorDoubleValue( - values: [accX.toDouble(), accY.toDouble(), accZ.toDouble()], - timestamp: DateTime.now().millisecondsSinceEpoch - startTime, - ), - ); - }); + streamController.onListen = () async { + try { + final sensorStream = await _bleManager.subscribe( + deviceId: _discoveredDevice.id, + serviceId: CosinussOne.ppgAndAccServiceUuid, + characteristicId: "0000a001-1212-efde-1523-785feabcd123", + ); + if (streamController.isClosed) { + return; + } + subscription = sensorStream.listen((data) { + Int8List bytes = Int8List.fromList(data); + + // description based on placing the earable into your right ear canal + int accX = bytes[14]; + int accY = bytes[16]; + int accZ = bytes[18]; + + streamController.add( + SensorDoubleValue( + values: [accX.toDouble(), accY.toDouble(), accZ.toDouble()], + timestamp: DateTime.now().millisecondsSinceEpoch - startTime, + ), + ); + }); + + await _bleManager.write( + deviceId: _discoveredDevice.id, + serviceId: CosinussOne.ppgAndAccServiceUuid, + characteristicId: "0000a001-1212-efde-1523-785feabcd123", + byteData: _sensorBluetoothCharacteristics, + ); + } catch (error, stack) { + if (!streamController.isClosed) { + streamController.addError(error, stack); + } + } + }; // Cancel BLE subscription when canceling stream streamController.onCancel = () { - subscription.cancel(); + subscription?.cancel(); }; return streamController.stream; @@ -243,60 +262,71 @@ class _CosinussOneSensor extends Sensor { Stream _createPpgStream() { StreamController streamController = StreamController(); + StreamSubscription>? subscription; int startTime = DateTime.now().millisecondsSinceEpoch; - _bleManager.write( - deviceId: _discoveredDevice.id, - serviceId: CosinussOne.ppgAndAccServiceUuid, - characteristicId: "0000a001-1212-efde-1523-785feabcd123", - byteData: _sensorBluetoothCharacteristics, - ); - - StreamSubscription subscription = _bleManager - .subscribe( - deviceId: _discoveredDevice.id, - serviceId: CosinussOne.ppgAndAccServiceUuid, - characteristicId: "0000a001-1212-efde-1523-785feabcd123", - ) - .listen((data) { - Uint8List bytes = Uint8List.fromList(data); - - // corresponds to the raw reading of the PPG sensor from which the heart rate is computed - // - // example plot https://e2e.ti.com/cfs-file/__key/communityserver-discussions-components-files/73/Screen-Shot-2019_2D00_01_2D00_24-at-19.30.24.png - // (image just for illustration purpose, obtained from a different sensor! Sensor value range differs.) - - var ppgRed = bytes[0] | - bytes[1] << 8 | - bytes[2] << 16 | - bytes[3] << 32; // raw green color value of PPG sensor - var ppgGreen = bytes[4] | - bytes[5] << 8 | - bytes[6] << 16 | - bytes[7] << 32; // raw red color value of PPG sensor - - var ppgGreenAmbient = bytes[8] | - bytes[9] << 8 | - bytes[10] << 16 | - bytes[11] << - 32; // ambient light sensor (e.g., if sensor is not placed correctly) - - streamController.add( - SensorDoubleValue( - values: [ - ppgRed.toDouble(), - ppgGreen.toDouble(), - ppgGreenAmbient.toDouble(), - ], - timestamp: DateTime.now().millisecondsSinceEpoch - startTime, - ), - ); - }); + streamController.onListen = () async { + try { + final sensorStream = await _bleManager.subscribe( + deviceId: _discoveredDevice.id, + serviceId: CosinussOne.ppgAndAccServiceUuid, + characteristicId: "0000a001-1212-efde-1523-785feabcd123", + ); + if (streamController.isClosed) { + return; + } + subscription = sensorStream.listen((data) { + Uint8List bytes = Uint8List.fromList(data); + + // corresponds to the raw reading of the PPG sensor from which the heart rate is computed + // + // example plot https://e2e.ti.com/cfs-file/__key/communityserver-discussions-components-files/73/Screen-Shot-2019_2D00_01_2D00_24-at-19.30.24.png + // (image just for illustration purpose, obtained from a different sensor! Sensor value range differs.) + + var ppgRed = bytes[0] | + bytes[1] << 8 | + bytes[2] << 16 | + bytes[3] << 32; // raw green color value of PPG sensor + var ppgGreen = bytes[4] | + bytes[5] << 8 | + bytes[6] << 16 | + bytes[7] << 32; // raw red color value of PPG sensor + + var ppgGreenAmbient = bytes[8] | + bytes[9] << 8 | + bytes[10] << 16 | + bytes[11] << + 32; // ambient light sensor (e.g., if sensor is not placed correctly) + + streamController.add( + SensorDoubleValue( + values: [ + ppgRed.toDouble(), + ppgGreen.toDouble(), + ppgGreenAmbient.toDouble(), + ], + timestamp: DateTime.now().millisecondsSinceEpoch - startTime, + ), + ); + }); + + await _bleManager.write( + deviceId: _discoveredDevice.id, + serviceId: CosinussOne.ppgAndAccServiceUuid, + characteristicId: "0000a001-1212-efde-1523-785feabcd123", + byteData: _sensorBluetoothCharacteristics, + ); + } catch (error, stack) { + if (!streamController.isClosed) { + streamController.addError(error, stack); + } + } + }; // Cancel BLE subscription when canceling stream streamController.onCancel = () { - subscription.cancel(); + subscription?.cancel(); }; return streamController.stream; @@ -304,39 +334,50 @@ class _CosinussOneSensor extends Sensor { Stream _createTempStream() { StreamController streamController = StreamController(); + StreamSubscription>? subscription; int startTime = DateTime.now().millisecondsSinceEpoch; - StreamSubscription subscription = _bleManager - .subscribe( - deviceId: _discoveredDevice.id, - serviceId: CosinussOne.temperatureServiceUuid, - characteristicId: "00002a1c-0000-1000-8000-00805f9b34fb", - ) - .listen((data) { - var flag = data[0]; - - // based on GATT standard - double temperature = _twosComplimentOfNegativeMantissa( - ((data[3] << 16) | (data[2] << 8) | data[1]) & 16777215, - ) / - 100.0; - if ((flag & 1) != 0) { - temperature = ((98.6 * temperature) - 32.0) * - (5.0 / 9.0); // convert Fahrenheit to Celsius + streamController.onListen = () async { + try { + final temperatureStream = await _bleManager.subscribe( + deviceId: _discoveredDevice.id, + serviceId: CosinussOne.temperatureServiceUuid, + characteristicId: "00002a1c-0000-1000-8000-00805f9b34fb", + ); + if (streamController.isClosed) { + return; + } + subscription = temperatureStream.listen((data) { + var flag = data[0]; + + // based on GATT standard + double temperature = _twosComplimentOfNegativeMantissa( + ((data[3] << 16) | (data[2] << 8) | data[1]) & 16777215, + ) / + 100.0; + if ((flag & 1) != 0) { + temperature = ((98.6 * temperature) - 32.0) * + (5.0 / 9.0); // convert Fahrenheit to Celsius + } + + streamController.add( + SensorDoubleValue( + values: [temperature], + timestamp: DateTime.now().millisecondsSinceEpoch - startTime, + ), + ); + }); + } catch (error, stack) { + if (!streamController.isClosed) { + streamController.addError(error, stack); + } } - - streamController.add( - SensorDoubleValue( - values: [temperature], - timestamp: DateTime.now().millisecondsSinceEpoch - startTime, - ), - ); - }); + }; // Cancel BLE subscription when canceling stream streamController.onCancel = () { - subscription.cancel(); + subscription?.cancel(); }; return streamController.stream; @@ -373,35 +414,46 @@ class _CosinussOneHeartRateSensor extends HeartRateSensor { Stream get sensorStream { StreamController streamController = StreamController(); + StreamSubscription>? subscription; int startTime = DateTime.now().millisecondsSinceEpoch; - StreamSubscription subscription = _bleManager - .subscribe( - deviceId: _discoveredDevice.id, - serviceId: CosinussOne.heartRateServiceUuid, - characteristicId: "00002a37-0000-1000-8000-00805f9b34fb", - ) - .listen((data) { - Uint8List bytes = Uint8List.fromList(data); - - // based on GATT standard - int bpm = bytes[1]; - if (!((bytes[0] & 0x01) == 0)) { - bpm = (((bpm >> 8) & 0xFF) | ((bpm << 8) & 0xFF00)); + streamController.onListen = () async { + try { + final heartRateStream = await _bleManager.subscribe( + deviceId: _discoveredDevice.id, + serviceId: CosinussOne.heartRateServiceUuid, + characteristicId: "00002a37-0000-1000-8000-00805f9b34fb", + ); + if (streamController.isClosed) { + return; + } + subscription = heartRateStream.listen((data) { + Uint8List bytes = Uint8List.fromList(data); + + // based on GATT standard + int bpm = bytes[1]; + if (!((bytes[0] & 0x01) == 0)) { + bpm = (((bpm >> 8) & 0xFF) | ((bpm << 8) & 0xFF00)); + } + + streamController.add( + HeartRateSensorValue( + heartRateBpm: bpm, + timestamp: DateTime.now().millisecondsSinceEpoch - startTime, + ), + ); + }); + } catch (error, stack) { + if (!streamController.isClosed) { + streamController.addError(error, stack); + } } - - streamController.add( - HeartRateSensorValue( - heartRateBpm: bpm, - timestamp: DateTime.now().millisecondsSinceEpoch - startTime, - ), - ); - }); + }; // Cancel BLE subscription when canceling stream streamController.onCancel = () { - subscription.cancel(); + subscription?.cancel(); }; return streamController.stream; diff --git a/lib/src/models/devices/esense_factory.dart b/lib/src/models/devices/esense_factory.dart index 1f798ad2..c0d4fe0d 100644 --- a/lib/src/models/devices/esense_factory.dart +++ b/lib/src/models/devices/esense_factory.dart @@ -128,35 +128,54 @@ class EsenseSensor extends Sensor { Stream get sensorStream { StreamController streamController = StreamController(); - _sensorHandler.subscribeToSensorData(_sensorId).listen( - (data) { - int timestamp = data["timestamp"]; - - List values = []; - for (var entry in (data[sensorName] as Map).entries) { - if (entry.key == 'units') { - continue; - } - - if (entry.value is int) { - values.add((entry.value as int).toDouble()); - } else if (entry.value is double) { - values.add(entry.value as double); - } else { - throw UnsupportedError( - "Unsupported sensor value type: ${entry.value.runtimeType}", - ); - } + StreamSubscription>? subscription; + + streamController.onListen = () async { + try { + final sensorDataStream = + await _sensorHandler.subscribeToSensorData(_sensorId); + if (streamController.isClosed) { + return; } + subscription = sensorDataStream.listen( + (data) { + int timestamp = data["timestamp"]; + + List values = []; + for (var entry in (data[sensorName] as Map).entries) { + if (entry.key == 'units') { + continue; + } + + if (entry.value is int) { + values.add((entry.value as int).toDouble()); + } else if (entry.value is double) { + values.add(entry.value as double); + } else { + throw UnsupportedError( + "Unsupported sensor value type: ${entry.value.runtimeType}", + ); + } + } + + SensorDoubleValue sensorValue = SensorDoubleValue( + values: values, + timestamp: timestamp, + ); - SensorDoubleValue sensorValue = SensorDoubleValue( - values: values, - timestamp: timestamp, + streamController.add(sensorValue); + }, ); + } catch (error, stack) { + if (!streamController.isClosed) { + streamController.addError(error, stack); + } + } + }; - streamController.add(sensorValue); - }, - ); + streamController.onCancel = () { + subscription?.cancel(); + }; return streamController.stream; } diff --git a/lib/src/models/devices/open_earable_factory.dart b/lib/src/models/devices/open_earable_factory.dart index 2dd59b4b..dc0f95a3 100644 --- a/lib/src/models/devices/open_earable_factory.dart +++ b/lib/src/models/devices/open_earable_factory.dart @@ -350,33 +350,47 @@ class _OpenEarableSensorV2 extends Sensor { ) { StreamController streamController = StreamController(); - StreamSubscription subscription = - _sensorManager.subscribeToSensorData(_sensorId).listen((data) { - int timestamp = data["timestamp"]; - logger.t("SensorData: $data"); - - logger.t("componentData of $componentName: ${data[componentName]}"); - - //TODO: use int for integer based values - List values = []; - for (var entry in (data[componentName] as Map).entries) { - if (entry.key == 'units') { - continue; + StreamSubscription>? subscription; + + streamController.onListen = () async { + try { + final sensorDataStream = + await _sensorManager.subscribeToSensorData(_sensorId); + if (streamController.isClosed) { + return; } + subscription = sensorDataStream.listen((data) { + int timestamp = data["timestamp"]; + logger.t("SensorData: $data"); - values.add(entry.value.toDouble()); - } + logger.t("componentData of $componentName: ${data[componentName]}"); - SensorDoubleValue sensorValue = SensorDoubleValue( - values: values, - timestamp: timestamp, - ); + //TODO: use int for integer based values + List values = []; + for (var entry in (data[componentName] as Map).entries) { + if (entry.key == 'units') { + continue; + } + + values.add(entry.value.toDouble()); + } - streamController.add(sensorValue); - }); + SensorDoubleValue sensorValue = SensorDoubleValue( + values: values, + timestamp: timestamp, + ); + + streamController.add(sensorValue); + }); + } catch (error, stack) { + if (!streamController.isClosed) { + streamController.addError(error, stack); + } + } + }; streamController.onCancel = () { - subscription.cancel(); + subscription?.cancel(); }; return streamController.stream; diff --git a/lib/src/models/devices/open_earable_v1.dart b/lib/src/models/devices/open_earable_v1.dart index 6366c52e..74d0baf8 100644 --- a/lib/src/models/devices/open_earable_v1.dart +++ b/lib/src/models/devices/open_earable_v1.dart @@ -484,25 +484,38 @@ class _OpenEarableSensor extends Sensor { q: 0.9, ); - StreamSubscription subscription = - _sensorManager.subscribeToSensorData(0).listen((data) { - int timestamp = data["timestamp"]; - - SensorDoubleValue sensorValue = SensorDoubleValue( - values: [ - kalmanX.filtered(data[sensorName]["X"]), - kalmanY.filtered(data[sensorName]["Y"]), - kalmanZ.filtered(data[sensorName]["Z"]), - ], - timestamp: timestamp, - ); + StreamSubscription>? subscription; - streamController.add(sensorValue); - }); + streamController.onListen = () async { + try { + final sensorDataStream = await _sensorManager.subscribeToSensorData(0); + if (streamController.isClosed) { + return; + } + subscription = sensorDataStream.listen((data) { + int timestamp = data["timestamp"]; + + SensorDoubleValue sensorValue = SensorDoubleValue( + values: [ + kalmanX.filtered(data[sensorName]["X"]), + kalmanY.filtered(data[sensorName]["Y"]), + kalmanZ.filtered(data[sensorName]["Z"]), + ], + timestamp: timestamp, + ); + + streamController.add(sensorValue); + }); + } catch (error, stack) { + if (!streamController.isClosed) { + streamController.addError(error, stack); + } + } + }; // Cancel BLE subscription when canceling stream streamController.onCancel = () { - subscription.cancel(); + subscription?.cancel(); }; return streamController.stream; @@ -513,21 +526,34 @@ class _OpenEarableSensor extends Sensor { ) { StreamController streamController = StreamController(); - StreamSubscription subscription = - _sensorManager.subscribeToSensorData(1).listen((data) { - int timestamp = data["timestamp"]; + StreamSubscription>? subscription; - SensorDoubleValue sensorValue = SensorDoubleValue( - values: [data[sensorName][componentName]], - timestamp: timestamp, - ); - - streamController.add(sensorValue); - }); + streamController.onListen = () async { + try { + final sensorDataStream = await _sensorManager.subscribeToSensorData(1); + if (streamController.isClosed) { + return; + } + subscription = sensorDataStream.listen((data) { + int timestamp = data["timestamp"]; + + SensorDoubleValue sensorValue = SensorDoubleValue( + values: [data[sensorName][componentName]], + timestamp: timestamp, + ); + + streamController.add(sensorValue); + }); + } catch (error, stack) { + if (!streamController.isClosed) { + streamController.addError(error, stack); + } + } + }; // Cancel BLE subscription when canceling stream streamController.onCancel = () { - subscription.cancel(); + subscription?.cancel(); }; return streamController.stream; diff --git a/lib/src/models/devices/open_earable_v2.dart b/lib/src/models/devices/open_earable_v2.dart index 235c4fd3..7067c5d2 100644 --- a/lib/src/models/devices/open_earable_v2.dart +++ b/lib/src/models/devices/open_earable_v2.dart @@ -109,30 +109,39 @@ class OpenEarableV2 extends BluetoothWearable controller = StreamController>(); - _sensorConfigSubscription?.cancel(); - - _sensorConfigSubscription = bleManager - .subscribe( - deviceId: deviceId, - serviceId: sensorServiceUuid, - characteristicId: sensorConfigStateCharacteristicUuid, - ) - .listen( - (data) { - controller.add(_parseConfigMap(data)); - }, - onError: (error) { - logger.e('Error in sensor configuration stream: $error'); - controller.addError(error); - }, - ); - controller.onCancel = () { _sensorConfigSubscription?.cancel(); _sensorConfigSubscription = null; }; - controller.onListen = () { + controller.onListen = () async { + _sensorConfigSubscription?.cancel(); + + try { + final configStream = await bleManager.subscribe( + deviceId: deviceId, + serviceId: sensorServiceUuid, + characteristicId: sensorConfigStateCharacteristicUuid, + ); + if (controller.isClosed) { + return; + } + _sensorConfigSubscription = configStream.listen( + (data) { + controller.add(_parseConfigMap(data)); + }, + onError: (error) { + logger.e('Error in sensor configuration stream: $error'); + controller.addError(error); + }, + ); + } catch (error, stack) { + logger.e('Error setting up sensor configuration stream: $error'); + if (!controller.isClosed) { + controller.addError(error, stack); + } + } + // Immediately read the current sensor configuration bleManager .read( @@ -222,37 +231,46 @@ class OpenEarableV2 extends BluetoothWearable Stream get buttonEvents { StreamController controller = StreamController(); - _buttonSubscription?.cancel(); - - _buttonSubscription = bleManager - .subscribe( - deviceId: deviceId, - serviceId: _buttonServiceUuid, - characteristicId: _buttonCharacteristicUuid, - ) - .listen( - (data) { - if (data.isNotEmpty) { - int buttonState = data[0]; - if (buttonState == 0) { - controller.add(ButtonEvent.released); - } else if (buttonState == 1) { - controller.add(ButtonEvent.pressed); - } - } - }, - onError: (error) { - logger.e('Error in button events stream: $error'); - controller.addError(error); - }, - ); - controller.onCancel = () { _buttonSubscription?.cancel(); _buttonSubscription = null; }; - controller.onListen = () { + controller.onListen = () async { + _buttonSubscription?.cancel(); + + try { + final buttonStream = await bleManager.subscribe( + deviceId: deviceId, + serviceId: _buttonServiceUuid, + characteristicId: _buttonCharacteristicUuid, + ); + if (controller.isClosed) { + return; + } + _buttonSubscription = buttonStream.listen( + (data) { + if (data.isNotEmpty) { + int buttonState = data[0]; + if (buttonState == 0) { + controller.add(ButtonEvent.released); + } else if (buttonState == 1) { + controller.add(ButtonEvent.pressed); + } + } + }, + onError: (error) { + logger.e('Error in button events stream: $error'); + controller.addError(error); + }, + ); + } catch (error, stack) { + logger.e('Error setting up button events stream: $error'); + if (!controller.isClosed) { + controller.addError(error, stack); + } + } + // Immediately read current button state bleManager .read( @@ -830,13 +848,12 @@ class OpenEarableV2TimeSyncImp implements TimeSynchronizable { // Subscribe to RTT responses late final StreamSubscription> rttSub; - rttSub = bleManager - .subscribe( + final rttStream = await bleManager.subscribe( deviceId: deviceId, serviceId: timeSynchronizationServiceUuid, characteristicId: _timeSyncRttCharacteristicUuid, - ) - .listen( + ); + rttSub = rttStream.listen( (data) async { final t4 = DateTime.now().microsecondsSinceEpoch; final pkt = _SyncTimePacket.fromBytes(Uint8List.fromList(data)); diff --git a/lib/src/models/devices/open_ring.dart b/lib/src/models/devices/open_ring.dart index cf5d99e4..96f41b20 100644 --- a/lib/src/models/devices/open_ring.dart +++ b/lib/src/models/devices/open_ring.dart @@ -192,13 +192,12 @@ class OpenRing extends Wearable final completer = Completer(); late final StreamSubscription> sub; - sub = _bleManager - .subscribe( + final batteryStream = await _bleManager.subscribe( deviceId: deviceId, serviceId: OpenRingGatt.service, characteristicId: OpenRingGatt.rxChar, - ) - .listen( + ); + sub = batteryStream.listen( (data) { final response = _parseBatteryResponse(data); if (response == null || !response.isRead) { @@ -330,30 +329,39 @@ class OpenRing extends Wearable unawaited(batteryPushSubscription?.cancel()); }; - controller.onListen = () { + controller.onListen = () async { final initialBatteryPercentage = _lastKnownBatteryPercentage; if (initialBatteryPercentage != null) { emitIfChanged(initialBatteryPercentage); } - batteryPushSubscription = _bleManager - .subscribe( - deviceId: deviceId, - serviceId: OpenRingGatt.service, - characteristicId: OpenRingGatt.rxChar, - ) - .listen( - (data) { - final response = _parseBatteryResponse(data); - if (response == null || !response.isPush) { - return; - } - emitIfChanged(response.batteryPercentage); - }, - onError: (error) { - logger.w('OpenRing battery push subscription error: $error'); - }, - ); + try { + final batteryStream = await _bleManager.subscribe( + deviceId: deviceId, + serviceId: OpenRingGatt.service, + characteristicId: OpenRingGatt.rxChar, + ); + if (controller.isClosed) { + return; + } + batteryPushSubscription = batteryStream.listen( + (data) { + final response = _parseBatteryResponse(data); + if (response == null || !response.isPush) { + return; + } + emitIfChanged(response.batteryPercentage); + }, + onError: (error) { + logger.w('OpenRing battery push subscription error: $error'); + }, + ); + } catch (error, stack) { + logger.w('OpenRing battery push subscription setup error: $error'); + if (!controller.isClosed) { + controller.addError(error, stack); + } + } batteryPollingTimer = Timer.periodic(const Duration(seconds: 5), (timer) { unawaited(pollBattery()); @@ -485,13 +493,12 @@ class OpenRingTimeSyncImp implements TimeSynchronizable { final completer = Completer(); late final StreamSubscription> sub; - sub = bleManager - .subscribe( + final timeSyncStream = await bleManager.subscribe( deviceId: deviceId, serviceId: OpenRingGatt.service, characteristicId: OpenRingGatt.rxChar, - ) - .listen( + ); + sub = timeSyncStream.listen( (data) { if (data.length < 4) { return; diff --git a/lib/src/models/devices/polar_factory.dart b/lib/src/models/devices/polar_factory.dart index c0ab8c6f..961e46ea 100644 --- a/lib/src/models/devices/polar_factory.dart +++ b/lib/src/models/devices/polar_factory.dart @@ -96,35 +96,46 @@ class _PolarHeartRateSensor extends HeartRateSensor { Stream get sensorStream { StreamController streamController = StreamController(); + StreamSubscription>? subscription; int startTime = DateTime.now().millisecondsSinceEpoch; - StreamSubscription subscription = _bleManager - .subscribe( - deviceId: _discoveredDevice.id, - serviceId: Polar.heartRateServiceUuid, - characteristicId: "00002a37-0000-1000-8000-00805f9b34fb", - ) - .listen((data) { - Uint8List bytes = Uint8List.fromList(data); - - int hrFormat = bytes[0] & 0x01; - - int heartRate = hrFormat == 1 - ? (bytes[1] & 0xFF) | ((bytes[2] & 0xFF) << 8) - : bytes[1] & 0xFF; - - streamController.add( - HeartRateSensorValue( - heartRateBpm: heartRate, - timestamp: DateTime.now().millisecondsSinceEpoch - startTime, - ), - ); - }); + streamController.onListen = () async { + try { + final heartRateStream = await _bleManager.subscribe( + deviceId: _discoveredDevice.id, + serviceId: Polar.heartRateServiceUuid, + characteristicId: "00002a37-0000-1000-8000-00805f9b34fb", + ); + if (streamController.isClosed) { + return; + } + subscription = heartRateStream.listen((data) { + Uint8List bytes = Uint8List.fromList(data); + + int hrFormat = bytes[0] & 0x01; + + int heartRate = hrFormat == 1 + ? (bytes[1] & 0xFF) | ((bytes[2] & 0xFF) << 8) + : bytes[1] & 0xFF; + + streamController.add( + HeartRateSensorValue( + heartRateBpm: heartRate, + timestamp: DateTime.now().millisecondsSinceEpoch - startTime, + ), + ); + }); + } catch (error, stack) { + if (!streamController.isClosed) { + streamController.addError(error, stack); + } + } + }; // Cancel BLE subscription when canceling stream streamController.onCancel = () { - subscription.cancel(); + subscription?.cancel(); }; return streamController.stream; @@ -145,41 +156,52 @@ class _PolarHeartRateVariabilitySensor extends HeartRateVariabilitySensor { DiscoveredDevice discoveredDevice, ) { StreamController> streamController = StreamController(); - - StreamSubscription subscription = bleManager - .subscribe( - deviceId: discoveredDevice.id, - serviceId: Polar.heartRateServiceUuid, - characteristicId: "00002a37-0000-1000-8000-00805f9b34fb", - ) - .listen((data) { - Uint8List bytes = Uint8List.fromList(data); - - int hrFormat = bytes[0] & 0x01; - bool rrPresent = (bytes[0] & 0x10) >> 4 == 1; - int energyExpendedFlag = (bytes[0] & 0x08) >> 3; - - int offset = hrFormat + 2; - if (energyExpendedFlag == 1) { - offset += 2; - } - - List rrIntervalsMs = []; - if (rrPresent) { - while (offset + 1 < bytes.length) { - int rrValue = - (bytes[offset] & 0xFF) | ((bytes[offset + 1] & 0xFF) << 8); - offset += 2; - rrIntervalsMs.add(_mapRr1024ToRrMs(rrValue)); + StreamSubscription>? subscription; + + streamController.onListen = () async { + try { + final heartRateStream = await bleManager.subscribe( + deviceId: discoveredDevice.id, + serviceId: Polar.heartRateServiceUuid, + characteristicId: "00002a37-0000-1000-8000-00805f9b34fb", + ); + if (streamController.isClosed) { + return; + } + subscription = heartRateStream.listen((data) { + Uint8List bytes = Uint8List.fromList(data); + + int hrFormat = bytes[0] & 0x01; + bool rrPresent = (bytes[0] & 0x10) >> 4 == 1; + int energyExpendedFlag = (bytes[0] & 0x08) >> 3; + + int offset = hrFormat + 2; + if (energyExpendedFlag == 1) { + offset += 2; + } + + List rrIntervalsMs = []; + if (rrPresent) { + while (offset + 1 < bytes.length) { + int rrValue = + (bytes[offset] & 0xFF) | ((bytes[offset + 1] & 0xFF) << 8); + offset += 2; + rrIntervalsMs.add(_mapRr1024ToRrMs(rrValue)); + } + + streamController.add(rrIntervalsMs); + } + }); + } catch (error, stack) { + if (!streamController.isClosed) { + streamController.addError(error, stack); } - - streamController.add(rrIntervalsMs); } - }); + }; // Cancel BLE subscription when canceling stream streamController.onCancel = () { - subscription.cancel(); + subscription?.cancel(); }; return streamController.stream; diff --git a/lib/src/utils/sensor_scheme_parser/v2_sensor_scheme_reader.dart b/lib/src/utils/sensor_scheme_parser/v2_sensor_scheme_reader.dart index 07b03df8..0e2d781f 100644 --- a/lib/src/utils/sensor_scheme_parser/v2_sensor_scheme_reader.dart +++ b/lib/src/utils/sensor_scheme_parser/v2_sensor_scheme_reader.dart @@ -60,7 +60,7 @@ class V2SensorSchemeReader extends SensorSchemeReader { } // Listen to the notification of the characteristic - final Stream> stream = _bleManager.subscribe( + final Stream> stream = await _bleManager.subscribe( deviceId: _deviceId, serviceId: parseInfoServiceUuid, characteristicId: sensorSchemeCharacteristicUuid, From 321a9a72cf3a665362beea1c26afab9aee31de39 Mon Sep 17 00:00:00 2001 From: Dennis Moschina <45356478+DennisMoschina@users.noreply.github.com> Date: Tue, 21 Jul 2026 15:53:05 +0200 Subject: [PATCH 02/11] chore(changelog): update unreleased section Document the breaking BLE subscription API changes and the notification setup race fix in the unreleased changelog. --- CHANGELOG.md | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 74629a5e..349c5cbd 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,9 @@ +## Unreleased + +* BREAKING CHANGE: `BleGattManager.subscribe` now returns `Future>>`, so callers must `await` subscription setup before listening to BLE notifications. +* BREAKING CHANGE: `SensorHandler.subscribeToSensorData` now returns `Future>>`, so callers must `await` sensor notification readiness before listening to sensor data. +* fixed BLE notification setup races by ensuring subscription futures complete only after the underlying GATT notification subscription is enabled. + ## 2.3.10 * added dynamic power saving mode capability for OpenEarable v2 devices From dbceb7b7b59ba49966bdafb41072a2ff352c112e Mon Sep 17 00:00:00 2001 From: Dennis Moschina <45356478+DennisMoschina@users.noreply.github.com> Date: Wed, 22 Jul 2026 13:36:29 +0200 Subject: [PATCH 03/11] fix(ble_manager): improve subscription teardown handling in BLE manager --- lib/src/managers/ble_manager.dart | 41 ++++++++++++++++++++++++++++--- 1 file changed, 38 insertions(+), 3 deletions(-) diff --git a/lib/src/managers/ble_manager.dart b/lib/src/managers/ble_manager.dart index e787e558..69fa8406 100644 --- a/lib/src/managers/ble_manager.dart +++ b/lib/src/managers/ble_manager.dart @@ -16,6 +16,7 @@ class BleManager extends BleGattManager { final Map>> _streamControllers = {}; final Map> _subscriptionSetups = {}; + final Map> _subscriptionTeardowns = {}; /// A stream of discovered devices during scanning. StreamController? _scanStreamController; @@ -51,6 +52,7 @@ class BleManager extends BleGattManager { logger.d("Closing stream for $key due to device disconnection"); _streamControllers.remove(key)?.close(); _subscriptionSetups.remove(key); + _subscriptionTeardowns.remove(key); } } @@ -326,6 +328,24 @@ class BleManager extends BleGattManager { deviceId, characteristicId, ); + final pendingTeardown = _subscriptionTeardowns[streamIdentifier]; + if (pendingTeardown != null) { + try { + await pendingTeardown; + } catch (error) { + logger.w( + "Previous unsubscribe for $streamIdentifier failed before resubscribe: $error", + ); + } finally { + if (identical( + _subscriptionTeardowns[streamIdentifier], + pendingTeardown, + )) { + _subscriptionTeardowns.remove(streamIdentifier); + } + } + } + StreamController>? streamController = _streamControllers[streamIdentifier]; streamController ??= StreamController>.broadcast(); @@ -347,15 +367,30 @@ class BleManager extends BleGattManager { rethrow; } - streamController.onCancel = () { + streamController.onCancel = () async { if (_streamControllers.containsKey(streamIdentifier)) { - _streamControllers.remove(streamIdentifier)?.close(); + final canceledController = _streamControllers.remove(streamIdentifier); _subscriptionSetups.remove(streamIdentifier); - UniversalBle.unsubscribe( + if (canceledController != null && !canceledController.isClosed) { + unawaited(canceledController.close()); + } + + final teardown = UniversalBle.unsubscribe( deviceId, serviceId, characteristicId, ); + _subscriptionTeardowns[streamIdentifier] = teardown; + + try { + await teardown; + } catch (error) { + logger.w("Unsubscribe failed for $streamIdentifier: $error"); + } finally { + if (identical(_subscriptionTeardowns[streamIdentifier], teardown)) { + _subscriptionTeardowns.remove(streamIdentifier); + } + } } }; From b576d73b9b0b6ebaa773ee136076a741c71fe238 Mon Sep 17 00:00:00 2001 From: Dennis Moschina <45356478+DennisMoschina@users.noreply.github.com> Date: Wed, 22 Jul 2026 13:52:29 +0200 Subject: [PATCH 04/11] fix(ble_manager): enhance connection handling by reusing pending connections --- lib/open_earable_flutter.dart | 27 +++++++++++++++++ lib/src/managers/ble_manager.dart | 50 ++++++++++++++++++++++--------- 2 files changed, 63 insertions(+), 14 deletions(-) diff --git a/lib/open_earable_flutter.dart b/lib/open_earable_flutter.dart index a9cc9266..80e0e7f6 100644 --- a/lib/open_earable_flutter.dart +++ b/lib/open_earable_flutter.dart @@ -105,6 +105,7 @@ class WearableManager { late final StreamController _connectingStreamController; final List _connectedIds = []; + final Map> _connectionFuturesByDeviceId = {}; List _autoConnectDeviceIds = []; StreamSubscription? _autoconnectScanSubscription; @@ -234,6 +235,32 @@ class WearableManager { logger.w('Device ${device.id} is already connected'); throw AlreadyConnectedException(); } + + final pendingConnection = _connectionFuturesByDeviceId[device.id]; + if (pendingConnection != null) { + logger.d('Reusing pending wearable connection for ${device.id}'); + return pendingConnection; + } + + final connectionFuture = _connectToDevice(device, options: options); + _connectionFuturesByDeviceId[device.id] = connectionFuture; + + try { + return await connectionFuture; + } finally { + if (identical( + _connectionFuturesByDeviceId[device.id], + connectionFuture, + )) { + _connectionFuturesByDeviceId.remove(device.id); + } + } + } + + Future _connectToDevice( + DiscoveredDevice device, { + required Set options, + }) async { _connectingStreamController.add(device); WearableDisconnectNotifier disconnectNotifier = diff --git a/lib/src/managers/ble_manager.dart b/lib/src/managers/ble_manager.dart index 69fa8406..21bfe7cb 100644 --- a/lib/src/managers/ble_manager.dart +++ b/lib/src/managers/ble_manager.dart @@ -26,7 +26,9 @@ class BleManager extends BleGattManager { String _getCharacteristicKey(String deviceId, String characteristicId) => "$deviceId||$characteristicId"; - final Map _connectionCompleters = {}; + final Map)>> _connectionCompleters = + {}; + final Map)>> _connectionFutures = {}; final Map _connectCallbacks = {}; final Map _disconnectCallbacks = {}; @@ -218,38 +220,58 @@ class BleManager extends BleGattManager { DiscoveredDevice device, VoidCallback onDisconnect, ) { + final pendingConnection = _connectionFutures[device.id]; + if (pendingConnection != null) { + logger.d("Reusing pending connection for ${device.id}"); + return pendingConnection; + } + // Multi-device setup: only reset stale streams for the device that is // being connected, not globally for all devices. _closeAndRemoveStreamsForDevice(device.id); - Completer<(bool, List)> completer = - Completer<(bool, List)>(); + final completer = Completer<(bool, List)>(); _connectionCompleters[device.id] = completer; + final connectionFuture = completer.future.whenComplete(() { + _connectionFutures.remove(device.id); + }); + _connectionFutures[device.id] = connectionFuture; _connectCallbacks[device.id] = () async { - if (!kIsWeb && !Platform.isLinux) { - _mtu = await UniversalBle.requestMtu(device.id, _desiredMtu); - } - bool connectionResult = false; - List services = []; + try { + if (!kIsWeb && !Platform.isLinux) { + _mtu = await UniversalBle.requestMtu(device.id, _desiredMtu); + } - services = await UniversalBle.discoverServices(device.id); - connectionResult = true; + final services = await UniversalBle.discoverServices(device.id); - _connectionCompleters[device.id]?.complete((connectionResult, services)); - _connectionCompleters.remove(device.id); + _connectionCompleters[device.id]?.complete((true, services)); + } catch (error, stack) { + _connectionCompleters[device.id]?.completeError(error, stack); + } finally { + _connectionCompleters.remove(device.id); + _connectCallbacks.remove(device.id); + } }; _disconnectCallbacks[device.id] = () { _connectionCompleters[device.id]?.complete((false, [])); _connectionCompleters.remove(device.id); + _connectCallbacks.remove(device.id); + _connectionFutures.remove(device.id); onDisconnect(); }; - UniversalBle.connect(device.id); + try { + UniversalBle.connect(device.id); + } catch (error, stack) { + if (!completer.isCompleted) { + completer.completeError(error, stack); + } + } - return completer.future; + return connectionFuture; } /// Checks if the connected device has a specific service. From f58a672c8ceeb95f49e15fd9de12589cdde830a3 Mon Sep 17 00:00:00 2001 From: Dennis Moschina <45356478+DennisMoschina@users.noreply.github.com> Date: Wed, 22 Jul 2026 13:56:47 +0200 Subject: [PATCH 05/11] fix(open_ring_sensor_handler): improve sensor data stream handling to allow retries on failure --- .../managers/open_ring_sensor_handler.dart | 27 ++++++++++++++++--- 1 file changed, 24 insertions(+), 3 deletions(-) diff --git a/lib/src/managers/open_ring_sensor_handler.dart b/lib/src/managers/open_ring_sensor_handler.dart index 6a461064..5b3098d6 100644 --- a/lib/src/managers/open_ring_sensor_handler.dart +++ b/lib/src/managers/open_ring_sensor_handler.dart @@ -60,15 +60,36 @@ class OpenRingSensorHandler extends SensorHandler { throw Exception("Can't subscribe to sensor data. Earable not connected"); } - _sensorDataStream ??= _createSensorDataStream(); - - final sensorDataStream = await _sensorDataStream!; + final sensorDataStream = await _getSensorDataStream(); return sensorDataStream.where((data) { final dynamic cmd = data['cmd']; return cmd is int && cmd == sensorId; }); } + /// Returns the shared OpenRing sensor data stream. + /// + /// If the underlying BLE subscription setup fails, the cached future is + /// cleared so a later call can retry instead of reusing a failed future. + Future>> _getSensorDataStream() async { + final existingStream = _sensorDataStream; + if (existingStream != null) { + return existingStream; + } + + final streamFuture = _createSensorDataStream(); + _sensorDataStream = streamFuture; + + try { + return await streamFuture; + } catch (_) { + if (identical(_sensorDataStream, streamFuture)) { + _sensorDataStream = null; + } + rethrow; + } + } + @override Future writeSensorConfig(OpenRingSensorConfig sensorConfig) async { if (!_bleManager.isConnected(_discoveredDevice.id)) { From a6b44a4c196ba5445ee2dff2bd17e36ba9a49245 Mon Sep 17 00:00:00 2001 From: Dennis Moschina <45356478+DennisMoschina@users.noreply.github.com> Date: Wed, 22 Jul 2026 14:02:21 +0200 Subject: [PATCH 06/11] fix: prevent adding data to closed stream controllers across sensor handlers --- .../managers/open_earable_sensor_manager.dart | 8 +- .../managers/open_ring_sensor_handler.dart | 8 +- lib/src/managers/v2_sensor_handler.dart | 8 +- .../battery_energy_status_gatt_reader.dart | 26 ++++--- .../battery_health_status_gatt_reader.dart | 26 ++++--- .../battery_level_status_gatt_reader.dart | 26 ++++--- ...tery_level_status_service_gatt_reader.dart | 26 ++++--- lib/src/models/devices/cosinuss_one.dart | 74 +++++++++++-------- lib/src/models/devices/esense_factory.dart | 4 +- .../models/devices/open_earable_factory.dart | 4 +- lib/src/models/devices/open_earable_v1.dart | 29 ++++++-- lib/src/models/devices/open_earable_v2.dart | 64 +++++++++------- lib/src/models/devices/polar_factory.dart | 18 +++-- 13 files changed, 196 insertions(+), 125 deletions(-) diff --git a/lib/src/managers/open_earable_sensor_manager.dart b/lib/src/managers/open_earable_sensor_manager.dart index 438d0703..761acc12 100644 --- a/lib/src/managers/open_earable_sensor_manager.dart +++ b/lib/src/managers/open_earable_sensor_manager.dart @@ -130,12 +130,16 @@ class OpenEarableSensorHandler extends SensorHandler { parsedData["EULER"] ["units"] = {"YAW": "rad", "PITCH": "rad", "ROLL": "rad"}; } - streamController.add(parsedData); + if (!streamController.isClosed) { + streamController.add(parsedData); + } } } }, onError: (error) { - streamController.addError(error); + if (!streamController.isClosed) { + streamController.addError(error); + } }, ); diff --git a/lib/src/managers/open_ring_sensor_handler.dart b/lib/src/managers/open_ring_sensor_handler.dart index 5b3098d6..ae66ca05 100644 --- a/lib/src/managers/open_ring_sensor_handler.dart +++ b/lib/src/managers/open_ring_sensor_handler.dart @@ -294,7 +294,9 @@ class OpenRingSensorHandler extends SensorHandler { final imuAlias = _createImuAliasFromPpg(filtered); _removeImuPayload(filtered); _emitIfSampleHasSensorPayload(streamController, filtered); - streamController.add(imuAlias); + if (!streamController.isClosed) { + streamController.add(imuAlias); + } return; } @@ -321,7 +323,9 @@ class OpenRingSensorHandler extends SensorHandler { if (!_hasAnySensorPayload(sample)) { return; } - streamController.add(sample); + if (!streamController.isClosed) { + streamController.add(sample); + } } bool _hasImuPayload(Map sample) { diff --git a/lib/src/managers/v2_sensor_handler.dart b/lib/src/managers/v2_sensor_handler.dart index 74bdefde..a8fd9427 100644 --- a/lib/src/managers/v2_sensor_handler.dart +++ b/lib/src/managers/v2_sensor_handler.dart @@ -50,13 +50,17 @@ class V2SensorHandler extends SensorHandler { if (data.isNotEmpty && data[0] == sensorId) { List> parsedData = await _parseData(data); for (var d in parsedData) { - streamController.add(d); + if (!streamController.isClosed) { + streamController.add(d); + } } } }, onError: (error) { logger.e("Error while subscribing to sensor data: $error"); - streamController.addError(error); + if (!streamController.isClosed) { + streamController.addError(error); + } }, ); diff --git a/lib/src/models/devices/battery_gatt_reader/battery_energy_status_gatt_reader.dart b/lib/src/models/devices/battery_gatt_reader/battery_energy_status_gatt_reader.dart index 83663826..bb95b28f 100644 --- a/lib/src/models/devices/battery_gatt_reader/battery_energy_status_gatt_reader.dart +++ b/lib/src/models/devices/battery_gatt_reader/battery_energy_status_gatt_reader.dart @@ -8,7 +8,8 @@ import '../bluetooth_wearable.dart'; const String _batteryEnergyStatusCharacteristicUuid = "2BF0"; const String _batteryServiceUuid = "180F"; -mixin BatteryEnergyStatusGattReader on BluetoothWearable implements BatteryEnergyStatusService { +mixin BatteryEnergyStatusGattReader on BluetoothWearable + implements BatteryEnergyStatusService { @override Future readEnergyStatus() async { List energyStatusList = await bleManager.read( @@ -63,24 +64,27 @@ mixin BatteryEnergyStatusGattReader on BluetoothWearable implements BatteryEnerg StreamController(); Timer? energyPollingTimer; + Future pollEnergyStatus() async { + try { + final energyStatus = await readEnergyStatus(); + if (!controller.isClosed) { + controller.add(energyStatus); + } + } catch (e) { + logger.e('Error reading energy status: $e'); + } + } + controller.onCancel = () { energyPollingTimer?.cancel(); }; controller.onListen = () { energyPollingTimer = Timer.periodic(const Duration(seconds: 5), (timer) { - readEnergyStatus().then((energyStatus) { - controller.add(energyStatus); - }).catchError((e) { - logger.e('Error reading energy status: $e'); - }); + unawaited(pollEnergyStatus()); }); - readEnergyStatus().then((energyStatus) { - controller.add(energyStatus); - }).catchError((e) { - logger.e('Error reading energy status: $e'); - }); + unawaited(pollEnergyStatus()); }; return controller.stream; diff --git a/lib/src/models/devices/battery_gatt_reader/battery_health_status_gatt_reader.dart b/lib/src/models/devices/battery_gatt_reader/battery_health_status_gatt_reader.dart index d652537a..83ac4424 100644 --- a/lib/src/models/devices/battery_gatt_reader/battery_health_status_gatt_reader.dart +++ b/lib/src/models/devices/battery_gatt_reader/battery_health_status_gatt_reader.dart @@ -7,7 +7,8 @@ import '../bluetooth_wearable.dart'; const String _batteryHealthStatusCharacteristicUuid = "2BEA"; const String _batteryServiceUuid = "180F"; -mixin BatteryHealthStatusGattReader on BluetoothWearable implements BatteryHealthStatusService { +mixin BatteryHealthStatusGattReader on BluetoothWearable + implements BatteryHealthStatusService { @override Future readHealthStatus() async { List healthStatusList = await bleManager.read( @@ -45,24 +46,27 @@ mixin BatteryHealthStatusGattReader on BluetoothWearable implements BatteryHealt StreamController(); Timer? healthPollingTimer; + Future pollHealthStatus() async { + try { + final healthStatus = await readHealthStatus(); + if (!controller.isClosed) { + controller.add(healthStatus); + } + } catch (e) { + logger.e('Error reading health status: $e'); + } + } + controller.onCancel = () { healthPollingTimer?.cancel(); }; controller.onListen = () { healthPollingTimer = Timer.periodic(const Duration(seconds: 5), (timer) { - readHealthStatus().then((healthStatus) { - controller.add(healthStatus); - }).catchError((e) { - logger.e('Error reading health status: $e'); - }); + unawaited(pollHealthStatus()); }); - readHealthStatus().then((healthStatus) { - controller.add(healthStatus); - }).catchError((e) { - logger.e('Error reading health status: $e'); - }); + unawaited(pollHealthStatus()); }; return controller.stream; diff --git a/lib/src/models/devices/battery_gatt_reader/battery_level_status_gatt_reader.dart b/lib/src/models/devices/battery_gatt_reader/battery_level_status_gatt_reader.dart index e648520d..4e1b6655 100644 --- a/lib/src/models/devices/battery_gatt_reader/battery_level_status_gatt_reader.dart +++ b/lib/src/models/devices/battery_gatt_reader/battery_level_status_gatt_reader.dart @@ -8,7 +8,8 @@ const String _batteryLevelCharacteristicUuid = "2A19"; const String _batteryServiceUuid = "180F"; /// Mixin that implements [BatteryLevelStatus] according to the GATT specification. -mixin BatteryLevelStatusGattReader on BluetoothWearable implements BatteryLevelStatus { +mixin BatteryLevelStatusGattReader on BluetoothWearable + implements BatteryLevelStatus { @override Future readBatteryPercentage() async { List batteryLevelList = await bleManager.read( @@ -33,24 +34,27 @@ mixin BatteryLevelStatusGattReader on BluetoothWearable implements BatteryLevelS StreamController controller = StreamController(); Timer? batteryPollingTimer; + Future pollBatteryPercentage() async { + try { + final batteryPercentage = await readBatteryPercentage(); + if (!controller.isClosed) { + controller.add(batteryPercentage); + } + } catch (e) { + logger.e('Error reading battery percentage: $e'); + } + } + controller.onCancel = () { batteryPollingTimer?.cancel(); }; controller.onListen = () { batteryPollingTimer = Timer.periodic(const Duration(seconds: 5), (timer) { - readBatteryPercentage().then((batteryPercentage) { - controller.add(batteryPercentage); - }).catchError((e) { - logger.e('Error reading battery percentage: $e'); - }); + unawaited(pollBatteryPercentage()); }); - readBatteryPercentage().then((batteryPercentage) { - controller.add(batteryPercentage); - }).catchError((e) { - logger.e('Error reading battery percentage: $e'); - }); + unawaited(pollBatteryPercentage()); }; return controller.stream; diff --git a/lib/src/models/devices/battery_gatt_reader/battery_level_status_service_gatt_reader.dart b/lib/src/models/devices/battery_gatt_reader/battery_level_status_service_gatt_reader.dart index 57e66e3d..0e31fa1a 100644 --- a/lib/src/models/devices/battery_gatt_reader/battery_level_status_service_gatt_reader.dart +++ b/lib/src/models/devices/battery_gatt_reader/battery_level_status_service_gatt_reader.dart @@ -7,7 +7,8 @@ import '../bluetooth_wearable.dart'; const String _batteryLevelStatusCharacteristicUuid = "2BED"; const String _batteryServiceUuid = "180F"; -mixin BatteryLevelStatusServiceGattReader on BluetoothWearable implements BatteryLevelStatusService { +mixin BatteryLevelStatusServiceGattReader on BluetoothWearable + implements BatteryLevelStatusService { @override Future readPowerStatus() async { List powerStateList = await bleManager.read( @@ -75,24 +76,27 @@ mixin BatteryLevelStatusServiceGattReader on BluetoothWearable implements Batter StreamController(); Timer? powerPollingTimer; + Future pollPowerStatus() async { + try { + final powerStatus = await readPowerStatus(); + if (!controller.isClosed) { + controller.add(powerStatus); + } + } catch (e) { + logger.e('Error reading power status: $e'); + } + } + controller.onCancel = () { powerPollingTimer?.cancel(); }; controller.onListen = () { powerPollingTimer = Timer.periodic(const Duration(seconds: 5), (timer) { - readPowerStatus().then((powerStatus) { - controller.add(powerStatus); - }).catchError((e) { - logger.e('Error reading power status: $e'); - }); + unawaited(pollPowerStatus()); }); - readPowerStatus().then((powerStatus) { - controller.add(powerStatus); - }).catchError((e) { - logger.e('Error reading power status: $e'); - }); + unawaited(pollPowerStatus()); }; return controller.stream; diff --git a/lib/src/models/devices/cosinuss_one.dart b/lib/src/models/devices/cosinuss_one.dart index 9609c13d..44a071b6 100644 --- a/lib/src/models/devices/cosinuss_one.dart +++ b/lib/src/models/devices/cosinuss_one.dart @@ -131,12 +131,16 @@ class CosinussOne extends Wearable return; } subscription = batteryStream.listen((data) { - streamController.add(data[0]); + if (!streamController.isClosed) { + streamController.add(data[0]); + } }); final percentage = await readBatteryPercentage(); - streamController.add(percentage); - await streamController.close(); + if (!streamController.isClosed) { + streamController.add(percentage); + await streamController.close(); + } } catch (error, stack) { if (!streamController.isClosed) { streamController.addError(error, stack); @@ -231,12 +235,14 @@ class _CosinussOneSensor extends Sensor { int accY = bytes[16]; int accZ = bytes[18]; - streamController.add( - SensorDoubleValue( - values: [accX.toDouble(), accY.toDouble(), accZ.toDouble()], - timestamp: DateTime.now().millisecondsSinceEpoch - startTime, - ), - ); + if (!streamController.isClosed) { + streamController.add( + SensorDoubleValue( + values: [accX.toDouble(), accY.toDouble(), accZ.toDouble()], + timestamp: DateTime.now().millisecondsSinceEpoch - startTime, + ), + ); + } }); await _bleManager.write( @@ -299,16 +305,18 @@ class _CosinussOneSensor extends Sensor { bytes[11] << 32; // ambient light sensor (e.g., if sensor is not placed correctly) - streamController.add( - SensorDoubleValue( - values: [ - ppgRed.toDouble(), - ppgGreen.toDouble(), - ppgGreenAmbient.toDouble(), - ], - timestamp: DateTime.now().millisecondsSinceEpoch - startTime, - ), - ); + if (!streamController.isClosed) { + streamController.add( + SensorDoubleValue( + values: [ + ppgRed.toDouble(), + ppgGreen.toDouble(), + ppgGreenAmbient.toDouble(), + ], + timestamp: DateTime.now().millisecondsSinceEpoch - startTime, + ), + ); + } }); await _bleManager.write( @@ -361,12 +369,14 @@ class _CosinussOneSensor extends Sensor { (5.0 / 9.0); // convert Fahrenheit to Celsius } - streamController.add( - SensorDoubleValue( - values: [temperature], - timestamp: DateTime.now().millisecondsSinceEpoch - startTime, - ), - ); + if (!streamController.isClosed) { + streamController.add( + SensorDoubleValue( + values: [temperature], + timestamp: DateTime.now().millisecondsSinceEpoch - startTime, + ), + ); + } }); } catch (error, stack) { if (!streamController.isClosed) { @@ -437,12 +447,14 @@ class _CosinussOneHeartRateSensor extends HeartRateSensor { bpm = (((bpm >> 8) & 0xFF) | ((bpm << 8) & 0xFF00)); } - streamController.add( - HeartRateSensorValue( - heartRateBpm: bpm, - timestamp: DateTime.now().millisecondsSinceEpoch - startTime, - ), - ); + if (!streamController.isClosed) { + streamController.add( + HeartRateSensorValue( + heartRateBpm: bpm, + timestamp: DateTime.now().millisecondsSinceEpoch - startTime, + ), + ); + } }); } catch (error, stack) { if (!streamController.isClosed) { diff --git a/lib/src/models/devices/esense_factory.dart b/lib/src/models/devices/esense_factory.dart index c0d4fe0d..f2518368 100644 --- a/lib/src/models/devices/esense_factory.dart +++ b/lib/src/models/devices/esense_factory.dart @@ -163,7 +163,9 @@ class EsenseSensor extends Sensor { timestamp: timestamp, ); - streamController.add(sensorValue); + if (!streamController.isClosed) { + streamController.add(sensorValue); + } }, ); } catch (error, stack) { diff --git a/lib/src/models/devices/open_earable_factory.dart b/lib/src/models/devices/open_earable_factory.dart index dc0f95a3..fd8db736 100644 --- a/lib/src/models/devices/open_earable_factory.dart +++ b/lib/src/models/devices/open_earable_factory.dart @@ -380,7 +380,9 @@ class _OpenEarableSensorV2 extends Sensor { timestamp: timestamp, ); - streamController.add(sensorValue); + if (!streamController.isClosed) { + streamController.add(sensorValue); + } }); } catch (error, stack) { if (!streamController.isClosed) { diff --git a/lib/src/models/devices/open_earable_v1.dart b/lib/src/models/devices/open_earable_v1.dart index 74d0baf8..b10350c0 100644 --- a/lib/src/models/devices/open_earable_v1.dart +++ b/lib/src/models/devices/open_earable_v1.dart @@ -406,17 +406,28 @@ class OpenEarableV1 extends Wearable pollingTimer = Timer.periodic(const Duration(seconds: 5), (timer) async { try { int batteryPercentage = await readBatteryPercentage(); - controller.add(batteryPercentage); + if (!controller.isClosed) { + controller.add(batteryPercentage); + } } catch (e) { + if (!controller.isClosed) { + controller.addError(e); + } + } + }); + + readBatteryPercentage().then((batteryPercentage) { + if (!controller.isClosed) { + controller.add(batteryPercentage); + } + }).catchError((e) { + logger.e('Error reading battery percentage: $e'); + if (!controller.isClosed) { controller.addError(e); } }); }; - readBatteryPercentage().then(controller.add).catchError((e) { - logger.e('Error reading battery percentage: $e'); - }); - return controller.stream; } @@ -504,7 +515,9 @@ class _OpenEarableSensor extends Sensor { timestamp: timestamp, ); - streamController.add(sensorValue); + if (!streamController.isClosed) { + streamController.add(sensorValue); + } }); } catch (error, stack) { if (!streamController.isClosed) { @@ -542,7 +555,9 @@ class _OpenEarableSensor extends Sensor { timestamp: timestamp, ); - streamController.add(sensorValue); + if (!streamController.isClosed) { + streamController.add(sensorValue); + } }); } catch (error, stack) { if (!streamController.isClosed) { diff --git a/lib/src/models/devices/open_earable_v2.dart b/lib/src/models/devices/open_earable_v2.dart index 7067c5d2..deb7a063 100644 --- a/lib/src/models/devices/open_earable_v2.dart +++ b/lib/src/models/devices/open_earable_v2.dart @@ -128,11 +128,15 @@ class OpenEarableV2 extends BluetoothWearable } _sensorConfigSubscription = configStream.listen( (data) { - controller.add(_parseConfigMap(data)); + if (!controller.isClosed) { + controller.add(_parseConfigMap(data)); + } }, onError: (error) { logger.e('Error in sensor configuration stream: $error'); - controller.addError(error); + if (!controller.isClosed) { + controller.addError(error); + } }, ); } catch (error, stack) { @@ -142,19 +146,21 @@ class OpenEarableV2 extends BluetoothWearable } } - // Immediately read the current sensor configuration - bleManager - .read( - deviceId: deviceId, - serviceId: sensorServiceUuid, - characteristicId: sensorConfigStateCharacteristicUuid, - ) - .then((data) { - controller.add(_parseConfigMap(data)); - }).catchError((error) { + try { + final data = await bleManager.read( + deviceId: deviceId, + serviceId: sensorServiceUuid, + characteristicId: sensorConfigStateCharacteristicUuid, + ); + if (!controller.isClosed) { + controller.add(_parseConfigMap(data)); + } + } catch (error, stack) { logger.e('Error reading initial sensor configuration: $error'); - controller.addError(error); - }); + if (!controller.isClosed) { + controller.addError(error, stack); + } + } }; return controller.stream; } @@ -250,7 +256,7 @@ class OpenEarableV2 extends BluetoothWearable } _buttonSubscription = buttonStream.listen( (data) { - if (data.isNotEmpty) { + if (!controller.isClosed && data.isNotEmpty) { int buttonState = data[0]; if (buttonState == 0) { controller.add(ButtonEvent.released); @@ -261,7 +267,9 @@ class OpenEarableV2 extends BluetoothWearable }, onError: (error) { logger.e('Error in button events stream: $error'); - controller.addError(error); + if (!controller.isClosed) { + controller.addError(error); + } }, ); } catch (error, stack) { @@ -271,15 +279,13 @@ class OpenEarableV2 extends BluetoothWearable } } - // Immediately read current button state - bleManager - .read( - deviceId: deviceId, - serviceId: _buttonServiceUuid, - characteristicId: _buttonCharacteristicUuid, - ) - .then((data) { - if (data.isNotEmpty) { + try { + final data = await bleManager.read( + deviceId: deviceId, + serviceId: _buttonServiceUuid, + characteristicId: _buttonCharacteristicUuid, + ); + if (!controller.isClosed && data.isNotEmpty) { int buttonState = data[0]; if (buttonState == 0) { controller.add(ButtonEvent.released); @@ -287,10 +293,12 @@ class OpenEarableV2 extends BluetoothWearable controller.add(ButtonEvent.pressed); } } - }).catchError((error) { + } catch (error, stack) { logger.e('Error reading initial button state: $error'); - controller.addError(error); - }); + if (!controller.isClosed) { + controller.addError(error, stack); + } + } }; return controller.stream; diff --git a/lib/src/models/devices/polar_factory.dart b/lib/src/models/devices/polar_factory.dart index 961e46ea..e2edb00d 100644 --- a/lib/src/models/devices/polar_factory.dart +++ b/lib/src/models/devices/polar_factory.dart @@ -119,12 +119,14 @@ class _PolarHeartRateSensor extends HeartRateSensor { ? (bytes[1] & 0xFF) | ((bytes[2] & 0xFF) << 8) : bytes[1] & 0xFF; - streamController.add( - HeartRateSensorValue( - heartRateBpm: heartRate, - timestamp: DateTime.now().millisecondsSinceEpoch - startTime, - ), - ); + if (!streamController.isClosed) { + streamController.add( + HeartRateSensorValue( + heartRateBpm: heartRate, + timestamp: DateTime.now().millisecondsSinceEpoch - startTime, + ), + ); + } }); } catch (error, stack) { if (!streamController.isClosed) { @@ -189,7 +191,9 @@ class _PolarHeartRateVariabilitySensor extends HeartRateVariabilitySensor { rrIntervalsMs.add(_mapRr1024ToRrMs(rrValue)); } - streamController.add(rrIntervalsMs); + if (!streamController.isClosed) { + streamController.add(rrIntervalsMs); + } } }); } catch (error, stack) { From 19087090fef662d95f00f6dd7736b716fe25be17 Mon Sep 17 00:00:00 2001 From: Dennis Moschina <45356478+DennisMoschina@users.noreply.github.com> Date: Wed, 22 Jul 2026 15:36:26 +0200 Subject: [PATCH 07/11] fix(v2_sensor_scheme_reader): enhance sensor scheme request handling with notification support and fallback --- .../v2_sensor_scheme_reader.dart | 437 ++++++++++++++++-- 1 file changed, 392 insertions(+), 45 deletions(-) diff --git a/lib/src/utils/sensor_scheme_parser/v2_sensor_scheme_reader.dart b/lib/src/utils/sensor_scheme_parser/v2_sensor_scheme_reader.dart index 0e2d781f..4b642320 100644 --- a/lib/src/utils/sensor_scheme_parser/v2_sensor_scheme_reader.dart +++ b/lib/src/utils/sensor_scheme_parser/v2_sensor_scheme_reader.dart @@ -7,17 +7,49 @@ import 'package:open_earable_flutter/open_earable_flutter.dart' show logger; import 'package:open_earable_flutter/src/constants.dart'; import '../../managers/ble_gatt_manager.dart'; +import 'edge_ml_sensor_scheme_reader.dart'; import 'sensor_scheme_reader.dart'; +/// Reads OpenEarable V2 sensor schemes from the parse-info GATT service. +/// +/// Newer firmware exposes a per-sensor request protocol: write a sensor id to +/// [requestSensorSchemeCharacteristicUuid], then receive the requested scheme +/// on [sensorSchemeCharacteristicUuid]. Some devices instead update the same +/// characteristic for synchronous reads, so this reader first waits for the +/// notification response and then falls back to a direct read if the +/// notification does not arrive. class V2SensorSchemeReader extends SensorSchemeReader { + /// Maximum time to wait for the notification-based scheme response. + static const Duration _sensorSchemeResponseTimeout = Duration(seconds: 5); + + /// Small delay after enabling notifications and before sending requests. + /// + /// Some BLE stacks report subscription setup complete before firmware is + /// ready to process the following request write. + static const Duration _sensorSchemeRequestSettleDelay = + Duration(milliseconds: 80); + + /// Limit for out-of-order notification responses kept in memory. + static const int _maxBufferedSensorSchemeResponses = 12; + final String _deviceId; final BleGattManager _bleManager; + /// Sensor schemes that were already read, keyed by sensor id. final Map _sensorSchemes = {}; + + /// Sensor ids reported by [sensorListCharacteristicUuid]. final List _sensorIds = []; + /// Creates a reader for the connected device represented by [_deviceId]. V2SensorSchemeReader(this._bleManager, this._deviceId); + /// Reads and caches the list of sensor ids exposed by the device. + /// + /// The first byte of the characteristic is the advertised count; following + /// bytes are the sensor ids. If the characteristic is shorter than reported, + /// the available ids are still used so partially compatible firmware can + /// continue to work. Future _readSensorIds() async { List sensorIdBuffer = await _bleManager.read( deviceId: _deviceId, @@ -59,46 +91,25 @@ class V2SensorSchemeReader extends SensorSchemeReader { return _sensorSchemes[sensorId]!; } - // Listen to the notification of the characteristic - final Stream> stream = await _bleManager.subscribe( + final notificationStream = await _bleManager.subscribe( deviceId: _deviceId, serviceId: parseInfoServiceUuid, characteristicId: sensorSchemeCharacteristicUuid, ); - - final Future> responseFuture = - stream.cast>().first.timeout(const Duration(seconds: 5)); - - // Request sensor value only after the listener/future is set up - await _bleManager.write( - deviceId: _deviceId, - serviceId: parseInfoServiceUuid, - characteristicId: requestSensorSchemeCharacteristicUuid, - byteData: [sensorId], + final responseReader = _SensorSchemeResponseReader( + notificationStream: notificationStream, + parseSensorScheme: _parseSensorScheme, ); try { - final value = await responseFuture; - logger.d( - "Received notification for sensor scheme of sensor $sensorId: $value", + final scheme = await _requestSensorScheme( + sensorId: sensorId, + responseReader: responseReader, ); - - final scheme = _parseSensorScheme(value); - if (scheme.sensorId == 0 && sensorId != 0) { - logger.w( - "Sensor scheme response for sensor $sensorId omitted the sensor id. Using the requested id.", - ); - scheme.sensorId = sensorId; - } else if (scheme.sensorId != sensorId) { - logger.w( - "Sensor scheme response for sensor $sensorId reported sensor id ${scheme.sensorId}. Using the returned scheme.", - ); - } - _sensorSchemes[scheme.sensorId] = scheme; return scheme; - } on TimeoutException catch (e) { - throw TimeoutException("Timeout while waiting for sensor scheme: $e"); + } finally { + await responseReader.dispose(); } } @@ -108,25 +119,50 @@ class V2SensorSchemeReader extends SensorSchemeReader { await _readSensorIds(); } - for (int sensorId in _sensorIds) { - if (!_sensorSchemes.containsKey(sensorId) || forceRead) { - try { - SensorScheme scheme = await getSchemeForSensor(sensorId); - _sensorSchemes[scheme.sensorId] = scheme; - } catch (e) { - logger.e( - "Failed to read sensor scheme for sensor $sensorId: $e${kIsWeb ? ' (on web platform)' : ''}", - ); - if (kIsWeb) { - logger.d( - "Skipping sensor $sensorId due to read failure on web. " - "This may be a BLE notification timeout or subscription issue.", + if (!forceRead && _sensorIds.every(_sensorSchemes.containsKey)) { + return _sensorSchemes.values.toList(); + } + + final sensorSchemeStream = await _bleManager.subscribe( + deviceId: _deviceId, + serviceId: parseInfoServiceUuid, + characteristicId: sensorSchemeCharacteristicUuid, + ); + final responseReader = _SensorSchemeResponseReader( + notificationStream: sensorSchemeStream, + parseSensorScheme: _parseSensorScheme, + ); + + try { + for (int sensorId in _sensorIds) { + if (!_sensorSchemes.containsKey(sensorId) || forceRead) { + try { + SensorScheme scheme = await _requestSensorScheme( + sensorId: sensorId, + responseReader: responseReader, + ); + _sensorSchemes[scheme.sensorId] = scheme; + } catch (e) { + logger.e( + "Failed to read sensor scheme for sensor $sensorId: $e${kIsWeb ? ' (on web platform)' : ''}", ); + if (kIsWeb) { + logger.d( + "Skipping sensor $sensorId due to read failure on web. " + "This may be a BLE notification timeout or subscription issue.", + ); + } + // Continue with next sensor instead of failing entirely + continue; } - // Continue with next sensor instead of failing entirely - continue; } } + } finally { + await responseReader.dispose(); + } + + if (_sensorSchemes.isEmpty) { + await _readLegacySensorSchemesFallback(forceRead: forceRead); } logger.d( @@ -136,6 +172,181 @@ class V2SensorSchemeReader extends SensorSchemeReader { return _sensorSchemes.values.toList(); } + /// Requests one sensor scheme using notification first, direct read second. + /// + /// The notification path is the preferred V2 protocol. If that path times out + /// or fails, the method sends the same request again and reads + /// [sensorSchemeCharacteristicUuid] synchronously. This supports firmware + /// variants that update the characteristic value but do not emit the expected + /// notification. + Future _requestSensorScheme({ + required int sensorId, + required _SensorSchemeResponseReader responseReader, + }) async { + try { + return await _requestSensorSchemeViaNotification( + sensorId: sensorId, + responseReader: responseReader, + ); + } catch (error) { + logger.w( + "Notification based sensor scheme request for sensor $sensorId failed: " + "$error. Trying synchronous read.", + ); + } + + return _requestSensorSchemeViaRead(sensorId); + } + + /// Sends a scheme request and waits for the matching notification response. + /// + /// [responseReader] owns the active notification subscription and matches the + /// next parsed scheme by sensor id. The pending response is explicitly + /// cancelled when write or timeout errors occur, so the response reader can be + /// reused for the next sensor id. + Future _requestSensorSchemeViaNotification({ + required int sensorId, + required _SensorSchemeResponseReader responseReader, + }) async { + final responseFuture = responseReader.nextMatchingScheme( + sensorId, + timeout: _sensorSchemeResponseTimeout, + ); + + try { + await Future.delayed(_sensorSchemeRequestSettleDelay); + + await _writeSensorSchemeRequest(sensorId); + + final scheme = await responseFuture; + logger.d( + "Received notification for sensor scheme of sensor $sensorId: $scheme", + ); + + return _normalizeRequestedSensorScheme( + requestedSensorId: sensorId, + scheme: scheme, + ); + } catch (error, stack) { + responseReader.cancelPendingResponse( + sensorId: sensorId, + error: error, + stackTrace: stack, + ); + try { + await responseFuture; + } catch (_) { + // The response future is intentionally drained after cancellation. + } + rethrow; + } + } + + /// Sends a scheme request and reads the scheme characteristic directly. + /// + /// This is the compatibility fallback for devices where the request updates + /// [sensorSchemeCharacteristicUuid] but no notification arrives. + Future _requestSensorSchemeViaRead(int sensorId) async { + await Future.delayed(_sensorSchemeRequestSettleDelay); + + await _writeSensorSchemeRequest(sensorId); + + final value = await _bleManager.read( + deviceId: _deviceId, + serviceId: parseInfoServiceUuid, + characteristicId: sensorSchemeCharacteristicUuid, + ); + final scheme = _parseSensorScheme(value); + logger.d("Read sensor scheme for sensor $sensorId synchronously: $scheme"); + + return _normalizeRequestedSensorScheme( + requestedSensorId: sensorId, + scheme: scheme, + allowOmittedSensorId: true, + throwOnMismatch: true, + ); + } + + /// Writes the requested sensor id to the V2 scheme request characteristic. + Future _writeSensorSchemeRequest(int sensorId) { + return _bleManager.write( + deviceId: _deviceId, + serviceId: parseInfoServiceUuid, + characteristicId: requestSensorSchemeCharacteristicUuid, + byteData: [sensorId], + ); + } + + /// Normalizes known firmware quirks in returned sensor schemes. + /// + /// Some responses omit the requested id and report `0`. In that case the + /// requested id is applied so callers can cache and match the scheme + /// correctly. Mismatched non-zero ids are logged, or rejected when + /// [throwOnMismatch] is enabled for stale-read-sensitive paths. + SensorScheme _normalizeRequestedSensorScheme({ + required int requestedSensorId, + required SensorScheme scheme, + bool allowOmittedSensorId = true, + bool throwOnMismatch = false, + }) { + if (scheme.sensorId == 0 && + requestedSensorId != 0 && + allowOmittedSensorId) { + logger.w( + "Sensor scheme response for sensor $requestedSensorId omitted the sensor id. Using the requested id.", + ); + scheme.sensorId = requestedSensorId; + } else if (scheme.sensorId != requestedSensorId) { + if (throwOnMismatch) { + throw StateError( + "Sensor scheme response for sensor $requestedSensorId reported " + "sensor id ${scheme.sensorId}. Refusing to cache mismatched scheme.", + ); + } + logger.w( + "Sensor scheme response for sensor $requestedSensorId reported sensor id ${scheme.sensorId}. Using the returned scheme.", + ); + } + + return scheme; + } + + /// Attempts the older full-scheme characteristic format as a final fallback. + /// + /// This is only used when no V2 per-sensor scheme could be read at all. + /// Devices with older firmware can still be initialized if they expose the + /// legacy [schemeCharacteristicUuid] payload. + Future _readLegacySensorSchemesFallback({ + required bool forceRead, + }) async { + logger.w( + "Per-sensor V2 scheme reads failed for $_deviceId. " + "Trying legacy full sensor scheme characteristic.", + ); + + try { + final legacySchemes = await EdgeMlSensorSchemeReader( + _bleManager, + _deviceId, + ).readSensorSchemes(forceRead: forceRead); + _sensorSchemes + ..clear() + ..addEntries( + legacySchemes.map((scheme) => MapEntry(scheme.sensorId, scheme)), + ); + logger.i( + "Loaded ${legacySchemes.length} sensor scheme(s) from legacy fallback.", + ); + } catch (error) { + logger.e("Legacy sensor scheme fallback failed: $error"); + } + } + + /// Parses a single V2 sensor scheme payload. + /// + /// The payload contains the sensor id, sensor name, component definitions, + /// and optional sensor configuration metadata such as supported features and + /// frequency definitions. SensorScheme _parseSensorScheme(List byteStream) { int currentIndex = 0; int sensorId = byteStream[currentIndex++]; @@ -224,3 +435,139 @@ class V2SensorSchemeReader extends SensorSchemeReader { return sensorScheme; } } + +/// Tracks a single outstanding notification response request. +class _PendingSensorSchemeResponse { + _PendingSensorSchemeResponse({ + required this.sensorId, + required this.completer, + }); + + final int sensorId; + final Completer completer; +} + +/// Reads and matches sensor scheme notifications from one active subscription. +/// +/// The V2 reader keeps one notification subscription open while requesting all +/// schemes. Notifications can arrive slightly out of order, so unmatched +/// schemes are buffered and checked before waiting for a new response. +class _SensorSchemeResponseReader { + _SensorSchemeResponseReader({ + required Stream> notificationStream, + required SensorScheme Function(List) parseSensorScheme, + }) : _parseSensorScheme = parseSensorScheme { + _subscription = notificationStream.listen( + _handleNotification, + onError: (error, stack) { + final pending = _pending; + if (pending != null && !pending.completer.isCompleted) { + pending.completer.completeError(error, stack); + } + }, + ); + } + + final SensorScheme Function(List) _parseSensorScheme; + + /// Parsed notifications that did not match the request pending at the time. + final List _bufferedSchemes = []; + + /// The currently awaited sensor scheme response, if any. + _PendingSensorSchemeResponse? _pending; + + /// Subscription to [sensorSchemeCharacteristicUuid] notifications. + late final StreamSubscription> _subscription; + + /// Returns the next buffered or future scheme matching [sensorId]. + Future nextMatchingScheme( + int sensorId, { + required Duration timeout, + }) { + final bufferedIndex = _bufferedSchemes.indexWhere( + (scheme) => _matchesSensorId(scheme, sensorId), + ); + if (bufferedIndex != -1) { + return Future.value(_bufferedSchemes.removeAt(bufferedIndex)); + } + + final completer = Completer(); + _pending = _PendingSensorSchemeResponse( + sensorId: sensorId, + completer: completer, + ); + + return completer.future.timeout( + timeout, + onTimeout: () { + if (identical(_pending?.completer, completer)) { + _pending = null; + } + throw TimeoutException( + "Timeout while waiting for sensor scheme response for sensor $sensorId", + timeout, + ); + }, + ); + } + + /// Cancels the notification subscription and clears transient state. + Future dispose() { + _pending = null; + _bufferedSchemes.clear(); + return _subscription.cancel(); + } + + /// Completes and clears the pending response after request-side failure. + /// + /// This prevents an abandoned pending completer from timing out later after + /// the caller has already moved to the synchronous read fallback. + void cancelPendingResponse({ + required int sensorId, + required Object error, + required StackTrace stackTrace, + }) { + final pending = _pending; + if (pending == null || pending.sensorId != sensorId) { + return; + } + _pending = null; + if (!pending.completer.isCompleted) { + pending.completer.completeError(error, stackTrace); + } + } + + /// Parses and routes one raw notification payload. + /// + /// Matching schemes complete the current pending request. Non-matching + /// schemes are buffered because they may belong to a later request. + void _handleNotification(List value) { + SensorScheme scheme; + try { + scheme = _parseSensorScheme(value); + } catch (error) { + logger.w("Ignoring malformed sensor scheme notification: $error"); + return; + } + + final pending = _pending; + if (pending != null && _matchesSensorId(scheme, pending.sensorId)) { + _pending = null; + if (!pending.completer.isCompleted) { + pending.completer.complete(scheme); + } + return; + } + + _bufferedSchemes.add(scheme); + if (_bufferedSchemes.length > + V2SensorSchemeReader._maxBufferedSensorSchemeResponses) { + _bufferedSchemes.removeAt(0); + } + } + + /// Returns whether [scheme] can satisfy a request for [requestedSensorId]. + bool _matchesSensorId(SensorScheme scheme, int requestedSensorId) { + return scheme.sensorId == requestedSensorId; + } +} From 41904bf85539839df731d58a72008146181f1e36 Mon Sep 17 00:00:00 2001 From: Dennis Moschina <45356478+DennisMoschina@users.noreply.github.com> Date: Wed, 22 Jul 2026 15:41:08 +0200 Subject: [PATCH 08/11] chore(changelog): document BLE race fixes --- CHANGELOG.md | 3 +++ 1 file changed, 3 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 349c5cbd..f86b5e61 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -3,6 +3,9 @@ * BREAKING CHANGE: `BleGattManager.subscribe` now returns `Future>>`, so callers must `await` subscription setup before listening to BLE notifications. * BREAKING CHANGE: `SensorHandler.subscribeToSensorData` now returns `Future>>`, so callers must `await` sensor notification readiness before listening to sensor data. * fixed BLE notification setup races by ensuring subscription futures complete only after the underlying GATT notification subscription is enabled. +* fixed OpenEarable V2 sensor scheme loading on devices that update the scheme characteristic without sending a notification. +* improved OpenEarable V2 sensor scheme loading reliability by keeping one notification subscription active while reading schemes and falling back to synchronous reads. +* fixed race conditions in BLE unsubscribe, duplicate device connection, and stream cancellation handling. ## 2.3.10 From d1fc2e23b7c460cd201b70933c4413daa2373bde Mon Sep 17 00:00:00 2001 From: Dennis <45356478+DennisMoschina@users.noreply.github.com> Date: Fri, 24 Jul 2026 11:41:16 +0200 Subject: [PATCH 09/11] fix(open_earable_sensor_manager): throw exception in subscribeToSensorData when device is not connected Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- lib/src/managers/open_earable_sensor_manager.dart | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/src/managers/open_earable_sensor_manager.dart b/lib/src/managers/open_earable_sensor_manager.dart index 761acc12..b02e2df2 100644 --- a/lib/src/managers/open_earable_sensor_manager.dart +++ b/lib/src/managers/open_earable_sensor_manager.dart @@ -68,7 +68,7 @@ class OpenEarableSensorHandler extends SensorHandler { int sensorId, ) async { if (!_bleManager.isConnected(deviceId)) { - Exception("Can't subscribe to sensor data. Earable not connected"); + throw Exception("Can't subscribe to sensor data. Earable not connected"); } StreamController> streamController = StreamController(); From 7133d7c957935022e3abeb2b18f02108c90f6f58 Mon Sep 17 00:00:00 2001 From: Dennis <45356478+DennisMoschina@users.noreply.github.com> Date: Fri, 24 Jul 2026 11:43:02 +0200 Subject: [PATCH 10/11] fix(ble_manager): clean up callbacks after connection error Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- lib/src/managers/ble_manager.dart | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/lib/src/managers/ble_manager.dart b/lib/src/managers/ble_manager.dart index 21bfe7cb..e5cd7725 100644 --- a/lib/src/managers/ble_manager.dart +++ b/lib/src/managers/ble_manager.dart @@ -266,6 +266,10 @@ class BleManager extends BleGattManager { try { UniversalBle.connect(device.id); } catch (error, stack) { + _connectCallbacks.remove(device.id); + _disconnectCallbacks.remove(device.id); + _connectionCompleters.remove(device.id); + _connectionFutures.remove(device.id); if (!completer.isCompleted) { completer.completeError(error, stack); } From 8113e33b255fda32fce74e9e25bbfdb33d49c591 Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Fri, 24 Jul 2026 09:54:20 +0000 Subject: [PATCH 11/11] docs(sensor-scheme): clarify sensor-id expectation for current firmware --- .../sensor_scheme_parser/v2_sensor_scheme_reader.dart | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/lib/src/utils/sensor_scheme_parser/v2_sensor_scheme_reader.dart b/lib/src/utils/sensor_scheme_parser/v2_sensor_scheme_reader.dart index 4b642320..f8102c6b 100644 --- a/lib/src/utils/sensor_scheme_parser/v2_sensor_scheme_reader.dart +++ b/lib/src/utils/sensor_scheme_parser/v2_sensor_scheme_reader.dart @@ -279,10 +279,10 @@ class V2SensorSchemeReader extends SensorSchemeReader { /// Normalizes known firmware quirks in returned sensor schemes. /// - /// Some responses omit the requested id and report `0`. In that case the - /// requested id is applied so callers can cache and match the scheme - /// correctly. Mismatched non-zero ids are logged, or rejected when - /// [throwOnMismatch] is enabled for stale-read-sensitive paths. + /// Current firmware is expected to always return the requested sensor id. + /// The `sensorId == 0` handling is kept as a defensive compatibility fallback + /// for older firmware variants. Mismatched non-zero ids are logged, or + /// rejected when [throwOnMismatch] is enabled for stale-read-sensitive paths. SensorScheme _normalizeRequestedSensorScheme({ required int requestedSensorId, required SensorScheme scheme,