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 _ctrl = StreamController.broadcast(); BesDeviceSession? _session; Worker? _statusWorker; bool _initialized = false; bool _disposed = false; @override String get vendorKey => DeviceVendors.bes; @override String get displayName => 'BES (恒玄)'; @override Set get capabilities => BesDeviceSession.caps; @override DevicePluginConfigSchema get configSchema => DevicePluginConfigSchema.empty; @override DeviceSession? get activeSession => _session; @override Stream get eventStream => _ctrl.stream; @override Future initialize(DevicePluginConfig config) async { _checkAlive(); if (_initialized) return; _initialized = true; if (!Get.isRegistered()) { Get.put(link, permanent: true); } _statusWorker = ever(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 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 stopScan() async {} @override Future isScanning() async => false; @override Future> 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 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 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 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 rememberDevice(DiscoveredDevice device) => link.savePairedDevice({'address': device.id, 'name': device.name}); @override Future tryAutoConnect() => link.tryAutoConnect(); @override Future forgetDevice(String deviceId) => link.removePairedDevice(deviceId); @override Future 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 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 _events = StreamController.broadcast(); final StreamController _wake = StreamController.broadcast(); final List _workers = []; StreamSubscription? _aiSub; DeviceConnectionState _state = DeviceConnectionState.ready; /// 通话翻译下行的引用计数:A/B 两路都关了才 stopCallDownlink int _callSinks = 0; BesOtaPort? _ota; void _bind() { _workers.add(ever(_link.lastCmdType, (type) { switch (type) { case 'startMicRecording': _feature(DeviceFeatureKeys.recordStarted); break; case 'stopMicRecording': _feature(DeviceFeatureKeys.recordStopped); break; } })); _workers.add(ever(_link.isCallOngoing, (on) { _feature(on ? DeviceFeatureKeys.callStarted : DeviceFeatureKeys.callEnded); // 同时刷一次快照:DeviceHub 把 info.metadata['inCall'] 当作 inCall 的兜底真相 // (事件会随 worker 一起被 dispose 掉,快照不会)。不发这一条,兜底就只能 // 等下次电量/版本变化时顺带对齐,慢一拍。 _infoUpdated(); })); void infoChanged(dynamic _) => _infoUpdated(); _workers.add(ever(_link.batteryLevel, (_) { _feature(DeviceFeatureKeys.batteryUpdated); _infoUpdated(); })); _workers.add(ever(_link.batteryCase, infoChanged)); _workers.add(ever(_link.chargingLeft, infoChanged)); _workers.add(ever(_link.chargingRight, infoChanged)); _workers.add(ever(_link.chargingCase, infoChanged)); _workers.add(ever(_link.deviceVersion, (_) { _feature(DeviceFeatureKeys.firmwareUpdated); _infoUpdated(); })); _workers.add(ever(_link.deviceMac, (_) { _feature(DeviceFeatureKeys.macUpdated); _infoUpdated(); })); _workers.add(ever(_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')); } @override Future> diagnostics() async { try { final dl = await _link.manager.getCallDownlinkStats(); return {'downlink': dl}; } catch (e) { return {'downlink': 'err:$e'}; } } void _feature(String key, [Map 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 get capabilities => caps; @override Stream get eventStream => _events.stream; @override Stream 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 readRssi() => throw DeviceException(DeviceErrorCode.notSupported, 'rssi'); @override Future readBattery() async { _requireReady(); await _link.queryBattery(); return info.battery.summary; } @override Future refreshInfo() async { _requireReady(); await _link.refreshDeviceState(); return info; } /// 取真实 MAC,iOS 上会主动发一次查询并等最多 [timeout]。 Future resolveMac({Duration timeout = const Duration(seconds: 3)}) => _link.resolveDeviceMac(timeout: timeout); // ---------- 音频 ---------- @override Future 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)'); } } /// 原生落盘录音。route → 原生的音频类型串: /// /// | route | 类型 | 说明 | /// |---|---|---| /// | mic | `mainMicData` | 现场录音的主麦 | /// | callLocal | `micData` | 通话本端麦 | /// | callPeer | `spkData` | 通话对端声 | /// /// ⚠️ `mic` 这里取的是**带类型的主麦帧**,而不是 [openMic] 那条 /// `_alternate(audioStream)`(裸流隔帧取)。隔帧取靠的是"订阅时恰好对上相位", /// 中途开始订阅就会永久取到前馈麦——原生早就把两路标好了 mainMicData / /// ffMicData,按类型取顺带把这类错相消掉。 /// 金测试 `test/devices/bes_recorder_sources_test.dart` 逐条钉着这张表: /// 串要与原生 `BluetoothManager` 发布的类型名逐字一致,改错不报任何错, /// 只表现为「录音起来了但一个字节都没写进去」。 static const Map recorderSources = { AudioRoute.mic: 'mainMicData', AudioRoute.callLocal: 'micData', AudioRoute.callPeer: 'spkData', }; @override Future openFileRecorder({ required String path, required List routes, }) async { _requireReady(); if (routes.isEmpty || routes.length > 2) return null; final sources = []; for (final r in routes) { final s = recorderSources[r]; if (s == null) return null; // 这一路没有对应的原生类型,让调用方退回 Dart 落盘 sources.add(s); } final id = 'bes_${DateTime.now().microsecondsSinceEpoch}'; final ok = await _link.startNativeFileRecorder( id: id, path: path, channels: sources.length, source: sources.length == 1 ? sources.first : null, left: sources.length == 2 ? sources[0] : null, right: sources.length == 2 ? sources[1] : null, ); if (!ok) return null; return _BesFileRecorder(id, path, _link); } static Stream _alternate(Stream raw) { var take = true; return raw.where((_) { final t = take; take = !take; return t; }); } @override Future 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 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> invokeFeature(String featureKey, [Map 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 DeviceFeatures.callPcmNativeBridge: // 不是发给耳机的命令,不进 featureCommands(金测试只管字节表)。 // 原生没这个方法(老包)会返回 false → 业务层退回 Dart 流。 final enabled = args['enabled'] == true; final ok = await _link.setNativeCallPcmBridge(enabled); if (!ok && enabled) { throw DeviceException(DeviceErrorCode.notSupported, featureKey); } return {'enabled': ok && enabled}; case DeviceFeatures.aiPcmNativeBridge: // 同 callPcmNativeBridge:不是发给耳机的命令,不进 featureCommands。 // 原生没这个方法(老包)会返回 false → 业务层退回 Dart 链路。 final aiOn = args['enabled'] == true; final aiOk = await _link.setNativeAiPcmBridge(aiOn); if (!aiOk && aiOn) { throw DeviceException(DeviceErrorCode.notSupported, featureKey); } return {'enabled': aiOk && aiOn}; 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() ?? const []; 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 disconnect() => _link.disconnectDevice(); } class _BesSource implements DeviceAudioSource { _BesSource(this.route, Stream src) { _sub = src.listen(_ctrl.add, onError: _ctrl.addError); } @override final AudioRoute route; final StreamController _ctrl = StreamController.broadcast(); StreamSubscription? _sub; @override AudioFormat get format => AudioFormat.pcm16kMono; @override Stream get pcm => _ctrl.stream; @override Future 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 write(Uint8List pcm) async { if (_closed) return; await _link.manager.pushAiPcm(pcm); } @override Future discardPending() async { if (_closed) return; await _link.manager.flushAiDownlink(); } @override Future 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 Function() _onClose; bool _closed = false; String get _leg => route == AudioRoute.callLegA ? 'A' : 'B'; @override AudioFormat get format => AudioFormat.pcm16kMono; @override Future write(Uint8List pcm) async { if (_closed) return; await _link.pushTranslationTtsPcm(_leg, pcm); } @override Future discardPending() async {} @override Future close() async { if (_closed) return; _closed = true; await _onClose(); } } /// [DeviceFileRecorder] 的恒玄实现:只是原生录制器的一个句柄。 /// 音频一个字节都不经过这里——它只转发统计、并在结束时让原生回填 WAV 头。 class _BesFileRecorder implements DeviceFileRecorder { _BesFileRecorder(this._id, this.path, this._link) { _sub = _link.nativeRecorderEvents .where((e) => e['id'] == _id) .listen((e) { final b = (e['bytes'] as num?)?.toInt() ?? _bytes; _bytes = b; if (!_ctrl.isClosed) { _ctrl.add(DeviceRecordStats( bytes: b, level: (e['level'] as num?)?.toDouble() ?? 0, )); } }); } final String _id; @override final String path; final BesBluetoothService _link; StreamSubscription? _sub; final StreamController _ctrl = StreamController.broadcast(); int _bytes = 0; bool _stopped = false; @override Stream get stats => _ctrl.stream; @override int get bytesWritten => _bytes; @override Future stop({bool save = true}) async { if (_stopped) return; _stopped = true; await _sub?.cancel(); _sub = null; final r = await _link.stopNativeFileRecorder(_id, save: save); // null = 原生已经自己收尾过(设备断连那条路径),这时最后一次统计里的 // 字节数就是最终值,别把它清零 final b = (r?['bytes'] as num?)?.toInt(); if (b != null) _bytes = b; await _ctrl.close(); } }