diff --git a/local_plugins/realtime/.gitignore b/local_plugins/realtime/.gitignore
new file mode 100644
index 000000000..f5af98475
--- /dev/null
+++ b/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
\ No newline at end of file
diff --git a/local_plugins/realtime/IMPLEMENTATION.md b/local_plugins/realtime/IMPLEMENTATION.md
new file mode 100644
index 000000000..a5d49a976
--- /dev/null
+++ b/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
+NSMicrophoneUsageDescription
+需要麦克风权限用于语音识别和录音功能
+```
+
+## 与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示例的功能。
\ No newline at end of file
diff --git a/local_plugins/realtime/README.md b/local_plugins/realtime/README.md
new file mode 100644
index 000000000..a5bc746f0
--- /dev/null
+++ b/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
+NSMicrophoneUsageDescription
+应用需要麦克风权限进行语音录制
+```
+
+## 注意事项
+
+1. 确保服务器地址正确且可访问
+2. 网络环境良好,避免频繁断线
+3. 音频格式与服务器保持一致
+4. 及时释放资源,避免内存泄漏
+
+## 错误处理
+
+插件会自动处理常见错误:
+- 网络连接失败
+- 音频设备不可用
+- 权限被拒绝
+- 服务器断开连接
+
+通过事件流可以监听这些错误并进行相应处理。
\ No newline at end of file
diff --git a/local_plugins/realtime/android/build.gradle.kts b/local_plugins/realtime/android/build.gradle.kts
new file mode 100644
index 000000000..b787e4124
--- /dev/null
+++ b/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")
+}
\ No newline at end of file
diff --git a/local_plugins/realtime/android/settings.gradle.kts b/local_plugins/realtime/android/settings.gradle.kts
new file mode 100644
index 000000000..dd4a80ad2
--- /dev/null
+++ b/local_plugins/realtime/android/settings.gradle.kts
@@ -0,0 +1 @@
+rootProject.name = "realtime"
\ No newline at end of file
diff --git a/local_plugins/realtime/android/src/main/AndroidManifest.xml b/local_plugins/realtime/android/src/main/AndroidManifest.xml
new file mode 100644
index 000000000..64df6a0ae
--- /dev/null
+++ b/local_plugins/realtime/android/src/main/AndroidManifest.xml
@@ -0,0 +1,11 @@
+
+
+
+
+
+
+
+
+
+
\ No newline at end of file
diff --git a/local_plugins/realtime/android/src/main/kotlin/com/yunqiinnovation/realtime/RealtimeAudioManager.kt b/local_plugins/realtime/android/src/main/kotlin/com/yunqiinnovation/realtime/RealtimeAudioManager.kt
new file mode 100644
index 000000000..711a85779
--- /dev/null
+++ b/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()
+ 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()
+ }
+}
\ No newline at end of file
diff --git a/local_plugins/realtime/android/src/main/kotlin/com/yunqiinnovation/realtime/RealtimePlugin.kt b/local_plugins/realtime/android/src/main/kotlin/com/yunqiinnovation/realtime/RealtimePlugin.kt
new file mode 100644
index 000000000..b8077c2cf
--- /dev/null
+++ b/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
+ 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
+ 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
+ ?: 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 {
+ val pcm16List = mutableListOf()
+
+ // 确保数据长度是偶数(每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帧
+}
\ No newline at end of file
diff --git a/local_plugins/realtime/android/src/main/kotlin/com/yunqiinnovation/realtime/RealtimeWebSocketManager.kt b/local_plugins/realtime/android/src/main/kotlin/com/yunqiinnovation/realtime/RealtimeWebSocketManager.kt
new file mode 100644
index 000000000..2726ceef7
--- /dev/null
+++ b/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)
+ }
+}
\ No newline at end of file
diff --git a/local_plugins/realtime/example.md b/local_plugins/realtime/example.md
new file mode 100644
index 000000000..74b8ebd44
--- /dev/null
+++ b/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 {
+ final RealtimeService _realtimeService = RealtimeService();
+ StreamSubscription? _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 _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 _connect() async {
+ try {
+ final success = await _realtimeService.connect();
+ if (!success) {
+ print('连接失败');
+ }
+ } catch (e) {
+ print('连接错误: $e');
+ }
+ }
+
+ Future _disconnect() async {
+ try {
+ await _realtimeService.disconnect();
+ } catch (e) {
+ print('断开连接错误: $e');
+ }
+ }
+
+ Future _startRecording() async {
+ try {
+ final success = await _realtimeService.startRecording();
+ if (!success) {
+ print('开始录音失败');
+ }
+ } catch (e) {
+ print('录音错误: $e');
+ }
+ }
+
+ Future _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? _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 _initializeService() async {
+ await _realtimeService.initialize(
+ serverUrl: 'wss://your-server/ws',
+ );
+ }
+
+ Future toggleConnection() async {
+ if (isConnected.value) {
+ await _realtimeService.disconnect();
+ } else {
+ await _realtimeService.connect();
+ }
+ }
+
+ Future 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
+NSMicrophoneUsageDescription
+应用需要麦克风权限进行语音录制
+```
+
+## 服务器配置示例
+
+需要一个支持WebSocket的Vocode服务器,具体实现可参考Vocode官方文档。
+
+服务器需要:
+1. 接收16kHz/16-bit/单声道的PCM音频数据
+2. 返回相同格式的音频数据
+3. 支持文本消息交换
+
+## 注意事项
+
+1. 确保网络连接稳定
+2. 音频格式必须与服务器保持一致
+3. 及时处理事件流中的错误
+4. 在适当时机释放资源
\ No newline at end of file
diff --git a/local_plugins/realtime/ios/realtime/Package.swift b/local_plugins/realtime/ios/realtime/Package.swift
new file mode 100644
index 000000000..514fbc1e1
--- /dev/null
+++ b/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"
+ )
+ ]
+)
\ No newline at end of file
diff --git a/local_plugins/realtime/ios/realtime/Sources/realtime/RealtimeAudioManager.swift b/local_plugins/realtime/ios/realtime/Sources/realtime/RealtimeAudioManager.swift
new file mode 100644
index 000000000..c2a9d424a
--- /dev/null
+++ b/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
+ }
+}
\ No newline at end of file
diff --git a/local_plugins/realtime/ios/realtime/Sources/realtime/RealtimePlugin.swift b/local_plugins/realtime/ios/realtime/Sources/realtime/RealtimePlugin.swift
new file mode 100644
index 000000000..23f6720b0
--- /dev/null
+++ b/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
+ }
+}
\ No newline at end of file
diff --git a/local_plugins/realtime/ios/realtime/Sources/realtime/RealtimeWebSocketManager.swift b/local_plugins/realtime/ios/realtime/Sources/realtime/RealtimeWebSocketManager.swift
new file mode 100644
index 000000000..039101982
--- /dev/null
+++ b/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)")
+ }
+ }
+ }
+}
\ No newline at end of file
diff --git a/local_plugins/realtime/lib/realtime.dart b/local_plugins/realtime/lib/realtime.dart
new file mode 100644
index 000000000..237e9e3e8
--- /dev/null
+++ b/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 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? _eventStreamController;
+ Stream? _eventStream;
+
+ /// 获取事件流
+ Stream get eventStream {
+ if (_eventStream == null) {
+ _eventStreamController = StreamController.broadcast();
+ _eventStream = _eventStreamController!.stream;
+
+ // 监听原生事件
+ _eventChannel.receiveBroadcastStream().listen(
+ (dynamic event) {
+ if (event is Map) {
+ final eventMap = Map.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 initialize({
+ required String serverUrl,
+ int sampleRate = 16000,
+ int channels = 1,
+ int bitsPerSample = 16,
+ }) async {
+ try {
+ final result = await _channel.invokeMethod(
+ 'initialize',
+ {
+ 'serverUrl': serverUrl,
+ 'sampleRate': sampleRate,
+ 'channels': channels,
+ 'bitsPerSample': bitsPerSample,
+ },
+ );
+
+ return result ?? false;
+ } catch (e) {
+ throw RealtimeException('初始化失败: $e');
+ }
+ }
+
+ /// 连接到服务器
+ Future connect() async {
+ try {
+ final result = await _channel.invokeMethod('connect');
+ return result ?? false;
+ } catch (e) {
+ throw RealtimeException('连接失败: $e');
+ }
+ }
+
+ /// 断开连接
+ Future disconnect() async {
+ try {
+ final result = await _channel.invokeMethod('disconnect');
+ return result ?? false;
+ } catch (e) {
+ throw RealtimeException('断开连接失败: $e');
+ }
+ }
+
+ /// 开始录音
+ Future startRecording() async {
+ try {
+ final result = await _channel.invokeMethod('startRecording');
+ return result ?? false;
+ } catch (e) {
+ throw RealtimeException('开始录音失败: $e');
+ }
+ }
+
+ /// 停止录音
+ Future stopRecording() async {
+ try {
+ final result = await _channel.invokeMethod('stopRecording');
+ return result ?? false;
+ } catch (e) {
+ throw RealtimeException('停止录音失败: $e');
+ }
+ }
+
+ /// 停止播放
+ Future stopPlaying() async {
+ try {
+ final result = await _channel.invokeMethod('stopPlaying');
+ return result ?? false;
+ } catch (e) {
+ throw RealtimeException('停止播放失败: $e');
+ }
+ }
+
+ /// 获取连接状态
+ Future getConnectionStatus() async {
+ try {
+ final result = await _channel.invokeMethod('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 getVoiceStatus() async {
+ try {
+ final result = await _channel.invokeMethod('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 sendTextMessage(String message) async {
+ try {
+ final result = await _channel.invokeMethod(
+ 'sendTextMessage',
+ {'message': message},
+ );
+ return result ?? false;
+ } catch (e) {
+ throw RealtimeException('发送文本消息失败: $e');
+ }
+ }
+
+ /// 设置音频参数
+ Future setAudioConfig({
+ int? sampleRate,
+ int? channels,
+ int? bitsPerSample,
+ }) async {
+ try {
+ final params = {};
+ if (sampleRate != null) params['sampleRate'] = sampleRate;
+ if (channels != null) params['channels'] = channels;
+ if (bitsPerSample != null) params['bitsPerSample'] = bitsPerSample;
+
+ final result = await _channel.invokeMethod('setAudioConfig', params);
+ return result ?? false;
+ } catch (e) {
+ throw RealtimeException('设置音频参数失败: $e');
+ }
+ }
+
+ /// 释放资源
+ Future dispose() async {
+ try {
+ await _channel.invokeMethod('dispose');
+ _eventStreamController?.close();
+ _eventStreamController = null;
+ _eventStream = null;
+ } catch (e) {
+ throw RealtimeException('释放资源失败: $e');
+ }
+ }
+}
\ No newline at end of file
diff --git a/local_plugins/realtime/pubspec.yaml b/local_plugins/realtime/pubspec.yaml
new file mode 100644
index 000000000..0eed88a7d
--- /dev/null
+++ b/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
\ No newline at end of file