Browse Source

feat: 飞书群聊自动转发机器人 + Web 管理面板

- Bot 核心:WebSocket 长连接 + 原生 Forward API 转发
- 防死循环:过滤 sender_type=bot 消息
- 幂等去重:message_id LRU 缓存
- 限流缓冲:asyncio.Queue 单 Worker 匀速消费 (<5 QPS)
- Web 面板:FastAPI + 单文件前端,配置/状态/日志三 Tab
- 子进程隔离:Bot 崩溃自动恢复,配置变更自动重连
- 密码登录鉴权 + SSE 实时日志推送

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
lihetongxue 2 tháng trước cách đây
commit
4a23bdbf80
17 tập tin đã thay đổi với 1900 bổ sung và 0 xóa
  1. 4 0
      .env.example
  2. 14 0
      .gitignore
  3. 18 0
      Dockerfile
  4. 115 0
      README.md
  5. 169 0
      bot.py
  6. 35 0
      config.example.yaml
  7. 110 0
      config.py
  8. 50 0
      dedup.py
  9. 140 0
      forwarder.py
  10. 69 0
      lark2lark.md
  11. 59 0
      main.py
  12. 98 0
      ratelimit.py
  13. 6 0
      requirements.txt
  14. 0 0
      web/__init__.py
  15. 51 0
      web/auth.py
  16. 502 0
      web/index.html
  17. 460 0
      web/supervisor.py

+ 4 - 0
.env.example

@@ -0,0 +1,4 @@
+# 复制为 .env 并填入实际凭据(生产环境推荐)
+# 也可不使用 .env,直接编辑 config.yaml
+APP_ID=cli_xxxxxxxxxxxxxxxx
+APP_SECRET=xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx

+ 14 - 0
.gitignore

@@ -0,0 +1,14 @@
+# Python
+__pycache__/
+*.py[cod]
+*.egg-info/
+.venv/
+venv/
+
+# IDE
+.vscode/
+.idea/
+
+# 日志
+*.log
+logs/

+ 18 - 0
Dockerfile

@@ -0,0 +1,18 @@
+FROM python:3.11-slim
+
+WORKDIR /app
+
+# 依赖(利用层缓存)
+COPY requirements.txt .
+RUN pip install --no-cache-dir -r requirements.txt
+
+# 代码
+COPY *.py ./
+COPY config.example.yaml ./
+
+# 配置通过 volume 挂载 config.yaml,或通过环境变量注入
+# docker run -d --name lark2lark \
+#   -v $(pwd)/config.yaml:/app/config.yaml \
+#   --restart unless-stopped \
+#   lark2lark
+CMD ["python", "main.py"]

+ 115 - 0
README.md

@@ -0,0 +1,115 @@
+# lark2lark - 飞书群聊自动跨群转发机器人
+
+基于飞书企业自建应用 + WebSocket 长连接,实现 7x24 小时无人值守的消息自动转发。
+
+## 架构
+
+```
+源群 ──┐
+       ├─ WebSocket 事件 ──> 白名单过滤 ──> 防死循环 ──> 去重 ──> 限流队列 ──> Forward API ──> 目标群
+源群 ──┘                                                                  (匀速 <5 QPS)
+```
+
+| 模块 | 职责 |
+|------|------|
+| `config.py` | 加载 YAML / 环境变量配置 |
+| `dedup.py` | message_id LRU 去重(飞书重发保护) |
+| `ratelimit.py` | asyncio.Queue + 单 Worker 匀速消费(<5 QPS) |
+| `forwarder.py` | 事件处理 + 白名单 + sender_type=bot 过滤 + Forward API 调用 |
+| `bot.py` | Bot 子进程入口:WebSocket 长连接 + 限流 Worker |
+| `main.py` | 入口,启动 Web 面板(含 Bot 子进程管理) |
+| `web/supervisor.py` | FastAPI 面板后端:配置/状态/日志 API + 子进程管理 |
+| `web/auth.py` | 面板密码登录 + Cookie 签名鉴权 |
+| `web/index.html` | 单文件前端:配置/状态/日志三 Tab |
+
+## 快速开始
+
+### 1. 创建飞书自建应用
+
+1. 登录 [飞书开放平台](https://open.feishu.cn/) → 创建企业自建应用
+2. **凭证与基础信息**:获取 `App ID` 和 `App Secret`
+3. **权限管理**:申请以下 Scopes
+   - `im:message.group_msg`(获取群组所有消息)
+   - `im:message:send_as_bot`(以机器人身份发消息)
+   - `im:resource`(可选,建议申请)
+4. **事件与回调** → 事件配置 → 订阅方式选 **使用长连接接收事件**
+5. 订阅事件:`im.message.receive_v1`
+6. **机器人** 菜单:启用机器人
+7. 发布版本 → 管理员审批
+
+### 2. 把机器人加入源群和目标群
+
+在每个群设置 → 群机器人 → 添加机器人 → 选你的应用。
+
+### 3. 获取 chat_id
+
+```bash
+# 可用飞书开放平台调试台,或调用获取群列表 API
+# https://open.feishu.cn/document/uAjLw4CM/ukTMukTMukTM/reference/im-v1/chat/list
+```
+
+### 4. 配置
+
+```bash
+cp config.example.yaml config.yaml
+# 编辑 config.yaml 填入 app_id / app_secret / source_chat_ids / target_chat_ids
+```
+
+### 5. 设置面板密码(重要)
+
+```bash
+# Linux
+export LARK_PANEL_PASSWORD="your-strong-password"
+# Windows PowerShell
+$env:LARK_PANEL_PASSWORD="your-strong-password"
+```
+
+### 6. 运行
+
+```bash
+pip install -r requirements.txt
+python main.py
+```
+
+启动后访问 `http://服务器IP:8080`,用密码登录面板。
+
+> 首次启动可只用面板配置:把 `config.example.yaml` 复制为 `config.yaml`(不必填真实凭据),启动后通过面板的"配置"页填入并保存即可,Bot 会自动重连。
+
+## Web 管理面板
+
+面板提供三个功能页:
+
+- **配置**:在线编辑 app_id / app_secret / 源群 / 目标群 / QPS 等参数,保存后 Bot 自动重启生效(无需手动重启进程)
+- **状态**:实时显示 Bot 进程存活、WebSocket 连接状态、队列积压、去重缓存、运行时长(每 3 秒刷新)
+- **日志**:SSE 实时推送 Bot 日志,支持级别过滤,最多保留 500 条历史
+
+面板密码通过环境变量 `LARK_PANEL_PASSWORD` 设置(推荐),或在 `config.yaml` 的 `panel.password` 配置(仅本地调试用)。
+
+## Docker 部署
+
+```bash
+docker build -t lark2lark .
+docker run -d --name lark2lark \
+  -p 8080:8080 \
+  -e LARK_PANEL_PASSWORD="your-strong-password" \
+  -v $(pwd)/config.yaml:/app/config.yaml \
+  --restart unless-stopped \
+  lark2lark
+```
+
+## 设计要点
+
+- **WebSocket 长连接**:免公网 IP / 域名 / Webhook,SDK 自动心跳保活
+- **防死循环**:丢弃 `sender_type == "bot"` 的消息,避免双向转发风暴
+- **幂等去重**:本地 LRU 缓存 message_id(不依赖 event_id),防止网络重发
+- **限流缓冲**:asyncio.Queue 单 Worker 匀速消费,间隔 `1/max_qps` 秒,保守默认 4 QPS
+- **重试退避**:单条失败按指数退避重试,最终失败仅告警不阻塞队列
+- **原生 Forward API**:无需下载二进制再上传,支持图片/文件/富文本无损转发
+- **Web 面板**:FastAPI 后端 + 单文件前端,在线配置/状态/日志,密码登录保护
+- **子进程隔离**:Bot 运行在子进程中,配置变更自动重启,崩溃自动恢复
+
+## 约束(飞书官方限制)
+
+1. 不支持转发:红包、投票、语音、日程转让、端到端加密消息
+2. 源消息被设置"禁止转发"时 API 无法越权转发
+3. "合并转发"消息包无法二次拆分转发

+ 169 - 0
bot.py

@@ -0,0 +1,169 @@
+"""Bot 子进程入口:启动 WebSocket 长连接 + 限流 Worker。
+
+被 supervisor.py 以子进程方式启动,stdout/stderr 被 supervisor 捕获用于日志和状态采集。
+定期向 stdout 输出 __STATUS__{...} JSON 行,supervisor 解析后暴露给面板。
+
+也可独立运行:python bot.py
+"""
+from __future__ import annotations
+
+import argparse
+import asyncio
+import json
+import logging
+import signal
+import sys
+import threading
+import time
+from typing import Optional
+
+import lark_oapi as lark
+
+from config import Config, load_config
+from dedup import MessageDedup
+from forwarder import Forwarder, build_event_handler
+from ratelimit import RateLimitedWorker
+
+logger = logging.getLogger("lark2lark")
+
+# 状态行前缀,supervisor 据此识别
+STATUS_PREFIX = "__STATUS__"
+_STATUS_INTERVAL = 5.0  # 秒
+
+# 全局引用,供状态输出线程读取
+_dedup: Optional[MessageDedup] = None
+_worker: Optional[RateLimitedWorker] = None
+_ws_connected = threading.Event()
+
+
+def setup_logging(level: str) -> None:
+    logging.basicConfig(
+        level=getattr(logging, level, logging.INFO),
+        format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
+        datefmt="%Y-%m-%d %H:%M:%S",
+        stream=sys.stdout,  # 关键:输出到 stdout 供 supervisor 捕获
+    )
+    logging.getLogger("lark-oapi").setLevel(logging.WARNING)
+
+
+def parse_args() -> argparse.Namespace:
+    p = argparse.ArgumentParser(description="飞书群聊自动转发机器人(子进程)")
+    p.add_argument("--config", "-c", default=None, help="配置文件路径")
+    return p.parse_args()
+
+
+def start_ws_client(cfg: Config, event_handler, ready: threading.Event) -> threading.Thread:
+    """在子线程启动 lark WebSocket 客户端。"""
+    def _run() -> None:
+        cli = lark.ws.Client(
+            cfg.app_id,
+            cfg.app_secret,
+            event_handler=event_handler,
+            log_level=lark.LogLevel.INFO,
+        )
+        ready.set()
+        _ws_connected.set()
+        logger.info("WebSocket 客户端启动")
+        try:
+            cli.start()
+        except Exception as e:
+            logger.exception("WebSocket 客户端异常退出: %s", e)
+        finally:
+            _ws_connected.clear()
+
+    t = threading.Thread(target=_run, name="lark-ws", daemon=True)
+    t.start()
+    return t
+
+
+def _status_loop(start_time: float, stop_event: threading.Event) -> None:
+    """每 5 秒向 stdout 输出 __STATUS__{...} 行,供 supervisor 采集。"""
+    while not stop_event.wait(_STATUS_INTERVAL):
+        status = {
+            "ts": int(time.time()),
+            "uptime": int(time.time() - start_time),
+            "ws_connected": _ws_connected.is_set(),
+            "queue_size": 0,
+            "dedup_size": 0,
+        }
+        if _worker is not None:
+            # asyncio.Queue 没有线程安全的 qsize 跨线程读,但底层有 _queue
+            try:
+                status["queue_size"] = _worker._queue.qsize()  # type: ignore[union-attr]
+            except Exception:
+                pass
+        if _dedup is not None:
+            stats = _dedup.stats()
+            status["dedup_size"] = stats["size"]
+            status["dedup_max"] = stats["max_size"]
+        sys.stdout.write(STATUS_PREFIX + json.dumps(status, ensure_ascii=False) + "\n")
+        sys.stdout.flush()
+
+
+async def main_async(cfg: Config) -> None:
+    global _dedup, _worker
+    loop = asyncio.get_running_loop()
+
+    client = lark.Client.builder().app_id(cfg.app_id).app_secret(cfg.app_secret).build()
+    _dedup = MessageDedup(max_size=cfg.dedup_cache_size)
+    _worker = RateLimitedWorker(max_qps=cfg.max_qps, max_retry=cfg.max_retry)
+    forwarder = Forwarder(cfg, client, _dedup, _worker, loop)
+    _worker.set_forward_fn(forwarder.forward)
+
+    stop_event = asyncio.Event()
+    worker_task = asyncio.create_task(run_worker(_worker, stop_event))
+
+    event_handler = build_event_handler(forwarder)
+    ws_ready = threading.Event()
+    ws_thread = start_ws_client(cfg, event_handler, ws_ready)
+    ws_ready.wait(timeout=5)
+
+    # 状态输出线程
+    start_time = time.time()
+    status_stop = threading.Event()
+    status_thread = threading.Thread(
+        target=_status_loop, args=(start_time, status_stop),
+        name="status", daemon=True,
+    )
+    status_thread.start()
+
+    def _on_signal(*_):
+        logger.info("收到退出信号,开始关闭...")
+        status_stop.set()
+        stop_event.set()
+
+    try:
+        signal.signal(signal.SIGINT, _on_signal)
+        signal.signal(signal.SIGTERM, _on_signal)
+    except (ValueError, AttributeError):
+        pass
+
+    await worker_task
+    status_stop.set()
+    logger.info("Bot 已停止")
+
+
+async def run_worker(worker: RateLimitedWorker, stop_event: asyncio.Event) -> None:
+    task = asyncio.create_task(worker.run())
+    await stop_event.wait()
+    task.cancel()
+    try:
+        await task
+    except asyncio.CancelledError:
+        pass
+
+
+def main() -> None:
+    args = parse_args()
+    cfg = load_config(args.config)
+    setup_logging(cfg.log_level)
+    logger.info("Bot 启动:源群 %d 个,目标群 %d 个,QPS=%d",
+                len(cfg.source_chat_ids), len(cfg.target_chat_ids), cfg.max_qps)
+    try:
+        asyncio.run(main_async(cfg))
+    except KeyboardInterrupt:
+        logger.info("已退出")
+
+
+if __name__ == "__main__":
+    main()

+ 35 - 0
config.example.yaml

@@ -0,0 +1,35 @@
+# ===== 飞书应用凭据 =====
+# 在飞书开发者后台 -> 凭证与基础信息 获取
+app_id: "cli_xxxxxxxxxxxxxxxx"
+app_secret: "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"
+
+# ===== 源群与目标群 =====
+# 源群白名单:只转发这些 chat_id 的消息(至少 1 个)
+source_chat_ids:
+  - "oc_source_group_aaaaaaaa"
+  # - "oc_source_group_bbbbbbbb"
+
+# 目标群:消息转发到这些 chat_id(至少 1 个)
+target_chat_ids:
+  - "oc_target_group_cccccccc"
+  # - "oc_target_group_dddddddd"
+
+# ===== 去重 / 限流 =====
+# message_id 去重缓存大小(条数)
+dedup_cache_size: 2000
+
+# 目标群发送 QPS 上限(飞书限制 5 QPS/群,保守取 4)
+max_qps: 4
+
+# 单条转发失败的重试次数(不含首次)
+max_retry: 2
+
+# ===== 日志 =====
+log_level: "INFO"   # DEBUG / INFO / WARNING / ERROR
+
+# ===== Web 面板 =====
+panel:
+  host: "0.0.0.0"
+  port: 8080
+  # 面板登录密码:优先用环境变量 LARK_PANEL_PASSWORD,勿明文写入此文件
+  # password: "changeme"   # 仅本地调试用

+ 110 - 0
config.py

@@ -0,0 +1,110 @@
+"""配置加载:优先环境变量,其次 config.yaml,最后 config.example.yaml。"""
+from __future__ import annotations
+
+import os
+from dataclasses import dataclass, field
+from pathlib import Path
+from typing import List
+
+import yaml
+
+try:
+    from dotenv import load_dotenv
+    load_dotenv()  # 若存在 .env 则自动加载
+except ImportError:
+    pass  # python-dotenv 非必须
+
+
+@dataclass
+class PanelConfig:
+    host: str = "0.0.0.0"
+    port: int = 8080
+    password: str = ""
+
+
+@dataclass
+class Config:
+    app_id: str
+    app_secret: str
+    source_chat_ids: List[str]
+    target_chat_ids: List[str]
+    dedup_cache_size: int = 2000
+    max_qps: int = 4
+    max_retry: int = 2
+    log_level: str = "INFO"
+    panel: PanelConfig = field(default_factory=PanelConfig)
+
+    def validate(self) -> None:
+        if not self.app_id or not self.app_secret:
+            raise ValueError("app_id / app_secret 不能为空(请在 config.yaml 或环境变量中配置)")
+        if not self.source_chat_ids:
+            raise ValueError("source_chat_ids 至少配置 1 个源群 chat_id")
+        if not self.target_chat_ids:
+            raise ValueError("target_chat_ids 至少配置 1 个目标群 chat_id")
+        if self.max_qps <= 0 or self.max_qps > 5:
+            raise ValueError(f"max_qps 必须在 (0, 5] 区间,当前 {self.max_qps}")
+
+    @property
+    def panel_password(self) -> str:
+        """面板密码:环境变量优先,其次配置文件。"""
+        return os.getenv("LARK_PANEL_PASSWORD") or self.panel.password
+
+
+def _load_yaml(path: Path) -> dict:
+    if not path.exists():
+        return {}
+    with path.open("r", encoding="utf-8") as f:
+        return yaml.safe_load(f) or {}
+
+
+def load_config(config_path: str | Path | None = None) -> Config:
+    """加载顺序:环境变量 > config.yaml > config.example.yaml。"""
+    base = Path(__file__).resolve().parent
+    candidates = []
+    if config_path:
+        candidates.append(Path(config_path))
+    candidates.append(base / "config.yaml")
+    candidates.append(base / "config.example.yaml")
+
+    merged: dict = {}
+    for p in candidates:
+        merged.update(_load_yaml(p))
+
+    # 环境变量覆盖(优先级最高)
+    env_app_id = os.getenv("APP_ID")
+    env_app_secret = os.getenv("APP_SECRET")
+    if env_app_id:
+        merged["app_id"] = env_app_id
+    if env_app_secret:
+        merged["app_secret"] = env_app_secret
+
+    # 标准化白名单:去空白、去空串
+    for key in ("source_chat_ids", "target_chat_ids"):
+        val = merged.get(key) or []
+        merged[key] = [str(x).strip() for x in val if str(x).strip()]
+
+    # panel 段
+    panel_raw = merged.get("panel") or {}
+    panel = PanelConfig(
+        host=str(panel_raw.get("host", "0.0.0.0")),
+        port=int(panel_raw.get("port", 8080)),
+        password=str(panel_raw.get("password", "")),
+    )
+
+    try:
+        cfg = Config(
+            app_id=merged.get("app_id", ""),
+            app_secret=merged.get("app_secret", ""),
+            source_chat_ids=merged["source_chat_ids"],
+            target_chat_ids=merged["target_chat_ids"],
+            dedup_cache_size=int(merged.get("dedup_cache_size", 2000)),
+            max_qps=int(merged.get("max_qps", 4)),
+            max_retry=int(merged.get("max_retry", 2)),
+            log_level=str(merged.get("log_level", "INFO")).upper(),
+            panel=panel,
+        )
+    except KeyError as e:
+        raise ValueError(f"配置缺失字段: {e}") from e
+
+    cfg.validate()
+    return cfg

+ 50 - 0
dedup.py

@@ -0,0 +1,50 @@
+"""message_id 幂等去重:基于 OrderedDict 实现 LRU + 容量上限。
+
+飞书在网络抖动时会重发事件推送,不能依赖 event_id,必须用 message_id 本地去重。
+"""
+from __future__ import annotations
+
+import logging
+import threading
+from collections import OrderedDict
+from typing import Optional
+
+logger = logging.getLogger(__name__)
+
+
+class MessageDedup:
+    """线程安全的 LRU message_id 去重器。
+
+    飞书事件推送在网络抖动时会重发,必须用 message_id 在本地去重,
+    不能依赖 event_id(官方文档明确说明 event_id 不保证幂等)。
+    """
+
+    def __init__(self, max_size: int = 2000) -> None:
+        if max_size <= 0:
+            raise ValueError("max_size 必须为正数")
+        self._max_size = max_size
+        self._seen: OrderedDict[str, None] = OrderedDict()
+        self._lock = threading.Lock()
+
+    def check_and_mark(self, message_id: str) -> bool:
+        """若 message_id 首次出现则记录并返回 True;已存在则返回 False(重复)。"""
+        if not message_id:
+            return False
+        with self._lock:
+            if message_id in self._seen:
+                # 命中:移到末尾(LRU)
+                self._seen.move_to_end(message_id)
+                return False
+            self._seen[message_id] = None
+            if len(self._seen) > self._max_size:
+                evicted_k, _ = self._seen.popitem(last=False)
+                logger.debug("去重缓存淘汰最旧 message_id: %s", evicted_k)
+            return True
+
+    def __len__(self) -> int:
+        with self._lock:
+            return len(self._seen)
+
+    def stats(self) -> dict:
+        with self._lock:
+            return {"size": len(self._seen), "max_size": self._max_size}

+ 140 - 0
forwarder.py

@@ -0,0 +1,140 @@
+"""事件处理 + 转发核心:白名单过滤、防死循环、去重、调用 Forward API。
+
+lark.ws.Client 在独立线程中运行事件循环并回调 handler。
+本模块的事件回调把任务投递到 asyncio.Queue(跨线程安全),
+由 ratelimit.Worker 在主事件循环中匀速消费。
+"""
+from __future__ import annotations
+
+import asyncio
+import logging
+import threading
+from typing import Optional
+
+import lark_oapi as lark
+from lark_oapi.api.im.v1 import (
+    ForwardMessageRequest,
+    ForwardMessageRequestBody,
+)
+
+from config import Config
+from dedup import MessageDedup
+from ratelimit import ForwardTask, RateLimitedWorker
+
+logger = logging.getLogger(__name__)
+
+
+class Forwarder:
+    """绑定配置、去重器、限流 Worker,并提供事件回调。"""
+
+    def __init__(
+        self,
+        cfg: Config,
+        client: lark.Client,
+        dedup: MessageDedup,
+        worker: RateLimitedWorker,
+        loop: asyncio.AbstractEventLoop,
+    ) -> None:
+        self._cfg = cfg
+        self._client = client
+        self._dedup = dedup
+        self._worker = worker
+        self._loop = loop
+        self._source_set = set(cfg.source_chat_ids)
+
+    # ---------- 事件回调(lark SDK 在子线程调用) ----------
+    def on_message_receive(self, data: lark.im.v1.P2ImMessageReceiveV1) -> None:
+        try:
+            self._handle_event(data)
+        except Exception as e:
+            logger.exception("事件处理异常: %s", e)
+
+    def _handle_event(self, data: lark.im.v1.P2ImMessageReceiveV1) -> None:
+        event = data.event
+        if event is None:
+            return
+
+        message = getattr(event, "message", None)
+        sender = getattr(event, "sender", None)
+        if message is None:
+            return
+
+        message_id: str = message.message_id or ""
+        chat_id: str = message.chat_id or ""
+        sender_type: str = (sender.sender_type if sender else "") or ""
+
+        # 1) 防死循环:丢弃机器人自己发出的消息
+        if sender_type == "bot":
+            logger.debug("跳过 bot 消息 message_id=%s", message_id)
+            return
+
+        # 2) 白名单:只处理源群
+        if chat_id not in self._source_set:
+            logger.debug("跳过非白名单群 chat_id=%s", chat_id)
+            return
+
+        # 3) 去重
+        if not self._dedup.check_and_mark(message_id):
+            logger.debug("重复消息已丢弃 message_id=%s", message_id)
+            return
+
+        logger.info("接收消息 message_id=%s chat_id=%s type=%s",
+                    message_id, chat_id, message.message_type)
+
+        # 4) 投递转发任务(每个目标群一个)
+        for target in self._cfg.target_chat_ids:
+            if target == chat_id:
+                # 源==目标,跳过避免无效转发
+                continue
+            task = ForwardTask(
+                message_id=message_id,
+                target_chat_id=target,
+                source_chat_id=chat_id,
+            )
+            # 跨线程安全投递到主事件循环的队列
+            asyncio.run_coroutine_threadsafe(
+                self._worker.enqueue(task), self._loop
+            )
+
+    # ---------- 实际转发(Worker 调用,运行在主事件循环) ----------
+    async def forward(self, task: ForwardTask) -> bool:
+        """调用飞书 Forward API。返回 True 表示成功。"""
+        req = (
+            ForwardMessageRequest.builder()
+            .message_id(task.message_id)
+            .receive_id_type("chat_id")
+            .request_body(
+                ForwardMessageRequestBody.builder()
+                .receive_id(task.target_chat_id)
+                .build()
+            )
+            .build()
+        )
+
+        def _call() -> bool:
+            try:
+                resp = self._client.im.v1.message.forward(req)
+            except Exception as e:
+                logger.error("Forward API 异常 message_id=%s -> %s: %s",
+                             task.message_id, task.target_chat_id, e)
+                return False
+            if not resp.success():
+                logger.error("Forward API 失败 message_id=%s -> %s code=%s msg=%s",
+                             task.message_id, task.target_chat_id,
+                             resp.code, resp.msg)
+                return False
+            logger.info("转发成功 message_id=%s -> %s",
+                        task.message_id, task.target_chat_id)
+            return True
+
+        # lark-oapi 的 Client 是同步阻塞调用,放到默认 executor 执行
+        return await asyncio.get_event_loop().run_in_executor(None, _call)
+
+
+def build_event_handler(forwarder: Forwarder):
+    """构造 lark EventDispatcherHandler,注册 im.message.receive_v1。"""
+    return (
+        lark.EventDispatcherHandler.builder("", "")
+        .register_p2_im_message_receive_v1(forwarder.on_message_receive)
+        .build()
+    )

+ 69 - 0
lark2lark.md

@@ -0,0 +1,69 @@
+# 飞书群聊全自动跨群转发机器人架构需求文档
+
+## 1. 项目背景与目标
+
+* **业务需求**:实现飞书源群聊到目标群聊的消息全自动、低延迟转发。
+* **功能要求**:支持包括文本、图片、超链接、富文本以及常见附件(Word、PDF 等)在内的多种消息载体无损流转。
+* **运行环境要求**:系统需实现“云端代挂”,完全脱离用户的本地 PC,不需要用户保持个人飞书客户端的登录状态,实现 7x24 小时无人值守。
+
+## 2. 整体系统架构选型
+
+* **应用形态**:飞书企业自建应用(Custom App),以“机器人(Bot)”身份运行。这使得程序拥有独立的 `Tenant Access Token` 进行鉴权,彻底与个人飞书账号解绑。
+* **事件订阅通道**:采用 **WebSocket 长连接模式(Long Connection)**。
+* *优势*:极大地简化了网络配置,无需配置公网 IP、无需域名、无需繁琐的 Webhook 签名校验,只需服务器能访问外网即可建立安全的双向通信。
+
+
+* **部署架构**:Linux 云服务器(如阿里云、腾讯云的基础款轻量应用服务器,1核1G配置即可满足)结合 Docker 容器化部署,实现全天候云端运行。
+
+## 3. 核心功能模块设计
+
+### 3.1 监听与事件接收模块
+
+* **事件订阅**:在飞书开发者后台订阅 `im.message.receive_v1`(接收消息事件)。
+* **逻辑过滤**:
+* **群组白名单**:解析接收到的 JSON 载荷,提取 `chat_id`,仅当该 ID 匹配预设的“源群组”时才放行。
+* **防死循环机制**:严格校验 `sender_type` 字段,丢弃所有 `sender_type == "bot"`(机器人发出)的消息,防止 A 群与 B 群之间产生无限转发风暴。
+
+
+
+### 3.2 消息转发与处理模块
+
+根据业务需求,消息转发建议采用**官方原生转发接口 (Forward API)**,此方案开发成本最低且支持格式最全。
+
+* **转发链路**:提取源事件中的 `message_id`,直接调用 `POST /open-apis/im/v1/messages/:message_id/forward` 接口。
+* **格式支持**:该原生接口不仅支持文本、链接,还天然支持包含图片、视频、Word、PDF 等在内的富文本与文件卡片,无需开发者在云服务器上先下载二进制文件再重新上传,极大节省了服务器带宽与 I/O 成本。
+
+### 3.3 幂等性与去重模块
+
+* **去重机制**:由于网络波动时飞书云端可能会重发推送,必须建立幂等性保护。程序不应依赖外部的 `event_id`,而应在本地内存或轻量级数据库(如 SQLite / Redis)中缓存已处理过的 `message_id`。接收到新推送时,若 `message_id` 已存在则直接丢弃。
+
+### 3.4 流量控制模块(限频缓冲)
+
+* 飞书官方对于发送消息有严格的频率限制:**向同一群组发送消息的上限为 5 QPS**(即每秒 5 条)。
+* **削峰设计**:在代码内部引入本地消息队列(如 Python 的 `asyncio.Queue` 或 Go 的 `Channel`)。监听模块收到消息后即刻推入队列,由一个单独的 Worker 消费线程以低于 5 QPS 的速度匀速向目标群组调用发送接口,防止因群聊消息瞬时并发而触发 API 封禁错误。
+
+## 4. 飞书开放平台权限清单 (Scopes)
+
+在开发前,需在飞书开发者后台向租户管理员申请并发布以下权限:
+
+1. **`im:message.group_msg`**:获取群组中所有消息(核心权限,用于静默监听源群组内所有人的聊天内容)。
+2. **`im:message:send_as_bot`**:以应用的身份发消息(用于在目标群组中执行推送)。
+3. **`im:resource`**(备选):如果未来需要自行下载、解析并重新组装图片或 Word/PDF 文件,必须申请此“获取消息中的资源文件”权限。若仅调用原生 Forward 接口则非必须,但建议一并申请。
+
+## 5. 开发语言与 SDK 推荐
+
+* **推荐语言**:Python 或 Go。
+* **官方 SDK**:
+* Python 使用 `lark-oapi`。
+* Go 使用 `oapi-sdk-go`。
+
+
+* 官方 SDK 已将底层的 Token 获取、刷新以及 WebSocket 长连接的心跳保活机制全部封装,开发者只需实现极简的事件回调函数即可。
+
+## 6. 约束与系统局限性说明
+
+需要提前规划或在业务层面注意以下飞书官方的系统级限制:
+
+1. **不支持的类型**:官方转发接口不支持转发红包、投票、语音、日程转让及端到端加密消息。
+2. **禁止转发限制**:如果源消息的发送者或源群主将特定消息设置了“禁止转发”,API 将无法越权转发。
+3. **合并转发的二次限制**:如果源群里有人发了一条“合并转发”的消息包,你的机器人无法通过 API 剥离该包内的子消息并进行二次转发。

+ 59 - 0
main.py

@@ -0,0 +1,59 @@
+"""入口:启动 Web 管理面板(含 Bot 子进程管理)。
+
+运行:
+    python main.py
+    python main.py --config /path/to/config.yaml
+
+环境变量:
+    LARK_PANEL_PASSWORD  面板登录密码(推荐,优先于配置文件)
+"""
+from __future__ import annotations
+
+import argparse
+import logging
+import sys
+from pathlib import Path
+
+# 确保项目根在 sys.path
+sys.path.insert(0, str(Path(__file__).resolve().parent))
+
+from config import load_config  # noqa: E402
+
+
+def parse_args() -> argparse.Namespace:
+    p = argparse.ArgumentParser(description="飞书群聊自动转发机器人 - 管理面板")
+    p.add_argument("--config", "-c", default=None, help="配置文件路径")
+    return p.parse_args()
+
+
+def main() -> None:
+    args = parse_args()
+    cfg = load_config(args.config)
+
+    logging.basicConfig(
+        level=logging.INFO,
+        format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
+        datefmt="%Y-%m-%d %H:%M:%S",
+    )
+
+    if not cfg.panel_password:
+        print("=" * 60, file=sys.stderr)
+        print("警告:面板未设置密码!", file=sys.stderr)
+        print("请设置环境变量 LARK_PANEL_PASSWORD 或在 config.yaml 的", file=sys.stderr)
+        print("panel.password 字段配置密码后再启动。", file=sys.stderr)
+        print("=" * 60, file=sys.stderr)
+
+    import uvicorn
+    from web.supervisor import app
+
+    print(f"面板启动:http://{cfg.panel.host}:{cfg.panel.port}", flush=True)
+    uvicorn.run(
+        app,
+        host=cfg.panel.host,
+        port=cfg.panel.port,
+        log_level="warning",  # 减少 uvicorn 自己的日志噪音
+    )
+
+
+if __name__ == "__main__":
+    main()

+ 98 - 0
ratelimit.py

@@ -0,0 +1,98 @@
+"""限流缓冲队列:单 Worker 匀速消费,保证 < max_qps QPS。
+
+飞书限制:同一群发送消息上限 5 QPS。本模块用 asyncio.Queue 削峰,
+独立 Worker 协程以 1/max_qps 秒/条的速率匀速消费,避免瞬时并发触发封禁。
+"""
+from __future__ import annotations
+
+import asyncio
+import logging
+from dataclasses import dataclass
+from typing import Awaitable, Callable, Optional
+
+logger = logging.getLogger(__name__)
+
+
+@dataclass
+class ForwardTask:
+    """待转发的单条任务:源消息 id -> 单个目标群。"""
+    message_id: str
+    target_chat_id: str
+    # 可选:携带源信息用于日志
+    source_chat_id: str = ""
+
+
+class RateLimitedWorker:
+    """单消费协程 + asyncio.Queue,匀速向目标群转发。
+
+    - 入队即返回(非阻塞),由 Worker 串行消费
+    - 每条间隔 1/max_qps 秒,确保 QPS 不超限
+    - 单条失败按 max_retry 退避重试
+    """
+
+    def __init__(
+        self,
+        max_qps: int,
+        max_retry: int,
+        forward_fn: Optional[Callable[[ForwardTask], Awaitable[bool]]] = None,
+        queue: Optional[asyncio.Queue] = None,
+    ) -> None:
+        if max_qps <= 0:
+            raise ValueError("max_qps 必须为正数")
+        self._interval = 1.0 / max_qps
+        self._max_retry = max_retry
+        self._forward_fn = forward_fn
+        self._queue: asyncio.Queue = queue or asyncio.Queue()
+
+    def set_forward_fn(self, fn: Callable[[ForwardTask], Awaitable[bool]]) -> None:
+        """后置注入转发函数(解决与 Forwarder 的循环依赖)。"""
+        self._forward_fn = fn
+
+    async def enqueue(self, task: ForwardTask) -> None:
+        await self._queue.put(task)
+
+    async def run(self) -> None:
+        logger.info("转发 Worker 启动,目标速率 %.2f QPS(间隔 %.3fs)",
+                    1.0 / self._interval, self._interval)
+        while True:
+            task: ForwardTask = await self._queue.get()
+            try:
+                await self._process(task)
+            except asyncio.CancelledError:
+                logger.info("Worker 收到取消信号,退出")
+                raise
+            except Exception as e:
+                logger.exception("Worker 处理异常: %s", e)
+            finally:
+                self._queue.task_done()
+
+    async def _process(self, task: ForwardTask) -> None:
+        if self._forward_fn is None:
+            logger.error("forward_fn 未设置,丢弃任务 message_id=%s", task.message_id)
+            return
+        for attempt in range(self._max_retry + 1):
+            ok = await self._forward_fn(task)
+            if ok:
+                if attempt > 0:
+                    logger.info("重试成功 message_id=%s -> %s(第 %d 次)",
+                                task.message_id, task.target_chat_id, attempt)
+                break
+            # 失败:退避后重试
+            backoff = self._interval * (attempt + 1) * 2
+            logger.warning("转发失败 message_id=%s -> %s,%ds 后重试(%d/%d)",
+                           task.message_id, task.target_chat_id,
+                           int(backoff), attempt + 1, self._max_retry)
+            await asyncio.sleep(backoff)
+        else:
+            logger.error("转发最终失败 message_id=%s -> %s,已放弃",
+                         task.message_id, task.target_chat_id)
+        # 成功或失败都遵守速率间隔
+        await asyncio.sleep(self._interval)
+
+    async def wait_drained(self, timeout: Optional[float] = None) -> bool:
+        """等待队列清空,返回是否在 timeout 内清空。"""
+        try:
+            await asyncio.wait_for(self._queue.join(), timeout=timeout)
+            return True
+        except asyncio.TimeoutError:
+            return False

+ 6 - 0
requirements.txt

@@ -0,0 +1,6 @@
+lark-oapi>=1.4.0
+PyYAML>=6.0
+fastapi>=0.110
+uvicorn[standard]>=0.27
+itsdangerous>=2.1
+python-multipart>=0.0.9

+ 0 - 0
web/__init__.py


+ 51 - 0
web/auth.py

@@ -0,0 +1,51 @@
+"""面板鉴权:itsdangerous 签名 Cookie。
+
+密码来源优先级:环境变量 LARK_PANEL_PASSWORD > config.panel.password。
+未设置密码时面板无法登录(拒绝空密码)。
+"""
+from __future__ import annotations
+
+import os
+import time
+from typing import Optional
+
+from itsdangerous import BadSignature, SignatureExpired, URLSafeTimedSerializer
+
+COOKIE_NAME = "lark2lark_session"
+COOKIE_MAX_AGE = 7 * 24 * 3600  # 7 天
+
+
+def _get_secret_key(panel_password: str) -> str:
+    """用面板密码派生签名密钥(固定 salt 保证进程间一致)。"""
+    # 密码本身就是准入凭据,直接作为签名密钥源;加固定前缀防止空密码
+    return f"lark2lark::{panel_password}"
+
+
+def has_password(panel_password: str) -> bool:
+    return bool(panel_password and panel_password.strip())
+
+
+def sign_session(panel_password: str) -> str:
+    """登录成功后签发会话 token。"""
+    serializer = URLSafeTimedSerializer(_get_secret_key(panel_password), salt="session")
+    return serializer.dumps({"login_at": int(time.time())})
+
+
+def verify_session(token: str, panel_password: str) -> bool:
+    """校验会话 token。失败返回 False。"""
+    if not token:
+        return False
+    serializer = URLSafeTimedSerializer(_get_secret_key(panel_password), salt="session")
+    try:
+        serializer.loads(token, max_age=COOKIE_MAX_AGE)
+        return True
+    except (BadSignature, SignatureExpired):
+        return False
+
+
+def verify_password(input_password: str, expected_password: str) -> bool:
+    """恒定时间比较,防侧信道计时攻击。"""
+    if not expected_password:
+        return False
+    import hmac
+    return hmac.compare_digest(input_password or "", expected_password)

+ 502 - 0
web/index.html

@@ -0,0 +1,502 @@
+<!DOCTYPE html>
+<html lang="zh-CN">
+<head>
+<meta charset="UTF-8">
+<meta name="viewport" content="width=device-width, initial-scale=1.0">
+<title>lark2lark 管理面板</title>
+<style>
+  :root {
+    --bg: #0f1115;
+    --panel: #1a1d24;
+    --panel-2: #22262f;
+    --border: #2d3139;
+    --text: #e4e6eb;
+    --text-dim: #8a8f9a;
+    --accent: #4a9eff;
+    --accent-hover: #3a8eef;
+    --green: #3ec47e;
+    --red: #ef5350;
+    --yellow: #f5a623;
+    --code-bg: #0a0c10;
+  }
+  * { box-sizing: border-box; margin: 0; padding: 0; }
+  body {
+    font-family: -apple-system, "Segoe UI", "PingFang SC", "Microsoft YaHei", sans-serif;
+    background: var(--bg);
+    color: var(--text);
+    font-size: 14px;
+    line-height: 1.5;
+  }
+  /* 登录页 */
+  #login-view {
+    max-width: 360px;
+    margin: 120px auto;
+    padding: 32px;
+    background: var(--panel);
+    border: 1px solid var(--border);
+    border-radius: 8px;
+  }
+  #login-view h1 { font-size: 20px; margin-bottom: 24px; text-align: center; }
+  #login-view input {
+    width: 100%; padding: 10px 12px;
+    background: var(--panel-2); border: 1px solid var(--border);
+    border-radius: 4px; color: var(--text); font-size: 14px;
+  }
+  #login-view input:focus { outline: none; border-color: var(--accent); }
+  #login-view button {
+    width: 100%; margin-top: 16px; padding: 10px;
+    background: var(--accent); color: #fff; border: none;
+    border-radius: 4px; cursor: pointer; font-size: 14px;
+  }
+  #login-view button:hover { background: var(--accent-hover); }
+  #login-error { color: var(--red); margin-top: 12px; min-height: 18px; font-size: 13px; }
+
+  /* 主面板 */
+  #main-view { display: none; }
+  .topbar {
+    display: flex; align-items: center; justify-content: space-between;
+    padding: 12px 24px; background: var(--panel); border-bottom: 1px solid var(--border);
+  }
+  .topbar h1 { font-size: 16px; font-weight: 600; }
+  .topbar .actions { display: flex; gap: 12px; align-items: center; }
+  .tabs { display: flex; gap: 4px; padding: 0 24px; background: var(--panel); border-bottom: 1px solid var(--border); }
+  .tab {
+    padding: 12px 20px; cursor: pointer; color: var(--text-dim);
+    border-bottom: 2px solid transparent; font-size: 14px;
+  }
+  .tab.active { color: var(--text); border-bottom-color: var(--accent); }
+  .tab:hover { color: var(--text); }
+  .content { padding: 24px; max-width: 960px; margin: 0 auto; }
+  .tab-pane { display: none; }
+  .tab-pane.active { display: block; }
+
+  /* 配置表单 */
+  .form-group { margin-bottom: 18px; }
+  .form-group label { display: block; margin-bottom: 6px; color: var(--text-dim); font-size: 13px; }
+  .form-group input, .form-group textarea, .form-group select {
+    width: 100%; padding: 8px 10px;
+    background: var(--panel-2); border: 1px solid var(--border);
+    border-radius: 4px; color: var(--text); font-size: 14px;
+    font-family: "Consolas", "Monaco", monospace;
+  }
+  .form-group textarea { min-height: 70px; resize: vertical; }
+  .form-group input:focus, .form-group textarea:focus, .form-group select:focus {
+    outline: none; border-color: var(--accent);
+  }
+  .form-row { display: grid; grid-template-columns: 1fr 1fr; gap: 16px; }
+  .form-hint { font-size: 12px; color: var(--text-dim); margin-top: 4px; }
+  .btn-primary {
+    padding: 8px 20px; background: var(--accent); color: #fff;
+    border: none; border-radius: 4px; cursor: pointer; font-size: 14px;
+  }
+  .btn-primary:hover { background: var(--accent-hover); }
+  .btn-primary:disabled { opacity: 0.5; cursor: not-allowed; }
+  .btn-secondary {
+    padding: 6px 14px; background: var(--panel-2); color: var(--text);
+    border: 1px solid var(--border); border-radius: 4px; cursor: pointer; font-size: 13px;
+  }
+  .btn-secondary:hover { border-color: var(--accent); }
+
+  /* 状态卡片 */
+  .status-grid { display: grid; grid-template-columns: repeat(auto-fill, minmax(200px, 1fr)); gap: 16px; }
+  .status-card {
+    padding: 16px; background: var(--panel); border: 1px solid var(--border);
+    border-radius: 6px;
+  }
+  .status-card .label { font-size: 12px; color: var(--text-dim); margin-bottom: 8px; }
+  .status-card .value { font-size: 22px; font-weight: 600; }
+  .status-card .value.ok { color: var(--green); }
+  .status-card .value.err { color: var(--red); }
+  .status-card .value.warn { color: var(--yellow); }
+  .badge {
+    display: inline-block; padding: 2px 8px; border-radius: 10px;
+    font-size: 12px; font-weight: 600;
+  }
+  .badge.ok { background: rgba(62,196,126,0.15); color: var(--green); }
+  .badge.err { background: rgba(239,83,80,0.15); color: var(--red); }
+  .badge.warn { background: rgba(245,166,35,0.15); color: var(--yellow); }
+
+  /* 日志 */
+  #log-filter { display: flex; gap: 8px; margin-bottom: 12px; align-items: center; }
+  #log-filter label { color: var(--text-dim); font-size: 13px; }
+  #log-container {
+    background: var(--code-bg); border: 1px solid var(--border);
+    border-radius: 6px; padding: 12px; height: 60vh; overflow-y: auto;
+    font-family: "Consolas", "Monaco", monospace; font-size: 12.5px; line-height: 1.6;
+  }
+  .log-line { white-space: pre-wrap; word-break: break-all; padding: 1px 0; }
+  .log-line .level-DEBUG { color: #6a737d; }
+  .log-line .level-INFO { color: var(--text); }
+  .log-line .level-WARNING { color: var(--yellow); }
+  .log-line .level-ERROR { color: var(--red); }
+  .log-line .level-CRITICAL { color: var(--red); font-weight: bold; }
+
+  .toast {
+    position: fixed; bottom: 24px; right: 24px; padding: 12px 20px;
+    background: var(--panel); border: 1px solid var(--border); border-radius: 6px;
+    box-shadow: 0 4px 12px rgba(0,0,0,0.3); z-index: 1000;
+    transition: opacity 0.3s; opacity: 0;
+  }
+  .toast.show { opacity: 1; }
+  .toast.ok { border-left: 3px solid var(--green); }
+  .toast.err { border-left: 3px solid var(--red); }
+</style>
+</head>
+<body>
+
+<div id="login-view">
+  <h1>lark2lark 管理面板</h1>
+  <form id="login-form">
+    <input type="password" id="login-password" placeholder="面板密码" autocomplete="current-password">
+    <button type="submit">登录</button>
+    <div id="login-error"></div>
+  </form>
+</div>
+
+<div id="main-view">
+  <div class="topbar">
+    <h1>lark2lark 管理面板</h1>
+    <div class="actions">
+      <button class="btn-secondary" onclick="restartBot()">重启 Bot</button>
+      <button class="btn-secondary" onclick="logout()">退出登录</button>
+    </div>
+  </div>
+  <div class="tabs">
+    <div class="tab active" data-tab="config">配置</div>
+    <div class="tab" data-tab="status">状态</div>
+    <div class="tab" data-tab="logs">日志</div>
+  </div>
+  <div class="content">
+    <div class="tab-pane active" id="pane-config">
+      <form id="config-form">
+        <div class="form-row">
+          <div class="form-group">
+            <label>App ID</label>
+            <input type="text" id="cfg-app_id" autocomplete="off">
+          </div>
+          <div class="form-group">
+            <label>App Secret(留 *** 表示不修改)</label>
+            <input type="text" id="cfg-app_secret" autocomplete="off">
+          </div>
+        </div>
+        <div class="form-group">
+          <label>源群 chat_id 列表(每行一个)</label>
+          <textarea id="cfg-source_chat_ids" placeholder="oc_xxx&#10;oc_yyy"></textarea>
+          <div class="form-hint">只转发这些群的消息</div>
+        </div>
+        <div class="form-group">
+          <label>目标群 chat_id 列表(每行一个)</label>
+          <textarea id="cfg-target_chat_ids" placeholder="oc_zzz"></textarea>
+          <div class="form-hint">消息转发到这些群</div>
+        </div>
+        <div class="form-row">
+          <div class="form-group">
+            <label>去重缓存大小</label>
+            <input type="number" id="cfg-dedup_cache_size" min="100">
+          </div>
+          <div class="form-group">
+            <label>最大 QPS(1-5)</label>
+            <input type="number" id="cfg-max_qps" min="1" max="5">
+          </div>
+        </div>
+        <div class="form-row">
+          <div class="form-group">
+            <label>失败重试次数</label>
+            <input type="number" id="cfg-max_retry" min="0">
+          </div>
+          <div class="form-group">
+            <label>日志级别</label>
+            <select id="cfg-log_level">
+              <option>DEBUG</option>
+              <option>INFO</option>
+              <option>WARNING</option>
+              <option>ERROR</option>
+            </select>
+          </div>
+        </div>
+        <button type="submit" class="btn-primary" id="save-btn">保存并重启</button>
+      </form>
+    </div>
+
+    <div class="tab-pane" id="pane-status">
+      <div class="status-grid" id="status-grid"></div>
+    </div>
+
+    <div class="tab-pane" id="pane-logs">
+      <div id="log-filter">
+        <label>级别:</label>
+        <label><input type="checkbox" id="lvl-DEBUG" checked> DEBUG</label>
+        <label><input type="checkbox" id="lvl-INFO" checked> INFO</label>
+        <label><input type="checkbox" id="lvl-WARNING" checked> WARNING</label>
+        <label><input type="checkbox" id="lvl-ERROR" checked> ERROR</label>
+        <button class="btn-secondary" onclick="clearLogs()">清屏</button>
+        <label style="margin-left:auto"><input type="checkbox" id="autoscroll" checked> 自动滚动</label>
+      </div>
+      <div id="log-container"></div>
+    </div>
+  </div>
+</div>
+
+<div class="toast" id="toast"></div>
+
+<script>
+const $ = (id) => document.getElementById(id);
+let sse = null;
+let statusTimer = null;
+let savedConfig = {};
+
+// ---------- 登录 ----------
+async function checkAuth() {
+  try {
+    const r = await fetch('/api/config');
+    if (r.status === 401) {
+      showLogin();
+      return;
+    }
+    showMain();
+  } catch {
+    showLogin();
+  }
+}
+
+function showLogin() {
+  $('login-view').style.display = 'block';
+  $('main-view').style.display = 'none';
+}
+
+function showMain() {
+  $('login-view').style.display = 'none';
+  $('main-view').style.display = 'block';
+  loadConfig();
+  startStatusPolling();
+  startLogStream();
+}
+
+$('login-form').addEventListener('submit', async (e) => {
+  e.preventDefault();
+  $('login-error').textContent = '';
+  const pwd = $('login-password').value;
+  try {
+    const r = await fetch('/api/login', {
+      method: 'POST',
+      headers: {'Content-Type': 'application/json'},
+      body: JSON.stringify({password: pwd}),
+    });
+    if (r.ok) {
+      showMain();
+    } else {
+      const d = await r.json().catch(() => ({}));
+      $('login-error').textContent = d.detail || '登录失败';
+    }
+  } catch (err) {
+    $('login-error').textContent = '网络错误';
+  }
+});
+
+async function logout() {
+  await fetch('/api/logout', {method: 'POST'});
+  showLogin();
+  if (sse) { sse.close(); sse = null; }
+  if (statusTimer) { clearInterval(statusTimer); statusTimer = null; }
+}
+
+// ---------- Tab 切换 ----------
+document.querySelectorAll('.tab').forEach(tab => {
+  tab.addEventListener('click', () => {
+    document.querySelectorAll('.tab').forEach(t => t.classList.remove('active'));
+    document.querySelectorAll('.tab-pane').forEach(p => p.classList.remove('active'));
+    tab.classList.add('active');
+    $('pane-' + tab.dataset.tab).classList.add('active');
+  });
+});
+
+// ---------- 配置 ----------
+async function loadConfig() {
+  try {
+    const r = await fetch('/api/config');
+    if (!r.ok) return;
+    const cfg = await r.json();
+    savedConfig = cfg;
+    $('cfg-app_id').value = cfg.app_id || '';
+    $('cfg-app_secret').value = cfg.app_secret || '';  // ***
+    $('cfg-source_chat_ids').value = (cfg.source_chat_ids || []).join('\n');
+    $('cfg-target_chat_ids').value = (cfg.target_chat_ids || []).join('\n');
+    $('cfg-dedup_cache_size').value = cfg.dedup_cache_size ?? 2000;
+    $('cfg-max_qps').value = cfg.max_qps ?? 4;
+    $('cfg-max_retry').value = cfg.max_retry ?? 2;
+    $('cfg-log_level').value = cfg.log_level || 'INFO';
+  } catch (err) {
+    toast('加载配置失败: ' + err.message, 'err');
+  }
+}
+
+$('config-form').addEventListener('submit', async (e) => {
+  e.preventDefault();
+  $('save-btn').disabled = true;
+  const payload = {
+    app_id: $('cfg-app_id').value.trim(),
+    app_secret: $('cfg-app_secret').value.trim(),
+    source_chat_ids: $('cfg-source_chat_ids').value.split('\n').map(s => s.trim()).filter(Boolean),
+    target_chat_ids: $('cfg-target_chat_ids').value.split('\n').map(s => s.trim()).filter(Boolean),
+    dedup_cache_size: parseInt($('cfg-dedup_cache_size').value, 10),
+    max_qps: parseInt($('cfg-max_qps').value, 10),
+    max_retry: parseInt($('cfg-max_retry').value, 10),
+    log_level: $('cfg-log_level').value,
+  };
+  try {
+    const r = await fetch('/api/config', {
+      method: 'PUT',
+      headers: {'Content-Type': 'application/json'},
+      body: JSON.stringify(payload),
+    });
+    const d = await r.json().catch(() => ({}));
+    if (r.ok) {
+      toast(d.note || '已保存', 'ok');
+      setTimeout(loadConfig, 3000);
+    } else {
+      toast('保存失败: ' + (d.detail || r.status), 'err');
+    }
+  } catch (err) {
+    toast('保存失败: ' + err.message, 'err');
+  } finally {
+    $('save-btn').disabled = false;
+  }
+});
+
+// ---------- 状态 ----------
+function startStatusPolling() {
+  fetchStatus();
+  statusTimer = setInterval(fetchStatus, 3000);
+}
+
+async function fetchStatus() {
+  try {
+    const r = await fetch('/api/status');
+    if (!r.ok) return;
+    const s = await r.json();
+    const grid = $('status-grid');
+    const alive = s.bot_alive;
+    const wsClass = s.ws_connected ? 'ok' : (alive ? 'warn' : 'err');
+    const wsText = s.ws_connected ? '已连接' : (alive ? '连接中' : '未运行');
+    const upFmt = s.uptime ? formatUptime(s.uptime) : '-';
+    grid.innerHTML = `
+      <div class="status-card">
+        <div class="label">Bot 进程</div>
+        <div class="value ${alive ? 'ok' : 'err'}">${alive ? '运行中' : '已停止'}</div>
+      </div>
+      <div class="status-card">
+        <div class="label">WebSocket</div>
+        <div class="value ${wsClass}">${wsText}</div>
+      </div>
+      <div class="status-card">
+        <div class="label">运行时长</div>
+        <div class="value">${upFmt}</div>
+      </div>
+      <div class="status-card">
+        <div class="label">队列积压</div>
+        <div class="value ${s.queue_size > 50 ? 'warn' : ''}">${s.queue_size}</div>
+      </div>
+      <div class="status-card">
+        <div class="label">去重缓存</div>
+        <div class="value">${s.dedup_size}<span style="font-size:14px;color:var(--text-dim)"> / ${s.dedup_max}</span></div>
+      </div>
+      <div class="status-card">
+        <div class="label">状态新鲜度</div>
+        <div class="value ${s.status_fresh ? 'ok' : 'warn'}">${s.status_fresh ? '正常' : '失联'}</div>
+      </div>
+    `;
+  } catch {}
+}
+
+function formatUptime(sec) {
+  const d = Math.floor(sec / 86400);
+  const h = Math.floor((sec % 86400) / 3600);
+  const m = Math.floor((sec % 3600) / 60);
+  if (d > 0) return `${d}d ${h}h`;
+  if (h > 0) return `${h}h ${m}m`;
+  return `${m}m`;
+}
+
+// ---------- 日志 SSE ----------
+function startLogStream() {
+  if (sse) sse.close();
+  // 先拉历史
+  fetch('/api/logs/history').then(r => r.json()).then(entries => {
+    entries.forEach(appendLog);
+    // 再建立 SSE
+    const lastSeq = entries.length ? entries[entries.length - 1].seq : 0;
+    sse = new EventSource('/api/logs?since=' + lastSeq);
+    sse.onmessage = (ev) => {
+      try {
+        const entry = JSON.parse(ev.data);
+        appendLog(entry);
+      } catch {}
+    };
+    sse.onerror = () => {
+      // 浏览器会自动重连
+    };
+  });
+}
+
+function appendLog(entry) {
+  const container = $('log-container');
+  const line = document.createElement('div');
+  line.className = 'log-line';
+  line.dataset.level = entry.level;
+  const ts = new Date(entry.ts * 1000).toLocaleTimeString('zh-CN', {hour12: false});
+  line.innerHTML = `<span class="level-${entry.level}">[${ts}] [${entry.level}]</span> ${escapeHtml(entry.text)}`;
+  container.appendChild(line);
+  // 限制 DOM 条数
+  while (container.children.length > 1000) container.removeChild(container.firstChild);
+  if ($('autoscroll').checked) container.scrollTop = container.scrollHeight;
+}
+
+function clearLogs() {
+  $('log-container').innerHTML = '';
+}
+
+function escapeHtml(s) {
+  return s.replace(/[&<>"']/g, c => ({'&':'&amp;','<':'&lt;','>':'&gt;','"':'&quot;',"'":'&#39;'}[c]));
+}
+
+// 日志级别过滤
+['DEBUG','INFO','WARNING','ERROR'].forEach(lvl => {
+  $('lvl-' + lvl).addEventListener('change', applyLogFilter);
+});
+function applyLogFilter() {
+  const enabled = {};
+  ['DEBUG','INFO','WARNING','ERROR'].forEach(lvl => {
+    enabled[lvl] = $('lvl-' + lvl).checked;
+  });
+  document.querySelectorAll('.log-line').forEach(line => {
+    const lvl = line.dataset.level;
+    // CRITICAL 归到 ERROR
+    const key = lvl === 'CRITICAL' ? 'ERROR' : lvl;
+    line.style.display = enabled[key] === false ? 'none' : '';
+  });
+}
+
+// ---------- 操作 ----------
+async function restartBot() {
+  if (!confirm('确认重启 Bot?')) return;
+  try {
+    const r = await fetch('/api/bot/restart', {method: 'POST'});
+    const d = await r.json();
+    toast(d.note || '已重启', 'ok');
+  } catch (err) {
+    toast('重启失败: ' + err.message, 'err');
+  }
+}
+
+function toast(msg, type) {
+  const t = $('toast');
+  t.textContent = msg;
+  t.className = 'toast show ' + (type || '');
+  setTimeout(() => { t.className = 'toast ' + (type || ''); }, 3000);
+}
+
+// 启动
+checkAuth();
+</script>
+</body>
+</html>

+ 460 - 0
web/supervisor.py

@@ -0,0 +1,460 @@
+"""Web 面板后端:FastAPI + Bot 子进程管理 + 日志缓冲 + SSE。
+
+被 main.py 通过 uvicorn 启动。
+"""
+from __future__ import annotations
+
+import asyncio
+import json
+import logging
+import os
+import re
+import signal
+import subprocess
+import sys
+import threading
+import time
+from collections import deque
+from dataclasses import dataclass, field
+from pathlib import Path
+from typing import AsyncGenerator, Optional
+
+from fastapi import FastAPI, Request, Response, HTTPException, Depends
+from fastapi.responses import (
+    HTMLResponse,
+    JSONResponse,
+    StreamingResponse,
+    RedirectResponse,
+)
+from fastapi.staticfiles import StaticFiles
+from pydantic import BaseModel
+import yaml
+
+# 确保能 import 项目根的 config
+sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
+from config import Config, load_config  # noqa: E402
+from web.auth import (  # noqa: E402
+    COOKIE_NAME,
+    COOKIE_MAX_AGE,
+    has_password,
+    sign_session,
+    verify_password,
+    verify_session,
+)
+
+logger = logging.getLogger("lark2lark.supervisor")
+
+PROJECT_ROOT = Path(__file__).resolve().parent.parent
+BOT_SCRIPT = PROJECT_ROOT / "bot.py"
+CONFIG_PATH = PROJECT_ROOT / "config.yaml"
+INDEX_HTML = Path(__file__).resolve().parent / "index.html"
+
+LOG_BUFFER_MAX = 500
+STATUS_PREFIX = "__STATUS__"
+BOT_RESTART_DELAY = 1.0  # 秒,重启间隔
+
+
+# ---------- 数据模型 ----------
+
+class LoginForm(BaseModel):
+    password: str
+
+
+class ConfigUpdate(BaseModel):
+    app_id: Optional[str] = None
+    app_secret: Optional[str] = None
+    source_chat_ids: Optional[list[str]] = None
+    target_chat_ids: Optional[list[str]] = None
+    dedup_cache_size: Optional[int] = None
+    max_qps: Optional[int] = None
+    max_retry: Optional[int] = None
+    log_level: Optional[str] = None
+
+
+@dataclass
+class LogEntry:
+    seq: int
+    ts: float
+    level: str
+    text: str
+
+
+class LogBuffer:
+    """线程安全环形缓冲 + SSE 订阅广播。"""
+
+    def __init__(self, max_size: int = LOG_BUFFER_MAX) -> None:
+        self._lock = threading.Lock()
+        self._entries: deque[LogEntry] = deque(maxlen=max_size)
+        self._seq = 0
+        self._cond = threading.Condition(self._lock)
+
+    def append(self, text: str, level: str = "INFO") -> None:
+        with self._cond:
+            self._seq += 1
+            entry = LogEntry(seq=self._seq, ts=time.time(), level=level, text=text.rstrip("\n"))
+            self._entries.append(entry)
+            self._cond.notify_all()
+
+    def since(self, seq: int) -> list[LogEntry]:
+        with self._lock:
+            return [e for e in self._entries if e.seq > seq]
+
+    def all(self) -> list[LogEntry]:
+        with self._lock:
+            return list(self._entries)
+
+    def wait_for_new(self, last_seq: int, timeout: float = 25.0) -> list[LogEntry]:
+        """阻塞等待新日志,超时返回空列表(SSE keepalive)。"""
+        with self._cond:
+            if not self._cond.wait_for(
+                lambda: any(e.seq > last_seq for e in self._entries), timeout=timeout
+            ):
+                return []
+            return [e for e in self._entries if e.seq > last_seq]
+
+
+# ---------- Bot 子进程管理 ----------
+
+class BotManager:
+    """管理 bot.py 子进程:启动、停止、重启、日志采集、状态采集。"""
+
+    def __init__(self, log_buffer: LogBuffer) -> None:
+        self._log_buffer = log_buffer
+        self._proc: Optional[subprocess.Popen] = None
+        self._lock = threading.Lock()
+        self._reader_thread: Optional[threading.Thread] = None
+        self._watch_thread: Optional[threading.Thread] = None
+        self._stop_requested = False
+        self._last_status: dict = {}
+        self._last_status_ts: float = 0.0
+        self._start_time: float = 0.0
+
+    def start(self) -> None:
+        with self._lock:
+            if self._proc and self._proc.poll() is None:
+                return
+            self._stop_requested = False
+            self._start_time = time.time()
+            env = os.environ.copy()
+            env["PYTHONUNBUFFERED"] = "1"
+            env["PYTHONIOENCODING"] = "utf-8"
+            try:
+                self._proc = subprocess.Popen(
+                    [sys.executable, str(BOT_SCRIPT), "--config", str(CONFIG_PATH)],
+                    stdout=subprocess.PIPE,
+                    stderr=subprocess.STDOUT,
+                    cwd=str(PROJECT_ROOT),
+                    env=env,
+                    encoding="utf-8",
+                    errors="replace",
+                    bufsize=1,  # 行缓冲
+                )
+            except Exception as e:
+                self._log_buffer.append(f"启动 bot 子进程失败: {e}", "ERROR")
+                raise
+            self._reader_thread = threading.Thread(
+                target=self._read_output, name="bot-reader", daemon=True
+            )
+            self._reader_thread.start()
+            self._watch_thread = threading.Thread(
+                target=self._watch, name="bot-watch", daemon=True
+            )
+            self._watch_thread.start()
+            self._log_buffer.append("Bot 子进程已启动", "INFO")
+
+    def stop(self) -> None:
+        with self._lock:
+            self._stop_requested = True
+            proc = self._proc
+        if proc and proc.poll() is None:
+            self._log_buffer.append("正在停止 Bot 子进程...", "INFO")
+            try:
+                proc.terminate()
+                try:
+                    proc.wait(timeout=5)
+                except subprocess.TimeoutExpired:
+                    proc.kill()
+                    proc.wait(timeout=3)
+            except Exception as e:
+                self._log_buffer.append(f"停止 Bot 异常: {e}", "WARNING")
+
+    def restart(self) -> None:
+        self._log_buffer.append("重启 Bot 子进程...", "INFO")
+        self.stop()
+        time.sleep(BOT_RESTART_DELAY)
+        self.start()
+
+    def is_alive(self) -> bool:
+        with self._lock:
+            return self._proc is not None and self._proc.poll() is None
+
+    def uptime(self) -> int:
+        if not self.is_alive():
+            return 0
+        return int(time.time() - self._start_time)
+
+    @property
+    def last_status(self) -> dict:
+        return dict(self._last_status) if self._last_status else {}
+
+    def _read_output(self) -> None:
+        """读取 bot stdout,识别 __STATUS__ 行,其余写入日志缓冲。"""
+        proc = self._proc
+        if proc is None or proc.stdout is None:
+            return
+        try:
+            for line in proc.stdout:
+                line = line.rstrip("\n")
+                if not line:
+                    continue
+                if line.startswith(STATUS_PREFIX):
+                    self._parse_status(line[len(STATUS_PREFIX):])
+                    continue
+                level = self._detect_level(line)
+                self._log_buffer.append(line, level)
+        except Exception as e:
+            self._log_buffer.append(f"读取 Bot 输出异常: {e}", "ERROR")
+
+    @staticmethod
+    def _detect_level(line: str) -> str:
+        m = re.search(r"\[(DEBUG|INFO|WARNING|ERROR|CRITICAL)\]", line)
+        return m.group(1) if m else "INFO"
+
+    def _parse_status(self, json_str: str) -> None:
+        try:
+            data = json.loads(json_str)
+            self._last_status = data
+            self._last_status_ts = time.time()
+        except json.JSONDecodeError:
+            pass
+
+    def _watch(self) -> None:
+        """监控子进程存活,异常退出自动重启。"""
+        proc = self._proc
+        if proc is None:
+            return
+        while True:
+            rc = proc.wait()
+            self._log_buffer.append(f"Bot 子进程退出,返回码={rc}", "WARNING")
+            if self._stop_requested:
+                break
+            self._log_buffer.append(f"{BOT_RESTART_DELAY}s 后自动重启 Bot...", "INFO")
+            time.sleep(BOT_RESTART_DELAY)
+            try:
+                self.start()
+            except Exception as e:
+                self._log_buffer.append(f"自动重启失败: {e},10s 后再试", "ERROR")
+                time.sleep(10)
+                continue
+            break  # 新进程已由新 watch 线程接管,本线程退出
+
+
+# ---------- FastAPI 应用 ----------
+
+app = FastAPI(title="lark2lark 面板", docs_url=None, redoc_url=None)
+
+log_buffer = LogBuffer()
+bot_mgr = BotManager(log_buffer)
+
+
+def get_panel_password() -> str:
+    """从当前配置读取面板密码。每次调用都重新加载,支持热更新。"""
+    try:
+        cfg = load_config()
+        return cfg.panel_password
+    except Exception:
+        return ""
+
+
+def require_auth(request: Request) -> None:
+    """FastAPI 依赖:校验登录 Cookie。"""
+    token = request.cookies.get(COOKIE_NAME)
+    if not verify_session(token or "", get_panel_password()):
+        raise HTTPException(status_code=401, detail="未登录")
+
+
+@app.on_event("startup")
+async def _startup() -> None:
+    logging.basicConfig(
+        level=logging.INFO,
+        format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
+        datefmt="%Y-%m-%d %H:%M:%S",
+    )
+    logging.getLogger("lark-oapi").setLevel(logging.WARNING)
+    logging.getLogger("uvicorn.access").setLevel(logging.WARNING)
+    log_buffer.append("Supervisor 启动", "INFO")
+    try:
+        bot_mgr.start()
+    except Exception as e:
+        log_buffer.append(f"启动 Bot 失败: {e}", "ERROR")
+
+
+@app.on_event("shutdown")
+async def _shutdown() -> None:
+    log_buffer.append("Supervisor 关闭中", "INFO")
+    bot_mgr.stop()
+
+
+# ---------- 路由 ----------
+
+@app.get("/", response_class=HTMLResponse)
+async def index() -> HTMLResponse:
+    if not INDEX_HTML.exists():
+        return HTMLResponse("<h1>index.html 缺失</h1>", status_code=500)
+    return HTMLResponse(INDEX_HTML.read_text(encoding="utf-8"))
+
+
+@app.post("/api/login")
+async def login(form: LoginForm, response: Response):
+    pwd = get_panel_password()
+    if not has_password(pwd):
+        raise HTTPException(status_code=503, detail="面板未设置密码(LARK_PANEL_PASSWORD 环境变量或 config.panel.password)")
+    if not verify_password(form.password, pwd):
+        raise HTTPException(status_code=401, detail="密码错误")
+    token = sign_session(pwd)
+    response.set_cookie(
+        COOKIE_NAME, token,
+        max_age=COOKIE_MAX_AGE,
+        httponly=True,
+        samesite="lax",
+    )
+    return {"ok": True}
+
+
+@app.post("/api/logout")
+async def logout(response: Response):
+    response.delete_cookie(COOKIE_NAME)
+    return {"ok": True}
+
+
+@app.get("/api/config")
+async def get_config(_: None = Depends(require_auth)):
+    """返回当前配置,app_secret 脱敏。"""
+    if not CONFIG_PATH.exists():
+        raise HTTPException(status_code=404, detail="config.yaml 不存在")
+    with CONFIG_PATH.open("r", encoding="utf-8") as f:
+        data = yaml.safe_load(f) or {}
+    # 脱敏
+    if data.get("app_secret"):
+        data["app_secret"] = "***"
+    return data
+
+
+@app.put("/api/config")
+async def update_config(payload: ConfigUpdate, _: None = Depends(require_auth)):
+    """更新 config.yaml 并重启 Bot。app_secret 为 '***' 时保留原值。"""
+    if not CONFIG_PATH.exists():
+        raise HTTPException(status_code=404, detail="config.yaml 不存在")
+    with CONFIG_PATH.open("r", encoding="utf-8") as f:
+        current = yaml.safe_load(f) or {}
+
+    updates = payload.model_dump(exclude_none=True)
+    for key, val in updates.items():
+        if key == "app_secret" and val == "***":
+            continue  # 保留原值
+        current[key] = val
+
+    # 校验:写前快速检查
+    try:
+        _validate_partial(current)
+    except ValueError as e:
+        raise HTTPException(status_code=400, detail=str(e))
+
+    with CONFIG_PATH.open("w", encoding="utf-8") as f:
+        yaml.safe_dump(current, f, allow_unicode=True, sort_keys=False)
+
+    log_buffer.append("配置已更新,重启 Bot 使其生效", "INFO")
+    # 异步重启,避免阻塞 HTTP 响应
+    threading.Thread(target=bot_mgr.restart, daemon=True).start()
+    return {"ok": True, "note": "配置已保存,Bot 正在重启"}
+
+
+def _validate_partial(data: dict) -> None:
+    """对待写入的配置做基本校验。"""
+    if "max_qps" in data:
+        qps = data["max_qps"]
+        if not isinstance(qps, int) or qps <= 0 or qps > 5:
+            raise ValueError(f"max_qps 必须在 (0, 5] 区间,当前 {qps}")
+    if "source_chat_ids" in data:
+        if not isinstance(data["source_chat_ids"], list) or not data["source_chat_ids"]:
+            raise ValueError("source_chat_ids 至少 1 个")
+    if "target_chat_ids" in data:
+        if not isinstance(data["target_chat_ids"], list) or not data["target_chat_ids"]:
+            raise ValueError("target_chat_ids 至少 1 个")
+
+
+@app.get("/api/status")
+async def get_status(_: None = Depends(require_auth)):
+    """返回 Bot 运行状态。"""
+    status = bot_mgr.last_status
+    # 状态超时判定:超过 15s 未更新视为失联
+    status_fresh = (time.time() - bot_mgr._last_status_ts) < 15 if bot_mgr._last_status_ts else False
+    return {
+        "bot_alive": bot_mgr.is_alive(),
+        "ws_connected": status.get("ws_connected", False) and status_fresh,
+        "uptime": status.get("uptime", 0) if status_fresh else 0,
+        "queue_size": status.get("queue_size", 0),
+        "dedup_size": status.get("dedup_size", 0),
+        "dedup_max": status.get("dedup_max", 0),
+        "status_fresh": status_fresh,
+        "last_status_ts": bot_mgr._last_status_ts,
+    }
+
+
+@app.get("/api/logs")
+async def get_logs(request: Request, since: int = 0, _: None = Depends(require_auth)):
+    """SSE 流:实时推送新日志。"""
+    async def event_stream() -> AsyncGenerator[bytes, None]:
+        last_seq = since
+        # 先发送历史
+        for entry in log_buffer.since(last_seq):
+            last_seq = entry.seq
+            yield _format_sse(entry)
+        # 再订阅新日志
+        while True:
+            if await request.is_disconnected():
+                break
+            new_entries = await asyncio.get_event_loop().run_in_executor(
+                None, log_buffer.wait_for_new, last_seq, 25.0
+            )
+            if not new_entries:
+                # keepalive
+                yield b": ping\n\n"
+                continue
+            for entry in new_entries:
+                last_seq = entry.seq
+                yield _format_sse(entry)
+
+    return StreamingResponse(
+        event_stream(),
+        media_type="text/event-stream",
+        headers={
+            "Cache-Control": "no-cache",
+            "X-Accel-Buffering": "no",  # nginx 不缓冲
+        },
+    )
+
+
+def _format_sse(entry: LogEntry) -> bytes:
+    data = json.dumps({
+        "seq": entry.seq,
+        "ts": entry.ts,
+        "level": entry.level,
+        "text": entry.text,
+    }, ensure_ascii=False)
+    return f"data: {data}\n\n".encode("utf-8")
+
+
+@app.get("/api/logs/history")
+async def get_logs_history(_: None = Depends(require_auth)):
+    """返回全部历史日志(一次性)。"""
+    return [
+        {"seq": e.seq, "ts": e.ts, "level": e.level, "text": e.text}
+        for e in log_buffer.all()
+    ]
+
+
+@app.post("/api/bot/restart")
+async def restart_bot(_: None = Depends(require_auth)):
+    threading.Thread(target=bot_mgr.restart, daemon=True).start()
+    return {"ok": True, "note": "Bot 正在重启"}