mirror of
https://github.com/zhenxun-org/zhenxun_bot.git
synced 2026-10-07 21:00:21 +08:00
⚡ 添加数据库查询次数记录
This commit is contained in:
@@ -1,7 +1,10 @@
|
||||
import asyncio
|
||||
from collections import deque
|
||||
from typing import Any
|
||||
|
||||
from nonebot.adapters import Bot, Message
|
||||
from nonebot.adapters.onebot.v11 import MessageSegment
|
||||
from nonebot_plugin_apscheduler import scheduler
|
||||
|
||||
from zhenxun.configs.config import Config
|
||||
from zhenxun.models.bot_message_store import BotMessageStore
|
||||
@@ -13,6 +16,45 @@ from zhenxun.utils.platform import PlatformUtils
|
||||
LOG_COMMAND = "MessageHook"
|
||||
|
||||
|
||||
_BOT_MSG_BUFFER: deque[dict[str, Any]] = deque()
|
||||
_BOT_MSG_BUFFER_LOCK = asyncio.Lock()
|
||||
_BOT_MSG_BULK_SIZE = 50
|
||||
_PENDING_TASKS: set[asyncio.Task] = set()
|
||||
|
||||
|
||||
async def _flush_bot_messages():
|
||||
async with _BOT_MSG_BUFFER_LOCK:
|
||||
if not _BOT_MSG_BUFFER:
|
||||
return
|
||||
items: list[dict[str, Any]] = []
|
||||
while _BOT_MSG_BUFFER:
|
||||
items.append(_BOT_MSG_BUFFER.popleft())
|
||||
try:
|
||||
await BotMessageStore.bulk_create([BotMessageStore(**it) for it in items])
|
||||
except Exception as e:
|
||||
logger.warning("批量写入BotMessageStore失败", LOG_COMMAND, e=e)
|
||||
# 尝试降级逐条写入,避免数据全部丢失
|
||||
try:
|
||||
for it in items:
|
||||
await BotMessageStore.create(**it)
|
||||
except Exception as e2:
|
||||
logger.warning("逐条写入BotMessageStore失败", LOG_COMMAND, e=e2)
|
||||
|
||||
|
||||
async def _enqueue_bot_message(item: dict[str, Any]):
|
||||
async with _BOT_MSG_BUFFER_LOCK:
|
||||
_BOT_MSG_BUFFER.append(item)
|
||||
if len(_BOT_MSG_BUFFER) >= _BOT_MSG_BULK_SIZE:
|
||||
task = asyncio.create_task(_flush_bot_messages())
|
||||
_PENDING_TASKS.add(task)
|
||||
task.add_done_callback(_PENDING_TASKS.discard)
|
||||
|
||||
|
||||
@scheduler.scheduled_job("interval", seconds=10)
|
||||
async def _flush_bot_messages_job():
|
||||
await _flush_bot_messages()
|
||||
|
||||
|
||||
def replace_message(message: Message) -> str:
|
||||
"""将消息中的at、image、record、face替换为字符串
|
||||
|
||||
@@ -95,18 +137,20 @@ async def handle_api_result(
|
||||
if not Config.get_config("hook", "RECORD_BOT_SENT_MESSAGES"):
|
||||
return
|
||||
try:
|
||||
await BotMessageStore.create(
|
||||
bot_id=bot.self_id,
|
||||
user_id=user_id,
|
||||
group_id=group_id,
|
||||
sent_type=BotSentType.GROUP
|
||||
if message_type == "group"
|
||||
else BotSentType.PRIVATE,
|
||||
text=replace_message(message),
|
||||
plain_text=message.extract_plain_text()
|
||||
if isinstance(message, Message)
|
||||
else replace_message(message),
|
||||
platform=PlatformUtils.get_platform(bot),
|
||||
await _enqueue_bot_message(
|
||||
{
|
||||
"bot_id": bot.self_id,
|
||||
"user_id": user_id,
|
||||
"group_id": group_id,
|
||||
"sent_type": BotSentType.GROUP
|
||||
if message_type == "group"
|
||||
else BotSentType.PRIVATE,
|
||||
"text": replace_message(message),
|
||||
"plain_text": message.extract_plain_text()
|
||||
if isinstance(message, Message)
|
||||
else replace_message(message),
|
||||
"platform": PlatformUtils.get_platform(bot),
|
||||
}
|
||||
)
|
||||
logger.debug(f"消息发送记录,message: {format_message_for_log(message)}")
|
||||
except Exception as e:
|
||||
|
||||
Reference in New Issue
Block a user