v0.5.0
重构了消息发送系统,口牙
This commit is contained in:
1
.gitignore
vendored
1
.gitignore
vendored
@@ -3,6 +3,7 @@ mongodb/
|
||||
NapCat.Framework.Windows.Once/
|
||||
log/
|
||||
/test
|
||||
/src/test
|
||||
message_queue_content.txt
|
||||
message_queue_content.bat
|
||||
message_queue_window.bat
|
||||
|
||||
22
README.md
22
README.md
@@ -66,25 +66,27 @@
|
||||
- 针对每个用户创建"关系",可以对不同用户进行个性化回复,目前只有极其简单的好感度(WIP)
|
||||
- 针对每个群创建"群印象",可以对不同群进行个性化回复(WIP)
|
||||
|
||||
## 🚧 开发中功能
|
||||
|
||||
|
||||
## 开发计划TODO:LIST
|
||||
- 人格功能:WIP
|
||||
- 群氛围功能:WIP
|
||||
- 图片发送,转发功能:WIP
|
||||
- 幽默和meme功能:WIP的WIP
|
||||
- 让麦麦玩mc:WIP的WIP的WIP
|
||||
|
||||
## 开发计划TODO:LIST
|
||||
|
||||
- 兼容gif的解析和保存
|
||||
- 小程序转发链接解析
|
||||
- 对思考链长度限制
|
||||
- 修复已知bug
|
||||
- 完善文档
|
||||
- ~~完善文档~~
|
||||
- 修复转发
|
||||
- config自动生成和检测
|
||||
- log别用print
|
||||
- 给发送消息写专门的类
|
||||
- ~~config自动生成和检测~~
|
||||
- ~~log别用print~~
|
||||
- ~~给发送消息写专门的类~~
|
||||
- 改进表情包发送逻辑
|
||||
- 自动生成的回复逻辑,例如自生成的回复方向,回复风格
|
||||
- 采用截断生成加快麦麦的反应速度
|
||||
- 改进发送消息的触发:
|
||||
|
||||
## 📌 注意事项
|
||||
纯编程外行,面向cursor编程,很多代码史一样多多包涵
|
||||
@@ -99,7 +101,9 @@
|
||||
|
||||
感谢各位大佬!
|
||||
|
||||
[](https://github.com/SengokuCola/MaiMBot/graphs/contributors)
|
||||
<a href="https://github.com/SengokuCola/MaiMBot/graphs/contributors">
|
||||
<img src="https://contrib.rocks/image?repo=SengokuCola/MaiMBot&time=true" />
|
||||
</a>
|
||||
|
||||
|
||||
## Stargazers over time
|
||||
|
||||
@@ -1,6 +0,0 @@
|
||||
@echo off
|
||||
echo 正在查找并结束所有 MongoDB 进程...
|
||||
taskkill /F /IM mongod.exe
|
||||
taskkill /F /IM mongo.exe
|
||||
echo MongoDB 进程已结束
|
||||
pause
|
||||
11
setup.py
Normal file
11
setup.py
Normal file
@@ -0,0 +1,11 @@
|
||||
from setuptools import setup, find_packages
|
||||
|
||||
setup(
|
||||
name="maimai-bot",
|
||||
version="0.1",
|
||||
packages=find_packages(),
|
||||
install_requires=[
|
||||
'python-dotenv',
|
||||
'pymongo',
|
||||
],
|
||||
)
|
||||
@@ -24,4 +24,25 @@ class Database:
|
||||
def get_instance(cls) -> "Database":
|
||||
if cls._instance is None:
|
||||
raise RuntimeError("Database not initialized")
|
||||
return cls._instance
|
||||
return cls._instance
|
||||
|
||||
|
||||
#测试用
|
||||
|
||||
def get_random_group_messages(self, group_id: str, limit: int = 5):
|
||||
# 先随机获取一条消息
|
||||
random_message = list(self.db.messages.aggregate([
|
||||
{"$match": {"group_id": group_id}},
|
||||
{"$sample": {"size": 1}}
|
||||
]))[0]
|
||||
|
||||
# 获取该消息之后的消息
|
||||
subsequent_messages = list(self.db.messages.find({
|
||||
"group_id": group_id,
|
||||
"time": {"$gt": random_message["time"]}
|
||||
}).sort("time", 1).limit(limit))
|
||||
|
||||
# 将随机消息和后续消息合并
|
||||
messages = [random_message] + subsequent_messages
|
||||
|
||||
return messages
|
||||
@@ -10,6 +10,9 @@ import random
|
||||
from .relationship_manager import relationship_manager
|
||||
from ..schedule.schedule_generator import bot_schedule
|
||||
from .willing_manager import willing_manager
|
||||
from nonebot.rule import to_me
|
||||
from .bot import chat_bot
|
||||
from .emoji_manager import emoji_manager
|
||||
|
||||
|
||||
# 获取驱动器
|
||||
@@ -30,8 +33,9 @@ print("\033[1;32m[初始化数据库完成]\033[0m")
|
||||
# 导入其他模块
|
||||
from .bot import ChatBot
|
||||
from .emoji_manager import emoji_manager
|
||||
from .message_send_control import message_sender
|
||||
# from .message_send_control import message_sender
|
||||
from .relationship_manager import relationship_manager
|
||||
from .message_sender import message_manager,message_sender
|
||||
from ..memory_system.memory import memory_graph,hippocampus
|
||||
|
||||
# 初始化表情管理器
|
||||
@@ -40,8 +44,8 @@ emoji_manager.initialize()
|
||||
print(f"\033[1;32m正在唤醒{global_config.BOT_NICKNAME}......\033[0m")
|
||||
# 创建机器人实例
|
||||
chat_bot = ChatBot()
|
||||
# 注册消息处理器
|
||||
group_msg = on_message()
|
||||
# 注册群消息处理器
|
||||
group_msg = on_message(priority=5)
|
||||
# 创建定时任务
|
||||
scheduler = require("nonebot_plugin_apscheduler").scheduler
|
||||
|
||||
@@ -66,10 +70,13 @@ async def init_relationships():
|
||||
async def _(bot: Bot):
|
||||
"""Bot连接成功时的处理"""
|
||||
print(f"\033[1;38;5;208m-----------{global_config.BOT_NICKNAME}成功连接!-----------\033[0m")
|
||||
message_sender.set_bot(bot)
|
||||
asyncio.create_task(message_sender.start_processor(bot))
|
||||
await willing_manager.ensure_started()
|
||||
|
||||
|
||||
message_sender.set_bot(bot)
|
||||
print("\033[1;38;5;208m-----------消息发送器已启动!-----------\033[0m")
|
||||
asyncio.create_task(message_manager.start_processor())
|
||||
print("\033[1;38;5;208m-----------消息处理器已启动!-----------\033[0m")
|
||||
|
||||
asyncio.create_task(emoji_manager._periodic_scan(interval_MINS=global_config.EMOJI_REGISTER_INTERVAL))
|
||||
print("\033[1;38;5;208m-----------开始偷表情包!-----------\033[0m")
|
||||
|
||||
@@ -1,16 +1,16 @@
|
||||
from nonebot.adapters.onebot.v11 import GroupMessageEvent, Message as EventMessage, Bot
|
||||
from .message import Message,MessageSet
|
||||
from .message import Message, MessageSet, Message_Sending
|
||||
from .config import BotConfig, global_config
|
||||
from .storage import MessageStorage
|
||||
from .llm_generator import ResponseGenerator
|
||||
from .message_stream import MessageStream, MessageStreamContainer
|
||||
# from .message_stream import MessageStream, MessageStreamContainer
|
||||
from .topic_identifier import topic_identifier
|
||||
from random import random, choice
|
||||
from .emoji_manager import emoji_manager # 导入表情包管理器
|
||||
import time
|
||||
import os
|
||||
from .cq_code import CQCode # 导入CQCode模块
|
||||
from .message_send_control import message_sender # 导入消息发送控制器
|
||||
from .message_sender import message_manager # 导入新的消息管理器
|
||||
from .message import Message_Thinking # 导入 Message_Thinking 类
|
||||
from .relationship_manager import relationship_manager
|
||||
from .willing_manager import willing_manager # 导入意愿管理器
|
||||
@@ -25,15 +25,12 @@ class ChatBot:
|
||||
self._started = False
|
||||
|
||||
self.emoji_chance = 0.2 # 发送表情包的基础概率
|
||||
self.message_streams = MessageStreamContainer()
|
||||
self.message_sender = message_sender
|
||||
# self.message_streams = MessageStreamContainer()
|
||||
|
||||
async def _ensure_started(self):
|
||||
"""确保所有任务已启动"""
|
||||
if not self._started:
|
||||
# 只保留必要的任务
|
||||
self._started = True
|
||||
|
||||
|
||||
async def handle_message(self, event: GroupMessageEvent, bot: Bot) -> None:
|
||||
"""处理收到的群消息"""
|
||||
@@ -44,46 +41,12 @@ class ChatBot:
|
||||
|
||||
if event.user_id in global_config.ban_user_id:
|
||||
return
|
||||
|
||||
# 打印原始消息内容
|
||||
'''
|
||||
print(f"\n\033[1;33m[消息详情]\033[0m")
|
||||
# print(f"- 原始消息: {str(event.raw_message)}")
|
||||
print(f"- post_type: {event.post_type}")
|
||||
print(f"- sub_type: {event.sub_type}")
|
||||
print(f"- user_id: {event.user_id}")
|
||||
print(f"- message_type: {event.message_type}")
|
||||
# print(f"- message_id: {event.message_id}")
|
||||
# print(f"- message: {event.message}")
|
||||
print(f"- original_message: {event.original_message}")
|
||||
print(f"- raw_message: {event.raw_message}")
|
||||
# print(f"- font: {event.font}")
|
||||
print(f"- sender: {event.sender}")
|
||||
# print(f"- to_me: {event.to_me}")
|
||||
|
||||
if event.reply:
|
||||
print(f"\n\033[1;33m[回复消息详情]\033[0m")
|
||||
# print(f"- message_id: {event.reply.message_id}")
|
||||
print(f"- message_type: {event.reply.message_type}")
|
||||
print(f"- sender: {event.reply.sender}")
|
||||
# print(f"- time: {event.reply.time}")
|
||||
print(f"- message: {event.reply.message}")
|
||||
print(f"- raw_message: {event.reply.raw_message}")
|
||||
# print(f"- original_message: {event.reply.original_message}")
|
||||
'''
|
||||
|
||||
|
||||
group_info = await bot.get_group_info(group_id=event.group_id)
|
||||
|
||||
|
||||
|
||||
sender_info = await bot.get_group_member_info(group_id=event.group_id, user_id=event.user_id, no_cache=True)
|
||||
|
||||
|
||||
await relationship_manager.update_relationship(user_id = event.user_id, data = sender_info)
|
||||
await relationship_manager.update_relationship_value(user_id = event.user_id, relationship_value = 0.5)
|
||||
# print(f"\033[1;32m[关系管理]\033[0m 更新关系值: {relationship_manager.get_relationship(event.user_id).relationship_value}")
|
||||
|
||||
|
||||
message = Message(
|
||||
group_id=event.group_id,
|
||||
@@ -104,11 +67,11 @@ class ChatBot:
|
||||
|
||||
current_time = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(message.time))
|
||||
|
||||
topic1 = topic_identifier.identify_topic_jieba(message.processed_plain_text)
|
||||
topic2 = await topic_identifier.identify_topic_llm(message.processed_plain_text)
|
||||
# topic1 = topic_identifier.identify_topic_jieba(message.processed_plain_text)
|
||||
# topic2 = await topic_identifier.identify_topic_llm(message.processed_plain_text)
|
||||
topic3 = topic_identifier.identify_topic_snownlp(message.processed_plain_text)
|
||||
print(f"\033[1;32m[主题识别]\033[0m 使用jieba主题: {topic1}")
|
||||
print(f"\033[1;32m[主题识别]\033[0m 使用llm主题: {topic2}")
|
||||
# print(f"\033[1;32m[主题识别]\033[0m 使用jieba主题: {topic1}")
|
||||
# print(f"\033[1;32m[主题识别]\033[0m 使用llm主题: {topic2}")
|
||||
print(f"\033[1;32m[主题识别]\033[0m 使用snownlp主题: {topic3}")
|
||||
topic = topic3
|
||||
|
||||
@@ -123,9 +86,7 @@ class ChatBot:
|
||||
print(f"\033[1;32m[前额叶]\033[0m 对|{current_topic}|有印象")
|
||||
interested_rate = interested_num / all_num if all_num > 0 else 0
|
||||
|
||||
|
||||
await self.storage.store_message(message, topic[0] if topic else None)
|
||||
|
||||
|
||||
is_mentioned = is_mentioned_bot_in_txt(message.processed_plain_text)
|
||||
reply_probability = willing_manager.change_reply_willing_received(
|
||||
@@ -139,42 +100,53 @@ class ChatBot:
|
||||
)
|
||||
current_willing = willing_manager.get_willing(event.group_id)
|
||||
|
||||
|
||||
print(f"\033[1;32m[{current_time}][{message.group_name}]{message.user_nickname}:\033[0m {message.processed_plain_text}\033[1;36m[回复意愿:{current_willing:.2f}][概率:{reply_probability:.1f}]\033[0m")
|
||||
response = ""
|
||||
# 创建思考消息
|
||||
|
||||
if random() < reply_probability:
|
||||
|
||||
|
||||
tinking_time_point = round(time.time(), 2)
|
||||
think_id = 'mt' + str(tinking_time_point)
|
||||
thinking_message = Message_Thinking(message=message,message_id=think_id)
|
||||
message_sender.send_temp_container.add_message(thinking_message)
|
||||
|
||||
message_manager.add_message(thinking_message)
|
||||
|
||||
willing_manager.change_reply_willing_sent(thinking_message.group_id)
|
||||
|
||||
response, emotion = await self.gpt.generate_response(message)
|
||||
|
||||
if response is None:
|
||||
thinking_message.interupt=True
|
||||
|
||||
# 如果生成了回复,发送并记录
|
||||
|
||||
'''
|
||||
生成回复后的内容
|
||||
|
||||
'''
|
||||
# if response is None:
|
||||
# thinking_message.interupt=True
|
||||
|
||||
if response:
|
||||
message_set = MessageSet(event.group_id, global_config.BOT_QQ, think_id)
|
||||
# print(f"\033[1;32m[思考结束]\033[0m 思考结束,已得到回复,开始回复")
|
||||
# 找到并删除对应的thinking消息
|
||||
container = message_manager.get_container(event.group_id)
|
||||
thinking_message = None
|
||||
# 找到message,删除
|
||||
for msg in container.messages:
|
||||
if isinstance(msg, Message_Thinking) and msg.message_id == think_id:
|
||||
thinking_message = msg
|
||||
container.messages.remove(msg)
|
||||
print(f"\033[1;32m[思考消息删除]\033[0m 已找到思考消息对象,开始删除")
|
||||
break
|
||||
|
||||
#记录开始思考的时间,避免从思考到回复的时间太久
|
||||
thinking_start_time = thinking_message.thinking_start_time
|
||||
message_set = MessageSet(event.group_id, global_config.BOT_QQ, think_id) # 发送消息的id和产生发送消息的message_thinking是一致的
|
||||
#计算打字时间,1是为了模拟打字,2是避免多条回复乱序
|
||||
accu_typing_time = 0
|
||||
|
||||
# print(f"\033[1;32m[开始回复]\033[0m 开始将回复1载入发送容器")
|
||||
for msg in response:
|
||||
print(f"当前消息: {msg}")
|
||||
# print(f"\033[1;32m[回复内容]\033[0m {msg}")
|
||||
#通过时间改变时间戳
|
||||
typing_time = calculate_typing_time(msg)
|
||||
accu_typing_time += typing_time
|
||||
timepoint = tinking_time_point+accu_typing_time
|
||||
# print(f"\033[1;32m[调试]\033[0m 消息: {msg},添加!, 累计打字时间: {accu_typing_time:.2f}秒")
|
||||
timepoint = tinking_time_point + accu_typing_time
|
||||
|
||||
bot_message = Message(
|
||||
bot_message = Message_Sending(
|
||||
group_id=event.group_id,
|
||||
user_id=global_config.BOT_QQ,
|
||||
message_id=think_id,
|
||||
@@ -183,13 +155,15 @@ class ChatBot:
|
||||
processed_plain_text=msg,
|
||||
user_nickname=global_config.BOT_NICKNAME,
|
||||
group_name=message.group_name,
|
||||
time=timepoint
|
||||
time=timepoint, #记录了回复生成的时间
|
||||
thinking_start_time=thinking_start_time, #记录了思考开始的时间
|
||||
reply_message_id=message.message_id
|
||||
)
|
||||
message_set.add_message(bot_message)
|
||||
|
||||
message_sender.send_temp_container.update_thinking_message(message_set)
|
||||
|
||||
|
||||
#message_set 可以直接加入 message_manager
|
||||
print(f"\033[1;32m[回复]\033[0m 将回复载入发送容器")
|
||||
message_manager.add_message(message_set)
|
||||
|
||||
bot_response_time = tinking_time_point
|
||||
if random() < global_config.emoji_chance:
|
||||
@@ -202,20 +176,24 @@ class ChatBot:
|
||||
else:
|
||||
bot_response_time = bot_response_time + 1
|
||||
|
||||
bot_message = Message(
|
||||
group_id=event.group_id,
|
||||
user_id=global_config.BOT_QQ,
|
||||
message_id=0,
|
||||
raw_message=emoji_cq,
|
||||
plain_text=emoji_cq,
|
||||
processed_plain_text=emoji_cq,
|
||||
user_nickname=global_config.BOT_NICKNAME,
|
||||
group_name=message.group_name,
|
||||
time=bot_response_time,
|
||||
is_emoji=True,
|
||||
translate_cq=False
|
||||
)
|
||||
message_sender.send_temp_container.add_message(bot_message)
|
||||
bot_message = Message_Sending(
|
||||
group_id=event.group_id,
|
||||
user_id=global_config.BOT_QQ,
|
||||
message_id=0,
|
||||
raw_message=emoji_cq,
|
||||
plain_text=emoji_cq,
|
||||
processed_plain_text=emoji_cq,
|
||||
user_nickname=global_config.BOT_NICKNAME,
|
||||
group_name=message.group_name,
|
||||
time=bot_response_time,
|
||||
is_emoji=True,
|
||||
translate_cq=False,
|
||||
thinking_start_time=thinking_start_time,
|
||||
# reply_message_id=message.message_id
|
||||
)
|
||||
message_manager.add_message(bot_message)
|
||||
|
||||
# 如果收到新消息,提高回复意愿
|
||||
willing_manager.change_reply_willing_after_sent(event.group_id)
|
||||
willing_manager.change_reply_willing_after_sent(event.group_id)
|
||||
|
||||
# 创建全局ChatBot实例
|
||||
chat_bot = ChatBot()
|
||||
@@ -171,10 +171,6 @@ class MessageSendControl:
|
||||
except(NameError):
|
||||
pass
|
||||
|
||||
def set_bot(self, bot: Bot):
|
||||
"""设置当前bot实例"""
|
||||
self._current_bot = bot
|
||||
|
||||
async def process_group_messages(self, group_id: int):
|
||||
queue = self.send_temp_container.get_queue(group_id)
|
||||
if queue.has_messages():
|
||||
@@ -252,4 +248,4 @@ class MessageSendControl:
|
||||
self.typing_speed = (min_speed, max_speed)
|
||||
|
||||
# 创建全局实例
|
||||
message_sender = MessageSendControl()
|
||||
message_sender_control = MessageSendControl()
|
||||
@@ -21,9 +21,9 @@ config = driver.config
|
||||
|
||||
class ResponseGenerator:
|
||||
def __init__(self):
|
||||
self.model_r1 = LLM_request(model=global_config.llm_reasoning, temperature=0.7)
|
||||
self.model_v3 = LLM_request(model=global_config.llm_normal, temperature=0.7)
|
||||
self.model_r1_distill = LLM_request(model=global_config.llm_reasoning_minor, temperature=0.7)
|
||||
self.model_r1 = LLM_request(model=global_config.llm_reasoning, temperature=0.7,max_tokens=1000)
|
||||
self.model_v3 = LLM_request(model=global_config.llm_normal, temperature=0.7,max_tokens=1000)
|
||||
self.model_r1_distill = LLM_request(model=global_config.llm_reasoning_minor, temperature=0.7,max_tokens=1000)
|
||||
self.db = Database.get_instance()
|
||||
self.current_model_type = 'r1' # 默认使用 R1
|
||||
|
||||
@@ -77,22 +77,22 @@ class ResponseGenerator:
|
||||
group_id=message.group_id
|
||||
)
|
||||
|
||||
# 读空气模块
|
||||
if global_config.enable_kuuki_read:
|
||||
content_check, reasoning_content_check = await self.model_v3.generate_response(prompt_check)
|
||||
print(f"\033[1;32m[读空气]\033[0m 读空气结果为{content_check}")
|
||||
if 'yes' not in content_check.lower() and random.random() < 0.3:
|
||||
self._save_to_db(
|
||||
message=message,
|
||||
sender_name=sender_name,
|
||||
prompt=prompt,
|
||||
prompt_check=prompt_check,
|
||||
content="",
|
||||
content_check=content_check,
|
||||
reasoning_content="",
|
||||
reasoning_content_check=reasoning_content_check
|
||||
)
|
||||
return None
|
||||
# 读空气模块 简化逻辑,先停用
|
||||
# if global_config.enable_kuuki_read:
|
||||
# content_check, reasoning_content_check = await self.model_v3.generate_response(prompt_check)
|
||||
# print(f"\033[1;32m[读空气]\033[0m 读空气结果为{content_check}")
|
||||
# if 'yes' not in content_check.lower() and random.random() < 0.3:
|
||||
# self._save_to_db(
|
||||
# message=message,
|
||||
# sender_name=sender_name,
|
||||
# prompt=prompt,
|
||||
# prompt_check=prompt_check,
|
||||
# content="",
|
||||
# content_check=content_check,
|
||||
# reasoning_content="",
|
||||
# reasoning_content_check=reasoning_content_check
|
||||
# )
|
||||
# return None
|
||||
|
||||
# 生成回复
|
||||
content, reasoning_content = await model.generate_response(prompt)
|
||||
@@ -104,15 +104,17 @@ class ResponseGenerator:
|
||||
prompt=prompt,
|
||||
prompt_check=prompt_check,
|
||||
content=content,
|
||||
content_check=content_check if global_config.enable_kuuki_read else "",
|
||||
# content_check=content_check if global_config.enable_kuuki_read else "",
|
||||
reasoning_content=reasoning_content,
|
||||
reasoning_content_check=reasoning_content_check if global_config.enable_kuuki_read else ""
|
||||
# reasoning_content_check=reasoning_content_check if global_config.enable_kuuki_read else ""
|
||||
)
|
||||
|
||||
return content
|
||||
|
||||
# def _save_to_db(self, message: Message, sender_name: str, prompt: str, prompt_check: str,
|
||||
# content: str, content_check: str, reasoning_content: str, reasoning_content_check: str):
|
||||
def _save_to_db(self, message: Message, sender_name: str, prompt: str, prompt_check: str,
|
||||
content: str, content_check: str, reasoning_content: str, reasoning_content_check: str):
|
||||
content: str, reasoning_content: str,):
|
||||
"""保存对话记录到数据库"""
|
||||
self.db.db.reasoning_logs.insert_one({
|
||||
'time': time.time(),
|
||||
@@ -120,8 +122,8 @@ class ResponseGenerator:
|
||||
'user': sender_name,
|
||||
'message': message.processed_plain_text,
|
||||
'model': self.current_model_type,
|
||||
'reasoning_check': reasoning_content_check,
|
||||
'response_check': content_check,
|
||||
# 'reasoning_check': reasoning_content_check,
|
||||
# 'response_check': content_check,
|
||||
'reasoning': reasoning_content,
|
||||
'response': content,
|
||||
'prompt': prompt,
|
||||
|
||||
@@ -77,21 +77,6 @@ class Message:
|
||||
name = self.user_nickname or f"用户{self.user_id}"
|
||||
content = self.processed_plain_text
|
||||
self.detailed_plain_text = f"[{time_str}] {name}: {content}\n"
|
||||
|
||||
|
||||
def get_groupname(self, group_id: int) -> str:
|
||||
if not group_id:
|
||||
return "未知群"
|
||||
group_id = int(group_id)
|
||||
# 使用数据库单例
|
||||
db = Database.get_instance()
|
||||
# 查找用户,打印查询条件和结果
|
||||
query = {'group_id': group_id}
|
||||
group = db.db.group_info.find_one(query)
|
||||
if group:
|
||||
return group.get('group_name')
|
||||
else:
|
||||
return f"群{group_id}"
|
||||
|
||||
def parse_message_segments(self, message: str) -> List[CQCode]:
|
||||
"""
|
||||
@@ -168,45 +153,52 @@ class Message_Thinking:
|
||||
self.message_id = message_id
|
||||
|
||||
# 思考状态相关属性
|
||||
self.thinking_text = "正在思考..."
|
||||
self.time = int(time.time())
|
||||
self.thinking_start_time = int(time.time())
|
||||
self.thinking_time = 0
|
||||
self.interupt=False
|
||||
|
||||
def update_thinking_time(self):
|
||||
self.thinking_time = round(time.time(), 2) - self.time
|
||||
self.thinking_time = round(time.time(), 2) - self.thinking_start_time
|
||||
|
||||
@property
|
||||
def processed_plain_text(self) -> str:
|
||||
"""获取处理后的文本"""
|
||||
return self.thinking_text
|
||||
|
||||
@dataclass
|
||||
class Message_Sending(Message):
|
||||
"""发送中的消息类"""
|
||||
thinking_start_time: float = None # 思考开始时间
|
||||
thinking_time: float = None # 思考时间
|
||||
|
||||
def __str__(self) -> str:
|
||||
return f"[思考中] 群:{self.group_id} 用户:{self.user_nickname} 时间:{self.time} 消息ID:{self.message_id}"
|
||||
|
||||
|
||||
reply_message_id: int = None # 存储 回复的 源消息ID
|
||||
|
||||
def update_thinking_time(self):
|
||||
self.thinking_time = round(time.time(), 2) - self.thinking_start_time
|
||||
return self.thinking_time
|
||||
|
||||
|
||||
|
||||
class MessageSet:
|
||||
"""消息集合类,可以存储多个相关的消息"""
|
||||
"""消息集合类,可以存储多个发送消息"""
|
||||
def __init__(self, group_id: int, user_id: int, message_id: str):
|
||||
self.group_id = group_id
|
||||
self.user_id = user_id
|
||||
self.message_id = message_id
|
||||
self.messages: List[Message] = []
|
||||
self.messages: List[Message_Sending] = [] # 修改类型标注
|
||||
self.time = round(time.time(), 2)
|
||||
|
||||
def add_message(self, message: Message) -> None:
|
||||
"""添加消息到集合"""
|
||||
def add_message(self, message: Message_Sending) -> None:
|
||||
"""添加消息到集合,只接受Message_Sending类型"""
|
||||
if not isinstance(message, Message_Sending):
|
||||
raise TypeError("MessageSet只能添加Message_Sending类型的消息")
|
||||
self.messages.append(message)
|
||||
# 按时间排序
|
||||
self.messages.sort(key=lambda x: x.time)
|
||||
|
||||
def get_message_by_index(self, index: int) -> Optional[Message]:
|
||||
def get_message_by_index(self, index: int) -> Optional[Message_Sending]:
|
||||
"""通过索引获取消息"""
|
||||
if 0 <= index < len(self.messages):
|
||||
return self.messages[index]
|
||||
return None
|
||||
|
||||
def get_message_by_time(self, target_time: float) -> Optional[Message]:
|
||||
def get_message_by_time(self, target_time: float) -> Optional[Message_Sending]:
|
||||
"""获取最接近指定时间的消息"""
|
||||
if not self.messages:
|
||||
return None
|
||||
@@ -227,7 +219,7 @@ class MessageSet:
|
||||
"""清空所有消息"""
|
||||
self.messages.clear()
|
||||
|
||||
def remove_message(self, message: Message) -> bool:
|
||||
def remove_message(self, message: Message_Sending) -> bool:
|
||||
"""移除指定消息"""
|
||||
if message in self.messages:
|
||||
self.messages.remove(message)
|
||||
@@ -241,40 +233,4 @@ class MessageSet:
|
||||
return len(self.messages)
|
||||
|
||||
|
||||
@dataclass
|
||||
class Message_Sending(Message):
|
||||
"""发送消息数据类,继承自Message类"""
|
||||
|
||||
priority: int = 0 # 发送优先级,数字越大优先级越高
|
||||
wait_until: float = None # 等待发送的时间戳
|
||||
continue_thinking: bool = False # 是否继续思考
|
||||
|
||||
def __post_init__(self):
|
||||
super().__post_init__()
|
||||
if self.wait_until is None:
|
||||
self.wait_until = self.time
|
||||
|
||||
@property
|
||||
def can_send(self) -> bool:
|
||||
"""检查是否可以发送消息"""
|
||||
return time.time() >= self.wait_until
|
||||
|
||||
def set_wait_time(self, seconds: float) -> None:
|
||||
"""设置等待发送时间"""
|
||||
self.wait_until = time.time() + seconds
|
||||
|
||||
def set_priority(self, priority: int) -> None:
|
||||
"""设置发送优先级"""
|
||||
self.priority = priority
|
||||
|
||||
def __lt__(self, other):
|
||||
"""重写小于比较,用于优先级排序"""
|
||||
if not isinstance(other, Message_Sending):
|
||||
return NotImplemented
|
||||
return (self.priority, -self.wait_until) < (other.priority, -other.wait_until)
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,220 @@
|
||||
from typing import Union, List, Optional, Dict
|
||||
from collections import deque
|
||||
from .message import Message, Message_Thinking, MessageSet, Message_Sending
|
||||
import time
|
||||
import asyncio
|
||||
from nonebot.adapters.onebot.v11 import Bot
|
||||
from .config import global_config
|
||||
from .storage import MessageStorage
|
||||
from .cq_code import cq_code_tool
|
||||
import random
|
||||
from .utils import calculate_typing_time
|
||||
|
||||
class Message_Sender:
|
||||
"""发送器"""
|
||||
def __init__(self):
|
||||
self.message_interval = (0.5, 1) # 消息间隔时间范围(秒)
|
||||
self.last_send_time = 0
|
||||
self._current_bot = None
|
||||
|
||||
def set_bot(self, bot: Bot):
|
||||
"""设置当前bot实例"""
|
||||
self._current_bot = bot
|
||||
|
||||
async def send_group_message(
|
||||
self,
|
||||
group_id: int,
|
||||
send_text: str,
|
||||
auto_escape: bool = False,
|
||||
reply_message_id: int = None,
|
||||
at_user_id: int = None
|
||||
) -> None:
|
||||
|
||||
if not self._current_bot:
|
||||
raise RuntimeError("Bot未设置,请先调用set_bot方法设置bot实例")
|
||||
|
||||
message = send_text
|
||||
|
||||
# 如果需要回复
|
||||
if reply_message_id:
|
||||
reply_cq = cq_code_tool.create_reply_cq(reply_message_id)
|
||||
message = reply_cq + message
|
||||
|
||||
# 如果需要at
|
||||
# if at_user_id:
|
||||
# at_cq = cq_code_tool.create_at_cq(at_user_id)
|
||||
# message = at_cq + " " + message
|
||||
|
||||
|
||||
typing_time = calculate_typing_time(message)
|
||||
if typing_time > 10:
|
||||
typing_time = 10
|
||||
await asyncio.sleep(typing_time)
|
||||
|
||||
# 发送消息
|
||||
await self._current_bot.send_group_msg(
|
||||
group_id=group_id,
|
||||
message=message,
|
||||
auto_escape=auto_escape
|
||||
)
|
||||
print(f"\033[1;34m[调试]\033[0m 发送消息{message}成功")
|
||||
|
||||
|
||||
class MessageContainer:
|
||||
"""单个群的发送/思考消息容器"""
|
||||
def __init__(self, group_id: int, max_size: int = 100):
|
||||
self.group_id = group_id
|
||||
self.max_size = max_size
|
||||
self.messages = []
|
||||
self.last_send_time = 0
|
||||
self.thinking_timeout = 20 # 思考超时时间(秒)
|
||||
|
||||
def get_timeout_messages(self) -> List[Message_Sending]:
|
||||
"""获取所有超时的Message_Sending对象(思考时间超过30秒),按thinking_start_time排序"""
|
||||
current_time = time.time()
|
||||
timeout_messages = []
|
||||
|
||||
for msg in self.messages:
|
||||
if isinstance(msg, Message_Sending):
|
||||
if current_time - msg.thinking_start_time > self.thinking_timeout:
|
||||
timeout_messages.append(msg)
|
||||
|
||||
# 按thinking_start_time排序,时间早的在前面
|
||||
timeout_messages.sort(key=lambda x: x.thinking_start_time)
|
||||
|
||||
return timeout_messages
|
||||
|
||||
def get_earliest_message(self) -> Optional[Union[Message_Thinking, Message_Sending]]:
|
||||
"""获取thinking_start_time最早的消息对象"""
|
||||
if not self.messages:
|
||||
return None
|
||||
earliest_time = float('inf')
|
||||
earliest_message = None
|
||||
for msg in self.messages:
|
||||
msg_time = msg.thinking_start_time
|
||||
if msg_time < earliest_time:
|
||||
earliest_time = msg_time
|
||||
earliest_message = msg
|
||||
return earliest_message
|
||||
|
||||
def add_message(self, message: Union[Message_Thinking, Message_Sending]) -> None:
|
||||
"""添加消息到队列"""
|
||||
print(f"\033[1;32m[添加消息]\033[0m 添加消息到对应群")
|
||||
if isinstance(message, MessageSet):
|
||||
for single_message in message.messages:
|
||||
self.messages.append(single_message)
|
||||
else:
|
||||
self.messages.append(message)
|
||||
|
||||
def remove_message(self, message: Union[Message_Thinking, Message_Sending]) -> bool:
|
||||
"""移除消息,如果消息存在则返回True,否则返回False"""
|
||||
try:
|
||||
if message in self.messages:
|
||||
self.messages.remove(message)
|
||||
return True
|
||||
return False
|
||||
except Exception as e:
|
||||
print(f"\033[1;31m[错误]\033[0m 移除消息时发生错误: {e}")
|
||||
return False
|
||||
|
||||
def has_messages(self) -> bool:
|
||||
"""检查是否有待发送的消息"""
|
||||
return bool(self.messages)
|
||||
|
||||
def get_all_messages(self) -> List[Union[Message, Message_Thinking]]:
|
||||
"""获取所有消息"""
|
||||
return list(self.messages)
|
||||
|
||||
|
||||
class MessageManager:
|
||||
"""管理所有群的消息容器"""
|
||||
def __init__(self):
|
||||
self.containers: Dict[int, MessageContainer] = {}
|
||||
self.storage = MessageStorage()
|
||||
self._running = True
|
||||
|
||||
def get_container(self, group_id: int) -> MessageContainer:
|
||||
"""获取或创建群的消息容器"""
|
||||
if group_id not in self.containers:
|
||||
self.containers[group_id] = MessageContainer(group_id)
|
||||
return self.containers[group_id]
|
||||
|
||||
def add_message(self, message: Union[Message_Thinking, Message_Sending, MessageSet]) -> None:
|
||||
container = self.get_container(message.group_id)
|
||||
container.add_message(message)
|
||||
|
||||
async def process_group_messages(self, group_id: int):
|
||||
"""处理群消息"""
|
||||
print(f"\033[1;34m[调试]\033[0m 开始处理群{group_id}的消息")
|
||||
container = self.get_container(group_id)
|
||||
if container.has_messages():
|
||||
#最早的对象,可能是思考消息,也可能是发送消息
|
||||
message_earliest = container.get_earliest_message() #一个message_thinking or message_sending
|
||||
|
||||
#一个月后删了
|
||||
if not message_earliest:
|
||||
print(f"\033[1;34m[BUG,如果出现这个,说明有BUG,3月4日留]\033[0m ")
|
||||
return
|
||||
|
||||
#如果是思考消息
|
||||
if isinstance(message_earliest, Message_Thinking):
|
||||
#优先等待这条消息
|
||||
message_earliest.update_thinking_time()
|
||||
thinking_time = message_earliest.thinking_time
|
||||
print(f"\033[1;34m[调试]\033[0m 消息正在思考中,已思考{int(thinking_time)}秒")
|
||||
else:# 如果不是message_thinking就只能是message_sending
|
||||
print(f"\033[1;34m[调试]\033[0m 消息'{message_earliest.processed_plain_text}'正在发送中")
|
||||
#直接发,等什么呢
|
||||
if message_earliest.update_thinking_time() < 30:
|
||||
await message_sender.send_group_message(group_id, message_earliest.processed_plain_text, auto_escape=False)
|
||||
else:
|
||||
await message_sender.send_group_message(group_id, message_earliest.processed_plain_text, auto_escape=False, reply_message_id=message_earliest.reply_message_id)
|
||||
|
||||
#移除消息
|
||||
if message_earliest.is_emoji:
|
||||
message_earliest.processed_plain_text = "[表情包]"
|
||||
await self.storage.store_message(message_earliest, None)
|
||||
|
||||
container.remove_message(message_earliest)
|
||||
|
||||
#获取并处理超时消息
|
||||
message_timeout = container.get_timeout_messages() #也许是一堆message_sending
|
||||
if message_timeout:
|
||||
print(f"\033[1;34m[调试]\033[0m 发现{len(message_timeout)}条超时消息")
|
||||
for msg in message_timeout:
|
||||
if msg == message_earliest:
|
||||
continue # 跳过已经处理过的消息
|
||||
|
||||
try:
|
||||
#发送
|
||||
if msg.update_thinking_time() < 30:
|
||||
await message_sender.send_group_message(group_id, msg.processed_plain_text, auto_escape=False)
|
||||
else:
|
||||
await message_sender.send_group_message(group_id, msg.processed_plain_text, auto_escape=False, reply_message_id=msg.reply_message_id)
|
||||
|
||||
#如果是表情包,则替换为"[表情包]"
|
||||
if msg.is_emoji:
|
||||
msg.processed_plain_text = "[表情包]"
|
||||
await self.storage.store_message(msg, None)
|
||||
|
||||
# 安全地移除消息
|
||||
if not container.remove_message(msg):
|
||||
print(f"\033[1;33m[警告]\033[0m 尝试删除不存在的消息")
|
||||
except Exception as e:
|
||||
print(f"\033[1;31m[错误]\033[0m 处理超时消息时发生错误: {e}")
|
||||
continue
|
||||
|
||||
async def start_processor(self):
|
||||
"""启动消息处理器"""
|
||||
while self._running:
|
||||
await asyncio.sleep(1)
|
||||
tasks = []
|
||||
for group_id in self.containers.keys():
|
||||
tasks.append(self.process_group_messages(group_id))
|
||||
|
||||
await asyncio.gather(*tasks)
|
||||
|
||||
# 创建全局消息管理器实例
|
||||
message_manager = MessageManager()
|
||||
# 创建全局发送器实例
|
||||
message_sender = Message_Sender()
|
||||
|
||||
@@ -255,15 +255,9 @@ class PromptBuilder:
|
||||
|
||||
def get_prompt_info(self,message:str,threshold:float):
|
||||
related_info = ''
|
||||
if len(message) > 10:
|
||||
message_segments = [message[i:i+10] for i in range(0, len(message), 10)]
|
||||
for segment in message_segments:
|
||||
embedding = get_embedding(segment)
|
||||
related_info += self.get_info_from_db(embedding,threshold=threshold)
|
||||
|
||||
else:
|
||||
embedding = get_embedding(message)
|
||||
related_info += self.get_info_from_db(embedding,threshold=threshold)
|
||||
print(f"\033[1;34m[调试]\033[0m 获取知识库内容,元消息:{message[:30]}...,消息长度: {len(message)}")
|
||||
embedding = get_embedding(message)
|
||||
related_info += self.get_info_from_db(embedding,threshold=threshold)
|
||||
|
||||
return related_info
|
||||
|
||||
|
||||
@@ -69,10 +69,10 @@ class LLM_request:
|
||||
await asyncio.sleep(wait_time)
|
||||
else:
|
||||
logger.critical(f"请求失败: {str(e)}", exc_info=True)
|
||||
return f"请求失败: {str(e)}", ""
|
||||
raise RuntimeError(f"API请求失败: {str(e)}")
|
||||
|
||||
logger.error("达到最大重试次数,请求仍然失败")
|
||||
return "达到最大重试次数,请求仍然失败", ""
|
||||
raise RuntimeError("达到最大重试次数,API请求仍然失败")
|
||||
|
||||
async def generate_response_for_image(self, prompt: str, image_base64: str) -> Tuple[str, str]:
|
||||
"""根据输入的提示和图片生成模型的异步响应"""
|
||||
@@ -137,10 +137,10 @@ class LLM_request:
|
||||
await asyncio.sleep(wait_time)
|
||||
else:
|
||||
logger.critical(f"请求失败: {str(e)}", exc_info=True)
|
||||
return f"请求失败: {str(e)}", ""
|
||||
raise RuntimeError(f"API请求失败: {str(e)}")
|
||||
|
||||
logger.error("达到最大重试次数,请求仍然失败")
|
||||
return "达到最大重试次数,请求仍然失败", ""
|
||||
raise RuntimeError("达到最大重试次数,API请求仍然失败")
|
||||
|
||||
def generate_response_for_image_sync(self, prompt: str, image_base64: str) -> Tuple[str, str]:
|
||||
"""同步方法:根据输入的提示和图片生成模型的响应"""
|
||||
@@ -205,10 +205,10 @@ class LLM_request:
|
||||
time.sleep(wait_time)
|
||||
else:
|
||||
logger.critical(f"请求失败: {str(e)}", exc_info=True)
|
||||
return f"请求失败: {str(e)}", ""
|
||||
raise RuntimeError(f"API请求失败: {str(e)}")
|
||||
|
||||
logger.error("达到最大重试次数,请求仍然失败")
|
||||
return "达到最大重试次数,请求仍然失败", ""
|
||||
raise RuntimeError("达到最大重试次数,API请求仍然失败")
|
||||
|
||||
def get_embedding_sync(self, text: str, model: str = "BAAI/bge-m3") -> Union[list, None]:
|
||||
"""同步方法:获取文本的embedding向量
|
||||
|
||||
@@ -5,6 +5,7 @@ from ...common.database import Database # 使用正确的导入语法
|
||||
from src.plugins.chat.config import global_config
|
||||
from nonebot import get_driver
|
||||
from ..models.utils_model import LLM_request
|
||||
from loguru import logger
|
||||
|
||||
driver = get_driver()
|
||||
config = driver.config
|
||||
@@ -42,8 +43,6 @@ class ScheduleGenerator:
|
||||
self.yesterday_schedule_text, self.yesterday_schedule = await self.generate_daily_schedule(target_date=yesterday,read_only=True)
|
||||
|
||||
async def generate_daily_schedule(self, target_date: datetime.datetime = None,read_only:bool = False) -> Dict[str, str]:
|
||||
if target_date is None:
|
||||
target_date = datetime.datetime.now()
|
||||
|
||||
date_str = target_date.strftime("%Y-%m-%d")
|
||||
weekday = target_date.strftime("%A")
|
||||
@@ -65,7 +64,11 @@ class ScheduleGenerator:
|
||||
3. 晚上的计划和休息时间
|
||||
请按照时间顺序列出具体时间点和对应的活动,用一个时间点而不是时间段来表示时间,用逗号,隔开时间与活动,格式为"时间,活动",例如"08:00,起床"。"""
|
||||
|
||||
schedule_text, _ = await self.llm_scheduler.generate_response(prompt)
|
||||
try:
|
||||
schedule_text, _ = await self.llm_scheduler.generate_response(prompt)
|
||||
except Exception as e:
|
||||
logger.error(f"生成日程失败: {str(e)}")
|
||||
schedule_text = "生成日程时出错了"
|
||||
# print(self.schedule_text)
|
||||
self.db.db.schedule.insert_one({"date": date_str, "schedule": schedule_text})
|
||||
else:
|
||||
|
||||
Reference in New Issue
Block a user