308 lines
12 KiB
Python
308 lines
12 KiB
Python
"""NoneBot OneBot V11 适配器:负责拉取群历史消息、群组/成员信息与发送。"""
|
|||
|
|
|
||
|
|
from __future__ import annotations
|
||
|
|
|
||
|
|
import asyncio
|
||
|
|
import base64
|
||
|
|
import time
|
||
|
|
from datetime import datetime, timedelta
|
||
|
|
from pathlib import Path
|
||
|
|
from typing import Any
|
||
|
|
|
||
|
|
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
|
||
|
|
|
||
|
|
|
||
|
|
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))
|
||
|
|
|
||
|
|
# —— 消息拉取 ——
|
||
|
|
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()
|
||
|
|
)
|
||
|
|
|
||
|
|
current_anchor = before_id
|
||
|
|
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:
|
||
|
|
break
|
||
|
|
|
||
|
|
first = messages[0]
|
||
|
|
last = messages[-1]
|
||
|
|
earliest = first if first.get("time", 0) <= last.get("time", 0) else last
|
||
|
|
chunk_earliest_ts = earliest.get("time", 0)
|
||
|
|
|
||
|
|
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)
|
||
|
|
|
||
|
|
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:
|
||
|
|
break
|
||
|
|
if current_anchor and str(new_anchor) == str(current_anchor):
|
||
|
|
break
|
||
|
|
current_anchor = 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 []
|
||
|
|
|
||
|
|
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
|