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.
246 lines
9.6 KiB
246 lines
9.6 KiB
import 'dart:io' show Platform;
|
|
import 'package:flutter/services.dart';
|
|
import 'package:flutter/foundation.dart' show kIsWeb;
|
|
import 'dart:async';
|
|
import 'dart:typed_data';
|
|
|
|
import 'package:record/record.dart';
|
|
|
|
import '../../core/utils/logger.dart';
|
|
import 'emai_audio_session.dart';
|
|
import 'pcm_stream_player.dart';
|
|
|
|
/// 手机麦克风的 PCM 流采集,专供百炼上行使用。
|
|
///
|
|
/// `record` 的 `startStream` 回调块大小由平台决定(Android 通常几百字节到 2KB),
|
|
/// 百炼建议每包 ~100ms(16k/16bit/单声道 = 3200 字节),所以这里做一层重新分包:
|
|
/// 攒够 [_packetBytes] 再发一包,尾巴不足一包的在 [stop] 时补发。
|
|
class MicPcmStreamer {
|
|
MicPcmStreamer({this.sampleRate = 16000});
|
|
|
|
static const String _tag = 'MicPcmStreamer';
|
|
|
|
final int sampleRate;
|
|
|
|
/// 100ms 一包
|
|
int get _packetBytes => sampleRate * 2 ~/ 10;
|
|
|
|
final AudioRecorder _recorder = AudioRecorder();
|
|
StreamSubscription<Uint8List>? _sub;
|
|
final _buffer = BytesBuilder(copy: false);
|
|
|
|
final _packets = StreamController<Uint8List>.broadcast();
|
|
|
|
/// 重新分包后的上行 PCM
|
|
Stream<Uint8List> get packets => _packets.stream;
|
|
|
|
final _levelCtl = StreamController<double>.broadcast();
|
|
|
|
/// 麦克风实时音量(0~1),UI 音浪用
|
|
Stream<double> get level => _levelCtl.stream;
|
|
|
|
bool _running = false;
|
|
bool get isRunning => _running;
|
|
|
|
Future<bool> hasPermission() => _recorder.hasPermission();
|
|
|
|
/// 正在进行中的 [start]。并发调用共用同一个 future。
|
|
///
|
|
/// ⚠️ 这个去重不能省。`_running` 是在 `await startStream()` **返回之后**才置位的,
|
|
/// 两个调用方几乎同时进来时都会看到 `_running == false`、都往下走,于是同一个
|
|
/// `AudioRecorder` 被 startStream 两遍——第二次把第一次那条流顶掉,上行当场断供。
|
|
/// 真机表现是「麦克风采集启动」打两行、上行只出几包就没了,服务端 30 秒后回
|
|
/// `ClientAudioTimeout`,客户端不报任何错。
|
|
Future<bool>? _startFuture;
|
|
|
|
Future<bool> start() {
|
|
if (_running) return Future.value(true);
|
|
return _startFuture ??= _start().whenComplete(() => _startFuture = null);
|
|
}
|
|
|
|
Future<bool> _start() async {
|
|
try {
|
|
// ⚠️ 必须先把会话配好并激活再开麦:会话停在 soloAmbient 且未激活时,
|
|
// startStream 会成功返回、引擎也起得来,但输入 tap 一个 buffer 都不回调。
|
|
// 返回值是「输出此刻走不走 A2DP」:走的话 iOS 不能开 voice-processing,
|
|
// 否则系统把 A2DP 从输出里剔掉、回复又回到扬声器(见 EmaiAudioSession.ensure)。
|
|
final a2dpOutput = await EmaiAudioSession.ensure();
|
|
final bool voiceProcessing = !(Platform.isIOS && a2dpOutput);
|
|
if (!await _recorder.hasPermission()) {
|
|
Logger.w(_tag, '没有麦克风权限');
|
|
return false;
|
|
}
|
|
// ⚠️ iOS:record 拿 inputNode.inputFormat 去 installTap,格式是 0Hz(会话刚
|
|
// 切类别 / 路由正在切换)时那是 ObjC 断言,这里的 try/catch 接不住,直接崩。
|
|
// 先问原生一声输入格式是否就绪,不就绪等一小会儿,还不行就报失败而不是崩。
|
|
if (!await _waitIosInputReady()) {
|
|
Logger.e(_tag, '输入音频格式无效(会话/路由切换中),本次不开麦');
|
|
return false;
|
|
}
|
|
final stream = await _recorder.startStream(
|
|
RecordConfig(
|
|
encoder: AudioEncoder.pcm16bits,
|
|
sampleRate: sampleRate,
|
|
numChannels: 1,
|
|
// 语音助手场景:开回声消除和降噪,否则外放播报会被自己的麦克风收回去,
|
|
// duplex 模式下会直接把服务端 VAD 打乱。
|
|
// ⚠️ iOS 上回复走 A2DP 耳机时**必须关**:这两项会开 voice-processing,
|
|
// 而语音处理模式下 A2DP 不可作输出。耳机里放的声音手机麦本来就收不到多少。
|
|
echoCancel: voiceProcessing,
|
|
noiseSuppress: voiceProcessing,
|
|
// ⚠️ iOS 上不让 record 碰 AVAudioSession(会话由 EmaiAudioSession 统一配)。
|
|
// 它自己那套是 `setCategory(.playAndRecord, options:)` → `setActive(true)`,
|
|
// 而 setCategory 会把 mode 重置成 default;紧接着 echoCancel 触发的
|
|
// `setVoiceProcessingEnabled(true)` 又把 mode 拽回 voiceChat。这一来一回
|
|
// 是一次 AVAudioEngine 配置变更,而 record 没有监听
|
|
// `AVAudioEngineConfigurationChange` —— 引擎「还在运行」,输入 tap 却
|
|
// 从此不再回调。真机表现:上行送了几十包(两三秒)后彻底停住,
|
|
// 没有任何报错,服务端随后回 ClientAudioTimeout。
|
|
// 我们先把会话配成 playAndRecord + voiceChat 并激活,voice-processing
|
|
// 打开时 mode 已经就位,就不会再触发那次变更。
|
|
// (Android 侧忽略 iosConfig,这里不影响。)
|
|
// ignore: deprecated_member_use
|
|
iosConfig: const IosRecordConfig(manageAudioSession: false),
|
|
),
|
|
);
|
|
_running = true;
|
|
_buffer.clear();
|
|
_startWatchdog();
|
|
_sub = stream.listen(
|
|
_onChunk,
|
|
onError: (e) => Logger.e(_tag, '录音流错误: $e'),
|
|
);
|
|
Logger.i(_tag, '麦克风采集启动 ${sampleRate}Hz');
|
|
return true;
|
|
} catch (e) {
|
|
Logger.e(_tag, '麦克风采集启动失败: $e');
|
|
_running = false;
|
|
return false;
|
|
}
|
|
}
|
|
|
|
/// 见 [_start] 里的说明。非 iOS 直接放行;原生不认这个方法(老插件)也放行。
|
|
static const MethodChannel _iosProbe = MethodChannel('azure_speech/pcm_player');
|
|
|
|
Future<bool> _waitIosInputReady() async {
|
|
if (kIsWeb || !Platform.isIOS) return true;
|
|
for (var i = 0; i < 12; i++) {
|
|
try {
|
|
final r = await _iosProbe.invokeMethod<Map>('inputFormat');
|
|
final rate = (r?['sampleRate'] as num?)?.toDouble() ?? 0;
|
|
final ch = (r?['channels'] as num?)?.toInt() ?? 0;
|
|
if (rate > 0 && ch > 0) return true;
|
|
Logger.w(_tag, '输入格式未就绪 rate=$rate ch=$ch,等待重试 #${i + 1}');
|
|
} on MissingPluginException {
|
|
return true;
|
|
} catch (e) {
|
|
Logger.w(_tag, '探测输入格式失败(放行): $e');
|
|
return true;
|
|
}
|
|
await Future.delayed(const Duration(milliseconds: 50));
|
|
}
|
|
return false;
|
|
}
|
|
|
|
/// 最近一次收到原始音频块的时刻。用来发现「流无声死亡」——
|
|
/// 麦克风即使没人说话也会持续出包,所以静默超过 [_silenceLimit] 一定是流断了。
|
|
DateTime? _lastChunkAt;
|
|
Timer? _watchdog;
|
|
|
|
static const Duration _silenceLimit = Duration(seconds: 3);
|
|
|
|
void _startWatchdog() {
|
|
_watchdog?.cancel();
|
|
_lastChunkAt = DateTime.now();
|
|
_watchdog = Timer.periodic(const Duration(seconds: 1), (_) {
|
|
if (!_running) return;
|
|
final last = _lastChunkAt;
|
|
if (last == null) return;
|
|
if (DateTime.now().difference(last) < _silenceLimit) return;
|
|
Logger.e(_tag, '录音流已静默 ${_silenceLimit.inSeconds}s,判定断流,尝试重启');
|
|
_lastChunkAt = DateTime.now();
|
|
unawaited(_restart());
|
|
});
|
|
}
|
|
|
|
Future<void> _restart() async {
|
|
await stop();
|
|
await start();
|
|
}
|
|
|
|
/// 音浪发布间隔与窗口峰值。原始块由平台决定大小(Android 几百字节~2KB
|
|
/// = 15~60ms 一块),块块都发就是每秒几十次 Rx 通知,而屏幕最多 60Hz、
|
|
/// 音浪控件也没那么灵敏。取窗口峰值而不是最后一块:说话的尖峰常常只占
|
|
/// 一两块,按最后一块取样会把音浪削平。
|
|
static const Duration _levelEmitInterval = Duration(milliseconds: 100);
|
|
double _levelPeak = 0;
|
|
DateTime? _levelEmitAt;
|
|
|
|
void _emitLevel(Uint8List chunk) {
|
|
if (_levelCtl.isClosed) return;
|
|
final v = PcmStreamPlayer.rmsLevel(chunk);
|
|
if (v > _levelPeak) _levelPeak = v;
|
|
final now = DateTime.now();
|
|
final last = _levelEmitAt;
|
|
if (last != null && now.difference(last) < _levelEmitInterval) return;
|
|
_levelEmitAt = now;
|
|
_levelCtl.add(_levelPeak);
|
|
_levelPeak = 0;
|
|
}
|
|
|
|
void _onChunk(Uint8List chunk) {
|
|
if (chunk.isEmpty) return;
|
|
_lastChunkAt = DateTime.now();
|
|
_emitLevel(chunk);
|
|
|
|
_buffer.add(chunk);
|
|
while (_buffer.length >= _packetBytes) {
|
|
// BytesBuilder 只能整体取出,取出后把余量放回去
|
|
final all = _buffer.takeBytes();
|
|
var offset = 0;
|
|
while (all.length - offset >= _packetBytes) {
|
|
_emit(Uint8List.sublistView(all, offset, offset + _packetBytes));
|
|
offset += _packetBytes;
|
|
}
|
|
if (offset < all.length) {
|
|
_buffer.add(Uint8List.sublistView(all, offset));
|
|
}
|
|
}
|
|
}
|
|
|
|
void _emit(Uint8List packet) {
|
|
if (_packets.isClosed) return;
|
|
// sublistView 是原缓冲的视图,下游可能异步持有,复制一份避免被后续写入覆盖
|
|
_packets.add(Uint8List.fromList(packet));
|
|
}
|
|
|
|
Future<void> stop() async {
|
|
if (!_running) return;
|
|
_running = false;
|
|
_watchdog?.cancel();
|
|
_watchdog = null;
|
|
await _sub?.cancel();
|
|
_sub = null;
|
|
try {
|
|
await _recorder.stop();
|
|
} catch (e) {
|
|
Logger.w(_tag, '停止录音异常: $e');
|
|
}
|
|
// 尾包补发,否则最后不足 100ms 的语音会被丢掉
|
|
if (_buffer.length > 0) {
|
|
_emit(_buffer.takeBytes());
|
|
}
|
|
_buffer.clear();
|
|
// 下一轮从零开始计峰值,并让首块立刻出一次电平
|
|
_levelPeak = 0;
|
|
_levelEmitAt = null;
|
|
if (!_levelCtl.isClosed) _levelCtl.add(0);
|
|
}
|
|
|
|
Future<void> dispose() async {
|
|
await stop();
|
|
await _packets.close();
|
|
await _levelCtl.close();
|
|
await _recorder.dispose();
|
|
}
|
|
}
|
|
|