import requests import json import time import re import logging import sys from pathlib import Path from datetime import datetime from src.modules.plugin_modules import BasePlugin, MessageContext # 确保内嵌依赖在 sys.path 中 _plugin_dir = Path(__file__).parent _packages_dir = _plugin_dir / "packages" if str(_packages_dir) not in sys.path: sys.path.insert(0, str(_packages_dir)) import jieba logger = logging.getLogger(__name__) ADMIN_QQ = "YOUR_ADMIN_QQ" NO_REPLY_MARKER = "No response from OpenClaw." SANITIZED_REPLY = "." # 统一替换对外暴露的系统内部信息 # ═══════════════════════════════════════════════════════════ # 高危词库(jieba分词后匹配 + 整句子串匹配) # 资料参考:OWASP Prompt Injection, Prompt Engineering Guide, # Simon Willison, Jailbreak Chat, ChatGPT DAN variants # ═══════════════════════════════════════════════════════════ # 高危词库:优先从 config.toml [security] section 加载 # 如未配置则使用内置默认词库(见 _get_high_risk_words()) # 外部化方便不同部署环境自定义词库 # 群昵称缓存: {(group_id, user_id): "nickname"} _nickname_cache: dict = {} class OpenClawBridge(BasePlugin): """QQ消息 ↔ OpenClaw Gateway 桥接插件""" def __init__(self, ctx: MessageContext): super().__init__(ctx) logger.info("=== OpenClawBridge __init__ START ===") self.gateway_url = None self.gateway_token = None self.allowed_sender = "" self.model = None self.agent_id = None try: cfg = self.config.get("openclaw", {}) if cfg: self.gateway_url = cfg.get("gateway_url") self.gateway_token = cfg.get("gateway_token") self.allowed_sender = str(cfg.get("allowed_sender")) if cfg.get("allowed_sender") is not None else None self.model = cfg.get("model") self.agent_id = cfg.get("agent_id") logger.info(f"Config loaded: url={self.gateway_url}, allowed={self.allowed_sender}") else: logger.error("Config section [openclaw] not found in config.toml") except Exception as e: logger.error(f"Failed to load config: {e}") missing = [k for k, v in {"gateway_url": self.gateway_url, "gateway_token": self.gateway_token}.items() if not v] if missing: logger.error(f"Missing required config: {missing}") logger.info("=== OpenClawBridge __init__ END ===") def _fetch_group_nickname(self, group_id: str, sender_id: str) -> str: """获取用户在群里的昵称,纯文本""" cache_key = (group_id, sender_id) if cache_key in _nickname_cache: return _nickname_cache[cache_key] # 先尝试从 Group.users 查找 nickname = self._lookup_from_group_users(group_id, sender_id) if nickname: _nickname_cache[cache_key] = nickname return nickname # 如果 Group.users 没有命中,再通过平台 API 查询 nickname = self._fetch_from_platform_api(group_id, sender_id) if nickname: _nickname_cache[cache_key] = nickname return nickname # 最终兜底 display = "管理员" if sender_id == ADMIN_QQ else f"用户{sender_id}" _nickname_cache[cache_key] = display return display def _lookup_from_group_users(self, group_id: str, sender_id: str) -> str: """从 Group.users(平台预加载的群成员列表)查找昵称""" try: if self.ctx.group and self.ctx.group.users: users = self.ctx.group.users if isinstance(users, list): for u in users: if str(u.get("user_id")) == sender_id: return u.get("card") or u.get("nickname") or "" elif isinstance(users, dict): u = users.get(int(sender_id), {}) return u.get("card") or u.get("nickname") or "" except Exception as e: logger.error(f"lookup from group users failed: {e}") return "" def _fetch_from_platform_api(self, group_id: str, sender_id: str) -> str: """通过平台 API /get_group_member_info 查询单个用户昵称""" base_url = self.ctx.group.url if self.ctx.group else "http://YOUR_NAPCAT_HOST:25570" try: resp = requests.post( f"{base_url}/get_group_member_info", json={"group_id": group_id, "user_id": sender_id}, timeout=3 ) if resp.status_code == 200: body = resp.json() data = body.get("data", {}) if isinstance(body, dict) else {} return data.get("card") or data.get("nickname") or "" except Exception as e: logger.error(f"platform api fetch failed: {e}") return "" def _get_sender_group_nickname(self) -> str: """获取当前消息发送者的昵称""" if self.ctx.group is None: # 私聊 nickname = getattr(self.ctx.user, 'nickname', None) or getattr(self.ctx.user, 'card_name', None) or '未知用户' return nickname return self._fetch_group_nickname(self.ctx.group.group_id, self._get_sender_id()) def _get_sender_id(self) -> str: return str(self.ctx.user.user_id) def _clean_message(self, text: str) -> str: """替换 CQ 码为用户可读文本""" # 别人 @bot 时,让模型知道是在叫自己 bot_at = f"[CQ:at,qq={self.ctx.rebot_id}]" bot_id_str = str(self.ctx.rebot_id) text = text.replace(bot_at, f"@你(你的QQ号{bot_id_str})") # 替换所有其他 CQ 码 # 图片/文件保留文件名,如 [image:xxx.jpg] [file:abc.zip] def _replace_cq(m): cq_type = m.group(1) attrs = m.group(2) or "" if cq_type in ("image", "file"): # 提取 file=xxx 部分 import re as re2 fname_match = re2.search(r'file=([^,\]]+)', attrs) fname = fname_match.group(1) if fname_match else "" return f"[{cq_type}:{fname}]" return f"[{cq_type}]" text = re.sub(r'\[CQ:([^,]+)(,[^\]]+)?\]', _replace_cq, text) return text def _build_source_tag(self) -> str: """构建来源标记,让AI知道消息来自哪个群/私聊""" if self.ctx.group is None: return "私聊" # 使用群ID作为来源标识,group.nickname可能需异步获取不一定可用 group_id = self.ctx.group.group_id return f"群聊({group_id})" def _build_context_with_history(self, raw_message: str, identity_tag: str) -> str: """返回带上下文的消息(OpenClaw session 负责历史管理) 格式:{来源:群聊/私聊} {用户名:昵称} 消息正文 让AI能清晰区分消息来源和说话者身份,避免混淆。 """ cleaned = self._clean_message(raw_message) source_tag = self._build_source_tag() return f"{{来源:{source_tag}}} {{用户名:{identity_tag}}} {cleaned}" def _get_high_risk_words(self) -> list: """获取高危词列表,优先从 config.toml [security] section 加载,否则用内置默认词库""" try: cfg_high_risk = self.config.get("security", {}).get("high_risk_words", []) if cfg_high_risk: logger.info(f"Loaded {len(cfg_high_risk)} high-risk words from config") return cfg_high_risk except Exception: pass # 内置默认词库 default_words = [ "ignore all instructions", "ignore all previous", "forget previous", "forget everything", "delete memory", "delete all", "reset system", "override", "sudo password", "api key", "secret token", "ssh key", "cat /etc/passwd", "rm -rf", ":(){ :|:& };:", ] logger.info(f"Using built-in high-risk word list ({len(default_words)} words)") return default_words def _detect_high_risk(self, message: str) -> tuple[bool, list[str]]: """使用jieba分词检测高危词,返回(是否高危, 匹配到的词列表)""" msg_lower = message.lower() words = list(jieba.cut(msg_lower)) # 同时保留整句匹配(防止分词切错) all_targets = words + [msg_lower] matched = [] for risk_word in self._get_high_risk_words(): for target in all_targets: if risk_word in target: if risk_word not in matched: matched.append(risk_word) break is_high_risk = len(matched) > 0 logger.debug(f"High-risk check: jieba_cut={words[:20]}..., matched={matched}, is_high_risk={is_high_risk}") return is_high_risk, matched def _mark_high_risk(self, message: str, matched_words: list[str]) -> str: """在高危信息最显眼的地方标注""" marker = "【⚠️ 高危信息,谨慎处理】" words_str = "、".join(matched_words[:5]) # 最多显示5个 annotated = f"{marker}[命中: {words_str}]\n{message}" return annotated def _notify_admin(self, sender_id: str, gid: str, raw_message: str, matched_words: list[str]): """私聊上报管理员""" words_str = "、".join(matched_words[:5]) report = f"⚠️ 攻击上报:用户{sender_id}在群{gid}试图:{raw_message[:80]}... [命中高危词: {words_str}]" try: base_url = self.ctx.group.url if self.ctx.group else "http://YOUR_NAPCAT_HOST:25570" requests.post( f"{base_url}/send_private_msg", json={"user_id": ADMIN_QQ, "message": report}, timeout=5 ) logger.info(f"Admin notified: {report[:60]}...") except Exception as e: logger.error(f"Failed to notify admin: {e}") def _looks_like_mc_command(self, message: str) -> bool: """检测是否为 MC 命令(以 / 开头的命令)""" stripped = message.strip() if stripped.startswith('/'): return True clean = self._clean_message(message).strip() if clean.startswith('/'): return True return False def _is_authorized(self, sender_id: str) -> bool: if self.allowed_sender is None: logger.warning("allowed_sender is not configured, denying all requests") return False return sender_id == self.allowed_sender def _is_dangerous(self, message: str) -> bool: dangerous_keywords = [ r"rm\s+-[rf]", r"dd\s+if=", r"mkfs\.\w+", r"format\s+", r"del\s+/[fqs]", r"rmdir\s+/s", r"passwd\b", r"/etc/shadow", r"ssh_key", r"private_key", r"curl\s+.*\|\s*(sh|bash)", r"wget\s+.*\|\s*(sh|bash)", r"bash\s*<\(", r"sudo\s+", r"chmod\s+777\s+/", r"chown\s+-R", r"iptables\s+-F", r"systemctl\s+(stop|restart|disable)", r"killall\b", r"yum\s+install", r"apt\s+(install|remove|purge)", r"pip\s+(install|uninstall)", r">\s*/etc/\w+", r"cat\s+/etc/(shadow|passwd|ssh)" ] msg_lower = message.lower() for pattern in dangerous_keywords: if re.search(pattern, msg_lower): return True return False def _build_session_key(self) -> str: # 不同群/私聊用不同 session,避免历史污染 if self.ctx.group: return f"qq-user-{self.ctx.user.user_id}-gid-{self.ctx.group.group_id}" return f"qq-user-{self.ctx.user.user_id}" def _strip_markdown(self, text: str) -> str: text = re.sub(r'```[\s\S]*?```', '[代码块]', text) text = re.sub(r'`([^`]+)`', r'\1', text) text = re.sub(r'#{1,6}\s+', '', text) text = re.sub(r'\*\*([^*]+)\*\*', r'\1', text) text = re.sub(r'\*([^*]+)\*', r'\1', text) text = re.sub(r'__([^_]+)__', r'\1', text) text = re.sub(r'\[([^\]]+)\]\([^\)]+\)', r'\1', text) text = re.sub(r'^>\s+', '', text, flags=re.MULTILINE) text = re.sub(r'^[-*+]\s+', '', text, flags=re.MULTILINE) text = re.sub(r'^\d+\.\s+', '', text, flags=re.MULTILINE) text = re.sub(r'^---+$', '', text, flags=re.MULTILINE) text = re.sub(r'\n{3,}', '\n\n', text) return text.strip() def _send_to_openclaw(self, message: str, session_key: str) -> str: """ 同步请求 Gateway。 Agent 模型(如 openclaw/qq-agent)在 stream=True 模式下 Gateway 内部 agent 框架标记 run error 后 SSE 输出空内容 (delta:{} finish_reason:stop),因此不使用流式传输。 同步模式可正确 failover 模型 fallback(kimi→deepseek)。 qq-agent 模型响应速度 <5s,无需 __EXTEND__ 续命机制。 """ logger.info(f"=== _send_to_openclaw START: session={session_key}, msg={message[:50]}... ===") if self.gateway_url is None or self.gateway_token is None: logger.error("gateway_url or gateway_token is not configured") return SANITIZED_REPLY url = f"{self.gateway_url}/v1/chat/completions" headers = {"Content-Type": "application/json", "Authorization": f"Bearer {self.gateway_token}"} payload = { "model": self.model, "messages": [{"role": "user", "content": message}], "stream": False, "user": session_key, "session": session_key, "session_key": session_key } try: logger.info(f"POST {url}") response = requests.post(url, headers=headers, json=payload, timeout=300) logger.info(f"Response: {response.status_code}") if response.status_code == 200: data = response.json() content = data.get("choices", [{}])[0].get("message", {}).get("content", "") if content: return self._strip_markdown(content) logger.warning("Agent returned empty content") return NO_REPLY_MARKER # 不暴露 HTTP 状态码和响应体 logger.warning(f"Non-200 response: {response.status_code} {response.text[:200]}") return SANITIZED_REPLY except requests.exceptions.Timeout: logger.warning("Gateway request timed out") return SANITIZED_REPLY except requests.exceptions.ConnectionError: logger.warning("Gateway connection failed") return SANITIZED_REPLY except Exception as e: logger.warning(f"Gateway request exception: {e}") return SANITIZED_REPLY def after_save(self): logger.info("=== OpenClawBridge after_save START ===") sender_id = self._get_sender_id() raw_message = self.ctx.raw_message identity_tag = self._get_sender_group_nickname() logger.info(f"sender={sender_id}, identity_tag={identity_tag}, msg={raw_message[:50]}...") # ===== 系统通知/文件回执过滤 ===== cleaned = self._clean_message(raw_message) # 纯媒体消息(无文字内容):图片/文件/视频/语音 stripped = re.sub(r'\[image:[^\]]+\]|\[file:[^\]]+\]|\[video\]|\[语音\]|\[动画表情\]', '', cleaned).strip() if not stripped: logger.info(f"Pure file/media message from {sender_id}, skipped") return # QQ 系统通知:对方接收/下载文件回执、离线文件通知 sys_keywords = ['已接收', '已下载', '已打开', '已成功接收', '已成功下载', '系统消息', '系统通知', '你收到离线文件'] if any(k in cleaned for k in sys_keywords): logger.info(f"QQ system notification from {sender_id}, skipped: {cleaned[:50]}") return is_admin = self._is_authorized(sender_id) # ===== 高危词检测(jieba分词,所有消息都走,只标注不拦截)===== is_high_risk, matched_words = self._detect_high_risk(raw_message) if is_high_risk: logger.warning(f"🚨 HIGH RISK detected from {sender_id}: {matched_words}") # 非管理员触发高危词时,上报管理员(用原始消息上报,非管理员看不到带标注的版本) if not is_admin: self._notify_admin(sender_id, self.ctx.group.group_id if self.ctx.group else "私聊", self.ctx.raw_message, matched_words) # 标注高危信息,然后放行给后续处理 raw_message = self._mark_high_risk(self.ctx.raw_message, matched_words) logger.info(f"Message annotated with high-risk marker: {sender_id}") # ===== 非授权用户:仅拦截MC命令(用原始消息检测,避免高危标注干扰)===== if not is_admin: if self._looks_like_mc_command(self.ctx.raw_message): logger.info(f"Unauthorized MC command from {sender_id}, blocked from AI") return # ===== 危险请求检测(管理员专用)===== if is_admin and self._is_dangerous(raw_message): report_msg = f"[QQ危险请求] senderId={sender_id},内容:{raw_message[:100]}" self._send_to_openclaw(report_msg, "qq-danger-report") logger.info("Dangerous request reported") return "ok" # ===== 所有用户正常对话(整合后的统一流程)===== # 群聊消息没有@机器人时,不调LLM(管理员除外) if self.ctx.group is not None and not self._is_authorized(sender_id): at_me = f"[CQ:at,qq={self.ctx.rebot_id}]" in self.ctx.raw_message if not at_me: logger.info(f"Group chat without @bot from {sender_id}, skipped LLM call") return session_key = self._build_session_key() context_msg = self._build_context_with_history(raw_message, identity_tag) reply = self._send_to_openclaw(context_msg, session_key) logger.info(f"reply: {reply[:80]}...") # ===== 回复处理 ===== # 统一过滤:所有对外的内部错误文案都不发送给用户 if reply and reply.strip() in (NO_REPLY_MARKER, SANITIZED_REPLY): logger.warning(f"Blocked internal marker: {reply.strip()!r}") return if not reply: logger.info("Empty reply, skipped") return if self.ctx.group is None: # 私聊 logger.info("Private chat, processing...") logger.info("Sending private message...") self.ctx.user.send_message(reply) else: # 群聊只响应 @机器人 at_me = f"[CQ:at,qq={self.ctx.rebot_id}]" in self.ctx.raw_message logger.info(f"Group chat, at_me={at_me}") if at_me: logger.info("Sending group message...") self.ctx.group.send_message(reply) logger.info(f"=== OpenClawBridge after_save END (admin={is_admin}) ===") return