From 9f90e166619ba116afa57b9e5809c24605215387 Mon Sep 17 00:00:00 2001 From: v6ole Date: Sun, 10 May 2026 13:59:57 +0800 Subject: [PATCH] =?UTF-8?q?feat(wechat):=20=E6=96=B0=E5=A2=9E=E4=BC=81?= =?UTF-8?q?=E5=BE=AE=E8=8F=9C=E5=8D=95=E5=8A=9F=E8=83=BD=E5=B9=B6=E4=BF=AE?= =?UTF-8?q?=E5=A4=8D=E9=87=8D=E5=A4=8D=E8=A7=A6=E5=8F=91=E9=97=AE=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 新增「查询」菜单组:监控配置、系统状态、今日工作日、最新公告 - 新增「系统管理」子菜单:暂停/恢复定时任务处理器 - 修复立即爬取被企微重试机制触发多次的问题(60秒防重入锁) - docker-compose 补充 image 名称和 container_name Co-Authored-By: Claude Sonnet 4.6 --- app/wechat/handler.py | 150 ++++++++++++++++++++++++++++++++++++++ app/wechat/menu.py | 25 +++++++ docker/docker-compose.yml | 2 + 3 files changed, 177 insertions(+) diff --git a/app/wechat/handler.py b/app/wechat/handler.py index 20ab613..70df7f2 100644 --- a/app/wechat/handler.py +++ b/app/wechat/handler.py @@ -1,4 +1,6 @@ +import asyncio import logging +import time import xml.etree.ElementTree as ET from app.config import settings @@ -7,6 +9,10 @@ from app.wechat.client import WeChatClient logger = logging.getLogger(__name__) +# 防重入:记录最近一次触发爬取的时间戳,60秒内不重复执行 +_last_crawl_time: float = 0.0 +_crawl_lock = asyncio.Lock() + class WeChatMessageHandler: def __init__(self): @@ -65,6 +71,18 @@ class WeChatMessageHandler: return await self._handle_trigger_crawl(from_user) elif event_key == "sync_holidays": return await self._handle_sync_holidays(from_user) + elif event_key == "monitor_config": + return await self._handle_monitor_config(from_user) + elif event_key == "system_status": + return await self._handle_system_status(from_user) + elif event_key == "workday_status": + return await self._handle_workday_status(from_user) + elif event_key == "latest_announcements": + return await self._handle_latest_announcements(from_user) + elif event_key == "pause_scheduler": + return await self._handle_pause_scheduler(from_user) + elif event_key == "resume_scheduler": + return await self._handle_resume_scheduler(from_user) return None async def handle_text(self, content: str, from_user: str) -> str | None: @@ -97,8 +115,17 @@ class WeChatMessageHandler: await self.client.send_text("查询失败,请稍后再试", from_user) async def _handle_trigger_crawl(self, from_user: str) -> str | None: + global _last_crawl_time from app.api.deps import get_crawl_service + # 防重入:企业微信会对同一事件重试多次,60秒内只执行一次 + async with _crawl_lock: + now = time.monotonic() + if now - _last_crawl_time < 60: + logger.info(f"爬取请求被忽略(防重入),距上次 {now - _last_crawl_time:.1f}s") + return None + _last_crawl_time = now + await self.client.send_text("开始爬取,请稍候...", from_user) try: @@ -136,3 +163,126 @@ class WeChatMessageHandler: except Exception as e: logger.error(f"同步节假日失败: {e}") await self.client.send_text(f"同步失败: {e}", from_user) + + async def _handle_monitor_config(self, from_user: str) -> str | None: + import json + try: + keywords = settings.crawler_keywords + sources = json.loads(settings.announcement_sources) + source_names = "、".join(v["name"] for v in sources.values()) + text = ( + f"监控关键词: {', '.join(keywords)}\n" + f"爬取页数: {settings.crawler_max_pages} 页\n" + f"定时规则: {settings.scheduler_cron}\n" + f"公告来源: {source_names}" + ) + await self.client.send_text(text, from_user) + except Exception as e: + logger.error(f"查询监控配置失败: {e}") + await self.client.send_text("查询失败,请稍后再试", from_user) + + async def _handle_system_status(self, from_user: str) -> str | None: + from app.api.deps import get_db + from sqlalchemy import func, select + from app.models.announcement import Announcement + + try: + async for db in get_db(): + total_result = await db.execute( + select(func.count()).select_from(Announcement) + ) + total = total_result.scalar() or 0 + + today_result = await db.execute( + select(func.count()).where( + func.date(Announcement.publish_date) == func.current_date() + ).select_from(Announcement) + ) + today = today_result.scalar() or 0 + + unsent_result = await db.execute( + select(func.count()).where( + Announcement.is_sent == False, # noqa: E712 + Announcement.keyword_matched == True, # noqa: E712 + ).select_from(Announcement) + ) + unsent = unsent_result.scalar() or 0 + + scheduler_status = "已启用" if settings.scheduler_enabled else "已禁用" + text = ( + f"累计公告: {total} 条\n" + f"今日新增: {today} 条\n" + f"待推送: {unsent} 条\n" + f"定时任务: {scheduler_status}\n" + f"定时规则: {settings.scheduler_cron}" + ) + await self.client.send_text(text, from_user) + except Exception as e: + logger.error(f"查询系统状态失败: {e}") + await self.client.send_text("查询失败,请稍后再试", from_user) + + async def _handle_workday_status(self, from_user: str) -> str | None: + from app.services.holiday_service import now_in_china, is_workday + from app.api.deps import get_db + + try: + async for db in get_db(): + today = now_in_china() + workday = await is_workday(db, today) + status = "工作日,正常爬取" if workday else "非工作日,跳过爬取" + text = f"今天 {today.strftime('%Y-%m-%d %A')}\n{status}" + await self.client.send_text(text, from_user) + except Exception as e: + logger.error(f"查询工作日状态失败: {e}") + await self.client.send_text("查询失败,请稍后再试", from_user) + + async def _handle_latest_announcements(self, from_user: str) -> str | None: + from app.api.deps import get_db + from sqlalchemy import select, desc + from app.models.announcement import Announcement + + try: + async for db in get_db(): + stmt = ( + select(Announcement) + .where(Announcement.keyword_matched == True) # noqa: E712 + .order_by(desc(Announcement.publish_date)) + .limit(5) + ) + result = await db.execute(stmt) + items = result.scalars().all() + + if not items: + await self.client.send_text("暂无匹配关键词的公告", from_user) + return None + + lines = ["最新匹配公告(最近5条):"] + for i, ann in enumerate(items, 1): + date_str = ann.publish_date.strftime("%m-%d") if ann.publish_date else "?" + title = ann.title[:30] + "..." if len(ann.title) > 30 else ann.title + lines.append(f"{i}. [{date_str}] {title}") + await self.client.send_text("\n".join(lines), from_user) + except Exception as e: + logger.error(f"查询最新公告失败: {e}") + await self.client.send_text("查询失败,请稍后再试", from_user) + + async def _handle_pause_scheduler(self, from_user: str) -> str | None: + try: + from app.scheduler.jobs import scheduler + if scheduler.running: + scheduler.pause() + await self.client.send_text("定时任务已暂停", from_user) + else: + await self.client.send_text("定时任务未在运行", from_user) + except Exception as e: + logger.error(f"暂停定时任务失败: {e}") + await self.client.send_text(f"操作失败: {e}", from_user) + + async def _handle_resume_scheduler(self, from_user: str) -> str | None: + try: + from app.scheduler.jobs import scheduler + scheduler.resume() + await self.client.send_text("定时任务已恢复", from_user) + except Exception as e: + logger.error(f"恢复定时任务失败: {e}") + await self.client.send_text(f"操作失败: {e}", from_user) diff --git a/app/wechat/menu.py b/app/wechat/menu.py index dce8c95..643b7ce 100644 --- a/app/wechat/menu.py +++ b/app/wechat/menu.py @@ -13,6 +13,31 @@ MENU = { "type": "click", "key": "today_stats", }, + { + "name": "查询", + "sub_button": [ + { + "name": "监控配置", + "type": "click", + "key": "monitor_config", + }, + { + "name": "系统状态", + "type": "click", + "key": "system_status", + }, + { + "name": "今日工作日", + "type": "click", + "key": "workday_status", + }, + { + "name": "最新公告", + "type": "click", + "key": "latest_announcements", + }, + ], + }, { "name": "系统管理", "sub_button": [ diff --git a/docker/docker-compose.yml b/docker/docker-compose.yml index 2c09cc2..aa9f8c1 100644 --- a/docker/docker-compose.yml +++ b/docker/docker-compose.yml @@ -3,6 +3,8 @@ services: build: context: .. dockerfile: docker/Dockerfile + image: gx-gp-notify:latest + container_name: gx-gp-notify ports: - "18001:8000" env_file: