Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,12 @@
## Unreleased

* BREAKING CHANGE: `BleGattManager.subscribe` now returns `Future<Stream<List<int>>>`, so callers must `await` subscription setup before listening to BLE notifications.
* BREAKING CHANGE: `SensorHandler.subscribeToSensorData` now returns `Future<Stream<Map<String, dynamic>>>`, 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

* added dynamic power saving mode capability for OpenEarable v2 devices
Expand Down
27 changes: 27 additions & 0 deletions lib/open_earable_flutter.dart
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,7 @@ class WearableManager {
late final StreamController<DiscoveredDevice> _connectingStreamController;

final List<String> _connectedIds = [];
final Map<String, Future<Wearable>> _connectionFuturesByDeviceId = {};

List<String> _autoConnectDeviceIds = [];
StreamSubscription<DiscoveredDevice>? _autoconnectScanSubscription;
Expand Down Expand Up @@ -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<Wearable> _connectToDevice(
DiscoveredDevice device, {
required Set<ConnectionOption> options,
}) async {
_connectingStreamController.add(device);

WearableDisconnectNotifier disconnectNotifier =
Expand Down
5 changes: 4 additions & 1 deletion lib/src/managers/ble_gatt_manager.dart
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,10 @@ abstract class BleGattManager {
});

/// Subscribes to a specific characteristic of the connected device.
Stream<List<int>> subscribe({
///
/// The returned future completes only after the underlying GATT
/// notification subscription has been enabled.
Future<Stream<List<int>>> subscribe({
required String deviceId,
required String serviceId,
required String characteristicId,
Expand Down
126 changes: 99 additions & 27 deletions lib/src/managers/ble_manager.dart
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,8 @@ class BleManager extends BleGattManager {
int get mtu => _mtu;

final Map<String, StreamController<List<int>>> _streamControllers = {};
final Map<String, Future<void>> _subscriptionSetups = {};
final Map<String, Future<void>> _subscriptionTeardowns = {};

/// A stream of discovered devices during scanning.
StreamController<DiscoveredDevice>? _scanStreamController;
Expand All @@ -24,7 +26,9 @@ class BleManager extends BleGattManager {
String _getCharacteristicKey(String deviceId, String characteristicId) =>
"$deviceId||$characteristicId";

final Map<String, Completer> _connectionCompleters = {};
final Map<String, Completer<(bool, List<BleService>)>> _connectionCompleters =
{};
final Map<String, Future<(bool, List<BleService>)>> _connectionFutures = {};
final Map<String, VoidCallback> _connectCallbacks = {};
final Map<String, VoidCallback> _disconnectCallbacks = {};

Expand All @@ -49,6 +53,8 @@ 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);
_subscriptionTeardowns.remove(key);
}
}

Expand Down Expand Up @@ -214,38 +220,62 @@ 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<BleService>)> completer =
Completer<(bool, List<BleService>)>();
final completer = Completer<(bool, List<BleService>)>();
_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<BleService> 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, <BleService>[]));
_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) {
_connectCallbacks.remove(device.id);
_disconnectCallbacks.remove(device.id);
_connectionCompleters.remove(device.id);
_connectionFutures.remove(device.id);
if (!completer.isCompleted) {
completer.completeError(error, stack);
}
}
Comment thread
Copilot marked this conversation as resolved.

return completer.future;
return connectionFuture;
}

/// Checks if the connected device has a specific service.
Expand Down Expand Up @@ -312,39 +342,81 @@ class BleManager extends BleGattManager {

/// Subscribes to a specific characteristic of the connected Earable device.
@override
Stream<List<int>> subscribe({
Future<Stream<List<int>>> subscribe({
required String deviceId,
required String serviceId,
required String characteristicId,
}) {
}) async {
logger.d(
"Subscribing to $deviceId, service $serviceId, characteristic $characteristicId",
);
String streamIdentifier = _getCharacteristicKey(
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<List<int>>? streamController =
_streamControllers[streamIdentifier];
streamController ??= StreamController<List<int>>.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 = () {
streamController.onCancel = () async {
if (_streamControllers.containsKey(streamIdentifier)) {
_streamControllers.remove(streamIdentifier)?.close();
UniversalBle.unsubscribe(
final canceledController = _streamControllers.remove(streamIdentifier);
_subscriptionSetups.remove(streamIdentifier);
if (canceledController != null && !canceledController.isClosed) {
unawaited(canceledController.close());
}

final teardown = UniversalBle.unsubscribe(
deviceId,
serviceId,
characteristicId,
);
_streamControllers.remove(streamIdentifier);
_subscriptionTeardowns[streamIdentifier] = teardown;

try {
await teardown;
} catch (error) {
logger.w("Unsubscribe failed for $streamIdentifier: $error");
} finally {
if (identical(_subscriptionTeardowns[streamIdentifier], teardown)) {
_subscriptionTeardowns.remove(streamIdentifier);
}
}
}
};

Expand Down
24 changes: 14 additions & 10 deletions lib/src/managers/esense_sensor_handler.dart
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,9 @@ class EsenseSensorHandler extends SensorHandler<EsenseSensorConfig> {
}

@override
Stream<Map<String, dynamic>> subscribeToSensorData(int sensorId) {
Future<Stream<Map<String, dynamic>>> subscribeToSensorData(
int sensorId,
) async {
if (!_bleGattManager.isConnected(_discoveredDevice.id)) {
throw Exception("Can't subscribe to sensor data. Earable not connected");
}
Expand All @@ -68,13 +70,13 @@ class EsenseSensorHandler extends SensorHandler<EsenseSensorConfig> {

final streamController = StreamController<Map<String, dynamic>>();

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;

Expand Down Expand Up @@ -310,7 +312,9 @@ class EsenseSensorHandler extends SensorHandler<EsenseSensorConfig> {
);
}

logger.t("Loaded IMU ranges: Accel=$_cachedAccelRange, Gyro=$_cachedGyroRange");
logger.t(
"Loaded IMU ranges: Accel=$_cachedAccelRange, Gyro=$_cachedGyroRange",
);

return (_cachedAccelRange!, _cachedGyroRange!);
}
Expand Down Expand Up @@ -347,7 +351,7 @@ class EsenseSensorHandler extends SensorHandler<EsenseSensorConfig> {

// Make *new* mutable, dynamic-typed inner maps
final accel = Map<String, dynamic>.from(result['Accelerometer'] as Map);
final gyro = Map<String, dynamic>.from(result['Gyroscope'] as Map);
final gyro = Map<String, dynamic>.from(result['Gyroscope'] as Map);

// Accelerometer to g
for (final key in const ['x', 'y', 'z']) {
Expand Down
27 changes: 17 additions & 10 deletions lib/src/managers/open_earable_sensor_manager.dart
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,8 @@ class OpenEarableSensorHandler extends SensorHandler<OpenEarableSensorConfig> {
SensorSchemeReader? sensorSchemeParser,
SensorValueParser? sensorValueParser,
}) : _bleManager = bleManager,
_sensorSchemeParser = sensorSchemeParser ?? EdgeMlSensorSchemeReader(bleManager, deviceId),
_sensorSchemeParser = sensorSchemeParser ??
EdgeMlSensorSchemeReader(bleManager, deviceId),
_sensorValueParser = sensorValueParser ?? EdgeMlSensorValueParser() {
_readSensorScheme();
}
Expand Down Expand Up @@ -63,20 +64,22 @@ class OpenEarableSensorHandler extends SensorHandler<OpenEarableSensorConfig> {
/// - 1: Barometer data
/// Returns a [Stream] of sensor data as a [Map] of sensor values.
@override
Stream<Map<String, dynamic>> subscribeToSensorData(int sensorId) {
Future<Stream<Map<String, dynamic>>> subscribeToSensorData(
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");
}
Comment thread
Copilot marked this conversation as resolved.
StreamController<Map<String, dynamic>> 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<Map<String, dynamic>> parsedDataList = await _parseData(data);
Expand Down Expand Up @@ -127,12 +130,16 @@ class OpenEarableSensorHandler extends SensorHandler<OpenEarableSensorConfig> {
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);
}
},
);

Expand All @@ -144,7 +151,7 @@ class OpenEarableSensorHandler extends SensorHandler<OpenEarableSensorConfig> {
/// Parses raw sensor data bytes into a [Map] of sensor values.
Future<List<Map<String, dynamic>>> _parseData(List<int> data) async {
ByteData byteData = ByteData.sublistView(Uint8List.fromList(data));

return _sensorValueParser.parse(byteData, _sensorSchemes!);
}

Expand Down
Loading
Loading