Files

436 lines
18 KiB
Python
Raw Permalink Normal View History

"""NoneBot OneBot V11 适配器:负责拉取群历史消息、群组/成员信息与发送。"""
from __future__ import annotations
import asyncio
from datetime import datetime, timedelta
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,
MessageContentType,
UnifiedMessage,
)
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 的最小适配器。"""
platform_name = "onebot"
USER_AVATAR_TEMPLATE = "https://q1.qlogo.cn/g?b=qq&nk={user_id}&s=640"
def __init__(self, bot: Any, config: dict | None = None):
self.bot = bot
self.config = config or {}
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(
self,
group_id: str,
days: int = 1,
max_count: int = 1000,
before_id: str | None = None,
since_ts: int | None = None,
) -> list[UnifiedMessage]:
try:
chunk_size = 100
all_raw = []
seen_raw_ids: set[str] = set()
if since_ts and since_ts > 0:
start_ts = since_ts
else:
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] = {
"group_id": int(group_id),
"count": fetch_count,
}
if current_anchor:
params["message_seq"] = current_anchor
params["reverseOrder"] = True
result = None
for attempt in range(1, 4):
try:
result = await self.bot.call_api(
"get_group_msg_history", **params
)
break
except Exception as exc:
if attempt < 3:
logger.warning(
f"OneBot 分页拉取失败(第{attempt}次): {exc}"
)
await asyncio.sleep(attempt)
else:
logger.warning(
f"OneBot 分页拉取重试耗尽: {exc}"
)
if not result or "messages" not in result:
break
messages = result.get("messages", [])
if not messages:
# 空页可能意味着当前锚点字段不被后端识别(如 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", ""))
if not msg_id or msg_id in seen_raw_ids:
continue
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
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] = []
seen: set[str] = set()
for raw in all_raw:
mid = str(raw.get("message_id", ""))
if not mid or mid in seen:
continue
u = self._convert_message(raw, group_id)
if u:
unified.append(u)
seen.add(mid)
unified.sort(key=lambda m: m.timestamp)
return unified
except Exception as e:
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", {})
chain = raw.get("message", [])
if isinstance(chain, str):
chain = [{"type": "text", "data": {"text": chain}}]
contents: list[MessageContent] = []
text_parts: list[str] = []
for seg in chain:
seg_t = seg.get("type", "")
seg_d = seg.get("data", {})
if seg_t == "text":
text = seg_d.get("text", "")
text_parts.append(text)
contents.append(MessageContent(type=MessageContentType.TEXT, text=text))
elif seg_t == "image":
sub_type = seg_d.get("subType", seg_d.get("sub_type"))
try:
is_sticker = int(sub_type) == 1
except (TypeError, ValueError):
is_sticker = False
raw_data: dict[str, Any] = {"summary": seg_d.get("summary", "")}
if sub_type is not None:
try:
raw_data["sub_type"] = int(sub_type)
except (TypeError, ValueError):
pass
contents.append(
MessageContent(
type=MessageContentType.EMOJI if is_sticker else MessageContentType.IMAGE,
url=seg_d.get("url", seg_d.get("file", "")),
raw_data=raw_data,
)
)
elif seg_t == "at":
contents.append(
MessageContent(type=MessageContentType.AT, at_user_id=str(seg_d.get("qq", "")))
)
elif seg_t in ("face", "mface", "bface", "sface"):
contents.append(
MessageContent(
type=MessageContentType.EMOJI,
emoji_id=str(seg_d.get("id", "")),
raw_data={"face_type": seg_t},
)
)
elif seg_t == "reply":
contents.append(
MessageContent(type=MessageContentType.REPLY, raw_data={"reply_id": seg_d.get("id", "")})
)
else:
contents.append(MessageContent(type=MessageContentType.UNKNOWN, raw_data=seg))
reply_to = None
for c in contents:
if c.type == MessageContentType.REPLY and c.raw_data:
reply_to = str(c.raw_data.get("reply_id", ""))
break
return UnifiedMessage(
message_id=str(raw.get("message_id", "")),
sender_id=str(sender.get("user_id", "")),
sender_name=sender.get("nickname", ""),
sender_card=sender.get("card", "") or None,
group_id=group_id,
text_content="".join(text_parts),
contents=tuple(contents),
timestamp=raw.get("time", 0),
platform="onebot",
reply_to_id=reply_to,
)
except Exception as e:
logger.debug(f"OneBot _convert_message error: {e}")
return None
# —— 群/成员信息 ——
async def get_group_info(self, group_id: str) -> UnifiedGroup | None:
try:
gid = int(group_id)
data = await self.bot.call_api("get_group_info", group_id=gid)
if isinstance(data, list):
data = next((d for d in data if str(d.get("group_id")) == str(gid)), data[0] if data else {})
return UnifiedGroup(
group_id=str(gid),
group_name=data.get("group_name", ""),
member_count=int(data.get("member_count", 0) or 0),
platform="onebot",
)
except Exception as e:
logger.debug(f"get_group_info({group_id}) 失败: {e}")
return None
async def get_user_avatar_url(self, user_id: str) -> str | None:
return self.USER_AVATAR_TEMPLATE.format(user_id=user_id)
async def get_member_info(self, group_id: str, user_id: str) -> UnifiedMember | None:
try:
data = await self.bot.call_api(
"get_group_member_info",
group_id=int(group_id),
user_id=int(user_id),
)
return UnifiedMember(
user_id=str(user_id),
nickname=data.get("nickname", ""),
card=data.get("card", "") or None,
role=str(data.get("role", "member")),
)
except Exception as e:
logger.debug(f"get_member_info({group_id},{user_id}) 失败: {e}")
return None
async def is_group_muted(self, group_id: str) -> bool:
return False
def get_platform_name(self) -> str:
return self.platform_name
# —— 发送(在手动/自动报告中可用;手动指令会走 SAA 发送) ——
async def send_text(self, group_id: str, text: str) -> bool:
try:
await self.bot.call_api("send_group_msg", group_id=int(group_id), message=str(text))
return True
except Exception as e:
logger.warning(f"send_text 失败: {e}")
return False
async def send_image(self, group_id: str, image_url: str, caption: str = "") -> bool:
try:
if image_url.startswith("base64://"):
b64 = image_url.split("base64://", 1)[1]
from nonebot.adapters.onebot.v11 import MessageSegment
msg = MessageSegment.image(file=f"base64://{b64}")
else:
from nonebot.adapters.onebot.v11 import MessageSegment
msg = MessageSegment.image(file=image_url)
await self.bot.call_api("send_group_msg", group_id=int(group_id), message=msg)
return True
except Exception as e:
logger.warning(f"send_image 失败: {e}")
return False
async def send_file(self, group_id: str, file_path: str, caption: str = "") -> bool:
return False
async def send_forward_msg(self, group_id: str, nodes: list[dict]) -> bool:
"""发送 OneBot v11 合并转发消息 (send_forward_msg)。"""
try:
from nonebot.adapters.onebot.v11 import MessageSegment
normalized = []
for n in nodes:
data = dict(n.get("data", {}))
content = data.get("content")
if isinstance(content, str):
content = [MessageSegment.text(content)]
elif content is None:
content = []
data["content"] = content
normalized.append({"type": "node", "data": data})
await self.bot.call_api(
"send_forward_msg", group_id=int(group_id), messages=normalized
)
return True
except Exception as e:
logger.warning(f"send_forward_msg 失败: {e}")
return False
async def set_reaction(self, group_id: str, message_id: str, emoji: str | int, is_add: bool = True) -> bool:
return False
async def prepare_group_member_cache(self, group_id: str):
return True, None