|
|
|
@ -1,6 +1,8 @@ |
|
|
|
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 |
|
|
|
@ -9,6 +11,7 @@ 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 |
|
|
|
|
|
|
|
/** |
|
|
|
* 实时语音聊天插件 |
|
|
|
@ -28,12 +31,15 @@ class RealtimePlugin: FlutterPlugin, MethodCallHandler, EventChannel.StreamHandl |
|
|
|
private lateinit var eventChannel: EventChannel |
|
|
|
private var eventSink: EventChannel.EventSink? = null |
|
|
|
|
|
|
|
// 核心组件 |
|
|
|
private lateinit var audioManager: RealtimeAudioManager |
|
|
|
private lateinit var webSocketManager: RealtimeWebSocketManager |
|
|
|
// 核心组件:设为可空,因为它们的生命周期与页面会话绑定 |
|
|
|
private var audioManager: RealtimeAudioManager? = null |
|
|
|
private var webSocketManager: RealtimeWebSocketManager? = null |
|
|
|
|
|
|
|
// 协程作用域 |
|
|
|
private val scope = CoroutineScope(Dispatchers.Main + SupervisorJob()) |
|
|
|
// 协程作用域:需要可变,因为在重用插件时需要重新创建 |
|
|
|
private var scope = CoroutineScope(Dispatchers.Default + SupervisorJob()) |
|
|
|
|
|
|
|
// 主线程Handler,用于向Flutter发送事件与回调 |
|
|
|
private val mainHandler = Handler(Looper.getMainLooper()) |
|
|
|
|
|
|
|
// 配置参数 |
|
|
|
private var serverUrl: String = "" |
|
|
|
@ -47,6 +53,11 @@ class RealtimePlugin: FlutterPlugin, MethodCallHandler, EventChannel.StreamHandl |
|
|
|
@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 |
|
|
|
|
|
|
|
@ -55,68 +66,96 @@ class RealtimePlugin: FlutterPlugin, MethodCallHandler, EventChannel.StreamHandl |
|
|
|
|
|
|
|
eventChannel = EventChannel(flutterPluginBinding.binaryMessenger, EVENT_CHANNEL) |
|
|
|
eventChannel.setStreamHandler(this) |
|
|
|
|
|
|
|
// 初始化核心组件 |
|
|
|
audioManager = RealtimeAudioManager(context) |
|
|
|
webSocketManager = RealtimeWebSocketManager() |
|
|
|
|
|
|
|
setupCallbacks() |
|
|
|
|
|
|
|
Log.i(TAG, "Realtime插件已附加到引擎") |
|
|
|
} |
|
|
|
|
|
|
|
override fun onDetachedFromEngine(binding: FlutterPlugin.FlutterPluginBinding) { |
|
|
|
Log.i(TAG, "插件正在从引擎分离,执行清理...") |
|
|
|
performCleanup() |
|
|
|
methodChannel.setMethodCallHandler(null) |
|
|
|
eventChannel.setStreamHandler(null) |
|
|
|
|
|
|
|
// 清理资源 |
|
|
|
scope.cancel() |
|
|
|
audioManager.dispose() |
|
|
|
webSocketManager.dispose() |
|
|
|
|
|
|
|
Log.i(TAG, "Realtime插件已从引擎分离") |
|
|
|
} |
|
|
|
|
|
|
|
/// 设置各组件的回调 |
|
|
|
private fun setupCallbacks() { |
|
|
|
Log.d(TAG, "setupCallbacks: 正在为新的管理器实例设置回调...") |
|
|
|
|
|
|
|
// 音频管理器回调 |
|
|
|
audioManager.onAudioData = { data -> |
|
|
|
webSocketManager.sendAudioData(data) |
|
|
|
audioManager?.onAudioData = { data -> |
|
|
|
// 发送到WebSocket |
|
|
|
webSocketManager?.sendAudioData(data) |
|
|
|
// 发送录音PCM数据到Flutter用于可视化 |
|
|
|
val pcm16Data = convertByteArrayToPCM16(data) |
|
|
|
sendEvent("recordingPcmData", pcm16Data) |
|
|
|
} |
|
|
|
|
|
|
|
audioManager.onStatusChanged = { status -> |
|
|
|
audioManager?.onStatusChanged = { status -> |
|
|
|
voiceStatus = status |
|
|
|
sendEvent("voiceStatusChanged", status.value) |
|
|
|
} |
|
|
|
|
|
|
|
// 新增RMS回调 |
|
|
|
audioManager?.onRmsChanged = { rms -> |
|
|
|
sendEvent("audioRms", rms) |
|
|
|
} |
|
|
|
|
|
|
|
// WebSocket管理器回调 |
|
|
|
webSocketManager.onConnectionStatusChanged = { status -> |
|
|
|
webSocketManager?.onConnectionStatusChanged = { status -> |
|
|
|
connectionStatus = status |
|
|
|
sendEvent("connectionStatusChanged", status.value) |
|
|
|
} |
|
|
|
|
|
|
|
webSocketManager.onAudioReceived = { data -> |
|
|
|
audioManager.playAudioData(data) |
|
|
|
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 -> |
|
|
|
webSocketManager?.onTextReceived = { text -> |
|
|
|
// 直接发送原始JSON文本到Flutter层处理 |
|
|
|
sendEvent("textReceived", text) |
|
|
|
} |
|
|
|
|
|
|
|
webSocketManager.onError = { error -> |
|
|
|
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 |
|
|
|
) |
|
|
|
|
|
|
|
scope.launch { |
|
|
|
eventSink?.success(event) |
|
|
|
try { |
|
|
|
mainHandler.post { |
|
|
|
try { |
|
|
|
sink.success(event) |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.w(TAG, "发送事件失败: \\${e.message}") |
|
|
|
} |
|
|
|
} |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.w(TAG, "发送事件时出现异常: \\${e.message}") |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ -139,21 +178,46 @@ class RealtimePlugin: FlutterPlugin, MethodCallHandler, EventChannel.StreamHandl |
|
|
|
|
|
|
|
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 |
|
|
|
this.sampleRate = args["sampleRate"] as? Int ?: 16000 |
|
|
|
this.channels = args["channels"] as? Int ?: 1 |
|
|
|
this.bitsPerSample = args["bitsPerSample"] as? Int ?: 16 |
|
|
|
|
|
|
|
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) |
|
|
|
val success = audioManager?.initialize(audioConfig) ?: false |
|
|
|
|
|
|
|
if (success) { |
|
|
|
webSocketManager.initialize(serverUrl) |
|
|
|
webSocketManager?.initialize(serverUrl) |
|
|
|
Log.i(TAG, "实时语音插件初始化成功") |
|
|
|
} else { |
|
|
|
Log.e(TAG, "实时语音插件初始化失败") |
|
|
|
@ -167,13 +231,14 @@ class RealtimePlugin: FlutterPlugin, MethodCallHandler, EventChannel.StreamHandl |
|
|
|
} |
|
|
|
|
|
|
|
private fun handleConnect(result: Result) { |
|
|
|
Log.d(TAG, "handleConnect: 收到连接请求。scope是否活跃? ${scope.isActive}") |
|
|
|
scope.launch { |
|
|
|
try { |
|
|
|
val success = webSocketManager.connect() |
|
|
|
result.success(success) |
|
|
|
val success = webSocketManager?.connect() ?: false |
|
|
|
mainHandler.post { result.success(success) } |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "连接失败: ${e.message}", e) |
|
|
|
result.error("CONNECTION_ERROR", "连接失败: ${e.message}", null) |
|
|
|
mainHandler.post { result.error("CONNECTION_ERROR", "连接失败: ${e.message}", null) } |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
@ -181,11 +246,11 @@ class RealtimePlugin: FlutterPlugin, MethodCallHandler, EventChannel.StreamHandl |
|
|
|
private fun handleDisconnect(result: Result) { |
|
|
|
scope.launch { |
|
|
|
try { |
|
|
|
webSocketManager.disconnect() |
|
|
|
result.success(true) |
|
|
|
webSocketManager?.disconnect() |
|
|
|
mainHandler.post { result.success(true) } |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "断开连接失败: ${e.message}", e) |
|
|
|
result.error("DISCONNECTION_ERROR", "断开连接失败: ${e.message}", null) |
|
|
|
mainHandler.post { result.error("DISCONNECTION_ERROR", "断开连接失败: ${e.message}", null) } |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
@ -193,11 +258,11 @@ class RealtimePlugin: FlutterPlugin, MethodCallHandler, EventChannel.StreamHandl |
|
|
|
private fun handleStartRecording(result: Result) { |
|
|
|
scope.launch { |
|
|
|
try { |
|
|
|
val success = audioManager.startRecording() |
|
|
|
result.success(success) |
|
|
|
val success = audioManager?.startRecording() ?: false |
|
|
|
mainHandler.post { result.success(success) } |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "开始录音失败: ${e.message}", e) |
|
|
|
result.error("RECORDING_ERROR", "开始录音失败: ${e.message}", null) |
|
|
|
mainHandler.post { result.error("RECORDING_ERROR", "开始录音失败: ${e.message}", null) } |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
@ -205,11 +270,11 @@ class RealtimePlugin: FlutterPlugin, MethodCallHandler, EventChannel.StreamHandl |
|
|
|
private fun handleStopRecording(result: Result) { |
|
|
|
scope.launch { |
|
|
|
try { |
|
|
|
val success = audioManager.stopRecording() |
|
|
|
result.success(success) |
|
|
|
val success = audioManager?.stopRecording() ?: false |
|
|
|
mainHandler.post { result.success(success) } |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "停止录音失败: ${e.message}", e) |
|
|
|
result.error("RECORDING_ERROR", "停止录音失败: ${e.message}", null) |
|
|
|
mainHandler.post { result.error("RECORDING_ERROR", "停止录音失败: ${e.message}", null) } |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
@ -217,11 +282,11 @@ class RealtimePlugin: FlutterPlugin, MethodCallHandler, EventChannel.StreamHandl |
|
|
|
private fun handleStopPlaying(result: Result) { |
|
|
|
scope.launch { |
|
|
|
try { |
|
|
|
val success = audioManager.stopPlaying() |
|
|
|
result.success(success) |
|
|
|
val success = audioManager?.stopPlaying() ?: false |
|
|
|
mainHandler.post { result.success(success) } |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "停止播放失败: ${e.message}", e) |
|
|
|
result.error("PLAYBACK_ERROR", "停止播放失败: ${e.message}", null) |
|
|
|
mainHandler.post { result.error("PLAYBACK_ERROR", "停止播放失败: ${e.message}", null) } |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
@ -232,7 +297,7 @@ class RealtimePlugin: FlutterPlugin, MethodCallHandler, EventChannel.StreamHandl |
|
|
|
val message = args?.get("message") as? String |
|
|
|
?: return result.error("INVALID_ARGUMENTS", "消息内容不能为空", null) |
|
|
|
|
|
|
|
val success = webSocketManager.sendTextMessage(message) |
|
|
|
val success = webSocketManager?.sendTextMessage(message) ?: false |
|
|
|
result.success(success) |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "发送文本消息失败: ${e.message}", e) |
|
|
|
@ -245,12 +310,23 @@ class RealtimePlugin: FlutterPlugin, MethodCallHandler, EventChannel.StreamHandl |
|
|
|
val args = call.arguments as? Map<String, Any> |
|
|
|
?: return result.error("INVALID_ARGUMENTS", "参数无效", null) |
|
|
|
|
|
|
|
args["sampleRate"]?.let { this.sampleRate = it as Int } |
|
|
|
args["channels"]?.let { this.channels = it as Int } |
|
|
|
args["bitsPerSample"]?.let { this.bitsPerSample = it as Int } |
|
|
|
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) |
|
|
|
val success = audioManager?.updateConfig(audioConfig) ?: false |
|
|
|
result.success(success) |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "设置音频配置失败: ${e.message}", e) |
|
|
|
@ -259,14 +335,84 @@ class RealtimePlugin: FlutterPlugin, MethodCallHandler, EventChannel.StreamHandl |
|
|
|
} |
|
|
|
|
|
|
|
private fun handleDispose(result: Result) { |
|
|
|
scope.launch { |
|
|
|
try { |
|
|
|
audioManager.dispose() |
|
|
|
webSocketManager.dispose() |
|
|
|
result.success(null) |
|
|
|
} catch (e: Exception) { |
|
|
|
Log.e(TAG, "释放资源失败: ${e.message}", e) |
|
|
|
result.error("DISPOSE_ERROR", "释放资源失败: ${e.message}", null) |
|
|
|
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, "核心资源清理完成。") |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|