Browse Source

Merge branch 'new_dev' of https://github.com/deepcloud2048/deep_voice into new_dev

weicu
tanlongsheng 1 year ago
parent
commit
533c74e13c
  1. 54
      lib/modules/agent/controllers/agent_controller.dart
  2. 508
      lib/modules/agent/views/agent_view.dart
  3. 190
      lib/modules/agent/widgets/full_screen_editor.dart
  4. 6
      lib/modules/profile/views/profile_view.dart
  5. 2
      lib/modules/settings/views/settings_view.dart
  6. 33
      local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt
  7. BIN
      local_plugins/agent_service/android/src/main/res/raw/await.mp3
  8. 1
      local_plugins/agent_service/ios/agent_service/Package.swift
  9. 176
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt
  10. 63
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt

54
lib/modules/agent/controllers/agent_controller.dart

@ -68,6 +68,13 @@ class AgentController extends GetxController {
final isImageInputActive = false.obs; // 是否处于图片输入模式
final selectedImagePath = Rx<String?>(null); // 当前选择的图片路径
// 输入状态控制 - 用于动态显示发送按钮
final hasTextInput = false.obs; // 是否有文本输入
// 全屏编辑相关状态
final isFullScreenEditMode = false.obs; // 是否处于全屏编辑模式
final shouldShowExpandButton = false.obs; // 是否显示展开按钮
// 流事件订阅
StreamSubscription<AgentServiceEvent>? _eventSubscription;
@ -107,6 +114,7 @@ class AgentController extends GetxController {
_subscribeToEvents(); //订阅代理服务事件
_initMusicStatus(); //初始化音乐状态
_setupMediaSessionCallbacks(); //设置媒体会话回调
_setupTextInputListener(); //设置文本输入监听器
}
// 初始化音乐播放状态
@ -178,6 +186,52 @@ class AgentController extends GetxController {
}
}
// 设置文本输入监听器
void _setupTextInputListener() {
textController.addListener(() {
final text = textController.text.trim();
hasTextInput.value = text.isNotEmpty;
// 检查是否需要显示展开按钮
_checkShouldShowExpandButton();
});
}
// 检查是否需要显示展开按钮
void _checkShouldShowExpandButton() {
final text = textController.text;
if (text.isEmpty) {
shouldShowExpandButton.value = false;
return;
}
// 计算文本行数
final lines = text.split('\n');
final maxLines = isImageInputActive.value ? 3 : 5;
// 简单估算:如果换行符数量超过maxLines-1,或者单行文本过长,显示展开按钮
final hasMultipleLines = lines.length > maxLines;
final hasLongLine = lines.any((line) => line.length > 50); // 估算单行字符数
shouldShowExpandButton.value = hasMultipleLines || hasLongLine;
}
// 打开全屏编辑模式
void openFullScreenEdit() {
isFullScreenEditMode.value = true;
}
// 关闭全屏编辑模式
void closeFullScreenEdit() {
isFullScreenEditMode.value = false;
}
// 从全屏编辑模式发送消息
Future<void> sendMessageFromFullScreen() async {
await sendMessage();
closeFullScreenEdit();
}
// 加载聊天历史记录
Future<void> _loadChatHistory() async {
try {

508
lib/modules/agent/views/agent_view.dart

@ -4,6 +4,7 @@ import 'package:get/get.dart';
import '../controllers/agent_controller.dart';
import '../controllers/report_controller.dart';
import 'message_bubble.dart';
import '../widgets/full_screen_editor.dart';
import 'package:flutter/services.dart';
import 'package:flutter_screenutil/flutter_screenutil.dart';
import '../../../core/theme/app_colors.dart';
@ -287,135 +288,52 @@ class AgentView extends StatelessWidget {
}
}),
// 输入控件行
Row(
children: [
// 键盘/语音切换按钮
Obx(() {
final isTextMode = controller.isTextInputMode.value;
return GestureDetector(
onTap: () {
controller.toggleInputMode();
},
child: Container(
width: 44,
height: 44,
decoration: BoxDecoration(
color: isDarkMode
? Colors.grey[800]
: Colors.grey[200],
borderRadius: BorderRadius.circular(22),
),
child: Icon(
isTextMode ? Icons.mic_none : Icons.keyboard,
color: isDarkMode
? Colors.grey[300]
: Colors.grey[700],
size: 22,
),
),
);
}),
const SizedBox(width: 8),
// 文本输入框或语音波形图
Expanded(
child: Obx(() {
// 输入控件行 - 使用IntrinsicHeight确保按钮与输入框底部对齐
IntrinsicHeight(
child: Row(
crossAxisAlignment: CrossAxisAlignment.end, // 所有元素底部对齐
children: [
// 键盘/语音切换按钮
Obx(() {
final isTextMode = controller.isTextInputMode.value;
if (isTextMode) {
// 文本输入模式 - 注意观察图片模式时的高度变化
return Container(
//height: 44,
padding: const EdgeInsets.symmetric(horizontal: 16),
return GestureDetector(
onTap: () {
controller.toggleInputMode();
},
child: Container(
width: 44,
height: 44,
decoration: BoxDecoration(
color: isDarkMode
? Colors.grey[800]
: Colors.grey[200],
borderRadius: BorderRadius.circular(22),
border: Border.all(
color: isDarkMode
? Colors.grey[
800]! // 修改: 从 Colors.grey[700] 改为 Colors.grey[800]
: Colors.grey.withOpacity(0.1),
width: 1.0,
),
),
child: Row(
children: [
Expanded(
child: Obx(() {
// 图片模式下使用更高的输入框
final isImageActive =
controller.isImageInputActive.value;
return TextField(
controller: controller.textController,
maxLines: isImageActive ? 2 : null,
minLines: isImageActive ? 2 : 1,
keyboardType:
TextInputType.multiline, // 启用多行键盘
textInputAction:
TextInputAction.newline, // 支持回车换行
style: TextStyle(
fontSize: 15,
color: isDarkMode
? Colors.white
: Colors.grey[800],
),
decoration: InputDecoration(
hintText: isImageActive
? '添加对图片的问题或描述...'
: '输入消息...',
hintStyle: TextStyle(
color: Colors.grey[500],
fontSize: 15,
),
filled: false,
border: InputBorder.none,
focusedBorder: OutlineInputBorder(
borderSide: BorderSide(
color: isDarkMode
? Colors.grey[800]!
: Colors.grey[200]!,
width: 1.0,
),
),
contentPadding: EdgeInsets.symmetric(
vertical: isImageActive ? 12 : 0),
),
onSubmitted: (_) =>
controller.sendMessage(),
);
}),
),
GestureDetector(
onTap: () => controller.sendMessage(),
child: Container(
padding: const EdgeInsets.all(6),
child: Icon(
Icons.send,
color: primaryColor,
size: 22,
),
),
),
],
child: Icon(
isTextMode ? Icons.mic_none : Icons.keyboard,
color: isDarkMode
? Colors.grey[300]
: Colors.grey[700],
size: 22,
),
);
} else {
// 语音输入模式 - 始终显示波形图
return GestureDetector(
onTap: () {
if (controller.isListening.value) {
controller.stopVoiceInput();
} else {
controller.startVoiceInput();
}
},
child: Container(
height: 50,
),
);
}),
const SizedBox(width: 8),
// 文本输入框或语音波形图
Expanded(
child: Obx(() {
final isTextMode = controller.isTextInputMode.value;
if (isTextMode) {
// 文本输入模式 - 限制最大高度并只向上扩展
return Container(
constraints: const BoxConstraints(
minHeight: 44, // 最小高度
maxHeight: 120, // 最大高度限制,防止挤压按钮
),
padding:
const EdgeInsets.symmetric(horizontal: 16),
decoration: BoxDecoration(
@ -425,118 +343,278 @@ class AgentView extends StatelessWidget {
borderRadius: BorderRadius.circular(22),
border: Border.all(
color: isDarkMode
? Colors.grey[800]!
? Colors.grey[
800]! // 修改: 从 Colors.grey[700] 改为 Colors.grey[800]
: Colors.grey.withOpacity(0.1),
width: 1.0,
),
),
child: Row(
crossAxisAlignment:
CrossAxisAlignment.end, // 对齐到底部,实现向上扩展
children: [
// 语音波形图 - 始终显示
Expanded(
child: Row(
mainAxisAlignment:
MainAxisAlignment.start,
children: [
...List.generate(10, (index) {
return _buildSoundBar(index);
}),
const Spacer(),
// "正在聆听"文字
Obx(() {
final isListening =
controller.isListening.value;
return Row(
children: [
Icon(
Icons.mic,
size: 14,
color: isListening
? primaryColor
: Colors.grey[500],
),
const SizedBox(width: 4),
Text(
'正在聆听',
style: TextStyle(
fontSize: 13,
fontWeight: FontWeight.w500,
child: Obx(() {
// 图片模式下使用更高的输入框
final isImageActive =
controller.isImageInputActive.value;
return TextField(
controller: controller.textController,
maxLines:
isImageActive ? 3 : 5, // 设置合理的最大行数
minLines: isImageActive ? 2 : 1,
keyboardType:
TextInputType.multiline, // 启用多行键盘
textInputAction:
TextInputAction.newline, // 支持回车换行
style: TextStyle(
fontSize: 15,
color: isDarkMode
? Colors.white
: Colors.grey[800],
),
decoration: InputDecoration(
hintText: isImageActive
? '添加对图片的问题或描述...'
: '输入消息...',
hintStyle: TextStyle(
color: Colors.grey[500],
fontSize: 15,
),
filled: false,
border: InputBorder.none,
focusedBorder: OutlineInputBorder(
borderSide: BorderSide(
color: isDarkMode
? Colors.grey[800]!
: Colors.grey[200]!,
width: 1.0,
),
),
contentPadding: EdgeInsets.symmetric(
vertical: isImageActive
? 12
: 8), // 调整垂直内边距
),
onSubmitted: (_) =>
controller.sendMessage(),
);
}),
),
],
),
);
} else {
// 语音输入模式 - 始终显示波形图
return GestureDetector(
onTap: () {
if (controller.isListening.value) {
controller.stopVoiceInput();
} else {
controller.startVoiceInput();
}
},
child: Container(
height: 50,
padding:
const EdgeInsets.symmetric(horizontal: 16),
decoration: BoxDecoration(
color: isDarkMode
? Colors.grey[800]
: Colors.grey[200],
borderRadius: BorderRadius.circular(22),
border: Border.all(
color: isDarkMode
? Colors.grey[800]!
: Colors.grey.withOpacity(0.1),
width: 1.0,
),
),
child: Row(
children: [
// 语音波形图 - 始终显示
Expanded(
child: Row(
mainAxisAlignment:
MainAxisAlignment.start,
children: [
...List.generate(10, (index) {
return _buildSoundBar(index);
}),
const Spacer(),
// "正在聆听"文字
Obx(() {
final isListening =
controller.isListening.value;
return Row(
children: [
Icon(
Icons.mic,
size: 14,
color: isListening
? primaryColor
: Colors.grey[500],
),
),
],
);
}),
],
const SizedBox(width: 4),
Text(
'正在聆听',
style: TextStyle(
fontSize: 13,
fontWeight: FontWeight.w500,
color: isListening
? primaryColor
: Colors.grey[500],
),
),
],
);
}),
],
),
),
),
],
],
),
),
),
);
}
}),
),
);
}
}),
),
const SizedBox(width: 8),
const SizedBox(width: 8),
// 图片选择按钮
Container(
width: 44,
height: 44,
decoration: BoxDecoration(
color: isDarkMode ? Colors.grey[800] : Colors.grey[200],
borderRadius: BorderRadius.circular(22),
),
child: PopupMenuButton<String>(
icon: Icon(
Icons.add_photo_alternate_outlined,
color:
isDarkMode ? Colors.grey[300] : Colors.grey[700],
size: 22,
),
padding: EdgeInsets.zero,
shape: RoundedRectangleBorder(
borderRadius: BorderRadius.circular(12),
),
onSelected: (value) async {
if (value == 'camera') {
await controller.takePhoto();
} else if (value == 'gallery') {
await controller.pickImage();
}
},
itemBuilder: (context) => [
PopupMenuItem<String>(
value: 'camera',
child: Row(
children: [
Icon(Icons.camera_alt,
color: primaryColor, size: 20),
const SizedBox(width: 8),
const Text('拍照',
style: TextStyle(fontSize: 14)),
],
// 动态按钮区域:展开按钮 + 发送按钮 或 图片选择按钮
Obx(() {
final hasText = controller.hasTextInput.value;
final hasImage = controller.isImageInputActive.value;
final shouldShowSend = hasText || hasImage;
final shouldShowExpand =
controller.shouldShowExpandButton.value;
if (shouldShowSend) {
// 显示发送按钮(可能还有展开按钮)
return Column(
mainAxisSize: MainAxisSize.min,
children: [
// 展开按钮(在发送按钮上方)
if (shouldShowExpand && hasText)
Padding(
padding: const EdgeInsets.only(bottom: 40),
child: GestureDetector(
onTap: () {
// 打开底部弹出编辑模式
controller.openFullScreenEdit();
showModalBottomSheet(
context: context,
isScrollControlled: true,
backgroundColor: Colors.transparent,
builder: (context) =>
const FullScreenEditor(),
).then((_) {
// 弹出页关闭时的回调
controller.closeFullScreenEdit();
});
},
child: Container(
width: 32,
height: 32,
decoration: BoxDecoration(
color: isDarkMode
? Colors.grey[700]
: Colors.grey[300],
borderRadius: BorderRadius.circular(16),
),
child: Icon(
Icons.expand_less,
color: isDarkMode
? Colors.grey[300]
: Colors.grey[700],
size: 20,
),
),
),
),
// 发送按钮
GestureDetector(
onTap: () => controller.sendMessage(),
child: Container(
width: 44,
height: 44,
decoration: BoxDecoration(
color: primaryColor,
borderRadius: BorderRadius.circular(22),
),
child: const Icon(
Icons.send,
color: Colors.white,
size: 22,
),
),
),
],
);
} else {
// 显示图片选择按钮
return Container(
width: 44,
height: 44,
decoration: BoxDecoration(
color: isDarkMode
? Colors.grey[800]
: Colors.grey[200],
borderRadius: BorderRadius.circular(22),
),
),
PopupMenuItem<String>(
value: 'gallery',
child: Row(
children: [
Icon(Icons.photo_library,
color: primaryColor, size: 20),
const SizedBox(width: 8),
const Text('从相册选择',
style: TextStyle(fontSize: 14)),
child: PopupMenuButton<String>(
icon: Icon(
Icons.add_photo_alternate_outlined,
color: isDarkMode
? Colors.grey[300]
: Colors.grey[700],
size: 22,
),
padding: EdgeInsets.zero,
shape: RoundedRectangleBorder(
borderRadius: BorderRadius.circular(12),
),
onSelected: (value) async {
if (value == 'camera') {
await controller.takePhoto();
} else if (value == 'gallery') {
await controller.pickImage();
}
},
itemBuilder: (context) => [
PopupMenuItem<String>(
value: 'camera',
child: Row(
children: [
Icon(Icons.camera_alt,
color: primaryColor, size: 20),
const SizedBox(width: 8),
const Text('拍照',
style: TextStyle(fontSize: 14)),
],
),
),
PopupMenuItem<String>(
value: 'gallery',
child: Row(
children: [
Icon(Icons.photo_library,
color: primaryColor, size: 20),
const SizedBox(width: 8),
const Text('从相册选择',
style: TextStyle(fontSize: 14)),
],
),
),
],
),
),
],
),
),
],
);
}
}),
],
),
),
],
),

190
lib/modules/agent/widgets/full_screen_editor.dart

@ -0,0 +1,190 @@
import 'package:flutter/material.dart';
import 'package:flutter_screenutil/flutter_screenutil.dart';
import 'package:get/get.dart';
import '../controllers/agent_controller.dart';
/// 底部弹出式文本编辑器组件
class FullScreenEditor extends StatefulWidget {
const FullScreenEditor({super.key});
@override
State<FullScreenEditor> createState() => _FullScreenEditorState();
}
class _FullScreenEditorState extends State<FullScreenEditor> {
late TextEditingController _textController;
late FocusNode _focusNode;
final AgentController controller = AgentController.to;
@override
void initState() {
super.initState();
// 使用与主界面相同的文本控制器,保持内容同步
_textController = controller.textController;
_focusNode = FocusNode();
// 延迟聚焦,确保界面完全加载后再聚焦
WidgetsBinding.instance.addPostFrameCallback((_) {
_focusNode.requestFocus();
});
}
@override
void dispose() {
_focusNode.dispose();
super.dispose();
}
@override
Widget build(BuildContext context) {
final isDarkMode = Theme.of(context).brightness == Brightness.dark;
final primaryColor = Theme.of(context).primaryColor;
final keyboardHeight = MediaQuery.of(context).viewInsets.bottom; // 获取键盘高度
return Container(
height: MediaQuery.of(context).size.height * 0.95, // 占屏幕高度的90%
decoration: BoxDecoration(
color: isDarkMode ? Colors.grey[900] : Colors.white,
borderRadius: BorderRadius.only(
topLeft: Radius.circular(20.r),
topRight: Radius.circular(20.r),
),
),
child: Stack(
children: [
// 主要内容区域
Column(
children: [
// 主要编辑区域
Expanded(
child: Container(
padding: EdgeInsets.fromLTRB(16.w, 60.w, 16.w, 16.w),
child: TextField(
controller: _textController,
focusNode: _focusNode,
maxLines: null,
expands: true,
keyboardType: TextInputType.multiline,
textInputAction: TextInputAction.newline,
style: TextStyle(
fontSize: 16.sp,
color: isDarkMode ? Colors.white : Colors.black87,
height: 1.5,
),
decoration: InputDecoration(
hintText: '在这里输入你的消息...',
hintStyle: TextStyle(
color: isDarkMode ? Colors.white54 : Colors.grey[500],
fontSize: 16.sp,
),
border: InputBorder.none,
contentPadding: EdgeInsets.zero,
),
),
),
),
// 底部操作区域
Container(
margin:
EdgeInsets.only(bottom: keyboardHeight), // 关键修改:添加键盘高度的边距
padding: EdgeInsets.all(16.w),
decoration: BoxDecoration(
color: isDarkMode ? Colors.grey[800] : Colors.grey[50],
border: Border(
top: BorderSide(
color: isDarkMode ? Colors.grey[700]! : Colors.grey[200]!,
width: 1,
),
),
),
child: Row(
children: [
// 字符计数
Expanded(
child: ValueListenableBuilder<TextEditingValue>(
valueListenable: _textController,
builder: (context, value, child) {
return Text(
'${value.text.length} 字符',
style: TextStyle(
color: isDarkMode
? Colors.white60
: Colors.grey[600],
fontSize: 14.sp,
),
);
},
),
),
// 发送按钮 - 保持与主界面一致的样式
Obx(() {
final hasText = controller.hasTextInput.value;
return GestureDetector(
onTap: hasText
? () async {
await controller.sendMessageFromFullScreen();
if (context.mounted) {
Navigator.of(context).pop();
}
}
: null,
child: Container(
width: 44.w,
height: 44.h,
decoration: BoxDecoration(
color: hasText
? primaryColor
: (isDarkMode
? Colors.grey[800]
: Colors.grey[300]),
borderRadius: BorderRadius.circular(22.r),
),
child: Icon(
Icons.send,
color: hasText
? Colors.white
: (isDarkMode
? Colors.grey[600]
: Colors.grey[500]),
size: 22.sp,
),
),
);
}),
],
),
),
],
),
// 右上角关闭按钮
Positioned(
top: 16.h,
right: 16.w,
child: GestureDetector(
onTap: () {
controller.closeFullScreenEdit();
Navigator.of(context).pop();
},
child: Container(
width: 32.w,
height: 32.h,
decoration: BoxDecoration(
color: isDarkMode ? Colors.grey[800] : Colors.grey[200],
borderRadius: BorderRadius.circular(16.r),
),
child: Icon(
Icons.close,
color: isDarkMode ? Colors.white : Colors.black87,
size: 20.sp,
),
),
),
),
],
),
);
}
}

6
lib/modules/profile/views/profile_view.dart

@ -327,7 +327,7 @@ class EditProfileView extends GetView<ProfileController> {
Row(
children: [
Obx(() => CircleAvatar(
radius: 30.r, // 调整头像大小以适应行高
radius: 26.r, // 调整头像大小以适应行高
backgroundImage: controller.localAvatar.value.isNotEmpty
? FileImage(File(controller.localAvatar.value))
: (controller.userInfo.value.avatar.isNotEmpty
@ -444,7 +444,7 @@ class EditProfileView extends GetView<ProfileController> {
borderRadius: BorderRadius.circular(16.r),
),
child: Icon(
Icons.edit,
Icons.chevron_right,
size: 16.w,
color: isDarkMode ? Colors.white70 : Colors.grey[500],
),
@ -504,7 +504,7 @@ class EditProfileView extends GetView<ProfileController> {
child: TextField(
controller: nameController,
autofocus: true,
textAlign: TextAlign.right, // 右对齐,与原昵称显示保持一致
textAlign: TextAlign.center, // 右对齐,与原昵称显示保持一致
style: TextStyle(
fontSize: 14.sp,
color: isDarkMode ? Colors.white70 : Colors.grey[600],

2
lib/modules/settings/views/settings_view.dart

@ -977,7 +977,7 @@ class SettingsView extends GetView<SettingsController> {
//------------测试使用标记 start---------------//
SizedBox(height: 4.h),
Text(
'${Platform.operatingSystem}-test-${DateTime.now().month.toString().padLeft(2, '0')}${DateTime.now().day.toString().padLeft(2, '0')}',
'${Platform.operatingSystem}-test-0626',
style: TextStyle(
fontSize: 12.sp,
color: const Color.fromARGB(255, 143, 94, 94),

33
local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt

@ -668,7 +668,7 @@ object AgentService : CoroutineScope {
try {
// 设置状态为正在流式输出
_isAiStreaming.set(true)
audioPlayer?.playAudio(R.raw.await, true,0.3f)
// 使用历史记录作为上下文发送到OpenAI
val responseBuilder = StringBuilder()
var aiMetadata:String = ""
@ -736,6 +736,9 @@ object AgentService : CoroutineScope {
responseBuilder.append(token)
if (speakResponse && broadcast) {
ttsService?.speakStream(token)
if (token.length > 0){
audioPlayer?.stopAudio()
}
}
if (broadcast){
// 发送流式回复token
@ -783,7 +786,7 @@ object AgentService : CoroutineScope {
override fun onFunctionCall(call: JSONObject) {
try {
audioPlayer?.playAudio(R.raw.calling, true)
// audioPlayer?.playAudio(R.raw.calling, true)
val name = call.getString("name")
sendEvent("function_call", mapOf(
"name" to name,
@ -1120,34 +1123,32 @@ object AgentService : CoroutineScope {
/**
* 播放音频资源
* @param resId 资源ID
* @param isLooping 是否循环播放
* @param volume 音量大小,范围0.0-1.0,默认1.0
*/
fun playAudio(resId: Int, isLooping: Boolean = false) {
fun playAudio(resId: Int, isLooping: Boolean = false, volume: Float = 1.0f) {
try {
// 释放之前的资源
release()
// 创建播放器并设置资源
mediaPlayer = MediaPlayer().apply {
// 设置资源
context.resources.openRawResourceFd(resId)?.use { fd ->
setDataSource(fd.fileDescriptor, fd.startOffset, fd.length)
}
// 播放完成后自动释放资源
this.isLooping = isLooping // 设置循环属性
setVolume(volume, volume) // 设置音量(左声道,右声道)
setOnCompletionListener {
release()
}
// 准备并播放
prepare()
if (isLooping) {
start()
} else {
start()
setOnCompletionListener {
if (!isLooping) {
release()
}
}
prepare()
start()
}
} catch (e: Exception) {
Log.e(TAG, "播放音频资源异常: ${e.message}", e)

BIN
local_plugins/agent_service/android/src/main/res/raw/await.mp3

Binary file not shown.

1
local_plugins/agent_service/ios/agent_service/Package.swift

@ -31,6 +31,7 @@ let package = Package(
path: "Sources/agent_service",
resources: [
.copy("Resources/calling.mp3"),
.copy("Resources/await.mp3"),
.copy("Resources/start.mp3"),
.copy("Resources/stop.mp3")
]

176
local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt

@ -25,6 +25,7 @@ class CustomSseClientTransport(
private val urlString: String?,
private val reconnectionTime: Duration? = null,
private val requestBuilder: HttpRequestBuilder.() -> Unit = {},
private val onConnectionLost: (() -> Unit)? = null
) : AbstractTransport() {
private val TAG = "CustomSseClientTransport"
@ -35,8 +36,10 @@ class CustomSseClientTransport(
private val initialized = AtomicBoolean(false)
private var session: ClientSSESession by Delegates.notNull()
private val endpoint = CompletableDeferred<String>()
private val isConnected = AtomicBoolean(false)
private var job: Job? = null
private var connectionMonitorJob: Job? = null
// 创建JSON解析器
private val json = Json {
@ -95,74 +98,118 @@ class CustomSseClientTransport(
*/
private suspend fun collectEvents() {
job = scope.launch(CoroutineName("CustomSseMcpClientTransport.collect#${hashCode()}")) {
session.incoming.collect { event ->
when (event.event) {
"error" -> {
val e = IllegalStateException("SSE error: ${event.data}")
Log.e(TAG, "SSE错误: ${event.data}")
_onError(e)
throw e
}
try {
session.incoming.collect { event ->
when (event.event) {
"error" -> {
Log.e(TAG, "SSE错误: ${event.data}")
isConnected.set(false)
val exception = Exception("SSE Error: ${event.data}")
_onError(exception)
onConnectionLost?.invoke()
throw exception
}
"open" -> {
// SSE连接已打开
}
"ping" -> {
// 心跳
}
"endpoint" -> {
try {
val eventData = event.data ?: ""
"open" -> {
// SSE连接已打开
Log.d(TAG, "SSE连接已打开")
isConnected.set(true)
}
"ping" -> {
// 心跳
}
"endpoint" -> {
try {
val eventData = event.data ?: ""
// 构建完整的端点URL
val fullEndpoint = if (eventData.contains(hostPart)) {
eventData
} else if (eventData.startsWith("/")) {
"$hostPart$eventData"
} else {
eventData
}
// 添加查询参数
val endpointWithParams = if (queryParams.isNotEmpty()) {
if (fullEndpoint.contains("?")) {
val queryString = queryParams.entries.joinToString("&") { "${it.key}=${it.value}" }
"$fullEndpoint&$queryString"
// 构建完整的端点URL
val fullEndpoint = if (eventData.contains(hostPart)) {
eventData
} else if (eventData.startsWith("/")) {
"$hostPart$eventData"
} else {
val queryString = queryParams.entries.joinToString("&") { "${it.key}=${it.value}" }
"$fullEndpoint?$queryString"
eventData
}
} else {
fullEndpoint
// 添加查询参数
val endpointWithParams = if (queryParams.isNotEmpty()) {
if (fullEndpoint.contains("?")) {
val queryString = queryParams.entries.joinToString("&") { "${it.key}=${it.value}" }
"$fullEndpoint&$queryString"
} else {
val queryString = queryParams.entries.joinToString("&") { "${it.key}=${it.value}" }
"$fullEndpoint?$queryString"
}
} else {
fullEndpoint
}
endpoint.complete(endpointWithParams)
} catch (e: Exception) {
Log.e(TAG, "处理endpoint事件失败: ${e.message}", e)
_onError(e)
close()
error(e)
}
endpoint.complete(endpointWithParams)
} catch (e: Exception) {
Log.e(TAG, "处理endpoint事件失败: ${e.message}", e)
_onError(e)
close()
error(e)
}
}
else -> {
try {
val data = event.data
if (data != null) {
try {
val message = json.decodeFromString<JSONRPCMessage>(data)
_onMessage(message)
} catch (e: Exception) {
Log.e(TAG, "解析JSON-RPC消息失败: ${e.message}", e)
_onError(e)
else -> {
try {
val data = event.data
if (data != null) {
try {
val message = json.decodeFromString<JSONRPCMessage>(data)
_onMessage(message)
} catch (e: Exception) {
Log.e(TAG, "解析JSON-RPC消息失败: ${e.message}", e)
_onError(e)
}
}
} catch (e: Exception) {
Log.e(TAG, "处理事件失败: ${e.message}", e)
_onError(e)
}
} catch (e: Exception) {
Log.e(TAG, "处理事件失败: ${e.message}", e)
_onError(e)
}
}
}
} catch (e: CancellationException) {
Log.d(TAG, "SSE事件收集被取消")
throw e
} catch (e: Exception) {
Log.e(TAG, "SSE连接异常断开: ${e.message}", e)
isConnected.set(false)
_onError(e)
onConnectionLost?.invoke()
throw e
}
}
// 启动连接监控
startConnectionMonitor()
}
/**
* 启动连接监控,定期检查连接状态
*/
private fun startConnectionMonitor() {
connectionMonitorJob = scope.launch {
while (isActive && isConnected.get()) {
try {
delay(10000) // 每10秒检查一次
// 检查session是否仍然活跃
if (session.coroutineContext[Job]?.isCancelled == true) {
Log.w(TAG, "检测到SSE会话已取消")
isConnected.set(false)
onConnectionLost?.invoke()
break
}
} catch (e: Exception) {
Log.e(TAG, "连接监控异常: ${e.message}", e)
isConnected.set(false)
onConnectionLost?.invoke()
break
}
}
}
}
@ -212,6 +259,8 @@ class CustomSseClientTransport(
// 等待endpoint就绪
endpoint.await()
Log.d(TAG, "CustomSseClientTransport启动完成")
}
/**
@ -245,6 +294,13 @@ class CustomSseClientTransport(
}
}
/**
* 检查连接状态
*/
fun isConnectionActive(): Boolean {
return isConnected.get() && session.coroutineContext[Job]?.isActive == true
}
/**
* 关闭传输层
*/
@ -253,8 +309,14 @@ class CustomSseClientTransport(
Log.e(TAG, "关闭失败: 传输层未初始化")
error("CustomSseClientTransport is not initialized!")
}
isConnected.set(false)
connectionMonitorJob?.cancel()
session.cancel()
_onClose()
job?.cancelAndJoin()
connectionMonitorJob?.cancelAndJoin()
Log.d(TAG, "CustomSseClientTransport已关闭")
}
}
}

63
local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt

@ -44,13 +44,14 @@ class MCPSubClient(
private var mcpClient: Client? = null
private var isConnected = false
private var availableTools = mutableListOf<Tool>()
private var transport: CustomSseClientTransport? = null
/**
* 连接到MCP服务器
*/
suspend fun connect(): Boolean = connectionMutex.withLock {
if (isConnected) return true
Log.e(TAG, "[$serverId] 开始连接mcp服务器: $serverUrl")
return try {
// 创建MCP客户端实例
val client = Client(
@ -61,13 +62,20 @@ class MCPSubClient(
)
// 根据URL类型选择传输方式
val transport = when {
val newTransport = when {
serverUrl.startsWith("http://") || serverUrl.startsWith("https://") -> {
// SSE传输 - 使用自定义的CustomSseClientTransport
val mcpHttpClient = httpClient ?: createMcpHttpClient()
CustomSseClientTransport(
client = mcpHttpClient,
urlString = serverUrl
urlString = serverUrl,
onConnectionLost = {
// 连接断开回调
Log.w(TAG, "[$serverId] 检测到连接断开")
scope.launch {
handleConnectionLost()
}
}
)
}
else -> {
@ -76,8 +84,10 @@ class MCPSubClient(
}
}
transport = newTransport
// 连接到服务器
client.connect(transport)
client.connect(newTransport)
// 获取可用工具列表
try {
@ -101,9 +111,9 @@ class MCPSubClient(
retryCount = 0
currentReconnectDelay = initialReconnectDelay
// 启动心跳检测
startHeartbeat()
// startHeartbeat()
Log.e(TAG, "[$serverId] 连接mcp服务器成功: $serverUrl")
true
} catch (e: Exception) {
Log.e(TAG, "[$serverId] MCP连接失败: ${e.message}", e)
false
@ -321,7 +331,15 @@ class MCPSubClient(
*/
suspend fun checkConnection(): Boolean {
if (!isConnected) {
Log.d(TAG, "当前未连接,尝试重新连接...")
Log.d(TAG, "[$serverId] 当前未连接,尝试重新连接...")
return connect()
}
// 检查传输层连接状态
val transportActive = transport?.isConnectionActive() ?: false
if (!transportActive) {
Log.w(TAG, "[$serverId] 传输层连接已断开")
isConnected = false
return connect()
}
@ -330,11 +348,34 @@ class MCPSubClient(
mcpClient?.ping()
return true
} catch (e: Exception) {
Log.e(TAG, "连接检查失败: ${e.message}")
Log.e(TAG, "[$serverId] 连接检查失败: ${e.message}")
isConnected = false
return false
}
}
/**
* 处理连接断开事件
*/
private suspend fun handleConnectionLost() {
connectionMutex.withLock {
if (isConnected) {
Log.w(TAG, "[$serverId] 连接已断开,更新状态")
isConnected = false
stopHeartbeat()
// 可以在这里添加自动重连逻辑
// 或者通知上层应用连接已断开
}
}
}
/**
* 获取连接状态
*/
fun getConnectionStatus(): Boolean {
return isConnected && (transport?.isConnectionActive() ?: false)
}
/**
* 停止心跳检测
*/
@ -393,10 +434,14 @@ class MCPSubClient(
scope.launch {
connectionMutex.withLock {
try {
stopHeartbeat()
mcpClient?.close()
transport?.close()
mcpClient = null
transport = null
isConnected = false
availableTools.clear()
Log.d(TAG, "[$serverId] MCP连接已关闭")
} catch (e: Exception) {
Log.e(TAG, "[$serverId] 关闭MCP连接时出错: ${e.message}", e)
}
@ -404,4 +449,4 @@ class MCPSubClient(
}
scope.cancel()
}
}
}
Loading…
Cancel
Save