You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

637 lines
19 KiB

import 'dart:async';
import 'dart:io';
import 'dart:typed_data';
import 'package:device_plugin_interface/device_plugin_interface.dart';
import 'package:get/get.dart';
import '../../core/utils/logger.dart';
import '../device_vendors.dart';
import 'bes_bluetooth_service.dart';
import 'bes_ota_port.dart';
import 'bes_protocol.dart';
/// 恒玄(BES) 耳机的厂商插件——把 [BesBluetoothService](协议实现,冻结)
/// 适配成 [DevicePlugin] / [DeviceSession]。
///
/// **这一层只做映射,不碰字节。** feature key → 命令字节的对应表在
/// [BesDeviceSession.invokeFeature],由 `test/devices/bes_session_test.dart`
/// 逐条断言发出的字节。
class BesDevicePlugin implements DevicePlugin {
static const String _tag = 'BesDevicePlugin';
/// 协议实现。仍注册进 GetX 只是为了走 GetxService 的生命周期;
/// 业务层不允许 import 它(见 device_hub_boundary_test)。
final BesBluetoothService link;
BesDevicePlugin({BesBluetoothService? link})
: link = link ?? BesBluetoothService();
final StreamController<DevicePluginEvent> _ctrl =
StreamController<DevicePluginEvent>.broadcast();
BesDeviceSession? _session;
Worker? _statusWorker;
bool _initialized = false;
bool _disposed = false;
@override
String get vendorKey => DeviceVendors.bes;
@override
String get displayName => 'BES (恒玄)';
@override
Set<DeviceCapability> get capabilities => BesDeviceSession.caps;
@override
DevicePluginConfigSchema get configSchema => DevicePluginConfigSchema.empty;
@override
DeviceSession? get activeSession => _session;
@override
Stream<DevicePluginEvent> get eventStream => _ctrl.stream;
@override
Future<void> initialize(DevicePluginConfig config) async {
_checkAlive();
if (_initialized) return;
_initialized = true;
if (!Get.isRegistered<BesBluetoothService>()) {
Get.put<BesBluetoothService>(link, permanent: true);
}
_statusWorker = ever<int>(link.deviceStatus, _onStatus);
if (link.deviceStatus.value == 2) _onStatus(2);
_emit(const DevicePluginEvent(type: DevicePluginEventType.pluginReady));
}
void _onStatus(int v) {
switch (v) {
case 2:
if (_session != null) return;
final s = BesDeviceSession._(link);
_session = s;
_emit(DevicePluginEvent(
type: DevicePluginEventType.connectionStateChanged,
deviceId: s.deviceId,
connectionState: DeviceConnectionState.ready,
));
break;
case 1:
_emit(DevicePluginEvent(
type: DevicePluginEventType.connectionStateChanged,
deviceId: link.connectedDeviceId.value,
connectionState: DeviceConnectionState.connecting,
));
break;
default:
final s = _session;
_session = null;
s?._markDisconnected();
_emit(DevicePluginEvent(
type: DevicePluginEventType.connectionStateChanged,
deviceId: s?.deviceId,
connectionState: DeviceConnectionState.disconnected,
));
}
}
void _emit(DevicePluginEvent e) {
if (!_ctrl.isClosed) _ctrl.add(e);
}
// ---------- 扫描(恒玄没有扫描:只从系统配对列表里挑) ----------
@override
Future<void> startScan({DeviceScanFilter? filter, Duration? timeout}) async {
_emit(const DevicePluginEvent(type: DevicePluginEventType.scanStarted));
for (final d in await bondedDevices()) {
_emit(DevicePluginEvent(
type: DevicePluginEventType.deviceDiscovered,
deviceId: d.id,
discovered: d));
}
_emit(const DevicePluginEvent(type: DevicePluginEventType.scanStopped));
}
@override
Future<void> stopScan() async {}
@override
Future<bool> isScanning() async => false;
@override
Future<List<DiscoveredDevice>> bondedDevices() async {
final list = await link.getBondedDevices();
return list
.map((d) => DiscoveredDevice(
id: (d['address'] ?? '') as String,
name: (d['name'] ?? '') as String,
vendor: vendorKey,
))
.where((d) => d.id.isNotEmpty)
.toList(growable: false);
}
@override
Future<bool> ensureReady() async {
try {
return await link.initialize();
} catch (e) {
Logger.e(_tag, '恒玄蓝牙初始化失败: $e');
return false;
}
}
// ---------- 连接 ----------
/// 连接并等到「已连接」。
///
/// 进过「连接中」又掉回 0 就提前判失败(原生 connectBlocking 失败会推 0);
/// 但仍要保留超时兜底:原生若把这次 connect 丢弃了(isConnecting 门禁),
/// 状态会停在 connectToDevice 自己写的 1,等不到 0。
@override
Future<DeviceSession> connect(String deviceId,
{DeviceConnectOptions? options}) async {
_checkAlive();
final name = (options?.extra['name'] as String?) ?? '';
final timeout = options?.timeout ?? const Duration(seconds: 8);
await link.connectToDevice(deviceId, name);
final end = DateTime.now().add(timeout);
var sawConnecting = false;
while (DateTime.now().isBefore(end)) {
final s = link.deviceStatus.value;
if (s == 2) break;
if (s == 1) sawConnecting = true;
if (sawConnecting && s == 0) {
throw DeviceException(DeviceErrorCode.connectFailed, deviceId);
}
await Future.delayed(const Duration(milliseconds: 300));
}
if (link.deviceStatus.value != 2) {
throw DeviceException(DeviceErrorCode.connectTimeout, deviceId);
}
// 状态流回调与这里是同一个事件循环,session 已经建好
_onStatus(2);
return _session!;
}
@override
List<DiscoveredDevice> get rememberedDevices => link
.getPairedDevices()
.map((d) => DiscoveredDevice(
id: (d['address'] ?? '') as String,
name: (d['name'] ?? '') as String,
vendor: vendorKey,
))
.where((d) => d.id.isNotEmpty)
.toList(growable: false);
@override
DiscoveredDevice? get lastKnownDevice {
final n = link.lastDeviceName;
final a = link.lastDeviceAddress;
if (n.isEmpty && a.isEmpty) return null;
return DiscoveredDevice(id: a, name: n, vendor: vendorKey);
}
@override
Future<void> rememberDevice(DiscoveredDevice device) =>
link.savePairedDevice({'address': device.id, 'name': device.name});
@override
Future<void> tryAutoConnect() => link.tryAutoConnect();
@override
Future<void> forgetDevice(String deviceId) =>
link.removePairedDevice(deviceId);
@override
Future<void> dispose() async {
_disposed = true;
_statusWorker?.dispose();
_session?._markDisconnected();
_session = null;
await _ctrl.close();
}
void _checkAlive() {
if (_disposed) throw StateError('BesDevicePlugin disposed');
}
}
/// 恒玄耳机会话。
class BesDeviceSession extends DeviceSession {
BesDeviceSession._(this._link) {
_bind();
}
static const String _tag = 'BesSession';
static const Set<DeviceCapability> caps = {
DeviceCapability.connect,
DeviceCapability.bond,
DeviceCapability.battery,
DeviceCapability.macAddress,
DeviceCapability.micUplink,
DeviceCapability.speakerDownlink,
DeviceCapability.micRecording,
DeviceCapability.callAudioTap,
DeviceCapability.callState,
DeviceCapability.aiWake,
DeviceCapability.translationDownlink,
DeviceCapability.faceToFace,
DeviceCapability.ota,
DeviceCapability.customCommand,
};
final BesBluetoothService _link;
final StreamController<DeviceSessionEvent> _events =
StreamController<DeviceSessionEvent>.broadcast();
final StreamController<DeviceWakeEvent> _wake =
StreamController<DeviceWakeEvent>.broadcast();
final List<Worker> _workers = [];
StreamSubscription? _aiSub;
DeviceConnectionState _state = DeviceConnectionState.ready;
/// 通话翻译下行的引用计数:A/B 两路都关了才 stopCallDownlink
int _callSinks = 0;
BesOtaPort? _ota;
void _bind() {
_workers.add(ever<String>(_link.lastCmdType, (type) {
switch (type) {
case 'startMicRecording':
_feature(DeviceFeatureKeys.recordStarted);
break;
case 'stopMicRecording':
_feature(DeviceFeatureKeys.recordStopped);
break;
}
}));
_workers.add(ever<bool>(_link.isCallOngoing, (on) {
_feature(on ? DeviceFeatureKeys.callStarted : DeviceFeatureKeys.callEnded);
}));
void infoChanged(dynamic _) => _infoUpdated();
_workers.add(ever<int>(_link.batteryLevel, (_) {
_feature(DeviceFeatureKeys.batteryUpdated);
_infoUpdated();
}));
_workers.add(ever<int>(_link.batteryCase, infoChanged));
_workers.add(ever<bool>(_link.chargingLeft, infoChanged));
_workers.add(ever<bool>(_link.chargingRight, infoChanged));
_workers.add(ever<bool>(_link.chargingCase, infoChanged));
_workers.add(ever<String>(_link.deviceVersion, (_) {
_feature(DeviceFeatureKeys.firmwareUpdated);
_infoUpdated();
}));
_workers.add(ever<String>(_link.deviceMac, (_) {
_feature(DeviceFeatureKeys.macUpdated);
_infoUpdated();
}));
_workers.add(ever<String>(_link.connectedDeviceName, infoChanged));
// AI 通道:BB 64 唤醒 / BB E5 结束。应答 AA 6A 由业务层通过 ai.accept 发,
// 与原先 BesAiSessionService 的判断顺序一致(未配置时不应答、发 ai.exit)。
_aiSub = _link.manager.onDeviceAiEvent().listen((event) {
switch (event['type']) {
case 'startAI':
_wake.add(DeviceWakeEvent(deviceId: deviceId, reason: WakeReason.ptt));
break;
case 'stopAI':
_wake.add(
DeviceWakeEvent(deviceId: deviceId, reason: WakeReason.hangup));
break;
}
}, onError: (e) => Logger.e(_tag, 'AI 事件流错误: $e'));
}
void _feature(String key, [Map<String, Object?> data = const {}]) {
if (_events.isClosed) return;
_events.add(DeviceSessionEvent(
type: DeviceSessionEventType.feature,
deviceId: deviceId,
feature: DeviceFeatureEvent(key: key, data: data),
));
}
void _infoUpdated() {
if (_events.isClosed) return;
_events.add(DeviceSessionEvent(
type: DeviceSessionEventType.deviceInfoUpdated,
deviceId: deviceId,
deviceInfo: info,
));
}
void _markDisconnected() {
if (_state == DeviceConnectionState.disconnected) return;
_state = DeviceConnectionState.disconnected;
for (final w in _workers) {
w.dispose();
}
_workers.clear();
_aiSub?.cancel();
_aiSub = null;
_ota?.onSessionGone();
if (!_events.isClosed) {
_events.add(DeviceSessionEvent(
type: DeviceSessionEventType.connectionStateChanged,
deviceId: deviceId,
connectionState: DeviceConnectionState.disconnected,
));
_events.close();
}
if (!_wake.isClosed) _wake.close();
}
void _requireReady() {
if (_state != DeviceConnectionState.ready) {
throw DeviceException(DeviceErrorCode.noActiveSession, 'bes');
}
}
// ---------- DeviceSession ----------
@override
String get deviceId => _link.connectedDeviceId.value;
@override
String get vendor => DeviceVendors.bes;
@override
DeviceConnectionState get state => _state;
@override
Set<DeviceCapability> get capabilities => caps;
@override
Stream<DeviceSessionEvent> get eventStream => _events.stream;
@override
Stream<DeviceWakeEvent> get wakeEvents => _wake.stream;
static int? _nz(int v) => v >= 0 ? v : null;
/// 真实 MAC:Android 系统给的就是;iOS 只能靠耳机自报(`BB 0A` 属性 7)。
String? get macAddress {
if (Platform.isAndroid) {
final id = _link.connectedDeviceId.value;
return id.isEmpty ? null : id.toUpperCase();
}
final m = _link.deviceMac.value;
return m.isEmpty ? null : m;
}
@override
DeviceInfo get info {
final battery = DeviceBattery(
left: _nz(_link.batteryLeft.value),
right: _nz(_link.batteryRight.value),
caseBox: _nz(_link.batteryCase.value),
single: _nz(_link.batteryLevel.value),
chargingLeft: _link.chargingLeft.value,
chargingRight: _link.chargingRight.value,
chargingCase: _link.chargingCase.value,
);
final fw = _link.deviceVersion.value;
return DeviceInfo(
id: deviceId,
name: _link.connectedDeviceName.value,
vendor: vendor,
firmwareVersion: fw.isEmpty ? null : fw,
batteryPercent: battery.summary,
battery: battery,
macAddress: macAddress,
metadata: {
'inCall': _link.isCallOngoing.value,
'localId': _link.connectedDeviceId.value,
},
);
}
@override
Future<int> readRssi() =>
throw DeviceException(DeviceErrorCode.notSupported, 'rssi');
@override
Future<int?> readBattery() async {
_requireReady();
await _link.queryBattery();
return info.battery.summary;
}
@override
Future<DeviceInfo> refreshInfo() async {
_requireReady();
await _link.refreshDeviceState();
return info;
}
/// 取真实 MAC,iOS 上会主动发一次查询并等最多 [timeout]。
Future<String?> resolveMac({Duration timeout = const Duration(seconds: 3)}) =>
_link.resolveDeviceMac(timeout: timeout);
// ---------- 音频 ----------
@override
Future<DeviceAudioSource> openMic({AudioRoute route = AudioRoute.mic}) async {
_requireReady();
switch (route) {
case AudioRoute.mic:
// 现场录音:主麦与前馈麦交替上行,隔帧取主麦
return _BesSource(route, _alternate(_link.audioStream));
case AudioRoute.callLocal:
return _BesSource(route, _link.micPcmStream);
case AudioRoute.callPeer:
return _BesSource(route, _link.spkPcmStream);
case AudioRoute.ai:
return _BesSource(
route,
_link.manager
.onDeviceAiEvent()
.where((e) => e['type'] == 'micData' && e['data'] is Uint8List)
.map((e) => e['data'] as Uint8List)
.where((d) => d.isNotEmpty));
default:
throw DeviceException(DeviceErrorCode.notSupported, 'openMic($route)');
}
}
static Stream<Uint8List> _alternate(Stream<Uint8List> raw) {
var take = true;
return raw.where((_) {
final t = take;
take = !take;
return t;
});
}
@override
Future<DeviceAudioSink> openSpeaker(
{AudioRoute route = AudioRoute.ai}) async {
_requireReady();
switch (route) {
case AudioRoute.ai:
await _link.manager.startAiDownlink();
return _BesAiSink(_link);
case AudioRoute.callLegA:
case AudioRoute.callLegB:
if (_callSinks == 0) await _link.startTranslationDownlink();
_callSinks++;
return _BesCallSink(_link, route, () async {
_callSinks--;
if (_callSinks <= 0) {
_callSinks = 0;
await _link.stopTranslationDownlink();
}
});
default:
throw DeviceException(
DeviceErrorCode.notSupported, 'openSpeaker($route)');
}
}
// ---------- feature ----------
/// feature → 字节的对应表。**冻结**,改动看金测试。
static const Map<String, int> featureCommands = {
DeviceFeatures.recordStart: BesCmd.micRecordStart,
DeviceFeatures.recordStop: BesCmd.micRecordStop,
DeviceFeatures.callStart: BesCmd.callModeStart,
DeviceFeatures.callStop: BesCmd.callModeStop,
DeviceFeatures.callQuery: BesCmd.callState,
DeviceFeatures.asrStarted: BesCmd.asrStarted,
DeviceFeatures.asrStopped: BesCmd.asrStopped,
DeviceFeatures.aiAccept: BesCmd.aiAccept,
DeviceFeatures.aiExit: BesCmd.aiExit,
DeviceFeatures.f2fStart: BesCmd.faceToFaceStart,
DeviceFeatures.f2fStop: BesCmd.faceToFaceStop,
DeviceFeatures.batteryQuery: BesCmd.status,
'bes.firmware.query': BesCmd.firmwareVersion,
};
@override
Future<Map<String, Object?>> invokeFeature(String featureKey,
[Map<String, Object?> args = const {}]) async {
_requireReady();
final cmd = featureCommands[featureKey];
if (cmd != null) {
await _link.sendData(BesCmd.frame(cmd));
return {'sent': BesCmd.frame(cmd)};
}
switch (featureKey) {
case DeviceFeatures.stateRefresh:
await _link.refreshDeviceState();
return const {};
case 'bes.probe':
await _link.probeCommands(
from: (args['from'] as int?) ?? 0x21,
to: (args['to'] as int?) ?? 0x7F,
);
return const {};
case 'bes.send':
final bytes = (args['bytes'] as List?)?.cast<int>() ?? const <int>[];
if (bytes.isEmpty) {
throw DeviceException(DeviceErrorCode.invalidArgument, 'bytes');
}
await _link.sendData(bytes);
return {'sent': bytes};
default:
throw DeviceException(DeviceErrorCode.notSupported, featureKey);
}
}
@override
DeviceOtaPort? otaPort() => _ota ??= BesOtaPort(this, _link);
@override
Future<void> disconnect() => _link.disconnectDevice();
}
class _BesSource implements DeviceAudioSource {
_BesSource(this.route, Stream<Uint8List> src) {
_sub = src.listen(_ctrl.add, onError: _ctrl.addError);
}
@override
final AudioRoute route;
final StreamController<Uint8List> _ctrl = StreamController.broadcast();
StreamSubscription? _sub;
@override
AudioFormat get format => AudioFormat.pcm16kMono;
@override
Stream<Uint8List> get pcm => _ctrl.stream;
@override
Future<void> close() async {
await _sub?.cancel();
_sub = null;
await _ctrl.close();
}
}
/// AI 下行:G.722 编码与 60ms 组包、0xD6/0xE9 调速全在原生 `AiAudioDownlink`。
class _BesAiSink implements DeviceAudioSink {
_BesAiSink(this._link);
final BesBluetoothService _link;
bool _closed = false;
@override
AudioRoute get route => AudioRoute.ai;
@override
AudioFormat get format => AudioFormat.pcm16kMono;
@override
Future<void> write(Uint8List pcm) async {
if (_closed) return;
await _link.manager.pushAiPcm(pcm);
}
@override
Future<void> discardPending() async {
if (_closed) return;
await _link.manager.flushAiDownlink();
}
@override
Future<void> close() async {
if (_closed) return;
_closed = true;
await _link.manager.stopAiDownlink();
}
}
/// 通话翻译下行:`AA 56` 84 字节包,G.722 + 20ms 节流在原生 `CallTranslationDownlink`。
class _BesCallSink implements DeviceAudioSink {
_BesCallSink(this._link, this.route, this._onClose);
final BesBluetoothService _link;
@override
final AudioRoute route;
final Future<void> Function() _onClose;
bool _closed = false;
String get _leg => route == AudioRoute.callLegA ? 'A' : 'B';
@override
AudioFormat get format => AudioFormat.pcm16kMono;
@override
Future<void> write(Uint8List pcm) async {
if (_closed) return;
await _link.pushTranslationTtsPcm(_leg, pcm);
}
@override
Future<void> discardPending() async {}
@override
Future<void> close() async {
if (_closed) return;
_closed = true;
await _onClose();
}
}