wolfplus 1 year ago
parent
commit
feab4c0d35
  1. 51
      local_plugins/realtime/.gitignore
  2. 190
      local_plugins/realtime/IMPLEMENTATION.md
  3. 125
      local_plugins/realtime/README.md
  4. 47
      local_plugins/realtime/android/build.gradle.kts
  5. 1
      local_plugins/realtime/android/settings.gradle.kts
  6. 11
      local_plugins/realtime/android/src/main/AndroidManifest.xml
  7. 519
      local_plugins/realtime/android/src/main/kotlin/com/yunqiinnovation/realtime/RealtimeAudioManager.kt
  8. 455
      local_plugins/realtime/android/src/main/kotlin/com/yunqiinnovation/realtime/RealtimePlugin.kt
  9. 303
      local_plugins/realtime/android/src/main/kotlin/com/yunqiinnovation/realtime/RealtimeWebSocketManager.kt
  10. 274
      local_plugins/realtime/example.md
  11. 18
      local_plugins/realtime/ios/realtime/Package.swift
  12. 344
      local_plugins/realtime/ios/realtime/Sources/realtime/RealtimeAudioManager.swift
  13. 319
      local_plugins/realtime/ios/realtime/Sources/realtime/RealtimePlugin.swift
  14. 323
      local_plugins/realtime/ios/realtime/Sources/realtime/RealtimeWebSocketManager.swift
  15. 243
      local_plugins/realtime/lib/realtime.dart
  16. 29
      local_plugins/realtime/pubspec.yaml

51
local_plugins/realtime/.gitignore

@ -0,0 +1,51 @@
# Miscellaneous
*.class
*.log
*.pyc
*.swp
.DS_Store
.atom/
.buildlog/
.history
.svn/
# IntelliJ related
*.iml
*.ipr
*.iws
.idea/
# Visual Studio Code related
.vscode/
# Flutter/Dart/Pub related
**/doc/api/
.dart_tool/
.flutter-plugins
.flutter-plugins-dependencies
.packages
.pub-cache/
.pub/
build/
# iOS related
**/ios/.generated/
**/ios/Flutter/Flutter.framework
**/ios/Flutter/Flutter.podspec
**/ios/Flutter/Generated.xcconfig
**/ios/Flutter/app.flx
**/ios/Flutter/app.zip
**/ios/Flutter/flutter_assets/
**/ios/Flutter/flutter_export_environment.sh
**/ios/ServiceDefinitions.json
**/ios/Runner/GeneratedPluginRegistrant.*
.build/
# Coverage
coverage/
# Exceptions to above rules.
!**/ios/**/default.mode1v3
!**/ios/**/default.mode2v3
!**/ios/**/default.pbxuser
!**/ios/**/default.perspectivev3

190
local_plugins/realtime/IMPLEMENTATION.md

@ -0,0 +1,190 @@
# Realtime Plugin 实现总结
## 项目概述
基于您提供的Android原生示例,我为DeepVoice项目创建了一个realtime插件,用于通过WebSocket接入Vocode服务器实现实时语音聊天功能。该插件已完成Android和iOS双平台实现。
## 架构设计
### 插件结构
```
local_plugins/realtime/
├── android/ # Android实现
│ ├── build.gradle # Android构建配置
│ ├── src/main/
│ │ ├── AndroidManifest.xml # Android权限配置
│ │ └── kotlin/com/yunqiinnovation/realtime/
│ │ ├── RealtimePlugin.kt # 主插件入口
│ │ ├── RealtimeAudioManager.kt # 音频管理器(AudioRecord/AudioTrack)
│ │ └── RealtimeWebSocketManager.kt # WebSocket管理器(OkHttp)
├── ios/ # iOS实现
│ ├── Classes/ # Objective-C桥接文件
│ │ ├── RealtimePlugin.h
│ │ └── RealtimePlugin.m
│ ├── realtime/ # Swift Package
│ │ ├── Package.swift
│ │ └── Sources/realtime/
│ │ ├── RealtimePlugin.swift # 主插件入口
│ │ ├── RealtimeAudioManager.swift # 音频管理器(AVAudioEngine)
│ │ └── RealtimeWebSocketManager.swift # WebSocket管理器(URLSession)
│ └── realtime.podspec
├── lib/
│ └── realtime.dart # Flutter接口
├── pubspec.yaml
└── README.md
```
### 核心组件
1. **RealtimePlugin (SwiftRealtimePlugin)**
- 主插件入口,处理Flutter方法调用
- 管理各组件间的协调
- 处理事件分发
2. **RealtimeAudioManager**
- 使用AVAudioEngine进行音频录制
- 使用AVAudioPlayerNode进行音频播放
- 支持16kHz/16-bit/单声道格式
- 实现20ms帧长的音频处理
3. **RealtimeWebSocketManager**
- 使用URLSessionWebSocketTask实现WebSocket通信
- 支持音频数据和文本消息的收发
- 自动连接状态管理
## 技术实现对照
| 功能 | Android示例 | iOS实现 |
|------|-------------|---------|
| 音频录制 | AudioRecord | AVAudioEngine + AVAudioInputNode |
| 音频播放 | AudioTrack | AVAudioPlayerNode |
| WebSocket | OkHttp WebSocket | URLSessionWebSocketTask |
| 线程管理 | Kotlin协程 | DispatchQueue |
| 音频格式 | 16kHz/16-bit/单声道 | 16kHz/16-bit/单声道 |
| 帧长 | 20ms (320字节) | 20ms (320字节) |
## Flutter接口
### 主要类
- `RealtimeService`: 主服务类
- `RealtimeEvent`: 事件类
- `ConnectionStatus`: 连接状态枚举
- `VoiceStatus`: 语音状态枚举
- `RealtimeException`: 异常类
### 主要方法
- `initialize()`: 初始化服务
- `connect()`: 连接服务器
- `disconnect()`: 断开连接
- `startRecording()`: 开始录音
- `stopRecording()`: 停止录音
- `eventStream`: 事件流
## 使用方式
### 1. 在RealtimeController中集成
已更新`lib/modules/realtime/controllers/realtime_controller.dart`使用真实的realtime插件,取代了原来的模拟实现。
### 2. 事件监听
```dart
_realtimeService.eventStream.listen((event) {
switch (event.type) {
case 'connectionStatusChanged':
// 处理连接状态变化
case 'voiceStatusChanged':
// 处理语音状态变化
case 'textReceived':
// 处理文本消息
case 'error':
// 处理错误
}
});
```
## 配置要求
### 依赖配置
已添加到主项目的`pubspec.yaml`:
```yaml
realtime:
path: local_plugins/realtime
```
### 权限配置
iOS的`Info.plist`已包含必要的麦克风权限:
```xml
<key>NSMicrophoneUsageDescription</key>
<string>需要麦克风权限用于语音识别和录音功能</string>
```
## 与Android示例的对应关系
1. **初始化对应**
- Android: 创建AudioRecord, AudioTrack, OkHttpClient
- iOS: 创建AVAudioEngine, AVAudioPlayerNode, URLSession
2. **录音线程对应**
- Android: `loopRecordSend()` 协程
- iOS: AVAudioInputNode的installTap回调
3. **播放线程对应**
- Android: `loopPlayback()` 协程 + LinkedBlockingQueue
- iOS: DispatchQueue + AVAudioPCMBuffer队列
4. **WebSocket对应**
- Android: OkHttp WebSocketListener
- iOS: URLSessionWebSocketDelegate
## 特性支持
✅ **已实现**
- Android和iOS双平台支持
- 实时音频录制和播放(16kHz/16-bit/单声道)
- WebSocket双向通信
- 状态管理和事件通知
- 错误处理
- 资源管理和清理
- 20ms帧长处理
- Kotlin协程和DispatchQueue线程管理
❌ **未实现(可扩展)**
- 音频编码(Opus等)
- 噪声消除/回声抑制
- 自动重连机制
- 音频质量自适应
## 使用注意事项
1. **服务器地址配置**
```dart
// 需要替换为实际的Vocode服务器地址
serverUrl: 'wss://your-vocode-server/ws'
```
2. **音频格式一致性**
- 确保服务器支持16kHz/16-bit/单声道格式
- 帧长固定为20ms (320字节)
3. **错误处理**
- 监听事件流中的错误
- 处理网络断线和重连
4. **资源管理**
- 及时调用dispose()释放资源
- 避免内存泄漏
## 测试建议
1. **本地测试**
- 先用echo服务器测试WebSocket连接
- 验证音频录制和播放功能
2. **集成测试**
- 与真实Vocode服务器集成
- 测试端到端语音交互
3. **性能测试**
- 测试长时间使用的稳定性
- 监控内存和CPU使用
这个实现为DeepVoice项目提供了完整的实时语音交互能力,可以直接与Vocode服务器进行通信,实现类似Android示例的功能。

125
local_plugins/realtime/README.md

@ -0,0 +1,125 @@
# Realtime Plugin
实时语音聊天插件,通过WebSocket连接Vocode服务器实现实时语音交互。
## 功能特性
- **实时音频录制**: 16kHz/16-bit/单声道格式录音
- **WebSocket通信**: 与Vocode服务器进行实时数据传输
- **音频播放**: 播放服务器返回的语音数据
- **状态管理**: 连接状态和语音状态监控
- **事件回调**: 支持各种事件的监听和处理
## 技术实现
### Android端实现
- **音频录制**: 使用AudioRecord进行16kHz/16-bit/单声道录音
- **音频播放**: 使用AudioTrack进行音频播放
- **WebSocket**: 使用OkHttp WebSocket客户端
- **线程管理**: 使用Kotlin协程处理录音和播放
### iOS端实现
- **音频录制**: 使用AVAudioEngine和AVAudioInputNode
- **音频播放**: 使用AVAudioPlayerNode进行播放
- **WebSocket**: 使用URLSessionWebSocketTask
- **线程管理**: 使用DispatchQueue管理录音和播放线程
### 核心组件
- `RealtimePlugin`: 主插件入口,处理Flutter方法调用
- `RealtimeAudioManager`: 音频管理器,负责录音和播放
- `RealtimeWebSocketManager`: WebSocket管理器,负责网络通信
**Android端组件**:
- AudioRecord + AudioTrack + OkHttp + Kotlin协程
- 对应您提供的Android原生示例功能
**iOS端组件**:
- AVAudioEngine + AVAudioPlayerNode + URLSessionWebSocketTask + DispatchQueue
## 使用方法
```dart
import 'package:realtime/realtime.dart';
final realtimeService = RealtimeService();
// 初始化
await realtimeService.initialize(
serverUrl: 'wss://your-server/ws',
sampleRate: 16000,
channels: 1,
bitsPerSample: 16,
);
// 监听事件
realtimeService.eventStream.listen((event) {
switch (event.type) {
case 'connectionStatusChanged':
// 处理连接状态变化
break;
case 'voiceStatusChanged':
// 处理语音状态变化
break;
case 'textReceived':
// 处理收到的文本消息
break;
case 'error':
// 处理错误
break;
}
});
// 连接服务器
await realtimeService.connect();
// 开始录音
await realtimeService.startRecording();
// 停止录音
await realtimeService.stopRecording();
// 断开连接
await realtimeService.disconnect();
// 释放资源
await realtimeService.dispose();
```
## 音频格式
- **采样率**: 16000 Hz
- **声道数**: 1(单声道)
- **位深**: 16-bit
- **帧长**: 20ms (320字节)
## 权限要求
### Android
权限已自动包含在插件中:
- `RECORD_AUDIO`: 音频录制权限
- `INTERNET`: 网络访问权限
- `ACCESS_NETWORK_STATE`: 网络状态权限
### iOS
在Info.plist中添加麦克风权限:
```xml
<key>NSMicrophoneUsageDescription</key>
<string>应用需要麦克风权限进行语音录制</string>
```
## 注意事项
1. 确保服务器地址正确且可访问
2. 网络环境良好,避免频繁断线
3. 音频格式与服务器保持一致
4. 及时释放资源,避免内存泄漏
## 错误处理
插件会自动处理常见错误:
- 网络连接失败
- 音频设备不可用
- 权限被拒绝
- 服务器断开连接
通过事件流可以监听这些错误并进行相应处理。

47
local_plugins/realtime/android/build.gradle.kts

@ -0,0 +1,47 @@
plugins {
id("com.android.library")
id("org.jetbrains.kotlin.android")
kotlin("plugin.serialization") version "1.9.24"
}
android {
namespace = "com.yunqiinnovation.realtime"
compileSdk = 35
defaultConfig {
minSdk = 21
targetSdk = 33
}
compileOptions {
sourceCompatibility = JavaVersion.VERSION_11
targetCompatibility = JavaVersion.VERSION_11
}
kotlinOptions {
jvmTarget = "11"
}
packaging {
resources {
excludes.add("META-INF/DEPENDENCIES")
excludes.add("META-INF/LICENSE")
excludes.add("META-INF/LICENSE.txt")
excludes.add("META-INF/license.txt")
excludes.add("META-INF/NOTICE")
excludes.add("META-INF/NOTICE.txt")
excludes.add("META-INF/notice.txt")
excludes.add("META-INF/*.kotlin_module")
}
}
}
dependencies {
implementation("org.jetbrains.kotlin:kotlin-stdlib-jdk8:1.9.10")
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-core:1.7.3")
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-android:1.7.3")
// OkHttp for WebSocket
implementation(platform("com.squareup.okhttp3:okhttp-bom:4.12.0"))
implementation("com.squareup.okhttp3:okhttp")
}

1
local_plugins/realtime/android/settings.gradle.kts

@ -0,0 +1 @@
rootProject.name = "realtime"

11
local_plugins/realtime/android/src/main/AndroidManifest.xml

@ -0,0 +1,11 @@
<manifest xmlns:android="http://schemas.android.com/apk/res/android"
package="com.yunqiinnovation.realtime">
<!-- 音频录制权限 -->
<uses-permission android:name="android.permission.RECORD_AUDIO" />
<!-- 网络权限 -->
<uses-permission android:name="android.permission.INTERNET" />
<!-- 网络状态权限 -->
<uses-permission android:name="android.permission.ACCESS_NETWORK_STATE" />
</manifest>

519
local_plugins/realtime/android/src/main/kotlin/com/yunqiinnovation/realtime/RealtimeAudioManager.kt

@ -0,0 +1,519 @@
package com.yunqiinnovation.realtime
import android.content.Context
import android.media.*
import android.media.audiofx.AcousticEchoCanceler
import android.media.audiofx.AutomaticGainControl
import android.media.audiofx.NoiseSuppressor
import android.util.Log
import kotlinx.coroutines.*
import java.io.ByteArrayOutputStream
import java.nio.ByteBuffer
import java.nio.ByteOrder
import java.util.concurrent.LinkedBlockingQueue
import kotlin.math.sqrt
/**
* 实时音频管理器
*
* 负责音频的录制和播放,对应Android示例中的AudioRecord和AudioTrack功能
*/
class RealtimeAudioManager(private val context: Context) {
companion object {
private const val TAG = "RealtimeAudioManager"
}
// 音频配置
private var audioConfig: AudioConfig? = null
// 录音相关
private var audioRecord: AudioRecord? = null
private var isRecording = false
private var recordJob: Job? = null
// 播放相关
private var audioTrack: AudioTrack? = null
private var isPlaying = false
private var playJob: Job? = null
private val playbackQueue = LinkedBlockingQueue<ByteArray>()
private val playbackInternalBuffer = ByteArrayOutputStream()
// 回音消除相关
private var acousticEchoCanceler: AcousticEchoCanceler? = null
private var noiseSuppressor: NoiseSuppressor? = null
private var automaticGainControl: AutomaticGainControl? = null
private var audioManager: AudioManager? = null
private var originalAudioMode = AudioManager.MODE_NORMAL
// 回调
var onAudioData: ((ByteArray) -> Unit)? = null
var onStatusChanged: ((VoiceStatus) -> Unit)? = null
var onRmsChanged: ((Float) -> Unit)? = null // 新增RMS回调
// 协程作用域
private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob())
/// 初始化音频管理器
fun initialize(config: AudioConfig): Boolean {
this.audioConfig = config
// 初始化系统音频管理器
audioManager = context.getSystemService(Context.AUDIO_SERVICE) as AudioManager
// 保存原始音频模式
originalAudioMode = audioManager?.mode ?: AudioManager.MODE_NORMAL
return setupAudioRecord(config) && setupAudioTrack(config)
}
/// 设置AudioRecord(录音)
private fun setupAudioRecord(config: AudioConfig): Boolean {
try {
val channelConfig = if (config.channels == 1) {
AudioFormat.CHANNEL_IN_MONO
} else {
AudioFormat.CHANNEL_IN_STEREO
}
val audioFormat = when (config.bitsPerSample) {
16 -> AudioFormat.ENCODING_PCM_16BIT
8 -> AudioFormat.ENCODING_PCM_8BIT
else -> AudioFormat.ENCODING_PCM_16BIT
}
val bufferSize = AudioRecord.getMinBufferSize(
config.sampleRate,
channelConfig,
audioFormat
)
if (bufferSize == AudioRecord.ERROR || bufferSize == AudioRecord.ERROR_BAD_VALUE) {
Log.e(TAG, "无法获取音频录制缓冲区大小")
return false
}
audioRecord = AudioRecord(
MediaRecorder.AudioSource.MIC,
config.sampleRate,
channelConfig,
audioFormat,
bufferSize
)
if (audioRecord?.state != AudioRecord.STATE_INITIALIZED) {
Log.e(TAG, "AudioRecord初始化失败")
return false
}
Log.i(TAG, "AudioRecord初始化成功,采样率: ${config.sampleRate}, 声道: ${config.channels}, 位深: ${config.bitsPerSample}")
return true
} catch (e: Exception) {
Log.e(TAG, "设置AudioRecord失败: ${e.message}", e)
return false
}
}
/// 设置AudioTrack(播放)
private fun setupAudioTrack(config: AudioConfig): Boolean {
try {
val channelConfig = if (config.channels == 1) {
AudioFormat.CHANNEL_OUT_MONO
} else {
AudioFormat.CHANNEL_OUT_STEREO
}
val audioFormat = when (config.bitsPerSample) {
16 -> AudioFormat.ENCODING_PCM_16BIT
8 -> AudioFormat.ENCODING_PCM_8BIT
else -> AudioFormat.ENCODING_PCM_16BIT
}
val bufferSize = AudioTrack.getMinBufferSize(
config.sampleRate,
channelConfig,
audioFormat
)
if (bufferSize == AudioTrack.ERROR || bufferSize == AudioTrack.ERROR_BAD_VALUE) {
Log.e(TAG, "无法获取音频播放缓冲区大小")
return false
}
// 总是使用VOICE_COMMUNICATION模式以获得最佳语音处理效果
val audioAttributes = AudioAttributes.Builder()
.setUsage(AudioAttributes.USAGE_VOICE_COMMUNICATION)
.setContentType(AudioAttributes.CONTENT_TYPE_SPEECH)
.build()
audioTrack = AudioTrack.Builder()
.setAudioAttributes(audioAttributes)
.setAudioFormat(
AudioFormat.Builder()
.setSampleRate(config.sampleRate)
.setChannelMask(channelConfig)
.setEncoding(audioFormat)
.build()
)
.setBufferSizeInBytes(bufferSize)
.setTransferMode(AudioTrack.MODE_STREAM)
.build()
if (audioTrack?.state != AudioTrack.STATE_INITIALIZED) {
Log.e(TAG, "AudioTrack初始化失败")
return false
}
Log.i(TAG, "AudioTrack初始化成功, 启用语音通信模式")
return true
} catch (e: Exception) {
Log.e(TAG, "设置AudioTrack失败: ${e.message}", e)
return false
}
}
/// 开始录音
fun startRecording(): Boolean {
if (isRecording || audioRecord == null) {
return false
}
try {
// 设置回音消除(总是启用)
setupEchoCancellation()
audioRecord?.startRecording()
isRecording = true
onStatusChanged?.invoke(VoiceStatus.RECORDING)
// 启动录音协程
recordJob = scope.launch {
loopRecordSend()
}
Log.i(TAG, "开始录音, 回音消除: ${acousticEchoCanceler?.enabled}")
return true
} catch (e: Exception) {
Log.e(TAG, "开始录音失败: ${e.message}", e)
return false
}
}
/// 录音循环(对应Android示例中的loopRecordSend)
private suspend fun loopRecordSend() {
val config = audioConfig ?: return
val frameSize = config.frameSize
val buffer = ByteArray(frameSize)
while (isRecording && audioRecord != null) {
try {
val bytesRead = audioRecord!!.read(buffer, 0, frameSize)
if (bytesRead > 0) {
// 确保读取的数据长度正确
val audioData = if (bytesRead == frameSize) {
buffer
} else {
buffer.copyOf(bytesRead)
}
// 计算RMS值
val rms = calculateRMS(audioData)
onRmsChanged?.invoke(rms)
// 直接回调音频数据到Flutter层
onAudioData?.invoke(audioData)
} else {
Log.w(TAG, "录音读取数据失败: $bytesRead")
}
} catch (e: Exception) {
Log.e(TAG, "录音循环异常: ${e.message}", e)
break
}
}
}
/// 停止录音
fun stopRecording(): Boolean {
if (!isRecording) {
return false
}
try {
isRecording = false
recordJob?.cancel()
audioRecord?.stop()
// 清理回音消除
cleanupEchoCancellation()
onStatusChanged?.invoke(VoiceStatus.IDLE)
Log.i(TAG, "录音已停止")
return true
} catch (e: Exception) {
Log.e(TAG, "停止录音失败: ${e.message}", e)
return false
}
}
/// 播放音频数据
fun playAudioData(data: ByteArray) {
if (data.isNotEmpty()) {
// 计算播放数据的RMS值
val rms = calculateRMS(data)
onRmsChanged?.invoke(rms)
playbackQueue.offer(data)
if (!isPlaying) {
startPlayback()
}
}
}
/// 开始播放
private fun startPlayback() {
if (isPlaying || audioTrack == null) {
return
}
try {
isPlaying = true
onStatusChanged?.invoke(VoiceStatus.PLAYING)
audioTrack?.play()
playJob = scope.launch {
loopPlayback()
}
Log.i(TAG, "开始播放")
} catch (e: Exception) {
Log.e(TAG, "开始播放失败: ${e.message}", e)
}
}
/// 播放循环
private suspend fun loopPlayback() {
Log.d(TAG, "启动播放循环")
while (isPlaying) {
try {
// take() 会一直等待直到队列中有数据
val data = playbackQueue.take()
// 直接播放整个数据块,让AudioTrack自己处理时序
audioTrack?.write(data, 0, data.size)
} catch (e: InterruptedException) {
Log.d(TAG, "播放循环被中断")
break
} catch (e: Exception) {
Log.e(TAG, "播放循环异常: ${e.message}", e)
}
}
Log.d(TAG, "播放循环结束")
stopPlaybackInternal()
}
/// 内部停止播放方法
private fun stopPlaybackInternal() {
if (!isPlaying) return
isPlaying = false
playbackQueue.clear()
// 使用try-catch保护,因为AudioTrack状态可能不稳定
try {
if (audioTrack?.playState == AudioTrack.PLAYSTATE_PLAYING) {
audioTrack?.pause()
audioTrack?.flush()
}
// 不要在这里调用stop,stop是用来释放资源的,pause/flush更适合流式播放的暂停
} catch (e: Exception) {
Log.e(TAG, "暂停和清空AudioTrack失败: ${e.message}", e)
}
onStatusChanged?.invoke(VoiceStatus.IDLE)
Log.i(TAG, "播放已通过内部调用停止")
}
/// 停止播放(外部调用)
fun stopPlaying(): Boolean {
if (!isPlaying) {
return true
}
try {
isPlaying = false
playJob?.cancel() // 取消协程
playJob = null
stopPlaybackInternal()
// 恢复音频模式
restoreAudioMode()
return true
} catch (e: Exception) {
Log.e(TAG, "停止播放失败: ${e.message}", e)
return false
}
}
/// 更新配置
fun updateConfig(config: AudioConfig): Boolean {
// 如果正在录音或播放,先停止
if (isRecording) {
stopRecording()
}
if (isPlaying) {
stopPlaying()
}
// 释放旧的音频对象
releaseAudioObjects()
return initialize(config)
}
/// 释放音频对象
private fun releaseAudioObjects() {
try {
audioRecord?.release()
audioRecord = null
audioTrack?.release()
audioTrack = null
} catch (e: Exception) {
Log.e(TAG, "释放音频对象失败: ${e.message}", e)
}
}
/// 释放资源
fun dispose() {
stopRecording()
stopPlaying()
// 清理回音消除和恢复音频模式
cleanupEchoCancellation()
restoreAudioMode()
scope.cancel()
releaseAudioObjects()
playbackQueue.clear()
Log.i(TAG, "音频管理器已释放")
}
/// 设置回音消除(总是启用)
private fun setupEchoCancellation() {
val record = audioRecord ?: return
try {
// 设置音频模式为通信模式
audioManager?.mode = AudioManager.MODE_IN_COMMUNICATION
// 设置声学回音消除
if (AcousticEchoCanceler.isAvailable()) {
acousticEchoCanceler = AcousticEchoCanceler.create(record.audioSessionId)
acousticEchoCanceler?.enabled = true
Log.i(TAG, "声学回音消除已启用")
} else {
Log.w(TAG, "设备不支持声学回音消除")
}
// 设置噪声抑制
if (NoiseSuppressor.isAvailable()) {
noiseSuppressor = NoiseSuppressor.create(record.audioSessionId)
noiseSuppressor?.enabled = true
Log.i(TAG, "噪声抑制已启用")
}
// 设置自动增益控制
if (AutomaticGainControl.isAvailable()) {
automaticGainControl = AutomaticGainControl.create(record.audioSessionId)
automaticGainControl?.enabled = true
Log.i(TAG, "自动增益控制已启用")
}
Log.i(TAG, "回音消除和音频增强功能已全部启用")
} catch (e: Exception) {
Log.e(TAG, "设置回音消除失败: ${e.message}", e)
}
}
/// 清理回音消除
private fun cleanupEchoCancellation() {
try {
acousticEchoCanceler?.let {
it.enabled = false
it.release()
acousticEchoCanceler = null
Log.d(TAG, "声学回音消除已释放")
}
noiseSuppressor?.let {
it.enabled = false
it.release()
noiseSuppressor = null
Log.d(TAG, "噪声抑制已释放")
}
automaticGainControl?.let {
it.enabled = false
it.release()
automaticGainControl = null
Log.d(TAG, "自动增益控制已释放")
}
} catch (e: Exception) {
Log.e(TAG, "清理回音消除失败: ${e.message}", e)
}
}
/// 恢复音频模式
private fun restoreAudioMode() {
try {
audioManager?.mode = originalAudioMode
Log.d(TAG, "音频模式已恢复为: $originalAudioMode")
} catch (e: Exception) {
Log.e(TAG, "恢复音频模式失败: ${e.message}", e)
}
}
fun destroy() {
Log.d(TAG, "销毁AudioManager")
stopRecording()
stopPlaying()
scope.cancel() // 确保所有协程都已取消
audioRecord?.release()
audioTrack?.release()
releaseAudioObjects()
playbackQueue.clear()
Log.i(TAG, "音频管理器已释放")
}
/// 计算音频数据的RMS值
private fun calculateRMS(audioData: ByteArray): Float {
if (audioData.isEmpty()) return 0f
var sum = 0.0
// 处理16位PCM数据
for (i in 0 until audioData.size step 2) {
if (i + 1 < audioData.size) {
// Little-endian: 低字节在前
val sample = (audioData[i].toInt() and 0xFF) or
(audioData[i + 1].toInt() shl 8)
// 转换为有符号16位整数
val signedSample = sample.toShort().toFloat()
// 归一化到[-1, 1]
val normalizedSample = signedSample / 32768f
sum += normalizedSample * normalizedSample
}
}
val mean = sum / (audioData.size / 2)
return sqrt(mean).toFloat()
}
}

455
local_plugins/realtime/android/src/main/kotlin/com/yunqiinnovation/realtime/RealtimePlugin.kt

@ -0,0 +1,455 @@
package com.yunqiinnovation.realtime
import android.content.Context
import android.os.Handler
import android.os.Looper
import android.util.Log
import io.flutter.embedding.engine.plugins.FlutterPlugin
import io.flutter.plugin.common.EventChannel
import io.flutter.plugin.common.MethodCall
import io.flutter.plugin.common.MethodChannel
import io.flutter.plugin.common.MethodChannel.MethodCallHandler
import io.flutter.plugin.common.MethodChannel.Result
import kotlinx.coroutines.*
import org.json.JSONObject
/**
* 实时语音聊天插件
*
* 通过WebSocket连接Vocode服务器实现实时语音交互
* 包含录音、播放、WebSocket通信等功能
*/
class RealtimePlugin: FlutterPlugin, MethodCallHandler, EventChannel.StreamHandler {
companion object {
private const val TAG = "RealtimePlugin"
private const val METHOD_CHANNEL = "realtime/methods"
private const val EVENT_CHANNEL = "realtime/events"
}
private lateinit var context: Context
private lateinit var methodChannel: MethodChannel
private lateinit var eventChannel: EventChannel
private var eventSink: EventChannel.EventSink? = null
// 核心组件:设为可空,因为它们的生命周期与页面会话绑定
private var audioManager: RealtimeAudioManager? = null
private var webSocketManager: RealtimeWebSocketManager? = null
// 协程作用域:需要可变,因为在重用插件时需要重新创建
private var scope = CoroutineScope(Dispatchers.Default + SupervisorJob())
// 主线程Handler,用于向Flutter发送事件与回调
private val mainHandler = Handler(Looper.getMainLooper())
// 配置参数
private var serverUrl: String = ""
private var sampleRate: Int = 16000
private var channels: Int = 1
private var bitsPerSample: Int = 16
// 状态
@Volatile
private var connectionStatus: ConnectionStatus = ConnectionStatus.DISCONNECTED
@Volatile
private var voiceStatus: VoiceStatus = VoiceStatus.IDLE
// 清理状态标志,确保清理逻辑只执行一次
@Volatile
private var isCleanedUp = false
private val cleanupLock = Any()
override fun onAttachedToEngine(flutterPluginBinding: FlutterPlugin.FlutterPluginBinding) {
context = flutterPluginBinding.applicationContext
methodChannel = MethodChannel(flutterPluginBinding.binaryMessenger, METHOD_CHANNEL)
methodChannel.setMethodCallHandler(this)
eventChannel = EventChannel(flutterPluginBinding.binaryMessenger, EVENT_CHANNEL)
eventChannel.setStreamHandler(this)
Log.i(TAG, "Realtime插件已附加到引擎")
}
override fun onDetachedFromEngine(binding: FlutterPlugin.FlutterPluginBinding) {
Log.i(TAG, "插件正在从引擎分离,执行清理...")
performCleanup()
methodChannel.setMethodCallHandler(null)
eventChannel.setStreamHandler(null)
Log.i(TAG, "Realtime插件已从引擎分离")
}
/// 设置各组件的回调
private fun setupCallbacks() {
Log.d(TAG, "setupCallbacks: 正在为新的管理器实例设置回调...")
// 音频管理器回调
audioManager?.onAudioData = { data ->
// 发送到WebSocket
webSocketManager?.sendAudioData(data)
// 发送录音PCM数据到Flutter用于可视化
val pcm16Data = convertByteArrayToPCM16(data)
sendEvent("recordingPcmData", pcm16Data)
}
audioManager?.onStatusChanged = { status ->
voiceStatus = status
sendEvent("voiceStatusChanged", status.value)
}
// 新增RMS回调
audioManager?.onRmsChanged = { rms ->
sendEvent("audioRms", rms)
}
// WebSocket管理器回调
webSocketManager?.onConnectionStatusChanged = { status ->
connectionStatus = status
sendEvent("connectionStatusChanged", status.value)
}
webSocketManager?.onAudioReceived = { data ->
// 添加音频数据前50字节的日志输出
val dataPreview = if (data.size > 50) {
data.sliceArray(0..49).joinToString(" ") { "%02x".format(it) } + "..."
} else {
data.joinToString(" ") { "%02x".format(it) }
}
Log.d(TAG, "收到音频数据前50字节: $dataPreview")
Log.d(TAG, "音频数据总长度: ${data.size} 字节")
audioManager?.playAudioData(data)
// 发送播放PCM数据到Flutter用于可视化
val pcm16Data = convertByteArrayToPCM16(data)
sendEvent("playbackPcmData", pcm16Data)
}
webSocketManager?.onTextReceived = { text ->
// 直接发送原始JSON文本到Flutter层处理
sendEvent("textReceived", text)
}
webSocketManager?.onError = { error ->
sendEvent("error", error)
}
}
/// 发送事件到Flutter端
private fun sendEvent(type: String, data: Any?) {
val sink = eventSink
if (sink == null) {
Log.d(TAG, "事件流已关闭,忽略事件: $type")
return
}
val event = mapOf(
"type" to type,
"data" to data
)
try {
mainHandler.post {
try {
sink.success(event)
} catch (e: Exception) {
Log.w(TAG, "发送事件失败: \\${e.message}")
}
}
} catch (e: Exception) {
Log.w(TAG, "发送事件时出现异常: \\${e.message}")
}
}
override fun onMethodCall(call: MethodCall, result: Result) {
when (call.method) {
"initialize" -> handleInitialize(call, result)
"connect" -> handleConnect(result)
"disconnect" -> handleDisconnect(result)
"startRecording" -> handleStartRecording(result)
"stopRecording" -> handleStopRecording(result)
"stopPlaying" -> handleStopPlaying(result)
"getConnectionStatus" -> result.success(connectionStatus.value)
"getVoiceStatus" -> result.success(voiceStatus.value)
"sendTextMessage" -> handleSendTextMessage(call, result)
"setAudioConfig" -> handleSetAudioConfig(call, result)
"dispose" -> handleDispose(result)
else -> result.notImplemented()
}
}
private fun handleInitialize(call: MethodCall, result: Result) {
try {
Log.d(TAG, "handleInitialize: 开始新一轮的初始化...")
// 1. 如果之前的协程作用域已被取消,则创建一个全新的
if (!scope.isActive) {
scope = CoroutineScope(Dispatchers.Default + SupervisorJob())
Log.i(TAG, "handleInitialize: 检测到协程作用域已失效,已创建新实例。")
}
// 2. 重置清理状态标志,允许下一次的清理操作
synchronized(cleanupLock) {
if (isCleanedUp) {
isCleanedUp = false
Log.i(TAG, "handleInitialize: 清理状态已重置,插件可再次使用。")
}
}
// 3. 为新会话创建全新的管理器实例
Log.d(TAG, "handleInitialize: 正在创建新的AudioManager和WebSocketManager实例...")
audioManager = RealtimeAudioManager(context)
webSocketManager = RealtimeWebSocketManager()
// 4. 为新实例设置回调
setupCallbacks()
val args = call.arguments as? Map<String, Any>
val serverUrl = args?.get("serverUrl") as? String
?: return result.error("INVALID_ARGUMENTS", "服务器地址不能为空", null)
this.serverUrl = serverUrl
args["sampleRate"]?.let { this.sampleRate = (it as Number).toInt() }
args["channels"]?.let { this.channels = (it as Number).toInt() }
args["bitsPerSample"]?.let { this.bitsPerSample = (it as Number).toInt() }
// 初始化音频管理器
val audioConfig = AudioConfig(sampleRate, channels, bitsPerSample)
val success = audioManager?.initialize(audioConfig) ?: false
if (success) {
webSocketManager?.initialize(serverUrl)
Log.i(TAG, "实时语音插件初始化成功")
} else {
Log.e(TAG, "实时语音插件初始化失败")
}
result.success(success)
} catch (e: Exception) {
Log.e(TAG, "初始化失败: ${e.message}", e)
result.error("INITIALIZATION_ERROR", "初始化失败: ${e.message}", null)
}
}
private fun handleConnect(result: Result) {
Log.d(TAG, "handleConnect: 收到连接请求。scope是否活跃? ${scope.isActive}")
scope.launch {
try {
val success = webSocketManager?.connect() ?: false
mainHandler.post { result.success(success) }
} catch (e: Exception) {
Log.e(TAG, "连接失败: ${e.message}", e)
mainHandler.post { result.error("CONNECTION_ERROR", "连接失败: ${e.message}", null) }
}
}
}
private fun handleDisconnect(result: Result) {
scope.launch {
try {
webSocketManager?.disconnect()
mainHandler.post { result.success(true) }
} catch (e: Exception) {
Log.e(TAG, "断开连接失败: ${e.message}", e)
mainHandler.post { result.error("DISCONNECTION_ERROR", "断开连接失败: ${e.message}", null) }
}
}
}
private fun handleStartRecording(result: Result) {
scope.launch {
try {
val success = audioManager?.startRecording() ?: false
mainHandler.post { result.success(success) }
} catch (e: Exception) {
Log.e(TAG, "开始录音失败: ${e.message}", e)
mainHandler.post { result.error("RECORDING_ERROR", "开始录音失败: ${e.message}", null) }
}
}
}
private fun handleStopRecording(result: Result) {
scope.launch {
try {
val success = audioManager?.stopRecording() ?: false
mainHandler.post { result.success(success) }
} catch (e: Exception) {
Log.e(TAG, "停止录音失败: ${e.message}", e)
mainHandler.post { result.error("RECORDING_ERROR", "停止录音失败: ${e.message}", null) }
}
}
}
private fun handleStopPlaying(result: Result) {
scope.launch {
try {
val success = audioManager?.stopPlaying() ?: false
mainHandler.post { result.success(success) }
} catch (e: Exception) {
Log.e(TAG, "停止播放失败: ${e.message}", e)
mainHandler.post { result.error("PLAYBACK_ERROR", "停止播放失败: ${e.message}", null) }
}
}
}
private fun handleSendTextMessage(call: MethodCall, result: Result) {
try {
val args = call.arguments as? Map<String, Any>
val message = args?.get("message") as? String
?: return result.error("INVALID_ARGUMENTS", "消息内容不能为空", null)
val success = webSocketManager?.sendTextMessage(message) ?: false
result.success(success)
} catch (e: Exception) {
Log.e(TAG, "发送文本消息失败: ${e.message}", e)
result.error("MESSAGE_ERROR", "发送文本消息失败: ${e.message}", null)
}
}
private fun handleSetAudioConfig(call: MethodCall, result: Result) {
try {
val args = call.arguments as? Map<String, Any>
?: return result.error("INVALID_ARGUMENTS", "参数无效", null)
Log.d(TAG, "收到的音频配置参数: $args")
args["sampleRate"]?.let {
Log.d(TAG, "sampleRate: $it, type: ${it.javaClass.name}")
this.sampleRate = (it as Number).toInt()
}
args["channels"]?.let {
Log.d(TAG, "channels: $it, type: ${it.javaClass.name}")
this.channels = (it as Number).toInt()
}
args["bitsPerSample"]?.let {
Log.d(TAG, "bitsPerSample: $it, type: ${it.javaClass.name}")
this.bitsPerSample = (it as Number).toInt()
}
val audioConfig = AudioConfig(sampleRate, channels, bitsPerSample)
val success = audioManager?.updateConfig(audioConfig) ?: false
result.success(success)
} catch (e: Exception) {
Log.e(TAG, "设置音频配置失败: ${e.message}", e)
result.error("CONFIG_ERROR", "设置音频配置失败: ${e.message}", null)
}
}
private fun handleDispose(result: Result) {
Log.i(TAG, "Flutter端请求释放资源...")
performCleanup()
// 立即返回,不等待异步清理完成
result.success(null)
}
/**
* 将字节数组转换为PCM16整数数组
* PCM16格式:16-bit signed integers, little-endian
*/
private fun convertByteArrayToPCM16(data: ByteArray): List<Int> {
val pcm16List = mutableListOf<Int>()
// 确保数据长度是偶数(每2个字节一个采样)
val validSize = data.size and 0xFFFFFFFE.toInt() // 清除最低位,确保偶数
// 每2个字节组成一个16-bit采样
for (i in 0 until validSize step 2) {
if (i + 1 < data.size) { // 安全检查
// Little-endian: 低字节在前
val low = data[i].toInt() and 0xFF
val high = data[i + 1].toInt() shl 8
val sample = (high or low).toShort().toInt() // 转换为有符号16位整数
pcm16List.add(sample)
}
}
return pcm16List
}
/**
* 执行核心清理操作。
* 此方法是幂等的(可重复调用)和线程安全的。
*/
private fun performCleanup() {
synchronized(cleanupLock) {
if (isCleanedUp) {
Log.d(TAG, "资源已清理,跳过。")
return
}
Log.i(TAG, "开始执行核心资源清理...")
// 立即标记,防止重入
isCleanedUp = true
// 立即停止接收事件
eventSink = null
// 在后台IO线程中执行所有耗时操作,避免阻塞主线程
CoroutineScope(Dispatchers.IO).launch {
runCatching { audioManager?.stopRecording() }
.onFailure { Log.w(TAG, "停止录音时异常: ${it.message}") }
runCatching { audioManager?.stopPlaying() }
.onFailure { Log.w(TAG, "停止播放时异常: ${it.message}") }
runCatching { webSocketManager?.disconnect() }
.onFailure { Log.w(TAG, "断开WebSocket时异常: ${it.message}") }
// 短暂延迟,以确保挂起的操作有时间完成
// delay(100L)
runCatching { audioManager?.dispose() }
.onFailure { Log.w(TAG, "释放音频管理器时异常: ${it.message}") }
runCatching { webSocketManager?.dispose() }
.onFailure { Log.w(TAG, "释放WebSocket管理器时异常: ${it.message}") }
runCatching {
scope.cancel()
Log.i(TAG, "插件协程作用域已取消。")
}.onFailure { Log.w(TAG, "取消协程作用域时异常: ${it.message}") }
// 清理完成,将实例置空
audioManager = null
webSocketManager = null
Log.i(TAG, "核心资源清理完成。")
}
}
}
// EventChannel.StreamHandler 实现
override fun onListen(arguments: Any?, events: EventChannel.EventSink?) {
eventSink = events
Log.i(TAG, "事件流监听已开始")
}
override fun onCancel(arguments: Any?) {
eventSink = null
Log.i(TAG, "事件流监听已取消")
}
}
// 枚举定义
enum class ConnectionStatus(val value: String) {
DISCONNECTED("disconnected"),
CONNECTING("connecting"),
CONNECTED("connected"),
ERROR("error")
}
enum class VoiceStatus(val value: String) {
IDLE("idle"),
RECORDING("recording"),
PROCESSING("processing"),
PLAYING("playing")
}
// 音频配置
data class AudioConfig(
val sampleRate: Int,
val channels: Int,
val bitsPerSample: Int
) {
val frameSize: Int
get() = (sampleRate * channels * bitsPerSample / 8 * 20) / 1000 // 20ms帧
}

303
local_plugins/realtime/android/src/main/kotlin/com/yunqiinnovation/realtime/RealtimeWebSocketManager.kt

@ -0,0 +1,303 @@
package com.yunqiinnovation.realtime
import android.util.Log
import kotlinx.coroutines.*
import okhttp3.*
import okio.ByteString
import java.util.concurrent.TimeUnit
/**
* 实时WebSocket管理器
*
* 负责与Vocode服务器的WebSocket通信,对应Android示例中的OkHttp WebSocket功能
*/
class RealtimeWebSocketManager : WebSocketListener() {
companion object {
private const val TAG = "RealtimeWebSocketManager"
private const val CONNECT_TIMEOUT = 30L
private const val READ_TIMEOUT = 60L
private const val WRITE_TIMEOUT = 30L
}
// OkHttp客户端和WebSocket
private var okHttpClient: OkHttpClient? = null
private var webSocket: WebSocket? = null
// 服务器URL
private var serverUrl: String = ""
// 连接状态
@Volatile
private var isConnected = false
@Volatile
private var isConnecting = false
// 回调
var onConnectionStatusChanged: ((ConnectionStatus) -> Unit)? = null
var onAudioReceived: ((ByteArray) -> Unit)? = null
var onTextReceived: ((String) -> Unit)? = null
var onError: ((String) -> Unit)? = null
// 协程作用域
private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob())
/// 初始化WebSocket管理器
fun initialize(serverUrl: String) {
Log.d(TAG, "initialize: 开始初始化。当前httpClient是否为null? ${okHttpClient == null}")
this.serverUrl = serverUrl
// 创建OkHttp客户端
okHttpClient = OkHttpClient.Builder()
.connectTimeout(CONNECT_TIMEOUT, TimeUnit.SECONDS)
.readTimeout(READ_TIMEOUT, TimeUnit.SECONDS)
.writeTimeout(WRITE_TIMEOUT, TimeUnit.SECONDS)
.retryOnConnectionFailure(true)
.build()
Log.i(TAG, "initialize: WebSocket管理器初始化完成。服务器地址: $serverUrl")
}
/// 连接到服务器
suspend fun connect(): Boolean = withContext(Dispatchers.IO) {
Log.d(TAG, "connect: 尝试连接。当前状态 - isConnected: $isConnected, isConnecting: $isConnecting")
if (isConnected || isConnecting) {
Log.w(TAG, "connect: 连接请求被拒绝,因为已有连接或正在连接中。")
return@withContext false
}
val client = okHttpClient ?: run {
Log.e(TAG, "connect: 连接失败,因为OkHttp客户端为null。请确保已调用initialize()。")
onError?.invoke("客户端未初始化")
return@withContext false
}
Log.d(TAG, "connect: OkHttpClient实例存在,准备发起新连接。")
try {
isConnecting = true
onConnectionStatusChanged?.invoke(ConnectionStatus.CONNECTING)
val request = Request.Builder()
.url(serverUrl)
.build()
Log.d(TAG, "connect: 正在创建新的WebSocket实例...")
webSocket = client.newWebSocket(request, this@RealtimeWebSocketManager)
// 等待连接建立
var waitTime = 0
while (isConnecting && waitTime < CONNECT_TIMEOUT * 1000) {
delay(100)
waitTime += 100
}
if (isConnected) {
Log.i(TAG, "connect: WebSocket连接成功")
return@withContext true
} else {
Log.e(TAG, "connect: WebSocket连接超时")
isConnecting = false
onConnectionStatusChanged?.invoke(ConnectionStatus.ERROR)
return@withContext false
}
} catch (e: Exception) {
Log.e(TAG, "WebSocket连接失败: ${e.message}", e)
isConnecting = false
onConnectionStatusChanged?.invoke(ConnectionStatus.ERROR)
onError?.invoke("连接失败: ${e.message}")
return@withContext false
}
}
/// 断开连接
fun disconnect() {
try {
isConnected = false
isConnecting = false
webSocket?.close(1000, "正常关闭")
webSocket = null
onConnectionStatusChanged?.invoke(ConnectionStatus.DISCONNECTED)
Log.i(TAG, "WebSocket已断开连接")
} catch (e: Exception) {
Log.e(TAG, "断开WebSocket连接失败: ${e.message}", e)
}
}
/// 发送音频数据
fun sendAudioData(data: ByteArray) {
if (!isConnected) {
return
}
try {
// 将音频数据编码为Base64并包装成Vocode WebSocket格式
val base64AudioData = android.util.Base64.encodeToString(data, android.util.Base64.NO_WRAP)
val audioMessage = org.json.JSONObject().apply {
put("type", "websocket_audio")
put("data", base64AudioData)
}
val success = webSocket?.send(audioMessage.toString()) ?: false
if (!success) {
Log.w(TAG, "发送音频数据失败")
}
} catch (e: Exception) {
Log.e(TAG, "发送音频数据异常: ${e.message}", e)
}
}
/// 发送文本消息
fun sendTextMessage(message: String): Boolean {
if (!isConnected) {
Log.w(TAG, "WebSocket未连接,无法发送文本消息")
return false
}
return try {
val success = webSocket?.send(message) ?: false
if (!success) {
Log.w(TAG, "发送文本消息失败")
}
success
} catch (e: Exception) {
Log.e(TAG, "发送文本消息异常: ${e.message}", e)
false
}
}
/// 释放资源
fun dispose() {
disconnect()
// 立即取消scope,不等待
try {
scope.cancel()
} catch (e: Exception) {
Log.w(TAG, "取消scope时出现异常: ${e.message}")
}
// 立即强制关闭WebSocket
try {
webSocket?.close(1001, "强制关闭")
webSocket = null
} catch (e: Exception) {
Log.w(TAG, "强制关闭WebSocket时出现异常: ${e.message}")
}
// 异步关闭线程池,设置更短的超时时间
okHttpClient?.dispatcher?.executorService?.let { executor ->
Thread {
try {
executor.shutdownNow() // 立即强制关闭
// 最多等待1秒,避免长时间等待
if (!executor.awaitTermination(1, TimeUnit.SECONDS)) {
Log.w(TAG, "线程池强制关闭超时,放弃等待")
}
} catch (e: Exception) {
Log.w(TAG, "关闭线程池时出现异常: ${e.message}")
}
}.start()
}
okHttpClient = null
Log.i(TAG, "WebSocket管理器已释放")
}
// WebSocketListener 回调实现
override fun onOpen(webSocket: WebSocket, response: Response) {
Log.i(TAG, "WebSocket连接已建立")
isConnected = true
isConnecting = false
onConnectionStatusChanged?.invoke(ConnectionStatus.CONNECTED)
}
override fun onMessage(webSocket: WebSocket, text: String) {
// 添加收到文本消息的日志输出,显示消息类型和前20字符
val textPreview = if (text.length > 50) text.substring(0, 50) + "..." else text
Log.d(TAG, "收到文本消息前20字符: $textPreview")
try {
// 尝试解析为JSON来检查消息类型
val json = org.json.JSONObject(text)
val messageType = json.optString("type", "")
// 添加消息类型的日志输出
Log.d(TAG, "消息类型: $messageType")
when (messageType) {
"websocket_audio" -> {
// 处理音频消息
val audioDataString = json.optString("data", "")
if (audioDataString.isNotEmpty()) {
// 将Base64编码的音频数据解码为字节数组
val audioData = android.util.Base64.decode(audioDataString, android.util.Base64.DEFAULT)
onAudioReceived?.invoke(audioData)
}
}
"websocket_transcript" -> {
// 处理统一的文本消息格式
// 包含sender、is_final、timestamp等字段
onTextReceived?.invoke(text)
}
"websocket_ready" -> {
// 连接就绪消息
Log.i(TAG, "WebSocket连接就绪")
}
"websocket_stop" -> {
// 会话停止消息
Log.i(TAG, "会话已结束")
}
"websocket_start" -> {
// 会话开始消息
Log.i(TAG, "会话已开始")
}
else -> {
// 其他消息类型,当作普通文本消息处理
onTextReceived?.invoke(text)
}
}
} catch (e: Exception) {
// 如果不是JSON格式,当作普通文本消息处理
Log.d(TAG, "非JSON格式消息,当作普通文本处理")
onTextReceived?.invoke(text)
}
}
override fun onMessage(webSocket: WebSocket, bytes: ByteString) {
val audioData = bytes.toByteArray()
onAudioReceived?.invoke(audioData)
}
override fun onClosing(webSocket: WebSocket, code: Int, reason: String) {
Log.i(TAG, "WebSocket正在关闭,代码: $code, 原因: $reason")
webSocket.close(1000, null)
}
override fun onClosed(webSocket: WebSocket, code: Int, reason: String) {
Log.i(TAG, "WebSocket已关闭,代码: $code, 原因: $reason")
isConnected = false
isConnecting = false
onConnectionStatusChanged?.invoke(ConnectionStatus.DISCONNECTED)
}
override fun onFailure(webSocket: WebSocket, t: Throwable, response: Response?) {
val errorMessage = when {
t.message?.contains("Unable to parse TLS packet header") == true ->
"SSL错误:请检查服务器是否支持wss://,本地开发请使用ws://"
t.message?.contains("Connection refused") == true ->
"连接被拒绝:请检查服务器是否运行在指定端口"
t.message?.contains("timeout") == true ->
"连接超时:请检查网络连接和服务器地址"
else -> "连接失败: ${t.message}"
}
Log.e(TAG, "WebSocket连接失败: ${t.message}", t)
isConnected = false
isConnecting = false
onConnectionStatusChanged?.invoke(ConnectionStatus.ERROR)
onError?.invoke(errorMessage)
}
}

274
local_plugins/realtime/example.md

@ -0,0 +1,274 @@
# Realtime Plugin 使用示例
## 基本用法
```dart
import 'package:realtime/realtime.dart';
class VoiceCallPage extends StatefulWidget {
@override
_VoiceCallPageState createState() => _VoiceCallPageState();
}
class _VoiceCallPageState extends State<VoiceCallPage> {
final RealtimeService _realtimeService = RealtimeService();
StreamSubscription<RealtimeEvent>? _eventSubscription;
ConnectionStatus _connectionStatus = ConnectionStatus.disconnected;
VoiceStatus _voiceStatus = VoiceStatus.idle;
@override
void initState() {
super.initState();
_setupEventListener();
_initializeService();
}
void _setupEventListener() {
_eventSubscription = _realtimeService.eventStream.listen((event) {
switch (event.type) {
case 'connectionStatusChanged':
setState(() {
_connectionStatus = _parseConnectionStatus(event.data);
});
break;
case 'voiceStatusChanged':
setState(() {
_voiceStatus = _parseVoiceStatus(event.data);
});
break;
case 'textReceived':
print('收到文本: ${event.data}');
break;
case 'error':
print('错误: ${event.data}');
break;
}
});
}
Future<void> _initializeService() async {
try {
final success = await _realtimeService.initialize(
serverUrl: 'wss://your-vocode-server.com/ws',
sampleRate: 16000,
channels: 1,
bitsPerSample: 16,
);
if (success) {
print('Realtime服务初始化成功');
} else {
print('Realtime服务初始化失败');
}
} catch (e) {
print('初始化错误: $e');
}
}
Future<void> _connect() async {
try {
final success = await _realtimeService.connect();
if (!success) {
print('连接失败');
}
} catch (e) {
print('连接错误: $e');
}
}
Future<void> _disconnect() async {
try {
await _realtimeService.disconnect();
} catch (e) {
print('断开连接错误: $e');
}
}
Future<void> _startRecording() async {
try {
final success = await _realtimeService.startRecording();
if (!success) {
print('开始录音失败');
}
} catch (e) {
print('录音错误: $e');
}
}
Future<void> _stopRecording() async {
try {
final success = await _realtimeService.stopRecording();
if (!success) {
print('停止录音失败');
}
} catch (e) {
print('停止录音错误: $e');
}
}
ConnectionStatus _parseConnectionStatus(String status) {
switch (status) {
case 'connected': return ConnectionStatus.connected;
case 'connecting': return ConnectionStatus.connecting;
case 'disconnected': return ConnectionStatus.disconnected;
case 'error': return ConnectionStatus.error;
default: return ConnectionStatus.disconnected;
}
}
VoiceStatus _parseVoiceStatus(String status) {
switch (status) {
case 'recording': return VoiceStatus.recording;
case 'processing': return VoiceStatus.processing;
case 'playing': return VoiceStatus.playing;
case 'idle': return VoiceStatus.idle;
default: return VoiceStatus.idle;
}
}
@override
Widget build(BuildContext context) {
return Scaffold(
appBar: AppBar(
title: Text('实时语音聊天'),
),
body: Center(
child: Column(
mainAxisAlignment: MainAxisAlignment.center,
children: [
// 连接状态显示
Text('连接状态: ${_connectionStatus.name}'),
SizedBox(height: 20),
// 语音状态显示
Text('语音状态: ${_voiceStatus.name}'),
SizedBox(height: 40),
// 连接按钮
ElevatedButton(
onPressed: _connectionStatus == ConnectionStatus.disconnected
? _connect
: _disconnect,
child: Text(_connectionStatus == ConnectionStatus.disconnected
? '连接'
: '断开'),
),
SizedBox(height: 20),
// 录音按钮
ElevatedButton(
onPressed: _connectionStatus == ConnectionStatus.connected
? (_voiceStatus == VoiceStatus.recording
? _stopRecording
: _startRecording)
: null,
child: Text(_voiceStatus == VoiceStatus.recording
? '停止录音'
: '开始录音'),
),
],
),
),
);
}
@override
void dispose() {
_eventSubscription?.cancel();
_realtimeService.dispose();
super.dispose();
}
}
```
## GetX控制器中的使用
```dart
import 'package:get/get.dart';
import 'package:realtime/realtime.dart';
class VoiceController extends GetxController {
final RealtimeService _realtimeService = RealtimeService();
final RxBool isConnected = false.obs;
final RxBool isRecording = false.obs;
StreamSubscription<RealtimeEvent>? _eventSubscription;
@override
void onInit() {
super.onInit();
_setupEventListener();
_initializeService();
}
void _setupEventListener() {
_eventSubscription = _realtimeService.eventStream.listen((event) {
switch (event.type) {
case 'connectionStatusChanged':
isConnected.value = event.data == 'connected';
break;
case 'voiceStatusChanged':
isRecording.value = event.data == 'recording';
break;
}
});
}
Future<void> _initializeService() async {
await _realtimeService.initialize(
serverUrl: 'wss://your-server/ws',
);
}
Future<void> toggleConnection() async {
if (isConnected.value) {
await _realtimeService.disconnect();
} else {
await _realtimeService.connect();
}
}
Future<void> toggleRecording() async {
if (isRecording.value) {
await _realtimeService.stopRecording();
} else {
await _realtimeService.startRecording();
}
}
@override
void onClose() {
_eventSubscription?.cancel();
_realtimeService.dispose();
super.onClose();
}
}
```
## 权限配置
### iOS
在 `ios/Runner/Info.plist` 中添加:
```xml
<key>NSMicrophoneUsageDescription</key>
<string>应用需要麦克风权限进行语音录制</string>
```
## 服务器配置示例
需要一个支持WebSocket的Vocode服务器,具体实现可参考Vocode官方文档。
服务器需要:
1. 接收16kHz/16-bit/单声道的PCM音频数据
2. 返回相同格式的音频数据
3. 支持文本消息交换
## 注意事项
1. 确保网络连接稳定
2. 音频格式必须与服务器保持一致
3. 及时处理事件流中的错误
4. 在适当时机释放资源

18
local_plugins/realtime/ios/realtime/Package.swift

@ -0,0 +1,18 @@
// swift-tools-version: 5.9
import PackageDescription
let package = Package(
name: "realtime",
platforms: [.iOS("18.0")],
products: [
.library(name: "realtime", targets: ["realtime"])
],
dependencies: [],
targets: [
.target(
name: "realtime",
dependencies: [],
path: "Sources/realtime"
)
]
)

344
local_plugins/realtime/ios/realtime/Sources/realtime/RealtimeAudioManager.swift

@ -0,0 +1,344 @@
import AVFoundation
import os.log
import Accelerate
/**
* 实时音频管理器
*
* 负责音频的录制和播放,对应Android示例中的AudioRecord和AudioTrack功能
*/
class RealtimeAudioManager {
private let log = OSLog(subsystem: "com.yunqiinnovation.realtime", category: "RealtimeAudioManager")
// 音频配置
private var audioConfig: AudioConfig?
// 录音相关
private var audioEngine: AVAudioEngine?
private var inputNode: AVAudioInputNode?
private var recordingFormat: AVAudioFormat?
private var isRecording = false
// 播放相关
private var audioPlayer: AVAudioPlayerNode?
private var playbackFormat: AVAudioFormat?
private var playbackQueue: DispatchQueue
private var audioBufferQueue: [AVAudioPCMBuffer] = []
private var isPlaying = false
// 回调
var onAudioData: ((Data) -> Void)?
var onStatusChanged: ((VoiceStatus) -> Void)?
var onRmsChanged: ((Float) -> Void)?
// 录音缓冲区
private var recordingBuffer: AVAudioPCMBuffer?
private var frameSize: Int = 320 // 20ms @ 16kHz
init() {
playbackQueue = DispatchQueue(label: "com.realtime.playback", qos: .userInitiated)
setupAudioSession()
}
/// 设置音频会话
private func setupAudioSession() {
do {
let session = AVAudioSession.sharedInstance()
try session.setCategory(.playAndRecord,
mode: .default,
options: [.defaultToSpeaker, .allowBluetooth])
try session.setActive(true)
os_log("音频会话设置成功", log: log, type: .info)
} catch {
os_log("音频会话设置失败: %@", log: log, type: .error, error.localizedDescription)
}
}
/// 初始化音频管理器
func initialize(config: AudioConfig) -> Bool {
self.audioConfig = config
self.frameSize = config.frameSize
guard setupAudioEngine(config: config) else {
os_log("音频引擎初始化失败", log: log, type: .error)
return false
}
guard setupAudioPlayer(config: config) else {
os_log("音频播放器初始化失败", log: log, type: .error)
return false
}
os_log("音频管理器初始化成功", log: log, type: .info)
return true
}
/// 设置音频引擎(录音)
private func setupAudioEngine(config: AudioConfig) -> Bool {
audioEngine = AVAudioEngine()
guard let audioEngine = audioEngine else { return false }
inputNode = audioEngine.inputNode
// 创建录音格式:16kHz, 16-bit, 单声道
recordingFormat = AVAudioFormat(commonFormat: .pcmFormatInt16,
sampleRate: Double(config.sampleRate),
channels: AVAudioChannelCount(config.channels),
interleaved: true)
guard let recordingFormat = recordingFormat else {
os_log("无法创建录音格式", log: log, type: .error)
return false
}
// 创建录音缓冲区
let frameCount = AVAudioFrameCount(config.frameSize / (config.bitsPerSample / 8))
recordingBuffer = AVAudioPCMBuffer(pcmFormat: recordingFormat, frameCapacity: frameCount)
return true
}
/// 设置音频播放器
private func setupAudioPlayer(config: AudioConfig) -> Bool {
guard let audioEngine = audioEngine else { return false }
audioPlayer = AVAudioPlayerNode()
// 创建播放格式:16kHz, 16-bit, 单声道
playbackFormat = AVAudioFormat(commonFormat: .pcmFormatInt16,
sampleRate: Double(config.sampleRate),
channels: AVAudioChannelCount(config.channels),
interleaved: true)
guard let audioPlayer = audioPlayer,
let playbackFormat = playbackFormat else {
os_log("无法创建播放格式", log: log, type: .error)
return false
}
// 连接播放器到音频引擎
audioEngine.attach(audioPlayer)
audioEngine.connect(audioPlayer, to: audioEngine.outputNode, format: playbackFormat)
return true
}
/// 开始录音
func startRecording() -> Bool {
guard let audioEngine = audioEngine,
let inputNode = inputNode,
let recordingFormat = recordingFormat,
!isRecording else {
return false
}
// 安装录音tap
inputNode.installTap(onBus: 0, bufferSize: AVAudioFrameCount(frameSize / 2), format: recordingFormat) { [weak self] buffer, time in
self?.processRecordedAudio(buffer: buffer)
}
do {
try audioEngine.start()
isRecording = true
onStatusChanged?(.recording)
os_log("开始录音", log: log, type: .info)
return true
} catch {
os_log("开始录音失败: %@", log: log, type: .error, error.localizedDescription)
return false
}
}
/// 处理录制的音频数据
private func processRecordedAudio(buffer: AVAudioPCMBuffer) {
guard let channelData = buffer.int16ChannelData?[0],
buffer.frameLength > 0 else {
return
}
// 确保数据长度匹配frameSize
let dataLength = min(Int(buffer.frameLength) * 2, frameSize) // 2 bytes per sample
// 转换为Data
let data = Data(bytes: channelData, count: dataLength)
// 计算RMS值
let rms = calculateRMS(buffer: buffer)
onRmsChanged?(rms)
// 回调音频数据
onAudioData?(data)
}
/// 停止录音
func stopRecording() -> Bool {
guard let audioEngine = audioEngine,
let inputNode = inputNode,
isRecording else {
return false
}
inputNode.removeTap(onBus: 0)
audioEngine.stop()
isRecording = false
onStatusChanged?(.idle)
os_log("停止录音", log: log, type: .info)
return true
}
/// 播放音频数据
func playAudioData(_ data: Data) {
guard let playbackFormat = playbackFormat,
let audioPlayer = audioPlayer else {
return
}
playbackQueue.async { [weak self] in
self?.enqueueAudioData(data, format: playbackFormat, player: audioPlayer)
}
}
/// 将音频数据加入播放队列
private func enqueueAudioData(_ data: Data, format: AVAudioFormat, player: AVAudioPlayerNode) {
// 创建PCM缓冲区
let frameCount = AVAudioFrameCount(data.count / 2) // 16-bit = 2 bytes per sample
guard let buffer = AVAudioPCMBuffer(pcmFormat: format, frameCapacity: frameCount) else {
os_log("无法创建播放缓冲区", log: log, type: .error)
return
}
buffer.frameLength = frameCount
// 复制数据到缓冲区
guard let channelData = buffer.int16ChannelData?[0] else {
return
}
data.withUnsafeBytes { bytes in
let int16Pointer = bytes.bindMemory(to: Int16.self)
channelData.update(from: int16Pointer.baseAddress!, count: Int(frameCount))
}
// 计算播放数据的RMS值
let rms = calculateRMS(buffer: buffer)
DispatchQueue.main.async { [weak self] in
self?.onRmsChanged?(rms)
}
// 加入播放队列
audioBufferQueue.append(buffer)
// 如果没在播放,开始播放
if !isPlaying {
startPlayback()
}
}
/// 开始播放
private func startPlayback() {
guard let audioEngine = audioEngine,
let audioPlayer = audioPlayer,
!audioBufferQueue.isEmpty else {
return
}
if !audioEngine.isRunning {
do {
try audioEngine.start()
} catch {
os_log("启动音频引擎失败: %@", log: log, type: .error, error.localizedDescription)
return
}
}
if !audioPlayer.isPlaying {
audioPlayer.play()
}
isPlaying = true
onStatusChanged?(.playing)
// 播放队列中的音频
playNextBuffer()
}
/// 播放下一个缓冲区
private func playNextBuffer() {
guard let audioPlayer = audioPlayer,
!audioBufferQueue.isEmpty else {
isPlaying = false
onStatusChanged?(.idle)
return
}
let buffer = audioBufferQueue.removeFirst()
audioPlayer.scheduleBuffer(buffer) { [weak self] in
DispatchQueue.main.async {
self?.playNextBuffer()
}
}
}
/// 停止播放
func stopPlaying() -> Bool {
guard let audioPlayer = audioPlayer else {
return false
}
audioPlayer.stop()
audioBufferQueue.removeAll()
isPlaying = false
onStatusChanged?(.idle)
os_log("停止播放", log: log, type: .info)
return true
}
/// 更新配置
func updateConfig(config: AudioConfig) -> Bool {
// 如果正在录音或播放,先停止
if isRecording {
_ = stopRecording()
}
if isPlaying {
_ = stopPlaying()
}
return initialize(config: config)
}
/// 释放资源
func dispose() {
_ = stopRecording()
_ = stopPlaying()
audioEngine?.stop()
audioEngine = nil
audioPlayer = nil
audioBufferQueue.removeAll()
os_log("音频管理器已释放", log: log, type: .info)
}
/// 计算音频缓冲区的RMS值
private func calculateRMS(buffer: AVAudioPCMBuffer) -> Float {
guard let channelData = buffer.int16ChannelData?[0],
buffer.frameLength > 0 else {
return 0.0
}
let frameCount = Int(buffer.frameLength)
var rms: Float = 0.0
// 使用Accelerate框架高效计算RMS
vDSP_measqv(channelData, 1, &rms, vDSP_Length(frameCount))
// 计算平方根并归一化
rms = sqrt(rms / Float(frameCount)) / 32768.0
return rms
}
}

319
local_plugins/realtime/ios/realtime/Sources/realtime/RealtimePlugin.swift

@ -0,0 +1,319 @@
import Flutter
import UIKit
import AVFoundation
import os.log
/**
* 实时语音聊天插件
*
* 通过WebSocket连接Vocode服务器实现实时语音交互
* 包含录音、播放、WebSocket通信等功能
*/
@objc public class RealtimePlugin: NSObject, FlutterPlugin {
private let log = OSLog(subsystem: "com.yunqiinnovation.realtime", category: "RealtimePlugin")
// Flutter通道
private var methodChannel: FlutterMethodChannel?
private var eventChannel: FlutterEventChannel?
private var eventSink: FlutterEventSink?
// 核心组件
private let audioManager = RealtimeAudioManager()
private let websocketManager = RealtimeWebSocketManager()
// 配置参数
private var serverUrl: String = ""
private var sampleRate: Int = 16000
private var channels: Int = 1
private var bitsPerSample: Int = 16
// 状态
private var connectionStatus: ConnectionStatus = .disconnected
private var voiceStatus: VoiceStatus = .idle
/// 插件注册
public static func register(with registrar: FlutterPluginRegistrar) {
let instance = RealtimePlugin()
// 方法通道
let methodChannel = FlutterMethodChannel(name: "realtime/methods", binaryMessenger: registrar.messenger())
registrar.addMethodCallDelegate(instance, channel: methodChannel)
instance.methodChannel = methodChannel
// 事件通道
let eventChannel = FlutterEventChannel(name: "realtime/events", binaryMessenger: registrar.messenger())
eventChannel.setStreamHandler(instance)
instance.eventChannel = eventChannel
// 设置组件回调
instance.setupCallbacks()
}
/// 设置各组件的回调
private func setupCallbacks() {
// 音频管理器回调
audioManager.onAudioData = { [weak self] data in
// 发送到WebSocket
self?.websocketManager.sendAudioData(data)
// 同时发送PCM数据到Flutter用于可视化
self?.sendPCMDataToFlutter(data)
}
audioManager.onStatusChanged = { [weak self] status in
self?.voiceStatus = status
self?.sendEvent(type: "voiceStatusChanged", data: status.rawValue)
}
// 新增RMS回调
audioManager.onRmsChanged = { [weak self] rms in
self?.sendEvent(type: "audioRms", data: rms)
}
// WebSocket管理器回调
websocketManager.onConnectionStatusChanged = { [weak self] status in
self?.connectionStatus = status
self?.sendEvent(type: "connectionStatusChanged", data: status.rawValue)
}
websocketManager.onAudioReceived = { [weak self] data in
// 添加音频数据前50字节的日志输出
let dataCount = min(data.count, 50)
let dataPreview = data.prefix(dataCount).map { String(format: "%02x", $0) }.joined(separator: " ")
let previewText = data.count > 50 ? dataPreview + "..." : dataPreview
self?.os_log("收到音频数据前50字节: %@", log: self?.log ?? OSLog.default, type: .debug, previewText)
self?.os_log("音频数据总长度: %d 字节", log: self?.log ?? OSLog.default, type: .debug, data.count)
self?.audioManager.playAudioData(data)
// 同时发送接收到的PCM数据到Flutter用于可视化
self?.sendPCMDataToFlutter(data)
}
websocketManager.onTextReceived = { [weak self] text in
// 直接发送原始JSON文本到Flutter层处理
self?.sendEvent(type: "textReceived", data: text)
}
websocketManager.onError = { [weak self] error in
self?.sendEvent(type: "error", data: error)
}
}
/// 发送事件到Flutter端
private func sendEvent(type: String, data: Any?) {
guard let eventSink = eventSink else { return }
let event: [String: Any] = [
"type": type,
"data": data ?? NSNull()
]
DispatchQueue.main.async {
eventSink(event)
}
}
/// 将PCM数据发送到Flutter用于可视化
private func sendPCMDataToFlutter(_ data: Data) {
// 确保数据长度是偶数(每2个字节一个16位采样)
let validSize = data.count & ~1 // 清除最低位,确保偶数
guard validSize > 0 else { return }
// 将Data转换为Int16数组 (PCM16格式)
let pcm16Array = data.prefix(validSize).withUnsafeBytes { bytes in
let int16Pointer = bytes.bindMemory(to: Int16.self)
return Array(UnsafeBufferPointer(start: int16Pointer.baseAddress, count: validSize / 2))
}
// 转换为Int数组以便传递给Flutter
let pcmData = pcm16Array.map { Int($0) }
sendEvent(type: "pcmData", data: pcmData)
}
/// 辅助日志方法
private func os_log(_ message: String, log: OSLog, type: OSLogType, _ args: CVarArg...) {
if args.isEmpty {
os_log("%@", log: log, type: type, message)
} else {
os_log(message, log: log, type: type, args)
}
}
}
// MARK: - FlutterPlugin
extension RealtimePlugin {
public func handle(_ call: FlutterMethodCall, result: @escaping FlutterResult) {
switch call.method {
case "initialize":
handleInitialize(call, result)
case "connect":
handleConnect(result)
case "disconnect":
handleDisconnect(result)
case "startRecording":
handleStartRecording(result)
case "stopRecording":
handleStopRecording(result)
case "stopPlaying":
handleStopPlaying(result)
case "getConnectionStatus":
result(connectionStatus.rawValue)
case "getVoiceStatus":
result(voiceStatus.rawValue)
case "sendTextMessage":
handleSendTextMessage(call, result)
case "setAudioConfig":
handleSetAudioConfig(call, result)
case "dispose":
handleDispose(result)
default:
result(FlutterMethodNotImplemented)
}
}
private func handleInitialize(_ call: FlutterMethodCall, _ result: @escaping FlutterResult) {
guard let args = call.arguments as? [String: Any],
let serverUrl = args["serverUrl"] as? String else {
result(FlutterError(code: "INVALID_ARGUMENTS", message: "服务器地址不能为空", details: nil))
return
}
self.serverUrl = serverUrl
self.sampleRate = args["sampleRate"] as? Int ?? 16000
self.channels = args["channels"] as? Int ?? 1
self.bitsPerSample = args["bitsPerSample"] as? Int ?? 16
// 初始化音频管理器
let audioConfig = AudioConfig(
sampleRate: sampleRate,
channels: channels,
bitsPerSample: bitsPerSample
)
let success = audioManager.initialize(config: audioConfig)
if success {
websocketManager.initialize(serverUrl: serverUrl)
os_log("实时语音插件初始化成功", log: log, type: .info)
} else {
os_log("实时语音插件初始化失败", log: log, type: .error)
}
result(success)
}
private func handleConnect(_ result: @escaping FlutterResult) {
websocketManager.connect { [weak self] success in
DispatchQueue.main.async {
result(success)
}
}
}
private func handleDisconnect(_ result: @escaping FlutterResult) {
websocketManager.disconnect()
result(true)
}
private func handleStartRecording(_ result: @escaping FlutterResult) {
let success = audioManager.startRecording()
result(success)
}
private func handleStopRecording(_ result: @escaping FlutterResult) {
let success = audioManager.stopRecording()
result(success)
}
private func handleStopPlaying(_ result: @escaping FlutterResult) {
let success = audioManager.stopPlaying()
result(success)
}
private func handleSendTextMessage(_ call: FlutterMethodCall, _ result: @escaping FlutterResult) {
guard let args = call.arguments as? [String: Any],
let message = args["message"] as? String else {
result(FlutterError(code: "INVALID_ARGUMENTS", message: "消息内容不能为空", details: nil))
return
}
let success = websocketManager.sendTextMessage(message)
result(success)
}
private func handleSetAudioConfig(_ call: FlutterMethodCall, _ result: @escaping FlutterResult) {
guard let args = call.arguments as? [String: Any] else {
result(FlutterError(code: "INVALID_ARGUMENTS", message: "参数无效", details: nil))
return
}
if let sampleRate = args["sampleRate"] as? Int {
self.sampleRate = sampleRate
}
if let channels = args["channels"] as? Int {
self.channels = channels
}
if let bitsPerSample = args["bitsPerSample"] as? Int {
self.bitsPerSample = bitsPerSample
}
let audioConfig = AudioConfig(
sampleRate: sampleRate,
channels: channels,
bitsPerSample: bitsPerSample
)
let success = audioManager.updateConfig(config: audioConfig)
result(success)
}
private func handleDispose(_ result: @escaping FlutterResult) {
audioManager.dispose()
websocketManager.dispose()
result(nil)
}
}
// MARK: - FlutterStreamHandler
extension RealtimePlugin: FlutterStreamHandler {
public func onListen(withArguments arguments: Any?, eventSink events: @escaping FlutterEventSink) -> FlutterError? {
self.eventSink = events
return nil
}
public func onCancel(withArguments arguments: Any?) -> FlutterError? {
self.eventSink = nil
return nil
}
}
// MARK: - 枚举定义
enum ConnectionStatus: String {
case disconnected = "disconnected"
case connecting = "connecting"
case connected = "connected"
case error = "error"
}
enum VoiceStatus: String {
case idle = "idle"
case recording = "recording"
case processing = "processing"
case playing = "playing"
}
// MARK: - 音频配置
struct AudioConfig {
let sampleRate: Int
let channels: Int
let bitsPerSample: Int
var frameSize: Int {
// 20ms的帧大小
return (sampleRate * channels * bitsPerSample / 8 * 20) / 1000
}
}

323
local_plugins/realtime/ios/realtime/Sources/realtime/RealtimeWebSocketManager.swift

@ -0,0 +1,323 @@
import Foundation
import Network
import os.log
/**
* 实时WebSocket管理器
*
* 负责与Vocode服务器的WebSocket通信,对应Android示例中的OkHttp WebSocket功能
*/
class RealtimeWebSocketManager: NSObject {
private let log = OSLog(subsystem: "com.yunqiinnovation.realtime", category: "RealtimeWebSocketManager")
// WebSocket连接
private var webSocketTask: URLSessionWebSocketTask?
private var urlSession: URLSession?
// 服务器URL
private var serverUrl: String = ""
// 连接状态
private var isConnected = false
private var isConnecting = false
// 回调
var onConnectionStatusChanged: ((ConnectionStatus) -> Void)?
var onAudioReceived: ((Data) -> Void)?
var onTextReceived: ((String) -> Void)?
var onError: ((String) -> Void)?
// 发送队列
private let sendQueue = DispatchQueue(label: "com.realtime.websocket.send", qos: .userInitiated)
private let receiveQueue = DispatchQueue(label: "com.realtime.websocket.receive", qos: .userInitiated)
override init() {
super.init()
setupURLSession()
}
/// 设置URL会话
private func setupURLSession() {
let config = URLSessionConfiguration.default
config.timeoutIntervalForRequest = 30
config.timeoutIntervalForResource = 60
urlSession = URLSession(configuration: config, delegate: self, delegateQueue: nil)
}
/// 初始化WebSocket管理器
func initialize(serverUrl: String) {
self.serverUrl = serverUrl
os_log("WebSocket管理器初始化,服务器地址: %@", log: log, type: .info, serverUrl)
}
/// 连接到服务器
func connect(completion: @escaping (Bool) -> Void) {
guard !isConnected && !isConnecting else {
completion(false)
return
}
guard let url = URL(string: serverUrl) else {
os_log("无效的服务器地址: %@", log: log, type: .error, serverUrl)
onError?("无效的服务器地址")
completion(false)
return
}
isConnecting = true
onConnectionStatusChanged?(.connecting)
// 创建WebSocket任务
var request = URLRequest(url: url)
request.timeoutInterval = 30
webSocketTask = urlSession?.webSocketTask(with: request)
// 开始监听消息
startListening()
// 开始连接
webSocketTask?.resume()
// 等待连接建立
DispatchQueue.global().asyncAfter(deadline: .now() + 1.0) { [weak self] in
self?.checkConnection(completion: completion)
}
}
/// 检查连接状态
private func checkConnection(completion: @escaping (Bool) -> Void) {
guard let webSocketTask = webSocketTask else {
isConnecting = false
onConnectionStatusChanged?(.error)
completion(false)
return
}
switch webSocketTask.state {
case .running:
isConnected = true
isConnecting = false
onConnectionStatusChanged?(.connected)
os_log("WebSocket连接成功", log: log, type: .info)
completion(true)
case .suspended:
// 等待连接完成
DispatchQueue.global().asyncAfter(deadline: .now() + 0.5) { [weak self] in
self?.checkConnection(completion: completion)
}
default:
isConnecting = false
onConnectionStatusChanged?(.error)
os_log("WebSocket连接失败", log: log, type: .error)
completion(false)
}
}
/// 开始监听消息
private func startListening() {
receiveMessage()
}
/// 接收消息
private func receiveMessage() {
webSocketTask?.receive { [weak self] result in
switch result {
case .success(let message):
self?.handleReceivedMessage(message)
// 继续监听下一条消息
self?.receiveMessage()
case .failure(let error):
self?.handleReceiveError(error)
}
}
}
/// 处理接收到的消息
private func handleReceivedMessage(_ message: URLSessionWebSocketTask.Message) {
switch message {
case .data(let data):
// 保留对二进制数据的支持以防万一
onAudioReceived?(data)
case .string(let text):
// 添加收到文本消息的日志输出,显示前20字符
let textPreview = text.count > 20 ? String(text.prefix(20)) + "..." : text
os_log("收到文本消息前20字符: %@", log: log, type: .debug, textPreview)
// 尝试解析为JSON来检查消息类型
do {
if let jsonData = text.data(using: .utf8),
let json = try JSONSerialization.jsonObject(with: jsonData) as? [String: Any],
let messageType = json["type"] as? String {
// 添加消息类型的日志输出
os_log("消息类型: %@", log: log, type: .debug, messageType)
switch messageType {
case "websocket_audio":
// 处理音频消息
if let audioDataString = json["data"] as? String,
let audioData = Data(base64Encoded: audioDataString) {
onAudioReceived?(audioData)
}
case "websocket_transcript":
// 处理统一的文本消息格式
// 包含sender、is_final、timestamp等字段
onTextReceived?(text)
case "websocket_ready":
// 连接就绪消息
os_log("WebSocket连接就绪", log: log, type: .info)
case "websocket_stop":
// 会话停止消息
os_log("会话已结束", log: log, type: .info)
case "websocket_start":
// 会话开始消息
os_log("会话已开始", log: log, type: .info)
default:
// 其他消息类型,当作普通文本消息处理
onTextReceived?(text)
}
} else {
// 如果不是JSON格式,当作普通文本消息处理
os_log("非JSON格式消息,当作普通文本处理", log: log, type: .debug)
onTextReceived?(text)
}
} catch {
// 如果解析失败,当作普通文本消息处理
os_log("非JSON格式消息,当作普通文本处理", log: log, type: .debug)
onTextReceived?(text)
}
@unknown default:
os_log("收到未知类型消息", log: log, type: .error)
}
}
/// 处理接收错误
private func handleReceiveError(_ error: Error) {
os_log("WebSocket接收错误: %@", log: log, type: .error, error.localizedDescription)
// 如果连接断开,尝试重连
if isConnected {
isConnected = false
onConnectionStatusChanged?(.error)
onError?("连接断开: \(error.localizedDescription)")
}
}
/// 发送音频数据
func sendAudioData(_ data: Data) {
guard isConnected else { return }
sendQueue.async { [weak self] in
// 将音频数据编码为Base64并包装成Vocode WebSocket格式
let base64AudioData = data.base64EncodedString()
let audioMessage: [String: Any] = [
"type": "websocket_audio",
"data": base64AudioData
]
do {
let jsonData = try JSONSerialization.data(withJSONObject: audioMessage)
if let jsonString = String(data: jsonData, encoding: .utf8) {
self?.webSocketTask?.send(.string(jsonString)) { error in
if let error = error {
self?.os_log("发送音频数据失败: %@", log: self?.log ?? OSLog.default, type: .error, error.localizedDescription)
}
}
}
} catch {
self?.os_log("音频数据JSON序列化失败: %@", log: self?.log ?? OSLog.default, type: .error, error.localizedDescription)
}
}
}
/// 发送文本消息
func sendTextMessage(_ message: String) -> Bool {
guard isConnected else {
os_log("WebSocket未连接,无法发送文本消息", log: log, type: .error)
return false
}
sendQueue.async { [weak self] in
self?.webSocketTask?.send(.string(message)) { error in
if let error = error {
self?.os_log("发送文本消息失败: %@", log: self?.log ?? OSLog.default, type: .error, error.localizedDescription)
}
}
}
return true
}
/// 断开连接
func disconnect() {
guard isConnected || isConnecting else { return }
webSocketTask?.cancel(with: .goingAway, reason: nil)
webSocketTask = nil
isConnected = false
isConnecting = false
onConnectionStatusChanged?(.disconnected)
os_log("WebSocket已断开连接", log: log, type: .info)
}
/// 释放资源
func dispose() {
disconnect()
urlSession?.invalidateAndCancel()
urlSession = nil
os_log("WebSocket管理器已释放", log: log, type: .info)
}
// 修复方法访问性问题
private func os_log(_ message: String, log: OSLog, type: OSLogType, _ args: CVarArg...) {
if args.isEmpty {
os_log("%@", log: log, type: type, message)
} else {
os_log(message, log: log, type: type, args)
}
}
}
// MARK: - URLSessionWebSocketDelegate
extension RealtimeWebSocketManager: URLSessionWebSocketDelegate {
func urlSession(_ session: URLSession, webSocketTask: URLSessionWebSocketTask, didOpenWithProtocol protocol: String?) {
os_log("WebSocket连接已建立", log: log, type: .info)
isConnected = true
isConnecting = false
onConnectionStatusChanged?(.connected)
}
func urlSession(_ session: URLSession, webSocketTask: URLSessionWebSocketTask, didCloseWith closeCode: URLSessionWebSocketTask.CloseCode, reason: Data?) {
os_log("WebSocket连接已关闭,关闭码: %d", log: log, type: .info, closeCode.rawValue)
isConnected = false
isConnecting = false
onConnectionStatusChanged?(.disconnected)
}
}
// MARK: - URLSessionDelegate
extension RealtimeWebSocketManager: URLSessionDelegate {
func urlSession(_ session: URLSession, didBecomeInvalidWithError error: Error?) {
if let error = error {
os_log("URLSession失效: %@", log: log, type: .error, error.localizedDescription)
onError?("会话失效: \(error.localizedDescription)")
}
}
func urlSession(_ session: URLSession, task: URLSessionTask, didCompleteWithError error: Error?) {
if let error = error {
os_log("URLSessionTask完成时出错: %@", log: log, type: .error, error.localizedDescription)
if isConnected || isConnecting {
isConnected = false
isConnecting = false
onConnectionStatusChanged?(.error)
onError?("连接错误: \(error.localizedDescription)")
}
}
}
}

243
local_plugins/realtime/lib/realtime.dart

@ -0,0 +1,243 @@
import 'dart:async';
import 'dart:typed_data';
import 'package:flutter/services.dart';
/// 实时语音服务异常
class RealtimeException implements Exception {
final String message;
RealtimeException(this.message);
@override
String toString() => 'RealtimeException: $message';
}
/// 连接状态
enum ConnectionStatus {
disconnected,
connecting,
connected,
error,
}
/// 语音状态
enum VoiceStatus {
idle,
recording,
processing,
playing,
}
/// 实时语音事件
class RealtimeEvent {
final String type;
final dynamic data;
RealtimeEvent({required this.type, this.data});
factory RealtimeEvent.fromMap(Map<String, dynamic> map) {
return RealtimeEvent(
type: map['type'] as String,
data: map['data'],
);
}
}
/// 实时语音聊天服务
class RealtimeService {
static const MethodChannel _channel = MethodChannel('realtime/methods');
static const EventChannel _eventChannel = EventChannel('realtime/events');
StreamController<RealtimeEvent>? _eventStreamController;
Stream<RealtimeEvent>? _eventStream;
/// 获取事件流
Stream<RealtimeEvent> get eventStream {
if (_eventStream == null) {
_eventStreamController = StreamController<RealtimeEvent>.broadcast();
_eventStream = _eventStreamController!.stream;
// 监听原生事件
_eventChannel.receiveBroadcastStream().listen(
(dynamic event) {
if (event is Map<dynamic, dynamic>) {
final eventMap = Map<String, dynamic>.from(event);
final realtimeEvent = RealtimeEvent.fromMap(eventMap);
_eventStreamController!.add(realtimeEvent);
}
},
onError: (error) {
_eventStreamController!.addError(RealtimeException('事件流错误: $error'));
},
);
}
return _eventStream!;
}
/// 初始化实时语音服务
///
/// [serverUrl] WebSocket服务器地址
/// [sampleRate] 采样率,默认16000
/// [channels] 声道数,默认1(单声道)
/// [bitsPerSample] 位深,默认16
Future<bool> initialize({
required String serverUrl,
int sampleRate = 16000,
int channels = 1,
int bitsPerSample = 16,
}) async {
try {
final result = await _channel.invokeMethod<bool>(
'initialize',
{
'serverUrl': serverUrl,
'sampleRate': sampleRate,
'channels': channels,
'bitsPerSample': bitsPerSample,
},
);
return result ?? false;
} catch (e) {
throw RealtimeException('初始化失败: $e');
}
}
/// 连接到服务器
Future<bool> connect() async {
try {
final result = await _channel.invokeMethod<bool>('connect');
return result ?? false;
} catch (e) {
throw RealtimeException('连接失败: $e');
}
}
/// 断开连接
Future<bool> disconnect() async {
try {
final result = await _channel.invokeMethod<bool>('disconnect');
return result ?? false;
} catch (e) {
throw RealtimeException('断开连接失败: $e');
}
}
/// 开始录音
Future<bool> startRecording() async {
try {
final result = await _channel.invokeMethod<bool>('startRecording');
return result ?? false;
} catch (e) {
throw RealtimeException('开始录音失败: $e');
}
}
/// 停止录音
Future<bool> stopRecording() async {
try {
final result = await _channel.invokeMethod<bool>('stopRecording');
return result ?? false;
} catch (e) {
throw RealtimeException('停止录音失败: $e');
}
}
/// 停止播放
Future<bool> stopPlaying() async {
try {
final result = await _channel.invokeMethod<bool>('stopPlaying');
return result ?? false;
} catch (e) {
throw RealtimeException('停止播放失败: $e');
}
}
/// 获取连接状态
Future<ConnectionStatus> getConnectionStatus() async {
try {
final result = await _channel.invokeMethod<String>('getConnectionStatus');
switch (result) {
case 'disconnected':
return ConnectionStatus.disconnected;
case 'connecting':
return ConnectionStatus.connecting;
case 'connected':
return ConnectionStatus.connected;
case 'error':
return ConnectionStatus.error;
default:
return ConnectionStatus.disconnected;
}
} catch (e) {
throw RealtimeException('获取连接状态失败: $e');
}
}
/// 获取语音状态
Future<VoiceStatus> getVoiceStatus() async {
try {
final result = await _channel.invokeMethod<String>('getVoiceStatus');
switch (result) {
case 'idle':
return VoiceStatus.idle;
case 'recording':
return VoiceStatus.recording;
case 'processing':
return VoiceStatus.processing;
case 'playing':
return VoiceStatus.playing;
default:
return VoiceStatus.idle;
}
} catch (e) {
throw RealtimeException('获取语音状态失败: $e');
}
}
/// 发送文本消息
Future<bool> sendTextMessage(String message) async {
try {
final result = await _channel.invokeMethod<bool>(
'sendTextMessage',
{'message': message},
);
return result ?? false;
} catch (e) {
throw RealtimeException('发送文本消息失败: $e');
}
}
/// 设置音频参数
Future<bool> setAudioConfig({
int? sampleRate,
int? channels,
int? bitsPerSample,
}) async {
try {
final params = <String, dynamic>{};
if (sampleRate != null) params['sampleRate'] = sampleRate;
if (channels != null) params['channels'] = channels;
if (bitsPerSample != null) params['bitsPerSample'] = bitsPerSample;
final result = await _channel.invokeMethod<bool>('setAudioConfig', params);
return result ?? false;
} catch (e) {
throw RealtimeException('设置音频参数失败: $e');
}
}
/// 释放资源
Future<void> dispose() async {
try {
await _channel.invokeMethod<void>('dispose');
_eventStreamController?.close();
_eventStreamController = null;
_eventStream = null;
} catch (e) {
throw RealtimeException('释放资源失败: $e');
}
}
}

29
local_plugins/realtime/pubspec.yaml

@ -0,0 +1,29 @@
name: realtime
description: 实时语音聊天插件,通过WebSocket连接Vocode服务器实现语音交互
version: 0.0.1
homepage:
environment:
sdk: ">=2.17.0 <3.0.0"
flutter: ">=2.5.0"
dependencies:
flutter:
sdk: flutter
dev_dependencies:
flutter_test:
sdk: flutter
flutter_lints: ^2.0.0
# Flutter插件配置
flutter:
plugin:
platforms:
android:
package: com.yunqiinnovation.realtime
pluginClass: RealtimePlugin
ios:
pluginClass: RealtimePlugin
swiftPackage:
path: ios/realtime
Loading…
Cancel
Save