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.
 
 
 
 
 
 

223 lines
8.7 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();
}
void _onChunk(Uint8List chunk) {
if (chunk.isEmpty) return;
_lastChunkAt = DateTime.now();
if (!_levelCtl.isClosed) _levelCtl.add(PcmStreamPlayer.rmsLevel(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();
if (!_levelCtl.isClosed) _levelCtl.add(0);
}
Future<void> dispose() async {
await stop();
await _packets.close();
await _levelCtl.close();
await _recorder.dispose();
}
}