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.
 
 
 
 
 
 

766 lines
24 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);
// 同时刷一次快照:DeviceHub 把 info.metadata['inCall'] 当作 inCall 的兜底真相
// (事件会随 worker 一起被 dispose 掉,快照不会)。不发这一条,兜底就只能
// 等下次电量/版本变化时顺带对齐,慢一拍。
_infoUpdated();
}));
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'));
}
@override
Future<Map<String, Object?>> diagnostics() async {
try {
final dl = await _link.manager.getCallDownlinkStats();
return {'downlink': dl};
} catch (e) {
return {'downlink': 'err:$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)');
}
}
/// 原生落盘录音。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<AudioRoute, String> recorderSources = {
AudioRoute.mic: 'mainMicData',
AudioRoute.callLocal: 'micData',
AudioRoute.callPeer: 'spkData',
};
@override
Future<DeviceFileRecorder?> openFileRecorder({
required String path,
required List<AudioRoute> routes,
}) async {
_requireReady();
if (routes.isEmpty || routes.length > 2) return null;
final sources = <String>[];
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<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 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<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();
}
}
/// [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<DeviceRecordStats> _ctrl =
StreamController<DeviceRecordStats>.broadcast();
int _bytes = 0;
bool _stopped = false;
@override
Stream<DeviceRecordStats> get stats => _ctrl.stream;
@override
int get bytesWritten => _bytes;
@override
Future<void> 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();
}
}