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.
1003 lines
31 KiB
1003 lines
31 KiB
import 'dart:async';
|
|
import 'dart:convert';
|
|
import 'dart:io';
|
|
import 'dart:math';
|
|
import 'dart:typed_data';
|
|
import 'package:crypto/crypto.dart';
|
|
import 'package:flutter/foundation.dart' show kIsWeb;
|
|
import 'package:flutter_dotenv/flutter_dotenv.dart';
|
|
import 'package:get/get.dart';
|
|
import 'package:uuid/uuid.dart';
|
|
import 'package:web_socket_channel/io.dart';
|
|
import 'package:web_socket_channel/web_socket_channel.dart';
|
|
import '../../core/utils/logger.dart';
|
|
import 'package:record/record.dart';
|
|
|
|
/// 识别事件类型
|
|
enum RecognitionEventType {
|
|
started,
|
|
recognizing,
|
|
finalResult,
|
|
error,
|
|
}
|
|
|
|
/// 识别事件
|
|
class RecognitionEvent {
|
|
final RecognitionEventType type;
|
|
final String text;
|
|
final String? error;
|
|
|
|
RecognitionEvent({
|
|
required this.type,
|
|
this.text = '',
|
|
this.error,
|
|
});
|
|
}
|
|
|
|
/// 火山语音识别API服务
|
|
///
|
|
/// 该服务提供了通过WebSocket与火山语音大模型流式识别API交互的接口
|
|
class VolcanoAsrApiService extends GetxService {
|
|
// 协议常量
|
|
static const int _protocolVersion = 0x01;
|
|
static const int _defaultHeaderSize = 0x01;
|
|
|
|
// 消息类型
|
|
static const int _fullClientRequest = 0x01;
|
|
static const int _audioOnlyRequest = 0x02;
|
|
static const int _fullServerResponse = 0x09;
|
|
static const int _serverAck = 0x0B;
|
|
static const int _serverErrorResponse = 0x0F;
|
|
|
|
// 消息类型特定标志
|
|
static const int _noSequence = 0x00;
|
|
static const int _posSequence = 0x01;
|
|
static const int _negSequence = 0x02;
|
|
static const int _negWithSequence = 0x03;
|
|
|
|
// 消息序列化方法
|
|
static const int _noSerialization = 0x00;
|
|
static const int _json = 0x01;
|
|
|
|
// 消息压缩方法
|
|
static const int _noCompression = 0x00;
|
|
static const int _gzip = 0x01;
|
|
|
|
// 配置参数
|
|
late final String _appId;
|
|
late final String _token;
|
|
late final String _resourceId;
|
|
late final String _apiUrl;
|
|
|
|
// WebSocket连接
|
|
WebSocketChannel? _webSocketChannel;
|
|
|
|
// 状态变量
|
|
final _isListening = false.obs;
|
|
final _isConnected = false.obs;
|
|
final _errorMessage = ''.obs;
|
|
final _recognitionResults = <String>[].obs;
|
|
|
|
// 序列号
|
|
int _sequence = 0;
|
|
|
|
// 识别结果流控制器
|
|
final _recognitionStreamController = StreamController<RecognitionEvent>.broadcast();
|
|
|
|
// 最新的识别结果
|
|
final _latestRecognizedText = ''.obs;
|
|
String get latestRecognizedText => _latestRecognizedText.value;
|
|
|
|
// 获取可观察状态
|
|
RxBool get isListening => _isListening;
|
|
bool get isConnected => _isConnected.value;
|
|
String get errorMessage => _errorMessage.value;
|
|
List<String> get recognitionResults => _recognitionResults;
|
|
|
|
// 识别结果流
|
|
Stream<RecognitionEvent> get recognitionStream => _recognitionStreamController.stream;
|
|
|
|
// 音频录制相关
|
|
final _audioRecorder = AudioRecorder();
|
|
final _isRecording = false.obs;
|
|
bool get isRecording => _isRecording.value;
|
|
|
|
// 音频数据缓冲区
|
|
final List<int> _audioBuffer = [];
|
|
StreamSubscription<Uint8List>? _audioStreamSubscription;
|
|
|
|
// 音频发送控制
|
|
DateTime _lastAudioSendTime = DateTime.now();
|
|
static const Duration _minAudioSendInterval = Duration(milliseconds: 100);
|
|
static const int _minAudioBufferSize = 3200; // 约200ms的16kHz 16bit音频
|
|
|
|
// 音频录制配置
|
|
static const int _defaultSampleRate = 16000;
|
|
static const int _defaultBitsPerSample = 16;
|
|
static const int _defaultChannels = 1;
|
|
|
|
VolcanoAsrApiService() {
|
|
// 从环境变量获取配置
|
|
_appId = dotenv.env['VOLCANO_ASR_APP_ID'] ?? '';
|
|
_token = dotenv.env['VOLCANO_ASR_APP_TOKEN'] ?? '';
|
|
_resourceId = dotenv.env['VOLCANO_ASR_RESOURCE_ID'] ?? 'volc.bigasr.sauc.duration';
|
|
_apiUrl = dotenv.env['VOLCANO_ASR_API_URL'] ?? 'wss://openspeech.bytedance.com/api/v3/sauc/bigmodel';
|
|
|
|
Logger.info('火山语音识别API配置: APP_ID=${_appId.isNotEmpty ? "已设置" : "未设置"}, APP_TOKEN=${_token.isNotEmpty ? "已设置" : "未设置"}');
|
|
Logger.info('资源ID: $_resourceId');
|
|
Logger.info('API URL: $_apiUrl');
|
|
|
|
if (_appId.isEmpty || _token.isEmpty) {
|
|
_errorMessage.value = '火山语音识别API配置不完整,请检查环境变量';
|
|
Logger.error(_errorMessage.value);
|
|
}
|
|
}
|
|
|
|
@override
|
|
void onInit() {
|
|
super.onInit();
|
|
}
|
|
|
|
/// 连接到WebSocket服务器
|
|
Future<bool> connect({
|
|
required int sampleRate,
|
|
required int bitsPerSample,
|
|
required int channels,
|
|
String format = 'pcm',
|
|
String codec = 'raw',
|
|
String modelName = 'bigmodel',
|
|
bool enablePunc = true,
|
|
}) async {
|
|
if (_isConnected.value) {
|
|
Logger.warning('已经连接到WebSocket服务器');
|
|
return true;
|
|
}
|
|
|
|
if (_appId.isEmpty || _token.isEmpty) {
|
|
_errorMessage.value = '火山语音识别API配置不完整,请检查环境变量';
|
|
Logger.error(_errorMessage.value);
|
|
return false;
|
|
}
|
|
Logger.info('开始连接到ASR WebSocket服务器');
|
|
try {
|
|
_errorMessage.value = '';
|
|
_recognitionResults.clear();
|
|
_sequence = 0;
|
|
|
|
// 构建请求头
|
|
final Map<String, dynamic> headers = {
|
|
'X-Api-App-Key': _appId,
|
|
'X-Api-Access-Key': _token,
|
|
'X-Api-Resource-Id': _resourceId,
|
|
'X-Api-Connect-Id': const Uuid().v4(),
|
|
};
|
|
|
|
Logger.info('连接到WebSocket服务器: $_apiUrl');
|
|
|
|
// 创建WebSocket连接
|
|
if (kIsWeb) {
|
|
// Web平台不支持在连接时添加headers,需要使用其他方式
|
|
// 可以考虑将认证信息添加到URL中
|
|
final Uri uri = Uri.parse(_apiUrl).replace(
|
|
queryParameters: {
|
|
'app_key': _appId,
|
|
'access_key': _token,
|
|
'resource_id': _resourceId,
|
|
'connect_id': const Uuid().v4(),
|
|
},
|
|
);
|
|
Logger.info('Web平台连接URL: ${uri.toString()}');
|
|
_webSocketChannel = WebSocketChannel.connect(uri);
|
|
} else {
|
|
// 非Web平台使用IOWebSocketChannel
|
|
_webSocketChannel = IOWebSocketChannel.connect(
|
|
Uri.parse(_apiUrl),
|
|
headers: headers,
|
|
pingInterval: const Duration(seconds: 10), // 添加ping间隔,保持连接活跃
|
|
);
|
|
}
|
|
|
|
// 监听WebSocket消息
|
|
_webSocketChannel!.stream.listen(
|
|
_handleWebSocketMessage,
|
|
onError: _handleWebSocketError,
|
|
onDone: _handleWebSocketDone,
|
|
);
|
|
|
|
// 发送全客户端请求
|
|
await _sendFullClientRequest(
|
|
sampleRate: sampleRate,
|
|
bitsPerSample: bitsPerSample,
|
|
channels: channels,
|
|
format: format,
|
|
codec: codec,
|
|
modelName: modelName,
|
|
enablePunc: enablePunc,
|
|
);
|
|
|
|
_isConnected.value = true;
|
|
_isListening.value = true;
|
|
|
|
_recognitionStreamController.add(RecognitionEvent(
|
|
type: RecognitionEventType.started,
|
|
text: '',
|
|
));
|
|
|
|
Logger.info('ASR WebSocket连接成功建立');
|
|
|
|
return true;
|
|
} catch (e) {
|
|
_errorMessage.value = '连接WebSocket服务器失败: $e';
|
|
Logger.error('连接WebSocket服务器失败: $e');
|
|
_cleanupWebSocket();
|
|
return false;
|
|
}
|
|
}
|
|
|
|
/// 发送全客户端请求
|
|
Future<void> _sendFullClientRequest({
|
|
required int sampleRate,
|
|
required int bitsPerSample,
|
|
required int channels,
|
|
required String format,
|
|
required String codec,
|
|
required String modelName,
|
|
required bool enablePunc,
|
|
}) async {
|
|
// 构建payload
|
|
final Map<String, dynamic> payload = {
|
|
'user': {
|
|
'uid': 'user_${DateTime.now().millisecondsSinceEpoch}',
|
|
},
|
|
'audio': {
|
|
'format': format,
|
|
'sample_rate': sampleRate,
|
|
'bits': bitsPerSample,
|
|
'channel': channels,
|
|
'codec': codec,
|
|
},
|
|
'request': {
|
|
'model_name': modelName,
|
|
'enable_punc': enablePunc,
|
|
'show_utterances': true, // 启用分句信息,用于区分中间结果和最终结果
|
|
'result_type': 'all', // 返回所有分句结果
|
|
'enable_words': true, // 启用词级别时间戳
|
|
'enable_itn': false, // 启用数字转换
|
|
'continuous_decoding': true, // 启用连续解码,支持连续监听
|
|
'vad_silence_end': 800, // 设置静音结束阈值,单位为毫秒,用于判断一句话的结束
|
|
'vad_max_duration': 60000, // 设置最大语音段时长,单位为毫秒
|
|
},
|
|
};
|
|
|
|
// 序列化并压缩payload
|
|
final String payloadStr = jsonEncode(payload);
|
|
// Logger.info('发送全客户端请求: $payloadStr');
|
|
final List<int> payloadBytes = await _gzipCompress(utf8.encode(payloadStr));
|
|
|
|
// 增加序列号
|
|
_sequence = 1;
|
|
|
|
// 构建消息
|
|
final List<int> message = _buildMessage(
|
|
messageType: _fullClientRequest,
|
|
messageTypeSpecificFlags: _posSequence,
|
|
serializationMethod: _json,
|
|
compressionType: _gzip,
|
|
sequence: _sequence,
|
|
payload: payloadBytes,
|
|
);
|
|
|
|
// 发送消息
|
|
_webSocketChannel!.sink.add(Uint8List.fromList(message));
|
|
}
|
|
|
|
/// 发送音频数据
|
|
Future<bool> sendAudioData(List<int> audioData, {bool isLast = false}) async {
|
|
if (!_isConnected.value || _webSocketChannel == null) {
|
|
return false;
|
|
}
|
|
|
|
try {
|
|
// 如果是最后一个包,直接发送
|
|
if (isLast) {
|
|
// 先发送缓冲区中的数据
|
|
if (_audioBuffer.isNotEmpty) {
|
|
await _sendAudioDataInternal(_audioBuffer, false);
|
|
_audioBuffer.clear();
|
|
}
|
|
|
|
// 发送最后一个包
|
|
return await _sendAudioDataInternal([0, 0, 0, 0, 0, 0, 0, 0], true);
|
|
}
|
|
|
|
// 将新的音频数据添加到缓冲区
|
|
_audioBuffer.addAll(audioData);
|
|
|
|
// 检查是否应该发送缓冲区中的数据
|
|
final now = DateTime.now();
|
|
final timeSinceLastSend = now.difference(_lastAudioSendTime);
|
|
|
|
if (_audioBuffer.length >= _minAudioBufferSize || timeSinceLastSend >= _minAudioSendInterval) {
|
|
// 创建缓冲区的副本并清空缓冲区
|
|
final bufferToSend = List<int>.from(_audioBuffer);
|
|
_audioBuffer.clear();
|
|
|
|
// 更新最后发送时间
|
|
_lastAudioSendTime = now;
|
|
|
|
// 发送缓冲区中的数据
|
|
return await _sendAudioDataInternal(bufferToSend, false);
|
|
}
|
|
|
|
return true;
|
|
} catch (e) {
|
|
Logger.error('发送音频数据失败: $e');
|
|
return false;
|
|
}
|
|
}
|
|
|
|
/// 内部方法:实际发送音频数据
|
|
Future<bool> _sendAudioDataInternal(List<int> audioData, bool isLast) async {
|
|
if (!_isConnected.value || _webSocketChannel == null) {
|
|
return false;
|
|
}
|
|
|
|
try {
|
|
// 增加序列号
|
|
_sequence++;
|
|
|
|
// 如果是最后一个音频包,使用负序列号
|
|
final int sequence = isLast ? -_sequence : _sequence;
|
|
final int messageTypeSpecificFlags = isLast ? _negWithSequence : _posSequence;
|
|
|
|
// 压缩音频数据
|
|
List<int> compressedAudio;
|
|
if (audioData.isEmpty && isLast) {
|
|
// 对于最后一个空包,使用一个有效的GZIP数据
|
|
compressedAudio = [31, 139, 8, 0, 0, 0, 0, 0, 0, 3, 3, 0, 0, 0, 0, 0, 0, 0, 0, 0];
|
|
} else {
|
|
compressedAudio = await _gzipCompress(audioData);
|
|
}
|
|
|
|
// 构建消息
|
|
final List<int> message = _buildMessage(
|
|
messageType: _audioOnlyRequest,
|
|
messageTypeSpecificFlags: messageTypeSpecificFlags,
|
|
serializationMethod: _noSerialization,
|
|
compressionType: _gzip,
|
|
sequence: sequence,
|
|
payload: compressedAudio,
|
|
);
|
|
|
|
// 发送消息
|
|
_webSocketChannel!.sink.add(Uint8List.fromList(message));
|
|
|
|
|
|
return true;
|
|
} catch (e) {
|
|
Logger.error('发送音频数据失败: $e');
|
|
return false;
|
|
}
|
|
}
|
|
|
|
/// 构建消息
|
|
List<int> _buildMessage({
|
|
required int messageType,
|
|
required int messageTypeSpecificFlags,
|
|
required int serializationMethod,
|
|
required int compressionType,
|
|
required int sequence,
|
|
required List<int> payload,
|
|
}) {
|
|
// 构建头部
|
|
final List<int> header = [
|
|
(_protocolVersion << 4) | _defaultHeaderSize,
|
|
(messageType << 4) | messageTypeSpecificFlags,
|
|
(serializationMethod << 4) | compressionType,
|
|
0, // 保留字段
|
|
];
|
|
|
|
// 序列号
|
|
final List<int> sequenceBytes = _intToBytes(sequence);
|
|
|
|
// payload大小
|
|
final List<int> payloadSizeBytes = _intToBytes(payload.length);
|
|
|
|
// 组装消息
|
|
final List<int> message = [
|
|
...header,
|
|
...sequenceBytes,
|
|
...payloadSizeBytes,
|
|
...payload,
|
|
];
|
|
|
|
return message;
|
|
}
|
|
|
|
/// 处理WebSocket消息
|
|
void _handleWebSocketMessage(dynamic message) {
|
|
if (message is! List<int>) {
|
|
// 尝试解析非二进制消息作为JSON错误
|
|
try {
|
|
if (message is String) {
|
|
final Map<String, dynamic> errorJson = jsonDecode(message);
|
|
if (errorJson.containsKey('error')) {
|
|
final String errorMessage = errorJson['error'];
|
|
Logger.warning('收到WebSocket错误消息: $errorMessage');
|
|
|
|
_errorMessage.value = '服务器错误: $errorMessage';
|
|
_recognitionStreamController.add(RecognitionEvent(
|
|
type: RecognitionEventType.error,
|
|
error: _errorMessage.value,
|
|
));
|
|
|
|
// 如果是严重错误,关闭连接
|
|
if (errorMessage.contains('decode ws request failed') ||
|
|
errorMessage.contains('unable to ungzip payload')) {
|
|
_isListening.value = false;
|
|
_isConnected.value = false;
|
|
_cleanupWebSocket();
|
|
}
|
|
|
|
return;
|
|
}
|
|
}
|
|
} catch (e) {
|
|
// 解析失败,使用原始消息
|
|
Logger.warning('收到非二进制WebSocket消息: $message');
|
|
}
|
|
return;
|
|
}
|
|
|
|
try {
|
|
final List<int> data = message;
|
|
|
|
// 确保消息长度足够
|
|
if (data.length < 12) {
|
|
Logger.warning('WebSocket消息长度不足: ${data.length}字节');
|
|
return;
|
|
}
|
|
|
|
// 解析头部
|
|
final int protocolVersion = (data[0] >> 4) & 0x0F;
|
|
final int headerSize = data[0] & 0x0F;
|
|
final int messageType = (data[1] >> 4) & 0x0F;
|
|
final int messageTypeSpecificFlags = data[1] & 0x0F;
|
|
final int serializationMethod = (data[2] >> 4) & 0x0F;
|
|
final int compressionType = data[2] & 0x0F;
|
|
|
|
// 解析序列号
|
|
final List<int> sequenceBytes = data.sublist(4, 8);
|
|
final int sequence = _bytesToInt(sequenceBytes);
|
|
|
|
// 解析payload大小
|
|
final List<int> payloadSizeBytes = data.sublist(8, 12);
|
|
final int payloadSize = _bytesToInt(payloadSizeBytes);
|
|
|
|
// 确保payload长度正确
|
|
if (data.length < 12 + payloadSize) {
|
|
Logger.warning('WebSocket消息payload长度不足: 预期${payloadSize}字节,实际${data.length - 12}字节');
|
|
return;
|
|
}
|
|
|
|
// 解析payload
|
|
final List<int> payload = data.sublist(12, 12 + payloadSize);
|
|
|
|
// 检查是否是最后一个包(负序列号)
|
|
final bool isLastPackage = sequence < 0;
|
|
|
|
// 处理不同类型的消息
|
|
if (messageType == _fullServerResponse) {
|
|
_handleFullServerResponse(payload, compressionType, isLastPackage);
|
|
} else if (messageType == _serverAck) {
|
|
_handleServerAck(payload);
|
|
} else if (messageType == _serverErrorResponse) {
|
|
_handleServerErrorResponse(sequence, payload);
|
|
}
|
|
|
|
// 如果是最后一个包,记录日志但不断开连接
|
|
if (isLastPackage) {
|
|
Logger.info('收到最后一个包,序列号: $sequence,继续监听');
|
|
}
|
|
} catch (e) {
|
|
Logger.error('处理WebSocket消息失败: $e');
|
|
|
|
_errorMessage.value = '处理WebSocket消息失败: $e';
|
|
_recognitionStreamController.add(RecognitionEvent(
|
|
type: RecognitionEventType.error,
|
|
error: _errorMessage.value,
|
|
));
|
|
}
|
|
}
|
|
|
|
/// 处理全服务器响应
|
|
Future<void> _handleFullServerResponse(List<int> payload, int compressionType, bool isLastPackage) async {
|
|
try {
|
|
String payloadStr;
|
|
|
|
if (compressionType == _gzip) {
|
|
final List<int> decompressed = await _gzipDecompress(payload);
|
|
payloadStr = utf8.decode(decompressed);
|
|
} else {
|
|
payloadStr = utf8.decode(payload);
|
|
}
|
|
|
|
// 解析JSON
|
|
final Map<String, dynamic> response = jsonDecode(payloadStr);
|
|
|
|
|
|
|
|
// 处理识别结果
|
|
if (response.containsKey('result')) {
|
|
// 处理分句信息
|
|
if (response['result'] is Map &&
|
|
response['result'].containsKey('utterances') &&
|
|
response['result']['utterances'] is List &&
|
|
response['result']['utterances'].isNotEmpty) {
|
|
|
|
// 遍历所有分句
|
|
for (final utterance in response['result']['utterances']) {
|
|
final String text = utterance['text'] ?? '';
|
|
final bool isDefinite = utterance['definite'] ?? false;
|
|
|
|
if (text.isNotEmpty) {
|
|
if (isDefinite) {
|
|
// 当definite为true时,表示这是一个完整的句子,作为最终结果处理
|
|
if (!_recognitionResults.contains(text)) {
|
|
_recognitionResults.add(text);
|
|
_latestRecognizedText.value = text;
|
|
|
|
_recognitionStreamController.add(RecognitionEvent(
|
|
type: RecognitionEventType.finalResult,
|
|
text: text,
|
|
));
|
|
|
|
// Logger.info('添加最终识别结果: $text');
|
|
|
|
// 重要:在这里不要断开连接,而是继续监听下一句话
|
|
// 但可以清空当前的中间结果状态,准备接收新的语音输入
|
|
_latestRecognizedText.value = '';
|
|
|
|
// 关键修改:向服务器发送一个短暂的静音包,触发服务器端的分段
|
|
// 这有助于确保下一次识别从新的内容开始
|
|
if (_isConnected.value && _webSocketChannel != null && _isRecording.value) {
|
|
// 发送一个极短的静音包,不会中断录音,但会触发服务器重置状态
|
|
sendAudioData([0, 0, 0, 0, 0, 0, 0, 0], isLast: false);
|
|
}
|
|
}
|
|
} else {
|
|
// 当definite为false时,表示这是一个中间结果
|
|
// 只有当文本与最新的不同时才发送,避免重复
|
|
if (text != _latestRecognizedText.value) {
|
|
_latestRecognizedText.value = text;
|
|
|
|
_recognitionStreamController.add(RecognitionEvent(
|
|
type: RecognitionEventType.recognizing,
|
|
text: text,
|
|
));
|
|
|
|
// Logger.info('添加中间识别结果: $text');
|
|
}
|
|
}
|
|
}
|
|
}
|
|
} else if (response['result'] is Map && response['result'].containsKey('text')) {
|
|
// 如果没有分句信息但有完整文本,则根据isLastPackage判断是否作为最终结果
|
|
final String fullText = response['result']['text'] ?? '';
|
|
|
|
if (fullText.isNotEmpty && fullText != _latestRecognizedText.value) {
|
|
_latestRecognizedText.value = fullText;
|
|
|
|
// 如果是最后一个包,则作为最终结果处理
|
|
if (isLastPackage) {
|
|
if (!_recognitionResults.contains(fullText)) {
|
|
_recognitionResults.add(fullText);
|
|
|
|
_recognitionStreamController.add(RecognitionEvent(
|
|
type: RecognitionEventType.finalResult,
|
|
text: fullText,
|
|
));
|
|
|
|
Logger.info('添加最终识别结果(最后一个包): $fullText');
|
|
|
|
// 重要:清空当前的中间结果状态,准备接收新的语音输入
|
|
_latestRecognizedText.value = '';
|
|
|
|
// 关键修改:向服务器发送一个短暂的静音包,触发服务器端的分段
|
|
if (_isConnected.value && _webSocketChannel != null && _isRecording.value) {
|
|
// 发送一个极短的静音包,不会中断录音,但会触发服务器重置状态
|
|
sendAudioData([0, 0, 0, 0, 0, 0, 0, 0], isLast: false);
|
|
}
|
|
}
|
|
} else {
|
|
// 否则作为中间结果处理
|
|
_recognitionStreamController.add(RecognitionEvent(
|
|
type: RecognitionEventType.recognizing,
|
|
text: fullText,
|
|
));
|
|
|
|
Logger.info('添加中间识别结果(无分句): $fullText');
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// 如果是最后一个包,记录日志
|
|
if (isLastPackage) {
|
|
Logger.info('收到最后一个包,继续监听');
|
|
}
|
|
} catch (e) {
|
|
Logger.error('处理全服务器响应失败: $e');
|
|
}
|
|
}
|
|
|
|
/// 处理服务器确认
|
|
void _handleServerAck(List<int> payload) {
|
|
try {
|
|
final String payloadStr = utf8.decode(payload);
|
|
Logger.info('收到服务器确认: $payloadStr');
|
|
} catch (e) {
|
|
Logger.error('处理服务器确认失败: $e');
|
|
}
|
|
}
|
|
|
|
/// 处理服务器错误响应
|
|
void _handleServerErrorResponse(int errorCode, List<int> payload) {
|
|
try {
|
|
final String errorMessage = utf8.decode(payload);
|
|
Logger.error('服务器错误: 错误码=$errorCode, 错误信息=$errorMessage');
|
|
|
|
_errorMessage.value = '服务器错误: $errorMessage (错误码: $errorCode)';
|
|
|
|
_recognitionStreamController.add(RecognitionEvent(
|
|
type: RecognitionEventType.error,
|
|
error: _errorMessage.value,
|
|
));
|
|
|
|
// 对于任何服务器错误,主动关闭连接
|
|
_isListening.value = false;
|
|
_isConnected.value = false;
|
|
_cleanupWebSocket();
|
|
|
|
// 如果正在录音,也停止录音
|
|
if (_isRecording.value) {
|
|
stopRecognition();
|
|
}
|
|
} catch (e) {
|
|
Logger.error('处理服务器错误响应失败: $e');
|
|
}
|
|
}
|
|
|
|
/// 处理最后一个包
|
|
void _handleLastPackage() {
|
|
|
|
// 确保最后的文本被作为最终结果发送
|
|
if (_latestRecognizedText.value.isNotEmpty &&
|
|
!_recognitionResults.contains(_latestRecognizedText.value)) {
|
|
_recognitionResults.add(_latestRecognizedText.value);
|
|
|
|
_recognitionStreamController.add(RecognitionEvent(
|
|
type: RecognitionEventType.finalResult,
|
|
text: _latestRecognizedText.value,
|
|
));
|
|
|
|
Logger.info('添加最终识别结果(最后一个包): ${_latestRecognizedText.value}');
|
|
}
|
|
|
|
_isListening.value = false;
|
|
_cleanupWebSocket();
|
|
}
|
|
|
|
/// 处理WebSocket错误
|
|
void _handleWebSocketError(dynamic error) {
|
|
// 避免重复记录相同的错误
|
|
if (_errorMessage.value.contains(error.toString())) {
|
|
return;
|
|
}
|
|
|
|
// 简化错误消息处理
|
|
final String errorStr = error.toString().toLowerCase();
|
|
String errorMessage;
|
|
|
|
if (errorStr.contains('not upgraded to websocket') || errorStr.contains('status code: 400')) {
|
|
errorMessage = '认证失败,请检查API密钥';
|
|
} else if (errorStr.contains('connection refused') || errorStr.contains('failed host lookup')) {
|
|
errorMessage = '无法连接到服务器,请检查网络';
|
|
} else if (errorStr.contains('timeout')) {
|
|
errorMessage = '连接超时';
|
|
} else {
|
|
errorMessage = '连接错误: $error';
|
|
}
|
|
|
|
_errorMessage.value = errorMessage;
|
|
Logger.error('WebSocket错误: $errorMessage');
|
|
|
|
_recognitionStreamController.add(RecognitionEvent(
|
|
type: RecognitionEventType.error,
|
|
error: errorMessage,
|
|
));
|
|
|
|
// 清理资源
|
|
_isListening.value = false;
|
|
_isConnected.value = false;
|
|
_cleanupWebSocket();
|
|
}
|
|
|
|
/// 处理WebSocket关闭
|
|
void _handleWebSocketDone() {
|
|
Logger.info('ASR WebSocket连接已关闭');
|
|
|
|
_isListening.value = false;
|
|
_isConnected.value = false;
|
|
_cleanupWebSocket();
|
|
}
|
|
|
|
/// 清理WebSocket连接
|
|
void _cleanupWebSocket() {
|
|
if (_webSocketChannel == null) {
|
|
return;
|
|
}
|
|
|
|
try {
|
|
// 先关闭sink,然后等待stream自然关闭
|
|
_webSocketChannel?.sink.close(WebSocketStatus.normalClosure, '客户端主动关闭连接');
|
|
} catch (e) {
|
|
Logger.error('关闭WebSocket连接失败: $e');
|
|
} finally {
|
|
_webSocketChannel = null;
|
|
}
|
|
}
|
|
|
|
/// 检查麦克风权限
|
|
Future<bool> checkMicrophonePermission() async {
|
|
try {
|
|
final hasPermission = await _audioRecorder.hasPermission();
|
|
return hasPermission;
|
|
} catch (e) {
|
|
Logger.error('检查麦克风权限失败: $e');
|
|
return false;
|
|
}
|
|
}
|
|
|
|
/// 开始录音并识别
|
|
Future<bool> startRecognition({
|
|
int sampleRate = _defaultSampleRate,
|
|
int bitsPerSample = _defaultBitsPerSample,
|
|
int channels = _defaultChannels,
|
|
String format = 'pcm',
|
|
String codec = 'raw',
|
|
String modelName = 'bigmodel',
|
|
bool enablePunc = true,
|
|
}) async {
|
|
try {
|
|
// 如果已经在录音,先停止
|
|
if (_isRecording.value) {
|
|
await stopRecognition();
|
|
// 添加短暂延迟,确保之前的会话完全关闭
|
|
await Future.delayed(const Duration(milliseconds: 500));
|
|
}
|
|
|
|
// 检查麦克风权限
|
|
final hasPermission = await checkMicrophonePermission();
|
|
if (!hasPermission) {
|
|
_errorMessage.value = '没有麦克风权限';
|
|
Logger.error(_errorMessage.value);
|
|
return false;
|
|
}
|
|
|
|
// 清空识别结果
|
|
_recognitionResults.clear();
|
|
_latestRecognizedText.value = '';
|
|
|
|
// 清空音频缓冲区
|
|
_audioBuffer.clear();
|
|
_lastAudioSendTime = DateTime.now();
|
|
|
|
// 连接到ASR服务
|
|
final connected = await connect(
|
|
sampleRate: sampleRate,
|
|
bitsPerSample: bitsPerSample,
|
|
channels: channels,
|
|
format: format,
|
|
codec: codec,
|
|
modelName: modelName,
|
|
enablePunc: enablePunc,
|
|
);
|
|
|
|
if (!connected) {
|
|
return false;
|
|
}
|
|
|
|
// 配置录音
|
|
final config = RecordConfig(
|
|
encoder: AudioEncoder.pcm16bits,
|
|
sampleRate: sampleRate,
|
|
numChannels: channels,
|
|
bitRate: bitsPerSample * sampleRate * channels,
|
|
);
|
|
|
|
Logger.info('开始录音,采样率: $sampleRate Hz, 位深: $bitsPerSample bits, 通道数: $channels');
|
|
|
|
// 开始录音流
|
|
final stream = await _audioRecorder.startStream(config);
|
|
|
|
_isRecording.value = true;
|
|
|
|
// 订阅音频流
|
|
_audioStreamSubscription = stream.listen(
|
|
(data) {
|
|
if (_isConnected.value && _isRecording.value && _webSocketChannel != null) {
|
|
sendAudioData(data.toList());
|
|
}
|
|
},
|
|
onError: (error) {
|
|
Logger.error('音频流错误: $error');
|
|
_errorMessage.value = '音频流错误: $error';
|
|
_recognitionStreamController.add(RecognitionEvent(
|
|
type: RecognitionEventType.error,
|
|
error: _errorMessage.value,
|
|
));
|
|
},
|
|
onDone: () {
|
|
Logger.info('音频流结束');
|
|
if (_isRecording.value) {
|
|
stopRecognition();
|
|
}
|
|
},
|
|
cancelOnError: false,
|
|
);
|
|
|
|
return true;
|
|
} catch (e) {
|
|
_errorMessage.value = '开始录音失败: $e';
|
|
Logger.error(_errorMessage.value);
|
|
|
|
_recognitionStreamController.add(RecognitionEvent(
|
|
type: RecognitionEventType.error,
|
|
error: _errorMessage.value,
|
|
));
|
|
|
|
return false;
|
|
}
|
|
}
|
|
|
|
/// 停止录音并完成识别
|
|
Future<bool> stopRecognition() async {
|
|
if (!_isRecording.value) {
|
|
return true;
|
|
}
|
|
|
|
try {
|
|
_isRecording.value = false;
|
|
|
|
// 取消音频流订阅
|
|
await _audioStreamSubscription?.cancel();
|
|
_audioStreamSubscription = null;
|
|
|
|
// 停止录音
|
|
await _audioRecorder.stop();
|
|
|
|
// 发送最后一个音频包,标记为结束
|
|
if (_isConnected.value && _webSocketChannel != null) {
|
|
try {
|
|
// 发送最后一个包并等待处理
|
|
await sendAudioData([0, 0, 0, 0, 0, 0, 0, 0], isLast: true);
|
|
await Future.delayed(const Duration(milliseconds: 300));
|
|
} catch (e) {
|
|
Logger.error('发送最后一个音频包失败: $e');
|
|
}
|
|
|
|
// 清空当前状态
|
|
_latestRecognizedText.value = '';
|
|
}
|
|
|
|
// 断开连接
|
|
await disconnect();
|
|
|
|
return true;
|
|
} catch (e) {
|
|
Logger.error('停止录音失败: $e');
|
|
await disconnect();
|
|
return false;
|
|
}
|
|
}
|
|
|
|
/// 断开连接
|
|
Future<void> disconnect() async {
|
|
Logger.info('断开ASR连接');
|
|
// 避免递归调用
|
|
if (_isConnected.value && _isRecording.value) {
|
|
_isRecording.value = false;
|
|
|
|
// 取消音频流订阅
|
|
await _audioStreamSubscription?.cancel();
|
|
_audioStreamSubscription = null;
|
|
|
|
// 停止录音
|
|
try {
|
|
await _audioRecorder.stop();
|
|
} catch (e) {
|
|
Logger.error('停止录音失败: $e');
|
|
}
|
|
}
|
|
|
|
_cleanupWebSocket();
|
|
_isConnected.value = false;
|
|
_isListening.value = false;
|
|
|
|
// 确保在断开连接时不会有未处理的事件
|
|
if (!_recognitionStreamController.isClosed && _recognitionStreamController.hasListener) {
|
|
try {
|
|
// 发送一个最终事件,表示连接已断开
|
|
_recognitionStreamController.add(RecognitionEvent(
|
|
type: RecognitionEventType.error,
|
|
error: '连接已断开',
|
|
));
|
|
} catch (e) {
|
|
Logger.error('发送断开连接事件失败: $e');
|
|
}
|
|
}
|
|
}
|
|
|
|
/// 整数转字节数组
|
|
List<int> _intToBytes(int value) {
|
|
return [
|
|
(value >> 24) & 0xFF,
|
|
(value >> 16) & 0xFF,
|
|
(value >> 8) & 0xFF,
|
|
value & 0xFF,
|
|
];
|
|
}
|
|
|
|
/// 字节数组转整数
|
|
int _bytesToInt(List<int> bytes) {
|
|
if (bytes.length != 4) {
|
|
throw ArgumentError('字节数组长度必须为4');
|
|
}
|
|
|
|
return ((bytes[0] & 0xFF) << 24) |
|
|
((bytes[1] & 0xFF) << 16) |
|
|
((bytes[2] & 0xFF) << 8) |
|
|
(bytes[3] & 0xFF);
|
|
}
|
|
|
|
/// GZIP压缩
|
|
Future<List<int>> _gzipCompress(List<int> data) async {
|
|
if (data.isEmpty) {
|
|
// 返回一个有效的空GZIP数据流
|
|
return [31, 139, 8, 0, 0, 0, 0, 0, 0, 3, 3, 0, 0, 0, 0, 0, 0, 0, 0, 0];
|
|
}
|
|
|
|
try {
|
|
return gzip.encode(data);
|
|
} catch (e) {
|
|
Logger.error('GZIP压缩失败: $e');
|
|
// 返回原始数据,但这可能导致协议错误
|
|
throw Exception('GZIP压缩失败: $e');
|
|
}
|
|
}
|
|
|
|
/// GZIP解压缩
|
|
Future<List<int>> _gzipDecompress(List<int> data) async {
|
|
if (data.isEmpty) {
|
|
return [];
|
|
}
|
|
|
|
try {
|
|
return gzip.decode(data);
|
|
} catch (e) {
|
|
Logger.error('GZIP解压缩失败: $e,数据长度: ${data.length}');
|
|
// 无法解压缩,返回原始数据
|
|
throw Exception('GZIP解压缩失败: $e');
|
|
}
|
|
}
|
|
|
|
@override
|
|
void onClose() {
|
|
stopRecognition();
|
|
disconnect();
|
|
_audioRecorder.dispose();
|
|
_recognitionStreamController.close();
|
|
super.onClose();
|
|
}
|
|
}
|