Compare commits

..
18 Commits
Author SHA1 Message Date
hrzandGitHub 14a63c4b97 update:补齐ddos攻击提示音 (#896) 2025-04-19 13:24:03 +08:00
hrzandGitHub 098d13a34b update:更新0.3.7 (#892) 2025-04-19 00:11:12 +08:00
488f247744 update:增加单台设备每天最多聊天字数,防止被ddos
* 添加清空redis所有库的接口
--AdminController.java 添加了清除所有的接口
--RedisUtils.java 添加了清除redis所有key的方法,redisTemplate提供的清空方法已经被标记为弃用了,所有选择用执行lua脚本方式

* fix:修复使用本地配置时忘记附带提示词

* fix:修复意图识别插件名称格式bug

* update:版本升级后强制刷新redis

* update:增加单台设备每天最多聊天字数,防止被ddos

---------

Co-authored-by: 剑雨 <2375294554@qq.com>
2025-04-18 23:39:56 +08:00
hrzandGitHub c74dbf03cd fix:修复意图识别插件名称格式bug (#885)
* fix:修复使用本地配置时忘记附带提示词

* fix:修复意图识别插件名称格式bug
2025-04-18 16:39:18 +08:00
hrzandGitHub 1f836c3235 fix:修复使用本地配置时忘记附带提示词 (#880) 2025-04-18 11:38:31 +08:00
hrzandGitHub bd2e2e77d5 update:更新版本号 (#874) 2025-04-18 00:51:09 +08:00
6282ef14e8 fix:修复第一句话是默认配置的bug (#873)
* update: 增加对httpClient的统一管理

* update: 增加retry机制

* update: 增加retry机制

* fix:websocket连接后不说话的备用关闭方法

* fix:修复第一句话是默认配置的bug

---------

Co-authored-by: haotian <haotian@codemao.cn>
Co-authored-by: GoodyHao <865700600@qq.com>
2025-04-18 00:37:40 +08:00
558f23688f fix:websocket连接后不说话的备用关闭方法 (#872)
* update: 增加对httpClient的统一管理

* update: 增加retry机制

* update: 增加retry机制

* fix:websocket连接后不说话的备用关闭方法

---------

Co-authored-by: haotian <haotian@codemao.cn>
Co-authored-by: GoodyHao <865700600@qq.com>
2025-04-17 23:57:04 +08:00
CGDandGitHub 541f2de599 Merge pull request #856 from xinnan-tech/web-page-check
解决切换页面样式变化的问题
2025-04-17 13:26:26 +08:00
hrzandGitHub 68116254cd dify、coze对话模式支持functioncall
Function call v2
2025-04-17 11:54:04 +08:00
hrzandGitHub 77ff4599ea Merge branch 'main' into function-call-v2 2025-04-17 11:53:25 +08:00
玄凤科技 652f5a3247 Merge branch 'function-call-v2' of https://github.com/xinnan-tech/xiaozhi-esp32-server into function-call-v2 2025-04-17 10:00:23 +08:00
玄凤科技 9b6e57b143 修改提示词,增强tool调用约束 2025-04-17 10:00:13 +08:00
Ran_Chen bdc19256bf 解决切换页面样式变化的问题 2025-04-17 09:34:42 +08:00
hrz 9fc164a69d udpate:增加iot消息properties和methods可能为空的情况 2025-04-16 23:30:46 +08:00
玄凤科技 208c045d3d 触发function call后,清空回复队列 2025-04-16 11:48:40 +08:00
玄凤科技 8c7d129089 同步coze llm 2025-04-16 08:58:55 +08:00
玄凤科技 a3a9b98a1d 增加系统提示词,支持dify使用function call 2025-04-16 08:56:23 +08:00
28 changed files with 660 additions and 203 deletions
+1
View File
@@ -158,6 +158,7 @@ my_wakeup_words.mp3
!main/xiaozhi-server/config/assets/bind_code.wav !main/xiaozhi-server/config/assets/bind_code.wav
!main/xiaozhi-server/config/assets/bind_not_found.wav !main/xiaozhi-server/config/assets/bind_not_found.wav
!main/xiaozhi-server/config/assets/bind_code/*.wav !main/xiaozhi-server/config/assets/bind_code/*.wav
!main/xiaozhi-server/config/assets/max_output_size.wav
main/manager-api/.vscode main/manager-api/.vscode
# Ignore webpack cache directory # Ignore webpack cache directory
@@ -163,4 +163,9 @@ public interface Constant {
return value; return value;
} }
} }
/**
* 版本号
*/
public static final String VERSION = "0.3.7";
} }
@@ -75,4 +75,11 @@ public class RedisKeys {
public static String getTimbreDetailsKey(String id) { public static String getTimbreDetailsKey(String id) {
return "timbre:details:" + id; return "timbre:details:" + id;
} }
/**
* 获取版本号Key
*/
public static String getVersionKey() {
return "system:version";
}
} }
@@ -1,11 +1,14 @@
package xiaozhi.common.redis; package xiaozhi.common.redis;
import java.util.Collection; import java.util.Collection;
import java.util.Collections;
import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import org.springframework.data.redis.core.HashOperations; import org.springframework.data.redis.core.HashOperations;
import org.springframework.data.redis.core.RedisTemplate; import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.data.redis.core.script.DefaultRedisScript;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import jakarta.annotation.Resource; import jakarta.annotation.Resource;
@@ -27,7 +30,7 @@ public class RedisUtils {
/** /**
* 过期时长为1小时,单位:秒 * 过期时长为1小时,单位:秒
*/ */
public final static long HOUR_ONE_EXPIRE = 60 * 60 * 1L; public final static long HOUR_ONE_EXPIRE = (long) 60 * 60;
/** /**
* 过期时长为6小时,单位:秒 * 过期时长为6小时,单位:秒
*/ */
@@ -124,4 +127,24 @@ public class RedisUtils {
public Object rightPop(String key) { public Object rightPop(String key) {
return redisTemplate.opsForList().rightPop(key); return redisTemplate.opsForList().rightPop(key);
} }
/**
* 清空所有 Redis 数据库中的所有键
*/
public void emptyAll() {
// Lua 脚本 FLUSHALL是redis清空所有库的命令
String luaScript ="redis.call('FLUSHALL')";
// 创建 DefaultRedisScript 对象
DefaultRedisScript<Void> redisScript = new DefaultRedisScript<>();
redisScript.setScriptText(luaScript); // 设置 Lua 脚本内容
redisScript.setResultType(Void.class); // 设置返回值类型
// 执行 Lua 脚本
List<String> keys = Collections.emptyList(); // 如果脚本不依赖 key,可以传入空列表
redisTemplate.execute(redisScript, keys);
}
} }
@@ -5,6 +5,9 @@ import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.DependsOn; import org.springframework.context.annotation.DependsOn;
import jakarta.annotation.PostConstruct; import jakarta.annotation.PostConstruct;
import xiaozhi.common.constant.Constant;
import xiaozhi.common.redis.RedisKeys;
import xiaozhi.common.redis.RedisUtils;
import xiaozhi.modules.config.service.ConfigService; import xiaozhi.modules.config.service.ConfigService;
import xiaozhi.modules.sys.service.SysParamsService; import xiaozhi.modules.sys.service.SysParamsService;
@@ -18,8 +21,20 @@ public class SystemInitConfig {
@Autowired @Autowired
private ConfigService configService; private ConfigService configService;
@Autowired
private RedisUtils redisUtils;
@PostConstruct @PostConstruct
public void init() { public void init() {
// 检查版本号
String redisVersion = (String) redisUtils.get(RedisKeys.getVersionKey());
if (!Constant.VERSION.equals(redisVersion)) {
// 如果版本不一致,清空Redis
redisUtils.emptyAll();
// 存储新版本号
redisUtils.set(RedisKeys.getVersionKey(), Constant.VERSION);
}
sysParamsService.initServerSecret(); sysParamsService.initServerSecret();
configService.getConfig(false); configService.getConfig(false);
} }
@@ -91,6 +91,7 @@ public class ConfigServiceImpl implements ConfigService {
} }
throw new RenException(ErrorCode.OTA_DEVICE_NOT_FOUND, "not found device"); throw new RenException(ErrorCode.OTA_DEVICE_NOT_FOUND, "not found device");
} }
// 获取智能体信息 // 获取智能体信息
AgentEntity agent = agentService.getAgentById(device.getAgentId()); AgentEntity agent = agentService.getAgentById(device.getAgentId());
if (agent == null) { if (agent == null) {
@@ -104,7 +105,9 @@ public class ConfigServiceImpl implements ConfigService {
} }
// 构建返回数据 // 构建返回数据
Map<String, Object> result = new HashMap<>(); Map<String, Object> result = new HashMap<>();
// 获取单台设备每天最多输出字数
String deviceMaxOutputSize = sysParamsService.getValue("device_max_output_size", true);
result.put("device_max_output_size", deviceMaxOutputSize);
// 如果客户端已实例化模型,则不返回 // 如果客户端已实例化模型,则不返回
String alreadySelectedVadModelId = (String) selectedModule.get("VAD"); String alreadySelectedVadModelId = (String) selectedModule.get("VAD");
if (alreadySelectedVadModelId != null && alreadySelectedVadModelId.equals(agent.getVadModelId())) { if (alreadySelectedVadModelId != null && alreadySelectedVadModelId.equals(agent.getVadModelId())) {
@@ -267,6 +270,12 @@ public class ConfigServiceImpl implements ConfigService {
if (intentLLMModelId != null && intentLLMModelId.equals(llmModelId)) { if (intentLLMModelId != null && intentLLMModelId.equals(llmModelId)) {
intentLLMModelId = null; intentLLMModelId = null;
} }
} else if ("function_call".equals(map.get("type"))) {
String functionStr = (String) map.get("functions");
if (StringUtils.isNotBlank(functionStr)) {
String[] functions = functionStr.split("\\;");
map.put("functions", functions);
}
} }
} }
// 如果是LLM类型,且intentLLMModelId不为空,则添加附加模型 // 如果是LLM类型,且intentLLMModelId不为空,则添加附加模型
@@ -15,6 +15,7 @@ import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.tags.Tag; import io.swagger.v3.oas.annotations.tags.Tag;
import jakarta.servlet.http.HttpServletResponse; import jakarta.servlet.http.HttpServletResponse;
import lombok.AllArgsConstructor; import lombok.AllArgsConstructor;
import xiaozhi.common.constant.Constant;
import xiaozhi.common.exception.ErrorCode; import xiaozhi.common.exception.ErrorCode;
import xiaozhi.common.exception.RenException; import xiaozhi.common.exception.RenException;
import xiaozhi.common.page.TokenDTO; import xiaozhi.common.page.TokenDTO;
@@ -121,7 +122,7 @@ public class LoginController {
@Operation(summary = "公共配置") @Operation(summary = "公共配置")
public Result<Map<String, Object>> pubConfig() { public Result<Map<String, Object>> pubConfig() {
Map<String, Object> config = new HashMap<>(); Map<String, Object> config = new HashMap<>();
config.put("version", "0.3.5"); config.put("version", Constant.VERSION);
config.put("allowUserRegister", sysUserService.getAllowUserRegister()); config.put("allowUserRegister", sysUserService.getAllowUserRegister());
return new Result<Map<String, Object>>().ok(config); return new Result<Map<String, Object>>().ok(config);
} }
@@ -0,0 +1,2 @@
delete from `ai_model_config` where id = 'Intent_function_call';
INSERT INTO `ai_model_config` VALUES ('Intent_function_call', 'Intent', 'function_call', '函数调用意图识别', 0, 1, '{\"type\": \"function_call\", \"functions\": \"change_role;get_weather;get_news;play_music\"}', NULL, NULL, 3, NULL, NULL, NULL, NULL);
@@ -0,0 +1,7 @@
-- 调整意图识别配置
delete from `ai_model_config` where id = 'Intent_function_call';
INSERT INTO `ai_model_config` VALUES ('Intent_function_call', 'Intent', 'function_call', '函数调用意图识别', 0, 1, '{\"type\": \"function_call\", \"functions\": \"change_role;get_weather;get_news;play_music\"}', NULL, NULL, 3, NULL, NULL, NULL, NULL);
-- 增加单台设备每天最多聊天句数
delete from `sys_params` where id = 105;
INSERT INTO `sys_params` (id, param_code, param_value, value_type, param_type, remark) VALUES (105, 'device_max_output_size', '0', 'number', 1, '单台设备每天最多输出字数,0表示不限制');
@@ -58,3 +58,10 @@ databaseChangeLog:
- sqlFile: - sqlFile:
encoding: utf8 encoding: utf8
path: classpath:db/changelog/202504151206.sql path: classpath:db/changelog/202504151206.sql
- changeSet:
id: 202504181536
author: John
changes:
- sqlFile:
encoding: utf8
path: classpath:db/changelog/202504181536.sql
+18 -13
View File
@@ -512,20 +512,25 @@ export default {
align-items: center; align-items: center;
} }
.save-btn { .action-bar {
background: #5778ff; .el-button.save-btn {
color: white; background: #5778ff;
border: none; color: white;
border-radius: 18px; border: none;
padding: 10px 20px; border-radius: 18px;
} padding: 10px 20px;
width: 100px;
height: 35px;
font-size: 14px;
}
.reset-btn { .el-button.reset-btn {
background: #e6ebff; background: #e6ebff;
color: #5778ff; color: #5778ff;
border: 1px solid #adbdff; border: 1px solid #adbdff;
border-radius: 18px; border-radius: 18px;
padding: 10px 20px; padding: 10px 20px;
}
} }
.hint-text { .hint-text {
+2 -2
View File
@@ -147,7 +147,7 @@ selected_module:
Memory: nomem Memory: nomem
# 意图识别模块开启后,可以播放音乐、控制音量、识别退出指令。 # 意图识别模块开启后,可以播放音乐、控制音量、识别退出指令。
# 不想开通意图识别,就设置成:nointent # 不想开通意图识别,就设置成:nointent
# 意图识别可使用intent_llm,如果你的LLM是DifyLLM或CozeLLM,建议使用这个。优点:通用性强,缺点:增加串行前置意图识别模块,会增加处理时间,这个意图识别暂时不支持控制音量大小等iot操作 # 意图识别可使用intent_llm。优点:通用性强,缺点:增加串行前置意图识别模块,会增加处理时间,这个意图识别暂时不支持控制音量大小等iot操作
# 意图识别可使用function_call,缺点:需要所选择的LLM支持function_call,优点:按需调用工具、速度快,理论上能全部操作所有iot指令 # 意图识别可使用function_call,缺点:需要所选择的LLM支持function_call,优点:按需调用工具、速度快,理论上能全部操作所有iot指令
# 默认免费的ChatGLMLLM就已经支持function_call,但是如果像追求稳定建议把LLM设置成:DoubaoLLM,使用的具体model_name是:doubao-pro-32k-functioncall-241028 # 默认免费的ChatGLMLLM就已经支持function_call,但是如果像追求稳定建议把LLM设置成:DoubaoLLM,使用的具体model_name是:doubao-pro-32k-functioncall-241028
Intent: function_call Intent: function_call
@@ -163,7 +163,7 @@ Intent:
type: intent_llm type: intent_llm
# 配备意图识别独立的思考模型 # 配备意图识别独立的思考模型
# 如果这里不填,则会默认使用selected_module.LLM的模型作为意图识别的思考模型 # 如果这里不填,则会默认使用selected_module.LLM的模型作为意图识别的思考模型
# 如果你的selected_module.LLM选择了DifyLLM或CozeLLM,这里最好使用独立的LLM作为意图识别,例如使用免费的ChatGLMLLM # 如果你的不想使用selected_module.LLM意图识别,这里最好使用独立的LLM作为意图识别,例如使用免费的ChatGLMLLM
llm: ChatGLMLLM llm: ChatGLMLLM
function_call: function_call:
# 不需要动type # 不需要动type
+11 -92
View File
@@ -1,18 +1,7 @@
import os import os
import argparse import argparse
import requests
import yaml import yaml
import time from config.manage_api_client import init_service, get_server_config, get_agent_models
class DeviceNotFoundException(Exception):
pass
class DeviceBindException(Exception):
def __init__(self, bind_code):
self.bind_code = bind_code
super().__init__(f"设备绑定异常,绑定码: {bind_code}")
# 添加全局配置缓存 # 添加全局配置缓存
@@ -65,97 +54,27 @@ def get_config_file():
return config_file return config_file
def _make_api_request(api_url, secret, endpoint, json_data=None):
"""执行API请求的通用函数
Args:
api_url: API的基础URL
secret: API密钥
endpoint: API端点
json_data: 请求的JSON数据
Returns:
dict: API返回的数据
Raises:
Exception: 当请求失败时抛出异常
"""
if not api_url or not secret:
raise Exception("manager-api的url或secret配置错误")
if "" in secret:
raise Exception("请先配置manager-api的secret")
max_retries = 10
retry_delay = 2 # 秒
for attempt in range(max_retries):
try:
response = requests.post(f"{api_url}{endpoint}", json=json_data)
if response.status_code == 200:
result = response.json()
if result.get("code") == 10041:
raise DeviceNotFoundException(result.get("msg"))
elif result.get("code") == 10042:
raise DeviceBindException(result.get("msg"))
elif result.get("code") != 0:
raise Exception(f"API返回错误: {result.get('msg', '未知错误')}")
return result.get("data")
error_msg = f"manager-api请求失败,状态码: {response.status_code}"
try:
error_data = response.json()
if "msg" in error_data:
error_msg = f"{error_msg}, 错误信息: {error_data['msg']}"
except:
error_msg = f"{error_msg}, 响应内容: {response.text}"
if attempt < max_retries - 1:
print(f"请求manager-api失败,正在重试 ({attempt + 1}/{max_retries})...")
time.sleep(retry_delay)
else:
raise Exception(error_msg)
except requests.exceptions.RequestException as e:
if attempt < max_retries - 1:
print(f"请求manager-api异常,正在重试 ({attempt + 1}/{max_retries})...")
time.sleep(retry_delay)
else:
raise Exception(f"manager-api请求异常: {str(e)}")
def get_config_from_api(config): def get_config_from_api(config):
"""从Java API获取配置""" """从Java API获取配置"""
api_url = config["manager-api"].get("url", "") # 初始化API客户端
secret = config["manager-api"].get("secret", "") init_service(config)
# 获取服务器配置
config_data = get_server_config()
if config_data is None:
raise Exception("Failed to fetch server config from API")
config_data = _make_api_request(
api_url, secret, "/config/server-base", {"secret": secret}
)
config_data["read_config_from_api"] = True config_data["read_config_from_api"] = True
config_data["manager-api"] = { config_data["manager-api"] = {
"url": api_url, "url": config["manager-api"].get("url", ""),
"secret": secret, "secret": config["manager-api"].get("secret", ""),
} }
return config_data return config_data
def get_private_config_from_api(config, device_id, client_id): def get_private_config_from_api(config, device_id, client_id):
"""从Java API获取私有配置""" """从Java API获取私有配置"""
api_url = config["manager-api"].get("url", "") return get_agent_models(device_id, client_id, config["selected_module"])
secret = config["manager-api"].get("secret", "")
return _make_api_request(
api_url,
secret,
"/config/agent-models",
{
"secret": secret,
"macAddress": device_id,
"clientId": client_id,
"selectedModule": config["selected_module"],
},
)
def ensure_directories(config): def ensure_directories(config):
+1 -1
View File
@@ -3,7 +3,7 @@ import sys
from loguru import logger from loguru import logger
from config.config_loader import load_config from config.config_loader import load_config
SERVER_VERSION = "0.3.5" SERVER_VERSION = "0.3.7"
def get_module_abbreviation(module_name, module_dict): def get_module_abbreviation(module_name, module_dict):
@@ -0,0 +1,154 @@
import os
import time
from typing import Optional, Dict
import httpx
TAG = __name__
class DeviceNotFoundException(Exception):
pass
class DeviceBindException(Exception):
def __init__(self, bind_code):
self.bind_code = bind_code
super().__init__(f"设备绑定异常,绑定码: {bind_code}")
class ManageApiClient:
_instance = None
_client = None
_secret = None
def __new__(cls, config):
"""单例模式确保全局唯一实例,并支持传入配置参数"""
if cls._instance is None:
cls._instance = super().__new__(cls)
cls._init_client(config)
return cls._instance
@classmethod
def _init_client(cls, config):
"""初始化持久化连接池"""
cls.config = config.get("manager-api")
if not cls.config:
raise Exception("manager-api配置错误")
if not cls.config.get("url") or not cls.config.get("secret"):
raise Exception("manager-api的url或secret配置错误")
if "" in cls.config.get("secret"):
raise Exception("请先配置manager-api的secret")
cls._secret = cls.config.get("secret")
cls.max_retries = cls.config.get("max_retries", 6) # 最大重试次数
cls.retry_delay = cls.config.get("retry_delay", 10) # 初始重试延迟(秒)
# NOTE(goody): 2025/4/16 http相关资源统一管理,后续可以增加线程池或者超时
# 后续也可以统一配置apiToken之类的走通用的Auth
cls._client = httpx.Client(
base_url=cls.config.get("url"),
headers={
"User-Agent": f"PythonClient/2.0 (PID:{os.getpid()})",
"Accept": "application/json",
},
timeout=cls.config.get("timeout", 30), # 默认超时时间30秒
)
@classmethod
def _request(cls, method: str, endpoint: str, **kwargs) -> Dict:
"""发送单次HTTP请求并处理响应"""
endpoint = endpoint.lstrip("/")
response = cls._client.request(method, endpoint, **kwargs)
response.raise_for_status()
result = response.json()
# 处理API返回的业务错误
if result.get("code") == 10041:
raise DeviceNotFoundException(result.get("msg"))
elif result.get("code") == 10042:
raise DeviceBindException(result.get("msg"))
elif result.get("code") != 0:
raise Exception(f"API返回错误: {result.get('msg', '未知错误')}")
# 返回成功数据
return result.get("data") if result.get("code") == 0 else None
@classmethod
def _should_retry(cls, exception: Exception) -> bool:
"""判断异常是否应该重试"""
# 网络连接相关错误
if isinstance(
exception, (httpx.ConnectError, httpx.TimeoutException, httpx.NetworkError)
):
return True
# HTTP状态码错误
if isinstance(exception, httpx.HTTPStatusError):
status_code = exception.response.status_code
return status_code in [408, 429, 500, 502, 503, 504]
return False
@classmethod
def _execute_request(cls, method: str, endpoint: str, **kwargs) -> Dict:
"""带重试机制的请求执行器"""
retry_count = 0
while retry_count <= cls.max_retries:
try:
# 执行请求
return cls._request(method, endpoint, **kwargs)
except Exception as e:
# 判断是否应该重试
if retry_count < cls.max_retries and cls._should_retry(e):
retry_count += 1
print(
f"{method} {endpoint} 请求失败,将在 {cls.retry_delay:.1f} 秒后进行第 {retry_count} 次重试"
)
time.sleep(cls.retry_delay)
continue
else:
# 不重试,直接抛出异常
raise
@classmethod
def safe_close(cls):
"""安全关闭连接池"""
if cls._client:
cls._client.close()
cls._instance = None
def get_server_config() -> Optional[Dict]:
"""获取服务器基础配置"""
return ManageApiClient._instance._execute_request(
"POST", "/config/server-base", json={"secret": ManageApiClient._secret}
)
def get_agent_models(
mac_address: str, client_id: str, selected_module: Dict
) -> Optional[Dict]:
"""获取代理模型配置"""
return ManageApiClient._instance._execute_request(
"POST",
"/config/agent-models",
json={
"secret": ManageApiClient._secret,
"macAddress": mac_address,
"clientId": client_id,
"selectedModule": selected_module,
},
)
def init_service(config):
ManageApiClient(config)
def manage_api_http_safe_close():
ManageApiClient.safe_close()
+95 -43
View File
@@ -27,11 +27,9 @@ from core.handle.functionHandler import FunctionHandler
from plugins_func.register import Action, ActionResponse from plugins_func.register import Action, ActionResponse
from core.auth import AuthMiddleware, AuthenticationError from core.auth import AuthMiddleware, AuthenticationError
from core.mcp.manager import MCPManager from core.mcp.manager import MCPManager
from config.config_loader import ( from config.config_loader import get_private_config_from_api
get_private_config_from_api, from config.manage_api_client import DeviceNotFoundException, DeviceBindException
DeviceNotFoundException, from core.utils.output_counter import add_device_output
DeviceBindException,
)
TAG = __name__ TAG = __name__
@@ -60,6 +58,7 @@ class ConnectionHandler:
self.session_id = None self.session_id = None
self.prompt = None self.prompt = None
self.welcome_msg = None self.welcome_msg = None
self.max_output_size = 0
# 客户端状态相关 # 客户端状态相关
self.client_abort = False self.client_abort = False
@@ -112,6 +111,11 @@ class ConnectionHandler:
self.close_after_chat = False # 是否在聊天结束后关闭连接 self.close_after_chat = False # 是否在聊天结束后关闭连接
self.use_function_call_mode = False self.use_function_call_mode = False
self.timeout_task = None
self.timeout_seconds = (
int(self.config.get("close_connection_no_voice_time", 120)) + 60
) # 在原来第一道关闭的基础上加60秒,进行二道关闭
async def handle_connection(self, ws): async def handle_connection(self, ws):
try: try:
# 获取并验证headers # 获取并验证headers
@@ -150,12 +154,17 @@ class ConnectionHandler:
self.websocket = ws self.websocket = ws
self.session_id = str(uuid.uuid4()) self.session_id = str(uuid.uuid4())
# 启动超时检查任务
self.timeout_task = asyncio.create_task(self._check_timeout())
self.welcome_msg = self.config["xiaozhi"] self.welcome_msg = self.config["xiaozhi"]
self.welcome_msg["session_id"] = self.session_id self.welcome_msg["session_id"] = self.session_id
await self.websocket.send(json.dumps(self.welcome_msg)) await self.websocket.send(json.dumps(self.welcome_msg))
# 获取差异化配置
private_config = self._initialize_private_config()
# 异步初始化 # 异步初始化
self.executor.submit(self._initialize_components) self.executor.submit(self._initialize_components, private_config)
# tts 消化线程 # tts 消化线程
self.tts_priority_thread = threading.Thread( self.tts_priority_thread = threading.Thread(
target=self._tts_priority_thread, daemon=True target=self._tts_priority_thread, daemon=True
@@ -195,18 +204,23 @@ class ConnectionHandler:
async def _route_message(self, message): async def _route_message(self, message):
"""消息路由""" """消息路由"""
# 重置超时计时器
if self.timeout_task:
self.timeout_task.cancel()
self.timeout_task = asyncio.create_task(self._check_timeout())
if isinstance(message, str): if isinstance(message, str):
await handleTextMessage(self, message) await handleTextMessage(self, message)
elif isinstance(message, bytes): elif isinstance(message, bytes):
await handleAudioMessage(self, message) await handleAudioMessage(self, message)
def _initialize_components(self): def _initialize_components(self, private_config):
"""初始化组件""" """初始化组件"""
self._initialize_models() if private_config is not None:
self._initialize_models(private_config)
"""加载提示词""" else:
self.dialogue.put(Message(role="system", content=self.prompt)) self.prompt = self.config["prompt"]
self.change_system_prompt(self.prompt)
"""加载记忆""" """加载记忆"""
self._initialize_memory() self._initialize_memory()
"""加载意图识别""" """加载意图识别"""
@@ -219,20 +233,23 @@ class ConnectionHandler:
self.dialogue.update_system_message(self.prompt) self.dialogue.update_system_message(self.prompt)
def _initialize_models(self): def _initialize_private_config(self):
read_config_from_api = self.config.get("read_config_from_api", False) read_config_from_api = self.config.get("read_config_from_api", False)
"""如果是从配置文件获取,则进行二次实例化""" """如果是从配置文件获取,则进行二次实例化"""
if not read_config_from_api: if not read_config_from_api:
return return
"""从接口获取差异化的配置进行二次实例化,非全量重新实例化""" """从接口获取差异化的配置进行二次实例化,非全量重新实例化"""
try: try:
begin_time = time.time()
private_config = get_private_config_from_api( private_config = get_private_config_from_api(
self.config, self.config,
self.headers.get("device-id", None), self.headers.get("device-id", None),
self.headers.get("client-id", None), self.headers.get("client-id", None),
) )
private_config["delete_audio"] = bool(self.config.get("delete_audio", True)) private_config["delete_audio"] = bool(self.config.get("delete_audio", True))
self.logger.bind(tag=TAG).info(f"获取差异化配置成功: {private_config}") self.logger.bind(tag=TAG).info(
f"{time.time() - begin_time} 秒,获取差异化配置成功: {private_config}"
)
except DeviceNotFoundException as e: except DeviceNotFoundException as e:
self.need_bind = True self.need_bind = True
private_config = {} private_config = {}
@@ -245,8 +262,37 @@ class ConnectionHandler:
self.logger.bind(tag=TAG).error(f"获取差异化配置失败: {e}") self.logger.bind(tag=TAG).error(f"获取差异化配置失败: {e}")
private_config = {} private_config = {}
init_vad, init_asr, init_llm, init_tts, init_memory, init_intent = ( init_tts = False
False, if private_config.get("TTS", None) is not None:
init_tts = True
self.config["TTS"] = private_config["TTS"]
self.config["selected_module"]["TTS"] = private_config["selected_module"][
"TTS"
]
try:
modules = initialize_modules(
self.logger,
private_config,
False,
False,
False,
init_tts,
False,
False,
)
except Exception as e:
self.logger.bind(tag=TAG).error(f"初始化组件失败: {e}")
modules = {}
if modules.get("tts", None) is not None:
self.tts = modules["tts"]
if modules.get("prompt", None) is not None:
self.change_system_prompt(modules["prompt"])
private_config["prompt"] = None
return private_config
def _initialize_models(self, private_config):
init_vad, init_asr, init_llm, init_memory, init_intent = (
False, False,
False, False,
False, False,
@@ -271,12 +317,6 @@ class ConnectionHandler:
self.config["selected_module"]["LLM"] = private_config["selected_module"][ self.config["selected_module"]["LLM"] = private_config["selected_module"][
"LLM" "LLM"
] ]
if private_config.get("TTS", None) is not None:
init_tts = True
self.config["TTS"] = private_config["TTS"]
self.config["selected_module"]["TTS"] = private_config["selected_module"][
"TTS"
]
if private_config.get("Memory", None) is not None: if private_config.get("Memory", None) is not None:
init_memory = True init_memory = True
self.config["Memory"] = private_config["Memory"] self.config["Memory"] = private_config["Memory"]
@@ -289,6 +329,8 @@ class ConnectionHandler:
self.config["selected_module"]["Intent"] = private_config[ self.config["selected_module"]["Intent"] = private_config[
"selected_module" "selected_module"
]["Intent"] ]["Intent"]
if private_config.get("device_max_output_size", None) is not None:
self.max_output_size = int(private_config["device_max_output_size"])
try: try:
modules = initialize_modules( modules = initialize_modules(
self.logger, self.logger,
@@ -296,7 +338,7 @@ class ConnectionHandler:
init_vad, init_vad,
init_asr, init_asr,
init_llm, init_llm,
init_tts, False,
init_memory, init_memory,
init_intent, init_intent,
) )
@@ -307,16 +349,12 @@ class ConnectionHandler:
self.vad = modules["vad"] self.vad = modules["vad"]
if modules.get("asr", None) is not None: if modules.get("asr", None) is not None:
self.asr = modules["asr"] self.asr = modules["asr"]
if modules.get("tts", None) is not None:
self.tts = modules["tts"]
if modules.get("llm", None) is not None: if modules.get("llm", None) is not None:
self.llm = modules["llm"] self.llm = modules["llm"]
if modules.get("intent", None) is not None: if modules.get("intent", None) is not None:
self.intent = modules["intent"] self.intent = modules["intent"]
if modules.get("memory", None) is not None: if modules.get("memory", None) is not None:
self.memory = modules["memory"] self.memory = modules["memory"]
if modules.get("prompt", None) is not None:
self.change_system_prompt(modules["prompt"])
def _initialize_memory(self): def _initialize_memory(self):
"""初始化记忆模块""" """初始化记忆模块"""
@@ -497,16 +535,19 @@ class ConnectionHandler:
function_id = None function_id = None
function_arguments = "" function_arguments = ""
content_arguments = "" content_arguments = ""
for response in llm_responses: for response in llm_responses:
content, tools_call = response content, tools_call = response
if "content" in response: if "content" in response:
content = response["content"] content = response["content"]
tools_call = None tools_call = None
if content is not None and len(content) > 0: if content is not None and len(content) > 0:
if len(response_message) <= 0 and ( content_arguments += content
content == "```" or "<tool_call>" in content
): if not tool_call_flag and content_arguments.startswith("<tool_call>"):
tool_call_flag = True # print("content_arguments", content_arguments)
tool_call_flag = True
if tools_call is not None: if tools_call is not None:
tool_call_flag = True tool_call_flag = True
@@ -518,9 +559,7 @@ class ConnectionHandler:
function_arguments += tools_call[0].function.arguments function_arguments += tools_call[0].function.arguments
if content is not None and len(content) > 0: if content is not None and len(content) > 0:
if tool_call_flag: if not tool_call_flag:
content_arguments += content
else:
response_message.append(content) response_message.append(content)
if self.client_abort: if self.client_abort:
@@ -581,9 +620,8 @@ class ConnectionHandler:
self.logger.bind(tag=TAG).error( self.logger.bind(tag=TAG).error(
f"function call error: {content_arguments}" f"function call error: {content_arguments}"
) )
else:
function_arguments = json.loads(function_arguments)
if not bHasError: if not bHasError:
response_message.clear()
self.logger.bind(tag=TAG).info( self.logger.bind(tag=TAG).info(
f"function_name={function_name}, function_id={function_id}, function_arguments={function_arguments}" f"function_name={function_name}, function_id={function_id}, function_arguments={function_arguments}"
) )
@@ -679,7 +717,6 @@ class ConnectionHandler:
self.tts_queue.put(future) self.tts_queue.put(future)
self.dialogue.put(Message(role="assistant", content=text)) self.dialogue.put(Message(role="assistant", content=text))
elif result.action == Action.REQLLM: # 调用函数后再请求llm生成回复 elif result.action == Action.REQLLM: # 调用函数后再请求llm生成回复
text = result.result text = result.result
if text is not None and len(text) > 0: if text is not None and len(text) > 0:
function_id = function_call_data["id"] function_id = function_call_data["id"]
@@ -706,18 +743,14 @@ class ConnectionHandler:
Message(role="tool", tool_call_id=function_id, content=text) Message(role="tool", tool_call_id=function_id, content=text)
) )
self.chat_with_function_calling(text, tool_call=True) self.chat_with_function_calling(text, tool_call=True)
elif result.action == Action.NOTFOUND: elif result.action == Action.NOTFOUND or result.action == Action.ERROR:
text = result.result text = result.result
self.recode_first_last_text(text, text_index) self.recode_first_last_text(text, text_index)
future = self.executor.submit(self.speak_and_play, text, text_index) future = self.executor.submit(self.speak_and_play, text, text_index)
self.tts_queue.put(future) self.tts_queue.put(future)
self.dialogue.put(Message(role="assistant", content=text)) self.dialogue.put(Message(role="assistant", content=text))
else: else:
text = result.result pass
self.recode_first_last_text(text, text_index)
future = self.executor.submit(self.speak_and_play, text, text_index)
self.tts_queue.put(future)
self.dialogue.put(Message(role="assistant", content=text))
def _tts_priority_thread(self): def _tts_priority_thread(self):
while not self.stop_event.is_set(): while not self.stop_event.is_set():
@@ -815,6 +848,8 @@ class ConnectionHandler:
self.logger.bind(tag=TAG).error(f"tts转换失败,{text}") self.logger.bind(tag=TAG).error(f"tts转换失败,{text}")
return None, text, text_index return None, text, text_index
self.logger.bind(tag=TAG).debug(f"TTS 文件生成完毕: {tts_file}") self.logger.bind(tag=TAG).debug(f"TTS 文件生成完毕: {tts_file}")
if self.max_output_size > 0:
add_device_output(self.headers.get("device-id"), len(text))
return tts_file, text, text_index return tts_file, text, text_index
def clearSpeakStatus(self): def clearSpeakStatus(self):
@@ -831,6 +866,11 @@ class ConnectionHandler:
async def close(self, ws=None): async def close(self, ws=None):
"""资源清理方法""" """资源清理方法"""
# 取消超时任务
if self.timeout_task:
self.timeout_task.cancel()
self.timeout_task = None
# 清理MCP资源 # 清理MCP资源
if hasattr(self, "mcp_manager") and self.mcp_manager: if hasattr(self, "mcp_manager") and self.mcp_manager:
await self.mcp_manager.cleanup_all() await self.mcp_manager.cleanup_all()
@@ -884,3 +924,15 @@ class ConnectionHandler:
self.close_after_chat = True self.close_after_chat = True
except Exception as e: except Exception as e:
self.logger.bind(tag=TAG).error(f"Chat and close error: {str(e)}") self.logger.bind(tag=TAG).error(f"Chat and close error: {str(e)}")
async def _check_timeout(self):
"""检查连接超时"""
try:
while not self.stop_event.is_set():
await asyncio.sleep(self.timeout_seconds)
if not self.stop_event.is_set():
self.logger.bind(tag=TAG).info("连接超时,准备关闭")
await self.close(self.websocket)
break
except Exception as e:
self.logger.bind(tag=TAG).error(f"超时检查任务出错: {e}")
@@ -1,9 +1,9 @@
from config.logger import setup_logging from config.logger import setup_logging
import time import time
import asyncio
from core.utils.util import remove_punctuation_and_length from core.utils.util import remove_punctuation_and_length
from core.handle.sendAudioHandle import send_stt_message from core.handle.sendAudioHandle import send_stt_message
from core.handle.intentHandler import handle_user_intent from core.handle.intentHandler import handle_user_intent
from core.utils.output_counter import check_device_output_limit
TAG = __name__ TAG = __name__
logger = setup_logging() logger = setup_logging()
@@ -51,6 +51,15 @@ async def startToChat(conn, text):
if conn.need_bind: if conn.need_bind:
await check_bind_device(conn) await check_bind_device(conn)
return return
# 如果当日的输出字数大于限定的字数
if conn.max_output_size > 0:
if check_device_output_limit(
conn.headers.get("device-id"), conn.max_output_size
):
await max_out_size(conn)
return
# 首先进行意图分析 # 首先进行意图分析
intent_handled = await handle_user_intent(conn, text) intent_handled = await handle_user_intent(conn, text)
@@ -89,6 +98,18 @@ async def no_voice_close_connect(conn):
await startToChat(conn, prompt) await startToChat(conn, prompt)
async def max_out_size(conn):
text = "不好意思,我现在有点事情要忙,明天这个时候我们再聊,约好了哦!明天不见不散,拜拜!"
await send_stt_message(conn, text)
conn.tts_first_text_index = 0
conn.tts_last_text_index = 0
conn.llm_finish_task = True
file_path = "config/assets/max_output_size.wav"
opus_packets, _ = conn.tts.audio_to_opus_data(file_path)
conn.audio_play_queue.put((opus_packets, text, 0))
conn.close_after_chat = True
async def check_bind_device(conn): async def check_bind_device(conn):
if conn.bind_code: if conn.bind_code:
# 确保bind_code是6位数字 # 确保bind_code是6位数字
@@ -35,4 +35,5 @@ class LLMProviderBase(ABC):
""" """
# For providers that don't support functions, just return regular response # For providers that don't support functions, just return regular response
for token in self.response(session_id, dialogue): for token in self.response(session_id, dialogue):
yield {"type": "content", "content": token} yield token, None
@@ -7,14 +7,8 @@ import os
# official coze sdk for Python [cozepy](https://github.com/coze-dev/coze-py) # official coze sdk for Python [cozepy](https://github.com/coze-dev/coze-py)
from cozepy import COZE_CN_BASE_URL from cozepy import COZE_CN_BASE_URL
from cozepy import ( from cozepy import Coze, TokenAuth, Message, ChatStatus, MessageContentType, ChatEventType # noqa
Coze, from core.providers.llm.system_prompt import get_system_prompt_for_function
TokenAuth,
Message,
ChatStatus,
MessageContentType,
ChatEventType,
) # noqa
TAG = __name__ TAG = __name__
logger = setup_logging() logger = setup_logging()
@@ -53,3 +47,23 @@ class LLMProvider(LLMProviderBase):
if event.event == ChatEventType.CONVERSATION_MESSAGE_DELTA: if event.event == ChatEventType.CONVERSATION_MESSAGE_DELTA:
print(event.message.content, end="", flush=True) print(event.message.content, end="", flush=True)
yield event.message.content yield event.message.content
def response_with_functions(self, session_id, dialogue, functions=None):
if len(dialogue) == 2 and functions is not None and len(functions) > 0:
# 第一次调用llm, 取最后一条用户消息,附加tool提示词
last_msg = dialogue[-1]["content"]
function_str = json.dumps(functions, ensure_ascii=False)
modify_msg = get_system_prompt_for_function(function_str) + last_msg
dialogue[-1]["content"] = modify_msg
# 如果最后一个是 role="tool",附加到user上
if len(dialogue) > 1 and dialogue[-1]["role"] == "tool":
assistant_msg = "\ntool call result: " + dialogue[-1]["content"] + "\n\n"
while len(dialogue) > 1 :
if dialogue[-1]["role"] == "user":
dialogue[-1]["content"] = assistant_msg + dialogue[-1]["content"]
break
dialogue.pop()
for token in self.response(session_id, dialogue):
yield token, None
@@ -2,6 +2,7 @@ import json
from config.logger import setup_logging from config.logger import setup_logging
import requests import requests
from core.providers.llm.base import LLMProviderBase from core.providers.llm.base import LLMProviderBase
from core.providers.llm.system_prompt import get_system_prompt_for_function
TAG = __name__ TAG = __name__
logger = setup_logging() logger = setup_logging()
@@ -81,3 +82,23 @@ class LLMProvider(LLMProviderBase):
except Exception as e: except Exception as e:
logger.bind(tag=TAG).error(f"Error in response generation: {e}") logger.bind(tag=TAG).error(f"Error in response generation: {e}")
yield "【服务响应异常】" yield "【服务响应异常】"
def response_with_functions(self, session_id, dialogue, functions=None):
if len(dialogue) == 2 and functions is not None and len(functions) > 0:
# 第一次调用llm, 取最后一条用户消息,附加tool提示词
last_msg = dialogue[-1]["content"]
function_str = json.dumps(functions, ensure_ascii=False)
modify_msg = get_system_prompt_for_function(function_str) + last_msg
dialogue[-1]["content"] = modify_msg
# 如果最后一个是 role="tool",附加到user上
if len(dialogue) > 1 and dialogue[-1]["role"] == "tool":
assistant_msg = "\ntool call result: " + dialogue[-1]["content"] + "\n\n"
while len(dialogue) > 1 :
if dialogue[-1]["role"] == "user":
dialogue[-1]["content"] = assistant_msg + dialogue[-1]["content"]
break
dialogue.pop()
for token in self.response(session_id, dialogue):
yield token, None
@@ -63,4 +63,4 @@ class LLMProvider(LLMProviderBase):
except Exception as e: except Exception as e:
logger.bind(tag=TAG).error(f"Error in Ollama function call: {e}") logger.bind(tag=TAG).error(f"Error in Ollama function call: {e}")
yield {"type": "content", "content": f"【Ollama服务响应异常: {str(e)}"} yield f"【Ollama服务响应异常: {str(e)}", None
@@ -74,4 +74,4 @@ class LLMProvider(LLMProviderBase):
except Exception as e: except Exception as e:
logger.bind(tag=TAG).error(f"Error in function call streaming: {e}") logger.bind(tag=TAG).error(f"Error in function call streaming: {e}")
yield {"type": "content", "content": f"【OpenAI服务响应异常: {e}"} yield f"【OpenAI服务响应异常: {e}", None
@@ -0,0 +1,103 @@
def get_system_prompt_for_function(functions: str) -> str:
"""
生成系统提示信息
:param functions: 可用的函数列表
:return: 系统提示信息
"""
SYSTEM_PROMPT = f"""
====
TOOL USE
You have access to a set of tools that are executed upon the user's approval. You can use one tool per message, and will receive the result of that tool use in the user's response.
You use tools step-by-step to accomplish a given task, with each tool use informed by the result of the previous tool use.
# Tool Use Formatting
Tool use is formatted using JSON-style tags. The tool name is enclosed in opening and closing tags, and each parameter is similarly enclosed within its own set of tags.
Here's the structure:
<tool_call>
{{
"name": "function name",
"arguments": {{
"param1": "value1",
"param2": "value2",
// Add more parameters as needed, if parameters are required, you must provide them
}}
}}
<tool_call>
For example:
if you got tool as follow
{{
"type": "function",
"function": {{
"name": "handle_exit_intent",
"description": "当用户想结束对话或需要退出系统时调用",
"parameters": {{
"type": "object",
"properties": {{
"say_goodbye": {{
"type": "string",
"description": "和用户友好结束对话的告别语",
}}
}},
"required": ["say_goodbye"],
}},
}},
}}
you should respond with the following format:
<tool_call>
{{
"name": "handle_exit_intent",
"arguments": {{
"say_goodbye": "再见,祝您生活愉快!"
}}
}}
</tool_call>
Always adhere to this format for the tool use to ensure proper parsing and execution.
# Tools
{functions}
# Tool Use Guidelines
1. Tools must be called in a separate message, Do not add thoughts when calling tools. The message must start with <tool_call> and end with </tool_call>, with the tool invocation JSON data in between. No additional response content is needed.
2. Choose the most appropriate tool based on the task and the tool descriptions provided. Assess if you need additional information to proceed, and which of the available tools would be most effective for gathering this information.
For example using the list_files tool is more effective than running a command like \`ls\` in the terminal. It's critical that you think about each available tool and use the one that best fits the current step in the task.
3. If multiple actions are needed, use one tool at a time per message to accomplish the task iteratively, with each tool use being informed by the result of the previous tool use. Do not assume the outcome of any tool use.
Each step must be informed by the previous step's result.
4. Formulate your tool use using the JSON format specified for each tool.
5. After each tool use, the user will respond with the result of that tool use. This result will provide you with the necessary information to continue your task or make further decisions. This response may include:
- Information about whether the tool succeeded or failed, along with any reasons for failure.
- Linter errors that may have arisen due to the changes you made, which you'll need to address.
- New terminal output in reaction to the changes, which you may need to consider or act upon.
- Any other relevant feedback or information related to the tool use.
6. ALWAYS wait for user confirmation after each tool use before proceeding. Never assume the success of a tool use without explicit confirmation of the result from the user.
7. Tool calls should contain no extra information. Only after receiving the tool's response should you integrate it into a complete reply.
It is crucial to proceed step-by-step, waiting for the user's message after each tool use before moving forward with the task. This approach allows you to:
1. Confirm the success of each step before proceeding.
2. Address any issues or errors that arise immediately.
3. Adapt your approach based on new information or unexpected results.
4. Ensure that each action builds correctly on the previous ones.
By waiting for and carefully considering the user's response after each tool use, you can react accordingly and make informed decisions about how to proceed with the task. This iterative process helps ensure the overall success and accuracy of your work.
====
USER CHAT CONTENT
The following additional message is the user's chat message, and should be followed to the best of your ability without interfering with the TOOL USE guidelines.
"""
return SYSTEM_PROMPT
@@ -0,0 +1,50 @@
import datetime
from typing import Dict, Tuple
# 全局字典,用于存储每个设备的每日输出字数
_device_daily_output: Dict[Tuple[str, datetime.date], int] = {}
# 记录最后一次检查的日期
_last_check_date: datetime.date = None
def reset_device_output():
"""
重置所有设备的每日输出字数
每天0点调用此函数
"""
_device_daily_output.clear()
def get_device_output(device_id: str) -> int:
"""
获取设备当日的输出字数
"""
current_date = datetime.datetime.now().date()
return _device_daily_output.get((device_id, current_date), 0)
def add_device_output(device_id: str, char_count: int):
"""
增加设备的输出字数
"""
current_date = datetime.datetime.now().date()
global _last_check_date
# 如果是第一次调用或者日期发生变化,清空计数器
if _last_check_date is None or _last_check_date != current_date:
_device_daily_output.clear()
_last_check_date = current_date
current_count = _device_daily_output.get((device_id, current_date), 0)
_device_daily_output[(device_id, current_date)] = current_count + char_count
def check_device_output_limit(device_id: str, max_output_size: int) -> bool:
"""
检查设备是否超过输出限制
:return: True 如果超过限制,False 如果未超过
"""
if not device_id:
return False
current_output = get_device_output(device_id)
return current_output >= max_output_size
+1 -1
View File
@@ -211,7 +211,7 @@ def check_ffmpeg_installed():
def extract_json_from_string(input_string): def extract_json_from_string(input_string):
"""提取字符串中的 JSON 部分""" """提取字符串中的 JSON 部分"""
pattern = r"(\{.*\})" pattern = r"(\{.*\})"
match = re.search(pattern, input_string) match = re.search(pattern, input_string, re.DOTALL) #添加 re.DOTALL
if match: if match:
return match.group(1) # 返回提取的 JSON 字符串 return match.group(1) # 返回提取的 JSON 字符串
return None return None
@@ -42,7 +42,6 @@ services:
- "8002:8002" - "8002:8002"
environment: environment:
- TZ=Asia/Shanghai - TZ=Asia/Shanghai
##记得改mysql和redis IP 密码
- SPRING_DATASOURCE_DRUID_URL=jdbc:mysql://xiaozhi-esp32-server-db:3306/xiaozhi_esp32_server?useUnicode=true&characterEncoding=UTF-8&serverTimezone=Asia/Shanghai&nullCatalogMeansCurrent=true&connectTimeout=30000&socketTimeout=30000&autoReconnect=true&failOverReadOnly=false&maxReconnects=10 - SPRING_DATASOURCE_DRUID_URL=jdbc:mysql://xiaozhi-esp32-server-db:3306/xiaozhi_esp32_server?useUnicode=true&characterEncoding=UTF-8&serverTimezone=Asia/Shanghai&nullCatalogMeansCurrent=true&connectTimeout=30000&socketTimeout=30000&autoReconnect=true&failOverReadOnly=false&maxReconnects=10
- SPRING_DATASOURCE_DRUID_USERNAME=root - SPRING_DATASOURCE_DRUID_USERNAME=root
- SPRING_DATASOURCE_DRUID_PASSWORD=123456 - SPRING_DATASOURCE_DRUID_PASSWORD=123456
@@ -9,7 +9,7 @@ import traceback
from pathlib import Path from pathlib import Path
from core.utils import p3 from core.utils import p3
from core.handle.sendAudioHandle import send_stt_message from core.handle.sendAudioHandle import send_stt_message
from plugins_func.register import register_function,ToolType, ActionResponse, Action from plugins_func.register import register_function, ToolType, ActionResponse, Action
TAG = __name__ TAG = __name__
@@ -18,38 +18,41 @@ logger = setup_logging()
MUSIC_CACHE = {} MUSIC_CACHE = {}
play_music_function_desc = { play_music_function_desc = {
"type": "function", "type": "function",
"function": { "function": {
"name": "play_music", "name": "play_music",
"description": "唱歌、听歌、播放音乐的方法。", "description": "唱歌、听歌、播放音乐的方法。",
"parameters": { "parameters": {
"type": "object", "type": "object",
"properties": { "properties": {
"song_name": { "song_name": {
"type": "string", "type": "string",
"description": "歌曲名称,如果用户没有指定具体歌名则为'random', 明确指定的时返回音乐的名字 示例: ```用户:播放两只老虎\n参数:两只老虎``` ```用户:播放音乐 \n参数:random ```" "description": "歌曲名称,如果用户没有指定具体歌名则为'random', 明确指定的时返回音乐的名字 示例: ```用户:播放两只老虎\n参数:两只老虎``` ```用户:播放音乐 \n参数:random ```",
}
},
"required": ["song_name"]
}
} }
} },
"required": ["song_name"],
},
},
}
@register_function('play_music', play_music_function_desc, ToolType.SYSTEM_CTL) @register_function("play_music", play_music_function_desc, ToolType.SYSTEM_CTL)
def play_music(conn, song_name: str): def play_music(conn, song_name: str):
try: try:
music_intent = f"播放音乐 {song_name}" if song_name != "random" else "随机播放音乐" music_intent = (
f"播放音乐 {song_name}" if song_name != "random" else "随机播放音乐"
)
# 检查事件循环状态 # 检查事件循环状态
if not conn.loop.is_running(): if not conn.loop.is_running():
logger.bind(tag=TAG).error("事件循环未运行,无法提交任务") logger.bind(tag=TAG).error("事件循环未运行,无法提交任务")
return ActionResponse(action=Action.RESPONSE, result="系统繁忙", response="请稍后再试") return ActionResponse(
action=Action.RESPONSE, result="系统繁忙", response="请稍后再试"
)
# 提交异步任务 # 提交异步任务
future = asyncio.run_coroutine_threadsafe( future = asyncio.run_coroutine_threadsafe(
handle_music_command(conn, music_intent), handle_music_command(conn, music_intent), conn.loop
conn.loop
) )
# 非阻塞回调处理 # 非阻塞回调处理
@@ -62,10 +65,14 @@ def play_music(conn, song_name: str):
future.add_done_callback(handle_done) future.add_done_callback(handle_done)
return ActionResponse(action=Action.RESPONSE, result="指令已接收", response="正在为您播放音乐") return ActionResponse(
action=Action.NONE, result="指令已接收", response="正在为您播放音乐"
)
except Exception as e: except Exception as e:
logger.bind(tag=TAG).error(f"处理音乐意图错误: {e}") logger.bind(tag=TAG).error(f"处理音乐意图错误: {e}")
return ActionResponse(action=Action.RESPONSE, result=str(e), response="播放音乐时出错了") return ActionResponse(
action=Action.RESPONSE, result=str(e), response="播放音乐时出错了"
)
def _extract_song_name(text): def _extract_song_name(text):
@@ -105,7 +112,9 @@ def get_music_files(music_dir, music_ext):
if ext in music_ext: if ext in music_ext:
# 添加相对路径 # 添加相对路径
music_files.append(str(file.relative_to(music_dir))) music_files.append(str(file.relative_to(music_dir)))
music_file_names.append(os.path.splitext(str(file.relative_to(music_dir)))[0]) music_file_names.append(
os.path.splitext(str(file.relative_to(music_dir)))[0]
)
return music_files, music_file_names return music_files, music_file_names
@@ -117,15 +126,20 @@ def initialize_music_handler(conn):
MUSIC_CACHE["music_dir"] = os.path.abspath( MUSIC_CACHE["music_dir"] = os.path.abspath(
MUSIC_CACHE["music_config"].get("music_dir", "./music") # 默认路径修改 MUSIC_CACHE["music_config"].get("music_dir", "./music") # 默认路径修改
) )
MUSIC_CACHE["music_ext"] = MUSIC_CACHE["music_config"].get("music_ext", (".mp3", ".wav", ".p3")) MUSIC_CACHE["music_ext"] = MUSIC_CACHE["music_config"].get(
MUSIC_CACHE["refresh_time"] = MUSIC_CACHE["music_config"].get("refresh_time", 60) "music_ext", (".mp3", ".wav", ".p3")
)
MUSIC_CACHE["refresh_time"] = MUSIC_CACHE["music_config"].get(
"refresh_time", 60
)
else: else:
MUSIC_CACHE["music_dir"] = os.path.abspath("./music") MUSIC_CACHE["music_dir"] = os.path.abspath("./music")
MUSIC_CACHE["music_ext"] = (".mp3", ".wav", ".p3") MUSIC_CACHE["music_ext"] = (".mp3", ".wav", ".p3")
MUSIC_CACHE["refresh_time"] = 60 MUSIC_CACHE["refresh_time"] = 60
# 获取音乐文件列表 # 获取音乐文件列表
MUSIC_CACHE["music_files"], MUSIC_CACHE["music_file_names"] = get_music_files(MUSIC_CACHE["music_dir"], MUSIC_CACHE["music_files"], MUSIC_CACHE["music_file_names"] = get_music_files(
MUSIC_CACHE["music_ext"]) MUSIC_CACHE["music_dir"], MUSIC_CACHE["music_ext"]
)
MUSIC_CACHE["scan_time"] = time.time() MUSIC_CACHE["scan_time"] = time.time()
return MUSIC_CACHE return MUSIC_CACHE
@@ -135,15 +149,16 @@ async def handle_music_command(conn, text):
global MUSIC_CACHE global MUSIC_CACHE
"""处理音乐播放指令""" """处理音乐播放指令"""
clean_text = re.sub(r'[^\w\s]', '', text).strip() clean_text = re.sub(r"[^\w\s]", "", text).strip()
logger.bind(tag=TAG).debug(f"检查是否是音乐命令: {clean_text}") logger.bind(tag=TAG).debug(f"检查是否是音乐命令: {clean_text}")
# 尝试匹配具体歌名 # 尝试匹配具体歌名
if os.path.exists(MUSIC_CACHE["music_dir"]): if os.path.exists(MUSIC_CACHE["music_dir"]):
if time.time() - MUSIC_CACHE["scan_time"] > MUSIC_CACHE["refresh_time"]: if time.time() - MUSIC_CACHE["scan_time"] > MUSIC_CACHE["refresh_time"]:
# 刷新音乐文件列表 # 刷新音乐文件列表
MUSIC_CACHE["music_files"], MUSIC_CACHE["music_file_names"] = get_music_files(MUSIC_CACHE["music_dir"], MUSIC_CACHE["music_files"], MUSIC_CACHE["music_file_names"] = (
MUSIC_CACHE["music_ext"]) get_music_files(MUSIC_CACHE["music_dir"], MUSIC_CACHE["music_ext"])
)
MUSIC_CACHE["scan_time"] = time.time() MUSIC_CACHE["scan_time"] = time.time()
potential_song = _extract_song_name(clean_text) potential_song = _extract_song_name(clean_text)
@@ -158,6 +173,23 @@ async def handle_music_command(conn, text):
return True return True
def _get_random_play_prompt(song_name):
"""生成随机播放引导语"""
# 移除文件扩展名
clean_name = os.path.splitext(song_name)[0]
prompts = [
f"正在为您播放,{clean_name}",
f"请欣赏歌曲,{clean_name}",
f"即将为您播放,{clean_name}",
f"为您带来,{clean_name}",
f"让我们聆听,{clean_name}",
f"接下来请欣赏,{clean_name}",
f"为您献上,{clean_name}",
]
# 直接使用random.choice,不设置seed
return random.choice(prompts)
async def play_local_music(conn, specific_file=None): async def play_local_music(conn, specific_file=None):
global MUSIC_CACHE global MUSIC_CACHE
"""播放本地音乐文件""" """播放本地音乐文件"""
@@ -180,16 +212,25 @@ async def play_local_music(conn, specific_file=None):
if not os.path.exists(music_path): if not os.path.exists(music_path):
logger.bind(tag=TAG).error(f"选定的音乐文件不存在: {music_path}") logger.bind(tag=TAG).error(f"选定的音乐文件不存在: {music_path}")
return return
text = f"正在播放{selected_music}" text = _get_random_play_prompt(selected_music)
await send_stt_message(conn, text) await send_stt_message(conn, text)
conn.tts_first_text_index = 0 conn.tts_first_text_index = 0
conn.tts_last_text_index = 0 conn.tts_last_text_index = 0
tts_file = await asyncio.to_thread(conn.tts.to_tts, text)
if tts_file is not None and os.path.exists(tts_file):
conn.tts_last_text_index = 1
opus_packets, _ = conn.tts.audio_to_opus_data(tts_file)
conn.audio_play_queue.put((opus_packets, None, 0))
os.remove(tts_file)
conn.llm_finish_task = True conn.llm_finish_task = True
if music_path.endswith(".p3"): if music_path.endswith(".p3"):
opus_packets, duration = p3.decode_opus_from_file(music_path) opus_packets, _ = p3.decode_opus_from_file(music_path)
else: else:
opus_packets, duration = conn.tts.audio_to_opus_data(music_path) opus_packets, _ = conn.tts.audio_to_opus_data(music_path)
conn.audio_play_queue.put((opus_packets, selected_music, 0)) conn.audio_play_queue.put((opus_packets, None, conn.tts_last_text_index))
except Exception as e: except Exception as e:
logger.bind(tag=TAG).error(f"播放音乐失败: {str(e)}") logger.bind(tag=TAG).error(f"播放音乐失败: {str(e)}")