feat: 群分析本地消息库 / 处理中占位图 / Web 日志窗口重构
- 群分析: 新增 history_store 只读本地消息源(读 learning_chat 落库,含 uninfo 昵称补齐/去重/截断判定),适配器优先读本地库、失败回退 OneBot 分页;分页 锚点字段回退链修复 NapCat 传 message_seq 翻页断裂;新增「本地记录」开关 - core/message_utils: 新增 common_proc_reply 占位图通用回复(引用消息 + 处理中 动图,支持后台任务显式指定 target),群分析/战况/倒放改用 - web_hub + web: 日志页改固定窗口滚动 + 翻页锚定 + 自动换行,SSE 日志轮转发 reset 帧,入口 HTML no-cache,行数统计增量缓存,CPU 改非阻塞采样,登录信息缓存 - 插件内 CLAUDE.md / DESIGN.md 不入库(.gitignore),galgame_card 两份文档取消 跟踪(文件保留在磁盘) Co-Authored-By: Claude Code <noreply@anthropic.com>
This commit is contained in:
@@ -17,6 +17,8 @@ from nonebot.params import CommandArg
|
||||
from nonebot.plugin import PluginMetadata
|
||||
from nonebot.rule import to_me
|
||||
|
||||
from ...core import message_utils
|
||||
|
||||
require("nonebot_plugin_alconna")
|
||||
from nonebot_plugin_alconna import UniMessage # noqa: E402
|
||||
|
||||
@@ -31,7 +33,6 @@ from .service import get_services, register_bot_adapter # noqa: E402
|
||||
from .renderer import html_render # noqa: E402
|
||||
from .templates import list_templates, template_exists # noqa: E402
|
||||
|
||||
|
||||
# ---------- 指令定义(全部要求 @ 机器人) ----------
|
||||
|
||||
analysis_cmd = on_command(
|
||||
@@ -192,7 +193,7 @@ async def _(bot: Bot, event: GroupMessageEvent):
|
||||
|
||||
group_id = str(event.group_id)
|
||||
days = _parse_days(event)
|
||||
await UniMessage.text("正在启动分析引擎,正在拉取最近消息...").send()
|
||||
await message_utils.common_proc_reply(event.message_id)
|
||||
|
||||
try:
|
||||
svc = get_services()
|
||||
@@ -347,6 +348,7 @@ async def _(bot: Bot, event: GroupMessageEvent):
|
||||
f"用户称号: {'开' if cm.get_user_title_analysis_enabled() else '关'}",
|
||||
f"金句分析: {'开' if cm.get_golden_quote_analysis_enabled() else '关'}",
|
||||
f"聊天质量: {'开' if cm.get_chat_quality_analysis_enabled() else '关'}",
|
||||
f"本地记录: {'开' if cm.get_use_local_history() else '关'}",
|
||||
"用法:设置分析 [参数] [值],如:设置分析 天数 3",
|
||||
]
|
||||
await UniMessage.text(chr(10).join(lines)).send()
|
||||
@@ -362,8 +364,12 @@ _SETTING_MAP = {
|
||||
"称号": ("set_user_title_analysis_enabled", str, "用户称号已{}"),
|
||||
"金句": ("set_golden_quote_analysis_enabled", str, "金句分析已{}"),
|
||||
"聊天质量": ("set_chat_quality_analysis_enabled", str, "聊天质量已{}"),
|
||||
"本地记录": ("set_use_local_history", str, "本地记录已{}"),
|
||||
}
|
||||
|
||||
# 走开/关布尔语义的键(值为 _BOOL_WORDS 里的词)
|
||||
_BOOL_KEYS = {"话题", "称号", "金句", "聊天质量", "本地记录"}
|
||||
|
||||
_BOOL_WORDS = {"开": True, "on": True, "true": True, "关": False, "off": False, "false": False}
|
||||
|
||||
|
||||
@@ -375,7 +381,9 @@ async def _(bot: Bot, event: GroupMessageEvent, args: tuple = CommandArg()):
|
||||
tokens = _cmd_tokens(args)
|
||||
if len(tokens) < 2:
|
||||
await UniMessage.text(
|
||||
"用法:设置分析 [参数] [值]" + chr(10) + "参数:天数/窗口/最大消息/最小消息/输出格式/话题/称号/金句/聊天质量"
|
||||
"用法:设置分析 [参数] [值]"
|
||||
+ chr(10)
|
||||
+ "参数:天数/窗口/最大消息/最小消息/输出格式/话题/称号/金句/聊天质量/本地记录"
|
||||
).send()
|
||||
return
|
||||
key = tokens[0]
|
||||
@@ -396,7 +404,7 @@ async def _(bot: Bot, event: GroupMessageEvent, args: tuple = CommandArg()):
|
||||
await UniMessage.text(f"{key} 已设为 {v}").send()
|
||||
elif conv is str:
|
||||
v = val.lower()
|
||||
if key in ("话题", "称号", "金句", "聊天质量"):
|
||||
if key in _BOOL_KEYS:
|
||||
if v not in _BOOL_WORDS:
|
||||
await UniMessage.text("布尔值请填:开/关 或 on/off").send()
|
||||
return
|
||||
|
||||
@@ -3,12 +3,10 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import base64
|
||||
import time
|
||||
from datetime import datetime, timedelta
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from . import history_store
|
||||
from .core.domain.value_objects.unified_group import UnifiedGroup, UnifiedMember
|
||||
from .core.domain.value_objects.unified_message import (
|
||||
MessageContent,
|
||||
@@ -17,6 +15,29 @@ from .core.domain.value_objects.unified_message import (
|
||||
)
|
||||
from .core.utils.logger import logger
|
||||
|
||||
# 分页锚点字段链(后端差异实测):NapCat 的 message_seq 参数按 OB11
|
||||
# message_id(shortId 短 ID 映射) 解析,传消息内的 message_seq(内核 msgSeq)
|
||||
# 命中不了映射 → NapCat 抛"消息不存在",翻页断裂;go-cqhttp/LLOneBot 等
|
||||
# 则识别 message_seq。优先 message_id,翻页无进度时依次回退。
|
||||
ANCHOR_FIELDS = ("message_id", "message_seq", "seq", "real_id")
|
||||
|
||||
|
||||
def _fmt_ts(ts: int) -> str:
|
||||
"""时间戳转日志用的可读时间。"""
|
||||
try:
|
||||
return datetime.fromtimestamp(ts).strftime("%Y-%m-%d %H:%M")
|
||||
except (OSError, OverflowError, ValueError):
|
||||
return str(ts)
|
||||
|
||||
|
||||
def _pick_anchor(raw: dict[str, Any], fields: tuple[str, ...], start: int):
|
||||
"""从消息中按字段优先级提取分页锚点值"""
|
||||
for field in fields[start:]:
|
||||
val = raw.get(field)
|
||||
if val is not None:
|
||||
return val
|
||||
return None
|
||||
|
||||
|
||||
class OneBotAdapter:
|
||||
"""面向 NoneBot OneBot V11 的最小适配器。"""
|
||||
@@ -31,6 +52,8 @@ class OneBotAdapter:
|
||||
self.platform_id = str(self.config.get("platform_id") or "onebot")
|
||||
self.bot_self_ids = [str(x) for x in self.config.get("bot_self_ids", [])]
|
||||
self.filter_bot_messages = bool(self.config.get("filter_bot_messages", True))
|
||||
# 优先读本地消息库(learning_chat 落库),取不到再回退接口分页
|
||||
self.use_local_history = bool(self.config.get("use_local_history", True))
|
||||
|
||||
# —— 消息拉取 ——
|
||||
async def fetch_messages(
|
||||
@@ -52,8 +75,21 @@ class OneBotAdapter:
|
||||
start_ts = int(
|
||||
(datetime.now() - timedelta(days=days)).timestamp()
|
||||
)
|
||||
end_ts = int(datetime.now().timestamp())
|
||||
|
||||
# 本地消息库优先(完整且毫秒级);before_id 无法映射到时间戳,故跳过
|
||||
if self.use_local_history and not before_id:
|
||||
local = await self._fetch_from_local_store(
|
||||
group_id, start_ts, end_ts, max_count
|
||||
)
|
||||
if local:
|
||||
return local
|
||||
|
||||
current_anchor = before_id
|
||||
anchor_idx = 0
|
||||
no_progress_pages = 0
|
||||
last_earliest: dict[str, Any] | None = None
|
||||
|
||||
while len(all_raw) < max_count:
|
||||
fetch_count = min(chunk_size, max_count - len(all_raw))
|
||||
params: dict[str, Any] = {
|
||||
@@ -85,13 +121,34 @@ class OneBotAdapter:
|
||||
break
|
||||
messages = result.get("messages", [])
|
||||
if not messages:
|
||||
break
|
||||
# 空页可能意味着当前锚点字段不被后端识别(如 NapCat 外的
|
||||
# 实现收到 message_id),换下一字段重试一次,仍空则结束。
|
||||
if current_anchor is None or last_earliest is None:
|
||||
break
|
||||
if no_progress_pages >= 1 or anchor_idx >= len(ANCHOR_FIELDS) - 1:
|
||||
logger.warning(
|
||||
"OneBot 分页拉取: 返回空页且锚点字段已耗尽,停止回溯"
|
||||
)
|
||||
break
|
||||
no_progress_pages += 1
|
||||
anchor_idx += 1
|
||||
current_anchor = _pick_anchor(
|
||||
last_earliest, ANCHOR_FIELDS, anchor_idx
|
||||
)
|
||||
if current_anchor is None:
|
||||
break
|
||||
logger.warning(
|
||||
f"OneBot 分页拉取: 空页,锚点字段切换为 {ANCHOR_FIELDS[anchor_idx]}"
|
||||
)
|
||||
continue
|
||||
|
||||
first = messages[0]
|
||||
last = messages[-1]
|
||||
earliest = first if first.get("time", 0) <= last.get("time", 0) else last
|
||||
last_earliest = earliest
|
||||
chunk_earliest_ts = earliest.get("time", 0)
|
||||
|
||||
prev_len = len(all_raw)
|
||||
for raw in messages:
|
||||
msg_time = raw.get("time", 0)
|
||||
msg_id = str(raw.get("message_id", ""))
|
||||
@@ -100,19 +157,42 @@ class OneBotAdapter:
|
||||
if start_ts <= msg_time <= int(datetime.now().timestamp()):
|
||||
all_raw.append(raw)
|
||||
seen_raw_ids.add(msg_id)
|
||||
added = len(all_raw) - prev_len
|
||||
|
||||
seq_val = (
|
||||
earliest.get("message_seq")
|
||||
or earliest.get("real_id")
|
||||
or earliest.get("seq")
|
||||
)
|
||||
mid_val = earliest.get("message_id")
|
||||
new_anchor = seq_val if seq_val is not None else mid_val
|
||||
if chunk_earliest_ts <= start_ts:
|
||||
logger.info(
|
||||
f"OneBot 分页拉取: 已到达起始时间,共 {len(all_raw)} 条"
|
||||
)
|
||||
break
|
||||
|
||||
if added == 0:
|
||||
# 本页没有新增消息:锚点不生效(后端不支持该字段)或数据已取尽。
|
||||
# 先切换锚点字段再试一次,仍无新增则结束。
|
||||
no_progress_pages += 1
|
||||
if no_progress_pages >= 2:
|
||||
logger.warning(
|
||||
"OneBot 分页拉取: 连续 2 页无新增,停止回溯"
|
||||
)
|
||||
break
|
||||
if anchor_idx < len(ANCHOR_FIELDS) - 1:
|
||||
anchor_idx += 1
|
||||
logger.warning(
|
||||
f"OneBot 分页拉取: 锚点字段切换为 {ANCHOR_FIELDS[anchor_idx]}"
|
||||
)
|
||||
else:
|
||||
no_progress_pages = 0
|
||||
|
||||
new_anchor = _pick_anchor(earliest, ANCHOR_FIELDS, anchor_idx)
|
||||
if new_anchor is None:
|
||||
break
|
||||
if current_anchor and str(new_anchor) == str(current_anchor):
|
||||
logger.info("OneBot 分页拉取: 锚点未位移,历史已取尽")
|
||||
break
|
||||
current_anchor = new_anchor
|
||||
logger.info(
|
||||
f"OneBot 分页拉取进度: {len(all_raw)} 条,"
|
||||
f"锚点({ANCHOR_FIELDS[anchor_idx]}): {new_anchor}"
|
||||
)
|
||||
await asyncio.sleep(0.05)
|
||||
|
||||
unified: list[UnifiedMessage] = []
|
||||
@@ -131,6 +211,54 @@ class OneBotAdapter:
|
||||
logger.warning(f"OneBot 分页获取消息失败: {e}")
|
||||
return []
|
||||
|
||||
async def _fetch_from_local_store(
|
||||
self, group_id: str, start_ts: int, end_ts: int, max_count: int
|
||||
) -> list[UnifiedMessage]:
|
||||
"""从本地消息库(learning_chat 落库)取群历史;不可用时返回空以回退分页。"""
|
||||
try:
|
||||
result = await asyncio.to_thread(
|
||||
history_store.fetch_group_messages,
|
||||
group_id,
|
||||
start_ts,
|
||||
end_ts,
|
||||
max_count,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning(f"本地历史读取失败,回退接口分页: {e}")
|
||||
return []
|
||||
|
||||
if not result:
|
||||
logger.info(
|
||||
f"本地历史无数据({result.error or '窗口内无消息'}),"
|
||||
"回退 OneBot 分页拉取"
|
||||
)
|
||||
return []
|
||||
|
||||
logger.info(
|
||||
f"本地历史记录拉取: group={group_id}, source={history_store.MESSAGE_TABLE}, "
|
||||
f"count={len(result.messages)}, "
|
||||
f"窗口=[{_fmt_ts(start_ts)}..{_fmt_ts(end_ts)}], "
|
||||
f"名称覆盖={result.names_resolved}/{len(result.messages)}, "
|
||||
f"去重={result.duplicates}"
|
||||
)
|
||||
if result.truncated:
|
||||
window_total = (
|
||||
result.window_total if result.window_total is not None else "?"
|
||||
)
|
||||
earliest = result.messages[0]["time"] if result.messages else end_ts
|
||||
logger.warning(
|
||||
f"本地历史截断: group={group_id}, max_messages={max_count}, "
|
||||
f"窗口内共 {window_total} 条, 实际取最近 {len(result.messages)} 条, "
|
||||
f"最早={_fmt_ts(earliest)}, 窗口起点={_fmt_ts(start_ts)}"
|
||||
)
|
||||
|
||||
unified: list[UnifiedMessage] = []
|
||||
for raw in result.messages:
|
||||
converted = self._convert_message(raw, group_id)
|
||||
if converted:
|
||||
unified.append(converted)
|
||||
return unified
|
||||
|
||||
def _convert_message(self, raw: dict, group_id: str) -> UnifiedMessage | None:
|
||||
try:
|
||||
sender = raw.get("sender", {})
|
||||
|
||||
@@ -98,6 +98,8 @@ def _default_config() -> dict:
|
||||
cfg.setdefault("basic", {}).setdefault("max_messages", 1000)
|
||||
cfg.setdefault("basic", {}).setdefault("min_messages_threshold", 50)
|
||||
cfg.setdefault("basic", {}).setdefault("filter_bot_messages", True)
|
||||
# 优先读本地消息库(learning_chat 落库),取不到再回退接口分页
|
||||
cfg.setdefault("basic", {}).setdefault("use_local_history", True)
|
||||
cfg.setdefault("analysis_features", {}).setdefault("chat_quality_analysis_enabled", False)
|
||||
cfg.setdefault("incremental", {}).setdefault("incremental_enabled", False)
|
||||
return cfg
|
||||
|
||||
+9
@@ -770,6 +770,15 @@ class ConfigManager:
|
||||
self._ensure_group("basic")["filter_bot_messages"] = enabled
|
||||
self.config.save_config()
|
||||
|
||||
def get_use_local_history(self) -> bool:
|
||||
"""获取是否优先从本地消息库读取群历史。"""
|
||||
return self._get_group("basic").get("use_local_history", True)
|
||||
|
||||
def set_use_local_history(self, enabled: bool):
|
||||
"""设置是否优先从本地消息库读取群历史(关闭则始终走接口分页)。"""
|
||||
self._ensure_group("basic")["use_local_history"] = enabled
|
||||
self.config.save_config()
|
||||
|
||||
def get_html_output_dir(self) -> str:
|
||||
"""获取HTML输出目录"""
|
||||
|
||||
|
||||
@@ -0,0 +1,281 @@
|
||||
"""本地群历史记录源:只读读取 learning_chat 落库的群消息。
|
||||
|
||||
`nonebot_plugin_learning_chat` 会把机器人收到的每条群消息写入
|
||||
`learning_chat_message`(与该群是否开启学习无关),比 OneBot
|
||||
`get_group_msg_history` 分页更完整——部分后端只能取回一页就被截断。
|
||||
昵称/群名片从同库的 `nonebot_plugin_uninfo_*` 表离线补齐,不额外调接口。
|
||||
|
||||
模块只依赖标准库(logger 做守卫导入),便于单测用 importlib 裸加载;
|
||||
读连接一律 `mode=ro`(绝不建库),任何失败都返回带 error 的空结果。
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import sqlite3
|
||||
from dataclasses import dataclass, field
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
try: # 插件包内导入:走插件 logger
|
||||
from .core.utils.logger import logger
|
||||
except Exception: # 单测裸加载(无包上下文)时降级为标准库 logger
|
||||
import logging
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
MESSAGE_TABLE = "learning_chat_message"
|
||||
# nonebot_plugin_uninfo 的 SceneType.GROUP
|
||||
SCENE_TYPE_GROUP = 1
|
||||
# 显式指定本地消息库位置(测试或非默认部署用)
|
||||
DB_PATH_ENV = "HEXI_GROUP_DAILY_HISTORY_DB"
|
||||
|
||||
_MESSAGE_SQL = (
|
||||
"SELECT id, user_id, message_id, raw_message, message, plain_text, time "
|
||||
f"FROM {MESSAGE_TABLE} "
|
||||
"WHERE group_id = ? AND time >= ? AND time <= ? "
|
||||
"ORDER BY time DESC, id DESC LIMIT ?"
|
||||
)
|
||||
_COUNT_SQL = (
|
||||
f"SELECT COUNT(*) FROM {MESSAGE_TABLE} "
|
||||
"WHERE group_id = ? AND time >= ? AND time <= ?"
|
||||
)
|
||||
# 同库 uninfo 三表:拿本群成员的 QQ 昵称与群名片
|
||||
_NAMES_SQL = (
|
||||
"SELECT u.user_id AS user_id, "
|
||||
"u.user_data AS user_data, "
|
||||
"s.member_data AS member_data "
|
||||
"FROM nonebot_plugin_uninfo_scenemodel AS sc "
|
||||
"JOIN nonebot_plugin_uninfo_sessionmodel AS s "
|
||||
"ON s.scene_persist_id = sc.id "
|
||||
"JOIN nonebot_plugin_uninfo_usermodel AS u "
|
||||
"ON u.id = s.user_persist_id "
|
||||
"WHERE sc.scene_id = ? AND sc.scene_type = ?"
|
||||
)
|
||||
|
||||
|
||||
@dataclass
|
||||
class LocalHistoryResult:
|
||||
"""本地历史读取结果;空 messages 表示需要回退到接口分页。"""
|
||||
|
||||
messages: list[dict[str, Any]] = field(default_factory=list)
|
||||
truncated: bool = False
|
||||
window_total: int | None = None
|
||||
names_resolved: int = 0
|
||||
duplicates: int = 0
|
||||
error: str | None = None
|
||||
|
||||
def __bool__(self) -> bool:
|
||||
return bool(self.messages)
|
||||
|
||||
|
||||
def _parse_sqlite_url(raw: Any) -> Path | None:
|
||||
"""从 SQLAlchemy 数据库配置里解析出 sqlite 文件路径。"""
|
||||
if raw is None:
|
||||
return None
|
||||
text = str(raw).strip()
|
||||
if not text.startswith("sqlite") or "memory" in text:
|
||||
return None
|
||||
# sqlite+aiosqlite:///D:/path/db.sqlite3 或 sqlite:///./rel/db.sqlite3
|
||||
_, _, tail = text.partition("///")
|
||||
tail = tail.strip()
|
||||
if not tail:
|
||||
return None
|
||||
try:
|
||||
return Path(tail).expanduser()
|
||||
except (OSError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def _driver_db_path() -> Path | None:
|
||||
"""读取 `SQLALCHEMY_DATABASE_URL`(环境变量优先,其次 .env 配置)。"""
|
||||
for key in ("SQLALCHEMY_DATABASE_URL", "sqlalchemy_database_url"):
|
||||
raw = os.environ.get(key)
|
||||
if raw:
|
||||
return _parse_sqlite_url(raw)
|
||||
try:
|
||||
from nonebot import get_driver
|
||||
|
||||
raw = getattr(get_driver().config, "sqlalchemy_database_url", None)
|
||||
except Exception:
|
||||
return None
|
||||
return _parse_sqlite_url(raw)
|
||||
|
||||
|
||||
def _localstore_db_path() -> Path | None:
|
||||
"""nonebot_plugin_orm 默认库:localstore 数据目录下的 db.sqlite3。"""
|
||||
try:
|
||||
from nonebot_plugin_localstore import get_data_dir
|
||||
|
||||
return get_data_dir("nonebot_plugin_orm") / "db.sqlite3"
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
def resolve_db_path() -> Path:
|
||||
"""解析本地消息库路径:环境变量 → ORM 配置 → localstore → 仓库默认位置。"""
|
||||
override = os.environ.get(DB_PATH_ENV)
|
||||
if override:
|
||||
return Path(override).expanduser()
|
||||
for candidate in (_driver_db_path(), _localstore_db_path()):
|
||||
if candidate is not None:
|
||||
return candidate
|
||||
# hexi/plugins/<plugin>/history_store.py -> hexi/
|
||||
return (
|
||||
Path(__file__).resolve().parents[2]
|
||||
/ "data"
|
||||
/ "nonebot_plugin_orm"
|
||||
/ "db.sqlite3"
|
||||
)
|
||||
|
||||
|
||||
def _open_readonly(path: Path) -> sqlite3.Connection:
|
||||
"""只读打开。路径可能含空格,必须用 as_uri 转义后的 file URI。"""
|
||||
con = sqlite3.connect(f"{path.resolve().as_uri()}?mode=ro", uri=True, timeout=3.0)
|
||||
con.row_factory = sqlite3.Row
|
||||
return con
|
||||
|
||||
|
||||
def _loads(raw: Any) -> dict[str, Any]:
|
||||
if not raw:
|
||||
return {}
|
||||
if isinstance(raw, dict):
|
||||
return raw
|
||||
try:
|
||||
data = json.loads(raw)
|
||||
except (TypeError, ValueError):
|
||||
return {}
|
||||
return data if isinstance(data, dict) else {}
|
||||
|
||||
|
||||
def _load_display_names(con: sqlite3.Connection, group_id: int) -> dict[str, tuple[str, str | None]]:
|
||||
"""取本群成员的 (昵称, 群名片);uninfo 表缺失时返回空表。"""
|
||||
names: dict[str, tuple[str, str | None]] = {}
|
||||
for row in con.execute(_NAMES_SQL, (str(group_id), SCENE_TYPE_GROUP)):
|
||||
user_data = _loads(row["user_data"])
|
||||
member_data = _loads(row["member_data"])
|
||||
nickname = str(user_data.get("name") or "")
|
||||
card = str(member_data.get("nick") or "") or None
|
||||
if nickname or card:
|
||||
names[str(row["user_id"])] = (nickname, card)
|
||||
return names
|
||||
|
||||
|
||||
def _parse_cq(cq: str) -> list[dict[str, Any]]:
|
||||
"""CQ 串转 OneBot 消息段(复用适配器自带解析,含反转义)。"""
|
||||
try:
|
||||
from nonebot.adapters.onebot.v11 import Message
|
||||
|
||||
return [{"type": seg.type, "data": dict(seg.data)} for seg in Message(cq)]
|
||||
except Exception as exc:
|
||||
logger.debug(f"本地历史 CQ 解析失败,降级为纯文本: {exc}")
|
||||
return [{"type": "text", "data": {"text": cq}}]
|
||||
|
||||
|
||||
def _build_raw_message(
|
||||
row: sqlite3.Row, names: dict[str, tuple[str, str | None]]
|
||||
) -> dict[str, Any] | None:
|
||||
"""组装成 adapter._convert_message 能直接消费的 OneBot 原始结构。"""
|
||||
message_id = row["message_id"]
|
||||
if message_id is None:
|
||||
return None
|
||||
user_id = str(row["user_id"])
|
||||
nickname, card = names.get(user_id, ("", None))
|
||||
cq = row["raw_message"] or row["message"] or row["plain_text"] or ""
|
||||
return {
|
||||
"message_id": message_id,
|
||||
"time": int(row["time"] or 0),
|
||||
"sender": {"user_id": user_id, "nickname": nickname, "card": card or ""},
|
||||
"message": _parse_cq(cq),
|
||||
}
|
||||
|
||||
|
||||
def _fetch_group_messages(
|
||||
group_id: str | int,
|
||||
start_ts: int,
|
||||
end_ts: int,
|
||||
limit: int,
|
||||
db_path: Path | str | None,
|
||||
) -> LocalHistoryResult:
|
||||
try:
|
||||
gid = int(group_id)
|
||||
except (TypeError, ValueError):
|
||||
return LocalHistoryResult(error=f"群号非法: {group_id!r}")
|
||||
try:
|
||||
row_limit = max(1, int(limit))
|
||||
except (TypeError, ValueError):
|
||||
row_limit = 1000
|
||||
|
||||
path = Path(db_path) if db_path is not None else resolve_db_path()
|
||||
if not path.exists():
|
||||
return LocalHistoryResult(error=f"本地消息库不存在: {path}")
|
||||
|
||||
start, end = int(start_ts), int(end_ts)
|
||||
con = _open_readonly(path)
|
||||
try:
|
||||
# 多取一条用于判定截断(多出来的那条正是窗口内最旧的消息)
|
||||
rows = con.execute(_MESSAGE_SQL, (gid, start, end, row_limit + 1)).fetchall()
|
||||
truncated = len(rows) > row_limit
|
||||
if truncated:
|
||||
rows = rows[:row_limit]
|
||||
window_total = None
|
||||
if truncated:
|
||||
try:
|
||||
window_total = int(con.execute(_COUNT_SQL, (gid, start, end)).fetchone()[0])
|
||||
except (sqlite3.Error, TypeError, ValueError):
|
||||
window_total = None
|
||||
try:
|
||||
names = _load_display_names(con, gid)
|
||||
except sqlite3.Error as exc:
|
||||
logger.debug(f"本地历史昵称补齐失败(忽略): {exc}")
|
||||
names = {}
|
||||
finally:
|
||||
con.close()
|
||||
|
||||
messages: list[dict[str, Any]] = []
|
||||
seen: set[str] = set()
|
||||
duplicates = 0
|
||||
names_resolved = 0
|
||||
for row in rows:
|
||||
raw = _build_raw_message(row, names)
|
||||
if raw is None:
|
||||
continue
|
||||
message_id = str(raw["message_id"])
|
||||
if not message_id:
|
||||
continue
|
||||
if message_id in seen:
|
||||
duplicates += 1
|
||||
continue
|
||||
seen.add(message_id)
|
||||
sender = raw["sender"]
|
||||
if sender["nickname"] or sender["card"]:
|
||||
names_resolved += 1
|
||||
messages.append(raw)
|
||||
# 查询按 (time, id) 倒序取最近 N 条,翻回时间升序
|
||||
messages.reverse()
|
||||
return LocalHistoryResult(
|
||||
messages=messages,
|
||||
truncated=truncated,
|
||||
window_total=window_total,
|
||||
names_resolved=names_resolved,
|
||||
duplicates=duplicates,
|
||||
)
|
||||
|
||||
|
||||
def fetch_group_messages(
|
||||
group_id: str | int,
|
||||
start_ts: int,
|
||||
end_ts: int,
|
||||
limit: int = 1000,
|
||||
db_path: Path | str | None = None,
|
||||
) -> LocalHistoryResult:
|
||||
"""读取 [start_ts, end_ts] 闭区间内的群消息(时间升序,最多 limit 条)。
|
||||
|
||||
同步函数,调用方用 `asyncio.to_thread` 包;失败只返回带 error 的空结果。
|
||||
"""
|
||||
try:
|
||||
return _fetch_group_messages(group_id, start_ts, end_ts, limit, db_path)
|
||||
except Exception as exc: # 兜底:本地库问题不该影响分析主链路
|
||||
logger.warning(f"本地历史读取异常: {exc}", exc_info=True)
|
||||
return LocalHistoryResult(error=str(exc))
|
||||
@@ -119,6 +119,7 @@ def register_bot_adapter(bot: Any, platform_id: str | None = None) -> OneBotAdap
|
||||
"platform_id": platform_id or "onebot",
|
||||
"bot_self_ids": bot_self_ids,
|
||||
"filter_bot_messages": config_manager.get_filter_bot_messages(),
|
||||
"use_local_history": config_manager.get_use_local_history(),
|
||||
},
|
||||
)
|
||||
bot_manager.register_adapter(adapter, platform_id or "onebot")
|
||||
|
||||
Reference in New Issue
Block a user