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

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': 1000, // 设置静音结束阈值,单位为毫秒,用于判断一句话的结束
'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();
}
}