Browse Source

内容: 小红书数据机器人双向监听(飞书长连接WebSocket版)

Sisyphus Agent 1 month ago
parent
commit
940b3fa430
1 changed files with 304 additions and 0 deletions
  1. 304 0
      运营文案/xhs_bot_listener.py

+ 304 - 0
运营文案/xhs_bot_listener.py

@@ -0,0 +1,304 @@
+# -*- coding: utf-8 -*-
+"""
+小红书数据机器人(飞书长连接双向版)
+=====================================
+监听飞书群的 @机器人 消息,按指令采集小红书数据并回复。
+
+指令:
+    看数据 / 日报 / 数据   → 通过 CDP 采集最新小红书数据,推送报告卡片
+    help / 帮助           → 显示帮助
+
+技术栈:
+    - lark_oapi.ws.Client(飞书官方长连接 WebSocket,免公网回调)
+    - cdp_publish.XiaohongshuPublisher(复用已登录 Chrome,端口 9222)
+    - xhs_daily_report.collect_overview / collect_notes / screenshot_current
+    - xhs_feishu.send_daily_report
+
+用法:
+    python xhs_bot_listener.py
+    python xhs_bot_listener.py --test   # 只跑一次连接测试,不常驻
+
+先决条件:
+    1. 本机 Chrome 已以 --remote-debugging-port=9222 启动且已登录小红书
+    2. 运营文案/xhs_daily_config.json 已填入飞书 app_id/secret/群ID
+"""
+import io
+import json
+import os
+import re
+import sys
+import threading
+import time
+
+sys.stdout = io.TextIOWrapper(sys.stdout.buffer, encoding='utf-8', errors='replace')
+sys.stderr = sys.stdout
+
+import requests
+
+SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__))
+OPERATION_DIR = SCRIPT_DIR  # 本脚本位于 运营文案/ 目录,配置文件也在这
+CFC_ROOT = os.path.dirname(OPERATION_DIR)  # C:\code\cfc
+
+CONFIG_PATH = os.path.join(OPERATION_DIR, "xhs_daily_config.json")
+SCREENSHOT_DIR = os.path.join(OPERATION_DIR, "小红书发布", "_过程文件", "截图")
+
+# 复用 xhs_daily_report 与 xhs_feishu 的路径注入
+XHS_DAILY_DIR = os.path.join(OPERATION_DIR, "小红书发布", "_过程脚本")
+for _p in (OPERATION_DIR, XHS_DAILY_DIR, r"C:\code\XiaohongshuSkills\scripts"):
+    if _p not in sys.path:
+        sys.path.insert(0, _p)
+
+FEISHU_TOKEN_URL = "https://open.feishu.cn/open-apis/auth/v3/tenant_access_token/internal"
+FEISHU_REPLY_URL = "https://open.feishu.cn/open-apis/im/v1/messages/{msg_id}/reply"
+FEISHU_MSG_URL = "https://open.feishu.cn/open-apis/im/v1/messages"
+
+
+def _log(msg):
+    print(f"[{time.strftime('%H:%M:%S')}] {msg}", flush=True)
+
+
+def _load_config() -> dict:
+    with open(CONFIG_PATH, "r", encoding="utf-8") as f:
+        cfg = json.load(f)
+    for k in ("feishu_app_id", "feishu_app_secret", "feishu_chat_id"):
+        if not cfg.get(k):
+            raise ValueError(f"配置缺少必填字段: {k}")
+    return cfg
+
+
+class FeishuChat:
+    """封装飞书 token、回复、发消息。"""
+
+    def __init__(self, cfg):
+        self.cfg = cfg
+        self._token = None
+
+    def token(self):
+        if self._token:
+            return self._token
+        r = requests.post(
+            FEISHU_TOKEN_URL,
+            json={"app_id": self.cfg["feishu_app_id"], "app_secret": self.cfg["feishu_app_secret"]},
+            timeout=15,
+        )
+        r.raise_for_status()
+        d = r.json()
+        if d.get("code") != 0:
+            raise RuntimeError(f"获取飞书 token 失败: {d.get('msg')}")
+        self._token = d["tenant_access_token"]
+        return self._token
+
+    def reply_text(self, message_id: str, text: str):
+        """回复某一 message_id(线程内回复)。"""
+        url = FEISHU_REPLY_URL.format(msg_id=message_id)
+        body = {
+            "receive_id_type": "open_id",
+            "msg_type": "text",
+            "content": json.dumps({"text": text}, ensure_ascii=False),
+        }
+        r = requests.post(
+            url,
+            headers={"Authorization": f"Bearer {self.token()}"},
+            json=body,
+            timeout=20,
+        )
+        r.raise_for_status()
+        d = r.json()
+        if d.get("code") != 0:
+            _log(f"回复失败: {d.get('msg')} (message_id={message_id[:20]})")
+        return d
+
+    def send_text(self, text: str):
+        rtype = self.cfg.get("receive_id_type", "chat_id")
+        url = f"{FEISHU_MSG_URL}?receive_id_type={rtype}"
+        body = {
+            "receive_id": self.cfg["feishu_chat_id"],
+            "msg_type": "text",
+            "content": json.dumps({"text": text}, ensure_ascii=False),
+        }
+        r = requests.post(url, headers={"Authorization": f"Bearer {self.token()}"}, json=body, timeout=20)
+        r.raise_for_status()
+        d = r.json()
+        if d.get("code") != 0:
+            raise RuntimeError(f"发送文本失败: {d.get('msg')}")
+        return d
+
+
+def _mention_text(msg) -> str:
+    """还原消息中提到机器人之后的纯文本指令。
+
+    EventMessage.content 是 JSON 字符串,解析出 text 字段后
+    剥离 <at> 标签,返回可读文本。
+    """
+    if not (msg and msg.content):
+        return ""
+    content = msg.content
+    text = ""
+    try:
+        c = json.loads(content)
+        text = c.get("text", "") or ""
+    except Exception:
+        text = content
+    # 剥离 at 标签与首尾空白
+    text = re.sub(r"<at\b[^>]*>.*?</at>", "", text)
+    text = re.sub(r"@\S+\s*", "", text, flags=re.UNICODE)
+    return text.strip()
+
+
+def _is_bot_mentioned(msg) -> bool:
+    """判断消息是否 @了机器人。
+
+    EventMessage.mentions 是 List[MentionEvent],其中 mentioned_type='app'
+    表示 @了应用(机器人)。群聊里只有 @机器人 的消息才需要响应。
+    """
+    if not (msg and msg.mentions):
+        return False
+    for mt in msg.mentions:
+        if getattr(mt, "mentioned_type", "") == "app":
+            return True
+    return False  # 只要 mention 了,就看作对该 bot 的指令入口(群内一般是 @机器人)
+
+
+def _handle_data_command(message_id: str, feishu: FeishuChat):
+    """执行『看数据』指令:采集最新数据并推送报告卡片。"""
+    _log("[指令] 看数据 → 开始采集")
+    publisher = None
+    try:
+        # 复用 xhs_daily_report 的采集逻辑
+        from xhs_daily_report import (
+            XiaohongshuPublisher,
+            collect_overview,
+            collect_notes,
+            screenshot_current,
+            CREATOR_HOME,
+            NOTES_MANAGE_URL,
+        )
+
+        cfg = _load_config()
+        date = time.strftime("%Y-%m-%d")
+
+        publisher = XiaohongshuPublisher()
+        publisher.connect(reuse_existing_tab=True)
+        if not publisher.check_login():
+            raise RuntimeError("未登录小红书,请先运行 --login")
+
+        _log("[采集] 账号总览...")
+        publisher._navigate(CREATOR_HOME)
+        time.sleep(4)
+        overview = collect_overview(publisher)
+
+        _log("[采集] 笔记列表...")
+        publisher._navigate(NOTES_MANAGE_URL)
+        time.sleep(4)
+        notes = collect_notes(publisher, 10)
+
+        shot_path = os.path.join(SCREENSHOT_DIR, f"xhs_bot_{int(time.time())}.jpg")
+        try:
+            screenshot_current(publisher, shot_path)
+        except Exception as e:
+            _log(f"[WARN] 截图失败: {e}")
+            shot_path = None
+        publisher.disconnect()
+        publisher = None
+
+        _log("[推送] 发送日报卡片...")
+        try:
+            sys.path.insert(0, OPERATION_DIR)
+            from xhs_feishu import send_daily_report
+            send_daily_report(CONFIG_PATH, date, overview, notes, screenshot_path=shot_path)
+            _log("[完成] 日报已推送")
+        except Exception as e:
+            _log(f"[推送失败] {e}")
+            # 退化为文本回复
+            summary = (f"粉丝 {overview.get('fans_count', 0):,} | "
+                       f"阅读 {overview.get('total_read', 0):,} | "
+                       f"互动 {overview.get('total_interact', 0):,} | "
+                       f"笔记 {len(notes)} 篇")
+            feishu.reply_text(message_id, f"数据已采集({date}):{summary}\n卡片推送失败,详见日志。")
+            return
+
+        feishu.reply_text(message_id, f"✅ 日报已生成并推送到群({date})")
+    except Exception as e:
+        _log(f"[指令失败] {e}")
+        feishu.reply_text(message_id, f"❌ 采集失败:{e}")
+    finally:
+        try:
+            publisher.disconnect()
+        except Exception:
+            pass
+
+
+HELP_TEXT = (
+    "📌 小红书数据机器人可用指令:\n"
+    "· 看数据 / 日报 / 数据 → 采集最新小红书数据并推送报告\n"
+    "· help / 帮助 → 显示本帮助\n\n"
+    "请 @机器人 后发送指令。"
+)
+
+
+def on_message_receive(data):
+    """飞书消息接收事件回调(在 SDK 线程内执行,耗时操作放后台线程)。"""
+    try:
+        ev = data.event
+        msg = ev.message
+        if not msg:
+            return
+        chat_id = msg.chat_id
+        chat_type = getattr(msg, "chat_type", "")
+        msg_type = msg.message_type
+        _log(f"[收到消息] chat={chat_id} type={msg_type} chat_type={chat_type} mentions={len(msg.mentions or [])}")
+
+        # 仅响应 @机器人 的消息(群聊),p2p 私聊单独处理
+        if chat_type not in ("group", "p2p"):
+            return
+        if chat_type == "group" and not _is_bot_mentioned(msg):
+            _log("  未@机器人,忽略")
+            return
+
+        # 后台线程执行指令,避免阻塞事件循环
+        def _worker():
+            try:
+                feishu = FeishuChat(_load_config())
+                cmd = _mention_text(msg).lower()
+                # 去掉可能的 at 前缀后解析指令
+                if any(k in cmd for k in ("看数据", "日报", "数据", "stats", "status")):
+                    _handle_data_command(msg.message_id, feishu)
+                elif any(k in cmd for k in ("help", "帮助", "?")):
+                    feishu.reply_text(msg.message_id, HELP_TEXT)
+                elif cmd:
+                    feishu.reply_text(msg.message_id, f"未识别指令: {cmd}\n发送 help 查看可用命令。")
+                else:
+                    feishu.reply_text(msg.message_id, "你好!@我并发送「看数据」即可查看最新小红书数据。")
+            except Exception as e:
+                _log(f"[worker] {e}")
+
+        threading.Thread(target=_worker, daemon=True).start()
+    except Exception as e:
+        _log(f"[on_message_receive] {e}")
+
+
+def main():
+    parser = __import__("argparse").ArgumentParser(description="小红书飞书数据机器人(长连接)")
+    parser.add_argument("--test", action="store_true", help="连接后保持运行,便于验证(默认也常驻)")
+    args = parser.parse_args()
+
+    cfg = _load_config()
+    _log(f"启动小红书数据机器人监听...")
+    _log(f"app_id={cfg['feishu_app_id']}, 群={cfg['feishu_chat_id']}")
+
+    from lark_oapi.ws import Client
+    from lark_oapi.ws.client import EventDispatcherHandler
+
+    handler = (
+        EventDispatcherHandler.builder("", "")
+        .register_p2_im_message_receive_v1(on_message_receive)
+        .build()
+    )
+    client = Client(cfg["feishu_app_id"], cfg["feishu_app_secret"], event_handler=handler)
+    _log("连接飞书长连接 (WebSocket)...")
+    client.start()
+    _log("连接已断开")
+
+
+if __name__ == "__main__":
+    main()